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