From 54ee9d216dd03d5c27d3e0a3625931de2a7e1b32 Mon Sep 17 00:00:00 2001 From: Tommy Falkowski <86825018+WismutHansen@users.noreply.github.com> Date: Sun, 6 Jul 2025 10:23:41 +0000 Subject: [PATCH] Add pause/resume WebSocket control --- README.md | 10 +++++++++- src/lib.rs | 37 ++++++++++++++++++++++++++++++++++--- websocket_example.html | 25 +++++++++++++++++++++++++ 3 files changed, 68 insertions(+), 4 deletions(-) diff --git a/README.md b/README.md index 304aa4a1..dce7f61b 100644 --- a/README.md +++ b/README.md @@ -142,13 +142,21 @@ Send from client to restart transcription after timeout or final message: } ``` +#### Pause/Resume Transcription +Toggle live inference without disconnecting: +```json +{ "type": "pause" } +{ "type": "resume" } +``` + ### Usage Pattern 1. Connect to WebSocket endpoint 2. Receive real-time word messages during transcription 3. Receive final message when session ends (timeout or silence) 4. Send restart command to begin new transcription session -5. Repeat as needed +5. Optionally send pause/resume commands to temporarily stop inference +6. Repeat as needed ## Model diff --git a/src/lib.rs b/src/lib.rs index d91a6dd1..c7456e38 100644 --- a/src/lib.rs +++ b/src/lib.rs @@ -175,6 +175,10 @@ pub enum WebSocketMessage { pub enum WebSocketCommand { #[serde(rename = "restart")] Restart, + #[serde(rename = "pause")] + Pause, + #[serde(rename = "resume")] + Resume, } pub struct TranscriptionOptions { @@ -414,7 +418,7 @@ impl Model { use futures::{SinkExt, StreamExt}; use std::io::Write; use std::sync::Arc; - use tokio::sync::{broadcast, mpsc}; + use tokio::sync::{broadcast, mpsc, watch}; use tokio_tungstenite::{accept_async, tungstenite::Message}; // WebSocket broadcast channel @@ -425,14 +429,20 @@ impl Model { let (restart_tx, mut restart_rx) = mpsc::unbounded_channel(); let restart_tx = Arc::new(restart_tx); + // Watch channel used to pause or resume transcription + let (pause_tx, _pause_rx) = watch::channel(false); + let pause_tx = Arc::new(pause_tx); + // Spawn WebSocket server let listener = tokio::net::TcpListener::bind(format!("127.0.0.1:{}", ws_port)).await?; let ws_tx_clone = ws_tx.clone(); let restart_tx_clone = restart_tx.clone(); + let pause_tx_clone = pause_tx.clone(); tokio::spawn(async move { while let Ok((stream, _)) = listener.accept().await { let ws_tx = ws_tx_clone.clone(); let restart_tx = restart_tx_clone.clone(); + let pause_tx = pause_tx_clone.clone(); tokio::spawn(async move { let ws_stream = match accept_async(stream).await { Ok(ws) => ws, @@ -452,8 +462,16 @@ impl Model { Ok(Message::Text(text)) => { if let Ok(cmd) = serde_json::from_str::(&text) { - if let WebSocketCommand::Restart = cmd { - let _ = restart_tx.send(()); + match cmd { + WebSocketCommand::Restart => { + let _ = restart_tx.send(()); + } + WebSocketCommand::Pause => { + let _ = pause_tx.send(true); + } + WebSocketCommand::Resume => { + let _ = pause_tx.send(false); + } } } } @@ -482,6 +500,7 @@ impl Model { // Bridge blocking audio receiver to async channel let (pcm_tx, mut pcm_rx) = mpsc::unbounded_channel(); + let mut pause_rx = pause_tx.subscribe(); std::thread::spawn(move || { while let Ok(chunk) = audio_rx.recv() { if pcm_tx.send(chunk).is_err() { @@ -501,6 +520,7 @@ impl Model { let mut printed_eot = false; let mut last_voice_activity: Option = None; let mut restart = false; + let mut paused = *pause_rx.borrow(); eprintln!("Starting transcription session..."); @@ -511,7 +531,18 @@ impl Model { restart = true; break; } + _ = pause_rx.changed() => { + paused = *pause_rx.borrow(); + if paused { + eprintln!("Transcription paused"); + } else { + eprintln!("Transcription resumed"); + } + } Some(pcm_chunk) = pcm_rx.recv() => { + if paused { + continue; + } if save_audio.is_some() { all_audio.extend_from_slice(&pcm_chunk); } diff --git a/websocket_example.html b/websocket_example.html index ce3c136b..cb55e7eb 100644 --- a/websocket_example.html +++ b/websocket_example.html @@ -37,6 +37,7 @@

eaRS WebSocket Transcription Client

+
Status: Disconnected
@@ -54,8 +55,10 @@

Instructions: