fix(coding-agent): queued extension sendUserMessage as steer while streaming

Extension sendUserMessage() without deliverAs fell through to prompt(),
which throws AgentBusyError during an active stream; the message was
dropped and surfaced as 'Extension sendUserMessage failed'. Route the
omitted-deliverAs path through prompt() with streamingBehavior 'steer'
so streaming queues a steer with normal prompt-flow side effects
(keyword notices, advisor auto-resume reset) and idle still starts a
turn.

ACP skill-command prompts now pass streamingBehavior 'steer'; the RPC
skill fast-path honors the prompt command's streamingBehavior field
(default steer) like the plain-prompt path already did. Documented the
extension-facing delivery semantics.

Synthesized from PR #4942 (prompt-flow steer routing, docs, tests) and
PR #4922 (RPC streamingBehavior threading, steer regression test);
dropped PR #4942's unrelated workflow-notice.md ellipsis churn.

Fixes #4923

Co-authored-by: roboomp <omp@can.ac>
Co-authored-by: metaphorics <metaphorics@users.noreply.github.com>
This commit is contained in:
can1357
2026-07-09 18:12:18 +02:00
parent 48877de084
commit e429166673
10 changed files with 161 additions and 29 deletions
+1 -1
View File
@@ -142,7 +142,7 @@ Also exposed:
- `deliverAs: "nextTurn"` — stored and injected on the next user prompt
- `triggerTurn: true` — starts a turn when idle (also honored with `deliverAs: "nextTurn"`: idle prompts immediately; while streaming the queued message schedules an internal continuation)
`pi.sendUserMessage(content, { deliverAs })` always goes through prompt flow; while streaming it queues as steer/follow-up.
`pi.sendUserMessage(content, { deliverAs })` always goes through prompt flow. Omit `deliverAs` to start a normal prompt when idle; while streaming, omitted `deliverAs` queues the message as a steer. Set `deliverAs: "followUp"` to wait until the current run finishes.
## 2) Handler context (`ExtensionContext`)
+3 -2
View File
@@ -215,8 +215,9 @@ Behavior:
1. optional command/template expansion (`/` commands, custom commands, file slash commands, prompt templates)
2. if currently streaming:
- requires `streamingBehavior: "steer" | "followUp"`
- queues instead of throwing work away
- `streamingBehavior: "steer" | "followUp"` chooses how `prompt()` queues
- extension `sendUserMessage(content)` defaults to steer when `deliverAs` is omitted
- queued messages are preserved instead of throwing work away
3. if idle:
- validates model + API key
- appends user message
+1
View File
@@ -23,6 +23,7 @@
- Fixed first-run setup ignoring a pre-seeded `config.yaml`: the settings loader now treats `config.yml` and `config.yaml` as equivalent existing main config files, writes back to the existing extension, and only creates canonical `config.yml` for fresh installs. ([#4914](https://github.com/can1357/oh-my-pi/issues/4914))
- Fixed extension `sendUserMessage()` without `deliverAs` surfacing `AgentBusyError` during active streams; omitted `deliverAs` now queues a steer through the normal prompt flow, and ACP/RPC skill-command prompts queue while streaming (RPC honors the prompt command's `streamingBehavior`, defaulting to steer) ([#4923](https://github.com/can1357/oh-my-pi/issues/4923)).
- Improved handling of unawaited promises in JS eval cells to prevent process crashes
- Added warning logs for unhandled rejections originating from finished eval cells
- Improved advisor robustness by blocking exhausted accounts during consecutive turn failures
@@ -1107,7 +1107,7 @@ export interface ExtensionAPI {
options?: { triggerTurn?: boolean; deliverAs?: "steer" | "followUp" | "nextTurn" },
): void;
/** Send a user message to the agent, or queue it when deliverAs is set. */
/** Send a user prompt: idle starts a turn; streaming queues as steer unless deliverAs is set. */
sendUserMessage(
content: string | (TextContent | ImageContent)[],
options?: { deliverAs?: "steer" | "followUp" },
@@ -843,13 +843,16 @@ export class AcpAgent implements Agent {
return false;
}
const built = await buildSkillPromptMessage(skill, parsed.args, "user");
await record.session.promptCustomMessage({
customType: SKILL_PROMPT_MESSAGE_TYPE,
content: built.message,
display: true,
details: built.details,
attribution: "user",
});
await record.session.promptCustomMessage(
{
customType: SKILL_PROMPT_MESSAGE_TYPE,
content: built.message,
display: true,
details: built.details,
attribution: "user",
},
{ streamingBehavior: "steer" },
);
return true;
}
@@ -88,6 +88,7 @@ export type RpcSkillCommandResult = { agentInvoked: true };
export async function tryRunRpcSkillCommand(
session: RpcSkillCommandSession,
text: string,
streamingBehavior: "steer" | "followUp" = "steer",
): Promise<RpcSkillCommandResult | false> {
if (!session.skillsSettings?.enableSkillCommands) return false;
const parsed = parseSkillInvocation(text);
@@ -95,13 +96,16 @@ export async function tryRunRpcSkillCommand(
const skill = session.skills.find(candidate => candidate.name === parsed.name);
if (!skill) return false;
const built = await buildSkillPromptMessage(skill, parsed.args, "user");
await session.promptCustomMessage({
customType: SKILL_PROMPT_MESSAGE_TYPE,
content: built.message,
display: true,
details: built.details,
attribution: "user",
});
await session.promptCustomMessage(
{
customType: SKILL_PROMPT_MESSAGE_TYPE,
content: built.message,
display: true,
details: built.details,
attribution: "user",
},
{ streamingBehavior },
);
return { agentInvoked: true };
}
@@ -842,7 +846,7 @@ export async function runRpcMode(
// =================================================================
case "prompt": {
const skillResult = await tryRunRpcSkillCommand(session, command.message);
const skillResult = await tryRunRpcSkillCommand(session, command.message, command.streamingBehavior);
if (skillResult) {
return success(id, "prompt", skillResult);
}
@@ -8309,11 +8309,10 @@ export class AgentSession {
}
/**
* Send a user message to the agent.
* When deliverAs is set, queue the message instead of starting a new turn.
* Send a user message through the prompt flow.
*
* @param content User message content (string or content array)
* @param options.deliverAs Delivery mode: "steer" or "followUp"
* Omitted `deliverAs` starts a turn when idle and queues as a steer while streaming.
* Explicit `deliverAs` queues without starting a turn in either state.
*/
async sendUserMessage(
content: string | (TextContent | ImageContent)[],
@@ -8348,10 +8347,13 @@ export class AgentSession {
return;
}
// Use prompt() with expandPromptTemplates: false to skip command handling and template expansion
// Use prompt() with expandPromptTemplates: false to skip command handling and template expansion.
// `streamingBehavior: "steer"` preserves prompt-flow side effects during streaming while
// covering the narrow race where a stream starts before prompt() acquires the turn.
await this.prompt(text, {
expandPromptTemplates: false,
images,
streamingBehavior: "steer",
});
}
+7 -1
View File
@@ -124,6 +124,7 @@ class FakeAgentSession {
}
promptCalls: string[] = [];
customMessages: Array<{ customType: string; content: string; details?: unknown }> = [];
customMessageOptions: Array<{ streamingBehavior?: "steer" | "followUp"; queueChipText?: string } | undefined> = [];
skillsSettings = { enableSkillCommands: true };
skills: Array<{ name: string; description: string; filePath: string; baseDir: string; source: string }> = [];
planModeState: PlanModeState | undefined;
@@ -235,8 +236,12 @@ class FakeAgentSession {
this.isStreaming = false;
}
async promptCustomMessage(message: { customType: string; content: string; details?: unknown }): Promise<void> {
async promptCustomMessage(
message: { customType: string; content: string; details?: unknown },
options?: { streamingBehavior?: "steer" | "followUp"; queueChipText?: string },
): Promise<void> {
this.customMessages.push(message);
this.customMessageOptions.push(options);
this.isStreaming = true;
const assistantMessage = makeAssistantMessage("skill pong");
for (const listener of this.#listeners) {
@@ -1443,6 +1448,7 @@ describe("ACP agent", () => {
expect(customMessage.content).toContain(`[Skill directory: ${skillDir}]`);
expect(customMessage.content).toMatch(/[Rr]esolve any relative paths/);
expect(customMessage.content).toContain("User: extra context");
expect(session.customMessageOptions[0]).toEqual({ streamingBehavior: "steer" });
harness.abortController.abort();
await Bun.sleep(0);
@@ -15,7 +15,7 @@ import { getBundledModel } from "@oh-my-pi/pi-catalog/models";
import { AsyncJobManager } from "@oh-my-pi/pi-coding-agent/async";
import type { Rule } from "@oh-my-pi/pi-coding-agent/capability/rule";
import { ModelRegistry } from "@oh-my-pi/pi-coding-agent/config/model-registry";
import { Settings } from "@oh-my-pi/pi-coding-agent/config/settings";
import { type SettingPath, Settings } from "@oh-my-pi/pi-coding-agent/config/settings";
import { TtsrManager } from "@oh-my-pi/pi-coding-agent/export/ttsr";
import type { ExtensionRunner } from "@oh-my-pi/pi-coding-agent/extensibility/extensions";
import { AgentSession } from "@oh-my-pi/pi-coding-agent/session/agent-session";
@@ -68,7 +68,7 @@ describe("AgentSession concurrent prompt guard", () => {
AsyncJobManager.resetForTests();
});
async function createSession() {
async function createSession(settingsOverrides?: Partial<Record<SettingPath, unknown>>) {
const model = getBundledModel("anthropic", "claude-sonnet-4-5")!;
let abortSignal: AbortSignal | undefined;
@@ -100,7 +100,7 @@ describe("AgentSession concurrent prompt guard", () => {
});
const sessionManager = SessionManager.inMemory();
const settings = Settings.isolated();
const settings = Settings.isolated(settingsOverrides);
const authStorage = await AuthStorage.create(path.join(tempDir, "testauth.db"));
authStorages.push(authStorage);
const modelRegistry = new ModelRegistry(authStorage, path.join(tempDir, "models.yml"));
@@ -176,6 +176,86 @@ describe("AgentSession concurrent prompt guard", () => {
await firstPrompt.catch(() => {});
});
it("queues sendUserMessage as steer while streaming without AgentBusyError", async () => {
await createSession();
const firstPrompt = session.prompt("First message");
await waitFor(() => session.isStreaming);
// The first agent loop may dequeue a steer before the assertion runs, so
// observe agent.steer itself rather than the residual queue length.
const steered: AgentMessage[] = [];
const originalSteer = session.agent.steer.bind(session.agent);
session.agent.steer = (message: AgentMessage) => {
steered.push(message);
originalSteer(message);
};
// Extension path: no deliverAs while busy must queue, not throw.
await expect(session.sendUserMessage("hello from extension")).resolves.toBeUndefined();
expect(steered).toHaveLength(1);
const queued = steered[0];
expect(queued?.role).toBe("user");
if (queued?.role === "user") {
expect(queued.content).toEqual([{ type: "text", text: "hello from extension" }]);
expect(queued.steering).toBe(true);
}
session.agent.clearAllQueues();
await session.abort();
await firstPrompt.catch(() => {});
});
it("sendUserMessage without deliverAs preserves prompt-flow keyword notices while streaming", async () => {
await createSession({ "magicKeywords.enabled": true, "magicKeywords.ultrathink": true });
const firstPrompt = session.prompt("First message");
await waitFor(() => session.isStreaming);
try {
await session.sendUserMessage("ultrathink fix via extension");
const queuedShape = session.agent
.peekSteeringQueue()
.map(message => (message.role === "custom" ? message.customType : message.role));
expect(queuedShape).toEqual(["ultrathink-notice", "user"]);
expect(session.getQueuedMessages()).toEqual({
steering: ["ultrathink fix via extension"],
followUp: [],
});
} finally {
session.agent.clearAllQueues();
await session.abort();
await firstPrompt.catch(() => {});
}
});
it("sendUserMessage without deliverAs starts a normal prompt when idle", async () => {
await createSession();
let rejected: unknown;
let settled = false;
const turn = session
.sendUserMessage("Idle extension message")
.catch(error => {
rejected = error;
})
.finally(() => {
settled = true;
});
try {
await waitFor(() => session.isStreaming || settled);
if (rejected) throw rejected;
expect(session.isStreaming).toBe(true);
expect(settled).toBe(false);
expect(session.getQueuedMessages()).toEqual({ steering: [], followUp: [] });
} finally {
await session.abort();
await turn;
}
});
it("delivers hidden nextTurn stop reactions through the next LLM call without exposing them in the visible queue", async () => {
const model = getBundledModel("anthropic", "claude-sonnet-4-5")!;
let firstStream: AssistantMessageEventStream | undefined;
@@ -16,6 +16,7 @@ describe("tryRunRpcSkillCommand", () => {
);
let message: Pick<CustomMessage, "attribution" | "content" | "customType" | "details" | "display"> | undefined;
let options: { streamingBehavior?: "steer" | "followUp" } | undefined;
const handled = await tryRunRpcSkillCommand(
{
@@ -23,8 +24,9 @@ describe("tryRunRpcSkillCommand", () => {
skills: [
{ name: "reviewer", description: "Review code", filePath: skillPath, baseDir: dir, source: "project" },
],
async promptCustomMessage(nextMessage: typeof message) {
async promptCustomMessage(nextMessage: typeof message, nextOptions?: typeof options) {
message = nextMessage;
options = nextOptions;
},
},
"/skill:reviewer focus on risks",
@@ -39,10 +41,43 @@ describe("tryRunRpcSkillCommand", () => {
expect(message?.content).toContain("User: focus on risks");
expect(message?.display).toBe(true);
expect(message?.attribution).toBe("user");
expect(options).toEqual({ streamingBehavior: "steer" });
await removeWithRetries(dir);
});
test("honors the RPC prompt streaming behavior for registered /skill commands", async () => {
const dir = await fs.mkdtemp(path.join(os.tmpdir(), `omp-rpc-skill-${Snowflake.next()}-`));
const skillPath = path.join(dir, "SKILL.md");
await Bun.write(
skillPath,
"---\nname: reviewer\ndescription: Review code\n---\n\nReview the supplied code carefully.\n",
);
let options: { streamingBehavior?: "steer" | "followUp" } | undefined;
try {
const handled = await tryRunRpcSkillCommand(
{
skillsSettings: { enableSkillCommands: true },
skills: [
{ name: "reviewer", description: "Review code", filePath: skillPath, baseDir: dir, source: "project" },
],
async promptCustomMessage(nextMessage, nextOptions) {
expect(nextMessage.customType).toBe(SKILL_PROMPT_MESSAGE_TYPE);
options = nextOptions;
},
},
"/skill:reviewer wait for the current turn",
"followUp",
);
expect(handled).toEqual({ agentInvoked: true });
expect(options?.streamingBehavior).toBe("followUp");
} finally {
await removeWithRetries(dir);
}
});
test("ignores unknown skill commands so normal prompt handling can continue", async () => {
const handled = await tryRunRpcSkillCommand(
{