diff --git a/apps/rustnzb/src/handlers.rs b/apps/rustnzb/src/handlers.rs index 785d230..472abc4 100644 --- a/apps/rustnzb/src/handlers.rs +++ b/apps/rustnzb/src/handlers.rs @@ -153,6 +153,7 @@ pub struct HistoryResponseEntry { pub output_dir: String, pub stages: Vec, pub error_message: Option, + pub failure_code: Option, pub server_stats: Vec, pub has_nzb_data: bool, pub duration_secs: f64, @@ -186,6 +187,7 @@ impl From for HistoryResponseEntry { output_dir: e.output_dir.to_string_lossy().to_string(), stages: e.stages, error_message: e.error_message, + failure_code: e.failure_code, server_stats: e.server_stats, has_nzb_data: has_nzb, duration_secs, diff --git a/crates/nzb-core/src/db.rs b/crates/nzb-core/src/db.rs index 1f25acf..0a93ee1 100644 --- a/crates/nzb-core/src/db.rs +++ b/crates/nzb-core/src/db.rs @@ -369,6 +369,21 @@ impl Database { )?; } + if version < 12 { + info!("Applying database migration v12: typed terminal failure codes"); + self.conn.execute_batch( + " + ALTER TABLE history ADD COLUMN failure_code TEXT; + UPDATE history + SET failure_code = 'download_failed' + WHERE LOWER(status) = 'failed' AND failure_code IS NULL; + + DELETE FROM schema_version; + INSERT INTO schema_version (version) VALUES (12); + ", + )?; + } + Ok(()) } @@ -504,46 +519,50 @@ impl Database { return Ok(None); }; - let queued_state = self + // The durable terminal history row is authoritative over the lingering + // queue view, so it is consulted first: a job that has completed or + // failed carries its typed failure code here. + let history_state = self .conn .query_row( - "SELECT status, total_bytes, downloaded_bytes, output_dir, error_message - FROM queue WHERE id = ?1", + "SELECT status, total_bytes, downloaded_bytes, completed_at, output_dir, + error_message, failure_code FROM history WHERE id = ?1", [&admission.job_id], |row| { - Ok(QueueAdmissionState::Queue { + Ok(QueueAdmissionState::History { status: parse_status(&row.get::<_, String>(0)?), total_bytes: row.get::<_, i64>(1)? as u64, downloaded_bytes: row.get::<_, i64>(2)? as u64, - output_dir: row.get::<_, String>(3)?.into(), - error_message: row.get(4)?, + completed_at: parse_datetime(&row.get::<_, String>(3)?), + output_dir: row.get::<_, String>(4)?.into(), + error_message: row.get(5)?, + failure_code: parse_failure_code(row.get(6)?)?, }) }, ) .optional()?; - if let Some(state) = queued_state { + if let Some(state) = history_state { return Ok(Some(QueueAdmissionObservation { admission, state })); } - let history_state = self + let queued_state = self .conn .query_row( - "SELECT status, total_bytes, downloaded_bytes, completed_at, output_dir, - error_message FROM history WHERE id = ?1", + "SELECT status, total_bytes, downloaded_bytes, output_dir, error_message + FROM queue WHERE id = ?1", [&admission.job_id], |row| { - Ok(QueueAdmissionState::History { + Ok(QueueAdmissionState::Queue { status: parse_status(&row.get::<_, String>(0)?), total_bytes: row.get::<_, i64>(1)? as u64, downloaded_bytes: row.get::<_, i64>(2)? as u64, - completed_at: parse_datetime(&row.get::<_, String>(3)?), - output_dir: row.get::<_, String>(4)?.into(), - error_message: row.get(5)?, + output_dir: row.get::<_, String>(3)?.into(), + error_message: row.get(4)?, }) }, ) .optional()?; - if let Some(state) = history_state { + if let Some(state) = queued_state { return Ok(Some(QueueAdmissionObservation { admission, state })); } @@ -684,12 +703,19 @@ impl Database { /// Move a completed/failed job to history. pub fn history_insert(&self, entry: &HistoryEntry) -> Result<(), NzbError> { + // Invariant: a Failed row carries exactly one typed failure code, and a + // non-Failed row carries none. This keeps the terminal history row a + // trustworthy source of the failure reason. + if (entry.status == JobStatus::Failed) != entry.failure_code.is_some() { + return Err(NzbError::Other("terminal_failure_code_invalid".to_string())); + } let stages_json = serde_json::to_string(&entry.stages).unwrap_or_default(); let server_stats_json = serde_json::to_string(&entry.server_stats).unwrap_or_default(); self.conn.execute( "INSERT INTO history (id, name, category, status, total_bytes, downloaded_bytes, - added_at, completed_at, download_time_secs, output_dir, stages, error_message, nzb_data, server_stats, retry_data) - VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10, ?11, ?12, ?13, ?14, ?15)", + added_at, completed_at, download_time_secs, output_dir, stages, error_message, + nzb_data, server_stats, retry_data, failure_code) + VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10, ?11, ?12, ?13, ?14, ?15, ?16)", params![ entry.id, entry.name, @@ -706,6 +732,7 @@ impl Database { entry.nzb_data, server_stats_json, entry.retry_data, + entry.failure_code.map(|code| code.to_string()), ], )?; @@ -770,7 +797,7 @@ impl Database { let mut stmt = self.conn.prepare( "SELECT id, name, category, status, total_bytes, downloaded_bytes, added_at, completed_at, download_time_secs, output_dir, stages, error_message, server_stats, - CASE WHEN nzb_data IS NOT NULL THEN 1 ELSE 0 END as has_nzb + CASE WHEN nzb_data IS NOT NULL THEN 1 ELSE 0 END as has_nzb, failure_code FROM history ORDER BY completed_at DESC LIMIT ?1", )?; @@ -797,6 +824,7 @@ impl Database { output_dir: row.get::<_, String>(9)?.into(), stages, error_message: row.get(11)?, + failure_code: parse_failure_code(row.get(14)?)?, server_stats, // Don't load actual blob in list - just note if it exists nzb_data: if has_nzb != 0 { Some(Vec::new()) } else { None }, @@ -859,7 +887,8 @@ impl Database { pub fn history_get(&self, id: &str) -> Result, NzbError> { let mut stmt = self.conn.prepare( "SELECT id, name, category, status, total_bytes, downloaded_bytes, - added_at, completed_at, download_time_secs, output_dir, stages, error_message, server_stats + added_at, completed_at, download_time_secs, output_dir, stages, error_message, server_stats, + failure_code FROM history WHERE id = ?1", )?; @@ -883,6 +912,7 @@ impl Database { output_dir: row.get::<_, String>(9)?.into(), stages, error_message: row.get(11)?, + failure_code: parse_failure_code(row.get(13)?)?, server_stats, nzb_data: None, retry_data: None, @@ -1251,6 +1281,24 @@ fn parse_status(s: &str) -> JobStatus { } } +fn parse_failure_code(value: Option) -> rusqlite::Result> { + value + .map(|value| { + value.parse().map_err(|_| { + rusqlite::Error::FromSqlConversionFailure( + 0, + rusqlite::types::Type::Text, + std::io::Error::new( + std::io::ErrorKind::InvalidData, + "invalid terminal failure code", + ) + .into(), + ) + }) + }) + .transpose() +} + fn parse_priority(v: i32) -> Priority { match v { 0 => Priority::Low, @@ -1316,6 +1364,7 @@ mod tests { duration_secs: 2.5, }], error_message: None, + failure_code: None, server_stats: Vec::new(), nzb_data: None, retry_data: None, @@ -1439,6 +1488,97 @@ mod tests { assert!(db.queue_list().unwrap().is_empty()); } + #[test] + fn typed_failure_migration_backfills_only_failed_history() { + let directory = tempfile::tempdir().unwrap(); + let path = directory.path().join("queue.db"); + let connection = Connection::open(&path).unwrap(); + connection + .execute_batch( + " + CREATE TABLE schema_version (version INTEGER NOT NULL); + INSERT INTO schema_version (version) VALUES (11); + CREATE TABLE history ( + id TEXT PRIMARY KEY, + status TEXT NOT NULL + ); + INSERT INTO history (id, status) VALUES + ('failed-job', 'Failed'), + ('completed-job', 'Completed'); + ", + ) + .unwrap(); + drop(connection); + + let db = Database::open(&path).unwrap(); + let failed: Option = db + .conn + .query_row( + "SELECT failure_code FROM history WHERE id = 'failed-job'", + [], + |row| row.get(0), + ) + .unwrap(); + let completed: Option = db + .conn + .query_row( + "SELECT failure_code FROM history WHERE id = 'completed-job'", + [], + |row| row.get(0), + ) + .unwrap(); + assert_eq!(failed.as_deref(), Some("download_failed")); + assert_eq!(completed, None); + } + + #[test] + fn queue_admission_prefers_typed_terminal_history_over_lingering_queue_view() { + let mut db = Database::open_memory().unwrap(); + db.queue_admit( + &make_job("terminal-job", "Terminal"), + b"payload", + "terminal-key", + "sha256:terminal", + ) + .unwrap(); + let mut history = make_history("terminal-job", "Terminal"); + history.status = JobStatus::Failed; + history.failure_code = Some(JobFailureCode::ArchiveInvalid); + history.error_message = Some("private extractor diagnostic".to_string()); + db.history_insert(&history).unwrap(); + + let observation = db.queue_admission_observe("terminal-key").unwrap().unwrap(); + assert!(matches!( + observation.state, + QueueAdmissionState::History { + status: JobStatus::Failed, + failure_code: Some(JobFailureCode::ArchiveInvalid), + .. + } + )); + } + + #[test] + fn terminal_history_requires_failure_code_exactly_on_failure() { + let db = Database::open_memory().unwrap(); + let mut failed = make_history("failed-job", "Failed"); + failed.status = JobStatus::Failed; + assert!(matches!( + db.history_insert(&failed), + Err(NzbError::Other(code)) if code == "terminal_failure_code_invalid" + )); + + failed.failure_code = Some(JobFailureCode::DownloadFailed); + db.history_insert(&failed).unwrap(); + + let mut completed = make_history("completed-job", "Completed"); + completed.failure_code = Some(JobFailureCode::DownloadFailed); + assert!(matches!( + db.history_insert(&completed), + Err(NzbError::Other(code)) if code == "terminal_failure_code_invalid" + )); + } + #[test] fn test_queue_update_progress() { let db = Database::open_memory().unwrap(); diff --git a/crates/nzb-core/src/models.rs b/crates/nzb-core/src/models.rs index de8d9c6..10c8d48 100644 --- a/crates/nzb-core/src/models.rs +++ b/crates/nzb-core/src/models.rs @@ -37,6 +37,50 @@ impl std::fmt::Display for JobStatus { } } +/// Engine-owned reason a job reached a terminal Failed state. +/// +/// These carry stable snake_case wire values so consumers key off a typed code +/// instead of parsing human-readable diagnostic prose. +#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)] +#[serde(rename_all = "snake_case")] +pub enum JobFailureCode { + ArticlesUnavailable, + RepairFailed, + ArchivePasswordRequired, + ArchiveInvalid, + StorageUnavailable, + DownloadFailed, +} + +impl std::fmt::Display for JobFailureCode { + fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { + match self { + Self::ArticlesUnavailable => write!(f, "articles_unavailable"), + Self::RepairFailed => write!(f, "repair_failed"), + Self::ArchivePasswordRequired => write!(f, "archive_password_required"), + Self::ArchiveInvalid => write!(f, "archive_invalid"), + Self::StorageUnavailable => write!(f, "storage_unavailable"), + Self::DownloadFailed => write!(f, "download_failed"), + } + } +} + +impl std::str::FromStr for JobFailureCode { + type Err = String; + + fn from_str(value: &str) -> Result { + match value { + "articles_unavailable" => Ok(Self::ArticlesUnavailable), + "repair_failed" => Ok(Self::RepairFailed), + "archive_password_required" => Ok(Self::ArchivePasswordRequired), + "archive_invalid" => Ok(Self::ArchiveInvalid), + "storage_unavailable" => Ok(Self::StorageUnavailable), + "download_failed" => Ok(Self::DownloadFailed), + _ => Err(value.to_string()), + } + } +} + // --------------------------------------------------------------------------- // Priority // --------------------------------------------------------------------------- @@ -199,6 +243,7 @@ pub struct HistoryEntry { /// Post-processing stages with results pub stages: Vec, pub error_message: Option, + pub failure_code: Option, /// Per-server download statistics #[serde(default)] pub server_stats: Vec, @@ -254,6 +299,7 @@ pub enum QueueAdmissionState { completed_at: DateTime, output_dir: PathBuf, error_message: Option, + failure_code: Option, }, Unobserved, } @@ -424,6 +470,21 @@ mod tests { assert_eq!(JobStatus::Failed.to_string(), "Failed"); } + #[test] + fn terminal_failure_codes_have_stable_wire_values() { + let code = JobFailureCode::ArchivePasswordRequired; + assert_eq!(code.to_string(), "archive_password_required"); + assert_eq!( + serde_json::to_string(&code).unwrap(), + "\"archive_password_required\"" + ); + assert_eq!( + "archive_password_required".parse(), + Ok(JobFailureCode::ArchivePasswordRequired) + ); + assert!("private extractor prose".parse::().is_err()); + } + #[test] fn test_job_status_serde_roundtrip() { let statuses = [ diff --git a/crates/nzb-postproc/src/pipeline.rs b/crates/nzb-postproc/src/pipeline.rs index 1f5b271..ee298c0 100644 --- a/crates/nzb-postproc/src/pipeline.rs +++ b/crates/nzb-postproc/src/pipeline.rs @@ -11,7 +11,7 @@ use std::path::{Path, PathBuf}; use std::sync::Arc; use std::time::Instant; -use nzb_core::models::{StageResult, StageStatus}; +use nzb_core::models::{JobFailureCode, StageResult, StageStatus}; use tracing::{debug, error, info, warn}; use crate::detect::{ @@ -19,7 +19,7 @@ use crate::detect::{ }; use crate::par2::par2_repair; use crate::resources::PostProcResourcePool; -use crate::unpack::{extract_7z, extract_rar, extract_tar, extract_zip}; +use crate::unpack::{ArchivePasswordRequired, extract_7z, extract_rar, extract_tar, extract_zip}; fn increment_counter(name: &'static str) { opentelemetry::global::meter_provider() @@ -114,6 +114,8 @@ pub struct PostProcResult { pub stages: Vec, /// Error message if the pipeline failed. pub error: Option, + /// Typed terminal failure code when the pipeline failed. + pub failure_code: Option, } /// Configuration for the post-processing pipeline. @@ -190,6 +192,7 @@ pub async fn run_pipeline_with_cleanup( ) -> PostProcResult { let mut stages: Vec = Vec::new(); let mut pipeline_ok = true; + let mut failure_code = None; info!(dir = %job_dir.display(), "Starting post-processing pipeline"); @@ -214,6 +217,7 @@ pub async fn run_pipeline_with_cleanup( if par2_files.is_empty() { if config.content_articles_failed > 0 { pipeline_ok = false; + failure_code = Some(JobFailureCode::ArticlesUnavailable); stages.push(StageResult { name: "Verify".to_string(), status: StageStatus::Failed, @@ -271,6 +275,7 @@ pub async fn run_pipeline_with_cleanup( }); if repair_result.status == StageStatus::Failed { pipeline_ok = false; + failure_code = Some(JobFailureCode::RepairFailed); } stages.push(repair_result); } @@ -381,6 +386,7 @@ pub async fn run_pipeline_with_cleanup( ); if !result.success { pipeline_ok = false; + failure_code = Some(JobFailureCode::RepairFailed); } stages.push(StageResult { name: "Repair".to_string(), @@ -409,6 +415,7 @@ pub async fn run_pipeline_with_cleanup( "Native PAR2 repair failed" ); pipeline_ok = false; + failure_code = Some(JobFailureCode::RepairFailed); stages.push(StageResult { name: "Repair".to_string(), status: StageStatus::Failed, @@ -426,6 +433,7 @@ pub async fn run_pipeline_with_cleanup( Err(e) => { error!(error = %e, "Verify/repair task panicked"); pipeline_ok = false; + failure_code = Some(JobFailureCode::RepairFailed); stages.push(StageResult { name: "Verify".to_string(), status: StageStatus::Failed, @@ -461,7 +469,7 @@ pub async fn run_pipeline_with_cleanup( } else { None }; - let (result, processed_archives) = run_extract_stage( + let (result, processed_archives, extract_failure_code) = run_extract_stage( source_dir, output_dir, config.password.as_deref(), @@ -471,6 +479,7 @@ pub async fn run_pipeline_with_cleanup( extracted_archives = processed_archives; if result.status == StageStatus::Failed { pipeline_ok = false; + failure_code = extract_failure_code.or(Some(JobFailureCode::ArchiveInvalid)); } else if result.status == StageStatus::Success { // Extraction succeeded despite verify/repair failure — recover pipeline_ok = true; @@ -515,6 +524,11 @@ pub async fn run_pipeline_with_cleanup( success: pipeline_ok, stages, error, + failure_code: if pipeline_ok { + None + } else { + failure_code.or(Some(JobFailureCode::DownloadFailed)) + }, } } @@ -659,7 +673,7 @@ async fn run_extract_stage( output_dir: &Path, password: Option<&str>, max_nested_archive_depth: u8, -) -> (StageResult, Vec) { +) -> (StageResult, Vec, Option) { let start = Instant::now(); let mut all_ok = true; let mut messages: Vec = Vec::new(); @@ -678,6 +692,7 @@ async fn run_extract_stage( let mut extracted_archives = Vec::new(); let mut scan_dir = source_dir; let mut extracted_any = false; + let mut failure_code = None; for depth in 0..=max_nested_archive_depth { let archives: Vec<_> = find_archives(scan_dir) @@ -705,6 +720,7 @@ async fn run_extract_stage( } Ok(unpack_result) => { all_ok = false; + failure_code = Some(JobFailureCode::ArchiveInvalid); let detail = unpack_result .error_output .trim() @@ -717,6 +733,11 @@ async fn run_extract_stage( } Err(e) => { all_ok = false; + failure_code = Some(if e.downcast_ref::().is_some() { + JobFailureCode::ArchivePasswordRequired + } else { + JobFailureCode::ArchiveInvalid + }); error!(depth, kind = %archive_type, file = %path.display(), error = %e, "Extraction error"); messages.push(format!("depth {depth} {archive_type}: {e}")); } @@ -735,6 +756,7 @@ async fn run_extract_stage( .all(|(_, path)| processed.contains(&path)) { all_ok = false; + failure_code = Some(JobFailureCode::ArchiveInvalid); messages.push(format!( "nested archive depth limit ({max_nested_archive_depth}) reached; source files retained" )); @@ -750,6 +772,7 @@ async fn run_extract_stage( duration_secs: start.elapsed().as_secs_f64(), }, Vec::new(), + None, ); } @@ -765,6 +788,7 @@ async fn run_extract_stage( duration_secs: start.elapsed().as_secs_f64(), }, extracted_archives, + failure_code, ) } @@ -948,10 +972,12 @@ mod tests { success: true, stages: vec![], error: None, + failure_code: None, }; assert!(result.success); assert!(result.stages.is_empty()); assert!(result.error.is_none()); + assert!(result.failure_code.is_none()); } #[test] @@ -1043,9 +1069,11 @@ mod tests { write_zip(&outer, &[("inner.zip", &fs::read(&inner).unwrap())]); fs::remove_file(inner).unwrap(); - let (result, _) = run_extract_stage(source.path(), output.path(), None, 1).await; + let (result, _, failure_code) = + run_extract_stage(source.path(), output.path(), None, 1).await; assert_eq!(result.status, StageStatus::Success, "{result:?}"); + assert_eq!(failure_code, None); assert_eq!( fs::read(output.path().join("payload.txt")).unwrap(), b"nested payload" @@ -1062,9 +1090,11 @@ mod tests { write_zip(&outer, &[("inner.zip", &fs::read(&inner).unwrap())]); fs::remove_file(inner).unwrap(); - let (result, _) = run_extract_stage(source.path(), output.path(), None, 0).await; + let (result, _, failure_code) = + run_extract_stage(source.path(), output.path(), None, 0).await; assert_eq!(result.status, StageStatus::Failed, "{result:?}"); + assert_eq!(failure_code, Some(JobFailureCode::ArchiveInvalid)); assert!(output.path().join("inner.zip").exists()); assert!(!output.path().join("payload.txt").exists()); } @@ -1208,6 +1238,10 @@ mod tests { }; let result = run_pipeline(dir.path(), &config).await; assert!(!result.success); + assert_eq!( + result.failure_code, + Some(JobFailureCode::ArticlesUnavailable) + ); let verify_stage = result.stages.iter().find(|s| s.name == "Verify").unwrap(); assert_eq!( diff --git a/crates/nzb-postproc/src/unpack.rs b/crates/nzb-postproc/src/unpack.rs index e01fba7..e63c7c3 100644 --- a/crates/nzb-postproc/src/unpack.rs +++ b/crates/nzb-postproc/src/unpack.rs @@ -12,6 +12,10 @@ use std::process::Stdio; use tokio::process::Command; use tracing::{info, warn}; +#[derive(Debug, thiserror::Error)] +#[error("archive password required")] +pub(crate) struct ArchivePasswordRequired; + /// Result of an unpack operation. #[derive(Debug)] pub struct UnpackResult { @@ -247,7 +251,7 @@ pub async fn extract_rar( file = %rar_file.display(), "RAR extraction failed — archive is password-protected" ); - anyhow::bail!("archive is password-protected"); + return Err(ArchivePasswordRequired.into()); } warn!( file = %rar_file.display(), @@ -318,7 +322,7 @@ pub async fn extract_7z( file = %archive_file.display(), "7z extraction failed — archive is password-protected" ); - anyhow::bail!("archive is password-protected"); + return Err(ArchivePasswordRequired.into()); } warn!( file = %archive_file.display(), @@ -789,4 +793,10 @@ mod tests { let args = sevenz_extract_args(archive, out, Some("secret")); assert!(args.iter().any(|arg| arg == "-psecret")); } + + #[test] + fn password_failure_survives_as_a_typed_error() { + let error: anyhow::Error = ArchivePasswordRequired.into(); + assert!(error.downcast_ref::().is_some()); + } } diff --git a/crates/nzb-web/src/queue_manager.rs b/crates/nzb-web/src/queue_manager.rs index d6eb8e7..93cca9b 100644 --- a/crates/nzb-web/src/queue_manager.rs +++ b/crates/nzb-web/src/queue_manager.rs @@ -891,6 +891,9 @@ struct JobState { hopeless_tracker: Option, /// Active worker-pool duration captured at terminal download resolution. download_time_secs: Option, + /// Typed terminal failure code, tracked in memory until it is persisted to + /// the terminal history row. + failure_code: Option, } /// Notification fired immediately when a job is accepted into the queue. @@ -1515,6 +1518,7 @@ impl QueueManager { direct_unpacker: None, hopeless_tracker: None, download_time_secs: None, + failure_code: None, }; self.jobs.lock().insert(job_id.clone(), state); self.globally_paused_jobs.lock().insert(job_id.clone()); @@ -1534,6 +1538,7 @@ impl QueueManager { direct_unpacker: None, hopeless_tracker: None, download_time_secs: None, + failure_code: None, }; self.jobs.lock().insert(job_id.clone(), state); self.job_order.lock().push(job_id); @@ -1648,6 +1653,7 @@ impl QueueManager { if let Some(state) = jobs.get_mut(job_id) { state.job.status = JobStatus::Paused; state.job.error_message = Some("Paused: low disk space".to_string()); + state.failure_code = Some(JobFailureCode::StorageUnavailable); } return; } @@ -1758,6 +1764,7 @@ impl QueueManager { }) { state.job.error_message = None; + state.failure_code = None; cleared = true; } } @@ -2184,6 +2191,7 @@ impl QueueManager { state.download_time_secs = Some(download_time_secs); state.job.status = JobStatus::Failed; state.job.error_message = Some(reason.clone()); + state.failure_code = Some(JobFailureCode::ArticlesUnavailable); state.job.articles_failed = state.job.articles_failed.max(articles_failed); state.job.completed_at = Some(chrono::Utc::now()); @@ -2403,6 +2411,8 @@ impl QueueManager { { state.job.status = JobStatus::Failed; state.job.error_message = result.error.clone(); + state.failure_code = + result.failure_code.or(Some(JobFailureCode::DownloadFailed)); } } @@ -2416,6 +2426,7 @@ impl QueueManager { state.job.status = JobStatus::Failed; state.job.error_message = Some(format!("{articles_failed} article(s) failed to download")); + state.failure_code = Some(JobFailureCode::ArticlesUnavailable); } } Vec::new() @@ -2616,6 +2627,7 @@ impl QueueManager { warn!(job_id = %state.job.id, output_dir = %state.job.output_dir.display(), "{message}"); final_status = JobStatus::Failed; state.job.error_message = Some(message.to_string()); + state.failure_code = Some(JobFailureCode::ArchiveInvalid); stages.push(StageResult { name: "Output".to_string(), status: StageStatus::Failed, @@ -2701,6 +2713,8 @@ impl QueueManager { output_dir: state.job.output_dir.clone(), stages, error_message: state.job.error_message.clone(), + failure_code: (final_status == JobStatus::Failed) + .then(|| state.failure_code.unwrap_or(JobFailureCode::DownloadFailed)), server_stats: state.job.server_stats.clone(), nzb_data: state.nzb_data.clone(), retry_data, @@ -3045,6 +3059,7 @@ impl QueueManager { // Job context still lives in the pool — just unpause it. state.job.status = JobStatus::Downloading; state.job.error_message = None; + state.failure_code = None; // The no-progress watchdog measures wall-clock idle time and // does not stop while paused, so restart its clock here or the // paused interval counts toward the article timeout and aborts @@ -3072,6 +3087,7 @@ impl QueueManager { } else { state.job.status = JobStatus::Downloading; state.job.error_message = None; + state.failure_code = None; true } } @@ -3140,6 +3156,9 @@ impl QueueManager { output_dir: state.job.output_dir.clone(), stages: Vec::new(), error_message: state.job.error_message.clone(), + failure_code: Some( + state.failure_code.unwrap_or(JobFailureCode::DownloadFailed), + ), server_stats: state.job.server_stats.clone(), nzb_data: state.nzb_data.clone(), retry_data: None, @@ -3341,6 +3360,7 @@ impl QueueManager { }; if state.job.status == JobStatus::Paused { state.job.error_message = None; + state.failure_code = None; if self.dispatch.has_job(&id) { state.job.status = JobStatus::Downloading; to_unpause.push(id); @@ -3399,6 +3419,7 @@ impl QueueManager { for (id, state) in jobs.iter_mut() { if state.job.status == JobStatus::Paused && state.job.error_message.is_some() { state.job.error_message = None; + state.failure_code = None; if self.dispatch.has_job(id) { state.job.status = JobStatus::Downloading; to_unpause.push(id.clone()); @@ -4089,6 +4110,7 @@ impl QueueManager { direct_unpacker: None, hopeless_tracker: None, download_time_secs: None, + failure_code: None, }; self.jobs.lock().insert(job_id.clone(), state); self.job_order.lock().push(job_id); @@ -4441,6 +4463,7 @@ mod global_pause_tests { direct_unpacker: None, hopeless_tracker: None, download_time_secs: None, + failure_code: None, }, ); manager.job_order.lock().push(id); @@ -4798,6 +4821,16 @@ mod global_pause_tests { .is_some() ); assert!(!work_dir.exists()); + assert_eq!( + manager + .db + .lock() + .history_get("failed-cleanup") + .unwrap() + .unwrap() + .failure_code, + Some(JobFailureCode::DownloadFailed) + ); } #[tokio::test] @@ -4856,6 +4889,7 @@ mod global_pause_tests { .unwrap() .unwrap(); assert_eq!(entry.status, JobStatus::Failed); + assert_eq!(entry.failure_code, Some(JobFailureCode::ArchiveInvalid)); assert_eq!( entry.error_message.as_deref(), Some("No usable output produced; only archive or PAR2 artifacts remain") diff --git a/crates/nzb-web/src/sabnzbd_compat.rs b/crates/nzb-web/src/sabnzbd_compat.rs index f2939e2..c0cf737 100644 --- a/crates/nzb-web/src/sabnzbd_compat.rs +++ b/crates/nzb-web/src/sabnzbd_compat.rs @@ -2093,6 +2093,7 @@ mod tests { duration_secs: 3.6, }], error_message: (status == JobStatus::Failed).then(|| "broken archive".into()), + failure_code: (status == JobStatus::Failed).then_some(JobFailureCode::ArchiveInvalid), server_stats: Vec::new(), nzb_data: (status == JobStatus::Failed).then(Vec::new), retry_data: None, @@ -2274,6 +2275,7 @@ mod tests { output_dir: "/downloads/complete".into(), stages: Vec::new(), error_message: None, + failure_code: None, server_stats: Vec::new(), nzb_data: None, retry_data: None, @@ -2521,6 +2523,7 @@ mod tests { .join("contract-history-job"), stages: Vec::new(), error_message: None, + failure_code: None, server_stats: Vec::new(), nzb_data: None, retry_data: None, @@ -2909,6 +2912,7 @@ mod tests { output_dir, stages: Vec::new(), error_message: None, + failure_code: None, server_stats: Vec::new(), nzb_data: None, retry_data: None,