From f583440650edb09b1ca4aaed8324ee498a18d622 Mon Sep 17 00:00:00 2001 From: thedancingdeveloper <306930456+thedancingdeveloper@users.noreply.github.com> Date: Sun, 27 Sep 2026 09:53:24 +0000 Subject: [PATCH 1/3] feat(nzb-web): idempotent queue admission via Idempotency-Key Add a durable, at-most-once admission path for queue adds so a caller that retries an ambiguous `POST /api/queue/add` cannot create a duplicate job. - New `queue_admissions` table (DB migration v11) keyed by idempotency key, binding it to a payload digest and job id. It intentionally has no foreign key to queue/history, so it stays authoritative after a job moves or bounded history is pruned. - `Idempotency-Key` request header (parsed/validated, path-safe, <=128 bytes) and SHA-256 payload digest (`apps/rustnzb/src/admissions.rs`), plus `GET /api/queue/admissions/{key}` to resolve one admission and its current engine location (queue / history / unobserved). - `Database::queue_admit` (transactional bind + queue insert) and `queue_admission_observe`; `QueueManager::add_job_idempotent` / `queue_admission_observe`, with `add_job` refactored so the shared activation path is `activate_admitted_job`. - New `NzbError::AdmissionConflict` -> HTTP 409 `admission_conflict` when a key is reused with a different payload. A keyed add must carry exactly one NZB. Ported from MrVampy/rustnzb (fe7c35d), adapted to main: migration renumbered v9 -> v11 (main's v9/v10 are retry_data / damage_ledger); the fork's flake.nix / workspace version bumps are excluded; digests use `hex::encode` because main is on sha2 0.11 (whose output no longer implements LowerHex); and the integration test now completes first-run auth setup and sends a bearer token, since main gates `/api` behind auth middleware. Tests: db `queue_admission_*` unit tests (replay/conflict, reopen/outlive); `admissions` key + digest unit tests; and an end-to-end `idempotent_admission` suite (exact replay returns one job, conflicting payload 409s, keyed multi-payload upload 400s). `cargo test` across nzb-core, nzb-web and rustnzb all pass; clippy/fmt clean. Co-Authored-By: MrVampy <4302946+MrVampy@users.noreply.github.com> Co-Authored-By: Claude Opus 4.8 --- Cargo.lock | 2 + apps/rustnzb/Cargo.toml | 2 + apps/rustnzb/src/admissions.rs | 90 ++++++++ apps/rustnzb/src/handlers.rs | 106 +++++---- apps/rustnzb/src/lib.rs | 1 + apps/rustnzb/src/server.rs | 7 +- apps/rustnzb/tests/idempotent_admission.rs | 182 +++++++++++++++ crates/nzb-core/src/db.rs | 256 ++++++++++++++++++++- crates/nzb-core/src/error.rs | 3 + crates/nzb-core/src/models.rs | 46 ++++ crates/nzb-web/src/error.rs | 10 + crates/nzb-web/src/queue_manager.rs | 65 +++++- 12 files changed, 725 insertions(+), 45 deletions(-) create mode 100644 apps/rustnzb/src/admissions.rs create mode 100644 apps/rustnzb/tests/idempotent_admission.rs diff --git a/Cargo.lock b/Cargo.lock index f3110170..a1393a00 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -2816,6 +2816,7 @@ dependencies = [ "clap", "crc32fast", "flate2", + "hex", "http", "libc", "mime_guess", @@ -2839,6 +2840,7 @@ dependencies = [ "serde", "serde_json", "serial_test", + "sha2 0.11.0", "tempfile", "tokio", "tokio-util", diff --git a/apps/rustnzb/Cargo.toml b/apps/rustnzb/Cargo.toml index 28f5565c..a48273db 100644 --- a/apps/rustnzb/Cargo.toml +++ b/apps/rustnzb/Cargo.toml @@ -42,6 +42,8 @@ reqwest = { workspace = true, features = ["json"] } regex = { workspace = true } uuid = { workspace = true } base64 = { workspace = true } +sha2 = "0.11" +hex = { workspace = true } rust-embed = { version = "8", features = ["debug-embed", "interpolate-folder-path"] } mime_guess = "2" libc = "0.2" diff --git a/apps/rustnzb/src/admissions.rs b/apps/rustnzb/src/admissions.rs new file mode 100644 index 00000000..e0d350c8 --- /dev/null +++ b/apps/rustnzb/src/admissions.rs @@ -0,0 +1,90 @@ +use std::sync::Arc; + +use axum::Json; +use axum::extract::{Path, State}; +use axum::http::HeaderMap; +use sha2::{Digest, Sha256}; + +use nzb_web::error::ApiError; +use nzb_web::nzb_core::models::QueueAdmissionObservation; +use nzb_web::state::AppState; + +pub const IDEMPOTENCY_KEY_HEADER: &str = "idempotency-key"; +const MAX_IDEMPOTENCY_KEY_BYTES: usize = 128; + +#[derive(Debug, Clone)] +pub struct IdempotencyKey(String); + +impl IdempotencyKey { + pub fn from_headers(headers: &HeaderMap) -> Result, ApiError> { + let mut values = headers.get_all(IDEMPOTENCY_KEY_HEADER).iter(); + let Some(value) = values.next() else { + return Ok(None); + }; + if values.next().is_some() { + return Err(ApiError::bad_request( + "Idempotency-Key must be supplied exactly once", + )); + } + let value = value + .to_str() + .map_err(|_| ApiError::bad_request("Idempotency-Key is not valid ASCII"))?; + Self::parse(value).map(Some) + } + + pub fn parse(value: &str) -> Result { + if value.is_empty() + || value.len() > MAX_IDEMPOTENCY_KEY_BYTES + || !value.bytes().all(|byte| { + byte.is_ascii_alphanumeric() || matches!(byte, b'-' | b'.' | b'_' | b':') + }) + { + return Err(ApiError::bad_request("Idempotency-Key is invalid")); + } + Ok(Self(value.to_string())) + } + + pub fn as_str(&self) -> &str { + &self.0 + } +} + +pub fn payload_digest(payload: &[u8]) -> String { + format!("sha256:{}", hex::encode(Sha256::digest(payload))) +} + +/// GET /api/queue/admissions/{idempotency_key} -- Resolve one exact admission. +pub async fn h_queue_admission_get( + State(state): State>, + Path(idempotency_key): Path, +) -> Result, ApiError> { + let idempotency_key = IdempotencyKey::parse(&idempotency_key)?; + let observation = state + .queue_manager + .queue_admission_observe(idempotency_key.as_str()) + .map_err(ApiError::from)? + .ok_or_else(|| ApiError::not_found("Queue admission not found"))?; + Ok(Json(observation)) +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn keys_are_bounded_and_path_safe() { + assert!(IdempotencyKey::parse("acquisition:018f-abc_DEF.2").is_ok()); + assert!(IdempotencyKey::parse("").is_err()); + assert!(IdempotencyKey::parse("contains/slash").is_err()); + assert!(IdempotencyKey::parse(&"a".repeat(129)).is_err()); + } + + #[test] + fn digest_binds_exact_payload_bytes() { + assert_eq!( + payload_digest(b"nzb"), + "sha256:5099941fc6e5440244a41b3f6e466d8933f73ea1647f042bf051b380a43acdcc" + ); + assert_ne!(payload_digest(b"nzb"), payload_digest(b"NZB")); + } +} diff --git a/apps/rustnzb/src/handlers.rs b/apps/rustnzb/src/handlers.rs index e1f4f208..785d2306 100644 --- a/apps/rustnzb/src/handlers.rs +++ b/apps/rustnzb/src/handlers.rs @@ -3,6 +3,7 @@ use std::sync::Arc; use axum::Json; use axum::extract::{Multipart, Path, Query, State}; +use axum::http::HeaderMap; use axum::response::IntoResponse; use flate2::read::GzDecoder; use http::StatusCode; @@ -38,6 +39,8 @@ use nzb_web::fetch_guard::{ use nzb_web::log_buffer::LogEntry; use nzb_web::state::AppState; +use crate::admissions::{IdempotencyKey, payload_digest}; + // --------------------------------------------------------------------------- // Priority helpers // --------------------------------------------------------------------------- @@ -359,11 +362,33 @@ fn extract_nzbs(file_name: &str, data: &[u8]) -> Result)>, } /// Enqueue a single NZB from raw bytes, applying category/priority from query params. +async fn next_uploaded_file( + multipart: &mut Multipart, +) -> Result)>, ApiError> { + let Some(field) = multipart + .next_field() + .await + .map_err(|error| ApiError::from(anyhow::anyhow!("Multipart error: {error}")))? + else { + return Ok(None); + }; + let file_name = field + .file_name() + .map(str::to_string) + .unwrap_or_else(|| "unknown.nzb".to_string()); + let data = field + .bytes() + .await + .map_err(|error| ApiError::from(anyhow::anyhow!("Read error: {error}")))?; + Ok(Some((file_name, data.to_vec()))) +} + fn enqueue_nzb( state: &AppState, q: &AddNzbQuery, file_name: &str, data: Vec, + idempotency_key: Option<&IdempotencyKey>, ) -> Result { let name = q.name.clone().unwrap_or_else(|| { file_name @@ -387,24 +412,21 @@ fn enqueue_nzb( .output_dir_for(&job.category, &job.name) .map_err(ApiError::from)?; - std::fs::create_dir_all(&job.work_dir).map_err(|e| { - ApiError::from(anyhow::anyhow!( - "Failed to create work dir '{}': {}", - job.work_dir.display(), - e - )) - })?; + if let Some(idempotency_key) = idempotency_key { + let digest = payload_digest(&data); + let outcome = qm + .add_job_idempotent(job, data, idempotency_key.as_str(), &digest) + .map_err(|error| match error { + nzb_web::nzb_core::NzbError::AdmissionConflict => ApiError::admission_conflict(), + error => ApiError::from(error), + })?; + return Ok(match outcome { + QueueAdmissionOutcome::Inserted(admission) + | QueueAdmissionOutcome::Existing(admission) => admission.job_id, + }); + } let id = job.id.clone(); - - tracing::info!( - name = %job.name, - id = %job.id, - files = job.file_count, - articles = job.article_count, - "NZB added to queue" - ); - qm.add_job(job, Some(data)).map_err(ApiError::from)?; Ok(id) } @@ -415,31 +437,37 @@ fn enqueue_nzb( pub async fn h_queue_add( State(state): State>, Query(q): Query, + headers: HeaderMap, mut multipart: Multipart, ) -> Result { + let idempotency_key = IdempotencyKey::from_headers(&headers)?; let mut nzo_ids = Vec::new(); - - while let Some(field) = multipart - .next_field() - .await - .map_err(|e| ApiError::from(anyhow::anyhow!("Multipart error: {e}")))? - { - let file_name = field - .file_name() - .map(|s| s.to_string()) - .unwrap_or_else(|| "unknown.nzb".into()); - - let data = field - .bytes() - .await - .map_err(|e| ApiError::from(anyhow::anyhow!("Read error: {e}")))?; - - // Extract NZBs (handles zip/gz archives or plain .nzb) - let nzbs = extract_nzbs(&file_name, &data).map_err(ApiError::from)?; - - for (nzb_name, nzb_data) in nzbs { - let id = enqueue_nzb(&state, &q, &nzb_name, nzb_data)?; - nzo_ids.push(id); + if let Some(idempotency_key) = idempotency_key.as_ref() { + // A keyed admission binds exactly one payload, so the request must carry + // exactly one NZB — a multi-NZB upload has no single job to replay. + let mut uploaded_nzbs = Vec::new(); + while let Some((file_name, data)) = next_uploaded_file(&mut multipart).await? { + uploaded_nzbs.extend(extract_nzbs(&file_name, &data).map_err(ApiError::from)?); + } + if uploaded_nzbs.len() != 1 { + return Err(ApiError::bad_request( + "Idempotency-Key requires exactly one NZB payload", + )); + } + let (nzb_name, nzb_data) = uploaded_nzbs.pop().expect("exactly one NZB payload"); + nzo_ids.push(enqueue_nzb( + &state, + &q, + &nzb_name, + nzb_data, + Some(idempotency_key), + )?); + } else { + while let Some((file_name, data)) = next_uploaded_file(&mut multipart).await? { + // Extract NZBs (handles zip/gz archives or plain .nzb) + for (nzb_name, nzb_data) in extract_nzbs(&file_name, &data).map_err(ApiError::from)? { + nzo_ids.push(enqueue_nzb(&state, &q, &nzb_name, nzb_data, None)?); + } } } @@ -543,7 +571,7 @@ pub async fn h_queue_add_url( let nzbs = extract_nzbs(&file_name, &data).map_err(ApiError::from)?; let mut nzo_ids = Vec::new(); for (nzb_name, nzb_data) in nzbs { - let id = enqueue_nzb(&state, &q, &nzb_name, nzb_data)?; + let id = enqueue_nzb(&state, &q, &nzb_name, nzb_data, None)?; nzo_ids.push(id); } diff --git a/apps/rustnzb/src/lib.rs b/apps/rustnzb/src/lib.rs index a6905f14..2e8f8be1 100644 --- a/apps/rustnzb/src/lib.rs +++ b/apps/rustnzb/src/lib.rs @@ -1,3 +1,4 @@ +pub mod admissions; pub mod group_handlers; pub mod handlers; pub mod server; diff --git a/apps/rustnzb/src/server.rs b/apps/rustnzb/src/server.rs index b6c6b78f..ae3748c6 100644 --- a/apps/rustnzb/src/server.rs +++ b/apps/rustnzb/src/server.rs @@ -17,8 +17,7 @@ use tracing::info; use utoipa::OpenApi; use utoipa_swagger_ui::SwaggerUi; -use crate::group_handlers; -use crate::handlers; +use crate::{admissions, group_handlers, handlers}; use nzb_web::auth; use nzb_web::error::ApiError; use nzb_web::sabnzbd_compat; @@ -155,6 +154,10 @@ pub fn build_router(state: Arc) -> Router { .route("/queue", get(handlers::h_queue_list)) .route("/queue/add", post(handlers::h_queue_add)) .route("/queue/add-url", post(handlers::h_queue_add_url)) + .route( + "/queue/admissions/{idempotency_key}", + get(admissions::h_queue_admission_get), + ) .route("/queue/pause", post(handlers::h_queue_pause_all)) .route("/queue/resume", post(handlers::h_queue_resume_all)) .route("/queue/pause-for", post(handlers::h_queue_pause_for)) diff --git a/apps/rustnzb/tests/idempotent_admission.rs b/apps/rustnzb/tests/idempotent_admission.rs new file mode 100644 index 00000000..f0b5ef48 --- /dev/null +++ b/apps/rustnzb/tests/idempotent_admission.rs @@ -0,0 +1,182 @@ +mod support; + +use reqwest::StatusCode; +use rustnzb::admissions::payload_digest; +use support::{sample_nzb_bytes, sample_nzb_variant_bytes, start_test_server}; + +const ADMISSION_KEY: &str = "newsgroups-acquisition-018f"; + +/// Complete first-run auth setup and return a bearer access token. Main gates +/// the `/api` routes behind auth middleware, so every request below carries it. +async fn setup_auth(client: &reqwest::Client, base_url: &str) -> String { + let setup = client + .post(format!("{base_url}/api/auth/setup")) + .json(&serde_json::json!({ + "username": "admission-test", + "password": "admission-test-password" + })) + .send() + .await + .expect("auth setup failed"); + assert_eq!(setup.status(), StatusCode::OK); + setup.json::().await.unwrap()["access_token"] + .as_str() + .expect("auth setup should return an access token") + .to_string() +} + +async fn upload( + client: &reqwest::Client, + base_url: &str, + access: &str, + key: &str, + name: &str, + payload: Vec, +) -> reqwest::Response { + let part = reqwest::multipart::Part::bytes(payload) + .file_name(name.to_string()) + .mime_str("application/x-nzb") + .unwrap(); + client + .post(format!("{base_url}/api/queue/add")) + .bearer_auth(access) + .header("Idempotency-Key", key) + .multipart(reqwest::multipart::Form::new().part("file", part)) + .send() + .await + .unwrap() +} + +#[tokio::test] +async fn exact_replay_returns_one_job_and_conflicting_payload_fails_closed() { + let app = start_test_server(Vec::new()).await; + let client = reqwest::Client::new(); + let access = setup_auth(&client, &app.base_url).await; + client + .post(format!("{}/api/queue/pause", app.base_url)) + .bearer_auth(&access) + .send() + .await + .unwrap() + .error_for_status() + .unwrap(); + + let payload = sample_nzb_bytes(); + let first = upload( + &client, + &app.base_url, + &access, + ADMISSION_KEY, + "first.nzb", + payload.clone(), + ) + .await; + assert_eq!(first.status(), StatusCode::OK); + let first = first.json::().await.unwrap(); + let first_job = first["nzo_ids"][0].as_str().unwrap(); + + let replay = upload( + &client, + &app.base_url, + &access, + ADMISSION_KEY, + "renamed.nzb", + payload.clone(), + ) + .await; + assert_eq!(replay.status(), StatusCode::OK); + let replay = replay.json::().await.unwrap(); + assert_eq!(replay["nzo_ids"][0], first_job); + + let queue = client + .get(format!("{}/api/queue", app.base_url)) + .bearer_auth(&access) + .send() + .await + .unwrap() + .json::() + .await + .unwrap(); + assert_eq!(queue["total"], 1); + + let observation = client + .get(format!( + "{}/api/queue/admissions/{ADMISSION_KEY}", + app.base_url + )) + .bearer_auth(&access) + .send() + .await + .unwrap(); + assert_eq!(observation.status(), StatusCode::OK); + let observation = observation.json::().await.unwrap(); + assert_eq!(observation["admission"]["idempotency_key"], ADMISSION_KEY); + assert_eq!(observation["admission"]["job_id"], first_job); + assert_eq!( + observation["admission"]["payload_digest"], + payload_digest(&payload) + ); + assert_eq!(observation["state"]["location"], "queue"); + assert_eq!(observation["state"]["status"], "queued"); + + let conflict = upload( + &client, + &app.base_url, + &access, + ADMISSION_KEY, + "different.nzb", + sample_nzb_variant_bytes(), + ) + .await; + assert_eq!(conflict.status(), StatusCode::CONFLICT); + let conflict = conflict.json::().await.unwrap(); + assert_eq!(conflict["error_kind"], "admission_conflict"); + + let queue = client + .get(format!("{}/api/queue", app.base_url)) + .bearer_auth(&access) + .send() + .await + .unwrap() + .json::() + .await + .unwrap(); + assert_eq!(queue["total"], 1); +} + +#[tokio::test] +async fn keyed_multi_payload_upload_is_rejected_before_admission() { + let app = start_test_server(Vec::new()).await; + let client = reqwest::Client::new(); + let access = setup_auth(&client, &app.base_url).await; + let form = reqwest::multipart::Form::new() + .part( + "first", + reqwest::multipart::Part::bytes(sample_nzb_bytes()).file_name("first.nzb"), + ) + .part( + "second", + reqwest::multipart::Part::bytes(sample_nzb_variant_bytes()).file_name("second.nzb"), + ); + + let response = client + .post(format!("{}/api/queue/add", app.base_url)) + .bearer_auth(&access) + .header("Idempotency-Key", ADMISSION_KEY) + .multipart(form) + .send() + .await + .unwrap(); + assert_eq!(response.status(), StatusCode::BAD_REQUEST); + + let missing = client + .get(format!( + "{}/api/queue/admissions/{ADMISSION_KEY}", + app.base_url + )) + .bearer_auth(&access) + .send() + .await + .unwrap(); + assert_eq!(missing.status(), StatusCode::NOT_FOUND); +} diff --git a/crates/nzb-core/src/db.rs b/crates/nzb-core/src/db.rs index ca128834..1f25acf8 100644 --- a/crates/nzb-core/src/db.rs +++ b/crates/nzb-core/src/db.rs @@ -1,7 +1,7 @@ use std::path::Path; use chrono::Utc; -use rusqlite::{Connection, params}; +use rusqlite::{Connection, OptionalExtension, params}; use tracing::info; use crate::error::NzbError; @@ -349,6 +349,26 @@ impl Database { )?; } + if version < 11 { + info!("Applying database migration v11: idempotent queue admissions"); + self.conn.execute_batch( + " + CREATE TABLE queue_admissions ( + idempotency_key TEXT PRIMARY KEY, + payload_digest TEXT NOT NULL, + job_id TEXT NOT NULL UNIQUE, + accepted_at TEXT NOT NULL + ); + + CREATE INDEX idx_queue_admissions_job + ON queue_admissions(job_id); + + DELETE FROM schema_version; + INSERT INTO schema_version (version) VALUES (11); + ", + )?; + } + Ok(()) } @@ -374,6 +394,165 @@ impl Database { // Queue operations // ----------------------------------------------------------------------- + /// Atomically bind an idempotency key to one payload and queued job. + /// + /// The admission row intentionally has no foreign key to queue or history: + /// it remains authoritative after a job moves or bounded history is pruned. + pub fn queue_admit( + &mut self, + job: &NzbJob, + nzb_data: &[u8], + idempotency_key: &str, + payload_digest: &str, + ) -> Result { + let transaction = self.conn.transaction()?; + let existing = transaction + .query_row( + "SELECT idempotency_key, payload_digest, job_id, accepted_at + FROM queue_admissions WHERE idempotency_key = ?1", + [idempotency_key], + |row| { + Ok(QueueAdmission { + idempotency_key: row.get(0)?, + payload_digest: row.get(1)?, + job_id: row.get(2)?, + accepted_at: parse_datetime(&row.get::<_, String>(3)?), + }) + }, + ) + .optional()?; + + if let Some(existing) = existing { + if existing.payload_digest != payload_digest { + return Err(NzbError::AdmissionConflict); + } + transaction.commit()?; + return Ok(QueueAdmissionOutcome::Existing(existing)); + } + + transaction.execute( + "INSERT INTO queue (id, name, category, status, priority, total_bytes, + downloaded_bytes, file_count, files_completed, article_count, + articles_downloaded, articles_failed, added_at, work_dir, output_dir, password, + nzb_raw) + VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10, ?11, ?12, ?13, ?14, ?15, + ?16, ?17)", + params![ + job.id, + job.name, + job.category, + job.status.to_string(), + job.priority as i32, + job.total_bytes as i64, + job.downloaded_bytes as i64, + job.file_count as i64, + job.files_completed as i64, + job.article_count as i64, + job.articles_downloaded as i64, + job.articles_failed as i64, + job.added_at.to_rfc3339(), + job.work_dir.to_string_lossy().to_string(), + job.output_dir.to_string_lossy().to_string(), + job.password, + nzb_data, + ], + )?; + + let admission = QueueAdmission { + idempotency_key: idempotency_key.to_string(), + payload_digest: payload_digest.to_string(), + job_id: job.id.clone(), + accepted_at: job.added_at, + }; + transaction.execute( + "INSERT INTO queue_admissions ( + idempotency_key, payload_digest, job_id, accepted_at + ) VALUES (?1, ?2, ?3, ?4)", + params![ + admission.idempotency_key, + admission.payload_digest, + admission.job_id, + admission.accepted_at.to_rfc3339(), + ], + )?; + transaction.commit()?; + Ok(QueueAdmissionOutcome::Inserted(admission)) + } + + /// Resolve one admission without relying on bounded queue/history lists. + pub fn queue_admission_observe( + &self, + idempotency_key: &str, + ) -> Result, NzbError> { + let admission = self + .conn + .query_row( + "SELECT idempotency_key, payload_digest, job_id, accepted_at + FROM queue_admissions WHERE idempotency_key = ?1", + [idempotency_key], + |row| { + Ok(QueueAdmission { + idempotency_key: row.get(0)?, + payload_digest: row.get(1)?, + job_id: row.get(2)?, + accepted_at: parse_datetime(&row.get::<_, String>(3)?), + }) + }, + ) + .optional()?; + let Some(admission) = admission else { + return Ok(None); + }; + + let queued_state = self + .conn + .query_row( + "SELECT status, total_bytes, downloaded_bytes, output_dir, error_message + FROM queue WHERE id = ?1", + [&admission.job_id], + |row| { + 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, + output_dir: row.get::<_, String>(3)?.into(), + error_message: row.get(4)?, + }) + }, + ) + .optional()?; + if let Some(state) = queued_state { + return Ok(Some(QueueAdmissionObservation { admission, state })); + } + + let history_state = self + .conn + .query_row( + "SELECT status, total_bytes, downloaded_bytes, completed_at, output_dir, + error_message FROM history WHERE id = ?1", + [&admission.job_id], + |row| { + 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, + completed_at: parse_datetime(&row.get::<_, String>(3)?), + output_dir: row.get::<_, String>(4)?.into(), + error_message: row.get(5)?, + }) + }, + ) + .optional()?; + if let Some(state) = history_state { + return Ok(Some(QueueAdmissionObservation { admission, state })); + } + + Ok(Some(QueueAdmissionObservation { + admission, + state: QueueAdmissionState::Unobserved, + })) + } + /// Insert a new job into the queue. pub fn queue_insert(&self, job: &NzbJob) -> Result<(), NzbError> { self.conn.execute( @@ -1185,6 +1364,81 @@ mod tests { assert_eq!(jobs[0].total_bytes, 1_000_000); } + #[test] + fn queue_admission_replays_exact_payload_and_rejects_conflict() { + let mut db = Database::open_memory().unwrap(); + let job = make_job("admitted-job", "Admitted"); + let inserted = db + .queue_admit(&job, b"payload", "request-key", "sha256:first") + .unwrap(); + assert!(matches!(inserted, QueueAdmissionOutcome::Inserted(_))); + + let replay = db + .queue_admit( + &make_job("unused-job", "Replay"), + b"payload", + "request-key", + "sha256:first", + ) + .unwrap(); + let QueueAdmissionOutcome::Existing(replay) = replay else { + panic!("exact replay should return the existing admission"); + }; + assert_eq!(replay.job_id, "admitted-job"); + assert_eq!(db.queue_list().unwrap().len(), 1); + + let conflict = db.queue_admit( + &make_job("conflicting-job", "Conflict"), + b"different", + "request-key", + "sha256:different", + ); + assert!(matches!(conflict, Err(NzbError::AdmissionConflict))); + assert_eq!(db.queue_list().unwrap().len(), 1); + } + + #[test] + fn queue_admission_survives_reopen_and_outlives_queue_observation() { + let directory = tempfile::tempdir().unwrap(); + let path = directory.path().join("queue.db"); + { + let mut db = Database::open(&path).unwrap(); + db.queue_admit( + &make_job("durable-job", "Durable"), + b"payload", + "durable-key", + "sha256:durable", + ) + .unwrap(); + } + + let mut db = Database::open(&path).unwrap(); + let observation = db.queue_admission_observe("durable-key").unwrap().unwrap(); + assert_eq!(observation.admission.job_id, "durable-job"); + assert!(matches!( + observation.state, + QueueAdmissionState::Queue { .. } + )); + + db.queue_remove("durable-job").unwrap(); + let observation = db.queue_admission_observe("durable-key").unwrap().unwrap(); + assert_eq!(observation.state, QueueAdmissionState::Unobserved); + + let replay = db + .queue_admit( + &make_job("new-job", "New"), + b"payload", + "durable-key", + "sha256:durable", + ) + .unwrap(); + let QueueAdmissionOutcome::Existing(replay) = replay else { + panic!("durable replay should preserve the original job identity"); + }; + assert_eq!(replay.job_id, "durable-job"); + assert!(db.queue_list().unwrap().is_empty()); + } + #[test] fn test_queue_update_progress() { let db = Database::open_memory().unwrap(); diff --git a/crates/nzb-core/src/error.rs b/crates/nzb-core/src/error.rs index a3cfa47b..e95b92c1 100644 --- a/crates/nzb-core/src/error.rs +++ b/crates/nzb-core/src/error.rs @@ -11,6 +11,9 @@ pub enum NzbError { #[error("Job not found: {0}")] JobNotFound(String), + #[error("Idempotency key is already bound to another payload")] + AdmissionConflict, + #[error("Server not found: {0}")] ServerNotFound(String), diff --git a/crates/nzb-core/src/models.rs b/crates/nzb-core/src/models.rs index bbffaeaa..de8d9c6c 100644 --- a/crates/nzb-core/src/models.rs +++ b/crates/nzb-core/src/models.rs @@ -225,6 +225,52 @@ pub struct DownloadStatistic { pub server_stats: Vec, } +/// Durable binding between a caller idempotency key and one admitted job. +/// +/// This row intentionally outlives queue and history retention so a caller can +/// resolve an ambiguous admission response without creating another job. +#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] +pub struct QueueAdmission { + pub idempotency_key: String, + pub payload_digest: String, + pub job_id: String, + pub accepted_at: DateTime, +} + +#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] +#[serde(tag = "location", rename_all = "snake_case")] +pub enum QueueAdmissionState { + Queue { + status: JobStatus, + total_bytes: u64, + downloaded_bytes: u64, + output_dir: PathBuf, + error_message: Option, + }, + History { + status: JobStatus, + total_bytes: u64, + downloaded_bytes: u64, + completed_at: DateTime, + output_dir: PathBuf, + error_message: Option, + }, + Unobserved, +} + +/// Exact observation of a durable admission and its current engine location. +#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] +pub struct QueueAdmissionObservation { + pub admission: QueueAdmission, + pub state: QueueAdmissionState, +} + +#[derive(Debug, Clone, PartialEq, Eq)] +pub enum QueueAdmissionOutcome { + Inserted(QueueAdmission), + Existing(QueueAdmission), +} + #[derive(Debug, Clone, Serialize, Deserialize)] pub struct StageResult { pub name: String, diff --git a/crates/nzb-web/src/error.rs b/crates/nzb-web/src/error.rs index 839de9d8..ec41ca17 100644 --- a/crates/nzb-web/src/error.rs +++ b/crates/nzb-web/src/error.rs @@ -67,6 +67,13 @@ impl ApiError { } } + pub const fn admission_conflict() -> Self { + Self { + status: Some(StatusCode::CONFLICT), + kind: ApiErrorKind::AdmissionConflict, + } + } + pub fn status(&self) -> StatusCode { self.status.unwrap_or(StatusCode::INTERNAL_SERVER_ERROR) } @@ -80,6 +87,8 @@ pub enum ApiErrorKind { ServerNotFound(String), #[error("unauthorized")] Unauthorized, + #[error("idempotency key is already bound to another NZB payload")] + AdmissionConflict, #[error("{0}")] Text(&'static str), #[error("{0}")] @@ -119,6 +128,7 @@ impl Serialize for ApiError { ApiErrorKind::JobNotFound(_) => "job_not_found", ApiErrorKind::ServerNotFound(_) => "server_not_found", ApiErrorKind::Unauthorized => "unauthorized", + ApiErrorKind::AdmissionConflict => "admission_conflict", _ => "internal_error", }, human_readable: format!("{:#}", self.kind), diff --git a/crates/nzb-web/src/queue_manager.rs b/crates/nzb-web/src/queue_manager.rs index 4aff63c2..d6eb8e76 100644 --- a/crates/nzb-web/src/queue_manager.rs +++ b/crates/nzb-web/src/queue_manager.rs @@ -1386,7 +1386,7 @@ impl QueueManager { /// the `max_active_downloads` limit. pub fn add_job( self: &Arc, - mut job: NzbJob, + job: NzbJob, nzb_data: Option>, ) -> crate::nzb_core::Result<()> { if crate::nzb_core::path::safe_component(&job.category).is_none() { @@ -1428,6 +1428,66 @@ impl QueueManager { } } + self.activate_admitted_job(job, nzb_data); + Ok(()) + } + + /// Admit one NZB exactly once for a caller-provided idempotency key. + /// + /// A first admission inserts the queue row, the durable admission binding, + /// and activates the job. A replay with the same key and identical payload + /// returns the existing admission without creating a second job; a replay + /// with a different payload is a conflict. + pub fn add_job_idempotent( + self: &Arc, + job: NzbJob, + nzb_data: Vec, + idempotency_key: &str, + payload_digest: &str, + ) -> crate::nzb_core::Result { + if let Some(observation) = self.db.lock().queue_admission_observe(idempotency_key)? { + if observation.admission.payload_digest != payload_digest { + return Err(crate::nzb_core::NzbError::AdmissionConflict); + } + return Ok(QueueAdmissionOutcome::Existing(observation.admission)); + } + + std::fs::create_dir_all(&job.work_dir)?; + let outcome = + self.db + .lock() + .queue_admit(&job, &nzb_data, idempotency_key, payload_digest)?; + match &outcome { + QueueAdmissionOutcome::Inserted(_) => { + self.activate_admitted_job(job, Some(nzb_data)); + } + QueueAdmissionOutcome::Existing(_) => { + if let Err(error) = std::fs::remove_dir(&job.work_dir) + && error.kind() != std::io::ErrorKind::NotFound + { + warn!( + work_dir = %job.work_dir.display(), + "Unable to remove an unused replay work directory: {error}" + ); + } + } + } + Ok(outcome) + } + + /// Observe one durable admission without scanning bounded queue or history lists. + pub fn queue_admission_observe( + &self, + idempotency_key: &str, + ) -> crate::nzb_core::Result> { + self.db.lock().queue_admission_observe(idempotency_key) + } + + /// Activate an already-persisted job: emit the added event and place it in + /// the in-memory queue (or paused set), starting a download slot if free. + /// The queue row must already exist — `add_job` and `add_job_idempotent` + /// persist it before calling this. + fn activate_admitted_job(self: &Arc, mut job: NzbJob, nzb_data: Option>) { let job_id = job.id.clone(); info!( job_id = %job_id, @@ -1460,7 +1520,7 @@ impl QueueManager { self.globally_paused_jobs.lock().insert(job_id.clone()); self.persist_globally_paused_jobs(); self.job_order.lock().push(job_id); - return Ok(()); + return; } // Insert as Queued — start_next_queued will atomically claim a @@ -1480,7 +1540,6 @@ impl QueueManager { // Try to start this or other queued jobs self.start_next_queued(); - Ok(()) } /// Rebuild a history job for retry. Newer history rows carry a checkpoint From 7f7a0449a67756ef71b3aab85c4582f1c0a34af0 Mon Sep 17 00:00:00 2001 From: thedancingdeveloper <306930456+thedancingdeveloper@users.noreply.github.com> Date: Sun, 27 Sep 2026 10:08:48 +0000 Subject: [PATCH 2/3] feat(nzb-postproc): publish typed terminal failure codes Persist an engine-owned failure code with each terminal history row instead of requiring consumers to parse human-readable diagnostic prose, and make the durable terminal row authoritative over a job's lingering queue view. - New `JobFailureCode` enum (`nzb-core`) with stable snake_case wire values (Display + FromStr): articles_unavailable, repair_failed, archive_password_required, archive_invalid, storage_unavailable, download_failed. - Post-processing assigns codes from causal stage state: `run_pipeline` threads a `failure_code` (articles-unavailable, repair-failed) and `run_extract_stage` returns a typed code, including a typed `ArchivePasswordRequired` error (replacing the `anyhow::bail!("archive is password-protected")` prose that consumers had to string-match). - Persist `HistoryEntry.failure_code` (DB migration v12 adds the column and backfills existing failed rows to `download_failed`); `history_insert` enforces the invariant that a Failed row carries exactly one code and a non-Failed row carries none. - `queue_manager` tracks the code in memory across pause/resume/clear and writes it on terminal transition; `queue_admission_observe` now consults the authoritative history row before the queue view, and reports the code via `QueueAdmissionState::History`. Exposed on the history API response. Ported from MrVampy/rustnzb (e0870d8), adapted to main: DB migration renumbered v10 -> v12 (main's v10/v11 are damage_ledger / queue_admissions); history SQL column indices adjusted for main's `retry_data` column; the pipeline changes reapplied onto main's `run_pipeline_with_cleanup` structure; and `QueueAdmissionState::History` extends the variant added in the idempotent queue-admission change this branch is stacked on. Fork flake/workspace bumps excluded. Tests: JobFailureCode wire-value roundtrip; typed ArchivePasswordRequired survives as a downcastable error; run_extract_stage / run_pipeline code assignment; DB migration backfill, history invariant, and terminal-history authority over the queue view; queue_manager terminal-code persistence (ArchiveInvalid / DownloadFailed). Full nzb-core / nzb-postproc / nzb-web / rustnzb suites pass; clippy/fmt clean. Co-Authored-By: MrVampy <4302946+MrVampy@users.noreply.github.com> Co-Authored-By: Claude Opus 4.8 --- apps/rustnzb/src/handlers.rs | 2 + crates/nzb-core/src/db.rs | 178 ++++++++++++++++++++++++--- crates/nzb-core/src/models.rs | 61 +++++++++ crates/nzb-postproc/src/pipeline.rs | 46 ++++++- crates/nzb-postproc/src/unpack.rs | 14 ++- crates/nzb-web/src/queue_manager.rs | 34 +++++ crates/nzb-web/src/sabnzbd_compat.rs | 4 + 7 files changed, 312 insertions(+), 27 deletions(-) diff --git a/apps/rustnzb/src/handlers.rs b/apps/rustnzb/src/handlers.rs index 785d2306..472abc4a 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 1f25acf8..0a93ee10 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 de8d9c6c..10c8d485 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 17a5e568..3feeb6ed 100644 --- a/crates/nzb-postproc/src/pipeline.rs +++ b/crates/nzb-postproc/src/pipeline.rs @@ -11,13 +11,13 @@ 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::{ArchiveType, find_archives, find_cleanup_files, find_par2_files}; 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() @@ -53,6 +53,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. @@ -129,6 +131,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"); @@ -153,6 +156,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, @@ -317,6 +321,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(), @@ -340,6 +345,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, @@ -352,6 +358,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, @@ -393,6 +400,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); } @@ -424,7 +432,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(), @@ -434,6 +442,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; @@ -478,6 +487,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)) + }, } } @@ -622,7 +636,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(); @@ -641,6 +655,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) @@ -668,6 +683,7 @@ async fn run_extract_stage( } Ok(unpack_result) => { all_ok = false; + failure_code = Some(JobFailureCode::ArchiveInvalid); let detail = unpack_result .error_output .trim() @@ -680,6 +696,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}")); } @@ -698,6 +719,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" )); @@ -713,6 +735,7 @@ async fn run_extract_stage( duration_secs: start.elapsed().as_secs_f64(), }, Vec::new(), + None, ); } @@ -728,6 +751,7 @@ async fn run_extract_stage( duration_secs: start.elapsed().as_secs_f64(), }, extracted_archives, + failure_code, ) } @@ -911,10 +935,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] @@ -1006,9 +1032,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" @@ -1025,9 +1053,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()); } @@ -1171,6 +1201,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 e01fba7b..e63c7c32 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 d6eb8e76..93cca9b8 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 f2939e28..c0cf7374 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, From f7955af307d11c7aa58d66553d48c83c86042939 Mon Sep 17 00:00:00 2001 From: thedancingdeveloper <306930456+thedancingdeveloper@users.noreply.github.com> Date: Sun, 27 Sep 2026 11:11:27 +0000 Subject: [PATCH 3/3] chore(desktop): reconcile src-tauri Cargo.lock for ported deps The desktop app depends on `rustnzb`/`nzb-web` by path; the new dependencies added by this change must be reflected in desktop/src-tauri/Cargo.lock so the `desktop` CI job (built with --locked) accepts it. Co-Authored-By: Claude Opus 4.8 --- desktop/src-tauri/Cargo.lock | 2 ++ 1 file changed, 2 insertions(+) diff --git a/desktop/src-tauri/Cargo.lock b/desktop/src-tauri/Cargo.lock index 670fc184..98eb47f0 100644 --- a/desktop/src-tauri/Cargo.lock +++ b/desktop/src-tauri/Cargo.lock @@ -4363,6 +4363,7 @@ dependencies = [ "chrono", "clap", "flate2", + "hex", "http", "libc", "mime_guess", @@ -4378,6 +4379,7 @@ dependencies = [ "rustls", "serde", "serde_json", + "sha2 0.11.0", "tokio", "tokio-util", "tower-http 0.7.0",