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