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 <steipete@gmail.com>
This commit is contained in:
Dallin Romney
2026-07-23 06:40:32 +09:00
committed by GitHub
parent e194979830
commit 635d396755
3 changed files with 113 additions and 85 deletions
+69 -78
View File
@@ -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<typeof setTimeout>;
@@ -160,8 +147,6 @@ type RelaySession = {
forcedTerminalProviderResults: Map<string, ForcedTerminalProviderResult>;
// 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<void> {
const pending = session.forcedConsults
const pending = session.harness.forcedConsults
.nativeCallIds(handle)
.map((callId) => session.pendingProviderToolResults.get(callId))
.filter((submission): submission is Promise<void> => 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 },
});
+31
View File
@@ -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 = {
+13 -7
View File
@@ -107,6 +107,8 @@ export function createRealtimeVoiceSessionHarness<TForcedConsultContext = unknow
talkback?: Omit<RealtimeVoiceAgentTalkbackQueueParams, "isStopped">;
forcedConsults?: RealtimeVoiceForcedConsultCoordinatorOptions;
echoSuppression?: RealtimeVoiceSessionHarnessEchoSuppression;
transcriptLookbackMs?: number;
captureBridgeEvents?: boolean;
}): RealtimeVoiceSessionHarness<TForcedConsultContext> {
let closed = false;
let bridge: RealtimeVoiceBridgeSession | undefined;
@@ -121,11 +123,13 @@ export function createRealtimeVoiceSessionHarness<TForcedConsultContext = unknow
const transcript: RealtimeVoiceTranscriptEntry[] = [];
const bridgeEvents: RealtimeVoiceBridgeEventLogEntry[] = [];
const outputActivity = createRealtimeVoiceOutputActivityTracker();
const transcriptLookbackMs =
params.transcriptLookbackMs ?? params.echoSuppression?.transcriptLookbackMs;
const forcedConsults = createRealtimeVoiceForcedConsultCoordinator<TForcedConsultContext>(
params.forcedConsults,
);
const talk = createTalkSessionController(
{ ...params.talk, maxRecentEvents: 40 },
{ maxRecentEvents: 40, ...params.talk },
{
onEvent: (event) => {
recordTalkObservabilityEvent(event);
@@ -173,7 +177,9 @@ export function createRealtimeVoiceSessionHarness<TForcedConsultContext = unknow
bridgeParams.onTranscript?.(role, text, isFinal);
},
onEvent: (event) => {
recordRealtimeVoiceBridgeEvent(bridgeEvents, event);
if (params.captureBridgeEvents !== false) {
recordRealtimeVoiceBridgeEvent(bridgeEvents, event);
}
bridgeParams.onEvent?.(event);
},
});
@@ -223,13 +229,13 @@ export function createRealtimeVoiceSessionHarness<TForcedConsultContext = unknow
}
},
isLikelyAssistantEchoTranscript(text) {
return params.echoSuppression
? isLikelyRealtimeVoiceAssistantEchoTranscript({
return transcriptLookbackMs === undefined
? false
: isLikelyRealtimeVoiceAssistantEchoTranscript({
transcript,
text,
lookbackMs: params.echoSuppression.transcriptLookbackMs,
})
: false;
lookbackMs: transcriptLookbackMs,
});
},
isOutputPlaybackWindowActive() {
return Date.now() <= Math.max(lastOutputPlayableUntilMs, suppressInputUntilMs);