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