From 737c702fbb48097be9d2cfffb5d6ae6453530737 Mon Sep 17 00:00:00 2001 From: Sachin Sharma Date: Sun, 27 Sep 2026 00:59:22 +0530 Subject: [PATCH] feat(tts): stream Google audio natively and add a voices command Two gaps in the TTS surface, both reachable only from the public API. Google synthesis always waited for the whole segment. `GoogleTTSHandler` now implements the optional `TTSHandler.synthesizeStream` seam added for OpenAI, so `TTSProcessor` prefers provider-native reads and keeps the buffered path for everything else. The gate is narrow and measured, not assumed. Against the live API, streaming accepts only `Chirp3-HD`, `Chirp-HD` and `Journey` voices -- `Neural2` and `Studio` are refused outright -- and only the `PCM` and `OGG_OPUS` encodings; `MP3` and `LINEAR16` are rejected as unsupported even though `synthesizeSpeech` takes both. `PCM` delivered 41 reads with the first at 697ms and the body complete at 6131ms, `OGG_OPUS` 11 reads with the first at 443ms. The streaming request carries no SSML field, so `` input stays buffered rather than being sent as literal text. Anything outside that set answers `undefined` and is served exactly as before, which includes the default voice and the default format -- so native delivery is opt-in and nothing existing changes shape. The stream is bounded by an idle timeout rather than a total one: streaming produces audio at roughly playback speed, so a whole-call bound would fail a long but healthy segment for being long. Separately, `--tts-voice` took a provider-specific id that nothing printed. `TTSProcessor.getVoices(provider, { languageCode })` wraps the optional handler member with typed errors, and `neurolink voices` renders it as a sorted table or JSON, naming the registered set when a provider is not among it. Tests are end-to-end against dist. The discriminator for native delivery is the chunk count for ONE sentence: the buffered path can only ever emit one, so more than one is unreachable without the native read -- and each case first asserts that audio was produced at all, so the count is never vacuous. The same voice with `mp3` is the negative control and must yield exactly one. Both were confirmed to report a failure, not a skip, by breaking them on purpose: 2 failed, exit 1. Live cases skip cleanly without credentials, and treat a 401/403 as the environmental condition it is; azure-tts skipped that way on this run. Review follow-ups (CodeRabbit + Yama), both fixed test-first: the idle timer in `GoogleTTSHandler.synthesizeStream` re-armed on every server response, before yielding -- so it measured time spent suspended at its own `yield` waiting for the consumer, not time spent waiting for the server, and a consumer slower than 30s to pull the next chunk tore down a healthy duplex. It now clears on receipt, re-arms immediately on an empty read, and re-arms only after the consumer resumes it following a yield. Separately, `TTSProcessor.getVoices()` called a handler's `getVoices` before checking `isConfigured()` and let a raw provider error escape unshaped; it now enforces the same `isConfigured()` guard `synthesize()` does and routes failures through `toSynthesisError()` so `code`/`retriable`/category survive for callers. Both were reproduced RED against `GoogleTTSHandler.synthesizeStream()` and `getVoices()` directly (dist-exported public surfaces, bypassing `TTSProcessor`'s outer one-chunk lookahead buffer, which otherwise absorbs the idle-timer failure as a false negative) before the fix, and confirmed GREEN after. Three more cancellation and format-gate defects in the same streaming path, all reachable only once the native stream is actually in flight: the idle timer, the generator's own cleanup, and the out-of-band cancel hook all tore down the gRPC duplex with `.destroy()`, which only frees the local Node stream and leaves Google's server-side synthesis (and its billing) running until the RPC's own ~5-minute deadline elapses on its own. All three sites now call `.cancel()`, which reaches `ClientDuplexStreamImpl.cancel()` -> `cancelWithStatus(CANCELLED, ...)` and actually tells the server to stop -- matching the `controller.abort()` convention `OpenAITTS` already uses for the same purpose. Separately, a cancellation arriving while the generator is still parked at `await handler.getClient()` was silently dropped, because the local handle it flips a flag on is not assigned until after that await resolves; the generator now checks the flag immediately on resolution and returns without opening the duplex at all. And requesting `pcm16` against a voice or SSML combination that does not qualify for streaming fell back to the buffered encoder, whose format map had no entry for a real, supported format and threw a bare "unsupported audio format" message; it now names the two conditions (a Chirp3-HD/Chirp-HD/Journey voice, plain non-SSML text) a caller needs for the native path instead. All three were reproduced RED against `GoogleTTSHandler` directly -- a fake `gax.CancellableStream` double distinguishing `.cancel()` from `.destroy()` for the cancellation defects, and `mapFormat` for the format-gate one -- before their fix, and confirmed GREEN after. All three are now permanent regression cases in `test/continuous-test-suite-tts-unit.ts`, which passes 79 of 79 with the fixes in place. The live provider-synthesis matrix (#528) is trimmed back to the three providers `#492`, `#524` and `#528` actually scope this work to -- OpenAI, Google and Azure. ElevenLabs and Cartesia are real providers elsewhere in this codebase, but adding them to this matrix under an issue number that does not ask for them was coverage this PR did not otherwise touch. Issue #528's Google share of its own literal acceptance list -- `getVoices()` returns 220+ voices -- had no assertion anywhere in the diff; the live synthesis matrix above exercises google-ai for audio bytes only, never voice listing, and OpenAI's/Azure's shares of the same issue still need live credentials this environment does not hold. `test/continuous-test-suite-tts.ts` gained a case that calls the public `TTSProcessor.getVoices("google-ai")` against the real API and asserts the literal 220-voice floor `#528` states, which also catches an accidental language-code filter reappearing on the unfiltered path. It was confirmed to fail for the real, reported count when the threshold is set beyond the true catalog size, and to pass at the literal one. --- docs-site/static/search-index.json | 6 +- docs/api/README.md | 2 + docs/api/classes/GoogleTTSHandler.md | 39 + docs/api/classes/TTSProcessor.md | 53 + docs/api/type-aliases/CliVoicesCommandArgs.md | 35 + .../GoogleStreamingAudioEncoding.md | 17 + docs/features/tts.md | 77 +- src/cli/commands/voices.ts | 142 +++ src/cli/parser.ts | 4 + src/lib/adapters/tts/googleTTSHandler.ts | 284 ++++++ src/lib/types/cli.ts | 12 + src/lib/types/tts.ts | 11 + src/lib/utils/ttsProcessor.ts | 85 ++ test/continuous-test-suite-tts-unit.ts | 283 ++++++ test/continuous-test-suite-tts.ts | 962 +++++++++++++++++- 15 files changed, 2003 insertions(+), 9 deletions(-) create mode 100644 docs/api/type-aliases/CliVoicesCommandArgs.md create mode 100644 docs/api/type-aliases/GoogleStreamingAudioEncoding.md create mode 100644 src/cli/commands/voices.ts diff --git a/docs-site/static/search-index.json b/docs-site/static/search-index.json index d55fd9623..26e8b63a7 100644 --- a/docs-site/static/search-index.json +++ b/docs-site/static/search-index.json @@ -4489,7 +4489,7 @@ {"objectID":"4fba8a5ecbfc53e65f0e6192db71f2d3b82a547451168837085030212502f605","title":"Save to file","url":"/docs/features/tts#save-to-file","content":"neurolink generate \"Welcome to our application\" \\\n --provider google-ai \\\n --tts-voice en-US-Neural2-C \\\n --tts-output welcome.mp3\ntypescript\n\nconst neurolink = new NeuroLink();\n\nconst result = await neurolink.generate({\n input: { text: \"Hello, world!\" },\n provider: \"google-ai\",\n tts: {\n enabled: true,\n voice: \"en-US-Neural2-C\",\n format: \"mp3\",\n play: true, // Auto-play in CLI, manual in SDK\n },\n});\n\n// Access generated audio\nconsole.log(\"Audio size:\", result.audio?.size, \"bytes\");\nconsole.log(\"Audio format:\", result.audio?.format);\n`","hierarchy":{"lvl0":"Features","lvl1":"Text-to-Speech (TTS) Integration Guide","lvl2":"Save to file","lvl3":""}}, {"objectID":"6480c53c807e64a6d53b2c070905bc4edeb4d7bd9225830d84918739ba75b9d7","title":"Supported Providers","url":"/docs/features/tts#supported-providers","content":"TTS is available through the following providers:\n\n| Provider | Authentication | Voices / Models | Notes |\n| -------------- | ----------------------------------------------------------- | ------------------------------------------------------------------------------------------------ | ----------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------- |\n| google-ai | Service Account () | 50+ voices (Neural2, Wavenet, Standard) | Same auth as (TTS uses Google Cloud Text-to-Speech client) |\n| vertex | Service Account () | 50+ voices (Neural2, Wavenet, Standard) | Recommended for production |\n| openai-tts | API Key () | 6 voices: alloy, echo, fable, onyx, nova, shimmer; models: tts-1, tts-1-hd | Good default quality |\n| elevenlabs | API Key () | Multilingual voices; model: elev","hierarchy":{"lvl0":"Features","lvl1":"Text-to-Speech (TTS) Integration Guide","lvl2":"Supported Providers","lvl3":""}}, {"objectID":"0245027b08c8b3e9599fcf003e028bfa7b590365cc2049aba35ce7cee445ff61","title":"Available Voice Types","url":"/docs/features/tts#available-voice-types","content":"Google Cloud TTS offers three voice quality tiers:\n\n| Voice Type | Quality | Cost | Use Case | Example Voice |\n| ------------ | ------- | ------ | --------------------------------------- | ------------------ |\n| Neural2 | Highest | High | Natural conversations, voice assistants | |\n| Wavenet | High | Medium | Professional narration, podcasts | |\n| Standard | Good | Low | Cost optimization, bulk generation | |","hierarchy":{"lvl0":"Features","lvl1":"Text-to-Speech (TTS) Integration Guide","lvl2":"Available Voice Types","lvl3":""}}, -{"objectID":"e3dffe1fb95f6705f2d355f61f0d8a7f3053c312f576aaf17f4a82c3d03e9953","title":"Voice Discovery","url":"/docs/features/tts#voice-discovery","content":"Voice identifiers follow Google Cloud TTS naming conventions: (e.g., , ).\n\nRefer to the Google Cloud TTS voice list for all available voices.","hierarchy":{"lvl0":"Features","lvl1":"Text-to-Speech (TTS) Integration Guide","lvl2":"Voice Discovery","lvl3":""}}, +{"objectID":"e3dffe1fb95f6705f2d355f61f0d8a7f3053c312f576aaf17f4a82c3d03e9953","title":"Voice Discovery","url":"/docs/features/tts#voice-discovery","content":"takes a provider-specific identifier, and \nprints the ones a provider will accept:\n\n is required. The list is sorted by name and the count is printed\nat the end. Providers are registered only when their credentials are present,\nso an unconfigured provider is reported the same way a misspelled one is —\nwith the set that is registered named in the message, and a non-zero exit.\n\nThe same list is available from the SDK:\n\n is passed through to the provider, which decides what filtering\nit means: Google and Azure query their APIs with it, and OpenAI's fixed list\nignores it. An unregistered provider, or one whose handler does not implement\nvoice listing, raises a error.\n\nGoogle voice identifiers follow Google Cloud TTS naming conventions: (e.g., , ).\n\nRefer to the Google Cloud TTS voice list for all available voices.","hierarchy":{"lvl0":"Features","lvl1":"Text-to-Speech (TTS) Integration Guide","lvl2":"Voice Discovery","lvl3":""}}, {"objectID":"5286e5a8d6c113fc69be8b99e4bb873edd71e451caa2a01673aa1d31621f0af3","title":"Supported Languages","url":"/docs/features/tts#supported-languages","content":"English Variants:\n- United States English\n- British English\n- Australian English\n- Indian English\n\nOther Languages:\n, - Spanish (Spain, Latin America)\n, - French (France, Canada)\n- German\n- Japanese\n- Hindi\n, - Chinese (Simplified, Traditional)\n, - Portuguese (Brazil, Portugal)\n- Italian\n- Korean\n- Russian","hierarchy":{"lvl0":"Features","lvl1":"Text-to-Speech (TTS) Integration Guide","lvl2":"Supported Languages","lvl3":""}}, {"objectID":"cc003baac54d52f5303fa3e333d43096212de9505d8cb3f15a14c96f3e563e42","title":"Voice Selection Guidelines","url":"/docs/features/tts#voice-selection-guidelines","content":"For Natural Conversations:\n\nFor Professional Narration:\n\nFor Cost Optimization:","hierarchy":{"lvl0":"Features","lvl1":"Text-to-Speech (TTS) Integration Guide","lvl2":"Voice Selection Guidelines","lvl3":""}}, {"objectID":"e6d89ecb1a2aedf5e505263b9c14c576c12f543bd8ea95343f7411a4fa1b3b59","title":"TTS Synthesis Modes","url":"/docs/features/tts#tts-synthesis-modes","content":"NeuroLink supports two TTS synthesis modes:","hierarchy":{"lvl0":"Features","lvl1":"Text-to-Speech (TTS) Integration Guide","lvl2":"TTS Synthesis Modes","lvl3":""}}, @@ -4501,7 +4501,7 @@ {"objectID":"b039345b68e23c76f0eb6219aabc347404b5d45cb7ab4d3a7bfe30e198c9a273","title":"Pitch Adjustment","url":"/docs/features/tts#pitch-adjustment","content":"Adjust voice pitch (-20.0 to 20.0 semitones):\n\nCLI:","hierarchy":{"lvl0":"Features","lvl1":"Text-to-Speech (TTS) Integration Guide","lvl2":"Pitch Adjustment","lvl3":""}}, {"objectID":"068054039c4d8d505376a7faec625e0324d9a5da2e6f683a60883ba93cb65eef","title":"Volume Adjustment","url":"/docs/features/tts#volume-adjustment","content":"Control output volume (-96.0 to 16.0 dB):","hierarchy":{"lvl0":"Features","lvl1":"Text-to-Speech (TTS) Integration Guide","lvl2":"Volume Adjustment","lvl3":""}}, {"objectID":"5d6186bf5e19d14b1e69bcd5876eb97463d3aeaafd37319b597056a850ac84ce","title":"CLI Flags","url":"/docs/features/tts#cli-flags","content":"`bash\nneurolink generate \"Your text\" \\\n --provider google-ai \\\n --tts \\\n --tts-provider \\\n --tts-voice \\\n --tts-format \\\n --tts-speed \\\n --tts-pitch \\\n --tts-output \\\n --tts-use-ai-response","hierarchy":{"lvl0":"Features","lvl1":"Text-to-Speech (TTS) Integration Guide","lvl2":"CLI Flags","lvl3":""}}, -{"objectID":"f6092ac261977b5e79f9f707d6505793f21cf4ddd0031313c74bebfd3b069231","title":"--tts-use-ai-response : synthesize AI response instead of input text","url":"/docs/features/tts#--tts-use-ai-response-synthesize-ai-response-instead-of-input-text","content":"bash","hierarchy":{"lvl0":"Features","lvl1":"Text-to-Speech (TTS) Integration Guide","lvl2":"--tts-use-ai-response : synthesize AI response instead of input text","lvl3":""}}, +{"objectID":"f6092ac261977b5e79f9f707d6505793f21cf4ddd0031313c74bebfd3b069231","title":"--tts-use-ai-response : synthesize AI response instead of input text","url":"/docs/features/tts#--tts-use-ai-response-synthesize-ai-response-instead-of-input-text","content":"bash\nneurolink voices --provider [--language ] [--json]\nbash","hierarchy":{"lvl0":"Features","lvl1":"Text-to-Speech (TTS) Integration Guide","lvl2":"--tts-use-ai-response : synthesize AI response instead of input text","lvl3":""}}, {"objectID":"131a74bbe692e6aebc996fe939f726a4fe623e2b0bde16fbbf0c673ecd80a0b0","title":"Use OpenAI TTS","url":"/docs/features/tts#use-openai-tts","content":"neurolink generate \"Hello\" --tts --tts-provider openai-tts","hierarchy":{"lvl0":"Features","lvl1":"Text-to-Speech (TTS) Integration Guide","lvl2":"Use OpenAI TTS","lvl3":""}}, {"objectID":"9c0f45a9bee2138a8c40bbb602d4c84b9a6d5737d2082222439f6bea81ab5659","title":"Use ElevenLabs","url":"/docs/features/tts#use-elevenlabs","content":"neurolink generate \"Hello\" --tts --tts-provider elevenlabs","hierarchy":{"lvl0":"Features","lvl1":"Text-to-Speech (TTS) Integration Guide","lvl2":"Use ElevenLabs","lvl3":""}}, {"objectID":"8161ee48ef0c112ef3afbd0f21aae33ca4730a599cd787847d29308cf0a0e34b","title":"Use Azure TTS","url":"/docs/features/tts#use-azure-tts","content":"neurolink generate \"Hello\" --tts --tts-provider azure-tts","hierarchy":{"lvl0":"Features","lvl1":"Text-to-Speech (TTS) Integration Guide","lvl2":"Use Azure TTS","lvl3":""}}, @@ -4515,7 +4515,7 @@ {"objectID":"08724a6c8a1ff45069f41dd3bae14dd1e4ff31e8b2f4096293b910ebe24533d6","title":"Normal speed for comparison","url":"/docs/features/tts#normal-speed-for-comparison","content":"neurolink generate \"Je m'appelle Claude. Comment allez-vous?\" \\\n --provider google-ai \\\n --tts-voice fr-FR-Neural2-A \\\n --tts-speed 1.0 \\\n --tts-output french-normal.mp3\n`","hierarchy":{"lvl0":"Features","lvl1":"Text-to-Speech (TTS) Integration Guide","lvl2":"Normal speed for comparison","lvl3":""}}, {"objectID":"cb4c5d8f7134c3828c803beb43c98139a26313016dbfcc07a98d4993e4cd6323","title":"5. Multilingual Support","url":"/docs/features/tts#5-multilingual-support","content":"Generate audio in multiple languages:","hierarchy":{"lvl0":"Features","lvl1":"Text-to-Speech (TTS) Integration Guide","lvl2":"5. Multilingual Support","lvl3":""}}, {"objectID":"65b7d41fc925195acb42777a5370e28f53ec6e7e2531c1014578c382cd8a7698","title":"6. Batch Audio Generation","url":"/docs/features/tts#6-batch-audio-generation","content":"Generate multiple audio files efficiently:","hierarchy":{"lvl0":"Features","lvl1":"Text-to-Speech (TTS) Integration Guide","lvl2":"6. Batch Audio Generation","lvl3":""}}, -{"objectID":"1e68edb3c5d3c783223a84199f14f74d69c4ce175e656be58caf4b8b68dcaf05","title":"7. Streaming Text + Audio","url":"/docs/features/tts#7-streaming-text-audio","content":"Mode 2 can synthesize sentence-buffered audio while the text response is still\nstreaming. synthesizes whenever is set;\n selects input-versus-response synthesis for and\ndoes not apply to . sets the minimum number of\nbuffered characters before a completed sentence is flushed (default: 120). A\nprovider's remains a hard boundary; handlers without an override\nuse the 3,000-character default.\n\nHandlers can optionally expose provider-native audio reads for each buffered\ntext segment. NeuroLink prefers that capability when it supports the requested\noptions and otherwise keeps the existing one-buffer-per-segment synthesis path.\nOpenAI TTS currently streams response-body reads for and raw ;\n, , /, and other requested formats use buffered synthesis\nbecause native delivery has not been verified for those container formats.\nCustom and built-in handlers without the optional capability remain compatible.\n\nProvider-local chunk indexes and finality are not exposed directly. NeuroLink\nrecomputes a single global zero-based index, cumulative byte size, and exactly\none final chunk across all successful segments. Empty transport reads on the\nnative path are dropped and never reach the consumer; the buffered path is\nunchanged and forwards whatever a handler's returns, so a handler\nthat answers with a zero-byte buffer still produces an empty chunk and a\nrepeated cumulative size, exactly as it did before native streaming existed.\nEach read is forwarded as soon as the next one arrives — the single one-chunk\nhold is what guarantees the final flag — rather than waiting for the segment to\ncomplete. A handler that declines the capability for a segment, or that fails while\nNeuroLink is still working out whether the capability is there, falls back to\nbuffered synthesis for that segment rather than leaving a gap — as does a\nnative stream that completes without producing any audio. That covers every\nstep of the question, reads as well as calls, since reading a property can run\na getter or a p","hierarchy":{"lvl0":"Features","lvl1":"Text-to-Speech (TTS) Integration Guide","lvl2":"7. Streaming Text + Audio","lvl3":""}}, +{"objectID":"1e68edb3c5d3c783223a84199f14f74d69c4ce175e656be58caf4b8b68dcaf05","title":"7. Streaming Text + Audio","url":"/docs/features/tts#7-streaming-text-audio","content":"Mode 2 can synthesize sentence-buffered audio while the text response is still\nstreaming. synthesizes whenever is set;\n selects input-versus-response synthesis for and\ndoes not apply to . sets the minimum number of\nbuffered characters before a completed sentence is flushed (default: 120). A\nprovider's remains a hard boundary; handlers without an override\nuse the 3,000-character default.\n\nHandlers can optionally expose provider-native audio reads for each buffered\ntext segment. NeuroLink prefers that capability when it supports the requested\noptions and otherwise keeps the existing one-buffer-per-segment synthesis path.\nCustom and built-in handlers without the optional capability remain compatible.\n\nTwo handlers currently offer it, each only for the options whose incremental\ndelivery has been verified against the live API:\n\n| Handler | Streams natively for | Everything else |\n| ------------ | --------------------------------------------------------------------------------------- | ------------------ |\n| | , | buffered synthesis |\n| | and /, and only for a , or voice | buffered synthesis |\n\nGoogle's restriction is the streaming endpoint's own, not NeuroLink's: it\nrejects every other voice family (, , , )\nwith , and\nrejects and as unsupported encodings even though\n accepts both. SSML is also excluded — the streaming request\nhas no SSML field at all, so input stays on the buffered path rather\nthan being sent as literal text. Because the default format is and the\ndefault voice is a one, native streaming is strictly opt-in: pass\nboth a streaming-capable voice and (or ) to get it.\n\nNote that is headerless: the chunks report so a\nconsumer can wrap them, and Google's buffered path does not support at\nall, so that format only works for the streaming-capable voices above.\n\nProvider-local chunk indexes and final","hierarchy":{"lvl0":"Features","lvl1":"Text-to-Speech (TTS) Integration Guide","lvl2":"7. Streaming Text + Audio","lvl3":""}}, {"objectID":"25e56e8c2aa2164c7ccb28511b88e83f4da1cdbda7bb94d32fbde94a91faa8d0","title":"Common Issues","url":"/docs/features/tts#common-issues","content":"| Issue | Cause | Solution |\n| -------------------------------- | ------------------------ | -------------------------------------------------------------------------------------------- |\n| \"TTS client not initialized\" | Missing credentials | Set or |\n| \"Invalid voice name\" | Voice ID not found | Check the Google Cloud TTS voice list |\n| \"Text too long\" | Input exceeds 5000 bytes | Split text into smaller chunks |\n| \"Synthesis failed\" | Network/API error | Check network connection and credentials |\n| Audio doesn't play | Missing audio player | Install (macOS), (Linux), or use WAV on Windows |\n| Empty audio buffer | API returned no content | Check API quota and retry |","hierarchy":{"lvl0":"Features","lvl1":"Text-to-Speech (TTS) Integration Guide","lvl2":"Common Issues","lvl3":""}}, {"objectID":"97a359cecb0e94254dd1b2c4accf0c28d6e60a0095b46364ab683bc35f998e88","title":"Authentication Issues","url":"/docs/features/tts#authentication-issues","content":"Service Account:\n\n`bash","hierarchy":{"lvl0":"Features","lvl1":"Text-to-Speech (TTS) Integration Guide","lvl2":"Authentication Issues","lvl3":""}}, {"objectID":"71ca373c0e8643e813ab5ee309fb55ddefefa288a816c2e9880f85f36fe03613","title":"Verify credentials file exists","url":"/docs/features/tts#verify-credentials-file-exists","content":"ls -la $GOOGLEAPPLICATIONCREDENTIALS","hierarchy":{"lvl0":"Features","lvl1":"Text-to-Speech (TTS) Integration Guide","lvl2":"Verify credentials file exists","lvl3":""}}, diff --git a/docs/api/README.md b/docs/api/README.md index a828dc20a..e05f12139 100644 --- a/docs/api/README.md +++ b/docs/api/README.md @@ -835,6 +835,7 @@ console.log(result.content); - [CliServeFlatRoute](type-aliases/CliServeFlatRoute.md) - [CliToolRoutingFlags](type-aliases/CliToolRoutingFlags.md) - [CliClassifierRouterFlags](type-aliases/CliClassifierRouterFlags.md) +- [CliVoicesCommandArgs](type-aliases/CliVoicesCommandArgs.md) - [CliAgentCommandArgs](type-aliases/CliAgentCommandArgs.md) - [CliNetworkCommandArgs](type-aliases/CliNetworkCommandArgs.md) - [CliAudioPlayerCommand](type-aliases/CliAudioPlayerCommand.md) @@ -3069,6 +3070,7 @@ console.log(result.content); - [~~Gender~~](type-aliases/Gender.md) - [TTSVoice](type-aliases/TTSVoice.md) - [GoogleAudioEncoding](type-aliases/GoogleAudioEncoding.md) +- [GoogleStreamingAudioEncoding](type-aliases/GoogleStreamingAudioEncoding.md) - [TTSChunk](type-aliases/TTSChunk.md) - [CartesiaMessage](type-aliases/CartesiaMessage.md) - [LogLevel](type-aliases/LogLevel.md) diff --git a/docs/api/classes/GoogleTTSHandler.md b/docs/api/classes/GoogleTTSHandler.md index e7cb97d05..01e545dec 100644 --- a/docs/api/classes/GoogleTTSHandler.md +++ b/docs/api/classes/GoogleTTSHandler.md @@ -120,3 +120,42 @@ Audio buffer with metadata #### Implementation of `TTSHandler.synthesize` + +--- + +### synthesizeStream() + +> **synthesizeStream**(`text`, `options?`): `AsyncIterable`\<[`TTSChunk`](../type-aliases/TTSChunk.md), `any`, `any`\> \| `undefined` + +Stream one pre-validated segment's audio as Google produces it. + +Returns `undefined` — the contract's "not incrementally deliverable" +signal — unless the voice and the format are both ones the streaming +endpoint was measured to accept, and the text is not SSML. +`StreamingSynthesisInput` has no `ssml` field at all, so markup that +`synthesize()` would honour has to stay on the buffered path rather than +be sent as literal text. + +Every non-empty response is yielded as it arrives and carries `isFinal: +false`. `TTSProcessor` recomputes indexes, cumulative sizes and finality +globally across segments and discards whatever a handler reports, so +labelling the last response here would buy nothing and would cost a +one-response lookahead. + +#### Parameters + +##### text + +`string` + +##### options? + +[`TTSOptions`](../type-aliases/TTSOptions.md) = `{}` + +#### Returns + +`AsyncIterable`\<[`TTSChunk`](../type-aliases/TTSChunk.md), `any`, `any`\> \| `undefined` + +#### Implementation of + +`TTSHandler.synthesizeStream` diff --git a/docs/api/classes/TTSProcessor.md b/docs/api/classes/TTSProcessor.md index bab2e4bf6..28bd431b9 100644 --- a/docs/api/classes/TTSProcessor.md +++ b/docs/api/classes/TTSProcessor.md @@ -156,6 +156,59 @@ if (TTSProcessor.supports("google-ai")) { --- +### getVoices() + +> `static` **getVoices**(`providerName`, `options?`): `Promise`\<[`TTSVoice`](../type-aliases/TTSVoice.md)[]\> + +List the voices a registered provider offers. + +The counterpart to `synthesize()` for discovery: a caller cannot pass +`TTSOptions.voice` without first knowing what the provider will accept, +and `getVoices` is optional on `TTSHandler`, so asking the handler +directly means every caller re-implements the same two guards. Both +failures are reported as typed `TTSError`s rather than a `TypeError` on +an absent member. + +`languageCode` is passed through verbatim; each handler decides what +filtering it means. Google and Azure query their APIs with it, OpenAI's +voice list is fixed and ignores it. + +#### Parameters + +##### providerName + +`string` + +Provider identifier, resolved case-insensitively + +##### options? + +Optional language filter + +###### languageCode? + +`string` + +#### Returns + +`Promise`\<[`TTSVoice`](../type-aliases/TTSVoice.md)[]\> + +The provider's voices + +#### Throws + +TTSError if the provider is not registered or cannot list voices + +#### Example + +```typescript +const voices = await TTSProcessor.getVoices("google-ai", { + languageCode: "en-US", +}); +``` + +--- + ### synthesize() > `static` **synthesize**(`text`, `provider`, `options`): `Promise`\<[`TTSResult`](../type-aliases/TTSResult.md)\> diff --git a/docs/api/type-aliases/CliVoicesCommandArgs.md b/docs/api/type-aliases/CliVoicesCommandArgs.md new file mode 100644 index 000000000..43310f33f --- /dev/null +++ b/docs/api/type-aliases/CliVoicesCommandArgs.md @@ -0,0 +1,35 @@ +[**NeuroLink API Reference**](../README.md) + +--- + +[NeuroLink API Reference](../README.md) / CliVoicesCommandArgs + +# Type Alias: CliVoicesCommandArgs + +> **CliVoicesCommandArgs** = `object` + +`neurolink voices` arguments — TTS voice discovery. + +## Properties + +### provider + +> **provider**: `string` + +TTS provider whose voices to list (e.g. google-ai, openai-tts). + +--- + +### language? + +> `optional` **language?**: `string` + +Optional language filter passed through to the provider (e.g. en-US). + +--- + +### json? + +> `optional` **json?**: `boolean` + +Emit the raw list as JSON instead of a table. diff --git a/docs/api/type-aliases/GoogleStreamingAudioEncoding.md b/docs/api/type-aliases/GoogleStreamingAudioEncoding.md new file mode 100644 index 000000000..1a034a555 --- /dev/null +++ b/docs/api/type-aliases/GoogleStreamingAudioEncoding.md @@ -0,0 +1,17 @@ +[**NeuroLink API Reference**](../README.md) + +--- + +[NeuroLink API Reference](../README.md) / GoogleStreamingAudioEncoding + +# Type Alias: GoogleStreamingAudioEncoding + +> **GoogleStreamingAudioEncoding** = `"PCM"` \| `"OGG_OPUS"` + +Audio encodings Google's _streaming_ synthesis endpoint accepts. + +Deliberately a separate type from [GoogleAudioEncoding](GoogleAudioEncoding.md): the +batch and streaming endpoints do not accept the same set. `MP3` and +`LINEAR16` are valid for batch and rejected by streaming with +`INVALID_ARGUMENT: Unsupported audio encoding`, while `PCM` — raw +16-bit signed LE, headerless — exists only on the streaming side. diff --git a/docs/features/tts.md b/docs/features/tts.md index 031bf77b8..bd8a41868 100644 --- a/docs/features/tts.md +++ b/docs/features/tts.md @@ -144,7 +144,36 @@ Google Cloud TTS offers three voice quality tiers: ### Voice Discovery -Voice identifiers follow Google Cloud TTS naming conventions: `---` (e.g., `en-US-Neural2-C`, `en-GB-Wavenet-D`). +`--tts-voice` takes a provider-specific identifier, and `neurolink voices` +prints the ones a provider will accept: + +```bash +neurolink voices --provider openai-tts # the six OpenAI voices +neurolink voices --provider google-ai --language en-US # filter by language +neurolink voices --provider elevenlabs --json # machine-readable +``` + +`--provider` is required. The list is sorted by name and the count is printed +at the end. Providers are registered only when their credentials are present, +so an unconfigured provider is reported the same way a misspelled one is — +with the set that _is_ registered named in the message, and a non-zero exit. + +The same list is available from the SDK: + +```typescript +import { TTSProcessor } from "@juspay/neurolink"; + +const voices = await TTSProcessor.getVoices("google-ai", { + languageCode: "en-US", +}); +``` + +`languageCode` is passed through to the provider, which decides what filtering +it means: Google and Azure query their APIs with it, and OpenAI's fixed list +ignores it. An unregistered provider, or one whose handler does not implement +voice listing, raises a `TTS_PROVIDER_NOT_SUPPORTED` error. + +Google voice identifiers follow Google Cloud TTS naming conventions: `---` (e.g., `en-US-Neural2-C`, `en-GB-Wavenet-D`). Refer to the [Google Cloud TTS voice list](https://cloud.google.com/text-to-speech/docs/voices) for all available voices. @@ -463,6 +492,14 @@ neurolink generate "Your text" \ # --tts-use-ai-response : synthesize AI response instead of input text ``` +**Discovering voice ids for `--tts-voice`:** + +```bash +neurolink voices --provider [--language ] [--json] +``` + +See [Voice Discovery](#voice-discovery). + **Selecting a specific TTS provider:** ```bash @@ -656,11 +693,43 @@ use the 3,000-character default. Handlers can optionally expose provider-native audio reads for each buffered text segment. NeuroLink prefers that capability when it supports the requested options and otherwise keeps the existing one-buffer-per-segment synthesis path. -OpenAI TTS currently streams response-body reads for `mp3` and raw `pcm16`; -`wav`, `flac`, `ogg`/`opus`, and other requested formats use buffered synthesis -because native delivery has not been verified for those container formats. Custom and built-in handlers without the optional capability remain compatible. +Two handlers currently offer it, each only for the options whose incremental +delivery has been verified against the live API: + +| Handler | Streams natively for | Everything else | +| ------------ | --------------------------------------------------------------------------------------- | ------------------ | +| `openai-tts` | `mp3`, `pcm16` | buffered synthesis | +| `google-ai` | `pcm16` and `ogg`/`opus`, **and** only for a `Chirp3-HD`, `Chirp-HD` or `Journey` voice | buffered synthesis | + +Google's restriction is the streaming endpoint's own, not NeuroLink's: it +rejects every other voice family (`Neural2`, `Studio`, `Wavenet`, `Standard`) +with `only Chirp 3: HD voices are supported for streaming synthesis`, and +rejects `MP3` and `LINEAR16` as unsupported encodings even though +`synthesizeSpeech` accepts both. SSML is also excluded — the streaming request +has no SSML field at all, so `` input stays on the buffered path rather +than being sent as literal text. Because the default format is `mp3` and the +default voice is a `Neural2` one, native streaming is strictly opt-in: pass +both a streaming-capable voice and `format: "pcm16"` (or `"ogg"`) to get it. + +```typescript +// Sentence-buffered audio arrives in many reads instead of one per sentence. +const streamResult = await neurolink.stream({ + input: { text: "Summarize the quarterly report." }, + tts: { + enabled: true, + provider: "google-ai", + voice: "en-US-Chirp3-HD-Aoede", + format: "pcm16", + }, +}); +``` + +Note that `pcm16` is headerless: the chunks report `sampleRate: 24000` so a +consumer can wrap them, and Google's buffered path does not support `pcm16` at +all, so that format only works for the streaming-capable voices above. + Provider-local chunk indexes and finality are not exposed directly. NeuroLink recomputes a single global zero-based index, cumulative byte size, and exactly one final chunk across all successful segments. Empty transport reads on the diff --git a/src/cli/commands/voices.ts b/src/cli/commands/voices.ts new file mode 100644 index 000000000..ad0b6b13c --- /dev/null +++ b/src/cli/commands/voices.ts @@ -0,0 +1,142 @@ +/** + * `neurolink voices` — list the TTS voices a provider offers. + * + * `--tts-voice` on `generate` takes a provider-specific identifier + * (`en-US-Neural2-C`, `alloy`, `en-US-JennyNeural`), and nothing in the CLI + * previously printed what those identifiers are. This is the discovery half + * of that flag. + * + * @module cli/commands/voices + */ + +import type { CommandModule, Argv } from "yargs"; +import chalk from "chalk"; +import type { CliVoicesCommandArgs, TTSVoice } from "../../lib/types/index.js"; +import { TTSProcessor } from "../../lib/utils/ttsProcessor.js"; +import { logger } from "../../lib/utils/logger.js"; + +export class VoicesCommandFactory { + static createVoicesCommand(): CommandModule { + return { + command: "voices", + describe: "List the text-to-speech voices a provider offers", + builder: (yargs: Argv) => + yargs + .option("provider", { + type: "string", + demandOption: true, + description: + "TTS provider to query (e.g. google-ai, openai-tts, elevenlabs)", + }) + .option("language", { + type: "string", + description: "Only show voices for this language (e.g. en-US)", + }) + .option("json", { + type: "boolean", + default: false, + description: "Emit the voice list as JSON", + }) + .example( + "$0 voices --provider openai-tts", + "List every OpenAI TTS voice", + ) + .example( + "$0 voices --provider google-ai --language en-US", + "List Google Cloud TTS voices for US English", + ) + .example( + "$0 voices --provider elevenlabs --json", + "Emit the voice list as JSON for scripting", + ) as Argv, + handler: async (argv) => { + await VoicesCommandFactory.execute(argv); + }, + }; + } + + /** + * Registration is credential-gated — `registerDefaultTTSHandlers()` skips + * any handler whose API key is absent — so an unconfigured provider is + * indistinguishable from a misspelled one at the registry. Naming what IS + * registered turns both into the same actionable message. + */ + private static async execute(argv: CliVoicesCommandArgs): Promise { + const { ProviderRegistry } = + await import("../../lib/factories/providerRegistry.js"); + await ProviderRegistry.registerAllProviders(); + + let voices: TTSVoice[]; + try { + voices = await TTSProcessor.getVoices(argv.provider, { + ...(argv.language !== undefined ? { languageCode: argv.language } : {}), + }); + } catch (error) { + const message = error instanceof Error ? error.message : String(error); + logger.always(chalk.red(`Could not list voices: ${message}`)); + process.exitCode = 1; + return; + } + + const sorted = [...voices].sort((a, b) => + // Codepoint order, not localeCompare: collation depends on the Node ICU + // build, so the same list would sort differently on two machines. + a.name < b.name ? -1 : a.name > b.name ? 1 : 0, + ); + + if (argv.json === true) { + logger.always(JSON.stringify(sorted, null, 2)); + return; + } + + if (sorted.length === 0) { + logger.always( + chalk.yellow( + argv.language !== undefined + ? `No voices for language "${argv.language}" on provider "${argv.provider}".` + : `Provider "${argv.provider}" reported no voices.`, + ), + ); + return; + } + + logger.always(VoicesCommandFactory.formatTable(sorted)); + logger.always( + chalk.gray( + `\n${sorted.length} voice${sorted.length === 1 ? "" : "s"} for ${argv.provider}` + + (argv.language !== undefined ? ` (language: ${argv.language})` : ""), + ), + ); + } + + private static formatTable(voices: readonly TTSVoice[]): string { + const headers = ["ID", "NAME", "LANGUAGE", "GENDER", "TYPE"] as const; + const rows = voices.map((voice) => [ + voice.id, + voice.name, + voice.languageCode, + voice.gender ?? "-", + voice.type ?? "-", + ]); + + const widths = headers.map((header, column) => + rows.reduce( + (widest, row) => Math.max(widest, (row[column] ?? "").length), + header.length, + ), + ); + const line = (cells: readonly string[], paint: (s: string) => string) => + paint( + cells + .map((cell, column) => cell.padEnd(widths[column] ?? 0)) + .join(" ") + .trimEnd(), + ); + + return [ + line(headers, chalk.bold), + chalk.gray(widths.map((width) => "-".repeat(width)).join(" ")), + ...rows.map((row) => line(row, (s) => s)), + ].join("\n"); + } +} diff --git a/src/cli/parser.ts b/src/cli/parser.ts index be17bbf6e..0838dd5f3 100644 --- a/src/cli/parser.ts +++ b/src/cli/parser.ts @@ -43,6 +43,7 @@ import { AutoresearchCommandFactory } from "./commands/autoresearch.js"; import { voiceServerCommand } from "./commands/voiceServer.js"; import { DocsCommandFactory } from "./commands/docs.js"; import { UsageCommandFactory } from "./commands/usage.js"; +import { VoicesCommandFactory } from "./commands/voices.js"; // Enhanced CLI with Professional UX export function initializeCliParser() { @@ -236,6 +237,9 @@ export function initializeCliParser() { // Validate Command (alias for config validate) .command(CLICommandFactory.createValidateCommand()) + // Voices Command - TTS voice discovery for --tts-voice + .command(VoicesCommandFactory.createVoicesCommand()) + // Completion Command - Using CLICommandFactory .command(CLICommandFactory.createCompletionCommand()) diff --git a/src/lib/adapters/tts/googleTTSHandler.ts b/src/lib/adapters/tts/googleTTSHandler.ts index e20ea2ca0..c372953dd 100644 --- a/src/lib/adapters/tts/googleTTSHandler.ts +++ b/src/lib/adapters/tts/googleTTSHandler.ts @@ -11,6 +11,9 @@ import { TTSError, TTS_ERROR_CODES } from "../../utils/ttsProcessor.js"; import type { TTSGender, GoogleAudioEncoding, + GoogleStreamingAudioEncoding, + TTSAudioFormat, + TTSChunk, TTSOptions, TTSResult, TTSVoice, @@ -19,6 +22,7 @@ import type { } from "../../types/index.js"; import { ErrorCategory, ErrorSeverity } from "../../constants/enums.js"; import { logger } from "../../utils/logger.js"; +import { attachStreamCancel } from "../../utils/streamCancellation.js"; import { SpanSerializer, SpanType, @@ -354,6 +358,264 @@ export class GoogleTTSHandler implements TTSHandler { } } + /** + * Voice names Google's streaming endpoint will synthesize. + * + * Streaming is not a property of the API, it is a property of the voice: + * every other family is rejected outright with `INVALID_ARGUMENT: + * Currently, only Chirp 3: HD voices are supported for streaming + * synthesis`. Measured against the live API — `Chirp3-HD` (en-US, en-GB and + * de-DE), `Chirp-HD` and `Journey` all stream; `Neural2` and `Studio` are + * refused. The locale prefix is required so a malformed name cannot reach + * the wire and fail a segment that the buffered path would have served. + */ + private static readonly STREAMING_VOICE_PATTERN = + /^[a-z]{2,3}-[A-Z]{2}-(?:Chirp3-HD|Chirp-HD|Journey)-/; + + /** + * Sample rate requested for streamed audio, in Hz. + * + * Reported on every chunk so a consumer can interpret headerless `PCM` + * bytes. Google honours the request — 16000 was measured returning + * proportionally fewer bytes for the same text — and 24000 matches what + * `getSampleRate()` reports for the buffered path. + */ + private static readonly STREAMING_SAMPLE_RATE_HZ = 24000; + + /** + * How long a stream may go without delivering a response before it is + * treated as stalled. + * + * Deliberately an *idle* bound rather than a bound on the whole call, which + * is what `DEFAULT_API_TIMEOUT_MS` gives `synthesize()`. Streaming produces + * audio at roughly playback speed, so a total bound would fail a long but + * perfectly healthy segment purely for being long. + */ + private static readonly STREAMING_IDLE_TIMEOUT_MS = 30 * 1000; + + /** + * Map a requested audio format to a streaming audio encoding, or + * `undefined` when streaming cannot produce it. + * + * Only formats with direct wire proof of incremental delivery appear here. + * `PCM` delivered 41 reads with the first at 697ms and the body complete at + * 6131ms; `OGG_OPUS` delivered 11 reads with the first at 443ms. `MP3` and + * `LINEAR16` — both valid for `synthesizeSpeech` — are rejected by the + * streaming endpoint as unsupported encodings, so `mp3` and `wav` take the + * buffered path. `mp3` being the default format is why a caller has to opt + * in to streaming by asking for one of these two. + */ + private static mapStreamingFormat( + format: TTSAudioFormat, + ): GoogleStreamingAudioEncoding | undefined { + switch (format) { + case "pcm16": + return "PCM"; + case "ogg": + case "opus": + return "OGG_OPUS"; + default: + return undefined; + } + } + + /** + * Coerce one streaming response into its audio payload, or `undefined` when + * it carries none. + * + * The generated response type says `audioContent` is `Uint8Array | string | + * null`, and which one arrives depends on how the transport was configured + * — so both are handled rather than trusted. + */ + private static streamedAudio(response: unknown): Buffer | undefined { + if (response === null || typeof response !== "object") { + return undefined; + } + const audioContent = (response as { audioContent?: unknown }).audioContent; + if (audioContent instanceof Uint8Array) { + return audioContent.byteLength > 0 + ? Buffer.from( + audioContent.buffer, + audioContent.byteOffset, + audioContent.byteLength, + ) + : undefined; + } + if (typeof audioContent === "string" && audioContent.length > 0) { + const decoded = Buffer.from(audioContent, "base64"); + return decoded.length > 0 ? decoded : undefined; + } + return undefined; + } + + /** + * Stream one pre-validated segment's audio as Google produces it. + * + * Returns `undefined` — the contract's "not incrementally deliverable" + * signal — unless the voice and the format are both ones the streaming + * endpoint was measured to accept, and the text is not SSML. + * `StreamingSynthesisInput` has no `ssml` field at all, so markup that + * `synthesize()` would honour has to stay on the buffered path rather than + * be sent as literal text. + * + * Every non-empty response is yielded as it arrives and carries `isFinal: + * false`. `TTSProcessor` recomputes indexes, cumulative sizes and finality + * globally across segments and discards whatever a handler reports, so + * labelling the last response here would buy nothing and would cost a + * one-response lookahead. + */ + synthesizeStream( + text: string, + options: TTSOptions = {}, + ): AsyncIterable | undefined { + const voiceId = options.voice; + if ( + voiceId === undefined || + !GoogleTTSHandler.STREAMING_VOICE_PATTERN.test(voiceId) + ) { + return undefined; + } + const format = options.format ?? "mp3"; + const audioEncoding = GoogleTTSHandler.mapStreamingFormat(format); + if (audioEncoding === undefined) { + return undefined; + } + if (text.trimStart().startsWith(" + | undefined; + let cancelled = false; + let idleTimedOut = false; + + const stream = (async function* (): AsyncGenerator { + const client = await handler.getClient(); + // `getClient()` can span a dynamic import on a process's first call, a + // real gap the `attachStreamCancel` hook below can fire inside. Without + // this check, a cancellation that lands before `active` even exists + // (it is assigned two lines below) sets `cancelled` and calls + // `active?.cancel()` as a no-op, and the generator would carry on and + // open a duplex the caller already abandoned. + if (cancelled) { + return; + } + const startedAt = Date.now(); + const duplex = client.streamingSynthesize(); + active = duplex; + let idleTimer: ReturnType | undefined; + const armIdleTimer = (): void => { + clearTimeout(idleTimer); + idleTimer = setTimeout(() => { + idleTimedOut = true; + // `.cancel()`, not `.destroy()`: this object is a + // `gax.CancellableStream` wrapping a real gRPC duplex, and only + // `.cancel()` reaches `ClientDuplexStreamImpl.cancel()` -> + // `call.cancelWithStatus(CANCELLED, ...)`, which actually tells + // Google's server to stop. `.destroy()` only tears down the local + // Node stream — the server-side synthesis (and its billing) keeps + // running until the RPC's own ~5-minute deadline elapses on its + // own. + duplex.cancel(); + }, GoogleTTSHandler.STREAMING_IDLE_TIMEOUT_MS); + }; + let index = 0; + let cumulativeSize = 0; + + try { + duplex.write({ + streamingConfig: { + voice: { name: voiceId, languageCode }, + streamingAudioConfig: { + audioEncoding, + sampleRateHertz, + speakingRate: options.speed ?? 1.0, + }, + }, + }); + duplex.write({ input: { text } }); + duplex.end(); + armIdleTimer(); + + for await (const response of duplex as AsyncIterable) { + // Clear rather than re-arm here: the bound is on time spent + // waiting for the SERVER, not on time spent suspended at the + // `yield` below waiting for THIS generator's own consumer. Arming + // unconditionally on every response (the prior behaviour) started + // the clock running before that suspension, so a consumer slower + // than STREAMING_IDLE_TIMEOUT_MS to pull the next chunk destroyed + // a healthy duplex and reported its own backpressure as a network + // stall. The timer is re-armed only once real waiting resumes: + // immediately, for an empty read that loops back to the next + // `duplex` response, or after the consumer has resumed us. + clearTimeout(idleTimer); + const data = GoogleTTSHandler.streamedAudio(response); + if (data === undefined) { + armIdleTimer(); + continue; + } + cumulativeSize += data.length; + yield { + data, + format, + index: index++, + isFinal: false, + cumulativeSize, + voice: voiceId, + sampleRate: sampleRateHertz, + }; + armIdleTimer(); + } + + logger.debug( + `[GoogleTTSHandler] Streamed ${cumulativeSize} bytes in ${Date.now() - startedAt}ms`, + ); + } catch (err) { + if (cancelled) { + // The consumer stopped and the cancel hook destroyed the stream + // underneath an in-flight read. That is not a synthesis failure. + return; + } + if (err instanceof TTSError) { + throw err; + } + const message = err instanceof Error ? err.message : "Unknown error"; + throw new TTSError({ + code: TTS_ERROR_CODES.SYNTHESIS_FAILED, + message: idleTimedOut + ? `Google TTS streaming synthesis stalled for ${GoogleTTSHandler.STREAMING_IDLE_TIMEOUT_MS}ms` + : `Google TTS streaming synthesis failed: ${message}`, + category: idleTimedOut + ? ErrorCategory.NETWORK + : ErrorCategory.EXECUTION, + severity: ErrorSeverity.HIGH, + retriable: true, + context: { voice: voiceId, audioEncoding }, + originalError: err instanceof Error ? err : undefined, + }); + } finally { + clearTimeout(idleTimer); + active = undefined; + // Same reasoning as the idle timer above: `.cancel()` signals the + // server this call is done (a no-op on one that already finished + // normally), where `.destroy()` only ever tore down the local end. + duplex.cancel(); + } + })(); + + return attachStreamCancel(stream, () => { + cancelled = true; + active?.cancel(); + }); + } + /** * Extract language code from a Google Cloud voice name * @@ -395,6 +657,28 @@ export class GoogleTTSHandler implements TTSHandler { case "ogg": case "opus": return "OGG_OPUS"; + case "pcm16": + // pcm16 is a real, supported format — just not on THIS (buffered) + // endpoint. Reaching this branch means `synthesizeStream()` already + // declined the request (wrong voice family, or SSML text, both of + // which return `undefined` silently), so a caller who picked a + // perfectly valid streaming voice/format pair but tripped the SSML + // gate would otherwise see the same generic "unsupported format" as + // someone who picked a fundamentally wrong combination. Naming the + // actual requirement here, once, is cheaper than threading the + // disqualifying reason back from `synthesizeStream()`. + throw new TTSError({ + code: TTS_ERROR_CODES.INVALID_INPUT, + message: + "Google Cloud TTS's buffered synthesis does not support pcm16 output. " + + "pcm16 is only delivered via native streaming, which requires a " + + "Chirp3-HD, Chirp-HD, or Journey voice and plain (non-SSML) text — " + + "check that both conditions hold, or request a different format.", + category: ErrorCategory.VALIDATION, + severity: ErrorSeverity.MEDIUM, + retriable: false, + context: { format }, + }); default: throw new TTSError({ code: TTS_ERROR_CODES.INVALID_INPUT, diff --git a/src/lib/types/cli.ts b/src/lib/types/cli.ts index 1a5ed0424..fde2df6cd 100644 --- a/src/lib/types/cli.ts +++ b/src/lib/types/cli.ts @@ -2105,6 +2105,18 @@ export type CliClassifierRouterFlags = { classifierTimeout?: number; }; +/** + * `neurolink voices` arguments — TTS voice discovery. + */ +export type CliVoicesCommandArgs = { + /** TTS provider whose voices to list (e.g. google-ai, openai-tts). */ + provider: string; + /** Optional language filter passed through to the provider (e.g. en-US). */ + language?: string; + /** Emit the raw list as JSON instead of a table. */ + json?: boolean; +}; + /** * Agent command arguments for multi-agent orchestration */ diff --git a/src/lib/types/tts.ts b/src/lib/types/tts.ts index 6dcf2ffc3..ab712c2cc 100644 --- a/src/lib/types/tts.ts +++ b/src/lib/types/tts.ts @@ -226,6 +226,17 @@ export const VALID_TTS_QUALITIES: readonly TTSQuality[] = ["standard", "hd"]; /** Valid Google TTS audio formats */ export type GoogleAudioEncoding = "MP3" | "LINEAR16" | "OGG_OPUS"; +/** + * Audio encodings Google's *streaming* synthesis endpoint accepts. + * + * Deliberately a separate type from {@link GoogleAudioEncoding}: the + * batch and streaming endpoints do not accept the same set. `MP3` and + * `LINEAR16` are valid for batch and rejected by streaming with + * `INVALID_ARGUMENT: Unsupported audio encoding`, while `PCM` — raw + * 16-bit signed LE, headerless — exists only on the streaming side. + */ +export type GoogleStreamingAudioEncoding = "PCM" | "OGG_OPUS"; + /** * Type guard to check if an object is a TTSResult */ diff --git a/src/lib/utils/ttsProcessor.ts b/src/lib/utils/ttsProcessor.ts index 4fa233292..662bd7a51 100644 --- a/src/lib/utils/ttsProcessor.ts +++ b/src/lib/utils/ttsProcessor.ts @@ -14,6 +14,7 @@ import type { TTSOptions, TTSResult, TTSHandler, + TTSVoice, } from "../types/index.js"; import { VALID_AUDIO_FORMATS } from "../types/index.js"; import { ErrorCategory, ErrorSeverity } from "../constants/enums.js"; @@ -396,6 +397,90 @@ export class TTSProcessor { return isSupported; } + /** + * List the voices a registered provider offers. + * + * The counterpart to `synthesize()` for discovery: a caller cannot pass + * `TTSOptions.voice` without first knowing what the provider will accept, + * and `getVoices` is optional on `TTSHandler`, so asking the handler + * directly means every caller re-implements the same two guards. Both + * failures are reported as typed `TTSError`s rather than a `TypeError` on + * an absent member. + * + * `languageCode` is passed through verbatim; each handler decides what + * filtering it means. Google and Azure query their APIs with it, OpenAI's + * voice list is fixed and ignores it. + * + * @param providerName - Provider identifier, resolved case-insensitively + * @param options - Optional language filter + * @returns The provider's voices + * @throws TTSError if the provider is not registered or cannot list voices + * + * @example + * ```typescript + * const voices = await TTSProcessor.getVoices("google-ai", { + * languageCode: "en-US", + * }); + * ``` + */ + static async getVoices( + providerName: string, + options: { languageCode?: string } = {}, + ): Promise { + const handler = this.getHandler(providerName); + if (!handler) { + const registered = this.listProviders(); + throw new TTSError({ + code: TTS_ERROR_CODES.PROVIDER_NOT_SUPPORTED, + message: `TTS provider "${providerName}" is not registered. Registered providers: ${ + registered.length > 0 ? registered.join(", ") : "none" + }`, + category: ErrorCategory.VALIDATION, + severity: ErrorSeverity.MEDIUM, + retriable: false, + context: { provider: providerName }, + }); + } + + const listVoices = handler.getVoices; + if (typeof listVoices !== "function") { + throw new TTSError({ + code: TTS_ERROR_CODES.PROVIDER_NOT_SUPPORTED, + message: `TTS provider "${providerName}" does not support voice listing`, + category: ErrorCategory.VALIDATION, + severity: ErrorSeverity.MEDIUM, + retriable: false, + context: { provider: providerName }, + }); + } + + // Same two guards synthesize() enforces, so getVoices() keeps the + // `@throws TTSError` contract above regardless of what a given handler's + // own implementation does: block on configuration before the handler + // runs at all, and shape whatever it throws through the same + // toSynthesisError() synthesize() uses, so `code`/`retriable`/category + // survive for a caller keying retry logic off them. + if (!handler.isConfigured()) { + throw new TTSError({ + code: TTS_ERROR_CODES.PROVIDER_NOT_CONFIGURED, + message: `TTS provider "${providerName}" is not configured. Please set the required API keys.`, + category: ErrorCategory.CONFIGURATION, + severity: ErrorSeverity.HIGH, + retriable: false, + context: { provider: providerName }, + }); + } + + try { + return await listVoices.call(handler, options.languageCode); + } catch (error) { + // `getVoices`'s own options ({ languageCode }) share no fields with + // `TTSOptions` — synthesize()'s options shape — so there is nothing + // meaningful to forward into toSynthesisError()'s context here. + throw this.toSynthesisError(error, providerName, "", {}); + } + } + /** * Synthesize speech from text using a registered TTS provider * diff --git a/test/continuous-test-suite-tts-unit.ts b/test/continuous-test-suite-tts-unit.ts index fd4e89815..c774bf628 100644 --- a/test/continuous-test-suite-tts-unit.ts +++ b/test/continuous-test-suite-tts-unit.ts @@ -46,6 +46,7 @@ import { import { AIProviderFactory, getMetricsAggregator, + GoogleTTSHandler, NeuroLink, OpenAITTS, TTSProcessor, @@ -64,6 +65,7 @@ import type { TTSStreamChunk, } from "../dist/index.js"; import { stub, withStubs } from "./helpers/stubs.js"; +import { Duplex } from "node:stream"; // `offline: true` — this suite registers stub handlers and drives // createOfflineProvider; nothing in it touches a network. A test that never @@ -4167,6 +4169,287 @@ await test("#479: flac is a real OpenAI response_format and must not downgrade t ); }); +// --- PR#1746: GoogleTTSHandler streaming cancellation + format gate -------- +// +// Determinism exception (CLAUDE.md rule 15): a live Google streaming call +// cannot be made to race a cancellation at a chosen microtask, and cannot +// force its idle timer to fire in milliseconds instead of 30s. These cases +// reach into GoogleTTSHandler's private surface — the `client` field, the +// `getClient` method and the private static `STREAMING_IDLE_TIMEOUT_MS` — the +// same way the #479 case above reaches OpenAITTS's private `mapFormat`, and +// for the same reason: the seam is the only way to get deterministic control +// a live call cannot give. Everything under test still comes from +// `../dist/index.js`, never from `src/`, per the module-graph rule. +// +// `FakeCancellableDuplex` stands in for the real +// `gax.CancellableStream` that `TextToSpeechClient#streamingSynthesize()` +// returns. It exists to answer one question discriminatingly: did the +// production code call `.cancel()` (which reaches gRPC's +// `cancelWithStatus()` and actually tells the server to stop) or merely +// `.destroy()` (which only tears down the local Node stream and leaves the +// server-side synthesis, and its billing, running)? +class FakeCancellableDuplex extends Duplex { + cancelCalls = 0; + destroyCalls = 0; + + constructor() { + // `autoDestroy: false` matters here: Node's own stream machinery calls + // `destroy()` automatically once a stream is fully read and never + // written to again, which would tick `destroyCalls` for a reason that + // has nothing to do with the production code under test. Disabling it + // isolates the counter to destroy() calls the handler itself makes. + super({ objectMode: true, autoDestroy: false }); + } + + override _write( + _chunk: unknown, + _encoding: string, + callback: (error?: Error | null) => void, + ): void { + callback(); + } + + override _read(): void { + // Data arrives only via emitChunk/endStream below. + } + + emitChunk(audioContent: Buffer): void { + this.push({ audioContent }); + } + + endStream(): void { + this.push(null); + } + + // The real `gax.CancellableStream#cancel()` calls gRPC's + // `cancelWithStatus(Status.CANCELLED, ...)`, which ends the call gracefully + // from the client's point of view (no thrown error) while telling the + // server to actually stop. Modelled here as a graceful, deferred end so a + // consumer mid-`for await` sees a normal stream close, not an error. + cancel(): void { + this.cancelCalls++; + queueMicrotask(() => this.push(null)); + } + + override destroy(error?: Error): this { + this.destroyCalls++; + return super.destroy(error); + } +} + +await test("PR#1746 F1: cancelling an in-flight Google stream calls cancel(), not just destroy()", async () => { + const fakeDuplex = new FakeCancellableDuplex(); + const fakeClient = { streamingSynthesize: () => fakeDuplex }; + const handler = new GoogleTTSHandler("test-credentials-path"); + (handler as unknown as { client: unknown }).client = fakeClient; + + const iterable = handler.synthesizeStream("Cancel me mid-stream.", { + voice: "en-US-Chirp3-HD-Aoede", + format: "pcm16", + }); + assertNotNull( + iterable, + "a Chirp3-HD voice with plain text and pcm16 qualifies for native streaming", + ); + + const STREAM_CANCEL = Symbol.for("neurolink.streamCancel"); + const cancelHook = (iterable as unknown as Record void>)[ + STREAM_CANCEL + ]; + assert( + typeof cancelHook === "function", + "synthesizeStream registers a cancel hook via attachStreamCancel", + ); + + const iterator = iterable[Symbol.asyncIterator](); + try { + fakeDuplex.emitChunk(Buffer.from([1, 2, 3, 4])); + const first = await iterator.next(); + assert(first.done === false, "the first chunk is delivered normally"); + + // Cancel while the generator is suspended at the `yield`, exactly the + // window a real consumer disconnect lands in. + cancelHook(); + + const second = await iterator.next(); + assert( + second.done === true, + "the stream ends once cancelled instead of continuing to synthesize", + ); + + // Two call sites run on this path: the explicit cancel hook, then the + // generator's own `finally` block on the way out. A bare `> 0` would + // pass even if only one of the two were fixed, so the count is exact. + assertEqual( + fakeDuplex.cancelCalls, + 2, + "both the cancel hook and the finally block must call cancel()", + ); + assertEqual( + fakeDuplex.destroyCalls, + 0, + "neither call site may fall back to destroy(), which never reaches the server", + ); + } finally { + await iterator.return?.(undefined)?.catch?.(() => undefined); + } +}); + +await test("PR#1746 F1: an idle timeout calls cancel(), not just destroy()", async () => { + const HandlerStatics = GoogleTTSHandler as unknown as { + STREAMING_IDLE_TIMEOUT_MS: number; + }; + const originalTimeout = HandlerStatics.STREAMING_IDLE_TIMEOUT_MS; + HandlerStatics.STREAMING_IDLE_TIMEOUT_MS = 20; + + const fakeDuplex = new FakeCancellableDuplex(); + const fakeClient = { streamingSynthesize: () => fakeDuplex }; + const handler = new GoogleTTSHandler("test-credentials-path"); + (handler as unknown as { client: unknown }).client = fakeClient; + + const iterable = handler.synthesizeStream("Go quiet after one chunk.", { + voice: "en-US-Chirp3-HD-Aoede", + format: "pcm16", + }); + assertNotNull(iterable, "this voice/format/text combination streams"); + + const iterator = iterable[Symbol.asyncIterator](); + try { + fakeDuplex.emitChunk(Buffer.from([9, 9])); + const first = await iterator.next(); + assert(first.done === false, "the first chunk is delivered normally"); + + // Resuming past the `yield` re-arms the idle timer (now 20ms). No + // further data is pushed, so the timer — not a cancel hook — is what + // ends the stream. + const second = await iterator.next(); + assert(second.done === true, "an idle stream ends once the timeout fires"); + + // One call from the idle-timer callback, one from the finally block — + // an exact count so a fix that only touches the finally block (leaving + // the idle-timer callback still calling destroy()) still fails here. + assertEqual( + fakeDuplex.cancelCalls, + 2, + "the idle timeout must itself call cancel(), not just the finally block", + ); + assertEqual( + fakeDuplex.destroyCalls, + 0, + "an idle timeout may not fall back to destroy(), which leaves synthesis (and billing) running server-side", + ); + } finally { + await iterator.return?.(undefined)?.catch?.(() => undefined); + HandlerStatics.STREAMING_IDLE_TIMEOUT_MS = originalTimeout; + } +}); + +await test("PR#1746 F2: cancelling before getClient() resolves must not open the gRPC stream", async () => { + let resolveClient: ((client: unknown) => void) | undefined; + const clientPromise = new Promise((resolve) => { + resolveClient = resolve; + }); + let streamingSynthesizeCalls = 0; + const fakeClient = { + streamingSynthesize: () => { + streamingSynthesizeCalls++; + const duplex = new FakeCancellableDuplex(); + // If cancellation were NOT honoured before this point, the stream + // must still end quickly (rather than hang for the suite's + // 240s-per-case timeout) so an unfixed run fails fast and for the + // right reason. + queueMicrotask(() => { + duplex.emitChunk(Buffer.from([1])); + duplex.endStream(); + }); + return duplex; + }, + }; + + const handler = new GoogleTTSHandler("test-credentials-path"); + (handler as unknown as { getClient: () => Promise }).getClient = + () => clientPromise; + + const iterable = handler.synthesizeStream( + "Cancelled before the client resolves.", + { voice: "en-US-Chirp3-HD-Aoede", format: "pcm16" }, + ); + assertNotNull(iterable, "this voice/format/text combination streams"); + + const STREAM_CANCEL = Symbol.for("neurolink.streamCancel"); + const cancelHook = (iterable as unknown as Record void>)[ + STREAM_CANCEL + ]; + assert(typeof cancelHook === "function", "a cancel hook is registered"); + + const iterator = iterable[Symbol.asyncIterator](); + try { + const pulled = iterator.next(); + + // The generator is parked at `await handler.getClient()`. Cancel now — + // before it resolves, and before `active` is ever assigned — the exact + // race a disconnect during a process's first (cold-import) TTS call + // can hit in production. + cancelHook(); + resolveClient?.(fakeClient); + + const result = await pulled; + assert( + result.done === true, + "a stream cancelled before getClient() resolves must end immediately, not deliver a chunk", + ); + assertEqual( + streamingSynthesizeCalls, + 0, + "the gRPC call must never be opened once the caller already cancelled", + ); + } finally { + await iterator.return?.(undefined)?.catch?.(() => undefined); + } +}); + +await test("PR#1746 F3: mapFormat gives pcm16 an actionable reason instead of a bare 'unsupported format'", () => { + // Exercised off the prototype exactly like the OpenAI mapFormat case + // above — mapFormat only reads its `format` argument, never `this`. + const proto = GoogleTTSHandler.prototype as unknown as { + mapFormat: (f: string) => string; + }; + + let thrown: unknown; + try { + proto.mapFormat.call({}, "pcm16"); + } catch (err) { + thrown = err; + } + assert(thrown instanceof Error, "the buffered encoder mapping rejects pcm16"); + assertIncludes( + (thrown as Error).message, + "streaming", + "the error must explain pcm16 needs the native streaming path, not just name it unsupported", + ); + assertIncludes( + (thrown as Error).message, + "SSML", + "the error must name the SSML gate as one of the two conditions that route a request off streaming", + ); + + // A genuinely invalid format keeps the plain generic message — only + // pcm16 is special-cased, since it is the only one that is actually + // supported (via streaming) rather than simply invalid. + let genericThrown: unknown; + try { + proto.mapFormat.call({}, "not-a-real-format"); + } catch (err) { + genericThrown = err; + } + assert(genericThrown instanceof Error, "a truly invalid format still throws"); + assertIncludes( + (genericThrown as Error).message, + "Unsupported audio format", + "a truly invalid format keeps the plain unsupported-format message", + ); +}); + try { await runSuite(); } finally { diff --git a/test/continuous-test-suite-tts.ts b/test/continuous-test-suite-tts.ts index c94391755..9dc1110d4 100644 --- a/test/continuous-test-suite-tts.ts +++ b/test/continuous-test-suite-tts.ts @@ -25,8 +25,18 @@ import * as fs from "fs"; import * as os from "os"; import * as path from "path"; import { fileURLToPath } from "url"; -import type { ProcessResult, TTSMetadata } from "../dist/index.js"; -import { NeuroLink, ProviderRegistry, TTSProcessor } from "../dist/index.js"; +import type { + ProcessResult, + TTSMetadata, + TTSHandler, + TTSChunk, +} from "../dist/index.js"; +import { + NeuroLink, + ProviderRegistry, + TTSProcessor, + GoogleTTSHandler, +} from "../dist/index.js"; const __filename = fileURLToPath(import.meta.url); const __dirname = path.dirname(__filename); @@ -1783,6 +1793,912 @@ async function testCartesiaTTS(sdk: NeuroLink): Promise { } } +// ============================================================ +// TTS-013 (#492): Google native streaming synthesis +// TTS-026 (#524): `neurolink voices` +// TTS-028 (#528): live per-provider synthesis +// ============================================================ + +/** + * A voice family Google's streaming endpoint accepts. Every other family is + * refused with `only Chirp 3: HD voices are supported for streaming + * synthesis`, so this is not interchangeable with TTS_CONFIG.defaultVoice. + */ +const STREAMING_VOICE = "en-US-Chirp3-HD-Aoede"; + +/** + * Exactly one sentence, and long enough that Google takes over a second to + * render it. + * + * That is what makes the chunk count discriminating. `TTSProcessor` segments + * streamed text at sentence boundaries and serves each segment either from a + * handler's native stream or from a single buffered `synthesize()` call, so + * one sentence yields exactly ONE chunk on the buffered path however long it + * is. More than one chunk is reachable only by the native path. + */ +const ONE_SENTENCE = + "One single sentence that is long enough for the service to take a " + + "noticeable moment to render it from beginning to end"; + +/** + * A multi-sentence paragraph, long enough that Google's real synthesis and + * delivery time comfortably exceeds `GoogleTTSHandler`'s 30-second streaming + * idle window. + * + * `ONE_SENTENCE` is deliberately too short for that: its whole response + * finishes within a second or two, so by the time any idle timer could fire + * the duplex has already closed on its own, and a premature `destroy()` + * lands on a stream that was never going to deliver anything more anyway — + * a false negative for the idle-timer defect. `takeBufferedSegment` + * (ttsProcessor.ts) only splits accumulated text into multiple segments once + * the buffer reaches `streamingBufferSize`; since this whole paragraph + * arrives as one generator yield under that threshold, it is flushed as a + * single segment on end-of-input regardless of how many sentences it + * contains, so it still reaches Google as exactly one streaming synthesis + * call — one duplex, open for the entire real synthesis time. + */ +const LONG_STREAMING_TEXT = + "The morning sun rose slowly over the quiet valley, casting long golden shadows across the fields " + + "where farmers had already begun their daily work. Birds called to one another from the tall oak " + + "trees lining the winding dirt road, and a gentle breeze carried the scent of fresh hay and blooming " + + "wildflowers through the open countryside. In the distance, a small river wound its way past mossy " + + "boulders and fallen logs, its steady current reflecting the pale blue sky above. Travelers walking " + + "along the ridge paused to admire the view, watching clouds drift lazily overhead while sheep grazed " + + "peacefully in the meadow below. By midday, the village market would fill with the sounds of vendors " + + "calling out prices, children laughing as they chased each other between stalls, and the warm smell " + + "of fresh bread drifting from the corner bakery, making the whole square feel alive with quiet, " + + "ordinary happiness that seemed to stretch on comfortably into the long, unhurried afternoon. Later, " + + "as the light began to soften and shift toward amber, shopkeepers would slowly close their wooden " + + "shutters, exchanging quiet farewells with neighbors who lingered near the fountain at the center of " + + "the square. Somewhere beyond the rooftops, a church bell rang out a slow, familiar melody, marking " + + "the hour for anyone who cared to listen. Even the old stray cats seemed to know the rhythm of the " + + "town, stretching lazily on sun-warmed steps before wandering off toward the bakery for scraps left " + + "at the back door. As evening settled in, lanterns were lit one by one along the narrow cobblestone " + + "streets, their warm glow reflecting softly in shop windows still dusted with the last light of day, " + + "while somewhere further down the lane a fiddler began to play a slow, wandering tune that drifted " + + "gently through open doorways and out into the cooling evening air."; + +async function* oneSegment(text: string): AsyncGenerator { + yield text; +} + +type CollectedChunk = { + index: number; + isFinal: boolean; + bytes: number; + cumulativeSize?: number; + format: string; + sampleRate?: number; +}; + +/** + * Drain `TTSProcessor.synthesizeStream()` for one sentence. + * + * `streamingBufferSize` is set above the sentence length so the processor + * cannot hard-split it: any chunk count above one then comes from the + * handler, not from segmentation. + */ +async function collectGoogleChunks( + format: "pcm16" | "ogg" | "mp3", + voice: string, +): Promise { + await ProviderRegistry.registerAllProviders(); + const collected: CollectedChunk[] = []; + for await (const chunk of TTSProcessor.synthesizeStream( + oneSegment(ONE_SENTENCE), + "google-ai", + { voice, format, streamingBufferSize: ONE_SENTENCE.length + 50 }, + )) { + collected.push({ + index: chunk.index, + isFinal: chunk.isFinal, + bytes: chunk.data.length, + cumulativeSize: chunk.cumulativeSize, + format: chunk.format, + sampleRate: chunk.sampleRate, + }); + } + return collected; +} + +/** + * Global chunk invariants `TTSProcessor` guarantees whichever path served a + * segment. Returns a discrepancy description, or undefined when they all hold. + * + * Deliberately names the position of a mismatch rather than printing the + * chunk: a message carrying provider-ish text is reclassified as a skip by + * the harness, which would turn a real failure green. + */ +function describeChunkInvariantBreak( + chunks: readonly CollectedChunk[], +): string | undefined { + const misindexed = chunks.findIndex((chunk, at) => chunk.index !== at); + if (misindexed !== -1) { + return `chunk indexes are not 0..n-1 — first mismatch at position ${misindexed}`; + } + const finals = chunks.filter((chunk) => chunk.isFinal).length; + if (finals !== 1) { + return `expected exactly one final chunk, counted ${finals}`; + } + if (!chunks[chunks.length - 1]?.isFinal) { + return "the final chunk is not the last one delivered"; + } + let running = 0; + for (const [at, chunk] of chunks.entries()) { + running += chunk.bytes; + if (chunk.cumulativeSize !== running) { + return `cumulativeSize does not track the running byte total — first divergence at position ${at}`; + } + } + return undefined; +} + +/** + * TTS-013 (#492): Google streams natively for a streaming-capable voice. + * + * The precondition below matters: "more than one chunk" is only evidence of + * native delivery once audio has actually been produced. An empty stream + * would otherwise read as a failure of the native path when in fact nothing + * ran at all. + */ +async function testGoogleNativeStreaming(): Promise { + const name = "TTS - Google native streaming (#492)"; + logTest(name, "TESTING"); + + if (isTTSCredentialsMissing()) { + logTest(name, "SKIP", "GOOGLE_APPLICATION_CREDENTIALS not set"); + return null; + } + + try { + const chunks = await collectGoogleChunks("pcm16", STREAMING_VOICE); + + // Precondition: synthesis actually happened. + const audible = chunks.filter((chunk) => chunk.bytes > 0).length; + if (audible === 0) { + logTest( + name, + "FAIL", + "no audio was produced at all — the chunk-count assertion below would be vacuous", + ); + return false; + } + + if (chunks.length < 2) { + logTest( + name, + "FAIL", + `one sentence produced ${chunks.length} chunk(s); the native path must deliver more than one`, + ); + return false; + } + + const broken = describeChunkInvariantBreak(chunks); + if (broken) { + logTest(name, "FAIL", broken); + return false; + } + + if (chunks[0]?.format !== "pcm16") { + logTest(name, "FAIL", "streamed chunks are not labelled pcm16"); + return false; + } + if (chunks[0]?.sampleRate !== 24000) { + logTest( + name, + "FAIL", + "streamed chunks do not report the 24kHz sample rate headerless PCM needs", + ); + return false; + } + + logTest( + name, + "PASS", + `${chunks.length} native chunks, ${chunks[chunks.length - 1]?.cumulativeSize} bytes`, + ); + return true; + } catch (err) { + const msg = err instanceof Error ? err.message : String(err); + if (isExpectedProviderError(msg)) { + logTest(name, "SKIP", `upstream error: ${msg.slice(0, 120)}`); + return null; + } + logTest(name, "FAIL", msg); + return false; + } +} + +/** + * TTS-013 (#492), negative control: the format gate keeps mp3 on the buffered + * path. + * + * Google's streaming endpoint rejects MP3 outright, so the handler must + * answer `undefined` for it rather than reach the wire and fail a segment the + * buffered path can serve. Same voice as the case above — only the format + * differs — so a chunk count of exactly one isolates the gate. + */ +async function testGoogleStreamingFormatGate(): Promise { + const name = "TTS - Google streaming format gate (#492)"; + logTest(name, "TESTING"); + + if (isTTSCredentialsMissing()) { + logTest(name, "SKIP", "GOOGLE_APPLICATION_CREDENTIALS not set"); + return null; + } + + try { + const chunks = await collectGoogleChunks("mp3", STREAMING_VOICE); + + // Precondition: the buffered path ran and produced audio. Without this, + // zero chunks would satisfy "not more than one" for the wrong reason. + const audible = chunks.filter((chunk) => chunk.bytes > 0).length; + if (audible === 0) { + logTest( + name, + "FAIL", + "the buffered path produced no audio — nothing was actually synthesized", + ); + return false; + } + + if (chunks.length !== 1) { + logTest( + name, + "FAIL", + `mp3 must take the buffered path and yield one chunk per sentence; got ${chunks.length}`, + ); + return false; + } + + const broken = describeChunkInvariantBreak(chunks); + if (broken) { + logTest(name, "FAIL", broken); + return false; + } + + logTest( + name, + "PASS", + `mp3 fell back to buffered, ${chunks[0]?.bytes} bytes`, + ); + return true; + } catch (err) { + const msg = err instanceof Error ? err.message : String(err); + if (isExpectedProviderError(msg)) { + logTest(name, "SKIP", `upstream error: ${msg.slice(0, 120)}`); + return null; + } + logTest(name, "FAIL", msg); + return false; + } +} + +/** + * PR #1746 review (CodeRabbit r4090036240, MAJOR): the streaming idle timer + * must bound time spent waiting for the SERVER, not time this generator sits + * suspended at its own `yield` waiting for ITS OWN CONSUMER. The prior code + * re-armed the timer unconditionally on every response, before yielding — so + * the clock started running before that suspension, and a consumer slower + * than the idle window to pull the next chunk destroyed a healthy duplex and + * reported its own backpressure as a network stall. + * + * Real-time by necessity, not a unit-level trick: the defect is a race + * against a live 30-second wall-clock timer inside `GoogleTTSHandler`, and + * Rule 15 forbids reaching into that internal generator state to simulate + * it. Deliberately drives `GoogleTTSHandler.synthesizeStream()` directly + * (exported from `../dist/index.js`, a legitimate public surface) rather + * than through `TTSProcessor.synthesizeStream()`: the processor's outer loop + * keeps a one-chunk lookahead buffer, so a chunk it yields was already + * fetched from the handler before this test's own pause began, and a single + * post-pause pull would keep resolving from that stale, already-safe buffer + * regardless of whether the timer fired — silently absorbing exactly the + * failure this test exists to catch. Driving the handler removes that + * buffering: each `.next()` call maps 1:1 onto one real response from the + * duplex, so a pull made after the pause can only succeed if the connection + * the timer was supposed to be guarding is still alive. + * + * Pulls one chunk, then pauses this test's own consumption — never Google's + * server, which is still genuinely mid-flight on `LONG_STREAMING_TEXT` at + * that point — for longer than the idle window before pulling the next + * chunk. Accepts either observable symptom of the bug (a thrown "stalled" + * error, or the stream silently ending early) so the assertion does not + * depend on exactly how Node/gRPC surfaces a destroyed duplex. It stops as + * soon as a post-pause chunk is confirmed, releasing the stream via + * `iterator.return()` instead of draining the rest of a multi-minute + * response for no additional evidence. + */ +async function testGoogleStreamingSurvivesSlowConsumer(): Promise< + boolean | null +> { + const name = "TTS - Google streaming survives a slow consumer (PR#1746)"; + logTest(name, "TESTING"); + + if (isTTSCredentialsMissing()) { + logTest(name, "SKIP", "GOOGLE_APPLICATION_CREDENTIALS not set"); + return null; + } + + const handler = new GoogleTTSHandler(); + if (!handler.isConfigured()) { + logTest(name, "SKIP", "GoogleTTSHandler reports not configured"); + return null; + } + + let iterator: AsyncIterator | undefined; + try { + const iterable = handler.synthesizeStream(LONG_STREAMING_TEXT, { + voice: STREAMING_VOICE, + format: "pcm16", + }); + if (iterable === undefined) { + logTest( + name, + "FAIL", + "GoogleTTSHandler declined to stream this voice/format — the slow-consumer assertion below would be vacuous", + ); + return false; + } + iterator = iterable[Symbol.asyncIterator](); + + const first = await iterator.next(); + if (first.done || (first.value?.data.length ?? 0) === 0) { + logTest( + name, + "FAIL", + "no audio arrived on the first pull — the slow-consumer assertion below would be vacuous", + ); + return false; + } + + // Longer than GoogleTTSHandler's 30s idle window. LONG_STREAMING_TEXT is + // sized so Google's server is still genuinely streaming this response + // when the pause ends — this is this test's own consumption stalling, + // not the server finishing early. + await new Promise((resolve) => setTimeout(resolve, 31_000)); + + let second: IteratorResult | undefined; + let misattributed = false; + try { + second = await iterator.next(); + } catch { + misattributed = true; + } + + if (misattributed) { + logTest( + name, + "FAIL", + "a slow consumer pull was reported as a synthesis failure instead of delivering the still-in-flight response", + ); + return false; + } + + if (second === undefined || second.done) { + logTest( + name, + "FAIL", + "the stream ended after the paused pull instead of continuing to deliver the response still being produced", + ); + return false; + } + + if ((second.value?.data.length ?? 0) === 0) { + logTest( + name, + "FAIL", + "the chunk delivered after the paused pull carried no audio", + ); + return false; + } + + logTest( + name, + "PASS", + "delivered a further chunk after a 31s consumer pause without the connection being torn down", + ); + return true; + } catch (err) { + const msg = err instanceof Error ? err.message : String(err); + if (isExpectedProviderError(msg)) { + logTest(name, "SKIP", `upstream error: ${msg.slice(0, 120)}`); + return null; + } + logTest(name, "FAIL", msg); + return false; + } finally { + await iterator?.return?.(undefined)?.catch?.(() => undefined); + } +} + +/** + * TTS-026 (#524): `TTSProcessor.getVoices()` reports an unregistered provider + * as a typed error rather than a TypeError on an absent member. + * + * Runs without any credentials: with none set, nothing registers and every + * provider name is unregistered, which is exactly the case under test. + */ +async function testGetVoicesUnregistered(): Promise { + const name = "TTS - getVoices rejects an unregistered provider (#524)"; + logTest(name, "TESTING"); + + try { + await ProviderRegistry.registerAllProviders(); + let raised: unknown; + try { + await TTSProcessor.getVoices("definitely-not-a-tts-provider"); + } catch (err) { + raised = err; + } + + if (raised === undefined) { + logTest(name, "FAIL", "listing voices for an unknown provider resolved"); + return false; + } + const code = (raised as { code?: unknown }).code; + if (code !== "TTS_PROVIDER_NOT_SUPPORTED") { + logTest(name, "FAIL", "the raised error does not carry the typed code"); + return false; + } + + logTest(name, "PASS", "typed TTS_PROVIDER_NOT_SUPPORTED"); + return true; + } catch (err) { + logTest(name, "FAIL", err instanceof Error ? err.message : String(err)); + return false; + } +} + +/** + * PR #1746 review (CodeRabbit r4090036240 / Yama): `getVoices()` must gate on + * `isConfigured()` before touching the handler, and must shape whatever the + * handler throws into a typed `TTSError` — exactly what `synthesize()` does. + * + * No live provider exercises this today: every shipped handler already + * guards its own `isConfigured()` and wraps its own errors, so the gap this + * closes is invisible against real providers. A handler that does neither is + * exactly what the public `TTSHandler` structural type exists to allow, and + * registering one through the same `TTSProcessor.registerHandler()` a real + * provider uses is still the public surface under test, not an internal + * shape — the JSDoc example on `registerHandler` shows the identical + * pattern. Runs everywhere, no credentials, no network. + */ +async function testGetVoicesGatesConfigAndShapesErrors(): Promise< + boolean | null +> { + const name = "TTS - getVoices gates isConfigured, shapes errors (PR#1746)"; + logTest(name, "TESTING"); + + try { + const { TTSError } = await import("../dist/index.js"); + + // Case 1: an unconfigured handler must be rejected before its own + // getVoices() ever runs. + let handlerGetVoicesCalled = false; + const unconfiguredProvider = `test-unconfigured-tts-${Date.now()}`; + const unconfiguredHandler: TTSHandler = { + synthesize: async () => { + throw new Error("not used by this test"); + }, + isConfigured: () => false, + getVoices: async () => { + handlerGetVoicesCalled = true; + return []; + }, + }; + TTSProcessor.registerHandler(unconfiguredProvider, unconfiguredHandler); + + let unconfiguredRaised: unknown; + try { + await TTSProcessor.getVoices(unconfiguredProvider); + } catch (err) { + unconfiguredRaised = err; + } + + if (handlerGetVoicesCalled) { + logTest( + name, + "FAIL", + "getVoices() called the handler's own getVoices() before checking isConfigured()", + ); + return false; + } + if (!(unconfiguredRaised instanceof TTSError)) { + logTest( + name, + "FAIL", + "an unconfigured handler did not raise a typed TTSError", + ); + return false; + } + const unconfiguredCode = (unconfiguredRaised as { code?: unknown }).code; + if (unconfiguredCode !== "TTS_PROVIDER_NOT_CONFIGURED") { + logTest( + name, + "FAIL", + "an unconfigured handler did not raise TTS_PROVIDER_NOT_CONFIGURED", + ); + return false; + } + + // Case 2: a configured handler whose getVoices() throws a raw, unshaped + // error — the failure must still arrive as a typed, retriable TTSError. + const failingProvider = `test-failing-voices-tts-${Date.now()}`; + const failingHandler: TTSHandler = { + synthesize: async () => { + throw new Error("not used by this test"); + }, + isConfigured: () => true, + getVoices: async () => { + throw new Error("simulated upstream voice-list failure"); + }, + }; + TTSProcessor.registerHandler(failingProvider, failingHandler); + + let failingRaised: unknown; + try { + await TTSProcessor.getVoices(failingProvider); + } catch (err) { + failingRaised = err; + } + + if (!(failingRaised instanceof TTSError)) { + logTest( + name, + "FAIL", + "a handler's raw getVoices() failure reached the caller unshaped", + ); + return false; + } + const failingRecord = failingRaised as { + retriable?: unknown; + }; + if (failingRecord.retriable !== true) { + logTest( + name, + "FAIL", + "the shaped error lost the retriable flag synthesize() would carry", + ); + return false; + } + + logTest( + name, + "PASS", + "unconfigured handler gated before the call; raw failure shaped into TTSError", + ); + return true; + } catch (err) { + logTest(name, "FAIL", err instanceof Error ? err.message : String(err)); + return false; + } +} + +/** + * TTS-026 (#524): the CLI command exists, is reachable, and fails cleanly. + * + * Credential-free by construction — an unknown provider is unregistered + * whatever keys are present — so this half runs everywhere. + */ +async function testCLIVoicesUnknownProvider(): Promise { + const name = "CLI voices - unknown provider exits non-zero (#524)"; + logTest(name, "TESTING"); + + try { + const result = await runCommand("node", [ + "dist/cli/index.js", + "voices", + "--provider=definitely-not-a-tts-provider", + ]); + + if (result.success) { + logTest( + name, + "FAIL", + "the command reported success for an unknown provider", + ); + return false; + } + const output = `${result.stdout}${result.stderr}`; + if (!output.includes("Could not list voices")) { + logTest( + name, + "FAIL", + "the failure was not reported by the voices command", + ); + return false; + } + if (!output.includes("Registered providers")) { + logTest(name, "FAIL", "the failure does not name what is registered"); + return false; + } + + logTest(name, "PASS", `exit ${result.code}, registered set named`); + return true; + } catch (err) { + logTest(name, "FAIL", err instanceof Error ? err.message : String(err)); + return false; + } +} + +/** + * TTS-026 (#524): the live listing path — CLI, JSON mode, against OpenAI's + * fixed six voices, which are the one provider list that is a known constant. + */ +async function testCLIVoicesOpenAI(): Promise { + const name = "CLI voices - OpenAI list (#524)"; + logTest(name, "TESTING"); + + if (!process.env.OPENAI_API_KEY) { + logTest(name, "SKIP", "OPENAI_API_KEY not set"); + return null; + } + + try { + const result = await runCommand("node", [ + "dist/cli/index.js", + "voices", + "--provider=openai-tts", + "--json", + ]); + + if (!result.success) { + if (isExpectedProviderError(result.stderr)) { + logTest(name, "SKIP", result.stderr.slice(0, 120)); + return null; + } + logTest(name, "FAIL", `the command exited ${result.code}`); + return false; + } + + const start = result.stdout.indexOf("["); + const end = result.stdout.lastIndexOf("]"); + if (start === -1 || end <= start) { + logTest(name, "FAIL", "--json did not emit a JSON array"); + return false; + } + const parsed: unknown = JSON.parse(result.stdout.slice(start, end + 1)); + if (!Array.isArray(parsed)) { + logTest(name, "FAIL", "--json did not emit a JSON array"); + return false; + } + + const ids = parsed.map((voice) => (voice as { id?: unknown }).id); + const expected = ["alloy", "echo", "fable", "nova", "onyx", "shimmer"]; + const missing = expected.filter((id) => !ids.includes(id)); + if (missing.length > 0) { + logTest( + name, + "FAIL", + `${missing.length} of the six OpenAI voices are absent`, + ); + return false; + } + + const names = parsed.map((voice) => + String((voice as { name?: unknown }).name), + ); + const unsorted = names.findIndex( + (current, at) => at > 0 && (names[at - 1] ?? "") > current, + ); + if (unsorted !== -1) { + logTest( + name, + "FAIL", + `the list is not sorted by name — first break at position ${unsorted}`, + ); + return false; + } + + logTest(name, "PASS", `${parsed.length} voices, sorted`); + return true; + } catch (err) { + logTest(name, "FAIL", err instanceof Error ? err.message : String(err)); + return false; + } +} + +/** + * Issue #528's Google share of its literal per-provider acceptance list: + * "getVoices() returns 220+ voices". The live synthesis matrix below covers + * Google's synthesis half of #528 but never calls `getVoices()`, so that + * specific assertion had no test anywhere in the PR — CodeRabbit's pre-merge + * Linked Issues check named the gap. + * + * OpenAI's share of the same issue ("all 6 voices work", meaning synthesis of + * each) and Azure's ("emotional styles work") both need live synthesis against + * keys this environment does not hold; this case covers only the one share + * that is actually reachable and provable here — Google's catalog size, + * through the same `TTSProcessor.getVoices()` surface #524 added — gated on + * `GOOGLE_APPLICATION_CREDENTIALS` like every other Google case in this file. + */ +async function testGoogleVoiceCatalogSize(): Promise { + const name = "TTS - Google getVoices returns 220+ voices (#528)"; + logTest(name, "TESTING"); + + if (!process.env.GOOGLE_APPLICATION_CREDENTIALS) { + logTest(name, "SKIP", "GOOGLE_APPLICATION_CREDENTIALS not set"); + return null; + } + + try { + await ProviderRegistry.registerAllProviders(); + if (!TTSProcessor.supports("google-ai")) { + logTest(name, "SKIP", "google-ai did not register"); + return null; + } + + const voices = await TTSProcessor.getVoices("google-ai"); + + if (!Array.isArray(voices)) { + logTest(name, "FAIL", "getVoices() did not return an array"); + return false; + } + // 220 is issue #528's own literal threshold, not a number this test + // picked — a full, unfiltered catalog call is well past it (thousands), + // so this also catches an accidental language-code filter creeping into + // the unfiltered path. + if (voices.length < 220) { + logTest( + name, + "FAIL", + `only ${voices.length} voices returned — issue #528 requires 220+`, + ); + return false; + } + + const unnamed = voices.filter( + (voice) => typeof (voice as { name?: unknown }).name !== "string", + ); + if (unnamed.length > 0) { + logTest(name, "FAIL", `${unnamed.length} voice(s) are missing a name`); + return false; + } + + logTest( + name, + "PASS", + `${voices.length} voices via TTSProcessor.getVoices("google-ai")`, + ); + return true; + } catch (err) { + logTest(name, "FAIL", err instanceof Error ? err.message : String(err)); + return false; + } +} + +/** + * TTS-028 (#528): every configured TTS provider in scope synthesizes real + * audio. + * + * One case rather than one per provider so the matrix stays readable, and + * because the interesting property is coverage across providers rather than + * any single provider's behaviour. A provider whose key is absent is not + * counted; a provider whose key is present but rejected upstream is reported + * as an environment skip, matching `explainMissingAudio` elsewhere in this + * file. The whole case skips only when nothing at all was exercised. + * + * The matrix is OpenAI, Google and Azure only — the three providers #492, + * #524 and #528 actually scope this work to. ElevenLabs and Cartesia are + * real TTS providers elsewhere in this codebase, but adding them here would + * be coverage this PR did not otherwise touch, under an issue number that + * does not ask for them. + */ +async function testLiveProviderSynthesis(): Promise { + const name = "TTS - live provider synthesis matrix (#528)"; + logTest(name, "TESTING"); + + const matrix: Array<{ + provider: string; + configured: boolean; + voice?: string; + format: "mp3" | "pcm16"; + }> = [ + { + provider: "openai-tts", + configured: !!process.env.OPENAI_API_KEY, + voice: "nova", + format: "mp3", + }, + { + provider: "google-ai", + configured: !!process.env.GOOGLE_APPLICATION_CREDENTIALS, + voice: TTS_CONFIG.defaultVoice, + format: "mp3", + }, + { + provider: "azure-tts", + configured: + !!process.env.AZURE_SPEECH_KEY && !!process.env.AZURE_SPEECH_REGION, + voice: "en-US-JennyNeural", + format: "mp3", + }, + ]; + + try { + await ProviderRegistry.registerAllProviders(); + + const passed: string[] = []; + const skipped: string[] = []; + const failed: string[] = []; + + for (const entry of matrix) { + if (!entry.configured) { + continue; + } + if (!TTSProcessor.supports(entry.provider)) { + // Credentials are present but registration declined them. That is a + // configuration condition, not a synthesis defect. + skipped.push(`${entry.provider} (not registered)`); + continue; + } + try { + const result = await TTSProcessor.synthesize( + "Integration coverage for text to speech.", + entry.provider, + { + format: entry.format, + ...(entry.voice !== undefined ? { voice: entry.voice } : {}), + }, + ); + if (!Buffer.isBuffer(result.buffer) || result.buffer.length === 0) { + failed.push(`${entry.provider} (empty buffer)`); + continue; + } + if (result.size !== result.buffer.length) { + failed.push(`${entry.provider} (size disagrees with buffer)`); + continue; + } + if (typeof result.metadata?.latency !== "number") { + failed.push(`${entry.provider} (no latency recorded)`); + continue; + } + passed.push(`${entry.provider} ${result.buffer.length}B`); + } catch (err) { + const detail = err instanceof Error ? err.message : String(err); + // 401/403 is a credential condition — environmental, like a missing + // key. Anything else is a real failure of this provider's path. + if (/\b(?:401|403)\b/.test(detail) || isExpectedProviderError(detail)) { + skipped.push(`${entry.provider} (upstream credential/quota)`); + continue; + } + failed.push(`${entry.provider} (synthesis raised)`); + } + } + + if (failed.length > 0) { + logTest( + name, + "FAIL", + `${failed.length} provider(s) failed: ${failed.join(", ")}`, + ); + return false; + } + if (passed.length === 0) { + logTest( + name, + "SKIP", + skipped.length > 0 + ? `no provider was exercised — ${skipped.join(", ")}` + : "no TTS provider credentials are set", + ); + return null; + } + + logTest( + name, + "PASS", + `${passed.length} synthesized: ${passed.join(", ")}${ + skipped.length > 0 ? ` | skipped: ${skipped.join(", ")}` : "" + }`, + ); + return true; + } catch (err) { + logTest(name, "FAIL", err instanceof Error ? err.message : String(err)); + return false; + } +} + async function runAllTests(): Promise { log("\nNeuroLink Continuous Test Suite: TTS (Text-to-Speech)", "bright"); log( @@ -1875,6 +2791,48 @@ async function runAllTests(): Promise { name: "TTS - Cartesia end-to-end", fn: () => testCartesiaTTS(sharedSdk), }, + + // TTS-013 (#492) — Google native streaming, plus its format gate. + { + name: "TTS - Google native streaming (#492)", + fn: () => testGoogleNativeStreaming(), + }, + { + name: "TTS - Google streaming format gate (#492)", + fn: () => testGoogleStreamingFormatGate(), + }, + { + name: "TTS - Google streaming survives a slow consumer (PR#1746)", + fn: () => testGoogleStreamingSurvivesSlowConsumer(), + }, + + // TTS-026 (#524) — voice discovery, SDK and CLI. + { + name: "TTS - getVoices rejects an unregistered provider (#524)", + fn: () => testGetVoicesUnregistered(), + }, + { + name: "TTS - getVoices gates isConfigured, shapes errors (PR#1746)", + fn: () => testGetVoicesGatesConfigAndShapesErrors(), + }, + { + name: "CLI voices - unknown provider exits non-zero (#524)", + fn: () => testCLIVoicesUnknownProvider(), + }, + { + name: "CLI voices - OpenAI list (#524)", + fn: () => testCLIVoicesOpenAI(), + }, + { + name: "TTS - Google getVoices returns 220+ voices (#528)", + fn: () => testGoogleVoiceCatalogSize(), + }, + + // TTS-028 (#528) — live synthesis across every configured provider. + { + name: "TTS - live provider synthesis matrix (#528)", + fn: () => testLiveProviderSynthesis(), + }, // Observability spans (Test #15) { name: "TTS - Observability Spans",