1use 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#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
20pub struct EntryRef {
21 pub entry_key: [u8; 32],
22 pub discriminator: u16,
23}
24
25#[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
40pub 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#[derive(Debug, Clone)]
90pub struct DeltaStat {
91 pub added: usize,
92 pub replaced: usize,
93 pub removed: usize,
94 pub chain_length_warning: bool,
96}
97
98#[derive(Debug, Clone)]
100pub struct ApplyResult {
101 pub bytes: Vec<u8>,
103 pub entry_count: usize,
104 pub final_epoch: u64,
106 pub deltas: Vec<DeltaStat>,
107}
108
109pub const MAX_CHAIN_LENGTH: u64 = 24;
111
112pub fn apply(base: &Bundle, deltas: &[Bundle]) -> Result<ApplyResult> {
119 apply_sealed(base, deltas, |_| Ok(Vec::new()))
120}
121
122pub 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 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 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 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 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 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 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 let a = full_bundle(1, &[(1, b"one"), (2, b"same")]);
377 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 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 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 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 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 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}