diff --git a/crates/ingest/Cargo.toml b/crates/ingest/Cargo.toml index 7eead91..26563f9 100644 --- a/crates/ingest/Cargo.toml +++ b/crates/ingest/Cargo.toml @@ -8,6 +8,7 @@ rust-version.workspace = true [dependencies] bobbin-types = { workspace = true } bobbin-edge-index = { workspace = true } +bobbin-record-lru = { workspace = true } jacquard-common = { workspace = true } futures = { workspace = true } diff --git a/crates/ingest/examples/smoke.rs b/crates/ingest/examples/smoke.rs index 18689aa..e7584a3 100644 --- a/crates/ingest/examples/smoke.rs +++ b/crates/ingest/examples/smoke.rs @@ -3,6 +3,8 @@ use std::time::Duration; use bobbin_edge_index::{CoverageWatch, EdgeStore}; use bobbin_ingest::{IngestConfig, NoopResolver, run}; +use bobbin_record_lru::{NoopRecordStore, RecordStore}; +use bobbin_types::search::NoopSearchSink; use futures::stream::{self, StreamExt}; use url::Url; @@ -28,8 +30,10 @@ async fn main() { let store_runner = store.clone(); let coverage_runner = coverage.clone(); let resolver = Arc::new(NoopResolver); + let search = Arc::new(NoopSearchSink); + let records: Arc = Arc::new(NoopRecordStore); let task = tokio::spawn(async move { - let _ = run(cfg, store_runner, coverage_runner, resolver).await; + let _ = run(cfg, store_runner, coverage_runner, resolver, search, records).await; }); stream::iter(0..ticks) diff --git a/crates/ingest/src/lib.rs b/crates/ingest/src/lib.rs index 751bb47..a6204b0 100644 --- a/crates/ingest/src/lib.rs +++ b/crates/ingest/src/lib.rs @@ -2,7 +2,9 @@ use std::sync::Arc; use std::time::{Duration, SystemTime, UNIX_EPOCH}; use bobbin_edge_index::{Coverage, CoverageWatch, EdgeStore, HydrantCursor, PromotionSignal}; +use bobbin_record_lru::RecordStore; use bobbin_types::edges::{Edge, ExtractError, Record}; +use bobbin_types::search::{SearchSink, SearchableRecord}; use futures::SinkExt; use futures::stream::{self, StreamExt, TryStreamExt}; use jacquard_common::DefaultStr; @@ -120,11 +122,13 @@ struct SessionEnd { error: Option, } -pub async fn run( +pub async fn run( config: IngestConfig, store: Arc, coverage: Arc, resolver: Arc, + search: Arc, + records: Arc, ) -> Result<(), IngestError> { let mut backoff = RECONNECT_INITIAL_DELAY; loop { @@ -135,6 +139,8 @@ pub async fn run( store.clone(), coverage.clone(), resolver.clone(), + search.clone(), + records.clone(), ) .await; match (outcome, &error) { @@ -173,12 +179,14 @@ fn jittered(base: Duration) -> Duration { base + Duration::from_millis(entropy % cap_ms) } -async fn run_session( +async fn run_session( config: &IngestConfig, cursor: HydrantCursor, store: Arc, coverage: Arc, resolver: Arc, + search: Arc, + records: Arc, ) -> SessionEnd { let url = match config.stream_url(cursor) { Ok(u) => u, @@ -209,6 +217,8 @@ async fn run_session( let processor_store = store.clone(); let processor_coverage = coverage.clone(); let processor_resolver = resolver.clone(); + let processor_search = search.clone(); + let processor_records = records.clone(); let processor = tokio::spawn(async move { while let Some(frame) = frame_rx.recv().await { handle_frame( @@ -216,6 +226,8 @@ async fn run_session( &processor_store, &processor_coverage, &*processor_resolver, + &*processor_search, + &*processor_records, ) .await; } @@ -289,11 +301,13 @@ async fn wait_until(deadline: Option) { } } -async fn handle_frame( +async fn handle_frame( frame: HydrantFrame, store: &EdgeStore, coverage: &CoverageWatch, resolver: &R, + search: &S, + records: &dyn RecordStore, ) { let cursor = HydrantCursor::new(frame.id); let signal = promotion_signal(frame.record.as_ref(), now_micros()); @@ -311,7 +325,7 @@ async fn handle_frame( if !record.collection.as_ref().starts_with(TANGLED_PREFIX) { return; } - if let Err(err) = apply_record(record, store, resolver).await { + if let Err(err) = apply_record(record, store, resolver, search, records).await { warn!(?err, "dropped record"); } } @@ -338,14 +352,17 @@ fn now_micros() -> u64 { .as_micros() as u64 } -async fn apply_record( +async fn apply_record( record: RecordFrame, store: &EdgeStore, resolver: &R, + search: &S, + records: &dyn RecordStore, ) -> Result<(), IngestError> { let source = build_source_uri(&record)?; match record.action { RecordAction::Create | RecordAction::Update => { + records.remove(&source); let Some(value) = record.record else { debug!(collection = %record.collection.as_ref(), "create/update missing record body, skipping"); return Ok(()); @@ -365,9 +382,14 @@ async fn apply_record( .try_collect() .await?; store.upsert_source(&source, normalized); + if let Some(searchable) = SearchableRecord::try_from_record(parsed) { + search.upsert(searchable.to_search_doc(&source)).await; + } } RecordAction::Delete => { store.remove_source(&source); + search.remove(&source).await; + records.remove(&source); } RecordAction::Other => { debug!(collection = %record.collection.as_ref(), "ignoring unknown record action"); @@ -419,6 +441,8 @@ fn build_source_uri(r: &RecordFrame) -> Result, IngestError> { mod tests { use super::*; use bobbin_edge_index::Coverage; + use bobbin_record_lru::NoopRecordStore; + use bobbin_types::search::NoopSearchSink; use jacquard_common::types::nsid::Nsid; use jacquard_common::types::tid::Tid; use serde_json::json; @@ -452,7 +476,7 @@ mod tests { "record": {"$type": "app.bsky.feed.post", "text": "hi"} } })); - handle_frame(frame, &store, &cov, &NoopResolver).await; + handle_frame(frame, &store, &cov, &NoopResolver, &NoopSearchSink, &NoopRecordStore).await; assert_eq!(store.key_count(), 0); assert_eq!(cov.snapshot().events_processed(), 1); } @@ -477,7 +501,7 @@ mod tests { } } })); - handle_frame(create, &store, &cov, &NoopResolver).await; + handle_frame(create, &store, &cov, &NoopResolver, &NoopSearchSink, &NoopRecordStore).await; let key = bobbin_types::ids::EdgeKey::new( Nsid::new_static("sh.tangled.feed.star").unwrap(), AtUri::new_owned("at://did:plc:abalone").unwrap(), @@ -497,7 +521,7 @@ mod tests { "record": null } })); - handle_frame(delete, &store, &cov, &NoopResolver).await; + handle_frame(delete, &store, &cov, &NoopResolver, &NoopSearchSink, &NoopRecordStore).await; assert_eq!(store.count(&key), 0); } @@ -523,8 +547,24 @@ mod tests { } })) }; - handle_frame(mk("did:plc:abalone", 1), &store, &cov, &NoopResolver).await; - handle_frame(mk("did:plc:uni", 2), &store, &cov, &NoopResolver).await; + handle_frame( + mk("did:plc:abalone", 1), + &store, + &cov, + &NoopResolver, + &NoopSearchSink, + &NoopRecordStore, + ) + .await; + handle_frame( + mk("did:plc:uni", 2), + &store, + &cov, + &NoopResolver, + &NoopSearchSink, + &NoopRecordStore, + ) + .await; let kind = Nsid::new_static("sh.tangled.feed.star").unwrap(); let old = bobbin_types::ids::EdgeKey::new( @@ -558,7 +598,7 @@ mod tests { } } })); - handle_frame(frame, &store, &cov, &NoopResolver).await; + handle_frame(frame, &store, &cov, &NoopResolver, &NoopSearchSink, &NoopRecordStore).await; assert!(cov.snapshot().is_ready()); assert_eq!(cov.snapshot().last_cursor(), HydrantCursor::new(99)); } @@ -584,7 +624,7 @@ mod tests { } } })); - handle_frame(frame, &store, &cov, &NoopResolver).await; + handle_frame(frame, &store, &cov, &NoopResolver, &NoopSearchSink, &NoopRecordStore).await; assert!(!cov.snapshot().is_ready()); assert!(matches!(cov.snapshot(), Coverage::Warming { .. })); } @@ -644,7 +684,7 @@ mod tests { "handle": "olaren.dev" } })); - handle_frame(frame, &store, &cov, &NoopResolver).await; + handle_frame(frame, &store, &cov, &NoopResolver, &NoopSearchSink, &NoopRecordStore).await; assert_eq!(store.key_count(), 0); assert_eq!(cov.snapshot().last_cursor(), HydrantCursor::new(5)); assert!(!cov.snapshot().is_ready()); @@ -658,7 +698,7 @@ mod tests { "type": "account", "account": {"did": "did:plc:olaren", "active": true} })); - handle_frame(frame, &store, &cov, &NoopResolver).await; + handle_frame(frame, &store, &cov, &NoopResolver, &NoopSearchSink, &NoopRecordStore).await; assert_eq!(store.key_count(), 0); assert_eq!(cov.snapshot().last_cursor(), HydrantCursor::new(6)); } @@ -668,7 +708,7 @@ mod tests { let (store, cov) = fresh(); let frame: HydrantFrame = parse_frame(json!({"id": 8, "type": "future_event"})); assert_eq!(frame.kind, FrameKind::Other); - handle_frame(frame, &store, &cov, &NoopResolver).await; + handle_frame(frame, &store, &cov, &NoopResolver, &NoopSearchSink, &NoopRecordStore).await; assert_eq!(cov.snapshot().last_cursor(), HydrantCursor::new(8)); }