import { generateKeyPairSync, randomUUID, sign, verify } from "node:crypto"; import type { IncomingMessage, ServerResponse } from "node:http"; import { definePluginEntry } from "openclaw/plugin-sdk/plugin-entry"; const AUDIO_FORMAT = { encoding: "pcm16", sampleRateHz: 24000, channels: 1 } as const; type ProviderConfig = { modelProvider?: string; baseUrl?: string; apiKey?: string; voice?: string; language?: string; vadThreshold?: number; silenceDurationMs?: number; prefixPaddingMs?: number; maxSpeechSeconds?: number; realtimeUpstreamUrl?: string; }; const BROWSER_OFFER_PATH = "/plugins/athena-talk/realtime/calls"; const BROWSER_KEY_PATH = "/plugins/athena-talk/realtime/public-key"; const MAX_OFFER_BYTES = 64 * 1024; const browserKeys = generateKeyPairSync("ed25519"); const publicKeyPem = browserKeys.publicKey.export({ format: "pem", type: "spki" }).toString(); function base64url(value: string): string { return Buffer.from(value).toString("base64url"); } function createBrowserToken(): { token: string; expiresAt: number } { const expiresAt = Date.now() + 60_000; const payload = base64url(JSON.stringify({ exp: Math.floor(expiresAt / 1000), jti: randomUUID() })); const signature = sign(null, Buffer.from(payload), browserKeys.privateKey).toString("base64url"); return { token: `${payload}.${signature}`, expiresAt }; } function validBrowserToken(value: string): boolean { const [payload, signature, extra] = value.split("."); if (!payload || !signature || extra !== undefined) return false; try { if (!verify(null, Buffer.from(payload), browserKeys.publicKey, Buffer.from(signature, "base64url"))) { return false; } const claims = record(JSON.parse(Buffer.from(payload, "base64url").toString("utf8"))); return typeof claims.jti === "string" && Number.isInteger(claims.exp) && claims.exp > Math.floor(Date.now() / 1000) && claims.exp <= Math.floor(Date.now() / 1000) + 65; } catch { return false; } } async function handleBrowserOffer( req: IncomingMessage, res: ServerResponse, upstream: string, ): Promise { if (req.method !== "POST") { res.writeHead(405).end("Method not allowed"); return true; } const token = req.headers.authorization?.replace(/^Bearer /i, "") || ""; if (!validBrowserToken(token)) { res.writeHead(401).end("Invalid realtime session token"); return true; } if (!upstream.startsWith("http://") && !upstream.startsWith("https://")) { res.writeHead(503).end("Athena realtime upstream is not configured"); return true; } const chunks: Buffer[] = []; let total = 0; for await (const part of req) { const chunk = Buffer.from(part); total += chunk.length; if (total > MAX_OFFER_BYTES) { res.writeHead(413).end("WebRTC offer too large"); return true; } chunks.push(chunk); } try { const response = await fetch(upstream, { method: "POST", headers: { Authorization: `Bearer ${token}`, "Content-Type": "application/sdp" }, body: Uint8Array.from(Buffer.concat(chunks)), signal: AbortSignal.timeout(30_000), }); const answer = Buffer.from(await response.arrayBuffer()); if (answer.length > MAX_OFFER_BYTES) throw new Error("WebRTC answer too large"); res.writeHead(response.status, { "Content-Type": response.headers.get("content-type") || "text/plain" }); res.end(answer); } catch { res.writeHead(502).end("Athena realtime bridge unavailable"); } return true; } function record(value: unknown): Record { return value && typeof value === "object" && !Array.isArray(value) ? value as Record : {}; } function resolveModelProvider(cfg: unknown, requested?: unknown): Record { 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: any): Required { const raw = record(req.providerConfig); const modelProvider = resolveModelProvider(req.cfg, raw.modelProvider); const baseUrl = String(raw.baseUrl || modelProvider.baseUrl || "http://192.168.1.212:8081/v1").replace(/\/$/, ""); const tts = record(record(record(req.cfg).tts).providers).openai; let sharedTtsKey = ""; try { if (new URL(String(record(tts).baseUrl)).origin === new URL(baseUrl).origin) { sharedTtsKey = typeof record(tts).apiKey === "string" ? record(tts).apiKey : ""; } } catch { /* The TTS provider may be absent or use another URL. */ } return { modelProvider: String(raw.modelProvider || ""), baseUrl, apiKey: String(raw.apiKey || modelProvider.apiKey || sharedTtsKey || ""), 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), realtimeUpstreamUrl: String(raw.realtimeUpstreamUrl || ""), }; } function wavFromPcm16(pcm: Buffer, sampleRate = 24000): Buffer { 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]); } // The browser transcription relay sends 8 kHz G.711 mu-law, while Athena's // existing Whisper endpoint accepts PCM WAV uploads. function wavFromMulaw8k(audio: Buffer): Buffer { const pcm = Buffer.allocUnsafe(audio.length * 2); for (let i = 0; i < audio.length; i += 1) { const sample = ~audio[i] & 0xff; const magnitude = (((sample & 15) << 3) + 132) << ((sample >> 4) & 7); pcm.writeInt16LE(Math.max(-32768, Math.min(32767, (sample & 128) ? 132 - magnitude : magnitude - 132)), i * 2); } return wavFromPcm16(pcm, 8000); } function resolveTranscriptionConfig(cfg: unknown, rawConfig: unknown): Required { const talkConfig = record(record(record(cfg).talk).realtime); const talkProvider = record(record(talkConfig.providers)["athena-talk"]); return resolveConfig({ cfg, providerConfig: { ...talkProvider, ...record(rawConfig) } }); } class AthenaTranscriptionSession { private connected = false; private closed = false; private readonly audio: Buffer[] = []; private bytes = 0; constructor( private readonly req: { cfg?: unknown; providerConfig: Record; onSpeechStart?: () => void; onTranscript?: (text: string) => void; onError?: (error: Error) => void; }, private readonly config: Required, ) {} async connect(): Promise { this.connected = true; } isConnected(): boolean { return this.connected && !this.closed; } sendAudio(audio: Buffer): void { if (!this.isConnected() || audio.length === 0) return; if (this.bytes === 0) this.req.onSpeechStart?.(); const maxBytes = this.config.maxSpeechSeconds * 8000; if (this.bytes + audio.length > maxBytes) { this.req.onError?.(new Error(`Athena dictation is limited to ${this.config.maxSpeechSeconds} seconds`)); this.close(); return; } this.audio.push(Buffer.from(audio)); this.bytes += audio.length; } close(): void { if (this.closed) return; this.closed = true; this.connected = false; if (!this.bytes) return; const audio = Buffer.concat(this.audio); this.audio.length = 0; void this.transcribe(audio); } private async transcribe(audio: Buffer): Promise { try { const form = new FormData(); form.append("file", new Blob([Uint8Array.from(wavFromMulaw8k(audio))], { type: "audio/wav" }), "dictation.wav"); form.append("model", "whisper-1"); form.append("language", this.config.language); const response = await fetch(`${this.config.baseUrl}/audio/transcriptions`, { method: "POST", headers: this.config.apiKey ? { Authorization: `Bearer ${this.config.apiKey}` } : {}, body: form, signal: AbortSignal.timeout(4500), }); if (!response.ok) throw new Error(`Athena STT failed (HTTP ${response.status})`); const text = String(record(await response.json()).text || "").trim(); if (text) this.req.onTranscript?.(text); } catch (error) { this.req.onError?.(error instanceof Error ? error : new Error(String(error))); } } } function pcmRms(pcm: Buffer): number { 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: 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"); } let offset = 12; let sampleRate = 0; let channels = 0; let bits = 0; let data: Buffer | undefined; 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 { readonly supportsToolResultContinuation = false; readonly supportsToolResultSuppression = false; private connected = false; private closed = false; private speaking = false; private speech: Buffer[] = []; private prefix: Buffer[] = []; private prefixBytes = 0; private silenceTimer?: NodeJS.Timeout; private active?: AbortController; private transcribing = false; private generation = 0; constructor(private readonly req: any, private readonly cfg: Required) {} async connect(): Promise { this.connected = true; this.req.onEvent?.({ direction: "server", type: "session.created" }); this.req.onReady?.(); } isConnected(): boolean { return this.connected && !this.closed; } setMediaTimestamp(_timestamp: number): void {} sendAudio(chunk: Buffer): void { 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: string): void { const trimmed = text.trim(); if (trimmed) void this.answer(trimmed); } triggerGreeting(instructions?: string): void { void this.answer(instructions?.trim() || "Begrüße mich kurz auf Deutsch."); } submitToolResult(_callId: string, _result: unknown): void {} acknowledgeMark(_markName: string): void {} close(): void { this.closed = true; this.connected = false; this.generation += 1; this.active?.abort(); if (this.silenceTimer) clearTimeout(this.silenceTimer); this.req.onClose?.("completed"); } private async finishSpeech(): Promise { 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 as Error).name !== "AbortError") this.req.onError?.(error as Error); } finally { this.transcribing = false; } } private async transcribe(pcm: Buffer): Promise { 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(); } private async answer(prompt: string): Promise { 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 as Error).name === "AbortError") { this.req.onResponseDone?.({ status: "cancelled", responseId }); } else { this.req.onError?.(error as Error); this.req.onResponseDone?.({ status: "failed", responseId, error: String(error) }); } } finally { if (this.active === controller) this.active = undefined; } } private async synthesize(text: string, signal: AbortSignal): Promise { 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())); } private authHeaders(): Record { 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.registerRealtimeTranscriptionProvider({ id: "athena-talk", label: "Athena Whisper (Diktieren)", defaultModel: "whisper-1", models: ["whisper-1"], autoSelectOrder: 1, resolveConfig: ({ cfg, rawConfig }) => resolveTranscriptionConfig(cfg, rawConfig), isConfigured: ({ providerConfig }) => Boolean(record(providerConfig).baseUrl), createSession: (req) => new AthenaTranscriptionSession(req, resolveTranscriptionConfig(req.cfg, req.providerConfig)), }); api.registerHttpRoute({ path: BROWSER_KEY_PATH, auth: "plugin", match: "exact", handler: (req: IncomingMessage, res: ServerResponse) => { if (req.method !== "GET") { res.writeHead(405).end("Method not allowed"); } else { res.writeHead(200, { "Content-Type": "application/x-pem-file", "Cache-Control": "no-store" }); res.end(publicKeyPem); } return true; }, }); api.registerHttpRoute({ path: BROWSER_OFFER_PATH, auth: "plugin", match: "exact", handler: (req: IncomingMessage, res: ServerResponse) => { const raw = record(record(record(api.runtime.config.current()).talk).realtime); const provider = record(record(raw.providers)["athena-talk"]); return handleBrowserOffer(req, res, String(provider.realtimeUpstreamUrl || "")); }, }); api.registerRealtimeVoiceProvider({ id: "athena-talk", label: "Athena Local Talk", aliases: ["athena", "local-athena"], defaultModel: "athena-local", voices: ["alloy"], autoSelectOrder: 1, capabilities: { transports: ["gateway-relay", "webrtc"], inputAudioFormats: [AUDIO_FORMAT], outputAudioFormats: [AUDIO_FORMAT], supportsBrowserSession: true, supportsBargeIn: false, handlesInputAudioBargeIn: false, supportsToolCalls: true, supportsSessionResumption: false, }, resolveConfig: ({ rawConfig }: { rawConfig: unknown }) => record(rawConfig), isConfigured: ({ providerConfig, cfg }: { providerConfig: unknown; cfg: unknown }) => { const raw = record(providerConfig); return Boolean(raw.baseUrl || resolveModelProvider(cfg, raw.modelProvider).baseUrl); }, createBridge: (req: any) => new AthenaTalkBridge(req, resolveConfig(req)), async createBrowserSession(req: any) { const config = resolveConfig(req); if (!config.realtimeUpstreamUrl) { throw new Error("athena-talk realtimeUpstreamUrl is required for browser Talk"); } const { token, expiresAt } = createBrowserToken(); return { provider: "athena-talk", transport: "webrtc", clientSecret: token, offerUrl: BROWSER_OFFER_PATH, audio: { inputEncoding: "pcm16", inputSampleRateHz: 24000, outputEncoding: "pcm16", outputSampleRateHz: 24000, }, model: req.model || "athena-local", voice: req.voice || config.voice, expiresAt, }; }, } as any); }, });