use concinnity_core::ecs::PayloadLocator;
use concinnity_core::result::CnResult;
enum BlobSlot {
Unloaded(String),
Loaded(Vec<u8>),
Released,
}
pub struct BlobData {
slots: Vec<BlobSlot>,
disk_backed: bool,
}
impl BlobData {
pub fn new(payload_sections: Vec<Option<Vec<u8>>>) -> Self {
let slots = payload_sections
.into_iter()
.map(|s| match s {
Some(bytes) => BlobSlot::Loaded(bytes),
None => BlobSlot::Released,
})
.collect();
Self {
slots,
disk_backed: false,
}
}
pub fn empty() -> Self {
Self {
slots: Vec::new(),
disk_backed: false,
}
}
pub(super) fn from_blob_files(blob0_payload: Vec<u8>, overflow_paths: Vec<String>) -> Self {
let mut slots = Vec::with_capacity(overflow_paths.len() + 1);
slots.push(BlobSlot::Loaded(blob0_payload));
slots.extend(overflow_paths.into_iter().map(BlobSlot::Unloaded));
Self {
slots,
disk_backed: true,
}
}
pub fn disk_backed(&self) -> bool {
self.disk_backed
}
pub fn read(&mut self, locator: &PayloadLocator) -> Result<&[u8], CnResult> {
let idx = locator.blob_index as usize;
let slot = self.slots.get_mut(idx).ok_or_else(|| {
tracing::error!("BlobData: blob {} is out of range", locator.blob_index);
CnResult::FileIo
})?;
if let BlobSlot::Unloaded(path) = slot {
tracing::debug!(
"BlobData: lazily loading overflow blob {}",
locator.blob_index
);
let bytes = super::read_payload_section(&path.clone())?;
*slot = BlobSlot::Loaded(bytes);
}
let section = match &self.slots[idx] {
BlobSlot::Loaded(bytes) => bytes,
BlobSlot::Released => {
tracing::error!("BlobData: blob {} has been released", locator.blob_index);
return Err(CnResult::FileIo);
}
BlobSlot::Unloaded(_) => return Err(CnResult::FileIo),
};
let start = locator.offset as usize;
let end = start.checked_add(locator.len as usize).ok_or_else(|| {
tracing::error!(
"BlobData: payload slice offset {} + len {} overflows in blob {}",
start,
locator.len,
locator.blob_index
);
CnResult::FileIo
})?;
section.get(start..end).ok_or_else(|| {
tracing::error!(
"BlobData: payload slice [{}, {}) out of bounds in blob {} (len={})",
start,
end,
locator.blob_index,
section.len()
);
CnResult::FileIo
})
}
pub fn release(&mut self, blob_index: u32) {
if let Some(slot) = self.slots.get_mut(blob_index as usize)
&& !matches!(slot, BlobSlot::Released)
{
tracing::debug!("BlobData: releasing payload for blob {}", blob_index);
*slot = BlobSlot::Released;
}
}
pub fn release_all_resident(&mut self) -> usize {
let mut freed = 0;
for slot in &mut self.slots {
if let BlobSlot::Loaded(bytes) = slot {
freed += bytes.len();
*slot = BlobSlot::Released;
}
}
freed
}
#[cfg(test)]
pub(crate) fn is_loaded(&self, blob_index: u32) -> bool {
matches!(
self.slots.get(blob_index as usize),
Some(BlobSlot::Loaded(_))
)
}
}
impl concinnity_core::ecs::PayloadStore for BlobData {
fn read(&mut self, locator: &PayloadLocator) -> Result<&[u8], CnResult> {
BlobData::read(self, locator)
}
fn release(&mut self, blob_index: u32) {
BlobData::release(self, blob_index)
}
fn disk_backed(&self) -> bool {
BlobData::disk_backed(self)
}
fn release_all_resident(&mut self) -> usize {
BlobData::release_all_resident(self)
}
}
#[cfg(test)]
mod tests {
use super::*;
use concinnity_core::SCHEMA_VERSION;
use concinnity_core::blob::{BlobMeta, encode_cnb};
fn locator(blob_index: u32, offset: u64, len: u64) -> PayloadLocator {
PayloadLocator {
blob_index,
offset,
len,
}
}
#[test]
fn disk_backed_defaults_false() {
assert!(!BlobData::empty().disk_backed());
assert!(!BlobData::new(vec![Some(vec![1, 2, 3])]).disk_backed());
}
#[test]
fn from_blob_files_is_disk_backed_with_blob0_resident() {
let bd = BlobData::from_blob_files(b"primary".to_vec(), vec!["1".into(), "2".into()]);
assert!(bd.disk_backed());
assert!(bd.is_loaded(0));
assert!(!bd.is_loaded(1));
assert!(!bd.is_loaded(2));
}
#[test]
fn read_lazily_loads_an_unloaded_overflow_blob() {
let dir = tempfile::tempdir().unwrap();
let path = dir.path().join("1").to_string_lossy().into_owned();
let image = encode_cnb(SCHEMA_VERSION, &BlobMeta::default(), b"hello world").unwrap();
std::fs::write(&path, image).expect("write blob");
let mut bd = BlobData::from_blob_files(Vec::new(), vec![path]);
assert!(!bd.is_loaded(1));
assert_eq!(bd.read(&locator(1, 6, 5)).expect("read ok"), b"world");
assert!(bd.is_loaded(1));
}
#[test]
fn read_errors_when_a_deferred_overflow_blob_is_missing() {
let mut bd = BlobData::from_blob_files(Vec::new(), vec!["/nonexistent/cn/1".into()]);
assert_eq!(bd.read(&locator(1, 0, 1)), Err(CnResult::FileIo));
}
#[test]
fn read_errors_on_released_blob() {
let mut bd = BlobData::new(vec![None]);
assert!(bd.read(&locator(0, 0, 1)).is_err());
}
#[test]
fn release_then_read_errors() {
let mut bd = BlobData::new(vec![Some(b"abcd".to_vec())]);
assert_eq!(bd.read(&locator(0, 0, 2)).expect("read ok"), b"ab");
bd.release(0);
assert!(!bd.is_loaded(0));
assert!(bd.read(&locator(0, 0, 2)).is_err());
}
#[test]
fn release_all_resident_frees_loaded_sections() {
let mut bd = BlobData::from_blob_files(b"abcd".to_vec(), vec!["/nonexistent/cn/1".into()]);
assert!(bd.is_loaded(0));
assert!(!bd.is_loaded(1));
let freed = bd.release_all_resident();
assert_eq!(freed, 4, "blob 0's four bytes were freed");
assert!(!bd.is_loaded(0));
assert!(bd.read(&locator(0, 0, 1)).is_err());
assert_eq!(bd.release_all_resident(), 0);
}
#[test]
fn read_errors_on_out_of_range_blob() {
let mut bd = BlobData::empty();
assert!(bd.read(&locator(3, 0, 1)).is_err());
}
#[test]
fn read_errors_when_the_locator_runs_past_the_section() {
let mut bd = BlobData::new(vec![Some(b"abcd".to_vec())]);
assert!(bd.read(&locator(0, 2, 99)).is_err());
assert!(bd.read(&locator(0, u64::MAX, 1)).is_err());
}
}