diff --git a/changelog.d/membrane-floor-0-5-81.changed.md b/changelog.d/membrane-floor-0-5-81.changed.md new file mode 100644 index 00000000..ed6f9e30 --- /dev/null +++ b/changelog.d/membrane-floor-0-5-81.changed.md @@ -0,0 +1,5 @@ +- **Hosts:** the `@animalabs/membrane` floor is now `^0.5.81`. Earlier + releases reported reason `user` from a broad abort catch for any error + whose message contained "abort", which the new cancellation path would + have recorded as a deliberate stop; from 0.5.81 `user` means exactly that + the request's signal was aborted. diff --git a/changelog.d/shutdown-cancel-provenance.fixed.md b/changelog.d/shutdown-cancel-provenance.fixed.md new file mode 100644 index 00000000..ca684d28 --- /dev/null +++ b/changelog.d/shutdown-cancel-provenance.fixed.md @@ -0,0 +1,9 @@ +- A graceful framework shutdown (`AgentFramework.stop()`) with an inference + still streaming records its own provenance before cancelling — the same + `frameworkCancelledStreams` track that `endTurn`, budget restarts and + quiesce use — so the stream driver settles the turn as a shutdown: no + `[turn-interrupted]` marker (nobody stopped the agent; the host went + away), no `inference:exhausted`, one `inference:aborted` with reason + `shutdown`, an inference-log terminal, and the agent settled so an + in-flight `runEphemeralToCompletion` rejects immediately instead of + waiting out its idle watchdog. diff --git a/changelog.d/user-interrupt-not-failure.fixed.md b/changelog.d/user-interrupt-not-failure.fixed.md new file mode 100644 index 00000000..9381db13 --- /dev/null +++ b/changelog.d/user-interrupt-not-failure.fixed.md @@ -0,0 +1,28 @@ +- A cancelled stream (host Stop button, `agent.cancelStream()`, + `framework.abortInference()`) is no longer recorded as an inference + failure. It emits `inference:aborted` instead of `inference:exhausted`, + leaves the consecutive-failure streak and ops alerts untouched, and writes + a `[turn-interrupted] … a deliberate stop, not a failure` chronicle marker + in place of the `[inference-failed] the model call failed…` text with + remediation advice for a failure that never happened. The marker names the + act, not an actor (membrane's `user` reason means "the signal was aborted", + not "the user did it"), and makes no claim about delivery. Callers can pass + their own provenance — `cancelStream(reason)` / `abortInference(reason)` — + which the trace and the marker's metadata carry; `abortInference` no longer + emits a second `inference:aborted` on top of the stream driver's. This + holds on both of a stream's cancel twins: a stream implementation that + reports `cancel()` through `error` rather than `aborted` reaches the same + terminal (one `inference:aborted`, the marker, no `inference:failed`, no + `errorPolicy` retry of the inference that was just stopped), and the + quiesce/shutdown twins no longer emit a contradictory `inference:failed` + before settling as aborted. A third shape — an implementation whose + `cancel()` simply closes the iterator, with no terminal event — reaches the + same terminal too, at the loop's end, instead of leaving the turn unsettled + under a `completed` lifecycle terminal. (#134, by Lari; reworked after + review.) +- Speech-route failures with no delivery locus (headless/WebUI turns with no + home or trigger channel) now read `[send-undeliverable] … had no channel + to go to` instead of claiming a Discord delivery failure to "the channel", + and say that this route delivered nowhere rather than that the reply + reached no one — another `dispatchSpeech` handler may have shown it. The + machine-readable marker `kind` is unchanged. diff --git a/package.json b/package.json index 9293665f..cad52453 100644 --- a/package.json +++ b/package.json @@ -48,7 +48,7 @@ "dependencies": { "@animalabs/chronicle": "^0.4.0", "@animalabs/context-manager": "^0.10.0", - "@animalabs/membrane": "^0.5.78", + "@animalabs/membrane": "^0.5.81", "chokidar": "^4.0.3", "discord.js": "^14.25.1", "ws": "^8.18.0" diff --git a/src/agent.ts b/src/agent.ts index 0855aafd..7a7184e0 100644 --- a/src/agent.ts +++ b/src/agent.ts @@ -818,6 +818,13 @@ export class Agent { this.failOpenKvSubmissions(); this._streamId++; + // A pending cancel belongs to the stream that was live when cancelStream() + // ran. The stream driver collects it on that stream's terminal event, or + // at the loop's end when the iterator closed without one; this clear is + // the belt for a driver that never got there (a stream that never + // started iterating). A fresh stream must not inherit it and read its + // own later error as a deliberate stop. + this._pendingCancel = undefined; this._inferenceStartedAt = Date.now(); this.lastStreamInputTokens = 0; // Reset the cache-inclusive counters too: a new stream must not inherit @@ -958,10 +965,24 @@ export class Agent { this._state = { status: 'streaming', stream }; } + /** The most recent cancelStream() of a live stream, until the stream + * driver collects it on that stream's terminal event. Its presence is + * the signal (a reasonless Stop is still a deliberate stop); its reason + * is the caller's own word. Membrane reports every cancel() as reason + * 'user' — the call, not the actor — and a stream implementation may + * report the cancel as `error` with no reason at all, so this record is + * the only provenance a host-side cancel has. */ + private _pendingCancel: { reason?: string } | undefined; + /** - * Cancel any active stream and reset to idle. + * Cancel any active stream and reset to idle. `reason` is provenance for + * the framework's inference:aborted trace and marker metadata (e.g. + * 'zombie_reclaim', 'subagent_cancel'); it never reaches the provider. */ - cancelStream(): void { + cancelStream(reason?: string): void { + const hadStream = this._state.status === 'streaming' || + (this._state.status === 'waiting_for_tools' && !!this._state.stream); + if (hadStream) this._pendingCancel = { reason }; if (this._state.status === 'streaming') { this._state.stream.cancel(); } else if (this._state.status === 'waiting_for_tools' && this._state.stream) { @@ -970,6 +991,17 @@ export class Agent { this._state = { status: 'idle' }; } + /** Collect (and clear) the cancel that ended the current stream, if this + * side issued one. Called once by the stream driver on the stream's + * terminal event — `aborted` normally, `error` for implementations that + * report cancel() that way — so it is consumed exactly once, by the + * stream it ended. */ + takeCancel(): { reason?: string } | undefined { + const c = this._pendingCancel; + this._pendingCancel = undefined; + return c; + } + /** * Check if agent has pending tool calls. */ @@ -1030,7 +1062,7 @@ export class Agent { if (this._state.status === 'streaming' || (this._state.status === 'waiting_for_tools' && this._state.stream)) { const durationMs = Date.now() - this._inferenceStartedAt; - this.cancelStream(); + this.cancelStream(reason); return { aborted: true, durationMs }; } diff --git a/src/framework.ts b/src/framework.ts index 87f9c2a4..b28ef1b8 100644 --- a/src/framework.ts +++ b/src/framework.ts @@ -1046,7 +1046,7 @@ export class AgentFramework { * `inference:exhausted` (which also pollutes the failure streak). Kept * separate from ephemeralRuns deliberately: endTurn/budget cancels happen * for resident agents too, and the key is per-stream, not per-agent. */ - private frameworkCancelledStreams: Map = new Map(); + private frameworkCancelledStreams: Map = new Map(); /** Active runEphemeralToCompletion runs, keyed by agent name. */ private ephemeralRuns: Map = new Map(); /** Ephemeral namespaces/names are single-generation for this framework @@ -1814,10 +1814,17 @@ export class AgentFramework { // A stopped host must never hang behind its own cooldown. this.cancelProviderAdmission(); - // Cancel all active streams + // Cancel all active streams. Membrane reports every stream.cancel() as + // reason 'user' — it names the call, not the actor — so the provenance + // is recorded HERE, before the cancel, the same way endTurn, budget + // restarts and quiesce do. driveStream's tracked branch then settles the + // turn as a shutdown: no [turn-interrupted] marker (nobody stopped the + // agent; the host went away), no inference:exhausted (not a failure), + // one inference:aborted with reason 'shutdown'. for (const agent of this.agents.values()) { if (agent.state.status === 'streaming' || (agent.state.status === 'waiting_for_tools' && agent.state.stream)) { + this.frameworkCancelledStreams.set(`${agent.name}:${agent.streamId}`, 'shutdown'); agent.cancelStream(); } } @@ -2263,13 +2270,102 @@ export class AgentFramework { if (!agent) { return false; } + const hadStream = agent.state.status === 'streaming' || + (agent.state.status === 'waiting_for_tools' && !!agent.state.stream); const result = agent.abortInference(reason); - if (result) { + // A streaming abort is reported by driveStream when the stream's own + // `aborted` event lands — once, carrying this reason (Agent hands it + // over via takeCancel). Emitting here as well produced two + // inference:aborted traces per abort. Only the non-streaming inference + // path has no stream driver to report for it. + if (result && !hadStream) { this.emitTrace({ type: 'inference:aborted', agentName, reason, durationMs: result.durationMs }); } return !!result; } + /** + * Terminal for a deliberate cancellation of a live stream: the host's Stop + * button, an admin abort, a subagent reclaim — anything that went through + * Agent.cancelStream(reason) / abortInference(reason). Both of a stream's + * cancel twins land here (`aborted` with wire reason 'user', or `error` + * for implementations that report cancel() that way — see + * host-quiesce.test.ts), so the two cannot drift: settle the turn, ONE + * inference:aborted carrying the caller's reason, the neutral + * [turn-interrupted] marker, an attributable inference-log terminal, gate + * release. Deliberately none of the failure accounting: inference:failed / + * inference:exhausted feed the consecutive-failure streak, hard-down ops + * alerts, the poison-history breaker and the "[inference-failed]" marker — + * all of which told a resident, in the user's voice, that its turn failed + * and advised remediation for a failure that never happened — and + * errorPolicy would relaunch the very inference someone just stopped. + * + * Only called while `agent` still owns stream `myStreamId`. + */ + private settleDeliberateCancel(args: { + agent: Agent; + myStreamId: number; + startTime: number; + requestId: string; + compiledRequest: NormalizedRequest | undefined; + reason: string; + }): void { + const { agent, myStreamId, startTime, requestId, compiledRequest, reason } = args; + const durationMs = Date.now() - startTime; + this.abortAgentScript(agent.name, 'stream cancelled'); + agent.reset(); + this.settleAgent(agent.name, { + stopReason: 'exhausted', + speech: '', + error: 'Stream cancelled', + }); + // Distinct trace type: emitTrace funnels every inference:exhausted into + // noteInferenceExhausted (streak, failures.log, marker); inference:aborted + // carries the honest cause without any of that. + this.emitTrace({ + type: 'inference:aborted', + agentName: agent.name, + reason, + durationMs, + }); + // Agent-facing marker in the same system envelope as the failure marker + // so surfaces render it the same way. The text names the act and not an + // actor (the wire reason names the call, not who made it) and says + // nothing about delivery: earlier rounds of this turn may have been + // live-routed already, so "your output was not delivered" would be a + // false claim. addMessage alone requests no inference: no loop. + try { + agent.getContextManager().addMessage( + 'user', + [{ + type: 'text', + text: + `[turn-interrupted] Your previous turn was cancelled mid-stream — ` + + `a deliberate stop, not a failure. Whatever you were still ` + + `producing when it stopped was cut off there.`, + }], + { system: true, kind: 'turn-interrupted', reason }, + ); + } catch (err) { + console.error(`[turn-interrupted] could not record chronicle marker for ${agent.name}:`, err); + } + // Postmortem 2026-05-28 P2 #7: persist the abort to the inference log so + // future investigations can attribute the terminal cause without relying + // on live in-memory reducer state. Without this, abort-terminated + // inferences are invisible to forensic queries (only request-side + // telemetry via llm-calls.jsonl shows them, and only by absence). + this.logInference({ + timestamp: startTime, + agentName: agent.name, + requestId, + success: false, + error: `Stream cancelled (${reason})`, + request: compiledRequest ?? { note: 'streaming request aborted' }, + durationMs, + }); + if (this.agents.get(agent.name) === agent && agent.streamId === myStreamId) this.eventGate?.onInferenceEnded(agent.name); + } + /** * Get all registered modules. */ @@ -8574,6 +8670,9 @@ export class AgentFramework { // failed, a framework-cancelled stream sets aborted. Best-effort // notifications — consumers dedupe by inferenceId and keep a timeout. let lifecyclePhase: 'completed' | 'aborted' | 'failed' = 'completed'; + // Set by the complete / aborted / error cases: the stream SAID how it + // ended. An iterator that just closes says nothing (see after the loop). + let streamTerminal = false; let preserveEventGateForSuccessor = false; this.hookOrchestrator?.emitLifecycle({ inferenceId: outgoingInferenceId, @@ -8845,6 +8944,7 @@ export class AgentFramework { } case 'complete': { + streamTerminal = true; adoptInjectedRound(); const durationMs = Date.now() - startTime; const response = event.response; @@ -9339,8 +9439,85 @@ export class AgentFramework { } case 'error': { + streamTerminal = true; const err = event.error; const durationMs = Date.now() - startTime; + + // A cancel may surface as `error` instead of `aborted` depending + // on how the stream implementation reports the cancellation + // (host-quiesce.test.ts pins that shape). Read the cancel + // provenance BEFORE any failure accounting: a deliberate stop + // must not first emit a contradictory inference:failed / failed + // inference-log row, and the caller's pending cancel is consumed + // here exactly as the `aborted` twin consumes it — once, by the + // stream it ended. + const cancelKey = `${agent.name}:${myStreamId}`; + const cancelKind = this.frameworkCancelledStreams.get(cancelKey); + const callerCancel = agent.takeCancel(); + + if (cancelKind === 'quiesce_abandoned' || cancelKind === 'shutdown') { + // Same contract as the aborted branch: settle honestly, no + // errorPolicy retry (which would relaunch inference + // mid-maintenance-window), no inference:exhausted (which feeds + // the failure streak / hard-down / poison-history accounting), + // and — a framework-owned cancel — no [turn-interrupted] + // marker: nobody stopped the agent. §10.5 lifecycle reads + // 'aborted' (not 'completed'), and the gate is released only + // for THIS agent instance — a name re-registered meanwhile + // (ephemeral disposal, conversation replacement) owns its own + // gate liveness. + this.frameworkCancelledStreams.delete(cancelKey); + lifecyclePhase = 'aborted'; + if (agent.streamId === myStreamId) { + this.abortAgentScript(agent.name, cancelKind === 'shutdown' ? 'framework shutting down' : 'turn abandoned by quiesce'); + agent.reset(); + this.settleAgent(agent.name, { + stopReason: 'exhausted', + speech: '', + error: cancelKind === 'shutdown' ? 'Framework shutting down' : 'Turn abandoned by operator quiesce', + }); + } + // Postmortem 2026-05-28 P2 #7: the abort is a terminal the + // inference log must be able to attribute later. + this.logInference({ + timestamp: startTime, + agentName: agent.name, + requestId, + success: false, + error: `Stream aborted: ${cancelKind}`, + request: compiledRequest ?? { note: `streaming request aborted by ${cancelKind}` }, + durationMs, + }); + this.emitTrace({ + type: 'inference:aborted', + agentName: agent.name, + durationMs, + reason: cancelKind, + }); + if (this.agents.get(agent.name) === agent && agent.streamId === myStreamId) { + this.eventGate?.onInferenceEnded(agent.name); + } + break; + } + + if (callerCancel !== undefined) { + // The caller-cancel twin: Agent.cancelStream(reason) / + // framework.abortInference(reason) on a stream that reports + // cancel() as `error`. Same terminal as `aborted` with wire + // reason 'user' — settleDeliberateCancel is the one place both + // twins end. Without this the Stop fell through to errorPolicy + // and relaunched the inference it had just stopped + // (inference:started → failed → started; Sol, 2026-09-23). + lifecyclePhase = 'aborted'; // §10.5 — terminal emitted in finally + if (agent.streamId === myStreamId) { + this.settleDeliberateCancel({ + agent, myStreamId, startTime, requestId, compiledRequest, + reason: callerCancel.reason ?? 'user', + }); + } + break; + } + this.emitTrace({ type: 'inference:failed', agentName: agent.name, @@ -9364,41 +9541,6 @@ export class AgentFramework { this.abortAgentScript(agent.name, 'stream error'); agent.reset(); - // A quiesce-abandoned cancel may surface as `error` instead of - // `aborted` depending on how the stream implementation reports - // the cancellation. Same contract as the aborted branch: settle - // honestly, no errorPolicy retry (which would relaunch inference - // mid-maintenance-window), no inference:exhausted (which feeds - // the failure streak / hard-down / poison-history accounting). - { - const cancelKey = `${agent.name}:${myStreamId}`; - if (this.frameworkCancelledStreams.get(cancelKey) === 'quiesce_abandoned') { - this.frameworkCancelledStreams.delete(cancelKey); - // Same terminal as the `aborted` twin: §10.5 lifecycle reads - // 'aborted' (not 'completed'), and the gate is released only - // for THIS agent instance — a name re-registered meanwhile - // (ephemeral disposal, conversation replacement) owns its own - // gate liveness. The inference-log terminal was already - // written at the top of this case (`Stream error`). - lifecyclePhase = 'aborted'; - this.settleAgent(agent.name, { - stopReason: 'exhausted', - speech: '', - error: 'Turn abandoned by operator quiesce', - }); - this.emitTrace({ - type: 'inference:aborted', - agentName: agent.name, - durationMs, - reason: 'quiesce_abandoned', - }); - if (this.agents.get(agent.name) === agent && agent.streamId === myStreamId) { - this.eventGate?.onInferenceEnded(agent.name); - } - break; - } - } - if (ownsProviderGate && this.holdProviderAcceleration(agent, err, trigger)) { lifecyclePhase = 'failed'; if (this.agents.get(agent.name) === agent && agent.streamId === myStreamId) this.eventGate?.onInferenceEnded(agent.name); @@ -9434,6 +9576,11 @@ export class AgentFramework { } case 'aborted': { + streamTerminal = true; + // Consumed once per terminal event, whichever path follows: a + // framework-owned cancel (stop() also goes through cancelStream) + // must not leave a caller record behind for a later stream. + const callerCancel = agent.takeCancel(); // Framework-initiated, non-terminal cancels (endTurn tool result, // context-budget restart) also surface here as `aborted` — but // the turn either already settled (endTurn) or a replacement @@ -9478,6 +9625,42 @@ export class AgentFramework { }); return; } + if (cancelKind === 'shutdown') { + // AgentFramework.stop() with this stream still live. Same + // contract as quiesce: settle honestly so an ephemeral + // run's promise rejects now instead of riding out its + // idle watchdog, no inference:exhausted, and no chronicle + // marker of any kind — nobody stopped the agent, and a + // resident must not read after restart that someone did. + const durationMs = Date.now() - startTime; + if (agent.streamId === myStreamId) { + this.abortAgentScript(agent.name, 'framework shutting down'); + agent.reset(); + this.settleAgent(agent.name, { + stopReason: 'exhausted', + speech: '', + error: 'Framework shutting down', + }); + } + // Postmortem 2026-05-28 P2 #7: the abort is a terminal the + // inference log must be able to attribute later. + this.logInference({ + timestamp: startTime, + agentName: agent.name, + requestId, + success: false, + error: 'Stream aborted: shutdown', + request: compiledRequest ?? { note: 'streaming request aborted by shutdown' }, + durationMs, + }); + this.emitTrace({ + type: 'inference:aborted', + agentName: agent.name, + durationMs, + reason: 'shutdown', + }); + return; + } // endTurn IS a logical turn end — earlier rounds may have // live-routed prose (narrate → skip_reply is a real shape), // so settle the delivery chain and drop the receipt. A @@ -9493,38 +9676,66 @@ export class AgentFramework { } } const reason = event.reason ?? 'unknown'; + // Membrane (>= 0.5.81) reports reason 'user' exactly when the + // request's signal was aborted, i.e. someone called cancel(): + // the host's Stop button, an admin abort, a subagent reclaim. + // That is a DELIBERATE cancellation, not a failure, and it must + // not feed the consecutive-failure streak, ops alerts, or the + // "[inference-failed] the model call failed" marker — all of + // which told the agent, in the user's voice, that its turn + // failed and advised remediation for a failure that never + // happened. For a long-lived resident whose transcript is + // memory, those accumulate as false self-knowledge. + // + // What 'user' does NOT say is WHO cancelled: the wire reason + // names the call, not the actor. Framework-owned cancels record + // their provenance in frameworkCancelledStreams and returned + // above; a caller of Agent.cancelStream(reason) / + // framework.abortInference(reason) hands its reason over here. + // The marker text therefore names the act and not an actor, and + // says nothing about delivery: earlier rounds of this turn may + // have been live-routed already, so "your output was not + // delivered" would be a false claim. + const deliberate = reason === 'user'; // Only reset if this is still the active stream (a budget restart // may have already started a new stream, bumping streamId) if (agent.streamId === myStreamId) { - const durationMs = Date.now() - startTime; - this.abortAgentScript(agent.name, `stream aborted (${reason})`); - agent.reset(); - this.settleAgent(agent.name, { - stopReason: 'exhausted', - speech: '', - error: `Stream aborted: ${reason}`, - }); - this.emitTrace({ - type: 'inference:exhausted', - agentName: agent.name, - error: `Stream aborted: ${reason}`, - }); - // Postmortem 2026-05-28 P2 #7: persist the abort to the - // inference log so future investigations can attribute the - // terminal cause without relying on live in-memory reducer - // state. Without this, abort-terminated inferences are - // invisible to forensic queries (only request-side telemetry - // via llm-calls.jsonl shows them, and only by absence). - this.logInference({ - timestamp: startTime, - agentName: agent.name, - requestId, - success: false, - error: `Stream aborted: ${reason}`, - request: compiledRequest ?? { note: 'streaming request aborted' }, - durationMs, - }); - if (this.agents.get(agent.name) === agent && agent.streamId === myStreamId) this.eventGate?.onInferenceEnded(agent.name); + if (deliberate) { + lifecyclePhase = 'aborted'; // §10.5 — terminal emitted in finally + this.settleDeliberateCancel({ + agent, myStreamId, startTime, requestId, compiledRequest, + reason: callerCancel?.reason ?? 'user', + }); + } else { + const durationMs = Date.now() - startTime; + const terminal = `Stream aborted: ${reason}`; + this.abortAgentScript(agent.name, `stream aborted (${reason})`); + agent.reset(); + this.settleAgent(agent.name, { + stopReason: 'exhausted', + speech: '', + error: terminal, + }); + this.emitTrace({ + type: 'inference:exhausted', + agentName: agent.name, + error: terminal, + }); + // Postmortem 2026-05-28 P2 #7: persist the abort to the + // inference log so future investigations can attribute the + // terminal cause without relying on live in-memory reducer + // state. + this.logInference({ + timestamp: startTime, + agentName: agent.name, + requestId, + success: false, + error: terminal, + request: compiledRequest ?? { note: 'streaming request aborted' }, + durationMs, + }); + if (this.agents.get(agent.name) === agent && agent.streamId === myStreamId) this.eventGate?.onInferenceEnded(agent.name); + } } break; } @@ -9617,6 +9828,26 @@ export class AgentFramework { } } } + // The third cancel shape: a stream whose cancel() just closes the + // iterator, with no `aborted` or `error` event at all. The loop ends + // having been told nothing, so nothing above settled the turn — the + // lifecycle terminal would read `completed`, no inference:aborted would + // be emitted, and an in-flight runEphemeralToCompletion would ride out + // its watchdog. If THIS side cancelled the stream, that cancel is the + // terminal: the same end as the `aborted` and `error` twins, the + // caller's reason consumed here by the stream it ended. A silent end + // nobody asked for keeps the prior behaviour (the finally still + // releases the gate). Greptile on #172, 2026-09-28. + if (!streamTerminal && this.agents.get(agent.name) === agent && agent.streamId === myStreamId) { + const callerCancel = agent.takeCancel(); + if (callerCancel !== undefined) { + lifecyclePhase = 'aborted'; // §10.5 — terminal emitted in finally + this.settleDeliberateCancel({ + agent, myStreamId, startTime, requestId, compiledRequest, + reason: callerCancel.reason ?? 'user', + }); + } + } } catch (error) { // A superseded physical stream may throw after its replacement starts. // It has no authority to settle, reset, gate-release, or publish failure. @@ -11873,12 +12104,21 @@ export class AgentFramework { ? `${label.startsWith('#') ? label : `#${label}`} (${channelId})` : channelId : 'the channel'; + // No-locus failures (headless / WebUI turns with no home or + // trigger channel) are not Discord failures: "[discord-send- + // failed] could not be delivered to the channel" sent the agent + // debugging a Discord problem that doesn't exist. The text names + // the real situation; the machine-readable `kind` is kept stable + // for downstream consumers (gate intents, surfaces) that key on + // it. Neither text claims the reply reached no one: only channel + // routing failed here, and other dispatchSpeech handlers may + // have shown it. + const text = channelId + ? `[discord-send-failed] Your previous reply (${textLen} chars) could not be delivered to ${where} (${reason}). It was saved to your archive but the human did not receive it.` + : `[send-undeliverable] Your previous reply (${textLen} chars) had no channel to go to — ${reason}. This is a routing/configuration situation, not a channel failure. This route delivered it nowhere; if another surface showed it (a reply thread, the console), that copy stands. It is saved in your archive.`; this.addMessage( 'user', - [{ - type: 'text', - text: `[discord-send-failed] Your previous reply (${textLen} chars) could not be delivered to ${where} (${reason}). It was saved to your archive but the human did not receive it.`, - }], + [{ type: 'text', text }], { system: true, kind: 'discord-send-failed', channelId: channelId ?? '', reason }, ); } catch (err) { diff --git a/test/cancel-surfaces-as-error.test.ts b/test/cancel-surfaces-as-error.test.ts new file mode 100644 index 00000000..aa10ffdb --- /dev/null +++ b/test/cancel-surfaces-as-error.test.ts @@ -0,0 +1,223 @@ +import { describe, it } from 'node:test'; +import assert from 'node:assert/strict'; +import { mkdtempSync, rmSync } from 'node:fs'; +import { tmpdir } from 'node:os'; +import { join } from 'node:path'; +import type { NormalizedRequest, StreamEvent, YieldingStream } from '@animalabs/membrane'; +import type { TraceEvent } from '../src/index.js'; +import { AgentFramework } from '../src/index.js'; +import { MockMembrane } from './helpers/mock-membrane.js'; + +/** + * A cancel that the stream reports as `error` is still a deliberate stop. + * + * Membrane's yielding stream answers cancel() with an `aborted` event, but + * that is one implementation's choice: host-quiesce.test.ts already pins + * that a stream may report cancel() through `error` instead, and driveStream + * honours quiesce/shutdown provenance on that twin. The caller-cancel twin + * was missed (Sol's review of #172, 2026-09-23): abortInference(reason) / + * cancelStream() on an error-surfacing stream fell through the failure + * pipeline — inference:failed, errorPolicy, and a retry that relaunched the + * very inference someone had just stopped: + * + * inference:started → inference:failed → inference:started + * + * Both twins must reach the same terminal: one inference:aborted carrying + * the caller's reason, the [turn-interrupted] marker, no inference:failed, + * no inference:exhausted, no retry. + */ + +/** Parks until cancel(), then reports the cancel as an error event. */ +class ErroringOnCancelStream implements YieldingStream { + private release: (() => void) | null = null; + private cancelled = false; + cancel(): void { this.cancelled = true; this.release?.(); } + provideToolResults(): void {} + get isWaitingForTools() { return false; } + get pendingToolCallIds(): string[] { return []; } + get toolDepth() { return 0; } + async *[Symbol.asyncIterator](): AsyncIterator { + if (!this.cancelled) await new Promise((resolve) => { this.release = resolve; }); + yield { type: 'error', error: new Error('stream cancelled') } as StreamEvent; + } +} + +/** Parks until cancel(), then just ENDS — no `aborted`, no `error`. The + * third shape: an implementation whose cancel() closes the iterator. */ +class SilentlyEndingOnCancelStream extends ErroringOnCancelStream { + override async *[Symbol.asyncIterator](): AsyncIterator { + await new Promise((resolve) => { (this as unknown as { release: () => void }).release = resolve; }); + } +} + +/** Fails spontaneously — the genuine provider error the failure pipeline is for. */ +class SpontaneouslyErroringStream extends ErroringOnCancelStream { + override async *[Symbol.asyncIterator](): AsyncIterator { + await new Promise((r) => setTimeout(r, 20)); + yield { type: 'error', error: new Error('provider exploded') } as StreamEvent; + } +} + +class ErroringMembrane extends MockMembrane { + constructor(private readonly make: () => YieldingStream) { super(); } + override streamYielding(request: NormalizedRequest): YieldingStream { + this.calls.push(request); + return this.make(); + } +} + +async function waitFor(cond: () => boolean, ms = 2000): Promise { + const start = Date.now(); + while (!cond()) { + if (Date.now() - start > ms) throw new Error('timeout waiting for condition'); + await new Promise((r) => setTimeout(r, 10)); + } +} + +async function boot(prefix: string, make: () => YieldingStream) { + const tempDir = mkdtempSync(join(tmpdir(), prefix)); + const membrane = new ErroringMembrane(make); + const framework = await AgentFramework.create({ + storePath: join(tempDir, 'test.chronicle'), + membrane: membrane.asMembrane(), + agents: [{ name: 'assistant', model: 'test-model', systemPrompt: 'Assist.' }], + modules: [], + }); + const traces: TraceEvent[] = []; + framework.onTrace((t) => { traces.push(t); }); + framework.start(); + framework.nudgeAgent('assistant', 'operator'); + await waitFor(() => membrane.calls.length === 1); + const agent = framework.getAgent('assistant')!; + await waitFor(() => agent.state.status === 'streaming'); + const teardown = async () => { + await framework.stop(); + rmSync(tempDir, { recursive: true, force: true }); + }; + return { framework, membrane, agent, traces, teardown }; +} + +function contextTexts(framework: AgentFramework): string[] { + const { messages } = framework.getAgent('assistant')!.getContextManager().queryMessages({}); + return messages.flatMap((m) => + m.content.filter((b): b is { type: 'text'; text: string } => b.type === 'text').map((b) => b.text)); +} + +function types(traces: TraceEvent[], ...wanted: string[]): string[] { + return traces.map((t) => String(t.type)).filter((t) => wanted.includes(t)); +} + +describe('a cancel that the stream reports as `error` is still a deliberate stop', () => { + it('abortInference(reason) → ONE inference:aborted with the reason; no inference:failed, no retry', async () => { + const { framework, membrane, agent, traces, teardown } = + await boot('cancel-error-abort-', () => new ErroringOnCancelStream()); + try { + assert.equal(framework.abortInference('assistant', 'operator_reclaim'), true); + await waitFor(() => traces.some((t) => t.type === 'inference:aborted')); + // Give the failure pipeline's retry every chance to show up (the + // observed relaunch landed ~1.2s after the stop). + await new Promise((r) => setTimeout(r, 1500)); + + const aborted = traces.filter((t) => t.type === 'inference:aborted') as Array<{ reason?: string }>; + assert.deepEqual(aborted.map((t) => t.reason), ['operator_reclaim'], + `expected exactly one inference:aborted with the caller's reason, got ${JSON.stringify(aborted)}`); + assert.deepEqual(types(traces, 'inference:failed', 'inference:exhausted'), [], + `a deliberate stop is neither a failure nor an exhaustion: ${JSON.stringify(types(traces, 'inference:started', 'inference:failed', 'inference:exhausted', 'inference:aborted'))}`); + assert.equal(membrane.calls.length, 1, 'the stopped inference must not be relaunched by errorPolicy'); + assert.equal(types(traces, 'inference:started').length, 1); + assert.equal(agent.state.status, 'idle'); + + // Same marker as the `aborted` twin: neutral text, reason in metadata. + await waitFor(() => contextTexts(framework).some((t) => t.includes('[turn-interrupted]'))); + const { messages } = agent.getContextManager().queryMessages({}); + const markerMsg = messages.find((m) => + m.content.some((b) => b.type === 'text' && b.text.includes('[turn-interrupted]')))!; + assert.equal((markerMsg.metadata as { reason?: string } | undefined)?.reason, 'operator_reclaim'); + assert.ok(!contextTexts(framework).some((t) => t.includes('[inference-failed]'))); + } finally { + await teardown(); + } + }); + + it('cancelStream() with no reason → inference:aborted reason user, marker, no failure accounting', async () => { + const { framework, membrane, agent, traces, teardown } = + await boot('cancel-error-plain-', () => new ErroringOnCancelStream()); + try { + agent.cancelStream(); + await waitFor(() => traces.some((t) => t.type === 'inference:aborted')); + await new Promise((r) => setTimeout(r, 1500)); + + const aborted = traces.filter((t) => t.type === 'inference:aborted') as Array<{ reason?: string }>; + assert.deepEqual(aborted.map((t) => t.reason), ['user']); + assert.deepEqual(types(traces, 'inference:failed', 'inference:exhausted'), []); + assert.equal(membrane.calls.length, 1); + await waitFor(() => contextTexts(framework).some((t) => t.includes('[turn-interrupted]'))); + } finally { + await teardown(); + } + }); + + it('a stream that errors on its own (nobody cancelled) still takes the failure pipeline', async () => { + const { membrane, traces, teardown } = + await boot('cancel-error-genuine-', () => new SpontaneouslyErroringStream()); + try { + await waitFor(() => traces.some((t) => t.type === 'inference:failed')); + // errorPolicy decides retry vs exhausted; either way it is a failure, + // never a deliberate stop. + await waitFor(() => traces.some((t) => t.type === 'inference:exhausted') || membrane.calls.length > 1, 3000); + assert.deepEqual(types(traces, 'inference:aborted'), [], 'a provider error is not a cancellation'); + } finally { + await teardown(); + } + }); + + it('a cancel whose stream ends with NO terminal event still ends as a deliberate stop', async () => { + // Greptile on #172 (2026-09-28): an iterator that closes on cancel() + // without `aborted` or `error` told the driver nothing, so nothing + // settled the turn — no inference:aborted, the lifecycle terminal read + // `completed`, and an ephemeral run would ride out its watchdog. The + // caller's cancel IS the terminal for that shape too. + const { framework, membrane, agent, traces, teardown } = await boot('cancel-silent-', () => new SilentlyEndingOnCancelStream()); + try { + framework.abortInference('assistant', 'silent_stop'); + await waitFor(() => traces.some((t) => t.type === 'inference:aborted')); + await new Promise((r) => setTimeout(r, 150)); + assert.deepEqual(types(traces, 'inference:aborted', 'inference:failed', 'inference:exhausted'), ['inference:aborted'], + `terminals: ${JSON.stringify(types(traces, 'inference:aborted', 'inference:failed', 'inference:exhausted'))}`); + const aborted = traces.find((t) => t.type === 'inference:aborted') as { reason?: string }; + assert.equal(aborted.reason, 'silent_stop', 'the caller\'s reason rides the terminal'); + assert.equal(agent.state.status, 'idle'); + assert.equal(membrane.calls.length, 1, 'no retry of the inference that was stopped'); + assert.ok(contextTexts(framework).some((t) => t.startsWith('[turn-interrupted]')), 'the neutral marker is written'); + } finally { + await teardown(); + } + }); + + it('the reason of a cancel is consumed by the stream it ended, never by the next one', async () => { + // A cancel whose stream reports NO terminal event at all (iterator just + // ends) is collected at that stream's loop end — one inference:aborted, + // the first stream's. The next stream must start clean: an error there + // is a genuine failure, not a stale stop, and adds no second abort. + let n = 0; + const { framework, membrane, agent, traces, teardown } = await boot('cancel-error-stale-', () => + n++ === 0 ? new SilentlyEndingOnCancelStream() : new SpontaneouslyErroringStream()); + try { + agent.cancelStream('stale_reason'); + await waitFor(() => agent.state.status === 'idle'); + await waitFor(() => traces.some((t) => t.type === 'inference:aborted')); + const firstStreamAborts = traces.filter((t) => t.type === 'inference:aborted').length; + assert.equal(firstStreamAborts, 1, 'the silent end is settled by the cancel that ended it'); + await new Promise((r) => setTimeout(r, 100)); + framework.nudgeAgent('assistant', 'operator'); + await waitFor(() => membrane.calls.length === 2); + await waitFor(() => traces.some((t) => t.type === 'inference:failed'), 3000); + await new Promise((r) => setTimeout(r, 100)); + const aborted = traces.filter((t) => t.type === 'inference:aborted') as Array<{ reason?: string }>; + assert.equal(aborted.length, firstStreamAborts, + `the first stream's reason leaked into the second: ${JSON.stringify(aborted)}`); + } finally { + await teardown(); + } + }); +}); diff --git a/test/host-quiesce.test.ts b/test/host-quiesce.test.ts index cae60190..a65fe5cb 100644 --- a/test/host-quiesce.test.ts +++ b/test/host-quiesce.test.ts @@ -579,6 +579,10 @@ test('a cancel that surfaces as a stream ERROR still settles as an operator abor assert.ok(traces.some((t) => t.type === 'inference:aborted' && (t as { reason?: string }).reason === 'quiesce_abandoned')); assert.ok(!traces.some((t) => t.type === 'inference:exhausted')); + // The cancel provenance is read BEFORE the failure accounting: the same + // stream must not be reported failed and then aborted. + assert.ok(!traces.some((t) => t.type === 'inference:failed'), + 'an abandoned turn is not first recorded as a provider failure'); const lifecycle = traces .filter((t) => String(t.type).includes('lifecycle')) .map((t) => (t as { phase?: string }).phase); diff --git a/test/shutdown-not-user-attributed.test.ts b/test/shutdown-not-user-attributed.test.ts new file mode 100644 index 00000000..9c29a0a2 --- /dev/null +++ b/test/shutdown-not-user-attributed.test.ts @@ -0,0 +1,230 @@ +import { describe, it } from 'node:test'; +import assert from 'node:assert/strict'; +import { mkdtempSync, rmSync } from 'node:fs'; +import { tmpdir } from 'node:os'; +import { join } from 'node:path'; +import type { + EventResponse, + Module, + ModuleContext, + ProcessEvent, + ProcessState, + ToolCall, + ToolDefinition, + ToolResult, + TraceEvent, +} from '../src/index.js'; +import type { NormalizedRequest, StreamEvent, YieldingStream } from '@animalabs/membrane'; +import { AgentFramework } from '../src/index.js'; +import { createMockResponse, MockMembrane } from './helpers/mock-membrane.js'; + +/** + * A graceful framework shutdown with an active stream is neither the user's + * act nor a failure. Membrane reports every `stream.cancel()` as reason + * `user` — it names the call, not the actor — so `AgentFramework.stop()` + * would otherwise take the deliberate-cancellation branch and write a + * "[turn-interrupted]" marker into the agent's durable context. A resident + * reading that after restart would learn that they were stopped; nobody + * stopped them, the host went away. + * + * stop() records shutdown provenance in `frameworkCancelledStreams` before + * the cancel, as endTurn / budget restarts / quiesce do, and driveStream's + * tracked branch settles the turn like a quiesce: no marker, no + * inference:exhausted, one inference:aborted with reason 'shutdown', and the + * agent settled — so an ephemeral run's promise rejects now instead of + * riding out its 15-minute idle watchdog (the first cut of this fix returned + * without settling; Sol's review repro, lifted below). + */ + +/** Module whose tool call hangs — keeps the stream open so the shutdown + * arrives mid-turn. */ +class HangingToolModule implements Module { + readonly name = 'test'; + async start(_ctx: ModuleContext): Promise {} + async stop(): Promise {} + getTools(): ToolDefinition[] { + return [{ name: 'hang', description: 'Hangs', inputSchema: { type: 'object', properties: {} } }]; + } + async handleToolCall(_call: ToolCall): Promise { + await new Promise(() => {}); + return { success: true, data: {} }; + } + async onProcess(event: ProcessEvent, _state: ProcessState): Promise { + if (event.type === 'external-message') { + return { + addMessages: [{ participant: 'User', content: [{ type: 'text', text: String(event.content) }] }], + requestInference: true, + }; + } + return {}; + } +} + +async function waitFor(cond: () => boolean, ms = 3000): Promise { + const start = Date.now(); + while (!cond()) { + if (Date.now() - start > ms) throw new Error('timeout waiting for condition'); + await new Promise((r) => setTimeout(r, 10)); + } +} + +function hungMembrane(): MockMembrane { + const membrane = new MockMembrane(); + membrane.pushResponse(createMockResponse( + [{ type: 'tool_use', id: 't1', name: 'test--hang', input: {} } as never], + 'tool_use', + )); + return membrane; +} + +describe('graceful shutdown is not attributed to the user', () => { + it('stop() with an active resident stream: no marker, one inference:aborted (reason shutdown), agent settled', async () => { + const tempDir = mkdtempSync(join(tmpdir(), 'shutdown-provenance-')); + const framework = await AgentFramework.create({ + storePath: join(tempDir, 'test.chronicle'), + membrane: hungMembrane().asMembrane(), + agents: [{ name: 'assistant', model: 'test-model', systemPrompt: 'Assist.' }], + modules: [new HangingToolModule()], + }); + const traces: TraceEvent[] = []; + framework.onTrace((t) => { traces.push(t); }); + let stopped = false; + try { + framework.pushEvent({ type: 'external-message', source: 'test', content: 'go', metadata: {} }); + framework.start(); + const agent = framework.getAgent('assistant')!; + await waitFor(() => agent.state.status === 'waiting_for_tools'); + + // Intercept context writes: the store is closed by the time stop() + // returns, so capture any marker at write time. + const cm = agent.getContextManager(); + const written: string[] = []; + const orig = cm.addMessage.bind(cm); + (cm as unknown as { addMessage: unknown }).addMessage = (role: never, content: Array<{ type: string; text?: string }>, meta: never) => { + for (const b of content) if (b.type === 'text' && b.text) written.push(b.text); + return orig(role, content as never, meta); + }; + + // No user action anywhere: the host process is shutting down while + // the stream is still active. + stopped = true; + await framework.stop(); + + const markers = written.filter((t) => t.includes('[turn-interrupted]') || t.includes('[inference-failed]')); + assert.deepEqual(markers, [], `graceful shutdown wrote a marker: ${JSON.stringify(markers)}`); + + const aborted = traces.filter((t): t is Extract => t.type === 'inference:aborted'); + assert.equal(aborted.length, 1, `expected exactly one inference:aborted trace, got ${JSON.stringify(aborted)}`); + assert.equal(aborted[0].reason, 'shutdown', 'the trace carries the recorded provenance, not the wire reason'); + assert.equal(traces.filter((t) => t.type === 'inference:exhausted').length, 0, + 'a shutdown is not a failure: no inference:exhausted'); + assert.equal(agent.state.status, 'idle', 'the shutdown settles the agent instead of leaving it waiting_for_tools'); + } finally { + if (!stopped) await framework.stop(); + rmSync(tempDir, { recursive: true, force: true }); + } + }); + + // The error twin: a stream implementation may report cancel() through + // `error` rather than `aborted` (host-quiesce.test.ts pins the shape). + // Shutdown provenance must be read there BEFORE any failure accounting — + // no inference:failed, no failure marker, no errorPolicy retry — and the + // agent still settles. + it('stop() with a stream that reports cancel() as `error`: same terminal, no failure trace first', async () => { + class ErroringOnCancelStream implements YieldingStream { + private release: (() => void) | null = null; + private cancelled = false; + cancel(): void { this.cancelled = true; this.release?.(); } + provideToolResults(): void {} + get isWaitingForTools() { return false; } + get pendingToolCallIds(): string[] { return []; } + get toolDepth() { return 0; } + async *[Symbol.asyncIterator](): AsyncIterator { + if (!this.cancelled) await new Promise((resolve) => { this.release = resolve; }); + yield { type: 'error', error: new Error('stream cancelled') } as StreamEvent; + } + } + class ErroringMembrane extends MockMembrane { + override streamYielding(request: NormalizedRequest): YieldingStream { + this.calls.push(request); + return new ErroringOnCancelStream(); + } + } + const tempDir = mkdtempSync(join(tmpdir(), 'shutdown-error-twin-')); + const membrane = new ErroringMembrane(); + const framework = await AgentFramework.create({ + storePath: join(tempDir, 'test.chronicle'), + membrane: membrane.asMembrane(), + agents: [{ name: 'assistant', model: 'test-model', systemPrompt: 'Assist.' }], + modules: [], + }); + const traces: TraceEvent[] = []; + framework.onTrace((t) => { traces.push(t); }); + let stopped = false; + try { + framework.start(); + framework.nudgeAgent('assistant', 'operator'); + const agent = framework.getAgent('assistant')!; + await waitFor(() => agent.state.status === 'streaming'); + + const cm = agent.getContextManager(); + const written: string[] = []; + const orig = cm.addMessage.bind(cm); + (cm as unknown as { addMessage: unknown }).addMessage = (role: never, content: Array<{ type: string; text?: string }>, meta: never) => { + for (const b of content) if (b.type === 'text' && b.text) written.push(b.text); + return orig(role, content as never, meta); + }; + + stopped = true; + await framework.stop(); + + const markers = written.filter((t) => t.includes('[turn-interrupted]') || t.includes('[inference-failed]')); + assert.deepEqual(markers, [], `graceful shutdown wrote a marker: ${JSON.stringify(markers)}`); + const kinds = traces.map((t) => String(t.type)) + .filter((t) => ['inference:started', 'inference:failed', 'inference:exhausted', 'inference:aborted'].includes(t)); + assert.deepEqual(kinds, ['inference:started', 'inference:aborted'], `expected started → aborted only, got ${JSON.stringify(kinds)}`); + const aborted = traces.find((t): t is Extract => t.type === 'inference:aborted')!; + assert.equal(aborted.reason, 'shutdown'); + assert.equal(membrane.calls.length, 1, 'no errorPolicy relaunch of the cancelled inference'); + assert.equal(agent.state.status, 'idle'); + } finally { + if (!stopped) await framework.stop(); + rmSync(tempDir, { recursive: true, force: true }); + } + }); + + // Review repro by Sol (Codex, via Anarchid) on #134, 2026-09-21: the first + // shutdown branch emitted its trace and returned without settling, so the + // ephemeral run stayed pending until its 15-minute idle watchdog rejected + // it with a false "stalled" error after the framework was gone. + it('stop() settles an in-flight ephemeral run promptly', async () => { + const tempDir = mkdtempSync(join(tmpdir(), 'rws-eph-')); + const framework = await AgentFramework.create({ + storePath: join(tempDir, 'test.chronicle'), + membrane: hungMembrane().asMembrane(), + agents: [{ name: 'assistant', model: 'test-model', systemPrompt: 'Assist.' }], + modules: [new HangingToolModule()], + }); + let stopped = false; + try { + framework.start(); + const created = await framework.createEphemeralAgent({ name: 'eph', model: 'test-model', systemPrompt: 'Run.' }); + created.contextManager.addMessage('user', [{ type: 'text', text: 'Run once.' }]); + let outcome = 'pending'; + const run = framework.runEphemeralToCompletion(created.agent, created.contextManager) + .then(() => { outcome = 'resolved'; }, (e) => { outcome = `rejected: ${(e as Error).message}`; }); + await waitFor(() => created.agent.state.status === 'waiting_for_tools'); + + stopped = true; + await framework.stop(); + await Promise.race([run, new Promise((r) => setTimeout(r, 1500))]); + + assert.notEqual(outcome, 'pending', 'ephemeral run still pending after stop()'); + assert.match(outcome, /^rejected: Framework shutting down/, `expected the shutdown terminal, got "${outcome}"`); + assert.equal(created.agent.state.status, 'idle'); + } finally { + if (!stopped) await framework.stop(); + rmSync(tempDir, { recursive: true, force: true }); + } + }); +}); diff --git a/test/user-interrupt-not-failure.test.ts b/test/user-interrupt-not-failure.test.ts new file mode 100644 index 00000000..0597f1fe --- /dev/null +++ b/test/user-interrupt-not-failure.test.ts @@ -0,0 +1,216 @@ +import { describe, it } from 'node:test'; +import assert from 'node:assert/strict'; +import { mkdtempSync, rmSync } from 'node:fs'; +import { tmpdir } from 'node:os'; +import { join } from 'node:path'; +import type { + EventResponse, + Module, + ModuleContext, + ProcessEvent, + ProcessState, + ToolCall, + ToolDefinition, + ToolResult, + TraceEvent, +} from '../src/index.js'; +import { AgentFramework } from '../src/index.js'; +import { createMockResponse, MockMembrane } from './helpers/mock-membrane.js'; + +/** + * A cancelled stream is a deliberate stop, not a model failure. + * + * Before the fix, every cancel (the host's Stop button, an admin abort, a + * subagent reclaim — all `agent.cancelStream()`) fell through driveStream's + * generic abort handling into the failure pipeline: an inference:exhausted + * trace, a bumped consecutive-failure streak (three Stops → hard-down ops + * alert), and an "[inference-failed] the model call failed and produced no + * response … drop an oversized attachment" chronicle marker attributed to + * the user — three inaccuracies (wrong cause, wrong speaker, irrelevant + * advice) accumulating as false self-knowledge in a resident's transcript. + * + * Membrane's wire reason 'user' means "the request's signal was aborted"; + * it does not say who did it. So the marker names the act, not an actor, + * and the caller's own reason (cancelStream(reason) / abortInference(reason)) + * is what the trace carries. + * + * Originally by Lari (#134); reworked after review (Sol, 2026-09-21). + */ + +/** Module whose tool call hangs until released — keeps the stream open so + * the test can cancel mid-turn, exactly as the TUI/WebUI Stop button does. */ +class HangingToolModule implements Module { + readonly name = 'test'; + release!: () => void; + private readonly gate = new Promise((resolve) => { this.release = resolve; }); + + async start(_ctx: ModuleContext): Promise {} + async stop(): Promise {} + + getTools(): ToolDefinition[] { + return [{ + name: 'hang', + description: 'Hangs until released', + inputSchema: { type: 'object', properties: {} }, + }]; + } + + async handleToolCall(_call: ToolCall): Promise { + await this.gate; + return { success: true, data: {} }; + } + + async onProcess(event: ProcessEvent, _state: ProcessState): Promise { + if (event.type === 'external-message') { + return { + addMessages: [{ participant: 'User', content: [{ type: 'text', text: String(event.content) }] }], + requestInference: true, + }; + } + return {}; + } +} + +async function waitFor(cond: () => boolean, ms = 2000): Promise { + const start = Date.now(); + while (!cond()) { + if (Date.now() - start > ms) throw new Error('timeout waiting for condition'); + await new Promise((r) => setTimeout(r, 10)); + } +} + +/** Boot a framework whose one agent is parked in waiting_for_tools on a + * hanging tool call. Everything after create() runs under the caller's + * try/finally so a waitFor timeout still stops the framework. */ +async function bootHung(prefix: string) { + const tempDir = mkdtempSync(join(tmpdir(), prefix)); + const membrane = new MockMembrane(); + membrane.pushResponse(createMockResponse( + [{ type: 'tool_use', id: 't1', name: 'test--hang', input: {} } as never], + 'tool_use', + )); + const module = new HangingToolModule(); + const framework = await AgentFramework.create({ + storePath: join(tempDir, 'test.chronicle'), + membrane: membrane.asMembrane(), + agents: [{ name: 'assistant', model: 'test-model', systemPrompt: 'Assist.' }], + modules: [module], + }); + const traces: TraceEvent[] = []; + framework.onTrace((t) => { traces.push(t); }); + const teardown = async () => { + // Release the hung tool BEFORE stopping: its completion pushes a + // tool-result event, which must land while the queue is still open. + module.release(); + await new Promise((r) => setTimeout(r, 50)); + await framework.stop(); + rmSync(tempDir, { recursive: true, force: true }); + }; + return { framework, membrane, module, traces, teardown }; +} + +function contextTexts(framework: AgentFramework): string[] { + const { messages } = framework.getAgent('assistant')!.getContextManager().queryMessages({}); + return messages.flatMap((m) => + m.content.filter((b): b is { type: 'text'; text: string } => b.type === 'text').map((b) => b.text)); +} + +describe('a cancelled stream is not recorded as a failure', () => { + it('cancelStream() mid-turn → neutral [turn-interrupted] marker, one inference:aborted (reason user), no failure streak', async () => { + const { framework, traces, teardown } = await bootHung('interrupt-'); + try { + framework.pushEvent({ type: 'external-message', source: 'test', content: 'go', metadata: {} }); + framework.start(); + const agent = framework.getAgent('assistant')!; + await waitFor(() => agent.state.status === 'waiting_for_tools'); + + // What the TUI / WebUI Stop button does. + agent.cancelStream(); + + await waitFor(() => traces.some((t) => t.type === 'inference:aborted')); + await waitFor(() => contextTexts(framework).some((t) => t.includes('[turn-interrupted]'))); + + const texts = contextTexts(framework); + const marker = texts.find((t) => t.includes('[turn-interrupted]'))!; + assert.match(marker, /deliberate stop, not a failure/); + // The wire reason names the call, not the actor: the marker must not + // attribute the stop to anyone … + assert.doesNotMatch(marker, /by the user|the user stopped/i, 'marker must not name an actor'); + // … and must not claim non-delivery: earlier rounds may have been + // live-routed, and the agent would be baited into a duplicate send. + assert.doesNotMatch(marker, /not delivered|did not receive/i, 'marker must not assert non-delivery'); + assert.ok(!texts.some((t) => t.includes('[inference-failed]')), + 'a cancel must not produce an [inference-failed] marker'); + + const aborted = traces.filter((t) => t.type === 'inference:aborted') as Array<{ reason?: string }>; + assert.equal(aborted.length, 1, `expected exactly one inference:aborted, got ${JSON.stringify(aborted)}`); + assert.equal(aborted[0].reason, 'user'); + assert.ok(!traces.some((t) => t.type === 'inference:exhausted'), + 'a cancel must not emit inference:exhausted (feeds streak + ops alerts)'); + assert.equal(agent.state.status, 'idle', 'the agent settles to idle'); + } finally { + await teardown(); + } + }); + + it('framework.abortInference(reason) mid-stream → ONE inference:aborted carrying the caller\'s reason', async () => { + const { framework, traces, teardown } = await bootHung('interrupt-abort-'); + try { + framework.pushEvent({ type: 'external-message', source: 'test', content: 'go', metadata: {} }); + framework.start(); + const agent = framework.getAgent('assistant')!; + await waitFor(() => agent.state.status === 'waiting_for_tools'); + + // A host-side abort with its own provenance (zombie reclaim, subagent + // cancel, operator). Before: abortInference emitted inference:aborted + // with this reason AND driveStream emitted a second one as 'user'. + assert.equal(framework.abortInference('assistant', 'operator_reclaim'), true); + + await waitFor(() => contextTexts(framework).some((t) => t.includes('[turn-interrupted]'))); + // Give a second (wrong) trace every chance to show up. + await new Promise((r) => setTimeout(r, 100)); + + const aborted = traces.filter((t) => t.type === 'inference:aborted') as Array<{ reason?: string }>; + assert.deepEqual(aborted.map((t) => t.reason), ['operator_reclaim'], + `expected exactly one inference:aborted with the caller's reason, got ${JSON.stringify(aborted)}`); + assert.ok(!traces.some((t) => t.type === 'inference:exhausted')); + + // The marker's metadata carries the reason; its text stays neutral. + const { messages } = agent.getContextManager().queryMessages({}); + const markerMsg = messages.find((m) => + m.content.some((b) => b.type === 'text' && b.text.includes('[turn-interrupted]')))!; + assert.equal((markerMsg.metadata as { reason?: string } | undefined)?.reason, 'operator_reclaim'); + const markerText = contextTexts(framework).find((t) => t.includes('[turn-interrupted]'))!; + assert.doesNotMatch(markerText, /operator_reclaim|by the user/); + } finally { + await teardown(); + } + }); + + it('a provider-side abort (reason error) still goes through the failure pipeline', async () => { + const { framework, membrane, traces, teardown } = await bootHung('interrupt-real-'); + try { + framework.pushEvent({ type: 'external-message', source: 'test', content: 'go', metadata: {} }); + framework.start(); + const agent = framework.getAgent('assistant')!; + await waitFor(() => agent.state.status === 'waiting_for_tools'); + + // Membrane's abort reasons are 'user' | 'timeout' | 'error'; only + // 'user' is a cancel. Emit a provider-side one directly on the live + // mock stream. + const stream = membrane.lastStream!; + (stream as unknown as { events: unknown[] }).events.push({ type: 'aborted', reason: 'error' }); + const pr = (stream as unknown as { pendingResolve: (() => void) | null }).pendingResolve; + if (pr) { (stream as unknown as { pendingResolve: null }).pendingResolve = null; pr(); } + + await waitFor(() => traces.some((t) => t.type === 'inference:exhausted')); + const exhausted = traces.find((t) => t.type === 'inference:exhausted') as { error?: string }; + assert.match(exhausted?.error ?? '', /Stream aborted: error/); + assert.ok(!traces.some((t) => t.type === 'inference:aborted'), + 'a provider abort is not a deliberate cancellation'); + assert.ok(!contextTexts(framework).some((t) => t.includes('[turn-interrupted]'))); + } finally { + await teardown(); + } + }); +});