#![cfg_attr(not(feature = "duplex"), allow(dead_code))]
#[cfg(any(feature = "duplex", test))]
use super::operations::*;
use crate::{clock, out::Out};
use anyhow::{bail, Context, Result};
use std::collections::VecDeque;
#[cfg(any(feature = "duplex", test))]
use std::io::Write as _;
use std::path::{Path, PathBuf};
use std::sync::atomic::{AtomicBool, Ordering};
use std::sync::mpsc::{self, SyncSender};
use std::sync::{Arc, Condvar, Mutex};
use std::time::{Duration, Instant};
pub const FRAME_SAMPLES: usize = soma_client::FRAME_SAMPLES;
pub const SAMPLE_RATE: u32 = soma_client::SAMPLE_RATE;
const _: () = assert!(FRAME_SAMPLES * 1_000 == 80 * SAMPLE_RATE as usize);
pub const DEFAULT_SOMA: &str = soma_client::DEFAULT_BASE;
const MAX_BACKLOG_FRAMES: usize = 8;
const N_TEXT_SPECIALS: i64 = 4;
const TEXT_PAD: i64 = 3;
const TEXT_EPAD: i64 = 0;
pub const DEFAULT_UTTERANCE_GAP: usize = 10;
const VOICE_ONSET_FRAMES: usize = 3;
const VOICE_RELEASE_FRAMES: usize = 12;
const VOICE_FLOOR: f32 = 0.012;
pub const DEFAULT_SYSTEM: &str = "You are a warm, direct conversational partner. \
Keep replies short and spoken — no lists, no markdown. Answer in the language \
you were addressed in.";
#[derive(Clone, Debug)]
pub struct RunOptions {
pub weights: PathBuf,
pub voice_prompt: PathBuf,
pub spm: Option<PathBuf>,
pub soma: String,
pub output: String,
pub fmt: WeightFormat,
pub floor: Floor,
pub system: String,
pub temp: f32,
pub seed: u64,
pub pace: usize,
pub trace: Option<PathBuf>,
pub min_word_gap: usize,
pub nudge_after: usize,
pub cadence: Cadence,
pub gate: bool,
pub utterance_gap: usize,
pub pile: Option<PathBuf>,
pub key: Option<PathBuf>,
pub frames: Option<usize>,
pub wav: Option<PathBuf>,
pub decode_context: usize,
pub decode_hop: usize,
pub no_input: bool,
pub lead: Option<usize>,
pub pause_file: Option<PathBuf>,
}
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
pub enum WeightFormat {
Q4,
Q8,
F16,
}
impl WeightFormat {
pub fn parse(value: &str) -> Result<Self> {
Ok(match value {
"q4" => Self::Q4,
"q8" => Self::Q8,
"f16" => Self::F16,
_ => bail!("unknown weight format {value} (expected q4, q8 or f16)"),
})
}
}
impl Floor {
pub fn parse(value: &str) -> Result<Self> {
Ok(match value {
"listen" => Self::Listen,
"converse" => Self::Converse,
_ => bail!("unknown floor policy {value} (expected listen or converse)"),
})
}
}
impl Cadence {
pub fn parse(value: &str) -> Result<Self> {
Ok(match value {
"model" => Self::Model,
"word-onset" => Self::WordOnset,
"uniform" => Self::Uniform,
"dense" => Self::Dense,
_ => bail!("unknown cadence {value} (expected model, word-onset, uniform or dense)"),
})
}
pub fn name(self) -> &'static str {
match self {
Self::Model => "model",
Self::WordOnset => "word-onset",
Self::Uniform => "uniform",
Self::Dense => "dense",
}
}
}
impl RunOptions {
pub fn new(weights: PathBuf, voice_prompt: PathBuf, output: String) -> Self {
Self {
weights,
voice_prompt,
output,
spm: None,
soma: DEFAULT_SOMA.into(),
fmt: WeightFormat::Q8,
floor: Floor::Listen,
system: DEFAULT_SYSTEM.into(),
temp: 0.8,
seed: 12_345_678,
pace: 3,
trace: None,
min_word_gap: 0,
nudge_after: 25,
cadence: Cadence::Model,
gate: false,
utterance_gap: DEFAULT_UTTERANCE_GAP,
pile: None,
key: None,
frames: None,
wav: None,
decode_context: DEFAULT_DECODE_CONTEXT,
decode_hop: DEFAULT_DECODE_HOP,
no_input: false,
lead: None,
pause_file: None,
}
}
pub fn validate(&self) -> Result<()> {
if self.output.trim().is_empty() {
bail!("an exact named output device is required");
}
if !self.temp.is_finite() {
bail!("sampling temperature must be finite");
}
if self.decode_hop == 0 {
bail!("decode_hop must be positive");
}
if self.utterance_gap == 0 {
bail!("utterance_gap must be positive");
}
self.decode_context
.checked_add(self.decode_hop)
.context("decode window overflow")?;
if self.lead.is_none() {
self.decode_hop
.checked_mul(2)
.context("default speaker lead overflow")?;
}
Ok(())
}
}
#[derive(Clone, Debug)]
pub struct EarReport {
pub frames: usize,
pub wall_seconds: f64,
pub first_frame: Option<u64>,
pub last_frame: u64,
pub skipped: usize,
pub mean_rms: Option<f64>,
pub peak_rms: Option<f32>,
pub ended: Option<String>,
}
#[derive(Clone, Debug)]
pub struct RunReport {
pub frames: usize,
pub wall_seconds: f64,
pub step_times_ms: Vec<f64>,
pub over_budget: usize,
pub playback_underruns: u64,
pub speaker_stalls: usize,
pub output_dropped: u64,
pub input_skipped: usize,
pub forced_onset_ranks: Vec<usize>,
pub forced_continuation_ranks: Vec<usize>,
pub forced_padding_ranks: Vec<usize>,
pub word_gaps: Vec<usize>,
pub model_chose: usize,
pub nudged: usize,
pub held_back: usize,
}
#[derive(Clone, Debug)]
pub struct PlaybackDevice {
pub name: String,
pub configuration: std::result::Result<PlaybackConfiguration, String>,
}
#[derive(Clone, Debug)]
pub struct PlaybackConfiguration {
pub channels: u16,
pub sample_rate: u32,
pub sample_format: String,
}
#[cfg(feature = "audio")]
pub fn devices() -> Result<Vec<PlaybackDevice>> {
use rodio::cpal::traits::{DeviceTrait, HostTrait};
let mut devices = Vec::new();
for device in rodio::cpal::default_host()
.output_devices()
.context("enumerate playback devices")?
{
devices.push(PlaybackDevice {
name: device.name().unwrap_or_else(|_| "<unnamed>".into()),
configuration: device
.default_output_config()
.map(|config| PlaybackConfiguration {
channels: config.channels(),
sample_rate: config.sample_rate(),
sample_format: format!("{:?}", config.sample_format()),
})
.map_err(|error| error.to_string()),
});
}
Ok(devices)
}
#[cfg(not(feature = "audio"))]
pub fn devices() -> Result<Vec<PlaybackDevice>> {
bail!("audio device support is not compiled into this build (enable the `audio` feature)")
}
struct Ear {
ring: Arc<(Mutex<EarRing>, Condvar)>,
stop: Arc<AtomicBool>,
thread: Option<std::thread::JoinHandle<()>>,
soma: String,
}
#[derive(Default)]
struct EarRing {
frames: VecDeque<[f32; FRAME_SAMPLES]>,
skipped: usize,
first_frame: Option<u64>,
last_frame: u64,
ended: Option<String>,
}
impl Ear {
fn open(soma: &str) -> Result<Self> {
let mut capture = soma_client::SomaCapture::open(soma)
.with_context(|| format!("subscribe to Soma's microphone at {soma}"))?;
let ring: Arc<(Mutex<EarRing>, Condvar)> =
Arc::new((Mutex::new(EarRing::default()), Condvar::new()));
let stop = Arc::new(AtomicBool::new(false));
let thread = {
let ring = Arc::clone(&ring);
let stop = Arc::clone(&stop);
std::thread::Builder::new()
.name("duplex-ear".into())
.spawn(move || {
let (lock, ready) = &*ring;
while !stop.load(Ordering::Relaxed) {
match capture.next_frame() {
Ok(frame) => {
let mut ring = lock.lock().expect("ear ring");
ring.first_frame.get_or_insert(frame.frame_index);
ring.last_frame = frame.frame_index;
ring.frames.push_back(frame.samples);
drop(ring);
ready.notify_all();
}
Err(error) => {
let mut ring = lock.lock().expect("ear ring");
ring.ended = Some(format!("{error:#}"));
drop(ring);
ready.notify_all();
return;
}
}
}
let mut ring = lock.lock().expect("ear ring");
ring.ended.get_or_insert_with(|| "ear closed".into());
drop(ring);
ready.notify_all();
})
.context("spawn ear thread")?
};
Ok(Self {
ring,
stop,
thread: Some(thread),
soma: soma.to_string(),
})
}
fn next_frame(&self) -> Option<[f32; FRAME_SAMPLES]> {
let (lock, ready) = &*self.ring;
let mut ring = lock.lock().expect("ear ring");
loop {
let backlog = ring.frames.len();
if backlog > MAX_BACKLOG_FRAMES {
let drop_frames = backlog - 2;
ring.frames.drain(..drop_frames);
ring.skipped += drop_frames;
}
if let Some(frame) = ring.frames.pop_front() {
return Some(frame);
}
if ring.ended.is_some() {
return None;
}
ring = ready.wait(ring).expect("ear ring");
}
}
fn skipped(&self) -> usize {
self.ring.0.lock().expect("ear ring").skipped
}
fn clock(&self) -> (Option<u64>, u64) {
let ring = self.ring.0.lock().expect("ear ring");
(ring.first_frame, ring.last_frame)
}
fn ended(&self) -> Option<String> {
self.ring.0.lock().expect("ear ring").ended.clone()
}
}
impl Drop for Ear {
fn drop(&mut self) {
self.stop.store(true, Ordering::Relaxed);
if let Some(thread) = self.thread.take() {
let _ = thread.join();
}
}
}
fn rms(samples: &[f32]) -> f32 {
if samples.is_empty() {
return 0.0;
}
(samples.iter().map(|s| s * s).sum::<f32>() / samples.len() as f32).sqrt()
}
pub fn ear(soma: &str, frames: usize, stop: &AtomicBool, out: &mut Out<'_>) -> Result<EarReport> {
let ear = Ear::open(soma)?;
out.line(format!(
"duplex: ear open on {} — {} Hz, {} sample frames, this process owns no device",
ear.soma,
soma_client::SAMPLE_RATE,
soma_client::FRAME_SAMPLES
))?;
let start = Instant::now();
let mut read = 0usize;
let mut peak = 0f32;
let mut total = 0f64;
while (frames == 0 || read < frames) && !stop.load(Ordering::Relaxed) {
let Some(frame) = ear.next_frame() else {
break;
};
let level = rms(&frame);
peak = peak.max(level);
total += level as f64;
read += 1;
if read % 12 == 0 {
let (first, last) = ear.clock();
out.line(format!(
" [ear] body frame {last} (joined at {}) | {read} read | \
{} skipped | rms {level:.4}",
first.unwrap_or(0),
ear.skipped()
))?;
}
}
let wall = start.elapsed().as_secs_f64();
let (first, last) = ear.clock();
out.line(format!(
"duplex: {read} frames in {wall:.1}s — {:.2}x realtime | body frames {}..{last} | \
{} skipped | ear rms mean {:.4} peak {:.4}",
(read as f64 * 0.08) / wall.max(1e-9),
first.unwrap_or(0),
ear.skipped(),
total / read.max(1) as f64,
peak
))?;
if let Some(reason) = ear.ended() {
out.line(format!("duplex: the body's stream ended — {reason}"))?;
}
Ok(EarReport {
frames: read,
wall_seconds: wall,
first_frame: first,
last_frame: last,
skipped: ear.skipped(),
mean_rms: (read > 0).then_some(total / read as f64),
peak_rms: (read > 0).then_some(peak),
ended: ear.ended(),
})
}
pub const DEFAULT_DECODE_CONTEXT: usize = 64;
pub const DEFAULT_DECODE_HOP: usize = 4;
const DECODE_QUEUE_FRAMES: usize = 16;
const PREBUFFER_FRAMES: usize = 5;
const PACE_DEADLINE: Duration = Duration::from_millis(500);
#[derive(Default)]
struct AudibleWindow {
tail: usize,
}
impl AudibleWindow {
fn observe(&mut self, speaking: bool, in_flight: usize) -> bool {
if speaking {
self.tail = in_flight + PREBUFFER_FRAMES;
return true;
}
let sounding = self.tail > 0;
self.tail = self.tail.saturating_sub(1);
sounding
}
}
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
pub enum Cadence {
WordOnset,
Uniform,
Dense,
Model,
}
#[cfg(feature = "duplex")]
fn cadence_schedule(
spm: &mary::models::personaplex::spm::SpmTokenizer,
line: &str,
gap: usize,
cadence: Cadence,
) -> Vec<i64> {
const WORD_MARK: &[u8] = "\u{2581}".as_bytes();
let ids = spm.encode(line);
if cadence == Cadence::Model {
return ids;
}
let mut out = Vec::with_capacity(ids.len() * (1 + gap));
for (k, &t) in ids.iter().enumerate() {
out.push(t);
let pad_here = match cadence {
Cadence::Model | Cadence::Dense => false,
Cadence::Uniform => true,
Cadence::WordOnset => ids
.get(k + 1)
.map(|&n| spm.piece_bytes(n).starts_with(WORD_MARK))
.unwrap_or(false),
};
if pad_here {
for g in 0..gap {
out.push(if g + 1 == gap { TEXT_EPAD } else { TEXT_PAD });
}
}
}
out
}
#[cfg(feature = "duplex")]
type Codes = [u32; mary::models::personaplex::mimi::config::NUM_CODEBOOKS];
#[cfg(feature = "duplex")]
use std::sync::atomic::AtomicU64;
#[cfg(feature = "duplex")]
use std::sync::mpsc::{Receiver, TrySendError};
#[cfg(feature = "duplex")]
struct Mouth {
frames: Option<SyncSender<Codes>>,
worker: Option<std::thread::JoinHandle<Result<()>>>,
dropped: Arc<AtomicU64>,
underruns: Arc<AtomicU64>,
played: Arc<AtomicU64>,
pushed: std::cell::Cell<u64>,
playing: Arc<AtomicBool>,
cancel: Arc<AtomicBool>,
}
#[cfg(feature = "duplex")]
impl Mouth {
fn spawn(
weights: PathBuf,
device: String,
wav: Option<PathBuf>,
context_frames: usize,
hop_frames: usize,
) -> Result<Self> {
let (frame_tx, frame_rx) = mpsc::sync_channel(DECODE_QUEUE_FRAMES);
let (ready_tx, ready_rx) = mpsc::channel::<std::result::Result<(), String>>();
let dropped = Arc::new(AtomicU64::new(0));
let underruns = Arc::new(AtomicU64::new(0));
let played = Arc::new(AtomicU64::new(0));
let playing = Arc::new(AtomicBool::new(false));
let cancel = Arc::new(AtomicBool::new(false));
let worker = {
let underruns = Arc::clone(&underruns);
let played = Arc::clone(&played);
let playing = Arc::clone(&playing);
let cancel = Arc::clone(&cancel);
std::thread::Builder::new()
.name("duplex-mouth".into())
.spawn(move || {
mouth_worker(
weights,
device,
wav,
frame_rx,
ready_tx,
underruns,
played,
playing,
cancel,
context_frames,
hop_frames,
)
})
.context("spawn playback thread")?
};
let mouth = Self {
frames: Some(frame_tx),
worker: Some(worker),
dropped,
underruns,
played,
pushed: std::cell::Cell::new(0),
playing,
cancel,
};
match ready_rx.recv_timeout(Duration::from_secs(900)) {
Ok(Ok(())) => {}
Ok(Err(message)) => bail!("playback: {message}"),
Err(error) => bail!("playback did not come up: {error}"),
}
Ok(mouth)
}
fn push(&self, codes: Codes) {
let Some(sender) = self.frames.as_ref() else {
return;
};
match sender.try_send(codes) {
Ok(()) => self.pushed.set(self.pushed.get() + 1),
Err(TrySendError::Full(_)) | Err(TrySendError::Disconnected(_)) => {
self.dropped.fetch_add(1, Ordering::Relaxed);
}
}
}
fn dropped(&self) -> u64 {
self.dropped.load(Ordering::Relaxed)
}
fn underruns(&self) -> u64 {
self.underruns.load(Ordering::Relaxed)
}
fn lead(&self) -> u64 {
self.pushed
.get()
.saturating_sub(self.played.load(Ordering::Relaxed))
}
fn pace(&self, target: u64) -> bool {
if !self.playing.load(Ordering::Relaxed) {
return true;
}
let deadline = Instant::now() + PACE_DEADLINE;
while self.lead() >= target {
if Instant::now() >= deadline {
return false;
}
std::thread::sleep(Duration::from_millis(2));
}
true
}
fn finish(mut self) -> Result<()> {
self.frames.take();
match self.worker.take() {
Some(worker) => match worker.join() {
Ok(result) => result,
Err(_) => bail!("playback thread panicked"),
},
None => Ok(()),
}
}
}
#[cfg(feature = "duplex")]
impl Drop for Mouth {
fn drop(&mut self) {
self.cancel.store(true, Ordering::Relaxed);
self.frames.take();
if let Some(worker) = self.worker.take() {
let _ = worker.join();
}
}
}
#[cfg(feature = "duplex")]
#[allow(clippy::too_many_arguments)]
fn mouth_worker(
weights: PathBuf,
device_name: String,
wav: Option<PathBuf>,
frames: Receiver<Codes>,
ready: mpsc::Sender<std::result::Result<(), String>>,
underruns: Arc<AtomicU64>,
played: Arc<AtomicU64>,
playing: Arc<AtomicBool>,
cancel: Arc<AtomicBool>,
context_frames: usize,
hop_frames: usize,
) -> Result<()> {
use mary::models::personaplex::mimi::config as codec_cfg;
use mary::models::personaplex::mimi::MimiDecoder;
use rodio::buffer::SamplesBuffer;
use std::num::NonZero;
set_interactive_qos();
let setup = || -> Result<(MimiDecoder, rodio::MixerDeviceSink, rodio::Player)> {
let loader = mary::persist::personaplex_loader(&weights)
.with_context(|| format!("load the codec from {}", weights.display()))?;
let decoder = MimiDecoder::load(&loader);
let _ = decoder.decode(&vec![
[0; codec_cfg::NUM_CODEBOOKS];
context_frames + hop_frames
]);
let (sink, player) = open_named_sink(&device_name)?;
Ok((decoder, sink, player))
};
let (decoder, _sink, player) = match setup() {
Ok(value) => {
let _ = ready.send(Ok(()));
value
}
Err(error) => {
let _ = ready.send(Err(format!("{error:#}")));
return Err(error);
}
};
let mono = NonZero::new(1u16).expect("1 is nonzero");
let rate = NonZero::new(SAMPLE_RATE).expect("24000 is nonzero");
let mut receipt = wav.map(WavWriter::create).transpose()?;
let mut history: VecDeque<Codes> = VecDeque::new();
let mut pending = 0usize;
let mut emitted = 0usize;
let mut started = false;
player.pause();
let mut emit = |history: &VecDeque<Codes>, pending: usize| {
let chunk: Vec<Codes> = history.iter().copied().collect();
let context = chunk.len() - pending;
let pcm = decoder.decode(&chunk);
let from = context * FRAME_SAMPLES;
let to = (from + pending * FRAME_SAMPLES).min(pcm.len());
let slice = pcm[from..to].to_vec();
if let Some(receipt) = receipt.as_mut() {
let _ = receipt.write(&slice);
}
player.append(SamplesBuffer::new(mono, rate, slice));
};
let publish = |player: &rodio::Player, emitted: usize| {
played.store(
(emitted.saturating_sub(player.len() * hop_frames)) as u64,
Ordering::Relaxed,
);
};
loop {
if cancel.load(Ordering::Relaxed) {
player.stop();
return Ok(());
}
let frame = match frames.recv_timeout(Duration::from_millis(2)) {
Ok(frame) => frame,
Err(mpsc::RecvTimeoutError::Timeout) => {
publish(&player, emitted);
continue;
}
Err(mpsc::RecvTimeoutError::Disconnected) => break,
};
history.push_back(frame);
pending += 1;
if pending < hop_frames {
publish(&player, emitted);
continue;
}
if started && player.empty() {
underruns.fetch_add(1, Ordering::Relaxed);
}
emit(&history, pending);
emitted += pending;
pending = 0;
while history.len() > context_frames {
history.pop_front();
}
if !started && emitted >= PREBUFFER_FRAMES {
player.play();
started = true;
playing.store(true, Ordering::Relaxed);
}
publish(&player, emitted);
}
if pending > 0 {
emit(&history, pending);
}
if !started {
player.play();
}
let deadline = Instant::now() + Duration::from_secs(10);
while !player.empty() && Instant::now() < deadline && !cancel.load(Ordering::Relaxed) {
std::thread::sleep(Duration::from_millis(20));
}
if cancel.load(Ordering::Relaxed) {
player.stop();
}
drop(receipt);
Ok(())
}
#[cfg(feature = "audio")]
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 playback devices")?
.find(|candidate| candidate.name().map(|n| n == name).unwrap_or(false))
.with_context(|| format!("playback device '{name}' not found (disconnected?)"))?;
let mut sink = rodio::DeviceSinkBuilder::from_device(device)
.with_context(|| format!("prepare playback device '{name}'"))?
.open_stream()
.with_context(|| format!("open playback device '{name}'"))?;
let player = rodio::Player::connect_new(sink.mixer());
sink.log_on_drop(false);
Ok((sink, player))
}
struct Ledger {
lines: Option<SyncSender<String>>,
worker: Option<std::thread::JoinHandle<()>>,
warnings: mpsc::Receiver<String>,
dropped_warnings: Arc<std::sync::atomic::AtomicUsize>,
dropped_lines: std::cell::Cell<usize>,
}
impl Ledger {
fn spawn(pile: Option<PathBuf>, key: Option<PathBuf>) -> Result<Self> {
let (warning_tx, warnings) = mpsc::sync_channel(16);
let dropped_warnings = Arc::new(std::sync::atomic::AtomicUsize::new(0));
let Some(pile) = pile else {
return Ok(Self {
lines: None,
worker: None,
warnings,
dropped_warnings,
dropped_lines: std::cell::Cell::new(0),
});
};
let (tx, rx) = mpsc::sync_channel::<String>(64);
let dropped = Arc::clone(&dropped_warnings);
let worker = std::thread::Builder::new()
.name("duplex-ledger".into())
.spawn(move || {
while let Ok(text) = rx.recv() {
if let Err(error) = record_utterance(&pile, key.as_deref(), &text) {
if warning_tx
.try_send(format!("duplex: could not record utterance: {error:#}"))
.is_err()
{
dropped.fetch_add(1, Ordering::Relaxed);
}
}
}
})
.context("spawn transcript ledger")?;
Ok(Self {
lines: Some(tx),
worker: Some(worker),
warnings,
dropped_warnings,
dropped_lines: std::cell::Cell::new(0),
})
}
fn record(&self, text: &str) {
if let Some(lines) = &self.lines {
if lines.try_send(text.to_owned()).is_err() {
self.dropped_lines.set(self.dropped_lines.get() + 1);
}
}
}
fn report(&self, out: &mut Out<'_>) -> Result<()> {
for warning in self.warnings.try_iter() {
out.line(warning)?;
}
let dropped = self.dropped_warnings.swap(0, Ordering::Relaxed);
if dropped > 0 {
out.line(format!(
"duplex: {dropped} additional ledger errors exceeded the diagnostic queue"
))?;
}
let dropped = self.dropped_lines.replace(0);
if dropped > 0 {
out.line(format!(
"duplex: {dropped} utterances could not enter the durable transcript queue"
))?;
}
Ok(())
}
fn finish(mut self, out: &mut Out<'_>) -> Result<()> {
self.lines.take();
if let Some(worker) = self.worker.take() {
worker
.join()
.map_err(|_| anyhow::anyhow!("transcript ledger thread panicked"))?;
}
self.report(out)
}
}
impl Drop for Ledger {
fn drop(&mut self) {
self.lines.take();
if let Some(worker) = self.worker.take() {
let _ = worker.join();
}
}
}
#[cfg(any(feature = "duplex", test))]
fn spoken_line(session: &Path, line: &Line, ledger: &Ledger, out: &mut Out<'_>) -> Result<()> {
let transcript = append_line(session, line);
ledger.record(&line.text);
transcript.context("record generated utterance in the session transcript")?;
out.line(format!(" [{}] {}", line.speaker, line.text))
}
#[cfg(any(feature = "duplex", test))]
fn accept_inject_drain(
drained: InjectDrain,
last_cleanup_failure: &mut Option<String>,
mut schedule: impl FnMut(&str),
out: &mut Out<'_>,
) -> Result<()> {
for line in &drained.lines {
schedule(line);
}
for line in drained.lines {
out.line(format!("duplex: to say — {line}"))?;
}
let failure = drained.cleanup_failure.map(|error| format!("{error:#}"));
if failure != *last_cleanup_failure {
if let Some(error) = &failure {
out.line(format!("duplex: inject queue cleanup failed — {error}"))?;
}
*last_cleanup_failure = failure;
}
Ok(())
}
fn record_utterance(pile_path: &Path, key: Option<&Path>, text: &str) -> Result<()> {
use crate::schemas::voice::{CHANNEL_SHOUT, COLLECTION_SCOPE_ID};
use triblespace::core::collection::CollectionStoreExt;
use triblespace::core::metadata;
use triblespace::prelude::*;
let stamp = clock::point_now()?;
let mut fragment = crate::voice::utterance_fragment(CHANNEL_SHOUT, text, None, stamp)?;
crate::storage::Storage::new(pile_path.to_owned(), key.map(Path::to_owned))
.with_store(|pile, signer, runtime| {
let collection = crate::collection_names::open_configured_acquiring(
pile, COLLECTION_SCOPE_ID, signer.verifying_key(), runtime,
)?;
crate::voice::validate_staged_payloads(&mut fragment)?;
fragment.describe_with(entity! { metadata::description: "duplex spoke" });
crate::collection_names::require_command_write_admission_acquiring(
pile,
collection,
signer,
"Duplex",
"voice route show",
runtime,
)?;
pile.commit(collection, signer, fragment)
.context("commit the utterance")?;
drop(
runtime.block_on(crate::storage::ensure_downstream(
pile, collection, signer,
))
.context("Duplex utterance was committed, but ensuring its derived views failed")?,
);
Ok(())
})
}
#[cfg(not(feature = "duplex"))]
pub fn run(
_session: &Path,
_args: &RunOptions,
_stop: &AtomicBool,
_out: &mut Out<'_>,
) -> Result<RunReport> {
bail!("the channel is not compiled into this build (build with --features duplex)")
}
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
pub enum Floor {
Listen,
Converse,
}
#[cfg(feature = "duplex")]
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
enum PersonaPlexStepApi {
Duplex,
DuplexArbitrated,
OutputOnly,
OutputOnlyArbitrated,
}
#[cfg(feature = "duplex")]
impl PersonaPlexStepApi {
fn select(output_only: bool, arbitrated: bool) -> Self {
match (output_only, arbitrated) {
(false, false) => Self::Duplex,
(false, true) => Self::DuplexArbitrated,
(true, false) => Self::OutputOnly,
(true, true) => Self::OutputOnlyArbitrated,
}
}
}
#[cfg(feature = "duplex")]
pub fn run(
session: &Path,
args: &RunOptions,
stop: &AtomicBool,
out: &mut Out<'_>,
) -> Result<RunReport> {
use mary::models::personaplex::config as model_cfg;
use mary::models::personaplex::pipeline::{agent_codes, RealtimePipeline, SILENCE};
use mary::models::personaplex::prompt::Prompt;
use mary::models::personaplex::sampling::SamplingConfig;
use mary::models::personaplex::temporal_metal::WeightFmt;
args.validate()?;
let fmt = match args.fmt {
WeightFormat::Q4 => WeightFmt::Q4,
WeightFormat::Q8 => WeightFmt::Q8,
WeightFormat::F16 => WeightFmt::F16,
};
let cadence = args.cadence;
let floor_policy = args.floor;
ensure_session(session)?;
let _ = std::fs::remove_file(session.join(TRANSCRIPT_FILE));
let _ = std::fs::remove_file(session.join(HOLD_FILE));
write_cursor(session, 0)?;
set_interactive_qos();
out.line(format!("duplex: bringing up '{}' …", args.output))?;
let mouth = Mouth::spawn(
args.weights.clone(),
args.output.clone(),
args.wav.clone(),
args.decode_context,
args.decode_hop,
)?;
out.line(format!(
"duplex: loading the model from {} …",
args.weights.display()
))?;
let load_start = Instant::now();
let source = mary::persist::personaplex_bundle(&args.weights)
.with_context(|| format!("load the model from {}", args.weights.display()))?
.into_runtime_source();
let mut pipeline = RealtimePipeline::load_auto(&source, fmt, true);
if args.temp <= 0.0 {
pipeline.set_greedy();
} else {
pipeline.set_sampling(
SamplingConfig {
temp: args.temp,
top_k: 250,
top_p: 0.95,
},
args.seed,
);
}
let spm = match args.spm.as_deref() {
Some(path) => mary::models::personaplex::spm::SpmTokenizer::load(path),
None => mary::persist::load_spm_tokenizer_from_pile(&args.weights)
.context("load the text tokenizer from the weight pile (or pass --spm)")?,
};
if spm.vocab_size() != model_cfg::TEXT_CARD {
bail!(
"tokenizer vocabulary {} does not match the model's {} — wrong tokenizer",
spm.vocab_size(),
model_cfg::TEXT_CARD
);
}
let prompt = Prompt::build(&args.voice_prompt, &spm, &args.system);
pipeline.run_prompt(&prompt);
out.line(format!(
"duplex: model ready in {:.1}s ({} prompt steps)",
load_start.elapsed().as_secs_f64(),
prompt.total_steps()
))?;
let nudge_after = args.nudge_after;
let min_word_gap = args.min_word_gap;
let mut trace_out = match args.trace.as_ref() {
Some(path) => {
let mut f = std::fs::File::create(path)
.with_context(|| format!("open the frame trace {}", path.display()))?;
use std::io::Write;
writeln!(f, "frame\ttoken\tclass\tsource")?;
Some(f)
}
None => None,
};
let lead_target = args.lead.unwrap_or(2 * args.decode_hop).max(1) as u64;
let mut stalled = 0usize;
let ear = if args.no_input {
out.line(format!(
"duplex: generation only — no ear, the speaker is the clock \
(lead {lead_target} frames = {} ms)",
lead_target * 80
))?;
None
} else {
out.line(format!(
"duplex: subscribing to the body's microphone at {} …",
args.soma
))?;
let ear = Ear::open(&args.soma)?;
out.line(format!(
"duplex: ear open — {} Hz, {} sample frames; this process owns no device, so \
`hear` can read the same frames",
SAMPLE_RATE, FRAME_SAMPLES
))?;
Some(ear)
};
let ledger = Ledger::spawn(args.pile.clone(), args.key.clone())?;
let mut encoder_state = pipeline.encoder.stream_state();
let mut queued: VecDeque<i64> = VecDeque::new();
let mut injected_text: VecDeque<String> = VecDeque::new();
let mut spoken = String::new();
let mut spoken_was_injected = false;
let mut pad_run = 0usize;
let mut speaking_hangover = 0usize;
if let Some(path) = args.pause_file.as_deref() {
if crate::turntaking::clear_stale(path) {
out.line(format!(
"duplex: cleared a stale pause file at {}",
path.display()
))?;
}
out.line(format!(
"duplex: holding {} while audible, so a `hear` on the same body stays deaf to us",
path.display()
))?;
}
let mut audible: Option<crate::turntaking::PauseGuard> = None;
let mut audible_window = AudibleWindow::default();
let mut seq = 0u64;
let mut held = false;
let mut far_loud = 0usize;
let mut far_quiet = 0usize;
let mut far_talking = false;
let mut far_started_ms = 0u64;
let mut heard_peak = 0f32;
let mut heard_total = 0f64;
let mut heard_frames = 0usize;
let mut rank_onset: Vec<usize> = Vec::new();
let mut rank_cont: Vec<usize> = Vec::new();
let mut rank_pad: Vec<usize> = Vec::new();
let mut word_gaps: Vec<usize> = Vec::new();
let mut since_word = 0usize;
let mut model_chose = 0usize;
let mut nudged = 0usize;
let mut wait_frames = 0usize;
let mut held_back = 0usize;
let mut since_onset = usize::MAX / 2;
let mut frame_index = 0usize;
let mut step_total = 0f64;
let mut step_max = 0f64;
let mut over_budget = 0usize;
let mut step_times: Vec<f64> = Vec::new();
let mut last_control_poll = Instant::now();
let mut last_inject_cleanup_failure = None;
out.line(format!(
"duplex: live. `duplex read` to catch up, `duplex say` to speak, Ctrl-C to stop."
))?;
let session_start = Instant::now();
loop {
if stop.load(Ordering::Relaxed) {
break;
}
if let Some(limit) = args.frames {
if frame_index >= limit {
break;
}
}
if last_control_poll.elapsed() >= Duration::from_millis(250) {
last_control_poll = Instant::now();
ledger.report(out)?;
let now_held = floor_held(session)?;
if now_held != held {
out.line(format!(
"duplex: floor {}",
if now_held {
"TAKEN — silent"
} else {
"released"
}
))?;
held = now_held;
}
accept_inject_drain(
drain_inject(session)?,
&mut last_inject_cleanup_failure,
|line| {
queued.extend(cadence_schedule(&spm, line, args.pace, cadence));
injected_text.push_back(line.to_owned());
},
out,
)?;
}
let samples = match ear.as_ref() {
Some(ear) => match ear.next_frame() {
Some(frame) => Some(frame),
None => {
match ear.ended() {
Some(reason) => {
out.line(format!("duplex: the body's stream ended — {reason}"))?
}
None => out.line("duplex: the body stopped producing frames")?,
}
break;
}
},
None => {
if !mouth.pace(lead_target) {
stalled += 1;
}
None
}
};
if let Some(frame) = samples.as_ref() {
let level = rms(frame);
heard_peak = heard_peak.max(level);
heard_total += level as f64;
heard_frames += 1;
if level >= VOICE_FLOOR {
far_loud += 1;
far_quiet = 0;
} else {
far_quiet += 1;
far_loud = 0;
}
let _ = level;
if !far_talking && far_loud >= VOICE_ONSET_FRAMES {
far_talking = true;
far_started_ms = now_millis()?;
} else if far_talking && far_quiet >= VOICE_RELEASE_FRAMES {
far_talking = false;
let seconds = (now_millis()?.saturating_sub(far_started_ms)) as f64 / 1000.0;
seq += 1;
let line = Line {
seq,
at_ms: far_started_ms,
speaker: Speaker::Far.tag().into(),
text: format!("[spoke for {seconds:.1}s — no transcription available]"),
};
let _ = append_line(session, &line);
}
}
let step_start = Instant::now();
let heard: Option<[i64; 8]> = match samples {
Some(_) if args.gate && speaking_hangover > 0 => Some(SILENCE),
Some(frame) => {
let codes = pipeline
.encoder
.encode_stream_frame(&mut encoder_state, &frame);
Some(std::array::from_fn(|q| codes[q] as i64))
}
None => None,
};
let was_speaking_injected = !queued.is_empty();
let free_timing = cadence == Cadence::Model && was_speaking_injected && !held;
let (forced_text, forced_audio) = if held {
(Some(TEXT_PAD), Some(&SILENCE))
} else if free_timing {
(None, None)
} else if !queued.is_empty() {
(queued.pop_front(), None)
} else {
match floor_policy {
Floor::Listen => (Some(TEXT_PAD), None),
Floor::Converse => (None, None),
}
};
let step_api = PersonaPlexStepApi::select(args.no_input, free_timing);
let trace = if free_timing {
let queue = &mut queued;
let waited = &mut wait_frames;
let chose = &mut model_chose;
let nudges = &mut nudged;
let held = &mut held_back;
let gap = &mut since_onset;
let spm_ref = &spm;
let onset = move |t: i64| {
t >= N_TEXT_SPECIALS && spm_ref.piece_bytes(t).starts_with("\u{2581}".as_bytes())
};
let mut decide = |_logits: &[f32], sampled: i64| -> i64 {
if sampled >= N_TEXT_SPECIALS {
let next_onset = queue.front().copied().map(onset).unwrap_or(false);
if min_word_gap > 0 && next_onset && *gap < min_word_gap {
*gap += 1;
*held += 1;
return TEXT_PAD;
}
*waited = 0;
*chose += 1;
let t = queue.pop_front().unwrap_or(sampled);
if onset(t) {
*gap = 0;
} else {
*gap += 1;
}
t
} else if *waited >= nudge_after {
*waited = 0;
*nudges += 1;
let t = queue.pop_front().unwrap_or(sampled);
if onset(t) {
*gap = 0;
} else {
*gap += 1;
}
t
} else {
*waited += 1;
*gap += 1;
sampled
}
};
match (step_api, heard.as_ref()) {
(PersonaPlexStepApi::DuplexArbitrated, Some(heard)) => {
pipeline.step_arbitrated(Some(heard), forced_audio, &mut decide)
}
(PersonaPlexStepApi::OutputOnlyArbitrated, None) => {
pipeline.step_output_only_arbitrated(forced_audio, &mut decide)
}
_ => unreachable!("PersonaPlex input protocol changed within one frame"),
}
} else {
match (step_api, heard.as_ref()) {
(PersonaPlexStepApi::Duplex, Some(heard)) => {
pipeline.step(Some(heard), forced_audio, forced_text)
}
(PersonaPlexStepApi::OutputOnly, None) => {
pipeline.step_output_only(forced_audio, forced_text)
}
_ => unreachable!("PersonaPlex input protocol changed within one frame"),
}
};
let chosen: Option<i64> = if free_timing {
Some(trace.next_text)
} else {
forced_text
};
let elapsed = step_start.elapsed().as_secs_f64() * 1e3;
step_total += elapsed;
step_max = step_max.max(elapsed);
step_times.push(elapsed);
if elapsed > 80.0 {
over_budget += 1;
}
if was_speaking_injected {
if let Some(t) = chosen {
if !trace.text_logits.is_empty() && (t as usize) < trace.text_logits.len() {
let lv = trace.text_logits[t as usize];
let rank = trace.text_logits.iter().filter(|&&v| v > lv).count();
if t < N_TEXT_SPECIALS {
rank_pad.push(rank);
} else if spm.piece_bytes(t).starts_with("\u{2581}".as_bytes()) {
rank_onset.push(rank);
} else {
rank_cont.push(rank);
}
}
}
}
if let Some(f) = trace_out.as_mut() {
use std::io::Write;
let t = chosen.unwrap_or(-1);
let class = if t < 0 {
"none"
} else if t == TEXT_EPAD {
"epad"
} else if t < N_TEXT_SPECIALS {
"pad"
} else if spm.piece_bytes(t).starts_with("\u{2581}".as_bytes()) {
"onset"
} else {
"cont"
};
let source = if free_timing { "model" } else { "schedule" };
let _ = writeln!(f, "{frame_index}\t{t}\t{class}\t{source}");
}
if was_speaking_injected {
match chosen {
Some(t) if t >= N_TEXT_SPECIALS => {
word_gaps.push(since_word);
since_word = 0;
}
_ => since_word += 1,
}
}
if let Some(out) = trace.out.as_ref() {
mouth.push(agent_codes(out));
let token = out[0];
if token >= N_TEXT_SPECIALS {
let piece = spm.decode_token(token);
if !piece.is_empty() {
spoken.push_str(&piece);
spoken_was_injected |= was_speaking_injected;
pad_run = 0;
speaking_hangover = args.utterance_gap;
}
} else {
pad_run += 1;
speaking_hangover = speaking_hangover.saturating_sub(1);
}
}
if let Some(path) = args.pause_file.as_deref() {
let sounding = audible_window.observe(speaking_hangover > 0, mouth.lead() as usize);
if sounding && audible.is_none() {
audible = Some(crate::turntaking::PauseGuard::hold(path));
} else if !sounding {
audible = None;
}
}
if queued.is_empty() && pad_run >= args.utterance_gap && !spoken.trim().is_empty() {
let line = spoken.trim().to_owned();
let speaker = if spoken_was_injected {
Speaker::Agent
} else {
Speaker::Model
};
seq += 1;
spoken_line(
session,
&Line {
seq,
at_ms: now_millis()?,
speaker: speaker.tag().into(),
text: line.clone(),
},
&ledger,
out,
)?;
let _ = injected_text.pop_front();
spoken.clear();
spoken_was_injected = false;
pad_run = 0;
}
frame_index += 1;
if frame_index % 125 == 0 {
let skipped = ear.as_ref().map(Ear::skipped).unwrap_or(0);
let body = ear
.as_ref()
.map(|ear| ear.clock().1.to_string())
.unwrap_or_else(|| "-".into());
out.line(format!(
" [clock] {:.1}s | body frame {body} | step mean {:.1} ms max {:.1} ms | \
{} over budget | {} in skipped | {} out dropped | {} underruns | lead {} | \
ear rms mean {:.4} peak {:.4}",
session_start.elapsed().as_secs_f64(),
step_total / frame_index as f64,
step_max,
over_budget,
skipped,
mouth.dropped(),
mouth.underruns(),
mouth.lead(),
heard_total / heard_frames.max(1) as f64,
heard_peak
))?;
}
}
if far_talking {
let seconds = (now_millis()?.saturating_sub(far_started_ms)) as f64 / 1000.0;
seq += 1;
let _ = append_line(
session,
&Line {
seq,
at_ms: far_started_ms,
speaker: Speaker::Far.tag().into(),
text: format!(
"[speaking for {seconds:.1}s when the session ended — no \
transcription available]"
),
},
);
}
if !spoken.trim().is_empty() {
let line = spoken.trim().to_owned();
seq += 1;
let speaker = if spoken_was_injected {
Speaker::Agent
} else {
Speaker::Model
};
spoken_line(
session,
&Line {
seq,
at_ms: now_millis()?,
speaker: speaker.tag().into(),
text: line.clone(),
},
&ledger,
out,
)?;
}
let wall = session_start.elapsed().as_secs_f64();
step_times.sort_by(|a, b| a.partial_cmp(b).expect("finite"));
let pct = |p: f64| -> f64 {
if step_times.is_empty() {
return 0.0;
}
let index = ((step_times.len() - 1) as f64 * p).round() as usize;
step_times[index]
};
out.line(format!(
"duplex: {} frames of audio in {:.1}s wall — {:.2}x realtime\n\
\x20 step ms: p50 {:.1} p90 {:.1} p99 {:.1} max {:.1} mean {:.1} \
(budget 80)\n\
\x20 {} of {} frames over budget ({:.0}%), {} playback underruns, \
{} speaker stalls",
frame_index,
wall,
(frame_index as f64 * 0.08) / wall.max(1e-9),
pct(0.50),
pct(0.90),
pct(0.99),
step_max,
step_total / frame_index.max(1) as f64,
over_budget,
frame_index,
100.0 * over_budget as f64 / frame_index.max(1) as f64,
mouth.underruns(),
stalled
))?;
let quantiles = |v: &mut Vec<usize>| -> (usize, usize, usize) {
v.sort_unstable();
let at = |p: f64| v[((v.len() - 1) as f64 * p).round() as usize];
(at(0.50), at(0.90), at(0.99))
};
if rank_onset.is_empty() && rank_cont.is_empty() {
out.line(format!(" forced-token rank: nothing was forced"))?;
} else {
let total = rank_onset.len() + rank_cont.len() + rank_pad.len();
out.line(format!(
" forced-token rank in the model's own logits ({} cadence, gap {}, {} forced frames, {:.0}% PAD):",
args.cadence.name(),
args.pace,
total,
100.0 * rank_pad.len() as f64 / total.max(1) as f64,
))?;
for (what, v) in [
("word continuation", &mut rank_cont),
("word onset ", &mut rank_onset),
("pad ", &mut rank_pad),
] {
if v.is_empty() {
continue;
}
let n = v.len();
let (a, b, c) = quantiles(v);
out.line(format!(
" {what} p50 {a:>6} p90 {b:>6} p99 {c:>6} (n={n})"
))?;
}
}
if !word_gaps.is_empty() {
let n = word_gaps.len();
let spread = {
let mut u: Vec<usize> = word_gaps.clone();
u.sort_unstable();
u.dedup();
u.len()
};
let mean = word_gaps.iter().sum::<usize>() as f64 / n as f64;
let (g50, g90, g99) = quantiles(&mut word_gaps);
out.line(format!(
" rhythm — frames from one word to the next: p50 {g50} p90 {g90} p99 {g99} mean {mean:.2} ({spread} distinct values over {n} words)"
))?;
}
if model_chose + nudged > 0 {
out.line(format!(
" timing: the model chose the moment {model_chose} times, we started it {nudged} times ({:.0}% ours)",
100.0 * nudged as f64 / (model_chose + nudged) as f64
))?;
if held_back > 0 {
out.line(format!(
" and held it back {held_back} frames for --min-word-gap {min_word_gap}"
))?;
}
}
let report = RunReport {
frames: frame_index,
wall_seconds: wall,
step_times_ms: step_times,
over_budget,
playback_underruns: mouth.underruns(),
speaker_stalls: stalled,
output_dropped: mouth.dropped(),
input_skipped: ear.as_ref().map(Ear::skipped).unwrap_or(0),
forced_onset_ranks: rank_onset,
forced_continuation_ranks: rank_cont,
forced_padding_ranks: rank_pad,
word_gaps,
model_chose,
nudged,
held_back,
};
drop(audible);
drop(ear);
mouth.finish()?;
ledger.finish(out)?;
Ok(report)
}
fn set_interactive_qos() {
#[cfg(target_os = "macos")]
unsafe {
extern "C" {
fn pthread_set_qos_class_self_np(qos_class: u32, relative_priority: i32) -> i32;
}
let _ = pthread_set_qos_class_self_np(0x21, 0);
}
}
#[cfg(any(feature = "duplex", test))]
struct WavWriter {
file: std::fs::File,
data_bytes: u32,
}
#[cfg(any(feature = "duplex", test))]
impl WavWriter {
fn create(path: PathBuf) -> Result<Self> {
let mut file =
std::fs::File::create(&path).with_context(|| format!("create {}", path.display()))?;
file.write_all(&wav_header(0))?;
Ok(Self {
file,
data_bytes: 0,
})
}
fn write(&mut self, samples: &[f32]) -> Result<()> {
let mut bytes = Vec::with_capacity(samples.len() * 2);
for sample in samples {
let value = (sample.clamp(-1.0, 1.0) * i16::MAX as f32).round() as i16;
bytes.extend_from_slice(&value.to_le_bytes());
}
self.file.write_all(&bytes)?;
self.data_bytes = self.data_bytes.saturating_add(bytes.len() as u32);
Ok(())
}
}
#[cfg(any(feature = "duplex", test))]
impl Drop for WavWriter {
fn drop(&mut self) {
use std::io::{Seek, SeekFrom};
if self.file.seek(SeekFrom::Start(0)).is_ok() {
let _ = self.file.write_all(&wav_header(self.data_bytes));
let _ = self.file.flush();
}
}
}
#[cfg(any(feature = "duplex", test))]
fn wav_header(data_bytes: u32) -> Vec<u8> {
let mut header = Vec::with_capacity(44);
header.extend_from_slice(b"RIFF");
header.extend_from_slice(&(36u32.saturating_add(data_bytes)).to_le_bytes());
header.extend_from_slice(b"WAVEfmt ");
header.extend_from_slice(&16u32.to_le_bytes());
header.extend_from_slice(&1u16.to_le_bytes());
header.extend_from_slice(&1u16.to_le_bytes());
header.extend_from_slice(&SAMPLE_RATE.to_le_bytes());
header.extend_from_slice(&(SAMPLE_RATE * 2).to_le_bytes());
header.extend_from_slice(&2u16.to_le_bytes());
header.extend_from_slice(&16u16.to_le_bytes());
header.extend_from_slice(b"data");
header.extend_from_slice(&data_bytes.to_le_bytes());
header
}
#[cfg(test)]
#[path = "tests.rs"]
mod tests;