diff --git a/packages/coding-agent/src/cli.ts b/packages/coding-agent/src/cli.ts index 22e2874cc..5eb3155ec 100755 --- a/packages/coding-agent/src/cli.ts +++ b/packages/coding-agent/src/cli.ts @@ -173,14 +173,20 @@ async function runWorkerEntrypoint(arg: string | undefined): Promise { * worker is idle, and hard-kills the process on parent `disconnect`. */ async function runIpcSubprocessWorker( - start: (transport: { send(message: Out): void; onMessage(handler: (message: In) => void): () => void }) => void, + start: (transport: { + send(message: Out): void; + sendAndFlush(message: Out): Promise; + onMessage(handler: (message: In) => void): () => void; + }) => void, ): Promise { const { promise: shuttingDown, resolve: shutdown } = Promise.withResolvers(); + type IpcSend = (this: NodeJS.Process, message: unknown, callback?: (error: Error | null) => void) => boolean; + // `process.send` only exists when spawned with an IPC channel; the parent + // always spawns us that way. If it's missing, the parent vanished and + // there's no one to talk to. + const ipcSend = (): IpcSend | undefined => (process as NodeJS.Process & { send?: IpcSend }).send; const send = (message: Out): void => { - // `process.send` only exists when spawned with an IPC channel; the - // parent always spawns us that way. If it's missing, the parent - // vanished and there's no one to talk to. - const sender = (process as NodeJS.Process & { send?: (m: unknown) => boolean }).send; + const sender = ipcSend(); if (!sender) { shutdown(); return; @@ -191,8 +197,24 @@ async function runIpcSubprocessWorker( shutdown(); } }; + const sendAndFlush = (message: Out): Promise => { + const sender = ipcSend(); + if (!sender) { + shutdown(); + return Promise.resolve(); + } + const { promise, resolve } = Promise.withResolvers(); + try { + sender.call(process, message, () => resolve()); + } catch { + shutdown(); + resolve(); + } + return promise; + }; start({ send, + sendAndFlush, onMessage(handler) { const wrap = (data: unknown): void => handler(data as In); process.on("message", wrap); diff --git a/packages/coding-agent/src/tts/index.ts b/packages/coding-agent/src/tts/index.ts index f33fb795d..7f114f9a2 100644 --- a/packages/coding-agent/src/tts/index.ts +++ b/packages/coding-agent/src/tts/index.ts @@ -1,6 +1,7 @@ export * from "./downloader"; export * from "./models"; export * from "./runtime"; +export * from "./speakable"; export * from "./tts-client"; export * from "./tts-protocol"; export * from "./tts-worker"; diff --git a/packages/coding-agent/src/tts/speakable.ts b/packages/coding-agent/src/tts/speakable.ts new file mode 100644 index 000000000..47ad3d77f --- /dev/null +++ b/packages/coding-agent/src/tts/speakable.ts @@ -0,0 +1,382 @@ +/** + * Streaming markdown → speakable-segment transform for assistant speech. + * + * Sits between the assistant's raw streaming text deltas and the TTS engine, + * deciding both *what* is worth speaking and *when* a piece of text is ready + * to synthesize. Three passes: + * + * 1. Block pass (per character, stateful): drops fenced code blocks and table + * rows, strips heading/bullet/blockquote markers (numbered-list markers are + * spoken as "1, …"), and turns newlines into hard segment breaks. + * 2. Segmentation (stateful): emits a segment the moment a sentence boundary + * appears — no next-sentence confirmation, which is what made the previous + * engine-side splitter stall a full sentence behind generation. The first + * segment cuts early at a clause boundary for fast time-to-first-audio, and + * over-long unpunctuated runs are force-split so no segment exceeds the + * synthesizer's input budget. + * 3. Inline normalization (per segment): markdown links speak their label, + * bare URLs speak their host, inline-code ticks and emphasis markers are + * stripped, multi-directory file paths collapse to their basename, HTML + * tags are dropped, and whitespace is collapsed. Segments with no letters + * or digits left are not spoken at all. + * + * Pure and synchronous — the vocalizer owns timers (idle flush) and the + * session lifecycle, so this class stays trivially unit-testable. + */ + +/** Minimum length before the very first segment may cut at a sentence boundary. */ +const FIRST_SEGMENT_MIN = 12; +/** Buffer length past which the first segment may cut at a clause boundary instead. */ +const FIRST_CLAUSE_MIN = 40; +/** Hard cap for the first segment: force a word cut for fast time-to-first-audio. */ +const FIRST_FORCED_MAX = 140; +/** Minimum segment length once speech has started (merges stubby sentences). */ +const MIN_SEGMENT = 24; +/** + * Mid-stream soft cut: once this much unpunctuated text is buffered, split at + * a clause boundary instead of waiting for the sentence to end. Long sentences + * synthesize as clause-sized pieces, keeping the playback pipeline fed (a + * 280-char segment costs ~6s of synthesis — enough to drain the player dry). + */ +const SOFT_CLAUSE_LEN = 160; +/** + * Hard cap per segment. Kokoro's `generate()` truncates past ~510 phoneme + * tokens rather than splitting, so every emission path must stay well under it. + */ +const MAX_SEGMENT = 280; + +/** Sentence-ending punctuation, optional closers, then whitespace. */ +const SENTENCE_BOUNDARY_RE = /[.!?…]+[)\]"'»”’]*\s/g; +/** Clause punctuation followed by whitespace — early-cut and force-split points. */ +const CLAUSE_BOUNDARY_RE = /[,;:—–]\s/g; +/** Abbreviations whose trailing dot is not a sentence boundary. */ +const ABBREVIATION_RE = /(?:^|\s)(?:e\.g|i\.e|etc|vs|Mr|Mrs|Ms|Dr|St|No)\.$/i; + +/** Line-start prefixes that may still grow into a block marker. */ +const UNDECIDED_PREFIX_RE = /^(?:#{1,6}|[-*+]|-{2,}|\*{2,}|_{2,}|\d{1,3}|\d{1,3}[.)]|>+|`{1,2}|~{1,2})$/; +/** A whole line that is a horizontal rule (or setext underline) — silence. */ +const HR_LINE_RE = /^(?:-{3,}|\*{3,}|_{3,})\s*$/; + +const IMAGE_RE = /!\[([^\]]*)\]\(([^()]*)\)/g; +const LINK_RE = /\[([^\]]+)\]\(([^()]*)\)/g; +const AUTOLINK_RE = /<(https?:\/\/[^\s>]+)>/g; +const BARE_URL_RE = /\bhttps?:\/\/[^\s<>()"'\]]+|\bwww\.[\w-]+(?:\.[\w-]+)+[^\s<>()"'\]]*/g; +const INLINE_CODE_RE = /`{1,2}([^`]+)`{1,2}/g; +const BOLD_STRIKE_RE = /\*\*|__|~~/g; +const EMPHASIS_ASTERISK_RE = /\*(?=\S)|(?<=\S)\*/g; +const EMPHASIS_UNDERSCORE_RE = /(^|\s)_+|_+(?=\s|$)/g; +const HTML_TAG_RE = /<\/?[a-zA-Z][^<>]*>/g; +const HR_INLINE_RE = /(^|\s)[-*_]{3,}(?=\s|$)/g; +const PATH_RE = /(^|[\s("'`])((?:~|\.{1,2})?\/?[\w.@+-]+(?:\/[\w.@+-]+){2,}\/?)/g; +const HAS_SPEAKABLE_RE = /[\p{L}\p{N}]/u; + +/** "https://github.com/foo/bar?x#y" → "github.com". */ +function speakableUrl(url: string): string { + return url + .replace(/^[a-z][\w+.-]*:\/\//i, "") + .replace(/^www\./i, "") + .replace(/[/?#].*$/, ""); +} + +/** + * Collapse one raw segment to its speakable form; empty string when nothing + * in it is worth vocalizing (pure markup, URLs-only, whitespace). + */ +function normalizeSpeakable(raw: string): string { + const spoken = raw + .replace(IMAGE_RE, "$1") + .replace(LINK_RE, "$1") + .replace(AUTOLINK_RE, (_match, url: string) => speakableUrl(url)) + .replace(BARE_URL_RE, match => speakableUrl(match)) + .replace(INLINE_CODE_RE, "$1") + .replace(BOLD_STRIKE_RE, "") + .replace(EMPHASIS_ASTERISK_RE, "") + .replace(EMPHASIS_UNDERSCORE_RE, "$1") + .replace(HTML_TAG_RE, " ") + .replace(HR_INLINE_RE, "$1") + .replace(PATH_RE, (_match, lead: string, path: string) => { + // "packages/coding-agent/src/tts/vocalizer.ts" → "vocalizer.ts". + const parts = path.split("/").filter(part => part.length > 0); + return lead + (parts[parts.length - 1] ?? path); + }) + .replace(/\s+/g, " ") + .trim(); + return HAS_SPEAKABLE_RE.test(spoken) ? spoken : ""; +} + +/** + * Earliest sentence boundary at or past `min` chars; -1 when none. Skips cuts + * that would strand an unclosed inline-code span or split an abbreviation. + */ +function findSentenceCut(text: string, min: number): number { + SENTENCE_BOUNDARY_RE.lastIndex = 0; + for (let match = SENTENCE_BOUNDARY_RE.exec(text); match; match = SENTENCE_BOUNDARY_RE.exec(text)) { + const cut = match.index + match[0].length; + if (cut < min) continue; + const head = text.slice(0, cut); + if (ABBREVIATION_RE.test(head.trimEnd())) continue; + if ((head.match(/`/g)?.length ?? 0) % 2 !== 0) continue; + return cut; + } + return -1; +} + +/** Earliest clause boundary at or past `min` chars; -1 when none. */ +function findClauseCut(text: string, min: number): number { + CLAUSE_BOUNDARY_RE.lastIndex = 0; + for (let match = CLAUSE_BOUNDARY_RE.exec(text); match; match = CLAUSE_BOUNDARY_RE.exec(text)) { + const cut = match.index + match[0].length; + if (cut >= min) return cut; + } + return -1; +} + +/** + * Latest clause boundary in `[min, max]` chars; -1 when none. Keeps soft-cut + * segments grouped near the target length instead of shaving off the earliest + * stale clause. + */ +function findLastClauseCut(text: string, min: number, max: number): number { + CLAUSE_BOUNDARY_RE.lastIndex = 0; + let best = -1; + for (let match = CLAUSE_BOUNDARY_RE.exec(text); match; match = CLAUSE_BOUNDARY_RE.exec(text)) { + const cut = match.index + match[0].length; + if (cut > max) break; + if (cut >= min) best = cut; + } + return best; +} + +/** Word-level cut for text with no usable punctuation: last space at or before `max`. */ +function findForcedCut(text: string, max: number): number { + const space = text.lastIndexOf(" ", max); + return space > 0 ? space + 1 : Math.min(max, text.length); +} + +/** How a line-start prefix resolved. */ +type PrefixDecision = + | { kind: "undecided" } + | { kind: "prose"; text: string } + | { kind: "marker"; spoken: string } + | { kind: "swallow" } + | { kind: "fence"; fence: string }; + +function classifyPrefix(prefix: string): PrefixDecision { + if (prefix === "|") return { kind: "swallow" }; + if (/^(?:`{3}|~{3})/.test(prefix)) return { kind: "fence", fence: prefix.slice(0, 3) }; + if (/^#{1,6}[ \t]/.test(prefix)) return { kind: "marker", spoken: "" }; + if (/^[-*+][ \t]/.test(prefix)) return { kind: "marker", spoken: "" }; + const numbered = /^(\d{1,3})[.)][ \t]/.exec(prefix); + if (numbered) return { kind: "marker", spoken: `${numbered[1]}, ` }; + if (/^>+/.test(prefix) && !/^>+$/.test(prefix)) { + return { kind: "prose", text: prefix.replace(/^>+[ \t]?/, "") }; + } + if (UNDECIDED_PREFIX_RE.test(prefix)) return { kind: "undecided" }; + return { kind: "prose", text: prefix }; +} + +/** Block-pass state: where the current character lands. */ +type BlockMode = "linestart" | "prose" | "swallow" | "code"; + +/** + * One per utterance. Feed raw assistant deltas through {@link push}; each call + * returns the segments that became ready to speak. {@link flush} drains the + * remainder at message end; {@link flushIdle} drains it when generation stalls + * mid-sentence so speech doesn't sit on buffered text through a tool call. + */ +export class SpeakableStream { + #mode: BlockMode = "linestart"; + /** Pending line-start characters while the block marker is still ambiguous. */ + #prefix = ""; + /** Opening fence of the code block being swallowed (``` or ~~~). */ + #fence = ""; + /** First characters of the current line inside a code block (fence-close probe). */ + #codeLine = ""; + /** Mode to enter after the current swallowed line ends (code for an opening fence). */ + #afterSwallow: BlockMode = "linestart"; + /** Prose accumulator the segmenter cuts from. */ + #buf = ""; + /** Whether anything has been emitted yet (enables the fast first segment). */ + #spoke = false; + + /** Consume a raw delta; returns segments now ready to speak, in order. */ + push(delta: string): string[] { + const out: string[] = []; + for (const ch of delta) this.#consume(ch, out); + this.#extract(out); + return out; + } + + /** Message end: drain everything left, including a trailing partial sentence. */ + flush(): string[] { + const out: string[] = []; + if (this.#mode === "linestart" && this.#prefix.length > 0 && !HR_LINE_RE.test(this.#prefix)) { + this.#buf += this.#prefix; + } + this.#prefix = ""; + this.#mode = "linestart"; + this.#drain(out); + return out; + } + + /** + * Generation stalled (tool call, thinking block): speak what we have rather + * than sit silent on buffered text. Keeps block state so the stream resumes + * afterwards, and refuses stubby mid-sentence fragments — the buffer must be + * a complete thought (trailing sentence punctuation) or at least + * {@link MIN_SEGMENT} long, so a stall right after "The" stays silent + * instead of turning into choppy one-word speech. + */ + flushIdle(): string[] { + const out: string[] = []; + const pending = this.#buf.trimEnd(); + const completeThought = /[.!?…][)\]"'»”’]*$/.test(pending); + if (!completeThought && pending.length < MIN_SEGMENT) return out; + this.#drain(out); + return out; + } + + #consume(ch: string, out: string[]): void { + switch (this.#mode) { + case "linestart": + this.#consumeLineStart(ch, out); + return; + case "prose": + if (ch === "\n") this.#hardBreak(out); + else this.#buf += ch; + return; + case "swallow": + if (ch === "\n") this.#mode = this.#afterSwallow; + return; + case "code": + this.#consumeCode(ch); + return; + } + } + + #consumeLineStart(ch: string, out: string[]): void { + if (ch === "\n") { + // The whole line fit in the prefix: an hr/blank line is silence; a + // short undecided prefix ("Hi.", "OK") was prose all along. + const line = this.#prefix; + this.#prefix = ""; + if (line.length > 0 && !HR_LINE_RE.test(line)) this.#buf += line; + this.#hardBreak(out); + return; + } + this.#prefix += ch; + const decision = classifyPrefix(this.#prefix); + if (decision.kind === "undecided") { + if (this.#prefix.length > 8) { + this.#buf += this.#prefix; + this.#prefix = ""; + this.#mode = "prose"; + } + return; + } + this.#prefix = ""; + switch (decision.kind) { + case "prose": + this.#buf += decision.text; + this.#mode = "prose"; + return; + case "marker": + this.#buf += decision.spoken; + this.#mode = "prose"; + return; + case "swallow": + this.#mode = "swallow"; + this.#afterSwallow = "linestart"; + return; + case "fence": + this.#fence = decision.fence; + this.#codeLine = ""; + this.#mode = "swallow"; + this.#afterSwallow = "code"; + return; + } + } + + #consumeCode(ch: string): void { + if (ch === "\n") { + this.#codeLine = ""; + return; + } + if (this.#codeLine.length < 3) { + this.#codeLine += ch; + if (this.#codeLine === this.#fence) { + // Closing fence: swallow the rest of its line, then resume prose. + this.#mode = "swallow"; + this.#afterSwallow = "linestart"; + } + } + } + + /** Newline in prose: everything buffered is a complete unit — emit it now. */ + #hardBreak(out: string[]): void { + this.#mode = "linestart"; + this.#drain(out); + } + + /** Emit every buffered character, force-splitting anything over the cap. */ + #drain(out: string[]): void { + let text = this.#buf; + this.#buf = ""; + while (text.length > MAX_SEGMENT) { + const cut = findForcedCut(text, MAX_SEGMENT); + this.#emit(text.slice(0, cut), out); + text = text.slice(cut); + } + this.#emit(text, out); + } + + /** Cut ready segments off the front of the buffer (streaming path). */ + #extract(out: string[]): void { + for (;;) { + const buf = this.#buf; + const min = this.#spoke ? MIN_SEGMENT : FIRST_SEGMENT_MIN; + const sentence = findSentenceCut(buf, min); + if (sentence !== -1) { + this.#cut(sentence, out); + continue; + } + if (!this.#spoke && buf.length >= FIRST_CLAUSE_MIN) { + const clause = findClauseCut(buf, FIRST_SEGMENT_MIN); + if (clause !== -1) { + this.#cut(clause, out); + continue; + } + if (buf.length >= FIRST_FORCED_MAX) { + this.#cut(findForcedCut(buf, FIRST_FORCED_MAX), out); + continue; + } + } + if (this.#spoke && buf.length >= SOFT_CLAUSE_LEN) { + const clause = findLastClauseCut(buf, MIN_SEGMENT, SOFT_CLAUSE_LEN); + if (clause !== -1) { + this.#cut(clause, out); + continue; + } + } + if (buf.length > MAX_SEGMENT) { + const clause = findLastClauseCut(buf, MIN_SEGMENT, MAX_SEGMENT); + this.#cut(clause !== -1 ? clause : findForcedCut(buf, MAX_SEGMENT), out); + continue; + } + return; + } + } + + #cut(at: number, out: string[]): void { + const head = this.#buf.slice(0, at); + this.#buf = this.#buf.slice(at); + this.#emit(head, out); + } + + #emit(raw: string, out: string[]): void { + const spoken = normalizeSpeakable(raw); + if (!spoken) return; + out.push(spoken); + this.#spoke = true; + } +} diff --git a/packages/coding-agent/src/tts/streaming-player.ts b/packages/coding-agent/src/tts/streaming-player.ts index 832d40631..586bca885 100644 --- a/packages/coding-agent/src/tts/streaming-player.ts +++ b/packages/coding-agent/src/tts/streaming-player.ts @@ -4,14 +4,17 @@ * Replaces the spawn-`afplay`-per-sentence approach (a fresh process per chunk * meant audible gaps, per-spawn latency, and no way to interrupt a clip mid-play) * with a single persistent player process fed raw 32-bit-float mono PCM over - * stdin. Chunks are queued and drained by one writer so sentences play back to + * stdin. Chunks are queued and drained by one writer so segments play back to * back; writes are paced to stay only {@link LEAD_SECONDS} ahead of realtime so * ducking and stop take effect promptly instead of after seconds of buffered * audio. {@link StreamingAudioPlayer.stop} kills the process for instant silence. * - * Where no streaming backend exists (Windows, or macOS without the bundled - * ffmpeg), it degrades to the per-file {@link playAudioFile} path so speech still - * works — just without gapless playback or mid-clip interruption. + * Where no streaming backend exists (Windows, or a host without ffmpeg/sox), it + * degrades to the per-file {@link playAudioFile} path so speech still works — + * just without gapless playback or mid-clip interruption. A backend that spawns + * but dies early (e.g. an ffmpeg built without its platform audio device) is + * detected via its exit and the session downgrades to per-file playback without + * dropping the chunk being played. */ import * as fs from "node:fs/promises"; import * as os from "node:os"; @@ -41,8 +44,7 @@ export interface StreamingPlayerLookup { * and plays it to the default output device. An empty list means no streaming * backend is available and the caller should fall back to per-file playback. * - * - darwin: none; `afplay` is file-only, so macOS uses the interruptible - * per-file fallback. + * - darwin: `ffmpeg` (AudioToolbox output device) → sox's `play` (coreaudio). * - linux/other POSIX: `ffmpeg` (`-f pulse` then `-f alsa`) → `paplay`/`aplay` * raw fallbacks. * - win32: none (PowerShell `SoundPlayer` is file-only). @@ -57,7 +59,19 @@ export function streamingPlayerCommandsFor( const rate = String(sampleRate > 0 ? sampleRate : DEFAULT_SAMPLE_RATE); const input = ["-loglevel", "error", "-nostdin", "-f", "f32le", "-ar", rate, "-ac", "1", "-i", "pipe:0"]; - if (platform === "darwin") return []; + if (platform === "darwin") { + const commands: PlayerCommand[] = []; + const ffmpegBin = ffmpeg(); + if (ffmpegBin) commands.push({ cmd: ffmpegBin, args: [...input, "-f", "audiotoolbox", "default"] }); + const play = which("play"); + if (play) { + commands.push({ + cmd: play, + args: ["-q", "-t", "raw", "-e", "floating-point", "-b", "32", "-r", rate, "-c", "1", "-"], + }); + } + return commands; + } if (platform === "win32") { return []; } @@ -87,6 +101,8 @@ export class StreamingAudioPlayer { #mode: "stream" | "file" = "file"; #proc: Subprocess<"pipe", "ignore", "ignore"> | null = null; #sink: FileSink | null = null; + /** Streaming backends not yet tried; consumed head-first by {@link #spawnStream}. */ + #candidates: PlayerCommand[] | null = null; #writtenSec = 0; #startedAt = 0; #started = false; @@ -140,20 +156,37 @@ export class StreamingAudioPlayer { } catch {} } + /** + * Spawn the next untried streaming backend; false once the list is + * exhausted. A backend that spawns but dies early (e.g. an ffmpeg built + * without this platform's audio output device) would otherwise swallow PCM + * into a dead pipe, so its exit advances to the next candidate — or to + * per-file playback — and #writeStream's failure path replays the + * in-flight chunk. + */ #spawnStream(): boolean { - for (const command of streamingPlayerCommandsFor(process.platform, this.#sampleRate)) { + this.#candidates ??= streamingPlayerCommandsFor(process.platform, this.#sampleRate); + for (let command = this.#candidates.shift(); command; command = this.#candidates.shift()) { + const { cmd, args } = command; try { - const proc = Bun.spawn([command.cmd, ...command.args], { + const proc = Bun.spawn([cmd, ...args], { stdin: "pipe", stdout: "ignore", stderr: "ignore", }); this.#proc = proc; this.#sink = proc.stdin; + void proc.exited.then(code => { + if (this.#proc !== proc || this.#stopped || this.#inputClosed) return; + logger.debug("tts: streaming backend exited early; trying next backend", { cmd, code }); + this.#proc = null; + this.#sink = null; + this.#mode = this.#spawnStream() ? "stream" : "file"; + }); return true; } catch (error) { logger.debug("tts: streaming player spawn failed", { - cmd: command.cmd, + cmd, error: error instanceof Error ? error.message : String(error), }); } @@ -184,8 +217,20 @@ export class StreamingAudioPlayer { await Bun.sleep((ahead - LEAD_SECONDS) * 1000); if (this.#stopped) return; } - this.#writeStream(chunk); - this.#writtenSec += chunk.length / this.#sampleRate; + if (this.#writeStream(chunk)) { + this.#writtenSec += chunk.length / this.#sampleRate; + continue; + } + // Backend died mid-write: move to the next streaming candidate + // (or the file path) and replay this exact chunk so nothing is + // dropped. + this.#mode = this.#spawnStream() ? "stream" : "file"; + if (this.#mode === "stream" && this.#writeStream(chunk)) { + this.#writtenSec += chunk.length / this.#sampleRate; + } else { + this.#mode = "file"; + await this.#playFile(chunk); + } } else { await this.#playFile(chunk); } @@ -219,16 +264,19 @@ export class StreamingAudioPlayer { return promise; } - #writeStream(pcm: Float32Array): void { + /** Write one chunk into the backend's stdin; false when the sink is gone or broken. */ + #writeStream(pcm: Float32Array): boolean { const sink = this.#sink; - if (!sink) return; + if (!sink) return false; try { sink.write(this.#bytes(pcm)); sink.flush(); + return true; } catch (error) { logger.debug("tts: streaming write failed", { error: error instanceof Error ? error.message : String(error), }); + return false; } } diff --git a/packages/coding-agent/src/tts/tts-client.ts b/packages/coding-agent/src/tts/tts-client.ts index df452c797..4dfebdc0f 100644 --- a/packages/coding-agent/src/tts/tts-client.ts +++ b/packages/coding-agent/src/tts/tts-client.ts @@ -42,7 +42,7 @@ export interface TtsStreamOptions { signal?: AbortSignal; } -/** One synthesized sentence of a streaming session, in emission order. */ +/** One synthesized segment of a streaming session, in emission order. */ export interface TtsAudioChunk { index: number; text: string; @@ -51,10 +51,10 @@ export interface TtsAudioChunk { } /** - * A live streaming-synthesis session. Feed text incrementally with {@link push} - * and close the input with {@link end}; `chunks` yields each synthesized - * sentence's audio as soon as it is ready, then completes once the worker - * finishes draining the closed input. + * A live streaming-synthesis session. Feed complete speakable segments with + * {@link push} (the worker synthesizes each push as-is) and close the input + * with {@link end}; `chunks` yields each segment's audio as soon as it is + * ready, then completes once the worker finishes draining the closed input. */ export interface TtsStreamHandle { push(text: string): void; @@ -239,11 +239,12 @@ export class TtsClient { } /** - * Open a streaming-synthesis session. Text is fed incrementally through the - * returned handle's `push`/`end`; audio is emitted one synthesized sentence at - * a time via `chunks`, so playback can begin before the full text is known. - * Returns an inert handle (immediately-ended `chunks`) for unknown models or - * an already-aborted signal, and fails the iterator if the worker cannot spawn. + * Open a streaming-synthesis session. Complete speakable segments are fed + * through the returned handle's `push`/`end`; audio is emitted one segment + * at a time via `chunks`, so playback can begin before the full text is + * known. Returns an inert handle (immediately-ended `chunks`) for unknown + * models or an already-aborted signal, and fails the iterator if the worker + * cannot spawn. */ synthesizeStream(modelKey: string, options: TtsStreamOptions = {}): TtsStreamHandle { if (!isTtsLocalModelKey(modelKey) || options.signal?.aborted) { diff --git a/packages/coding-agent/src/tts/tts-protocol.ts b/packages/coding-agent/src/tts/tts-protocol.ts index 3d9613d2e..5e88f4dbe 100644 --- a/packages/coding-agent/src/tts/tts-protocol.ts +++ b/packages/coding-agent/src/tts/tts-protocol.ts @@ -24,10 +24,12 @@ export type TtsWorkerInbound = | { type: "ping"; id: string } | { type: "synthesize"; id: string; modelKey: TtsLocalModelKey; text: string; voice?: string } | { type: "download"; id: string; modelKey: TtsLocalModelKey } - // Streaming synthesis: a session is opened with `stream-start`, fed incrementally - // with `stream-push`, and closed with `stream-end`. `stream-cancel` interrupts - // without a final drain. The worker emits an `audio-chunk` per synthesized - // sentence and a final `stream-done` only for non-cancelled sessions. + // Streaming synthesis: a session is opened with `stream-start`, fed complete + // speakable segments with `stream-push` (the parent's SpeakableStream does all + // splitting/normalization; the worker synthesizes each push as-is), and closed + // with `stream-end`. `stream-cancel` interrupts without a final drain. The + // worker emits an `audio-chunk` per segment and a final `stream-done` only for + // non-cancelled sessions. | { type: "stream-start"; id: string; modelKey: TtsLocalModelKey; voice?: string } | { type: "stream-push"; id: string; text: string } | { type: "stream-end"; id: string } @@ -40,7 +42,7 @@ export type TtsWorkerOutbound = | { type: "error"; id: string; error: string } | { type: "progress"; id: string; event: TtsProgressEvent } | { type: "log"; level: "debug" | "warn" | "error"; msg: string; meta?: Record } - // One synthesized sentence of a streaming session, in emission order, followed + // One synthesized segment of a streaming session, in emission order, followed // by a single `stream-done` once the input stream is closed and drained. | { type: "audio-chunk"; id: string; index: number; text: string; pcm: Float32Array; sampleRate: number } | { type: "stream-done"; id: string }; @@ -56,5 +58,12 @@ export type TtsWorkerOutbound = */ export interface TtsTransport { send(message: TtsWorkerOutbound): void; + /** + * Send and resolve once the message has drained into the IPC channel. + * Streaming synthesis awaits this per audio chunk: ONNX inference blocks + * the worker's event loop for seconds at a time, so fire-and-forget sends + * queue unflushed until the session ends and arrive as one burst. + */ + sendAndFlush(message: TtsWorkerOutbound): Promise; onMessage(handler: (message: TtsWorkerInbound) => void): () => void; } diff --git a/packages/coding-agent/src/tts/tts-worker.ts b/packages/coding-agent/src/tts/tts-worker.ts index 04a659bef..00892173e 100644 --- a/packages/coding-agent/src/tts/tts-worker.ts +++ b/packages/coding-agent/src/tts/tts-worker.ts @@ -39,20 +39,6 @@ type KokoroDevice = "cpu" | "wasm" | "webgpu"; /** A loaded Kokoro voice synthesizer (subset of `kokoro-js`'s `KokoroTTS`). */ interface KokoroTtsInstance { generate(text: string, options: { voice: string }): Promise; - stream( - text: string | TextSplitterStreamInstance, - options: { voice: string }, - ): AsyncGenerator<{ text: string; phonemes: string; audio: RawAudio }, void, void>; -} - -/** - * Incremental text source for {@link KokoroTtsInstance.stream} (subset of - * `kokoro-js`'s `TextSplitterStream`). Text pushed at any time is split into - * complete sentences; `close` flushes the trailing buffer and ends the stream. - */ -interface TextSplitterStreamInstance { - push(...texts: string[]): void; - close(): void; } /** `KokoroTTS` static surface used to load a model from the Hugging Face Hub. */ @@ -67,7 +53,6 @@ interface KokoroRuntime { }, ): Promise; }; - TextSplitterStream: new () => TextSplitterStreamInstance; } /** @@ -96,15 +81,18 @@ const kokoroRuntime = new MemoizedRuntime(); /** * In-flight streaming sessions keyed by request id. A session is created on - * `stream-start` and torn down when its generator finishes. Text pushed before - * the model finishes loading is held in `buffered` and flushed into the splitter - * once it exists; pushes after that go straight to the live splitter. + * `stream-start` and torn down when its run loop finishes. Each `stream-push` + * carries one complete speakable segment (the parent's `SpeakableStream` does + * all splitting and normalization); segments queue here and the run loop + * synthesizes them in arrival order, waking via `wake` when idle. */ interface StreamSession { modelKey: TtsLocalModelKey; voice: string | undefined; - buffered: string[]; - splitter: TextSplitterStreamInstance | null; + /** Speakable segments awaiting synthesis, in arrival order. */ + queue: string[]; + /** Resolves the run loop's idle wait when a push/end/cancel arrives. */ + wake: (() => void) | null; ended: boolean; cancelled: boolean; } @@ -316,43 +304,51 @@ async function handleQueuedRequest( } /** - * Drive one streaming session to completion: load the model, create the - * splitter, flush any text pushed before the model was ready, then emit one - * `audio-chunk` per synthesized sentence followed by a single `stream-done`. - * Serialized through {@link synthesizeQueue} so it never interleaves model - * access with a batch synthesize/download. + * Drive one streaming session to completion: load the model, then synthesize + * each queued segment in arrival order — one `audio-chunk` per segment, + * followed by a single `stream-done`. Chunk sends are drained before the next + * segment's inference (see the comment at the send site). Serialized through + * {@link synthesizeQueue} so it never interleaves model access with a batch + * synthesize/download. */ async function runStreamSession(transport: TtsTransport, id: string, session: StreamSession): Promise { try { - if (session.cancelled) return; - const runtime = await loadKokoroRuntime(transport, id, session.modelKey); if (session.cancelled) return; const synthesizer = await loadModel(session.modelKey, transport, id); if (session.cancelled) return; const spec = getTtsLocalModelSpec(session.modelKey); - const splitter = new runtime.TextSplitterStream(); - // Flush buffered text before exposing the splitter so a push racing this - // block can't slip ahead of the already-queued fragments. - for (const text of session.buffered) { - if (session.cancelled) return; - splitter.push(text); - } - session.buffered = []; - session.splitter = splitter; - if (session.ended || session.cancelled) splitter.close(); const voice = resolveTtsVoice(session.modelKey, session.voice); let index = 0; - for await (const chunk of synthesizer.stream(splitter, { voice })) { + while (!session.cancelled) { + const segment = session.queue.shift(); + if (segment === undefined) { + if (session.ended) break; + const { promise, resolve } = Promise.withResolvers(); + session.wake = resolve; + // Re-check after arming: a push/end/cancel racing the empty shift. + if (session.queue.length > 0 || session.ended || session.cancelled) { + session.wake = null; + resolve(); + } + await promise; + continue; + } + const output = await synthesizer.generate(segment, { voice }); if (session.cancelled) break; - const audio = Array.isArray(chunk.audio.audio) ? chunk.audio.audio[0] : chunk.audio.audio; + const audio = Array.isArray(output.audio) ? output.audio[0] : output.audio; if (!audio) continue; - transport.send({ + // Drain the IPC write before the next segment's inference: ONNX + // blocks this event loop for seconds at a time, so a fire-and-forget + // send would sit in the pipe queue until the session ends and every + // chunk would arrive in one burst (long silence, then all segments + // at once) instead of streaming per-segment. + await transport.sendAndFlush({ type: "audio-chunk", id, index: index++, - text: chunk.text, + text: segment, pcm: audio, - sampleRate: chunk.audio.sampling_rate || spec?.sampleRate || 24_000, + sampleRate: output.sampling_rate || spec?.sampleRate || 24_000, }); } if (!session.cancelled) transport.send({ type: "stream-done", id }); @@ -370,8 +366,8 @@ function startStreamSession( const session: StreamSession = { modelKey: message.modelKey, voice: message.voice, - buffered: [], - splitter: null, + queue: [], + wake: null, ended: false, cancelled: false, }; @@ -382,26 +378,33 @@ function startStreamSession( ); } +/** Wake the session's run loop if it is parked on an empty queue. */ +function wakeStreamSession(session: StreamSession): void { + const wake = session.wake; + session.wake = null; + wake?.(); +} + function pushToStreamSession(id: string, text: string): void { const session = streamSessions.get(id); if (!session || session.cancelled) return; - if (session.splitter) session.splitter.push(text); - else session.buffered.push(text); + session.queue.push(text); + wakeStreamSession(session); } function endStreamSession(id: string): void { const session = streamSessions.get(id); if (!session || session.cancelled) return; session.ended = true; - session.splitter?.close(); + wakeStreamSession(session); } function cancelStreamSession(id: string): void { const session = streamSessions.get(id); if (!session) return; session.cancelled = true; - session.buffered = []; - session.splitter?.close(); + session.queue.length = 0; + wakeStreamSession(session); streamSessions.delete(id); } diff --git a/packages/coding-agent/src/tts/vocalizer.ts b/packages/coding-agent/src/tts/vocalizer.ts index 35507e266..89eee0e95 100644 --- a/packages/coding-agent/src/tts/vocalizer.ts +++ b/packages/coding-agent/src/tts/vocalizer.ts @@ -2,12 +2,17 @@ * Streaming assistant speech-vocalization. * * The vocalizer turns the assistant's STREAMING output into spoken audio as a - * side effect of the normal turn. Text deltas are streamed *straight into the - * TTS engine* ({@link Vocalizer.pushDelta} → the worker's incremental text - * input): the engine splits the running text at sentence boundaries and emits - * one audio chunk per sentence, which a single {@link StreamingAudioPlayer} - * plays back gaplessly. So the assistant starts speaking sentence 1 while later - * sentences are still being generated — low latency, never overlapping. + * side effect of the normal turn. Text deltas run through a + * {@link SpeakableStream} — which drops code/tables/markup, speaks link labels + * and URL hosts instead of raw URLs, and cuts speakable segments the moment a + * boundary appears — and each ready segment is pushed to the TTS worker, which + * synthesizes it into one audio chunk. A single {@link StreamingAudioPlayer} + * plays the chunks back gaplessly, so the assistant starts speaking the first + * clause while later sentences are still being generated. + * + * An idle timer covers generation stalls: when no delta arrives for + * {@link IDLE_FLUSH_MS} (tool call, thinking block) the buffered partial + * sentence is spoken rather than held silent. * * Overspeech control: * - {@link clear} stops playback instantly (kills the player) and aborts @@ -25,9 +30,13 @@ import { logger } from "@oh-my-pi/pi-utils"; import { settings } from "../config/settings"; import { DEFAULT_TTS_VOICE } from "./models"; +import { SpeakableStream } from "./speakable"; import { createStreamingPlayer, DUCK_GAIN } from "./streaming-player"; import { type TtsStreamHandle, ttsClient } from "./tts-client"; +/** Quiet time on the delta stream before the buffered partial is spoken. */ +const IDLE_FLUSH_MS = 1000; + export interface VocalizerPlayer { start(sampleRate: number): void; write(pcm: Float32Array): void; @@ -39,6 +48,10 @@ export interface VocalizerPlayer { export class Vocalizer { /** Open stream session for the current utterance; null when none is active. */ #handle: TtsStreamHandle | null = null; + /** Markdown → speakable-segment transform for the current utterance. */ + #speakable: SpeakableStream | null = null; + /** Fires when the delta stream goes quiet mid-sentence; speaks the partial. */ + #idleTimer: NodeJS.Timeout | null = null; /** Aborts the in-flight session on {@link clear}; replaced per session. */ #abort: AbortController | null = null; /** The current session's player; stopped on {@link clear}, gain-tracked for ducking. */ @@ -54,22 +67,30 @@ export class Vocalizer { } /** - * Stream a delta of assistant text into the engine. No-op when vocalization - * is disabled. The engine buffers the running text and emits audio for each - * complete sentence; the trailing partial is flushed by {@link flush}. + * Stream a delta of assistant text into the pipeline. No-op when + * vocalization is disabled. The synthesis session (worker, player) is only + * opened once the first speakable segment exists, so a reply that + * normalizes to silence (pure code, tables, URLs) costs nothing. The + * trailing partial is flushed by {@link flush} or the idle timer. */ pushDelta(text: string): void { if (!settings.get("speech.enabled")) return; if (!text) return; - this.#ensureSession().push(text); + this.#speakable ??= new SpeakableStream(); + this.#pushSegments(this.#speakable.push(text)); + this.#armIdleFlush(this.#speakable); } /** - * Close the current input stream (call at message/turn end). The engine - * flushes its trailing partial as a final chunk; the player keeps draining - * queued audio until it completes. + * Close the current input stream (call at message/turn end). Drains the + * trailing partial as final segments; the player keeps draining queued + * audio until it completes. */ flush(): void { + this.#clearIdleTimer(); + const speakable = this.#speakable; + this.#speakable = null; + if (speakable) this.#pushSegments(speakable.flush()); this.#handle?.end(); this.#handle = null; } @@ -79,9 +100,7 @@ export class Vocalizer { * message): stream it in and immediately close the input. No-op when disabled. */ speak(text: string): void { - if (!settings.get("speech.enabled")) return; - if (!text) return; - this.#ensureSession().push(text); + this.pushDelta(text); this.flush(); } @@ -90,6 +109,8 @@ export class Vocalizer { * synthesis (new turn / user message / Esc interrupt). Audio stops at once. */ clear(): void { + this.#clearIdleTimer(); + this.#speakable = null; this.#handle = null; this.#abort?.abort(); this.#abort = null; @@ -114,9 +135,17 @@ export class Vocalizer { return this.#chain; } + /** Feed ready segments to the synthesizer, opening the session lazily. */ + #pushSegments(segments: string[]): void { + if (segments.length === 0) return; + const handle = this.#ensureSession(); + for (const segment of segments) handle.push(segment); + } + /** - * Open a streaming-synthesis session lazily on the first delta and chain its - * playback after any prior session's, so sequential utterances never overlap. + * Open a streaming-synthesis session lazily on the first speakable segment + * and chain its playback after any prior session's, so sequential + * utterances never overlap. */ #ensureSession(): TtsStreamHandle { if (this.#handle) return this.#handle; @@ -133,6 +162,28 @@ export class Vocalizer { return handle; } + /** + * (Re)arm the stall timer: if no delta arrives for {@link IDLE_FLUSH_MS}, + * speak the buffered partial sentence instead of holding it through a tool + * call or thinking block. No-op by the time it fires if the utterance moved on. + */ + #armIdleFlush(speakable: SpeakableStream): void { + this.#clearIdleTimer(); + const timer = setTimeout(() => { + this.#idleTimer = null; + if (this.#speakable !== speakable) return; + this.#pushSegments(speakable.flushIdle()); + }, IDLE_FLUSH_MS); + timer.unref?.(); + this.#idleTimer = timer; + } + + #clearIdleTimer(): void { + if (this.#idleTimer === null) return; + clearTimeout(this.#idleTimer); + this.#idleTimer = null; + } + /** Feed each synthesized sentence into the player in arrival order; abort stops it. */ async #play(handle: TtsStreamHandle, player: VocalizerPlayer, signal: AbortSignal): Promise { let started = false; diff --git a/packages/coding-agent/test/tts/speakable.test.ts b/packages/coding-agent/test/tts/speakable.test.ts new file mode 100644 index 000000000..26d3dcbf1 --- /dev/null +++ b/packages/coding-agent/test/tts/speakable.test.ts @@ -0,0 +1,133 @@ +import { describe, expect, it } from "bun:test"; +import { SpeakableStream } from "@oh-my-pi/pi-coding-agent/tts/speakable"; + +/** Push each delta in order, then flush; returns per-push segments plus the flush tail. */ +function speak(...deltas: string[]): { pushed: string[][]; all: string[] } { + const stream = new SpeakableStream(); + const pushed = deltas.map(delta => stream.push(delta)); + const flushed = stream.flush(); + return { pushed, all: [...pushed.flat(), ...flushed] }; +} + +describe("SpeakableStream code fences", () => { + it("silences a fenced block while speaking the prose around it", () => { + const { all } = speak("Here is the code:\n```ts\nconst x = 1;\nconsole.log(x);\n```\nAnd that is all of it now.\n"); + expect(all).toEqual(["Here is the code:", "And that is all of it now."]); + }); + + it("silences a block whose opening and closing fences are split across deltas", () => { + const { pushed, all } = speak("Intro line here.\n``", "`ts\nconst hidden = 42;\n``", "`\nOutro after code block.\n"); + expect(pushed[0]).toEqual(["Intro line here."]); + expect(pushed[1]).toEqual([]); + expect(all).toEqual(["Intro line here.", "Outro after code block."]); + }); +}); + +describe("SpeakableStream tables", () => { + it("silences | rows while speaking surrounding prose", () => { + const { all } = speak("Results below.\n| a | b |\n| --- | --- |\n| 1 | 2 |\nDone with the table now.\n"); + expect(all).toEqual(["Results below.", "Done with the table now."]); + }); +}); + +describe("SpeakableStream links and URLs", () => { + it("speaks only the label of a markdown link split across deltas, with no early mid-link cut", () => { + const { pushed, all } = speak("See [the do", "cs](https://exam", "ple.com/path) for details.\n"); + expect(pushed[0]).toEqual([]); + expect(pushed[1]).toEqual([]); + expect(all).toEqual(["See the docs for details."]); + }); + + it.each([ + ["bare https URL speaks only the host", "Repo lives at https://github.com/foo/bar?x for now.\n", "Repo lives at github.com for now."], + ["www URL speaks the host without the www prefix or path", "Visit www.example.com/path when you can.\n", "Visit example.com when you can."], + ])("%s", (_name, input, spoken) => { + expect(speak(input).all).toEqual([spoken]); + }); +}); + +describe("SpeakableStream inline markup and line markers", () => { + it.each([ + ["inline code speaks the identifier without ticks", "Call `parseConfig` before use today.\n", ["Call parseConfig before use today."]], + ["bold, italic, and strikethrough markers are stripped", "This **bold** and *ital* and ~~struck~~ text stays.\n", ["This bold and ital and struck text stays."]], + ["heading markers are stripped but the title speaks", "## Release Notes\nBody text follows here.\n", ["Release Notes", "Body text follows here."]], + ["bullet markers are stripped", "- item one is ready\n- item two is ready\n", ["item one is ready", "item two is ready"]], + ["numbered list markers speak as a numeric prefix", "1. First\n", ["1, First"]], + ])("%s", (_name, input, spoken) => { + expect(speak(input).all).toEqual(spoken); + }); +}); + +describe("SpeakableStream file paths", () => { + it("collapses a multi-directory path to its basename", () => { + const { all } = speak("Edit packages/coding-agent/src/tts/vocalizer.ts to fix it.\n"); + expect(all).toEqual(["Edit vocalizer.ts to fix it."]); + }); + + it("leaves two-component tokens like and/or untouched", () => { + const { all } = speak("Use and/or as needed today.\n"); + expect(all).toEqual(["Use and/or as needed today."]); + }); +}); + +describe("SpeakableStream streaming latency", () => { + it("emits a completed sentence from push() itself, without waiting for the next sentence", () => { + const stream = new SpeakableStream(); + expect(stream.push("First sentence is long enough here. ")).toEqual(["First sentence is long enough here."]); + expect(stream.flush()).toEqual([]); + }); +}); + +describe("SpeakableStream silent replies", () => { + it("yields zero segments for a reply that is only markup, whitespace, and label-less link markup", () => { + const stream = new SpeakableStream(); + expect(stream.push("---\n\n \n**\n![](https://x.com/a.png)\n```\nlet a = 1;\n```\n| a |\n")).toEqual([]); + expect(stream.flush()).toEqual([]); + }); +}); + +describe("SpeakableStream flush and flushIdle", () => { + it("flush() drains a trailing partial sentence that push() held back", () => { + const stream = new SpeakableStream(); + expect(stream.push("This is a trailing partial")).toEqual([]); + expect(stream.flush()).toEqual(["This is a trailing partial"]); + }); + + it("flushIdle() refuses a short mid-sentence fragment, which a later flush() still drains", () => { + const stream = new SpeakableStream(); + expect(stream.push("The")).toEqual([]); + expect(stream.flushIdle()).toEqual([]); + expect(stream.flush()).toEqual(["The"]); + }); + + it("flushIdle() drains a short but complete thought", () => { + const stream = new SpeakableStream(); + expect(stream.push("Done here now.")).toEqual([]); + expect(stream.flushIdle()).toEqual(["Done here now."]); + expect(stream.flush()).toEqual([]); + }); +}); + +describe("SpeakableStream segment length cap", () => { + it("force-splits an unpunctuated 1000+ char run into <=280-char segments that preserve every word", () => { + const run = Array.from({ length: 200 }, (_, i) => `word${i}`).join(" "); + expect(run.length).toBeGreaterThan(1000); + const { all } = speak(run); + expect(all.length).toBeGreaterThan(1); + for (const segment of all) expect(segment.length).toBeLessThanOrEqual(280); + expect(all.join(" ")).toBe(run); + }); +}); + +describe("SpeakableStream abbreviations", () => { + it.each([ + ['"e.g. " near the start does not end the first segment', "See e.g. the docs for more. ", "See e.g. the docs for more."], + // The abbreviation here sits past the first-segment minimum, so only the + // abbreviation guard (not the length floor) prevents a cut after "e.g. ". + ['"e.g. " past the minimum cut length still does not split the sentence', "See the docs e.g. the guide for more info. ", "See the docs e.g. the guide for more info."], + ])("%s", (_name, input, spoken) => { + const stream = new SpeakableStream(); + expect(stream.push(input)).toEqual([spoken]); + expect(stream.flush()).toEqual([]); + }); +});