From bd96752b83401b4fd6119bac79f06a856fd4ca80 Mon Sep 17 00:00:00 2001 From: can1357 Date: Fri, 31 Jul 2026 19:21:56 +0200 Subject: [PATCH] fix(rpc): serialize IRC wake finalization --- .../coding-agent/src/session/agent-session.ts | 37 +++++++------ packages/coding-agent/src/task/executor.ts | 54 ++++++++++--------- .../test/task/executor-soft-budget.test.ts | 6 ++- .../test/task/persisted-revive.test.ts | 4 +- 4 files changed, 56 insertions(+), 45 deletions(-) diff --git a/packages/coding-agent/src/session/agent-session.ts b/packages/coding-agent/src/session/agent-session.ts index 513fd9805..358a935b9 100644 --- a/packages/coding-agent/src/session/agent-session.ts +++ b/packages/coding-agent/src/session/agent-session.ts @@ -504,7 +504,9 @@ export class AgentSession { #asyncDeliveryEpoch = 0; readonly #irc: IrcBridge; - #ircWakeTurnObserver: ((records: CustomMessage[]) => ((error?: unknown) => void) | undefined) | undefined; + #ircWakeTurnObserver: + | ((records: CustomMessage[]) => ((error?: unknown) => void | Promise) | undefined) + | undefined; // Agent identity (registry id) used for IRC routing and job ownership. #agentId: string | undefined; #agentKind: "main" | "sub" = "main"; @@ -575,7 +577,7 @@ export class AgentSession { // on `agent_end` can fire its next `prompt` before #promptWithMessage's finally #promptGeneration = 0; #pendingAgentEndEmit: AgentSessionEvent | undefined; - #inFlightSettledCallbacks: Array<() => void> = []; + #inFlightSettledCallbacks: Array<() => void | Promise> = []; #sessionStopContinuationCount = 0; #sessionStopHookActive = false; #obfuscator: SecretObfuscator | undefined; @@ -641,23 +643,25 @@ export class AgentSession { } } - #endInFlight(onSettled?: () => void): void { + #endInFlight(onSettled?: () => void | Promise): void { if (onSettled) this.#inFlightSettledCallbacks.push(onSettled); this.#promptInFlightCount = Math.max(0, this.#promptInFlightCount - 1); - if (this.#promptInFlightCount === 0) { - this.#releasePowerAssertion(); - this.#flushPendingAgentEnd(); - this.#flushInFlightSettledCallbacks(); + if (this.#promptInFlightCount !== 0) return; + this.#releasePowerAssertion(); + this.#flushPendingAgentEnd(); + if (this.#inFlightSettledCallbacks.length === 0) { this.#drainStrandedQueuedMessages(); + return; } + void this.#flushInFlightSettledCallbacks().finally(() => this.#drainStrandedQueuedMessages()); } - #flushInFlightSettledCallbacks(): void { + async #flushInFlightSettledCallbacks(): Promise { const callbacks = this.#inFlightSettledCallbacks; this.#inFlightSettledCallbacks = []; for (const callback of callbacks) { try { - callback(); + await callback(); } catch (error) { logger.warn("In-flight settle callback failed", { error: String(error) }); } @@ -744,7 +748,7 @@ export class AgentSession { if (parkedFollowUps.length > 0) { this.agent.replaceQueues([...this.agent.peekSteeringQueue()], []); } - let finishObservation: ((error?: unknown) => void) | undefined; + let finishObservation: ((error?: unknown) => void | Promise) | undefined; try { finishObservation = this.#ircWakeTurnObserver?.(records); } catch (error) { @@ -779,9 +783,9 @@ export class AgentSession { [...parkedFollowUps, ...this.agent.peekFollowUpQueue()], ); } - this.#endInFlight(() => { + this.#endInFlight(async () => { try { - finishObservation?.(turnError); + await finishObservation?.(turnError); } catch (error) { logger.warn("IRC wake turn observer failed to finish", { error: String(error) }); } @@ -826,8 +830,11 @@ export class AgentSession { this.#promptInFlightCount = 0; this.#releasePowerAssertion(); this.#flushPendingAgentEnd(); - this.#flushInFlightSettledCallbacks(); - this.#drainStrandedQueuedMessages(); + if (this.#inFlightSettledCallbacks.length === 0) { + this.#drainStrandedQueuedMessages(); + return; + } + void this.#flushInFlightSettledCallbacks().finally(() => this.#drainStrandedQueuedMessages()); } #flushPendingAgentEnd(): void { @@ -6953,7 +6960,7 @@ export class AgentSession { /** Installs task-executor monitoring around autonomous IRC wake turns. */ setIrcWakeTurnObserver( - observer: ((records: CustomMessage[]) => ((error?: unknown) => void) | undefined) | undefined, + observer: ((records: CustomMessage[]) => ((error?: unknown) => void | Promise) | undefined) | undefined, ): void { this.#ircWakeTurnObserver = observer; } diff --git a/packages/coding-agent/src/task/executor.ts b/packages/coding-agent/src/task/executor.ts index 5daed097d..854d2489d 100644 --- a/packages/coding-agent/src/task/executor.ts +++ b/packages/coding-agent/src/task/executor.ts @@ -2277,7 +2277,7 @@ export function attachIrcWakeTurnMonitor(session: AgentSession, options: IrcWake turnMonitor.setActiveSession(session); const unsubscribeTurn = turnMonitor.attach(session); - return turnError => { + return async turnError => { unsubscribeTurn(); const activeSession = turnMonitor.takeActiveSession(); if (activeSession) turnMonitor.captureSalvage(activeSession); @@ -2294,35 +2294,37 @@ export function attachIrcWakeTurnMonitor(session: AgentSession, options: IrcWake : String(turnError) : undefined; turnMonitor.finish(); - void finalizeRunResult({ - monitor: turnMonitor, - done: { - exitCode: aborted || error ? 1 : 0, - error, - aborted, - abortReason: aborted ? turnMonitor.resolveAbortReasonText() : undefined, - durationMs: Date.now() - turnStartTime, - }, - index, - id, - agent, - task: ircTask, - modelOverride: options.modelOverride, - outputSchema: options.outputSchema, - outputSchemaMode: options.outputSchemaMode, - outputSchemaSource: options.outputSchemaSource, - artifactsDir: options.artifactsDir, - eventBus: options.eventBus, - parentToolCallId: options.parentToolCallId, - detached: true, - sessionFile, - startTime: turnStartTime, - }).catch(finalizeError => { + try { + await finalizeRunResult({ + monitor: turnMonitor, + done: { + exitCode: aborted || error ? 1 : 0, + error, + aborted, + abortReason: aborted ? turnMonitor.resolveAbortReasonText() : undefined, + durationMs: Date.now() - turnStartTime, + }, + index, + id, + agent, + task: ircTask, + modelOverride: options.modelOverride, + outputSchema: options.outputSchema, + outputSchemaMode: options.outputSchemaMode, + outputSchemaSource: options.outputSchemaSource, + artifactsDir: options.artifactsDir, + eventBus: options.eventBus, + parentToolCallId: options.parentToolCallId, + detached: true, + sessionFile, + startTime: turnStartTime, + }); + } catch (finalizeError) { logger.warn("IRC subagent turn finalization failed", { id, error: finalizeError instanceof Error ? finalizeError.message : String(finalizeError), }); - }); + } }; }); } diff --git a/packages/coding-agent/test/task/executor-soft-budget.test.ts b/packages/coding-agent/test/task/executor-soft-budget.test.ts index aab9b5698..7fbd28c57 100644 --- a/packages/coding-agent/test/task/executor-soft-budget.test.ts +++ b/packages/coding-agent/test/task/executor-soft-budget.test.ts @@ -54,7 +54,9 @@ function createMockSession( let abortCount = 0; let disposeCount = 0; let promptIndex = 0; - let ircWakeTurnObserver: ((records: CustomMessage[]) => ((error?: unknown) => void) | undefined) | undefined; + let ircWakeTurnObserver: + | ((records: CustomMessage[]) => ((error?: unknown) => void | Promise) | undefined) + | undefined; const emit = (event: AgentSessionEvent) => { for (const listener of [...listeners]) listener(event); @@ -125,7 +127,7 @@ function createMockSession( isError: false, } as AgentSessionEvent); emit({ type: "agent_end", messages: [yieldMessage] } as unknown as AgentSessionEvent); - finishObservation?.(); + await finishObservation?.(); return "woken"; }, abort: async () => { diff --git a/packages/coding-agent/test/task/persisted-revive.test.ts b/packages/coding-agent/test/task/persisted-revive.test.ts index 17d366678..eede151d1 100644 --- a/packages/coding-agent/test/task/persisted-revive.test.ts +++ b/packages/coding-agent/test/task/persisted-revive.test.ts @@ -39,7 +39,7 @@ function createRef(sessionFile: string): AgentRef { }; } -type IrcWakeObserver = (records: CustomMessage[]) => ((error?: unknown) => void) | undefined; +type IrcWakeObserver = (records: CustomMessage[]) => ((error?: unknown) => void | Promise) | undefined; interface RevivedSessionHandle { session: AgentSession; @@ -223,7 +223,7 @@ describe("persisted subagent revival", () => { timestamp: Date.now(), }; const finish = observer?.([record]); - finish?.(); + await finish?.(); await terminal.promise; expect(frames[0]).toMatchObject({