Federation: replicate content, re-derive judgement (UR-008)
Implements §9a. The replication surface is four reads and no writes: a change feed, fetch by content_id, a batch have, and a human-facing peer directory — plus a capabilities endpoint carrying the accepted envelope versions, which lets a client discover a schema mismatch in one request instead of a 400 per manifest across a library sweep. Pull, never push: a pulling server chooses what it ingests and when. Push would let any peer inject work into the validation queue — the same abuse surface as anonymous upload, at higher volume. Nothing inherits a peer's judgement. A pulled manifest runs the full §6 stage 1 and 2 validation and this server's own cast check, and the fetched body must hash to the content_id that was asked for — the check that stops an intermediary or a misbehaving peer substituting content under a trusted id. A peer's retraction flags for review rather than delisting, because auto-delisting would hand every peer a remote delete primitive; only the opt-in per-peer abuse channel delists, because a takedown propagating at the speed of manual review is the wrong failure mode for that one case. A test caught a real bug in the first cut: the feed cursor was a ULID, and ULIDs are only monotonic *between* milliseconds — two generated in the same millisecond carry independent random components and can sort opposite to write order. A peer resuming from `seq > cursor` would then silently skip an entry: replication losing manifests with no error anywhere. The cursor is now an AUTOINCREMENT integer, and the test asserts strict monotonicity rather than merely sortedness. Peer administration is deliberately not an API. §9a requires that a peering exist only because an operator typed a URL, so nothing a remote server returns can establish or widen one; there_is_no_endpoint_that_creates_a_peering asserts that absence rather than trusting it. 212 tests. Coverage 25/32 (78%). Co-Authored-By: Claude Opus 5 <noreply@anthropic.com> TRACES: UR-008 | PR-006
This commit is contained in:
@@ -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<i64>,
|
||||
#[serde(default)]
|
||||
changes: Vec<ChangeBody>,
|
||||
}
|
||||
|
||||
#[derive(Debug, Deserialize)]
|
||||
struct ChangeBody {
|
||||
content_id: String,
|
||||
op: String,
|
||||
#[serde(default)]
|
||||
reason: Option<String>,
|
||||
#[serde(default)]
|
||||
origin: Option<String>,
|
||||
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<T: serde::de::DeserializeOwned>(&self, url: &str) -> Result<T, PullError> {
|
||||
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<PullOutcome, PullError> {
|
||||
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<bool> {
|
||||
// §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);
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user