@livekit/agents-plugin-fishaudio
Advanced tools
+1
-1
@@ -25,3 +25,3 @@ "use strict"; | ||
| title: "fishaudio", | ||
| version: "1.5.0", | ||
| version: "1.5.1", | ||
| package: "@livekit/agents-plugin-fishaudio" | ||
@@ -28,0 +28,0 @@ }); |
+1
-1
@@ -8,3 +8,3 @@ import { Plugin } from "@livekit/agents"; | ||
| title: "fishaudio", | ||
| version: "1.5.0", | ||
| version: "1.5.1", | ||
| package: "@livekit/agents-plugin-fishaudio" | ||
@@ -11,0 +11,0 @@ }); |
+1
-0
@@ -291,2 +291,3 @@ "use strict"; | ||
| if (!sentence) continue; | ||
| this.markStarted(); | ||
| ws.send(Buffer.from((0, import_msgpack.encode)({ event: "text", text: sentence + " " }))); | ||
@@ -293,0 +294,0 @@ ws.send(Buffer.from((0, import_msgpack.encode)({ event: "flush" }))); |
+1
-1
@@ -1,1 +0,1 @@ | ||
| {"version":3,"sources":["../src/tts.ts"],"sourcesContent":["// SPDX-FileCopyrightText: 2026 LiveKit, Inc.\n//\n// SPDX-License-Identifier: Apache-2.0\nimport {\n type APIConnectOptions,\n APIConnectionError,\n APIStatusError,\n AudioByteStream,\n Future,\n log,\n shortuuid,\n tokenize,\n tts,\n} from '@livekit/agents';\nimport type { AudioFrame } from '@livekit/rtc-node';\nimport { decode, encode } from '@msgpack/msgpack';\nimport { request } from 'node:https';\nimport { type RawData, WebSocket } from 'ws';\nimport type { LatencyMode, TTSModels } from './models.js';\n\nconst DEFAULT_MODEL: TTSModels = 's2.1-pro';\nconst DEFAULT_VOICE_ID = '933563129e564b19a115bedd57b7406a';\nconst DEFAULT_BASE_URL = 'https://api.fish.audio';\nconst NUM_CHANNELS = 1;\n// Fish Audio's default sample rate for raw PCM output.\nconst DEFAULT_SAMPLE_RATE = 24000;\n\nexport interface TTSOptions {\n apiKey?: string;\n model?: TTSModels | string;\n voiceId?: string;\n sampleRate?: number;\n baseURL?: string;\n latencyMode?: LatencyMode;\n /**\n * Upper bound on the number of characters Fish buffers before auto-synthesizing.\n * Must be between 100 and 300. With sentence-level flushing this is only hit by\n * sentences longer than `chunkLength`; otherwise audio is produced as soon as\n * each sentence is flushed. Defaults to 100.\n */\n chunkLength?: number;\n /**\n * Speaking rate multiplier for Fish `prosody.speed`. `1.0` is normal; below\n * 1.0 is slower, above is faster. Unset uses the voice's natural pace.\n */\n speed?: number;\n /**\n * Loudness adjustment in decibels for Fish `prosody.volume`. `0` is the\n * voice's natural level. Unset leaves it unchanged.\n */\n volume?: number;\n tokenizer?: tokenize.SentenceTokenizer;\n}\n\ninterface ResolvedTTSOptions {\n apiKey: string;\n model: TTSModels | string;\n voiceId?: string;\n sampleRate: number;\n baseURL: string;\n latencyMode: LatencyMode;\n chunkLength: number;\n speed?: number;\n volume?: number;\n tokenizer: tokenize.SentenceTokenizer;\n}\n\nconst DEFAULT_OPTS: Omit<ResolvedTTSOptions, 'apiKey' | 'tokenizer'> = {\n model: DEFAULT_MODEL,\n voiceId: DEFAULT_VOICE_ID,\n sampleRate: DEFAULT_SAMPLE_RATE,\n baseURL: DEFAULT_BASE_URL,\n latencyMode: 'balanced',\n chunkLength: 100,\n};\n\nconst validateChunkLength = (chunkLength: number) => {\n if (!Number.isFinite(chunkLength) || chunkLength < 100 || chunkLength > 300) {\n throw new Error('chunkLength must be between 100 and 300');\n }\n};\n\n// Fish Audio's wire format mirrors the upstream Python SDK so the server\n// doesn't fall back to its own larger defaults — in particular the docs default\n// of `chunk_length=300` produces large bursts that leave audible gaps between\n// chunk boundaries.\nconst buildTtsRequest = (opts: ResolvedTTSOptions, text: string = ''): Record<string, unknown> => {\n const prosody =\n opts.speed !== undefined || opts.volume !== undefined\n ? {\n ...(opts.speed !== undefined ? { speed: opts.speed } : {}),\n ...(opts.volume !== undefined ? { volume: opts.volume } : {}),\n }\n : null;\n\n return {\n text,\n chunk_length: opts.chunkLength,\n format: 'pcm',\n sample_rate: opts.sampleRate,\n mp3_bitrate: 64,\n opus_bitrate: 64000,\n references: [],\n // Fish Audio's wire field is `reference_id`; we expose it as `voiceId` on\n // the plugin for consistency with other TTS plugins.\n reference_id: opts.voiceId ?? null,\n normalize: true,\n latency: opts.latencyMode,\n prosody,\n top_p: 0.7,\n temperature: 0.7,\n };\n};\n\nexport class TTS extends tts.TTS {\n #opts: ResolvedTTSOptions;\n label = 'fishaudio.TTS';\n\n constructor(opts: TTSOptions = {}) {\n const apiKey = opts.apiKey ?? process.env.FISH_API_KEY;\n if (!apiKey) {\n throw new Error(\n 'Fish Audio API key is required, either as argument or set FISH_API_KEY environment variable',\n );\n }\n\n const chunkLength = opts.chunkLength ?? DEFAULT_OPTS.chunkLength;\n validateChunkLength(chunkLength);\n\n const sampleRate = opts.sampleRate ?? DEFAULT_OPTS.sampleRate;\n\n super(sampleRate, NUM_CHANNELS, { streaming: true });\n\n // min_sentence_len=1 emits each sentence as soon as the next one starts,\n // rather than batching short sentences together — minimizes TTFB on the\n // first sentence and keeps Fish synthesizing continuously.\n const tokenizer =\n opts.tokenizer ?? new tokenize.basic.SentenceTokenizer({ minSentenceLength: 1 });\n\n this.#opts = {\n apiKey,\n model: opts.model ?? DEFAULT_OPTS.model,\n voiceId: opts.voiceId ?? DEFAULT_OPTS.voiceId,\n sampleRate,\n baseURL: opts.baseURL ?? DEFAULT_OPTS.baseURL,\n latencyMode: opts.latencyMode ?? DEFAULT_OPTS.latencyMode,\n chunkLength,\n speed: opts.speed,\n volume: opts.volume,\n tokenizer,\n };\n }\n\n get model(): string {\n return this.#opts.model;\n }\n\n get provider(): string {\n return 'FishAudio';\n }\n\n updateOptions(opts: {\n model?: TTSModels | string;\n voiceId?: string;\n latencyMode?: LatencyMode;\n chunkLength?: number;\n speed?: number;\n volume?: number;\n }): void {\n if (opts.model !== undefined) this.#opts.model = opts.model;\n if (opts.voiceId !== undefined) this.#opts.voiceId = opts.voiceId;\n if (opts.latencyMode !== undefined) this.#opts.latencyMode = opts.latencyMode;\n if (opts.chunkLength !== undefined) {\n validateChunkLength(opts.chunkLength);\n this.#opts.chunkLength = opts.chunkLength;\n }\n if (opts.speed !== undefined) this.#opts.speed = opts.speed;\n if (opts.volume !== undefined) this.#opts.volume = opts.volume;\n }\n\n synthesize(\n text: string,\n connOptions?: APIConnectOptions,\n abortSignal?: AbortSignal,\n ): tts.ChunkedStream {\n return new ChunkedStream(this, text, this.#opts, connOptions, abortSignal);\n }\n\n stream(options?: { connOptions?: APIConnectOptions }): tts.SynthesizeStream {\n return new SynthesizeStream(this, this.#opts, options?.connOptions);\n }\n}\n\nexport class ChunkedStream extends tts.ChunkedStream {\n label = 'fishaudio.ChunkedStream';\n #logger = log();\n #opts: ResolvedTTSOptions;\n #text: string;\n\n constructor(\n tts: TTS,\n text: string,\n opts: ResolvedTTSOptions,\n connOptions?: APIConnectOptions,\n abortSignal?: AbortSignal,\n ) {\n super(text, tts, connOptions, abortSignal);\n this.#text = text;\n this.#opts = opts;\n }\n\n protected async run() {\n const requestId = shortuuid();\n const bstream = new AudioByteStream(this.#opts.sampleRate, NUM_CHANNELS);\n const payload = encode(buildTtsRequest(this.#opts, this.#text));\n\n const baseUrl = new URL(this.#opts.baseURL);\n const isHttps = baseUrl.protocol === 'https:';\n if (!isHttps) {\n // The plugin only supports https; fall back via Node's http module is\n // intentionally not implemented to keep the code path simple.\n throw new APIConnectionError({\n message: `Fish Audio base URL must use https (got ${this.#opts.baseURL})`,\n });\n }\n\n const doneFut = new Future<void>();\n\n const req = request(\n {\n hostname: baseUrl.hostname,\n port: parseInt(baseUrl.port) || 443,\n path: '/v1/tts',\n method: 'POST',\n headers: {\n Authorization: `Bearer ${this.#opts.apiKey}`,\n 'Content-Type': 'application/msgpack',\n model: this.#opts.model,\n 'Content-Length': payload.byteLength,\n },\n signal: this.abortSignal,\n },\n (res) => {\n const status = res.statusCode ?? -1;\n if (status < 200 || status >= 300) {\n const chunks: Buffer[] = [];\n res.on('data', (c: Buffer) => chunks.push(c));\n res.on('end', () => {\n const body = Buffer.concat(chunks).toString();\n if (!doneFut.done) {\n doneFut.reject(\n new APIStatusError({\n message: `Fish Audio TTS request failed: ${body}`,\n options: { statusCode: status, body: { raw: body } },\n }),\n );\n }\n });\n res.on('error', (err) => {\n if (err.message === 'aborted') return;\n this.#logger.error({ err }, 'Fish Audio TTS error response stream error');\n if (!doneFut.done) {\n doneFut.reject(\n new APIStatusError({\n message: `Fish Audio TTS request failed (status ${status})`,\n options: { statusCode: status },\n }),\n );\n }\n });\n return;\n }\n\n res.on('data', (chunk: Buffer) => {\n for (const frame of bstream.write(chunk)) {\n this.queue.put({\n requestId,\n segmentId: requestId,\n frame,\n final: false,\n });\n }\n });\n res.on('close', () => {\n for (const frame of bstream.flush()) {\n this.queue.put({\n requestId,\n segmentId: requestId,\n frame,\n final: false,\n });\n }\n if (!this.queue.closed) this.queue.close();\n if (!doneFut.done) doneFut.resolve();\n });\n res.on('error', (err) => {\n if (err.message === 'aborted') return;\n this.#logger.error({ err }, 'Fish Audio TTS response error');\n if (!doneFut.done) doneFut.reject(err);\n });\n },\n );\n\n req.on('error', (err) => {\n if (err.name === 'AbortError') return;\n this.#logger.error({ err }, 'Fish Audio TTS request error');\n if (!doneFut.done) doneFut.reject(err);\n });\n req.write(payload);\n req.end();\n\n try {\n await doneFut.await;\n } catch (e) {\n if (this.abortSignal.aborted) return;\n if (!this.queue.closed) this.queue.close();\n if (e instanceof APIStatusError || e instanceof APIConnectionError) {\n throw e;\n }\n throw new APIConnectionError({\n message: `Fish Audio connection failed: ${(e as Error).message ?? 'unknown error'}`,\n });\n }\n }\n}\n\nexport class SynthesizeStream extends tts.SynthesizeStream {\n label = 'fishaudio.SynthesizeStream';\n #logger = log();\n #opts: ResolvedTTSOptions;\n\n constructor(tts: TTS, opts: ResolvedTTSOptions, connOptions?: APIConnectOptions) {\n super(tts, connOptions);\n this.#opts = opts;\n }\n\n protected async run() {\n const requestId = shortuuid();\n const bstream = new AudioByteStream(this.#opts.sampleRate, NUM_CHANNELS);\n\n // Tokenize incoming text by sentence and flush after each sentence so Fish\n // synthesizes immediately at sentence boundaries instead of waiting for\n // `chunkLength` characters to accumulate. The result is much smoother\n // audio: gaps line up with sentence breaks (where pauses are natural)\n // rather than mid-clause.\n const sentStream = this.#opts.tokenizer.stream();\n\n const wsUrl = `${this.#opts.baseURL.replace(/^http/, 'ws')}/v1/tts/live`;\n let ws: WebSocket | undefined;\n try {\n ws = await connectWebSocket({\n url: wsUrl,\n headers: {\n Authorization: `Bearer ${this.#opts.apiKey}`,\n model: this.#opts.model,\n },\n timeoutMs: this.connOptions.timeoutMs,\n abortSignal: this.abortSignal,\n });\n } catch (e) {\n throw new APIConnectionError({\n message: `Fish Audio websocket connect failed: ${(e as Error).message ?? 'unknown error'}`,\n });\n }\n\n const finished = new Future<void>();\n\n const inputTask = async () => {\n try {\n for await (const data of this.input) {\n if (this.abortController.signal.aborted) break;\n if (data === SynthesizeStream.FLUSH_SENTINEL) {\n sentStream.flush();\n continue;\n }\n if (!data) continue;\n sentStream.pushText(data);\n }\n } finally {\n if (!sentStream.closed) sentStream.endInput();\n }\n };\n\n const sendTask = async () => {\n const startMsg = { event: 'start', request: buildTtsRequest(this.#opts) };\n ws!.send(Buffer.from(encode(startMsg)));\n\n for await (const ev of sentStream) {\n if (this.abortController.signal.aborted) break;\n const sentence = ev.token;\n if (!sentence) continue;\n ws!.send(Buffer.from(encode({ event: 'text', text: sentence + ' ' })));\n ws!.send(Buffer.from(encode({ event: 'flush' })));\n }\n\n if (!this.abortController.signal.aborted) {\n ws!.send(Buffer.from(encode({ event: 'stop' })));\n }\n };\n\n let lastFrame: AudioFrame | undefined;\n const sendLastFrame = (final: boolean) => {\n if (lastFrame) {\n this.queue.put({ requestId, segmentId: requestId, frame: lastFrame, final });\n lastFrame = undefined;\n }\n };\n\n const recvTask = async () => {\n // No per-receive timeout: Fish has natural inter-sentence gaps that can\n // exceed connOptions.timeoutMs when the LLM is slow.\n const onMessage = (raw: RawData) => {\n let frame: Buffer;\n if (Buffer.isBuffer(raw)) {\n frame = raw;\n } else if (Array.isArray(raw)) {\n frame = Buffer.concat(raw);\n } else {\n frame = Buffer.from(raw as ArrayBuffer);\n }\n\n let parsed: Record<string, unknown>;\n try {\n parsed = decode(frame) as Record<string, unknown>;\n } catch (err) {\n this.#logger.warn({ err }, 'Fish Audio failed to decode message');\n return;\n }\n\n const event = parsed.event as string | undefined;\n if (event === 'audio') {\n const audio = parsed.audio as Uint8Array | undefined;\n if (audio && audio.byteLength > 0) {\n for (const f of bstream.write(audio)) {\n sendLastFrame(false);\n lastFrame = f;\n }\n }\n } else if (event === 'finish') {\n const reason = parsed.reason as string | undefined;\n if (reason === 'error') {\n finished.reject(\n new APIStatusError({\n message: 'Fish Audio TTS reported an error',\n options: { body: { raw: JSON.stringify(parsed) } },\n }),\n );\n return;\n }\n for (const f of bstream.flush()) {\n sendLastFrame(false);\n lastFrame = f;\n }\n sendLastFrame(true);\n if (!this.queue.closed) {\n this.queue.put(SynthesizeStream.END_OF_STREAM);\n }\n if (!finished.done) finished.resolve();\n } else {\n this.#logger.debug({ event }, 'unknown Fish Audio event');\n }\n };\n\n const onClose = (code: number, reason: Buffer) => {\n if (!finished.done) {\n finished.reject(\n new APIStatusError({\n message: 'Fish Audio websocket connection closed unexpectedly',\n options: {\n statusCode: code || -1,\n body: { reason: reason.toString() },\n },\n }),\n );\n }\n };\n\n const onError = (err: Error) => {\n if (!finished.done) finished.reject(err);\n };\n\n ws!.on('message', onMessage);\n ws!.on('close', onClose);\n ws!.on('error', onError);\n\n try {\n await finished.await;\n } finally {\n ws!.off('message', onMessage);\n ws!.off('close', onClose);\n ws!.off('error', onError);\n }\n };\n\n try {\n await Promise.all([inputTask(), sendTask(), recvTask()]);\n } catch (e) {\n if (this.abortSignal.aborted) return;\n if (e instanceof APIStatusError || e instanceof APIConnectionError) {\n throw e;\n }\n throw new APIConnectionError({\n message: `Fish Audio websocket failed: ${(e as Error).message ?? 'unknown error'}`,\n });\n } finally {\n if (!sentStream.closed) sentStream.close();\n if (ws && ws.readyState !== WebSocket.CLOSED) {\n try {\n ws.close();\n } catch {\n // ignore\n }\n }\n }\n }\n}\n\nconst connectWebSocket = async ({\n url,\n headers,\n timeoutMs,\n abortSignal,\n}: {\n url: string;\n headers: Record<string, string>;\n timeoutMs: number;\n abortSignal: AbortSignal;\n}): Promise<WebSocket> => {\n const ws = new WebSocket(url, { headers, handshakeTimeout: timeoutMs });\n const fut = new Future<void>();\n\n let timeout: NodeJS.Timeout | undefined;\n const cleanup = () => {\n if (timeout) clearTimeout(timeout);\n ws.off('open', onOpen);\n ws.off('error', onError);\n ws.off('close', onClose);\n abortSignal.removeEventListener('abort', onAbort);\n };\n\n const onOpen = () => fut.resolve();\n const onError = (err: Error) => fut.reject(err);\n const onClose = (code: number, reason: Buffer) =>\n fut.reject(\n new Error(`websocket closed before open (code=${code}, reason=${reason.toString()})`),\n );\n const onAbort = () => fut.reject(new Error('aborted'));\n\n ws.on('open', onOpen);\n ws.on('error', onError);\n ws.on('close', onClose);\n abortSignal.addEventListener('abort', onAbort, { once: true });\n\n if (timeoutMs > 0) {\n timeout = setTimeout(() => fut.reject(new Error('connect timeout')), timeoutMs);\n }\n\n try {\n await fut.await;\n return ws;\n } catch (e) {\n try {\n ws.on('error', () => {});\n if (ws.readyState === WebSocket.CONNECTING) {\n ws.close();\n } else {\n ws.terminate();\n }\n } catch {\n // ignore\n }\n throw e;\n } finally {\n cleanup();\n }\n};\n"],"mappings":";;;;;;;;;;;;;;;;;;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAGA,oBAUO;AAEP,qBAA+B;AAC/B,wBAAwB;AACxB,gBAAwC;AAGxC,MAAM,gBAA2B;AACjC,MAAM,mBAAmB;AACzB,MAAM,mBAAmB;AACzB,MAAM,eAAe;AAErB,MAAM,sBAAsB;AA0C5B,MAAM,eAAiE;AAAA,EACrE,OAAO;AAAA,EACP,SAAS;AAAA,EACT,YAAY;AAAA,EACZ,SAAS;AAAA,EACT,aAAa;AAAA,EACb,aAAa;AACf;AAEA,MAAM,sBAAsB,CAAC,gBAAwB;AACnD,MAAI,CAAC,OAAO,SAAS,WAAW,KAAK,cAAc,OAAO,cAAc,KAAK;AAC3E,UAAM,IAAI,MAAM,yCAAyC;AAAA,EAC3D;AACF;AAMA,MAAM,kBAAkB,CAAC,MAA0B,OAAe,OAAgC;AAChG,QAAM,UACJ,KAAK,UAAU,UAAa,KAAK,WAAW,SACxC;AAAA,IACE,GAAI,KAAK,UAAU,SAAY,EAAE,OAAO,KAAK,MAAM,IAAI,CAAC;AAAA,IACxD,GAAI,KAAK,WAAW,SAAY,EAAE,QAAQ,KAAK,OAAO,IAAI,CAAC;AAAA,EAC7D,IACA;AAEN,SAAO;AAAA,IACL;AAAA,IACA,cAAc,KAAK;AAAA,IACnB,QAAQ;AAAA,IACR,aAAa,KAAK;AAAA,IAClB,aAAa;AAAA,IACb,cAAc;AAAA,IACd,YAAY,CAAC;AAAA;AAAA;AAAA,IAGb,cAAc,KAAK,WAAW;AAAA,IAC9B,WAAW;AAAA,IACX,SAAS,KAAK;AAAA,IACd;AAAA,IACA,OAAO;AAAA,IACP,aAAa;AAAA,EACf;AACF;AAEO,MAAM,YAAY,kBAAI,IAAI;AAAA,EAC/B;AAAA,EACA,QAAQ;AAAA,EAER,YAAY,OAAmB,CAAC,GAAG;AACjC,UAAM,SAAS,KAAK,UAAU,QAAQ,IAAI;AAC1C,QAAI,CAAC,QAAQ;AACX,YAAM,IAAI;AAAA,QACR;AAAA,MACF;AAAA,IACF;AAEA,UAAM,cAAc,KAAK,eAAe,aAAa;AACrD,wBAAoB,WAAW;AAE/B,UAAM,aAAa,KAAK,cAAc,aAAa;AAEnD,UAAM,YAAY,cAAc,EAAE,WAAW,KAAK,CAAC;AAKnD,UAAM,YACJ,KAAK,aAAa,IAAI,uBAAS,MAAM,kBAAkB,EAAE,mBAAmB,EAAE,CAAC;AAEjF,SAAK,QAAQ;AAAA,MACX;AAAA,MACA,OAAO,KAAK,SAAS,aAAa;AAAA,MAClC,SAAS,KAAK,WAAW,aAAa;AAAA,MACtC;AAAA,MACA,SAAS,KAAK,WAAW,aAAa;AAAA,MACtC,aAAa,KAAK,eAAe,aAAa;AAAA,MAC9C;AAAA,MACA,OAAO,KAAK;AAAA,MACZ,QAAQ,KAAK;AAAA,MACb;AAAA,IACF;AAAA,EACF;AAAA,EAEA,IAAI,QAAgB;AAClB,WAAO,KAAK,MAAM;AAAA,EACpB;AAAA,EAEA,IAAI,WAAmB;AACrB,WAAO;AAAA,EACT;AAAA,EAEA,cAAc,MAOL;AACP,QAAI,KAAK,UAAU,OAAW,MAAK,MAAM,QAAQ,KAAK;AACtD,QAAI,KAAK,YAAY,OAAW,MAAK,MAAM,UAAU,KAAK;AAC1D,QAAI,KAAK,gBAAgB,OAAW,MAAK,MAAM,cAAc,KAAK;AAClE,QAAI,KAAK,gBAAgB,QAAW;AAClC,0BAAoB,KAAK,WAAW;AACpC,WAAK,MAAM,cAAc,KAAK;AAAA,IAChC;AACA,QAAI,KAAK,UAAU,OAAW,MAAK,MAAM,QAAQ,KAAK;AACtD,QAAI,KAAK,WAAW,OAAW,MAAK,MAAM,SAAS,KAAK;AAAA,EAC1D;AAAA,EAEA,WACE,MACA,aACA,aACmB;AACnB,WAAO,IAAI,cAAc,MAAM,MAAM,KAAK,OAAO,aAAa,WAAW;AAAA,EAC3E;AAAA,EAEA,OAAO,SAAqE;AAC1E,WAAO,IAAI,iBAAiB,MAAM,KAAK,OAAO,mCAAS,WAAW;AAAA,EACpE;AACF;AAEO,MAAM,sBAAsB,kBAAI,cAAc;AAAA,EACnD,QAAQ;AAAA,EACR,cAAU,mBAAI;AAAA,EACd;AAAA,EACA;AAAA,EAEA,YACEA,MACA,MACA,MACA,aACA,aACA;AACA,UAAM,MAAMA,MAAK,aAAa,WAAW;AACzC,SAAK,QAAQ;AACb,SAAK,QAAQ;AAAA,EACf;AAAA,EAEA,MAAgB,MAAM;AACpB,UAAM,gBAAY,yBAAU;AAC5B,UAAM,UAAU,IAAI,8BAAgB,KAAK,MAAM,YAAY,YAAY;AACvE,UAAM,cAAU,uBAAO,gBAAgB,KAAK,OAAO,KAAK,KAAK,CAAC;AAE9D,UAAM,UAAU,IAAI,IAAI,KAAK,MAAM,OAAO;AAC1C,UAAM,UAAU,QAAQ,aAAa;AACrC,QAAI,CAAC,SAAS;AAGZ,YAAM,IAAI,iCAAmB;AAAA,QAC3B,SAAS,2CAA2C,KAAK,MAAM,OAAO;AAAA,MACxE,CAAC;AAAA,IACH;AAEA,UAAM,UAAU,IAAI,qBAAa;AAEjC,UAAM,UAAM;AAAA,MACV;AAAA,QACE,UAAU,QAAQ;AAAA,QAClB,MAAM,SAAS,QAAQ,IAAI,KAAK;AAAA,QAChC,MAAM;AAAA,QACN,QAAQ;AAAA,QACR,SAAS;AAAA,UACP,eAAe,UAAU,KAAK,MAAM,MAAM;AAAA,UAC1C,gBAAgB;AAAA,UAChB,OAAO,KAAK,MAAM;AAAA,UAClB,kBAAkB,QAAQ;AAAA,QAC5B;AAAA,QACA,QAAQ,KAAK;AAAA,MACf;AAAA,MACA,CAAC,QAAQ;AACP,cAAM,SAAS,IAAI,cAAc;AACjC,YAAI,SAAS,OAAO,UAAU,KAAK;AACjC,gBAAM,SAAmB,CAAC;AAC1B,cAAI,GAAG,QAAQ,CAAC,MAAc,OAAO,KAAK,CAAC,CAAC;AAC5C,cAAI,GAAG,OAAO,MAAM;AAClB,kBAAM,OAAO,OAAO,OAAO,MAAM,EAAE,SAAS;AAC5C,gBAAI,CAAC,QAAQ,MAAM;AACjB,sBAAQ;AAAA,gBACN,IAAI,6BAAe;AAAA,kBACjB,SAAS,kCAAkC,IAAI;AAAA,kBAC/C,SAAS,EAAE,YAAY,QAAQ,MAAM,EAAE,KAAK,KAAK,EAAE;AAAA,gBACrD,CAAC;AAAA,cACH;AAAA,YACF;AAAA,UACF,CAAC;AACD,cAAI,GAAG,SAAS,CAAC,QAAQ;AACvB,gBAAI,IAAI,YAAY,UAAW;AAC/B,iBAAK,QAAQ,MAAM,EAAE,IAAI,GAAG,4CAA4C;AACxE,gBAAI,CAAC,QAAQ,MAAM;AACjB,sBAAQ;AAAA,gBACN,IAAI,6BAAe;AAAA,kBACjB,SAAS,yCAAyC,MAAM;AAAA,kBACxD,SAAS,EAAE,YAAY,OAAO;AAAA,gBAChC,CAAC;AAAA,cACH;AAAA,YACF;AAAA,UACF,CAAC;AACD;AAAA,QACF;AAEA,YAAI,GAAG,QAAQ,CAAC,UAAkB;AAChC,qBAAW,SAAS,QAAQ,MAAM,KAAK,GAAG;AACxC,iBAAK,MAAM,IAAI;AAAA,cACb;AAAA,cACA,WAAW;AAAA,cACX;AAAA,cACA,OAAO;AAAA,YACT,CAAC;AAAA,UACH;AAAA,QACF,CAAC;AACD,YAAI,GAAG,SAAS,MAAM;AACpB,qBAAW,SAAS,QAAQ,MAAM,GAAG;AACnC,iBAAK,MAAM,IAAI;AAAA,cACb;AAAA,cACA,WAAW;AAAA,cACX;AAAA,cACA,OAAO;AAAA,YACT,CAAC;AAAA,UACH;AACA,cAAI,CAAC,KAAK,MAAM,OAAQ,MAAK,MAAM,MAAM;AACzC,cAAI,CAAC,QAAQ,KAAM,SAAQ,QAAQ;AAAA,QACrC,CAAC;AACD,YAAI,GAAG,SAAS,CAAC,QAAQ;AACvB,cAAI,IAAI,YAAY,UAAW;AAC/B,eAAK,QAAQ,MAAM,EAAE,IAAI,GAAG,+BAA+B;AAC3D,cAAI,CAAC,QAAQ,KAAM,SAAQ,OAAO,GAAG;AAAA,QACvC,CAAC;AAAA,MACH;AAAA,IACF;AAEA,QAAI,GAAG,SAAS,CAAC,QAAQ;AACvB,UAAI,IAAI,SAAS,aAAc;AAC/B,WAAK,QAAQ,MAAM,EAAE,IAAI,GAAG,8BAA8B;AAC1D,UAAI,CAAC,QAAQ,KAAM,SAAQ,OAAO,GAAG;AAAA,IACvC,CAAC;AACD,QAAI,MAAM,OAAO;AACjB,QAAI,IAAI;AAER,QAAI;AACF,YAAM,QAAQ;AAAA,IAChB,SAAS,GAAG;AACV,UAAI,KAAK,YAAY,QAAS;AAC9B,UAAI,CAAC,KAAK,MAAM,OAAQ,MAAK,MAAM,MAAM;AACzC,UAAI,aAAa,gCAAkB,aAAa,kCAAoB;AAClE,cAAM;AAAA,MACR;AACA,YAAM,IAAI,iCAAmB;AAAA,QAC3B,SAAS,iCAAkC,EAAY,WAAW,eAAe;AAAA,MACnF,CAAC;AAAA,IACH;AAAA,EACF;AACF;AAEO,MAAM,yBAAyB,kBAAI,iBAAiB;AAAA,EACzD,QAAQ;AAAA,EACR,cAAU,mBAAI;AAAA,EACd;AAAA,EAEA,YAAYA,MAAU,MAA0B,aAAiC;AAC/E,UAAMA,MAAK,WAAW;AACtB,SAAK,QAAQ;AAAA,EACf;AAAA,EAEA,MAAgB,MAAM;AACpB,UAAM,gBAAY,yBAAU;AAC5B,UAAM,UAAU,IAAI,8BAAgB,KAAK,MAAM,YAAY,YAAY;AAOvE,UAAM,aAAa,KAAK,MAAM,UAAU,OAAO;AAE/C,UAAM,QAAQ,GAAG,KAAK,MAAM,QAAQ,QAAQ,SAAS,IAAI,CAAC;AAC1D,QAAI;AACJ,QAAI;AACF,WAAK,MAAM,iBAAiB;AAAA,QAC1B,KAAK;AAAA,QACL,SAAS;AAAA,UACP,eAAe,UAAU,KAAK,MAAM,MAAM;AAAA,UAC1C,OAAO,KAAK,MAAM;AAAA,QACpB;AAAA,QACA,WAAW,KAAK,YAAY;AAAA,QAC5B,aAAa,KAAK;AAAA,MACpB,CAAC;AAAA,IACH,SAAS,GAAG;AACV,YAAM,IAAI,iCAAmB;AAAA,QAC3B,SAAS,wCAAyC,EAAY,WAAW,eAAe;AAAA,MAC1F,CAAC;AAAA,IACH;AAEA,UAAM,WAAW,IAAI,qBAAa;AAElC,UAAM,YAAY,YAAY;AAC5B,UAAI;AACF,yBAAiB,QAAQ,KAAK,OAAO;AACnC,cAAI,KAAK,gBAAgB,OAAO,QAAS;AACzC,cAAI,SAAS,iBAAiB,gBAAgB;AAC5C,uBAAW,MAAM;AACjB;AAAA,UACF;AACA,cAAI,CAAC,KAAM;AACX,qBAAW,SAAS,IAAI;AAAA,QAC1B;AAAA,MACF,UAAE;AACA,YAAI,CAAC,WAAW,OAAQ,YAAW,SAAS;AAAA,MAC9C;AAAA,IACF;AAEA,UAAM,WAAW,YAAY;AAC3B,YAAM,WAAW,EAAE,OAAO,SAAS,SAAS,gBAAgB,KAAK,KAAK,EAAE;AACxE,SAAI,KAAK,OAAO,SAAK,uBAAO,QAAQ,CAAC,CAAC;AAEtC,uBAAiB,MAAM,YAAY;AACjC,YAAI,KAAK,gBAAgB,OAAO,QAAS;AACzC,cAAM,WAAW,GAAG;AACpB,YAAI,CAAC,SAAU;AACf,WAAI,KAAK,OAAO,SAAK,uBAAO,EAAE,OAAO,QAAQ,MAAM,WAAW,IAAI,CAAC,CAAC,CAAC;AACrE,WAAI,KAAK,OAAO,SAAK,uBAAO,EAAE,OAAO,QAAQ,CAAC,CAAC,CAAC;AAAA,MAClD;AAEA,UAAI,CAAC,KAAK,gBAAgB,OAAO,SAAS;AACxC,WAAI,KAAK,OAAO,SAAK,uBAAO,EAAE,OAAO,OAAO,CAAC,CAAC,CAAC;AAAA,MACjD;AAAA,IACF;AAEA,QAAI;AACJ,UAAM,gBAAgB,CAAC,UAAmB;AACxC,UAAI,WAAW;AACb,aAAK,MAAM,IAAI,EAAE,WAAW,WAAW,WAAW,OAAO,WAAW,MAAM,CAAC;AAC3E,oBAAY;AAAA,MACd;AAAA,IACF;AAEA,UAAM,WAAW,YAAY;AAG3B,YAAM,YAAY,CAAC,QAAiB;AAClC,YAAI;AACJ,YAAI,OAAO,SAAS,GAAG,GAAG;AACxB,kBAAQ;AAAA,QACV,WAAW,MAAM,QAAQ,GAAG,GAAG;AAC7B,kBAAQ,OAAO,OAAO,GAAG;AAAA,QAC3B,OAAO;AACL,kBAAQ,OAAO,KAAK,GAAkB;AAAA,QACxC;AAEA,YAAI;AACJ,YAAI;AACF,uBAAS,uBAAO,KAAK;AAAA,QACvB,SAAS,KAAK;AACZ,eAAK,QAAQ,KAAK,EAAE,IAAI,GAAG,qCAAqC;AAChE;AAAA,QACF;AAEA,cAAM,QAAQ,OAAO;AACrB,YAAI,UAAU,SAAS;AACrB,gBAAM,QAAQ,OAAO;AACrB,cAAI,SAAS,MAAM,aAAa,GAAG;AACjC,uBAAW,KAAK,QAAQ,MAAM,KAAK,GAAG;AACpC,4BAAc,KAAK;AACnB,0BAAY;AAAA,YACd;AAAA,UACF;AAAA,QACF,WAAW,UAAU,UAAU;AAC7B,gBAAM,SAAS,OAAO;AACtB,cAAI,WAAW,SAAS;AACtB,qBAAS;AAAA,cACP,IAAI,6BAAe;AAAA,gBACjB,SAAS;AAAA,gBACT,SAAS,EAAE,MAAM,EAAE,KAAK,KAAK,UAAU,MAAM,EAAE,EAAE;AAAA,cACnD,CAAC;AAAA,YACH;AACA;AAAA,UACF;AACA,qBAAW,KAAK,QAAQ,MAAM,GAAG;AAC/B,0BAAc,KAAK;AACnB,wBAAY;AAAA,UACd;AACA,wBAAc,IAAI;AAClB,cAAI,CAAC,KAAK,MAAM,QAAQ;AACtB,iBAAK,MAAM,IAAI,iBAAiB,aAAa;AAAA,UAC/C;AACA,cAAI,CAAC,SAAS,KAAM,UAAS,QAAQ;AAAA,QACvC,OAAO;AACL,eAAK,QAAQ,MAAM,EAAE,MAAM,GAAG,0BAA0B;AAAA,QAC1D;AAAA,MACF;AAEA,YAAM,UAAU,CAAC,MAAc,WAAmB;AAChD,YAAI,CAAC,SAAS,MAAM;AAClB,mBAAS;AAAA,YACP,IAAI,6BAAe;AAAA,cACjB,SAAS;AAAA,cACT,SAAS;AAAA,gBACP,YAAY,QAAQ;AAAA,gBACpB,MAAM,EAAE,QAAQ,OAAO,SAAS,EAAE;AAAA,cACpC;AAAA,YACF,CAAC;AAAA,UACH;AAAA,QACF;AAAA,MACF;AAEA,YAAM,UAAU,CAAC,QAAe;AAC9B,YAAI,CAAC,SAAS,KAAM,UAAS,OAAO,GAAG;AAAA,MACzC;AAEA,SAAI,GAAG,WAAW,SAAS;AAC3B,SAAI,GAAG,SAAS,OAAO;AACvB,SAAI,GAAG,SAAS,OAAO;AAEvB,UAAI;AACF,cAAM,SAAS;AAAA,MACjB,UAAE;AACA,WAAI,IAAI,WAAW,SAAS;AAC5B,WAAI,IAAI,SAAS,OAAO;AACxB,WAAI,IAAI,SAAS,OAAO;AAAA,MAC1B;AAAA,IACF;AAEA,QAAI;AACF,YAAM,QAAQ,IAAI,CAAC,UAAU,GAAG,SAAS,GAAG,SAAS,CAAC,CAAC;AAAA,IACzD,SAAS,GAAG;AACV,UAAI,KAAK,YAAY,QAAS;AAC9B,UAAI,aAAa,gCAAkB,aAAa,kCAAoB;AAClE,cAAM;AAAA,MACR;AACA,YAAM,IAAI,iCAAmB;AAAA,QAC3B,SAAS,gCAAiC,EAAY,WAAW,eAAe;AAAA,MAClF,CAAC;AAAA,IACH,UAAE;AACA,UAAI,CAAC,WAAW,OAAQ,YAAW,MAAM;AACzC,UAAI,MAAM,GAAG,eAAe,oBAAU,QAAQ;AAC5C,YAAI;AACF,aAAG,MAAM;AAAA,QACX,QAAQ;AAAA,QAER;AAAA,MACF;AAAA,IACF;AAAA,EACF;AACF;AAEA,MAAM,mBAAmB,OAAO;AAAA,EAC9B;AAAA,EACA;AAAA,EACA;AAAA,EACA;AACF,MAK0B;AACxB,QAAM,KAAK,IAAI,oBAAU,KAAK,EAAE,SAAS,kBAAkB,UAAU,CAAC;AACtE,QAAM,MAAM,IAAI,qBAAa;AAE7B,MAAI;AACJ,QAAM,UAAU,MAAM;AACpB,QAAI,QAAS,cAAa,OAAO;AACjC,OAAG,IAAI,QAAQ,MAAM;AACrB,OAAG,IAAI,SAAS,OAAO;AACvB,OAAG,IAAI,SAAS,OAAO;AACvB,gBAAY,oBAAoB,SAAS,OAAO;AAAA,EAClD;AAEA,QAAM,SAAS,MAAM,IAAI,QAAQ;AACjC,QAAM,UAAU,CAAC,QAAe,IAAI,OAAO,GAAG;AAC9C,QAAM,UAAU,CAAC,MAAc,WAC7B,IAAI;AAAA,IACF,IAAI,MAAM,sCAAsC,IAAI,YAAY,OAAO,SAAS,CAAC,GAAG;AAAA,EACtF;AACF,QAAM,UAAU,MAAM,IAAI,OAAO,IAAI,MAAM,SAAS,CAAC;AAErD,KAAG,GAAG,QAAQ,MAAM;AACpB,KAAG,GAAG,SAAS,OAAO;AACtB,KAAG,GAAG,SAAS,OAAO;AACtB,cAAY,iBAAiB,SAAS,SAAS,EAAE,MAAM,KAAK,CAAC;AAE7D,MAAI,YAAY,GAAG;AACjB,cAAU,WAAW,MAAM,IAAI,OAAO,IAAI,MAAM,iBAAiB,CAAC,GAAG,SAAS;AAAA,EAChF;AAEA,MAAI;AACF,UAAM,IAAI;AACV,WAAO;AAAA,EACT,SAAS,GAAG;AACV,QAAI;AACF,SAAG,GAAG,SAAS,MAAM;AAAA,MAAC,CAAC;AACvB,UAAI,GAAG,eAAe,oBAAU,YAAY;AAC1C,WAAG,MAAM;AAAA,MACX,OAAO;AACL,WAAG,UAAU;AAAA,MACf;AAAA,IACF,QAAQ;AAAA,IAER;AACA,UAAM;AAAA,EACR,UAAE;AACA,YAAQ;AAAA,EACV;AACF;","names":["tts"]} | ||
| {"version":3,"sources":["../src/tts.ts"],"sourcesContent":["// SPDX-FileCopyrightText: 2026 LiveKit, Inc.\n//\n// SPDX-License-Identifier: Apache-2.0\nimport {\n type APIConnectOptions,\n APIConnectionError,\n APIStatusError,\n AudioByteStream,\n Future,\n log,\n shortuuid,\n tokenize,\n tts,\n} from '@livekit/agents';\nimport type { AudioFrame } from '@livekit/rtc-node';\nimport { decode, encode } from '@msgpack/msgpack';\nimport { request } from 'node:https';\nimport { type RawData, WebSocket } from 'ws';\nimport type { LatencyMode, TTSModels } from './models.js';\n\nconst DEFAULT_MODEL: TTSModels = 's2.1-pro';\nconst DEFAULT_VOICE_ID = '933563129e564b19a115bedd57b7406a';\nconst DEFAULT_BASE_URL = 'https://api.fish.audio';\nconst NUM_CHANNELS = 1;\n// Fish Audio's default sample rate for raw PCM output.\nconst DEFAULT_SAMPLE_RATE = 24000;\n\nexport interface TTSOptions {\n apiKey?: string;\n model?: TTSModels | string;\n voiceId?: string;\n sampleRate?: number;\n baseURL?: string;\n latencyMode?: LatencyMode;\n /**\n * Upper bound on the number of characters Fish buffers before auto-synthesizing.\n * Must be between 100 and 300. With sentence-level flushing this is only hit by\n * sentences longer than `chunkLength`; otherwise audio is produced as soon as\n * each sentence is flushed. Defaults to 100.\n */\n chunkLength?: number;\n /**\n * Speaking rate multiplier for Fish `prosody.speed`. `1.0` is normal; below\n * 1.0 is slower, above is faster. Unset uses the voice's natural pace.\n */\n speed?: number;\n /**\n * Loudness adjustment in decibels for Fish `prosody.volume`. `0` is the\n * voice's natural level. Unset leaves it unchanged.\n */\n volume?: number;\n tokenizer?: tokenize.SentenceTokenizer;\n}\n\ninterface ResolvedTTSOptions {\n apiKey: string;\n model: TTSModels | string;\n voiceId?: string;\n sampleRate: number;\n baseURL: string;\n latencyMode: LatencyMode;\n chunkLength: number;\n speed?: number;\n volume?: number;\n tokenizer: tokenize.SentenceTokenizer;\n}\n\nconst DEFAULT_OPTS: Omit<ResolvedTTSOptions, 'apiKey' | 'tokenizer'> = {\n model: DEFAULT_MODEL,\n voiceId: DEFAULT_VOICE_ID,\n sampleRate: DEFAULT_SAMPLE_RATE,\n baseURL: DEFAULT_BASE_URL,\n latencyMode: 'balanced',\n chunkLength: 100,\n};\n\nconst validateChunkLength = (chunkLength: number) => {\n if (!Number.isFinite(chunkLength) || chunkLength < 100 || chunkLength > 300) {\n throw new Error('chunkLength must be between 100 and 300');\n }\n};\n\n// Fish Audio's wire format mirrors the upstream Python SDK so the server\n// doesn't fall back to its own larger defaults — in particular the docs default\n// of `chunk_length=300` produces large bursts that leave audible gaps between\n// chunk boundaries.\nconst buildTtsRequest = (opts: ResolvedTTSOptions, text: string = ''): Record<string, unknown> => {\n const prosody =\n opts.speed !== undefined || opts.volume !== undefined\n ? {\n ...(opts.speed !== undefined ? { speed: opts.speed } : {}),\n ...(opts.volume !== undefined ? { volume: opts.volume } : {}),\n }\n : null;\n\n return {\n text,\n chunk_length: opts.chunkLength,\n format: 'pcm',\n sample_rate: opts.sampleRate,\n mp3_bitrate: 64,\n opus_bitrate: 64000,\n references: [],\n // Fish Audio's wire field is `reference_id`; we expose it as `voiceId` on\n // the plugin for consistency with other TTS plugins.\n reference_id: opts.voiceId ?? null,\n normalize: true,\n latency: opts.latencyMode,\n prosody,\n top_p: 0.7,\n temperature: 0.7,\n };\n};\n\nexport class TTS extends tts.TTS {\n #opts: ResolvedTTSOptions;\n label = 'fishaudio.TTS';\n\n constructor(opts: TTSOptions = {}) {\n const apiKey = opts.apiKey ?? process.env.FISH_API_KEY;\n if (!apiKey) {\n throw new Error(\n 'Fish Audio API key is required, either as argument or set FISH_API_KEY environment variable',\n );\n }\n\n const chunkLength = opts.chunkLength ?? DEFAULT_OPTS.chunkLength;\n validateChunkLength(chunkLength);\n\n const sampleRate = opts.sampleRate ?? DEFAULT_OPTS.sampleRate;\n\n super(sampleRate, NUM_CHANNELS, { streaming: true });\n\n // min_sentence_len=1 emits each sentence as soon as the next one starts,\n // rather than batching short sentences together — minimizes TTFB on the\n // first sentence and keeps Fish synthesizing continuously.\n const tokenizer =\n opts.tokenizer ?? new tokenize.basic.SentenceTokenizer({ minSentenceLength: 1 });\n\n this.#opts = {\n apiKey,\n model: opts.model ?? DEFAULT_OPTS.model,\n voiceId: opts.voiceId ?? DEFAULT_OPTS.voiceId,\n sampleRate,\n baseURL: opts.baseURL ?? DEFAULT_OPTS.baseURL,\n latencyMode: opts.latencyMode ?? DEFAULT_OPTS.latencyMode,\n chunkLength,\n speed: opts.speed,\n volume: opts.volume,\n tokenizer,\n };\n }\n\n get model(): string {\n return this.#opts.model;\n }\n\n get provider(): string {\n return 'FishAudio';\n }\n\n updateOptions(opts: {\n model?: TTSModels | string;\n voiceId?: string;\n latencyMode?: LatencyMode;\n chunkLength?: number;\n speed?: number;\n volume?: number;\n }): void {\n if (opts.model !== undefined) this.#opts.model = opts.model;\n if (opts.voiceId !== undefined) this.#opts.voiceId = opts.voiceId;\n if (opts.latencyMode !== undefined) this.#opts.latencyMode = opts.latencyMode;\n if (opts.chunkLength !== undefined) {\n validateChunkLength(opts.chunkLength);\n this.#opts.chunkLength = opts.chunkLength;\n }\n if (opts.speed !== undefined) this.#opts.speed = opts.speed;\n if (opts.volume !== undefined) this.#opts.volume = opts.volume;\n }\n\n synthesize(\n text: string,\n connOptions?: APIConnectOptions,\n abortSignal?: AbortSignal,\n ): tts.ChunkedStream {\n return new ChunkedStream(this, text, this.#opts, connOptions, abortSignal);\n }\n\n stream(options?: { connOptions?: APIConnectOptions }): tts.SynthesizeStream {\n return new SynthesizeStream(this, this.#opts, options?.connOptions);\n }\n}\n\nexport class ChunkedStream extends tts.ChunkedStream {\n label = 'fishaudio.ChunkedStream';\n #logger = log();\n #opts: ResolvedTTSOptions;\n #text: string;\n\n constructor(\n tts: TTS,\n text: string,\n opts: ResolvedTTSOptions,\n connOptions?: APIConnectOptions,\n abortSignal?: AbortSignal,\n ) {\n super(text, tts, connOptions, abortSignal);\n this.#text = text;\n this.#opts = opts;\n }\n\n protected async run() {\n const requestId = shortuuid();\n const bstream = new AudioByteStream(this.#opts.sampleRate, NUM_CHANNELS);\n const payload = encode(buildTtsRequest(this.#opts, this.#text));\n\n const baseUrl = new URL(this.#opts.baseURL);\n const isHttps = baseUrl.protocol === 'https:';\n if (!isHttps) {\n // The plugin only supports https; fall back via Node's http module is\n // intentionally not implemented to keep the code path simple.\n throw new APIConnectionError({\n message: `Fish Audio base URL must use https (got ${this.#opts.baseURL})`,\n });\n }\n\n const doneFut = new Future<void>();\n\n const req = request(\n {\n hostname: baseUrl.hostname,\n port: parseInt(baseUrl.port) || 443,\n path: '/v1/tts',\n method: 'POST',\n headers: {\n Authorization: `Bearer ${this.#opts.apiKey}`,\n 'Content-Type': 'application/msgpack',\n model: this.#opts.model,\n 'Content-Length': payload.byteLength,\n },\n signal: this.abortSignal,\n },\n (res) => {\n const status = res.statusCode ?? -1;\n if (status < 200 || status >= 300) {\n const chunks: Buffer[] = [];\n res.on('data', (c: Buffer) => chunks.push(c));\n res.on('end', () => {\n const body = Buffer.concat(chunks).toString();\n if (!doneFut.done) {\n doneFut.reject(\n new APIStatusError({\n message: `Fish Audio TTS request failed: ${body}`,\n options: { statusCode: status, body: { raw: body } },\n }),\n );\n }\n });\n res.on('error', (err) => {\n if (err.message === 'aborted') return;\n this.#logger.error({ err }, 'Fish Audio TTS error response stream error');\n if (!doneFut.done) {\n doneFut.reject(\n new APIStatusError({\n message: `Fish Audio TTS request failed (status ${status})`,\n options: { statusCode: status },\n }),\n );\n }\n });\n return;\n }\n\n res.on('data', (chunk: Buffer) => {\n for (const frame of bstream.write(chunk)) {\n this.queue.put({\n requestId,\n segmentId: requestId,\n frame,\n final: false,\n });\n }\n });\n res.on('close', () => {\n for (const frame of bstream.flush()) {\n this.queue.put({\n requestId,\n segmentId: requestId,\n frame,\n final: false,\n });\n }\n if (!this.queue.closed) this.queue.close();\n if (!doneFut.done) doneFut.resolve();\n });\n res.on('error', (err) => {\n if (err.message === 'aborted') return;\n this.#logger.error({ err }, 'Fish Audio TTS response error');\n if (!doneFut.done) doneFut.reject(err);\n });\n },\n );\n\n req.on('error', (err) => {\n if (err.name === 'AbortError') return;\n this.#logger.error({ err }, 'Fish Audio TTS request error');\n if (!doneFut.done) doneFut.reject(err);\n });\n req.write(payload);\n req.end();\n\n try {\n await doneFut.await;\n } catch (e) {\n if (this.abortSignal.aborted) return;\n if (!this.queue.closed) this.queue.close();\n if (e instanceof APIStatusError || e instanceof APIConnectionError) {\n throw e;\n }\n throw new APIConnectionError({\n message: `Fish Audio connection failed: ${(e as Error).message ?? 'unknown error'}`,\n });\n }\n }\n}\n\nexport class SynthesizeStream extends tts.SynthesizeStream {\n label = 'fishaudio.SynthesizeStream';\n #logger = log();\n #opts: ResolvedTTSOptions;\n\n constructor(tts: TTS, opts: ResolvedTTSOptions, connOptions?: APIConnectOptions) {\n super(tts, connOptions);\n this.#opts = opts;\n }\n\n protected async run() {\n const requestId = shortuuid();\n const bstream = new AudioByteStream(this.#opts.sampleRate, NUM_CHANNELS);\n\n // Tokenize incoming text by sentence and flush after each sentence so Fish\n // synthesizes immediately at sentence boundaries instead of waiting for\n // `chunkLength` characters to accumulate. The result is much smoother\n // audio: gaps line up with sentence breaks (where pauses are natural)\n // rather than mid-clause.\n const sentStream = this.#opts.tokenizer.stream();\n\n const wsUrl = `${this.#opts.baseURL.replace(/^http/, 'ws')}/v1/tts/live`;\n let ws: WebSocket | undefined;\n try {\n ws = await connectWebSocket({\n url: wsUrl,\n headers: {\n Authorization: `Bearer ${this.#opts.apiKey}`,\n model: this.#opts.model,\n },\n timeoutMs: this.connOptions.timeoutMs,\n abortSignal: this.abortSignal,\n });\n } catch (e) {\n throw new APIConnectionError({\n message: `Fish Audio websocket connect failed: ${(e as Error).message ?? 'unknown error'}`,\n });\n }\n\n const finished = new Future<void>();\n\n const inputTask = async () => {\n try {\n for await (const data of this.input) {\n if (this.abortController.signal.aborted) break;\n if (data === SynthesizeStream.FLUSH_SENTINEL) {\n sentStream.flush();\n continue;\n }\n if (!data) continue;\n sentStream.pushText(data);\n }\n } finally {\n if (!sentStream.closed) sentStream.endInput();\n }\n };\n\n const sendTask = async () => {\n const startMsg = { event: 'start', request: buildTtsRequest(this.#opts) };\n ws!.send(Buffer.from(encode(startMsg)));\n\n for await (const ev of sentStream) {\n if (this.abortController.signal.aborted) break;\n const sentence = ev.token;\n if (!sentence) continue;\n this.markStarted();\n ws!.send(Buffer.from(encode({ event: 'text', text: sentence + ' ' })));\n ws!.send(Buffer.from(encode({ event: 'flush' })));\n }\n\n if (!this.abortController.signal.aborted) {\n ws!.send(Buffer.from(encode({ event: 'stop' })));\n }\n };\n\n let lastFrame: AudioFrame | undefined;\n const sendLastFrame = (final: boolean) => {\n if (lastFrame) {\n this.queue.put({ requestId, segmentId: requestId, frame: lastFrame, final });\n lastFrame = undefined;\n }\n };\n\n const recvTask = async () => {\n // No per-receive timeout: Fish has natural inter-sentence gaps that can\n // exceed connOptions.timeoutMs when the LLM is slow.\n const onMessage = (raw: RawData) => {\n let frame: Buffer;\n if (Buffer.isBuffer(raw)) {\n frame = raw;\n } else if (Array.isArray(raw)) {\n frame = Buffer.concat(raw);\n } else {\n frame = Buffer.from(raw as ArrayBuffer);\n }\n\n let parsed: Record<string, unknown>;\n try {\n parsed = decode(frame) as Record<string, unknown>;\n } catch (err) {\n this.#logger.warn({ err }, 'Fish Audio failed to decode message');\n return;\n }\n\n const event = parsed.event as string | undefined;\n if (event === 'audio') {\n const audio = parsed.audio as Uint8Array | undefined;\n if (audio && audio.byteLength > 0) {\n for (const f of bstream.write(audio)) {\n sendLastFrame(false);\n lastFrame = f;\n }\n }\n } else if (event === 'finish') {\n const reason = parsed.reason as string | undefined;\n if (reason === 'error') {\n finished.reject(\n new APIStatusError({\n message: 'Fish Audio TTS reported an error',\n options: { body: { raw: JSON.stringify(parsed) } },\n }),\n );\n return;\n }\n for (const f of bstream.flush()) {\n sendLastFrame(false);\n lastFrame = f;\n }\n sendLastFrame(true);\n if (!this.queue.closed) {\n this.queue.put(SynthesizeStream.END_OF_STREAM);\n }\n if (!finished.done) finished.resolve();\n } else {\n this.#logger.debug({ event }, 'unknown Fish Audio event');\n }\n };\n\n const onClose = (code: number, reason: Buffer) => {\n if (!finished.done) {\n finished.reject(\n new APIStatusError({\n message: 'Fish Audio websocket connection closed unexpectedly',\n options: {\n statusCode: code || -1,\n body: { reason: reason.toString() },\n },\n }),\n );\n }\n };\n\n const onError = (err: Error) => {\n if (!finished.done) finished.reject(err);\n };\n\n ws!.on('message', onMessage);\n ws!.on('close', onClose);\n ws!.on('error', onError);\n\n try {\n await finished.await;\n } finally {\n ws!.off('message', onMessage);\n ws!.off('close', onClose);\n ws!.off('error', onError);\n }\n };\n\n try {\n await Promise.all([inputTask(), sendTask(), recvTask()]);\n } catch (e) {\n if (this.abortSignal.aborted) return;\n if (e instanceof APIStatusError || e instanceof APIConnectionError) {\n throw e;\n }\n throw new APIConnectionError({\n message: `Fish Audio websocket failed: ${(e as Error).message ?? 'unknown error'}`,\n });\n } finally {\n if (!sentStream.closed) sentStream.close();\n if (ws && ws.readyState !== WebSocket.CLOSED) {\n try {\n ws.close();\n } catch {\n // ignore\n }\n }\n }\n }\n}\n\nconst connectWebSocket = async ({\n url,\n headers,\n timeoutMs,\n abortSignal,\n}: {\n url: string;\n headers: Record<string, string>;\n timeoutMs: number;\n abortSignal: AbortSignal;\n}): Promise<WebSocket> => {\n const ws = new WebSocket(url, { headers, handshakeTimeout: timeoutMs });\n const fut = new Future<void>();\n\n let timeout: NodeJS.Timeout | undefined;\n const cleanup = () => {\n if (timeout) clearTimeout(timeout);\n ws.off('open', onOpen);\n ws.off('error', onError);\n ws.off('close', onClose);\n abortSignal.removeEventListener('abort', onAbort);\n };\n\n const onOpen = () => fut.resolve();\n const onError = (err: Error) => fut.reject(err);\n const onClose = (code: number, reason: Buffer) =>\n fut.reject(\n new Error(`websocket closed before open (code=${code}, reason=${reason.toString()})`),\n );\n const onAbort = () => fut.reject(new Error('aborted'));\n\n ws.on('open', onOpen);\n ws.on('error', onError);\n ws.on('close', onClose);\n abortSignal.addEventListener('abort', onAbort, { once: true });\n\n if (timeoutMs > 0) {\n timeout = setTimeout(() => fut.reject(new Error('connect timeout')), timeoutMs);\n }\n\n try {\n await fut.await;\n return ws;\n } catch (e) {\n try {\n ws.on('error', () => {});\n if (ws.readyState === WebSocket.CONNECTING) {\n ws.close();\n } else {\n ws.terminate();\n }\n } catch {\n // ignore\n }\n throw e;\n } finally {\n cleanup();\n }\n};\n"],"mappings":";;;;;;;;;;;;;;;;;;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAGA,oBAUO;AAEP,qBAA+B;AAC/B,wBAAwB;AACxB,gBAAwC;AAGxC,MAAM,gBAA2B;AACjC,MAAM,mBAAmB;AACzB,MAAM,mBAAmB;AACzB,MAAM,eAAe;AAErB,MAAM,sBAAsB;AA0C5B,MAAM,eAAiE;AAAA,EACrE,OAAO;AAAA,EACP,SAAS;AAAA,EACT,YAAY;AAAA,EACZ,SAAS;AAAA,EACT,aAAa;AAAA,EACb,aAAa;AACf;AAEA,MAAM,sBAAsB,CAAC,gBAAwB;AACnD,MAAI,CAAC,OAAO,SAAS,WAAW,KAAK,cAAc,OAAO,cAAc,KAAK;AAC3E,UAAM,IAAI,MAAM,yCAAyC;AAAA,EAC3D;AACF;AAMA,MAAM,kBAAkB,CAAC,MAA0B,OAAe,OAAgC;AAChG,QAAM,UACJ,KAAK,UAAU,UAAa,KAAK,WAAW,SACxC;AAAA,IACE,GAAI,KAAK,UAAU,SAAY,EAAE,OAAO,KAAK,MAAM,IAAI,CAAC;AAAA,IACxD,GAAI,KAAK,WAAW,SAAY,EAAE,QAAQ,KAAK,OAAO,IAAI,CAAC;AAAA,EAC7D,IACA;AAEN,SAAO;AAAA,IACL;AAAA,IACA,cAAc,KAAK;AAAA,IACnB,QAAQ;AAAA,IACR,aAAa,KAAK;AAAA,IAClB,aAAa;AAAA,IACb,cAAc;AAAA,IACd,YAAY,CAAC;AAAA;AAAA;AAAA,IAGb,cAAc,KAAK,WAAW;AAAA,IAC9B,WAAW;AAAA,IACX,SAAS,KAAK;AAAA,IACd;AAAA,IACA,OAAO;AAAA,IACP,aAAa;AAAA,EACf;AACF;AAEO,MAAM,YAAY,kBAAI,IAAI;AAAA,EAC/B;AAAA,EACA,QAAQ;AAAA,EAER,YAAY,OAAmB,CAAC,GAAG;AACjC,UAAM,SAAS,KAAK,UAAU,QAAQ,IAAI;AAC1C,QAAI,CAAC,QAAQ;AACX,YAAM,IAAI;AAAA,QACR;AAAA,MACF;AAAA,IACF;AAEA,UAAM,cAAc,KAAK,eAAe,aAAa;AACrD,wBAAoB,WAAW;AAE/B,UAAM,aAAa,KAAK,cAAc,aAAa;AAEnD,UAAM,YAAY,cAAc,EAAE,WAAW,KAAK,CAAC;AAKnD,UAAM,YACJ,KAAK,aAAa,IAAI,uBAAS,MAAM,kBAAkB,EAAE,mBAAmB,EAAE,CAAC;AAEjF,SAAK,QAAQ;AAAA,MACX;AAAA,MACA,OAAO,KAAK,SAAS,aAAa;AAAA,MAClC,SAAS,KAAK,WAAW,aAAa;AAAA,MACtC;AAAA,MACA,SAAS,KAAK,WAAW,aAAa;AAAA,MACtC,aAAa,KAAK,eAAe,aAAa;AAAA,MAC9C;AAAA,MACA,OAAO,KAAK;AAAA,MACZ,QAAQ,KAAK;AAAA,MACb;AAAA,IACF;AAAA,EACF;AAAA,EAEA,IAAI,QAAgB;AAClB,WAAO,KAAK,MAAM;AAAA,EACpB;AAAA,EAEA,IAAI,WAAmB;AACrB,WAAO;AAAA,EACT;AAAA,EAEA,cAAc,MAOL;AACP,QAAI,KAAK,UAAU,OAAW,MAAK,MAAM,QAAQ,KAAK;AACtD,QAAI,KAAK,YAAY,OAAW,MAAK,MAAM,UAAU,KAAK;AAC1D,QAAI,KAAK,gBAAgB,OAAW,MAAK,MAAM,cAAc,KAAK;AAClE,QAAI,KAAK,gBAAgB,QAAW;AAClC,0BAAoB,KAAK,WAAW;AACpC,WAAK,MAAM,cAAc,KAAK;AAAA,IAChC;AACA,QAAI,KAAK,UAAU,OAAW,MAAK,MAAM,QAAQ,KAAK;AACtD,QAAI,KAAK,WAAW,OAAW,MAAK,MAAM,SAAS,KAAK;AAAA,EAC1D;AAAA,EAEA,WACE,MACA,aACA,aACmB;AACnB,WAAO,IAAI,cAAc,MAAM,MAAM,KAAK,OAAO,aAAa,WAAW;AAAA,EAC3E;AAAA,EAEA,OAAO,SAAqE;AAC1E,WAAO,IAAI,iBAAiB,MAAM,KAAK,OAAO,mCAAS,WAAW;AAAA,EACpE;AACF;AAEO,MAAM,sBAAsB,kBAAI,cAAc;AAAA,EACnD,QAAQ;AAAA,EACR,cAAU,mBAAI;AAAA,EACd;AAAA,EACA;AAAA,EAEA,YACEA,MACA,MACA,MACA,aACA,aACA;AACA,UAAM,MAAMA,MAAK,aAAa,WAAW;AACzC,SAAK,QAAQ;AACb,SAAK,QAAQ;AAAA,EACf;AAAA,EAEA,MAAgB,MAAM;AACpB,UAAM,gBAAY,yBAAU;AAC5B,UAAM,UAAU,IAAI,8BAAgB,KAAK,MAAM,YAAY,YAAY;AACvE,UAAM,cAAU,uBAAO,gBAAgB,KAAK,OAAO,KAAK,KAAK,CAAC;AAE9D,UAAM,UAAU,IAAI,IAAI,KAAK,MAAM,OAAO;AAC1C,UAAM,UAAU,QAAQ,aAAa;AACrC,QAAI,CAAC,SAAS;AAGZ,YAAM,IAAI,iCAAmB;AAAA,QAC3B,SAAS,2CAA2C,KAAK,MAAM,OAAO;AAAA,MACxE,CAAC;AAAA,IACH;AAEA,UAAM,UAAU,IAAI,qBAAa;AAEjC,UAAM,UAAM;AAAA,MACV;AAAA,QACE,UAAU,QAAQ;AAAA,QAClB,MAAM,SAAS,QAAQ,IAAI,KAAK;AAAA,QAChC,MAAM;AAAA,QACN,QAAQ;AAAA,QACR,SAAS;AAAA,UACP,eAAe,UAAU,KAAK,MAAM,MAAM;AAAA,UAC1C,gBAAgB;AAAA,UAChB,OAAO,KAAK,MAAM;AAAA,UAClB,kBAAkB,QAAQ;AAAA,QAC5B;AAAA,QACA,QAAQ,KAAK;AAAA,MACf;AAAA,MACA,CAAC,QAAQ;AACP,cAAM,SAAS,IAAI,cAAc;AACjC,YAAI,SAAS,OAAO,UAAU,KAAK;AACjC,gBAAM,SAAmB,CAAC;AAC1B,cAAI,GAAG,QAAQ,CAAC,MAAc,OAAO,KAAK,CAAC,CAAC;AAC5C,cAAI,GAAG,OAAO,MAAM;AAClB,kBAAM,OAAO,OAAO,OAAO,MAAM,EAAE,SAAS;AAC5C,gBAAI,CAAC,QAAQ,MAAM;AACjB,sBAAQ;AAAA,gBACN,IAAI,6BAAe;AAAA,kBACjB,SAAS,kCAAkC,IAAI;AAAA,kBAC/C,SAAS,EAAE,YAAY,QAAQ,MAAM,EAAE,KAAK,KAAK,EAAE;AAAA,gBACrD,CAAC;AAAA,cACH;AAAA,YACF;AAAA,UACF,CAAC;AACD,cAAI,GAAG,SAAS,CAAC,QAAQ;AACvB,gBAAI,IAAI,YAAY,UAAW;AAC/B,iBAAK,QAAQ,MAAM,EAAE,IAAI,GAAG,4CAA4C;AACxE,gBAAI,CAAC,QAAQ,MAAM;AACjB,sBAAQ;AAAA,gBACN,IAAI,6BAAe;AAAA,kBACjB,SAAS,yCAAyC,MAAM;AAAA,kBACxD,SAAS,EAAE,YAAY,OAAO;AAAA,gBAChC,CAAC;AAAA,cACH;AAAA,YACF;AAAA,UACF,CAAC;AACD;AAAA,QACF;AAEA,YAAI,GAAG,QAAQ,CAAC,UAAkB;AAChC,qBAAW,SAAS,QAAQ,MAAM,KAAK,GAAG;AACxC,iBAAK,MAAM,IAAI;AAAA,cACb;AAAA,cACA,WAAW;AAAA,cACX;AAAA,cACA,OAAO;AAAA,YACT,CAAC;AAAA,UACH;AAAA,QACF,CAAC;AACD,YAAI,GAAG,SAAS,MAAM;AACpB,qBAAW,SAAS,QAAQ,MAAM,GAAG;AACnC,iBAAK,MAAM,IAAI;AAAA,cACb;AAAA,cACA,WAAW;AAAA,cACX;AAAA,cACA,OAAO;AAAA,YACT,CAAC;AAAA,UACH;AACA,cAAI,CAAC,KAAK,MAAM,OAAQ,MAAK,MAAM,MAAM;AACzC,cAAI,CAAC,QAAQ,KAAM,SAAQ,QAAQ;AAAA,QACrC,CAAC;AACD,YAAI,GAAG,SAAS,CAAC,QAAQ;AACvB,cAAI,IAAI,YAAY,UAAW;AAC/B,eAAK,QAAQ,MAAM,EAAE,IAAI,GAAG,+BAA+B;AAC3D,cAAI,CAAC,QAAQ,KAAM,SAAQ,OAAO,GAAG;AAAA,QACvC,CAAC;AAAA,MACH;AAAA,IACF;AAEA,QAAI,GAAG,SAAS,CAAC,QAAQ;AACvB,UAAI,IAAI,SAAS,aAAc;AAC/B,WAAK,QAAQ,MAAM,EAAE,IAAI,GAAG,8BAA8B;AAC1D,UAAI,CAAC,QAAQ,KAAM,SAAQ,OAAO,GAAG;AAAA,IACvC,CAAC;AACD,QAAI,MAAM,OAAO;AACjB,QAAI,IAAI;AAER,QAAI;AACF,YAAM,QAAQ;AAAA,IAChB,SAAS,GAAG;AACV,UAAI,KAAK,YAAY,QAAS;AAC9B,UAAI,CAAC,KAAK,MAAM,OAAQ,MAAK,MAAM,MAAM;AACzC,UAAI,aAAa,gCAAkB,aAAa,kCAAoB;AAClE,cAAM;AAAA,MACR;AACA,YAAM,IAAI,iCAAmB;AAAA,QAC3B,SAAS,iCAAkC,EAAY,WAAW,eAAe;AAAA,MACnF,CAAC;AAAA,IACH;AAAA,EACF;AACF;AAEO,MAAM,yBAAyB,kBAAI,iBAAiB;AAAA,EACzD,QAAQ;AAAA,EACR,cAAU,mBAAI;AAAA,EACd;AAAA,EAEA,YAAYA,MAAU,MAA0B,aAAiC;AAC/E,UAAMA,MAAK,WAAW;AACtB,SAAK,QAAQ;AAAA,EACf;AAAA,EAEA,MAAgB,MAAM;AACpB,UAAM,gBAAY,yBAAU;AAC5B,UAAM,UAAU,IAAI,8BAAgB,KAAK,MAAM,YAAY,YAAY;AAOvE,UAAM,aAAa,KAAK,MAAM,UAAU,OAAO;AAE/C,UAAM,QAAQ,GAAG,KAAK,MAAM,QAAQ,QAAQ,SAAS,IAAI,CAAC;AAC1D,QAAI;AACJ,QAAI;AACF,WAAK,MAAM,iBAAiB;AAAA,QAC1B,KAAK;AAAA,QACL,SAAS;AAAA,UACP,eAAe,UAAU,KAAK,MAAM,MAAM;AAAA,UAC1C,OAAO,KAAK,MAAM;AAAA,QACpB;AAAA,QACA,WAAW,KAAK,YAAY;AAAA,QAC5B,aAAa,KAAK;AAAA,MACpB,CAAC;AAAA,IACH,SAAS,GAAG;AACV,YAAM,IAAI,iCAAmB;AAAA,QAC3B,SAAS,wCAAyC,EAAY,WAAW,eAAe;AAAA,MAC1F,CAAC;AAAA,IACH;AAEA,UAAM,WAAW,IAAI,qBAAa;AAElC,UAAM,YAAY,YAAY;AAC5B,UAAI;AACF,yBAAiB,QAAQ,KAAK,OAAO;AACnC,cAAI,KAAK,gBAAgB,OAAO,QAAS;AACzC,cAAI,SAAS,iBAAiB,gBAAgB;AAC5C,uBAAW,MAAM;AACjB;AAAA,UACF;AACA,cAAI,CAAC,KAAM;AACX,qBAAW,SAAS,IAAI;AAAA,QAC1B;AAAA,MACF,UAAE;AACA,YAAI,CAAC,WAAW,OAAQ,YAAW,SAAS;AAAA,MAC9C;AAAA,IACF;AAEA,UAAM,WAAW,YAAY;AAC3B,YAAM,WAAW,EAAE,OAAO,SAAS,SAAS,gBAAgB,KAAK,KAAK,EAAE;AACxE,SAAI,KAAK,OAAO,SAAK,uBAAO,QAAQ,CAAC,CAAC;AAEtC,uBAAiB,MAAM,YAAY;AACjC,YAAI,KAAK,gBAAgB,OAAO,QAAS;AACzC,cAAM,WAAW,GAAG;AACpB,YAAI,CAAC,SAAU;AACf,aAAK,YAAY;AACjB,WAAI,KAAK,OAAO,SAAK,uBAAO,EAAE,OAAO,QAAQ,MAAM,WAAW,IAAI,CAAC,CAAC,CAAC;AACrE,WAAI,KAAK,OAAO,SAAK,uBAAO,EAAE,OAAO,QAAQ,CAAC,CAAC,CAAC;AAAA,MAClD;AAEA,UAAI,CAAC,KAAK,gBAAgB,OAAO,SAAS;AACxC,WAAI,KAAK,OAAO,SAAK,uBAAO,EAAE,OAAO,OAAO,CAAC,CAAC,CAAC;AAAA,MACjD;AAAA,IACF;AAEA,QAAI;AACJ,UAAM,gBAAgB,CAAC,UAAmB;AACxC,UAAI,WAAW;AACb,aAAK,MAAM,IAAI,EAAE,WAAW,WAAW,WAAW,OAAO,WAAW,MAAM,CAAC;AAC3E,oBAAY;AAAA,MACd;AAAA,IACF;AAEA,UAAM,WAAW,YAAY;AAG3B,YAAM,YAAY,CAAC,QAAiB;AAClC,YAAI;AACJ,YAAI,OAAO,SAAS,GAAG,GAAG;AACxB,kBAAQ;AAAA,QACV,WAAW,MAAM,QAAQ,GAAG,GAAG;AAC7B,kBAAQ,OAAO,OAAO,GAAG;AAAA,QAC3B,OAAO;AACL,kBAAQ,OAAO,KAAK,GAAkB;AAAA,QACxC;AAEA,YAAI;AACJ,YAAI;AACF,uBAAS,uBAAO,KAAK;AAAA,QACvB,SAAS,KAAK;AACZ,eAAK,QAAQ,KAAK,EAAE,IAAI,GAAG,qCAAqC;AAChE;AAAA,QACF;AAEA,cAAM,QAAQ,OAAO;AACrB,YAAI,UAAU,SAAS;AACrB,gBAAM,QAAQ,OAAO;AACrB,cAAI,SAAS,MAAM,aAAa,GAAG;AACjC,uBAAW,KAAK,QAAQ,MAAM,KAAK,GAAG;AACpC,4BAAc,KAAK;AACnB,0BAAY;AAAA,YACd;AAAA,UACF;AAAA,QACF,WAAW,UAAU,UAAU;AAC7B,gBAAM,SAAS,OAAO;AACtB,cAAI,WAAW,SAAS;AACtB,qBAAS;AAAA,cACP,IAAI,6BAAe;AAAA,gBACjB,SAAS;AAAA,gBACT,SAAS,EAAE,MAAM,EAAE,KAAK,KAAK,UAAU,MAAM,EAAE,EAAE;AAAA,cACnD,CAAC;AAAA,YACH;AACA;AAAA,UACF;AACA,qBAAW,KAAK,QAAQ,MAAM,GAAG;AAC/B,0BAAc,KAAK;AACnB,wBAAY;AAAA,UACd;AACA,wBAAc,IAAI;AAClB,cAAI,CAAC,KAAK,MAAM,QAAQ;AACtB,iBAAK,MAAM,IAAI,iBAAiB,aAAa;AAAA,UAC/C;AACA,cAAI,CAAC,SAAS,KAAM,UAAS,QAAQ;AAAA,QACvC,OAAO;AACL,eAAK,QAAQ,MAAM,EAAE,MAAM,GAAG,0BAA0B;AAAA,QAC1D;AAAA,MACF;AAEA,YAAM,UAAU,CAAC,MAAc,WAAmB;AAChD,YAAI,CAAC,SAAS,MAAM;AAClB,mBAAS;AAAA,YACP,IAAI,6BAAe;AAAA,cACjB,SAAS;AAAA,cACT,SAAS;AAAA,gBACP,YAAY,QAAQ;AAAA,gBACpB,MAAM,EAAE,QAAQ,OAAO,SAAS,EAAE;AAAA,cACpC;AAAA,YACF,CAAC;AAAA,UACH;AAAA,QACF;AAAA,MACF;AAEA,YAAM,UAAU,CAAC,QAAe;AAC9B,YAAI,CAAC,SAAS,KAAM,UAAS,OAAO,GAAG;AAAA,MACzC;AAEA,SAAI,GAAG,WAAW,SAAS;AAC3B,SAAI,GAAG,SAAS,OAAO;AACvB,SAAI,GAAG,SAAS,OAAO;AAEvB,UAAI;AACF,cAAM,SAAS;AAAA,MACjB,UAAE;AACA,WAAI,IAAI,WAAW,SAAS;AAC5B,WAAI,IAAI,SAAS,OAAO;AACxB,WAAI,IAAI,SAAS,OAAO;AAAA,MAC1B;AAAA,IACF;AAEA,QAAI;AACF,YAAM,QAAQ,IAAI,CAAC,UAAU,GAAG,SAAS,GAAG,SAAS,CAAC,CAAC;AAAA,IACzD,SAAS,GAAG;AACV,UAAI,KAAK,YAAY,QAAS;AAC9B,UAAI,aAAa,gCAAkB,aAAa,kCAAoB;AAClE,cAAM;AAAA,MACR;AACA,YAAM,IAAI,iCAAmB;AAAA,QAC3B,SAAS,gCAAiC,EAAY,WAAW,eAAe;AAAA,MAClF,CAAC;AAAA,IACH,UAAE;AACA,UAAI,CAAC,WAAW,OAAQ,YAAW,MAAM;AACzC,UAAI,MAAM,GAAG,eAAe,oBAAU,QAAQ;AAC5C,YAAI;AACF,aAAG,MAAM;AAAA,QACX,QAAQ;AAAA,QAER;AAAA,MACF;AAAA,IACF;AAAA,EACF;AACF;AAEA,MAAM,mBAAmB,OAAO;AAAA,EAC9B;AAAA,EACA;AAAA,EACA;AAAA,EACA;AACF,MAK0B;AACxB,QAAM,KAAK,IAAI,oBAAU,KAAK,EAAE,SAAS,kBAAkB,UAAU,CAAC;AACtE,QAAM,MAAM,IAAI,qBAAa;AAE7B,MAAI;AACJ,QAAM,UAAU,MAAM;AACpB,QAAI,QAAS,cAAa,OAAO;AACjC,OAAG,IAAI,QAAQ,MAAM;AACrB,OAAG,IAAI,SAAS,OAAO;AACvB,OAAG,IAAI,SAAS,OAAO;AACvB,gBAAY,oBAAoB,SAAS,OAAO;AAAA,EAClD;AAEA,QAAM,SAAS,MAAM,IAAI,QAAQ;AACjC,QAAM,UAAU,CAAC,QAAe,IAAI,OAAO,GAAG;AAC9C,QAAM,UAAU,CAAC,MAAc,WAC7B,IAAI;AAAA,IACF,IAAI,MAAM,sCAAsC,IAAI,YAAY,OAAO,SAAS,CAAC,GAAG;AAAA,EACtF;AACF,QAAM,UAAU,MAAM,IAAI,OAAO,IAAI,MAAM,SAAS,CAAC;AAErD,KAAG,GAAG,QAAQ,MAAM;AACpB,KAAG,GAAG,SAAS,OAAO;AACtB,KAAG,GAAG,SAAS,OAAO;AACtB,cAAY,iBAAiB,SAAS,SAAS,EAAE,MAAM,KAAK,CAAC;AAE7D,MAAI,YAAY,GAAG;AACjB,cAAU,WAAW,MAAM,IAAI,OAAO,IAAI,MAAM,iBAAiB,CAAC,GAAG,SAAS;AAAA,EAChF;AAEA,MAAI;AACF,UAAM,IAAI;AACV,WAAO;AAAA,EACT,SAAS,GAAG;AACV,QAAI;AACF,SAAG,GAAG,SAAS,MAAM;AAAA,MAAC,CAAC;AACvB,UAAI,GAAG,eAAe,oBAAU,YAAY;AAC1C,WAAG,MAAM;AAAA,MACX,OAAO;AACL,WAAG,UAAU;AAAA,MACf;AAAA,IACF,QAAQ;AAAA,IAER;AACA,UAAM;AAAA,EACR,UAAE;AACA,YAAQ;AAAA,EACV;AACF;","names":["tts"]} |
@@ -1,1 +0,1 @@ | ||
| {"version":3,"file":"tts.d.ts","sourceRoot":"","sources":["../src/tts.ts"],"names":[],"mappings":"AAGA,OAAO,EACL,KAAK,iBAAiB,EAOtB,QAAQ,EACR,GAAG,EACJ,MAAM,iBAAiB,CAAC;AAKzB,OAAO,KAAK,EAAE,WAAW,EAAE,SAAS,EAAE,MAAM,aAAa,CAAC;AAS1D,MAAM,WAAW,UAAU;IACzB,MAAM,CAAC,EAAE,MAAM,CAAC;IAChB,KAAK,CAAC,EAAE,SAAS,GAAG,MAAM,CAAC;IAC3B,OAAO,CAAC,EAAE,MAAM,CAAC;IACjB,UAAU,CAAC,EAAE,MAAM,CAAC;IACpB,OAAO,CAAC,EAAE,MAAM,CAAC;IACjB,WAAW,CAAC,EAAE,WAAW,CAAC;IAC1B;;;;;OAKG;IACH,WAAW,CAAC,EAAE,MAAM,CAAC;IACrB;;;OAGG;IACH,KAAK,CAAC,EAAE,MAAM,CAAC;IACf;;;OAGG;IACH,MAAM,CAAC,EAAE,MAAM,CAAC;IAChB,SAAS,CAAC,EAAE,QAAQ,CAAC,iBAAiB,CAAC;CACxC;AAED,UAAU,kBAAkB;IAC1B,MAAM,EAAE,MAAM,CAAC;IACf,KAAK,EAAE,SAAS,GAAG,MAAM,CAAC;IAC1B,OAAO,CAAC,EAAE,MAAM,CAAC;IACjB,UAAU,EAAE,MAAM,CAAC;IACnB,OAAO,EAAE,MAAM,CAAC;IAChB,WAAW,EAAE,WAAW,CAAC;IACzB,WAAW,EAAE,MAAM,CAAC;IACpB,KAAK,CAAC,EAAE,MAAM,CAAC;IACf,MAAM,CAAC,EAAE,MAAM,CAAC;IAChB,SAAS,EAAE,QAAQ,CAAC,iBAAiB,CAAC;CACvC;AAiDD,qBAAa,GAAI,SAAQ,GAAG,CAAC,GAAG;;IAE9B,KAAK,SAAmB;gBAEZ,IAAI,GAAE,UAAe;IAmCjC,IAAI,KAAK,IAAI,MAAM,CAElB;IAED,IAAI,QAAQ,IAAI,MAAM,CAErB;IAED,aAAa,CAAC,IAAI,EAAE;QAClB,KAAK,CAAC,EAAE,SAAS,GAAG,MAAM,CAAC;QAC3B,OAAO,CAAC,EAAE,MAAM,CAAC;QACjB,WAAW,CAAC,EAAE,WAAW,CAAC;QAC1B,WAAW,CAAC,EAAE,MAAM,CAAC;QACrB,KAAK,CAAC,EAAE,MAAM,CAAC;QACf,MAAM,CAAC,EAAE,MAAM,CAAC;KACjB,GAAG,IAAI;IAYR,UAAU,CACR,IAAI,EAAE,MAAM,EACZ,WAAW,CAAC,EAAE,iBAAiB,EAC/B,WAAW,CAAC,EAAE,WAAW,GACxB,GAAG,CAAC,aAAa;IAIpB,MAAM,CAAC,OAAO,CAAC,EAAE;QAAE,WAAW,CAAC,EAAE,iBAAiB,CAAA;KAAE,GAAG,GAAG,CAAC,gBAAgB;CAG5E;AAED,qBAAa,aAAc,SAAQ,GAAG,CAAC,aAAa;;IAClD,KAAK,SAA6B;gBAMhC,GAAG,EAAE,GAAG,EACR,IAAI,EAAE,MAAM,EACZ,IAAI,EAAE,kBAAkB,EACxB,WAAW,CAAC,EAAE,iBAAiB,EAC/B,WAAW,CAAC,EAAE,WAAW;cAOX,GAAG;CAiHpB;AAED,qBAAa,gBAAiB,SAAQ,GAAG,CAAC,gBAAgB;;IACxD,KAAK,SAAgC;gBAIzB,GAAG,EAAE,GAAG,EAAE,IAAI,EAAE,kBAAkB,EAAE,WAAW,CAAC,EAAE,iBAAiB;cAK/D,GAAG;CAmLpB"} | ||
| {"version":3,"file":"tts.d.ts","sourceRoot":"","sources":["../src/tts.ts"],"names":[],"mappings":"AAGA,OAAO,EACL,KAAK,iBAAiB,EAOtB,QAAQ,EACR,GAAG,EACJ,MAAM,iBAAiB,CAAC;AAKzB,OAAO,KAAK,EAAE,WAAW,EAAE,SAAS,EAAE,MAAM,aAAa,CAAC;AAS1D,MAAM,WAAW,UAAU;IACzB,MAAM,CAAC,EAAE,MAAM,CAAC;IAChB,KAAK,CAAC,EAAE,SAAS,GAAG,MAAM,CAAC;IAC3B,OAAO,CAAC,EAAE,MAAM,CAAC;IACjB,UAAU,CAAC,EAAE,MAAM,CAAC;IACpB,OAAO,CAAC,EAAE,MAAM,CAAC;IACjB,WAAW,CAAC,EAAE,WAAW,CAAC;IAC1B;;;;;OAKG;IACH,WAAW,CAAC,EAAE,MAAM,CAAC;IACrB;;;OAGG;IACH,KAAK,CAAC,EAAE,MAAM,CAAC;IACf;;;OAGG;IACH,MAAM,CAAC,EAAE,MAAM,CAAC;IAChB,SAAS,CAAC,EAAE,QAAQ,CAAC,iBAAiB,CAAC;CACxC;AAED,UAAU,kBAAkB;IAC1B,MAAM,EAAE,MAAM,CAAC;IACf,KAAK,EAAE,SAAS,GAAG,MAAM,CAAC;IAC1B,OAAO,CAAC,EAAE,MAAM,CAAC;IACjB,UAAU,EAAE,MAAM,CAAC;IACnB,OAAO,EAAE,MAAM,CAAC;IAChB,WAAW,EAAE,WAAW,CAAC;IACzB,WAAW,EAAE,MAAM,CAAC;IACpB,KAAK,CAAC,EAAE,MAAM,CAAC;IACf,MAAM,CAAC,EAAE,MAAM,CAAC;IAChB,SAAS,EAAE,QAAQ,CAAC,iBAAiB,CAAC;CACvC;AAiDD,qBAAa,GAAI,SAAQ,GAAG,CAAC,GAAG;;IAE9B,KAAK,SAAmB;gBAEZ,IAAI,GAAE,UAAe;IAmCjC,IAAI,KAAK,IAAI,MAAM,CAElB;IAED,IAAI,QAAQ,IAAI,MAAM,CAErB;IAED,aAAa,CAAC,IAAI,EAAE;QAClB,KAAK,CAAC,EAAE,SAAS,GAAG,MAAM,CAAC;QAC3B,OAAO,CAAC,EAAE,MAAM,CAAC;QACjB,WAAW,CAAC,EAAE,WAAW,CAAC;QAC1B,WAAW,CAAC,EAAE,MAAM,CAAC;QACrB,KAAK,CAAC,EAAE,MAAM,CAAC;QACf,MAAM,CAAC,EAAE,MAAM,CAAC;KACjB,GAAG,IAAI;IAYR,UAAU,CACR,IAAI,EAAE,MAAM,EACZ,WAAW,CAAC,EAAE,iBAAiB,EAC/B,WAAW,CAAC,EAAE,WAAW,GACxB,GAAG,CAAC,aAAa;IAIpB,MAAM,CAAC,OAAO,CAAC,EAAE;QAAE,WAAW,CAAC,EAAE,iBAAiB,CAAA;KAAE,GAAG,GAAG,CAAC,gBAAgB;CAG5E;AAED,qBAAa,aAAc,SAAQ,GAAG,CAAC,aAAa;;IAClD,KAAK,SAA6B;gBAMhC,GAAG,EAAE,GAAG,EACR,IAAI,EAAE,MAAM,EACZ,IAAI,EAAE,kBAAkB,EACxB,WAAW,CAAC,EAAE,iBAAiB,EAC/B,WAAW,CAAC,EAAE,WAAW;cAOX,GAAG;CAiHpB;AAED,qBAAa,gBAAiB,SAAQ,GAAG,CAAC,gBAAgB;;IACxD,KAAK,SAAgC;gBAIzB,GAAG,EAAE,GAAG,EAAE,IAAI,EAAE,kBAAkB,EAAE,WAAW,CAAC,EAAE,iBAAiB;cAK/D,GAAG;CAoLpB"} |
+1
-0
@@ -275,2 +275,3 @@ import { | ||
| if (!sentence) continue; | ||
| this.markStarted(); | ||
| ws.send(Buffer.from(encode({ event: "text", text: sentence + " " }))); | ||
@@ -277,0 +278,0 @@ ws.send(Buffer.from(encode({ event: "flush" }))); |
+1
-1
@@ -1,1 +0,1 @@ | ||
| {"version":3,"sources":["../src/tts.ts"],"sourcesContent":["// SPDX-FileCopyrightText: 2026 LiveKit, Inc.\n//\n// SPDX-License-Identifier: Apache-2.0\nimport {\n type APIConnectOptions,\n APIConnectionError,\n APIStatusError,\n AudioByteStream,\n Future,\n log,\n shortuuid,\n tokenize,\n tts,\n} from '@livekit/agents';\nimport type { AudioFrame } from '@livekit/rtc-node';\nimport { decode, encode } from '@msgpack/msgpack';\nimport { request } from 'node:https';\nimport { type RawData, WebSocket } from 'ws';\nimport type { LatencyMode, TTSModels } from './models.js';\n\nconst DEFAULT_MODEL: TTSModels = 's2.1-pro';\nconst DEFAULT_VOICE_ID = '933563129e564b19a115bedd57b7406a';\nconst DEFAULT_BASE_URL = 'https://api.fish.audio';\nconst NUM_CHANNELS = 1;\n// Fish Audio's default sample rate for raw PCM output.\nconst DEFAULT_SAMPLE_RATE = 24000;\n\nexport interface TTSOptions {\n apiKey?: string;\n model?: TTSModels | string;\n voiceId?: string;\n sampleRate?: number;\n baseURL?: string;\n latencyMode?: LatencyMode;\n /**\n * Upper bound on the number of characters Fish buffers before auto-synthesizing.\n * Must be between 100 and 300. With sentence-level flushing this is only hit by\n * sentences longer than `chunkLength`; otherwise audio is produced as soon as\n * each sentence is flushed. Defaults to 100.\n */\n chunkLength?: number;\n /**\n * Speaking rate multiplier for Fish `prosody.speed`. `1.0` is normal; below\n * 1.0 is slower, above is faster. Unset uses the voice's natural pace.\n */\n speed?: number;\n /**\n * Loudness adjustment in decibels for Fish `prosody.volume`. `0` is the\n * voice's natural level. Unset leaves it unchanged.\n */\n volume?: number;\n tokenizer?: tokenize.SentenceTokenizer;\n}\n\ninterface ResolvedTTSOptions {\n apiKey: string;\n model: TTSModels | string;\n voiceId?: string;\n sampleRate: number;\n baseURL: string;\n latencyMode: LatencyMode;\n chunkLength: number;\n speed?: number;\n volume?: number;\n tokenizer: tokenize.SentenceTokenizer;\n}\n\nconst DEFAULT_OPTS: Omit<ResolvedTTSOptions, 'apiKey' | 'tokenizer'> = {\n model: DEFAULT_MODEL,\n voiceId: DEFAULT_VOICE_ID,\n sampleRate: DEFAULT_SAMPLE_RATE,\n baseURL: DEFAULT_BASE_URL,\n latencyMode: 'balanced',\n chunkLength: 100,\n};\n\nconst validateChunkLength = (chunkLength: number) => {\n if (!Number.isFinite(chunkLength) || chunkLength < 100 || chunkLength > 300) {\n throw new Error('chunkLength must be between 100 and 300');\n }\n};\n\n// Fish Audio's wire format mirrors the upstream Python SDK so the server\n// doesn't fall back to its own larger defaults — in particular the docs default\n// of `chunk_length=300` produces large bursts that leave audible gaps between\n// chunk boundaries.\nconst buildTtsRequest = (opts: ResolvedTTSOptions, text: string = ''): Record<string, unknown> => {\n const prosody =\n opts.speed !== undefined || opts.volume !== undefined\n ? {\n ...(opts.speed !== undefined ? { speed: opts.speed } : {}),\n ...(opts.volume !== undefined ? { volume: opts.volume } : {}),\n }\n : null;\n\n return {\n text,\n chunk_length: opts.chunkLength,\n format: 'pcm',\n sample_rate: opts.sampleRate,\n mp3_bitrate: 64,\n opus_bitrate: 64000,\n references: [],\n // Fish Audio's wire field is `reference_id`; we expose it as `voiceId` on\n // the plugin for consistency with other TTS plugins.\n reference_id: opts.voiceId ?? null,\n normalize: true,\n latency: opts.latencyMode,\n prosody,\n top_p: 0.7,\n temperature: 0.7,\n };\n};\n\nexport class TTS extends tts.TTS {\n #opts: ResolvedTTSOptions;\n label = 'fishaudio.TTS';\n\n constructor(opts: TTSOptions = {}) {\n const apiKey = opts.apiKey ?? process.env.FISH_API_KEY;\n if (!apiKey) {\n throw new Error(\n 'Fish Audio API key is required, either as argument or set FISH_API_KEY environment variable',\n );\n }\n\n const chunkLength = opts.chunkLength ?? DEFAULT_OPTS.chunkLength;\n validateChunkLength(chunkLength);\n\n const sampleRate = opts.sampleRate ?? DEFAULT_OPTS.sampleRate;\n\n super(sampleRate, NUM_CHANNELS, { streaming: true });\n\n // min_sentence_len=1 emits each sentence as soon as the next one starts,\n // rather than batching short sentences together — minimizes TTFB on the\n // first sentence and keeps Fish synthesizing continuously.\n const tokenizer =\n opts.tokenizer ?? new tokenize.basic.SentenceTokenizer({ minSentenceLength: 1 });\n\n this.#opts = {\n apiKey,\n model: opts.model ?? DEFAULT_OPTS.model,\n voiceId: opts.voiceId ?? DEFAULT_OPTS.voiceId,\n sampleRate,\n baseURL: opts.baseURL ?? DEFAULT_OPTS.baseURL,\n latencyMode: opts.latencyMode ?? DEFAULT_OPTS.latencyMode,\n chunkLength,\n speed: opts.speed,\n volume: opts.volume,\n tokenizer,\n };\n }\n\n get model(): string {\n return this.#opts.model;\n }\n\n get provider(): string {\n return 'FishAudio';\n }\n\n updateOptions(opts: {\n model?: TTSModels | string;\n voiceId?: string;\n latencyMode?: LatencyMode;\n chunkLength?: number;\n speed?: number;\n volume?: number;\n }): void {\n if (opts.model !== undefined) this.#opts.model = opts.model;\n if (opts.voiceId !== undefined) this.#opts.voiceId = opts.voiceId;\n if (opts.latencyMode !== undefined) this.#opts.latencyMode = opts.latencyMode;\n if (opts.chunkLength !== undefined) {\n validateChunkLength(opts.chunkLength);\n this.#opts.chunkLength = opts.chunkLength;\n }\n if (opts.speed !== undefined) this.#opts.speed = opts.speed;\n if (opts.volume !== undefined) this.#opts.volume = opts.volume;\n }\n\n synthesize(\n text: string,\n connOptions?: APIConnectOptions,\n abortSignal?: AbortSignal,\n ): tts.ChunkedStream {\n return new ChunkedStream(this, text, this.#opts, connOptions, abortSignal);\n }\n\n stream(options?: { connOptions?: APIConnectOptions }): tts.SynthesizeStream {\n return new SynthesizeStream(this, this.#opts, options?.connOptions);\n }\n}\n\nexport class ChunkedStream extends tts.ChunkedStream {\n label = 'fishaudio.ChunkedStream';\n #logger = log();\n #opts: ResolvedTTSOptions;\n #text: string;\n\n constructor(\n tts: TTS,\n text: string,\n opts: ResolvedTTSOptions,\n connOptions?: APIConnectOptions,\n abortSignal?: AbortSignal,\n ) {\n super(text, tts, connOptions, abortSignal);\n this.#text = text;\n this.#opts = opts;\n }\n\n protected async run() {\n const requestId = shortuuid();\n const bstream = new AudioByteStream(this.#opts.sampleRate, NUM_CHANNELS);\n const payload = encode(buildTtsRequest(this.#opts, this.#text));\n\n const baseUrl = new URL(this.#opts.baseURL);\n const isHttps = baseUrl.protocol === 'https:';\n if (!isHttps) {\n // The plugin only supports https; fall back via Node's http module is\n // intentionally not implemented to keep the code path simple.\n throw new APIConnectionError({\n message: `Fish Audio base URL must use https (got ${this.#opts.baseURL})`,\n });\n }\n\n const doneFut = new Future<void>();\n\n const req = request(\n {\n hostname: baseUrl.hostname,\n port: parseInt(baseUrl.port) || 443,\n path: '/v1/tts',\n method: 'POST',\n headers: {\n Authorization: `Bearer ${this.#opts.apiKey}`,\n 'Content-Type': 'application/msgpack',\n model: this.#opts.model,\n 'Content-Length': payload.byteLength,\n },\n signal: this.abortSignal,\n },\n (res) => {\n const status = res.statusCode ?? -1;\n if (status < 200 || status >= 300) {\n const chunks: Buffer[] = [];\n res.on('data', (c: Buffer) => chunks.push(c));\n res.on('end', () => {\n const body = Buffer.concat(chunks).toString();\n if (!doneFut.done) {\n doneFut.reject(\n new APIStatusError({\n message: `Fish Audio TTS request failed: ${body}`,\n options: { statusCode: status, body: { raw: body } },\n }),\n );\n }\n });\n res.on('error', (err) => {\n if (err.message === 'aborted') return;\n this.#logger.error({ err }, 'Fish Audio TTS error response stream error');\n if (!doneFut.done) {\n doneFut.reject(\n new APIStatusError({\n message: `Fish Audio TTS request failed (status ${status})`,\n options: { statusCode: status },\n }),\n );\n }\n });\n return;\n }\n\n res.on('data', (chunk: Buffer) => {\n for (const frame of bstream.write(chunk)) {\n this.queue.put({\n requestId,\n segmentId: requestId,\n frame,\n final: false,\n });\n }\n });\n res.on('close', () => {\n for (const frame of bstream.flush()) {\n this.queue.put({\n requestId,\n segmentId: requestId,\n frame,\n final: false,\n });\n }\n if (!this.queue.closed) this.queue.close();\n if (!doneFut.done) doneFut.resolve();\n });\n res.on('error', (err) => {\n if (err.message === 'aborted') return;\n this.#logger.error({ err }, 'Fish Audio TTS response error');\n if (!doneFut.done) doneFut.reject(err);\n });\n },\n );\n\n req.on('error', (err) => {\n if (err.name === 'AbortError') return;\n this.#logger.error({ err }, 'Fish Audio TTS request error');\n if (!doneFut.done) doneFut.reject(err);\n });\n req.write(payload);\n req.end();\n\n try {\n await doneFut.await;\n } catch (e) {\n if (this.abortSignal.aborted) return;\n if (!this.queue.closed) this.queue.close();\n if (e instanceof APIStatusError || e instanceof APIConnectionError) {\n throw e;\n }\n throw new APIConnectionError({\n message: `Fish Audio connection failed: ${(e as Error).message ?? 'unknown error'}`,\n });\n }\n }\n}\n\nexport class SynthesizeStream extends tts.SynthesizeStream {\n label = 'fishaudio.SynthesizeStream';\n #logger = log();\n #opts: ResolvedTTSOptions;\n\n constructor(tts: TTS, opts: ResolvedTTSOptions, connOptions?: APIConnectOptions) {\n super(tts, connOptions);\n this.#opts = opts;\n }\n\n protected async run() {\n const requestId = shortuuid();\n const bstream = new AudioByteStream(this.#opts.sampleRate, NUM_CHANNELS);\n\n // Tokenize incoming text by sentence and flush after each sentence so Fish\n // synthesizes immediately at sentence boundaries instead of waiting for\n // `chunkLength` characters to accumulate. The result is much smoother\n // audio: gaps line up with sentence breaks (where pauses are natural)\n // rather than mid-clause.\n const sentStream = this.#opts.tokenizer.stream();\n\n const wsUrl = `${this.#opts.baseURL.replace(/^http/, 'ws')}/v1/tts/live`;\n let ws: WebSocket | undefined;\n try {\n ws = await connectWebSocket({\n url: wsUrl,\n headers: {\n Authorization: `Bearer ${this.#opts.apiKey}`,\n model: this.#opts.model,\n },\n timeoutMs: this.connOptions.timeoutMs,\n abortSignal: this.abortSignal,\n });\n } catch (e) {\n throw new APIConnectionError({\n message: `Fish Audio websocket connect failed: ${(e as Error).message ?? 'unknown error'}`,\n });\n }\n\n const finished = new Future<void>();\n\n const inputTask = async () => {\n try {\n for await (const data of this.input) {\n if (this.abortController.signal.aborted) break;\n if (data === SynthesizeStream.FLUSH_SENTINEL) {\n sentStream.flush();\n continue;\n }\n if (!data) continue;\n sentStream.pushText(data);\n }\n } finally {\n if (!sentStream.closed) sentStream.endInput();\n }\n };\n\n const sendTask = async () => {\n const startMsg = { event: 'start', request: buildTtsRequest(this.#opts) };\n ws!.send(Buffer.from(encode(startMsg)));\n\n for await (const ev of sentStream) {\n if (this.abortController.signal.aborted) break;\n const sentence = ev.token;\n if (!sentence) continue;\n ws!.send(Buffer.from(encode({ event: 'text', text: sentence + ' ' })));\n ws!.send(Buffer.from(encode({ event: 'flush' })));\n }\n\n if (!this.abortController.signal.aborted) {\n ws!.send(Buffer.from(encode({ event: 'stop' })));\n }\n };\n\n let lastFrame: AudioFrame | undefined;\n const sendLastFrame = (final: boolean) => {\n if (lastFrame) {\n this.queue.put({ requestId, segmentId: requestId, frame: lastFrame, final });\n lastFrame = undefined;\n }\n };\n\n const recvTask = async () => {\n // No per-receive timeout: Fish has natural inter-sentence gaps that can\n // exceed connOptions.timeoutMs when the LLM is slow.\n const onMessage = (raw: RawData) => {\n let frame: Buffer;\n if (Buffer.isBuffer(raw)) {\n frame = raw;\n } else if (Array.isArray(raw)) {\n frame = Buffer.concat(raw);\n } else {\n frame = Buffer.from(raw as ArrayBuffer);\n }\n\n let parsed: Record<string, unknown>;\n try {\n parsed = decode(frame) as Record<string, unknown>;\n } catch (err) {\n this.#logger.warn({ err }, 'Fish Audio failed to decode message');\n return;\n }\n\n const event = parsed.event as string | undefined;\n if (event === 'audio') {\n const audio = parsed.audio as Uint8Array | undefined;\n if (audio && audio.byteLength > 0) {\n for (const f of bstream.write(audio)) {\n sendLastFrame(false);\n lastFrame = f;\n }\n }\n } else if (event === 'finish') {\n const reason = parsed.reason as string | undefined;\n if (reason === 'error') {\n finished.reject(\n new APIStatusError({\n message: 'Fish Audio TTS reported an error',\n options: { body: { raw: JSON.stringify(parsed) } },\n }),\n );\n return;\n }\n for (const f of bstream.flush()) {\n sendLastFrame(false);\n lastFrame = f;\n }\n sendLastFrame(true);\n if (!this.queue.closed) {\n this.queue.put(SynthesizeStream.END_OF_STREAM);\n }\n if (!finished.done) finished.resolve();\n } else {\n this.#logger.debug({ event }, 'unknown Fish Audio event');\n }\n };\n\n const onClose = (code: number, reason: Buffer) => {\n if (!finished.done) {\n finished.reject(\n new APIStatusError({\n message: 'Fish Audio websocket connection closed unexpectedly',\n options: {\n statusCode: code || -1,\n body: { reason: reason.toString() },\n },\n }),\n );\n }\n };\n\n const onError = (err: Error) => {\n if (!finished.done) finished.reject(err);\n };\n\n ws!.on('message', onMessage);\n ws!.on('close', onClose);\n ws!.on('error', onError);\n\n try {\n await finished.await;\n } finally {\n ws!.off('message', onMessage);\n ws!.off('close', onClose);\n ws!.off('error', onError);\n }\n };\n\n try {\n await Promise.all([inputTask(), sendTask(), recvTask()]);\n } catch (e) {\n if (this.abortSignal.aborted) return;\n if (e instanceof APIStatusError || e instanceof APIConnectionError) {\n throw e;\n }\n throw new APIConnectionError({\n message: `Fish Audio websocket failed: ${(e as Error).message ?? 'unknown error'}`,\n });\n } finally {\n if (!sentStream.closed) sentStream.close();\n if (ws && ws.readyState !== WebSocket.CLOSED) {\n try {\n ws.close();\n } catch {\n // ignore\n }\n }\n }\n }\n}\n\nconst connectWebSocket = async ({\n url,\n headers,\n timeoutMs,\n abortSignal,\n}: {\n url: string;\n headers: Record<string, string>;\n timeoutMs: number;\n abortSignal: AbortSignal;\n}): Promise<WebSocket> => {\n const ws = new WebSocket(url, { headers, handshakeTimeout: timeoutMs });\n const fut = new Future<void>();\n\n let timeout: NodeJS.Timeout | undefined;\n const cleanup = () => {\n if (timeout) clearTimeout(timeout);\n ws.off('open', onOpen);\n ws.off('error', onError);\n ws.off('close', onClose);\n abortSignal.removeEventListener('abort', onAbort);\n };\n\n const onOpen = () => fut.resolve();\n const onError = (err: Error) => fut.reject(err);\n const onClose = (code: number, reason: Buffer) =>\n fut.reject(\n new Error(`websocket closed before open (code=${code}, reason=${reason.toString()})`),\n );\n const onAbort = () => fut.reject(new Error('aborted'));\n\n ws.on('open', onOpen);\n ws.on('error', onError);\n ws.on('close', onClose);\n abortSignal.addEventListener('abort', onAbort, { once: true });\n\n if (timeoutMs > 0) {\n timeout = setTimeout(() => fut.reject(new Error('connect timeout')), timeoutMs);\n }\n\n try {\n await fut.await;\n return ws;\n } catch (e) {\n try {\n ws.on('error', () => {});\n if (ws.readyState === WebSocket.CONNECTING) {\n ws.close();\n } else {\n ws.terminate();\n }\n } catch {\n // ignore\n }\n throw e;\n } finally {\n cleanup();\n }\n};\n"],"mappings":"AAGA;AAAA,EAEE;AAAA,EACA;AAAA,EACA;AAAA,EACA;AAAA,EACA;AAAA,EACA;AAAA,EACA;AAAA,EACA;AAAA,OACK;AAEP,SAAS,QAAQ,cAAc;AAC/B,SAAS,eAAe;AACxB,SAAuB,iBAAiB;AAGxC,MAAM,gBAA2B;AACjC,MAAM,mBAAmB;AACzB,MAAM,mBAAmB;AACzB,MAAM,eAAe;AAErB,MAAM,sBAAsB;AA0C5B,MAAM,eAAiE;AAAA,EACrE,OAAO;AAAA,EACP,SAAS;AAAA,EACT,YAAY;AAAA,EACZ,SAAS;AAAA,EACT,aAAa;AAAA,EACb,aAAa;AACf;AAEA,MAAM,sBAAsB,CAAC,gBAAwB;AACnD,MAAI,CAAC,OAAO,SAAS,WAAW,KAAK,cAAc,OAAO,cAAc,KAAK;AAC3E,UAAM,IAAI,MAAM,yCAAyC;AAAA,EAC3D;AACF;AAMA,MAAM,kBAAkB,CAAC,MAA0B,OAAe,OAAgC;AAChG,QAAM,UACJ,KAAK,UAAU,UAAa,KAAK,WAAW,SACxC;AAAA,IACE,GAAI,KAAK,UAAU,SAAY,EAAE,OAAO,KAAK,MAAM,IAAI,CAAC;AAAA,IACxD,GAAI,KAAK,WAAW,SAAY,EAAE,QAAQ,KAAK,OAAO,IAAI,CAAC;AAAA,EAC7D,IACA;AAEN,SAAO;AAAA,IACL;AAAA,IACA,cAAc,KAAK;AAAA,IACnB,QAAQ;AAAA,IACR,aAAa,KAAK;AAAA,IAClB,aAAa;AAAA,IACb,cAAc;AAAA,IACd,YAAY,CAAC;AAAA;AAAA;AAAA,IAGb,cAAc,KAAK,WAAW;AAAA,IAC9B,WAAW;AAAA,IACX,SAAS,KAAK;AAAA,IACd;AAAA,IACA,OAAO;AAAA,IACP,aAAa;AAAA,EACf;AACF;AAEO,MAAM,YAAY,IAAI,IAAI;AAAA,EAC/B;AAAA,EACA,QAAQ;AAAA,EAER,YAAY,OAAmB,CAAC,GAAG;AACjC,UAAM,SAAS,KAAK,UAAU,QAAQ,IAAI;AAC1C,QAAI,CAAC,QAAQ;AACX,YAAM,IAAI;AAAA,QACR;AAAA,MACF;AAAA,IACF;AAEA,UAAM,cAAc,KAAK,eAAe,aAAa;AACrD,wBAAoB,WAAW;AAE/B,UAAM,aAAa,KAAK,cAAc,aAAa;AAEnD,UAAM,YAAY,cAAc,EAAE,WAAW,KAAK,CAAC;AAKnD,UAAM,YACJ,KAAK,aAAa,IAAI,SAAS,MAAM,kBAAkB,EAAE,mBAAmB,EAAE,CAAC;AAEjF,SAAK,QAAQ;AAAA,MACX;AAAA,MACA,OAAO,KAAK,SAAS,aAAa;AAAA,MAClC,SAAS,KAAK,WAAW,aAAa;AAAA,MACtC;AAAA,MACA,SAAS,KAAK,WAAW,aAAa;AAAA,MACtC,aAAa,KAAK,eAAe,aAAa;AAAA,MAC9C;AAAA,MACA,OAAO,KAAK;AAAA,MACZ,QAAQ,KAAK;AAAA,MACb;AAAA,IACF;AAAA,EACF;AAAA,EAEA,IAAI,QAAgB;AAClB,WAAO,KAAK,MAAM;AAAA,EACpB;AAAA,EAEA,IAAI,WAAmB;AACrB,WAAO;AAAA,EACT;AAAA,EAEA,cAAc,MAOL;AACP,QAAI,KAAK,UAAU,OAAW,MAAK,MAAM,QAAQ,KAAK;AACtD,QAAI,KAAK,YAAY,OAAW,MAAK,MAAM,UAAU,KAAK;AAC1D,QAAI,KAAK,gBAAgB,OAAW,MAAK,MAAM,cAAc,KAAK;AAClE,QAAI,KAAK,gBAAgB,QAAW;AAClC,0BAAoB,KAAK,WAAW;AACpC,WAAK,MAAM,cAAc,KAAK;AAAA,IAChC;AACA,QAAI,KAAK,UAAU,OAAW,MAAK,MAAM,QAAQ,KAAK;AACtD,QAAI,KAAK,WAAW,OAAW,MAAK,MAAM,SAAS,KAAK;AAAA,EAC1D;AAAA,EAEA,WACE,MACA,aACA,aACmB;AACnB,WAAO,IAAI,cAAc,MAAM,MAAM,KAAK,OAAO,aAAa,WAAW;AAAA,EAC3E;AAAA,EAEA,OAAO,SAAqE;AAC1E,WAAO,IAAI,iBAAiB,MAAM,KAAK,OAAO,mCAAS,WAAW;AAAA,EACpE;AACF;AAEO,MAAM,sBAAsB,IAAI,cAAc;AAAA,EACnD,QAAQ;AAAA,EACR,UAAU,IAAI;AAAA,EACd;AAAA,EACA;AAAA,EAEA,YACEA,MACA,MACA,MACA,aACA,aACA;AACA,UAAM,MAAMA,MAAK,aAAa,WAAW;AACzC,SAAK,QAAQ;AACb,SAAK,QAAQ;AAAA,EACf;AAAA,EAEA,MAAgB,MAAM;AACpB,UAAM,YAAY,UAAU;AAC5B,UAAM,UAAU,IAAI,gBAAgB,KAAK,MAAM,YAAY,YAAY;AACvE,UAAM,UAAU,OAAO,gBAAgB,KAAK,OAAO,KAAK,KAAK,CAAC;AAE9D,UAAM,UAAU,IAAI,IAAI,KAAK,MAAM,OAAO;AAC1C,UAAM,UAAU,QAAQ,aAAa;AACrC,QAAI,CAAC,SAAS;AAGZ,YAAM,IAAI,mBAAmB;AAAA,QAC3B,SAAS,2CAA2C,KAAK,MAAM,OAAO;AAAA,MACxE,CAAC;AAAA,IACH;AAEA,UAAM,UAAU,IAAI,OAAa;AAEjC,UAAM,MAAM;AAAA,MACV;AAAA,QACE,UAAU,QAAQ;AAAA,QAClB,MAAM,SAAS,QAAQ,IAAI,KAAK;AAAA,QAChC,MAAM;AAAA,QACN,QAAQ;AAAA,QACR,SAAS;AAAA,UACP,eAAe,UAAU,KAAK,MAAM,MAAM;AAAA,UAC1C,gBAAgB;AAAA,UAChB,OAAO,KAAK,MAAM;AAAA,UAClB,kBAAkB,QAAQ;AAAA,QAC5B;AAAA,QACA,QAAQ,KAAK;AAAA,MACf;AAAA,MACA,CAAC,QAAQ;AACP,cAAM,SAAS,IAAI,cAAc;AACjC,YAAI,SAAS,OAAO,UAAU,KAAK;AACjC,gBAAM,SAAmB,CAAC;AAC1B,cAAI,GAAG,QAAQ,CAAC,MAAc,OAAO,KAAK,CAAC,CAAC;AAC5C,cAAI,GAAG,OAAO,MAAM;AAClB,kBAAM,OAAO,OAAO,OAAO,MAAM,EAAE,SAAS;AAC5C,gBAAI,CAAC,QAAQ,MAAM;AACjB,sBAAQ;AAAA,gBACN,IAAI,eAAe;AAAA,kBACjB,SAAS,kCAAkC,IAAI;AAAA,kBAC/C,SAAS,EAAE,YAAY,QAAQ,MAAM,EAAE,KAAK,KAAK,EAAE;AAAA,gBACrD,CAAC;AAAA,cACH;AAAA,YACF;AAAA,UACF,CAAC;AACD,cAAI,GAAG,SAAS,CAAC,QAAQ;AACvB,gBAAI,IAAI,YAAY,UAAW;AAC/B,iBAAK,QAAQ,MAAM,EAAE,IAAI,GAAG,4CAA4C;AACxE,gBAAI,CAAC,QAAQ,MAAM;AACjB,sBAAQ;AAAA,gBACN,IAAI,eAAe;AAAA,kBACjB,SAAS,yCAAyC,MAAM;AAAA,kBACxD,SAAS,EAAE,YAAY,OAAO;AAAA,gBAChC,CAAC;AAAA,cACH;AAAA,YACF;AAAA,UACF,CAAC;AACD;AAAA,QACF;AAEA,YAAI,GAAG,QAAQ,CAAC,UAAkB;AAChC,qBAAW,SAAS,QAAQ,MAAM,KAAK,GAAG;AACxC,iBAAK,MAAM,IAAI;AAAA,cACb;AAAA,cACA,WAAW;AAAA,cACX;AAAA,cACA,OAAO;AAAA,YACT,CAAC;AAAA,UACH;AAAA,QACF,CAAC;AACD,YAAI,GAAG,SAAS,MAAM;AACpB,qBAAW,SAAS,QAAQ,MAAM,GAAG;AACnC,iBAAK,MAAM,IAAI;AAAA,cACb;AAAA,cACA,WAAW;AAAA,cACX;AAAA,cACA,OAAO;AAAA,YACT,CAAC;AAAA,UACH;AACA,cAAI,CAAC,KAAK,MAAM,OAAQ,MAAK,MAAM,MAAM;AACzC,cAAI,CAAC,QAAQ,KAAM,SAAQ,QAAQ;AAAA,QACrC,CAAC;AACD,YAAI,GAAG,SAAS,CAAC,QAAQ;AACvB,cAAI,IAAI,YAAY,UAAW;AAC/B,eAAK,QAAQ,MAAM,EAAE,IAAI,GAAG,+BAA+B;AAC3D,cAAI,CAAC,QAAQ,KAAM,SAAQ,OAAO,GAAG;AAAA,QACvC,CAAC;AAAA,MACH;AAAA,IACF;AAEA,QAAI,GAAG,SAAS,CAAC,QAAQ;AACvB,UAAI,IAAI,SAAS,aAAc;AAC/B,WAAK,QAAQ,MAAM,EAAE,IAAI,GAAG,8BAA8B;AAC1D,UAAI,CAAC,QAAQ,KAAM,SAAQ,OAAO,GAAG;AAAA,IACvC,CAAC;AACD,QAAI,MAAM,OAAO;AACjB,QAAI,IAAI;AAER,QAAI;AACF,YAAM,QAAQ;AAAA,IAChB,SAAS,GAAG;AACV,UAAI,KAAK,YAAY,QAAS;AAC9B,UAAI,CAAC,KAAK,MAAM,OAAQ,MAAK,MAAM,MAAM;AACzC,UAAI,aAAa,kBAAkB,aAAa,oBAAoB;AAClE,cAAM;AAAA,MACR;AACA,YAAM,IAAI,mBAAmB;AAAA,QAC3B,SAAS,iCAAkC,EAAY,WAAW,eAAe;AAAA,MACnF,CAAC;AAAA,IACH;AAAA,EACF;AACF;AAEO,MAAM,yBAAyB,IAAI,iBAAiB;AAAA,EACzD,QAAQ;AAAA,EACR,UAAU,IAAI;AAAA,EACd;AAAA,EAEA,YAAYA,MAAU,MAA0B,aAAiC;AAC/E,UAAMA,MAAK,WAAW;AACtB,SAAK,QAAQ;AAAA,EACf;AAAA,EAEA,MAAgB,MAAM;AACpB,UAAM,YAAY,UAAU;AAC5B,UAAM,UAAU,IAAI,gBAAgB,KAAK,MAAM,YAAY,YAAY;AAOvE,UAAM,aAAa,KAAK,MAAM,UAAU,OAAO;AAE/C,UAAM,QAAQ,GAAG,KAAK,MAAM,QAAQ,QAAQ,SAAS,IAAI,CAAC;AAC1D,QAAI;AACJ,QAAI;AACF,WAAK,MAAM,iBAAiB;AAAA,QAC1B,KAAK;AAAA,QACL,SAAS;AAAA,UACP,eAAe,UAAU,KAAK,MAAM,MAAM;AAAA,UAC1C,OAAO,KAAK,MAAM;AAAA,QACpB;AAAA,QACA,WAAW,KAAK,YAAY;AAAA,QAC5B,aAAa,KAAK;AAAA,MACpB,CAAC;AAAA,IACH,SAAS,GAAG;AACV,YAAM,IAAI,mBAAmB;AAAA,QAC3B,SAAS,wCAAyC,EAAY,WAAW,eAAe;AAAA,MAC1F,CAAC;AAAA,IACH;AAEA,UAAM,WAAW,IAAI,OAAa;AAElC,UAAM,YAAY,YAAY;AAC5B,UAAI;AACF,yBAAiB,QAAQ,KAAK,OAAO;AACnC,cAAI,KAAK,gBAAgB,OAAO,QAAS;AACzC,cAAI,SAAS,iBAAiB,gBAAgB;AAC5C,uBAAW,MAAM;AACjB;AAAA,UACF;AACA,cAAI,CAAC,KAAM;AACX,qBAAW,SAAS,IAAI;AAAA,QAC1B;AAAA,MACF,UAAE;AACA,YAAI,CAAC,WAAW,OAAQ,YAAW,SAAS;AAAA,MAC9C;AAAA,IACF;AAEA,UAAM,WAAW,YAAY;AAC3B,YAAM,WAAW,EAAE,OAAO,SAAS,SAAS,gBAAgB,KAAK,KAAK,EAAE;AACxE,SAAI,KAAK,OAAO,KAAK,OAAO,QAAQ,CAAC,CAAC;AAEtC,uBAAiB,MAAM,YAAY;AACjC,YAAI,KAAK,gBAAgB,OAAO,QAAS;AACzC,cAAM,WAAW,GAAG;AACpB,YAAI,CAAC,SAAU;AACf,WAAI,KAAK,OAAO,KAAK,OAAO,EAAE,OAAO,QAAQ,MAAM,WAAW,IAAI,CAAC,CAAC,CAAC;AACrE,WAAI,KAAK,OAAO,KAAK,OAAO,EAAE,OAAO,QAAQ,CAAC,CAAC,CAAC;AAAA,MAClD;AAEA,UAAI,CAAC,KAAK,gBAAgB,OAAO,SAAS;AACxC,WAAI,KAAK,OAAO,KAAK,OAAO,EAAE,OAAO,OAAO,CAAC,CAAC,CAAC;AAAA,MACjD;AAAA,IACF;AAEA,QAAI;AACJ,UAAM,gBAAgB,CAAC,UAAmB;AACxC,UAAI,WAAW;AACb,aAAK,MAAM,IAAI,EAAE,WAAW,WAAW,WAAW,OAAO,WAAW,MAAM,CAAC;AAC3E,oBAAY;AAAA,MACd;AAAA,IACF;AAEA,UAAM,WAAW,YAAY;AAG3B,YAAM,YAAY,CAAC,QAAiB;AAClC,YAAI;AACJ,YAAI,OAAO,SAAS,GAAG,GAAG;AACxB,kBAAQ;AAAA,QACV,WAAW,MAAM,QAAQ,GAAG,GAAG;AAC7B,kBAAQ,OAAO,OAAO,GAAG;AAAA,QAC3B,OAAO;AACL,kBAAQ,OAAO,KAAK,GAAkB;AAAA,QACxC;AAEA,YAAI;AACJ,YAAI;AACF,mBAAS,OAAO,KAAK;AAAA,QACvB,SAAS,KAAK;AACZ,eAAK,QAAQ,KAAK,EAAE,IAAI,GAAG,qCAAqC;AAChE;AAAA,QACF;AAEA,cAAM,QAAQ,OAAO;AACrB,YAAI,UAAU,SAAS;AACrB,gBAAM,QAAQ,OAAO;AACrB,cAAI,SAAS,MAAM,aAAa,GAAG;AACjC,uBAAW,KAAK,QAAQ,MAAM,KAAK,GAAG;AACpC,4BAAc,KAAK;AACnB,0BAAY;AAAA,YACd;AAAA,UACF;AAAA,QACF,WAAW,UAAU,UAAU;AAC7B,gBAAM,SAAS,OAAO;AACtB,cAAI,WAAW,SAAS;AACtB,qBAAS;AAAA,cACP,IAAI,eAAe;AAAA,gBACjB,SAAS;AAAA,gBACT,SAAS,EAAE,MAAM,EAAE,KAAK,KAAK,UAAU,MAAM,EAAE,EAAE;AAAA,cACnD,CAAC;AAAA,YACH;AACA;AAAA,UACF;AACA,qBAAW,KAAK,QAAQ,MAAM,GAAG;AAC/B,0BAAc,KAAK;AACnB,wBAAY;AAAA,UACd;AACA,wBAAc,IAAI;AAClB,cAAI,CAAC,KAAK,MAAM,QAAQ;AACtB,iBAAK,MAAM,IAAI,iBAAiB,aAAa;AAAA,UAC/C;AACA,cAAI,CAAC,SAAS,KAAM,UAAS,QAAQ;AAAA,QACvC,OAAO;AACL,eAAK,QAAQ,MAAM,EAAE,MAAM,GAAG,0BAA0B;AAAA,QAC1D;AAAA,MACF;AAEA,YAAM,UAAU,CAAC,MAAc,WAAmB;AAChD,YAAI,CAAC,SAAS,MAAM;AAClB,mBAAS;AAAA,YACP,IAAI,eAAe;AAAA,cACjB,SAAS;AAAA,cACT,SAAS;AAAA,gBACP,YAAY,QAAQ;AAAA,gBACpB,MAAM,EAAE,QAAQ,OAAO,SAAS,EAAE;AAAA,cACpC;AAAA,YACF,CAAC;AAAA,UACH;AAAA,QACF;AAAA,MACF;AAEA,YAAM,UAAU,CAAC,QAAe;AAC9B,YAAI,CAAC,SAAS,KAAM,UAAS,OAAO,GAAG;AAAA,MACzC;AAEA,SAAI,GAAG,WAAW,SAAS;AAC3B,SAAI,GAAG,SAAS,OAAO;AACvB,SAAI,GAAG,SAAS,OAAO;AAEvB,UAAI;AACF,cAAM,SAAS;AAAA,MACjB,UAAE;AACA,WAAI,IAAI,WAAW,SAAS;AAC5B,WAAI,IAAI,SAAS,OAAO;AACxB,WAAI,IAAI,SAAS,OAAO;AAAA,MAC1B;AAAA,IACF;AAEA,QAAI;AACF,YAAM,QAAQ,IAAI,CAAC,UAAU,GAAG,SAAS,GAAG,SAAS,CAAC,CAAC;AAAA,IACzD,SAAS,GAAG;AACV,UAAI,KAAK,YAAY,QAAS;AAC9B,UAAI,aAAa,kBAAkB,aAAa,oBAAoB;AAClE,cAAM;AAAA,MACR;AACA,YAAM,IAAI,mBAAmB;AAAA,QAC3B,SAAS,gCAAiC,EAAY,WAAW,eAAe;AAAA,MAClF,CAAC;AAAA,IACH,UAAE;AACA,UAAI,CAAC,WAAW,OAAQ,YAAW,MAAM;AACzC,UAAI,MAAM,GAAG,eAAe,UAAU,QAAQ;AAC5C,YAAI;AACF,aAAG,MAAM;AAAA,QACX,QAAQ;AAAA,QAER;AAAA,MACF;AAAA,IACF;AAAA,EACF;AACF;AAEA,MAAM,mBAAmB,OAAO;AAAA,EAC9B;AAAA,EACA;AAAA,EACA;AAAA,EACA;AACF,MAK0B;AACxB,QAAM,KAAK,IAAI,UAAU,KAAK,EAAE,SAAS,kBAAkB,UAAU,CAAC;AACtE,QAAM,MAAM,IAAI,OAAa;AAE7B,MAAI;AACJ,QAAM,UAAU,MAAM;AACpB,QAAI,QAAS,cAAa,OAAO;AACjC,OAAG,IAAI,QAAQ,MAAM;AACrB,OAAG,IAAI,SAAS,OAAO;AACvB,OAAG,IAAI,SAAS,OAAO;AACvB,gBAAY,oBAAoB,SAAS,OAAO;AAAA,EAClD;AAEA,QAAM,SAAS,MAAM,IAAI,QAAQ;AACjC,QAAM,UAAU,CAAC,QAAe,IAAI,OAAO,GAAG;AAC9C,QAAM,UAAU,CAAC,MAAc,WAC7B,IAAI;AAAA,IACF,IAAI,MAAM,sCAAsC,IAAI,YAAY,OAAO,SAAS,CAAC,GAAG;AAAA,EACtF;AACF,QAAM,UAAU,MAAM,IAAI,OAAO,IAAI,MAAM,SAAS,CAAC;AAErD,KAAG,GAAG,QAAQ,MAAM;AACpB,KAAG,GAAG,SAAS,OAAO;AACtB,KAAG,GAAG,SAAS,OAAO;AACtB,cAAY,iBAAiB,SAAS,SAAS,EAAE,MAAM,KAAK,CAAC;AAE7D,MAAI,YAAY,GAAG;AACjB,cAAU,WAAW,MAAM,IAAI,OAAO,IAAI,MAAM,iBAAiB,CAAC,GAAG,SAAS;AAAA,EAChF;AAEA,MAAI;AACF,UAAM,IAAI;AACV,WAAO;AAAA,EACT,SAAS,GAAG;AACV,QAAI;AACF,SAAG,GAAG,SAAS,MAAM;AAAA,MAAC,CAAC;AACvB,UAAI,GAAG,eAAe,UAAU,YAAY;AAC1C,WAAG,MAAM;AAAA,MACX,OAAO;AACL,WAAG,UAAU;AAAA,MACf;AAAA,IACF,QAAQ;AAAA,IAER;AACA,UAAM;AAAA,EACR,UAAE;AACA,YAAQ;AAAA,EACV;AACF;","names":["tts"]} | ||
| {"version":3,"sources":["../src/tts.ts"],"sourcesContent":["// SPDX-FileCopyrightText: 2026 LiveKit, Inc.\n//\n// SPDX-License-Identifier: Apache-2.0\nimport {\n type APIConnectOptions,\n APIConnectionError,\n APIStatusError,\n AudioByteStream,\n Future,\n log,\n shortuuid,\n tokenize,\n tts,\n} from '@livekit/agents';\nimport type { AudioFrame } from '@livekit/rtc-node';\nimport { decode, encode } from '@msgpack/msgpack';\nimport { request } from 'node:https';\nimport { type RawData, WebSocket } from 'ws';\nimport type { LatencyMode, TTSModels } from './models.js';\n\nconst DEFAULT_MODEL: TTSModels = 's2.1-pro';\nconst DEFAULT_VOICE_ID = '933563129e564b19a115bedd57b7406a';\nconst DEFAULT_BASE_URL = 'https://api.fish.audio';\nconst NUM_CHANNELS = 1;\n// Fish Audio's default sample rate for raw PCM output.\nconst DEFAULT_SAMPLE_RATE = 24000;\n\nexport interface TTSOptions {\n apiKey?: string;\n model?: TTSModels | string;\n voiceId?: string;\n sampleRate?: number;\n baseURL?: string;\n latencyMode?: LatencyMode;\n /**\n * Upper bound on the number of characters Fish buffers before auto-synthesizing.\n * Must be between 100 and 300. With sentence-level flushing this is only hit by\n * sentences longer than `chunkLength`; otherwise audio is produced as soon as\n * each sentence is flushed. Defaults to 100.\n */\n chunkLength?: number;\n /**\n * Speaking rate multiplier for Fish `prosody.speed`. `1.0` is normal; below\n * 1.0 is slower, above is faster. Unset uses the voice's natural pace.\n */\n speed?: number;\n /**\n * Loudness adjustment in decibels for Fish `prosody.volume`. `0` is the\n * voice's natural level. Unset leaves it unchanged.\n */\n volume?: number;\n tokenizer?: tokenize.SentenceTokenizer;\n}\n\ninterface ResolvedTTSOptions {\n apiKey: string;\n model: TTSModels | string;\n voiceId?: string;\n sampleRate: number;\n baseURL: string;\n latencyMode: LatencyMode;\n chunkLength: number;\n speed?: number;\n volume?: number;\n tokenizer: tokenize.SentenceTokenizer;\n}\n\nconst DEFAULT_OPTS: Omit<ResolvedTTSOptions, 'apiKey' | 'tokenizer'> = {\n model: DEFAULT_MODEL,\n voiceId: DEFAULT_VOICE_ID,\n sampleRate: DEFAULT_SAMPLE_RATE,\n baseURL: DEFAULT_BASE_URL,\n latencyMode: 'balanced',\n chunkLength: 100,\n};\n\nconst validateChunkLength = (chunkLength: number) => {\n if (!Number.isFinite(chunkLength) || chunkLength < 100 || chunkLength > 300) {\n throw new Error('chunkLength must be between 100 and 300');\n }\n};\n\n// Fish Audio's wire format mirrors the upstream Python SDK so the server\n// doesn't fall back to its own larger defaults — in particular the docs default\n// of `chunk_length=300` produces large bursts that leave audible gaps between\n// chunk boundaries.\nconst buildTtsRequest = (opts: ResolvedTTSOptions, text: string = ''): Record<string, unknown> => {\n const prosody =\n opts.speed !== undefined || opts.volume !== undefined\n ? {\n ...(opts.speed !== undefined ? { speed: opts.speed } : {}),\n ...(opts.volume !== undefined ? { volume: opts.volume } : {}),\n }\n : null;\n\n return {\n text,\n chunk_length: opts.chunkLength,\n format: 'pcm',\n sample_rate: opts.sampleRate,\n mp3_bitrate: 64,\n opus_bitrate: 64000,\n references: [],\n // Fish Audio's wire field is `reference_id`; we expose it as `voiceId` on\n // the plugin for consistency with other TTS plugins.\n reference_id: opts.voiceId ?? null,\n normalize: true,\n latency: opts.latencyMode,\n prosody,\n top_p: 0.7,\n temperature: 0.7,\n };\n};\n\nexport class TTS extends tts.TTS {\n #opts: ResolvedTTSOptions;\n label = 'fishaudio.TTS';\n\n constructor(opts: TTSOptions = {}) {\n const apiKey = opts.apiKey ?? process.env.FISH_API_KEY;\n if (!apiKey) {\n throw new Error(\n 'Fish Audio API key is required, either as argument or set FISH_API_KEY environment variable',\n );\n }\n\n const chunkLength = opts.chunkLength ?? DEFAULT_OPTS.chunkLength;\n validateChunkLength(chunkLength);\n\n const sampleRate = opts.sampleRate ?? DEFAULT_OPTS.sampleRate;\n\n super(sampleRate, NUM_CHANNELS, { streaming: true });\n\n // min_sentence_len=1 emits each sentence as soon as the next one starts,\n // rather than batching short sentences together — minimizes TTFB on the\n // first sentence and keeps Fish synthesizing continuously.\n const tokenizer =\n opts.tokenizer ?? new tokenize.basic.SentenceTokenizer({ minSentenceLength: 1 });\n\n this.#opts = {\n apiKey,\n model: opts.model ?? DEFAULT_OPTS.model,\n voiceId: opts.voiceId ?? DEFAULT_OPTS.voiceId,\n sampleRate,\n baseURL: opts.baseURL ?? DEFAULT_OPTS.baseURL,\n latencyMode: opts.latencyMode ?? DEFAULT_OPTS.latencyMode,\n chunkLength,\n speed: opts.speed,\n volume: opts.volume,\n tokenizer,\n };\n }\n\n get model(): string {\n return this.#opts.model;\n }\n\n get provider(): string {\n return 'FishAudio';\n }\n\n updateOptions(opts: {\n model?: TTSModels | string;\n voiceId?: string;\n latencyMode?: LatencyMode;\n chunkLength?: number;\n speed?: number;\n volume?: number;\n }): void {\n if (opts.model !== undefined) this.#opts.model = opts.model;\n if (opts.voiceId !== undefined) this.#opts.voiceId = opts.voiceId;\n if (opts.latencyMode !== undefined) this.#opts.latencyMode = opts.latencyMode;\n if (opts.chunkLength !== undefined) {\n validateChunkLength(opts.chunkLength);\n this.#opts.chunkLength = opts.chunkLength;\n }\n if (opts.speed !== undefined) this.#opts.speed = opts.speed;\n if (opts.volume !== undefined) this.#opts.volume = opts.volume;\n }\n\n synthesize(\n text: string,\n connOptions?: APIConnectOptions,\n abortSignal?: AbortSignal,\n ): tts.ChunkedStream {\n return new ChunkedStream(this, text, this.#opts, connOptions, abortSignal);\n }\n\n stream(options?: { connOptions?: APIConnectOptions }): tts.SynthesizeStream {\n return new SynthesizeStream(this, this.#opts, options?.connOptions);\n }\n}\n\nexport class ChunkedStream extends tts.ChunkedStream {\n label = 'fishaudio.ChunkedStream';\n #logger = log();\n #opts: ResolvedTTSOptions;\n #text: string;\n\n constructor(\n tts: TTS,\n text: string,\n opts: ResolvedTTSOptions,\n connOptions?: APIConnectOptions,\n abortSignal?: AbortSignal,\n ) {\n super(text, tts, connOptions, abortSignal);\n this.#text = text;\n this.#opts = opts;\n }\n\n protected async run() {\n const requestId = shortuuid();\n const bstream = new AudioByteStream(this.#opts.sampleRate, NUM_CHANNELS);\n const payload = encode(buildTtsRequest(this.#opts, this.#text));\n\n const baseUrl = new URL(this.#opts.baseURL);\n const isHttps = baseUrl.protocol === 'https:';\n if (!isHttps) {\n // The plugin only supports https; fall back via Node's http module is\n // intentionally not implemented to keep the code path simple.\n throw new APIConnectionError({\n message: `Fish Audio base URL must use https (got ${this.#opts.baseURL})`,\n });\n }\n\n const doneFut = new Future<void>();\n\n const req = request(\n {\n hostname: baseUrl.hostname,\n port: parseInt(baseUrl.port) || 443,\n path: '/v1/tts',\n method: 'POST',\n headers: {\n Authorization: `Bearer ${this.#opts.apiKey}`,\n 'Content-Type': 'application/msgpack',\n model: this.#opts.model,\n 'Content-Length': payload.byteLength,\n },\n signal: this.abortSignal,\n },\n (res) => {\n const status = res.statusCode ?? -1;\n if (status < 200 || status >= 300) {\n const chunks: Buffer[] = [];\n res.on('data', (c: Buffer) => chunks.push(c));\n res.on('end', () => {\n const body = Buffer.concat(chunks).toString();\n if (!doneFut.done) {\n doneFut.reject(\n new APIStatusError({\n message: `Fish Audio TTS request failed: ${body}`,\n options: { statusCode: status, body: { raw: body } },\n }),\n );\n }\n });\n res.on('error', (err) => {\n if (err.message === 'aborted') return;\n this.#logger.error({ err }, 'Fish Audio TTS error response stream error');\n if (!doneFut.done) {\n doneFut.reject(\n new APIStatusError({\n message: `Fish Audio TTS request failed (status ${status})`,\n options: { statusCode: status },\n }),\n );\n }\n });\n return;\n }\n\n res.on('data', (chunk: Buffer) => {\n for (const frame of bstream.write(chunk)) {\n this.queue.put({\n requestId,\n segmentId: requestId,\n frame,\n final: false,\n });\n }\n });\n res.on('close', () => {\n for (const frame of bstream.flush()) {\n this.queue.put({\n requestId,\n segmentId: requestId,\n frame,\n final: false,\n });\n }\n if (!this.queue.closed) this.queue.close();\n if (!doneFut.done) doneFut.resolve();\n });\n res.on('error', (err) => {\n if (err.message === 'aborted') return;\n this.#logger.error({ err }, 'Fish Audio TTS response error');\n if (!doneFut.done) doneFut.reject(err);\n });\n },\n );\n\n req.on('error', (err) => {\n if (err.name === 'AbortError') return;\n this.#logger.error({ err }, 'Fish Audio TTS request error');\n if (!doneFut.done) doneFut.reject(err);\n });\n req.write(payload);\n req.end();\n\n try {\n await doneFut.await;\n } catch (e) {\n if (this.abortSignal.aborted) return;\n if (!this.queue.closed) this.queue.close();\n if (e instanceof APIStatusError || e instanceof APIConnectionError) {\n throw e;\n }\n throw new APIConnectionError({\n message: `Fish Audio connection failed: ${(e as Error).message ?? 'unknown error'}`,\n });\n }\n }\n}\n\nexport class SynthesizeStream extends tts.SynthesizeStream {\n label = 'fishaudio.SynthesizeStream';\n #logger = log();\n #opts: ResolvedTTSOptions;\n\n constructor(tts: TTS, opts: ResolvedTTSOptions, connOptions?: APIConnectOptions) {\n super(tts, connOptions);\n this.#opts = opts;\n }\n\n protected async run() {\n const requestId = shortuuid();\n const bstream = new AudioByteStream(this.#opts.sampleRate, NUM_CHANNELS);\n\n // Tokenize incoming text by sentence and flush after each sentence so Fish\n // synthesizes immediately at sentence boundaries instead of waiting for\n // `chunkLength` characters to accumulate. The result is much smoother\n // audio: gaps line up with sentence breaks (where pauses are natural)\n // rather than mid-clause.\n const sentStream = this.#opts.tokenizer.stream();\n\n const wsUrl = `${this.#opts.baseURL.replace(/^http/, 'ws')}/v1/tts/live`;\n let ws: WebSocket | undefined;\n try {\n ws = await connectWebSocket({\n url: wsUrl,\n headers: {\n Authorization: `Bearer ${this.#opts.apiKey}`,\n model: this.#opts.model,\n },\n timeoutMs: this.connOptions.timeoutMs,\n abortSignal: this.abortSignal,\n });\n } catch (e) {\n throw new APIConnectionError({\n message: `Fish Audio websocket connect failed: ${(e as Error).message ?? 'unknown error'}`,\n });\n }\n\n const finished = new Future<void>();\n\n const inputTask = async () => {\n try {\n for await (const data of this.input) {\n if (this.abortController.signal.aborted) break;\n if (data === SynthesizeStream.FLUSH_SENTINEL) {\n sentStream.flush();\n continue;\n }\n if (!data) continue;\n sentStream.pushText(data);\n }\n } finally {\n if (!sentStream.closed) sentStream.endInput();\n }\n };\n\n const sendTask = async () => {\n const startMsg = { event: 'start', request: buildTtsRequest(this.#opts) };\n ws!.send(Buffer.from(encode(startMsg)));\n\n for await (const ev of sentStream) {\n if (this.abortController.signal.aborted) break;\n const sentence = ev.token;\n if (!sentence) continue;\n this.markStarted();\n ws!.send(Buffer.from(encode({ event: 'text', text: sentence + ' ' })));\n ws!.send(Buffer.from(encode({ event: 'flush' })));\n }\n\n if (!this.abortController.signal.aborted) {\n ws!.send(Buffer.from(encode({ event: 'stop' })));\n }\n };\n\n let lastFrame: AudioFrame | undefined;\n const sendLastFrame = (final: boolean) => {\n if (lastFrame) {\n this.queue.put({ requestId, segmentId: requestId, frame: lastFrame, final });\n lastFrame = undefined;\n }\n };\n\n const recvTask = async () => {\n // No per-receive timeout: Fish has natural inter-sentence gaps that can\n // exceed connOptions.timeoutMs when the LLM is slow.\n const onMessage = (raw: RawData) => {\n let frame: Buffer;\n if (Buffer.isBuffer(raw)) {\n frame = raw;\n } else if (Array.isArray(raw)) {\n frame = Buffer.concat(raw);\n } else {\n frame = Buffer.from(raw as ArrayBuffer);\n }\n\n let parsed: Record<string, unknown>;\n try {\n parsed = decode(frame) as Record<string, unknown>;\n } catch (err) {\n this.#logger.warn({ err }, 'Fish Audio failed to decode message');\n return;\n }\n\n const event = parsed.event as string | undefined;\n if (event === 'audio') {\n const audio = parsed.audio as Uint8Array | undefined;\n if (audio && audio.byteLength > 0) {\n for (const f of bstream.write(audio)) {\n sendLastFrame(false);\n lastFrame = f;\n }\n }\n } else if (event === 'finish') {\n const reason = parsed.reason as string | undefined;\n if (reason === 'error') {\n finished.reject(\n new APIStatusError({\n message: 'Fish Audio TTS reported an error',\n options: { body: { raw: JSON.stringify(parsed) } },\n }),\n );\n return;\n }\n for (const f of bstream.flush()) {\n sendLastFrame(false);\n lastFrame = f;\n }\n sendLastFrame(true);\n if (!this.queue.closed) {\n this.queue.put(SynthesizeStream.END_OF_STREAM);\n }\n if (!finished.done) finished.resolve();\n } else {\n this.#logger.debug({ event }, 'unknown Fish Audio event');\n }\n };\n\n const onClose = (code: number, reason: Buffer) => {\n if (!finished.done) {\n finished.reject(\n new APIStatusError({\n message: 'Fish Audio websocket connection closed unexpectedly',\n options: {\n statusCode: code || -1,\n body: { reason: reason.toString() },\n },\n }),\n );\n }\n };\n\n const onError = (err: Error) => {\n if (!finished.done) finished.reject(err);\n };\n\n ws!.on('message', onMessage);\n ws!.on('close', onClose);\n ws!.on('error', onError);\n\n try {\n await finished.await;\n } finally {\n ws!.off('message', onMessage);\n ws!.off('close', onClose);\n ws!.off('error', onError);\n }\n };\n\n try {\n await Promise.all([inputTask(), sendTask(), recvTask()]);\n } catch (e) {\n if (this.abortSignal.aborted) return;\n if (e instanceof APIStatusError || e instanceof APIConnectionError) {\n throw e;\n }\n throw new APIConnectionError({\n message: `Fish Audio websocket failed: ${(e as Error).message ?? 'unknown error'}`,\n });\n } finally {\n if (!sentStream.closed) sentStream.close();\n if (ws && ws.readyState !== WebSocket.CLOSED) {\n try {\n ws.close();\n } catch {\n // ignore\n }\n }\n }\n }\n}\n\nconst connectWebSocket = async ({\n url,\n headers,\n timeoutMs,\n abortSignal,\n}: {\n url: string;\n headers: Record<string, string>;\n timeoutMs: number;\n abortSignal: AbortSignal;\n}): Promise<WebSocket> => {\n const ws = new WebSocket(url, { headers, handshakeTimeout: timeoutMs });\n const fut = new Future<void>();\n\n let timeout: NodeJS.Timeout | undefined;\n const cleanup = () => {\n if (timeout) clearTimeout(timeout);\n ws.off('open', onOpen);\n ws.off('error', onError);\n ws.off('close', onClose);\n abortSignal.removeEventListener('abort', onAbort);\n };\n\n const onOpen = () => fut.resolve();\n const onError = (err: Error) => fut.reject(err);\n const onClose = (code: number, reason: Buffer) =>\n fut.reject(\n new Error(`websocket closed before open (code=${code}, reason=${reason.toString()})`),\n );\n const onAbort = () => fut.reject(new Error('aborted'));\n\n ws.on('open', onOpen);\n ws.on('error', onError);\n ws.on('close', onClose);\n abortSignal.addEventListener('abort', onAbort, { once: true });\n\n if (timeoutMs > 0) {\n timeout = setTimeout(() => fut.reject(new Error('connect timeout')), timeoutMs);\n }\n\n try {\n await fut.await;\n return ws;\n } catch (e) {\n try {\n ws.on('error', () => {});\n if (ws.readyState === WebSocket.CONNECTING) {\n ws.close();\n } else {\n ws.terminate();\n }\n } catch {\n // ignore\n }\n throw e;\n } finally {\n cleanup();\n }\n};\n"],"mappings":"AAGA;AAAA,EAEE;AAAA,EACA;AAAA,EACA;AAAA,EACA;AAAA,EACA;AAAA,EACA;AAAA,EACA;AAAA,EACA;AAAA,OACK;AAEP,SAAS,QAAQ,cAAc;AAC/B,SAAS,eAAe;AACxB,SAAuB,iBAAiB;AAGxC,MAAM,gBAA2B;AACjC,MAAM,mBAAmB;AACzB,MAAM,mBAAmB;AACzB,MAAM,eAAe;AAErB,MAAM,sBAAsB;AA0C5B,MAAM,eAAiE;AAAA,EACrE,OAAO;AAAA,EACP,SAAS;AAAA,EACT,YAAY;AAAA,EACZ,SAAS;AAAA,EACT,aAAa;AAAA,EACb,aAAa;AACf;AAEA,MAAM,sBAAsB,CAAC,gBAAwB;AACnD,MAAI,CAAC,OAAO,SAAS,WAAW,KAAK,cAAc,OAAO,cAAc,KAAK;AAC3E,UAAM,IAAI,MAAM,yCAAyC;AAAA,EAC3D;AACF;AAMA,MAAM,kBAAkB,CAAC,MAA0B,OAAe,OAAgC;AAChG,QAAM,UACJ,KAAK,UAAU,UAAa,KAAK,WAAW,SACxC;AAAA,IACE,GAAI,KAAK,UAAU,SAAY,EAAE,OAAO,KAAK,MAAM,IAAI,CAAC;AAAA,IACxD,GAAI,KAAK,WAAW,SAAY,EAAE,QAAQ,KAAK,OAAO,IAAI,CAAC;AAAA,EAC7D,IACA;AAEN,SAAO;AAAA,IACL;AAAA,IACA,cAAc,KAAK;AAAA,IACnB,QAAQ;AAAA,IACR,aAAa,KAAK;AAAA,IAClB,aAAa;AAAA,IACb,cAAc;AAAA,IACd,YAAY,CAAC;AAAA;AAAA;AAAA,IAGb,cAAc,KAAK,WAAW;AAAA,IAC9B,WAAW;AAAA,IACX,SAAS,KAAK;AAAA,IACd;AAAA,IACA,OAAO;AAAA,IACP,aAAa;AAAA,EACf;AACF;AAEO,MAAM,YAAY,IAAI,IAAI;AAAA,EAC/B;AAAA,EACA,QAAQ;AAAA,EAER,YAAY,OAAmB,CAAC,GAAG;AACjC,UAAM,SAAS,KAAK,UAAU,QAAQ,IAAI;AAC1C,QAAI,CAAC,QAAQ;AACX,YAAM,IAAI;AAAA,QACR;AAAA,MACF;AAAA,IACF;AAEA,UAAM,cAAc,KAAK,eAAe,aAAa;AACrD,wBAAoB,WAAW;AAE/B,UAAM,aAAa,KAAK,cAAc,aAAa;AAEnD,UAAM,YAAY,cAAc,EAAE,WAAW,KAAK,CAAC;AAKnD,UAAM,YACJ,KAAK,aAAa,IAAI,SAAS,MAAM,kBAAkB,EAAE,mBAAmB,EAAE,CAAC;AAEjF,SAAK,QAAQ;AAAA,MACX;AAAA,MACA,OAAO,KAAK,SAAS,aAAa;AAAA,MAClC,SAAS,KAAK,WAAW,aAAa;AAAA,MACtC;AAAA,MACA,SAAS,KAAK,WAAW,aAAa;AAAA,MACtC,aAAa,KAAK,eAAe,aAAa;AAAA,MAC9C;AAAA,MACA,OAAO,KAAK;AAAA,MACZ,QAAQ,KAAK;AAAA,MACb;AAAA,IACF;AAAA,EACF;AAAA,EAEA,IAAI,QAAgB;AAClB,WAAO,KAAK,MAAM;AAAA,EACpB;AAAA,EAEA,IAAI,WAAmB;AACrB,WAAO;AAAA,EACT;AAAA,EAEA,cAAc,MAOL;AACP,QAAI,KAAK,UAAU,OAAW,MAAK,MAAM,QAAQ,KAAK;AACtD,QAAI,KAAK,YAAY,OAAW,MAAK,MAAM,UAAU,KAAK;AAC1D,QAAI,KAAK,gBAAgB,OAAW,MAAK,MAAM,cAAc,KAAK;AAClE,QAAI,KAAK,gBAAgB,QAAW;AAClC,0BAAoB,KAAK,WAAW;AACpC,WAAK,MAAM,cAAc,KAAK;AAAA,IAChC;AACA,QAAI,KAAK,UAAU,OAAW,MAAK,MAAM,QAAQ,KAAK;AACtD,QAAI,KAAK,WAAW,OAAW,MAAK,MAAM,SAAS,KAAK;AAAA,EAC1D;AAAA,EAEA,WACE,MACA,aACA,aACmB;AACnB,WAAO,IAAI,cAAc,MAAM,MAAM,KAAK,OAAO,aAAa,WAAW;AAAA,EAC3E;AAAA,EAEA,OAAO,SAAqE;AAC1E,WAAO,IAAI,iBAAiB,MAAM,KAAK,OAAO,mCAAS,WAAW;AAAA,EACpE;AACF;AAEO,MAAM,sBAAsB,IAAI,cAAc;AAAA,EACnD,QAAQ;AAAA,EACR,UAAU,IAAI;AAAA,EACd;AAAA,EACA;AAAA,EAEA,YACEA,MACA,MACA,MACA,aACA,aACA;AACA,UAAM,MAAMA,MAAK,aAAa,WAAW;AACzC,SAAK,QAAQ;AACb,SAAK,QAAQ;AAAA,EACf;AAAA,EAEA,MAAgB,MAAM;AACpB,UAAM,YAAY,UAAU;AAC5B,UAAM,UAAU,IAAI,gBAAgB,KAAK,MAAM,YAAY,YAAY;AACvE,UAAM,UAAU,OAAO,gBAAgB,KAAK,OAAO,KAAK,KAAK,CAAC;AAE9D,UAAM,UAAU,IAAI,IAAI,KAAK,MAAM,OAAO;AAC1C,UAAM,UAAU,QAAQ,aAAa;AACrC,QAAI,CAAC,SAAS;AAGZ,YAAM,IAAI,mBAAmB;AAAA,QAC3B,SAAS,2CAA2C,KAAK,MAAM,OAAO;AAAA,MACxE,CAAC;AAAA,IACH;AAEA,UAAM,UAAU,IAAI,OAAa;AAEjC,UAAM,MAAM;AAAA,MACV;AAAA,QACE,UAAU,QAAQ;AAAA,QAClB,MAAM,SAAS,QAAQ,IAAI,KAAK;AAAA,QAChC,MAAM;AAAA,QACN,QAAQ;AAAA,QACR,SAAS;AAAA,UACP,eAAe,UAAU,KAAK,MAAM,MAAM;AAAA,UAC1C,gBAAgB;AAAA,UAChB,OAAO,KAAK,MAAM;AAAA,UAClB,kBAAkB,QAAQ;AAAA,QAC5B;AAAA,QACA,QAAQ,KAAK;AAAA,MACf;AAAA,MACA,CAAC,QAAQ;AACP,cAAM,SAAS,IAAI,cAAc;AACjC,YAAI,SAAS,OAAO,UAAU,KAAK;AACjC,gBAAM,SAAmB,CAAC;AAC1B,cAAI,GAAG,QAAQ,CAAC,MAAc,OAAO,KAAK,CAAC,CAAC;AAC5C,cAAI,GAAG,OAAO,MAAM;AAClB,kBAAM,OAAO,OAAO,OAAO,MAAM,EAAE,SAAS;AAC5C,gBAAI,CAAC,QAAQ,MAAM;AACjB,sBAAQ;AAAA,gBACN,IAAI,eAAe;AAAA,kBACjB,SAAS,kCAAkC,IAAI;AAAA,kBAC/C,SAAS,EAAE,YAAY,QAAQ,MAAM,EAAE,KAAK,KAAK,EAAE;AAAA,gBACrD,CAAC;AAAA,cACH;AAAA,YACF;AAAA,UACF,CAAC;AACD,cAAI,GAAG,SAAS,CAAC,QAAQ;AACvB,gBAAI,IAAI,YAAY,UAAW;AAC/B,iBAAK,QAAQ,MAAM,EAAE,IAAI,GAAG,4CAA4C;AACxE,gBAAI,CAAC,QAAQ,MAAM;AACjB,sBAAQ;AAAA,gBACN,IAAI,eAAe;AAAA,kBACjB,SAAS,yCAAyC,MAAM;AAAA,kBACxD,SAAS,EAAE,YAAY,OAAO;AAAA,gBAChC,CAAC;AAAA,cACH;AAAA,YACF;AAAA,UACF,CAAC;AACD;AAAA,QACF;AAEA,YAAI,GAAG,QAAQ,CAAC,UAAkB;AAChC,qBAAW,SAAS,QAAQ,MAAM,KAAK,GAAG;AACxC,iBAAK,MAAM,IAAI;AAAA,cACb;AAAA,cACA,WAAW;AAAA,cACX;AAAA,cACA,OAAO;AAAA,YACT,CAAC;AAAA,UACH;AAAA,QACF,CAAC;AACD,YAAI,GAAG,SAAS,MAAM;AACpB,qBAAW,SAAS,QAAQ,MAAM,GAAG;AACnC,iBAAK,MAAM,IAAI;AAAA,cACb;AAAA,cACA,WAAW;AAAA,cACX;AAAA,cACA,OAAO;AAAA,YACT,CAAC;AAAA,UACH;AACA,cAAI,CAAC,KAAK,MAAM,OAAQ,MAAK,MAAM,MAAM;AACzC,cAAI,CAAC,QAAQ,KAAM,SAAQ,QAAQ;AAAA,QACrC,CAAC;AACD,YAAI,GAAG,SAAS,CAAC,QAAQ;AACvB,cAAI,IAAI,YAAY,UAAW;AAC/B,eAAK,QAAQ,MAAM,EAAE,IAAI,GAAG,+BAA+B;AAC3D,cAAI,CAAC,QAAQ,KAAM,SAAQ,OAAO,GAAG;AAAA,QACvC,CAAC;AAAA,MACH;AAAA,IACF;AAEA,QAAI,GAAG,SAAS,CAAC,QAAQ;AACvB,UAAI,IAAI,SAAS,aAAc;AAC/B,WAAK,QAAQ,MAAM,EAAE,IAAI,GAAG,8BAA8B;AAC1D,UAAI,CAAC,QAAQ,KAAM,SAAQ,OAAO,GAAG;AAAA,IACvC,CAAC;AACD,QAAI,MAAM,OAAO;AACjB,QAAI,IAAI;AAER,QAAI;AACF,YAAM,QAAQ;AAAA,IAChB,SAAS,GAAG;AACV,UAAI,KAAK,YAAY,QAAS;AAC9B,UAAI,CAAC,KAAK,MAAM,OAAQ,MAAK,MAAM,MAAM;AACzC,UAAI,aAAa,kBAAkB,aAAa,oBAAoB;AAClE,cAAM;AAAA,MACR;AACA,YAAM,IAAI,mBAAmB;AAAA,QAC3B,SAAS,iCAAkC,EAAY,WAAW,eAAe;AAAA,MACnF,CAAC;AAAA,IACH;AAAA,EACF;AACF;AAEO,MAAM,yBAAyB,IAAI,iBAAiB;AAAA,EACzD,QAAQ;AAAA,EACR,UAAU,IAAI;AAAA,EACd;AAAA,EAEA,YAAYA,MAAU,MAA0B,aAAiC;AAC/E,UAAMA,MAAK,WAAW;AACtB,SAAK,QAAQ;AAAA,EACf;AAAA,EAEA,MAAgB,MAAM;AACpB,UAAM,YAAY,UAAU;AAC5B,UAAM,UAAU,IAAI,gBAAgB,KAAK,MAAM,YAAY,YAAY;AAOvE,UAAM,aAAa,KAAK,MAAM,UAAU,OAAO;AAE/C,UAAM,QAAQ,GAAG,KAAK,MAAM,QAAQ,QAAQ,SAAS,IAAI,CAAC;AAC1D,QAAI;AACJ,QAAI;AACF,WAAK,MAAM,iBAAiB;AAAA,QAC1B,KAAK;AAAA,QACL,SAAS;AAAA,UACP,eAAe,UAAU,KAAK,MAAM,MAAM;AAAA,UAC1C,OAAO,KAAK,MAAM;AAAA,QACpB;AAAA,QACA,WAAW,KAAK,YAAY;AAAA,QAC5B,aAAa,KAAK;AAAA,MACpB,CAAC;AAAA,IACH,SAAS,GAAG;AACV,YAAM,IAAI,mBAAmB;AAAA,QAC3B,SAAS,wCAAyC,EAAY,WAAW,eAAe;AAAA,MAC1F,CAAC;AAAA,IACH;AAEA,UAAM,WAAW,IAAI,OAAa;AAElC,UAAM,YAAY,YAAY;AAC5B,UAAI;AACF,yBAAiB,QAAQ,KAAK,OAAO;AACnC,cAAI,KAAK,gBAAgB,OAAO,QAAS;AACzC,cAAI,SAAS,iBAAiB,gBAAgB;AAC5C,uBAAW,MAAM;AACjB;AAAA,UACF;AACA,cAAI,CAAC,KAAM;AACX,qBAAW,SAAS,IAAI;AAAA,QAC1B;AAAA,MACF,UAAE;AACA,YAAI,CAAC,WAAW,OAAQ,YAAW,SAAS;AAAA,MAC9C;AAAA,IACF;AAEA,UAAM,WAAW,YAAY;AAC3B,YAAM,WAAW,EAAE,OAAO,SAAS,SAAS,gBAAgB,KAAK,KAAK,EAAE;AACxE,SAAI,KAAK,OAAO,KAAK,OAAO,QAAQ,CAAC,CAAC;AAEtC,uBAAiB,MAAM,YAAY;AACjC,YAAI,KAAK,gBAAgB,OAAO,QAAS;AACzC,cAAM,WAAW,GAAG;AACpB,YAAI,CAAC,SAAU;AACf,aAAK,YAAY;AACjB,WAAI,KAAK,OAAO,KAAK,OAAO,EAAE,OAAO,QAAQ,MAAM,WAAW,IAAI,CAAC,CAAC,CAAC;AACrE,WAAI,KAAK,OAAO,KAAK,OAAO,EAAE,OAAO,QAAQ,CAAC,CAAC,CAAC;AAAA,MAClD;AAEA,UAAI,CAAC,KAAK,gBAAgB,OAAO,SAAS;AACxC,WAAI,KAAK,OAAO,KAAK,OAAO,EAAE,OAAO,OAAO,CAAC,CAAC,CAAC;AAAA,MACjD;AAAA,IACF;AAEA,QAAI;AACJ,UAAM,gBAAgB,CAAC,UAAmB;AACxC,UAAI,WAAW;AACb,aAAK,MAAM,IAAI,EAAE,WAAW,WAAW,WAAW,OAAO,WAAW,MAAM,CAAC;AAC3E,oBAAY;AAAA,MACd;AAAA,IACF;AAEA,UAAM,WAAW,YAAY;AAG3B,YAAM,YAAY,CAAC,QAAiB;AAClC,YAAI;AACJ,YAAI,OAAO,SAAS,GAAG,GAAG;AACxB,kBAAQ;AAAA,QACV,WAAW,MAAM,QAAQ,GAAG,GAAG;AAC7B,kBAAQ,OAAO,OAAO,GAAG;AAAA,QAC3B,OAAO;AACL,kBAAQ,OAAO,KAAK,GAAkB;AAAA,QACxC;AAEA,YAAI;AACJ,YAAI;AACF,mBAAS,OAAO,KAAK;AAAA,QACvB,SAAS,KAAK;AACZ,eAAK,QAAQ,KAAK,EAAE,IAAI,GAAG,qCAAqC;AAChE;AAAA,QACF;AAEA,cAAM,QAAQ,OAAO;AACrB,YAAI,UAAU,SAAS;AACrB,gBAAM,QAAQ,OAAO;AACrB,cAAI,SAAS,MAAM,aAAa,GAAG;AACjC,uBAAW,KAAK,QAAQ,MAAM,KAAK,GAAG;AACpC,4BAAc,KAAK;AACnB,0BAAY;AAAA,YACd;AAAA,UACF;AAAA,QACF,WAAW,UAAU,UAAU;AAC7B,gBAAM,SAAS,OAAO;AACtB,cAAI,WAAW,SAAS;AACtB,qBAAS;AAAA,cACP,IAAI,eAAe;AAAA,gBACjB,SAAS;AAAA,gBACT,SAAS,EAAE,MAAM,EAAE,KAAK,KAAK,UAAU,MAAM,EAAE,EAAE;AAAA,cACnD,CAAC;AAAA,YACH;AACA;AAAA,UACF;AACA,qBAAW,KAAK,QAAQ,MAAM,GAAG;AAC/B,0BAAc,KAAK;AACnB,wBAAY;AAAA,UACd;AACA,wBAAc,IAAI;AAClB,cAAI,CAAC,KAAK,MAAM,QAAQ;AACtB,iBAAK,MAAM,IAAI,iBAAiB,aAAa;AAAA,UAC/C;AACA,cAAI,CAAC,SAAS,KAAM,UAAS,QAAQ;AAAA,QACvC,OAAO;AACL,eAAK,QAAQ,MAAM,EAAE,MAAM,GAAG,0BAA0B;AAAA,QAC1D;AAAA,MACF;AAEA,YAAM,UAAU,CAAC,MAAc,WAAmB;AAChD,YAAI,CAAC,SAAS,MAAM;AAClB,mBAAS;AAAA,YACP,IAAI,eAAe;AAAA,cACjB,SAAS;AAAA,cACT,SAAS;AAAA,gBACP,YAAY,QAAQ;AAAA,gBACpB,MAAM,EAAE,QAAQ,OAAO,SAAS,EAAE;AAAA,cACpC;AAAA,YACF,CAAC;AAAA,UACH;AAAA,QACF;AAAA,MACF;AAEA,YAAM,UAAU,CAAC,QAAe;AAC9B,YAAI,CAAC,SAAS,KAAM,UAAS,OAAO,GAAG;AAAA,MACzC;AAEA,SAAI,GAAG,WAAW,SAAS;AAC3B,SAAI,GAAG,SAAS,OAAO;AACvB,SAAI,GAAG,SAAS,OAAO;AAEvB,UAAI;AACF,cAAM,SAAS;AAAA,MACjB,UAAE;AACA,WAAI,IAAI,WAAW,SAAS;AAC5B,WAAI,IAAI,SAAS,OAAO;AACxB,WAAI,IAAI,SAAS,OAAO;AAAA,MAC1B;AAAA,IACF;AAEA,QAAI;AACF,YAAM,QAAQ,IAAI,CAAC,UAAU,GAAG,SAAS,GAAG,SAAS,CAAC,CAAC;AAAA,IACzD,SAAS,GAAG;AACV,UAAI,KAAK,YAAY,QAAS;AAC9B,UAAI,aAAa,kBAAkB,aAAa,oBAAoB;AAClE,cAAM;AAAA,MACR;AACA,YAAM,IAAI,mBAAmB;AAAA,QAC3B,SAAS,gCAAiC,EAAY,WAAW,eAAe;AAAA,MAClF,CAAC;AAAA,IACH,UAAE;AACA,UAAI,CAAC,WAAW,OAAQ,YAAW,MAAM;AACzC,UAAI,MAAM,GAAG,eAAe,UAAU,QAAQ;AAC5C,YAAI;AACF,aAAG,MAAM;AAAA,QACX,QAAQ;AAAA,QAER;AAAA,MACF;AAAA,IACF;AAAA,EACF;AACF;AAEA,MAAM,mBAAmB,OAAO;AAAA,EAC9B;AAAA,EACA;AAAA,EACA;AAAA,EACA;AACF,MAK0B;AACxB,QAAM,KAAK,IAAI,UAAU,KAAK,EAAE,SAAS,kBAAkB,UAAU,CAAC;AACtE,QAAM,MAAM,IAAI,OAAa;AAE7B,MAAI;AACJ,QAAM,UAAU,MAAM;AACpB,QAAI,QAAS,cAAa,OAAO;AACjC,OAAG,IAAI,QAAQ,MAAM;AACrB,OAAG,IAAI,SAAS,OAAO;AACvB,OAAG,IAAI,SAAS,OAAO;AACvB,gBAAY,oBAAoB,SAAS,OAAO;AAAA,EAClD;AAEA,QAAM,SAAS,MAAM,IAAI,QAAQ;AACjC,QAAM,UAAU,CAAC,QAAe,IAAI,OAAO,GAAG;AAC9C,QAAM,UAAU,CAAC,MAAc,WAC7B,IAAI;AAAA,IACF,IAAI,MAAM,sCAAsC,IAAI,YAAY,OAAO,SAAS,CAAC,GAAG;AAAA,EACtF;AACF,QAAM,UAAU,MAAM,IAAI,OAAO,IAAI,MAAM,SAAS,CAAC;AAErD,KAAG,GAAG,QAAQ,MAAM;AACpB,KAAG,GAAG,SAAS,OAAO;AACtB,KAAG,GAAG,SAAS,OAAO;AACtB,cAAY,iBAAiB,SAAS,SAAS,EAAE,MAAM,KAAK,CAAC;AAE7D,MAAI,YAAY,GAAG;AACjB,cAAU,WAAW,MAAM,IAAI,OAAO,IAAI,MAAM,iBAAiB,CAAC,GAAG,SAAS;AAAA,EAChF;AAEA,MAAI;AACF,UAAM,IAAI;AACV,WAAO;AAAA,EACT,SAAS,GAAG;AACV,QAAI;AACF,SAAG,GAAG,SAAS,MAAM;AAAA,MAAC,CAAC;AACvB,UAAI,GAAG,eAAe,UAAU,YAAY;AAC1C,WAAG,MAAM;AAAA,MACX,OAAO;AACL,WAAG,UAAU;AAAA,MACf;AAAA,IACF,QAAQ;AAAA,IAER;AACA,UAAM;AAAA,EACR,UAAE;AACA,YAAQ;AAAA,EACV;AACF;","names":["tts"]} |
+7
-7
| { | ||
| "name": "@livekit/agents-plugin-fishaudio", | ||
| "version": "1.5.0", | ||
| "version": "1.5.1", | ||
| "description": "Fish Audio plugin for LiveKit Node Agents", | ||
@@ -28,3 +28,3 @@ "main": "dist/index.js", | ||
| "devDependencies": { | ||
| "@livekit/rtc-node": "^0.13.30", | ||
| "@livekit/rtc-node": "^0.13.31", | ||
| "@microsoft/api-extractor": "^7.35.0", | ||
@@ -34,5 +34,5 @@ "@types/ws": "^8.5.10", | ||
| "typescript": "^5.0.0", | ||
| "@livekit/agents": "1.5.0", | ||
| "@livekit/agents-plugin-openai": "1.5.0", | ||
| "@livekit/agents-plugins-test": "1.5.0" | ||
| "@livekit/agents": "1.5.1", | ||
| "@livekit/agents-plugin-openai": "1.5.1", | ||
| "@livekit/agents-plugins-test": "1.5.1" | ||
| }, | ||
@@ -44,4 +44,4 @@ "dependencies": { | ||
| "peerDependencies": { | ||
| "@livekit/rtc-node": "^0.13.30", | ||
| "@livekit/agents": "1.5.0" | ||
| "@livekit/rtc-node": "^0.13.31", | ||
| "@livekit/agents": "1.5.1" | ||
| }, | ||
@@ -48,0 +48,0 @@ "scripts": { |
+1
-0
@@ -392,2 +392,3 @@ // SPDX-FileCopyrightText: 2026 LiveKit, Inc. | ||
| if (!sentence) continue; | ||
| this.markStarted(); | ||
| ws!.send(Buffer.from(encode({ event: 'text', text: sentence + ' ' }))); | ||
@@ -394,0 +395,0 @@ ws!.send(Buffer.from(encode({ event: 'flush' }))); |
132152
0.13%1631
0.18%