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/2] 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 971c12c13891f215d68a2374b7779c8be2eeacbf Mon Sep 17 00:00:00 2001 From: thedancingdeveloper <306930456+thedancingdeveloper@users.noreply.github.com> Date: Sun, 27 Sep 2026 11:11:24 +0000 Subject: [PATCH 2/2] 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",