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