From 648b135193348a3708a4725ad3b99b47f9d67ece Mon Sep 17 00:00:00 2001 From: John Tennant Date: Mon, 27 Jul 2026 19:34:09 -0400 Subject: [PATCH] refactor(desktop): split Pocket voice transitions Signed-off-by: John Tennant --- desktop/src-tauri/src/huddle/tts.rs | 168 +---------------- .../src/huddle/tts_voice_transition.rs | 175 ++++++++++++++++++ 2 files changed, 179 insertions(+), 164 deletions(-) create mode 100644 desktop/src-tauri/src/huddle/tts_voice_transition.rs diff --git a/desktop/src-tauri/src/huddle/tts.rs b/desktop/src-tauri/src/huddle/tts.rs index 5f6c6e373..9d93f94a6 100644 --- a/desktop/src-tauri/src/huddle/tts.rs +++ b/desktop/src-tauri/src/huddle/tts.rs @@ -47,49 +47,12 @@ use std::{ time::Duration, }; -use super::pocket::{ - load_text_to_speech, load_voice_style, VoiceStyle, SAMPLE_RATE, VOICE_FILE_EXT, -}; +use super::pocket::{load_text_to_speech, load_voice_style, SAMPLE_RATE, VOICE_FILE_EXT}; use super::preprocessing::{preprocess_for_tts, split_sentences}; -#[derive(Debug)] -struct PendingVoiceChange { - generation: u64, - acknowledged: tokio::sync::oneshot::Sender<()>, -} - -type VoiceChangeAck = Arc>>; -type WorkerVoiceState = (Arc>, Arc, VoiceChangeAck); -type WorkerCancelSignals = (Arc, Arc); -type CancelTextState<'a> = ( - &'a mpsc::Receiver, - &'a mut VecDeque, - &'a mut Option, -); -type CancelSignals<'a> = (&'a AtomicBool, &'a AtomicBool); - -#[derive(Debug)] -struct QueuedText { - generation: u64, - text: String, -} - -#[derive(Clone, Debug)] -pub(crate) struct TtsTextSender { - text_tx: SyncSender, - generation: u64, -} - -impl TtsTextSender { - pub(crate) fn send(&self, text: String) -> Result<(), String> { - self.text_tx - .send(QueuedText { - generation: self.generation, - text, - }) - .map_err(|error| error.to_string()) - } -} +#[path = "tts_voice_transition.rs"] +mod voice_transition; +use voice_transition::*; // ── Constants ───────────────────────────────────────────────────────────────── @@ -747,58 +710,6 @@ fn tts_worker( // ── Helpers ─────────────────────────────────────────────────────────────────── -fn begin_voice_change( - selected_voice: &Mutex, - voice_generation: &AtomicU64, - voice_cancel: &AtomicBool, - voice_change_ack: &VoiceChangeAck, - voice: &str, -) -> Option> { - let mut pending_ack = voice_change_ack - .lock() - .unwrap_or_else(|error| error.into_inner()); - let mut selected = selected_voice - .lock() - .unwrap_or_else(|error| error.into_inner()); - if selected.as_str() == voice { - return None; - } - - let (sender, receiver) = tokio::sync::oneshot::channel(); - voice_cancel.store(true, Ordering::Release); - let generation = voice_generation.fetch_add(1, Ordering::AcqRel) + 1; - if let Some(superseded) = pending_ack.replace(PendingVoiceChange { - generation, - acknowledged: sender, - }) { - let _ = superseded.acknowledged.send(()); - } - *selected = voice.to_string(); - Some(receiver) -} - -fn acknowledge_voice_change(voice_change_ack: &VoiceChangeAck, voice_cancel: &AtomicBool) { - let mut pending_ack = voice_change_ack - .lock() - .unwrap_or_else(|error| error.into_inner()); - if voice_cancel.load(Ordering::Acquire) { - return; - } - if let Some(pending) = pending_ack.take() { - let _ = pending.acknowledged.send(()); - } -} - -fn finish_voice_change_ack(voice_change_ack: &VoiceChangeAck) { - if let Some(pending) = voice_change_ack - .lock() - .unwrap_or_else(|error| error.into_inner()) - .take() - { - let _ = pending.acknowledged.send(()); - } -} - fn drain_tts_until_shutdown( text_rx: mpsc::Receiver, shutdown: &AtomicBool, @@ -822,52 +733,6 @@ fn drain_tts_until_shutdown( finish_voice_change_ack(voice_change_ack); } -fn reconcile_selected_voice( - model_dir: &std::path::Path, - selected_voice: &Mutex, - voice_name: &mut String, - style: &mut VoiceStyle, -) -> bool { - let requested_voice = selected_voice - .lock() - .unwrap_or_else(|error| error.into_inner()) - .clone(); - if requested_voice == *voice_name { - return true; - } - - let requested_path = model_dir.join(format!("{requested_voice}.{VOICE_FILE_EXT}")); - match load_voice_style(&requested_path) { - Ok(requested_style) => { - *style = requested_style; - *voice_name = requested_voice; - true - } - Err(error) => { - use super::pocket::DEFAULT_VOICE; - eprintln!( - "buzz-desktop: Pocket voice {requested_voice} is unavailable ({error}); falling back to Mary" - ); - let fallback_path = model_dir.join(format!("{DEFAULT_VOICE}.{VOICE_FILE_EXT}")); - match load_voice_style(&fallback_path) { - Ok(fallback_style) => { - *style = fallback_style; - *voice_name = DEFAULT_VOICE.to_string(); - *selected_voice - .lock() - .unwrap_or_else(|lock_error| lock_error.into_inner()) = - DEFAULT_VOICE.to_string(); - true - } - Err(fallback_error) => { - eprintln!("buzz-desktop: Mary voice fallback is unavailable: {fallback_error}"); - false - } - } - } - } -} - /// Check for cancel or shutdown. Returns `true` if the caller should break/continue. /// On cancel: drains the text queue and clears the cancel flag. /// @@ -930,31 +795,6 @@ fn handle_cancel_or_shutdown( false } -fn retain_cancelled_text( - deferred_text: &mut VecDeque, - current_text: &mut Option, - text_rx: &mpsc::Receiver, - preserve_generation: Option, -) { - if let Some(generation) = preserve_generation { - deferred_text.retain(|text| text.generation >= generation); - if let Some(text) = current_text.take() { - if text.generation >= generation { - deferred_text.push_front(text); - } - } - while let Ok(text) = text_rx.try_recv() { - if text.generation >= generation { - deferred_text.push_back(text); - } - } - } else { - deferred_text.clear(); - current_text.take(); - while text_rx.try_recv().is_ok() {} - } -} - /// Acquire the `player_ops` lock, recovering from poison. /// /// The data under the mutex is `()` — it only serializes Player mutations — diff --git a/desktop/src-tauri/src/huddle/tts_voice_transition.rs b/desktop/src-tauri/src/huddle/tts_voice_transition.rs new file mode 100644 index 000000000..1839db0ee --- /dev/null +++ b/desktop/src-tauri/src/huddle/tts_voice_transition.rs @@ -0,0 +1,175 @@ +use std::{ + collections::VecDeque, + path::Path, + sync::{ + atomic::{AtomicBool, AtomicU64, Ordering}, + mpsc::{self, SyncSender}, + Arc, Mutex, + }, +}; + +use crate::huddle::pocket::{load_voice_style, VoiceStyle, DEFAULT_VOICE, VOICE_FILE_EXT}; + +#[derive(Debug)] +pub(super) struct PendingVoiceChange { + pub(super) generation: u64, + acknowledged: tokio::sync::oneshot::Sender<()>, +} + +pub(super) type VoiceChangeAck = Arc>>; +pub(super) type WorkerVoiceState = (Arc>, Arc, VoiceChangeAck); +pub(super) type WorkerCancelSignals = (Arc, Arc); +pub(super) type CancelTextState<'a> = ( + &'a mpsc::Receiver, + &'a mut VecDeque, + &'a mut Option, +); +pub(super) type CancelSignals<'a> = (&'a AtomicBool, &'a AtomicBool); + +#[derive(Debug)] +pub(super) struct QueuedText { + pub(super) generation: u64, + pub(super) text: String, +} + +#[derive(Clone, Debug)] +pub(crate) struct TtsTextSender { + pub(super) text_tx: SyncSender, + pub(super) generation: u64, +} + +impl TtsTextSender { + pub(crate) fn send(&self, text: String) -> Result<(), String> { + self.text_tx + .send(QueuedText { + generation: self.generation, + text, + }) + .map_err(|error| error.to_string()) + } +} + +pub(super) fn begin_voice_change( + selected_voice: &Mutex, + voice_generation: &AtomicU64, + voice_cancel: &AtomicBool, + voice_change_ack: &VoiceChangeAck, + voice: &str, +) -> Option> { + let mut pending_ack = voice_change_ack + .lock() + .unwrap_or_else(|error| error.into_inner()); + let mut selected = selected_voice + .lock() + .unwrap_or_else(|error| error.into_inner()); + if selected.as_str() == voice { + return None; + } + + let (sender, receiver) = tokio::sync::oneshot::channel(); + voice_cancel.store(true, Ordering::Release); + let generation = voice_generation.fetch_add(1, Ordering::AcqRel) + 1; + if let Some(superseded) = pending_ack.replace(PendingVoiceChange { + generation, + acknowledged: sender, + }) { + let _ = superseded.acknowledged.send(()); + } + *selected = voice.to_string(); + Some(receiver) +} + +pub(super) fn acknowledge_voice_change( + voice_change_ack: &VoiceChangeAck, + voice_cancel: &AtomicBool, +) { + let mut pending_ack = voice_change_ack + .lock() + .unwrap_or_else(|error| error.into_inner()); + if voice_cancel.load(Ordering::Acquire) { + return; + } + if let Some(pending) = pending_ack.take() { + let _ = pending.acknowledged.send(()); + } +} + +pub(super) fn finish_voice_change_ack(voice_change_ack: &VoiceChangeAck) { + if let Some(pending) = voice_change_ack + .lock() + .unwrap_or_else(|error| error.into_inner()) + .take() + { + let _ = pending.acknowledged.send(()); + } +} + +pub(super) fn reconcile_selected_voice( + model_dir: &Path, + selected_voice: &Mutex, + voice_name: &mut String, + style: &mut VoiceStyle, +) -> bool { + let requested_voice = selected_voice + .lock() + .unwrap_or_else(|error| error.into_inner()) + .clone(); + if requested_voice == *voice_name { + return true; + } + + let requested_path = model_dir.join(format!("{requested_voice}.{VOICE_FILE_EXT}")); + match load_voice_style(&requested_path) { + Ok(requested_style) => { + *style = requested_style; + *voice_name = requested_voice; + true + } + Err(error) => { + eprintln!( + "buzz-desktop: Pocket voice {requested_voice} is unavailable ({error}); falling back to Mary" + ); + let fallback_path = model_dir.join(format!("{DEFAULT_VOICE}.{VOICE_FILE_EXT}")); + match load_voice_style(&fallback_path) { + Ok(fallback_style) => { + *style = fallback_style; + *voice_name = DEFAULT_VOICE.to_string(); + *selected_voice + .lock() + .unwrap_or_else(|lock_error| lock_error.into_inner()) = + DEFAULT_VOICE.to_string(); + true + } + Err(fallback_error) => { + eprintln!("buzz-desktop: Mary voice fallback is unavailable: {fallback_error}"); + false + } + } + } + } +} + +pub(super) fn retain_cancelled_text( + deferred_text: &mut VecDeque, + current_text: &mut Option, + text_rx: &mpsc::Receiver, + preserve_generation: Option, +) { + if let Some(generation) = preserve_generation { + deferred_text.retain(|text| text.generation >= generation); + if let Some(text) = current_text.take() { + if text.generation >= generation { + deferred_text.push_front(text); + } + } + while let Ok(text) = text_rx.try_recv() { + if text.generation >= generation { + deferred_text.push_back(text); + } + } + } else { + deferred_text.clear(); + current_text.take(); + while text_rx.try_recv().is_ok() {} + } +}