#![cfg(feature = "blob")]
use std::sync::{Arc, Mutex};
use std::time::Duration;
use zenkey_fleet::zblob::{self, BlobServer, BlobSpec, MemoryBlobSource};
use zenkey_fleet::{BlobFetchSpec, BlobTarget, blob_fetch, blob_probe};
use zenoh::qos::Priority;
const ID: &str = "01jqz3demo0001";
const BASE: &str = "";
const TIMEOUT: Duration = Duration::from_secs(5);
mod util;
use util::peer_pair;
fn payload(len: usize, seed: u64) -> Vec<u8> {
let mut state = seed.wrapping_mul(6364136223846793005).wrapping_add(1);
(0..len)
.map(|_| {
state = state.wrapping_mul(6364136223846793005).wrapping_add(1);
(state >> 33) as u8
})
.collect()
}
fn target() -> BlobTarget {
BlobTarget::parse(ID).expect("the demo id is a plain chunk")
}
fn prefix_at(origin: &str) -> String {
let origin = zenkey::grammar::Origin::Host(zenkey::HostId::parse(origin).unwrap());
zenkey::grammar::blob_tier_prefix(&origin, zenkey::grammar::BlobTier::Artifact)
.as_str()
.to_string()
}
async fn wait_routable(session: &zenoh::Session, key: &str) {
for _ in 0..50 {
let answers = zenkey_fleet::fleet_get(
&zenkey_fleet::Fleet::new(session, BASE),
key,
&zenkey_fleet::GetOpts::new(TIMEOUT),
)
.await
.expect("settle get");
if !answers.is_empty() {
return;
}
tokio::time::sleep(Duration::from_millis(100)).await;
}
panic!("fixture never became routable: {key}");
}
async fn serve(
session: &zenoh::Session,
origin: &str,
data: Vec<u8>,
) -> (zblob::ServerHandle, zblob::Manifest) {
let server = BlobServer::new(
session,
zblob::ServePrefix::new(prefix_at(origin)).expect("concrete prefix"),
);
let manifest = server
.register_source(
BlobSpec::new(ID).filename("demo-bundle.bin"),
Arc::new(MemoryBlobSource::new(data)),
)
.await
.expect("register");
(server.spawn().await.expect("spawn"), manifest)
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn a_probe_names_every_holder_and_a_fetch_names_one() {
let (serving, asking) = peer_pair().await;
let data = payload(200_000, 7);
let (a, manifest) = serve(&serving, "h-aaaaaaaaaaaa", data.clone()).await;
let (b, _) = serve(&serving, "h-bbbbbbbbbbbb", data.clone()).await;
for origin in ["h-aaaaaaaaaaaa", "h-bbbbbbbbbbbb"] {
wait_routable(&asking, &zblob::keys::manifest_key(&prefix_at(origin), ID)).await;
}
let report = blob_probe(
&zenkey_fleet::Fleet::new(&asking, BASE),
&target(),
&[],
TIMEOUT,
)
.await
.expect("probe");
let origins: Vec<&str> = report.holders.iter().map(|h| h.origin.as_str()).collect();
assert_eq!(
origins,
vec!["h-aaaaaaaaaaaa", "h-bbbbbbbbbbbb"],
"a probe names every holder: {report:?}"
);
assert_eq!(
report.asked.len(),
2,
"have and manifest: {:?}",
report.asked
);
assert!(report.asked.iter().any(|s| s.ends_with("/have")));
assert!(report.asked.iter().any(|s| s.ends_with("/manifest")));
assert!(report.not_probed.is_none());
assert_eq!(report.roots, vec![manifest.root.to_string()]);
for holder in &report.holders {
let avail = holder.availability.as_ref().expect("a have reply");
assert!(avail.complete, "a full server answers all-ones: {avail:?}");
assert_eq!(
holder.manifest.as_ref().map(|m| m.total_len),
Some(data.len() as u64)
);
assert!(holder.unreadable.is_none(), "{:?}", holder.unreadable);
assert!(holder.key.contains(&holder.origin), "{}", holder.key);
}
let dir = tempfile::tempdir().unwrap();
let dest = dir.path().join("demo.bin");
let spec = BlobFetchSpec {
timeout: TIMEOUT,
root: Some(zenkey::ContentHash::parse(&manifest.root.to_string()).unwrap()),
..Default::default()
};
let fetched = blob_fetch(
&zenkey_fleet::Fleet::new(&asking, BASE),
"h-bbbbbbbbbbbb",
&target(),
&dest,
&spec,
&|_| {},
)
.await
.expect("fetch");
assert_eq!(
fetched.origin, "h-bbbbbbbbbbbb",
"one origin, the chosen one"
);
assert_eq!(fetched.rejected, 0, "nothing failed verification");
assert!(fetched.root_pinned, "the root was pinned, not TOFU");
assert_eq!(fetched.priority, "data-low");
assert_eq!(
std::fs::read(&dest).unwrap(),
data,
"the bytes on disk are the bytes served"
);
a.shutdown().await.unwrap();
b.shutdown().await.unwrap();
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn a_disagreeing_root_is_named_not_averaged() {
let (serving, asking) = peer_pair().await;
let honest = payload(120_000, 11);
let rogue = payload(120_000, 12);
let (a, manifest) = serve(&serving, "h-aaaaaaaaaaaa", honest.clone()).await;
let (c, rogue_manifest) = serve(&serving, "h-cccccccccccc", rogue).await;
for origin in ["h-aaaaaaaaaaaa", "h-cccccccccccc"] {
wait_routable(&asking, &zblob::keys::manifest_key(&prefix_at(origin), ID)).await;
}
assert_ne!(
manifest.root, rogue_manifest.root,
"the fixture must disagree"
);
let report = blob_probe(
&zenkey_fleet::Fleet::new(&asking, BASE),
&target(),
&[],
TIMEOUT,
)
.await
.expect("probe");
assert_eq!(report.holders.len(), 2);
assert_eq!(
report.roots.len(),
2,
"two roots under one id, both reported: {:?}",
report.roots
);
let dir = tempfile::tempdir().unwrap();
let dest = dir.path().join("demo.bin");
let spec = BlobFetchSpec {
timeout: TIMEOUT,
root: Some(zenkey::ContentHash::parse(&manifest.root.to_string()).unwrap()),
..Default::default()
};
let err = blob_fetch(
&zenkey_fleet::Fleet::new(&asking, BASE),
"h-cccccccccccc",
&target(),
&dest,
&spec,
&|_| {},
)
.await
.expect_err("a pinned root must refuse the wrong content")
.to_string();
assert!(
err.contains("h-cccccccccccc"),
"the failure must name the origin that produced it: {err}"
);
assert!(
!dest.exists(),
"verification happens before disk (RFC 07 §2.1), so nothing was written"
);
a.shutdown().await.unwrap();
c.shutdown().await.unwrap();
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn every_blob_get_rides_at_data_low() {
let (serving, asking) = peer_pair().await;
let seen: Arc<Mutex<Vec<Priority>>> = Arc::new(Mutex::new(Vec::new()));
let manifest = zblob::Manifest {
version: zblob::wire::WIRE_VERSION,
id: zblob::BlobId::new(ID).expect("valid id"),
filename: None,
total_len: 65_536 * 4,
chunk_size: 65_536,
root: zblob::Hash::of(b"not the content"),
created_ms: 0,
ext: zblob::wire::Ext::default(),
};
let manifest_bytes = zblob::wire::encode(&manifest).unwrap();
let prefix = prefix_at("h-dddddddddddd");
let manifest_key = zblob::keys::manifest_key(&prefix, ID);
let cancel = zblob::CancelToken::new();
let recorder = {
let seen = seen.clone();
let cancel = cancel.clone();
let manifest_key = manifest_key.clone();
serving
.declare_queryable(format!("{prefix}/**"))
.callback(move |query| {
seen.lock().unwrap().push(query.priority());
if query.key_expr().as_str().ends_with("/manifest")
|| query.key_expr().as_str().ends_with("/**")
{
let bytes = manifest_bytes.clone();
let key = manifest_key.clone();
tokio::spawn(async move {
let _ = query
.reply(key, bytes)
.encoding(&zblob::wire::ENC_MANIFEST)
.await;
});
} else {
cancel.cancel();
}
})
.await
.expect("recording queryable")
};
wait_routable(&asking, &manifest_key).await;
seen.lock().unwrap().clear();
let probe = blob_probe(
&zenkey_fleet::Fleet::new(&asking, BASE),
&target(),
&[],
Duration::from_secs(2),
)
.await
.expect("probe");
assert_eq!(probe.answered, 1, "the recorder answered the manifest GET");
let dir = tempfile::tempdir().unwrap();
let spec = BlobFetchSpec {
timeout: Duration::from_secs(2),
cancel,
..Default::default()
};
let _ = tokio::time::timeout(
Duration::from_secs(30),
blob_fetch(
&zenkey_fleet::Fleet::new(&asking, BASE),
"h-dddddddddddd",
&target(),
&dir.path().join("demo.bin"),
&spec,
&|_| {},
),
)
.await;
let priorities = seen.lock().unwrap().clone();
assert!(
priorities.len() >= 3,
"expected have + manifest + at least one slice query, saw {}",
priorities.len()
);
assert!(
priorities.iter().all(|p| *p == Priority::DataLow),
"every @blob GET must ride at data-low (RFC 07 §2.6): {priorities:?}"
);
recorder.undeclare().await.unwrap();
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn tier_two_is_probed_and_a_foreign_algo_says_why_not() {
let (_serving, asking) = peer_pair().await;
let hash = "0123456789abcdef0123456789abcdef0123456789abcdef0123456789abcdef";
let target = BlobTarget::parse(&format!("store/blake3/{hash}")).unwrap();
let report = blob_probe(
&zenkey_fleet::Fleet::new(&asking, BASE),
&target,
&[],
Duration::from_secs(1),
)
.await
.expect("probe");
assert!(
report.not_probed.is_none(),
"blake3 must be probed: {report:?}"
);
assert_eq!(
report.asked,
vec!["v1/*/@blob/store/blake3/have".to_string()],
"the §2.4 probe key under the empty base, and no wider"
);
assert!(report.holders.is_empty());
let target = BlobTarget::parse(&format!("store/sha256/{hash}")).unwrap();
let report = blob_probe(
&zenkey_fleet::Fleet::new(&asking, BASE),
&target,
&[],
Duration::from_secs(1),
)
.await
.expect("probe");
assert!(report.asked.is_empty(), "nothing was asked");
assert!(report.holders.is_empty());
let why = report.not_probed.expect("an unasked probe must say why");
assert!(why.contains("blake3"), "{why}");
assert!(why.contains("2.4"), "{why}");
}