From 5eab32bb5060b1c25a488fcf2bf9b0acda2a2f5b Mon Sep 17 00:00:00 2001 From: dawn Date: Tue, 8 Sep 2026 03:47:50 +0900 Subject: [PATCH] bobbin/{ingest,resolver}: settle repoDid claims by record createdAt Signed-off-by: dawn --- bobbin/crates/ingest/src/lib.rs | 132 +++++++++++-- bobbin/crates/resolver/src/legacy_upgrade.rs | 17 +- bobbin/crates/resolver/src/lib.rs | 189 +++++++++++++++---- bobbin/crates/xrpc/tests/cold_start.rs | 8 +- 4 files changed, 298 insertions(+), 48 deletions(-) diff --git a/bobbin/crates/ingest/src/lib.rs b/bobbin/crates/ingest/src/lib.rs index fe7bb1dce..3cd229fbc 100644 --- a/bobbin/crates/ingest/src/lib.rs +++ b/bobbin/crates/ingest/src/lib.rs @@ -9,7 +9,9 @@ use bobbin_edge_index::{ }; use bobbin_knot_ingest::{CapabilityGate, KnotRegistry}; use bobbin_record_lru::RecordStore; -use bobbin_resolver::{NormalizeRepoRefs, decode_canon_or_upgrade_bytes, synthesize_created_at}; +use bobbin_resolver::{ + NormalizeRepoRefs, RepoClaim, decode_canon_or_upgrade_bytes, synthesize_created_at, +}; use bobbin_runtime::{ Clock, Entropy, NetworkError, RuntimeHasher, UnixMicros, WsConn, WsMessage, WsStream, WsTransport, @@ -843,6 +845,11 @@ enum PendingOp { source: AtUri, nsid: Nsid, }, + // somebody else owns this record, so any copy we indexed earlier has to go too + NotIndexed { + source: AtUri, + nsid: Nsid, + }, } struct Prepared { @@ -862,7 +869,7 @@ struct Resolved { fn pending_nsid(op: &PendingOp) -> Option<&Nsid> { match op { PendingOp::Upsert(pieces) => Some(&pieces.nsid), - PendingOp::Delete { nsid, .. } => Some(nsid), + PendingOp::Delete { nsid, .. } | PendingOp::NotIndexed { nsid, .. } => Some(nsid), PendingOp::Parked { nsid, .. } => Some(nsid), PendingOp::Noop | PendingOp::ClearCache { .. } => None, } @@ -911,11 +918,20 @@ async fn claim_pending( mut pending: Pending, ctx: &PipelineCtx<'_, S>, ) -> Pending { + let mut refused = None; match &mut pending.op { PendingOp::Upsert(pieces) => { evict_from_buffer(ctx.buffer, &pieces.source).await; if let Record::Repo(repo) = &pieces.parsed { - pieces.supersedes = claim_repo(&pieces.source, repo, ctx).await; + match claim_repo(&pieces.source, repo, ctx).await { + RepoClaim::Current { displaced } => pieces.supersedes = displaced, + RepoClaim::Superseded { .. } => { + refused = Some(PendingOp::NotIndexed { + source: pieces.source.clone(), + nsid: pieces.nsid.clone(), + }); + } + } } } PendingOp::ClearCache { source } => evict_from_buffer(ctx.buffer, source).await, @@ -927,30 +943,38 @@ async fn claim_pending( ctx.resolver.forget(&ident.owner, &ident.rkey).await; } } + PendingOp::NotIndexed { source, .. } => evict_from_buffer(ctx.buffer, source).await, PendingOp::Noop | PendingOp::Parked { .. } => {} } + if let Some(op) = refused { + pending.op = op; + } pending } -/// takes the rkey for this repo and hands back the ident it displaced, if any async fn claim_repo( source: &AtUri, repo: &RepoRecord, ctx: &PipelineCtx<'_, S>, -) -> Option { - let ident = repo_ident_of(source)?; +) -> RepoClaim { + let Some(ident) = repo_ident_of(source) else { + return RepoClaim::Current { displaced: None }; + }; if let Some(shadow) = ctx.shadow { shadow.note_observed(&ident.owner, &ident.rkey).await; } - let supersedes = ctx + let claim = ctx .resolver .observe( ident.owner.clone(), ident.rkey.clone(), repo.repo_did.clone(), + &repo.created_at, ) .await; - if let Some(prior) = supersedes.as_ref() + if let RepoClaim::Current { + displaced: Some(prior), + } = &claim && let Some(prior_uri) = repo_ident_uri(prior) { evict_from_buffer(ctx.buffer, &prior_uri).await; @@ -963,7 +987,7 @@ async fn claim_repo( None => registry.observe_host(&host), } } - supersedes + claim } async fn resolve_stage( @@ -1428,7 +1452,7 @@ async fn commit_pending( log_unknown_state_variant(outcome, &source); index_search(search, resolver, &source, parsed).await; } - PendingOp::Delete { source, nsid } => { + PendingOp::Delete { source, nsid } | PendingOp::NotIndexed { source, nsid } => { store.remove_source(&source); apply_delete_to_state_index(issue_states, pull_statuses, &source, &nsid); records.remove(&source); @@ -1641,6 +1665,7 @@ mod tests { use bobbin_record_lru::{CacheCapacity, LruRecordStore, NoopRecordStore, RecordStore}; use bobbin_runtime::{OsEntropy, RuntimeHasher, SystemClock, TungsteniteWs}; use bobbin_types::search::NoopSearchSink; + use jacquard_common::types::datetime::Datetime; use jacquard_common::types::nsid::Nsid; use jacquard_common::types::tid::Tid; use serde_json::json; @@ -2823,6 +2848,83 @@ mod tests { ); } + #[tokio::test] + async fn a_backfilled_rename_alias_cannot_steal_the_repo_did() { + let (store, issue_states, pull_statuses, cov, resolver) = fresh(); + let search = NoopSearchSink; + let records = NoopRecordStore; + let ctx = PipelineCtx { + resolver: &resolver, + store: &store, + issue_states: &issue_states, + pull_statuses: &pull_statuses, + coverage: &cov, + records: &records, + search: &search, + shadow: None, + buffer: None, + knot_registry: None, + knot_gate: None, + }; + let repo_frame = |id: u64, rkey: &str, created_at: &str| { + parse_frame(json!({ + "id": id, + "type": "record", + "record": { + "live": false, + "did": "did:plc:oppi", + "rev": fresh_tid().as_str(), + "collection": "sh.tangled.repo", + "rkey": rkey, + "action": "create", + "record": { + "$type": "sh.tangled.repo", + "createdAt": created_at, + "knot": "knot1.tangled.sh", + "repoDid": "did:plc:abalone" + } + } + })) + }; + // a backfill replays a repo in rkey order, so the alias trails what replaced it + for frame in [ + repo_frame(1, "stinkpot", "2026-07-24T14:33:33Z"), + repo_frame(2, "tortu", "2026-07-20T19:26:39Z"), + ] { + let pending = claim_pending(prepare_frame(frame, &ctx, now()).await, &ctx).await; + let pending = resolve_pending(pending, &ctx).await; + commit_pending( + pending, + &store, + &issue_states, + &pull_statuses, + &cov, + &search, + &records, + &resolver, + ) + .await; + } + + let owner = Did::new_owned("did:plc:oppi").unwrap(); + assert_eq!( + resolver + .lookup_by_repo_did(&Did::new_owned("did:plc:abalone").unwrap()) + .await, + Some(RepoIdent::new(owner, rkey("stinkpot"))), + "the newest record owns the repoDid, whichever rkey the replay ends on", + ); + let key = bobbin_types::ids::EdgeKey::new( + Nsid::new_static("sh.tangled.repo").unwrap(), + did_subj("did:plc:oppi"), + ); + assert_eq!( + store.sources_for(&key), + vec![AtUri::new_owned("at://did:plc:oppi/sh.tangled.repo/stinkpot").unwrap()], + "only the renamed repo belongs in the owner's index", + ); + } + #[tokio::test] async fn unresolvable_repo_subject_drops_edge() { let (store, issue_states, pull_statuses, cov, resolver) = fresh(); @@ -3007,6 +3109,7 @@ mod tests { Did::new_owned("did:plc:nel").unwrap(), Rkey::new_owned("abcabcabcabcz").unwrap(), Some(Did::new_owned("did:plc:abalone").unwrap()), + &Datetime::raw_str("2026-05-01T00:00:00Z"), ) .await; let issue: HydrantFrame = parse_frame(json!({ @@ -3067,6 +3170,7 @@ mod tests { Did::new_owned("did:plc:nel").unwrap(), Rkey::new_owned("abcabcabcabcz").unwrap(), Some(Did::new_owned("did:plc:abalone").unwrap()), + &Datetime::raw_str("2026-05-01T00:00:00Z"), ) .await; let search = RecordingSearchSink::default(); @@ -3120,6 +3224,7 @@ mod tests { owner.clone(), rkey.clone(), Some(Did::new_owned("did:plc:abalone").unwrap()), + &Datetime::raw_str("2026-05-01T00:00:00Z"), ) .await; assert!( @@ -3306,7 +3411,12 @@ mod tests { let slow_rkeys: [Rkey; 2] = [rkey("slowrkeyaa01"), rkey("slowrkeyaa02")]; for r in &fast_rkeys { resolver - .observe(owner.clone(), r.clone(), Some(abalone.clone())) + .observe( + owner.clone(), + r.clone(), + Some(abalone.clone()), + &Datetime::raw_str("2026-05-01T00:00:00Z"), + ) .await; } diff --git a/bobbin/crates/resolver/src/legacy_upgrade.rs b/bobbin/crates/resolver/src/legacy_upgrade.rs index 3032294ad..a847dff78 100644 --- a/bobbin/crates/resolver/src/legacy_upgrade.rs +++ b/bobbin/crates/resolver/src/legacy_upgrade.rs @@ -512,6 +512,7 @@ mod tests { use bobbin_runtime::RuntimeHasher; use bobbin_types::edges::Record; use jacquard_common::DefaultStr; + use jacquard_common::types::datetime::Datetime; use jacquard_common::types::did::Did; use jacquard_common::types::recordkey::Rkey; @@ -752,6 +753,7 @@ mod tests { did("did:plc:scallop"), rkey("limpet"), Some(did("did:plc:scallop")), + &Datetime::raw_str("2026-05-01T00:00:00Z"), ) .await; let canon = match decoded { @@ -783,7 +785,12 @@ mod tests { let owner = did("did:plc:nel"); let key = rkey("abcabcabcabcz"); resolver - .observe(owner.clone(), key.clone(), Some(did("did:plc:scallop"))) + .observe( + owner.clone(), + key.clone(), + Some(did("did:plc:scallop")), + &Datetime::raw_str("2026-05-01T00:00:00Z"), + ) .await; let json = br#"{"$type":"sh.tangled.repo.issue","repo":"at://did:plc:nel/sh.tangled.repo/abcabcabcabcz","title":"t","createdAt":"2026-05-01T00:00:00Z"}"#; let legacy = @@ -920,6 +927,7 @@ mod tests { did("did:plc:nel"), rkey("abcabcabcabcz"), Some(did("did:plc:scallop")), + &Datetime::raw_str("2026-05-01T00:00:00Z"), ) .await; let json = br#"{"$type":"sh.tangled.repo.pull","title":"t","createdAt":"2026-05-01T00:00:00Z","rounds":[],"target":{"branch":"main","repo":"at://did:plc:nel/sh.tangled.repo/abcabcabcabcz","repoDid":""}}"#; @@ -1034,7 +1042,12 @@ mod tests { let owner = did("did:plc:nel"); let key = rkey("abcabcabcabcz"); resolver - .observe(owner.clone(), key.clone(), Some(did("did:plc:scallop"))) + .observe( + owner.clone(), + key.clone(), + Some(did("did:plc:scallop")), + &Datetime::raw_str("2026-05-01T00:00:00Z"), + ) .await; let json = br#"{"$type":"sh.tangled.feed.star","createdAt":"2026-05-01T00:00:00Z","subject":"at://did:plc:nel/sh.tangled.repo/abcabcabcabcz"}"#; let legacy = diff --git a/bobbin/crates/resolver/src/lib.rs b/bobbin/crates/resolver/src/lib.rs index 3d1674a02..a48b95b3e 100644 --- a/bobbin/crates/resolver/src/lib.rs +++ b/bobbin/crates/resolver/src/lib.rs @@ -18,6 +18,7 @@ use bobbin_slingshot_client::{SlingshotClient, SlingshotError}; use bobbin_types::edges::{ExtractError, Record}; use bobbin_types::ids::{RepoIdent, nsid_static}; use jacquard_common::DefaultStr; +use jacquard_common::types::datetime::Datetime; use jacquard_common::types::did::Did; use jacquard_common::types::nsid::Nsid; use jacquard_common::types::recordkey::Rkey; @@ -50,6 +51,25 @@ impl AuthoritativeResolution { } } +#[derive(Clone, Debug, Eq, PartialEq)] +struct RepoDidHolder { + ident: RepoIdent, + created_at: Datetime, +} + +impl RepoDidHolder { + // a backfill replays a repo in rkey order, so the stamp decides and not arrival + fn yields_to(&self, incoming: &Self) -> bool { + self.ident == incoming.ident || incoming.created_at >= self.created_at + } +} + +#[derive(Clone, Debug, Eq, PartialEq)] +pub enum RepoClaim { + Current { displaced: Option }, + Superseded { by: RepoIdent }, +} + #[derive(Clone, Debug, Eq, PartialEq)] enum CacheEntry { Authoritative(AuthoritativeResolution), @@ -168,7 +188,7 @@ struct SlingshotProbe { pub struct RepoIdResolver { cache: SccMap, - by_repo_did: SccMap, RepoIdent, RuntimeHasher>, + by_repo_did: SccMap, RepoDidHolder, RuntimeHasher>, in_flight: SccMap>, RuntimeHasher>, probe: Option, stats: ResolverStats, @@ -223,7 +243,7 @@ impl RepoIdResolver { self.by_repo_did .get_async(repo_did) .await - .map(|e| e.get().clone()) + .map(|e| e.get().ident.clone()) } pub async fn observe( @@ -231,7 +251,8 @@ impl RepoIdResolver { owner: Did, rkey: Rkey, repo_did: Option>, - ) -> Option { + created_at: &Datetime, + ) -> RepoClaim { let ident = RepoIdent::new(owner, rkey); let entry = CacheEntry::Authoritative(AuthoritativeResolution::from_repo_did(repo_did.clone())); @@ -241,19 +262,32 @@ impl RepoIdResolver { .and_modify(|existing| *existing = entry.clone()) .or_insert(entry); - let repo_did = repo_did?; - let mut prior: Option = None; + let Some(repo_did) = repo_did else { + return RepoClaim::Current { displaced: None }; + }; + let incoming = RepoDidHolder { + ident, + created_at: created_at.clone(), + }; + let mut outcome = RepoClaim::Current { displaced: None }; self.by_repo_did .entry_async(repo_did) .await - .and_modify(|existing| { - if *existing != ident { - prior = Some(existing.clone()); - *existing = ident.clone(); + .and_modify(|held| { + let yields = held.yields_to(&incoming); + outcome = yields + .then(|| RepoClaim::Current { + displaced: (held.ident != incoming.ident).then(|| held.ident.clone()), + }) + .unwrap_or_else(|| RepoClaim::Superseded { + by: held.ident.clone(), + }); + if yields { + *held = incoming.clone(); } }) - .or_insert(ident); - prior + .or_insert(incoming); + outcome } pub async fn forget(&self, owner: &Did, rkey: &Rkey) { @@ -265,7 +299,7 @@ impl RepoIdResolver { .map(|(_, entry)| entry.into_resolution()); if let Some(Resolution::Mapped(repo_did)) = prior_resolution { self.by_repo_did - .remove_if_async(&repo_did, |existing| *existing == ident) + .remove_if_async(&repo_did, |held| held.ident == ident) .await; } } @@ -446,40 +480,98 @@ mod tests { Arc::new(SystemClock::new()) } + fn made(day: u32) -> Datetime { + Datetime::raw_str(format!("2026-05-{day:02}T00:00:00Z")) + } + #[tokio::test] async fn observation_returns_prior_ident_when_repo_did_moves() { let resolver = RepoIdResolver::detached(RuntimeHasher::default()); - let prior = resolver + let claim = resolver .observe( did("did:plc:nel"), rkey("3liuighjy2h22"), Some(did("did:plc:clam")), + &made(1), ) .await; - assert!(prior.is_none(), "first observation has no prior"); + assert_eq!( + claim, + RepoClaim::Current { displaced: None }, + "first observation has no prior" + ); - let prior = resolver - .observe(did("did:plc:nel"), rkey("core"), Some(did("did:plc:clam"))) + let claim = resolver + .observe( + did("did:plc:nel"), + rkey("core"), + Some(did("did:plc:clam")), + &made(2), + ) .await; assert_eq!( - prior, - Some(RepoIdent::new(did("did:plc:nel"), rkey("3liuighjy2h22"))), + claim, + RepoClaim::Current { + displaced: Some(RepoIdent::new(did("did:plc:nel"), rkey("3liuighjy2h22"))), + }, "same repoDID at a new (owner, rkey) returns the prior ident so callers can evict the stale at-uri", ); - let prior = resolver - .observe(did("did:plc:nel"), rkey("core"), Some(did("did:plc:clam"))) + let claim = resolver + .observe( + did("did:plc:nel"), + rkey("core"), + Some(did("did:plc:clam")), + &made(2), + ) .await; - assert!(prior.is_none(), "re-observing the same ident is a no-op"); + assert_eq!( + claim, + RepoClaim::Current { displaced: None }, + "re-observing the same ident is a no-op" + ); + } + + #[tokio::test] + async fn a_record_older_than_the_holder_cannot_take_the_repo_did() { + let resolver = RepoIdResolver::detached(RuntimeHasher::default()); + let owner = did("did:plc:oppi"); + let repo_did = did("did:plc:clam"); + resolver + .observe( + owner.clone(), + rkey("stinkpot"), + Some(repo_did.clone()), + &made(24), + ) + .await; + let claim = resolver + .observe( + owner.clone(), + rkey("tortu"), + Some(repo_did.clone()), + &made(20), + ) + .await; + + let stinkpot = RepoIdent::new(owner, rkey("stinkpot")); + assert_eq!( + claim, + RepoClaim::Superseded { + by: stinkpot.clone() + }, + "the older record is the alias a rename left behind", + ); + assert_eq!(resolver.lookup_by_repo_did(&repo_did).await, Some(stinkpot)); } #[tokio::test] async fn observation_without_repo_did_does_not_track_reverse() { let resolver = RepoIdResolver::detached(RuntimeHasher::default()); - let prior = resolver - .observe(did("did:plc:nel"), rkey("abcabcabcabcz"), None) + let claim = resolver + .observe(did("did:plc:nel"), rkey("abcabcabcabcz"), None, &made(1)) .await; - assert!(prior.is_none()); + assert_eq!(claim, RepoClaim::Current { displaced: None }); } #[tokio::test] @@ -490,26 +582,35 @@ mod tests { did("did:plc:nel"), rkey("3liuighjy2h22"), Some(did("did:plc:clam")), + &made(1), ) .await; resolver - .observe(did("did:plc:nel"), rkey("core"), Some(did("did:plc:clam"))) + .observe( + did("did:plc:nel"), + rkey("core"), + Some(did("did:plc:clam")), + &made(1), + ) .await; resolver .forget(&did("did:plc:nel"), &rkey("3liuighjy2h22")) .await; - let prior = resolver + let claim = resolver .observe( did("did:plc:nel"), rkey("core-renamed"), Some(did("did:plc:clam")), + &made(1), ) .await; assert_eq!( - prior, - Some(RepoIdent::new(did("did:plc:nel"), rkey("core"))), + claim, + RepoClaim::Current { + displaced: Some(RepoIdent::new(did("did:plc:nel"), rkey("core"))), + }, "stale at-uri's forget must not displace the live owner of did:plc:clam", ); } @@ -522,6 +623,7 @@ mod tests { did("did:plc:nel"), rkey("abcabcabcabcz"), Some(did("did:plc:clam")), + &made(1), ) .await; let got = resolver @@ -534,7 +636,7 @@ mod tests { async fn observation_without_repo_did_resolves_no_repo_did() { let resolver = RepoIdResolver::detached(RuntimeHasher::default()); resolver - .observe(did("did:plc:nel"), rkey("abcabcabcabcz"), None) + .observe(did("did:plc:nel"), rkey("abcabcabcabcz"), None, &made(1)) .await; let got = resolver .resolve(&did("did:plc:nel"), &rkey("abcabcabcabcz")) @@ -563,6 +665,7 @@ mod tests { did("did:plc:nel"), rkey("abcabcabcabcz"), Some(did("did:plc:limpet")), + &made(1), ) .await; let got = resolver.lookup_by_repo_did(&did("did:plc:limpet")).await; @@ -583,7 +686,7 @@ mod tests { async fn lookup_by_repo_did_misses_when_repo_did_was_none() { let resolver = RepoIdResolver::detached(RuntimeHasher::default()); resolver - .observe(did("did:plc:nel"), rkey("abcabcabcabcz"), None) + .observe(did("did:plc:nel"), rkey("abcabcabcabcz"), None, &made(1)) .await; let got = resolver.lookup_by_repo_did(&did("did:plc:limpet")).await; assert_eq!(got, None); @@ -597,6 +700,7 @@ mod tests { did("did:plc:nel"), rkey("abcabcabcabcz"), Some(did("did:plc:limpet")), + &made(1), ) .await; resolver @@ -604,6 +708,7 @@ mod tests { did("did:plc:olaren"), rkey("xyzxyzxyzxyzx"), Some(did("did:plc:limpet")), + &made(1), ) .await; let got = resolver.lookup_by_repo_did(&did("did:plc:limpet")).await; @@ -621,6 +726,7 @@ mod tests { did("did:plc:nel"), rkey("abcabcabcabcz"), Some(did("did:plc:clam")), + &made(1), ) .await; resolver @@ -628,6 +734,7 @@ mod tests { did("did:plc:nel"), rkey("abcabcabcabcz"), Some(did("did:plc:uni")), + &made(1), ) .await; let got = resolver @@ -642,7 +749,12 @@ mod tests { let owner = did("did:plc:nel"); let key = rkey("abcabcabcabcz"); resolver - .observe(owner.clone(), key.clone(), Some(did("did:plc:clam"))) + .observe( + owner.clone(), + key.clone(), + Some(did("did:plc:clam")), + &made(1), + ) .await; resolver .fill_provisional( @@ -663,7 +775,9 @@ mod tests { let resolver = RepoIdResolver::detached(RuntimeHasher::default()); let owner = did("did:plc:nel"); let key = rkey("abcabcabcabcz"); - resolver.observe(owner.clone(), key.clone(), None).await; + resolver + .observe(owner.clone(), key.clone(), None, &made(1)) + .await; resolver .fill_provisional( RepoIdent::new(owner.clone(), key.clone()), @@ -1003,7 +1117,9 @@ mod tests { Resolution::Mapped(did("did:plc:clam")), ) .await; - resolver.observe(owner.clone(), key.clone(), None).await; + resolver + .observe(owner.clone(), key.clone(), None, &made(1)) + .await; let got = resolver.resolve(&owner, &key).await; assert_eq!( got, @@ -1018,7 +1134,12 @@ mod tests { let owner = did("did:plc:nel"); let key = rkey("abcabcabcabcz"); resolver - .observe(owner.clone(), key.clone(), Some(did("did:plc:clam"))) + .observe( + owner.clone(), + key.clone(), + Some(did("did:plc:clam")), + &made(1), + ) .await; assert_eq!( resolver.cached_resolution(&owner, &key).await, diff --git a/bobbin/crates/xrpc/tests/cold_start.rs b/bobbin/crates/xrpc/tests/cold_start.rs index 30bc3bea4..d58d7f690 100644 --- a/bobbin/crates/xrpc/tests/cold_start.rs +++ b/bobbin/crates/xrpc/tests/cold_start.rs @@ -12,6 +12,7 @@ use bobbin_xrpc::{AppState, router}; use futures::stream::{self, StreamExt}; use http::{Request, StatusCode}; use jacquard_common::DefaultStr; +use jacquard_common::types::datetime::Datetime; use jacquard_common::types::did::Did; use jacquard_common::types::nsid::Nsid; use jacquard_common::types::recordkey::Rkey; @@ -640,7 +641,12 @@ async fn get_repo_by_repo_did_returns_observed_record() { let state = fresh_app(&Url::parse(&server.uri()).unwrap()).await; state .resolver - .observe(owner_did.clone(), rk.clone(), Some(repo_did.clone())) + .observe( + owner_did.clone(), + rk.clone(), + Some(repo_did.clone()), + &Datetime::raw_str("2026-05-01T00:00:00Z"), + ) .await; let app = router(state); -- 2.51.2