From 635d3967552c94feb58c8adcb6eba52cd5fe2511 Mon Sep 17 00:00:00 2001 From: Dallin Romney Date: Thu, 23 Jul 2026 06:40:32 +0900 Subject: [PATCH] refactor(talk): run the Gateway realtime relay through the shared session harness (#112590) * refactor(talk): adopt session harness in gateway relay * fix(talk): preserve relay harness behavior * style(talk): format relay barge-in call --------- Co-authored-by: Peter Steinberger --- src/gateway/talk-realtime-relay.ts | 147 ++++++++++------------ src/talk/realtime-session-harness.test.ts | 31 +++++ src/talk/realtime-session-harness.ts | 20 +-- 3 files changed, 113 insertions(+), 85 deletions(-) diff --git a/src/gateway/talk-realtime-relay.ts b/src/gateway/talk-realtime-relay.ts index 225c756f36bc..2add83376eaa 100644 --- a/src/gateway/talk-realtime-relay.ts +++ b/src/gateway/talk-realtime-relay.ts @@ -24,12 +24,7 @@ import { registerClientVoiceConsultRun, } from "../talk/client-voice-session.js"; import { readSpeakableRealtimeVoiceToolResult } from "../talk/consult-question.js"; -import { - createRealtimeVoiceForcedConsultCoordinator, - type RealtimeVoiceForcedConsultCoordinator, - type RealtimeVoiceForcedConsultHandle, -} from "../talk/forced-consult-coordinator.js"; -import { recordTalkObservabilityEvent } from "../talk/observability.js"; +import type { RealtimeVoiceForcedConsultHandle } from "../talk/forced-consult-coordinator.js"; import { REALTIME_VOICE_AUDIO_FORMAT_PCM16_24KHZ, type RealtimeVoiceBrowserAudioContract, @@ -39,20 +34,11 @@ import { type RealtimeVoiceToolResultOptions, } from "../talk/provider-types.js"; import { - isLikelyRealtimeVoiceAssistantEchoTranscript, - recordRealtimeVoiceTranscript, - type RealtimeVoiceTranscriptEntry, -} from "../talk/session-log-runtime.js"; -import { - createRealtimeVoiceBridgeSession, - type RealtimeVoiceBridgeSession, -} from "../talk/session-runtime.js"; -import { - type TalkEvent, - type TalkEventInput, - type TalkSessionController, - createTalkSessionController, -} from "../talk/talk-session-controller.js"; + createRealtimeVoiceSessionHarness, + type RealtimeVoiceSessionHarness, +} from "../talk/realtime-session-harness.js"; +import type { RealtimeVoiceBridgeSession } from "../talk/session-runtime.js"; +import type { TalkEvent, TalkEventInput } from "../talk/talk-session-controller.js"; import { abortChatRunById } from "./chat-abort.js"; import type { GatewayRequestContext } from "./server-methods/shared-types.js"; import { @@ -73,6 +59,7 @@ const RELAY_EVENT = "talk.event"; const RELAY_TRANSCRIPT_ECHO_LOOKBACK_MS = 12_000; const FORCED_CONSULT_FALLBACK_DELAY_MS = 200; const FORCED_CONSULT_RESULT_MAX_CHARS = 1_800; +const noFallbackRelayOutputFlush = () => {}; type TalkRealtimeRelayEventPayload = | { relaySessionId: string; type: "ready" } @@ -136,7 +123,7 @@ type RelaySession = { connId: string; context: GatewayRequestContext; bridge: RealtimeVoiceBridgeSession; - talk: TalkSessionController; + harness: RealtimeVoiceSessionHarness; sessionKey?: string; expiresAtMs: number; cleanupTimer: ReturnType; @@ -160,8 +147,6 @@ type RelaySession = { forcedTerminalProviderResults: Map; // Turn cancellation invalidates async acceptance callbacks from the prior turn. toolResultEpoch: number; - forcedConsults: RealtimeVoiceForcedConsultCoordinator; - transcript: RealtimeVoiceTranscriptEntry[]; voiceConfig?: OpenClawConfig; voiceSessionCreated: boolean; voiceTranscriptSeq: number; @@ -328,14 +313,7 @@ function isWorkingToolResult(result: unknown): boolean { } function isRelayAssistantEchoTranscript(session: RelaySession | undefined, text: string): boolean { - if (!session) { - return false; - } - return isLikelyRealtimeVoiceAssistantEchoTranscript({ - transcript: session.transcript, - text, - lookbackMs: RELAY_TRANSCRIPT_ECHO_LOOKBACK_MS, - }); + return session?.harness.isLikelyAssistantEchoTranscript(text) ?? false; } function buildForcedConsultCheckingPrompt(): string { return [ @@ -369,8 +347,8 @@ function suppressedToolResultOptions( } function cancelForcedConsults(session: RelaySession): void { - for (const handle of session.forcedConsults.handles()) { - session.forcedConsults.markCancelled(handle); + for (const handle of session.harness.forcedConsults.handles()) { + session.harness.forcedConsults.markCancelled(handle); } } @@ -437,7 +415,7 @@ function broadcastToolResultToOwner( relaySessionId: session.id, type: "toolResult", callId: params.callId, - talkEvent: session.talk.emit({ + talkEvent: session.harness.talk.emit({ type: "tool.result", callId: params.callId, turnId: params.turnId, @@ -608,7 +586,7 @@ function submitRelayAgentControlProviderResults( return; } if (forcedConsult) { - session.forcedConsults.markCancelled(forcedConsult); + session.harness.forcedConsults.markCancelled(forcedConsult); } broadcastToolResultToOwner(session, { callId, @@ -620,9 +598,11 @@ function submitRelayAgentControlProviderResults( session.completedAgentToolCalls.add(callId); }; for (const callId of activeCallIds) { - const forcedConsult = session.forcedConsults.handles().find((handle) => handle.id === callId); + const forcedConsult = session.harness.forcedConsults + .handles() + .find((handle) => handle.id === callId); if (forcedConsult) { - const nativeCallIds = session.forcedConsults.nativeCallIds(forcedConsult); + const nativeCallIds = session.harness.forcedConsults.nativeCallIds(forcedConsult); providerResponseStarted ||= toolResultOptions === undefined && nativeCallIds.length > 0; const terminal: ForcedTerminalProviderResult = { result: providerResult, @@ -667,7 +647,7 @@ function submitRelayAgentControlProviderResults( } function closeRelaySession(session: RelaySession, reason: "completed" | "error"): void { - session.forcedConsults.clear(); + session.harness.close(); relaySessions.delete(session.id); forgetUnifiedTalkSession(session.id); clearTimeout(session.cleanupTimer); @@ -678,7 +658,7 @@ function closeRelaySession(session: RelaySession, reason: "completed" | "error") relaySessionId: session.id, type: "close", reason, - talkEvent: session.talk.emit({ + talkEvent: session.harness.talk.emit({ type: "session.closed", payload: { reason }, final: true, @@ -725,23 +705,34 @@ export function createTalkRealtimeRelaySession( if (expiresAtMs === undefined) { throw new Error("Realtime relay session expiry is outside the supported Date range"); } - const talk = createTalkSessionController( - { + const harness = createRealtimeVoiceSessionHarness({ + talk: { sessionId: relaySessionId, mode: "realtime", transport: "gateway-relay", brain: "agent-consult", provider: params.provider.id, + // Keep the pre-harness steering window; other harness consumers use the shared default. + maxRecentEvents: 20, }, - { onEvent: recordTalkObservabilityEvent }, - ); + talkPayloads: { + turnStarted: () => ({}), + turnEnded: (reason) => ({ reason }), + inputAudioDelta: (audio) => ({ byteLength: audio.byteLength }), + outputAudioStarted: () => ({}), + outputAudioDelta: (audio) => ({ byteLength: audio.byteLength }), + outputAudioDone: (reason) => ({ reason }), + }, + transcriptLookbackMs: RELAY_TRANSCRIPT_ECHO_LOOKBACK_MS, + captureBridgeEvents: false, + }); const emit = (event: TalkRealtimeRelayEventPayload, talkEvent?: TalkEventInput) => broadcastToOwner( params.context, params.connId, { ...event, - ...(talkEvent ? { talkEvent: talk.emit(talkEvent) } : {}), + ...(talkEvent ? { talkEvent: harness.emit(talkEvent) } : {}), }, relayEventDeliveryOptions(event), ); @@ -750,7 +741,7 @@ export function createTalkRealtimeRelaySession( let ready = false; let failureEmitted = false; const relayRef: { current?: RelaySession } = {}; - const bridge = createRealtimeVoiceBridgeSession({ + const bridge = harness.createBridge({ provider: params.provider, cfg: params.cfg, providerConfig: params.providerConfig, @@ -843,7 +834,6 @@ export function createTalkRealtimeRelaySession( const relay = relayRef.current; const turnId = relay ? ensureRelayTurn(relay) : undefined; if (final && relay) { - recordRealtimeVoiceTranscript(relay.transcript, role, text); enqueueRelayVoiceTranscript(relay, role, text); } const eventType = @@ -907,13 +897,13 @@ export function createTalkRealtimeRelaySession( const relay = relayRef.current; let shouldSubmitWorkingResult = false; if (relay && toolCall.name === REALTIME_VOICE_AGENT_CONSULT_TOOL_NAME) { - const forcedConsult = relay.forcedConsults.recordNativeConsult( + const forcedConsult = relay.harness.forcedConsults.recordNativeConsult( toolCall.args, toolCall.callId, ); if (forcedConsult.kind === "in_flight" || forcedConsult.kind === "already_delivered") { if (forcedConsult.kind === "already_delivered") { - const result = relay.forcedConsults.isCancelled(forcedConsult.handle) + const result = relay.harness.forcedConsults.isCancelled(forcedConsult.handle) ? buildRealtimeVoiceAgentCancelProviderResult( "OpenClaw cancelled this consult before completion. Do not restart it.", ) @@ -977,7 +967,7 @@ export function createTalkRealtimeRelaySession( if (!active) { return; } - active.forcedConsults.clear(); + active.harness.close(); relaySessions.delete(relaySessionId); forgetUnifiedTalkSession(relaySessionId); clearTimeout(active.cleanupTimer); @@ -1007,7 +997,7 @@ export function createTalkRealtimeRelaySession( connId: params.connId, context: params.context, bridge, - talk, + harness, sessionKey: params.sessionKey?.trim() || undefined, expiresAtMs, cleanupTimer: setTimeout(() => { @@ -1027,8 +1017,6 @@ export function createTalkRealtimeRelaySession( pendingWorkingToolResults: new Map(), forcedTerminalProviderResults: new Map(), toolResultEpoch: 0, - forcedConsults: createRealtimeVoiceForcedConsultCoordinator(), - transcript: [], ...(params.cfg ? { voiceConfig: params.cfg } : {}), voiceSessionCreated: false, voiceTranscriptSeq: 0, @@ -1077,23 +1065,26 @@ function scheduleForcedAgentConsult(session: RelaySession | undefined, question: if (!session || !question.trim()) { return; } - if (session.forcedConsults.hasRecentNativeConsult(question)) { + if (session.harness.forcedConsults.hasRecentNativeConsult(question)) { return; } - session.forcedConsults.clearPending(); - const handle = session.forcedConsults.prepare(question); + session.harness.forcedConsults.clearPending(); + const handle = session.harness.forcedConsults.prepare(question); if (!handle) { return; } - session.forcedConsults.schedule(handle, FORCED_CONSULT_FALLBACK_DELAY_MS, () => { + session.harness.forcedConsults.schedule(handle, FORCED_CONSULT_FALLBACK_DELAY_MS, () => { if (!relaySessions.has(session.id)) { return; } const turnId = ensureRelayTurn(session); const callId = handle.id; const itemId = `forced-consult-item-${randomUUID()}`; - session.forcedConsults.markStarted(handle); - session.bridge.handleBargeIn({ audioPlaybackActive: true, force: true }); + session.harness.forcedConsults.markStarted(handle); + session.harness.handleBargeIn( + { audioPlaybackActive: true, force: true }, + noFallbackRelayOutputFlush, + ); broadcastToOwner(session.context, session.connId, { relaySessionId: session.id, type: "toolCall", @@ -1107,7 +1098,7 @@ function scheduleForcedAgentConsult(session: RelaySession | undefined, question: "The realtime provider produced a final user transcript without invoking openclaw_agent_consult, so OpenClaw is forcing the consult for realtime Talk.", responseStyle: "Reply in a concise spoken tone.", }, - talkEvent: session.talk.emit({ + talkEvent: session.harness.talk.emit({ type: "tool.call", itemId, callId, @@ -1144,7 +1135,7 @@ function drainForcedTerminalProviderResults( if (session.forcedTerminalProviderResults.get(handle.id) !== terminal) { return; } - const submissions = session.forcedConsults + const submissions = session.harness.forcedConsults .nativeCallIds(handle) .map((callId) => submitForcedConsultProviderResult(session, callId, terminal.result, terminal.options), @@ -1157,7 +1148,7 @@ function drainForcedTerminalProviderResults( drainForcedTerminalProviderResults(session, handle, terminal), ); } - const hasUnsubmittedCall = session.forcedConsults + const hasUnsubmittedCall = session.harness.forcedConsults .nativeCallIds(handle) .some((callId) => !session.completedProviderToolResults.has(callId)); if (hasUnsubmittedCall) { @@ -1170,7 +1161,7 @@ function drainForcedTerminalProviderResultsAfterPending( handle: RealtimeVoiceForcedConsultHandle, terminal: ForcedTerminalProviderResult, ): void | Promise { - const pending = session.forcedConsults + const pending = session.harness.forcedConsults .nativeCallIds(handle) .map((callId) => session.pendingProviderToolResults.get(callId)) .filter((submission): submission is Promise => submission !== undefined); @@ -1204,7 +1195,7 @@ function submitRealtimeAgentConsultWorkingResponse( relaySessionId: session.id, type: "toolResult", callId, - talkEvent: session.talk.emit({ + talkEvent: session.harness.talk.emit({ type: "tool.progress", callId, turnId, @@ -1216,7 +1207,7 @@ function submitRealtimeAgentConsultWorkingResponse( } function ensureRelayTurn(session: RelaySession): string { - const turn = session.talk.ensureTurn(); + const turn = session.harness.talk.ensureTurn(); if (turn.event) { broadcastToOwner(session.context, session.connId, { relaySessionId: session.id, @@ -1256,7 +1247,7 @@ export function sendTalkRealtimeRelayAudio(params: { relaySessionId: session.id, type: "inputAudio", byteLength: audio.byteLength, - talkEvent: session.talk.emit({ + talkEvent: session.harness.talk.emit({ type: "input.audio.delta", turnId, payload: { byteLength: audio.byteLength }, @@ -1294,13 +1285,13 @@ export function submitTalkRealtimeRelayToolResult(params: { if (pendingFinal && !cancelledAgentCall) { return pendingFinal; } - const forcedConsult = session.forcedConsults + const forcedConsult = session.harness.forcedConsults .handles() .find((handle) => handle.id === params.callId); if (forcedConsult) { - const cancelled = session.forcedConsults.isCancelled(forcedConsult); + const cancelled = session.harness.forcedConsults.isCancelled(forcedConsult); const turnId = cancelled - ? (session.cancelledAgentToolCalls.get(params.callId) ?? session.talk.activeTurnId) + ? (session.cancelledAgentToolCalls.get(params.callId) ?? session.harness.talk.activeTurnId) : ensureRelayTurn(session); if (!turnId) { throw new Error("Cancelled realtime consult is missing its original turn"); @@ -1331,7 +1322,7 @@ export function submitTalkRealtimeRelayToolResult(params: { if (session.toolResultEpoch !== terminal.epoch) { return; } - session.forcedConsults.markCancelled(forcedConsult); + session.harness.forcedConsults.markCancelled(forcedConsult); clearRelayAgentToolCall(session, params.callId); session.cancelledAgentToolCalls.delete(params.callId); session.completedAgentToolCalls.add(params.callId); @@ -1382,10 +1373,10 @@ export function submitTalkRealtimeRelayToolResult(params: { if (session.toolResultEpoch !== terminal.epoch) { return; } - session.forcedConsults.markDelivered(forcedConsult); + session.harness.forcedConsults.markDelivered(forcedConsult); clearRelayAgentToolCall(session, params.callId); session.completedAgentToolCalls.add(params.callId); - const hasNativeCalls = session.forcedConsults.nativeCallIds(forcedConsult).length > 0; + const hasNativeCalls = session.harness.forcedConsults.nativeCallIds(forcedConsult).length > 0; if (text && (!hasNativeCalls || providerOptions)) { session.bridge.sendUserMessage(buildForcedConsultSpeechPrompt(text)); } @@ -1537,7 +1528,7 @@ export async function steerTalkRealtimeRelayAgentRun(params: { sessionKey, text: params.text, mode: params.mode, - recentEvents: session.talk.recentEvents, + recentEvents: session.harness.talk.recentEvents, }); if (relaySessions.get(session.id) !== session) { throw new Error("Realtime relay session closed while steering the agent run"); @@ -1557,7 +1548,7 @@ export async function steerTalkRealtimeRelayAgentRun(params: { relaySessionId: session.id, type: "toolProgress", result: finalResult, - talkEvent: session.talk.emit({ + talkEvent: session.harness.talk.emit({ type: "tool.progress", turnId, payload: { @@ -1586,17 +1577,17 @@ export function cancelTalkRealtimeRelayTurn(params: { for (const callId of session.activeAgentToolCalls.keys()) { session.cancelledAgentToolCalls.set(callId, turnId); } - for (const forcedConsult of session.forcedConsults.handles()) { - if (session.forcedConsults.isCancelled(forcedConsult)) { + for (const forcedConsult of session.harness.forcedConsults.handles()) { + if (session.harness.forcedConsults.isCancelled(forcedConsult)) { session.cancelledAgentToolCalls.set(forcedConsult.id, turnId); - for (const nativeCallId of session.forcedConsults.nativeCallIds(forcedConsult)) { + for (const nativeCallId of session.harness.forcedConsults.nativeCallIds(forcedConsult)) { session.cancelledAgentToolCalls.set(nativeCallId, turnId); } } } - session.bridge.handleBargeIn({ audioPlaybackActive: true }); + session.harness.handleBargeIn({ audioPlaybackActive: true }, noFallbackRelayOutputFlush); abortRelayAgentRuns(session, reason); - const cancelled = session.talk.cancelTurn({ + const cancelled = session.harness.talk.cancelTurn({ turnId, payload: { reason }, }); diff --git a/src/talk/realtime-session-harness.test.ts b/src/talk/realtime-session-harness.test.ts index 56c5f68f4338..3769f050ff22 100644 --- a/src/talk/realtime-session-harness.test.ts +++ b/src/talk/realtime-session-harness.test.ts @@ -64,6 +64,28 @@ describe("realtime voice session harness", () => { expect(harness.talk.recentEvents.map((event) => event.seq)).toEqual([1, 2, 3, 4, 5, 6]); }); + it("honors a caller-specific recent Talk event limit", () => { + const harness = createHarness({ + talk: { + sessionId: "limited-session", + mode: "realtime", + transport: "gateway-relay", + brain: "agent-consult", + provider: "test", + maxRecentEvents: 2, + }, + }); + + harness.emit({ type: "session.started", payload: {} }); + harness.emit({ type: "session.ready", payload: {} }); + harness.emit({ type: "session.closed", payload: {}, final: true }); + + expect(harness.talk.recentEvents.map((event) => event.type)).toEqual([ + "session.ready", + "session.closed", + ]); + }); + it("suppresses input through queued output playback plus the echo tail", () => { vi.useFakeTimers(); vi.setSystemTime(1_000); @@ -120,6 +142,15 @@ describe("realtime voice session harness", () => { expect(deliver).toHaveBeenCalledWith("answer:first\nsecond"); }); + it("detects assistant transcript echo without enabling audio suppression", () => { + const harness = createHarness({ transcriptLookbackMs: 12_000 }); + + harness.recordTranscript("assistant", "I found the shopping list"); + + expect(harness.isLikelyAssistantEchoTranscript("I found the shopping list")).toBe(true); + expect(harness.recordInputAudio(Buffer.from([1, 2]))).toBe(true); + }); + it("flushes transport output when provider barge-in does not clear it", () => { const handleBargeIn = vi.fn(); const provider: RealtimeVoiceProviderPlugin = { diff --git a/src/talk/realtime-session-harness.ts b/src/talk/realtime-session-harness.ts index 9fd51bf4fd31..94041a5bdafc 100644 --- a/src/talk/realtime-session-harness.ts +++ b/src/talk/realtime-session-harness.ts @@ -107,6 +107,8 @@ export function createRealtimeVoiceSessionHarness; forcedConsults?: RealtimeVoiceForcedConsultCoordinatorOptions; echoSuppression?: RealtimeVoiceSessionHarnessEchoSuppression; + transcriptLookbackMs?: number; + captureBridgeEvents?: boolean; }): RealtimeVoiceSessionHarness { let closed = false; let bridge: RealtimeVoiceBridgeSession | undefined; @@ -121,11 +123,13 @@ export function createRealtimeVoiceSessionHarness( params.forcedConsults, ); const talk = createTalkSessionController( - { ...params.talk, maxRecentEvents: 40 }, + { maxRecentEvents: 40, ...params.talk }, { onEvent: (event) => { recordTalkObservabilityEvent(event); @@ -173,7 +177,9 @@ export function createRealtimeVoiceSessionHarness { - recordRealtimeVoiceBridgeEvent(bridgeEvents, event); + if (params.captureBridgeEvents !== false) { + recordRealtimeVoiceBridgeEvent(bridgeEvents, event); + } bridgeParams.onEvent?.(event); }, }); @@ -223,13 +229,13 @@ export function createRealtimeVoiceSessionHarness