use std::collections::HashMap;
use std::sync::atomic::{AtomicUsize, Ordering};
use std::sync::Arc;
use async_trait::async_trait;
use dig_dht::{CandidateAddr, ContentId, PeerId, ProviderRecord};
use dig_nat::{AvailabilityAnswer, AvailabilityItem, AvailabilityResponse};
use dig_rpc_protocol::types::ModuleInfo;
use tokio::sync::Mutex;
use crate::error::DownloadError;
use crate::locate::ProviderLocator;
use crate::select::{RangeOutcome, SelectPlan, SelectRequest, SourceSelector};
use crate::source::{FetchedRange, RangeMeta, RangeTransport};
#[derive(Debug, Clone)]
pub struct MockContent {
pub bytes: Vec<u8>,
pub chunk_lens: Vec<u64>,
pub root: String,
pub inclusion_proof: Option<String>,
offsets: Vec<u64>,
declared_layout: Option<(u64, Vec<u64>)>,
}
impl MockContent {
pub fn new(bytes: Vec<u8>, chunk_lens: Vec<u64>) -> Self {
assert_eq!(
bytes.len() as u64,
chunk_lens.iter().sum::<u64>(),
"chunk_lens must sum to bytes.len()"
);
let mut offsets = Vec::with_capacity(chunk_lens.len() + 1);
let mut acc = 0u64;
offsets.push(0);
for &l in &chunk_lens {
acc += l;
offsets.push(acc);
}
MockContent {
bytes,
chunk_lens,
root: "ab".repeat(32),
inclusion_proof: Some("mock-proof".into()),
offsets,
declared_layout: None,
}
}
pub fn declaring(mut self, total_length: u64, chunk_lens: Vec<u64>) -> Self {
self.declared_layout = Some((total_length, chunk_lens));
self
}
pub fn even(n: usize, chunks: usize) -> Self {
let chunks = chunks.max(1);
let base = n / chunks;
let mut lens = vec![base as u64; chunks];
let assigned: u64 = lens.iter().sum();
if let Some(last) = lens.last_mut() {
*last += n as u64 - assigned;
}
let bytes: Vec<u8> = (0..n).map(|i| (i % 251) as u8).collect();
MockContent::new(bytes, lens)
}
fn chunk_index_at(&self, offset: u64) -> u64 {
self.offsets.iter().position(|&o| o == offset).unwrap_or(0) as u64
}
fn meta(&self, offset: u64) -> RangeMeta {
let (total_length, chunk_lens) = match &self.declared_layout {
Some((declared_total, declared_lens)) => (*declared_total, declared_lens.clone()),
None => (self.bytes.len() as u64, self.chunk_lens.clone()),
};
let chunk_count = chunk_lens.len() as u64;
RangeMeta::default()
.declaring_layout(total_length, chunk_lens, chunk_count)
.declaring_anchor(self.root.clone(), self.inclusion_proof.clone())
.declaring_chunk_index(self.chunk_index_at(offset))
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum Behavior {
Honest,
Corrupt,
Truncate,
ShortAligned,
Unavailable,
DropAfter(usize),
AlwaysFail,
WrongRoot,
ShortLayout {
total_length: u64,
chunk_lens: Vec<u64>,
availability: AvailabilityClaim,
},
UnderDeliveredPrologue {
declared_chunk_count: u64,
},
NoRoot,
AvailableThenRefuses,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum AvailabilityClaim {
Honest,
OwnShortShape,
Silent,
}
pub struct MockRangeTransport {
content: MockContent,
behaviors: Mutex<HashMap<String, Behavior>>,
provider_attempts: Mutex<HashMap<String, usize>>,
offset_attempts: Mutex<HashMap<u64, usize>>,
probe_count: Mutex<usize>,
delay: Mutex<Option<std::time::Duration>>,
delay_probes: Mutex<bool>,
}
impl MockRangeTransport {
pub fn new(content: MockContent) -> Self {
MockRangeTransport {
content,
behaviors: Mutex::new(HashMap::new()),
provider_attempts: Mutex::new(HashMap::new()),
offset_attempts: Mutex::new(HashMap::new()),
probe_count: Mutex::new(0),
delay: Mutex::new(None),
delay_probes: Mutex::new(true),
}
}
pub async fn set_delay(&self, delay: std::time::Duration) {
*self.delay.lock().await = Some(delay);
*self.delay_probes.lock().await = true;
}
pub async fn set_range_delay(&self, delay: std::time::Duration) {
*self.delay.lock().await = Some(delay);
*self.delay_probes.lock().await = false;
}
pub async fn set_behavior(&self, peer_id: &str, behavior: Behavior) {
self.behaviors
.lock()
.await
.insert(peer_id.to_string(), behavior);
}
pub async fn attempts_for(&self, peer_id: &str) -> usize {
self.provider_attempts
.lock()
.await
.get(peer_id)
.copied()
.unwrap_or(0)
}
pub async fn metadata_probes(&self) -> usize {
*self.probe_count.lock().await
}
pub async fn attempts_at(&self, offset: u64) -> usize {
self.offset_attempts
.lock()
.await
.get(&offset)
.copied()
.unwrap_or(0)
}
async fn behavior(&self, peer_id: &str) -> Behavior {
self.behaviors
.lock()
.await
.get(peer_id)
.cloned()
.unwrap_or(Behavior::Honest)
}
}
#[async_trait]
impl RangeTransport for MockRangeTransport {
async fn query_availability(
&self,
provider: &ProviderRecord,
items: Vec<AvailabilityItem>,
) -> Result<AvailabilityResponse, DownloadError> {
let behavior = self.behavior(&provider.provider_peer_id).await;
let held = !matches!(behavior, Behavior::Unavailable | Behavior::AlwaysFail);
let claimed: Option<(u64, u64)> = match &behavior {
Behavior::ShortLayout {
availability: AvailabilityClaim::Silent,
..
} => None,
Behavior::ShortLayout {
total_length,
chunk_lens,
availability: AvailabilityClaim::OwnShortShape,
} => Some((*total_length, chunk_lens.len() as u64)),
_ => Some((
self.content.bytes.len() as u64,
self.content.chunk_lens.len() as u64,
)),
};
let answers = items
.iter()
.map(|_| {
let answer = if held {
AvailabilityAnswer::available()
} else {
AvailabilityAnswer::unavailable()
};
let answer = answer
.with_roots(vec![self.content.root.clone()])
.with_complete(true);
match claimed {
Some((total_length, chunk_count)) => answer
.with_total_length(total_length)
.with_chunk_count(chunk_count),
None => answer,
}
})
.collect();
Ok(AvailabilityResponse::new(answers))
}
async fn fetch_range(
&self,
provider: &ProviderRecord,
req: &dig_nat::RangeRequest,
) -> Result<FetchedRange, DownloadError> {
let peer = provider.provider_peer_id.clone();
let attempts = {
let mut a = self.provider_attempts.lock().await;
let n = a.entry(peer.clone()).or_insert(0);
*n += 1;
*n
};
*self
.offset_attempts
.lock()
.await
.entry(req.offset)
.or_insert(0) += 1;
if req.length == 1 {
*self.probe_count.lock().await += 1;
}
let is_probe = req.length == 1;
if let Some(d) = *self.delay.lock().await {
if !is_probe || *self.delay_probes.lock().await {
tokio::time::sleep(d).await;
}
}
let behavior = self.behavior(&peer).await;
let fail = || DownloadError::transport(&peer, "mock: source failed");
match behavior {
Behavior::Unavailable | Behavior::AlwaysFail | Behavior::AvailableThenRefuses => {
return Err(fail())
}
Behavior::DropAfter(n) if attempts > n => return Err(fail()),
_ => {}
}
let start = req.offset as usize;
let end = (req.offset + req.length).min(self.content.bytes.len() as u64) as usize;
let mut bytes = self.content.bytes[start..end].to_vec();
match behavior {
Behavior::Truncate => {
bytes.pop(); }
Behavior::Corrupt => {
for b in bytes.iter_mut() {
*b ^= 0xFF; }
}
Behavior::ShortAligned => {
let first_chunk_idx = self.content.chunk_index_at(req.offset) as usize;
if let Some(&first_len) = self.content.chunk_lens.get(first_chunk_idx) {
let keep = (first_len as usize).min(bytes.len());
if keep < bytes.len() {
bytes.truncate(keep);
}
}
}
_ => {}
}
let mut meta = self.content.meta(req.offset);
match &behavior {
Behavior::WrongRoot => {
meta.root = Some("cd".repeat(32));
}
Behavior::ShortLayout {
total_length,
chunk_lens,
..
} => {
meta = meta.declaring_layout(
*total_length,
chunk_lens.clone(),
chunk_lens.len() as u64,
);
}
Behavior::UnderDeliveredPrologue {
declared_chunk_count,
} => {
meta.chunk_count = Some(*declared_chunk_count);
}
Behavior::NoRoot => {
meta.root = None;
}
_ => {}
}
Ok(FetchedRange {
request_offset: req.offset,
bytes,
meta,
})
}
}
pub struct MockProviderLocator {
batches: Vec<Vec<ProviderRecord>>,
calls: Mutex<usize>,
}
impl MockProviderLocator {
pub fn fixed(providers: Vec<ProviderRecord>) -> Self {
MockProviderLocator {
batches: vec![providers],
calls: Mutex::new(0),
}
}
pub fn scripted(batches: Vec<Vec<ProviderRecord>>) -> Self {
MockProviderLocator {
batches: if batches.is_empty() {
vec![vec![]]
} else {
batches
},
calls: Mutex::new(0),
}
}
pub async fn call_count(&self) -> usize {
*self.calls.lock().await
}
}
#[async_trait]
impl ProviderLocator for MockProviderLocator {
async fn find_providers(
&self,
_content: &ContentId,
) -> Result<Vec<ProviderRecord>, DownloadError> {
let mut calls = self.calls.lock().await;
let idx = (*calls).min(self.batches.len() - 1);
*calls += 1;
Ok(self.batches[idx].clone())
}
}
#[derive(Default)]
pub struct MockSelector {
select_calls: AtomicUsize,
recorded: std::sync::Mutex<Vec<RangeOutcome>>,
forced_order: std::sync::Mutex<Option<Vec<String>>>,
}
impl MockSelector {
pub fn new() -> Arc<Self> {
Arc::new(MockSelector::default())
}
pub fn with_order(order: Vec<String>) -> Arc<Self> {
let sel = MockSelector::default();
*sel.forced_order.lock().unwrap() = Some(order);
Arc::new(sel)
}
pub fn select_call_count(&self) -> usize {
self.select_calls.load(Ordering::Relaxed)
}
pub fn outcomes(&self) -> Vec<RangeOutcome> {
self.recorded.lock().unwrap().clone()
}
}
impl SourceSelector for MockSelector {
fn select(&self, req: &SelectRequest) -> SelectPlan {
self.select_calls.fetch_add(1, Ordering::Relaxed);
let offered: Vec<String> = req.candidates.iter().map(|c| c.peer_id.clone()).collect();
match self.forced_order.lock().unwrap().as_ref() {
Some(order) => {
let mut plan: Vec<String> = order
.iter()
.filter(|p| offered.contains(p))
.cloned()
.collect();
for p in &offered {
if !plan.contains(p) {
plan.push(p.clone());
}
}
SelectPlan::ordered(plan)
}
None => SelectPlan::ordered(offered),
}
}
fn record(&self, outcome: &RangeOutcome) {
self.recorded.lock().unwrap().push(outcome.clone());
}
}
pub fn mock_provider(n: u8, content: &ContentId) -> ProviderRecord {
ProviderRecord::new(
&content.to_key(),
&PeerId::from_bytes([n; 32]),
vec![CandidateAddr::direct(format!("10.0.0.{n}"), 9444)],
u64::MAX,
)
}
pub fn mock_peer_hex(n: u8) -> String {
PeerId::from_bytes([n; 32]).to_hex()
}
pub fn mock_provider_with_peer_id(peer_id: &str, content: &ContentId) -> ProviderRecord {
let mut record = mock_provider(1, content);
record.provider_peer_id = peer_id.to_string();
record
}
pub fn mock_content_id() -> ContentId {
ContentId::resource([1; 32], [0xAB; 32], [3; 32])
}
pub struct MockModuleTransport {
store_id: String,
root: String,
blob: Vec<u8>,
chunk_size: usize,
tamper_peer: Option<String>,
corrupt_module_hash: bool,
overserve: bool,
declared_total_size: Option<u64>,
inflating_peer: Option<(String, u64)>,
wrong_descriptor_for: Option<String>,
alternate_module_for: Option<(String, Vec<u8>)>,
fabricated_chunk_hashes_for: Option<String>,
one_byte_liar: Option<String>,
budget: Option<Arc<AtomicUsize>>,
severed_beyond: Option<u64>,
info_failures_left: Option<Arc<AtomicUsize>>,
fetches: Mutex<Vec<(String, u64)>>,
info_calls: Mutex<Vec<String>>,
}
impl MockModuleTransport {
pub fn serving(store_id: &str, root: &str, blob: Vec<u8>, chunk_size: usize) -> Self {
MockModuleTransport {
store_id: store_id.to_string(),
root: root.to_string(),
blob,
chunk_size,
tamper_peer: None,
corrupt_module_hash: false,
overserve: false,
declared_total_size: None,
inflating_peer: None,
wrong_descriptor_for: None,
alternate_module_for: None,
fabricated_chunk_hashes_for: None,
one_byte_liar: None,
budget: None,
severed_beyond: None,
info_failures_left: None,
fetches: Mutex::new(Vec::new()),
info_calls: Mutex::new(Vec::new()),
}
}
pub fn overserving(mut self) -> Self {
self.overserve = true;
self
}
pub fn declaring_total_size(mut self, total_size: u64) -> Self {
self.declared_total_size = Some(total_size);
self
}
pub fn inflating_total_size_from(mut self, peer_id: &str, total_size: u64) -> Self {
self.inflating_peer = Some((peer_id.to_string(), total_size));
self
}
pub fn tampering(mut self, peer_id: &str) -> Self {
self.tamper_peer = Some(peer_id.to_string());
self
}
pub fn with_corrupt_module_hash(mut self) -> Self {
self.corrupt_module_hash = true;
self
}
pub fn lying_descriptor_from(mut self, peer_id: &str) -> Self {
self.wrong_descriptor_for = Some(peer_id.to_string());
self
}
pub fn serving_alternate_module_from(mut self, peer_id: &str, module: Vec<u8>) -> Self {
self.alternate_module_for = Some((peer_id.to_string(), module));
self
}
pub fn fabricating_chunk_hashes_from(mut self, peer_id: &str) -> Self {
self.fabricated_chunk_hashes_for = Some(peer_id.to_string());
self
}
pub fn serving_one_byte_then_refusing_from(mut self, peer_id: &str) -> Self {
self.one_byte_liar = Some(peer_id.to_string());
self
}
pub fn with_success_budget(mut self, n: usize) -> Self {
self.budget = Some(Arc::new(AtomicUsize::new(n)));
self
}
pub fn severing_beyond(mut self, offset: u64) -> Self {
self.severed_beyond = Some(offset);
self
}
pub fn failing_the_first_info_asks(mut self, n: usize) -> Self {
self.info_failures_left = Some(Arc::new(AtomicUsize::new(n)));
self
}
pub async fn fetches(&self) -> Vec<(String, u64)> {
self.fetches.lock().await.clone()
}
pub async fn module_info_calls(&self) -> Vec<String> {
self.info_calls.lock().await.clone()
}
fn descriptor_of(&self, module: &[u8]) -> ModuleInfo {
let mut chunk_hashes = Vec::new();
let mut chunk_lens = Vec::new();
for chunk in module.chunks(self.chunk_size.max(1)) {
chunk_hashes.push(hex_sha256(chunk));
chunk_lens.push(chunk.len() as u64);
}
ModuleInfo {
total_size: module.len() as u64,
module_hash: hex_sha256(module),
chunk_hashes,
chunk_lens,
}
}
fn descriptor(&self) -> ModuleInfo {
let ModuleInfo {
chunk_hashes,
mut chunk_lens,
..
} = self.descriptor_of(&self.blob);
let module_hash = if self.corrupt_module_hash {
"0".repeat(64)
} else {
hex_sha256(&self.blob)
};
let mut total_size = self.blob.len() as u64;
if let Some(inflated) = self.declared_total_size {
let last = chunk_lens.last_mut().expect("blob yields >= 1 chunk");
*last += inflated.saturating_sub(total_size);
total_size = inflated;
}
ModuleInfo {
total_size,
module_hash,
chunk_hashes,
chunk_lens,
}
}
fn descriptor_inflated_to(&self, inflated: u64) -> ModuleInfo {
let mut info = self.descriptor();
let real_total = info.total_size;
if let Some(last) = info.chunk_lens.last_mut() {
*last += inflated.saturating_sub(real_total);
}
info.total_size = inflated;
info
}
fn lying_descriptor(&self) -> ModuleInfo {
ModuleInfo {
module_hash: "0".repeat(64),
..self.descriptor()
}
}
fn fabricated_chunk_hash_descriptor(&self) -> ModuleInfo {
let mut info = self.descriptor();
info.chunk_hashes = info
.chunk_hashes
.iter()
.enumerate()
.map(|(i, _)| hex_sha256(format!("fabricated-chunk-{i}").as_bytes()))
.collect();
info
}
fn one_byte_then_unsatisfiable_descriptor(&self) -> ModuleInfo {
let rest = self.blob.len().saturating_sub(1) as u64;
ModuleInfo {
total_size: self.blob.len() as u64,
module_hash: hex_sha256(&self.blob),
chunk_hashes: vec![
hex_sha256(&self.blob[..1]),
hex_sha256(b"fabricated-second-chunk"),
],
chunk_lens: vec![1, rest],
}
}
fn served_blob(&self, provider_peer_id: &str) -> &[u8] {
match &self.alternate_module_for {
Some((peer, module)) if peer == provider_peer_id => module,
_ => &self.blob,
}
}
}
#[async_trait]
impl crate::module::ModuleTransport for MockModuleTransport {
async fn get_module_info(
&self,
provider_peer_id: &str,
store_id: &str,
root: &str,
) -> Result<ModuleInfo, DownloadError> {
self.info_calls
.lock()
.await
.push(provider_peer_id.to_string());
if let Some(left) = &self.info_failures_left {
if left.fetch_sub(1, Ordering::SeqCst) == 0 {
left.fetch_add(1, Ordering::SeqCst); } else {
return Err(DownloadError::transport(
provider_peer_id,
"descriptor ask timed out",
));
}
}
if store_id != self.store_id || root != self.root {
return Err(DownloadError::transport(
provider_peer_id,
"unknown (store_id, root)",
));
}
if let Some(wrong) = self.wrong_descriptor_for.as_deref() {
if wrong == provider_peer_id {
return Ok(self.lying_descriptor());
}
}
if let Some(fabricator) = self.fabricated_chunk_hashes_for.as_deref() {
if fabricator == provider_peer_id {
return Ok(self.fabricated_chunk_hash_descriptor());
}
}
if let Some((peer, module)) = &self.alternate_module_for {
if peer == provider_peer_id {
return Ok(self.descriptor_of(module));
}
}
if self.one_byte_liar.as_deref() == Some(provider_peer_id) {
return Ok(self.one_byte_then_unsatisfiable_descriptor());
}
if let Some((peer, inflated)) = &self.inflating_peer {
if peer == provider_peer_id {
return Ok(self.descriptor_inflated_to(*inflated));
}
}
Ok(self.descriptor())
}
async fn fetch_module_range(
&self,
provider_peer_id: &str,
_store_id: &str,
_root: &str,
offset: u64,
length: u64,
) -> Result<Vec<u8>, DownloadError> {
if let Some(budget) = &self.budget {
if budget.fetch_sub(1, Ordering::SeqCst) == 0 {
budget.fetch_add(1, Ordering::SeqCst); return Err(DownloadError::transport(provider_peer_id, "budget spent"));
}
}
if let Some(severed) = self.severed_beyond {
if offset >= severed && self.info_calls.lock().await.len() < 2 {
return Err(DownloadError::transport(
provider_peer_id,
"connection reset mid-transfer",
));
}
}
self.fetches
.lock()
.await
.push((provider_peer_id.to_string(), offset));
if self.fabricated_chunk_hashes_for.as_deref() == Some(provider_peer_id) {
return Err(DownloadError::transport(
provider_peer_id,
"serving nothing",
));
}
if self.one_byte_liar.as_deref() == Some(provider_peer_id) {
return if offset == 0 && length == 1 {
Ok(self.blob[..1].to_vec())
} else {
Err(DownloadError::transport(
provider_peer_id,
"refusing everything after the first byte",
))
};
}
let served = self.served_blob(provider_peer_id);
let start = (offset as usize).min(served.len());
let mut window = length as usize;
if self.overserve {
window += self.chunk_size.max(1);
}
let end = (start + window).min(served.len());
let mut bytes = served[start..end].to_vec();
if self.tamper_peer.as_deref() == Some(provider_peer_id) && !bytes.is_empty() {
bytes[0] ^= 0xFF;
}
Ok(bytes)
}
}
fn hex_sha256(bytes: &[u8]) -> String {
use sha2::{Digest, Sha256};
let digest = Sha256::digest(bytes);
let mut out = String::with_capacity(64);
for b in digest {
out.push(char::from_digit((b >> 4) as u32, 16).unwrap());
out.push(char::from_digit((b & 0x0f) as u32, 16).unwrap());
}
out
}
#[derive(Debug, Clone, Copy, Default)]
pub struct RejectAllModuleAnchor;
#[async_trait]
impl crate::module::ModuleAnchorVerifier for RejectAllModuleAnchor {
async fn verify_module_anchor(
&self,
_module: &dyn crate::module::ModuleReader,
_store_id: &str,
_root: &str,
) -> crate::module::ModuleAnchor {
crate::module::ModuleAnchor::NotAnchored
}
}
#[derive(Debug, Clone, Copy, Default)]
pub struct UnreachableChainAnchor;
#[async_trait]
impl crate::module::ModuleAnchorVerifier for UnreachableChainAnchor {
async fn verify_module_anchor(
&self,
_module: &dyn crate::module::ModuleReader,
_store_id: &str,
_root: &str,
) -> crate::module::ModuleAnchor {
crate::module::ModuleAnchor::Unavailable("chain source unreachable".into())
}
}
#[derive(Debug, Clone)]
pub struct OnlyThisModuleAnchor(Vec<u8>);
impl OnlyThisModuleAnchor {
pub fn new(module: Vec<u8>) -> Self {
OnlyThisModuleAnchor(module)
}
}
const ANCHOR_READ_WINDOW: u64 = 7;
#[async_trait]
impl crate::module::ModuleAnchorVerifier for OnlyThisModuleAnchor {
async fn verify_module_anchor(
&self,
module: &dyn crate::module::ModuleReader,
_store_id: &str,
_root: &str,
) -> crate::module::ModuleAnchor {
let mut seen: Vec<u8> = Vec::new();
while (seen.len() as u64) < module.len() {
let want = ANCHOR_READ_WINDOW.min(module.len() - seen.len() as u64);
match module.read_at(seen.len() as u64, want).await {
Ok(bytes) => seen.extend_from_slice(&bytes),
Err(e) => return crate::module::ModuleAnchor::Unavailable(e.to_string()),
}
}
if seen == self.0 {
crate::module::ModuleAnchor::Anchored
} else {
crate::module::ModuleAnchor::NotAnchored
}
}
}
pub fn mock_providers(n: u8, content: &ContentId) -> Vec<ProviderRecord> {
(1..=n).map(|i| mock_provider(i, content)).collect()
}
#[cfg(test)]
mod tests {
use super::*;
#[tokio::test]
async fn honest_transport_serves_correct_slice() {
let content = MockContent::even(30, 3);
let t = MockRangeTransport::new(content.clone());
let cid = mock_content_id();
let p = mock_provider(1, &cid);
let req = dig_nat::RangeRequest::resource("s", "r", 10, 10);
let got = t.fetch_range(&p, &req).await.unwrap();
assert_eq!(got.bytes, content.bytes[10..20]);
assert_eq!(got.meta.chunk_lens, Some(content.chunk_lens.clone()));
assert_eq!(t.attempts_at(10).await, 1);
assert_eq!(t.attempts_for(&mock_peer_hex(1)).await, 1);
}
#[tokio::test]
async fn behaviors_corrupt_and_truncate() {
let content = MockContent::even(30, 3);
let cid = mock_content_id();
let p = mock_provider(2, &cid);
let hex = mock_peer_hex(2);
let t = MockRangeTransport::new(content.clone());
t.set_behavior(&hex, Behavior::Truncate).await;
let req = dig_nat::RangeRequest::resource("s", "r", 0, 10);
let got = t.fetch_range(&p, &req).await.unwrap();
assert_eq!(got.bytes.len(), 9);
let t2 = MockRangeTransport::new(content.clone());
t2.set_behavior(&hex, Behavior::Corrupt).await;
let got2 = t2.fetch_range(&p, &req).await.unwrap();
assert_eq!(got2.bytes.len(), 10);
assert_ne!(got2.bytes, content.bytes[0..10]);
}
#[tokio::test]
async fn drop_after_fails_late() {
let content = MockContent::even(30, 3);
let cid = mock_content_id();
let p = mock_provider(3, &cid);
let hex = mock_peer_hex(3);
let t = MockRangeTransport::new(content);
t.set_behavior(&hex, Behavior::DropAfter(1)).await;
let req = dig_nat::RangeRequest::resource("s", "r", 0, 10);
assert!(t.fetch_range(&p, &req).await.is_ok()); assert!(t.fetch_range(&p, &req).await.is_err()); }
#[tokio::test]
async fn unavailable_source_reports_not_held() {
let content = MockContent::even(30, 3);
let cid = mock_content_id();
let p = mock_provider(4, &cid);
let hex = mock_peer_hex(4);
let t = MockRangeTransport::new(content);
t.set_behavior(&hex, Behavior::Unavailable).await;
let resp = t
.query_availability(&p, vec![AvailabilityItem::store("s")])
.await
.unwrap();
assert!(!resp.items[0].available);
}
#[tokio::test]
async fn scripted_locator_advances_batches() {
let cid = mock_content_id();
let loc = MockProviderLocator::scripted(vec![
vec![mock_provider(1, &cid)],
vec![mock_provider(1, &cid), mock_provider(2, &cid)],
]);
assert_eq!(loc.find_providers(&cid).await.unwrap().len(), 1);
assert_eq!(loc.find_providers(&cid).await.unwrap().len(), 2);
assert_eq!(loc.find_providers(&cid).await.unwrap().len(), 2); assert_eq!(loc.call_count().await, 3);
}
}