feat(coding-agent/tts): redesigned streaming speech vocalization

- Added `SpeakableStream` to strip markdown noise, silence code blocks and tables, normalize links and paths, and emit sentence/clause segments.
- Reworked `Vocalizer` to segment assistant deltas in the parent process, lazily open TTS streams, idle-flush partial thoughts, and chain playback sessions.
- Added gapless streaming playback with ffmpeg/sox backends, ducking-aware pacing, fallback file playback, and immediate stop handling.
- Added IPC `sendAndFlush` support and used it in the TTS worker so audio chunks drain before blocking ONNX inference resumes.
- Added speakable-stream coverage for markdown filtering, segmentation latency, idle flushing, and forced long-segment splits.
This commit is contained in:
can1357
2026-07-02 08:30:58 +02:00
parent 95b91c7f73
commit 41cc57c238
9 changed files with 751 additions and 101 deletions
+27 -5
View File
@@ -173,14 +173,20 @@ async function runWorkerEntrypoint(arg: string | undefined): Promise<boolean> {
* worker is idle, and hard-kills the process on parent `disconnect`.
*/
async function runIpcSubprocessWorker<In, Out>(
start: (transport: { send(message: Out): void; onMessage(handler: (message: In) => void): () => void }) => void,
start: (transport: {
send(message: Out): void;
sendAndFlush(message: Out): Promise<void>;
onMessage(handler: (message: In) => void): () => void;
}) => void,
): Promise<void> {
const { promise: shuttingDown, resolve: shutdown } = Promise.withResolvers<void>();
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<In, Out>(
shutdown();
}
};
const sendAndFlush = (message: Out): Promise<void> => {
const sender = ipcSend();
if (!sender) {
shutdown();
return Promise.resolve();
}
const { promise, resolve } = Promise.withResolvers<void>();
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);
+1
View File
@@ -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";
+382
View File
@@ -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;
}
}
@@ -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;
}
}
+11 -10
View File
@@ -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) {
+14 -5
View File
@@ -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<string, unknown> }
// 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<void>;
onMessage(handler: (message: TtsWorkerInbound) => void): () => void;
}
+52 -49
View File
@@ -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<RawAudio>;
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<KokoroTtsInstance>;
};
TextSplitterStream: new () => TextSplitterStreamInstance;
}
/**
@@ -96,15 +81,18 @@ const kokoroRuntime = new MemoizedRuntime<KokoroRuntime>();
/**
* 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<void> {
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<void>();
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);
}
+69 -18
View File
@@ -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<void> {
let started = false;
@@ -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([]);
});
});