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