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 3516 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 repo.name.clone(), 1055 ) 1056 .await; 1057 if let Some(prior) = superseded { 1058 let prior_uri = format!( 1059 "at://{}/sh.tangled.repo/{}", 1060 prior.owner.as_ref(), 1061 prior.rkey.as_ref(), 1062 ); 1063 if let Ok(prior_at_uri) = AtUri::<DefaultStr>::new_owned(&prior_uri) { 1064 evict_from_buffer(ctx.buffer, &prior_at_uri).await; 1065 ctx.store.remove_source(&prior_at_uri); 1066 ctx.records.remove(&prior_at_uri); 1067 ctx.search.remove(&prior_at_uri).await; 1068 } 1069 } 1070 if let Some(buffer) = ctx.buffer { 1071 let drained = buffer.take_observed(&record.did, &record.rkey).await; 1072 if !drained.is_empty() { 1073 finalize_drained(ctx, drained).await; 1074 } 1075 } 1076 if let Some(registry) = ctx.knot_registry { 1077 let host = KnotHostKey::new(repo.knot.as_ref()); 1078 match repo.repo_did.clone() { 1079 Some(repo_did) => registry.observe_repo(&host, repo_did), 1080 None => registry.observe_host(&host), 1081 } 1082 } 1083 } 1084 match acl_disposition(&parsed, ctx.knot_gate, ctx.knot_registry) { 1085 AclDisposition::NativeSkip => { 1086 if let Some(registry) = ctx.knot_registry { 1087 registry.forget_legacy_member(&source); 1088 } 1089 return PendingOp::Delete { source, nsid }; 1090 } 1091 AclDisposition::LegacyMember { host } => { 1092 if let Some(registry) = ctx.knot_registry { 1093 registry.observe_host(&host); 1094 registry.note_legacy_member(source.clone(), &host); 1095 } 1096 } 1097 AclDisposition::Other => {} 1098 } 1099 let edges = match parsed.extract_edges(&source) { 1100 Ok(es) => es, 1101 Err(e) => { 1102 warn!(?e, "edge extraction failed, clearing cache"); 1103 return PendingOp::ClearCache { source }; 1104 } 1105 }; 1106 let _ = ctx.store.intern_source(&source); 1107 PendingOp::Upsert { 1108 source, 1109 nsid, 1110 parsed: Box::new(parsed), 1111 bytes, 1112 cid: record.cid, 1113 edges, 1114 } 1115 } 1116 RecordAction::Delete => { 1117 evict_from_buffer(ctx.buffer, &source).await; 1118 if nsid.as_ref() == "sh.tangled.repo" { 1119 ctx.resolver.forget(&record.did, &record.rkey).await; 1120 } 1121 if nsid.as_ref() == "sh.tangled.knot.member" 1122 && let Some(registry) = ctx.knot_registry 1123 { 1124 registry.forget_legacy_member(&source); 1125 } 1126 PendingOp::Delete { source, nsid } 1127 } 1128 RecordAction::Other => { 1129 debug!(collection = %nsid, "ignoring unknown record action"); 1130 PendingOp::Noop 1131 } 1132 } 1133} 1134 1135enum AclDisposition { 1136 Other, 1137 NativeSkip, 1138 LegacyMember { host: KnotHostKey }, 1139} 1140 1141fn acl_disposition( 1142 parsed: &Record, 1143 gate: Option<&CapabilityGate>, 1144 registry: Option<&KnotRegistry>, 1145) -> AclDisposition { 1146 let Some(gate) = gate else { 1147 return AclDisposition::Other; 1148 }; 1149 match parsed { 1150 Record::KnotMember(member) => { 1151 let host = KnotHostKey::new(member.domain.as_ref()); 1152 if gate.is_native(&host) { 1153 AclDisposition::NativeSkip 1154 } else { 1155 AclDisposition::LegacyMember { host } 1156 } 1157 } 1158 Record::Collaborator(collaborator) => { 1159 let native = registry 1160 .and_then(|registry| registry.host_of_repo(&collaborator.repo)) 1161 .is_some_and(|host| gate.is_native(&host)); 1162 if native { 1163 AclDisposition::NativeSkip 1164 } else { 1165 AclDisposition::Other 1166 } 1167 } 1168 _ => AclDisposition::Other, 1169 } 1170} 1171 1172fn fallback_rfc3339( 1173 rkey: &Rkey<DefaultStr>, 1174 rev: &jacquard_common::types::tid::Tid, 1175) -> Option<String> { 1176 let tid = jacquard_common::types::tid::Tid::new(rkey.as_ref()) 1177 .ok() 1178 .unwrap_or_else(|| rev.clone()); 1179 let micros = i64::try_from(tid.timestamp()).ok()?; 1180 let dt = chrono::DateTime::<chrono::Utc>::from_timestamp_micros(micros)?; 1181 Some(dt.to_rfc3339_opts(chrono::SecondsFormat::Micros, true)) 1182} 1183 1184async fn evict_from_buffer(buffer: Option<&WarmingBuffer>, source: &AtUri<DefaultStr>) { 1185 if let Some(buffer) = buffer 1186 && !buffer.is_sealed() 1187 { 1188 buffer.evict_source(source).await; 1189 } 1190} 1191 1192async fn resolve_pending<S: SearchSink + 'static>( 1193 pending: Pending, 1194 ctx: &PipelineCtx<'_, S>, 1195) -> Pending { 1196 let Pending { 1197 cursor, 1198 signal, 1199 regime, 1200 op, 1201 } = pending; 1202 let op = match op { 1203 PendingOp::Upsert { 1204 source, 1205 nsid, 1206 parsed, 1207 bytes, 1208 cid, 1209 edges, 1210 } => { 1211 let pieces = UpsertPieces { 1212 source, 1213 nsid, 1214 parsed: *parsed, 1215 bytes, 1216 cid, 1217 edges, 1218 }; 1219 let pieces = match try_park_warming(ctx, cursor, pieces).await { 1220 ParkOutcome::Parked { nsid } => { 1221 return Pending { 1222 cursor, 1223 signal, 1224 regime, 1225 op: PendingOp::Parked { nsid }, 1226 }; 1227 } 1228 ParkOutcome::Passthrough(pieces) => *pieces, 1229 }; 1230 let UpsertPieces { 1231 source, 1232 nsid, 1233 parsed, 1234 bytes, 1235 cid, 1236 edges, 1237 } = pieces; 1238 let edges = normalize_subjects(edges, ctx.resolver, ctx.coverage, ctx.shadow).await; 1239 PendingOp::Upsert { 1240 source, 1241 nsid, 1242 parsed: Box::new(parsed), 1243 bytes, 1244 cid, 1245 edges, 1246 } 1247 } 1248 other => other, 1249 }; 1250 Pending { 1251 cursor, 1252 signal, 1253 regime, 1254 op, 1255 } 1256} 1257 1258struct UpsertPieces { 1259 source: AtUri<DefaultStr>, 1260 nsid: Nsid<DefaultStr>, 1261 parsed: Record, 1262 bytes: Bytes, 1263 cid: Option<Cid<DefaultStr>>, 1264 edges: Vec<Edge>, 1265} 1266 1267impl From<ParkedUpsert> for UpsertPieces { 1268 fn from(u: ParkedUpsert) -> Self { 1269 Self { 1270 source: u.source, 1271 nsid: u.nsid, 1272 parsed: u.parsed, 1273 bytes: u.bytes, 1274 cid: u.cid, 1275 edges: u.edges, 1276 } 1277 } 1278} 1279 1280enum ParkOutcome { 1281 Parked { nsid: Nsid<DefaultStr> }, 1282 Passthrough(Box<UpsertPieces>), 1283} 1284 1285async fn try_park_warming<S: SearchSink + 'static>( 1286 ctx: &PipelineCtx<'_, S>, 1287 cursor: HydrantCursor, 1288 pieces: UpsertPieces, 1289) -> ParkOutcome { 1290 let Some(buffer) = ctx.buffer else { 1291 return ParkOutcome::Passthrough(Box::new(pieces)); 1292 }; 1293 if ctx.coverage.snapshot().is_ready() || buffer.is_sealed() { 1294 return ParkOutcome::Passthrough(Box::new(pieces)); 1295 } 1296 let deps = collect_unresolved_deps(&pieces.edges, ctx.resolver).await; 1297 if deps.is_empty() { 1298 return ParkOutcome::Passthrough(Box::new(pieces)); 1299 } 1300 let nsid = pieces.nsid.clone(); 1301 let upsert = ParkedUpsert { 1302 cursor, 1303 source: pieces.source, 1304 nsid: pieces.nsid, 1305 parsed: pieces.parsed, 1306 bytes: pieces.bytes, 1307 cid: pieces.cid, 1308 edges: pieces.edges, 1309 }; 1310 let deps_for_shadow = ctx.shadow.is_some().then(|| deps.clone()); 1311 match buffer.try_park(upsert, deps).await { 1312 Ok(()) => { 1313 if let Some((shadow, noted)) = ctx.shadow.zip(deps_for_shadow) { 1314 let _ = futures::future::join_all(noted.into_iter().map(|dep| async move { 1315 shadow.note_unresolved(dep.owner, dep.rkey).await; 1316 })) 1317 .await; 1318 } 1319 ParkOutcome::Parked { nsid } 1320 } 1321 Err(returned) => ParkOutcome::Passthrough(Box::new(returned.into())), 1322 } 1323} 1324 1325async fn collect_unresolved_deps(edges: &[Edge], resolver: &RepoIdResolver) -> Vec<RepoIdent> { 1326 futures::stream::iter(edges) 1327 .fold(Vec::new(), |mut acc, edge| async move { 1328 let Some(uri) = edge.subject.as_uri() else { 1329 return acc; 1330 }; 1331 let Some((owner, rkey)) = parse_repo_subject_uri(uri) else { 1332 return acc; 1333 }; 1334 if resolver.cached_resolution(&owner, &rkey).await.is_some() { 1335 return acc; 1336 } 1337 let candidate = RepoIdent::new(owner, rkey); 1338 if !acc.contains(&candidate) { 1339 acc.push(candidate); 1340 } 1341 acc 1342 }) 1343 .await 1344} 1345 1346async fn finalize_drained<S: SearchSink + 'static>( 1347 ctx: &PipelineCtx<'_, S>, 1348 drained: Vec<ParkedUpsert>, 1349) { 1350 for upsert in drained { 1351 let ParkedUpsert { 1352 cursor: _, 1353 source, 1354 nsid: _, 1355 parsed, 1356 bytes, 1357 cid, 1358 edges, 1359 } = upsert; 1360 let edges = normalize_subjects(edges, ctx.resolver, ctx.coverage, None).await; 1361 cache_body(ctx.records, &source, cid, bytes); 1362 let outcome = upsert_record_indexes( 1363 ctx.store, 1364 ctx.issue_states, 1365 ctx.pull_statuses, 1366 &source, 1367 edges, 1368 &parsed, 1369 ); 1370 log_unknown_state_variant(outcome, &source); 1371 index_search(ctx.search, ctx.resolver, &source, parsed).await; 1372 } 1373} 1374 1375#[allow(clippy::too_many_arguments)] 1376async fn commit_pending<S: SearchSink>( 1377 pending: Pending, 1378 store: &EdgeStore, 1379 issue_states: &StateIndex<IssueStateKind>, 1380 pull_statuses: &StateIndex<PullStatusKind>, 1381 coverage: &CoverageWatch, 1382 search: &S, 1383 records: &dyn RecordStore, 1384 resolver: &RepoIdResolver, 1385) { 1386 let Pending { 1387 cursor, 1388 signal, 1389 regime: _, 1390 op, 1391 } = pending; 1392 match op { 1393 PendingOp::Noop | PendingOp::Parked { .. } => {} 1394 PendingOp::ClearCache { source } => records.remove(&source), 1395 PendingOp::Upsert { 1396 source, 1397 nsid: _, 1398 parsed, 1399 bytes, 1400 cid, 1401 edges, 1402 } => { 1403 cache_body(records, &source, cid, bytes); 1404 let outcome = 1405 upsert_record_indexes(store, issue_states, pull_statuses, &source, edges, &parsed); 1406 log_unknown_state_variant(outcome, &source); 1407 index_search(search, resolver, &source, *parsed).await; 1408 } 1409 PendingOp::Delete { source, nsid } => { 1410 delete_record_indexes(store, issue_states, pull_statuses, &source, &nsid); 1411 records.remove(&source); 1412 search.remove(&source).await; 1413 } 1414 } 1415 coverage.update(|c| c.advance(cursor).maybe_promote(signal)); 1416} 1417 1418async fn index_search<S: SearchSink>( 1419 search: &S, 1420 resolver: &RepoIdResolver, 1421 source: &AtUri<DefaultStr>, 1422 parsed: Record, 1423) { 1424 let Some(searchable) = SearchableRecord::try_from_record(parsed) else { 1425 return; 1426 }; 1427 let Some(searchable) = searchable.normalize(resolver).await else { 1428 return; 1429 }; 1430 search.upsert(searchable.to_search_doc(source)).await; 1431} 1432 1433#[cfg(test)] 1434#[allow(clippy::too_many_arguments)] 1435async fn handle_frame<S: SearchSink + 'static>( 1436 frame: HydrantFrame, 1437 store: &EdgeStore, 1438 issue_states: &StateIndex<IssueStateKind>, 1439 pull_statuses: &StateIndex<PullStatusKind>, 1440 coverage: &CoverageWatch, 1441 search: &S, 1442 records: &dyn RecordStore, 1443 resolver: &RepoIdResolver, 1444 clock: &dyn Clock, 1445 now: UnixMicros, 1446) { 1447 let _ = clock; 1448 let ctx = PipelineCtx { 1449 resolver, 1450 store, 1451 issue_states, 1452 pull_statuses, 1453 coverage, 1454 records, 1455 search, 1456 shadow: None, 1457 buffer: None, 1458 knot_registry: None, 1459 knot_gate: None, 1460 }; 1461 let pending = prepare_frame(frame, &ctx, now).await; 1462 let pending = resolve_pending(pending, &ctx).await; 1463 commit_pending( 1464 pending, 1465 store, 1466 issue_states, 1467 pull_statuses, 1468 coverage, 1469 search, 1470 records, 1471 resolver, 1472 ) 1473 .await; 1474} 1475 1476fn log_unknown_state_variant(outcome: ApplyOutcome, source: &AtUri<DefaultStr>) { 1477 if matches!(outcome, ApplyOutcome::UnknownVariant) { 1478 warn!( 1479 target: "bobbin_ingest::state_index", 1480 %source, 1481 "state record has unknown wire variant, skipping index update", 1482 ); 1483 } 1484} 1485 1486fn promotion_signal(record: Option<&RecordFrame>, now: UnixMicros) -> PromotionSignal { 1487 PromotionSignal { 1488 rev_micros: record.map(|r| r.rev.timestamp()), 1489 now_micros: now.raw(), 1490 skew_micros: READY_SKEW.as_micros() as u64, 1491 } 1492} 1493 1494fn cache_body( 1495 records: &dyn RecordStore, 1496 source: &AtUri<DefaultStr>, 1497 cid: Option<Cid<DefaultStr>>, 1498 bytes: Bytes, 1499) { 1500 match cid { 1501 Some(cid) => records.put( 1502 source.clone(), 1503 Arc::new(RecordBody { 1504 uri: source.clone(), 1505 cid, 1506 value: bytes, 1507 }), 1508 ), 1509 None => records.remove(source), 1510 } 1511} 1512 1513async fn normalize_subjects( 1514 edges: Vec<Edge>, 1515 resolver: &RepoIdResolver, 1516 coverage: &CoverageWatch, 1517 shadow: Option<&WarmingShadowBuffer>, 1518) -> Vec<Edge> { 1519 let warming = shadow.is_some() && !coverage.snapshot().is_ready(); 1520 futures::stream::iter(edges) 1521 .filter_map(|edge| async move { 1522 let Some(uri) = edge.subject.as_uri() else { 1523 return Some(edge); 1524 }; 1525 let Some((owner, rkey)) = parse_repo_subject_uri(uri) else { 1526 return Some(edge); 1527 }; 1528 if warming 1529 && let Some(shadow) = shadow 1530 && resolver.cached_resolution(&owner, &rkey).await.is_none() 1531 { 1532 shadow 1533 .note_unresolved(owner.clone(), rkey.clone()) 1534 .await; 1535 } 1536 match resolver.resolve(&owner, &rkey).await { 1537 Resolution::Mapped(repo_did) => Some(Edge { 1538 subject: SubjectRef::Did(repo_did), 1539 ..edge 1540 }), 1541 Resolution::NoRepoDid => { 1542 warn!( 1543 target: "bobbin_ingest::normalize", 1544 kind = %edge.kind, 1545 owner = owner.as_ref(), 1546 rkey = rkey.as_ref(), 1547 source = edge.source.as_ref(), 1548 "dropping edge: target repo has no repoDid, no canonical DID subject available", 1549 ); 1550 None 1551 } 1552 Resolution::Unresolvable => { 1553 warn!( 1554 target: "bobbin_ingest::normalize", 1555 kind = %edge.kind, 1556 owner = owner.as_ref(), 1557 rkey = rkey.as_ref(), 1558 source = edge.source.as_ref(), 1559 "dropping edge: repo unresolvable, rkey-form subject will not match bare-DID queries", 1560 ); 1561 None 1562 } 1563 } 1564 }) 1565 .collect() 1566 .await 1567} 1568 1569fn parse_repo_subject_uri(uri: &AtUri<DefaultStr>) -> Option<(Did<DefaultStr>, Rkey<DefaultStr>)> { 1570 let collection = uri.collection()?; 1571 if collection.as_ref() != "sh.tangled.repo" { 1572 return None; 1573 } 1574 let AtIdentifier::Did(authority) = uri.authority() else { 1575 return None; 1576 }; 1577 let rkey = uri.rkey()?; 1578 let owner = Did::new_owned(authority.as_ref()).ok()?; 1579 let rkey = Rkey::new_owned(rkey.as_ref()).ok()?; 1580 Some((owner, rkey)) 1581} 1582 1583fn build_source_uri(r: &RecordFrame) -> Result<AtUri<DefaultStr>, IngestError> { 1584 Ok(AtUri::from_parts_owned( 1585 r.did.as_ref(), 1586 r.collection.as_ref(), 1587 r.rkey.as_ref(), 1588 )?) 1589} 1590 1591#[cfg(test)] 1592mod tests { 1593 use super::*; 1594 use bobbin_edge_index::Coverage; 1595 use bobbin_record_lru::{CacheCapacity, LruRecordStore, NoopRecordStore, RecordStore}; 1596 use bobbin_runtime::{OsEntropy, RuntimeHasher, SystemClock, TungsteniteWs}; 1597 use bobbin_types::search::NoopSearchSink; 1598 use jacquard_common::types::nsid::Nsid; 1599 use jacquard_common::types::tid::Tid; 1600 use serde_json::json; 1601 1602 const VALID_CID: &str = "bafyreieqygohnz2zqyvtvktbjpvhutphobcmbsnt4q5lc36ri7vpcmoz4i"; 1603 1604 fn did_subj(s: &str) -> SubjectRef { 1605 SubjectRef::Did(Did::new_owned(s).unwrap()) 1606 } 1607 1608 fn uri_subj(s: &str) -> SubjectRef { 1609 SubjectRef::Uri(AtUri::new_owned(s).unwrap()) 1610 } 1611 1612 fn rkey(s: &str) -> Rkey<DefaultStr> { 1613 Rkey::new_owned(s).unwrap() 1614 } 1615 1616 #[allow(clippy::type_complexity)] 1617 fn fresh() -> ( 1618 Arc<EdgeStore>, 1619 Arc<StateIndex<IssueStateKind>>, 1620 Arc<StateIndex<PullStatusKind>>, 1621 Arc<CoverageWatch>, 1622 Arc<RepoIdResolver>, 1623 ) { 1624 ( 1625 Arc::new(EdgeStore::new(RuntimeHasher::default())), 1626 Arc::new(StateIndex::new(RuntimeHasher::default())), 1627 Arc::new(StateIndex::new(RuntimeHasher::default())), 1628 Arc::new(CoverageWatch::new()), 1629 Arc::new(RepoIdResolver::detached(RuntimeHasher::default())), 1630 ) 1631 } 1632 1633 fn now() -> UnixMicros { 1634 SystemClock::new().now_unix_micros() 1635 } 1636 1637 fn sys_clock() -> SystemClock { 1638 SystemClock::new() 1639 } 1640 1641 fn parse_frame(value: serde_json::Value) -> HydrantFrame { 1642 let text = serde_json::to_string(&value).expect("serialize fixture"); 1643 serde_json::from_str(&text).expect("deserialize fixture") 1644 } 1645 1646 fn fresh_tid() -> Tid { 1647 Tid::now_0() 1648 } 1649 1650 #[tokio::test] 1651 async fn ignores_non_tangled_collections() { 1652 let (store, issue_states, pull_statuses, cov, resolver) = fresh(); 1653 let frame: HydrantFrame = parse_frame(json!({ 1654 "id": 1, 1655 "type": "record", 1656 "record": { 1657 "live": false, 1658 "did": "did:plc:nel", 1659 "rev": fresh_tid().as_str(), 1660 "collection": "app.bsky.feed.post", 1661 "rkey": "abcabcabcabcz", 1662 "action": "create", 1663 "record": {"$type": "app.bsky.feed.post", "text": "hi"} 1664 } 1665 })); 1666 handle_frame( 1667 frame, 1668 &store, 1669 &issue_states, 1670 &pull_statuses, 1671 &cov, 1672 &NoopSearchSink, 1673 &NoopRecordStore, 1674 &resolver, 1675 &sys_clock(), 1676 now(), 1677 ) 1678 .await; 1679 assert_eq!(store.key_count(), 0); 1680 assert_eq!(cov.snapshot().events_processed(), 1); 1681 } 1682 1683 #[tokio::test] 1684 async fn native_knot_member_skipped_legacy_indexed() { 1685 use bobbin_knot_ingest::{CapabilityGate, KnotClient, KnotRegistry}; 1686 use wiremock::matchers::{method, path}; 1687 use wiremock::{Mock, MockServer, ResponseTemplate}; 1688 1689 let server = MockServer::start().await; 1690 Mock::given(method("GET")) 1691 .and(path("/xrpc/sh.tangled.knot.version")) 1692 .respond_with(ResponseTemplate::new(200).set_body_json(json!({ 1693 "version": "1.1.0", 1694 "capabilities": ["knot-acl"] 1695 }))) 1696 .mount(&server) 1697 .await; 1698 let url = url::Url::parse(&server.uri()).unwrap(); 1699 let native_host = format!("{}:{}", url.host_str().unwrap(), url.port().unwrap()); 1700 1701 let gate = CapabilityGate::new( 1702 KnotClient::with_default_http(true).unwrap(), 1703 Arc::new(SystemClock::new()), 1704 true, 1705 true, 1706 ); 1707 assert!(gate.has_knot_acl(&KnotHostKey::new(&native_host)).await); 1708 1709 let registry = KnotRegistry::new(); 1710 let (store, issue_states, pull_statuses, cov, resolver) = fresh(); 1711 let ctx = PipelineCtx { 1712 resolver: &resolver, 1713 store: &store, 1714 issue_states: &issue_states, 1715 pull_statuses: &pull_statuses, 1716 coverage: &cov, 1717 records: &NoopRecordStore, 1718 search: &NoopSearchSink, 1719 shadow: None, 1720 buffer: None, 1721 knot_registry: Some(&registry), 1722 knot_gate: Some(&gate), 1723 }; 1724 1725 let member_frame = |id: u64, rkey: &str, domain: &str| { 1726 parse_frame(json!({ 1727 "id": id, 1728 "type": "record", 1729 "record": { 1730 "live": false, 1731 "did": "did:plc:akshay", 1732 "rev": fresh_tid().as_str(), 1733 "collection": "sh.tangled.knot.member", 1734 "rkey": rkey, 1735 "action": "create", 1736 "record": { 1737 "$type": "sh.tangled.knot.member", 1738 "subject": "did:plc:boltless", 1739 "domain": domain, 1740 "createdAt": "2026-06-01T00:00:00Z" 1741 } 1742 } 1743 })) 1744 }; 1745 1746 let native = 1747 prepare_frame(member_frame(1, "aaaaaaaaaaaaz", &native_host), &ctx, now()).await; 1748 assert!( 1749 matches!(native.op, PendingOp::Delete { .. }), 1750 "member record for a native knot must be dropped" 1751 ); 1752 1753 let legacy = 1754 prepare_frame(member_frame(2, "bbbbbbbbbbbbz", "legacy.knot"), &ctx, now()).await; 1755 assert!( 1756 matches!(legacy.op, PendingOp::Upsert { .. }), 1757 "member record for a legacy knot must be ingested" 1758 ); 1759 assert!( 1760 registry.hosts().contains(&KnotHostKey::new("legacy.knot")), 1761 "a member record seeds host discovery even before any repo is seen" 1762 ); 1763 assert_eq!( 1764 registry 1765 .drain_legacy_members(&KnotHostKey::new("legacy.knot")) 1766 .len(), 1767 1, 1768 "legacy member edge is indexed for later purge once the knot upgrades" 1769 ); 1770 } 1771 1772 #[tokio::test] 1773 async fn create_then_delete_round_trips_a_star() { 1774 let (store, issue_states, pull_statuses, cov, resolver) = fresh(); 1775 let create: HydrantFrame = parse_frame(json!({ 1776 "id": 10, 1777 "type": "record", 1778 "record": { 1779 "live": false, 1780 "did": "did:plc:olaren", 1781 "rev": fresh_tid().as_str(), 1782 "collection": "sh.tangled.feed.star", 1783 "rkey": "abcabcabcabcz", 1784 "action": "create", 1785 "record": { 1786 "$type": "sh.tangled.feed.star", 1787 "createdAt": "2026-05-01T00:00:00Z", 1788 "subjectDid": "did:plc:abalone" 1789 } 1790 } 1791 })); 1792 handle_frame( 1793 create, 1794 &store, 1795 &issue_states, 1796 &pull_statuses, 1797 &cov, 1798 &NoopSearchSink, 1799 &NoopRecordStore, 1800 &resolver, 1801 &sys_clock(), 1802 now(), 1803 ) 1804 .await; 1805 let key = bobbin_types::ids::EdgeKey::new( 1806 Nsid::new_static("sh.tangled.feed.star").unwrap(), 1807 did_subj("did:plc:abalone"), 1808 ); 1809 assert_eq!(store.count(&key), 1); 1810 1811 let delete: HydrantFrame = parse_frame(json!({ 1812 "id": 11, 1813 "type": "record", 1814 "record": { 1815 "live": false, 1816 "did": "did:plc:olaren", 1817 "rev": fresh_tid().as_str(), 1818 "collection": "sh.tangled.feed.star", 1819 "rkey": "abcabcabcabcz", 1820 "action": "delete", 1821 "record": null 1822 } 1823 })); 1824 handle_frame( 1825 delete, 1826 &store, 1827 &issue_states, 1828 &pull_statuses, 1829 &cov, 1830 &NoopSearchSink, 1831 &NoopRecordStore, 1832 &resolver, 1833 &sys_clock(), 1834 now(), 1835 ) 1836 .await; 1837 assert_eq!(store.count(&key), 0); 1838 } 1839 1840 #[tokio::test] 1841 async fn update_replaces_prior_edges() { 1842 let (store, issue_states, pull_statuses, cov, resolver) = fresh(); 1843 let mk = |subject_did: &Did<DefaultStr>, id: u64| -> HydrantFrame { 1844 parse_frame(json!({ 1845 "id": id, 1846 "type": "record", 1847 "record": { 1848 "live": false, 1849 "did": "did:plc:olaren", 1850 "rev": fresh_tid().as_str(), 1851 "collection": "sh.tangled.feed.star", 1852 "rkey": "abcabcabcabcz", 1853 "action": "update", 1854 "record": { 1855 "$type": "sh.tangled.feed.star", 1856 "createdAt": "2026-05-01T00:00:00Z", 1857 "subjectDid": subject_did.as_ref() 1858 } 1859 } 1860 })) 1861 }; 1862 handle_frame( 1863 mk(&Did::new_owned("did:plc:abalone").unwrap(), 1), 1864 &store, 1865 &issue_states, 1866 &pull_statuses, 1867 &cov, 1868 &NoopSearchSink, 1869 &NoopRecordStore, 1870 &resolver, 1871 &sys_clock(), 1872 now(), 1873 ) 1874 .await; 1875 handle_frame( 1876 mk(&Did::new_owned("did:plc:uni").unwrap(), 2), 1877 &store, 1878 &issue_states, 1879 &pull_statuses, 1880 &cov, 1881 &NoopSearchSink, 1882 &NoopRecordStore, 1883 &resolver, 1884 &sys_clock(), 1885 now(), 1886 ) 1887 .await; 1888 1889 let kind = Nsid::new_static("sh.tangled.feed.star").unwrap(); 1890 let old = bobbin_types::ids::EdgeKey::new(kind.clone(), did_subj("did:plc:abalone")); 1891 let new = bobbin_types::ids::EdgeKey::new(kind, did_subj("did:plc:uni")); 1892 assert_eq!(store.count(&old), 0); 1893 assert_eq!(store.count(&new), 1); 1894 } 1895 1896 #[tokio::test] 1897 async fn prepare_frame_tags_regime_from_live_flag() { 1898 let (store, issue_states, pull_statuses, cov, resolver) = fresh(); 1899 let search = NoopSearchSink; 1900 let records = NoopRecordStore; 1901 let ctx = PipelineCtx { 1902 resolver: &resolver, 1903 store: &store, 1904 issue_states: &issue_states, 1905 pull_statuses: &pull_statuses, 1906 coverage: &cov, 1907 records: &records, 1908 search: &search, 1909 shadow: None, 1910 buffer: None, 1911 knot_registry: None, 1912 knot_gate: None, 1913 }; 1914 let mk = |live: bool| -> HydrantFrame { 1915 parse_frame(json!({ 1916 "id": 1, 1917 "type": "record", 1918 "record": { 1919 "live": live, 1920 "did": "did:plc:olaren", 1921 "rev": fresh_tid().as_str(), 1922 "collection": "sh.tangled.feed.star", 1923 "rkey": "abcabcabcabcz", 1924 "action": "create", 1925 "record": { 1926 "$type": "sh.tangled.feed.star", 1927 "createdAt": "2026-05-01T00:00:00Z", 1928 "subjectDid": "did:plc:abalone" 1929 } 1930 } 1931 })) 1932 }; 1933 let live_pending = prepare_frame(mk(true), &ctx, now()).await; 1934 assert_eq!(live_pending.regime, Regime::Live); 1935 let replay_pending = prepare_frame(mk(false), &ctx, now()).await; 1936 assert_eq!(replay_pending.regime, Regime::Replay); 1937 1938 let identity: HydrantFrame = parse_frame(json!({ 1939 "id": 9, 1940 "type": "identity", 1941 })); 1942 let id_pending = prepare_frame(identity, &ctx, now()).await; 1943 assert_eq!(id_pending.regime, Regime::NonRecord); 1944 } 1945 1946 #[tokio::test] 1947 async fn create_with_cid_warms_record_lru() { 1948 let (store, issue_states, pull_statuses, cov, resolver) = fresh(); 1949 let lru = LruRecordStore::new(CacheCapacity::from_bytes(64 * 1024)); 1950 let frame: HydrantFrame = parse_frame(json!({ 1951 "id": 1, 1952 "type": "record", 1953 "record": { 1954 "live": false, 1955 "did": "did:plc:olaren", 1956 "rev": fresh_tid().as_str(), 1957 "collection": "sh.tangled.feed.star", 1958 "rkey": "abcabcabcabcz", 1959 "action": "create", 1960 "cid": VALID_CID, 1961 "record": { 1962 "$type": "sh.tangled.feed.star", 1963 "createdAt": "2026-05-01T00:00:00Z", 1964 "subjectDid": "did:plc:abalone" 1965 } 1966 } 1967 })); 1968 let source = 1969 AtUri::new_owned("at://did:plc:olaren/sh.tangled.feed.star/abcabcabcabcz").unwrap(); 1970 handle_frame( 1971 frame, 1972 &store, 1973 &issue_states, 1974 &pull_statuses, 1975 &cov, 1976 &NoopSearchSink, 1977 &lru, 1978 &resolver, 1979 &sys_clock(), 1980 now(), 1981 ) 1982 .await; 1983 let cached = lru.get(&source).expect("hydrant cid must seed the lru"); 1984 assert_eq!(cached.cid.as_ref(), VALID_CID); 1985 let parsed: serde_json::Value = serde_json::from_slice(&cached.value).unwrap(); 1986 assert_eq!( 1987 parsed["subject"]["did"], "did:plc:abalone", 1988 "legacy wire is upgraded to canon shape before caching so downstream readers see canonical fields" 1989 ); 1990 } 1991 1992 #[tokio::test] 1993 async fn create_without_cid_clears_record_lru() { 1994 let (store, issue_states, pull_statuses, cov, resolver) = fresh(); 1995 let lru = LruRecordStore::new(CacheCapacity::from_bytes(64 * 1024)); 1996 let source = 1997 AtUri::new_owned("at://did:plc:olaren/sh.tangled.feed.star/abcabcabcabcz").unwrap(); 1998 let cid: Cid<DefaultStr> = VALID_CID.parse().unwrap(); 1999 lru.put( 2000 source.clone(), 2001 Arc::new(RecordBody { 2002 uri: source.clone(), 2003 cid, 2004 value: bytes::Bytes::from_static(b"{\"stale\":true}"), 2005 }), 2006 ); 2007 let frame: HydrantFrame = parse_frame(json!({ 2008 "id": 1, 2009 "type": "record", 2010 "record": { 2011 "live": false, 2012 "did": "did:plc:olaren", 2013 "rev": fresh_tid().as_str(), 2014 "collection": "sh.tangled.feed.star", 2015 "rkey": "abcabcabcabcz", 2016 "action": "update", 2017 "record": { 2018 "$type": "sh.tangled.feed.star", 2019 "createdAt": "2026-05-01T00:00:00Z", 2020 "subjectDid": "did:plc:abalone" 2021 } 2022 } 2023 })); 2024 handle_frame( 2025 frame, 2026 &store, 2027 &issue_states, 2028 &pull_statuses, 2029 &cov, 2030 &NoopSearchSink, 2031 &lru, 2032 &resolver, 2033 &sys_clock(), 2034 now(), 2035 ) 2036 .await; 2037 assert!( 2038 lru.get(&source).is_none(), 2039 "missing cid means we cannot trust the body, so the lru must be cleared", 2040 ); 2041 } 2042 2043 #[tokio::test] 2044 async fn delete_evicts_record_lru() { 2045 let (store, issue_states, pull_statuses, cov, resolver) = fresh(); 2046 let lru = LruRecordStore::new(CacheCapacity::from_bytes(64 * 1024)); 2047 let source = 2048 AtUri::new_owned("at://did:plc:olaren/sh.tangled.feed.star/abcabcabcabcz").unwrap(); 2049 let cid: Cid<DefaultStr> = VALID_CID.parse().unwrap(); 2050 lru.put( 2051 source.clone(), 2052 Arc::new(RecordBody { 2053 uri: source.clone(), 2054 cid, 2055 value: bytes::Bytes::from_static(b"{}"), 2056 }), 2057 ); 2058 let frame: HydrantFrame = parse_frame(json!({ 2059 "id": 1, 2060 "type": "record", 2061 "record": { 2062 "live": false, 2063 "did": "did:plc:olaren", 2064 "rev": fresh_tid().as_str(), 2065 "collection": "sh.tangled.feed.star", 2066 "rkey": "abcabcabcabcz", 2067 "action": "delete", 2068 "record": null 2069 } 2070 })); 2071 handle_frame( 2072 frame, 2073 &store, 2074 &issue_states, 2075 &pull_statuses, 2076 &cov, 2077 &NoopSearchSink, 2078 &lru, 2079 &resolver, 2080 &sys_clock(), 2081 now(), 2082 ) 2083 .await; 2084 assert!(lru.get(&source).is_none()); 2085 } 2086 2087 #[tokio::test] 2088 async fn live_recent_event_promotes_coverage_to_ready() { 2089 let (store, issue_states, pull_statuses, cov, resolver) = fresh(); 2090 assert!(!cov.snapshot().is_ready()); 2091 let frame: HydrantFrame = parse_frame(json!({ 2092 "id": 99, 2093 "type": "record", 2094 "record": { 2095 "live": true, 2096 "did": "did:plc:olaren", 2097 "rev": fresh_tid().as_str(), 2098 "collection": "sh.tangled.feed.star", 2099 "rkey": "abcabcabcabcz", 2100 "action": "create", 2101 "record": { 2102 "$type": "sh.tangled.feed.star", 2103 "createdAt": "2026-05-01T00:00:00Z", 2104 "subjectDid": "did:plc:abalone" 2105 } 2106 } 2107 })); 2108 handle_frame( 2109 frame, 2110 &store, 2111 &issue_states, 2112 &pull_statuses, 2113 &cov, 2114 &NoopSearchSink, 2115 &NoopRecordStore, 2116 &resolver, 2117 &sys_clock(), 2118 now(), 2119 ) 2120 .await; 2121 assert!(cov.snapshot().is_ready()); 2122 assert_eq!(cov.snapshot().last_cursor(), HydrantCursor::new(99)); 2123 } 2124 2125 #[tokio::test] 2126 async fn live_but_stale_rev_does_not_promote() { 2127 let (store, issue_states, pull_statuses, cov, resolver) = fresh(); 2128 let stale_tid = Tid::from_time(1_000_000, 0); 2129 let frame: HydrantFrame = parse_frame(json!({ 2130 "id": 7, 2131 "type": "record", 2132 "record": { 2133 "live": true, 2134 "did": "did:plc:olaren", 2135 "rev": stale_tid.as_str(), 2136 "collection": "sh.tangled.feed.star", 2137 "rkey": "abcabcabcabcz", 2138 "action": "create", 2139 "record": { 2140 "$type": "sh.tangled.feed.star", 2141 "createdAt": "2026-05-01T00:00:00Z", 2142 "subjectDid": "did:plc:abalone" 2143 } 2144 } 2145 })); 2146 handle_frame( 2147 frame, 2148 &store, 2149 &issue_states, 2150 &pull_statuses, 2151 &cov, 2152 &NoopSearchSink, 2153 &NoopRecordStore, 2154 &resolver, 2155 &sys_clock(), 2156 now(), 2157 ) 2158 .await; 2159 assert!(!cov.snapshot().is_ready()); 2160 assert!(matches!(cov.snapshot(), Coverage::Warming { .. })); 2161 } 2162 2163 #[test] 2164 fn classify_consumer_too_slow_frame_returns_typed_variant() { 2165 let text = r#"{"type":"error","error":"ConsumerTooSlow","message":"stream socket send blocked for at least 30 seconds"}"#; 2166 match classify_text_frame(text) { 2167 Err(IngestError::ConsumerTooSlow { message }) => { 2168 assert!( 2169 message.contains("30 seconds"), 2170 "message field preserved verbatim, got: {message}" 2171 ); 2172 } 2173 other => panic!("expected ConsumerTooSlow variant, got: {other:?}"), 2174 } 2175 } 2176 2177 #[test] 2178 fn classify_unknown_hydrant_error_falls_back_to_generic_variant() { 2179 let text = r#"{"type":"error","error":"NewFutureCode","message":"some new failure mode"}"#; 2180 match classify_text_frame(text) { 2181 Err(IngestError::HydrantStream { code, message }) => { 2182 assert_eq!(code, "NewFutureCode"); 2183 assert_eq!(message, "some new failure mode"); 2184 } 2185 other => panic!("expected HydrantStream variant, got: {other:?}"), 2186 } 2187 } 2188 2189 #[test] 2190 fn classify_error_frame_without_message_uses_empty_string() { 2191 let text = r#"{"type":"error","error":"ConsumerTooSlow"}"#; 2192 match classify_text_frame(text) { 2193 Err(IngestError::ConsumerTooSlow { message }) => assert!(message.is_empty()), 2194 other => panic!("expected ConsumerTooSlow with empty message, got: {other:?}"), 2195 } 2196 } 2197 2198 #[test] 2199 fn classify_normal_record_frame_unchanged() { 2200 let text = r#"{"id":42,"type":"record"}"#; 2201 let frame = classify_text_frame(text).expect("normal record frame must parse"); 2202 assert_eq!(frame.id, 42); 2203 assert_eq!(frame.kind, FrameKind::Record); 2204 } 2205 2206 #[test] 2207 fn classify_garbage_object_returns_decode_error() { 2208 let text = r#"{"random":"object","without":"required fields"}"#; 2209 match classify_text_frame(text) { 2210 Err(IngestError::Decode(_)) => {} 2211 other => panic!("expected Decode error, got: {other:?}"), 2212 } 2213 } 2214 2215 struct ScriptedWsStream { 2216 messages: std::collections::VecDeque<Result<WsMessage, NetworkError>>, 2217 } 2218 2219 impl WsStream for ScriptedWsStream { 2220 fn next<'a>(&'a mut self) -> bobbin_runtime::WsMessageFuture<'a> { 2221 let msg = self.messages.pop_front(); 2222 Box::pin(async move { msg }) 2223 } 2224 } 2225 2226 #[tokio::test] 2227 async fn reader_loop_surfaces_consumer_too_slow_from_error_frame() { 2228 let mut messages = std::collections::VecDeque::new(); 2229 messages.push_back(Ok(WsMessage::Text( 2230 r#"{"type":"error","error":"ConsumerTooSlow","message":"stream delivery blocked"}"# 2231 .to_owned(), 2232 ))); 2233 let stream: Box<dyn WsStream> = Box::new(ScriptedWsStream { messages }); 2234 let (frame_tx, _frame_rx) = tokio::sync::mpsc::channel::<HydrantFrame>(FRAME_CHANNEL_DEPTH); 2235 let (control_tx, _control_rx) = 2236 tokio::sync::mpsc::channel::<WsEvent>(CONTROL_CHANNEL_DEPTH); 2237 let cancel = CancellationToken::new(); 2238 2239 let end = reader_loop(stream, frame_tx, control_tx, cancel).await; 2240 2241 match end.error { 2242 Some(IngestError::ConsumerTooSlow { message }) => { 2243 assert_eq!(message, "stream delivery blocked"); 2244 } 2245 other => panic!("expected ConsumerTooSlow SessionEnd error, got: {other:?}"), 2246 } 2247 assert_eq!(end.outcome, SessionOutcome::Empty); 2248 } 2249 2250 #[tokio::test] 2251 async fn reader_loop_surfaces_unknown_hydrant_error_distinctly() { 2252 let mut messages = std::collections::VecDeque::new(); 2253 messages.push_back(Ok(WsMessage::Text( 2254 r#"{"type":"error","error":"NewFutureCode","message":"new mode"}"#.to_owned(), 2255 ))); 2256 let stream: Box<dyn WsStream> = Box::new(ScriptedWsStream { messages }); 2257 let (frame_tx, _frame_rx) = tokio::sync::mpsc::channel::<HydrantFrame>(FRAME_CHANNEL_DEPTH); 2258 let (control_tx, _control_rx) = 2259 tokio::sync::mpsc::channel::<WsEvent>(CONTROL_CHANNEL_DEPTH); 2260 let cancel = CancellationToken::new(); 2261 2262 let end = reader_loop(stream, frame_tx, control_tx, cancel).await; 2263 2264 match end.error { 2265 Some(IngestError::HydrantStream { code, message }) => { 2266 assert_eq!(code, "NewFutureCode"); 2267 assert_eq!(message, "new mode"); 2268 } 2269 other => panic!("expected HydrantStream SessionEnd error, got: {other:?}"), 2270 } 2271 } 2272 2273 struct ChannelWsStream { 2274 rx: tokio::sync::mpsc::Receiver<Result<WsMessage, NetworkError>>, 2275 } 2276 2277 impl WsStream for ChannelWsStream { 2278 fn next<'a>(&'a mut self) -> bobbin_runtime::WsMessageFuture<'a> { 2279 Box::pin(async move { self.rx.recv().await }) 2280 } 2281 } 2282 2283 fn star_frame_text(id: u64, rkey: &Rkey<DefaultStr>) -> String { 2284 json!({ 2285 "id": id, 2286 "type": "record", 2287 "record": { 2288 "live": false, 2289 "did": "did:plc:olaren", 2290 "rev": fresh_tid().as_str(), 2291 "collection": "sh.tangled.feed.star", 2292 "rkey": rkey.as_ref(), 2293 "action": "create", 2294 "record": { 2295 "$type": "sh.tangled.feed.star", 2296 "createdAt": "2026-05-01T00:00:00Z", 2297 "subjectDid": "did:plc:abalone" 2298 } 2299 } 2300 }) 2301 .to_string() 2302 } 2303 2304 #[tokio::test] 2305 async fn pong_forwards_promptly_when_frame_channel_is_full() { 2306 let (ws_tx, ws_rx) = tokio::sync::mpsc::channel::<Result<WsMessage, NetworkError>>(8); 2307 let stream: Box<dyn WsStream> = Box::new(ChannelWsStream { rx: ws_rx }); 2308 let (frame_tx, frame_rx) = tokio::sync::mpsc::channel::<HydrantFrame>(1); 2309 let (control_tx, mut control_rx) = 2310 tokio::sync::mpsc::channel::<WsEvent>(CONTROL_CHANNEL_DEPTH); 2311 let cancel = CancellationToken::new(); 2312 2313 let prefill: HydrantFrame = parse_frame(json!({ 2314 "id": 0, 2315 "type": "record", 2316 "record": { 2317 "live": false, 2318 "did": "did:plc:olaren", 2319 "rev": fresh_tid().as_str(), 2320 "collection": "sh.tangled.feed.star", 2321 "rkey": "prefilrkey001", 2322 "action": "create", 2323 "record": { 2324 "$type": "sh.tangled.feed.star", 2325 "createdAt": "2026-05-01T00:00:00Z", 2326 "subjectDid": "did:plc:abalone" 2327 } 2328 } 2329 })); 2330 frame_tx 2331 .try_send(prefill) 2332 .expect("depth-1 frame channel must accept the prefill"); 2333 2334 let reader_handle = tokio::spawn(reader_loop( 2335 stream, 2336 frame_tx.clone(), 2337 control_tx.clone(), 2338 cancel.clone(), 2339 )); 2340 2341 ws_tx 2342 .send(Ok(WsMessage::Text(star_frame_text( 2343 1, 2344 &rkey("starrkeyaa001"), 2345 )))) 2346 .await 2347 .unwrap(); 2348 ws_tx.send(Ok(WsMessage::Pong(Bytes::new()))).await.unwrap(); 2349 2350 let pong_event = tokio::time::timeout(Duration::from_millis(500), control_rx.recv()) 2351 .await 2352 .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") 2353 .expect("control_tx was not closed"); 2354 assert!( 2355 matches!(pong_event, WsEvent::IncomingPong), 2356 "first control event must be the pong, not a held text", 2357 ); 2358 2359 let too_slow = 2360 "{\"type\":\"error\",\"error\":\"ConsumerTooSlow\",\"message\":\"saturated\"}" 2361 .to_string(); 2362 ws_tx.send(Ok(WsMessage::Text(too_slow))).await.unwrap(); 2363 2364 let end = tokio::time::timeout(Duration::from_secs(1), reader_handle) 2365 .await 2366 .expect("reader must exit promptly once ConsumerTooSlow is read") 2367 .expect("reader task must not panic"); 2368 match end.error { 2369 Some(IngestError::ConsumerTooSlow { message }) => assert_eq!(message, "saturated"), 2370 other => panic!("expected ConsumerTooSlow disconnect, got: {other:?}"), 2371 } 2372 2373 drop(frame_rx); 2374 } 2375 2376 #[tokio::test] 2377 async fn reader_drains_held_frames_once_processor_catches_up() { 2378 let (ws_tx, ws_rx) = tokio::sync::mpsc::channel::<Result<WsMessage, NetworkError>>(8); 2379 let stream: Box<dyn WsStream> = Box::new(ChannelWsStream { rx: ws_rx }); 2380 let (frame_tx, mut frame_rx) = tokio::sync::mpsc::channel::<HydrantFrame>(1); 2381 let (control_tx, _control_rx) = 2382 tokio::sync::mpsc::channel::<WsEvent>(CONTROL_CHANNEL_DEPTH); 2383 let cancel = CancellationToken::new(); 2384 2385 let prefill: HydrantFrame = parse_frame(json!({ 2386 "id": 0, 2387 "type": "record", 2388 "record": { 2389 "live": false, 2390 "did": "did:plc:olaren", 2391 "rev": fresh_tid().as_str(), 2392 "collection": "sh.tangled.feed.star", 2393 "rkey": "prefilrkey001", 2394 "action": "create", 2395 "record": { 2396 "$type": "sh.tangled.feed.star", 2397 "createdAt": "2026-05-01T00:00:00Z", 2398 "subjectDid": "did:plc:abalone" 2399 } 2400 } 2401 })); 2402 frame_tx.try_send(prefill).unwrap(); 2403 2404 let reader_handle = tokio::spawn(reader_loop( 2405 stream, 2406 frame_tx.clone(), 2407 control_tx, 2408 cancel.clone(), 2409 )); 2410 2411 ws_tx 2412 .send(Ok(WsMessage::Text(star_frame_text( 2413 1, 2414 &rkey("heldrkeyaa001"), 2415 )))) 2416 .await 2417 .unwrap(); 2418 ws_tx 2419 .send(Ok(WsMessage::Text(star_frame_text( 2420 2, 2421 &rkey("heldrkeyaa002"), 2422 )))) 2423 .await 2424 .unwrap(); 2425 2426 let _drained_prefill = frame_rx.recv().await.expect("prefilled frame drains"); 2427 let first = tokio::time::timeout(Duration::from_millis(500), frame_rx.recv()) 2428 .await 2429 .expect("first held frame must reach frame_rx after slot opens") 2430 .expect("frame_tx still open"); 2431 assert_eq!(first.id, 1); 2432 let second = tokio::time::timeout(Duration::from_millis(500), frame_rx.recv()) 2433 .await 2434 .expect("second held frame must reach frame_rx after slot opens") 2435 .expect("frame_tx still open"); 2436 assert_eq!(second.id, 2); 2437 2438 cancel.cancel(); 2439 let _ = tokio::time::timeout(Duration::from_secs(1), reader_handle) 2440 .await 2441 .expect("reader must stop after cancel"); 2442 } 2443 2444 #[test] 2445 fn first_connect_uses_configured_start_cursor() { 2446 let start = HydrantCursor::new(42); 2447 assert_eq!(next_connect_cursor(Coverage::default(), start), start); 2448 } 2449 2450 #[test] 2451 fn reconnect_resumes_strictly_after_last_seen() { 2452 let snap = Coverage::default().advance(HydrantCursor::new(7)); 2453 assert_eq!( 2454 next_connect_cursor(snap, HydrantCursor::new(0)), 2455 HydrantCursor::new(8), 2456 ); 2457 } 2458 2459 #[test] 2460 fn reconnect_overrides_configured_start() { 2461 let snap = Coverage::default().advance(HydrantCursor::new(100)); 2462 assert_eq!( 2463 next_connect_cursor(snap, HydrantCursor::new(50)), 2464 HydrantCursor::new(101), 2465 ); 2466 } 2467 2468 #[test] 2469 fn first_connect_uses_start_even_when_first_frame_id_would_be_zero() { 2470 let start = HydrantCursor::new(7); 2471 let snap = Coverage::default(); 2472 assert_eq!(snap.last_cursor(), HydrantCursor::new(0)); 2473 assert_eq!(snap.events_processed(), 0); 2474 assert_eq!(next_connect_cursor(snap, start), start); 2475 } 2476 2477 #[test] 2478 fn reconnect_after_processing_id_zero_advances_to_one() { 2479 let snap = Coverage::default().advance(HydrantCursor::new(0)); 2480 assert_eq!(snap.events_processed(), 1); 2481 assert_eq!( 2482 next_connect_cursor(snap, HydrantCursor::new(99)), 2483 HydrantCursor::new(1), 2484 "events_processed disambiguates 'never seen' from 'saw id 0'", 2485 ); 2486 } 2487 2488 #[tokio::test] 2489 async fn identity_frame_advances_cursor_only() { 2490 let (store, issue_states, pull_statuses, cov, resolver) = fresh(); 2491 let frame: HydrantFrame = parse_frame(json!({ 2492 "id": 5, 2493 "type": "identity", 2494 "identity": { 2495 "did": "did:plc:olaren", 2496 "handle": "olaren.dev" 2497 } 2498 })); 2499 handle_frame( 2500 frame, 2501 &store, 2502 &issue_states, 2503 &pull_statuses, 2504 &cov, 2505 &NoopSearchSink, 2506 &NoopRecordStore, 2507 &resolver, 2508 &sys_clock(), 2509 now(), 2510 ) 2511 .await; 2512 assert_eq!(store.key_count(), 0); 2513 assert_eq!(cov.snapshot().last_cursor(), HydrantCursor::new(5)); 2514 assert!(!cov.snapshot().is_ready()); 2515 } 2516 2517 #[tokio::test] 2518 async fn account_frame_advances_cursor_only() { 2519 let (store, issue_states, pull_statuses, cov, resolver) = fresh(); 2520 let frame: HydrantFrame = parse_frame(json!({ 2521 "id": 6, 2522 "type": "account", 2523 "account": {"did": "did:plc:olaren", "active": true} 2524 })); 2525 handle_frame( 2526 frame, 2527 &store, 2528 &issue_states, 2529 &pull_statuses, 2530 &cov, 2531 &NoopSearchSink, 2532 &NoopRecordStore, 2533 &resolver, 2534 &sys_clock(), 2535 now(), 2536 ) 2537 .await; 2538 assert_eq!(store.key_count(), 0); 2539 assert_eq!(cov.snapshot().last_cursor(), HydrantCursor::new(6)); 2540 } 2541 2542 #[tokio::test] 2543 async fn unknown_frame_kind_advances_cursor_without_panic() { 2544 let (store, issue_states, pull_statuses, cov, resolver) = fresh(); 2545 let frame: HydrantFrame = parse_frame(json!({"id": 8, "type": "future_event"})); 2546 assert_eq!(frame.kind, FrameKind::Other); 2547 handle_frame( 2548 frame, 2549 &store, 2550 &issue_states, 2551 &pull_statuses, 2552 &cov, 2553 &NoopSearchSink, 2554 &NoopRecordStore, 2555 &resolver, 2556 &sys_clock(), 2557 now(), 2558 ) 2559 .await; 2560 assert_eq!(cov.snapshot().last_cursor(), HydrantCursor::new(8)); 2561 } 2562 2563 fn fresh_runtime(cancel: CancellationToken) -> IngestRuntime<NoopSearchSink> { 2564 IngestRuntime { 2565 store: Arc::new(EdgeStore::new(RuntimeHasher::default())), 2566 issue_states: Arc::new(StateIndex::new(RuntimeHasher::default())), 2567 pull_statuses: Arc::new(StateIndex::new(RuntimeHasher::default())), 2568 coverage: Arc::new(CoverageWatch::new()), 2569 search: Arc::new(NoopSearchSink), 2570 records: Arc::new(NoopRecordStore) as Arc<dyn RecordStore>, 2571 resolver: Arc::new(RepoIdResolver::detached(RuntimeHasher::default())), 2572 clock: Arc::new(SystemClock::new()), 2573 entropy: Arc::new(OsEntropy), 2574 ws: TungsteniteWs::shared(), 2575 cancel, 2576 disconnects: None, 2577 warming_shadow: None, 2578 warming_buffer: None, 2579 knot_registry: None, 2580 knot_gate: None, 2581 } 2582 } 2583 2584 #[tokio::test(start_paused = true)] 2585 async fn cancel_token_short_circuits_reconnect_sleep() { 2586 let cfg = IngestConfig::new(Url::parse("ws://127.0.0.1:1").unwrap()); 2587 let cancel = CancellationToken::new(); 2588 let runtime = fresh_runtime(cancel.clone()); 2589 let task = tokio::spawn(async move { run(cfg, runtime).await }); 2590 tokio::time::advance(Duration::from_millis(10)).await; 2591 cancel.cancel(); 2592 let outcome = tokio::time::timeout(Duration::from_secs(1), task) 2593 .await 2594 .expect("ingest must stop within timeout once cancel fires"); 2595 assert!(matches!(outcome, Ok(Ok(()))), "got {outcome:?}"); 2596 } 2597 2598 #[test] 2599 fn jittered_stays_within_one_quarter_of_base() { 2600 let base = Duration::from_secs(1); 2601 let cap = base + Duration::from_millis(250); 2602 let entropy = OsEntropy; 2603 (0..50).for_each(|_| { 2604 let j = jittered(base, &entropy); 2605 assert!(j >= base, "jitter must not undershoot"); 2606 assert!(j <= cap, "jitter must not exceed +25%, got {:?}", j); 2607 }); 2608 } 2609 2610 #[tokio::test] 2611 async fn star_after_observed_repo_keys_on_repo_did() { 2612 let (store, issue_states, pull_statuses, cov, resolver) = fresh(); 2613 let repo: HydrantFrame = parse_frame(json!({ 2614 "id": 1, 2615 "type": "record", 2616 "record": { 2617 "live": false, 2618 "did": "did:plc:nel", 2619 "rev": fresh_tid().as_str(), 2620 "collection": "sh.tangled.repo", 2621 "rkey": "abcabcabcabcz", 2622 "action": "create", 2623 "record": { 2624 "$type": "sh.tangled.repo", 2625 "createdAt": "2026-05-01T00:00:00Z", 2626 "knot": "oyster.cafe", 2627 "name": "abalone", 2628 "repoDid": "did:plc:abalone" 2629 } 2630 } 2631 })); 2632 handle_frame( 2633 repo, 2634 &store, 2635 &issue_states, 2636 &pull_statuses, 2637 &cov, 2638 &NoopSearchSink, 2639 &NoopRecordStore, 2640 &resolver, 2641 &sys_clock(), 2642 now(), 2643 ) 2644 .await; 2645 2646 let star: HydrantFrame = parse_frame(json!({ 2647 "id": 2, 2648 "type": "record", 2649 "record": { 2650 "live": false, 2651 "did": "did:plc:olaren", 2652 "rev": fresh_tid().as_str(), 2653 "collection": "sh.tangled.feed.star", 2654 "rkey": "abcabcabcabcz", 2655 "action": "create", 2656 "record": { 2657 "$type": "sh.tangled.feed.star", 2658 "createdAt": "2026-05-01T00:00:00Z", 2659 "subject": "at://did:plc:nel/sh.tangled.repo/abcabcabcabcz" 2660 } 2661 } 2662 })); 2663 handle_frame( 2664 star, 2665 &store, 2666 &issue_states, 2667 &pull_statuses, 2668 &cov, 2669 &NoopSearchSink, 2670 &NoopRecordStore, 2671 &resolver, 2672 &sys_clock(), 2673 now(), 2674 ) 2675 .await; 2676 2677 let nsid = Nsid::new_static("sh.tangled.feed.star").unwrap(); 2678 let repo_keyed = bobbin_types::ids::EdgeKey::new(nsid.clone(), did_subj("did:plc:abalone")); 2679 let owner_keyed = bobbin_types::ids::EdgeKey::new(nsid, did_subj("did:plc:nel")); 2680 assert_eq!( 2681 store.count(&repo_keyed), 2682 1, 2683 "star should be keyed on repoDID once the repo is observed" 2684 ); 2685 assert_eq!( 2686 store.count(&owner_keyed), 2687 0, 2688 "owner DID should not collect the edge" 2689 ); 2690 } 2691 2692 #[tokio::test] 2693 async fn unresolvable_repo_subject_drops_edge() { 2694 let (store, issue_states, pull_statuses, cov, resolver) = fresh(); 2695 let star: HydrantFrame = parse_frame(json!({ 2696 "id": 1, 2697 "type": "record", 2698 "record": { 2699 "live": false, 2700 "did": "did:plc:olaren", 2701 "rev": fresh_tid().as_str(), 2702 "collection": "sh.tangled.feed.star", 2703 "rkey": "abcabcabcabcz", 2704 "action": "create", 2705 "record": { 2706 "$type": "sh.tangled.feed.star", 2707 "createdAt": "2026-05-01T00:00:00Z", 2708 "subject": "at://did:plc:nel/sh.tangled.repo/abcabcabcabcz" 2709 } 2710 } 2711 })); 2712 handle_frame( 2713 star, 2714 &store, 2715 &issue_states, 2716 &pull_statuses, 2717 &cov, 2718 &NoopSearchSink, 2719 &NoopRecordStore, 2720 &resolver, 2721 &sys_clock(), 2722 now(), 2723 ) 2724 .await; 2725 2726 let nsid = Nsid::new_static("sh.tangled.feed.star").unwrap(); 2727 let owner_keyed = bobbin_types::ids::EdgeKey::new(nsid.clone(), did_subj("did:plc:nel")); 2728 let uri_keyed = bobbin_types::ids::EdgeKey::new( 2729 nsid, 2730 uri_subj("at://did:plc:nel/sh.tangled.repo/abcabcabcabcz"), 2731 ); 2732 assert_eq!( 2733 store.count(&owner_keyed), 2734 0, 2735 "must not silently misfile under the authoring DID", 2736 ); 2737 assert_eq!( 2738 store.count(&uri_keyed), 2739 0, 2740 "unresolvable rkey-form subject must drop the edge; keeping it would index against a key that never matches bare-DID queries", 2741 ); 2742 } 2743 2744 #[tokio::test] 2745 async fn repo_without_repo_did_drops_edge() { 2746 let (store, issue_states, pull_statuses, cov, resolver) = fresh(); 2747 let repo: HydrantFrame = parse_frame(json!({ 2748 "id": 1, 2749 "type": "record", 2750 "record": { 2751 "live": false, 2752 "did": "did:plc:nel", 2753 "rev": fresh_tid().as_str(), 2754 "collection": "sh.tangled.repo", 2755 "rkey": "abcabcabcabcz", 2756 "action": "create", 2757 "record": { 2758 "$type": "sh.tangled.repo", 2759 "createdAt": "2026-05-01T00:00:00Z", 2760 "knot": "oyster.cafe", 2761 "name": "abalone" 2762 } 2763 } 2764 })); 2765 handle_frame( 2766 repo, 2767 &store, 2768 &issue_states, 2769 &pull_statuses, 2770 &cov, 2771 &NoopSearchSink, 2772 &NoopRecordStore, 2773 &resolver, 2774 &sys_clock(), 2775 now(), 2776 ) 2777 .await; 2778 2779 let star: HydrantFrame = parse_frame(json!({ 2780 "id": 2, 2781 "type": "record", 2782 "record": { 2783 "live": false, 2784 "did": "did:plc:olaren", 2785 "rev": fresh_tid().as_str(), 2786 "collection": "sh.tangled.feed.star", 2787 "rkey": "abcabcabcabcz", 2788 "action": "create", 2789 "record": { 2790 "$type": "sh.tangled.feed.star", 2791 "createdAt": "2026-05-01T00:00:00Z", 2792 "subject": "at://did:plc:nel/sh.tangled.repo/abcabcabcabcz" 2793 } 2794 } 2795 })); 2796 handle_frame( 2797 star, 2798 &store, 2799 &issue_states, 2800 &pull_statuses, 2801 &cov, 2802 &NoopSearchSink, 2803 &NoopRecordStore, 2804 &resolver, 2805 &sys_clock(), 2806 now(), 2807 ) 2808 .await; 2809 2810 let nsid = Nsid::new_static("sh.tangled.feed.star").unwrap(); 2811 let uri_keyed = bobbin_types::ids::EdgeKey::new( 2812 nsid.clone(), 2813 uri_subj("at://did:plc:nel/sh.tangled.repo/abcabcabcabcz"), 2814 ); 2815 let owner_keyed = bobbin_types::ids::EdgeKey::new(nsid, did_subj("did:plc:nel")); 2816 assert_eq!( 2817 store.count(&uri_keyed), 2818 0, 2819 "no canonical DID exists for a repo without repoDID, so the edge must be dropped", 2820 ); 2821 assert_eq!( 2822 store.count(&owner_keyed), 2823 0, 2824 "the authoring DID is not the canonical repo identity", 2825 ); 2826 } 2827 2828 #[tokio::test] 2829 async fn explicit_subject_did_skips_normalization() { 2830 let (store, issue_states, pull_statuses, cov, resolver) = fresh(); 2831 let star: HydrantFrame = parse_frame(json!({ 2832 "id": 1, 2833 "type": "record", 2834 "record": { 2835 "live": false, 2836 "did": "did:plc:olaren", 2837 "rev": fresh_tid().as_str(), 2838 "collection": "sh.tangled.feed.star", 2839 "rkey": "abcabcabcabcz", 2840 "action": "create", 2841 "record": { 2842 "$type": "sh.tangled.feed.star", 2843 "createdAt": "2026-05-01T00:00:00Z", 2844 "subjectDid": "did:plc:abalone" 2845 } 2846 } 2847 })); 2848 handle_frame( 2849 star, 2850 &store, 2851 &issue_states, 2852 &pull_statuses, 2853 &cov, 2854 &NoopSearchSink, 2855 &NoopRecordStore, 2856 &resolver, 2857 &sys_clock(), 2858 now(), 2859 ) 2860 .await; 2861 let key = bobbin_types::ids::EdgeKey::new( 2862 Nsid::new_static("sh.tangled.feed.star").unwrap(), 2863 did_subj("did:plc:abalone"), 2864 ); 2865 assert_eq!(store.count(&key), 1); 2866 } 2867 2868 #[tokio::test] 2869 async fn issue_with_repo_uri_resolves_to_repo_did() { 2870 let (store, issue_states, pull_statuses, cov, resolver) = fresh(); 2871 resolver 2872 .observe( 2873 Did::new_owned("did:plc:nel").unwrap(), 2874 Rkey::new_owned("abcabcabcabcz").unwrap(), 2875 Some(Did::new_owned("did:plc:abalone").unwrap()), 2876 None, 2877 ) 2878 .await; 2879 let issue: HydrantFrame = parse_frame(json!({ 2880 "id": 1, 2881 "type": "record", 2882 "record": { 2883 "live": false, 2884 "did": "did:plc:olaren", 2885 "rev": fresh_tid().as_str(), 2886 "collection": "sh.tangled.repo.issue", 2887 "rkey": "abcabcabcabcz", 2888 "action": "create", 2889 "record": { 2890 "$type": "sh.tangled.repo.issue", 2891 "createdAt": "2026-05-01T00:00:00Z", 2892 "title": "bug", 2893 "repo": "at://did:plc:nel/sh.tangled.repo/abcabcabcabcz" 2894 } 2895 } 2896 })); 2897 handle_frame( 2898 issue, 2899 &store, 2900 &issue_states, 2901 &pull_statuses, 2902 &cov, 2903 &NoopSearchSink, 2904 &NoopRecordStore, 2905 &resolver, 2906 &sys_clock(), 2907 now(), 2908 ) 2909 .await; 2910 let key = bobbin_types::ids::EdgeKey::new( 2911 Nsid::new_static("sh.tangled.repo.issue").unwrap(), 2912 did_subj("did:plc:abalone"), 2913 ); 2914 assert_eq!(store.count(&key), 1); 2915 assert_eq!( 2916 store 2917 .count_issue_state(&key, IssueStateKind::Open, None) 2918 .count 2919 .get(), 2920 1, 2921 "hydrant ingest materializes the default state count", 2922 ); 2923 } 2924 2925 #[derive(Default)] 2926 struct RecordingSearchSink { 2927 docs: tokio::sync::Mutex<Vec<bobbin_types::search::SearchDoc>>, 2928 } 2929 2930 impl SearchSink for RecordingSearchSink { 2931 async fn upsert(&self, doc: bobbin_types::search::SearchDoc) { 2932 self.docs.lock().await.push(doc); 2933 } 2934 async fn remove(&self, _uri: &AtUri<DefaultStr>) {} 2935 } 2936 2937 #[tokio::test] 2938 async fn search_index_hydrates_repo_did_via_resolver() { 2939 let (store, issue_states, pull_statuses, cov, resolver) = fresh(); 2940 resolver 2941 .observe( 2942 Did::new_owned("did:plc:nel").unwrap(), 2943 Rkey::new_owned("abcabcabcabcz").unwrap(), 2944 Some(Did::new_owned("did:plc:abalone").unwrap()), 2945 None, 2946 ) 2947 .await; 2948 let search = RecordingSearchSink::default(); 2949 let issue: HydrantFrame = parse_frame(json!({ 2950 "id": 1, 2951 "type": "record", 2952 "record": { 2953 "live": false, 2954 "did": "did:plc:olaren", 2955 "rev": fresh_tid().as_str(), 2956 "collection": "sh.tangled.repo.issue", 2957 "rkey": "abcabcabcabcz", 2958 "action": "create", 2959 "record": { 2960 "$type": "sh.tangled.repo.issue", 2961 "createdAt": "2026-05-01T00:00:00Z", 2962 "title": "bug", 2963 "repo": "at://did:plc:nel/sh.tangled.repo/abcabcabcabcz" 2964 } 2965 } 2966 })); 2967 handle_frame( 2968 issue, 2969 &store, 2970 &issue_states, 2971 &pull_statuses, 2972 &cov, 2973 &search, 2974 &NoopRecordStore, 2975 &resolver, 2976 &sys_clock(), 2977 now(), 2978 ) 2979 .await; 2980 let docs = search.docs.lock().await; 2981 assert_eq!(docs.len(), 1, "issue should produce one search doc"); 2982 assert_eq!( 2983 docs[0].repo, 2984 Some(Did::new_owned("did:plc:abalone").unwrap()), 2985 "search doc repo field must be resolved from the observed repo, not left as None", 2986 ); 2987 } 2988 2989 #[tokio::test] 2990 async fn repo_record_indexes_its_rkey_and_name() { 2991 let (store, issue_states, pull_statuses, cov, resolver) = fresh(); 2992 let repo: HydrantFrame = parse_frame(json!({ 2993 "id": 1, 2994 "type": "record", 2995 "record": { 2996 "live": false, 2997 "did": "did:plc:nel", 2998 "rev": fresh_tid().as_str(), 2999 "collection": "sh.tangled.repo", 3000 "rkey": "abcabcabcabcz", 3001 "action": "create", 3002 "record": { 3003 "$type": "sh.tangled.repo", 3004 "createdAt": "2026-05-01T00:00:00Z", 3005 "knot": "oyster.cafe", 3006 "name": "abalone" 3007 } 3008 } 3009 })); 3010 handle_frame( 3011 repo, 3012 &store, 3013 &issue_states, 3014 &pull_statuses, 3015 &cov, 3016 &NoopSearchSink, 3017 &NoopRecordStore, 3018 &resolver, 3019 &sys_clock(), 3020 now(), 3021 ) 3022 .await; 3023 let owner = Did::new_owned("did:plc:nel").unwrap(); 3024 assert_eq!( 3025 resolver.lookup_by_name(&owner, "abcabcabcabcz").await, 3026 Some(bobbin_types::ids::RepoIdent::new( 3027 owner.clone(), 3028 Rkey::new_owned("abcabcabcabcz").unwrap() 3029 )), 3030 ); 3031 assert_eq!( 3032 resolver.lookup_by_name(&owner, "abalone").await, 3033 Some(bobbin_types::ids::RepoIdent::new( 3034 owner.clone(), 3035 Rkey::new_owned("abcabcabcabcz").unwrap() 3036 )), 3037 "the record's cosmetic name resolves to its rkey ident", 3038 ); 3039 } 3040 3041 #[tokio::test] 3042 async fn fork_record_edges_against_its_source() { 3043 let (store, issue_states, pull_statuses, cov, resolver) = fresh(); 3044 let fork: HydrantFrame = parse_frame(json!({ 3045 "id": 1, 3046 "type": "record", 3047 "record": { 3048 "live": false, 3049 "did": "did:plc:olaren", 3050 "rev": fresh_tid().as_str(), 3051 "collection": "sh.tangled.repo", 3052 "rkey": "abcabcabcabcz", 3053 "action": "create", 3054 "record": { 3055 "$type": "sh.tangled.repo", 3056 "createdAt": "2026-05-01T00:00:00Z", 3057 "knot": "oyster.cafe", 3058 "name": "abalone", 3059 "source": "did:plc:abalone" 3060 } 3061 } 3062 })); 3063 handle_frame( 3064 fork, 3065 &store, 3066 &issue_states, 3067 &pull_statuses, 3068 &cov, 3069 &NoopSearchSink, 3070 &NoopRecordStore, 3071 &resolver, 3072 &sys_clock(), 3073 now(), 3074 ) 3075 .await; 3076 let key = bobbin_types::ids::EdgeKey::new( 3077 Nsid::new_static(bobbin_types::edges::REPO_SOURCE_EDGE_KIND).unwrap(), 3078 did_subj("did:plc:abalone"), 3079 ); 3080 assert_eq!(store.count(&key), 1); 3081 } 3082 3083 fn fork_frame(source: &str) -> HydrantFrame { 3084 parse_frame(json!({ 3085 "id": 1, 3086 "type": "record", 3087 "record": { 3088 "live": false, 3089 "did": "did:plc:olaren", 3090 "rev": fresh_tid().as_str(), 3091 "collection": "sh.tangled.repo", 3092 "rkey": "abcabcabcabcz", 3093 "action": "create", 3094 "record": { 3095 "$type": "sh.tangled.repo", 3096 "createdAt": "2026-05-01T00:00:00Z", 3097 "knot": "oyster.cafe", 3098 "name": "abalone", 3099 "source": source 3100 } 3101 } 3102 })) 3103 } 3104 3105 fn fork_edge_key(subject: SubjectRef) -> bobbin_types::ids::EdgeKey { 3106 bobbin_types::ids::EdgeKey::new( 3107 Nsid::new_static(bobbin_types::edges::REPO_SOURCE_EDGE_KIND).unwrap(), 3108 subject, 3109 ) 3110 } 3111 3112 #[tokio::test] 3113 async fn fork_by_at_uri_edges_against_the_source_did() { 3114 let (store, issue_states, pull_statuses, cov, resolver) = fresh(); 3115 resolver 3116 .observe( 3117 Did::new_owned("did:plc:nel").unwrap(), 3118 Rkey::new_owned("core").unwrap(), 3119 Some(Did::new_owned("did:plc:abalone").unwrap()), 3120 None, 3121 ) 3122 .await; 3123 handle_frame( 3124 fork_frame("at://did:plc:nel/sh.tangled.repo/core"), 3125 &store, 3126 &issue_states, 3127 &pull_statuses, 3128 &cov, 3129 &NoopSearchSink, 3130 &NoopRecordStore, 3131 &resolver, 3132 &sys_clock(), 3133 now(), 3134 ) 3135 .await; 3136 assert_eq!(store.count(&fork_edge_key(did_subj("did:plc:abalone"))), 1); 3137 assert_eq!( 3138 store.count(&fork_edge_key(uri_subj( 3139 "at://did:plc:nel/sh.tangled.repo/core" 3140 ))), 3141 0, 3142 "the rkey form would never match a bare-DID query", 3143 ); 3144 } 3145 3146 #[tokio::test] 3147 async fn fork_of_a_repo_without_a_did_has_nothing_to_count_against() { 3148 let (store, issue_states, pull_statuses, cov, resolver) = fresh(); 3149 handle_frame( 3150 fork_frame("at://did:plc:nel/sh.tangled.repo/core"), 3151 &store, 3152 &issue_states, 3153 &pull_statuses, 3154 &cov, 3155 &NoopSearchSink, 3156 &NoopRecordStore, 3157 &resolver, 3158 &sys_clock(), 3159 now(), 3160 ) 3161 .await; 3162 assert_eq!( 3163 store.count(&fork_edge_key(uri_subj( 3164 "at://did:plc:nel/sh.tangled.repo/core" 3165 ))), 3166 0, 3167 ); 3168 } 3169 3170 #[tokio::test] 3171 async fn repo_record_without_a_name_indexes_its_rkey() { 3172 let (store, issue_states, pull_statuses, cov, resolver) = fresh(); 3173 let repo: HydrantFrame = parse_frame(json!({ 3174 "id": 1, 3175 "type": "record", 3176 "record": { 3177 "live": false, 3178 "did": "did:plc:nel", 3179 "rev": fresh_tid().as_str(), 3180 "collection": "sh.tangled.repo", 3181 "rkey": "abalone", 3182 "action": "create", 3183 "record": { 3184 "$type": "sh.tangled.repo", 3185 "createdAt": "2026-05-01T00:00:00Z", 3186 "knot": "oyster.cafe" 3187 } 3188 } 3189 })); 3190 handle_frame( 3191 repo, 3192 &store, 3193 &issue_states, 3194 &pull_statuses, 3195 &cov, 3196 &NoopSearchSink, 3197 &NoopRecordStore, 3198 &resolver, 3199 &sys_clock(), 3200 now(), 3201 ) 3202 .await; 3203 let owner = Did::new_owned("did:plc:nel").unwrap(); 3204 assert_eq!( 3205 resolver.lookup_by_name(&owner, "abalone").await, 3206 Some(bobbin_types::ids::RepoIdent::new( 3207 owner, 3208 Rkey::new_owned("abalone").unwrap() 3209 )), 3210 "repos made before the name field are only reachable by rkey", 3211 ); 3212 } 3213 3214 #[tokio::test] 3215 async fn delete_repo_record_evicts_resolver_cache() { 3216 let (store, issue_states, pull_statuses, cov, resolver) = fresh(); 3217 let owner = Did::new_owned("did:plc:nel").unwrap(); 3218 let rkey = Rkey::new_owned("abcabcabcabcz").unwrap(); 3219 resolver 3220 .observe( 3221 owner.clone(), 3222 rkey.clone(), 3223 Some(Did::new_owned("did:plc:abalone").unwrap()), 3224 None, 3225 ) 3226 .await; 3227 resolver.observe_rkey(owner.clone(), rkey.clone()).await; 3228 assert!( 3229 resolver.cached_resolution(&owner, &rkey).await.is_some(), 3230 "observe must seed the cache", 3231 ); 3232 let delete: HydrantFrame = parse_frame(json!({ 3233 "id": 1, 3234 "type": "record", 3235 "record": { 3236 "live": false, 3237 "did": owner.as_ref(), 3238 "rev": fresh_tid().as_str(), 3239 "collection": "sh.tangled.repo", 3240 "rkey": rkey.as_ref(), 3241 "action": "delete" 3242 } 3243 })); 3244 handle_frame( 3245 delete, 3246 &store, 3247 &issue_states, 3248 &pull_statuses, 3249 &cov, 3250 &NoopSearchSink, 3251 &NoopRecordStore, 3252 &resolver, 3253 &sys_clock(), 3254 now(), 3255 ) 3256 .await; 3257 assert!( 3258 resolver.cached_resolution(&owner, &rkey).await.is_none(), 3259 "deleting the repo record must clear the resolver cache so future observes are not blocked by a stale Authoritative entry", 3260 ); 3261 assert_eq!( 3262 resolver.lookup_by_name(&owner, "abcabcabcabcz").await, 3263 None, 3264 "a deleted repo must stop answering on its url", 3265 ); 3266 } 3267 3268 #[tokio::test] 3269 async fn cancel_short_circuits_a_hung_ws_connect() { 3270 let _listener = tokio::net::TcpListener::bind("127.0.0.1:0") 3271 .await 3272 .expect("bind sink listener"); 3273 let port = _listener.local_addr().expect("local addr").port(); 3274 let cfg = 3275 IngestConfig::new(Url::parse(&format!("ws://127.0.0.1:{port}")).expect("hydrant url")); 3276 let cancel = CancellationToken::new(); 3277 let runtime = fresh_runtime(cancel.clone()); 3278 let task = tokio::spawn(async move { run(cfg, runtime).await }); 3279 tokio::time::sleep(Duration::from_millis(100)).await; 3280 cancel.cancel(); 3281 let outcome = tokio::time::timeout(Duration::from_secs(2), task) 3282 .await 3283 .expect("cancel must short-circuit the hung ws connect"); 3284 assert!(matches!(outcome, Ok(Ok(()))), "got {outcome:?}"); 3285 } 3286 3287 struct CloseOnConnectTransport { 3288 used: std::sync::Mutex<bool>, 3289 } 3290 impl bobbin_runtime::WsTransport for CloseOnConnectTransport { 3291 fn connect(&self, _url: Url) -> bobbin_runtime::WsConnectFuture { 3292 let mut used = self.used.lock().unwrap(); 3293 if *used { 3294 return Box::pin(async move { 3295 Err(NetworkError::Connect("only one connect allowed".to_owned())) 3296 }); 3297 } 3298 *used = true; 3299 Box::pin(async move { 3300 let mut q = std::collections::VecDeque::new(); 3301 q.push_back(Ok(WsMessage::Close { 3302 code: 1000, 3303 reason: "bye".to_owned(), 3304 })); 3305 let stream: Box<dyn WsStream> = Box::new(ScriptedWsStream { messages: q }); 3306 struct NoopSink; 3307 impl bobbin_runtime::WsSink for NoopSink { 3308 fn send<'a>(&'a mut self, _m: WsMessage) -> bobbin_runtime::WsSendFuture<'a> { 3309 Box::pin(async move { Ok(()) }) 3310 } 3311 } 3312 let sink: Box<dyn bobbin_runtime::WsSink> = Box::new(NoopSink); 3313 Ok(bobbin_runtime::WsConn { sink, stream }) 3314 }) 3315 } 3316 } 3317 3318 #[tokio::test] 3319 async fn run_session_returns_after_remote_close_when_outer_cancel_unfired() { 3320 let cfg = IngestConfig::new(Url::parse("ws://127.0.0.1:1").unwrap()); 3321 let cancel = CancellationToken::new(); 3322 let mut runtime = fresh_runtime(cancel.clone()); 3323 runtime.ws = Arc::new(CloseOnConnectTransport { 3324 used: std::sync::Mutex::new(false), 3325 }); 3326 3327 let res = tokio::time::timeout( 3328 Duration::from_secs(3), 3329 run_session(&cfg, HydrantCursor::new(0), &runtime), 3330 ) 3331 .await; 3332 assert!( 3333 res.is_ok(), 3334 "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", 3335 ); 3336 } 3337 3338 struct HangingSink; 3339 impl bobbin_runtime::WsSink for HangingSink { 3340 fn send<'a>(&'a mut self, _: WsMessage) -> bobbin_runtime::WsSendFuture<'a> { 3341 Box::pin(std::future::pending()) 3342 } 3343 } 3344 3345 #[tokio::test(start_paused = true)] 3346 async fn timed_send_surfaces_send_timeout_when_sink_pends_forever() { 3347 let mut sink: Box<dyn bobbin_runtime::WsSink> = Box::new(HangingSink); 3348 let task = 3349 tokio::spawn(async move { timed_send(&mut sink, WsMessage::Ping(Bytes::new())).await }); 3350 tokio::time::advance(SEND_TIMEOUT + Duration::from_secs(1)).await; 3351 let result = task.await.expect("task panicked"); 3352 match result { 3353 Err(IngestError::SendTimeout(d)) => assert_eq!(d, SEND_TIMEOUT), 3354 other => panic!( 3355 "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", 3356 ), 3357 } 3358 } 3359 3360 struct OkSink; 3361 impl bobbin_runtime::WsSink for OkSink { 3362 fn send<'a>(&'a mut self, _: WsMessage) -> bobbin_runtime::WsSendFuture<'a> { 3363 Box::pin(async move { Ok(()) }) 3364 } 3365 } 3366 3367 #[tokio::test] 3368 async fn timed_send_returns_ok_when_sink_succeeds_promptly() { 3369 let mut sink: Box<dyn bobbin_runtime::WsSink> = Box::new(OkSink); 3370 let result = timed_send(&mut sink, WsMessage::Ping(Bytes::new())).await; 3371 assert!(matches!(result, Ok(())), "got {result:?}"); 3372 } 3373 3374 #[test] 3375 fn classify_error_frame_with_id_dispatches_as_error_not_unknown_kind() { 3376 let text = r#"{"id":42,"type":"error","error":"ConsumerTooSlow","message":"slow"}"#; 3377 match classify_text_frame(text) { 3378 Err(IngestError::ConsumerTooSlow { message }) => assert_eq!(message, "slow"), 3379 other => panic!( 3380 "type=\"error\" must dispatch to the error path even when id is present, got: {other:?}" 3381 ), 3382 } 3383 } 3384 3385 #[tokio::test] 3386 async fn buffered_pipeline_preserves_cursor_order_under_resolve_latency_skew() { 3387 use bobbin_record_lru::RecordStore; 3388 use bobbin_types::record::RecordBody; 3389 use std::sync::Mutex; 3390 3391 let server = wiremock::MockServer::start().await; 3392 wiremock::Mock::given(wiremock::matchers::method("GET")) 3393 .and(wiremock::matchers::path("/xrpc/com.atproto.repo.getRecord")) 3394 .respond_with( 3395 wiremock::ResponseTemplate::new(404).set_delay(Duration::from_millis(150)), 3396 ) 3397 .mount(&server) 3398 .await; 3399 3400 let client = bobbin_slingshot_client::SlingshotClient::with_default_http( 3401 Url::parse(&server.uri()).unwrap(), 3402 ) 3403 .unwrap(); 3404 let clock: Arc<dyn Clock> = Arc::new(SystemClock::new()); 3405 let resolver = Arc::new(RepoIdResolver::with_slingshot( 3406 client, 3407 clock.clone(), 3408 RuntimeHasher::default(), 3409 )); 3410 3411 let owner = Did::new_owned("did:plc:nel").unwrap(); 3412 let abalone = Did::new_owned("did:plc:abalone").unwrap(); 3413 let fast_rkeys: [Rkey<DefaultStr>; 2] = [rkey("fastrkeyaa01"), rkey("fastrkeyaa02")]; 3414 let slow_rkeys: [Rkey<DefaultStr>; 2] = [rkey("slowrkeyaa01"), rkey("slowrkeyaa02")]; 3415 for r in &fast_rkeys { 3416 resolver 3417 .observe(owner.clone(), r.clone(), Some(abalone.clone()), None) 3418 .await; 3419 } 3420 3421 #[derive(Default)] 3422 struct Capturing { 3423 urls: Mutex<Vec<AtUri<DefaultStr>>>, 3424 } 3425 impl RecordStore for Capturing { 3426 fn get(&self, _uri: &AtUri<DefaultStr>) -> Option<Arc<RecordBody>> { 3427 None 3428 } 3429 fn put(&self, uri: AtUri<DefaultStr>, _body: Arc<RecordBody>) { 3430 self.urls.lock().unwrap().push(uri); 3431 } 3432 fn remove(&self, _uri: &AtUri<DefaultStr>) {} 3433 } 3434 let capturing = Arc::new(Capturing::default()); 3435 3436 let runtime: IngestRuntime<NoopSearchSink> = IngestRuntime { 3437 store: Arc::new(EdgeStore::new(RuntimeHasher::default())), 3438 issue_states: Arc::new(StateIndex::new(RuntimeHasher::default())), 3439 pull_statuses: Arc::new(StateIndex::new(RuntimeHasher::default())), 3440 coverage: Arc::new(CoverageWatch::new()), 3441 search: Arc::new(NoopSearchSink), 3442 records: capturing.clone() as Arc<dyn RecordStore>, 3443 resolver, 3444 clock, 3445 entropy: Arc::new(OsEntropy), 3446 ws: TungsteniteWs::shared(), 3447 cancel: CancellationToken::new(), 3448 disconnects: None, 3449 warming_shadow: None, 3450 warming_buffer: None, 3451 knot_registry: None, 3452 knot_gate: None, 3453 }; 3454 3455 let parallelism = 4usize; 3456 let (frame_tx, frame_rx) = tokio::sync::mpsc::channel::<HydrantFrame>(64); 3457 let pipeline_rt = runtime.clone(); 3458 let pipeline = tokio::spawn(async move { 3459 let prep_rt = pipeline_rt.clone(); 3460 let resolve_rt = pipeline_rt.clone(); 3461 let commit_rt = pipeline_rt; 3462 ReceiverStream::new(frame_rx) 3463 .then(move |frame| prep_stage(frame, prep_rt.clone())) 3464 .map(move |staged| resolve_stage(staged, resolve_rt.clone())) 3465 .buffered(parallelism) 3466 .for_each(move |staged| commit_stage(staged, commit_rt.clone(), parallelism)) 3467 .await; 3468 }); 3469 3470 let mk = |id: u64, idx: u64, repo_rkey: &Rkey<DefaultStr>| -> HydrantFrame { 3471 parse_frame(json!({ 3472 "id": id, 3473 "type": "record", 3474 "record": { 3475 "live": false, 3476 "did": "did:plc:olaren", 3477 "rev": fresh_tid().as_str(), 3478 "collection": "sh.tangled.feed.star", 3479 "rkey": format!("starrkeya{idx:04}"), 3480 "action": "create", 3481 "cid": VALID_CID, 3482 "record": { 3483 "$type": "sh.tangled.feed.star", 3484 "createdAt": "2026-05-01T00:00:00Z", 3485 "subject": format!("at://did:plc:nel/sh.tangled.repo/{}", repo_rkey.as_ref()), 3486 } 3487 } 3488 })) 3489 }; 3490 3491 frame_tx.send(mk(1, 1, &slow_rkeys[0])).await.unwrap(); 3492 frame_tx.send(mk(2, 2, &fast_rkeys[0])).await.unwrap(); 3493 frame_tx.send(mk(3, 3, &slow_rkeys[1])).await.unwrap(); 3494 frame_tx.send(mk(4, 4, &fast_rkeys[1])).await.unwrap(); 3495 drop(frame_tx); 3496 3497 pipeline.await.unwrap(); 3498 3499 let captured = capturing.urls.lock().unwrap(); 3500 let rkeys: Vec<String> = captured 3501 .iter() 3502 .filter_map(|uri| uri.rkey().map(|r| r.as_ref().to_owned())) 3503 .collect(); 3504 assert_eq!( 3505 rkeys, 3506 vec![ 3507 "starrkeya0001".to_owned(), 3508 "starrkeya0002".to_owned(), 3509 "starrkeya0003".to_owned(), 3510 "starrkeya0004".to_owned(), 3511 ], 3512 "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", 3513 ); 3514 assert_eq!(runtime.coverage.snapshot().last_cursor().raw(), 4); 3515 } 3516}