use std::sync::Arc;
use async_trait::async_trait;
use dig_dht::ContentId;
use dig_rpc_protocol::types::ModuleInfo;
use sha2::{Digest, Sha256};
use crate::error::{
hex64_or_sentinel, sanitize_untrusted_text, DownloadError, VerifyError, MAX_ERROR_REASON_CHARS,
};
use crate::locate::ProviderLocator;
use crate::progress::{DownloadState, StateStore};
use crate::sink::{promote_verified, Sink};
#[async_trait]
pub trait ModuleTransport: Send + Sync {
async fn get_module_info(
&self,
provider_peer_id: &str,
store_id: &str,
root: &str,
) -> Result<ModuleInfo, DownloadError>;
async fn fetch_module_range(
&self,
provider_peer_id: &str,
store_id: &str,
root: &str,
offset: u64,
length: u64,
) -> Result<Vec<u8>, DownloadError>;
}
#[async_trait]
pub trait ModuleReader: Send + Sync {
fn len(&self) -> u64;
fn is_empty(&self) -> bool {
self.len() == 0
}
async fn read_at(&self, offset: u64, len: u64) -> Result<Vec<u8>, DownloadError>;
}
#[async_trait]
pub trait ModuleAnchorVerifier: Send + Sync {
async fn verify_module_anchor(
&self,
module: &dyn ModuleReader,
store_id: &str,
root: &str,
) -> ModuleAnchor;
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum ModuleAnchor {
Anchored,
NotAnchored,
Unavailable(String),
}
#[cfg(any(test, feature = "testkit"))]
#[doc(hidden)]
#[derive(Debug, Clone, Copy, Default)]
pub struct AcceptAnyModuleAnchor;
#[cfg(any(test, feature = "testkit"))]
#[async_trait]
impl ModuleAnchorVerifier for AcceptAnyModuleAnchor {
async fn verify_module_anchor(
&self,
_module: &dyn ModuleReader,
_store_id: &str,
_root: &str,
) -> ModuleAnchor {
ModuleAnchor::Anchored
}
}
pub const DEFAULT_MAX_MODULE_SIZE: u64 = 512 * 1024 * 1024;
pub const MAX_MODULE_CHUNK_COUNT: usize = 1024 * 1024;
pub const MAX_DESCRIPTOR_ATTEMPTS: usize = 3;
#[derive(Debug, Clone)]
pub struct ModuleDownloadConfig {
pub range_timeout: std::time::Duration,
pub max_module_size: u64,
}
impl Default for ModuleDownloadConfig {
fn default() -> Self {
ModuleDownloadConfig {
range_timeout: std::time::Duration::from_secs(30),
max_module_size: DEFAULT_MAX_MODULE_SIZE,
}
}
}
pub struct ModuleDownloader {
locator: Arc<dyn ProviderLocator>,
transport: Arc<dyn ModuleTransport>,
anchor: Arc<dyn ModuleAnchorVerifier>,
state_store: Arc<dyn StateStore>,
config: ModuleDownloadConfig,
}
impl ModuleDownloader {
pub fn new(
locator: Arc<dyn ProviderLocator>,
transport: Arc<dyn ModuleTransport>,
anchor: Arc<dyn ModuleAnchorVerifier>,
state_store: Arc<dyn StateStore>,
config: ModuleDownloadConfig,
) -> Self {
ModuleDownloader {
locator,
transport,
anchor,
state_store,
config,
}
}
pub async fn download(
&self,
store_id: &str,
root: &str,
sink: &dyn Sink,
) -> Result<u64, DownloadError> {
let content = module_content_id(store_id, root).ok_or(DownloadError::NotDownloadable)?;
let mut providers = self.locator.find_providers(&content).await?;
if providers.is_empty() {
return Err(DownloadError::NotFound {
content: module_download_key(store_id, root),
});
}
let key = module_download_key(store_id, root);
let mut exclusions = DescriptorExclusions::new(self.remembered_verdicts(&key).await);
let mut attempts = 0usize;
loop {
if exclusions.usable_holders(&providers) == 0 && exclusions.stop_trusting_memory() {
tracing::warn!(
holders = providers.len(),
"no holder is left to ask for a descriptor once past verdicts are honoured; \
re-asking every holder rather than letting reputation deny the pull"
);
continue;
}
let (source, info) = self
.fetch_module_info(&providers, &exclusions.excluded(), store_id, root)
.await?;
attempts += 1;
let failure = match self
.pull_with_descriptor(&info, store_id, root, sink, &mut providers)
.await
{
Ok(len) => return Ok(len),
Err(PullFailure::Terminal(e)) => return Err(e),
Err(failure) => failure,
};
let proven_false = failure.is_proven_false();
tracing::warn!(
peer = %hex64_or_sentinel(&source, "peer-id"),
error = %failure.error(),
proven_false,
"module pull: descriptor attempt failed; demoting this source and re-handshaking \
with another holder"
);
if proven_false {
if let Err(store_err) = self.state_store.record_bad_descriptor(&key, &source).await
{
tracing::debug!(error = %store_err, "could not persist a bad-descriptor verdict");
}
}
exclusions.demote(source);
if attempts >= MAX_DESCRIPTOR_ATTEMPTS
|| (exclusions.usable_holders(&providers) == 0 && !exclusions.trusts_memory())
{
return Err(match failure {
PullFailure::BadDescriptor(e) | PullFailure::UnsatisfiableDescriptor(e) => e,
PullFailure::Terminal(e) => e,
});
}
self.state_store.clear(&key).await?;
sink.truncate(0).await?;
}
}
async fn pull_with_descriptor(
&self,
info: &ModuleInfo,
store_id: &str,
root: &str,
sink: &dyn Sink,
providers: &mut Vec<dig_dht::ProviderRecord>,
) -> Result<u64, PullFailure> {
let layout = ChunkPlan::from_info(info, self.config.max_module_size)?;
let key = module_download_key(store_id, root);
let Resume {
mut state,
resumes_staging,
} = self.load_or_fresh_state(&key, &layout).await?;
if !resumes_staging {
sink.truncate(0).await?;
}
let checkpointed = std::mem::take(&mut state.done_ranges);
let mut hasher = Sha256::new();
let mut any_chunk_verified = false;
for index in 0..layout.chunk_count() {
let (offset, len) = layout.chunk_span(index);
let staged = if checkpointed.contains(&index) {
self.read_back_verified_chunk(sink, info, index, offset, len)
.await
} else {
None
};
let bytes = match staged {
Some(bytes) => bytes,
None => {
let bytes = match self
.fetch_verified_chunk(providers, info, &layout, index, store_id, root)
.await
{
Ok(bytes) => bytes,
Err(e)
if e.is_recoverable()
|| matches!(e, DownloadError::NotFound { .. }) =>
{
return Err(PullFailure::UnsatisfiableDescriptor(
describe_chunk_exhaustion(e, any_chunk_verified),
))
}
Err(e) => return Err(PullFailure::Terminal(e)),
};
sink.write_at(offset, &bytes).await?;
state.mark_done(index);
self.state_store.save(&state).await?;
bytes
}
};
any_chunk_verified = true;
hasher.update(&bytes);
state.mark_done(index);
}
let assembled_hash = hex_of(hasher.finalize());
if assembled_hash != info.module_hash {
return Err(PullFailure::BadDescriptor(DownloadError::Verify(
VerifyError::Metadata(format!(
"assembled module_hash {assembled_hash} != declared {}",
hex64_or_sentinel(&info.module_hash, "module-hash")
)),
)));
}
if !sink.supports_read_back() {
return Err(PullFailure::Terminal(DownloadError::sink(
"this sink cannot read back its staged bytes, so the chain-anchor gate has nothing \
to read and the module could never be promoted; implement Sink::read_at + \
Sink::supports_read_back",
)));
}
let reader = StagedModuleReader::new(sink, &layout, &info.chunk_hashes);
match self
.anchor
.verify_module_anchor(&reader, store_id, root)
.await
{
ModuleAnchor::Anchored => {}
ModuleAnchor::NotAnchored => {
return Err(PullFailure::BadDescriptor(DownloadError::Verify(
VerifyError::Metadata(format!(
"assembled module is not chain-anchored under ({store_id}, {root})"
)),
)))
}
ModuleAnchor::Unavailable(reason) => {
return Err(PullFailure::Terminal(DownloadError::state(format!(
"cannot verify the chain anchor for ({store_id}, {root}): {}",
sanitize_untrusted_text(&reason, MAX_ERROR_REASON_CHARS)
))))
}
}
if let Err(e) = promote_verified(sink, layout.total_size).await {
let _ = self.state_store.clear(&key).await;
let _ = sink.truncate(0).await;
return Err(PullFailure::Terminal(e));
}
self.state_store.clear(&key).await?;
Ok(layout.total_size)
}
async fn remembered_verdicts(&self, key: &str) -> Vec<String> {
match self.state_store.bad_descriptor_peers(key).await {
Ok(peers) => peers,
Err(e) => {
tracing::debug!(error = %e, "could not read remembered descriptor verdicts");
Vec::new()
}
}
}
async fn fetch_module_info(
&self,
providers: &[dig_dht::ProviderRecord],
excluded: &[String],
store_id: &str,
root: &str,
) -> Result<(String, ModuleInfo), DownloadError> {
let mut reasons = HolderReasons::default();
let mut tried = 0usize;
for provider in providers {
let peer = &provider.provider_peer_id;
if excluded.iter().any(|d| d == peer) {
continue; }
tried += 1;
match self.transport.get_module_info(peer, store_id, root).await {
Ok(info) => return Ok((peer.clone(), info)),
Err(e) if e.is_recoverable() => reasons.record(peer, e),
Err(e) => return Err(e),
}
}
Err(DownloadError::NotFound {
content: format!(
"getModuleInfo failed on all {tried} usable holder(s) ({} demoted) for module {} — \
{reasons}",
excluded.len(),
module_download_key(store_id, root),
),
})
}
async fn fetch_verified_chunk(
&self,
providers: &mut Vec<dig_dht::ProviderRecord>,
info: &ModuleInfo,
layout: &ChunkPlan,
index: usize,
store_id: &str,
root: &str,
) -> Result<Vec<u8>, DownloadError> {
let (offset, len) = layout.chunk_span(index);
let expected_hash = &info.chunk_hashes[index];
let mut reasons = HolderReasons::default();
let mut relocated = false;
loop {
let count = providers.len();
for step in 0..count {
let peer = providers[(index + step) % count].provider_peer_id.clone();
match self
.fetch_chunk_from(&peer, store_id, root, offset, len, expected_hash)
.await
{
Ok(bytes) => return Ok(bytes),
Err(reason) => reasons.record(&peer, reason),
}
}
if relocated {
return Err(DownloadError::NotFound {
content: format!(
"fetchModuleRange failed for chunk {index} ([{offset}, {}) of module {}) on \
all {count} known holder(s) — {reasons}",
offset + len,
module_download_key(store_id, root),
),
});
}
let content =
module_content_id(store_id, root).ok_or(DownloadError::NotDownloadable)?;
let refreshed = self.locator.find_providers(&content).await?;
merge_new_providers(providers, refreshed);
relocated = true;
}
}
async fn fetch_chunk_from(
&self,
peer: &str,
store_id: &str,
root: &str,
offset: u64,
len: u64,
expected_hash: &str,
) -> Result<Vec<u8>, String> {
let fetched = tokio::time::timeout(
self.config.range_timeout,
self.transport
.fetch_module_range(peer, store_id, root, offset, len),
)
.await;
let mut bytes = match fetched {
Ok(Ok(bytes)) => bytes,
Ok(Err(e)) => return Err(format!("transport: {e}")),
Err(_) => return Err(format!("timed out after {:?}", self.config.range_timeout)),
};
if bytes.len() as u64 > len {
bytes.truncate(len as usize); }
if bytes.len() as u64 != len {
return Err(format!(
"short range: wanted {len} bytes, got {}",
bytes.len()
));
}
if sha256_hex(&bytes) != expected_hash {
return Err("chunk hash mismatch".to_string());
}
Ok(bytes)
}
async fn load_or_fresh_state(
&self,
key: &str,
layout: &ChunkPlan,
) -> Result<Resume, DownloadError> {
let fresh = || {
let mut s = DownloadState::new(key);
s.total_length = layout.total_size;
s.chunk_lens = layout.chunk_lens.clone();
Resume {
state: s,
resumes_staging: false,
}
};
match self.state_store.load(key).await? {
Some(prev) if prev.chunk_lens == layout.chunk_lens => Ok(Resume {
state: prev,
resumes_staging: true,
}),
_ => Ok(fresh()),
}
}
async fn read_back_verified_chunk(
&self,
sink: &dyn Sink,
info: &ModuleInfo,
index: usize,
offset: u64,
len: u64,
) -> Option<Vec<u8>> {
let bytes = sink.read_at(offset, len).await.ok()?;
if bytes.len() as u64 != len || sha256_hex(&bytes) != info.chunk_hashes[index] {
tracing::warn!(
chunk = index,
offset,
"staged chunk failed re-attribution on resume; re-fetching"
);
return None;
}
Some(bytes)
}
}
struct StagedModuleReader<'a> {
sink: &'a dyn Sink,
layout: &'a ChunkPlan,
chunk_hashes: &'a [String],
}
impl<'a> StagedModuleReader<'a> {
fn new(sink: &'a dyn Sink, layout: &'a ChunkPlan, chunk_hashes: &'a [String]) -> Self {
StagedModuleReader {
sink,
layout,
chunk_hashes,
}
}
async fn verified_chunk(&self, index: usize) -> Result<Vec<u8>, DownloadError> {
let (offset, len) = self.layout.chunk_span(index);
let bytes = self.sink.read_at(offset, len).await?;
if bytes.len() as u64 != len {
return Err(DownloadError::sink(format!(
"staged chunk {index} reads {} bytes, expected {len}",
bytes.len()
)));
}
if sha256_hex(&bytes) != self.chunk_hashes[index] {
return Err(DownloadError::sink(format!(
"staged chunk {index} no longer matches its verified hash"
)));
}
Ok(bytes)
}
}
#[async_trait]
impl ModuleReader for StagedModuleReader<'_> {
fn len(&self) -> u64 {
self.layout.total_size
}
async fn read_at(&self, offset: u64, len: u64) -> Result<Vec<u8>, DownloadError> {
if len == 0 {
return Ok(Vec::new());
}
let end = offset.checked_add(len).filter(|e| *e <= self.len());
let Some(end) = end else {
return Err(DownloadError::sink(format!(
"read [{offset}, {offset}+{len}) falls outside the {}-byte module",
self.len()
)));
};
let mut out = Vec::with_capacity(usize::try_from(len).map_err(|_| {
DownloadError::sink(format!(
"read of {len} bytes exceeds this platform's address space"
))
})?);
let mut index = self
.layout
.offsets
.partition_point(|&start| start <= offset)
.saturating_sub(1);
while (out.len() as u64) < len {
if index >= self.layout.chunk_count() {
return Err(DownloadError::sink(format!(
"the staged chunk plan does not cover [{offset}, {end})"
)));
}
let (chunk_offset, chunk_len) = self.layout.chunk_span(index);
if chunk_len == 0 {
index += 1;
continue;
}
let chunk = self.verified_chunk(index).await?;
let want_from = offset + out.len() as u64;
let start = usize::try_from(want_from - chunk_offset).unwrap_or(usize::MAX);
let take = chunk
.len()
.saturating_sub(start)
.min(usize::try_from(len - out.len() as u64).unwrap_or(usize::MAX));
out.extend_from_slice(&chunk[start..start + take]);
index += 1;
}
Ok(out)
}
}
struct Resume {
state: DownloadState,
resumes_staging: bool,
}
enum PullFailure {
BadDescriptor(DownloadError),
UnsatisfiableDescriptor(DownloadError),
Terminal(DownloadError),
}
struct DescriptorExclusions {
demoted_here: Vec<String>,
remembered: Vec<String>,
trusts_memory: bool,
}
impl DescriptorExclusions {
fn new(remembered: Vec<String>) -> Self {
let trusts_memory = !remembered.is_empty();
DescriptorExclusions {
demoted_here: Vec::new(),
remembered,
trusts_memory,
}
}
fn excluded(&self) -> Vec<String> {
let mut excluded = self.demoted_here.clone();
if self.trusts_memory {
excluded.extend(self.remembered.iter().cloned());
}
excluded
}
fn usable_holders(&self, providers: &[dig_dht::ProviderRecord]) -> usize {
let excluded = self.excluded();
providers
.iter()
.filter(|p| !excluded.contains(&p.provider_peer_id))
.count()
}
fn demote(&mut self, peer: String) {
self.demoted_here.push(peer);
}
fn stop_trusting_memory(&mut self) -> bool {
let was_trusting = self.trusts_memory;
self.trusts_memory = false;
was_trusting && !self.remembered.is_empty()
}
fn trusts_memory(&self) -> bool {
self.trusts_memory
}
}
fn describe_chunk_exhaustion(e: DownloadError, any_chunk_verified: bool) -> DownloadError {
let diagnosis = if any_chunk_verified {
"some chunk(s) had already verified under this descriptor, so the missing bytes are more \
likely genuinely unavailable than fabricated"
} else {
"no chunk ever verified under this descriptor, so it is more likely fabricated than the \
bytes unavailable"
};
match e {
DownloadError::NotFound { content } => DownloadError::NotFound {
content: format!("{content} — {diagnosis}"),
},
other => other,
}
}
impl PullFailure {
fn error(&self) -> &DownloadError {
match self {
PullFailure::BadDescriptor(e)
| PullFailure::UnsatisfiableDescriptor(e)
| PullFailure::Terminal(e) => e,
}
}
fn is_proven_false(&self) -> bool {
matches!(self, PullFailure::BadDescriptor(_))
}
}
impl From<DownloadError> for PullFailure {
fn from(e: DownloadError) -> Self {
PullFailure::Terminal(e)
}
}
#[derive(Debug, Default)]
struct HolderReasons(Vec<String>);
impl HolderReasons {
fn record(&mut self, peer: &str, reason: impl std::fmt::Display) {
let peer = hex64_or_sentinel(peer, "peer-id");
let reason = sanitize_untrusted_text(&reason.to_string(), MAX_ERROR_REASON_CHARS);
tracing::debug!(%peer, %reason, "module pull: holder rejected");
self.0.push(format!("{peer}: {reason}"));
}
}
impl std::fmt::Display for HolderReasons {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
if self.0.is_empty() {
return f.write_str("no holder reasons recorded");
}
write!(f, "reasons: [{}]", self.0.join("; "))
}
}
#[derive(Debug)]
struct ChunkPlan {
total_size: u64,
chunk_lens: Vec<u64>,
offsets: Vec<u64>,
}
impl ChunkPlan {
fn from_info(info: &ModuleInfo, max_module_size: u64) -> Result<Self, PullFailure> {
let false_descriptor = |reason: String| {
PullFailure::BadDescriptor(DownloadError::Verify(VerifyError::Metadata(reason)))
};
if info.total_size > max_module_size {
return Err(false_descriptor(format!(
"declared module total_size {} exceeds the maximum {max_module_size}",
info.total_size
)));
}
if info.chunk_lens.is_empty() {
return Err(false_descriptor(
"ModuleInfo carries no chunk_lens (cannot map ranges to chunk hashes)".into(),
));
}
if info.chunk_lens.len() > MAX_MODULE_CHUNK_COUNT {
return Err(false_descriptor(format!(
"declared chunk_lens count {} exceeds the maximum {MAX_MODULE_CHUNK_COUNT}",
info.chunk_lens.len()
)));
}
if info.chunk_lens.len() != info.chunk_hashes.len() {
return Err(false_descriptor(format!(
"chunk_lens ({}) != chunk_hashes ({})",
info.chunk_lens.len(),
info.chunk_hashes.len()
)));
}
let sum = info
.chunk_lens
.iter()
.try_fold(0u64, |acc, &len| acc.checked_add(len))
.ok_or_else(|| {
false_descriptor("chunk_lens sum overflows u64 (hostile descriptor)".into())
})?;
if sum != info.total_size {
return Err(false_descriptor(format!(
"chunk_lens sum {sum} != total_size {}",
info.total_size
)));
}
let chunk_lens = info.chunk_lens.clone();
let mut offsets = Vec::new();
offsets.try_reserve_exact(chunk_lens.len()).map_err(|e| {
PullFailure::UnsatisfiableDescriptor(DownloadError::sink(format!(
"this host cannot allocate the {}-entry chunk plan this descriptor declares: {e}",
chunk_lens.len()
)))
})?;
let mut acc = 0u64;
for &len in &chunk_lens {
offsets.push(acc);
acc = acc.checked_add(len).ok_or_else(|| {
false_descriptor("chunk offsets overflow u64 (hostile descriptor)".into())
})?;
}
Ok(ChunkPlan {
total_size: info.total_size,
chunk_lens,
offsets,
})
}
fn chunk_count(&self) -> usize {
self.chunk_lens.len()
}
fn chunk_span(&self, index: usize) -> (u64, u64) {
(self.offsets[index], self.chunk_lens[index])
}
}
fn merge_new_providers(
known: &mut Vec<dig_dht::ProviderRecord>,
fresh: Vec<dig_dht::ProviderRecord>,
) {
for p in fresh {
if !known
.iter()
.any(|k| k.provider_peer_id == p.provider_peer_id)
{
known.push(p);
}
}
}
fn sha256_hex(bytes: &[u8]) -> String {
hex_of(Sha256::digest(bytes))
}
pub(crate) fn hex_of(digest: impl AsRef<[u8]>) -> String {
let mut out = String::with_capacity(64);
for b in digest.as_ref() {
out.push(char::from_digit((b >> 4) as u32, 16).unwrap());
out.push(char::from_digit((b & 0x0f) as u32, 16).unwrap());
}
out
}
pub fn module_download_key(store_id: &str, root: &str) -> String {
format!("module:{store_id}:{root}")
}
pub fn module_content_id(store_id: &str, root: &str) -> Option<ContentId> {
Some(ContentId::root(hex32(store_id)?, hex32(root)?))
}
fn hex32(s: &str) -> Option<[u8; 32]> {
if s.len() != 64 {
return None;
}
let mut out = [0u8; 32];
for (i, byte) in out.iter_mut().enumerate() {
*byte = u8::from_str_radix(&s[i * 2..i * 2 + 2], 16).ok()?;
}
Some(out)
}
#[cfg(test)]
mod tests {
use super::*;
use crate::progress::InMemoryStateStore;
use crate::sink::InMemorySink;
use crate::testkit::{
mock_providers, MockModuleTransport, MockProviderLocator, RejectAllModuleAnchor,
};
use std::collections::BTreeSet;
fn hex_id(byte: u8) -> String {
format!("{byte:02x}").repeat(32)
}
fn locator_with(n: u8, store_id: &str, root: &str) -> Arc<MockProviderLocator> {
let content = module_content_id(store_id, root).unwrap();
Arc::new(MockProviderLocator::fixed(mock_providers(n, &content)))
}
#[tokio::test]
async fn happy_path_assembles_verified_module_from_multiple_sources() {
let store_id = hex_id(0x11);
let root = hex_id(0x22);
let module = b"the whole .dig module blob".to_vec();
let transport = Arc::new(MockModuleTransport::serving(
&store_id,
&root,
module.clone(),
8,
));
let downloader = ModuleDownloader::new(
locator_with(3, &store_id, &root),
transport.clone(),
Arc::new(AcceptAnyModuleAnchor),
Arc::new(InMemoryStateStore::new()),
ModuleDownloadConfig::default(),
);
let sink = InMemorySink::new();
let len = downloader
.download(&store_id, &root, &sink)
.await
.expect("pull succeeds");
assert_eq!(len, module.len() as u64);
assert_eq!(
sink.contents().await,
module,
"reassembled blob is byte-exact"
);
assert!(sink.is_finalized().await, "verified module is finalized");
let distinct: BTreeSet<String> = transport
.fetches()
.await
.into_iter()
.map(|(p, _)| p)
.collect();
assert!(
distinct.len() > 1,
"chunks came from multiple holders: {distinct:?}"
);
}
#[tokio::test]
async fn resume_after_interrupt_refetches_only_missing_chunks() {
let store_id = hex_id(0x33);
let root = hex_id(0x44);
let module = (0u8..40).collect::<Vec<u8>>(); let state_store = Arc::new(InMemoryStateStore::new());
let sink = InMemorySink::new();
let interrupted = Arc::new(
MockModuleTransport::serving(&store_id, &root, module.clone(), 8)
.with_success_budget(2),
);
let first = ModuleDownloader::new(
locator_with(1, &store_id, &root),
interrupted.clone(),
Arc::new(AcceptAnyModuleAnchor),
state_store.clone(),
ModuleDownloadConfig::default(),
);
let err = first
.download(&store_id, &root, &sink)
.await
.expect_err("interrupted pull fails before finalize");
assert!(
matches!(err, DownloadError::NotFound { .. }),
"exhaustion is terminal and names its step: {err}"
);
assert!(
!sink.is_finalized().await,
"an incomplete pull is never finalized"
);
let healthy = Arc::new(MockModuleTransport::serving(
&store_id,
&root,
module.clone(),
8,
));
let second = ModuleDownloader::new(
locator_with(1, &store_id, &root),
healthy.clone(),
Arc::new(AcceptAnyModuleAnchor),
state_store,
ModuleDownloadConfig::default(),
);
let len = second
.download(&store_id, &root, &sink)
.await
.expect("resumed pull succeeds");
assert_eq!(len, module.len() as u64);
assert_eq!(sink.contents().await, module);
assert!(sink.is_finalized().await);
let resumed_offsets: BTreeSet<u64> = healthy
.fetches()
.await
.into_iter()
.map(|(_, o)| o)
.collect();
assert!(!resumed_offsets.contains(&0), "chunk 0 not re-fetched");
assert!(!resumed_offsets.contains(&8), "chunk 1 not re-fetched");
assert_eq!(
resumed_offsets,
BTreeSet::from([16, 24, 32]),
"only missing chunks fetched"
);
}
#[tokio::test]
async fn tampered_range_is_rejected_and_routed_around() {
let store_id = hex_id(0x55);
let root = hex_id(0x66);
let module = b"honest bytes across several chunks here".to_vec();
let peer1 = crate::testkit::mock_peer_hex(1);
let transport = Arc::new(
MockModuleTransport::serving(&store_id, &root, module.clone(), 8).tampering(&peer1),
);
let downloader = ModuleDownloader::new(
locator_with(3, &store_id, &root),
transport,
Arc::new(AcceptAnyModuleAnchor),
Arc::new(InMemoryStateStore::new()),
ModuleDownloadConfig::default(),
);
let sink = InMemorySink::new();
let len = downloader
.download(&store_id, &root, &sink)
.await
.expect("pull recovers around the tampering holder");
assert_eq!(
sink.contents().await,
module,
"only honest bytes were accepted"
);
assert_eq!(len, module.len() as u64);
}
#[tokio::test]
async fn all_sources_tampering_fails_closed_without_finalize() {
let store_id = hex_id(0x77);
let root = hex_id(0x88);
let module = b"content nobody serves honestly".to_vec();
let peer1 = crate::testkit::mock_peer_hex(1);
let transport =
Arc::new(MockModuleTransport::serving(&store_id, &root, module, 8).tampering(&peer1));
let downloader = ModuleDownloader::new(
locator_with(1, &store_id, &root),
transport,
Arc::new(AcceptAnyModuleAnchor),
Arc::new(InMemoryStateStore::new()),
ModuleDownloadConfig::default(),
);
let sink = InMemorySink::new();
let err = downloader
.download(&store_id, &root, &sink)
.await
.unwrap_err();
assert!(
matches!(err, DownloadError::NotFound { .. }),
"no honest source left is terminal: {err}"
);
assert!(
!sink.is_finalized().await,
"tampered content is never written through"
);
}
#[tokio::test]
async fn anchor_rejection_fails_closed_without_finalize() {
let store_id = hex_id(0x99);
let root = hex_id(0xAA);
let module = b"assembles cleanly but is not chain-anchored".to_vec();
let transport = Arc::new(MockModuleTransport::serving(&store_id, &root, module, 8));
let downloader = ModuleDownloader::new(
locator_with(2, &store_id, &root),
transport,
Arc::new(RejectAllModuleAnchor),
Arc::new(InMemoryStateStore::new()),
ModuleDownloadConfig::default(),
);
let sink = InMemorySink::new();
let err = downloader
.download(&store_id, &root, &sink)
.await
.unwrap_err();
assert!(
matches!(err, DownloadError::Verify(_)),
"anchor rejection is a verify failure"
);
assert!(
!sink.is_finalized().await,
"an unanchored module is never finalized"
);
}
#[tokio::test]
async fn wrong_whole_module_hash_fails_closed() {
let store_id = hex_id(0xBB);
let root = hex_id(0xCC);
let module = b"chunks are honest, module_hash lies".to_vec();
let transport = Arc::new(
MockModuleTransport::serving(&store_id, &root, module, 8).with_corrupt_module_hash(),
);
let downloader = ModuleDownloader::new(
locator_with(2, &store_id, &root),
transport,
Arc::new(AcceptAnyModuleAnchor),
Arc::new(InMemoryStateStore::new()),
ModuleDownloadConfig::default(),
);
let sink = InMemorySink::new();
let err = downloader
.download(&store_id, &root, &sink)
.await
.unwrap_err();
assert!(matches!(err, DownloadError::Verify(_)));
assert!(!sink.is_finalized().await);
}
#[tokio::test]
async fn no_holders_located_is_not_found() {
let store_id = hex_id(0x01);
let root = hex_id(0x02);
let transport = Arc::new(MockModuleTransport::serving(
&store_id,
&root,
vec![1, 2, 3],
8,
));
let downloader = ModuleDownloader::new(
Arc::new(MockProviderLocator::fixed(vec![])),
transport,
Arc::new(AcceptAnyModuleAnchor),
Arc::new(InMemoryStateStore::new()),
ModuleDownloadConfig::default(),
);
let sink = InMemorySink::new();
let err = downloader
.download(&store_id, &root, &sink)
.await
.unwrap_err();
assert!(matches!(err, DownloadError::NotFound { .. }));
}
#[tokio::test]
async fn an_over_long_range_is_clipped_not_rejected() {
let store_id = hex_id(0xD1);
let root = hex_id(0xD2);
let module = b"a chunk-granular holder overserves every window".to_vec();
let transport = Arc::new(
MockModuleTransport::serving(&store_id, &root, module.clone(), 8).overserving(),
);
let downloader = ModuleDownloader::new(
locator_with(1, &store_id, &root),
transport,
Arc::new(AcceptAnyModuleAnchor),
Arc::new(InMemoryStateStore::new()),
ModuleDownloadConfig::default(),
);
let sink = InMemorySink::new();
let len = downloader
.download(&store_id, &root, &sink)
.await
.expect("an over-long frame is clipped, so the pull completes");
assert_eq!(len, module.len() as u64);
assert_eq!(
sink.contents().await,
module,
"clipped to exactly the requested window — no bleed-through of the extra bytes"
);
assert!(sink.is_finalized().await);
}
#[tokio::test]
async fn exhausted_holders_name_the_failing_step_and_the_reasons() {
let store_id = hex_id(0xE1);
let root = hex_id(0xE2);
let peer1 = crate::testkit::mock_peer_hex(1);
let transport = Arc::new(
MockModuleTransport::serving(
&store_id,
&root,
b"nobody serves this honestly".to_vec(),
8,
)
.tampering(&peer1),
);
let downloader = ModuleDownloader::new(
locator_with(1, &store_id, &root),
transport,
Arc::new(AcceptAnyModuleAnchor),
Arc::new(InMemoryStateStore::new()),
ModuleDownloadConfig::default(),
);
let sink = InMemorySink::new();
let message = downloader
.download(&store_id, &root, &sink)
.await
.unwrap_err()
.to_string();
assert!(
message.contains("fetchModuleRange"),
"names the step that failed: {message}"
);
assert!(
message.contains("chunk 0"),
"names the chunk that could not be fetched: {message}"
);
assert!(
message.contains("chunk hash mismatch"),
"carries the per-holder reason instead of swallowing it: {message}"
);
assert!(
message.contains(&peer1),
"attributes the reason to a holder"
);
}
#[tokio::test]
async fn an_oversized_declared_module_is_refused_before_staging() {
let store_id = hex_id(0xF1);
let root = hex_id(0xF2);
let transport = Arc::new(
MockModuleTransport::serving(&store_id, &root, b"small blob, huge lie".to_vec(), 8)
.declaring_total_size(64 * 1024 * 1024 * 1024),
);
let downloader = ModuleDownloader::new(
locator_with(1, &store_id, &root),
transport,
Arc::new(AcceptAnyModuleAnchor),
Arc::new(InMemoryStateStore::new()),
ModuleDownloadConfig {
max_module_size: 1024,
..ModuleDownloadConfig::default()
},
);
let sink = InMemorySink::new();
let err = downloader
.download(&store_id, &root, &sink)
.await
.unwrap_err();
assert!(
matches!(err, DownloadError::Verify(_)),
"an over-cap descriptor is a verify failure: {err}"
);
assert!(
err.to_string().contains("exceeds the maximum"),
"names the bound it broke: {err}"
);
assert!(!sink.is_finalized().await);
}
#[tokio::test]
async fn a_corrupted_staged_chunk_is_re_fetched_on_resume() {
let store_id = hex_id(0xA1);
let root = hex_id(0xA2);
let module = (0u8..40).collect::<Vec<u8>>(); let state_store = Arc::new(InMemoryStateStore::new());
let sink = InMemorySink::new();
let interrupted = Arc::new(
MockModuleTransport::serving(&store_id, &root, module.clone(), 8)
.with_success_budget(2),
);
ModuleDownloader::new(
locator_with(1, &store_id, &root),
interrupted,
Arc::new(AcceptAnyModuleAnchor),
state_store.clone(),
ModuleDownloadConfig::default(),
)
.download(&store_id, &root, &sink)
.await
.expect_err("interrupted pull fails before finalize");
sink.write_at(0, &[0xFF; 8])
.await
.expect("staging is writable");
let healthy = Arc::new(MockModuleTransport::serving(
&store_id,
&root,
module.clone(),
8,
));
let len = ModuleDownloader::new(
locator_with(1, &store_id, &root),
healthy.clone(),
Arc::new(AcceptAnyModuleAnchor),
state_store,
ModuleDownloadConfig::default(),
)
.download(&store_id, &root, &sink)
.await
.expect("resume detects the corrupt staged chunk and re-fetches it");
assert_eq!(len, module.len() as u64);
assert_eq!(
sink.contents().await,
module,
"the corrupted staged chunk was replaced with honest bytes"
);
let refetched: BTreeSet<u64> = healthy
.fetches()
.await
.into_iter()
.map(|(_, o)| o)
.collect();
assert!(
refetched.contains(&0),
"the corrupt chunk was re-fetched: {refetched:?}"
);
assert!(
!refetched.contains(&8),
"the still-valid staged chunk was NOT re-fetched: {refetched:?}"
);
}
#[tokio::test]
async fn a_non_canonical_peer_id_is_sentinelled_not_echoed() {
let store_id = hex_id(0xB1);
let root = hex_id(0xB2);
let hostile = "not-hex <script>alert(1)</script>\n[FATAL] forged log line";
let content = module_content_id(&store_id, &root).unwrap();
let locator = Arc::new(MockProviderLocator::fixed(vec![
crate::testkit::mock_provider_with_peer_id(hostile, &content),
]));
let transport = Arc::new(MockModuleTransport::serving(
"unrelated-store",
&root,
vec![1, 2, 3],
8,
));
let downloader = ModuleDownloader::new(
locator,
transport,
Arc::new(AcceptAnyModuleAnchor),
Arc::new(InMemoryStateStore::new()),
ModuleDownloadConfig::default(),
);
let sink = InMemorySink::new();
let message = downloader
.download(&store_id, &root, &sink)
.await
.unwrap_err()
.to_string();
assert!(
!message.contains("<script>") && !message.contains("[FATAL]"),
"peer-supplied text is never echoed: {message}"
);
assert!(
message.contains("non-canonical-peer-id"),
"a sentinel stands in for it: {message}"
);
assert!(
!message.contains('\n'),
"the whole record is ONE line — a forged log line cannot ride in on the reason: {message}"
);
}
#[test]
fn untrusted_hex_is_sentinelled() {
let canonical = "ab".repeat(32);
assert_eq!(hex64_or_sentinel(&canonical, "peer-id"), canonical);
assert_eq!(
hex64_or_sentinel("AB".repeat(32).as_str(), "peer-id"),
"ab".repeat(32),
"canonical form is lowercase"
);
assert_eq!(
hex64_or_sentinel("short", "peer-id"),
"<non-canonical-peer-id>"
);
assert_eq!(
hex64_or_sentinel("zz".repeat(32).as_str(), "hash"),
"<non-canonical-hash>"
);
}
#[test]
fn malformed_ids_are_not_downloadable() {
assert!(module_content_id("too-short", &hex_id(1)).is_none());
assert!(module_content_id(&hex_id(1), "zz").is_none());
assert!(module_content_id(&hex_id(1), &hex_id(2)).is_some());
}
#[test]
fn download_key_is_module_scoped() {
let k = module_download_key(&hex_id(1), &hex_id(2));
assert!(k.starts_with("module:"));
}
#[tokio::test]
async fn a_lying_descriptor_source_is_demoted_and_an_honest_holder_completes_the_pull() {
let store_id = hex_id(0xC1);
let root = hex_id(0xC2);
let module = b"honest bytes, one lying descriptor source".to_vec();
let liar = crate::testkit::mock_peer_hex(1);
let transport = Arc::new(
MockModuleTransport::serving(&store_id, &root, module.clone(), 8)
.lying_descriptor_from(&liar),
);
let downloader = ModuleDownloader::new(
locator_with(2, &store_id, &root),
transport.clone(),
Arc::new(AcceptAnyModuleAnchor),
Arc::new(InMemoryStateStore::new()),
ModuleDownloadConfig::default(),
);
let sink = InMemorySink::new();
let len = downloader
.download(&store_id, &root, &sink)
.await
.expect("the honest holder's descriptor completes the pull");
assert_eq!(len, module.len() as u64);
assert_eq!(sink.contents().await, module);
assert!(sink.is_finalized().await);
let handshakes = transport.module_info_calls().await;
assert_eq!(
handshakes.len(),
2,
"the liar's descriptor was demoted and another holder re-handshaked: {handshakes:?}"
);
assert_eq!(handshakes[0], liar, "the liar answered first");
assert_ne!(
handshakes[1], liar,
"the demoted source is never re-asked: {handshakes:?}"
);
}
#[tokio::test]
async fn a_pull_does_not_re_adopt_a_demoted_descriptor_within_the_same_call() {
let store_id = hex_id(0xC5);
let root = hex_id(0xC6);
let module = (0u8..40).collect::<Vec<u8>>(); let liar = crate::testkit::mock_peer_hex(1);
let state_store = Arc::new(InMemoryStateStore::new());
let sink = InMemorySink::new();
let interrupted = Arc::new(
MockModuleTransport::serving(&store_id, &root, module.clone(), 8)
.lying_descriptor_from(&liar)
.with_success_budget(2),
);
ModuleDownloader::new(
locator_with(2, &store_id, &root),
interrupted,
Arc::new(AcceptAnyModuleAnchor),
state_store.clone(),
ModuleDownloadConfig::default(),
)
.download(&store_id, &root, &sink)
.await
.expect_err("the interrupted pull fails before finalize");
let healthy = Arc::new(
MockModuleTransport::serving(&store_id, &root, module.clone(), 8)
.lying_descriptor_from(&liar),
);
let len = ModuleDownloader::new(
locator_with(2, &store_id, &root),
healthy,
Arc::new(AcceptAnyModuleAnchor),
state_store,
ModuleDownloadConfig::default(),
)
.download(&store_id, &root, &sink)
.await
.expect("the resumed pull completes via the honest holder");
assert_eq!(len, module.len() as u64);
assert_eq!(sink.contents().await, module);
assert!(sink.is_finalized().await);
}
#[tokio::test]
async fn every_descriptor_source_lying_is_terminal_and_never_finalizes() {
let store_id = hex_id(0xC3);
let root = hex_id(0xC4);
let liar = crate::testkit::mock_peer_hex(1);
let transport = Arc::new(
MockModuleTransport::serving(&store_id, &root, b"only a liar holds this".to_vec(), 8)
.lying_descriptor_from(&liar),
);
let downloader = ModuleDownloader::new(
locator_with(1, &store_id, &root),
transport.clone(),
Arc::new(AcceptAnyModuleAnchor),
Arc::new(InMemoryStateStore::new()),
ModuleDownloadConfig::default(),
);
let sink = InMemorySink::new();
let err = downloader
.download(&store_id, &root, &sink)
.await
.unwrap_err();
assert!(
matches!(err, DownloadError::Verify(_)),
"fail-closed: {err}"
);
assert!(!sink.is_finalized().await);
assert!(
transport.module_info_calls().await.len() <= MAX_DESCRIPTOR_ATTEMPTS,
"descriptor attempts are bounded"
);
}
#[test]
fn a_wrapping_chunk_len_sum_is_rejected_not_panicked() {
let hostile = ModuleInfo {
total_size: 0,
module_hash: "ab".repeat(32),
chunk_hashes: vec!["cd".repeat(32), "ef".repeat(32)],
chunk_lens: vec![1, u64::MAX],
};
let err = ChunkPlan::from_info(&hostile, DEFAULT_MAX_MODULE_SIZE)
.expect_err("a wrapping descriptor is refused");
assert!(
err.is_proven_false(),
"a hostile descriptor is attributable to the holder that supplied it"
);
let err = err.error();
assert!(
matches!(err, DownloadError::Verify(_)),
"a hostile descriptor is a verify failure: {err}"
);
assert!(
err.to_string().contains("overflow"),
"names the arithmetic it broke: {err}"
);
}
#[test]
fn an_absurd_chunk_count_is_refused() {
let hostile = ModuleInfo {
total_size: 0,
module_hash: "ab".repeat(32),
chunk_hashes: Vec::new(),
chunk_lens: vec![0; MAX_MODULE_CHUNK_COUNT + 1],
};
let err = ChunkPlan::from_info(&hostile, DEFAULT_MAX_MODULE_SIZE)
.expect_err("an over-count descriptor is refused");
assert!(err.is_proven_false(), "attributable to its source");
let err = err.error();
assert!(
err.to_string().contains("chunk_lens"),
"names the bound it broke: {err}"
);
}
#[test]
fn chunk_plan_rejects_inconsistent_descriptor() {
let bad = ModuleInfo {
total_size: 99,
module_hash: "ab".repeat(32),
chunk_hashes: vec!["cd".repeat(32)],
chunk_lens: vec![5],
};
assert!(ChunkPlan::from_info(&bad, DEFAULT_MAX_MODULE_SIZE).is_err());
let no_lens = ModuleInfo {
total_size: 5,
module_hash: "ab".repeat(32),
chunk_hashes: vec!["cd".repeat(32)],
chunk_lens: vec![],
};
assert!(ChunkPlan::from_info(&no_lens, DEFAULT_MAX_MODULE_SIZE).is_err());
let mismatched = ModuleInfo {
total_size: 5,
module_hash: "ab".repeat(32),
chunk_hashes: vec!["cd".repeat(32), "ef".repeat(32)],
chunk_lens: vec![5],
};
assert!(ChunkPlan::from_info(&mismatched, DEFAULT_MAX_MODULE_SIZE).is_err());
}
fn temp_dir(tag: &str) -> std::path::PathBuf {
let d = std::env::temp_dir().join(format!(
"dig-download-module-{tag}-{}-{}",
std::process::id(),
std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap()
.as_nanos()
));
std::fs::create_dir_all(&d).unwrap();
d
}
#[tokio::test]
async fn the_promoted_artifact_is_byte_equal_to_the_verified_one_after_a_shorter_retry() {
let dir = temp_dir("shrinking-lie");
let final_path = dir.join("module.dig");
let store_id = hex_id(0xE1);
let root = hex_id(0xE2);
let honest = b"honest!!".to_vec(); let fabricated = vec![0xAA; 32]; let liar = crate::testkit::mock_peer_hex(1);
let transport = Arc::new(
MockModuleTransport::serving(&store_id, &root, honest.clone(), 8)
.serving_alternate_module_from(&liar, fabricated.clone()),
);
let downloader = ModuleDownloader::new(
locator_with(3, &store_id, &root),
transport,
Arc::new(crate::testkit::OnlyThisModuleAnchor::new(honest.clone())),
Arc::new(InMemoryStateStore::new()),
ModuleDownloadConfig::default(),
);
let sink = crate::sink::FileSink::new(&final_path);
let len = downloader
.download(&store_id, &root, &sink)
.await
.expect("the honest holder's descriptor completes the pull");
assert_eq!(
len,
honest.len() as u64,
"the VERIFIED length is the honest one"
);
let promoted = std::fs::read(&final_path).expect("the module was promoted");
assert_eq!(
promoted,
honest,
"the promoted artifact carries the attacker's tail: {} promoted bytes vs {} verified",
promoted.len(),
honest.len()
);
let _ = std::fs::remove_dir_all(&dir);
}
#[tokio::test]
async fn leftover_staging_of_another_shape_never_survives_into_the_promotion() {
let dir = temp_dir("stale-staging");
let final_path = dir.join("module.dig");
let store_id = hex_id(0xE3);
let root = hex_id(0xE4);
let honest = b"honest!!".to_vec();
let staging = crate::sink::staging_path_for(&final_path);
std::fs::write(&staging, vec![0xAA; 32]).unwrap();
let state_store = Arc::new(InMemoryStateStore::new());
let key = module_download_key(&store_id, &root);
let mut stale = DownloadState::new(&key);
stale.total_length = 32;
stale.chunk_lens = vec![8, 8, 8, 8];
stale.mark_done(0);
state_store.save(&stale).await.unwrap();
let downloader = ModuleDownloader::new(
locator_with(2, &store_id, &root),
Arc::new(MockModuleTransport::serving(
&store_id,
&root,
honest.clone(),
8,
)),
Arc::new(crate::testkit::OnlyThisModuleAnchor::new(honest.clone())),
state_store,
ModuleDownloadConfig::default(),
);
let sink = crate::sink::FileSink::new(&final_path);
let len = downloader.download(&store_id, &root, &sink).await.unwrap();
assert_eq!(len, honest.len() as u64);
assert_eq!(
std::fs::read(&final_path).unwrap(),
honest,
"the stale longer staging tail was promoted with the verified bytes"
);
let _ = std::fs::remove_dir_all(&dir);
}
struct UnshrinkableSink(InMemorySink);
#[async_trait]
impl Sink for UnshrinkableSink {
async fn write_at(&self, offset: u64, bytes: &[u8]) -> Result<(), DownloadError> {
self.0.write_at(offset, bytes).await
}
async fn truncate(&self, _len: u64) -> Result<(), DownloadError> {
Ok(()) }
fn supports_read_back(&self) -> bool {
true
}
async fn read_at(&self, offset: u64, len: u64) -> Result<Vec<u8>, DownloadError> {
self.0.read_at(offset, len).await
}
async fn finalize(&self) -> Result<(), DownloadError> {
self.0.finalize().await
}
}
#[tokio::test]
async fn a_staging_area_that_cannot_shrink_is_never_promoted() {
let store_id = hex_id(0xE7);
let root = hex_id(0xE8);
let honest = b"honest!!".to_vec();
let fabricated = vec![0xAA; 32];
let liar = crate::testkit::mock_peer_hex(1);
let downloader = ModuleDownloader::new(
locator_with(3, &store_id, &root),
Arc::new(
MockModuleTransport::serving(&store_id, &root, honest.clone(), 8)
.serving_alternate_module_from(&liar, fabricated),
),
Arc::new(crate::testkit::OnlyThisModuleAnchor::new(honest.clone())),
Arc::new(InMemoryStateStore::new()),
ModuleDownloadConfig::default(),
);
let sink = UnshrinkableSink(InMemorySink::new());
let err = downloader
.download(&store_id, &root, &sink)
.await
.expect_err("a staging area that still holds the demoted tail is not promoted");
assert!(
err.to_string().contains("past the verified length"),
"names the promotion invariant it refused: {err}"
);
assert!(!sink.0.is_finalized().await, "and it never finalized");
}
struct TwoDefaultSink(InMemorySink);
#[async_trait]
impl Sink for TwoDefaultSink {
async fn write_at(&self, offset: u64, bytes: &[u8]) -> Result<(), DownloadError> {
self.0.write_at(offset, bytes).await
}
async fn finalize(&self) -> Result<(), DownloadError> {
self.0.finalize().await
}
}
#[tokio::test]
async fn a_sink_on_both_truncate_and_read_at_defaults_fails_closed_at_the_first_reset() {
let store_id = hex_id(0xE9);
let root = hex_id(0xEA);
let honest = b"honest!!".to_vec();
let fabricated = vec![0xAA; 32];
let liar = crate::testkit::mock_peer_hex(1);
let downloader = ModuleDownloader::new(
locator_with(3, &store_id, &root),
Arc::new(
MockModuleTransport::serving(&store_id, &root, honest.clone(), 8)
.serving_alternate_module_from(&liar, fabricated),
),
Arc::new(crate::testkit::OnlyThisModuleAnchor::new(honest.clone())),
Arc::new(InMemoryStateStore::new()),
ModuleDownloadConfig::default(),
);
let sink = TwoDefaultSink(InMemorySink::new());
let err = downloader
.download(&store_id, &root, &sink)
.await
.expect_err("a sink on both defaults must fail closed, never promote blind");
assert!(
err.to_string().contains("truncation unsupported"),
"it is the fail-closed truncate DEFAULT that refuses, at the plan reset: {err}"
);
assert!(!sink.0.is_finalized().await, "and it never finalized");
}
#[tokio::test]
async fn a_fabricated_chunk_hash_descriptor_source_is_demoted_and_the_pull_completes() {
let store_id = hex_id(0xF1);
let root = hex_id(0xF2);
let module = b"honest bytes behind a zero-byte liar".to_vec();
let liar = crate::testkit::mock_peer_hex(1);
let transport = Arc::new(
MockModuleTransport::serving(&store_id, &root, module.clone(), 8)
.fabricating_chunk_hashes_from(&liar),
);
let downloader = ModuleDownloader::new(
locator_with(3, &store_id, &root),
transport.clone(),
Arc::new(crate::testkit::OnlyThisModuleAnchor::new(module.clone())),
Arc::new(InMemoryStateStore::new()),
ModuleDownloadConfig::default(),
);
let sink = InMemorySink::new();
let len = downloader
.download(&store_id, &root, &sink)
.await
.expect("an honest holder's descriptor completes the pull");
assert_eq!(len, module.len() as u64);
assert_eq!(sink.contents().await, module);
let handshakes = transport.module_info_calls().await;
assert_eq!(handshakes[0], liar, "the liar answered first");
assert!(
handshakes.len() >= 2 && handshakes[1] != liar,
"the fabricating source was demoted and another holder re-handshaked: {handshakes:?}"
);
}
#[tokio::test]
async fn a_liar_that_serves_one_byte_then_refuses_is_still_demoted() {
let store_id = hex_id(0x5A);
let root = hex_id(0x5B);
let module = b"an honest module blob of some length".to_vec();
let liar = crate::testkit::mock_peer_hex(1);
let transport = Arc::new(
MockModuleTransport::serving(&store_id, &root, module.clone(), 8)
.serving_one_byte_then_refusing_from(&liar),
);
let downloader = ModuleDownloader::new(
locator_with(3, &store_id, &root),
transport.clone(),
Arc::new(AcceptAnyModuleAnchor),
Arc::new(InMemoryStateStore::new()),
ModuleDownloadConfig::default(),
);
let sink = InMemorySink::new();
let len = downloader
.download(&store_id, &root, &sink)
.await
.expect("an honest holder's descriptor completes the pull");
assert_eq!(len, module.len() as u64);
assert_eq!(sink.contents().await, module, "only honest bytes promoted");
let handshakes = transport.module_info_calls().await;
assert!(
handshakes.len() > 1,
"the one-byte liar was demoted and another holder's descriptor tried: {handshakes:?}"
);
assert!(
handshakes.len() <= MAX_DESCRIPTOR_ATTEMPTS + 1,
"descriptor retries stay bounded by the attempt budget: {handshakes:?}"
);
}
#[test]
fn exhaustion_diagnosis_names_whether_any_chunk_had_verified() {
let base = DownloadError::NotFound {
content: "fetchModuleRange failed for chunk 3".into(),
};
let with_progress = describe_chunk_exhaustion(base, true).to_string();
let without_progress = describe_chunk_exhaustion(
DownloadError::NotFound {
content: "fetchModuleRange failed for chunk 3".into(),
},
false,
)
.to_string();
assert!(
with_progress.contains("chunk 3"),
"keeps the original reason: {with_progress}"
);
assert_ne!(
with_progress, without_progress,
"the two exhaustion cases read differently"
);
}
#[tokio::test]
async fn a_liar_demoted_in_one_call_is_not_re_asked_in_the_next() {
let store_id = hex_id(0x6A);
let root = hex_id(0x6B);
let module = b"an honest module across chunks".to_vec();
let liar = crate::testkit::mock_peer_hex(1);
let state_store = Arc::new(InMemoryStateStore::new());
let first_transport = Arc::new(
MockModuleTransport::serving(&store_id, &root, module.clone(), 8)
.lying_descriptor_from(&liar),
);
let first = ModuleDownloader::new(
locator_with(3, &store_id, &root),
first_transport.clone(),
Arc::new(AcceptAnyModuleAnchor),
state_store.clone(),
ModuleDownloadConfig::default(),
);
first
.download(&store_id, &root, &InMemorySink::new())
.await
.expect("an honest holder completes call 1");
assert!(
first_transport.module_info_calls().await.contains(&liar),
"call 1 did ask the liar (that is how it learned)"
);
let second_transport = Arc::new(
MockModuleTransport::serving(&store_id, &root, module.clone(), 8)
.lying_descriptor_from(&liar),
);
let second = ModuleDownloader::new(
locator_with(3, &store_id, &root),
second_transport.clone(),
Arc::new(AcceptAnyModuleAnchor),
state_store.clone(),
ModuleDownloadConfig::default(),
);
let len = second
.download(&store_id, &root, &InMemorySink::new())
.await
.expect("call 2 completes");
assert_eq!(len, module.len() as u64);
assert!(
!second_transport.module_info_calls().await.contains(&liar),
"the remembered liar is never asked for a descriptor again: {:?}",
second_transport.module_info_calls().await
);
assert_eq!(
second_transport.module_info_calls().await.len(),
1,
"and exactly one honest handshake was needed"
);
assert!(
second_transport
.fetches()
.await
.iter()
.any(|(peer, _)| peer == &liar),
"a demoted descriptor source is still used for CHUNK fetches"
);
assert_eq!(
state_store
.bad_descriptor_peers(&module_download_key(&store_id, &root))
.await
.unwrap(),
vec![liar],
"the verdict is what the store persisted"
);
}
#[tokio::test]
async fn reputation_never_denies_a_pull_when_every_holder_is_remembered() {
let store_id = hex_id(0x6C);
let root = hex_id(0x6D);
let module = b"honest bytes from a once-bad holder".to_vec();
let key = module_download_key(&store_id, &root);
let state_store = Arc::new(InMemoryStateStore::new());
state_store
.record_bad_descriptor(&key, &crate::testkit::mock_peer_hex(1))
.await
.unwrap();
let downloader = ModuleDownloader::new(
locator_with(1, &store_id, &root),
Arc::new(MockModuleTransport::serving(
&store_id,
&root,
module.clone(),
8,
)),
Arc::new(AcceptAnyModuleAnchor),
state_store,
ModuleDownloadConfig::default(),
);
let sink = InMemorySink::new();
let len = downloader
.download(&store_id, &root, &sink)
.await
.expect("a remembered holder is still asked when it is the only one");
assert_eq!(len, module.len() as u64);
assert_eq!(sink.contents().await, module);
}
#[tokio::test]
async fn reputation_never_denies_a_pull_when_only_the_honest_holders_are_remembered() {
let store_id = hex_id(0x7A);
let root = hex_id(0x7B);
let module = b"honest bytes the network can still serve".to_vec();
let key = module_download_key(&store_id, &root);
let liar = crate::testkit::mock_peer_hex(1);
let state_store = Arc::new(InMemoryStateStore::new());
for honest in [
crate::testkit::mock_peer_hex(2),
crate::testkit::mock_peer_hex(3),
] {
state_store
.record_bad_descriptor(&key, &honest)
.await
.unwrap();
}
let downloader = ModuleDownloader::new(
locator_with(3, &store_id, &root),
Arc::new(
MockModuleTransport::serving(&store_id, &root, module.clone(), 8)
.lying_descriptor_from(&liar),
),
Arc::new(AcceptAnyModuleAnchor),
state_store,
ModuleDownloadConfig::default(),
);
let sink = InMemorySink::new();
let len = downloader
.download(&store_id, &root, &sink)
.await
.expect("remembered HONEST holders are re-asked rather than denying the pull");
assert_eq!(len, module.len() as u64);
assert_eq!(sink.contents().await, module);
}
#[tokio::test]
async fn chunk_exhaustion_demotes_for_this_call_but_records_no_durable_verdict() {
let store_id = hex_id(0x7C);
let root = hex_id(0x7D);
let key = module_download_key(&store_id, &root);
let state_store = Arc::new(InMemoryStateStore::new());
let downloader = ModuleDownloader::new(
locator_with(1, &store_id, &root),
Arc::new(
MockModuleTransport::serving(&store_id, &root, b"unavailable bytes".to_vec(), 8)
.with_success_budget(0),
),
Arc::new(AcceptAnyModuleAnchor),
state_store.clone(),
ModuleDownloadConfig::default(),
);
let sink = InMemorySink::new();
downloader
.download(&store_id, &root, &sink)
.await
.expect_err("unavailable chunks fail the pull");
assert!(
state_store
.bad_descriptor_peers(&key)
.await
.unwrap()
.is_empty(),
"no durable verdict from mere unavailability — else sybils can brand an honest holder"
);
assert!(!sink.is_finalized().await);
}
#[tokio::test]
async fn the_documented_whole_commit_sink_recipe_cannot_promote_unproven_bytes() {
struct RecipeSink(InMemorySink);
#[async_trait]
impl Sink for RecipeSink {
async fn write_at(&self, offset: u64, bytes: &[u8]) -> Result<(), DownloadError> {
self.0.write_at(offset, bytes).await
}
async fn truncate(&self, _len: u64) -> Result<(), DownloadError> {
Ok(()) }
async fn finalize(&self) -> Result<(), DownloadError> {
self.0.finalize().await
}
}
let store_id = hex_id(0x7E);
let root = hex_id(0x7F);
let honest = b"honest!!".to_vec();
let liar = crate::testkit::mock_peer_hex(1);
let downloader = ModuleDownloader::new(
locator_with(3, &store_id, &root),
Arc::new(
MockModuleTransport::serving(&store_id, &root, honest.clone(), 8)
.serving_alternate_module_from(&liar, vec![0xAA; 32]),
),
Arc::new(crate::testkit::OnlyThisModuleAnchor::new(honest)),
Arc::new(InMemoryStateStore::new()),
ModuleDownloadConfig::default(),
);
let sink = RecipeSink(InMemorySink::new());
let err = downloader
.download(&store_id, &root, &sink)
.await
.expect_err("a sink that cannot prove its staged length is never promoted");
assert!(
err.to_string().contains("cannot read back"),
"names WHY it refused — an unprovable promotion, not a length verdict: {err}"
);
assert!(!sink.0.is_finalized().await, "and it never finalized");
}
#[test]
fn an_18_exbibyte_declaration_costs_no_allocation() {
let hostile = ModuleInfo {
total_size: u64::MAX,
module_hash: hex_id(0x01),
chunk_hashes: vec![hex_id(0x02)],
chunk_lens: vec![u64::MAX],
};
let Ok(plan) = ChunkPlan::from_info(&hostile, u64::MAX) else {
panic!("the plan is derived, not allocated — an 18 EiB claim is now cheap to hold")
};
assert_eq!(plan.total_size, u64::MAX);
assert_eq!(plan.chunk_count(), 1);
}
#[tokio::test]
async fn an_honest_holder_completes_the_pull_after_an_impossible_descriptor() {
let store_id = hex_id(0x8A);
let root = hex_id(0x8B);
let key = module_download_key(&store_id, &root);
let module = b"a real module served honestly".to_vec();
let liar = crate::testkit::mock_peer_hex(1);
let state_store = Arc::new(InMemoryStateStore::new());
let transport = Arc::new(
MockModuleTransport::serving(&store_id, &root, module.clone(), 8)
.inflating_total_size_from(&liar, u64::MAX),
);
let downloader = ModuleDownloader::new(
locator_with(3, &store_id, &root),
transport.clone(),
Arc::new(AcceptAnyModuleAnchor),
state_store.clone(),
ModuleDownloadConfig {
max_module_size: u64::MAX,
..ModuleDownloadConfig::default()
},
);
let sink = InMemorySink::new();
let len = downloader
.download(&store_id, &root, &sink)
.await
.expect("an honest holder's descriptor completes the pull");
assert_eq!(len, module.len() as u64);
assert_eq!(sink.contents().await, module);
assert!(
transport.module_info_calls().await.len() > 1,
"the unallocatable descriptor's source was demoted and another holder asked: {:?}",
transport.module_info_calls().await
);
assert!(
state_store
.bad_descriptor_peers(&key)
.await
.unwrap()
.is_empty(),
"and NOBODY is branded: an unsatisfiable descriptor is not PROOF the holder lied — the \
bytes it declares may simply be unavailable"
);
let _ = &liar;
}
#[tokio::test]
async fn every_holder_declaring_an_impossible_module_fails_closed() {
let store_id = hex_id(0x8E);
let root = hex_id(0x8F);
let key = module_download_key(&store_id, &root);
let state_store = Arc::new(InMemoryStateStore::new());
let downloader = ModuleDownloader::new(
locator_with(2, &store_id, &root),
Arc::new(
MockModuleTransport::serving(&store_id, &root, b"a real module".to_vec(), 8)
.declaring_total_size(u64::MAX),
),
Arc::new(AcceptAnyModuleAnchor),
state_store.clone(),
ModuleDownloadConfig {
max_module_size: u64::MAX,
..ModuleDownloadConfig::default()
},
);
let sink = InMemorySink::new();
downloader
.download(&store_id, &root, &sink)
.await
.expect_err("no holder offers a module that could exist");
assert!(
state_store
.bad_descriptor_peers(&key)
.await
.unwrap()
.is_empty(),
"no holder is branded for a descriptor merely unsatisfiable"
);
assert!(!sink.is_finalized().await, "and nothing is promoted");
}
#[tokio::test]
async fn an_unreachable_chain_anchor_is_terminal_and_brands_nobody() {
let store_id = hex_id(0x8C);
let root = hex_id(0x8D);
let key = module_download_key(&store_id, &root);
let state_store = Arc::new(InMemoryStateStore::new());
let downloader = ModuleDownloader::new(
locator_with(3, &store_id, &root),
Arc::new(MockModuleTransport::serving(
&store_id,
&root,
b"a correct, genuinely anchored module".to_vec(),
8,
)),
Arc::new(crate::testkit::UnreachableChainAnchor),
state_store.clone(),
ModuleDownloadConfig::default(),
);
let sink = InMemorySink::new();
let err = downloader
.download(&store_id, &root, &sink)
.await
.expect_err("an unverifiable anchor is fail-closed");
assert!(
err.to_string().contains("cannot verify the chain anchor"),
"it reports an unfinished CHECK, not a verdict on the module: {err}"
);
assert!(
state_store
.bad_descriptor_peers(&key)
.await
.unwrap()
.is_empty(),
"and no honest holder is branded by this node's own outage"
);
assert!(
!sink.is_finalized().await,
"fail-closed: nothing is promoted while the anchor is unproven"
);
}
#[tokio::test]
async fn the_whole_module_hash_is_taken_in_chunk_order_not_arrival_order() {
let store_id = hex_id(0x90);
let root = hex_id(0x91);
let key = module_download_key(&store_id, &root);
let module = (0u8..40).collect::<Vec<u8>>();
let sink = InMemorySink::new();
sink.write_at(16, &module[16..24]).await.unwrap();
let state_store = Arc::new(InMemoryStateStore::new());
let mut state = DownloadState::new(&key);
state.total_length = module.len() as u64;
state.chunk_lens = vec![8; 5];
state.mark_done(2);
state_store.save(&state).await.unwrap();
let downloader = ModuleDownloader::new(
locator_with(1, &store_id, &root),
Arc::new(MockModuleTransport::serving(
&store_id,
&root,
module.clone(),
8,
)),
Arc::new(crate::testkit::OnlyThisModuleAnchor::new(module.clone())),
state_store,
ModuleDownloadConfig::default(),
);
let len = downloader
.download(&store_id, &root, &sink)
.await
.expect("a mid-module checkpoint resumes and still passes both gates");
assert_eq!(len, module.len() as u64);
assert_eq!(sink.contents().await, module);
assert!(sink.is_finalized().await);
}
#[tokio::test]
async fn the_anchor_gate_cannot_read_past_the_verified_module() {
struct ReadsPastTheEnd;
#[async_trait]
impl ModuleAnchorVerifier for ReadsPastTheEnd {
async fn verify_module_anchor(
&self,
module: &dyn ModuleReader,
_store_id: &str,
_root: &str,
) -> ModuleAnchor {
match module.read_at(module.len().saturating_sub(1), 2).await {
Ok(_) => ModuleAnchor::Anchored,
Err(e) => ModuleAnchor::Unavailable(e.to_string()),
}
}
}
let store_id = hex_id(0x92);
let root = hex_id(0x93);
let module = (0u8..40).collect::<Vec<u8>>();
let downloader = ModuleDownloader::new(
locator_with(1, &store_id, &root),
Arc::new(MockModuleTransport::serving(
&store_id,
&root,
module.clone(),
8,
)),
Arc::new(ReadsPastTheEnd),
Arc::new(InMemoryStateStore::new()),
ModuleDownloadConfig::default(),
);
let sink = InMemorySink::new();
let err = downloader
.download(&store_id, &root, &sink)
.await
.expect_err("a read past the verified end is refused, so the gate reaches no answer");
assert!(
err.to_string().contains("falls outside"),
"the refusal names the out-of-range window: {err}"
);
assert!(!sink.is_finalized().await, "and nothing is promoted");
}
#[tokio::test]
async fn staging_corrupted_between_the_two_gates_fails_closed_and_brands_nobody() {
struct CorruptingReadBack(InMemorySink);
#[async_trait]
impl Sink for CorruptingReadBack {
async fn write_at(&self, offset: u64, bytes: &[u8]) -> Result<(), DownloadError> {
self.0.write_at(offset, bytes).await
}
async fn truncate(&self, len: u64) -> Result<(), DownloadError> {
self.0.truncate(len).await
}
fn supports_read_back(&self) -> bool {
true
}
async fn read_at(&self, offset: u64, len: u64) -> Result<Vec<u8>, DownloadError> {
let mut bytes = self.0.read_at(offset, len).await?;
if let Some(first) = bytes.first_mut() {
*first ^= 0xFF;
}
Ok(bytes)
}
async fn finalize(&self) -> Result<(), DownloadError> {
self.0.finalize().await
}
}
let store_id = hex_id(0x94);
let root = hex_id(0x95);
let key = module_download_key(&store_id, &root);
let module = (0u8..40).collect::<Vec<u8>>();
let state_store = Arc::new(InMemoryStateStore::new());
let downloader = ModuleDownloader::new(
locator_with(2, &store_id, &root),
Arc::new(MockModuleTransport::serving(
&store_id,
&root,
module.clone(),
8,
)),
Arc::new(crate::testkit::OnlyThisModuleAnchor::new(module.clone())),
state_store.clone(),
ModuleDownloadConfig::default(),
);
let sink = CorruptingReadBack(InMemorySink::new());
let err = downloader
.download(&store_id, &root, &sink)
.await
.expect_err("the gate must not run on bytes nothing has attributed");
assert!(
err.to_string()
.contains("no longer matches its verified hash"),
"the failure names the staging area, not the chain or a holder: {err}"
);
assert!(
state_store
.bad_descriptor_peers(&key)
.await
.unwrap()
.is_empty(),
"local corruption is never evidence against a holder that served correct bytes"
);
assert!(!sink.0.is_finalized().await, "and nothing is promoted");
}
}