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