mirror of
https://github.com/block/buzz.git
synced 2026-08-18 06:50:31 +02:00
refactor(desktop): satisfy file size ratchets
Signed-off-by: John Tennant <jtennant@squareup.com>
This commit is contained in:
committed by
John Tennant
parent
8d31bf79a9
commit
8a25584a78
@@ -48,18 +48,13 @@ pub struct AppState {
|
||||
pub channel_templates_store_lock: Mutex<()>,
|
||||
pub managed_agent_processes: Mutex<HashMap<ManagedAgentRuntimeKey, ManagedAgentPairRuntime>>,
|
||||
pub huddle_state: Mutex<HuddleState>,
|
||||
pub tts_settings: Mutex<crate::huddle::tts_settings::TtsSettings>,
|
||||
pub tts_settings_load_error: Mutex<Option<String>>,
|
||||
pub tts_settings_transition: tokio::sync::Mutex<()>,
|
||||
pub huddle_audio: crate::huddle::tts_settings::HuddleAudioSettingsState,
|
||||
/// Tauri app handle — stored after setup so huddle commands can emit
|
||||
/// `huddle-state-changed` events without needing the handle threaded
|
||||
/// through every call site.
|
||||
///
|
||||
/// Set once during `setup()` in `lib.rs`; never cleared.
|
||||
pub app_handle: Mutex<Option<AppHandle>>,
|
||||
/// Selected audio output device name. `None` = system default.
|
||||
/// Used by `connect_audio_relay` and TTS pipeline when opening sinks.
|
||||
pub audio_output_device: Mutex<Option<String>>,
|
||||
/// Port of the localhost media streaming proxy (set during setup).
|
||||
pub media_proxy_port: AtomicU16,
|
||||
/// Set when identity resolution detected a "keyring-locked" state: the
|
||||
@@ -216,11 +211,8 @@ pub fn build_app_state() -> AppState {
|
||||
managed_agent_processes: Mutex::new(HashMap::new()),
|
||||
session_config_cache: Mutex::new(HashMap::new()),
|
||||
huddle_state: Mutex::new(HuddleState::default()),
|
||||
tts_settings: Mutex::new(Default::default()),
|
||||
tts_settings_load_error: Mutex::new(None),
|
||||
tts_settings_transition: tokio::sync::Mutex::new(()),
|
||||
huddle_audio: Default::default(),
|
||||
app_handle: Mutex::new(None),
|
||||
audio_output_device: Mutex::new(None),
|
||||
media_proxy_port: AtomicU16::new(0),
|
||||
prevent_sleep: Arc::new(Mutex::new(
|
||||
crate::prevent_sleep::PreventSleepState::default(),
|
||||
|
||||
@@ -39,7 +39,8 @@ fn list_audio_output_devices_blocking() -> Result<Vec<AudioOutputDevice>, String
|
||||
#[tauri::command]
|
||||
pub fn set_audio_output_device(name: String, state: State<'_, AppState>) -> Result<(), String> {
|
||||
let mut guard = state
|
||||
.audio_output_device
|
||||
.huddle_audio
|
||||
.output_device
|
||||
.lock()
|
||||
.map_err(|e| e.to_string())?;
|
||||
*guard = if name.is_empty() { None } else { Some(name) };
|
||||
@@ -50,7 +51,8 @@ pub fn set_audio_output_device(name: String, state: State<'_, AppState>) -> Resu
|
||||
#[tauri::command]
|
||||
pub fn get_audio_output_device(state: State<'_, AppState>) -> Result<String, String> {
|
||||
let guard = state
|
||||
.audio_output_device
|
||||
.huddle_audio
|
||||
.output_device
|
||||
.lock()
|
||||
.map_err(|e| e.to_string())?;
|
||||
Ok(guard.clone().unwrap_or_default())
|
||||
|
||||
@@ -220,12 +220,14 @@ pub(crate) async fn maybe_start_tts_pipeline(state: &AppState) -> Result<bool, S
|
||||
// Construct outside the lock — this spawns the TTS worker thread and
|
||||
// loads ONNX sessions (~200ms). If this fails, clear the sentinel.
|
||||
let output_device = state
|
||||
.audio_output_device
|
||||
.huddle_audio
|
||||
.output_device
|
||||
.lock()
|
||||
.unwrap_or_else(|e| e.into_inner())
|
||||
.clone();
|
||||
let initial_voice = state
|
||||
.tts_settings
|
||||
.huddle_audio
|
||||
.tts
|
||||
.lock()
|
||||
.map_err(|error| format!("text-to-speech settings lock poisoned: {error}"))
|
||||
.map(|settings| {
|
||||
@@ -315,7 +317,8 @@ fn finalize_tts_pipeline_start(
|
||||
return Ok(false);
|
||||
}
|
||||
let voice = state
|
||||
.tts_settings
|
||||
.huddle_audio
|
||||
.tts
|
||||
.lock()
|
||||
.map_err(|error| format!("text-to-speech settings lock poisoned: {error}"))
|
||||
.map(|settings| {
|
||||
@@ -525,7 +528,8 @@ mod tts_start_race_tests {
|
||||
.tts_starting
|
||||
.load(Ordering::Acquire));
|
||||
state
|
||||
.tts_settings
|
||||
.huddle_audio
|
||||
.tts
|
||||
.lock()
|
||||
.expect("text-to-speech settings")
|
||||
.voice_preferences = vec!["pocket:eve".to_string()];
|
||||
|
||||
@@ -164,7 +164,8 @@ pub(crate) async fn connect_audio_relay(
|
||||
let cancel_clone = cancel.clone();
|
||||
let (pcm_tx, pcm_rx) = tokio::sync::mpsc::channel::<Vec<u8>>(50);
|
||||
let output_device_name = state
|
||||
.audio_output_device
|
||||
.huddle_audio
|
||||
.output_device
|
||||
.lock()
|
||||
.unwrap_or_else(|e| e.into_inner())
|
||||
.clone();
|
||||
|
||||
@@ -8,7 +8,7 @@
|
||||
|
||||
use std::{
|
||||
path::{Path, PathBuf},
|
||||
sync::Arc,
|
||||
sync::{Arc, Mutex},
|
||||
time::Duration,
|
||||
};
|
||||
|
||||
@@ -37,6 +37,16 @@ type VoiceChangeWait = (
|
||||
const VOICE_AVAILABILITY_BUNDLED: &str = "bundled";
|
||||
const VOICE_AVAILABILITY_INSTALLED: &str = "installed";
|
||||
|
||||
/// Installation-global huddle audio and speech preferences.
|
||||
#[derive(Default)]
|
||||
pub struct HuddleAudioSettingsState {
|
||||
pub tts: Mutex<TtsSettings>,
|
||||
pub tts_load_error: Mutex<Option<String>>,
|
||||
pub tts_transition: tokio::sync::Mutex<()>,
|
||||
/// Selected huddle output device. `None` uses the system default.
|
||||
pub output_device: Mutex<Option<String>>,
|
||||
}
|
||||
|
||||
#[derive(Debug, Clone, Serialize, PartialEq, Eq)]
|
||||
#[serde(rename_all = "camelCase")]
|
||||
pub struct VoiceRegistryEntry {
|
||||
@@ -261,7 +271,8 @@ pub fn load_for_app(app: &AppHandle) -> (TtsSettings, Option<String>) {
|
||||
#[tauri::command]
|
||||
pub fn get_tts_settings(state: State<'_, AppState>) -> Result<TtsSettings, String> {
|
||||
if let Some(error) = state
|
||||
.tts_settings_load_error
|
||||
.huddle_audio
|
||||
.tts_load_error
|
||||
.lock()
|
||||
.map_err(|lock_error| format!("text-to-speech settings lock poisoned: {lock_error}"))?
|
||||
.clone()
|
||||
@@ -271,7 +282,8 @@ pub fn get_tts_settings(state: State<'_, AppState>) -> Result<TtsSettings, Strin
|
||||
));
|
||||
}
|
||||
state
|
||||
.tts_settings
|
||||
.huddle_audio
|
||||
.tts
|
||||
.lock()
|
||||
.map(|settings| settings.clone())
|
||||
.map_err(|error| format!("text-to-speech settings lock poisoned: {error}"))
|
||||
@@ -284,7 +296,8 @@ pub fn list_voice_registry() -> Vec<VoiceRegistryEntry> {
|
||||
|
||||
fn ensure_settings_writable(state: &AppState) -> Result<(), String> {
|
||||
if let Some(error) = state
|
||||
.tts_settings_load_error
|
||||
.huddle_audio
|
||||
.tts_load_error
|
||||
.lock()
|
||||
.map_err(|lock_error| format!("text-to-speech settings lock poisoned: {lock_error}"))?
|
||||
.as_ref()
|
||||
@@ -321,7 +334,8 @@ fn disable_tts_runtime(state: &AppState) -> Result<(), String> {
|
||||
|
||||
fn commit_effective_off(state: &AppState) -> Result<(), String> {
|
||||
state
|
||||
.tts_settings
|
||||
.huddle_audio
|
||||
.tts
|
||||
.lock()
|
||||
.map_err(|error| format!("text-to-speech settings lock poisoned: {error}"))?
|
||||
.agent_text_to_speech = false;
|
||||
@@ -370,7 +384,8 @@ async fn apply_tts_settings(
|
||||
save_to_path(&settings_path(app)?, &settings)?;
|
||||
|
||||
*state
|
||||
.tts_settings
|
||||
.huddle_audio
|
||||
.tts
|
||||
.lock()
|
||||
.map_err(|error| format!("text-to-speech settings lock poisoned: {error}"))? =
|
||||
settings.clone();
|
||||
@@ -399,7 +414,8 @@ async fn apply_tts_settings(
|
||||
|
||||
fn current_settings(state: &AppState) -> Result<TtsSettings, String> {
|
||||
state
|
||||
.tts_settings
|
||||
.huddle_audio
|
||||
.tts
|
||||
.lock()
|
||||
.map_err(|error| format!("text-to-speech settings lock poisoned: {error}"))
|
||||
.map(|settings| settings.clone())
|
||||
@@ -448,9 +464,10 @@ pub async fn set_tts_enabled(
|
||||
app: AppHandle,
|
||||
state: State<'_, AppState>,
|
||||
) -> Result<TtsSettings, String> {
|
||||
let transition = state.tts_settings_transition.lock().await;
|
||||
let transition = state.huddle_audio.tts_transition.lock().await;
|
||||
let mut settings = state
|
||||
.tts_settings
|
||||
.huddle_audio
|
||||
.tts
|
||||
.lock()
|
||||
.map_err(|error| format!("text-to-speech settings lock poisoned: {error}"))?
|
||||
.clone();
|
||||
@@ -491,9 +508,10 @@ pub async fn set_pocket_voice(
|
||||
app: AppHandle,
|
||||
state: State<'_, AppState>,
|
||||
) -> Result<TtsSettings, String> {
|
||||
let transition = state.tts_settings_transition.lock().await;
|
||||
let transition = state.huddle_audio.tts_transition.lock().await;
|
||||
let settings = state
|
||||
.tts_settings
|
||||
.huddle_audio
|
||||
.tts
|
||||
.lock()
|
||||
.map_err(|error| format!("text-to-speech settings lock poisoned: {error}"))?
|
||||
.clone();
|
||||
@@ -525,7 +543,8 @@ pub async fn preview_pocket_voice(
|
||||
}
|
||||
let model_dir = models::tts_model_dir().ok_or("Pocket voice files are unavailable")?;
|
||||
let output_device = state
|
||||
.audio_output_device
|
||||
.huddle_audio
|
||||
.output_device
|
||||
.lock()
|
||||
.unwrap_or_else(|error| error.into_inner())
|
||||
.clone();
|
||||
@@ -785,7 +804,7 @@ mod tests {
|
||||
|
||||
// This models the next command after the OFF save fails: it must merge
|
||||
// from effective memory state, not the stale last-persisted ON value.
|
||||
let current = state.tts_settings.lock().expect("settings").clone();
|
||||
let current = state.huddle_audio.tts.lock().expect("settings").clone();
|
||||
let voice_update =
|
||||
settings_with_pocket_voice(current, EVE_VOICE_KEY).expect("available voice");
|
||||
assert!(!voice_update.agent_text_to_speech);
|
||||
@@ -795,16 +814,17 @@ mod tests {
|
||||
fn failed_disabled_voice_save_does_not_change_the_remembered_voice() {
|
||||
let state = crate::app_state::build_app_state();
|
||||
state
|
||||
.tts_settings
|
||||
.huddle_audio
|
||||
.tts
|
||||
.lock()
|
||||
.expect("settings")
|
||||
.agent_text_to_speech = false;
|
||||
let current = state.tts_settings.lock().expect("settings").clone();
|
||||
let current = state.huddle_audio.tts.lock().expect("settings").clone();
|
||||
let unsaved = settings_with_pocket_voice(current, EVE_VOICE_KEY).expect("available voice");
|
||||
|
||||
// This is the only pre-persistence mutation for an OFF candidate.
|
||||
commit_effective_off(&state).expect("commit effective OFF state");
|
||||
let remembered = state.tts_settings.lock().expect("settings").clone();
|
||||
let remembered = state.huddle_audio.tts.lock().expect("settings").clone();
|
||||
assert_eq!(remembered.voice_preferences, vec![MARY_VOICE_KEY]);
|
||||
assert_eq!(unsaved.voice_preferences, vec![EVE_VOICE_KEY]);
|
||||
}
|
||||
|
||||
@@ -454,10 +454,10 @@ pub fn run() {
|
||||
|
||||
let (tts_settings, tts_settings_load_error) =
|
||||
huddle::tts_settings::load_for_app(&app_handle);
|
||||
if let Ok(mut guard) = state.tts_settings.lock() {
|
||||
if let Ok(mut guard) = state.huddle_audio.tts.lock() {
|
||||
*guard = tts_settings.clone();
|
||||
}
|
||||
if let Ok(mut guard) = state.tts_settings_load_error.lock() {
|
||||
if let Ok(mut guard) = state.huddle_audio.tts_load_error.lock() {
|
||||
*guard = tts_settings_load_error;
|
||||
}
|
||||
if let Ok(mut huddle) = state.huddle_state.lock() {
|
||||
|
||||
@@ -12,7 +12,7 @@ import {
|
||||
AUTH_TIMEOUT_MS,
|
||||
HISTORY_TIMEOUT_MS,
|
||||
PUBLISH_TIMEOUT_MS,
|
||||
} from "@/shared/api/relayClientSession";
|
||||
} from "@/shared/api/relayClientTimings";
|
||||
|
||||
type PendingHistory = {
|
||||
events: RelayEvent[];
|
||||
|
||||
@@ -51,34 +51,19 @@ import {
|
||||
shouldScheduleReconnect,
|
||||
} from "@/shared/api/relayReconnectPolicy";
|
||||
import { RelayStallWatchdog } from "@/shared/api/relayStallWatchdog";
|
||||
import {
|
||||
AUTH_TIMEOUT_MS,
|
||||
BACKOFF_RESET_STABLE_MS,
|
||||
EVENT_BATCH_MS,
|
||||
HISTORY_TIMEOUT_MS,
|
||||
PUBLISH_TIMEOUT_MS,
|
||||
RECONNECT_BASE_DELAY_MS,
|
||||
RECONNECT_MAX_DELAY_MS,
|
||||
STALL_CHECK_INTERVAL_MS,
|
||||
STALL_IDLE_TIMEOUT_MS,
|
||||
} from "@/shared/api/relayClientTimings";
|
||||
import { closeWebSocket } from "@/shared/api/relayWebSocketClose";
|
||||
import { buildThreadReferenceTags } from "@/features/messages/lib/threading";
|
||||
const RECONNECT_BASE_DELAY_MS = 1_000,
|
||||
RECONNECT_MAX_DELAY_MS = 30_000,
|
||||
EVENT_BATCH_MS = 16;
|
||||
|
||||
/**
|
||||
* Op-level timeout constants. Raised from 8 s to 25 s to survive degraded
|
||||
* networks where TLS handshakes and DNS resolution can take 3–10 s.
|
||||
*/
|
||||
export const AUTH_TIMEOUT_MS = 25_000;
|
||||
export const HISTORY_TIMEOUT_MS = 25_000;
|
||||
export const PUBLISH_TIMEOUT_MS = 25_000;
|
||||
|
||||
/**
|
||||
* The connection must remain stable for this long after a successful AUTH
|
||||
* before the reconnect backoff delay resets to its base value. Stability-
|
||||
* gated reset prevents repeated fast reconnects (flapping) from erasing the
|
||||
* backoff that throttles them.
|
||||
*/
|
||||
export const BACKOFF_RESET_STABLE_MS = 60_000;
|
||||
|
||||
/**
|
||||
* Passive liveness check. The relay sends heartbeat pings every 30s; if no
|
||||
* inbound frame arrives for two heartbeat windows, treat the socket as stalled.
|
||||
*/
|
||||
const STALL_CHECK_INTERVAL_MS = 10_000;
|
||||
const STALL_IDLE_TIMEOUT_MS = 60_000;
|
||||
|
||||
export class RelayClient {
|
||||
private wsId: number | null = null;
|
||||
|
||||
@@ -0,0 +1,20 @@
|
||||
export const RECONNECT_BASE_DELAY_MS = 1_000;
|
||||
export const RECONNECT_MAX_DELAY_MS = 30_000;
|
||||
export const EVENT_BATCH_MS = 16;
|
||||
|
||||
/**
|
||||
* Op-level timeouts tolerate degraded networks where TLS handshakes and DNS
|
||||
* resolution can take several seconds.
|
||||
*/
|
||||
export const AUTH_TIMEOUT_MS = 25_000;
|
||||
export const HISTORY_TIMEOUT_MS = 25_000;
|
||||
export const PUBLISH_TIMEOUT_MS = 25_000;
|
||||
|
||||
/**
|
||||
* A stability-gated reset prevents reconnect flapping from erasing backoff.
|
||||
*/
|
||||
export const BACKOFF_RESET_STABLE_MS = 60_000;
|
||||
|
||||
/** Passive liveness thresholds for the relay heartbeat stream. */
|
||||
export const STALL_CHECK_INTERVAL_MS = 10_000;
|
||||
export const STALL_IDLE_TIMEOUT_MS = 60_000;
|
||||
Reference in New Issue
Block a user