Merge pull request #117682 from openclaw/fix/talk-late-audio-admission-20260801

* commit '76bafba7ee44c1a3c8f772f7f37c106312ab8666':
  fix(talk): preserve provisional terminal state
  test(talk): keep provisional close monotonic
  fix(talk): keep reconnect ownership provider-local
  test(talk): reject post-terminal facade reconnect
  fix(talk): separate session and connection closure
  test(talk): cover reconnect lifecycle ownership
  fix(talk): stop late audio after session close
  test(talk): cover late audio after session close
This commit is contained in:
Vincent Koc
2026-08-02 07:29:11 +08:00
2 changed files with 208 additions and 9 deletions
+161
View File
@@ -343,6 +343,167 @@ describe("realtime voice bridge session runtime", () => {
expect(onError).not.toHaveBeenCalled();
});
it("permanently closes once while preserving synchronous transcript flush", async () => {
let callbacks: Parameters<RealtimeVoiceProviderPlugin["createBridge"]>[0] | undefined;
const close = vi.fn(() => {
callbacks?.onTranscript?.("assistant", "final transcript", true);
});
const connect = vi.fn(async () => {});
const sendProviderAudio = vi.fn();
const sendSinkAudio = vi.fn();
const onTranscript = vi.fn();
const provider: RealtimeVoiceProviderPlugin = {
id: "test",
label: "Test",
isConfigured: () => true,
createBridge: (request) => {
callbacks = request;
return makeBridge({ close, connect, sendAudio: sendProviderAudio });
},
};
const session = createRealtimeVoiceBridgeSession({
provider,
providerConfig: {},
audioSink: { sendAudio: sendSinkAudio },
onTranscript,
});
session.close();
session.close();
session.sendAudio(Buffer.from("late-input"));
callbacks?.onAudio(Buffer.from("late-output"));
await expect(session.connect()).rejects.toThrow("Realtime voice session is closed");
expect(close).toHaveBeenCalledTimes(1);
expect(connect).not.toHaveBeenCalled();
expect(sendProviderAudio).not.toHaveBeenCalled();
expect(sendSinkAudio).not.toHaveBeenCalled();
expect(onTranscript).toHaveBeenCalledExactlyOnceWith("assistant", "final transcript", true);
});
it("stops audio admission after provider close and still closes the provider once", () => {
let callbacks: Parameters<RealtimeVoiceProviderPlugin["createBridge"]>[0] | undefined;
const close = vi.fn();
const sendProviderAudio = vi.fn();
const sendSinkAudio = vi.fn();
const onClose = vi.fn();
const provider: RealtimeVoiceProviderPlugin = {
id: "test",
label: "Test",
isConfigured: () => true,
createBridge: (request) => {
callbacks = request;
return makeBridge({ close, sendAudio: sendProviderAudio });
},
};
const session = createRealtimeVoiceBridgeSession({
provider,
providerConfig: {},
audioSink: { sendAudio: sendSinkAudio },
onClose,
});
callbacks?.onClose?.("completed");
callbacks?.onClose?.("completed");
session.sendAudio(Buffer.from("late-input"));
callbacks?.onAudio(Buffer.from("late-output"));
session.close();
session.close();
expect(onClose).toHaveBeenCalledTimes(1);
expect(close).toHaveBeenCalledTimes(1);
expect(sendProviderAudio).not.toHaveBeenCalled();
expect(sendSinkAudio).not.toHaveBeenCalled();
});
it("reopens audio and close reporting for an explicit connection generation", async () => {
let callbacks: Parameters<RealtimeVoiceProviderPlugin["createBridge"]>[0] | undefined;
const connect = vi.fn(async () => {});
const sendProviderAudio = vi.fn();
const sendSinkAudio = vi.fn();
const onClose = vi.fn();
const onReady = vi.fn();
const provider: RealtimeVoiceProviderPlugin = {
id: "test",
label: "Test",
isConfigured: () => true,
createBridge: (request) => {
callbacks = request;
request.onClose?.("error");
return makeBridge({ connect, sendAudio: sendProviderAudio });
},
};
const session = createRealtimeVoiceBridgeSession({
provider,
providerConfig: {},
audioSink: { sendAudio: sendSinkAudio },
onClose,
onReady,
});
session.sendAudio(Buffer.from("closed-input"));
callbacks?.onAudio(Buffer.from("closed-output"));
callbacks?.onClose?.("error");
callbacks?.onReady?.();
expect(onReady).not.toHaveBeenCalled();
await session.connect();
callbacks?.onReady?.();
session.sendAudio(Buffer.from("next-input"));
callbacks?.onAudio(Buffer.from("next-output"));
callbacks?.onClose?.("completed");
callbacks?.onClose?.("completed");
await expect(session.connect()).rejects.toThrow("Realtime voice connection is closed");
expect(connect).toHaveBeenCalledTimes(1);
expect(onReady).toHaveBeenCalledWith(session);
expect(onClose).toHaveBeenNthCalledWith(1, "error");
expect(onClose).toHaveBeenNthCalledWith(2, "completed");
expect(onClose).toHaveBeenCalledTimes(2);
expect(sendProviderAudio).toHaveBeenCalledExactlyOnceWith(Buffer.from("next-input"));
expect(sendSinkAudio).toHaveBeenCalledExactlyOnceWith(Buffer.from("next-output"));
});
it("rejects reconnect and ignores tool failures after an established provider close", async () => {
let callbacks: Parameters<RealtimeVoiceProviderPlugin["createBridge"]>[0] | undefined;
let rejectToolCall: ((error: Error) => void) | undefined;
const connect = vi.fn(async () => {});
const onError = vi.fn();
const provider: RealtimeVoiceProviderPlugin = {
id: "test",
label: "Test",
isConfigured: () => true,
createBridge: (request) => {
callbacks = request;
return makeBridge({ connect });
},
};
const session = createRealtimeVoiceBridgeSession({
provider,
providerConfig: {},
audioSink: { sendAudio: vi.fn() },
onToolCall: () =>
new Promise<void>((_resolve, reject) => {
rejectToolCall = reject;
}),
onError,
});
callbacks?.onToolCall?.({
itemId: "item-1",
callId: "call-1",
name: "lookup",
args: {},
});
callbacks?.onClose?.("error");
await expect(session.connect()).rejects.toThrow("Realtime voice connection is closed");
rejectToolCall?.(new Error("late tool callback failure"));
await Promise.resolve();
expect(connect).not.toHaveBeenCalled();
expect(onError).not.toHaveBeenCalled();
});
it("forwards tool result continuation options and async acceptance to the provider bridge", () => {
const acceptance = Promise.resolve();
const submitToolResult = vi.fn(() => acceptance);
+47 -9
View File
@@ -78,6 +78,8 @@ export type RealtimeVoiceBridgeSessionParams = {
onClose?: (reason: RealtimeVoiceCloseReason) => void;
};
type RealtimeVoiceSessionPhase = "admitting" | "provider-terminal" | "disposed";
/**
* Creates a realtime voice bridge session and wires provider events to the configured audio sink.
*/
@@ -85,7 +87,12 @@ export function createRealtimeVoiceBridgeSession(
params: RealtimeVoiceBridgeSessionParams,
): RealtimeVoiceBridgeSession {
const bridgeRef: { current?: RealtimeVoiceBridge } = {};
let isActive = true;
// Local disposal owns provider cleanup. Only a terminal callback fired before bridge
// adoption may reopen; adopted bridges own reconnects and stale-event fencing internally.
let phase: RealtimeVoiceSessionPhase = "admitting";
let terminalBeforeBridgeAdoption = false;
let closeReported = false;
const isAdmitting = () => phase === "admitting";
const requireBridge = () => {
if (!bridgeRef.current) {
throw new Error("Realtime voice bridge is not ready");
@@ -100,12 +107,32 @@ export function createRealtimeVoiceBridgeSession(
},
acknowledgeMark: (markName) => requireBridge().acknowledgeMark(markName),
close: () => {
if (phase === "disposed") {
return;
}
const bridge = requireBridge();
isActive = false;
phase = "disposed";
bridge.close();
},
connect: () => requireBridge().connect(),
sendAudio: (audio) => requireBridge().sendAudio(audio),
connect: () => {
if (phase === "disposed") {
return Promise.reject(new Error("Realtime voice session is closed"));
}
if (phase === "provider-terminal") {
if (!terminalBeforeBridgeAdoption) {
return Promise.reject(new Error("Realtime voice connection is closed"));
}
terminalBeforeBridgeAdoption = false;
phase = "admitting";
closeReported = false;
}
return requireBridge().connect();
},
sendAudio: (audio) => {
if (isAdmitting()) {
requireBridge().sendAudio(audio);
}
},
sendUserMessage: (text) => requireBridge().sendUserMessage?.(text),
handleBargeIn: (options) => requireBridge().handleBargeIn?.(options),
setMediaTimestamp: (ts) => requireBridge().setMediaTimestamp(ts),
@@ -118,11 +145,13 @@ export function createRealtimeVoiceBridgeSession(
},
triggerGreeting: (instructions) => requireBridge().triggerGreeting?.(instructions),
};
const canSendAudio = () => params.audioSink.isOpen?.() ?? true;
// Session inactivity is the shared admission boundary for both audio directions.
// Provider and transport callbacks may still race after close, but cannot retain new audio.
const canSendAudio = () => isAdmitting() && (params.audioSink.isOpen?.() ?? true);
const reportCallbackError = (error: unknown) => {
// Async tool handlers can settle after the provider closes. Once inactive, no
// callback may report stale failures into the next session lifecycle.
if (!isActive) {
if (!isAdmitting()) {
return;
}
try {
@@ -167,7 +196,7 @@ export function createRealtimeVoiceBridgeSession(
onTranscript: params.onTranscript,
onEvent: params.onEvent,
onToolCall: (event) => {
if (!bridgeRef.current || !isActive) {
if (!bridgeRef.current || !isAdmitting()) {
return;
}
try {
@@ -180,7 +209,7 @@ export function createRealtimeVoiceBridgeSession(
}
},
onReady: () => {
if (!bridgeRef.current) {
if (!bridgeRef.current || !isAdmitting()) {
return;
}
if (params.triggerGreetingOnReady) {
@@ -190,7 +219,16 @@ export function createRealtimeVoiceBridgeSession(
},
onError: params.onError,
onClose: (reason) => {
isActive = false;
if (!bridgeRef.current) {
terminalBeforeBridgeAdoption = true;
}
if (phase !== "disposed") {
phase = "provider-terminal";
}
if (closeReported) {
return;
}
closeReported = true;
params.onClose?.(reason);
},
});