From 998b04643495fb013bed51e5e379a6ee0ec8ac9e Mon Sep 17 00:00:00 2001 From: Harshil Goel Date: Tue, 29 Sep 2026 13:32:15 +0530 Subject: [PATCH] Add some new metrics for errors and per table stuff --- src/bin/stream/bootstrap.rs | 1 + src/bin/stream/metrics_publish.rs | 70 ++++++++++++++ src/bin/stream/session.rs | 7 ++ src/bin/stream/tracing_setup.rs | 1 + src/emit/ch_emitter.rs | 10 ++ src/emit/pipeline/batcher.rs | 39 +++++--- src/emit/pipeline/inserter.rs | 8 ++ src/emit/pipeline/mod.rs | 1 + src/emit/pipeline/reorder.rs | 55 ++++++----- src/emit/pipeline/row_ledger.rs | 72 ++++++++++++++ src/lib.rs | 3 +- src/ops/control.rs | 12 +++ src/ops/log_events.rs | 151 ++++++++++++++++++++++++++++++ src/ops/metrics.rs | 118 +++++++++++++++++++++++ src/ops/mod.rs | 1 + src/schema.rs | 4 +- tests/control_plane_e2e.rs | 58 +++++++++++- 17 files changed, 568 insertions(+), 43 deletions(-) create mode 100644 src/emit/pipeline/row_ledger.rs create mode 100644 src/ops/log_events.rs diff --git a/src/bin/stream/bootstrap.rs b/src/bin/stream/bootstrap.rs index 6b6fe7e5..79f5f06e 100644 --- a/src/bin/stream/bootstrap.rs +++ b/src/bin/stream/bootstrap.rs @@ -428,6 +428,7 @@ pub(crate) async fn run_bootstrap( by_database: vec![boot_db.series()], ..stage_gauges(&StageCounters { emitter: Some(&stats), + db_name: &|_| boot_db.database.clone(), oracle: [None, Some(&oracle_stats)], bootstrap: Some(&progress), bootstrap_attempt: attempt, diff --git a/src/bin/stream/metrics_publish.rs b/src/bin/stream/metrics_publish.rs index 7ba5c8ad..cd718ec7 100644 --- a/src/bin/stream/metrics_publish.rs +++ b/src/bin/stream/metrics_publish.rs @@ -10,6 +10,7 @@ use walshadow::boundary_hold::BoundaryHoldStats; use walshadow::ch_emitter::EmitterStats; use walshadow::config::ConfigResolver; use walshadow::metrics::{DbSeries, MetricsRegistry, MetricsSnapshot}; +use walshadow::pipeline::row_ledger::TableRowCounts; use walshadow::pos::{ Drain, EmitterAck, FilterDispatched, Floor, Pos, ShadowReplay, SourceReceived, }; @@ -416,6 +417,8 @@ pub(crate) async fn populate_pipeline_metrics( /// counters sum into one series pub(crate) struct StageCounters<'a> { pub(crate) emitter: Option<&'a walshadow::ch_emitter::EmitterStats>, + /// `database=` label for a source database oid + pub(crate) db_name: &'a (dyn Fn(u32) -> String + Sync), /// Live first, bootstrap's second. Only one of the pair is ever serving pub(crate) oracle: [Option<&'a walshadow::oracle::OracleStats>; 2], pub(crate) bootstrap: Option<&'a BootstrapProgress>, @@ -432,12 +435,26 @@ pub(crate) fn emitter_counts( picks.map(|pick| stats.map_or(0, |s| pick(s).load(Ordering::Relaxed))) } +/// Stop conditions beyond the crossing and shadow views pub(crate) fn stage_gauges(v: &StageCounters<'_>) -> MetricsSnapshot { stage_gauges_on(v, MetricsSnapshot::default()) } +const UNMAPPED_STATUS_ROWS: usize = 20; + pub(crate) fn stage_gauges_on(v: &StageCounters<'_>, base: MetricsSnapshot) -> MetricsSnapshot { let (proc_cpu, proc_rss, proc_threads) = read_process_stats(); + let by_table = |pick: fn(&EmitterStats) -> &TableRowCounts| { + v.emitter.map_or_else(Vec::new, |s| { + pick(s) + .snapshot() + .into_iter() + .map(|(key, rows)| ((v.db_name)(key.db_oid), key.rel.to_string(), rows)) + .collect::>() + }) + }; + let mut unmapped_rows_by_table = by_table(|s| &s.unmapped_rows_by_table); + unmapped_rows_by_table.truncate(UNMAPPED_STATUS_ROWS); let emitter = |pick: fn(&EmitterStats) -> &AtomicU64| -> u64 { v.emitter.map_or(0, |s| pick(s).load(Ordering::Relaxed)) }; @@ -538,6 +555,17 @@ pub(crate) fn stage_gauges_on(v: &StageCounters<'_>, base: MetricsSnapshot) -> M emitter_xacts_total: emitter(|s| &s.xacts_committed), emitter_unsupported_relations: emitter(|s| &s.unsupported_relations), emitter_deletes_discarded: emitter(|s| &s.deletes_discarded), + unmapped_rows_by_table, + rows_inserted_by_table: by_table(|s| &s.rows_inserted_by_table), + cells_inserted_by_type: v.emitter.map_or_else(Vec::new, |s| { + s.cells_inserted_by_type + .snapshot() + .into_iter() + .map(|((ty, enc), n)| (ty, enc, n)) + .collect() + }), + emitter_retries_attempted_total: emitter(|s| &s.retries_attempted), + emitter_reconnects_total: emitter(|s| &s.reconnects), oracle_local_columns_total: emitter(|s| &s.oracle_local_columns), oracle_blocks_total: oracle(|s| &s.blocks), oracle_rows_total: oracle(|s| &s.rows), @@ -643,6 +671,7 @@ mod tests { boundary_hold: &boundary, by_database: Vec::new(), counters: StageCounters { + db_name: &|_| String::new(), emitter: Some(&emitter), oracle: [None, None], bootstrap: None, @@ -667,11 +696,50 @@ mod tests { } } + /// The publisher names each counted table by resolving its `db_oid`, so a + /// ledger with rows has to reach the snapshot under the database's name + #[test] + fn per_table_counters_resolve_their_database_name() { + use walshadow::schema::{RelName, TableKey}; + let emitter = EmitterStats::default(); + let orders = TableKey::new(5, RelName::new("public", "orders")); + let noise = TableKey::new(9, RelName::new("public", "noise")); + emitter + .rows_inserted_by_table + .counter(&orders) + .fetch_add(7, Ordering::Relaxed); + emitter + .unmapped_rows_by_table + .counter(&noise) + .fetch_add(3, Ordering::Relaxed); + let snap = stage_gauges(&StageCounters { + db_name: &|oid| match oid { + 5 => "app".to_owned(), + _ => format!("db{oid}"), + }, + emitter: Some(&emitter), + oracle: [None, None], + bootstrap: None, + bootstrap_attempt: 0, + uptime_secs: 0, + }); + assert_eq!( + snap.rows_inserted_by_table, + vec![("app".to_owned(), "public.orders".to_owned(), 7)], + ); + assert_eq!( + snap.unmapped_rows_by_table, + vec![("db9".to_owned(), "public.noise".to_owned(), 3)], + "an oid the config does not name still has to be reportable", + ); + } + #[test] fn archive_stage_metrics_refresh_without_resetting_recovery_state() { use walshadow::record::WAL_SEG_SIZE; let emitter = EmitterStats::default(); let counters = StageCounters { + db_name: &|_| String::new(), emitter: Some(&emitter), oracle: [None, None], bootstrap: None, @@ -726,6 +794,7 @@ mod tests { boot_oracle.rows.fetch_add(7, Ordering::Relaxed); let snap = stage_gauges(&StageCounters { + db_name: &|_| String::new(), emitter: None, oracle: [Some(&live_oracle), Some(&boot_oracle)], bootstrap: None, @@ -737,6 +806,7 @@ mod tests { assert_eq!(snap.bootstrap_attempt, 2); let boot_only = stage_gauges(&StageCounters { + db_name: &|_| String::new(), emitter: None, oracle: [None, Some(&boot_oracle)], bootstrap: None, diff --git a/src/bin/stream/session.rs b/src/bin/stream/session.rs index 911edd8e..4df06958 100644 --- a/src/bin/stream/session.rs +++ b/src/bin/stream/session.rs @@ -1895,6 +1895,13 @@ pub(crate) async fn run_session( metrics_dbs.iter().map(DbMetricSources::series).collect(), StageCounters { emitter: emitter_stats, + db_name: &|oid| { + db_conns + .iter() + .find(|c| c.oid == oid) + .map(|c| c.name.clone()) + .unwrap_or_default() + }, oracle: [oracle_stats, bootstrap_metrics.as_ref().map(|b| &*b.oracle)], bootstrap: bootstrap_metrics.as_ref().map(|b| &b.progress), bootstrap_attempt: 0, diff --git a/src/bin/stream/tracing_setup.rs b/src/bin/stream/tracing_setup.rs index 6c4bb169..dcf4b9d9 100644 --- a/src/bin/stream/tracing_setup.rs +++ b/src/bin/stream/tracing_setup.rs @@ -81,6 +81,7 @@ pub(crate) fn init_tracing( .with(filter) .with(fmt_layer) .with(otel_layer) + .with(walshadow::log_events::LogEventLayer) .try_init(); provider } diff --git a/src/emit/ch_emitter.rs b/src/emit/ch_emitter.rs index 20828f88..325bd4b4 100644 --- a/src/emit/ch_emitter.rs +++ b/src/emit/ch_emitter.rs @@ -1394,6 +1394,13 @@ pub(crate) enum ColumnEncoding { } impl ColumnEncoding { + pub(crate) fn label(&self) -> &'static str { + match self { + Self::Local => "local", + Self::Oracle { .. } => "oracle", + } + } + fn choose(att: Option<&RelAttr>, ast: &TypeAst) -> Self { let Some(att) = att else { // Predeclared post-ALTER column has only default cells @@ -2551,6 +2558,9 @@ crate::atomic_stats! { /// (memoised) pub route_snapshots_mapped, pub route_snapshots_unmapped, + pub unmapped_rows_by_table: crate::emit::pipeline::row_ledger::TableRowCounts, + pub rows_inserted_by_table: crate::emit::pipeline::row_ledger::TableRowCounts, + pub cells_inserted_by_type: crate::emit::pipeline::row_ledger::CellCounts, /// Commit-resolve raw decode: records by verdict kind, per op /// (`raw_decode_records_total{kind,op}`) pub raw_decode_toast_ops: crate::decode::heap_decoder::OpCounters, diff --git a/src/emit/pipeline/batcher.rs b/src/emit/pipeline/batcher.rs index 6a469cb0..4f29c1c4 100644 --- a/src/emit/pipeline/batcher.rs +++ b/src/emit/pipeline/batcher.rs @@ -20,7 +20,7 @@ use std::collections::hash_map::Entry; use std::sync::Arc; -use std::sync::atomic::{AtomicUsize, Ordering}; +use std::sync::atomic::{AtomicU64, AtomicUsize, Ordering}; use std::time::Duration; use clickhouse_c::Allocator; @@ -32,7 +32,8 @@ use tokio::time::Instant; use crate::config::ResolvedConfig; use crate::decode::heap_decoder::{CommittedTuple, HeapOp}; use crate::emit::ch_emitter::{ - Append, ColumnBuf, EmitterStats, OP_DELETE, OP_INSERT, OP_UPDATE, TableEncoder, TablePlan, + Append, ColumnBuf, ColumnEncoding, EmitterStats, OP_DELETE, OP_INSERT, OP_UPDATE, TableEncoder, + TablePlan, }; use crate::emit::pipeline::{DEFAULT_PIPELINE_FLUSH, Fatal}; use crate::emit::route::RouteSnapshot; @@ -67,6 +68,7 @@ pub struct RowChunk { pub struct ColMeta { pub name: String, pub type_repr: String, + pub cells_inserted: Arc, } /// Immutable per-table block shape, shared by every batch of that table @@ -78,16 +80,26 @@ pub struct BatchMeta { /// synthetic ones (lsn, xid, commit_ts, delete marker when configured). pub columns: Vec, pub schema_epoch: u64, + pub rows_inserted: Arc, } impl BatchMeta { - fn from_plan(plan: &TablePlan, table_key: TableKey, schema_epoch: u64) -> Self { + fn from_plan( + plan: &TablePlan, + table_key: TableKey, + schema_epoch: u64, + stats: &EmitterStats, + ) -> Self { + let col_meta = |name: &String, type_repr: &String, encoding: &ColumnEncoding| ColMeta { + name: name.clone(), + type_repr: type_repr.clone(), + cells_inserted: stats + .cells_inserted_by_type + .counter(&(type_repr.clone(), encoding.label())), + }; let mut columns = Vec::with_capacity(plan.columns.len() + 4); for c in &plan.columns { - columns.push(ColMeta { - name: c.name.clone(), - type_repr: c.type_repr.clone(), - }); + columns.push(col_meta(&c.name, &c.type_repr, &c.encoding)); } for synth in [ Some(&plan.synth_lsn), @@ -98,12 +110,10 @@ impl BatchMeta { .into_iter() .flatten() { - columns.push(ColMeta { - name: synth.name.clone(), - type_repr: synth.type_repr.clone(), - }); + columns.push(col_meta(&synth.name, &synth.type_repr, &synth.encoding)); } Self { + rows_inserted: stats.rows_inserted_by_table.counter(&table_key), table_key, insert_sql: plan.insert_sql.clone(), columns, @@ -333,7 +343,12 @@ async fn handle_row( row.route.system_columns(), ) .map_err(|e| e.to_string())?; - let meta = Arc::new(BatchMeta::from_plan(&plan, e.key().clone(), ctx.epoch)); + let meta = Arc::new(BatchMeta::from_plan( + &plan, + e.key().clone(), + ctx.epoch, + ctx.stats, + )); let enc = TableEncoder::new(plan).map_err(|e| e.to_string())?; e.insert(Table { enc, diff --git a/src/emit/pipeline/inserter.rs b/src/emit/pipeline/inserter.rs index f0c17ac0..59f399e2 100644 --- a/src/emit/pipeline/inserter.rs +++ b/src/emit/pipeline/inserter.rs @@ -172,6 +172,14 @@ impl Inserter { self.stats .rows_emitted .fetch_add(batch.n_rows as u64, Ordering::Relaxed); + batch + .meta + .rows_inserted + .fetch_add(batch.n_rows as u64, Ordering::Relaxed); + for col in &batch.meta.columns { + col.cells_inserted + .fetch_add(batch.n_rows as u64, Ordering::Relaxed); + } self.stats.blocks_sent.fetch_add(1, Ordering::Relaxed); self.stats .inserter_batches_in diff --git a/src/emit/pipeline/mod.rs b/src/emit/pipeline/mod.rs index 930090da..e6915bf4 100644 --- a/src/emit/pipeline/mod.rs +++ b/src/emit/pipeline/mod.rs @@ -19,6 +19,7 @@ pub mod plan_spool; pub mod planner; pub mod reorder; pub mod resolver; +pub mod row_ledger; pub mod tail; use std::sync::Arc; diff --git a/src/emit/pipeline/reorder.rs b/src/emit/pipeline/reorder.rs index d5e5d006..397a3490 100644 --- a/src/emit/pipeline/reorder.rs +++ b/src/emit/pipeline/reorder.rs @@ -30,7 +30,7 @@ use crate::decode::visibility::{PgXactPatch, PgXactView, read_pg_xact}; use crate::emit::ch_ddl::DdlApplicator; use crate::emit::ch_emitter::EmitterStats; use crate::record::{Record, RecordSink, SinkError}; -use crate::schema::{RelDescriptor, RelName, SchemaEvent}; +use crate::schema::{RelDescriptor, RelName, SchemaEvent, TableKey}; use ahash::{HashMap, HashMapExt, HashSet, HashSetExt}; use tracing::Instrument; @@ -1161,31 +1161,40 @@ impl<'a> ReorderRouteView<'a> { impl PlanRouteView for ReorderRouteView<'_> { fn route_for(&mut self, heap: &DescribedHeap) -> Option> { let rel_name = &heap.descriptor.rel_name; - if let Some(r) = self.memo.get(rel_name) { - return r.clone(); - } - let mapped = match self.overlay.get(rel_name) { - Some(o) => o.clone(), - None => self.mapping.as_ref().and_then(|m| m.get(rel_name)).cloned(), + let route = if let Some(r) = self.memo.get(rel_name) { + r.clone() + } else { + let mapped = match self.overlay.get(rel_name) { + Some(o) => o.clone(), + None => self.mapping.as_ref().and_then(|m| m.get(rel_name)).cloned(), + }; + let route = mapped.map(|m| { + let rules = self + .config + .as_ref() + .map_or_else(Arc::default, |rc| rc.column_rules.clone()); + let policy = self.row_policy.for_rel(self.config.as_deref(), rel_name); + RouteSnapshot::freeze(Arc::new(m), rules, policy) + }); + let result = if route.is_none() { + self.stats + .unsupported_relations + .fetch_add(1, Ordering::Relaxed); + &self.stats.route_snapshots_unmapped + } else { + &self.stats.route_snapshots_mapped + }; + result.fetch_add(1, Ordering::Relaxed); + self.memo.insert(rel_name.clone(), route.clone()); + route }; - let route = mapped.map(|m| { - let rules = self - .config - .as_ref() - .map_or_else(Arc::default, |rc| rc.column_rules.clone()); - let policy = self.row_policy.for_rel(self.config.as_deref(), rel_name); - RouteSnapshot::freeze(Arc::new(m), rules, policy) - }); - let result = if route.is_none() { + if route.is_none() { + let key = TableKey::new(heap.descriptor.rfn.db_node, rel_name.clone()); self.stats - .unsupported_relations + .unmapped_rows_by_table + .counter(&key) .fetch_add(1, Ordering::Relaxed); - &self.stats.route_snapshots_unmapped - } else { - &self.stats.route_snapshots_mapped - }; - result.fetch_add(1, Ordering::Relaxed); - self.memo.insert(rel_name.clone(), route.clone()); + } route } diff --git a/src/emit/pipeline/row_ledger.rs b/src/emit/pipeline/row_ledger.rs new file mode 100644 index 00000000..c4b4d6de --- /dev/null +++ b/src/emit/pipeline/row_ledger.rs @@ -0,0 +1,72 @@ +use std::hash::Hash; +use std::sync::atomic::{AtomicU64, Ordering}; +use std::sync::{Arc, RwLock}; + +use ahash::HashMap; + +use crate::schema::TableKey; + +pub type TableRowCounts = KeyedCounts; +/// Keyed by ClickHouse type and `ColumnEncoding::label` +pub type CellCounts = KeyedCounts<(String, &'static str)>; + +/// Counter per label set. Hot paths hold the returned handle, so a series +/// exists at zero before its first increment +#[derive(Debug)] +pub struct KeyedCounts { + counts: RwLock>>, +} + +impl Default for KeyedCounts { + fn default() -> Self { + Self { + counts: RwLock::default(), + } + } +} + +impl KeyedCounts { + pub fn counter(&self, key: &K) -> Arc { + if let Some(n) = self.counts.read().expect("keyed counts lock").get(key) { + return n.clone(); + } + self.counts + .write() + .expect("keyed counts lock") + .entry(key.clone()) + .or_default() + .clone() + } + + /// Heaviest first + pub fn snapshot(&self) -> Vec<(K, u64)> { + let mut rows: Vec<_> = self + .counts + .read() + .expect("keyed counts lock") + .iter() + .map(|(k, n)| (k.clone(), n.load(Ordering::Relaxed))) + .collect(); + rows.sort_unstable_by(|l, r| r.1.cmp(&l.1).then_with(|| l.0.cmp(&r.0))); + rows + } +} + +#[cfg(test)] +mod tests { + use super::*; + use crate::schema::RelName; + + #[test] + fn counts_accumulate_per_key_and_rank_by_weight() { + let counts = TableRowCounts::default(); + let a = TableKey::new(1, RelName::new("public", "a")); + let b = TableKey::new(1, RelName::new("public", "b")); + let other_db = TableKey::new(2, RelName::new("public", "a")); + counts.counter(&b); + counts.counter(&a).fetch_add(3, Ordering::Relaxed); + counts.counter(&a).fetch_add(4, Ordering::Relaxed); + counts.counter(&other_db).fetch_add(9, Ordering::Relaxed); + assert_eq!(counts.snapshot(), vec![(other_db, 9), (a, 7), (b, 0)]); + } +} diff --git a/src/lib.rs b/src/lib.rs index 7287af10..c1226052 100644 --- a/src/lib.rs +++ b/src/lib.rs @@ -71,7 +71,8 @@ pub use emit::{ch_ddl, ch_emitter, pipeline}; pub use filter::{catalog_tracker, classify, main_data, pg_class_decoder, rewrite}; #[doc(hidden)] pub use ops::{ - bridge, control, ctl, init, introspect, metrics, oracle, preflight, retention, trace, + bridge, control, ctl, init, introspect, log_events, metrics, oracle, preflight, retention, + trace, }; #[doc(hidden)] pub use source::{ diff --git a/src/ops/control.rs b/src/ops/control.rs index d8ed26d9..8ab63860 100644 --- a/src/ops/control.rs +++ b/src/ops/control.rs @@ -503,6 +503,18 @@ async fn stream_status(ctx: &SharedCtx) -> Result { "backfills_pending".into(), (snap.config_backfills_pending as i64).into(), ); + let unmapped: Vec = snap + .unmapped_rows_by_table + .iter() + .map(|(database, table, rows)| { + let mut row = Table::new(); + row.insert("database".into(), database.clone().into()); + row.insert("table".into(), table.clone().into()); + row.insert("rows".into(), (*rows as i64).into()); + Value::Table(row) + }) + .collect(); + out.insert("unmapped_tables".into(), Value::Array(unmapped)); out.insert( "lag_bytes".into(), (snap.shadow_apply_lag_bytes as i64).into(), diff --git a/src/ops/log_events.rs b/src/ops/log_events.rs new file mode 100644 index 00000000..8f07aea5 --- /dev/null +++ b/src/ops/log_events.rs @@ -0,0 +1,151 @@ +use std::fmt; +use std::sync::atomic::{AtomicU64, Ordering}; +use std::sync::{LazyLock, PoisonError, RwLock}; + +use ahash::{HashMap, HashMapExt}; +use prometheus_client::collector::Collector; +use prometheus_client::encoding::DescriptorEncoder; + +use crate::ops::metrics::{counter_pairs, counter_series}; +use tracing::{Event, Level, Subscriber}; +use tracing_subscriber::layer::{Context, Layer}; + +type Key = (&'static str, &'static str); + +static EVENTS: LazyLock>> = + LazyLock::new(|| RwLock::new(HashMap::new())); + +const LEVELS: [Level; 5] = [ + Level::ERROR, + Level::WARN, + Level::INFO, + Level::DEBUG, + Level::TRACE, +]; + +fn record(key: Key) { + if let Some(n) = EVENTS + .read() + .unwrap_or_else(PoisonError::into_inner) + .get(&key) + { + n.fetch_add(1, Ordering::Relaxed); + return; + } + EVENTS + .write() + .unwrap_or_else(PoisonError::into_inner) + .entry(key) + .or_default() + .fetch_add(1, Ordering::Relaxed); +} + +#[derive(Debug, Default, Clone, Copy)] +pub struct LogEventLayer; + +impl Layer for LogEventLayer { + fn on_event(&self, event: &Event<'_>, _ctx: Context<'_, S>) { + let meta = event.metadata(); + record((meta.level().as_str(), meta.target())); + } +} + +#[derive(Debug)] +pub struct LogEventCollector; + +impl Collector for LogEventCollector { + fn encode(&self, mut enc: DescriptorEncoder) -> fmt::Result { + let map = EVENTS.read().unwrap_or_else(PoisonError::into_inner); + let mut rows: Vec<_> = map + .iter() + .map(|(&(level, target), n)| (level, target, n.load(Ordering::Relaxed))) + .collect(); + rows.sort_unstable(); + + counter_series( + &mut enc, + "walshadow_log_level_events", + "Log events emitted by this process, summed over targets. Always present, so a first event is an increase rather than a new series.", + "level", + LEVELS.map(|l| l.as_str()).into_iter().map(|level| { + let n: u64 = rows + .iter() + .filter(|(l, _, _)| *l == level) + .map(|(_, _, n)| n) + .sum(); + (level, n) + }), + )?; + counter_pairs( + &mut enc, + "walshadow_log_events", + "Log events emitted by this process, labeled by level + target.", + ["level", "target"], + rows.iter() + .map(|(level, target, n)| ([*level, *target], *n)), + ) + } +} + +#[cfg(test)] +mod tests { + use tracing_subscriber::prelude::*; + + use super::*; + use crate::ops::metrics::{MetricsSnapshot, render}; + + const TARGET: &str = "walshadow::log_events_test"; + + fn count(level: &str, target: &str) -> u64 { + EVENTS + .read() + .expect("counter map") + .iter() + .find(|((l, t), _)| *l == level && *t == target) + .map_or(0, |(_, n)| n.load(Ordering::Relaxed)) + } + + #[test] + fn every_level_total_is_present_before_any_event() { + let out = render(MetricsSnapshot::default()); + for level in LEVELS.map(|l| l.as_str()) { + assert!( + out.contains(&format!( + "walshadow_log_level_events_total{{level=\"{level}\"}}" + )), + "{level} missing from {out}" + ); + } + } + + #[test] + fn counts_by_level_and_target_then_renders_them() { + tracing::subscriber::with_default( + tracing_subscriber::registry().with(LogEventLayer), + || { + tracing::warn!(target: TARGET, "one"); + tracing::info!(target: TARGET, "two"); + tracing::info!(target: TARGET, "three"); + }, + ); + assert_eq!(count("WARN", TARGET), 1); + assert_eq!(count("INFO", TARGET), 2); + assert_eq!(count("ERROR", TARGET), 0); + + let out = render(MetricsSnapshot::default()); + assert!( + out.contains(&format!( + "walshadow_log_events_total{{level=\"INFO\",target=\"{TARGET}\"}} 2" + )), + "{out}" + ); + let warns = count("WARN", TARGET); + assert!( + out.lines().any(|l| l + .strip_prefix("walshadow_log_level_events_total{level=\"WARN\"} ") + .and_then(|n| n.trim().parse::().ok()) + .is_some_and(|n| n >= warns)), + "per-level total must cover this target's warns: {out}" + ); + } +} diff --git a/src/ops/metrics.rs b/src/ops/metrics.rs index 184c40bf..7123668f 100644 --- a/src/ops/metrics.rs +++ b/src/ops/metrics.rs @@ -183,6 +183,9 @@ snapshot! { /// across databases. `ctl status` reads it; the scrape renders the /// per-database split instead config_backfills_pending: u64, + unmapped_rows_by_table: Vec<(String, String, u64)>, + rows_inserted_by_table: Vec<(String, String, u64)>, + cells_inserted_by_type: Vec<(String, &'static str, u64)>, /// Refusal a crossing parked on, rendered as `crossing_wedged`. The /// pump keeps publishing rather than exiting into a restart that /// re-crosses and re-fails, so this is how the daemon says it is @@ -363,6 +366,10 @@ snapshot! { "Tuples skipped because the source relation has no mapping in --ch-config.", counter emitter_deletes_discarded: u64 = "DELETE rows dropped because is_deleted = false leaves no marker column.", + counter emitter_retries_attempted_total: u64 = + "Failed ClickHouse operations retried, one per failing operation rather than per attempt.", + counter emitter_reconnects_total: u64 = + "ClickHouse client dials after the first, from retries and live endpoint or compression changes.", gauge pump_queue_depth: u64 = "Records buffered between the WAL pump and the queueing worker.", counter queue_records_out_total: u64 = "Records the queueing/reorder worker has dequeued and dispatched. rate() is the worker's throughput; with pump_queue_depth it tells deep-and-draining from deep-and-stalled.", @@ -766,6 +773,22 @@ pub(crate) fn counter_series<'a, V: EncodeCounterValue>( Ok(()) } +pub(crate) fn counter_pairs<'a, V: EncodeCounterValue>( + enc: &mut DescriptorEncoder<'_>, + name: &str, + help: &str, + keys: [&str; 2], + series: impl IntoIterator, +) -> fmt::Result { + let mut family = declare(enc, name, help, MetricType::Counter)?; + for (values, value) in series { + family + .encode_family(&[(keys[0], values[0]), (keys[1], values[1])])? + .encode_counter::(&value, None)?; + } + Ok(()) +} + pub(crate) fn gauge_series<'a, V: EncodeGaugeValue>( enc: &mut DescriptorEncoder<'_>, name: &str, @@ -886,6 +909,26 @@ impl Collector for SnapshotCollector { fn encode_snapshot(snap: &MetricsSnapshot, enc: &mut DescriptorEncoder<'_>) -> fmt::Result { encode_fields(snap, enc)?; + counter_pairs( + enc, + "walshadow_cells_inserted", + "Cells inserted per ClickHouse column type, split by whether the shadow oracle encoded them.", + ["type", "encoding"], + snap.cells_inserted_by_type + .iter() + .map(|(ty, encoding, n)| ([ty.as_str(), *encoding], *n)), + )?; + + counter_pairs( + enc, + "walshadow_table_rows_inserted", + "Rows inserted per source table, initial load and CDC alike.", + ["database", "table"], + snap.rows_inserted_by_table + .iter() + .map(|(database, table, n)| ([database.as_str(), table.as_str()], *n)), + )?; + gauge( enc, "walshadow_crossing_wedged", @@ -940,6 +983,7 @@ pub fn render(snap: MetricsSnapshot) -> String { let mut registry = Registry::default(); registry.register_collector(Box::new(SnapshotCollector(snap))); registry.register_collector(Box::new(crate::ops::stages::StageCollector)); + registry.register_collector(Box::new(crate::ops::log_events::LogEventCollector)); let mut out = String::with_capacity(16 << 10); text::encode(&mut out, ®istry).expect("String write cannot fail"); out @@ -1015,6 +1059,59 @@ async fn handle_client( mod tests { use super::*; + /// Every family encodes through `?`, so a writer that gives out partway + /// has to surface as an error rather than a panic or a truncated scrape + #[test] + fn encode_propagates_a_failing_writer_at_every_family() { + struct FailAfter { + writes: usize, + out: String, + } + impl fmt::Write for FailAfter { + fn write_str(&mut self, s: &str) -> fmt::Result { + if self.writes == 0 { + return Err(fmt::Error); + } + self.writes -= 1; + self.out.push_str(s); + Ok(()) + } + } + + let mut snap = MetricsSnapshot { + crossing_blocked_on: "slot_missing", + source_endpoint_swap_blocked_on: "slot_missing", + promotion_blocked_on: "not_paused", + unmapped_rows_by_table: vec![("app".into(), "public.noise".into(), 5)], + rows_inserted_by_table: vec![("app".into(), "public.orders".into(), 9)], + cells_inserted_by_type: vec![("String".into(), "local", 9)], + by_database: vec![DbSeries { + database: "app".into(), + ..DbSeries::default() + }], + ..MetricsSnapshot::default() + }; + snap.records_by_rm_route + .insert(("Heap".into(), "to_decoder"), 3); + + let full = render(snap.clone()).len(); + let mut errors = 0; + for writes in 0..64 { + let mut w = FailAfter { + writes, + out: String::new(), + }; + let mut registry = Registry::default(); + registry.register_collector(Box::new(SnapshotCollector(snap.clone()))); + registry.register_collector(Box::new(crate::ops::log_events::LogEventCollector)); + if text::encode(&mut w, ®istry).is_err() { + errors += 1; + assert!(w.out.len() < full, "a failed encode cannot be complete"); + } + } + assert_eq!(errors, 64, "every prefix length must surface the failure"); + } + #[test] fn render_includes_help_type_lines() { let mut snap = MetricsSnapshot { @@ -1039,6 +1136,27 @@ mod tests { assert!(body.contains("walshadow_uptime_seconds_total 42")); } + #[test] + fn render_labels_per_table_and_cell_series() { + let body = render(MetricsSnapshot::default()); + assert!( + !body.contains("walshadow_table_rows_inserted_total{"), + "{body}" + ); + + let body = render(MetricsSnapshot { + cells_inserted_by_type: vec![("Decimal(38, 9)".into(), "oracle", 7)], + rows_inserted_by_table: vec![("app".into(), "public.orders".into(), 4096)], + ..MetricsSnapshot::default() + }); + for want in [ + "walshadow_cells_inserted_total{type=\"Decimal(38, 9)\",encoding=\"oracle\"} 7", + "walshadow_table_rows_inserted_total{database=\"app\",table=\"public.orders\"} 4096", + ] { + assert!(body.contains(want), "{want} missing from {body}"); + } + } + #[test] fn render_emits_plan_families_labelled() { let snap = MetricsSnapshot { diff --git a/src/ops/mod.rs b/src/ops/mod.rs index 62f919af..c750e531 100644 --- a/src/ops/mod.rs +++ b/src/ops/mod.rs @@ -3,6 +3,7 @@ pub mod control; pub mod ctl; pub mod init; pub mod introspect; +pub mod log_events; pub mod metrics; pub mod oracle; pub mod preflight; diff --git a/src/schema.rs b/src/schema.rs index 2f6bf334..a8fa3f68 100644 --- a/src/schema.rs +++ b/src/schema.rs @@ -37,7 +37,7 @@ pub const NUMERICOID: u32 = 1700; pub const UUIDOID: u32 = 2950; pub const JSONBOID: u32 = 3802; -#[derive(Debug, Clone, PartialEq, Eq, Hash)] +#[derive(Debug, Clone, PartialEq, Eq, Hash, PartialOrd, Ord)] pub struct RelName { pub namespace: Arc, pub name: Arc, @@ -60,7 +60,7 @@ impl std::fmt::Display for RelName { /// Batch identity down the insert path: the same relation name in two source /// databases routes to two destinations, so the database is part of the key -#[derive(Debug, Clone, PartialEq, Eq, Hash)] +#[derive(Debug, Clone, PartialEq, Eq, Hash, PartialOrd, Ord)] pub struct TableKey { pub db_oid: u32, pub rel: RelName, diff --git a/tests/control_plane_e2e.rs b/tests/control_plane_e2e.rs index 7e4900f8..d673f463 100644 --- a/tests/control_plane_e2e.rs +++ b/tests/control_plane_e2e.rs @@ -930,6 +930,50 @@ async fn ch_password_rotation_via_ctl_keeps_streaming() { } } +#[tokio::test(flavor = "multi_thread", worker_threads = 4)] +async fn status_names_tables_whose_wal_is_not_needed() { + if !gated() { + return; + } + let mut h = Harness::up(&fx::Ports::alloc()) + .await + .expect("bring up harness"); + h.wait_ready(Duration::from_secs(60)) + .await + .expect("pump never started"); + + let result = async { + h.psql("CREATE TABLE demo.noise (id bigint PRIMARY KEY)")?; + h.psql("INSERT INTO demo.noise SELECT generate_series(1, 5)")?; + let deadline = Instant::now() + Duration::from_secs(30); + loop { + let status: toml::Table = h.ctl(&["status"])?.parse()?; + let unmapped = status + .get("unmapped_tables") + .and_then(toml::Value::as_array) + .context("no unmapped_tables in status")?; + let noise_rows = unmapped.iter().find_map(|row| { + (row.get("table")?.as_str()? == "demo.noise") + .then(|| row.get("rows")?.as_integer())? + }); + if noise_rows.is_some_and(|n| n >= 5) { + return Ok::<(), anyhow::Error>(()); + } + ensure!( + Instant::now() < deadline, + "status never counted demo.noise inserts, said {unmapped:?}", + ); + tokio::time::sleep(Duration::from_millis(250)).await; + } + } + .await; + + let stderr = h.teardown(); + if let Err(e) = result { + panic!("{e:#}\n--- daemon stderr ---\n{stderr}"); + } +} + /// Regression: applying one table used to opt pinned tables out #[tokio::test(flavor = "multi_thread", worker_threads = 4)] async fn apply_preserves_previously_pinned_table() { @@ -2389,11 +2433,15 @@ async fn slot_proofs_name_the_missing_slot_and_park_the_crossing() { // decision back and gives it again h.ctl_body(&["apply"], "[source]\nslot = \"sw\"")?; h.ctl_body(&["apply"], "[stream]\npaused = true")?; - tokio::time::sleep(Duration::from_millis(600)).await; - ensure!( - h.metric("walshadow_crossing_wedged")? == 0, - "pause left the crossing parked", - ); + let deadline = Instant::now() + Duration::from_secs(30); + loop { + if h.metric("walshadow_crossing_wedged")? == 0 { + break; + } + ensure!(h.alive(), "daemon exited while unparking the crossing"); + ensure!(Instant::now() < deadline, "pause left the crossing parked"); + tokio::time::sleep(Duration::from_millis(250)).await; + } h.ctl_body(&["apply"], "[stream]\npaused = false")?; h.wait_metric( "walshadow_timeline_switches_total",