Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 2 additions & 0 deletions apps/rustnzb/src/handlers.rs
Original file line number Diff line number Diff line change
Expand Up @@ -153,6 +153,7 @@ pub struct HistoryResponseEntry {
pub output_dir: String,
pub stages: Vec<StageResult>,
pub error_message: Option<String>,
pub failure_code: Option<JobFailureCode>,
pub server_stats: Vec<ServerArticleStats>,
pub has_nzb_data: bool,
pub duration_secs: f64,
Expand Down Expand Up @@ -186,6 +187,7 @@ impl From<HistoryEntry> 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,
Expand Down
178 changes: 159 additions & 19 deletions crates/nzb-core/src/db.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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(())
}

Expand Down Expand Up @@ -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 }));
}

Expand Down Expand Up @@ -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,
Expand All @@ -706,6 +732,7 @@ impl Database {
entry.nzb_data,
server_stats_json,
entry.retry_data,
entry.failure_code.map(|code| code.to_string()),
],
)?;

Expand Down Expand Up @@ -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",
)?;

Expand All @@ -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 },
Expand Down Expand Up @@ -859,7 +887,8 @@ impl Database {
pub fn history_get(&self, id: &str) -> Result<Option<HistoryEntry>, 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",
)?;

Expand All @@ -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,
Expand Down Expand Up @@ -1251,6 +1281,24 @@ fn parse_status(s: &str) -> JobStatus {
}
}

fn parse_failure_code(value: Option<String>) -> rusqlite::Result<Option<JobFailureCode>> {
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,
Expand Down Expand Up @@ -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,
Expand Down Expand Up @@ -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<String> = db
.conn
.query_row(
"SELECT failure_code FROM history WHERE id = 'failed-job'",
[],
|row| row.get(0),
)
.unwrap();
let completed: Option<String> = 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();
Expand Down
61 changes: 61 additions & 0 deletions crates/nzb-core/src/models.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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<Self, Self::Err> {
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
// ---------------------------------------------------------------------------
Expand Down Expand Up @@ -199,6 +243,7 @@ pub struct HistoryEntry {
/// Post-processing stages with results
pub stages: Vec<StageResult>,
pub error_message: Option<String>,
pub failure_code: Option<JobFailureCode>,
/// Per-server download statistics
#[serde(default)]
pub server_stats: Vec<ServerArticleStats>,
Expand Down Expand Up @@ -254,6 +299,7 @@ pub enum QueueAdmissionState {
completed_at: DateTime<Utc>,
output_dir: PathBuf,
error_message: Option<String>,
failure_code: Option<JobFailureCode>,
},
Unobserved,
}
Expand Down Expand Up @@ -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::<JobFailureCode>().is_err());
}

#[test]
fn test_job_status_serde_roundtrip() {
let statuses = [
Expand Down
Loading
Loading