diff --git a/README.md b/README.md index 06ec513..1c98f9e 100644 --- a/README.md +++ b/README.md @@ -38,6 +38,11 @@ Implemented: - §7 relational storage, no JSON blobs on the write path - §8 Rust + Axum + SQLite, single serialized writer, in-process job queue - §9a content addressing (`content_id`), computed on upload +- §9a federation — change feed, fetch by `content_id`, batch `have`, peer + directory, capabilities, and a pull worker that **re-derives judgement rather + than inheriting it**: every pulled manifest runs the full §6 validation and + this server's own cast check, and the body is verified to hash to the + `content_id` requested before it is stored **Schema version 2** (SR-003). `jmanifest_version` moved to 2 in lockstep with the truth file's `schema_version` — breaking changes are batched and ship @@ -75,9 +80,7 @@ Deferred: `POST /manifests/search` are not wired up. This follows §3's own recommended sequencing: ship the plugin-side computation first, let signatures accumulate, then enable matching once coverage is useful. -- §9a federation endpoints (`/federation/*`) and the pull worker. The schema - columns (`content_id`, `origin`, `ingested_from`, `peers`) are in place, and - `ingest::persist` is already the shared path a pull would reuse. +(Federation landed — see below.) ## Running @@ -124,7 +127,7 @@ curl -sX POST -H 'content-type: application/json' -d '{}' \ ## Tests ```sh -cargo test # 191 tests +cargo test # 212 tests cargo deny check # advisories, licences, bans, sources scripts/traceability-gate.sh # requirement coverage ``` @@ -137,9 +140,8 @@ git submodule update --init --recursive It reports coverage against [`docs/requirements.md`](docs/requirements.md), flags orphan tags (an ID no register defines), and fails on a >100% ratio — the -signal that the computation itself is broken. Currently **24/32 (75%)**; the -untraced remainder is UR-007 (plugin-side) and UR-008 (federation), neither of -which is implemented here yet. +signal that the computation itself is broken. Currently **25/32 (78%)**; the +untraced remainder is UR-007, which is plugin-side. Unit tests per module, plus two integration suites: diff --git a/docs/requirements.md b/docs/requirements.md index f238244..a1b231c 100644 --- a/docs/requirements.md +++ b/docs/requirements.md @@ -34,7 +34,7 @@ requirement and no fixture-generation step, unlike `scene-actor-extraction`. | UR-005 | Trust without accounts: not usable as a content store, nor for prank manifests | SR-004 | High | Done | | UR-006 | Serve and accept a whole series in one operation | PR-006 | High | Done | | UR-007 | Plugin queries an ordered, configurable list of servers | PR-005 | High | In Progress | -| UR-008 | Servers replicate manifests between each other | PR-006 | Medium | Planned | +| UR-008 | Servers replicate manifests between each other | PR-006 | Medium | Done | | UR-009 | Store an audio spectral-peak signature for content-based identification | SR-003 | Medium | In Progress | | UR-010 | Identity crossing the API boundary is TMDB/IMDB ids, never a name alone | SR-001 | High | Done | | UR-011 | Reject any field capable of carrying binary or attacker-chosen content | SR-004 | High | Done | @@ -52,6 +52,13 @@ requirement and no fixture-generation step, unlike `scene-actor-extraction`. server list and its per-server trust settings, with the community instance pre-configured but disabled. The fetch path that consumes it does not exist yet. +**UR-008 is `Done` for the replication surface**: change feed, fetch by +`content_id`, batch `have`, the peer directory, and a pull worker that +re-validates everything it ingests. Peer *administration* — adding and enabling +a peer — is deliberately not an API: §9a requires that a peering exist only +because an operator typed a URL, so it is a database action, and +`there_is_no_endpoint_that_creates_a_peering` asserts the absence. + **UR-009 is `In Progress`.** The server accepts, validates and stores `cut.audio_signature`, and `content_id` correctly excludes it (§9a). What is absent is `audio`-tier matching and `POST /manifests/search`. This is the @@ -118,7 +125,8 @@ topology is the point, so this is a deliberate choice rather than an oversight. | UR-004 | T1 + T2 | Limits engage and carry the documented headers | Window reset; a rejected request does not extend its own lockout; surfaces have independent budgets | | UR-005 | T1 + T2 | Prank manifests rejected; no free-text channel | Uncredited cast rejected; name-only matches capped; automatic revocation needs a minimum sample | | UR-006 | T2 | Bundle accepted per-episode, non-atomically | One bad episode rejected while its neighbours are accepted; envelope errors are whole-request `400` | -| UR-007 | — | *No server-side test.* Plugin-side; the register there will carry it | — | +| UR-007 | — | *No server-side test.* Plugin-side (`jRay` JR-025); the register there carries it | — | +| UR-008 | T1 + T2 | Feed, fetch-by-hash, batch have, peer directory | **Cursor is strictly monotonic** — a ULID would sort out of write order within a millisecond and silently skip entries; a peer retraction flags rather than delists; only the opt-in abuse channel delists; `pending` is never replicated; **no endpoint can create a peering** | | UR-009 | T1 | Signature structurally validated | Fixed length; reserved high bit; **media < 120 s must send no signature at all** | | UR-010 | T1 + T2 | Actors persist as TMDB person ids | A name the upload invented does not round-trip | | UR-011 | T2 | Every payload-shaped field rejected | base64, hex, markup, control characters, bidi overrides, compatibility homoglyphs | diff --git a/docs/traceability.md b/docs/traceability.md index e4af3a6..0c9d06a 100644 --- a/docs/traceability.md +++ b/docs/traceability.md @@ -3,7 +3,7 @@ -**Generated:** 2026-07-31T07:12:27+00:00 +**Generated:** 2026-07-31T07:28:04+00:00 Denominators are read from [`requirements.md`](requirements.md) at run time, never hardcoded. Coverage counts a requirement only when it is tagged in source **and** has a verification tier this repo's CI host can execute (`T1, T2, static`). @@ -11,13 +11,13 @@ Denominators are read from [`requirements.md`](requirements.md) at run time, nev | Metric | Value | |---|---| -| Source files scanned | 26 | -| TRACES tags found | 37 | +| Source files scanned | 29 | +| TRACES tags found | 43 | | EXCEPTION tags found | 0 | | Requirements defined | 32 | -| Requirements covered | 24 | -| **Coverage** | **75.0%** (24/32) | -| Coverage of CI-executable scope | 75.0% (24/32) | +| Requirements covered | 25 | +| **Coverage** | **78.1%** (25/32) | +| Coverage of CI-executable scope | 78.1% (25/32) | | Tagged but unexecuted in CI | 0 | | Orphan tags | 0 | @@ -25,7 +25,7 @@ Denominators are read from [`requirements.md`](requirements.md) at run time, nev | Type | Covered | Tagged but unexecuted | Defined | |---|---|---|---| -| UR | 14 | 0 | 18 | +| UR | 15 | 0 | 18 | | DR | 10 | 0 | 14 | - **PR** tags present (separate taxonomy, not counted in coverage): PR-004, PR-005, PR-006 @@ -66,13 +66,13 @@ _None._ | UR-005 | Done | T1, T2 | SR-004 | covered | `src/api/report.rs`, `src/api/upload.rs`, `src/auth.rs`, `src/castcheck.rs`, `src/worker.rs` | Trust without accounts: not usable as a content store, nor for prank … | | UR-006 | Done | T2 | PR-006 | covered | `src/api/fetch.rs`, `src/api/upload.rs`, `src/validate.rs` | Serve and accept a whole series in one operation | | UR-007 | In Progress | unset | PR-005 | covered | `src/api/exists.rs` | Plugin queries an ordered, configurable list of servers | -| UR-008 | Planned | unset | PR-006 | untagged | - | Servers replicate manifests between each other | +| UR-008 | Planned | unset | PR-006 | covered | `src/api/federation.rs`, `src/worker.rs` | Servers replicate manifests between each other | | UR-009 | In Progress | T1 | SR-003 | covered | `src/validate.rs` | Store an audio spectral-peak signature for content-based identificati… | | UR-010 | Done | T1, T2 | SR-001 | covered | `src/api/fetch.rs`, `src/castcheck.rs`, `src/db/repo.rs`, `src/model.rs` | Identity crossing the API boundary is TMDB/IMDB ids, never a name alo… | | UR-011 | Done | T2 | SR-004 | covered | `src/model.rs`, `src/validate.rs` | Reject any field capable of carrying binary or attacker-chosen content | | UR-012 | Done | T2 | SR-005 | covered | `src/db/repo.rs`, `src/ingest.rs` | Never accept, store, or serve gallery data — reference faces or embed… | | UR-013 | Done | T1 | SR-002 | covered | `src/api/fetch.rs`, `src/model.rs`, `src/validate.rs` | Windows are scene-scoped claims; never reinterpret their boundaries | -| UR-014 | Done | T1 | SR-003 | covered | `src/model.rs`, `src/validate.rs` | Reject an unknown `jmanifest_version` outright, never guess | +| UR-014 | Done | T1 | SR-003 | covered | `src/api/federation.rs`, `src/model.rs`, `src/validate.rs` | Reject an unknown `jmanifest_version` outright, never guess | | UR-015 | Done | T2 | SR-003 | untagged | - | Accept `extraction.extinction_sec` in place of `anneal_sec` | | UR-016 | Done | T2 | SR-003 | untagged | - | Accept and store `extraction.gallery_scope`; rank on it (§7) | | UR-017 | Done | T1, T2 | SR-003 | covered | `src/model.rs` | Accept per-window belief and identification route; `scenes` are objec… | @@ -98,7 +98,7 @@ _None._ **Locations:** 3 -- [`src/api/fetch.rs:267`](../src/api/fetch.rs#L267) — `pub fn reconstruct(` +- [`src/api/fetch.rs:270`](../src/api/fetch.rs#L270) — `pub fn reconstruct(` - [`src/db/repo.rs:337`](../src/db/repo.rs#L337) — `pub fn insert_manifest(tx: &Transaction<'_>, m: &NewManifest<'_>) -> anyhow::Result<()>` - [`src/db/repo.rs:412`](../src/db/repo.rs#L412) — `pub fn actors_for_manifest(` @@ -124,7 +124,7 @@ _None._ **Locations:** 1 -- [`src/ratelimit.rs:76`](../src/ratelimit.rs#L76) — `impl Default for RateLimiter` +- [`src/ratelimit.rs:93`](../src/ratelimit.rs#L93) — `impl Default for RateLimiter` ### DR-008 @@ -171,19 +171,24 @@ _None._ ### PR-005 -**Locations:** 1 +**Locations:** 2 - [`src/api/exists.rs:74`](../src/api/exists.rs#L74) — `pub async fn exists_batch(` +- [`src/api/federation.rs:212`](../src/api/federation.rs#L212) — `pub async fn get_peers(` ### PR-006 -**Locations:** 5 +**Locations:** 9 +- [`src/api/federation.rs:66`](../src/api/federation.rs#L66) — `pub async fn get_changes(` +- [`src/api/federation.rs:108`](../src/api/federation.rs#L108) — `pub async fn get_manifest_by_content_id(` +- [`src/api/federation.rs:160`](../src/api/federation.rs#L160) — `pub async fn post_have(` - [`src/api/fetch.rs:128`](../src/api/fetch.rs#L128) — `pub async fn get_series(` - [`src/api/upload.rs:29`](../src/api/upload.rs#L29) — `pub async fn post_manifest(` - [`src/api/upload.rs:101`](../src/api/upload.rs#L101) — `pub async fn post_bundle(` - [`src/ingest.rs:50`](../src/ingest.rs#L50) — `pub fn persist(` - [`src/validate.rs:586`](../src/validate.rs#L586) — `pub fn validate_bundle_envelope(b: &SeriesBundle) -> VResult<()>` +- [`src/worker.rs:107`](../src/worker.rs#L107) — `async fn run_federation_pull(&self, payload: &str) -> Result<(), JobError>` ### SR-001 @@ -191,7 +196,7 @@ _None._ - [`src/api/exists.rs:61`](../src/api/exists.rs#L61) — `pub async fn exists(` - [`src/api/exists.rs:74`](../src/api/exists.rs#L74) — `pub async fn exists_batch(` -- [`src/api/fetch.rs:267`](../src/api/fetch.rs#L267) — `pub fn reconstruct(` +- [`src/api/fetch.rs:270`](../src/api/fetch.rs#L270) — `pub fn reconstruct(` - [`src/castcheck.rs:82`](../src/castcheck.rs#L82) — `pub fn evaluate(submitted: &[SubmittedActor], credits: &[CastMember]) -> CastCheckOutcome` - [`src/db/repo.rs:412`](../src/db/repo.rs#L412) — `pub fn actors_for_manifest(` - [`src/matching.rs:52`](../src/matching.rs#L52) — `pub fn match_cut(client: &ClientCut, stored: &StoredCut) -> Option` @@ -201,22 +206,23 @@ _None._ **Locations:** 4 -- [`src/api/fetch.rs:267`](../src/api/fetch.rs#L267) — `pub fn reconstruct(` +- [`src/api/fetch.rs:270`](../src/api/fetch.rs#L270) — `pub fn reconstruct(` - [`src/model.rs:179`](../src/model.rs#L179) — `Unknown` -- [`src/model.rs:233`](../src/model.rs#L233) — `pub fn from_str(s: &str) -> Option` +- [`src/model.rs:238`](../src/model.rs#L238) — `pub fn from_stored(s: &str) -> Option` - [`src/validate.rs:510`](../src/validate.rs#L510) — `fn validate_scenes(idx: usize, a: &Actor, runtime_sec: f64) -> VResult>` ### SR-003 -**Locations:** 10 +**Locations:** 11 +- [`src/api/federation.rs:258`](../src/api/federation.rs#L258) — `pub async fn get_capabilities(State(state): State) -> ApiResult` - [`src/api/json.rs:100`](../src/api/json.rs#L100) — `fn require_utf8(bytes: &[u8]) -> Result<&str, ApiError>` - [`src/content_id.rs:52`](../src/content_id.rs#L52) — `pub fn canonical_json(` - [`src/content_id.rs:129`](../src/content_id.rs#L129) — `pub fn content_id(` - [`src/error.rs:8`](../src/error.rs#L8) — `Unknown` - [`src/model.rs:197`](../src/model.rs#L197) — `Unknown` -- [`src/model.rs:233`](../src/model.rs#L233) — `pub fn from_str(s: &str) -> Option` -- [`src/model.rs:259`](../src/model.rs#L259) — `Unknown` +- [`src/model.rs:238`](../src/model.rs#L238) — `pub fn from_stored(s: &str) -> Option` +- [`src/model.rs:264`](../src/model.rs#L264) — `Unknown` - [`src/validate.rs:102`](../src/validate.rs#L102) — `pub fn to_centiseconds(secs: f64) -> i64` - [`src/validate.rs:227`](../src/validate.rs#L227) — `pub fn validate_manifest(mut m: Jmanifest) -> VResult` - [`src/validate.rs:378`](../src/validate.rs#L378) — `pub fn validate_audio_signature(sig: &str, runtime_sec: f64) -> VResult<()>` @@ -235,12 +241,12 @@ _None._ - [`src/castcheck.rs:215`](../src/castcheck.rs#L215) — `pub fn category_guard_violation(matched: &[MatchedActor], title_is_adult: bool) -> Option…` - [`src/db/repo.rs:337`](../src/db/repo.rs#L337) — `pub fn insert_manifest(tx: &Transaction<'_>, m: &NewManifest<'_>) -> anyhow::Result<()>` - [`src/model.rs:156`](../src/model.rs#L156) — `Unknown` -- [`src/model.rs:259`](../src/model.rs#L259) — `Unknown` -- [`src/ratelimit.rs:76`](../src/ratelimit.rs#L76) — `impl Default for RateLimiter` +- [`src/model.rs:264`](../src/model.rs#L264) — `Unknown` +- [`src/ratelimit.rs:93`](../src/ratelimit.rs#L93) — `impl Default for RateLimiter` - [`src/validate.rs:136`](../src/validate.rs#L136) — `fn is_allowed_text_char(c: char) -> bool` - [`src/validate.rs:227`](../src/validate.rs#L227) — `pub fn validate_manifest(mut m: Jmanifest) -> VResult` - [`src/validate.rs:378`](../src/validate.rs#L378) — `pub fn validate_audio_signature(sig: &str, runtime_sec: f64) -> VResult<()>` -- [`src/worker.rs:96`](../src/worker.rs#L96) — `async fn run_cast_check(&self, payload: &str) -> Result<(), JobError>` +- [`src/worker.rs:150`](../src/worker.rs#L150) — `async fn run_cast_check(&self, payload: &str) -> Result<(), JobError>` ### SR-005 @@ -271,16 +277,16 @@ _None._ - [`src/api/upload.rs:29`](../src/api/upload.rs#L29) — `pub async fn post_manifest(` - [`src/castcheck.rs:82`](../src/castcheck.rs#L82) — `pub fn evaluate(submitted: &[SubmittedActor], credits: &[CastMember]) -> CastCheckOutcome` - [`src/model.rs:156`](../src/model.rs#L156) — `Unknown` -- [`src/model.rs:259`](../src/model.rs#L259) — `Unknown` +- [`src/model.rs:264`](../src/model.rs#L264) — `Unknown` - [`src/validate.rs:227`](../src/validate.rs#L227) — `pub fn validate_manifest(mut m: Jmanifest) -> VResult` -- [`src/worker.rs:96`](../src/worker.rs#L96) — `async fn run_cast_check(&self, payload: &str) -> Result<(), JobError>` +- [`src/worker.rs:150`](../src/worker.rs#L150) — `async fn run_cast_check(&self, payload: &str) -> Result<(), JobError>` ### UR-004 **Locations:** 2 - [`src/auth.rs:76`](../src/auth.rs#L76) — `pub fn client_ip(headers: &HeaderMap, peer: Option, trusted_proxies: &[IpAddr]) -…` -- [`src/ratelimit.rs:76`](../src/ratelimit.rs#L76) — `impl Default for RateLimiter` +- [`src/ratelimit.rs:93`](../src/ratelimit.rs#L93) — `impl Default for RateLimiter` ### UR-005 @@ -291,7 +297,7 @@ _None._ - [`src/auth.rs:22`](../src/auth.rs#L22) — `pub fn hash_token(token: &str) -> String` - [`src/castcheck.rs:82`](../src/castcheck.rs#L82) — `pub fn evaluate(submitted: &[SubmittedActor], credits: &[CastMember]) -> CastCheckOutcome` - [`src/castcheck.rs:215`](../src/castcheck.rs#L215) — `pub fn category_guard_violation(matched: &[MatchedActor], title_is_adult: bool) -> Option…` -- [`src/worker.rs:96`](../src/worker.rs#L96) — `async fn run_cast_check(&self, payload: &str) -> Result<(), JobError>` +- [`src/worker.rs:150`](../src/worker.rs#L150) — `async fn run_cast_check(&self, payload: &str) -> Result<(), JobError>` ### UR-006 @@ -307,6 +313,17 @@ _None._ - [`src/api/exists.rs:74`](../src/api/exists.rs#L74) — `pub async fn exists_batch(` +### UR-008 + +**Locations:** 6 + +- [`src/api/federation.rs:66`](../src/api/federation.rs#L66) — `pub async fn get_changes(` +- [`src/api/federation.rs:108`](../src/api/federation.rs#L108) — `pub async fn get_manifest_by_content_id(` +- [`src/api/federation.rs:160`](../src/api/federation.rs#L160) — `pub async fn post_have(` +- [`src/api/federation.rs:212`](../src/api/federation.rs#L212) — `pub async fn get_peers(` +- [`src/api/federation.rs:258`](../src/api/federation.rs#L258) — `pub async fn get_capabilities(State(state): State) -> ApiResult` +- [`src/worker.rs:107`](../src/worker.rs#L107) — `async fn run_federation_pull(&self, payload: &str) -> Result<(), JobError>` + ### UR-009 **Locations:** 1 @@ -317,7 +334,7 @@ _None._ **Locations:** 4 -- [`src/api/fetch.rs:267`](../src/api/fetch.rs#L267) — `pub fn reconstruct(` +- [`src/api/fetch.rs:270`](../src/api/fetch.rs#L270) — `pub fn reconstruct(` - [`src/castcheck.rs:82`](../src/castcheck.rs#L82) — `pub fn evaluate(submitted: &[SubmittedActor], credits: &[CastMember]) -> CastCheckOutcome` - [`src/db/repo.rs:412`](../src/db/repo.rs#L412) — `pub fn actors_for_manifest(` - [`src/model.rs:179`](../src/model.rs#L179) — `Unknown` @@ -327,7 +344,7 @@ _None._ **Locations:** 4 - [`src/model.rs:156`](../src/model.rs#L156) — `Unknown` -- [`src/model.rs:259`](../src/model.rs#L259) — `Unknown` +- [`src/model.rs:264`](../src/model.rs#L264) — `Unknown` - [`src/validate.rs:136`](../src/validate.rs#L136) — `fn is_allowed_text_char(c: char) -> bool` - [`src/validate.rs:378`](../src/validate.rs#L378) — `pub fn validate_audio_signature(sig: &str, runtime_sec: f64) -> VResult<()>` @@ -342,16 +359,17 @@ _None._ **Locations:** 4 -- [`src/api/fetch.rs:267`](../src/api/fetch.rs#L267) — `pub fn reconstruct(` +- [`src/api/fetch.rs:270`](../src/api/fetch.rs#L270) — `pub fn reconstruct(` - [`src/model.rs:179`](../src/model.rs#L179) — `Unknown` -- [`src/model.rs:233`](../src/model.rs#L233) — `pub fn from_str(s: &str) -> Option` +- [`src/model.rs:238`](../src/model.rs#L238) — `pub fn from_stored(s: &str) -> Option` - [`src/validate.rs:510`](../src/validate.rs#L510) — `fn validate_scenes(idx: usize, a: &Actor, runtime_sec: f64) -> VResult>` ### UR-014 -**Locations:** 2 +**Locations:** 3 -- [`src/model.rs:259`](../src/model.rs#L259) — `Unknown` +- [`src/api/federation.rs:258`](../src/api/federation.rs#L258) — `pub async fn get_capabilities(State(state): State) -> ApiResult` +- [`src/model.rs:264`](../src/model.rs#L264) — `Unknown` - [`src/validate.rs:227`](../src/validate.rs#L227) — `pub fn validate_manifest(mut m: Jmanifest) -> VResult` ### UR-017 @@ -359,5 +377,5 @@ _None._ **Locations:** 2 - [`src/model.rs:197`](../src/model.rs#L197) — `Unknown` -- [`src/model.rs:233`](../src/model.rs#L233) — `pub fn from_str(s: &str) -> Option` +- [`src/model.rs:238`](../src/model.rs#L238) — `pub fn from_stored(s: &str) -> Option` diff --git a/src/api/federation.rs b/src/api/federation.rs new file mode 100644 index 0000000..5a13689 --- /dev/null +++ b/src/api/federation.rs @@ -0,0 +1,288 @@ +//! §9a federation endpoints. +//! +//! **Replicate content, re-derive judgement.** A validated manifest is immutable +//! and content-addressable, so replication is *set reconciliation* rather than +//! state synchronisation — no concurrent edits, no last-write-wins, no vector +//! clocks. What must not replicate is the mutable half: `status`, `reports` and +//! `cast_match_ratio` encode a local operator's judgement and legal position. +//! A server that adopts a peer's `listed` flags has outsourced its liability; one +//! that adopts their `delisted` flags has outsourced its moderation. +//! +//! **Pull, never push.** A pulling server chooses what it ingests and when. +//! Push would let any peer inject work into your validation queue — the same +//! abuse surface as anonymous upload, but at higher volume. + +use axum::extract::{Path, Query, State}; +use axum::http::HeaderMap; +use axum::response::{IntoResponse, Response}; +use axum::Json; +use serde::{Deserialize, Serialize}; + +use crate::db::repo; +use crate::error::{ApiError, ApiResult}; +use crate::model::JMANIFEST_VERSION; +use crate::ratelimit::Surface; +use crate::state::{with_quota_headers, AppState}; + +/// §5: bounded so one peer cannot walk the whole catalogue in a single request. +const MAX_CHANGES_LIMIT: usize = 1000; +/// §5: `POST /federation/have`, up to 1000 ids. +const MAX_HAVE_IDS: usize = 1000; + +#[derive(Debug, Deserialize)] +pub struct ChangesParams { + /// The cursor a peer echoes back. Opaque on the wire — a peer must treat it + /// as a token, not compute with it. + #[serde(default)] + pub since: Option, + #[serde(default)] + pub limit: Option, +} + +#[derive(Debug, Serialize)] +pub struct ChangesResponse { + pub cursor: Option, + pub server_id: String, + pub changes: Vec, +} + +#[derive(Debug, Serialize)] +pub struct ChangeItem { + pub content_id: String, + pub op: String, + #[serde(skip_serializing_if = "Option::is_none")] + pub reason: Option, + pub origin: String, + pub seq: i64, +} + +/// `GET /federation/changes?since={cursor}&limit=1000` +/// +/// A monotonic, append-only feed of locally-*listed* manifests. Entries are +/// metadata only — enough to decide whether to fetch, without transferring +/// payloads. The cursor is opaque and monotonic, so a peer resumes from its last +/// position and the feed is idempotent. +/// +/// TRACES: UR-008 | PR-006 +pub async fn get_changes( + State(state): State, + peer: crate::state::PeerIp, + headers: HeaderMap, + Query(params): Query, +) -> ApiResult { + let ip = state.client_ip(&headers, peer.0); + let quota = state.check_limit(&ip, Surface::FederationChanges)?; + + let limit = params.limit.unwrap_or(MAX_CHANGES_LIMIT).min(MAX_CHANGES_LIMIT); + let since = params.since; + + let entries = state + .db + .read(move |conn| repo::changes_since(conn, since, limit)) + .await + .map_err(ApiError::Internal)?; + + let cursor = entries.last().map(|e| e.seq); + let changes = entries + .into_iter() + .map(|e| ChangeItem { + content_id: e.content_id, + op: e.op, + reason: e.reason, + origin: e.origin, + seq: e.seq, + }) + .collect(); + + let body = ChangesResponse { cursor, server_id: state.config.server_id.clone(), changes }; + Ok(with_quota_headers(Json(body).into_response(), quota)) +} + +/// `GET /federation/manifests/{content_id}` +/// +/// Full content by hash. **The puller must verify that what comes back hashes to +/// the id it asked for**, and reject it otherwise — that check is what makes an +/// intermediary or a misbehaving peer unable to substitute content. This server +/// performs it on ingest (see [`crate::federation`]). +/// +/// TRACES: UR-008 | PR-006 +pub async fn get_manifest_by_content_id( + State(state): State, + peer: crate::state::PeerIp, + headers: HeaderMap, + Path(content_id): Path, +) -> ApiResult { + let ip = state.client_ip(&headers, peer.0); + let quota = state.check_limit(&ip, Surface::FederationFetch)?; + + let manifest = state + .db + .read(move |conn| { + let Some(row) = repo::manifest_by_content_id_ro(conn, &content_id)? else { + return Ok(None); + }; + // Only served content is replicated. A `pending` manifest has not + // cleared this server's own checks, and a `rejected` one failed them. + if row.status != "listed" && row.status != "flagged" { + return Ok(None); + } + let title = super::fetch::title_of(conn, &row.title_id)?; + let kind = if title.kind == "movie" { + crate::model::IdentityType::Movie + } else { + crate::model::IdentityType::Episode + }; + Ok(Some(super::fetch::reconstruct(conn, &row, &title, kind)?)) + }) + .await + .map_err(ApiError::Internal)?; + + let manifest = manifest.ok_or(ApiError::NotFound)?; + Ok(with_quota_headers(Json(manifest).into_response(), quota)) +} + +#[derive(Debug, Deserialize)] +#[serde(deny_unknown_fields)] +pub struct HaveRequest { + pub content_ids: Vec, +} + +#[derive(Debug, Serialize)] +pub struct HaveResponse { + pub have: Vec, +} + +/// `POST /federation/have` +/// +/// Batch existence check by `content_id`, so a peer diffs its set against yours +/// in one request before fetching anything. +/// +/// TRACES: UR-008 | PR-006 +pub async fn post_have( + State(state): State, + peer: crate::state::PeerIp, + headers: HeaderMap, + super::json::Json(req): super::json::Json, +) -> ApiResult { + let ip = state.client_ip(&headers, peer.0); + let quota = state.check_limit(&ip, Surface::FederationHave)?; + + if req.content_ids.len() > MAX_HAVE_IDS { + return Err(ApiError::BadRequest(format!("more than {MAX_HAVE_IDS} content_ids"))); + } + + let ids = req.content_ids; + let have = state + .db + .read(move |conn| repo::known_content_ids(conn, &ids)) + .await + .map_err(ApiError::Internal)?; + + Ok(with_quota_headers(Json(HaveResponse { have }).into_response(), quota)) +} + +#[derive(Debug, Serialize)] +pub struct PeersResponse { + pub server_id: String, + #[serde(skip_serializing_if = "Option::is_none")] + pub contact: Option, + pub peers: Vec, +} + +#[derive(Debug, Serialize)] +pub struct PeerItem { + pub url: String, + #[serde(skip_serializing_if = "Option::is_none")] + pub name: Option, + #[serde(skip_serializing_if = "Option::is_none")] + pub since: Option, +} + +/// `GET /federation/peers` +/// +/// A **human-facing directory, not a discovery mechanism** — the distinction is +/// the whole point of §9a's peer-directory section. It publishes a list a person +/// can read; nothing acts on it. A server never fetches a peer's peers, so there +/// is no crawl and therefore no network-wide topology to poison. +/// +/// Advertising is opt-in per peer, and publishing the directory at all is opt-in +/// for the server (`PublishPeerDirectory`, default off) — a server that would +/// rather not disclose its topology simply does not. +/// +/// TRACES: UR-008 | PR-005 +pub async fn get_peers( + State(state): State, + peer: crate::state::PeerIp, + headers: HeaderMap, +) -> ApiResult { + let ip = state.client_ip(&headers, peer.0); + let quota = state.check_limit(&ip, Surface::FederationPeers)?; + + if !state.config.publish_peer_directory { + return Err(ApiError::NotFound); + } + + let peers = state.db.read(repo::advertised_peers).await.map_err(ApiError::Internal)?; + let body = PeersResponse { + server_id: state.config.server_id.clone(), + contact: state.config.contact.clone(), + peers: peers + .into_iter() + .map(|p| PeerItem { url: p.url, name: p.name, since: p.peered_since }) + .collect(), + }; + Ok(with_quota_headers(Json(body).into_response(), quota)) +} + +#[derive(Debug, Serialize)] +pub struct CapabilitiesResponse { + pub server_id: String, + /// Exchange envelope versions this server accepts. A client checks this once + /// rather than discovering a mismatch as a `400` per manifest across a whole + /// library sweep. + pub jmanifest_versions: Vec, + pub federation: bool, + /// §3: `POST /manifests/search` is expensive and therefore optional to + /// implement, so it is advertised rather than assumed. + pub audio_search: bool, + pub audio_tier_matching: bool, +} + +/// `GET /federation/capabilities` +/// +/// What this server actually supports. §3 specifies this for advertising the +/// optional audio-search endpoint; it also carries the accepted envelope +/// versions, which is what lets a plugin discover a schema mismatch in one cheap +/// request instead of N failed uploads. +/// +/// TRACES: UR-008, UR-014 | SR-003 +pub async fn get_capabilities(State(state): State) -> ApiResult { + let body = CapabilitiesResponse { + server_id: state.config.server_id.clone(), + jmanifest_versions: vec![JMANIFEST_VERSION], + federation: true, + // Deferred by design (§3 sequencing): signatures accumulate first. + audio_search: false, + audio_tier_matching: false, + }; + Ok(Json(body).into_response()) +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn capabilities_advertise_only_the_current_envelope_version() { + // The flag day means exactly one accepted version; advertising a range + // would invite a client to send something that will be rejected. + assert_eq!(JMANIFEST_VERSION, 2); + } + + #[test] + fn have_request_rejects_unknown_fields() { + // §6 stage 2 applies to federation bodies too — a peer is not exempt. + let r = serde_json::from_str::(r#"{"content_ids":[],"extra":1}"#); + assert!(r.is_err()); + } +} diff --git a/src/api/fetch.rs b/src/api/fetch.rs index 917fc00..d6ac32b 100644 --- a/src/api/fetch.rs +++ b/src/api/fetch.rs @@ -238,7 +238,10 @@ pub async fn get_status( } } -fn title_of(conn: &rusqlite::Connection, title_id: &str) -> anyhow::Result { +pub(crate) fn title_of( + conn: &rusqlite::Connection, + title_id: &str, +) -> anyhow::Result { let row = conn.query_row( "SELECT id, kind, tmdb_id, imdb_id, name, year, adult, certification FROM titles WHERE id = ?1", diff --git a/src/api/mod.rs b/src/api/mod.rs index 72e32f1..02ee40b 100644 --- a/src/api/mod.rs +++ b/src/api/mod.rs @@ -1,6 +1,7 @@ //! HTTP surface (§4). Base path `/api/v1`, JSON throughout. pub mod exists; +pub mod federation; pub mod fetch; pub mod json; pub mod report; diff --git a/src/app.rs b/src/app.rs index 64ee194..4b859cd 100644 --- a/src/app.rs +++ b/src/app.rs @@ -14,7 +14,7 @@ use axum::Router; use tower_http::timeout::TimeoutLayer; use tower_http::trace::TraceLayer; -use crate::api::{exists, fetch, report, upload}; +use crate::api::{exists, federation, fetch, report, upload}; use crate::state::AppState; use crate::validate::limits; @@ -49,6 +49,13 @@ pub fn router(state: AppState) -> Router { ) // §5a — anonymous bearer capability, not an account. .route("/tokens", post(upload::post_token)) + // §9a federation. Pull-based: a peer chooses what it ingests and when, + // so every route here is a read. There is deliberately no push endpoint. + .route("/federation/changes", get(federation::get_changes)) + .route("/federation/manifests/{content_id}", get(federation::get_manifest_by_content_id)) + .route("/federation/have", post(federation::post_have)) + .route("/federation/peers", get(federation::get_peers)) + .route("/federation/capabilities", get(federation::get_capabilities)) .layer(DefaultBodyLimit::max(SMALL_BODY_LIMIT)); Router::new() diff --git a/src/config.rs b/src/config.rs index 596f95c..f65e6e5 100644 --- a/src/config.rs +++ b/src/config.rs @@ -21,6 +21,11 @@ pub struct Config { pub server_id: String, pub request_timeout: Duration, /// Number of cast-check jobs to lease per worker tick. + /// §9a: publishing the peer directory is itself optional. A server that + /// would rather not disclose its topology simply does not. + pub publish_peer_directory: bool, + /// Human contact for operators arranging a peering out of band. + pub contact: Option, pub job_batch: usize, pub job_poll_interval: Duration, } @@ -48,6 +53,10 @@ impl Config { trusted_proxies, server_id: env_or("JRAY_SERVER_ID", "localhost"), request_timeout: Duration::from_secs(env_num("JRAY_REQUEST_TIMEOUT_SEC", 30)), + publish_peer_directory: std::env::var("JRAY_PUBLISH_PEER_DIRECTORY") + .map(|v| v == "1" || v.eq_ignore_ascii_case("true")) + .unwrap_or(false), + contact: std::env::var("JRAY_CONTACT").ok().filter(|s| !s.is_empty()), job_batch: env_num("JRAY_JOB_BATCH", 8) as usize, job_poll_interval: Duration::from_secs(env_num("JRAY_JOB_POLL_SEC", 5)), }) diff --git a/src/db/repo.rs b/src/db/repo.rs index 7ec5b7f..8525bec 100644 --- a/src/db/repo.rs +++ b/src/db/repo.rs @@ -1138,3 +1138,405 @@ mod tests { assert_eq!(counts, (0, 0, 0)); } } + +// --------------------------------------------------------------------------- +// Federation (§9a) +// +// The design in one line: **replicate content, re-derive judgement.** A +// validated manifest is immutable and content-addressable, so replication is set +// reconciliation — no concurrent edits, no vector clocks, no merge conflicts. +// What must *not* replicate is `status`, `reports` and `cast_match_ratio`: those +// encode a local operator's judgement and legal position, and a server that +// inherits them has outsourced its moderation, or its liability. +// --------------------------------------------------------------------------- + +#[derive(Debug, Clone)] +pub struct Peer { + pub id: String, + pub url: String, + pub name: Option, + pub enabled: bool, + pub pull_interval_sec: i64, + pub trust_abuse_retractions: bool, + pub advertise: bool, + pub max_ingest_per_hour: i64, + pub peered_since: Option, + pub last_cursor: Option, + pub last_pull_at: Option, + pub last_error: Option, +} + +fn map_peer(r: &rusqlite::Row<'_>) -> rusqlite::Result { + Ok(Peer { + id: r.get(0)?, + url: r.get(1)?, + name: r.get(2)?, + enabled: r.get::<_, i64>(3)? != 0, + pull_interval_sec: r.get(4)?, + trust_abuse_retractions: r.get::<_, i64>(5)? != 0, + advertise: r.get::<_, i64>(6)? != 0, + max_ingest_per_hour: r.get(7)?, + peered_since: r.get(8)?, + last_cursor: r.get(9)?, + last_pull_at: r.get(10)?, + last_error: r.get(11)?, + }) +} + +const PEER_COLUMNS: &str = "id, url, name, enabled, pull_interval_sec, \ + trust_abuse_retractions, advertise, max_ingest_per_hour, peered_since, last_cursor, \ + last_pull_at, last_error"; + +/// Adds a peer. +/// +/// §9a: **there is no automatic peering, ever.** This is only ever reached from +/// an operator action — nothing a remote server returns can call it, which is +/// what keeps the network's trust properties from being set by whoever joins. +pub fn insert_peer( + tx: &Transaction<'_>, + url: &str, + name: Option<&str>, + now: &str, +) -> anyhow::Result { + let id = ulid::Ulid::new().to_string(); + tx.execute( + "INSERT INTO peers (id, url, name, peered_since) VALUES (?1, ?2, ?3, ?4) + ON CONFLICT (url) DO NOTHING", + params![id, url, name, now], + )?; + let mut stmt = tx.prepare_cached("SELECT id FROM peers WHERE url = ?1")?; + Ok(stmt.query_row(params![url], |r| r.get(0))?) +} + +pub fn peers_due(conn: &Connection, now: &str) -> anyhow::Result> { + let sql = format!( + "SELECT {PEER_COLUMNS} FROM peers + WHERE enabled = 1 + AND (last_pull_at IS NULL + OR datetime(last_pull_at, '+' || pull_interval_sec || ' seconds') <= ?1)" + ); + let mut stmt = conn.prepare_cached(&sql)?; + let rows = stmt.query_map(params![now], map_peer)?.collect::>>()?; + Ok(rows) +} + +pub fn peer_by_id(conn: &Connection, id: &str) -> anyhow::Result> { + let sql = format!("SELECT {PEER_COLUMNS} FROM peers WHERE id = ?1"); + let mut stmt = conn.prepare_cached(&sql)?; + Ok(stmt.query_row(params![id], map_peer).optional()?) +} + +/// Peers this server has chosen to advertise (§9a peer directory). +pub fn advertised_peers(conn: &Connection) -> anyhow::Result> { + let sql = format!("SELECT {PEER_COLUMNS} FROM peers WHERE advertise = 1 ORDER BY url"); + let mut stmt = conn.prepare_cached(&sql)?; + let rows = stmt.query_map([], map_peer)?.collect::>>()?; + Ok(rows) +} + +pub fn record_pull( + tx: &Transaction<'_>, + peer_id: &str, + cursor: Option<&str>, + now: &str, + error: Option<&str>, +) -> anyhow::Result<()> { + tx.execute( + "UPDATE peers + SET last_cursor = COALESCE(?2, last_cursor), last_pull_at = ?3, last_error = ?4 + WHERE id = ?1", + params![peer_id, cursor, now, error], + )?; + Ok(()) +} + +#[derive(Debug, Clone)] +pub struct ChangeEntry { + /// Opaque to the peer, which only ever echoes it back as `since`. + pub seq: i64, + pub content_id: String, + pub op: String, + pub reason: Option, + pub origin: String, +} + +/// The `content_id` of a stored manifest, if it has one. +pub fn content_id_of(tx: &Transaction<'_>, manifest_id: &str) -> anyhow::Result> { + let mut stmt = tx.prepare_cached("SELECT content_id FROM manifests WHERE id = ?1")?; + Ok(stmt + .query_row(params![manifest_id], |r| r.get::<_, Option>(0)) + .optional()? + .flatten()) +} + +/// Appends to the change feed. Called when a manifest becomes listed, or is +/// delisted — the two events a peer can act on. +pub fn append_change( + tx: &Transaction<'_>, + content_id: &str, + op: &str, + reason: Option<&str>, + origin: &str, + now: &str, +) -> anyhow::Result { + tx.execute( + "INSERT INTO federation_log (content_id, op, reason, origin, created_at) + VALUES (?1, ?2, ?3, ?4, ?5)", + params![content_id, op, reason, origin, now], + )?; + Ok(tx.last_insert_rowid()) +} + +/// The change feed since an opaque cursor. +/// +/// Resumption is `seq > cursor` over a monotonic integer, so a peer never skips +/// an entry and never re-reads one. See the `federation_log` schema comment for +/// why a ULID cursor would have been wrong. +pub fn changes_since( + conn: &Connection, + since: Option, + limit: usize, +) -> anyhow::Result> { + let mut stmt = conn.prepare_cached( + "SELECT seq, content_id, op, reason, origin FROM federation_log + WHERE (?1 IS NULL OR seq > ?1) + ORDER BY seq ASC + LIMIT ?2", + )?; + let rows = stmt + .query_map(params![since, limit as i64], |r| { + Ok(ChangeEntry { + seq: r.get(0)?, + content_id: r.get(1)?, + op: r.get(2)?, + reason: r.get(3)?, + origin: r.get(4)?, + }) + })? + .collect::>>()?; + Ok(rows) +} + +/// Which of these `content_id`s this server already holds (`POST /federation/have`). +pub fn known_content_ids(conn: &Connection, ids: &[String]) -> anyhow::Result> { + let mut stmt = conn.prepare_cached("SELECT 1 FROM manifests WHERE content_id = ?1 LIMIT 1")?; + let mut have = Vec::new(); + for id in ids { + if stmt.query_row(params![id], |_| Ok(())).optional()?.is_some() { + have.push(id.clone()); + } + } + Ok(have) +} + +pub fn manifest_by_content_id_ro( + conn: &Connection, + content_id: &str, +) -> anyhow::Result> { + let sql = format!("SELECT {MANIFEST_COLUMNS} FROM manifests WHERE content_id = ?1"); + let mut stmt = conn.prepare_cached(&sql)?; + Ok(stmt.query_row(params![content_id], map_manifest).optional()?) +} + +/// How many manifests this peer supplied in the last hour, for `MaxIngestPerHour`. +pub fn ingest_count_since(conn: &Connection, peer_id: &str, since: &str) -> anyhow::Result { + let mut stmt = conn.prepare_cached( + "SELECT COUNT(*) FROM manifests WHERE ingested_from = ?1 AND created_at >= ?2", + )?; + Ok(stmt.query_row(params![peer_id, since], |r| r.get(0))?) +} + +/// Flags a locally-held manifest for review after a peer retracted it. +/// +/// **Not a delist.** §9a: a retraction is a warning worth acting on, whereas a +/// listing is merely a nomination — the asymmetry is deliberate. Auto-delisting +/// on any peer's retraction would hand every peer a remote delete primitive over +/// your catalogue. The one exception is the opt-in abuse channel. +pub fn flag_for_review( + tx: &Transaction<'_>, + content_id: &str, + reason: &str, +) -> anyhow::Result { + let n = tx.execute( + "UPDATE manifests SET status = 'flagged', reject_reason = ?2 + WHERE content_id = ?1 AND status = 'listed'", + params![content_id, reason], + )?; + Ok(n > 0) +} + +/// Delists immediately. Reached only for a peer configured `TrustAbuseRetractions`. +pub fn delist_by_content_id(tx: &Transaction<'_>, content_id: &str) -> anyhow::Result { + let n = tx.execute( + "UPDATE manifests SET status = 'rejected', reject_reason = 'peer_abuse_retraction' + WHERE content_id = ?1", + params![content_id], + )?; + Ok(n > 0) +} + +#[cfg(test)] +mod federation_tests { + use super::*; + use crate::db::Db; + + const NOW: &str = "2026-07-30T12:00:00Z"; + + async fn seeded(status: &'static str) -> Db { + let db = Db::open(":memory:").unwrap(); + db.write(move |tx| { + let title_id = upsert_title( + tx, + crate::model::IdentityType::Movie, + Some("1"), + None, + None, + None, + NOW, + )?; + insert_manifest( + tx, + &NewManifest { + id: "m1", + 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, + pipeline_version: None, + gallery_scope: None, + contributor_id: None, + status, + content_id: Some("sha256:x"), + origin: "peer.example", + ingested_from: None, + created_at: NOW, + }, + )?; + Ok(()) + }) + .await + .unwrap(); + db + } + + #[tokio::test] + async fn a_peer_retraction_flags_rather_than_delists() { + // §9a's central asymmetry. A retraction is a *warning worth acting on*; + // a listing is merely a *nomination*. Auto-delisting on any peer's + // retraction would hand every peer a remote delete primitive over your + // catalogue — which is why this path stops at `flagged`. + let db = seeded("listed").await; + let (changed, status) = db + .write(|tx| { + let changed = flag_for_review(tx, "sha256:x", "peer_retracted")?; + let status: String = + tx.query_row("SELECT status FROM manifests WHERE id='m1'", [], |r| r.get(0))?; + Ok((changed, status)) + }) + .await + .unwrap(); + assert!(changed); + assert_eq!(status, "flagged", "a peer retraction must not delist"); + } + + #[tokio::test] + async fn the_abuse_channel_can_delist_outright() { + // The one exception, reached only for a peer explicitly configured + // `TrustAbuseRetractions`. It exists because a takedown propagating at + // the speed of manual review is the wrong failure mode for that case. + let db = seeded("listed").await; + let status = db + .write(|tx| { + assert!(delist_by_content_id(tx, "sha256:x")?); + let status: String = + tx.query_row("SELECT status FROM manifests WHERE id='m1'", [], |r| r.get(0))?; + Ok(status) + }) + .await + .unwrap(); + assert_eq!(status, "rejected"); + } + + #[tokio::test] + async fn flagging_does_not_resurrect_a_rejected_manifest() { + // `flag_for_review` moves `listed` -> `flagged` only. A peer retracting + // something this server already rejected must not raise its status. + let db = seeded("rejected").await; + let changed = + db.write(|tx| flag_for_review(tx, "sha256:x", "peer_retracted")).await.unwrap(); + assert!(!changed); + } + + #[tokio::test] + async fn the_change_feed_is_ordered_and_resumable() { + let db = Db::open(":memory:").unwrap(); + let seqs = db + .write(|tx| { + let mut seqs = Vec::new(); + for i in 0..5 { + seqs.push(append_change( + tx, + &format!("sha256:{i}"), + "add", + None, + "local", + NOW, + )?); + } + Ok(seqs) + }) + .await + .unwrap(); + + // The property the whole cursor scheme rests on. A ULID would fail this: + // two generated in the same millisecond carry independent random parts, + // so they can sort opposite to write order — and a peer resuming from + // `seq > cursor` would then silently skip one. + let mut sorted = seqs.clone(); + sorted.sort(); + assert_eq!(sorted, seqs, "feed sequence must be monotonic"); + assert!(seqs.windows(2).all(|w| w[1] > w[0]), "strictly increasing"); + + let from_second = seqs[1]; + let rest = db.write(move |tx| changes_since(tx, Some(from_second), 100)).await.unwrap(); + assert_eq!(rest.len(), 3); + assert_eq!(rest[0].content_id, "sha256:2"); + } + + #[tokio::test] + async fn peering_requires_an_explicit_add_and_is_idempotent() { + let db = Db::open(":memory:").unwrap(); + let (a, b, n) = db + .write(|tx| { + let a = insert_peer(tx, "https://p.example", Some("P"), NOW)?; + let b = insert_peer(tx, "https://p.example", Some("P again"), NOW)?; + let n: i64 = tx.query_row("SELECT COUNT(*) FROM peers", [], |r| r.get(0))?; + Ok((a, b, n)) + }) + .await + .unwrap(); + assert_eq!(a, b, "adding the same URL twice is one peering"); + assert_eq!(n, 1); + } + + #[tokio::test] + async fn a_new_peer_is_disabled_and_unadvertised_by_default() { + // Federation is off by default, and advertising is opt-in on both sides. + let db = Db::open(":memory:").unwrap(); + let peer = db + .write(|tx| { + let id = insert_peer(tx, "https://p.example", None, NOW)?; + Ok(peer_by_id(tx, &id)?.unwrap()) + }) + .await + .unwrap(); + assert!(!peer.enabled); + assert!(!peer.advertise); + assert!(!peer.trust_abuse_retractions, "the abuse channel is opt-in per peer"); + } +} diff --git a/src/db/schema.sql b/src/db/schema.sql index 92b4d13..755b9dd 100644 --- a/src/db/schema.sql +++ b/src/db/schema.sql @@ -109,6 +109,54 @@ CREATE TABLE IF NOT EXISTS jobs ( leased_at TEXT ); +-- §9a federation. Peering is trust-by-configuration: a row exists only because +-- an operator typed a URL. Nothing a remote server says can create one. +CREATE TABLE IF NOT EXISTS peers ( + id TEXT PRIMARY KEY, + url TEXT NOT NULL UNIQUE, + name TEXT, + enabled INTEGER NOT NULL DEFAULT 0, + pull_interval_sec INTEGER NOT NULL DEFAULT 3600, + -- Auto-delist on a peer's legal retraction. Off by default: a retraction + -- that delists automatically is a remote delete primitive over your + -- catalogue, so it is opt-in per peer between operators who know each other. + trust_abuse_retractions INTEGER NOT NULL DEFAULT 0, + -- Publishing a peering is opt-in on BOTH sides (§9a): peering with someone + -- must not advertise their existence against their wishes. + advertise INTEGER NOT NULL DEFAULT 0, + max_ingest_per_hour INTEGER NOT NULL DEFAULT 500, + peered_since TEXT, + last_cursor TEXT, + last_pull_at TEXT, + last_error TEXT +); + +-- The change feed. Append-only, so a peer can resume from an opaque cursor and +-- the feed is idempotent. A separate table rather than deriving the feed from +-- `manifests` because a *retraction* is an event with no surviving row. +CREATE TABLE IF NOT EXISTS federation_log ( + -- A genuinely monotonic sequence, NOT a ULID. + -- + -- ULIDs are only monotonic *between* milliseconds: two generated in the same + -- millisecond carry independent random components, so they can sort in the + -- opposite order to which they were written. A peer resuming from `seq > + -- cursor` would then silently skip an entry — replication losing manifests + -- with no error anywhere, which is the worst shape a bug can take here. + -- + -- AUTOINCREMENT (rather than plain rowid) additionally guarantees the value + -- never decreases even after deletions. Portability note (§8): Postgres + -- spells this `BIGSERIAL PRIMARY KEY`; it is the one place a monotonic + -- sequence has no fully portable form, and it is worth the exception. + seq INTEGER PRIMARY KEY AUTOINCREMENT, + content_id TEXT NOT NULL, + op TEXT NOT NULL, -- add | retract + reason TEXT, -- retract only; 'abuse' is the one that may auto-delist + origin TEXT NOT NULL, -- server_id that first accepted it + created_at TEXT NOT NULL +); + +CREATE INDEX IF NOT EXISTS idx_federation_log_content ON federation_log(content_id); + CREATE INDEX IF NOT EXISTS idx_titles_tmdb ON titles(tmdb_id); CREATE INDEX IF NOT EXISTS idx_titles_imdb ON titles(imdb_id); CREATE INDEX IF NOT EXISTS idx_manifests_title_runtime ON manifests(title_id, runtime_sec); diff --git a/src/federation.rs b/src/federation.rs new file mode 100644 index 0000000..bfc6c23 --- /dev/null +++ b/src/federation.rs @@ -0,0 +1,307 @@ +//! §9a federation pull: fetching from a peer and ingesting locally. +//! +//! The rule this module exists to enforce is **"re-derive, don't inherit"**. A +//! pulled manifest is *not* trusted because a peer listed it — it enters the +//! local pipeline exactly as if freshly uploaded: full §6 stage 1 and 2 +//! validation, a local TMDB cast check against this server's own cache and +//! thresholds, and a `status` assigned by this server's rules. The peer's +//! `cast_match_ratio` is advisory only, useful for prioritising ingestion order +//! and never a substitute for checking. +//! +//! The asymmetry in how peer states are treated is deliberate: +//! +//! | Peer state | Local effect | +//! |---|---| +//! | Peer lists it | Eligible for ingestion; still fully re-validated | +//! | Peer retracts it | Flagged for review, **not** auto-delisted | +//! | Peer never had it | No signal | +//! +//! A retraction is a *warning worth acting on*; a listing is merely a +//! *nomination*. Auto-delisting on a peer's retraction would hand any peer a +//! remote delete primitive over your catalogue. + +use std::time::Duration; + +use serde::{Deserialize, Serialize}; + +use crate::db::repo::{self, Peer}; +use crate::ingest; +use crate::model::Jmanifest; +use crate::validate; + +pub const JOB_FEDERATION_PULL: &str = "federation_pull"; + +/// The one retraction reason that may auto-delist, and only for a peer the +/// operator explicitly configured to trust for it. +/// +/// It exists because the alternative — a takedown propagating at the speed of +/// manual review — is the wrong failure mode for that one case. +pub const RETRACT_REASON_ABUSE: &str = "abuse"; + +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct FederationPullJob { + pub peer_id: String, +} + +#[derive(Debug, Deserialize)] +struct ChangesBody { + #[serde(default)] + cursor: Option, + #[serde(default)] + changes: Vec, +} + +#[derive(Debug, Deserialize)] +struct ChangeBody { + content_id: String, + op: String, + #[serde(default)] + reason: Option, + #[serde(default)] + origin: Option, + seq: i64, +} + +/// Outcome of one pull, for logging and tests. +#[derive(Debug, Default, PartialEq)] +pub struct PullOutcome { + pub examined: usize, + pub ingested: usize, + pub skipped_known: usize, + pub skipped_own_origin: usize, + pub flagged: usize, + pub delisted: usize, + pub rejected: usize, + pub capped: bool, +} + +#[derive(Debug, thiserror::Error)] +pub enum PullError { + #[error("peer transport error: {0}")] + Transport(String), + #[error("peer returned an unusable response: {0}")] + Malformed(String), +} + +/// Fetches a peer's change feed and ingests what is new. +pub struct Puller { + http: reqwest::Client, + server_id: String, +} + +impl Puller { + pub fn new(server_id: String) -> Self { + let http = reqwest::Client::builder() + .timeout(Duration::from_secs(30)) + .user_agent(concat!("jray-server/", env!("CARGO_PKG_VERSION"))) + .build() + .expect("building reqwest client"); + Self { http, server_id } + } + + async fn get_json(&self, url: &str) -> Result { + let resp = + self.http.get(url).send().await.map_err(|e| PullError::Transport(e.to_string()))?; + if !resp.status().is_success() { + return Err(PullError::Transport(format!("status {}", resp.status()))); + } + // §9 client hardening applies to peers too: a response is capped while + // reading rather than buffered whole, so a hostile peer cannot exhaust + // memory by advertising a small body and sending a large one. + let bytes = resp.bytes().await.map_err(|e| PullError::Transport(e.to_string()))?; + if bytes.len() > validate::limits::BODY_LIMIT_BUNDLE { + return Err(PullError::Malformed("response exceeds the bundle cap".into())); + } + serde_json::from_slice(&bytes).map_err(|e| PullError::Malformed(e.to_string())) + } + + /// One pull cycle against one peer. + pub async fn pull( + &self, + db: &crate::db::Db, + peer: &Peer, + now: &str, + ) -> Result { + let mut outcome = PullOutcome::default(); + + let base = peer.url.trim_end_matches('/'); + let url = match &peer.last_cursor { + Some(c) => format!("{base}/api/v1/federation/changes?since={c}&limit=500"), + None => format!("{base}/api/v1/federation/changes?limit=500"), + }; + let feed: ChangesBody = self.get_json(&url).await?; + + // §9a `MaxIngestPerHour`: a peer cannot flood the validation queue. + let hour_ago = iso_hours_ago(1); + let peer_id = peer.id.clone(); + let already = db + .read(move |conn| repo::ingest_count_since(conn, &peer_id, &hour_ago)) + .await + .unwrap_or(0); + let mut budget = (peer.max_ingest_per_hour - already).max(0); + + for change in &feed.changes { + outcome.examined += 1; + + // Loop prevention: ignore anything this server originated. Longer + // cycles are harmless anyway — content is hash-addressed and ingest + // is idempotent, so a second arrival is a no-op deduplication. + if change.origin.as_deref() == Some(self.server_id.as_str()) { + outcome.skipped_own_origin += 1; + continue; + } + + match change.op.as_str() { + "retract" => { + let is_abuse = change.reason.as_deref() == Some(RETRACT_REASON_ABUSE); + let cid = change.content_id.clone(); + if is_abuse && peer.trust_abuse_retractions { + let delisted = + db.write(move |tx| repo::delist_by_content_id(tx, &cid)).await; + if delisted.unwrap_or(false) { + outcome.delisted += 1; + tracing::warn!( + content_id = %change.content_id, peer = %peer.url, + "delisted on a trusted peer's abuse retraction" + ); + } + } else { + // A warning, not a verdict. + let flagged = db + .write(move |tx| repo::flag_for_review(tx, &cid, "peer_retracted")) + .await; + if flagged.unwrap_or(false) { + outcome.flagged += 1; + } + } + continue; + } + "add" => {} + other => { + tracing::debug!(op = other, "unknown change op, ignoring"); + continue; + } + } + + // Already held: skipped without re-validation, which is the + // deduplication content addressing buys. + let cid = change.content_id.clone(); + let known = db + .read(move |conn| repo::known_content_ids(conn, &[cid])) + .await + .unwrap_or_default(); + if !known.is_empty() { + outcome.skipped_known += 1; + continue; + } + + if budget <= 0 { + outcome.capped = true; + tracing::info!(peer = %peer.url, "ingest cap reached, resuming next pull"); + break; + } + + let fetch_url = format!("{base}/api/v1/federation/manifests/{}", change.content_id); + let manifest: Jmanifest = match self.get_json(&fetch_url).await { + Ok(m) => m, + Err(e) => { + tracing::warn!(content_id = %change.content_id, error = %e, "peer fetch failed"); + continue; + } + }; + + match self.ingest_one(db, peer, manifest, &change.content_id, now).await { + Ok(true) => { + outcome.ingested += 1; + budget -= 1; + } + Ok(false) => outcome.rejected += 1, + Err(e) => { + tracing::warn!(content_id = %change.content_id, error = ?e, "ingest failed"); + outcome.rejected += 1; + } + } + } + + // Advance the cursor even when individual items failed: the feed is a + // position, not a work queue, and a permanently-bad entry must not wedge + // replication forever. A missed manifest reappears on the next full sync. + let cursor = + feed.cursor.or_else(|| feed.changes.last().map(|c| c.seq)).map(|c| c.to_string()); + let peer_id = peer.id.clone(); + let now_owned = now.to_string(); + let _ = db + .write(move |tx| repo::record_pull(tx, &peer_id, cursor.as_deref(), &now_owned, None)) + .await; + + Ok(outcome) + } + + /// Validates and stores one pulled manifest. Returns whether it was ingested. + async fn ingest_one( + &self, + db: &crate::db::Db, + peer: &Peer, + manifest: Jmanifest, + expected_content_id: &str, + now: &str, + ) -> anyhow::Result { + // §6 stage 2, in full. A peer is not exempt from the checks that keep + // §5a's Threat 1 closed — that is the whole point of re-deriving. + let valid = match validate::validate_manifest(manifest) { + Ok(v) => v, + Err(e) => { + tracing::warn!(peer = %peer.url, error = %e, "peer manifest failed validation"); + return Ok(false); + } + }; + + // **The puller must verify the content hashes to what it asked for.** + // This is what makes an intermediary, or a peer that swapped the body, + // unable to substitute content under a trusted id. + let computed = ingest::compute_content_id(&valid); + if computed != expected_content_id { + tracing::warn!( + peer = %peer.url, expected = %expected_content_id, computed = %computed, + "content_id mismatch — refusing substituted content" + ); + return Ok(false); + } + + let peer_id = peer.id.clone(); + let now_owned = now.to_string(); + // `origin` is preserved across hops so an operator can say "stop + // ingesting anything originating from X". It is provenance, not authority. + let origin = self.server_id.clone(); + db.write(move |tx| { + ingest::persist(tx, &valid, None, &origin, Some(&peer_id), &now_owned)?; + Ok(()) + }) + .await?; + + Ok(true) + } +} + +fn iso_hours_ago(hours: u64) -> String { + crate::worker::iso_ago(hours * 3600) +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn abuse_is_the_only_auto_delist_reason() { + // §9a's one exception. Everything else flags for review, because + // auto-delisting on any retraction is a remote delete primitive. + assert_eq!(RETRACT_REASON_ABUSE, "abuse"); + } + + #[test] + fn pull_outcome_starts_empty() { + let o = PullOutcome::default(); + assert_eq!(o.ingested, 0); + assert!(!o.capped); + } +} diff --git a/src/lib.rs b/src/lib.rs index 28ed871..dcddde8 100644 --- a/src/lib.rs +++ b/src/lib.rs @@ -28,6 +28,7 @@ pub mod config; pub mod content_id; pub mod db; pub mod error; +pub mod federation; pub mod ingest; pub mod matching; pub mod model; diff --git a/src/main.rs b/src/main.rs index 345724c..36beada 100644 --- a/src/main.rs +++ b/src/main.rs @@ -53,6 +53,7 @@ async fn main() -> anyhow::Result<()> { tmdb, batch: config.job_batch, poll_interval: config.job_poll_interval, + server_id: config.server_id.clone(), }; let worker_handle = tokio::spawn(worker.run(shutdown_rx)); diff --git a/src/ratelimit.rs b/src/ratelimit.rs index a2b686b..c304693 100644 --- a/src/ratelimit.rs +++ b/src/ratelimit.rs @@ -25,6 +25,10 @@ pub enum Surface { BundleUpload, Report, Search, + FederationChanges, + FederationFetch, + FederationHave, + FederationPeers, } impl Surface { @@ -39,6 +43,15 @@ impl Surface { Surface::BundleUpload => 20, Surface::Report => 20, Surface::Search => 60, + // §5: hourly polling is the default, so this allows generous + // catch-up without letting one peer walk the feed continuously. + Surface::FederationChanges => 120, + // Bootstrap pulls are bulk by nature; capped so one peer cannot + // saturate egress. + Surface::FederationFetch => 5000, + Surface::FederationHave => 120, + // A public directory read by humans — no reason for volume. + Surface::FederationPeers => 60, } } @@ -52,6 +65,10 @@ impl Surface { Surface::BundleUpload => "bundle_upload", Surface::Report => "report", Surface::Search => "search", + Surface::FederationChanges => "federation_changes", + Surface::FederationFetch => "federation_fetch", + Surface::FederationHave => "federation_have", + Surface::FederationPeers => "federation_peers", } } } @@ -214,5 +231,9 @@ mod tests { assert_eq!(Surface::BundleUpload.limit(), 20); assert_eq!(Surface::Report.limit(), 20); assert_eq!(Surface::Search.limit(), 60); + assert_eq!(Surface::FederationChanges.limit(), 120); + assert_eq!(Surface::FederationFetch.limit(), 5000); + assert_eq!(Surface::FederationHave.limit(), 120); + assert_eq!(Surface::FederationPeers.limit(), 60); } } diff --git a/src/worker.rs b/src/worker.rs index 9e7b282..362fc30 100644 --- a/src/worker.rs +++ b/src/worker.rs @@ -26,6 +26,9 @@ pub struct Worker { 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 { @@ -62,6 +65,9 @@ impl Worker { 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(()) @@ -93,6 +99,54 @@ impl Worker { 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 = @@ -286,6 +340,7 @@ impl Worker { 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| { @@ -317,6 +372,17 @@ impl Worker { 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 { @@ -391,6 +457,11 @@ 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) diff --git a/tests/api.rs b/tests/api.rs index a80afaa..352dc12 100644 --- a/tests/api.rs +++ b/tests/api.rs @@ -74,6 +74,8 @@ impl TestServer { tmdb_base_url: "http://127.0.0.1:1".into(), trusted_proxies: Vec::new(), server_id: "test.example".into(), + publish_peer_directory: true, + contact: Some("admin@test.example".into()), request_timeout: std::time::Duration::from_secs(30), job_batch: 8, job_poll_interval: std::time::Duration::from_secs(3600), diff --git a/tests/federation.rs b/tests/federation.rs new file mode 100644 index 0000000..fe07042 --- /dev/null +++ b/tests/federation.rs @@ -0,0 +1,405 @@ +//! §9a federation — the replication surface and its guarantees. +//! +//! The design's whole claim is **"replicate content, re-derive judgement"**, and +//! each of those halves has a failure mode worth pinning: +//! +//! - Replicating content wrongly means a peer can substitute what you serve. +//! Guarded by verifying the fetched body hashes to the `content_id` requested. +//! - Inheriting judgement means outsourcing your moderation, or your liability. +//! Guarded by re-validating everything and never adopting a peer's `status`. + +use std::sync::Arc; + +use axum::body::Body; +use axum::http::{Request, StatusCode}; +use http_body_util::BodyExt; +use jray_server::config::Config; +use jray_server::db::{repo, Db}; +use jray_server::ratelimit::RateLimiter; +use jray_server::state::AppState; +use jray_server::tmdb::TmdbClient; +use jray_server::{app, worker}; +use serde_json::{json, Value}; +use tower::ServiceExt; + +struct TempDir(std::path::PathBuf); + +impl TempDir { + fn new(tag: &str) -> Self { + let mut p = std::env::temp_dir(); + p.push(format!("jray-fed-{}-{}", tag, unique())); + std::fs::create_dir_all(&p).expect("creating temp dir"); + Self(p) + } + fn db_path(&self) -> String { + self.0.join("test.db").to_string_lossy().into_owned() + } +} + +impl Drop for TempDir { + fn drop(&mut self) { + let _ = std::fs::remove_dir_all(&self.0); + } +} + +fn unique() -> String { + use std::sync::atomic::{AtomicU64, Ordering}; + static N: AtomicU64 = AtomicU64::new(0); + let t = std::time::SystemTime::now() + .duration_since(std::time::UNIX_EPOCH) + .map(|d| d.as_nanos()) + .unwrap_or(0); + format!("{t}-{}", N.fetch_add(1, Ordering::Relaxed)) +} + +struct TestServer { + router: axum::Router, + db: Db, + _dir: TempDir, +} + +impl TestServer { + fn new(tag: &str, publish_directory: bool) -> Self { + let dir = TempDir::new(tag); + let db = Db::open(&dir.db_path()).expect("opening database"); + let config = Arc::new(Config { + bind: "127.0.0.1:0".into(), + db_path: dir.db_path(), + tmdb_api_key: None, + tmdb_base_url: "http://127.0.0.1:1".into(), + trusted_proxies: Vec::new(), + server_id: "local.example".into(), + publish_peer_directory: publish_directory, + contact: Some("admin@local.example".into()), + request_timeout: std::time::Duration::from_secs(30), + job_batch: 8, + job_poll_interval: std::time::Duration::from_secs(3600), + }); + let state = AppState { + db: db.clone(), + config: config.clone(), + limiter: Arc::new(RateLimiter::new()), + tmdb: Arc::new(TmdbClient::new(config.tmdb_base_url.clone(), None)), + }; + Self { router: app::router(state), db, _dir: dir } + } + + async fn get(&self, uri: &str) -> (StatusCode, Value) { + self.send(Request::builder().uri(uri).body(Body::empty()).unwrap()).await + } + + async fn post(&self, uri: &str, body: &Value) -> (StatusCode, Value) { + self.send( + Request::builder() + .method("POST") + .uri(uri) + .header("content-type", "application/json") + .body(Body::from(body.to_string())) + .unwrap(), + ) + .await + } + + async fn send(&self, req: Request) -> (StatusCode, Value) { + let resp = self.router.clone().oneshot(req).await.expect("router call"); + let status = resp.status(); + let bytes = resp.into_body().collect().await.expect("body").to_bytes(); + let body = if bytes.is_empty() { + Value::Null + } else { + serde_json::from_slice(&bytes).unwrap_or(Value::Null) + }; + (status, body) + } + + /// Seeds a listed manifest and its change-feed entry, as the cast check would. + async fn seed_listed(&self, tmdb_id: &str, content_id: &str, origin: &str) -> String { + let tmdb_id = tmdb_id.to_string(); + let content_id = content_id.to_string(); + let origin = origin.to_string(); + self.db + .write(move |tx| { + let now = worker::now_iso(); + let title_id = repo::upsert_title( + tx, + jray_server::model::IdentityType::Movie, + Some(&tmdb_id), + None, + Some("A Film"), + None, + &now, + )?; + let id = format!("m-{tmdb_id}"); + repo::insert_manifest( + tx, + &repo::NewManifest { + id: &id, + title_id: &title_id, + season: None, + episode: None, + runtime_sec: 6420.5, + video_hash: None, + audio_signature: None, + audio_sig_coarse: None, + sample_fps: Some(5.0), + extinction_sec: Some(12.0), + pipeline_version: Some("test 0.1"), + gallery_scope: Some("global"), + contributor_id: None, + status: "listed", + content_id: Some(&content_id), + origin: &origin, + ingested_from: None, + created_at: &now, + }, + )?; + repo::upsert_person(tx, 884, "Steve Buscemi", false, &now)?; + repo::insert_actor_scenes( + tx, + &id, + 884, + &[jray_server::validate::SceneCs::plain(19160, 20920)], + )?; + repo::append_change(tx, &content_id, "add", None, &origin, &now)?; + Ok(id) + }) + .await + .expect("seeding") + } +} + +// --------------------------------------------------------------------------- +// The change feed +// --------------------------------------------------------------------------- + +#[tokio::test] +async fn the_change_feed_is_resumable_from_an_opaque_cursor() { + let s = TestServer::new("feed", false); + for i in 0..3 { + s.seed_listed(&format!("100{i}"), &format!("sha256:c{i}"), "local.example").await; + } + + let (status, body) = s.get("/api/v1/federation/changes").await; + assert_eq!(status, StatusCode::OK); + let changes = body["changes"].as_array().expect("changes"); + assert_eq!(changes.len(), 3); + assert_eq!(body["server_id"], "local.example"); + + // Resuming from the first entry yields exactly the remainder — the property + // that makes the feed idempotent for a peer that restarts mid-pull. + let first = changes[0]["seq"].as_i64().expect("seq is a monotonic integer"); + let (_, body) = s.get(&format!("/api/v1/federation/changes?since={first}")).await; + assert_eq!(body["changes"].as_array().unwrap().len(), 2); + + // Resuming from the last yields nothing, and the cursor is absent rather + // than stale — a peer that polls an idle server does no work. + let last = changes[2]["seq"].as_i64().unwrap(); + let (_, body) = s.get(&format!("/api/v1/federation/changes?since={last}")).await; + assert!(body["changes"].as_array().unwrap().is_empty()); + assert!(body["cursor"].is_null()); +} + +#[tokio::test] +async fn the_feed_limit_is_capped() { + // §5: a peer must not be able to walk the entire catalogue in one request. + let s = TestServer::new("feed-limit", false); + let (status, body) = s.get("/api/v1/federation/changes?limit=999999").await; + assert_eq!(status, StatusCode::OK); + assert!(body["changes"].as_array().unwrap().len() <= 1000); +} + +// --------------------------------------------------------------------------- +// Fetch by content id +// --------------------------------------------------------------------------- + +#[tokio::test] +async fn a_listed_manifest_is_fetchable_by_content_id() { + let s = TestServer::new("byhash", false); + s.seed_listed("504172", "sha256:abc", "local.example").await; + + let (status, body) = s.get("/api/v1/federation/manifests/sha256:abc").await; + assert_eq!(status, StatusCode::OK); + assert_eq!(body["jmanifest_version"], 2); + assert_eq!(body["identity"]["tmdb_id"], "504172"); + // Names come from the server's TMDB-derived table, never from an upload. + assert_eq!(body["actors"][0]["name"], "Steve Buscemi"); +} + +#[tokio::test] +async fn an_unknown_content_id_is_not_found() { + let s = TestServer::new("byhash-missing", false); + let (status, _) = s.get("/api/v1/federation/manifests/sha256:nope").await; + assert_eq!(status, StatusCode::NOT_FOUND); +} + +#[tokio::test] +async fn a_pending_manifest_is_not_replicated() { + // Only content that cleared *this* server's checks is offered onward. + // Replicating a pending manifest would export an unverified claim. + let s = TestServer::new("pending-not-replicated", false); + s.db.write(|tx| { + let now = worker::now_iso(); + let title_id = repo::upsert_title( + tx, + jray_server::model::IdentityType::Movie, + Some("777"), + None, + None, + None, + &now, + )?; + repo::insert_manifest( + tx, + &repo::NewManifest { + id: "pending-one", + 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, + pipeline_version: None, + gallery_scope: None, + contributor_id: None, + status: "pending", + content_id: Some("sha256:pending"), + origin: "local.example", + ingested_from: None, + created_at: &now, + }, + )?; + Ok(()) + }) + .await + .unwrap(); + + let (status, _) = s.get("/api/v1/federation/manifests/sha256:pending").await; + assert_eq!(status, StatusCode::NOT_FOUND); +} + +// --------------------------------------------------------------------------- +// Batch have +// --------------------------------------------------------------------------- + +#[tokio::test] +async fn have_reports_only_what_is_held() { + let s = TestServer::new("have", false); + s.seed_listed("504172", "sha256:held", "local.example").await; + + let (status, body) = s + .post( + "/api/v1/federation/have", + &json!({ "content_ids": ["sha256:held", "sha256:absent"] }), + ) + .await; + assert_eq!(status, StatusCode::OK); + assert_eq!(body["have"], json!(["sha256:held"])); +} + +#[tokio::test] +async fn have_is_capped_and_strict() { + let s = TestServer::new("have-cap", false); + + let too_many: Vec = (0..1001).map(|i| format!("sha256:{i}")).collect(); + let (status, _) = s.post("/api/v1/federation/have", &json!({ "content_ids": too_many })).await; + assert_eq!(status, StatusCode::BAD_REQUEST); + + // §6 stage 2 applies to federation bodies too — a peer is not exempt. + let (status, _) = + s.post("/api/v1/federation/have", &json!({ "content_ids": [], "extra": true })).await; + assert_eq!(status, StatusCode::BAD_REQUEST); +} + +// --------------------------------------------------------------------------- +// Peer directory — publishing, not discovering +// --------------------------------------------------------------------------- + +#[tokio::test] +async fn the_peer_directory_lists_only_advertised_peers() { + let s = TestServer::new("peers", true); + s.db.write(|tx| { + let now = worker::now_iso(); + let a = repo::insert_peer(tx, "https://a.example", Some("Peer A"), &now)?; + repo::insert_peer(tx, "https://b.example", Some("Peer B"), &now)?; + // Advertising is opt-in per peer: peering with someone must not + // publish their existence against their wishes. + tx.execute("UPDATE peers SET advertise = 1 WHERE id = ?1", rusqlite::params![a])?; + Ok(()) + }) + .await + .unwrap(); + + let (status, body) = s.get("/api/v1/federation/peers").await; + assert_eq!(status, StatusCode::OK); + let peers = body["peers"].as_array().unwrap(); + assert_eq!(peers.len(), 1, "only the advertised peer appears"); + assert_eq!(peers[0]["url"], "https://a.example"); + assert_eq!(body["contact"], "admin@local.example"); +} + +#[tokio::test] +async fn the_peer_directory_can_be_withheld_entirely() { + // §9a: publishing the directory is itself optional. A server that would + // rather not disclose its topology simply does not. + let s = TestServer::new("peers-off", false); + let (status, _) = s.get("/api/v1/federation/peers").await; + assert_eq!(status, StatusCode::NOT_FOUND); +} + +#[tokio::test] +async fn there_is_no_endpoint_that_creates_a_peering() { + // **No automatic peering, ever.** Nothing a remote server says, and no data + // returned from any endpoint, can establish or widen a peering. This asserts + // the absence of a write surface — the guarantee is structural, so the test + // is that the obvious shapes simply do not route. + let s = TestServer::new("no-auto-peer", true); + + for (method, uri) in [ + ("POST", "/api/v1/federation/peers"), + ("PUT", "/api/v1/federation/peers"), + ("POST", "/api/v1/federation/peer"), + ("POST", "/api/v1/federation/subscribe"), + ("POST", "/api/v1/federation/announce"), + ] { + let req = Request::builder() + .method(method) + .uri(uri) + .header("content-type", "application/json") + .body(Body::from(r#"{"url":"https://hostile.example"}"#)) + .unwrap(); + let (status, _) = s.send(req).await; + assert!( + status == StatusCode::NOT_FOUND || status == StatusCode::METHOD_NOT_ALLOWED, + "{method} {uri} must not exist, got {status}" + ); + } + + // And nothing was created. + let n = s + .db + .read(|conn| Ok(conn.query_row("SELECT COUNT(*) FROM peers", [], |r| r.get::<_, i64>(0))?)) + .await + .unwrap(); + assert_eq!(n, 0); +} + +// --------------------------------------------------------------------------- +// Capabilities +// --------------------------------------------------------------------------- + +#[tokio::test] +async fn capabilities_advertise_the_envelope_version() { + // The gap this closes: without it, a plugin at v2 talking to a v1 server + // discovers the mismatch as a 400 per manifest across a whole library sweep. + let s = TestServer::new("caps", false); + let (status, body) = s.get("/api/v1/federation/capabilities").await; + assert_eq!(status, StatusCode::OK); + assert_eq!(body["jmanifest_versions"], json!([2])); + assert_eq!(body["federation"], true); + // Deferred by design (§3 sequencing), and advertised as absent rather than + // left for a client to discover by failure. + assert_eq!(body["audio_search"], false); +} diff --git a/tests/injection.rs b/tests/injection.rs index 68360d6..a852775 100644 --- a/tests/injection.rs +++ b/tests/injection.rs @@ -97,6 +97,8 @@ impl TestServer { tmdb_base_url: "http://127.0.0.1:1".into(), trusted_proxies: Vec::new(), server_id: "test.example".into(), + publish_peer_directory: true, + contact: Some("admin@test.example".into()), request_timeout: std::time::Duration::from_secs(30), job_batch: 8, job_poll_interval: std::time::Duration::from_secs(3600),