diff --git a/crates/bobbin/Cargo.toml b/crates/bobbin/Cargo.toml index 55b76ed..652000f 100644 --- a/crates/bobbin/Cargo.toml +++ b/crates/bobbin/Cargo.toml @@ -14,6 +14,7 @@ bobbin-edge-index = { workspace = true } bobbin-ingest = { workspace = true } bobbin-knot-proxy = { workspace = true } bobbin-record-lru = { workspace = true } +bobbin-search = { workspace = true } bobbin-slingshot-client = { workspace = true } bobbin-types = { workspace = true } bobbin-xrpc = { workspace = true } diff --git a/crates/bobbin/src/main.rs b/crates/bobbin/src/main.rs index 026ea57..e50c17b 100644 --- a/crates/bobbin/src/main.rs +++ b/crates/bobbin/src/main.rs @@ -7,6 +7,7 @@ use bobbin_edge_index::{CoverageWatch, EdgeStore, HydrantCursor}; use bobbin_ingest::{IngestConfig, RepoDidResolver, ResolveError, run as run_ingest}; use bobbin_knot_proxy::{KnotProxy, KnotProxyConfig}; use bobbin_record_lru::{CacheCapacity, LruRecordStore, RecordStore}; +use bobbin_search::{DEFAULT_WRITER_HEAP_BYTES, SearchIndex}; use bobbin_slingshot_client::{SlingshotClient, SlingshotError}; use bobbin_types::record::RecordBody; use bobbin_xrpc::{AppState, router}; @@ -39,6 +40,10 @@ async fn main() -> anyhow::Result<()> { .ok() .and_then(|v| v.parse().ok()) .unwrap_or(64 * 1024 * 1024); + let search_heap_bytes: usize = env::var("BOBBIN_SEARCH_HEAP_BYTES") + .ok() + .and_then(|v| v.parse().ok()) + .unwrap_or(DEFAULT_WRITER_HEAP_BYTES); let start_cursor: u64 = env::var("BOBBIN_START_CURSOR") .ok() .and_then(|v| v.parse().ok()) @@ -50,6 +55,7 @@ async fn main() -> anyhow::Result<()> { let edges = Arc::new(EdgeStore::new()); let coverage = Arc::new(CoverageWatch::new()); let knots = Arc::new(KnotProxy::new(KnotProxyConfig::default())?); + let search = Arc::new(SearchIndex::new(search_heap_bytes)?); let ingest_cfg = IngestConfig { hydrant_base: Url::parse(&hydrant_url)?, @@ -62,11 +68,21 @@ async fn main() -> anyhow::Result<()> { let ingest_edges = edges.clone(); let ingest_coverage = coverage.clone(); let ingest_resolver = resolver.clone(); + let ingest_search = search.clone(); + let ingest_records = records.clone(); let ingest_handle = tokio::spawn(async move { - run_ingest(ingest_cfg, ingest_edges, ingest_coverage, ingest_resolver).await + run_ingest( + ingest_cfg, + ingest_edges, + ingest_coverage, + ingest_resolver, + ingest_search, + ingest_records, + ) + .await }); - let state = AppState::new(records, slingshot, edges, coverage, knots); + let state = AppState::new(records, slingshot, edges, coverage, knots, search); let app = router(state); tracing::info!(%bind, %hydrant_url, %slingshot_url, "bobbin listening"); diff --git a/crates/knot-proxy/src/dns.rs b/crates/knot-proxy/src/dns.rs index 7fa5d30..5d3ea64 100644 --- a/crates/knot-proxy/src/dns.rs +++ b/crates/knot-proxy/src/dns.rs @@ -22,9 +22,8 @@ impl Resolve for PrivateAddressFilter { let allow_private = self.allow_private; let host = name.as_str().to_owned(); Box::pin(async move { - let resolved: Vec = tokio::net::lookup_host((host.as_str(), 0)) - .await? - .collect(); + let resolved: Vec = + tokio::net::lookup_host((host.as_str(), 0)).await?.collect(); partition_safe(host, allow_private, resolved) }) } @@ -35,9 +34,7 @@ fn partition_safe( allow_private: bool, resolved: Vec, ) -> Result> { - if !allow_private - && let Some(reason) = resolved.iter().find_map(|sa| classify_ip(&sa.ip())) - { + if !allow_private && let Some(reason) = resolved.iter().find_map(|sa| classify_ip(&sa.ip())) { return Err(Box::new(BlockedAddressError { host, reason })); } if resolved.is_empty() {