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