diff --git a/src-tauri/src/engine.rs b/src-tauri/src/engine.rs index 3e8a0a1..c59d6c3 100644 --- a/src-tauri/src/engine.rs +++ b/src-tauri/src/engine.rs @@ -256,27 +256,46 @@ impl Engine { } } - // TODO(#63): split legacy recording startup into capped helpers. - #[allow(clippy::too_many_lines)] pub fn start(self: &Arc, app: &AppHandle) { - let support = platform::current(); - if let Some(message) = support.unsupported_dictation_message() { - if AppConfig::load().general.sounds { - sounds::play(sounds::Cue::Error); - } - self.set_state(app, |s| { - s.stage = Stage::Idle; - s.recording_started_ms = None; - s.segments.clear(); - s.error = Some(message); - s.message = None; - }); + if self.reject_unsupported_platform(app) { return; } let cfg = AppConfig::load(); - let recording = match recorder::start(&cfg.stt) { - Ok(rec) => rec, + let Some(recording) = self.start_recorder(app, &cfg) else { + return; + }; + if cfg.general.sounds { + sounds::play(sounds::Cue::Start); + } + + let (audio_path, session_id) = self.begin_recording_session(app, &cfg, recording); + + self.spawn_recorder_warmup_check(app, session_id); + self.spawn_level_meter(app, audio_path); + } + + fn reject_unsupported_platform(&self, app: &AppHandle) -> bool { + let support = platform::current(); + let Some(message) = support.unsupported_dictation_message() else { + return false; + }; + if AppConfig::load().general.sounds { + sounds::play(sounds::Cue::Error); + } + self.set_state(app, |s| { + s.stage = Stage::Idle; + s.recording_started_ms = None; + s.segments.clear(); + s.error = Some(message); + s.message = None; + }); + true + } + + fn start_recorder(&self, app: &AppHandle, cfg: &AppConfig) -> Option { + match recorder::start(&cfg.stt) { + Ok(rec) => Some(rec), Err(err) => { if cfg.general.sounds { sounds::play(sounds::Cue::Error); @@ -288,12 +307,17 @@ impl Engine { s.error = Some(format!("{err:#}")); s.message = None; }); - return; + None } - }; - if cfg.general.sounds { - sounds::play(sounds::Cue::Start); } + } + + fn begin_recording_session( + self: &Arc, + app: &AppHandle, + cfg: &AppConfig, + recording: recorder::Recording, + ) -> (PathBuf, String) { let audio_path = recording.audio_path.clone(); let session_id = format!("{}-{}", now_ms(), std::process::id()); let cancel_token = CancelToken::new(); @@ -301,7 +325,7 @@ impl Engine { let temp_dir = recorder::state_dir().join("incremental").join(&session_id); match self.prepare_incremental_worker( app, - &cfg, + cfg, audio_path.clone(), session_id.clone(), cancel_token.clone(), @@ -363,41 +387,45 @@ impl Engine { }); } - // Recorder warm-up check, off the command thread so toggling stays - // responsive: if the recorder exited immediately, surface the error - // and tear the session down. - { - let engine = Arc::clone(self); - let check_app = app.clone(); - std::thread::spawn(move || { - std::thread::sleep(Duration::from_millis(250)); - let err = { - let mut guard = engine.recording.lock().unwrap(); - match guard.as_mut() { - Some(active) if active.session_id == session_id => { - active.recording.exit_error() - } - _ => None, + (audio_path, session_id) + } + + // Recorder warm-up check, off the command thread so toggling stays + // responsive: if the recorder exited immediately, surface the error + // and tear the session down. + fn spawn_recorder_warmup_check(self: &Arc, app: &AppHandle, session_id: String) { + let engine = Arc::clone(self); + let check_app = app.clone(); + std::thread::spawn(move || { + std::thread::sleep(Duration::from_millis(250)); + let err = { + let mut guard = engine.recording.lock().unwrap(); + match guard.as_mut() { + Some(active) if active.session_id == session_id => { + active.recording.exit_error() } - }; - let Some(err) = err else { - return; - }; - engine.cancel(&check_app); - if AppConfig::load().general.sounds { - sounds::play(sounds::Cue::Error); + _ => None, } - engine.set_state(&check_app, |s| { - s.stage = Stage::Idle; - s.recording_started_ms = None; - s.segments.clear(); - s.error = Some(err); - s.message = None; - }); + }; + let Some(err) = err else { + return; + }; + engine.cancel(&check_app); + if AppConfig::load().general.sounds { + sounds::play(sounds::Cue::Error); + } + engine.set_state(&check_app, |s| { + s.stage = Stage::Idle; + s.recording_started_ms = None; + s.segments.clear(); + s.error = Some(err); + s.message = None; }); - } + }); + } - // Live level meter for the waveform. + // Live level meter for the waveform. + fn spawn_level_meter(self: &Arc, app: &AppHandle, audio_path: PathBuf) { self.levels_running.store(true, Ordering::SeqCst); let running = Arc::clone(&self.levels_running); let level_app = app.clone(); @@ -578,8 +606,6 @@ impl Engine { } } - // TODO(#63): split legacy pipeline orchestration into capped helpers. - #[allow(clippy::cognitive_complexity, clippy::too_many_lines)] fn run_pipeline(self: Arc, app: &AppHandle, cfg: AppConfig, active: ActiveRecording) { let ActiveRecording { session_id, @@ -588,47 +614,91 @@ impl Engine { incremental, } = active; - let fail = |err: String| { - if !self.is_session_current(&session_id, &cancel_token) { - return; - } - if cfg.general.sounds { - sounds::play(sounds::Cue::Error); - } - let (error, message) = lifecycle::failure_outcome(err); + let Some((duration_ms, raw)) = + self.acquire_transcript(app, &cfg, &session_id, &cancel_token, incremental, recording) + else { + return; + }; + + if raw.is_empty() { + let (error, message) = lifecycle::no_speech_outcome(); let _ = self.finish_session(app, &session_id, &cancel_token, |s| { s.stage = Stage::Idle; s.recording_started_ms = None; s.segments.clear(); - s.error = error; s.message = message; + s.error = error; }); - }; + return; + } + + self.clean_and_deliver(app, &cfg, &session_id, &cancel_token, duration_ms, raw); + } + + fn fail_session( + &self, + app: &AppHandle, + cfg: &AppConfig, + session_id: &str, + cancel_token: &CancelToken, + err: String, + ) { + if !self.is_session_current(session_id, cancel_token) { + return; + } + if cfg.general.sounds { + sounds::play(sounds::Cue::Error); + } + let (error, message) = lifecycle::failure_outcome(err); + let _ = self.finish_session(app, session_id, cancel_token, |s| { + s.stage = Stage::Idle; + s.recording_started_ms = None; + s.segments.clear(); + s.error = error; + s.message = message; + }); + } + /// Stop the recorder and resolve the final transcript text, handling the + /// incremental-result/fallback-tail dispatch and every stale-session + /// guard along the way. Returns `None` once the session has already been + /// failed or torn down and the caller should just return. + fn acquire_transcript( + &self, + app: &AppHandle, + cfg: &AppConfig, + session_id: &str, + cancel_token: &CancelToken, + incremental: Option, + recording: recorder::Recording, + ) -> Option<(u64, String)> { let (audio_path, duration_ms) = match recording.stop() { Ok(result) => result, - Err(err) => return fail(format!("{err:#}")), + Err(err) => { + self.fail_session(app, cfg, session_id, cancel_token, format!("{err:#}")); + return None; + } }; - if !self.is_session_current(&session_id, &cancel_token) { + if !self.is_session_current(session_id, cancel_token) { if !cfg.general.keep_audio { let _ = fs::remove_file(&audio_path); } - return; + return None; } - let raw = match self.incremental_result(&cfg, incremental, &session_id, &cancel_token) { + let raw = match self.incremental_result(cfg, incremental, session_id, cancel_token) { IncrementalResult::Complete(text) => Ok(text), outcome => { - if !self.is_session_current(&session_id, &cancel_token) { + if !self.is_session_current(session_id, cancel_token) { if !cfg.general.keep_audio { let _ = fs::remove_file(&audio_path); } - return; + return None; } - let is_cancelled = || !self.is_session_current(&session_id, &cancel_token); + let is_cancelled = || !self.is_session_current(session_id, cancel_token); match outcome { IncrementalResult::Partial(prefix) => { - match slice_fallback_tail(&cfg, &audio_path, &prefix) { + match slice_fallback_tail(cfg, &audio_path, &prefix) { // Recording ended inside the already-transcribed // prefix — nothing left to transcribe. Ok(None) => Ok(prefix.raw_text), @@ -657,38 +727,39 @@ impl Engine { if !cfg.general.keep_audio { let _ = fs::remove_file(&audio_path); } - return fail(format!("{err:#}")); + self.fail_session(app, cfg, session_id, cancel_token, format!("{err:#}")); + return None; } }; if !cfg.general.keep_audio { let _ = fs::remove_file(&audio_path); } - if !self.is_session_current(&session_id, &cancel_token) { - return; - } - if raw.is_empty() { - let (error, message) = lifecycle::no_speech_outcome(); - let _ = self.finish_session(app, &session_id, &cancel_token, |s| { - s.stage = Stage::Idle; - s.recording_started_ms = None; - s.segments.clear(); - s.message = message; - s.error = error; - }); - return; + if !self.is_session_current(session_id, cancel_token) { + return None; } + Some((duration_ms, raw)) + } - if !self.set_state_for_session(app, &session_id, &cancel_token, |s| { + fn clean_and_deliver( + &self, + app: &AppHandle, + cfg: &AppConfig, + session_id: &str, + cancel_token: &CancelToken, + duration_ms: u64, + raw: String, + ) { + if !self.set_state_for_session(app, session_id, cancel_token, |s| { s.stage = Stage::Cleaning }) { return; } - let outcome = cleanup::clean(&cfg, &raw); - if !self.is_session_current(&session_id, &cancel_token) { + let outcome = cleanup::clean(cfg, &raw); + if !self.is_session_current(session_id, cancel_token) { return; } - if !self.set_state_for_session(app, &session_id, &cancel_token, |s| { + if !self.set_state_for_session(app, session_id, cancel_token, |s| { s.stage = Stage::Pasting }) { return; @@ -702,7 +773,7 @@ impl Engine { .into_result() .err() .map(|err| format!("paste failed (text copied if possible): {err:#}")); - if !self.is_session_current(&session_id, &cancel_token) { + if !self.is_session_current(session_id, cancel_token) { return; } @@ -741,7 +812,7 @@ impl Engine { }; let delivered = lifecycle::delivery_outcome(paste_error, &outcome, last_entry); - let _ = self.finish_session(app, &session_id, &cancel_token, |s| { + let _ = self.finish_session(app, session_id, cancel_token, |s| { s.stage = Stage::Idle; s.recording_started_ms = None; s.segments.clear(); diff --git a/src-tauri/src/file_job.rs b/src-tauri/src/file_job.rs index 02e7d18..c64a577 100644 --- a/src-tauri/src/file_job.rs +++ b/src-tauri/src/file_job.rs @@ -123,8 +123,6 @@ fn validate_input_path(path: &Path) -> Result<()> { Ok(()) } -// TODO(#63): split legacy file-job orchestration into capped helpers. -#[allow(clippy::too_many_lines)] fn run_file_job( engine: Arc, app: AppHandle, @@ -138,97 +136,15 @@ fn run_file_job( let dir = create_temp_dir()?; guard.set_temp_dir(dir.clone()); let wav = dir.join("audio.wav"); - let cfg = AppConfig::load(); - - if cancel_token.is_cancelled() { - bail!("file transcription cancelled"); - } - emit_state(&app, FileStage::Converting, 0, &source_file, None, None); - media::convert_to_wav_16k_mono(Path::new(&source_file), &wav, || { - cancel_token.is_cancelled() - })?; - if cancel_token.is_cancelled() { - bail!("file transcription cancelled"); - } - - let duration_ms = media::wav_duration_ms(&wav)?; - emit_state(&app, FileStage::Transcribing, 0, &source_file, None, None); - let progress_for_callback = Arc::clone(&progress); - let app_for_callback = app.clone(); - let source_for_callback = source_file.clone(); - let segments = stt::transcribe_file_with_cancel( - &cfg.stt, + execute_file_job( + &engine, + &app, &wav, - || cancel_token.is_cancelled(), - move |percentage| { - if progress_for_callback.swap(percentage, Ordering::Relaxed) != percentage { - emit_state( - &app_for_callback, - FileStage::Transcribing, - percentage, - &source_for_callback, - None, - None, - ); - } - }, - )?; - if cancel_token.is_cancelled() { - bail!("file transcription cancelled"); - } - - let raw_text = transcript::to_txt(&segments); - if raw_text.is_empty() { - bail!("no speech detected in the file"); - } - let (cleaned_text, provider, model, cleanup_error) = if cleanup_requested { - emit_state( - &app, - FileStage::Cleaning, - progress.load(Ordering::Relaxed), - &source_file, - None, - None, - ); - let outcome = cleanup::clean(&cfg, &raw_text); - if cancel_token.is_cancelled() { - bail!("file transcription cancelled"); - } - ( - (outcome.text != raw_text).then_some(outcome.text), - outcome.provider, - outcome.model, - outcome - .error - .map(|error| format!("cleanup skipped: {error}")), - ) - } else { - (None, "none".to_string(), String::new(), None) - }; - let segments_json = if segments.is_empty() { - None - } else { - Some(serde_json::to_string(&segments).context("serializing file segments")?) - }; - if cancel_token.is_cancelled() { - bail!("file transcription cancelled"); - } - - let entry_id = engine.history.insert(&NewEntry { - duration_ms, - raw_text, - cleaned_text, - provider, - model, - language: cfg.stt.language, - source_file: Some(source_file.clone()), - segments_json, - })?; - let _ = app.emit(EVENT_HISTORY, ()); - Ok(FileJobComplete { - entry_id, - cleanup_error, - }) + &source_file, + cleanup_requested, + &cancel_token, + &progress, + ) })(); match result { @@ -264,6 +180,138 @@ struct FileJobComplete { cleanup_error: Option, } +fn execute_file_job( + engine: &Engine, + app: &AppHandle, + wav: &Path, + source_file: &str, + cleanup_requested: bool, + cancel_token: &CancelToken, + progress: &Arc, +) -> Result { + let cfg = AppConfig::load(); + let (duration_ms, segments) = + convert_and_transcribe(&cfg, wav, source_file, cancel_token, app, progress)?; + + let raw_text = transcript::to_txt(&segments); + if raw_text.is_empty() { + bail!("no speech detected in the file"); + } + let (cleaned_text, provider, model, cleanup_error) = clean_transcript_if_requested( + &cfg, + &raw_text, + cleanup_requested, + cancel_token, + app, + source_file, + progress, + )?; + let segments_json = if segments.is_empty() { + None + } else { + Some(serde_json::to_string(&segments).context("serializing file segments")?) + }; + if cancel_token.is_cancelled() { + bail!("file transcription cancelled"); + } + + let entry_id = engine.history.insert(&NewEntry { + duration_ms, + raw_text, + cleaned_text, + provider, + model, + language: cfg.stt.language, + source_file: Some(source_file.to_string()), + segments_json, + })?; + let _ = app.emit(EVENT_HISTORY, ()); + Ok(FileJobComplete { + entry_id, + cleanup_error, + }) +} + +fn convert_and_transcribe( + cfg: &AppConfig, + wav: &Path, + source_file: &str, + cancel_token: &CancelToken, + app: &AppHandle, + progress: &Arc, +) -> Result<(i64, Vec)> { + if cancel_token.is_cancelled() { + bail!("file transcription cancelled"); + } + emit_state(app, FileStage::Converting, 0, source_file, None, None); + media::convert_to_wav_16k_mono(Path::new(source_file), wav, || cancel_token.is_cancelled())?; + if cancel_token.is_cancelled() { + bail!("file transcription cancelled"); + } + + let duration_ms = media::wav_duration_ms(wav)?; + emit_state(app, FileStage::Transcribing, 0, source_file, None, None); + let progress_for_callback = Arc::clone(progress); + let app_for_callback = app.clone(); + let source_for_callback = source_file.to_string(); + let segments = stt::transcribe_file_with_cancel( + &cfg.stt, + wav, + || cancel_token.is_cancelled(), + move |percentage| { + if progress_for_callback.swap(percentage, Ordering::Relaxed) != percentage { + emit_state( + &app_for_callback, + FileStage::Transcribing, + percentage, + &source_for_callback, + None, + None, + ); + } + }, + )?; + if cancel_token.is_cancelled() { + bail!("file transcription cancelled"); + } + Ok((duration_ms, segments)) +} + +fn clean_transcript_if_requested( + cfg: &AppConfig, + raw_text: &str, + cleanup_requested: bool, + cancel_token: &CancelToken, + app: &AppHandle, + source_file: &str, + progress: &Arc, +) -> Result<(Option, String, String, Option)> { + if cleanup_requested { + emit_state( + app, + FileStage::Cleaning, + progress.load(Ordering::Relaxed), + source_file, + None, + None, + ); + let outcome = cleanup::clean(cfg, raw_text); + if cancel_token.is_cancelled() { + bail!("file transcription cancelled"); + } + Ok(( + (outcome.text != raw_text).then_some(outcome.text), + outcome.provider, + outcome.model, + outcome + .error + .map(|error| format!("cleanup skipped: {error}")), + )) + } else { + Ok((None, "none".to_string(), String::new(), None)) + } +} + fn create_temp_dir() -> Result { let parent = recorder::prepare_state_dir()?; let stamp = SystemTime::now() diff --git a/src-tauri/src/lib.rs b/src-tauri/src/lib.rs index f0d0af2..96b5004 100644 --- a/src-tauri/src/lib.rs +++ b/src-tauri/src/lib.rs @@ -849,10 +849,9 @@ fn clamp_float_window_size(window: &tauri::WebviewWindow) { )); } -// TODO(#63): split legacy application setup into capped helpers. -#[allow(clippy::too_many_lines)] -pub fn run() { - let context = tauri::generate_context!(); +fn init_sentry( + context: &tauri::Context, +) -> (sentry::ClientInitGuard, tauri::plugin::TauriPlugin) { let cfg = AppConfig::load(); let sentry_enabled = settings::sentry_client_enabled(&cfg); let release = format!( @@ -890,6 +889,57 @@ pub fn run() { } else { tauri_plugin_sentry::init_with_no_injection(&sentry_client) }; + (sentry_client, sentry_plugin) +} + +fn setup_app(app: &mut tauri::App) -> Result<(), Box> { + tray::setup(app)?; + let cfg = AppConfig::load(); + ensure_float_window(app.handle(), cfg.general.float_button); + if let Err(err) = shortcut::register_startup(app.handle(), &cfg.shortcut.toggle) { + eprintln!("{err}"); + } + #[cfg(target_os = "linux")] + if let Some(window) = app.get_webview_window("main") { + fix_csd_titlebar_input(&window); + } + let args: Vec = std::env::args().collect(); + if args.iter().any(|a| a == "--hidden") + && let Some(window) = app.get_webview_window("main") + { + let _ = window.hide(); + } + if args.iter().any(|a| a == "--toggle") { + let handle = app.handle().clone(); + let engine = engine::engine(&handle); + engine.set_chord_override(parse_chord_arg(&args)); + engine.toggle(&handle); + } + Ok(()) +} + +fn handle_run_event(app_handle: &tauri::AppHandle, event: tauri::RunEvent) { + match event { + tauri::RunEvent::WindowEvent { + label, + event: WindowEvent::CloseRequested { api, .. }, + .. + } if label == "main" => { + api.prevent_close(); + if let Some(window) = app_handle.get_webview_window("main") { + let _ = window.hide(); + } + } + tauri::RunEvent::ExitRequested { code: None, api, .. } => { + api.prevent_exit(); + } + _ => {} + } +} + +pub fn run() { + let context = tauri::generate_context!(); + let (_sentry_client, sentry_plugin) = init_sentry(&context); let engine = Arc::new(Engine::new().expect("failed to open PickScribe data directory")); let builder = tauri::Builder::default() @@ -952,47 +1002,8 @@ pub fn run() { copy_text, show_main_window, ]) - .setup(|app| { - tray::setup(app)?; - let cfg = AppConfig::load(); - ensure_float_window(app.handle(), cfg.general.float_button); - if let Err(err) = shortcut::register_startup(app.handle(), &cfg.shortcut.toggle) { - eprintln!("{err}"); - } - #[cfg(target_os = "linux")] - if let Some(window) = app.get_webview_window("main") { - fix_csd_titlebar_input(&window); - } - let args: Vec = std::env::args().collect(); - if args.iter().any(|a| a == "--hidden") - && let Some(window) = app.get_webview_window("main") - { - let _ = window.hide(); - } - if args.iter().any(|a| a == "--toggle") { - let handle = app.handle().clone(); - let engine = engine::engine(&handle); - engine.set_chord_override(parse_chord_arg(&args)); - engine.toggle(&handle); - } - Ok(()) - }) + .setup(setup_app) .build(context) .expect("error while building PickScribe") - .run(|app_handle, event| match event { - tauri::RunEvent::WindowEvent { - label, - event: WindowEvent::CloseRequested { api, .. }, - .. - } if label == "main" => { - api.prevent_close(); - if let Some(window) = app_handle.get_webview_window("main") { - let _ = window.hide(); - } - } - tauri::RunEvent::ExitRequested { code: None, api, .. } => { - api.prevent_exit(); - } - _ => {} - }); + .run(handle_run_event); } diff --git a/src/bin/pickscribe.rs b/src/bin/pickscribe.rs index fb19ac4..1e4d358 100644 --- a/src/bin/pickscribe.rs +++ b/src/bin/pickscribe.rs @@ -1087,8 +1087,6 @@ fn transcribe(args: &Args, audio_path: &Path) -> Result { } } -// TODO(#63): extract the legacy whisper process orchestration below the cap. -#[allow(clippy::too_many_lines)] fn transcribe_incremental_segment( args: &Args, audio_path: &Path, @@ -1101,62 +1099,87 @@ fn transcribe_incremental_segment( .as_deref() .filter(|value| !value.is_empty()) { - let model = args - .whisper_model - .as_ref() - .map(|path| path.display().to_string()) - .unwrap_or_default(); - let command = custom - .replace("{audio}", &shell_escape(&audio_path.display().to_string())) - .replace("{model}", &shell_escape(&model)) - .replace( - "{output}", - &shell_escape(&output_prefix.display().to_string()), - ); - let stdout_path = output_prefix.with_extension("transcript.stdout.log"); - let stderr_path = output_prefix.with_extension("transcript.stderr.log"); - let stdout_file = File::create(&stdout_path) - .with_context(|| format!("failed to create {}", stdout_path.display()))?; - let stderr_file = File::create(&stderr_path) - .with_context(|| format!("failed to create {}", stderr_path.display()))?; - let (mut cmd, process_group) = cancellable_shell_command(); - cmd.arg("-lc") - .arg(&command) - .stdout(Stdio::from(stdout_file)) - .stderr(Stdio::from(stderr_file)); - let mut child = cmd - .spawn() - .with_context(|| format!("failed to run custom STT command: {command}"))?; - let status = match wait_for_cancellable_child(&mut child, process_group, is_cancelled) { - Ok(status) => status, - Err(err) => { - cleanup_segment_transcript_files(audio_path); - return Err(err); - } - }; - if !status.success() { - let stderr = fs::read_to_string(&stderr_path).unwrap_or_default(); + return transcribe_with_custom_command( + args, + audio_path, + &output_prefix, + custom, + is_cancelled, + ); + } + + transcribe_with_whisper_cpp(args, audio_path, &output_prefix, is_cancelled) +} + +fn transcribe_with_custom_command( + args: &Args, + audio_path: &Path, + output_prefix: &Path, + custom: &str, + is_cancelled: impl Fn() -> bool, +) -> Result { + let model = args + .whisper_model + .as_ref() + .map(|path| path.display().to_string()) + .unwrap_or_default(); + let command = custom + .replace("{audio}", &shell_escape(&audio_path.display().to_string())) + .replace("{model}", &shell_escape(&model)) + .replace( + "{output}", + &shell_escape(&output_prefix.display().to_string()), + ); + let stdout_path = output_prefix.with_extension("transcript.stdout.log"); + let stderr_path = output_prefix.with_extension("transcript.stderr.log"); + let stdout_file = File::create(&stdout_path) + .with_context(|| format!("failed to create {}", stdout_path.display()))?; + let stderr_file = File::create(&stderr_path) + .with_context(|| format!("failed to create {}", stderr_path.display()))?; + let (mut cmd, process_group) = cancellable_shell_command(); + cmd.arg("-lc") + .arg(&command) + .stdout(Stdio::from(stdout_file)) + .stderr(Stdio::from(stderr_file)); + let mut child = cmd + .spawn() + .with_context(|| format!("failed to run custom STT command: {command}"))?; + let status = match wait_for_cancellable_child(&mut child, process_group, is_cancelled) { + Ok(status) => status, + Err(err) => { cleanup_segment_transcript_files(audio_path); - bail!( - "custom STT command failed with {}:\n{}", - status, - stderr.trim() - ); + return Err(err); } - - let txt_path = transcript_txt_path_for(audio_path); - let output = if txt_path.exists() { - fs::read_to_string(&txt_path) - .with_context(|| format!("failed to read {}", txt_path.display()))? - } else { - fs::read_to_string(&stdout_path) - .with_context(|| format!("failed to read {}", stdout_path.display()))? - }; - let _ = fs::remove_file(&stdout_path); - let _ = fs::remove_file(&stderr_path); - return Ok(output); + }; + if !status.success() { + let stderr = fs::read_to_string(&stderr_path).unwrap_or_default(); + cleanup_segment_transcript_files(audio_path); + bail!( + "custom STT command failed with {}:\n{}", + status, + stderr.trim() + ); } + let txt_path = transcript_txt_path_for(audio_path); + let output = if txt_path.exists() { + fs::read_to_string(&txt_path) + .with_context(|| format!("failed to read {}", txt_path.display()))? + } else { + fs::read_to_string(&stdout_path) + .with_context(|| format!("failed to read {}", stdout_path.display()))? + }; + let _ = fs::remove_file(&stdout_path); + let _ = fs::remove_file(&stderr_path); + Ok(output) +} + +fn transcribe_with_whisper_cpp( + args: &Args, + audio_path: &Path, + output_prefix: &Path, + is_cancelled: impl Fn() -> bool, +) -> Result { let whisper = resolve_whisper_command(args)?; let stderr_path = output_prefix.with_extension("transcript.stderr.log"); let stderr_file = File::create(&stderr_path) @@ -1171,7 +1194,7 @@ fn transcribe_incremental_segment( .arg(audio_path) .arg("--output-txt") .arg("--output-file") - .arg(&output_prefix) + .arg(output_prefix) .arg("--no-prints"); if let Some(language) = args.language.as_deref().filter(|value| !value.is_empty()) { diff --git a/src/engine/incremental.rs b/src/engine/incremental.rs index 5f88846..26d4ef1 100644 --- a/src/engine/incremental.rs +++ b/src/engine/incremental.rs @@ -367,8 +367,6 @@ pub enum RunResult { /// progress state (the `RecordingSession` and when to publish it), the /// fallback/completion decision, and the final drain. STT execution, live /// segment cleanup, and progress transport are the host's job. -// TODO(#63): split the legacy session driver into capped orchestration helpers. -#[allow(clippy::too_many_lines)] pub fn run( host: &mut impl IncrementalHost, audio_path: &Path, @@ -379,138 +377,24 @@ pub fn run( let mut next_start_ms = 0u64; let mut segment_id = 0u64; - let stop = 'run: loop { + let stop = loop { if drain_and_publish(host, &mut session) { // published below alongside other state changes too; draining // alone is enough reason to republish. } - match host.control() { - Control::Abandoned => { - host.cleanup_artifacts(); - return RunResult::Abandoned; - } - Control::Cancelled => break 'run StopOutcome::cancelled(), - control => { - let final_requested = matches!(control, Control::Stopping); - let available_ms = audio_segments::duration_ms(audio_path).unwrap_or(0); - - let (start_ms, desired_end_ms) = - match next_step(next_start_ms, available_ms, final_requested, &cfg) { - Step::Wait => { - std::thread::sleep(POLL_INTERVAL); - continue 'run; - } - Step::Stop(stop) => break 'run stop, - Step::Produce { - start_ms, - desired_end_ms, - } => (start_ms, desired_end_ms), - }; - - let end_ms = - choose_segment_end(audio_path, start_ms, desired_end_ms, available_ms, final_requested); - match after_boundary(start_ms, end_ms, available_ms, final_requested) { - BoundaryStep::Wait => { - std::thread::sleep(POLL_INTERVAL); - continue 'run; - } - BoundaryStep::Stop(stop) => break 'run stop, - BoundaryStep::Proceed => {} - } - - segment_id = segment_id.saturating_add(1); - let slice_start_ms = start_ms.saturating_sub(cfg.overlap_ms); - let segment_path = host.segment_path(segment_id); - let slice_result = - audio_segments::slice_wav(audio_path, &segment_path, slice_start_ms, end_ms); - - let slice = match classify_slice(slice_result, final_requested, available_ms, start_ms) { - SliceStep::Wait => { - std::thread::sleep(POLL_INTERVAL); - continue 'run; - } - SliceStep::Stop(stop) => break 'run stop, - SliceStep::Failed(err) => { - session.upsert_segment(TranscriptSegment::failed( - segment_id, - slice_start_ms, - end_ms, - err, - )); - host.publish(&session); - break 'run StopOutcome::fallback(); - } - SliceStep::Use(slice) => slice, - }; - - session.upsert_segment(TranscriptSegment { - id: segment_id, - start_ms: slice.start_ms, - end_ms: slice.end_ms, - status: TranscriptSegmentStatus::Transcribing, - raw_text: String::new(), - cleaned_text: None, - error: None, - }); - host.publish(&session); - - let job = SegmentJob { - segment_id, - audio_path: segment_path.clone(), - start_ms: slice.start_ms, - end_ms: slice.end_ms, - }; - let result = host.transcribe(&job); - if !host.keep_audio() { - let _ = fs::remove_file(&segment_path); - } - - match host.control() { - Control::Abandoned => { - host.cleanup_artifacts(); - return RunResult::Abandoned; - } - Control::Cancelled => break 'run StopOutcome::cancelled(), - _ => {} - } - - match result { - Ok(text) => { - let raw = TranscriptSegment::raw_ready( - segment_id, - slice.start_ms, - slice.end_ms, - text, - ); - session.upsert_segment(raw.clone()); - host.publish(&session); - - if !raw.raw_text.trim().is_empty() && host.try_queue_cleanup(raw.clone()) { - session.upsert_segment(TranscriptSegment { - status: TranscriptSegmentStatus::Cleaning, - ..raw - }); - host.publish(&session); - } - } - Err(err) => { - session.upsert_segment(TranscriptSegment::failed( - segment_id, - slice.start_ms, - slice.end_ms, - format!("{err:#}"), - )); - host.publish(&session); - break 'run StopOutcome::fallback(); - } - } - - next_start_ms = slice.end_ms; - if final_requested && next_start_ms >= available_ms { - break 'run StopOutcome::finished(); - } - } + match advance( + host, + audio_path, + &cfg, + &mut session, + &mut next_start_ms, + &mut segment_id, + ) { + Progress::Abandoned => return RunResult::Abandoned, + Progress::Stop(stop) => break stop, + Progress::Wait => std::thread::sleep(POLL_INTERVAL), + Progress::Advanced => {} } }; @@ -523,6 +407,210 @@ pub fn run( }) } +/// One tick of the driver loop: check the host's control signal, decide the +/// next segment's bounds, and (if any) produce and dispatch it. +enum Progress { + /// Not enough new audio yet, or a transient block; sleep and re-check. + Wait, + /// A segment was handled; loop again immediately. + Advanced, + /// The host was abandoned mid-tick; `run` must return without publishing. + Abandoned, + /// Terminal outcome reached. + Stop(StopOutcome), +} + +fn advance( + host: &mut impl IncrementalHost, + audio_path: &Path, + cfg: &SchedulingConfig, + session: &mut RecordingSession, + next_start_ms: &mut u64, + segment_id: &mut u64, +) -> Progress { + match host.control() { + Control::Abandoned => { + host.cleanup_artifacts(); + Progress::Abandoned + } + Control::Cancelled => Progress::Stop(StopOutcome::cancelled()), + control => { + let bounds = match compute_segment_bounds(audio_path, cfg, *next_start_ms, control) { + BoundsOutcome::Wait => return Progress::Wait, + BoundsOutcome::Stop(stop) => return Progress::Stop(stop), + BoundsOutcome::Ready(bounds) => bounds, + }; + + produce_segment( + host, + audio_path, + cfg, + session, + segment_id, + next_start_ms, + &bounds, + ) + } + } +} + +/// A segment's resolved cut points, and the request context they were +/// resolved under. +struct Bounds { + start_ms: u64, + end_ms: u64, + available_ms: u64, + final_requested: bool, +} + +enum BoundsOutcome { + Wait, + Stop(StopOutcome), + Ready(Bounds), +} + +fn compute_segment_bounds( + audio_path: &Path, + cfg: &SchedulingConfig, + next_start_ms: u64, + control: Control, +) -> BoundsOutcome { + let final_requested = matches!(control, Control::Stopping); + let available_ms = audio_segments::duration_ms(audio_path).unwrap_or(0); + + let (start_ms, desired_end_ms) = + match next_step(next_start_ms, available_ms, final_requested, cfg) { + Step::Wait => return BoundsOutcome::Wait, + Step::Stop(stop) => return BoundsOutcome::Stop(stop), + Step::Produce { + start_ms, + desired_end_ms, + } => (start_ms, desired_end_ms), + }; + + let end_ms = choose_segment_end( + audio_path, + start_ms, + desired_end_ms, + available_ms, + final_requested, + ); + match after_boundary(start_ms, end_ms, available_ms, final_requested) { + BoundaryStep::Wait => BoundsOutcome::Wait, + BoundaryStep::Stop(stop) => BoundsOutcome::Stop(stop), + BoundaryStep::Proceed => BoundsOutcome::Ready(Bounds { + start_ms, + end_ms, + available_ms, + final_requested, + }), + } +} + +/// Slice, dispatch, and record the outcome of one segment; advances +/// `next_start_ms` (or reaches a terminal [`Progress::Stop`]) on success. +fn produce_segment( + host: &mut impl IncrementalHost, + audio_path: &Path, + cfg: &SchedulingConfig, + session: &mut RecordingSession, + segment_id: &mut u64, + next_start_ms: &mut u64, + bounds: &Bounds, +) -> Progress { + *segment_id = segment_id.saturating_add(1); + let segment_id = *segment_id; + let slice_start_ms = bounds.start_ms.saturating_sub(cfg.overlap_ms); + let segment_path = host.segment_path(segment_id); + let slice_result = + audio_segments::slice_wav(audio_path, &segment_path, slice_start_ms, bounds.end_ms); + + let slice = match classify_slice( + slice_result, + bounds.final_requested, + bounds.available_ms, + bounds.start_ms, + ) { + SliceStep::Wait => return Progress::Wait, + SliceStep::Stop(stop) => return Progress::Stop(stop), + SliceStep::Failed(err) => { + session.upsert_segment(TranscriptSegment::failed( + segment_id, + slice_start_ms, + bounds.end_ms, + err, + )); + host.publish(session); + return Progress::Stop(StopOutcome::fallback()); + } + SliceStep::Use(slice) => slice, + }; + + session.upsert_segment(TranscriptSegment { + id: segment_id, + start_ms: slice.start_ms, + end_ms: slice.end_ms, + status: TranscriptSegmentStatus::Transcribing, + raw_text: String::new(), + cleaned_text: None, + error: None, + }); + host.publish(session); + + let job = SegmentJob { + segment_id, + audio_path: segment_path.clone(), + start_ms: slice.start_ms, + end_ms: slice.end_ms, + }; + let result = host.transcribe(&job); + if !host.keep_audio() { + let _ = fs::remove_file(&segment_path); + } + + match host.control() { + Control::Abandoned => { + host.cleanup_artifacts(); + return Progress::Abandoned; + } + Control::Cancelled => return Progress::Stop(StopOutcome::cancelled()), + _ => {} + } + + match result { + Ok(text) => { + let raw = TranscriptSegment::raw_ready(segment_id, slice.start_ms, slice.end_ms, text); + session.upsert_segment(raw.clone()); + host.publish(session); + + if !raw.raw_text.trim().is_empty() && host.try_queue_cleanup(raw.clone()) { + session.upsert_segment(TranscriptSegment { + status: TranscriptSegmentStatus::Cleaning, + ..raw + }); + host.publish(session); + } + } + Err(err) => { + session.upsert_segment(TranscriptSegment::failed( + segment_id, + slice.start_ms, + slice.end_ms, + format!("{err:#}"), + )); + host.publish(session); + return Progress::Stop(StopOutcome::fallback()); + } + } + + *next_start_ms = slice.end_ms; + if bounds.final_requested && *next_start_ms >= bounds.available_ms { + Progress::Stop(StopOutcome::finished()) + } else { + Progress::Advanced + } +} + fn drain_and_publish(host: &mut impl IncrementalHost, session: &mut RecordingSession) -> bool { let drained = host.drain_cleanup(); if drained.is_empty() {