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