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