use std::path::Path;
use std::sync::atomic::{AtomicU64, Ordering};
use std::sync::{Arc, RwLock};
#[cfg(feature = "parallel")]
use rayon::prelude::*;
use crate::blocks::{
CaBlock, CaStorage, CcBlock, CgBlock, ChElement, ChannelType, CnBlock, CompressionType,
Conversion, DataType, DgBlock, DlBlock, DzBlock, HdBlock, HlBlock, LdBlock, ParseBlock,
SiBlock, SourceInfo, TableEntry, UnfinalizedFlags, BLOCK_HEADER_SIZE,
};
use crate::cache::BlockCache;
use crate::channels_db::{ChannelLocation, ChannelsDB, MastersDB, SearchMode};
use crate::data_index::{
CompressionInfo, DataBlockIndex, DataBlockInfo, DataBlockType, RecordIndex,
};
use crate::error::{Mf4Error, Result};
use crate::io::{ByteSource, IoBackend};
use crate::model::signal::row_major_to_stored;
use crate::model::{
ArrayElement, Attachment, Channel, ChannelGroup, ChannelHierarchyNode, DataGroup, Event,
FileHistoryEntry, FileStatistics, Metadata, MlsdLength, RecordLayout, RecordingTime,
ReductionKind, SampleReduction, Signal, UnreadableReason, VlsdPayloads,
};
use crate::parser::links::{LinkChain, MAX_COMPOSITION_DEPTH};
use crate::parser::{self, parse_hd_block, parse_id_block, Mf4Version};
#[derive(Debug, Clone)]
pub struct OpenOptions {
pub build_channels_db: bool,
pub parallel_parsing: bool,
pub parallel_threshold: usize,
pub max_alloc: usize,
pub max_decompressed: u64,
}
#[derive(Debug, Clone, Copy)]
pub(crate) struct Limits {
pub max_alloc: usize,
pub max_decompressed: u64,
}
impl From<&OpenOptions> for Limits {
fn from(options: &OpenOptions) -> Self {
Limits {
max_alloc: options.max_alloc,
max_decompressed: options.max_decompressed,
}
}
}
impl Default for Limits {
fn default() -> Self {
Limits {
max_alloc: DEFAULT_MAX_ALLOC,
max_decompressed: DEFAULT_MAX_DECOMPRESSED,
}
}
}
pub const DEFAULT_MAX_ALLOC: usize = 64 * 1024 * 1024;
pub const DEFAULT_MAX_DECOMPRESSED: u64 = 1024 * 1024 * 1024;
impl Default for OpenOptions {
fn default() -> Self {
Self {
build_channels_db: true,
parallel_parsing: true,
parallel_threshold: 100,
max_alloc: DEFAULT_MAX_ALLOC,
max_decompressed: DEFAULT_MAX_DECOMPRESSED,
}
}
}
pub struct Mf4File {
source: Arc<IoBackend>,
version: Mf4Version,
unfinalized: Option<UnfinalizedFlags>,
start_time: RecordingTime,
comment: Arc<str>,
metadata: Metadata,
data_groups: Vec<DataGroup>,
channels_db: ChannelsDB,
masters_db: MastersDB,
file_size: u64,
cache: BlockCache,
record_cache: RwLock<BoundedLru<(usize, usize), CachedRecords>>,
limits: Limits,
payload_cache: RwLock<BoundedLru<u64, Arc<VlsdPayloads>>>,
attachments: Vec<Attachment>,
events: Vec<Event>,
hierarchy: Vec<ChannelHierarchyNode>,
file_history: Vec<FileHistoryEntry>,
}
#[derive(Clone)]
struct CachedRecords {
data: Arc<Vec<u8>>,
layout: RecordLayout,
sample_count: usize,
}
const CACHE_ENTRIES: usize = 4;
struct LruEntry<K, V> {
key: K,
value: V,
size: usize,
last_used: AtomicU64,
}
struct BoundedLru<K, V> {
entries: Vec<LruEntry<K, V>>,
max_entries: usize,
max_bytes: usize,
clock: AtomicU64,
}
impl<K: PartialEq, V: Clone> BoundedLru<K, V> {
fn new(max_entries: usize, max_bytes: usize) -> Self {
BoundedLru {
entries: Vec::new(),
max_entries,
max_bytes,
clock: AtomicU64::new(0),
}
}
fn get(&self, key: &K) -> Option<V> {
let entry = self.entries.iter().find(|e| &e.key == key)?;
entry.last_used.store(
self.clock.fetch_add(1, Ordering::Relaxed),
Ordering::Relaxed,
);
Some(entry.value.clone())
}
fn insert(&mut self, key: K, value: V, size: usize) {
self.entries.retain(|e| e.key != key);
let tick = self.clock.fetch_add(1, Ordering::Relaxed);
self.entries.push(LruEntry {
key,
value,
size,
last_used: AtomicU64::new(tick),
});
let mut total: usize = self.entries.iter().map(|e| e.size).sum();
while self.entries.len() > 1
&& (self.entries.len() > self.max_entries || total > self.max_bytes)
{
let oldest = self
.entries
.iter()
.enumerate()
.min_by_key(|(_, e)| e.last_used.load(Ordering::Relaxed))
.map(|(i, _)| i)
.expect("entries is non-empty");
total -= self.entries.remove(oldest).size;
}
}
}
impl Mf4File {
pub fn open<P: AsRef<Path>>(path: P) -> Result<Self> {
Self::open_with_options(path, OpenOptions::default())
}
pub fn open_with_options<P: AsRef<Path>>(path: P, options: OpenOptions) -> Result<Self> {
let source = IoBackend::open(path)?;
Self::from_source_with_options(source, options)
}
pub fn from_bytes(bytes: Vec<u8>) -> Result<Self> {
Self::from_bytes_with_options(bytes, OpenOptions::default())
}
pub fn from_bytes_with_options(bytes: Vec<u8>, options: OpenOptions) -> Result<Self> {
let source = IoBackend::from_bytes(bytes);
Self::from_source_with_options(source, options)
}
#[cfg(feature = "mmap")]
pub fn open_mmap<P: AsRef<Path>>(path: P) -> Result<Self> {
let source = IoBackend::open_mmap(path)?;
Self::from_source_with_options(source, OpenOptions::default())
}
pub fn open_buffered<P: AsRef<Path>>(path: P) -> Result<Self> {
let source = IoBackend::open_buffered(path)?;
Self::from_source_with_options(source, OpenOptions::default())
}
fn from_source_with_options(source: IoBackend, options: OpenOptions) -> Result<Self> {
let source = Arc::new(source);
let file_size = source.len();
let id_block = parse_id_block(&source)?;
let version = Mf4Version::from_id_block(&id_block);
let is_unfinished = id_block.is_unfinished();
let unfinalized = id_block.unfinalized();
version.validate()?;
let hd_block = parse_hd_block(&source, 64)?;
let mut cache = BlockCache::new();
let start_time = RecordingTime::new(
hd_block.start_time_ns,
hd_block.tz_offset_min,
hd_block.dst_offset_min,
);
let comment = cache.get_or_parse_text(&source, hd_block.md_comment)?;
let metadata = parser::read_metadata(&source, hd_block.md_comment)?.unwrap_or_default();
let data_groups = Self::parse_data_groups(
&source,
&hd_block,
&mut cache,
&options,
is_unfinished,
file_size,
)?;
let attachments = Self::parse_attachments(&source, hd_block.at_first, &mut cache)?;
let events = Self::parse_events(&source, hd_block.ev_first, &mut cache)?;
let hierarchy = Self::parse_hierarchy(
&source,
hd_block.ch_first,
&mut cache,
&mut std::collections::HashSet::new(),
)?;
let file_history = Self::parse_file_history(&source, hd_block.fh_first, &mut cache)?;
let (channels_db, masters_db) = if options.build_channels_db {
Self::build_channel_indices(&data_groups)
} else {
(ChannelsDB::new(), MastersDB::new())
};
Ok(Mf4File {
source,
version,
unfinalized,
start_time,
comment,
data_groups,
channels_db,
masters_db,
file_size,
cache,
metadata,
limits: Limits::from(&options),
record_cache: RwLock::new(BoundedLru::new(CACHE_ENTRIES, options.max_alloc)),
payload_cache: RwLock::new(BoundedLru::new(CACHE_ENTRIES, options.max_alloc)),
attachments,
events,
hierarchy,
file_history,
})
}
fn build_channel_indices(data_groups: &[DataGroup]) -> (ChannelsDB, MastersDB) {
let total_channels: usize = data_groups
.iter()
.flat_map(|dg| &dg.channel_groups)
.map(|cg| cg.channels.len())
.sum();
let mut channels_db = ChannelsDB::with_capacity(total_channels);
let mut masters_db = MastersDB::new();
for (dg_idx, dg) in data_groups.iter().enumerate() {
for (cg_idx, cg) in dg.channel_groups.iter().enumerate() {
for (ch_idx, ch) in cg.channels.iter().enumerate() {
channels_db.insert(&ch.name, ChannelLocation::new(dg_idx, cg_idx, ch_idx));
if ch.is_master() {
masters_db.insert(dg_idx, cg_idx, ch_idx);
}
}
}
}
(channels_db, masters_db)
}
fn parse_data_groups(
source: &IoBackend,
hd: &HdBlock,
cache: &mut BlockCache,
options: &OpenOptions,
is_unfinished: bool,
file_size: u64,
) -> Result<Vec<DataGroup>> {
let mut data_groups = Vec::new();
let mut dg_offset = hd.dg_first;
let mut dg_index = 0;
let mut chain = LinkChain::new();
while dg_offset != 0 {
chain.visit(dg_offset, "dg_next")?;
let dg_block = parser::parse_dg_block(source, dg_offset)?;
let channel_groups =
Self::parse_channel_groups(source, &dg_block, dg_index, cache, options)?;
let comment = cache.get_or_parse_text(source, dg_block.md_comment)?;
let data_block_index = if dg_block.data != 0 {
Self::build_data_block_index(
source,
dg_block.data,
is_unfinished,
file_size,
&channel_groups,
)?
} else {
DataBlockIndex::new()
};
let mut channel_groups = channel_groups;
let record_index = Self::index_records(
source,
&dg_block,
&data_block_index,
&mut channel_groups,
Limits::from(options),
)?;
for cg in &mut channel_groups {
for ch in &mut cg.channels {
ch.sample_count = cg.sample_count;
}
}
let data_group = DataGroup {
id: dg_index,
index: dg_index,
channel_groups,
comment: comment.to_string(),
dg_offset,
data_offset: dg_block.data,
rec_id_size: dg_block.rec_id_size,
data_block_index,
record_index,
};
data_groups.push(data_group);
dg_offset = dg_block.dg_next;
dg_index += 1;
}
for dg_idx in 0..data_groups.len() {
for cg_idx in 0..data_groups[dg_idx].channel_groups.len() {
for ch_idx in 0..data_groups[dg_idx].channel_groups[cg_idx].channels.len() {
let ch = &data_groups[dg_idx].channel_groups[cg_idx].channels[ch_idx];
let mut resolved_sc = None;
if let Some(ref elem) = ch.array_element {
if elem.storage == CaStorage::CgTemplate {
if let Some(&first_link) = elem.group_links.first() {
for dg in &data_groups {
for cg in &dg.channel_groups {
if cg.cg_offset == first_link {
resolved_sc = Some(cg.sample_count);
break;
}
}
if resolved_sc.is_some() {
break;
}
}
}
} else if elem.storage == CaStorage::DgTemplate {
if let Some(&first_link) = elem.group_links.first() {
for dg in &data_groups {
if dg.dg_offset == first_link {
if let Some(cg) = dg.channel_groups.first() {
resolved_sc = Some(cg.sample_count);
}
break;
}
}
}
}
}
if let Some(sc) = resolved_sc {
data_groups[dg_idx].channel_groups[cg_idx].channels[ch_idx].sample_count =
sc;
}
}
}
}
Ok(data_groups)
}
fn index_records(
source: &IoBackend,
dg: &DgBlock,
data_index: &DataBlockIndex,
channel_groups: &mut [ChannelGroup],
limits: Limits,
) -> Result<Option<RecordIndex>> {
let rec_id_size = dg.rec_id_size;
let unsorted = rec_id_size > 0 && channel_groups.len() > 1;
if data_index.is_empty() {
return Ok(unsorted.then(|| RecordIndex::with_groups(channel_groups.len())));
}
if !unsorted {
let data_size = data_index.total_size();
for cg in channel_groups.iter_mut() {
let record_size = cg.record_size(rec_id_size);
if record_size == 0 {
cg.sample_count = 0;
continue;
}
let capacity = data_size / record_size as u64;
if cg.sample_count == 0 || cg.sample_count > capacity {
cg.sample_count = capacity;
}
}
return Ok(None);
}
let mut record_map: std::collections::HashMap<u64, (usize, usize, bool)> =
std::collections::HashMap::with_capacity(channel_groups.len());
for (idx, cg) in channel_groups.iter().enumerate() {
record_map.insert(cg.record_id, (idx, cg.record_size(rec_id_size), cg.is_vlsd));
}
let mut index = RecordIndex::with_groups(channel_groups.len());
let mut base: u64 = 0;
for (_offset, block_info) in data_index.iter() {
let block_data = Self::read_block_payload(source, block_info, limits)?;
let data: &[u8] = &block_data;
let mut pos: usize = 0;
while pos < data.len() {
let Some(rec_id) = read_record_id(data, pos, rec_id_size) else {
break;
};
let Some(&(cg_idx, record_size, is_vlsd)) = record_map.get(&rec_id) else {
break;
};
let next = if is_vlsd {
let len_at = pos + rec_id_size as usize;
if len_at + 4 > data.len() {
break;
}
let vlsd_len = u32::from_le_bytes([
data[len_at],
data[len_at + 1],
data[len_at + 2],
data[len_at + 3],
]) as usize;
len_at + 4 + vlsd_len
} else {
if record_size == 0 || pos + record_size > data.len() {
break;
}
pos + record_size
};
index.push(cg_idx, base + pos as u64);
pos = next;
}
base += block_data.len() as u64;
}
for (idx, cg) in channel_groups.iter_mut().enumerate() {
cg.sample_count = index.count(idx) as u64;
}
Ok(Some(index))
}
fn read_block_payload(
source: &IoBackend,
info: &DataBlockInfo,
limits: Limits,
) -> Result<Vec<u8>> {
let read_single = |inf: &DataBlockInfo| -> Result<Vec<u8>> {
if let Some(compression) = &inf.compression {
let compressed =
source.read_bytes(compression.data_offset, inf.compressed_size as usize)?;
Self::decompress(&compressed, compression, inf.original_size as usize, limits)
} else {
let at = inf.offset + BLOCK_HEADER_SIZE as u64;
Ok(source
.read_bytes(at, inf.original_size as usize)?
.into_owned())
}
};
if let Some(inval_info) = &info.invalidation_block {
let data_buf = read_single(info)?;
let inval_buf = read_single(inval_info)?;
let data_bytes = info.data_bytes as usize;
let inval_bytes = info.inval_bytes as usize;
if data_bytes > 0 && inval_bytes > 0 {
let cycles = (data_buf.len() / data_bytes).min(inval_buf.len() / inval_bytes);
let mut out = Vec::with_capacity(cycles * (data_bytes + inval_bytes));
for i in 0..cycles {
let d_start = i * data_bytes;
out.extend_from_slice(&data_buf[d_start..d_start + data_bytes]);
let i_start = i * inval_bytes;
out.extend_from_slice(&inval_buf[i_start..i_start + inval_bytes]);
}
Ok(out)
} else {
Err(Mf4Error::unsupported(
"linked-data invalidation",
"the block carries invalidation bytes, but the record layout \
needed to interleave them is not available on this path",
))
}
} else {
read_single(info)
}
}
pub(crate) fn build_data_block_index_at(&self, offset: u64) -> Result<DataBlockIndex> {
Self::build_data_block_index(&self.source, offset, false, self.file_size, &[])
}
fn build_data_block_index(
source: &IoBackend,
offset: u64,
is_unfinished: bool,
file_size: u64,
channel_groups: &[ChannelGroup],
) -> Result<DataBlockIndex> {
if offset == 0 {
return Ok(DataBlockIndex::new());
}
let header = parser::parse_block_header(source, offset)?;
let block_id = source.read_bytes(offset, 4)?;
match &block_id[..] {
b"##DT" | b"##SD" | b"##RD" | b"##DV" | b"##DI" => {
let mut index = DataBlockIndex::new();
let mut data_size = header.length.saturating_sub(BLOCK_HEADER_SIZE as u64);
if is_unfinished && data_size == 0 {
let data_start = offset + BLOCK_HEADER_SIZE as u64;
data_size = file_size.saturating_sub(data_start);
}
let block_type = match &block_id[..] {
b"##SD" => DataBlockType::SortedData,
b"##DV" => DataBlockType::DataValues,
b"##DI" => DataBlockType::DataInvalidation,
_ => DataBlockType::Data,
};
index.push(DataBlockInfo::uncompressed(offset, block_type, data_size));
Ok(index)
}
b"##DZ" => {
let dz_data = source.read_bytes(offset, header.length as usize)?;
let dz = DzBlock::parse(&dz_data, offset)?;
let mut index = DataBlockIndex::new();
let compression = CompressionInfo {
algorithm: dz.zip_type,
parameter: dz.zip_parameter,
data_offset: dz.compressed_data_offset,
};
index.push(DataBlockInfo::compressed(
offset,
dz.original_size,
dz.compressed_size,
compression,
));
Ok(index)
}
b"##DL" => Self::build_data_list_index(source, offset),
b"##LD" => Self::build_list_data_index(source, offset, channel_groups),
b"##HL" => {
let hl_data = source.read_bytes(offset, header.length as usize)?;
let hl = HlBlock::parse(&hl_data, offset)?;
if hl.dl_first != 0 {
let target_id = source.read_bytes(hl.dl_first, 4)?;
if &target_id[..] == b"##LD" {
Self::build_list_data_index(source, hl.dl_first, channel_groups)
} else {
Self::build_data_list_index(source, hl.dl_first)
}
} else {
Ok(DataBlockIndex::new())
}
}
other => Err(Mf4Error::unsupported(
"data block",
format!(
"block '{}' at offset {offset} is not a data block this build can read",
String::from_utf8_lossy(other)
),
)),
}
}
fn build_data_list_index(source: &IoBackend, dl_offset: u64) -> Result<DataBlockIndex> {
let mut index = DataBlockIndex::new();
let mut current_dl = dl_offset;
let mut chain = LinkChain::new();
while current_dl != 0 {
chain.visit(current_dl, "dl_next")?;
let header = parser::parse_block_header(source, current_dl)?;
let dl_data = source.read_bytes(current_dl, header.length as usize)?;
let dl = DlBlock::parse(&dl_data, current_dl)?;
for &data_link in &dl.data_links {
if data_link == 0 {
continue;
}
let block_header = parser::parse_block_header(source, data_link)?;
let block_id = source.read_bytes(data_link, 4)?;
match &block_id[..] {
b"##DT" | b"##SD" | b"##RD" | b"##DV" | b"##DI" => {
let data_size =
block_header.length.saturating_sub(BLOCK_HEADER_SIZE as u64);
let block_type = match &block_id[..] {
b"##SD" => DataBlockType::SortedData,
b"##DV" => DataBlockType::DataValues,
b"##DI" => DataBlockType::DataInvalidation,
_ => DataBlockType::Data,
};
index.push(DataBlockInfo::uncompressed(
data_link, block_type, data_size,
));
}
b"##DZ" => {
let dz_data = source.read_bytes(data_link, block_header.length as usize)?;
let dz = DzBlock::parse(&dz_data, data_link)?;
let compression = CompressionInfo {
algorithm: dz.zip_type,
parameter: dz.zip_parameter,
data_offset: dz.compressed_data_offset,
};
index.push(DataBlockInfo::compressed(
data_link,
dz.original_size,
dz.compressed_size,
compression,
));
}
other => {
return Err(Mf4Error::unsupported(
"data block",
format!(
"data list at offset {current_dl} refers to block '{}' at \
offset {data_link}, which this build cannot read",
String::from_utf8_lossy(other)
),
));
}
}
}
current_dl = dl.dl_next;
}
Ok(index)
}
fn build_list_data_index(
source: &IoBackend,
ld_offset: u64,
channel_groups: &[ChannelGroup],
) -> Result<DataBlockIndex> {
let mut index = DataBlockIndex::new();
let mut current_ld = ld_offset;
let mut chain = LinkChain::new();
let (data_bytes, inval_bytes) = if let Some(cg) = channel_groups.first() {
(cg.data_bytes, cg.inval_bytes)
} else {
(0, 0)
};
let layouts_agree = channel_groups
.windows(2)
.all(|w| w[0].data_bytes == w[1].data_bytes && w[0].inval_bytes == w[1].inval_bytes);
let mut layout_checked = false;
while current_ld != 0 {
chain.visit(current_ld, "ld_next")?;
let header = parser::parse_block_header(source, current_ld)?;
let ld_data = source.read_bytes(current_ld, header.length as usize)?;
let ld = LdBlock::parse(&ld_data, current_ld)?;
let mut previous_size: Option<u64> = None;
for (i, &data_link) in ld.data_links.iter().enumerate() {
if data_link == 0 {
previous_size = None;
continue;
}
let block_header = parser::parse_block_header(source, data_link)?;
let block_id = source.read_bytes(data_link, 4)?;
let mut data_info = match &block_id[..] {
b"##DT" | b"##SD" | b"##RD" | b"##DV" => {
let data_size =
block_header.length.saturating_sub(BLOCK_HEADER_SIZE as u64);
let block_type = match &block_id[..] {
b"##SD" => DataBlockType::SortedData,
b"##DV" => DataBlockType::DataValues,
_ => DataBlockType::Data,
};
DataBlockInfo::uncompressed(data_link, block_type, data_size)
}
b"##DZ" => {
let dz_data = source.read_bytes(data_link, block_header.length as usize)?;
let dz = DzBlock::parse(&dz_data, data_link)?;
let compression = CompressionInfo {
algorithm: dz.zip_type,
parameter: dz.zip_parameter,
data_offset: dz.compressed_data_offset,
};
DataBlockInfo::compressed(
data_link,
dz.original_size,
dz.compressed_size,
compression,
)
}
other => {
return Err(Mf4Error::unsupported(
"data block",
format!(
"list data at offset {current_ld} refers to block '{}' at \
offset {data_link}, which this build cannot read",
String::from_utf8_lossy(other)
),
));
}
};
if i > 0 {
if let (None, Some(&offset), Some(&earlier), Some(prev_size)) = (
ld.equal_length,
ld.offsets.get(i),
ld.offsets.get(i - 1),
previous_size,
) {
let delta = offset.checked_sub(earlier).ok_or_else(|| {
Mf4Error::unsupported(
"list data",
format!(
"list data at offset {current_ld} declares offset {offset} \
for data block {i} before {earlier} for block {}; the \
offsets run backwards, which no contiguous stream can do",
i - 1
),
)
})?;
if delta != prev_size {
return Err(Mf4Error::unsupported(
"list data",
format!(
"list data at offset {current_ld} declares data block {} \
starts {delta} bytes after block {}, but the previous \
block holds {prev_size} bytes; the declared offsets \
contradict the stream this build reads",
i,
i - 1
),
));
}
}
}
if i < ld.invalidation_links.len() {
let inval_link = ld.invalidation_links[i];
if inval_link != 0 {
if !layout_checked {
layout_checked = true;
if !layouts_agree {
return Err(Mf4Error::unsupported(
"list data",
"the channel groups of this data group disagree on the \
record layout, and the invalidation blocks in this LD \
chain can only be interleaved with a single split of \
data and invalidation bytes; refusing to guess which \
layout applies to which block",
));
}
}
let inval_header = parser::parse_block_header(source, inval_link)?;
let inval_id = source.read_bytes(inval_link, 4)?;
let inval_info = match &inval_id[..] {
b"##DI" | b"##DT" | b"##DV" => {
let inval_size =
inval_header.length.saturating_sub(BLOCK_HEADER_SIZE as u64);
DataBlockInfo::uncompressed(
inval_link,
DataBlockType::DataInvalidation,
inval_size,
)
}
b"##DZ" => {
let dz_data =
source.read_bytes(inval_link, inval_header.length as usize)?;
let dz = DzBlock::parse(&dz_data, inval_link)?;
let compression = CompressionInfo {
algorithm: dz.zip_type,
parameter: dz.zip_parameter,
data_offset: dz.compressed_data_offset,
};
DataBlockInfo::compressed(
inval_link,
dz.original_size,
dz.compressed_size,
compression,
)
}
other => {
return Err(Mf4Error::unsupported(
"invalidation block",
format!(
"list data at offset {current_ld} refers to invalidation block '{}' at \
offset {inval_link}, which this build cannot read",
String::from_utf8_lossy(other)
),
));
}
};
data_info.invalidation_block = Some(Box::new(inval_info));
data_info.data_bytes = data_bytes;
data_info.inval_bytes = inval_bytes;
}
}
previous_size = Some(data_info.original_size);
index.push(data_info);
}
current_ld = ld.ld_next;
}
Ok(index)
}
fn parse_channel_groups(
source: &IoBackend,
dg: &DgBlock,
dg_index: usize,
cache: &mut BlockCache,
options: &OpenOptions,
) -> Result<Vec<ChannelGroup>> {
let mut channel_groups = Vec::new();
let mut cg_offset = dg.cg_first;
let mut cg_index = 0;
let mut chain = LinkChain::new();
while cg_offset != 0 {
chain.visit(cg_offset, "cg_next")?;
let cg_block = parser::parse_cg_block(source, cg_offset)?;
let channels =
Self::parse_channels(source, &cg_block, dg_index, cg_index, cache, options)?;
let acquisition_name = cache.get_or_parse_text(source, cg_block.tx_acq_name)?;
let comment = cache.get_or_parse_text(source, cg_block.md_comment)?;
let source_info =
if let Some(si_arc) = cache.get_or_parse_si(source, cg_block.si_acq_source)? {
Some(build_source_info(source, &si_arc, cache)?)
} else {
None
};
let sample_reductions = Self::parse_sample_reductions(source, cg_block.sr_first)?;
let channel_group = ChannelGroup {
id: cg_index,
index: cg_index,
data_group_index: dg_index,
acquisition_name: acquisition_name.to_string(),
sample_count: cg_block.cycle_count,
channels,
source: source_info,
comment: comment.to_string(),
record_id: cg_block.record_id,
data_bytes: cg_block.data_bytes,
inval_bytes: cg_block.inval_bytes,
cg_offset,
is_vlsd: cg_block.flags.vlsd,
bus_event: cg_block.flags.bus_event,
plain_bus_event: cg_block.flags.plain_bus_event,
sample_reductions,
};
channel_groups.push(channel_group);
cg_offset = cg_block.cg_next;
cg_index += 1;
}
Ok(channel_groups)
}
fn parse_channels(
source: &IoBackend,
cg: &CgBlock,
dg_index: usize,
cg_index: usize,
cache: &mut BlockCache,
options: &OpenOptions,
) -> Result<Vec<Channel>> {
let mut cn_offsets = Vec::new();
let mut cn_offset = cg.cn_first;
let mut chain = LinkChain::new();
while cn_offset != 0 {
chain.visit(cn_offset, "cn_next")?;
cn_offsets.push(cn_offset);
let header = parser::parse_block_header(source, cn_offset)?;
let block_data = source.read_bytes(cn_offset, header.length.min(80) as usize)?;
if block_data.len() >= 32 {
let mut bytes = [0u8; 8];
bytes.copy_from_slice(&block_data[24..32]);
cn_offset = u64::from_le_bytes(bytes);
} else {
break;
}
}
let channel_count = cn_offsets.len();
let limits = Limits::from(options);
let mut channels =
Vec::with_capacity(channel_count.saturating_mul(2).min(limits.max_alloc));
let mut process_cn_block = |cn_block: CnBlock| -> Result<()> {
let parent_name = cache
.get_or_parse_text(source, cn_block.tx_name)?
.to_string();
let composition_offset = cn_block.composition;
let channel =
Self::build_channel(source, cn_block, dg_index, cg_index, channels.len(), cache)?;
channels.push(channel);
if composition_offset != 0 {
let parent = channels.len() - 1;
let outcome = Self::expand_composition_channels(
source,
composition_offset,
&parent_name,
dg_index,
cg_index,
&mut channels,
cache,
0,
limits.max_alloc,
)?;
if let CompositionOutcome::UnsupportedArray(reason) = outcome {
channels[parent].unreadable = Some(reason);
}
}
Ok(())
};
#[cfg(feature = "parallel")]
if options.parallel_parsing && channel_count >= options.parallel_threshold {
let parsed_blocks: Vec<Result<(CnBlock, usize)>> = cn_offsets
.par_iter()
.enumerate()
.map(|(idx, &offset)| {
let cn_block = parser::parse_cn_block(source, offset)?;
Ok((cn_block, idx))
})
.collect();
for result in parsed_blocks {
let (cn_block, _cn_index) = result?;
process_cn_block(cn_block)?;
}
} else {
for &cn_offset in cn_offsets.iter() {
let cn_block = parser::parse_cn_block(source, cn_offset)?;
process_cn_block(cn_block)?;
}
}
#[cfg(not(feature = "parallel"))]
{
for &cn_offset in cn_offsets.iter() {
let cn_block = parser::parse_cn_block(source, cn_offset)?;
process_cn_block(cn_block)?;
}
}
for (idx, ch) in channels.iter_mut().enumerate() {
ch.id = idx;
ch.index = idx;
}
Ok(channels)
}
#[allow(clippy::too_many_arguments)]
fn expand_composition_channels(
source: &IoBackend,
composition_offset: u64,
parent_name: &str,
dg_index: usize,
cg_index: usize,
channels: &mut Vec<Channel>,
cache: &mut BlockCache,
depth: usize,
max_alloc: usize,
) -> Result<CompositionOutcome> {
if depth >= MAX_COMPOSITION_DEPTH {
return Err(Mf4Error::CyclicLink {
chain: format!("cn_composition (depth limit {MAX_COMPOSITION_DEPTH})"),
offset: composition_offset,
});
}
let _header = parser::parse_block_header(source, composition_offset)?;
let block_id = source.read_bytes(composition_offset, 4)?;
match &block_id[..] {
b"##CN" => {
let mut cn_offset = composition_offset;
let mut chain = LinkChain::new();
while cn_offset != 0 {
chain.visit(cn_offset, "cn_next (composition)")?;
let cn_block = parser::parse_cn_block(source, cn_offset)?;
let child_name = cache.get_or_parse_text(source, cn_block.tx_name)?;
let next_offset = cn_block.cn_next;
let nested_composition = cn_block.composition;
let qualified_name = qualify_channel_name(parent_name, &child_name);
let channel = Self::build_channel_with_name(
source,
cn_block,
dg_index,
cg_index,
channels.len(),
cache,
qualified_name,
)?;
channels.push(channel);
if nested_composition != 0 {
let last_name = channels.last().map(|c| c.name.clone()).unwrap_or_default();
let _ = Self::expand_composition_channels(
source,
nested_composition,
&last_name,
dg_index,
cg_index,
channels,
cache,
depth + 1,
max_alloc,
)?;
}
cn_offset = next_offset;
}
}
b"##CA" => {
let ca_block = parser::parse_ca_block(source, composition_offset)?;
let stride = if ca_block.ca_byte_offset_base > 0 {
ca_block.ca_byte_offset_base as usize
} else {
0 };
if ca_block.flags.dynamic_size {
if ca_block.ca_storage != CaStorage::CnTemplate {
return Ok(CompositionOutcome::UnsupportedArray(
UnreadableReason::ArrayDynamicSize,
));
}
let single = (ca_block.ca_ndim == 1)
.then(|| ca_block.ca_dynamic_size.first().copied())
.flatten();
let Some(size_ref) = single else {
return Ok(CompositionOutcome::UnsupportedArray(
UnreadableReason::ArrayDynamicSize,
));
};
let element = if ca_block.ca_composition == 0 {
channels
.last()
.map(|p| (p.data_type, p.bit_count, p.bit_offset, 0u32))
} else {
let id = source.read_bytes(ca_block.ca_composition, 4)?;
if &id[..] != b"##CN" {
return Ok(CompositionOutcome::UnsupportedArray(
UnreadableReason::ArrayDynamicSize,
));
}
let t = parser::parse_cn_block(source, ca_block.ca_composition)?;
Some((t.data_type, t.bit_count, t.bit_offset, t.byte_offset))
};
let Some((data_type, bit_count, bit_offset, byte_offset)) = element else {
return Ok(CompositionOutcome::UnsupportedArray(
UnreadableReason::ArrayDynamicSize,
));
};
let width = (bit_count as usize).div_ceil(8).max(1);
if let Some(parent) = channels.last_mut() {
parent.array_shape = Some(ca_block.ca_dim_size.clone());
parent.array_element = Some(ArrayElement {
data_type,
bit_count,
bit_offset,
byte_offset,
inverse_layout: ca_block.flags.inverse_layout,
stride: if stride > 0 { stride } else { width },
element_offsets: None,
storage: ca_block.ca_storage,
group_links: ca_block.ca_data.clone(),
});
parent.array_dynamic_size = Some(size_ref);
parent.unreadable = None;
}
return Ok(CompositionOutcome::Expanded);
}
if ca_block.ca_composition == 0 {
if let Some(parent) = channels.last_mut() {
let width = (parent.bit_count as usize).div_ceil(8).max(1);
parent.array_shape = Some(ca_block.ca_dim_size.clone());
parent.array_element = Some(ArrayElement {
data_type: parent.data_type,
bit_count: parent.bit_count,
bit_offset: parent.bit_offset,
byte_offset: 0,
inverse_layout: ca_block.flags.inverse_layout,
stride: if stride > 0 { stride } else { width },
element_offsets: None,
storage: ca_block.ca_storage,
group_links: ca_block.ca_data.clone(),
});
parent.unreadable = None;
}
return Ok(CompositionOutcome::Expanded);
}
let id = source.read_bytes(ca_block.ca_composition, 4)?;
match &id[..] {
b"##CN" => {
let template_cn = parser::parse_cn_block(source, ca_block.ca_composition)?;
if let Some(parent) = channels.last_mut() {
parent.array_shape = Some(ca_block.ca_dim_size.clone());
let width = (template_cn.bit_count as usize).div_ceil(8).max(1);
parent.array_element = Some(ArrayElement {
data_type: template_cn.data_type,
bit_count: template_cn.bit_count,
bit_offset: template_cn.bit_offset,
byte_offset: template_cn.byte_offset,
inverse_layout: ca_block.flags.inverse_layout,
stride: if stride > 0 { stride } else { width },
element_offsets: None,
storage: ca_block.ca_storage,
group_links: ca_block.ca_data.clone(),
});
parent.unreadable = None;
}
return Ok(CompositionOutcome::Expanded);
}
b"##CA" => {
if ca_block.ca_storage != CaStorage::CnTemplate {
return Ok(CompositionOutcome::UnsupportedArray(
UnreadableReason::ArrayComposition,
));
}
let (parent_type, parent_bits, parent_bit_off) = match channels.last() {
Some(p) => (p.data_type, p.bit_count, p.bit_offset),
None => {
return Ok(CompositionOutcome::UnsupportedArray(
UnreadableReason::ArrayComposition,
))
}
};
let resolved = resolve_ca_chain(
source,
parent_type,
parent_bits,
parent_bit_off,
&ca_block,
max_alloc,
)?;
let Some(res) = resolved else {
return Ok(CompositionOutcome::UnsupportedArray(
UnreadableReason::ArrayComposition,
));
};
if let Some(parent) = channels.last_mut() {
parent.array_shape = Some(res.shape);
parent.array_element = Some(ArrayElement {
data_type: res.data_type,
bit_count: res.bit_count,
bit_offset: res.bit_offset,
byte_offset: res.byte_offset,
inverse_layout: false,
stride: 1,
element_offsets: Some(res.offsets),
storage: CaStorage::CnTemplate,
group_links: Vec::new(),
});
parent.unreadable = None;
}
return Ok(CompositionOutcome::Expanded);
}
_ => {
return Ok(CompositionOutcome::UnsupportedArray(
UnreadableReason::ArrayComposition,
))
}
}
}
_ => {
}
}
Ok(CompositionOutcome::Expanded)
}
fn build_channel_with_name(
source: &IoBackend,
cn_block: CnBlock,
dg_index: usize,
cg_index: usize,
cn_index: usize,
cache: &mut BlockCache,
name: String,
) -> Result<Channel> {
let unit = cache.get_or_parse_text(source, cn_block.md_unit)?;
let comment = cache.get_or_parse_text(source, cn_block.md_comment)?;
let conversion =
if let Some(cc_arc) = cache.get_or_parse_cc(source, cn_block.cc_conversion)? {
build_conversion(source, &cc_arc, cache)?
} else {
Conversion::None
};
let source_info = if let Some(si_arc) = cache.get_or_parse_si(source, cn_block.si_source)? {
Some(build_source_info(source, &si_arc, cache)?)
} else {
None
};
let (min_value, max_value) = if cn_block.flags.range_valid {
(Some(cn_block.val_range_min), Some(cn_block.val_range_max))
} else {
(None, None)
};
Ok(Channel {
id: cn_index,
index: cn_index,
channel_group_index: cg_index,
data_group_index: dg_index,
name,
unit: unit.to_string(),
channel_type: cn_block.channel_type,
sync_type: cn_block.sync_type,
data_type: cn_block.data_type,
conversion,
bit_count: cn_block.bit_count,
byte_offset: cn_block.byte_offset,
bit_offset: cn_block.bit_offset,
all_invalid: cn_block.flags.all_invalid,
invalidation_bit: cn_block.flags.invalidation_bit,
inval_bit_pos: cn_block.inval_bit_pos,
comment: comment.to_string(),
source: source_info,
min_value,
max_value,
sample_count: 0,
cn_offset: cn_block.header.offset,
data_link: cn_block.data,
unreadable: Self::unreadable_reason(&cn_block),
array_shape: None,
array_element: None,
array_dynamic_size: None,
})
}
fn build_channel(
source: &IoBackend,
cn_block: CnBlock,
dg_index: usize,
cg_index: usize,
cn_index: usize,
cache: &mut BlockCache,
) -> Result<Channel> {
let name = cache.get_or_parse_text(source, cn_block.tx_name)?;
let unit = cache.get_or_parse_text(source, cn_block.md_unit)?;
let comment = cache.get_or_parse_text(source, cn_block.md_comment)?;
let conversion =
if let Some(cc_arc) = cache.get_or_parse_cc(source, cn_block.cc_conversion)? {
build_conversion(source, &cc_arc, cache)?
} else {
Conversion::None
};
let source_info = if let Some(si_arc) = cache.get_or_parse_si(source, cn_block.si_source)? {
Some(build_source_info(source, &si_arc, cache)?)
} else {
None
};
let (min_value, max_value) = if cn_block.flags.range_valid {
(Some(cn_block.val_range_min), Some(cn_block.val_range_max))
} else {
(None, None)
};
Ok(Channel {
id: cn_index,
index: cn_index,
channel_group_index: cg_index,
data_group_index: dg_index,
name: name.to_string(),
unit: unit.to_string(),
channel_type: cn_block.channel_type,
sync_type: cn_block.sync_type,
data_type: cn_block.data_type,
conversion,
bit_count: cn_block.bit_count,
byte_offset: cn_block.byte_offset,
bit_offset: cn_block.bit_offset,
all_invalid: cn_block.flags.all_invalid,
invalidation_bit: cn_block.flags.invalidation_bit,
inval_bit_pos: cn_block.inval_bit_pos,
comment: comment.to_string(),
source: source_info,
min_value,
max_value,
sample_count: 0,
cn_offset: cn_block.header.offset,
data_link: cn_block.data,
unreadable: Self::unreadable_reason(&cn_block),
array_shape: None,
array_element: None,
array_dynamic_size: None,
})
}
fn unreadable_reason(_cn_block: &CnBlock) -> Option<UnreadableReason> {
None
}
pub fn version(&self) -> Mf4Version {
self.version
}
pub fn unfinalized(&self) -> Option<UnfinalizedFlags> {
self.unfinalized
}
pub fn start_time(&self) -> &RecordingTime {
&self.start_time
}
pub fn comment(&self) -> &str {
&self.comment
}
pub fn data_groups(&self) -> &[DataGroup] {
&self.data_groups
}
pub fn data_group_count(&self) -> usize {
self.data_groups.len()
}
pub fn channel_count(&self) -> usize {
self.data_groups
.iter()
.flat_map(|dg| dg.channel_groups.iter())
.map(|cg| cg.channels.len())
.sum()
}
pub fn channels(&self) -> impl Iterator<Item = &Channel> {
self.data_groups
.iter()
.flat_map(|dg| dg.channel_groups.iter())
.flat_map(|cg| cg.channels.iter())
}
pub fn find_channel(&self, name: &str) -> Option<&Channel> {
if let Some(loc) = self.channels_db.find_first(name) {
return Some(
&self.data_groups[loc.data_group_index].channel_groups[loc.channel_group_index]
.channels[loc.channel_index],
);
}
if self.channels_db.is_empty() {
return self.channels().find(|ch| ch.name == name);
}
None
}
pub fn find_channels(&self, name: &str) -> Vec<&Channel> {
if self.channels_db.is_empty() {
return self.channels().filter(|ch| ch.name == name).collect();
}
self.channels_db
.find_all(name)
.iter()
.map(|loc| {
&self.data_groups[loc.data_group_index].channel_groups[loc.channel_group_index]
.channels[loc.channel_index]
})
.collect()
}
pub fn has_channel(&self, name: &str) -> bool {
if self.channels_db.is_empty() {
return self.channels().any(|ch| ch.name == name);
}
self.channels_db.contains(name)
}
pub fn channel_names(&self) -> Vec<&str> {
if self.channels_db.is_empty() {
let mut names: Vec<&str> = self.channels().map(|ch| ch.name.as_str()).collect();
names.sort_unstable();
names.dedup();
return names;
}
self.channels_db.names().collect()
}
pub fn channels_matching<F>(&self, mut predicate: F) -> Vec<&Channel>
where
F: FnMut(&str) -> bool,
{
let mut out: Vec<&Channel> = self.channels().filter(|ch| predicate(&ch.name)).collect();
out.sort_by(|a, b| {
a.name
.cmp(&b.name)
.then(a.data_group_index.cmp(&b.data_group_index))
.then(a.channel_group_index.cmp(&b.channel_group_index))
.then(a.index.cmp(&b.index))
});
out
}
pub fn find_channel_locations(&self, name: &str) -> Vec<ChannelLocation> {
if self.channels_db.is_empty() {
return self
.channels()
.filter(|ch| ch.name == name)
.map(|ch| {
ChannelLocation::new(ch.data_group_index, ch.channel_group_index, ch.index)
})
.collect();
}
self.channels_db.find_all(name).to_vec()
}
pub fn search_channels(&self, pattern: &str, mode: SearchMode) -> Vec<String> {
let all_names: Vec<&str> = if self.channels_db.is_empty() {
self.channel_names()
} else {
self.channels_db.names().collect()
};
let matches: Vec<String> = all_names
.into_iter()
.filter(|name| match mode {
SearchMode::Plain => name.contains(pattern),
SearchMode::CaseInsensitive => {
name.to_lowercase().contains(&pattern.to_lowercase())
}
SearchMode::Wildcard => crate::channels_db::wildcard_match(name, pattern),
})
.map(String::from)
.collect();
matches
}
pub fn master_channel(
&self,
data_group_index: usize,
channel_group_index: usize,
) -> Option<&Channel> {
self.masters_db
.find(data_group_index, channel_group_index)
.map(|ch_idx| {
&self.data_groups[data_group_index].channel_groups[channel_group_index].channels
[ch_idx]
})
}
pub fn can_frame_groups(&self) -> Vec<&ChannelGroup> {
self.data_groups
.iter()
.flat_map(|dg| dg.channel_groups.iter())
.filter(|cg| crate::bus::is_can_frame_group(cg))
.collect()
}
pub fn lin_frame_groups(&self) -> Vec<&ChannelGroup> {
self.data_groups
.iter()
.flat_map(|dg| dg.channel_groups.iter())
.filter(|cg| crate::lin::is_lin_frame_group(cg))
.collect()
}
pub fn eth_frame_groups(&self) -> Vec<&ChannelGroup> {
self.data_groups
.iter()
.flat_map(|dg| dg.channel_groups.iter())
.filter(|cg| crate::eth::is_eth_frame_group(cg))
.collect()
}
pub fn flexray_frame_groups(&self) -> Vec<&ChannelGroup> {
self.data_groups
.iter()
.flat_map(|dg| dg.channel_groups.iter())
.filter(|cg| crate::flexray::is_flexray_frame_group(cg))
.collect()
}
pub fn can_frames(&self, group: &ChannelGroup) -> Result<crate::bus::CanFrames> {
crate::bus::read_can_frames(self, group)
}
pub fn lin_frames(&self, group: &ChannelGroup) -> Result<crate::lin::LinFrames> {
crate::lin::read_lin_frames(self, group)
}
pub fn eth_frames(&self, group: &ChannelGroup) -> Result<crate::eth::EthFrames> {
crate::eth::read_eth_frames(self, group)
}
pub fn flexray_frames(&self, group: &ChannelGroup) -> Result<crate::flexray::FlexRayFrames> {
crate::flexray::read_flexray_frames(self, group)
}
pub fn decode_bus<'a>(
&self,
database: &'a crate::candb::CanDatabase,
) -> Result<crate::bus::BusSignals<'a>> {
crate::bus::decode_bus_signals(self, database)
}
pub fn decode_lin<'a>(
&self,
database: &'a crate::candb::CanDatabase,
) -> Result<crate::bus::BusSignals<'a>> {
crate::lin::decode_lin_signals(self, database)
}
#[cfg(feature = "asc")]
pub fn export_asc<W: std::io::Write>(&self, out: &mut W) -> Result<()> {
crate::export::write_asc(self, out)
}
pub fn statistics(&self) -> FileStatistics {
FileStatistics::from_data_groups(&self.data_groups, self.file_size)
}
pub fn signal(&self, channel: &Channel) -> Result<Signal> {
let key = (channel.data_group_index, channel.channel_group_index);
let records = self.records_for(key)?;
self.signal_over(
channel,
records.data.clone(),
records.layout,
records.sample_count,
)
}
pub fn signals<C: std::borrow::Borrow<Channel>>(&self, channels: &[C]) -> Result<Vec<Signal>> {
if channels.is_empty() {
return Ok(Vec::new());
}
let mut result: Vec<Option<Signal>> = (0..channels.len()).map(|_| None).collect();
let mut groups: std::collections::BTreeMap<(usize, usize), Vec<(usize, &Channel)>> =
std::collections::BTreeMap::new();
for (i, ch) in channels.iter().enumerate() {
let ch = ch.borrow();
let key = (ch.data_group_index, ch.channel_group_index);
groups.entry(key).or_default().push((i, ch));
}
for (key, group_channels) in groups {
let records = self.records_for(key)?;
for (idx, ch) in group_channels {
let signal = self.signal_over(
ch,
records.data.clone(),
records.layout,
records.sample_count,
)?;
result[idx] = Some(signal);
}
}
Ok(result
.into_iter()
.map(|s| s.expect("all channel slots populated"))
.collect())
}
pub fn time_series(&self, channel: &Channel) -> Result<crate::time_ops::SignalSeries> {
let signal = self.signal(channel)?;
let timestamps = self.channel_timestamps(channel)?;
let values = signal.values()?;
let validity = signal.validity();
crate::time_ops::SignalSeries::new(channel.clone(), timestamps, values, validity)
}
fn channel_timestamps(&self, channel: &Channel) -> Result<Vec<f64>> {
let dg = &self.data_groups[channel.data_group_index];
let cg = &dg.channel_groups[channel.channel_group_index];
if let Some(master) = cg.master_channel() {
let sig = self.signal(master)?;
sig.values_f64()
} else {
let sample_count = if channel.sample_count > 0 {
channel.sample_count as usize
} else {
cg.sample_count as usize
};
Ok((0..sample_count).map(|i| i as f64).collect())
}
}
pub(crate) fn series_for<C: std::borrow::Borrow<Channel>>(
&self,
channels: &[C],
) -> Result<Vec<crate::time_ops::SignalSeries>> {
if channels.is_empty() {
return Ok(Vec::new());
}
let signals = self.signals(channels)?;
let mut time_cache: std::collections::HashMap<(usize, usize), Vec<f64>> =
std::collections::HashMap::new();
let mut out = Vec::with_capacity(channels.len());
for (ch, signal) in channels.iter().zip(signals) {
let ch = ch.borrow();
let key = (ch.data_group_index, ch.channel_group_index);
let timestamps = if let Some(ts) = time_cache.get(&key) {
ts.clone()
} else {
let ts = self.channel_timestamps(ch)?;
time_cache.insert(key, ts.clone());
ts
};
let values = signal.values()?;
let validity = signal.validity();
out.push(crate::time_ops::SignalSeries::new(
ch.clone(),
timestamps,
values,
validity,
)?);
}
Ok(out)
}
pub fn filter(
&self,
selectors: &[crate::multi_ops::ChannelSelector],
) -> Result<Vec<crate::time_ops::SignalSeries>> {
let channels = selectors
.iter()
.map(|s| s.resolve(self))
.collect::<Result<Vec<&Channel>>>()?;
self.series_for(&channels)
}
pub fn cut<C: std::borrow::Borrow<Channel>>(
&self,
channels: &[C],
start: f64,
end: f64,
) -> Result<Vec<crate::time_ops::SignalSeries>> {
Ok(self
.series_for(channels)?
.iter()
.map(|s| s.cut(start, end))
.collect())
}
pub fn cut_channel(
&self,
channel: &Channel,
start: f64,
end: f64,
) -> Result<crate::time_ops::SignalSeries> {
let series = self.time_series(channel)?;
Ok(series.cut(start, end))
}
pub fn resample<C: std::borrow::Borrow<Channel>>(
&self,
channels: &[C],
raster: impl Into<crate::time_ops::Raster>,
mode: crate::time_ops::InterpolationMode,
) -> Result<Vec<crate::time_ops::SignalSeries>> {
if channels.is_empty() {
return Ok(Vec::new());
}
let raster = raster.into();
let series_list = self.series_for(channels)?;
let target_timestamps = match raster {
crate::time_ops::Raster::Step(dt) => {
if dt <= 0.0 || !dt.is_finite() {
return Err(Mf4Error::parse_error(format!(
"resample raster step must be positive and finite, got {dt}"
)));
}
let mut global_min = f64::INFINITY;
let mut global_max = f64::NEG_INFINITY;
for s in &series_list {
if let Some(&first) = s.timestamps.first() {
if first < global_min {
global_min = first;
}
}
if let Some(&last) = s.timestamps.last() {
if last > global_max {
global_max = last;
}
}
}
if global_min.is_infinite() || global_max.is_infinite() {
Vec::new()
} else {
crate::time_ops::generate_raster_grid(global_min, global_max, dt)
}
}
crate::time_ops::Raster::Timestamps(ts) => ts,
};
let mut out = Vec::with_capacity(series_list.len());
for s in series_list {
let resampled_values = crate::time_ops::resample_values(
&s.values,
&s.timestamps,
&target_timestamps,
mode,
);
let resampled_validity = crate::time_ops::resample_validity(
s.validity.as_deref(),
&s.timestamps,
&target_timestamps,
);
let resampled = crate::time_ops::SignalSeries::new(
s.channel,
target_timestamps.clone(),
resampled_values,
resampled_validity,
)?;
out.push(resampled);
}
Ok(out)
}
pub fn resample_channel(
&self,
channel: &Channel,
raster: impl Into<crate::time_ops::Raster>,
mode: crate::time_ops::InterpolationMode,
) -> Result<crate::time_ops::SignalSeries> {
let series = self.time_series(channel)?;
series.resample(raster, mode)
}
pub(crate) fn signal_over(
&self,
channel: &Channel,
data: Arc<Vec<u8>>,
layout: RecordLayout,
sample_count: usize,
) -> Result<Signal> {
let mut signal = Signal::new(channel.clone(), data, layout, sample_count);
if channel.channel_type == ChannelType::VariableLength {
signal.attach_payloads(self.payloads_for(channel)?);
}
if channel.channel_type == ChannelType::MaxLength {
if let Some(length) = self.mlsd_length_for(channel)? {
signal.attach_mlsd_length(length);
}
}
if let Some(size_ref) = channel.array_dynamic_size {
if let Some(length) = self.dynamic_array_length_for(channel, size_ref)? {
signal.attach_dynamic_array_length(length);
}
}
if let Some(ref elem) = channel.array_element {
if elem.storage == CaStorage::CgTemplate || elem.storage == CaStorage::DgTemplate {
let members = self.resolve_array_group_members(channel, elem)?;
signal.attach_array_group_members(members);
}
}
Ok(signal)
}
fn resolve_array_group_members(
&self,
channel: &Channel,
elem: &crate::model::ArrayElement,
) -> Result<Vec<crate::model::signal::ArrayGroupMember>> {
let expected_elements = channel
.array_shape
.as_ref()
.map(|s| s.iter().copied().fold(1u64, |acc, d| acc.saturating_mul(d)) as usize)
.unwrap_or(0);
if elem.group_links.len() != expected_elements {
return Err(Mf4Error::parse_error(format!(
"channel '{}' declares {} array elements, but its CA block link list holds {}",
channel.name,
expected_elements,
elem.group_links.len()
)));
}
let mut member_keys = Vec::with_capacity(elem.group_links.len());
match elem.storage {
CaStorage::CgTemplate => {
for (i, &link) in elem.group_links.iter().enumerate() {
if link == 0 {
return Err(Mf4Error::parse_error(format!(
"channel '{}' CG-template array member {} has null link",
channel.name, i
)));
}
let mut found = None;
for (dg_idx, dg) in self.data_groups.iter().enumerate() {
for (cg_idx, cg) in dg.channel_groups.iter().enumerate() {
if cg.cg_offset == link {
found = Some((dg_idx, cg_idx));
break;
}
}
if found.is_some() {
break;
}
}
let key = found.ok_or_else(|| {
Mf4Error::parse_error(format!(
"channel '{}' CG-template member channel group at offset {:#x} is missing or not found",
channel.name, link
))
})?;
member_keys.push(key);
}
}
CaStorage::DgTemplate => {
for (i, &link) in elem.group_links.iter().enumerate() {
if link == 0 {
return Err(Mf4Error::parse_error(format!(
"channel '{}' DG-template array member {} has null link",
channel.name, i
)));
}
let mut found = None;
for (dg_idx, dg) in self.data_groups.iter().enumerate() {
if dg.dg_offset == link {
found = Some(dg_idx);
break;
}
}
let dg_idx = found.ok_or_else(|| {
Mf4Error::parse_error(format!(
"channel '{}' DG-template member data group at offset {:#x} is missing or not found",
channel.name, link
))
})?;
if self.data_groups[dg_idx].channel_groups.is_empty() {
return Err(Mf4Error::parse_error(format!(
"channel '{}' DG-template member data group at offset {:#x} has no channel groups",
channel.name, link
)));
}
member_keys.push((dg_idx, 0));
}
}
_ => return Ok(Vec::new()),
}
let mut members = Vec::with_capacity(member_keys.len());
let mut expected_samples: Option<usize> = None;
for (k, &key) in member_keys.iter().enumerate() {
let records = self.records_for(key)?;
let sc = records.sample_count;
if let Some(exp) = expected_samples {
if sc != exp {
return Err(Mf4Error::parse_error(format!(
"channel '{}' member group {} has {} samples, which disagrees with sibling count {}",
channel.name, k, sc, exp
)));
}
} else {
expected_samples = Some(sc);
}
members.push(crate::model::signal::ArrayGroupMember {
raw_data: records.data.clone(),
layout: records.layout,
sample_count: sc,
});
}
Ok(members)
}
fn mlsd_length_for(&self, channel: &Channel) -> Result<Option<MlsdLength>> {
if channel.data_link == 0 {
return Ok(None);
}
let cn = parser::parse_cn_block(&self.source, channel.data_link)?;
Ok(Some(MlsdLength {
byte_offset: cn.byte_offset,
bit_offset: cn.bit_offset,
bit_count: cn.bit_count,
little_endian: cn.data_type.is_little_endian(),
}))
}
fn dynamic_array_length_for(
&self,
channel: &Channel,
size_ref: crate::blocks::AxisRef,
) -> Result<Option<MlsdLength>> {
let dg = &self.data_groups[channel.data_group_index];
let cg = &dg.channel_groups[channel.channel_group_index];
if dg.dg_offset != size_ref.dg || cg.cg_offset != size_ref.cg {
return Ok(None);
}
let cn = parser::parse_cn_block(&self.source, size_ref.cn)?;
Ok(Some(MlsdLength {
byte_offset: cn.byte_offset,
bit_offset: cn.bit_offset,
bit_count: cn.bit_count,
little_endian: cn.data_type.is_little_endian(),
}))
}
fn records_for(&self, key: (usize, usize)) -> Result<CachedRecords> {
if let Ok(guard) = self.record_cache.read() {
if let Some(hit) = guard.get(&key) {
return Ok(hit);
}
}
let built = self.build_records(key)?;
if let Ok(mut guard) = self.record_cache.write() {
guard.insert(key, built.clone(), built.data.len());
}
Ok(built)
}
fn build_records(&self, key: (usize, usize)) -> Result<CachedRecords> {
let (dg_index, cg_index) = key;
let dg = &self.data_groups[dg_index];
let cg = &dg.channel_groups[cg_index];
let raw_data = self.read_raw_data_indexed(dg)?;
if let Some(index) = &dg.record_index {
let payload = cg.payload_size();
let (records, sample_count) = Self::gather_records(
&raw_data,
index.offsets(cg_index),
dg.rec_id_size as usize,
payload,
self.limits,
);
return Ok(CachedRecords {
data: Arc::new(records),
layout: RecordLayout {
record_size: payload,
record_offset: 0,
inval_start: cg.data_bytes_len(),
inval_bytes: cg.inval_bytes_len(),
},
sample_count,
});
}
let record_size = cg.record_size(dg.rec_id_size);
let record_offset = dg.rec_id_size as usize;
let sample_count = if cg.sample_count > 0 {
cg.sample_count as usize
} else {
raw_data.len().checked_div(record_size).unwrap_or(0)
};
Ok(CachedRecords {
data: Arc::new(raw_data),
layout: RecordLayout {
record_size,
record_offset,
inval_start: cg.data_bytes_len(),
inval_bytes: cg.inval_bytes_len(),
},
sample_count,
})
}
fn payloads_for(&self, channel: &Channel) -> Result<Arc<VlsdPayloads>> {
let link = channel.data_link();
if link != 0 {
if let Ok(guard) = self.payload_cache.read() {
if let Some(hit) = guard.get(&link) {
return Ok(hit);
}
}
}
let built = Arc::new(self.vlsd_payloads(channel)?);
if link != 0 {
if let Ok(mut guard) = self.payload_cache.write() {
guard.insert(link, built.clone(), built.total_bytes());
}
}
Ok(built)
}
fn vlsd_payloads(&self, channel: &Channel) -> Result<VlsdPayloads> {
let dg = &self.data_groups[channel.data_group_index];
let link = channel.data_link();
if link == 0 {
return Err(Mf4Error::unsupported(
"variable-length signal data (VLSD)",
format!("channel '{}' has no signal-data link", channel.name),
));
}
if let Some(cg_index) = dg
.channel_groups
.iter()
.position(|cg| cg.matches_offset(link))
{
let Some(index) = &dg.record_index else {
return Err(Mf4Error::unsupported(
"variable-length signal data (VLSD)",
format!(
"channel '{}' points at a channel group in a sorted data group",
channel.name
),
));
};
let raw = self.read_raw_data_indexed(dg)?;
return Ok(VlsdPayloads::from_records(
&raw,
index.offsets(cg_index),
dg.rec_id_size as usize,
));
}
let block_id = self.source.read_bytes(link, 4)?;
match &block_id[..] {
b"##SD" | b"##DT" | b"##DV" | b"##DZ" | b"##DL" | b"##LD" | b"##HL" => {
let index =
Self::build_data_block_index(&self.source, link, false, self.file_size, &[])?;
let mut stream = Vec::new();
for (_offset, info) in index.iter() {
stream.extend(Self::read_block_payload(&self.source, info, self.limits)?);
}
Ok(VlsdPayloads::from_stream(&stream))
}
other => Err(Mf4Error::unsupported(
"variable-length signal data (VLSD)",
format!(
"channel '{}' points at an unexpected block '{}'",
channel.name,
String::from_utf8_lossy(other)
),
)),
}
}
fn gather_records(
raw: &[u8],
offsets: &[u64],
rec_id_size: usize,
payload: usize,
limits: Limits,
) -> (Vec<u8>, usize) {
if payload == 0 {
return (Vec::new(), 0);
}
let mut out =
Vec::with_capacity(offsets.len().saturating_mul(payload).min(limits.max_alloc));
for &offset in offsets {
let start = offset as usize + rec_id_size;
let Some(slice) = raw.get(start..start + payload) else {
break;
};
out.extend_from_slice(slice);
}
let count = out.len() / payload;
(out, count)
}
fn read_raw_data_indexed(&self, dg: &DataGroup) -> Result<Vec<u8>> {
if dg.data_block_index.is_empty() {
return Ok(Vec::new());
}
let total_size = dg.data_block_index.total_size() as usize;
let capacity = total_size.min(self.file_size as usize);
let mut all_data = Vec::with_capacity(capacity);
for (_offset, block_info) in dg.data_block_index.iter() {
self.append_data_block(block_info, &mut all_data)?;
}
Ok(all_data)
}
pub(crate) fn read_block_range(
&self,
info: &DataBlockInfo,
start: usize,
len: usize,
) -> Result<Vec<u8>> {
if info.compression.is_some() || info.invalidation_block.is_some() {
let mut whole = Vec::new();
self.append_data_block(info, &mut whole)?;
let end = start.saturating_add(len).min(whole.len());
let from = start.min(end);
return Ok(whole[from..end].to_vec());
}
let offset = info.offset + BLOCK_HEADER_SIZE as u64 + start as u64;
let available = (info.original_size as usize).saturating_sub(start);
Ok(self.source.read_bytes(offset, len.min(available))?.to_vec())
}
pub(crate) fn read_data_index_range(
&self,
index: &DataBlockIndex,
start: u64,
len: usize,
) -> Result<Vec<u8>> {
if len == 0 || start >= index.total_size() {
return Ok(Vec::new());
}
let end = (start + len as u64).min(index.total_size());
let actual_len = (end - start) as usize;
let mut out = Vec::with_capacity(actual_len);
let mut current_offset = start;
while current_offset < end {
let Some((_block_idx, block_info, local_start)) =
index.block_for_offset(current_offset)
else {
break;
};
let block_size = block_info.effective_size();
let remaining_in_block = (block_size - local_start) as usize;
let bytes_to_read = remaining_in_block.min((end - current_offset) as usize);
let block_bytes =
self.read_block_range(block_info, local_start as usize, bytes_to_read)?;
if block_bytes.is_empty() {
break;
}
let read_len = block_bytes.len();
out.extend_from_slice(&block_bytes);
current_offset += read_len as u64;
}
Ok(out)
}
pub(crate) fn append_data_block(&self, info: &DataBlockInfo, out: &mut Vec<u8>) -> Result<()> {
let payload = Self::read_block_payload(&self.source, info, self.limits)?;
out.extend_from_slice(&payload);
Ok(())
}
#[doc(hidden)]
pub fn un_transpose(transposed: &[u8], column_size: usize) -> Result<Vec<u8>> {
if column_size == 0 {
return Err(Mf4Error::Decompression(
"Invalid transposition parameter".to_string(),
));
}
let lines = transposed.len() / column_size;
if lines == 0 {
return Ok(transposed.to_vec());
}
let prefix_len = lines * column_size;
let mut result = vec![0u8; transposed.len()];
for (src_idx, &byte) in transposed[..prefix_len].iter().enumerate() {
let col = src_idx / lines;
let line = src_idx % lines;
let dst_idx = line * column_size + col;
result[dst_idx] = byte;
}
result[prefix_len..].copy_from_slice(&transposed[prefix_len..]);
Ok(result)
}
fn decompress_deflate(
compressed: &[u8],
original_size: usize,
limits: Limits,
) -> Result<Vec<u8>> {
use flate2::read::ZlibDecoder;
use std::io::Read;
let mut decoder = ZlibDecoder::new(compressed).take(limits.max_decompressed);
let mut decompressed = Vec::with_capacity(original_size.min(limits.max_alloc));
decoder
.read_to_end(&mut decompressed)
.map_err(|e| Mf4Error::Decompression(e.to_string()))?;
if decompressed.len() != original_size {
return Err(Mf4Error::Decompression(format!(
"Decompressed size mismatch: expected {original_size} bytes, got {}",
decompressed.len()
)));
}
Ok(decompressed)
}
fn decompress_zstd(compressed: &[u8], original_size: usize, limits: Limits) -> Result<Vec<u8>> {
#[cfg(feature = "zstd")]
{
use std::io::Read;
let mut decoder = ruzstd::StreamingDecoder::new(compressed)
.map_err(|e| Mf4Error::Decompression(format!("zstd initialization error: {e:?}")))?
.take(limits.max_decompressed);
let mut decompressed = Vec::with_capacity(original_size.min(limits.max_alloc));
decoder
.read_to_end(&mut decompressed)
.map_err(|e| Mf4Error::Decompression(format!("zstd decoding error: {e:?}")))?;
if decompressed.len() != original_size {
return Err(Mf4Error::Decompression(format!(
"Decompressed size mismatch: expected {original_size} bytes, got {}",
decompressed.len()
)));
}
Ok(decompressed)
}
#[cfg(not(feature = "zstd"))]
{
let _ = (compressed, original_size, limits);
Err(Mf4Error::unsupported(
"zstd compression",
"enable the 'zstd' cargo feature of falcon_mdf to read zstd-compressed DZ blocks",
))
}
}
fn decompress_lz4(compressed: &[u8], original_size: usize, limits: Limits) -> Result<Vec<u8>> {
#[cfg(feature = "lz4")]
{
use std::io::Read;
let mut decoder =
lz4_flex::frame::FrameDecoder::new(compressed).take(limits.max_decompressed);
let mut decompressed = Vec::with_capacity(original_size.min(limits.max_alloc));
decoder
.read_to_end(&mut decompressed)
.map_err(|e| Mf4Error::Decompression(format!("lz4 decoding error: {e:?}")))?;
if decompressed.len() != original_size {
return Err(Mf4Error::Decompression(format!(
"Decompressed size mismatch: expected {original_size} bytes, got {}",
decompressed.len()
)));
}
Ok(decompressed)
}
#[cfg(not(feature = "lz4"))]
{
let _ = (compressed, original_size, limits);
Err(Mf4Error::unsupported(
"lz4 compression",
"enable the 'lz4' cargo feature of falcon_mdf to read lz4-compressed DZ blocks",
))
}
}
fn decompress(
compressed: &[u8],
compression: &CompressionInfo,
original_size: usize,
limits: Limits,
) -> Result<Vec<u8>> {
match compression.algorithm {
CompressionType::Deflate => Self::decompress_deflate(compressed, original_size, limits),
CompressionType::TransposedDeflate => {
let raw = Self::decompress_deflate(compressed, original_size, limits)?;
Self::un_transpose(&raw, compression.parameter as usize)
}
CompressionType::Zstd => Self::decompress_zstd(compressed, original_size, limits),
CompressionType::TransposedZstd => {
let raw = Self::decompress_zstd(compressed, original_size, limits)?;
Self::un_transpose(&raw, compression.parameter as usize)
}
CompressionType::Lz4 => Self::decompress_lz4(compressed, original_size, limits),
CompressionType::TransposedLz4 => {
let raw = Self::decompress_lz4(compressed, original_size, limits)?;
Self::un_transpose(&raw, compression.parameter as usize)
}
CompressionType::Unknown(t) => Err(Mf4Error::Decompression(format!(
"Unknown compression type: {}",
t
))),
}
}
pub fn reduced_signal(
&self,
channel: &Channel,
reduction: &SampleReduction,
kind: ReductionKind,
) -> Result<Signal> {
let dg = &self.data_groups[channel.data_group_index];
let cg = &dg.channel_groups[channel.channel_group_index];
if reduction.data_link == 0 {
return Err(Mf4Error::unsupported(
"sample reduction",
format!(
"the reduction attached to '{}' names no data block",
channel.name
),
));
}
let block = cg.data_bytes_len();
if block == 0 {
return Err(Mf4Error::unsupported(
"sample reduction",
format!("channel group of '{}' has no record bytes", channel.name),
));
}
let index = Self::build_data_block_index(
&self.source,
reduction.data_link,
false,
self.file_size,
&[],
)?;
let mut records = Vec::new();
for (_offset, info) in index.iter() {
self.append_data_block(info, &mut records)?;
}
let available = records.len() / (block * 3);
let sample_count = (reduction.cycle_count as usize).min(available);
Ok(Signal::new(
channel.clone(),
Arc::new(records),
RecordLayout {
record_size: block * 3,
record_offset: kind.index() * block,
inval_start: cg.data_bytes_len(),
inval_bytes: 0,
},
sample_count,
))
}
pub fn metadata(&self) -> &Metadata {
&self.metadata
}
pub fn file_size(&self) -> u64 {
self.file_size
}
pub fn to_writer(&self) -> Result<crate::write::Mf4Writer> {
crate::write::Mf4Writer::from_file(self)
}
pub fn block_map(&self) -> crate::inspect::BlockMap {
crate::inspect::BlockMap::scan(self.source.as_ref())
}
pub fn read_raw(&self, offset: u64, len: usize) -> Result<Vec<u8>> {
Ok(self.source.read_bytes(offset, len)?.to_vec())
}
pub fn channels_db_stats(&self) -> (usize, usize) {
(
self.channels_db.unique_name_count(),
self.channels_db.total_channel_count(),
)
}
pub fn cache_stats(&self) -> &crate::cache::CacheStats {
self.cache.stats()
}
pub fn attachments(&self) -> &[Attachment] {
&self.attachments
}
pub fn events(&self) -> &[Event] {
&self.events
}
pub fn file_history(&self) -> &[FileHistoryEntry] {
&self.file_history
}
pub fn channel_hierarchy(&self) -> &[ChannelHierarchyNode] {
&self.hierarchy
}
pub fn channel_at(&self, element: &ChElement) -> Option<&Channel> {
let dg = self
.data_groups
.iter()
.find(|g| g.dg_offset == element.data_group)?;
let cg = dg
.channel_groups
.iter()
.find(|g| g.cg_offset == element.channel_group)?;
cg.channels.iter().find(|c| c.cn_offset == element.channel)
}
pub fn attachment_data(&self, attachment: &Attachment) -> Result<Option<Vec<u8>>> {
self.attachment_data_with_password(attachment, None)
}
pub fn attachment_data_with_password(
&self,
attachment: &Attachment,
password: Option<&str>,
) -> Result<Option<Vec<u8>>> {
if !attachment.is_embedded || attachment.embedded_offset == 0 {
return Ok(None);
}
let data = self.source.read_bytes(
attachment.embedded_offset,
attachment.embedded_size as usize,
)?;
let raw_payload = if !attachment.is_compressed {
data.to_vec()
} else {
let compression = CompressionInfo {
algorithm: CompressionType::Deflate,
parameter: 0,
data_offset: attachment.embedded_offset,
};
Self::decompress(
&data,
&compression,
attachment.original_size as usize,
self.limits,
)?
};
if let Some(enc_info) = attachment.encryption_info() {
if enc_info.encrypted {
let Some(pw) = password else {
return Err(Mf4Error::unsupported(
"encrypted attachment",
"password must be provided for encrypted attachments",
));
};
if !enc_info.algorithm.eq_ignore_ascii_case("aes256") {
return Err(Mf4Error::unsupported(
"attachment encryption",
format!(
"not implemented attachment encryption algorithm <{}>",
enc_info.algorithm
),
));
}
let computed_md5 = crate::crypto::md5_hex(&raw_payload);
if !computed_md5.eq_ignore_ascii_case(&enc_info.original_md5_sum) {
return Err(Mf4Error::parse_error(format!(
"MD5 sum mismatch for encrypted attachment: original={} and computed={}",
enc_info.original_md5_sum, computed_md5
)));
}
if raw_payload.len() < 16 {
return Err(Mf4Error::parse_error(
"encrypted attachment data is shorter than IV length (16 bytes)",
));
}
let iv: [u8; 16] = raw_payload[..16].try_into().unwrap();
let ciphertext = &raw_payload[16..];
if ciphertext.len() % 16 != 0 {
return Err(Mf4Error::parse_error(
"encrypted attachment ciphertext length is not a multiple of 16 bytes",
));
}
let key = crate::crypto::derive_aes256_key(pw.as_bytes());
let mut decrypted = crate::crypto::aes256_cbc_decrypt(ciphertext, &key, &iv);
decrypted.truncate(enc_info.original_size.min(decrypted.len()));
return Ok(Some(decrypted));
}
}
Ok(Some(raw_payload))
}
fn parse_attachments(
source: &IoBackend,
first_at: u64,
cache: &mut BlockCache,
) -> Result<Vec<Attachment>> {
let mut attachments = Vec::new();
let mut offset = first_at;
let mut chain = LinkChain::new();
while offset != 0 {
chain.visit(offset, "at_next")?;
let at_block = parser::parse_at_block(source, offset)?;
let file_name = cache
.get_or_parse_text(source, at_block.tx_file_name)?
.to_string();
let file_path = cache
.get_or_parse_text(source, at_block.tx_file_path)?
.to_string();
let comment = cache
.get_or_parse_text(source, at_block.md_comment)?
.to_string();
attachments.push(Attachment {
file_name,
file_path,
comment,
is_embedded: at_block.is_embedded(),
is_compressed: at_block.flags.compressed,
original_size: at_block.original_size,
md5_checksum: at_block.md5_checksum,
md5_valid: at_block.flags.md5_valid,
embedded_offset: at_block.embedded_data_offset(),
embedded_size: at_block.embedded_size,
});
offset = at_block.at_next;
}
Ok(attachments)
}
fn parse_events(
source: &IoBackend,
first_ev: u64,
cache: &mut BlockCache,
) -> Result<Vec<Event>> {
let mut events = Vec::new();
let mut offset = first_ev;
let mut chain = LinkChain::new();
while offset != 0 {
chain.visit(offset, "ev_next")?;
let ev_block = parser::parse_ev_block(source, offset)?;
let comment = cache
.get_or_parse_text(source, ev_block.md_comment)?
.to_string();
let name = cache
.get_or_parse_text(source, ev_block.tx_name)?
.to_string();
events.push(Event {
event_type: ev_block.ev_type,
sync_type: ev_block.ev_sync_type,
range_type: ev_block.ev_range_type,
cause: ev_block.ev_cause,
sync_base_value: ev_block.ev_sync_base_value,
sync_factor: ev_block.ev_sync_factor,
scope_count: ev_block.ev_scope_count,
attachment_count: ev_block.ev_attachment_count,
comment,
name,
});
offset = ev_block.ev_next;
}
Ok(events)
}
fn parse_sample_reductions(source: &IoBackend, first_sr: u64) -> Result<Vec<SampleReduction>> {
let mut levels = Vec::new();
let mut offset = first_sr;
let mut chain = LinkChain::new();
while offset != 0 {
chain.visit(offset, "sr_next")?;
let sr = parser::parse_sr_block(source, offset)?;
levels.push(SampleReduction {
cycle_count: sr.sr_cycle_count,
interval: sr.sr_interval,
sync_type: sr.sr_sync_type,
flags: sr.sr_flags,
data_link: sr.sr_data,
});
offset = sr.sr_next;
}
Ok(levels)
}
fn parse_file_history(
source: &IoBackend,
first_fh: u64,
cache: &mut BlockCache,
) -> Result<Vec<FileHistoryEntry>> {
let mut entries = Vec::new();
let mut offset = first_fh;
let mut chain = LinkChain::new();
while offset != 0 {
chain.visit(offset, "fh_next")?;
let fh = parser::parse_fh_block(source, offset)?;
let comment = cache.get_or_parse_text(source, fh.md_comment)?.to_string();
let metadata = parser::read_metadata(source, fh.md_comment)?.unwrap_or_default();
entries.push(FileHistoryEntry {
time: RecordingTime::new(fh.time_ns as i64, fh.tz_offset_min, fh.dst_offset_min),
comment,
metadata,
});
offset = fh.fh_next;
}
Ok(entries)
}
fn parse_hierarchy(
source: &IoBackend,
first_ch: u64,
cache: &mut BlockCache,
visited: &mut std::collections::HashSet<u64>,
) -> Result<Vec<ChannelHierarchyNode>> {
let mut nodes = Vec::new();
let mut offset = first_ch;
let mut chain = LinkChain::new();
while offset != 0 {
chain.visit(offset, "ch_next")?;
if !visited.insert(offset) {
break;
}
let ch_block = parser::parse_ch_block(source, offset)?;
let name = cache
.get_or_parse_text(source, ch_block.tx_name)?
.to_string();
let comment = cache
.get_or_parse_text(source, ch_block.md_comment)?
.to_string();
let children = if ch_block.ch_first != 0 {
Self::parse_hierarchy(source, ch_block.ch_first, cache, visited)?
} else {
Vec::new()
};
nodes.push(ChannelHierarchyNode {
name,
comment,
hierarchy_type: ch_block.ch_type,
elements: ch_block.ch_element.clone(),
has_children: ch_block.ch_first != 0,
children,
});
offset = ch_block.ch_next;
}
Ok(nodes)
}
}
fn build_source_info(
source: &IoBackend,
si: &SiBlock,
cache: &mut BlockCache,
) -> Result<SourceInfo> {
let name = cache.get_or_parse_text(source, si.tx_name)?;
let path = cache.get_or_parse_text(source, si.tx_path)?;
Ok(SourceInfo {
name: name.to_string(),
path: path.to_string(),
source_type: Some(si.source_type),
bus_type: Some(si.bus_type),
simulated: si.flags.simulated,
})
}
fn build_conversion(
source: &IoBackend,
cc: &CcBlock,
cache: &mut BlockCache,
) -> Result<Conversion> {
build_conversion_at_depth(source, cc, cache, 0)
}
const MAX_CONVERSION_DEPTH: u32 = 4;
fn build_conversion_at_depth(
source: &IoBackend,
cc: &CcBlock,
cache: &mut BlockCache,
depth: u32,
) -> Result<Conversion> {
use crate::blocks::{ConversionType as Ct, Expr};
let text_at = |cache: &mut BlockCache, idx: usize| -> Result<Option<String>> {
let Some(&link) = cc.references.get(idx) else {
return Ok(None);
};
if link == 0 {
return Ok(None);
}
let s = cache.get_or_parse_text(source, link)?;
Ok(Some(s.to_string()))
};
let unsupported = |kind: Ct, reason: &str| Conversion::Unsupported {
kind,
reason: reason.to_string(),
};
let v = &cc.values;
Ok(match cc.conversion_type {
Ct::Identity => Conversion::None,
Ct::Linear => {
if v.len() < 2 {
unsupported(Ct::Linear, "linear conversion needs 2 parameters")
} else {
Conversion::Linear {
offset: v[0],
factor: v[1],
}
}
}
Ct::Rational => {
if v.len() < 6 {
unsupported(Ct::Rational, "rational conversion needs 6 parameters")
} else {
Conversion::Rational {
coefficients: [v[0], v[1], v[2], v[3], v[4], v[5]],
}
}
}
Ct::Algebraic => match text_at(cache, 0)? {
Some(formula) => match Expr::parse(&formula) {
Ok(expr) => Conversion::Algebraic { formula, expr },
Err(e) => unsupported(Ct::Algebraic, &e.to_string()),
},
None => unsupported(Ct::Algebraic, "formula text is missing"),
},
Ct::TabInterpolation | Ct::TabLookup => {
let n = v.len() / 2;
if n == 0 {
unsupported(cc.conversion_type, "conversion table has no entries")
} else {
let keys = (0..n).map(|i| v[i * 2]).collect();
let values = (0..n).map(|i| v[i * 2 + 1]).collect();
if cc.conversion_type == Ct::TabInterpolation {
Conversion::TableInterpolated { keys, values }
} else {
Conversion::TableLookup { keys, values }
}
}
}
Ct::TabRangeLookup => {
let n = v.len() / 3;
let default = if v.len() % 3 == 1 {
v.last().copied()
} else {
None
};
if n == 0 && default.is_none() {
return Ok(unsupported(
Ct::TabRangeLookup,
"range table has no entries and no default",
));
}
Conversion::RangeTable {
lower: (0..n).map(|i| v[i * 3]).collect(),
upper: (0..n).map(|i| v[i * 3 + 1]).collect(),
values: (0..n).map(|i| v[i * 3 + 2]).collect(),
default,
}
}
Ct::TabValueToText => {
let n = v.len();
let mut entries = Vec::with_capacity(n);
for i in 0..n {
let link = cc.references.get(i).copied().unwrap_or(0);
entries.push(
table_entry(source, cache, link, depth)?
.unwrap_or_else(|| TableEntry::Text(String::new())),
);
}
let default_link = cc.references.get(n).copied().unwrap_or(0);
Conversion::ValueToText {
keys: v.clone(),
entries,
default: table_entry(source, cache, default_link, depth)?,
}
}
Ct::TabRangeToText => {
let n = v.len() / 2;
let mut entries = Vec::with_capacity(n);
for i in 0..n {
let link = cc.references.get(i).copied().unwrap_or(0);
entries.push(
table_entry(source, cache, link, depth)?
.unwrap_or_else(|| TableEntry::Text(String::new())),
);
}
let default_link = cc.references.get(n).copied().unwrap_or(0);
Conversion::RangeToText {
lower: (0..n).map(|i| v[i * 2]).collect(),
upper: (0..n).map(|i| v[i * 2 + 1]).collect(),
entries,
default: table_entry(source, cache, default_link, depth)?,
}
}
Ct::TabTextToValue => {
let n = cc.references.len();
if v.len() < n {
unsupported(
Ct::TabTextToValue,
"text-to-value table has fewer values than keys",
)
} else {
let mut keys = Vec::with_capacity(n);
for i in 0..n {
keys.push(text_at(cache, i)?.unwrap_or_default());
}
Conversion::TextToValue {
keys,
values: v[..n].to_vec(),
default: v.get(n).copied(),
}
}
}
Ct::TabTextToText => {
let refs = cc.references.len();
if refs.is_multiple_of(2) {
unsupported(
Ct::TabTextToText,
"text-to-text table needs an odd number of references: key/text pairs plus a default",
)
} else {
let n = (refs - 1) / 2;
let mut keys = Vec::with_capacity(n);
let mut texts = Vec::with_capacity(n);
for i in 0..n {
keys.push(text_at(cache, i * 2)?.unwrap_or_default());
texts.push(text_at(cache, i * 2 + 1)?.unwrap_or_default());
}
Conversion::TextToText {
keys,
texts,
default: text_at(cache, refs - 1)?,
}
}
}
Ct::BitfieldToText => {
if depth >= MAX_CONVERSION_DEPTH {
unsupported(
Ct::BitfieldToText,
"nested conversions are more than four deep, which no valid file needs",
)
} else {
let masks: Vec<u64> = v.iter().map(|x| x.to_bits()).collect();
let n = masks.len().min(cc.references.len());
let mut entries = Vec::with_capacity(n);
for i in 0..n {
entries.push(bitfield_entry(source, cache, cc.references[i], depth + 1)?);
}
Conversion::Bitfield {
masks: masks[..n].to_vec(),
entries,
}
}
}
Ct::Unknown(code) => Conversion::Unsupported {
kind: Ct::Unknown(code),
reason: format!("unknown conversion type {code}"),
},
})
}
fn table_entry(
source: &IoBackend,
cache: &mut BlockCache,
link: u64,
depth: u32,
) -> Result<Option<TableEntry>> {
if link == 0 {
return Ok(None);
}
let id = match source.read_bytes(link, 4) {
Ok(bytes) => {
let mut id = [0u8; 4];
id.copy_from_slice(&bytes);
id
}
Err(_) => return Ok(None),
};
match &id {
b"##CC" => {
if depth >= MAX_CONVERSION_DEPTH {
return Ok(None);
}
let Some(nested) = cache.get_or_parse_cc(source, link)? else {
return Ok(None);
};
let conversion = build_conversion_at_depth(source, &nested, cache, depth + 1)?;
Ok(Some(TableEntry::Nested(Box::new(conversion))))
}
b"##TX" | b"##MD" => Ok(Some(TableEntry::Text(
cache.get_or_parse_text(source, link)?.to_string(),
))),
_ => Ok(None),
}
}
fn bitfield_entry(
source: &IoBackend,
cache: &mut BlockCache,
link: u64,
depth: u32,
) -> Result<crate::blocks::BitfieldEntry> {
use crate::blocks::BitfieldEntry;
if link == 0 {
return Ok(BitfieldEntry::Unresolved);
}
let id = match source.read_bytes(link, 4) {
Ok(bytes) => {
let mut id = [0u8; 4];
id.copy_from_slice(&bytes);
id
}
Err(_) => return Ok(BitfieldEntry::Unresolved),
};
match &id {
b"##CC" => {
let Some(nested) = cache.get_or_parse_cc(source, link)? else {
return Ok(BitfieldEntry::Unresolved);
};
let name = cache.get_or_parse_text(source, nested.tx_name)?.to_string();
let conversion = build_conversion_at_depth(source, &nested, cache, depth)?;
Ok(BitfieldEntry::Nested {
name,
conversion: Box::new(conversion),
})
}
b"##TX" | b"##MD" => Ok(BitfieldEntry::Flag(
cache.get_or_parse_text(source, link)?.to_string(),
)),
_ => Ok(BitfieldEntry::Unresolved),
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
enum CompositionOutcome {
Expanded,
UnsupportedArray(UnreadableReason),
}
struct CaChainLevel {
dims: Vec<u64>,
raw_stride: i32,
inverse_layout: bool,
}
struct CaChainResolution {
shape: Vec<u64>,
data_type: DataType,
bit_count: u32,
bit_offset: u8,
byte_offset: u32,
offsets: Vec<usize>,
}
fn resolve_ca_chain(
source: &IoBackend,
parent_data_type: DataType,
parent_bit_count: u32,
parent_bit_offset: u8,
first: &CaBlock,
max_alloc: usize,
) -> Result<Option<CaChainResolution>> {
let elem_cap = max_alloc / std::mem::size_of::<f64>();
let mut levels: Vec<CaChainLevel> = Vec::new();
let mut current = first.clone();
let (data_type, bit_count, bit_offset, byte_offset) = loop {
if levels.len() >= MAX_COMPOSITION_DEPTH {
return Ok(None);
}
levels.push(CaChainLevel {
dims: current.ca_dim_size.clone(),
raw_stride: current.ca_byte_offset_base,
inverse_layout: current.flags.inverse_layout,
});
if current.ca_composition == 0 {
break (parent_data_type, parent_bit_count, parent_bit_offset, 0u32);
}
let id = source.read_bytes(current.ca_composition, 4)?;
match &id[..] {
b"##CN" => {
let t = parser::parse_cn_block(source, current.ca_composition)?;
break (t.data_type, t.bit_count, t.bit_offset, t.byte_offset);
}
b"##CA" => {
let next = parser::parse_ca_block(source, current.ca_composition)?;
if next.ca_storage != CaStorage::CnTemplate || next.flags.dynamic_size {
return Ok(None);
}
current = next;
}
_ => return Ok(None),
}
};
let elem_width = (bit_count as usize).div_ceil(8).max(1);
let total: u64 = levels
.iter()
.flat_map(|l| l.dims.iter().copied())
.try_fold(1u64, |acc, d| acc.checked_mul(d))
.ok_or_else(|| Mf4Error::invalid_block_size("CA", u64::MAX, 1))?;
if total as usize > elem_cap {
return Ok(None);
}
let mut offsets = vec![0usize];
for level in &levels {
let stride = if level.raw_stride > 0 {
level.raw_stride as usize
} else {
elem_width
};
let level_total = level.dims.iter().copied().product::<u64>() as usize;
let level_offsets: Vec<usize> = if level.inverse_layout && level.dims.len() > 1 {
row_major_to_stored(&level.dims, level_total)
.into_iter()
.map(|s| s * stride)
.collect()
} else {
(0..level_total).map(|k| k * stride).collect()
};
let mut next = Vec::with_capacity(offsets.len() * level_offsets.len().max(1));
for &base in &offsets {
for &lo in &level_offsets {
next.push(base + lo);
}
}
offsets = next;
}
let shape: Vec<u64> = levels.into_iter().flat_map(|l| l.dims).collect();
Ok(Some(CaChainResolution {
shape,
data_type,
bit_count,
bit_offset,
byte_offset,
offsets,
}))
}
fn read_record_id(data: &[u8], pos: usize, rec_id_size: u8) -> Option<u64> {
let end = pos.checked_add(rec_id_size as usize)?;
let bytes = data.get(pos..end)?;
Some(match rec_id_size {
1 => bytes[0] as u64,
2 => u16::from_le_bytes(bytes.try_into().ok()?) as u64,
4 => u32::from_le_bytes(bytes.try_into().ok()?) as u64,
8 => u64::from_le_bytes(bytes.try_into().ok()?),
_ => return None,
})
}
fn qualify_channel_name(parent: &str, child: &str) -> String {
if parent.is_empty() {
return child.to_string();
}
if child.is_empty() {
return parent.to_string();
}
if child == parent {
return child.to_string();
}
if let Some(rest) = child.strip_prefix(parent) {
if rest.starts_with('.') {
return child.to_string();
}
}
format!("{parent}.{child}")
}
impl std::fmt::Debug for Mf4File {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("Mf4File")
.field("version", &self.version)
.field("data_group_count", &self.data_groups.len())
.field("channel_count", &self.channel_count())
.field(
"unique_channel_names",
&self.channels_db.unique_name_count(),
)
.field("file_size", &self.file_size)
.finish()
}
}
#[cfg(test)]
mod name_tests {
use super::qualify_channel_name;
#[test]
fn qualifies_bare_child_names() {
assert_eq!(
qualify_channel_name("CAN_DataFrame", "BusChannel"),
"CAN_DataFrame.BusChannel"
);
}
#[test]
fn leaves_already_qualified_names_alone() {
assert_eq!(
qualify_channel_name("CAN_DataFrame", "CAN_DataFrame.BusChannel"),
"CAN_DataFrame.BusChannel"
);
}
#[test]
fn does_not_treat_a_shared_prefix_as_qualification() {
assert_eq!(
qualify_channel_name("CAN_DataFrame", "CAN_DataFrameExtra"),
"CAN_DataFrame.CAN_DataFrameExtra"
);
}
#[test]
fn handles_nested_qualification() {
assert_eq!(qualify_channel_name("A.B", "A.B.C"), "A.B.C");
assert_eq!(qualify_channel_name("A.B", "C"), "A.B.C");
}
#[test]
fn handles_empty_and_identical_parts() {
assert_eq!(qualify_channel_name("", "X"), "X");
assert_eq!(qualify_channel_name("X", ""), "X");
assert_eq!(qualify_channel_name("X", "X"), "X");
}
}
#[cfg(test)]
mod demux_tests {
use super::{read_record_id, Limits, Mf4File};
#[test]
fn reads_record_ids_of_each_permitted_width() {
let data = [0xAA, 0xBB, 0xCC, 0xDD, 0xEE, 0xFF, 0x11, 0x22];
assert_eq!(read_record_id(&data, 0, 1), Some(0xAA));
assert_eq!(read_record_id(&data, 0, 2), Some(0xBBAA));
assert_eq!(read_record_id(&data, 0, 4), Some(0xDDCC_BBAA));
assert_eq!(read_record_id(&data, 0, 8), Some(0x2211_FFEE_DDCC_BBAA));
}
#[test]
fn reads_record_id_at_an_offset() {
let data = [0x00, 0x00, 0x07, 0x00];
assert_eq!(read_record_id(&data, 2, 2), Some(7));
}
#[test]
fn rejects_ids_running_past_the_buffer() {
let data = [0x01, 0x02, 0x03];
assert_eq!(read_record_id(&data, 2, 2), None);
assert_eq!(read_record_id(&data, 3, 1), None);
assert_eq!(read_record_id(&data, 0, 8), None);
}
#[test]
fn rejects_widths_the_format_does_not_define() {
let data = [0u8; 16];
assert_eq!(read_record_id(&data, 0, 0), None);
assert_eq!(read_record_id(&data, 0, 3), None);
assert_eq!(read_record_id(&data, 0, 255), None);
}
#[test]
fn gathers_interleaved_records_into_a_dense_buffer() {
let raw = [
1, 0xAA, 0xBB, 2, 0x11, 0x22, 0x33, 1, 0xCC, 0xDD, 2, 0x44, 0x55, 0x66,
];
let group1 = [0u64, 7];
let group2 = [3u64, 10];
let (out, n) = Mf4File::gather_records(&raw, &group1, 1, 2, Limits::default());
assert_eq!(n, 2);
assert_eq!(out, vec![0xAA, 0xBB, 0xCC, 0xDD]);
let (out, n) = Mf4File::gather_records(&raw, &group2, 1, 3, Limits::default());
assert_eq!(n, 2);
assert_eq!(out, vec![0x11, 0x22, 0x33, 0x44, 0x55, 0x66]);
}
#[test]
fn drops_a_record_truncated_by_the_end_of_the_stream() {
let raw = [1, 0xAA, 0xBB, 0xCC, 0xDD, 1, 0x11, 0x22];
let (out, n) = Mf4File::gather_records(&raw, &[0, 5], 1, 4, Limits::default());
assert_eq!(
n, 1,
"the truncated tail record must be dropped, not padded"
);
assert_eq!(out, vec![0xAA, 0xBB, 0xCC, 0xDD]);
}
#[test]
fn handles_empty_inputs() {
assert_eq!(
Mf4File::gather_records(&[], &[], 1, 4, Limits::default()),
(Vec::new(), 0)
);
assert_eq!(
Mf4File::gather_records(&[1, 2, 3], &[0], 1, 0, Limits::default()),
(Vec::new(), 0)
);
}
#[test]
fn skips_the_record_id_when_gathering() {
let raw = [0xFF, 0xFF, 0x42, 0x43];
let (out, n) = Mf4File::gather_records(&raw, &[0], 2, 2, Limits::default());
assert_eq!(n, 1);
assert_eq!(out, vec![0x42, 0x43], "record ID bytes must not be copied");
}
}
#[cfg(test)]
mod lru_tests {
use super::BoundedLru;
#[test]
fn a_hit_returns_the_cached_value_without_evicting_it() {
let mut cache = BoundedLru::new(4, 1024);
cache.insert(1, "one", 10);
assert_eq!(cache.get(&1), Some("one"));
assert_eq!(
cache.get(&1),
Some("one"),
"a hit must not consume the entry"
);
}
#[test]
fn a_miss_returns_none() {
let cache: BoundedLru<i32, &str> = BoundedLru::new(4, 1024);
assert_eq!(cache.get(&1), None);
}
#[test]
fn evicts_the_least_recently_used_entry_once_the_byte_budget_is_exceeded() {
let mut cache = BoundedLru::new(10, 30);
cache.insert(1, "one", 10);
cache.insert(2, "two", 10);
cache.insert(3, "three", 10);
assert_eq!(cache.get(&1), Some("one"));
cache.insert(4, "four", 10);
assert_eq!(cache.get(&2), None, "the untouched entry should be evicted");
assert_eq!(cache.get(&1), Some("one"));
assert_eq!(cache.get(&3), Some("three"));
assert_eq!(cache.get(&4), Some("four"));
}
#[test]
fn evicts_by_entry_count_even_when_under_the_byte_budget() {
let mut cache = BoundedLru::new(2, 1_000_000);
cache.insert(1, "one", 1);
cache.insert(2, "two", 1);
cache.insert(3, "three", 1);
assert_eq!(cache.get(&1), None);
assert_eq!(cache.get(&2), Some("two"));
assert_eq!(cache.get(&3), Some("three"));
}
#[test]
fn an_entry_larger_than_the_whole_budget_is_still_retained_alone() {
let mut cache = BoundedLru::new(4, 100);
cache.insert(1, "big", 500);
assert_eq!(cache.get(&1), Some("big"));
}
#[test]
fn re_inserting_a_key_replaces_rather_than_duplicates_it() {
let mut cache = BoundedLru::new(4, 1024);
cache.insert(1, "old", 10);
cache.insert(1, "new", 10);
assert_eq!(cache.get(&1), Some("new"));
assert_eq!(cache.entries.len(), 1);
}
}
#[cfg(test)]
mod un_transpose_tests {
use super::*;
#[test]
fn un_transpose_matches_reference_algorithm_on_tail_and_exact_cases() {
let in_9_3: Vec<u8> = (0..9).collect();
assert_eq!(
Mf4File::un_transpose(&in_9_3, 3).unwrap(),
vec![0, 3, 6, 1, 4, 7, 2, 5, 8]
);
let in_8_3: Vec<u8> = (0..8).collect();
assert_eq!(
Mf4File::un_transpose(&in_8_3, 3).unwrap(),
vec![0, 2, 4, 1, 3, 5, 6, 7]
);
let in_7_3: Vec<u8> = (0..7).collect();
assert_eq!(
Mf4File::un_transpose(&in_7_3, 3).unwrap(),
vec![0, 2, 4, 1, 3, 5, 6]
);
let in_12_4: Vec<u8> = (0..12).collect();
assert_eq!(
Mf4File::un_transpose(&in_12_4, 4).unwrap(),
vec![0, 3, 6, 9, 1, 4, 7, 10, 2, 5, 8, 11]
);
let in_10_4: Vec<u8> = (0..10).collect();
assert_eq!(
Mf4File::un_transpose(&in_10_4, 4).unwrap(),
vec![0, 2, 4, 6, 1, 3, 5, 7, 8, 9]
);
let in_2_4: Vec<u8> = (0..2).collect();
assert_eq!(Mf4File::un_transpose(&in_2_4, 4).unwrap(), vec![0, 1]);
assert!(Mf4File::un_transpose(&in_9_3, 0).is_err());
}
}