#[cfg(not(feature = "std"))]
use alloc::{string::String, string::ToString, vec, vec::Vec};
#[cfg(not(feature = "std"))]
use alloc::format;
#[cfg(not(feature = "std"))]
use alloc::collections::BTreeMap as HashMap;
#[cfg(feature = "std")]
use std::collections::HashMap;
use crate::address::BaseAddress;
use crate::attribute::AttributeMessage;
use crate::btree_v2_write::{self, BTreeV2Plan};
use crate::chunked_write::{
ByteSink, ChunkOptions, ChunkProvider, ChunkedMeasure, CompressedChunkSet, StorageAllocation,
VerbatimLayout, VerbatimPlan, assemble_chunked_at, compress_chunks, emit_chunked_data_verbatim,
measure_chunked_at, plan_chunked_data_verbatim,
};
use crate::convert::TryToUsize;
use crate::dataspace::{Dataspace, DataspaceType};
use crate::error::{FormatError, OBJECT_HEADER_MESSAGE_MAX};
use crate::file_create_properties::FileCreateProperties;
use crate::file_space_info::{
DEFAULT_PAGE_SIZE, DEFAULT_THRESHOLD, FileSpaceInfo, FileSpaceStrategy, NUM_FILE_FSM_MANAGERS,
};
use crate::fractal_heap_write::{self, ManagedPlan, PlanRefusal};
use crate::free_space_manager::{
FreeSection, SECT_CLASS_LARGE, SECT_CLASS_SMALL, fshd_len, fsse_len, serialize_file_fsm,
};
use crate::libver::LibVer;
use crate::link_message::{LinkMessage, LinkTarget};
use crate::message_type::MessageType;
use crate::object_header_writer::ObjectHeaderWriter;
use crate::shared_message::DatatypeLocation;
use crate::superblock::Superblock;
use crate::type_builders::{
AttrSpec, CommittedDatatype, DatasetBuilder, FinishedGroup, GroupBuilder, VlStringStaging,
build_global_heap_collections, patch_vl_refs, patch_vl_refs_masked, write_reference_address,
};
pub use crate::type_builders::AttrValue;
use crate::datatype::{CharacterSet, Datatype};
pub(crate) const OFFSET_SIZE: u8 = 8;
pub(crate) const LENGTH_SIZE: u8 = 8;
const SUPERBLOCK_SIZE: usize = 48;
const MSG_CONSTANT: u8 = 0x01;
const MSG_SHARED: u8 = 0x02;
const MSG_DONTSHARE: u8 = 0x04;
pub(crate) const DENSE_ATTR_THRESHOLD: usize = 8;
fn align_up(value: u64, page: u64) -> u64 {
value.div_ceil(page) * page
}
pub(crate) fn build_chunked_dataset_oh(
dt: &Datatype,
dt_location: &DatatypeLocation,
ds: &Dataspace,
layout_message: &[u8],
pipeline_message: Option<&[u8]>,
attrs: &[AttributeMessage],
attr_info: Option<&[u8]>,
fill: Option<&[u8]>,
) -> Result<Vec<u8>, FormatError> {
let mut w = ObjectHeaderWriter::new();
add_datatype(&mut w, dt, dt_location);
w.add_message(MessageType::Dataspace, ds.serialize(LENGTH_SIZE));
w.add_message_with_flags(
MessageType::FillValue,
crate::fill_value::fill_value_message_v3(fill),
0x01,
);
w.add_message(MessageType::DataLayout, layout_message.to_vec());
if let Some(pm) = pipeline_message {
w.add_message(MessageType::FilterPipeline, pm.to_vec());
}
add_attributes(&mut w, attrs, attr_info);
w.serialize()
}
pub(crate) fn superblock_version(libver: LibVer) -> u8 {
if libver >= LibVer::V110 { 3 } else { 2 }
}
pub(crate) fn contiguous_layout_version(libver: LibVer) -> u8 {
if libver >= LibVer::V110 { 4 } else { 3 }
}
pub(crate) fn build_dataset_oh(
dt: &Datatype,
dt_location: &DatatypeLocation,
ds: &Dataspace,
data_addr: u64,
data_size: u64,
attrs: &[AttributeMessage],
attr_info: Option<&[u8]>,
fill: Option<&[u8]>,
libver: LibVer,
) -> Result<Vec<u8>, FormatError> {
let mut w = ObjectHeaderWriter::new();
add_datatype(&mut w, dt, dt_location);
w.add_message(MessageType::Dataspace, ds.serialize(LENGTH_SIZE));
w.add_message_with_flags(
MessageType::FillValue,
crate::fill_value::fill_value_message_v3(fill),
0x01,
);
let mut dl = Vec::new();
dl.push(contiguous_layout_version(libver));
dl.push(1); dl.extend_from_slice(&data_addr.to_le_bytes());
dl.extend_from_slice(&data_size.to_le_bytes());
w.add_message(MessageType::DataLayout, dl);
add_attributes(&mut w, attrs, attr_info);
w.serialize()
}
fn add_datatype(w: &mut ObjectHeaderWriter, dt: &Datatype, location: &DatatypeLocation) {
match location.reference_bytes(OFFSET_SIZE) {
Some(reference) => {
w.add_message_with_flags(MessageType::Datatype, reference, MSG_CONSTANT | MSG_SHARED);
}
None => w.add_message_with_flags(MessageType::Datatype, dt.serialize(), MSG_CONSTANT),
}
}
pub(crate) fn build_committed_datatype_oh(
dt: &Datatype,
references: u32,
) -> Result<Vec<u8>, FormatError> {
let mut w = ObjectHeaderWriter::new();
w.add_message_with_flags(
MessageType::Datatype,
dt.serialize(),
MSG_CONSTANT | MSG_DONTSHARE,
);
if references > 1 {
let mut refcount = Vec::with_capacity(5);
refcount.push(0); refcount.extend_from_slice(&references.to_le_bytes());
w.add_message_with_flags(MessageType::ObjectReferenceCount, refcount, MSG_DONTSHARE);
}
w.serialize()
}
pub(crate) fn build_group_oh(
links: &[LinkMessage],
attrs: &[AttributeMessage],
attr_info: Option<&[u8]>,
) -> Result<Vec<u8>, FormatError> {
let mut w = ObjectHeaderWriter::new();
let mut li = Vec::new();
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);
w.add_message(MessageType::GroupInfo, vec![0, 0]);
for link in links {
w.add_message(MessageType::Link, link.serialize(OFFSET_SIZE));
}
add_attributes(&mut w, attrs, attr_info);
w.serialize()
}
fn add_attributes(
w: &mut ObjectHeaderWriter,
attrs: &[AttributeMessage],
attr_info: Option<&[u8]>,
) {
if let Some(message) = attr_info {
w.add_message(MessageType::AttributeInfo, message.to_vec());
} else {
if !attrs.is_empty() {
w.add_message(MessageType::AttributeInfo, compact_attribute_info_message());
}
for attr in attrs {
w.add_message(MessageType::Attribute, attr.serialize(LENGTH_SIZE));
}
}
}
pub(crate) fn make_link(name: &str, addr: u64) -> LinkMessage {
LinkMessage {
name: name.to_string(),
link_target: LinkTarget::Hard {
object_header_address: addr,
},
creation_order: None,
charset: CharacterSet::Ascii,
}
}
pub(crate) struct DenseAttrBlob {
pub(crate) attr_info_message: Vec<u8>,
pub(crate) blob: Vec<u8>,
}
const DENSE_ATTR_MAX_HEAP_SIZE_BITS: u16 = fractal_heap_write::MAX_HEAP_SIZE_BITS;
const DENSE_ATTR_BLOCK_OFFSET_BYTES: usize = fractal_heap_write::BLOCK_OFFSET_BYTES;
pub(crate) fn needs_dense_attrs(attrs: &[AttributeMessage]) -> bool {
attrs.len() > DENSE_ATTR_THRESHOLD
|| attrs
.iter()
.any(|a| a.serialize(LENGTH_SIZE).len() > OBJECT_HEADER_MESSAGE_MAX)
}
pub(crate) const DENSE_ATTR_MAX_MANAGED_OBJECT: usize =
fractal_heap_write::max_managed_object(OFFSET_SIZE);
const DENSE_ATTR_BTREE_RECORD: u16 = 8 + 1 + 4 + 4;
const DENSE_ATTR_HUGE_BTREE_RECORD: u16 =
OFFSET_SIZE as u16 + LENGTH_SIZE as u16 + LENGTH_SIZE as u16;
const DENSE_ATTR_NAME_BTREE_TYPE: u8 = 8;
const DENSE_ATTR_HUGE_BTREE_TYPE: u8 = 1;
pub(crate) fn dense_attrs_check(attrs: &[AttributeMessage]) -> Result<(), FormatError> {
let mut managed = Vec::new();
for a in attrs {
if let Some((field, size)) = a.v3_header_field_overflow(LENGTH_SIZE) {
return Err(FormatError::AttributeFieldTooLong {
name: a.name.clone(),
field,
size,
limit: u16::MAX as usize,
});
}
let size = a.serialize_v3(LENGTH_SIZE).len();
if size <= DENSE_ATTR_MAX_MANAGED_OBJECT {
managed.push(size as u64);
}
}
match ManagedPlan::new(&managed, OFFSET_SIZE) {
Ok(_) => Ok(()),
Err(PlanRefusal::HeapSpace) => Err(FormatError::DenseAttributeHeapTooLarge {
limit: fractal_heap_write::MAX_HEAP_SPACE,
}),
Err(PlanRefusal::Host { bytes }) => Err(FormatError::ValueTooLargeForPlatform {
value: bytes,
target: "usize",
}),
}
}
pub(crate) struct DenseAttrPlan {
serialized: Vec<Vec<u8>>,
huge_id_of: Vec<Option<u64>>,
huge_count: usize,
huge_total: u64,
managed_plan: ManagedPlan,
name_plan: BTreeV2Plan,
huge_plan: Option<BTreeV2Plan>,
order: Vec<(u32, u32)>,
managed_off: u64,
btree_off: u64,
name_nodes_off: u64,
huge_bthd_off: u64,
huge_nodes_off: u64,
huge_data_off: u64,
total_len: u64,
}
pub(crate) fn dense_attrs_plan(attrs: &[AttributeMessage]) -> DenseAttrPlan {
let serialized: Vec<Vec<u8>> = attrs.iter().map(|a| a.serialize_v3(LENGTH_SIZE)).collect();
let name_hashes: Vec<u32> = attrs
.iter()
.map(|a| crate::checksum::jenkins_lookup3(a.name.as_bytes()))
.collect();
let mut huge_id_of: Vec<Option<u64>> = vec![None; attrs.len()];
let mut huge_count: usize = 0;
for (slot, bytes) in huge_id_of.iter_mut().zip(&serialized) {
if bytes.len() > DENSE_ATTR_MAX_MANAGED_OBJECT {
huge_count += 1;
*slot = Some(huge_count as u64);
}
}
let huge_total: u64 = serialized
.iter()
.zip(&huge_id_of)
.filter(|(_, id)| id.is_some())
.map(|(s, _)| s.len() as u64)
.sum();
let managed_sizes: Vec<u64> = serialized
.iter()
.zip(&huge_id_of)
.filter(|(_, id)| id.is_none())
.map(|(s, _)| s.len() as u64)
.collect();
let managed_plan = ManagedPlan::new(&managed_sizes, OFFSET_SIZE)
.expect("dense_attrs_check, which every caller must run first, plans the same layout");
let bthd_size = btree_v2_write::header_size(OFFSET_SIZE, LENGTH_SIZE);
debug_assert_eq!(
bthd_size,
4 + 1 + 1 + 4 + 2 + 2 + 1 + 1 + OFFSET_SIZE as usize + 2 + LENGTH_SIZE as usize + 4
);
let name_plan = BTreeV2Plan::new(
DENSE_ATTR_NAME_BTREE_TYPE,
attrs.len(),
DENSE_ATTR_BTREE_RECORD,
btree_v2_write::NODE_SIZE,
OFFSET_SIZE,
)
.expect("a 512-byte node holds 29 name records, enough to plan any count");
let huge_plan = (huge_count > 0).then(|| {
BTreeV2Plan::new(
DENSE_ATTR_HUGE_BTREE_TYPE,
huge_count,
DENSE_ATTR_HUGE_BTREE_RECORD,
btree_v2_write::NODE_SIZE,
OFFSET_SIZE,
)
.expect("a 512-byte node holds 20 huge records, enough to plan any count")
});
let managed_off = DENSE_ATTR_FRHP_SIZE as u64;
let btree_off = managed_off + managed_plan.region_size();
let name_nodes_off = btree_off + bthd_size as u64;
let huge_bthd_off = name_nodes_off + name_plan.nodes_size();
let huge_nodes_off = huge_bthd_off + bthd_size as u64;
let huge_data_off = huge_nodes_off + huge_plan.as_ref().map_or(0, BTreeV2Plan::nodes_size);
let total_len = if huge_count > 0 {
huge_data_off + huge_total
} else {
huge_bthd_off
};
#[expect(
clippy::cast_possible_truncation,
reason = "i is an attribute index bounded by the attribute count, far below u32::MAX"
)]
let mut order: Vec<(u32, u32)> = (0..attrs.len())
.map(|i| (name_hashes[i], i as u32))
.collect();
order.sort_unstable_by(|a, b| {
a.0.cmp(&b.0)
.then_with(|| attrs[a.1 as usize].name.cmp(&attrs[b.1 as usize].name))
});
DenseAttrPlan {
serialized,
huge_id_of,
huge_count,
huge_total,
managed_plan,
name_plan,
huge_plan,
order,
managed_off,
btree_off,
name_nodes_off,
huge_bthd_off,
huge_nodes_off,
huge_data_off,
total_len,
}
}
const DENSE_ATTR_FRHP_SIZE: usize = {
let os = OFFSET_SIZE as usize;
let ls = LENGTH_SIZE as usize;
4 + 1
+ 2
+ 2
+ 1
+ 4
+ ls
+ os
+ ls
+ os
+ ls
+ ls
+ ls
+ ls
+ ls
+ ls
+ ls
+ ls
+ 2
+ ls
+ ls
+ 2
+ 2
+ os
+ 2
+ 4
};
impl DenseAttrPlan {
pub(crate) fn blob_len(&self) -> u64 {
self.total_len
}
pub(crate) fn attr_info_message(&self, heap_address: u64) -> Vec<u8> {
serialize_attribute_info(heap_address, heap_address + self.btree_off)
}
pub(crate) fn build(&self, heap_address: u64) -> DenseAttrBlob {
let max_heap_size: u16 = DENSE_ATTR_MAX_HEAP_SIZE_BITS;
let block_offset_bytes = DENSE_ATTR_BLOCK_OFFSET_BYTES; let heap_id_length: u16 = 8;
let Self {
serialized,
huge_id_of,
huge_count,
huge_total,
managed_plan,
name_plan,
huge_plan,
order,
..
} = self;
let huge_count = *huge_count;
let frhp_addr = heap_address;
let managed_addr = heap_address + self.managed_off;
let btree_addr = heap_address + self.btree_off;
let name_nodes_addr = heap_address + self.name_nodes_off;
let huge_bthd_addr = heap_address + self.huge_bthd_off;
let huge_nodes_addr = heap_address + self.huge_nodes_off;
let huge_data_addr = heap_address + self.huge_data_off;
debug_assert_eq!(
1 + block_offset_bytes + encoded_size_width(DENSE_ATTR_MAX_MANAGED_OBJECT as u64),
heap_id_length as usize,
"managed heap ID width must match what the declared maximum managed object size implies"
);
let mut frhp = Vec::with_capacity(DENSE_ATTR_FRHP_SIZE);
frhp.extend_from_slice(b"FRHP");
frhp.push(0); frhp.extend_from_slice(&heap_id_length.to_le_bytes());
frhp.extend_from_slice(&0u16.to_le_bytes()); frhp.push(0x02); #[expect(
clippy::cast_possible_truncation,
reason = "DENSE_ATTR_MAX_MANAGED_OBJECT is 65,514, well inside the 4-byte \
max-managed-object-size field"
)]
let max_managed = DENSE_ATTR_MAX_MANAGED_OBJECT as u32;
frhp.extend_from_slice(&max_managed.to_le_bytes());
write_length(&mut frhp, huge_count as u64, LENGTH_SIZE); if huge_count == 0 {
write_undef_offset(&mut frhp, OFFSET_SIZE); } else {
write_offset(&mut frhp, huge_bthd_addr, OFFSET_SIZE);
}
write_length(&mut frhp, managed_plan.free_space(), LENGTH_SIZE); write_undef_offset(&mut frhp, OFFSET_SIZE); write_length(&mut frhp, managed_plan.managed_space(), LENGTH_SIZE); write_length(&mut frhp, managed_plan.allocated_space(), LENGTH_SIZE); write_length(&mut frhp, managed_plan.allocation_iterator(), LENGTH_SIZE); let managed_count = (serialized.len() - huge_count) as u64;
write_length(&mut frhp, managed_count, LENGTH_SIZE); write_length(&mut frhp, *huge_total, LENGTH_SIZE); write_length(&mut frhp, huge_count as u64, LENGTH_SIZE); write_length(&mut frhp, 0, LENGTH_SIZE); write_length(&mut frhp, 0, LENGTH_SIZE); frhp.extend_from_slice(&fractal_heap_write::TABLE_WIDTH.to_le_bytes()); write_length(
&mut frhp,
fractal_heap_write::STARTING_BLOCK_SIZE,
LENGTH_SIZE,
);
write_length(
&mut frhp,
fractal_heap_write::MAX_DIRECT_BLOCK_SIZE,
LENGTH_SIZE,
); frhp.extend_from_slice(&max_heap_size.to_le_bytes());
frhp.extend_from_slice(&fractal_heap_write::START_ROOT_ROWS.to_le_bytes()); write_offset(
&mut frhp,
managed_plan.root_address(managed_addr),
OFFSET_SIZE,
);
frhp.extend_from_slice(&managed_plan.root_rows().to_le_bytes());
let frhp_checksum = crate::checksum::jenkins_lookup3(&frhp);
frhp.extend_from_slice(&frhp_checksum.to_le_bytes());
debug_assert_eq!(frhp.len(), DENSE_ATTR_FRHP_SIZE);
let mut heap_ids: Vec<Vec<u8>> = Vec::with_capacity(serialized.len());
let mut huge_records: Vec<(u64, u64, u64)> = Vec::with_capacity(huge_count);
let mut managed: Vec<&[u8]> = Vec::with_capacity(serialized.len() - huge_count);
let mut next_huge_addr = huge_data_addr;
for (s, huge_id) in serialized.iter().zip(huge_id_of) {
match huge_id {
Some(id) => {
huge_records.push((*id, next_huge_addr, s.len() as u64));
next_huge_addr += s.len() as u64;
heap_ids.push(encode_huge_id(*id, heap_id_length));
}
None => {
heap_ids.push(encode_managed_id(
managed_plan.heap_offset(managed.len()),
s.len() as u64,
max_heap_size,
heap_id_length,
));
managed.push(s.as_slice());
}
}
}
let managed_blocks = managed_plan.serialize(&managed, managed_addr, frhp_addr);
debug_assert_eq!(managed_blocks.len() as u64, managed_plan.region_size());
let record_size: u16 = heap_id_length + 1 + 4 + 4;
debug_assert_eq!(record_size, DENSE_ATTR_BTREE_RECORD);
let mut name_records = Vec::with_capacity(serialized.len() * record_size as usize);
for &(hash, i) in order {
name_records.extend_from_slice(&heap_ids[i as usize]);
name_records.push(0); name_records.extend_from_slice(&i.to_le_bytes()); name_records.extend_from_slice(&hash.to_le_bytes()); }
let bthd_addr = btree_addr;
let name_tree =
name_plan.serialize(&name_records, name_nodes_addr, OFFSET_SIZE, LENGTH_SIZE);
let mut blob = Vec::with_capacity(self.total_len.to_usize().unwrap_or(0));
blob.extend_from_slice(&frhp);
blob.extend_from_slice(&managed_blocks);
debug_assert_eq!(blob.len() as u64, bthd_addr - heap_address);
blob.extend_from_slice(&name_tree.header);
blob.extend_from_slice(&name_tree.nodes);
if let Some(huge_plan) = huge_plan {
let mut huge_bytes =
Vec::with_capacity(huge_records.len() * DENSE_ATTR_HUGE_BTREE_RECORD as usize);
for (id, addr, len) in &huge_records {
write_offset(&mut huge_bytes, *addr, OFFSET_SIZE);
write_length(&mut huge_bytes, *len, LENGTH_SIZE);
write_length(&mut huge_bytes, *id, LENGTH_SIZE);
}
let huge_tree =
huge_plan.serialize(&huge_bytes, huge_nodes_addr, OFFSET_SIZE, LENGTH_SIZE);
debug_assert_eq!(blob.len() as u64, huge_bthd_addr - heap_address);
blob.extend_from_slice(&huge_tree.header);
blob.extend_from_slice(&huge_tree.nodes);
debug_assert_eq!(blob.len() as u64, huge_data_addr - heap_address);
for (s, huge_id) in serialized.iter().zip(huge_id_of) {
if huge_id.is_some() {
blob.extend_from_slice(s);
}
}
}
debug_assert_eq!(
blob.len() as u64,
self.total_len,
"a dense attribute heap must fill the length its plan promised"
);
DenseAttrBlob {
attr_info_message: serialize_attribute_info(frhp_addr, bthd_addr),
blob,
}
}
}
const DUMMY_DENSE_BASE: u64 = 0;
pub(crate) fn build_dense_attrs(attrs: &[AttributeMessage], heap_address: u64) -> DenseAttrBlob {
dense_attrs_plan(attrs).build(heap_address)
}
fn encoded_size_width(value: u64) -> usize {
(64 - value.leading_zeros() as usize).div_ceil(8).max(1)
}
fn encode_huge_id(huge_id: u64, id_length: u16) -> Vec<u8> {
let payload_len = (id_length as usize) - 1;
debug_assert!(
payload_len >= 8 || huge_id < (1u64 << (payload_len * 8)),
"huge object ID overflows the heap ID payload"
);
let mut id = vec![0u8; id_length as usize];
id[0] = 0x10; for i in 0..payload_len.min(8) {
id[1 + i] = ((huge_id >> (i * 8)) & 0xFF) as u8;
}
id
}
fn encode_managed_id(offset: u64, length: u64, max_heap_size: u16, id_length: u16) -> Vec<u8> {
debug_assert!(length <= DENSE_ATTR_MAX_MANAGED_OBJECT as u64);
debug_assert_eq!(
offset >> max_heap_size,
0,
"heap offset overflows its field"
);
let mut id = vec![0u8; id_length as usize];
id[0] = 0x00; let combined = offset | (length << max_heap_size);
let payload_len = (id_length as usize) - 1;
for i in 0..payload_len.min(8) {
id[1 + i] = ((combined >> (i * 8)) & 0xFF) as u8;
}
id
}
pub(crate) fn compact_attribute_info_message() -> Vec<u8> {
serialize_attribute_info(u64::MAX, u64::MAX)
}
fn serialize_attribute_info(fh_addr: u64, btree_name_addr: u64) -> Vec<u8> {
let mut data = Vec::new();
data.push(0); data.push(0x00); data.extend_from_slice(&fh_addr.to_le_bytes());
data.extend_from_slice(&btree_name_addr.to_le_bytes());
data
}
pub(crate) fn write_offset(buf: &mut Vec<u8>, val: u64, offset_size: u8) {
#[expect(
clippy::cast_possible_truncation,
reason = "each arm narrows to offset_size, the on-disk address width chosen for this file"
)]
match offset_size {
2 => buf.extend_from_slice(&(val as u16).to_le_bytes()),
4 => buf.extend_from_slice(&(val as u32).to_le_bytes()),
8 => buf.extend_from_slice(&val.to_le_bytes()),
_ => {}
}
}
fn write_length(buf: &mut Vec<u8>, val: u64, length_size: u8) {
write_offset(buf, val, length_size);
}
pub(crate) fn write_undef_offset(buf: &mut Vec<u8>, offset_size: u8) {
for _ in 0..offset_size {
buf.push(0xFF);
}
}
pub struct FileWriter {
root_datasets: Vec<DatasetBuilder>,
root_attrs: Vec<(String, AttrSpec)>,
root_committed: Vec<CommittedDatatype>,
groups: Vec<FinishedGroup>,
userblock_size: u64,
userblock_content: Vec<u8>,
libver_bounds: Option<(LibVer, LibVer)>,
file_space_strategy: Option<(FileSpaceStrategy, bool, u64)>,
file_space_page_size: Option<u64>,
}
impl Default for FileWriter {
fn default() -> Self {
Self::new()
}
}
impl FileWriter {
pub fn new() -> Self {
Self {
root_datasets: Vec::new(),
root_attrs: Vec::new(),
root_committed: Vec::new(),
groups: Vec::new(),
userblock_size: 0,
userblock_content: Vec::new(),
libver_bounds: None,
file_space_strategy: None,
file_space_page_size: None,
}
}
pub fn with_libver_bounds(&mut self, low: LibVer, high: LibVer) -> &mut Self {
self.libver_bounds = Some((low, high));
self
}
pub(crate) fn apply_create_properties(&mut self, properties: &FileCreateProperties) {
self.userblock_size = properties.userblock();
self.libver_bounds = properties.libver_bounds();
self.file_space_strategy = properties.file_space_strategy();
self.file_space_page_size = properties.file_space_page_size();
}
fn resolve_libver(&self) -> Result<LibVer, FormatError> {
LibVer::resolve_writable(self.libver_bounds)
}
pub fn with_userblock(&mut self, size: u64) -> &mut Self {
self.userblock_size = size;
self
}
pub fn with_userblock_content(&mut self, content: &[u8]) -> &mut Self {
self.userblock_content = content.to_vec();
self
}
fn put_userblock<S: ByteSink>(
sink: &mut S,
ub: usize,
content: &[u8],
) -> Result<(), FormatError> {
if ub == 0 {
return Ok(());
}
sink.put(content)?;
sink.put_zeros(ub - content.len())
}
pub fn with_file_space_strategy(
&mut self,
strategy: FileSpaceStrategy,
persist: bool,
threshold: u64,
) -> &mut Self {
self.file_space_strategy = Some((strategy, persist, threshold));
self
}
pub fn with_file_space_page_size(&mut self, page_size: u64) -> &mut Self {
self.file_space_page_size = Some(page_size);
self
}
pub(crate) fn needs_latest_format(&self) -> bool {
fn any_chunked(datasets: &[DatasetBuilder]) -> bool {
datasets.iter().any(|d| {
d.chunk_options.is_chunked() || d.maxshape.is_some() || d.raw_chunks.is_some()
})
}
fn group_needs(group: &FinishedGroup) -> bool {
any_chunked(&group.datasets) || group.sub_groups.iter().any(group_needs)
}
self.file_space_info().is_some()
|| any_chunked(&self.root_datasets)
|| self.groups.iter().any(group_needs)
}
fn file_space_info(&self) -> Option<FileSpaceInfo> {
if self.file_space_strategy.is_none() && self.file_space_page_size.is_none() {
return None;
}
let (strategy, persist, threshold) = self.file_space_strategy.unwrap_or((
FileSpaceStrategy::FsmAggr,
false,
DEFAULT_THRESHOLD,
));
let page_size = self.file_space_page_size.unwrap_or(DEFAULT_PAGE_SIZE);
Some(if persist {
FileSpaceInfo::persistent_empty(strategy, threshold, page_size)
} else {
FileSpaceInfo::non_persistent(strategy, threshold, page_size)
})
}
fn file_space_extension_oh(&self) -> Result<Option<Vec<u8>>, FormatError> {
self.file_space_info()
.map(|info| {
let mut oh = ObjectHeaderWriter::new();
oh.add_message_with_flags(MessageType::FileSpaceInfo, info.serialize(), 0x14);
oh.serialize()
})
.transpose()
}
pub fn create_group(&mut self, name: &str) -> GroupBuilder {
GroupBuilder::new(name)
}
pub fn add_group(&mut self, group: FinishedGroup) {
self.groups.push(group);
}
pub fn create_dataset(&mut self, name: &str) -> &mut DatasetBuilder {
self.root_datasets.push(DatasetBuilder::new(name));
self.root_datasets.last_mut().unwrap()
}
pub fn commit_datatype(&mut self, name: &str, datatype: Datatype) {
self.root_committed.push(CommittedDatatype {
name: name.to_string(),
datatype,
});
}
pub fn set_root_attr_committed(&mut self, name: &str, value: AttrValue, path: &str) {
self.root_attrs.push((
name.to_string(),
crate::type_builders::committed_attr_spec(name, &value, path),
));
}
pub fn set_root_attr(&mut self, name: &str, value: AttrValue) {
self.root_attrs
.push((name.to_string(), AttrSpec::Value(value)));
}
pub(crate) fn set_root_attr_verbatim(&mut self, message: crate::attribute::AttributeMessage) {
self.root_attrs
.push((message.name.clone(), AttrSpec::Verbatim(message)));
}
pub(crate) fn set_root_attr_var_len_verbatim(
&mut self,
mut message: crate::attribute::AttributeMessage,
strings: Vec<String>,
) {
message.raw_data = crate::type_builders::vl_string_reference_bytes(&strings);
self.root_attrs.push((
message.name.clone(),
AttrSpec::VerbatimVarLen { message, strings },
));
}
pub fn finish(self) -> Result<Vec<u8>, FormatError> {
let mut buf = Vec::new();
self.finish_to_sink(&mut buf)?;
Ok(buf)
}
pub(crate) fn finish_to_sink<S: ByteSink>(self, sink: &mut S) -> Result<(), FormatError> {
let libver = self.resolve_libver()?;
if self.file_space_info().is_some() && libver < LibVer::V110 {
return Err(FormatError::LibverTooOldForContent {
content: "a file-space setting",
needs: LibVer::V110.name(),
writing: libver.name(),
});
}
if self.userblock_size != 0
&& (self.userblock_size < 512 || !self.userblock_size.is_power_of_two())
{
return Err(FormatError::InvalidUserblockSize(self.userblock_size));
}
if self.userblock_content.len() as u64 > self.userblock_size {
return Err(FormatError::UserblockContentTooLarge {
content: self.userblock_content.len() as u64,
userblock: self.userblock_size,
});
}
let (paged, persist_paged, page_size, fs_threshold) = match self.file_space_strategy {
Some((FileSpaceStrategy::Page, persist, threshold)) => {
let ps = self.file_space_page_size.unwrap_or(DEFAULT_PAGE_SIZE);
if ps < 512 || !ps.is_power_of_two() {
return Err(FormatError::InvalidFileSpacePageSize(ps));
}
if self.userblock_size % ps != 0 {
return Err(FormatError::UserblockNotPageAligned(
self.userblock_size,
ps,
));
}
(true, persist, ps, threshold)
}
_ => (false, false, 0, DEFAULT_THRESHOLD),
};
let ext_oh = self.file_space_extension_oh()?;
let nonpaged_persist: Option<(FileSpaceStrategy, u64, u64)> = match self.file_space_strategy
{
Some((strategy, true, threshold)) if strategy != FileSpaceStrategy::Page => Some((
strategy,
threshold,
self.file_space_page_size.unwrap_or(DEFAULT_PAGE_SIZE),
)),
_ => None,
};
struct DsFlat {
name: String,
dt: Datatype,
dt_location: DatatypeLocation,
ds: Dataspace,
raw: Vec<u8>,
attrs: Vec<AttributeMessage>,
chunk_options: ChunkOptions,
maxshape: Option<Vec<u64>>,
raw_chunks: Option<crate::type_builders::RawChunkPayload>,
reference_targets: Option<Vec<crate::type_builders::ObjectRefPatch>>,
vl_string_staging: Option<VlStringStaging>,
fill: Option<Vec<u8>>,
produced: Option<crate::type_builders::ProducedPayload>,
allocation: StorageAllocation,
declared_contiguous_len: u64,
}
impl DsFlat {
fn contiguous_len(&self) -> u64 {
match &self.produced {
Some(p) => p.total_bytes,
None => self.raw.len() as u64,
}
}
fn take_contiguous(&mut self) -> DsData {
match &self.produced {
Some(p) => DsData::Produced {
total_bytes: p.total_bytes,
block_bytes: p.block_bytes,
},
None => DsData::InMemory(core::mem::take(&mut self.raw)),
}
}
}
enum DsData {
InMemory(Vec<u8>),
Streamed(VerbatimPlan),
Produced {
total_bytes: u64,
block_bytes: u64,
},
}
impl DsData {
fn len(&self) -> u64 {
match self {
DsData::InMemory(v) => v.len() as u64,
DsData::Streamed(plan) => plan.total_len,
DsData::Produced { total_bytes, .. } => *total_bytes,
}
}
}
fn emit_ds_data<Sk: ByteSink>(
sink: &mut Sk,
data: &DsData,
raw_chunks: Option<&crate::type_builders::RawChunkPayload>,
produced: Option<&crate::type_builders::ProducedPayload>,
) -> Result<(), FormatError> {
match data {
DsData::InMemory(bytes) => sink.put(bytes),
DsData::Streamed(plan) => {
let provider = raw_chunks
.expect("a streamed data region implies a raw-chunk payload")
.provider
.0
.as_ref();
emit_chunked_data_verbatim(sink, plan, provider)
}
DsData::Produced {
total_bytes,
block_bytes,
} => {
let payload =
produced.expect("a produced data region implies a produced payload");
emit_produced_data(
sink,
payload.provider.0.as_ref(),
*total_bytes,
*block_bytes,
)
}
}
}
fn emit_produced_data<Sk: ByteSink>(
sink: &mut Sk,
provider: &dyn ChunkProvider,
total_bytes: u64,
block_bytes: u64,
) -> Result<(), FormatError> {
if block_bytes == 0 && total_bytes > 0 {
return Err(FormatError::SerializationError(
"a produced dataset declared a zero-length block".into(),
));
}
let mut block = Vec::new();
let mut written = 0u64;
let mut index = 0usize;
while written < total_bytes {
let expected = block_bytes.min(total_bytes - written);
block.clear();
provider.chunk_bytes(index, &mut block)?;
if block.len() as u64 != expected {
return Err(FormatError::SerializationError(
"a produced dataset's block does not match its planned size".into(),
));
}
sink.put(&block)?;
written += expected;
index += 1;
}
Ok(())
}
fn emit_collections<Sk: ByteSink>(
sink: &mut Sk,
collections: &[Vec<u8>],
) -> Result<(), FormatError> {
for collection in collections {
sink.put(collection)?;
}
Ok(())
}
struct ChunkedBuilt {
layout_message: Vec<u8>,
pipeline_message: Option<Vec<u8>>,
data: DsData,
}
fn measure_chunked(
d: &DsFlat,
data_address: u64,
chunk_set: Option<&CompressedChunkSet>,
) -> Result<ChunkedMeasure, FormatError> {
if let Some(rc) = &d.raw_chunks {
let VerbatimLayout {
plan,
layout_message,
pipeline_message,
} = plan_chunked_data_verbatim(
&rc.meta,
&d.ds.dimensions,
&rc.chunk_dims,
rc.element_size,
rc.pipeline_message.as_deref(),
data_address,
d.maxshape.as_deref(),
)?;
Ok(ChunkedMeasure {
data_len: plan.total_len,
layout_message,
pipeline_message,
})
} else {
let set = chunk_set
.expect("an encode-path chunked dataset must have a precomputed chunk set");
measure_chunked_at(set, data_address)
}
}
fn build_chunked(
d: &DsFlat,
data_address: u64,
chunk_set: Option<&CompressedChunkSet>,
) -> Result<ChunkedBuilt, FormatError> {
if let Some(rc) = &d.raw_chunks {
let VerbatimLayout {
plan,
layout_message,
pipeline_message,
} = plan_chunked_data_verbatim(
&rc.meta,
&d.ds.dimensions,
&rc.chunk_dims,
rc.element_size,
rc.pipeline_message.as_deref(),
data_address,
d.maxshape.as_deref(),
)?;
Ok(ChunkedBuilt {
layout_message,
pipeline_message,
data: DsData::Streamed(plan),
})
} else {
let set = chunk_set
.expect("an encode-path chunked dataset must have a precomputed chunk set");
let result = assemble_chunked_at(set, data_address)?;
Ok(ChunkedBuilt {
layout_message: result.layout_message,
pipeline_message: result.pipeline_message,
data: DsData::InMemory(result.data_bytes),
})
}
}
struct GrpFlat {
name: String,
attrs: Vec<AttributeMessage>,
ds_indices: Vec<usize>,
sub_group_indices: Vec<usize>,
committed_indices: Vec<usize>,
}
struct CtFlat {
name: String,
dt: Datatype,
references: u32,
}
let mut all_ds: Vec<DsFlat> = Vec::new();
let mut groups: Vec<GrpFlat> = Vec::new();
let mut committed: Vec<CtFlat> = Vec::new();
let mut root_ds_indices: Vec<usize> = Vec::new();
let mut root_group_indices: Vec<usize> = Vec::new();
fn flatten_dataset(
db: DatasetBuilder,
all_ds: &mut Vec<DsFlat>,
ds_vl: &mut Vec<Vec<VlPatch>>,
) -> Result<usize, FormatError> {
let dt = db.datatype.ok_or(FormatError::DatasetMissingData)?;
let shape = db.shape.ok_or(FormatError::DatasetMissingShape)?;
let raw_chunks = db.raw_chunks;
let produced = db.produced;
let allocation = db.allocation;
let unallocated = allocation == StorageAllocation::Unallocated;
debug_assert!(
!(unallocated && (db.data.is_some() || produced.is_some() || raw_chunks.is_some())),
"dataset {:?} declares unallocated storage and stages data for it",
db.name,
);
let is_empty = shape.contains(&0);
let raw = if is_empty || unallocated || raw_chunks.is_some() || produced.is_some() {
db.data.unwrap_or_default()
} else {
db.data.ok_or(FormatError::DatasetMissingData)?
};
let elem_size = dt.element_size_usize()?;
let elem_bytes = elem_size.get() as u64;
let extent = shape
.iter()
.copied()
.try_fold(1u64, |acc, d| acc.checked_mul(d))
.unwrap_or(u64::MAX)
.saturating_mul(elem_bytes);
let region_len = produced
.as_ref()
.map_or(raw.len() as u64, |p| p.total_bytes);
let declared_contiguous_len = if unallocated { extent } else { region_len };
if !unallocated && raw_chunks.is_none() && region_len != extent {
#[expect(
clippy::cast_possible_truncation,
reason = "byte counts reported in a shape-mismatch error; display-only"
)]
return Err(FormatError::ShapeDataMismatch {
expected: extent as usize,
actual: region_len as usize,
element_size: elem_size,
});
}
if db.chunk_options.is_chunked() || db.maxshape.is_some() {
db.chunk_options
.validate_geometry(&shape, db.maxshape.as_deref())
.map_err(FormatError::InvalidChunkGeometry)?;
}
let max_dimensions = db.maxshape.clone();
let dspace = Dataspace {
space_type: if shape.is_empty() {
DataspaceType::Scalar
} else {
DataspaceType::Simple
},
#[expect(
clippy::cast_possible_truncation,
reason = "dataspace rank fits the 1-byte dimensionality field (HDF5 caps \
rank at 32)"
)]
rank: shape.len() as u8,
dimensions: shape,
max_dimensions,
};
let patches = collect_vl_patches(&db.attrs);
let mut attrs = Vec::new();
for (n, v) in &db.attrs {
attrs.push(v.to_message(n));
}
#[cfg(feature = "provenance")]
if let Some(ref prov) = db.provenance {
let p = crate::provenance::Provenance {
creator: prov.creator.clone(),
timestamp: prov.timestamp.clone(),
source: prov.source.clone(),
};
attrs.extend(p.build_attrs(&raw));
}
if let Some(fill) = &db.fill {
if fill.len() != elem_size.get() {
return Err(FormatError::FillValueSizeMismatch {
expected: elem_size.get(),
actual: fill.len(),
});
}
}
let idx = all_ds.len();
all_ds.push(DsFlat {
name: db.name,
dt,
dt_location: db.datatype_location,
ds: dspace,
raw,
attrs,
chunk_options: db.chunk_options,
maxshape: db.maxshape,
raw_chunks,
produced,
reference_targets: db.reference_targets,
vl_string_staging: db.vl_string_staging,
fill: db.fill,
allocation,
declared_contiguous_len,
});
ds_vl.push(patches);
Ok(idx)
}
fn flatten_committed(
types: Vec<CommittedDatatype>,
committed: &mut Vec<CtFlat>,
) -> Result<Vec<usize>, FormatError> {
types
.into_iter()
.map(|ct| {
ct.datatype.element_size()?;
committed.push(CtFlat {
name: ct.name,
dt: ct.datatype,
references: 1,
});
Ok(committed.len() - 1)
})
.collect()
}
fn flatten_group(
g: FinishedGroup,
all_ds: &mut Vec<DsFlat>,
groups: &mut Vec<GrpFlat>,
committed: &mut Vec<CtFlat>,
grp_vl: &mut Vec<Vec<VlPatch>>,
ds_vl: &mut Vec<Vec<VlPatch>>,
) -> Result<usize, FormatError> {
let patches = collect_vl_patches(&g.attrs);
let mut gattrs = Vec::new();
for (n, v) in &g.attrs {
gattrs.push(v.to_message(n));
}
let committed_idx = flatten_committed(g.committed, committed)?;
let mut ds_idx = Vec::new();
for db in g.datasets {
ds_idx.push(flatten_dataset(db, all_ds, ds_vl)?);
}
let mut sub_grp_idx = Vec::new();
for sg in g.sub_groups {
sub_grp_idx.push(flatten_group(sg, all_ds, groups, committed, grp_vl, ds_vl)?);
}
let gi = groups.len();
groups.push(GrpFlat {
name: g.name,
attrs: gattrs,
ds_indices: ds_idx,
sub_group_indices: sub_grp_idx,
committed_indices: committed_idx,
});
grp_vl.push(patches);
Ok(gi)
}
let mut grp_vl: Vec<Vec<VlPatch>> = Vec::new();
let mut ds_vl: Vec<Vec<VlPatch>> = Vec::new();
let root_committed_indices = flatten_committed(self.root_committed, &mut committed)?;
for db in self.root_datasets {
root_ds_indices.push(flatten_dataset(db, &mut all_ds, &mut ds_vl)?);
}
for g in self.groups.into_iter() {
root_group_indices.push(flatten_group(
g,
&mut all_ds,
&mut groups,
&mut committed,
&mut grp_vl,
&mut ds_vl,
)?);
}
struct VlPatch {
collections: Vec<Vec<u8>>,
attr_index: usize, }
fn place_collections(collections: &[Vec<u8>], cursor: &mut u64) -> Vec<u64> {
collections
.iter()
.map(|c| {
let addr = *cursor;
*cursor += c.len() as u64;
addr
})
.collect()
}
fn collect_vl_patches(attrs_raw: &[(String, AttrSpec)]) -> Vec<VlPatch> {
let mut patches = Vec::new();
for (i, (_n, v)) in attrs_raw.iter().enumerate() {
if let Some(strings) = v.var_len_strings() {
patches.push(VlPatch {
collections: build_global_heap_collections(strings),
attr_index: i,
});
}
}
patches
}
let vl_root = collect_vl_patches(&self.root_attrs);
let mut root_attrs: Vec<AttributeMessage> = Vec::new();
for (n, v) in &self.root_attrs {
root_attrs.push(v.to_message(n));
}
let committed_by_path: HashMap<String, usize> = {
fn walk(
prefix: &str,
gi: usize,
groups: &[GrpFlat],
committed: &[CtFlat],
out: &mut HashMap<String, usize>,
) {
for &ci in &groups[gi].committed_indices {
out.insert(format!("{prefix}/{}", committed[ci].name), ci);
}
for &sgi in &groups[gi].sub_group_indices {
walk(
&format!("{prefix}/{}", groups[sgi].name),
sgi,
groups,
committed,
out,
);
}
}
let mut out: HashMap<String, usize> = root_committed_indices
.iter()
.map(|&ci| (committed[ci].name.clone(), ci))
.collect();
for &gi in &root_group_indices {
walk(&groups[gi].name, gi, &groups, &committed, &mut out);
}
out
};
fn register_committed_use(
committed: &mut [CtFlat],
by_path: &HashMap<String, usize>,
location: &DatatypeLocation,
dt: &Datatype,
user: impl FnOnce() -> String,
) -> Result<(), FormatError> {
let Some(path) = location.unresolved_path() else {
return Ok(());
};
let Some(&ci) = by_path.get(path) else {
return Err(FormatError::UnknownCommittedDatatype(path.to_string()));
};
if committed[ci].dt.serialize() != dt.serialize() {
return Err(FormatError::CommittedDatatypeMismatch {
path: path.to_string(),
user: user(),
});
}
committed[ci].references = committed[ci].references.saturating_add(1);
Ok(())
}
for attr in &root_attrs {
register_committed_use(
&mut committed,
&committed_by_path,
&attr.datatype_location,
&attr.datatype,
|| format!("root attribute {:?}", attr.name),
)?;
}
for g in &groups {
for attr in &g.attrs {
register_committed_use(
&mut committed,
&committed_by_path,
&attr.datatype_location,
&attr.datatype,
|| format!("attribute {:?} of group {:?}", attr.name, g.name),
)?;
}
}
for d in &all_ds {
register_committed_use(
&mut committed,
&committed_by_path,
&d.dt_location,
&d.dt,
|| format!("dataset {:?}", d.name),
)?;
for attr in &d.attrs {
register_committed_use(
&mut committed,
&committed_by_path,
&attr.datatype_location,
&attr.datatype,
|| format!("attribute {:?} of dataset {:?}", attr.name, d.name),
)?;
}
}
let committed_oh: Vec<Vec<u8>> = committed
.iter()
.map(|ct| build_committed_datatype_oh(&ct.dt, ct.references))
.collect::<Result<_, _>>()?;
struct LinkTables<'a> {
all_ds: &'a [DsFlat],
committed: &'a [CtFlat],
groups: &'a [GrpFlat],
ds_addrs: &'a [u64],
committed_addrs: &'a [u64],
group_addrs: &'a [u64],
}
impl LinkTables<'_> {
fn links(
&self,
ds_indices: &[usize],
committed_indices: &[usize],
sub_group_indices: &[usize],
) -> Vec<LinkMessage> {
let mut links = Vec::with_capacity(
ds_indices.len() + committed_indices.len() + sub_group_indices.len(),
);
for &i in ds_indices {
links.push(make_link(&self.all_ds[i].name, self.ds_addrs[i]));
}
for &ci in committed_indices {
links.push(make_link(
&self.committed[ci].name,
self.committed_addrs[ci],
));
}
for &gi in sub_group_indices {
links.push(make_link(&self.groups[gi].name, self.group_addrs[gi]));
}
links
}
}
let dummy_ds_addrs = vec![0u64; all_ds.len()];
let dummy_ct_addrs = vec![0u64; committed.len()];
let dummy_grp_addrs = vec![0u64; groups.len()];
let root_dense = needs_dense_attrs(&root_attrs);
let group_dense: Vec<bool> = groups.iter().map(|g| needs_dense_attrs(&g.attrs)).collect();
let ds_dense: Vec<bool> = all_ds.iter().map(|d| needs_dense_attrs(&d.attrs)).collect();
fn check_compact_attrs(attrs: &[AttributeMessage]) -> Result<(), FormatError> {
for a in attrs {
let size = a.serialize(LENGTH_SIZE).len();
if size > OBJECT_HEADER_MESSAGE_MAX {
return Err(FormatError::AttributeMessageTooLarge {
name: a.name.clone(),
size,
});
}
}
Ok(())
}
if root_dense {
dense_attrs_check(&root_attrs)?;
} else {
check_compact_attrs(&root_attrs)?;
}
for (gi, g) in groups.iter().enumerate() {
if group_dense[gi] {
dense_attrs_check(&g.attrs)?;
} else {
check_compact_attrs(&g.attrs)?;
}
}
for (i, d) in all_ds.iter().enumerate() {
if ds_dense[i] {
dense_attrs_check(&d.attrs)?;
} else {
check_compact_attrs(&d.attrs)?;
}
}
let is_chunked: Vec<bool> = all_ds
.iter()
.map(|d| d.chunk_options.is_chunked() || d.maxshape.is_some() || d.raw_chunks.is_some())
.collect();
if libver < LibVer::V110 && is_chunked.iter().any(|&c| c) {
return Err(FormatError::LibverTooOldForContent {
content: "a chunked, filtered, or resizable dataset",
needs: LibVer::V110.name(),
writing: libver.name(),
});
}
let mut early_gcol: Vec<usize> = Vec::new();
let early_gcol_size = {
let mut cursor = SUPERBLOCK_SIZE as u64;
for i in 0..all_ds.len() {
if !is_chunked[i] || all_ds[i].raw_chunks.is_some() {
continue;
}
let Some(staging) = all_ds[i].vl_string_staging.take() else {
continue;
};
let addrs = place_collections(&staging.collections, &mut cursor);
patch_vl_refs_masked(&mut all_ds[i].raw, &staging.patch_offsets, &addrs);
all_ds[i].vl_string_staging = Some(staging);
early_gcol.push(i);
}
(cursor - SUPERBLOCK_SIZE as u64).to_usize()?
};
let chunk_sets: Vec<Option<CompressedChunkSet>> = all_ds
.iter()
.enumerate()
.map(|(i, d)| {
if is_chunked[i] && d.raw_chunks.is_none() {
let chunk_dims = d.chunk_options.resolve_chunk_dims(&d.ds.dimensions);
let ctx = crate::filters::ChunkContext::from_datatype(&chunk_dims, &d.dt)?;
let elem = crate::convert::nonzero_usize_from(ctx.element_size)?;
let fill = crate::fill_value::FillPattern::new(d.fill.as_deref(), elem);
Ok(Some(compress_chunks(
&d.raw,
&d.ds.dimensions,
ctx,
&d.chunk_options,
d.maxshape.as_deref(),
fill,
d.allocation,
)?))
} else {
Ok(None)
}
})
.collect::<Result<_, FormatError>>()?;
struct DenseSpans {
root: Option<(u64, usize)>,
groups: Vec<Option<(u64, usize)>>,
datasets: Vec<Option<(u64, usize)>>,
}
struct DenseBlobs {
root: Option<DenseAttrBlob>,
groups: Vec<Option<DenseAttrBlob>>,
datasets: Vec<Option<DenseAttrBlob>>,
}
impl DenseSpans {
fn build(
&self,
root_attrs: &[AttributeMessage],
groups: &[GrpFlat],
datasets: &[DsFlat],
) -> Result<DenseBlobs, FormatError> {
fn one(
attrs: &[AttributeMessage],
span: Option<(u64, usize)>,
) -> Result<Option<DenseAttrBlob>, FormatError> {
let Some((address, reserved)) = span else {
return Ok(None);
};
let blob = build_dense_attrs(attrs, address);
if blob.blob.len() != reserved {
return Err(FormatError::SerializationError(format!(
"a dense attribute heap built {} bytes into a span of {reserved} \
reserved for it; every address after it would be wrong",
blob.blob.len(),
)));
}
Ok(Some(blob))
}
Ok(DenseBlobs {
root: one(root_attrs, self.root)?,
groups: groups
.iter()
.zip(&self.groups)
.map(|(g, &span)| one(&g.attrs, span))
.collect::<Result<_, _>>()?,
datasets: datasets
.iter()
.zip(&self.datasets)
.map(|(d, &span)| one(&d.attrs, span))
.collect::<Result<_, _>>()?,
})
}
}
let dummy_tables = LinkTables {
all_ds: &all_ds,
committed: &committed,
groups: &groups,
ds_addrs: &dummy_ds_addrs,
committed_addrs: &dummy_ct_addrs,
group_addrs: &dummy_grp_addrs,
};
let mut group_oh_sizes: Vec<usize> = Vec::with_capacity(groups.len());
let mut group_dense_lens: Vec<Option<usize>> = Vec::with_capacity(groups.len());
for (gi, g) in groups.iter().enumerate() {
let dummy_links =
dummy_tables.links(&g.ds_indices, &g.committed_indices, &g.sub_group_indices);
let (oh, dense_len) = if group_dense[gi] {
let plan = dense_attrs_plan(&g.attrs);
(
build_group_oh(
&dummy_links,
&g.attrs,
Some(&plan.attr_info_message(DUMMY_DENSE_BASE)),
)?,
Some(plan.blob_len().to_usize()?),
)
} else {
(build_group_oh(&dummy_links, &g.attrs, None)?, None)
};
group_oh_sizes.push(oh.len());
group_dense_lens.push(dense_len);
}
let root_dummy_links = dummy_tables.links(
&root_ds_indices,
&root_committed_indices,
&root_group_indices,
);
let (root_oh_size, root_dense_len) = if root_dense {
let plan = dense_attrs_plan(&root_attrs);
(
build_group_oh(
&root_dummy_links,
&root_attrs,
Some(&plan.attr_info_message(DUMMY_DENSE_BASE)),
)?
.len(),
Some(plan.blob_len().to_usize()?),
)
} else {
(
build_group_oh(&root_dummy_links, &root_attrs, None)?.len(),
None,
)
};
let mut actual_ds_oh_sizes: Vec<usize> = Vec::with_capacity(all_ds.len());
let mut ds_data_lens: Vec<u64> = Vec::with_capacity(all_ds.len());
let mut ds_dense_lens: Vec<Option<usize>> = Vec::with_capacity(all_ds.len());
let mut dummy_cursor = 0u64;
for (i, d) in all_ds.iter().enumerate() {
let dense_plan = if ds_dense[i] {
Some(dense_attrs_plan(&d.attrs))
} else {
None
};
let dense_attr_info = dense_plan
.as_ref()
.map(|plan| plan.attr_info_message(DUMMY_DENSE_BASE));
ds_dense_lens.push(
dense_plan
.as_ref()
.map(|plan| plan.blob_len().to_usize())
.transpose()?,
);
let oh = if is_chunked[i] {
let measured = measure_chunked(d, dummy_cursor, chunk_sets[i].as_ref())?;
dummy_cursor += measured.data_len;
ds_data_lens.push(measured.data_len);
build_chunked_dataset_oh(
&d.dt,
&d.dt_location,
&d.ds,
&measured.layout_message,
measured.pipeline_message.as_deref(),
&d.attrs,
dense_attr_info.as_deref(),
d.fill.as_deref(),
)?
} else {
ds_data_lens.push(d.contiguous_len());
build_dataset_oh(
&d.dt,
&d.dt_location,
&d.ds,
0,
d.declared_contiguous_len,
&d.attrs,
dense_attr_info.as_deref(),
d.fill.as_deref(),
libver,
)?
};
actual_ds_oh_sizes.push(oh.len());
}
#[expect(
clippy::cast_possible_truncation,
reason = "userblock_size is a small power-of-two header size used as an in-memory \
buffer offset; it fits usize on every supported target"
)]
let ub = self.userblock_size as usize;
let root_group_addr = (SUPERBLOCK_SIZE + early_gcol_size) as u64;
let mut cursor2 = SUPERBLOCK_SIZE + early_gcol_size + root_oh_size;
let reserve = |cursor2: &mut usize, len: Option<usize>| {
len.map(|len| {
let addr = *cursor2 as u64;
*cursor2 += len;
(addr, len)
})
};
let root_dense_span = reserve(&mut cursor2, root_dense_len);
let mut group_dense_spans: Vec<Option<(u64, usize)>> = Vec::with_capacity(groups.len());
let group_addrs2: Vec<u64> = group_oh_sizes
.iter()
.enumerate()
.map(|(gi, &sz)| {
let addr = cursor2 as u64;
cursor2 += sz;
group_dense_spans.push(reserve(&mut cursor2, group_dense_lens[gi]));
addr
})
.collect();
let committed_addrs: Vec<u64> = committed_oh
.iter()
.map(|oh| {
let addr = cursor2 as u64;
cursor2 += oh.len();
addr
})
.collect();
let mut ds_dense_spans: Vec<Option<(u64, usize)>> = Vec::with_capacity(all_ds.len());
let ds_oh_addrs2: Vec<u64> = actual_ds_oh_sizes
.iter()
.enumerate()
.map(|(i, &sz)| {
let addr = cursor2 as u64;
cursor2 += sz;
ds_dense_spans.push(reserve(&mut cursor2, ds_dense_lens[i]));
addr
})
.collect();
let dense_spans = DenseSpans {
root: root_dense_span,
groups: group_dense_spans,
datasets: ds_dense_spans,
};
{
let mut path_map = HashMap::<String, u64>::new();
path_map.insert(String::new(), root_group_addr);
for &i in &root_ds_indices {
path_map.insert(all_ds[i].name.clone(), ds_oh_addrs2[i]);
}
for (path, &ci) in &committed_by_path {
path_map.insert(path.clone(), committed_addrs[ci]);
}
for &gi in &root_group_indices {
fn register_group(
prefix: &str,
gi: usize,
groups: &[GrpFlat],
ds_addrs: &[u64],
grp_addrs: &[u64],
all_ds: &[DsFlat],
map: &mut HashMap<String, u64>,
) {
map.insert(prefix.to_string(), grp_addrs[gi]);
for &di in &groups[gi].ds_indices {
map.insert(format!("{}/{}", prefix, all_ds[di].name), ds_addrs[di]);
}
for &sgi in &groups[gi].sub_group_indices {
register_group(
&format!("{}/{}", prefix, groups[sgi].name),
sgi,
groups,
ds_addrs,
grp_addrs,
all_ds,
map,
);
}
}
register_group(
&groups[gi].name,
gi,
&groups,
&ds_oh_addrs2,
&group_addrs2,
&all_ds,
&mut path_map,
);
}
for d in all_ds.iter_mut() {
let Some(ref patches) = d.reference_targets else {
continue;
};
for patch in patches {
let addr = match &patch.target {
crate::type_builders::ObjectRefTarget::Path(path) => {
path_map.get(path).copied().unwrap_or(u64::MAX)
}
crate::type_builders::ObjectRefTarget::Raw(addr) => *addr,
};
write_reference_address(&mut d.raw, patch.byte_offset, addr);
}
}
}
{
let resolve = |location: &mut DatatypeLocation| {
let ci = match location.unresolved_path() {
Some(path) => *committed_by_path.get(path).expect(
"every committed-datatype path was resolved against this same map \
before any header was sized",
),
None => return,
};
*location = DatatypeLocation::Committed(committed_addrs[ci]);
};
for attr in &mut root_attrs {
resolve(&mut attr.datatype_location);
}
for g in &mut groups {
for attr in &mut g.attrs {
resolve(&mut attr.datatype_location);
}
}
for d in &mut all_ds {
resolve(&mut d.dt_location);
for attr in &mut d.attrs {
resolve(&mut attr.datatype_location);
}
}
}
struct DsLayout {
data: DsData,
data_addr: u64,
chunked_msgs: Option<(Vec<u8>, Option<Vec<u8>>)>,
}
fn build_ds_ohs(
all_ds: &[DsFlat],
ds_layouts: &[DsLayout],
ds_dense_blobs: &[Option<DenseAttrBlob>],
libver: LibVer,
) -> Result<Vec<Vec<u8>>, FormatError> {
let mut oh_bytes: Vec<Vec<u8>> = Vec::with_capacity(all_ds.len());
for (i, d) in all_ds.iter().enumerate() {
let layout = &ds_layouts[i];
let oh = if let Some((ref lm, ref pm)) = layout.chunked_msgs {
build_chunked_dataset_oh(
&d.dt,
&d.dt_location,
&d.ds,
lm,
pm.as_deref(),
&d.attrs,
ds_dense_blobs[i]
.as_ref()
.map(|b| b.attr_info_message.as_slice()),
d.fill.as_deref(),
)?
} else {
build_dataset_oh(
&d.dt,
&d.dt_location,
&d.ds,
layout.data_addr,
d.declared_contiguous_len,
&d.attrs,
ds_dense_blobs[i]
.as_ref()
.map(|b| b.attr_info_message.as_slice()),
d.fill.as_deref(),
libver,
)?
};
oh_bytes.push(oh);
}
Ok(oh_bytes)
}
if paged {
let os = OFFSET_SIZE;
let base = BaseAddress::new(ub as u64);
let mut meta = cursor2 as u64;
let gcol_start = meta;
let mut gcol_cursor = meta;
let mut elem_gcol: Vec<(usize, Vec<u64>)> = Vec::new();
{
for patch in &vl_root {
let addrs = place_collections(&patch.collections, &mut gcol_cursor);
patch_vl_refs(&mut root_attrs[patch.attr_index].raw_data, &addrs);
}
for (gi, patches) in grp_vl.iter().enumerate() {
for patch in patches {
let addrs = place_collections(&patch.collections, &mut gcol_cursor);
patch_vl_refs(&mut groups[gi].attrs[patch.attr_index].raw_data, &addrs);
}
}
for (di, patches) in ds_vl.iter().enumerate() {
for patch in patches {
let addrs = place_collections(&patch.collections, &mut gcol_cursor);
patch_vl_refs(&mut all_ds[di].attrs[patch.attr_index].raw_data, &addrs);
}
}
for (i, d) in all_ds.iter().enumerate() {
if early_gcol.contains(&i) {
continue;
}
if let Some(staging) = &d.vl_string_staging {
elem_gcol
.push((i, place_collections(&staging.collections, &mut gcol_cursor)));
}
}
}
let gcol_total_size = gcol_cursor - gcol_start;
meta += gcol_total_size;
let DenseBlobs {
root: root_dense_blob,
groups: group_dense_blobs,
datasets: ds_dense_blobs,
} = dense_spans.build(&root_attrs, &groups, &all_ds)?;
let ext_addr = meta;
let ext_len = ext_oh.as_ref().map_or(0, |b| b.len()) as u64;
meta += ext_len;
let meta_content_end = meta;
let mut empty_indices: Vec<usize> = Vec::new();
let mut small_indices: Vec<usize> = Vec::new();
let mut large_indices: Vec<usize> = Vec::new();
for i in 0..all_ds.len() {
let len = ds_data_lens[i];
if !is_chunked[i] && len == 0 {
empty_indices.push(i);
} else if len < page_size {
small_indices.push(i);
} else {
large_indices.push(i);
}
}
let small_raw_total: u64 = small_indices.iter().map(|&i| ds_data_lens[i]).sum();
let large_frag_sizes: Vec<u64> = large_indices
.iter()
.map(|&i| align_up(ds_data_lens[i], page_size) - ds_data_lens[i])
.filter(|&f| f > 0)
.collect();
let draw_active =
small_raw_total > 0 && align_up(small_raw_total, page_size) != small_raw_total;
let large_active = !large_frag_sizes.is_empty();
let mut slots = [u64::MAX; NUM_FILE_FSM_MANAGERS];
let mut super_fsm: Option<(u64, u64)> = None;
let mut draw_fsm: Option<(u64, u64)> = None;
let mut large_fsm: Option<(u64, u64)> = None;
let super_block_len = fshd_len(os) + fsse_len(&[0], os);
let draw_block_len = if draw_active {
fshd_len(os) + fsse_len(&[0], os)
} else {
0
};
let large_block_len = if large_active {
fshd_len(os) + fsse_len(&large_frag_sizes, os)
} else {
0
};
let super_active = if persist_paged {
let with = meta_content_end + super_block_len + draw_block_len + large_block_len;
align_up(with, page_size) > with
} else {
false
};
if persist_paged {
if super_active {
let fshd_addr = meta;
meta += fshd_len(os);
let fsse_addr = meta;
meta += fsse_len(&[0], os);
slots[0] = fshd_addr;
super_fsm = Some((fshd_addr, fsse_addr));
}
if draw_active {
let fshd_addr = meta;
meta += fshd_len(os);
let fsse_addr = meta;
meta += fsse_len(&[0], os);
slots[2] = fshd_addr;
draw_fsm = Some((fshd_addr, fsse_addr));
}
if large_active {
let fshd_addr = meta;
meta += fshd_len(os);
let fsse_addr = meta;
meta += fsse_len(&large_frag_sizes, os);
slots[6] = fshd_addr;
large_fsm = Some((fshd_addr, fsse_addr));
}
}
let meta_end = meta;
let raw_start = align_up(meta_end, page_size);
let super_section = super_fsm.map(|_| FreeSection {
addr: meta_end,
size: raw_start - meta_end,
});
let mut layouts: Vec<Option<DsLayout>> = (0..all_ds.len()).map(|_| None).collect();
for &i in &empty_indices {
let raw = core::mem::take(&mut all_ds[i].raw);
layouts[i] = Some(DsLayout {
data: DsData::InMemory(raw),
data_addr: u64::MAX,
chunked_msgs: None,
});
}
let mut c = raw_start;
for &i in &small_indices {
let base_addr = c;
let layout = if is_chunked[i] {
let built = build_chunked(&all_ds[i], base_addr, chunk_sets[i].as_ref())?;
debug_assert_eq!(built.data.len(), ds_data_lens[i]);
c += built.data.len();
DsLayout {
data: built.data,
data_addr: base_addr,
chunked_msgs: Some((built.layout_message, built.pipeline_message)),
}
} else {
let data = all_ds[i].take_contiguous();
c += data.len();
DsLayout {
data,
data_addr: base_addr,
chunked_msgs: None,
}
};
layouts[i] = Some(layout);
}
let small_raw_end = c;
let draw_section = if draw_active {
let padded = align_up(small_raw_end, page_size);
c = padded;
Some(FreeSection {
addr: small_raw_end,
size: padded - small_raw_end,
})
} else {
if small_raw_total > 0 {
c = align_up(small_raw_end, page_size);
}
None
};
let mut large_sections: Vec<FreeSection> = Vec::new();
for &i in &large_indices {
c = align_up(c, page_size);
let data_addr = c;
let built_len;
let layout = if is_chunked[i] {
let built = build_chunked(&all_ds[i], data_addr, chunk_sets[i].as_ref())?;
debug_assert_eq!(built.data.len(), ds_data_lens[i]);
built_len = built.data.len();
DsLayout {
data: built.data,
data_addr,
chunked_msgs: Some((built.layout_message, built.pipeline_message)),
}
} else {
let data = all_ds[i].take_contiguous();
built_len = data.len();
DsLayout {
data,
data_addr,
chunked_msgs: None,
}
};
layouts[i] = Some(layout);
let data_end = data_addr + built_len;
let frag = align_up(data_end, page_size) - data_end;
if frag > 0 {
large_sections.push(FreeSection {
addr: data_end,
size: frag,
});
}
c = align_up(data_end, page_size);
}
let eoa_rel = c; let eof_addr2 = base.absolute(eoa_rel)?;
let eoa_pre_fsm = eoa_rel;
for (i, gaddrs) in &elem_gcol {
let staging = all_ds[*i]
.vl_string_staging
.as_ref()
.expect("elem_gcol only holds datasets with VL staging");
let Some(DsLayout {
data: DsData::InMemory(bytes),
..
}) = layouts[*i].as_mut()
else {
unreachable!(
"a staged VL-string dataset holds its element bytes in memory: the \
chunked path patches before encoding, and a produced region refuses \
to carry VL staging at all"
)
};
patch_vl_refs_masked(bytes, &staging.patch_offsets, gaddrs);
}
let ds_layouts: Vec<DsLayout> = layouts
.into_iter()
.map(|o| o.expect("every dataset placed"))
.collect();
let ds_oh_bytes = build_ds_ohs(&all_ds, &ds_layouts, &ds_dense_blobs, libver)?;
debug_assert_eq!(
ds_oh_bytes.iter().map(|b| b.len()).collect::<Vec<_>>(),
actual_ds_oh_sizes
);
let real_ext_oh = if persist_paged {
let info = FileSpaceInfo::persistent_managers(
FileSpaceStrategy::Page,
fs_threshold,
page_size,
slots,
eoa_pre_fsm,
);
let mut oh = ObjectHeaderWriter::new();
oh.add_message_with_flags(MessageType::FileSpaceInfo, info.serialize(), 0x14);
oh.serialize()?
} else {
ext_oh
.clone()
.expect("a paged file always emits a File Space Info message")
};
debug_assert_eq!(real_ext_oh.len() as u64, ext_len);
let super_blocks = super_fsm.map(|(fshd_addr, fsse_addr)| {
serialize_file_fsm(
&[super_section.expect("SUPER active implies a section")],
fshd_addr,
fsse_addr,
os,
SECT_CLASS_SMALL,
)
});
let draw_blocks = draw_fsm.map(|(fshd_addr, fsse_addr)| {
serialize_file_fsm(
&[draw_section.expect("DRAW active implies a section")],
fshd_addr,
fsse_addr,
os,
SECT_CLASS_SMALL,
)
});
let large_blocks = large_fsm.map(|(fshd_addr, fsse_addr)| {
serialize_file_fsm(&large_sections, fshd_addr, fsse_addr, os, SECT_CLASS_LARGE)
});
sink.reserve(eof_addr2.to_usize()?);
Self::put_userblock(sink, ub, &self.userblock_content)?;
let sb = Superblock {
version: superblock_version(libver),
offset_size: OFFSET_SIZE,
length_size: LENGTH_SIZE,
base_address: base,
eof_address: eof_addr2,
root_group_address: root_group_addr,
group_leaf_node_k: None,
group_internal_node_k: None,
indexed_storage_internal_node_k: None,
free_space_address: None,
driver_info_address: None,
consistency_flags: 0,
superblock_extension_address: Some(ext_addr),
checksum: None,
};
sink.put(&sb.serialize())?;
for &i in &early_gcol {
let staging = all_ds[i]
.vl_string_staging
.as_ref()
.expect("early_gcol only holds datasets with VL staging");
emit_collections(sink, &staging.collections)?;
}
let tables = LinkTables {
all_ds: &all_ds,
committed: &committed,
groups: &groups,
ds_addrs: &ds_oh_addrs2,
committed_addrs: &committed_addrs,
group_addrs: &group_addrs2,
};
let root_links = tables.links(
&root_ds_indices,
&root_committed_indices,
&root_group_indices,
);
let root_oh = build_group_oh(
&root_links,
&root_attrs,
root_dense_blob
.as_ref()
.map(|b| b.attr_info_message.as_slice()),
)?;
debug_assert_eq!(root_oh.len(), root_oh_size);
sink.put(&root_oh)?;
if let Some(ref blob) = root_dense_blob {
sink.put(&blob.blob)?;
}
for (gi, g) in groups.iter().enumerate() {
let links = tables.links(&g.ds_indices, &g.committed_indices, &g.sub_group_indices);
let oh = build_group_oh(
&links,
&g.attrs,
group_dense_blobs[gi]
.as_ref()
.map(|b| b.attr_info_message.as_slice()),
)?;
debug_assert_eq!(oh.len(), group_oh_sizes[gi]);
sink.put(&oh)?;
if let Some(ref blob) = group_dense_blobs[gi] {
sink.put(&blob.blob)?;
}
}
for oh in &committed_oh {
sink.put(oh)?;
}
for (i, oh) in ds_oh_bytes.iter().enumerate() {
sink.put(oh)?;
if let Some(ref dense) = ds_dense_blobs[i] {
sink.put(&dense.blob)?;
}
}
for patch in &vl_root {
emit_collections(sink, &patch.collections)?;
}
for patches in &grp_vl {
for patch in patches {
emit_collections(sink, &patch.collections)?;
}
}
for patches in &ds_vl {
for patch in patches {
emit_collections(sink, &patch.collections)?;
}
}
for (i, d) in all_ds.iter().enumerate() {
if early_gcol.contains(&i) {
continue; }
if let Some(staging) = &d.vl_string_staging {
emit_collections(sink, &staging.collections)?;
}
}
debug_assert_eq!(sink.position(), base.get() + ext_addr);
sink.put(&real_ext_oh)?;
for blocks in [&super_blocks, &draw_blocks, &large_blocks]
.into_iter()
.flatten()
{
sink.put(&blocks.0)?;
sink.put(&blocks.1)?;
}
debug_assert_eq!(sink.position(), base.get() + meta_end);
sink.put_zeros((raw_start - meta_end).to_usize()?)?;
for &i in &small_indices {
debug_assert_eq!(sink.position(), base.get() + ds_layouts[i].data_addr);
emit_ds_data(
sink,
&ds_layouts[i].data,
all_ds[i].raw_chunks.as_ref(),
all_ds[i].produced.as_ref(),
)?;
}
if small_raw_total > 0 {
sink.put_zeros((align_up(small_raw_end, page_size) - small_raw_end).to_usize()?)?;
}
for &i in &large_indices {
let data_addr = ds_layouts[i].data_addr;
let gap = base.absolute(data_addr)? - sink.position();
sink.put_zeros(gap.to_usize()?)?;
emit_ds_data(
sink,
&ds_layouts[i].data,
all_ds[i].raw_chunks.as_ref(),
all_ds[i].produced.as_ref(),
)?;
let end_rel = base.relative(sink.position())?;
sink.put_zeros((align_up(end_rel, page_size) - end_rel).to_usize()?)?;
}
let final_pad = eof_addr2 - sink.position();
sink.put_zeros(final_pad.to_usize()?)?;
debug_assert_eq!(sink.position(), eof_addr2);
return Ok(());
}
let mut ds_layouts: Vec<DsLayout> = Vec::new();
for (i, d) in all_ds.iter_mut().enumerate() {
if is_chunked[i] {
let data_address = cursor2 as u64;
let built = build_chunked(d, data_address, chunk_sets[i].as_ref())?;
cursor2 += built.data.len().to_usize()?;
ds_layouts.push(DsLayout {
data: built.data,
data_addr: data_address,
chunked_msgs: Some((built.layout_message, built.pipeline_message)),
});
} else {
let data = d.take_contiguous();
let addr = if data.len() == 0 {
u64::MAX
} else {
let a = cursor2 as u64;
cursor2 += data.len().to_usize()?;
a
};
ds_layouts.push(DsLayout {
data,
data_addr: addr,
chunked_msgs: None,
});
}
}
let has_vl = !vl_root.is_empty()
|| grp_vl.iter().any(|v| !v.is_empty())
|| ds_vl.iter().any(|v| !v.is_empty())
|| all_ds.iter().any(|d| d.vl_string_staging.is_some());
let mut gcol_total_size = 0usize;
if has_vl {
let mut gcol_cursor = cursor2 as u64;
for patch in &vl_root {
let addrs = place_collections(&patch.collections, &mut gcol_cursor);
patch_vl_refs(&mut root_attrs[patch.attr_index].raw_data, &addrs);
}
for (gi, patches) in grp_vl.iter().enumerate() {
for patch in patches {
let addrs = place_collections(&patch.collections, &mut gcol_cursor);
patch_vl_refs(&mut groups[gi].attrs[patch.attr_index].raw_data, &addrs);
}
}
for (di, patches) in ds_vl.iter().enumerate() {
for patch in patches {
let addrs = place_collections(&patch.collections, &mut gcol_cursor);
patch_vl_refs(&mut all_ds[di].attrs[patch.attr_index].raw_data, &addrs);
}
}
for (i, d) in all_ds.iter().enumerate() {
if early_gcol.contains(&i) {
continue;
}
if let Some(staging) = &d.vl_string_staging {
let DsData::InMemory(ref mut bytes) = ds_layouts[i].data else {
unreachable!(
"a chunked VL-string dataset is patched before encoding, so a \
dataset patched here always has its data in memory"
);
};
let addrs = place_collections(&staging.collections, &mut gcol_cursor);
patch_vl_refs_masked(bytes, &staging.patch_offsets, &addrs);
}
}
#[expect(
clippy::cast_possible_truncation,
reason = "global-heap total size is an in-memory output span bounded by \
addressable memory on the target"
)]
{
gcol_total_size = (gcol_cursor - cursor2 as u64) as usize;
}
}
let DenseBlobs {
root: root_dense_blob,
groups: group_dense_blobs,
datasets: ds_dense_blobs,
} = dense_spans.build(&root_attrs, &groups, &all_ds)?;
let ds_oh_bytes2 = build_ds_ohs(&all_ds, &ds_layouts, &ds_dense_blobs, libver)?;
let actual_ds_oh_sizes2: Vec<usize> = ds_oh_bytes2.iter().map(|b| b.len()).collect();
debug_assert_eq!(actual_ds_oh_sizes, actual_ds_oh_sizes2);
let ext_addr = ext_oh.as_ref().map(|_| (cursor2 + gcol_total_size) as u64);
let ext_len = ext_oh.as_ref().map_or(0, |b| b.len());
let eof_addr2 = (ub + cursor2 + gcol_total_size + ext_len) as u64;
sink.reserve(eof_addr2.to_usize()?);
Self::put_userblock(sink, ub, &self.userblock_content)?;
let sb = Superblock {
version: superblock_version(libver),
offset_size: OFFSET_SIZE,
length_size: LENGTH_SIZE,
base_address: BaseAddress::new(ub as u64),
eof_address: eof_addr2,
root_group_address: root_group_addr,
group_leaf_node_k: None,
group_internal_node_k: None,
indexed_storage_internal_node_k: None,
free_space_address: None,
driver_info_address: None,
consistency_flags: 0,
superblock_extension_address: Some(ext_addr.unwrap_or(u64::MAX)),
checksum: None,
};
sink.put(&sb.serialize())?;
for &i in &early_gcol {
let staging = all_ds[i]
.vl_string_staging
.as_ref()
.expect("early_gcol only holds datasets with VL staging");
emit_collections(sink, &staging.collections)?;
}
debug_assert_eq!(
sink.position(),
ub as u64 + root_group_addr,
"early VL collections must occupy exactly the space reserved for them"
);
let tables = LinkTables {
all_ds: &all_ds,
committed: &committed,
groups: &groups,
ds_addrs: &ds_oh_addrs2,
committed_addrs: &committed_addrs,
group_addrs: &group_addrs2,
};
let root_links = tables.links(
&root_ds_indices,
&root_committed_indices,
&root_group_indices,
);
let root_oh = build_group_oh(
&root_links,
&root_attrs,
root_dense_blob
.as_ref()
.map(|b| b.attr_info_message.as_slice()),
)?;
debug_assert_eq!(root_oh.len(), root_oh_size);
sink.put(&root_oh)?;
if let Some(ref blob) = root_dense_blob {
sink.put(&blob.blob)?;
}
for (gi, g) in groups.iter().enumerate() {
let links = tables.links(&g.ds_indices, &g.committed_indices, &g.sub_group_indices);
let oh = build_group_oh(
&links,
&g.attrs,
group_dense_blobs[gi]
.as_ref()
.map(|b| b.attr_info_message.as_slice()),
)?;
debug_assert_eq!(oh.len(), group_oh_sizes[gi]);
sink.put(&oh)?;
if let Some(ref blob) = group_dense_blobs[gi] {
sink.put(&blob.blob)?;
}
}
for oh in &committed_oh {
sink.put(oh)?;
}
for (i, oh) in ds_oh_bytes2.iter().enumerate() {
sink.put(oh)?;
if let Some(ref dense) = ds_dense_blobs[i] {
sink.put(&dense.blob)?;
}
}
for (i, layout) in ds_layouts.iter().enumerate() {
emit_ds_data(
sink,
&layout.data,
all_ds[i].raw_chunks.as_ref(),
all_ds[i].produced.as_ref(),
)?;
}
for patch in &vl_root {
emit_collections(sink, &patch.collections)?;
}
for patches in &grp_vl {
for patch in patches {
emit_collections(sink, &patch.collections)?;
}
}
for patches in &ds_vl {
for patch in patches {
emit_collections(sink, &patch.collections)?;
}
}
for (i, d) in all_ds.iter().enumerate() {
if early_gcol.contains(&i) {
continue; }
if let Some(staging) = &d.vl_string_staging {
emit_collections(sink, &staging.collections)?;
}
}
let real_ext_oh = match (&ext_oh, nonpaged_persist) {
(Some(_), Some((strategy, threshold, np_page_size))) => {
let mut info = FileSpaceInfo::persistent_empty(strategy, threshold, np_page_size);
info.eoa_pre_fsm = eof_addr2 - ub as u64;
let mut oh = ObjectHeaderWriter::new();
oh.add_message_with_flags(MessageType::FileSpaceInfo, info.serialize(), 0x14);
Some(oh.serialize()?)
}
(other, _) => other.clone(),
};
debug_assert_eq!(
real_ext_oh.as_ref().map_or(0, |b| b.len()),
ext_len,
"rebuilt extension header length must match the reserved length"
);
if let Some(bytes) = &real_ext_oh {
debug_assert_eq!(
sink.position(),
ub as u64 + ext_addr.unwrap(),
"extension header must land at its recorded base-relative address"
);
sink.put(bytes)?;
}
Ok(())
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::group_v2::resolve_path_any;
use crate::link_info::LinkInfoMessage;
use crate::object_header::ObjectHeader;
use crate::signature;
use crate::type_builders::{build_attr_message, make_i32_type};
#[test]
fn a_committed_datatype_header_holds_only_its_type() {
let bytes = build_committed_datatype_oh(&make_i32_type(), 1).unwrap();
let hdr = ObjectHeader::parse(&bytes, 0, OFFSET_SIZE, LENGTH_SIZE).unwrap();
let types: Vec<MessageType> = hdr.messages.iter().map(|m| m.msg_type).collect();
assert_eq!(
types,
vec![MessageType::Datatype],
"a singly referenced committed type carries its datatype and nothing else"
);
assert_eq!(
Datatype::parse(&hdr.messages[0].data).unwrap().0,
make_i32_type()
);
assert_eq!(
hdr.messages[0].flags,
MSG_CONSTANT | MSG_DONTSHARE,
"the type of a committed object never changes and must not be shared onward"
);
}
#[test]
fn a_committed_datatype_header_records_a_count_above_one() {
let bytes = build_committed_datatype_oh(&make_i32_type(), 4).unwrap();
let hdr = ObjectHeader::parse(&bytes, 0, OFFSET_SIZE, LENGTH_SIZE).unwrap();
let refcount = hdr
.messages
.iter()
.find(|m| m.msg_type == MessageType::ObjectReferenceCount)
.expect("a count above one is stored");
assert_eq!(
refcount.data,
vec![0, 4, 0, 0, 0],
"version 0 followed by the 4-byte count"
);
}
#[test]
fn the_reference_count_is_the_link_plus_every_user() {
let mut w = FileWriter::new();
w.commit_datatype("mytype", make_i32_type());
w.set_root_attr_committed("root_attr", AttrValue::I32(1), "mytype");
let ds = w.create_dataset("typed");
ds.with_i32_data(&[1, 2]);
ds.with_committed_datatype("mytype");
ds.set_attr_committed("shared_attr", AttrValue::I32(2), "mytype");
let bytes = w.finish().unwrap();
assert_eq!(committed_reference_count(&bytes, "mytype"), Some(4));
}
#[test]
fn an_unreferenced_committed_type_stores_no_count() {
let mut w = FileWriter::new();
w.commit_datatype("mytype", make_i32_type());
let bytes = w.finish().unwrap();
assert_eq!(committed_reference_count(&bytes, "mytype"), None);
}
fn committed_reference_count(bytes: &[u8], path: &str) -> Option<u32> {
let sig = signature::find_signature(bytes).unwrap();
let sb = Superblock::parse(bytes, sig).unwrap();
let addr = resolve_path_any(bytes, &sb, path).unwrap();
let hdr =
ObjectHeader::parse(bytes, addr as usize, sb.offset_size, sb.length_size).unwrap();
let msg = hdr
.messages
.iter()
.find(|m| m.msg_type == MessageType::ObjectReferenceCount)?;
Some(u32::from_le_bytes(msg.data[1..5].try_into().unwrap()))
}
#[test]
fn a_dataset_naming_a_committed_type_stores_a_reference_to_it() {
let mut w = FileWriter::new();
w.commit_datatype("mytype", make_i32_type());
let ds = w.create_dataset("typed");
ds.with_i32_data(&[1, 2]);
ds.with_committed_datatype("mytype");
let bytes = w.finish().unwrap();
let sig = signature::find_signature(&bytes).unwrap();
let sb = Superblock::parse(&bytes, sig).unwrap();
let type_addr = resolve_path_any(&bytes, &sb, "mytype").unwrap();
let ds_addr = resolve_path_any(&bytes, &sb, "typed").unwrap();
let hdr =
ObjectHeader::parse(&bytes, ds_addr as usize, sb.offset_size, sb.length_size).unwrap();
let msg = hdr
.messages
.iter()
.find(|m| m.msg_type == MessageType::Datatype)
.expect("a dataset always has a datatype message");
assert!(
crate::shared_message::is_shared(msg.flags),
"the record must say its body is a reference, or the reference decodes as a type"
);
assert_eq!(
msg.data,
crate::shared_message::encode_committed_ref(type_addr, OFFSET_SIZE),
"the reference must name the object the link resolves to"
);
}
fn parse_file(bytes: &[u8]) -> (Superblock, ObjectHeader) {
let sig = signature::find_signature(bytes).unwrap();
let sb = Superblock::parse(bytes, sig).unwrap();
let oh = ObjectHeader::parse(
bytes,
sb.root_group_address as usize,
sb.offset_size,
sb.length_size,
)
.unwrap();
(sb, oh)
}
fn read_dataset_f64(bytes: &[u8], path: &str) -> Vec<f64> {
let sig = signature::find_signature(bytes).unwrap();
let sb = Superblock::parse(bytes, sig).unwrap();
let addr = resolve_path_any(bytes, &sb, path).unwrap();
let hdr =
ObjectHeader::parse(bytes, addr as usize, sb.offset_size, sb.length_size).unwrap();
let dt_data = &hdr
.messages
.iter()
.find(|m| m.msg_type == MessageType::Datatype)
.unwrap()
.data;
let ds_data = &hdr
.messages
.iter()
.find(|m| m.msg_type == MessageType::Dataspace)
.unwrap()
.data;
let dl_data = &hdr
.messages
.iter()
.find(|m| m.msg_type == MessageType::DataLayout)
.unwrap()
.data;
let (dt, _) = Datatype::parse(dt_data).unwrap();
let ds = Dataspace::parse(ds_data, sb.length_size).unwrap();
let dl =
crate::data_layout::DataLayout::parse(dl_data, sb.offset_size, sb.length_size).unwrap();
let raw = crate::data_read::read_raw_data(bytes, &dl, &ds, &dt).unwrap();
crate::data_read::read_as_f64(&raw, &dt).unwrap()
}
#[test]
fn empty_file_root_group_only() {
let fw = FileWriter::new();
let bytes = fw.finish().unwrap();
let (sb, oh) = parse_file(&bytes);
assert_eq!(sb.version, 3);
assert_eq!(oh.version, 2);
}
#[test]
fn file_with_f64_dataset() {
let mut fw = FileWriter::new();
fw.create_dataset("data").with_f64_data(&[1.0, 2.0, 3.0]);
let bytes = fw.finish().unwrap();
assert_eq!(read_dataset_f64(&bytes, "data"), vec![1.0, 2.0, 3.0]);
}
#[test]
fn file_with_dataset_attrs() {
let mut fw = FileWriter::new();
fw.create_dataset("data")
.with_f64_data(&[1.0, 2.0])
.set_attr("scale", AttrValue::F64(0.5));
let bytes = fw.finish().unwrap();
assert_eq!(read_dataset_f64(&bytes, "data"), vec![1.0, 2.0]);
let sig = signature::find_signature(&bytes).unwrap();
let sb = Superblock::parse(&bytes, sig).unwrap();
let addr = resolve_path_any(&bytes, &sb, "data").unwrap();
let hdr =
ObjectHeader::parse(&bytes, addr as usize, sb.offset_size, sb.length_size).unwrap();
let attrs = crate::attribute::extract_attributes(&hdr, sb.length_size).unwrap();
assert_eq!(attrs.len(), 1);
assert_eq!(attrs[0].name, "scale");
}
#[test]
fn file_with_group_and_dataset() {
let mut fw = FileWriter::new();
let mut gb = fw.create_group("grp");
gb.create_dataset("vals").with_f64_data(&[10.0, 20.0]);
fw.add_group(gb.finish());
let bytes = fw.finish().unwrap();
assert_eq!(read_dataset_f64(&bytes, "grp/vals"), vec![10.0, 20.0]);
}
#[test]
fn root_group_carries_no_timestamps() {
let fw = FileWriter::new();
let bytes = fw.finish().unwrap();
let (_, oh) = parse_file(&bytes);
assert_eq!(oh.flags & 0x20, 0, "times-stored flag must be clear");
assert!(oh.modification_time.is_none());
assert!(oh.access_time.is_none());
assert!(oh.change_time.is_none());
assert!(oh.birth_time.is_none());
}
#[test]
fn sub_group_carries_no_timestamps() {
let mut fw = FileWriter::new();
let mut gb = fw.create_group("grp");
gb.create_dataset("vals").with_f64_data(&[1.0]);
fw.add_group(gb.finish());
let bytes = fw.finish().unwrap();
let sig = signature::find_signature(&bytes).unwrap();
let sb = Superblock::parse(&bytes, sig).unwrap();
let addr = resolve_path_any(&bytes, &sb, "grp").unwrap();
let hdr =
ObjectHeader::parse(&bytes, addr as usize, sb.offset_size, sb.length_size).unwrap();
assert_eq!(hdr.flags & 0x20, 0, "times-stored flag must be clear");
assert!(hdr.modification_time.is_none());
}
#[test]
fn group_links_stay_compact_regardless_of_child_count() {
let mut fw = FileWriter::new();
let mut gb = fw.create_group("grp");
for i in 0..20 {
gb.create_dataset(&format!("d{i}"))
.with_f64_data(&[i as f64]);
}
fw.add_group(gb.finish());
let bytes = fw.finish().unwrap();
let sig = signature::find_signature(&bytes).unwrap();
let sb = Superblock::parse(&bytes, sig).unwrap();
let addr = resolve_path_any(&bytes, &sb, "grp").unwrap();
let hdr =
ObjectHeader::parse(&bytes, addr as usize, sb.offset_size, sb.length_size).unwrap();
let link_info_msg = hdr
.messages
.iter()
.find(|m| m.msg_type == MessageType::LinkInfo)
.unwrap();
let link_info = LinkInfoMessage::parse(&link_info_msg.data, sb.offset_size).unwrap();
assert!(
link_info.fractal_heap_address.is_none(),
"no dense link storage is ever used"
);
let group_info_msg = hdr
.messages
.iter()
.find(|m| m.msg_type == MessageType::GroupInfo)
.unwrap();
assert_eq!(group_info_msg.data, vec![0, 0]);
let link_count = hdr
.messages
.iter()
.filter(|m| m.msg_type == MessageType::Link)
.count();
assert_eq!(link_count, 20);
}
#[test]
fn file_with_root_attr() {
let mut fw = FileWriter::new();
fw.set_root_attr("version", AttrValue::I64(42));
let bytes = fw.finish().unwrap();
let (sb, oh) = parse_file(&bytes);
let attrs = crate::attribute::extract_attributes(&oh, sb.length_size).unwrap();
assert_eq!(attrs[0].name, "version");
}
#[test]
fn dense_attrs_self_roundtrip() {
let mut fw = FileWriter::new();
let ds = fw.create_dataset("data");
ds.with_f64_data(&[1.0, 2.0, 3.0]);
for i in 0..20 {
ds.set_attr(&format!("attr_{i:03}"), AttrValue::F64(i as f64 * 1.5));
}
let bytes = fw.finish().unwrap();
let sig = signature::find_signature(&bytes).unwrap();
let sb = Superblock::parse(&bytes, sig).unwrap();
let addr = resolve_path_any(&bytes, &sb, "data").unwrap();
let hdr =
ObjectHeader::parse(&bytes, addr as usize, sb.offset_size, sb.length_size).unwrap();
let attrs =
crate::attribute::extract_attributes_full(&bytes, &hdr, sb.offset_size, sb.length_size)
.unwrap();
assert_eq!(attrs.len(), 20);
for i in 0..20 {
let attr = attrs
.iter()
.find(|a| a.name == format!("attr_{i:03}"))
.unwrap();
let v = attr.read_as_f64().unwrap();
assert!((v[0] - i as f64 * 1.5).abs() < 1e-10);
}
assert_eq!(read_dataset_f64(&bytes, "data"), vec![1.0, 2.0, 3.0]);
}
#[test]
fn dense_attrs_root_group_self_roundtrip() {
let mut fw = FileWriter::new();
fw.create_dataset("dummy").with_f64_data(&[0.0]);
for i in 0..15 {
fw.set_root_attr(&format!("root_{i:02}"), AttrValue::F64(i as f64 * 2.0));
}
let bytes = fw.finish().unwrap();
let sig = signature::find_signature(&bytes).unwrap();
let sb = Superblock::parse(&bytes, sig).unwrap();
let oh = ObjectHeader::parse(
&bytes,
sb.root_group_address as usize,
sb.offset_size,
sb.length_size,
)
.unwrap();
let attrs =
crate::attribute::extract_attributes_full(&bytes, &oh, sb.offset_size, sb.length_size)
.unwrap();
assert_eq!(attrs.len(), 15);
}
#[test]
fn inline_attrs_below_threshold() {
let mut fw = FileWriter::new();
let ds = fw.create_dataset("data");
ds.with_f64_data(&[1.0]);
for i in 0..5 {
ds.set_attr(&format!("a{i}"), AttrValue::F64(i as f64));
}
let bytes = fw.finish().unwrap();
let sig = signature::find_signature(&bytes).unwrap();
let sb = Superblock::parse(&bytes, sig).unwrap();
let addr = resolve_path_any(&bytes, &sb, "data").unwrap();
let hdr =
ObjectHeader::parse(&bytes, addr as usize, sb.offset_size, sb.length_size).unwrap();
let info = hdr
.messages
.iter()
.find(|m| m.msg_type == MessageType::AttributeInfo)
.map(|m| {
crate::attribute_info::AttributeInfoMessage::parse(&m.data, sb.offset_size).unwrap()
})
.expect("a compact attribute set still carries an Attribute Info message");
assert_eq!(info.fractal_heap_address, None);
assert_eq!(info.btree_name_index_address, None);
let attrs = crate::attribute::extract_attributes(&hdr, sb.length_size).unwrap();
assert_eq!(attrs.len(), 5);
}
#[test]
fn encode_decode_managed_id_roundtrip() {
let id = encode_managed_id(100, 42, 40, 8);
let fh = crate::fractal_heap::FractalHeapHeader {
heap_id_length: 8,
io_filter_encoded_length: 0,
max_managed_object_size: 1024,
btree_huge_objects_address: u64::MAX,
table_width: 4,
starting_block_size: 4096,
max_direct_block_size: 65536,
max_heap_size: 40,
start_root_rows: 1,
root_block_address: 0,
current_rows_in_root_indirect_block: 0,
managed_objects_count: 0,
};
let (off, len) = fh.decode_managed_id(&id).unwrap();
assert_eq!(off, 100);
assert_eq!(len, 42);
}
fn dense_attr_of_size(name: &str, size: usize) -> AttributeMessage {
let probe = build_attr_message(name, &AttrValue::AsciiString("y".to_string()));
let overhead = probe.serialize_v3(LENGTH_SIZE).len() - 1;
let attr = build_attr_message(name, &AttrValue::AsciiString("y".repeat(size - overhead)));
assert_eq!(attr.serialize_v3(LENGTH_SIZE).len(), size);
attr
}
#[test]
fn dense_attrs_check_bounds_each_attribute_not_the_total() {
let many: Vec<AttributeMessage> = (0..40)
.map(|i| dense_attr_of_size(&format!("a{i}"), 60_000))
.collect();
let total: usize = many.iter().map(|a| a.serialize_v3(LENGTH_SIZE).len()).sum();
assert!(total > 2_000_000, "expected a multi-megabyte set");
assert_eq!(dense_attrs_check(&many), Ok(()));
let at_limit = vec![dense_attr_of_size("edge", DENSE_ATTR_MAX_MANAGED_OBJECT)];
assert_eq!(dense_attrs_check(&at_limit), Ok(()));
let past = vec![dense_attr_of_size(
"edge",
DENSE_ATTR_MAX_MANAGED_OBJECT + 1,
)];
assert_eq!(dense_attrs_check(&past), Ok(()));
assert_eq!(huge_object_count(&build_dense_attrs(&past, 0).blob), 1);
assert_eq!(huge_object_count(&build_dense_attrs(&at_limit, 0).blob), 0);
}
#[test]
fn dense_attr_plan_length_matches_what_it_builds() {
let mut shapes: Vec<(String, Vec<AttributeMessage>)> = Vec::new();
for n in 0..=64usize {
shapes.push((
format!("{n} x 512B"),
(0..n)
.map(|i| dense_attr_of_size(&format!("a{i}"), 512))
.collect(),
));
}
for n in [100usize, 200] {
shapes.push((
format!("{n} x 512B"),
(0..n)
.map(|i| dense_attr_of_size(&format!("a{i}"), 512))
.collect(),
));
}
shapes.push((
"40 x 60,000B".to_string(),
(0..40)
.map(|i| dense_attr_of_size(&format!("big{i}"), 60_000))
.collect(),
));
shapes.push((
"at the managed limit".to_string(),
vec![dense_attr_of_size("edge", DENSE_ATTR_MAX_MANAGED_OBJECT)],
));
shapes.push((
"one past the managed limit".to_string(),
vec![dense_attr_of_size(
"edge",
DENSE_ATTR_MAX_MANAGED_OBJECT + 1,
)],
));
let mut mixed: Vec<AttributeMessage> = (0..10)
.map(|i| dense_attr_of_size(&format!("small{i}"), 300))
.collect();
mixed.push(dense_attr_of_size(
"huge0",
DENSE_ATTR_MAX_MANAGED_OBJECT + 1,
));
mixed.push(dense_attr_of_size(
"huge1",
DENSE_ATTR_MAX_MANAGED_OBJECT * 4,
));
mixed.extend((10..20).map(|i| dense_attr_of_size(&format!("small{i}"), 300)));
shapes.push(("managed and huge mixed".to_string(), mixed));
for (label, attrs) in shapes {
assert_eq!(dense_attrs_check(&attrs), Ok(()), "{label}");
let plan = dense_attrs_plan(&attrs);
for base in [0u64, 0x1000, 0x8000_0000] {
let built = plan.build(base);
assert_eq!(
plan.blob_len(),
built.blob.len() as u64,
"planned length must match the emitted heap for {label} at base {base:#x}"
);
}
}
}
fn name_index_order(attrs: &[AttributeMessage]) -> Vec<String> {
const RECORD: usize = 8 + 1 + 4 + 4;
let blob = build_dense_attrs(attrs, 0).blob;
let header = blob
.windows(4)
.position(|w| w == b"BTHD")
.expect("a name index has a header");
let depth = u16::from_le_bytes(blob[header + 12..header + 14].try_into().expect("2 bytes"));
assert_eq!(depth, 0, "the fixture outgrew a single leaf");
let leaf = blob
.windows(4)
.position(|w| w == b"BTLF")
.expect("a name index has a leaf node");
(0..attrs.len())
.map(|i| {
let at = leaf + 6 + i * RECORD + 9;
let order = u32::from_le_bytes(blob[at..at + 4].try_into().expect("4 bytes"));
attrs[order as usize].name.clone()
})
.collect()
}
#[test]
fn the_name_index_order_does_not_depend_on_insertion_order() {
let names = [
"k69209", "k155448", "zeta", "alpha", "m", "beta", "gamma", "delta", "epsilon", "eta",
];
let forward: Vec<AttributeMessage> = names
.iter()
.map(|n| build_attr_message(n, &AttrValue::I64(1)))
.collect();
let mut reversed = forward.clone();
reversed.reverse();
let order = name_index_order(&forward);
assert_eq!(
order,
name_index_order(&reversed),
"the index order changed with the insertion order"
);
assert_eq!(
crate::checksum::jenkins_lookup3(b"k69209"),
crate::checksum::jenkins_lookup3(b"k155448"),
"the fixture names no longer hash alike; pick a new colliding pair"
);
let at = |name: &str| {
order
.iter()
.position(|n| n == name)
.expect("every attribute is indexed")
};
assert_eq!(
at("k155448") + 1,
at("k69209"),
"the colliding pair is not indexed in name order"
);
}
fn huge_object_count(blob: &[u8]) -> u64 {
let ls = LENGTH_SIZE as usize;
let os = OFFSET_SIZE as usize;
assert_eq!(&blob[..4], b"FRHP");
let at = 4 + 1 + 2 + 2 + 1 + 4 + ls + os + ls + os + ls + ls + ls + ls + ls;
u64::from_le_bytes(blob[at..at + 8].try_into().unwrap())
}
#[test]
fn a_fixed_node_size_keeps_the_derived_count_width_at_one_byte() {
for record_size in [DENSE_ATTR_BTREE_RECORD, DENSE_ATTR_HUGE_BTREE_RECORD] {
let (info, depth) = crate::btree_v2::NodeInfo::for_record_count(
btree_v2_write::NODE_SIZE,
record_size,
OFFSET_SIZE,
0,
)
.expect("the emitted geometry is plannable");
assert_eq!(depth, 0);
let capacity = info.max_nrec(0);
assert!(
capacity >= 20,
"a {record_size}-byte record should leave room for 20 per node, got {capacity}"
);
assert_eq!(
encoded_size_width(capacity),
1,
"a {record_size}-byte record's leaf capacity must stay inside one byte"
);
}
}
#[test]
fn the_managed_object_limit_is_where_storage_changes_class() {
let mixed = vec![
dense_attr_of_size("at", DENSE_ATTR_MAX_MANAGED_OBJECT),
dense_attr_of_size("past", DENSE_ATTR_MAX_MANAGED_OBJECT + 1),
];
assert_eq!(dense_attrs_check(&mixed), Ok(()));
let blob = build_dense_attrs(&mixed, 0).blob;
assert_eq!(huge_object_count(&blob), 1);
assert_eq!(managed_object_count(&blob), 1);
}
fn managed_object_count(blob: &[u8]) -> u64 {
let ls = LENGTH_SIZE as usize;
let os = OFFSET_SIZE as usize;
assert_eq!(&blob[..4], b"FRHP");
let at = 4 + 1 + 2 + 2 + 1 + 4 + ls + os + ls + os + ls + ls + ls;
u64::from_le_bytes(blob[at..at + 8].try_into().unwrap())
}
fn root_block(blob: &[u8]) -> (usize, u16) {
let ls = LENGTH_SIZE as usize;
let os = OFFSET_SIZE as usize;
assert_eq!(&blob[..4], b"FRHP");
let at = 4 + 1 + 2 + 2 + 1 + 4 + ls + os + ls + os + 8 * ls + 2 + ls + ls + 2 + 2;
let address = u64::from_le_bytes(blob[at..at + os].try_into().unwrap());
let rows = u16::from_le_bytes(blob[at + os..at + os + 2].try_into().unwrap());
(address as usize, rows)
}
#[test]
fn the_root_grows_into_an_indirect_block_rather_than_a_bigger_direct_one() {
let fits = vec![dense_attr_of_size("a", 400)];
let blob = build_dense_attrs(&fits, 0).blob;
let (address, rows) = root_block(&blob);
assert_eq!(rows, 0, "one starting-size block still holds this heap");
assert_eq!(&blob[address..address + 4], b"FHDB");
let spills: Vec<AttributeMessage> = (0..8)
.map(|i| dense_attr_of_size(&format!("a{i}"), 400))
.collect();
let blob = build_dense_attrs(&spills, 0).blob;
let (address, rows) = root_block(&blob);
assert!(rows >= 1, "content past one block needs an indirect root");
assert_eq!(&blob[address..address + 4], b"FHIB");
}
#[test]
fn attributes_survive_a_heap_with_nested_indirect_blocks() {
let mut fw = FileWriter::new();
let ds = fw.create_dataset("data");
ds.with_f64_data(&[1.0]);
let count = DENSE_ATTR_THRESHOLD + 2;
let value = |i: usize| char::from(b'a' + i as u8).to_string().repeat(60_000);
for i in 0..count {
ds.set_attr(&format!("big{i}"), AttrValue::AsciiString(value(i)));
}
let bytes = fw.finish().unwrap();
let nested = bytes.windows(4).filter(|w| *w == b"FHIB").count();
assert!(
nested >= 2,
"expected a nested indirect block, got {nested}"
);
let sig = signature::find_signature(&bytes).unwrap();
let sb = Superblock::parse(&bytes, sig).unwrap();
let addr = resolve_path_any(&bytes, &sb, "data").unwrap();
let hdr =
ObjectHeader::parse(&bytes, addr as usize, sb.offset_size, sb.length_size).unwrap();
let attrs =
crate::attribute::extract_attributes_full(&bytes, &hdr, sb.offset_size, sb.length_size)
.unwrap();
assert_eq!(attrs.len(), count);
for i in 0..count {
let attr = attrs.iter().find(|a| a.name == format!("big{i}")).unwrap();
assert_eq!(attr.read_as_string().unwrap(), value(i));
}
}
fn largest_fitting_i64_attr_elements() -> usize {
let one = build_attr_message("boundary", &AttrValue::I64Array(vec![0i64; 1]));
let overhead = one.serialize(LENGTH_SIZE).len() - 8;
(OBJECT_HEADER_MESSAGE_MAX - overhead) / 8
}
#[test]
fn compact_attr_at_the_message_size_limit_is_written() {
let n = largest_fitting_i64_attr_elements();
let attr = build_attr_message("boundary", &AttrValue::I64Array(vec![7i64; n]));
let size = attr.serialize(LENGTH_SIZE).len();
assert!(size <= OBJECT_HEADER_MESSAGE_MAX);
assert!(
size + 8 > OBJECT_HEADER_MESSAGE_MAX,
"probe is not at the limit (got {size})"
);
let mut fw = FileWriter::new();
fw.set_root_attr("boundary", AttrValue::I64Array(vec![7i64; n]));
fw.create_dataset("d").with_f64_data(&[1.0]);
let bytes = fw.finish().unwrap();
let file = crate::reader::File::from_bytes(bytes).unwrap();
let attrs = file.root().attrs().unwrap();
assert_eq!(attrs.len(), 1);
}
#[test]
fn an_attr_past_the_message_size_limit_moves_to_dense_storage() {
let n = largest_fitting_i64_attr_elements() + 1;
let attrs = vec![build_attr_message(
"boundary",
&AttrValue::I64Array(vec![0i64; n]),
)];
assert!(attrs[0].serialize(LENGTH_SIZE).len() > OBJECT_HEADER_MESSAGE_MAX);
assert!(
needs_dense_attrs(&attrs),
"one oversized attribute must select dense storage by itself"
);
let mut fw = FileWriter::new();
fw.set_root_attr("boundary", AttrValue::I64Array(vec![7i64; n]));
fw.create_dataset("d").with_f64_data(&[1.0]);
let bytes = fw.finish().expect("written, not refused");
let file = crate::reader::File::from_bytes(bytes).unwrap();
let attrs = file.root().attrs().unwrap();
assert_eq!(attrs.len(), 1);
match attrs.get("boundary") {
Some(AttrValue::I64Array(v)) => assert_eq!(v.len(), n),
other => panic!("expected the attribute back, got {other:?}"),
}
}
#[test]
fn an_attr_at_the_message_size_limit_stays_compact() {
let n = largest_fitting_i64_attr_elements();
let attrs = vec![build_attr_message(
"boundary",
&AttrValue::I64Array(vec![0i64; n]),
)];
assert!(!needs_dense_attrs(&attrs));
}
fn read_vl_bytes(bytes: Vec<u8>, path: &str) -> Vec<crate::vl_data::VlByteObject> {
let file = crate::reader::File::from_bytes(bytes).unwrap();
file.dataset(path)
.unwrap()
.read_vlen_string_bytes(crate::vl_data::VlenStringReadOptions::default())
.unwrap()
}
#[test]
fn vlen_string_dataset_roundtrips_values() {
let mut fw = FileWriter::new();
fw.create_dataset("labels")
.with_vlen_strings(&["alpha", "beta", "gamma"]);
let bytes = fw.finish().unwrap();
let objs = read_vl_bytes(bytes, "labels");
let got: Vec<_> = objs
.iter()
.map(|o| match o {
crate::vl_data::VlByteObject::Bytes(b) => String::from_utf8(b.clone()).unwrap(),
crate::vl_data::VlByteObject::Null => "<null>".to_string(),
})
.collect();
assert_eq!(got, vec!["alpha", "beta", "gamma"]);
}
#[test]
fn vlen_string_dataset_preserves_null_vs_empty() {
use crate::type_builders::VlStringElement;
use crate::vl_data::VlByteObject;
let dt = crate::type_builders::make_vlen_string_type(CharacterSet::Utf8);
let elements = vec![
VlStringElement::Bytes(b"hi".to_vec()),
VlStringElement::Null,
VlStringElement::Bytes(Vec::new()), VlStringElement::Bytes(b"end".to_vec()),
];
let mut fw = FileWriter::new();
fw.create_dataset("mixed")
.with_vlen_string_elements(dt, &elements)
.unwrap();
let bytes = fw.finish().unwrap();
let objs = read_vl_bytes(bytes, "mixed");
assert_eq!(
objs,
vec![
VlByteObject::Bytes(b"hi".to_vec()),
VlByteObject::Null,
VlByteObject::Bytes(Vec::new()),
VlByteObject::Bytes(b"end".to_vec()),
]
);
}
#[test]
fn vlen_string_dataset_preserves_embedded_nul() {
use crate::type_builders::VlStringElement;
use crate::vl_data::VlByteObject;
let dt = crate::type_builders::make_vlen_string_type(CharacterSet::Ascii);
let payload = b"a\0b\0c".to_vec();
let elements = vec![VlStringElement::Bytes(payload.clone())];
let mut fw = FileWriter::new();
fw.create_dataset("nul")
.with_vlen_string_elements(dt, &elements)
.unwrap();
let bytes = fw.finish().unwrap();
let objs = read_vl_bytes(bytes, "nul");
assert_eq!(objs, vec![VlByteObject::Bytes(payload)]);
}
#[test]
fn vlen_string_dataset_preserves_non_utf8_bytes() {
use crate::type_builders::VlStringElement;
use crate::vl_data::VlByteObject;
let dt = crate::type_builders::make_vlen_string_type(CharacterSet::Ascii);
let payload = vec![0xffu8, 0xfe, 0x80, 0x00, 0x41];
let elements = vec![VlStringElement::Bytes(payload.clone())];
let mut fw = FileWriter::new();
fw.create_dataset("raw")
.with_vlen_string_elements(dt, &elements)
.unwrap();
let bytes = fw.finish().unwrap();
let objs = read_vl_bytes(bytes, "raw");
assert_eq!(objs, vec![VlByteObject::Bytes(payload)]);
}
#[test]
fn vlen_string_dataset_2d_shape_roundtrips() {
let mut fw = FileWriter::new();
fw.create_dataset("grid")
.with_vlen_strings(&["a", "bb", "ccc", "dddd"])
.with_shape(&[2, 2]);
let bytes = fw.finish().unwrap();
let file = crate::reader::File::from_bytes(bytes).unwrap();
let ds = file.dataset("grid").unwrap();
assert_eq!(ds.shape().unwrap(), vec![2, 2]);
assert_eq!(
ds.read_vlen_strings(crate::vl_data::VlenStringReadOptions::default())
.unwrap(),
vec!["a", "bb", "ccc", "dddd"]
);
}
#[test]
fn vlen_string_dataset_all_null_no_heap() {
use crate::type_builders::VlStringElement;
use crate::vl_data::VlByteObject;
let dt = crate::type_builders::make_vlen_string_type(CharacterSet::Utf8);
let elements = vec![VlStringElement::Null, VlStringElement::Null];
let mut fw = FileWriter::new();
fw.create_dataset("nulls")
.with_vlen_string_elements(dt, &elements)
.unwrap();
let bytes = fw.finish().unwrap();
let objs = read_vl_bytes(bytes, "nulls");
assert_eq!(objs, vec![VlByteObject::Null, VlByteObject::Null]);
}
#[test]
fn vlen_string_dataset_with_nulls_spans_multiple_heap_collections() {
use crate::type_builders::VlStringElement;
use crate::vl_data::VlByteObject;
let count = 100_000;
let elements: Vec<VlStringElement> = (0..count)
.map(|i| {
if i % 3 == 0 {
VlStringElement::Null
} else {
VlStringElement::Bytes(format!("s{i}").into_bytes())
}
})
.collect();
let dt = crate::type_builders::make_vlen_string_type(CharacterSet::Utf8);
let mut fw = FileWriter::new();
fw.create_dataset("mixed")
.with_vlen_string_elements(dt, &elements)
.unwrap();
let bytes = fw.finish().unwrap();
let objs = read_vl_bytes(bytes, "mixed");
assert_eq!(objs.len(), count);
for i in [0, 1, 65_535, 98_302, 98_303, 98_304, count - 1] {
let expected = if i % 3 == 0 {
VlByteObject::Null
} else {
VlByteObject::Bytes(format!("s{i}").into_bytes())
};
assert_eq!(objs[i], expected, "element {i} did not round-trip");
}
}
#[test]
fn chunked_vlen_string_dataset_roundtrips() {
let mut fw = FileWriter::new();
fw.create_dataset("chunked")
.with_vlen_strings(&["a", "bb", "ccc", "dddd"])
.with_chunks(&[2]);
let bytes = fw.finish().unwrap();
let f = crate::reader::File::from_bytes(bytes).unwrap();
let ds = f.dataset("chunked").unwrap();
assert_eq!(ds.read_string().unwrap(), ["a", "bb", "ccc", "dddd"]);
}
#[test]
#[cfg(feature = "deflate")]
fn filtered_chunked_vlen_string_dataset_roundtrips() {
let mut fw = FileWriter::new();
fw.create_dataset("filtered")
.with_vlen_strings(&["alpha", "beta", "gamma", "delta"])
.with_chunks(&[2])
.with_deflate(6);
let bytes = fw.finish().unwrap();
let f = crate::reader::File::from_bytes(bytes).unwrap();
let ds = f.dataset("filtered").unwrap();
assert_eq!(
ds.read_string().unwrap(),
["alpha", "beta", "gamma", "delta"]
);
}
#[test]
fn resizable_vlen_string_dataset_roundtrips() {
let mut fw = FileWriter::new();
fw.create_dataset("growable")
.with_vlen_strings(&["one", "two", "three"])
.with_shape(&[3])
.with_maxshape(&[u64::MAX])
.with_chunks(&[2]);
let bytes = fw.finish().unwrap();
let f = crate::reader::File::from_bytes(bytes).unwrap();
let ds = f.dataset("growable").unwrap();
assert_eq!(ds.read_string().unwrap(), ["one", "two", "three"]);
}
#[test]
fn chunked_vlen_string_dataset_preserves_nulls() {
use crate::type_builders::VlStringElement;
use crate::vl_data::VlByteObject;
let dt = crate::type_builders::make_vlen_string_type(CharacterSet::Utf8);
let elements = vec![
VlStringElement::Bytes(b"set".to_vec()),
VlStringElement::Null,
VlStringElement::Bytes(Vec::new()), VlStringElement::Bytes(b"tail".to_vec()),
];
let mut fw = FileWriter::new();
fw.create_dataset("mixed")
.with_vlen_string_elements(dt, &elements)
.unwrap()
.with_chunks(&[2]);
let bytes = fw.finish().unwrap();
assert_eq!(
read_vl_bytes(bytes, "mixed"),
vec![
VlByteObject::Bytes(b"set".to_vec()),
VlByteObject::Null,
VlByteObject::Bytes(Vec::new()),
VlByteObject::Bytes(b"tail".to_vec()),
]
);
}
#[test]
fn contiguous_vlen_layout_is_unchanged_by_the_early_path() {
let build = || {
let mut fw = FileWriter::new();
fw.create_dataset("plain")
.with_vlen_strings(&["a", "bb", "ccc"]);
fw.finish().unwrap()
};
let bytes = build();
assert_eq!(
&bytes[SUPERBLOCK_SIZE..SUPERBLOCK_SIZE + 4],
b"OHDR",
"a contiguous-only VL file must still open with the root OH at SUPERBLOCK_SIZE"
);
}
#[test]
fn vlen_sequence_dataset_roundtrips_i32() {
use crate::type_builders::VlStringElement;
use crate::vl_data::{VlByteObject, VlenStringReadOptions};
let dt = Datatype::VariableLength {
is_string: false,
padding: None,
charset: None,
base_type: Box::new(crate::type_builders::make_i32_type()),
};
let seqs: Vec<Vec<i32>> = vec![vec![1, 2, 3], vec![], vec![-7, 42]];
let elements: Vec<VlStringElement> = seqs
.iter()
.map(|s| VlStringElement::Bytes(s.iter().flat_map(|v| v.to_le_bytes()).collect()))
.collect();
let mut fw = FileWriter::new();
fw.create_dataset("seq")
.with_vlen_sequence_elements(dt, &elements)
.unwrap();
let bytes = fw.finish().unwrap();
let file = crate::reader::File::from_bytes(bytes).unwrap();
let ds = file.dataset("seq").unwrap();
assert!(
matches!(
ds.datatype().unwrap(),
Datatype::VariableLength {
is_string: false,
..
}
),
"datatype must stay a non-string variable-length sequence"
);
let (objs, elem_size) = ds
.read_vlen_sequence_bytes(VlenStringReadOptions::default())
.unwrap();
assert_eq!(elem_size, 4);
let got: Vec<Vec<i32>> = objs
.iter()
.map(|o| match o {
VlByteObject::Null => Vec::new(),
VlByteObject::Bytes(b) => b
.as_chunks::<4>()
.0
.iter()
.map(|c| i32::from_le_bytes(*c))
.collect(),
})
.collect();
assert_eq!(got, seqs);
}
#[test]
fn vlen_sequence_rejects_string_datatype() {
use crate::type_builders::VlStringElement;
let dt = crate::type_builders::make_vlen_string_type(CharacterSet::Utf8);
let mut fw = FileWriter::new();
let res = fw
.create_dataset("x")
.with_vlen_sequence_elements(dt, &[VlStringElement::Bytes(b"hi".to_vec())]);
assert!(matches!(res, Err(FormatError::TypeMismatch { .. })));
}
fn write_bounded(bounds: Option<(LibVer, LibVer)>) -> Result<Vec<u8>, FormatError> {
let mut fw = FileWriter::new();
if let Some((low, high)) = bounds {
fw.with_libver_bounds(low, high);
}
fw.create_dataset("values").with_f64_data(&[1.0, 2.0, 3.0]);
fw.finish()
}
fn layout_message_version(bytes: &[u8], path: &str) -> u8 {
let sig = signature::find_signature(bytes).unwrap();
let sb = Superblock::parse(bytes, sig).unwrap();
let addr = resolve_path_any(bytes, &sb, path).unwrap();
let oh = ObjectHeader::parse(bytes, addr as usize, sb.offset_size, sb.length_size).unwrap();
oh.messages
.iter()
.find(|m| m.msg_type == MessageType::DataLayout)
.expect("a dataset carries a data-layout message")
.data[0]
}
#[test]
fn libver_bounds_select_the_superblock_and_layout_versions() {
for (bounds, superblock, layout) in [
(None, 3, 4),
(Some((LibVer::Earliest, LibVer::LATEST)), 3, 4),
(Some((LibVer::Earliest, LibVer::V110)), 3, 4),
(Some((LibVer::V110, LibVer::LATEST)), 3, 4),
(Some((LibVer::Earliest, LibVer::V18)), 2, 3),
(Some((LibVer::V18, LibVer::V18)), 2, 3),
] {
let bytes = write_bounded(bounds).expect("these bounds are satisfiable");
let sig = signature::find_signature(&bytes).unwrap();
assert_eq!(
bytes[sig + 8],
superblock,
"superblock version under bounds {bounds:?}"
);
assert_eq!(
layout_message_version(&bytes, "values"),
layout,
"data-layout message version under bounds {bounds:?}"
);
}
}
#[test]
fn the_two_formats_differ_only_in_their_version_bytes() {
let newer = write_bounded(Some((LibVer::Earliest, LibVer::V110))).unwrap();
let older = write_bounded(Some((LibVer::Earliest, LibVer::V18))).unwrap();
assert_eq!(newer.len(), older.len(), "the two formats differ in size");
let differing: Vec<usize> = (0..newer.len()).filter(|&i| newer[i] != older[i]).collect();
let sig = signature::find_signature(&newer).unwrap();
assert_eq!(
differing.len(),
10,
"expected two version bytes and the two checksums over them, got {differing:?}"
);
assert_eq!(
&differing[..5],
&[sig + 8, sig + 44, sig + 45, sig + 46, sig + 47],
"the superblock version byte and its checksum must be the first to differ"
);
assert_eq!(
newer[differing[5]], 4,
"the sixth differing byte is the data-layout message version"
);
assert_eq!(older[differing[5]], 3);
let trailing = &differing[6..];
assert!(
trailing[0] > differing[5] && trailing.windows(2).all(|w| w[1] == w[0] + 1),
"the last four differing bytes must be the contiguous checksum trailing \
the object header that holds the layout message, got {differing:?}"
);
}
#[test]
fn bounds_admitting_no_format_this_crate_writes_are_refused() {
for (low, high) in [
(LibVer::Earliest, LibVer::Earliest),
(LibVer::V112, LibVer::LATEST),
] {
let err = write_bounded(Some((low, high))).unwrap_err();
assert!(
matches!(err, FormatError::LibverBoundsUnsatisfiable { .. }),
"bounds [{}, {}] gave {err:?}",
low.name(),
high.name()
);
}
}
#[test]
fn content_the_older_format_cannot_carry_is_refused_not_upgraded() {
let mut fw = FileWriter::new();
fw.with_libver_bounds(LibVer::Earliest, LibVer::V18);
fw.create_dataset("d")
.with_i32_data(&(0..100).collect::<Vec<_>>())
.with_chunks(&[10]);
assert!(matches!(
fw.finish().unwrap_err(),
FormatError::LibverTooOldForContent { .. }
));
let mut fw = FileWriter::new();
fw.with_libver_bounds(LibVer::Earliest, LibVer::V18);
fw.with_file_space_strategy(FileSpaceStrategy::Page, true, 1);
fw.create_dataset("d").with_i32_data(&[1]);
assert!(matches!(
fw.finish().unwrap_err(),
FormatError::LibverTooOldForContent { .. }
));
let mut fw = FileWriter::new();
fw.with_libver_bounds(LibVer::Earliest, LibVer::V18);
fw.with_file_space_page_size(4096);
fw.create_dataset("d").with_i32_data(&[1]);
assert!(matches!(
fw.finish().unwrap_err(),
FormatError::LibverTooOldForContent { .. }
));
let mut fw = FileWriter::new();
fw.with_libver_bounds(LibVer::Earliest, LibVer::V18);
fw.create_dataset("d")
.with_i32_data(&[1, 2, 3])
.with_shape(&[3])
.with_maxshape(&[u64::MAX]);
assert!(matches!(
fw.finish().unwrap_err(),
FormatError::LibverTooOldForContent { .. }
));
}
}