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: [] } });