From e6e4bf895f35785839b46a32bfbb17db4cae8fab Mon Sep 17 00:00:00 2001 From: dawn <90008@klbr.net> Date: Sun, 4 Oct 2026 23:39:07 +0300 Subject: [PATCH] [config] build config from any key lookup, not just the process env Signed-off-by: dawn <90008@klbr.net> --- src/config/env.rs | 219 ++++++++++++++++++----------- src/main.rs | 5 +- tests/representation_inventory.tsv | 5 +- 3 files changed, 145 insertions(+), 84 deletions(-) diff --git a/src/config/env.rs b/src/config/env.rs index 722b4aa..150f412 100644 --- a/src/config/env.rs +++ b/src/config/env.rs @@ -11,18 +11,16 @@ use crate::pds_meta::{TierPolicy, TierRule}; #[doc(hidden)] #[macro_export] macro_rules! __cfg { - (@val $key:expr) => { - std::env::var(concat!("HYDRANT_", $key)) + (@val $lookup:ident, $key:expr) => { + $lookup(concat!("HYDRANT_", $key)) }; - ($key:expr, $default:expr, sec) => { - cfg!(@val $key) - .ok() + ($lookup:ident, $key:expr, $default:expr, sec) => { + $crate::__cfg!(@val $lookup, $key) .and_then(|s| humantime::parse_duration(&s).ok()) .unwrap_or($default) }; - ($key:expr, $default:expr) => { - cfg!(@val $key) - .ok() + ($lookup:ident, $key:expr, $default:expr) => { + $crate::__cfg!(@val $lookup, $key) .and_then(|s| s.parse().ok()) .unwrap_or($default.to_owned()) .into() @@ -60,13 +58,6 @@ fn load_dotenv() { } } -fn parse_seed_poll_interval(default: Option) -> Option { - parse_seed_poll_interval_value( - std::env::var("HYDRANT_SEED_POLL_INTERVAL").ok().as_deref(), - default, - ) -} - fn parse_seed_poll_interval_value( value: Option<&str>, default: Option, @@ -81,13 +72,6 @@ fn parse_seed_poll_interval_value( } } -fn parse_new_host_limit(default: Option) -> Option { - parse_new_host_limit_value( - std::env::var("HYDRANT_NEW_HOST_LIMIT").ok().as_deref(), - default, - ) -} - fn parse_new_host_limit_value(value: Option<&str>, default: Option) -> Option { match value.map(str::trim) { Some(value) if value.eq_ignore_ascii_case("none") => None, @@ -100,20 +84,25 @@ impl Config { /// reads and builds the config from environment variables, loading `.env` first if present. pub fn from_env() -> Result { load_dotenv(); + Self::from_lookup(|key| std::env::var(key).ok()) + } + /// builds the config from `HYDRANT_*` keys resolved by `lookup` instead of the process + /// environment, so an embedder can configure hydrant without touching its own env or a `.env`. + pub fn from_lookup(lookup: impl Fn(&str) -> Option) -> Result { // full_network is read first since it determines which defaults to use. // relay mode defaults to true so that the network is indexed by default. #[cfg(feature = "relay")] let default_full_network = true; #[cfg(not(feature = "relay"))] let default_full_network = false; - let full_network: bool = cfg!("FULL_NETWORK", default_full_network); + let full_network: bool = cfg!(lookup, "FULL_NETWORK", default_full_network); let defaults = full_network .then(Self::full_network) .unwrap_or_else(Self::default); - let relay_hosts = match std::env::var("HYDRANT_RELAY_HOSTS") { - Ok(hosts) if !hosts.trim().is_empty() => hosts + let relay_hosts = match lookup("HYDRANT_RELAY_HOSTS") { + Some(hosts) if !hosts.trim().is_empty() => hosts .split(',') .filter_map(|s| { let s = s.trim(); @@ -128,18 +117,17 @@ impl Config { }) .collect(), // HYDRANT_RELAY_HOSTS explicitly set to "" - Ok(_) => vec![], + Some(_) => vec![], // not set at all, fall back to RELAY_HOST (bare URL, no pds:: prefix support here) - Err(_) => match std::env::var("HYDRANT_RELAY_HOST") { - Ok(s) if !s.trim().is_empty() => { + None => match lookup("HYDRANT_RELAY_HOST") { + Some(s) if !s.trim().is_empty() => { FirehoseSource::parse(s.trim()).into_iter().collect() } _ => defaults.relays.clone(), }, }; - let plc_urls: Vec = std::env::var("HYDRANT_PLC_URL") - .ok() + let plc_urls: Vec = lookup("HYDRANT_PLC_URL") .map(|s| { s.split(',') .map(|s| Url::parse(s.trim())) @@ -148,45 +136,57 @@ impl Config { }) .unwrap_or_else(|| Ok(defaults.plc_urls.clone()))?; - let cursor_save_interval = cfg!("CURSOR_SAVE_INTERVAL", defaults.cursor_save_interval, sec); - let repo_fetch_timeout = cfg!("REPO_FETCH_TIMEOUT", defaults.repo_fetch_timeout, sec); - let maintenance_timeout = cfg!("MAINTENANCE_TIMEOUT", defaults.maintenance_timeout, sec); - let max_car_body_bytes = cfg!("MAX_CAR_BODY_BYTES", defaults.max_car_body_bytes); + let cursor_save_interval = cfg!( + lookup, + "CURSOR_SAVE_INTERVAL", + defaults.cursor_save_interval, + sec + ); + let repo_fetch_timeout = cfg!( + lookup, + "REPO_FETCH_TIMEOUT", + defaults.repo_fetch_timeout, + sec + ); + let maintenance_timeout = cfg!( + lookup, + "MAINTENANCE_TIMEOUT", + defaults.maintenance_timeout, + sec + ); + let max_car_body_bytes = cfg!(lookup, "MAX_CAR_BODY_BYTES", defaults.max_car_body_bytes); - let ephemeral: bool = cfg!("EPHEMERAL", defaults.ephemeral); - let ephemeral_ttl = cfg!("EPHEMERAL_TTL", defaults.ephemeral_ttl, sec); - let history_ttl: Option = std::env::var("HYDRANT_HISTORY_TTL") - .ok() + let ephemeral: bool = cfg!(lookup, "EPHEMERAL", defaults.ephemeral); + let ephemeral_ttl = cfg!(lookup, "EPHEMERAL_TTL", defaults.ephemeral_ttl, sec); + let history_ttl: Option = lookup("HYDRANT_HISTORY_TTL") .filter(|s| !s.trim().is_empty()) .map(|s| { humantime::parse_duration(s.trim()) .map_err(|e| miette::miette!("invalid HYDRANT_HISTORY_TTL: {e}")) }) .transpose()?; - let database_path = cfg!("DATABASE_PATH", defaults.database_path); - let cache_size = cfg!("CACHE_SIZE", defaults.cache_size); - let data_compression = cfg!("DATA_COMPRESSION", defaults.data_compression); - let journal_compression = cfg!("JOURNAL_COMPRESSION", defaults.journal_compression); - - let verify_signatures = cfg!("VERIFY_SIGNATURES", defaults.verify_signatures); - let identity_cache_size = cfg!("IDENTITY_CACHE_SIZE", defaults.identity_cache_size); - let verify_mst: bool = cfg!("VERIFY_MST", defaults.verify_mst); - let rev_clock_skew_secs: i64 = cfg!("REV_CLOCK_SKEW", defaults.rev_clock_skew_secs); - let enable_firehose = cfg!("ENABLE_FIREHOSE", defaults.enable_firehose); - let enable_backfill = cfg!("ENABLE_BACKFILL", defaults.enable_backfill); - let enable_crawler = std::env::var("HYDRANT_ENABLE_CRAWLER") - .ok() - .and_then(|s| s.parse().ok()); + let database_path = cfg!(lookup, "DATABASE_PATH", defaults.database_path); + let cache_size = cfg!(lookup, "CACHE_SIZE", defaults.cache_size); + let data_compression = cfg!(lookup, "DATA_COMPRESSION", defaults.data_compression); + let journal_compression = cfg!(lookup, "JOURNAL_COMPRESSION", defaults.journal_compression); + + let verify_signatures = cfg!(lookup, "VERIFY_SIGNATURES", defaults.verify_signatures); + let identity_cache_size = cfg!(lookup, "IDENTITY_CACHE_SIZE", defaults.identity_cache_size); + let verify_mst: bool = cfg!(lookup, "VERIFY_MST", defaults.verify_mst); + let rev_clock_skew_secs: i64 = cfg!(lookup, "REV_CLOCK_SKEW", defaults.rev_clock_skew_secs); + let enable_firehose = cfg!(lookup, "ENABLE_FIREHOSE", defaults.enable_firehose); + let enable_backfill = cfg!(lookup, "ENABLE_BACKFILL", defaults.enable_backfill); + let enable_crawler = lookup("HYDRANT_ENABLE_CRAWLER").and_then(|s| s.parse().ok()); let backfill_concurrency_limit = cfg!( + lookup, "BACKFILL_CONCURRENCY_LIMIT", defaults.backfill_concurrency_limit ); - let backfill_strategy = cfg!("BACKFILL_STRATEGY", defaults.backfill_strategy); + let backfill_strategy = cfg!(lookup, "BACKFILL_STRATEGY", defaults.backfill_strategy); // comma-separated proxy URLs to spread backfill fetches across egress IPs. - let backfill_proxies: Vec = std::env::var("HYDRANT_BACKFILL_PROXIES") - .ok() + let backfill_proxies: Vec = lookup("HYDRANT_BACKFILL_PROXIES") .map(|s| { s.split(',') .filter_map(|u| { @@ -205,91 +205,114 @@ impl Config { .collect() }) .unwrap_or_else(|| defaults.backfill_proxies.clone()); - let per_pds_concurrency = cfg!("PER_PDS_CONCURRENCY", defaults.per_pds_concurrency); - let firehose_workers = cfg!("FIREHOSE_WORKERS", defaults.firehose_workers); - let firehose_max_failures = cfg!("FIREHOSE_MAX_FAILURES", defaults.firehose_max_failures); + let per_pds_concurrency = cfg!(lookup, "PER_PDS_CONCURRENCY", defaults.per_pds_concurrency); + let firehose_workers = cfg!(lookup, "FIREHOSE_WORKERS", defaults.firehose_workers); + let firehose_max_failures = cfg!( + lookup, + "FIREHOSE_MAX_FAILURES", + defaults.firehose_max_failures + ); - let db_worker_threads = cfg!("DB_WORKER_THREADS", defaults.db_worker_threads); + let db_worker_threads = cfg!(lookup, "DB_WORKER_THREADS", defaults.db_worker_threads); let db_max_journaling_size_mb = cfg!( + lookup, "DB_MAX_JOURNALING_SIZE_MB", defaults.db_max_journaling_size_mb ); let db_events_memtable_size_mb = cfg!( + lookup, "DB_EVENTS_MEMTABLE_SIZE_MB", defaults.db_events_memtable_size_mb ); let db_records_memtable_size_mb = cfg!( + lookup, "DB_RECORDS_MEMTABLE_SIZE_MB", defaults.db_records_memtable_size_mb ); let db_records_bloom_filters: bool = cfg!( + lookup, "DB_RECORDS_BLOOM_FILTERS", defaults.db_records_bloom_filters ); let db_repos_memtable_size_mb = cfg!( + lookup, "DB_REPOS_MEMTABLE_SIZE_MB", defaults.db_repos_memtable_size_mb ); let stream_replay_chunk_size = cfg!( + lookup, "STREAM_REPLAY_CHUNK_SIZE", defaults.stream_replay_chunk_size ); let stream_replay_chunk_pause = cfg!( + lookup, "STREAM_REPLAY_CHUNK_PAUSE", defaults.stream_replay_chunk_pause, sec ); let stream_pending_event_limit = cfg!( + lookup, "STREAM_PENDING_EVENT_LIMIT", defaults.stream_pending_event_limit ); - let stream_send_timeout = cfg!("STREAM_SEND_TIMEOUT", defaults.stream_send_timeout, sec); + let stream_send_timeout = cfg!( + lookup, + "STREAM_SEND_TIMEOUT", + defaults.stream_send_timeout, + sec + ); let crawler_max_pending_repos = cfg!( + lookup, "CRAWLER_MAX_PENDING_REPOS", defaults.crawler_max_pending_repos ); let crawler_resume_pending_repos = cfg!( + lookup, "CRAWLER_RESUME_PENDING_REPOS", defaults.crawler_resume_pending_repos ); - let filter_signals = std::env::var("HYDRANT_FILTER_SIGNALS").ok().map(|s| { + let filter_signals = lookup("HYDRANT_FILTER_SIGNALS").map(|s| { s.split(',') .map(|s| s.trim().to_string()) .filter(|s| !s.is_empty()) .collect() }); - let filter_collections = std::env::var("HYDRANT_FILTER_COLLECTIONS").ok().map(|s| { + let filter_collections = lookup("HYDRANT_FILTER_COLLECTIONS").map(|s| { s.split(',') .map(|s| s.trim().to_string()) .filter(|s| !s.is_empty()) .collect() }); - let filter_excludes = std::env::var("HYDRANT_FILTER_EXCLUDES").ok().map(|s| { + let filter_excludes = lookup("HYDRANT_FILTER_EXCLUDES").map(|s| { s.split(',') .map(|s| s.trim().to_string()) .filter(|s| !s.is_empty()) .collect() }); - let enable_backlinks: bool = cfg!("ENABLE_BACKLINKS", defaults.enable_backlinks); + let enable_backlinks: bool = cfg!(lookup, "ENABLE_BACKLINKS", defaults.enable_backlinks); let get_repo_concurrency_limit = cfg!( + lookup, "GET_REPO_CONCURRENCY_LIMIT", defaults.get_repo_concurrency_limit ); - let verify_cids: bool = cfg!("VERIFY_CIDS", defaults.verify_cids); - let only_index_links: bool = cfg!("ONLY_INDEX_LINKS", defaults.only_index_links); - let max_pds_added_per_day = parse_new_host_limit(defaults.new_host_limit); - let seed_poll_interval = parse_seed_poll_interval(defaults.seed_poll_interval); + let verify_cids: bool = cfg!(lookup, "VERIFY_CIDS", defaults.verify_cids); + let only_index_links: bool = cfg!(lookup, "ONLY_INDEX_LINKS", defaults.only_index_links); + let max_pds_added_per_day = parse_new_host_limit_value( + lookup("HYDRANT_NEW_HOST_LIMIT").as_deref(), + defaults.new_host_limit, + ); + let seed_poll_interval = parse_seed_poll_interval_value( + lookup("HYDRANT_SEED_POLL_INTERVAL").as_deref(), + defaults.seed_poll_interval, + ); let offline_retry_interval: Option = - match std::env::var("HYDRANT_OFFLINE_HOST_RETRY_INTERVAL") - .ok() - .as_deref() - { + match lookup("HYDRANT_OFFLINE_HOST_RETRY_INTERVAL").as_deref() { None => defaults.offline_host_retry_interval, Some("none") => None, Some(s) => humantime::parse_duration(s) @@ -300,7 +323,7 @@ impl Config { // start with built-in tier definitions, then layer in any env-defined overrides. // format: HYDRANT_RATE_TIERS=name:base/mul/hourly/daily,... let mut tiers = defaults.tier_policy.tiers.clone(); - if let Ok(s) = std::env::var("HYDRANT_RATE_TIERS") { + if let Some(s) = lookup("HYDRANT_RATE_TIERS") { for entry in s.split(',') { let entry = entry.trim(); if let Some((name, spec)) = entry.split_once(':') { @@ -316,8 +339,7 @@ impl Config { } } - let seed_hosts: Vec = std::env::var("HYDRANT_SEED_HOSTS") - .ok() + let seed_hosts: Vec = lookup("HYDRANT_SEED_HOSTS") .map(|s| { s.split(',') .filter_map(|u| { @@ -337,7 +359,7 @@ impl Config { // build ordered glob rules from HYDRANT_TIER_RULES let mut rules: Vec = vec![]; let mut tier_rules: Vec<(String, String)> = vec![]; - if let Ok(s) = std::env::var("HYDRANT_TIER_RULES") { + if let Some(s) = lookup("HYDRANT_TIER_RULES") { for entry in s.split(',') { let entry = entry.trim(); if entry.is_empty() { @@ -365,14 +387,14 @@ impl Config { let tier_policy = TierPolicy { tiers, rules }; let default_mode = CrawlerMode::default_for(full_network); - let crawler_sources = match std::env::var("HYDRANT_CRAWLER_URLS") { - Ok(s) => s + let crawler_sources = match lookup("HYDRANT_CRAWLER_URLS") { + Some(s) => s .split(',') .map(|s| s.trim()) .filter(|s| !s.is_empty()) .filter_map(|s| CrawlerSource::parse(s, default_mode)) .collect(), - Err(_) => match default_mode { + None => match default_mode { CrawlerMode::ListRepos => relay_hosts .iter() .map(|source| CrawlerSource { @@ -500,4 +522,43 @@ mod tests { default_interval ); } + + #[test] + fn from_lookup_reads_keys_from_the_lookup() { + let keys = std::collections::HashMap::from([ + ("HYDRANT_DATABASE_PATH", "/tmp/hydrant-lookup-test"), + ( + "HYDRANT_FILTER_SIGNALS", + "sh.tangled.repo, sh.tangled.spindle.member", + ), + ("HYDRANT_RELAY_HOSTS", ""), + ("HYDRANT_HISTORY_TTL", "2h"), + ("HYDRANT_ENABLE_CRAWLER", "false"), + ]); + let config = Config::from_lookup(|key| keys.get(key).map(|v| v.to_string())).unwrap(); + + assert_eq!( + config.database_path, + std::path::PathBuf::from("/tmp/hydrant-lookup-test") + ); + assert_eq!( + config.filter_signals, + Some(vec![ + "sh.tangled.repo".to_string(), + "sh.tangled.spindle.member".to_string() + ]) + ); + assert!(config.relays.is_empty()); + assert_eq!(config.history_ttl, Some(Duration::from_secs(2 * 60 * 60))); + assert_eq!(config.enable_crawler, Some(false)); + } + + #[test] + fn from_lookup_rejects_invalid_values_it_cannot_default() { + let err = Config::from_lookup(|key| { + (key == "HYDRANT_HISTORY_TTL").then(|| "not a duration".to_string()) + }) + .unwrap_err(); + assert!(err.to_string().contains("HYDRANT_HISTORY_TTL"), "{err}"); + } } diff --git a/src/main.rs b/src/main.rs index a44924b..6fea487 100644 --- a/src/main.rs +++ b/src/main.rs @@ -25,8 +25,9 @@ struct AppConfig { impl AppConfig { fn from_env() -> miette::Result { use hydrant::__cfg as cfg; + let env = |key: &str| std::env::var(key).ok(); let api_binds = parse_api_binds()?; - let enable_debug = cfg!("ENABLE_DEBUG", false); + let enable_debug = cfg!(env, "ENABLE_DEBUG", false); let debug_port_default = api_binds .iter() .next() @@ -34,7 +35,7 @@ impl AppConfig { .port() .checked_add(1) .unwrap_or(DEFAULT_DEBUG_PORT); - let debug_port = cfg!("DEBUG_PORT", debug_port_default); + let debug_port = cfg!(env, "DEBUG_PORT", debug_port_default); Ok(Self { api_binds, enable_debug, diff --git a/tests/representation_inventory.tsv b/tests/representation_inventory.tsv index 18d2065..6808cd3 100644 --- a/tests/representation_inventory.tsv +++ b/tests/representation_inventory.tsv @@ -112,7 +112,7 @@ codec-call:src/backlinks/store.rs::delete_repo::rmp_serde::from_slice rust:src/b codec-call:src/backlinks/store.rs::index_record::rmp_serde::from_slice rust:src/backlinks/store.rs::forward_value_persisted_golden_fixture codec-call:src/backlinks/store.rs::index_record::rmp_serde::to_vec rust:src/backlinks/store.rs::forward_value_persisted_golden_fixture codec-call:src/car.rs::parse_car::iroh_car::CarHeader::decode rust:src/car.rs::parses_rooted_car_with_zero_copy_payloads -codec-call:src/config/env.rs::Config::from_env::.parse rust:src/config/env.rs::env_text_parser_contracts +codec-call:src/config/env.rs::Config::from_lookup::.parse rust:src/config/env.rs::from_lookup_reads_keys_from_the_lookup codec-call:src/config/env.rs::parse_new_host_limit_value::.parse rust:src/config/env.rs::env_text_parser_contracts codec-call:src/config/types.rs::CrawlerMode::deserialize::String::deserialize rust:src/config/types.rs::crawler_mode_string_and_msgpack_contracts codec-call:src/control/crawler.rs::CrawlerHandle::add_source::rmp_serde::to_vec rust:src/control/crawler.rs::reset_cursor_clears_relay_first_progress_for_only_that_source @@ -253,9 +253,9 @@ codec-call:src/sparse_mst.rs::mst_node_layer::serde_ipld_dagcbor::from_slice rus codec-call:src/util/mod.rs::deser_status_code::Option::deserialize rust:src/util/mod.rs::http_status_serde_contract codec-call:src/util/mod.rs::did_key_deserialize_str::multibase::decode rust:src/util/mod.rs::xrpc_output_string_serializers_are_stable config-parser:src/config/env.rs::Config::from_env rust:src/config/env.rs::env_text_parser_contracts +config-parser:src/config/env.rs::Config::from_lookup rust:src/config/env.rs::from_lookup_reads_keys_from_the_lookup config-parser:src/config/env.rs::dotenv_assignments rust:src/config/env.rs::env_text_parser_contracts config-parser:src/config/env.rs::load_dotenv rust:src/config/env.rs::env_text_parser_contracts -config-parser:src/config/env.rs::parse_new_host_limit rust:src/config/env.rs::env_text_parser_contracts config-parser:src/config/env.rs::parse_new_host_limit_value rust:src/config/env.rs::env_text_parser_contracts config-parser:src/config/types.rs::CrawlerSource::parse rust:src/config/types.rs::source_and_rate_tier_parser_contracts config-parser:src/config/types.rs::FirehoseSource::parse rust:src/config/types.rs::source_and_rate_tier_parser_contracts @@ -780,7 +780,6 @@ codec-call:src/control/firehose.rs::FirehoseHandle::persist_validated_source::rm codec-call:src/control/repos/mod.rs::read_mini_doc_snapshot::rmp_serde::from_slice rust:src/control/repos/mod.rs::identity_reconcile_trigger_matrix_is_narrow codec-call:src/db/migration/v11.rs::status::rmp_serde::from_slice rust:src/db/migration/v11.rs::splits_legacy_bans_before_canonicalizing_other_status_aliases codec-call:src/db/mod.rs::stage_pds_daily_adds::.to_be_bytes rust:src/db/mod.rs::legacy_firehose_row_without_trusted_marker_is_public -config-parser:src/config/env.rs::parse_seed_poll_interval rust:src/config/env.rs::env_text_parser_contracts config-parser:src/config/env.rs::parse_seed_poll_interval_value rust:src/config/env.rs::env_text_parser_contracts db-key-codec:src/db/keys/mod.rs::trusted_firehose_source_key rust:src/db/keys/mod.rs::shared_key_codec_inventory_contracts db-key-codec:src/db/pds_meta.rs::v5::pds_ban_key nu:tests/pds_host_status_transitions.nu -- 2.51.2