import { randomUUID } from "node:crypto"; import { definePluginEntry } from "openclaw/plugin-sdk/plugin-entry"; const AUDIO_FORMAT = { encoding: "pcm16", sampleRateHz: 24000, channels: 1 }; function record(value) { return value && typeof value === "object" && !Array.isArray(value) ? value : {}; } function resolveModelProvider(cfg, requested) { const providers = record(record(cfg).models).providers; const requestedId = typeof requested === "string" ? requested.trim() : ""; if (requestedId) return record(record(providers)[requestedId]); for (const id of ["athena", "llama-cpp", "openai"]) { const candidate = record(record(providers)[id]); if (typeof candidate.baseUrl === "string" && candidate.baseUrl.trim()) return candidate; } return {}; } function resolveConfig(req) { const raw = record(req.providerConfig); const modelProvider = resolveModelProvider(req.cfg, raw.modelProvider); return { modelProvider: String(raw.modelProvider || ""), baseUrl: String(raw.baseUrl || modelProvider.baseUrl || "http://192.168.1.212:8081/v1").replace(/\/$/, ""), apiKey: String(raw.apiKey || modelProvider.apiKey || ""), voice: String(raw.voice || req.voice || "alloy"), language: String(raw.language || req.language || "de"), vadThreshold: Number(raw.vadThreshold ?? 0.018), silenceDurationMs: Number(raw.silenceDurationMs ?? 750), prefixPaddingMs: Number(raw.prefixPaddingMs ?? 300), maxSpeechSeconds: Number(raw.maxSpeechSeconds ?? 45), }; } function wavFromPcm16(pcm, sampleRate = 24000) { const header = Buffer.alloc(44); header.write("RIFF", 0); header.writeUInt32LE(36 + pcm.length, 4); header.write("WAVEfmt ", 8); header.writeUInt32LE(16, 16); header.writeUInt16LE(1, 20); header.writeUInt16LE(1, 22); header.writeUInt32LE(sampleRate, 24); header.writeUInt32LE(sampleRate * 2, 28); header.writeUInt16LE(2, 32); header.writeUInt16LE(16, 34); header.write("data", 36); header.writeUInt32LE(pcm.length, 40); return Buffer.concat([header, pcm]); } function pcmRms(pcm) { if (pcm.length < 2) return 0; let sum = 0; const count = Math.floor(pcm.length / 2); for (let i = 0; i < count; i += 1) { const value = pcm.readInt16LE(i * 2) / 32768; sum += value * value; } return Math.sqrt(sum / count); } function splitForIncrementalSpeech(text, limit = 80) { const words = text.replace(/\s+/g, " ").trim().split(" ").filter(Boolean); const chunks = []; 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) { if (wav.length < 44 || wav.toString("ascii", 0, 4) !== "RIFF") { throw new Error("Athena TTS did not return PCM WAV audio"); } let offset = 12; let sampleRate = 0; let channels = 0; let bits = 0; let data; while (offset + 8 <= wav.length) { const id = wav.toString("ascii", offset, offset + 4); const size = wav.readUInt32LE(offset + 4); const start = offset + 8; if (id === "fmt " && size >= 16) { if (wav.readUInt16LE(start) !== 1) throw new Error("Athena TTS WAV is not PCM"); channels = wav.readUInt16LE(start + 2); sampleRate = wav.readUInt32LE(start + 4); bits = wav.readUInt16LE(start + 14); } else if (id === "data") { data = wav.subarray(start, Math.min(start + size, wav.length)); } offset = start + size + (size % 2); } if (!data || !sampleRate || bits !== 16 || channels < 1) { throw new Error("Unsupported Athena TTS WAV format"); } const frames = Math.floor(data.length / (2 * channels)); const mono = new Int16Array(frames); for (let i = 0; i < frames; i += 1) mono[i] = data.readInt16LE(i * channels * 2); if (sampleRate === 24000) return Buffer.from(mono.buffer); const outFrames = Math.max(1, Math.round(frames * 24000 / sampleRate)); const out = Buffer.alloc(outFrames * 2); for (let i = 0; i < outFrames; i += 1) { const source = i * sampleRate / 24000; const left = Math.min(frames - 1, Math.floor(source)); const right = Math.min(frames - 1, left + 1); const fraction = source - left; const value = Math.round(mono[left] * (1 - fraction) + mono[right] * fraction); out.writeInt16LE(Math.max(-32768, Math.min(32767, value)), i * 2); } return out; } class AthenaTalkBridge { req; cfg; supportsToolResultContinuation = false; supportsToolResultSuppression = false; connected = false; closed = false; speaking = false; speech = []; prefix = []; prefixBytes = 0; silenceTimer; active; transcribing = false; generation = 0; constructor(req, cfg) { this.req = req; this.cfg = cfg; } async connect() { this.connected = true; this.req.onEvent?.({ direction: "server", type: "session.created" }); this.req.onReady?.(); } isConnected() { return this.connected && !this.closed; } setMediaTimestamp(_timestamp) { } sendAudio(chunk) { if (!this.isConnected() || chunk.length === 0) return; // Half-duplex by design: while Athena is transcribing, consulting the // agent, synthesizing, or playing a reply, microphone input is ignored. // This prevents speaker feedback and background noise from cancelling the // response that is currently being delivered. if (this.transcribing || this.active) return; const voiced = pcmRms(chunk) >= this.cfg.vadThreshold; const maxPrefix = Math.round(24000 * 2 * this.cfg.prefixPaddingMs / 1000); if (!this.speaking) { this.prefix.push(Buffer.from(chunk)); this.prefixBytes += chunk.length; while (this.prefixBytes > maxPrefix && this.prefix.length > 1) { this.prefixBytes -= this.prefix.shift().length; } if (!voiced) return; this.speaking = true; this.speech = this.prefix; this.prefix = []; this.prefixBytes = 0; } else { this.speech.push(Buffer.from(chunk)); } const maxBytes = this.cfg.maxSpeechSeconds * 24000 * 2; if (this.speech.reduce((sum, part) => sum + part.length, 0) >= maxBytes) { void this.finishSpeech(); return; } if (voiced) { if (this.silenceTimer) clearTimeout(this.silenceTimer); this.silenceTimer = setTimeout(() => void this.finishSpeech(), this.cfg.silenceDurationMs); } } sendUserMessage(text) { const trimmed = text.trim(); if (trimmed) void this.answer(trimmed); } triggerGreeting(instructions) { void this.answer(instructions?.trim() || "Begrüße mich kurz auf Deutsch."); } submitToolResult(_callId, _result) { } acknowledgeMark(_markName) { } close() { this.closed = true; this.connected = false; this.generation += 1; this.active?.abort(); if (this.silenceTimer) clearTimeout(this.silenceTimer); this.req.onClose?.("completed"); } async finishSpeech() { if (!this.speaking) return; this.speaking = false; if (this.silenceTimer) clearTimeout(this.silenceTimer); this.silenceTimer = undefined; const pcm = Buffer.concat(this.speech); this.speech = []; if (pcm.length < 24000 * 2 * 0.25) return; this.transcribing = true; try { const text = await this.transcribe(pcm); if (!text || !this.isConnected()) return; this.req.onTranscript?.("user", text, true); this.transcribing = false; await this.answer(text); } catch (error) { if (error.name !== "AbortError") this.req.onError?.(error); } finally { this.transcribing = false; } } async transcribe(pcm) { const form = new FormData(); const wavBytes = Uint8Array.from(wavFromPcm16(pcm)); form.append("file", new Blob([wavBytes], { type: "audio/wav" }), "talk.wav"); form.append("model", "whisper-1"); form.append("language", this.cfg.language); const response = await fetch(`${this.cfg.baseUrl}/audio/transcriptions`, { method: "POST", headers: this.authHeaders(), body: form, }); if (!response.ok) throw new Error(`Athena STT failed (HTTP ${response.status})`); const payload = record(await response.json()); return String(payload.text || "").trim(); } async answer(prompt) { if (!this.req.runAgentConsult) throw new Error("OpenClaw agent-consult is unavailable"); if (this.active) return; const controller = new AbortController(); this.active = controller; const generation = ++this.generation; const responseId = `athena-${randomUUID()}`; try { 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 speechChunks = splitForIncrementalSpeech(text); if (speechChunks.length === 0) return; const frameBytes = 24000 * 2 / 50; 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)); } } 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 }); } catch (error) { if (error.name === "AbortError") { this.req.onResponseDone?.({ status: "cancelled", responseId }); } else { this.req.onError?.(error); this.req.onResponseDone?.({ status: "failed", responseId, error: String(error) }); } } finally { if (this.active === controller) this.active = undefined; } } async synthesize(text, signal) { const response = await fetch(`${this.cfg.baseUrl}/audio/speech`, { method: "POST", headers: { "Content-Type": "application/json", ...this.authHeaders() }, // Athena normalizes the request and serves it through Qwen3-TTS. body: JSON.stringify({ voice: this.cfg.voice, input: text, response_format: "wav" }), signal, }); if (!response.ok) throw new Error(`Athena TTS failed (HTTP ${response.status})`); return readWavPcm24k(Buffer.from(await response.arrayBuffer())); } authHeaders() { return this.cfg.apiKey ? { Authorization: `Bearer ${this.cfg.apiKey}` } : {}; } } export default definePluginEntry({ id: "athena-talk", name: "Athena Local Talk", description: "Private voice loop using Athena STT and TTS with the normal OpenClaw agent.", register(api) { api.registerRealtimeVoiceProvider({ id: "athena-talk", label: "Athena Local Talk", aliases: ["athena", "local-athena"], defaultModel: "athena-local", voices: ["alloy"], autoSelectOrder: 1, capabilities: { transports: ["gateway-relay"], inputAudioFormats: [AUDIO_FORMAT], outputAudioFormats: [AUDIO_FORMAT], supportsBargeIn: false, handlesInputAudioBargeIn: false, supportsToolCalls: true, supportsSessionResumption: false, }, resolveConfig: ({ rawConfig }) => record(rawConfig), isConfigured: ({ providerConfig, cfg }) => { const raw = record(providerConfig); return Boolean(raw.baseUrl || resolveModelProvider(cfg, raw.modelProvider).baseUrl); }, createBridge: (req) => new AthenaTalkBridge(req, resolveConfig(req)), }); }, });