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