//! Background worker for the §6 stage 3 cast check. //! //! §8: this runs as a Tokio background task in the same binary, with the job //! queue as a SQLite table so state survives restart — replacing an external //! broker entirely. The check needs an outbound TMDB call and so cannot run //! inside the request without coupling upload latency to a third party (§6). use std::sync::Arc; use std::time::Duration; use anyhow::Context; use crate::castcheck::{self, SubmittedActor, Verdict}; use crate::db::{repo, Db}; use crate::ingest::{CastCheckJob, JOB_CAST_CHECK}; use crate::tmdb::{CastMember, Credits, TmdbClient, TmdbError}; /// §6: TMDB responses are cached for 24h, so a burst of episode uploads for one /// series costs a single upstream call. const CACHE_TTL: Duration = Duration::from_secs(24 * 3600); /// Cap on retry backoff for a persistent TMDB outage. const MAX_BACKOFF_SECS: u64 = 3600; pub struct Worker { pub db: Db, pub tmdb: Arc, pub batch: usize, pub poll_interval: Duration, /// Stamped as `origin` on this server's own change-feed entries (§9a), so a /// peer can tell what a manifest originated from and ignore its own echoes. pub server_id: String, } impl Worker { /// Runs until `shutdown` resolves. pub async fn run(self, mut shutdown: tokio::sync::watch::Receiver) { // A process that died mid-job would otherwise leave work stranded. match self.db.write(repo::release_all_leases).await { Ok(n) if n > 0 => tracing::info!(released = n, "released stranded job leases"), Ok(_) => {} Err(e) => tracing::error!(error = ?e, "failed to release job leases at startup"), } loop { tokio::select! { _ = shutdown.changed() => { tracing::info!("worker shutting down"); return; } _ = tokio::time::sleep(self.poll_interval) => { if let Err(e) = self.tick().await { tracing::error!(error = ?e, "worker tick failed"); } } } } } async fn tick(&self) -> anyhow::Result<()> { let now = now_iso(); let batch = self.batch; let leased_at = now.clone(); let jobs = self.db.write(move |tx| repo::lease_jobs(tx, &leased_at, batch)).await?; for job in jobs { let result = match job.kind.as_str() { JOB_CAST_CHECK => self.run_cast_check(&job.payload).await, crate::federation::JOB_FEDERATION_PULL => { self.run_federation_pull(&job.payload).await } other => { tracing::warn!(kind = other, "unknown job kind, dropping"); Ok(()) } }; let job_id = job.id.clone(); match result { Ok(()) => { self.db.write(move |tx| repo::delete_job(tx, &job_id)).await?; } Err(JobError::Retry(msg)) => { // §6: TMDB unreachable or rate-limited means retry with // backoff; the manifest stays unlisted, not rejected. let delay = backoff_secs(job.attempts); let run_after = iso_in(delay); tracing::warn!(job = %job_id, attempts = job.attempts, delay, reason = %msg, "rescheduling job"); self.db .write(move |tx| repo::reschedule_job(tx, &job_id, &run_after, &msg)) .await?; } Err(JobError::Fatal(e)) => { tracing::error!(job = %job_id, error = ?e, "dropping job after fatal error"); self.db.write(move |tx| repo::delete_job(tx, &job_id)).await?; } } } Ok(()) } /// Pulls one peer's change feed and ingests what is new (§9a). /// /// Transport failures are retryable: a peer being down is not a verdict on /// its content, exactly as a TMDB outage is not a verdict on an upload. /// /// TRACES: UR-008 | PR-006 async fn run_federation_pull(&self, payload: &str) -> Result<(), JobError> { let job: crate::federation::FederationPullJob = serde_json::from_str(payload).map_err(|e| JobError::Fatal(e.into()))?; let peer_id = job.peer_id.clone(); let peer = self .db .read(move |conn| repo::peer_by_id(conn, &peer_id)) .await .map_err(JobError::Fatal)?; // A peer removed or disabled between enqueue and run is not an error. let Some(peer) = peer.filter(|p| p.enabled) else { return Ok(()); }; let puller = crate::federation::Puller::new(self.server_id.clone()); let now = now_iso(); match puller.pull(&self.db, &peer, &now).await { Ok(outcome) => { tracing::info!( peer = %peer.url, examined = outcome.examined, ingested = outcome.ingested, known = outcome.skipped_known, flagged = outcome.flagged, rejected = outcome.rejected, capped = outcome.capped, "federation pull complete" ); Ok(()) } Err(e) => { let msg = e.to_string(); let pid = peer.id.clone(); let now2 = now.clone(); let m2 = msg.clone(); let _ = self .db .write(move |tx| repo::record_pull(tx, &pid, None, &now2, Some(&m2))) .await; Err(JobError::Retry(msg)) } } } /// TRACES: UR-003, UR-005 | SR-004 async fn run_cast_check(&self, payload: &str) -> Result<(), JobError> { let job: CastCheckJob = serde_json::from_str(payload).map_err(|e| JobError::Fatal(e.into()))?; let manifest_id = job.manifest_id; // Load what the check needs. let mid = manifest_id.clone(); let loaded = self .db .read(move |conn| { let Some(m) = repo::manifest_by_id(conn, &mid)? else { return Ok(None) }; let title = conn .query_row( "SELECT kind, tmdb_id, imdb_id, adult, certification FROM titles WHERE id = ?1", rusqlite::params![m.title_id], |r| { Ok(( r.get::<_, String>(0)?, r.get::<_, Option>(1)?, r.get::<_, Option>(2)?, r.get::<_, i64>(3)? != 0, r.get::<_, Option>(4)?, )) }, ) .map_err(anyhow::Error::from)?; let actor_ids = repo::manifest_actor_ids(conn, &mid)?; Ok(Some((m, title, actor_ids))) }) .await .map_err(JobError::Fatal)?; // The manifest may have been deleted (contributor revoked, §5a) between // enqueue and now; that is not an error. let Some((manifest, (kind, tmdb_id, _imdb_id, title_adult, certification), actor_ids)) = loaded else { return Ok(()); }; if manifest.status != "pending" { return Ok(()); } let Some(tmdb_id) = tmdb_id else { // No TMDB id means the cast check cannot run at all. §6 treats absent // reference data as flagged, not rejected. self.finalise(&manifest_id, Verdict::Flagged, 0.0, Some("no_tmdb_id"), &[], &[]) .await .map_err(JobError::Fatal)?; return Ok(()); }; let credits = self .credits_for(&kind, &tmdb_id, manifest.season, manifest.episode) .await .map_err(|e| { if e.is_retryable() { JobError::Retry(e.to_string()) } else { // A genuinely absent title is a verdict, not a transport // failure — handled below via empty credits. JobError::Retry(format!("non-retryable tmdb error treated as absent: {e}")) } }); let credits = match credits { Ok(c) => c, Err(JobError::Retry(msg)) if msg.starts_with("non-retryable") => { tracing::info!(manifest = %manifest_id, "tmdb has no such title; flagging"); Credits::default() } Err(e) => return Err(e), }; let reference: Vec = credits.all().cloned().collect(); let submitted: Vec = actor_ids .iter() .map(|id| SubmittedActor { tmdb_id: Some(*id), imdb_id: None, name: None }) .collect(); let mut outcome = castcheck::evaluate(&submitted, &reference); // §5a layer 1 — category guard. if let Some(offender) = castcheck::category_guard_violation(&outcome.matched, title_adult) { tracing::warn!(manifest = %manifest_id, person = offender, "category guard: adult-flagged performer on a non-adult title"); outcome.verdict = Verdict::Rejected; outcome.reason = Some("category_guard".into()); } // §5a layer 2 — age-appropriateness guard. if let Some(cert) = &certification { if castcheck::is_childrens_certification(cert) { castcheck::apply_childrens_guard(&mut outcome, submitted.len()); } } self.finalise( &manifest_id, outcome.verdict, outcome.ratio, outcome.reason.as_deref(), &outcome.matched, &outcome.unmatched_person_ids, ) .await .map_err(JobError::Fatal)?; Ok(()) } /// Fetches credits, using the 24h cache (§6). /// /// For episodes this is the **union** of TMDB's per-episode credits (cast + /// guest stars) and the series' aggregate credits: per-episode alone would /// reject recurring cast TMDB lists only at series level, series-wide alone /// would reject legitimate guest stars. async fn credits_for( &self, kind: &str, tmdb_id: &str, season: Option, episode: Option, ) -> Result { if kind == "movie" { return self.cached("movie", tmdb_id, || self.tmdb.movie_credits(tmdb_id)).await; } let series = self.cached("series", tmdb_id, || self.tmdb.series_credits(tmdb_id)).await?; let mut combined = series; if let (Some(s), Some(e)) = (season, episode) { let key = format!("{tmdb_id}:{s}:{e}"); match self.cached("episode", &key, || self.tmdb.episode_credits(tmdb_id, s, e)).await { Ok(ep) => { combined.cast.extend(ep.cast); combined.guest_stars.extend(ep.guest_stars); } // A missing episode entry is normal; the series set still applies. Err(TmdbError::NotFound) => {} Err(e) if e.is_retryable() => return Err(e), Err(e) => tracing::warn!(error = ?e, "ignoring episode credits error"), } } Ok(combined) } async fn cached(&self, kind: &str, key: &str, fetch: F) -> Result where F: FnOnce() -> Fut, Fut: std::future::Future>, { let (k, kk) = (key.to_string(), kind.to_string()); let cached = self .db .read(move |conn| repo::cached_credits(conn, &k, &kk)) .await .map_err(|e| TmdbError::Transport(e.to_string()))?; if let Some((json, fetched_at)) = cached { if !is_stale(&fetched_at, CACHE_TTL) { if let Ok(c) = serde_json::from_str::(&json) { return Ok(c); } } } let fresh = fetch().await?; let json = serde_json::to_string(&SerializableCredits::from(&fresh)) .map_err(|e| TmdbError::Malformed(e.to_string()))?; let (k, kk, now) = (key.to_string(), kind.to_string(), now_iso()); let _ = self.db.write(move |tx| repo::put_credits(tx, &k, &kk, &json, &now)).await; Ok(fresh) } /// Applies the verdict: resolves names into `people`, drops unmatched actors, /// and updates status and contributor counters — in one transaction. async fn finalise( &self, manifest_id: &str, verdict: Verdict, ratio: f64, reason: Option<&str>, matched: &[castcheck::MatchedActor], unmatched: &[u64], ) -> anyhow::Result<()> { let id = manifest_id.to_string(); let reason = reason.map(str::to_string); let matched: Vec<(u64, String, bool)> = matched.iter().map(|m| (m.tmdb_person_id, m.name.clone(), m.adult)).collect(); let unmatched = unmatched.to_vec(); let now = now_iso(); let server_id = self.server_id.clone(); self.db .write(move |tx| { let contributor: Option = tx .query_row( "SELECT contributor_id FROM manifests WHERE id = ?1", rusqlite::params![id], |r| r.get(0), ) .map_err(anyhow::Error::from)?; if verdict == Verdict::Rejected { // §6: the manifest is deleted and the contributor notified // (via `GET /manifests/{id}/status` until it is gone). repo::delete_manifest(tx, &id)?; } else { // Names come from TMDB, never from the upload (§5a, §7). for (person_id, name, adult) in &matched { repo::upsert_person(tx, *person_id, name, *adult, &now)?; } // §6: unmatched actors are dropped rather than stored. for person_id in &unmatched { repo::delete_manifest_actor(tx, &id, *person_id)?; } repo::set_manifest_status( tx, &id, verdict.status(), reason.as_deref(), Some(ratio), )?; // §9a: the change feed carries locally-*listed* manifests // only. A `flagged` one is served with reduced ranking here // but is not offered for replication — nominating something // this server itself doubts would push a local judgement // call outward, which is exactly what federation must not do. if verdict == Verdict::Listed { if let Some(cid) = repo::content_id_of(tx, &id)? { repo::append_change(tx, &cid, "add", None, &server_id, &now)?; } } } if let Some(c) = contributor { let counter = match verdict { Verdict::Listed => "accepted", Verdict::Flagged => "flagged", Verdict::Rejected => "rejected", }; repo::bump_contributor_counter(tx, &c, counter)?; if repo::maybe_revoke_contributor(tx, &c, &now)? { tracing::warn!(contributor = %c, "revoked token for excessive rejections"); } } Ok(()) }) .await .context("finalising cast check") } } /// Serialisable projection of `Credits` for the cache. #[derive(serde::Serialize)] struct SerializableCredits { cast: Vec, guest_stars: Vec, } #[derive(serde::Serialize)] struct SerializableMember { id: u64, name: String, adult: bool, } impl From<&Credits> for SerializableCredits { fn from(c: &Credits) -> Self { let f = |m: &CastMember| SerializableMember { id: m.id, name: m.name.clone(), adult: m.adult }; Self { cast: c.cast.iter().map(f).collect(), guest_stars: c.guest_stars.iter().map(&f).collect(), } } } enum JobError { Retry(String), Fatal(anyhow::Error), } /// Exponential backoff, capped (§5: the client must back off exponentially /// rather than retrying tightly; the same discipline applies to our own /// outbound calls). fn backoff_secs(attempts: i64) -> u64 { let base = 30u64; base.saturating_mul(1u64 << attempts.clamp(0, 8) as u32).min(MAX_BACKOFF_SECS) } fn is_stale(fetched_at: &str, ttl: Duration) -> bool { let Some(then) = parse_iso(fetched_at) else { return true }; let now = unix_now(); now.saturating_sub(then) > ttl.as_secs() } /// Current time as an RFC 3339 UTC string, which is what every timestamp column /// stores. Kept in one place so the format cannot drift. pub fn now_iso() -> String { iso_in(0) } pub fn iso_in(secs: u64) -> String { format_unix(unix_now() + secs) } /// A timestamp `secs` in the past, for windowed counts (§9a `MaxIngestPerHour`). pub fn iso_ago(secs: u64) -> String { format_unix(unix_now().saturating_sub(secs)) } fn unix_now() -> u64 { std::time::SystemTime::now() .duration_since(std::time::UNIX_EPOCH) .map(|d| d.as_secs()) .unwrap_or(0) } /// Formats a Unix timestamp as `YYYY-MM-DDTHH:MM:SSZ`. /// /// Hand-rolled rather than pulling in `chrono`/`time`: the only requirement is a /// lexicographically-sortable UTC string, which is what the `jobs.run_after` /// comparison relies on. pub fn format_unix(mut secs: u64) -> String { let days = secs / 86_400; secs %= 86_400; let (h, m, s) = (secs / 3600, (secs % 3600) / 60, secs % 60); // Civil-from-days, Howard Hinnant's algorithm. let z = days as i64 + 719_468; let era = z.div_euclid(146_097); let doe = z.rem_euclid(146_097); let yoe = (doe - doe / 1460 + doe / 36_524 - doe / 146_096) / 365; let y = yoe + era * 400; let doy = doe - (365 * yoe + yoe / 4 - yoe / 100); let mp = (5 * doy + 2) / 153; let d = doy - (153 * mp + 2) / 5 + 1; let mo = if mp < 10 { mp + 3 } else { mp - 9 }; let y = if mo <= 2 { y + 1 } else { y }; format!("{y:04}-{mo:02}-{d:02}T{h:02}:{m:02}:{s:02}Z") } fn parse_iso(s: &str) -> Option { // Parses the format `format_unix` produces. let b = s.as_bytes(); if b.len() < 20 { return None; } let num = |from: usize, to: usize| s.get(from..to)?.parse::().ok(); let (y, mo, d) = (num(0, 4)?, num(5, 7)?, num(8, 10)?); let (h, mi, se) = (num(11, 13)?, num(14, 16)?, num(17, 19)?); let y_adj = if mo <= 2 { y - 1 } else { y }; let era = y_adj.div_euclid(400); let yoe = y_adj - era * 400; let mp = if mo > 2 { mo - 3 } else { mo + 9 }; let doy = (153 * mp + 2) / 5 + d - 1; let doe = yoe * 365 + yoe / 4 - yoe / 100 + doy; let days = era * 146_097 + doe - 719_468; Some((days * 86_400 + h * 3600 + mi * 60 + se).max(0) as u64) } #[cfg(test)] mod tests { use super::*; #[test] fn timestamp_roundtrips() { for t in [0u64, 1, 1_000_000, 1_700_000_000, 1_785_000_000, 4_000_000_000] { let s = format_unix(t); assert_eq!(parse_iso(&s), Some(t), "roundtrip failed for {t} => {s}"); } } #[test] fn timestamp_format_is_sortable() { // `jobs.run_after <= ?1` is a string comparison, so lexical order must // match chronological order. let a = format_unix(1_700_000_000); let b = format_unix(1_700_000_001); let c = format_unix(1_800_000_000); assert!(a < b && b < c, "{a} {b} {c}"); assert_eq!(format_unix(0), "1970-01-01T00:00:00Z"); } #[test] fn known_dates_format_correctly() { // 2026-07-30T12:00:00Z assert_eq!(format_unix(1_785_412_800), "2026-07-30T12:00:00Z"); // A leap day, since the civil-from-days algorithm is where this would // break. assert_eq!(format_unix(1_709_164_800), "2024-02-29T00:00:00Z"); } #[test] fn backoff_grows_and_is_capped() { assert_eq!(backoff_secs(0), 30); assert_eq!(backoff_secs(1), 60); assert_eq!(backoff_secs(4), 480); assert_eq!(backoff_secs(50), MAX_BACKOFF_SECS, "must not overflow or grow unbounded"); } #[test] fn staleness_uses_the_ttl() { let fresh = format_unix(unix_now()); assert!(!is_stale(&fresh, CACHE_TTL)); let old = format_unix(unix_now() - 25 * 3600); assert!(is_stale(&old, CACHE_TTL)); assert!(is_stale("not-a-timestamp", CACHE_TTL)); } }