use super::routing::*;
#[cfg(feature = "voice")]
use super::synthesis::prebuffer_target_secs;
use super::synthesis::Synthesizer;
use super::{Channel, Voice};
use crate::out::Out;
#[cfg(feature = "voice")]
use crate::schemas::voice::CHANNEL_SAY;
use anyhow::{bail, Context, Result};
use rand_core::OsRng;
use std::path::{Path, PathBuf};
mod reporting;
mod retention;
pub use reporting::SpeechReportingError;
#[cfg(feature = "voice")]
use reporting::{report, retain_before_report};
use retention::retain_after_playback;
pub use retention::SpeechRetentionError;
pub const DEFAULT_DAEMON: &str = "http://localhost:8000";
#[cfg(feature = "voice")]
fn http() -> reqwest::blocking::Client {
reqwest::blocking::Client::builder()
.timeout(std::time::Duration::from_secs(30))
.build()
.expect("build voice HTTP client")
}
fn probe_http() -> reqwest::blocking::Client {
reqwest::blocking::Client::builder()
.connect_timeout(std::time::Duration::from_secs(1))
.timeout(std::time::Duration::from_secs(2))
.build()
.expect("build voice endpoint probe")
}
#[cfg(feature = "audio")]
pub fn detect_output_devices() -> Result<Vec<AudioDevice>> {
use rodio::cpal::traits::{DeviceTrait, HostTrait};
let host = rodio::cpal::default_host();
let default_name = host
.default_output_device()
.and_then(|d| d.description().ok().map(|desc| desc.name().to_string()));
let mut devices = Vec::new();
for dev in host
.output_devices()
.context("enumerate audio output devices (cpal)")?
{
let Ok(desc) = dev.description() else {
continue;
};
let name = desc.name().to_string();
if name.is_empty() {
continue;
}
devices.push(AudioDevice {
is_default_output: Some(&name) == default_name.as_ref(),
name,
});
}
Ok(devices)
}
#[cfg(not(feature = "audio"))]
pub fn detect_output_devices() -> Result<Vec<AudioDevice>> {
anyhow::bail!("audio device support not compiled into this build (enable the `audio` feature)")
}
#[cfg(feature = "voice")]
fn open_named_sink(name: &str) -> Result<(rodio::MixerDeviceSink, rodio::Player)> {
use rodio::cpal::traits::{DeviceTrait, HostTrait};
let host = rodio::cpal::default_host();
let device = host
.output_devices()
.context("enumerate audio output devices (cpal)")?
.find(|d| {
d.description()
.map(|desc| desc.name() == name)
.unwrap_or(false)
})
.with_context(|| format!("output device '{name}' not found (disconnected?)"))?;
let mut sink = rodio::DeviceSinkBuilder::from_device(device)
.and_then(|b| b.open_stream())
.map_err(|e| anyhow::anyhow!("open audio stream on '{name}': {e}"))?;
sink.log_on_drop(false); let player = rodio::Player::connect_new(sink.mixer());
Ok((sink, player))
}
#[cfg(feature = "voice")]
fn play_on_reachy(daemon: &str, wav: &Path) -> Result<()> {
let bytes = std::fs::read(wav)?;
let fname = wav.file_name().unwrap().to_string_lossy().to_string();
let part = reqwest::blocking::multipart::Part::bytes(bytes)
.file_name(fname.clone())
.mime_str("audio/wav")?;
let form = reqwest::blocking::multipart::Form::new().part("file", part);
let resp = http()
.post(format!("{daemon}/api/media/sounds/upload"))
.multipart(form)
.send()
.context("upload to Reachy daemon")?;
if !resp.status().is_success() {
bail!("Reachy upload failed: {}", resp.text().unwrap_or_default());
}
let resp = http()
.post(format!("{daemon}/api/media/play_sound"))
.json(&serde_json::json!({ "file": fname }))
.send()
.context("Reachy play_sound")?;
if !resp.status().is_success() {
bail!(
"Reachy play_sound failed: {}",
resp.text().unwrap_or_default()
);
}
Ok(())
}
fn reachy_reachable(daemon: &str) -> bool {
probe_http()
.get(format!("{daemon}/api/daemon/status"))
.send()
.map(|r| r.status().is_success())
.unwrap_or(false)
}
pub(super) fn soma_reachable(soma: &str) -> bool {
probe_http()
.get(format!("{}/state", soma.trim_end_matches('/')))
.send()
.map(|response| response.status().is_success())
.unwrap_or(false)
}
#[cfg_attr(not(feature = "voice"), allow(dead_code))]
enum Spoken {
Played(SpeechDisposition),
PlaybackFailed(anyhow::Error),
ReportingFailed(anyhow::Error),
}
#[cfg(feature = "voice")]
fn speak_and_play(
routed: &Routed,
daemon: &str,
channel: &str,
text: &str,
out: &Path,
synthesizer: &Synthesizer,
output: &mut Out<'_>,
) -> Result<Spoken> {
let super::synthesis::PreparedSpeech {
mut stream,
sample_rate: sr,
estimated_seconds: est_secs,
started: t_call,
} = synthesizer.start(text)?;
let mut samples: Vec<f32> = Vec::new();
let played: Result<()> = match routed {
Routed::Reachy => {
for chunk in stream.by_ref() {
samples.extend_from_slice(&chunk);
}
Ok(())
}
Routed::Soma(endpoint) => stream_to_soma(&mut stream, &mut samples, endpoint, sr, output),
Routed::Devices(ladder) => stream_to_device(
&mut stream,
&mut samples,
t_call,
channel,
ladder,
sr,
est_secs,
output,
),
Routed::Text(_) => Ok(()), };
for chunk in stream.by_ref() {
samples.extend_from_slice(&chunk);
}
stream.finish()?;
super::synthesis::AudioClip::from_samples(&samples, sr)?.save(out)?;
if let Err(e) = played {
return Ok(if e.is::<SpeechReportingError>() {
Spoken::ReportingFailed(e)
} else {
Spoken::PlaybackFailed(e)
});
}
if matches!(routed, Routed::Reachy) {
if let Err(e) = play_on_reachy(daemon, out) {
return Ok(Spoken::PlaybackFailed(e));
}
}
Ok(Spoken::Played(match routed {
Routed::Soma(_) => SpeechDisposition::SomaPlaybackDrained,
Routed::Reachy => SpeechDisposition::ReachyRequestAccepted,
Routed::Devices(_) => SpeechDisposition::LocalPlaybackDrained,
Routed::Text(_) => unreachable!("text fallback never reaches playback"),
}))
}
#[cfg(feature = "voice")]
fn stream_to_soma(
stream: &mut mary::speak::SpeakStream,
samples: &mut Vec<f32>,
endpoint: &str,
sample_rate: u32,
output: &mut Out<'_>,
) -> Result<()> {
let mut playback =
soma_client::SomaPlayback::open(endpoint, soma_client::PlaybackSpec::mono(sample_rate))?;
for chunk in stream.by_ref() {
samples.extend_from_slice(&chunk);
playback.push_f32(&chunk)?;
}
let receipt = playback.finish()?;
let audio_seconds = receipt.samples as f64 / sample_rate as f64;
report(
output,
format!(
" [stream] Soma drained {audio_seconds:.1}s ({} samples; {} underrun; \
{} callbacks, {} late, max {:.2}ms)",
receipt.samples,
receipt.underrun_samples,
receipt.callbacks,
receipt.late_callbacks,
receipt.max_callback_ms,
),
Some(SpeechDisposition::SomaPlaybackDrained),
)?;
Ok(())
}
#[cfg(feature = "voice")]
#[allow(clippy::too_many_arguments)]
fn stream_to_device(
stream: &mut mary::speak::SpeakStream,
samples: &mut Vec<f32>,
t_call: std::time::Instant,
channel: &str,
ladder: &[String],
sr: u32,
est_secs: f32,
output: &mut Out<'_>,
) -> Result<()> {
use rodio::buffer::SamplesBuffer;
use std::num::NonZero;
let mut opened = None;
for name in ladder {
if channel == CHANNEL_SAY && classify(name) != DeviceClass::Private {
eprintln!(
" [stream] refusing non-private device '{name}' on the say channel \
(privacy invariant)"
);
continue;
}
match open_named_sink(name) {
Ok(sink) => {
opened = Some((sink, name.as_str()));
break;
}
Err(e) => eprintln!(" [stream] could not open '{name}': {e:#} — trying next device"),
}
}
let Some(((_device_sink, player), device)) = opened else {
bail!(
"no device in the routing ladder could be opened: {}",
ladder.join(" → ")
);
};
let mono = NonZero::new(1).expect("1 is nonzero");
let sr_nz = NonZero::new(sr).context("PCM sample rate must be nonzero")?;
let secs = |n: usize| n as f32 / sr as f32;
player.pause();
let mut appended = 0usize;
let mut started = false;
let mut chunks = 0usize;
let mut measure_from = 0usize; let mut t_first: Option<std::time::Instant> = None;
let mut prod_rate = 1.0f32; let mut announced = false;
let mut rebuffer_from: Option<usize> = None; for chunk in stream.by_ref() {
retain_before_report(samples, &chunk, || {
if started && rebuffer_from.is_none() && player.empty() {
player.pause();
let target = prebuffer_target_secs((est_secs - secs(appended)).max(0.0), prod_rate);
report(output, format!(
" [stream] underrun at {:.1}s — rebuffering {:.1}s (synthesis at {:.2}x realtime)",
secs(appended), target, prod_rate
), None)?;
rebuffer_from = Some(appended);
}
Ok(())
})?;
appended += chunk.len();
player.append(SamplesBuffer::new(mono, sr_nz, chunk));
chunks += 1;
match t_first {
None => {
t_first = Some(std::time::Instant::now());
measure_from = appended;
}
Some(t0) => {
prod_rate = secs(appended - measure_from) / t0.elapsed().as_secs_f32().max(1e-3);
}
}
if !started && chunks >= 2 {
let target = prebuffer_target_secs(est_secs, prod_rate);
if secs(appended) >= target {
player.play();
started = true;
report(
output,
format!(
" [stream] TTFA {:.2}s → {device}",
t_call.elapsed().as_secs_f32()
),
None,
)?;
} else if !announced {
report(
output,
format!(
" [stream] buffering {:.1}s of ~{:.0}s (synthesis at {:.2}x realtime)",
target, est_secs, prod_rate
),
None,
)?;
announced = true;
}
}
if let Some(from) = rebuffer_from {
let target = prebuffer_target_secs((est_secs - secs(from)).max(0.0), prod_rate);
if secs(appended - from) >= target {
player.play();
rebuffer_from = None;
report(
output,
format!(
" [stream] resumed with {:.1}s rebuffered",
secs(appended - from)
),
None,
)?;
}
}
}
if appended == 0 {
bail!("no audio chunks arrived to play on '{device}'");
}
if !started {
player.play();
report(
output,
format!(
" [stream] TTFA {:.2}s → {device}",
t_call.elapsed().as_secs_f32()
),
None,
)?;
} else if rebuffer_from.is_some() {
player.play(); }
let audio_secs = appended as f32 / sr as f32;
let deadline = std::time::Instant::now() + std::time::Duration::from_secs_f32(audio_secs + 5.0);
while !player.empty() {
if std::time::Instant::now() > deadline {
bail!("playback stalled on '{device}' ({audio_secs:.1}s of audio never drained)");
}
std::thread::sleep(std::time::Duration::from_millis(50));
}
std::thread::sleep(std::time::Duration::from_millis(100));
report(
output,
format!(" [stream] played {audio_secs:.1}s on {device} ({appended} samples)"),
Some(SpeechDisposition::LocalPlaybackDrained),
)?;
Ok(())
}
#[cfg(not(feature = "voice"))]
fn speak_and_play(
_routed: &Routed,
_daemon: &str,
_channel: &str,
_text: &str,
_out: &Path,
_synthesizer: &Synthesizer,
_output: &mut Out<'_>,
) -> Result<Spoken> {
bail!(
"voice was built without the `voice` feature — rebuild with \
`cargo build --release --features voice --bin voice` (pulls mary's \
Qwen3-TTS Burn voice pipeline). Routing (`voice route`/`voice devices`) \
and the text-fallback path work without it."
);
}
fn unique_voice_tmp() -> Result<PathBuf> {
use rand_core::RngCore;
for _ in 0..16 {
let mut r = [0u8; 8];
OsRng.fill_bytes(&mut r);
let path = std::env::temp_dir().join(format!(
"voice_out_{}_{:016x}.wav",
std::process::id(),
u64::from_le_bytes(r)
));
match std::fs::OpenOptions::new()
.write(true)
.create_new(true)
.open(&path)
{
Ok(_) => return Ok(path),
Err(e) if e.kind() == std::io::ErrorKind::AlreadyExists => continue,
Err(e) => {
return Err(e).with_context(|| format!("create temp wav {}", path.display()));
}
}
}
bail!(
"could not mint a unique temp wav in {}",
std::env::temp_dir().display()
);
}
#[derive(Clone, Debug)]
pub struct RouteReport {
pub devices: Vec<AudioDevice>,
pub daemon_up: bool,
pub soma: Option<String>,
pub soma_up: bool,
pub routes: Vec<(super::RoutePolicy, Routed)>,
}
impl RouteReport {
pub fn emit(&self, out: &mut Out<'_>) -> Result<()> {
out.line(format!(
"Reachy daemon: {}",
if self.daemon_up { "reachable" } else { "down" }
))?;
match &self.soma {
Some(soma) => out.line(format!(
"Soma at {soma}: {}",
if self.soma_up { "reachable" } else { "down" }
))?,
None => out.line("Soma: not configured")?,
}
out.line("")?;
out.line("connected output devices:")?;
for device in &self.devices {
out.line(format!(" {}", describe_device(device)))?;
}
out.line("")?;
for (policy, routed) in &self.routes {
out.line(format!(
"{} policy (priority order): {}",
policy.channel.name(),
policy.devices.join(" → ")
))?;
out.line(format!(" would route to: {}", routed.describe()))?;
}
Ok(())
}
}
pub fn describe_device(device: &AudioDevice) -> String {
format!(
"{:<28} {}{}",
device.name,
device.class().label(),
if device.is_default_output {
" [default output]"
} else {
""
}
)
}
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
pub enum SpeechDisposition {
DryRun,
TextFallback,
LocalPlaybackDrained,
SomaPlaybackDrained,
ReachyRequestAccepted,
}
#[derive(Clone, Debug)]
pub struct SpeechReceipt {
pub route: Routed,
pub disposition: SpeechDisposition,
pub utterance: Option<triblespace::prelude::Id>,
}
impl SpeechReceipt {
pub fn emit(&self, channel: Channel, out: &mut Out<'_>) -> Result<()> {
if self.disposition == SpeechDisposition::ReachyRequestAccepted {
out.line(" Reachy accepted the playback request; completion was not observed")?;
}
if let Some(id) = self.utterance {
out.line(format!(" logged utterance {:x} [{}]", id, channel.name()))?;
}
Ok(())
}
}
#[derive(Clone, Debug)]
pub struct Device {
pub voice: Voice,
pub daemon: String,
pub soma: Option<String>,
pub synthesizer: Synthesizer,
}
impl Device {
pub fn new(
voice: Voice,
daemon: String,
soma: Option<String>,
synthesizer: Synthesizer,
) -> Self {
Self {
voice,
daemon,
soma,
synthesizer,
}
}
pub fn routes(&self) -> Result<RouteReport> {
let policies = self.voice.routes()?;
let devices = detect_output_devices()?;
let daemon_up = reachy_reachable(&self.daemon);
let soma_up = self.soma.as_deref().is_some_and(soma_reachable);
let routes = policies
.into_iter()
.map(|policy| {
let routed = match policy.channel {
Channel::Say => route_say(&policy.devices, &devices),
Channel::Shout => route_shout(
&policy.devices,
&devices,
daemon_up,
self.soma.as_deref().filter(|_| soma_up),
),
};
(policy, routed)
})
.collect();
Ok(RouteReport {
devices,
daemon_up,
soma: self.soma.clone(),
soma_up,
routes,
})
}
pub fn speak(
&self,
channel: Channel,
text: &str,
dry_run: bool,
pause_file: Option<&Path>,
out: &mut Out<'_>,
) -> Result<SpeechReceipt> {
super::operations::validate_text(text)?;
let prefs = self.voice.route(channel)?;
let routed = match channel {
Channel::Say => route_say(&prefs, &detect_output_devices()?),
Channel::Shout => {
let soma = match self.soma.as_deref() {
Some(soma) if soma_reachable(soma) => Some(soma),
Some(soma) => {
out.line(format!(
" [route] Soma at {soma} is not answering; falling through"
))?;
None
}
None => None,
};
if soma.is_some() {
route_shout(&prefs, &[], false, soma)
} else {
route_shout(
&prefs,
&detect_output_devices()?,
reachy_reachable(&self.daemon),
None,
)
}
}
};
out.line(format!("[{}] → {}", channel.name(), routed.describe()))?;
if dry_run {
return Ok(SpeechReceipt {
route: routed,
disposition: SpeechDisposition::DryRun,
utterance: None,
});
}
if matches!(routed, Routed::Text(_)) {
out.line(text)?;
let utterance = self.voice.record(channel, text, None, "voice spoke")?;
return Ok(SpeechReceipt {
route: routed,
disposition: SpeechDisposition::TextFallback,
utterance: Some(utterance),
});
}
if let Some(path) = pause_file {
out.line(format!(" [half-duplex] holding {}", path.display()))?;
}
let _pause = pause_file.map(crate::turntaking::PauseGuard::hold);
let path = unique_voice_tmp()?;
let outcome = speak_and_play(
&routed,
&self.daemon,
channel.name(),
text,
&path,
&self.synthesizer,
out,
);
match outcome {
Err(error) => {
let _ = std::fs::remove_file(&path);
if let Err(log_error) = self.voice.record(
channel,
text,
None,
"voice spoke (synthesis or WAV retention FAILED; text-only, no complete audio)",
) {
return Err(error.context(format!(
"recording failed utterance also failed: {log_error:#}"
)));
}
Err(error)
}
Ok(outcome) => {
let bytes = std::fs::read(&path)
.with_context(|| format!("read generated WAV {}", path.display()));
let _ = std::fs::remove_file(&path);
let logged = bytes
.and_then(|bytes| self.voice.record(channel, text, Some(bytes), "voice spoke"));
match outcome {
Spoken::PlaybackFailed(error) => {
if channel == Channel::Say {
if let Err(emission) = out.line(text) {
return Err(error.context(format!("private text fallback failed: {emission:#}; recording: {logged:?}")));
}
}
match logged {
Ok(id) => {
Err(error.context(format!("audio retained as utterance {id:x}")))
}
Err(log_error) => Err(error.context(format!(
"recording synthesized audio also failed: {log_error:#}"
))),
}
}
Spoken::ReportingFailed(error) => Err(match logged {
Ok(id) => error.context(format!("audio retained as utterance {id:x}")),
Err(log_error) => error.context(format!(
"recording synthesized audio also failed: {log_error:#}"
)),
}),
Spoken::Played(disposition) => Ok(SpeechReceipt {
route: routed,
disposition,
utterance: Some(retain_after_playback(disposition, logged)?),
}),
}
}
}
}
}