v0.8: TCP transport via gatuna tunnel, session IDs, pre-synth playback, server Rust rewrite

This commit is contained in:
2026-08-13 09:58:33 +00:00
parent ecf00c1bc4
commit facbfe6a5c
16 changed files with 955 additions and 638 deletions
+68 -98
View File
@@ -1,136 +1,106 @@
# Robovoice DHCP Tunnel Protocol
# Robovoice STT Protocol
## Overview
Robovoice communicates with a remote STT server by tunneling through the
DHCP UDP ports (68→67). This exploits a common killswitch exception: VPN
software (e.g. WireGuard) blocks all traffic except DHCP, which is allowed
for network connectivity maintenance.
Robovoice connects to the STT server over TCP (typically through a gatuna
L2 tunnel). The server captures audio, runs Moonshine STT, and sends
transcript segments back. The client pre-synthesizes TTS on segments and
plays audio on final.
```
[Robovoice client] --broadcast UDP :68→:67--> [STT server]
[Robovoice client] <--unicast UDP :67→:68-- [STT server]
[Robovoice client] --TCP--> [STT server 127.0.0.1:6996]
│ │
├── ON <session>\n ────────►│ (abort old, start new session)
├── OFF <session>\n ────────►│ (stop, final STT pass)
│◄── P <session> <text>\n ──┤ (completed segment)
│◄── F <session> <text>\n ──┤ (all done; text may be empty)
```
The client broadcasts NOP heartbeats while PTT is held. The server starts
recording on the first NOP and stops when it receives OFF or when 150ms
pass with no NOPs.
## Transport
- **Protocol:** UDP (connectionless, unreliable)
- **Client → Server:** broadcast, source port 68, dest port 67
- **Server → Client:** unicast, source port 67, dest port 68
- **Client binds:** to a specific LAN interface IP on port 68 (with
`SO_REUSEADDR` to coexist with the Windows DHCP service)
- **No connection state** — purely fire-and-forget datagrams
- **Protocol:** TCP (reliable, ordered, connection-oriented)
- **Server:** `127.0.0.1:6996` (hardcoded loopback)
- **Framing:** newline-delimited text (`\n`), UTF-8
- **Auto-reconnect:** client retries every 3s if connection drops
## Wire format
All messages are plain text, newline-terminated (`\n`). Every message starts
with the 6-byte magic `HKMSTR` to distinguish our traffic from real DHCP.
### Client → Server
**NOP (heartbeat while PTT held):**
**ON (PTT pressed):**
```
HKMSTR <session> <nonce>\n
ON <session>\n
```
Sent every 50ms while PTT is held. The session is an incrementing unsigned
integer that identifies the current PTT utterance (incremented on each PTT
press). The nonce is an incrementing unsigned integer that makes each
datagram unique. The server should echo the session back in replies. Both
are discarded by the server for protocol logic — the server tracks liveness
via "did anything arrive recently."
Starts a new STT session. The server aborts any active session and starts
recording. `<session>` is an incrementing unsigned integer chosen by the
client. Replies from the server echo this session ID.
**OFF (PTT released):**
```
HKMSTR:OFF <session> <nonce>\n
OFF <session>\n
```
Sent once when PTT is released. This is the fast-stop signal. If lost, the
150ms timeout acts as a backstop.
Stops the session. The server does a final STT pass on remaining audio and
sends any new segments followed by `F`.
### Server → Client
The server echoes the session ID from the NOPs in all replies. The client
drops any reply with a stale session ID.
**Partial segment (completed VAD segment):**
**Segment (completed VAD segment):**
```
HKMSTR:P <session> <text>\n
P <session> <text>\n
```
A completed, VAD-separated utterance segment. The client starts TTS
synthesis immediately and buffers the audio output, but does **not** play
it yet. Playback starts when `:F` arrives (or timeout).
synthesis immediately and buffers the audio (does not play yet).
**Final (all done):**
```
HKMSTR:F <session> <text>\n
F <session> <text>\n
```
Signals that all segments have been sent. May be empty
(`HKMSTR:F <session>\n`). Triggers playback of all buffered audio on the
client. If `<text>` is non-empty, the client synthesizes it before playing.
Signals all segments have been sent. `<text>` may be empty (`F <session>\n`).
Triggers playback of all buffered audio on the client. If text is non-empty,
the client synthesizes it before playing.
The purpose of this design is to minimize latency: TTS synthesis runs in
parallel with recording, so by the time `:F` arrives, audio is already
buffered and playback starts immediately.
## Session IDs
- Client increments session ID on each PTT press
- Server echoes the session ID in all replies for that session
- Client drops any reply with a stale session ID (handles the race where
stale segments from an aborted session are still in the TCP buffer)
- Server aborts old session on receiving `ON` with a new session ID
## Client playback model
1. `P` arrives → start TTS synthesis immediately, buffer audio (don't play)
2. More `P` arrive → keep synthesizing and buffering
3. `F` arrives → play all buffered audio immediately
4. PTT pressed → flush: stop playback, cancel synthesis, clear buffers
The purpose of pre-synthesis is to minimize latency between PTT release
and audio playback. By the time `F` arrives, audio is already buffered.
## Server state machine
```
┌──────────────────────────────────────────┐
▼ │
┌──────────┐ first NOP ┌──────────────┐
IDLE │ ──────────► │ RECORDING
└──────────┘ └──────────────┘
OFF │ │ 150ms │
recv'd │ silence
▼ ▼ │
┌─────────────┐
PROCESSING │
└─────────────┘
│ │
send │ │
:P/:F │ │
▼ │
back to IDLE ──────────┘
┌──────────┐ ON <session> ┌──────────────┐
IDLE──────────────► │ RECORDING
└──────────┘ └──────────────┘
OFF │ │
recv'd │
┌─────────────┐
│ PROCESSING
└─────────────┘
send
P/F
back to IDLE
```
- **IDLE → RECORDING:** first NOP received, start mic capture
- **RECORDING:** VAD detects completed segments → send `HKMSTR:P <text>`
- **RECORDING → PROCESSING:** OFF received, OR 150ms since last NOP
- **PROCESSING → IDLE:** send remaining segments as `:P`, then `HKMSTR:F`
- **IDLE → RECORDING:** `ON <session>` received, start mic capture
- **RECORDING:** Moonshine streaming produces completed segments → send `P`
- **RECORDING → PROCESSING:** `OFF <session>` received
- **PROCESSING → IDLE:** final STT pass, send remaining `P` + `F`
## Client playback model
1. `:P` arrives → start TTS synthesis immediately, buffer audio (don't play).
Reset 250ms segment timer.
2. More `:P` arrive → keep synthesizing and buffering, reset timer each time.
3. `:F` arrives → play all buffered audio immediately, cancel timer.
4. If `:F` doesn't arrive within 250ms of the last `:P` → play buffered audio
early. If more `:P` arrive after early playback, synthesis continues and
new audio is appended to the output — not a failure.
5. `:F` may be empty — it just signals "all segments sent, start/confirm playback."
## Timing
| Parameter | Value | Purpose |
|-----------|-------|---------|
| NOP interval | 50ms | Heartbeat frequency while PTT held |
| Silence timeout | 150ms | Stop recording if no NOPs (3 missed = lost OFF) |
| NOP bandwidth | ~20 msg/s × ~20 bytes | ~400 bytes/s — negligible |
## Why this works
1. **Outbound broadcast `:68→:67` to `255.255.255.255`** passes the
WireGuard WFP killswitch (DHCP exception matches this exact pattern)
2. **Inbound `:67→:68`** has no address restriction in the WFP rule, so
unicast replies pass through
3. **Binding to a specific interface IP** (not `0.0.0.0`) wins unicast
delivery over the Windows DHCP client service
4. **NOP spam** ensures the ON message gets through even at 5% packet loss
(3 consecutive NOPs = ~0.01% drop probability)
5. **150ms timeout** is the backstop for lost OFF — at 50ms intervals, 3
consecutive NOPs must all be lost to false-stop
If `ON` arrives while recording, the current session is aborted (no final
flush) and a new session starts immediately.
+18
View File
@@ -0,0 +1,18 @@
[package]
name = "robovoice-stt-server"
version = "0.1.0"
edition = "2021"
[[bin]]
name = "robovoice-stt-server"
path = "src/main.rs"
[dependencies]
anyhow = "1"
cpal = "0.15"
[build-dependencies]
bindgen = "0.71"
[profile.release]
opt-level = 3
+27
View File
@@ -0,0 +1,27 @@
use std::path::PathBuf;
fn main() {
let manifest_dir = PathBuf::from(env!("CARGO_MANIFEST_DIR"));
let lib_dir = manifest_dir.join("..").join("bench").join("moonshine-voice").join("lib");
let include_dir = manifest_dir.join("..").join("bench").join("moonshine-voice").join("include");
let header = include_dir.join("moonshine-c-api.h");
println!("cargo:rerun-if-changed={}", header.display());
println!("cargo:rustc-link-search=native={}", lib_dir.display());
println!("cargo:rustc-link-lib=dylib=moonshine");
println!("cargo:rustc-link-arg=-Wl,-rpath,{}", lib_dir.display());
let bindings = bindgen::Builder::default()
.header(header.to_str().unwrap())
.allowlist_function("moonshine_.*")
.allowlist_var("MOONSHINE_.*")
.allowlist_type("transcript.*|moonshine_option_t|speaker_span_t|transcript_word_t")
.derive_default(true)
.generate()
.expect("Unable to generate moonshine bindings");
let out_path = PathBuf::from(std::env::var("OUT_DIR").unwrap());
bindings
.write_to_file(out_path.join("moonshine_bindings.rs"))
.expect("Couldn't write bindings");
}
+478
View File
@@ -0,0 +1,478 @@
use anyhow::{anyhow, Result};
use cpal::traits::{DeviceTrait, HostTrait, StreamTrait};
use cpal::{SampleFormat, SampleRate};
use std::collections::HashSet;
use std::ffi::CStr;
use std::io::{BufRead, BufReader, Write};
use std::net::{TcpListener, TcpStream};
use std::sync::atomic::{AtomicBool, Ordering};
use std::sync::{Arc, Mutex};
use std::thread;
use std::time::{Duration, SystemTime, UNIX_EPOCH};
include!(concat!(env!("OUT_DIR"), "/moonshine_bindings.rs"));
const SAMPLE_RATE: i32 = 16000;
const HEADER_VERSION: i32 = 30000;
const ARCH: u32 = 5; // MOONSHINE_MODEL_ARCH_MEDIUM_STREAMING
const BIND_ADDR: &str = "127.0.0.1:6996";
const MAX_TEXT_BYTES: usize = 1380;
// ─── helpers ──────────────────────────────────────────────────────────────
fn ts() -> String {
let now = SystemTime::now()
.duration_since(UNIX_EPOCH)
.unwrap_or_default();
let secs = now.as_secs() % 86400;
let h = secs / 3600;
let m = (secs % 3600) / 60;
let s = secs % 60;
let ms = now.subsec_millis();
format!("{:02}:{:02}:{:02}.{:03}", h, m, s, ms)
}
fn log(msg: &str) {
eprintln!("[{}] {}", ts(), msg);
}
fn err_str(code: i32) -> String {
unsafe {
let s = moonshine_error_to_string(code);
if s.is_null() {
format!("error {}", code)
} else {
CStr::from_ptr(s).to_string_lossy().into_owned()
}
}
}
fn truncate_to_word(text: &str, max_bytes: usize) -> &str {
if text.len() <= max_bytes {
return text;
}
let cut = &text[..max_bytes.min(text.len())];
match cut.rfind(' ') {
Some(pos) => &text[..pos],
None => cut,
}
}
fn line_text(line: &transcript_line_t) -> String {
if line.text.is_null() {
return String::new();
}
unsafe { CStr::from_ptr(line.text) }
.to_string_lossy()
.into_owned()
}
// ─── shared state ─────────────────────────────────────────────────────────
struct Shared {
writer: Mutex<TcpStream>,
session_id: u64,
transcriber_handle: i32,
}
impl Shared {
fn send_msg(&self, prefix: &str, text: &str) {
let text = truncate_to_word(text, MAX_TEXT_BYTES);
let line = if text.is_empty() {
format!("{} {}\n", prefix, self.session_id)
} else {
format!("{} {} {}\n", prefix, self.session_id, text)
};
let mut writer = self.writer.lock().unwrap();
match writer.write_all(line.as_bytes()) {
Ok(_) => log(&format!("TX {} {} {}", prefix, self.session_id, text)),
Err(e) => log(&format!("TX failed: {}", e)),
}
}
}
// ─── session ──────────────────────────────────────────────────────────────
struct Session {
shared: Arc<Shared>,
stop_signal: Arc<AtomicBool>,
aborted: Arc<AtomicBool>,
transcriber: thread::JoinHandle<()>,
cpal_stream: cpal::Stream,
stream_handle: i32,
}
impl Session {
fn stop(self) {
self.stop_signal.store(true, Ordering::SeqCst);
drop(self.cpal_stream);
self.transcriber.join().ok();
unsafe { moonshine_free_stream(self.shared.transcriber_handle, self.stream_handle) };
}
fn abort(self) {
self.stop_signal.store(true, Ordering::SeqCst);
self.aborted.store(true, Ordering::SeqCst);
drop(self.cpal_stream);
self.transcriber.join().ok();
unsafe { moonshine_free_stream(self.shared.transcriber_handle, self.stream_handle) };
}
}
fn start_session(shared: Arc<Shared>) -> Option<Session> {
let stream_handle = unsafe { moonshine_create_stream(shared.transcriber_handle, 0) };
if stream_handle < 0 {
log(&format!("create_stream failed: {}", err_str(stream_handle)));
return None;
}
let rc = unsafe { moonshine_start_stream(shared.transcriber_handle, stream_handle) };
if rc != 0 {
log(&format!("start_stream failed: {}", err_str(rc)));
unsafe { moonshine_free_stream(shared.transcriber_handle, stream_handle) };
return None;
}
let audio_buf: Arc<Mutex<Vec<f32>>> = Arc::new(Mutex::new(Vec::new()));
let stop_signal = Arc::new(AtomicBool::new(false));
let aborted = Arc::new(AtomicBool::new(false));
let cpal_stream = match start_cpal(audio_buf.clone(), stop_signal.clone()) {
Ok(s) => s,
Err(e) => {
log(&format!("cpal failed: {}", e));
unsafe { moonshine_free_stream(shared.transcriber_handle, stream_handle) };
return None;
}
};
let shared_clone = shared.clone();
let stop_signal_clone = stop_signal.clone();
let aborted_clone = aborted.clone();
let transcriber = thread::spawn(move || {
transcriber_loop(shared_clone, audio_buf, stop_signal_clone, aborted_clone, stream_handle);
});
Some(Session {
shared,
stop_signal,
aborted,
transcriber,
cpal_stream,
stream_handle,
})
}
fn transcriber_loop(
shared: Arc<Shared>,
audio_buf: Arc<Mutex<Vec<f32>>>,
stop_signal: Arc<AtomicBool>,
aborted: Arc<AtomicBool>,
stream_handle: i32,
) {
let handle = shared.transcriber_handle;
let mut sent_ids: HashSet<u64> = HashSet::new();
while !stop_signal.load(Ordering::SeqCst) {
let chunk = {
let mut buf = audio_buf.lock().unwrap();
if buf.is_empty() {
drop(buf);
thread::sleep(Duration::from_millis(5));
continue;
}
std::mem::take(&mut *buf)
};
unsafe {
moonshine_transcribe_add_audio_to_stream(
handle, stream_handle,
chunk.as_ptr(), chunk.len() as u64,
SAMPLE_RATE, 0,
);
}
let mut t_ptr: *mut transcript_t = std::ptr::null_mut();
let rc = unsafe { moonshine_transcribe_stream(handle, stream_handle, 0, &mut t_ptr) };
if rc != 0 || t_ptr.is_null() {
continue;
}
send_new_segments(&shared, t_ptr, &mut sent_ids, "P");
}
// If aborted (new session took over), skip final flush entirely
if aborted.load(Ordering::SeqCst) {
unsafe { moonshine_stop_stream(handle, stream_handle) };
log(&format!("Session {} aborted, skipping final flush", shared.session_id));
return;
}
// Drain remaining audio
let remaining = {
let mut buf = audio_buf.lock().unwrap();
std::mem::take(&mut *buf)
};
if !remaining.is_empty() {
unsafe {
moonshine_transcribe_add_audio_to_stream(
handle, stream_handle,
remaining.as_ptr(), remaining.len() as u64,
SAMPLE_RATE, 0,
);
}
}
// Final flush
unsafe { moonshine_stop_stream(handle, stream_handle) };
let mut t_ptr: *mut transcript_t = std::ptr::null_mut();
let rc = unsafe { moonshine_transcribe_stream(handle, stream_handle, 0, &mut t_ptr) };
if rc == 0 && !t_ptr.is_null() {
let t = unsafe { &*t_ptr };
let mut new_segments: Vec<String> = Vec::new();
for i in 0..t.line_count as usize {
let line = unsafe { &*t.lines.add(i) };
if line.text.is_null() || line.is_complete == 0 {
continue;
}
if !sent_ids.insert(line.id) {
continue;
}
let text = line_text(line);
if !text.is_empty() {
new_segments.push(text);
}
}
if new_segments.is_empty() {
shared.send_msg("F", "");
} else {
let last = new_segments.len() - 1;
for (i, text) in new_segments.iter().enumerate() {
let prefix = if i == last { "F" } else { "P" };
shared.send_msg(prefix, text);
}
}
} else {
shared.send_msg("F", "");
}
}
fn send_new_segments(
shared: &Shared,
t_ptr: *const transcript_t,
sent_ids: &mut HashSet<u64>,
prefix: &str,
) {
let t = unsafe { &*t_ptr };
for i in 0..t.line_count as usize {
let line = unsafe { &*t.lines.add(i) };
if line.text.is_null() || line.is_complete == 0 {
continue;
}
if !sent_ids.insert(line.id) {
continue;
}
let text = line_text(line);
if text.is_empty() {
continue;
}
shared.send_msg(prefix, &text);
}
}
// ─── cpal ─────────────────────────────────────────────────────────────────
fn start_cpal(
audio_buf: Arc<Mutex<Vec<f32>>>,
stop_signal: Arc<AtomicBool>,
) -> Result<cpal::Stream> {
let host = cpal::default_host();
let dev = host
.default_input_device()
.ok_or_else(|| anyhow!("no input device"))?;
let supported = dev
.supported_input_configs()?
.filter(|c| c.channels() <= 2 && c.min_sample_rate().0 <= 16000)
.min_by_key(|c| match c.sample_format() {
SampleFormat::F32 => 0,
SampleFormat::I16 => 1,
SampleFormat::U8 => 2,
_ => 99,
})
.ok_or_else(|| anyhow!("no suitable input config"))?;
let fmt = supported.sample_format();
let mut config = supported.with_max_sample_rate().config();
if config.channels > 1 {
config.channels = 1;
}
config.sample_rate = SampleRate(16000);
let err_fn = |e: cpal::StreamError| log(&format!("cpal error: {}", e));
let stream = match fmt {
SampleFormat::F32 => dev.build_input_stream(
&config,
move |data: &[f32], _: &_| {
if !stop_signal.load(Ordering::Relaxed) {
audio_buf.lock().unwrap().extend_from_slice(data);
}
},
err_fn,
None,
)?,
SampleFormat::I16 => dev.build_input_stream(
&config,
move |data: &[i16], _: &_| {
if !stop_signal.load(Ordering::Relaxed) {
audio_buf.lock().unwrap().extend(data.iter().map(|&x| x as f32 / 32768.0));
}
},
err_fn,
None,
)?,
SampleFormat::U8 => dev.build_input_stream(
&config,
move |data: &[u8], _: &_| {
if !stop_signal.load(Ordering::Relaxed) {
audio_buf.lock().unwrap().extend(data.iter().map(|&x| (x as f32 - 128.0) / 128.0));
}
},
err_fn,
None,
)?,
_ => return Err(anyhow!("unsupported sample format {:?}", fmt)),
};
stream.play()?;
Ok(stream)
}
// ─── main ─────────────────────────────────────────────────────────────────
fn main() -> Result<()> {
let mut args = std::env::args().skip(1);
let mut model_dir = String::from("../bench/medium-streaming-en");
while let Some(a) = args.next() {
match a.as_str() {
"--model-dir" | "-m" => {
model_dir = args.next().unwrap_or(model_dir);
}
"--help" | "-h" => {
println!("Usage: robovoice-stt-server [--model-dir DIR]");
println!("Listens on TCP {}", BIND_ADDR);
println!("Model: medium-streaming (Moonshine)");
return Ok(());
}
_ => return Err(anyhow!("unknown arg: {}", a)),
}
}
let model_path = std::fs::canonicalize(&model_dir)
.unwrap_or_else(|_| std::path::PathBuf::from(&model_dir));
log(&format!("Loading model from {}...", model_path.display()));
let c_dir = std::ffi::CString::new(model_path.to_str().unwrap()).unwrap();
let transcriber_handle = unsafe {
moonshine_load_transcriber_from_files(
c_dir.as_ptr(),
ARCH,
std::ptr::null(),
0,
HEADER_VERSION,
)
};
if transcriber_handle < 0 {
return Err(anyhow!("failed to load model: {}", err_str(transcriber_handle)));
}
log(&format!("Model loaded (handle {})", transcriber_handle));
let listener = TcpListener::bind(BIND_ADDR)?;
log(&format!("STT server listening on TCP {}", BIND_ADDR));
let mut session_id_counter: u64 = 0;
let mut current_session: Option<Session> = None;
for stream in listener.incoming() {
let stream = match stream {
Ok(s) => s,
Err(e) => {
log(&format!("accept failed: {}", e));
continue;
}
};
stream.set_nodelay(true).ok();
log(&format!("Client connected: {}", stream.peer_addr().unwrap_or_default()));
let writer_stream = stream.try_clone()?;
let reader = BufReader::new(stream);
for line in reader.lines() {
let line = match line {
Ok(l) => l,
Err(_) => break,
};
let line = line.trim();
log(&format!("RX {}", line));
// ON <session>
if let Some(rest) = line.strip_prefix("ON ") {
let new_session_id: u64 = rest.parse().unwrap_or(0);
if let Some(s) = current_session.take() {
log(&format!("Aborting session {} for new session {}", s.shared.session_id, new_session_id));
s.abort();
}
session_id_counter = new_session_id;
let shared = Arc::new(Shared {
writer: Mutex::new(writer_stream.try_clone()?),
session_id: session_id_counter,
transcriber_handle,
});
log(&format!("PTT on session {}", session_id_counter));
match start_session(shared) {
Some(s) => current_session = Some(s),
None => log("Failed to start session"),
}
}
// OFF <session>
else if let Some(rest) = line.strip_prefix("OFF ") {
let off_session: u64 = rest.parse().unwrap_or(0);
if let Some(s) = current_session.as_ref() {
if s.shared.session_id == off_session {
log(&format!("OFF session {}", off_session));
if let Some(s) = current_session.take() {
s.stop();
}
} else {
log(&format!("OFF session {} (stale, current={}), ignoring", off_session, s.shared.session_id));
}
} else {
log(&format!("OFF session {} (no active session), ignoring", off_session));
}
}
else {
log(&format!("Unknown command: {}", line));
}
}
log("Client disconnected");
if let Some(s) = current_session.take() {
s.abort();
}
}
Ok(())
}