Skip to main content

faucet_core/
diff.rs

1//! Content verification primitives for `faucet verify` (#701).
2//!
3//! A pipeline can report green while its destination quietly diverges from the
4//! source: a missed change event, a hand edit downstream, a retried page that
5//! landed twice. Row counts (`reconcile:`) cannot see a row that exists on both
6//! sides with different values. This module holds the **pure** building blocks
7//! the verifier composes:
8//!
9//! - [`KeyRange`] — a half-open slice of an integer key space, bisectable, and
10//!   convertible into the PK-range [`ShardSpec`] every SQL source already
11//!   understands (so a range read needs no new connector code);
12//! - [`Normalizer`] + [`row_hash`] — a canonical, type-tolerant fingerprint of
13//!   one row's compared columns, so `1` and `1.0`, or two spellings of the same
14//!   instant, hash alike;
15//! - [`DigestAccumulator`] / [`ContentDigest`] — an order-independent digest of
16//!   a set of rows (count + folded row hashes + observed key bounds), computed
17//!   client-side while streaming;
18//! - [`ServerDigest`] — the same shape computed *inside* the backend (only
19//!   comparable when both sides report the same `algorithm`);
20//! - [`diff_rows`] — the row-level comparison of two leaf ranges, keyed.
21//!
22//! Nothing here performs I/O; the CLI (`cli/src/verify/`) drives the sources.
23
24use crate::shard::ShardSpec;
25use serde::{Deserialize, Serialize};
26use serde_json::{Map, Value};
27use std::collections::{BTreeMap, HashMap};
28
29/// A half-open integer key range `[lo, hi)`; `None` = unbounded on that side.
30#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
31pub struct KeyRange {
32    pub lo: Option<i64>,
33    pub hi: Option<i64>,
34}
35
36impl KeyRange {
37    /// The whole key space.
38    pub const ALL: KeyRange = KeyRange { lo: None, hi: None };
39
40    /// Whether `k` falls inside this range.
41    pub fn contains(&self, k: i64) -> bool {
42        self.lo.is_none_or(|lo| k >= lo) && self.hi.is_none_or(|hi| k < hi)
43    }
44
45    /// Number of integer keys the range spans, when both ends are bounded.
46    pub fn width(&self) -> Option<u128> {
47        match (self.lo, self.hi) {
48            (Some(lo), Some(hi)) if hi > lo => Some((hi as i128 - lo as i128) as u128),
49            (Some(_), Some(_)) => Some(0),
50            _ => None,
51        }
52    }
53
54    /// The PK-range shard descriptor a SQL source's `apply_shard` narrows to —
55    /// the same shape [`plan_pk_shards`](crate::shard::plan_pk_shards) emits, so
56    /// the range predicate is built by one shared, tested routine. NULL keys are
57    /// owned by the unbounded-above range, so a NULL-key row is read exactly
58    /// once across a partition of the key space.
59    pub fn to_shard(&self, key: &str) -> ShardSpec {
60        ShardSpec::new(
61            format!(
62                "{}..{}",
63                self.lo
64                    .map(|v| v.to_string())
65                    .unwrap_or_else(|| "-inf".into()),
66                self.hi
67                    .map(|v| v.to_string())
68                    .unwrap_or_else(|| "+inf".into())
69            ),
70            serde_json::json!({
71                "key": key,
72                "lo": self.lo.unwrap_or(0),
73                "hi": self.hi.unwrap_or(0),
74                "lo_unbounded": self.lo.is_none(),
75                "hi_unbounded": self.hi.is_none(),
76                "include_null": self.hi.is_none(),
77            }),
78        )
79    }
80
81    /// Split the range in two at its midpoint. An unbounded side is first
82    /// clamped to the observed key bounds (`obs`, the union of both sides'
83    /// `key_min..=key_max`), because a digest can only tell us where rows *are*.
84    /// Returns `None` when the range cannot be split further (one key wide, or
85    /// no observed bounds to clamp with) — the caller then treats it as a leaf.
86    pub fn bisect(&self, obs: Option<(i64, i64)>) -> Option<(KeyRange, KeyRange)> {
87        let lo = self.lo.or(obs.map(|(min, _)| min))?;
88        let hi = self.hi.or(obs.map(|(_, max)| max.saturating_add(1)))?;
89        if hi - lo < 2 {
90            return None;
91        }
92        let mid = lo + (hi - lo) / 2;
93        // Keep the outer edges as they were (possibly unbounded) so the two
94        // halves still tile exactly the parent range.
95        Some((
96            KeyRange {
97                lo: self.lo,
98                hi: Some(mid),
99            },
100            KeyRange {
101                lo: Some(mid),
102                hi: self.hi,
103            },
104        ))
105    }
106}
107
108impl std::fmt::Display for KeyRange {
109    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
110        match (self.lo, self.hi) {
111            (None, None) => write!(f, "all"),
112            (Some(lo), None) => write!(f, "[{lo}, +inf)"),
113            (None, Some(hi)) => write!(f, "(-inf, {hi})"),
114            (Some(lo), Some(hi)) => write!(f, "[{lo}, {hi})"),
115        }
116    }
117}
118
119/// Split `[min, max]` into `target` contiguous ranges tiling the whole key
120/// space (the first is unbounded below, the last unbounded above), via the same
121/// planner PK-range sharding uses.
122pub fn plan_ranges(min: i64, max: i64, target: usize) -> Vec<KeyRange> {
123    crate::shard::plan_pk_shards("k", min, max, target)
124        .iter()
125        .filter_map(crate::shard::PkShardBounds::from_spec)
126        .map(|b| KeyRange {
127            lo: (!b.lo_unbounded).then_some(b.lo),
128            hi: (!b.hi_unbounded).then_some(b.hi),
129        })
130        .collect()
131}
132
133/// How values are canonicalised before hashing and comparing, so two backends'
134/// spellings of the same value agree.
135#[derive(Debug, Clone, Serialize, Deserialize, schemars::JsonSchema, PartialEq)]
136#[serde(deny_unknown_fields)]
137pub struct Normalizer {
138    /// Round floats to this absolute tolerance before comparing (`0` = exact).
139    #[serde(default)]
140    pub float_tolerance: f64,
141    /// Parse strings that look like timestamps (RFC 3339, or
142    /// `YYYY-MM-DD HH:MM:SS[.fff]`) and compare them as UTC microseconds.
143    #[serde(default = "default_true")]
144    pub timestamps: bool,
145    /// Treat a string holding a number (`"42"`, `"1.5"`) as that number.
146    #[serde(default)]
147    pub numeric_strings: bool,
148}
149
150fn default_true() -> bool {
151    true
152}
153
154impl Default for Normalizer {
155    fn default() -> Self {
156        Self {
157            float_tolerance: 0.0,
158            timestamps: true,
159            numeric_strings: false,
160        }
161    }
162}
163
164impl Normalizer {
165    /// Fail-fast validation: a finite, non-negative tolerance.
166    pub fn validate(&self) -> Result<(), crate::FaucetError> {
167        if !self.float_tolerance.is_finite() || self.float_tolerance < 0.0 {
168            return Err(crate::FaucetError::Config(format!(
169                "verify: float_tolerance must be a finite number >= 0, got {}",
170                self.float_tolerance
171            )));
172        }
173        Ok(())
174    }
175
176    /// The canonical text of one value. Missing and `null` both canonicalise to
177    /// the same token, so a column absent on one side equals `null` on the other.
178    pub fn canonical(&self, v: Option<&Value>) -> String {
179        match v {
180            None | Some(Value::Null) => "\u{0}null".into(),
181            Some(Value::Bool(b)) => (if *b { "t" } else { "f" }).into(),
182            Some(Value::Number(n)) => self.canonical_number(n.as_f64(), n.as_i64(), n.as_u64()),
183            Some(Value::String(s)) => self.canonical_string(s),
184            Some(Value::Array(a)) => {
185                let parts: Vec<String> = a.iter().map(|x| self.canonical(Some(x))).collect();
186                format!("[{}]", parts.join(","))
187            }
188            Some(Value::Object(o)) => {
189                // Sorted keys: object key order is not content.
190                let sorted: BTreeMap<&String, &Value> = o.iter().collect();
191                let parts: Vec<String> = sorted
192                    .into_iter()
193                    .map(|(k, x)| format!("{}:{}", k, self.canonical(Some(x))))
194                    .collect();
195                format!("{{{}}}", parts.join(","))
196            }
197        }
198    }
199
200    fn canonical_number(&self, f: Option<f64>, i: Option<i64>, u: Option<u64>) -> String {
201        if let Some(i) = i {
202            return format!("n{i}");
203        }
204        if let Some(u) = u {
205            return format!("n{u}");
206        }
207        let Some(f) = f else {
208            return "n?".into();
209        };
210        // An integral float is the same number as the integer (`1.0` == `1`).
211        if f.fract() == 0.0 && f.abs() < 9.0e15 {
212            return format!("n{}", f as i64);
213        }
214        if self.float_tolerance > 0.0 {
215            let rounded = (f / self.float_tolerance).round() * self.float_tolerance;
216            return format!("f{rounded:e}");
217        }
218        format!("f{f:e}")
219    }
220
221    fn canonical_string(&self, s: &str) -> String {
222        if self.timestamps
223            && let Some(micros) = parse_timestamp_micros(s)
224        {
225            return format!("ts{micros}");
226        }
227        if self.numeric_strings {
228            let t = s.trim();
229            if let Ok(i) = t.parse::<i64>() {
230                return format!("n{i}");
231            }
232            if let Ok(f) = t.parse::<f64>()
233                && f.is_finite()
234            {
235                return self.canonical_number(Some(f), None, None);
236            }
237        }
238        format!("s{s}")
239    }
240}
241
242/// UTC microseconds for a timestamp string in RFC 3339 form or the
243/// `YYYY-MM-DD HH:MM:SS[.fff][Z|±hh:mm]` form SQL backends emit (a naive
244/// timestamp is read as UTC). `None` when the string is not a timestamp.
245pub fn parse_timestamp_micros(s: &str) -> Option<i64> {
246    let s = s.trim();
247    // Cheap pre-check: every timestamp we accept starts with a 4-digit year and
248    // a dash, so ordinary strings never pay for the parse attempts.
249    let b = s.as_bytes();
250    if b.len() < 10 || !b[..4].iter().all(u8::is_ascii_digit) || b[4] != b'-' {
251        return None;
252    }
253    if let Ok(dt) = chrono::DateTime::parse_from_rfc3339(s) {
254        return Some(dt.timestamp_micros());
255    }
256    let with_space = s.replacen(' ', "T", 1);
257    if let Ok(dt) = chrono::DateTime::parse_from_rfc3339(&with_space) {
258        return Some(dt.timestamp_micros());
259    }
260    if let Ok(naive) = chrono::NaiveDateTime::parse_from_str(&with_space, "%Y-%m-%dT%H:%M:%S%.f") {
261        return Some(naive.and_utc().timestamp_micros());
262    }
263    if let Ok(date) = chrono::NaiveDate::parse_from_str(s, "%Y-%m-%d") {
264        return Some(date.and_hms_opt(0, 0, 0)?.and_utc().timestamp_micros());
265    }
266    None
267}
268
269/// FNV-1a over `bytes` with a caller-chosen offset basis, so two independent
270/// 64-bit hashes of the same text can be folded into one 128-bit fingerprint.
271fn fnv1a_64(bytes: &[u8], basis: u64) -> u64 {
272    let mut h = basis;
273    for b in bytes {
274        h ^= u64::from(*b);
275        h = h.wrapping_mul(0x0000_0100_0000_01b3);
276    }
277    h
278}
279
280const FNV_BASIS_A: u64 = 0xcbf2_9ce4_8422_2325;
281const FNV_BASIS_B: u64 = 0x84222325_cbf29ce4;
282
283/// A 128-bit fingerprint of `record`'s compared columns: the key columns plus
284/// `columns` (every non-excluded field when `columns` is `None`), each
285/// canonicalised through `norm`. Column order does not matter — the fields are
286/// hashed in sorted name order.
287pub fn row_hash(
288    record: &Value,
289    key: &[String],
290    columns: Option<&[String]>,
291    exclude: &[String],
292    norm: &Normalizer,
293) -> u128 {
294    let obj = record.as_object();
295    let mut text = String::new();
296    let mut names: Vec<&str> = match columns {
297        Some(cols) => cols.iter().map(String::as_str).collect(),
298        None => obj
299            .map(|o| o.keys().map(String::as_str).collect())
300            .unwrap_or_default(),
301    };
302    for k in key {
303        if !names.contains(&k.as_str()) {
304            names.push(k.as_str());
305        }
306    }
307    names.retain(|n| !is_excluded(n, exclude));
308    names.sort_unstable();
309    names.dedup();
310    for name in names {
311        text.push_str(name);
312        text.push('=');
313        text.push_str(&norm.canonical(obj.and_then(|o| o.get(name))));
314        text.push('\u{1f}');
315    }
316    let a = fnv1a_64(text.as_bytes(), FNV_BASIS_A);
317    let b = fnv1a_64(text.as_bytes(), FNV_BASIS_B);
318    (u128::from(a) << 64) | u128::from(b)
319}
320
321/// Whether a column is excluded: an exact name, or a `prefix*` glob.
322pub fn is_excluded(name: &str, exclude: &[String]) -> bool {
323    exclude.iter().any(|e| match e.strip_suffix('*') {
324        Some(prefix) => name.starts_with(prefix),
325        None => e == name,
326    })
327}
328
329/// The canonical text of a record's key tuple (`k1=…;k2=…`), used to match
330/// rows across the two sides.
331pub fn key_text(record: &Value, key: &[String], norm: &Normalizer) -> String {
332    let obj = record.as_object();
333    let mut out = String::new();
334    for k in key {
335        out.push_str(k);
336        out.push('=');
337        out.push_str(&norm.canonical(obj.and_then(|o| o.get(k))));
338        out.push(';');
339    }
340    out
341}
342
343/// The integer value of a single-column key, when it is one.
344pub fn key_int(record: &Value, key: &[String]) -> Option<i64> {
345    if key.len() != 1 {
346        return None;
347    }
348    match record.get(&key[0])? {
349        Value::Number(n) => n.as_i64(),
350        Value::String(s) => s.trim().parse().ok(),
351        _ => None,
352    }
353}
354
355/// An order-independent digest of a set of rows.
356#[derive(Debug, Clone, Copy, PartialEq, Eq, Default, Serialize, Deserialize)]
357pub struct ContentDigest {
358    pub rows: u64,
359    /// Wrapping sum of the row fingerprints.
360    pub sum: u128,
361    /// XOR of the row fingerprints (catches a sum collision from swapped rows).
362    pub xor: u128,
363    /// Smallest / largest integer key seen (single integer keys only).
364    pub key_min: Option<i64>,
365    pub key_max: Option<i64>,
366}
367
368impl ContentDigest {
369    /// Whether two digests describe the same row set.
370    pub fn same(&self, other: &ContentDigest) -> bool {
371        self.rows == other.rows && self.sum == other.sum && self.xor == other.xor
372    }
373
374    /// The union of the observed key bounds of two digests.
375    pub fn bounds_union(&self, other: &ContentDigest) -> Option<(i64, i64)> {
376        let min = match (self.key_min, other.key_min) {
377            (Some(a), Some(b)) => Some(a.min(b)),
378            (a, b) => a.or(b),
379        }?;
380        let max = match (self.key_max, other.key_max) {
381            (Some(a), Some(b)) => Some(a.max(b)),
382            (a, b) => a.or(b),
383        }?;
384        Some((min, max))
385    }
386}
387
388/// Streams rows into a [`ContentDigest`].
389#[derive(Debug, Default)]
390pub struct DigestAccumulator {
391    digest: ContentDigest,
392}
393
394impl DigestAccumulator {
395    pub fn new() -> Self {
396        Self::default()
397    }
398
399    /// Fold one row in.
400    pub fn add(&mut self, hash: u128, key: Option<i64>) {
401        self.digest.rows += 1;
402        self.digest.sum = self.digest.sum.wrapping_add(hash);
403        self.digest.xor ^= hash;
404        if let Some(k) = key {
405            self.digest.key_min = Some(self.digest.key_min.map_or(k, |m| m.min(k)));
406            self.digest.key_max = Some(self.digest.key_max.map_or(k, |m| m.max(k)));
407        }
408    }
409
410    pub fn finish(self) -> ContentDigest {
411        self.digest
412    }
413}
414
415/// A digest a backend computed itself, without shipping rows. Two server
416/// digests are comparable only when their `algorithm` ids match — each backend
417/// hashes its own text rendering of a row, so a Postgres digest says nothing
418/// about a SQLite table. When the algorithms differ the verifier streams both
419/// sides and uses [`ContentDigest`] instead.
420#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
421pub struct ServerDigest {
422    pub algorithm: String,
423    pub rows: u64,
424    /// Opaque digest text (a decimal sum, a hex string — whatever the
425    /// algorithm defines).
426    pub digest: String,
427    pub key_min: Option<i64>,
428    pub key_max: Option<i64>,
429}
430
431impl ServerDigest {
432    /// Whether two server digests agree. Different algorithms never agree —
433    /// the caller must not treat that as a difference but as "not comparable".
434    pub fn comparable(&self, other: &ServerDigest) -> bool {
435        self.algorithm == other.algorithm
436    }
437
438    pub fn same(&self, other: &ServerDigest) -> bool {
439        self.comparable(other) && self.rows == other.rows && self.digest == other.digest
440    }
441
442    pub fn bounds_union(&self, other: &ServerDigest) -> Option<(i64, i64)> {
443        let min = match (self.key_min, other.key_min) {
444            (Some(a), Some(b)) => Some(a.min(b)),
445            (a, b) => a.or(b),
446        }?;
447        let max = match (self.key_max, other.key_max) {
448            (Some(a), Some(b)) => Some(a.max(b)),
449            (a, b) => a.or(b),
450        }?;
451        Some((min, max))
452    }
453}
454
455/// Why one key differs between the two sides.
456#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
457#[serde(tag = "kind", rename_all = "snake_case")]
458pub enum DifferenceKind {
459    /// The source has the key, the destination does not.
460    MissingInDest,
461    /// The destination has the key, the source does not.
462    ExtraInDest,
463    /// Both have the key; these columns differ.
464    Changed { columns: Vec<String> },
465    /// The key appears more than once on one side, so rows cannot be matched.
466    Duplicate { side: Side, count: u64 },
467}
468
469#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
470#[serde(rename_all = "snake_case")]
471pub enum Side {
472    Source,
473    Destination,
474}
475
476/// One differing key.
477#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
478pub struct Difference {
479    /// The key columns and their values (from whichever side has the row).
480    pub key: Value,
481    #[serde(flatten)]
482    pub kind: DifferenceKind,
483}
484
485impl Difference {
486    /// Whether repairing this difference means writing the source row.
487    pub fn needs_upsert(&self) -> bool {
488        matches!(
489            self.kind,
490            DifferenceKind::MissingInDest | DifferenceKind::Changed { .. }
491        )
492    }
493
494    /// Whether repairing this difference means deleting the destination row.
495    pub fn needs_delete(&self) -> bool {
496        matches!(self.kind, DifferenceKind::ExtraInDest)
497    }
498}
499
500/// The columns whose canonical values differ between two rows.
501fn changed_columns(
502    a: &Value,
503    b: &Value,
504    columns: Option<&[String]>,
505    exclude: &[String],
506    norm: &Normalizer,
507) -> Vec<String> {
508    let names: Vec<String> = match columns {
509        Some(cols) => cols.to_vec(),
510        None => {
511            let mut set: Vec<String> = Vec::new();
512            for side in [a, b] {
513                if let Some(o) = side.as_object() {
514                    for k in o.keys() {
515                        if !set.contains(k) {
516                            set.push(k.clone());
517                        }
518                    }
519                }
520            }
521            set
522        }
523    };
524    let mut out: Vec<String> = names
525        .into_iter()
526        .filter(|n| !is_excluded(n, exclude))
527        .filter(|n| norm.canonical(a.get(n.as_str())) != norm.canonical(b.get(n.as_str())))
528        .collect();
529    out.sort();
530    out
531}
532
533fn key_object(record: &Value, key: &[String]) -> Value {
534    let mut m = Map::new();
535    for k in key {
536        m.insert(k.clone(), record.get(k).cloned().unwrap_or(Value::Null));
537    }
538    Value::Object(m)
539}
540
541/// Compare two leaf ranges row by row. Returns every differing key, ordered by
542/// key text so the report is stable. Rows on either side are matched on the
543/// canonical key text; a key present more than once on a side is reported as a
544/// [`DifferenceKind::Duplicate`] and not matched.
545pub fn diff_rows(
546    source: &[Value],
547    dest: &[Value],
548    key: &[String],
549    columns: Option<&[String]>,
550    exclude: &[String],
551    norm: &Normalizer,
552) -> Vec<Difference> {
553    let mut by_key: HashMap<String, (u64, &Value)> = HashMap::with_capacity(dest.len());
554    for d in dest {
555        let e = by_key.entry(key_text(d, key, norm)).or_insert((0, d));
556        e.0 += 1;
557    }
558    // A key repeated on the source side cannot be matched row-for-row: it is
559    // reported once as a duplicate and takes no part in the comparison.
560    let mut source_counts: HashMap<String, u64> = HashMap::with_capacity(source.len());
561    for s in source {
562        *source_counts.entry(key_text(s, key, norm)).or_insert(0) += 1;
563    }
564    let mut out: Vec<(String, Difference)> = Vec::new();
565    let mut reported: std::collections::HashSet<String> = std::collections::HashSet::new();
566    for s in source {
567        let kt = key_text(s, key, norm);
568        let count = source_counts[&kt];
569        if count > 1 {
570            if reported.insert(kt.clone()) {
571                by_key.remove(&kt);
572                out.push((
573                    kt,
574                    Difference {
575                        key: key_object(s, key),
576                        kind: DifferenceKind::Duplicate {
577                            side: Side::Source,
578                            count,
579                        },
580                    },
581                ));
582            }
583            continue;
584        }
585        match by_key.remove(&kt) {
586            None => out.push((
587                kt,
588                Difference {
589                    key: key_object(s, key),
590                    kind: DifferenceKind::MissingInDest,
591                },
592            )),
593            Some((count, _)) if count > 1 => out.push((
594                kt,
595                Difference {
596                    key: key_object(s, key),
597                    kind: DifferenceKind::Duplicate {
598                        side: Side::Destination,
599                        count,
600                    },
601                },
602            )),
603            Some((_, d)) => {
604                let cols = changed_columns(s, d, columns, exclude, norm);
605                if !cols.is_empty() {
606                    out.push((
607                        kt,
608                        Difference {
609                            key: key_object(s, key),
610                            kind: DifferenceKind::Changed { columns: cols },
611                        },
612                    ));
613                }
614            }
615        }
616    }
617    for (kt, (count, d)) in by_key {
618        out.push((
619            kt,
620            Difference {
621                key: key_object(d, key),
622                kind: if count > 1 {
623                    DifferenceKind::Duplicate {
624                        side: Side::Destination,
625                        count,
626                    }
627                } else {
628                    DifferenceKind::ExtraInDest
629                },
630            },
631        ));
632    }
633    out.sort_by(|a, b| a.0.cmp(&b.0));
634    out.into_iter().map(|(_, d)| d).collect()
635}
636
637/// Summary of one verification, serialisable for `--json` and the HTTP API.
638#[derive(Debug, Clone, Default, Serialize, Deserialize)]
639pub struct VerifyReport {
640    /// Ranges compared by digest (both passes).
641    pub ranges_compared: u64,
642    /// Ranges whose digests disagreed and were bisected/fetched.
643    pub ranges_differing: u64,
644    /// Rows the source side produced across the leaf fetches.
645    pub rows_fetched_source: u64,
646    /// Rows the destination side produced across the leaf fetches.
647    pub rows_fetched_dest: u64,
648    /// Whether digests were computed inside the backends (no rows shipped for
649    /// matching ranges) rather than by streaming.
650    pub server_digests: bool,
651    pub differences: Vec<Difference>,
652    /// The scan stopped at `max_rows_scanned`; `differences` is a lower bound.
653    pub truncated: bool,
654    /// Rows written back by `--repair` (upserts), if repair ran.
655    #[serde(skip_serializing_if = "Option::is_none")]
656    pub repaired_upserts: Option<u64>,
657    /// Rows deleted by `--repair --allow-delete`, if repair ran.
658    #[serde(skip_serializing_if = "Option::is_none")]
659    pub repaired_deletes: Option<u64>,
660}
661
662impl VerifyReport {
663    /// Counts per difference kind, for the human summary.
664    pub fn tally(&self) -> (usize, usize, usize, usize) {
665        let mut missing = 0;
666        let mut extra = 0;
667        let mut changed = 0;
668        let mut dup = 0;
669        for d in &self.differences {
670            match d.kind {
671                DifferenceKind::MissingInDest => missing += 1,
672                DifferenceKind::ExtraInDest => extra += 1,
673                DifferenceKind::Changed { .. } => changed += 1,
674                DifferenceKind::Duplicate { .. } => dup += 1,
675            }
676        }
677        (missing, extra, changed, dup)
678    }
679
680    /// Whether the two sides matched (no differences and nothing skipped).
681    pub fn equal(&self) -> bool {
682        self.differences.is_empty() && !self.truncated
683    }
684
685    /// Whether a repair re-synced every reported difference (and the report is
686    /// complete), so the destination now matches.
687    pub fn healed(&self) -> bool {
688        if self.truncated || self.differences.is_empty() {
689            return false;
690        }
691        let repaired = self.repaired_upserts.unwrap_or(0) + self.repaired_deletes.unwrap_or(0);
692        repaired >= self.differences.len() as u64
693    }
694}
695
696#[cfg(test)]
697mod tests {
698    use super::*;
699    use serde_json::json;
700
701    fn norm() -> Normalizer {
702        Normalizer::default()
703    }
704
705    #[test]
706    fn ranges_tile_the_key_space_and_bisect_to_leaves() {
707        let ranges = plan_ranges(1, 100, 4);
708        assert_eq!(ranges.len(), 4);
709        assert_eq!(ranges[0].lo, None, "first is unbounded below");
710        assert_eq!(ranges[3].hi, None, "last is unbounded above");
711        for k in [-5, 1, 50, 100, 1_000] {
712            assert_eq!(ranges.iter().filter(|r| r.contains(k)).count(), 1, "{k}");
713        }
714        let (a, b) = KeyRange {
715            lo: Some(0),
716            hi: Some(10),
717        }
718        .bisect(None)
719        .unwrap();
720        assert_eq!(
721            (a.lo, a.hi, b.lo, b.hi),
722            (Some(0), Some(5), Some(5), Some(10))
723        );
724        // An unbounded side clamps to the observed bounds.
725        let (a, b) = KeyRange::ALL.bisect(Some((10, 19))).unwrap();
726        assert_eq!((a.lo, a.hi), (None, Some(15)));
727        assert_eq!((b.lo, b.hi), (Some(15), None));
728        assert!(KeyRange::ALL.bisect(None).is_none());
729        assert!(
730            KeyRange {
731                lo: Some(3),
732                hi: Some(4)
733            }
734            .bisect(None)
735            .is_none(),
736            "one key wide is a leaf"
737        );
738        assert_eq!(
739            KeyRange {
740                lo: Some(3),
741                hi: Some(4)
742            }
743            .width(),
744            Some(1)
745        );
746        assert_eq!(KeyRange::ALL.width(), None);
747        assert_eq!(KeyRange::ALL.to_string(), "all");
748        assert_eq!(
749            KeyRange {
750                lo: Some(1),
751                hi: None
752            }
753            .to_string(),
754            "[1, +inf)"
755        );
756        assert_eq!(
757            KeyRange {
758                lo: None,
759                hi: Some(9)
760            }
761            .to_string(),
762            "(-inf, 9)"
763        );
764    }
765
766    #[test]
767    fn range_shard_is_the_pk_range_shape_sql_sources_parse() {
768        let shard = KeyRange {
769            lo: Some(5),
770            hi: None,
771        }
772        .to_shard("id");
773        let bounds = crate::shard::PkShardBounds::from_spec(&shard).unwrap();
774        assert_eq!(
775            (bounds.lo, bounds.hi_unbounded, bounds.include_null),
776            (5, true, true)
777        );
778        assert_eq!(shard.id, "5..+inf");
779        let sql = bounds.wrap("SELECT * FROM t", |s| format!("\"{s}\""));
780        assert!(
781            sql.contains("\"id\" >= 5") && sql.contains("IS NULL"),
782            "{sql}"
783        );
784    }
785
786    #[test]
787    fn canonical_values_agree_across_spellings() {
788        let n = norm();
789        assert_eq!(n.canonical(Some(&json!(1))), n.canonical(Some(&json!(1.0))));
790        assert_eq!(n.canonical(None), n.canonical(Some(&json!(null))));
791        assert_ne!(n.canonical(Some(&json!("1"))), n.canonical(Some(&json!(1))));
792        assert_eq!(
793            n.canonical(Some(&json!("2026-09-25T10:00:00Z"))),
794            n.canonical(Some(&json!("2026-09-25 12:00:00+02:00")))
795        );
796        assert_eq!(
797            n.canonical(Some(&json!("2026-09-25 10:00:00.000"))),
798            n.canonical(Some(&json!("2026-09-25T10:00:00Z")))
799        );
800        assert_eq!(
801            n.canonical(Some(&json!("2026-09-25"))),
802            n.canonical(Some(&json!("2026-09-25T00:00:00Z")))
803        );
804        assert_eq!(
805            n.canonical(Some(&json!({"b": 1, "a": [true, null]}))),
806            n.canonical(Some(&json!({"a": [true, null], "b": 1})))
807        );
808        // Plain strings never pay for a timestamp parse and stay themselves.
809        assert_eq!(n.canonical(Some(&json!("hello"))), "shello");
810        assert_eq!(n.canonical(Some(&json!("2026-x"))), "s2026-x");
811        assert!(parse_timestamp_micros("not a date").is_none());
812    }
813
814    #[test]
815    fn float_tolerance_and_numeric_strings_are_opt_in() {
816        let exact = norm();
817        assert_ne!(
818            exact.canonical(Some(&json!(1.00000001))),
819            exact.canonical(Some(&json!(1.00000002)))
820        );
821        let loose = Normalizer {
822            float_tolerance: 1e-6,
823            ..norm()
824        };
825        assert_eq!(
826            loose.canonical(Some(&json!(1.00000001))),
827            loose.canonical(Some(&json!(1.00000002)))
828        );
829        let strings = Normalizer {
830            numeric_strings: true,
831            ..norm()
832        };
833        assert_eq!(
834            strings.canonical(Some(&json!("42"))),
835            strings.canonical(Some(&json!(42)))
836        );
837        assert_eq!(
838            strings.canonical(Some(&json!("1.5"))),
839            strings.canonical(Some(&json!(1.5)))
840        );
841        assert_eq!(strings.canonical(Some(&json!("abc"))), "sabc");
842        assert!(loose.validate().is_ok());
843        assert!(
844            Normalizer {
845                float_tolerance: -1.0,
846                ..norm()
847            }
848            .validate()
849            .is_err()
850        );
851        assert!(
852            Normalizer {
853                float_tolerance: f64::NAN,
854                ..norm()
855            }
856            .validate()
857            .is_err()
858        );
859    }
860
861    #[test]
862    fn row_hash_ignores_field_order_and_excluded_columns() {
863        let key = vec!["id".to_string()];
864        let a = json!({"id": 1, "name": "x", "_faucet_run_id": "r1"});
865        let b = json!({"name": "x", "id": 1, "_faucet_run_id": "r2"});
866        let ex = vec!["_faucet_*".to_string()];
867        assert_eq!(
868            row_hash(&a, &key, None, &ex, &norm()),
869            row_hash(&b, &key, None, &ex, &norm())
870        );
871        assert_ne!(
872            row_hash(&a, &key, None, &[], &norm()),
873            row_hash(&b, &key, None, &[], &norm()),
874            "without the exclusion the run id differs"
875        );
876        // An explicit column list hashes only those (plus the key).
877        let c = json!({"id": 1, "name": "y", "other": 5});
878        let cols = vec!["other".to_string()];
879        assert_eq!(
880            row_hash(&a, &key, Some(&cols), &[], &norm()),
881            row_hash(
882                &json!({"id": 1, "name": "zzz"}),
883                &key,
884                Some(&cols),
885                &[],
886                &norm()
887            )
888        );
889        assert_ne!(
890            row_hash(&c, &key, Some(&cols), &[], &norm()),
891            row_hash(&a, &key, Some(&cols), &[], &norm())
892        );
893        assert!(is_excluded("_faucet_loaded_at", &ex));
894        assert!(!is_excluded("faucet", &ex));
895        assert!(is_excluded("name", &["name".to_string()]));
896        assert_eq!(key_int(&a, &key), Some(1));
897        assert_eq!(key_int(&json!({"id": "7"}), &key), Some(7));
898        assert_eq!(key_int(&a, &["id".into(), "name".into()]), None);
899        assert_eq!(key_int(&json!({"id": true}), &key), None);
900        assert_eq!(key_text(&a, &key, &norm()), "id=n1;");
901    }
902
903    #[test]
904    fn digest_is_order_independent_and_tracks_bounds() {
905        let key = vec!["id".to_string()];
906        let rows = [json!({"id": 3, "v": "a"}), json!({"id": 1, "v": "b"})];
907        let mut fwd = DigestAccumulator::new();
908        let mut rev = DigestAccumulator::new();
909        for r in &rows {
910            fwd.add(row_hash(r, &key, None, &[], &norm()), key_int(r, &key));
911        }
912        for r in rows.iter().rev() {
913            rev.add(row_hash(r, &key, None, &[], &norm()), key_int(r, &key));
914        }
915        let (f, r) = (fwd.finish(), rev.finish());
916        assert!(f.same(&r));
917        assert_eq!((f.rows, f.key_min, f.key_max), (2, Some(1), Some(3)));
918        let mut other = DigestAccumulator::new();
919        other.add(
920            row_hash(&json!({"id": 3, "v": "a"}), &key, None, &[], &norm()),
921            Some(3),
922        );
923        let o = other.finish();
924        assert!(!f.same(&o));
925        assert_eq!(f.bounds_union(&o), Some((1, 3)));
926        assert_eq!(
927            ContentDigest::default().bounds_union(&ContentDigest::default()),
928            None
929        );
930        assert_eq!(ContentDigest::default().bounds_union(&o), Some((3, 3)));
931    }
932
933    #[test]
934    fn server_digests_compare_only_within_one_algorithm() {
935        let pg = ServerDigest {
936            algorithm: "postgres:v1".into(),
937            rows: 2,
938            digest: "12".into(),
939            key_min: Some(1),
940            key_max: Some(5),
941        };
942        let same = ServerDigest {
943            key_min: Some(2),
944            key_max: Some(9),
945            ..pg.clone()
946        };
947        let other_algo = ServerDigest {
948            algorithm: "sqlite:v1".into(),
949            ..pg.clone()
950        };
951        assert!(pg.same(&same));
952        assert!(!pg.comparable(&other_algo));
953        assert!(!pg.same(&other_algo));
954        assert!(!pg.same(&ServerDigest {
955            digest: "13".into(),
956            ..pg.clone()
957        }));
958        assert_eq!(pg.bounds_union(&same), Some((1, 9)));
959        assert_eq!(
960            pg.bounds_union(&ServerDigest {
961                key_min: None,
962                key_max: None,
963                ..pg.clone()
964            }),
965            Some((1, 5))
966        );
967    }
968
969    #[test]
970    fn diff_rows_reports_every_kind_in_key_order() {
971        let key = vec!["id".to_string()];
972        let source = vec![
973            json!({"id": 1, "v": "a"}),
974            json!({"id": 2, "v": "b"}),
975            json!({"id": 3, "v": "c"}),
976            json!({"id": 5, "v": "e"}),
977            json!({"id": 5, "v": "e2"}),
978        ];
979        let dest = vec![
980            json!({"id": 1, "v": "a", "_faucet_run_id": "r"}),
981            json!({"id": 2, "v": "B"}),
982            json!({"id": 4, "v": "d"}),
983            json!({"id": 6, "v": "f"}),
984            json!({"id": 6, "v": "f"}),
985        ];
986        let ex = vec!["_faucet_*".to_string()];
987        let diffs = diff_rows(&source, &dest, &key, None, &ex, &norm());
988        let kinds: Vec<String> = diffs
989            .iter()
990            .map(|d| format!("{}:{:?}", d.key["id"], d.kind))
991            .collect();
992        assert_eq!(diffs.len(), 5, "{kinds:?}");
993        assert_eq!(
994            diffs[0].kind,
995            DifferenceKind::Changed {
996                columns: vec!["v".into()]
997            }
998        );
999        assert_eq!(diffs[0].key, json!({"id": 2}));
1000        assert_eq!(diffs[1].kind, DifferenceKind::MissingInDest);
1001        assert_eq!(diffs[1].key, json!({"id": 3}));
1002        assert_eq!(diffs[2].kind, DifferenceKind::ExtraInDest);
1003        assert_eq!(diffs[2].key, json!({"id": 4}));
1004        assert_eq!(
1005            diffs[3].kind,
1006            DifferenceKind::Duplicate {
1007                side: Side::Source,
1008                count: 2
1009            }
1010        );
1011        assert_eq!(
1012            diffs[4].kind,
1013            DifferenceKind::Duplicate {
1014                side: Side::Destination,
1015                count: 2
1016            }
1017        );
1018        assert!(diffs[0].needs_upsert() && diffs[1].needs_upsert());
1019        assert!(diffs[2].needs_delete() && !diffs[2].needs_upsert());
1020        assert!(!diffs[3].needs_upsert() && !diffs[3].needs_delete());
1021        // Identical sides: nothing.
1022        assert!(diff_rows(&source[..1], &dest[..1], &key, None, &ex, &norm()).is_empty());
1023        // A key repeated on both sides is one report, on the source side, with
1024        // the source's count; it is never also "missing" or "extra".
1025        let three = vec![json!({"id": 9}), json!({"id": 9}), json!({"id": 9})];
1026        let d = diff_rows(
1027            &three,
1028            &[json!({"id": 9}), json!({"id": 9})],
1029            &key,
1030            None,
1031            &[],
1032            &norm(),
1033        );
1034        assert_eq!(d.len(), 1, "a duplicated key is reported once: {d:?}");
1035        assert_eq!(
1036            d[0].kind,
1037            DifferenceKind::Duplicate {
1038                side: Side::Source,
1039                count: 3
1040            }
1041        );
1042    }
1043
1044    #[test]
1045    fn diff_rows_with_explicit_columns_ignores_the_rest() {
1046        let key = vec!["id".to_string()];
1047        let cols = vec!["amount".to_string()];
1048        let s = vec![json!({"id": 1, "amount": 10, "note": "x"})];
1049        let d = vec![json!({"id": 1, "amount": 10, "note": "y"})];
1050        assert!(diff_rows(&s, &d, &key, Some(&cols), &[], &norm()).is_empty());
1051        let d2 = vec![json!({"id": 1, "amount": 11})];
1052        let diffs = diff_rows(&s, &d2, &key, Some(&cols), &[], &norm());
1053        assert_eq!(
1054            diffs[0].kind,
1055            DifferenceKind::Changed {
1056                columns: vec!["amount".into()]
1057            }
1058        );
1059    }
1060
1061    #[test]
1062    fn range_width_display_and_number_canonical_forms() {
1063        let inverted = KeyRange {
1064            lo: Some(5),
1065            hi: Some(5),
1066        };
1067        assert_eq!(inverted.width(), Some(0));
1068        assert_eq!(format!("{inverted}"), "[5, 5)");
1069        assert_eq!(format!("{}", KeyRange::ALL), "all");
1070        let norm = Normalizer::default();
1071        assert_eq!(
1072            norm.canonical(Some(&json!(u64::MAX))),
1073            format!("n{}", u64::MAX)
1074        );
1075        assert_eq!(norm.canonical(Some(&json!(-3))), "n-3");
1076        assert_eq!(norm.canonical(Some(&json!(2.0))), "n2");
1077    }
1078
1079    #[test]
1080    fn timestamps_parse_every_supported_spelling_or_none() {
1081        assert!(parse_timestamp_micros("2026-01-01T00:00:00Z").is_some());
1082        assert!(parse_timestamp_micros("2026-01-01 00:00:00.5").is_some());
1083        assert!(parse_timestamp_micros("2026-01-01 00:00:00+02:00").is_some());
1084        assert!(parse_timestamp_micros("2026-01-01").is_some());
1085        assert_eq!(parse_timestamp_micros("2026-13-99"), None);
1086        assert_eq!(parse_timestamp_micros("not a date"), None);
1087        assert_eq!(parse_timestamp_micros("20260101"), None);
1088    }
1089
1090    #[test]
1091    fn bounds_union_needs_both_ends() {
1092        let a = ContentDigest {
1093            key_min: Some(1),
1094            ..Default::default()
1095        };
1096        let b = ContentDigest::default();
1097        assert_eq!(a.bounds_union(&b), None, "no max anywhere");
1098        let sd = ServerDigest {
1099            algorithm: "x".into(),
1100            rows: 0,
1101            digest: "0".into(),
1102            key_min: None,
1103            key_max: Some(9),
1104        };
1105        assert_eq!(sd.bounds_union(&sd), None, "no min anywhere");
1106        let other = ServerDigest {
1107            key_min: Some(2),
1108            ..sd.clone()
1109        };
1110        assert_eq!(sd.bounds_union(&other), Some((2, 9)));
1111    }
1112
1113    #[test]
1114    fn duplicate_destination_keys_are_reported_once_on_that_side() {
1115        let source = vec![json!({"id": 1, "v": "a"})];
1116        let dest = vec![json!({"id": 1, "v": "a"}), json!({"id": 1, "v": "b"})];
1117        let d = diff_rows(
1118            &source,
1119            &dest,
1120            &["id".into()],
1121            None,
1122            &[],
1123            &Normalizer::default(),
1124        );
1125        assert_eq!(d.len(), 1, "{d:?}");
1126        assert_eq!(
1127            d[0].kind,
1128            DifferenceKind::Duplicate {
1129                side: Side::Destination,
1130                count: 2
1131            }
1132        );
1133        assert!(!d[0].needs_upsert() && !d[0].needs_delete());
1134    }
1135
1136    #[test]
1137    fn report_healed_only_when_every_difference_was_repaired() {
1138        let mut r = VerifyReport::default();
1139        assert!(!r.healed(), "nothing to heal");
1140        r.differences.push(Difference {
1141            key: json!({"id": 1}),
1142            kind: DifferenceKind::MissingInDest,
1143        });
1144        r.differences.push(Difference {
1145            key: json!({"id": 2}),
1146            kind: DifferenceKind::ExtraInDest,
1147        });
1148        assert!(!r.healed(), "no repair ran");
1149        r.repaired_upserts = Some(1);
1150        r.repaired_deletes = Some(0);
1151        assert!(!r.healed(), "the extra row was not deleted");
1152        r.repaired_deletes = Some(1);
1153        assert!(r.healed());
1154        r.truncated = true;
1155        assert!(!r.healed(), "a truncated report cannot claim to be healed");
1156    }
1157
1158    #[test]
1159    fn report_tallies_and_equality() {
1160        let mut r = VerifyReport::default();
1161        assert!(r.equal());
1162        r.differences.push(Difference {
1163            key: json!({"id": 1}),
1164            kind: DifferenceKind::MissingInDest,
1165        });
1166        r.differences.push(Difference {
1167            key: json!({"id": 2}),
1168            kind: DifferenceKind::ExtraInDest,
1169        });
1170        r.differences.push(Difference {
1171            key: json!({"id": 3}),
1172            kind: DifferenceKind::Changed { columns: vec![] },
1173        });
1174        r.differences.push(Difference {
1175            key: json!({"id": 4}),
1176            kind: DifferenceKind::Duplicate {
1177                side: Side::Source,
1178                count: 2,
1179            },
1180        });
1181        assert_eq!(r.tally(), (1, 1, 1, 1));
1182        assert!(!r.equal());
1183        let truncated = VerifyReport {
1184            truncated: true,
1185            ..Default::default()
1186        };
1187        assert!(!truncated.equal(), "a truncated scan proves nothing");
1188        let json = serde_json::to_value(&r).unwrap();
1189        assert_eq!(json["differences"][0]["kind"], "missing_in_dest");
1190        assert_eq!(json["differences"][2]["kind"], "changed");
1191    }
1192}