diff --git a/Cargo.lock b/Cargo.lock index 20ed1d0cf..b0b83d979 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -3069,6 +3069,7 @@ dependencies = [ "gix-transport", "jacquard-common", "jacquard-identity", + "metrics", "reqwest 0.13.5", "scc", "sqlx", diff --git a/gitmirror/crates/gitmirror-db/.sqlx/query-6a8ebbb538b3254515d113c27135f20cb2a40f8a2f3c89b0a07115caf920c876.json b/gitmirror/crates/gitmirror-db/.sqlx/query-6a8ebbb538b3254515d113c27135f20cb2a40f8a2f3c89b0a07115caf920c876.json new file mode 100644 index 000000000..67e34aab3 --- /dev/null +++ b/gitmirror/crates/gitmirror-db/.sqlx/query-6a8ebbb538b3254515d113c27135f20cb2a40f8a2f3c89b0a07115caf920c876.json @@ -0,0 +1,38 @@ +{ + "db_name": "PostgreSQL", + "query": "select case when state = 'error' and retry_count >= $2 then 'error'\n when state = 'error' then 'backoff'\n else state end as \"state!\",\n (state <> 'active' and (retry_after_us is null or retry_after_us <= $1)) as \"due!\",\n count(*) as \"count!\"\n from repos\n where deleted_at is null\n group by 1, 2", + "describe": { + "columns": [ + { + "ordinal": 0, + "name": "state!", + "type_info": "Text", + "origin": "Expression" + }, + { + "ordinal": 1, + "name": "due!", + "type_info": "Bool", + "origin": "Expression" + }, + { + "ordinal": 2, + "name": "count!", + "type_info": "Int8", + "origin": "Expression" + } + ], + "parameters": { + "Left": [ + "Int8", + "Int4" + ] + }, + "nullable": [ + null, + null, + null + ] + }, + "hash": "6a8ebbb538b3254515d113c27135f20cb2a40f8a2f3c89b0a07115caf920c876" +} diff --git a/gitmirror/crates/gitmirror-db/src/lib.rs b/gitmirror/crates/gitmirror-db/src/lib.rs index f5743ecd9..6153c2add 100644 --- a/gitmirror/crates/gitmirror-db/src/lib.rs +++ b/gitmirror/crates/gitmirror-db/src/lib.rs @@ -192,6 +192,30 @@ impl Db { .collect()) } + /// Live repo counts by `(state, due)`, with `error` split into `backoff` / `error` at + /// `max_retry`. + pub async fn repo_state_counts( + &self, + now_us: i64, + max_retry: i32, + ) -> Result, DbError> { + let rows = sqlx::query!( + r#"select case when state = 'error' and retry_count >= $2 then 'error' + when state = 'error' then 'backoff' + else state end as "state!", + (state <> 'active' and (retry_after_us is null or retry_after_us <= $1)) as "due!", + count(*) as "count!" + from repos + where deleted_at is null + group by 1, 2"#, + now_us, + max_retry, + ) + .fetch_all(&self.0) + .await?; + Ok(rows.into_iter().map(|r| (r.state, r.due, r.count)).collect()) + } + pub async fn mark_synced(&self, did: &Did, refs_cid: &str) -> Result<(), DbError> { let done = sqlx::query!( "update repos diff --git a/gitmirror/crates/gitmirror-sync/Cargo.toml b/gitmirror/crates/gitmirror-sync/Cargo.toml index 0e4e05172..db02fd74f 100644 --- a/gitmirror/crates/gitmirror-sync/Cargo.toml +++ b/gitmirror/crates/gitmirror-sync/Cargo.toml @@ -20,6 +20,7 @@ gix-protocol = { workspace = true } gix-transport = { workspace = true, features = ["http-client-reqwest-rust-tls"] } jacquard-common = { workspace = true } jacquard-identity = { workspace = true } +metrics = { workspace = true } reqwest = { workspace = true } scc = { workspace = true } thiserror = { workspace = true } diff --git a/gitmirror/crates/gitmirror-sync/src/lib.rs b/gitmirror/crates/gitmirror-sync/src/lib.rs index 25e72c3e8..154ace68a 100644 --- a/gitmirror/crates/gitmirror-sync/src/lib.rs +++ b/gitmirror/crates/gitmirror-sync/src/lib.rs @@ -1,7 +1,7 @@ //! sync git repository objects (packfile) use std::sync::Arc; -use std::time::Duration; +use std::time::{Duration, Instant}; use bobbin_runtime::Clock; @@ -13,7 +13,8 @@ use jacquard_common::types::did::Did; use thiserror::Error; use tokio::sync::{Notify, Semaphore}; use tokio_util::sync::CancellationToken; -use tracing::{error, info, Instrument as _}; +use metrics::{gauge, histogram}; +use tracing::{error, info, warn, Instrument as _}; mod admission; mod task; @@ -38,6 +39,11 @@ const MAX_BACKOFF_SECS: i64 = 60 * 60; /// not a backoff: it only has to be long enough to get the row out of the scan window. const DEFER_HOLD_US: i64 = 5 * 1_000_000; +/// Metric-only split between `backoff` and `error`; retries continue past it. +const MAX_RETRY: i32 = 12; + +const STATE_METRICS_INTERVAL: Duration = Duration::from_secs(30); + pub const DEFAULT_CONCURRENCY: usize = 8; pub const DEFAULT_PER_KNOT_CONCURRENCY: usize = 4; pub const DEFAULT_POLL_INTERVAL: Duration = Duration::from_secs(5); @@ -102,6 +108,7 @@ struct InFlightGuard { impl Drop for InFlightGuard { fn drop(&mut self) { let _ = self.set.remove_sync(&self.did); + gauge!("gitmirror_sync_in_flight").decrement(1); } } @@ -149,6 +156,7 @@ impl BackfillWorker { break; } + gauge!("gitmirror_sync_in_flight").increment(1); let guard = InFlightGuard { did: row.did.clone(), set: self.in_flight.clone(), @@ -163,10 +171,15 @@ impl BackfillWorker { let admission = self.admission.clone(); let drained = self.drained.clone(); - let span = tracing::info_span!("backfill", did = %row.did); + let span = tracing::info_span!( + "backfill", + did = %row.did, + knot = tracing::field::Empty + ); tokio::spawn( async move { let _guard = guard; + let started = Instant::now(); let res = task::did_task( &http, &directory, @@ -178,6 +191,22 @@ impl BackfillWorker { ) .await; + let duration = started.elapsed().as_secs_f64(); + let outcome = match &res { + Ok(task::TaskDisposition::Finished) => "finished", + Ok(task::TaskDisposition::Deferred) => "deferred", + Err(_) => "failed", + }; + histogram!("gitmirror_sync_task_duration_seconds", "outcome" => outcome) + .record(duration); + info!( + outcome, + duration_ms = duration * 1e3, + retry_count = row.retry_count, + error = res.as_ref().err().map(tracing::field::display), + "sync finished" + ); + let write = match &res { Ok(task::TaskDisposition::Finished) => { db.mark_synced(&row.did, &row.refs_cid).await @@ -229,3 +258,29 @@ impl BackfillWorker { Ok(()) } } + +/// Publishes `gitmirror_sync_repos` from the DB until cancelled. +pub async fn report_repo_states(db: Db, clock: Arc, cancel: CancellationToken) { + loop { + match db.repo_state_counts(now_us(&*clock), MAX_RETRY).await { + Ok(rows) => { + for state in ["pending", "active", "backoff", "error"] { + for due in [true, false] { + let count = rows + .iter() + .find(|(s, d, _)| s == state && *d == due) + .map_or(0, |(_, _, c)| *c); + gauge!("gitmirror_sync_repos", "state" => state, "due" => due.to_string()) + .set(count as f64); + } + } + } + Err(e) => warn!(error = %e, "counting repo states failed"), + } + + tokio::select! { + () = cancel.cancelled() => break, + () = clock.sleep(STATE_METRICS_INTERVAL) => {} + } + } +} diff --git a/gitmirror/crates/gitmirror-sync/src/task.rs b/gitmirror/crates/gitmirror-sync/src/task.rs index 10e46df4b..8cdf28e8e 100644 --- a/gitmirror/crates/gitmirror-sync/src/task.rs +++ b/gitmirror/crates/gitmirror-sync/src/task.rs @@ -74,6 +74,7 @@ pub async fn did_task( let knot = knot_host(directory, &repo.did, policy).await?; let knot = knot.url(); + tracing::Span::current().record("knot", knot.as_str()); let Some(_knot_permit) = admission.try_acquire(knot).await else { return Ok(TaskDisposition::Deferred); diff --git a/gitmirror/crates/gitmirror-xrpc/src/metrics.rs b/gitmirror/crates/gitmirror-xrpc/src/metrics.rs index ae3643ef2..cb18c09e4 100644 --- a/gitmirror/crates/gitmirror-xrpc/src/metrics.rs +++ b/gitmirror/crates/gitmirror-xrpc/src/metrics.rs @@ -26,6 +26,10 @@ const BLOCKING_WAIT_BUCKETS: &[f64] = &[ 1.0, ]; +const SYNC_TASK_DURATION_BUCKETS: &[f64] = &[ + 0.05, 0.1, 0.25, 0.5, 1.0, 2.5, 5.0, 10.0, 30.0, 60.0, 120.0, 300.0, 600.0, 1800.0, +]; + pub fn init_metrics() -> PrometheusHandle { let builder = PrometheusBuilder::new() .set_buckets_for_metric( @@ -37,7 +41,12 @@ pub fn init_metrics() -> PrometheusHandle { Matcher::Full("gitmirror_blocking_queue_wait_seconds".to_owned()), BLOCKING_WAIT_BUCKETS, ) - .expect("blocking wait buckets should not be empty"); + .expect("blocking wait buckets should not be empty") + .set_buckets_for_metric( + Matcher::Full("gitmirror_sync_task_duration_seconds".to_owned()), + SYNC_TASK_DURATION_BUCKETS, + ) + .expect("sync task duration buckets should not be empty"); let handle = builder .install_recorder() .expect("failed to install Prometheus recorder"); @@ -73,6 +82,18 @@ fn describe_metrics() { "gitmirror_blocking_queue_wait_seconds", "Time a task spent waiting for a blocking-pool thread" ); + metrics::describe_gauge!( + "gitmirror_sync_repos", + "Repos by sync state and whether they are due for a sync" + ); + metrics::describe_gauge!( + "gitmirror_sync_in_flight", + "Sync tasks currently running" + ); + metrics::describe_histogram!( + "gitmirror_sync_task_duration_seconds", + "Sync task duration in seconds" + ); } #[derive(Clone, Copy)] diff --git a/gitmirror/crates/gitmirror/src/main.rs b/gitmirror/crates/gitmirror/src/main.rs index 00b7d2b15..bf74a70aa 100644 --- a/gitmirror/crates/gitmirror/src/main.rs +++ b/gitmirror/crates/gitmirror/src/main.rs @@ -231,7 +231,7 @@ async fn serve(cfg: MirrorConfig) -> anyhow::Result<()> { checkpoint_interval: Duration::from_secs(cfg.ingest.checkpoint_interval_secs), }; let ingest_runtime = IngestRuntime { - db, + db: db.clone(), clock: clock.clone(), cancel: cancel.clone(), }; @@ -270,6 +270,11 @@ async fn serve(cfg: MirrorConfig) -> anyhow::Result<()> { tracing::info!("metrics endpoint disabled"); } else { gitmirror_xrpc::metrics::init_metrics(); + let (clock, cancel) = (clock.clone(), cancel.clone()); + tasks.spawn(async move { + gitmirror_sync::report_repo_states(db, clock, cancel).await; + ("sync-metrics", Ok(())) + }); tracing::info!(binds = %bind_display(&metrics_binds), "gitmirror metrics listening"); } diff --git a/nix/Cargo.nix b/nix/Cargo.nix index fab810ca3..0b1b5bcf0 100644 --- a/nix/Cargo.nix +++ b/nix/Cargo.nix @@ -10702,6 +10702,10 @@ rec { packageId = "jacquard-identity"; features = [ "cache" ]; } + { + name = "metrics"; + packageId = "metrics"; + } { name = "reqwest"; packageId = "reqwest 0.13.5";