1use std::collections::BTreeMap;
43use std::path::Path;
44
45use crate::container::{Descriptor, ObjectSource};
46use crate::error::Result;
47use crate::store::{EmbeddedStore, Id, ObjectStore, externalize};
48
49pub const KIND_OBJECT: &str = "object";
51pub const KIND_CHANNEL_PAYLOAD: &str = "channel_payload";
53pub const KIND_CHANNEL_HEADER: &str = "channel_header";
55pub const KIND_MODEL: &str = "model";
57
58#[derive(Debug, Clone, Copy, PartialEq, Eq)]
63pub struct ShareUnit {
64 pub kind: &'static str,
66 pub id: [u8; 32],
68 pub len: u64,
70}
71
72#[derive(Debug, Clone, Default, PartialEq, Eq)]
80pub struct ShareReport {
81 pub total_bytes: u64,
83 pub unique_bytes: u64,
85 pub unit_count: u64,
87 pub unique_count: u64,
89 pub by_kind: Vec<(String, u64, u64)>,
91}
92
93impl ShareReport {
94 pub fn to_json(&self) -> String {
96 let kinds: Vec<String> = self
97 .by_kind
98 .iter()
99 .map(|(k, t, u)| format!("{{\"kind\":\"{k}\",\"total\":{t},\"unique\":{u}}}"))
100 .collect();
101 format!(
102 concat!(
103 "{{",
104 "\"total_bytes\":{},",
105 "\"unique_bytes\":{},",
106 "\"unit_count\":{},",
107 "\"unique_count\":{},",
108 "\"by_kind\":[{}]",
109 "}}"
110 ),
111 self.total_bytes,
112 self.unique_bytes,
113 self.unit_count,
114 self.unique_count,
115 kinds.join(",")
116 )
117 }
118}
119
120fn unit(kind: &'static str, bytes: &[u8]) -> ShareUnit {
122 ShareUnit {
123 kind,
124 id: *Id::of(bytes).as_bytes(),
125 len: bytes.len() as u64,
126 }
127}
128
129pub fn extract_units(descriptor: &Descriptor) -> Result<Vec<ShareUnit>> {
136 let mut units = Vec::new();
137 for model in &descriptor.models {
138 let bytes = model.encode()?;
139 units.push(unit(KIND_MODEL, &bytes));
140 }
141 for channel in &descriptor.channels {
142 units.push(unit(KIND_CHANNEL_HEADER, &channel.header_bytes()?));
143 units.push(unit(KIND_CHANNEL_PAYLOAD, &channel.payload));
144 }
145 for obj in &descriptor.objects {
146 match obj {
147 ObjectSource::Inline(bytes) => units.push(unit(KIND_OBJECT, bytes)),
148 ObjectSource::External { id, len } => units.push(ShareUnit {
149 kind: KIND_OBJECT,
150 id: *id.as_bytes(),
151 len: *len,
152 }),
153 }
154 }
155 Ok(units)
156}
157
158pub fn cohort_report(descriptors: &[Descriptor]) -> Result<ShareReport> {
164 let mut report = ShareReport::default();
165 let mut seen: BTreeMap<[u8; 32], u64> = BTreeMap::new();
167 let mut kind_total: BTreeMap<&'static str, u64> = BTreeMap::new();
169 let mut kind_unique: BTreeMap<&'static str, BTreeMap<[u8; 32], u64>> = BTreeMap::new();
171
172 for descriptor in descriptors {
173 for u in extract_units(descriptor)? {
174 report.total_bytes += u.len;
175 report.unit_count += 1;
176 *kind_total.entry(u.kind).or_default() += u.len;
177 kind_unique.entry(u.kind).or_default().insert(u.id, u.len);
178 if seen.insert(u.id, u.len).is_none() {
179 report.unique_bytes += u.len;
180 report.unique_count += 1;
181 }
182 }
183 }
184
185 report.by_kind = kind_total
186 .into_iter()
187 .map(|(kind, total)| {
188 let unique: u64 = kind_unique.get(kind).map(|m| m.values().sum()).unwrap_or(0);
189 (kind.to_string(), total, unique)
190 })
191 .collect();
192 Ok(report)
193}
194
195pub fn share_store(field_root: &Path) -> Result<EmbeddedStore> {
198 EmbeddedStore::open(field_root.join("share"))
199}
200
201pub fn externalize_objects(field_root: &Path, descriptor: &mut Descriptor) -> Result<u64> {
208 let mut store = share_store(field_root)?;
209 let resolver = store.clone();
210 externalize(descriptor, &resolver, &mut store)?;
211 Ok(descriptor.objects.len() as u64)
212}
213
214pub fn store_units(field_root: &Path, descriptor: &Descriptor) -> Result<u64> {
224 let mut store = share_store(field_root)?;
225 let mut offered: u64 = 0;
226 for model in &descriptor.models {
227 store.put(&model.encode()?)?;
228 offered += 1;
229 }
230 for channel in &descriptor.channels {
231 store.put(&channel.header_bytes()?)?;
232 store.put(&channel.payload)?;
233 offered += 2;
234 }
235 for obj in &descriptor.objects {
236 if let ObjectSource::Inline(bytes) = obj {
237 store.put(bytes)?;
238 offered += 1;
239 }
240 }
241 Ok(offered)
242}
243
244#[cfg(test)]
245mod tests {
246 use super::*;
247 use crate::container::UNIVERSE;
248 use crate::dra::Program;
249 use crate::entropy::codec::{
250 CODER_ORDER0_BYTE_RANS, CODER_VERSION_1, EntropyChannelDescriptor,
251 };
252 use crate::entropy::model::EntropyModel;
253 use crate::{SOURCE_FORMAT_OPAQUE, integrity};
254
255 fn channel(
256 payload: Vec<u8>,
257 initial_state: u32,
258 decoded_length: u64,
259 ) -> EntropyChannelDescriptor {
260 EntropyChannelDescriptor {
261 coder: CODER_ORDER0_BYTE_RANS,
262 coder_version: CODER_VERSION_1,
263 scale_bits: 8,
264 lane_count: 1,
265 model_id: 0,
266 symbol_count: decoded_length,
267 decoded_length,
268 initial_state,
269 payload,
270 }
271 }
272
273 fn descriptor(
274 objects: Vec<ObjectSource>,
275 channels: Vec<EntropyChannelDescriptor>,
276 ) -> Descriptor {
277 Descriptor {
278 universe: UNIVERSE.to_string(),
279 source_format: SOURCE_FORMAT_OPAQUE,
280 format_basis: "opaque:share-test".to_string(),
281 models: vec![EntropyModel::uniform(8).unwrap()],
282 channels,
283 objects,
284 program: Program::new(vec![]),
285 observation_index: None,
286 seek_directory: false,
287 source_sha256: [0u8; 32],
288 source_len: 0,
289 }
290 }
291
292 #[test]
293 fn repeated_objects_are_counted_twice_and_unique_once() {
294 let payload = b"the same object bytes".to_vec();
295 let distinct = b"distinct".to_vec();
296 let d = descriptor(
297 vec![
298 ObjectSource::Inline(payload.clone()),
299 ObjectSource::Inline(payload.clone()),
300 ObjectSource::Inline(distinct.clone()),
301 ],
302 vec![],
303 );
304 let report = cohort_report(std::slice::from_ref(&d)).unwrap();
305 assert_eq!(report.unit_count, 4);
307 let model_len = EntropyModel::uniform(8).unwrap().encode().unwrap().len() as u64;
309 let distinct_len = distinct.len() as u64;
310 assert_eq!(
311 report.total_bytes,
312 model_len + 2 * payload.len() as u64 + distinct_len
313 );
314 assert_eq!(report.unique_count, 3);
315 assert_eq!(
316 report.unique_bytes,
317 model_len + payload.len() as u64 + distinct_len
318 );
319 assert!(report.unique_bytes <= report.total_bytes);
320 }
321
322 #[test]
323 fn identical_channel_payloads_share_the_payload_not_the_header() {
324 let payload = vec![7u8; 64];
325 let a = channel(payload.clone(), 1, 64);
326 let b = channel(payload, 2, 64);
327 let d = descriptor(vec![], vec![a, b]);
328 let report = cohort_report(std::slice::from_ref(&d)).unwrap();
329 let by = |kind: &str| {
330 report
331 .by_kind
332 .iter()
333 .find(|(k, _, _)| k == kind)
334 .cloned()
335 .unwrap()
336 };
337 let (_, payload_total, payload_unique) = by(KIND_CHANNEL_PAYLOAD);
338 assert_eq!(payload_total, 128);
339 assert_eq!(payload_unique, 64, "identical payloads share once");
340 let (_, header_total, header_unique) = by(KIND_CHANNEL_HEADER);
341 assert_eq!(header_total, 2 * 33);
342 assert_eq!(
343 header_unique,
344 2 * 33,
345 "distinct initial_state -> distinct header"
346 );
347 }
348
349 #[test]
350 fn identical_headers_share_when_only_payload_differs() {
351 let a = channel(vec![1u8; 16], 9, 16);
352 let b = channel(vec![2u8; 16], 9, 16);
353 let d = descriptor(vec![], vec![a, b]);
354 let report = cohort_report(std::slice::from_ref(&d)).unwrap();
355 let header = report
356 .by_kind
357 .iter()
358 .find(|(k, _, _)| k == KIND_CHANNEL_HEADER)
359 .unwrap();
360 assert_eq!(header.1, 66);
361 assert_eq!(header.2, 33, "equal-length payloads -> one shared header");
362 }
363
364 #[test]
365 fn report_is_deterministic_and_order_independent() {
366 let d1 = descriptor(
367 vec![ObjectSource::Inline(b"alpha".to_vec())],
368 vec![channel(vec![1, 2, 3], 1, 3)],
369 );
370 let d2 = descriptor(
371 vec![ObjectSource::Inline(b"alpha".to_vec())],
372 vec![channel(vec![4, 5, 6], 2, 3)],
373 );
374 let a = cohort_report(&[d1.clone(), d2.clone()]).unwrap();
375 let b = cohort_report(&[d2, d1]).unwrap();
376 assert_eq!(a, b);
377 assert!(a.unique_bytes <= a.total_bytes);
378 }
379
380 #[test]
381 fn external_objects_contribute_id_and_length_without_a_store() {
382 let inline = descriptor(vec![ObjectSource::Inline(b"0123456789".to_vec())], vec![]);
383 let id = Id::of(b"0123456789");
384 let external = descriptor(vec![ObjectSource::External { id, len: 10 }], vec![]);
385 assert_eq!(
386 extract_units(&inline).unwrap(),
387 extract_units(&external).unwrap()
388 );
389 }
390
391 #[test]
392 fn object_ids_and_source_sha_are_distinct_namespaces() {
393 let bytes = b"namespace check";
395 let u = unit(KIND_OBJECT, bytes);
396 assert_eq!(u.id, *Id::of(bytes).as_bytes());
397 assert_ne!(u.id, integrity::sha256(bytes));
398 }
399
400 #[test]
401 fn store_units_persists_every_unique_unit_exactly_once() {
402 let root = std::env::temp_dir().join(format!(
407 "vole-share-store-{}-{}",
408 std::process::id(),
409 std::time::SystemTime::now()
410 .duration_since(std::time::UNIX_EPOCH)
411 .unwrap()
412 .as_nanos()
413 ));
414 let shared_payload = vec![9u8; 32];
415 let a = descriptor(
416 vec![ObjectSource::Inline(b"shared-object".to_vec())],
417 vec![channel(shared_payload.clone(), 1, 32)],
418 );
419 let b = descriptor(
420 vec![ObjectSource::Inline(b"shared-object".to_vec())],
421 vec![channel(shared_payload, 1, 32)],
422 );
423 let offered = store_units(&root, &a).unwrap() + store_units(&root, &b).unwrap();
424 assert_eq!(
425 offered, 8,
426 "2 descriptors x (model + header + payload + object)"
427 );
428 let report = cohort_report(&[a, b]).unwrap();
429 let stats = share_store(&root).unwrap().stats().unwrap();
430 assert_eq!(stats.stored_bytes, report.unique_bytes);
431 assert_eq!(stats.object_count, report.unique_count);
432 std::fs::remove_dir_all(&root).ok();
433 }
434
435 #[test]
436 fn externalize_objects_resolves_and_stays_exact() {
437 let root = std::env::temp_dir().join(format!(
438 "vole-share-{}-{}",
439 std::process::id(),
440 std::time::SystemTime::now()
441 .duration_since(std::time::UNIX_EPOCH)
442 .unwrap()
443 .as_nanos()
444 ));
445 let source = b"externalized exact bytes";
446 let mut d = Descriptor {
447 universe: UNIVERSE.to_string(),
448 source_format: SOURCE_FORMAT_OPAQUE,
449 format_basis: "opaque:share-test".to_string(),
450 models: vec![],
451 channels: vec![],
452 objects: vec![
453 ObjectSource::Inline(source.to_vec()),
454 ObjectSource::Inline(source.to_vec()),
455 ],
456 program: Program::new(vec![crate::dra::Op::EmitObject { object_id: 0 }]),
457 observation_index: None,
458 seek_directory: false,
459 source_sha256: integrity::sha256(source),
460 source_len: source.len() as u64,
461 };
462 let offered = store_units(&root, &d).unwrap();
463 assert_eq!(offered, 2);
464 externalize_objects(&root, &mut d).unwrap();
465 assert!(
466 d.objects
467 .iter()
468 .all(|o| matches!(o, ObjectSource::External { .. }))
469 );
470 let (bytes, _) = d.serialize().unwrap();
471 let parsed = Descriptor::parse(&bytes, crate::limits::Limits::DEFAULT).unwrap();
472 let store = share_store(&root).unwrap();
473 let out =
474 crate::materialize::materialize_with(&parsed, &store, crate::limits::Limits::DEFAULT)
475 .unwrap();
476 assert_eq!(out, source);
477 std::fs::remove_dir_all(&root).ok();
478 }
479}