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