From b59ea8e34487f48cf12689bd4acaf7a4d0a344c9 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Korm=C3=A1kur?= Date: Sat, 13 Jun 2026 00:46:45 +0000 Subject: [PATCH] fix(coding-agent): handle queued rpc extension messages --- packages/coding-agent/CHANGELOG.md | 4 + .../coding-agent/src/modes/acp/acp-agent.ts | 11 +- .../coding-agent/src/modes/rpc/rpc-mode.ts | 46 +++++++- .../coding-agent/src/modes/runtime-init.ts | 28 ++++- packages/coding-agent/src/task/executor.ts | 13 ++- packages/coding-agent/test/acp-agent.test.ts | 4 +- .../test/acp-initialize-conformance.test.ts | 4 +- .../test/acp-lazy-startup.test.ts | 4 +- .../test/rpc-prompt-result.test.ts | 106 ++++++++++++++++++ .../task/executor-subagent-reminders.test.ts | 1 + .../test/task/task-guards.test.ts | 1 + 11 files changed, 204 insertions(+), 18 deletions(-) diff --git a/packages/coding-agent/CHANGELOG.md b/packages/coding-agent/CHANGELOG.md index 41e8264ed..3b3a021be 100644 --- a/packages/coding-agent/CHANGELOG.md +++ b/packages/coding-agent/CHANGELOG.md @@ -163,6 +163,10 @@ - Fixed misaligned box borders in Mermaid ASCII rendering for CJK (Korean/Japanese/Chinese) and emoji labels — affects both fenced `mermaid` code blocks in assistant messages and the `render_mermaid` tool. `beautiful-mermaid@1.1.3` measures label width in UTF-16 code units while terminals render East Asian characters 2 columns wide; a `patchedDependencies` entry rebuilds its ASCII renderer to measure terminal display columns (grapheme-cluster aware, wcwidth-style policy). The patch mirrors the upstream PR (lukilabs/beautiful-mermaid#128) and should be dropped once it ships in a release. - Fixed interrupt loader state getting stuck after queued-message aborts by removing the session-layer flush/latch path; empty Enter now aborts the active turn and lets the existing post-unwind queue drain resume normally. - Fixed `/goal ` and `/goal set ` during streaming so goal context is steered immediately but objective submission waits for the active turn to finish instead of spamming `AgentBusyError` ([#2454](https://github.com/can1357/oh-my-pi/issues/2454)). +- Fixed RPC local-only prompt completion for extension `pi.sendUserMessage(..., { deliverAs: "followUp" | "steer" })` calls that only queue messages, so hosts no longer wait for an `agent_end` that will not be emitted. + +### Fixed + - Fixed concurrent `omp --session` startups (e.g. cmux pane restore after an unclean shutdown) crashing with `SQLITE_BUSY_RECOVERY` while the agent SQLite databases were still under WAL recovery. The auth credential store and `AgentStorage.open()` retry the `SQLITE_BUSY` family with bounded backoff, and every shared SQLite open path (`AgentStorage`, history, autoresearch, memories, github cache, auto-QA grievances, catalog model cache, stats) now installs the busy handler before the first lock-taking statement so transient WAL recovery contention waits instead of crashing ([#2421](https://github.com/can1357/oh-my-pi/issues/2421)). - Mnemopi `per-project` / `per-project-tagged` bank derivation is now stable for one cwd, ignoring the surrounding git layout. Previously the bank id was hashed from `git.repo.resolveSync(cwd)?.repoRoot ?? path.resolve(cwd)`, so adding or removing a `.git` anywhere above the working directory silently repointed the same conversation to a new bank and stranded its memories (e.g. `/home/x/projects/repo` flipping between `projects-…` and `repo-…`). The derivation in `packages/coding-agent/src/mnemopi/config.ts` now hashes `path.resolve(cwd)` directly, and session startup widens the recall set with any sibling bank under `/banks/` whose `working_memory` rows already carry the active cwd in `metadata_json.$.cwd`, so memories stranded by the old, less-stable derivation become visible again on the next session without manual migration ([#2412](https://github.com/can1357/oh-my-pi/issues/2412)). - Fixed model switching (Ctrl+P role cycling and the alt+p / `/switch` / `/models` selector) intermittently freezing the UI for several seconds. `AgentSession.setModel`/`setModelTemporary` ran an eager `await modelRegistry.getApiKey(model)` purely as an existence pre-flight and discarded the value — but `getApiKey` does real work: it synchronously executes command-backed key programs (`apiKey: "!cmd"`, `execSync` with a 10s timeout, blocking the event loop) and refreshes OAuth tokens over the network when one crosses the expiry window (the "fine for a few switches, then a multi-second stall" symptom). Switching now uses the synchronous, side-effect-free `ModelRegistry.hasConfiguredAuth` check; the concrete key (command execution + OAuth refresh) is still resolved lazily per request via the existing resolver, so an unconfigured provider still fails fast with `No API key` while a healthy switch never touches the network or spawns a subprocess. `hasConfiguredAuth` no longer runs the command program or refreshes tokens either, matching its documented "probe without resolving an API key" contract. diff --git a/packages/coding-agent/src/modes/acp/acp-agent.ts b/packages/coding-agent/src/modes/acp/acp-agent.ts index a2fc34824..856486196 100644 --- a/packages/coding-agent/src/modes/acp/acp-agent.ts +++ b/packages/coding-agent/src/modes/acp/acp-agent.ts @@ -685,10 +685,13 @@ export class AcpAgent implements Agent { } } - #trackExtensionUserMessage(record: ManagedSessionRecord, task: Promise): void { - const tracked = task.catch((error: unknown) => { - logger.warn("ACP extension sendUserMessage failed", { error }); - }); + #trackExtensionUserMessage(record: ManagedSessionRecord, task: Promise): void { + const tracked = task.then( + () => undefined, + (error: unknown) => { + logger.warn("ACP extension sendUserMessage failed", { error }); + }, + ); record.extensionUserMessageTasks.add(tracked); void tracked.finally(() => { record.extensionUserMessageTasks.delete(tracked); diff --git a/packages/coding-agent/src/modes/rpc/rpc-mode.ts b/packages/coding-agent/src/modes/rpc/rpc-mode.ts index 237edfa63..52d92a3a3 100644 --- a/packages/coding-agent/src/modes/rpc/rpc-mode.ts +++ b/packages/coding-agent/src/modes/rpc/rpc-mode.ts @@ -111,10 +111,13 @@ export function reportLocalOnlyPromptResult(input: { output: (obj: object) => void; onError: (error: Error) => void; hasExtensionAgentMessageTask?: () => boolean; + waitForExtensionAgentMessageTasks?: () => Promise; }): void { void input.prompt - .then(agentInvoked => { - if (!agentInvoked && !input.hasExtensionAgentMessageTask?.()) { + .then(async agentInvoked => { + if (agentInvoked) return; + await input.waitForExtensionAgentMessageTasks?.(); + if (!input.hasExtensionAgentMessageTask?.()) { input.output({ type: "prompt_result", id: input.id, agentInvoked: false }); } }) @@ -125,6 +128,7 @@ export function reportLocalOnlyPromptResult(input: { type RpcExtensionUserMessageScope = { hasAgentMessageTask: boolean; + pendingAgentMessageTasks: Set>; }; /** @@ -142,11 +146,42 @@ export class RpcExtensionUserMessageTracker { } } + trackAgentMessageTask(task: Promise): void { + for (const scope of this.#activePromptScopes) { + this.#trackAgentMessageTaskForScope(scope, task); + } + } + + #trackAgentMessageTaskForScope(scope: RpcExtensionUserMessageScope, task: Promise): void { + const scopedTask = task.then( + agentInvoked => { + if (agentInvoked) { + scope.hasAgentMessageTask = true; + } + }, + () => {}, + ); + scope.pendingAgentMessageTasks.add(scopedTask); + void scopedTask.finally(() => { + scope.pendingAgentMessageTasks.delete(scopedTask); + }); + } + + async #waitForAgentMessageTasks(scope: RpcExtensionUserMessageScope): Promise { + while (scope.pendingAgentMessageTasks.size > 0) { + await Promise.allSettled(Array.from(scope.pendingAgentMessageTasks)); + } + } + watchPrompt(startPrompt: () => Promise): { prompt: Promise; hasAgentMessageTask: () => boolean; + waitForAgentMessageTasks: () => Promise; } { - const scope: RpcExtensionUserMessageScope = { hasAgentMessageTask: false }; + const scope: RpcExtensionUserMessageScope = { + hasAgentMessageTask: false, + pendingAgentMessageTasks: new Set(), + }; this.#activePromptScopes.add(scope); let prompt: Promise; try { @@ -160,6 +195,7 @@ export class RpcExtensionUserMessageTracker { this.#activePromptScopes.delete(scope); }), hasAgentMessageTask: () => scope.hasAgentMessageTask, + waitForAgentMessageTasks: () => this.#waitForAgentMessageTasks(scope), }; } } @@ -178,6 +214,7 @@ export function watchAndReportLocalOnlyPromptResult(input: { output: input.output, onError: input.onError, hasExtensionAgentMessageTask: trackedPrompt.hasAgentMessageTask, + waitForExtensionAgentMessageTasks: trackedPrompt.waitForAgentMessageTasks, }); } @@ -619,6 +656,9 @@ export async function runRpcMode( markAgentInvokingMessage: () => { extensionUserMessageTracker.markAgentMessageTask(); }, + trackAgentInvokingUserMessage: task => { + extensionUserMessageTracker.trackAgentMessageTask(task); + }, uiContext: rpcUiContext, }); diff --git a/packages/coding-agent/src/modes/runtime-init.ts b/packages/coding-agent/src/modes/runtime-init.ts index d78196b29..8167dc048 100644 --- a/packages/coding-agent/src/modes/runtime-init.ts +++ b/packages/coding-agent/src/modes/runtime-init.ts @@ -25,6 +25,8 @@ export interface InitializeExtensionsOptions { uiContext?: ExtensionUIContext; /** Optional lifecycle hook for extension-originated messages that can start an agent turn. */ markAgentInvokingMessage?: () => void; + /** Optional lifecycle hook for extension-originated user messages with async delivery outcomes. */ + trackAgentInvokingUserMessage?: (task: Promise) => void; } /** @@ -37,7 +39,14 @@ export async function initializeExtensions(session: AgentSession, options: Initi const runner = session.extensionRunner; if (!runner) return; - const { reportSendError, reportRuntimeError, onShutdown, uiContext, markAgentInvokingMessage } = options; + const { + reportSendError, + reportRuntimeError, + onShutdown, + uiContext, + markAgentInvokingMessage, + trackAgentInvokingUserMessage, + } = options; const shutdown = onShutdown ?? (() => {}); runner.initialize( @@ -52,8 +61,21 @@ export async function initializeExtensions(session: AgentSession, options: Initi }); }, sendUserMessage: (content, sendOptions) => { - markAgentInvokingMessage?.(); - session.sendUserMessage(content, sendOptions).catch(e => { + const sendTask = session.sendUserMessage(content, sendOptions); + if (trackAgentInvokingUserMessage) { + trackAgentInvokingUserMessage(sendTask); + } else if (!sendOptions?.deliverAs) { + markAgentInvokingMessage?.(); + } else { + void sendTask + .then(agentInvoked => { + if (agentInvoked) { + markAgentInvokingMessage?.(); + } + }) + .catch(() => {}); + } + void sendTask.catch(e => { reportSendError("extension_send_user", e instanceof Error ? e : new Error(String(e))); }); }, diff --git a/packages/coding-agent/src/task/executor.ts b/packages/coding-agent/src/task/executor.ts index 5d7f4b93d..effe0ed23 100644 --- a/packages/coding-agent/src/task/executor.ts +++ b/packages/coding-agent/src/task/executor.ts @@ -2018,11 +2018,14 @@ export async function runSubprocess(options: ExecutorOptions): Promise { - const sendPromise = session.sendUserMessage(content, options).catch(e => { - logger.error("Extension sendUserMessage failed", { - error: e instanceof Error ? e.message : String(e), - }); - }); + const sendPromise = session.sendUserMessage(content, options).then( + () => undefined, + e => { + logger.error("Extension sendUserMessage failed", { + error: e instanceof Error ? e.message : String(e), + }); + }, + ); pendingExtensionMessages.push(sendPromise); }, appendEntry: (customType, data) => { diff --git a/packages/coding-agent/test/acp-agent.test.ts b/packages/coding-agent/test/acp-agent.test.ts index 1dc44230f..e3ca80b34 100644 --- a/packages/coding-agent/test/acp-agent.test.ts +++ b/packages/coding-agent/test/acp-agent.test.ts @@ -330,7 +330,9 @@ class FakeAgentSession { async sendCustomMessage(_message: string, _options?: unknown): Promise {} - async sendUserMessage(_content: string, _options?: unknown): Promise {} + async sendUserMessage(_content: string, _options?: unknown): Promise { + return false; + } async compact(_instructions?: string, _options?: unknown): Promise {} diff --git a/packages/coding-agent/test/acp-initialize-conformance.test.ts b/packages/coding-agent/test/acp-initialize-conformance.test.ts index 8df0a2a37..fbe3e687f 100644 --- a/packages/coding-agent/test/acp-initialize-conformance.test.ts +++ b/packages/coding-agent/test/acp-initialize-conformance.test.ts @@ -110,7 +110,9 @@ class FakeAgentSession { } setPlanModeState(): void {} async sendCustomMessage(): Promise {} - async sendUserMessage(): Promise {} + async sendUserMessage(): Promise { + return false; + } async compact(): Promise {} async fork(): Promise { return false; diff --git a/packages/coding-agent/test/acp-lazy-startup.test.ts b/packages/coding-agent/test/acp-lazy-startup.test.ts index d8368d4ae..a95fc8c5f 100644 --- a/packages/coding-agent/test/acp-lazy-startup.test.ts +++ b/packages/coding-agent/test/acp-lazy-startup.test.ts @@ -133,7 +133,9 @@ class LazyFakeSession { } setPlanModeState(): void {} async sendCustomMessage(): Promise {} - async sendUserMessage(): Promise {} + async sendUserMessage(): Promise { + return false; + } async compact(): Promise {} async fork(): Promise { return false; diff --git a/packages/coding-agent/test/rpc-prompt-result.test.ts b/packages/coding-agent/test/rpc-prompt-result.test.ts index a0a9b61a4..9167e5865 100644 --- a/packages/coding-agent/test/rpc-prompt-result.test.ts +++ b/packages/coding-agent/test/rpc-prompt-result.test.ts @@ -140,6 +140,112 @@ describe("reportLocalOnlyPromptResult", () => { expect(sentOptions).toEqual({ triggerTurn: true }); }); + test("emits prompt_result when extension sendUserMessage only queues follow-up", async () => { + let extensionActions: ExtensionActions | undefined; + let sentOptions: { deliverAs?: "steer" | "followUp" } | undefined; + const delivery = Promise.withResolvers(); + const extensionUserMessages = new RpcExtensionUserMessageTracker(); + const session = { + extensionRunner: { + initialize: (actions: ExtensionActions) => { + extensionActions = actions; + }, + onError: () => {}, + emit: async () => {}, + }, + sendUserMessage: async (_content: unknown, options?: { deliverAs?: "steer" | "followUp" }) => { + sentOptions = options; + return await delivery.promise; + }, + } as unknown as AgentSession; + + await initializeExtensions(session, { + reportSendError: (_action, error) => { + throw error; + }, + reportRuntimeError: error => { + throw error.error; + }, + trackAgentInvokingUserMessage: task => { + extensionUserMessages.trackAgentMessageTask(task); + }, + }); + + const output: object[] = []; + const trackedPrompt = extensionUserMessages.watchPrompt(() => { + if (!extensionActions) throw new Error("extensions not initialized"); + extensionActions.sendUserMessage("queued locally", { deliverAs: "followUp" }); + return Promise.resolve(false); + }); + reportLocalOnlyPromptResult({ + id: "req_queued", + prompt: trackedPrompt.prompt, + output: frame => output.push(frame), + onError: error => { + throw error; + }, + hasExtensionAgentMessageTask: trackedPrompt.hasAgentMessageTask, + waitForExtensionAgentMessageTasks: trackedPrompt.waitForAgentMessageTasks, + }); + + await waitForPromptHandlers(trackedPrompt.prompt); + expect(output).toEqual([]); + + delivery.resolve(false); + await waitForPromptHandlers(delivery.promise); + await waitForPromptHandlers(trackedPrompt.prompt); + + expect(sentOptions).toEqual({ deliverAs: "followUp" }); + expect(output).toEqual([{ type: "prompt_result", id: "req_queued", agentInvoked: false }]); + }); + + test("suppresses prompt_result when extension sendUserMessage starts agent work", async () => { + let extensionActions: ExtensionActions | undefined; + const extensionUserMessages = new RpcExtensionUserMessageTracker(); + const session = { + extensionRunner: { + initialize: (actions: ExtensionActions) => { + extensionActions = actions; + }, + onError: () => {}, + emit: async () => {}, + }, + sendUserMessage: async () => true, + } as unknown as AgentSession; + + await initializeExtensions(session, { + reportSendError: (_action, error) => { + throw error; + }, + reportRuntimeError: error => { + throw error.error; + }, + trackAgentInvokingUserMessage: task => { + extensionUserMessages.trackAgentMessageTask(task); + }, + }); + + const output: object[] = []; + const trackedPrompt = extensionUserMessages.watchPrompt(() => { + if (!extensionActions) throw new Error("extensions not initialized"); + extensionActions.sendUserMessage("start work"); + return Promise.resolve(false); + }); + reportLocalOnlyPromptResult({ + id: "req_agent", + prompt: trackedPrompt.prompt, + output: frame => output.push(frame), + onError: error => { + throw error; + }, + hasExtensionAgentMessageTask: trackedPrompt.hasAgentMessageTask, + waitForExtensionAgentMessageTasks: trackedPrompt.waitForAgentMessageTasks, + }); + await waitForPromptHandlers(trackedPrompt.prompt); + + expect(output).toEqual([]); + }); + test("does not emit when prompt invokes the agent", async () => { const output: object[] = []; const prompt = Promise.resolve(true); diff --git a/packages/coding-agent/test/task/executor-subagent-reminders.test.ts b/packages/coding-agent/test/task/executor-subagent-reminders.test.ts index 7f69af243..166ed05c9 100644 --- a/packages/coding-agent/test/task/executor-subagent-reminders.test.ts +++ b/packages/coding-agent/test/task/executor-subagent-reminders.test.ts @@ -146,6 +146,7 @@ describe("runSubprocess yield reminders", () => { messageInFlight = true; await Bun.sleep(20); messageInFlight = false; + return true; }; mutableSession.extensionRunner = { initialize: (actions: ExtensionActions) => { diff --git a/packages/coding-agent/test/task/task-guards.test.ts b/packages/coding-agent/test/task/task-guards.test.ts index 37134a2ec..e3029f4d6 100644 --- a/packages/coding-agent/test/task/task-guards.test.ts +++ b/packages/coding-agent/test/task/task-guards.test.ts @@ -98,6 +98,7 @@ function createFakeSession(config: FakeSessionConfig = {}): FakeSessionHandle { }, sendUserMessage: async (content, options) => { steerCalls.push({ content: String(content), options }); + return false; }, getLastAssistantMessage: () => (config.lastAssistantMessage ?? undefined) as never, abort: async () => {