use std::collections::BTreeMap;
use crate::Error;
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum McapSummarySource {
Embedded,
Reconstructed,
}
#[derive(Debug, Clone, PartialEq)]
pub struct McapInfo {
pub profile: String,
pub library: String,
pub message_count: Option<u64>,
pub message_start_time_ns: Option<u64>,
pub message_end_time_ns: Option<u64>,
pub duration_ns: Option<u64>,
pub schema_count: usize,
pub channel_count: usize,
pub attachment_count: usize,
pub metadata_count: usize,
pub statistics_present: bool,
pub summary_source: McapSummarySource,
pub chunks: McapChunkInfo,
pub compression: Vec<McapCompressionInfo>,
pub channels: Vec<McapChannelInfo>,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct McapChunkInfo {
pub count: usize,
pub max_uncompressed_size_bytes: Option<u64>,
pub max_compressed_size_bytes: Option<u64>,
pub has_overlapping_time_ranges: bool,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct McapCompressionInfo {
pub codec: String,
pub chunk_count: usize,
pub compressed_size_bytes: u64,
pub uncompressed_size_bytes: u64,
}
impl McapCompressionInfo {
pub fn savings_ratio(&self) -> Option<f64> {
(self.uncompressed_size_bytes > 0)
.then(|| 1.0 - self.compressed_size_bytes as f64 / self.uncompressed_size_bytes as f64)
}
}
#[derive(Debug, Clone, PartialEq)]
pub struct McapChannelInfo {
pub id: u16,
pub topic: String,
pub message_encoding: String,
pub metadata: BTreeMap<String, String>,
pub schema: Option<McapSchemaInfo>,
pub message_count: Option<u64>,
pub frequency_hz: Option<(f64, f64)>,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct McapSchemaInfo {
pub id: u16,
pub name: String,
pub encoding: String,
pub data_size_bytes: usize,
}
pub(crate) fn read_header(mcap: &[u8]) -> Result<mcap::records::Header, Error> {
if !mcap.starts_with(mcap::MAGIC) {
return Err(mcap::McapError::BadMagic.into());
}
let record_start = mcap::MAGIC.len();
let prefix_end = record_start + crate::RECORD_HEADER_LEN;
let Some(prefix) = mcap.get(record_start..prefix_end) else {
return Err(mcap::McapError::UnexpectedEof.into());
};
let opcode = prefix[0];
let record_len = u64::from_le_bytes(prefix[1..].try_into().expect("eight-byte record length"));
let record_len = usize::try_from(record_len)
.map_err(|_err| Error::other(anyhow::anyhow!("MCAP header record is too large")))?;
let body_end = prefix_end
.checked_add(record_len)
.ok_or_else(|| Error::other(anyhow::anyhow!("MCAP header record length overflow")))?;
let Some(body) = mcap.get(prefix_end..body_end) else {
return Err(mcap::McapError::UnexpectedEof.into());
};
match mcap::parse_record(opcode, body)? {
mcap::records::Record::Header(header) => Ok(header),
record => Err(Error::other(anyhow::anyhow!(
"Expected MCAP header as the first record, found opcode 0x{:02x}",
record.opcode()
))),
}
}
impl McapInfo {
pub(crate) fn from_summary(
header: &mcap::records::Header,
summary: &mcap::Summary,
mcap: &[u8],
summary_source: McapSummarySource,
) -> Self {
let statistics_present =
summary_source == McapSummarySource::Embedded && summary.stats.is_some();
let channel_message_counts = channel_message_counts(summary, mcap);
let message_count = summary
.stats
.as_ref()
.map(|stats| stats.message_count)
.or_else(|| {
channel_message_counts
.as_ref()
.map(|counts| counts.values().copied().sum())
});
let (message_start_time_ns, message_end_time_ns) = time_bounds(summary, message_count);
let duration_ns = Option::zip(message_start_time_ns, message_end_time_ns)
.map(|(start, end)| end.saturating_sub(start));
let mut channels = summary
.channels
.values()
.map(|channel| {
let message_count = channel_message_counts
.as_ref()
.and_then(|counts| counts.get(&channel.id).copied());
let frequency_hz = Option::zip(message_count, duration_ns)
.and_then(|(count, duration)| frequency_hz(count, duration));
McapChannelInfo {
id: channel.id,
topic: channel.topic.clone(),
message_encoding: channel.message_encoding.clone(),
metadata: channel.metadata.clone(),
schema: channel.schema.as_ref().map(|schema| McapSchemaInfo {
id: schema.id,
name: schema.name.clone(),
encoding: schema.encoding.clone(),
data_size_bytes: schema.data.len(),
}),
message_count,
frequency_hz,
}
})
.collect::<Vec<_>>();
channels.sort_by_key(|channel| channel.id);
Self {
profile: header.profile.clone(),
library: header.library.clone(),
message_count,
message_start_time_ns,
message_end_time_ns,
duration_ns,
schema_count: summary.schemas.len(),
channel_count: summary.channels.len(),
attachment_count: summary.attachment_indexes.len(),
metadata_count: summary.metadata_indexes.len(),
statistics_present,
summary_source,
chunks: chunk_info(summary),
compression: compression_info(summary),
channels,
}
}
}
fn channel_message_counts(summary: &mcap::Summary, mcap: &[u8]) -> Option<BTreeMap<u16, u64>> {
if let Some(stats) = &summary.stats {
let mut counts = summary
.channels
.keys()
.copied()
.map(|id| (id, 0))
.collect::<BTreeMap<_, _>>();
counts.extend(&stats.channel_message_counts);
return Some(counts);
}
if summary.chunk_indexes.is_empty() {
return None;
}
let mut counts = summary
.channels
.keys()
.copied()
.map(|id| (id, 0))
.collect::<BTreeMap<_, _>>();
for chunk in &summary.chunk_indexes {
let indexes = summary.read_message_indexes(mcap, chunk).ok()?;
for (channel, messages) in indexes {
*counts.entry(channel.id).or_default() += messages.len() as u64;
}
}
Some(counts)
}
fn time_bounds(summary: &mcap::Summary, message_count: Option<u64>) -> (Option<u64>, Option<u64>) {
if message_count == Some(0) {
return (None, None);
}
if let Some(stats) = &summary.stats
&& stats.message_count > 0
{
return (Some(stats.message_start_time), Some(stats.message_end_time));
}
let start = summary
.chunk_indexes
.iter()
.map(|chunk| chunk.message_start_time)
.min();
let end = summary
.chunk_indexes
.iter()
.map(|chunk| chunk.message_end_time)
.max();
(start, end)
}
fn frequency_hz(message_count: u64, duration_ns: u64) -> Option<(f64, f64)> {
if message_count < 2 || duration_ns == 0 {
return None;
}
let duration_seconds = duration_ns as f64 / 1_000_000_000.0;
Some((
(message_count - 1) as f64 / duration_seconds,
message_count as f64 / duration_seconds,
))
}
fn chunk_info(summary: &mcap::Summary) -> McapChunkInfo {
let max_uncompressed_size_bytes = summary
.chunk_indexes
.iter()
.map(|chunk| chunk.uncompressed_size)
.max();
let max_compressed_size_bytes = summary
.chunk_indexes
.iter()
.map(|chunk| chunk.compressed_size)
.max();
let mut chunks = summary.chunk_indexes.iter().collect::<Vec<_>>();
chunks.sort_by_key(|chunk| chunk.message_start_time);
let mut running_end = None;
let mut has_overlapping_time_ranges = false;
for chunk in chunks {
if running_end.is_some_and(|end| chunk.message_start_time < end) {
has_overlapping_time_ranges = true;
break;
}
running_end = Some(running_end.map_or(chunk.message_end_time, |end: u64| {
end.max(chunk.message_end_time)
}));
}
McapChunkInfo {
count: summary.chunk_indexes.len(),
max_uncompressed_size_bytes,
max_compressed_size_bytes,
has_overlapping_time_ranges,
}
}
fn compression_info(summary: &mcap::Summary) -> Vec<McapCompressionInfo> {
let mut by_codec: BTreeMap<String, McapCompressionInfo> = BTreeMap::new();
for chunk in &summary.chunk_indexes {
let info = by_codec
.entry(chunk.compression.clone())
.or_insert_with(|| McapCompressionInfo {
codec: chunk.compression.clone(),
chunk_count: 0,
compressed_size_bytes: 0,
uncompressed_size_bytes: 0,
});
info.chunk_count += 1;
info.compressed_size_bytes = info
.compressed_size_bytes
.saturating_add(chunk.compressed_size);
info.uncompressed_size_bytes = info
.uncompressed_size_bytes
.saturating_add(chunk.uncompressed_size);
}
by_codec.into_values().collect()
}
#[cfg(test)]
mod tests {
use std::io::Cursor;
use mcap::records::MessageHeader;
use super::*;
fn fixture() -> Vec<u8> {
let mut writer = mcap::Writer::with_options(
Cursor::new(Vec::new()),
mcap::WriteOptions::new()
.profile("test-profile")
.library("test-library"),
)
.expect("create writer");
let schema_id = writer
.add_schema("example.Message", "protobuf", b"schema")
.expect("add schema");
let channel_id = writer
.add_channel(schema_id, "/example", "protobuf", &BTreeMap::new())
.expect("add channel");
writer
.add_channel(schema_id, "/empty", "protobuf", &BTreeMap::new())
.expect("add empty channel");
for (sequence, time) in [1_000_000_000, 2_000_000_000, 3_000_000_000]
.into_iter()
.enumerate()
{
writer
.write_to_known_channel(
&MessageHeader {
channel_id,
sequence: sequence as u32,
log_time: time,
publish_time: time,
},
b"message",
)
.expect("write message");
}
writer.finish().expect("finish writer");
writer.into_inner().into_inner()
}
#[test]
fn info_is_derived_from_header_and_summary() {
let bytes = fixture();
let header = read_header(&bytes).expect("read header");
let summary = mcap::Summary::read(&bytes)
.expect("read summary")
.expect("summary present");
let info = McapInfo::from_summary(&header, &summary, &bytes, McapSummarySource::Embedded);
assert_eq!(info.profile, "test-profile");
assert_eq!(info.library, "test-library");
assert_eq!(info.message_count, Some(3));
assert_eq!(info.message_start_time_ns, Some(1_000_000_000));
assert_eq!(info.message_end_time_ns, Some(3_000_000_000));
assert_eq!(info.duration_ns, Some(2_000_000_000));
assert_eq!(info.schema_count, 1);
assert_eq!(info.channel_count, 2);
let example = info
.channels
.iter()
.find(|channel| channel.topic == "/example")
.expect("example channel");
let empty = info
.channels
.iter()
.find(|channel| channel.topic == "/empty")
.expect("empty channel");
assert_eq!(example.message_count, Some(3));
assert_eq!(example.frequency_hz, Some((1.0, 1.5)));
assert_eq!(empty.message_count, Some(0));
assert_eq!(empty.frequency_hz, None);
assert_eq!(
example.schema.as_ref().map(|schema| schema.name.as_str()),
Some("example.Message")
);
}
#[test]
fn empty_file_has_no_time_bounds() {
let mut writer = mcap::Writer::with_options(
Cursor::new(Vec::new()),
mcap::WriteOptions::new()
.profile("empty")
.library("test-library"),
)
.expect("create writer");
writer.finish().expect("finish writer");
let bytes = writer.into_inner().into_inner();
let header = read_header(&bytes).expect("read header");
let summary = mcap::Summary::read(&bytes)
.expect("read summary")
.expect("summary present");
let info = McapInfo::from_summary(&header, &summary, &bytes, McapSummarySource::Embedded);
assert_eq!(info.message_count, Some(0));
assert_eq!(info.message_start_time_ns, None);
assert_eq!(info.message_end_time_ns, None);
assert_eq!(info.duration_ns, None);
assert_eq!(info.chunks.max_compressed_size_bytes, None);
}
}