use crate::format::bytes::read_le_uint as read_uint;
use crate::format::fractal_heap::{
collect_managed_blocks, read_heap_object, FractalHeapHeader, HeapId,
};
use crate::format::messages::attribute::{ATTR_FLAG_SPACE_SHARED, ATTR_FLAG_TYPE_SHARED};
use crate::format::messages::datatype::{DatatypeMessage, DatatypeNodeVersion};
use crate::format::messages::shared::{MessageStorage, MSG_FLAG_SHARED};
use crate::format::messages::MSG_FLAG_SHAREABLE;
use crate::format::messages::{
MSG_ATTRIBUTE, MSG_DATASPACE, MSG_DATATYPE, MSG_LINK, MSG_LINK_INFO, MSG_NULL,
MSG_OBJ_HEADER_CONTINUATION, MSG_SYMBOL_TABLE,
};
use crate::format::object_header::{ObjectHeader, ObjectHeaderMessage};
use crate::format::sohm::{SharedLocation, SharedMessagePointer};
use crate::format::BlockReader;
use crate::format::{FormatError, UNDEF_ADDR};
use crate::io::file_handle::FileHandle;
use crate::io::{FileMeta, IoResult};
const MAX_CONT_BLOCKS: usize = 4096;
const HEADER_PROBE: usize = 8192;
const MAX_SHARED_DEPTH: usize = 8;
pub(crate) fn read_object_header_full(
handle: &mut FileHandle,
meta: &FileMeta,
addr: u64,
) -> IoResult<ObjectHeader> {
read_object_header_at(handle, meta, addr, 0)
}
pub(crate) type HeaderBlocks = Vec<(u64, u64)>;
pub(crate) fn read_object_header_with_blocks(
handle: &mut FileHandle,
meta: &FileMeta,
addr: u64,
) -> IoResult<(ObjectHeader, HeaderBlocks)> {
let (mut header, blocks) = read_header_chain(handle, meta, addr)?;
resolve_shared_messages(handle, meta, &mut header, addr, 0)?;
Ok((header, blocks))
}
pub(crate) fn read_header_message_storage(
handle: &mut FileHandle,
meta: &FileMeta,
addr: u64,
) -> IoResult<Vec<(u8, MessageStorage)>> {
let (header, _) = read_header_chain(handle, meta, addr)?;
let mut out = Vec::new();
for msg in &header.messages {
let storage = if msg.flags & MSG_FLAG_SHARED != 0 {
MessageStorage::Shared(SharedMessagePointer::decode(&msg.data, &meta.ctx)?.location)
} else if msg.flags & MSG_FLAG_SHAREABLE != 0 {
MessageStorage::Shareable
} else {
continue;
};
out.push((msg.msg_type, storage));
}
Ok(out)
}
pub(crate) fn read_header_message_flags(
handle: &mut FileHandle,
meta: &FileMeta,
addr: u64,
) -> IoResult<Vec<(u8, u8)>> {
let (header, _) = read_header_chain(handle, meta, addr)?;
Ok(header
.messages
.iter()
.filter(|m| !matches!(m.msg_type, MSG_NULL | MSG_OBJ_HEADER_CONTINUATION))
.map(|m| (m.msg_type, m.flags))
.collect())
}
pub(crate) fn read_header_datatype_versions(
handle: &mut FileHandle,
meta: &FileMeta,
addr: u64,
) -> IoResult<Vec<DatatypeNodeVersion>> {
let header = read_object_header_full(handle, meta, addr)?;
let msg = header
.messages
.iter()
.find(|m| m.msg_type == MSG_DATATYPE)
.ok_or_else(|| {
crate::io::IoError::NotFound(format!(
"object header at {addr:#x} holds no datatype message"
))
})?;
Ok(DatatypeMessage::decode_versions(&msg.data)?)
}
pub(crate) fn read_header_recorded_times(
handle: &mut FileHandle,
meta: &FileMeta,
addr: u64,
) -> IoResult<Option<crate::format::object_header::ObjectTimes>> {
let (header, _) = read_header_chain(handle, meta, addr)?;
Ok(header.recorded_times())
}
#[derive(Clone)]
pub(crate) struct ExtensionMessage {
pub msg_type: u8,
pub flags: u8,
pub body: Vec<u8>,
}
pub(crate) fn superblock_extension_messages(
handle: &mut FileHandle,
meta: &FileMeta,
addr: u64,
) -> IoResult<(Vec<ExtensionMessage>, HeaderBlocks)> {
use crate::format::messages::MSG_SHARED_MESSAGE_TABLE;
let (header, blocks) = read_header_chain(handle, meta, addr)?;
let carried = header
.messages
.iter()
.filter(|m| {
!matches!(
m.msg_type,
MSG_NULL | MSG_OBJ_HEADER_CONTINUATION | MSG_SHARED_MESSAGE_TABLE
)
})
.map(|m| ExtensionMessage {
msg_type: m.msg_type,
flags: m.flags,
body: m.data.clone(),
})
.collect();
Ok((carried, blocks))
}
fn read_object_header_at(
handle: &mut FileHandle,
meta: &FileMeta,
addr: u64,
depth: usize,
) -> IoResult<ObjectHeader> {
let (mut header, _) = read_header_chain(handle, meta, addr)?;
resolve_shared_messages(handle, meta, &mut header, addr, depth)?;
Ok(header)
}
fn read_header_chain(
handle: &mut FileHandle,
meta: &FileMeta,
addr: u64,
) -> IoResult<(ObjectHeader, HeaderBlocks)> {
let ctx = &meta.ctx;
let mut buf = handle.read_at_most(addr, HEADER_PROBE)?;
if let Err(FormatError::BufferTooShort { needed, .. }) = ObjectHeader::decode_any(&buf) {
if needed > buf.len() {
buf = handle.read_at_most(addr, needed)?;
}
}
let (mut header, chunk0_len) = ObjectHeader::decode_any(&buf)?;
let mut blocks = vec![(addr, chunk0_len as u64)];
let is_v2 = buf.len() >= 4 && buf[0..4] == crate::format::object_header::OHDR_SIGNATURE;
let track_creation_order = is_v2 && (header.flags & 0x04) != 0;
let sa = ctx.sizeof_addr as usize;
let ss = ctx.sizeof_size as usize;
let collect = |msgs: &[ObjectHeaderMessage], out: &mut HeaderBlocks| {
for msg in msgs {
if msg.msg_type == MSG_OBJ_HEADER_CONTINUATION && msg.data.len() >= sa + ss {
let cont_addr = read_uint(&msg.data, sa);
let cont_len = read_uint(&msg.data[sa..], ss);
out.push((cont_addr, cont_len));
}
}
};
let mut pending: HeaderBlocks = Vec::new();
collect(&header.messages, &mut pending);
let mut visited = std::collections::HashSet::new();
let mut blocks_read = 0usize;
while let Some((cont_addr, cont_len)) = pending.pop() {
if cont_addr == UNDEF_ADDR || cont_addr == 0 || cont_len == 0 {
continue;
}
if !visited.insert(cont_addr) {
continue; }
blocks_read += 1;
if blocks_read > MAX_CONT_BLOCKS {
break;
}
let cont_buf = handle.read_at(cont_addr, cont_len as usize)?;
blocks.push((cont_addr, cont_len));
let mut new_msgs = Vec::new();
parse_continuation_block(&cont_buf, is_v2, track_creation_order, &mut new_msgs)?;
collect(&new_msgs, &mut pending);
header.messages.extend(new_msgs);
}
Ok((header, blocks))
}
pub(crate) fn blocks_shared_message_rebuild(
handle: &mut FileHandle,
meta: &FileMeta,
addr: u64,
) -> IoResult<Option<&'static str>> {
let (header, _) = read_header_chain(handle, meta, addr)?;
let in_heap = |bytes: &[u8]| {
SharedMessagePointer::decode(bytes, &meta.ctx)
.is_ok_and(|p| p.location == SharedLocation::Sohm)
};
for msg in &header.messages {
if msg.flags & MSG_FLAG_SHARED != 0 && in_heap(&msg.data) {
return Ok(Some("holds a shared object header message"));
}
if msg.msg_type == MSG_ATTRIBUTE
&& shared_attribute_fields(&msg.data)
.into_iter()
.flatten()
.any(in_heap)
{
return Ok(Some(
"holds an attribute whose datatype or dataspace is a shared object header \
message",
));
}
if matches!(msg.msg_type, MSG_LINK | MSG_LINK_INFO | MSG_SYMBOL_TABLE) {
return Ok(Some(
"names objects of its own, which would keep their bytes with it",
));
}
}
Ok(None)
}
pub(crate) fn committed_datatype_address(
handle: &mut FileHandle,
meta: &FileMeta,
addr: u64,
) -> IoResult<Option<u64>> {
let (header, _) = read_header_chain(handle, meta, addr)?;
for msg in &header.messages {
if msg.msg_type != MSG_DATATYPE || msg.flags & MSG_FLAG_SHARED == 0 {
continue;
}
let ptr = SharedMessagePointer::decode(&msg.data, &meta.ctx)?;
if ptr.location == SharedLocation::Committed {
return Ok(Some(ptr.oh_addr));
}
}
Ok(None)
}
fn shared_attribute_fields(body: &[u8]) -> [Option<&[u8]>; 2] {
let fields = || {
if body.len() < 8 || body[0] < 2 {
return None;
}
let flags = body[1];
if flags & (ATTR_FLAG_TYPE_SHARED | ATTR_FLAG_SPACE_SHARED) == 0 {
return None;
}
let name_size = u16::from_le_bytes([body[2], body[3]]) as usize;
let dt_size = u16::from_le_bytes([body[4], body[5]]) as usize;
let ds_size = u16::from_le_bytes([body[6], body[7]]) as usize;
let hdr_len: usize = if body[0] >= 3 { 9 } else { 8 };
let name_end = hdr_len.checked_add(name_size)?;
let dt_end = name_end.checked_add(dt_size)?;
let ds_end = dt_end.checked_add(ds_size)?;
if body.len() < ds_end {
return None;
}
Some([
(flags & ATTR_FLAG_TYPE_SHARED != 0).then(|| &body[name_end..dt_end]),
(flags & ATTR_FLAG_SPACE_SHARED != 0).then(|| &body[dt_end..ds_end]),
])
};
fields().unwrap_or([None, None])
}
fn resolve_shared_messages(
handle: &mut FileHandle,
meta: &FileMeta,
header: &mut ObjectHeader,
self_addr: u64,
depth: usize,
) -> IoResult<()> {
let mut resolved: Vec<(usize, Vec<u8>)> = Vec::new();
for (i, msg) in header.messages.iter().enumerate() {
if msg.flags & MSG_FLAG_SHARED == 0 {
continue;
}
let body = resolve_shared_body(
handle,
meta,
msg.msg_type,
&msg.data,
&header.messages,
self_addr,
depth,
)?;
resolved.push((i, body));
}
for (i, body) in resolved {
header.messages[i].data = body;
header.messages[i].flags &= !MSG_FLAG_SHARED;
}
let mut rewritten: Vec<(usize, Vec<u8>)> = Vec::new();
for (i, msg) in header.messages.iter().enumerate() {
if msg.msg_type != MSG_ATTRIBUTE {
continue;
}
if let Some(body) = resolve_shared_attribute_fields(
handle,
meta,
&msg.data,
&header.messages,
self_addr,
depth,
)? {
rewritten.push((i, body));
}
}
for (i, body) in rewritten {
header.messages[i].data = body;
}
Ok(())
}
fn resolve_shared_body(
handle: &mut FileHandle,
meta: &FileMeta,
msg_type: u8,
ptr_bytes: &[u8],
siblings: &[ObjectHeaderMessage],
self_addr: u64,
depth: usize,
) -> IoResult<Vec<u8>> {
let invalid = |msg: String| crate::io::IoError::Format(FormatError::InvalidData(msg));
if depth >= MAX_SHARED_DEPTH {
return Err(invalid(format!(
"shared message indirection deeper than {MAX_SHARED_DEPTH} levels"
)));
}
let ptr = SharedMessagePointer::decode(ptr_bytes, &meta.ctx)?;
match ptr.location {
SharedLocation::Sohm => {
let table = meta.sohm.as_ref().ok_or_else(|| {
invalid(format!(
"message type {msg_type:#04x} is shared in the heap but the file has no \
shared message table"
))
})?;
let heap_addr = table.heap_addr(msg_type).ok_or_else(|| {
invalid(format!(
"no shared-message index covers message type {msg_type:#04x}"
))
})?;
let hdr_buf = handle.read_at_most(heap_addr, 512)?;
let fh_header = FractalHeapHeader::decode(&hdr_buf, &meta.ctx)?;
let mut br = HandleBlockReader { handle };
let id = HeapId::parse(&ptr.heap_id, &fh_header, &meta.ctx)?;
let blocks = collect_managed_blocks(&fh_header, &meta.ctx, &mut br)?;
Ok(read_heap_object(
&id, &fh_header, &meta.ctx, &blocks, &mut br,
)?)
}
SharedLocation::Committed => {
let pick = |msgs: &[ObjectHeaderMessage]| {
msgs.iter()
.find(|m| m.msg_type == msg_type && m.flags & MSG_FLAG_SHARED == 0)
.map(|m| m.data.clone())
};
let found = if ptr.oh_addr == self_addr {
pick(siblings)
} else {
let target = read_object_header_at(handle, meta, ptr.oh_addr, depth + 1)?;
pick(&target.messages)
};
found.ok_or_else(|| {
invalid(format!(
"object header at {:#x} holds no message of type {msg_type:#04x} for a \
committed shared message",
ptr.oh_addr
))
})
}
SharedLocation::Unshared | SharedLocation::Here => Err(invalid(format!(
"message type {msg_type:#04x} is flagged shared but its pointer says it is not \
stored shared"
))),
}
}
fn resolve_shared_attribute_fields(
handle: &mut FileHandle,
meta: &FileMeta,
body: &[u8],
siblings: &[ObjectHeaderMessage],
self_addr: u64,
depth: usize,
) -> IoResult<Option<Vec<u8>>> {
if body.len() < 8 || body[0] < 2 {
return Ok(None);
}
let flags = body[1];
if flags & (ATTR_FLAG_TYPE_SHARED | ATTR_FLAG_SPACE_SHARED) == 0 {
return Ok(None);
}
let name_size = u16::from_le_bytes([body[2], body[3]]) as usize;
let dt_size = u16::from_le_bytes([body[4], body[5]]) as usize;
let ds_size = u16::from_le_bytes([body[6], body[7]]) as usize;
let hdr_len = if body[0] >= 3 { 9 } else { 8 };
let name_end = hdr_len + name_size;
let dt_end = name_end + dt_size;
let ds_end = dt_end + ds_size;
if body.len() < ds_end {
return Err(crate::io::IoError::Format(FormatError::BufferTooShort {
needed: ds_end,
available: body.len(),
}));
}
let resolve = |handle: &mut FileHandle, msg_type: u8, range: std::ops::Range<usize>| {
resolve_shared_body(
handle,
meta,
msg_type,
&body[range],
siblings,
self_addr,
depth,
)
};
let datatype = if flags & ATTR_FLAG_TYPE_SHARED != 0 {
resolve(handle, MSG_DATATYPE, name_end..dt_end)?
} else {
body[name_end..dt_end].to_vec()
};
let dataspace = if flags & ATTR_FLAG_SPACE_SHARED != 0 {
resolve(handle, MSG_DATASPACE, dt_end..ds_end)?
} else {
body[dt_end..ds_end].to_vec()
};
let (Ok(dt_len), Ok(ds_len)) = (
u16::try_from(datatype.len()),
u16::try_from(dataspace.len()),
) else {
return Err(crate::io::IoError::Format(FormatError::InvalidData(
"shared attribute datatype/dataspace does not fit an attribute message".into(),
)));
};
let mut out = Vec::with_capacity(hdr_len + name_size + datatype.len() + dataspace.len());
out.extend_from_slice(&body[..hdr_len]);
out[1] = 0;
out[4..6].copy_from_slice(&dt_len.to_le_bytes());
out[6..8].copy_from_slice(&ds_len.to_le_bytes());
out.extend_from_slice(&body[hdr_len..name_end]);
out.extend_from_slice(&datatype);
out.extend_from_slice(&dataspace);
out.extend_from_slice(&body[ds_end..]);
Ok(Some(out))
}
fn parse_continuation_block(
cont_buf: &[u8],
is_v2: bool,
track_creation_order: bool,
out: &mut Vec<ObjectHeaderMessage>,
) -> crate::format::FormatResult<()> {
if is_v2 {
if cont_buf.len() < 8 {
return Err(FormatError::BufferTooShort {
needed: 8,
available: cont_buf.len(),
});
}
if cont_buf[0..4] != *b"OCHK" {
return Err(FormatError::InvalidSignature);
}
let msgs_end = cont_buf.len() - 4; let stored = u32::from_le_bytes([
cont_buf[msgs_end],
cont_buf[msgs_end + 1],
cont_buf[msgs_end + 2],
cont_buf[msgs_end + 3],
]);
let computed = crate::format::checksum::checksum_metadata(&cont_buf[..msgs_end]);
if stored != computed {
return Err(FormatError::ChecksumMismatch {
expected: stored,
computed,
});
}
let mut pos = 4; let hdr_size = if track_creation_order { 6 } else { 4 };
while pos + hdr_size <= msgs_end {
let msg_type = cont_buf[pos];
let data_size = u16::from_le_bytes([cont_buf[pos + 1], cont_buf[pos + 2]]) as usize;
let msg_flags = cont_buf[pos + 3];
let creation_index = if track_creation_order {
u16::from_le_bytes([cont_buf[pos + 4], cont_buf[pos + 5]])
} else {
0
};
pos += hdr_size;
if pos + data_size > msgs_end {
break;
}
if msg_type != 0 {
out.push(ObjectHeaderMessage {
msg_type,
flags: msg_flags,
creation_index,
data: cont_buf[pos..pos + data_size].to_vec(),
});
}
pos += data_size;
}
} else {
let mut pos = 0;
while pos + 8 <= cont_buf.len() {
let msg_type = u16::from_le_bytes([cont_buf[pos], cont_buf[pos + 1]]);
let data_size = u16::from_le_bytes([cont_buf[pos + 2], cont_buf[pos + 3]]) as usize;
let msg_flags = cont_buf[pos + 4];
pos += 8; if pos + data_size > cont_buf.len() {
break;
}
if msg_type != 0 {
out.push(ObjectHeaderMessage {
msg_type: msg_type as u8,
flags: msg_flags,
creation_index: 0,
data: cont_buf[pos..pos + data_size].to_vec(),
});
}
pos += data_size;
pos = (pos + 7) & !7; }
}
Ok(())
}
pub(crate) fn read_datatype_message(
handle: &mut FileHandle,
meta: &FileMeta,
msg: &ObjectHeaderMessage,
) -> IoResult<DatatypeMessage> {
if msg.flags & MSG_FLAG_SHARED != 0 {
let body = resolve_shared_body(handle, meta, msg.msg_type, &msg.data, &[], UNDEF_ADDR, 0)?;
let (dt, _) = DatatypeMessage::decode(&body, &meta.ctx)?;
return Ok(dt);
}
let (dt, _) = DatatypeMessage::decode(&msg.data, &meta.ctx)?;
Ok(dt)
}
struct HandleBlockReader<'a> {
handle: &'a mut FileHandle,
}
impl BlockReader for HandleBlockReader<'_> {
fn read_block(&mut self, offset: u64, len: usize) -> crate::format::FormatResult<Vec<u8>> {
self.handle.read_at(offset, len).map_err(|e| {
FormatError::InvalidData(format!(
"fractal heap block read failed at {offset:#x}: {e}"
))
})
}
}