mirror of
https://github.com/hansjone/dsh-im-ops.git
synced 2026-10-09 01:53:21 +08:00
fix(feishu): renew the reply wait window while the turn is active
A long-running Harness task was falsely reported as MODEL_REPLY_TIMEOUT after the fixed 10-minute deadline even though the backend was still working. Replace the hard wall-clock deadline with an activity-based stall window: the wait renews whenever the turn produces new events, the control ownership is active, or (mid-turn with no new events) Harness still reports the session as running. Only a genuine stall with no progress for the whole timeoutMs window fails fast. This is the shared harness-client layer, exercised end-to-end by the Feishu channel (which already exposes HARNESS_REPLY_TIMEOUT_MS); other channels inherit the same fix. Adds two regression tests: a long-running turn with periodic activity completes instead of timing out, and a genuinely stalled turn still fails with harness-reply-timeout.
This commit is contained in:
parent
c2be2389b0
commit
955c2217fd
3 changed files with 132 additions and 18 deletions
32
lib/index.js
32
lib/index.js
File diff suppressed because one or more lines are too long
|
|
@ -428,6 +428,11 @@ export class HarnessReplyTracker {
|
|||
return this.#finished;
|
||||
}
|
||||
|
||||
/** The highest event seq consumed so far; advances as the turn produces events. */
|
||||
get lastSeq() {
|
||||
return this.#lastSeq;
|
||||
}
|
||||
|
||||
get answer() {
|
||||
return this.#latestText.trim();
|
||||
}
|
||||
|
|
@ -1387,8 +1392,13 @@ export class HarnessClient {
|
|||
promptAccepted = true;
|
||||
|
||||
try {
|
||||
const deadline = Date.now() + timeoutMs;
|
||||
while (Date.now() < deadline) {
|
||||
// Renew the wait window whenever the turn is still making progress, so
|
||||
// a long-running task is not falsely reported as MODEL_REPLY_TIMEOUT
|
||||
// while Harness is still working. A genuine stall (no new events and
|
||||
// no active turn for the whole timeoutMs window) still fails fast.
|
||||
let lastProgressAt = Date.now();
|
||||
let lastPollSeq = tracker.lastSeq;
|
||||
while (Date.now() - lastProgressAt < timeoutMs) {
|
||||
await sleep(300, signal);
|
||||
const history = await this.rpc(
|
||||
'session.history',
|
||||
|
|
@ -1402,6 +1412,26 @@ export class HarnessClient {
|
|||
if (!wasActive && ownership.active) ownership.reconnect?.();
|
||||
}
|
||||
const updates = tracker.consumeAll(history.events ?? []);
|
||||
// Progress = the turn produced new events, or the ownership is still
|
||||
// active. Either refreshes the stall window.
|
||||
const seqAdvanced = tracker.lastSeq > lastPollSeq;
|
||||
lastPollSeq = tracker.lastSeq;
|
||||
if (seqAdvanced) {
|
||||
lastProgressAt = Date.now();
|
||||
} else if (ownership?.active) {
|
||||
lastProgressAt = Date.now();
|
||||
} else if (tracker.tracking) {
|
||||
// Mid-turn but no new events: ask Harness whether the session is
|
||||
// still running before declaring a stall. A failing probe must not
|
||||
// mask a real stall, so it is not treated as progress.
|
||||
try {
|
||||
if (await this.isSessionRunning(sessionId, { signal })) {
|
||||
lastProgressAt = Date.now();
|
||||
}
|
||||
} catch {
|
||||
// ignore probe failure; the deadline still applies
|
||||
}
|
||||
}
|
||||
if (onUpdate) {
|
||||
const visibleUpdates = progressMode === 'all' ? updates : updates.slice(-1);
|
||||
for (const update of visibleUpdates) {
|
||||
|
|
|
|||
|
|
@ -974,3 +974,87 @@ test('HarnessReplyTracker keeps every frame of a batched turn in order', () => {
|
|||
]);
|
||||
assert.equal(tracker.answer, '先创建再观察:');
|
||||
});
|
||||
|
||||
test('a turn that keeps producing events is not falsely reported as MODEL_REPLY_TIMEOUT', async () => {
|
||||
const client = new HarnessClient({
|
||||
baseUrl: 'http://127.0.0.1:3080',
|
||||
workspace: '/tmp/dsh-feishu-workspace',
|
||||
});
|
||||
client.ensureRunning = async () => undefined;
|
||||
let promptRpcId;
|
||||
let prompted = false;
|
||||
let seq = 0;
|
||||
const startedAt = Date.now();
|
||||
client.rpc = async (method, _payload, _timeoutMs, options) => {
|
||||
if (method === 'session.history') {
|
||||
if (!prompted) return { events: [] };
|
||||
const events = [];
|
||||
if (!promptRpcId) promptRpcId = options?.rpcId;
|
||||
// The turn "works" for ~250ms: each poll appends a fresh chunk so the
|
||||
// seq keeps advancing (the renewal signal). After the window it ends.
|
||||
const elapsed = Date.now() - startedAt;
|
||||
if (elapsed < 250) {
|
||||
events.push(
|
||||
{ event: { type: 'turn/start', seq: ++seq, data: { turn: 1 } } },
|
||||
{ event: { type: 'user/message', seq: ++seq, data: { turn: 1, source: { rpcId: promptRpcId } } } },
|
||||
{ event: { type: 'assistant/chunk', seq: ++seq, data: { turn: 1, step: 1, chunk: { type: 'text-delta', index: 0, text: `chunk-${seq}` } } } },
|
||||
);
|
||||
} else {
|
||||
events.push(
|
||||
{ event: { type: 'turn/start', seq: ++seq, data: { turn: 1 } } },
|
||||
{ event: { type: 'user/message', seq: ++seq, data: { turn: 1, source: { rpcId: promptRpcId } } } },
|
||||
{ event: { type: 'assistant/message', seq: ++seq, data: { turn: 1, step: 1, message: { content: [{ type: 'text', text: '最终结果' }] } } } },
|
||||
{ event: { type: 'turn/end', seq: ++seq, data: { turn: 1, reason: { kind: 'completed' } } } },
|
||||
);
|
||||
}
|
||||
return { events };
|
||||
}
|
||||
if (method === 'session.prompt') {
|
||||
prompted = true;
|
||||
promptRpcId = options.rpcId;
|
||||
return {};
|
||||
}
|
||||
if (method === 'session.list') {
|
||||
return { items: [{ sessionId: 'session-renewal', running: true }] };
|
||||
}
|
||||
throw new Error(`unexpected rpc ${method}`);
|
||||
};
|
||||
|
||||
// timeoutMs is tiny (100ms) but the turn keeps producing events, so the
|
||||
// stall window keeps renewing and the ask must complete instead of timing out.
|
||||
const answer = await client.ask('session-renewal', 'long task', { timeoutMs: 100 });
|
||||
assert.equal(answer, '最终结果');
|
||||
});
|
||||
|
||||
test('a turn with no activity and no running session still times out', async () => {
|
||||
const client = new HarnessClient({
|
||||
baseUrl: 'http://127.0.0.1:3080',
|
||||
workspace: '/tmp/dsh-feishu-workspace',
|
||||
});
|
||||
client.ensureRunning = async () => undefined;
|
||||
let promptRpcId;
|
||||
let prompted = false;
|
||||
client.rpc = async (method, _payload, _timeoutMs, options) => {
|
||||
if (method === 'session.history') {
|
||||
if (!prompted) return { events: [] };
|
||||
// The turn started but never emits another event, and the session is not
|
||||
// reported as running: a genuine stall.
|
||||
return { events: [] };
|
||||
}
|
||||
if (method === 'session.prompt') {
|
||||
prompted = true;
|
||||
promptRpcId = options.rpcId;
|
||||
return {};
|
||||
}
|
||||
if (method === 'session.list') {
|
||||
return { items: [{ sessionId: 'session-stall', running: false }] };
|
||||
}
|
||||
throw new Error(`unexpected rpc ${method}`);
|
||||
};
|
||||
|
||||
await assert.rejects(
|
||||
client.ask('session-stall', 'never runs', { timeoutMs: 120 }),
|
||||
(error) => error?.code === 'harness-reply-timeout',
|
||||
'a genuinely stalled turn must still fail with MODEL_REPLY_TIMEOUT',
|
||||
);
|
||||
});
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue