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
16 changes: 11 additions & 5 deletions packages/server/src/harness/pi/process.ts
Original file line number Diff line number Diff line change
Expand Up @@ -559,11 +559,17 @@ export const makePiProcessWithDependencies = <R>(
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);
Expand Down
7 changes: 6 additions & 1 deletion packages/server/src/harness/pi/rpc/rpc-mode.ts
Original file line number Diff line number Diff line change
Expand Up @@ -445,7 +445,12 @@ export async function runRpcMode(runtimeHost: AgentSessionRuntime): Promise<neve
// 0.99 reports disposition instead of a boolean. isStreaming is
// already true for a prompt that just started, so it cannot mean
// "queued".
output(success(id, "prompt", { started: disposition === "started" }));
output(
success(id, "prompt", {
started: disposition === "started",
disposition,
}),
);
},
})
.catch((cause: unknown) => {
Expand Down
2 changes: 1 addition & 1 deletion packages/server/src/harness/pi/rpc/rpc-types.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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 }
Expand Down
51 changes: 51 additions & 0 deletions packages/server/test/harness/pi/agent.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -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: [] } });
Expand Down
Loading