Skip to content

Commit 7099921

Browse files
committed
Allow tool-only turns to finalize successfully
1 parent bff7738 commit 7099921

4 files changed

Lines changed: 206 additions & 18 deletions

File tree

src/orchestrator/router.test.ts

Lines changed: 140 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -533,6 +533,10 @@ describe("AgentRouter.runEvent — JSONL advancement guarantees", () => {
533533
.fn()
534534
.mockReturnValueOnce(0)
535535
.mockReturnValueOnce(0),
536+
latestAgentOutcome: vi
537+
.fn()
538+
.mockReturnValueOnce(undefined)
539+
.mockReturnValueOnce(undefined),
536540
latestAgentMessage: vi
537541
.fn()
538542
.mockReturnValueOnce(undefined)
@@ -576,7 +580,7 @@ describe("AgentRouter.runEvent — JSONL advancement guarantees", () => {
576580
.fn()
577581
.mockReturnValueOnce(0)
578582
.mockReturnValueOnce(1),
579-
latestAgentMessage: vi
583+
latestAgentOutcome: vi
580584
.fn()
581585
.mockReturnValueOnce(undefined)
582586
.mockReturnValueOnce({
@@ -589,6 +593,18 @@ describe("AgentRouter.runEvent — JSONL advancement guarantees", () => {
589593
tokensOut: 7,
590594
costUsd: 0.12,
591595
}),
596+
latestAgentMessage: vi
597+
.fn()
598+
.mockReturnValueOnce({
599+
eventId: "evt_new",
600+
sessionId: "ses_unused",
601+
type: "agent.message",
602+
content: "done",
603+
createdAt: Date.now(),
604+
tokensIn: 11,
605+
tokensOut: 7,
606+
costUsd: 0.12,
607+
}),
592608
};
593609
const { router, store } = makeRouter({
594610
poolStub: {
@@ -612,7 +628,7 @@ describe("AgentRouter.runEvent — JSONL advancement guarantees", () => {
612628
expect(finished?.costUsd).toBe(0.12);
613629
});
614630

615-
it("bakes first-turn model and thinking overrides into spawn options", async () => {
631+
it("keeps the turn successful when a tool result advances but no final agent.message is written", async () => {
616632
vi.stubGlobal(
617633
"fetch",
618634
vi.fn().mockResolvedValue(
@@ -631,7 +647,63 @@ describe("AgentRouter.runEvent — JSONL advancement guarantees", () => {
631647
.fn()
632648
.mockReturnValueOnce(0)
633649
.mockReturnValueOnce(1),
650+
latestAgentOutcome: vi
651+
.fn()
652+
.mockReturnValueOnce(undefined)
653+
.mockReturnValueOnce({
654+
eventId: "evt_tool_result",
655+
sessionId: "ses_unused",
656+
type: "agent.tool_result",
657+
content: "Fri Apr 24 17:58:01 UTC 2026",
658+
createdAt: Date.now(),
659+
toolName: "exec",
660+
toolCallId: "call-date",
661+
}),
634662
latestAgentMessage: vi
663+
.fn()
664+
.mockReturnValueOnce(undefined),
665+
};
666+
const { router, store } = makeRouter({
667+
poolStub: {
668+
acquireForSession: async () =>
669+
({ baseUrl: "http://container.test", token: "tok" }) as any,
670+
evictSession: async () => {},
671+
},
672+
eventReaderStub: fakeEvents as unknown as PiJsonlEventReader,
673+
});
674+
const agent = seedAgent(store);
675+
const session = router.createSession(agent.agentId);
676+
677+
await router.runEvent({ sessionId: session.sessionId, content: "what time is it?" });
678+
await waitForSessionToStopRunning(store, session.sessionId);
679+
680+
const finished = store.sessions.get(session.sessionId);
681+
expect(finished?.status).toBe("idle");
682+
expect(finished?.error).toBeNull();
683+
expect(finished?.tokensIn).toBe(11);
684+
expect(finished?.tokensOut).toBe(7);
685+
});
686+
687+
it("bakes first-turn model and thinking overrides into spawn options", async () => {
688+
vi.stubGlobal(
689+
"fetch",
690+
vi.fn().mockResolvedValue(
691+
new Response(
692+
JSON.stringify({
693+
choices: [{ message: { content: "done" } }],
694+
usage: { prompt_tokens: 11, completion_tokens: 7 },
695+
}),
696+
{ status: 200, headers: { "content-type": "application/json" } },
697+
),
698+
),
699+
);
700+
const fakeEvents = {
701+
stateRoot: "/tmp/test-state",
702+
countUserTurns: vi
703+
.fn()
704+
.mockReturnValueOnce(0)
705+
.mockReturnValueOnce(1),
706+
latestAgentOutcome: vi
635707
.fn()
636708
.mockReturnValueOnce(undefined)
637709
.mockReturnValueOnce({
@@ -643,6 +715,17 @@ describe("AgentRouter.runEvent — JSONL advancement guarantees", () => {
643715
tokensIn: 11,
644716
tokensOut: 7,
645717
}),
718+
latestAgentMessage: vi
719+
.fn()
720+
.mockReturnValueOnce({
721+
eventId: "evt_new",
722+
sessionId: "ses_unused",
723+
type: "agent.message",
724+
content: "done",
725+
createdAt: Date.now(),
726+
tokensIn: 11,
727+
tokensOut: 7,
728+
}),
646729
};
647730
let capturedSpawnOptions: unknown;
648731
const { router, store } = makeRouter({
@@ -701,9 +784,26 @@ describe("AgentRouter.runEvent — JSONL advancement guarantees", () => {
701784
.fn()
702785
.mockReturnValueOnce(1)
703786
.mockReturnValueOnce(2),
787+
latestAgentOutcome: vi
788+
.fn()
789+
.mockReturnValueOnce({
790+
eventId: "evt_old",
791+
sessionId: "ses_unused",
792+
type: "agent.message",
793+
content: "old",
794+
createdAt: 1,
795+
})
796+
.mockReturnValueOnce({
797+
eventId: "evt_new",
798+
sessionId: "ses_unused",
799+
type: "agent.message",
800+
content: "done",
801+
createdAt: Date.now(),
802+
tokensIn: 11,
803+
tokensOut: 7,
804+
}),
704805
latestAgentMessage: vi
705806
.fn()
706-
.mockReturnValueOnce({ eventId: "evt_old", sessionId: "ses_unused", type: "agent.message", content: "old", createdAt: 1 })
707807
.mockReturnValueOnce({
708808
eventId: "evt_new",
709809
sessionId: "ses_unused",
@@ -782,7 +882,7 @@ describe("AgentRouter.runEvent — JSONL advancement guarantees", () => {
782882
.fn()
783883
.mockReturnValueOnce(0)
784884
.mockReturnValueOnce(1),
785-
latestAgentMessage: vi
885+
latestAgentOutcome: vi
786886
.fn()
787887
.mockReturnValueOnce(undefined)
788888
.mockReturnValueOnce({
@@ -795,6 +895,18 @@ describe("AgentRouter.runEvent — JSONL advancement guarantees", () => {
795895
tokensOut: 45,
796896
costUsd: 0.42,
797897
}),
898+
latestAgentMessage: vi
899+
.fn()
900+
.mockReturnValueOnce({
901+
eventId: "evt_new",
902+
sessionId: "ses_unused",
903+
type: "agent.message",
904+
content: "done",
905+
createdAt: Date.now(),
906+
tokensIn: 321,
907+
tokensOut: 45,
908+
costUsd: 0.42,
909+
}),
798910
};
799911
const { router, store } = makeRouter({
800912
poolStub: {
@@ -836,7 +948,7 @@ describe("AgentRouter.runEvent — JSONL advancement guarantees", () => {
836948
.fn()
837949
.mockReturnValueOnce(0)
838950
.mockReturnValueOnce(1),
839-
latestAgentMessage: vi
951+
latestAgentOutcome: vi
840952
.fn()
841953
.mockReturnValueOnce(undefined)
842954
.mockReturnValueOnce({
@@ -847,6 +959,16 @@ describe("AgentRouter.runEvent — JSONL advancement guarantees", () => {
847959
createdAt: Date.now(),
848960
costUsd: 0.05,
849961
}),
962+
latestAgentMessage: vi
963+
.fn()
964+
.mockReturnValueOnce({
965+
eventId: "evt_new",
966+
sessionId: "ses_unused",
967+
type: "agent.message",
968+
content: "done",
969+
createdAt: Date.now(),
970+
costUsd: 0.05,
971+
}),
850972
};
851973
const { router, store } = makeRouter({
852974
poolStub: {
@@ -909,7 +1031,7 @@ describe("AgentRouter.runEvent — JSONL advancement guarantees", () => {
9091031
.fn()
9101032
.mockReturnValueOnce(0)
9111033
.mockReturnValueOnce(1),
912-
latestAgentMessage: vi
1034+
latestAgentOutcome: vi
9131035
.fn()
9141036
.mockReturnValueOnce(undefined)
9151037
.mockReturnValueOnce({
@@ -922,6 +1044,18 @@ describe("AgentRouter.runEvent — JSONL advancement guarantees", () => {
9221044
tokensOut: 45,
9231045
model: "zenmux/openai/gpt-5.4",
9241046
}),
1047+
latestAgentMessage: vi
1048+
.fn()
1049+
.mockReturnValueOnce({
1050+
eventId: "evt_new",
1051+
sessionId: "ses_unused",
1052+
type: "agent.message",
1053+
content: "done",
1054+
createdAt: Date.now(),
1055+
tokensIn: 321,
1056+
tokensOut: 45,
1057+
model: "zenmux/openai/gpt-5.4",
1058+
}),
9251059
};
9261060
const { router, store } = makeRouter({
9271061
poolStub: {

src/orchestrator/router.ts

Lines changed: 13 additions & 12 deletions
Original file line numberDiff line numberDiff line change
@@ -148,7 +148,7 @@ export type PendingApproval = {
148148

149149
type TurnProgressSnapshot = {
150150
userTurns: number;
151-
latestAgentMessageId?: string;
151+
latestAgentOutcomeId?: string;
152152
};
153153

154154
export class AgentRouter {
@@ -1638,10 +1638,11 @@ export class AgentRouter {
16381638
}
16391639

16401640
// Fast path: turn may have completed while the orchestrator was down.
1641-
// Compare the latest agent.message's createdAt to the session's
1642-
// lastEventAt (which was set by beginRun before the crash). A newer
1643-
// agent.message means Pi finished the turn without us watching.
1644-
const latest = this.events.latestAgentMessage(agent.agentId, sessionId);
1641+
// Compare the latest assistant-side outcome event's createdAt to the
1642+
// session's lastEventAt (which was set by beginRun before the crash).
1643+
// A newer agent.message OR agent.tool_result means Pi finished the
1644+
// turn without us watching.
1645+
const latest = this.events.latestAgentOutcome(agent.agentId, sessionId);
16451646
const startedAt = session.lastEventAt ?? session.createdAt;
16461647
if (latest && latest.createdAt > startedAt) {
16471648
log.info(
@@ -1721,8 +1722,8 @@ export class AgentRouter {
17211722
): TurnProgressSnapshot {
17221723
return {
17231724
userTurns: this.events.countUserTurns(agentId, sessionId),
1724-
latestAgentMessageId:
1725-
this.events.latestAgentMessage(agentId, sessionId)?.eventId,
1725+
latestAgentOutcomeId:
1726+
this.events.latestAgentOutcome(agentId, sessionId)?.eventId,
17261727
};
17271728
}
17281729

@@ -1757,22 +1758,22 @@ export class AgentRouter {
17571758
agentId: string,
17581759
sessionId: string,
17591760
before: TurnProgressSnapshot,
1760-
): Event {
1761+
): Event | undefined {
17611762
const afterUserTurns = this.events.countUserTurns(agentId, sessionId);
17621763
if (afterUserTurns <= before.userTurns) {
17631764
throw new RouterError(
17641765
"chat_completions_failed",
17651766
"turn returned but no new user.message was written to JSONL",
17661767
);
17671768
}
1768-
const latestAgent = this.events.latestAgentMessage(agentId, sessionId);
1769-
if (!latestAgent || latestAgent.eventId === before.latestAgentMessageId) {
1769+
const latestAgentOutcome = this.events.latestAgentOutcome(agentId, sessionId);
1770+
if (!latestAgentOutcome || latestAgentOutcome.eventId === before.latestAgentOutcomeId) {
17701771
throw new RouterError(
17711772
"chat_completions_failed",
1772-
"turn returned but no new agent.message was written to JSONL",
1773+
"turn returned but no new agent.message or agent.tool_result was written to JSONL",
17731774
);
17741775
}
1775-
return latestAgent;
1776+
return this.events.latestAgentMessage(agentId, sessionId);
17761777
}
17771778

17781779
/**

src/store/pi-jsonl.test.ts

Lines changed: 41 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -55,6 +55,7 @@ describe("PiJsonlEventReader", () => {
5555
const reader = new PiJsonlEventReader(root);
5656
expect(reader.listBySession("no-agent", "no-session")).toEqual([]);
5757
expect(reader.latestAgentMessage("no-agent", "no-session")).toBeUndefined();
58+
expect(reader.latestAgentOutcome("no-agent", "no-session")).toBeUndefined();
5859
});
5960

6061
it("returns [] when the JSONL file is missing but sessions.json maps the key", () => {
@@ -263,6 +264,46 @@ describe("PiJsonlEventReader", () => {
263264
expect(latest?.eventId).toBe("evt-c");
264265
});
265266

267+
it("latestAgentOutcome returns the newest agent.message or agent.tool_result", () => {
268+
const f = makeFixture([
269+
{
270+
type: "message",
271+
id: "evt-a",
272+
message: { role: "user", content: [{ type: "text", text: "what time is it?" }] },
273+
},
274+
{
275+
type: "message",
276+
id: "evt-b",
277+
message: {
278+
role: "assistant",
279+
content: [
280+
{
281+
type: "toolCall",
282+
id: "call-date",
283+
name: "exec",
284+
arguments: { command: "date" },
285+
},
286+
],
287+
},
288+
},
289+
{
290+
type: "message",
291+
id: "evt-c",
292+
message: {
293+
role: "toolResult",
294+
toolCallId: "call-date",
295+
toolName: "exec",
296+
content: [{ type: "text", text: "Fri Apr 24 17:58:01 UTC 2026" }],
297+
},
298+
},
299+
]);
300+
fixtures.push(f);
301+
const reader = new PiJsonlEventReader(f.root);
302+
const latest = reader.latestAgentOutcome(f.agentId, f.sessionId);
303+
expect(latest?.type).toBe("agent.tool_result");
304+
expect(latest?.eventId).toBe("evt-c");
305+
});
306+
266307
it("drops empty-content assistant messages (Pi auto-retry noise)", () => {
267308
const f = makeFixture([
268309
{

src/store/pi-jsonl.ts

Lines changed: 12 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -121,6 +121,18 @@ export class PiJsonlEventReader {
121121
return undefined;
122122
}
123123

124+
latestAgentOutcome(agentId: string, sessionId: string): Event | undefined {
125+
const events = this.listBySession(agentId, sessionId);
126+
for (let i = events.length - 1; i >= 0; i--) {
127+
const e = events[i];
128+
if (!e) continue;
129+
if (e.type === "agent.message" || e.type === "agent.tool_result") {
130+
return e;
131+
}
132+
}
133+
return undefined;
134+
}
135+
124136
/**
125137
* Count "turns" = user.message events for this session. Used by the
126138
* orchestrator's session-list endpoint to expose a turns field on the

0 commit comments

Comments
 (0)