diff --git a/README.md b/README.md index 92fe15c8..853793ee 100644 --- a/README.md +++ b/README.md @@ -53,6 +53,11 @@ await server.start(); An agent wraps an LLM identity: model, system prompt, context strategy, and tool permissions. Multiple agents can coexist, each with independent context and inference state. +Agents can opt into the [tool result guard](docs/tool-result-guard.md) through +`agent_settings` with `{"action":"update","tool_result_guard":true}`. The +setting persists across restarts; withheld results remain recoverable in +Chronicle's audit history. + ```typescript { name: 'researcher', diff --git a/changelog.d/tool-result-guard.added.md b/changelog.d/tool-result-guard.added.md new file mode 100644 index 00000000..895427ed --- /dev/null +++ b/changelog.d/tool-result-guard.added.md @@ -0,0 +1,8 @@ +- Add the opt-in, durable `agent_settings.tool_result_guard` setting + (programmatic `AgentConfig.toolResultGuard`; hosts must forward it from + recipes). On a provider refusal after tool output, withhold the + latest result batch and retry inference once without rerunning tools or + automatically rewinding older messages. Full originals remain in an + append-only Chronicle audit log (synced before submission); pending output + stays out of speculative compression, guard effects apply only to a batch + actually submitted in the current turn, and disabling the guard does not restore withheld results. diff --git a/docs/tool-result-guard.md b/docs/tool-result-guard.md new file mode 100644 index 00000000..1e447e36 --- /dev/null +++ b/docs/tool-result-guard.md @@ -0,0 +1,141 @@ +# Tool result guard + +The tool result guard is off by default. An agent can enable it through +`agent_settings`: + +```json +{"action":"update","tool_result_guard":true} +``` + +Use `{"action":"get"}` to inspect the effective boolean and its source. +The setting persists across turns and process restarts. Set it to `false` +to disable it, or explicitly reset `tool_result_guard` to restore the recipe +default. Programmatic hosts can set `AgentConfig.toolResultGuard: true`. +A host that builds `AgentConfig` from a recipe must forward that key itself; +connectome-host does not forward it yet, so in that host a recipe +`toolResultGuard` line is currently a silent no-op — use `agent_settings`. +Disabling the setting never restores previously withheld output. + +## Behavior + +The guard applies to the batch of results returned together by the agent's +latest tool round, including text, images, and errors. It recognizes a +structured provider `stopReason: 'refusal'`, not keywords in output, ordinary +provider errors, or natural-language refusal text. + +On that signal, every result in the batch is withheld. Tool calls, result IDs, +and error flags remain intact; each entire payload becomes: + +> Tool result withheld by the guard. The tool has already executed. + +The refused attempt's partial assistant output is discarded. Inference is +retried once on the same model, within the same logical turn; executed tools +are never automatically run again. Refusal details stay in operational logs, +outside the agent-facing notice and setting description. + +The guard takes precedence over Membrane's unchanged-input refusal retries +for a pending batch and its recovery attempt. A second refusal stops this +recovery without automatic rewind of older exchanges or human messages. +Explicit operator `/unstick` remains a separate action. A later clean tool +round can stage a new batch with its own single recovery allowance. + +Guard effects apply only to a batch that was actually **submitted** to the +provider in the current turn. A batch whose turn ends before submission +(an `endTurn` tool such as `end_turn`/`skip_reply`, or results that arrive +after the response completed) is admitted immediately, recorded as +`accepted` with `reason: "turn_ended"`. A batch that compilation omitted +from the request is recorded as `withheld` with `reason: "unsubmitted"` when +a refusal arrives; that refusal then goes through ordinary refusal handling +(`refusalHandling.retries`/`autoRewind`), since the guarded output was never +on the wire. + +While a submitted batch is pending, compilation reserves the originals' real +wire cost (chars/4 plus a flat per-image cost) in `reserveForResponse`, so +substituting them for the placeholders cannot push a restart's request past +its budget. + +While a submitted batch is pending (or a recovery is running), streamed text +and its `inference:tokens` trace events are held at one publication +boundary; they are released in order on a clean round and discarded on a +guard refusal. + +Normal successful rounds admit the preceding results to memory. This works +for framework yielding streams (including ephemeral agents and context-budget +restarts) and the backward-compatible direct `Agent.runInference` API. + +## Chronicle and memory + +Withholding is non-destructive. Before submitting a guarded batch, the host +appends a `staged` record to the Chronicle append-log state +`framework/tool-result-guard`. It contains the full `originals` (including +pre-truncation data, error strings, and image bytes), the serialized history +`content`, and the `wireResults`. A `linked` record connects its `batchId` to +the context message's `messageId`; later `accepted` or `withheld` records +record the outcome. No guard operation deletes or overwrites these records. +Payload fields larger than 10 KB use Chronicle blobs (`{blobId}`), following +the inference log convention; resolve them with `store.getBlob(blobId)` and +parse the JSON. The append-log snapshots retain the blob references. + +The context manager initially receives only placeholders, so speculative +compression cannot incorporate output that is subsequently withheld. The +converse is a known limit: if compression summarizes the placeholder +*before* acceptance (e.g. deferred messages flushed behind the batch push it +out of the protected tail), the later acceptance edit does not invalidate +that summary, and the accepted output is absent from compressed memory. +Closing this needs a context-manager hook (a transient compression hold +for pending messages, or edit-aware invalidation of derived entries). Raw +pending output goes directly to the provider. A clean following response +promotes the history payload through Chronicle's versioned message-edit API. +On a refusal, the placeholders remain. The audit slot is not a context or +compression source. + +The host calls `store.sync()` after the `staged`/`linked` records and before +the originals can be submitted. A failed sync fails closed: it is logged +(`[tool-result-guard] ... audit sync failed`), the provider receives the +placeholders instead of the originals, and the batch settles as `withheld` +(`reason: "unsubmitted"`). Output whose audit is not durable never reaches a +provider. + +A stream that ends without a clean round and without a successor (abort, +exhausted error retries; also an aborted direct `Agent.runInference`) +settles its batch as `withheld` with `reason: "aborted"`, so a later turn +neither resubmits the originals nor treats its own refusal as the guard's. + +If the process stops before acceptance, placeholders remain after restart; +the full pending originals are still available in the audit. This is +deliberately conservative: an interrupted submission does not establish that +the output was accepted. Similarly, if context compilation omits a pending +exchange, its unsubmitted payload is not admitted to memory. + +For example, an operator can inspect records without changing the agent's +view: + +```ts +const store = framework.getStore(); +const records = store.getStateJson('framework/tool-result-guard'); +// For large logs, use getStateLen/getStateItemJson instead of loading all. +const historical = store.getStateJsonAt('framework/tool-result-guard', sequence); +``` + +This is reactive recovery, not pre-submission screening: the provider sees +the original batch once before returning the signal. The setting is not +retroactive and does not rewrite results already accepted into memory. + +## Known limits + +- **Compression before acceptance** — see above; needs a context-manager + change. +- **Explicit-send suppression across a budget restart.** A guard recovery + carries the same-turn send suppression, so a text-only recovery after a + successful explicit send is not routed (locus/hybrid). Context-budget + restarts still re-enter without it (inherited, unchanged here). +- **Audit growth.** Every guarded batch, accepted ones included, archives + `originals`, `content`, and `wireResults` indefinitely. There is no + retention policy and no restore tool; large payloads are blobs, but a full + `getStateJson` read materializes the whole log. +- **Token stats.** The context manager's token-stats cache is not + write-through on edit, so a stats read taken while a batch was pending can + keep pricing that message as the placeholder. +- **Held stream text on abort.** Text held while a submitted batch is + pending is dropped (not previewed) if the stream aborts or errors before + the round settles. diff --git a/src/agent.ts b/src/agent.ts index 0855aafd..488ef28c 100644 --- a/src/agent.ts +++ b/src/agent.ts @@ -2,6 +2,7 @@ import type { Membrane, NormalizedMessage, NormalizedRequest, ContentBlock, Yiel import { isAbortedResponse } from '@animalabs/membrane'; import { createHash } from 'node:crypto'; import type { CacheWireReceipt, KvUnifiedRequestHooks } from './kv-unified-wire.js'; +import { ToolResultGuard, TOOL_RESULT_GUARD_NOTICE } from './tool-result-guard.js'; import { toolResultDataToHistoryString, truncateForHistory, @@ -82,6 +83,7 @@ export class Agent { readonly thinking: AgentConfig['thinking']; /** Refusal auto-rewind policy (see AgentConfig.refusalHandling). */ readonly refusalHandling: AgentConfig['refusalHandling']; + readonly toolResultGuard: ToolResultGuard; /** Prose delivery mode (see AgentConfig.proseRouting). Default 'locus'. */ readonly proseRouting: 'locus' | 'explicit' | 'hybrid' | 'disabled'; /** Exact whole-response known-tool wrapper containment (default off). */ @@ -143,6 +145,7 @@ export class Agent { this.temperature = config.temperature; this.thinking = config.thinking; this.refusalHandling = config.refusalHandling; + this.toolResultGuard = new ToolResultGuard(config.name, contextManager, config.toolResultGuard); this.proseRouting = config.proseRouting ?? 'locus'; this.toolWrapperProseGuard = config.toolWrapperProseGuard ?? false; this.cacheTtl = config.cacheTtl ?? '1h'; @@ -225,6 +228,20 @@ export class Agent { return { maxTokens: this.contextBudgetTokens, reserveForResponse: this.maxTokens }; } + /** Compile budget with the tool-result guard's reservation applied. While + * a batch is pending the strategy selects against short placeholders, but + * prepareRequest substitutes the originals on the wire; reserve their real + * cost so the final request still fits (a budget restart that recompiles + * small must not re-inflate past the same window). */ + private compileBudget(budget?: TokenBudget): TokenBudget | undefined { + const resolved = this.resolveBudget(budget); + const reserve = this.toolResultGuard.pendingWireReserveTokens; + if (reserve === 0) return resolved; + // Mirrors ContextManager.compile's own default when no budget is set. + const base = resolved ?? { maxTokens: 100_000, reserveForResponse: 4_000 }; + return { ...base, reserveForResponse: base.reserveForResponse + reserve }; + } + /** Structural compatibility keeps lightweight test/host ContextManager * doubles working while making the new capability optional at runtime. */ private getHotContextSettings(): HotContextSettingsStatus | null { @@ -523,7 +540,7 @@ export class Agent { } async compileContext(budget?: TokenBudget): Promise { - const result = await this.contextManager.compile(this.resolveBudget(budget)); + const result = await this.contextManager.compile(this.compileBudget(budget)); if (!budget) this.settleRuntimeSettingsTransition(); return result; } @@ -539,7 +556,7 @@ export class Agent { opts?: { kvUnifiedImmutablePrefixHash?: string }, ): Promise { const result = await this.contextManager.compile( - this.resolveBudget(budget), injections, opts as never, + this.compileBudget(budget), injections, opts as never, ); if (!budget) this.settleRuntimeSettingsTransition(); return result; @@ -582,17 +599,28 @@ export class Agent { // Filter tools to only allowed ones const tools = availableTools.filter((t) => this.canUseTool(t.name)); + // The direct Agent API uses the same admission rules as framework-driven + // streams. Stage before compiling so compression only sees placeholders. + const guardedResults = this._state.status === 'ready' && this.toolResultGuard.enabled; + if (guardedResults && this._state.status === 'ready') { + const content = this.buildToolResultMessages(this._state.toolResults)[0].content; + const wireResults = content.flatMap((block) => block.type === 'tool_result' + // buildToolResultMessages serializes every payload to a string. + ? [{ toolUseId: block.toolUseId, content: block.content as string, isError: block.isError }] : []); + this.toolResultGuard.storeResults(content, wireResults, this._state.toolResults); + } + // Compile context (with optional injections) const { messages, systemInjections } = await this.compileWithInjections(budget, injections); // If we have pending tool results, add them - if (this._state.status === 'ready') { + if (this._state.status === 'ready' && !guardedResults) { const toolResultMessages = this.buildToolResultMessages(this._state.toolResults); messages.push(...toolResultMessages); } const request: NormalizedRequest = { - messages, + messages: this.toolResultGuard.prepareRequest(messages, true), system: this.buildSystemPrompt(systemInjections), config: { model: this.model, @@ -648,6 +676,7 @@ export class Agent { } catch (error) { // On error, go back to idle this._state = { status: 'idle' }; + this.toolResultGuard.recovering = false; throw error; } } @@ -773,7 +802,7 @@ export class Agent { } return { - messages, + messages: this.toolResultGuard.prepareRequest(messages), system: this.buildSystemPrompt(systemInjections), config: { model: this.model, @@ -830,6 +859,7 @@ export class Agent { this.lastStreamOutputTokens = 0; const request = await this.buildActivationRequest(availableTools, injections, budget); + request.messages = this.toolResultGuard.prepareRequest(request.messages, true); const receiptAware = (this.contextManager as unknown as { getStrategy?: () => unknown }) .getStrategy?.() as { @@ -854,6 +884,7 @@ export class Agent { }; } + const agent = this; const stream = this.membrane.streamYielding(request, { emitTokens: true, emitBlocks: false, @@ -863,7 +894,13 @@ export class Agent { // (cache-warm), where a framework-level requeue would recompile and // land a different window. The framework's driveStream handles the // resulting `retrying` event by discarding the abandoned attempt. - ...(this.refusalHandling?.retries ? { refusalRetries: this.refusalHandling.retries } : {}), + // Membrane reads this per physical round, including after settings + // tools run. A guarded result must reach the host on its FIRST refusal; + // never spend retries resubmitting unchanged guarded content. + get refusalRetries() { + return agent.toolResultGuard.hasSubmittedPending || agent.toolResultGuard.recovering + ? 0 : agent.refusalHandling?.retries ?? 0; + }, }); this._state = { status: 'streaming', stream }; @@ -1041,9 +1078,28 @@ export class Agent { request: NormalizedRequest, signal?: AbortSignal ): Promise { - const response = await this.membrane.stream(request, { signal }); + let response = await this.membrane.stream(request, { signal }); + // Usage of a refused attempt abandoned by the guard; it was billed and + // is added to the retry's usage below rather than overwritten. + let abandonedUsage: { inputTokens: number; outputTokens: number } | undefined; + if (!isAbortedResponse(response) && response.stopReason === 'refusal') { + // A staged batch never on the wire cannot have caused the refusal. + this.toolResultGuard.settleUnsubmitted('unknown'); + const ids = this.toolResultGuard.withhold('unknown'); + if (ids) { + abandonedUsage = response.usage; + const withheld = new Set(ids); + request = { ...request, messages: request.messages.map((message) => ({ + ...message, content: message.content.map((block) => block.type === 'tool_result' && withheld.has(block.toolUseId) + ? { type: 'tool_result', toolUseId: block.toolUseId, content: TOOL_RESULT_GUARD_NOTICE, isError: block.isError } : block), + })) }; + response = await this.membrane.stream(request, { signal }); + } + } if (isAbortedResponse(response)) { + // Interrupted before a clean response: settle, never leave pending. + this.toolResultGuard.abandon('aborted'); const partialContent = response.partialContent ?? []; const { toolCalls, speechContent } = this.extractToolCallsAndSpeech(partialContent); return { @@ -1056,16 +1112,24 @@ export class Agent { }; } - const { toolCalls, speechContent } = this.extractToolCallsAndSpeech(response.content); + const guardedRefusal = response.stopReason === 'refusal' && this.toolResultGuard.recovering; + if (response.stopReason !== 'refusal') this.toolResultGuard.accept(); + this.toolResultGuard.recovering = false; + const content = guardedRefusal ? [] : response.content; + const { toolCalls, speechContent } = this.extractToolCallsAndSpeech(content); // Add assistant response to context - this.contextManager.addMessage(this.name, response.content); + if (!guardedRefusal) this.contextManager.addMessage(this.name, content); return { toolCalls, speechContent, raw: response.raw, - usage: response.usage, + usage: abandonedUsage && response.usage ? { + ...response.usage, + inputTokens: response.usage.inputTokens + abandonedUsage.inputTokens, + outputTokens: response.usage.outputTokens + abandonedUsage.outputTokens, + } : response.usage, stopReason: response.stopReason, }; } diff --git a/src/framework.ts b/src/framework.ts index d918a022..1a4a888e 100644 --- a/src/framework.ts +++ b/src/framework.ts @@ -438,7 +438,8 @@ function withDeferredWriteId(metadata: MessageMetadata | undefined, id: string): } function isTurnContinuation(reason: string): boolean { - return reason === 'context_budget_restart' || reason === 'tool_results_ready'; + return reason === 'context_budget_restart' || reason === 'tool_results_ready' + || reason === 'tool_result_guard_retry'; } const CONVERSATION_ROUTER_STATE_ID = 'framework/conversation-router'; const INFERENCE_LOG_ID = 'framework/inference-log'; @@ -1303,6 +1304,8 @@ export class AgentFramework { // Session-level token usage tracking (always-on) private usageTracker: UsageTracker; + /** Explicit-send suppression carried into a tool-result-guard retry. */ + private guardRetryTurnSilenced = new Map(); /** Presentation-only wall-clock zone; persistence remains UTC/epoch. */ private readonly timeZone: string; @@ -2398,6 +2401,11 @@ export class AgentFramework { ext.keys.forEach((k) => taken.add(k)); result.set('_framework', ext); } + { + const ext = this.toolResultGuardSettingsExtension(); + ext.keys.forEach((key) => taken.add(key)); + result.set('_toolResultGuard', ext); + } for (const module of this.moduleRegistry.getAllModules()) { const ext = module.getAgentSettingsExtension?.(); if (!ext) continue; @@ -2468,6 +2476,56 @@ export class AgentFramework { }; } + private toolResultGuardSettingsExtension(): AgentSettingsExtension { + const get = (name: string): Record => { + const guard = this.agents.get(name)?.toolResultGuard; + return { + tool_result_guard: guard?.enabled ?? false, + tool_result_guard_source: guard?.settingOverride !== undefined ? 'runtime_override' : 'recipe_default', + }; + }; + const set = (name: string, enabled: boolean | undefined): Record => { + const agent = this.agents.get(name); + if (!agent) throw new Error(`Unknown agent: ${name}`); + const data = this.store.getStateJson(FRAMEWORK_STATE_ID); + const state = (data && typeof data === 'object' ? data : {}) as Record; + const values = { ...((state.toolResultGuards as Record | undefined) ?? {}) }; + if (enabled === undefined) delete values[name]; + else values[name] = enabled; + state.toolResultGuards = values; + this.store.setStateJson(FRAMEWORK_STATE_ID, state); + agent.toolResultGuard.setOverride(enabled); + return get(name); + }; + return { + properties: { + tool_result_guard: { + type: 'boolean', + description: 'Enable the tool result guard. It can withhold a batch of tool output and retry inference once. ' + + 'Tools are not re-executed. Persists across turns and restarts until explicitly changed. ' + + 'Disabling affects future results only; previously withheld results stay withheld.', + }, + }, + keys: ['tool_result_guard'], + get, + update: (name, patch) => { + if (typeof patch.tool_result_guard !== 'boolean') throw new Error('tool_result_guard must be a boolean'); + return set(name, patch.tool_result_guard); + }, + reset: (name) => set(name, undefined), + }; + } + + private restoreToolResultGuardSetting(agent: Agent): void { + const enabled = (this.store.getStateJson(FRAMEWORK_STATE_ID) as { + toolResultGuards?: Record; + } | null)?.toolResultGuards?.[agent.name]; + if (enabled !== undefined && typeof enabled !== 'boolean') { + throw new Error(`Invalid persisted tool_result_guard for ${agent.name}: expected a boolean`); + } + agent.toolResultGuard.setOverride(enabled); + } + private buildThinkTool( policy: SameRoundThinkTextPolicy, ): import('./types/index.js').ToolDefinition { @@ -3849,6 +3907,7 @@ export class AgentFramework { }); const agent = new Agent(config, contextManager, this.membrane); + this.restoreToolResultGuardSetting(agent); this.ephemeralCandidates.set(agent, contextManager); const cleanup = () => { @@ -6153,6 +6212,7 @@ export class AgentFramework { }); const agent = new Agent(config, contextManager, this.membrane); + this.restoreToolResultGuardSetting(agent); const restoredSettings = this.readAgentRuntimeSettings(config.name); if (restoredSettings) { agent.restoreRuntimeSettings( @@ -6227,6 +6287,7 @@ export class AgentFramework { allowedTools: [...SUBCONSCIOUS_TOOL_NAMES, 'think', 'skip_reply', 'end_turn'], }; const agent = new Agent(agentConfig, contextManager, this.membrane); + this.restoreToolResultGuardSetting(agent); this.agents.set(name, agent); this.agentConfigs.set(name, agentConfig); this.subconsciousAgentName = name; @@ -6348,7 +6409,10 @@ export class AgentFramework { // divergence breaks the compile prefix). const { blocks: toolResultContent, spilled } = await this.buildStoredToolResultContent(agent.name, currentState.toolResults, maxChars); - agent.getContextManager().addMessage('user', toolResultContent); + const membraneResults = currentState.toolResults.map(tc => + this.toMembraneToolResult(tc.id, tc.result, maxChars, spilled.get(tc.id)) + ); + agent.toolResultGuard.storeResults(toolResultContent, membraneResults, currentState.toolResults); // Flush any messages that were deferred while this turn was in // flight. Route to the PRIMARY agent — deferred messages are @@ -6531,6 +6595,9 @@ export class AgentFramework { if (shouldEndTurn) { // endTurn: messages already stored above, cancel stream, reset to idle. + // The batch is never submitted, so nothing can refuse it: admit + // it now instead of leaving it pending across the idle gap. + agent.toolResultGuard.settleTurnEnded(); if (currentState.stream) { this.frameworkCancelledStreams.set(`${agent.name}:${agent.streamId}`, 'turn_ended'); currentState.stream.cancel(); @@ -6586,11 +6653,8 @@ export class AgentFramework { // Mid-turn messages collected above ride along as injected user // messages (membrane ≥0.5.72) — appended after the tool_result // envelope so the next round of THIS turn hears them. - const membraneResults = currentState.toolResults.map(tc => - this.toMembraneToolResult(tc.id, tc.result, maxChars, spilled.get(tc.id)) - ); currentState.stream.provideToolResults( - membraneResults, + agent.toolResultGuard.submissionResults(membraneResults), midTurnInjections.length > 0 ? { injectedMessages: midTurnInjections } : undefined, ); agent.setStreaming(currentState.stream); @@ -6999,6 +7063,7 @@ export class AgentFramework { const config: AgentConfig = { ...templateConfig, name, strategy: undefined }; const agent = new Agent(config, contextManager, this.membrane); + this.restoreToolResultGuardSetting(agent); this.agents.set(name, agent); this.agentConfigs.set(name, config); this.conversationAgentHomes.set(name, channelId); @@ -8140,6 +8205,8 @@ export class AgentFramework { turnToken: number, ownsProviderGate: boolean, ): Promise { + const continuingTurn = trigger?.reason === 'context_budget_restart' + || trigger?.reason === 'tool_result_guard_retry'; // Flush messages deferred during the PREVIOUS turn — before the // checkpoint, the locus announcement, and the compile — so a turn started // by a queued wake actually CONTAINS the message that woke it. (2026-07-31 @@ -8156,7 +8223,7 @@ export class AgentFramework { // products — undoing the turn must not destroy them. if ( attempt === 0 && - trigger?.reason !== 'context_budget_restart' && + !continuingTurn && this.deferredMessages.length > 0 ) { { @@ -8183,7 +8250,7 @@ export class AgentFramework { } // Record turn checkpoint before inference (only on first attempt, not retries) - if (attempt === 0) { + if (attempt === 0 && trigger?.reason !== 'tool_result_guard_retry') { this.recordTurnCheckpoint(agent.name); this.redoStacks.delete(agent.name); // new work invalidates redo } @@ -8201,7 +8268,7 @@ export class AgentFramework { // after the restart fell through to the hours-old defaultPublishChannel. // Keep the turn's trigger channel across restarts; only real new turns // reset it. - if (trigger?.reason !== 'context_budget_restart') { + if (!continuingTurn) { if (trigger?.channelId) { this.activeTriggerChannels.set(agent.name, trigger.channelId); } else { @@ -8221,12 +8288,12 @@ export class AgentFramework { // BEFORE this turn compiles, so the agent always knows where its voice // goes (announce-on-change only — no per-turn chatter, append-only for // KV stability). - if (trigger?.reason === 'context_budget_restart') { + if (continuingTurn) { const previousLogicalToolState = this.logicalTurnToolCalls.get(agent); this.logicalTurnToolCalls.set(agent, { turnToken, count: previousLogicalToolState?.count ?? 0 }); } - if (trigger?.reason !== 'context_budget_restart') { + if (!continuingTurn) { if (attempt === 0) this.maybePrimeProseMode(agent); const previousLogicalToolState = this.logicalTurnToolCalls.get(agent); if (attempt === 0) { @@ -8271,7 +8338,7 @@ export class AgentFramework { }); // A budget restart continues the same logical inference window. The // predecessor keeps EventGate liveness until its successor terminates. - if (trigger?.reason !== 'context_budget_restart') { + if (!continuingTurn) { this.eventGate?.onInferenceStarted(agent.name); } this.lastInferenceAt.set(agent.name, { ...this.lastInferenceAt.get(agent.name), startedAt: Date.now() }); @@ -8523,7 +8590,9 @@ export class AgentFramework { // Sticky explicit-send suppression: prose after send_message stays quiet // to prevent a redundant "sent it" postscript. Fresh injected input // clears it, because the following prose is a reply to a new message. - let turnSilenced = false; + let turnSilenced = trigger?.reason === 'tool_result_guard_retry' + && this.guardRetryTurnSilenced.get(agent.name) === true; + this.guardRetryTurnSilenced.delete(agent.name); // Live routing is only trusted when the membrane provides verbatim // round-scoped blocks (roundContent, native tool mode, membrane ≥0.5.64). @@ -8605,6 +8674,23 @@ export class AgentFramework { } }; + // Tool-result-guard publication boundary. While a SUBMITTED batch is + // pending (or a guard recovery is running), a refusal can still abandon + // this physical round, so its streamed text must not reach any surface + // yet — neither the prose router nor the inference:tokens trace (TTS and + // other trace consumers voice it). Hold both; release in order on a clean + // round (tool-calls / non-guard completion), discard on a guard refusal. + let guardHeld: Array<{ trace: Parameters[0]; text?: string }> = []; + const releaseGuardHeld = (): void => { + const held = guardHeld; + guardHeld = []; + for (const item of held) { + this.emitTrace(item.trace); + if (proseStream && item.text !== undefined) emitOutgoing(proseStream.feed(item.text)); + } + }; + const discardGuardHeld = (): void => { guardHeld = []; }; + const adoptInjectedRound = (): void => { if (!this.midTurnInputSignals.has(agent.name)) return; this.midTurnInputSignals.delete(agent.name); @@ -8627,8 +8713,8 @@ export class AgentFramework { } this.touchEphemeralRun(agent.name, true); switch (event.type) { - case 'tokens': - this.emitTrace({ + case 'tokens': { + const trace: Parameters[0] = { type: 'inference:tokens', agentName: agent.name, content: event.content, @@ -8638,11 +8724,16 @@ export class AgentFramework { // lets trace consumers tag every chunk with the channel this // turn's prose is bound for, without re-deriving routing. channelId: typingChannel ?? undefined, - }); - if (proseStream && event.meta.type === 'text') { - emitOutgoing(proseStream.feed(event.content)); + }; + const text = event.meta.type === 'text' ? event.content : undefined; + if (agent.toolResultGuard.hasSubmittedPending || agent.toolResultGuard.recovering) { + guardHeld.push({ trace, text }); + } else { + this.emitTrace(trace); + if (proseStream && text !== undefined) emitOutgoing(proseStream.feed(text)); } break; + } case 'retrying': { // Membrane is re-issuing after a content-policy refusal. Per the @@ -8664,6 +8755,7 @@ export class AgentFramework { ); outgoingInferenceId = newOutgoingInferenceId(); outgoingIndex = 0; + discardGuardHeld(); proseStream?.reset(); this.emitTrace({ type: 'inference:tokens', @@ -8690,6 +8782,10 @@ export class AgentFramework { } case 'tool-calls': { + // A tool-call response is a clean physical round: admit its + // preceding results before storing/dispatching this new round. + agent.toolResultGuard.accept(); + releaseGuardHeld(); adoptInjectedRound(); hadToolCalls = true; this.recordLogicalTurnToolCalls(agent, myTurnToken ?? -1, event.calls.length); @@ -8847,7 +8943,86 @@ export class AgentFramework { case 'complete': { adoptInjectedRound(); const durationMs = Date.now() - startTime; - const response = event.response; + let response = event.response; + const guardCategory = (response.raw?.response as { + stop_details?: { category?: string }; + } | undefined)?.stop_details?.category ?? 'unknown'; + // A staged batch the strategy omitted from the request cannot + // have caused this refusal: settle it and fall through to the + // ordinary refusal handling (rewind/reaction) below. + if (response.stopReason === 'refusal') agent.toolResultGuard.settleUnsubmitted(guardCategory); + const guardRefusal = response.stopReason === 'refusal' + && (agent.toolResultGuard.hasSubmittedPending || agent.toolResultGuard.recovering); + if (guardRefusal) { + discardGuardHeld(); + const category = guardCategory; + const withheld = agent.toolResultGuard.withhold(category); + // Nothing from the abandoned physical attempt is published or + // persisted as assistant speech, including a terminal refusal + // after the single recovery retry. + proseStream?.reset(); + if (withheld) { + const usage = response.details?.usage ?? response.usage; + const tokenUsage = usage ? { + input: usage.inputTokens, output: usage.outputTokens, + cacheCreation: usage.cacheCreationTokens, cacheRead: usage.cacheReadTokens, + } : undefined; + this.noteRefusal(agent.name, category, tokenUsage); + // The abandoned stream was billed (membrane usage here is + // cumulative across this physical tool loop): count it before + // the retry opens a fresh stream with its own usage. + const abandonedUsage = response.details?.usage; + if (abandonedUsage) { + this.usageTracker.onInferenceCompleted(agent.name, { + inputTokens: abandonedUsage.inputTokens, + outputTokens: abandonedUsage.outputTokens, + cacheCreationTokens: abandonedUsage.cacheCreationTokens, + cacheReadTokens: abandonedUsage.cacheReadTokens, + }, abandonedUsage.estimatedCost + ? { total: abandonedUsage.estimatedCost.total, currency: abandonedUsage.estimatedCost.currency } + : undefined); + this.persistUsageState(); + } + this.logInference({ + timestamp: startTime, agentName: agent.name, requestId, + success: false, error: 'Tool output withheld by the guard', + request: compiledRequest ?? {}, response, durationMs, tokenUsage, stopReason: 'refusal', + }); + console.error(`[tool-result-guard] agent=${agent.name} withheld ${withheld.length} result(s); retrying inference`); + this.emitTrace({ + type: 'inference:stream_restarted', agentName: agent.name, + reason: 'tool_result_guard', inputTokens: agent.lastStreamInputTokens, + budget: agent.maxStreamTokens, + }); + await turnSpeechChain; + if (this.agents.get(agent.name) !== agent || agent.streamId !== myStreamId) { + generationLost = true; + lifecyclePhase = 'aborted'; + return; + } + agent.reset(); + preserveEventGateForSuccessor = true; + lifecyclePhase = 'aborted'; + // Carry same-turn explicit-send suppression into the retry: + // a text-only recovery after a successful send must not post + // a postscript to a message already delivered. + if (turnSilenced) this.guardRetryTurnSilenced.set(agent.name, true); + // Restart inference inside this logical turn. No tool is + // executed again and no settle/checkpoint/locus reset occurs. + await this.startAgentStream(agent, { + ...trigger, agentName: agent.name, reason: 'tool_result_guard_retry', + source: trigger?.source ?? 'framework', timestamp: Date.now(), + }, attempt); + return; + } + const lastResult = response.content.reduce((last, block, index) => + block.type === 'tool_result' ? index : last, -1); + response = { ...response, content: response.content.slice(0, lastResult + 1), rawAssistantText: '' }; + agent.toolResultGuard.recovering = false; + } else if (response.stopReason !== 'refusal') { + agent.toolResultGuard.accept(); + } + if (!guardRefusal) releaseGuardHeld(); // If the agent is still waiting_for_tools when 'complete' fires // (shouldn't happen after incomplete-tool-call fix, but guard anyway), @@ -8876,12 +9051,17 @@ export class AgentFramework { // one door a giant blob still walks through). const readyState = agent.state as AgentState; if (readyState.status === 'ready') { - const { blocks: toolResultContent } = await this.buildStoredToolResultContent( + const cap = this.resolveToolResultInlineCap(agent).cap; + const { blocks: toolResultContent, spilled } = await this.buildStoredToolResultContent( agent.name, readyState.toolResults, - this.resolveToolResultInlineCap(agent).cap, + cap, ); - agent.getContextManager().addMessage('user', toolResultContent); + agent.toolResultGuard.storeResults(toolResultContent, readyState.toolResults.map((tc) => + this.toMembraneToolResult(tc.id, tc.result, cap, spilled.get(tc.id))), readyState.toolResults); + // The response is already complete: this batch is never + // submitted in this turn, so it cannot be refused. Admit it. + agent.toolResultGuard.settleTurnEnded(); } } @@ -9097,7 +9277,7 @@ export class AgentFramework { }, () => this.finishUnstick(agent.name, false, category), ); - } else if (rh?.autoRewind) { + } else if (rh?.autoRewind && !guardRefusal) { const cap = Math.max(1, rh.maxRewinds ?? 3); const used = this.refusalRewinds.get(agent.name) ?? 0; doRewind( @@ -9179,6 +9359,12 @@ export class AgentFramework { if (agent.proseRouting === 'disabled') { console.error(`[routing] ${agent.name}: text-only prose NOT routed (proseRouting=disabled)`); this.recordProseSuppression(agent.name, 1); + } else if (turnSilenced && agent.proseRouting !== 'explicit') { + // Only reachable on a guard recovery that carried the + // same-turn send suppression (a fresh text-only turn starts + // unsilenced). Explicit mode is exempt, as mid-turn. + console.error(`[routing] ${agent.name}: text-only recovery prose NOT routed (turn silenced)`); + this.recordProseSuppression(agent.name, 1); } else if (agent.proseRouting === 'hybrid') { const locus = resolveTurnLocus(); console.error(`[prose] ${agent.name}: text-only turn -> hybrid prose gateway`); @@ -9719,6 +9905,15 @@ export class AgentFramework { // (typing still stops, compression still runs — matching the observed // wedge). onInferenceEnded is idempotent, so a redundant call is safe. if (ownsPhysicalStream && !preserveEventGateForSuccessor) { + // A frame that ends without handing off to a successor (abort, + // exhausted errors) must not leave its batch pending: a later turn + // would resubmit the originals and claim its own refusal. Successor + // frames (error-policy retry, budget/guard restart) bump streamId or + // set preserveEventGateForSuccessor and keep the batch. + if (agent.streamId === myStreamId && agent.toolResultGuard.hasPending) { + agent.toolResultGuard.abandon('aborted'); + } + agent.toolResultGuard.recovering = false; this.eventGate?.onInferenceEnded(agent.name); } if (!generationLost && ownsPhysicalStream) { diff --git a/src/tool-result-guard.ts b/src/tool-result-guard.ts new file mode 100644 index 00000000..5c692aa6 --- /dev/null +++ b/src/tool-result-guard.ts @@ -0,0 +1,225 @@ +import { randomUUID } from 'node:crypto'; +import type { ContextManager, MessageId } from '@animalabs/context-manager'; +import type { ContentBlock, NormalizedMessage, ToolResult } from '@animalabs/membrane'; +import type { CompletedToolCall } from './types/index.js'; +import { isStateExistsError } from './module-registry.js'; + +export const TOOL_RESULT_GUARD_NOTICE = 'Tool result withheld by the guard. The tool has already executed.'; +export const TOOL_RESULT_GUARD_AUDIT_STATE = 'framework/tool-result-guard'; + +interface PendingBatch { + id: string; + messageId: MessageId; + content: ContentBlock[]; + wireResults: ToolResult[]; + submitted: boolean; + /** False when the audit sync failed: originals then never go on the wire. */ + durable: boolean; +} + +/** + * Admission of newly returned tool output to durable model-facing memory. + * + * Raw output is appended to a separate Chronicle audit slot BEFORE the + * placeholder enters the context manager. In particular, onNewMessage and + * speculative compression never see unaccepted output. Acceptance edits the + * placeholder through CM's versioned edit API; withholding only appends an + * audit event. Neither operation erases the original output or its blobs. + */ +export class ToolResultGuard { + private pending: PendingBatch | undefined; + private registered = false; + private override: boolean | undefined; + /** True until a recovery produces a clean response/new tool round. */ + recovering = false; + + constructor( + private readonly agentName: string, + private readonly cm: ContextManager, + private readonly configured = false, + ) {} + + get enabled(): boolean { return this.override ?? this.configured; } + setOverride(value: boolean | undefined): void { this.override = value; } + get settingOverride(): boolean | undefined { return this.override; } + get hasPending(): boolean { return this.pending !== undefined; } + /** True only once the pending batch's originals were actually put on the + * wire. Guard effects (retry suppression, refusal claim, prose buffering) + * are scoped to this state; a merely staged batch never claims a refusal. */ + get hasSubmittedPending(): boolean { return this.pending?.submitted === true; } + + /** Extra input tokens the pending batch costs on the wire beyond the + * placeholder the strategy selected against (chars/4 + flat per image, + * the same heuristic as the physical-window projection). Compilation + * reserves this so substituting originals cannot exceed the budget. */ + get pendingWireReserveTokens(): number { + const pending = this.pending; + if (!pending) return 0; + let chars = 0; + let images = 0; + for (const result of pending.wireResults) { + if (typeof result.content === 'string') chars += result.content.length; + else for (const block of result.content as ContentBlock[]) { + if (block.type === 'image') images += 1; + else chars += JSON.stringify(block).length; + } + chars -= TOOL_RESULT_GUARD_NOTICE.length; + } + return Math.max(0, Math.ceil(chars / 4) + images * 1600); + } + + private append(record: Record): void { + const store = this.cm.getStore(); + if (!this.registered) { + try { + store.registerState({ id: TOOL_RESULT_GUARD_AUDIT_STATE, strategy: 'append_log' }); + } catch (error) { + if (!isStateExistsError(error)) throw error; + } + this.registered = true; + } + store.appendToStateJson(TOOL_RESULT_GUARD_AUDIT_STATE, { + agentName: this.agentName, timestamp: Date.now(), ...record, + }); + } + + private archive(value: unknown): unknown { + const json = JSON.stringify(value); + // Match inference-log storage: large payloads (especially images and + // pre-spill output) must not be copied into every append-log snapshot. + return json.length > 10_000 + ? { blobId: this.cm.getStore().storeBlob(Buffer.from(json), 'application/json') } + : JSON.parse(json); + } + + storeResults(content: ContentBlock[], wireResults: ToolResult[], originals: CompletedToolCall[]): MessageId { + if (!this.enabled) return this.cm.addMessage('user', content); + if (this.pending) { + // Idempotent for the same not-yet-submitted batch: a caller that + // retries after a failure between staging and submission (e.g. a + // transient compile error on the direct Agent API) must not wedge. + const same = !this.pending.submitted + && this.pending.wireResults.map((r) => r.toolUseId).join('\0') + === wireResults.map((r) => r.toolUseId).join('\0'); + if (same) return this.pending.messageId; + throw new Error('Tool result guard already has a pending batch'); + } + const id = randomUUID(); + // Includes full pre-truncation/error/image payloads, not just the wire + // preview. This slot is audit data, never a context/compression source. + this.append({ type: 'staged', batchId: id, + originals: this.archive(originals), content: this.archive(content), wireResults: this.archive(wireResults) }); + const withheld: ContentBlock[] = content.map((block) => block.type === 'tool_result' + ? { type: 'tool_result', toolUseId: block.toolUseId, content: TOOL_RESULT_GUARD_NOTICE, isError: block.isError } + : block); + const messageId = this.cm.addMessage('user', withheld); + this.pending = { id, messageId, content, wireResults, submitted: false, durable: false }; + this.append({ type: 'linked', batchId: id, messageId }); + // Durability barrier: the audit must reach Chronicle's chain heads before + // the originals can go to a provider. On a failed sync the batch fails + // CLOSED: the placeholders go on the wire instead (the turn continues, + // nothing is stranded), and the batch later settles as unsubmitted. + try { + this.cm.getStore().sync(); + this.pending.durable = true; + } catch (error) { + console.error(`[tool-result-guard] agent=${this.agentName} audit sync failed; ` + + 'submitting placeholders instead of originals:', error); + } + return messageId; + } + + /** Results for a live continuation (provideToolResults). Marks the batch + * submitted and returns the originals only when the audit is durable; + * otherwise returns placeholder results and leaves it unsubmitted. */ + submissionResults(results: ToolResult[]): ToolResult[] { + const pending = this.pending; + if (!pending) return results; + if (pending.durable) { pending.submitted = true; return results; } + const ids = new Set(pending.wireResults.map((result) => result.toolUseId)); + return results.map((result) => ids.has(result.toolUseId) + ? { ...result, content: TOOL_RESULT_GUARD_NOTICE } : result); + } + + /** The stream carrying this batch ended without a clean round (abort or + * exhausted errors). Settle conservatively — an interrupted submission + * does not establish acceptance — so a later turn neither resubmits the + * originals nor attributes its own refusal to this batch. */ + abandon(reason: string): void { + const pending = this.pending; + this.recovering = false; + if (!pending) return; + this.pending = undefined; + this.append({ type: 'withheld', batchId: pending.id, messageId: pending.messageId, + reason, submitted: pending.submitted }); + } + + /** A budget/error restart compiles placeholders; restore pending output + * only in this provider request, never in the strategy's view. */ + prepareRequest(messages: NormalizedMessage[], recordSubmission = false): NormalizedMessage[] { + const pending = this.pending; + if (!pending || !pending.durable) return messages; + const byId = new Map(pending.wireResults.map((result) => [result.toolUseId, result])); + const present = new Set(messages.flatMap((message) => message.content + .filter((block) => block.type === 'tool_result' && byId.has(block.toolUseId)) + .map((block) => (block as ContentBlock & { toolUseId: string }).toolUseId))); + // A strategy may have folded the entire exchange away. Do not release + // content that was never submitted. Its originals remain in the audit. + const submitted = present.size === byId.size; + if (recordSubmission) pending.submitted = submitted; + if (!submitted) return messages; + return messages.map((message) => ({ ...message, content: message.content.map((block) => { + const result = block.type === 'tool_result' ? byId.get(block.toolUseId) : undefined; + return result ? { type: 'tool_result', toolUseId: result.toolUseId, content: result.content, isError: result.isError } : block; + }) })); + } + + /** The turn ended (endTurn/skip_reply) before the batch was submitted: + * nothing was refused, so admit it exactly as an unguarded agent would. + * Leaving it pending would turn it into a permanent withheld notice on + * restart and disable ordinary refusal handling on the next turn. */ + settleTurnEnded(): void { + const pending = this.pending; + if (!pending) return; + this.pending = undefined; + this.append({ type: 'accepted', batchId: pending.id, messageId: pending.messageId, reason: 'turn_ended' }); + this.cm.editMessage(pending.messageId, pending.content); + } + + /** A refusal arrived while the pending batch was never on the wire (the + * strategy omitted the exchange). The guard cannot have caused it: record + * the batch as withheld, clear it, and let ordinary refusal handling run. + * Returns true when it settled such a batch. */ + settleUnsubmitted(category: string): boolean { + const pending = this.pending; + if (!pending || pending.submitted) return false; + this.pending = undefined; + this.append({ type: 'withheld', batchId: pending.id, messageId: pending.messageId, reason: 'unsubmitted', category }); + return true; + } + + /** A clean physical response accepts precisely the last submitted batch. */ + accept(): void { + const pending = this.pending; + if (pending) { + this.append({ type: pending.submitted ? 'accepted' : 'withheld', batchId: pending.id, messageId: pending.messageId, + ...(pending.submitted ? {} : { reason: 'unsubmitted' }) }); + if (pending.submitted) this.cm.editMessage(pending.messageId, pending.content); + this.pending = undefined; + } + this.recovering = false; + } + + /** At most one recovery per batch; no scanning/deleting older history. */ + withhold(category: string): string[] | null { + const pending = this.pending; + if (!pending?.submitted) return null; + const ids = pending.wireResults.map((result) => result.toolUseId); + // Even a failed outcome-log write must never re-arm rejected output for + // a later submission. Its originals were archived before admission. + this.pending = undefined; + this.recovering = true; + this.append({ type: 'withheld', batchId: pending.id, messageId: pending.messageId, toolUseIds: ids, category }); + return ids; + } +} diff --git a/src/types/agent.ts b/src/types/agent.ts index b5cd7471..e7009cc2 100644 --- a/src/types/agent.ts +++ b/src/types/agent.ts @@ -170,6 +170,12 @@ export interface AgentConfig { announceHumanTurns?: boolean; }; + /** Opt-in tool-output admission and one guarded retry after a provider + * refusal. Agents can persistently override this via agent_settings + * tool_result_guard. Full originals remain in Chronicle's audit history. + * Default false. Disabling does not restore previously withheld output. */ + toolResultGuard?: boolean; + /** * How the agent's PLAIN PROSE (non-tool output) reaches channels. * - 'locus' (default): host-inferred — the turn-frozen locus machinery. diff --git a/test/tool-result-guard.test.ts b/test/tool-result-guard.test.ts new file mode 100644 index 00000000..dbd082af --- /dev/null +++ b/test/tool-result-guard.test.ts @@ -0,0 +1,600 @@ +import { afterEach, test } 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 { PassthroughStrategy, type StoredMessage, type StrategyContext } from '@animalabs/context-manager'; +import type { ContentBlock, Membrane, NormalizedRequest, NormalizedResponse, YieldingStreamOptions } from '@animalabs/membrane'; +import { Membrane as RealMembrane, NativeFormatter, type ProviderAdapter, type ProviderRequest, + type ProviderResponse, type StreamCallbacks } from '@animalabs/membrane'; +import { AgentFramework, type AgentConfig, type AgentSettingsExtension, type Module, type ModuleContext, + type ToolCall, type ToolResult, type ProcessEvent, type ProcessState } from '../src/index.js'; +import { TOOL_RESULT_GUARD_AUDIT_STATE, TOOL_RESULT_GUARD_NOTICE } from '../src/tool-result-guard.js'; +import { MockYieldingStream, createMockResponse } from './helpers/mock-membrane.js'; + +const dirs: string[] = []; +afterEach(() => { while (dirs.length) rmSync(dirs.pop()!, { recursive: true, force: true }); }); +const answer = () => createMockResponse([{ type: 'text', text: 'continued' }]); +const refused = () => ({ + ...createMockResponse([{ type: 'text', text: 'discard-this-partial-output' }], 'refusal'), + raw: { request: {}, response: { stop_details: { category: 'test-category' } } }, +}) as NormalizedResponse; +const calls = (...ids: string[]) => createMockResponse(ids.map((id) => ({ + type: 'tool_use', id, name: 'test--read', input: {}, +})), 'tool_use'); + +class ScriptMembrane { + requests: NormalizedRequest[] = []; + streams: MockYieldingStream[] = []; + retriesAtSubmission: number[] = []; + onSubmit?: () => void; + constructor(readonly scripts: NormalizedResponse[][]) {} + streamYielding(request: NormalizedRequest, options: YieldingStreamOptions = {}) { + this.requests.push(structuredClone({ ...request, onCacheWireReceipt: undefined })); + const script = this.scripts.shift(); + assert.ok(script, 'unexpected extra inference'); + const stream = new MockYieldingStream(script); + const provide = stream.provideToolResults.bind(stream); + stream.provideToolResults = (...args) => { + this.retriesAtSubmission.push(options.refusalRetries ?? 0); + this.onSubmit?.(); + provide(...args); + }; + this.streams.push(stream); + return stream; + } + asMembrane() { return this as unknown as Membrane; } +} + +class ReadModule implements Module { + readonly name = 'test'; + readonly calls: string[] = []; + readonly speeches: string[] = []; + constructor(readonly results: Record = {}) {} + async start(ctx: ModuleContext) { ctx.registerSpeechHandler('*'); } + async stop() {} + getTools() { + return [ + { name: 'read', description: 'Read a result', inputSchema: { type: 'object' as const, properties: {} } }, + { name: 'send_message', description: 'Explicit send', inputSchema: { type: 'object' as const, properties: {} } }, + ]; + } + async handleToolCall(call: ToolCall): Promise { + this.calls.push(call.id); + return this.results[call.id] ?? { success: true, data: `payload-${call.id}` }; + } + async onProcess(event: ProcessEvent, _state: ProcessState) { + return event.type === 'external-message' + ? { addMessages: [{ participant: 'user', content: [{ type: 'text' as const, text: String(event.content) }] }], requestInference: true } + : {}; + } + async onAgentSpeech(_name: string, content: ContentBlock[]) { + this.speeches.push(...content.flatMap((block) => block.type === 'text' ? [block.text] : [])); + } +} + +class IngressObserver extends PassthroughStrategy { + snapshots: string[] = []; + async onNewMessage(_message: StoredMessage, ctx: StrategyContext) { + this.snapshots.push(JSON.stringify(ctx.messageStore.getAll())); + } +} + +async function harness(scripts: NormalizedResponse[][], config: Partial = {}, results?: Record) { + const dir = mkdtempSync(join(tmpdir(), 'af-tool-result-guard-')); dirs.push(dir); + const membrane = new ScriptMembrane(scripts); + const module = new ReadModule(results); + const base = { storePath: join(dir, 'store'), membrane: membrane.asMembrane(), + agents: [{ name: 'assistant', model: 'test', systemPrompt: 'system', ...config }], modules: [module], syncIntervalMs: 0 }; + const framework = await AgentFramework.create(base); + const run = async () => { + framework.pushEvent({ type: 'external-message', source: 'test', content: 'read', metadata: {} }); + await framework.runUntilIdle(); + }; + return { framework, membrane, module, base, run }; +} + +function extension(framework: AgentFramework): AgentSettingsExtension { + const extensions = (framework as unknown as { + collectAgentSettingsExtensions(): Map; + }).collectAgentSettingsExtensions(); + return [...extensions.values()].find((ext) => ext.keys.includes('tool_result_guard'))!; +} + +function toolResults(framework: AgentFramework) { + return framework.getAgent('assistant')!.getContextManager().getAllMessages() + .flatMap((message) => message.content.filter((block) => block.type === 'tool_result')); +} + +test('default off retains existing refusal behavior and original tool output', async () => { + const h = await harness([[calls('one'), refused()]]); + try { + await h.run(); + assert.equal(h.membrane.requests.length, 1); + assert.match(JSON.stringify(toolResults(h.framework)), /payload-one/); + assert.equal(extension(h.framework).get('assistant').tool_result_guard, false); + } finally { await h.framework.stop(); } +}); + +test('withholds the entire latest batch, retries inference once, and keeps originals in Chronicle', async () => { + const strategy = new IngressObserver(); + const image = Buffer.from('original-image-bytes').toString('base64'); + const h = await harness([[calls('text', 'image', 'error'), refused()], [answer()]], + { toolResultGuard: true, strategy, refusalHandling: { retries: 3 } }, { + text: { success: true, data: 'original-text-payload' }, + image: { success: true, data: [{ type: 'image', mimeType: 'image/png', data: image }] }, + error: { success: false, isError: true, error: 'original-error-payload' }, + }); + let stagedSequence = 0; + let originalStopped = false; + h.membrane.onSubmit = () => { stagedSequence = h.framework.getStore().currentSequence(); }; + try { + await h.run(); + assert.deepEqual(h.module.calls.sort(), ['error', 'image', 'text']); + assert.equal(h.membrane.requests.length, 2); + assert.deepEqual(h.membrane.retriesAtSubmission, [0], 'first refusal must reach guard before plain retries'); + const retry = h.membrane.requests[1]; + assert.equal(retry.config.model, 'test'); + assert.doesNotMatch(JSON.stringify(retry), /original-text-payload|original-error-payload|discard-this-partial-output|test-category/); + assert.ok(!JSON.stringify(retry).includes(image)); + assert.equal(retry.messages.flatMap((m) => m.content).filter((b) => b.type === 'tool_use').length, 3); + const guarded = toolResults(h.framework); + assert.equal(guarded.length, 3); + assert.ok(guarded.every((block) => block.content === TOOL_RESULT_GUARD_NOTICE)); + assert.equal(guarded.find((b) => b.toolUseId === 'error')?.isError, true); + assert.deepEqual(h.module.speeches, ['continued']); + assert.ok(strategy.snapshots.every((s) => !s.includes('original-text-payload') && !s.includes(image)), + 'background strategy ingress must never see withheld payloads'); + + const store = h.framework.getStore(); + const audit = store.getStateJson(TOOL_RESULT_GUARD_AUDIT_STATE) as Array>; + assert.match(JSON.stringify(audit[0]), /original-text-payload|original-error-payload/); + assert.ok(JSON.stringify(audit[0]).includes(image)); + assert.equal(audit.at(-1)?.type, 'withheld'); + const historical = store.getStateJsonAt(TOOL_RESULT_GUARD_AUDIT_STATE, stagedSequence) as unknown[]; + assert.deepEqual(audit[0], historical[0], 'redaction only appends; original Chronicle record is unchanged'); + const original = structuredClone(audit[0]); + extension(h.framework).update('assistant', { tool_result_guard: false }); + assert.ok(toolResults(h.framework).every((block) => block.content === TOOL_RESULT_GUARD_NOTICE)); + await h.framework.stop(); + originalStopped = true; + const restarted = await AgentFramework.create(h.base); + try { + assert.equal(extension(restarted).get('assistant').tool_result_guard, false, 'explicit disable persists'); + assert.ok(toolResults(restarted).every((block) => block.content === TOOL_RESULT_GUARD_NOTICE)); + assert.deepEqual((restarted.getStore().getStateJson(TOOL_RESULT_GUARD_AUDIT_STATE) as unknown[])[0], original); + const preview = await restarted.previewActivation('assistant'); + assert.doesNotMatch(JSON.stringify(preview), /original-text-payload|original-error-payload/); + } finally { await restarted.stop(); } + } finally { if (!originalStopped) await h.framework.stop(); } +}); + +test('successful physical rounds admit results; refusal affects only the newest batch', async () => { + const h = await harness([[calls('accepted'), calls('withheld'), refused()], [answer()]], { toolResultGuard: true }); + try { + await h.run(); + const results = toolResults(h.framework); + assert.match(String(results.find((b) => b.toolUseId === 'accepted')?.content), /payload-accepted/); + assert.equal(results.find((b) => b.toolUseId === 'withheld')?.content, TOOL_RESULT_GUARD_NOTICE); + assert.deepEqual(h.module.calls, ['accepted', 'withheld']); + assert.equal(h.framework.getAgent('assistant')!.toolResultGuard.enabled, true); + } finally { await h.framework.stop(); } +}); + +test('a second refusal stops recovery without auto-rewinding older or human messages', async () => { + const h = await harness([[calls('one'), refused()], [refused()]], + { toolResultGuard: true, refusalHandling: { autoRewind: true, retries: 3 } }); + try { + await h.run(); + assert.equal(h.membrane.requests.length, 2); + assert.deepEqual(h.module.calls, ['one']); + assert.deepEqual(h.module.speeches, []); + const messages = h.framework.getAgent('assistant')!.getContextManager().getAllMessages(); + assert.ok(messages.some((m) => m.content.some((b) => b.type === 'text' && b.text === 'read'))); + assert.doesNotMatch(JSON.stringify(messages), /discard-this-partial-output|refusal-rewind/); + assert.equal(toolResults(h.framework)[0].content, TOOL_RESULT_GUARD_NOTICE); + } finally { await h.framework.stop(); } +}); + +test('durable typed agent setting can be enabled, disabled, and explicitly reset', async () => { + const h = await harness([]); + let originalStopped = false; + try { + const ext = extension(h.framework); + for (const value of ['true', 1, null, {}]) { + assert.throws(() => ext.update('assistant', { tool_result_guard: value }), /must be a boolean/); + } + ext.update('assistant', { tool_result_guard: true }); + assert.equal(ext.get('assistant').tool_result_guard, true); + await h.framework.stop(); + originalStopped = true; + const restarted = await AgentFramework.create(h.base); + try { + const restored = extension(restarted); + assert.equal(restored.get('assistant').tool_result_guard, true); + assert.equal(restored.get('assistant').tool_result_guard_source, 'runtime_override'); + restored.reset!('assistant'); + assert.equal(restored.get('assistant').tool_result_guard, false); + const tool = restarted.getAllTools().find((t) => t.name === 'agent_settings')!; + assert.equal((tool.inputSchema.properties as Record).tool_result_guard.type, 'boolean'); + assert.doesNotMatch(JSON.stringify(tool), /classifier/i); + } finally { await restarted.stop(); } + } finally { if (!originalStopped) await h.framework.stop(); } +}); + +test('a normal successful response releases the pending output into versioned history', async () => { + const h = await harness([[calls('one'), answer()]], { toolResultGuard: true }); + try { + await h.run(); + assert.equal(h.membrane.requests.length, 1); + assert.match(String(toolResults(h.framework)[0].content), /payload-one/); + const records = h.framework.getStore().getStateJson(TOOL_RESULT_GUARD_AUDIT_STATE) as Array>; + assert.equal(records.at(-1)?.type, 'accepted'); + assert.equal(h.framework.getAgent('assistant')!.toolResultGuard.enabled, true); + } finally { await h.framework.stop(); } +}); + +test('refusal without a new tool result does not invoke tool guard recovery', async () => { + const h = await harness([[refused()]], { toolResultGuard: true }); + try { await h.run(); assert.equal(h.membrane.requests.length, 1); } + finally { await h.framework.stop(); } +}); + +test('budget restart submits staged originals, then recovers without re-executing tools', async () => { + const h = await harness([[calls('one')], [refused()], [answer()]], { toolResultGuard: true, maxStreamTokens: 1 }); + try { + await h.run(); + assert.equal(h.membrane.requests.length, 3); + assert.match(JSON.stringify(h.membrane.requests[1]), /payload-one/); + assert.doesNotMatch(JSON.stringify(h.membrane.requests[2]), /payload-one/); + assert.deepEqual(h.module.calls, ['one']); + assert.equal(toolResults(h.framework)[0].content, TOOL_RESULT_GUARD_NOTICE); + } finally { await h.framework.stop(); } +}); + +test('ephemeral run settles only after recovery and counts each tool once', async () => { + const h = await harness([[calls('one'), refused()], [answer()]]); + try { + const created = await h.framework.createEphemeralAgent({ + name: 'ephemeral', model: 'test', systemPrompt: 'system', toolResultGuard: true, + }); + created.contextManager.addMessage('user', [{ type: 'text', text: 'read' }]); + const completion = h.framework.runEphemeralToCompletion(created.agent, created.contextManager); + h.framework.start(); + const result = await completion; + assert.deepEqual(result, { speech: 'continued', toolCallsCount: 1 }); + assert.deepEqual(h.module.calls, ['one']); + assert.equal(h.membrane.requests.length, 2); + } finally { await h.framework.stop(); } +}); + +test('quiesced scheduler admits a queued guard recovery as a continuation', async () => { + const h = await harness([[answer()]], { toolResultGuard: true }); + try { + await h.framework.quiesce(); + h.framework.getAgent('assistant')!.getContextManager().addMessage('user', [{ + type: 'text', text: TOOL_RESULT_GUARD_NOTICE, + }]); + // A recovery can be requeued while waiting for provider admission. It + // must finish the held turn even after the host stops admitting new work. + (h.framework as unknown as { pendingRequests: Array> }).pendingRequests.push({ + agentName: 'assistant', reason: 'tool_result_guard_retry', source: 'framework', timestamp: Date.now(), + }); + await h.framework.runUntilIdle(); + assert.equal(h.membrane.requests.length, 1); + assert.equal(h.framework.getHostModeStatus().quiesced, true); + } finally { await h.framework.stop(); } +}); + +test('native Membrane observes the first refusal even when guard is enabled by a tool mid-stream', async () => { + const h = await harness([]); + await h.framework.stop(); + const requests: ProviderRequest[] = []; + const adapter: ProviderAdapter = { + name: 'test', usageCacheConvention: 'cache-excluded', supportsModel: () => true, + complete: async () => { throw new Error('unexpected complete'); }, + stream: async (request: ProviderRequest, callbacks: StreamCallbacks): Promise => { + requests.push(structuredClone(request)); + const index = requests.length; + assert.ok(index <= 3, 'must not retry unchanged refused input'); + const content = index === 1 ? [ + { type: 'tool_use', id: 'enable', name: 'agent_settings', input: { action: 'update', tool_result_guard: true } }, + { type: 'tool_use', id: 'one', name: 'test--read', input: {} }, + ] : [{ type: 'text', text: index === 2 ? 'discard-native-partial' : 'continued' }]; + if (index > 1) callbacks.onChunk?.(index === 2 ? 'discard-native-partial' : 'continued'); + return { + content, stopReason: index === 1 ? 'tool_use' : index === 2 ? 'refusal' : 'end_turn', + usage: { inputTokens: 20, outputTokens: 5 }, model: 'test', + raw: { response: { stop_details: { category: 'test-category' } } }, + } as ProviderResponse; + }, + }; + const framework = await AgentFramework.create({ ...h.base, + agents: [{ name: 'assistant', model: 'test', systemPrompt: 'system', refusalHandling: { retries: 4 } }], + membrane: new RealMembrane(adapter, { formatter: new NativeFormatter() }), + }); + const routed: string[] = []; + const outgoing: string[] = []; + (framework as unknown as { channelRegistry: unknown }).channelRegistry = new Proxy({ + resolveLocus: () => 'world:test', + routeSpeech: async (_agent: string, speech: string) => { + routed.push(speech); return { delivered: true, channelId: 'world:test' }; + }, + sendOutgoingChunk: (_channel: string, _agent: string, _id: string, _index: number, delta: string) => { outgoing.push(delta); }, + getDefaultPublishChannel: () => null, isChannelOpen: () => true, + getDescriptor: () => undefined, getChannelTools: () => [], + }, { get: (target, key: string) => key in target ? (target as Record)[key] : () => undefined }); + const tokenTraces: string[] = []; + framework.onTrace((event) => { + if (event.type === 'inference:tokens') tokenTraces.push(String((event as { content?: unknown }).content)); + }); + try { + framework.pushEvent({ type: 'external-message', source: 'test', content: 'read', metadata: {} }); + await framework.runUntilIdle(); + assert.equal(requests.length, 3); + assert.match(JSON.stringify(requests[1]), /payload-one/); + assert.doesNotMatch(JSON.stringify(requests[2]), /payload-one|discard-native-partial|test-category/); + assert.match(JSON.stringify(requests[2]), /Tool result withheld by the guard/); + assert.deepEqual(h.module.calls, ['one']); + assert.equal(extension(framework).get('assistant').tool_result_guard, true); + assert.deepEqual(h.module.speeches, ['continued']); + assert.deepEqual(routed, ['continued']); + assert.doesNotMatch(outgoing.join(''), /discard-native-partial/); + assert.match(outgoing.join(''), /continued/, 'accepted answer must reach outgoing-stream consumers'); + assert.doesNotMatch(tokenTraces.join(''), /discard-native-partial/, 'inference:tokens must not leak refused text'); + assert.match(tokenTraces.join(''), /continued/); + } finally { await framework.stop(); } +}); + +test('backward-compatible direct Agent inference also guards tool results', async () => { + const h = await harness([], { toolResultGuard: true }); + const responses = [calls('one'), refused(), answer()]; + const requests: NormalizedRequest[] = []; + (h.membrane as unknown as { stream: (request: NormalizedRequest) => Promise }).stream = async (request) => { + requests.push(structuredClone(request)); + const response = responses.shift(); + assert.ok(response); + return response; + }; + try { + const agent = h.framework.getAgent('assistant')!; + agent.getContextManager().addMessage('user', [{ type: 'text', text: 'read' }]); + const first = await agent.runInference(h.framework.getAllTools()); + assert.equal(first.toolCalls.length, 1); + agent.provideToolResult('one', { success: true, data: 'direct-original' }); + const final = await agent.runInference(h.framework.getAllTools()); + assert.deepEqual(final.speechContent, [{ type: 'text', text: 'continued' }]); + assert.equal(requests.length, 3); + assert.equal(final.usage?.inputTokens, 20, 'abandoned refused attempt usage is included'); + assert.match(JSON.stringify(requests[1]), /direct-original/); + assert.doesNotMatch(JSON.stringify(requests[2]), /direct-original|discard-this/); + assert.equal(toolResults(h.framework)[0].content, TOOL_RESULT_GUARD_NOTICE); + } finally { await h.framework.stop(); } +}); + +test('full oversized output survives withholding and reopening as a Chronicle blob', async () => { + const original = 'original-large-'.repeat(8_000) + 'end-of-original'; + const h = await harness([[calls('large'), refused()], [answer()]], { toolResultGuard: true }, { + large: { success: true, data: original }, + }); + let originalStopped = false; + try { + await h.run(); + const audit = h.framework.getStore().getStateJson(TOOL_RESULT_GUARD_AUDIT_STATE) as Array<{ + originals: { blobId: string }; + }>; + const blobId = audit[0].originals.blobId; + assert.equal(typeof blobId, 'string'); + const blob = h.framework.getStore().getBlob(blobId)!; + assert.equal(JSON.parse(blob.toString())[0].result.data, original, 'keep pre-truncation bytes'); + await h.framework.stop(); originalStopped = true; + const restarted = await AgentFramework.create(h.base); + try { + assert.deepEqual(restarted.getStore().getBlob(blobId), blob); + assert.equal(toolResults(restarted)[0].content, TOOL_RESULT_GUARD_NOTICE); + } finally { await restarted.stop(); } + } finally { if (!originalStopped) await h.framework.stop(); } +}); + +// --------------------------------------------------------------------------- +// Review regressions (PR #159): guard effects apply only to a batch that was +// actually SUBMITTED to the provider in the current turn. +// --------------------------------------------------------------------------- + +class OmitToolExchange extends PassthroughStrategy { + select(...args: Parameters) { + return super.select(...args).filter((entry) => + !entry.content.some((block) => block.type === 'tool_use' || block.type === 'tool_result')); + } +} + +test('a staged batch omitted by compilation does not claim the refusal; autoRewind proceeds', async () => { + const h = await harness([[calls('one')], [refused()], [answer()]], { + toolResultGuard: true, maxStreamTokens: 1, strategy: new OmitToolExchange(), + refusalHandling: { autoRewind: true }, + }); + try { + await h.run(); + const guard = h.framework.getAgent('assistant')!.toolResultGuard; + assert.equal(guard.hasPending, false, 'unsubmitted batch must be settled, not stranded'); + assert.equal(guard.recovering, false); + assert.doesNotMatch(JSON.stringify(h.membrane.requests[1]), /payload-one/); + assert.equal(h.membrane.requests.length, 3, 'ordinary refusal rewind retry must run'); + const messages = h.framework.getAgent('assistant')!.getContextManager().getAllMessages(); + assert.match(JSON.stringify(messages), /\[refusal-rewind\]/); + const audit = h.framework.getStore().getStateJson(TOOL_RESULT_GUARD_AUDIT_STATE) as Array>; + assert.ok(audit.some((r) => r.type === 'withheld' && r.reason === 'unsubmitted')); + } finally { await h.framework.stop(); } +}); + +test('a turn ended by an endTurn tool settles (accepts) its batch', async () => { + const h = await harness([[calls('one')]], { toolResultGuard: true }, + { one: { success: true, data: 'payload-one', endTurn: true } }); + try { + await h.run(); + const guard = h.framework.getAgent('assistant')!.toolResultGuard; + assert.equal(guard.hasPending, false); + assert.match(String(toolResults(h.framework)[0].content), /payload-one/); + const audit = h.framework.getStore().getStateJson(TOOL_RESULT_GUARD_AUDIT_STATE) as Array>; + assert.equal(audit.at(-1)?.type, 'accepted'); + assert.equal(audit.at(-1)?.reason, 'turn_ended'); + } finally { await h.framework.stop(); } +}); + +test('a refusal on the next turn after endTurn uses ordinary refusal handling', async () => { + const h = await harness([[calls('one')], [refused()], [answer()]], { + toolResultGuard: true, refusalHandling: { autoRewind: true, retries: 2 }, + }, { one: { success: true, data: 'payload-one', endTurn: true } }); + try { + await h.run(); + await h.run(); + const messages = h.framework.getAgent('assistant')!.getContextManager().getAllMessages(); + assert.match(JSON.stringify(messages), /\[refusal-rewind\]/, 'autoRewind must not be suppressed'); + assert.equal(h.membrane.requests.length, 3); + } finally { await h.framework.stop(); } +}); + +test('budget restart reserves the real wire cost of the pending batch', async () => { + const big = 'x'.repeat(40_000); // ~10k tokens on the wire, notice is ~20 + const budget = 14_000; + const h = await harness([[calls('one')], [answer()]], { + toolResultGuard: true, maxStreamTokens: 1, contextBudgetTokens: budget, maxTokens: 1_000, + }, { one: { success: true, data: big } }); + try { + const cm = h.framework.getAgent('assistant')!.getContextManager(); + for (let i = 0; i < 20; i++) { + cm.addMessage('user', [{ type: 'text', text: `filler-${i} ` + 'y'.repeat(2_000) }]); + } + await h.run(); + const rebuilt = h.membrane.requests[1]; + assert.match(JSON.stringify(rebuilt), /xxxxxxxx/, 'pending original still submitted'); + const chars = rebuilt.messages.flatMap((m) => m.content).reduce((n, b) => + n + JSON.stringify(b).length, 0); + assert.ok(Math.ceil(chars / 4) <= budget - 1_000, + `rebuilt request ~${Math.ceil(chars / 4)} tokens exceeds budget ${budget - 1_000}`); + } finally { await h.framework.stop(); } +}); + +test('audit is synced before the pending batch is submitted', async () => { + const h = await harness([[calls('one'), answer()]], { toolResultGuard: true }); + const store = h.framework.getStore(); + const sync = store.sync.bind(store); + let syncedWhilePending = false; + (store as { sync: () => void }).sync = () => { + if (h.framework.getAgent('assistant')!.toolResultGuard.hasPending) syncedWhilePending = true; + sync(); + }; + let durableAtSubmit = false; + h.membrane.onSubmit = () => { durableAtSubmit = syncedWhilePending; }; + try { + await h.run(); + assert.equal(durableAtSubmit, true); + } finally { (store as { sync: () => void }).sync = sync; await h.framework.stop(); } +}); + +test('direct API: a transient compile failure does not wedge the guard', async () => { + const h = await harness([], { toolResultGuard: true }); + const responses = [calls('one'), answer()]; + (h.membrane as unknown as { stream: (request: NormalizedRequest) => Promise }).stream = + async () => responses.shift()!; + try { + const agent = h.framework.getAgent('assistant')!; + agent.getContextManager().addMessage('user', [{ type: 'text', text: 'read' }]); + await agent.runInference(h.framework.getAllTools()); + agent.provideToolResult('one', { success: true, data: 'direct-original' }); + const compile = agent.compileWithInjections.bind(agent); + let failed = false; + agent.compileWithInjections = async (...args) => { + if (!failed) { failed = true; throw new Error('transient compile failure'); } + return compile(...args); + }; + await assert.rejects(agent.runInference(h.framework.getAllTools()), /transient compile failure/); + const final = await agent.runInference(h.framework.getAllTools()); + assert.deepEqual(final.speechContent, [{ type: 'text', text: 'continued' }]); + assert.equal(toolResults(h.framework).length, 1, 'batch staged exactly once'); + assert.match(String(toolResults(h.framework)[0].content), /direct-original/); + } finally { await h.framework.stop(); } +}); + +test('usage of an abandoned guarded round is counted in session totals', async () => { + const withUsage = (r: NormalizedResponse, input: number) => + ({ ...r, details: { ...(r.details ?? {}), usage: { inputTokens: input, outputTokens: 1 } } }) as NormalizedResponse; + const h = await harness([[calls('one'), withUsage(refused(), 1_000)], [withUsage(answer(), 7)]], { toolResultGuard: true }); + try { + await h.run(); + const totals = h.framework.getSessionUsage().totals as unknown as Record; + assert.ok(totals.inputTokens >= 1_007, `refused round usage missing: ${JSON.stringify(totals)}`); + } finally { await h.framework.stop(); } +}); + +// --------------------------------------------------------------------------- +// Greptile review regressions (PR #159, head 647f081). +// --------------------------------------------------------------------------- + +function fakeRegistry(framework: AgentFramework) { + const routed: string[] = []; + (framework as unknown as { channelRegistry: unknown }).channelRegistry = new Proxy({ + resolveLocus: () => 'world:test', + routeSpeech: async (_agent: string, speech: string) => { + routed.push(speech); return { delivered: true, channelId: 'world:test' }; + }, + sendOutgoingChunk: () => {}, + getDefaultPublishChannel: () => null, isChannelOpen: () => true, + getDescriptor: () => undefined, getChannelTools: () => [], + }, { get: (target, key: string) => key in target ? (target as Record)[key] : () => undefined }); + return routed; +} + +test('failed audit sync fails closed: originals are never submitted', async () => { + const h = await harness([[calls('one'), answer()]], { toolResultGuard: true }); + const store = h.framework.getStore(); + const sync = store.sync.bind(store); + (store as { sync: () => void }).sync = () => { + if (h.framework.getAgent('assistant')!.toolResultGuard.hasPending) throw new Error('disk full'); + sync(); + }; + try { + await h.run(); + const provided = JSON.stringify(h.membrane.streams[0].receivedToolResults); + assert.doesNotMatch(provided, /payload-one/, 'non-durable audit must not release originals'); + assert.match(provided, /Tool result withheld by the guard/); + assert.equal(toolResults(h.framework)[0].content, TOOL_RESULT_GUARD_NOTICE); + const guard = h.framework.getAgent('assistant')!.toolResultGuard; + assert.equal(guard.hasPending, false); + } finally { (store as { sync: () => void }).sync = sync; await h.framework.stop(); } +}); + +test('an aborted stream settles its submitted batch; the next turn neither resubmits nor claims it', async () => { + const h = await harness([[calls('one')], [refused()], [answer()]], { + toolResultGuard: true, refusalHandling: { autoRewind: true }, + }); + h.membrane.onSubmit = () => { + const stream = h.membrane.streams.at(-1)!; + queueMicrotask(() => stream.cancel()); + }; + try { + await h.run(); + const guard = h.framework.getAgent('assistant')!.toolResultGuard; + assert.equal(guard.hasPending, false, 'aborted batch must not stay pending'); + h.membrane.onSubmit = undefined; + await h.run(); + assert.doesNotMatch(JSON.stringify(h.membrane.requests[1]), /payload-one/, 'old originals must not be resubmitted'); + const messages = h.framework.getAgent('assistant')!.getContextManager().getAllMessages(); + assert.match(JSON.stringify(messages), /\[refusal-rewind\]/, 'ordinary autoRewind must run'); + const audit = h.framework.getStore().getStateJson(TOOL_RESULT_GUARD_AUDIT_STATE) as Array>; + assert.ok(audit.some((r) => r.type === 'withheld' && r.reason === 'aborted')); + } finally { await h.framework.stop(); } +}); + +test('guard recovery keeps same-turn explicit-send suppression', async () => { + const h = await harness([[createMockResponse([ + { type: 'tool_use', id: 'send', name: 'test--send_message', input: {} }, + { type: 'tool_use', id: 'one', name: 'test--read', input: {} }, + ], 'tool_use'), refused()], [createMockResponse([{ type: 'text', text: 'postscript' }])]], { toolResultGuard: true }); + const routed = fakeRegistry(h.framework); + try { + await h.run(); + assert.equal(h.membrane.requests.length, 2); + assert.ok(!routed.some((text) => text.includes('postscript')), `postscript routed: ${JSON.stringify(routed)}`); + } finally { await h.framework.stop(); } +});