Files
JRay-public-server/src/worker.rs
dtourolleandClaude Opus 5 545c7d92a2
CI / fmt, clippy, test (push) Failing after 1m20s
CI / static musl binary (push) Has been skipped
CI / advisories and licences (push) Successful in 25s
Federation: replicate content, re-derive judgement (UR-008)
Implements §9a. The replication surface is four reads and no writes: a change
feed, fetch by content_id, a batch have, and a human-facing peer directory —
plus a capabilities endpoint carrying the accepted envelope versions, which
lets a client discover a schema mismatch in one request instead of a 400 per
manifest across a library sweep.

Pull, never push: a pulling server chooses what it ingests and when. Push would
let any peer inject work into the validation queue — the same abuse surface as
anonymous upload, at higher volume.

Nothing inherits a peer's judgement. A pulled manifest runs the full §6 stage 1
and 2 validation and this server's own cast check, and the fetched body must
hash to the content_id that was asked for — the check that stops an
intermediary or a misbehaving peer substituting content under a trusted id.
A peer's retraction flags for review rather than delisting, because
auto-delisting would hand every peer a remote delete primitive; only the opt-in
per-peer abuse channel delists, because a takedown propagating at the speed of
manual review is the wrong failure mode for that one case.

A test caught a real bug in the first cut: the feed cursor was a ULID, and
ULIDs are only monotonic *between* milliseconds — two generated in the same
millisecond carry independent random components and can sort opposite to write
order. A peer resuming from `seq > cursor` would then silently skip an entry:
replication losing manifests with no error anywhere. The cursor is now an
AUTOINCREMENT integer, and the test asserts strict monotonicity rather than
merely sortedness.

Peer administration is deliberately not an API. §9a requires that a peering
exist only because an operator typed a URL, so nothing a remote server returns
can establish or widen one; there_is_no_endpoint_that_creates_a_peering asserts
that absence rather than trusting it.

212 tests. Coverage 25/32 (78%).

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>

TRACES: UR-008 | PR-006
2026-07-31 09:28:32 +02:00

567 lines
21 KiB
Rust

//! 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,
/// 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<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,
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<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();
let server_id = self.server_id.clone();
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),
)?;
// §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<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)
}
/// 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<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));
}
}