use std::sync::Arc;
use parking_lot::Mutex;
use crate::{Error, McapInfo, McapSummarySource, Summary};
struct SummaryWithSource {
summary: Arc<Summary>,
source: McapSummarySource,
}
pub struct McapFile<BytesSource> {
bytes: BytesSource,
recover: bool,
summary: Mutex<Option<SummaryWithSource>>,
info: Mutex<Option<Arc<McapInfo>>>,
}
impl<BytesSource> McapFile<BytesSource>
where
BytesSource: AsRef<[u8]>,
{
pub fn new(bytes: BytesSource, recover: bool) -> Self {
Self {
bytes,
recover,
summary: Mutex::new(None),
info: Mutex::new(None),
}
}
pub fn bytes(&self) -> &[u8] {
self.bytes.as_ref()
}
pub fn recover(&self) -> bool {
self.recover
}
pub fn cached_summary(&self) -> Option<Arc<Summary>> {
self.summary
.lock()
.as_ref()
.map(|cached| Arc::clone(&cached.summary))
}
pub fn summary(&self) -> Result<Arc<Summary>, Error> {
self.summary_with_source().map(|(summary, _source)| summary)
}
fn summary_with_source(&self) -> Result<(Arc<Summary>, McapSummarySource), Error> {
let mut cached = self.summary.lock();
if let Some(cached) = cached.as_ref() {
return Ok((Arc::clone(&cached.summary), cached.source));
}
let (summary, source) =
crate::recover::read_or_reconstruct_summary_with_source(self.bytes(), self.recover)?;
let summary = Arc::new(summary);
*cached = Some(SummaryWithSource {
summary: Arc::clone(&summary),
source,
});
Ok((summary, source))
}
pub fn info(&self) -> Result<Arc<McapInfo>, Error> {
let mut cached = self.info.lock();
if let Some(info) = cached.as_ref() {
return Ok(Arc::clone(info));
}
let (summary, summary_source) = self.summary_with_source()?;
let header = crate::info::read_header(self.bytes())?;
let info = Arc::new(McapInfo::from_summary(
&header,
&summary,
self.bytes(),
summary_source,
));
*cached = Some(Arc::clone(&info));
Ok(info)
}
}
#[cfg(test)]
mod tests {
use std::io::Cursor;
use super::*;
#[test]
fn summary_and_info_are_initialized_once_across_threads() {
let mut writer = mcap::Writer::new(Cursor::new(Vec::new())).expect("create writer");
writer.finish().expect("finish writer");
let bytes = writer.into_inner().into_inner();
let file = Arc::new(McapFile::new(bytes, false));
let handles = (0..8)
.map(|index| {
let file = Arc::clone(&file);
std::thread::Builder::new()
.name(format!("test-mcap-file-cache-{index}"))
.spawn(move || (file.summary().unwrap(), file.info().unwrap()))
.expect("spawn cache test thread")
})
.collect::<Vec<_>>();
let results = handles
.into_iter()
.map(|handle| handle.join().expect("thread completed"))
.collect::<Vec<_>>();
for (summary, info) in &results[1..] {
assert!(Arc::ptr_eq(&results[0].0, summary));
assert!(Arc::ptr_eq(&results[0].1, info));
}
}
}