From 22ec0bbe33a292ecdeb2dee42cbc73125b50414b Mon Sep 17 00:00:00 2001 From: phil Date: Sat, 7 Mar 2026 10:31:04 -0500 Subject: [PATCH] busy logs => trace, don't die after backfill --- src/main.rs | 16 +++++++++-- src/server/hello.rs | 2 +- src/storage/collection_index.rs | 5 ++++ src/storage/resync_queue.rs | 4 +-- src/sync/backfill.rs | 4 +-- src/sync/firehose/commit_event.rs | 31 +++++++++++++++++++++ src/sync/firehose/event_dispatcher.rs | 40 +++++++++++++++++++-------- src/sync/firehose/mod.rs | 31 +++++++++++++++------ src/sync/resync/dispatcher.rs | 6 ++-- 9 files changed, 110 insertions(+), 29 deletions(-) diff --git a/src/main.rs b/src/main.rs index d38e40e..38c1245 100644 --- a/src/main.rs +++ b/src/main.rs @@ -4,6 +4,7 @@ use std::path::PathBuf; use clap::Parser; use tokio::task::JoinSet; use tokio_util::sync::CancellationToken; +use tracing::info; use lightrail::error::{Error, Result}; use lightrail::identity; @@ -45,6 +46,11 @@ struct Args { /// Max identities kept in in-process identity cache #[arg(long, env = "LIGHTRAIL_IDENT_CACHE_SIZE", default_value_t = 1_000_000)] ident_cache_size: u64, + + /// Log an error when a commit claims a collection birth but the index + /// already has that collection for the DID (temporary diagnostic flag). + #[arg(long, env = "LIGHTRAIL_VALIDATE_BIRTHS")] + validate_births: bool, } #[tokio::main] @@ -86,8 +92,9 @@ async fn main() -> Result<()> { let db = db.clone(); let host = subscribe_host.clone(); let resolver = resolver.clone(); + let validate_births = args.validate_births; async move { - let mut sub = firehose::Subscriber::new(host, db, resolver); + let mut sub = firehose::Subscriber::new(host, db, resolver, validate_births); sub.run(token).await } }); @@ -96,7 +103,12 @@ async fn main() -> Result<()> { let token = token.clone(); let db = db.clone(); let host = subscribe_host.clone(); - async move { backfill::run(host, db, token).await } + async move { + backfill::run(host, db, token.clone()).await?; + info!("backfill task idle"); + token.cancelled().await; + Ok(()) + } }); tasks.spawn({ diff --git a/src/server/hello.rs b/src/server/hello.rs index 66affff..8943c9d 100644 --- a/src/server/hello.rs +++ b/src/server/hello.rs @@ -13,7 +13,7 @@ this is an atproto backfill-by-collection assister. available endpoints: - /xrpc/com.atproto.sync.listRepos - - /xrpc/com.atproto.sync.listReposByCollection <=the real deal + - /xrpc/com.atproto.sync.listReposByCollection <= the real deal - /xrpc/com.atproto.sync.getRepoStatus source: https://tangled.org/microcosm.blue/lightrail diff --git a/src/storage/collection_index.rs b/src/storage/collection_index.rs index ecc16c2..60c4c66 100644 --- a/src/storage/collection_index.rs +++ b/src/storage/collection_index.rs @@ -207,6 +207,11 @@ pub fn remove_into( batch.remove(&db.ks, cbr(did, collection)); } +/// Return `true` if `did` has `collection` in the cbr index. +pub fn has_collection(db: &DbRef, did: Did<'_>, collection: Nsid<'_>) -> StorageResult { + Ok(db.ks.get(cbr(did, collection))?.is_some()) +} + /// Insert a `(collection, did)` pair into both the rbc and cbr indexes. pub fn insert(db: &DbRef, did: Did<'_>, collection: Nsid<'_>) -> StorageResult<()> { let mut batch = db.database.batch(); diff --git a/src/storage/resync_queue.rs b/src/storage/resync_queue.rs index dd6ac7a..b7f29b1 100644 --- a/src/storage/resync_queue.rs +++ b/src/storage/resync_queue.rs @@ -7,7 +7,7 @@ use std::collections::HashSet; use fjall::util::prefixed_range; use jacquard_common::types::string::Did; -use tracing::{debug, info, warn}; +use tracing::{debug, trace, warn}; use crate::storage::{ DbRef, PREFIX_RESYNC_QUEUE, @@ -280,7 +280,7 @@ pub fn claim_resync( batch.insert(&db.ks, &repo_key, repo::encode_repo_info(&new_info)); batch.commit()?; - info!( + trace!( did = item.did.as_str(), reason = %item.retry_reason, retry = item.retry_count, diff --git a/src/sync/backfill.rs b/src/sync/backfill.rs index 70c8118..2be335a 100644 --- a/src/sync/backfill.rs +++ b/src/sync/backfill.rs @@ -9,7 +9,7 @@ use jacquard_api::com_atproto::sync::list_repos::ListRepos; use jacquard_common::url::Host; use jacquard_common::{IntoStatic, xrpc::XrpcExt}; -use tracing::{info, warn}; +use tracing::{info, trace, warn}; use crate::error::Result; use crate::storage::{ @@ -141,7 +141,7 @@ pub async fn run(host: Host, db: DbRef, token: tokio_util::sync::CancellationTok total_queued += page_queued; - info!( + trace!( host = %host, page_repos = page_len, page_queued, diff --git a/src/sync/firehose/commit_event.rs b/src/sync/firehose/commit_event.rs index ea0ba76..561a542 100644 --- a/src/sync/firehose/commit_event.rs +++ b/src/sync/firehose/commit_event.rs @@ -42,6 +42,7 @@ pub(super) async fn process_commit_event( commit: Box>, resolver: &Resolver, db: &DbRef, + validate_births: bool, ) -> crate::error::Result<()> { let did = commit.repo.clone(); @@ -127,6 +128,7 @@ pub(super) async fn process_commit_event( new_mst_root_bytes, born, died, + validate_births, ) }) .await??; @@ -406,6 +408,7 @@ fn process_blocking( new_mst_root_bytes: Vec, born: Vec>, died: Vec>, + validate_births: bool, ) -> crate::error::Result<()> { // Load the current repo state and chain tip (may be absent for new repos). let (info, prev) = match storage::repo::get(db, did.clone())? { @@ -445,6 +448,34 @@ fn process_blocking( } } + // Temporary birth validation: verify that collections we think are newly + // born aren't already in the index (which would mean our CAR-adjacency + // heuristic missed a pre-existing key). Remove this block once we are + // confident in the birth detection logic. + if validate_births { + for coll in &born { + match storage::collection_index::has_collection(db, did.clone(), coll.clone()) { + Ok(true) => { + tracing::error!( + did = %did.as_str(), + collection = coll.as_str(), + "birth validation: detected spurious birth — \ + collection already exists in index", + ); + } + Ok(false) => {} + Err(e) => { + tracing::warn!( + did = %did.as_str(), + collection = coll.as_str(), + error = %e, + "birth validation: index lookup failed", + ); + } + } + } + } + // All checks passed — update the chain tip for next-commit validation. storage::repo::put_prev( db, diff --git a/src/sync/firehose/event_dispatcher.rs b/src/sync/firehose/event_dispatcher.rs index b000552..4afd6ea 100644 --- a/src/sync/firehose/event_dispatcher.rs +++ b/src/sync/firehose/event_dispatcher.rs @@ -28,7 +28,7 @@ use std::time::Instant; use jacquard_api::com_atproto::sync::subscribe_repos::{Account, Commit, Identity, Sync}; use jacquard_common::types::string::Did; use tokio::task::{Id as TaskId, JoinError, JoinSet}; -use tracing::{debug, error, warn}; +use tracing::{error, trace, warn}; use crate::storage::DbRef; @@ -50,6 +50,7 @@ pub(crate) struct CommitDispatcher { max_concurrent: usize, resolver: Arc, db: DbRef, + validate_births: bool, } struct PendingCommit { @@ -101,7 +102,12 @@ pub(crate) struct CommitWorkerResult { // --------------------------------------------------------------------------- impl CommitDispatcher { - pub fn new(resolver: Arc, db: DbRef, max_concurrent: usize) -> Self { + pub fn new( + resolver: Arc, + db: DbRef, + max_concurrent: usize, + validate_births: bool, + ) -> Self { Self { queues: HashMap::new(), busy: HashSet::new(), @@ -111,6 +117,7 @@ impl CommitDispatcher { max_concurrent, resolver, db, + validate_births, } } @@ -204,10 +211,19 @@ impl CommitDispatcher { let db = self.db.clone(); let did_for_result = did.clone(); let seq = pending.seq(); + let validate_births = self.validate_births; let handle = self.workers.spawn(async move { match pending { PendingWork::Commit(p) => { - run_commit_event_worker(p.commit, p.seq, did_for_result, resolver, db).await + run_commit_event_worker( + p.commit, + p.seq, + did_for_result, + resolver, + db, + validate_births, + ) + .await } PendingWork::Sync(p) => { run_sync_event_worker(p.sync, p.seq, did_for_result, resolver, db).await @@ -221,7 +237,7 @@ impl CommitDispatcher { } }); self.task_did_seq.insert(handle.id(), (did, seq)); - debug!(seq, "spawned worker"); + trace!(seq, "spawned worker"); } metrics::gauge!("lightrail_firehose_workers").set(self.workers.len() as f64); @@ -328,10 +344,12 @@ async fn run_commit_event_worker( did: Did<'static>, resolver: Arc, db: DbRef, + validate_births: bool, ) -> CommitWorkerResult { - let outcome = super::commit_event::process_commit_event(commit, &resolver, &db) - .await - .map_err(|e| e.to_string()); + let outcome = + super::commit_event::process_commit_event(commit, &resolver, &db, validate_births) + .await + .map_err(|e| e.to_string()); CommitWorkerResult { did, seq, outcome } } @@ -434,7 +452,7 @@ mod tests { async fn commits_for_same_did_are_sequential() { let db = crate::storage::open_temporary().unwrap(); let resolver = make_resolver(); - let mut d = CommitDispatcher::new(resolver, db, 4); + let mut d = CommitDispatcher::new(resolver, db, 4, false); let did: Did<'static> = Did::new_owned("did:plc:testsequential").unwrap(); let c1 = { @@ -474,7 +492,7 @@ mod tests { async fn commits_for_different_dids_run_in_parallel() { let db = crate::storage::open_temporary().unwrap(); let resolver = make_resolver(); - let mut d = CommitDispatcher::new(resolver, db, 4); + let mut d = CommitDispatcher::new(resolver, db, 4, false); let did_a: Did<'static> = Did::new_owned("did:plc:testa").unwrap(); let did_b: Did<'static> = Did::new_owned("did:plc:testb").unwrap(); @@ -491,7 +509,7 @@ mod tests { async fn watermark_advances_after_completion() { let db = crate::storage::open_temporary().unwrap(); let resolver = make_resolver(); - let mut d = CommitDispatcher::new(resolver, db, 4); + let mut d = CommitDispatcher::new(resolver, db, 4, false); let did_a: Did<'static> = Did::new_owned("did:plc:testwma").unwrap(); let did_b: Did<'static> = Did::new_owned("did:plc:testwmb").unwrap(); @@ -512,7 +530,7 @@ mod tests { async fn stalled_seq_evicted_from_watermark() { let db = crate::storage::open_temporary().unwrap(); let resolver = make_resolver(); - let mut d = CommitDispatcher::new(resolver, db, 4); + let mut d = CommitDispatcher::new(resolver, db, 4, false); // Manually inject an old entry into outstanding without spawning a worker. let stale_instant = Instant::now() - std::time::Duration::from_secs(STALL_EVICT_SECS + 1); diff --git a/src/sync/firehose/mod.rs b/src/sync/firehose/mod.rs index 6edccbc..bc25ef3 100644 --- a/src/sync/firehose/mod.rs +++ b/src/sync/firehose/mod.rs @@ -34,7 +34,7 @@ use jacquard_api::com_atproto::sync::subscribe_repos::{SubscribeRepos, Subscribe use jacquard_common::url::Host; use jacquard_common::xrpc::SubscriptionExt; use jacquard_common::{StreamErrorKind, TungsteniteClient}; -use tracing::{debug, info, warn}; +use tracing::{debug, info, trace, warn}; use crate::error::Result; use crate::storage::{self, DbRef}; @@ -54,11 +54,22 @@ pub struct Subscriber { host: Host, db: DbRef, resolver: Arc, + validate_births: bool, } impl Subscriber { - pub fn new(host: Host, db: DbRef, resolver: Arc) -> Self { - Self { host, db, resolver } + pub fn new( + host: Host, + db: DbRef, + resolver: Arc, + validate_births: bool, + ) -> Self { + Self { + host, + db, + resolver, + validate_births, + } } /// Connect and run the subscriber loop, reconnecting on disconnect. @@ -75,8 +86,12 @@ impl Subscriber { // The dispatcher survives reconnects so in-flight workers keep running // and the watermark doesn't regress. - let mut dispatcher = - CommitDispatcher::new(self.resolver.clone(), self.db.clone(), MAX_COMMIT_WORKERS); + let mut dispatcher = CommitDispatcher::new( + self.resolver.clone(), + self.db.clone(), + MAX_COMMIT_WORKERS, + self.validate_births, + ); let mut last_seq: i64 = 0; let mut cursor_tick = Instant::now(); @@ -237,18 +252,18 @@ impl Subscriber { // #info and unknown frames (e.g. server error frames like FutureCursor) // carry no sequence number and don't advance the cursor. SubscribeReposMessage::Info(info) => { - debug!(name = %info.name, message = ?info.message, "firehose #info"); + info!(name = %info.name, message = ?info.message, "firehose #info"); return None; } SubscribeReposMessage::Unknown(_) => { - debug!("firehose unknown frame"); + info!("firehose unknown frame"); return None; } }; match msg { SubscribeReposMessage::Commit(commit) => { - debug!(did = %commit.repo, rev = %commit.rev, seq, "firehose #commit"); + trace!(did = %commit.repo, rev = %commit.rev, seq, "firehose #commit"); dispatcher.enqueue(commit, seq); } SubscribeReposMessage::Sync(sync) => { diff --git a/src/sync/resync/dispatcher.rs b/src/sync/resync/dispatcher.rs index 8658824..6890f60 100644 --- a/src/sync/resync/dispatcher.rs +++ b/src/sync/resync/dispatcher.rs @@ -16,7 +16,7 @@ use std::time::{Duration, Instant}; use super::{GetCollectionsError, ResyncError}; use tokio::task::{Id as TaskId, JoinSet}; -use tracing::{debug, error, info, warn}; +use tracing::{debug, error, info, trace, warn}; use crate::error::Result; use crate::storage::{ @@ -121,7 +121,7 @@ pub async fn run( .spawn(async move { run_worker(item, &resolver, &client, &db).await }); task_dids.insert(handle.id(), did_str.clone()); metrics::gauge!("lightrail_resync_workers").set(workers.len() as f64); - debug!(did = %did_str, running = workers.len(), "spawned resync worker"); + trace!(did = %did_str, running = workers.len(), "spawned resync worker"); } Ok(Ok(None)) => break, // queue empty or all ready items busy Ok(Err(e)) => { @@ -236,7 +236,7 @@ async fn run_worker( let did_str = item.did.as_str().to_string(); match super::index_repo(client, resolver, item.did, db).await { Ok(()) => { - info!(did = %did_str, "resync completed"); + trace!(did = %did_str, "resync completed"); WorkerOutcome::Success } Err(ResyncError::Fetch(GetCollectionsError::RateLimited(host))) => { -- 2.51.2