//! §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, 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