use serde_json::Value;
use std::collections::HashMap;
use super::scenario::{FuzzOutcome, Scenario};
use crate::zakura::testkit::TraceReader;
#[derive(Clone, Copy, Debug, Default)]
pub(crate) struct InvariantReport {
pub(crate) state_samples: usize,
pub(crate) max_outstanding: u64,
pub(crate) peak_budget_reserved: u64,
pub(crate) peak_retained_pipeline_wire_bytes: u64,
pub(crate) peak_active_pipeline_decoded_attributed_memory_bytes: u64,
pub(crate) final_active_pipeline_decoded_attributed_memory_bytes: u64,
pub(crate) decoded_stage_totals_match: bool,
pub(crate) final_budget_reserved: u64,
pub(crate) protocol_rejects: usize,
pub(crate) total_requests: usize,
pub(crate) max_requests_without_block_progress: u64,
pub(crate) max_unproven_requests_without_block_progress: u64,
pub(crate) floor_bypass_requests: usize,
pub(crate) peak_cwnd_bytes: u64,
pub(crate) peak_inflight_bytes: u64,
pub(crate) peak_cwnd_requests: u64,
pub(crate) min_reliability_permille: u64,
pub(crate) final_cwnd_bytes: u64,
pub(crate) final_reliability_permille: u64,
}
pub(crate) fn report(reader: &TraceReader) -> InvariantReport {
let state_rows: Vec<&Value> = reader
.table("block_sync")
.rows()
.into_iter()
.filter(|row| event(row) == Some("block_sync_state"))
.collect();
let max_outstanding = state_rows
.iter()
.filter_map(|row| u64_field(row, "outstanding"))
.max()
.unwrap_or(0);
let peak_budget_reserved = state_rows
.iter()
.filter_map(|row| u64_field(row, "budget_reserved"))
.max()
.unwrap_or(0);
let final_budget_reserved = state_rows
.iter()
.rev()
.find_map(|row| u64_field(row, "budget_reserved"))
.unwrap_or(0);
let peak_retained_pipeline_wire_bytes = state_rows
.iter()
.filter_map(|row| u64_field(row, "retained_pipeline_wire_bytes"))
.max()
.unwrap_or(0);
let peak_active_pipeline_decoded_attributed_memory_bytes = state_rows
.iter()
.filter_map(|row| u64_field(row, "active_pipeline_decoded_attributed_memory_bytes"))
.max()
.unwrap_or(0);
let final_active_pipeline_decoded_attributed_memory_bytes = state_rows
.iter()
.rev()
.find_map(|row| u64_field(row, "active_pipeline_decoded_attributed_memory_bytes"))
.unwrap_or(0);
let decoded_stage_totals_match = state_rows.iter().all(|row| {
let stage_total = u64_field(row, "sequencer_input_decoded_attributed_memory_bytes")
.unwrap_or(0)
.saturating_add(u64_field(row, "reorder_decoded_attributed_memory_bytes").unwrap_or(0))
.saturating_add(
u64_field(row, "applying_decoded_attributed_memory_bytes").unwrap_or(0),
);
u64_field(row, "active_pipeline_decoded_attributed_memory_bytes") == Some(stage_total)
});
let protocol_rejects = reader
.table("block_sync")
.count("block_peer_protocol_reject");
let total_requests = reader.table("block_sync").count("block_get_blocks_sent");
let max_requests_without_block_progress = max_requests_without_block_progress(reader);
let max_unproven_requests_without_block_progress =
max_unproven_requests_without_block_progress(reader);
let body_rows: Vec<&Value> = reader
.table("block_sync")
.rows()
.into_iter()
.filter(|row| event(row) == Some("block_body_received"))
.collect();
let floor_bypass_requests = reader
.table("block_sync")
.rows()
.into_iter()
.filter(|row| event(row) == Some("block_get_blocks_sent"))
.filter(|row| u64_field(row, "floor_bypass") == Some(1))
.count();
let peak_cwnd_bytes = body_rows
.iter()
.filter_map(|row| u64_field(row, "bbr_cwnd_bytes"))
.max()
.unwrap_or(0);
let peak_inflight_bytes = body_rows
.iter()
.filter_map(|row| u64_field(row, "bbr_inflight_bytes"))
.max()
.unwrap_or(0);
let peak_cwnd_requests = body_rows
.iter()
.filter_map(|row| u64_field(row, "bbr_cwnd"))
.max()
.unwrap_or(0);
let min_reliability_permille = reader
.table("block_sync")
.rows()
.into_iter()
.filter_map(|row| u64_field(row, "bbr_reliability_permille"))
.min()
.unwrap_or(1000);
let final_cwnd_bytes = reader
.table("block_sync")
.rows()
.into_iter()
.rev()
.find_map(|row| u64_field(row, "bbr_cwnd_bytes"))
.unwrap_or(0);
let final_reliability_permille = reader
.table("block_sync")
.rows()
.into_iter()
.rev()
.find_map(|row| u64_field(row, "bbr_reliability_permille"))
.unwrap_or(1000);
InvariantReport {
state_samples: state_rows.len(),
max_outstanding,
peak_budget_reserved,
peak_retained_pipeline_wire_bytes,
peak_active_pipeline_decoded_attributed_memory_bytes,
final_active_pipeline_decoded_attributed_memory_bytes,
decoded_stage_totals_match,
final_budget_reserved,
protocol_rejects,
total_requests,
max_requests_without_block_progress,
max_unproven_requests_without_block_progress,
floor_bypass_requests,
peak_cwnd_bytes,
peak_inflight_bytes,
peak_cwnd_requests,
min_reliability_permille,
final_cwnd_bytes,
final_reliability_permille,
}
}
fn max_unproven_requests_without_block_progress(reader: &TraceReader) -> u64 {
reader
.table("block_sync")
.rows()
.into_iter()
.filter(|row| event(row) == Some("block_get_blocks_sent"))
.filter(|row| u64_field(row, "block_progress_proven") == Some(0))
.filter_map(|row| u64_field(row, "requests_without_block_progress"))
.max()
.unwrap_or(0)
}
fn max_requests_without_block_progress(reader: &TraceReader) -> u64 {
let mut streaks: HashMap<String, u64> = HashMap::new();
let mut max_streak = 0u64;
for row in reader.table("block_sync").rows() {
let Some(peer) = str_field(row, "peer") else {
continue;
};
match event(row) {
Some("block_peer_connected") => {
streaks.insert(peer.to_string(), 0);
}
Some("block_get_blocks_sent") => {
let streak = streaks
.entry(peer.to_string())
.and_modify(|streak| *streak = streak.saturating_add(1))
.or_insert(1);
max_streak = max_streak.max(*streak);
}
Some("block_body_received")
| Some("block_peer_disconnected")
| Some("block_peer_protocol_reject") => {
streaks.insert(peer.to_string(), 0);
}
_ => {}
}
}
max_streak
}
pub(crate) fn assert_core(
scenario: &Scenario,
outcome: &FuzzOutcome,
report: &InvariantReport,
outstanding_slack: u64,
) {
assert!(
outcome.reached_target(),
"sync stalled at {} of {} (state_samples={}, max_outstanding={}, rejects={})",
outcome.committed_tip.0,
outcome.target.0,
report.state_samples,
report.max_outstanding,
report.protocol_rejects,
);
assert!(
report.state_samples > 0,
"run emitted no block_sync_state rows",
);
assert!(
report.decoded_stage_totals_match,
"decoded pipeline aggregate drifted from its stage totals",
);
assert!(
report.peak_active_pipeline_decoded_attributed_memory_bytes > 0,
"successful run never observed active decoded pipeline bytes",
);
let outstanding_bound: u64 = scenario
.peers
.iter()
.map(|peer| u64::from(peer.max_inflight_requests))
.sum::<u64>()
.saturating_add(outstanding_slack);
assert!(
report.max_outstanding <= outstanding_bound,
"aggregate outstanding {} exceeded the advertised-inflight bound {}",
report.max_outstanding,
outstanding_bound,
);
let overdraft_slack = scenario.config.floor_request_byte_reservation();
assert!(
report.peak_budget_reserved
<= scenario
.config
.max_inflight_block_bytes
.saturating_add(overdraft_slack),
"peak reserved bytes {} exceeded the in-flight request budget {} (+{} floor-overdraft slack)",
report.peak_budget_reserved,
scenario.config.max_inflight_block_bytes,
overdraft_slack,
);
assert!(
report.final_budget_reserved <= overdraft_slack,
"reserved request bytes {} must drain to zero at quiescence (allowing one \
in-flight receipt release, <= {})",
report.final_budget_reserved,
overdraft_slack,
);
}
fn event(row: &Value) -> Option<&str> {
row.get("event").and_then(Value::as_str)
}
fn u64_field(row: &Value, field: &str) -> Option<u64> {
row.get(field).and_then(Value::as_u64)
}
fn str_field<'a>(row: &'a Value, field: &str) -> Option<&'a str> {
row.get(field).and_then(Value::as_str)
}