diff --git a/src/config.rs b/src/config.rs index 35e5575..d5a2f5a 100644 --- a/src/config.rs +++ b/src/config.rs @@ -104,7 +104,6 @@ impl Config { .unwrap_or_else(|| Ok(vec![Url::parse("https://plc.wtf").unwrap()]))?; let full_network: bool = cfg!("FULL_NETWORK", false); - let backfill_concurrency_limit = cfg!("BACKFILL_CONCURRENCY_LIMIT", 128usize); let cursor_save_interval = cfg!("CURSOR_SAVE_INTERVAL", 5, sec); let repo_fetch_timeout = cfg!("REPO_FETCH_TIMEOUT", 300, sec); @@ -119,7 +118,12 @@ impl Config { let identity_cache_size = cfg!("IDENTITY_CACHE_SIZE", 1_000_000u64); let disable_firehose = cfg!("DISABLE_FIREHOSE", false); let disable_backfill = cfg!("DISABLE_BACKFILL", false); - let firehose_workers = cfg!("FIREHOSE_WORKERS", 32usize); + + let backfill_concurrency_limit = cfg!("BACKFILL_CONCURRENCY_LIMIT", 128usize); + let firehose_workers = cfg!( + "FIREHOSE_WORKERS", + full_network.then_some(32usize).unwrap_or(8usize) + ); let ( default_db_worker_threads, diff --git a/src/ingest/mod.rs b/src/ingest/mod.rs index b126ea4..d386b44 100644 --- a/src/ingest/mod.rs +++ b/src/ingest/mod.rs @@ -12,8 +12,6 @@ pub enum IngestMessage { BackfillFinished(Did<'static>), } -pub type BufferedMessage = IngestMessage; - -pub type BufferTx = mpsc::UnboundedSender; +pub type BufferTx = mpsc::UnboundedSender; #[allow(dead_code)] -pub type BufferRx = mpsc::UnboundedReceiver; +pub type BufferRx = mpsc::UnboundedReceiver; diff --git a/src/ingest/worker.rs b/src/ingest/worker.rs index 1904e28..6d5ca15 100644 --- a/src/ingest/worker.rs +++ b/src/ingest/worker.rs @@ -1,5 +1,5 @@ use crate::db::{self, keys}; -use crate::ingest::{BufferedMessage, IngestMessage}; +use crate::ingest::{BufferRx, IngestMessage}; use crate::ops; use crate::resolver::{NoSigningKeyError, ResolverError}; use crate::state::AppState; @@ -56,7 +56,7 @@ enum RepoProcessResult<'s, 'c> { pub struct FirehoseWorker { state: Arc, - rx: mpsc::UnboundedReceiver, + rx: BufferRx, verify_signatures: bool, num_shards: usize, } @@ -75,7 +75,7 @@ struct WorkerContext<'a> { impl FirehoseWorker { pub fn new( state: Arc, - rx: mpsc::UnboundedReceiver, + rx: BufferRx, verify_signatures: bool, num_shards: usize, ) -> Self { @@ -148,7 +148,7 @@ impl FirehoseWorker { // enters the tokio runtime only when necessary (key resolution) fn worker_thread( id: usize, - mut rx: mpsc::UnboundedReceiver, + mut rx: mpsc::UnboundedReceiver, state: Arc, verify_signatures: bool, handle: tokio::runtime::Handle,