Adds jray-project as a submodule at scripts/vendor/jray-project, so this repo runs the same extractor as every other component rather than its own copy, and gains the system spec that defines the PR/SR requirements its register traces up to. scripts/traceability-gate.sh is a thin wrapper holding only what is specific to this repo: UR/DR prefixes, .rs sources, and REPO_ROOT — which the shared gate cannot infer once vendored, since its default resolves to the submodule itself. Each override fails silently in a way that looks like "no work done" rather than "misconfigured", so the wrapper documents why each is needed. Annotates 35 units with TRACES tags, on the code that decides rather than every helper it calls. Coverage is 23/32 (71.9%) with no orphan tags. The nine untraced are genuinely unimplemented: UR-007 is plugin-side, UR-008 is federation, and UR-015..018 are the pending SR-003 schema bump. The gate caught a real error in the first pass: several tags separated IDs of different types with commas. A comma joins IDs within one type; a pipe separates types. Fixed, and the diagnostics are now clean. MIN_COVERAGE stays 0 deliberately. The gate still fails on orphan tags, a >100% ratio, a register parsing to nothing, or an empty source scan — raise the threshold as a ratchet once the remaining work lands. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
1114 lines
38 KiB
Rust
1114 lines
38 KiB
Rust
//! Repository queries.
|
|
//!
|
|
//! §8: all database access lives behind this layer rather than being scattered
|
|
//! through handlers — that is what keeps the Turso/Postgres options cheap and
|
|
//! localises the single-writer serialization in one place.
|
|
//!
|
|
//! Functions here take `&Connection` or `&Transaction` and are synchronous; the
|
|
//! `spawn_blocking` boundary is [`crate::db::Db`]'s concern.
|
|
|
|
use anyhow::Context;
|
|
use rusqlite::{params, Connection, OptionalExtension, Transaction};
|
|
|
|
use crate::model::IdentityType;
|
|
|
|
/// A stored manifest row, as needed to serve reads.
|
|
#[derive(Debug, Clone)]
|
|
pub struct ManifestRow {
|
|
pub id: String,
|
|
pub title_id: String,
|
|
pub season: Option<i64>,
|
|
pub episode: Option<i64>,
|
|
pub runtime_sec: f64,
|
|
pub video_hash: Option<String>,
|
|
pub audio_signature: Option<Vec<u8>>,
|
|
pub sample_fps: Option<f64>,
|
|
pub extinction_sec: Option<f64>,
|
|
pub pipeline_version: Option<String>,
|
|
pub gallery_scope: Option<String>,
|
|
pub status: String,
|
|
pub cast_match_ratio: Option<f64>,
|
|
pub content_id: Option<String>,
|
|
}
|
|
|
|
const MANIFEST_COLUMNS: &str = "id, title_id, season, episode, runtime_sec, video_hash, \
|
|
audio_signature, sample_fps, extinction_sec, pipeline_version, gallery_scope, \
|
|
status, cast_match_ratio, \
|
|
content_id";
|
|
|
|
fn map_manifest(row: &rusqlite::Row<'_>) -> rusqlite::Result<ManifestRow> {
|
|
Ok(ManifestRow {
|
|
id: row.get(0)?,
|
|
title_id: row.get(1)?,
|
|
season: row.get(2)?,
|
|
episode: row.get(3)?,
|
|
runtime_sec: row.get(4)?,
|
|
video_hash: row.get(5)?,
|
|
audio_signature: row.get(6)?,
|
|
sample_fps: row.get(7)?,
|
|
extinction_sec: row.get(8)?,
|
|
pipeline_version: row.get(9)?,
|
|
gallery_scope: row.get(10)?,
|
|
status: row.get(11)?,
|
|
cast_match_ratio: row.get(12)?,
|
|
content_id: row.get(13)?,
|
|
})
|
|
}
|
|
|
|
/// Statuses that are served to clients (§7).
|
|
pub const SERVED_STATUSES: &str = "('listed','flagged')";
|
|
|
|
// ---------------------------------------------------------------------------
|
|
// Titles
|
|
// ---------------------------------------------------------------------------
|
|
|
|
#[derive(Debug, Clone)]
|
|
pub struct TitleRow {
|
|
pub id: String,
|
|
pub kind: String,
|
|
pub tmdb_id: Option<String>,
|
|
pub imdb_id: Option<String>,
|
|
pub name: Option<String>,
|
|
pub year: Option<i64>,
|
|
pub adult: bool,
|
|
pub certification: Option<String>,
|
|
}
|
|
|
|
pub fn find_title(
|
|
conn: &Connection,
|
|
kind: IdentityType,
|
|
tmdb_id: Option<&str>,
|
|
imdb_id: Option<&str>,
|
|
) -> anyhow::Result<Option<TitleRow>> {
|
|
let kind_str = title_kind(kind);
|
|
let mut stmt = conn.prepare_cached(
|
|
"SELECT id, kind, tmdb_id, imdb_id, name, year, adult, certification
|
|
FROM titles
|
|
WHERE kind = ?1
|
|
AND ( (?2 IS NOT NULL AND tmdb_id = ?2) OR (?3 IS NOT NULL AND imdb_id = ?3) )
|
|
LIMIT 1",
|
|
)?;
|
|
let row = stmt
|
|
.query_row(params![kind_str, tmdb_id, imdb_id], |r| {
|
|
Ok(TitleRow {
|
|
id: r.get(0)?,
|
|
kind: r.get(1)?,
|
|
tmdb_id: r.get(2)?,
|
|
imdb_id: r.get(3)?,
|
|
name: r.get(4)?,
|
|
year: r.get(5)?,
|
|
adult: r.get::<_, i64>(6)? != 0,
|
|
certification: r.get(7)?,
|
|
})
|
|
})
|
|
.optional()?;
|
|
Ok(row)
|
|
}
|
|
|
|
/// `kind` for the `titles` table: an episode manifest hangs off its *series*
|
|
/// title, since §2 keys episodes on series coordinates.
|
|
pub fn title_kind(kind: IdentityType) -> &'static str {
|
|
match kind {
|
|
IdentityType::Movie => "movie",
|
|
IdentityType::Episode => "series",
|
|
}
|
|
}
|
|
|
|
/// Finds or creates the title row, returning its id.
|
|
pub fn upsert_title(
|
|
tx: &Transaction<'_>,
|
|
kind: IdentityType,
|
|
tmdb_id: Option<&str>,
|
|
imdb_id: Option<&str>,
|
|
name: Option<&str>,
|
|
year: Option<i64>,
|
|
now: &str,
|
|
) -> anyhow::Result<String> {
|
|
if let Some(existing) = find_title(tx, kind, tmdb_id, imdb_id)? {
|
|
// Backfill identifiers a later upload supplied but an earlier one lacked.
|
|
tx.execute(
|
|
"UPDATE titles
|
|
SET tmdb_id = COALESCE(tmdb_id, ?2),
|
|
imdb_id = COALESCE(imdb_id, ?3),
|
|
year = COALESCE(year, ?4),
|
|
updated_at = ?5
|
|
WHERE id = ?1",
|
|
params![existing.id, tmdb_id, imdb_id, year, now],
|
|
)?;
|
|
return Ok(existing.id);
|
|
}
|
|
|
|
let id = ulid::Ulid::new().to_string();
|
|
tx.execute(
|
|
"INSERT INTO titles (id, kind, tmdb_id, imdb_id, name, year, adult, certification, updated_at)
|
|
VALUES (?1, ?2, ?3, ?4, ?5, ?6, 0, NULL, ?7)",
|
|
params![id, title_kind(kind), tmdb_id, imdb_id, name, year, now],
|
|
)?;
|
|
Ok(id)
|
|
}
|
|
|
|
/// Records TMDB-derived title attributes used by the §5a guards.
|
|
pub fn set_title_attributes(
|
|
tx: &Transaction<'_>,
|
|
title_id: &str,
|
|
adult: bool,
|
|
certification: Option<&str>,
|
|
name: Option<&str>,
|
|
now: &str,
|
|
) -> anyhow::Result<()> {
|
|
tx.execute(
|
|
"UPDATE titles
|
|
SET adult = ?2,
|
|
certification = COALESCE(?3, certification),
|
|
name = COALESCE(?4, name),
|
|
updated_at = ?5
|
|
WHERE id = ?1",
|
|
params![title_id, adult as i64, certification, name, now],
|
|
)?;
|
|
Ok(())
|
|
}
|
|
|
|
// ---------------------------------------------------------------------------
|
|
// People — TMDB-derived, never from an upload (§5a, §7)
|
|
// ---------------------------------------------------------------------------
|
|
|
|
pub fn upsert_person(
|
|
tx: &Transaction<'_>,
|
|
tmdb_person_id: u64,
|
|
name: &str,
|
|
adult: bool,
|
|
now: &str,
|
|
) -> anyhow::Result<()> {
|
|
// Portable `INSERT ... ON CONFLICT` rather than `INSERT OR REPLACE` (§8).
|
|
tx.execute(
|
|
"INSERT INTO people (tmdb_person_id, name, adult, updated_at)
|
|
VALUES (?1, ?2, ?3, ?4)
|
|
ON CONFLICT (tmdb_person_id) DO UPDATE
|
|
SET name = excluded.name, adult = excluded.adult, updated_at = excluded.updated_at",
|
|
params![tmdb_person_id as i64, name, adult as i64, now],
|
|
)?;
|
|
Ok(())
|
|
}
|
|
|
|
pub fn person_names(
|
|
conn: &Connection,
|
|
ids: &[u64],
|
|
) -> anyhow::Result<std::collections::HashMap<u64, String>> {
|
|
let mut out = std::collections::HashMap::new();
|
|
let mut stmt = conn.prepare_cached("SELECT name FROM people WHERE tmdb_person_id = ?1")?;
|
|
for &id in ids {
|
|
if let Some(name) =
|
|
stmt.query_row(params![id as i64], |r| r.get::<_, String>(0)).optional()?
|
|
{
|
|
out.insert(id, name);
|
|
}
|
|
}
|
|
Ok(out)
|
|
}
|
|
|
|
// ---------------------------------------------------------------------------
|
|
// Manifests
|
|
// ---------------------------------------------------------------------------
|
|
|
|
/// Every candidate manifest for a title (and episode coordinates, if given)
|
|
/// that is eligible to be served.
|
|
///
|
|
/// Cut matching is done in [`crate::matching`] rather than in SQL: the tier
|
|
/// logic is spec-critical and belongs somewhere testable.
|
|
pub fn candidates_for_title(
|
|
conn: &Connection,
|
|
title_id: &str,
|
|
season: Option<i64>,
|
|
episode: Option<i64>,
|
|
) -> anyhow::Result<Vec<ManifestRow>> {
|
|
let sql = format!(
|
|
"SELECT {MANIFEST_COLUMNS} FROM manifests
|
|
WHERE title_id = ?1
|
|
AND status IN {SERVED_STATUSES}
|
|
AND ( (?2 IS NULL AND season IS NULL) OR season = ?2 )
|
|
AND ( (?3 IS NULL AND episode IS NULL) OR episode = ?3 )
|
|
ORDER BY cast_match_ratio DESC NULLS LAST, sample_fps DESC NULLS LAST, created_at ASC"
|
|
);
|
|
let mut stmt = conn.prepare_cached(&sql)?;
|
|
let rows = stmt
|
|
.query_map(params![title_id, season, episode], map_manifest)?
|
|
.collect::<rusqlite::Result<Vec<_>>>()?;
|
|
Ok(rows)
|
|
}
|
|
|
|
/// All servable episode manifests for a series, for bundle assembly (§2).
|
|
///
|
|
/// A bundle is assembled per request; there is no "series manifest" row.
|
|
pub fn episodes_for_series(
|
|
conn: &Connection,
|
|
title_id: &str,
|
|
season: Option<i64>,
|
|
) -> anyhow::Result<Vec<ManifestRow>> {
|
|
let sql = format!(
|
|
"SELECT {MANIFEST_COLUMNS} FROM manifests
|
|
WHERE title_id = ?1
|
|
AND status IN {SERVED_STATUSES}
|
|
AND season IS NOT NULL AND episode IS NOT NULL
|
|
AND (?2 IS NULL OR season = ?2)
|
|
ORDER BY season ASC, episode ASC,
|
|
cast_match_ratio DESC NULLS LAST, sample_fps DESC NULLS LAST"
|
|
);
|
|
let mut stmt = conn.prepare_cached(&sql)?;
|
|
let rows = stmt
|
|
.query_map(params![title_id, season], map_manifest)?
|
|
.collect::<rusqlite::Result<Vec<_>>>()?;
|
|
Ok(rows)
|
|
}
|
|
|
|
pub fn manifest_by_id(conn: &Connection, id: &str) -> anyhow::Result<Option<ManifestRow>> {
|
|
let sql = format!("SELECT {MANIFEST_COLUMNS} FROM manifests WHERE id = ?1");
|
|
let mut stmt = conn.prepare_cached(&sql)?;
|
|
Ok(stmt.query_row(params![id], map_manifest).optional()?)
|
|
}
|
|
|
|
pub fn manifest_status(
|
|
conn: &Connection,
|
|
id: &str,
|
|
) -> anyhow::Result<Option<(String, Option<String>)>> {
|
|
let mut stmt =
|
|
conn.prepare_cached("SELECT status, reject_reason FROM manifests WHERE id = ?1")?;
|
|
Ok(stmt.query_row(params![id], |r| Ok((r.get(0)?, r.get(1)?))).optional()?)
|
|
}
|
|
|
|
/// §4 `409`: an identical `(identity, cut)` manifest already exists from this
|
|
/// contributor.
|
|
pub fn duplicate_from_contributor(
|
|
tx: &Transaction<'_>,
|
|
title_id: &str,
|
|
season: Option<i64>,
|
|
episode: Option<i64>,
|
|
runtime_sec: f64,
|
|
video_hash: Option<&str>,
|
|
contributor_id: &str,
|
|
) -> anyhow::Result<Option<String>> {
|
|
let mut stmt = tx.prepare_cached(
|
|
"SELECT id FROM manifests
|
|
WHERE title_id = ?1 AND contributor_id = ?6
|
|
AND status <> 'rejected'
|
|
AND ( (?2 IS NULL AND season IS NULL) OR season = ?2 )
|
|
AND ( (?3 IS NULL AND episode IS NULL) OR episode = ?3 )
|
|
AND ABS(runtime_sec - ?4) < 0.001
|
|
AND ( (?5 IS NULL AND video_hash IS NULL) OR video_hash = ?5 )
|
|
LIMIT 1",
|
|
)?;
|
|
Ok(stmt
|
|
.query_row(
|
|
params![title_id, season, episode, runtime_sec, video_hash, contributor_id],
|
|
|r| r.get::<_, String>(0),
|
|
)
|
|
.optional()?)
|
|
}
|
|
|
|
pub fn manifest_by_content_id(
|
|
tx: &Transaction<'_>,
|
|
content_id: &str,
|
|
) -> anyhow::Result<Option<String>> {
|
|
let mut stmt = tx.prepare_cached("SELECT id FROM manifests WHERE content_id = ?1")?;
|
|
Ok(stmt.query_row(params![content_id], |r| r.get::<_, String>(0)).optional()?)
|
|
}
|
|
|
|
/// Everything needed to insert one manifest.
|
|
pub struct NewManifest<'a> {
|
|
pub id: &'a str,
|
|
pub title_id: &'a str,
|
|
pub season: Option<i64>,
|
|
pub episode: Option<i64>,
|
|
pub runtime_sec: f64,
|
|
pub video_hash: Option<&'a str>,
|
|
pub audio_signature: Option<&'a [u8]>,
|
|
pub audio_sig_coarse: Option<&'a [u8]>,
|
|
pub sample_fps: Option<f64>,
|
|
pub extinction_sec: Option<f64>,
|
|
pub pipeline_version: Option<&'a str>,
|
|
pub gallery_scope: Option<&'a str>,
|
|
pub contributor_id: Option<&'a str>,
|
|
pub status: &'a str,
|
|
pub content_id: Option<&'a str>,
|
|
pub origin: &'a str,
|
|
pub ingested_from: Option<&'a str>,
|
|
pub created_at: &'a str,
|
|
}
|
|
|
|
/// TRACES: UR-012 | DR-002, DR-004 | SR-004, SR-005
|
|
pub fn insert_manifest(tx: &Transaction<'_>, m: &NewManifest<'_>) -> anyhow::Result<()> {
|
|
tx.execute(
|
|
"INSERT INTO manifests
|
|
(id, title_id, season, episode, runtime_sec, video_hash,
|
|
audio_signature, audio_sig_coarse, sample_fps, extinction_sec, pipeline_version,
|
|
gallery_scope, contributor_id, status, content_id, origin, ingested_from,
|
|
created_at)
|
|
VALUES (?1,?2,?3,?4,?5,?6,?7,?8,?9,?10,?11,?12,?13,?14,?15,?16,?17,?18)",
|
|
params![
|
|
m.id,
|
|
m.title_id,
|
|
m.season,
|
|
m.episode,
|
|
m.runtime_sec,
|
|
m.video_hash,
|
|
m.audio_signature,
|
|
m.audio_sig_coarse,
|
|
m.sample_fps,
|
|
m.extinction_sec,
|
|
m.pipeline_version,
|
|
m.gallery_scope,
|
|
m.contributor_id,
|
|
m.status,
|
|
m.content_id,
|
|
m.origin,
|
|
m.ingested_from,
|
|
m.created_at,
|
|
],
|
|
)
|
|
.context("inserting manifest")?;
|
|
Ok(())
|
|
}
|
|
|
|
/// Inserts one actor's rows. Batched within the caller's transaction — one
|
|
/// transaction per manifest, not per row (§8).
|
|
pub fn insert_actor_scenes(
|
|
tx: &Transaction<'_>,
|
|
manifest_id: &str,
|
|
tmdb_person_id: u64,
|
|
scenes_cs: &[(i64, i64)],
|
|
) -> anyhow::Result<()> {
|
|
tx.execute(
|
|
"INSERT INTO manifest_actors (manifest_id, tmdb_person_id) VALUES (?1, ?2)
|
|
ON CONFLICT (manifest_id, tmdb_person_id) DO NOTHING",
|
|
params![manifest_id, tmdb_person_id as i64],
|
|
)?;
|
|
let mut stmt = tx.prepare_cached(
|
|
"INSERT INTO scenes (manifest_id, tmdb_person_id, start_cs, end_cs)
|
|
VALUES (?1, ?2, ?3, ?4)",
|
|
)?;
|
|
for (start, end) in scenes_cs {
|
|
stmt.execute(params![manifest_id, tmdb_person_id as i64, start, end])?;
|
|
}
|
|
Ok(())
|
|
}
|
|
|
|
/// The actor payload of a stored manifest, reconstructed from rows.
|
|
///
|
|
/// §7: the submitted JSON is discarded; the document served to clients is
|
|
/// *reconstructed*, never echoed. Names come from `people`, populated from TMDB.
|
|
#[derive(Debug, Clone)]
|
|
pub struct StoredActor {
|
|
pub tmdb_person_id: u64,
|
|
pub name: Option<String>,
|
|
pub scenes_cs: Vec<(i64, i64)>,
|
|
}
|
|
|
|
/// TRACES: UR-010 | DR-002 | SR-001
|
|
pub fn actors_for_manifest(
|
|
conn: &Connection,
|
|
manifest_id: &str,
|
|
) -> anyhow::Result<Vec<StoredActor>> {
|
|
let mut stmt = conn.prepare_cached(
|
|
"SELECT ma.tmdb_person_id, p.name
|
|
FROM manifest_actors ma
|
|
LEFT JOIN people p ON p.tmdb_person_id = ma.tmdb_person_id
|
|
WHERE ma.manifest_id = ?1
|
|
ORDER BY ma.tmdb_person_id ASC",
|
|
)?;
|
|
let people = stmt
|
|
.query_map(params![manifest_id], |r| {
|
|
Ok((r.get::<_, i64>(0)? as u64, r.get::<_, Option<String>>(1)?))
|
|
})?
|
|
.collect::<rusqlite::Result<Vec<_>>>()?;
|
|
|
|
let mut scene_stmt = conn.prepare_cached(
|
|
"SELECT start_cs, end_cs FROM scenes
|
|
WHERE manifest_id = ?1 AND tmdb_person_id = ?2
|
|
ORDER BY start_cs ASC",
|
|
)?;
|
|
|
|
let mut out = Vec::with_capacity(people.len());
|
|
for (id, name) in people {
|
|
let scenes_cs = scene_stmt
|
|
.query_map(params![manifest_id, id as i64], |r| Ok((r.get(0)?, r.get(1)?)))?
|
|
.collect::<rusqlite::Result<Vec<_>>>()?;
|
|
out.push(StoredActor { tmdb_person_id: id, name, scenes_cs });
|
|
}
|
|
Ok(out)
|
|
}
|
|
|
|
pub fn set_manifest_status(
|
|
tx: &Transaction<'_>,
|
|
id: &str,
|
|
status: &str,
|
|
reason: Option<&str>,
|
|
cast_match_ratio: Option<f64>,
|
|
) -> anyhow::Result<()> {
|
|
tx.execute(
|
|
"UPDATE manifests
|
|
SET status = ?2, reject_reason = ?3, cast_match_ratio = COALESCE(?4, cast_match_ratio)
|
|
WHERE id = ?1",
|
|
params![id, status, reason, cast_match_ratio],
|
|
)?;
|
|
Ok(())
|
|
}
|
|
|
|
/// §6 stage 3: unmatched actors are dropped rather than stored, which is what
|
|
/// closes the free-text channel described in §5a.
|
|
pub fn delete_manifest_actor(
|
|
tx: &Transaction<'_>,
|
|
manifest_id: &str,
|
|
tmdb_person_id: u64,
|
|
) -> anyhow::Result<()> {
|
|
tx.execute(
|
|
"DELETE FROM scenes WHERE manifest_id = ?1 AND tmdb_person_id = ?2",
|
|
params![manifest_id, tmdb_person_id as i64],
|
|
)?;
|
|
tx.execute(
|
|
"DELETE FROM manifest_actors WHERE manifest_id = ?1 AND tmdb_person_id = ?2",
|
|
params![manifest_id, tmdb_person_id as i64],
|
|
)?;
|
|
Ok(())
|
|
}
|
|
|
|
/// §6 stage 3: a rejected manifest is deleted, not merely marked.
|
|
pub fn delete_manifest(tx: &Transaction<'_>, id: &str) -> anyhow::Result<()> {
|
|
tx.execute("DELETE FROM scenes WHERE manifest_id = ?1", params![id])?;
|
|
tx.execute("DELETE FROM manifest_actors WHERE manifest_id = ?1", params![id])?;
|
|
tx.execute("DELETE FROM manifests WHERE id = ?1", params![id])?;
|
|
Ok(())
|
|
}
|
|
|
|
pub fn manifest_actor_ids(conn: &Connection, manifest_id: &str) -> anyhow::Result<Vec<u64>> {
|
|
let mut stmt = stmt_actor_ids(conn)?;
|
|
let ids = stmt
|
|
.query_map(params![manifest_id], |r| Ok(r.get::<_, i64>(0)? as u64))?
|
|
.collect::<rusqlite::Result<Vec<_>>>()?;
|
|
Ok(ids)
|
|
}
|
|
|
|
fn stmt_actor_ids(conn: &Connection) -> rusqlite::Result<rusqlite::CachedStatement<'_>> {
|
|
conn.prepare_cached(
|
|
"SELECT tmdb_person_id FROM manifest_actors WHERE manifest_id = ?1 ORDER BY tmdb_person_id",
|
|
)
|
|
}
|
|
|
|
pub fn report_count(conn: &Connection, manifest_id: &str) -> anyhow::Result<i64> {
|
|
let mut stmt = conn.prepare_cached("SELECT COUNT(*) FROM reports WHERE manifest_id = ?1")?;
|
|
Ok(stmt.query_row(params![manifest_id], |r| r.get(0))?)
|
|
}
|
|
|
|
// ---------------------------------------------------------------------------
|
|
// Contributors
|
|
// ---------------------------------------------------------------------------
|
|
|
|
#[derive(Debug, Clone)]
|
|
pub struct Contributor {
|
|
pub id: String,
|
|
pub revoked: bool,
|
|
pub accepted_count: i64,
|
|
pub rejected_count: i64,
|
|
pub flagged_count: i64,
|
|
}
|
|
|
|
pub fn contributor_by_token_hash(
|
|
conn: &Connection,
|
|
token_hash: &str,
|
|
) -> anyhow::Result<Option<Contributor>> {
|
|
let mut stmt = conn.prepare_cached(
|
|
"SELECT id, revoked_at, accepted_count, rejected_count, flagged_count
|
|
FROM contributors WHERE token_hash = ?1",
|
|
)?;
|
|
Ok(stmt
|
|
.query_row(params![token_hash], |r| {
|
|
Ok(Contributor {
|
|
id: r.get(0)?,
|
|
revoked: r.get::<_, Option<String>>(1)?.is_some(),
|
|
accepted_count: r.get(2)?,
|
|
rejected_count: r.get(3)?,
|
|
flagged_count: r.get(4)?,
|
|
})
|
|
})
|
|
.optional()?)
|
|
}
|
|
|
|
pub fn insert_contributor(
|
|
tx: &Transaction<'_>,
|
|
token_hash: &str,
|
|
now: &str,
|
|
) -> anyhow::Result<String> {
|
|
let id = ulid::Ulid::new().to_string();
|
|
tx.execute(
|
|
"INSERT INTO contributors (id, token_hash, created_at) VALUES (?1, ?2, ?3)",
|
|
params![id, token_hash, now],
|
|
)?;
|
|
Ok(id)
|
|
}
|
|
|
|
/// §5a: the only reputational state is per-token counters.
|
|
pub fn bump_contributor_counter(
|
|
tx: &Transaction<'_>,
|
|
contributor_id: &str,
|
|
which: &str,
|
|
) -> anyhow::Result<()> {
|
|
// Column name is from a closed set, never from input.
|
|
let sql = match which {
|
|
"accepted" => "UPDATE contributors SET accepted_count = accepted_count + 1 WHERE id = ?1",
|
|
"rejected" => "UPDATE contributors SET rejected_count = rejected_count + 1 WHERE id = ?1",
|
|
"flagged" => "UPDATE contributors SET flagged_count = flagged_count + 1 WHERE id = ?1",
|
|
other => anyhow::bail!("unknown contributor counter {other}"),
|
|
};
|
|
tx.execute(sql, params![contributor_id])?;
|
|
Ok(())
|
|
}
|
|
|
|
/// §5a: a token whose rejection rate exceeds a threshold over a minimum sample
|
|
/// is revoked automatically, and its `pending`/`flagged` manifests are dropped.
|
|
/// No human is in the loop for the common case.
|
|
pub const REVOKE_MIN_SAMPLE: i64 = 20;
|
|
pub const REVOKE_REJECTION_RATE: f64 = 0.5;
|
|
|
|
pub fn maybe_revoke_contributor(
|
|
tx: &Transaction<'_>,
|
|
contributor_id: &str,
|
|
now: &str,
|
|
) -> anyhow::Result<bool> {
|
|
let (accepted, rejected, revoked): (i64, i64, Option<String>) = tx.query_row(
|
|
"SELECT accepted_count, rejected_count, revoked_at FROM contributors WHERE id = ?1",
|
|
params![contributor_id],
|
|
|r| Ok((r.get(0)?, r.get(1)?, r.get(2)?)),
|
|
)?;
|
|
if revoked.is_some() {
|
|
return Ok(false);
|
|
}
|
|
let total = accepted + rejected;
|
|
if total < REVOKE_MIN_SAMPLE {
|
|
return Ok(false);
|
|
}
|
|
if (rejected as f64) / (total as f64) <= REVOKE_REJECTION_RATE {
|
|
return Ok(false);
|
|
}
|
|
|
|
tx.execute(
|
|
"UPDATE contributors SET revoked_at = ?2 WHERE id = ?1",
|
|
params![contributor_id, now],
|
|
)?;
|
|
// Drop the token's pending/flagged manifests.
|
|
let ids: Vec<String> = {
|
|
let mut stmt = tx.prepare(
|
|
"SELECT id FROM manifests WHERE contributor_id = ?1 AND status IN ('pending','flagged')",
|
|
)?;
|
|
let rows = stmt
|
|
.query_map(params![contributor_id], |r| r.get::<_, String>(0))?
|
|
.collect::<rusqlite::Result<Vec<_>>>()?;
|
|
rows
|
|
};
|
|
for id in &ids {
|
|
delete_manifest(tx, id)?;
|
|
}
|
|
Ok(true)
|
|
}
|
|
|
|
// ---------------------------------------------------------------------------
|
|
// Reports
|
|
// ---------------------------------------------------------------------------
|
|
|
|
pub fn insert_report(
|
|
tx: &Transaction<'_>,
|
|
manifest_id: &str,
|
|
reason: &str,
|
|
note: Option<&str>,
|
|
source_ip_hash: &str,
|
|
now: &str,
|
|
) -> anyhow::Result<String> {
|
|
let id = ulid::Ulid::new().to_string();
|
|
tx.execute(
|
|
"INSERT INTO reports (id, manifest_id, reason, note, created_at, source_ip_hash)
|
|
VALUES (?1,?2,?3,?4,?5,?6)",
|
|
params![id, manifest_id, reason, note, now, source_ip_hash],
|
|
)?;
|
|
Ok(id)
|
|
}
|
|
|
|
// ---------------------------------------------------------------------------
|
|
// TMDB cache — the sole JSON column, holding TMDB's responses, not users' (§7)
|
|
// ---------------------------------------------------------------------------
|
|
|
|
pub fn cached_credits(
|
|
conn: &Connection,
|
|
tmdb_id: &str,
|
|
kind: &str,
|
|
) -> anyhow::Result<Option<(String, String)>> {
|
|
let mut stmt = conn.prepare_cached(
|
|
"SELECT credits, fetched_at FROM tmdb_cache WHERE tmdb_id = ?1 AND kind = ?2",
|
|
)?;
|
|
Ok(stmt.query_row(params![tmdb_id, kind], |r| Ok((r.get(0)?, r.get(1)?))).optional()?)
|
|
}
|
|
|
|
pub fn put_credits(
|
|
tx: &Transaction<'_>,
|
|
tmdb_id: &str,
|
|
kind: &str,
|
|
credits: &str,
|
|
now: &str,
|
|
) -> anyhow::Result<()> {
|
|
tx.execute(
|
|
"INSERT INTO tmdb_cache (tmdb_id, kind, credits, fetched_at) VALUES (?1,?2,?3,?4)
|
|
ON CONFLICT (tmdb_id, kind) DO UPDATE
|
|
SET credits = excluded.credits, fetched_at = excluded.fetched_at",
|
|
params![tmdb_id, kind, credits, now],
|
|
)?;
|
|
Ok(())
|
|
}
|
|
|
|
// ---------------------------------------------------------------------------
|
|
// Jobs — a table rather than an external broker, so work survives restart (§7)
|
|
// ---------------------------------------------------------------------------
|
|
|
|
#[derive(Debug, Clone)]
|
|
pub struct Job {
|
|
pub id: String,
|
|
pub kind: String,
|
|
pub payload: String,
|
|
pub attempts: i64,
|
|
}
|
|
|
|
pub fn enqueue_job(
|
|
tx: &Transaction<'_>,
|
|
kind: &str,
|
|
payload: &str,
|
|
run_after: &str,
|
|
) -> anyhow::Result<String> {
|
|
let id = ulid::Ulid::new().to_string();
|
|
tx.execute(
|
|
"INSERT INTO jobs (id, kind, payload, run_after) VALUES (?1,?2,?3,?4)",
|
|
params![id, kind, payload, run_after],
|
|
)?;
|
|
Ok(id)
|
|
}
|
|
|
|
/// Claims up to `limit` due jobs, marking them leased so a second worker tick
|
|
/// cannot pick up the same work.
|
|
/// TRACES: DR-005 | PR-004
|
|
pub fn lease_jobs(tx: &Transaction<'_>, now: &str, limit: usize) -> anyhow::Result<Vec<Job>> {
|
|
let jobs: Vec<Job> = {
|
|
let mut stmt = tx.prepare(
|
|
"SELECT id, kind, payload, attempts FROM jobs
|
|
WHERE run_after <= ?1 AND leased_at IS NULL
|
|
ORDER BY run_after ASC
|
|
LIMIT ?2",
|
|
)?;
|
|
let rows = stmt
|
|
.query_map(params![now, limit as i64], |r| {
|
|
Ok(Job { id: r.get(0)?, kind: r.get(1)?, payload: r.get(2)?, attempts: r.get(3)? })
|
|
})?
|
|
.collect::<rusqlite::Result<Vec<_>>>()?;
|
|
rows
|
|
};
|
|
for j in &jobs {
|
|
tx.execute("UPDATE jobs SET leased_at = ?2 WHERE id = ?1", params![j.id, now])?;
|
|
}
|
|
Ok(jobs)
|
|
}
|
|
|
|
pub fn delete_job(tx: &Transaction<'_>, id: &str) -> anyhow::Result<()> {
|
|
tx.execute("DELETE FROM jobs WHERE id = ?1", params![id])?;
|
|
Ok(())
|
|
}
|
|
|
|
/// Releases a leased job for a later retry with backoff (§6 stage 3: TMDB
|
|
/// unreachable means retry, not reject).
|
|
pub fn reschedule_job(
|
|
tx: &Transaction<'_>,
|
|
id: &str,
|
|
run_after: &str,
|
|
error: &str,
|
|
) -> anyhow::Result<()> {
|
|
tx.execute(
|
|
"UPDATE jobs SET attempts = attempts + 1, last_error = ?3, run_after = ?2, leased_at = NULL
|
|
WHERE id = ?1",
|
|
params![id, run_after, error],
|
|
)?;
|
|
Ok(())
|
|
}
|
|
|
|
/// Frees leases held at startup — a process that died mid-job would otherwise
|
|
/// leave work stranded.
|
|
pub fn release_all_leases(tx: &Transaction<'_>) -> anyhow::Result<usize> {
|
|
let n = tx.execute("UPDATE jobs SET leased_at = NULL WHERE leased_at IS NOT NULL", [])?;
|
|
Ok(n)
|
|
}
|
|
|
|
#[cfg(test)]
|
|
mod tests {
|
|
use super::*;
|
|
use crate::db::Db;
|
|
|
|
async fn db() -> Db {
|
|
Db::open(":memory:").unwrap()
|
|
}
|
|
|
|
const NOW: &str = "2026-07-30T12:00:00Z";
|
|
|
|
#[tokio::test]
|
|
async fn upsert_title_is_idempotent_and_backfills() {
|
|
let db = db().await;
|
|
let (a, b) = db
|
|
.write(|tx| {
|
|
let a = upsert_title(
|
|
tx,
|
|
IdentityType::Movie,
|
|
Some("504172"),
|
|
None,
|
|
Some("X"),
|
|
Some(2017),
|
|
NOW,
|
|
)?;
|
|
// Second upload supplies the IMDB id the first lacked.
|
|
let b = upsert_title(
|
|
tx,
|
|
IdentityType::Movie,
|
|
Some("504172"),
|
|
Some("tt4686844"),
|
|
None,
|
|
None,
|
|
NOW,
|
|
)?;
|
|
Ok((a, b))
|
|
})
|
|
.await
|
|
.unwrap();
|
|
assert_eq!(a, b, "same title must not be duplicated");
|
|
|
|
let found = db
|
|
.write(|tx| find_title(tx, IdentityType::Movie, Some("504172"), None))
|
|
.await
|
|
.unwrap()
|
|
.unwrap();
|
|
assert_eq!(found.imdb_id.as_deref(), Some("tt4686844"), "identifier should be backfilled");
|
|
}
|
|
|
|
#[tokio::test]
|
|
async fn movie_and_series_titles_do_not_collide_on_id() {
|
|
// TMDB numbers movies and series in separate spaces, so id 1396 is both
|
|
// a film and Breaking Bad.
|
|
let db = db().await;
|
|
let (m, s) = db
|
|
.write(|tx| {
|
|
let m = upsert_title(tx, IdentityType::Movie, Some("1396"), None, None, None, NOW)?;
|
|
let s =
|
|
upsert_title(tx, IdentityType::Episode, Some("1396"), None, None, None, NOW)?;
|
|
Ok((m, s))
|
|
})
|
|
.await
|
|
.unwrap();
|
|
assert_ne!(m, s);
|
|
}
|
|
|
|
#[tokio::test]
|
|
async fn manifest_roundtrips_through_rows() {
|
|
let db = db().await;
|
|
let id = db
|
|
.write(|tx| {
|
|
let title_id =
|
|
upsert_title(tx, IdentityType::Movie, Some("1"), None, None, None, NOW)?;
|
|
let cid = insert_contributor(tx, "hash", NOW)?;
|
|
let id = ulid::Ulid::new().to_string();
|
|
insert_manifest(
|
|
tx,
|
|
&NewManifest {
|
|
id: &id,
|
|
title_id: &title_id,
|
|
season: None,
|
|
episode: None,
|
|
runtime_sec: 6420.5,
|
|
video_hash: Some("opensubtitles:8e245d9679d31e12"),
|
|
audio_signature: None,
|
|
audio_sig_coarse: None,
|
|
sample_fps: Some(5.0),
|
|
extinction_sec: Some(12.0),
|
|
gallery_scope: Some("global"),
|
|
pipeline_version: Some("test 0.1"),
|
|
contributor_id: Some(&cid),
|
|
status: "listed",
|
|
content_id: Some("sha256:abc"),
|
|
origin: "local",
|
|
ingested_from: None,
|
|
created_at: NOW,
|
|
},
|
|
)?;
|
|
upsert_person(tx, 884, "Steve Buscemi", false, NOW)?;
|
|
insert_actor_scenes(tx, &id, 884, &[(19160, 20920), (43820, 46560)])?;
|
|
Ok(id)
|
|
})
|
|
.await
|
|
.unwrap();
|
|
|
|
let (row, actors) = db
|
|
.write(move |tx| {
|
|
let row = manifest_by_id(tx, &id)?.unwrap();
|
|
let actors = actors_for_manifest(tx, &id)?;
|
|
Ok((row, actors))
|
|
})
|
|
.await
|
|
.unwrap();
|
|
|
|
assert_eq!(row.runtime_sec, 6420.5);
|
|
assert_eq!(actors.len(), 1);
|
|
assert_eq!(actors[0].tmdb_person_id, 884);
|
|
// §7: names come from `people`, populated from TMDB, never from upload.
|
|
assert_eq!(actors[0].name.as_deref(), Some("Steve Buscemi"));
|
|
assert_eq!(actors[0].scenes_cs, vec![(19160, 20920), (43820, 46560)]);
|
|
}
|
|
|
|
#[tokio::test]
|
|
async fn only_served_statuses_are_candidates() {
|
|
let db = db().await;
|
|
let title_id = db
|
|
.write(|tx| {
|
|
let title_id =
|
|
upsert_title(tx, IdentityType::Movie, Some("7"), None, None, None, NOW)?;
|
|
for (i, status) in ["listed", "flagged", "pending", "rejected"].iter().enumerate() {
|
|
insert_manifest(
|
|
tx,
|
|
&NewManifest {
|
|
id: &format!("m{i}"),
|
|
title_id: &title_id,
|
|
season: None,
|
|
episode: None,
|
|
runtime_sec: 100.0,
|
|
video_hash: None,
|
|
audio_signature: None,
|
|
audio_sig_coarse: None,
|
|
sample_fps: None,
|
|
extinction_sec: None,
|
|
gallery_scope: None,
|
|
pipeline_version: None,
|
|
contributor_id: None,
|
|
status,
|
|
content_id: Some(&format!("c{i}")),
|
|
origin: "local",
|
|
ingested_from: None,
|
|
created_at: NOW,
|
|
},
|
|
)?;
|
|
}
|
|
Ok(title_id)
|
|
})
|
|
.await
|
|
.unwrap();
|
|
|
|
let rows =
|
|
db.write(move |tx| candidates_for_title(tx, &title_id, None, None)).await.unwrap();
|
|
// `pending` is held unlisted and not served to anyone (§6 stage 3);
|
|
// `rejected` never is.
|
|
assert_eq!(rows.len(), 2);
|
|
assert!(rows.iter().all(|r| r.status == "listed" || r.status == "flagged"));
|
|
}
|
|
|
|
#[tokio::test]
|
|
async fn duplicate_detection_is_per_contributor() {
|
|
let db = db().await;
|
|
let dup = db
|
|
.write(|tx| {
|
|
let title_id =
|
|
upsert_title(tx, IdentityType::Movie, Some("9"), None, None, None, NOW)?;
|
|
let a = insert_contributor(tx, "hash-a", NOW)?;
|
|
let b = insert_contributor(tx, "hash-b", NOW)?;
|
|
insert_manifest(
|
|
tx,
|
|
&NewManifest {
|
|
id: "m1",
|
|
title_id: &title_id,
|
|
season: None,
|
|
episode: None,
|
|
runtime_sec: 6420.5,
|
|
video_hash: None,
|
|
audio_signature: None,
|
|
audio_sig_coarse: None,
|
|
sample_fps: None,
|
|
extinction_sec: None,
|
|
gallery_scope: None,
|
|
pipeline_version: None,
|
|
contributor_id: Some(&a),
|
|
status: "listed",
|
|
content_id: Some("c1"),
|
|
origin: "local",
|
|
ingested_from: None,
|
|
created_at: NOW,
|
|
},
|
|
)?;
|
|
let same = duplicate_from_contributor(tx, &title_id, None, None, 6420.5, None, &a)?;
|
|
// §7: multiple manifests for the same cut from *different*
|
|
// contributors are allowed, and ranked.
|
|
let other =
|
|
duplicate_from_contributor(tx, &title_id, None, None, 6420.5, None, &b)?;
|
|
Ok((same.is_some(), other.is_some()))
|
|
})
|
|
.await
|
|
.unwrap();
|
|
assert_eq!(dup, (true, false));
|
|
}
|
|
|
|
#[tokio::test]
|
|
async fn revocation_needs_a_minimum_sample_then_drops_pending_work() {
|
|
let db = db().await;
|
|
let revoked_early = db
|
|
.write(|tx| {
|
|
let c = insert_contributor(tx, "h", NOW)?;
|
|
for _ in 0..5 {
|
|
bump_contributor_counter(tx, &c, "rejected")?;
|
|
}
|
|
// §5a: a threshold over a *minimum sample* — 5 rejections is not
|
|
// yet evidence.
|
|
let early = maybe_revoke_contributor(tx, &c, NOW)?;
|
|
|
|
for _ in 0..20 {
|
|
bump_contributor_counter(tx, &c, "rejected")?;
|
|
}
|
|
let title_id =
|
|
upsert_title(tx, IdentityType::Movie, Some("11"), None, None, None, NOW)?;
|
|
insert_manifest(
|
|
tx,
|
|
&NewManifest {
|
|
id: "pend",
|
|
title_id: &title_id,
|
|
season: None,
|
|
episode: None,
|
|
runtime_sec: 100.0,
|
|
video_hash: None,
|
|
audio_signature: None,
|
|
audio_sig_coarse: None,
|
|
sample_fps: None,
|
|
extinction_sec: None,
|
|
gallery_scope: None,
|
|
pipeline_version: None,
|
|
contributor_id: Some(&c),
|
|
status: "pending",
|
|
content_id: Some("cp"),
|
|
origin: "local",
|
|
ingested_from: None,
|
|
created_at: NOW,
|
|
},
|
|
)?;
|
|
let now_revoked = maybe_revoke_contributor(tx, &c, NOW)?;
|
|
let still_there = manifest_by_id(tx, "pend")?.is_some();
|
|
Ok((early, now_revoked, still_there))
|
|
})
|
|
.await
|
|
.unwrap();
|
|
assert_eq!(revoked_early, (false, true, false));
|
|
}
|
|
|
|
#[tokio::test]
|
|
async fn good_contributors_are_not_revoked() {
|
|
let db = db().await;
|
|
let revoked = db
|
|
.write(|tx| {
|
|
let c = insert_contributor(tx, "h", NOW)?;
|
|
for _ in 0..40 {
|
|
bump_contributor_counter(tx, &c, "accepted")?;
|
|
}
|
|
for _ in 0..3 {
|
|
bump_contributor_counter(tx, &c, "rejected")?;
|
|
}
|
|
maybe_revoke_contributor(tx, &c, NOW)
|
|
})
|
|
.await
|
|
.unwrap();
|
|
assert!(!revoked);
|
|
}
|
|
|
|
#[tokio::test]
|
|
async fn jobs_lease_once_and_reschedule() {
|
|
let db = db().await;
|
|
let (first, second, after_reschedule) = db
|
|
.write(|tx| {
|
|
enqueue_job(tx, "cast_check", "{}", NOW)?;
|
|
let first = lease_jobs(tx, NOW, 10)?;
|
|
// A leased job must not be handed out twice.
|
|
let second = lease_jobs(tx, NOW, 10)?;
|
|
reschedule_job(tx, &first[0].id, NOW, "tmdb unreachable")?;
|
|
let third = lease_jobs(tx, NOW, 10)?;
|
|
Ok((first.len(), second.len(), third))
|
|
})
|
|
.await
|
|
.unwrap();
|
|
assert_eq!((first, second), (1, 0));
|
|
assert_eq!(after_reschedule.len(), 1);
|
|
assert_eq!(after_reschedule[0].attempts, 1);
|
|
}
|
|
|
|
#[tokio::test]
|
|
async fn future_jobs_are_not_leased() {
|
|
let db = db().await;
|
|
let n = db
|
|
.write(|tx| {
|
|
enqueue_job(tx, "cast_check", "{}", "2099-01-01T00:00:00Z")?;
|
|
Ok(lease_jobs(tx, NOW, 10)?.len())
|
|
})
|
|
.await
|
|
.unwrap();
|
|
assert_eq!(n, 0);
|
|
}
|
|
|
|
#[tokio::test]
|
|
async fn startup_releases_stranded_leases() {
|
|
let db = db().await;
|
|
let n = db
|
|
.write(|tx| {
|
|
enqueue_job(tx, "cast_check", "{}", NOW)?;
|
|
lease_jobs(tx, NOW, 10)?;
|
|
release_all_leases(tx)?;
|
|
Ok(lease_jobs(tx, NOW, 10)?.len())
|
|
})
|
|
.await
|
|
.unwrap();
|
|
assert_eq!(n, 1);
|
|
}
|
|
|
|
#[tokio::test]
|
|
async fn deleting_a_manifest_removes_its_rows() {
|
|
let db = db().await;
|
|
let counts = db
|
|
.write(|tx| {
|
|
let title_id =
|
|
upsert_title(tx, IdentityType::Movie, Some("13"), None, None, None, NOW)?;
|
|
insert_manifest(
|
|
tx,
|
|
&NewManifest {
|
|
id: "m",
|
|
title_id: &title_id,
|
|
season: None,
|
|
episode: None,
|
|
runtime_sec: 100.0,
|
|
video_hash: None,
|
|
audio_signature: None,
|
|
audio_sig_coarse: None,
|
|
sample_fps: None,
|
|
extinction_sec: None,
|
|
gallery_scope: None,
|
|
pipeline_version: None,
|
|
contributor_id: None,
|
|
status: "listed",
|
|
content_id: Some("c"),
|
|
origin: "local",
|
|
ingested_from: None,
|
|
created_at: NOW,
|
|
},
|
|
)?;
|
|
upsert_person(tx, 1, "A", false, NOW)?;
|
|
insert_actor_scenes(tx, "m", 1, &[(0, 100)])?;
|
|
delete_manifest(tx, "m")?;
|
|
let scenes: i64 = tx.query_row("SELECT COUNT(*) FROM scenes", [], |r| r.get(0))?;
|
|
let actors: i64 =
|
|
tx.query_row("SELECT COUNT(*) FROM manifest_actors", [], |r| r.get(0))?;
|
|
let manifests: i64 =
|
|
tx.query_row("SELECT COUNT(*) FROM manifests", [], |r| r.get(0))?;
|
|
Ok((scenes, actors, manifests))
|
|
})
|
|
.await
|
|
.unwrap();
|
|
assert_eq!(counts, (0, 0, 0));
|
|
}
|
|
}
|