From 3c4f0bbe2bb36b5adbe2009ee2b7987eb58e51fe Mon Sep 17 00:00:00 2001 From: Vincent Koc Date: Sun, 2 Aug 2026 22:30:18 +0800 Subject: [PATCH] fix(talk): require successful transcript close barriers --- src/gateway/talk-realtime-relay-voice.test.ts | 23 +++++++ src/gateway/talk-realtime-relay-voice.ts | 2 +- src/shared/bounded-serial-queue.test.ts | 16 +++++ src/shared/bounded-serial-queue.ts | 19 ++++- src/talk/client-voice-session.test.ts | 69 +++++++++++++++++++ src/talk/voice-transcript.ts | 2 +- 6 files changed, 127 insertions(+), 4 deletions(-) diff --git a/src/gateway/talk-realtime-relay-voice.test.ts b/src/gateway/talk-realtime-relay-voice.test.ts index 8a7dfbfc934c..1fbbb107a29b 100644 --- a/src/gateway/talk-realtime-relay-voice.test.ts +++ b/src/gateway/talk-realtime-relay-voice.test.ts @@ -108,6 +108,29 @@ describe("realtime relay voice transcript persistence", () => { expect(enqueueRelayVoiceTranscript(session, "user", "too late")).toBe(false); }); + it("does not close the durable record after an accepted transcript fails", async () => { + vi.useFakeTimers(); + try { + voiceSessionMocks.appendRelayVoiceTranscript.mockRejectedValue( + new Error("transcript write failed"), + ); + const { session } = createRelaySession(); + + expect(enqueueRelayVoiceTranscript(session, "user", "persist me")).toBe(true); + const close = closeRelayVoiceSession(session); + await vi.runAllTimersAsync(); + await close; + + expect(voiceSessionMocks.appendRelayVoiceTranscript).toHaveBeenCalledTimes(3); + expect(voiceSessionMocks.closeClientVoiceSession).not.toHaveBeenCalled(); + expect(session.context.logGateway?.warn).toHaveBeenCalledWith( + expect.stringContaining("realtime relay voice session close failed"), + ); + } finally { + vi.useRealTimers(); + } + }); + it("normalizes the bounded pre-bind transcript buffer", () => { const { session } = createRelaySession(); session.sessionKey = undefined; diff --git a/src/gateway/talk-realtime-relay-voice.ts b/src/gateway/talk-realtime-relay-voice.ts index 5b6869503b45..8e936497984a 100644 --- a/src/gateway/talk-realtime-relay-voice.ts +++ b/src/gateway/talk-realtime-relay-voice.ts @@ -144,7 +144,7 @@ export function closeRelayVoiceSession(session: RelaySession): Promise { } const sessionKey = session.sessionKey; session.voiceSessionClose = session.voiceTranscriptQueue - .flush() + .flush({ requireSuccess: true }) .then(async () => { const config = session.voiceConfig ?? session.context.getRuntimeConfig(); await closeClientVoiceSession({ diff --git a/src/shared/bounded-serial-queue.test.ts b/src/shared/bounded-serial-queue.test.ts index b3d5081faba5..2b75b532fed9 100644 --- a/src/shared/bounded-serial-queue.test.ts +++ b/src/shared/bounded-serial-queue.test.ts @@ -126,6 +126,22 @@ describe("BoundedSerialQueue", () => { } }); + it("can require every task in the accepted prefix to succeed", async () => { + const failure = new Error("persistence failed"); + const queue = new BoundedSerialQueue({ maxPendingCount: 1, maxPendingWeight: 1 }); + const task = queue.enqueue(async () => { + throw failure; + }); + const ordinaryFlush = queue.flush(); + const strictFlush = queue.flush({ requireSuccess: true }); + + await expect(ordinaryFlush).resolves.toBeUndefined(); + await expect(strictFlush).rejects.toBe(failure); + if (task.accepted) { + await expect(task.completion).rejects.toBe(failure); + } + }); + it("seals idempotently while preserving accepted work", async () => { const first = deferred(); const queue = new BoundedSerialQueue({ maxPendingCount: 1, maxPendingWeight: 1 }); diff --git a/src/shared/bounded-serial-queue.ts b/src/shared/bounded-serial-queue.ts index af696062899d..e8f942596f6e 100644 --- a/src/shared/bounded-serial-queue.ts +++ b/src/shared/bounded-serial-queue.ts @@ -21,6 +21,8 @@ export class BoundedSerialQueue { private active = false; private sealed = false; private overflowed = false; + private failed = false; + private firstFailure: unknown; private settledPrefix: Promise = Promise.resolve(); constructor( @@ -104,9 +106,18 @@ export class BoundedSerialQueue { * * Later admissions do not extend this barrier, which keeps consult flushes * finite while close can seal first to drain the entire accepted prefix. + * Close owners can require that prefix to have completed without failures. */ - flush(): Promise { - return this.settledPrefix; + flush(options: { requireSuccess?: boolean } = {}): Promise { + const prefix = this.settledPrefix; + if (options.requireSuccess !== true) { + return prefix; + } + return prefix.then(() => { + if (this.failed) { + throw this.firstFailure; + } + }); } private startTask(task: BoundedSerialQueueTask): void { @@ -117,6 +128,10 @@ export class BoundedSerialQueue { try { task.resolve(await task.run()); } catch (error) { + if (!this.failed) { + this.failed = true; + this.firstFailure = error; + } task.reject(error); } finally { const next = this.pending.shift(); diff --git a/src/talk/client-voice-session.test.ts b/src/talk/client-voice-session.test.ts index dfe4f10cde84..fbd9d58eb2e2 100644 --- a/src/talk/client-voice-session.test.ts +++ b/src/talk/client-voice-session.test.ts @@ -403,6 +403,75 @@ describe("client voice session", () => { ); }); + it("keeps the session open when an accepted transcript fails during close", async () => { + await seedSession("agent:main:main"); + const voiceSessionId = createOrResumeClientVoiceSession({ + agentId: "main", + sessionKey: "agent:main:main", + origin: "client", + voiceSessionId: "voice-close-after-failure", + }); + const transcriptWrite = createDeferred(); + const failure = new Error("transcript write failed"); + const actualAppend = sessionAccessorMocks.actualAppendTranscriptMessage!; + sessionAccessorMocks.appendTranscriptMessage.mockImplementationOnce(async () => { + await transcriptWrite.promise; + throw failure; + }); + const append = appendClientVoiceTranscript({ + agentId: "main", + sessionKey: "agent:main:main", + voiceSessionId, + entryId: "retryable", + role: "user", + text: "persist me", + }); + const appendResult = append.then( + () => undefined, + (error: unknown) => error, + ); + await vi.waitFor(() => + expect(sessionAccessorMocks.appendTranscriptMessage).toHaveBeenCalledOnce(), + ); + const close = closeClientVoiceSession({ + agentId: "main", + sessionKey: "agent:main:main", + voiceSessionId, + config: {}, + now: 42, + }); + const closeResult = close.then( + () => undefined, + (error: unknown) => error, + ); + + transcriptWrite.resolve(); + expect(await appendResult).toBe(failure); + expect(await closeResult).toBe(failure); + expect(clientVoiceSessionTesting.readRecord("main", voiceSessionId)?.status).toBe("open"); + + sessionAccessorMocks.appendTranscriptMessage.mockImplementation(actualAppend); + await appendClientVoiceTranscript({ + agentId: "main", + sessionKey: "agent:main:main", + voiceSessionId, + entryId: "retryable", + role: "user", + text: "persist me", + }); + await closeClientVoiceSession({ + agentId: "main", + sessionKey: "agent:main:main", + voiceSessionId, + config: {}, + now: 99, + }); + expect(clientVoiceSessionTesting.readRecord("main", voiceSessionId)).toMatchObject({ + status: "closed", + closedAt: 99, + }); + }); + it("bounds stalled transcript operations and closes after the accepted prefix", async () => { await seedSession("agent:main:main"); const voiceSessionId = createOrResumeClientVoiceSession({ diff --git a/src/talk/voice-transcript.ts b/src/talk/voice-transcript.ts index ac9875f11d8e..f5864e5a16bb 100644 --- a/src/talk/voice-transcript.ts +++ b/src/talk/voice-transcript.ts @@ -106,7 +106,7 @@ class VoiceTranscriptOperationRegistry { if (!owner.closePromise) { // Seal synchronously so no transcript can enter behind the close barrier. owner.queue.seal(); - owner.closePromise = owner.queue.flush().then(operation); + owner.closePromise = owner.queue.flush({ requireSuccess: true }).then(operation); } try { await owner.closePromise;