Initial implementation: core vertical slice
Implements the core of SPEC.md — the manifest exchange, less audio-tier matching (§3) and federation (§9a), both of which the spec sequences as later work. - §2 Jmanifest format and series bundles - §3 cut matching: exact / runtime / loose tiers - §4 API, less POST /manifests/search - §5 rate limiting; §5a trust model, anonymous bearer tokens - §6 upload validation, all four stages - §7 relational storage, no JSON blob on the write path - §8 Rust + Axum + SQLite, single serialized writer, in-process job queue - §9a content addressing, computed on upload Reconciled against the system spec: - anneal_sec removed, withdrawn upstream by AR-012/AR-013. Presence follows track extent, so a track survives its own gaps and there is nothing to anneal. Its successor extinction_sec and the new gallery_scope are accepted and stored; scope enters the §7 ranking. A manifest still carrying anneal_sec is a hard 400, not silently ignored — it came from a pipeline whose window semantics differ from what this server assumes. - Audio signature: media under 120 s now emits no signature at all, matching scene-actor-extraction IR-007. The earlier §3 draft allowed a shortened window under 150 s, which was the weaker rule — a caller-varying length is the property SR-004 forbids. - UR IDs regularised to UR-nnn; docs/requirements.md registers 32 requirements, each tracing to an SR-nnn or PR-nnn. 189 tests: unit, end-to-end through the real router, and an injection suite covering SQL, JSON, header and Unicode payloads. Writing that suite found two real gaps, both fixed here: compatibility homoglyphs passed the §5a character class, and a one-frame audio signature was accepted on a feature-length item. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
This commit is contained in:
+494
@@ -0,0 +1,494 @@
|
||||
//! 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<TmdbClient>,
|
||||
pub batch: usize,
|
||||
pub poll_interval: Duration,
|
||||
}
|
||||
|
||||
impl Worker {
|
||||
/// Runs until `shutdown` resolves.
|
||||
pub async fn run(self, mut shutdown: tokio::sync::watch::Receiver<bool>) {
|
||||
// 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,
|
||||
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(())
|
||||
}
|
||||
|
||||
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<String>>(1)?,
|
||||
r.get::<_, Option<String>>(2)?,
|
||||
r.get::<_, i64>(3)? != 0,
|
||||
r.get::<_, Option<String>>(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<CastMember> = credits.all().cloned().collect();
|
||||
let submitted: Vec<SubmittedActor> = 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<i64>,
|
||||
episode: Option<i64>,
|
||||
) -> Result<Credits, TmdbError> {
|
||||
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<F, Fut>(&self, kind: &str, key: &str, fetch: F) -> Result<Credits, TmdbError>
|
||||
where
|
||||
F: FnOnce() -> Fut,
|
||||
Fut: std::future::Future<Output = Result<Credits, TmdbError>>,
|
||||
{
|
||||
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::<Credits>(&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();
|
||||
|
||||
self.db
|
||||
.write(move |tx| {
|
||||
let contributor: Option<String> = 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),
|
||||
)?;
|
||||
}
|
||||
|
||||
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<SerializableMember>,
|
||||
guest_stars: Vec<SerializableMember>,
|
||||
}
|
||||
|
||||
#[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)
|
||||
}
|
||||
|
||||
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<u64> {
|
||||
// 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::<i64>().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));
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user