661 lines
25 KiB
TypeScript
661 lines
25 KiB
TypeScript
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 DICTATION_SAMPLE_RATE_HZ = 8000;
|
|
const DICTATION_SEGMENT_BYTES = DICTATION_SAMPLE_RATE_HZ * 6;
|
|
const DICTATION_OVERLAP_BYTES = DICTATION_SAMPLE_RATE_HZ / 2;
|
|
const DICTATION_REQUEST_TIMEOUT_MS = 4500;
|
|
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<boolean> {
|
|
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<string, any> {
|
|
return value && typeof value === "object" && !Array.isArray(value)
|
|
? value as Record<string, any>
|
|
: {};
|
|
}
|
|
|
|
function resolveModelProvider(cfg: unknown, requested?: unknown): Record<string, any> {
|
|
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<ProviderConfig> {
|
|
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<ProviderConfig> {
|
|
const talkConfig = record(record(record(cfg).talk).realtime);
|
|
const talkProvider = record(record(talkConfig.providers)["athena-talk"]);
|
|
return resolveConfig({ cfg, providerConfig: { ...talkProvider, ...record(rawConfig) } });
|
|
}
|
|
|
|
function comparableWord(value: string): string {
|
|
return value.toLocaleLowerCase("de-DE").replace(/[^\p{L}\p{N}]+/gu, "");
|
|
}
|
|
|
|
function removeTranscriptOverlap(previous: string, current: string): string {
|
|
const priorWords = previous.trim().split(/\s+/).filter(Boolean);
|
|
const currentWords = current.trim().split(/\s+/).filter(Boolean);
|
|
const maximum = Math.min(12, priorWords.length, currentWords.length);
|
|
for (let count = maximum; count >= 1; count -= 1) {
|
|
const left = priorWords.slice(-count).map(comparableWord);
|
|
const right = currentWords.slice(0, count).map(comparableWord);
|
|
if (!left.every((word, index) => word && word === right[index])) continue;
|
|
// A single short word is too ambiguous to remove safely. Longer words are
|
|
// sufficient because the audio overlap is only half a second.
|
|
if (count === 1 && left[0].length < 5) continue;
|
|
return currentWords.slice(count).join(" ");
|
|
}
|
|
return currentWords.join(" ");
|
|
}
|
|
|
|
class AthenaTranscriptionSession {
|
|
private connected = false;
|
|
private closed = false;
|
|
private audio: Buffer[] = [];
|
|
private bufferedBytes = 0;
|
|
private totalBytes = 0;
|
|
private processing: Promise<void> = Promise.resolve();
|
|
private readonly completedTranscripts: string[] = [];
|
|
private emittedTranscripts = 0;
|
|
private processingError: Error | null = null;
|
|
|
|
constructor(
|
|
private readonly req: {
|
|
cfg?: unknown;
|
|
providerConfig: Record<string, unknown>;
|
|
onSpeechStart?: () => void;
|
|
onTranscript?: (text: string) => void;
|
|
onError?: (error: Error) => void;
|
|
},
|
|
private readonly config: Required<ProviderConfig>,
|
|
) {}
|
|
|
|
async connect(): Promise<void> { this.connected = true; }
|
|
isConnected(): boolean { return this.connected && !this.closed; }
|
|
|
|
sendAudio(audio: Buffer): void {
|
|
if (!this.isConnected() || audio.length === 0) return;
|
|
if (this.totalBytes === 0) this.req.onSpeechStart?.();
|
|
const maxBytes = this.config.maxSpeechSeconds * 8000;
|
|
if (this.totalBytes + 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.bufferedBytes += audio.length;
|
|
this.totalBytes += audio.length;
|
|
while (this.bufferedBytes >= DICTATION_SEGMENT_BYTES) {
|
|
this.queueFullSegment();
|
|
}
|
|
}
|
|
|
|
close(): void {
|
|
if (this.closed) return;
|
|
this.closed = true;
|
|
this.connected = false;
|
|
if (!this.totalBytes) return;
|
|
this.emitCompletedTranscripts();
|
|
const tail = Buffer.concat(this.audio);
|
|
this.audio = [];
|
|
this.bufferedBytes = 0;
|
|
if (tail.length > DICTATION_OVERLAP_BYTES || this.completedTranscripts.length === 0) {
|
|
this.queueTranscription(tail);
|
|
}
|
|
void this.processing.finally(() => {
|
|
this.emitCompletedTranscripts();
|
|
if (this.processingError && this.completedTranscripts.length === 0) {
|
|
this.req.onError?.(this.processingError);
|
|
}
|
|
});
|
|
}
|
|
|
|
private queueFullSegment(): void {
|
|
const buffered = Buffer.concat(this.audio);
|
|
const segment = Buffer.from(buffered.subarray(0, DICTATION_SEGMENT_BYTES));
|
|
const retained = Buffer.from(buffered.subarray(DICTATION_SEGMENT_BYTES - DICTATION_OVERLAP_BYTES));
|
|
this.audio = retained.length ? [retained] : [];
|
|
this.bufferedBytes = retained.length;
|
|
this.queueTranscription(segment);
|
|
}
|
|
|
|
private queueTranscription(audio: Buffer): void {
|
|
if (!audio.length) return;
|
|
this.processing = this.processing.then(async () => {
|
|
const previous = this.completedTranscripts.join(" ");
|
|
const text = await this.transcribe(audio, previous.slice(-240));
|
|
const novel = removeTranscriptOverlap(previous, text);
|
|
if (novel) this.completedTranscripts.push(novel);
|
|
if (this.closed) this.emitCompletedTranscripts();
|
|
}).catch((error: unknown) => {
|
|
this.processingError = error instanceof Error ? error : new Error(String(error));
|
|
});
|
|
}
|
|
|
|
private emitCompletedTranscripts(): void {
|
|
while (this.emittedTranscripts < this.completedTranscripts.length) {
|
|
this.req.onTranscript?.(this.completedTranscripts[this.emittedTranscripts]);
|
|
this.emittedTranscripts += 1;
|
|
}
|
|
}
|
|
|
|
private async transcribe(audio: Buffer, prompt: string): Promise<string> {
|
|
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);
|
|
if (prompt) form.append("prompt", prompt);
|
|
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(DICTATION_REQUEST_TIMEOUT_MS),
|
|
});
|
|
if (!response.ok) throw new Error(`Athena STT failed (HTTP ${response.status})`);
|
|
return String(record(await response.json()).text || "").trim();
|
|
}
|
|
}
|
|
|
|
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<ProviderConfig>) {}
|
|
|
|
async connect(): Promise<void> {
|
|
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<void> {
|
|
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<string> {
|
|
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<void> {
|
|
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<Buffer> {
|
|
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<string, string> {
|
|
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);
|
|
},
|
|
});
|