From b786f4a804f8ae050def290c8c439934b8eb718f Mon Sep 17 00:00:00 2001 From: Vincent Koc Date: Sat, 1 Aug 2026 13:21:05 +0800 Subject: [PATCH] fix(talk): expose one transcript queue policy --- src/talk/client-voice-session.ts | 5 ++++- src/talk/voice-transcript.test.ts | 6 +++--- src/talk/voice-transcript.ts | 31 ++++++++++++++++++++----------- 3 files changed, 27 insertions(+), 15 deletions(-) diff --git a/src/talk/client-voice-session.ts b/src/talk/client-voice-session.ts index ccda4eac7204..53835279c589 100644 --- a/src/talk/client-voice-session.ts +++ b/src/talk/client-voice-session.ts @@ -41,10 +41,13 @@ import { import { createVoiceTranscriptOperationRegistry, normalizeVoiceTranscriptText, + VOICE_TRANSCRIPT_QUEUE_POLICY, } from "./voice-transcript.js"; const voiceSessionByRunId = new Map(); -const voiceSessionOperations = createVoiceTranscriptOperationRegistry(); +const voiceSessionOperations = createVoiceTranscriptOperationRegistry( + VOICE_TRANSCRIPT_QUEUE_POLICY, +); const deferredDigestConfigs = new Map(); let unsubscribeToolEffects: (() => void) | undefined; let unsubscribeRunCompletion: (() => void) | undefined; diff --git a/src/talk/voice-transcript.test.ts b/src/talk/voice-transcript.test.ts index 4fddafe7e638..fab6a48ba797 100644 --- a/src/talk/voice-transcript.test.ts +++ b/src/talk/voice-transcript.test.ts @@ -1,7 +1,7 @@ import { describe, expect, it, vi } from "vitest"; import { createVoiceTranscriptOperationRegistry, - VOICE_TRANSCRIPT_QUEUE_MAX_PENDING, + VOICE_TRANSCRIPT_QUEUE_POLICY, } from "./voice-transcript.js"; function deferred(): { promise: Promise; resolve: () => void } { @@ -14,12 +14,12 @@ function deferred(): { promise: Promise; resolve: () => void } { describe("VoiceTranscriptOperationRegistry", () => { it("keeps overflow terminal through drain and releases it only on close", async () => { - const registry = createVoiceTranscriptOperationRegistry(); + const registry = createVoiceTranscriptOperationRegistry(VOICE_TRANSCRIPT_QUEUE_POLICY); const first = deferred(); const key = "agent\0voice-overflow"; const accepted = [ registry.run(key, async () => await first.promise), - ...Array.from({ length: VOICE_TRANSCRIPT_QUEUE_MAX_PENDING }, () => + ...Array.from({ length: VOICE_TRANSCRIPT_QUEUE_POLICY.maxPendingCount }, () => registry.run(key, async () => undefined), ), ]; diff --git a/src/talk/voice-transcript.ts b/src/talk/voice-transcript.ts index 8b2355ddbb8b..ac9875f11d8e 100644 --- a/src/talk/voice-transcript.ts +++ b/src/talk/voice-transcript.ts @@ -2,22 +2,25 @@ import { truncateUtf16Safe } from "@openclaw/normalization-core/utf16-slice"; import { BoundedSerialQueue } from "../shared/bounded-serial-queue.js"; const VOICE_TRANSCRIPT_MAX_CHARS = 8_000; -export const VOICE_TRANSCRIPT_QUEUE_MAX_PENDING = 40; +const VOICE_TRANSCRIPT_QUEUE_MAX_PENDING = 40; const VOICE_TRANSCRIPT_QUEUE_MAX_PENDING_CHARS = VOICE_TRANSCRIPT_QUEUE_MAX_PENDING * VOICE_TRANSCRIPT_MAX_CHARS; -export const VOICE_TRANSCRIPT_QUEUE_OVERFLOW_MESSAGE = +const VOICE_TRANSCRIPT_QUEUE_OVERFLOW_MESSAGE = "Voice transcript persistence could not keep up; the realtime session was stopped."; export function normalizeVoiceTranscriptText(text: string): string { return truncateUtf16Safe(text.trim(), VOICE_TRANSCRIPT_MAX_CHARS); } -export function createVoiceTranscriptQueue(): BoundedSerialQueue { - return new BoundedSerialQueue({ - maxPendingCount: VOICE_TRANSCRIPT_QUEUE_MAX_PENDING, - maxPendingWeight: VOICE_TRANSCRIPT_QUEUE_MAX_PENDING_CHARS, - }); -} +export const VOICE_TRANSCRIPT_QUEUE_POLICY = { + maxPendingCount: VOICE_TRANSCRIPT_QUEUE_MAX_PENDING, + overflowMessage: VOICE_TRANSCRIPT_QUEUE_OVERFLOW_MESSAGE, + createQueue: () => + new BoundedSerialQueue({ + maxPendingCount: VOICE_TRANSCRIPT_QUEUE_MAX_PENDING, + maxPendingWeight: VOICE_TRANSCRIPT_QUEUE_MAX_PENDING_CHARS, + }), +} as const; type VoiceTranscriptOperationOwner = { queue: BoundedSerialQueue; @@ -27,12 +30,16 @@ type VoiceTranscriptOperationOwner = { class VoiceTranscriptOperationRegistry { private readonly owners = new Map(); + constructor( + private readonly queuePolicy: Pick, + ) {} + private getOrCreate(key: string): VoiceTranscriptOperationOwner { const existing = this.owners.get(key); if (existing) { return existing; } - const created = { queue: createVoiceTranscriptQueue() }; + const created = { queue: this.queuePolicy.createQueue() }; this.owners.set(key, created); return created; } @@ -115,6 +122,8 @@ class VoiceTranscriptOperationRegistry { } } -export function createVoiceTranscriptOperationRegistry(): VoiceTranscriptOperationRegistry { - return new VoiceTranscriptOperationRegistry(); +export function createVoiceTranscriptOperationRegistry( + queuePolicy: Pick, +): VoiceTranscriptOperationRegistry { + return new VoiceTranscriptOperationRegistry(queuePolicy); }