diff --git a/changelog.d/145-puppet-turn-reservation.fixed.md b/changelog.d/145-puppet-turn-reservation.fixed.md new file mode 100644 index 0000000..b099656 --- /dev/null +++ b/changelog.d/145-puppet-turn-reservation.fixed.md @@ -0,0 +1,11 @@ +- `puppetToolCall` reserves the agent against turn start for its whole + duration (#145). The idle check was a point in time before two awaits + (tool execution, result build); a wake arriving in either gap started a + real turn, and the synthetic `tool_use`/`tool_result` pair then wrote + straight into the window under that turn — the wire-order corruption the + puppet exists to avoid. The puppet now holds the turn-alive marker the + scheduler and the `addMessage` deferral guard already respect, and refuses + when a turn is alive even if status reads idle. The one path that could + still start a turn over the reservation — a wake parked on provider + admission resuming after an auxiliary call — now re-tests turn-alive, gives + admission back and requeues the wake for the scheduler instead. diff --git a/src/framework.ts b/src/framework.ts index 8202319..cac759e 100644 --- a/src/framework.ts +++ b/src/framework.ts @@ -6110,6 +6110,23 @@ export class AgentFramework { this.releasePrimaryProviderGate(agent.name); return; } + // The scheduler's turn-alive test passed BEFORE the park; the wait + // is the one gap in which another turn-alive holder can take the + // agent (puppetToolCall's reservation, #145). Never replace its + // token — the re-entry would set ours over it, and a short turn + // here ends and flushes before the puppeted tool returns, leaving + // its pair queued behind a flush that already happened. Give + // admission back and hand the wake to the scheduler, which keeps + // it queued while the agent is turn-alive and starts it normally + // once the holder releases. + if (this.activeTurnTokens.has(agent.name)) { + this.releasePrimaryProviderGate(agent.name); + this.pendingRequests.push(trigger ?? { + agentName: agent.name, reason: 'provider-admission:requeue', source: 'scheduler', timestamp: Date.now(), + }); + console.error(`[provider-admission] ${agent.name}: turn-alive while parked — wake requeued, not started`); + return; + } await this.startAgentStream(agent, trigger, attempt, true); }).catch((error) => { this.releasePrimaryProviderGate(agent.name); @@ -7764,8 +7781,11 @@ export class AgentFramework { * affordance to a model that cannot find it (older models especially). * * Semantics: - * - Refused unless the agent is idle: puppeting mid-turn would corrupt - * the turn state machine and the live stream's wire ordering. + * - Refused unless the agent is idle with no turn alive: puppeting + * mid-turn would corrupt the turn state machine and the live stream's + * wire ordering. For its own duration the puppet HOLDS the turn-alive + * marker (wakes requeue, cross-turn writers defer), so a wake landing + * during the tool's execution cannot start a turn underneath the pair. * - Refused for tools outside the agent's own surface (canUseTool): the * stored pair must be an act the agent could genuinely have taken. * - The call executes FOR REAL through the shared dispatch (MCPL, @@ -7789,9 +7809,13 @@ export class AgentFramework { ): Promise<{ toolUseId: string; result: ToolResult }> { const agent = this.agents.get(agentName); if (!agent) throw new Error(`Unknown agent: ${agentName}`); - if (agent.state.status !== 'idle') { + // Idle AND no turn alive: status reads 'idle' from dequeue until the + // stream registers, and again while a turn's teardown is pending + // (the scheduler's own busy test, 'idle+turn-alive'). + if (agent.state.status !== 'idle' || this.activeTurnTokens.has(agentName)) { + const shown = agent.state.status === 'idle' ? 'idle+turn-alive' : agent.state.status; throw new Error( - `puppet refused: agent ${agentName} is ${agent.state.status} (requires idle — ` + + `puppet refused: agent ${agentName} is ${shown} (requires idle — ` + `injecting a turn under an active stream corrupts wire ordering)`, ); } @@ -7811,44 +7835,94 @@ export class AgentFramework { for (let i = 0; i < 22; i++) suffix += alphabet[Math.floor(Math.random() * alphabet.length)]; const toolUseId = `toolu_01${suffix}`; - const started = Date.now(); - const result = await this.executeToolCall({ - id: toolUseId, - name: toolName, - input, - callerAgentName: agentName, - }); - const durationMs = Date.now() - started; - - // Store the pair through the same shapes the ordinary path uses. Build - // the result blocks BEFORE storing the tool_use: the spill path awaits, - // and a message arriving during that await must land before the pair, - // never between tool_use and its tool_result. The two addMessage calls - // below are synchronous and adjacent — nothing can interleave. - const { blocks } = await this.buildStoredToolResultContent( - agentName, - [{ id: toolUseId, name: toolName, input, result, durationMs }], - this.resolveToolResultInlineCap(agent).cap, - ); - const cm = agent.getContextManager(); - cm.addMessage(agentName, [ - { type: 'tool_use', id: toolUseId, name: toolName, input } as ContentBlock, - ]); - cm.addMessage('user', blocks); + // Reserve the agent against turn start for the WHOLE operation, before + // the first await (#145). The idle check above is a point in time; the + // tool execution and the result build both await, and a wake arriving + // in either gap used to start a real turn — then the pair below wrote + // straight through the context manager into a turn that never saw it, + // the exact wire-order corruption this method exists to avoid. Holding + // a turn token is the same invariant the scheduler respects (turn-alive + // requeues wakes) and the addMessage guard reads (cross-turn writers + // defer), so for the duration this behaves like a turn with no stream. + // Token-matched release in finally, same leak-proofing as + // startAgentStream: a token nobody clears wedges the agent. + const turnToken = this.nextTurnToken++; + this.activeTurnTokens.set(agentName, turnToken); + let ownedToEnd = false; + try { + const started = Date.now(); + const result = await this.executeToolCall({ + id: toolUseId, + name: toolName, + input, + callerAgentName: agentName, + }); + const durationMs = Date.now() - started; + + // Store the pair through the same shapes the ordinary path uses. Build + // the result blocks BEFORE storing the tool_use: the spill path awaits. + // A message arriving during either await is deferred by the turn + // token (as during a real tool call) and lands AFTER the pair, never + // between tool_use and its tool_result. The two addMessage calls + // below are synchronous and adjacent — nothing can interleave. + const { blocks } = await this.buildStoredToolResultContent( + agentName, + [{ id: toolUseId, name: toolName, input, result, durationMs }], + this.resolveToolResultInlineCap(agent).cap, + ); + const toolUse: ContentBlock[] = [ + { type: 'tool_use', id: toolUseId, name: toolName, input } as ContentBlock, + ]; + // Defensive: no path replaces a live token any more (the provider- + // admission re-entry re-tests turn-alive and requeues instead), but + // if one ever does, the side effect has happened and the pair must + // still be stored — not under that turn: queue it as a unit for the + // next flush (drainDeferredFor keeps order; the tool_result entry is + // never split from its tool_use because both ride the same drain). + // The finally below flushes it itself if that turn is already over. + ownedToEnd = this.activeTurnTokens.get(agentName) === turnToken; + if (ownedToEnd) { + const cm = agent.getContextManager(); + cm.addMessage(agentName, toolUse); + cm.addMessage('user', blocks); + } else { + this.deferredMessages.push( + { participant: agentName, content: toolUse, forAgent: agentName }, + { participant: 'user', content: blocks, forAgent: agentName }, + ); + } - this.emitTrace({ - type: 'puppet:tool-call', - agentName, - toolName, - toolUseId, - isError: !!result.isError, - durationMs, - }); - console.log( - `[puppet] ${agentName}: ${toolName} → ${result.isError ? 'ERROR' : 'ok'} ` + - `in ${durationMs}ms (${toolUseId})`, - ); - return { toolUseId, result }; + this.emitTrace({ + type: 'puppet:tool-call', + agentName, + toolName, + toolUseId, + isError: !!result.isError, + durationMs, + ...(ownedToEnd ? {} : { deferred: true }), + }); + console.log( + `[puppet] ${agentName}: ${toolName} → ${result.isError ? 'ERROR' : 'ok'} ` + + `in ${durationMs}ms (${toolUseId})${ownedToEnd ? '' : ' — deferred behind a live turn'}`, + ); + return { toolUseId, result }; + } finally { + if (this.activeTurnTokens.get(agentName) === turnToken) this.activeTurnTokens.delete(agentName); + // Messages deferred while we held the token land now, after the pair + // — the same end-of-turn flush driveStream performs, under the same + // tool-cycle guard. (No wake is requested: the agent sees the pair, + // and anything that arrived meanwhile, on its next turn.) Keyed on + // "no turn alive", not on token ownership: if a turn did replace us + // and has already ended, its flush ran before our pair was queued — + // nobody else will flush it, so we do. + if (!this.activeTurnTokens.has(agentName) + && this.deferredMessages.length > 0 && this.pendingAssistantBlocks.size === 0) { + for (const msg of this.drainDeferredFor(agentName)) { + this.addMessage(msg.participant, msg.content, msg.metadata, + msg.forAgent ? { forAgent: msg.forAgent } : undefined); + } + } + } } private async executeToolCallFrom(call: ToolCall, origin: ChannelToolOrigin): Promise { diff --git a/src/types/trace.ts b/src/types/trace.ts index 1b1dfd1..9a4c753 100644 --- a/src/types/trace.ts +++ b/src/types/trace.ts @@ -351,6 +351,10 @@ export type TraceEvent = toolUseId: string; isError: boolean; durationMs: number; + /** The pair was queued behind a turn that replaced the puppet's + * reservation (provider-admission re-entry) instead of stored + * directly; it lands at that turn's end flush. */ + deferred?: boolean; }) // MCPL server connection lifecycle (spawn / handshake / reconnect health). diff --git a/test/puppet-tool-call.test.ts b/test/puppet-tool-call.test.ts index 7ffa4fa..6ae001d 100644 --- a/test/puppet-tool-call.test.ts +++ b/test/puppet-tool-call.test.ts @@ -43,6 +43,23 @@ function puppetHarness(opts?: { (framework as unknown as { agents: Map }).agents = new Map([['princess', agent]]); (framework as unknown as { toolImageLedgers: Map }).toolImageLedgers = new Map(); + // Turn-alive machinery the puppet reserves through (#145). + const fw = framework as unknown as Record; + fw.activeTurnTokens = new Map(); + fw.nextTurnToken = 1; + fw.deferredMessages = []; + fw.pendingAssistantBlocks = new Map(); + fw.primaryAgentName = 'princess'; + // Cross-turn writer path (what a channel message goes through): defers + // while a turn is alive, else stores — the real guard, reduced. + fw.addMessage = (participant: string, content: Array>, _m?: unknown, opts?: { forAgent?: string }) => { + if ((fw.activeTurnTokens as Map).has('princess')) { + (fw.deferredMessages as unknown[]).push({ participant, content, forAgent: opts?.forAgent }); + return ''; + } + stored.push({ participant, content }); + return `msg-${stored.length}`; + }; (framework as unknown as Record).getToolsForAgent = () => surface.map((name) => ({ name })); (framework as unknown as Record).executeToolCall = @@ -55,7 +72,27 @@ function puppetHarness(opts?: { (framework as unknown as Record).emitTrace = (e: Record) => { traces.push(e); }; - return { framework, stored, traces, executed }; + const turn = fw as unknown as { + activeTurnTokens: Map; + deferredMessages: Array<{ participant: string; content: Array>; forAgent?: string }>; + nextTurnToken: number; + }; + return { framework, stored, traces, executed, fw: turn }; +} + +/** A harness whose executeToolCall stays open until `release()` is called. */ +function heldHarness() { + const h = puppetHarness(); + let release!: () => void; + const gate = new Promise((r) => { release = r; }); + const executed: Array> = []; + (h.framework as unknown as Record).executeToolCall = + async (call: Record) => { + executed.push(call); + await gate; + return { success: true, data: [{ type: 'text', text: 'done' }], isError: false }; + }; + return { ...h, executed, release }; } test('puppetToolCall executes with agent provenance and stores the pair', async () => { @@ -137,3 +174,93 @@ test('puppetToolCall stores an error result as isError, still paired', async () console.log = quiet; } }); + +// --------------------------------------------------------------------------- +// #145 — the idle guard is a point in time; the reservation must span the +// awaits, or a wake starts a turn underneath the pair. +// --------------------------------------------------------------------------- + +test('puppetToolCall holds the turn-alive marker across tool execution and releases it after the pair is stored', async () => { + const { framework, stored, fw, release, executed } = heldHarness(); + const quiet = console.log; + console.log = () => {}; + try { + assert.equal(fw.activeTurnTokens.has('princess'), false, 'nothing reserved before the call'); + const p = framework.puppetToolCall('princess', 'mcpl--eido--look', {}); + await new Promise((r) => setImmediate(r)); + assert.equal(executed.length, 1, 'tool is executing'); + assert.equal(fw.activeTurnTokens.has('princess'), true, + 'while the tool runs the agent is turn-alive — the scheduler requeues wakes on exactly this'); + release(); + await p; + assert.equal(fw.activeTurnTokens.has('princess'), false, 'released once the pair is stored'); + assert.equal(stored.length, 2, 'the pair, stored directly'); + assert.equal(stored[0].content[0].type, 'tool_use'); + assert.equal(stored[1].content[0].type, 'tool_result'); + } finally { + console.log = quiet; + } +}); + +test('puppetToolCall releases the reservation when the tool throws', async () => { + const { framework, fw } = puppetHarness(); + (framework as unknown as Record).executeToolCall = + async () => { throw new Error('boom'); }; + await assert.rejects(() => framework.puppetToolCall('princess', 'mcpl--eido--look', {}), /boom/); + assert.equal(fw.activeTurnTokens.has('princess'), false, 'no leaked token (the idle+turn-alive wedge)'); +}); + +test('puppetToolCall refuses an idle agent whose turn is still alive (teardown pending)', async () => { + const { framework, stored, fw } = puppetHarness(); + fw.activeTurnTokens.set('princess', 41); + await assert.rejects( + () => framework.puppetToolCall('princess', 'mcpl--eido--look', {}), + /idle\+turn-alive.*requires idle/, + ); + assert.equal(stored.length, 0); + assert.equal(fw.activeTurnTokens.get('princess'), 41, 'someone else\'s token untouched'); +}); + +test('a message arriving while the puppet holds the agent lands AFTER the pair, never between', async () => { + const { framework, stored, fw, release } = heldHarness(); + const quiet = console.log; + console.log = () => {}; + try { + const p = framework.puppetToolCall('princess', 'mcpl--eido--look', {}); + await new Promise((r) => setImmediate(r)); + // a channel message for the agent, mid-execution: the cross-turn writer defers + (framework as unknown as { addMessage: (p: string, c: unknown[]) => unknown }) + .addMessage('user', [{ type: 'text', text: 'hey, you there?' }]); + assert.equal(stored.length, 0, 'deferred, not stored mid-puppet'); + assert.equal(fw.deferredMessages.length, 1); + release(); + await p; + assert.equal(fw.deferredMessages.length, 0, 'flushed at the puppet\'s end, like a turn\'s end'); + assert.deepEqual(stored.map((m) => m.content[0].type), ['tool_use', 'tool_result', 'text'], + 'pair first and adjacent; the deferred message follows'); + } finally { + console.log = quiet; + } +}); + +test('if a turn replaced the reservation mid-flight, the pair is queued as a unit behind it, not written into it', async () => { + const { framework, stored, traces, fw, release } = heldHarness(); + const quiet = console.log; + console.log = () => {}; + try { + const p = framework.puppetToolCall('princess', 'mcpl--eido--look', {}); + await new Promise((r) => setImmediate(r)); + // the provider-admission re-entry: a turn takes the marker without re-testing it + fw.activeTurnTokens.set('princess', 9999); + release(); + await p; + assert.equal(stored.length, 0, 'nothing written under the live turn'); + assert.equal(fw.activeTurnTokens.get('princess'), 9999, 'the turn\'s token is not clobbered'); + assert.deepEqual(fw.deferredMessages.map((m) => [m.forAgent, m.content[0].type]), + [['princess', 'tool_use'], ['princess', 'tool_result']], + 'both halves queued, adjacent, addressed to the agent'); + assert.equal(traces[0].deferred, true, 'trace says so'); + } finally { + console.log = quiet; + } +}); diff --git a/test/puppet-turn-reservation.test.ts b/test/puppet-turn-reservation.test.ts new file mode 100644 index 0000000..0405e2b --- /dev/null +++ b/test/puppet-turn-reservation.test.ts @@ -0,0 +1,156 @@ +/** + * #145 end-to-end: a wake that arrives while a puppeted tool is executing + * must not start a turn. Before the fix the scheduler saw an idle agent with + * no turn alive, started streaming, and the puppet's tool_use/tool_result + * pair then landed inside that turn. + */ +import { describe, it, beforeEach, afterEach } from 'node:test'; +import assert from 'node:assert/strict'; +import { tmpdir } from 'node:os'; +import { mkdtempSync, rmSync } from 'node:fs'; +import { join } from 'node:path'; +import { AgentFramework } from '../src/index.js'; +import type { InferenceRequest } from '../src/index.js'; +import { MockMembrane, createMockResponse } from './helpers/mock-membrane.js'; + +type Internals = { + pendingRequests: InferenceRequest[]; + processInferenceRequests(): Promise; + activeTurnTokens: Map; + executeToolCall: (call: Record) => Promise; + withAuxiliaryAdmission(name: string, run: () => Promise): Promise; + providerGates: Map; + deferredMessages: unknown[]; +}; +const settle = () => new Promise((r) => setTimeout(r, 30)); + +describe('puppetToolCall vs a concurrent wake (#145)', () => { + let tempDir: string; + let membrane: MockMembrane; + + beforeEach(() => { + tempDir = mkdtempSync(join(tmpdir(), 'puppet-reservation-')); + membrane = new MockMembrane(); + }); + afterEach(() => { + rmSync(tempDir, { recursive: true, force: true }); + }); + + it('a wake during a held puppet execution is requeued; the pair lands before the turn it would have corrupted', async () => { + membrane.pushResponse(createMockResponse([{ type: 'text', text: 'ok, I see it' }])); + const framework = await AgentFramework.create({ + storePath: join(tempDir, 'test.chronicle'), + membrane: membrane.asMembrane(), + agents: [{ name: 'scout', model: 'test-model', systemPrompt: 'You are scout.' }], + modules: [], + }); + const i = framework as unknown as Internals; + let release!: () => void; + const gate = new Promise((r) => { release = r; }); + i.executeToolCall = async () => { + await gate; + return { success: true, data: 'settings snapshot', isError: false }; + }; + const quiet = console.log; + console.log = () => {}; + try { + const puppet = framework.puppetToolCall('scout', 'agent_settings', { action: 'get' }); + await new Promise((r) => setImmediate(r)); + assert.equal(i.activeTurnTokens.has('scout'), true, 'reserved while the tool runs'); + + // the race: a wake arrives mid-execution and the scheduler runs + i.pendingRequests.push({ agentName: 'scout', reason: 'mcpl:channel-incoming', source: 'test', timestamp: Date.now() }); + await i.processInferenceRequests(); + assert.equal(membrane.calls.length, 0, 'no turn started under the puppet'); + assert.equal(i.pendingRequests.length, 1, 'the wake is requeued, not dropped'); + + release(); + await puppet; + assert.equal(i.activeTurnTokens.has('scout'), false, 'reservation released'); + + await i.processInferenceRequests(); + await framework.runUntilIdle(); + assert.equal(membrane.calls.length, 1, 'the requeued wake ran once the puppet was done'); + + const all = framework.getAgent('scout')!.getContextManager().getAllMessages() as Array<{ content: Array<{ type: string }> }>; + const types = all.map((m) => m.content[0]?.type); + const use = types.indexOf('tool_use'); + assert.ok(use >= 0, 'pair stored'); + assert.equal(types[use + 1], 'tool_result', 'pair adjacent'); + const turnText = types.lastIndexOf('text'); + assert.ok(turnText > use + 1, 'the turn the wake started comes AFTER the pair, and saw it'); + } finally { + console.log = quiet; + await framework.stop(); + } + }); + + it('a turn parked on provider admission that resumes during a puppet is requeued, not started over the reservation (review of #148)', async () => { + // The one gap the scheduler's busy test does not cover: a wake passed it, + // acquired primary admission, and parked behind an auxiliary call. If the + // auxiliary settles while a puppet holds the agent, the continuation used + // to re-enter startAgentStream and replace the puppet's token; a fast + // replacement turn then ended (and flushed) before the puppeted tool + // returned, and the pair queued behind it was never stored. The public + // promise still resolved "ok". + membrane.pushResponse(createMockResponse([{ type: 'text', text: 'replacement turn' }])); + const framework = await AgentFramework.create({ + storePath: join(tempDir, 'test.chronicle'), + membrane: membrane.asMembrane(), + agents: [{ name: 'scout', model: 'test-model', systemPrompt: 'You are scout.' }], + modules: [], + }); + const i = framework as unknown as Internals; + let releaseTool!: () => void; + const toolGate = new Promise((r) => { releaseTool = r; }); + i.executeToolCall = async () => { await toolGate; return { success: true, data: 'settings snapshot', isError: false }; }; + let releaseAux!: () => void; + const auxiliary = i.withAuxiliaryAdmission('scout', () => new Promise((r) => { releaseAux = r; })); + const quiet = console.log, quietErr = console.error; + console.log = () => {}; console.error = () => {}; + try { + // 1. a wake passes the busy test and parks on provider admission + i.pendingRequests.push({ agentName: 'scout', reason: 'mcpl:channel-incoming', source: 'test', timestamp: Date.now() }); + await i.processInferenceRequests(); + await settle(); + assert.equal(membrane.calls.length, 0, 'parked behind the auxiliary call'); + assert.equal(i.providerGates.get('scout')?.primaryDepth, 1, 'admission held while parked'); + assert.equal(i.activeTurnTokens.has('scout'), false, 'no token yet — the park is before the token'); + + // 2. a puppet takes the agent while the turn is parked + const puppet = framework.puppetToolCall('scout', 'agent_settings', { action: 'get' }); + await settle(); + assert.equal(i.activeTurnTokens.has('scout'), true, 'puppet holds the reservation'); + const puppetToken = i.activeTurnTokens.get('scout'); + + // 3. the auxiliary settles: the parked turn would now resume + releaseAux(); + await auxiliary; + await settle(); + assert.equal(membrane.calls.length, 0, 'no provider call starts during the puppeted tool'); + assert.equal(i.activeTurnTokens.get('scout'), puppetToken, 'the reservation is not replaced'); + assert.equal(i.pendingRequests.length, 1, 'the wake is requeued for the scheduler, not dropped'); + assert.equal(i.providerGates.get('scout')?.primaryDepth, 0, 'admission given back'); + + // 4. the tool returns: the pair is persisted before the promise resolves + releaseTool(); + await puppet; + const stored = () => (framework.getAgent('scout')!.getContextManager().getAllMessages() as Array<{ content: Array<{ type: string }> }>) + .map((m) => m.content[0]?.type); + assert.deepEqual(stored().filter((t) => t === 'tool_use' || t === 'tool_result'), ['tool_use', 'tool_result'], 'pair persisted when puppetToolCall resolves'); + assert.equal(i.deferredMessages.length, 0, 'nothing stranded in the deferred queue'); + assert.equal(i.activeTurnTokens.has('scout'), false, 'reservation released'); + + // 5. the retained wake runs afterwards and sees the pair + await i.processInferenceRequests(); + await framework.runUntilIdle(); + assert.equal(membrane.calls.length, 1, 'the requeued wake ran once the puppet was done'); + const types = stored(); + assert.ok(types.lastIndexOf('text') > types.indexOf('tool_result'), 'the turn comes after the pair'); + } finally { + console.log = quiet; console.error = quietErr; + releaseAux(); releaseTool(); + await framework.stop(); + } + }); +});