From 5fa2c2d98fac5b264d84251acef8298e8615aaf7 Mon Sep 17 00:00:00 2001 From: Vincent Koc Date: Tue, 11 Aug 2026 13:45:23 +0800 Subject: [PATCH] refactor(talk): internalize authoritative relay events Co-authored-by: Dallin Romney Punchcard-Session: golden-meadow-cedar-dv --- src/gateway/talk-realtime-relay-operations.ts | 17 +- .../talk-realtime-relay-session-create.ts | 50 +++--- src/gateway/talk-realtime-relay-state.ts | 31 ++-- src/gateway/talk-realtime-relay.test.ts | 77 +++++++- src/plugin-sdk/realtime-voice.test.ts | 41 ++++- src/talk/realtime-session-harness.test.ts | 83 ++++++--- src/talk/realtime-session-harness.ts | 169 +++++++++++++++--- .../chat/realtime-talk-gateway-relay-types.ts | 1 + .../pages/chat/realtime-talk-gateway-relay.ts | 4 +- 9 files changed, 374 insertions(+), 99 deletions(-) diff --git a/src/gateway/talk-realtime-relay-operations.ts b/src/gateway/talk-realtime-relay-operations.ts index 11395057fb76..44efdc44f871 100644 --- a/src/gateway/talk-realtime-relay-operations.ts +++ b/src/gateway/talk-realtime-relay-operations.ts @@ -26,6 +26,7 @@ import { MAX_AUDIO_BASE64_BYTES, MAX_RELAY_SESSIONS_GLOBAL, MAX_RELAY_SESSIONS_PER_CONN, + broadcastRelayTurnStarted, broadcastToOwner, drainingRelaySessions, ensureRelayTurn, @@ -186,18 +187,20 @@ export function sendTalkRealtimeRelayAudio(params: { } const session = getRelaySession(params.relaySessionId, params.connId); const audio = decodeTalkRelayAudioBase64(params.audioBase64, "Realtime relay"); - const turnId = ensureRelayTurn(session); - session.bridge.sendAudio(audio); + const recorded = session.harness.recordInputAudio(audio); + if (!recorded) { + return; + } + broadcastRelayTurnStarted(session, recorded.turn.event); broadcastToOwner(session.context, session.connId, { relaySessionId: session.id, type: "inputAudio", byteLength: audio.byteLength, - talkEvent: session.harness.talk.emit({ - type: "input.audio.delta", - turnId, - payload: { byteLength: audio.byteLength }, - }), + talkEvent: recorded.inputAudioDelta, }); + // Publish the recorded input before provider code can synchronously re-enter + // with output events, preserving the harness's authoritative sequence order. + session.bridge.sendAudio(audio); if (typeof params.timestamp === "number" && Number.isFinite(params.timestamp)) { session.bridge.setMediaTimestamp(params.timestamp); } diff --git a/src/gateway/talk-realtime-relay-session-create.ts b/src/gateway/talk-realtime-relay-session-create.ts index d0b88b719421..add5256fb1a2 100644 --- a/src/gateway/talk-realtime-relay-session-create.ts +++ b/src/gateway/talk-realtime-relay-session-create.ts @@ -8,8 +8,8 @@ import { REALTIME_VOICE_AUDIO_FORMAT_PCM16_24KHZ, type RealtimeVoiceCloseReason, } from "../talk/provider-types.js"; -import { createRealtimeVoiceSessionHarness } from "../talk/realtime-session-harness.js"; -import type { TalkEventInput } from "../talk/talk-session-controller.js"; +import { createRealtimeVoiceEventCapturingSessionHarness } from "../talk/realtime-session-harness.js"; +import type { TalkEvent, TalkEventInput } from "../talk/talk-session-controller.js"; import { VOICE_TRANSCRIPT_QUEUE_POLICY } from "../talk/voice-transcript.js"; import { createTalkClientAgentConsultRunner, @@ -41,6 +41,7 @@ import { RELAY_SESSION_TTL_MS, RELAY_TRANSCRIPT_ECHO_LOOKBACK_MS, adoptRelayProviderToolCallId, + broadcastRelayTurnStarted, broadcastToOwner, ensureRelayTurn, relaySessions, @@ -78,7 +79,7 @@ export function createTalkRealtimeRelaySession( if (expiresAtMs === undefined) { throw new Error("Realtime relay session expiry is outside the supported Date range"); } - const harness = createRealtimeVoiceSessionHarness({ + const harness = createRealtimeVoiceEventCapturingSessionHarness({ talk: { sessionId: relaySessionId, mode: "realtime", @@ -94,16 +95,19 @@ export function createTalkRealtimeRelaySession( inputAudioDelta: (audio) => ({ byteLength: audio.byteLength }), outputAudioStarted: () => ({}), outputAudioDelta: (audio) => ({ byteLength: audio.byteLength }), - outputAudioDone: (reason) => ({ reason }), + outputAudioDone: (reason, details) => + details?.markName ? { markName: details.markName } : { reason }, }, transcriptLookbackMs: RELAY_TRANSCRIPT_ECHO_LOOKBACK_MS, captureBridgeEvents: false, }); - const emit = (event: TalkRealtimeRelayEventPayload, talkEvent?: TalkEventInput) => + const broadcastEvent = (event: TalkRealtimeRelayEventPayload, talkEvent?: TalkEvent) => broadcastToOwner(params.context, params.connId, { ...event, - ...(talkEvent ? { talkEvent: harness.emit(talkEvent) } : {}), + ...(talkEvent ? { talkEvent } : {}), }); + const emit = (event: TalkRealtimeRelayEventPayload, talkEvent?: TalkEventInput) => + broadcastEvent(event, talkEvent ? harness.emit(talkEvent) : undefined); let currentOutputItemId: string | undefined; let currentOutputResponseId: string | undefined; let ready = false; @@ -217,8 +221,12 @@ export function createTalkRealtimeRelaySession( if (!relay) { return; } - const turnId = ensureRelayTurn(relay); - emit( + const recorded = relay.harness.recordOutputAudio(audio); + broadcastRelayTurnStarted(relay, recorded.turn.event); + if (recorded.outputAudioStarted) { + broadcastEvent({ relaySessionId, type: "audioStarted" }, recorded.outputAudioStarted); + } + broadcastEvent( { relaySessionId, type: "audio", @@ -226,11 +234,7 @@ export function createTalkRealtimeRelaySession( ...(currentOutputItemId ? { itemId: currentOutputItemId } : {}), ...(currentOutputResponseId ? { responseId: currentOutputResponseId } : {}), }, - { - type: "output.audio.delta", - turnId, - payload: { byteLength: audio.length }, - }, + recorded.outputAudioDelta, ); }, clearAudio: (reason) => { @@ -238,15 +242,9 @@ export function createTalkRealtimeRelaySession( if (!relay) { return; } - const turnId = ensureRelayTurn(relay); - emit( + broadcastEvent( { relaySessionId, type: "clear", ...(reason ? { reason } : {}) }, - { - type: "output.audio.done", - turnId, - payload: { reason: reason ?? "clear" }, - final: true, - }, + relay.harness.finishOutputAudio(reason ?? "clear"), ); }, sendMark: (markName) => { @@ -254,15 +252,9 @@ export function createTalkRealtimeRelaySession( if (!relay) { return; } - const turnId = ensureRelayTurn(relay); - emit( + broadcastEvent( { relaySessionId, type: "mark", markName }, - { - type: "output.audio.done", - turnId, - payload: { markName }, - final: true, - }, + relay.harness.finishOutputAudio("mark", { markName }), ); }, }, diff --git a/src/gateway/talk-realtime-relay-state.ts b/src/gateway/talk-realtime-relay-state.ts index af274ebd900d..f4e30c47f48d 100644 --- a/src/gateway/talk-realtime-relay-state.ts +++ b/src/gateway/talk-realtime-relay-state.ts @@ -10,7 +10,7 @@ import type { RealtimeVoiceTool, RealtimeVoiceToolResultOptions, } from "../talk/provider-types.js"; -import type { RealtimeVoiceSessionHarness } from "../talk/realtime-session-harness.js"; +import type { RealtimeVoiceEventCapturingSessionHarness } from "../talk/realtime-session-harness.js"; import type { RealtimeVoiceBridgeSession } from "../talk/session-runtime.js"; import type { TalkEvent } from "../talk/talk-session-controller.js"; import type { GatewayRequestContext } from "./server-methods/shared-types.js"; @@ -28,6 +28,7 @@ export const noFallbackRelayOutputFlush = () => {}; export type TalkRealtimeRelayEventPayload = | { relaySessionId: string; type: "ready" } | { relaySessionId: string; type: "inputAudio"; byteLength: number } + | { relaySessionId: string; type: "audioStarted" } | { relaySessionId: string; type: "audio"; @@ -88,7 +89,7 @@ export type RelaySession = { connId: string; context: GatewayRequestContext; bridge: RealtimeVoiceBridgeSession; - harness: RealtimeVoiceSessionHarness; + harness: RealtimeVoiceEventCapturingSessionHarness; sessionKey?: string; agentId?: string; expiresAtMs: number; @@ -214,14 +215,22 @@ function relayEventDeliveryOptions( } export function ensureRelayTurn(session: RelaySession): string { - const turn = session.harness.talk.ensureTurn(); - if (turn.event) { - broadcastToOwner(session.context, session.connId, { - relaySessionId: session.id, - type: "inputAudio", - byteLength: 0, - talkEvent: turn.event, - }); - } + const turn = session.harness.ensureTurn(); + broadcastRelayTurnStarted(session, turn.event); return turn.turnId; } + +export function broadcastRelayTurnStarted( + session: RelaySession, + event: TalkEvent | undefined, +): void { + if (!event) { + return; + } + broadcastToOwner(session.context, session.connId, { + relaySessionId: session.id, + type: "inputAudio", + byteLength: 0, + talkEvent: event, + }); +} diff --git a/src/gateway/talk-realtime-relay.test.ts b/src/gateway/talk-realtime-relay.test.ts index 6483f11b95e9..c5fbe95edcdd 100644 --- a/src/gateway/talk-realtime-relay.test.ts +++ b/src/gateway/talk-realtime-relay.test.ts @@ -1526,6 +1526,11 @@ describe("talk realtime gateway relay", () => { language: "de", }); await Promise.resolve(); + const relay = relaySessions.get(session.relaySessionId); + expect(relay).toBeDefined(); + if (!relay) { + throw new Error("expected active relay session"); + } const sessionFields = expectRecordFields(session, { provider: "relay-test", @@ -1566,13 +1571,28 @@ describe("talk realtime gateway relay", () => { expectRecordFields(readyEvent, { event: "talk.event", connIds: ["conn-1"] }); expectDelivery(readyPayload, false); + const audioStartedPayload = findEventPayload( + events, + (payload) => payload.type === "audioStarted", + ); + const audioStartedEvent = expectRecordFields(audioStartedPayload.talkEvent, { + type: "output.audio.started", + }); + expect(relay.harness.talk.recentEvents).toContain(audioStartedPayload.talkEvent); + expectDelivery(audioStartedPayload, false); + const audioPayload = findEventPayload(events, (payload) => payload.type === "audio"); expectRecordFields(audioPayload, { relaySessionId: session.relaySessionId, type: "audio", audioBase64: Buffer.from("audio-out").toString("base64"), }); - expectRecordFields(audioPayload.talkEvent, { type: "output.audio.delta" }); + const audioEvent = expectRecordFields(audioPayload.talkEvent, { type: "output.audio.delta" }); + if (typeof audioStartedEvent.seq !== "number" || typeof audioEvent.seq !== "number") { + throw new Error("Expected sequenced relay audio events"); + } + expect(audioEvent.seq).toBe(audioStartedEvent.seq + 1); + expect(relay.harness.talk.recentEvents).toContain(audioPayload.talkEvent); expectDelivery(audioPayload, true); const markPayload = findEventPayload(events, (payload) => payload.type === "mark"); @@ -1581,6 +1601,12 @@ describe("talk realtime gateway relay", () => { type: "mark", markName: "mark-1", }); + expectRecordFields(markPayload.talkEvent, { + type: "output.audio.done", + payload: { markName: "mark-1" }, + final: true, + }); + expect(relay.harness.talk.recentEvents).toContain(markPayload.talkEvent); expectDelivery(markPayload, false); const partialTranscript = findEventPayload( @@ -1718,6 +1744,7 @@ describe("talk realtime gateway relay", () => { byteLength: Buffer.from("audio-in").byteLength, }); expectRecordFields(inputAudioPayload.talkEvent, { type: "input.audio.delta" }); + expect(relay.harness.talk.recentEvents).toContain(inputAudioPayload.talkEvent); expectDelivery(inputAudioPayload, true); const clearPayload = findEventPayload(events, (payload) => payload.type === "clear"); @@ -1786,6 +1813,54 @@ describe("talk realtime gateway relay", () => { expectDelivery(closePayload, false); }); + it("publishes recorded input before synchronous provider output", () => { + const events: Array<{ payload: Record }> = []; + const provider: RealtimeVoiceProviderPlugin = { + id: "relay-test", + label: "Relay Test", + isConfigured: () => true, + createBridge: (request) => ({ + connect: vi.fn(async () => undefined), + sendAudio: vi.fn(() => request.onAudio(Buffer.from("provider-output"))), + setMediaTimestamp: vi.fn(), + handleBargeIn: vi.fn(), + submitToolResult: vi.fn(), + acknowledgeMark: vi.fn(), + close: vi.fn(), + isConnected: vi.fn(() => true), + }), + }; + const session = createTalkRealtimeRelaySession({ + context: { + broadcastToConnIds: (_event: string, payload: Record) => { + events.push({ payload }); + }, + } as never, + connId: "conn-1", + provider, + providerConfig: {}, + instructions: "brief", + tools: [], + }); + + sendTalkRealtimeRelayAudio({ + relaySessionId: session.relaySessionId, + connId: "conn-1", + audioBase64: Buffer.from("browser-input").toString("base64"), + }); + + const sequencedEvents = events + .map(({ payload }) => payload.talkEvent) + .filter((event): event is Record => Boolean(event)); + expect(sequencedEvents.map((event) => event.type)).toEqual([ + "turn.started", + "input.audio.delta", + "output.audio.started", + "output.audio.delta", + ]); + expect(sequencedEvents.map((event) => event.seq)).toEqual([1, 2, 3, 4]); + }); + it("emits generic issue details when relay connect fails", async () => { const provider: RealtimeVoiceProviderPlugin = { id: "openai", diff --git a/src/plugin-sdk/realtime-voice.test.ts b/src/plugin-sdk/realtime-voice.test.ts index 60501b77a067..25cd2d4c5616 100644 --- a/src/plugin-sdk/realtime-voice.test.ts +++ b/src/plugin-sdk/realtime-voice.test.ts @@ -1,9 +1,12 @@ -import { describe, expect, it, vi } from "vitest"; +import { describe, expect, expectTypeOf, it, vi } from "vitest"; +import * as realtimeVoiceSdk from "./realtime-voice.js"; import { createRealtimeVoiceAudioQueue, + createRealtimeVoiceSessionHarness, normalizeRealtimeVoiceResponseOutcome, RealtimeVoiceSessionLifecycle, type RealtimeVoiceSessionConnection, + type RealtimeVoiceSessionHarness, } from "./realtime-voice.js"; function connectLifecycle( @@ -237,6 +240,42 @@ describe("RealtimeVoiceSessionLifecycle", () => { }); }); +describe("realtime voice session harness SDK contract", () => { + it("keeps legacy helper results for plugin consumers", () => { + const harness: RealtimeVoiceSessionHarness = createRealtimeVoiceSessionHarness({ + talk: { + sessionId: "plugin-session", + mode: "realtime", + transport: "gateway-relay", + brain: "agent-consult", + provider: "test", + }, + talkPayloads: { + turnStarted: () => ({}), + turnEnded: (reason) => ({ reason }), + inputAudioDelta: (audio) => ({ byteLength: audio.byteLength }), + outputAudioStarted: () => ({}), + outputAudioDelta: (audio) => ({ byteLength: audio.byteLength }), + outputAudioDone: (reason) => ({ reason }), + }, + }); + + expect(harness.ensureTurn()).toBe("turn-1"); + expect(harness.recordInputAudio(Buffer.from([1]))).toBe(true); + expect(harness.recordOutputAudio(Buffer.from([2]))).toBeUndefined(); + expect(harness.finishOutputAudio("completed")).toBeUndefined(); + expect(harness.endTurn("completed")).toBeUndefined(); + expectTypeOf().returns.toEqualTypeOf(); + expectTypeOf< + RealtimeVoiceSessionHarness["recordInputAudio"] + >().returns.toEqualTypeOf(); + expectTypeOf().returns.toEqualTypeOf(); + expectTypeOf().returns.toEqualTypeOf(); + expectTypeOf().returns.toEqualTypeOf(); + expect("createRealtimeVoiceEventCapturingSessionHarness" in realtimeVoiceSdk).toBe(false); + }); +}); + describe("normalizeRealtimeVoiceResponseOutcome", () => { it.each([ [ diff --git a/src/talk/realtime-session-harness.test.ts b/src/talk/realtime-session-harness.test.ts index 1c6e33043a97..79f867c9c467 100644 --- a/src/talk/realtime-session-harness.test.ts +++ b/src/talk/realtime-session-harness.test.ts @@ -2,31 +2,38 @@ import { afterEach, describe, expect, it, vi } from "vitest"; import type { RealtimeVoiceProviderPlugin } from "../plugins/types.js"; import type { RealtimeVoiceBridge } from "./provider-types.js"; -import { createRealtimeVoiceSessionHarness } from "./realtime-session-harness.js"; +import { + createRealtimeVoiceEventCapturingSessionHarness, + createRealtimeVoiceSessionHarness, +} from "./realtime-session-harness.js"; afterEach(() => { vi.useRealTimers(); }); +const defaultHarnessParams: Parameters[0] = { + talk: { + sessionId: "test-session", + mode: "realtime", + transport: "gateway-relay", + brain: "agent-consult", + provider: "test", + }, + talkPayloads: { + turnStarted: () => ({ surface: "test" }), + turnEnded: (reason) => ({ reason }), + inputAudioDelta: (audio) => ({ byteLength: audio.byteLength }), + outputAudioStarted: () => ({ surface: "test" }), + outputAudioDelta: (audio) => ({ byteLength: audio.byteLength }), + outputAudioDone: (reason) => ({ reason }), + }, +}; + function createHarness( overrides: Partial[0]> = {}, ) { return createRealtimeVoiceSessionHarness({ - talk: { - sessionId: "test-session", - mode: "realtime", - transport: "gateway-relay", - brain: "agent-consult", - provider: "test", - }, - talkPayloads: { - turnStarted: () => ({ surface: "test" }), - turnEnded: (reason) => ({ reason }), - inputAudioDelta: (audio) => ({ byteLength: audio.byteLength }), - outputAudioStarted: () => ({ surface: "test" }), - outputAudioDelta: (audio) => ({ byteLength: audio.byteLength }), - outputAudioDone: (reason) => ({ reason }), - }, + ...defaultHarnessParams, ...overrides, }); } @@ -137,13 +144,14 @@ describe("realtime voice session harness", () => { 1, ); }); - it("keeps shared Talk events ordered across input, output, and turn completion", () => { + it("keeps legacy helpers compatible while ordering Talk events", () => { const harness = createHarness(); + expect(harness.ensureTurn()).toBe("turn-1"); expect(harness.recordInputAudio(Buffer.from([1, 2]))).toBe(true); - harness.recordOutputAudio(Buffer.from([3, 4, 5])); - harness.finishOutputAudio("response.done"); - harness.endTurn("response.done"); + expect(harness.recordOutputAudio(Buffer.from([3, 4, 5]))).toBeUndefined(); + expect(harness.finishOutputAudio("response.done")).toBeUndefined(); + expect(harness.endTurn("response.done")).toBeUndefined(); expect(harness.talk.recentEvents.map((event) => event.type)).toEqual([ "turn.started", @@ -156,6 +164,41 @@ describe("realtime voice session harness", () => { expect(harness.talk.recentEvents.map((event) => event.seq)).toEqual([1, 2, 3, 4, 5, 6]); }); + it("returns the exact Talk events emitted by audio helpers", () => { + const harness = createRealtimeVoiceEventCapturingSessionHarness({ + ...defaultHarnessParams, + talkPayloads: { + ...defaultHarnessParams.talkPayloads, + outputAudioDone: (reason, details) => + details?.markName ? { markName: details.markName } : { reason }, + }, + }); + + const input = harness.recordInputAudio(Buffer.from([1, 2])); + const output = harness.recordOutputAudio(Buffer.from([3, 4, 5])); + const done = harness.finishOutputAudio("mark", { markName: "played" }); + + expect(input).toMatchObject({ + turn: { turnId: "turn-1", event: { type: "turn.started" } }, + inputAudioDelta: { type: "input.audio.delta", turnId: "turn-1" }, + }); + expect(output).toMatchObject({ + turn: { turnId: "turn-1" }, + outputAudioStarted: { type: "output.audio.started", turnId: "turn-1" }, + outputAudioDelta: { type: "output.audio.delta", turnId: "turn-1" }, + }); + expect(done).toMatchObject({ + type: "output.audio.done", + turnId: "turn-1", + payload: { markName: "played" }, + }); + expect(input?.turn.event).toBe(harness.talk.recentEvents[0]); + expect(input?.inputAudioDelta).toBe(harness.talk.recentEvents[1]); + expect(output.outputAudioStarted).toBe(harness.talk.recentEvents[2]); + expect(output.outputAudioDelta).toBe(harness.talk.recentEvents[3]); + expect(done).toBe(harness.talk.recentEvents[4]); + }); + it("honors a caller-specific recent Talk event limit", () => { const harness = createHarness({ talk: { diff --git a/src/talk/realtime-session-harness.ts b/src/talk/realtime-session-harness.ts index 742ab5d7c376..051951d0a68b 100644 --- a/src/talk/realtime-session-harness.ts +++ b/src/talk/realtime-session-harness.ts @@ -38,6 +38,7 @@ import { import type { TalkEvent, TalkEventInput } from "./talk-events.js"; import { createTalkSessionController, + type TalkEnsureTurnResult, type TalkSessionController, type TalkSessionControllerParams, type TalkTurnResult, @@ -74,6 +75,17 @@ type RealtimeVoiceSessionHarnessTalkPayloads = { outputAudioDone: (reason: string) => unknown; }; +type RealtimeVoiceOutputAudioDoneDetails = { + markName?: string; +}; + +type RealtimeVoiceEventCapturingTalkPayloads = Omit< + RealtimeVoiceSessionHarnessTalkPayloads, + "outputAudioDone" +> & { + outputAudioDone: (reason: string, details?: RealtimeVoiceOutputAudioDoneDetails) => unknown; +}; + type RealtimeVoiceSessionHarnessEchoSuppression = { bytesPerMs: number; tailMs: number; @@ -103,7 +115,18 @@ type RealtimeVoiceSessionHarnessHealth = ReturnType; }; -export type RealtimeVoiceSessionHarness = { +type RealtimeVoiceInputAudioEvents = { + inputAudioDelta: TalkEvent; + turn: TalkEnsureTurnResult; +}; + +type RealtimeVoiceOutputAudioEvents = { + outputAudioDelta: TalkEvent; + outputAudioStarted?: TalkEvent; + turn: TalkEnsureTurnResult; +}; + +type RealtimeVoiceSessionHarnessBase = { readonly forcedConsults: RealtimeVoiceForcedConsultCoordinator; readonly outputActivity: RealtimeVoiceOutputActivityTracker; readonly talk: TalkSessionController; @@ -112,10 +135,7 @@ export type RealtimeVoiceSessionHarness = { close(): void; createBridge(params: RealtimeVoiceBridgeSessionParams): RealtimeVoiceBridgeSession; emit(input: TalkEventInput): TalkEvent; - ensureTurn(): string; - endTurn(reason?: string): void; finishResponse(outcome: RealtimeVoiceResponseOutcome): TalkTurnResult; - finishOutputAudio(reason: string): void; flushOutput(flush: () => void): void; getHealth(params: { providerConnected: boolean; @@ -124,12 +144,54 @@ export type RealtimeVoiceSessionHarness = { handleBargeIn(options: RealtimeVoiceBargeInOptions, flushOutput: () => void): void; isLikelyAssistantEchoTranscript(text: string): boolean; isOutputPlaybackWindowActive(): boolean; - recordInputAudio(audio: Buffer): boolean; - recordOutputAudio(audio: Buffer, activity?: RealtimeVoiceOutputActivityDelta): void; recordTranscript(role: RealtimeVoiceRole, text: string): RealtimeVoiceTranscriptEntry; }; -export function createRealtimeVoiceSessionHarness(params: { +type RealtimeVoiceSessionHarnessMethods = { + ensureTurn(): string; + endTurn(reason?: string): void; + finishOutputAudio(reason: string): void; + recordInputAudio(audio: Buffer): boolean; + recordOutputAudio(audio: Buffer, activity?: RealtimeVoiceOutputActivityDelta): void; +}; + +type RealtimeVoiceEventCapturingSessionHarnessMethods = { + ensureTurn(): TalkEnsureTurnResult; + endTurn(reason?: string): TalkTurnResult; + finishOutputAudio( + reason: string, + details?: RealtimeVoiceOutputAudioDoneDetails, + ): TalkEvent | undefined; + recordInputAudio(audio: Buffer): RealtimeVoiceInputAudioEvents | undefined; + recordOutputAudio( + audio: Buffer, + activity?: RealtimeVoiceOutputActivityDelta, + ): RealtimeVoiceOutputAudioEvents; +}; + +export type RealtimeVoiceSessionHarness = + RealtimeVoiceSessionHarnessBase & RealtimeVoiceSessionHarnessMethods; + +export type RealtimeVoiceEventCapturingSessionHarness = + RealtimeVoiceSessionHarnessBase & + RealtimeVoiceEventCapturingSessionHarnessMethods; + +type RealtimeVoiceSessionHarnessImplementation = + RealtimeVoiceSessionHarnessBase & { + ensureTurn(): string | TalkEnsureTurnResult; + endTurn(reason?: string): TalkTurnResult | undefined; + finishOutputAudio( + reason: string, + details?: RealtimeVoiceOutputAudioDoneDetails, + ): TalkEvent | undefined; + recordInputAudio(audio: Buffer): boolean | RealtimeVoiceInputAudioEvents | undefined; + recordOutputAudio( + audio: Buffer, + activity?: RealtimeVoiceOutputActivityDelta, + ): RealtimeVoiceOutputAudioEvents | undefined; + }; + +type RealtimeVoiceSessionHarnessParams = { talk: TalkSessionControllerParams; talkPayloads: RealtimeVoiceSessionHarnessTalkPayloads; onTalkEvent?: (event: TalkEvent) => void; @@ -138,7 +200,38 @@ export function createRealtimeVoiceSessionHarness { +}; + +type RealtimeVoiceEventCapturingSessionHarnessParams = Omit< + RealtimeVoiceSessionHarnessParams, + "talkPayloads" +> & { + talkPayloads: RealtimeVoiceEventCapturingTalkPayloads; +}; + +export function createRealtimeVoiceSessionHarness( + params: RealtimeVoiceSessionHarnessParams, +): RealtimeVoiceSessionHarness { + return createRealtimeVoiceSessionHarnessImplementation( + params, + false, + ) as RealtimeVoiceSessionHarness; +} + +/** Gateway-internal factory that returns the exact Talk events emitted by each helper. */ +export function createRealtimeVoiceEventCapturingSessionHarness( + params: RealtimeVoiceEventCapturingSessionHarnessParams, +): RealtimeVoiceEventCapturingSessionHarness { + return createRealtimeVoiceSessionHarnessImplementation( + params, + true, + ) as RealtimeVoiceEventCapturingSessionHarness; +} + +function createRealtimeVoiceSessionHarnessImplementation( + params: RealtimeVoiceEventCapturingSessionHarnessParams, + captureEvents: boolean, +): RealtimeVoiceSessionHarnessImplementation { let closed = false; let bridge: RealtimeVoiceBridgeSession | undefined; let lastInputAt: string | undefined; @@ -178,10 +271,10 @@ export function createRealtimeVoiceSessionHarness { - const turnId = talk.ensureTurn({ payload: params.talkPayloads.turnStarted() }).turnId; - responseOwnerTurnId ??= turnId; - return turnId; + const ensureTurnWithEvents = () => { + const result = talk.ensureTurn({ payload: params.talkPayloads.turnStarted() }); + responseOwnerTurnId ??= result.turnId; + return result; }; const rememberSettledResponse = (responseId: string | undefined): void => { @@ -202,7 +295,7 @@ export function createRealtimeVoiceSessionHarness = { + const harness: RealtimeVoiceSessionHarnessImplementation = { forcedConsults, outputActivity, talk, @@ -331,19 +424,26 @@ export function createRealtimeVoiceSessionHarness talk.emit(input), - ensureTurn, + ensureTurn() { + const result = ensureTurnWithEvents(); + return captureEvents ? result : result.turnId; + }, endTurn(reason = "completed") { const result = talk.endTurn({ payload: params.talkPayloads.turnEnded(reason) }); if (result.ok) { responseOwnerTurnId = undefined; responseOwnerId = undefined; } + return captureEvents ? result : undefined; }, finishResponse(outcome) { return finishResponse(outcome, "typed"); }, - finishOutputAudio(reason) { - talk.finishOutputAudio({ payload: params.talkPayloads.outputAudioDone(reason) }); + finishOutputAudio(reason, details) { + const result = talk.finishOutputAudio({ + payload: params.talkPayloads.outputAudioDone(reason, details), + }); + return captureEvents ? result : undefined; }, flushOutput, getHealth(healthParams) { @@ -392,30 +492,31 @@ export function createRealtimeVoiceSessionHarness recordRealtimeVoiceTranscript(transcript, role, text), }; - harnessResponseOwners.set(harness, { claimResponseEvent, finishLegacyEvent }); + harnessResponseOwners.set(harness as RealtimeVoiceSessionHarness, { + claimResponseEvent, + finishLegacyEvent, + }); return harness; } diff --git a/ui/src/pages/chat/realtime-talk-gateway-relay-types.ts b/ui/src/pages/chat/realtime-talk-gateway-relay-types.ts index b44e428d52a9..e96e8dfccb27 100644 --- a/ui/src/pages/chat/realtime-talk-gateway-relay-types.ts +++ b/ui/src/pages/chat/realtime-talk-gateway-relay-types.ts @@ -5,6 +5,7 @@ export type GatewayRelayEvent = { talkEvent?: RealtimeTalkEvent; } & ( | { type?: "ready" } + | { type?: "audioStarted" } | { type?: "audio"; audioBase64?: string } | { type?: "clear"; reason?: "barge-in" } | { type?: "mark"; markName?: string } diff --git a/ui/src/pages/chat/realtime-talk-gateway-relay.ts b/ui/src/pages/chat/realtime-talk-gateway-relay.ts index 7aa6a3dee939..866253d926fa 100644 --- a/ui/src/pages/chat/realtime-talk-gateway-relay.ts +++ b/ui/src/pages/chat/realtime-talk-gateway-relay.ts @@ -316,7 +316,9 @@ export class GatewayRelayRealtimeTalkTransport implements RealtimeTalkTransport switch (event.type) { case "ready": this.ctx.callbacks.onStatus?.("listening"); - return; + break; + case "audioStarted": + break; case "audio": if (event.audioBase64 && !this.playbackOverflowed) { this.cancelRequestedForPlayback = false;