diff --git a/.gitignore b/.gitignore index 532bc62..01243b7 100644 --- a/.gitignore +++ b/.gitignore @@ -3,3 +3,4 @@ .envrc result .env +mock_debug.log diff --git a/flake.nix b/flake.nix index e5aa718..00f6668 100644 --- a/flake.nix +++ b/flake.nix @@ -27,6 +27,7 @@ go cmake websocat + http-nu ]; }; }; diff --git a/src/api/debug.rs b/src/api/debug.rs index 0843d05..701bb8e 100644 --- a/src/api/debug.rs +++ b/src/api/debug.rs @@ -106,6 +106,9 @@ fn deserialize_value(partition: &str, value: &[u8]) -> Value { if let Ok(arr) = value.try_into() { return Value::Number(u64::from_be_bytes(arr).into()); } + if let Ok(s) = String::from_utf8(value.to_vec()) { + return Value::String(s); + } } "blocks" => { if let Ok(val) = serde_ipld_dagcbor::from_slice::(value) { diff --git a/src/config.rs b/src/config.rs index 55de3be..7359852 100644 --- a/src/config.rs +++ b/src/config.rs @@ -52,6 +52,8 @@ pub struct Config { pub enable_debug: bool, pub verify_signatures: SignatureVerification, pub identity_cache_size: u64, + pub disable_firehose: bool, + pub disable_backfill: bool, } impl Config { @@ -97,6 +99,8 @@ impl Config { let debug_port = cfg!("DEBUG_PORT", 3001u16); let verify_signatures = cfg!("VERIFY_SIGNATURES", SignatureVerification::Full); let identity_cache_size = cfg!("IDENTITY_CACHE_SIZE", 100_000u64); + let disable_firehose = cfg!("DISABLE_FIREHOSE", false); + let disable_backfill = cfg!("DISABLE_BACKFILL", false); Ok(Self { database_path, @@ -114,6 +118,8 @@ impl Config { enable_debug, verify_signatures, identity_cache_size, + disable_firehose, + disable_backfill, }) } } diff --git a/src/main.rs b/src/main.rs index b8e4619..8d0e6f8 100644 --- a/src/main.rs +++ b/src/main.rs @@ -39,34 +39,23 @@ async fn main() -> miette::Result<()> { ); } - tokio::spawn({ - let state = state.clone(); - let timeout = cfg.repo_fetch_timeout; - BackfillWorker::new( - state, - backfill_rx, - timeout, - cfg.backfill_concurrency_limit, - matches!( - cfg.verify_signatures, - SignatureVerification::Full | SignatureVerification::BackfillOnly - ), - ) - .run() - }); - - let firehose_worker = std::thread::spawn({ - let state = state.clone(); - let handle = tokio::runtime::Handle::current(); - move || { - FirehoseWorker::new( + if !cfg.disable_backfill { + tokio::spawn({ + let state = state.clone(); + let timeout = cfg.repo_fetch_timeout; + BackfillWorker::new( state, - buffer_rx, - matches!(cfg.verify_signatures, SignatureVerification::Full), + backfill_rx, + timeout, + cfg.backfill_concurrency_limit, + matches!( + cfg.verify_signatures, + SignatureVerification::Full | SignatureVerification::BackfillOnly + ), ) - .run(handle) - } - }); + .run() + }); + } if let Err(e) = spawn_blocking({ let state = state.clone(); @@ -167,27 +156,49 @@ async fn main() -> miette::Result<()> { ); } - let ingestor = FirehoseIngestor::new( - state.clone(), - buffer_tx, - cfg.relay_host, - cfg.full_network, - matches!(cfg.verify_signatures, SignatureVerification::Full), - ); + let tasks = if !cfg.disable_firehose { + let firehose_worker = std::thread::spawn({ + let state = state.clone(); + let handle = tokio::runtime::Handle::current(); + move || { + FirehoseWorker::new( + state, + buffer_rx, + matches!(cfg.verify_signatures, SignatureVerification::Full), + ) + .run(handle) + } + }); + + let ingestor = FirehoseIngestor::new( + state.clone(), + buffer_tx, + cfg.relay_host, + cfg.full_network, + matches!(cfg.verify_signatures, SignatureVerification::Full), + ); - let res = futures::future::try_join_all::<[BoxFuture<_>; _]>([ - Box::pin( - tokio::task::spawn_blocking(move || { - firehose_worker - .join() - .map_err(|e| miette::miette!("buffer processor thread died: {e:?}")) - }) - .map(|r| r.into_diagnostic().flatten().flatten()), - ), - Box::pin(ingestor.run()), - ]); - if let Err(e) = res.await { - error!("ingestor or buffer processor died: {e}"); + vec![ + Box::pin( + tokio::task::spawn_blocking(move || { + firehose_worker + .join() + .map_err(|e| miette::miette!("buffer processor thread died: {e:?}")) + }) + .map(|r| r.into_diagnostic().flatten().flatten()), + ) as BoxFuture<_>, + Box::pin(ingestor.run()), + ] + } else { + info!("firehose ingestion disabled by config"); + // if firehose is disabled, we just wait indefinitely (or until signal) + // essentially we just want to keep the main thread alive for the other components + vec![Box::pin(futures::future::pending::>()) as BoxFuture<_>] + }; + + let res = futures::future::select_all(tasks); + if let (Err(e), _, _) = res.await { + error!("critical worker died: {e}"); db::check_poisoned_report(&e); } diff --git a/tests/mock_relay.nu b/tests/mock_relay.nu new file mode 100644 index 0000000..b86faca --- /dev/null +++ b/tests/mock_relay.nu @@ -0,0 +1,62 @@ +# mock_relay.nu + +# A closure that handles HTTP requests +{|req| + + + # check path + if ($req.path | str starts-with "/xrpc/com.atproto.sync.listRepos") { + + # parse query params if any + let query_string = ($req.path | split row "?" | get 1? | default "") + let params = if ($query_string | is-empty) { + [] + } else { + ($query_string | split row "&" | each { |it| $it | split row "=" }) + } + let cursor = ($params | where { |x| $x.0 == "cursor" } | get 0?.1?) + + # define some mock repos + let all_repos = [ + { did: "did:web:mock1.com", head: "bafyreidf747c4x3lps3k4n357l3a3r57k3k465743k573k465743k5", rev: "3j6s746574657" }, + { did: "did:web:mock2.com", head: "bafyreidf747c4x3lps3k4n357l3a3r57k3k465743k573k465743k5", rev: "3j6s746574657" }, + { did: "did:web:mock3.com", head: "bafyreidf747c4x3lps3k4n357l3a3r57k3k465743k573k465743k5", rev: "3j6s746574657" }, + { did: "did:web:mock4.com", head: "bafyreidf747c4x3lps3k4n357l3a3r57k3k465743k573k465743k5", rev: "3j6s746574657" }, + { did: "did:web:mock5.com", head: "bafyreidf747c4x3lps3k4n357l3a3r57k3k465743k573k465743k5", rev: "3j6s746574657" } + ] + + let repos = if ($cursor == "50") { + [] + } else { + $all_repos + } + + let next_cursor = if ($cursor == "50") { + null + } else { + "50" + } + + { + cursor: $next_cursor, + repos: $repos + } + | to json + | metadata set --merge { + http.response: { + headers: { + "Content-Type": "application/json" + } + } + } + + } else { + # 404 + "not found" + | metadata set --merge { + http.response: { + status: 404 + } + } + } +} diff --git a/tests/verify_crawler.nu b/tests/verify_crawler.nu new file mode 100644 index 0000000..5b762ad --- /dev/null +++ b/tests/verify_crawler.nu @@ -0,0 +1,131 @@ +#!/usr/bin/env nu +use common.nu * + +def main [] { + # 1. ensure http-nu is installed + if (which http-nu | is-empty) { + print "http-nu not found, installing..." + cargo install http-nu + } + + # 2. setup ports and paths + let port = 3006 + let mock_port = 3008 + let url = $"http://localhost:($port)" + let debug_url = $"http://localhost:($port + 1)" + let mock_url = $"http://localhost:($mock_port)" + let db_path = (mktemp -d -t hydrant_full_net.XXXXXX) + + print $"testing full network crawler..." + print $"database path: ($db_path)" + + # 3. start mock relay + print $"starting mock relay on ($mock_port)..." + let mock_pid = ( + bash -c $"http-nu :($mock_port) tests/mock_relay.nu > ($db_path)/mock.log 2>&1 & echo $!" + | str trim + | into int + ) + print $"mock relay pid: ($mock_pid)" + + # give mock relay a moment + sleep 1sec + + # 4. start hydrant in full network mode, firehose disabled + let binary = build-hydrant + + let log_file = $"($db_path)/hydrant.log" + print $"starting hydrant - logs at ($log_file)..." + + let hydrant_pid = ( + with-env { + HYDRANT_DATABASE_PATH: ($db_path), + HYDRANT_FULL_NETWORK: "true", + HYDRANT_RELAY_HOST: ($mock_url), + HYDRANT_DISABLE_FIREHOSE: "true", + HYDRANT_DISABLE_BACKFILL: "true", + HYDRANT_API_PORT: ($port | into string), + HYDRANT_ENABLE_DEBUG: "true", # for stats checking + HYDRANT_DEBUG_PORT: ($port + 1 | into string), + HYDRANT_LOG_LEVEL: "debug", + HYDRANT_CURSOR_SAVE_INTERVAL: "1" # faster save + } { + sh -c $"($binary) >($log_file) 2>&1 & echo $!" | str trim | into int + } + ) + print $"hydrant started with pid: ($hydrant_pid)" + + mut success = false + + try { + if (wait-for-api $url) { + print "hydrant api is up." + + # wait for crawler to run (it runs on startup) + print "waiting for crawler to fetch repos..." + + # retry check for 30s + for i in 1..30 { + let stats = (http get $"($url)/stats?accurate=true").counts + let pending = ($stats.pending | into int) + let repos = ($stats.repos | default 0 | into int) + + # we expect 5 repos from the mock + print $"[($i)/30] pending: ($pending), known_repos: ($repos)" + + if $repos >= 5 { + print "crawler successfully discovered repos!" + $success = true + break + } + + sleep 1sec + } + + if not $success { + print "timeout waiting for crawler." + } + + # check cursor persistence + print "verifying crawler cursor persistence..." + let cursor_check = try { + let cursor_res = (http get $"($debug_url)/debug/get?partition=cursors&key=crawler_cursor").value + print $"cursor value from debug: ($cursor_res)" + + if $cursor_res == "50" { + print "cursor verified." + true + } else { + print "cursor mismatch or missing." + false + } + } catch { + print "failed to get cursor from debug endpoint" + false + } + if not $cursor_check { $success = false } + + } else { + print "hydrant failed to start." + } + } catch { |e| + print $"test failed with error: ($e)" + } + + # cleanup + print "stopping processes..." + try { kill $hydrant_pid } + try { kill $mock_pid } + + if $success { + print "test passed!" + exit 0 + } else { + print "test failed!" + print "hydrant logs:" + open $log_file | tail -n 20 + print "mock logs:" + open $"($db_path)/mock.log" + exit 1 + } +}