use std::fs;
use std::io::{Read, Write};
use std::path::{Path, PathBuf};
use std::sync::atomic::{AtomicU64, Ordering};
use rabs_protocol::result_identity::TypedDigest;
use crate::collision_policy::REASON_OBJECT_COLLISION_QUARANTINED;
use crate::digest_set::{DigestError, DigestRequest, StreamingObjectWriter};
use crate::metadata_store::{QuarantineScope, RabsMetadataStore, StoreError, digest_key};
pub const RAW_PROFILE_V1: &str = "raw-v1";
pub const ENCODED_REPRESENTATION_DOMAIN: &str = "rabs.encoded-representation.sha256.v1";
pub type RepresentationDecoder<'a> = &'a dyn Fn(&[u8]) -> Result<Vec<u8>, String>;
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct StoredRepresentationId {
pub logical: TypedDigest,
pub storage_profile: String,
pub encoded_digest: TypedDigest,
pub encoded_size: u64,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub struct DurabilityPolicy {
pub fsync_file: bool,
pub fsync_directory: bool,
}
impl DurabilityPolicy {
pub const FULL: Self = Self {
fsync_file: true,
fsync_directory: true,
};
#[must_use]
pub const fn is_full(self) -> bool {
self.fsync_file && self.fsync_directory
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
pub struct PutLimits {
pub max_logical_bytes: Option<u64>,
pub expected_size: Option<u64>,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum PutError {
LogicalLimitExceeded {
limit: u64,
},
Digest(DigestError),
DeclaredDigestMismatch {
declared: String,
computed: String,
},
CollisionIncident {
digest: String,
existing_path: String,
preserved_incoming_path: String,
},
EncodingDecodeFailed {
profile: String,
error: String,
},
Store(StoreError),
Io {
step: &'static str,
error: String,
},
CrashInjected(FaultPoint),
}
impl From<StoreError> for PutError {
fn from(e: StoreError) -> Self {
Self::Store(e)
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum PutOutcome {
Stored {
path: String,
},
IdempotentDuplicate {
path: String,
},
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum FaultPoint {
StagingWritten,
StagingSynced,
Linked,
DirectorySynced,
MetadataRecorded,
}
#[derive(Debug, Clone)]
pub struct BlobStoreLayout {
root: PathBuf,
}
static PUT_COUNTER: AtomicU64 = AtomicU64::new(0);
impl BlobStoreLayout {
pub fn open(root: &Path) -> Result<Self, PutError> {
let layout = Self {
root: root.to_path_buf(),
};
for dir in [
layout.objects_dir(),
layout.staging_dir(),
layout.quarantine_dir(),
] {
fs::create_dir_all(&dir).map_err(|e| PutError::Io {
step: "create-layout",
error: e.to_string(),
})?;
}
Ok(layout)
}
#[must_use]
pub fn root(&self) -> &Path {
&self.root
}
fn objects_dir(&self) -> PathBuf {
self.root.join("objects")
}
fn staging_dir(&self) -> PathBuf {
self.root.join("staging")
}
fn quarantine_dir(&self) -> PathBuf {
self.root.join("quarantine")
}
#[must_use]
pub fn published_path(
&self,
logical: &TypedDigest,
profile: &str,
encoded: &TypedDigest,
) -> PathBuf {
let logical_hex = hex(&logical.bytes);
let encoded_hex = hex(&encoded.bytes);
self.objects_dir()
.join(&logical_hex[..2])
.join(format!("{logical_hex}.{profile}.{encoded_hex}"))
}
fn fresh_staging_path(&self) -> PathBuf {
let n = PUT_COUNTER.fetch_add(1, Ordering::SeqCst);
self.staging_dir()
.join(format!("put-{}-{n}.tmp", std::process::id()))
}
fn quarantine_path(&self, logical: &TypedDigest) -> PathBuf {
let n = PUT_COUNTER.fetch_add(1, Ordering::SeqCst);
self.quarantine_dir().join(format!(
"incoming-{}-{}-{n}",
hex(&logical.bytes),
std::process::id()
))
}
}
fn hex(bytes: &[u8]) -> String {
use std::fmt::Write as _;
let mut out = String::with_capacity(bytes.len() * 2);
for b in bytes {
let _ = write!(out, "{b:02x}");
}
out
}
pub(crate) fn io_err(step: &'static str) -> impl FnOnce(std::io::Error) -> PutError {
move |e| PutError::Io {
step,
error: e.to_string(),
}
}
fn hash_file_under_domain(path: &Path, domain: &'static str) -> Result<TypedDigest, PutError> {
use sha2::{Digest as _, Sha256};
let mut file = fs::File::open(path).map_err(io_err("open-existing"))?;
let mut hasher = Sha256::new();
hasher.update((domain.len() as u64).to_be_bytes());
hasher.update(domain.as_bytes());
let mut buffer = vec![0_u8; 64 * 1024];
loop {
let n = file.read(&mut buffer).map_err(io_err("read-existing"))?;
if n == 0 {
break;
}
hasher.update(&buffer[..n]);
}
Ok(TypedDigest {
algorithm: rabs_protocol::result_identity::DigestAlgorithm::Sha256V1,
domain,
bytes: hasher.finalize().into(),
})
}
pub(crate) fn recompute_file_digest(path: &Path) -> Result<TypedDigest, PutError> {
let mut file = fs::File::open(path).map_err(io_err("open-existing"))?;
let mut writer = StreamingObjectWriter::new(DigestRequest::default(), None);
let mut buffer = vec![0_u8; 64 * 1024];
loop {
let n = file.read(&mut buffer).map_err(io_err("read-existing"))?;
if n == 0 {
break;
}
writer.write(&buffer[..n]).map_err(PutError::Digest)?;
}
Ok(writer.finish().map_err(PutError::Digest)?.atp_content_id)
}
pub fn put_if_absent(
layout: &BlobStoreLayout,
store: &mut dyn RabsMetadataStore,
declared: &TypedDigest,
reader: &mut dyn Read,
limits: PutLimits,
durability: DurabilityPolicy,
) -> Result<PutOutcome, PutError> {
put_if_absent_with_fault(layout, store, declared, reader, limits, durability, None)
}
#[allow(clippy::too_many_lines)]
pub fn put_if_absent_with_fault(
layout: &BlobStoreLayout,
store: &mut dyn RabsMetadataStore,
declared: &TypedDigest,
reader: &mut dyn Read,
limits: PutLimits,
durability: DurabilityPolicy,
fault: Option<FaultPoint>,
) -> Result<PutOutcome, PutError> {
let staging = layout.fresh_staging_path();
let result = stream_to_staging(&staging, reader, limits);
let digests = match result {
Ok(digests) => digests,
Err(e) => {
let _ = fs::remove_file(&staging);
return Err(e);
}
};
if digests.atp_content_id != *declared {
let computed = digest_key(&digests.atp_content_id);
let _ = fs::remove_file(&staging);
return Err(PutError::DeclaredDigestMismatch {
declared: digest_key(declared),
computed,
});
}
if fault == Some(FaultPoint::StagingWritten) {
return Err(PutError::CrashInjected(FaultPoint::StagingWritten));
}
publish_staged_inner(
layout,
store,
declared,
&staging,
digests.logical_size,
durability,
fault,
RAW_PROFILE_V1,
declared,
)
}
pub fn publish_staged(
layout: &BlobStoreLayout,
store: &mut dyn RabsMetadataStore,
declared: &TypedDigest,
staging: &Path,
durability: DurabilityPolicy,
) -> Result<PutOutcome, PutError> {
let logical_size = fs::metadata(staging).map_err(io_err("stat-staging"))?.len();
publish_staged_inner(
layout,
store,
declared,
staging,
logical_size,
durability,
None,
RAW_PROFILE_V1,
declared,
)
}
#[allow(clippy::too_many_arguments)]
fn publish_staged_inner(
layout: &BlobStoreLayout,
store: &mut dyn RabsMetadataStore,
declared: &TypedDigest,
staging: &Path,
logical_size: u64,
durability: DurabilityPolicy,
fault: Option<FaultPoint>,
profile: &str,
file_digest: &TypedDigest,
) -> Result<PutOutcome, PutError> {
if durability.fsync_file {
let file = fs::File::open(staging).map_err(io_err("open-for-sync"))?;
file.sync_all().map_err(io_err("fsync-staging"))?;
}
if fault == Some(FaultPoint::StagingSynced) {
return Err(PutError::CrashInjected(FaultPoint::StagingSynced));
}
let target = layout.published_path(declared, profile, file_digest);
let target_dir = target
.parent()
.ok_or_else(|| PutError::Io {
step: "target-parent",
error: "published path has no parent".to_owned(),
})?
.to_path_buf();
fs::create_dir_all(&target_dir).map_err(io_err("create-fanout"))?;
match fs::hard_link(staging, &target) {
Ok(()) => {}
Err(e) if e.kind() == std::io::ErrorKind::AlreadyExists => {
return handle_existing(
layout,
store,
declared,
logical_size,
profile,
file_digest,
staging,
&target,
durability,
);
}
Err(e) => {
let _ = fs::remove_file(staging);
return Err(io_err("publish-link")(e));
}
}
if fault == Some(FaultPoint::Linked) {
return Err(PutError::CrashInjected(FaultPoint::Linked));
}
if durability.fsync_directory {
let dir = fs::File::open(&target_dir).map_err(io_err("open-dir"))?;
dir.sync_all().map_err(io_err("fsync-dir"))?;
}
if fault == Some(FaultPoint::DirectorySynced) {
return Err(PutError::CrashInjected(FaultPoint::DirectorySynced));
}
store.record_object(declared, logical_size)?;
store.add_location(
declared,
&target.to_string_lossy(),
None,
profile,
durability.is_full(),
)?;
if fault == Some(FaultPoint::MetadataRecorded) {
return Err(PutError::CrashInjected(FaultPoint::MetadataRecorded));
}
let _ = fs::remove_file(staging);
Ok(PutOutcome::Stored {
path: target.to_string_lossy().into_owned(),
})
}
#[allow(clippy::too_many_arguments)]
pub fn put_encoded_representation(
layout: &BlobStoreLayout,
store: &mut dyn RabsMetadataStore,
declared_logical: &TypedDigest,
profile: &str,
encoded: &mut dyn Read,
decoder: RepresentationDecoder<'_>,
limits: PutLimits,
durability: DurabilityPolicy,
) -> Result<(StoredRepresentationId, PutOutcome), PutError> {
let staging = layout.fresh_staging_path();
let stream = (|| -> Result<(TypedDigest, u64), PutError> {
use sha2::{Digest as _, Sha256};
let mut file = fs::File::create(&staging).map_err(io_err("create-staging"))?;
let mut hasher = Sha256::new();
hasher.update((ENCODED_REPRESENTATION_DOMAIN.len() as u64).to_be_bytes());
hasher.update(ENCODED_REPRESENTATION_DOMAIN.as_bytes());
let mut buffer = vec![0_u8; 64 * 1024];
let mut total: u64 = 0;
loop {
let n = encoded.read(&mut buffer).map_err(io_err("read-encoded"))?;
if n == 0 {
break;
}
total = total.saturating_add(n as u64);
hasher.update(&buffer[..n]);
file.write_all(&buffer[..n])
.map_err(io_err("write-staging"))?;
}
if let Some(expected) = limits.expected_size
&& expected != total
{
return Err(PutError::Digest(DigestError::LogicalSizeMismatch {
expected,
actual: total,
}));
}
Ok((
TypedDigest {
algorithm: rabs_protocol::result_identity::DigestAlgorithm::Sha256V1,
domain: ENCODED_REPRESENTATION_DOMAIN,
bytes: hasher.finalize().into(),
},
total,
))
})();
let (encoded_digest, encoded_size) = match stream {
Ok(v) => v,
Err(e) => {
let _ = fs::remove_file(&staging);
return Err(e);
}
};
let verify = (|| -> Result<u64, PutError> {
let encoded_bytes = fs::read(&staging).map_err(io_err("read-staged"))?;
let logical_bytes =
decoder(&encoded_bytes).map_err(|error| PutError::EncodingDecodeFailed {
profile: profile.to_owned(),
error,
})?;
if let Some(limit) = limits.max_logical_bytes
&& logical_bytes.len() as u64 > limit
{
return Err(PutError::LogicalLimitExceeded { limit });
}
let computed = crate::digest_set::digest_set(
&logical_bytes,
crate::digest_set::DigestRequest::default(),
None,
)
.map_err(PutError::Digest)?
.atp_content_id;
if computed != *declared_logical {
return Err(PutError::DeclaredDigestMismatch {
declared: digest_key(declared_logical),
computed: digest_key(&computed),
});
}
Ok(logical_bytes.len() as u64)
})();
let logical_size = match verify {
Ok(v) => v,
Err(e) => {
let _ = fs::remove_file(&staging);
return Err(e);
}
};
let outcome = publish_staged_inner(
layout,
store,
declared_logical,
&staging,
logical_size,
durability,
None,
profile,
&encoded_digest,
)?;
Ok((
StoredRepresentationId {
logical: declared_logical.clone(),
storage_profile: profile.to_owned(),
encoded_digest,
encoded_size,
},
outcome,
))
}
pub(crate) fn stream_to_staging(
staging: &Path,
reader: &mut dyn Read,
limits: PutLimits,
) -> Result<crate::digest_set::DigestSet, PutError> {
let mut file = fs::File::create(staging).map_err(io_err("create-staging"))?;
let mut writer = StreamingObjectWriter::new(DigestRequest::default(), limits.expected_size);
let mut buffer = vec![0_u8; 64 * 1024];
let mut total: u64 = 0;
loop {
let n = reader.read(&mut buffer).map_err(io_err("read-source"))?;
if n == 0 {
break;
}
total = total.saturating_add(n as u64);
if let Some(limit) = limits.max_logical_bytes
&& total > limit
{
return Err(PutError::LogicalLimitExceeded { limit });
}
writer.write(&buffer[..n]).map_err(PutError::Digest)?;
file.write_all(&buffer[..n])
.map_err(io_err("write-staging"))?;
}
writer.finish().map_err(PutError::Digest)
}
#[allow(clippy::too_many_arguments)]
fn handle_existing(
layout: &BlobStoreLayout,
store: &mut dyn RabsMetadataStore,
declared: &TypedDigest,
logical_size: u64,
profile: &str,
file_digest: &TypedDigest,
staging: &Path,
target: &Path,
durability: DurabilityPolicy,
) -> Result<PutOutcome, PutError> {
let existing_digest = hash_file_under_domain(target, file_digest.domain)?;
if existing_digest == *file_digest {
if durability.is_full() {
let file = fs::File::open(target).map_err(io_err("fsync-existing"))?;
file.sync_all().map_err(io_err("fsync-existing"))?;
if let Some(dir) = target.parent() {
let dir = fs::File::open(dir).map_err(io_err("fsync-existing-dir"))?;
dir.sync_all().map_err(io_err("fsync-existing-dir"))?;
}
}
store.record_object(declared, logical_size)?;
store.add_location(
declared,
&target.to_string_lossy(),
None,
profile,
durability.is_full(),
)?;
let _ = fs::remove_file(staging);
return Ok(PutOutcome::IdempotentDuplicate {
path: target.to_string_lossy().into_owned(),
});
}
let preserved = layout.quarantine_path(declared);
fs::rename(staging, &preserved).map_err(io_err("preserve-incoming"))?;
store.add_quarantine(
QuarantineScope::LogicalObject,
&digest_key(declared),
REASON_OBJECT_COLLISION_QUARANTINED,
)?;
store.set_location_quarantined(declared, &target.to_string_lossy(), true)?;
Err(PutError::CollisionIncident {
digest: digest_key(declared),
existing_path: target.to_string_lossy().into_owned(),
preserved_incoming_path: preserved.to_string_lossy().into_owned(),
})
}
#[cfg(test)]
mod tests {
use super::*;
use crate::digest_set::digest_set;
use crate::metadata_store::{RusqliteEngine, SqlMetadataStore};
use std::sync::atomic::AtomicU64 as TestCounter;
static DIR_COUNTER: TestCounter = TestCounter::new(0);
fn fresh_root(tag: &str) -> PathBuf {
let n = DIR_COUNTER.fetch_add(1, Ordering::SeqCst);
let root = std::env::temp_dir().join(format!("rabs-h003-{}-{tag}-{n}", std::process::id()));
fs::create_dir_all(&root).unwrap();
root
}
fn store() -> SqlMetadataStore<RusqliteEngine> {
SqlMetadataStore::open(RusqliteEngine::open_in_memory().unwrap()).unwrap()
}
fn id_of(bytes: &[u8]) -> TypedDigest {
digest_set(bytes, DigestRequest::default(), None)
.unwrap()
.atp_content_id
}
fn assert_no_partial_published(layout: &BlobStoreLayout) {
let objects = layout.objects_dir();
let mut stack = vec![objects];
while let Some(dir) = stack.pop() {
let Ok(entries) = fs::read_dir(&dir) else {
continue;
};
for entry in entries {
let path = entry.unwrap().path();
if path.is_dir() {
stack.push(path);
continue;
}
let name = path.file_name().unwrap().to_string_lossy().into_owned();
let segments: Vec<&str> = name.split('.').collect();
let [logical_hex, profile, encoded_hex] = segments.as_slice() else {
panic!("unexpected published name {name}");
};
let (domain, claimed) = if *profile == RAW_PROFILE_V1 {
(crate::digest_set::ATP_OBJECT_CONTENT_DOMAIN, *logical_hex)
} else {
(ENCODED_REPRESENTATION_DOMAIN, *encoded_hex)
};
let recomputed = hash_file_under_domain(&path, domain).unwrap();
assert_eq!(
hex(&recomputed.bytes),
claimed,
"partial or corrupt representation exposed at {}",
path.display()
);
}
}
}
#[test]
fn h003_put_stores_then_idempotent_and_staging_clean() {
let layout = BlobStoreLayout::open(&fresh_root("basic")).unwrap();
let mut store = store();
let bytes = b"h003 object bytes".to_vec();
let declared = id_of(&bytes);
let outcome = put_if_absent(
&layout,
&mut store,
&declared,
&mut bytes.as_slice(),
PutLimits {
max_logical_bytes: Some(1024),
expected_size: Some(bytes.len() as u64),
},
DurabilityPolicy::FULL,
)
.unwrap();
let PutOutcome::Stored { path } = outcome else {
panic!("expected Stored, got {outcome:?}");
};
assert_eq!(fs::read(&path).unwrap(), bytes);
assert!(store.object_located(&declared).unwrap());
let again = put_if_absent(
&layout,
&mut store,
&declared,
&mut bytes.as_slice(),
PutLimits::default(),
DurabilityPolicy::FULL,
)
.unwrap();
assert_eq!(again, PutOutcome::IdempotentDuplicate { path });
assert_eq!(
fs::read_dir(layout.staging_dir()).unwrap().count(),
0,
"race loser must clean its temp"
);
assert_no_partial_published(&layout);
}
#[test]
fn h003_refusals_publish_nothing_and_clean_staging() {
let layout = BlobStoreLayout::open(&fresh_root("refusals")).unwrap();
let mut store = store();
let bytes = b"refusal bytes".to_vec();
let declared = id_of(&bytes);
let wrong = id_of(b"other bytes");
assert!(matches!(
put_if_absent(
&layout,
&mut store,
&wrong,
&mut bytes.as_slice(),
PutLimits::default(),
DurabilityPolicy::FULL,
),
Err(PutError::DeclaredDigestMismatch { .. })
));
assert_eq!(
put_if_absent(
&layout,
&mut store,
&declared,
&mut bytes.as_slice(),
PutLimits {
max_logical_bytes: Some(4),
expected_size: None,
},
DurabilityPolicy::FULL,
),
Err(PutError::LogicalLimitExceeded { limit: 4 })
);
assert!(matches!(
put_if_absent(
&layout,
&mut store,
&declared,
&mut bytes.as_slice(),
PutLimits {
max_logical_bytes: None,
expected_size: Some(3),
},
DurabilityPolicy::FULL,
),
Err(PutError::Digest(DigestError::LogicalSizeMismatch { .. }))
));
assert!(!store.object_located(&declared).unwrap());
assert_eq!(fs::read_dir(layout.staging_dir()).unwrap().count(), 0);
assert_no_partial_published(&layout);
}
#[test]
fn duplicate_ingestion_preserves_location_quarantine_and_accepts_clean_replica() {
let layout = BlobStoreLayout::open(&fresh_root("quarantined-duplicate")).unwrap();
let alternate = BlobStoreLayout::open(&fresh_root("clean-replica")).unwrap();
let mut store = store();
let bytes = b"verified bytes on a disk under investigation";
let declared = id_of(bytes);
let PutOutcome::Stored { path } = put_if_absent(
&layout,
&mut store,
&declared,
&mut bytes.as_slice(),
PutLimits::default(),
DurabilityPolicy::FULL,
)
.unwrap() else {
panic!("first ingestion must store a new copy");
};
store
.set_location_quarantined(&declared, &path, true)
.unwrap();
assert_eq!(
put_if_absent(
&layout,
&mut store,
&declared,
&mut bytes.as_slice(),
PutLimits::default(),
DurabilityPolicy::FULL,
)
.unwrap(),
PutOutcome::IdempotentDuplicate { path: path.clone() }
);
assert_eq!(fs::read(&path).unwrap(), bytes);
assert!(!store.object_located(&declared).unwrap());
assert!(!store.object_durably_located(&declared).unwrap());
assert!(store.object_locations(&declared).unwrap().is_empty());
let PutOutcome::Stored {
path: alternate_path,
} = put_if_absent(
&alternate,
&mut store,
&declared,
&mut bytes.as_slice(),
PutLimits::default(),
DurabilityPolicy::FULL,
)
.unwrap()
else {
panic!("alternate storage must receive a new copy");
};
assert!(store.object_durably_located(&declared).unwrap());
assert_eq!(
store.object_locations(&declared).unwrap(),
vec![(alternate_path, RAW_PROFILE_V1.to_owned(), true)]
);
assert!(
store
.reconciliation_scan()
.unwrap()
.iter()
.any(|row| row.store_path == path && row.quarantined)
);
}
#[test]
fn h003_existing_digest_different_bytes_is_quarantined_incident() {
let layout = BlobStoreLayout::open(&fresh_root("collision")).unwrap();
let mut store = store();
let bytes = b"honest object".to_vec();
let declared = id_of(&bytes);
put_if_absent(
&layout,
&mut store,
&declared,
&mut bytes.as_slice(),
PutLimits::default(),
DurabilityPolicy::FULL,
)
.unwrap();
let target = layout.published_path(&declared, RAW_PROFILE_V1, &declared);
fs::write(&target, b"corrupted!").unwrap();
let result = put_if_absent(
&layout,
&mut store,
&declared,
&mut bytes.as_slice(),
PutLimits::default(),
DurabilityPolicy::FULL,
);
let Err(PutError::CollisionIncident {
digest,
existing_path,
preserved_incoming_path,
}) = result
else {
panic!("expected collision incident, got {result:?}");
};
assert_eq!(digest, digest_key(&declared));
assert_eq!(fs::read(&existing_path).unwrap(), b"corrupted!");
assert_eq!(fs::read(&preserved_incoming_path).unwrap(), bytes);
let snapshot = store.differential_snapshot().unwrap();
assert!(
snapshot
.iter()
.any(|l| l.starts_with("quarantines|logical-object|")
&& l.contains(REASON_OBJECT_COLLISION_QUARANTINED))
);
assert!(
store
.reconciliation_scan()
.unwrap()
.iter()
.any(|row| row.quarantined)
);
assert_eq!(fs::read_dir(layout.staging_dir()).unwrap().count(), 0);
}
#[test]
fn h003_crash_injection_at_every_step_never_exposes_a_partial_object() {
for fault in [
FaultPoint::StagingWritten,
FaultPoint::StagingSynced,
FaultPoint::Linked,
FaultPoint::DirectorySynced,
FaultPoint::MetadataRecorded,
] {
let layout = BlobStoreLayout::open(&fresh_root("crash")).unwrap();
let mut store = store();
let bytes = format!("crash object {fault:?}").into_bytes();
let declared = id_of(&bytes);
let result = put_if_absent_with_fault(
&layout,
&mut store,
&declared,
&mut bytes.as_slice(),
PutLimits::default(),
DurabilityPolicy::FULL,
Some(fault),
);
assert_eq!(result, Err(PutError::CrashInjected(fault)));
assert_no_partial_published(&layout);
for row in store.reconciliation_scan().unwrap() {
assert!(
Path::new(&row.store_path).exists(),
"{fault:?}: metadata points at missing path {}",
row.store_path
);
}
let retry = put_if_absent(
&layout,
&mut store,
&declared,
&mut bytes.as_slice(),
PutLimits::default(),
DurabilityPolicy::FULL,
)
.unwrap();
let path = match retry {
PutOutcome::Stored { path } | PutOutcome::IdempotentDuplicate { path } => path,
};
assert_eq!(fs::read(&path).unwrap(), bytes);
assert!(store.object_located(&declared).unwrap());
assert!(
fs::read_dir(layout.staging_dir()).unwrap().count() <= 1,
"{fault:?}: retry left its own staging temp behind"
);
assert_no_partial_published(&layout);
}
}
fn rev_encode(bytes: &[u8]) -> Vec<u8> {
bytes.iter().rev().copied().collect()
}
fn rev_decoder(encoded: &[u8]) -> Result<Vec<u8>, String> {
Ok(encoded.iter().rev().copied().collect())
}
#[test]
fn h030_encoded_and_raw_representations_coexist_without_ambiguity() {
let layout = BlobStoreLayout::open(&fresh_root("h030")).unwrap();
let mut store = store();
let bytes = b"multi-representation object".to_vec();
let declared = id_of(&bytes);
let PutOutcome::Stored { path: raw_path } = put_if_absent(
&layout,
&mut store,
&declared,
&mut bytes.as_slice(),
PutLimits::default(),
DurabilityPolicy::FULL,
)
.unwrap() else {
panic!("raw put must store");
};
let encoded = rev_encode(&bytes);
let (representation, outcome) = put_encoded_representation(
&layout,
&mut store,
&declared,
"rev-v1",
&mut encoded.as_slice(),
&rev_decoder,
PutLimits {
max_logical_bytes: Some(1024),
expected_size: Some(encoded.len() as u64),
},
DurabilityPolicy::FULL,
)
.unwrap();
let PutOutcome::Stored { path: rev_path } = outcome else {
panic!("encoded put must store");
};
assert_ne!(raw_path, rev_path);
assert_eq!(fs::read(&rev_path).unwrap(), encoded);
assert_eq!(representation.storage_profile, "rev-v1");
assert_eq!(representation.encoded_size, encoded.len() as u64);
assert_eq!(
representation.encoded_digest.domain,
ENCODED_REPRESENTATION_DOMAIN
);
assert_eq!(representation.logical, declared);
let encodings: Vec<String> = store
.reconciliation_scan()
.unwrap()
.into_iter()
.map(|row| row.encoding)
.collect();
assert!(encodings.contains(&"raw-v1".to_owned()));
assert!(encodings.contains(&"rev-v1".to_owned()));
let (_, again) = put_encoded_representation(
&layout,
&mut store,
&declared,
"rev-v1",
&mut encoded.as_slice(),
&rev_decoder,
PutLimits::default(),
DurabilityPolicy::FULL,
)
.unwrap();
assert_eq!(again, PutOutcome::IdempotentDuplicate { path: rev_path });
assert_eq!(fs::read_dir(layout.staging_dir()).unwrap().count(), 0);
}
#[test]
fn h030_decode_verification_refuses_wrong_and_bombing_representations() {
let layout = BlobStoreLayout::open(&fresh_root("h030-verify")).unwrap();
let mut store = store();
let bytes = b"verified object".to_vec();
let declared = id_of(&bytes);
let encoded = rev_encode(&bytes);
let bad_decoder = |_: &[u8]| -> Result<Vec<u8>, String> { Ok(b"other bytes".to_vec()) };
assert!(matches!(
put_encoded_representation(
&layout,
&mut store,
&declared,
"rev-v1",
&mut encoded.as_slice(),
&bad_decoder,
PutLimits::default(),
DurabilityPolicy::FULL,
),
Err(PutError::DeclaredDigestMismatch { .. })
));
let failing = |_: &[u8]| -> Result<Vec<u8>, String> { Err("truncated frame".to_owned()) };
assert!(matches!(
put_encoded_representation(
&layout,
&mut store,
&declared,
"rev-v1",
&mut encoded.as_slice(),
&failing,
PutLimits::default(),
DurabilityPolicy::FULL,
),
Err(PutError::EncodingDecodeFailed { .. })
));
let bomb = |_: &[u8]| -> Result<Vec<u8>, String> { Ok(vec![0_u8; 4096]) };
assert_eq!(
put_encoded_representation(
&layout,
&mut store,
&declared,
"bomb-v1",
&mut encoded.as_slice(),
&bomb,
PutLimits {
max_logical_bytes: Some(64),
expected_size: None,
},
DurabilityPolicy::FULL,
),
Err(PutError::LogicalLimitExceeded { limit: 64 })
);
assert!(!store.object_located(&declared).unwrap());
assert_eq!(fs::read_dir(layout.staging_dir()).unwrap().count(), 0);
assert_no_partial_published(&layout);
}
#[test]
fn h030_concurrent_profiles_publish_distinct_records_not_one_race() {
let layout = BlobStoreLayout::open(&fresh_root("h030-race")).unwrap();
let bytes = b"contended multi-encoding object".to_vec();
let declared = id_of(&bytes);
let workers: Vec<_> = (0..8)
.map(|i| {
let layout = layout.clone();
let bytes = bytes.clone();
let declared = declared.clone();
std::thread::spawn(move || {
let mut store = store();
if i % 2 == 0 {
put_if_absent(
&layout,
&mut store,
&declared,
&mut bytes.as_slice(),
PutLimits::default(),
DurabilityPolicy::FULL,
)
.map(|outcome| ("raw-v1", outcome))
} else {
let encoded = rev_encode(&bytes);
put_encoded_representation(
&layout,
&mut store,
&declared,
"rev-v1",
&mut encoded.as_slice(),
&rev_decoder,
PutLimits::default(),
DurabilityPolicy::FULL,
)
.map(|(_, outcome)| ("rev-v1", outcome))
}
})
})
.collect();
let outcomes: Vec<(&str, PutOutcome)> = workers
.into_iter()
.map(|t| t.join().unwrap().unwrap())
.collect();
for profile in ["raw-v1", "rev-v1"] {
let stored = outcomes
.iter()
.filter(|(p, o)| *p == profile && matches!(o, PutOutcome::Stored { .. }))
.count();
assert_eq!(stored, 1, "profile {profile}");
}
assert_eq!(fs::read_dir(layout.staging_dir()).unwrap().count(), 0);
assert_no_partial_published(&layout);
}
#[test]
fn h030_representation_ops_never_touch_action_keys() {
let layout = BlobStoreLayout::open(&fresh_root("h030-keys")).unwrap();
let mut store = store();
let bytes = b"identity-stable object".to_vec();
let declared = id_of(&bytes);
let before: Vec<String> = store
.differential_snapshot()
.unwrap()
.into_iter()
.filter(|l| l.starts_with("action_"))
.collect();
put_if_absent(
&layout,
&mut store,
&declared,
&mut bytes.as_slice(),
PutLimits::default(),
DurabilityPolicy::FULL,
)
.unwrap();
let encoded = rev_encode(&bytes);
put_encoded_representation(
&layout,
&mut store,
&declared,
"rev-v1",
&mut encoded.as_slice(),
&rev_decoder,
PutLimits::default(),
DurabilityPolicy::FULL,
)
.unwrap();
let after: Vec<String> = store
.differential_snapshot()
.unwrap()
.into_iter()
.filter(|l| l.starts_with("action_"))
.collect();
assert_eq!(before, after);
}
#[test]
fn h003_concurrent_writers_one_stores_rest_verify_idempotent() {
let layout = BlobStoreLayout::open(&fresh_root("race")).unwrap();
let bytes = b"contended object".to_vec();
let declared = id_of(&bytes);
let workers: Vec<_> = (0..8)
.map(|_| {
let layout = layout.clone();
let bytes = bytes.clone();
let declared = declared.clone();
std::thread::spawn(move || {
let mut store = store();
put_if_absent(
&layout,
&mut store,
&declared,
&mut bytes.as_slice(),
PutLimits::default(),
DurabilityPolicy::FULL,
)
})
})
.collect();
let outcomes: Vec<_> = workers
.into_iter()
.map(|t| t.join().unwrap().unwrap())
.collect();
let stored = outcomes
.iter()
.filter(|o| matches!(o, PutOutcome::Stored { .. }))
.count();
assert_eq!(stored, 1, "exactly one writer wins the publish");
assert_eq!(outcomes.len(), 8);
let target = layout.published_path(&declared, RAW_PROFILE_V1, &declared);
assert_eq!(fs::read(&target).unwrap(), bytes);
assert_eq!(fs::read_dir(layout.staging_dir()).unwrap().count(), 0);
assert_no_partial_published(&layout);
}
}