1use std::collections::HashMap;
32
33use crate::format::chunk_index::btree_v2::build_index;
34use crate::format::fractal_heap::HeapParams;
35use crate::format::fractal_heap_write::{plan_heap, HeapBlock};
36use crate::format::sohm::{
37 encode_list, list_size, message_hash, record_size, SohmIndexHeader, SohmMasterTable,
38 SohmRecord, SohmRecordLocation, BT2_TYPE_SOHM_INDEX, SOHM_B2_NODE_SIZE, SOHM_HEAP_ID_LEN,
39 SOHM_INDEX_BTREE, SOHM_INDEX_LIST,
40};
41use crate::format::{FormatContext, FormatError, FormatResult};
42
43#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
47pub struct SohmIndexSpec {
48 pub mesg_types: u16,
51 pub min_mesg_size: u32,
53 pub list_max: u16,
55 pub btree_min: u16,
57}
58
59pub type SharedKey = (u8, Vec<u8>);
62
63#[derive(Debug, Clone, PartialEq, Eq, Hash)]
77pub struct NestedShare {
78 pub heap_id_at: usize,
81 pub target: SharedKey,
84}
85
86#[derive(Debug, Clone, PartialEq, Eq)]
89pub struct SharedMessage {
90 pub msg_type: u8,
92 pub body: Vec<u8>,
95 pub nested: Vec<NestedShare>,
98 pub ref_count: u32,
100 pub ohdr_addr: Option<u64>,
110}
111
112impl SharedMessage {
113 pub fn kept_in_ohdr(&self) -> Option<u64> {
119 self.ohdr_addr.filter(|_| self.ref_count <= 1)
120 }
121}
122
123#[derive(Debug, Clone, PartialEq, Eq)]
125pub struct SohmIndexContent {
126 pub spec: SohmIndexSpec,
128 pub messages: Vec<SharedMessage>,
130}
131
132#[derive(Debug, Clone, PartialEq, Eq)]
135pub struct BuiltSharedMessages {
136 pub table_addr: u64,
139 pub blocks: Vec<HeapBlock>,
141 pub heap_ids: HashMap<(u8, Vec<u8>), [u8; SOHM_HEAP_ID_LEN]>,
144}
145
146pub fn build_shared_messages(
152 indexes: &[SohmIndexContent],
153 ctx: &FormatContext,
154 alloc: &mut dyn FnMut(u64) -> u64,
155) -> FormatResult<BuiltSharedMessages> {
156 let mut blocks = Vec::new();
157 let mut heap_ids: HashMap<SharedKey, [u8; SOHM_HEAP_ID_LEN]> = HashMap::new();
158 let mut headers = Vec::with_capacity(indexes.len());
159
160 let mut plans = Vec::with_capacity(indexes.len());
166 for index in indexes {
167 let lengths: Vec<usize> = index
168 .messages
169 .iter()
170 .filter(|m| m.kept_in_ohdr().is_none())
171 .map(|m| m.body.len())
172 .collect();
173 plans.push(plan_heap(
174 &HeapParams::object_header(),
175 ctx,
176 &lengths,
177 alloc,
178 )?);
179 }
180
181 let mut collected: HashMap<SharedKey, [u8; SOHM_HEAP_ID_LEN]> = HashMap::new();
186 for (index, plan) in indexes.iter().zip(&plans) {
187 let heaped = index.messages.iter().filter(|m| m.kept_in_ohdr().is_none());
188 for (message, id) in heaped.zip(plan.ids()) {
189 collected.insert((message.msg_type, message.body.clone()), heap_id(id)?);
190 }
191 }
192
193 for (index, plan) in indexes.iter().zip(plans) {
194 let num_messages = u16::try_from(index.messages.len()).map_err(|_| {
195 FormatError::InvalidData(format!(
196 "shared-message index holds {} messages, more than the count field takes",
197 index.messages.len()
198 ))
199 })?;
200
201 let mut bodies = Vec::with_capacity(index.messages.len());
202 for message in &index.messages {
203 bodies.push(resolve_nested(message, &collected)?);
204 }
205 let heap_addr = plan.header_addr();
206 let ids = plan.ids().to_vec();
207 let heaped: Vec<Vec<u8>> = index
210 .messages
211 .iter()
212 .zip(&bodies)
213 .filter(|(m, _)| m.kept_in_ohdr().is_none())
214 .map(|(_, body)| body.clone())
215 .collect();
216 blocks.extend(plan.finish(&heaped)?.blocks);
217
218 let mut records = Vec::with_capacity(index.messages.len());
219 let mut heaped_at = 0usize;
220 for (message, body) in index.messages.iter().zip(bodies) {
221 let location = match message.kept_in_ohdr() {
222 Some(oh_addr) => SohmRecordLocation::InObjectHeader {
223 msg_type: message.msg_type,
224 index: 0,
225 oh_addr,
226 },
227 None => {
228 let id = heap_id(&ids[heaped_at])?;
229 heaped_at += 1;
230 heap_ids.insert((message.msg_type, body.clone()), id);
234 SohmRecordLocation::InHeap {
235 ref_count: message.ref_count,
236 heap_id: id,
237 }
238 }
239 };
240 records.push((
241 SohmRecord {
242 hash: message_hash(&body, message.msg_type),
243 location,
244 },
245 body,
246 ));
247 }
248
249 let (index_type, index_addr) = if is_btree(&index.spec, num_messages) {
250 (
251 SOHM_INDEX_BTREE,
252 build_btree(&mut records, ctx, alloc, &mut blocks),
253 )
254 } else {
255 (
256 SOHM_INDEX_LIST,
257 build_list(&records, &index.spec, ctx, alloc, &mut blocks),
258 )
259 };
260
261 headers.push(SohmIndexHeader {
262 index_type,
263 mesg_types: index.spec.mesg_types,
264 min_mesg_size: index.spec.min_mesg_size,
265 list_max: index.spec.list_max,
266 btree_min: index.spec.btree_min,
267 num_messages,
268 index_addr,
269 heap_addr,
270 });
271 }
272
273 let table = SohmMasterTable { indexes: headers };
274 let nindexes = u8::try_from(indexes.len()).map_err(|_| {
275 FormatError::InvalidData(format!("{} shared-message indexes", indexes.len()))
276 })?;
277 let table_addr = alloc(SohmMasterTable::encoded_size(ctx, nindexes) as u64);
278 let image = table.encode(ctx);
279 blocks.push(HeapBlock {
280 addr: table_addr,
281 len: image.len() as u64,
282 image,
283 });
284
285 Ok(BuiltSharedMessages {
286 table_addr,
287 blocks,
288 heap_ids,
289 })
290}
291
292fn is_btree(spec: &SohmIndexSpec, num_messages: u16) -> bool {
300 spec.list_max == 0 || num_messages > spec.list_max
301}
302
303fn heap_id(id: &[u8]) -> FormatResult<[u8; SOHM_HEAP_ID_LEN]> {
306 id.try_into().map_err(|_| {
307 FormatError::InvalidData(format!(
308 "shared-message heap returned a {}-byte id, expected {SOHM_HEAP_ID_LEN}",
309 id.len()
310 ))
311 })
312}
313
314fn resolve_nested(
320 message: &SharedMessage,
321 placed: &HashMap<SharedKey, [u8; SOHM_HEAP_ID_LEN]>,
322) -> FormatResult<Vec<u8>> {
323 if message.nested.is_empty() {
324 return Ok(message.body.clone());
325 }
326 let mut body = message.body.clone();
327 for nested in &message.nested {
328 let Some(heap_id) = placed.get(&nested.target) else {
329 return Err(FormatError::InvalidData(
330 "a shared message points at a body no index holds".into(),
331 ));
332 };
333 let at = nested.heap_id_at;
334 if body.len() < at + SOHM_HEAP_ID_LEN {
335 return Err(FormatError::InvalidData(
336 "a shared message's nested pointer runs past the body holding it".into(),
337 ));
338 }
339 body[at..at + SOHM_HEAP_ID_LEN].copy_from_slice(heap_id);
340 }
341 Ok(body)
342}
343
344fn build_list(
346 records: &[(SohmRecord, Vec<u8>)],
347 spec: &SohmIndexSpec,
348 ctx: &FormatContext,
349 alloc: &mut dyn FnMut(u64) -> u64,
350 blocks: &mut Vec<HeapBlock>,
351) -> u64 {
352 let len = list_size(ctx, spec.list_max) as u64;
355 let addr = alloc(len);
356 let entries: Vec<SohmRecord> = records.iter().map(|(r, _)| *r).collect();
357 blocks.push(HeapBlock {
358 addr,
359 len,
360 image: encode_list(&entries, ctx),
361 });
362 addr
363}
364
365fn build_btree(
367 records: &mut [(SohmRecord, Vec<u8>)],
368 ctx: &FormatContext,
369 alloc: &mut dyn FnMut(u64) -> u64,
370 blocks: &mut Vec<HeapBlock>,
371) -> u64 {
372 records.sort_by(|a, b| a.0.hash.cmp(&b.0.hash).then_with(|| a.1.cmp(&b.1)));
376 let mut image = Vec::with_capacity(records.len() * record_size(ctx));
377 for (record, _) in records.iter() {
378 image.extend_from_slice(&record.encode(ctx));
379 }
380 let (addr, nodes) = build_index(
381 BT2_TYPE_SOHM_INDEX,
382 record_size(ctx) as u16,
383 SOHM_B2_NODE_SIZE,
384 &image,
385 ctx,
386 alloc,
387 );
388 blocks.extend(nodes.into_iter().map(|(addr, image)| HeapBlock {
389 addr,
390 len: image.len() as u64,
391 image,
392 }));
393 addr
394}
395
396#[cfg(test)]
399mod tests {
400 use super::*;
401 use crate::format::chunk_index::btree_v2::{collect_btree_v2_records, Bt2Header};
402 use crate::format::fractal_heap::{
403 collect_managed_blocks, read_heap_object, FractalHeapHeader, HeapId,
404 };
405 use crate::format::messages::{MSG_ATTRIBUTE, MSG_DATASPACE, MSG_DATATYPE};
406 use crate::format::sohm::{SharedLocation, SharedMessagePointer, SOHM_IN_HEAP};
407 use crate::format::{BlockReader, UNDEF_ADDR};
408
409 struct MemFile {
412 bytes: Vec<u8>,
413 }
414
415 impl MemFile {
416 fn new() -> Self {
417 Self { bytes: vec![0; 16] }
419 }
420 fn alloc(&mut self, len: u64) -> u64 {
421 let addr = self.bytes.len() as u64;
422 self.bytes.resize(self.bytes.len() + len as usize, 0);
423 addr
424 }
425 }
426
427 impl BlockReader for MemFile {
428 fn read_block(&mut self, offset: u64, len: usize) -> FormatResult<Vec<u8>> {
429 let start = offset as usize;
430 if start > self.bytes.len() {
431 return Err(FormatError::BufferTooShort {
432 needed: start,
433 available: self.bytes.len(),
434 });
435 }
436 let end = (start + len).min(self.bytes.len());
437 Ok(self.bytes[start..end].to_vec())
438 }
439 }
440
441 fn ctx() -> FormatContext {
442 FormatContext::default_v3()
443 }
444
445 fn spec(list_max: u16) -> SohmIndexSpec {
446 SohmIndexSpec {
447 mesg_types: (1 << MSG_DATASPACE) | (1 << MSG_DATATYPE) | (1 << MSG_ATTRIBUTE),
448 min_mesg_size: 0,
449 list_max,
450 btree_min: 40,
451 }
452 }
453
454 fn message(msg_type: u8, seed: u8, len: usize, ref_count: u32) -> SharedMessage {
455 SharedMessage {
456 msg_type,
457 body: (0..len).map(|i| seed.wrapping_add(i as u8)).collect(),
458 nested: Vec::new(),
459 ref_count,
460 ohdr_addr: None,
461 }
462 }
463
464 fn lay_out(indexes: &[SohmIndexContent]) -> (MemFile, BuiltSharedMessages) {
467 let mut file = MemFile::new();
468 let built = build_shared_messages(indexes, &ctx(), &mut |len| file.alloc(len)).unwrap();
469 for block in &built.blocks {
470 assert!(
471 block.image.len() as u64 <= block.len,
472 "block image overruns its allocation"
473 );
474 let at = block.addr as usize;
475 file.bytes[at..at + block.image.len()].copy_from_slice(&block.image);
476 }
477 (file, built)
478 }
479
480 fn read_records(file: &mut MemFile, header: &SohmIndexHeader) -> Vec<SohmRecord> {
482 let size = record_size(&ctx());
483 let raw = if header.index_type == SOHM_INDEX_LIST {
484 let buf = file
485 .read_block(header.index_addr, 4 + size * header.num_messages as usize)
486 .unwrap();
487 assert_eq!(&buf[..4], b"SMLI");
488 buf[4..].to_vec()
489 } else {
490 let bt2 = Bt2Header::decode(&file.read_block(header.index_addr, 256).unwrap(), &ctx())
491 .unwrap();
492 assert_eq!(bt2.record_type, BT2_TYPE_SOHM_INDEX);
493 assert_eq!(bt2.node_size, SOHM_B2_NODE_SIZE);
494 assert_eq!(bt2.record_size as usize, size);
495 collect_btree_v2_records(&bt2, &ctx(), file).unwrap()
496 };
497 raw.chunks_exact(size)
498 .map(|r| SohmRecord {
499 hash: u32::from_le_bytes(r[1..5].try_into().unwrap()),
500 location: if r[0] == SOHM_IN_HEAP {
501 SohmRecordLocation::InHeap {
502 ref_count: u32::from_le_bytes(r[5..9].try_into().unwrap()),
503 heap_id: r[9..17].try_into().unwrap(),
504 }
505 } else {
506 SohmRecordLocation::InObjectHeader {
507 msg_type: r[6],
508 index: u16::from_le_bytes(r[7..9].try_into().unwrap()),
509 oh_addr: u64::from_le_bytes(r[9..17].try_into().unwrap()),
510 }
511 },
512 })
513 .collect()
514 }
515
516 fn in_heap(record: &SohmRecord) -> (u32, [u8; SOHM_HEAP_ID_LEN]) {
518 match record.location {
519 SohmRecordLocation::InHeap { ref_count, heap_id } => (ref_count, heap_id),
520 other => panic!("expected a heap record, got {other:?}"),
521 }
522 }
523
524 fn read_body(file: &mut MemFile, heap_addr: u64, record: &SohmRecord) -> Vec<u8> {
526 let heap =
527 FractalHeapHeader::decode(&file.read_block(heap_addr, 512).unwrap(), &ctx()).unwrap();
528 let blocks = collect_managed_blocks(&heap, &ctx(), file).unwrap();
529 let id = HeapId::parse(&in_heap(record).1, &heap, &ctx()).unwrap();
530 read_heap_object(&id, &heap, &ctx(), &blocks, file).unwrap()
531 }
532
533 #[test]
536 fn a_list_index_round_trips_every_body() {
537 let messages = vec![
538 message(MSG_DATASPACE, 1, 24, 5),
539 message(MSG_DATATYPE, 40, 20, 1),
540 message(MSG_DATASPACE, 90, 24, 1),
541 message(MSG_ATTRIBUTE, 7, 56, 4),
542 ];
543 let (mut file, built) = lay_out(&[SohmIndexContent {
544 spec: spec(50),
545 messages: messages.clone(),
546 }]);
547
548 let table = SohmMasterTable::decode(
549 &file
550 .read_block(built.table_addr, SohmMasterTable::encoded_size(&ctx(), 1))
551 .unwrap(),
552 &ctx(),
553 1,
554 )
555 .unwrap();
556 let header = &table.indexes[0];
557 assert_eq!(header.index_type, SOHM_INDEX_LIST);
558 assert_eq!(header.num_messages, 4);
559 assert_ne!(header.index_addr, UNDEF_ADDR);
560
561 let records = read_records(&mut file, header);
562 assert_eq!(records.len(), 4);
563 for (message, record) in messages.iter().zip(&records) {
564 assert_eq!(record.hash, message_hash(&message.body, message.msg_type));
565 assert_eq!(in_heap(record).0, message.ref_count);
566 assert_eq!(read_body(&mut file, header.heap_addr, record), message.body);
567 }
568
569 for message in &messages {
571 let id = built.heap_ids[&(message.msg_type, message.body.clone())];
572 let pointer =
573 SharedMessagePointer::decode(&SharedMessagePointer::encode_sohm(id), &ctx())
574 .unwrap();
575 assert_eq!(pointer.location, SharedLocation::Sohm);
576 assert_eq!(pointer.heap_id, id);
577 }
578 }
579
580 #[test]
583 fn a_zero_list_max_index_is_a_btree() {
584 let messages = vec![
585 message(MSG_DATASPACE, 1, 24, 5),
586 message(MSG_DATATYPE, 40, 20, 1),
587 message(MSG_ATTRIBUTE, 7, 56, 4),
588 ];
589 let (mut file, built) = lay_out(&[SohmIndexContent {
590 spec: SohmIndexSpec {
591 list_max: 0,
592 btree_min: 0,
593 ..spec(0)
594 },
595 messages: messages.clone(),
596 }]);
597 let table = SohmMasterTable::decode(
598 &file
599 .read_block(built.table_addr, SohmMasterTable::encoded_size(&ctx(), 1))
600 .unwrap(),
601 &ctx(),
602 1,
603 )
604 .unwrap();
605 let header = &table.indexes[0];
606 assert_eq!(header.index_type, SOHM_INDEX_BTREE);
607
608 let records = read_records(&mut file, header);
609 assert_eq!(records.len(), 3);
610 assert!(
613 records.windows(2).all(|w| w[0].hash <= w[1].hash),
614 "records are not hash-ordered: {records:?}"
615 );
616 for message in &messages {
617 let hash = message_hash(&message.body, message.msg_type);
618 let record = records.iter().find(|r| r.hash == hash).unwrap();
619 assert_eq!(in_heap(record).0, message.ref_count);
620 assert_eq!(read_body(&mut file, header.heap_addr, record), message.body);
621 }
622 }
623
624 #[test]
628 fn an_index_past_its_list_maximum_is_a_btree() {
629 let messages: Vec<SharedMessage> = (0..200u32)
630 .map(|i| SharedMessage {
631 msg_type: MSG_DATASPACE,
632 body: i.to_le_bytes().repeat(6),
633 nested: Vec::new(),
634 ref_count: i + 1,
635 ohdr_addr: None,
636 })
637 .collect();
638 let (mut file, built) = lay_out(&[SohmIndexContent {
639 spec: spec(50),
640 messages: messages.clone(),
641 }]);
642 let table = SohmMasterTable::decode(
643 &file
644 .read_block(built.table_addr, SohmMasterTable::encoded_size(&ctx(), 1))
645 .unwrap(),
646 &ctx(),
647 1,
648 )
649 .unwrap();
650 let header = &table.indexes[0];
651 assert_eq!(header.index_type, SOHM_INDEX_BTREE);
652 assert_eq!(header.num_messages, 200);
653
654 let bt2 =
655 Bt2Header::decode(&file.read_block(header.index_addr, 256).unwrap(), &ctx()).unwrap();
656 assert!(bt2.depth > 0, "expected a multi-level index, got one leaf");
657
658 let records = read_records(&mut file, header);
659 assert_eq!(records.len(), 200);
660 for message in &messages {
661 let hash = message_hash(&message.body, message.msg_type);
662 let record = records.iter().find(|r| r.hash == hash).unwrap();
663 assert_eq!(read_body(&mut file, header.heap_addr, record), message.body);
664 }
665 }
666
667 #[test]
670 fn each_index_gets_its_own_heap() {
671 let (mut file, built) = lay_out(&[
672 SohmIndexContent {
673 spec: SohmIndexSpec {
674 mesg_types: 1 << MSG_ATTRIBUTE,
675 ..spec(50)
676 },
677 messages: vec![message(MSG_ATTRIBUTE, 3, 40, 2)],
678 },
679 SohmIndexContent {
680 spec: SohmIndexSpec {
681 mesg_types: (1 << MSG_DATATYPE) | (1 << MSG_DATASPACE),
682 ..spec(50)
683 },
684 messages: vec![message(MSG_DATATYPE, 9, 20, 3)],
685 },
686 ]);
687 let table = SohmMasterTable::decode(
688 &file
689 .read_block(built.table_addr, SohmMasterTable::encoded_size(&ctx(), 2))
690 .unwrap(),
691 &ctx(),
692 2,
693 )
694 .unwrap();
695 assert_eq!(table.indexes.len(), 2);
696 assert_ne!(table.indexes[0].heap_addr, table.indexes[1].heap_addr);
697 assert_eq!(
698 table.heap_addr(MSG_ATTRIBUTE),
699 Some(table.indexes[0].heap_addr)
700 );
701 assert_eq!(
702 table.heap_addr(MSG_DATASPACE),
703 Some(table.indexes[1].heap_addr)
704 );
705 }
706
707 #[test]
710 fn an_empty_index_is_still_laid_out() {
711 let (mut file, built) = lay_out(&[SohmIndexContent {
712 spec: spec(50),
713 messages: Vec::new(),
714 }]);
715 let table = SohmMasterTable::decode(
716 &file
717 .read_block(built.table_addr, SohmMasterTable::encoded_size(&ctx(), 1))
718 .unwrap(),
719 &ctx(),
720 1,
721 )
722 .unwrap();
723 assert_eq!(table.indexes[0].num_messages, 0);
724 assert_ne!(table.indexes[0].heap_addr, UNDEF_ADDR);
725 assert!(built.heap_ids.is_empty());
726 }
727}