diff --git a/Cargo.lock b/Cargo.lock index b5fb6a0f2..664fb4804 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -985,6 +985,7 @@ version = "0.0.1" dependencies = [ "axum", "bobbin-edge-index", + "bobbin-ingest", "bobbin-knot-proxy", "bobbin-record-lru", "bobbin-resolver", @@ -1003,6 +1004,7 @@ dependencies = [ "serde_json", "thiserror 2.0.18", "tokio", + "tokio-util", "tower", "tower-http 0.7.0", "tracing", diff --git a/bobbin/crates/bobbin-sim/src/runtime.rs b/bobbin/crates/bobbin-sim/src/runtime.rs index 2af009a66..2fb41c52a 100644 --- a/bobbin/crates/bobbin-sim/src/runtime.rs +++ b/bobbin/crates/bobbin-sim/src/runtime.rs @@ -126,6 +126,7 @@ impl Sim { warming_buffer: warming_buffer_enabled.then(|| warming_buffer.clone()), knot_registry: None, knot_gate: None, + settlements: None, }; let ingest_config = IngestConfig { hydrant_base, diff --git a/bobbin/crates/bobbin/src/main.rs b/bobbin/crates/bobbin/src/main.rs index c1f682077..5e24fbb05 100644 --- a/bobbin/crates/bobbin/src/main.rs +++ b/bobbin/crates/bobbin/src/main.rs @@ -6,7 +6,7 @@ use std::sync::Arc; use std::time::Duration; use anyhow::{Context, anyhow}; -use bobbin_edge_index::{CoverageWatch, EdgeStore, HydrantCursor, StateIndex}; +use bobbin_edge_index::{CoverageWatch, EdgeStore, HydrantCursor, Settlements, StateIndex}; use bobbin_ingest::{ IngestConfig, IngestRuntime, RepoIdResolver, WarmingBuffer, run as run_ingest, }; @@ -219,6 +219,7 @@ async fn run(cfg: BobbinConfig) -> anyhow::Result<()> { let issue_states = Arc::new(StateIndex::new(hasher.clone())); let pull_statuses = Arc::new(StateIndex::new(hasher.clone())); let coverage = Arc::new(CoverageWatch::new()); + let settlements = Arc::new(Settlements::new()); let warming_buffer = Arc::new(WarmingBuffer::new(hasher.clone())); let knot_registry = Arc::new(KnotRegistry::new()); let knots = Arc::new(KnotProxy::new( @@ -321,6 +322,7 @@ async fn run(cfg: BobbinConfig) -> anyhow::Result<()> { warming_buffer: Some(warming_buffer), knot_registry: Some(knot_registry.clone()), knot_gate: Some(knot_gate.clone()), + settlements: Some(settlements.clone()), }; let mut ingest_handle = tokio::spawn(run_ingest(ingest_cfg, ingest_runtime)); let _identity_warmer = tokio::spawn(identity.clone().run_warming(cancel.clone())); @@ -388,7 +390,8 @@ async fn run(cfg: BobbinConfig) -> anyhow::Result<()> { .with_limiter(limiter) .with_mirror(mirror) .with_mirror_v2(mirror_v2) - .with_proxies(trusted_proxies); + .with_proxies(trusted_proxies) + .with_settlements(settlements); let app = router(state); let _debug_server = match (debug_bind, mem_probe) { diff --git a/bobbin/crates/edge-index/Cargo.toml b/bobbin/crates/edge-index/Cargo.toml index bfdff6a6c..00777dfd5 100644 --- a/bobbin/crates/edge-index/Cargo.toml +++ b/bobbin/crates/edge-index/Cargo.toml @@ -18,3 +18,6 @@ smallvec = "1" serde = { workspace = true } thiserror = { workspace = true } tokio = { workspace = true } + +[dev-dependencies] +tokio = { workspace = true, features = ["macros", "rt", "test-util"] } diff --git a/bobbin/crates/edge-index/src/lib.rs b/bobbin/crates/edge-index/src/lib.rs index 19963fe91..f37d2294d 100644 --- a/bobbin/crates/edge-index/src/lib.rs +++ b/bobbin/crates/edge-index/src/lib.rs @@ -49,10 +49,12 @@ impl BucketKey { } pub mod coverage; +pub mod settle; mod sorted_blocks; pub mod state_index; mod state_view; pub use coverage::{Coverage, CoverageWatch, HydrantCursor, PromotionSignal}; +pub use settle::{ParsedCid, RecordOutcome, Rejection, Settlement, Settlements, Watch}; use sorted_blocks::{PageKey, SortedBlocks}; pub use state_index::{ ApplyOutcome, IssueStateKind, PullStatusKind, StateIndex, StateKind, apply_record_state, @@ -158,7 +160,17 @@ pub enum PageCursor { } #[derive( - Clone, Copy, Debug, Default, Eq, Hash, Ord, PartialEq, PartialOrd, serde::Deserialize, serde::Serialize, + Clone, + Copy, + Debug, + Default, + Eq, + Hash, + Ord, + PartialEq, + PartialOrd, + serde::Deserialize, + serde::Serialize, )] #[serde(transparent)] pub struct PageOffset(u64); @@ -213,7 +225,17 @@ impl From for PageRank { } #[derive( - Clone, Copy, Debug, Default, Eq, Hash, Ord, PartialEq, PartialOrd, serde::Deserialize, serde::Serialize, + Clone, + Copy, + Debug, + Default, + Eq, + Hash, + Ord, + PartialEq, + PartialOrd, + serde::Deserialize, + serde::Serialize, )] #[serde(transparent)] pub struct TotalCount(u64); @@ -1479,7 +1501,12 @@ mod tests { rows.into_iter().map(|(_, i)| source_uri(i)).collect() } - fn paginate_all(store: &EdgeStore, key: &EdgeKey, dir: SortDir, page: PageLimit) -> Vec { + fn paginate_all( + store: &EdgeStore, + key: &EdgeKey, + dir: SortDir, + page: PageLimit, + ) -> Vec { std::iter::successors( Some(store.list(key, PageCursor::Start, page, dir)), |prev| { @@ -1528,7 +1555,10 @@ mod tests { page: PageLimit, ) -> Vec { std::iter::successors( - Some((PageOffset::ZERO, store.list(key, PageOffset::ZERO, page, dir))), + Some(( + PageOffset::ZERO, + store.list(key, PageOffset::ZERO, page, dir), + )), |(offset, prev)| { prev.next.map(|_| { let next = PageOffset::from(offset.get() + page.get() as u64); @@ -1564,14 +1594,24 @@ mod tests { for o in [1u64, 5, 63, 255, 256, 257, 999] { let o = o.min(n as u64 - 1); - let page = store.list(&key, PageStart::Offset(PageOffset::from(o)), limit(10), SortDir::Asc); + let page = store.list( + &key, + PageStart::Offset(PageOffset::from(o)), + limit(10), + SortDir::Asc, + ); let want: Vec = asc.iter().skip(o as usize).take(10).cloned().collect(); let got: Vec = page.items.iter().map(|u| u.as_ref().to_owned()).collect(); assert_eq!(got, want, "asc offset {o} mismatch n={n}"); assert_eq!(page.total, Some(TotalCount::from(n as u64)), "total n={n}"); } - let page = store.list(&key, PageStart::Offset(PageOffset::from(n as u64)), limit(10), SortDir::Asc); + let page = store.list( + &key, + PageStart::Offset(PageOffset::from(n as u64)), + limit(10), + SortDir::Asc, + ); assert!(page.items.is_empty(), "offset==len must be empty n={n}"); assert!(page.next.is_none()); assert_eq!(page.total, Some(TotalCount::from(n as u64))); @@ -1585,9 +1625,15 @@ mod tests { let key = fill_subject(&store, n); let asc = reference(0..n); - let first = store.list(&key, PageStart::Offset(PageOffset::new(600)), limit(7), SortDir::Asc); + let first = store.list( + &key, + PageStart::Offset(PageOffset::new(600)), + limit(7), + SortDir::Asc, + ); let got: Vec = std::iter::successors(Some(first), |prev| { - prev.next.map(|tok| store.list(&key, PageCursor::After(tok), limit(7), SortDir::Asc)) + prev.next + .map(|tok| store.list(&key, PageCursor::After(tok), limit(7), SortDir::Asc)) }) .flat_map(|p| p.items.into_iter().map(|u| u.as_ref().to_owned())) .collect(); @@ -1599,9 +1645,15 @@ mod tests { let key = fill_subject(&store, 5); let desc: Vec = reference(0..5).into_iter().rev().collect(); - let first = store.list(&key, PageStart::Offset(PageOffset::new(2)), limit(3), SortDir::Desc); + let first = store.list( + &key, + PageStart::Offset(PageOffset::new(2)), + limit(3), + SortDir::Desc, + ); let got: Vec = std::iter::successors(Some(first), |prev| { - prev.next.map(|tok| store.list(&key, PageCursor::After(tok), limit(3), SortDir::Desc)) + prev.next + .map(|tok| store.list(&key, PageCursor::After(tok), limit(3), SortDir::Desc)) }) .flat_map(|p| p.items.into_iter().map(|u| u.as_ref().to_owned())) .collect(); @@ -1614,9 +1666,15 @@ mod tests { let key = fill_subject(&store, n); let desc: Vec = reference(0..n).into_iter().rev().collect(); - let first = store.list(&key, PageStart::Offset(PageOffset::new(600)), limit(7), SortDir::Desc); + let first = store.list( + &key, + PageStart::Offset(PageOffset::new(600)), + limit(7), + SortDir::Desc, + ); let got: Vec = std::iter::successors(Some(first), |prev| { - prev.next.map(|tok| store.list(&key, PageCursor::After(tok), limit(7), SortDir::Desc)) + prev.next + .map(|tok| store.list(&key, PageCursor::After(tok), limit(7), SortDir::Desc)) }) .flat_map(|p| p.items.into_iter().map(|u| u.as_ref().to_owned())) .collect(); @@ -1634,15 +1692,30 @@ mod tests { for &(dir, ref_order, label) in &[(SortDir::Asc, &asc, "asc"), (SortDir::Desc, &desc, "desc")] { - let p = store.list(&key, PageStart::Offset(PageOffset::ZERO), limit(n as u32), dir); + let p = store.list( + &key, + PageStart::Offset(PageOffset::ZERO), + limit(n as u32), + dir, + ); let got: Vec = p.items.iter().map(|u| u.as_ref().to_owned()).collect(); assert_eq!(got, ref_order[..n].to_vec(), "n={n} {label} offset=0"); - let p = store.list(&key, PageStart::Offset(PageOffset::new(1)), limit(n as u32), dir); + let p = store.list( + &key, + PageStart::Offset(PageOffset::new(1)), + limit(n as u32), + dir, + ); let got: Vec = p.items.iter().map(|u| u.as_ref().to_owned()).collect(); assert_eq!(got, ref_order[1..].to_vec(), "n={n} {label} offset=1"); - let p = store.list(&key, PageStart::Offset(PageOffset::from(n as u64 - 1)), limit(n as u32), dir); + let p = store.list( + &key, + PageStart::Offset(PageOffset::from(n as u64 - 1)), + limit(n as u32), + dir, + ); let got: Vec = p.items.iter().map(|u| u.as_ref().to_owned()).collect(); assert_eq!( got, @@ -1651,9 +1724,18 @@ mod tests { ); assert!(p.next.is_none(), "last item has no continuation"); - let p = store.list(&key, PageStart::Offset(PageOffset::from(n as u64)), limit(n as u32), dir); + let p = store.list( + &key, + PageStart::Offset(PageOffset::from(n as u64)), + limit(n as u32), + dir, + ); assert!(p.items.is_empty(), "n={n} {label} offset=len is empty"); - assert_eq!(p.total, Some(TotalCount::from(n as u64)), "past-end page still reports total"); + assert_eq!( + p.total, + Some(TotalCount::from(n as u64)), + "past-end page still reports total" + ); } } } @@ -1684,7 +1766,12 @@ mod tests { assert_eq!(got, expected, "descending cursor after index {target_idx}"); } - let page = store.list(&key, PageStart::Offset(PageOffset::new(256)), limit(10), SortDir::Desc); + let page = store.list( + &key, + PageStart::Offset(PageOffset::new(256)), + limit(10), + SortDir::Desc, + ); let got: Vec = page.items.iter().map(|u| u.as_ref().to_owned()).collect(); assert_eq!(got, desc[256..266].to_vec(), "descending offset 256"); } @@ -1757,9 +1844,11 @@ mod tests { for &(dir, ref_order, label) in &[(SortDir::Asc, &asc, "asc"), (SortDir::Desc, &desc, "desc")] { - let first = store.list(&key, PageStart::Offset(PageOffset::from(o)), limit(25), dir); + let first = + store.list(&key, PageStart::Offset(PageOffset::from(o)), limit(25), dir); let got: Vec = std::iter::successors(Some(first), |prev| { - prev.next.map(|tok| store.list(&key, PageCursor::After(tok), limit(25), dir)) + prev.next + .map(|tok| store.list(&key, PageCursor::After(tok), limit(25), dir)) }) .flat_map(|p| p.items.into_iter().map(|u| u.as_ref().to_owned())) .collect(); @@ -1816,7 +1905,11 @@ mod tests { "select rank {rank}" ); } - assert_eq!(blocks.select(PageRank::new(oracle.len())), None, "select past end"); + assert_eq!( + blocks.select(PageRank::new(oracle.len())), + None, + "select past end" + ); let asc: Vec = blocks.directed(PageCursor::Start, SortDir::Asc).collect(); assert_eq!(asc, oracle, "asc full iteration"); @@ -2149,10 +2242,7 @@ mod tests { SortDir::Asc, None, ); - assert_eq!( - page.total, - Some(TotalCount::new(5)), - ); + assert_eq!(page.total, Some(TotalCount::new(5)),); assert_eq!( page.items .iter() diff --git a/bobbin/crates/edge-index/src/settle.rs b/bobbin/crates/edge-index/src/settle.rs new file mode 100644 index 000000000..a03f13479 --- /dev/null +++ b/bobbin/crates/edge-index/src/settle.rs @@ -0,0 +1,540 @@ +use std::collections::{HashMap, VecDeque}; +use std::str::FromStr; +use std::sync::{Arc, Mutex, MutexGuard}; +use std::time::Duration; + +use jacquard_common::DefaultStr; +use jacquard_common::types::cid::IpldCid; +use jacquard_common::types::string::{AtUri, Cid}; +use serde::Deserialize; +use smallvec::SmallVec; +use tokio::sync::oneshot; +use tokio::time::Instant; + +// a settlement only has to outlive the round trip from the writer's pds back to us +const RETENTION: Duration = Duration::from_secs(30); +// best effort, counted in distinct at-uris so a writer rewriting one record +// cant push everyone else out +const CAPACITY: usize = 4096; +// enough revisions to still answer a writer whose record was superseded mid flight +const MAX_REVISIONS: usize = 8; +// the detail can echo caller-controlled bytes (parse errors) +const DETAIL_LIMIT: usize = 200; + +/// why a record bobbin saw will not show up in its index +#[derive(Clone, Copy, Debug, Eq, PartialEq)] +pub enum Rejection { + UnknownCollection, + Undecodable, + MissingBody, + NotIndexed, +} + +impl Rejection { + pub const fn as_str(self) -> &'static str { + match self { + Self::UnknownCollection => "unknownCollection", + Self::Undecodable => "undecodable", + Self::MissingBody => "missingBody", + Self::NotIndexed => "notIndexed", + } + } +} + +/// what bobbin did with a record once it came off the firehose +#[derive(Clone, Debug, Eq, PartialEq)] +pub enum RecordOutcome { + Indexed, + Deleted, + Rejected { + reason: Rejection, + detail: Option>, + }, +} + +impl RecordOutcome { + pub fn rejected(reason: Rejection, detail: Option<&str>) -> Self { + Self::Rejected { + reason, + detail: detail.map(clamp), + } + } + + pub const fn as_str(&self) -> &'static str { + match self { + Self::Indexed => "indexed", + Self::Deleted => "deleted", + Self::Rejected { .. } => "rejected", + } + } + + pub const fn rejection(&self) -> Option { + match self { + Self::Rejected { reason, .. } => Some(*reason), + Self::Indexed | Self::Deleted => None, + } + } + + pub fn detail(&self) -> Option<&str> { + match self { + Self::Rejected { detail, .. } => detail.as_deref(), + Self::Indexed | Self::Deleted => None, + } + } +} + +fn clamp(s: &str) -> Box { + let mut end = DETAIL_LIMIT.min(s.len()); + while !s.is_char_boundary(end) { + end -= 1; + } + s[..end].into() +} + +#[derive(Clone, Copy, Debug, Eq, Hash, PartialEq)] +pub struct ParsedCid(IpldCid); + +impl From for ParsedCid { + fn from(cid: IpldCid) -> Self { + Self(cid) + } +} + +impl FromStr for ParsedCid { + type Err = ::Err; + + fn from_str(s: &str) -> Result { + s.parse::().map(Self) + } +} + +impl serde::Serialize for ParsedCid { + fn serialize(&self, serializer: S) -> Result { + serializer.collect_str(&self.0) + } +} + +impl<'de> Deserialize<'de> for ParsedCid { + fn deserialize(deserializer: D) -> Result + where + D: serde::Deserializer<'de>, + { + let raw = Cid::::deserialize(deserializer)?; + raw.to_ipld().map(Self).map_err(serde::de::Error::custom) + } +} + +#[derive(Clone, Debug, Eq, PartialEq)] +pub struct Settlement { + pub source: AtUri, + pub cid: Option, + pub outcome: RecordOutcome, +} + +#[derive(Debug, Default)] +pub struct Settlements { + inner: Mutex, +} + +#[derive(Debug, Default)] +struct Inner { + // per-revision history so a superseded record can still be resolved + recent: HashMap, SmallVec<[Recent; 1]>>, + // first-seen order, exactly one entry per recent key + order: VecDeque>, + waiters: HashMap, Vec>, + next_seq: u64, +} + +#[derive(Debug)] +struct Recent { + at: Instant, + cid: Option, + outcome: RecordOutcome, +} + +#[derive(Debug)] +struct Waiter { + seq: u64, + cid: Option, + tx: oneshot::Sender, +} + +pub enum Watch { + Settled(Settlement), + /// the registration's Drop cancels the wait, so timeouts and + /// abandoned requests don't leave a waiter behind + Waiting(Registration, oneshot::Receiver), +} + +/// one entry in the waiter list. dropping it removes the entry, so a +/// settlement that arrives later is only remembered, not sent +pub struct Registration { + watch: Arc, + source: AtUri, + seq: u64, +} + +impl Drop for Registration { + fn drop(&mut self) { + let mut inner = self.watch.lock(); + let Some(waiting) = inner.waiters.get_mut(&self.source) else { + return; + }; + waiting.retain(|w| w.seq != self.seq); + if waiting.is_empty() { + inner.waiters.remove(&self.source); + } + } +} + +impl Settlements { + pub fn new() -> Self { + Self::default() + } + + fn lock(&self) -> MutexGuard<'_, Inner> { + self.inner.lock().unwrap_or_else(|e| e.into_inner()) + } + + pub fn settle(&self, settlement: Settlement) { + let now = Instant::now(); + let winners = { + let mut inner = self.lock(); + let winners = inner.take_matching(&settlement); + inner.remember(now, settlement.clone()); + inner.prune(now); + winners + }; + for tx in winners { + let _ = tx.send(settlement.clone()); + } + } + + pub fn watch(self: &Arc, source: &AtUri, cid: Option) -> Watch { + let now = Instant::now(); + let mut inner = self.lock(); + inner.prune(now); + if let Some(recent) = inner.find(source, cid) { + return Watch::Settled(Settlement { + source: source.clone(), + cid: recent.cid, + outcome: recent.outcome.clone(), + }); + } + let seq = inner.next_seq(); + let (tx, rx) = oneshot::channel(); + inner + .waiters + .entry(source.clone()) + .or_default() + .push(Waiter { seq, cid, tx }); + let registration = Registration { + watch: self.clone(), + source: source.clone(), + seq, + }; + Watch::Waiting(registration, rx) + } +} + +impl Inner { + fn next_seq(&mut self) -> u64 { + self.next_seq += 1; + self.next_seq + } + + fn take_matching(&mut self, s: &Settlement) -> Vec> { + if let Some(waiting) = self.waiters.remove(&s.source) { + let (hit, miss): (Vec<_>, Vec<_>) = waiting + .into_iter() + .partition(|w| w.cid.is_none_or(|want| s.cid == Some(want))); + if !miss.is_empty() { + self.waiters.insert(s.source.clone(), miss); + } + hit.into_iter().map(|w| w.tx).collect() + } else { + Vec::new() + } + } + + fn find(&self, source: &AtUri, cid: Option) -> Option<&Recent> { + self.recent + .get(source)? + .iter() + .rev() + .find(|r| cid.is_none_or(|want| r.cid == Some(want))) + } + + fn remember(&mut self, now: Instant, settlement: Settlement) { + let entries = match self.recent.entry(settlement.source) { + std::collections::hash_map::Entry::Vacant(v) => { + self.order.push_back(v.key().clone()); + v.insert(SmallVec::new()) + } + std::collections::hash_map::Entry::Occupied(o) => o.into_mut(), + }; + entries.push(Recent { + at: now, + cid: settlement.cid, + outcome: settlement.outcome, + }); + entries.retain(|r| now.duration_since(r.at) < RETENTION); + while entries.len() > MAX_REVISIONS { + entries.remove(0); + } + } + + fn prune(&mut self, now: Instant) { + while let Some(source) = self.order.pop_front() { + let over = self.recent.len() > CAPACITY; + let keep = match self.recent.get_mut(&source) { + None => false, + Some(entries) => { + entries.retain(|r| now.duration_since(r.at) < RETENTION); + !over && !entries.is_empty() + } + }; + if keep { + self.order.push_front(source); + break; + } + self.recent.remove(&source); + } + } +} + +#[cfg(test)] +mod tests { + use super::*; + + fn uri(rkey: &str) -> AtUri { + AtUri::new_owned(format!("at://did:plc:alice/sh.tangled.feed.star/{rkey}")).unwrap() + } + + const CID_A: &str = "bafyreieqygohnz2zqyvtvktbjpvhutphobcmbsnt4q5lc36ri7vpcmoz4i"; + const CID_B: &str = "bafkreigh2akiscaildc7gnvtklbsfhdgwz72eolmpckbqr5ej26byp3uli"; + + fn cid(s: &str) -> ParsedCid { + s.parse().unwrap() + } + + // a structurally valid cidv1 whose digest is a counter, so every mint differs + fn mint(i: u32) -> ParsedCid { + let mut raw = vec![0x01u8, 0x55, 0x12, 0x20]; + raw.extend_from_slice(&[0u8; 28]); + raw.extend_from_slice(&i.to_be_bytes()); + ParsedCid::from(IpldCid::try_from(raw.as_slice()).unwrap()) + } + + fn indexed(rkey: &str, c: Option) -> Settlement { + Settlement { + source: uri(rkey), + cid: c, + outcome: RecordOutcome::Indexed, + } + } + + #[test] + fn a_long_detail_is_clamped_on_a_char_boundary() { + let long = "\u{1f600}".repeat(200); + let outcome = RecordOutcome::rejected(Rejection::Undecodable, Some(&long)); + let detail = outcome.detail().expect("detail kept"); + assert!(detail.len() <= DETAIL_LIMIT); + assert!(long.starts_with(detail)); + } + + fn settled(watch: Watch) -> Option { + match watch { + Watch::Settled(s) => Some(s), + Watch::Waiting(..) => None, + } + } + + #[tokio::test(start_paused = true)] + async fn settlement_before_the_ask_is_still_there() { + let watch = Arc::new(Settlements::new()); + watch.settle(indexed("a", Some(cid(CID_A)))); + let got = settled(watch.watch(&uri("a"), Some(cid(CID_A)))).expect("already settled"); + assert_eq!(got.outcome, RecordOutcome::Indexed); + } + + #[tokio::test(start_paused = true)] + async fn a_waiter_wakes_on_its_own_cid() { + let watch = Arc::new(Settlements::new()); + let Watch::Waiting(_reg, rx) = watch.watch(&uri("a"), Some(cid(CID_B))) else { + panic!("nothing has settled yet"); + }; + watch.settle(indexed("a", Some(cid(CID_A)))); + watch.settle(Settlement { + source: uri("a"), + cid: Some(cid(CID_B)), + outcome: RecordOutcome::rejected(Rejection::Undecodable, None), + }); + let got = rx.await.expect("waiter woken"); + assert_eq!( + got.outcome, + RecordOutcome::rejected(Rejection::Undecodable, None) + ); + } + + #[tokio::test(start_paused = true)] + async fn a_cidless_waiter_takes_the_next_settlement() { + let watch = Arc::new(Settlements::new()); + let Watch::Waiting(_reg, rx) = watch.watch(&uri("a"), None) else { + panic!("nothing has settled yet"); + }; + watch.settle(Settlement { + source: uri("a"), + cid: None, + outcome: RecordOutcome::Deleted, + }); + assert_eq!(rx.await.unwrap().outcome, RecordOutcome::Deleted); + } + + // a second write to the same uri must not eat the answer the first + // writer is still waiting for + #[tokio::test(start_paused = true)] + async fn a_superseded_revision_keeps_its_own_outcome() { + let watch = Arc::new(Settlements::new()); + watch.settle(Settlement { + source: uri("a"), + cid: Some(cid(CID_A)), + outcome: RecordOutcome::rejected(Rejection::Undecodable, None), + }); + watch.settle(indexed("a", Some(cid(CID_B)))); + let old = settled(watch.watch(&uri("a"), Some(cid(CID_A)))).expect("older revision"); + assert_eq!( + old.outcome, + RecordOutcome::rejected(Rejection::Undecodable, None) + ); + let new = settled(watch.watch(&uri("a"), Some(cid(CID_B)))).expect("newer revision"); + assert_eq!(new.outcome, RecordOutcome::Indexed); + let any = settled(watch.watch(&uri("a"), None)).expect("newest wins without a cid"); + assert_eq!(any.cid, Some(cid(CID_B))); + } + + #[tokio::test(start_paused = true)] + async fn an_unrelated_revision_does_not_wake_a_waiter() { + let watch = Arc::new(Settlements::new()); + let Watch::Waiting(_reg, mut rx) = watch.watch(&uri("a"), Some(cid(CID_B))) else { + panic!("nothing has settled yet"); + }; + watch.settle(indexed("a", Some(cid(CID_A)))); + assert!(rx.try_recv().is_err()); + } + + #[tokio::test(start_paused = true)] + async fn dropping_the_registration_takes_the_waiter_back_out() { + let watch = Arc::new(Settlements::new()); + let Watch::Waiting(reg, rx) = watch.watch(&uri("a"), None) else { + panic!("nothing has settled yet"); + }; + drop(reg); + assert!(watch.lock().waiters.is_empty()); + drop(rx); + } + + #[tokio::test(start_paused = true)] + async fn settlements_expire() { + let watch = Arc::new(Settlements::new()); + watch.settle(indexed("a", None)); + tokio::time::advance(RETENTION + Duration::from_secs(1)).await; + watch.settle(indexed("b", None)); + assert!(settled(watch.watch(&uri("a"), None)).is_none()); + assert!(settled(watch.watch(&uri("b"), None)).is_some()); + } + + #[tokio::test(start_paused = true)] + async fn a_later_revision_outlives_an_expired_earlier_one() { + let watch = Arc::new(Settlements::new()); + watch.settle(indexed("a", Some(cid(CID_A)))); + tokio::time::advance(RETENTION - Duration::from_secs(1)).await; + watch.settle(indexed("a", Some(cid(CID_B)))); + tokio::time::advance(Duration::from_secs(2)).await; + watch.settle(indexed("z", None)); + assert!(settled(watch.watch(&uri("a"), Some(cid(CID_A)))).is_none()); + let got = settled(watch.watch(&uri("a"), Some(cid(CID_B)))).expect("newer survives"); + assert_eq!(got.outcome, RecordOutcome::Indexed); + } + + #[tokio::test(start_paused = true)] + async fn the_recent_map_stays_bounded() { + let watch = Arc::new(Settlements::new()); + let keys = vec!["a", "b", "c", "d", "e"]; + for key in keys { + watch.settle(indexed(key, None)); + } + for i in 0..(MAX_REVISIONS * 3) { + watch.settle(indexed("c", Some(mint(i as u32)))); + } + for i in 0..3 { + watch.settle(indexed("a", Some(mint(i as u32)))); + } + let inner = watch.lock(); + assert_eq!(inner.order.len(), inner.recent.len()); + for entries in inner.recent.values() { + assert!(entries.len() <= MAX_REVISIONS); + } + } + + #[tokio::test(start_paused = true)] + async fn a_rewrite_flood_does_not_evict_other_records() { + let watch = Arc::new(Settlements::new()); + watch.settle(indexed("b", None)); + for i in 0..(CAPACITY + 100) { + watch.settle(indexed("a", Some(mint(i as u32)))); + } + let watch_b = watch.watch(&uri("b"), None); + assert!(matches!(watch_b, Watch::Settled(_))); + let inner = watch.lock(); + assert_eq!(inner.order.len(), inner.recent.len()); + assert_eq!(inner.order.len(), 2); + } + + #[tokio::test(start_paused = true)] + async fn only_the_last_few_revisions_of_a_record_survive() { + let watch = Arc::new(Settlements::new()); + let u = uri("a"); + for i in 0..(MAX_REVISIONS + 5) { + watch.settle(indexed("a", Some(mint(i as u32)))); + } + for i in 0..5 { + assert!( + settled(watch.watch(&u, Some(mint(i as u32)))).is_none(), + "cid c{i} should be gone" + ); + } + for i in 5..(MAX_REVISIONS + 5) { + let s = settled(watch.watch(&u, Some(mint(i as u32)))).expect("should be findable"); + assert_eq!(s.cid, Some(mint(i as u32))); + } + } + + #[tokio::test(start_paused = true)] + async fn too_many_records_evicts_the_oldest() { + let watch = Arc::new(Settlements::new()); + for i in 0..(CAPACITY + 10) { + watch.settle(indexed(&format!("r{i}"), None)); + } + let inner = watch.lock(); + assert!(inner.recent.len() <= CAPACITY); + drop(inner); + assert!( + settled(watch.watch(&uri("r0"), None)).is_none(), + "first should be gone" + ); + assert!( + settled(watch.watch(&uri(&format!("r{}", CAPACITY + 9)), None)).is_some(), + "last should survive" + ); + } + + #[tokio::test(start_paused = true)] + async fn expires_after_retention_with_no_pressure() { + let watch = Arc::new(Settlements::new()); + watch.settle(indexed("a", None)); + tokio::time::advance(RETENTION + Duration::from_secs(1)).await; + assert!(settled(watch.watch(&uri("a"), None)).is_none()); + } +} diff --git a/bobbin/crates/edge-index/src/state_view.rs b/bobbin/crates/edge-index/src/state_view.rs index 579b52a94..92806cd10 100644 --- a/bobbin/crates/edge-index/src/state_view.rs +++ b/bobbin/crates/edge-index/src/state_view.rs @@ -1,4 +1,3 @@ - use std::collections::HashMap; use std::num::NonZeroU32; use std::sync::{RwLock, RwLockReadGuard, RwLockWriteGuard}; diff --git a/bobbin/crates/ingest/examples/smoke.rs b/bobbin/crates/ingest/examples/smoke.rs index dffccecdb..28b02d2bc 100644 --- a/bobbin/crates/ingest/examples/smoke.rs +++ b/bobbin/crates/ingest/examples/smoke.rs @@ -50,6 +50,7 @@ async fn main() { warming_buffer: None, knot_registry: None, knot_gate: None, + settlements: None, }; let task = tokio::spawn(async move { let _ = run(cfg, runtime).await; diff --git a/bobbin/crates/ingest/src/lib.rs b/bobbin/crates/ingest/src/lib.rs index d555ce98a..695d461cd 100644 --- a/bobbin/crates/ingest/src/lib.rs +++ b/bobbin/crates/ingest/src/lib.rs @@ -4,8 +4,9 @@ use std::sync::Arc; use std::time::Duration; use bobbin_edge_index::{ - ApplyOutcome, Coverage, CoverageWatch, EdgeStore, HydrantCursor, IssueStateKind, - PromotionSignal, PullStatusKind, StateIndex, delete_record_indexes, upsert_record_indexes, + ApplyOutcome, Coverage, CoverageWatch, EdgeStore, HydrantCursor, IssueStateKind, ParsedCid, + PromotionSignal, PullStatusKind, RecordOutcome, Rejection, Settlement, Settlements, StateIndex, + delete_record_indexes, upsert_record_indexes, }; use bobbin_knot_ingest::{CapabilityGate, KnotRegistry}; use bobbin_record_lru::RecordStore; @@ -223,6 +224,7 @@ pub struct IngestRuntime { pub warming_buffer: Option>, pub knot_registry: Option>, pub knot_gate: Option>, + pub settlements: Option>, } impl Clone for IngestRuntime { @@ -245,6 +247,7 @@ impl Clone for IngestRuntime { warming_buffer: self.warming_buffer.clone(), knot_registry: self.knot_registry.clone(), knot_gate: self.knot_gate.clone(), + settlements: self.settlements.clone(), } } } @@ -263,6 +266,7 @@ impl IngestRuntime { buffer: self.warming_buffer.as_deref(), knot_registry: self.knot_registry.as_deref(), knot_gate: self.knot_gate.as_deref(), + settlements: self.settlements.as_deref(), } } } @@ -279,6 +283,7 @@ struct PipelineCtx<'a, S: SearchSink + 'static> { buffer: Option<&'a WarmingBuffer>, knot_registry: Option<&'a KnotRegistry>, knot_gate: Option<&'a CapabilityGate>, + settlements: Option<&'a Settlements>, } pub async fn run( @@ -838,8 +843,10 @@ enum PendingOp { Noop, Identity(Option), Account(Option), - ClearCache { + Rejected { source: AtUri, + cid: Option>, + outcome: RecordOutcome, }, Upsert(Box), Parked { @@ -849,6 +856,13 @@ enum PendingOp { source: AtUri, nsid: Nsid, }, + // a native knot is authoritative for this record. any copy indexed + // while the knot was still legacy has to go too + KnotOwned { + source: AtUri, + nsid: Nsid, + cid: Option>, + }, } struct Prepared { @@ -868,12 +882,12 @@ 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::KnotOwned { nsid, .. } => Some(nsid), PendingOp::Parked { nsid, .. } => Some(nsid), PendingOp::Noop | PendingOp::Identity(_) | PendingOp::Account(_) - | PendingOp::ClearCache { .. } => None, + | PendingOp::Rejected { .. } => None, } } @@ -927,7 +941,7 @@ async fn claim_pending( pieces.supersedes = claim_repo(&pieces.source, repo, ctx).await; } } - PendingOp::ClearCache { source } => evict_from_buffer(ctx.buffer, source).await, + PendingOp::Rejected { source, .. } => evict_from_buffer(ctx.buffer, source).await, PendingOp::Delete { source, nsid } => { evict_from_buffer(ctx.buffer, source).await; if nsid.as_ref() == "sh.tangled.repo" @@ -936,6 +950,7 @@ async fn claim_pending( ctx.resolver.forget(&ident.owner, &ident.rkey).await; } } + PendingOp::KnotOwned { source, .. } => evict_from_buffer(ctx.buffer, source).await, PendingOp::Noop | PendingOp::Identity(_) | PendingOp::Account(_) @@ -1013,7 +1028,7 @@ async fn commit_stage( resolve_end, } = staged; let commit_start = rt.clock.now_instant(); - commit_pending( + let settled = commit_pending( pending, &rt.store, &rt.issue_states, @@ -1026,6 +1041,9 @@ async fn commit_stage( ) .await; let commit_end = rt.clock.now_instant(); + if let Some((watch, settled)) = rt.settlements.as_deref().zip(settled) { + watch.settle(settled); + } tracing::trace!( target: "bobbin_ingest::stage", cursor, @@ -1091,11 +1109,12 @@ async fn prepare_record( return PendingOp::Noop; } }; + let cid = record.cid.clone(); match record.action { RecordAction::Create | RecordAction::Update => { let Some(raw) = record.record else { debug!(collection = %nsid, "create/update missing record body, clearing cache"); - return PendingOp::ClearCache { source }; + return rejected(source, cid, Rejection::MissingBody, None); }; let raw_bytes = Bytes::copy_from_slice(raw.get().as_bytes()); let wire_bytes = match fallback_rfc3339(&record.rkey, &record.rev) @@ -1120,11 +1139,16 @@ async fn prepare_record( } Err(ExtractError::UnknownCollection(name)) => { debug!(collection = %name, "unknown sh.tangled.* collection, clearing cache"); - return PendingOp::ClearCache { source }; + return rejected(source, cid, Rejection::UnknownCollection, Some(&name)); + } + Err(ExtractError::UpgradeFailed(name)) => { + debug!(collection = %name, "legacy record could not be upgraded, clearing cache"); + let detail = format!("no upgrade path from the legacy {name} shape"); + return rejected(source, cid, Rejection::Undecodable, Some(&detail)); } Err(e) => { warn!(?e, collection = %record.collection, "record decode failed, clearing cache"); - return PendingOp::ClearCache { source }; + return rejected(source, cid, Rejection::Undecodable, Some(&e.to_string())); } }; match acl_disposition(&parsed, ctx.knot_gate, ctx.knot_registry) { @@ -1132,7 +1156,7 @@ async fn prepare_record( if let Some(registry) = ctx.knot_registry { registry.forget_legacy_member(&source); } - return PendingOp::Delete { source, nsid }; + return PendingOp::KnotOwned { source, nsid, cid }; } AclDisposition::LegacyMember { host } => { if let Some(registry) = ctx.knot_registry { @@ -1146,7 +1170,7 @@ async fn prepare_record( Ok(es) => es, Err(e) => { warn!(?e, "edge extraction failed, clearing cache"); - return PendingOp::ClearCache { source }; + return rejected(source, cid, Rejection::Undecodable, Some(&e.to_string())); } }; let _ = ctx.store.intern_source(&source); @@ -1175,6 +1199,19 @@ async fn prepare_record( } } +fn rejected( + source: AtUri, + cid: Option>, + reason: Rejection, + detail: Option<&str>, +) -> PendingOp { + PendingOp::Rejected { + source, + cid, + outcome: RecordOutcome::rejected(reason, detail), + } +} + enum AclDisposition { Other, NativeSkip, @@ -1251,7 +1288,16 @@ async fn resolve_pending( } = pending; let op = match op { PendingOp::Upsert(pieces) => match try_park_warming(ctx, cursor, pieces).await { - ParkOutcome::Parked { nsid } => PendingOp::Parked { nsid }, + ParkOutcome::Parked { nsid } => { + // nothing is visible yet, finalize_drained settles it once the + // deps land + return Pending { + cursor, + signal, + regime, + op: PendingOp::Parked { nsid }, + }; + } ParkOutcome::Passthrough(mut pieces) => { let edges = std::mem::take(&mut pieces.edges); pieces.edges = @@ -1392,6 +1438,7 @@ async fn finalize_drained( } = upsert; let edges = normalize_subjects(edges, ctx.resolver, ctx.coverage, None).await; remove_superseded_source(ctx.store, ctx.records, ctx.search, supersedes.as_ref()).await; + let settled_cid = settle_cid(cid.as_ref()); cache_body(ctx.records, &source, cid, bytes); let outcome = upsert_record_indexes( ctx.store, @@ -1403,9 +1450,26 @@ async fn finalize_drained( ); log_unknown_state_variant(outcome, &source); index_search(ctx.search, ctx.resolver, &source, parsed).await; + if let Some(watch) = ctx.settlements { + watch.settle(Settlement { + source, + cid: settled_cid, + outcome: RecordOutcome::Indexed, + }); + } } } +// hydrant already decoded these cids from car bytes, so garbage here means a bug upstream +fn settle_cid(cid: Option<&Cid>) -> Option { + cid.and_then(|c| { + c.to_ipld() + .map(ParsedCid::from) + .inspect_err(|e| warn!(?e, "unparseable frame cid, settling without it")) + .ok() + }) +} + #[allow(clippy::too_many_arguments)] async fn commit_pending( pending: Pending, @@ -1417,15 +1481,17 @@ async fn commit_pending( records: &dyn RecordStore, resolver: &RepoIdResolver, identity: &IdentityResolver, -) { +) -> Option { let Pending { cursor, signal, regime: _, op, } = pending; - match op { - PendingOp::Noop | PendingOp::Parked { .. } => {} + // the settlement goes out only after its op is applied, so a writer asking + // about this record sees the same thing a reader would + let settled = match op { + PendingOp::Noop | PendingOp::Parked { .. } => None, PendingOp::Identity(observed) => { if let Some(observed) = observed { match observed.handle { @@ -1433,6 +1499,7 @@ async fn commit_pending( None => identity.refresh(&observed.did), } } + None } PendingOp::Account(account) => { if let Some(account) = account @@ -1440,8 +1507,20 @@ async fn commit_pending( { identity.deactivate(account.did); } + None + } + PendingOp::Rejected { + source, + cid, + outcome, + } => { + records.remove(&source); + Some(Settlement { + source, + cid: settle_cid(cid.as_ref()), + outcome, + }) } - PendingOp::ClearCache { source } => records.remove(&source), PendingOp::Upsert(pieces) => { let UpsertPieces { source, @@ -1453,19 +1532,44 @@ async fn commit_pending( supersedes, } = *pieces; remove_superseded_source(store, records, search, supersedes.as_ref()).await; + let settled_cid = settle_cid(cid.as_ref()); cache_body(records, &source, cid, bytes); let outcome = upsert_record_indexes(store, issue_states, pull_statuses, &source, edges, &parsed); log_unknown_state_variant(outcome, &source); index_search(search, resolver, &source, parsed).await; + Some(Settlement { + source, + cid: settled_cid, + outcome: RecordOutcome::Indexed, + }) } PendingOp::Delete { source, nsid } => { delete_record_indexes(store, issue_states, pull_statuses, &source, &nsid); records.remove(&source); search.remove(&source).await; + Some(Settlement { + source, + cid: None, + outcome: RecordOutcome::Deleted, + }) } - } + PendingOp::KnotOwned { source, nsid, cid } => { + delete_record_indexes(store, issue_states, pull_statuses, &source, &nsid); + records.remove(&source); + search.remove(&source).await; + Some(Settlement { + source, + cid: settle_cid(cid.as_ref()), + outcome: RecordOutcome::rejected( + Rejection::NotIndexed, + Some("the knot is authoritative for this record"), + ), + }) + } + }; coverage.update(|c| c.advance(cursor).maybe_promote(signal)); + settled } async fn index_search( @@ -1510,6 +1614,7 @@ async fn handle_frame( buffer: None, knot_registry: None, knot_gate: None, + settlements: None, }; let pending = claim_pending(prepare_frame(frame, &ctx, now).await, &ctx).await; let pending = resolve_pending(pending, &ctx).await; @@ -1783,6 +1888,7 @@ mod tests { buffer: None, knot_registry: Some(®istry), knot_gate: Some(&gate), + settlements: None, }; let member_frame = |id: u64, rkey: &str, domain: &str| { @@ -1809,7 +1915,7 @@ mod tests { let native = prepare_frame(member_frame(1, "aaaaaaaaaaaaz", &native_host), &ctx, now()).await; assert!( - matches!(native.op, PendingOp::Delete { .. }), + matches!(native.op, PendingOp::KnotOwned { .. }), "member record for a native knot must be dropped" ); @@ -1973,6 +2079,7 @@ mod tests { buffer: None, knot_registry: None, knot_gate: None, + settlements: None, }; let mk = |live: bool| -> HydrantFrame { parse_frame(json!({ @@ -2611,6 +2718,7 @@ mod tests { buffer: None, knot_registry: None, knot_gate: None, + settlements: None, }; let olaren = Did::new_static("did:plc:olaren").unwrap(); let apply = async |frame: HydrantFrame| { @@ -2736,6 +2844,7 @@ mod tests { warming_buffer: None, knot_registry: None, knot_gate: None, + settlements: None, } } @@ -2865,6 +2974,7 @@ mod tests { buffer: None, knot_registry: None, knot_gate: None, + settlements: None, }; let repo_frame = |id: u64, rkey: &str| { parse_frame(json!({ @@ -3700,6 +3810,7 @@ mod tests { warming_buffer: None, knot_registry: None, knot_gate: None, + settlements: None, }; let parallelism = 4usize; @@ -3763,4 +3874,168 @@ mod tests { ); assert_eq!(runtime.coverage.snapshot().last_cursor().raw(), 4); } + + fn settling_runtime(watch: &Arc) -> IngestRuntime { + IngestRuntime { + settlements: Some(watch.clone()), + ..fresh_runtime(CancellationToken::new()) + } + } + + async fn run_stages(frame: HydrantFrame, rt: &IngestRuntime) { + let staged = prep_stage(frame, rt.clone()).await; + let staged = claim_stage(staged, rt.clone()).await; + let staged = resolve_stage(staged, rt.clone()).await; + commit_stage(staged, rt.clone(), 1).await; + } + + fn published(watch: &Arc, uri: &str) -> Option { + match watch.watch(&AtUri::new_owned(uri).unwrap(), None) { + bobbin_edge_index::Watch::Settled(s) => Some(s.outcome), + bobbin_edge_index::Watch::Waiting(..) => None, + } + } + + fn star_frame(collection: &str, action: &str, body: serde_json::Value) -> HydrantFrame { + parse_frame(json!({ + "id": 1, + "type": "record", + "record": { + "live": true, + "did": "did:plc:nel", + "rev": fresh_tid().as_str(), + "collection": collection, + "rkey": "abcabcabcabcz", + "action": action, + "cid": VALID_CID, + "record": body, + } + })) + } + + const STAR_URI: &str = "at://did:plc:nel/sh.tangled.feed.star/abcabcabcabcz"; + + #[tokio::test] + async fn an_indexed_record_is_published_as_indexed() { + let watch = Arc::new(Settlements::new()); + let rt = settling_runtime(&watch); + let frame = star_frame( + "sh.tangled.feed.star", + "create", + json!({ + "$type": "sh.tangled.feed.star", + "createdAt": "2026-05-01T00:00:00Z", + "subject": "at://did:plc:olaren/sh.tangled.repo/reporkeyabcde", + }), + ); + run_stages(frame, &rt).await; + assert_eq!(published(&watch, STAR_URI), Some(RecordOutcome::Indexed)); + } + + #[tokio::test] + async fn a_deleted_record_is_published_as_deleted() { + let watch = Arc::new(Settlements::new()); + let rt = settling_runtime(&watch); + let frame = star_frame("sh.tangled.feed.star", "delete", json!(null)); + run_stages(frame, &rt).await; + assert_eq!(published(&watch, STAR_URI), Some(RecordOutcome::Deleted)); + } + + #[tokio::test] + async fn a_collection_bobbin_does_not_index_is_published_as_rejected() { + let watch = Arc::new(Settlements::new()); + let rt = settling_runtime(&watch); + let frame = star_frame( + "sh.tangled.nonsense", + "create", + json!({"$type": "sh.tangled.nonsense"}), + ); + run_stages(frame, &rt).await; + let got = published(&watch, "at://did:plc:nel/sh.tangled.nonsense/abcabcabcabcz") + .expect("published"); + assert_eq!(got.rejection(), Some(Rejection::UnknownCollection)); + assert_eq!(got.detail(), Some("sh.tangled.nonsense")); + } + + #[tokio::test] + async fn a_record_that_will_not_decode_is_published_as_rejected() { + let watch = Arc::new(Settlements::new()); + let rt = settling_runtime(&watch); + let frame = star_frame( + "sh.tangled.feed.star", + "create", + json!({"$type": "sh.tangled.feed.star", "subject": 17}), + ); + run_stages(frame, &rt).await; + let got = published(&watch, STAR_URI).expect("published"); + assert_eq!(got.rejection(), Some(Rejection::Undecodable)); + // the writer gets the bad value back and what was expected there + let detail = got.detail().unwrap_or_default(); + assert!(detail.contains("17"), "detail was {detail:?}"); + assert!(detail.contains("StarSubject"), "detail was {detail:?}"); + } + + // bobbin never sees these, so a writer waiting on one has to be told no by + // the endpoint rather than left hanging + #[tokio::test] + async fn a_foreign_collection_publishes_nothing() { + let watch = Arc::new(Settlements::new()); + let rt = settling_runtime(&watch); + let frame = star_frame( + "app.bsky.feed.post", + "create", + json!({"$type": "app.bsky.feed.post", "text": "hi"}), + ); + run_stages(frame, &rt).await; + assert_eq!( + published(&watch, "at://did:plc:nel/app.bsky.feed.post/abcabcabcabcz"), + None, + ); + } + + // a parked record is not visible yet, so it must not claim to be indexed + // until its deps land + #[tokio::test] + async fn a_parked_record_publishes_nothing_until_it_drains() { + let watch = Arc::new(Settlements::new()); + let buffer = Arc::new(WarmingBuffer::new(RuntimeHasher::default())); + let rt = IngestRuntime { + warming_buffer: Some(buffer.clone()), + ..settling_runtime(&watch) + }; + let frame = star_frame( + "sh.tangled.feed.star", + "create", + json!({ + "$type": "sh.tangled.feed.star", + "createdAt": "2026-05-01T00:00:00Z", + "subject": "at://did:plc:olaren/sh.tangled.repo/reporkeyabcde", + }), + ); + run_stages(frame, &rt).await; + assert_eq!(published(&watch, STAR_URI), None); + + let repo = parse_frame(json!({ + "id": 2, + "type": "record", + "record": { + "live": true, + "did": "did:plc:olaren", + "rev": fresh_tid().as_str(), + "collection": "sh.tangled.repo", + "rkey": "reporkeyabcde", + "action": "create", + "cid": VALID_CID, + "record": { + "$type": "sh.tangled.repo", + "createdAt": "2026-05-01T00:00:00Z", + "name": "hydrant", + "knot": "knot.example", + "owner": "did:plc:olaren", + } + } + })); + run_stages(repo, &rt).await; + assert_eq!(published(&watch, STAR_URI), Some(RecordOutcome::Indexed)); + } } diff --git a/bobbin/crates/resolver/src/legacy_upgrade.rs b/bobbin/crates/resolver/src/legacy_upgrade.rs index 948b9e74d..dc10f9255 100644 --- a/bobbin/crates/resolver/src/legacy_upgrade.rs +++ b/bobbin/crates/resolver/src/legacy_upgrade.rs @@ -177,7 +177,7 @@ pub async fn upgrade_wire_bytes>( } fn upgrade_failed>(nsid: &Nsid) -> ExtractError { - ExtractError::UnknownCollection(alloc::format!("{}: legacy upgrade failed", nsid.as_ref())) + ExtractError::UpgradeFailed(alloc::string::String::from(nsid.as_ref())) } pub async fn decode_canon_or_upgrade>( diff --git a/bobbin/crates/types/src/edges.rs b/bobbin/crates/types/src/edges.rs index 982b358b8..2323ecb9b 100644 --- a/bobbin/crates/types/src/edges.rs +++ b/bobbin/crates/types/src/edges.rs @@ -49,6 +49,8 @@ pub enum ExtractError { InvalidAtUri(#[from] AtStrError), #[error("unknown collection NSID: {0}")] UnknownCollection(alloc::string::String), + #[error("legacy record could not be upgraded: {0}")] + UpgradeFailed(alloc::string::String), } #[derive(Debug)] diff --git a/bobbin/crates/xrpc/Cargo.toml b/bobbin/crates/xrpc/Cargo.toml index 539f069c4..7021b0d0e 100644 --- a/bobbin/crates/xrpc/Cargo.toml +++ b/bobbin/crates/xrpc/Cargo.toml @@ -30,12 +30,14 @@ tower-http = { workspace = true, features = ["trace"] } tracing = { workspace = true } trusted-proxies = { workspace = true } url = { workspace = true } +tokio = { workspace = true, features = ["time"] } reqwest = { workspace = true } [dev-dependencies] +bobbin-ingest = { workspace = true } +tokio-util = { workspace = true } bobbin-runtime = { workspace = true } http = { workspace = true } -tokio = { workspace = true, features = ["macros", "rt-multi-thread"] } -trusted-proxies = { workspace = true } +tokio = { workspace = true, features = ["macros", "rt-multi-thread", "test-util"] } url = { workspace = true } wiremock = { workspace = true } diff --git a/bobbin/crates/xrpc/src/backpressure.rs b/bobbin/crates/xrpc/src/backpressure.rs index f34c2126f..3cd2d42b9 100644 --- a/bobbin/crates/xrpc/src/backpressure.rs +++ b/bobbin/crates/xrpc/src/backpressure.rs @@ -2,6 +2,7 @@ use std::sync::Arc; use std::sync::atomic::{AtomicUsize, Ordering}; use bobbin_runtime::MemoryBudget; +use tokio::sync::{OwnedSemaphorePermit, Semaphore}; use crate::XrpcError; @@ -121,6 +122,51 @@ impl HeavyLimiter { } } +// a wait is idle, so this caps connection buildup rather than memory +#[derive(Clone, Copy, Debug, Eq, PartialEq)] +pub struct MaxAwaiting(usize); + +impl MaxAwaiting { + pub const fn new(limit: usize) -> Self { + Self(limit) + } + + pub const fn get(self) -> usize { + self.0 + } +} + +impl Default for MaxAwaiting { + fn default() -> Self { + Self(1024) + } +} + +// a wait is idle and must not spend the hydration budget, so it gets its own cap +pub struct AwaitLimiter { + sem: Arc, +} + +pub struct AwaitPermit { + _permit: OwnedSemaphorePermit, +} + +impl AwaitLimiter { + pub fn new(max: MaxAwaiting) -> Self { + Self { + sem: Arc::new(Semaphore::new(max.get())), + } + } + + pub fn try_enter(&self) -> Result { + self.sem + .clone() + .try_acquire_owned() + .map(|_permit| AwaitPermit { _permit }) + .map_err(|_| XrpcError::overloaded()) + } +} + #[cfg(test)] mod tests { use super::*; @@ -168,4 +214,13 @@ mod tests { (0..20).for_each(|_| limiter.adjust(PressureVerdict::Relieve)); assert_eq!(limiter.limit(), 8); } + + #[test] + fn await_limiter_limits_concurrency() { + let limiter = AwaitLimiter::new(MaxAwaiting::new(1)); + let held = limiter.try_enter().expect("first permit enters"); + assert!(matches!(limiter.try_enter(), Err(XrpcError::Overloaded))); + drop(held); + assert!(limiter.try_enter().is_ok()); + } } diff --git a/bobbin/crates/xrpc/src/lib.rs b/bobbin/crates/xrpc/src/lib.rs index 9c6964ef3..2c42c41a4 100644 --- a/bobbin/crates/xrpc/src/lib.rs +++ b/bobbin/crates/xrpc/src/lib.rs @@ -24,8 +24,8 @@ use axum::{ }; use bobbin_edge_index::{ Coverage, CoverageWatch, CursorParseError, EdgeItem, EdgePage, EdgeStore, IssueStateKind, - PageCursor, PageLimit, PageOffset, PageStart, PageToken, PullStatusKind, SortDir, StateIndex, - StateKind, TotalCount, + PageCursor, PageLimit, PageOffset, PageStart, PageToken, ParsedCid, PullStatusKind, Rejection, + Settlement, Settlements, SortDir, StateIndex, StateKind, TotalCount, Watch, }; use bobbin_knot_proxy::{ KnotHost, KnotProxy, KnotProxyError, MirrorNsid, MirrorProxy, ProxyResponse, RepoSlug, @@ -112,7 +112,8 @@ mod repo; mod view; pub use backpressure::{ - HeavyLimiter, HeavyPermit, MaxInFlight, PerRequestAnonBytes, PressureVerdict, ReservedFloor, + AwaitLimiter, AwaitPermit, HeavyLimiter, HeavyPermit, MaxAwaiting, MaxInFlight, + PerRequestAnonBytes, PressureVerdict, ReservedFloor, }; use client_address::X_FORWARDED_FOR; pub use client_address::{ClientAddress, SocketPeer}; @@ -121,6 +122,7 @@ use trusted_proxies::TrustedProxies; const DEFAULT_LIMIT: u32 = 50; const FETCH_CONCURRENCY: usize = 8; +const TANGLED_NSID_PREFIX: &str = "sh.tangled."; pub type Directory = JacquardResolver; @@ -144,7 +146,9 @@ pub struct AppState { pub identity: Arc, pub directory: Arc, pub limiter: Option>, + pub awaiting: Arc, pub client_address: Arc, + pub settlements: Arc, enrich_router: Arc>, } @@ -182,7 +186,9 @@ impl AppState { identity, directory, limiter: None, + awaiting: Arc::new(AwaitLimiter::new(MaxAwaiting::default())), client_address: Arc::new(ClientAddress::default()), + settlements: Arc::new(Settlements::new()), enrich_router: Arc::new(std::sync::OnceLock::new()), } } @@ -192,6 +198,11 @@ impl AppState { self } + pub fn with_max_awaiting(mut self, max: MaxAwaiting) -> Self { + self.awaiting = Arc::new(AwaitLimiter::new(max)); + self + } + pub fn with_mirror(mut self, mirror: Option>) -> Self { self.mirror = mirror; self @@ -212,6 +223,11 @@ impl AppState { self } + pub fn with_settlements(mut self, settlements: Arc) -> Self { + self.settlements = settlements; + self + } + /// for internal xrpc dispatch [`enrich`] pub fn self_router(&self) -> Router { self.enrich_router @@ -472,6 +488,7 @@ pub fn router(state: AppState) -> Router { axum::routing::post(enrich::enrich), ) .route("/xrpc/sh.tangled.bobbin.getCoverage", get(get_coverage)) + .route("/xrpc/sh.tangled.bobbin.awaitRecord", get(await_record)) .route( "/xrpc/blue.microcosm.identity.resolveMiniDoc", get(resolve_mini_doc), @@ -813,6 +830,8 @@ pub enum XrpcError { Internal(String), #[error("overloaded, shedding under memory pressure")] Overloaded, + #[error("record has not been ingested yet")] + NotSettled, } impl XrpcError { @@ -837,6 +856,7 @@ impl IntoResponse for XrpcError { Self::InvalidRecord(_) => (StatusCode::BAD_GATEWAY, "InvalidRecord"), Self::Internal(_) => (StatusCode::INTERNAL_SERVER_ERROR, "InternalError"), Self::Overloaded => (StatusCode::SERVICE_UNAVAILABLE, "Overloaded"), + Self::NotSettled => (StatusCode::GATEWAY_TIMEOUT, "NotSettled"), }; let body = ErrorBody { error, @@ -1695,7 +1715,9 @@ fn drop_unhydratable( Ok(None) } }, - Err(err @ (XrpcError::Internal(_) | XrpcError::Overloaded)) => Err(err), + Err(err @ (XrpcError::Internal(_) | XrpcError::Overloaded | XrpcError::NotSettled)) => { + Err(err) + } } } @@ -2703,6 +2725,75 @@ async fn get_coverage(State(state): State) -> Json { Json(state.coverage.snapshot().into()) } +// enough for a write to travel pds -> relay -> indexer -> bobbin, plus slack +// for a slow hop +const AWAIT_RECORD_TIMEOUT: Duration = Duration::from_secs(5); + +#[derive(Debug, Deserialize)] +struct AwaitRecordQuery { + uri: AtUri, + cid: Option, +} + +#[derive(Serialize)] +#[serde(rename_all = "camelCase")] +struct AwaitRecordResponse { + status: &'static str, + #[serde(skip_serializing_if = "Option::is_none")] + cid: Option, + #[serde(skip_serializing_if = "Option::is_none")] + reason: Option<&'static str>, + #[serde(skip_serializing_if = "Option::is_none")] + detail: Option>, +} + +impl From for AwaitRecordResponse { + fn from(s: Settlement) -> Self { + Self { + status: s.outcome.as_str(), + cid: s.cid, + reason: s.outcome.rejection().map(Rejection::as_str), + detail: s.outcome.detail().map(Box::from), + } + } +} + +async fn await_record( + State(state): State, + XrpcQuery(q): XrpcQuery, +) -> Result, XrpcError> { + let Some((collection, _)) = q.uri.collection().zip(q.uri.rkey()) else { + return Err(XrpcError::InvalidParams( + "uri must be at:////".into(), + )); + }; + if !collection.as_str().starts_with(TANGLED_NSID_PREFIX) { + return Err(XrpcError::InvalidParams(format!( + "bobbin never sees {collection} records" + ))); + } + // records are only ever indexed under their did, a handle would wait forever + if source_authority_did(&q.uri).is_none() { + return Err(XrpcError::InvalidParams( + "uri authority must be a did, not a handle".into(), + )); + } + // Drop releases the await slot when the handler returns + let _awaiting = state.awaiting.try_enter()?; + let settled = match state.settlements.watch(&q.uri, q.cid) { + Watch::Settled(s) => s, + // Drop removes the waiter from the list when the handler returns, + // so a timeout won't leave a channel no one is waiting on + Watch::Waiting(_registration, rx) => { + match tokio::time::timeout(AWAIT_RECORD_TIMEOUT, rx).await { + Ok(Ok(s)) => s, + Ok(Err(_)) | Err(_) => return Err(XrpcError::NotSettled), + } + } + }; + Ok(Json(settled.into())) +} + async fn search_query( State(state): State, XrpcQuery(q): XrpcQuery, diff --git a/bobbin/crates/xrpc/tests/aggregation.rs b/bobbin/crates/xrpc/tests/aggregation.rs index 67fc453a9..3af207e5b 100644 --- a/bobbin/crates/xrpc/tests/aggregation.rs +++ b/bobbin/crates/xrpc/tests/aggregation.rs @@ -813,10 +813,7 @@ async fn list_filtered_paginates_via_offset_with_total() { .unwrap(), ) .await; - assert_eq!( - page["total"], - json!(2), - ); + assert_eq!(page["total"], json!(2),); closed_uris.extend( page["items"] .as_array() diff --git a/bobbin/crates/xrpc/tests/await_record.rs b/bobbin/crates/xrpc/tests/await_record.rs new file mode 100644 index 000000000..b55e9909e --- /dev/null +++ b/bobbin/crates/xrpc/tests/await_record.rs @@ -0,0 +1,286 @@ +use std::sync::Arc; +use std::time::Duration; + +use axum::body::{Body, to_bytes}; +use bobbin_edge_index::{ + CoverageWatch, EdgeStore, RecordOutcome, Rejection, Settlement, Settlements, StateIndex, +}; +use bobbin_knot_proxy::{KnotHttpConfig, KnotProxy, KnotProxyConfig}; +use bobbin_record_lru::{CacheCapacity, LruRecordStore}; +use bobbin_resolver::RepoIdResolver; +use bobbin_runtime::{RuntimeHasher, SystemClock}; +use bobbin_search::{DEFAULT_WRITER_HEAP_BYTES, SearchIndex, SearchReader}; +use bobbin_slingshot_client::SlingshotClient; +use bobbin_xrpc::{AppState, MaxAwaiting, router}; +use http::{Request, StatusCode}; +use jacquard_common::DefaultStr; +use jacquard_common::types::string::{AtUri, Cid}; +use serde_json::{Value, json}; +use tower::ServiceExt; +use url::Url; +use wiremock::MockServer; + +const STAR: &str = "at://did:plc:alice/sh.tangled.feed.star/3l"; + +const CID: &str = "bafyreieqygohnz2zqyvtvktbjpvhutphobcmbsnt4q5lc36ri7vpcmoz4i"; +const CID2: &str = "bafkreigh2akiscaildc7gnvtklbsfhdgwz72eolmpckbqr5ej26byp3uli"; + +struct Harness { + settlements: Arc, + state: AppState, +} + +impl Harness { + async fn new() -> Self { + Self::with_max_awaiting(MaxAwaiting::default()).await + } + + async fn with_max_awaiting(max: MaxAwaiting) -> Self { + let server = MockServer::start().await; + let settlements = Arc::new(Settlements::new()); + let state = AppState::new( + Arc::new(LruRecordStore::new(CacheCapacity::from_bytes(64 * 1024))), + SlingshotClient::with_default_http(Url::parse(&server.uri()).unwrap()).unwrap(), + Arc::new(EdgeStore::new(RuntimeHasher::default())), + Arc::new(StateIndex::new(RuntimeHasher::default())), + Arc::new(StateIndex::new(RuntimeHasher::default())), + Arc::new(CoverageWatch::new()), + Arc::new( + KnotProxy::new( + KnotProxyConfig::default(), + KnotHttpConfig::default(), + Arc::new(SystemClock::new()), + RuntimeHasher::default(), + ) + .unwrap(), + ), + Arc::new( + SearchIndex::new(DEFAULT_WRITER_HEAP_BYTES, Arc::new(SystemClock::new())).unwrap(), + ) as Arc, + Arc::new(RepoIdResolver::detached(RuntimeHasher::default())), + Arc::new(bobbin_xrpc::default_directory()), + ) + .with_settlements(settlements.clone()) + .with_max_awaiting(max); + Self { settlements, state } + } + + async fn ask(&self, uri: &str, cid: Option<&str>) -> (StatusCode, Value) { + let resp = router(self.state.clone()) + .oneshot(request(uri, cid)) + .await + .unwrap(); + json_response(resp).await + } +} + +fn request(uri: &str, cid: Option<&str>) -> Request { + let query = match cid { + Some(cid) => format!("uri={}&cid={cid}", urlencoding(uri)), + None => format!("uri={}", urlencoding(uri)), + }; + Request::builder() + .uri(format!("/xrpc/sh.tangled.bobbin.awaitRecord?{query}")) + .body(Body::empty()) + .unwrap() +} + +fn urlencoding(raw: &str) -> String { + raw.replace(':', "%3A").replace('/', "%2F") +} + +async fn json_response(resp: axum::response::Response) -> (StatusCode, Value) { + let status = resp.status(); + let bytes = to_bytes(resp.into_body(), 1 << 20).await.unwrap(); + let parsed: Value = serde_json::from_slice(&bytes).expect("JSON body"); + (status, parsed) +} + +fn settlement(uri: &str, cid: Option<&str>, outcome: RecordOutcome) -> Settlement { + Settlement { + source: AtUri::::new_owned(uri).unwrap(), + cid: cid.map(|c| c.parse().unwrap()), + outcome, + } +} + +// the paused clock only moves when something parks, so a zero elapsed proves the +// recent cache answered without waiting at all +#[tokio::test(start_paused = true)] +async fn reports_a_record_that_already_landed() { + let h = Harness::new().await; + h.settlements + .settle(settlement(STAR, Some(CID), RecordOutcome::Indexed)); + let start = tokio::time::Instant::now(); + let (status, body) = h.ask(STAR, Some(CID)).await; + assert_eq!(status, StatusCode::OK); + assert_eq!(body, json!({ "status": "indexed", "cid": CID })); + assert_eq!(start.elapsed(), Duration::ZERO); +} + +#[tokio::test(start_paused = true)] +async fn reports_why_a_record_was_rejected() { + let h = Harness::new().await; + h.settlements.settle(settlement( + STAR, + Some(CID), + RecordOutcome::rejected(Rejection::Undecodable, None), + )); + let (status, body) = h.ask(STAR, Some(CID)).await; + assert_eq!(status, StatusCode::OK); + assert_eq!( + body, + json!({ "status": "rejected", "cid": CID, "reason": "undecodable" }) + ); +} + +#[tokio::test(start_paused = true)] +async fn a_rejection_carries_its_underlying_error() { + let h = Harness::new().await; + h.settlements.settle(settlement( + STAR, + Some(CID), + RecordOutcome::rejected(Rejection::Undecodable, Some("invalid type: integer `17`")), + )); + let (_, body) = h.ask(STAR, Some(CID)).await; + assert_eq!( + body, + json!({ + "status": "rejected", + "cid": CID, + "reason": "undecodable", + "detail": "invalid type: integer `17`", + }) + ); +} + +#[tokio::test(start_paused = true)] +async fn a_delete_reports_no_cid() { + let h = Harness::new().await; + h.settlements + .settle(settlement(STAR, None, RecordOutcome::Deleted)); + let (status, body) = h.ask(STAR, None).await; + assert_eq!(status, StatusCode::OK); + assert_eq!(body, json!({ "status": "deleted" })); +} + +// a delete has no cid anywhere in the protocol, so asking for one cannot be answered +#[tokio::test(start_paused = true)] +async fn a_cid_never_matches_a_delete() { + let h = Harness::new().await; + h.settlements + .settle(settlement(STAR, None, RecordOutcome::Deleted)); + let (status, _) = h.ask(STAR, Some(CID)).await; + assert_eq!(status, StatusCode::GATEWAY_TIMEOUT); +} + +// join! polls the request first, so it is parked on the waiter by the time the +// settlement lands +#[tokio::test(start_paused = true)] +async fn answers_a_record_that_lands_mid_request() { + let h = Harness::new().await; + let settlements = h.settlements.clone(); + let (answer, ()) = tokio::join!(h.ask(STAR, Some(CID)), async { + settlements.settle(settlement(STAR, Some(CID), RecordOutcome::Indexed)); + }); + let (status, body) = answer; + assert_eq!(status, StatusCode::OK); + assert_eq!(body, json!({ "status": "indexed", "cid": CID })); +} + +#[tokio::test(start_paused = true)] +async fn a_record_that_never_lands_runs_out() { + let h = Harness::new().await; + let start = tokio::time::Instant::now(); + let (status, body) = h.ask(STAR, Some(CID)).await; + assert_eq!(status, StatusCode::GATEWAY_TIMEOUT); + assert_eq!(body["error"], "NotSettled"); + // the record never settles, so the wait lasts the whole timeout + assert_eq!(start.elapsed(), Duration::from_secs(5)); +} + +// a different revision of the same record is not an answer +#[tokio::test(start_paused = true)] +async fn another_revision_does_not_answer() { + let h = Harness::new().await; + h.settlements + .settle(settlement(STAR, Some(CID), RecordOutcome::Indexed)); + let (status, _) = h.ask(STAR, Some(CID2)).await; + assert_eq!(status, StatusCode::GATEWAY_TIMEOUT); +} + +#[tokio::test(start_paused = true)] +async fn a_collection_bobbin_never_sees_fails_fast() { + let h = Harness::new().await; + let (status, body) = h + .ask("at://did:plc:alice/app.bsky.feed.post/3l", None) + .await; + assert_eq!(status, StatusCode::BAD_REQUEST); + assert_eq!(body["error"], "InvalidRequest"); +} + +#[tokio::test(start_paused = true)] +async fn a_uri_without_a_collection_fails_fast() { + let h = Harness::new().await; + let (status, body) = h.ask("at://did:plc:alice", None).await; + assert_eq!(status, StatusCode::BAD_REQUEST); + assert_eq!(body["error"], "InvalidRequest"); +} + +#[tokio::test(start_paused = true)] +async fn a_uri_without_an_rkey_fails_fast() { + let h = Harness::new().await; + let (status, body) = h.ask("at://did:plc:alice/sh.tangled.feed.star", None).await; + assert_eq!(status, StatusCode::BAD_REQUEST); + assert_eq!(body["error"], "InvalidRequest"); +} + +// the indexed uri is always did-based, so a handle would never match +#[tokio::test(start_paused = true)] +async fn a_handle_authority_fails_fast() { + let h = Harness::new().await; + let (status, body) = h + .ask("at://alice.example/sh.tangled.feed.star/3l", None) + .await; + assert_eq!(status, StatusCode::BAD_REQUEST); + assert_eq!(body["error"], "InvalidRequest"); +} + +// paused time turns a slow-path wait into a 504, so the 400 proves this never parked +#[tokio::test(start_paused = true)] +async fn a_bad_cid_does_not_burn_the_wait() { + let h = Harness::new().await; + let (status, body) = h.ask(STAR, Some("notacid")).await; + assert_eq!(status, StatusCode::BAD_REQUEST); + assert_eq!(body["error"], "InvalidRequest"); + assert!( + body["message"] + .as_str() + .unwrap_or_default() + .contains("cid:"), + "message should name the bad field, got {body}" + ); +} + +// join! polls the first ask onto the permit before the second one runs, so the +// overload lands without waiting on a clock +#[tokio::test(start_paused = true)] +async fn respects_concurrency_limit() { + let h = Harness::with_max_awaiting(MaxAwaiting::new(1)).await; + + let ((status1, body1), (status2, body2)) = tokio::join!(h.ask(STAR, Some(CID)), async { + let overloaded = h.ask(STAR, Some(CID2)).await; + h.settlements + .settle(settlement(STAR, Some(CID), RecordOutcome::Indexed)); + overloaded + }); + assert_eq!(status2, StatusCode::SERVICE_UNAVAILABLE); + assert_eq!(body2["error"], "Overloaded"); + assert_eq!(status1, StatusCode::OK); + assert_eq!(body1, json!({ "status": "indexed", "cid": CID })); + + // the first ask handed its permit back on the way out + let (status3, body3) = h.ask(STAR, Some(CID2)).await; + assert_eq!(status3, StatusCode::GATEWAY_TIMEOUT); + assert_eq!(body3["error"], "NotSettled"); +} diff --git a/bobbin/crates/xrpc/tests/await_record_e2e.rs b/bobbin/crates/xrpc/tests/await_record_e2e.rs new file mode 100644 index 000000000..60033b704 --- /dev/null +++ b/bobbin/crates/xrpc/tests/await_record_e2e.rs @@ -0,0 +1,162 @@ +//! a record arriving over the hydrant stream has to answer an awaitRecord call +//! through the same Settlements bobbin wires up in main + +use std::sync::Arc; +use std::time::Duration; + +use axum::body::{Body, to_bytes}; +use bobbin_edge_index::{CoverageWatch, EdgeStore, HydrantCursor, Settlements, StateIndex}; +use bobbin_ingest::{IngestConfig, IngestRuntime, RepoIdResolver, run as run_ingest}; +use bobbin_knot_proxy::{KnotHttpConfig, KnotProxy, KnotProxyConfig}; +use bobbin_record_lru::{CacheCapacity, LruRecordStore, RecordStore}; +use bobbin_resolver::IdentityResolver; +use bobbin_runtime::{ + MemWsResponder, MemWsServerFuture, MemWsTransport, OsEntropy, RuntimeHasher, SystemClock, + WsMessage, +}; +use bobbin_search::{DEFAULT_WRITER_HEAP_BYTES, SearchIndex, SearchReader}; +use bobbin_slingshot_client::SlingshotClient; +use bobbin_types::search::NoopSearchSink; +use bobbin_xrpc::{AppState, router}; +use http::{Request, StatusCode}; +use serde_json::{Value, json}; +use tokio::sync::mpsc; +use tokio_util::sync::CancellationToken; +use tower::ServiceExt; +use url::Url; +use wiremock::MockServer; + +const CID: &str = "bafyreieqygohnz2zqyvtvktbjpvhutphobcmbsnt4q5lc36ri7vpcmoz4i"; +const STAR: &str = "at://did:plc:woof/sh.tangled.feed.star/abcabcabcabcz"; + +struct Hydrant { + frames: Vec, +} + +impl MemWsResponder for Hydrant { + fn spawn_server( + &self, + _url: Url, + mut recv: mpsc::UnboundedReceiver, + send: mpsc::Sender, + ) -> MemWsServerFuture { + let frames = self.frames.clone(); + Box::pin(async move { + for frame in frames { + let _ = send.send(WsMessage::Text(frame)).await; + } + while recv.recv().await.is_some() {} + }) + } +} + +fn star_frame() -> String { + json!({ + "id": 1, + "type": "record", + "record": { + "live": true, + "did": "did:plc:woof", + "rev": "3lzzzzzzzzz2a", + "collection": "sh.tangled.feed.star", + "rkey": "abcabcabcabcz", + "action": "create", + "cid": CID, + "record": { + "$type": "sh.tangled.feed.star", + "createdAt": "2026-05-01T00:00:00Z", + "subject": "at://did:plc:meow/sh.tangled.repo/reporkeyabcde", + } + } + }) + .to_string() +} + +#[tokio::test] +async fn a_record_off_the_stream_answers_an_await() { + let slingshot_server = MockServer::start().await; + let slingshot = + SlingshotClient::with_default_http(Url::parse(&slingshot_server.uri()).unwrap()).unwrap(); + let settlements = Arc::new(Settlements::new()); + let records = Arc::new(LruRecordStore::new(CacheCapacity::from_bytes(64 * 1024))); + let edges = Arc::new(EdgeStore::new(RuntimeHasher::default())); + let issue_states = Arc::new(StateIndex::new(RuntimeHasher::default())); + let pull_statuses = Arc::new(StateIndex::new(RuntimeHasher::default())); + let coverage = Arc::new(CoverageWatch::new()); + let resolver = Arc::new(RepoIdResolver::detached(RuntimeHasher::default())); + let cancel = CancellationToken::new(); + + let ingest: IngestRuntime = IngestRuntime { + store: edges.clone(), + issue_states: issue_states.clone(), + pull_statuses: pull_statuses.clone(), + coverage: coverage.clone(), + search: Arc::new(NoopSearchSink), + records: records.clone() as Arc, + resolver: resolver.clone(), + identity: Arc::new(IdentityResolver::detached(RuntimeHasher::default())), + clock: Arc::new(SystemClock::new()), + entropy: Arc::new(OsEntropy), + ws: MemWsTransport::shared(Arc::new(Hydrant { + frames: vec![star_frame()], + })), + cancel: cancel.clone(), + disconnects: None, + warming_shadow: None, + warming_buffer: None, + knot_registry: None, + knot_gate: None, + settlements: Some(settlements.clone()), + }; + + let state = AppState::new( + records as Arc, + slingshot, + edges, + issue_states, + pull_statuses, + coverage, + Arc::new( + KnotProxy::new( + KnotProxyConfig::default(), + KnotHttpConfig::default(), + Arc::new(SystemClock::new()), + RuntimeHasher::default(), + ) + .unwrap(), + ), + Arc::new(SearchIndex::new(DEFAULT_WRITER_HEAP_BYTES, Arc::new(SystemClock::new())).unwrap()) + as Arc, + resolver, + Arc::new(bobbin_xrpc::default_directory()), + ) + .with_settlements(settlements); + + let request = Request::builder() + .uri(format!( + "/xrpc/sh.tangled.bobbin.awaitRecord?uri={}&cid={CID}", + STAR.replace(':', "%3A").replace('/', "%2F"), + )) + .body(Body::empty()) + .unwrap(); + + let config = IngestConfig { + hydrant_base: Url::parse("ws://hydrant.invalid").unwrap(), + start_cursor: HydrantCursor::new(0), + parallelism: bobbin_ingest::DEFAULT_INGEST_PARALLELISM, + }; + let ingesting = tokio::spawn(run_ingest(config, ingest)); + let answer = tokio::time::timeout(Duration::from_secs(10), router(state).oneshot(request)) + .await + .expect("the endpoint must answer well inside its own wait"); + + cancel.cancel(); + let _ = tokio::time::timeout(Duration::from_secs(5), ingesting).await; + + let resp = answer.unwrap(); + let status = resp.status(); + let bytes = to_bytes(resp.into_body(), 1 << 20).await.unwrap(); + let body: Value = serde_json::from_slice(&bytes).unwrap(); + assert_eq!(status, StatusCode::OK, "body was {body}"); + assert_eq!(body, json!({ "status": "indexed", "cid": CID })); +} diff --git a/bobbin/crates/xrpc/tests/search.rs b/bobbin/crates/xrpc/tests/search.rs index 84487112b..019e44224 100644 --- a/bobbin/crates/xrpc/tests/search.rs +++ b/bobbin/crates/xrpc/tests/search.rs @@ -414,13 +414,8 @@ async fn pagination_offset_matches_cursor_walk() { if let Some(c) = &cursor { params.push(("cursor", c)); } - let (_, page) = json_response( - app.clone() - .oneshot(search_request(¶ms)) - .await - .unwrap(), - ) - .await; + let (_, page) = + json_response(app.clone().oneshot(search_request(¶ms)).await.unwrap()).await; walked.extend( page["hits"] .as_array() @@ -437,7 +432,11 @@ async fn pagination_offset_matches_cursor_walk() { let (_, page) = json_response( app.clone() - .oneshot(search_request(&[("q", "anemone"), ("limit", "2"), ("offset", "2")])) + .oneshot(search_request(&[ + ("q", "anemone"), + ("limit", "2"), + ("offset", "2"), + ])) .await .unwrap(), ) diff --git a/lexicons/bobbin/awaitRecord.json b/lexicons/bobbin/awaitRecord.json new file mode 100644 index 000000000..df026b04e --- /dev/null +++ b/lexicons/bobbin/awaitRecord.json @@ -0,0 +1,70 @@ +{ + "lexicon": 1, + "id": "sh.tangled.bobbin.awaitRecord", + "defs": { + "main": { + "type": "query", + "description": "Wait for a record written to a PDS to reach this bobbin instance, then report what happened to it.", + "parameters": { + "type": "params", + "required": ["uri"], + "properties": { + "uri": { + "type": "string", + "format": "at-uri", + "description": "AT-URI of the record that was written." + }, + "cid": { + "type": "string", + "format": "cid", + "description": "CID of the revision that was written, as returned by createRecord, putRecord or applyWrites. Omit it only when awaiting a delete, which has no CID. any other responds with whatever revision of the record settles first." + } + } + }, + "output": { + "encoding": "application/json", + "schema": { + "type": "object", + "required": ["status"], + "properties": { + "status": { + "type": "string", + "knownValues": ["indexed", "deleted", "rejected"], + "description": "What bobbin did with the record." + }, + "cid": { + "type": "string", + "format": "cid", + "description": "CID of the revision this result is about." + }, + "reason": { + "type": "string", + "knownValues": [ + "unknownCollection", + "undecodable", + "missingBody", + "notIndexed" + ], + "description": "Why a rejected record will not appear." + }, + "detail": { + "type": "string", + "maxLength": 200, + "description": "Opaque error message which can be more detailed than the reason." + } + } + } + }, + "errors": [ + { + "name": "NotSettled", + "description": "The record had not reached the appview before the wait ran out" + }, + { + "name": "Overloaded", + "description": "Too many callers are already waiting, retry later or poll the record instead" + } + ] + } + } +} diff --git a/web/lex.config.ts b/web/lex.config.ts index 68057e7d2..07045d357 100644 --- a/web/lex.config.ts +++ b/web/lex.config.ts @@ -9,6 +9,7 @@ export default defineLexiconConfig({ imports: ["@atcute/atproto"], files: [ "../lexicons/*.json", + "../lexicons/bobbin/**/*.json", "../lexicons/actor/**/*.json", "../lexicons/ci/**/*.json", "../lexicons/feed/**/*.json", diff --git a/web/src/lib/api/lexicons/index.ts b/web/src/lib/api/lexicons/index.ts index 5e2da77d8..e0457359f 100644 --- a/web/src/lib/api/lexicons/index.ts +++ b/web/src/lib/api/lexicons/index.ts @@ -2,6 +2,7 @@ export * as ShTangledActorDefs from "./types/sh/tangled/actor/defs.js"; export * as ShTangledActorGetProfile from "./types/sh/tangled/actor/getProfile.js"; export * as ShTangledActorGetProfiles from "./types/sh/tangled/actor/getProfiles.js"; export * as ShTangledActorProfile from "./types/sh/tangled/actor/profile.js"; +export * as ShTangledBobbinAwaitRecord from "./types/sh/tangled/bobbin/awaitRecord.js"; export * as ShTangledCiCancelPipeline from "./types/sh/tangled/ci/cancelPipeline.js"; export * as ShTangledCiGetPipeline from "./types/sh/tangled/ci/getPipeline.js"; export * as ShTangledCiPipeline from "./types/sh/tangled/ci/pipeline.js"; diff --git a/web/src/lib/api/lexicons/types/sh/tangled/bobbin/awaitRecord.ts b/web/src/lib/api/lexicons/types/sh/tangled/bobbin/awaitRecord.ts new file mode 100644 index 000000000..242c977ad --- /dev/null +++ b/web/src/lib/api/lexicons/types/sh/tangled/bobbin/awaitRecord.ts @@ -0,0 +1,67 @@ +import type {} from "@atcute/lexicons"; +import * as v from "@atcute/lexicons/validations"; +import type {} from "@atcute/lexicons/ambient"; + +const _mainSchema = /*#__PURE__*/ v.query("sh.tangled.bobbin.awaitRecord", { + params: /*#__PURE__*/ v.object({ + /** + * CID of the revision that was written, as returned by createRecord, putRecord or applyWrites. Omit it only when awaiting a delete, which has no CID; any other omission answers with whatever revision of the record settles first. + */ + cid: /*#__PURE__*/ v.optional(/*#__PURE__*/ v.cidString()), + /** + * AT-URI of the record that was written. + */ + uri: /*#__PURE__*/ v.resourceUriString(), + }), + output: { + type: "lex", + schema: /*#__PURE__*/ v.object({ + /** + * CID of the revision this result is about. + */ + cid: /*#__PURE__*/ v.optional(/*#__PURE__*/ v.cidString()), + /** + * Opaque error message which can be more detailed than the reason. + * @maxLength 200 + */ + detail: /*#__PURE__*/ v.optional( + /*#__PURE__*/ v.constrain(/*#__PURE__*/ v.string(), [ + /*#__PURE__*/ v.stringLength(0, 200), + ]), + ), + /** + * Why a rejected record will not appear. + */ + reason: /*#__PURE__*/ v.optional( + /*#__PURE__*/ v.string< + | "missingBody" + | "notIndexed" + | "undecodable" + | "unknownCollection" + | (string & {}) + >(), + ), + /** + * What bobbin did with the record. + */ + status: /*#__PURE__*/ v.string< + "deleted" | "indexed" | "rejected" | (string & {}) + >(), + }), + }, +}); + +type main$schematype = typeof _mainSchema; + +export interface mainSchema extends main$schematype {} + +export const mainSchema = _mainSchema as mainSchema; + +export interface $params extends v.InferInput {} +export interface $output extends v.InferXRPCBodyInput {} + +declare module "@atcute/lexicons/ambient" { + interface XRPCQueries { + "sh.tangled.bobbin.awaitRecord": mainSchema; + } +}