fix(desktop): preserve TTS route outcomes

Signed-off-by: John Tennant <jtennant@squareup.com>
This commit is contained in:
John Tennant
2026-07-29 12:20:52 -04:00
committed by John Tennant
parent a20a671f9b
commit 18fec8c7bb
6 changed files with 348 additions and 202 deletions
+66 -191
View File
@@ -58,6 +58,9 @@ use voice_transition::*;
#[path = "tts_startup.rs"]
mod startup;
use startup::await_worker_startup;
#[path = "tts_audio.rs"]
mod audio;
use audio::*;
// ── Constants ─────────────────────────────────────────────────────────────────
@@ -501,6 +504,33 @@ fn tts_worker(
let mut first_append = true;
let mut last_route_id = 0;
let mut deferred_text = VecDeque::new();
let append_audio = |prepared: PreparedModelAudio, route_id: u64| {
let _ops = lock_player_ops(&player_ops);
if cancel.load(Ordering::Acquire)
|| voice_cancel.load(Ordering::Acquire)
|| shutdown.load(Ordering::Acquire)
{
let reason = if shutdown.load(Ordering::Acquire) {
"shutdown"
} else if cancel.load(Ordering::Acquire) {
"barge_in"
} else {
"voice_switch"
};
eprintln!(
"buzz-desktop: tts stage=synthesis status=cancelled reason={reason} route_id={route_id}"
);
return false;
}
player.append(SamplesBuffer::new(channels, rate, prepared.buffer));
eprintln!(
"buzz-desktop: tts stage=player status=append_accepted route_id={route_id} chunk_index={} sample_count={}",
prepared.chunk_index, prepared.sample_count
);
// Set this only after append so STT remains open during synthesis.
tts_active.store(true, Ordering::Release);
true
};
loop {
let mut no_current_text = None;
@@ -575,6 +605,10 @@ fn tts_worker(
continue;
};
if queued_text.generation < voice_generation.load(Ordering::Acquire) {
eprintln!(
"buzz-desktop: tts stage=queue status=dropped reason=voice_switch route_id={}",
queued_text.route_id
);
continue;
}
let raw_text = queued_text.text;
@@ -586,6 +620,9 @@ fn tts_worker(
// queued after an unpublished pipeline is installed cannot use the
// voice captured when construction began.
if !reconcile_selected_voice(&model_dir, &selected_voice, &mut voice_name, &mut style) {
eprintln!(
"buzz-desktop: tts stage=synthesis status=failed reason=voice_unavailable route_id={route_id}"
);
continue;
}
@@ -670,8 +707,8 @@ fn tts_worker(
);
continue;
}
let model_chunk_count = model_chunks.len();
for (model_chunk_index, model_chunk) in model_chunks.iter().enumerate() {
let mut playback_audio = PlaybackChunkAudio::new();
for model_chunk in &model_chunks {
let chunk_index = model_unit_index;
model_unit_index += 1;
let mut no_current_text = None;
@@ -689,7 +726,6 @@ fn tts_worker(
break 'playback_chunks;
}
let ends_playback_chunk = model_chunk_index + 1 == model_chunk_count;
let synthesis = engine.synth_chunk(model_chunk, "en", &style, SYNTH_STEPS);
if cancel.load(Ordering::Acquire)
|| voice_cancel.load(Ordering::Acquire)
@@ -715,67 +751,21 @@ fn tts_worker(
}
match synthesis {
Ok(samples) if !samples.is_empty() => {
let synthesized_samples = samples.len();
let mut audio = clamp_to_full_scale(samples);
if ends_playback_chunk {
// Fade only at the playback-chunk boundary. Applying
// it at the model's internal token boundary would
// create an audible dip between contiguous units.
apply_fade_out(&mut audio);
}
let buf = build_sentence_append_buffer(
if let Some(prepared) = playback_audio.push(
samples,
chunk_index,
&mut first_append,
audio,
silence_buf_len,
model_chunk_index == 0 || player.empty(),
ends_playback_chunk,
);
// Check-and-append under `player_ops`, serialized with
// the monitor: a barge-in may have arrived during
// synthesis (the blocking window the monitor thread
// exists for). Don't append the now-stale sentence — the
// human interrupted; speaking it anyway would talk over
// them. Holding the lock for the check + append means the
// monitor can never clear between our check passing and
// the buffer landing. The flag is deliberately NOT
// consumed here: the loop-top handle_cancel_or_shutdown
// does the full consume (drain queue, reset lead-in) on
// the next iteration.
let _ops = lock_player_ops(&player_ops);
if cancel.load(Ordering::Acquire)
|| voice_cancel.load(Ordering::Acquire)
|| shutdown.load(Ordering::Acquire)
{
// Nothing appended; the loop-top consume re-arms
// `first_append` (the flag is still set — the worker
// is its only consumer).
let reason = if shutdown.load(Ordering::Acquire) {
"shutdown"
} else if cancel.load(Ordering::Acquire) {
"barge_in"
} else {
"voice_switch"
};
eprintln!(
"buzz-desktop: tts stage=synthesis status=cancelled reason={reason} route_id={route_id}"
);
synthesis_outcome = "cancelled";
break 'playback_chunks;
player.empty(),
) {
if !append_audio(prepared, route_id) {
first_append = true;
synthesis_outcome = "cancelled";
break 'playback_chunks;
}
appended_audio = true;
last_route_id = route_id;
}
player.append(SamplesBuffer::new(channels, rate, buf));
appended_audio = true;
last_route_id = route_id;
eprintln!(
"buzz-desktop: tts stage=player status=append_accepted route_id={route_id} chunk_index={chunk_index} sample_count={synthesized_samples}"
);
// NOTE: tts_active is set AFTER player.append(), not
// before. Setting it before synthesis would cause STT to
// discard user speech during the synthesis window as
// "echo" even though no audio is actually playing yet.
// See crossfire review C3.
tts_active.store(true, Ordering::Release);
}
Ok(_) => {
eprintln!(
@@ -787,10 +777,24 @@ fn tts_worker(
"buzz-desktop: tts stage=synthesis status=failed reason=inference route_id={route_id} chunk_index={chunk_index}"
);
synthesis_outcome = "failed";
break 'playback_chunks;
break;
}
}
}
if let Some(prepared) =
playback_audio.finish(&mut first_append, silence_buf_len, player.empty())
{
if !append_audio(prepared, route_id) {
first_append = true;
synthesis_outcome = "cancelled";
break 'playback_chunks;
}
appended_audio = true;
last_route_id = route_id;
}
if synthesis_outcome == "failed" {
break 'playback_chunks;
}
}
if synthesis_outcome == "completed" && appended_audio {
eprintln!("buzz-desktop: tts stage=synthesis status=completed route_id={route_id}");
@@ -896,135 +900,6 @@ fn lock_player_ops(ops: &Mutex<()>) -> MutexGuard<'_, ()> {
ops.lock().unwrap_or_else(PoisonError::into_inner)
}
/// Hard-clamp samples to ±1.0 full scale.
///
/// No gain is applied because Pocket TTS already emits speech-level audio and
/// the reference pipeline applies no output scaling. Normalizing each sentence
/// would cause level pumping between chunks. The clamp remains only as a safety
/// net against outlier transients.
fn clamp_to_full_scale(samples: Vec<f32>) -> Vec<f32> {
samples.into_iter().map(|s| s.clamp(-1.0, 1.0)).collect()
}
/// Apply a short linear fade-out at the *end* of `samples`.
///
/// Uses `FADE_OUT_SAMPLES` (8 ms) or half the buffer length, whichever is
/// smaller. Eliminates the click that occurs when a non-zero waveform
/// terminates abruptly at a sentence boundary.
///
/// # Why no fade-in
///
/// A symmetric fade-in would attenuate the leading consonant attack because
/// Pocket TTS produces real audio energy inside the first millisecond. A
/// linear 0→1 ramp over 192 samples scales those onset samples by ≤50% for the
/// first ~4 ms, which can make the first phoneme sound clipped.
///
/// The first sample of Pocket output measures ≈ 0.0018 (≈ −54 dBFS) — well
/// below the threshold at which a DC-jump would be audible as a click — so
/// no fade-in is needed. The OS audio device gets its quiet ramp-up window
/// from `SENTENCE_LEAD_IN_SAMPLES` instead, inserted as pure silence before
/// each sentence buffer.
fn apply_fade_out(samples: &mut [f32]) {
let len = samples.len();
let fade = FADE_OUT_SAMPLES.min(len / 2);
for i in 0..fade {
samples[len - 1 - i] *= i as f32 / fade as f32;
}
}
/// Build one buffer appended to the rodio `Player` for a synthesis unit.
///
/// Every playback boundary gets a short lead-in pad immediately before its
/// audio. This matters for chunks that start with soft first phonemes (`I'm`,
/// `I've`): the synthesized buffer can begin with speech within the first
/// millisecond, so the playback layer must provide the device/mixer cushion.
/// To keep the audible gap unchanged, the trailing silence after this chunk is
/// shortened by the same amount (`silence_buf_len - SENTENCE_LEAD_IN_SAMPLES`):
/// sentence N contributes 80 ms of post-speech silence and sentence N+1
/// contributes the remaining 20 ms of pre-speech cushion.
///
/// The lead-in, audio, and trailing silence are concatenated into one
/// `SamplesBuffer` before appending. This keeps rodio's queue shape at one
/// tracked source per synthesized sentence, avoiding source-boundary/drain
/// regressions from enqueueing the lead-in, audio, and tail as separate sounds.
///
/// A playback chunk may contain several model-sized synthesis units. Only the
/// first unit receives the onset cushion and only the last receives the
/// remaining gap. If playback underruns while the next unit is synthesized,
/// that unit becomes a new playback boundary and receives a fresh cushion.
///
/// `first_append` is flipped on the first call after the player goes idle.
/// The worker uses it in the idle branch of the main loop to distinguish
/// "never queued anything since last drain" from "drained after speaking",
/// which controls when `tts_active` is released and the lead-in re-armed.
fn build_sentence_append_buffer(
first_append: &mut bool,
audio: Vec<f32>,
silence_buf_len: usize,
starts_playback_chunk: bool,
ends_playback_chunk: bool,
) -> Vec<f32> {
if *first_append {
*first_append = false;
}
let lead_in_len = if starts_playback_chunk {
SENTENCE_LEAD_IN_SAMPLES
} else {
0
};
let trailing_silence_len = if ends_playback_chunk {
silence_buf_len.saturating_sub(SENTENCE_LEAD_IN_SAMPLES)
} else {
0
};
let mut buf = Vec::with_capacity(lead_in_len + audio.len() + trailing_silence_len);
buf.extend(std::iter::repeat_n(0.0_f32, lead_in_len));
buf.extend(audio);
buf.extend(std::iter::repeat_n(0.0_f32, trailing_silence_len));
buf
}
/// Group sentences into synthesis chunks.
///
/// The first sentence always stands alone — it is what the listener hears
/// first, and synthesizing it by itself keeps time-to-first-audio at the
/// single-sentence cost. Subsequent sentences pack greedily: a sentence
/// joins the current chunk while the combined length stays within
/// `max_chars`; otherwise it starts a new chunk. A single sentence longer
/// than `max_chars` becomes its own chunk here, then the Pocket engine splits
/// it at the April bundle's exact token limit before synthesis.
///
/// Sentences within a chunk are joined with a single space; sentence-ending
/// punctuation is preserved by `split_sentences`, so the model sees natural
/// multi-sentence prose — the same shape upstream's ~50-token chunker feeds it.
fn group_sentences_into_chunks(sentences: &[String], max_chars: usize) -> Vec<String> {
let mut chunks: Vec<String> = Vec::new();
for (i, sentence) in sentences.iter().enumerate() {
let sentence = sentence.trim();
if sentence.is_empty() {
continue;
}
if i == 0 || chunks.is_empty() {
chunks.push(sentence.to_string());
continue;
}
// Never merge into the first chunk — it's the latency-critical one.
let can_merge = chunks.len() > 1
&& chunks
.last()
.is_some_and(|c| c.len() + 1 + sentence.len() <= max_chars);
if can_merge {
let last = chunks.last_mut().expect("non-empty checked above");
last.push(' ');
last.push_str(sentence);
} else {
chunks.push(sentence.to_string());
}
}
chunks
}
// ── Tests ─────────────────────────────────────────────────────────────────────
#[cfg(test)]
+235
View File
@@ -0,0 +1,235 @@
use super::{FADE_OUT_SAMPLES, SENTENCE_LEAD_IN_SAMPLES};
pub(super) struct PreparedModelAudio {
pub(super) buffer: Vec<f32>,
pub(super) sample_count: usize,
pub(super) chunk_index: usize,
}
/// Holds one synthesized model unit so playback-boundary decoration is based
/// on the first and last unit that actually produced audio.
pub(super) struct PlaybackChunkAudio {
pending: Option<(Vec<f32>, usize)>,
appended: bool,
}
impl PlaybackChunkAudio {
pub(super) fn new() -> Self {
Self {
pending: None,
appended: false,
}
}
pub(super) fn push(
&mut self,
samples: Vec<f32>,
chunk_index: usize,
first_append: &mut bool,
silence_buf_len: usize,
playback_idle: bool,
) -> Option<PreparedModelAudio> {
if samples.is_empty() {
return None;
}
let previous = self.pending.replace((samples, chunk_index))?;
let prepared = prepare_model_audio(
previous,
first_append,
silence_buf_len,
!self.appended || playback_idle,
false,
);
self.appended = true;
Some(prepared)
}
pub(super) fn finish(
&mut self,
first_append: &mut bool,
silence_buf_len: usize,
playback_idle: bool,
) -> Option<PreparedModelAudio> {
let pending = self.pending.take()?;
Some(prepare_model_audio(
pending,
first_append,
silence_buf_len,
!self.appended || playback_idle,
true,
))
}
}
fn prepare_model_audio(
(samples, chunk_index): (Vec<f32>, usize),
first_append: &mut bool,
silence_buf_len: usize,
starts_playback_chunk: bool,
ends_playback_chunk: bool,
) -> PreparedModelAudio {
let sample_count = samples.len();
let mut audio = clamp_to_full_scale(samples);
if ends_playback_chunk {
apply_fade_out(&mut audio);
}
PreparedModelAudio {
buffer: build_sentence_append_buffer(
first_append,
audio,
silence_buf_len,
starts_playback_chunk,
ends_playback_chunk,
),
sample_count,
chunk_index,
}
}
/// Hard-clamp samples to ±1.0 full scale.
pub(super) fn clamp_to_full_scale(samples: Vec<f32>) -> Vec<f32> {
samples.into_iter().map(|s| s.clamp(-1.0, 1.0)).collect()
}
/// Apply a short linear fade-out to avoid a discontinuity at playback boundaries.
pub(super) fn apply_fade_out(samples: &mut [f32]) {
let len = samples.len();
let fade = FADE_OUT_SAMPLES.min(len / 2);
for i in 0..fade {
samples[len - 1 - i] *= i as f32 / fade as f32;
}
}
pub(super) fn build_sentence_append_buffer(
first_append: &mut bool,
audio: Vec<f32>,
silence_buf_len: usize,
starts_playback_chunk: bool,
ends_playback_chunk: bool,
) -> Vec<f32> {
if *first_append {
*first_append = false;
}
let lead_in_len = if starts_playback_chunk {
SENTENCE_LEAD_IN_SAMPLES
} else {
0
};
let trailing_silence_len = if ends_playback_chunk {
silence_buf_len.saturating_sub(SENTENCE_LEAD_IN_SAMPLES)
} else {
0
};
let mut buffer = Vec::with_capacity(lead_in_len + audio.len() + trailing_silence_len);
buffer.extend(std::iter::repeat_n(0.0_f32, lead_in_len));
buffer.extend(audio);
buffer.extend(std::iter::repeat_n(0.0_f32, trailing_silence_len));
buffer
}
pub(super) fn group_sentences_into_chunks(sentences: &[String], max_chars: usize) -> Vec<String> {
let mut chunks: Vec<String> = Vec::new();
for (index, sentence) in sentences.iter().enumerate() {
let sentence = sentence.trim();
if sentence.is_empty() {
continue;
}
if index == 0 || chunks.is_empty() {
chunks.push(sentence.to_string());
continue;
}
let can_merge = chunks.len() > 1
&& chunks
.last()
.is_some_and(|chunk| chunk.len() + 1 + sentence.len() <= max_chars);
if can_merge {
if let Some(last) = chunks.last_mut() {
last.push(' ');
last.push_str(sentence);
}
} else {
chunks.push(sentence.to_string());
}
}
chunks
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn multi_unit_audio_decorates_only_outer_playback_boundaries() {
let mut chunk = PlaybackChunkAudio::new();
let mut first_append = true;
let silence = SENTENCE_LEAD_IN_SAMPLES + 100;
assert!(chunk
.push(vec![0.4; 16], 0, &mut first_append, silence, false)
.is_none());
let first = chunk
.push(vec![0.5; 16], 1, &mut first_append, silence, false)
.expect("first ready model unit");
assert_eq!(first.buffer.len(), SENTENCE_LEAD_IN_SAMPLES + 16);
assert!(first.buffer[..SENTENCE_LEAD_IN_SAMPLES]
.iter()
.all(|sample| *sample == 0.0));
assert_eq!(first.buffer[SENTENCE_LEAD_IN_SAMPLES], 0.4);
let last = chunk
.finish(&mut first_append, silence, false)
.expect("last ready model unit");
assert_eq!(last.buffer.len(), 16 + 100);
assert_eq!(last.buffer.last(), Some(&0.0));
}
#[test]
fn empty_edge_units_do_not_steal_lead_in_or_trailing_boundary() {
let mut chunk = PlaybackChunkAudio::new();
let mut first_append = true;
let silence = SENTENCE_LEAD_IN_SAMPLES + 100;
assert!(chunk
.push(Vec::new(), 0, &mut first_append, silence, false)
.is_none());
assert!(chunk
.push(vec![0.5; 16], 1, &mut first_append, silence, false)
.is_none());
assert!(chunk
.push(Vec::new(), 2, &mut first_append, silence, false)
.is_none());
let only = chunk
.finish(&mut first_append, silence, false)
.expect("only audible model unit");
assert_eq!(only.buffer.len(), SENTENCE_LEAD_IN_SAMPLES + 16 + 100);
assert!(only.buffer[..SENTENCE_LEAD_IN_SAMPLES]
.iter()
.all(|sample| *sample == 0.0));
assert_eq!(only.buffer.last(), Some(&0.0));
}
#[test]
fn playback_underrun_rearms_the_onset_cushion() {
let mut chunk = PlaybackChunkAudio::new();
let mut first_append = true;
let silence = SENTENCE_LEAD_IN_SAMPLES + 100;
assert!(chunk
.push(vec![0.4; 16], 0, &mut first_append, silence, false)
.is_none());
let first = chunk
.push(vec![0.5; 16], 1, &mut first_append, silence, false)
.expect("first model unit");
assert_eq!(first.buffer.len(), SENTENCE_LEAD_IN_SAMPLES + 16);
let after_underrun = chunk
.push(vec![0.6; 16], 2, &mut first_append, silence, true)
.expect("model unit after underrun");
assert_eq!(after_underrun.buffer.len(), SENTENCE_LEAD_IN_SAMPLES + 16);
assert!(after_underrun.buffer[..SENTENCE_LEAD_IN_SAMPLES]
.iter()
.all(|sample| *sample == 0.0));
}
}
@@ -158,20 +158,40 @@ pub(super) fn retain_cancelled_text(
preserve_generation: Option<u64>,
) {
if let Some(generation) = preserve_generation {
deferred_text.retain(|text| text.generation >= generation);
deferred_text.retain(|text| {
let preserve = text.generation >= generation;
if !preserve {
log_cancelled_route(text.route_id, "voice_switch");
}
preserve
});
if let Some(text) = current_text.take() {
if text.generation >= generation {
deferred_text.push_front(text);
} else {
log_cancelled_route(text.route_id, "voice_switch");
}
}
while let Ok(text) = text_rx.try_recv() {
if text.generation >= generation {
deferred_text.push_back(text);
} else {
log_cancelled_route(text.route_id, "voice_switch");
}
}
} else {
deferred_text.clear();
current_text.take();
while text_rx.try_recv().is_ok() {}
for text in deferred_text.drain(..) {
log_cancelled_route(text.route_id, "barge_in");
}
if let Some(text) = current_text.take() {
log_cancelled_route(text.route_id, "barge_in");
}
while let Ok(text) = text_rx.try_recv() {
log_cancelled_route(text.route_id, "barge_in");
}
}
}
fn log_cancelled_route(route_id: u64, reason: &str) {
eprintln!("buzz-desktop: tts stage=queue status=dropped reason={reason} route_id={route_id}");
}
@@ -187,17 +187,23 @@ test("queues agent messages in live thread arrival order", async () => {
test("disabling cancels queued speech and rejects new messages until enabled", async () => {
const invoked = [];
const dropped = [];
let releaseFirst;
const firstBlocked = new Promise((resolve) => {
releaseFirst = resolve;
});
const speaker = createOrderedSpeaker(async (text) => {
invoked.push(text);
if (text === "first") await firstBlocked;
}, assert.fail);
const speaker = createOrderedSpeaker(
async (text) => {
invoked.push(text);
if (text === "first") await firstBlocked;
},
assert.fail,
true,
(routeId, reason) => dropped.push([routeId, reason]),
);
speaker.enqueue("first");
speaker.enqueue("queued-before-off");
speaker.enqueue("first", 51);
speaker.enqueue("queued-before-off", 52);
await Promise.resolve();
speaker.setEnabled(false);
speaker.enqueue("while-off");
@@ -207,6 +213,7 @@ test("disabling cancels queued speech and rejects new messages until enabled", a
speaker.enqueue("after-on");
await new Promise((resolve) => setTimeout(resolve, 0));
assert.deepEqual(invoked, ["first", "after-on"]);
assert.deepEqual(dropped, [[52, "disabled"]]);
});
test("does not speak before the native enabled state is known", async () => {
@@ -112,6 +112,7 @@ export function createOrderedSpeaker(
speak: (text: string, routeId: number) => Promise<void>,
onError: (error: unknown) => void,
initiallyEnabled = true,
onDrop: (routeId: number, reason: "disabled") => void = () => {},
): {
enqueue: (text: string, routeId?: number) => "queued" | "disabled";
setEnabled: (enabled: boolean) => void;
@@ -125,7 +126,10 @@ export function createOrderedSpeaker(
const queuedGeneration = generation;
tail = tail
.then(() => {
if (!enabled || generation !== queuedGeneration) return;
if (!enabled || generation !== queuedGeneration) {
onDrop(routeId, "disabled");
return;
}
return speak(text, routeId);
})
.catch(onError);
@@ -71,6 +71,11 @@ export function useTtsSubscription(
},
() => {},
false,
(routeId, reason) => {
console.debug(
`[huddle] tts stage=queue status=dropped reason=${reason} route_id=${routeId}`,
);
},
);
const deliver = ({