From 74c2a180eecb8795240cdb1cca8223bbc2eb4bea Mon Sep 17 00:00:00 2001 From: GodsQuantum <193283917+GodsQuantum@users.noreply.github.com> Date: Wed, 16 Sep 2026 11:10:07 +0200 Subject: [PATCH 1/6] feat: add transcription chain model --- src/api/providers.rs | 14 +-- src/api/quick.rs | 1 + src/domain.rs | 51 +++++++++ src/jobs.rs | 28 ++--- src/lib.rs | 1 + src/quick.rs | 34 +++--- src/quick/tests.rs | 1 + src/transcription_chain.rs | 217 +++++++++++++++++++++++++++++++++++++ src/workflows.rs | 22 ++-- 9 files changed, 323 insertions(+), 46 deletions(-) create mode 100644 src/transcription_chain.rs diff --git a/src/api/providers.rs b/src/api/providers.rs index 17465d9..d799c24 100644 --- a/src/api/providers.rs +++ b/src/api/providers.rs @@ -113,13 +113,13 @@ pub async fn delete( State(state): State, Path(id): Path, ) -> AppResult { - if state - .workflows - .read() - .await - .iter() - .any(|workflow| workflow.provider_id == id) - { + if state.workflows.read().await.iter().any(|workflow| { + workflow.provider_id == id + || workflow + .transcription_chain + .iter() + .any(|route| route.provider_id == id) + }) { return Err(AppError::Conflict( "provider is still used by a workflow".into(), )); diff --git a/src/api/quick.rs b/src/api/quick.rs index a98774f..e695ba7 100644 --- a/src/api/quick.rs +++ b/src/api/quick.rs @@ -176,6 +176,7 @@ pub async fn upload( .filter(|value| !value.trim().is_empty()) .ok_or_else(|| AppError::BadRequest("providerId is required".into()))?, model: model.unwrap_or_default(), + transcription_chain: Vec::new(), language, output_kind: parse_output_kind(output_kind)?, output_dir, diff --git a/src/domain.rs b/src/domain.rs index 6161e1b..efc2d12 100644 --- a/src/domain.rs +++ b/src/domain.rs @@ -49,6 +49,39 @@ impl From<&Provider> for ProviderView { } } +#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)] +#[serde(rename_all = "camelCase")] +pub struct TranscriptionRoute { + pub provider_id: String, + pub model: String, + #[serde(default, skip_serializing_if = "Option::is_none")] + pub fallback_after_seconds: Option, +} + +#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)] +#[serde(rename_all = "snake_case")] +pub enum TranscriptionAttemptOutcome { + Success, + Failed, + TimedOut, + Cancelled, +} + +#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)] +#[serde(rename_all = "camelCase")] +pub struct TranscriptionAttempt { + pub run: u32, + pub route_index: usize, + pub provider_id: String, + pub provider_name: String, + pub model: String, + pub started_at_ms: u128, + pub finished_at_ms: u128, + pub outcome: TranscriptionAttemptOutcome, + #[serde(default, skip_serializing_if = "Option::is_none")] + pub error: Option, +} + #[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)] #[serde(rename_all = "camelCase")] pub struct MarkdownOptions { @@ -82,9 +115,12 @@ pub struct Workflow { pub archive_dir: String, #[serde(default)] pub tags: Vec, + #[serde(default)] pub provider_id: String, #[serde(default)] pub model: String, + #[serde(default)] + pub transcription_chain: Vec, #[serde(default, skip_serializing_if = "Option::is_none")] pub language: Option, #[serde(default)] @@ -162,6 +198,16 @@ pub struct Job { #[serde(default, skip_serializing_if = "Option::is_none")] pub quick: Option, pub provider_id: String, + #[serde(default)] + pub transcription_chain: Vec, + #[serde(default)] + pub transcription_attempts: Vec, + #[serde(default, skip_serializing_if = "Option::is_none")] + pub used_provider_id: Option, + #[serde(default, skip_serializing_if = "Option::is_none")] + pub used_provider_name: Option, + #[serde(default, skip_serializing_if = "Option::is_none")] + pub used_model: Option, pub original_name: String, pub source_path: PathBuf, pub source_size: u64, @@ -212,12 +258,17 @@ mod tests { let workflow: Workflow = serde_json::from_str(workflow_json).unwrap(); assert_eq!(workflow.output_dir, None); assert!(workflow.tags.is_empty()); + assert!(workflow.transcription_chain.is_empty()); let job_json = r#"{"id":"j","workflowId":"w","providerId":"p","originalName":"a.m4a","sourcePath":"/w/a.m4a","sourceSize":1,"sourceMtimeNs":1,"model":"m","status":"done","attempts":1,"markdownPublished":true,"createdAtMs":1,"updatedAtMs":1}"#; let job: Job = serde_json::from_str(job_json).unwrap(); assert_eq!(job.kind, JobKind::Workflow); assert_eq!(job.workflow_id.as_deref(), Some("w")); assert!(job.quick.is_none()); + assert!(job.transcription_chain.is_empty()); + assert!(job.transcription_attempts.is_empty()); + assert!(job.used_provider_id.is_none()); + assert!(job.used_model.is_none()); } #[test] diff --git a/src/jobs.rs b/src/jobs.rs index 4bfaa31..a9022ad 100644 --- a/src/jobs.rs +++ b/src/jobs.rs @@ -53,19 +53,13 @@ pub async fn create_job( source_size: u64, source_mtime_ns: i128, ) -> Result { - let provider = state - .providers - .read() - .await - .iter() - .find(|provider| provider.id == workflow.provider_id) + let providers = state.providers.read().await.clone(); + let transcription_chain = + crate::transcription_chain::normalize_workflow_chain(workflow, &providers)?; + let primary = transcription_chain + .first() .cloned() - .ok_or_else(|| anyhow!("provider not found: {}", workflow.provider_id))?; - let model = if workflow.model.trim().is_empty() { - provider.model - } else { - workflow.model.clone() - }; + .ok_or_else(|| anyhow!("transcription chain is empty"))?; let original_name = source_path .file_name() .and_then(|v| v.to_str()) @@ -77,12 +71,17 @@ pub async fn create_job( kind: crate::domain::JobKind::Workflow, workflow_id: Some(workflow.id.clone()), quick: None, - provider_id: workflow.provider_id.clone(), + provider_id: primary.provider_id.clone(), + transcription_chain, + transcription_attempts: Vec::new(), + used_provider_id: None, + used_provider_name: None, + used_model: None, original_name, source_path, source_size, source_mtime_ns, - model, + model: primary.model, language: workflow.language.clone(), status: JobStatus::Pending, attempts: 1, @@ -255,6 +254,7 @@ mod tests { tags: Vec::new(), provider_id: "provider".into(), model: String::new(), + transcription_chain: Vec::new(), language: None, markdown: MarkdownOptions::default(), enabled: true, diff --git a/src/lib.rs b/src/lib.rs index df115ee..7ab4969 100644 --- a/src/lib.rs +++ b/src/lib.rs @@ -10,4 +10,5 @@ mod pipeline_quick; pub mod provider; pub mod quick; pub mod state; +pub mod transcription_chain; pub mod workflows; diff --git a/src/quick.rs b/src/quick.rs index 2c716ab..4430146 100644 --- a/src/quick.rs +++ b/src/quick.rs @@ -1,5 +1,7 @@ use crate::{ - domain::{Job, JobKind, JobStatus, QuickJobMeta, QuickOutputKind, QuickSourceKind}, + domain::{ + Job, JobKind, JobStatus, QuickJobMeta, QuickOutputKind, QuickSourceKind, TranscriptionRoute, + }, jobs::{now_ms, persist_job}, state::AppState, workflows::is_audio_candidate, @@ -16,10 +18,13 @@ use uuid::Uuid; #[derive(Debug, Clone, Serialize, Deserialize)] #[serde(rename_all = "camelCase")] pub struct QuickOptions { + #[serde(default)] pub provider_id: String, #[serde(default)] pub model: String, #[serde(default)] + pub transcription_chain: Vec, + #[serde(default)] pub language: Option, pub output_kind: QuickOutputKind, #[serde(default)] @@ -50,19 +55,13 @@ async fn normalized_job( if !is_audio_candidate(Path::new(&original_name)) { bail!("unsupported audio file: {original_name}"); } - let provider = state - .providers - .read() - .await - .iter() - .find(|provider| provider.id == options.provider_id) + let providers = state.providers.read().await.clone(); + let transcription_chain = + crate::transcription_chain::normalize_quick_chain(&options, &providers)?; + let primary = transcription_chain + .first() .cloned() - .ok_or_else(|| anyhow!("provider not found: {}", options.provider_id))?; - let model = if options.model.trim().is_empty() { - provider.model.clone() - } else { - options.model.trim().to_owned() - }; + .ok_or_else(|| anyhow!("transcription chain is empty"))?; options.language = options .language .as_deref() @@ -103,12 +102,17 @@ async fn normalized_job( result_name: None, frontmatter: options.frontmatter, }), - provider_id: provider.id, + provider_id: primary.provider_id.clone(), + transcription_chain, + transcription_attempts: Vec::new(), + used_provider_id: None, + used_provider_name: None, + used_model: None, original_name, source_path, source_size: metadata.len(), source_mtime_ns: mtime_ns(&metadata), - model, + model: primary.model, language: options.language, status: JobStatus::Pending, attempts: 1, diff --git a/src/quick/tests.rs b/src/quick/tests.rs index c54342a..4f52de9 100644 --- a/src/quick/tests.rs +++ b/src/quick/tests.rs @@ -77,6 +77,7 @@ fn options(output_kind: QuickOutputKind, output_dir: Option) -> QuickOpt QuickOptions { provider_id: "provider".into(), model: String::new(), + transcription_chain: Vec::new(), language: None, output_kind, output_dir, diff --git a/src/transcription_chain.rs b/src/transcription_chain.rs new file mode 100644 index 0000000..fb31601 --- /dev/null +++ b/src/transcription_chain.rs @@ -0,0 +1,217 @@ +use crate::{ + domain::{Job, Provider, TranscriptionRoute, Workflow}, + quick::QuickOptions, +}; +use anyhow::{Result, anyhow, bail}; + +fn normalize_explicit_routes( + routes: &[TranscriptionRoute], + providers: &[Provider], +) -> Result> { + let mut normalized = Vec::with_capacity(routes.len()); + for (index, route) in routes.iter().enumerate() { + let provider_id = route.provider_id.trim(); + let model = route.model.trim(); + if provider_id.is_empty() { + bail!("transcription route {} providerId is required", index + 1); + } + if model.is_empty() { + bail!("transcription route {} model is required", index + 1); + } + if route.fallback_after_seconds == Some(0) { + bail!( + "transcription route {} fallback timeout must be greater than zero", + index + 1 + ); + } + if !providers.iter().any(|provider| provider.id == provider_id) { + bail!( + "unknown providerId in transcription route {}: {provider_id}", + index + 1 + ); + } + normalized.push(TranscriptionRoute { + provider_id: provider_id.to_owned(), + model: model.to_owned(), + fallback_after_seconds: route.fallback_after_seconds, + }); + } + Ok(normalized) +} + +fn normalize_legacy_route( + provider_id: &str, + model: &str, + providers: &[Provider], +) -> Result> { + let provider_id = provider_id.trim(); + if provider_id.is_empty() { + bail!("providerId is required"); + } + let provider = providers + .iter() + .find(|provider| provider.id == provider_id) + .ok_or_else(|| anyhow!("unknown providerId: {provider_id}"))?; + let model = if model.trim().is_empty() { + provider.model.trim() + } else { + model.trim() + }; + if model.is_empty() { + bail!("model is required"); + } + Ok(vec![TranscriptionRoute { + provider_id: provider_id.to_owned(), + model: model.to_owned(), + fallback_after_seconds: None, + }]) +} + +pub fn normalize_workflow_chain( + workflow: &Workflow, + providers: &[Provider], +) -> Result> { + if workflow.transcription_chain.is_empty() { + normalize_legacy_route(&workflow.provider_id, &workflow.model, providers) + } else { + normalize_explicit_routes(&workflow.transcription_chain, providers) + } +} + +pub fn normalize_quick_chain( + options: &QuickOptions, + providers: &[Provider], +) -> Result> { + if options.transcription_chain.is_empty() { + normalize_legacy_route(&options.provider_id, &options.model, providers) + } else { + normalize_explicit_routes(&options.transcription_chain, providers) + } +} + +pub fn effective_job_chain(job: &Job) -> Result> { + if !job.transcription_chain.is_empty() { + let mut normalized = Vec::with_capacity(job.transcription_chain.len()); + for (index, route) in job.transcription_chain.iter().enumerate() { + let provider_id = route.provider_id.trim(); + let model = route.model.trim(); + if provider_id.is_empty() || model.is_empty() { + bail!("persisted transcription route {} is invalid", index + 1); + } + if route.fallback_after_seconds == Some(0) { + bail!( + "persisted transcription route {} has invalid timeout", + index + 1 + ); + } + normalized.push(TranscriptionRoute { + provider_id: provider_id.to_owned(), + model: model.to_owned(), + fallback_after_seconds: route.fallback_after_seconds, + }); + } + return Ok(normalized); + } + let provider_id = job.provider_id.trim(); + let model = job.model.trim(); + if provider_id.is_empty() || model.is_empty() { + bail!("legacy persisted job has no usable provider/model"); + } + Ok(vec![TranscriptionRoute { + provider_id: provider_id.to_owned(), + model: model.to_owned(), + fallback_after_seconds: None, + }]) +} + +#[cfg(test)] +mod tests { + use super::*; + use crate::domain::{MarkdownOptions, Provider, TranscriptionRoute, Workflow}; + + fn provider(id: &str, model: &str) -> Provider { + Provider { + id: id.into(), + name: id.into(), + transcription_url: "http://localhost/v1/audio/transcriptions".into(), + model: model.into(), + api_key: String::new(), + timeout_seconds: 600, + enabled: true, + } + } + + fn route(provider_id: &str, model: &str) -> TranscriptionRoute { + TranscriptionRoute { + provider_id: provider_id.into(), + model: model.into(), + fallback_after_seconds: None, + } + } + + fn workflow_with_chain(transcription_chain: Vec) -> Workflow { + Workflow { + id: "w".into(), + name: "test".into(), + watch_dir: "/w".into(), + output_dir: None, + archive_dir: "/a".into(), + tags: vec![], + provider_id: "p".into(), + model: String::new(), + transcription_chain, + language: None, + markdown: MarkdownOptions::default(), + enabled: true, + } + } + + #[test] + fn legacy_workflow_becomes_one_route_with_provider_default_model() { + let workflow: Workflow = serde_json::from_str( + r#"{ + "id":"w","name":"old","watchDir":"/w","archiveDir":"/a", + "providerId":"p","model":"","enabled":true + }"#, + ) + .unwrap(); + let providers = vec![provider("p", "whisper-large-v3")]; + let chain = normalize_workflow_chain(&workflow, &providers).unwrap(); + assert_eq!( + chain, + vec![TranscriptionRoute { + provider_id: "p".into(), + model: "whisper-large-v3".into(), + fallback_after_seconds: None, + }] + ); + } + + #[test] + fn same_provider_can_appear_twice_with_different_models() { + let routes = vec![route("p", "large-v3"), route("p", "distil-large-v3")]; + let workflow = workflow_with_chain(routes.clone()); + assert_eq!( + normalize_workflow_chain(&workflow, &[provider("p", "default")]).unwrap(), + routes + ); + } + + #[test] + fn explicit_routes_are_strict_but_allow_disabled_provider() { + let mut p = provider("p", "default"); + p.enabled = false; + let mut workflow = workflow_with_chain(vec![route(" p ", " large-v3 ")]); + workflow.transcription_chain[0].fallback_after_seconds = Some(30); + assert_eq!( + normalize_workflow_chain(&workflow, &[p]).unwrap(), + vec![TranscriptionRoute { + provider_id: "p".into(), + model: "large-v3".into(), + fallback_after_seconds: Some(30), + }] + ); + workflow.transcription_chain[0].fallback_after_seconds = Some(0); + assert!(normalize_workflow_chain(&workflow, &[provider("p", "default")]).is_err()); + } +} diff --git a/src/workflows.rs b/src/workflows.rs index e65e0c4..a593b4c 100644 --- a/src/workflows.rs +++ b/src/workflows.rs @@ -317,20 +317,19 @@ pub async fn normalize_definition(state: &AppState, mut workflow: Workflow) -> R if archive == watch || archive.starts_with(&watch) { anyhow::bail!("archive directory must be outside the watch directory"); } - if !state - .providers - .read() - .await - .iter() - .any(|provider| provider.id == workflow.provider_id) - { - anyhow::bail!("unknown providerId"); - } + let providers = state.providers.read().await.clone(); + let chain = crate::transcription_chain::normalize_workflow_chain(&workflow, &providers)?; + let primary = chain + .first() + .cloned() + .ok_or_else(|| anyhow::anyhow!("transcription chain is empty"))?; + workflow.transcription_chain = chain; + workflow.provider_id = primary.provider_id; + workflow.model = primary.model; workflow.watch_dir = watch.to_string_lossy().into_owned(); workflow.output_dir = output.map(|path| path.to_string_lossy().into_owned()); workflow.archive_dir = archive.to_string_lossy().into_owned(); workflow.tags = crate::markdown::normalize_custom_tags(&workflow.tags); - workflow.model = workflow.model.trim().to_owned(); workflow.language = workflow .language .as_deref() @@ -417,6 +416,7 @@ mod validation_tests { tags: Vec::new(), provider_id: "provider".into(), model: String::new(), + transcription_chain: Vec::new(), language: Some("auto".into()), markdown: MarkdownOptions::default(), enabled: true, @@ -462,6 +462,7 @@ mod validation_tests { tags: Vec::new(), provider_id: "provider".into(), model: String::new(), + transcription_chain: Vec::new(), language: None, markdown: MarkdownOptions::default(), enabled: true, @@ -494,6 +495,7 @@ mod validation_tests { tags: Vec::new(), provider_id: "provider".into(), model: String::new(), + transcription_chain: Vec::new(), language: None, markdown: MarkdownOptions::default(), enabled: true, From 328f51dc1dc0a76a14abb7d2ae269a418e254d30 Mon Sep 17 00:00:00 2001 From: GodsQuantum <193283917+GodsQuantum@users.noreply.github.com> Date: Wed, 16 Sep 2026 11:13:03 +0200 Subject: [PATCH 2/6] feat: execute resilient transcription fallbacks --- src/provider.rs | 7 +- src/transcription_chain.rs | 567 ++++++++++++++++++++++++++++++++++++- 2 files changed, 567 insertions(+), 7 deletions(-) diff --git a/src/provider.rs b/src/provider.rs index 4df931d..1f989bb 100644 --- a/src/provider.rs +++ b/src/provider.rs @@ -51,12 +51,7 @@ pub async fn transcribe( form = form.text("language", language.to_owned()); } - let mut request = client - .post(&provider.transcription_url) - .timeout(std::time::Duration::from_secs( - provider.timeout_seconds.max(1), - )) - .multipart(form); + let mut request = client.post(&provider.transcription_url).multipart(form); if !provider.api_key.trim().is_empty() { request = request.bearer_auth(&provider.api_key); } diff --git a/src/transcription_chain.rs b/src/transcription_chain.rs index fb31601..3ea8493 100644 --- a/src/transcription_chain.rs +++ b/src/transcription_chain.rs @@ -1,8 +1,16 @@ use crate::{ - domain::{Job, Provider, TranscriptionRoute, Workflow}, + domain::{ + Job, JobStatus, Provider, TranscriptionAttempt, TranscriptionAttemptOutcome, + TranscriptionRoute, Workflow, + }, + jobs::{get_job, now_ms, update_job}, + provider::{self, TranscriptionError}, quick::QuickOptions, + state::AppState, }; use anyhow::{Result, anyhow, bail}; +use std::{path::Path, time::Duration}; +use tokio_util::sync::CancellationToken; fn normalize_explicit_routes( routes: &[TranscriptionRoute], @@ -124,6 +132,245 @@ pub fn effective_job_chain(job: &Job) -> Result> { }]) } +#[derive(Debug, Clone, PartialEq, Eq)] +pub struct TranscriptionSuccess { + pub transcript: String, + pub provider_id: String, + pub provider_name: String, + pub model: String, +} + +fn safe_attempt_error(value: impl ToString) -> String { + let flattened = value.to_string().replace(['\r', '\n'], " "); + flattened.chars().take(300).collect() +} + +fn append_attempt(state: &AppState, job_id: &str, attempt: TranscriptionAttempt) -> Result<()> { + update_job(state, job_id, |job| { + job.transcription_attempts.push(attempt) + })?; + Ok(()) +} + +fn finished_attempt( + run: u32, + route_index: usize, + route: &TranscriptionRoute, + provider_name: &str, + started_at_ms: u128, + outcome: TranscriptionAttemptOutcome, + error: Option, +) -> TranscriptionAttempt { + TranscriptionAttempt { + run, + route_index, + provider_id: route.provider_id.clone(), + provider_name: provider_name.to_owned(), + model: route.model.clone(), + started_at_ms, + finished_at_ms: now_ms(), + outcome, + error, + } +} + +pub async fn transcribe_job_chain( + state: &AppState, + job_id: &str, + source: &Path, + language: Option<&str>, + token: &CancellationToken, +) -> Result { + let job = get_job(state, job_id)?; + let routes = effective_job_chain(&job)?; + let run = job.attempts.max(1); + let mut failures = Vec::new(); + + for (route_index, route) in routes.iter().enumerate() { + if token.is_cancelled() { + bail!("cancelled"); + } + + let provider = state + .providers + .read() + .await + .iter() + .find(|provider| provider.id == route.provider_id) + .cloned(); + let started_at_ms = now_ms(); + let Some(provider) = provider else { + let message = "provider is no longer configured".to_owned(); + append_attempt( + state, + job_id, + finished_attempt( + run, + route_index, + route, + &route.provider_id, + started_at_ms, + TranscriptionAttemptOutcome::Failed, + Some(message.clone()), + ), + )?; + failures.push(format!("{}/{}: {message}", route.provider_id, route.model)); + continue; + }; + if !provider.enabled { + let message = "provider is disabled".to_owned(); + append_attempt( + state, + job_id, + finished_attempt( + run, + route_index, + route, + &provider.name, + started_at_ms, + TranscriptionAttemptOutcome::Failed, + Some(message.clone()), + ), + )?; + failures.push(format!("{}/{}: {message}", provider.name, route.model)); + continue; + } + + update_job(state, job_id, |job| { + job.status = JobStatus::Transcribing; + job.error = None; + })?; + + let request = provider::transcribe( + &provider, + &route.model, + language, + source, + &state.http, + token, + ); + let result = if let Some(seconds) = route.fallback_after_seconds { + match tokio::time::timeout(Duration::from_secs(seconds), request).await { + Ok(result) => result, + Err(_) => { + let message = format!("fallback timeout after {seconds}s"); + append_attempt( + state, + job_id, + finished_attempt( + run, + route_index, + route, + &provider.name, + started_at_ms, + TranscriptionAttemptOutcome::TimedOut, + Some(message.clone()), + ), + )?; + failures.push(format!("{}/{}: {message}", provider.name, route.model)); + continue; + } + } + } else { + request.await + }; + + match result { + Ok(transcript) => { + let attempt = finished_attempt( + run, + route_index, + route, + &provider.name, + started_at_ms, + TranscriptionAttemptOutcome::Success, + None, + ); + update_job(state, job_id, |job| { + job.transcription_attempts.push(attempt); + job.used_provider_id = Some(provider.id.clone()); + job.used_provider_name = Some(provider.name.clone()); + job.used_model = Some(route.model.clone()); + })?; + return Ok(TranscriptionSuccess { + transcript, + provider_id: provider.id, + provider_name: provider.name, + model: route.model.clone(), + }); + } + Err(TranscriptionError::Cancelled) => { + append_attempt( + state, + job_id, + finished_attempt( + run, + route_index, + route, + &provider.name, + started_at_ms, + TranscriptionAttemptOutcome::Cancelled, + Some("cancelled".into()), + ), + )?; + bail!("cancelled"); + } + Err(TranscriptionError::Io(error)) => { + let message = safe_attempt_error(&error); + append_attempt( + state, + job_id, + finished_attempt( + run, + route_index, + route, + &provider.name, + started_at_ms, + TranscriptionAttemptOutcome::Failed, + Some(message.clone()), + ), + )?; + return Err(error.into()); + } + Err(error) => { + if token.is_cancelled() { + append_attempt( + state, + job_id, + finished_attempt( + run, + route_index, + route, + &provider.name, + started_at_ms, + TranscriptionAttemptOutcome::Cancelled, + Some("cancelled".into()), + ), + )?; + bail!("cancelled"); + } + let message = safe_attempt_error(&error); + append_attempt( + state, + job_id, + finished_attempt( + run, + route_index, + route, + &provider.name, + started_at_ms, + TranscriptionAttemptOutcome::Failed, + Some(message.clone()), + ), + )?; + failures.push(format!("{}/{}: {message}", provider.name, route.model)); + } + } + } + + bail!("all transcription routes failed: {}", failures.join("; ")) +} + #[cfg(test)] mod tests { use super::*; @@ -214,4 +461,322 @@ mod tests { workflow.transcription_chain[0].fallback_after_seconds = Some(0); assert!(normalize_workflow_chain(&workflow, &[provider("p", "default")]).is_err()); } + + use crate::{ + config::Config, + domain::{Job, JobKind, JobStatus}, + jobs::{get_job, persist_job}, + state::AppState, + }; + use axum::{Json, Router, body::Bytes, http::StatusCode, routing::post}; + use serde_json::json; + use std::{ + path::{Path, PathBuf}, + sync::{ + Arc, + atomic::{AtomicUsize, Ordering}, + }, + time::Duration, + }; + use tokio_util::sync::CancellationToken; + + async fn mock_provider( + status: StatusCode, + text: &'static str, + delay_ms: u64, + calls: Arc, + ) -> String { + let app = Router::new().route("/v1/audio/transcriptions", post(move |_body: Bytes| { + let calls = calls.clone(); + async move { + calls.fetch_add(1, Ordering::SeqCst); + if delay_ms > 0 { tokio::time::sleep(Duration::from_millis(delay_ms)).await; } + (status, Json(json!({"text": text, "error": if status.is_success() { "" } else { text }}))) + } + })); + let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap(); + let address = listener.local_addr().unwrap(); + tokio::spawn(async move { + axum::serve(listener, app).await.unwrap(); + }); + format!("http://{address}/v1/audio/transcriptions") + } + + fn test_config(root: &Path) -> Config { + Config { + host: "127.0.0.1".into(), + port: 0, + config_dir: root.join("config"), + dist_dir: root.join("dist"), + data_dir: root.join("data"), + allowed_roots: vec![root.to_path_buf()], + scan_seconds: 1, + file_stability_ms: 20, + max_transcription_jobs: 1, + max_upload_bytes: 2_147_483_648, + quick_result_retention_hours: 24, + } + } + + async fn state_with_job( + root: &Path, + providers: Vec, + routes: Vec, + ) -> (AppState, Job, PathBuf) { + let state = AppState::load(test_config(root)).await.unwrap(); + for provider in &providers { + state.db.upsert("provider", &provider.id, provider).unwrap(); + } + *state.providers.write().await = providers; + let source = root.join("audio.m4a"); + std::fs::write(&source, b"audio").unwrap(); + let primary = routes.first().unwrap().clone(); + let now = crate::jobs::now_ms(); + let job = Job { + id: "chain-job".into(), + kind: JobKind::Quick, + workflow_id: None, + quick: None, + provider_id: primary.provider_id.clone(), + transcription_chain: routes, + transcription_attempts: vec![], + used_provider_id: None, + used_provider_name: None, + used_model: None, + original_name: "audio.m4a".into(), + source_path: source.clone(), + source_size: 5, + source_mtime_ns: 1, + model: primary.model, + language: None, + status: JobStatus::Pending, + attempts: 1, + error: None, + markdown_path: None, + archive_path: None, + markdown_published: false, + created_at_ms: now, + updated_at_ms: now, + }; + state.jobs.insert(job.id.clone(), job.clone()); + persist_job(&state, &job).unwrap(); + (state, job, source) + } + + fn test_provider(id: &str, url: String, enabled: bool) -> Provider { + Provider { + id: id.into(), + name: id.into(), + transcription_url: url, + model: "default".into(), + api_key: String::new(), + timeout_seconds: 1, + enabled, + } + } + + #[tokio::test] + async fn primary_failure_then_fallback_success_records_both_attempts() { + let temp = tempfile::tempdir().unwrap(); + let c1 = Arc::new(AtomicUsize::new(0)); + let c2 = Arc::new(AtomicUsize::new(0)); + let p1 = test_provider( + "p1", + mock_provider(StatusCode::BAD_GATEWAY, "down", 0, c1.clone()).await, + true, + ); + let p2 = test_provider( + "p2", + mock_provider(StatusCode::OK, "fallback works", 0, c2.clone()).await, + true, + ); + let routes = vec![route("p1", "m1"), route("p2", "m2")]; + let (state, job, source) = state_with_job(temp.path(), vec![p1, p2], routes).await; + let result = + transcribe_job_chain(&state, &job.id, &source, None, &CancellationToken::new()) + .await + .unwrap(); + assert_eq!(result.transcript, "fallback works"); + assert_eq!(result.provider_id, "p2"); + assert_eq!(result.model, "m2"); + let saved = get_job(&state, &job.id).unwrap(); + assert_eq!(saved.transcription_attempts.len(), 2); + assert_eq!( + saved.transcription_attempts[0].outcome, + crate::domain::TranscriptionAttemptOutcome::Failed + ); + assert_eq!( + saved.transcription_attempts[1].outcome, + crate::domain::TranscriptionAttemptOutcome::Success + ); + assert_eq!(c1.load(Ordering::SeqCst), 1); + assert_eq!(c2.load(Ordering::SeqCst), 1); + } + + #[tokio::test] + async fn disabled_or_missing_provider_falls_through() { + let temp = tempfile::tempdir().unwrap(); + let calls = Arc::new(AtomicUsize::new(0)); + let disabled = test_provider( + "disabled", + "http://127.0.0.1:1/v1/audio/transcriptions".into(), + false, + ); + let good = test_provider( + "good", + mock_provider(StatusCode::OK, "ok", 0, calls.clone()).await, + true, + ); + let routes = vec![ + route("missing", "m0"), + route("disabled", "m1"), + route("good", "m2"), + ]; + let (state, job, source) = state_with_job(temp.path(), vec![disabled, good], routes).await; + let result = + transcribe_job_chain(&state, &job.id, &source, None, &CancellationToken::new()) + .await + .unwrap(); + assert_eq!(result.provider_id, "good"); + assert_eq!(calls.load(Ordering::SeqCst), 1); + assert_eq!( + get_job(&state, &job.id) + .unwrap() + .transcription_attempts + .len(), + 3 + ); + } + + #[tokio::test] + async fn route_timeout_falls_through() { + let temp = tempfile::tempdir().unwrap(); + let slow = Arc::new(AtomicUsize::new(0)); + let fast = Arc::new(AtomicUsize::new(0)); + let p1 = test_provider( + "slow", + mock_provider(StatusCode::OK, "too late", 1500, slow.clone()).await, + true, + ); + let p2 = test_provider( + "fast", + mock_provider(StatusCode::OK, "fast", 0, fast.clone()).await, + true, + ); + let mut first = route("slow", "m1"); + first.fallback_after_seconds = Some(1); + let (state, job, source) = + state_with_job(temp.path(), vec![p1, p2], vec![first, route("fast", "m2")]).await; + let result = + transcribe_job_chain(&state, &job.id, &source, None, &CancellationToken::new()) + .await + .unwrap(); + assert_eq!(result.transcript, "fast"); + let saved = get_job(&state, &job.id).unwrap(); + assert_eq!( + saved.transcription_attempts[0].outcome, + crate::domain::TranscriptionAttemptOutcome::TimedOut + ); + assert_eq!(slow.load(Ordering::SeqCst), 1); + assert_eq!(fast.load(Ordering::SeqCst), 1); + } + + #[tokio::test] + async fn no_route_timeout_waits_for_slow_success() { + let temp = tempfile::tempdir().unwrap(); + let calls = Arc::new(AtomicUsize::new(0)); + let p = test_provider( + "slow", + mock_provider(StatusCode::OK, "eventually", 150, calls.clone()).await, + true, + ); + let (state, job, source) = + state_with_job(temp.path(), vec![p], vec![route("slow", "m1")]).await; + let result = + transcribe_job_chain(&state, &job.id, &source, None, &CancellationToken::new()) + .await + .unwrap(); + assert_eq!(result.transcript, "eventually"); + assert_eq!(calls.load(Ordering::SeqCst), 1); + } + + #[tokio::test] + async fn cancellation_stops_chain_before_second_provider() { + let temp = tempfile::tempdir().unwrap(); + let slow = Arc::new(AtomicUsize::new(0)); + let second = Arc::new(AtomicUsize::new(0)); + let p1 = test_provider( + "slow", + mock_provider(StatusCode::OK, "late", 1000, slow.clone()).await, + true, + ); + let p2 = test_provider( + "second", + mock_provider(StatusCode::OK, "should not run", 0, second.clone()).await, + true, + ); + let (state, job, source) = state_with_job( + temp.path(), + vec![p1, p2], + vec![route("slow", "m1"), route("second", "m2")], + ) + .await; + let token = CancellationToken::new(); + let cancel = token.clone(); + tokio::spawn(async move { + tokio::time::sleep(Duration::from_millis(50)).await; + cancel.cancel(); + }); + assert!( + transcribe_job_chain(&state, &job.id, &source, None, &token) + .await + .is_err() + ); + assert_eq!(second.load(Ordering::SeqCst), 0); + assert_eq!( + get_job(&state, &job.id) + .unwrap() + .transcription_attempts + .last() + .unwrap() + .outcome, + crate::domain::TranscriptionAttemptOutcome::Cancelled + ); + } + + #[tokio::test] + async fn all_routes_fail_returns_aggregate_error_and_history() { + let temp = tempfile::tempdir().unwrap(); + let c1 = Arc::new(AtomicUsize::new(0)); + let c2 = Arc::new(AtomicUsize::new(0)); + let p1 = test_provider( + "p1", + mock_provider(StatusCode::BAD_GATEWAY, "first down", 0, c1).await, + true, + ); + let p2 = test_provider( + "p2", + mock_provider(StatusCode::TOO_MANY_REQUESTS, "rate limited", 0, c2).await, + true, + ); + let (state, job, source) = state_with_job( + temp.path(), + vec![p1, p2], + vec![route("p1", "m1"), route("p2", "m2")], + ) + .await; + let error = transcribe_job_chain(&state, &job.id, &source, None, &CancellationToken::new()) + .await + .unwrap_err() + .to_string(); + assert!(error.contains("p1/m1")); + assert!(error.contains("p2/m2")); + assert_eq!( + get_job(&state, &job.id) + .unwrap() + .transcription_attempts + .len(), + 2 + ); + } } From 9411ed68f7b9d20ead0c3866684f54b7f08d58c8 Mon Sep 17 00:00:00 2001 From: GodsQuantum <193283917+GodsQuantum@users.noreply.github.com> Date: Wed, 16 Sep 2026 11:50:40 +0200 Subject: [PATCH 3/6] feat: wire fallback chains into transcription jobs --- src/jobs.rs | 151 +++++++++++++++++++++++++++++++++++++++++- src/pipeline.rs | 32 +++------ src/pipeline_quick.rs | 31 +++------ 3 files changed, 165 insertions(+), 49 deletions(-) diff --git a/src/jobs.rs b/src/jobs.rs index a9022ad..e33cd24 100644 --- a/src/jobs.rs +++ b/src/jobs.rs @@ -174,11 +174,18 @@ mod tests { use crate::pipeline::process_job; use crate::{ config::Config, - domain::{MarkdownOptions, Provider}, + domain::{MarkdownOptions, Provider, TranscriptionAttemptOutcome, TranscriptionRoute}, }; use axum::{Json, Router, body::Bytes, http::StatusCode, routing::post}; use serde_json::json; - use std::{path::Path, time::Duration}; + use std::{ + path::Path, + sync::{ + Arc, + atomic::{AtomicUsize, Ordering}, + }, + time::Duration, + }; async fn success_provider() -> String { let app = Router::new().route( @@ -200,6 +207,37 @@ mod tests { ); spawn_server(app).await } + + async fn counted_provider( + status: StatusCode, + text: &'static str, + ) -> (String, Arc) { + let calls = Arc::new(AtomicUsize::new(0)); + let seen = calls.clone(); + let app = Router::new().route( + "/v1/audio/transcriptions", + post(move |_body: Bytes| { + let seen = seen.clone(); + async move { + seen.fetch_add(1, Ordering::SeqCst); + if status.is_success() { + (status, Json(json!({"text": text}))) + } else { + (status, Json(json!({"error": text}))) + } + } + }), + ); + (spawn_server(app).await, calls) + } + + fn route(provider_id: &str, model: &str) -> TranscriptionRoute { + TranscriptionRoute { + provider_id: provider_id.into(), + model: model.into(), + fallback_after_seconds: None, + } + } async fn spawn_server(app: Router) -> String { let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap(); let address = listener.local_addr().unwrap(); @@ -412,6 +450,102 @@ mod tests { .is_file() ); } + #[tokio::test] + async fn markdown_uses_provider_and_model_that_actually_succeeded() { + let temp = tempfile::tempdir().unwrap(); + let (fail_url, _) = counted_provider(StatusCode::BAD_GATEWAY, "primary down").await; + let (success_url, success_calls) = + counted_provider(StatusCode::OK, "fallback transcript").await; + let (state, mut workflow) = setup(temp.path(), fail_url).await; + state.providers.write().await.push(Provider { + id: "fallback".into(), + name: "Fallback Provider".into(), + transcription_url: success_url, + model: "fallback-default".into(), + api_key: String::new(), + timeout_seconds: 5, + enabled: true, + }); + workflow.transcription_chain = vec![ + route("provider", "primary-model"), + route("fallback", "fallback-model"), + ]; + let job = make_job(&state, &workflow, "winner.m4a").await; + process_job(&state, &job.id, &CancellationToken::new()) + .await + .unwrap(); + assert_eq!(success_calls.load(Ordering::SeqCst), 1); + let finished = get_job(&state, &job.id).unwrap(); + assert_eq!(finished.used_provider_id.as_deref(), Some("fallback")); + assert_eq!(finished.used_model.as_deref(), Some("fallback-model")); + let markdown = std::fs::read_to_string(finished.markdown_path.unwrap()).unwrap(); + assert!(markdown.contains("provider: \"Fallback Provider\"")); + assert!(markdown.contains("model: \"fallback-model\"")); + assert!(markdown.contains("fallback transcript")); + } + + #[tokio::test] + async fn publication_failure_does_not_contact_another_route() { + let temp = tempfile::tempdir().unwrap(); + let (primary_url, primary_calls) = counted_provider(StatusCode::OK, "publish me").await; + let (fallback_url, fallback_calls) = + counted_provider(StatusCode::OK, "should not run").await; + let (state, mut workflow) = setup(temp.path(), primary_url).await; + state.providers.write().await.push(Provider { + id: "fallback".into(), + name: "Fallback".into(), + transcription_url: fallback_url, + model: "fallback-model".into(), + api_key: String::new(), + timeout_seconds: 5, + enabled: true, + }); + let notes = temp.path().join("notes"); + std::fs::create_dir_all(¬es).unwrap(); + workflow.output_dir = Some(notes.to_string_lossy().into_owned()); + workflow.transcription_chain = vec![ + route("provider", "primary-model"), + route("fallback", "fallback-model"), + ]; + state.workflows.write().await[0] = workflow.clone(); + state + .db + .upsert("workflow", &workflow.id, &workflow) + .unwrap(); + let job = make_job(&state, &workflow, "publish-fail.m4a").await; + std::fs::remove_dir(¬es).unwrap(); + assert!( + process_job(&state, &job.id, &CancellationToken::new()) + .await + .is_err() + ); + assert_eq!(primary_calls.load(Ordering::SeqCst), 1); + assert_eq!(fallback_calls.load(Ordering::SeqCst), 0); + } + + #[tokio::test] + async fn archive_failure_does_not_retranscribe() { + let temp = tempfile::tempdir().unwrap(); + let (primary_url, primary_calls) = counted_provider(StatusCode::OK, "archive me").await; + let (state, mut workflow) = setup(temp.path(), primary_url).await; + workflow.transcription_chain = vec![route("provider", "primary-model")]; + let job = make_job(&state, &workflow, "archive-fail.m4a").await; + std::fs::remove_dir(&workflow.archive_dir).unwrap(); + assert!( + process_job(&state, &job.id, &CancellationToken::new()) + .await + .is_err() + ); + assert_eq!(primary_calls.load(Ordering::SeqCst), 1); + assert_eq!( + get_job(&state, &job.id) + .unwrap() + .transcription_attempts + .len(), + 1 + ); + } + #[tokio::test] async fn retry_reuses_failed_job_and_completes() { let temp = tempfile::tempdir().unwrap(); @@ -427,6 +561,19 @@ mod tests { let finished = wait_for_terminal(&state, &job.id).await; assert_eq!(finished.status, JobStatus::Done); assert_eq!(finished.attempts, 2); + assert_eq!(finished.transcription_attempts.len(), 2); + assert_eq!(finished.transcription_attempts[0].run, 1); + assert_eq!(finished.transcription_attempts[0].route_index, 0); + assert_eq!( + finished.transcription_attempts[0].outcome, + TranscriptionAttemptOutcome::Failed + ); + assert_eq!(finished.transcription_attempts[1].run, 2); + assert_eq!(finished.transcription_attempts[1].route_index, 0); + assert_eq!( + finished.transcription_attempts[1].outcome, + TranscriptionAttemptOutcome::Success + ); assert!( Path::new(&workflow.watch_dir) .join("hello from provider.md") diff --git a/src/pipeline.rs b/src/pipeline.rs index 38fd376..a7a52f9 100644 --- a/src/pipeline.rs +++ b/src/pipeline.rs @@ -2,8 +2,8 @@ use crate::{ domain::{JobKind, JobStatus, Workflow}, jobs::{get_job, update_job}, markdown::{self, NoteContext}, - provider, state::AppState, + transcription_chain, }; use anyhow::{Context, Result, anyhow, bail}; use std::{ @@ -38,18 +38,6 @@ async fn process_workflow_job( .find(|workflow| workflow.id == workflow_id) .cloned() .ok_or_else(|| anyhow!("workflow not found: {workflow_id}"))?; - let provider = state - .providers - .read() - .await - .iter() - .find(|provider| provider.id == job.provider_id) - .cloned() - .ok_or_else(|| anyhow!("provider not found: {}", job.provider_id))?; - if !provider.enabled { - bail!("provider '{}' is disabled", provider.name); - } - let source = state.config.resolve_allowed_file(&job.source_path)?; let watch_dir = state .config @@ -71,20 +59,16 @@ async fn process_workflow_job( return finish_archive(state, &job.id, &workflow, &source, token).await; } - update_job(state, &job.id, |job| { - job.status = JobStatus::Transcribing; - job.error = None; - })?; - let transcript = provider::transcribe( - &provider, - &job.model, - job.language.as_deref(), + let success = transcription_chain::transcribe_job_chain( + state, + &job.id, &source, - &state.http, + job.language.as_deref(), token, ) .await .context("transcription")?; + let transcript = success.transcript; if token.is_cancelled() { bail!("cancelled"); } @@ -99,8 +83,8 @@ async fn process_workflow_job( title: &title, source_name: &job.original_name, workflow_name: Some(&workflow.name), - provider_name: &provider.name, - model: &job.model, + provider_name: &success.provider_name, + model: &success.model, language: job.language.as_deref(), tags: &workflow.tags, frontmatter: workflow.markdown.frontmatter, diff --git a/src/pipeline_quick.rs b/src/pipeline_quick.rs index 54cdb19..f216081 100644 --- a/src/pipeline_quick.rs +++ b/src/pipeline_quick.rs @@ -2,8 +2,8 @@ use crate::{ domain::{Job, JobStatus, QuickJobMeta, QuickOutputKind, QuickSourceKind}, jobs::update_job, markdown::{self, NoteContext}, - provider, state::AppState, + transcription_chain, }; use anyhow::{Context, Result, anyhow, bail}; use std::path::Path; @@ -45,34 +45,19 @@ async fn process_quick_inner( source: &Path, token: &CancellationToken, ) -> Result<()> { - let provider = state - .providers - .read() - .await - .iter() - .find(|provider| provider.id == job.provider_id) - .cloned() - .ok_or_else(|| anyhow!("provider not found: {}", job.provider_id))?; - if !provider.enabled { - bail!("provider '{}' is disabled", provider.name); - } if token.is_cancelled() { bail!("cancelled"); } - update_job(state, &job.id, |job| { - job.status = JobStatus::Transcribing; - job.error = None; - })?; - let transcript = provider::transcribe( - &provider, - &job.model, - job.language.as_deref(), + let success = transcription_chain::transcribe_job_chain( + state, + &job.id, source, - &state.http, + job.language.as_deref(), token, ) .await .context("transcription")?; + let transcript = success.transcript; if token.is_cancelled() { bail!("cancelled"); } @@ -82,8 +67,8 @@ async fn process_quick_inner( title: &title, source_name: &job.original_name, workflow_name: None, - provider_name: &provider.name, - model: &job.model, + provider_name: &success.provider_name, + model: &success.model, language: job.language.as_deref(), tags: &[], frontmatter: meta.frontmatter, From c7c349f34b92a19f042a03dc8ada3bfd12927f74 Mon Sep 17 00:00:00 2001 From: GodsQuantum <193283917+GodsQuantum@users.noreply.github.com> Date: Wed, 16 Sep 2026 11:54:43 +0200 Subject: [PATCH 4/6] feat: expose transcription chains through the API --- frontend/src/lib/api.ts | 5 +- frontend/src/lib/types.ts | 29 ++++++++- frontend/tests/scribewatch-surface.test.mjs | 11 ++++ src/api/quick.rs | 66 +++++++++++++++++++-- 4 files changed, 103 insertions(+), 8 deletions(-) diff --git a/frontend/src/lib/api.ts b/frontend/src/lib/api.ts index 83d2e73..9968e48 100644 --- a/frontend/src/lib/api.ts +++ b/frontend/src/lib/api.ts @@ -25,8 +25,9 @@ async function request(path: string, init?: RequestInit): Promise { } function appendQuickOptions(form: FormData, options: QuickOptions) { - form.append('providerId', options.providerId); - form.append('model', options.model ?? ''); + if (options.providerId) form.append('providerId', options.providerId); + if (options.model) form.append('model', options.model); + if (options.transcriptionChain) form.append('transcriptionChain', JSON.stringify(options.transcriptionChain)); form.append('language', options.language ?? ''); form.append('outputKind', options.outputKind); form.append('outputDir', options.outputDir ?? ''); diff --git a/frontend/src/lib/types.ts b/frontend/src/lib/types.ts index f5d1c0b..6225b29 100644 --- a/frontend/src/lib/types.ts +++ b/frontend/src/lib/types.ts @@ -5,6 +5,26 @@ export type JobKind = 'workflow' | 'quick'; export type QuickSourceKind = 'server' | 'upload'; export type QuickOutputKind = 'server' | 'client'; +export interface TranscriptionRoute { + providerId: string; + model: string; + fallbackAfterSeconds?: number; +} + +export type TranscriptionAttemptOutcome = 'success' | 'failed' | 'timed_out' | 'cancelled'; + +export interface TranscriptionAttempt { + run: number; + routeIndex: number; + providerId: string; + providerName: string; + model: string; + startedAtMs: number; + finishedAtMs: number; + outcome: TranscriptionAttemptOutcome; + error?: string; +} + export interface Provider { id: string; name: string; @@ -38,6 +58,7 @@ export interface Workflow { tags: string[]; providerId: string; model: string; + transcriptionChain?: TranscriptionRoute[]; language?: string; markdown: MarkdownOptions; enabled: boolean; @@ -57,6 +78,11 @@ export interface Job { workflowId?: string; quick?: QuickJobMeta; providerId: string; + transcriptionChain: TranscriptionRoute[]; + transcriptionAttempts: TranscriptionAttempt[]; + usedProviderId?: string; + usedProviderName?: string; + usedModel?: string; originalName: string; sourcePath: string; sourceSize: number; @@ -74,8 +100,9 @@ export interface Job { } export interface QuickOptions { - providerId: string; + providerId?: string; model?: string; + transcriptionChain?: TranscriptionRoute[]; language?: string; outputKind: QuickOutputKind; outputDir?: string; diff --git a/frontend/tests/scribewatch-surface.test.mjs b/frontend/tests/scribewatch-surface.test.mjs index f801fee..7d08c6f 100644 --- a/frontend/tests/scribewatch-surface.test.mjs +++ b/frontend/tests/scribewatch-surface.test.mjs @@ -40,3 +40,14 @@ test('product promise is visible on Home', async () => { const view = await read('lib/views/DashboardView.svelte'); assert.ok(view.includes('Turn voice notes into Markdown before you forget them.')); }); + + +test('quick transport exposes transcription chains', async () => { + const api = await read('lib/api.ts'); + const types = await read('lib/types.ts'); + assert.ok(api.includes("form.append('transcriptionChain'")); + assert.ok(types.includes('interface TranscriptionRoute')); + assert.ok(types.includes('transcriptionChain')); + assert.ok(types.includes('transcriptionAttempts')); + assert.ok(types.includes('usedProviderId')); +}); diff --git a/src/api/quick.rs b/src/api/quick.rs index e695ba7..29ccdae 100644 --- a/src/api/quick.rs +++ b/src/api/quick.rs @@ -1,5 +1,5 @@ use crate::{ - domain::{Job, QuickOutputKind}, + domain::{Job, QuickOutputKind, TranscriptionRoute}, error::{AppError, AppResult}, jobs, quick::{self, QuickOptions}, @@ -44,6 +44,27 @@ fn text_bool(value: Option, default: bool) -> Result { } } +fn parse_transcription_chain(value: Option) -> Result, AppError> { + let Some(value) = value else { + return Ok(Vec::new()); + }; + let value = value.trim(); + if value.is_empty() { + return Ok(Vec::new()); + } + let routes: Vec = serde_json::from_str(value) + .map_err(|error| AppError::BadRequest(format!("invalid transcriptionChain: {error}")))?; + if routes + .iter() + .any(|route| route.fallback_after_seconds == Some(0)) + { + return Err(AppError::BadRequest( + "fallbackAfterSeconds must be greater than zero".into(), + )); + } + Ok(routes) +} + fn parse_output_kind(value: Option) -> Result { match value.as_deref().map(str::trim) { Some("server") => Ok(QuickOutputKind::Server), @@ -67,6 +88,7 @@ pub async fn upload( let mut original_name: Option = None; let mut provider_id = None; let mut model = None; + let mut transcription_chain = None; let mut language = None; let mut output_kind = None; let mut output_dir = None; @@ -159,6 +181,7 @@ pub async fn upload( match name.as_str() { "providerId" => provider_id = Some(value), "model" => model = Some(value), + "transcriptionChain" => transcription_chain = Some(value), "language" => language = Some(value), "outputKind" => output_kind = Some(value), "outputDir" => output_dir = Some(value), @@ -171,12 +194,15 @@ pub async fn upload( staged.ok_or_else(|| AppError::BadRequest("audio file is required".into()))?; let original_name = original_name.unwrap_or_else(|| "audio".into()); let options_result: AppResult = (|| { + let transcription_chain = parse_transcription_chain(transcription_chain)?; + let provider_id = provider_id.unwrap_or_default(); + if transcription_chain.is_empty() && provider_id.trim().is_empty() { + return Err(AppError::BadRequest("providerId is required".into())); + } Ok(QuickOptions { - provider_id: provider_id - .filter(|value| !value.trim().is_empty()) - .ok_or_else(|| AppError::BadRequest("providerId is required".into()))?, + provider_id, model: model.unwrap_or_default(), - transcription_chain: Vec::new(), + transcription_chain, language, output_kind: parse_output_kind(output_kind)?, output_dir, @@ -231,4 +257,34 @@ mod tests { assert!(text_bool(Some("1".into()), false).unwrap()); assert!(text_bool(Some("sometimes".into()), true).is_err()); } + + #[test] + fn transcription_chain_parser_accepts_ordered_routes() { + let chain = parse_transcription_chain(Some( + r#"[ + {"providerId":"p","model":"large-v3"}, + {"providerId":"p","model":"distil-large-v3","fallbackAfterSeconds":1800} + ]"# + .into(), + )) + .unwrap(); + assert_eq!(chain.len(), 2); + assert_eq!(chain[0].provider_id, "p"); + assert_eq!(chain[0].model, "large-v3"); + assert_eq!(chain[0].fallback_after_seconds, None); + assert_eq!(chain[1].model, "distil-large-v3"); + assert_eq!(chain[1].fallback_after_seconds, Some(1800)); + } + + #[test] + fn transcription_chain_parser_rejects_bad_json_and_zero_timeout() { + assert!(parse_transcription_chain(Some("not-json".into())).is_err()); + assert!( + parse_transcription_chain(Some( + r#"[{"providerId":"p","model":"m","fallbackAfterSeconds":0}]"#.into() + )) + .is_err() + ); + assert!(parse_transcription_chain(None).unwrap().is_empty()); + } } From b0b7703fb86c1969730df2470a16c0849cd32e2c Mon Sep 17 00:00:00 2001 From: GodsQuantum <193283917+GodsQuantum@users.noreply.github.com> Date: Wed, 16 Sep 2026 12:01:28 +0200 Subject: [PATCH 5/6] feat: add fallback chain editor --- frontend/src/app.css | 5 + .../TranscriptionChainEditor.svelte | 127 ++++++++++++++++++ frontend/src/lib/transcription-chain.js | 30 +++++ frontend/src/lib/views/JobsView.svelte | 24 +++- frontend/src/lib/views/ProvidersView.svelte | 2 +- .../src/lib/views/QuickTranscribeView.svelte | 29 ++-- frontend/src/lib/views/WorkflowsView.svelte | 29 ++-- frontend/tests/scribewatch-surface.test.mjs | 17 +++ frontend/tests/transcription-chain.test.mjs | 32 +++++ 9 files changed, 269 insertions(+), 26 deletions(-) create mode 100644 frontend/src/lib/components/TranscriptionChainEditor.svelte create mode 100644 frontend/src/lib/transcription-chain.js create mode 100644 frontend/tests/transcription-chain.test.mjs diff --git a/frontend/src/app.css b/frontend/src/app.css index b474b5b..bd92bd0 100644 --- a/frontend/src/app.css +++ b/frontend/src/app.css @@ -47,3 +47,8 @@ button,input,select,textarea{font:inherit} button{color:inherit} a{color:inherit @media(max-width:1080px){.home-hero{grid-template-columns:1fr}.hero-visual{display:none}.quick-layout{grid-template-columns:1fr}.result-card{position:static;min-height:300px}.result-body{min-height:230px}.home-kpis{grid-template-columns:repeat(2,minmax(0,1fr))}} @media(max-width:680px){.home-hero{padding:27px 20px;border-radius:18px}.home-hero h1{font-size:39px}.hero-actions{flex-direction:column}.hero-actions .btn{width:100%}.hero-proof{gap:10px;flex-direction:column}.quick-card .card-body{padding:17px}.source-tabs{width:100%}.bottom-nav{grid-template-columns:repeat(5,1fr)}.bottom-nav button{font-size:9px}.folder-step{padding-left:45px}.step-no{left:11px}.home-kpis{grid-template-columns:repeat(2,minmax(0,1fr))}} @media(prefers-reduced-motion:reduce){*,*:before,*:after{scroll-behavior:auto!important;animation-duration:.001ms!important;animation-iteration-count:1!important;transition-duration:.001ms!important}.drop-zone:hover{transform:none}} + +/* resilient transcription chains */ +.chain-editor{display:flex;flex-direction:column;gap:10px;padding:14px;border:1px solid var(--line);border-radius:14px;background:#0c1513}.chain-head{display:flex;justify-content:space-between;gap:12px}.chain-head strong{font-size:13px}.chain-head p{margin:3px 0 0}.chain-route{display:grid;grid-template-columns:100px minmax(150px,.85fr) minmax(180px,1.15fr) minmax(130px,.55fr) auto;gap:10px;align-items:end;padding:12px;border:1px solid var(--line);border-radius:12px;background:#0e1816}.chain-route-label{align-self:center;display:flex;flex-direction:column;gap:3px}.chain-route-label strong{font-size:11px;color:var(--accent)}.chain-route-label span{font-size:9px;color:var(--muted-2)}.chain-actions{display:flex;gap:5px;align-items:end;padding-bottom:1px}.icon-btn{width:38px;min-width:38px;padding:0}.chain-add{align-self:flex-start}.chain-timeout .timeout-custom{display:grid;grid-template-columns:minmax(0,1fr) auto;gap:6px;align-items:center}.timeout-custom span{font-size:10px;color:var(--muted)}.error-inline{color:var(--danger)}.link-button{padding:0;border:0;background:transparent;color:var(--accent);font-weight:800;cursor:pointer}.attempt-details{margin:10px 0;border:1px solid var(--line);border-radius:10px;background:#0b1412}.attempt-details summary{cursor:pointer;padding:9px 11px;color:var(--muted);font-size:11px;font-weight:800}.attempt-list{display:flex;flex-direction:column;border-top:1px solid var(--line)}.attempt-row{display:grid;grid-template-columns:90px minmax(0,1fr);gap:7px 10px;padding:9px 11px;border-bottom:1px solid rgba(255,255,255,.04);font-size:10px;color:var(--muted)}.attempt-row:last-child{border-bottom:0}.attempt-outcome{font-weight:900;text-transform:uppercase;font-size:9px;letter-spacing:.05em}.attempt-outcome.success{color:var(--success)}.attempt-outcome.failed,.attempt-outcome.timed_out{color:var(--danger)}.attempt-outcome.cancelled{color:var(--warning)}.attempt-error{grid-column:2;color:var(--danger);overflow-wrap:anywhere} +@media(max-width:1080px){.chain-route{grid-template-columns:90px repeat(2,minmax(0,1fr));}.chain-timeout{grid-column:2}.chain-actions{grid-column:3;justify-content:flex-end}} +@media(max-width:680px){.chain-route{grid-template-columns:1fr}.chain-route-label,.chain-timeout,.chain-actions{grid-column:auto}.chain-actions{justify-content:flex-start}.chain-route-label{flex-direction:row;align-items:center;justify-content:space-between}.attempt-row{grid-template-columns:1fr}.attempt-error{grid-column:1}} diff --git a/frontend/src/lib/components/TranscriptionChainEditor.svelte b/frontend/src/lib/components/TranscriptionChainEditor.svelte new file mode 100644 index 0000000..faa2a57 --- /dev/null +++ b/frontend/src/lib/components/TranscriptionChainEditor.svelte @@ -0,0 +1,127 @@ + + +
+
+
Transcription chain

Choose the exact provider + model order. If one route fails, ScribeWatch tries the next.

+
+ {#each value as route,index (index)} +
+
{index===0?'Primary':`Fallback ${index}`}#{index+1}
+
+ + +
+
+ + + {#if errorsByProvider[route.providerId]} + Model discovery unavailable. + {/if} +
+
+ + + {#if timeoutChoice(route)==='custom'} +
customTimeoutChanged(index,event)} />min
+ {/if} +
+
+ + + +
+
+ {/each} + +
diff --git a/frontend/src/lib/transcription-chain.js b/frontend/src/lib/transcription-chain.js new file mode 100644 index 0000000..24beb2b --- /dev/null +++ b/frontend/src/lib/transcription-chain.js @@ -0,0 +1,30 @@ +// @ts-check +/** @typedef {import('./types').TranscriptionRoute} TranscriptionRoute */ + +/** @param {TranscriptionRoute[]} routes @returns {TranscriptionRoute[]} */ +export function addRoute(routes) { + return [...routes, { providerId: '', model: '' }]; +} + +/** @param {TranscriptionRoute[]} routes @param {number} index @returns {TranscriptionRoute[]} */ +export function removeRoute(routes, index) { + if (routes.length <= 1 || index < 0 || index >= routes.length) return [...routes]; + return routes.filter((_, current) => current !== index); +} + +/** @param {TranscriptionRoute[]} routes @param {number} from @param {number} to @returns {TranscriptionRoute[]} */ +export function moveRoute(routes, from, to) { + if (from < 0 || from >= routes.length || to < 0 || to >= routes.length || from === to) { + return [...routes]; + } + const next = [...routes]; + const [route] = next.splice(from, 1); + next.splice(to, 0, route); + return next; +} + +/** @param {TranscriptionRoute[]} routes @param {number} index @param {Partial} patch @returns {TranscriptionRoute[]} */ +export function updateRoute(routes, index, patch) { + if (index < 0 || index >= routes.length) return [...routes]; + return routes.map((route, current) => current === index ? { ...route, ...patch } : route); +} diff --git a/frontend/src/lib/views/JobsView.svelte b/frontend/src/lib/views/JobsView.svelte index 3339304..1c7ee40 100644 --- a/frontend/src/lib/views/JobsView.svelte +++ b/frontend/src/lib/views/JobsView.svelte @@ -1,7 +1,7 @@
-

HISTORY

Jobs

Every automatic and one-off transcription keeps a durable lifecycle, including retries and restart recovery.

+

HISTORY

Jobs

Every automatic and one-off transcription keeps a durable lifecycle, including retries, fallback attempts and restart recovery.

{#if visible.length===0}
No jobs in this view.
{/if} {#each visible as job} @@ -35,8 +39,22 @@ {#if job.archivePath}
Archive{job.archivePath}
{/if}
{#if job.error}
Processing error{job.error}
{/if} + {#if (job.transcriptionAttempts??[]).length>0} +
+ Transcription attempts ({job.transcriptionAttempts.length}) +
+ {#each job.transcriptionAttempts as attempt} +
+ {attempt.outcome.replace('_',' ')} + Run {attempt.run} · {attempt.providerName} / {attempt.model} · {elapsed(attempt)} + {#if attempt.error}{attempt.error}{/if} +
+ {/each} +
+
+ {/if}