From dda775459ac702f4cc240c8aa2e173f44fc7d3cd Mon Sep 17 00:00:00 2001 From: Vincent Koc Date: Tue, 18 Aug 2026 22:43:04 -0700 Subject: [PATCH] fix(talk): preserve accepted work on relay detach Co-authored-by: Dallin Romney Punchcard-Session: calm-brook-timber-r0 --- ...ltime-quicksilver-delegation-controller.ts | 36 +++++- ...icksilver-gateway-bridge-lifecycle.test.ts | 114 ++++++++++++++++++ .../realtime-quicksilver-gateway-bridge.ts | 24 ++-- src/gateway/talk-realtime-relay-operations.ts | 28 ++++- src/gateway/talk-realtime-relay.test.ts | 49 ++++++++ src/plugin-sdk/realtime-voice.ts | 2 + src/talk/provider-types.ts | 9 +- src/talk/session-runtime.test.ts | 19 +++ src/talk/session-runtime.ts | 7 +- 9 files changed, 266 insertions(+), 22 deletions(-) create mode 100644 extensions/openai/realtime-quicksilver-gateway-bridge-lifecycle.test.ts diff --git a/extensions/openai/realtime-quicksilver-delegation-controller.ts b/extensions/openai/realtime-quicksilver-delegation-controller.ts index 51138b913cc9..42e29edc76c3 100644 --- a/extensions/openai/realtime-quicksilver-delegation-controller.ts +++ b/extensions/openai/realtime-quicksilver-delegation-controller.ts @@ -54,12 +54,22 @@ function readWireEventType(payload: string): string | undefined { export class OpenAIQuicksilverDelegationController { private activeDelegationId: string | undefined; private consultController: AbortController | undefined; + private readonly onSessionAbort = () => { + const reason = this.options.signal.reason; + this.stop(reason instanceof Error ? reason : new Error("GPT-Live session stopped")); + }; private partialTranscriptRole: "user" | "assistant" | undefined; private pendingDelegation: PendingDelegation | undefined; private stopped = false; private transcript: OpenAIQuicksilverTranscriptEntry[] = []; - constructor(private readonly options: OpenAIQuicksilverDelegationControllerOptions) {} + constructor(private readonly options: OpenAIQuicksilverDelegationControllerOptions) { + if (options.signal.aborted) { + this.onSessionAbort(); + } else { + options.signal.addEventListener("abort", this.onSessionAbort, { once: true }); + } + } handleFrame(data: RawData, isBinary: boolean): void { if (isBinary) { @@ -122,13 +132,19 @@ export class OpenAIQuicksilverDelegationController { if (this.stopped) { return; } - this.stopped = true; - this.pendingDelegation = undefined; - this.activeDelegationId = undefined; + this.markStopped(); this.consultController?.abort(reason); this.consultController = undefined; } + /** Releases sideband ownership without canceling work already accepted by the host. */ + detach(): void { + if (this.stopped) { + return; + } + this.markStopped(); + } + private appendTranscript( event: Extract, ): void { @@ -180,8 +196,7 @@ export class OpenAIQuicksilverDelegationController { const controller = new AbortController(); this.consultController = controller; this.activeDelegationId = delegation.id; - const signal = AbortSignal.any([this.options.signal, controller.signal]); - void this.runDelegation(delegation, signal) + void this.runDelegation(delegation, controller.signal) .catch((error: unknown) => this.fail(toErrorObject(error, "OpenAI GPT-Live delegation failed")), ) @@ -200,6 +215,15 @@ export class OpenAIQuicksilverDelegationController { }); } + private markStopped(): void { + this.stopped = true; + this.options.signal.removeEventListener("abort", this.onSessionAbort); + this.pendingDelegation = undefined; + this.activeDelegationId = undefined; + this.partialTranscriptRole = undefined; + this.transcript = []; + } + private async runDelegation(delegation: PendingDelegation, signal: AbortSignal): Promise { let text: string; try { diff --git a/extensions/openai/realtime-quicksilver-gateway-bridge-lifecycle.test.ts b/extensions/openai/realtime-quicksilver-gateway-bridge-lifecycle.test.ts new file mode 100644 index 000000000000..d1bf2a8ab4ef --- /dev/null +++ b/extensions/openai/realtime-quicksilver-gateway-bridge-lifecycle.test.ts @@ -0,0 +1,114 @@ +import { describe, expect, it, vi } from "vitest"; +import { OpenAIQuicksilverGatewayBridge } from "./realtime-quicksilver-gateway-bridge.js"; +import { + createCallResponse, + emitSideband, + FakeSocket, + parseSent, +} from "./realtime-quicksilver.test-helpers.js"; + +function createBridge(params: { + runAgentConsult: (request: { prompt: string; signal?: AbortSignal }) => Promise<{ text: string }>; +}) { + let socket: FakeSocket | undefined; + const bridge = new OpenAIQuicksilverGatewayBridge({ + providerConfig: {}, + model: "gpt-live-1-boulder-alpha", + voice: "marin", + audioFormat: { encoding: "pcm16", sampleRateHz: 24_000, channels: 1 }, + onAudio: vi.fn(), + onClearAudio: vi.fn(), + runAgentConsult: params.runAgentConsult, + logger: { debug: vi.fn(), warn: vi.fn() }, + resolveAuth: vi.fn(async () => ({ + type: "api-key" as const, + token: "platform-key", + })), + createPeer: vi.fn(async () => ({ + createOffer: vi.fn(async () => "v=offer\r\n"), + applyAnswer: vi.fn(async () => undefined), + adoptPendingAudio: vi.fn(), + sendAudio: vi.fn(), + close: vi.fn(), + })), + fetchImpl: vi.fn(async () => createCallResponse("v=answer\r\n", "rtc_lifecycle")), + webSocketFactory: () => { + socket = new FakeSocket(); + return socket; + }, + }); + return { + bridge, + getSocket: () => { + if (!socket) { + throw new Error("expected sideband socket"); + } + return socket; + }, + }; +} + +function emitDelegation(socket: FakeSocket, id: string, text: string): void { + emitSideband(socket, { + type: "delegation.created", + item: { + type: "delegation", + target: "client", + id, + content: [{ type: "input_text", text }], + }, + }); +} + +describe("OpenAI Quicksilver gateway bridge lifecycle", () => { + it("aborts an accepted delegation when the bridge closes normally", async () => { + let consultSignal: AbortSignal | undefined; + const runAgentConsult = vi.fn(async ({ signal }: { prompt: string; signal?: AbortSignal }) => { + consultSignal = signal; + await new Promise((resolve) => { + signal?.addEventListener("abort", () => resolve(), { once: true }); + }); + return { text: "must not be delivered" }; + }); + const harness = createBridge({ runAgentConsult }); + + await harness.bridge.connect(); + const socket = harness.getSocket(); + emitDelegation(socket, "delegation-abort", "Cancel this on close"); + await vi.waitFor(() => expect(runAgentConsult).toHaveBeenCalledOnce()); + + harness.bridge.close(); + expect(consultSignal?.aborted).toBe(true); + await Promise.resolve(); + expect(parseSent(socket).filter((event) => event.type === "delegation.context.append")).toEqual( + [], + ); + }); + + it("detaches transport without aborting an accepted delegation", async () => { + let consultSignal: AbortSignal | undefined; + let resolveConsult!: (result: { text: string }) => void; + const consultResult = new Promise<{ text: string }>((resolve) => { + resolveConsult = resolve; + }); + const runAgentConsult = vi.fn(async ({ signal }: { prompt: string; signal?: AbortSignal }) => { + consultSignal = signal; + return await consultResult; + }); + const harness = createBridge({ runAgentConsult }); + + await harness.bridge.connect(); + const socket = harness.getSocket(); + emitDelegation(socket, "delegation-detach", "Finish after disconnect"); + await vi.waitFor(() => expect(runAgentConsult).toHaveBeenCalledOnce()); + + harness.bridge.close({ disposition: "detach" }); + expect(consultSignal?.aborted).toBe(false); + resolveConsult({ text: "finished after detach" }); + await Promise.resolve(); + await Promise.resolve(); + expect(parseSent(socket).filter((event) => event.type === "delegation.context.append")).toEqual( + [], + ); + }); +}); diff --git a/extensions/openai/realtime-quicksilver-gateway-bridge.ts b/extensions/openai/realtime-quicksilver-gateway-bridge.ts index e6740e0acfbc..21eaf4270452 100644 --- a/extensions/openai/realtime-quicksilver-gateway-bridge.ts +++ b/extensions/openai/realtime-quicksilver-gateway-bridge.ts @@ -5,6 +5,8 @@ import type { PluginLogger } from "openclaw/plugin-sdk/plugin-entry"; import type { RealtimeVoiceBridge, RealtimeVoiceBridgeCreateRequest, + RealtimeVoiceCloseDisposition, + RealtimeVoiceCloseOptions, } from "openclaw/plugin-sdk/realtime-voice"; import WebSocket, { type RawData } from "ws"; import { OpenAIQuicksilverPendingAudio } from "./realtime-quicksilver-audio-buffer.js"; @@ -180,8 +182,8 @@ export class OpenAIQuicksilverGatewayBridge implements RealtimeVoiceBridge { acknowledgeMark(): void {} - close(): void { - this.teardown("completed"); + close(options?: RealtimeVoiceCloseOptions): void { + this.teardown("completed", undefined, options?.disposition ?? "abort"); } isConnected(): boolean { @@ -298,7 +300,7 @@ export class OpenAIQuicksilverGatewayBridge implements RealtimeVoiceBridge { throw new Error(describeSidebandClose(terminalEvent.code, reason)); } } catch (error) { - this.releaseResources(); + this.releaseResources("abort"); throw toErrorObject(error, "OpenAI GPT-Live gateway relay failed"); } } @@ -335,14 +337,18 @@ export class OpenAIQuicksilverGatewayBridge implements RealtimeVoiceBridge { this.teardown("error", () => this.config.onError?.(error)); } - private teardown(reason: "completed" | "error", beforeClose?: () => void): void { + private teardown( + reason: "completed" | "error", + beforeClose?: () => void, + disposition: RealtimeVoiceCloseDisposition = "abort", + ): void { if (this.closed) { return; } // Claim terminal ownership and release resources before callbacks so reentrant close // cannot replace the outcome, while finally preserves error-before-close ordering. this.closed = true; - this.releaseResources(); + this.releaseResources(disposition); try { beforeClose?.(); } finally { @@ -353,12 +359,16 @@ export class OpenAIQuicksilverGatewayBridge implements RealtimeVoiceBridge { } } - private releaseResources(): void { + private releaseResources(disposition: RealtimeVoiceCloseDisposition): void { releaseOpenAIQuicksilverSession(this); this.connected = false; this.pendingAudio.clear(); + if (disposition === "detach") { + this.delegations?.detach(); + } else { + this.delegations?.stop(new Error("GPT-Live delegation stopped")); + } this.abortController.abort(new Error("GPT-Live gateway relay bridge closed")); - this.delegations?.stop(new Error("GPT-Live delegation stopped")); if (this.timer) { clearTimeout(this.timer); this.timer = undefined; diff --git a/src/gateway/talk-realtime-relay-operations.ts b/src/gateway/talk-realtime-relay-operations.ts index 11395057fb76..ba9999645102 100644 --- a/src/gateway/talk-realtime-relay-operations.ts +++ b/src/gateway/talk-realtime-relay-operations.ts @@ -4,7 +4,10 @@ import { type RealtimeVoiceAgentControlResult, } from "../talk/agent-run-control.js"; import { registerClientVoiceConsultRun } from "../talk/client-voice-session.js"; -import type { RealtimeVoiceToolResultOptions } from "../talk/provider-types.js"; +import type { + RealtimeVoiceCloseOptions, + RealtimeVoiceToolResultOptions, +} from "../talk/provider-types.js"; import type { TalkEvent } from "../talk/talk-session-controller.js"; import { abortChatRunById } from "./chat-abort.js"; import { formatError } from "./server-utils.js"; @@ -78,6 +81,12 @@ export function abortRelayAgentRuns(session: RelaySession, reason: string): void session.activeAgentToolCalls.clear(); } +/** Releases relay-local correlation without cancelling durable voice-bound agent runs. */ +function detachRelayAgentRuns(session: RelaySession): void { + session.activeAgentRuns.clear(); + session.activeAgentToolCalls.clear(); +} + export function pruneInactiveRelayAgentRuns(session: RelaySession): number { for (const runId of session.activeAgentRuns.keys()) { if (!session.context.chatAbortControllers.has(runId)) { @@ -92,14 +101,23 @@ export function pruneInactiveRelayAgentRuns(session: RelaySession): number { return session.activeAgentRuns.size; } -export function closeRelaySession(session: RelaySession, reason: "completed" | "error"): void { +export function closeRelaySession( + session: RelaySession, + reason: "completed" | "error", + options?: RealtimeVoiceCloseOptions, +): void { + const disposition = options?.disposition ?? "abort"; session.harness.close(); relaySessions.delete(session.id); forgetUnifiedTalkSession(session.id); clearTimeout(session.cleanupTimer); - abortRelayAgentRuns(session, reason === "error" ? "relay-error" : "relay-closed"); + if (disposition === "detach") { + detachRelayAgentRuns(session); + } else { + abortRelayAgentRuns(session, reason === "error" ? "relay-error" : "relay-closed"); + } try { - session.bridge.close(); + session.bridge.close({ disposition }); } finally { // Provider teardown may throw, but the relay must still reach its durable // voice and owner-visible terminal state before that error is surfaced. @@ -122,7 +140,7 @@ export function closeTalkRealtimeRelaySessionsForConnection(connId: string): voi closeTalkRelaySessionsForConnection({ sessions: relaySessions.values(), connId, - closeSession: (session) => closeRelaySession(session, "completed"), + closeSession: (session) => closeRelaySession(session, "completed", { disposition: "detach" }), onCloseError: (error, session) => { session.context.logGateway.warn( `failed to close realtime relay session after connection disconnect: ${formatError(error)}`, diff --git a/src/gateway/talk-realtime-relay.test.ts b/src/gateway/talk-realtime-relay.test.ts index 6483f11b95e9..9e138a71d637 100644 --- a/src/gateway/talk-realtime-relay.test.ts +++ b/src/gateway/talk-realtime-relay.test.ts @@ -3887,6 +3887,55 @@ describe("talk realtime gateway relay", () => { expectNodeAbortPayload(nodeSendToSession); }); + it("detaches linked agent consult runs when the gateway connection closes", () => { + const close = vi.fn(); + const provider = createIdleRelayProvider(); + provider.createBridge = () => makeRelayTransport({ close }); + const { abortController, session } = createAbortableRelayRunFixture(provider); + + cleanupTalkConnection("conn-1", { warn: vi.fn() }); + + expect(close).toHaveBeenCalledOnce(); + expect(abortController.signal.aborted).toBe(false); + expect(relaySessions.has(session.relaySessionId)).toBe(false); + }); + + it.each([ + { + name: "expires", + close: (session: { relaySessionId: string }) => { + const relay = relaySessions.get(session.relaySessionId); + if (!relay) { + throw new Error("expected active relay"); + } + relay.expiresAtMs = Date.now() - 1; + sendTalkRealtimeRelayAudio({ + relaySessionId: session.relaySessionId, + connId: "conn-1", + audioBase64: "AQI=", + }); + }, + }, + { + name: "fails connection ownership", + close: (session: { relaySessionId: string }) => { + sendTalkRealtimeRelayAudio({ + relaySessionId: session.relaySessionId, + connId: "conn-other", + audioBase64: "AQI=", + }); + }, + }, + ])("aborts linked agent consult runs when the relay $name", ({ close }) => { + const fixture = createAbortableRelayRunFixture(); + + expect(() => close(fixture.session)).toThrow("Unknown realtime relay session"); + + expect(fixture.abortController.signal.aborted).toBe(true); + expectChatAbortPayload(fixture.broadcast, "relay-closed"); + expectNodeAbortPayload(fixture.nodeSendToSession); + }); + it("aborts linked agent consult runs when the provider closes the relay", () => { const abortController = new AbortController(); let bridgeRequest: RealtimeVoiceBridgeCreateRequest | undefined; diff --git a/src/plugin-sdk/realtime-voice.ts b/src/plugin-sdk/realtime-voice.ts index 26e2b38a31f7..55234c1e53cb 100644 --- a/src/plugin-sdk/realtime-voice.ts +++ b/src/plugin-sdk/realtime-voice.ts @@ -8,6 +8,8 @@ export type { RealtimeVoiceBridge, RealtimeVoiceBridgeCallbacks, RealtimeVoiceBridgeEvent, + RealtimeVoiceCloseDisposition, + RealtimeVoiceCloseOptions, RealtimeVoiceBrowserSession, RealtimeVoiceBrowserSessionCreateRequest, RealtimeVoiceGatewayControl, diff --git a/src/talk/provider-types.ts b/src/talk/provider-types.ts index 4214333c16fc..ed44a7545d3a 100644 --- a/src/talk/provider-types.ts +++ b/src/talk/provider-types.ts @@ -79,6 +79,13 @@ export type RealtimeVoiceToolResultOptions = { willContinue?: boolean; }; +export type RealtimeVoiceCloseDisposition = "abort" | "detach"; + +export type RealtimeVoiceCloseOptions = { + /** Whether closing the transport also cancels work already accepted by the host. */ + disposition?: RealtimeVoiceCloseDisposition; +}; + export type RealtimeVoiceBridgeEvent = { direction: "client" | "server"; type: string; @@ -333,7 +340,7 @@ export type RealtimeVoiceBridge = { options?: RealtimeVoiceToolResultOptions, ): void | Promise; acknowledgeMark(markName?: string): void; - close(): void; + close(options?: RealtimeVoiceCloseOptions): void; isConnected(): boolean; }; diff --git a/src/talk/session-runtime.test.ts b/src/talk/session-runtime.test.ts index fdc9368aad26..4c8cc7eb4d4e 100644 --- a/src/talk/session-runtime.test.ts +++ b/src/talk/session-runtime.test.ts @@ -394,6 +394,25 @@ describe("realtime voice bridge session runtime", () => { expect(onError).not.toHaveBeenCalled(); }); + it("forwards the close disposition to the provider bridge", () => { + const close = vi.fn(); + const provider: RealtimeVoiceProviderPlugin = { + id: "test", + label: "Test", + isConfigured: () => true, + createBridge: () => makeBridge({ close }), + }; + const session = createRealtimeVoiceBridgeSession({ + provider, + providerConfig: {}, + audioSink: { sendAudio: vi.fn() }, + }); + + session.close({ disposition: "detach" }); + + expect(close).toHaveBeenCalledWith({ disposition: "detach" }); + }); + it("permanently closes once while preserving synchronous transcript flush", async () => { let callbacks: Parameters[0] | undefined; const close = vi.fn(() => { diff --git a/src/talk/session-runtime.ts b/src/talk/session-runtime.ts index b95bdd88a329..5117edd39c3d 100644 --- a/src/talk/session-runtime.ts +++ b/src/talk/session-runtime.ts @@ -6,6 +6,7 @@ import type { RealtimeVoiceAudioClearReason, RealtimeVoiceAudioFormat, RealtimeVoiceBargeInOptions, + RealtimeVoiceCloseOptions, RealtimeVoiceCloseReason, RealtimeVoiceBridgeEvent, RealtimeVoiceProviderConfig, @@ -37,7 +38,7 @@ export type RealtimeVoiceMarkStrategy = "transport" | "ack-immediately" | "ignor export type RealtimeVoiceBridgeSession = { bridge: RealtimeVoiceBridge; acknowledgeMark(markName?: string): void; - close(): void; + close(options?: RealtimeVoiceCloseOptions): void; connect(): Promise; sendAudio(audio: Buffer): void; sendUserMessage(text: string): void; @@ -110,13 +111,13 @@ export function createRealtimeVoiceBridgeSession( return requireBridge(); }, acknowledgeMark: (markName) => requireBridge().acknowledgeMark(markName), - close: () => { + close: (options) => { if (phase === "disposed") { return; } const bridge = requireBridge(); phase = "disposed"; - bridge.close(); + bridge.close(options); }, connect: () => { if (phase === "disposed") {