Skip to content

Commit 40cb0f7

Browse files
committed
Check node health in antithesis. Reduce metric granularity for buffered changes. Admin command for bookie consistency check
1 parent 89297e0 commit 40cb0f7

3 files changed

Lines changed: 205 additions & 18 deletions

File tree

crates/corro-admin/src/lib.rs

Lines changed: 114 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -1,12 +1,13 @@
11
use std::{
2+
collections::HashMap,
23
fmt::Display,
34
time::{Duration, Instant},
45
};
56

67
use camino::Utf8PathBuf;
78
use corro_types::{
89
actor::{ActorId, ClusterId},
9-
agent::{Agent, Bookie},
10+
agent::{Agent, BookedVersions, Bookie},
1011
base::{CrsqlDbVersion, CrsqlSeq},
1112
broadcast::{FocaCmd, FocaInput},
1213
sqlite::SqlitePoolError,
@@ -113,6 +114,7 @@ pub enum Command {
113114
pub enum SyncCommand {
114115
Generate,
115116
ReconcileGaps,
117+
CheckBookieConsistency,
116118
}
117119

118120
#[derive(Debug, Clone, Serialize, Deserialize)]
@@ -282,6 +284,54 @@ async fn handle_conn(
282284

283285
send_success(&mut stream).await;
284286
}
287+
Command::Sync(SyncCommand::CheckBookieConsistency) => {
288+
info_log(
289+
&mut stream,
290+
"checking in-memory and DB bookie consistency...",
291+
)
292+
.await;
293+
294+
let rw_conn = match agent.pool().write_low().await {
295+
Ok(conn) => conn,
296+
Err(e) => {
297+
send_error(&mut stream, e).await;
298+
continue;
299+
}
300+
};
301+
302+
let report = match block_in_place(|| check_bookie_consistency(&rw_conn, bookie))
303+
{
304+
Ok(report) => report,
305+
Err(e) => {
306+
send_error(&mut stream, e).await;
307+
continue;
308+
}
309+
};
310+
311+
match serde_json::to_value(&report) {
312+
Ok(json) => send(&mut stream, Response::Json(json)).await,
313+
Err(e) => {
314+
send_error(&mut stream, e).await;
315+
continue;
316+
}
317+
}
318+
319+
if !report.ok {
320+
send_error(
321+
&mut stream,
322+
format!(
323+
"bookie mismatch: value_mismatch={}, only_in_memory={}, only_in_db={}",
324+
report.value_mismatch_actors.len(),
325+
report.only_in_memory_actors.len(),
326+
report.only_in_db_actors.len()
327+
),
328+
)
329+
.await;
330+
continue;
331+
}
332+
333+
send_success(&mut stream).await;
334+
}
285335
Command::Cluster(ClusterCommand::Rejoin) => {
286336
let (cb_tx, cb_rx) = oneshot::channel();
287337

@@ -707,3 +757,66 @@ fn collapse_gaps(
707757

708758
Ok((deleted, inserted))
709759
}
760+
761+
#[derive(Debug, Serialize)]
762+
struct BookieConsistencyReport {
763+
ok: bool,
764+
in_memory_actors: usize,
765+
db_actors: usize,
766+
value_mismatch_actors: Vec<ActorId>,
767+
only_in_memory_actors: Vec<ActorId>,
768+
only_in_db_actors: Vec<ActorId>,
769+
}
770+
771+
fn check_bookie_consistency(
772+
conn: &rusqlite::Connection,
773+
bookie: &Bookie,
774+
) -> rusqlite::Result<BookieConsistencyReport> {
775+
let _bookie_write_guard = bookie.write_lock_blocking();
776+
777+
let db_map = BookedVersions::load_all_from_conn(conn)?;
778+
779+
let in_memory_map: HashMap<ActorId, BookedVersions> = {
780+
let guard = bookie.owned_guard();
781+
bookie
782+
.iter(&guard)
783+
.map(|(actor_id, booked)| {
784+
let read = booked.read();
785+
(*actor_id, (**read).clone())
786+
})
787+
.collect()
788+
};
789+
790+
let mut value_mismatch_actors = Vec::new();
791+
let mut only_in_memory_actors = Vec::new();
792+
let mut only_in_db_actors = Vec::new();
793+
794+
for (actor_id, mem) in &in_memory_map {
795+
match db_map.get(actor_id) {
796+
Some(db) if db != mem => value_mismatch_actors.push(*actor_id),
797+
Some(_) => {}
798+
None => only_in_memory_actors.push(*actor_id),
799+
}
800+
}
801+
802+
for actor_id in db_map.keys() {
803+
if !in_memory_map.contains_key(actor_id) {
804+
only_in_db_actors.push(*actor_id);
805+
}
806+
}
807+
808+
value_mismatch_actors.sort_unstable();
809+
only_in_memory_actors.sort_unstable();
810+
only_in_db_actors.sort_unstable();
811+
812+
Ok(BookieConsistencyReport {
813+
ok: value_mismatch_actors.is_empty()
814+
&& only_in_memory_actors.is_empty()
815+
&& only_in_db_actors.is_empty(),
816+
in_memory_actors: in_memory_map.len(),
817+
db_actors: db_map.len(),
818+
value_mismatch_actors,
819+
only_in_memory_actors,
820+
only_in_db_actors,
821+
})
822+
}

crates/corro-agent/src/agent/metrics.rs

Lines changed: 82 additions & 17 deletions
Original file line numberDiff line numberDiff line change
@@ -1,25 +1,28 @@
11
use crate::transport::Transport;
2-
use corro_types::{actor::ActorId, agent::Agent};
2+
use antithesis_sdk::assert_always;
3+
use corro_types::agent::Agent;
34
use metrics::gauge;
5+
use serde_json::json;
46
use std::time::Duration;
57
use tokio::task::block_in_place;
68
use tracing::error;
79
use tripwire::Tripwire;
810

911
pub async fn metrics_loop(agent: Agent, transport: Transport, mut tripwire: Tripwire) {
1012
let mut metrics_interval = tokio::time::interval(Duration::from_secs(10));
13+
let on_antithesis = std::env::var("ANTITHESIS_OUTPUT_DIR").is_ok();
1114

1215
loop {
1316
tokio::select! {
1417
_ = metrics_interval.tick() => {},
1518
_ = &mut tripwire => break,
1619
}
1720

18-
block_in_place(|| collect_metrics(&agent, &transport));
21+
block_in_place(|| collect_metrics(&agent, &transport, on_antithesis));
1922
}
2023
}
2124

22-
pub fn collect_metrics(agent: &Agent, transport: &Transport) {
25+
pub fn collect_metrics(agent: &Agent, transport: &Transport, on_antithesis: bool) {
2326
agent.pool().emit_metrics();
2427
transport.emit_metrics();
2528

@@ -50,38 +53,66 @@ pub fn collect_metrics(agent: &Agent, transport: &Transport) {
5053
}
5154
}
5255

53-
// TODO: collect from bookie?
56+
// Buffered changes stats across all actors
5457
match conn
55-
.prepare_cached("SELECT actor_id, (select count(site_id) FROM __corro_buffered_changes WHERE site_id = actor_id) FROM __corro_members")
58+
.prepare_cached(
59+
"
60+
SELECT
61+
COALESCE(json_extract(m.foca_state, '$.state'), 'orphan') AS peer_state,
62+
count(*) AS total,
63+
strftime('%s', 'now') - (min(bc.ts) >> 32) as oldest_age_seconds
64+
FROM __corro_buffered_changes bc
65+
LEFT JOIN __corro_members m ON m.actor_id = bc.site_id
66+
GROUP BY peer_state
67+
",
68+
)
5669
.and_then(|mut prepped| {
5770
prepped
58-
.query_map((), |row| {
59-
Ok((row.get::<_, ActorId>(0)?, row.get::<_, i64>(1)?))
71+
.query_map([], |row| {
72+
Ok((
73+
row.get::<_, String>(0)?,
74+
row.get::<_, i64>(1)?,
75+
row.get::<_, i64>(2)?,
76+
))
6077
})
6178
.and_then(|mapped| mapped.collect::<Result<Vec<_>, _>>())
6279
}) {
63-
Ok(mapped) => {
64-
for (actor_id, count) in mapped {
65-
gauge!("corro.db.buffered.changes.rows.total", "actor_id" => actor_id.to_string()).set(count as f64)
80+
Ok(rows) => {
81+
for (peer_state, total_count, oldest_age_seconds) in rows {
82+
// v2 has O(N) time series while the v1 metric had O(N^2) time series
83+
// The name got changed on purpose so the metric has a fresh namespace
84+
gauge!("corro.db.buffered.changes.v2.total", "peer_state" => peer_state.clone())
85+
.set(total_count as f64);
86+
gauge!("corro.db.buffered.changes.v2.oldest_age_seconds", "peer_state" => peer_state)
87+
.set(oldest_age_seconds as f64);
6688
}
6789
}
6890
Err(e) => {
69-
error!("could not query count for buffered changes: {e}");
91+
error!("could not query buffered changes v2 stats: {e}");
7092
}
7193
}
7294

7395
match conn
74-
.prepare_cached("select actor_id, sum((end - start) + 1) from __corro_bookkeeping_gaps group by actor_id")
96+
.prepare_cached(
97+
"
98+
SELECT
99+
COALESCE(json_extract(m.foca_state, '$.state'), 'orphan') AS peer_state,
100+
COALESCE(SUM((g.end - g.start) + 1), 0) AS gaps_sum
101+
FROM __corro_bookkeeping_gaps g
102+
LEFT JOIN __corro_members m ON m.actor_id = g.actor_id
103+
GROUP BY peer_state
104+
",
105+
)
75106
.and_then(|mut prepped| {
76107
prepped
77-
.query_map((), |row| {
78-
Ok((row.get::<_, ActorId>(0)?, row.get::<_, u64>(1)?))
108+
.query_map([], |row| {
109+
Ok((row.get::<_, String>(0)?, row.get::<_, u64>(1)?))
79110
})
80111
.and_then(|mapped| mapped.collect::<Result<Vec<_>, _>>())
81112
}) {
82-
Ok(mapped) => {
83-
for (actor_id, sum) in mapped {
84-
gauge!("corro.db.gaps.sum", "actor_id" => actor_id.to_string()).set(sum as f64)
113+
Ok(rows) => {
114+
for (peer_state, sum) in rows {
115+
gauge!("corro.db.gaps.sum", "peer_state" => peer_state).set(sum as f64)
85116
}
86117
}
87118
Err(e) => {
@@ -109,4 +140,38 @@ pub fn collect_metrics(agent: &Agent, transport: &Transport) {
109140
gauge!("corro.db.wal.size").set(meta.len() as f64);
110141
}
111142
}
143+
144+
if on_antithesis {
145+
let gaps = conn
146+
.prepare_cached(
147+
"SELECT COALESCE(SUM(end - start + 1), 0) FROM __corro_bookkeeping_gaps",
148+
)
149+
.and_then(|mut prepped| prepped.query_row([], |row| row.get::<_, i64>(0)));
150+
151+
if let Ok(gaps) = gaps {
152+
const MAX_GAPS: i64 = 100;
153+
const MAX_QUEUE: u64 = 16_000;
154+
const MAX_P99_LAG: f64 = 3.0;
155+
156+
let p99_lag = agent.metrics_tracker().quantile_lag(0.99).unwrap_or(0.0);
157+
let queue_size = agent.metrics_tracker().queue_size();
158+
let is_healthy = gaps <= MAX_GAPS && queue_size <= MAX_QUEUE && p99_lag <= MAX_P99_LAG;
159+
let details = json!({
160+
"gaps": gaps,
161+
"queue_size": queue_size,
162+
"p99_lag": p99_lag,
163+
"max_gaps": MAX_GAPS,
164+
"max_queue": MAX_QUEUE,
165+
"max_p99_lag": MAX_P99_LAG,
166+
});
167+
168+
assert_always!(
169+
is_healthy,
170+
"Corrosion node remains healthy in antithesis",
171+
&details
172+
);
173+
} else if let Err(e) = gaps {
174+
error!("could not query gaps for antithesis health assertion: {e}");
175+
}
176+
}
112177
}

crates/corrosion/src/main.rs

Lines changed: 9 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -499,6 +499,13 @@ async fn process_cli(cli: Cli) -> eyre::Result<()> {
499499
))
500500
.await?;
501501
}
502+
Command::Sync(SyncCommand::CheckBookieConsistency) => {
503+
let mut conn = AdminConn::connect(cli.admin_path()).await?;
504+
conn.send_command(corro_admin::Command::Sync(
505+
corro_admin::SyncCommand::CheckBookieConsistency,
506+
))
507+
.await?;
508+
}
502509
Command::Template { template, flags } => {
503510
command::tpl::run(cli.api_addr()?, template, flags).await?;
504511
}
@@ -780,6 +787,8 @@ enum ConsulCommand {
780787
enum SyncCommand {
781788
/// Generate a sync message from the current agent
782789
Generate,
790+
/// Check in-memory bookie state against DB-loaded bookie state
791+
CheckBookieConsistency,
783792
ReconcileGaps,
784793
}
785794

0 commit comments

Comments
 (0)