From cafa6209b7f13a42bf651593c99deb235e8e2ab8 Mon Sep 17 00:00:00 2001 From: dawn <90008@klbr.net> Date: Sat, 12 Sep 2026 17:52:42 +0300 Subject: [PATCH] [discovery] preserve local bans in seed snapshots --- src/control/seed.rs | 80 ++++++++++++++++++++++++++++++++++++++------- 1 file changed, 68 insertions(+), 12 deletions(-) diff --git a/src/control/seed.rs b/src/control/seed.rs index 5c1bd0b..3c50dc2 100644 --- a/src/control/seed.rs +++ b/src/control/seed.rs @@ -9,7 +9,7 @@ use tracing::{debug, info, warn}; use url::Url; use super::firehose::FirehoseHandle; -use crate::pds_discovery::list_hosts_url; +use crate::pds_discovery::{canonical_pds_host, list_hosts_url}; use crate::state::AppState; const MAX_CONCURRENT_SEEDS: usize = 4; @@ -139,7 +139,7 @@ async fn seed_one( total += body.hosts.len(); let page_hosts = body.hosts.len(); - if let Err(e) = apply_seed_snapshot(state, &body.hosts) { + if let Err(e) = apply_seed_snapshot(state, &body.hosts).await { warn!(err = %e, "failed to apply listHosts seed snapshot"); } @@ -218,19 +218,40 @@ fn body_preview(bytes: &[u8]) -> String { body } -fn apply_seed_snapshot(state: &Arc, hosts: &[Host]) -> miette::Result<()> { +async fn apply_seed_snapshot(state: &Arc, hosts: &[Host]) -> miette::Result<()> { + let _admission = state.pds_admission.lock().await; let mut batch = state.db.inner.batch(); let mut status_updates = Vec::with_capacity(hosts.len()); for host in hosts { - let hostname = host.hostname.as_ref(); + let hostname = match canonical_pds_host(&host.hostname) { + Ok(hostname) => hostname, + Err(error) => { + debug!(hostname = %host.hostname, %error, "invalid hostname in listHosts snapshot, skipping"); + continue; + } + }; let status = map_seed_status(host.status.as_ref()); + let current = state + .db + .filter + .get(crate::db::pds_meta::pds_status_key(hostname.as_str())) + .into_diagnostic()? + .map(|stored| { + rmp_serde::from_slice::(stored.as_ref()) + .into_diagnostic() + }) + .transpose()?; + if current == Some(crate::pds_meta::HostStatus::Banned) { + continue; + } - crate::db::pds_meta::set_status(&mut batch, &state.db.filter, hostname, status)?; - status_updates.push((hostname.to_string(), status)); + crate::db::pds_meta::set_status(&mut batch, &state.db.filter, hostname.as_str(), status)?; + status_updates.push((hostname.as_str().to_string(), status)); } batch.commit().into_diagnostic()?; + state.db.persist()?; debug!( hosts = hosts.len(), status_updates = status_updates.len(), @@ -263,6 +284,7 @@ fn map_seed_status(status: Option<&HostStatus>) -> crate::pds_meta::HostStatus { mod tests { use super::*; use crate::config::Config; + use crate::control::Hydrant; use crate::db::keys; use jacquard_common::CowStr; use std::sync::atomic::{AtomicI64, Ordering}; @@ -284,8 +306,9 @@ mod tests { .transpose() } - #[test] - fn apply_seed_snapshot_persists_statuses_and_preserves_absent_cursor() -> miette::Result<()> { + #[tokio::test] + async fn apply_seed_snapshot_persists_statuses_and_preserves_absent_cursor() + -> miette::Result<()> { let tmp = tempdir().into_diagnostic()?; let cfg = Config { database_path: tmp.path().to_path_buf(), @@ -310,7 +333,7 @@ mod tests { }, ]; - apply_seed_snapshot(&state, &hosts)?; + apply_seed_snapshot(&state, &hosts).await?; assert_eq!( state @@ -349,8 +372,8 @@ mod tests { Ok(()) } - #[test] - fn apply_seed_snapshot_preserves_existing_cursor() -> miette::Result<()> { + #[tokio::test] + async fn apply_seed_snapshot_preserves_existing_cursor() -> miette::Result<()> { let tmp = tempdir().into_diagnostic()?; let cfg = Config { database_path: tmp.path().to_path_buf(), @@ -372,7 +395,7 @@ mod tests { extra_data: None, }]; - apply_seed_snapshot(&state, &hosts)?; + apply_seed_snapshot(&state, &hosts).await?; assert_eq!(persisted_cursor(&state, "active.example")?, Some(150)); let in_memory = state @@ -382,6 +405,39 @@ mod tests { Ok(()) } + + #[tokio::test] + async fn seed_snapshot_cannot_clear_a_committed_ban_after_reopen() -> miette::Result<()> { + let tmp = tempdir().into_diagnostic()?; + let cfg = Config { + database_path: tmp.path().to_path_buf(), + ..Default::default() + }; + let hydrant = Hydrant::new(cfg.clone()).await?; + let hosts = vec![Host { + hostname: "BLOCKED.EXAMPLE.".into(), + account_count: None, + seq: None, + status: Some(HostStatus::Active), + extra_data: None, + }]; + + // the listHosts page was fetched before the local ban committed. + hydrant.pds.ban("blocked.example").await?; + apply_seed_snapshot(&hydrant.state, &hosts).await?; + assert_eq!( + hydrant.state.pds_meta.load().status("blocked.example"), + crate::pds_meta::HostStatus::Banned + ); + + drop(hydrant); + let reopened = AppState::new(&cfg)?; + assert_eq!( + reopened.pds_meta.load().status("blocked.example"), + crate::pds_meta::HostStatus::Banned + ); + Ok(()) + } } async fn read_limited_response( -- 2.51.2