Stream Athena Talk replies incrementally

This commit is contained in:
Mikei386
2026-09-03 23:36:51 +02:00
parent 384d81f6cd
commit 435c59da41
5 changed files with 90 additions and 31 deletions
+42 -19
View File
@@ -64,6 +64,27 @@ function pcmRms(pcm: Buffer): number {
return Math.sqrt(sum / count);
}
function splitForIncrementalSpeech(text: string, limit = 80): string[] {
const words = text.replace(/\s+/g, " ").trim().split(" ").filter(Boolean);
const chunks: string[] = [];
let current = "";
for (const word of words) {
const candidate = current ? `${current} ${word}` : word;
if (current && candidate.length > limit) {
chunks.push(current);
current = word;
} else {
current = candidate;
}
if (current.length >= 35 && /[.!?](?:["')\]_*]+)?$/.test(word)) {
chunks.push(current);
current = "";
}
}
if (current) chunks.push(current);
return chunks;
}
function readWavPcm24k(wav: Buffer): Buffer {
if (wav.length < 44 || wav.toString("ascii", 0, 4) !== "RIFF") {
throw new Error("Athena TTS did not return PCM WAV audio");
@@ -238,33 +259,35 @@ class AthenaTalkBridge {
const generation = ++this.generation;
const responseId = `athena-${randomUUID()}`;
try {
const result = await this.req.runAgentConsult({ prompt, signal: controller.signal });
const voicePrompt = `${prompt}\n\nAntwortregeln für diese Sprachantwort: Antworte ausschließlich auf Deutsch, kurz und direkt. Keine Analyse, keine Meta-Kommentare, kein Markdown und keine Wiederholung der Anfrage. Gib nur den Text aus, der gesprochen werden soll.`;
const result = await this.req.runAgentConsult({ prompt: voicePrompt, signal: controller.signal });
if (generation !== this.generation || !this.isConnected()) return;
const text = String(result?.text || "").trim();
if (!text) return;
this.req.onEvent?.({ direction: "server", type: "response.created", responseId });
this.req.onTranscript?.("assistant", text, true);
const pcm = await this.synthesize(text, controller.signal);
if (generation !== this.generation || !this.isConnected()) return;
// Feed OpenClaw in its native 20 ms frame size. Sending larger bursts
// made the relay split several frames at once; slow-client protection
// could then drop individual frames and produce audible holes.
const speechChunks = splitForIncrementalSpeech(text);
if (speechChunks.length === 0) return;
const frameBytes = 24000 * 2 / 50;
const playbackStartedAt = Date.now();
for (let offset = 0; offset < pcm.length; offset += frameBytes) {
if (controller.signal.aborted || generation !== this.generation || !this.isConnected()) {
return;
let playbackEndsAt = 0;
let nextAudio = this.synthesize(speechChunks[0], controller.signal);
for (let index = 0; index < speechChunks.length; index += 1) {
const pcm = await nextAudio;
if (generation !== this.generation || !this.isConnected()) return;
if (index + 1 < speechChunks.length) {
nextAudio = this.synthesize(speechChunks[index + 1], controller.signal);
}
const audioDurationMs = pcm.length / (24000 * 2) * 1000;
playbackEndsAt = Math.max(playbackEndsAt, Date.now()) + audioDurationMs;
for (let offset = 0; offset < pcm.length; offset += frameBytes) {
if (controller.signal.aborted || generation !== this.generation || !this.isConnected()) {
return;
}
this.req.onAudio?.(pcm.subarray(offset, Math.min(offset + frameBytes, pcm.length)));
await new Promise((resolve) => setTimeout(resolve, 18));
}
this.req.onAudio?.(pcm.subarray(offset, Math.min(offset + frameBytes, pcm.length)));
// Run just ahead of realtime so the client builds a small jitter
// buffer without flooding the WebSocket queue.
await new Promise((resolve) => setTimeout(resolve, 18));
}
// Transmission is intentionally a little faster than playback. Keep the
// microphone gated until the client has consumed that buffered tail;
// otherwise the last spoken words are transcribed again as user input.
const audioDurationMs = pcm.length / (24000 * 2) * 1000;
const remainingPlaybackMs = Math.max(0, audioDurationMs - (Date.now() - playbackStartedAt));
const remainingPlaybackMs = Math.max(0, playbackEndsAt - Date.now());
await new Promise((resolve) => setTimeout(resolve, remainingPlaybackMs + 300));
this.req.onResponseDone?.({ status: "completed", responseId });
this.req.onEvent?.({ direction: "server", type: "response.done", responseId });