Skip to main content

ahu/
ops.rs

1//! Pure bundle operations shared by the CLI and admin API: structural diff and
2//! delta-chain application. These functions take parsed `Bundle`s and return
3//! data structures — no printing, no filesystem, no `std::process::exit`. The
4//! CLI formats the result for a terminal; the server serializes it to JSON.
5
6use std::collections::{BTreeMap, HashSet};
7
8#[cfg(test)]
9use sha2::{Digest, Sha256};
10
11use crate::bundle::{Bundle, BundleBuilder};
12use crate::error::{AhuError, Result};
13use crate::index::{self, IndexFlags};
14use crate::manifest::BundleType;
15
16/// A single entry identified by its key plus algorithm discriminator. Dual-algorithm
17/// bundles hold the same `entry_key` under different discriminators, so both fields
18/// are needed to name an entry uniquely.
19#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
20pub struct EntryRef {
21    pub entry_key: [u8; 32],
22    pub discriminator: u16,
23}
24
25/// Structural difference between two bundles (A → B), computed over
26/// `(entry_key, discriminator)` pairs. "Changed" means the key exists in both
27/// but the response bytes differ.
28#[derive(Debug, Clone)]
29pub struct DiffResult {
30    pub a_entry_count: usize,
31    pub b_entry_count: usize,
32    pub a_epochs: Vec<u64>,
33    pub b_epochs: Vec<u64>,
34    pub added: Vec<EntryRef>,
35    pub removed: Vec<EntryRef>,
36    pub changed: Vec<EntryRef>,
37    pub unchanged: usize,
38}
39
40/// Compute the structural diff of two bundles. Never fails: it only reads the
41/// already-parsed index and data sections.
42pub fn diff(a: &Bundle, b: &Bundle) -> DiffResult {
43    let a_keys: HashSet<([u8; 32], u16)> = a
44        .index
45        .iter()
46        .map(|r| (r.entry_key, r.discriminator))
47        .collect();
48    let b_keys: HashSet<([u8; 32], u16)> = b
49        .index
50        .iter()
51        .map(|r| (r.entry_key, r.discriminator))
52        .collect();
53
54    let to_ref = |(entry_key, discriminator): &([u8; 32], u16)| EntryRef {
55        entry_key: *entry_key,
56        discriminator: *discriminator,
57    };
58
59    let added: Vec<EntryRef> = b_keys.difference(&a_keys).map(to_ref).collect();
60    let removed: Vec<EntryRef> = a_keys.difference(&b_keys).map(to_ref).collect();
61
62    let mut changed = Vec::new();
63    let mut unchanged = 0usize;
64    for pair @ (key, disc) in a_keys.intersection(&b_keys) {
65        let a_data = index::binary_search_with_discriminator(&a.index, key, *disc)
66            .and_then(|idx| a.entry_at(idx));
67        let b_data = index::binary_search_with_discriminator(&b.index, key, *disc)
68            .and_then(|idx| b.entry_at(idx));
69        if a_data != b_data {
70            changed.push(to_ref(pair));
71        } else {
72            unchanged += 1;
73        }
74    }
75
76    DiffResult {
77        a_entry_count: a.index.len(),
78        b_entry_count: b.index.len(),
79        a_epochs: a.manifest.ca_scopes.iter().map(|s| s.epoch).collect(),
80        b_epochs: b.manifest.ca_scopes.iter().map(|s| s.epoch).collect(),
81        added,
82        removed,
83        changed,
84        unchanged,
85    }
86}
87
88/// Per-delta application statistics.
89#[derive(Debug, Clone)]
90pub struct DeltaStat {
91    pub added: usize,
92    pub replaced: usize,
93    pub removed: usize,
94    /// The delta's `chain_length` exceeded the recommended maximum (24).
95    pub chain_length_warning: bool,
96}
97
98/// Result of materializing a base bundle plus an ordered chain of deltas.
99#[derive(Debug, Clone)]
100pub struct ApplyResult {
101    /// Serialized full bundle bytes.
102    pub bytes: Vec<u8>,
103    pub entry_count: usize,
104    /// Epoch assigned to the materialized bundle's scopes (max seen + 1).
105    pub final_epoch: u64,
106    pub deltas: Vec<DeltaStat>,
107}
108
109/// Recommended maximum delta-chain length before a full re-base is advised.
110pub const MAX_CHAIN_LENGTH: u64 = 24;
111
112/// Apply an ordered chain of delta bundles onto a full base bundle, producing a
113/// materialized full bundle. Verifies the continuity chain (base digest, then
114/// `prev_manifest_digest` links) and rejects type mismatches. Callers are
115/// responsible for authenticating inputs against their configured trust policy.
116/// Output is an UNSIGNED INTERMEDIATE, with an empty seal; it must be sealed
117/// before installation on an authenticated responder.
118pub fn apply(base: &Bundle, deltas: &[Bundle]) -> Result<ApplyResult> {
119    apply_sealed(base, deltas, |_| Ok(Vec::new()))
120}
121
122/// Materialize a delta chain and authenticate its final manifest using `seal_fn`.
123/// The caller must authenticate every input and authorize the output signer.
124pub fn apply_sealed<F>(base: &Bundle, deltas: &[Bundle], seal_fn: F) -> Result<ApplyResult>
125where
126    F: FnOnce(&[u8]) -> Result<Vec<u8>>,
127{
128    crate::verify_structure(base)?;
129    if base.manifest.bundle_type != BundleType::Full {
130        return Err(AhuError::InvalidOperation(
131            "base bundle must be a full bundle, not a delta".into(),
132        ));
133    }
134
135    // Working set keyed by (entry_key, discriminator) so dual-algorithm entries
136    // are tracked independently.
137    type RetainedEntry = (Vec<u8>, IndexFlags, crate::Window);
138    let mut working_set: BTreeMap<([u8; 32], u16), RetainedEntry> = BTreeMap::new();
139    for record in &base.index {
140        if let Some(data) = base.entry_bytes(record) {
141            working_set.insert(
142                (record.entry_key, record.discriminator),
143                (data.to_vec(), record.flags, base.manifest.window.clone()),
144            );
145        }
146    }
147
148    let base_manifest_digest = crate::manifest_digest(&base.manifest_bytes);
149    let mut prev_manifest_digest = base_manifest_digest;
150    let mut max_epoch = base
151        .manifest
152        .ca_scopes
153        .iter()
154        .map(|s| s.epoch)
155        .max()
156        .unwrap_or(0);
157
158    let mut stats = Vec::with_capacity(deltas.len());
159
160    for (i, delta) in deltas.iter().enumerate() {
161        if delta.manifest.bundle_type != BundleType::Delta {
162            return Err(AhuError::InvalidOperation(format!(
163                "delta {} is a full bundle, expected a delta",
164                i + 1
165            )));
166        }
167
168        crate::verify_structure(delta)?;
169        if delta.manifest.continuity.base_manifest_digest != Some(base_manifest_digest) {
170            return Err(AhuError::InvalidOperation(format!(
171                "delta {} base_manifest_digest does not match base bundle",
172                i + 1
173            )));
174        }
175        let scopes = |bundle: &Bundle| {
176            bundle
177                .manifest
178                .ca_scopes
179                .iter()
180                .map(|scope| {
181                    (
182                        scope.hash_algorithm.clone(),
183                        scope.issuer_name_hash.clone(),
184                        scope.issuer_key_hash.clone(),
185                    )
186                })
187                .collect::<std::collections::BTreeSet<_>>()
188        };
189        if delta.manifest.producer_id != base.manifest.producer_id || scopes(delta) != scopes(base)
190        {
191            return Err(AhuError::InvalidOperation(format!(
192                "delta {} producer/CA scopes do not match base",
193                i + 1
194            )));
195        }
196        // Every delta must chain from its predecessor.
197        if let Some(ref prev_digest) = delta.manifest.continuity.prev_manifest_digest {
198            if *prev_digest != prev_manifest_digest {
199                return Err(AhuError::InvalidOperation(format!(
200                    "delta {} prev_manifest_digest chain broken (expected {}, got {})",
201                    i + 1,
202                    hex::encode(prev_manifest_digest),
203                    hex::encode(prev_digest),
204                )));
205            }
206        }
207
208        let chain_length_warning = delta.manifest.continuity.chain_length > MAX_CHAIN_LENGTH;
209
210        let mut added = 0usize;
211        let mut replaced = 0usize;
212        let mut removed = 0usize;
213
214        for record in &delta.index {
215            let key = (record.entry_key, record.discriminator);
216            if record.flags.contains(IndexFlags::TOMBSTONE) {
217                if working_set.remove(&key).is_some() {
218                    removed += 1;
219                }
220            } else if let Some(data) = delta.entry_bytes(record) {
221                if working_set
222                    .insert(
223                        key,
224                        (data.to_vec(), record.flags, delta.manifest.window.clone()),
225                    )
226                    .is_some()
227                {
228                    replaced += 1;
229                } else {
230                    added += 1;
231                }
232            }
233        }
234
235        for scope in &delta.manifest.ca_scopes {
236            max_epoch = max_epoch.max(scope.epoch);
237        }
238        prev_manifest_digest = crate::manifest_digest(&delta.manifest_bytes);
239
240        stats.push(DeltaStat {
241            added,
242            replaced,
243            removed,
244            chain_length_warning,
245        });
246    }
247
248    // Build the materialized full bundle from the base manifest.
249    let mut manifest = base.manifest.clone();
250    manifest.bundle_type = BundleType::Full;
251    manifest.continuity.chain_length = 0;
252    manifest.continuity.prev_manifest_digest = Some(prev_manifest_digest);
253    manifest.continuity.base_manifest_digest = None;
254
255    let final_epoch = max_epoch
256        .checked_add(1)
257        .ok_or_else(|| AhuError::InvalidOperation("epoch exhausted".into()))?;
258    for scope in &mut manifest.ca_scopes {
259        scope.epoch = final_epoch;
260    }
261
262    // Retained payloads carry their source validity envelope. Taking the
263    // earliest expiry is conservative even for batched/multi-algorithm data;
264    // replaced and tombstoned payloads no longer constrain the result.
265    if !working_set.is_empty() {
266        let windows: Vec<_> = working_set.values().map(|(_, _, w)| w).collect();
267        if windows
268            .iter()
269            .any(|w| w.this_update_min > w.next_update_min || w.next_update_min > w.next_update_max)
270        {
271            return Err(AhuError::InvalidOperation(
272                "inconsistent source validity window".into(),
273            ));
274        }
275        manifest.window.this_update_min = windows.iter().map(|w| w.this_update_min).min().unwrap();
276        manifest.window.produced_at = windows.iter().map(|w| w.produced_at).max().unwrap();
277        let expires = windows.iter().map(|w| w.next_update_min).min().unwrap();
278        if manifest.window.produced_at >= expires {
279            return Err(AhuError::InvalidOperation(
280                "retained response validity windows do not overlap".into(),
281            ));
282        }
283        manifest.window.next_update_min = expires;
284        manifest.window.next_update_max = expires;
285    }
286    let mut builder = BundleBuilder::new(manifest);
287    for ((entry_key, disc), (data, _flags, _window)) in &working_set {
288        builder.add_entry_with_discriminator(*entry_key, *disc, data.clone());
289    }
290
291    let bytes = builder.build(seal_fn)?;
292    let entry_count = working_set.len();
293
294    Ok(ApplyResult {
295        bytes,
296        entry_count,
297        final_epoch,
298        deltas: stats,
299    })
300}
301
302#[cfg(test)]
303mod tests {
304    use super::*;
305    use crate::manifest::{
306        CaScope, Completeness, Continuity, Integrity, Manifest, ResponderId, ResponderIdType,
307        Window,
308    };
309    use uuid::Uuid;
310
311    /// Build a bare manifest of the given type/epoch. Continuity is filled in by
312    /// the caller for delta cases.
313    fn manifest(bundle_type: BundleType, epoch: u64) -> Manifest {
314        Manifest {
315            format_version: 1,
316            bundle_id: Uuid::nil(),
317            producer_id: "test".into(),
318            created_at: 1700000000,
319            bundle_type,
320            ca_scopes: vec![CaScope {
321                hash_algorithm: vec![0x01],
322                issuer_name_hash: vec![0xAA; 32],
323                issuer_key_hash: vec![0xBB; 32],
324                epoch,
325                responder_id: ResponderId {
326                    id_type: ResponderIdType::ByKey,
327                    value: vec![0xCC; 20],
328                },
329                responder_chain: None,
330                signature_algorithm: vec![0x02],
331                completeness: Completeness::AuthoritativeComplete,
332            }],
333            window: Window {
334                produced_at: 1700000000,
335                this_update_min: 1700000000,
336                next_update_min: 1700086400,
337                next_update_max: 1700093600,
338            },
339            integrity: Integrity {
340                index_digest: [0; 32],
341                data_digest: [0; 32],
342            },
343            entry_count: 0,
344            continuity: Continuity {
345                prev_manifest_digest: None,
346                base_manifest_digest: None,
347                chain_length: 0,
348            },
349            shard: None,
350            compression: None,
351            extensions: None,
352        }
353    }
354
355    fn seal(m: &[u8]) -> Result<Vec<u8>> {
356        Ok(Sha256::digest(m).to_vec())
357    }
358
359    fn key(byte: u8) -> [u8; 32] {
360        [byte; 32]
361    }
362
363    /// Build a full base bundle with the given (key_byte, response) entries.
364    fn full_bundle(epoch: u64, entries: &[(u8, &[u8])]) -> Bundle {
365        let mut builder = BundleBuilder::new(manifest(BundleType::Full, epoch));
366        for (k, resp) in entries {
367            builder.add_entry(key(*k), resp.to_vec());
368        }
369        let bytes = builder.build(seal).unwrap();
370        Bundle::from_bytes(&bytes).unwrap()
371    }
372
373    #[test]
374    fn diff_detects_added_removed_changed() {
375        // A: {1 => "old", 2 => "same"}   B: {2 => "same", 3 => "new"} with 1 changed dropped.
376        let a = full_bundle(1, &[(1, b"one"), (2, b"same")]);
377        // B replaces key 2's payload, drops key 1, adds key 3.
378        let b = full_bundle(2, &[(2, b"CHANGED"), (3, b"three")]);
379
380        let d = diff(&a, &b);
381        assert_eq!(d.a_entry_count, 2);
382        assert_eq!(d.b_entry_count, 2);
383        assert_eq!(d.a_epochs, vec![1]);
384        assert_eq!(d.b_epochs, vec![2]);
385
386        assert_eq!(d.added.len(), 1, "key 3 added");
387        assert_eq!(d.added[0].entry_key, key(3));
388        assert_eq!(d.removed.len(), 1, "key 1 removed");
389        assert_eq!(d.removed[0].entry_key, key(1));
390        assert_eq!(d.changed.len(), 1, "key 2 payload changed");
391        assert_eq!(d.changed[0].entry_key, key(2));
392        assert_eq!(d.unchanged, 0);
393    }
394
395    #[test]
396    fn diff_identical_bundles_all_unchanged() {
397        let a = full_bundle(1, &[(1, b"one"), (2, b"two")]);
398        let b = full_bundle(1, &[(1, b"one"), (2, b"two")]);
399        let d = diff(&a, &b);
400        assert!(d.added.is_empty());
401        assert!(d.removed.is_empty());
402        assert!(d.changed.is_empty());
403        assert_eq!(d.unchanged, 2);
404    }
405
406    /// Build a delta that chains from `prev` (which may be the base or a prior
407    /// delta), setting continuity digests correctly.
408    fn delta_from(
409        prev: &Bundle,
410        base: &Bundle,
411        epoch: u64,
412        chain_length: u64,
413        adds: &[(u8, &[u8])],
414        tombstones: &[u8],
415    ) -> Bundle {
416        let mut m = manifest(BundleType::Delta, epoch);
417        m.continuity.base_manifest_digest = Some(crate::manifest_digest(&base.manifest_bytes));
418        m.continuity.prev_manifest_digest = Some(crate::manifest_digest(&prev.manifest_bytes));
419        m.continuity.chain_length = chain_length;
420        let mut builder = BundleBuilder::new(m);
421        for (k, resp) in adds {
422            builder.add_entry(key(*k), resp.to_vec());
423        }
424        for k in tombstones {
425            builder.add_tombstone(key(*k), 0);
426        }
427        let bytes = builder.build(seal).unwrap();
428        Bundle::from_bytes(&bytes).unwrap()
429    }
430
431    #[test]
432    fn apply_single_delta_add_replace_remove() {
433        let base = full_bundle(5, &[(1, b"one"), (2, b"two")]);
434        // Delta: replace 1, remove 2, add 3.
435        let delta = delta_from(&base, &base, 6, 1, &[(1, b"ONE"), (3, b"three")], &[2]);
436
437        let result = apply(&base, &[delta]).unwrap();
438        assert_eq!(result.entry_count, 2, "one + three remain (two tombstoned)");
439        assert_eq!(result.final_epoch, 7, "max epoch (6) + 1");
440        assert_eq!(result.deltas.len(), 1);
441        assert_eq!(result.deltas[0].added, 1, "key 3");
442        assert_eq!(result.deltas[0].replaced, 1, "key 1");
443        assert_eq!(result.deltas[0].removed, 1, "key 2");
444        assert!(!result.deltas[0].chain_length_warning);
445
446        // The materialized bundle must parse and reflect the applied changes.
447        let materialized = Bundle::from_bytes(&result.bytes).unwrap();
448        assert_eq!(materialized.manifest.bundle_type, BundleType::Full);
449        assert_eq!(materialized.index.len(), 2);
450        let idx1 =
451            index::binary_search_with_discriminator(&materialized.index, &key(1), 0).unwrap();
452        assert_eq!(materialized.entry_at(idx1), Some(&b"ONE"[..]));
453        assert!(index::binary_search_with_discriminator(&materialized.index, &key(2), 0).is_none());
454    }
455
456    #[test]
457    fn apply_two_delta_chain() {
458        let base = full_bundle(5, &[(1, b"one")]);
459        let d1 = delta_from(&base, &base, 6, 1, &[(2, b"two")], &[]);
460        let d2 = delta_from(&d1, &base, 7, 2, &[(3, b"three")], &[]);
461
462        let result = apply(&base, &[d1, d2]).unwrap();
463        assert_eq!(result.entry_count, 3);
464        assert_eq!(result.final_epoch, 8);
465        assert_eq!(result.deltas.len(), 2);
466    }
467
468    #[test]
469    fn apply_rejects_full_bundle_as_delta() {
470        let base = full_bundle(5, &[(1, b"one")]);
471        let not_a_delta = full_bundle(6, &[(2, b"two")]);
472        let err = apply(&base, &[not_a_delta]).unwrap_err();
473        assert!(matches!(err, AhuError::InvalidOperation(_)));
474    }
475
476    #[test]
477    fn apply_rejects_delta_as_base() {
478        let base = full_bundle(5, &[(1, b"one")]);
479        let delta = delta_from(&base, &base, 6, 1, &[(2, b"two")], &[]);
480        // Passing the delta as the base must be rejected.
481        let err = apply(&delta, &[]).unwrap_err();
482        assert!(matches!(err, AhuError::InvalidOperation(_)));
483    }
484
485    #[test]
486    fn apply_rejects_broken_chain() {
487        let base = full_bundle(5, &[(1, b"one")]);
488        let d1 = delta_from(&base, &base, 6, 1, &[(2, b"two")], &[]);
489        // d2 chains from d1, but we apply [d2] directly onto base — prev digest mismatch.
490        let d2 = delta_from(&d1, &base, 7, 2, &[(3, b"three")], &[]);
491        let err = apply(&base, &[d2]).unwrap_err();
492        assert!(matches!(err, AhuError::InvalidOperation(_)));
493    }
494
495    #[test]
496    fn apply_flags_chain_length_warning() {
497        let base = full_bundle(5, &[(1, b"one")]);
498        let delta = delta_from(&base, &base, 6, MAX_CHAIN_LENGTH + 1, &[(2, b"two")], &[]);
499        let result = apply(&base, &[delta]).unwrap();
500        assert!(result.deltas[0].chain_length_warning);
501    }
502    #[test]
503    fn every_delta_must_reference_base_even_without_predecessor() {
504        let base = full_bundle(1, &[(1, b"one")]);
505        let first = delta_from(&base, &base, 2, 1, &[], &[]);
506        let mut second = delta_from(&first, &base, 3, 2, &[], &[]);
507        second.manifest.continuity.base_manifest_digest = Some([0xff; 32]);
508        second.manifest.continuity.prev_manifest_digest = None;
509        assert!(apply(&base, &[first, second]).is_err());
510    }
511
512    #[test]
513    fn delta_validity_tracks_retained_payloads_and_output_is_unsigned() {
514        let base = full_bundle(1, &[(1, b"one"), (2, b"two")]);
515        let mut delta = delta_from(&base, &base, 2, 1, &[(1, b"ONE")], &[]);
516        let expiry = base.manifest.window.produced_at + 10;
517        delta.manifest.window.next_update_min = expiry;
518        delta.manifest.window.next_update_max = expiry;
519        let output = apply(&base, &[delta]).unwrap();
520        let materialized = Bundle::from_bytes(&output.bytes).unwrap();
521        assert_eq!(materialized.manifest.window.next_update_max, expiry);
522        assert!(materialized.seal_bytes.is_empty());
523    }
524
525    #[test]
526    fn delta_cannot_change_producer_or_ca_scope() {
527        let base = full_bundle(1, &[(1, b"one")]);
528        let mut delta = delta_from(&base, &base, 2, 1, &[], &[]);
529        delta.manifest.producer_id = "other".into();
530        assert!(apply(&base, &[delta]).is_err());
531        let mut delta = delta_from(&base, &base, 2, 1, &[], &[]);
532        delta.manifest.ca_scopes[0].issuer_key_hash = vec![0xff; 32];
533        assert!(apply(&base, &[delta]).is_err());
534    }
535}