diff --git a/crates/ingest/examples/smoke.rs b/crates/ingest/examples/smoke.rs index f1a78a7..a5beeeb 100644 --- a/crates/ingest/examples/smoke.rs +++ b/crates/ingest/examples/smoke.rs @@ -2,10 +2,11 @@ use std::sync::Arc; use std::time::Duration; use bobbin_edge_index::{CoverageWatch, EdgeStore}; -use bobbin_ingest::{IngestConfig, run}; +use bobbin_ingest::{IngestConfig, IngestRuntime, RepoIdResolver, run}; use bobbin_record_lru::{NoopRecordStore, RecordStore}; use bobbin_types::search::NoopSearchSink; use futures::stream::{self, StreamExt}; +use tokio_util::sync::CancellationToken; use url::Url; const MIN_EVENTS_PER_RUN: u64 = 100; @@ -27,12 +28,17 @@ async fn main() { let coverage = Arc::new(CoverageWatch::new()); let cfg = IngestConfig::new(url); - let store_runner = store.clone(); - let coverage_runner = coverage.clone(); - let search = Arc::new(NoopSearchSink); - let records: Arc = Arc::new(NoopRecordStore); + let cancel = CancellationToken::new(); + let runtime = IngestRuntime { + store: store.clone(), + coverage: coverage.clone(), + search: Arc::new(NoopSearchSink), + records: Arc::new(NoopRecordStore) as Arc, + resolver: Arc::new(RepoIdResolver::detached()), + cancel: cancel.clone(), + }; let task = tokio::spawn(async move { - let _ = run(cfg, store_runner, coverage_runner, search, records).await; + let _ = run(cfg, runtime).await; }); stream::iter(0..ticks) @@ -51,7 +57,8 @@ async fn main() { }) .await; - task.abort(); + cancel.cancel(); + let _ = task.await; let final_snap = coverage.snapshot(); let events = final_snap.events_processed(); diff --git a/crates/ingest/src/frame.rs b/crates/ingest/src/frame.rs index 87d564c..6fb6e6d 100644 --- a/crates/ingest/src/frame.rs +++ b/crates/ingest/src/frame.rs @@ -2,6 +2,7 @@ use jacquard_common::DefaultStr; use jacquard_common::types::did::Did; use jacquard_common::types::nsid::Nsid; use jacquard_common::types::recordkey::Rkey; +use jacquard_common::types::string::Cid; use jacquard_common::types::tid::Tid; use serde::Deserialize; @@ -47,6 +48,8 @@ pub struct RecordFrame { pub action: RecordAction, #[serde(default)] pub record: Option, + #[serde(default)] + pub cid: Option>, } #[derive(Clone, Copy, Debug, Eq, PartialEq)] diff --git a/crates/ingest/src/lib.rs b/crates/ingest/src/lib.rs index 6ba2f46..bc819a1 100644 --- a/crates/ingest/src/lib.rs +++ b/crates/ingest/src/lib.rs @@ -3,19 +3,28 @@ use std::time::{Duration, SystemTime, UNIX_EPOCH}; use bobbin_edge_index::{Coverage, CoverageWatch, EdgeStore, HydrantCursor, PromotionSignal}; use bobbin_record_lru::RecordStore; -use bobbin_types::edges::{ExtractError, Record}; +use bobbin_types::edges::{Edge, ExtractError, Record}; +use bobbin_types::record::RecordBody; use bobbin_types::search::{SearchSink, SearchableRecord}; use futures::{SinkExt, StreamExt}; use jacquard_common::DefaultStr; -use jacquard_common::types::string::{AtStrError, AtUri}; +use jacquard_common::types::did::Did; +use jacquard_common::types::ident::AtIdentifier; +use jacquard_common::types::recordkey::Rkey; +use jacquard_common::types::string::{AtStrError, AtUri, Cid}; use thiserror::Error; use tokio::time::{Instant, MissedTickBehavior, interval}; -use tokio_tungstenite::tungstenite::{Bytes, Message}; +use tokio_tungstenite::tungstenite::{ + Bytes, Message, protocol::CloseFrame, protocol::frame::coding::CloseCode, +}; +use tokio_util::sync::CancellationToken; use tracing::{debug, info, warn}; use url::Url; mod frame; +mod resolver; pub use frame::{FrameKind, HydrantFrame, RecordAction, RecordFrame}; +pub use resolver::{RepoIdResolver, Resolution}; const TANGLED_PREFIX: &str = "sh.tangled."; const RECONNECT_INITIAL_DELAY: Duration = Duration::from_millis(500); @@ -24,6 +33,8 @@ const PING_INTERVAL: Duration = Duration::from_secs(20); const PONG_TIMEOUT: Duration = Duration::from_secs(15); const READY_SKEW: Duration = Duration::from_secs(60); const FRAME_CHANNEL_DEPTH: usize = 256; +const CONTROL_CHANNEL_DEPTH: usize = 16; +const NORMALIZE_CONCURRENCY: usize = 8; #[derive(Clone, Debug)] pub struct IngestConfig { @@ -89,25 +100,43 @@ struct SessionEnd { error: Option, } +pub struct IngestRuntime { + pub store: Arc, + pub coverage: Arc, + pub search: Arc, + pub records: Arc, + pub resolver: Arc, + pub cancel: CancellationToken, +} + +impl Clone for IngestRuntime { + fn clone(&self) -> Self { + Self { + store: self.store.clone(), + coverage: self.coverage.clone(), + search: self.search.clone(), + records: self.records.clone(), + resolver: self.resolver.clone(), + cancel: self.cancel.clone(), + } + } +} + pub async fn run( config: IngestConfig, - store: Arc, - coverage: Arc, - search: Arc, - records: Arc, + runtime: IngestRuntime, ) -> Result<(), IngestError> { let mut backoff = RECONNECT_INITIAL_DELAY; loop { - let cursor = next_connect_cursor(coverage.snapshot(), config.start_cursor); - let SessionEnd { outcome, error } = run_session( - &config, - cursor, - store.clone(), - coverage.clone(), - search.clone(), - records.clone(), - ) - .await; + let cursor = next_connect_cursor(runtime.coverage.snapshot(), config.start_cursor); + let SessionEnd { outcome, error } = run_session(&config, cursor, &runtime).await; + if runtime.cancel.is_cancelled() { + info!( + last_cursor = runtime.coverage.snapshot().last_cursor().raw(), + "ingest stopped after shutdown signal" + ); + return Ok(()); + } match (outcome, &error) { (SessionOutcome::Progressed, None) => { info!("hydrant stream closed after delivering frames, reconnecting") @@ -121,7 +150,10 @@ pub async fn run( if made_progress { backoff = RECONNECT_INITIAL_DELAY; } else { - tokio::time::sleep(jittered(backoff)).await; + tokio::select! { + _ = tokio::time::sleep(jittered(backoff)) => {} + _ = runtime.cancel.cancelled() => return Ok(()), + } backoff = (backoff * 2).min(RECONNECT_MAX_DELAY); } } @@ -147,10 +179,7 @@ fn jittered(base: Duration) -> Duration { async fn run_session( config: &IngestConfig, cursor: HydrantCursor, - store: Arc, - coverage: Arc, - search: Arc, - records: Arc, + runtime: &IngestRuntime, ) -> SessionEnd { let url = match config.stream_url(cursor) { Ok(u) => u, @@ -162,7 +191,14 @@ async fn run_session( } }; info!(%url, "connecting to hydrant /stream"); - let (mut ws, _resp) = match tokio_tungstenite::connect_async(url.as_str()).await { + let connect = tokio::select! { + biased; + _ = runtime.cancel.cancelled() => { + return SessionEnd { outcome: SessionOutcome::Empty, error: None }; + } + res = tokio_tungstenite::connect_async(url.as_str()) => res, + }; + let (ws, _resp) = match connect { Ok(pair) => pair, Err(e) => { return SessionEnd { @@ -171,36 +207,104 @@ async fn run_session( }; } }; - let mut outcome = SessionOutcome::Empty; - let mut pinger = interval(PING_INTERVAL); - pinger.set_missed_tick_behavior(MissedTickBehavior::Delay); - pinger.tick().await; - let mut pong_deadline: Option = None; + let (mut ws_sink, mut ws_stream) = ws.split(); let (frame_tx, mut frame_rx) = tokio::sync::mpsc::channel::(FRAME_CHANNEL_DEPTH); - let processor_store = store.clone(); - let processor_coverage = coverage.clone(); - let processor_search = search.clone(); - let processor_records = records.clone(); + let processor_runtime = runtime.clone(); let processor = tokio::spawn(async move { - while let Some(frame) = frame_rx.recv().await { - handle_frame( - frame, - &processor_store, - &processor_coverage, - &*processor_search, - &*processor_records, - ) - .await; + loop { + tokio::select! { + biased; + _ = processor_runtime.cancel.cancelled() => break, + next = frame_rx.recv() => { + let Some(frame) = next else { break }; + handle_frame( + frame, + &processor_runtime.store, + &processor_runtime.coverage, + &*processor_runtime.search, + &*processor_runtime.records, + &processor_runtime.resolver, + ) + .await; + } + } } }); - let error: Option = loop { + let (control_tx, mut control_rx) = tokio::sync::mpsc::channel::(CONTROL_CHANNEL_DEPTH); + let session_cancel = runtime.cancel.child_token(); + let reader_cancel = session_cancel.clone(); + let reader = tokio::spawn(async move { + let mut outcome = SessionOutcome::Empty; + let error: Option = loop { + tokio::select! { + biased; + _ = reader_cancel.cancelled() => break None, + msg = ws_stream.next() => { + let Some(msg) = msg else { break None; }; + let parsed = match msg { + Ok(m) => m, + Err(e) => break Some(IngestError::Transport(e)), + }; + match parsed { + Message::Text(text) => { + let frame: HydrantFrame = match serde_json::from_str(&text) { + Ok(f) => f, + Err(e) => break Some(IngestError::Decode(e)), + }; + tokio::select! { + biased; + _ = reader_cancel.cancelled() => break None, + res = frame_tx.send(frame) => { + if res.is_err() { break None; } + } + } + outcome = SessionOutcome::Progressed; + } + Message::Binary(_) => { + debug!("hydrant sent unexpected binary frame, ignoring"); + } + Message::Ping(payload) => { + if control_tx.send(WsEvent::IncomingPing(payload)).await.is_err() { + break None; + } + } + Message::Pong(_) => { + if control_tx.send(WsEvent::IncomingPong).await.is_err() { + break None; + } + } + Message::Frame(_) => {} + Message::Close(close) => { + debug!(?close, "hydrant closed stream"); + break None; + } + } + } + } + }; + SessionEnd { outcome, error } + }); + + let mut pinger = interval(PING_INTERVAL); + pinger.set_missed_tick_behavior(MissedTickBehavior::Delay); + pinger.tick().await; + let mut pong_deadline: Option = None; + + let writer_error: Option = loop { tokio::select! { biased; + _ = runtime.cancel.cancelled() => { + let close_frame = CloseFrame { code: CloseCode::Normal, reason: "bobbin shutdown".into() }; + if let Err(e) = ws_sink.send(Message::Close(Some(close_frame))).await { + debug!(?e, "could not send hydrant close frame on shutdown"); + } + break None; + } _ = pinger.tick() => { if pong_deadline.is_none() { - if let Err(e) = ws.send(Message::Ping(Bytes::new())).await { + if let Err(e) = ws_sink.send(Message::Ping(Bytes::new())).await { break Some(IngestError::Transport(e)); } pong_deadline = Some(Instant::now() + PONG_TIMEOUT); @@ -209,51 +313,43 @@ async fn run_session( _ = wait_until(pong_deadline) => { break Some(IngestError::PongTimeout(PONG_TIMEOUT)); } - msg = ws.next() => { - let Some(msg) = msg else { - break None; - }; - let parsed = match msg { - Ok(m) => m, - Err(e) => break Some(IngestError::Transport(e)), - }; - match parsed { - Message::Text(text) => { - let frame: HydrantFrame = match serde_json::from_str(&text) { - Ok(f) => f, - Err(e) => break Some(IngestError::Decode(e)), - }; - if frame_tx.send(frame).await.is_err() { - break None; - } - outcome = SessionOutcome::Progressed; - } - Message::Binary(_) => { - debug!("hydrant sent unexpected binary frame, ignoring"); - } - Message::Ping(payload) => { - if let Err(e) = ws.send(Message::Pong(payload)).await { + evt = control_rx.recv() => { + let Some(evt) = evt else { break None; }; + match evt { + WsEvent::IncomingPing(payload) => { + if let Err(e) = ws_sink.send(Message::Pong(payload)).await { break Some(IngestError::Transport(e)); } } - Message::Pong(_) => { + WsEvent::IncomingPong => { pong_deadline = None; } - Message::Frame(_) => {} - Message::Close(close) => { - debug!(?close, "hydrant closed stream"); - break None; - } } } } }; - drop(frame_tx); + session_cancel.cancel(); + drop(ws_sink); + drop(control_rx); + + let reader_end = reader.await.unwrap_or(SessionEnd { + outcome: SessionOutcome::Empty, + error: None, + }); if let Err(join) = processor.await { warn!(?join, "frame processor task panicked"); } - SessionEnd { outcome, error } + SessionEnd { + outcome: reader_end.outcome, + error: writer_error.or(reader_end.error), + } +} + +#[derive(Debug)] +enum WsEvent { + IncomingPing(Bytes), + IncomingPong, } async fn wait_until(deadline: Option) { @@ -269,6 +365,7 @@ async fn handle_frame( coverage: &CoverageWatch, search: &S, records: &dyn RecordStore, + resolver: &RepoIdResolver, ) { let cursor = HydrantCursor::new(frame.id); let signal = promotion_signal(frame.record.as_ref(), now_micros()); @@ -286,7 +383,7 @@ async fn handle_frame( if !record.collection.as_ref().starts_with(TANGLED_PREFIX) { return; } - if let Err(err) = apply_record(record, store, search, records).await { + if let Err(err) = apply_record(record, store, search, records, resolver).await { warn!(?err, "dropped record"); } } @@ -318,24 +415,42 @@ async fn apply_record( store: &EdgeStore, search: &S, records: &dyn RecordStore, + resolver: &RepoIdResolver, ) -> Result<(), IngestError> { let source = build_source_uri(&record)?; match record.action { RecordAction::Create | RecordAction::Update => { - records.remove(&source); let Some(value) = record.record else { + records.remove(&source); debug!(collection = %record.collection.as_ref(), "create/update missing record body, skipping"); return Ok(()); }; - let parsed = match Record::from_json_value(&record.collection, value) { + let bytes = serde_json::to_vec(&value) + .expect("serde_json::Value re-serialization is infallible"); + let parsed = match Record::from_json_bytes(record.collection.as_ref(), &bytes) { Ok(r) => r, Err(ExtractError::UnknownCollection(name)) => { + records.remove(&source); debug!(collection = %name, "unknown sh.tangled.* collection, skipping"); return Ok(()); } - Err(e) => return Err(e.into()), + Err(e) => { + records.remove(&source); + return Err(e.into()); + } }; - let edges = parsed.extract_edges(&source)?; + cache_body(records, &source, record.cid, bytes); + if let Record::Repo(repo) = &parsed { + resolver + .observe( + record.did.clone(), + record.rkey.clone(), + repo.repo_did.clone(), + ) + .await; + } + let raw_edges = parsed.extract_edges(&source)?; + let edges = normalize_subjects(raw_edges, resolver).await; store.upsert_source(&source, edges); if let Some(searchable) = SearchableRecord::try_from_record(parsed) { search.upsert(searchable.to_search_doc(&source)).await; @@ -353,6 +468,63 @@ async fn apply_record( Ok(()) } +fn cache_body( + records: &dyn RecordStore, + source: &AtUri, + cid: Option>, + bytes: Vec, +) { + match cid { + Some(cid) => records.put( + source.clone(), + Arc::new(RecordBody { + uri: source.clone(), + cid, + value: Bytes::from(bytes), + }), + ), + None => records.remove(source), + } +} + +async fn normalize_subjects(edges: Vec, resolver: &RepoIdResolver) -> Vec { + futures::stream::iter(edges) + .map(|edge| async move { + let Some((owner, rkey)) = parse_repo_subject(&edge.subject) else { + return edge; + }; + match resolver.resolve(&owner, &rkey).await { + Resolution::Mapped(repo_did) => Edge { + subject: owned_did_aturi(&repo_did), + ..edge + }, + Resolution::NoRepoDid | Resolution::Unresolvable => edge, + } + }) + .buffered(NORMALIZE_CONCURRENCY) + .collect() + .await +} + +fn parse_repo_subject(uri: &AtUri) -> Option<(Did, Rkey)> { + let collection = uri.collection()?; + if collection.as_ref() != "sh.tangled.repo" { + return None; + } + let AtIdentifier::Did(authority) = uri.authority() else { + return None; + }; + let rkey = uri.rkey()?; + let owner = Did::new_owned(authority.as_ref()).ok()?; + let rkey = Rkey::new_owned(rkey.as_ref()).ok()?; + Some((owner, rkey)) +} + +fn owned_did_aturi(did: &Did) -> AtUri { + AtUri::new_owned(format!("at://{}", did.as_ref())) + .expect("DID concatenated with at:// scheme is a valid AT-URI") +} + fn build_source_uri(r: &RecordFrame) -> Result, IngestError> { Ok(AtUri::from_parts_owned( r.did.as_ref(), @@ -365,14 +537,20 @@ fn build_source_uri(r: &RecordFrame) -> Result, IngestError> { mod tests { use super::*; use bobbin_edge_index::Coverage; - use bobbin_record_lru::NoopRecordStore; + use bobbin_record_lru::{CacheCapacity, LruRecordStore, NoopRecordStore, RecordStore}; use bobbin_types::search::NoopSearchSink; use jacquard_common::types::nsid::Nsid; use jacquard_common::types::tid::Tid; use serde_json::json; - fn fresh() -> (Arc, Arc) { - (Arc::new(EdgeStore::new()), Arc::new(CoverageWatch::new())) + const VALID_CID: &str = "bafyreieqygohnz2zqyvtvktbjpvhutphobcmbsnt4q5lc36ri7vpcmoz4i"; + + fn fresh() -> (Arc, Arc, Arc) { + ( + Arc::new(EdgeStore::new()), + Arc::new(CoverageWatch::new()), + Arc::new(RepoIdResolver::detached()), + ) } fn parse_frame(value: serde_json::Value) -> HydrantFrame { @@ -386,7 +564,7 @@ mod tests { #[tokio::test] async fn ignores_non_tangled_collections() { - let (store, cov) = fresh(); + let (store, cov, resolver) = fresh(); let frame: HydrantFrame = parse_frame(json!({ "id": 1, "type": "record", @@ -400,14 +578,22 @@ mod tests { "record": {"$type": "app.bsky.feed.post", "text": "hi"} } })); - handle_frame(frame, &store, &cov, &NoopSearchSink, &NoopRecordStore).await; + handle_frame( + frame, + &store, + &cov, + &NoopSearchSink, + &NoopRecordStore, + &resolver, + ) + .await; assert_eq!(store.key_count(), 0); assert_eq!(cov.snapshot().events_processed(), 1); } #[tokio::test] async fn create_then_delete_round_trips_a_star() { - let (store, cov) = fresh(); + let (store, cov, resolver) = fresh(); let create: HydrantFrame = parse_frame(json!({ "id": 10, "type": "record", @@ -425,7 +611,15 @@ mod tests { } } })); - handle_frame(create, &store, &cov, &NoopSearchSink, &NoopRecordStore).await; + handle_frame( + create, + &store, + &cov, + &NoopSearchSink, + &NoopRecordStore, + &resolver, + ) + .await; let key = bobbin_types::ids::EdgeKey::new( Nsid::new_static("sh.tangled.feed.star").unwrap(), AtUri::new_owned("at://did:plc:abalone").unwrap(), @@ -445,13 +639,21 @@ mod tests { "record": null } })); - handle_frame(delete, &store, &cov, &NoopSearchSink, &NoopRecordStore).await; + handle_frame( + delete, + &store, + &cov, + &NoopSearchSink, + &NoopRecordStore, + &resolver, + ) + .await; assert_eq!(store.count(&key), 0); } #[tokio::test] async fn update_replaces_prior_edges() { - let (store, cov) = fresh(); + let (store, cov, resolver) = fresh(); let mk = |subject_did: &str, id: u64| -> HydrantFrame { parse_frame(json!({ "id": id, @@ -477,6 +679,7 @@ mod tests { &cov, &NoopSearchSink, &NoopRecordStore, + &resolver, ) .await; handle_frame( @@ -485,6 +688,7 @@ mod tests { &cov, &NoopSearchSink, &NoopRecordStore, + &resolver, ) .await; @@ -499,9 +703,111 @@ mod tests { assert_eq!(store.count(&new), 1); } + #[tokio::test] + async fn create_with_cid_warms_record_lru() { + let (store, cov, resolver) = fresh(); + let lru = LruRecordStore::new(CacheCapacity::from_bytes(64 * 1024)); + let frame: HydrantFrame = parse_frame(json!({ + "id": 1, + "type": "record", + "record": { + "live": false, + "did": "did:plc:olaren", + "rev": fresh_tid().as_str(), + "collection": "sh.tangled.feed.star", + "rkey": "abcabcabcabcz", + "action": "create", + "cid": VALID_CID, + "record": { + "$type": "sh.tangled.feed.star", + "createdAt": "2026-05-01T00:00:00Z", + "subjectDid": "did:plc:abalone" + } + } + })); + let source = + AtUri::new_owned("at://did:plc:olaren/sh.tangled.feed.star/abcabcabcabcz").unwrap(); + handle_frame(frame, &store, &cov, &NoopSearchSink, &lru, &resolver).await; + let cached = lru.get(&source).expect("hydrant cid must seed the lru"); + assert_eq!(cached.cid.as_ref(), VALID_CID); + let parsed: serde_json::Value = serde_json::from_slice(&cached.value).unwrap(); + assert_eq!(parsed["subjectDid"], "did:plc:abalone"); + } + + #[tokio::test] + async fn create_without_cid_clears_record_lru() { + let (store, cov, resolver) = fresh(); + let lru = LruRecordStore::new(CacheCapacity::from_bytes(64 * 1024)); + let source = + AtUri::new_owned("at://did:plc:olaren/sh.tangled.feed.star/abcabcabcabcz").unwrap(); + let cid: Cid = VALID_CID.parse().unwrap(); + lru.put( + source.clone(), + Arc::new(RecordBody { + uri: source.clone(), + cid, + value: bytes::Bytes::from_static(b"{\"stale\":true}"), + }), + ); + let frame: HydrantFrame = parse_frame(json!({ + "id": 1, + "type": "record", + "record": { + "live": false, + "did": "did:plc:olaren", + "rev": fresh_tid().as_str(), + "collection": "sh.tangled.feed.star", + "rkey": "abcabcabcabcz", + "action": "update", + "record": { + "$type": "sh.tangled.feed.star", + "createdAt": "2026-05-01T00:00:00Z", + "subjectDid": "did:plc:abalone" + } + } + })); + handle_frame(frame, &store, &cov, &NoopSearchSink, &lru, &resolver).await; + assert!( + lru.get(&source).is_none(), + "missing cid means we cannot trust the body, so the lru must be cleared", + ); + } + + #[tokio::test] + async fn delete_evicts_record_lru() { + let (store, cov, resolver) = fresh(); + let lru = LruRecordStore::new(CacheCapacity::from_bytes(64 * 1024)); + let source = + AtUri::new_owned("at://did:plc:olaren/sh.tangled.feed.star/abcabcabcabcz").unwrap(); + let cid: Cid = VALID_CID.parse().unwrap(); + lru.put( + source.clone(), + Arc::new(RecordBody { + uri: source.clone(), + cid, + value: bytes::Bytes::from_static(b"{}"), + }), + ); + let frame: HydrantFrame = parse_frame(json!({ + "id": 1, + "type": "record", + "record": { + "live": false, + "did": "did:plc:olaren", + "rev": fresh_tid().as_str(), + "collection": "sh.tangled.feed.star", + "rkey": "abcabcabcabcz", + "action": "delete", + "record": null + } + })); + handle_frame(frame, &store, &cov, &NoopSearchSink, &lru, &resolver).await; + assert!(lru.get(&source).is_none()); + } + #[tokio::test] async fn live_recent_event_promotes_coverage_to_ready() { - let (store, cov) = fresh(); + let (store, cov, resolver) = fresh(); assert!(!cov.snapshot().is_ready()); let frame: HydrantFrame = parse_frame(json!({ "id": 99, @@ -520,14 +826,22 @@ mod tests { } } })); - handle_frame(frame, &store, &cov, &NoopSearchSink, &NoopRecordStore).await; + handle_frame( + frame, + &store, + &cov, + &NoopSearchSink, + &NoopRecordStore, + &resolver, + ) + .await; assert!(cov.snapshot().is_ready()); assert_eq!(cov.snapshot().last_cursor(), HydrantCursor::new(99)); } #[tokio::test] async fn live_but_stale_rev_does_not_promote() { - let (store, cov) = fresh(); + let (store, cov, resolver) = fresh(); let stale_tid = Tid::from_time(1_000_000, 0); let frame: HydrantFrame = parse_frame(json!({ "id": 7, @@ -546,7 +860,15 @@ mod tests { } } })); - handle_frame(frame, &store, &cov, &NoopSearchSink, &NoopRecordStore).await; + handle_frame( + frame, + &store, + &cov, + &NoopSearchSink, + &NoopRecordStore, + &resolver, + ) + .await; assert!(!cov.snapshot().is_ready()); assert!(matches!(cov.snapshot(), Coverage::Warming { .. })); } @@ -597,7 +919,7 @@ mod tests { #[tokio::test] async fn identity_frame_advances_cursor_only() { - let (store, cov) = fresh(); + let (store, cov, resolver) = fresh(); let frame: HydrantFrame = parse_frame(json!({ "id": 5, "type": "identity", @@ -606,7 +928,15 @@ mod tests { "handle": "olaren.dev" } })); - handle_frame(frame, &store, &cov, &NoopSearchSink, &NoopRecordStore).await; + handle_frame( + frame, + &store, + &cov, + &NoopSearchSink, + &NoopRecordStore, + &resolver, + ) + .await; assert_eq!(store.key_count(), 0); assert_eq!(cov.snapshot().last_cursor(), HydrantCursor::new(5)); assert!(!cov.snapshot().is_ready()); @@ -614,26 +944,67 @@ mod tests { #[tokio::test] async fn account_frame_advances_cursor_only() { - let (store, cov) = fresh(); + let (store, cov, resolver) = fresh(); let frame: HydrantFrame = parse_frame(json!({ "id": 6, "type": "account", "account": {"did": "did:plc:olaren", "active": true} })); - handle_frame(frame, &store, &cov, &NoopSearchSink, &NoopRecordStore).await; + handle_frame( + frame, + &store, + &cov, + &NoopSearchSink, + &NoopRecordStore, + &resolver, + ) + .await; assert_eq!(store.key_count(), 0); assert_eq!(cov.snapshot().last_cursor(), HydrantCursor::new(6)); } #[tokio::test] async fn unknown_frame_kind_advances_cursor_without_panic() { - let (store, cov) = fresh(); + let (store, cov, resolver) = fresh(); let frame: HydrantFrame = parse_frame(json!({"id": 8, "type": "future_event"})); assert_eq!(frame.kind, FrameKind::Other); - handle_frame(frame, &store, &cov, &NoopSearchSink, &NoopRecordStore).await; + handle_frame( + frame, + &store, + &cov, + &NoopSearchSink, + &NoopRecordStore, + &resolver, + ) + .await; assert_eq!(cov.snapshot().last_cursor(), HydrantCursor::new(8)); } + fn fresh_runtime(cancel: CancellationToken) -> IngestRuntime { + IngestRuntime { + store: Arc::new(EdgeStore::new()), + coverage: Arc::new(CoverageWatch::new()), + search: Arc::new(NoopSearchSink), + records: Arc::new(NoopRecordStore) as Arc, + resolver: Arc::new(RepoIdResolver::detached()), + cancel, + } + } + + #[tokio::test(start_paused = true)] + async fn cancel_token_short_circuits_reconnect_sleep() { + let cfg = IngestConfig::new(Url::parse("ws://127.0.0.1:1").unwrap()); + let cancel = CancellationToken::new(); + let runtime = fresh_runtime(cancel.clone()); + let task = tokio::spawn(async move { run(cfg, runtime).await }); + tokio::time::advance(Duration::from_millis(10)).await; + cancel.cancel(); + let outcome = tokio::time::timeout(Duration::from_secs(1), task) + .await + .expect("ingest must stop within timeout once cancel fires"); + assert!(matches!(outcome, Ok(Ok(()))), "got {outcome:?}"); + } + #[test] fn jittered_stays_within_one_quarter_of_base() { let base = Duration::from_secs(1); @@ -644,4 +1015,301 @@ mod tests { assert!(j <= cap, "jitter must not exceed +25%, got {:?}", j); }); } + + #[tokio::test] + async fn star_after_observed_repo_keys_on_repo_did() { + let (store, cov, resolver) = fresh(); + let repo: HydrantFrame = parse_frame(json!({ + "id": 1, + "type": "record", + "record": { + "live": false, + "did": "did:plc:nel", + "rev": fresh_tid().as_str(), + "collection": "sh.tangled.repo", + "rkey": "abcabcabcabcz", + "action": "create", + "record": { + "$type": "sh.tangled.repo", + "createdAt": "2026-05-01T00:00:00Z", + "knot": "oyster.cafe", + "name": "abalone", + "repoDid": "did:plc:abalone" + } + } + })); + handle_frame( + repo, + &store, + &cov, + &NoopSearchSink, + &NoopRecordStore, + &resolver, + ) + .await; + + let star: HydrantFrame = parse_frame(json!({ + "id": 2, + "type": "record", + "record": { + "live": false, + "did": "did:plc:olaren", + "rev": fresh_tid().as_str(), + "collection": "sh.tangled.feed.star", + "rkey": "abcabcabcabcz", + "action": "create", + "record": { + "$type": "sh.tangled.feed.star", + "createdAt": "2026-05-01T00:00:00Z", + "subject": "at://did:plc:nel/sh.tangled.repo/abcabcabcabcz" + } + } + })); + handle_frame( + star, + &store, + &cov, + &NoopSearchSink, + &NoopRecordStore, + &resolver, + ) + .await; + + let nsid = Nsid::new_static("sh.tangled.feed.star").unwrap(); + let repo_keyed = bobbin_types::ids::EdgeKey::new( + nsid.clone(), + AtUri::new_owned("at://did:plc:abalone").unwrap(), + ); + let owner_keyed = + bobbin_types::ids::EdgeKey::new(nsid, AtUri::new_owned("at://did:plc:nel").unwrap()); + assert_eq!( + store.count(&repo_keyed), + 1, + "star should be keyed on repoDID once the repo is observed" + ); + assert_eq!( + store.count(&owner_keyed), + 0, + "owner DID should not collect the edge" + ); + } + + #[tokio::test] + async fn unresolvable_repo_subject_keeps_at_uri() { + let (store, cov, resolver) = fresh(); + let star: HydrantFrame = parse_frame(json!({ + "id": 1, + "type": "record", + "record": { + "live": false, + "did": "did:plc:olaren", + "rev": fresh_tid().as_str(), + "collection": "sh.tangled.feed.star", + "rkey": "abcabcabcabcz", + "action": "create", + "record": { + "$type": "sh.tangled.feed.star", + "createdAt": "2026-05-01T00:00:00Z", + "subject": "at://did:plc:nel/sh.tangled.repo/abcabcabcabcz" + } + } + })); + handle_frame( + star, + &store, + &cov, + &NoopSearchSink, + &NoopRecordStore, + &resolver, + ) + .await; + + let nsid = Nsid::new_static("sh.tangled.feed.star").unwrap(); + let owner_keyed = bobbin_types::ids::EdgeKey::new( + nsid.clone(), + AtUri::new_owned("at://did:plc:nel").unwrap(), + ); + let uri_keyed = bobbin_types::ids::EdgeKey::new( + nsid, + AtUri::new_owned("at://did:plc:nel/sh.tangled.repo/abcabcabcabcz").unwrap(), + ); + assert_eq!( + store.count(&owner_keyed), + 0, + "must not silently misfile under the authoring DID", + ); + assert_eq!( + store.count(&uri_keyed), + 1, + "an unresolvable repo subject keeps the at-uri so the edge stays queryable", + ); + } + + #[tokio::test] + async fn repo_without_repo_did_keeps_full_uri_subject() { + let (store, cov, resolver) = fresh(); + let repo: HydrantFrame = parse_frame(json!({ + "id": 1, + "type": "record", + "record": { + "live": false, + "did": "did:plc:nel", + "rev": fresh_tid().as_str(), + "collection": "sh.tangled.repo", + "rkey": "abcabcabcabcz", + "action": "create", + "record": { + "$type": "sh.tangled.repo", + "createdAt": "2026-05-01T00:00:00Z", + "knot": "oyster.cafe", + "name": "abalone" + } + } + })); + handle_frame( + repo, + &store, + &cov, + &NoopSearchSink, + &NoopRecordStore, + &resolver, + ) + .await; + + let star: HydrantFrame = parse_frame(json!({ + "id": 2, + "type": "record", + "record": { + "live": false, + "did": "did:plc:olaren", + "rev": fresh_tid().as_str(), + "collection": "sh.tangled.feed.star", + "rkey": "abcabcabcabcz", + "action": "create", + "record": { + "$type": "sh.tangled.feed.star", + "createdAt": "2026-05-01T00:00:00Z", + "subject": "at://did:plc:nel/sh.tangled.repo/abcabcabcabcz" + } + } + })); + handle_frame( + star, + &store, + &cov, + &NoopSearchSink, + &NoopRecordStore, + &resolver, + ) + .await; + + let key = bobbin_types::ids::EdgeKey::new( + Nsid::new_static("sh.tangled.feed.star").unwrap(), + AtUri::new_owned("at://did:plc:nel/sh.tangled.repo/abcabcabcabcz").unwrap(), + ); + assert_eq!( + store.count(&key), + 1, + "the at-uri is the only canonical identity for a repo with no repoDID", + ); + } + + #[tokio::test] + async fn explicit_subject_did_skips_normalization() { + let (store, cov, resolver) = fresh(); + let star: HydrantFrame = parse_frame(json!({ + "id": 1, + "type": "record", + "record": { + "live": false, + "did": "did:plc:olaren", + "rev": fresh_tid().as_str(), + "collection": "sh.tangled.feed.star", + "rkey": "abcabcabcabcz", + "action": "create", + "record": { + "$type": "sh.tangled.feed.star", + "createdAt": "2026-05-01T00:00:00Z", + "subjectDid": "did:plc:abalone" + } + } + })); + handle_frame( + star, + &store, + &cov, + &NoopSearchSink, + &NoopRecordStore, + &resolver, + ) + .await; + let key = bobbin_types::ids::EdgeKey::new( + Nsid::new_static("sh.tangled.feed.star").unwrap(), + AtUri::new_owned("at://did:plc:abalone").unwrap(), + ); + assert_eq!(store.count(&key), 1); + } + + #[tokio::test] + async fn issue_with_repo_uri_resolves_to_repo_did() { + let (store, cov, resolver) = fresh(); + resolver + .observe( + Did::new_owned("did:plc:nel").unwrap(), + Rkey::new_owned("abcabcabcabcz").unwrap(), + Some(Did::new_owned("did:plc:abalone").unwrap()), + ) + .await; + let issue: HydrantFrame = parse_frame(json!({ + "id": 1, + "type": "record", + "record": { + "live": false, + "did": "did:plc:olaren", + "rev": fresh_tid().as_str(), + "collection": "sh.tangled.repo.issue", + "rkey": "abcabcabcabcz", + "action": "create", + "record": { + "$type": "sh.tangled.repo.issue", + "createdAt": "2026-05-01T00:00:00Z", + "title": "bug", + "repo": "at://did:plc:nel/sh.tangled.repo/abcabcabcabcz" + } + } + })); + handle_frame( + issue, + &store, + &cov, + &NoopSearchSink, + &NoopRecordStore, + &resolver, + ) + .await; + let key = bobbin_types::ids::EdgeKey::new( + Nsid::new_static("sh.tangled.repo.issue").unwrap(), + AtUri::new_owned("at://did:plc:abalone").unwrap(), + ); + assert_eq!(store.count(&key), 1); + } + + #[tokio::test] + async fn cancel_short_circuits_a_hung_ws_connect() { + let _listener = tokio::net::TcpListener::bind("127.0.0.1:0") + .await + .expect("bind sink listener"); + let port = _listener.local_addr().expect("local addr").port(); + let cfg = + IngestConfig::new(Url::parse(&format!("ws://127.0.0.1:{port}")).expect("hydrant url")); + let cancel = CancellationToken::new(); + let runtime = fresh_runtime(cancel.clone()); + let task = tokio::spawn(async move { run(cfg, runtime).await }); + tokio::time::sleep(Duration::from_millis(100)).await; + cancel.cancel(); + let outcome = tokio::time::timeout(Duration::from_secs(2), task) + .await + .expect("cancel must short-circuit the hung ws connect"); + assert!(matches!(outcome, Ok(Ok(()))), "got {outcome:?}"); + } }