mirror of
https://github.com/hansjone/dsh-im-ops.git
synced 2026-10-08 23:13:17 +08:00
Merge pull request #122 and harden reply stall detection
Resolve the current main conflict, treat reply timeouts as inactivity windows, and confirm liveness through Harness instead of sticky interaction ownership.
This commit is contained in:
commit
34f217cafc
4 changed files with 286 additions and 53 deletions
|
|
@ -6,6 +6,11 @@ This file records the notable changes in each dsh-im release. Its format follows
|
|||
|
||||
## [Unreleased]
|
||||
|
||||
### Fixed / 修复
|
||||
|
||||
- Harness 回复等待改为按活动续期:持续产生事件或 Harness 明确报告 Session 仍在运行的长任务不再被固定 10 分钟上限误报超时;已开始但停止推进且不再运行的任务仍会按停滞窗口超时。
|
||||
Harness reply waits now renew from activity: long-running turns that keep producing events or are still reported as running no longer hit a fixed ten-minute timeout, while started turns that stop progressing and are no longer running still time out after the stall window.
|
||||
|
||||
## [4.8.0] - 2026-09-02
|
||||
|
||||
### Added / 新增
|
||||
|
|
|
|||
74
lib/index.js
74
lib/index.js
File diff suppressed because one or more lines are too long
|
|
@ -434,6 +434,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();
|
||||
}
|
||||
|
|
@ -1435,8 +1440,14 @@ export class HarnessClient {
|
|||
promptAccepted = true;
|
||||
|
||||
try {
|
||||
const deadline = Date.now() + timeoutMs;
|
||||
while (Date.now() < deadline) {
|
||||
// Treat timeoutMs as a stall window rather than a hard runtime limit.
|
||||
// Durable events are direct progress. Once a full quiet window elapses,
|
||||
// confirm the Session is still running before renewing the wait.
|
||||
// Interaction ownership is intentionally not a liveness signal: it stays
|
||||
// active until turn/end and can therefore outlive a stalled turn.
|
||||
let lastProgressAt = Date.now();
|
||||
let lastPollSeq = tracker.lastSeq;
|
||||
while (true) {
|
||||
await sleep(300, signal);
|
||||
const history = await this.rpc(
|
||||
'session.history',
|
||||
|
|
@ -1450,6 +1461,9 @@ export class HarnessClient {
|
|||
if (!wasActive && ownership.active) ownership.reconnect?.();
|
||||
}
|
||||
const updates = tracker.consumeAll(history.events ?? []);
|
||||
const seqAdvanced = tracker.lastSeq > lastPollSeq;
|
||||
lastPollSeq = tracker.lastSeq;
|
||||
if (seqAdvanced) lastProgressAt = Date.now();
|
||||
if (onUpdate) {
|
||||
const visibleUpdates = progressMode === 'all' ? updates : updates.slice(-1);
|
||||
for (const update of visibleUpdates) {
|
||||
|
|
@ -1460,24 +1474,39 @@ export class HarnessClient {
|
|||
}
|
||||
}
|
||||
}
|
||||
if (!tracker.finished) continue;
|
||||
turnFinished = true;
|
||||
if (!ownership?.stopRequested && !harnessTurnSucceeded(tracker.reason)) {
|
||||
if (tracker.finished) {
|
||||
turnFinished = true;
|
||||
if (!ownership?.stopRequested && !harnessTurnSucceeded(tracker.reason)) {
|
||||
throw harnessTurnError(tracker.reason);
|
||||
}
|
||||
// An accepted /stop revokes attachment delivery even when Harness
|
||||
// preserved a useful partial text answer for the existing UX.
|
||||
const artifactCount = ownership?.stopRequested
|
||||
? 0
|
||||
: await deliverArtifacts();
|
||||
if (tracker.answer) {
|
||||
return tracker.answer;
|
||||
}
|
||||
if (artifactCount > 0) return '';
|
||||
if (ownership?.stopRequested) throw turnStoppedError();
|
||||
throw harnessTurnError(tracker.reason);
|
||||
}
|
||||
// An accepted /stop revokes attachment delivery even when Harness
|
||||
// preserved a useful partial text answer for the existing UX.
|
||||
const artifactCount = ownership?.stopRequested
|
||||
? 0
|
||||
: await deliverArtifacts();
|
||||
if (tracker.answer) {
|
||||
return tracker.answer;
|
||||
|
||||
if (Date.now() - lastProgressAt < timeoutMs) continue;
|
||||
|
||||
let running = false;
|
||||
try {
|
||||
running = await this.isSessionRunning(sessionId, { signal });
|
||||
} catch (error) {
|
||||
if (signal?.aborted) throw signal.reason ?? error;
|
||||
// A failed liveness probe is not evidence of progress.
|
||||
}
|
||||
if (artifactCount > 0) return '';
|
||||
if (ownership?.stopRequested) throw turnStoppedError();
|
||||
throw harnessTurnError(tracker.reason);
|
||||
if (running) {
|
||||
lastProgressAt = Date.now();
|
||||
continue;
|
||||
}
|
||||
throw new HarnessTurnError('harness-reply-timeout');
|
||||
}
|
||||
throw new HarnessTurnError('harness-reply-timeout');
|
||||
} catch (error) {
|
||||
// Once cancellation was accepted, transport/poll failures and timeouts
|
||||
// describe the convergence of that stop, not an unrelated ask failure.
|
||||
|
|
|
|||
|
|
@ -974,3 +974,202 @@ test('HarnessReplyTracker keeps every frame of a batched turn in order', () => {
|
|||
]);
|
||||
assert.equal(tracker.answer, '先创建再观察:');
|
||||
});
|
||||
|
||||
test('new turn events renew the stall window beyond the original fixed deadline', 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;
|
||||
let historyPolls = 0;
|
||||
let sessionListPolls = 0;
|
||||
const events = [];
|
||||
client.rpc = async (method, _payload, _timeoutMs, options) => {
|
||||
if (method === 'session.history') {
|
||||
if (!prompted) return { events: [] };
|
||||
historyPolls += 1;
|
||||
if (historyPolls === 1) {
|
||||
events.push(
|
||||
{ type: 'turn/start', seq: ++seq, data: { turn: 1 } },
|
||||
{ type: 'user/message', seq: ++seq, data: { turn: 1, source: { rpcId: promptRpcId } } },
|
||||
{ type: 'assistant/chunk', seq: ++seq, data: { turn: 1, step: 1, chunk: { type: 'text-delta', index: 0, text: '第一段' } } },
|
||||
);
|
||||
} else if (historyPolls === 2) {
|
||||
events.push(
|
||||
{ type: 'assistant/chunk', seq: ++seq, data: { turn: 1, step: 1, chunk: { type: 'text-delta', index: 0, text: '第二段' } } },
|
||||
);
|
||||
} else if (historyPolls === 3) {
|
||||
events.push(
|
||||
{ type: 'assistant/message', seq: ++seq, data: { turn: 1, step: 1, message: { content: [{ type: 'text', text: '最终结果' }] } } },
|
||||
{ type: 'turn/end', seq: ++seq, data: { turn: 1, reason: { kind: 'completed' } } },
|
||||
);
|
||||
}
|
||||
return { events: events.map((event) => ({ event })) };
|
||||
}
|
||||
if (method === 'session.prompt') {
|
||||
prompted = true;
|
||||
promptRpcId = options.rpcId;
|
||||
return {};
|
||||
}
|
||||
if (method === 'session.list') {
|
||||
sessionListPolls += 1;
|
||||
return { items: [{ sessionId: 'session-renewal-events', running: false }] };
|
||||
}
|
||||
throw new Error(`unexpected rpc ${method}`);
|
||||
};
|
||||
|
||||
// The third poll lands around 900ms. The old 450ms fixed deadline exits after
|
||||
// the second poll; the stall window reaches the third poll because seq advances.
|
||||
const answer = await client.ask('session-renewal-events', 'long task', {
|
||||
timeoutMs: 450,
|
||||
control: { owner: {}, key: 'route' },
|
||||
});
|
||||
assert.equal(answer, '最终结果');
|
||||
assert.equal(historyPolls, 3);
|
||||
assert.equal(sessionListPolls, 0);
|
||||
});
|
||||
|
||||
test('a production-owned turn that starts and then stalls 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;
|
||||
let historyPolls = 0;
|
||||
let sessionListPolls = 0;
|
||||
let events = [];
|
||||
client.rpc = async (method, _payload, _timeoutMs, options) => {
|
||||
if (method === 'session.history') {
|
||||
if (!prompted) return { events: [] };
|
||||
historyPolls += 1;
|
||||
if (historyPolls === 1) {
|
||||
events = [
|
||||
{ event: { type: 'turn/start', seq: 1, data: { turn: 1 } } },
|
||||
{ event: { type: 'user/message', seq: 2, data: { turn: 1, source: { rpcId: promptRpcId } } } },
|
||||
];
|
||||
}
|
||||
return { events };
|
||||
}
|
||||
if (method === 'session.prompt') {
|
||||
prompted = true;
|
||||
promptRpcId = options.rpcId;
|
||||
return {};
|
||||
}
|
||||
if (method === 'session.list') {
|
||||
sessionListPolls += 1;
|
||||
return { items: [{ sessionId: 'session-stalled-owned', running: false }] };
|
||||
}
|
||||
throw new Error(`unexpected rpc ${method}`);
|
||||
};
|
||||
|
||||
await assert.rejects(
|
||||
client.ask('session-stalled-owned', 'stall after starting', {
|
||||
timeoutMs: 120,
|
||||
control: { owner: {}, key: 'route' },
|
||||
}),
|
||||
(error) => error?.code === 'harness-reply-timeout',
|
||||
'stale control ownership must not keep a stalled turn alive',
|
||||
);
|
||||
assert.equal(historyPolls, 2);
|
||||
assert.equal(sessionListPolls, 1);
|
||||
});
|
||||
|
||||
test('a silent turn renews only after Harness confirms the Session is running', 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 historyPolls = 0;
|
||||
let sessionListPolls = 0;
|
||||
const events = [];
|
||||
client.rpc = async (method, _payload, _timeoutMs, options) => {
|
||||
if (method === 'session.history') {
|
||||
if (!prompted) return { events: [] };
|
||||
historyPolls += 1;
|
||||
if (historyPolls === 1) {
|
||||
events.push(
|
||||
{ event: { type: 'turn/start', seq: 1, data: { turn: 1 } } },
|
||||
{ event: { type: 'user/message', seq: 2, data: { turn: 1, source: { rpcId: promptRpcId } } } },
|
||||
);
|
||||
} else if (historyPolls === 3) {
|
||||
events.push(
|
||||
{ event: { type: 'assistant/message', seq: 3, data: { turn: 1, step: 1, message: { content: [{ type: 'text', text: '静默任务完成' }] } } } },
|
||||
{ event: { type: 'turn/end', seq: 4, data: { turn: 1, reason: { kind: 'completed' } } } },
|
||||
);
|
||||
}
|
||||
return { events };
|
||||
}
|
||||
if (method === 'session.prompt') {
|
||||
prompted = true;
|
||||
promptRpcId = options.rpcId;
|
||||
return {};
|
||||
}
|
||||
if (method === 'session.list') {
|
||||
sessionListPolls += 1;
|
||||
return { items: [{ sessionId: 'session-silent-running', running: true }] };
|
||||
}
|
||||
throw new Error(`unexpected rpc ${method}`);
|
||||
};
|
||||
|
||||
const answer = await client.ask('session-silent-running', 'silent long task', {
|
||||
timeoutMs: 120,
|
||||
control: { owner: {}, key: 'route' },
|
||||
});
|
||||
assert.equal(answer, '静默任务完成');
|
||||
assert.equal(historyPolls, 3);
|
||||
assert.equal(sessionListPolls, 1);
|
||||
});
|
||||
|
||||
test('a failed liveness probe does not renew a stalled turn', async () => {
|
||||
const client = new HarnessClient({
|
||||
baseUrl: 'http://127.0.0.1:3080',
|
||||
workspace: '/tmp/dsh-feishu-workspace',
|
||||
});
|
||||
client.ensureRunning = async () => undefined;
|
||||
let prompted = false;
|
||||
let promptRpcId;
|
||||
let historyPolls = 0;
|
||||
let sessionListPolls = 0;
|
||||
let events = [];
|
||||
client.rpc = async (method, _payload, _timeoutMs, options) => {
|
||||
if (method === 'session.history') {
|
||||
if (!prompted) return { events: [] };
|
||||
historyPolls += 1;
|
||||
if (historyPolls === 1) {
|
||||
events = [
|
||||
{ event: { type: 'turn/start', seq: 1, data: { turn: 1 } } },
|
||||
{ event: { type: 'user/message', seq: 2, data: { turn: 1, source: { rpcId: promptRpcId } } } },
|
||||
];
|
||||
}
|
||||
return { events };
|
||||
}
|
||||
if (method === 'session.prompt') {
|
||||
prompted = true;
|
||||
promptRpcId = options.rpcId;
|
||||
return {};
|
||||
}
|
||||
if (method === 'session.list') {
|
||||
sessionListPolls += 1;
|
||||
throw new Error('probe unavailable');
|
||||
}
|
||||
throw new Error(`unexpected rpc ${method}`);
|
||||
};
|
||||
|
||||
await assert.rejects(
|
||||
client.ask('session-probe-fails', 'stall after starting', {
|
||||
timeoutMs: 120,
|
||||
control: { owner: {}, key: 'route' },
|
||||
}),
|
||||
(error) => error?.code === 'harness-reply-timeout',
|
||||
);
|
||||
assert.equal(historyPolls, 2);
|
||||
assert.equal(sessionListPolls, 1);
|
||||
});
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue