This repository has no description
0

Configure Feed

Select the types of activity you want to include in your feed.

core / bobbin / crates / ingest / src / lib.rs
119 kB 3504 lines
1use std::collections::{HashSet, VecDeque}; 2use std::num::NonZeroUsize; 3use std::sync::Arc; 4use std::time::Duration; 5 6use bobbin_edge_index::{ 7 ApplyOutcome, Coverage, CoverageWatch, EdgeStore, HydrantCursor, IssueStateKind, 8 PromotionSignal, PullStatusKind, StateIndex, delete_record_indexes, upsert_record_indexes, 9}; 10use bobbin_knot_ingest::{CapabilityGate, KnotRegistry}; 11use bobbin_record_lru::RecordStore; 12use bobbin_resolver::{NormalizeRepoRefs, decode_canon_or_upgrade_bytes, synthesize_created_at}; 13use bobbin_runtime::{ 14 Clock, Entropy, NetworkError, RuntimeHasher, UnixMicros, WsConn, WsMessage, WsStream, 15 WsTransport, 16}; 17use bobbin_types::edges::{Edge, ExtractError, Record}; 18use bobbin_types::ids::{RepoIdent, SubjectRef}; 19use bobbin_types::knot_acl::KnotHostKey; 20use bobbin_types::record::RecordBody; 21use bobbin_types::search::{SearchSink, SearchableRecord}; 22use bytes::Bytes; 23use futures::StreamExt; 24use jacquard_common::DefaultStr; 25use jacquard_common::types::did::Did; 26use jacquard_common::types::ident::AtIdentifier; 27use jacquard_common::types::nsid::Nsid; 28use jacquard_common::types::recordkey::Rkey; 29use jacquard_common::types::string::{AtStrError, AtUri, Cid}; 30use thiserror::Error; 31use tokio::time::Instant; 32use tokio_stream::wrappers::ReceiverStream; 33use tokio_util::sync::CancellationToken; 34use tracing::{debug, info, warn}; 35use url::Url; 36 37mod frame; 38mod resolver; 39mod shadow; 40mod warming; 41use frame::HydrantStreamErrorFrame; 42pub use frame::{FrameKind, HydrantFrame, RecordAction, RecordFrame}; 43pub use resolver::{RepoIdResolver, Resolution}; 44pub use shadow::{WarmingShadowBuffer, WarmingShadowSnapshot}; 45pub use warming::{ParkedUpsert, WarmingBuffer, WarmingBufferSnapshot}; 46 47const TANGLED_PREFIX: &str = "sh.tangled."; 48const RECONNECT_INITIAL_DELAY: Duration = Duration::from_millis(500); 49const RECONNECT_MAX_DELAY: Duration = Duration::from_secs(30); 50const PING_INTERVAL: Duration = Duration::from_secs(20); 51const PONG_TIMEOUT: Duration = Duration::from_secs(15); 52const READY_SKEW: Duration = Duration::from_secs(60); 53const FRAME_CHANNEL_DEPTH: usize = 256; 54const CONTROL_CHANNEL_DEPTH: usize = 16; 55const READER_HOLD_LIMIT: usize = 64; 56const SEND_TIMEOUT: Duration = Duration::from_secs(10); 57const NORMAL_CLOSE: u16 = 1000; 58const METRICS_DUMP_INTERVAL: Duration = Duration::from_secs(10); 59const WARMING_FLUSH_PARALLELISM: usize = 64; 60pub const DEFAULT_INGEST_PARALLELISM: NonZeroUsize = match NonZeroUsize::new(16) { 61 Some(n) => n, 62 None => unreachable!(), 63}; 64 65#[derive(Clone, Debug)] 66pub struct IngestConfig { 67 pub hydrant_base: Url, 68 pub start_cursor: HydrantCursor, 69 pub parallelism: NonZeroUsize, 70} 71 72impl IngestConfig { 73 pub fn new(hydrant_base: Url) -> Self { 74 Self { 75 hydrant_base, 76 start_cursor: HydrantCursor::new(0), 77 parallelism: DEFAULT_INGEST_PARALLELISM, 78 } 79 } 80 81 fn stream_url(&self, cursor: HydrantCursor) -> Result<Url, IngestError> { 82 let mut url = self.hydrant_base.clone(); 83 match url.scheme() { 84 "http" => url 85 .set_scheme("ws") 86 .map_err(|_| IngestError::Url("set ws scheme"))?, 87 "https" => url 88 .set_scheme("wss") 89 .map_err(|_| IngestError::Url("set wss scheme"))?, 90 "ws" | "wss" => {} 91 other => return Err(IngestError::UnknownScheme(other.to_owned())), 92 } 93 url.set_path("/stream"); 94 url.query_pairs_mut() 95 .clear() 96 .append_pair("cursor", &cursor.raw().to_string()); 97 Ok(url) 98 } 99} 100 101#[derive(Debug, Error)] 102pub enum IngestError { 103 #[error("invalid hydrant url: {0}")] 104 Url(&'static str), 105 #[error("unsupported url scheme: {0}")] 106 UnknownScheme(String), 107 #[error("network: {0}")] 108 Network(#[from] NetworkError), 109 #[error("frame decode: {0}")] 110 Decode(#[from] serde_json::Error), 111 #[error("invalid at-uri synthesized from frame: {0}")] 112 InvalidAtUri(#[from] AtStrError), 113 #[error("record extraction: {0}")] 114 Extract(#[from] ExtractError), 115 #[error("hydrant did not respond to ping within {0:?}")] 116 PongTimeout(Duration), 117 #[error("websocket send blocked for at least {0:?}, treating link as dead")] 118 SendTimeout(Duration), 119 #[error("hydrant disconnected because bobbin's stream consumer fell behind: {message}")] 120 ConsumerTooSlow { message: String }, 121 #[error("hydrant signaled stream error {code}: {message}")] 122 HydrantStream { code: String, message: String }, 123} 124 125#[derive(Clone, Copy, Debug, Eq, PartialEq, Hash)] 126pub enum DisconnectKind { 127 Url, 128 UnknownScheme, 129 Network, 130 Decode, 131 InvalidAtUri, 132 Extract, 133 PongTimeout, 134 SendTimeout, 135 ConsumerTooSlow, 136 HydrantStream, 137} 138 139impl DisconnectKind { 140 pub fn from_error(err: &IngestError) -> Self { 141 match err { 142 IngestError::Url(_) => Self::Url, 143 IngestError::UnknownScheme(_) => Self::UnknownScheme, 144 IngestError::Network(_) => Self::Network, 145 IngestError::Decode(_) => Self::Decode, 146 IngestError::InvalidAtUri(_) => Self::InvalidAtUri, 147 IngestError::Extract(_) => Self::Extract, 148 IngestError::PongTimeout(_) => Self::PongTimeout, 149 IngestError::SendTimeout(_) => Self::SendTimeout, 150 IngestError::ConsumerTooSlow { .. } => Self::ConsumerTooSlow, 151 IngestError::HydrantStream { .. } => Self::HydrantStream, 152 } 153 } 154} 155 156#[derive(Clone, Debug, Eq, PartialEq)] 157pub struct DisconnectSnapshot { 158 pub kind: DisconnectKind, 159 pub message: String, 160 pub at_unix_micros: UnixMicros, 161 pub last_cursor: HydrantCursor, 162} 163 164#[derive(Default)] 165pub struct DisconnectSink { 166 last: std::sync::Mutex<Option<DisconnectSnapshot>>, 167 count: std::sync::atomic::AtomicU64, 168} 169 170impl DisconnectSink { 171 pub fn new() -> Self { 172 Self::default() 173 } 174 175 pub fn record(&self, snap: DisconnectSnapshot) { 176 *self.last.lock().expect("disconnect sink mutex poisoned") = Some(snap); 177 self.count 178 .fetch_add(1, std::sync::atomic::Ordering::Relaxed); 179 } 180 181 pub fn snapshot(&self) -> Option<DisconnectSnapshot> { 182 self.last 183 .lock() 184 .expect("disconnect sink mutex poisoned") 185 .clone() 186 } 187 188 pub fn count(&self) -> u64 { 189 self.count.load(std::sync::atomic::Ordering::Relaxed) 190 } 191} 192 193#[derive(Clone, Copy, Debug, Eq, PartialEq)] 194enum SessionOutcome { 195 Progressed, 196 Empty, 197} 198 199#[derive(Debug)] 200struct SessionEnd { 201 outcome: SessionOutcome, 202 error: Option<IngestError>, 203} 204 205pub struct IngestRuntime<S: SearchSink + 'static> { 206 pub store: Arc<EdgeStore>, 207 pub issue_states: Arc<StateIndex<IssueStateKind>>, 208 pub pull_statuses: Arc<StateIndex<PullStatusKind>>, 209 pub coverage: Arc<CoverageWatch>, 210 pub search: Arc<S>, 211 pub records: Arc<dyn RecordStore>, 212 pub resolver: Arc<RepoIdResolver>, 213 pub clock: Arc<dyn Clock>, 214 pub entropy: Arc<dyn Entropy>, 215 pub ws: Arc<dyn WsTransport>, 216 pub cancel: CancellationToken, 217 pub disconnects: Option<Arc<DisconnectSink>>, 218 pub warming_shadow: Option<Arc<WarmingShadowBuffer>>, 219 pub warming_buffer: Option<Arc<WarmingBuffer>>, 220 pub knot_registry: Option<Arc<KnotRegistry>>, 221 pub knot_gate: Option<Arc<CapabilityGate>>, 222} 223 224impl<S: SearchSink + 'static> Clone for IngestRuntime<S> { 225 fn clone(&self) -> Self { 226 Self { 227 store: self.store.clone(), 228 issue_states: self.issue_states.clone(), 229 pull_statuses: self.pull_statuses.clone(), 230 coverage: self.coverage.clone(), 231 search: self.search.clone(), 232 records: self.records.clone(), 233 resolver: self.resolver.clone(), 234 clock: self.clock.clone(), 235 entropy: self.entropy.clone(), 236 ws: self.ws.clone(), 237 cancel: self.cancel.clone(), 238 disconnects: self.disconnects.clone(), 239 warming_shadow: self.warming_shadow.clone(), 240 warming_buffer: self.warming_buffer.clone(), 241 knot_registry: self.knot_registry.clone(), 242 knot_gate: self.knot_gate.clone(), 243 } 244 } 245} 246 247impl<S: SearchSink + 'static> IngestRuntime<S> { 248 fn pipeline_ctx(&self) -> PipelineCtx<'_, S> { 249 PipelineCtx { 250 resolver: &self.resolver, 251 store: &self.store, 252 issue_states: &self.issue_states, 253 pull_statuses: &self.pull_statuses, 254 coverage: &self.coverage, 255 records: &*self.records, 256 search: &self.search, 257 shadow: self.warming_shadow.as_deref(), 258 buffer: self.warming_buffer.as_deref(), 259 knot_registry: self.knot_registry.as_deref(), 260 knot_gate: self.knot_gate.as_deref(), 261 } 262 } 263} 264 265struct PipelineCtx<'a, S: SearchSink + 'static> { 266 resolver: &'a RepoIdResolver, 267 store: &'a EdgeStore, 268 issue_states: &'a StateIndex<IssueStateKind>, 269 pull_statuses: &'a StateIndex<PullStatusKind>, 270 coverage: &'a CoverageWatch, 271 records: &'a dyn RecordStore, 272 search: &'a S, 273 shadow: Option<&'a WarmingShadowBuffer>, 274 buffer: Option<&'a WarmingBuffer>, 275 knot_registry: Option<&'a KnotRegistry>, 276 knot_gate: Option<&'a CapabilityGate>, 277} 278 279pub async fn run<S: SearchSink + 'static>( 280 config: IngestConfig, 281 runtime: IngestRuntime<S>, 282) -> Result<(), IngestError> { 283 let metrics_dumper = spawn_metrics_dumper(&runtime); 284 let idle_promoter = spawn_idle_promoter(&runtime); 285 let warming_flusher = spawn_warming_flusher(&runtime); 286 let result = run_inner(config, &runtime).await; 287 if let Err(join) = metrics_dumper.await { 288 warn!(?join, "metrics dumper task panicked"); 289 } 290 if let Err(join) = idle_promoter.await { 291 warn!(?join, "idle promoter task panicked"); 292 } 293 if let Some(handle) = warming_flusher 294 && let Err(join) = handle.await 295 { 296 warn!(?join, "warming flusher task panicked"); 297 } 298 result 299} 300 301async fn run_inner<S: SearchSink + 'static>( 302 config: IngestConfig, 303 runtime: &IngestRuntime<S>, 304) -> Result<(), IngestError> { 305 let mut backoff = RECONNECT_INITIAL_DELAY; 306 loop { 307 let cursor = next_connect_cursor(runtime.coverage.snapshot(), config.start_cursor); 308 let SessionEnd { outcome, error } = run_session(&config, cursor, runtime).await; 309 if runtime.cancel.is_cancelled() { 310 info!( 311 last_cursor = runtime.coverage.snapshot().last_cursor().raw(), 312 "ingest stopped after shutdown signal" 313 ); 314 return Ok(()); 315 } 316 match (outcome, &error) { 317 (SessionOutcome::Progressed, None) => { 318 info!("hydrant stream closed after delivering frames, reconnecting") 319 } 320 (SessionOutcome::Empty, None) => { 321 warn!("hydrant stream closed without delivering frames") 322 } 323 (_, Some(err)) => warn!(?err, "hydrant stream errored"), 324 } 325 if let (Some(sink), Some(err)) = (runtime.disconnects.as_ref(), error.as_ref()) { 326 sink.record(DisconnectSnapshot { 327 kind: DisconnectKind::from_error(err), 328 message: err.to_string(), 329 at_unix_micros: runtime.clock.now_unix_micros(), 330 last_cursor: runtime.coverage.snapshot().last_cursor(), 331 }); 332 } 333 let made_progress = matches!(outcome, SessionOutcome::Progressed); 334 if made_progress { 335 backoff = RECONNECT_INITIAL_DELAY; 336 } else { 337 tokio::select! { 338 biased; 339 _ = runtime.cancel.cancelled() => return Ok(()), 340 _ = runtime.clock.sleep(jittered(backoff, &*runtime.entropy)) => {} 341 } 342 backoff = (backoff * 2).min(RECONNECT_MAX_DELAY); 343 } 344 } 345} 346 347fn spawn_warming_flusher<S: SearchSink + 'static>( 348 runtime: &IngestRuntime<S>, 349) -> Option<tokio::task::JoinHandle<()>> { 350 runtime.warming_buffer.as_ref()?; 351 let rt = runtime.clone(); 352 Some(tokio::spawn(async move { 353 let buffer = rt 354 .warming_buffer 355 .as_deref() 356 .expect("warming flusher only spawns when buffer is set"); 357 let mut rx = rt.coverage.subscribe(); 358 let reason = loop { 359 if rx.borrow_and_update().is_ready() { 360 break FlushReason::Ready; 361 } 362 tokio::select! { 363 biased; 364 _ = rt.cancel.cancelled() => break FlushReason::Cancelled, 365 res = rx.changed() => match res { 366 Ok(()) => continue, 367 Err(_) => break FlushReason::CoverageDropped, 368 }, 369 } 370 }; 371 flush_warming_buffer(&rt, buffer, reason).await; 372 })) 373} 374 375#[derive(Clone, Copy, Debug, Eq, PartialEq)] 376enum FlushReason { 377 Ready, 378 Cancelled, 379 CoverageDropped, 380} 381 382impl FlushReason { 383 fn as_str(self) -> &'static str { 384 match self { 385 FlushReason::Ready => "ready", 386 FlushReason::Cancelled => "cancelled", 387 FlushReason::CoverageDropped => "coverage_dropped", 388 } 389 } 390} 391 392async fn flush_warming_buffer<S: SearchSink + 'static>( 393 runtime: &IngestRuntime<S>, 394 buffer: &WarmingBuffer, 395 reason: FlushReason, 396) { 397 let drained = buffer.drain_for_promote().await; 398 if drained.is_empty() { 399 return; 400 } 401 if reason != FlushReason::Ready { 402 info!( 403 target: "bobbin_ingest::warming", 404 abandoned_entries = drained.len(), 405 reason = reason.as_str(), 406 "abandoning parked items on non-ready flush", 407 ); 408 return; 409 } 410 let hasher = buffer.hasher().clone(); 411 let unique: HashSet<RepoIdent, RuntimeHasher> = drained 412 .iter() 413 .flat_map(|(_, deps)| deps.iter().cloned()) 414 .fold(HashSet::with_hasher(hasher), |mut acc, dep| { 415 acc.insert(dep); 416 acc 417 }); 418 if !unique.is_empty() { 419 let resolver = runtime.resolver.clone(); 420 let _: Vec<()> = futures::stream::iter(unique) 421 .map(|key| { 422 let resolver = resolver.clone(); 423 async move { 424 let _ = resolver.resolve(&key.owner, &key.rkey).await; 425 } 426 }) 427 .buffer_unordered(WARMING_FLUSH_PARALLELISM) 428 .collect() 429 .await; 430 } 431 let upserts: Vec<ParkedUpsert> = drained.into_iter().map(|(u, _)| u).collect(); 432 let count = upserts.len(); 433 let ctx = runtime.pipeline_ctx(); 434 finalize_drained(&ctx, upserts).await; 435 info!( 436 target: "bobbin_ingest::warming", 437 flushed_entries = count, 438 reason = reason.as_str(), 439 "drained warming buffer", 440 ); 441} 442 443const IDLE_PROMOTE_WINDOW: Duration = Duration::from_secs(15); 444const IDLE_PROMOTE_MIN_EVENTS: u64 = 256; 445 446fn spawn_idle_promoter<S: SearchSink + 'static>( 447 runtime: &IngestRuntime<S>, 448) -> tokio::task::JoinHandle<()> { 449 let rt = runtime.clone(); 450 tokio::spawn(async move { 451 let mut prev = rt.coverage.snapshot().events_processed(); 452 loop { 453 tokio::select! { 454 biased; 455 _ = rt.cancel.cancelled() => return, 456 _ = rt.clock.sleep(IDLE_PROMOTE_WINDOW) => {} 457 } 458 let snap = rt.coverage.snapshot(); 459 if snap.is_ready() { 460 return; 461 } 462 let processed = snap.events_processed(); 463 if processed >= IDLE_PROMOTE_MIN_EVENTS && processed == prev { 464 rt.coverage.update(|c| c.force_ready()); 465 info!( 466 target: "bobbin_ingest::coverage", 467 events_processed = processed, 468 last_cursor = snap.last_cursor().raw(), 469 "stream idle, promoting coverage to ready", 470 ); 471 return; 472 } 473 prev = processed; 474 } 475 }) 476} 477 478fn spawn_metrics_dumper<S: SearchSink + 'static>( 479 runtime: &IngestRuntime<S>, 480) -> tokio::task::JoinHandle<()> { 481 let rt = runtime.clone(); 482 tokio::spawn(async move { 483 loop { 484 tokio::select! { 485 biased; 486 _ = rt.cancel.cancelled() => break, 487 _ = rt.clock.sleep(METRICS_DUMP_INTERVAL) => { 488 let s = rt.resolver.stats(); 489 info!( 490 target: "bobbin_ingest::metrics", 491 resolver_hits = s.hits, 492 resolver_misses_mapped = s.misses_mapped, 493 resolver_misses_no_repo_did = s.misses_no_repo_did, 494 resolver_misses_unresolvable = s.misses_unresolvable, 495 resolver_misses_transient = s.misses_transient, 496 resolver_misses_no_client = s.misses_no_client, 497 resolver_miss_latency_micros_avg = s.miss_latency_micros_avg().unwrap_or(0), 498 resolver_miss_latency_micros_max = s.miss_latency_micros_max, 499 resolver_total = s.total(), 500 "resolver stats", 501 ); 502 } 503 } 504 } 505 }) 506} 507 508fn next_connect_cursor(snapshot: Coverage, start: HydrantCursor) -> HydrantCursor { 509 if snapshot.events_processed() == 0 { 510 start 511 } else { 512 HydrantCursor::new(snapshot.last_cursor().raw().saturating_add(1)) 513 } 514} 515 516fn jittered(base: Duration, entropy: &dyn Entropy) -> Duration { 517 let base_ms = u64::try_from(base.as_millis()).unwrap_or(u64::MAX); 518 let cap_ms = (base_ms / 4).max(1); 519 base + Duration::from_millis(entropy.next_u64() % cap_ms) 520} 521 522async fn run_session<S: SearchSink + 'static>( 523 config: &IngestConfig, 524 cursor: HydrantCursor, 525 runtime: &IngestRuntime<S>, 526) -> SessionEnd { 527 let url = match config.stream_url(cursor) { 528 Ok(u) => u, 529 Err(e) => { 530 return SessionEnd { 531 outcome: SessionOutcome::Empty, 532 error: Some(e), 533 }; 534 } 535 }; 536 info!(%url, "connecting to hydrant /stream"); 537 let connect = tokio::select! { 538 biased; 539 _ = runtime.cancel.cancelled() => { 540 return SessionEnd { outcome: SessionOutcome::Empty, error: None }; 541 } 542 res = runtime.ws.connect(url) => res, 543 }; 544 let WsConn { 545 sink: mut ws_sink, 546 stream: ws_stream, 547 } = match connect { 548 Ok(c) => c, 549 Err(e) => { 550 return SessionEnd { 551 outcome: SessionOutcome::Empty, 552 error: Some(IngestError::Network(e)), 553 }; 554 } 555 }; 556 557 let (frame_tx, frame_rx) = tokio::sync::mpsc::channel::<HydrantFrame>(FRAME_CHANNEL_DEPTH); 558 let parallelism = config.parallelism.get(); 559 let processor_runtime = runtime.clone(); 560 let processor = tokio::spawn(async move { 561 let cancel = processor_runtime.cancel.clone(); 562 let prep_rt = processor_runtime.clone(); 563 let resolve_rt = processor_runtime.clone(); 564 let commit_rt = processor_runtime; 565 566 let pipeline = ReceiverStream::new(frame_rx) 567 .map(move |frame| prep_stage(frame, prep_rt.clone())) 568 .buffered(parallelism) 569 .map(move |staged| resolve_stage(staged, resolve_rt.clone())) 570 .buffered(parallelism) 571 .for_each(move |staged| commit_stage(staged, commit_rt.clone(), parallelism)); 572 573 tokio::select! { 574 biased; 575 _ = cancel.cancelled() => {}, 576 _ = pipeline => {}, 577 } 578 }); 579 580 let (control_tx, mut control_rx) = tokio::sync::mpsc::channel::<WsEvent>(CONTROL_CHANNEL_DEPTH); 581 let session_cancel = runtime.cancel.child_token(); 582 let reader_cancel = session_cancel.clone(); 583 let reader = tokio::spawn(reader_loop(ws_stream, frame_tx, control_tx, reader_cancel)); 584 585 let mut next_ping = runtime.clock.now_instant() + PING_INTERVAL; 586 let mut pong_deadline: Option<Instant> = None; 587 588 let writer_error: Option<IngestError> = loop { 589 tokio::select! { 590 biased; 591 _ = runtime.cancel.cancelled() => { 592 let _ = timed_send( 593 &mut ws_sink, 594 WsMessage::Close { code: NORMAL_CLOSE, reason: "bobbin shutdown".to_owned() }, 595 ).await; 596 break None; 597 } 598 _ = runtime.clock.sleep_until(next_ping) => { 599 next_ping = runtime.clock.now_instant() + PING_INTERVAL; 600 if pong_deadline.is_none() { 601 if let Err(e) = timed_send(&mut ws_sink, WsMessage::Ping(Bytes::new())).await { 602 break Some(e); 603 } 604 pong_deadline = Some(runtime.clock.now_instant() + PONG_TIMEOUT); 605 } 606 } 607 _ = wait_until(pong_deadline, runtime.clock.as_ref()) => { 608 break Some(IngestError::PongTimeout(PONG_TIMEOUT)); 609 } 610 evt = control_rx.recv() => { 611 let Some(evt) = evt else { break None; }; 612 match evt { 613 WsEvent::IncomingPing(payload) => { 614 if let Err(e) = timed_send(&mut ws_sink, WsMessage::Pong(payload)).await { 615 break Some(e); 616 } 617 } 618 WsEvent::IncomingPong => { 619 pong_deadline = None; 620 } 621 } 622 } 623 } 624 }; 625 626 session_cancel.cancel(); 627 drop(ws_sink); 628 drop(control_rx); 629 630 let reader_end = reader.await.unwrap_or(SessionEnd { 631 outcome: SessionOutcome::Empty, 632 error: None, 633 }); 634 if let Err(join) = processor.await { 635 warn!(?join, "frame processor task panicked"); 636 } 637 SessionEnd { 638 outcome: reader_end.outcome, 639 error: writer_error.or(reader_end.error), 640 } 641} 642 643async fn timed_send( 644 sink: &mut Box<dyn bobbin_runtime::WsSink>, 645 msg: WsMessage, 646) -> Result<(), IngestError> { 647 match tokio::time::timeout(SEND_TIMEOUT, sink.send(msg)).await { 648 Ok(Ok(())) => Ok(()), 649 Ok(Err(e)) => Err(IngestError::Network(e)), 650 Err(_elapsed) => Err(IngestError::SendTimeout(SEND_TIMEOUT)), 651 } 652} 653 654async fn reader_loop( 655 mut ws_stream: Box<dyn WsStream>, 656 frame_tx: tokio::sync::mpsc::Sender<HydrantFrame>, 657 control_tx: tokio::sync::mpsc::Sender<WsEvent>, 658 cancel: CancellationToken, 659) -> SessionEnd { 660 let mut outcome = SessionOutcome::Empty; 661 let mut held: VecDeque<HydrantFrame> = VecDeque::new(); 662 663 let error: Option<IngestError> = loop { 664 let step = if held.is_empty() { 665 tokio::select! { 666 biased; 667 _ = cancel.cancelled() => ReaderStep::Cancelled, 668 msg = ws_stream.next() => ReaderStep::WsMessage(msg), 669 } 670 } else if held.len() < READER_HOLD_LIMIT { 671 tokio::select! { 672 biased; 673 _ = cancel.cancelled() => ReaderStep::Cancelled, 674 permit_result = frame_tx.reserve() => match permit_result { 675 Ok(permit) => { 676 let frame = held 677 .pop_front() 678 .expect("non-empty held when reserve succeeds"); 679 permit.send(frame); 680 ReaderStep::Sent 681 } 682 Err(_) => ReaderStep::FrameSinkClosed, 683 }, 684 msg = ws_stream.next() => ReaderStep::WsMessage(msg), 685 } 686 } else { 687 tokio::select! { 688 biased; 689 _ = cancel.cancelled() => ReaderStep::Cancelled, 690 permit_result = frame_tx.reserve() => match permit_result { 691 Ok(permit) => { 692 let frame = held 693 .pop_front() 694 .expect("held at limit when reserve succeeds"); 695 permit.send(frame); 696 ReaderStep::Sent 697 } 698 Err(_) => ReaderStep::FrameSinkClosed, 699 }, 700 } 701 }; 702 703 match step { 704 ReaderStep::Cancelled => break None, 705 ReaderStep::FrameSinkClosed => break None, 706 ReaderStep::Sent => { 707 outcome = SessionOutcome::Progressed; 708 } 709 ReaderStep::WsMessage(msg) => { 710 let Some(msg) = msg else { 711 break None; 712 }; 713 let parsed = match msg { 714 Ok(m) => m, 715 Err(e) => break Some(IngestError::Network(e)), 716 }; 717 match parsed { 718 WsMessage::Text(text) => { 719 let frame = match classify_text_frame(&text) { 720 Ok(f) => f, 721 Err(e) => break Some(e), 722 }; 723 held.push_back(frame); 724 } 725 WsMessage::Binary(_) => { 726 debug!("hydrant sent unexpected binary frame, ignoring"); 727 } 728 WsMessage::Ping(payload) => { 729 if control_tx 730 .send(WsEvent::IncomingPing(payload)) 731 .await 732 .is_err() 733 { 734 break None; 735 } 736 } 737 WsMessage::Pong(_) => { 738 if control_tx.send(WsEvent::IncomingPong).await.is_err() { 739 break None; 740 } 741 } 742 WsMessage::Close { code, reason } => { 743 debug!(code, %reason, "hydrant closed stream"); 744 break None; 745 } 746 } 747 } 748 } 749 }; 750 SessionEnd { outcome, error } 751} 752 753enum ReaderStep { 754 Cancelled, 755 FrameSinkClosed, 756 Sent, 757 WsMessage(Option<Result<WsMessage, NetworkError>>), 758} 759 760fn classify_text_frame(text: &str) -> Result<HydrantFrame, IngestError> { 761 #[derive(serde::Deserialize)] 762 struct PeekType<'a> { 763 #[serde(rename = "type", borrow)] 764 kind: Option<std::borrow::Cow<'a, str>>, 765 } 766 let is_error_frame = serde_json::from_str::<PeekType>(text) 767 .ok() 768 .and_then(|p| p.kind) 769 .as_deref() 770 == Some("error"); 771 if is_error_frame { 772 return match serde_json::from_str::<HydrantStreamErrorFrame>(text) { 773 Ok(err_frame) => Err(classify_hydrant_error(err_frame)), 774 Err(decode_err) => Err(IngestError::Decode(decode_err)), 775 }; 776 } 777 serde_json::from_str::<HydrantFrame>(text).map_err(IngestError::Decode) 778} 779 780fn classify_hydrant_error(frame: HydrantStreamErrorFrame) -> IngestError { 781 let HydrantStreamErrorFrame { error, message } = frame; 782 let message = message.unwrap_or_default(); 783 match error.as_str() { 784 "ConsumerTooSlow" => IngestError::ConsumerTooSlow { message }, 785 _ => IngestError::HydrantStream { 786 code: error, 787 message, 788 }, 789 } 790} 791 792#[derive(Debug)] 793enum WsEvent { 794 IncomingPing(Bytes), 795 IncomingPong, 796} 797 798async fn wait_until(deadline: Option<Instant>, clock: &dyn Clock) { 799 match deadline { 800 Some(d) => clock.sleep_until(d).await, 801 None => std::future::pending::<()>().await, 802 } 803} 804 805#[derive(Clone, Copy, Debug, Eq, PartialEq)] 806enum Regime { 807 Replay, 808 Live, 809 NonRecord, 810} 811 812impl Regime { 813 fn as_str(self) -> &'static str { 814 match self { 815 Self::Replay => "replay", 816 Self::Live => "live", 817 Self::NonRecord => "non_record", 818 } 819 } 820} 821 822struct Pending { 823 cursor: HydrantCursor, 824 signal: PromotionSignal, 825 regime: Regime, 826 op: PendingOp, 827} 828 829enum PendingOp { 830 Noop, 831 ClearCache { 832 source: AtUri<DefaultStr>, 833 }, 834 Upsert { 835 source: AtUri<DefaultStr>, 836 nsid: Nsid<DefaultStr>, 837 parsed: Box<Record>, 838 bytes: Bytes, 839 cid: Option<Cid<DefaultStr>>, 840 edges: Vec<Edge>, 841 }, 842 Parked { 843 nsid: Nsid<DefaultStr>, 844 }, 845 Delete { 846 source: AtUri<DefaultStr>, 847 nsid: Nsid<DefaultStr>, 848 }, 849} 850 851struct Prepared { 852 pending: Pending, 853 prepare_start: Instant, 854 prepare_end: Instant, 855} 856 857struct Resolved { 858 pending: Pending, 859 prepare_start: Instant, 860 prepare_end: Instant, 861 resolve_start: Instant, 862 resolve_end: Instant, 863} 864 865fn pending_nsid(op: &PendingOp) -> Option<&Nsid<DefaultStr>> { 866 match op { 867 PendingOp::Upsert { nsid, .. } => Some(nsid), 868 PendingOp::Delete { nsid, .. } => Some(nsid), 869 PendingOp::Parked { nsid, .. } => Some(nsid), 870 PendingOp::Noop | PendingOp::ClearCache { .. } => None, 871 } 872} 873 874fn pending_edge_count(op: &PendingOp) -> u64 { 875 match op { 876 PendingOp::Upsert { edges, .. } => edges.len() as u64, 877 _ => 0, 878 } 879} 880 881async fn prep_stage<S: SearchSink + 'static>( 882 frame: HydrantFrame, 883 rt: IngestRuntime<S>, 884) -> Prepared { 885 let now = rt.clock.now_unix_micros(); 886 let prepare_start = rt.clock.now_instant(); 887 let ctx = rt.pipeline_ctx(); 888 let pending = prepare_frame(frame, &ctx, now).await; 889 let prepare_end = rt.clock.now_instant(); 890 Prepared { 891 pending, 892 prepare_start, 893 prepare_end, 894 } 895} 896 897async fn resolve_stage<S: SearchSink + 'static>( 898 staged: Prepared, 899 rt: IngestRuntime<S>, 900) -> Resolved { 901 let resolve_start = rt.clock.now_instant(); 902 let ctx = rt.pipeline_ctx(); 903 let pending = resolve_pending(staged.pending, &ctx).await; 904 let resolve_end = rt.clock.now_instant(); 905 Resolved { 906 pending, 907 prepare_start: staged.prepare_start, 908 prepare_end: staged.prepare_end, 909 resolve_start, 910 resolve_end, 911 } 912} 913 914async fn commit_stage<S: SearchSink + 'static>( 915 staged: Resolved, 916 rt: IngestRuntime<S>, 917 parallelism: usize, 918) { 919 let nsid = pending_nsid(&staged.pending.op).cloned(); 920 let edge_count = pending_edge_count(&staged.pending.op); 921 let regime = staged.pending.regime; 922 let cursor = staged.pending.cursor.raw(); 923 let Resolved { 924 pending, 925 prepare_start, 926 prepare_end, 927 resolve_start, 928 resolve_end, 929 } = staged; 930 let commit_start = rt.clock.now_instant(); 931 commit_pending( 932 pending, 933 &rt.store, 934 &rt.issue_states, 935 &rt.pull_statuses, 936 &rt.coverage, 937 &*rt.search, 938 &*rt.records, 939 &rt.resolver, 940 ) 941 .await; 942 let commit_end = rt.clock.now_instant(); 943 tracing::trace!( 944 target: "bobbin_ingest::stage", 945 cursor, 946 regime = regime.as_str(), 947 nsid = %nsid.as_ref().map(Nsid::as_str).unwrap_or(""), 948 edge_count, 949 prepare_us = prepare_end.duration_since(prepare_start).as_micros() as u64, 950 queue_resolve_wait_us = resolve_start.duration_since(prepare_end).as_micros() as u64, 951 resolve_us = resolve_end.duration_since(resolve_start).as_micros() as u64, 952 queue_commit_wait_us = commit_start.duration_since(resolve_end).as_micros() as u64, 953 commit_us = commit_end.duration_since(commit_start).as_micros() as u64, 954 total_us = commit_end.duration_since(prepare_start).as_micros() as u64, 955 parallelism, 956 "pipeline stage timings", 957 ); 958} 959 960async fn prepare_frame<S: SearchSink + 'static>( 961 frame: HydrantFrame, 962 ctx: &PipelineCtx<'_, S>, 963 now: UnixMicros, 964) -> Pending { 965 let cursor = HydrantCursor::new(frame.id); 966 let signal = promotion_signal(frame.record.as_ref(), now); 967 let regime = match frame.record.as_ref() { 968 Some(r) if r.live => Regime::Live, 969 Some(_) => Regime::Replay, 970 None => Regime::NonRecord, 971 }; 972 let op = match frame.kind { 973 FrameKind::Record => prepare_record(frame.record, ctx).await, 974 FrameKind::Identity | FrameKind::Account => PendingOp::Noop, 975 FrameKind::Other => { 976 debug!(id = frame.id, "ignoring unknown hydrant frame kind"); 977 PendingOp::Noop 978 } 979 }; 980 Pending { 981 cursor, 982 signal, 983 regime, 984 op, 985 } 986} 987 988async fn prepare_record<S: SearchSink + 'static>( 989 record: Option<RecordFrame>, 990 ctx: &PipelineCtx<'_, S>, 991) -> PendingOp { 992 let Some(record) = record else { 993 debug!("record-typed frame missing payload, skipping"); 994 return PendingOp::Noop; 995 }; 996 if !record.collection.as_ref().starts_with(TANGLED_PREFIX) { 997 return PendingOp::Noop; 998 } 999 let nsid = record.collection.clone(); 1000 let source = match build_source_uri(&record) { 1001 Ok(s) => s, 1002 Err(e) => { 1003 warn!(?e, "invalid frame source, dropping record"); 1004 return PendingOp::Noop; 1005 } 1006 }; 1007 match record.action { 1008 RecordAction::Create | RecordAction::Update => { 1009 evict_from_buffer(ctx.buffer, &source).await; 1010 let Some(raw) = record.record else { 1011 debug!(collection = %nsid, "create/update missing record body, clearing cache"); 1012 return PendingOp::ClearCache { source }; 1013 }; 1014 let raw_bytes = Bytes::copy_from_slice(raw.get().as_bytes()); 1015 let wire_bytes = match fallback_rfc3339(&record.rkey, &record.rev) 1016 .and_then(|fallback| synthesize_created_at(&raw_bytes, &fallback)) 1017 { 1018 Some(patched) => Bytes::from(patched), 1019 None => raw_bytes, 1020 }; 1021 let (parsed, bytes) = match decode_canon_or_upgrade_bytes( 1022 &record.collection, 1023 &wire_bytes, 1024 ctx.resolver, 1025 ) 1026 .await 1027 { 1028 Ok((parsed, canon_bytes)) => { 1029 let bytes = match canon_bytes { 1030 std::borrow::Cow::Borrowed(_) => wire_bytes, 1031 std::borrow::Cow::Owned(v) => Bytes::from(v), 1032 }; 1033 (parsed, bytes) 1034 } 1035 Err(ExtractError::UnknownCollection(name)) => { 1036 debug!(collection = %name, "unknown sh.tangled.* collection, clearing cache"); 1037 return PendingOp::ClearCache { source }; 1038 } 1039 Err(e) => { 1040 warn!(?e, collection = %record.collection, "record decode failed, clearing cache"); 1041 return PendingOp::ClearCache { source }; 1042 } 1043 }; 1044 if let Record::Repo(repo) = &parsed { 1045 if let Some(shadow) = ctx.shadow { 1046 shadow.note_observed(&record.did, &record.rkey).await; 1047 } 1048 let superseded = ctx 1049 .resolver 1050 .observe( 1051 record.did.clone(), 1052 record.rkey.clone(), 1053 repo.repo_did.clone(), 1054 ) 1055 .await; 1056 if let Some(prior) = superseded { 1057 let prior_uri = format!( 1058 "at://{}/sh.tangled.repo/{}", 1059 prior.owner.as_ref(), 1060 prior.rkey.as_ref(), 1061 ); 1062 if let Ok(prior_at_uri) = AtUri::<DefaultStr>::new_owned(&prior_uri) { 1063 evict_from_buffer(ctx.buffer, &prior_at_uri).await; 1064 ctx.store.remove_source(&prior_at_uri); 1065 ctx.records.remove(&prior_at_uri); 1066 ctx.search.remove(&prior_at_uri).await; 1067 } 1068 } 1069 if let Some(buffer) = ctx.buffer { 1070 let drained = buffer.take_observed(&record.did, &record.rkey).await; 1071 if !drained.is_empty() { 1072 finalize_drained(ctx, drained).await; 1073 } 1074 } 1075 if let Some(registry) = ctx.knot_registry { 1076 let host = KnotHostKey::new(repo.knot.as_ref()); 1077 match repo.repo_did.clone() { 1078 Some(repo_did) => registry.observe_repo(&host, repo_did), 1079 None => registry.observe_host(&host), 1080 } 1081 } 1082 } 1083 match acl_disposition(&parsed, ctx.knot_gate, ctx.knot_registry) { 1084 AclDisposition::NativeSkip => { 1085 if let Some(registry) = ctx.knot_registry { 1086 registry.forget_legacy_member(&source); 1087 } 1088 return PendingOp::Delete { source, nsid }; 1089 } 1090 AclDisposition::LegacyMember { host } => { 1091 if let Some(registry) = ctx.knot_registry { 1092 registry.observe_host(&host); 1093 registry.note_legacy_member(source.clone(), &host); 1094 } 1095 } 1096 AclDisposition::Other => {} 1097 } 1098 let edges = match parsed.extract_edges(&source) { 1099 Ok(es) => es, 1100 Err(e) => { 1101 warn!(?e, "edge extraction failed, clearing cache"); 1102 return PendingOp::ClearCache { source }; 1103 } 1104 }; 1105 let _ = ctx.store.intern_source(&source); 1106 PendingOp::Upsert { 1107 source, 1108 nsid, 1109 parsed: Box::new(parsed), 1110 bytes, 1111 cid: record.cid, 1112 edges, 1113 } 1114 } 1115 RecordAction::Delete => { 1116 evict_from_buffer(ctx.buffer, &source).await; 1117 if nsid.as_ref() == "sh.tangled.repo" { 1118 ctx.resolver.forget(&record.did, &record.rkey).await; 1119 } 1120 if nsid.as_ref() == "sh.tangled.knot.member" 1121 && let Some(registry) = ctx.knot_registry 1122 { 1123 registry.forget_legacy_member(&source); 1124 } 1125 PendingOp::Delete { source, nsid } 1126 } 1127 RecordAction::Other => { 1128 debug!(collection = %nsid, "ignoring unknown record action"); 1129 PendingOp::Noop 1130 } 1131 } 1132} 1133 1134enum AclDisposition { 1135 Other, 1136 NativeSkip, 1137 LegacyMember { host: KnotHostKey }, 1138} 1139 1140fn acl_disposition( 1141 parsed: &Record, 1142 gate: Option<&CapabilityGate>, 1143 registry: Option<&KnotRegistry>, 1144) -> AclDisposition { 1145 let Some(gate) = gate else { 1146 return AclDisposition::Other; 1147 }; 1148 match parsed { 1149 Record::KnotMember(member) => { 1150 let host = KnotHostKey::new(member.domain.as_ref()); 1151 if gate.is_native(&host) { 1152 AclDisposition::NativeSkip 1153 } else { 1154 AclDisposition::LegacyMember { host } 1155 } 1156 } 1157 Record::Collaborator(collaborator) => { 1158 let native = registry 1159 .and_then(|registry| registry.host_of_repo(&collaborator.repo)) 1160 .is_some_and(|host| gate.is_native(&host)); 1161 if native { 1162 AclDisposition::NativeSkip 1163 } else { 1164 AclDisposition::Other 1165 } 1166 } 1167 _ => AclDisposition::Other, 1168 } 1169} 1170 1171fn fallback_rfc3339( 1172 rkey: &Rkey<DefaultStr>, 1173 rev: &jacquard_common::types::tid::Tid, 1174) -> Option<String> { 1175 let tid = jacquard_common::types::tid::Tid::new(rkey.as_ref()) 1176 .ok() 1177 .unwrap_or_else(|| rev.clone()); 1178 let micros = i64::try_from(tid.timestamp()).ok()?; 1179 let dt = chrono::DateTime::<chrono::Utc>::from_timestamp_micros(micros)?; 1180 Some(dt.to_rfc3339_opts(chrono::SecondsFormat::Micros, true)) 1181} 1182 1183async fn evict_from_buffer(buffer: Option<&WarmingBuffer>, source: &AtUri<DefaultStr>) { 1184 if let Some(buffer) = buffer 1185 && !buffer.is_sealed() 1186 { 1187 buffer.evict_source(source).await; 1188 } 1189} 1190 1191async fn resolve_pending<S: SearchSink + 'static>( 1192 pending: Pending, 1193 ctx: &PipelineCtx<'_, S>, 1194) -> Pending { 1195 let Pending { 1196 cursor, 1197 signal, 1198 regime, 1199 op, 1200 } = pending; 1201 let op = match op { 1202 PendingOp::Upsert { 1203 source, 1204 nsid, 1205 parsed, 1206 bytes, 1207 cid, 1208 edges, 1209 } => { 1210 let pieces = UpsertPieces { 1211 source, 1212 nsid, 1213 parsed: *parsed, 1214 bytes, 1215 cid, 1216 edges, 1217 }; 1218 let pieces = match try_park_warming(ctx, cursor, pieces).await { 1219 ParkOutcome::Parked { nsid } => { 1220 return Pending { 1221 cursor, 1222 signal, 1223 regime, 1224 op: PendingOp::Parked { nsid }, 1225 }; 1226 } 1227 ParkOutcome::Passthrough(pieces) => *pieces, 1228 }; 1229 let UpsertPieces { 1230 source, 1231 nsid, 1232 parsed, 1233 bytes, 1234 cid, 1235 edges, 1236 } = pieces; 1237 let edges = normalize_subjects(edges, ctx.resolver, ctx.coverage, ctx.shadow).await; 1238 PendingOp::Upsert { 1239 source, 1240 nsid, 1241 parsed: Box::new(parsed), 1242 bytes, 1243 cid, 1244 edges, 1245 } 1246 } 1247 other => other, 1248 }; 1249 Pending { 1250 cursor, 1251 signal, 1252 regime, 1253 op, 1254 } 1255} 1256 1257struct UpsertPieces { 1258 source: AtUri<DefaultStr>, 1259 nsid: Nsid<DefaultStr>, 1260 parsed: Record, 1261 bytes: Bytes, 1262 cid: Option<Cid<DefaultStr>>, 1263 edges: Vec<Edge>, 1264} 1265 1266impl From<ParkedUpsert> for UpsertPieces { 1267 fn from(u: ParkedUpsert) -> Self { 1268 Self { 1269 source: u.source, 1270 nsid: u.nsid, 1271 parsed: u.parsed, 1272 bytes: u.bytes, 1273 cid: u.cid, 1274 edges: u.edges, 1275 } 1276 } 1277} 1278 1279enum ParkOutcome { 1280 Parked { nsid: Nsid<DefaultStr> }, 1281 Passthrough(Box<UpsertPieces>), 1282} 1283 1284async fn try_park_warming<S: SearchSink + 'static>( 1285 ctx: &PipelineCtx<'_, S>, 1286 cursor: HydrantCursor, 1287 pieces: UpsertPieces, 1288) -> ParkOutcome { 1289 let Some(buffer) = ctx.buffer else { 1290 return ParkOutcome::Passthrough(Box::new(pieces)); 1291 }; 1292 if ctx.coverage.snapshot().is_ready() || buffer.is_sealed() { 1293 return ParkOutcome::Passthrough(Box::new(pieces)); 1294 } 1295 let deps = collect_unresolved_deps(&pieces.edges, ctx.resolver).await; 1296 if deps.is_empty() { 1297 return ParkOutcome::Passthrough(Box::new(pieces)); 1298 } 1299 let nsid = pieces.nsid.clone(); 1300 let upsert = ParkedUpsert { 1301 cursor, 1302 source: pieces.source, 1303 nsid: pieces.nsid, 1304 parsed: pieces.parsed, 1305 bytes: pieces.bytes, 1306 cid: pieces.cid, 1307 edges: pieces.edges, 1308 }; 1309 let deps_for_shadow = ctx.shadow.is_some().then(|| deps.clone()); 1310 match buffer.try_park(upsert, deps).await { 1311 Ok(()) => { 1312 if let Some((shadow, noted)) = ctx.shadow.zip(deps_for_shadow) { 1313 let _ = futures::future::join_all(noted.into_iter().map(|dep| async move { 1314 shadow.note_unresolved(dep.owner, dep.rkey).await; 1315 })) 1316 .await; 1317 } 1318 ParkOutcome::Parked { nsid } 1319 } 1320 Err(returned) => ParkOutcome::Passthrough(Box::new(returned.into())), 1321 } 1322} 1323 1324async fn collect_unresolved_deps(edges: &[Edge], resolver: &RepoIdResolver) -> Vec<RepoIdent> { 1325 futures::stream::iter(edges) 1326 .fold(Vec::new(), |mut acc, edge| async move { 1327 let Some(uri) = edge.subject.as_uri() else { 1328 return acc; 1329 }; 1330 let Some((owner, rkey)) = parse_repo_subject_uri(uri) else { 1331 return acc; 1332 }; 1333 if resolver.cached_resolution(&owner, &rkey).await.is_some() { 1334 return acc; 1335 } 1336 let candidate = RepoIdent::new(owner, rkey); 1337 if !acc.contains(&candidate) { 1338 acc.push(candidate); 1339 } 1340 acc 1341 }) 1342 .await 1343} 1344 1345async fn finalize_drained<S: SearchSink + 'static>( 1346 ctx: &PipelineCtx<'_, S>, 1347 drained: Vec<ParkedUpsert>, 1348) { 1349 for upsert in drained { 1350 let ParkedUpsert { 1351 cursor: _, 1352 source, 1353 nsid: _, 1354 parsed, 1355 bytes, 1356 cid, 1357 edges, 1358 } = upsert; 1359 let edges = normalize_subjects(edges, ctx.resolver, ctx.coverage, None).await; 1360 cache_body(ctx.records, &source, cid, bytes); 1361 let outcome = upsert_record_indexes( 1362 ctx.store, 1363 ctx.issue_states, 1364 ctx.pull_statuses, 1365 &source, 1366 edges, 1367 &parsed, 1368 ); 1369 log_unknown_state_variant(outcome, &source); 1370 index_search(ctx.search, ctx.resolver, &source, parsed).await; 1371 } 1372} 1373 1374#[allow(clippy::too_many_arguments)] 1375async fn commit_pending<S: SearchSink>( 1376 pending: Pending, 1377 store: &EdgeStore, 1378 issue_states: &StateIndex<IssueStateKind>, 1379 pull_statuses: &StateIndex<PullStatusKind>, 1380 coverage: &CoverageWatch, 1381 search: &S, 1382 records: &dyn RecordStore, 1383 resolver: &RepoIdResolver, 1384) { 1385 let Pending { 1386 cursor, 1387 signal, 1388 regime: _, 1389 op, 1390 } = pending; 1391 match op { 1392 PendingOp::Noop | PendingOp::Parked { .. } => {} 1393 PendingOp::ClearCache { source } => records.remove(&source), 1394 PendingOp::Upsert { 1395 source, 1396 nsid: _, 1397 parsed, 1398 bytes, 1399 cid, 1400 edges, 1401 } => { 1402 cache_body(records, &source, cid, bytes); 1403 let outcome = 1404 upsert_record_indexes(store, issue_states, pull_statuses, &source, edges, &parsed); 1405 log_unknown_state_variant(outcome, &source); 1406 index_search(search, resolver, &source, *parsed).await; 1407 } 1408 PendingOp::Delete { source, nsid } => { 1409 delete_record_indexes(store, issue_states, pull_statuses, &source, &nsid); 1410 records.remove(&source); 1411 search.remove(&source).await; 1412 } 1413 } 1414 coverage.update(|c| c.advance(cursor).maybe_promote(signal)); 1415} 1416 1417async fn index_search<S: SearchSink>( 1418 search: &S, 1419 resolver: &RepoIdResolver, 1420 source: &AtUri<DefaultStr>, 1421 parsed: Record, 1422) { 1423 let Some(searchable) = SearchableRecord::try_from_record(parsed) else { 1424 return; 1425 }; 1426 let Some(searchable) = searchable.normalize(resolver).await else { 1427 return; 1428 }; 1429 search.upsert(searchable.to_search_doc(source)).await; 1430} 1431 1432#[cfg(test)] 1433#[allow(clippy::too_many_arguments)] 1434async fn handle_frame<S: SearchSink + 'static>( 1435 frame: HydrantFrame, 1436 store: &EdgeStore, 1437 issue_states: &StateIndex<IssueStateKind>, 1438 pull_statuses: &StateIndex<PullStatusKind>, 1439 coverage: &CoverageWatch, 1440 search: &S, 1441 records: &dyn RecordStore, 1442 resolver: &RepoIdResolver, 1443 clock: &dyn Clock, 1444 now: UnixMicros, 1445) { 1446 let _ = clock; 1447 let ctx = PipelineCtx { 1448 resolver, 1449 store, 1450 issue_states, 1451 pull_statuses, 1452 coverage, 1453 records, 1454 search, 1455 shadow: None, 1456 buffer: None, 1457 knot_registry: None, 1458 knot_gate: None, 1459 }; 1460 let pending = prepare_frame(frame, &ctx, now).await; 1461 let pending = resolve_pending(pending, &ctx).await; 1462 commit_pending( 1463 pending, 1464 store, 1465 issue_states, 1466 pull_statuses, 1467 coverage, 1468 search, 1469 records, 1470 resolver, 1471 ) 1472 .await; 1473} 1474 1475fn log_unknown_state_variant(outcome: ApplyOutcome, source: &AtUri<DefaultStr>) { 1476 if matches!(outcome, ApplyOutcome::UnknownVariant) { 1477 warn!( 1478 target: "bobbin_ingest::state_index", 1479 %source, 1480 "state record has unknown wire variant, skipping index update", 1481 ); 1482 } 1483} 1484 1485fn promotion_signal(record: Option<&RecordFrame>, now: UnixMicros) -> PromotionSignal { 1486 PromotionSignal { 1487 rev_micros: record.map(|r| r.rev.timestamp()), 1488 now_micros: now.raw(), 1489 skew_micros: READY_SKEW.as_micros() as u64, 1490 } 1491} 1492 1493fn cache_body( 1494 records: &dyn RecordStore, 1495 source: &AtUri<DefaultStr>, 1496 cid: Option<Cid<DefaultStr>>, 1497 bytes: Bytes, 1498) { 1499 match cid { 1500 Some(cid) => records.put( 1501 source.clone(), 1502 Arc::new(RecordBody { 1503 uri: source.clone(), 1504 cid, 1505 value: bytes, 1506 }), 1507 ), 1508 None => records.remove(source), 1509 } 1510} 1511 1512async fn normalize_subjects( 1513 edges: Vec<Edge>, 1514 resolver: &RepoIdResolver, 1515 coverage: &CoverageWatch, 1516 shadow: Option<&WarmingShadowBuffer>, 1517) -> Vec<Edge> { 1518 let warming = shadow.is_some() && !coverage.snapshot().is_ready(); 1519 futures::stream::iter(edges) 1520 .filter_map(|edge| async move { 1521 let Some(uri) = edge.subject.as_uri() else { 1522 return Some(edge); 1523 }; 1524 let Some((owner, rkey)) = parse_repo_subject_uri(uri) else { 1525 return Some(edge); 1526 }; 1527 if warming 1528 && let Some(shadow) = shadow 1529 && resolver.cached_resolution(&owner, &rkey).await.is_none() 1530 { 1531 shadow 1532 .note_unresolved(owner.clone(), rkey.clone()) 1533 .await; 1534 } 1535 match resolver.resolve(&owner, &rkey).await { 1536 Resolution::Mapped(repo_did) => Some(Edge { 1537 subject: SubjectRef::Did(repo_did), 1538 ..edge 1539 }), 1540 Resolution::NoRepoDid => { 1541 warn!( 1542 target: "bobbin_ingest::normalize", 1543 kind = %edge.kind, 1544 owner = owner.as_ref(), 1545 rkey = rkey.as_ref(), 1546 source = edge.source.as_ref(), 1547 "dropping edge: target repo has no repoDid, no canonical DID subject available", 1548 ); 1549 None 1550 } 1551 Resolution::Unresolvable => { 1552 warn!( 1553 target: "bobbin_ingest::normalize", 1554 kind = %edge.kind, 1555 owner = owner.as_ref(), 1556 rkey = rkey.as_ref(), 1557 source = edge.source.as_ref(), 1558 "dropping edge: repo unresolvable, rkey-form subject will not match bare-DID queries", 1559 ); 1560 None 1561 } 1562 } 1563 }) 1564 .collect() 1565 .await 1566} 1567 1568fn parse_repo_subject_uri(uri: &AtUri<DefaultStr>) -> Option<(Did<DefaultStr>, Rkey<DefaultStr>)> { 1569 let collection = uri.collection()?; 1570 if collection.as_ref() != "sh.tangled.repo" { 1571 return None; 1572 } 1573 let AtIdentifier::Did(authority) = uri.authority() else { 1574 return None; 1575 }; 1576 let rkey = uri.rkey()?; 1577 let owner = Did::new_owned(authority.as_ref()).ok()?; 1578 let rkey = Rkey::new_owned(rkey.as_ref()).ok()?; 1579 Some((owner, rkey)) 1580} 1581 1582fn build_source_uri(r: &RecordFrame) -> Result<AtUri<DefaultStr>, IngestError> { 1583 Ok(AtUri::from_parts_owned( 1584 r.did.as_ref(), 1585 r.collection.as_ref(), 1586 r.rkey.as_ref(), 1587 )?) 1588} 1589 1590#[cfg(test)] 1591mod tests { 1592 use super::*; 1593 use bobbin_edge_index::Coverage; 1594 use bobbin_record_lru::{CacheCapacity, LruRecordStore, NoopRecordStore, RecordStore}; 1595 use bobbin_runtime::{OsEntropy, RuntimeHasher, SystemClock, TungsteniteWs}; 1596 use bobbin_types::search::NoopSearchSink; 1597 use jacquard_common::types::nsid::Nsid; 1598 use jacquard_common::types::tid::Tid; 1599 use serde_json::json; 1600 1601 const VALID_CID: &str = "bafyreieqygohnz2zqyvtvktbjpvhutphobcmbsnt4q5lc36ri7vpcmoz4i"; 1602 1603 fn did_subj(s: &str) -> SubjectRef { 1604 SubjectRef::Did(Did::new_owned(s).unwrap()) 1605 } 1606 1607 fn uri_subj(s: &str) -> SubjectRef { 1608 SubjectRef::Uri(AtUri::new_owned(s).unwrap()) 1609 } 1610 1611 fn rkey(s: &str) -> Rkey<DefaultStr> { 1612 Rkey::new_owned(s).unwrap() 1613 } 1614 1615 #[allow(clippy::type_complexity)] 1616 fn fresh() -> ( 1617 Arc<EdgeStore>, 1618 Arc<StateIndex<IssueStateKind>>, 1619 Arc<StateIndex<PullStatusKind>>, 1620 Arc<CoverageWatch>, 1621 Arc<RepoIdResolver>, 1622 ) { 1623 ( 1624 Arc::new(EdgeStore::new(RuntimeHasher::default())), 1625 Arc::new(StateIndex::new(RuntimeHasher::default())), 1626 Arc::new(StateIndex::new(RuntimeHasher::default())), 1627 Arc::new(CoverageWatch::new()), 1628 Arc::new(RepoIdResolver::detached(RuntimeHasher::default())), 1629 ) 1630 } 1631 1632 fn now() -> UnixMicros { 1633 SystemClock::new().now_unix_micros() 1634 } 1635 1636 fn sys_clock() -> SystemClock { 1637 SystemClock::new() 1638 } 1639 1640 fn parse_frame(value: serde_json::Value) -> HydrantFrame { 1641 let text = serde_json::to_string(&value).expect("serialize fixture"); 1642 serde_json::from_str(&text).expect("deserialize fixture") 1643 } 1644 1645 fn fresh_tid() -> Tid { 1646 Tid::now_0() 1647 } 1648 1649 #[tokio::test] 1650 async fn ignores_non_tangled_collections() { 1651 let (store, issue_states, pull_statuses, cov, resolver) = fresh(); 1652 let frame: HydrantFrame = parse_frame(json!({ 1653 "id": 1, 1654 "type": "record", 1655 "record": { 1656 "live": false, 1657 "did": "did:plc:nel", 1658 "rev": fresh_tid().as_str(), 1659 "collection": "app.bsky.feed.post", 1660 "rkey": "abcabcabcabcz", 1661 "action": "create", 1662 "record": {"$type": "app.bsky.feed.post", "text": "hi"} 1663 } 1664 })); 1665 handle_frame( 1666 frame, 1667 &store, 1668 &issue_states, 1669 &pull_statuses, 1670 &cov, 1671 &NoopSearchSink, 1672 &NoopRecordStore, 1673 &resolver, 1674 &sys_clock(), 1675 now(), 1676 ) 1677 .await; 1678 assert_eq!(store.key_count(), 0); 1679 assert_eq!(cov.snapshot().events_processed(), 1); 1680 } 1681 1682 #[tokio::test] 1683 async fn native_knot_member_skipped_legacy_indexed() { 1684 use bobbin_knot_ingest::{CapabilityGate, KnotClient, KnotRegistry}; 1685 use wiremock::matchers::{method, path}; 1686 use wiremock::{Mock, MockServer, ResponseTemplate}; 1687 1688 let server = MockServer::start().await; 1689 Mock::given(method("GET")) 1690 .and(path("/xrpc/sh.tangled.knot.version")) 1691 .respond_with(ResponseTemplate::new(200).set_body_json(json!({ 1692 "version": "1.1.0", 1693 "capabilities": ["knot-acl"] 1694 }))) 1695 .mount(&server) 1696 .await; 1697 let url = url::Url::parse(&server.uri()).unwrap(); 1698 let native_host = format!("{}:{}", url.host_str().unwrap(), url.port().unwrap()); 1699 1700 let gate = CapabilityGate::new( 1701 KnotClient::with_default_http(true).unwrap(), 1702 Arc::new(SystemClock::new()), 1703 true, 1704 true, 1705 ); 1706 assert!(gate.has_knot_acl(&KnotHostKey::new(&native_host)).await); 1707 1708 let registry = KnotRegistry::new(); 1709 let (store, issue_states, pull_statuses, cov, resolver) = fresh(); 1710 let ctx = PipelineCtx { 1711 resolver: &resolver, 1712 store: &store, 1713 issue_states: &issue_states, 1714 pull_statuses: &pull_statuses, 1715 coverage: &cov, 1716 records: &NoopRecordStore, 1717 search: &NoopSearchSink, 1718 shadow: None, 1719 buffer: None, 1720 knot_registry: Some(&registry), 1721 knot_gate: Some(&gate), 1722 }; 1723 1724 let member_frame = |id: u64, rkey: &str, domain: &str| { 1725 parse_frame(json!({ 1726 "id": id, 1727 "type": "record", 1728 "record": { 1729 "live": false, 1730 "did": "did:plc:akshay", 1731 "rev": fresh_tid().as_str(), 1732 "collection": "sh.tangled.knot.member", 1733 "rkey": rkey, 1734 "action": "create", 1735 "record": { 1736 "$type": "sh.tangled.knot.member", 1737 "subject": "did:plc:boltless", 1738 "domain": domain, 1739 "createdAt": "2026-06-01T00:00:00Z" 1740 } 1741 } 1742 })) 1743 }; 1744 1745 let native = 1746 prepare_frame(member_frame(1, "aaaaaaaaaaaaz", &native_host), &ctx, now()).await; 1747 assert!( 1748 matches!(native.op, PendingOp::Delete { .. }), 1749 "member record for a native knot must be dropped" 1750 ); 1751 1752 let legacy = 1753 prepare_frame(member_frame(2, "bbbbbbbbbbbbz", "legacy.knot"), &ctx, now()).await; 1754 assert!( 1755 matches!(legacy.op, PendingOp::Upsert { .. }), 1756 "member record for a legacy knot must be ingested" 1757 ); 1758 assert!( 1759 registry.hosts().contains(&KnotHostKey::new("legacy.knot")), 1760 "a member record seeds host discovery even before any repo is seen" 1761 ); 1762 assert_eq!( 1763 registry 1764 .drain_legacy_members(&KnotHostKey::new("legacy.knot")) 1765 .len(), 1766 1, 1767 "legacy member edge is indexed for later purge once the knot upgrades" 1768 ); 1769 } 1770 1771 #[tokio::test] 1772 async fn create_then_delete_round_trips_a_star() { 1773 let (store, issue_states, pull_statuses, cov, resolver) = fresh(); 1774 let create: HydrantFrame = parse_frame(json!({ 1775 "id": 10, 1776 "type": "record", 1777 "record": { 1778 "live": false, 1779 "did": "did:plc:olaren", 1780 "rev": fresh_tid().as_str(), 1781 "collection": "sh.tangled.feed.star", 1782 "rkey": "abcabcabcabcz", 1783 "action": "create", 1784 "record": { 1785 "$type": "sh.tangled.feed.star", 1786 "createdAt": "2026-05-01T00:00:00Z", 1787 "subjectDid": "did:plc:abalone" 1788 } 1789 } 1790 })); 1791 handle_frame( 1792 create, 1793 &store, 1794 &issue_states, 1795 &pull_statuses, 1796 &cov, 1797 &NoopSearchSink, 1798 &NoopRecordStore, 1799 &resolver, 1800 &sys_clock(), 1801 now(), 1802 ) 1803 .await; 1804 let key = bobbin_types::ids::EdgeKey::new( 1805 Nsid::new_static("sh.tangled.feed.star").unwrap(), 1806 did_subj("did:plc:abalone"), 1807 ); 1808 assert_eq!(store.count(&key), 1); 1809 1810 let delete: HydrantFrame = parse_frame(json!({ 1811 "id": 11, 1812 "type": "record", 1813 "record": { 1814 "live": false, 1815 "did": "did:plc:olaren", 1816 "rev": fresh_tid().as_str(), 1817 "collection": "sh.tangled.feed.star", 1818 "rkey": "abcabcabcabcz", 1819 "action": "delete", 1820 "record": null 1821 } 1822 })); 1823 handle_frame( 1824 delete, 1825 &store, 1826 &issue_states, 1827 &pull_statuses, 1828 &cov, 1829 &NoopSearchSink, 1830 &NoopRecordStore, 1831 &resolver, 1832 &sys_clock(), 1833 now(), 1834 ) 1835 .await; 1836 assert_eq!(store.count(&key), 0); 1837 } 1838 1839 #[tokio::test] 1840 async fn update_replaces_prior_edges() { 1841 let (store, issue_states, pull_statuses, cov, resolver) = fresh(); 1842 let mk = |subject_did: &Did<DefaultStr>, id: u64| -> HydrantFrame { 1843 parse_frame(json!({ 1844 "id": id, 1845 "type": "record", 1846 "record": { 1847 "live": false, 1848 "did": "did:plc:olaren", 1849 "rev": fresh_tid().as_str(), 1850 "collection": "sh.tangled.feed.star", 1851 "rkey": "abcabcabcabcz", 1852 "action": "update", 1853 "record": { 1854 "$type": "sh.tangled.feed.star", 1855 "createdAt": "2026-05-01T00:00:00Z", 1856 "subjectDid": subject_did.as_ref() 1857 } 1858 } 1859 })) 1860 }; 1861 handle_frame( 1862 mk(&Did::new_owned("did:plc:abalone").unwrap(), 1), 1863 &store, 1864 &issue_states, 1865 &pull_statuses, 1866 &cov, 1867 &NoopSearchSink, 1868 &NoopRecordStore, 1869 &resolver, 1870 &sys_clock(), 1871 now(), 1872 ) 1873 .await; 1874 handle_frame( 1875 mk(&Did::new_owned("did:plc:uni").unwrap(), 2), 1876 &store, 1877 &issue_states, 1878 &pull_statuses, 1879 &cov, 1880 &NoopSearchSink, 1881 &NoopRecordStore, 1882 &resolver, 1883 &sys_clock(), 1884 now(), 1885 ) 1886 .await; 1887 1888 let kind = Nsid::new_static("sh.tangled.feed.star").unwrap(); 1889 let old = bobbin_types::ids::EdgeKey::new(kind.clone(), did_subj("did:plc:abalone")); 1890 let new = bobbin_types::ids::EdgeKey::new(kind, did_subj("did:plc:uni")); 1891 assert_eq!(store.count(&old), 0); 1892 assert_eq!(store.count(&new), 1); 1893 } 1894 1895 #[tokio::test] 1896 async fn prepare_frame_tags_regime_from_live_flag() { 1897 let (store, issue_states, pull_statuses, cov, resolver) = fresh(); 1898 let search = NoopSearchSink; 1899 let records = NoopRecordStore; 1900 let ctx = PipelineCtx { 1901 resolver: &resolver, 1902 store: &store, 1903 issue_states: &issue_states, 1904 pull_statuses: &pull_statuses, 1905 coverage: &cov, 1906 records: &records, 1907 search: &search, 1908 shadow: None, 1909 buffer: None, 1910 knot_registry: None, 1911 knot_gate: None, 1912 }; 1913 let mk = |live: bool| -> HydrantFrame { 1914 parse_frame(json!({ 1915 "id": 1, 1916 "type": "record", 1917 "record": { 1918 "live": live, 1919 "did": "did:plc:olaren", 1920 "rev": fresh_tid().as_str(), 1921 "collection": "sh.tangled.feed.star", 1922 "rkey": "abcabcabcabcz", 1923 "action": "create", 1924 "record": { 1925 "$type": "sh.tangled.feed.star", 1926 "createdAt": "2026-05-01T00:00:00Z", 1927 "subjectDid": "did:plc:abalone" 1928 } 1929 } 1930 })) 1931 }; 1932 let live_pending = prepare_frame(mk(true), &ctx, now()).await; 1933 assert_eq!(live_pending.regime, Regime::Live); 1934 let replay_pending = prepare_frame(mk(false), &ctx, now()).await; 1935 assert_eq!(replay_pending.regime, Regime::Replay); 1936 1937 let identity: HydrantFrame = parse_frame(json!({ 1938 "id": 9, 1939 "type": "identity", 1940 })); 1941 let id_pending = prepare_frame(identity, &ctx, now()).await; 1942 assert_eq!(id_pending.regime, Regime::NonRecord); 1943 } 1944 1945 #[tokio::test] 1946 async fn create_with_cid_warms_record_lru() { 1947 let (store, issue_states, pull_statuses, cov, resolver) = fresh(); 1948 let lru = LruRecordStore::new(CacheCapacity::from_bytes(64 * 1024)); 1949 let frame: HydrantFrame = parse_frame(json!({ 1950 "id": 1, 1951 "type": "record", 1952 "record": { 1953 "live": false, 1954 "did": "did:plc:olaren", 1955 "rev": fresh_tid().as_str(), 1956 "collection": "sh.tangled.feed.star", 1957 "rkey": "abcabcabcabcz", 1958 "action": "create", 1959 "cid": VALID_CID, 1960 "record": { 1961 "$type": "sh.tangled.feed.star", 1962 "createdAt": "2026-05-01T00:00:00Z", 1963 "subjectDid": "did:plc:abalone" 1964 } 1965 } 1966 })); 1967 let source = 1968 AtUri::new_owned("at://did:plc:olaren/sh.tangled.feed.star/abcabcabcabcz").unwrap(); 1969 handle_frame( 1970 frame, 1971 &store, 1972 &issue_states, 1973 &pull_statuses, 1974 &cov, 1975 &NoopSearchSink, 1976 &lru, 1977 &resolver, 1978 &sys_clock(), 1979 now(), 1980 ) 1981 .await; 1982 let cached = lru.get(&source).expect("hydrant cid must seed the lru"); 1983 assert_eq!(cached.cid.as_ref(), VALID_CID); 1984 let parsed: serde_json::Value = serde_json::from_slice(&cached.value).unwrap(); 1985 assert_eq!( 1986 parsed["subject"]["did"], "did:plc:abalone", 1987 "legacy wire is upgraded to canon shape before caching so downstream readers see canonical fields" 1988 ); 1989 } 1990 1991 #[tokio::test] 1992 async fn create_without_cid_clears_record_lru() { 1993 let (store, issue_states, pull_statuses, cov, resolver) = fresh(); 1994 let lru = LruRecordStore::new(CacheCapacity::from_bytes(64 * 1024)); 1995 let source = 1996 AtUri::new_owned("at://did:plc:olaren/sh.tangled.feed.star/abcabcabcabcz").unwrap(); 1997 let cid: Cid<DefaultStr> = VALID_CID.parse().unwrap(); 1998 lru.put( 1999 source.clone(), 2000 Arc::new(RecordBody { 2001 uri: source.clone(), 2002 cid, 2003 value: bytes::Bytes::from_static(b"{\"stale\":true}"), 2004 }), 2005 ); 2006 let frame: HydrantFrame = parse_frame(json!({ 2007 "id": 1, 2008 "type": "record", 2009 "record": { 2010 "live": false, 2011 "did": "did:plc:olaren", 2012 "rev": fresh_tid().as_str(), 2013 "collection": "sh.tangled.feed.star", 2014 "rkey": "abcabcabcabcz", 2015 "action": "update", 2016 "record": { 2017 "$type": "sh.tangled.feed.star", 2018 "createdAt": "2026-05-01T00:00:00Z", 2019 "subjectDid": "did:plc:abalone" 2020 } 2021 } 2022 })); 2023 handle_frame( 2024 frame, 2025 &store, 2026 &issue_states, 2027 &pull_statuses, 2028 &cov, 2029 &NoopSearchSink, 2030 &lru, 2031 &resolver, 2032 &sys_clock(), 2033 now(), 2034 ) 2035 .await; 2036 assert!( 2037 lru.get(&source).is_none(), 2038 "missing cid means we cannot trust the body, so the lru must be cleared", 2039 ); 2040 } 2041 2042 #[tokio::test] 2043 async fn delete_evicts_record_lru() { 2044 let (store, issue_states, pull_statuses, cov, resolver) = fresh(); 2045 let lru = LruRecordStore::new(CacheCapacity::from_bytes(64 * 1024)); 2046 let source = 2047 AtUri::new_owned("at://did:plc:olaren/sh.tangled.feed.star/abcabcabcabcz").unwrap(); 2048 let cid: Cid<DefaultStr> = VALID_CID.parse().unwrap(); 2049 lru.put( 2050 source.clone(), 2051 Arc::new(RecordBody { 2052 uri: source.clone(), 2053 cid, 2054 value: bytes::Bytes::from_static(b"{}"), 2055 }), 2056 ); 2057 let frame: HydrantFrame = parse_frame(json!({ 2058 "id": 1, 2059 "type": "record", 2060 "record": { 2061 "live": false, 2062 "did": "did:plc:olaren", 2063 "rev": fresh_tid().as_str(), 2064 "collection": "sh.tangled.feed.star", 2065 "rkey": "abcabcabcabcz", 2066 "action": "delete", 2067 "record": null 2068 } 2069 })); 2070 handle_frame( 2071 frame, 2072 &store, 2073 &issue_states, 2074 &pull_statuses, 2075 &cov, 2076 &NoopSearchSink, 2077 &lru, 2078 &resolver, 2079 &sys_clock(), 2080 now(), 2081 ) 2082 .await; 2083 assert!(lru.get(&source).is_none()); 2084 } 2085 2086 #[tokio::test] 2087 async fn live_recent_event_promotes_coverage_to_ready() { 2088 let (store, issue_states, pull_statuses, cov, resolver) = fresh(); 2089 assert!(!cov.snapshot().is_ready()); 2090 let frame: HydrantFrame = parse_frame(json!({ 2091 "id": 99, 2092 "type": "record", 2093 "record": { 2094 "live": true, 2095 "did": "did:plc:olaren", 2096 "rev": fresh_tid().as_str(), 2097 "collection": "sh.tangled.feed.star", 2098 "rkey": "abcabcabcabcz", 2099 "action": "create", 2100 "record": { 2101 "$type": "sh.tangled.feed.star", 2102 "createdAt": "2026-05-01T00:00:00Z", 2103 "subjectDid": "did:plc:abalone" 2104 } 2105 } 2106 })); 2107 handle_frame( 2108 frame, 2109 &store, 2110 &issue_states, 2111 &pull_statuses, 2112 &cov, 2113 &NoopSearchSink, 2114 &NoopRecordStore, 2115 &resolver, 2116 &sys_clock(), 2117 now(), 2118 ) 2119 .await; 2120 assert!(cov.snapshot().is_ready()); 2121 assert_eq!(cov.snapshot().last_cursor(), HydrantCursor::new(99)); 2122 } 2123 2124 #[tokio::test] 2125 async fn live_but_stale_rev_does_not_promote() { 2126 let (store, issue_states, pull_statuses, cov, resolver) = fresh(); 2127 let stale_tid = Tid::from_time(1_000_000, 0); 2128 let frame: HydrantFrame = parse_frame(json!({ 2129 "id": 7, 2130 "type": "record", 2131 "record": { 2132 "live": true, 2133 "did": "did:plc:olaren", 2134 "rev": stale_tid.as_str(), 2135 "collection": "sh.tangled.feed.star", 2136 "rkey": "abcabcabcabcz", 2137 "action": "create", 2138 "record": { 2139 "$type": "sh.tangled.feed.star", 2140 "createdAt": "2026-05-01T00:00:00Z", 2141 "subjectDid": "did:plc:abalone" 2142 } 2143 } 2144 })); 2145 handle_frame( 2146 frame, 2147 &store, 2148 &issue_states, 2149 &pull_statuses, 2150 &cov, 2151 &NoopSearchSink, 2152 &NoopRecordStore, 2153 &resolver, 2154 &sys_clock(), 2155 now(), 2156 ) 2157 .await; 2158 assert!(!cov.snapshot().is_ready()); 2159 assert!(matches!(cov.snapshot(), Coverage::Warming { .. })); 2160 } 2161 2162 #[test] 2163 fn classify_consumer_too_slow_frame_returns_typed_variant() { 2164 let text = r#"{"type":"error","error":"ConsumerTooSlow","message":"stream socket send blocked for at least 30 seconds"}"#; 2165 match classify_text_frame(text) { 2166 Err(IngestError::ConsumerTooSlow { message }) => { 2167 assert!( 2168 message.contains("30 seconds"), 2169 "message field preserved verbatim, got: {message}" 2170 ); 2171 } 2172 other => panic!("expected ConsumerTooSlow variant, got: {other:?}"), 2173 } 2174 } 2175 2176 #[test] 2177 fn classify_unknown_hydrant_error_falls_back_to_generic_variant() { 2178 let text = r#"{"type":"error","error":"NewFutureCode","message":"some new failure mode"}"#; 2179 match classify_text_frame(text) { 2180 Err(IngestError::HydrantStream { code, message }) => { 2181 assert_eq!(code, "NewFutureCode"); 2182 assert_eq!(message, "some new failure mode"); 2183 } 2184 other => panic!("expected HydrantStream variant, got: {other:?}"), 2185 } 2186 } 2187 2188 #[test] 2189 fn classify_error_frame_without_message_uses_empty_string() { 2190 let text = r#"{"type":"error","error":"ConsumerTooSlow"}"#; 2191 match classify_text_frame(text) { 2192 Err(IngestError::ConsumerTooSlow { message }) => assert!(message.is_empty()), 2193 other => panic!("expected ConsumerTooSlow with empty message, got: {other:?}"), 2194 } 2195 } 2196 2197 #[test] 2198 fn classify_normal_record_frame_unchanged() { 2199 let text = r#"{"id":42,"type":"record"}"#; 2200 let frame = classify_text_frame(text).expect("normal record frame must parse"); 2201 assert_eq!(frame.id, 42); 2202 assert_eq!(frame.kind, FrameKind::Record); 2203 } 2204 2205 #[test] 2206 fn classify_garbage_object_returns_decode_error() { 2207 let text = r#"{"random":"object","without":"required fields"}"#; 2208 match classify_text_frame(text) { 2209 Err(IngestError::Decode(_)) => {} 2210 other => panic!("expected Decode error, got: {other:?}"), 2211 } 2212 } 2213 2214 struct ScriptedWsStream { 2215 messages: std::collections::VecDeque<Result<WsMessage, NetworkError>>, 2216 } 2217 2218 impl WsStream for ScriptedWsStream { 2219 fn next<'a>(&'a mut self) -> bobbin_runtime::WsMessageFuture<'a> { 2220 let msg = self.messages.pop_front(); 2221 Box::pin(async move { msg }) 2222 } 2223 } 2224 2225 #[tokio::test] 2226 async fn reader_loop_surfaces_consumer_too_slow_from_error_frame() { 2227 let mut messages = std::collections::VecDeque::new(); 2228 messages.push_back(Ok(WsMessage::Text( 2229 r#"{"type":"error","error":"ConsumerTooSlow","message":"stream delivery blocked"}"# 2230 .to_owned(), 2231 ))); 2232 let stream: Box<dyn WsStream> = Box::new(ScriptedWsStream { messages }); 2233 let (frame_tx, _frame_rx) = tokio::sync::mpsc::channel::<HydrantFrame>(FRAME_CHANNEL_DEPTH); 2234 let (control_tx, _control_rx) = 2235 tokio::sync::mpsc::channel::<WsEvent>(CONTROL_CHANNEL_DEPTH); 2236 let cancel = CancellationToken::new(); 2237 2238 let end = reader_loop(stream, frame_tx, control_tx, cancel).await; 2239 2240 match end.error { 2241 Some(IngestError::ConsumerTooSlow { message }) => { 2242 assert_eq!(message, "stream delivery blocked"); 2243 } 2244 other => panic!("expected ConsumerTooSlow SessionEnd error, got: {other:?}"), 2245 } 2246 assert_eq!(end.outcome, SessionOutcome::Empty); 2247 } 2248 2249 #[tokio::test] 2250 async fn reader_loop_surfaces_unknown_hydrant_error_distinctly() { 2251 let mut messages = std::collections::VecDeque::new(); 2252 messages.push_back(Ok(WsMessage::Text( 2253 r#"{"type":"error","error":"NewFutureCode","message":"new mode"}"#.to_owned(), 2254 ))); 2255 let stream: Box<dyn WsStream> = Box::new(ScriptedWsStream { messages }); 2256 let (frame_tx, _frame_rx) = tokio::sync::mpsc::channel::<HydrantFrame>(FRAME_CHANNEL_DEPTH); 2257 let (control_tx, _control_rx) = 2258 tokio::sync::mpsc::channel::<WsEvent>(CONTROL_CHANNEL_DEPTH); 2259 let cancel = CancellationToken::new(); 2260 2261 let end = reader_loop(stream, frame_tx, control_tx, cancel).await; 2262 2263 match end.error { 2264 Some(IngestError::HydrantStream { code, message }) => { 2265 assert_eq!(code, "NewFutureCode"); 2266 assert_eq!(message, "new mode"); 2267 } 2268 other => panic!("expected HydrantStream SessionEnd error, got: {other:?}"), 2269 } 2270 } 2271 2272 struct ChannelWsStream { 2273 rx: tokio::sync::mpsc::Receiver<Result<WsMessage, NetworkError>>, 2274 } 2275 2276 impl WsStream for ChannelWsStream { 2277 fn next<'a>(&'a mut self) -> bobbin_runtime::WsMessageFuture<'a> { 2278 Box::pin(async move { self.rx.recv().await }) 2279 } 2280 } 2281 2282 fn star_frame_text(id: u64, rkey: &Rkey<DefaultStr>) -> String { 2283 json!({ 2284 "id": id, 2285 "type": "record", 2286 "record": { 2287 "live": false, 2288 "did": "did:plc:olaren", 2289 "rev": fresh_tid().as_str(), 2290 "collection": "sh.tangled.feed.star", 2291 "rkey": rkey.as_ref(), 2292 "action": "create", 2293 "record": { 2294 "$type": "sh.tangled.feed.star", 2295 "createdAt": "2026-05-01T00:00:00Z", 2296 "subjectDid": "did:plc:abalone" 2297 } 2298 } 2299 }) 2300 .to_string() 2301 } 2302 2303 #[tokio::test] 2304 async fn pong_forwards_promptly_when_frame_channel_is_full() { 2305 let (ws_tx, ws_rx) = tokio::sync::mpsc::channel::<Result<WsMessage, NetworkError>>(8); 2306 let stream: Box<dyn WsStream> = Box::new(ChannelWsStream { rx: ws_rx }); 2307 let (frame_tx, frame_rx) = tokio::sync::mpsc::channel::<HydrantFrame>(1); 2308 let (control_tx, mut control_rx) = 2309 tokio::sync::mpsc::channel::<WsEvent>(CONTROL_CHANNEL_DEPTH); 2310 let cancel = CancellationToken::new(); 2311 2312 let prefill: HydrantFrame = parse_frame(json!({ 2313 "id": 0, 2314 "type": "record", 2315 "record": { 2316 "live": false, 2317 "did": "did:plc:olaren", 2318 "rev": fresh_tid().as_str(), 2319 "collection": "sh.tangled.feed.star", 2320 "rkey": "prefilrkey001", 2321 "action": "create", 2322 "record": { 2323 "$type": "sh.tangled.feed.star", 2324 "createdAt": "2026-05-01T00:00:00Z", 2325 "subjectDid": "did:plc:abalone" 2326 } 2327 } 2328 })); 2329 frame_tx 2330 .try_send(prefill) 2331 .expect("depth-1 frame channel must accept the prefill"); 2332 2333 let reader_handle = tokio::spawn(reader_loop( 2334 stream, 2335 frame_tx.clone(), 2336 control_tx.clone(), 2337 cancel.clone(), 2338 )); 2339 2340 ws_tx 2341 .send(Ok(WsMessage::Text(star_frame_text( 2342 1, 2343 &rkey("starrkeyaa001"), 2344 )))) 2345 .await 2346 .unwrap(); 2347 ws_tx.send(Ok(WsMessage::Pong(Bytes::new()))).await.unwrap(); 2348 2349 let pong_event = tokio::time::timeout(Duration::from_millis(500), control_rx.recv()) 2350 .await 2351 .expect("pong must surface inside 500ms even when frame_tx is saturated. A reader blocked on a frame send would never poll the next ws message") 2352 .expect("control_tx was not closed"); 2353 assert!( 2354 matches!(pong_event, WsEvent::IncomingPong), 2355 "first control event must be the pong, not a held text", 2356 ); 2357 2358 let too_slow = 2359 "{\"type\":\"error\",\"error\":\"ConsumerTooSlow\",\"message\":\"saturated\"}" 2360 .to_string(); 2361 ws_tx.send(Ok(WsMessage::Text(too_slow))).await.unwrap(); 2362 2363 let end = tokio::time::timeout(Duration::from_secs(1), reader_handle) 2364 .await 2365 .expect("reader must exit promptly once ConsumerTooSlow is read") 2366 .expect("reader task must not panic"); 2367 match end.error { 2368 Some(IngestError::ConsumerTooSlow { message }) => assert_eq!(message, "saturated"), 2369 other => panic!("expected ConsumerTooSlow disconnect, got: {other:?}"), 2370 } 2371 2372 drop(frame_rx); 2373 } 2374 2375 #[tokio::test] 2376 async fn reader_drains_held_frames_once_processor_catches_up() { 2377 let (ws_tx, ws_rx) = tokio::sync::mpsc::channel::<Result<WsMessage, NetworkError>>(8); 2378 let stream: Box<dyn WsStream> = Box::new(ChannelWsStream { rx: ws_rx }); 2379 let (frame_tx, mut frame_rx) = tokio::sync::mpsc::channel::<HydrantFrame>(1); 2380 let (control_tx, _control_rx) = 2381 tokio::sync::mpsc::channel::<WsEvent>(CONTROL_CHANNEL_DEPTH); 2382 let cancel = CancellationToken::new(); 2383 2384 let prefill: HydrantFrame = parse_frame(json!({ 2385 "id": 0, 2386 "type": "record", 2387 "record": { 2388 "live": false, 2389 "did": "did:plc:olaren", 2390 "rev": fresh_tid().as_str(), 2391 "collection": "sh.tangled.feed.star", 2392 "rkey": "prefilrkey001", 2393 "action": "create", 2394 "record": { 2395 "$type": "sh.tangled.feed.star", 2396 "createdAt": "2026-05-01T00:00:00Z", 2397 "subjectDid": "did:plc:abalone" 2398 } 2399 } 2400 })); 2401 frame_tx.try_send(prefill).unwrap(); 2402 2403 let reader_handle = tokio::spawn(reader_loop( 2404 stream, 2405 frame_tx.clone(), 2406 control_tx, 2407 cancel.clone(), 2408 )); 2409 2410 ws_tx 2411 .send(Ok(WsMessage::Text(star_frame_text( 2412 1, 2413 &rkey("heldrkeyaa001"), 2414 )))) 2415 .await 2416 .unwrap(); 2417 ws_tx 2418 .send(Ok(WsMessage::Text(star_frame_text( 2419 2, 2420 &rkey("heldrkeyaa002"), 2421 )))) 2422 .await 2423 .unwrap(); 2424 2425 let _drained_prefill = frame_rx.recv().await.expect("prefilled frame drains"); 2426 let first = tokio::time::timeout(Duration::from_millis(500), frame_rx.recv()) 2427 .await 2428 .expect("first held frame must reach frame_rx after slot opens") 2429 .expect("frame_tx still open"); 2430 assert_eq!(first.id, 1); 2431 let second = tokio::time::timeout(Duration::from_millis(500), frame_rx.recv()) 2432 .await 2433 .expect("second held frame must reach frame_rx after slot opens") 2434 .expect("frame_tx still open"); 2435 assert_eq!(second.id, 2); 2436 2437 cancel.cancel(); 2438 let _ = tokio::time::timeout(Duration::from_secs(1), reader_handle) 2439 .await 2440 .expect("reader must stop after cancel"); 2441 } 2442 2443 #[test] 2444 fn first_connect_uses_configured_start_cursor() { 2445 let start = HydrantCursor::new(42); 2446 assert_eq!(next_connect_cursor(Coverage::default(), start), start); 2447 } 2448 2449 #[test] 2450 fn reconnect_resumes_strictly_after_last_seen() { 2451 let snap = Coverage::default().advance(HydrantCursor::new(7)); 2452 assert_eq!( 2453 next_connect_cursor(snap, HydrantCursor::new(0)), 2454 HydrantCursor::new(8), 2455 ); 2456 } 2457 2458 #[test] 2459 fn reconnect_overrides_configured_start() { 2460 let snap = Coverage::default().advance(HydrantCursor::new(100)); 2461 assert_eq!( 2462 next_connect_cursor(snap, HydrantCursor::new(50)), 2463 HydrantCursor::new(101), 2464 ); 2465 } 2466 2467 #[test] 2468 fn first_connect_uses_start_even_when_first_frame_id_would_be_zero() { 2469 let start = HydrantCursor::new(7); 2470 let snap = Coverage::default(); 2471 assert_eq!(snap.last_cursor(), HydrantCursor::new(0)); 2472 assert_eq!(snap.events_processed(), 0); 2473 assert_eq!(next_connect_cursor(snap, start), start); 2474 } 2475 2476 #[test] 2477 fn reconnect_after_processing_id_zero_advances_to_one() { 2478 let snap = Coverage::default().advance(HydrantCursor::new(0)); 2479 assert_eq!(snap.events_processed(), 1); 2480 assert_eq!( 2481 next_connect_cursor(snap, HydrantCursor::new(99)), 2482 HydrantCursor::new(1), 2483 "events_processed disambiguates 'never seen' from 'saw id 0'", 2484 ); 2485 } 2486 2487 #[tokio::test] 2488 async fn identity_frame_advances_cursor_only() { 2489 let (store, issue_states, pull_statuses, cov, resolver) = fresh(); 2490 let frame: HydrantFrame = parse_frame(json!({ 2491 "id": 5, 2492 "type": "identity", 2493 "identity": { 2494 "did": "did:plc:olaren", 2495 "handle": "olaren.dev" 2496 } 2497 })); 2498 handle_frame( 2499 frame, 2500 &store, 2501 &issue_states, 2502 &pull_statuses, 2503 &cov, 2504 &NoopSearchSink, 2505 &NoopRecordStore, 2506 &resolver, 2507 &sys_clock(), 2508 now(), 2509 ) 2510 .await; 2511 assert_eq!(store.key_count(), 0); 2512 assert_eq!(cov.snapshot().last_cursor(), HydrantCursor::new(5)); 2513 assert!(!cov.snapshot().is_ready()); 2514 } 2515 2516 #[tokio::test] 2517 async fn account_frame_advances_cursor_only() { 2518 let (store, issue_states, pull_statuses, cov, resolver) = fresh(); 2519 let frame: HydrantFrame = parse_frame(json!({ 2520 "id": 6, 2521 "type": "account", 2522 "account": {"did": "did:plc:olaren", "active": true} 2523 })); 2524 handle_frame( 2525 frame, 2526 &store, 2527 &issue_states, 2528 &pull_statuses, 2529 &cov, 2530 &NoopSearchSink, 2531 &NoopRecordStore, 2532 &resolver, 2533 &sys_clock(), 2534 now(), 2535 ) 2536 .await; 2537 assert_eq!(store.key_count(), 0); 2538 assert_eq!(cov.snapshot().last_cursor(), HydrantCursor::new(6)); 2539 } 2540 2541 #[tokio::test] 2542 async fn unknown_frame_kind_advances_cursor_without_panic() { 2543 let (store, issue_states, pull_statuses, cov, resolver) = fresh(); 2544 let frame: HydrantFrame = parse_frame(json!({"id": 8, "type": "future_event"})); 2545 assert_eq!(frame.kind, FrameKind::Other); 2546 handle_frame( 2547 frame, 2548 &store, 2549 &issue_states, 2550 &pull_statuses, 2551 &cov, 2552 &NoopSearchSink, 2553 &NoopRecordStore, 2554 &resolver, 2555 &sys_clock(), 2556 now(), 2557 ) 2558 .await; 2559 assert_eq!(cov.snapshot().last_cursor(), HydrantCursor::new(8)); 2560 } 2561 2562 fn fresh_runtime(cancel: CancellationToken) -> IngestRuntime<NoopSearchSink> { 2563 IngestRuntime { 2564 store: Arc::new(EdgeStore::new(RuntimeHasher::default())), 2565 issue_states: Arc::new(StateIndex::new(RuntimeHasher::default())), 2566 pull_statuses: Arc::new(StateIndex::new(RuntimeHasher::default())), 2567 coverage: Arc::new(CoverageWatch::new()), 2568 search: Arc::new(NoopSearchSink), 2569 records: Arc::new(NoopRecordStore) as Arc<dyn RecordStore>, 2570 resolver: Arc::new(RepoIdResolver::detached(RuntimeHasher::default())), 2571 clock: Arc::new(SystemClock::new()), 2572 entropy: Arc::new(OsEntropy), 2573 ws: TungsteniteWs::shared(), 2574 cancel, 2575 disconnects: None, 2576 warming_shadow: None, 2577 warming_buffer: None, 2578 knot_registry: None, 2579 knot_gate: None, 2580 } 2581 } 2582 2583 #[tokio::test(start_paused = true)] 2584 async fn cancel_token_short_circuits_reconnect_sleep() { 2585 let cfg = IngestConfig::new(Url::parse("ws://127.0.0.1:1").unwrap()); 2586 let cancel = CancellationToken::new(); 2587 let runtime = fresh_runtime(cancel.clone()); 2588 let task = tokio::spawn(async move { run(cfg, runtime).await }); 2589 tokio::time::advance(Duration::from_millis(10)).await; 2590 cancel.cancel(); 2591 let outcome = tokio::time::timeout(Duration::from_secs(1), task) 2592 .await 2593 .expect("ingest must stop within timeout once cancel fires"); 2594 assert!(matches!(outcome, Ok(Ok(()))), "got {outcome:?}"); 2595 } 2596 2597 #[test] 2598 fn jittered_stays_within_one_quarter_of_base() { 2599 let base = Duration::from_secs(1); 2600 let cap = base + Duration::from_millis(250); 2601 let entropy = OsEntropy; 2602 (0..50).for_each(|_| { 2603 let j = jittered(base, &entropy); 2604 assert!(j >= base, "jitter must not undershoot"); 2605 assert!(j <= cap, "jitter must not exceed +25%, got {:?}", j); 2606 }); 2607 } 2608 2609 #[tokio::test] 2610 async fn star_after_observed_repo_keys_on_repo_did() { 2611 let (store, issue_states, pull_statuses, cov, resolver) = fresh(); 2612 let repo: HydrantFrame = parse_frame(json!({ 2613 "id": 1, 2614 "type": "record", 2615 "record": { 2616 "live": false, 2617 "did": "did:plc:nel", 2618 "rev": fresh_tid().as_str(), 2619 "collection": "sh.tangled.repo", 2620 "rkey": "abcabcabcabcz", 2621 "action": "create", 2622 "record": { 2623 "$type": "sh.tangled.repo", 2624 "createdAt": "2026-05-01T00:00:00Z", 2625 "knot": "oyster.cafe", 2626 "name": "abalone", 2627 "repoDid": "did:plc:abalone" 2628 } 2629 } 2630 })); 2631 handle_frame( 2632 repo, 2633 &store, 2634 &issue_states, 2635 &pull_statuses, 2636 &cov, 2637 &NoopSearchSink, 2638 &NoopRecordStore, 2639 &resolver, 2640 &sys_clock(), 2641 now(), 2642 ) 2643 .await; 2644 2645 let star: HydrantFrame = parse_frame(json!({ 2646 "id": 2, 2647 "type": "record", 2648 "record": { 2649 "live": false, 2650 "did": "did:plc:olaren", 2651 "rev": fresh_tid().as_str(), 2652 "collection": "sh.tangled.feed.star", 2653 "rkey": "abcabcabcabcz", 2654 "action": "create", 2655 "record": { 2656 "$type": "sh.tangled.feed.star", 2657 "createdAt": "2026-05-01T00:00:00Z", 2658 "subject": "at://did:plc:nel/sh.tangled.repo/abcabcabcabcz" 2659 } 2660 } 2661 })); 2662 handle_frame( 2663 star, 2664 &store, 2665 &issue_states, 2666 &pull_statuses, 2667 &cov, 2668 &NoopSearchSink, 2669 &NoopRecordStore, 2670 &resolver, 2671 &sys_clock(), 2672 now(), 2673 ) 2674 .await; 2675 2676 let nsid = Nsid::new_static("sh.tangled.feed.star").unwrap(); 2677 let repo_keyed = bobbin_types::ids::EdgeKey::new(nsid.clone(), did_subj("did:plc:abalone")); 2678 let owner_keyed = bobbin_types::ids::EdgeKey::new(nsid, did_subj("did:plc:nel")); 2679 assert_eq!( 2680 store.count(&repo_keyed), 2681 1, 2682 "star should be keyed on repoDID once the repo is observed" 2683 ); 2684 assert_eq!( 2685 store.count(&owner_keyed), 2686 0, 2687 "owner DID should not collect the edge" 2688 ); 2689 } 2690 2691 #[tokio::test] 2692 async fn unresolvable_repo_subject_drops_edge() { 2693 let (store, issue_states, pull_statuses, cov, resolver) = fresh(); 2694 let star: HydrantFrame = parse_frame(json!({ 2695 "id": 1, 2696 "type": "record", 2697 "record": { 2698 "live": false, 2699 "did": "did:plc:olaren", 2700 "rev": fresh_tid().as_str(), 2701 "collection": "sh.tangled.feed.star", 2702 "rkey": "abcabcabcabcz", 2703 "action": "create", 2704 "record": { 2705 "$type": "sh.tangled.feed.star", 2706 "createdAt": "2026-05-01T00:00:00Z", 2707 "subject": "at://did:plc:nel/sh.tangled.repo/abcabcabcabcz" 2708 } 2709 } 2710 })); 2711 handle_frame( 2712 star, 2713 &store, 2714 &issue_states, 2715 &pull_statuses, 2716 &cov, 2717 &NoopSearchSink, 2718 &NoopRecordStore, 2719 &resolver, 2720 &sys_clock(), 2721 now(), 2722 ) 2723 .await; 2724 2725 let nsid = Nsid::new_static("sh.tangled.feed.star").unwrap(); 2726 let owner_keyed = bobbin_types::ids::EdgeKey::new(nsid.clone(), did_subj("did:plc:nel")); 2727 let uri_keyed = bobbin_types::ids::EdgeKey::new( 2728 nsid, 2729 uri_subj("at://did:plc:nel/sh.tangled.repo/abcabcabcabcz"), 2730 ); 2731 assert_eq!( 2732 store.count(&owner_keyed), 2733 0, 2734 "must not silently misfile under the authoring DID", 2735 ); 2736 assert_eq!( 2737 store.count(&uri_keyed), 2738 0, 2739 "unresolvable rkey-form subject must drop the edge; keeping it would index against a key that never matches bare-DID queries", 2740 ); 2741 } 2742 2743 #[tokio::test] 2744 async fn repo_without_repo_did_drops_edge() { 2745 let (store, issue_states, pull_statuses, cov, resolver) = fresh(); 2746 let repo: HydrantFrame = parse_frame(json!({ 2747 "id": 1, 2748 "type": "record", 2749 "record": { 2750 "live": false, 2751 "did": "did:plc:nel", 2752 "rev": fresh_tid().as_str(), 2753 "collection": "sh.tangled.repo", 2754 "rkey": "abcabcabcabcz", 2755 "action": "create", 2756 "record": { 2757 "$type": "sh.tangled.repo", 2758 "createdAt": "2026-05-01T00:00:00Z", 2759 "knot": "oyster.cafe", 2760 "name": "abalone" 2761 } 2762 } 2763 })); 2764 handle_frame( 2765 repo, 2766 &store, 2767 &issue_states, 2768 &pull_statuses, 2769 &cov, 2770 &NoopSearchSink, 2771 &NoopRecordStore, 2772 &resolver, 2773 &sys_clock(), 2774 now(), 2775 ) 2776 .await; 2777 2778 let star: HydrantFrame = parse_frame(json!({ 2779 "id": 2, 2780 "type": "record", 2781 "record": { 2782 "live": false, 2783 "did": "did:plc:olaren", 2784 "rev": fresh_tid().as_str(), 2785 "collection": "sh.tangled.feed.star", 2786 "rkey": "abcabcabcabcz", 2787 "action": "create", 2788 "record": { 2789 "$type": "sh.tangled.feed.star", 2790 "createdAt": "2026-05-01T00:00:00Z", 2791 "subject": "at://did:plc:nel/sh.tangled.repo/abcabcabcabcz" 2792 } 2793 } 2794 })); 2795 handle_frame( 2796 star, 2797 &store, 2798 &issue_states, 2799 &pull_statuses, 2800 &cov, 2801 &NoopSearchSink, 2802 &NoopRecordStore, 2803 &resolver, 2804 &sys_clock(), 2805 now(), 2806 ) 2807 .await; 2808 2809 let nsid = Nsid::new_static("sh.tangled.feed.star").unwrap(); 2810 let uri_keyed = bobbin_types::ids::EdgeKey::new( 2811 nsid.clone(), 2812 uri_subj("at://did:plc:nel/sh.tangled.repo/abcabcabcabcz"), 2813 ); 2814 let owner_keyed = bobbin_types::ids::EdgeKey::new(nsid, did_subj("did:plc:nel")); 2815 assert_eq!( 2816 store.count(&uri_keyed), 2817 0, 2818 "no canonical DID exists for a repo without repoDID, so the edge must be dropped", 2819 ); 2820 assert_eq!( 2821 store.count(&owner_keyed), 2822 0, 2823 "the authoring DID is not the canonical repo identity", 2824 ); 2825 } 2826 2827 #[tokio::test] 2828 async fn explicit_subject_did_skips_normalization() { 2829 let (store, issue_states, pull_statuses, cov, resolver) = fresh(); 2830 let star: HydrantFrame = parse_frame(json!({ 2831 "id": 1, 2832 "type": "record", 2833 "record": { 2834 "live": false, 2835 "did": "did:plc:olaren", 2836 "rev": fresh_tid().as_str(), 2837 "collection": "sh.tangled.feed.star", 2838 "rkey": "abcabcabcabcz", 2839 "action": "create", 2840 "record": { 2841 "$type": "sh.tangled.feed.star", 2842 "createdAt": "2026-05-01T00:00:00Z", 2843 "subjectDid": "did:plc:abalone" 2844 } 2845 } 2846 })); 2847 handle_frame( 2848 star, 2849 &store, 2850 &issue_states, 2851 &pull_statuses, 2852 &cov, 2853 &NoopSearchSink, 2854 &NoopRecordStore, 2855 &resolver, 2856 &sys_clock(), 2857 now(), 2858 ) 2859 .await; 2860 let key = bobbin_types::ids::EdgeKey::new( 2861 Nsid::new_static("sh.tangled.feed.star").unwrap(), 2862 did_subj("did:plc:abalone"), 2863 ); 2864 assert_eq!(store.count(&key), 1); 2865 } 2866 2867 #[tokio::test] 2868 async fn issue_with_repo_uri_resolves_to_repo_did() { 2869 let (store, issue_states, pull_statuses, cov, resolver) = fresh(); 2870 resolver 2871 .observe( 2872 Did::new_owned("did:plc:nel").unwrap(), 2873 Rkey::new_owned("abcabcabcabcz").unwrap(), 2874 Some(Did::new_owned("did:plc:abalone").unwrap()), 2875 ) 2876 .await; 2877 let issue: HydrantFrame = parse_frame(json!({ 2878 "id": 1, 2879 "type": "record", 2880 "record": { 2881 "live": false, 2882 "did": "did:plc:olaren", 2883 "rev": fresh_tid().as_str(), 2884 "collection": "sh.tangled.repo.issue", 2885 "rkey": "abcabcabcabcz", 2886 "action": "create", 2887 "record": { 2888 "$type": "sh.tangled.repo.issue", 2889 "createdAt": "2026-05-01T00:00:00Z", 2890 "title": "bug", 2891 "repo": "at://did:plc:nel/sh.tangled.repo/abcabcabcabcz" 2892 } 2893 } 2894 })); 2895 handle_frame( 2896 issue, 2897 &store, 2898 &issue_states, 2899 &pull_statuses, 2900 &cov, 2901 &NoopSearchSink, 2902 &NoopRecordStore, 2903 &resolver, 2904 &sys_clock(), 2905 now(), 2906 ) 2907 .await; 2908 let key = bobbin_types::ids::EdgeKey::new( 2909 Nsid::new_static("sh.tangled.repo.issue").unwrap(), 2910 did_subj("did:plc:abalone"), 2911 ); 2912 assert_eq!(store.count(&key), 1); 2913 assert_eq!( 2914 store 2915 .count_issue_state(&key, IssueStateKind::Open, None) 2916 .count 2917 .get(), 2918 1, 2919 "hydrant ingest materializes the default state count", 2920 ); 2921 } 2922 2923 #[derive(Default)] 2924 struct RecordingSearchSink { 2925 docs: tokio::sync::Mutex<Vec<bobbin_types::search::SearchDoc>>, 2926 } 2927 2928 impl SearchSink for RecordingSearchSink { 2929 async fn upsert(&self, doc: bobbin_types::search::SearchDoc) { 2930 self.docs.lock().await.push(doc); 2931 } 2932 async fn remove(&self, _uri: &AtUri<DefaultStr>) {} 2933 } 2934 2935 #[tokio::test] 2936 async fn search_index_hydrates_repo_did_via_resolver() { 2937 let (store, issue_states, pull_statuses, cov, resolver) = fresh(); 2938 resolver 2939 .observe( 2940 Did::new_owned("did:plc:nel").unwrap(), 2941 Rkey::new_owned("abcabcabcabcz").unwrap(), 2942 Some(Did::new_owned("did:plc:abalone").unwrap()), 2943 ) 2944 .await; 2945 let search = RecordingSearchSink::default(); 2946 let issue: HydrantFrame = parse_frame(json!({ 2947 "id": 1, 2948 "type": "record", 2949 "record": { 2950 "live": false, 2951 "did": "did:plc:olaren", 2952 "rev": fresh_tid().as_str(), 2953 "collection": "sh.tangled.repo.issue", 2954 "rkey": "abcabcabcabcz", 2955 "action": "create", 2956 "record": { 2957 "$type": "sh.tangled.repo.issue", 2958 "createdAt": "2026-05-01T00:00:00Z", 2959 "title": "bug", 2960 "repo": "at://did:plc:nel/sh.tangled.repo/abcabcabcabcz" 2961 } 2962 } 2963 })); 2964 handle_frame( 2965 issue, 2966 &store, 2967 &issue_states, 2968 &pull_statuses, 2969 &cov, 2970 &search, 2971 &NoopRecordStore, 2972 &resolver, 2973 &sys_clock(), 2974 now(), 2975 ) 2976 .await; 2977 let docs = search.docs.lock().await; 2978 assert_eq!(docs.len(), 1, "issue should produce one search doc"); 2979 assert_eq!( 2980 docs[0].repo, 2981 Some(Did::new_owned("did:plc:abalone").unwrap()), 2982 "search doc repo field must be resolved from the observed repo, not left as None", 2983 ); 2984 } 2985 2986 #[tokio::test] 2987 async fn repo_record_indexes_its_rkey() { 2988 let (store, issue_states, pull_statuses, cov, resolver) = fresh(); 2989 let repo: HydrantFrame = parse_frame(json!({ 2990 "id": 1, 2991 "type": "record", 2992 "record": { 2993 "live": false, 2994 "did": "did:plc:nel", 2995 "rev": fresh_tid().as_str(), 2996 "collection": "sh.tangled.repo", 2997 "rkey": "abcabcabcabcz", 2998 "action": "create", 2999 "record": { 3000 "$type": "sh.tangled.repo", 3001 "createdAt": "2026-05-01T00:00:00Z", 3002 "knot": "oyster.cafe", 3003 "name": "abalone" 3004 } 3005 } 3006 })); 3007 handle_frame( 3008 repo, 3009 &store, 3010 &issue_states, 3011 &pull_statuses, 3012 &cov, 3013 &NoopSearchSink, 3014 &NoopRecordStore, 3015 &resolver, 3016 &sys_clock(), 3017 now(), 3018 ) 3019 .await; 3020 let owner = Did::new_owned("did:plc:nel").unwrap(); 3021 assert_eq!( 3022 resolver.lookup_by_name(&owner, "abcabcabcabcz").await, 3023 Some(bobbin_types::ids::RepoIdent::new( 3024 owner.clone(), 3025 Rkey::new_owned("abcabcabcabcz").unwrap() 3026 )), 3027 ); 3028 assert_eq!(resolver.lookup_by_name(&owner, "abalone").await, None); 3029 } 3030 3031 #[tokio::test] 3032 async fn fork_record_edges_against_its_source() { 3033 let (store, issue_states, pull_statuses, cov, resolver) = fresh(); 3034 let fork: HydrantFrame = parse_frame(json!({ 3035 "id": 1, 3036 "type": "record", 3037 "record": { 3038 "live": false, 3039 "did": "did:plc:olaren", 3040 "rev": fresh_tid().as_str(), 3041 "collection": "sh.tangled.repo", 3042 "rkey": "abcabcabcabcz", 3043 "action": "create", 3044 "record": { 3045 "$type": "sh.tangled.repo", 3046 "createdAt": "2026-05-01T00:00:00Z", 3047 "knot": "oyster.cafe", 3048 "name": "abalone", 3049 "source": "did:plc:abalone" 3050 } 3051 } 3052 })); 3053 handle_frame( 3054 fork, 3055 &store, 3056 &issue_states, 3057 &pull_statuses, 3058 &cov, 3059 &NoopSearchSink, 3060 &NoopRecordStore, 3061 &resolver, 3062 &sys_clock(), 3063 now(), 3064 ) 3065 .await; 3066 let key = bobbin_types::ids::EdgeKey::new( 3067 Nsid::new_static(bobbin_types::edges::REPO_SOURCE_EDGE_KIND).unwrap(), 3068 did_subj("did:plc:abalone"), 3069 ); 3070 assert_eq!(store.count(&key), 1); 3071 } 3072 3073 fn fork_frame(source: &str) -> HydrantFrame { 3074 parse_frame(json!({ 3075 "id": 1, 3076 "type": "record", 3077 "record": { 3078 "live": false, 3079 "did": "did:plc:olaren", 3080 "rev": fresh_tid().as_str(), 3081 "collection": "sh.tangled.repo", 3082 "rkey": "abcabcabcabcz", 3083 "action": "create", 3084 "record": { 3085 "$type": "sh.tangled.repo", 3086 "createdAt": "2026-05-01T00:00:00Z", 3087 "knot": "oyster.cafe", 3088 "name": "abalone", 3089 "source": source 3090 } 3091 } 3092 })) 3093 } 3094 3095 fn fork_edge_key(subject: SubjectRef) -> bobbin_types::ids::EdgeKey { 3096 bobbin_types::ids::EdgeKey::new( 3097 Nsid::new_static(bobbin_types::edges::REPO_SOURCE_EDGE_KIND).unwrap(), 3098 subject, 3099 ) 3100 } 3101 3102 #[tokio::test] 3103 async fn fork_by_at_uri_edges_against_the_source_did() { 3104 let (store, issue_states, pull_statuses, cov, resolver) = fresh(); 3105 resolver 3106 .observe( 3107 Did::new_owned("did:plc:nel").unwrap(), 3108 Rkey::new_owned("core").unwrap(), 3109 Some(Did::new_owned("did:plc:abalone").unwrap()), 3110 ) 3111 .await; 3112 handle_frame( 3113 fork_frame("at://did:plc:nel/sh.tangled.repo/core"), 3114 &store, 3115 &issue_states, 3116 &pull_statuses, 3117 &cov, 3118 &NoopSearchSink, 3119 &NoopRecordStore, 3120 &resolver, 3121 &sys_clock(), 3122 now(), 3123 ) 3124 .await; 3125 assert_eq!(store.count(&fork_edge_key(did_subj("did:plc:abalone"))), 1); 3126 assert_eq!( 3127 store.count(&fork_edge_key(uri_subj( 3128 "at://did:plc:nel/sh.tangled.repo/core" 3129 ))), 3130 0, 3131 "the rkey form would never match a bare-DID query", 3132 ); 3133 } 3134 3135 #[tokio::test] 3136 async fn fork_of_a_repo_without_a_did_has_nothing_to_count_against() { 3137 let (store, issue_states, pull_statuses, cov, resolver) = fresh(); 3138 handle_frame( 3139 fork_frame("at://did:plc:nel/sh.tangled.repo/core"), 3140 &store, 3141 &issue_states, 3142 &pull_statuses, 3143 &cov, 3144 &NoopSearchSink, 3145 &NoopRecordStore, 3146 &resolver, 3147 &sys_clock(), 3148 now(), 3149 ) 3150 .await; 3151 assert_eq!( 3152 store.count(&fork_edge_key(uri_subj( 3153 "at://did:plc:nel/sh.tangled.repo/core" 3154 ))), 3155 0, 3156 ); 3157 } 3158 3159 #[tokio::test] 3160 async fn repo_record_without_a_name_indexes_its_rkey() { 3161 let (store, issue_states, pull_statuses, cov, resolver) = fresh(); 3162 let repo: HydrantFrame = parse_frame(json!({ 3163 "id": 1, 3164 "type": "record", 3165 "record": { 3166 "live": false, 3167 "did": "did:plc:nel", 3168 "rev": fresh_tid().as_str(), 3169 "collection": "sh.tangled.repo", 3170 "rkey": "abalone", 3171 "action": "create", 3172 "record": { 3173 "$type": "sh.tangled.repo", 3174 "createdAt": "2026-05-01T00:00:00Z", 3175 "knot": "oyster.cafe" 3176 } 3177 } 3178 })); 3179 handle_frame( 3180 repo, 3181 &store, 3182 &issue_states, 3183 &pull_statuses, 3184 &cov, 3185 &NoopSearchSink, 3186 &NoopRecordStore, 3187 &resolver, 3188 &sys_clock(), 3189 now(), 3190 ) 3191 .await; 3192 let owner = Did::new_owned("did:plc:nel").unwrap(); 3193 assert_eq!( 3194 resolver.lookup_by_name(&owner, "abalone").await, 3195 Some(bobbin_types::ids::RepoIdent::new( 3196 owner, 3197 Rkey::new_owned("abalone").unwrap() 3198 )), 3199 "repos made before the name field are only reachable by rkey", 3200 ); 3201 } 3202 3203 #[tokio::test] 3204 async fn delete_repo_record_evicts_resolver_cache() { 3205 let (store, issue_states, pull_statuses, cov, resolver) = fresh(); 3206 let owner = Did::new_owned("did:plc:nel").unwrap(); 3207 let rkey = Rkey::new_owned("abcabcabcabcz").unwrap(); 3208 resolver 3209 .observe( 3210 owner.clone(), 3211 rkey.clone(), 3212 Some(Did::new_owned("did:plc:abalone").unwrap()), 3213 ) 3214 .await; 3215 resolver.observe_rkey(owner.clone(), rkey.clone()).await; 3216 assert!( 3217 resolver.cached_resolution(&owner, &rkey).await.is_some(), 3218 "observe must seed the cache", 3219 ); 3220 let delete: HydrantFrame = parse_frame(json!({ 3221 "id": 1, 3222 "type": "record", 3223 "record": { 3224 "live": false, 3225 "did": owner.as_ref(), 3226 "rev": fresh_tid().as_str(), 3227 "collection": "sh.tangled.repo", 3228 "rkey": rkey.as_ref(), 3229 "action": "delete" 3230 } 3231 })); 3232 handle_frame( 3233 delete, 3234 &store, 3235 &issue_states, 3236 &pull_statuses, 3237 &cov, 3238 &NoopSearchSink, 3239 &NoopRecordStore, 3240 &resolver, 3241 &sys_clock(), 3242 now(), 3243 ) 3244 .await; 3245 assert!( 3246 resolver.cached_resolution(&owner, &rkey).await.is_none(), 3247 "deleting the repo record must clear the resolver cache so future observes are not blocked by a stale Authoritative entry", 3248 ); 3249 assert_eq!( 3250 resolver.lookup_by_name(&owner, "abcabcabcabcz").await, 3251 None, 3252 "a deleted repo must stop answering on its url", 3253 ); 3254 } 3255 3256 #[tokio::test] 3257 async fn cancel_short_circuits_a_hung_ws_connect() { 3258 let _listener = tokio::net::TcpListener::bind("127.0.0.1:0") 3259 .await 3260 .expect("bind sink listener"); 3261 let port = _listener.local_addr().expect("local addr").port(); 3262 let cfg = 3263 IngestConfig::new(Url::parse(&format!("ws://127.0.0.1:{port}")).expect("hydrant url")); 3264 let cancel = CancellationToken::new(); 3265 let runtime = fresh_runtime(cancel.clone()); 3266 let task = tokio::spawn(async move { run(cfg, runtime).await }); 3267 tokio::time::sleep(Duration::from_millis(100)).await; 3268 cancel.cancel(); 3269 let outcome = tokio::time::timeout(Duration::from_secs(2), task) 3270 .await 3271 .expect("cancel must short-circuit the hung ws connect"); 3272 assert!(matches!(outcome, Ok(Ok(()))), "got {outcome:?}"); 3273 } 3274 3275 struct CloseOnConnectTransport { 3276 used: std::sync::Mutex<bool>, 3277 } 3278 impl bobbin_runtime::WsTransport for CloseOnConnectTransport { 3279 fn connect(&self, _url: Url) -> bobbin_runtime::WsConnectFuture { 3280 let mut used = self.used.lock().unwrap(); 3281 if *used { 3282 return Box::pin(async move { 3283 Err(NetworkError::Connect("only one connect allowed".to_owned())) 3284 }); 3285 } 3286 *used = true; 3287 Box::pin(async move { 3288 let mut q = std::collections::VecDeque::new(); 3289 q.push_back(Ok(WsMessage::Close { 3290 code: 1000, 3291 reason: "bye".to_owned(), 3292 })); 3293 let stream: Box<dyn WsStream> = Box::new(ScriptedWsStream { messages: q }); 3294 struct NoopSink; 3295 impl bobbin_runtime::WsSink for NoopSink { 3296 fn send<'a>(&'a mut self, _m: WsMessage) -> bobbin_runtime::WsSendFuture<'a> { 3297 Box::pin(async move { Ok(()) }) 3298 } 3299 } 3300 let sink: Box<dyn bobbin_runtime::WsSink> = Box::new(NoopSink); 3301 Ok(bobbin_runtime::WsConn { sink, stream }) 3302 }) 3303 } 3304 } 3305 3306 #[tokio::test] 3307 async fn run_session_returns_after_remote_close_when_outer_cancel_unfired() { 3308 let cfg = IngestConfig::new(Url::parse("ws://127.0.0.1:1").unwrap()); 3309 let cancel = CancellationToken::new(); 3310 let mut runtime = fresh_runtime(cancel.clone()); 3311 runtime.ws = Arc::new(CloseOnConnectTransport { 3312 used: std::sync::Mutex::new(false), 3313 }); 3314 3315 let res = tokio::time::timeout( 3316 Duration::from_secs(3), 3317 run_session(&cfg, HydrantCursor::new(0), &runtime), 3318 ) 3319 .await; 3320 assert!( 3321 res.is_ok(), 3322 "run_session must return after a remote Close even when outer cancel never fires - regression for a session-scoped task hanging on the parent token", 3323 ); 3324 } 3325 3326 struct HangingSink; 3327 impl bobbin_runtime::WsSink for HangingSink { 3328 fn send<'a>(&'a mut self, _: WsMessage) -> bobbin_runtime::WsSendFuture<'a> { 3329 Box::pin(std::future::pending()) 3330 } 3331 } 3332 3333 #[tokio::test(start_paused = true)] 3334 async fn timed_send_surfaces_send_timeout_when_sink_pends_forever() { 3335 let mut sink: Box<dyn bobbin_runtime::WsSink> = Box::new(HangingSink); 3336 let task = 3337 tokio::spawn(async move { timed_send(&mut sink, WsMessage::Ping(Bytes::new())).await }); 3338 tokio::time::advance(SEND_TIMEOUT + Duration::from_secs(1)).await; 3339 let result = task.await.expect("task panicked"); 3340 match result { 3341 Err(IngestError::SendTimeout(d)) => assert_eq!(d, SEND_TIMEOUT), 3342 other => panic!( 3343 "expected SendTimeout, got {other:?}. A bare ws_sink.send.await would hang forever on a half-dead socket and starve the writer's pong-deadline arm", 3344 ), 3345 } 3346 } 3347 3348 struct OkSink; 3349 impl bobbin_runtime::WsSink for OkSink { 3350 fn send<'a>(&'a mut self, _: WsMessage) -> bobbin_runtime::WsSendFuture<'a> { 3351 Box::pin(async move { Ok(()) }) 3352 } 3353 } 3354 3355 #[tokio::test] 3356 async fn timed_send_returns_ok_when_sink_succeeds_promptly() { 3357 let mut sink: Box<dyn bobbin_runtime::WsSink> = Box::new(OkSink); 3358 let result = timed_send(&mut sink, WsMessage::Ping(Bytes::new())).await; 3359 assert!(matches!(result, Ok(())), "got {result:?}"); 3360 } 3361 3362 #[test] 3363 fn classify_error_frame_with_id_dispatches_as_error_not_unknown_kind() { 3364 let text = r#"{"id":42,"type":"error","error":"ConsumerTooSlow","message":"slow"}"#; 3365 match classify_text_frame(text) { 3366 Err(IngestError::ConsumerTooSlow { message }) => assert_eq!(message, "slow"), 3367 other => panic!( 3368 "type=\"error\" must dispatch to the error path even when id is present, got: {other:?}" 3369 ), 3370 } 3371 } 3372 3373 #[tokio::test] 3374 async fn buffered_pipeline_preserves_cursor_order_under_resolve_latency_skew() { 3375 use bobbin_record_lru::RecordStore; 3376 use bobbin_types::record::RecordBody; 3377 use std::sync::Mutex; 3378 3379 let server = wiremock::MockServer::start().await; 3380 wiremock::Mock::given(wiremock::matchers::method("GET")) 3381 .and(wiremock::matchers::path("/xrpc/com.atproto.repo.getRecord")) 3382 .respond_with( 3383 wiremock::ResponseTemplate::new(404).set_delay(Duration::from_millis(150)), 3384 ) 3385 .mount(&server) 3386 .await; 3387 3388 let client = bobbin_slingshot_client::SlingshotClient::with_default_http( 3389 Url::parse(&server.uri()).unwrap(), 3390 ) 3391 .unwrap(); 3392 let clock: Arc<dyn Clock> = Arc::new(SystemClock::new()); 3393 let resolver = Arc::new(RepoIdResolver::with_slingshot( 3394 client, 3395 clock.clone(), 3396 RuntimeHasher::default(), 3397 )); 3398 3399 let owner = Did::new_owned("did:plc:nel").unwrap(); 3400 let abalone = Did::new_owned("did:plc:abalone").unwrap(); 3401 let fast_rkeys: [Rkey<DefaultStr>; 2] = [rkey("fastrkeyaa01"), rkey("fastrkeyaa02")]; 3402 let slow_rkeys: [Rkey<DefaultStr>; 2] = [rkey("slowrkeyaa01"), rkey("slowrkeyaa02")]; 3403 for r in &fast_rkeys { 3404 resolver 3405 .observe(owner.clone(), r.clone(), Some(abalone.clone())) 3406 .await; 3407 } 3408 3409 #[derive(Default)] 3410 struct Capturing { 3411 urls: Mutex<Vec<AtUri<DefaultStr>>>, 3412 } 3413 impl RecordStore for Capturing { 3414 fn get(&self, _uri: &AtUri<DefaultStr>) -> Option<Arc<RecordBody>> { 3415 None 3416 } 3417 fn put(&self, uri: AtUri<DefaultStr>, _body: Arc<RecordBody>) { 3418 self.urls.lock().unwrap().push(uri); 3419 } 3420 fn remove(&self, _uri: &AtUri<DefaultStr>) {} 3421 } 3422 let capturing = Arc::new(Capturing::default()); 3423 3424 let runtime: IngestRuntime<NoopSearchSink> = IngestRuntime { 3425 store: Arc::new(EdgeStore::new(RuntimeHasher::default())), 3426 issue_states: Arc::new(StateIndex::new(RuntimeHasher::default())), 3427 pull_statuses: Arc::new(StateIndex::new(RuntimeHasher::default())), 3428 coverage: Arc::new(CoverageWatch::new()), 3429 search: Arc::new(NoopSearchSink), 3430 records: capturing.clone() as Arc<dyn RecordStore>, 3431 resolver, 3432 clock, 3433 entropy: Arc::new(OsEntropy), 3434 ws: TungsteniteWs::shared(), 3435 cancel: CancellationToken::new(), 3436 disconnects: None, 3437 warming_shadow: None, 3438 warming_buffer: None, 3439 knot_registry: None, 3440 knot_gate: None, 3441 }; 3442 3443 let parallelism = 4usize; 3444 let (frame_tx, frame_rx) = tokio::sync::mpsc::channel::<HydrantFrame>(64); 3445 let pipeline_rt = runtime.clone(); 3446 let pipeline = tokio::spawn(async move { 3447 let prep_rt = pipeline_rt.clone(); 3448 let resolve_rt = pipeline_rt.clone(); 3449 let commit_rt = pipeline_rt; 3450 ReceiverStream::new(frame_rx) 3451 .then(move |frame| prep_stage(frame, prep_rt.clone())) 3452 .map(move |staged| resolve_stage(staged, resolve_rt.clone())) 3453 .buffered(parallelism) 3454 .for_each(move |staged| commit_stage(staged, commit_rt.clone(), parallelism)) 3455 .await; 3456 }); 3457 3458 let mk = |id: u64, idx: u64, repo_rkey: &Rkey<DefaultStr>| -> HydrantFrame { 3459 parse_frame(json!({ 3460 "id": id, 3461 "type": "record", 3462 "record": { 3463 "live": false, 3464 "did": "did:plc:olaren", 3465 "rev": fresh_tid().as_str(), 3466 "collection": "sh.tangled.feed.star", 3467 "rkey": format!("starrkeya{idx:04}"), 3468 "action": "create", 3469 "cid": VALID_CID, 3470 "record": { 3471 "$type": "sh.tangled.feed.star", 3472 "createdAt": "2026-05-01T00:00:00Z", 3473 "subject": format!("at://did:plc:nel/sh.tangled.repo/{}", repo_rkey.as_ref()), 3474 } 3475 } 3476 })) 3477 }; 3478 3479 frame_tx.send(mk(1, 1, &slow_rkeys[0])).await.unwrap(); 3480 frame_tx.send(mk(2, 2, &fast_rkeys[0])).await.unwrap(); 3481 frame_tx.send(mk(3, 3, &slow_rkeys[1])).await.unwrap(); 3482 frame_tx.send(mk(4, 4, &fast_rkeys[1])).await.unwrap(); 3483 drop(frame_tx); 3484 3485 pipeline.await.unwrap(); 3486 3487 let captured = capturing.urls.lock().unwrap(); 3488 let rkeys: Vec<String> = captured 3489 .iter() 3490 .filter_map(|uri| uri.rkey().map(|r| r.as_ref().to_owned())) 3491 .collect(); 3492 assert_eq!( 3493 rkeys, 3494 vec![ 3495 "starrkeya0001".to_owned(), 3496 "starrkeya0002".to_owned(), 3497 "starrkeya0003".to_owned(), 3498 "starrkeya0004".to_owned(), 3499 ], 3500 "buffered(N) must preserve cursor order even when resolves complete out of order, with ~150ms slow vs cache-hit fast as the latency skew here", 3501 ); 3502 assert_eq!(runtime.coverage.snapshot().last_cursor().raw(), 4); 3503 } 3504}