use std::collections::BTreeMap;
use crate::model::diff::{byte_diff, diff as value_diff};
use crate::model::facts::{KeyFacts, KeyShape, OriginKind};
use crate::model::origin_map::MapPlan;
use crate::report::{
Asked, Holder, KeyChange, Snapshot, SnapshotDiff, SnapshotRow, SubjectDelta, ZsnapHeader,
};
#[derive(Debug, Clone, Copy)]
pub struct DiffOpts {
pub max_changes: usize,
pub max_keys: usize,
pub stamps_alone: bool,
}
impl Default for DiffOpts {
fn default() -> Self {
DiffOpts {
max_changes: 20,
max_keys: crate::model::bounded::DEFAULT_MAX_KEYS,
stamps_alone: true,
}
}
}
enum Outcome {
Added,
Removed,
Changed(Box<KeyChange>),
Unchanged,
}
pub fn diff_snapshots(a: &Snapshot, b: &Snapshot, opts: DiffOpts) -> SnapshotDiff {
bound(&a.header, &b.header, compare(a, b, opts), opts.max_keys)
}
pub fn diff_normalized(a: &Snapshot, b: &Snapshot, plan: &MapPlan, opts: DiffOpts) -> SnapshotDiff {
let mut out = SnapshotDiff {
a: a.header.clone(),
b: b.header.clone(),
added: Vec::new(),
removed: Vec::new(),
changed: Vec::new(),
unchanged: 0,
truncated: 0,
origin_map: Asked::Asked(plan.pairs.clone()),
unmapped: plan.unmapped.clone(),
by_subject: Asked::NotAsked,
};
if !plan.is_complete() {
return out;
}
let opts = DiffOpts {
stamps_alone: false,
..opts
};
let b_to_a = plan.b_to_a();
let a_norm = Snapshot {
header: a.header.clone(),
rows: a
.rows
.iter()
.map(|r| canonical_bridge(r, &a.header.base, None))
.collect(),
};
let b_norm = Snapshot {
header: b.header.clone(),
rows: b
.rows
.iter()
.map(|r| rewrite_row(r, &b.header.base, &a.header.base, &b_to_a))
.collect(),
};
let outcomes = compare(&a_norm, &b_norm, opts);
let mut subjects: BTreeMap<String, SubjectDelta> = BTreeMap::new();
for (key, outcome) in &outcomes {
let subject = subject_of(&a.header.base, key);
let s = subjects
.entry(subject.clone())
.or_insert_with(|| SubjectDelta {
subject,
compared: 0,
differing: 0,
only_in_a: 0,
only_in_b: 0,
example: None,
});
match outcome {
Outcome::Added => s.only_in_b += 1,
Outcome::Removed => s.only_in_a += 1,
Outcome::Unchanged => s.compared += 1,
Outcome::Changed(c) => {
s.compared += 1;
s.differing += 1;
if s.example.is_none() {
s.example = Some((**c).clone());
}
}
}
}
let bounded = bound(&a.header, &b.header, outcomes, opts.max_keys);
out.added = bounded.added;
out.removed = bounded.removed;
out.changed = bounded.changed;
out.unchanged = bounded.unchanged;
out.truncated = bounded.truncated;
out.by_subject = Asked::Asked(subjects.into_values().collect());
out
}
fn subject_of(base: &str, key: &str) -> String {
let facts = KeyFacts::project(base, key);
let (KeyShape::V1(f), Some(relative)) = (&facts.shape, zenkey::grammar::strip_base(base, key))
else {
return key.to_string();
};
let mut chunks = relative.split('/');
chunks.next(); if f.origin_kind == OriginKind::Host {
chunks.next(); }
chunks.collect::<Vec<_>>().join("/")
}
fn rewrite_row(
row: &SnapshotRow,
from_base: &str,
to_base: &str,
b_to_a: &BTreeMap<&str, &str>,
) -> SnapshotRow {
let facts = KeyFacts::project(from_base, &row.key);
let (KeyShape::V1(f), Some(relative)) = (
&facts.shape,
zenkey::grammar::strip_base(from_base, &row.key),
) else {
return row.clone();
};
let mapped = if f.origin_kind == OriginKind::Host {
b_to_a.get(f.origin.as_str()).copied()
} else {
None
};
let mut chunks: Vec<&str> = relative.split('/').collect();
if let Some(to) = mapped
&& chunks.len() > 1
{
chunks[1] = to;
}
let mut out = row.clone();
out.key = zenkey::grammar::with_base(to_base, chunks.join("/"));
if let Some(to) = mapped {
out.holder = match &row.holder {
Holder::Live {
origin,
answered_by,
} if origin == &f.origin => Holder::Live {
origin: to.to_string(),
answered_by: *answered_by,
},
Holder::StorageOnly { origin } if origin == &f.origin => Holder::StorageOnly {
origin: to.to_string(),
},
other => other.clone(),
};
}
canonical_bridge(&out, to_base, mapped.map(|to| (f.origin.as_str(), to)))
}
fn canonical_bridge(row: &SnapshotRow, base: &str, rename: Option<(&str, &str)>) -> SnapshotRow {
use base64::Engine as _;
let facts = KeyFacts::project(base, &row.key);
let KeyShape::V1(f) = &facts.shape else {
return row.clone();
};
let is_bridge = f.origin_kind == OriginKind::Host
&& f.class == "state"
&& matches!(f.subject.as_slice(), [s] if s == "health" || s == "sensor");
if !is_bridge || row.delete {
return row.clone();
}
let Some(serde_json::Value::Object(mut doc)) = structural_of(row) else {
return row.clone();
};
let own = rename.map(|(from, _)| from).unwrap_or(f.origin.as_str());
if doc.get("host_id").and_then(|v| v.as_str()) != Some(own) {
return row.clone();
}
if let Some((_, to)) = rename {
doc.insert("host_id".into(), serde_json::Value::String(to.to_string()));
}
let mut out = row.clone();
let bytes = serde_json::to_vec(&serde_json::Value::Object(doc)).unwrap_or_default();
out.bytes = Some(base64::engine::general_purpose::STANDARD.encode(bytes));
out
}
fn compare(a: &Snapshot, b: &Snapshot, opts: DiffOpts) -> Vec<(String, Outcome)> {
fn by_key(s: &Snapshot) -> BTreeMap<&str, &SnapshotRow> {
s.rows.iter().map(|r| (r.key.as_str(), r)).collect()
}
let (ra, rb) = (by_key(a), by_key(b));
let mut out = Vec::with_capacity(ra.len() + rb.len());
for (key, row_a) in &ra {
let outcome = match rb.get(key) {
None => Outcome::Removed,
Some(row_b) => match key_change(row_a, row_b, opts) {
None => Outcome::Unchanged,
Some(change) => Outcome::Changed(Box::new(change)),
},
};
out.push(((*key).to_string(), outcome));
}
for key in rb.keys() {
if !ra.contains_key(key) {
out.push(((*key).to_string(), Outcome::Added));
}
}
out
}
fn bound(
a: &ZsnapHeader,
b: &ZsnapHeader,
outcomes: Vec<(String, Outcome)>,
max_keys: usize,
) -> SnapshotDiff {
let mut out = SnapshotDiff {
a: a.clone(),
b: b.clone(),
added: Vec::new(),
removed: Vec::new(),
changed: Vec::new(),
unchanged: 0,
truncated: 0,
origin_map: Asked::NotAsked,
unmapped: Vec::new(),
by_subject: Asked::NotAsked,
};
let mut listed = 0usize;
for (key, outcome) in outcomes {
if matches!(outcome, Outcome::Unchanged) {
out.unchanged += 1;
continue;
}
if listed >= max_keys {
out.truncated += 1;
continue;
}
listed += 1;
match outcome {
Outcome::Added => out.added.push(key),
Outcome::Removed => out.removed.push(key),
Outcome::Changed(c) => out.changed.push(*c),
Outcome::Unchanged => unreachable!("counted above"),
}
}
out
}
pub(crate) fn structural_of(row: &SnapshotRow) -> Option<serde_json::Value> {
let bytes = payload(row)?;
#[cfg(feature = "decode")]
{
crate::model::decode::structural_value(&bytes)
}
#[cfg(not(feature = "decode"))]
{
serde_json::from_slice(&bytes).ok()
}
}
fn payload(row: &SnapshotRow) -> Option<Vec<u8>> {
row.payload()
}
fn key_change(a: &SnapshotRow, b: &SnapshotRow, opts: DiffOpts) -> Option<KeyChange> {
let mut change = KeyChange {
key: a.key.clone(),
value: None,
bytes: None,
verdict: None,
registration: None,
holder: None,
timestamp: (a.timestamp.clone(), b.timestamp.clone()),
};
let mut moved = false;
if a.delete != b.delete || a.bytes != b.bytes {
moved = true;
match (structural_of(a), structural_of(b)) {
(Some(va), Some(vb)) => {
let d = value_diff(&va, &vb, opts.max_changes);
if d.is_empty() {
change.bytes = Some(byte_diff(
&payload(a).unwrap_or_default(),
&payload(b).unwrap_or_default(),
));
} else {
change.value = Some(d);
}
}
_ => {
change.bytes = Some(byte_diff(
&payload(a).unwrap_or_default(),
&payload(b).unwrap_or_default(),
));
}
}
}
if a.verdict != b.verdict {
moved = true;
change.verdict = Some((a.verdict.clone(), b.verdict.clone()));
}
if a.registration != b.registration {
moved = true;
change.registration = Some((a.registration, b.registration));
}
if a.holder != b.holder {
moved = true;
change.holder = Some((a.holder.clone(), b.holder.clone()));
}
if opts.stamps_alone && a.timestamp != b.timestamp {
moved = true;
}
moved.then_some(change)
}
#[cfg(test)]
mod tests {
use super::*;
use crate::report::{AnsweredBy, Holder, RegistrationWire, VerdictWire, ZsnapHeader};
fn header() -> ZsnapHeader {
ZsnapHeader {
zsnap: 1,
selectors: vec!["v1/**".into()],
base: String::new(),
collected_at: "2026-09-06T00:00:00Z".into(),
collection_span_s: 0.5,
asked: 1,
answered: 0,
elided: 0,
errors: 0,
superseded: 0,
roster: Asked::NotAsked,
}
}
fn row(key: &str, body: &[u8]) -> SnapshotRow {
use base64::Engine as _;
SnapshotRow {
key: key.into(),
delete: false,
bytes: Some(base64::engine::general_purpose::STANDARD.encode(body)),
encoding: Some("application/json".into()),
timestamp: None,
stamper: None,
source: None,
source_zid: None,
registration: RegistrationWire::RegistryNotLoaded,
verdict: VerdictWire::NotValidated {
reason: "no_registry".into(),
},
holder: Holder::Unattributed {
reason: "roster not asked".into(),
},
}
}
fn snap(rows: Vec<SnapshotRow>) -> Snapshot {
Snapshot {
header: header(),
rows,
}
}
#[test]
fn a_self_diff_is_empty() {
let s = snap(vec![
row("v1/h-aaaaaaaaaaaa/state/p/a", b"{\"x\":1}"),
row("v1/h-aaaaaaaaaaaa/state/p/b", b"text"),
]);
let d = diff_snapshots(&s, &s, DiffOpts::default());
assert!(!d.differs());
assert_eq!(d.unchanged, 2);
assert!(d.changed.is_empty() && d.added.is_empty() && d.removed.is_empty());
}
#[test]
fn a_non_json_payload_falls_back_to_bytes() {
let a = snap(vec![
row("k/json", b"{\"x\":1}"),
row("k/text", b"hello world"),
]);
let b = snap(vec![
row("k/json", b"{\"x\":2}"),
row("k/text", b"hello there"),
]);
let d = diff_snapshots(&a, &b, DiffOpts::default());
assert_eq!(d.changed.len(), 2);
let json = d.changed.iter().find(|c| c.key == "k/json").unwrap();
assert!(json.value.is_some() && json.bytes.is_none());
assert_eq!(json.value.as_ref().unwrap().changes[0].path(), "x");
let text = d.changed.iter().find(|c| c.key == "k/text").unwrap();
assert!(text.bytes.is_some() && text.value.is_none());
assert_eq!(text.bytes.unwrap().common_prefix, 6);
}
#[test]
fn a_facet_pair_rides_only_when_that_facet_moved() {
let mut live = row("k", b"{}");
live.holder = Holder::Live {
origin: "h-aaaaaaaaaaaa".into(),
answered_by: AnsweredBy::Stamper,
};
let a = snap(vec![row("k", b"{}")]);
let b = snap(vec![live]);
let d = diff_snapshots(&a, &b, DiffOpts::default());
let c = &d.changed[0];
assert!(c.holder.is_some());
assert!(c.value.is_none() && c.bytes.is_none());
assert!(c.verdict.is_none() && c.registration.is_none());
}
#[test]
fn added_and_removed_keys_are_listed_by_side() {
let a = snap(vec![row("only/a", b"1"), row("both", b"1")]);
let b = snap(vec![row("only/b", b"1"), row("both", b"1")]);
let d = diff_snapshots(&a, &b, DiffOpts::default());
assert_eq!(d.added, ["only/b"]);
assert_eq!(d.removed, ["only/a"]);
assert_eq!(d.unchanged, 1);
assert!(d.differs());
}
#[test]
fn differing_keys_past_the_bound_are_counted() {
let a = snap((0..5).map(|i| row(&format!("k/{i}"), b"1")).collect());
let b = snap((3..8).map(|i| row(&format!("k/{i}"), b"2")).collect());
let d = diff_snapshots(
&a,
&b,
DiffOpts {
max_keys: 3,
..DiffOpts::default()
},
);
let listed = d.added.len() + d.removed.len() + d.changed.len();
assert_eq!(listed, 3);
assert_eq!(
d.truncated, 5,
"3 removed + 2 changed + 3 added = 8, 3 listed"
);
assert!(d.differs());
}
#[test]
fn a_tombstone_against_a_value_is_a_byte_change() {
let mut gone = row("k", b"{}");
gone.delete = true;
gone.bytes = None;
let d = diff_snapshots(
&snap(vec![row("k", b"{}")]),
&snap(vec![gone]),
DiffOpts::default(),
);
let c = &d.changed[0];
assert!(c.bytes.is_some());
assert_eq!(c.bytes.unwrap().new_len, 0);
}
use crate::model::origin_map::tests::{host, snap as fleet};
use crate::model::origin_map::{origin_profiles, plan_map};
use crate::report::MapEvidence;
const A1: &str = "h-aaaaaaaaaaa1";
const A2: &str = "h-aaaaaaaaaaa2";
const B1: &str = "h-bbbbbbbbbbb1";
const B2: &str = "h-bbbbbbbbbbb2";
fn renamed(base_b: &str) -> (Snapshot, Snapshot) {
let a = fleet(
"prod",
[
host("prod", A1, "web", &[]),
host("prod", A2, "db", &["logs"]),
]
.concat(),
);
let mut b = fleet(
base_b,
[
host(base_b, B1, "web", &[]),
host(base_b, B2, "db", &["logs"]),
]
.concat(),
);
for r in &mut b.rows {
r.timestamp = Some("7f3b2a1c00000009/ef56".into());
}
(a, b)
}
fn plan(a: &Snapshot, b: &Snapshot) -> MapPlan {
plan_map(&origin_profiles(a), &origin_profiles(b), &[]).unwrap()
}
#[test]
fn a_renamed_fleet_diffs_to_zero_once_aligned() {
let (a, b) = renamed("stg");
let plain = diff_snapshots(&a, &b, DiffOpts::default());
assert!(plain.differs(), "verbatim, nothing lines up");
assert!(plain.origin_map.is_not_asked() && plain.by_subject.is_not_asked());
let d = diff_normalized(&a, &b, &plan(&a, &b), DiffOpts::default());
assert!(!d.differs(), "{d:?}");
assert_eq!(d.unchanged, 5);
assert!(d.unmapped.is_empty());
let pairs = d.origin_map.as_option().unwrap();
assert_eq!(pairs.len(), 2);
assert!(matches!(pairs[0].evidence, MapEvidence::Label { .. }));
let subjects = d.by_subject.as_option().unwrap();
assert!(
subjects
.iter()
.all(|s| s.differing == 0 && s.only_in_a == 0 && s.only_in_b == 0)
);
assert_eq!(
subjects
.iter()
.map(|s| s.subject.as_str())
.collect::<Vec<_>>(),
[
"state/logs/rotated",
"state/sysinfo/health",
"telemetry/sysinfo/disk/root/used"
]
);
assert_eq!(subjects[1].compared, 2, "one health document per origin");
assert_eq!(crate::report::judgement_exit_code(&d.to_judgement()), 0);
}
#[test]
fn a_changed_value_on_one_host_reads_as_one_of_n_on_its_subject() {
let (a, mut b) = renamed("prod");
let disk = b
.rows
.iter_mut()
.find(|r| r.key == format!("prod/v1/{B2}/telemetry/sysinfo/disk/root/used"))
.unwrap();
disk.bytes = Some({
use base64::Engine as _;
base64::engine::general_purpose::STANDARD.encode(r#"{"value":97.0}"#)
});
let d = diff_normalized(&a, &b, &plan(&a, &b), DiffOpts::default());
assert!(d.differs());
assert_eq!(d.changed.len(), 1);
assert_eq!(
d.changed[0].key,
format!("prod/v1/{A2}/telemetry/sysinfo/disk/root/used"),
"the changed key is spelled in a's origin"
);
let subjects = d.by_subject.as_option().unwrap();
let disk = subjects
.iter()
.find(|s| s.subject == "telemetry/sysinfo/disk/root/used")
.unwrap();
assert_eq!((disk.compared, disk.differing), (2, 1));
assert!(disk.example.as_ref().unwrap().value.is_some());
assert_eq!(crate::report::judgement_exit_code(&d.to_judgement()), 1);
}
#[test]
fn a_key_only_one_side_holds_rolls_up_as_only_in() {
let (a, mut b) = renamed("prod");
b.rows.retain(|r| !r.key.ends_with("/logs/rotated"));
let d = diff_normalized(&a, &b, &plan(&a, &b), DiffOpts::default());
assert_eq!(d.removed, [format!("prod/v1/{A2}/state/logs/rotated")]);
let logs = &d.by_subject.as_option().unwrap()[0];
assert_eq!(logs.subject, "state/logs/rotated");
assert_eq!((logs.compared, logs.only_in_a, logs.only_in_b), (0, 1, 0));
assert!(
d.changed.iter().all(|c| c.holder.is_none()),
"a renamed holder is not a moved holder: {:?}",
d.changed
);
}
#[test]
fn an_incomplete_plan_is_refused_not_compared_around() {
let (a, mut b) = renamed("prod");
b.rows
.extend(host("prod", "h-bbbbbbbbbbb3", "db", &["logs"]));
let plan = plan(&a, &b);
assert_eq!(
plan.unmapped.len(),
3,
"db is claimed twice in b, so a's db is unpaired too"
);
let d = diff_normalized(&a, &b, &plan, DiffOpts::default());
assert!(d.refused());
assert_eq!(d.unmapped.len(), plan.unmapped.len(), "never dropped");
assert_eq!(
d.origin_map.as_option().unwrap().len(),
1,
"web still paired"
);
assert!(d.by_subject.is_not_asked());
assert!(d.added.is_empty() && d.removed.is_empty() && d.changed.is_empty());
assert_eq!(d.unchanged, 0);
assert!(!d.differs());
assert_eq!(crate::report::judgement_exit_code(&d.to_judgement()), 2);
}
#[test]
fn a_foreign_host_id_is_not_rewritten() {
let (a, mut b) = renamed("prod");
let health = b
.rows
.iter_mut()
.find(|r| r.key == format!("prod/v1/{B1}/state/sysinfo/health"))
.unwrap();
health.bytes = Some({
use base64::Engine as _;
base64::engine::general_purpose::STANDARD
.encode(r#"{"host_id":"h-000000000000","source":"web","status":"ok"}"#)
});
let plan = plan_map(
&origin_profiles(&a),
&origin_profiles(&b),
&[(
zenkey::origin::HostId::parse(A1).unwrap(),
zenkey::origin::HostId::parse(B1).unwrap(),
)],
)
.unwrap();
let d = diff_normalized(&a, &b, &plan, DiffOpts::default());
let c = d
.changed
.iter()
.find(|c| c.key.ends_with("/health"))
.unwrap();
let v = c.value.as_ref().unwrap();
assert_eq!(v.changes[0].path(), "host_id");
}
#[test]
fn a_subject_keeps_a_service_origin_and_a_foreign_key_verbatim() {
assert_eq!(
subject_of("acme", "acme/v1/h-aaaaaaaaaaa1/state/sysinfo/health"),
"state/sysinfo/health"
);
assert_eq!(
subject_of("acme", "acme/v1/@catalog/state/entity/x"),
"@catalog/state/entity/x"
);
assert_eq!(subject_of("acme", "acme/plain/leak"), "acme/plain/leak");
assert_eq!(
subject_of("", "v1/h-aaaaaaaaaaa1/telemetry/sysinfo-2/cpu"),
"telemetry/sysinfo-2/cpu"
);
}
}