From e223e872444e341f7d3eef8fee3968970766f565 Mon Sep 17 00:00:00 2001 From: dawn <90008@klbr.net> Date: Sat, 3 Oct 2026 00:25:13 +0300 Subject: [PATCH] [seed] seed absent PDS cursors from listHosts again f83867b dropped seeding new PDS cursors from the relay's listHosts seq. the startup snapshot is taken before the crawler runs, so a host started from it also covers commits made between a repo's backfill and the host's jittered first connect, which a live tail would lose. unlike the old seeding this only fills in a missing cursor: a stored one, the -1 live tail marker, and a running ingestor's in-memory one all win. --- docs/concepts/relay.md | 2 + src/control/seed.rs | 101 +++++++++++++++++++++++++++++++++++++++-- 2 files changed, 98 insertions(+), 5 deletions(-) diff --git a/docs/concepts/relay.md b/docs/concepts/relay.md index 1d47b14..afd5634 100644 --- a/docs/concepts/relay.md +++ b/docs/concepts/relay.md @@ -36,6 +36,8 @@ each discovered host is added as a persistent PDS firehose source (`is_pds: true banned hosts (`status: "banned"`) are skipped. all other statuses are included since the firehose ingestor retries on disconnect and transiently-unavailable hosts will reconnect on their own. +a host we have no cursor for starts from the seq the seed relay reported for it, or from the live head when it reported none, so its retained history is not replayed. hosts that already have a cursor keep it. + seeding runs from latest cursor on restart so new PDS' added to the upstream relay since the last start are picked up automatically (if they haven't through firehose). sources that are already running are detected and skipped, so re-seeding is idempotent. ## crawler sources diff --git a/src/control/seed.rs b/src/control/seed.rs index 5a68285..3a43e10 100644 --- a/src/control/seed.rs +++ b/src/control/seed.rs @@ -9,6 +9,7 @@ use tracing::{debug, info, warn}; use url::Url; use super::firehose::FirehoseHandle; +use crate::db::keys; use crate::pds_discovery::{canonical_pds_host, list_hosts_url}; use crate::state::AppState; @@ -42,6 +43,7 @@ pub(crate) async fn seed_from_list_hosts( /// /// this runs before persisted sources are spawned so existing PDS tasks start /// with updated host status metadata while preserving local firehose cursors. +/// hosts without a cursor get the seed relay's seq for them. pub(crate) async fn refresh_seed_snapshots(seed_urls: &[Url], state: &Arc) { info!("will refresh seed snapshots..."); @@ -220,6 +222,7 @@ async fn apply_seed_snapshot(state: &Arc, hosts: &[Host]) -> miette::R let _status_write = state.pds_status_write.lock(); let mut batch = state.db.inner.batch(); let mut status_updates = Vec::with_capacity(hosts.len()); + let mut cursor_seeds = 0usize; for host in hosts { let hostname = match canonical_pds_host(&host.hostname) { @@ -232,6 +235,10 @@ async fn apply_seed_snapshot(state: &Arc, hosts: &[Host]) -> miette::R let status = map_seed_status(host.status.as_ref()); crate::db::pds_meta::set_status(&mut batch, &state.db.filter, hostname.as_str(), status)?; status_updates.push((hostname.as_str().to_string(), status)); + + if seed_absent_cursor(state, &mut batch, host)? { + cursor_seeds += 1; + } } batch.commit().into_diagnostic()?; @@ -239,6 +246,7 @@ async fn apply_seed_snapshot(state: &Arc, hosts: &[Host]) -> miette::R debug!( hosts = hosts.len(), status_updates = status_updates.len(), + cursor_seeds, "applied listHosts seed snapshot" ); @@ -253,6 +261,38 @@ async fn apply_seed_snapshot(state: &Arc, hosts: &[Host]) -> miette::R Ok(()) } +/// stage the seed relay's seq as the cursor of a host we have no cursor for. +/// +/// at startup this runs before the crawler, so a new host's stream starts from +/// before any of its repos are backfilled instead of from its jittered first +/// connect. a cursor we already hold, including the `-1` live tail marker, +/// always wins, and so does a running ingestor's in-memory one. +fn seed_absent_cursor( + state: &AppState, + batch: &mut fjall::OwnedWriteBatch, + host: &Host, +) -> miette::Result { + let Some(seq) = host + .seq + .and_then(|seq| i64::try_from(seq).ok()) + .filter(|seq| *seq > 0) + else { + return Ok(false); + }; + let Ok(url) = Url::parse(&format!("wss://{}/", host.hostname)) else { + return Ok(false); + }; + if state.firehose_cursors.contains(&url) { + return Ok(false); + } + let key = keys::firehose_cursor_key_from_url(&url); + if state.db.cursors.contains_key(&key).into_diagnostic()? { + return Ok(false); + } + batch.insert(&state.db.cursors, key, seq.to_be_bytes()); + Ok(true) +} + fn map_seed_status(status: Option<&HostStatus>) -> crate::pds_meta::HostStatus { match status { Some(HostStatus::Active) | None => crate::pds_meta::HostStatus::Active, @@ -269,7 +309,6 @@ 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}; use tempfile::tempdir; @@ -291,8 +330,7 @@ mod tests { } #[tokio::test] - async fn apply_seed_snapshot_persists_statuses_and_preserves_absent_cursor() - -> miette::Result<()> { + async fn apply_seed_snapshot_persists_statuses_and_seeds_cursors() -> miette::Result<()> { let tmp = tempdir().into_diagnostic()?; let cfg = Config { database_path: tmp.path().to_path_buf(), @@ -344,8 +382,8 @@ mod tests { assert!(meta.hosts.contains_key("active.example")); assert!(meta.hosts.contains_key("offline.example")); - assert_eq!(persisted_cursor(&state, "active.example")?, None); - assert_eq!(persisted_cursor(&state, "offline.example")?, None); + assert_eq!(persisted_cursor(&state, "active.example")?, Some(100)); + assert_eq!(persisted_cursor(&state, "offline.example")?, Some(5)); let active_url = Url::parse("wss://active.example/").into_diagnostic()?; let in_memory = state @@ -390,6 +428,59 @@ mod tests { Ok(()) } + #[tokio::test] + async fn apply_seed_snapshot_keeps_live_tail_marker() -> miette::Result<()> { + let tmp = tempdir().into_diagnostic()?; + let cfg = Config { + database_path: tmp.path().to_path_buf(), + ..Default::default() + }; + let state = Arc::new(AppState::new(&cfg)?); + let url = Url::parse("wss://active.example/").into_diagnostic()?; + crate::db::set_firehose_cursor(&state.db, &url, -1)?; + + let hosts = vec![Host { + hostname: "active.example".into(), + account_count: Some(42), + seq: Some(500), + status: Some(HostStatus::Active), + extra_data: None, + }]; + + apply_seed_snapshot(&state, &hosts).await?; + + assert_eq!(persisted_cursor(&state, "active.example")?, Some(-1)); + Ok(()) + } + + #[tokio::test] + async fn apply_seed_snapshot_leaves_running_ingestor_cursor() -> miette::Result<()> { + let tmp = tempdir().into_diagnostic()?; + let cfg = Config { + database_path: tmp.path().to_path_buf(), + ..Default::default() + }; + let state = Arc::new(AppState::new(&cfg)?); + let url = Url::parse("wss://active.example/").into_diagnostic()?; + // a live tailing ingestor that has not checkpointed an event yet + let _ = state + .firehose_cursors + .insert_sync(url.clone(), AtomicI64::new(0)); + + let hosts = vec![Host { + hostname: "active.example".into(), + account_count: Some(42), + seq: Some(500), + status: Some(HostStatus::Active), + extra_data: None, + }]; + + apply_seed_snapshot(&state, &hosts).await?; + + assert_eq!(persisted_cursor(&state, "active.example")?, None); + Ok(()) + } + #[tokio::test] async fn relay_banned_snapshot_blocks_admission_until_active() -> miette::Result<()> { let tmp = tempdir().into_diagnostic()?; -- 2.51.2