Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
11 changes: 11 additions & 0 deletions changelog.d/145-puppet-turn-reservation.fixed.md
Original file line number Diff line number Diff line change
@@ -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.
156 changes: 115 additions & 41 deletions src/framework.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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);
Expand Down Expand Up @@ -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,
Expand All @@ -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)`,
);
}
Expand All @@ -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<ToolResult> {
Expand Down
4 changes: 4 additions & 0 deletions src/types/trace.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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).
Expand Down
129 changes: 128 additions & 1 deletion test/puppet-tool-call.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -43,6 +43,23 @@ function puppetHarness(opts?: {
(framework as unknown as { agents: Map<string, unknown> }).agents =
new Map([['princess', agent]]);
(framework as unknown as { toolImageLedgers: Map<string, unknown> }).toolImageLedgers = new Map();
// Turn-alive machinery the puppet reserves through (#145).
const fw = framework as unknown as Record<string, unknown>;
fw.activeTurnTokens = new Map<string, number>();
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<Record<string, unknown>>, _m?: unknown, opts?: { forAgent?: string }) => {
if ((fw.activeTurnTokens as Map<string, number>).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<string, unknown>).getToolsForAgent =
() => surface.map((name) => ({ name }));
(framework as unknown as Record<string, unknown>).executeToolCall =
Expand All @@ -55,7 +72,27 @@ function puppetHarness(opts?: {
(framework as unknown as Record<string, unknown>).emitTrace =
(e: Record<string, unknown>) => { traces.push(e); };

return { framework, stored, traces, executed };
const turn = fw as unknown as {
activeTurnTokens: Map<string, number>;
deferredMessages: Array<{ participant: string; content: Array<Record<string, unknown>>; 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<void>((r) => { release = r; });
const executed: Array<Record<string, unknown>> = [];
(h.framework as unknown as Record<string, unknown>).executeToolCall =
async (call: Record<string, unknown>) => {
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 () => {
Expand Down Expand Up @@ -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<string, unknown>).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;
}
});
Loading