Files
dtourolleandClaude Opus 5 0ff1018bcc
CI / static musl binary (push) Has been skipped
CI / fmt, clippy, test (push) Failing after 2m0s
CI / advisories and licences (push) Successful in 27s
Withdraw the file-hash tier; document the legal posture
Removes `cut.video_hash` and the `exact` match tier on legal grounds. The
OpenSubtitles hash was the strongest technical signal available — it identifies
a specific file, so it cannot produce a false positive — and that is exactly
the problem.

Every tier must be a claim about a *cut*, never about a copy. A TMDB id
discloses "some copy of this film", which is what a library catalogue
discloses. A file hash discloses "this exact release": it made a read endpoint
into a release-level oracle, and made an instance's database a mapping from
file fingerprints to the instances holding them. That is a far more specific
disclosure than PR-005 permits, and a dataset no volunteer operator should be
asked to hold. The audio signature is the replacement: derived from content, it
identifies the cut rather than the copy, so two encodes of the same edit agree.

The field is deleted rather than kept as a vestigial null, on the same
reasoning §2 applied to `anneal_sec` — a key naming a signal the format no
longer has is actively misleading — so an upload carrying one is now an
unknown-field 400, with a test asserting it.

**Every content_id changes**, including for manifests that never carried a
hash, because the canonical `cut` object lost a key. The golden vector is
regenerated and re-verified against an independent Python implementation; the
plugin and extraction repos must adopt the new value or federation
deduplication silently breaks. Free now, pre-release; not free later.

Adds docs/legal-posture.md, the operator-facing half of what §5a asks for:
what an instance holds exhaustively, what it structurally cannot do, and how
that sits against the intermediary-liability regimes that plausibly apply.

208 tests. Coverage 25/32.

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

TRACES: UR-011 | SR-004, PR-005
2026-07-31 09:52:03 +02:00

404 lines
15 KiB
Rust

//! §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<Body>) -> (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,
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,
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<String> = (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);
}