From 7b9bd017f1227ec3ee094a4b5d175bb26b13fdfe Mon Sep 17 00:00:00 2001 From: dinq Date: Thu, 1 Oct 2026 02:20:08 -0400 Subject: [PATCH] fix(server): stream a prompt Pi queued after the server turn ended Pi can still be streaming after Pie has consumed the finish event. A follow-up was then rejected as "queued without an active server turn", so the UI froze while Pi ran the prompt unobserved. --- packages/server/src/harness/pi/process.ts | 16 ++++-- .../server/src/harness/pi/rpc/rpc-mode.ts | 7 ++- .../server/src/harness/pi/rpc/rpc-types.ts | 2 +- packages/server/test/harness/pi/agent.test.ts | 51 +++++++++++++++++++ 4 files changed, 69 insertions(+), 7 deletions(-) diff --git a/packages/server/src/harness/pi/process.ts b/packages/server/src/harness/pi/process.ts index 3fb1b3e27..811015db7 100644 --- a/packages/server/src/harness/pi/process.ts +++ b/packages/server/src/harness/pi/process.ts @@ -559,11 +559,17 @@ export const makePiProcessWithDependencies = ( output: Stream.empty, }; } - return yield* new AgentOperationError({ - sessionId: input.sessionId, - operation: "prompt-admission-state", - cause: new Error("Pi queued a prompt without an active server turn"), - }); + // Extension command already ran. No model turn follows. + if (admission?.disposition === "handled") { + return { + turnId: uuid(), + started: false, + output: Stream.empty, + }; + } + // Pi queued this prompt while we have no turn (its run + // outlived the finish we already consumed). Attach and + // stream that run instead of failing the RPC. } const previous = yield* Ref.get(session.turnState); diff --git a/packages/server/src/harness/pi/rpc/rpc-mode.ts b/packages/server/src/harness/pi/rpc/rpc-mode.ts index 2d0222226..548d567ba 100644 --- a/packages/server/src/harness/pi/rpc/rpc-mode.ts +++ b/packages/server/src/harness/pi/rpc/rpc-mode.ts @@ -445,7 +445,12 @@ export async function runRpcMode(runtimeHost: AgentSessionRuntime): Promise { diff --git a/packages/server/src/harness/pi/rpc/rpc-types.ts b/packages/server/src/harness/pi/rpc/rpc-types.ts index e22de8388..a83febf60 100644 --- a/packages/server/src/harness/pi/rpc/rpc-types.ts +++ b/packages/server/src/harness/pi/rpc/rpc-types.ts @@ -144,7 +144,7 @@ export type RpcResponse = type: "response"; command: "prompt"; success: true; - data: { started: boolean }; + data: { started: boolean; disposition?: "started" | "queued" | "handled" }; } | { id?: string; type: "response"; command: "steer"; success: true } | { id?: string; type: "response"; command: "follow_up"; success: true } diff --git a/packages/server/test/harness/pi/agent.test.ts b/packages/server/test/harness/pi/agent.test.ts index 2c6a128ee..7b1025457 100644 --- a/packages/server/test/harness/pi/agent.test.ts +++ b/packages/server/test/harness/pi/agent.test.ts @@ -94,6 +94,22 @@ rl.on("line", (line) => { if (msg.type !== "prompt") return; const text = msg.message; if (text === "fail") { send({ id: msg.id, type: "response", command: "prompt", success: false, error: "cannot prompt" }); return; } + if (text === "queue-idle") { + send({ id: msg.id, type: "response", command: "prompt", success: true, data: { started: false, disposition: "queued" } }); + send({ type: "agent_start" }); + send({ type: "message_start", message: assistant() }); + upd({ type: "start" }); + upd({ type: "text_start", contentIndex: 0 }); + upd({ type: "text_delta", contentIndex: 0, delta: "queued" }); + upd({ type: "text_end", contentIndex: 0, content: "queued" }); + send({ type: "message_end", message: assistant() }); + settle(); + return; + } + if (text === "handled") { + send({ id: msg.id, type: "response", command: "prompt", success: true, data: { started: false, disposition: "handled" } }); + return; + } if (holding && !msg.streamingBehavior) { send({ id: msg.id, type: "response", command: "prompt", success: false, error: "Agent is already processing. Specify streamingBehavior ('steer' or 'followUp') to queue the message." }); return; @@ -192,6 +208,41 @@ layer(NodeServices.layer)("PiAgent", (it) => { }), ); + it.effect("streams a prompt Pi queued while the server has no turn", () => + Effect.gen(function* () { + const agent = yield* makePiProcess({ executable: { command: makeFake(), prefixArgs: [] } }); + const { sessionId } = yield* agent.session.create({ cwd: "/tmp" }); + const prompt = yield* agent.session.prompt({ sessionId, text: "queue-idle" }); + assert.equal(prompt.started, true); + const chunks = yield* Stream.runCollect(prompt.output); + assert.deepEqual( + Array.from(chunks, (chunk) => chunk.type), + [ + "start", + "message-metadata", + "text-start", + "text-delta", + "text-end", + "message-metadata", + "finish", + ], + ); + yield* agent.session.abort(sessionId); + }), + ); + + it.effect("accepts a handled prompt without opening a turn", () => + Effect.gen(function* () { + const agent = yield* makePiProcess({ executable: { command: makeFake(), prefixArgs: [] } }); + const { sessionId } = yield* agent.session.create({ cwd: "/tmp" }); + const prompt = yield* agent.session.prompt({ sessionId, text: "handled" }); + assert.equal(prompt.started, false); + const chunks = yield* Stream.runCollect(prompt.output); + assert.equal(Array.from(chunks).length, 0); + yield* agent.session.abort(sessionId); + }), + ); + it.effect("resume keeps the caller-provided session id", () => Effect.gen(function* () { const agent = yield* makePiProcess({ executable: { command: makeFake(), prefixArgs: [] } });