Monorepo for Tangled
Something went wrong. Try again.
12345678910111213141516171819202122232425262728293031323334353637383940414243444546474849505152535455565758596061626364656667686970717273747576777879808182838485868788899091929394959697989910010110210310410510610710810911011111211311411511611711811912012112212312412512612712812913013113213313413513613713813914014114214314414514614714814915015115215315415515615715815916016116216316416516616716816917017117217317417517617717817918018118218318418518618718818919019119219319419519619719819920020120220320420520620720820921021121221321421521621721821922022122222322422522622722822923023123223323423523623723823924024124224324424524624724824925025125225325425525625725825926026126226326426526626726826927027127227327427527627727827928028128228328428528628728828929029129229329429529629729829930030130230330430530630730830931031131231331431531631731831932032132232332432532632732832933033133233333433533633733833934034134234334434534634734834935035135235335435535635735835936036136236336436536636736836937037137237337437537637737837938038138238338438538638738838939039139239339439539639739839940040140240340440540640740840941041141241341441541641741841942042142242342442542642742842943043143243343443543643743843944044144244344444544644744844945045145245345445545645745845946046146246346446546646746846947047147247347447547647747847948048148248348448548648748848949049149249349449549649749849950050150250350450550650750850951051151251351451551651751851952052152252352452552652752852953053153253353453553653753853954054154254354454554654754854955055155255355455555655755855956056156256356456556656756856957057157257357457557657757857958058158258358458558658758858959059159259359459559659759859960060160260360460560660760860961061161261361461561661761861962062162262362462562662762862963063163263363463563663763863964064164264364464564664764864965065165265365465565665765865966066166266366466566666766866967067167267367467567667767867968068168268368468568668768868969069169269369469569669769869970070170270370470570670770870971071171271371471571671771871972072172272372472572672772872973073173273373473573673773873974074174274374474574674774874975075175275375475575675775875976076176276376476576676776876977077177277377477577677777877978078178278378478578678778878979079179279379479579679779879980080180280380480580680780880981081181281381481581681781881982082182282382482582682782882983083183283383483583683783883984084184284384484584684784884985085185285385485585685785885986086186286386486586686786886987087187287387487587687787887988088188288388488588688788888989089189289389489589689789889990090190290390490590690790890991091191291391491591691791891992092192292392492592692792892993093193293393493593693793893994094194294394494594694794894995095195295395495595695795895996096196296396496596696796896997097197297397497597697797897998098198298398498598698798898999099199299399499599699799899910001001100210031004100510061007100810091010101110121013101410151016101710181019102010211022102310241025102610271028102910301031103210331034103510361037103810391040104110421043104410451046104710481049105010511052105310541055105610571058105910601061106210631064106510661067106810691070107110721073107410751076107710781079108010811082108310841085108610871088108910901091109210931094109510961097109810991100110111021103110411051106110711081109111011111112111311141115111611171118111911201121112211231124112511261127112811291130113111321133113411351136113711381139114011411142114311441145114611471148114911501151115211531154115511561157115811591160116111621163116411651166116711681169117011711172117311741175117611771178117911801181118211831184118511861187118811891190119111921193119411951196119711981199120012011202120312041205120612071208120912101211121212131214121512161217121812191220122112221223122412251226122712281229123012311232123312341235123612371238123912401241124212431244124512461247124812491250125112521253125412551256125712581259126012611262126312641265126612671268126912701271127212731274127512761277127812791280128112821283128412851286128712881289129012911292129312941295129612971298129913001301130213031304130513061307130813091310131113121313131413151316131713181319132013211322132313241325132613271328132913301331133213331334133513361337133813391340134113421343134413451346134713481349135013511352135313541355135613571358135913601361136213631364136513661367136813691370137113721373137413751376137713781379138013811382138313841385138613871388138913901391139213931394139513961397139813991400140114021403140414051406140714081409141014111412141314141415141614171418141914201421142214231424142514261427142814291430143114321433143414351436143714381439144014411442144314441445144614471448144914501451145214531454145514561457145814591460146114621463146414651466146714681469147014711472147314741475147614771478147914801481148214831484148514861487148814891490149114921493149414951496149714981499150015011502150315041505150615071508150915101511151215131514151515161517151815191520152115221523152415251526152715281529153015311532153315341535153615371538153915401541154215431544154515461547154815491550155115521553155415551556155715581559156015611562156315641565156615671568156915701571157215731574157515761577157815791580158115821583158415851586158715881589159015911592159315941595159615971598159916001601160216031604160516061607160816091610161116121613161416151616161716181619162016211622162316241625162616271628162916301631163216331634163516361637163816391640164116421643164416451646164716481649165016511652165316541655165616571658165916601661166216631664166516661667166816691670167116721673167416751676167716781679168016811682168316841685168616871688168916901691169216931694169516961697169816991700170117021703170417051706170717081709171017111712171317141715171617171718171917201721172217231724172517261727172817291730173117321733173417351736173717381739174017411742174317441745174617471748174917501751175217531754175517561757175817591760176117621763176417651766176717681769177017711772177317741775177617771778177917801781178217831784178517861787178817891790179117921793179417951796179717981799180018011802180318041805180618071808180918101811181218131814181518161817181818191820182118221823182418251826182718281829183018311832183318341835183618371838183918401841184218431844184518461847184818491850185118521853185418551856185718581859186018611862186318641865186618671868186918701871187218731874187518761877187818791880188118821883188418851886188718881889189018911892189318941895189618971898189919001901190219031904190519061907190819091910191119121913191419151916191719181919192019211922192319241925192619271928192919301931193219331934193519361937193819391940194119421943194419451946194719481949195019511952195319541955195619571958195919601961196219631964196519661967196819691970197119721973197419751976197719781979198019811982198319841985198619871988198919901991199219931994199519961997199819992000200120022003200420052006200720082009201020112012201320142015201620172018201920202021202220232024202520262027202820292030203120322033203420352036203720382039204020412042204320442045204620472048204920502051205220532054205520562057205820592060206120622063206420652066206720682069207020712072207320742075207620772078207920802081208220832084208520862087208820892090209120922093209420952096209720982099210021012102210321042105210621072108210921102111211221132114211521162117211821192120212121222123212421252126212721282129213021312132213321342135213621372138213921402141214221432144214521462147214821492150215121522153215421552156215721582159216021612162216321642165216621672168216921702171217221732174217521762177217821792180218121822183218421852186218721882189219021912192219321942195219621972198219922002201220222032204220522062207220822092210221122122213221422152216221722182219222022212222222322242225222622272228222922302231223222332234223522362237223822392240224122422243224422452246224722482249225022512252225322542255225622572258225922602261226222632264226522662267226822692270227122722273227422752276227722782279228022812282228322842285228622872288228922902291229222932294229522962297229822992300230123022303230423052306230723082309231023112312231323142315231623172318231923202321232223232324232523262327232823292330233123322333233423352336233723382339234023412342234323442345234623472348234923502351235223532354235523562357235823592360236123622363236423652366236723682369237023712372237323742375237623772378237923802381238223832384238523862387238823892390239123922393239423952396239723982399240024012402240324042405240624072408240924102411241224132414241524162417241824192420242124222423242424252426242724282429243024312432243324342435243624372438243924402441244224432444244524462447244824492450245124522453245424552456245724582459246024612462246324642465246624672468246924702471247224732474247524762477247824792480248124822483248424852486248724882489249024912492249324942495249624972498249925002501250225032504250525062507250825092510251125122513251425152516251725182519252025212522252325242525252625272528252925302531253225332534253525362537253825392540254125422543254425452546254725482549255025512552255325542555255625572558255925602561256225632564256525662567256825692570257125722573257425752576257725782579258025812582258325842585258625872588258925902591259225932594259525962597259825992600260126022603260426052606260726082609261026112612261326142615261626172618261926202621262226232624262526262627262826292630263126322633263426352636263726382639264026412642264326442645264626472648264926502651265226532654265526562657265826592660266126622663266426652666266726682669267026712672267326742675267626772678267926802681268226832684268526862687268826892690269126922693269426952696269726982699270027012702270327042705270627072708270927102711271227132714271527162717271827192720272127222723272427252726272727282729273027312732273327342735273627372738273927402741274227432744274527462747274827492750275127522753275427552756275727582759276027612762276327642765276627672768276927702771277227732774277527762777277827792780278127822783278427852786278727882789279027912792279327942795279627972798279928002801280228032804280528062807280828092810281128122813281428152816281728182819282028212822282328242825282628272828282928302831283228332834283528362837283828392840284128422843284428452846284728482849285028512852285328542855285628572858285928602861286228632864286528662867286828692870287128722873287428752876287728782879288028812882288328842885288628872888288928902891289228932894289528962897289828992900290129022903290429052906290729082909291029112912291329142915291629172918291929202921292229232924292529262927292829292930293129322933293429352936293729382939294029412942294329442945294629472948294929502951295229532954295529562957295829592960296129622963296429652966296729682969297029712972297329742975297629772978297929802981298229832984298529862987298829892990299129922993299429952996299729982999300030013002300330043005300630073008300930103011301230133014301530163017301830193020302130223023302430253026302730283029303030313032303330343035303630373038303930403041304230433044304530463047304830493050305130523053305430553056305730583059306030613062306330643065306630673068306930703071307230733074307530763077307830793080308130823083308430853086308730883089309030913092309330943095309630973098309931003101310231033104310531063107310831093110311131123113311431153116311731183119312031213122312331243125312631273128312931303131313231333134313531363137313831393140314131423143314431453146314731483149315031513152315331543155315631573158315931603161316231633164316531663167316831693170317131723173317431753176317731783179318031813182318331843185318631873188318931903191319231933194319531963197319831993200320132023203320432053206320732083209321032113212321332143215321632173218321932203221322232233224322532263227322832293230323132323233323432353236323732383239324032413242324332443245324632473248324932503251325232533254325532563257325832593260326132623263326432653266326732683269327032713272327332743275327632773278327932803281328232833284use std::collections::{HashSet, VecDeque};use std::num::NonZeroUsize;use std::sync::Arc;use std::time::Duration;
use bobbin_edge_index::{ ApplyOutcome, Coverage, CoverageWatch, EdgeStore, HydrantCursor, IssueStateKind, PromotionSignal, PullStatusKind, StateIndex, apply_record_state,};use bobbin_knot_ingest::{CapabilityGate, KnotRegistry};use bobbin_record_lru::RecordStore;use bobbin_resolver::{NormalizeRepoRefs, decode_canon_or_upgrade_bytes, synthesize_created_at};use bobbin_runtime::{ Clock, Entropy, NetworkError, RuntimeHasher, UnixMicros, WsConn, WsMessage, WsStream, WsTransport,};use bobbin_types::edges::{Edge, ExtractError, Record};use bobbin_types::ids::{RepoIdent, SubjectRef};use bobbin_types::knot_acl::KnotHostKey;use bobbin_types::record::RecordBody;use bobbin_types::search::{SearchSink, SearchableRecord};use bytes::Bytes;use futures::StreamExt;use jacquard_common::DefaultStr;use jacquard_common::types::did::Did;use jacquard_common::types::ident::AtIdentifier;use jacquard_common::types::nsid::Nsid;use jacquard_common::types::recordkey::Rkey;use jacquard_common::types::string::{AtStrError, AtUri, Cid};use thiserror::Error;use tokio::time::Instant;use tokio_stream::wrappers::ReceiverStream;use tokio_util::sync::CancellationToken;use tracing::{debug, info, warn};use url::Url;
mod frame;mod resolver;mod shadow;mod warming;use frame::HydrantStreamErrorFrame;pub use frame::{FrameKind, HydrantFrame, RecordAction, RecordFrame};pub use resolver::{RepoIdResolver, Resolution};pub use shadow::{WarmingShadowBuffer, WarmingShadowSnapshot};pub use warming::{ParkedUpsert, WarmingBuffer, WarmingBufferSnapshot};
const TANGLED_PREFIX: &str = "sh.tangled.";const RECONNECT_INITIAL_DELAY: Duration = Duration::from_millis(500);const RECONNECT_MAX_DELAY: Duration = Duration::from_secs(30);const PING_INTERVAL: Duration = Duration::from_secs(20);const PONG_TIMEOUT: Duration = Duration::from_secs(15);const READY_SKEW: Duration = Duration::from_secs(60);const FRAME_CHANNEL_DEPTH: usize = 256;const CONTROL_CHANNEL_DEPTH: usize = 16;const READER_HOLD_LIMIT: usize = 64;const SEND_TIMEOUT: Duration = Duration::from_secs(10);const NORMAL_CLOSE: u16 = 1000;const METRICS_DUMP_INTERVAL: Duration = Duration::from_secs(10);const WARMING_FLUSH_PARALLELISM: usize = 64;pub const DEFAULT_INGEST_PARALLELISM: NonZeroUsize = match NonZeroUsize::new(16) { Some(n) => n, None => unreachable!(),};
#[derive(Clone, Debug)]pub struct IngestConfig { pub hydrant_base: Url, pub start_cursor: HydrantCursor, pub parallelism: NonZeroUsize,}
impl IngestConfig { pub fn new(hydrant_base: Url) -> Self { Self { hydrant_base, start_cursor: HydrantCursor::new(0), parallelism: DEFAULT_INGEST_PARALLELISM, } }
fn stream_url(&self, cursor: HydrantCursor) -> Result<Url, IngestError> { let mut url = self.hydrant_base.clone(); match url.scheme() { "http" => url .set_scheme("ws") .map_err(|_| IngestError::Url("set ws scheme"))?, "https" => url .set_scheme("wss") .map_err(|_| IngestError::Url("set wss scheme"))?, "ws" | "wss" => {} other => return Err(IngestError::UnknownScheme(other.to_owned())), } url.set_path("/stream"); url.query_pairs_mut() .clear() .append_pair("cursor", &cursor.raw().to_string()); Ok(url) }}
#[derive(Debug, Error)]pub enum IngestError { #[error("invalid hydrant url: {0}")] Url(&'static str), #[error("unsupported url scheme: {0}")] UnknownScheme(String), #[error("network: {0}")] Network(#[from] NetworkError), #[error("frame decode: {0}")] Decode(#[from] serde_json::Error), #[error("invalid at-uri synthesized from frame: {0}")] InvalidAtUri(#[from] AtStrError), #[error("record extraction: {0}")] Extract(#[from] ExtractError), #[error("hydrant did not respond to ping within {0:?}")] PongTimeout(Duration), #[error("websocket send blocked for at least {0:?}, treating link as dead")] SendTimeout(Duration), #[error("hydrant disconnected because bobbin's stream consumer fell behind: {message}")] ConsumerTooSlow { message: String }, #[error("hydrant signaled stream error {code}: {message}")] HydrantStream { code: String, message: String },}
#[derive(Clone, Copy, Debug, Eq, PartialEq, Hash)]pub enum DisconnectKind { Url, UnknownScheme, Network, Decode, InvalidAtUri, Extract, PongTimeout, SendTimeout, ConsumerTooSlow, HydrantStream,}
impl DisconnectKind { pub fn from_error(err: &IngestError) -> Self { match err { IngestError::Url(_) => Self::Url, IngestError::UnknownScheme(_) => Self::UnknownScheme, IngestError::Network(_) => Self::Network, IngestError::Decode(_) => Self::Decode, IngestError::InvalidAtUri(_) => Self::InvalidAtUri, IngestError::Extract(_) => Self::Extract, IngestError::PongTimeout(_) => Self::PongTimeout, IngestError::SendTimeout(_) => Self::SendTimeout, IngestError::ConsumerTooSlow { .. } => Self::ConsumerTooSlow, IngestError::HydrantStream { .. } => Self::HydrantStream, } }}
#[derive(Clone, Debug, Eq, PartialEq)]pub struct DisconnectSnapshot { pub kind: DisconnectKind, pub message: String, pub at_unix_micros: UnixMicros, pub last_cursor: HydrantCursor,}
#[derive(Default)]pub struct DisconnectSink { last: std::sync::Mutex<Option<DisconnectSnapshot>>, count: std::sync::atomic::AtomicU64,}
impl DisconnectSink { pub fn new() -> Self { Self::default() }
pub fn record(&self, snap: DisconnectSnapshot) { *self.last.lock().expect("disconnect sink mutex poisoned") = Some(snap); self.count .fetch_add(1, std::sync::atomic::Ordering::Relaxed); }
pub fn snapshot(&self) -> Option<DisconnectSnapshot> { self.last .lock() .expect("disconnect sink mutex poisoned") .clone() }
pub fn count(&self) -> u64 { self.count.load(std::sync::atomic::Ordering::Relaxed) }}
#[derive(Clone, Copy, Debug, Eq, PartialEq)]enum SessionOutcome { Progressed, Empty,}
#[derive(Debug)]struct SessionEnd { outcome: SessionOutcome, error: Option<IngestError>,}
pub struct IngestRuntime<S: SearchSink + 'static> { pub store: Arc<EdgeStore>, pub issue_states: Arc<StateIndex<IssueStateKind>>, pub pull_statuses: Arc<StateIndex<PullStatusKind>>, pub coverage: Arc<CoverageWatch>, pub search: Arc<S>, pub records: Arc<dyn RecordStore>, pub resolver: Arc<RepoIdResolver>, pub clock: Arc<dyn Clock>, pub entropy: Arc<dyn Entropy>, pub ws: Arc<dyn WsTransport>, pub cancel: CancellationToken, pub disconnects: Option<Arc<DisconnectSink>>, pub warming_shadow: Option<Arc<WarmingShadowBuffer>>, pub warming_buffer: Option<Arc<WarmingBuffer>>, pub knot_registry: Option<Arc<KnotRegistry>>, pub knot_gate: Option<Arc<CapabilityGate>>,}
impl<S: SearchSink + 'static> Clone for IngestRuntime<S> { fn clone(&self) -> Self { Self { store: self.store.clone(), issue_states: self.issue_states.clone(), pull_statuses: self.pull_statuses.clone(), coverage: self.coverage.clone(), search: self.search.clone(), records: self.records.clone(), resolver: self.resolver.clone(), clock: self.clock.clone(), entropy: self.entropy.clone(), ws: self.ws.clone(), cancel: self.cancel.clone(), disconnects: self.disconnects.clone(), warming_shadow: self.warming_shadow.clone(), warming_buffer: self.warming_buffer.clone(), knot_registry: self.knot_registry.clone(), knot_gate: self.knot_gate.clone(), } }}
impl<S: SearchSink + 'static> IngestRuntime<S> { fn pipeline_ctx(&self) -> PipelineCtx<'_, S> { PipelineCtx { resolver: &self.resolver, store: &self.store, issue_states: &self.issue_states, pull_statuses: &self.pull_statuses, coverage: &self.coverage, records: &*self.records, search: &self.search, shadow: self.warming_shadow.as_deref(), buffer: self.warming_buffer.as_deref(), knot_registry: self.knot_registry.as_deref(), knot_gate: self.knot_gate.as_deref(), } }}
struct PipelineCtx<'a, S: SearchSink + 'static> { resolver: &'a RepoIdResolver, store: &'a EdgeStore, issue_states: &'a StateIndex<IssueStateKind>, pull_statuses: &'a StateIndex<PullStatusKind>, coverage: &'a CoverageWatch, records: &'a dyn RecordStore, search: &'a S, shadow: Option<&'a WarmingShadowBuffer>, buffer: Option<&'a WarmingBuffer>, knot_registry: Option<&'a KnotRegistry>, knot_gate: Option<&'a CapabilityGate>,}
pub async fn run<S: SearchSink + 'static>( config: IngestConfig, runtime: IngestRuntime<S>,) -> Result<(), IngestError> { let metrics_dumper = spawn_metrics_dumper(&runtime); let idle_promoter = spawn_idle_promoter(&runtime); let warming_flusher = spawn_warming_flusher(&runtime); let result = run_inner(config, &runtime).await; if let Err(join) = metrics_dumper.await { warn!(?join, "metrics dumper task panicked"); } if let Err(join) = idle_promoter.await { warn!(?join, "idle promoter task panicked"); } if let Some(handle) = warming_flusher && let Err(join) = handle.await { warn!(?join, "warming flusher task panicked"); } result}
async fn run_inner<S: SearchSink + 'static>( config: IngestConfig, runtime: &IngestRuntime<S>,) -> Result<(), IngestError> { let mut backoff = RECONNECT_INITIAL_DELAY; loop { let cursor = next_connect_cursor(runtime.coverage.snapshot(), config.start_cursor); let SessionEnd { outcome, error } = run_session(&config, cursor, runtime).await; if runtime.cancel.is_cancelled() { info!( last_cursor = runtime.coverage.snapshot().last_cursor().raw(), "ingest stopped after shutdown signal" ); return Ok(()); } match (outcome, &error) { (SessionOutcome::Progressed, None) => { info!("hydrant stream closed after delivering frames, reconnecting") } (SessionOutcome::Empty, None) => { warn!("hydrant stream closed without delivering frames") } (_, Some(err)) => warn!(?err, "hydrant stream errored"), } if let (Some(sink), Some(err)) = (runtime.disconnects.as_ref(), error.as_ref()) { sink.record(DisconnectSnapshot { kind: DisconnectKind::from_error(err), message: err.to_string(), at_unix_micros: runtime.clock.now_unix_micros(), last_cursor: runtime.coverage.snapshot().last_cursor(), }); } let made_progress = matches!(outcome, SessionOutcome::Progressed); if made_progress { backoff = RECONNECT_INITIAL_DELAY; } else { tokio::select! { biased; _ = runtime.cancel.cancelled() => return Ok(()), _ = runtime.clock.sleep(jittered(backoff, &*runtime.entropy)) => {} } backoff = (backoff * 2).min(RECONNECT_MAX_DELAY); } }}
fn spawn_warming_flusher<S: SearchSink + 'static>( runtime: &IngestRuntime<S>,) -> Option<tokio::task::JoinHandle<()>> { runtime.warming_buffer.as_ref()?; let rt = runtime.clone(); Some(tokio::spawn(async move { let buffer = rt .warming_buffer .as_deref() .expect("warming flusher only spawns when buffer is set"); let mut rx = rt.coverage.subscribe(); let reason = loop { if rx.borrow_and_update().is_ready() { break FlushReason::Ready; } tokio::select! { biased; _ = rt.cancel.cancelled() => break FlushReason::Cancelled, res = rx.changed() => match res { Ok(()) => continue, Err(_) => break FlushReason::CoverageDropped, }, } }; flush_warming_buffer(&rt, buffer, reason).await; }))}
#[derive(Clone, Copy, Debug, Eq, PartialEq)]enum FlushReason { Ready, Cancelled, CoverageDropped,}
impl FlushReason { fn as_str(self) -> &'static str { match self { FlushReason::Ready => "ready", FlushReason::Cancelled => "cancelled", FlushReason::CoverageDropped => "coverage_dropped", } }}
async fn flush_warming_buffer<S: SearchSink + 'static>( runtime: &IngestRuntime<S>, buffer: &WarmingBuffer, reason: FlushReason,) { let drained = buffer.drain_for_promote().await; if drained.is_empty() { return; } if reason != FlushReason::Ready { info!( target: "bobbin_ingest::warming", abandoned_entries = drained.len(), reason = reason.as_str(), "abandoning parked items on non-ready flush", ); return; } let hasher = buffer.hasher().clone(); let unique: HashSet<RepoIdent, RuntimeHasher> = drained .iter() .flat_map(|(_, deps)| deps.iter().cloned()) .fold(HashSet::with_hasher(hasher), |mut acc, dep| { acc.insert(dep); acc }); if !unique.is_empty() { let resolver = runtime.resolver.clone(); let _: Vec<()> = futures::stream::iter(unique) .map(|key| { let resolver = resolver.clone(); async move { let _ = resolver.resolve(&key.owner, &key.rkey).await; } }) .buffer_unordered(WARMING_FLUSH_PARALLELISM) .collect() .await; } let upserts: Vec<ParkedUpsert> = drained.into_iter().map(|(u, _)| u).collect(); let count = upserts.len(); let ctx = runtime.pipeline_ctx(); finalize_drained(&ctx, upserts).await; info!( target: "bobbin_ingest::warming", flushed_entries = count, reason = reason.as_str(), "drained warming buffer", );}
const IDLE_PROMOTE_WINDOW: Duration = Duration::from_secs(15);const IDLE_PROMOTE_MIN_EVENTS: u64 = 256;
fn spawn_idle_promoter<S: SearchSink + 'static>( runtime: &IngestRuntime<S>,) -> tokio::task::JoinHandle<()> { let rt = runtime.clone(); tokio::spawn(async move { let mut prev = rt.coverage.snapshot().events_processed(); loop { tokio::select! { biased; _ = rt.cancel.cancelled() => return, _ = rt.clock.sleep(IDLE_PROMOTE_WINDOW) => {} } let snap = rt.coverage.snapshot(); if snap.is_ready() { return; } let processed = snap.events_processed(); if processed >= IDLE_PROMOTE_MIN_EVENTS && processed == prev { rt.coverage.update(|c| c.force_ready()); info!( target: "bobbin_ingest::coverage", events_processed = processed, last_cursor = snap.last_cursor().raw(), "stream idle, promoting coverage to ready", ); return; } prev = processed; } })}
fn spawn_metrics_dumper<S: SearchSink + 'static>( runtime: &IngestRuntime<S>,) -> tokio::task::JoinHandle<()> { let rt = runtime.clone(); tokio::spawn(async move { loop { tokio::select! { biased; _ = rt.cancel.cancelled() => break, _ = rt.clock.sleep(METRICS_DUMP_INTERVAL) => { let s = rt.resolver.stats(); info!( target: "bobbin_ingest::metrics", resolver_hits = s.hits, resolver_misses_mapped = s.misses_mapped, resolver_misses_no_repo_did = s.misses_no_repo_did, resolver_misses_unresolvable = s.misses_unresolvable, resolver_misses_transient = s.misses_transient, resolver_misses_no_client = s.misses_no_client, resolver_miss_latency_micros_avg = s.miss_latency_micros_avg().unwrap_or(0), resolver_miss_latency_micros_max = s.miss_latency_micros_max, resolver_total = s.total(), "resolver stats", ); } } } })}
fn next_connect_cursor(snapshot: Coverage, start: HydrantCursor) -> HydrantCursor { if snapshot.events_processed() == 0 { start } else { HydrantCursor::new(snapshot.last_cursor().raw().saturating_add(1)) }}
fn jittered(base: Duration, entropy: &dyn Entropy) -> Duration { let base_ms = u64::try_from(base.as_millis()).unwrap_or(u64::MAX); let cap_ms = (base_ms / 4).max(1); base + Duration::from_millis(entropy.next_u64() % cap_ms)}
async fn run_session<S: SearchSink + 'static>( config: &IngestConfig, cursor: HydrantCursor, runtime: &IngestRuntime<S>,) -> SessionEnd { let url = match config.stream_url(cursor) { Ok(u) => u, Err(e) => { return SessionEnd { outcome: SessionOutcome::Empty, error: Some(e), }; } }; info!(%url, "connecting to hydrant /stream"); let connect = tokio::select! { biased; _ = runtime.cancel.cancelled() => { return SessionEnd { outcome: SessionOutcome::Empty, error: None }; } res = runtime.ws.connect(url) => res, }; let WsConn { sink: mut ws_sink, stream: ws_stream, } = match connect { Ok(c) => c, Err(e) => { return SessionEnd { outcome: SessionOutcome::Empty, error: Some(IngestError::Network(e)), }; } };
let (frame_tx, frame_rx) = tokio::sync::mpsc::channel::<HydrantFrame>(FRAME_CHANNEL_DEPTH); let parallelism = config.parallelism.get(); let processor_runtime = runtime.clone(); let processor = tokio::spawn(async move { let cancel = processor_runtime.cancel.clone(); let prep_rt = processor_runtime.clone(); let resolve_rt = processor_runtime.clone(); let commit_rt = processor_runtime;
let pipeline = ReceiverStream::new(frame_rx) .map(move |frame| prep_stage(frame, prep_rt.clone())) .buffered(parallelism) .map(move |staged| resolve_stage(staged, resolve_rt.clone())) .buffered(parallelism) .for_each(move |staged| commit_stage(staged, commit_rt.clone(), parallelism));
tokio::select! { biased; _ = cancel.cancelled() => {}, _ = pipeline => {}, } });
let (control_tx, mut control_rx) = tokio::sync::mpsc::channel::<WsEvent>(CONTROL_CHANNEL_DEPTH); let session_cancel = runtime.cancel.child_token(); let reader_cancel = session_cancel.clone(); let reader = tokio::spawn(reader_loop(ws_stream, frame_tx, control_tx, reader_cancel));
let mut next_ping = runtime.clock.now_instant() + PING_INTERVAL; let mut pong_deadline: Option<Instant> = None;
let writer_error: Option<IngestError> = loop { tokio::select! { biased; _ = runtime.cancel.cancelled() => { let _ = timed_send( &mut ws_sink, WsMessage::Close { code: NORMAL_CLOSE, reason: "bobbin shutdown".to_owned() }, ).await; break None; } _ = runtime.clock.sleep_until(next_ping) => { next_ping = runtime.clock.now_instant() + PING_INTERVAL; if pong_deadline.is_none() { if let Err(e) = timed_send(&mut ws_sink, WsMessage::Ping(Bytes::new())).await { break Some(e); } pong_deadline = Some(runtime.clock.now_instant() + PONG_TIMEOUT); } } _ = wait_until(pong_deadline, runtime.clock.as_ref()) => { break Some(IngestError::PongTimeout(PONG_TIMEOUT)); } evt = control_rx.recv() => { let Some(evt) = evt else { break None; }; match evt { WsEvent::IncomingPing(payload) => { if let Err(e) = timed_send(&mut ws_sink, WsMessage::Pong(payload)).await { break Some(e); } } WsEvent::IncomingPong => { pong_deadline = None; } } } } };
session_cancel.cancel(); drop(ws_sink); drop(control_rx);
let reader_end = reader.await.unwrap_or(SessionEnd { outcome: SessionOutcome::Empty, error: None, }); if let Err(join) = processor.await { warn!(?join, "frame processor task panicked"); } SessionEnd { outcome: reader_end.outcome, error: writer_error.or(reader_end.error), }}
async fn timed_send( sink: &mut Box<dyn bobbin_runtime::WsSink>, msg: WsMessage,) -> Result<(), IngestError> { match tokio::time::timeout(SEND_TIMEOUT, sink.send(msg)).await { Ok(Ok(())) => Ok(()), Ok(Err(e)) => Err(IngestError::Network(e)), Err(_elapsed) => Err(IngestError::SendTimeout(SEND_TIMEOUT)), }}
async fn reader_loop( mut ws_stream: Box<dyn WsStream>, frame_tx: tokio::sync::mpsc::Sender<HydrantFrame>, control_tx: tokio::sync::mpsc::Sender<WsEvent>, cancel: CancellationToken,) -> SessionEnd { let mut outcome = SessionOutcome::Empty; let mut held: VecDeque<HydrantFrame> = VecDeque::new();
let error: Option<IngestError> = loop { let step = if held.is_empty() { tokio::select! { biased; _ = cancel.cancelled() => ReaderStep::Cancelled, msg = ws_stream.next() => ReaderStep::WsMessage(msg), } } else if held.len() < READER_HOLD_LIMIT { tokio::select! { biased; _ = cancel.cancelled() => ReaderStep::Cancelled, permit_result = frame_tx.reserve() => match permit_result { Ok(permit) => { let frame = held .pop_front() .expect("non-empty held when reserve succeeds"); permit.send(frame); ReaderStep::Sent } Err(_) => ReaderStep::FrameSinkClosed, }, msg = ws_stream.next() => ReaderStep::WsMessage(msg), } } else { tokio::select! { biased; _ = cancel.cancelled() => ReaderStep::Cancelled, permit_result = frame_tx.reserve() => match permit_result { Ok(permit) => { let frame = held .pop_front() .expect("held at limit when reserve succeeds"); permit.send(frame); ReaderStep::Sent } Err(_) => ReaderStep::FrameSinkClosed, }, } };
match step { ReaderStep::Cancelled => break None, ReaderStep::FrameSinkClosed => break None, ReaderStep::Sent => { outcome = SessionOutcome::Progressed; } ReaderStep::WsMessage(msg) => { let Some(msg) = msg else { break None; }; let parsed = match msg { Ok(m) => m, Err(e) => break Some(IngestError::Network(e)), }; match parsed { WsMessage::Text(text) => { let frame = match classify_text_frame(&text) { Ok(f) => f, Err(e) => break Some(e), }; held.push_back(frame); } WsMessage::Binary(_) => { debug!("hydrant sent unexpected binary frame, ignoring"); } WsMessage::Ping(payload) => { if control_tx .send(WsEvent::IncomingPing(payload)) .await .is_err() { break None; } } WsMessage::Pong(_) => { if control_tx.send(WsEvent::IncomingPong).await.is_err() { break None; } } WsMessage::Close { code, reason } => { debug!(code, %reason, "hydrant closed stream"); break None; } } } } }; SessionEnd { outcome, error }}
enum ReaderStep { Cancelled, FrameSinkClosed, Sent, WsMessage(Option<Result<WsMessage, NetworkError>>),}
fn classify_text_frame(text: &str) -> Result<HydrantFrame, IngestError> { #[derive(serde::Deserialize)] struct PeekType<'a> { #[serde(rename = "type", borrow)] kind: Option<std::borrow::Cow<'a, str>>, } let is_error_frame = serde_json::from_str::<PeekType>(text) .ok() .and_then(|p| p.kind) .as_deref() == Some("error"); if is_error_frame { return match serde_json::from_str::<HydrantStreamErrorFrame>(text) { Ok(err_frame) => Err(classify_hydrant_error(err_frame)), Err(decode_err) => Err(IngestError::Decode(decode_err)), }; } serde_json::from_str::<HydrantFrame>(text).map_err(IngestError::Decode)}
fn classify_hydrant_error(frame: HydrantStreamErrorFrame) -> IngestError { let HydrantStreamErrorFrame { error, message } = frame; let message = message.unwrap_or_default(); match error.as_str() { "ConsumerTooSlow" => IngestError::ConsumerTooSlow { message }, _ => IngestError::HydrantStream { code: error, message, }, }}
#[derive(Debug)]enum WsEvent { IncomingPing(Bytes), IncomingPong,}
async fn wait_until(deadline: Option<Instant>, clock: &dyn Clock) { match deadline { Some(d) => clock.sleep_until(d).await, None => std::future::pending::<()>().await, }}
#[derive(Clone, Copy, Debug, Eq, PartialEq)]enum Regime { Replay, Live, NonRecord,}
impl Regime { fn as_str(self) -> &'static str { match self { Self::Replay => "replay", Self::Live => "live", Self::NonRecord => "non_record", } }}
struct Pending { cursor: HydrantCursor, signal: PromotionSignal, regime: Regime, op: PendingOp,}
enum PendingOp { Noop, ClearCache { source: AtUri<DefaultStr>, }, Upsert { source: AtUri<DefaultStr>, nsid: Nsid<DefaultStr>, parsed: Box<Record>, bytes: Bytes, cid: Option<Cid<DefaultStr>>, edges: Vec<Edge>, }, Parked { nsid: Nsid<DefaultStr>, }, Delete { source: AtUri<DefaultStr>, nsid: Nsid<DefaultStr>, },}
struct Prepared { pending: Pending, prepare_start: Instant, prepare_end: Instant,}
struct Resolved { pending: Pending, prepare_start: Instant, prepare_end: Instant, resolve_start: Instant, resolve_end: Instant,}
fn pending_nsid(op: &PendingOp) -> Option<&Nsid<DefaultStr>> { match op { PendingOp::Upsert { nsid, .. } => Some(nsid), PendingOp::Delete { nsid, .. } => Some(nsid), PendingOp::Parked { nsid, .. } => Some(nsid), PendingOp::Noop | PendingOp::ClearCache { .. } => None, }}
fn pending_edge_count(op: &PendingOp) -> u64 { match op { PendingOp::Upsert { edges, .. } => edges.len() as u64, _ => 0, }}
async fn prep_stage<S: SearchSink + 'static>( frame: HydrantFrame, rt: IngestRuntime<S>,) -> Prepared { let now = rt.clock.now_unix_micros(); let prepare_start = rt.clock.now_instant(); let ctx = rt.pipeline_ctx(); let pending = prepare_frame(frame, &ctx, now).await; let prepare_end = rt.clock.now_instant(); Prepared { pending, prepare_start, prepare_end, }}
async fn resolve_stage<S: SearchSink + 'static>( staged: Prepared, rt: IngestRuntime<S>,) -> Resolved { let resolve_start = rt.clock.now_instant(); let ctx = rt.pipeline_ctx(); let pending = resolve_pending(staged.pending, &ctx).await; let resolve_end = rt.clock.now_instant(); Resolved { pending, prepare_start: staged.prepare_start, prepare_end: staged.prepare_end, resolve_start, resolve_end, }}
async fn commit_stage<S: SearchSink + 'static>( staged: Resolved, rt: IngestRuntime<S>, parallelism: usize,) { let nsid = pending_nsid(&staged.pending.op).cloned(); let edge_count = pending_edge_count(&staged.pending.op); let regime = staged.pending.regime; let cursor = staged.pending.cursor.raw(); let Resolved { pending, prepare_start, prepare_end, resolve_start, resolve_end, } = staged; let commit_start = rt.clock.now_instant(); commit_pending( pending, &rt.store, &rt.issue_states, &rt.pull_statuses, &rt.coverage, &*rt.search, &*rt.records, &rt.resolver, ) .await; let commit_end = rt.clock.now_instant(); tracing::trace!( target: "bobbin_ingest::stage", cursor, regime = regime.as_str(), nsid = %nsid.as_ref().map(Nsid::as_str).unwrap_or(""), edge_count, prepare_us = prepare_end.duration_since(prepare_start).as_micros() as u64, queue_resolve_wait_us = resolve_start.duration_since(prepare_end).as_micros() as u64, resolve_us = resolve_end.duration_since(resolve_start).as_micros() as u64, queue_commit_wait_us = commit_start.duration_since(resolve_end).as_micros() as u64, commit_us = commit_end.duration_since(commit_start).as_micros() as u64, total_us = commit_end.duration_since(prepare_start).as_micros() as u64, parallelism, "pipeline stage timings", );}
async fn prepare_frame<S: SearchSink + 'static>( frame: HydrantFrame, ctx: &PipelineCtx<'_, S>, now: UnixMicros,) -> Pending { let cursor = HydrantCursor::new(frame.id); let signal = promotion_signal(frame.record.as_ref(), now); let regime = match frame.record.as_ref() { Some(r) if r.live => Regime::Live, Some(_) => Regime::Replay, None => Regime::NonRecord, }; let op = match frame.kind { FrameKind::Record => prepare_record(frame.record, ctx).await, FrameKind::Identity | FrameKind::Account => PendingOp::Noop, FrameKind::Other => { debug!(id = frame.id, "ignoring unknown hydrant frame kind"); PendingOp::Noop } }; Pending { cursor, signal, regime, op, }}
async fn prepare_record<S: SearchSink + 'static>( record: Option<RecordFrame>, ctx: &PipelineCtx<'_, S>,) -> PendingOp { let Some(record) = record else { debug!("record-typed frame missing payload, skipping"); return PendingOp::Noop; }; if !record.collection.as_ref().starts_with(TANGLED_PREFIX) { return PendingOp::Noop; } let nsid = record.collection.clone(); let source = match build_source_uri(&record) { Ok(s) => s, Err(e) => { warn!(?e, "invalid frame source, dropping record"); return PendingOp::Noop; } }; match record.action { RecordAction::Create | RecordAction::Update => { evict_from_buffer(ctx.buffer, &source).await; let Some(raw) = record.record else { debug!(collection = %nsid, "create/update missing record body, clearing cache"); return PendingOp::ClearCache { source }; }; let raw_bytes = Bytes::copy_from_slice(raw.get().as_bytes()); let wire_bytes = match fallback_rfc3339(&record.rkey, &record.rev) .and_then(|fallback| synthesize_created_at(&raw_bytes, &fallback)) { Some(patched) => Bytes::from(patched), None => raw_bytes, }; let (parsed, bytes) = match decode_canon_or_upgrade_bytes( &record.collection, &wire_bytes, ctx.resolver, ) .await { Ok((parsed, canon_bytes)) => { let bytes = match canon_bytes { std::borrow::Cow::Borrowed(_) => wire_bytes, std::borrow::Cow::Owned(v) => Bytes::from(v), }; (parsed, bytes) } Err(ExtractError::UnknownCollection(name)) => { debug!(collection = %name, "unknown sh.tangled.* collection, clearing cache"); return PendingOp::ClearCache { source }; } Err(e) => { warn!(?e, collection = %record.collection, "record decode failed, clearing cache"); return PendingOp::ClearCache { source }; } }; if let Record::Repo(repo) = &parsed { if let Some(shadow) = ctx.shadow { shadow.note_observed(&record.did, &record.rkey).await; } let superseded = ctx .resolver .observe( record.did.clone(), record.rkey.clone(), repo.repo_did.clone(), ) .await; if let Some(prior) = superseded { let prior_uri = format!( "at://{}/sh.tangled.repo/{}", prior.owner.as_ref(), prior.rkey.as_ref(), ); if let Ok(prior_at_uri) = AtUri::<DefaultStr>::new_owned(&prior_uri) { evict_from_buffer(ctx.buffer, &prior_at_uri).await; ctx.store.remove_source(&prior_at_uri); ctx.records.remove(&prior_at_uri); ctx.search.remove(&prior_at_uri).await; } } if let Some(buffer) = ctx.buffer { let drained = buffer.take_observed(&record.did, &record.rkey).await; if !drained.is_empty() { finalize_drained(ctx, drained).await; } } if let Some(registry) = ctx.knot_registry { let host = KnotHostKey::new(repo.knot.as_ref()); match repo.repo_did.clone() { Some(repo_did) => registry.observe_repo(&host, repo_did), None => registry.observe_host(&host), } } } match acl_disposition(&parsed, ctx.knot_gate, ctx.knot_registry) { AclDisposition::NativeSkip => { if let Some(registry) = ctx.knot_registry { registry.forget_legacy_member(&source); } return PendingOp::Delete { source, nsid }; } AclDisposition::LegacyMember { host } => { if let Some(registry) = ctx.knot_registry { registry.observe_host(&host); registry.note_legacy_member(source.clone(), &host); } } AclDisposition::Other => {} } let edges = match parsed.extract_edges(&source) { Ok(es) => es, Err(e) => { warn!(?e, "edge extraction failed, clearing cache"); return PendingOp::ClearCache { source }; } }; let _ = ctx.store.intern_source(&source); PendingOp::Upsert { source, nsid, parsed: Box::new(parsed), bytes, cid: record.cid, edges, } } RecordAction::Delete => { evict_from_buffer(ctx.buffer, &source).await; if nsid.as_ref() == "sh.tangled.repo" { ctx.resolver.forget(&record.did, &record.rkey).await; } if nsid.as_ref() == "sh.tangled.knot.member" && let Some(registry) = ctx.knot_registry { registry.forget_legacy_member(&source); } PendingOp::Delete { source, nsid } } RecordAction::Other => { debug!(collection = %nsid, "ignoring unknown record action"); PendingOp::Noop } }}
enum AclDisposition { Other, NativeSkip, LegacyMember { host: KnotHostKey },}
fn acl_disposition( parsed: &Record, gate: Option<&CapabilityGate>, registry: Option<&KnotRegistry>,) -> AclDisposition { let Some(gate) = gate else { return AclDisposition::Other; }; match parsed { Record::KnotMember(member) => { let host = KnotHostKey::new(member.domain.as_ref()); if gate.is_native(&host) { AclDisposition::NativeSkip } else { AclDisposition::LegacyMember { host } } } Record::Collaborator(collaborator) => { let native = registry .and_then(|registry| registry.host_of_repo(&collaborator.repo)) .is_some_and(|host| gate.is_native(&host)); if native { AclDisposition::NativeSkip } else { AclDisposition::Other } } _ => AclDisposition::Other, }}
fn fallback_rfc3339( rkey: &Rkey<DefaultStr>, rev: &jacquard_common::types::tid::Tid,) -> Option<String> { let tid = jacquard_common::types::tid::Tid::new(rkey.as_ref()) .ok() .unwrap_or_else(|| rev.clone()); let micros = i64::try_from(tid.timestamp()).ok()?; let dt = chrono::DateTime::<chrono::Utc>::from_timestamp_micros(micros)?; Some(dt.to_rfc3339_opts(chrono::SecondsFormat::Micros, true))}
async fn evict_from_buffer(buffer: Option<&WarmingBuffer>, source: &AtUri<DefaultStr>) { if let Some(buffer) = buffer && !buffer.is_sealed() { buffer.evict_source(source).await; }}
async fn resolve_pending<S: SearchSink + 'static>( pending: Pending, ctx: &PipelineCtx<'_, S>,) -> Pending { let Pending { cursor, signal, regime, op, } = pending; let op = match op { PendingOp::Upsert { source, nsid, parsed, bytes, cid, edges, } => { let pieces = UpsertPieces { source, nsid, parsed: *parsed, bytes, cid, edges, }; let pieces = match try_park_warming(ctx, cursor, pieces).await { ParkOutcome::Parked { nsid } => { return Pending { cursor, signal, regime, op: PendingOp::Parked { nsid }, }; } ParkOutcome::Passthrough(pieces) => *pieces, }; let UpsertPieces { source, nsid, parsed, bytes, cid, edges, } = pieces; let edges = normalize_subjects(edges, ctx.resolver, ctx.coverage, ctx.shadow).await; PendingOp::Upsert { source, nsid, parsed: Box::new(parsed), bytes, cid, edges, } } other => other, }; Pending { cursor, signal, regime, op, }}
struct UpsertPieces { source: AtUri<DefaultStr>, nsid: Nsid<DefaultStr>, parsed: Record, bytes: Bytes, cid: Option<Cid<DefaultStr>>, edges: Vec<Edge>,}
impl From<ParkedUpsert> for UpsertPieces { fn from(u: ParkedUpsert) -> Self { Self { source: u.source, nsid: u.nsid, parsed: u.parsed, bytes: u.bytes, cid: u.cid, edges: u.edges, } }}
enum ParkOutcome { Parked { nsid: Nsid<DefaultStr> }, Passthrough(Box<UpsertPieces>),}
async fn try_park_warming<S: SearchSink + 'static>( ctx: &PipelineCtx<'_, S>, cursor: HydrantCursor, pieces: UpsertPieces,) -> ParkOutcome { let Some(buffer) = ctx.buffer else { return ParkOutcome::Passthrough(Box::new(pieces)); }; if ctx.coverage.snapshot().is_ready() || buffer.is_sealed() { return ParkOutcome::Passthrough(Box::new(pieces)); } let deps = collect_unresolved_deps(&pieces.edges, ctx.resolver).await; if deps.is_empty() { return ParkOutcome::Passthrough(Box::new(pieces)); } let nsid = pieces.nsid.clone(); let upsert = ParkedUpsert { cursor, source: pieces.source, nsid: pieces.nsid, parsed: pieces.parsed, bytes: pieces.bytes, cid: pieces.cid, edges: pieces.edges, }; let deps_for_shadow = ctx.shadow.is_some().then(|| deps.clone()); match buffer.try_park(upsert, deps).await { Ok(()) => { if let Some((shadow, noted)) = ctx.shadow.zip(deps_for_shadow) { let _ = futures::future::join_all(noted.into_iter().map(|dep| async move { shadow.note_unresolved(dep.owner, dep.rkey).await; })) .await; } ParkOutcome::Parked { nsid } } Err(returned) => ParkOutcome::Passthrough(Box::new(returned.into())), }}
async fn collect_unresolved_deps(edges: &[Edge], resolver: &RepoIdResolver) -> Vec<RepoIdent> { futures::stream::iter(edges) .fold(Vec::new(), |mut acc, edge| async move { let Some(uri) = edge.subject.as_uri() else { return acc; }; let Some((owner, rkey)) = parse_repo_subject_uri(uri) else { return acc; }; if resolver.cached_resolution(&owner, &rkey).await.is_some() { return acc; } let candidate = RepoIdent::new(owner, rkey); if !acc.contains(&candidate) { acc.push(candidate); } acc }) .await}
async fn finalize_drained<S: SearchSink + 'static>( ctx: &PipelineCtx<'_, S>, drained: Vec<ParkedUpsert>,) { for upsert in drained { let ParkedUpsert { cursor: _, source, nsid: _, parsed, bytes, cid, edges, } = upsert; let edges = normalize_subjects(edges, ctx.resolver, ctx.coverage, None).await; cache_body(ctx.records, &source, cid, bytes); ctx.store.upsert_source(&source, edges); let outcome = apply_record_state(ctx.issue_states, ctx.pull_statuses, &source, &parsed); log_unknown_state_variant(outcome, &source); index_search(ctx.search, ctx.resolver, &source, parsed).await; }}
#[allow(clippy::too_many_arguments)]async fn commit_pending<S: SearchSink>( pending: Pending, store: &EdgeStore, issue_states: &StateIndex<IssueStateKind>, pull_statuses: &StateIndex<PullStatusKind>, coverage: &CoverageWatch, search: &S, records: &dyn RecordStore, resolver: &RepoIdResolver,) { let Pending { cursor, signal, regime: _, op, } = pending; match op { PendingOp::Noop | PendingOp::Parked { .. } => {} PendingOp::ClearCache { source } => records.remove(&source), PendingOp::Upsert { source, nsid: _, parsed, bytes, cid, edges, } => { cache_body(records, &source, cid, bytes); store.upsert_source(&source, edges); let outcome = apply_record_state(issue_states, pull_statuses, &source, &parsed); log_unknown_state_variant(outcome, &source); index_search(search, resolver, &source, *parsed).await; } PendingOp::Delete { source, nsid } => { store.remove_source(&source); apply_delete_to_state_index(issue_states, pull_statuses, &source, &nsid); records.remove(&source); search.remove(&source).await; } } coverage.update(|c| c.advance(cursor).maybe_promote(signal));}
async fn index_search<S: SearchSink>( search: &S, resolver: &RepoIdResolver, source: &AtUri<DefaultStr>, parsed: Record,) { let Some(searchable) = SearchableRecord::try_from_record(parsed) else { return; }; let Some(searchable) = searchable.normalize(resolver).await else { return; }; search.upsert(searchable.to_search_doc(source)).await;}
#[cfg(test)]#[allow(clippy::too_many_arguments)]async fn handle_frame<S: SearchSink + 'static>( frame: HydrantFrame, store: &EdgeStore, issue_states: &StateIndex<IssueStateKind>, pull_statuses: &StateIndex<PullStatusKind>, coverage: &CoverageWatch, search: &S, records: &dyn RecordStore, resolver: &RepoIdResolver, clock: &dyn Clock, now: UnixMicros,) { let _ = clock; let ctx = PipelineCtx { resolver, store, issue_states, pull_statuses, coverage, records, search, shadow: None, buffer: None, knot_registry: None, knot_gate: None, }; let pending = prepare_frame(frame, &ctx, now).await; let pending = resolve_pending(pending, &ctx).await; commit_pending( pending, store, issue_states, pull_statuses, coverage, search, records, resolver, ) .await;}
fn log_unknown_state_variant(outcome: ApplyOutcome, source: &AtUri<DefaultStr>) { if matches!(outcome, ApplyOutcome::UnknownVariant) { warn!( target: "bobbin_ingest::state_index", %source, "state record has unknown wire variant, skipping index update", ); }}
fn apply_delete_to_state_index( issue_states: &StateIndex<IssueStateKind>, pull_statuses: &StateIndex<PullStatusKind>, source: &AtUri<DefaultStr>, nsid: &Nsid<DefaultStr>,) { match nsid.as_ref() { "sh.tangled.repo.issue" => issue_states.remove_entity(source), "sh.tangled.repo.pull" => pull_statuses.remove_entity(source), "sh.tangled.repo.issue.state" => issue_states.remove_source(source), "sh.tangled.repo.pull.status" => pull_statuses.remove_source(source), _ => {} }}
fn promotion_signal(record: Option<&RecordFrame>, now: UnixMicros) -> PromotionSignal { PromotionSignal { rev_micros: record.map(|r| r.rev.timestamp()), now_micros: now.raw(), skew_micros: READY_SKEW.as_micros() as u64, }}
fn cache_body( records: &dyn RecordStore, source: &AtUri<DefaultStr>, cid: Option<Cid<DefaultStr>>, bytes: Bytes,) { match cid { Some(cid) => records.put( source.clone(), Arc::new(RecordBody { uri: source.clone(), cid, value: bytes, }), ), None => records.remove(source), }}
async fn normalize_subjects( edges: Vec<Edge>, resolver: &RepoIdResolver, coverage: &CoverageWatch, shadow: Option<&WarmingShadowBuffer>,) -> Vec<Edge> { let warming = shadow.is_some() && !coverage.snapshot().is_ready(); futures::stream::iter(edges) .filter_map(|edge| async move { let Some(uri) = edge.subject.as_uri() else { return Some(edge); }; let Some((owner, rkey)) = parse_repo_subject_uri(uri) else { return Some(edge); }; if warming && let Some(shadow) = shadow && resolver.cached_resolution(&owner, &rkey).await.is_none() { shadow .note_unresolved(owner.clone(), rkey.clone()) .await; } match resolver.resolve(&owner, &rkey).await { Resolution::Mapped(repo_did) => Some(Edge { subject: SubjectRef::Did(repo_did), ..edge }), Resolution::NoRepoDid => { warn!( target: "bobbin_ingest::normalize", kind = %edge.kind, owner = owner.as_ref(), rkey = rkey.as_ref(), source = edge.source.as_ref(), "dropping edge: target repo has no repoDid, no canonical DID subject available", ); None } Resolution::Unresolvable => { warn!( target: "bobbin_ingest::normalize", kind = %edge.kind, owner = owner.as_ref(), rkey = rkey.as_ref(), source = edge.source.as_ref(), "dropping edge: repo unresolvable, rkey-form subject will not match bare-DID queries", ); None } } }) .collect() .await}
fn parse_repo_subject_uri(uri: &AtUri<DefaultStr>) -> Option<(Did<DefaultStr>, Rkey<DefaultStr>)> { let collection = uri.collection()?; if collection.as_ref() != "sh.tangled.repo" { return None; } let AtIdentifier::Did(authority) = uri.authority() else { return None; }; let rkey = uri.rkey()?; let owner = Did::new_owned(authority.as_ref()).ok()?; let rkey = Rkey::new_owned(rkey.as_ref()).ok()?; Some((owner, rkey))}
fn build_source_uri(r: &RecordFrame) -> Result<AtUri<DefaultStr>, IngestError> { Ok(AtUri::from_parts_owned( r.did.as_ref(), r.collection.as_ref(), r.rkey.as_ref(), )?)}
#[cfg(test)]mod tests { use super::*; use bobbin_edge_index::Coverage; use bobbin_record_lru::{CacheCapacity, LruRecordStore, NoopRecordStore, RecordStore}; use bobbin_runtime::{OsEntropy, RuntimeHasher, SystemClock, TungsteniteWs}; use bobbin_types::search::NoopSearchSink; use jacquard_common::types::nsid::Nsid; use jacquard_common::types::tid::Tid; use serde_json::json;
const VALID_CID: &str = "bafyreieqygohnz2zqyvtvktbjpvhutphobcmbsnt4q5lc36ri7vpcmoz4i";
fn did_subj(s: &str) -> SubjectRef { SubjectRef::Did(Did::new_owned(s).unwrap()) }
fn uri_subj(s: &str) -> SubjectRef { SubjectRef::Uri(AtUri::new_owned(s).unwrap()) }
fn rkey(s: &str) -> Rkey<DefaultStr> { Rkey::new_owned(s).unwrap() }
#[allow(clippy::type_complexity)] fn fresh() -> ( Arc<EdgeStore>, Arc<StateIndex<IssueStateKind>>, Arc<StateIndex<PullStatusKind>>, Arc<CoverageWatch>, Arc<RepoIdResolver>, ) { ( Arc::new(EdgeStore::new(RuntimeHasher::default())), Arc::new(StateIndex::new(RuntimeHasher::default())), Arc::new(StateIndex::new(RuntimeHasher::default())), Arc::new(CoverageWatch::new()), Arc::new(RepoIdResolver::detached(RuntimeHasher::default())), ) }
fn now() -> UnixMicros { SystemClock::new().now_unix_micros() }
fn sys_clock() -> SystemClock { SystemClock::new() }
fn parse_frame(value: serde_json::Value) -> HydrantFrame { let text = serde_json::to_string(&value).expect("serialize fixture"); serde_json::from_str(&text).expect("deserialize fixture") }
fn fresh_tid() -> Tid { Tid::now_0() }
#[tokio::test] async fn ignores_non_tangled_collections() { let (store, issue_states, pull_statuses, cov, resolver) = fresh(); let frame: HydrantFrame = parse_frame(json!({ "id": 1, "type": "record", "record": { "live": false, "did": "did:plc:nel", "rev": fresh_tid().as_str(), "collection": "app.bsky.feed.post", "rkey": "abcabcabcabcz", "action": "create", "record": {"$type": "app.bsky.feed.post", "text": "hi"} } })); handle_frame( frame, &store, &issue_states, &pull_statuses, &cov, &NoopSearchSink, &NoopRecordStore, &resolver, &sys_clock(), now(), ) .await; assert_eq!(store.key_count(), 0); assert_eq!(cov.snapshot().events_processed(), 1); }
#[tokio::test] async fn native_knot_member_skipped_legacy_indexed() { use bobbin_knot_ingest::{CapabilityGate, KnotClient, KnotRegistry}; use wiremock::matchers::{method, path}; use wiremock::{Mock, MockServer, ResponseTemplate};
let server = MockServer::start().await; Mock::given(method("GET")) .and(path("/xrpc/sh.tangled.knot.version")) .respond_with(ResponseTemplate::new(200).set_body_json(json!({ "version": "1.1.0", "capabilities": ["knot-acl"] }))) .mount(&server) .await; let url = url::Url::parse(&server.uri()).unwrap(); let native_host = format!("{}:{}", url.host_str().unwrap(), url.port().unwrap());
let gate = CapabilityGate::new( KnotClient::with_default_http(true).unwrap(), Arc::new(SystemClock::new()), true, true, ); assert!(gate.has_knot_acl(&KnotHostKey::new(&native_host)).await);
let registry = KnotRegistry::new(); let (store, issue_states, pull_statuses, cov, resolver) = fresh(); let ctx = PipelineCtx { resolver: &resolver, store: &store, issue_states: &issue_states, pull_statuses: &pull_statuses, coverage: &cov, records: &NoopRecordStore, search: &NoopSearchSink, shadow: None, buffer: None, knot_registry: Some(®istry), knot_gate: Some(&gate), };
let member_frame = |id: u64, rkey: &str, domain: &str| { parse_frame(json!({ "id": id, "type": "record", "record": { "live": false, "did": "did:plc:akshay", "rev": fresh_tid().as_str(), "collection": "sh.tangled.knot.member", "rkey": rkey, "action": "create", "record": { "$type": "sh.tangled.knot.member", "subject": "did:plc:boltless", "domain": domain, "createdAt": "2026-06-01T00:00:00Z" } } })) };
let native = prepare_frame(member_frame(1, "aaaaaaaaaaaaz", &native_host), &ctx, now()).await; assert!( matches!(native.op, PendingOp::Delete { .. }), "member record for a native knot must be dropped" );
let legacy = prepare_frame(member_frame(2, "bbbbbbbbbbbbz", "legacy.knot"), &ctx, now()).await; assert!( matches!(legacy.op, PendingOp::Upsert { .. }), "member record for a legacy knot must be ingested" ); assert!( registry.hosts().contains(&KnotHostKey::new("legacy.knot")), "a member record seeds host discovery even before any repo is seen" ); assert_eq!( registry .drain_legacy_members(&KnotHostKey::new("legacy.knot")) .len(), 1, "legacy member edge is indexed for later purge once the knot upgrades" ); }
#[tokio::test] async fn create_then_delete_round_trips_a_star() { let (store, issue_states, pull_statuses, cov, resolver) = fresh(); let create: HydrantFrame = parse_frame(json!({ "id": 10, "type": "record", "record": { "live": false, "did": "did:plc:olaren", "rev": fresh_tid().as_str(), "collection": "sh.tangled.feed.star", "rkey": "abcabcabcabcz", "action": "create", "record": { "$type": "sh.tangled.feed.star", "createdAt": "2026-05-01T00:00:00Z", "subjectDid": "did:plc:abalone" } } })); handle_frame( create, &store, &issue_states, &pull_statuses, &cov, &NoopSearchSink, &NoopRecordStore, &resolver, &sys_clock(), now(), ) .await; let key = bobbin_types::ids::EdgeKey::new( Nsid::new_static("sh.tangled.feed.star").unwrap(), did_subj("did:plc:abalone"), ); assert_eq!(store.count(&key), 1);
let delete: HydrantFrame = parse_frame(json!({ "id": 11, "type": "record", "record": { "live": false, "did": "did:plc:olaren", "rev": fresh_tid().as_str(), "collection": "sh.tangled.feed.star", "rkey": "abcabcabcabcz", "action": "delete", "record": null } })); handle_frame( delete, &store, &issue_states, &pull_statuses, &cov, &NoopSearchSink, &NoopRecordStore, &resolver, &sys_clock(), now(), ) .await; assert_eq!(store.count(&key), 0); }
#[tokio::test] async fn update_replaces_prior_edges() { let (store, issue_states, pull_statuses, cov, resolver) = fresh(); let mk = |subject_did: &Did<DefaultStr>, id: u64| -> HydrantFrame { parse_frame(json!({ "id": id, "type": "record", "record": { "live": false, "did": "did:plc:olaren", "rev": fresh_tid().as_str(), "collection": "sh.tangled.feed.star", "rkey": "abcabcabcabcz", "action": "update", "record": { "$type": "sh.tangled.feed.star", "createdAt": "2026-05-01T00:00:00Z", "subjectDid": subject_did.as_ref() } } })) }; handle_frame( mk(&Did::new_owned("did:plc:abalone").unwrap(), 1), &store, &issue_states, &pull_statuses, &cov, &NoopSearchSink, &NoopRecordStore, &resolver, &sys_clock(), now(), ) .await; handle_frame( mk(&Did::new_owned("did:plc:uni").unwrap(), 2), &store, &issue_states, &pull_statuses, &cov, &NoopSearchSink, &NoopRecordStore, &resolver, &sys_clock(), now(), ) .await;
let kind = Nsid::new_static("sh.tangled.feed.star").unwrap(); let old = bobbin_types::ids::EdgeKey::new(kind.clone(), did_subj("did:plc:abalone")); let new = bobbin_types::ids::EdgeKey::new(kind, did_subj("did:plc:uni")); assert_eq!(store.count(&old), 0); assert_eq!(store.count(&new), 1); }
#[tokio::test] async fn prepare_frame_tags_regime_from_live_flag() { let (store, issue_states, pull_statuses, cov, resolver) = fresh(); let search = NoopSearchSink; let records = NoopRecordStore; let ctx = PipelineCtx { resolver: &resolver, store: &store, issue_states: &issue_states, pull_statuses: &pull_statuses, coverage: &cov, records: &records, search: &search, shadow: None, buffer: None, knot_registry: None, knot_gate: None, }; let mk = |live: bool| -> HydrantFrame { parse_frame(json!({ "id": 1, "type": "record", "record": { "live": live, "did": "did:plc:olaren", "rev": fresh_tid().as_str(), "collection": "sh.tangled.feed.star", "rkey": "abcabcabcabcz", "action": "create", "record": { "$type": "sh.tangled.feed.star", "createdAt": "2026-05-01T00:00:00Z", "subjectDid": "did:plc:abalone" } } })) }; let live_pending = prepare_frame(mk(true), &ctx, now()).await; assert_eq!(live_pending.regime, Regime::Live); let replay_pending = prepare_frame(mk(false), &ctx, now()).await; assert_eq!(replay_pending.regime, Regime::Replay);
let identity: HydrantFrame = parse_frame(json!({ "id": 9, "type": "identity", })); let id_pending = prepare_frame(identity, &ctx, now()).await; assert_eq!(id_pending.regime, Regime::NonRecord); }
#[tokio::test] async fn create_with_cid_warms_record_lru() { let (store, issue_states, pull_statuses, cov, resolver) = fresh(); let lru = LruRecordStore::new(CacheCapacity::from_bytes(64 * 1024)); let frame: HydrantFrame = parse_frame(json!({ "id": 1, "type": "record", "record": { "live": false, "did": "did:plc:olaren", "rev": fresh_tid().as_str(), "collection": "sh.tangled.feed.star", "rkey": "abcabcabcabcz", "action": "create", "cid": VALID_CID, "record": { "$type": "sh.tangled.feed.star", "createdAt": "2026-05-01T00:00:00Z", "subjectDid": "did:plc:abalone" } } })); let source = AtUri::new_owned("at://did:plc:olaren/sh.tangled.feed.star/abcabcabcabcz").unwrap(); handle_frame( frame, &store, &issue_states, &pull_statuses, &cov, &NoopSearchSink, &lru, &resolver, &sys_clock(), now(), ) .await; let cached = lru.get(&source).expect("hydrant cid must seed the lru"); assert_eq!(cached.cid.as_ref(), VALID_CID); let parsed: serde_json::Value = serde_json::from_slice(&cached.value).unwrap(); assert_eq!( parsed["subject"]["did"], "did:plc:abalone", "legacy wire is upgraded to canon shape before caching so downstream readers see canonical fields" ); }
#[tokio::test] async fn create_without_cid_clears_record_lru() { let (store, issue_states, pull_statuses, cov, resolver) = fresh(); let lru = LruRecordStore::new(CacheCapacity::from_bytes(64 * 1024)); let source = AtUri::new_owned("at://did:plc:olaren/sh.tangled.feed.star/abcabcabcabcz").unwrap(); let cid: Cid<DefaultStr> = VALID_CID.parse().unwrap(); lru.put( source.clone(), Arc::new(RecordBody { uri: source.clone(), cid, value: bytes::Bytes::from_static(b"{\"stale\":true}"), }), ); let frame: HydrantFrame = parse_frame(json!({ "id": 1, "type": "record", "record": { "live": false, "did": "did:plc:olaren", "rev": fresh_tid().as_str(), "collection": "sh.tangled.feed.star", "rkey": "abcabcabcabcz", "action": "update", "record": { "$type": "sh.tangled.feed.star", "createdAt": "2026-05-01T00:00:00Z", "subjectDid": "did:plc:abalone" } } })); handle_frame( frame, &store, &issue_states, &pull_statuses, &cov, &NoopSearchSink, &lru, &resolver, &sys_clock(), now(), ) .await; assert!( lru.get(&source).is_none(), "missing cid means we cannot trust the body, so the lru must be cleared", ); }
#[tokio::test] async fn delete_evicts_record_lru() { let (store, issue_states, pull_statuses, cov, resolver) = fresh(); let lru = LruRecordStore::new(CacheCapacity::from_bytes(64 * 1024)); let source = AtUri::new_owned("at://did:plc:olaren/sh.tangled.feed.star/abcabcabcabcz").unwrap(); let cid: Cid<DefaultStr> = VALID_CID.parse().unwrap(); lru.put( source.clone(), Arc::new(RecordBody { uri: source.clone(), cid, value: bytes::Bytes::from_static(b"{}"), }), ); let frame: HydrantFrame = parse_frame(json!({ "id": 1, "type": "record", "record": { "live": false, "did": "did:plc:olaren", "rev": fresh_tid().as_str(), "collection": "sh.tangled.feed.star", "rkey": "abcabcabcabcz", "action": "delete", "record": null } })); handle_frame( frame, &store, &issue_states, &pull_statuses, &cov, &NoopSearchSink, &lru, &resolver, &sys_clock(), now(), ) .await; assert!(lru.get(&source).is_none()); }
#[tokio::test] async fn live_recent_event_promotes_coverage_to_ready() { let (store, issue_states, pull_statuses, cov, resolver) = fresh(); assert!(!cov.snapshot().is_ready()); let frame: HydrantFrame = parse_frame(json!({ "id": 99, "type": "record", "record": { "live": true, "did": "did:plc:olaren", "rev": fresh_tid().as_str(), "collection": "sh.tangled.feed.star", "rkey": "abcabcabcabcz", "action": "create", "record": { "$type": "sh.tangled.feed.star", "createdAt": "2026-05-01T00:00:00Z", "subjectDid": "did:plc:abalone" } } })); handle_frame( frame, &store, &issue_states, &pull_statuses, &cov, &NoopSearchSink, &NoopRecordStore, &resolver, &sys_clock(), now(), ) .await; assert!(cov.snapshot().is_ready()); assert_eq!(cov.snapshot().last_cursor(), HydrantCursor::new(99)); }
#[tokio::test] async fn live_but_stale_rev_does_not_promote() { let (store, issue_states, pull_statuses, cov, resolver) = fresh(); let stale_tid = Tid::from_time(1_000_000, 0); let frame: HydrantFrame = parse_frame(json!({ "id": 7, "type": "record", "record": { "live": true, "did": "did:plc:olaren", "rev": stale_tid.as_str(), "collection": "sh.tangled.feed.star", "rkey": "abcabcabcabcz", "action": "create", "record": { "$type": "sh.tangled.feed.star", "createdAt": "2026-05-01T00:00:00Z", "subjectDid": "did:plc:abalone" } } })); handle_frame( frame, &store, &issue_states, &pull_statuses, &cov, &NoopSearchSink, &NoopRecordStore, &resolver, &sys_clock(), now(), ) .await; assert!(!cov.snapshot().is_ready()); assert!(matches!(cov.snapshot(), Coverage::Warming { .. })); }
#[test] fn classify_consumer_too_slow_frame_returns_typed_variant() { let text = r#"{"type":"error","error":"ConsumerTooSlow","message":"stream socket send blocked for at least 30 seconds"}"#; match classify_text_frame(text) { Err(IngestError::ConsumerTooSlow { message }) => { assert!( message.contains("30 seconds"), "message field preserved verbatim, got: {message}" ); } other => panic!("expected ConsumerTooSlow variant, got: {other:?}"), } }
#[test] fn classify_unknown_hydrant_error_falls_back_to_generic_variant() { let text = r#"{"type":"error","error":"NewFutureCode","message":"some new failure mode"}"#; match classify_text_frame(text) { Err(IngestError::HydrantStream { code, message }) => { assert_eq!(code, "NewFutureCode"); assert_eq!(message, "some new failure mode"); } other => panic!("expected HydrantStream variant, got: {other:?}"), } }
#[test] fn classify_error_frame_without_message_uses_empty_string() { let text = r#"{"type":"error","error":"ConsumerTooSlow"}"#; match classify_text_frame(text) { Err(IngestError::ConsumerTooSlow { message }) => assert!(message.is_empty()), other => panic!("expected ConsumerTooSlow with empty message, got: {other:?}"), } }
#[test] fn classify_normal_record_frame_unchanged() { let text = r#"{"id":42,"type":"record"}"#; let frame = classify_text_frame(text).expect("normal record frame must parse"); assert_eq!(frame.id, 42); assert_eq!(frame.kind, FrameKind::Record); }
#[test] fn classify_garbage_object_returns_decode_error() { let text = r#"{"random":"object","without":"required fields"}"#; match classify_text_frame(text) { Err(IngestError::Decode(_)) => {} other => panic!("expected Decode error, got: {other:?}"), } }
struct ScriptedWsStream { messages: std::collections::VecDeque<Result<WsMessage, NetworkError>>, }
impl WsStream for ScriptedWsStream { fn next<'a>(&'a mut self) -> bobbin_runtime::WsMessageFuture<'a> { let msg = self.messages.pop_front(); Box::pin(async move { msg }) } }
#[tokio::test] async fn reader_loop_surfaces_consumer_too_slow_from_error_frame() { let mut messages = std::collections::VecDeque::new(); messages.push_back(Ok(WsMessage::Text( r#"{"type":"error","error":"ConsumerTooSlow","message":"stream delivery blocked"}"# .to_owned(), ))); let stream: Box<dyn WsStream> = Box::new(ScriptedWsStream { messages }); let (frame_tx, _frame_rx) = tokio::sync::mpsc::channel::<HydrantFrame>(FRAME_CHANNEL_DEPTH); let (control_tx, _control_rx) = tokio::sync::mpsc::channel::<WsEvent>(CONTROL_CHANNEL_DEPTH); let cancel = CancellationToken::new();
let end = reader_loop(stream, frame_tx, control_tx, cancel).await;
match end.error { Some(IngestError::ConsumerTooSlow { message }) => { assert_eq!(message, "stream delivery blocked"); } other => panic!("expected ConsumerTooSlow SessionEnd error, got: {other:?}"), } assert_eq!(end.outcome, SessionOutcome::Empty); }
#[tokio::test] async fn reader_loop_surfaces_unknown_hydrant_error_distinctly() { let mut messages = std::collections::VecDeque::new(); messages.push_back(Ok(WsMessage::Text( r#"{"type":"error","error":"NewFutureCode","message":"new mode"}"#.to_owned(), ))); let stream: Box<dyn WsStream> = Box::new(ScriptedWsStream { messages }); let (frame_tx, _frame_rx) = tokio::sync::mpsc::channel::<HydrantFrame>(FRAME_CHANNEL_DEPTH); let (control_tx, _control_rx) = tokio::sync::mpsc::channel::<WsEvent>(CONTROL_CHANNEL_DEPTH); let cancel = CancellationToken::new();
let end = reader_loop(stream, frame_tx, control_tx, cancel).await;
match end.error { Some(IngestError::HydrantStream { code, message }) => { assert_eq!(code, "NewFutureCode"); assert_eq!(message, "new mode"); } other => panic!("expected HydrantStream SessionEnd error, got: {other:?}"), } }
struct ChannelWsStream { rx: tokio::sync::mpsc::Receiver<Result<WsMessage, NetworkError>>, }
impl WsStream for ChannelWsStream { fn next<'a>(&'a mut self) -> bobbin_runtime::WsMessageFuture<'a> { Box::pin(async move { self.rx.recv().await }) } }
fn star_frame_text(id: u64, rkey: &Rkey<DefaultStr>) -> String { json!({ "id": id, "type": "record", "record": { "live": false, "did": "did:plc:olaren", "rev": fresh_tid().as_str(), "collection": "sh.tangled.feed.star", "rkey": rkey.as_ref(), "action": "create", "record": { "$type": "sh.tangled.feed.star", "createdAt": "2026-05-01T00:00:00Z", "subjectDid": "did:plc:abalone" } } }) .to_string() }
#[tokio::test] async fn pong_forwards_promptly_when_frame_channel_is_full() { let (ws_tx, ws_rx) = tokio::sync::mpsc::channel::<Result<WsMessage, NetworkError>>(8); let stream: Box<dyn WsStream> = Box::new(ChannelWsStream { rx: ws_rx }); let (frame_tx, frame_rx) = tokio::sync::mpsc::channel::<HydrantFrame>(1); let (control_tx, mut control_rx) = tokio::sync::mpsc::channel::<WsEvent>(CONTROL_CHANNEL_DEPTH); let cancel = CancellationToken::new();
let prefill: HydrantFrame = parse_frame(json!({ "id": 0, "type": "record", "record": { "live": false, "did": "did:plc:olaren", "rev": fresh_tid().as_str(), "collection": "sh.tangled.feed.star", "rkey": "prefilrkey001", "action": "create", "record": { "$type": "sh.tangled.feed.star", "createdAt": "2026-05-01T00:00:00Z", "subjectDid": "did:plc:abalone" } } })); frame_tx .try_send(prefill) .expect("depth-1 frame channel must accept the prefill");
let reader_handle = tokio::spawn(reader_loop( stream, frame_tx.clone(), control_tx.clone(), cancel.clone(), ));
ws_tx .send(Ok(WsMessage::Text(star_frame_text( 1, &rkey("starrkeyaa001"), )))) .await .unwrap(); ws_tx.send(Ok(WsMessage::Pong(Bytes::new()))).await.unwrap();
let pong_event = tokio::time::timeout(Duration::from_millis(500), control_rx.recv()) .await .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") .expect("control_tx was not closed"); assert!( matches!(pong_event, WsEvent::IncomingPong), "first control event must be the pong, not a held text", );
let too_slow = "{\"type\":\"error\",\"error\":\"ConsumerTooSlow\",\"message\":\"saturated\"}" .to_string(); ws_tx.send(Ok(WsMessage::Text(too_slow))).await.unwrap();
let end = tokio::time::timeout(Duration::from_secs(1), reader_handle) .await .expect("reader must exit promptly once ConsumerTooSlow is read") .expect("reader task must not panic"); match end.error { Some(IngestError::ConsumerTooSlow { message }) => assert_eq!(message, "saturated"), other => panic!("expected ConsumerTooSlow disconnect, got: {other:?}"), }
drop(frame_rx); }
#[tokio::test] async fn reader_drains_held_frames_once_processor_catches_up() { let (ws_tx, ws_rx) = tokio::sync::mpsc::channel::<Result<WsMessage, NetworkError>>(8); let stream: Box<dyn WsStream> = Box::new(ChannelWsStream { rx: ws_rx }); let (frame_tx, mut frame_rx) = tokio::sync::mpsc::channel::<HydrantFrame>(1); let (control_tx, _control_rx) = tokio::sync::mpsc::channel::<WsEvent>(CONTROL_CHANNEL_DEPTH); let cancel = CancellationToken::new();
let prefill: HydrantFrame = parse_frame(json!({ "id": 0, "type": "record", "record": { "live": false, "did": "did:plc:olaren", "rev": fresh_tid().as_str(), "collection": "sh.tangled.feed.star", "rkey": "prefilrkey001", "action": "create", "record": { "$type": "sh.tangled.feed.star", "createdAt": "2026-05-01T00:00:00Z", "subjectDid": "did:plc:abalone" } } })); frame_tx.try_send(prefill).unwrap();
let reader_handle = tokio::spawn(reader_loop( stream, frame_tx.clone(), control_tx, cancel.clone(), ));
ws_tx .send(Ok(WsMessage::Text(star_frame_text( 1, &rkey("heldrkeyaa001"), )))) .await .unwrap(); ws_tx .send(Ok(WsMessage::Text(star_frame_text( 2, &rkey("heldrkeyaa002"), )))) .await .unwrap();
let _drained_prefill = frame_rx.recv().await.expect("prefilled frame drains"); let first = tokio::time::timeout(Duration::from_millis(500), frame_rx.recv()) .await .expect("first held frame must reach frame_rx after slot opens") .expect("frame_tx still open"); assert_eq!(first.id, 1); let second = tokio::time::timeout(Duration::from_millis(500), frame_rx.recv()) .await .expect("second held frame must reach frame_rx after slot opens") .expect("frame_tx still open"); assert_eq!(second.id, 2);
cancel.cancel(); let _ = tokio::time::timeout(Duration::from_secs(1), reader_handle) .await .expect("reader must stop after cancel"); }
#[test] fn first_connect_uses_configured_start_cursor() { let start = HydrantCursor::new(42); assert_eq!(next_connect_cursor(Coverage::default(), start), start); }
#[test] fn reconnect_resumes_strictly_after_last_seen() { let snap = Coverage::default().advance(HydrantCursor::new(7)); assert_eq!( next_connect_cursor(snap, HydrantCursor::new(0)), HydrantCursor::new(8), ); }
#[test] fn reconnect_overrides_configured_start() { let snap = Coverage::default().advance(HydrantCursor::new(100)); assert_eq!( next_connect_cursor(snap, HydrantCursor::new(50)), HydrantCursor::new(101), ); }
#[test] fn first_connect_uses_start_even_when_first_frame_id_would_be_zero() { let start = HydrantCursor::new(7); let snap = Coverage::default(); assert_eq!(snap.last_cursor(), HydrantCursor::new(0)); assert_eq!(snap.events_processed(), 0); assert_eq!(next_connect_cursor(snap, start), start); }
#[test] fn reconnect_after_processing_id_zero_advances_to_one() { let snap = Coverage::default().advance(HydrantCursor::new(0)); assert_eq!(snap.events_processed(), 1); assert_eq!( next_connect_cursor(snap, HydrantCursor::new(99)), HydrantCursor::new(1), "events_processed disambiguates 'never seen' from 'saw id 0'", ); }
#[tokio::test] async fn identity_frame_advances_cursor_only() { let (store, issue_states, pull_statuses, cov, resolver) = fresh(); let frame: HydrantFrame = parse_frame(json!({ "id": 5, "type": "identity", "identity": { "did": "did:plc:olaren", "handle": "olaren.dev" } })); handle_frame( frame, &store, &issue_states, &pull_statuses, &cov, &NoopSearchSink, &NoopRecordStore, &resolver, &sys_clock(), now(), ) .await; assert_eq!(store.key_count(), 0); assert_eq!(cov.snapshot().last_cursor(), HydrantCursor::new(5)); assert!(!cov.snapshot().is_ready()); }
#[tokio::test] async fn account_frame_advances_cursor_only() { let (store, issue_states, pull_statuses, cov, resolver) = fresh(); let frame: HydrantFrame = parse_frame(json!({ "id": 6, "type": "account", "account": {"did": "did:plc:olaren", "active": true} })); handle_frame( frame, &store, &issue_states, &pull_statuses, &cov, &NoopSearchSink, &NoopRecordStore, &resolver, &sys_clock(), now(), ) .await; assert_eq!(store.key_count(), 0); assert_eq!(cov.snapshot().last_cursor(), HydrantCursor::new(6)); }
#[tokio::test] async fn unknown_frame_kind_advances_cursor_without_panic() { let (store, issue_states, pull_statuses, cov, resolver) = fresh(); let frame: HydrantFrame = parse_frame(json!({"id": 8, "type": "future_event"})); assert_eq!(frame.kind, FrameKind::Other); handle_frame( frame, &store, &issue_states, &pull_statuses, &cov, &NoopSearchSink, &NoopRecordStore, &resolver, &sys_clock(), now(), ) .await; assert_eq!(cov.snapshot().last_cursor(), HydrantCursor::new(8)); }
fn fresh_runtime(cancel: CancellationToken) -> IngestRuntime<NoopSearchSink> { IngestRuntime { store: Arc::new(EdgeStore::new(RuntimeHasher::default())), issue_states: Arc::new(StateIndex::new(RuntimeHasher::default())), pull_statuses: Arc::new(StateIndex::new(RuntimeHasher::default())), coverage: Arc::new(CoverageWatch::new()), search: Arc::new(NoopSearchSink), records: Arc::new(NoopRecordStore) as Arc<dyn RecordStore>, resolver: Arc::new(RepoIdResolver::detached(RuntimeHasher::default())), clock: Arc::new(SystemClock::new()), entropy: Arc::new(OsEntropy), ws: TungsteniteWs::shared(), cancel, disconnects: None, warming_shadow: None, warming_buffer: None, knot_registry: None, knot_gate: None, } }
#[tokio::test(start_paused = true)] async fn cancel_token_short_circuits_reconnect_sleep() { let cfg = IngestConfig::new(Url::parse("ws://127.0.0.1:1").unwrap()); let cancel = CancellationToken::new(); let runtime = fresh_runtime(cancel.clone()); let task = tokio::spawn(async move { run(cfg, runtime).await }); tokio::time::advance(Duration::from_millis(10)).await; cancel.cancel(); let outcome = tokio::time::timeout(Duration::from_secs(1), task) .await .expect("ingest must stop within timeout once cancel fires"); assert!(matches!(outcome, Ok(Ok(()))), "got {outcome:?}"); }
#[test] fn jittered_stays_within_one_quarter_of_base() { let base = Duration::from_secs(1); let cap = base + Duration::from_millis(250); let entropy = OsEntropy; (0..50).for_each(|_| { let j = jittered(base, &entropy); assert!(j >= base, "jitter must not undershoot"); assert!(j <= cap, "jitter must not exceed +25%, got {:?}", j); }); }
#[tokio::test] async fn star_after_observed_repo_keys_on_repo_did() { let (store, issue_states, pull_statuses, cov, resolver) = fresh(); let repo: HydrantFrame = parse_frame(json!({ "id": 1, "type": "record", "record": { "live": false, "did": "did:plc:nel", "rev": fresh_tid().as_str(), "collection": "sh.tangled.repo", "rkey": "abcabcabcabcz", "action": "create", "record": { "$type": "sh.tangled.repo", "createdAt": "2026-05-01T00:00:00Z", "knot": "oyster.cafe", "name": "abalone", "repoDid": "did:plc:abalone" } } })); handle_frame( repo, &store, &issue_states, &pull_statuses, &cov, &NoopSearchSink, &NoopRecordStore, &resolver, &sys_clock(), now(), ) .await;
let star: HydrantFrame = parse_frame(json!({ "id": 2, "type": "record", "record": { "live": false, "did": "did:plc:olaren", "rev": fresh_tid().as_str(), "collection": "sh.tangled.feed.star", "rkey": "abcabcabcabcz", "action": "create", "record": { "$type": "sh.tangled.feed.star", "createdAt": "2026-05-01T00:00:00Z", "subject": "at://did:plc:nel/sh.tangled.repo/abcabcabcabcz" } } })); handle_frame( star, &store, &issue_states, &pull_statuses, &cov, &NoopSearchSink, &NoopRecordStore, &resolver, &sys_clock(), now(), ) .await;
let nsid = Nsid::new_static("sh.tangled.feed.star").unwrap(); let repo_keyed = bobbin_types::ids::EdgeKey::new(nsid.clone(), did_subj("did:plc:abalone")); let owner_keyed = bobbin_types::ids::EdgeKey::new(nsid, did_subj("did:plc:nel")); assert_eq!( store.count(&repo_keyed), 1, "star should be keyed on repoDID once the repo is observed" ); assert_eq!( store.count(&owner_keyed), 0, "owner DID should not collect the edge" ); }
#[tokio::test] async fn unresolvable_repo_subject_drops_edge() { let (store, issue_states, pull_statuses, cov, resolver) = fresh(); let star: HydrantFrame = parse_frame(json!({ "id": 1, "type": "record", "record": { "live": false, "did": "did:plc:olaren", "rev": fresh_tid().as_str(), "collection": "sh.tangled.feed.star", "rkey": "abcabcabcabcz", "action": "create", "record": { "$type": "sh.tangled.feed.star", "createdAt": "2026-05-01T00:00:00Z", "subject": "at://did:plc:nel/sh.tangled.repo/abcabcabcabcz" } } })); handle_frame( star, &store, &issue_states, &pull_statuses, &cov, &NoopSearchSink, &NoopRecordStore, &resolver, &sys_clock(), now(), ) .await;
let nsid = Nsid::new_static("sh.tangled.feed.star").unwrap(); let owner_keyed = bobbin_types::ids::EdgeKey::new(nsid.clone(), did_subj("did:plc:nel")); let uri_keyed = bobbin_types::ids::EdgeKey::new( nsid, uri_subj("at://did:plc:nel/sh.tangled.repo/abcabcabcabcz"), ); assert_eq!( store.count(&owner_keyed), 0, "must not silently misfile under the authoring DID", ); assert_eq!( store.count(&uri_keyed), 0, "unresolvable rkey-form subject must drop the edge; keeping it would index against a key that never matches bare-DID queries", ); }
#[tokio::test] async fn repo_without_repo_did_drops_edge() { let (store, issue_states, pull_statuses, cov, resolver) = fresh(); let repo: HydrantFrame = parse_frame(json!({ "id": 1, "type": "record", "record": { "live": false, "did": "did:plc:nel", "rev": fresh_tid().as_str(), "collection": "sh.tangled.repo", "rkey": "abcabcabcabcz", "action": "create", "record": { "$type": "sh.tangled.repo", "createdAt": "2026-05-01T00:00:00Z", "knot": "oyster.cafe", "name": "abalone" } } })); handle_frame( repo, &store, &issue_states, &pull_statuses, &cov, &NoopSearchSink, &NoopRecordStore, &resolver, &sys_clock(), now(), ) .await;
let star: HydrantFrame = parse_frame(json!({ "id": 2, "type": "record", "record": { "live": false, "did": "did:plc:olaren", "rev": fresh_tid().as_str(), "collection": "sh.tangled.feed.star", "rkey": "abcabcabcabcz", "action": "create", "record": { "$type": "sh.tangled.feed.star", "createdAt": "2026-05-01T00:00:00Z", "subject": "at://did:plc:nel/sh.tangled.repo/abcabcabcabcz" } } })); handle_frame( star, &store, &issue_states, &pull_statuses, &cov, &NoopSearchSink, &NoopRecordStore, &resolver, &sys_clock(), now(), ) .await;
let nsid = Nsid::new_static("sh.tangled.feed.star").unwrap(); let uri_keyed = bobbin_types::ids::EdgeKey::new( nsid.clone(), uri_subj("at://did:plc:nel/sh.tangled.repo/abcabcabcabcz"), ); let owner_keyed = bobbin_types::ids::EdgeKey::new(nsid, did_subj("did:plc:nel")); assert_eq!( store.count(&uri_keyed), 0, "no canonical DID exists for a repo without repoDID, so the edge must be dropped", ); assert_eq!( store.count(&owner_keyed), 0, "the authoring DID is not the canonical repo identity", ); }
#[tokio::test] async fn explicit_subject_did_skips_normalization() { let (store, issue_states, pull_statuses, cov, resolver) = fresh(); let star: HydrantFrame = parse_frame(json!({ "id": 1, "type": "record", "record": { "live": false, "did": "did:plc:olaren", "rev": fresh_tid().as_str(), "collection": "sh.tangled.feed.star", "rkey": "abcabcabcabcz", "action": "create", "record": { "$type": "sh.tangled.feed.star", "createdAt": "2026-05-01T00:00:00Z", "subjectDid": "did:plc:abalone" } } })); handle_frame( star, &store, &issue_states, &pull_statuses, &cov, &NoopSearchSink, &NoopRecordStore, &resolver, &sys_clock(), now(), ) .await; let key = bobbin_types::ids::EdgeKey::new( Nsid::new_static("sh.tangled.feed.star").unwrap(), did_subj("did:plc:abalone"), ); assert_eq!(store.count(&key), 1); }
#[tokio::test] async fn issue_with_repo_uri_resolves_to_repo_did() { let (store, issue_states, pull_statuses, cov, resolver) = fresh(); resolver .observe( Did::new_owned("did:plc:nel").unwrap(), Rkey::new_owned("abcabcabcabcz").unwrap(), Some(Did::new_owned("did:plc:abalone").unwrap()), ) .await; let issue: HydrantFrame = parse_frame(json!({ "id": 1, "type": "record", "record": { "live": false, "did": "did:plc:olaren", "rev": fresh_tid().as_str(), "collection": "sh.tangled.repo.issue", "rkey": "abcabcabcabcz", "action": "create", "record": { "$type": "sh.tangled.repo.issue", "createdAt": "2026-05-01T00:00:00Z", "title": "bug", "repo": "at://did:plc:nel/sh.tangled.repo/abcabcabcabcz" } } })); handle_frame( issue, &store, &issue_states, &pull_statuses, &cov, &NoopSearchSink, &NoopRecordStore, &resolver, &sys_clock(), now(), ) .await; let key = bobbin_types::ids::EdgeKey::new( Nsid::new_static("sh.tangled.repo.issue").unwrap(), did_subj("did:plc:abalone"), ); assert_eq!(store.count(&key), 1); }
#[derive(Default)] struct RecordingSearchSink { docs: tokio::sync::Mutex<Vec<bobbin_types::search::SearchDoc>>, }
impl SearchSink for RecordingSearchSink { async fn upsert(&self, doc: bobbin_types::search::SearchDoc) { self.docs.lock().await.push(doc); } async fn remove(&self, _uri: &AtUri<DefaultStr>) {} }
#[tokio::test] async fn search_index_hydrates_repo_did_via_resolver() { let (store, issue_states, pull_statuses, cov, resolver) = fresh(); resolver .observe( Did::new_owned("did:plc:nel").unwrap(), Rkey::new_owned("abcabcabcabcz").unwrap(), Some(Did::new_owned("did:plc:abalone").unwrap()), ) .await; let search = RecordingSearchSink::default(); let issue: HydrantFrame = parse_frame(json!({ "id": 1, "type": "record", "record": { "live": false, "did": "did:plc:olaren", "rev": fresh_tid().as_str(), "collection": "sh.tangled.repo.issue", "rkey": "abcabcabcabcz", "action": "create", "record": { "$type": "sh.tangled.repo.issue", "createdAt": "2026-05-01T00:00:00Z", "title": "bug", "repo": "at://did:plc:nel/sh.tangled.repo/abcabcabcabcz" } } })); handle_frame( issue, &store, &issue_states, &pull_statuses, &cov, &search, &NoopRecordStore, &resolver, &sys_clock(), now(), ) .await; let docs = search.docs.lock().await; assert_eq!(docs.len(), 1, "issue should produce one search doc"); assert_eq!( docs[0].repo, Some(Did::new_owned("did:plc:abalone").unwrap()), "search doc repo field must be resolved from the observed repo, not left as None", ); }
#[tokio::test] async fn delete_repo_record_evicts_resolver_cache() { let (store, issue_states, pull_statuses, cov, resolver) = fresh(); let owner = Did::new_owned("did:plc:nel").unwrap(); let rkey = Rkey::new_owned("abcabcabcabcz").unwrap(); resolver .observe( owner.clone(), rkey.clone(), Some(Did::new_owned("did:plc:abalone").unwrap()), ) .await; assert!( resolver.cached_resolution(&owner, &rkey).await.is_some(), "observe must seed the cache", ); let delete: HydrantFrame = parse_frame(json!({ "id": 1, "type": "record", "record": { "live": false, "did": owner.as_ref(), "rev": fresh_tid().as_str(), "collection": "sh.tangled.repo", "rkey": rkey.as_ref(), "action": "delete" } })); handle_frame( delete, &store, &issue_states, &pull_statuses, &cov, &NoopSearchSink, &NoopRecordStore, &resolver, &sys_clock(), now(), ) .await; assert!( resolver.cached_resolution(&owner, &rkey).await.is_none(), "deleting the repo record must clear the resolver cache so future observes are not blocked by a stale Authoritative entry", ); }
#[tokio::test] async fn cancel_short_circuits_a_hung_ws_connect() { let _listener = tokio::net::TcpListener::bind("127.0.0.1:0") .await .expect("bind sink listener"); let port = _listener.local_addr().expect("local addr").port(); let cfg = IngestConfig::new(Url::parse(&format!("ws://127.0.0.1:{port}")).expect("hydrant url")); let cancel = CancellationToken::new(); let runtime = fresh_runtime(cancel.clone()); let task = tokio::spawn(async move { run(cfg, runtime).await }); tokio::time::sleep(Duration::from_millis(100)).await; cancel.cancel(); let outcome = tokio::time::timeout(Duration::from_secs(2), task) .await .expect("cancel must short-circuit the hung ws connect"); assert!(matches!(outcome, Ok(Ok(()))), "got {outcome:?}"); }
struct CloseOnConnectTransport { used: std::sync::Mutex<bool>, } impl bobbin_runtime::WsTransport for CloseOnConnectTransport { fn connect(&self, _url: Url) -> bobbin_runtime::WsConnectFuture { let mut used = self.used.lock().unwrap(); if *used { return Box::pin(async move { Err(NetworkError::Connect("only one connect allowed".to_owned())) }); } *used = true; Box::pin(async move { let mut q = std::collections::VecDeque::new(); q.push_back(Ok(WsMessage::Close { code: 1000, reason: "bye".to_owned(), })); let stream: Box<dyn WsStream> = Box::new(ScriptedWsStream { messages: q }); struct NoopSink; impl bobbin_runtime::WsSink for NoopSink { fn send<'a>(&'a mut self, _m: WsMessage) -> bobbin_runtime::WsSendFuture<'a> { Box::pin(async move { Ok(()) }) } } let sink: Box<dyn bobbin_runtime::WsSink> = Box::new(NoopSink); Ok(bobbin_runtime::WsConn { sink, stream }) }) } }
#[tokio::test] async fn run_session_returns_after_remote_close_when_outer_cancel_unfired() { let cfg = IngestConfig::new(Url::parse("ws://127.0.0.1:1").unwrap()); let cancel = CancellationToken::new(); let mut runtime = fresh_runtime(cancel.clone()); runtime.ws = Arc::new(CloseOnConnectTransport { used: std::sync::Mutex::new(false), });
let res = tokio::time::timeout( Duration::from_secs(3), run_session(&cfg, HydrantCursor::new(0), &runtime), ) .await; assert!( res.is_ok(), "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", ); }
struct HangingSink; impl bobbin_runtime::WsSink for HangingSink { fn send<'a>(&'a mut self, _: WsMessage) -> bobbin_runtime::WsSendFuture<'a> { Box::pin(std::future::pending()) } }
#[tokio::test(start_paused = true)] async fn timed_send_surfaces_send_timeout_when_sink_pends_forever() { let mut sink: Box<dyn bobbin_runtime::WsSink> = Box::new(HangingSink); let task = tokio::spawn(async move { timed_send(&mut sink, WsMessage::Ping(Bytes::new())).await }); tokio::time::advance(SEND_TIMEOUT + Duration::from_secs(1)).await; let result = task.await.expect("task panicked"); match result { Err(IngestError::SendTimeout(d)) => assert_eq!(d, SEND_TIMEOUT), other => panic!( "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", ), } }
struct OkSink; impl bobbin_runtime::WsSink for OkSink { fn send<'a>(&'a mut self, _: WsMessage) -> bobbin_runtime::WsSendFuture<'a> { Box::pin(async move { Ok(()) }) } }
#[tokio::test] async fn timed_send_returns_ok_when_sink_succeeds_promptly() { let mut sink: Box<dyn bobbin_runtime::WsSink> = Box::new(OkSink); let result = timed_send(&mut sink, WsMessage::Ping(Bytes::new())).await; assert!(matches!(result, Ok(())), "got {result:?}"); }
#[test] fn classify_error_frame_with_id_dispatches_as_error_not_unknown_kind() { let text = r#"{"id":42,"type":"error","error":"ConsumerTooSlow","message":"slow"}"#; match classify_text_frame(text) { Err(IngestError::ConsumerTooSlow { message }) => assert_eq!(message, "slow"), other => panic!( "type=\"error\" must dispatch to the error path even when id is present, got: {other:?}" ), } }
#[tokio::test] async fn buffered_pipeline_preserves_cursor_order_under_resolve_latency_skew() { use bobbin_record_lru::RecordStore; use bobbin_types::record::RecordBody; use std::sync::Mutex;
let server = wiremock::MockServer::start().await; wiremock::Mock::given(wiremock::matchers::method("GET")) .and(wiremock::matchers::path("/xrpc/com.atproto.repo.getRecord")) .respond_with( wiremock::ResponseTemplate::new(404).set_delay(Duration::from_millis(150)), ) .mount(&server) .await;
let client = bobbin_slingshot_client::SlingshotClient::with_default_http( Url::parse(&server.uri()).unwrap(), ) .unwrap(); let clock: Arc<dyn Clock> = Arc::new(SystemClock::new()); let resolver = Arc::new(RepoIdResolver::with_slingshot( client, clock.clone(), RuntimeHasher::default(), ));
let owner = Did::new_owned("did:plc:nel").unwrap(); let abalone = Did::new_owned("did:plc:abalone").unwrap(); let fast_rkeys: [Rkey<DefaultStr>; 2] = [rkey("fastrkeyaa01"), rkey("fastrkeyaa02")]; let slow_rkeys: [Rkey<DefaultStr>; 2] = [rkey("slowrkeyaa01"), rkey("slowrkeyaa02")]; for r in &fast_rkeys { resolver .observe(owner.clone(), r.clone(), Some(abalone.clone())) .await; }
#[derive(Default)] struct Capturing { urls: Mutex<Vec<AtUri<DefaultStr>>>, } impl RecordStore for Capturing { fn get(&self, _uri: &AtUri<DefaultStr>) -> Option<Arc<RecordBody>> { None } fn put(&self, uri: AtUri<DefaultStr>, _body: Arc<RecordBody>) { self.urls.lock().unwrap().push(uri); } fn remove(&self, _uri: &AtUri<DefaultStr>) {} } let capturing = Arc::new(Capturing::default());
let runtime: IngestRuntime<NoopSearchSink> = IngestRuntime { store: Arc::new(EdgeStore::new(RuntimeHasher::default())), issue_states: Arc::new(StateIndex::new(RuntimeHasher::default())), pull_statuses: Arc::new(StateIndex::new(RuntimeHasher::default())), coverage: Arc::new(CoverageWatch::new()), search: Arc::new(NoopSearchSink), records: capturing.clone() as Arc<dyn RecordStore>, resolver, clock, entropy: Arc::new(OsEntropy), ws: TungsteniteWs::shared(), cancel: CancellationToken::new(), disconnects: None, warming_shadow: None, warming_buffer: None, knot_registry: None, knot_gate: None, };
let parallelism = 4usize; let (frame_tx, frame_rx) = tokio::sync::mpsc::channel::<HydrantFrame>(64); let pipeline_rt = runtime.clone(); let pipeline = tokio::spawn(async move { let prep_rt = pipeline_rt.clone(); let resolve_rt = pipeline_rt.clone(); let commit_rt = pipeline_rt; ReceiverStream::new(frame_rx) .then(move |frame| prep_stage(frame, prep_rt.clone())) .map(move |staged| resolve_stage(staged, resolve_rt.clone())) .buffered(parallelism) .for_each(move |staged| commit_stage(staged, commit_rt.clone(), parallelism)) .await; });
let mk = |id: u64, idx: u64, repo_rkey: &Rkey<DefaultStr>| -> HydrantFrame { parse_frame(json!({ "id": id, "type": "record", "record": { "live": false, "did": "did:plc:olaren", "rev": fresh_tid().as_str(), "collection": "sh.tangled.feed.star", "rkey": format!("starrkeya{idx:04}"), "action": "create", "cid": VALID_CID, "record": { "$type": "sh.tangled.feed.star", "createdAt": "2026-05-01T00:00:00Z", "subject": format!("at://did:plc:nel/sh.tangled.repo/{}", repo_rkey.as_ref()), } } })) };
frame_tx.send(mk(1, 1, &slow_rkeys[0])).await.unwrap(); frame_tx.send(mk(2, 2, &fast_rkeys[0])).await.unwrap(); frame_tx.send(mk(3, 3, &slow_rkeys[1])).await.unwrap(); frame_tx.send(mk(4, 4, &fast_rkeys[1])).await.unwrap(); drop(frame_tx);
pipeline.await.unwrap();
let captured = capturing.urls.lock().unwrap(); let rkeys: Vec<String> = captured .iter() .filter_map(|uri| uri.rkey().map(|r| r.as_ref().to_owned())) .collect(); assert_eq!( rkeys, vec![ "starrkeya0001".to_owned(), "starrkeya0002".to_owned(), "starrkeya0003".to_owned(), "starrkeya0004".to_owned(), ], "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", ); assert_eq!(runtime.coverage.snapshot().last_cursor().raw(), 4); }}