diff --git a/packages/stage-ui/src/composables/queues.ts b/packages/stage-ui/src/composables/queues.ts index fd23ed99d..a44967566 100644 --- a/packages/stage-ui/src/composables/queues.ts +++ b/packages/stage-ui/src/composables/queues.ts @@ -4,7 +4,8 @@ import type { UseQueueReturn } from './queue' import { sleep } from '@moeru/std' import { EMOTION_VALUES } from '../constants/emotions' -import { chunkTTSInput } from '../utils/tts' +import { createControllableStream } from '../utils/stream' +import { chunkToTTSQueue } from '../utils/tts' import { useQueue } from './queue' export function useEmotionsMessageQueue(emotionsQueue: UseQueueReturn) { @@ -102,28 +103,14 @@ export function useDelayMessageQueue() { export function useMessageContentQueue(ttsQueue: UseQueueReturn) { const encoder = new TextEncoder() - let enqueue: (data: Uint8Array) => void - const stream = new ReadableStream({ - start(controller) { - enqueue = data => controller.enqueue(data) - }, - }); + const { stream, controller } = createControllableStream() - (async () => { - try { - for await (const chunk of chunkTTSInput(stream.getReader())) { - await ttsQueue.add(chunk.text) - } - } - catch (e) { - console.error('Error chunking input stream for TTS:', e) - } - })() + chunkToTTSQueue(stream.getReader(), ttsQueue) return useQueue({ handlers: [ async (ctx) => { - enqueue(encoder.encode(ctx.data)) + controller.enqueue(encoder.encode(ctx.data)) }, ], }) diff --git a/packages/stage-ui/src/utils/stream.ts b/packages/stage-ui/src/utils/stream.ts new file mode 100644 index 000000000..8e5d3d226 --- /dev/null +++ b/packages/stage-ui/src/utils/stream.ts @@ -0,0 +1,16 @@ +export interface ControllableStream { + stream: ReadableStream + controller: ReadableStreamDefaultController +} + +export function createControllableStream(): ControllableStream { + // WHY!: ReadableStream.start is called synchronously and immediately + let controller!: ReadableStreamDefaultController + + const stream = new ReadableStream({ + start(ctrl) { + controller = ctrl + }, + }) + return { stream, controller } +} diff --git a/packages/stage-ui/src/utils/tts.ts b/packages/stage-ui/src/utils/tts.ts index 2a58729ba..e0c218a2b 100644 --- a/packages/stage-ui/src/utils/tts.ts +++ b/packages/stage-ui/src/utils/tts.ts @@ -1,5 +1,7 @@ import type { ReaderLike } from 'clustr' +import type { UseQueueReturn } from '../composables/queue' + import { readGraphemeClusters } from 'clustr' // A special character to instruct the TTS pipeline to flush @@ -144,3 +146,14 @@ export async function* chunkTTSInput(input: string | ReaderLike, options?: TTSIn } } } + +export async function chunkToTTSQueue(reader: ReaderLike, queue: UseQueueReturn) { + try { + for await (const chunk of chunkTTSInput(reader)) { + await queue.add(chunk.text) + } + } + catch (e) { + console.error('Error chunking stream to TTS queue:', e) + } +}