diff --git a/ui/src/pages/chat/realtime-talk-lifecycle.test.ts b/ui/src/pages/chat/realtime-talk-lifecycle.test.ts index ddaf3bdbda0c..6123a5fef255 100644 --- a/ui/src/pages/chat/realtime-talk-lifecycle.test.ts +++ b/ui/src/pages/chat/realtime-talk-lifecycle.test.ts @@ -308,6 +308,56 @@ describe("RealtimeTalkSession lifecycle", () => { session.stop(); }); + it("stops a pending replacement when transcript overflow closes the active call", async () => { + const warn = vi.spyOn(console, "warn").mockImplementation(() => undefined); + const replacementStart = createDeferred(); + const firstTranscript = createDeferred(); + const request = vi.fn(async (method: string, params?: { entryId?: string }) => { + if (method === "talk.client.create") { + return { + provider: "openai", + transport: "webrtc", + voiceSessionId: "voice-overflow-replacement", + clientSecret: "secret", + }; + } + if (method === "talk.client.transcript" && params?.entryId === "1") { + await firstTranscript.promise; + } + return { ok: true }; + }); + const onStatus = vi.fn(); + const session = new RealtimeTalkSession({ request } as never, "agent:main:main", { + onStatus, + }); + await session.start(); + const activeContext = transcriptContext(transportMock.webRtcContexts); + transportMock.start.mockImplementationOnce(async () => await replacementStart.promise); + + const replacement = session.start(); + await vi.waitFor(() => expect(transportMock.webRtcContexts).toHaveLength(2)); + for (let index = 0; index < 42; index += 1) { + activeContext.callbacks.onTranscript?.({ + role: "user", + text: "x".repeat(9_000), + final: true, + }); + } + + expect(onStatus).toHaveBeenCalledWith("error", expect.stringContaining("could not keep up")); + expect(warn).toHaveBeenCalledOnce(); + expect(transportMock.webRtcStops[0]).toHaveBeenCalledWith(); + expect(transportMock.webRtcStops[1]).toHaveBeenCalledWith({ emitClosed: false }); + + replacementStart.resolve(); + await replacement; + firstTranscript.resolve(); + await vi.waitFor(() => + expect(request.mock.calls.some(([method]) => method === "talk.client.close")).toBe(true), + ); + warn.mockRestore(); + }); + it("releases newly allocated owners after transport startup failures", async () => { let createCount = 0; const request = vi.fn(async (method: string) => { diff --git a/ui/src/pages/chat/realtime-talk.ts b/ui/src/pages/chat/realtime-talk.ts index 4651789b1217..479d11c2eae6 100644 --- a/ui/src/pages/chat/realtime-talk.ts +++ b/ui/src/pages/chat/realtime-talk.ts @@ -168,9 +168,7 @@ export class RealtimeTalkSession { let ownerTransferred = false; try { const lifecycleGeneration = ++this.lifecycleGeneration; - const supersededPendingTransport = this.pendingTransport; - this.pendingTransport = null; - supersededPendingTransport?.stop({ emitClosed: false }); + this.stopPendingTransport(); this.closed = false; this.callbacks.onStatus?.("connecting"); const existingTransport = this.transport; @@ -393,9 +391,7 @@ export class RealtimeTalkSession { this.videoEnabled = false; activeRealtimeTalkSessions.delete(this); this.callbacks.onStatus?.("idle"); - const pendingTransport = this.pendingTransport; - this.pendingTransport = null; - pendingTransport?.stop({ emitClosed: false }); + this.stopPendingTransport(); const detached = this.detachVoiceSession(); this.transport?.stop(); this.transport = null; @@ -404,6 +400,12 @@ export class RealtimeTalkSession { } } + private stopPendingTransport(): void { + const pendingTransport = this.pendingTransport; + this.pendingTransport = null; + pendingTransport?.stop({ emitClosed: false }); + } + private closeUnadoptedVoiceSession( voiceSessionId: string, transport: string, @@ -511,6 +513,7 @@ export class RealtimeTalkSession { this.videoOperation += 1; this.videoEnabled = false; activeRealtimeTalkSessions.delete(this); + this.stopPendingTransport(); const detached = this.detachVoiceSession(); // Retire the overflowing transport before accepted-write and close failures // settle so the first terminal persistence error keeps precedence.