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