From 3dd776db79c96c0b8f5560049933dbb2f2249ba4 Mon Sep 17 00:00:00 2001 From: dawn <90008@gaze.systems> Date: Sun, 15 Feb 2026 18:48:22 +0300 Subject: [PATCH] [db] refactor gauge updates to use shared state and macro --- src/api/repo.rs | 61 +++++++++++------------- src/backfill/mod.rs | 108 +++++++++++++++++++------------------------ src/db/mod.rs | 62 +++++++++++++++++++++++++ src/ingest/worker.rs | 16 +++++-- src/types.rs | 15 +++++- 5 files changed, 163 insertions(+), 99 deletions(-) diff --git a/src/api/repo.rs b/src/api/repo.rs index 0dce503..c4e128f 100644 --- a/src/api/repo.rs +++ b/src/api/repo.rs @@ -1,6 +1,6 @@ use crate::api::AppState; use crate::db::{Db, keys, ser_repo_state}; -use crate::types::RepoState; +use crate::types::{GaugeState, RepoState}; use axum::{Json, Router, extract::State, http::StatusCode, routing::post}; use jacquard::types::did::Did; use serde::Deserialize; @@ -51,7 +51,10 @@ pub async fn handle_repo_add( .map_err(|e| (StatusCode::INTERNAL_SERVER_ERROR, e))?; state.db.update_count_async("repos", added).await; - state.db.update_count_async("pending", added).await; + state + .db + .update_gauge_diff_async(&GaugeState::Synced, &GaugeState::Pending) + .await; // trigger backfill worker state.notify_backfill(); @@ -71,8 +74,6 @@ pub async fn handle_repo_remove( let db = &state.db; let mut batch = db.inner.batch(); let mut removed_repos = 0; - let mut removed_pending = 0; - let mut removed_resync = 0; for did_str in req.dids { let did = Did::new_owned(did_str.as_str()) @@ -95,37 +96,38 @@ pub async fn handle_repo_remove( | crate::types::RepoStatus::Suspended ); - batch.remove(&db.repos, &did_key); - - if was_pending { - batch.remove(&db.pending, &did_key); - removed_pending -= 1; - } - if let Some(resync_bytes) = Db::get(db.resync.clone(), &did_key) + let old_gauge = if was_pending { + GaugeState::Pending + } else if let Some(resync_bytes) = Db::get(db.resync.clone(), &did_key) .await .map_err(|e| (StatusCode::INTERNAL_SERVER_ERROR, e.to_string()))? { let resync_state: crate::types::ResyncState = rmp_serde::from_slice(&resync_bytes) .map_err(|e| (StatusCode::INTERNAL_SERVER_ERROR, e.to_string()))?; - if let crate::types::ResyncState::Error { kind, .. } = resync_state { - match kind { - crate::types::ResyncErrorKind::Ratelimited => { - state.db.update_count_async("error_ratelimited", -1).await - } - crate::types::ResyncErrorKind::Transport => { - state.db.update_count_async("error_transport", -1).await - } - crate::types::ResyncErrorKind::Generic => { - state.db.update_count_async("error_generic", -1).await - } - } - } + let kind = if let crate::types::ResyncState::Error { kind, .. } = resync_state { + Some(kind) + } else { + None + }; + GaugeState::Resync(kind) + } else { + GaugeState::Synced + }; + batch.remove(&db.repos, &did_key); + if was_pending { + batch.remove(&db.pending, &did_key); + } + if old_gauge.is_resync() { batch.remove(&db.resync, &did_key); - removed_resync -= 1; } + state + .db + .update_gauge_diff_async(&old_gauge, &GaugeState::Synced) + .await; + removed_repos -= 1; } } @@ -138,15 +140,6 @@ pub async fn handle_repo_remove( if removed_repos != 0 { state.db.update_count_async("repos", removed_repos).await; } - if removed_pending != 0 { - state - .db - .update_count_async("pending", removed_pending) - .await; - } - if removed_resync != 0 { - state.db.update_count_async("resync", removed_resync).await; - } Ok(StatusCode::OK) } diff --git a/src/backfill/mod.rs b/src/backfill/mod.rs index 4660980..0b3d884 100644 --- a/src/backfill/mod.rs +++ b/src/backfill/mod.rs @@ -3,7 +3,9 @@ use crate::db::{self, Db, keys, ser_repo_state}; use crate::ops; use crate::resolver::ResolverError; use crate::state::AppState; -use crate::types::{AccountEvt, BroadcastEvent, RepoState, RepoStatus, ResyncState, StoredEvent}; +use crate::types::{ + AccountEvt, BroadcastEvent, GaugeState, RepoState, RepoStatus, ResyncState, StoredEvent, +}; use fjall::Slice; use jacquard::api::com_atproto::sync::get_repo::{GetRepo, GetRepoError}; @@ -235,33 +237,50 @@ async fn did_task( Ok(previous_state) => { let did_key = keys::repo_key(&did); - let was_pending = matches!(previous_state.status, RepoStatus::Backfilling); - let was_resync = matches!( - previous_state.status, + // determine old gauge state + // if it was error/suspended etc, we need to know which error kind it was to decrement correctly. + // we have to peek at the resync state. `previous_state` is the repo state, which tells us the Status. + // we have to peek at the resync state. `previous_state` is the repo state, which tells us the Status. + let old_gauge = match previous_state.status { + RepoStatus::Backfilling => GaugeState::Pending, RepoStatus::Error(_) - | RepoStatus::Deactivated - | RepoStatus::Takendown - | RepoStatus::Suspended - ); + | RepoStatus::Deactivated + | RepoStatus::Takendown + | RepoStatus::Suspended => { + // we need to fetch the resync state to know the kind + // if it's missing, we assume Generic (or handle gracefully) + // this is an extra read, but necessary for accurate gauges. + let resync_state = Db::get(db.resync.clone(), &did_key).await.ok().flatten(); + let kind = resync_state.and_then(|b| { + rmp_serde::from_slice::(&b) + .ok() + .and_then(|s| match s { + ResyncState::Error { kind, .. } => Some(kind), + _ => None, + }) + }); + GaugeState::Resync(kind) + } + RepoStatus::Synced => GaugeState::Synced, + }; let mut batch = db.inner.batch(); // remove from pending - if was_pending { + if old_gauge == GaugeState::Pending { batch.remove(&db.pending, pending_key); } // remove from resync - if was_resync { + if old_gauge.is_resync() { batch.remove(&db.resync, &did_key); } tokio::task::spawn_blocking(move || batch.commit().into_diagnostic()) .await .into_diagnostic()??; - if was_pending { - state.db.update_count_async("pending", -1).await; - } - if was_resync { - state.db.update_count_async("resync", -1).await; - } + + state + .db + .update_gauge_diff_async(&old_gauge, &GaugeState::Synced) + .await; let state = state.clone(); tokio::task::spawn_blocking(move || { @@ -349,7 +368,7 @@ async fn did_task( let mut batch = state.db.inner.batch(); batch.insert(&state.db.resync, &did_key, serialized_resync_state); - batch.remove(&state.db.pending, pending_key); + batch.remove(&state.db.pending, pending_key.clone()); if let Some(state_bytes) = serialized_repo_state { batch.insert(&state.db.repos, &did_key, state_bytes); } @@ -359,50 +378,19 @@ async fn did_task( .await .into_diagnostic()??; - state.db.update_count_async("resync", 1).await; - state.db.update_count_async("pending", -1).await; - - // update gauges - if let Some(prev) = prev_kind { - if prev != error_kind { - match prev { - crate::types::ResyncErrorKind::Ratelimited => { - state.db.update_count_async("error_ratelimited", -1).await - } - crate::types::ResyncErrorKind::Transport => { - state.db.update_count_async("error_transport", -1).await - } - crate::types::ResyncErrorKind::Generic => { - state.db.update_count_async("error_generic", -1).await - } - } - match error_kind { - crate::types::ResyncErrorKind::Ratelimited => { - state.db.update_count_async("error_ratelimited", 1).await - } - crate::types::ResyncErrorKind::Transport => { - state.db.update_count_async("error_transport", 1).await - } - crate::types::ResyncErrorKind::Generic => { - state.db.update_count_async("error_generic", 1).await - } - } - } - // if same, do nothing (count already accurate) + let old_gauge = if let Some(k) = prev_kind { + GaugeState::Resync(Some(k)) } else { - // new error - match error_kind { - crate::types::ResyncErrorKind::Ratelimited => { - state.db.update_count_async("error_ratelimited", 1).await - } - crate::types::ResyncErrorKind::Transport => { - state.db.update_count_async("error_transport", 1).await - } - crate::types::ResyncErrorKind::Generic => { - state.db.update_count_async("error_generic", 1).await - } - } - } + GaugeState::Pending + }; + + let new_gauge = GaugeState::Resync(Some(error_kind)); + + state + .db + .update_gauge_diff_async(&old_gauge, &new_gauge) + .await; + Err(e) } } diff --git a/src/db/mod.rs b/src/db/mod.rs index 422749f..d271f3a 100644 --- a/src/db/mod.rs +++ b/src/db/mod.rs @@ -35,6 +35,52 @@ pub struct Db { pub counts_map: HashMap, } +macro_rules! update_gauge_diff_impl { + ($self:ident, $old:ident, $new:ident, $update_method:ident $(, $await:tt)?) => {{ + use crate::types::GaugeState; + + if $old == $new { + return; + } + + // pending + match ($old, $new) { + (GaugeState::Pending, GaugeState::Pending) => {} + (GaugeState::Pending, _) => $self.$update_method("pending", -1) $(.$await)?, + (_, GaugeState::Pending) => $self.$update_method("pending", 1) $(.$await)?, + _ => {} + } + + // resync + let old_resync = $old.is_resync(); + let new_resync = $new.is_resync(); + match (old_resync, new_resync) { + (true, false) => $self.$update_method("resync", -1) $(.$await)?, + (false, true) => $self.$update_method("resync", 1) $(.$await)?, + _ => {} + } + + // error kinds + if let GaugeState::Resync(Some(kind)) = $old { + let key = match kind { + crate::types::ResyncErrorKind::Ratelimited => "error_ratelimited", + crate::types::ResyncErrorKind::Transport => "error_transport", + crate::types::ResyncErrorKind::Generic => "error_generic", + }; + $self.$update_method(key, -1) $(.$await)?; + } + + if let GaugeState::Resync(Some(kind)) = $new { + let key = match kind { + crate::types::ResyncErrorKind::Ratelimited => "error_ratelimited", + crate::types::ResyncErrorKind::Transport => "error_transport", + crate::types::ResyncErrorKind::Generic => "error_generic", + }; + $self.$update_method(key, 1) $(.$await)?; + } + }}; +} + impl Db { pub fn open(cfg: &crate::config::Config) -> Result { let db = Database::builder(&cfg.database_path) @@ -191,6 +237,22 @@ impl Db { .unwrap_or(0) } + pub fn update_gauge_diff( + &self, + old: &crate::types::GaugeState, + new: &crate::types::GaugeState, + ) { + update_gauge_diff_impl!(self, old, new, update_count); + } + + pub async fn update_gauge_diff_async( + &self, + old: &crate::types::GaugeState, + new: &crate::types::GaugeState, + ) { + update_gauge_diff_impl!(self, old, new, update_count_async, await); + } + pub fn update_repo_state( batch: &mut OwnedWriteBatch, repos: &Keyspace, diff --git a/src/ingest/worker.rs b/src/ingest/worker.rs index fd8cc60..1904e28 100644 --- a/src/ingest/worker.rs +++ b/src/ingest/worker.rs @@ -3,7 +3,7 @@ use crate::ingest::{BufferedMessage, IngestMessage}; use crate::ops; use crate::resolver::{NoSigningKeyError, ResolverError}; use crate::state::AppState; -use crate::types::{AccountEvt, BroadcastEvent, IdentityEvt, RepoState, RepoStatus}; +use crate::types::{AccountEvt, BroadcastEvent, GaugeState, IdentityEvt, RepoState, RepoStatus}; use jacquard::api::com_atproto::sync::subscribe_repos::SubscribeReposMessage; use fjall::OwnedWriteBatch; @@ -361,7 +361,9 @@ impl FirehoseWorker { )?; batch.insert(&ctx.state.db.pending, keys::repo_key(did), &[]); batch.commit().into_diagnostic()?; - ctx.state.db.update_count("pending", 1); + ctx.state + .db + .update_gauge_diff(&GaugeState::Synced, &GaugeState::Pending); ctx.state.notify_backfill(); return Ok(RepoProcessResult::Ok(repo_state)); } @@ -507,7 +509,10 @@ impl FirehoseWorker { )?; batch.insert(&ctx.state.db.pending, keys::repo_key(did), &[]); batch.commit().into_diagnostic()?; - ctx.state.db.update_count("pending", 1); + ctx.state.db.update_gauge_diff( + &crate::types::GaugeState::Synced, + &crate::types::GaugeState::Pending, + ); ctx.repo_cache .insert(did.clone().into_static(), repo_state.clone().into_static()); ctx.state.notify_backfill(); @@ -571,7 +576,10 @@ impl FirehoseWorker { batch.commit().into_diagnostic()?; ctx.state.db.update_count("repos", 1); - ctx.state.db.update_count("pending", 1); + ctx.state.db.update_gauge_diff( + &crate::types::GaugeState::Synced, + &crate::types::GaugeState::Pending, + ); ctx.state.notify_backfill(); diff --git a/src/types.rs b/src/types.rs index 604fc00..f93d889 100644 --- a/src/types.rs +++ b/src/types.rs @@ -77,7 +77,7 @@ impl<'i> IntoStatic for RepoState<'i> { // from src/backfill/resync_state.rs -#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)] +#[derive(Debug, Clone, Copy, Serialize, Deserialize, PartialEq, Eq)] pub enum ResyncErrorKind { Ratelimited, Transport, @@ -185,3 +185,16 @@ pub struct StoredEvent<'i> { #[serde(skip_serializing_if = "Option::is_none")] pub cid: Option, } + +#[derive(Debug, PartialEq, Eq, Clone, Copy)] +pub enum GaugeState { + Synced, + Pending, + Resync(Option), +} + +impl GaugeState { + pub fn is_resync(&self) -> bool { + matches!(self, GaugeState::Resync(_)) + } +} -- 2.51.2