Streaming voice pipeline: sentence-level TTS, zero temp files
- playback.zig: sentence queue + playback thread — agent reply chunks are split on sentence boundaries and each sentence is synthesized+played as it completes, so voice starts after the first sentence, not after the reply - melo_server.py v2: WAV written to /dev/shm, raw s16le PCM + sample rate streamed over stdout (WAV n rate header + bytes), file deleted immediately; melo stdout chatter rerouted to stderr to keep the protocol clean - melo.zig: speakRaw returns in-memory PCM; pw-play fed via stdin pipe (no on-disk TTS files at all) - whisper: --language auto - fix: double-close of pw-play stdin panicked Threaded Io in debug - fix: quit-path frees for transcript cache, bubble lists, sentence queue, reply buffer (debug allocator leak panic) 👾 Generated with [Letta Code](https://letta.com) Co-Authored-By: Letta Code <noreply@letta.com>
This commit is contained in:
1 parent
2b2235ff86
commit
bcd48a3443
5 files changed
+244
-55
No files matched your search
+41
-13
@@ -5,17 +5,21 @@ Long-running subprocess: loads the MeloTTS model once, then reads lines of
|
||||
text on stdin and writes "<wav-path> <duration-seconds>" per line on stdout.
|
||||
stderr passes through for logging.
|
||||
|
||||
Protocol:
|
||||
Protocol (v2, raw PCM):
|
||||
-> one line of text to synthesize
|
||||
<- e.g. "/tmp/inferon/tts-1723.wav 3.42"
|
||||
<- header: "WAV <n_bytes> <sample_rate>\n" followed by <n_bytes> of raw s16le PCM
|
||||
|
||||
Kill with SIGTERM — no cleanup needed, WAVs live in /tmp.
|
||||
Synthesizes to /dev/shm (RAM-backed), streams bytes, deletes immediately.
|
||||
Kill with SIGTERM.
|
||||
"""
|
||||
|
||||
import sys
|
||||
import os
|
||||
import io
|
||||
import time
|
||||
import wave
|
||||
import signal
|
||||
import tempfile
|
||||
import warnings
|
||||
|
||||
warnings.filterwarnings("ignore")
|
||||
@@ -31,12 +35,25 @@ SPEED = 1.3 # speech rate; 1.0 = natural, higher = faster
|
||||
sys.path.insert(0, MELO_ROOT)
|
||||
|
||||
|
||||
def synth_bytes(path):
|
||||
"""Read a WAV, return (raw s16le PCM, sample_rate), delete the file."""
|
||||
with wave.open(path, "rb") as w:
|
||||
rate = w.getframerate()
|
||||
pcm = w.readframes(w.getnframes())
|
||||
os.unlink(path)
|
||||
return pcm, rate
|
||||
|
||||
|
||||
def main() -> None:
|
||||
signal.signal(signal.SIGTERM, lambda *_: sys.exit(0))
|
||||
signal.signal(signal.SIGINT, lambda *_: sys.exit(0))
|
||||
|
||||
os.makedirs(OUT_DIR, exist_ok=True)
|
||||
|
||||
# Protocol handle — melo's progress chatter goes to stdout, so we keep
|
||||
# this reference and then point stdout at stderr.
|
||||
proto_out = sys.stdout.buffer
|
||||
|
||||
# Heavy imports inside main so signal handlers are installed first.
|
||||
from melo.api import TTS
|
||||
|
||||
@@ -52,15 +69,24 @@ def main() -> None:
|
||||
else:
|
||||
speaker = speaker_ids[SPEAKER]
|
||||
|
||||
# melo's progress chatter goes to stdout — keep the protocol clean.
|
||||
sys.stdout = sys.stderr
|
||||
|
||||
# Warm up with a full sentence (longer text paths hit NLTK/num2words,
|
||||
# which is where missing-data errors surface). Fail before READY.
|
||||
warmup = "/tmp/inferon/tts-warmup.wav"
|
||||
model.tts_to_file("Warmup sentence number twelve, spoken on the first of January.",
|
||||
speaker, warmup, speed=SPEED)
|
||||
os.unlink(warmup)
|
||||
fd, wpath = tempfile.mkstemp(suffix=".wav", dir="/dev/shm")
|
||||
os.close(fd)
|
||||
try:
|
||||
model.tts_to_file("Warmup sentence number twelve, spoken on the first of January.",
|
||||
speaker, wpath, speed=SPEED)
|
||||
_ = synth_bytes(wpath) # exercise the full byte path before READY
|
||||
finally:
|
||||
if os.path.exists(wpath):
|
||||
os.unlink(wpath)
|
||||
|
||||
# Ready marker — inferon waits for this before declaring the backend up.
|
||||
print("READY", flush=True)
|
||||
proto_out.write(b"READY\n")
|
||||
proto_out.flush()
|
||||
|
||||
for line in sys.stdin:
|
||||
text = line.strip()
|
||||
@@ -68,12 +94,14 @@ def main() -> None:
|
||||
continue
|
||||
text = text[:2000] # protocol cap — extremely long replies get clipped
|
||||
|
||||
wav_path = os.path.join(OUT_DIR, f"tts-{time.time_ns()}.wav")
|
||||
try:
|
||||
start = time.monotonic()
|
||||
model.tts_to_file(text, speaker, wav_path, speed=SPEED)
|
||||
dur = time.monotonic() - start
|
||||
print(f"{wav_path} {dur:.2f}", flush=True)
|
||||
fd, wpath = tempfile.mkstemp(suffix=".wav", dir="/dev/shm")
|
||||
os.close(fd)
|
||||
model.tts_to_file(text, speaker, wpath, speed=SPEED)
|
||||
pcm, rate = synth_bytes(wpath)
|
||||
proto_out.write(f"WAV {len(pcm)} {rate}\n".encode())
|
||||
proto_out.write(pcm)
|
||||
proto_out.flush()
|
||||
except Exception as e: # noqa: BLE001 - protocol: report and survive
|
||||
print(f"melo_server: synthesis failed: {e}", file=sys.stderr, flush=True)
|
||||
print("ERROR", flush=True)
|
||||
|
||||
+35
-23
@@ -5,6 +5,7 @@ const std = @import("std");
|
||||
const Io = std.Io;
|
||||
const whisper = @import("whisper.zig");
|
||||
const melo = @import("melo.zig");
|
||||
const playback = @import("playback.zig");
|
||||
const qt6 = @import("libqt6zig");
|
||||
const QApplication = qt6.QApplication;
|
||||
const QSystemTrayIcon = qt6.QSystemTrayIcon;
|
||||
@@ -122,6 +123,12 @@ fn onLettaEvent(_: ?*anyopaque, ev: letta.Event) void {
|
||||
.chunk => |c| {
|
||||
const dup = allocator.dupe(u8, c) catch return;
|
||||
pushOverlay(.chunk, dup);
|
||||
// Stream voice: queue each sentence as it completes.
|
||||
sentence_buf.appendSlice(allocator, c) catch {};
|
||||
if (c.len > 0) {
|
||||
const last = c[c.len - 1];
|
||||
if (last == '.' or last == '!' or last == '?') flushSentence(false);
|
||||
}
|
||||
},
|
||||
}
|
||||
}
|
||||
@@ -186,20 +193,24 @@ fn setState(next: State) void {
|
||||
}
|
||||
}
|
||||
|
||||
/// Play a WAV via pw-play (blocking; audio length). Fine for v1 — same
|
||||
/// thread already blocks on whisper/letta. UI-thread decoupling comes later.
|
||||
fn playWav(path: []const u8) void {
|
||||
const io = io_ctx orelse return;
|
||||
var child = std.process.spawn(io, .{
|
||||
.argv = &.{ "pw-play", path },
|
||||
.stdin = .ignore,
|
||||
.stdout = .ignore,
|
||||
.stderr = .ignore,
|
||||
}) catch |e| {
|
||||
std.debug.print("inferon: pw-play spawn failed: {s}\n", .{@errorName(e)});
|
||||
/// Buffer for sentence-split streaming into the playback pipeline.
|
||||
var sentence_buf: std.ArrayList(u8) = .empty;
|
||||
|
||||
fn flushSentence(final_flush: bool) void {
|
||||
const a = allocator;
|
||||
if (sentence_buf.items.len == 0) {
|
||||
if (final_flush) playback.endReply();
|
||||
return;
|
||||
};
|
||||
_ = child.wait(io) catch {};
|
||||
}
|
||||
// Sentences end on . ! ? — hold back a bare trailing terminator-less tail
|
||||
// unless this is the end of the reply.
|
||||
if (!final_flush) {
|
||||
const last = sentence_buf.items[sentence_buf.items.len - 1];
|
||||
if (last != '.' and last != '!' and last != '?') return;
|
||||
}
|
||||
const s = sentence_buf.toOwnedSlice(a) catch return;
|
||||
playback.push(s);
|
||||
if (final_flush) playback.endReply();
|
||||
}
|
||||
|
||||
// --------------------------------------------------------------- capture ---
|
||||
@@ -321,17 +332,10 @@ fn processUtterance() void {
|
||||
defer allocator.free(response);
|
||||
std.debug.print("[INFERON]\n{s}\n\n", .{response});
|
||||
|
||||
// --- tts + playback ---
|
||||
// --- tts + playback: sentences already streaming; flush the tail ---
|
||||
setState(.speaking);
|
||||
var duration: f64 = 0;
|
||||
if (melo.speak(allocator, io, response, &duration)) |tts_wav| {
|
||||
defer allocator.free(tts_wav);
|
||||
std.debug.print("inferon: tts ready {s} ({d:.2}s synth)\n", .{ tts_wav, duration });
|
||||
playWav(tts_wav);
|
||||
cleanupFile(tts_wav);
|
||||
} else |e| {
|
||||
std.debug.print("inferon: tts failed: {s}\n", .{@errorName(e)});
|
||||
}
|
||||
flushSentence(true);
|
||||
playback.waitDrained();
|
||||
|
||||
// After speaking: back to idle-but-in-conversation. User clicks when
|
||||
// they want to talk again.
|
||||
@@ -476,12 +480,20 @@ pub fn main(init: std.process.Init) !void {
|
||||
};
|
||||
if (melo.waitReady(allocator, init.io)) {
|
||||
std.debug.print("inferon: melo ready\n", .{});
|
||||
playback.start(allocator, init.io);
|
||||
} else {
|
||||
std.debug.print("inferon: melo NOT ready\n", .{});
|
||||
}
|
||||
|
||||
_ = QApplication.exec();
|
||||
|
||||
// Session-lifetime buffers: free so the debug allocator doesn't squawk.
|
||||
if (last_user_text) |t| allocator.free(t);
|
||||
last_user_text = null;
|
||||
sentence_buf.deinit(allocator);
|
||||
overlay_mod.Overlay.shutdownMem(); // static, no instance needed
|
||||
playback.shutdown();
|
||||
|
||||
whisper.stop();
|
||||
melo.stop();
|
||||
}
|
||||
+50
-19
@@ -1,10 +1,9 @@
|
||||
//! MeloTTS subprocess lifecycle + line protocol client.
|
||||
//! MeloTTS subprocess lifecycle + raw-PCM protocol client.
|
||||
//!
|
||||
//! inferon owns the melo_server.py process: spawns it on startup (model load
|
||||
//! takes a few seconds), kills it on exit. Communication is line-based:
|
||||
//! takes a few seconds), kills it on exit. Protocol:
|
||||
//! inferon -> server: "<text>\n"
|
||||
//! server -> inferon: "<wav-path> <duration>\n" | "ERROR\n"
|
||||
//! The server prints "READY" once the model is loaded.
|
||||
//! server -> inferon: "WAV <n_bytes> <sample_rate>\n" + n raw s16le bytes
|
||||
|
||||
const std = @import("std");
|
||||
const Io = std.Io;
|
||||
@@ -101,14 +100,14 @@ fn sanitize(alloc: std.mem.Allocator, text: []const u8) ![]u8 {
|
||||
return out;
|
||||
}
|
||||
|
||||
/// Synthesize text: send one line, read one reply line.
|
||||
/// Returns the WAV path (allocated, caller frees) and sets `duration_s`.
|
||||
pub fn speak(alloc: std.mem.Allocator, io: Io, text: []const u8, duration_s: *f64) ![]u8 {
|
||||
pub const Audio = struct { pcm: []u8, rate: u32 }; // caller frees pcm
|
||||
|
||||
/// Synthesize text: send one line, read "WAV <n> <rate>" header + raw PCM.
|
||||
pub fn speakRaw(alloc: std.mem.Allocator, io: Io, text: []const u8) !Audio {
|
||||
const c = &(child orelse return error.NotRunning);
|
||||
const stdin = c.stdin orelse return error.NotRunning;
|
||||
const stdout = c.stdout orelse return error.NotRunning;
|
||||
|
||||
// Send request line (text sanitized: no embedded newlines/markdown).
|
||||
const clean = try sanitize(alloc, text);
|
||||
defer alloc.free(clean);
|
||||
var req_buf: [8192]u8 = undefined;
|
||||
@@ -118,26 +117,25 @@ pub fn speak(alloc: std.mem.Allocator, io: Io, text: []const u8, duration_s: *f6
|
||||
try writer.interface.writeAll(req);
|
||||
try writer.interface.flush();
|
||||
|
||||
// Read reply lines until we get a non-READY result line.
|
||||
var acc: [4096]u8 = undefined;
|
||||
// Header line.
|
||||
var acc: [256]u8 = undefined;
|
||||
var acc_len: usize = 0;
|
||||
while (true) {
|
||||
var header: []const u8 = "";
|
||||
while (header.len == 0) {
|
||||
var tmp: [128]u8 = undefined;
|
||||
const n = stdout.readStreaming(io, &.{&tmp}) catch return error.ReadFailed;
|
||||
if (n == 0) return error.ServerDied;
|
||||
for (tmp[0..n]) |ch| {
|
||||
var leftover: []const u8 = "";
|
||||
for (tmp[0..n], 0..) |ch, i| {
|
||||
if (ch == '\n') {
|
||||
const line = acc[0..acc_len];
|
||||
acc_len = 0;
|
||||
if (std.mem.eql(u8, line, "ERROR")) return error.SynthesisFailed;
|
||||
if (std.mem.startsWith(u8, line, "/tmp/")) {
|
||||
// "<path> <duration>"
|
||||
const sp = std.mem.lastIndexOfScalar(u8, line, ' ') orelse return error.BadReply;
|
||||
duration_s.* = std.fmt.parseFloat(f64, line[sp + 1 ..]) catch 0;
|
||||
return try alloc.dupe(u8, line[0..sp]);
|
||||
if (std.mem.startsWith(u8, line, "WAV ")) {
|
||||
header = try alloc.dupe(u8, line);
|
||||
leftover = tmp[i + 1 .. n]; // PCM may start in same read
|
||||
break;
|
||||
}
|
||||
// Anything else: log and keep reading.
|
||||
std.debug.print("inferon/melo: {s}\n", .{line});
|
||||
continue;
|
||||
}
|
||||
if (acc_len < acc.len) {
|
||||
@@ -145,5 +143,38 @@ pub fn speak(alloc: std.mem.Allocator, io: Io, text: []const u8, duration_s: *f6
|
||||
acc_len += 1;
|
||||
}
|
||||
}
|
||||
// PCM bytes already read past the header.
|
||||
if (header.len > 0 and leftover.len > 0) {
|
||||
return finishAudio(alloc, io, stdout, header, leftover);
|
||||
}
|
||||
}
|
||||
return finishAudio(alloc, io, stdout, header, "");
|
||||
}
|
||||
|
||||
fn finishAudio(
|
||||
alloc: std.mem.Allocator,
|
||||
io: Io,
|
||||
stdout: Io.File,
|
||||
header: []const u8,
|
||||
initial: []const u8,
|
||||
) !Audio {
|
||||
defer alloc.free(@constCast(header));
|
||||
// "WAV <n> <rate>"
|
||||
var it = std.mem.tokenizeScalar(u8, header, ' ');
|
||||
_ = it.next(); // "WAV"
|
||||
const n = std.fmt.parseInt(usize, it.next() orelse return error.BadReply, 10) catch return error.BadReply;
|
||||
const rate = std.fmt.parseInt(u32, it.next() orelse return error.BadReply, 10) catch return error.BadReply;
|
||||
|
||||
const pcm = try alloc.alloc(u8, n);
|
||||
var have: usize = 0;
|
||||
const pre = @min(initial.len, n);
|
||||
@memcpy(pcm[0..pre], initial[0..pre]);
|
||||
have = pre;
|
||||
while (have < n) {
|
||||
const got = stdout.readStreaming(io, &.{@constCast(pcm[have..])}) catch break;
|
||||
if (got == 0) break;
|
||||
have += got;
|
||||
}
|
||||
if (have != n) return error.ShortRead;
|
||||
return .{ .pcm = pcm, .rate = rate };
|
||||
}
|
||||
@@ -230,6 +230,15 @@ pub const Overlay = struct {
|
||||
self.window.hide();
|
||||
}
|
||||
|
||||
/// Free session-lifetime buffers (quit path).
|
||||
pub fn shutdownMem() void {
|
||||
const a = alloc orelse return;
|
||||
labels.clearAndFree(a);
|
||||
for (steps.items) |s| a.free(s.full);
|
||||
steps.clearAndFree(a);
|
||||
reply_buf.deinit(a);
|
||||
}
|
||||
|
||||
fn endReply() void {
|
||||
reply_label = null;
|
||||
reply_buf.clearRetainingCapacity();
|
||||
|
||||
@@ -0,0 +1,109 @@
|
||||
//! Sentence playback pipeline: synth + play sentences as they complete,
|
||||
//! so voice starts after the first sentence instead of the whole reply.
|
||||
|
||||
const std = @import("std");
|
||||
const Io = std.Io;
|
||||
const melo = @import("melo.zig");
|
||||
|
||||
var mutex: std.atomic.Mutex = .unlocked;
|
||||
|
||||
fn lock() void {
|
||||
while (!mutex.tryLock()) std.atomic.spinLoopHint();
|
||||
}
|
||||
var queue: std.ArrayList([]u8) = .empty; // owned sentences, FIFO
|
||||
var reply_done: bool = true; // current reply fully queued?
|
||||
var alloc: ?std.mem.Allocator = null;
|
||||
var io_ctx: ?Io = null;
|
||||
var thread: ?std.Thread = null;
|
||||
|
||||
pub fn start(a: std.mem.Allocator, io: Io) void {
|
||||
alloc = a;
|
||||
io_ctx = io;
|
||||
thread = std.Thread.spawn(.{}, run, .{}) catch null;
|
||||
}
|
||||
|
||||
/// Queue one owned sentence (caller allocates, we free).
|
||||
pub fn push(sentence: []u8) void {
|
||||
lock();
|
||||
defer mutex.unlock();
|
||||
queue.append(alloc.?, sentence) catch {
|
||||
alloc.?.free(sentence);
|
||||
return;
|
||||
};
|
||||
reply_done = false;
|
||||
}
|
||||
|
||||
/// Mark the current reply as fully queued.
|
||||
pub fn endReply() void {
|
||||
lock();
|
||||
defer mutex.unlock();
|
||||
reply_done = true;
|
||||
}
|
||||
|
||||
/// Free any queued-but-unplayed sentences (quit path).
|
||||
pub fn shutdown() void {
|
||||
const a = alloc orelse return;
|
||||
lock();
|
||||
defer mutex.unlock();
|
||||
for (queue.items) |s| a.free(s);
|
||||
queue.clearAndFree(a);
|
||||
}
|
||||
|
||||
/// Block until every queued sentence has been played.
|
||||
pub fn waitDrained() void {
|
||||
while (true) {
|
||||
lock();
|
||||
const empty = queue.items.len == 0 and reply_done;
|
||||
mutex.unlock();
|
||||
if (empty) return;
|
||||
io_ctx.?.sleep(.{ .nanoseconds = 50 * std.time.ns_per_ms }, .real) catch {};
|
||||
}
|
||||
}
|
||||
|
||||
fn run() void {
|
||||
const a = alloc orelse return;
|
||||
const io = io_ctx orelse return;
|
||||
while (true) {
|
||||
lock();
|
||||
const sentence: ?[]u8 = if (queue.items.len > 0) queue.orderedRemove(0) else null;
|
||||
mutex.unlock();
|
||||
|
||||
const s = sentence orelse {
|
||||
io.sleep(.{ .nanoseconds = 30 * std.time.ns_per_ms }, .real) catch {};
|
||||
continue;
|
||||
};
|
||||
defer a.free(s);
|
||||
|
||||
if (melo.speakRaw(a, io, s)) |audio| {
|
||||
defer a.free(audio.pcm);
|
||||
playRaw(io, audio.pcm, audio.rate);
|
||||
} else |e| {
|
||||
std.debug.print("inferon/playback: synth failed: {s}\n", .{@errorName(e)});
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/// Pipe raw s16le mono PCM straight into pw-play's stdin — no temp files.
|
||||
fn playRaw(io: Io, pcm: []const u8, rate: u32) void {
|
||||
var rate_buf: [16]u8 = undefined;
|
||||
const rate_str = std.fmt.bufPrint(&rate_buf, "{d}", .{rate}) catch return;
|
||||
|
||||
var child = std.process.spawn(io, .{
|
||||
.argv = &.{ "pw-play", "--raw", "--format", "s16", "--rate", rate_str, "--channels", "1", "-" },
|
||||
.stdin = .pipe,
|
||||
.stdout = .ignore,
|
||||
.stderr = .ignore,
|
||||
}) catch |e| {
|
||||
std.debug.print("inferon/playback: pw-play spawn failed: {s}\n", .{@errorName(e)});
|
||||
return;
|
||||
};
|
||||
const stdin = child.stdin orelse return;
|
||||
|
||||
var wbuf: [64 * 1024]u8 = undefined;
|
||||
var writer = stdin.writer(io, &wbuf);
|
||||
writer.interface.writeAll(pcm) catch {};
|
||||
writer.interface.flush() catch {};
|
||||
stdin.close(io); // EOF so pw-play finishes
|
||||
child.stdin = null; // already closed — stop wait() from closing again
|
||||
_ = child.wait(io) catch {};
|
||||
}
|
||||
Reference in new issue
Block a user