1#[cfg(not(feature = "std"))]
7use alloc::{string::String, string::ToString, vec, vec::Vec};
8
9use crate::attribute::AttributeMessage;
10use crate::chunked_write::{ChunkOptions, build_chunked_data_at_ext};
11use crate::dataspace::{Dataspace, DataspaceType};
12use crate::error::FormatError;
13use crate::link_message::{LinkMessage, LinkTarget};
14use crate::message_type::MessageType;
15use crate::metadata_index::{DatasetMetadata, MetadataBlock, MetadataIndex};
16use crate::object_header_writer::ObjectHeaderWriter;
17use crate::superblock::Superblock;
18use crate::type_builders::{
19 build_attr_message, DatasetBuilder, FinishedGroup, GroupBuilder,
20};
21
22pub use crate::type_builders::{AttrValue, CompoundTypeBuilder, EnumTypeBuilder};
24#[cfg(feature = "provenance")]
25pub use crate::type_builders::ProvenanceConfig;
26
27use crate::datatype::{CharacterSet, Datatype};
28
29pub(crate) const OFFSET_SIZE: u8 = 8;
30pub(crate) const LENGTH_SIZE: u8 = 8;
31const SUPERBLOCK_SIZE: usize = 48;
32
33const DENSE_ATTR_THRESHOLD: usize = 8;
35
36pub(crate) fn build_chunked_dataset_oh(
39 dt: &Datatype,
40 ds: &Dataspace,
41 layout_message: &[u8],
42 pipeline_message: Option<&[u8]>,
43 attrs: &[AttributeMessage],
44 dense_blob: Option<&DenseAttrBlob>,
45) -> Vec<u8> {
46 let mut w = ObjectHeaderWriter::new();
47 w.add_message_with_flags(MessageType::Datatype, dt.serialize(), 0x01);
48 w.add_message(MessageType::Dataspace, ds.serialize(LENGTH_SIZE));
49 w.add_message_with_flags(MessageType::FillValue, vec![3, 0x0a], 0x01);
50 w.add_message(MessageType::DataLayout, layout_message.to_vec());
51 if let Some(pm) = pipeline_message {
52 w.add_message(MessageType::FilterPipeline, pm.to_vec());
53 }
54 if let Some(blob) = dense_blob {
55 w.add_message(MessageType::AttributeInfo, blob.attr_info_message.clone());
56 } else {
57 for attr in attrs {
58 w.add_message(MessageType::Attribute, attr.serialize(LENGTH_SIZE));
59 }
60 }
61 w.serialize()
62}
63
64pub(crate) fn build_dataset_oh(
65 dt: &Datatype,
66 ds: &Dataspace,
67 data_addr: u64,
68 data_size: u64,
69 attrs: &[AttributeMessage],
70 dense_blob: Option<&DenseAttrBlob>,
71) -> Vec<u8> {
72 let mut w = ObjectHeaderWriter::new();
73 w.add_message_with_flags(MessageType::Datatype, dt.serialize(), 0x01);
74 w.add_message(MessageType::Dataspace, ds.serialize(LENGTH_SIZE));
75 w.add_message_with_flags(MessageType::FillValue, vec![3, 0x0a], 0x01);
76 let mut dl = Vec::new();
77 dl.push(4); dl.push(1); dl.extend_from_slice(&data_addr.to_le_bytes());
80 dl.extend_from_slice(&data_size.to_le_bytes());
81 w.add_message(MessageType::DataLayout, dl);
82 if let Some(blob) = dense_blob {
83 w.add_message(MessageType::AttributeInfo, blob.attr_info_message.clone());
84 } else {
85 for attr in attrs {
86 w.add_message(MessageType::Attribute, attr.serialize(LENGTH_SIZE));
87 }
88 }
89 w.serialize()
90}
91
92pub(crate) fn build_group_oh(
93 links: &[LinkMessage],
94 attrs: &[AttributeMessage],
95 dense_blob: Option<&DenseAttrBlob>,
96) -> Vec<u8> {
97 let mut w = ObjectHeaderWriter::new();
98 let mut li = Vec::new();
99 li.push(0); li.push(0); li.extend_from_slice(&u64::MAX.to_le_bytes()); li.extend_from_slice(&u64::MAX.to_le_bytes()); w.add_message(MessageType::LinkInfo, li);
104 for link in links {
105 w.add_message(MessageType::Link, link.serialize(OFFSET_SIZE));
106 }
107 if let Some(blob) = dense_blob {
108 w.add_message(MessageType::AttributeInfo, blob.attr_info_message.clone());
109 } else {
110 for attr in attrs {
111 w.add_message(MessageType::Attribute, attr.serialize(LENGTH_SIZE));
112 }
113 }
114 w.serialize()
115}
116
117pub(crate) fn make_link(name: &str, addr: u64) -> LinkMessage {
118 LinkMessage {
119 name: name.to_string(),
120 link_target: LinkTarget::Hard {
121 object_header_address: addr,
122 },
123 creation_order: None,
124 charset: CharacterSet::Ascii,
125 }
126}
127
128pub(crate) struct DenseAttrBlob {
132 pub(crate) attr_info_message: Vec<u8>,
134 pub(crate) blob: Vec<u8>,
136}
137
138pub(crate) fn build_dense_attrs(attrs: &[AttributeMessage], base_address: u64) -> DenseAttrBlob {
140 let serialized: Vec<Vec<u8>> = attrs.iter().map(|a| a.serialize_v3(LENGTH_SIZE)).collect();
142
143 let name_hashes: Vec<u32> = attrs
144 .iter()
145 .map(|a| crate::checksum::jenkins_lookup3(a.name.as_bytes()))
146 .collect();
147
148 let os = OFFSET_SIZE as usize;
149 let ls = LENGTH_SIZE as usize;
150 let max_heap_size: u16 = 40;
151 let block_offset_bytes = (max_heap_size as usize).div_ceil(8); let heap_id_length: u16 = 8;
153 let max_direct_block_size: u64 = 65536;
154
155 let dblock_header_size = 4 + 1 + os + block_offset_bytes + 4; let total_data_size: usize = serialized.iter().map(|s| s.len()).sum();
159 let dblock_content_size = dblock_header_size + total_data_size;
160 let starting_block_size = dblock_content_size.next_power_of_two().max(512) as u64;
161
162 let frhp_size = 4 + 1 + 2 + 2 + 1 + 4
164 + ls + os + ls + os + ls + ls + ls + ls + ls + ls + ls + ls
165 + 2 + ls + ls + 2 + 2 + os + 2 + 4;
166
167 let frhp_addr = base_address;
168 let dblock_addr = frhp_addr + frhp_size as u64;
169 let btree_addr = dblock_addr + starting_block_size;
170
171 let data_space = starting_block_size as usize - dblock_header_size;
172 let free_space = data_space - total_data_size;
173
174 let mut frhp = Vec::with_capacity(frhp_size);
176 frhp.extend_from_slice(b"FRHP");
177 frhp.push(0); frhp.extend_from_slice(&heap_id_length.to_le_bytes());
179 frhp.extend_from_slice(&0u16.to_le_bytes()); frhp.push(0x02); let max_managed = max_direct_block_size as u32 - dblock_header_size as u32;
182 frhp.extend_from_slice(&max_managed.to_le_bytes());
183 write_length(&mut frhp, 0, LENGTH_SIZE); write_undef_offset(&mut frhp, OFFSET_SIZE); write_length(&mut frhp, free_space as u64, LENGTH_SIZE); write_undef_offset(&mut frhp, OFFSET_SIZE); write_length(&mut frhp, starting_block_size, LENGTH_SIZE); write_length(&mut frhp, starting_block_size, LENGTH_SIZE); write_length(&mut frhp, 0, LENGTH_SIZE); write_length(&mut frhp, attrs.len() as u64, LENGTH_SIZE); write_length(&mut frhp, 0, LENGTH_SIZE); write_length(&mut frhp, 0, LENGTH_SIZE); write_length(&mut frhp, 0, LENGTH_SIZE); write_length(&mut frhp, 0, LENGTH_SIZE); frhp.extend_from_slice(&4u16.to_le_bytes()); write_length(&mut frhp, starting_block_size, LENGTH_SIZE);
197 write_length(&mut frhp, max_direct_block_size, LENGTH_SIZE); frhp.extend_from_slice(&max_heap_size.to_le_bytes());
199 let sri: u16 = 1;
200 frhp.extend_from_slice(&sri.to_le_bytes()); write_offset(&mut frhp, dblock_addr, OFFSET_SIZE);
202 frhp.extend_from_slice(&0u16.to_le_bytes()); let frhp_checksum = crate::checksum::jenkins_lookup3(&frhp);
204 frhp.extend_from_slice(&frhp_checksum.to_le_bytes());
205 debug_assert_eq!(frhp.len(), frhp_size);
206
207 let mut dblock = Vec::with_capacity(starting_block_size as usize);
209 dblock.extend_from_slice(b"FHDB");
210 dblock.push(0); write_offset(&mut dblock, frhp_addr, OFFSET_SIZE);
212 dblock.extend_from_slice(&vec![0u8; block_offset_bytes]); let cksum_pos = dblock.len();
214 dblock.extend_from_slice(&[0u8; 4]); debug_assert_eq!(dblock.len(), dblock_header_size);
216
217 let mut attr_offsets: Vec<(u64, u64)> = Vec::with_capacity(attrs.len());
219 for s in &serialized {
220 let offset_in_heap = dblock.len() as u64;
221 attr_offsets.push((offset_in_heap, s.len() as u64));
222 dblock.extend_from_slice(s);
223 }
224
225 dblock.resize(starting_block_size as usize, 0);
227
228 let dblock_checksum = crate::checksum::jenkins_lookup3(&dblock);
230 dblock[cksum_pos..cksum_pos + 4].copy_from_slice(&dblock_checksum.to_le_bytes());
231 debug_assert_eq!(dblock.len(), starting_block_size as usize);
232
233 let heap_ids: Vec<Vec<u8>> = attr_offsets
235 .iter()
236 .map(|(off, len)| encode_managed_id(*off, *len, max_heap_size, heap_id_length))
237 .collect();
238
239 let record_size: u16 = heap_id_length + 1 + 4 + 4;
241 let mut records: Vec<(u32, u32, Vec<u8>)> = Vec::with_capacity(attrs.len());
242 for (i, heap_id) in heap_ids.iter().enumerate() {
243 let mut rec = Vec::with_capacity(record_size as usize);
244 rec.extend_from_slice(heap_id);
245 rec.push(0); rec.extend_from_slice(&(i as u32).to_le_bytes()); rec.extend_from_slice(&name_hashes[i].to_le_bytes()); records.push((name_hashes[i], i as u32, rec));
249 }
250 records.sort_by(|a, b| a.0.cmp(&b.0).then(a.1.cmp(&b.1)));
251
252 let bthd_size = 4 + 1 + 1 + 4 + 2 + 2 + 1 + 1 + os + 2 + ls + 4;
253 let num_records = attrs.len();
254 let btlf_size = 4 + 1 + 1 + (num_records * record_size as usize) + 4;
255 let node_size = btlf_size.next_power_of_two().max(512) as u32;
256
257 let bthd_addr = btree_addr;
258 let btlf_addr = bthd_addr + bthd_size as u64;
259
260 let mut bthd = Vec::with_capacity(bthd_size);
261 bthd.extend_from_slice(b"BTHD");
262 bthd.push(0); bthd.push(8); bthd.extend_from_slice(&node_size.to_le_bytes());
265 bthd.extend_from_slice(&record_size.to_le_bytes());
266 bthd.extend_from_slice(&0u16.to_le_bytes()); bthd.push(100); bthd.push(40); write_offset(&mut bthd, btlf_addr, OFFSET_SIZE);
270 bthd.extend_from_slice(&(num_records as u16).to_le_bytes());
271 write_length(&mut bthd, num_records as u64, LENGTH_SIZE);
272 let bthd_checksum = crate::checksum::jenkins_lookup3(&bthd);
273 bthd.extend_from_slice(&bthd_checksum.to_le_bytes());
274 debug_assert_eq!(bthd.len(), bthd_size);
275
276 let mut btlf = Vec::with_capacity(node_size as usize);
277 btlf.extend_from_slice(b"BTLF");
278 btlf.push(0); btlf.push(8); for (_, _, rec) in &records {
281 btlf.extend_from_slice(rec);
282 }
283 let btlf_checksum = crate::checksum::jenkins_lookup3(&btlf);
286 btlf.extend_from_slice(&btlf_checksum.to_le_bytes());
287 btlf.resize(node_size as usize, 0);
289
290 let mut blob = Vec::with_capacity(frhp.len() + dblock.len() + bthd.len() + btlf.len());
291 blob.extend_from_slice(&frhp);
292 blob.extend_from_slice(&dblock);
293 blob.extend_from_slice(&bthd);
294 blob.extend_from_slice(&btlf);
295
296 let attr_info = serialize_attribute_info(frhp_addr, bthd_addr);
297
298 DenseAttrBlob {
299 attr_info_message: attr_info,
300 blob,
301 }
302}
303
304fn encode_managed_id(offset: u64, length: u64, max_heap_size: u16, id_length: u16) -> Vec<u8> {
305 let mut id = vec![0u8; id_length as usize];
306 id[0] = 0x00; let combined = offset | (length << max_heap_size);
308 let payload_len = (id_length as usize) - 1;
309 for i in 0..payload_len.min(8) {
310 id[1 + i] = ((combined >> (i * 8)) & 0xFF) as u8;
311 }
312 id
313}
314
315fn serialize_attribute_info(fh_addr: u64, btree_name_addr: u64) -> Vec<u8> {
316 let mut data = Vec::new();
317 data.push(0); data.push(0x00); data.extend_from_slice(&fh_addr.to_le_bytes());
320 data.extend_from_slice(&btree_name_addr.to_le_bytes());
321 data
322}
323
324fn write_offset(buf: &mut Vec<u8>, val: u64, offset_size: u8) {
325 match offset_size {
326 2 => buf.extend_from_slice(&(val as u16).to_le_bytes()),
327 4 => buf.extend_from_slice(&(val as u32).to_le_bytes()),
328 8 => buf.extend_from_slice(&val.to_le_bytes()),
329 _ => {}
330 }
331}
332
333fn write_length(buf: &mut Vec<u8>, val: u64, length_size: u8) {
334 write_offset(buf, val, length_size);
335}
336
337fn write_undef_offset(buf: &mut Vec<u8>, offset_size: u8) {
338 for _ in 0..offset_size {
339 buf.push(0xFF);
340 }
341}
342
343pub struct FileWriter {
347 root_datasets: Vec<DatasetBuilder>,
348 root_attrs: Vec<(String, AttrValue)>,
349 groups: Vec<FinishedGroup>,
350}
351
352impl Default for FileWriter {
353 fn default() -> Self {
354 Self::new()
355 }
356}
357
358impl FileWriter {
359 pub fn new() -> Self {
360 Self {
361 root_datasets: Vec::new(),
362 root_attrs: Vec::new(),
363 groups: Vec::new(),
364 }
365 }
366
367 pub fn create_group(&mut self, name: &str) -> GroupBuilder {
368 GroupBuilder::new(name)
369 }
370
371 pub fn add_group(&mut self, group: FinishedGroup) {
372 self.groups.push(group);
373 }
374
375 pub fn create_dataset(&mut self, name: &str) -> &mut DatasetBuilder {
376 self.root_datasets.push(DatasetBuilder::new(name));
377 self.root_datasets.last_mut().unwrap()
378 }
379
380 pub fn set_root_attr(&mut self, name: &str, value: AttrValue) {
381 self.root_attrs.push((name.to_string(), value));
382 }
383
384 pub fn finish(self) -> Result<Vec<u8>, FormatError> {
385 struct DsFlat {
386 name: String,
387 dt: Datatype,
388 ds: Dataspace,
389 raw: Vec<u8>,
390 attrs: Vec<AttributeMessage>,
391 chunk_options: ChunkOptions,
392 maxshape: Option<Vec<u64>>,
393 }
394 struct GrpFlat {
395 name: String,
396 attrs: Vec<AttributeMessage>,
397 ds_indices: Vec<usize>,
398 }
399
400 let mut all_ds: Vec<DsFlat> = Vec::new();
401 let mut groups: Vec<GrpFlat> = Vec::new();
402 let mut root_ds_indices: Vec<usize> = Vec::new();
403
404 for db in self.root_datasets {
405 let dt = db.datatype.ok_or(FormatError::DatasetMissingData)?;
406 let shape = db.shape.ok_or(FormatError::DatasetMissingShape)?;
407 let raw = db.data.ok_or(FormatError::DatasetMissingData)?;
408 let max_dimensions = db.maxshape.clone();
409 let dspace = Dataspace {
410 space_type: if shape.is_empty() { DataspaceType::Scalar } else { DataspaceType::Simple },
411 rank: shape.len() as u8, dimensions: shape, max_dimensions,
412 };
413 let mut attrs = Vec::new();
414 for (n, v) in &db.attrs { attrs.push(build_attr_message(n, v)); }
415 #[cfg(feature = "provenance")]
416 if let Some(ref prov) = db.provenance {
417 let p = crate::provenance::Provenance {
418 creator: prov.creator.clone(),
419 timestamp: prov.timestamp.clone(),
420 source: prov.source.clone(),
421 };
422 attrs.extend(p.build_attrs(&raw));
423 }
424 root_ds_indices.push(all_ds.len());
425 all_ds.push(DsFlat { name: db.name, dt, ds: dspace, raw, attrs, chunk_options: db.chunk_options, maxshape: db.maxshape });
426 }
427
428 for g in self.groups.into_iter() {
429 let mut gattrs = Vec::new();
430 for (n, v) in &g.attrs { gattrs.push(build_attr_message(n, v)); }
431 let mut ds_idx = Vec::new();
432 for db in g.datasets {
433 let dt = db.datatype.ok_or(FormatError::DatasetMissingData)?;
434 let shape = db.shape.ok_or(FormatError::DatasetMissingShape)?;
435 let raw = db.data.ok_or(FormatError::DatasetMissingData)?;
436 let max_dimensions = db.maxshape.clone();
437 let dspace = Dataspace {
438 space_type: if shape.is_empty() { DataspaceType::Scalar } else { DataspaceType::Simple },
439 rank: shape.len() as u8, dimensions: shape, max_dimensions,
440 };
441 let mut attrs = Vec::new();
442 for (n, v) in &db.attrs { attrs.push(build_attr_message(n, v)); }
443 #[cfg(feature = "provenance")]
444 if let Some(ref prov) = db.provenance {
445 let p = crate::provenance::Provenance {
446 creator: prov.creator.clone(),
447 timestamp: prov.timestamp.clone(),
448 source: prov.source.clone(),
449 };
450 attrs.extend(p.build_attrs(&raw));
451 }
452 ds_idx.push(all_ds.len());
453 all_ds.push(DsFlat { name: db.name, dt, ds: dspace, raw, attrs, chunk_options: db.chunk_options, maxshape: db.maxshape });
454 }
455 groups.push(GrpFlat { name: g.name, attrs: gattrs, ds_indices: ds_idx });
456 }
457
458 let mut root_attrs: Vec<AttributeMessage> = Vec::new();
459 for (n, v) in &self.root_attrs { root_attrs.push(build_attr_message(n, v)); }
460
461 let is_chunked: Vec<bool> = all_ds.iter().map(|d| d.chunk_options.is_chunked() || d.maxshape.is_some()).collect();
462 let root_dense = root_attrs.len() > DENSE_ATTR_THRESHOLD;
463 let group_dense: Vec<bool> = groups.iter().map(|g| g.attrs.len() > DENSE_ATTR_THRESHOLD).collect();
464 let ds_dense: Vec<bool> = all_ds.iter().map(|d| d.attrs.len() > DENSE_ATTR_THRESHOLD).collect();
465
466 let group_oh_sizes: Vec<usize> = groups.iter().enumerate().map(|(gi, g)| {
468 let dummy_links: Vec<LinkMessage> = g.ds_indices.iter().map(|&i| make_link(&all_ds[i].name, 0)).collect();
469 if group_dense[gi] {
470 let dummy_blob = build_dense_attrs(&g.attrs, 0);
471 build_group_oh(&dummy_links, &g.attrs, Some(&dummy_blob)).len()
472 } else {
473 build_group_oh(&dummy_links, &g.attrs, None).len()
474 }
475 }).collect();
476
477 let root_dummy_links: Vec<LinkMessage> = {
478 let mut links = Vec::new();
479 for &i in &root_ds_indices { links.push(make_link(&all_ds[i].name, 0)); }
480 for g in &groups { links.push(make_link(&g.name, 0)); }
481 links
482 };
483 let root_oh_size = if root_dense {
484 let dummy_blob = build_dense_attrs(&root_attrs, 0);
485 build_group_oh(&root_dummy_links, &root_attrs, Some(&dummy_blob)).len()
486 } else {
487 build_group_oh(&root_dummy_links, &root_attrs, None).len()
488 };
489
490 struct DataBlob { data: Vec<u8>, oh_bytes: Vec<u8> }
491
492 let mut dummy_blobs: Vec<DataBlob> = Vec::new();
493 let mut dummy_cursor = 0u64;
494 for (i, d) in all_ds.iter().enumerate() {
495 if is_chunked[i] {
496 let chunk_dims = d.chunk_options.resolve_chunk_dims(&d.ds.dimensions);
497 let elem_size = d.dt.type_size() as usize;
498 let result = build_chunked_data_at_ext(&d.raw, &d.ds.dimensions, &chunk_dims, elem_size, &d.chunk_options, dummy_cursor, d.maxshape.as_deref())?;
499 dummy_cursor += result.data_bytes.len() as u64;
500 let dense_blob = if ds_dense[i] { Some(build_dense_attrs(&d.attrs, 0)) } else { None };
501 let oh = build_chunked_dataset_oh(&d.dt, &d.ds, &result.layout_message, result.pipeline_message.as_deref(), &d.attrs, dense_blob.as_ref());
502 dummy_blobs.push(DataBlob { data: result.data_bytes, oh_bytes: oh });
503 } else {
504 let dense_blob = if ds_dense[i] { Some(build_dense_attrs(&d.attrs, 0)) } else { None };
505 let oh = build_dataset_oh(&d.dt, &d.ds, 0, d.raw.len() as u64, &d.attrs, dense_blob.as_ref());
506 dummy_blobs.push(DataBlob { data: d.raw.clone(), oh_bytes: oh });
507 }
508 }
509
510 let actual_ds_oh_sizes: Vec<usize> = dummy_blobs.iter().map(|b| b.oh_bytes.len()).collect();
511
512 let root_group_addr = SUPERBLOCK_SIZE as u64;
514 let mut cursor2 = SUPERBLOCK_SIZE + root_oh_size;
515
516 let root_dense_blob = if root_dense {
517 let blob = build_dense_attrs(&root_attrs, cursor2 as u64);
518 cursor2 += blob.blob.len();
519 Some(blob)
520 } else {
521 None
522 };
523
524 let mut group_dense_blobs: Vec<Option<DenseAttrBlob>> = Vec::new();
525 let group_addrs2: Vec<u64> = group_oh_sizes.iter().enumerate().map(|(gi, &sz)| {
526 let addr = cursor2 as u64;
527 cursor2 += sz;
528 if group_dense[gi] {
529 let blob = build_dense_attrs(&groups[gi].attrs, cursor2 as u64);
530 cursor2 += blob.blob.len();
531 group_dense_blobs.push(Some(blob));
532 } else {
533 group_dense_blobs.push(None);
534 }
535 addr
536 }).collect();
537
538 let mut ds_dense_blobs: Vec<Option<DenseAttrBlob>> = Vec::new();
539 let ds_oh_addrs2: Vec<u64> = actual_ds_oh_sizes.iter().enumerate().map(|(i, &sz)| {
540 let addr = cursor2 as u64;
541 cursor2 += sz;
542 if ds_dense[i] {
543 let blob = build_dense_attrs(&all_ds[i].attrs, cursor2 as u64);
544 cursor2 += blob.blob.len();
545 ds_dense_blobs.push(Some(blob));
546 } else {
547 ds_dense_blobs.push(None);
548 }
549 addr
550 }).collect();
551
552 let mut ds_blobs2: Vec<DataBlob> = Vec::new();
553 for (i, d) in all_ds.iter().enumerate() {
554 if is_chunked[i] {
555 let chunk_dims = d.chunk_options.resolve_chunk_dims(&d.ds.dimensions);
556 let elem_size = d.dt.type_size() as usize;
557 let base_address = cursor2 as u64;
558 let result = build_chunked_data_at_ext(&d.raw, &d.ds.dimensions, &chunk_dims, elem_size, &d.chunk_options, base_address, d.maxshape.as_deref())?;
559 cursor2 += result.data_bytes.len();
560 let oh = build_chunked_dataset_oh(&d.dt, &d.ds, &result.layout_message, result.pipeline_message.as_deref(), &d.attrs, ds_dense_blobs[i].as_ref());
561 ds_blobs2.push(DataBlob { data: result.data_bytes, oh_bytes: oh });
562 } else {
563 let oh = build_dataset_oh(&d.dt, &d.ds, cursor2 as u64, d.raw.len() as u64, &d.attrs, ds_dense_blobs[i].as_ref());
564 let data = d.raw.clone();
565 cursor2 += data.len();
566 ds_blobs2.push(DataBlob { data, oh_bytes: oh });
567 }
568 }
569
570 let actual_ds_oh_sizes2: Vec<usize> = ds_blobs2.iter().map(|b| b.oh_bytes.len()).collect();
571 debug_assert_eq!(actual_ds_oh_sizes, actual_ds_oh_sizes2);
572
573 let eof_addr2 = cursor2 as u64;
574 let mut buf = Vec::with_capacity(cursor2);
575
576 let sb = Superblock {
577 version: 3, offset_size: OFFSET_SIZE, length_size: LENGTH_SIZE,
578 base_address: 0, eof_address: eof_addr2, root_group_address: root_group_addr,
579 group_leaf_node_k: None, group_internal_node_k: None, indexed_storage_internal_node_k: None,
580 free_space_address: None, driver_info_address: None,
581 consistency_flags: 0, superblock_extension_address: Some(u64::MAX), checksum: None,
582 };
583 buf.extend_from_slice(&sb.serialize());
584
585 let mut root_links: Vec<LinkMessage> = Vec::new();
587 for &i in &root_ds_indices { root_links.push(make_link(&all_ds[i].name, ds_oh_addrs2[i])); }
588 for (gi, g) in groups.iter().enumerate() { root_links.push(make_link(&g.name, group_addrs2[gi])); }
589 buf.extend_from_slice(&build_group_oh(&root_links, &root_attrs, root_dense_blob.as_ref()));
590 if let Some(ref blob) = root_dense_blob { buf.extend_from_slice(&blob.blob); }
591
592 for (gi, g) in groups.iter().enumerate() {
594 let links: Vec<LinkMessage> = g.ds_indices.iter().map(|&i| make_link(&all_ds[i].name, ds_oh_addrs2[i])).collect();
595 buf.extend_from_slice(&build_group_oh(&links, &g.attrs, group_dense_blobs[gi].as_ref()));
596 if let Some(ref blob) = group_dense_blobs[gi] { buf.extend_from_slice(&blob.blob); }
597 }
598
599 for (i, blob) in ds_blobs2.iter().enumerate() {
601 buf.extend_from_slice(&blob.oh_bytes);
602 if let Some(ref dense) = ds_dense_blobs[i] { buf.extend_from_slice(&dense.blob); }
603 }
604
605 for blob in &ds_blobs2 { buf.extend_from_slice(&blob.data); }
607
608 debug_assert_eq!(buf.len(), cursor2);
609 Ok(buf)
610 }
611}
612
613pub struct IndependentDatasetBuilder {
623 block: MetadataBlock,
624}
625
626impl IndependentDatasetBuilder {
627 pub fn new(creator_id: u32) -> Self {
629 Self {
630 block: MetadataBlock::new(creator_id),
631 }
632 }
633
634 pub fn add_dataset(&mut self, meta: DatasetMetadata) {
636 self.block.add_dataset(meta);
637 }
638
639 pub fn finish(self) -> MetadataBlock {
641 self.block
642 }
643}
644
645pub fn finalize_parallel(blocks: Vec<MetadataBlock>) -> Result<Vec<u8>, FormatError> {
651 let index = MetadataIndex::merge_blocks(&blocks)?;
652 finalize_from_index(index)
653}
654
655fn finalize_from_index(index: MetadataIndex) -> Result<Vec<u8>, FormatError> {
657 let mut fw = FileWriter::new();
660 for ds_meta in &index.datasets {
661 let db = fw.create_dataset(&ds_meta.name);
662 db.datatype = Some(ds_meta.datatype.clone());
664 db.shape = Some(ds_meta.dataspace.dimensions.clone());
665 db.maxshape = ds_meta.maxshape.clone();
666 db.data = Some(ds_meta.raw_data.clone());
667 db.chunk_options = ds_meta.chunk_options.clone();
668 for (name, val) in &ds_meta.attrs {
669 db.set_attr(name, val.clone());
670 }
671 }
672 fw.finish()
673}
674
675#[cfg(test)]
676mod tests {
677 use super::*;
678 use crate::group_v2::resolve_path_any;
679 use crate::object_header::ObjectHeader;
680 use crate::signature;
681
682 fn parse_file(bytes: &[u8]) -> (Superblock, ObjectHeader) {
683 let sig = signature::find_signature(bytes).unwrap();
684 let sb = Superblock::parse(bytes, sig).unwrap();
685 let oh = ObjectHeader::parse(bytes, sb.root_group_address as usize, sb.offset_size, sb.length_size).unwrap();
686 (sb, oh)
687 }
688
689 fn read_dataset_f64(bytes: &[u8], path: &str) -> Vec<f64> {
690 let sig = signature::find_signature(bytes).unwrap();
691 let sb = Superblock::parse(bytes, sig).unwrap();
692 let addr = resolve_path_any(bytes, &sb, path).unwrap();
693 let hdr = ObjectHeader::parse(bytes, addr as usize, sb.offset_size, sb.length_size).unwrap();
694 let dt_data = &hdr.messages.iter().find(|m| m.msg_type == MessageType::Datatype).unwrap().data;
695 let ds_data = &hdr.messages.iter().find(|m| m.msg_type == MessageType::Dataspace).unwrap().data;
696 let dl_data = &hdr.messages.iter().find(|m| m.msg_type == MessageType::DataLayout).unwrap().data;
697 let (dt, _) = Datatype::parse(dt_data).unwrap();
698 let ds = Dataspace::parse(ds_data, sb.length_size).unwrap();
699 let dl = crate::data_layout::DataLayout::parse(dl_data, sb.offset_size, sb.length_size).unwrap();
700 let raw = crate::data_read::read_raw_data(bytes, &dl, &ds, &dt).unwrap();
701 crate::data_read::read_as_f64(&raw, &dt).unwrap()
702 }
703
704 #[test]
705 fn empty_file_root_group_only() {
706 let fw = FileWriter::new();
707 let bytes = fw.finish().unwrap();
708 let (sb, oh) = parse_file(&bytes);
709 assert_eq!(sb.version, 3);
710 assert_eq!(oh.version, 2);
711 }
712
713 #[test]
714 fn file_with_f64_dataset() {
715 let mut fw = FileWriter::new();
716 fw.create_dataset("data").with_f64_data(&[1.0, 2.0, 3.0]);
717 let bytes = fw.finish().unwrap();
718 assert_eq!(read_dataset_f64(&bytes, "data"), vec![1.0, 2.0, 3.0]);
719 }
720
721 #[test]
722 fn file_with_dataset_attrs() {
723 let mut fw = FileWriter::new();
724 fw.create_dataset("data").with_f64_data(&[1.0, 2.0]).set_attr("scale", AttrValue::F64(0.5));
725 let bytes = fw.finish().unwrap();
726 assert_eq!(read_dataset_f64(&bytes, "data"), vec![1.0, 2.0]);
727 let sig = signature::find_signature(&bytes).unwrap();
728 let sb = Superblock::parse(&bytes, sig).unwrap();
729 let addr = resolve_path_any(&bytes, &sb, "data").unwrap();
730 let hdr = ObjectHeader::parse(&bytes, addr as usize, sb.offset_size, sb.length_size).unwrap();
731 let attrs = crate::attribute::extract_attributes(&hdr, sb.length_size).unwrap();
732 assert_eq!(attrs.len(), 1);
733 assert_eq!(attrs[0].name, "scale");
734 }
735
736 #[test]
737 fn file_with_group_and_dataset() {
738 let mut fw = FileWriter::new();
739 let mut gb = fw.create_group("grp");
740 gb.create_dataset("vals").with_f64_data(&[10.0, 20.0]);
741 fw.add_group(gb.finish());
742 let bytes = fw.finish().unwrap();
743 assert_eq!(read_dataset_f64(&bytes, "grp/vals"), vec![10.0, 20.0]);
744 }
745
746 #[test]
747 fn file_with_root_attr() {
748 let mut fw = FileWriter::new();
749 fw.set_root_attr("version", AttrValue::I64(42));
750 let bytes = fw.finish().unwrap();
751 let (sb, oh) = parse_file(&bytes);
752 let attrs = crate::attribute::extract_attributes(&oh, sb.length_size).unwrap();
753 assert_eq!(attrs[0].name, "version");
754 }
755
756 #[test]
757 fn dense_attrs_self_roundtrip() {
758 let mut fw = FileWriter::new();
759 let ds = fw.create_dataset("data");
760 ds.with_f64_data(&[1.0, 2.0, 3.0]);
761 for i in 0..20 {
762 ds.set_attr(&format!("attr_{i:03}"), AttrValue::F64(i as f64 * 1.5));
763 }
764 let bytes = fw.finish().unwrap();
765 let sig = signature::find_signature(&bytes).unwrap();
766 let sb = Superblock::parse(&bytes, sig).unwrap();
767 let addr = resolve_path_any(&bytes, &sb, "data").unwrap();
768 let hdr = ObjectHeader::parse(&bytes, addr as usize, sb.offset_size, sb.length_size).unwrap();
769 let attrs = crate::attribute::extract_attributes_full(&bytes, &hdr, sb.offset_size, sb.length_size).unwrap();
770 assert_eq!(attrs.len(), 20);
771 for i in 0..20 {
772 let attr = attrs.iter().find(|a| a.name == format!("attr_{i:03}")).unwrap();
773 let v = attr.read_as_f64().unwrap();
774 assert!((v[0] - i as f64 * 1.5).abs() < 1e-10);
775 }
776 assert_eq!(read_dataset_f64(&bytes, "data"), vec![1.0, 2.0, 3.0]);
777 }
778
779 #[test]
780 fn dense_attrs_root_group_self_roundtrip() {
781 let mut fw = FileWriter::new();
782 fw.create_dataset("dummy").with_f64_data(&[0.0]);
783 for i in 0..15 {
784 fw.set_root_attr(&format!("root_{i:02}"), AttrValue::F64(i as f64 * 2.0));
785 }
786 let bytes = fw.finish().unwrap();
787 let sig = signature::find_signature(&bytes).unwrap();
788 let sb = Superblock::parse(&bytes, sig).unwrap();
789 let oh = ObjectHeader::parse(&bytes, sb.root_group_address as usize, sb.offset_size, sb.length_size).unwrap();
790 let attrs = crate::attribute::extract_attributes_full(&bytes, &oh, sb.offset_size, sb.length_size).unwrap();
791 assert_eq!(attrs.len(), 15);
792 }
793
794 #[test]
795 fn inline_attrs_below_threshold() {
796 let mut fw = FileWriter::new();
797 let ds = fw.create_dataset("data");
798 ds.with_f64_data(&[1.0]);
799 for i in 0..5 { ds.set_attr(&format!("a{i}"), AttrValue::F64(i as f64)); }
800 let bytes = fw.finish().unwrap();
801 let sig = signature::find_signature(&bytes).unwrap();
802 let sb = Superblock::parse(&bytes, sig).unwrap();
803 let addr = resolve_path_any(&bytes, &sb, "data").unwrap();
804 let hdr = ObjectHeader::parse(&bytes, addr as usize, sb.offset_size, sb.length_size).unwrap();
805 assert!(!hdr.messages.iter().any(|m| m.msg_type == MessageType::AttributeInfo));
806 let attrs = crate::attribute::extract_attributes(&hdr, sb.length_size).unwrap();
807 assert_eq!(attrs.len(), 5);
808 }
809
810 #[test]
811 fn encode_decode_managed_id_roundtrip() {
812 let id = encode_managed_id(100, 42, 40, 8);
813 let fh = crate::fractal_heap::FractalHeapHeader {
814 heap_id_length: 8, io_filter_encoded_length: 0,
815 max_managed_object_size: 1024, table_width: 4,
816 starting_block_size: 4096, max_direct_block_size: 65536,
817 max_heap_size: 40, starting_row_of_indirect_blocks: 1,
818 root_block_address: 0, current_rows_in_root_indirect_block: 0,
819 managed_objects_count: 0,
820 };
821 let (off, len) = fh.decode_managed_id(&id).unwrap();
822 assert_eq!(off, 100);
823 assert_eq!(len, 42);
824 }
825
826 #[test]
827 fn finalize_parallel_basic() {
828 use crate::metadata_index::{MetadataBlock, build_dataset_metadata};
829 use crate::chunked_write::ChunkOptions;
830 use crate::type_builders::make_f64_type;
831
832 let mut b0 = MetadataBlock::new(0);
833 let data_a: Vec<u8> = [1.0f64, 2.0, 3.0].iter().flat_map(|v| v.to_le_bytes()).collect();
834 b0.add_dataset(build_dataset_metadata(
835 "alpha", make_f64_type(), vec![3], data_a,
836 ChunkOptions::default(), None, vec![],
837 ));
838
839 let mut b1 = MetadataBlock::new(1);
840 let data_b: Vec<u8> = [10.0f64, 20.0].iter().flat_map(|v| v.to_le_bytes()).collect();
841 b1.add_dataset(build_dataset_metadata(
842 "beta", make_f64_type(), vec![2], data_b,
843 ChunkOptions::default(), None, vec![],
844 ));
845
846 let bytes = finalize_parallel(vec![b0, b1]).unwrap();
847 assert_eq!(read_dataset_f64(&bytes, "alpha"), vec![1.0, 2.0, 3.0]);
848 assert_eq!(read_dataset_f64(&bytes, "beta"), vec![10.0, 20.0]);
849 }
850
851 #[test]
852 fn finalize_parallel_duplicate_error() {
853 use crate::metadata_index::{MetadataBlock, build_dataset_metadata};
854 use crate::chunked_write::ChunkOptions;
855 use crate::type_builders::make_f64_type;
856
857 let mut b0 = MetadataBlock::new(0);
858 b0.add_dataset(build_dataset_metadata(
859 "dup", make_f64_type(), vec![1], vec![0u8; 8],
860 ChunkOptions::default(), None, vec![],
861 ));
862 let mut b1 = MetadataBlock::new(1);
863 b1.add_dataset(build_dataset_metadata(
864 "dup", make_f64_type(), vec![1], vec![0u8; 8],
865 ChunkOptions::default(), None, vec![],
866 ));
867 let err = finalize_parallel(vec![b0, b1]).unwrap_err();
868 assert!(matches!(err, FormatError::DuplicateDatasetName(_)));
869 }
870}