use std::path::Path;
use std::time::Duration;
use crate::{Error, Result};
use zenkey::grammar::{self, ContentHash, Origin};
use zenkey::{RegistrySlice, RemoteOrigin, ServiceOrigin};
use zenoh::qos::Priority;
use super::{BlobTarget, declared_by};
use crate::bus::query::{Answer, FleetAnswer, GetOpts, fleet_get};
use crate::report::{
BlobAvailability, BlobFetchReport, BlobHolder, BlobManifest, BlobProbeReport, BlobProgress,
CallError,
};
pub const FETCH_PRIORITY: Priority = Priority::DataLow;
pub struct BlobFetchSpec {
pub timeout: Duration,
pub overwrite: bool,
pub root: Option<ContentHash>,
pub cancel: zblob::CancelToken,
}
impl Default for BlobFetchSpec {
fn default() -> Self {
BlobFetchSpec {
timeout: Duration::from_secs(30),
overwrite: false,
root: None,
cancel: zblob::CancelToken::new(),
}
}
}
pub async fn blob_probe(
fleet: &crate::Fleet<'_>,
target: &BlobTarget,
slices: &[RegistrySlice],
timeout: Duration,
) -> Result<BlobProbeReport> {
let base = fleet.base();
let tier = target.tier();
let declared = declared_by(slices, tier);
let Some(id) = target.artifact_id() else {
return probe_tier2(fleet, target, declared, slices.len(), timeout).await;
};
let prefix = target.probe_prefix();
let have = grammar::with_base(base, zblob::keys::availability_key(prefix.as_str(), id));
let manifest = grammar::with_base(base, zblob::keys::manifest_key(prefix.as_str(), id));
let asked = vec![have.clone(), manifest.clone()];
let bulk = GetOpts::new(timeout).priority(FETCH_PRIORITY);
let (have_answers, manifest_answers) = tokio::join!(
fleet_get(fleet, &have, &bulk),
fleet_get(fleet, &manifest, &bulk),
);
let mut holders: Vec<BlobHolder> = Vec::new();
for (answers, kind) in [
(have_answers?, Endpoint::Have),
(manifest_answers?, Endpoint::Manifest),
] {
for answer in answers {
fold(&mut holders, base, kind, answer);
}
}
holders.sort_by(|a, b| a.origin.cmp(&b.origin));
let mut roots: Vec<String> = holders
.iter()
.filter_map(|h| h.manifest.as_ref().map(|m| m.root.clone()))
.collect();
roots.sort();
roots.dedup();
Ok(BlobProbeReport {
target: target.spelling(),
tier: tier.chunk().to_string(),
asked,
not_probed: None,
answered: holders.len(),
holders,
roots,
declared_by: declared,
slices_considered: slices.len(),
})
}
async fn probe_tier2(
fleet: &crate::Fleet<'_>,
target: &BlobTarget,
declared: Vec<String>,
slices_considered: usize,
timeout: Duration,
) -> Result<BlobProbeReport> {
let base = fleet.base();
let tier = target.tier();
let probe_prefix = grammar::with_base(base, target.probe_prefix().as_str());
let report = |asked: Vec<String>, not_probed: Option<String>, holders: Vec<BlobHolder>| {
BlobProbeReport {
target: target.spelling(),
tier: tier.chunk().to_string(),
asked,
not_probed,
answered: holders.len(),
holders,
roots: Vec::new(),
declared_by: declared.clone(),
slices_considered,
}
};
match target {
BlobTarget::Store { algo, hash } => {
if algo != zblob::Hash::ALGO {
return Ok(report(
Vec::new(),
Some(format!(
"the reference client speaks `{}` only, so a `{algo}` chunk cannot be probed by this build (RFC 07 §2.4 — dedup and probing are per-algorithm)",
zblob::Hash::ALGO
)),
Vec::new(),
));
}
let parsed: zblob::Hash = hash
.as_str()
.parse()
.map_err(|e| Error::unaskable_from(hash.to_string(), e))?;
let have_key = zblob::keys::store_have_key(&probe_prefix, zblob::HashAlgo::Blake3);
let want = zblob::wire::encode(&zblob::wire::WantList::new(vec![parsed]))
.map_err(|e| Error::Internal(format!("encoding the want-list: {e}")))?;
let answers = fleet_get(
fleet,
&have_key,
&GetOpts::new(timeout)
.payload(Some(want))
.priority(FETCH_PRIORITY),
)
.await?;
let holders = fold_tier2(base, answers, |bytes| {
let bits: zblob::wire::HaveBits = zblob::wire::decode(bytes)
.map_err(|e| format!("undecodable have bitfield: {e}"))?;
bits.validate(1)
.map_err(|e| format!("invalid have bitfield: {e}"))?;
let held = bits.is_set(0);
Ok((
BlobAvailability {
chunk_count: 1,
have: u32::from(held),
complete: held,
},
None,
))
});
Ok(report(vec![have_key], None, holders))
}
BlobTarget::Tree { root } => {
let parsed: zblob::Hash = root
.as_str()
.parse()
.map_err(|e| Error::unaskable_from(root.to_string(), e))?;
let have_key = zblob::keys::tree_have_key(&probe_prefix, &parsed.to_string());
let answers = fleet_get(
fleet,
&have_key,
&GetOpts::new(timeout).priority(FETCH_PRIORITY),
)
.await?;
let holders = fold_tier2(base, answers, |bytes| {
let probe: zblob::wire::TreeProbe = zblob::wire::decode(bytes)
.map_err(|e| format!("undecodable tree probe: {e}"))?;
probe
.validate()
.map_err(|e| format!("invalid tree probe: {e}"))?;
let note = (!probe.have_index && probe.chunks_present > 0).then(|| {
"holds chunks but not the index — an index fetch from this origin will fail"
.to_string()
});
Ok((
BlobAvailability {
chunk_count: probe.chunks_total,
have: probe.chunks_present,
complete: probe.have_index && probe.chunks_present == probe.chunks_total,
},
note,
))
});
Ok(report(vec![have_key], None, holders))
}
BlobTarget::Artifact { .. } => Err(Error::Internal(
"tier-1 target reached the tier-2 probe path — a bug in blob_probe".into(),
)),
}
}
fn fold_tier2(
base: &str,
answers: Vec<FleetAnswer>,
decode: impl Fn(&[u8]) -> Result<(BlobAvailability, Option<String>), String>,
) -> Vec<BlobHolder> {
let mut holders: Vec<BlobHolder> = Vec::new();
for answer in answers {
let origin = attribute(base, &answer);
if holders.iter().any(|h| h.origin == origin) {
continue;
}
let mut holder = BlobHolder {
origin,
key: answer.key.clone(),
availability: None,
manifest: None,
note: None,
unreadable: None,
error: None,
};
match answer.answer {
Answer::Error { name, message } => {
holder.error = Some(CallError { name, message });
}
Answer::Value(payload) => match decode(&payload.to_bytes()) {
Ok((availability, note)) => {
holder.availability = Some(availability);
holder.note = note;
}
Err(why) => {
let declared = answer.encoding.as_deref().unwrap_or("(none)");
holder.unreadable = Some(format!("{why} (encoding `{declared}`)"));
}
},
}
holders.push(holder);
}
holders.sort_by(|a, b| a.origin.cmp(&b.origin));
holders
}
async fn write_atomically(dest: &Path, bytes: &[u8]) -> Result<()> {
use tokio::io::AsyncWriteExt;
static TMP_SEQ: std::sync::atomic::AtomicU64 = std::sync::atomic::AtomicU64::new(0);
let name = dest
.file_name()
.map(|n| n.to_string_lossy().into_owned())
.ok_or_else(|| Error::unaskable(dest.display().to_string(), "names no file to write"))?;
if let Some(parent) = dest.parent().filter(|p| !p.as_os_str().is_empty()) {
tokio::fs::create_dir_all(parent)
.await
.map_err(|e| Error::io(parent, e))?;
}
let seq = TMP_SEQ.fetch_add(1, std::sync::atomic::Ordering::Relaxed);
let tmp = dest.with_file_name(format!(".{name}.{}.{seq}.zenkey-tmp", std::process::id()));
let write = async {
let mut f = tokio::fs::File::create(&tmp).await?;
f.write_all(bytes).await?;
f.sync_all().await
};
if let Err(e) = write.await {
let _ = tokio::fs::remove_file(&tmp).await;
return Err(Error::io(&tmp, e));
}
if let Err(e) = tokio::fs::rename(&tmp, dest).await {
let _ = tokio::fs::remove_file(&tmp).await;
return Err(Error::io(dest, e));
}
Ok(())
}
#[derive(Clone, Copy)]
enum Endpoint {
Have,
Manifest,
}
fn fold(holders: &mut Vec<BlobHolder>, base: &str, kind: Endpoint, answer: FleetAnswer) {
let origin = attribute(base, &answer);
let idx = match holders.iter().position(|h| h.origin == origin) {
Some(i) => i,
None => {
holders.push(BlobHolder {
origin: origin.clone(),
key: answer.key.clone(),
availability: None,
manifest: None,
note: None,
unreadable: None,
error: None,
});
holders.len() - 1
}
};
let holder = &mut holders[idx];
if holder.key.is_empty() {
holder.key = answer.key.clone();
}
match answer.answer {
Answer::Error { name, message } => {
holder.error = Some(CallError { name, message });
}
Answer::Value(payload) => {
let bytes = payload.to_bytes();
let (want, decoded) = match kind {
Endpoint::Have => (
&zblob::wire::ENC_AVAIL,
decode_have(&bytes).map(|a| holder.availability = Some(a)),
),
Endpoint::Manifest => (
&zblob::wire::ENC_MANIFEST,
decode_manifest(&bytes).map(|m| holder.manifest = Some(m)),
),
};
if let Err(why) = decoded {
let declared = answer.encoding.as_deref().unwrap_or("(none)");
holder.unreadable =
Some(format!("{why} (encoding `{declared}`, expected `{want}`)"));
}
}
}
}
fn attribute(base: &str, answer: &FleetAnswer) -> String {
if answer.origin != "?" {
return answer.origin.clone();
}
let stripped = answer
.key
.strip_prefix(base)
.map(|s| s.trim_start_matches('/'))
.unwrap_or(&answer.key);
stripped
.split('/')
.nth(1)
.filter(|c| !c.is_empty())
.unwrap_or("?")
.to_string()
}
fn decode_have(bytes: &[u8]) -> Result<BlobAvailability, String> {
let avail: zblob::wire::Availability =
zblob::wire::decode(bytes).map_err(|e| format!("undecodable availability: {e}"))?;
Ok(BlobAvailability {
chunk_count: avail.chunk_count,
have: avail.count(),
complete: avail.count() == avail.chunk_count,
})
}
fn decode_manifest(bytes: &[u8]) -> Result<BlobManifest, String> {
let m: zblob::Manifest =
zblob::wire::decode(bytes).map_err(|e| format!("undecodable manifest: {e}"))?;
let chunk_count = m.chunk_count().unwrap_or(0);
Ok(BlobManifest {
chunk_count,
id: m.id.to_string(),
filename: m.filename,
total_len: m.total_len,
chunk_size: m.chunk_size,
root: m.root.to_string(),
created_ms: m.created_ms,
})
}
pub async fn blob_fetch(
fleet: &crate::Fleet<'_>,
origin: &str,
target: &BlobTarget,
dest: &Path,
spec: &BlobFetchSpec,
on_progress: &(dyn Fn(BlobProgress) + Send + Sync),
) -> Result<BlobFetchReport> {
let (session, base) = (fleet.session(), fleet.base());
let origin = parse_origin(origin)?;
let Some(id) = target.artifact_id() else {
return fetch_tier2(fleet, &origin, target, dest, spec, on_progress).await;
};
let prefix = grammar::with_base(base, target.prefix_at(&origin).as_str());
let key = grammar::with_base(base, target.key_at(&origin)?.as_str());
let prefix = zblob::QueryPrefix::new(prefix).map_err(|e| {
Error::unaskable(
format!("{}'s artifact prefix", origin.chunk()),
format!("is not queryable: {e}"),
)
})?;
let client = zblob::BlobClient::builder(session, prefix)
.query_timeout(spec.timeout)
.overwrite(if spec.overwrite {
zblob::Overwrite::Replace
} else {
zblob::Overwrite::Refuse
})
.build();
let request = match &spec.root {
Some(root) => {
let parsed: zblob::Hash = root
.as_str()
.parse()
.map_err(|e| Error::unaskable_from(root.to_string(), e))?;
zblob::DownloadRequest::pinned(id, parsed)
}
None => zblob::DownloadRequest::new(id),
};
let root_pinned = request.expected_root.is_some();
let sink = move |p: zblob::Progress| on_progress(translate(p));
let stats = client
.download_to(&request, dest)
.progress(&sink)
.cancel(&spec.cancel)
.await
.map_err(|e| Error::bus("fetch", origin.chunk(), e.to_string()))?;
Ok(BlobFetchReport {
origin: origin.chunk().to_string(),
key,
dest: dest.display().to_string(),
bytes: stats.bytes_fetched,
chunks: stats.chunks_fetched,
chunks_resumed: stats.chunks_resumed,
rejected: stats.rejected,
retries: stats.retries,
elapsed_ms: stats.elapsed.as_millis() as u64,
root: request
.expected_root
.map(|r| r.to_string())
.unwrap_or_default(),
root_pinned,
priority: priority_name(FETCH_PRIORITY).to_string(),
})
}
async fn fetch_tier2(
fleet: &crate::Fleet<'_>,
origin: &Origin,
target: &BlobTarget,
dest: &Path,
spec: &BlobFetchSpec,
on_progress: &(dyn Fn(BlobProgress) + Send + Sync),
) -> Result<BlobFetchReport> {
let (session, base) = (fleet.session(), fleet.base());
let started = std::time::Instant::now();
match target {
BlobTarget::Store { algo, hash } => {
if algo != zblob::Hash::ALGO {
return Err(Error::unaskable(
target.spelling(),
format!(
"cannot be fetched by this build: the reference client \
speaks `{}` only (RFC 07 §2.4 — addressing is \
per-algorithm)",
zblob::Hash::ALGO
),
));
}
if let Some(pin) = &spec.root
&& pin != hash
{
return Err(Error::unaskable(
format!("the pinned root {pin}"),
format!(
"contradicts the content address {hash}: a store fetch \
is pinned by its key (RFC 07 §2.1) — drop the pin, or \
fetch the address you mean"
),
));
}
let parsed: zblob::Hash = hash
.as_str()
.parse()
.map_err(|e| Error::unaskable_from(hash.to_string(), e))?;
let prefix_str = grammar::with_base(
base,
grammar::blob_tier_prefix(origin, grammar::BlobTier::Store).as_str(),
);
let prefix = zblob::QueryPrefix::new(prefix_str.clone())
.map_err(|e| Error::unaskable_from(prefix_str.to_string(), e))?;
let key = zblob::keys::store_key(prefix.as_str(), zblob::HashAlgo::Blake3, &parsed);
if !spec.overwrite && tokio::fs::try_exists(dest).await.unwrap_or(false) {
return Err(Error::unaskable(
dest.display().to_string(),
"already exists — pass overwrite to replace it",
));
}
let client = zblob::StoreClient::builder(session, prefix)
.query_timeout(spec.timeout)
.priority(FETCH_PRIORITY)
.build();
let bytes = match spec
.cancel
.until_cancelled(client.fetch_chunk(&parsed))
.await
{
None => {
on_progress(BlobProgress::Cancelled {
received: 0,
total: 1,
});
return Err(Error::bus("fetch", origin.chunk(), "cancelled"));
}
Some(fetched) => {
fetched.map_err(|e| Error::bus("fetch", origin.chunk(), e.to_string()))?
}
};
on_progress(BlobProgress::Chunk {
index: 0,
received: 1,
total: 1,
bytes_received: bytes.len() as u64,
});
write_atomically(dest, &bytes).await?;
on_progress(BlobProgress::Completed {
path: dest.display().to_string(),
});
Ok(BlobFetchReport {
origin: origin.chunk().to_string(),
key,
dest: dest.display().to_string(),
bytes: bytes.len() as u64,
chunks: 1,
chunks_resumed: 0,
rejected: 0,
retries: 0,
elapsed_ms: started.elapsed().as_millis() as u64,
root: hash.to_string(),
root_pinned: true,
priority: priority_name(FETCH_PRIORITY).to_string(),
})
}
BlobTarget::Tree { .. } => Err(Error::unaskable(
target.spelling(),
"is inspected, not downloaded, by this explorer: a validated index \
summary needs no content store (RFC 07 §2.3, v1.17) — the \
frontends route tree targets to the tree-index report; \
materializing a tree is the reference client's `download_tree`, \
which needs a store this build deliberately does not keep",
)),
BlobTarget::Artifact { .. } => Err(Error::Internal(
"tier-1 target reached the tier-2 fetch path — a bug in blob_fetch".into(),
)),
}
}
pub async fn blob_tree_index(
fleet: &crate::Fleet<'_>,
origin: &str,
root: &ContentHash,
timeout: Duration,
) -> Result<crate::report::BlobTreeIndexReport> {
let (session, base) = (fleet.session(), fleet.base());
let started = std::time::Instant::now();
let origin = parse_origin(origin)?;
let tree_str = grammar::with_base(
base,
grammar::blob_tier_prefix(&origin, grammar::BlobTier::Tree).as_str(),
);
let store_str = grammar::with_base(
base,
grammar::blob_tier_prefix(&origin, grammar::BlobTier::Store).as_str(),
);
let tree_prefix = zblob::QueryPrefix::new(tree_str.clone())
.map_err(|e| Error::unaskable_from(tree_str.to_string(), e))?;
let store_prefix = zblob::QueryPrefix::new(store_str.clone())
.map_err(|e| Error::unaskable_from(store_str.to_string(), e))?;
let parsed: zblob::Hash = root
.as_str()
.parse()
.map_err(|e| Error::unaskable_from(root.to_string(), e))?;
let key = zblob::keys::tree_key(tree_prefix.as_str(), root.as_str());
let client = zblob::TreeClient::builder(session, store_prefix, tree_prefix)
.query_timeout(timeout)
.build();
let index = client
.fetch_index_by_root(&parsed)
.await
.map_err(|e| Error::bus("fetch", origin.chunk(), e.to_string()))?;
Ok(crate::report::BlobTreeIndexReport {
origin: origin.chunk().to_string(),
key,
root: root.to_string(),
entries: index.entries().len(),
files: index.file_count(),
total_size: index.total_size(),
chunks: index.needed_chunk_refs().len(),
elapsed_ms: started.elapsed().as_millis() as u64,
priority: priority_name(FETCH_PRIORITY).to_string(),
})
}
fn parse_origin(origin: &str) -> Result<Origin> {
if let Some(service) = origin.strip_prefix('@') {
let _ = service;
let svc =
ServiceOrigin::new(origin).map_err(|e| Error::unaskable_from(origin.to_string(), e))?;
return Ok(Origin::Service(svc));
}
let host = RemoteOrigin::parse(origin).map_err(|e| {
Error::unaskable(
origin.to_string(),
format!(
"is not one concrete origin: {e}. A fetch names exactly one \
holder (RFC 07 §2.5) — probe first, then fetch from an origin \
the probe reported."
),
)
})?;
Ok(Origin::Host(host.host_id().clone()))
}
fn translate(p: zblob::Progress) -> BlobProgress {
match p {
zblob::Progress::Started {
total_len,
chunk_count,
} => BlobProgress::Started {
total_len,
chunk_count,
},
zblob::Progress::Resumed { received, total } => BlobProgress::Resumed { received, total },
zblob::Progress::Chunk {
index,
received,
total,
bytes_received,
} => BlobProgress::Chunk {
index,
received,
total,
bytes_received,
},
zblob::Progress::Verifying => BlobProgress::Verifying,
zblob::Progress::Completed { path } => BlobProgress::Completed {
path: path.display().to_string(),
},
zblob::Progress::Cancelled { received, total } => {
BlobProgress::Cancelled { received, total }
}
zblob::Progress::Failed { error } => BlobProgress::Failed { error },
other => BlobProgress::Failed {
error: format!("unrecognised progress event from the reference client: {other:?}"),
},
}
}
fn priority_name(p: Priority) -> &'static str {
match p {
Priority::RealTime => "real-time",
Priority::InteractiveHigh => "interactive-high",
Priority::InteractiveLow => "interactive-low",
Priority::DataHigh => "data-high",
Priority::Data => "data",
Priority::DataLow => "data-low",
Priority::Background => "background",
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn a_wildcard_origin_is_not_an_origin() {
for spelled in ["*", "**", "h-*", "", "not-a-host"] {
assert!(
parse_origin(spelled).is_err(),
"`{spelled}` must not parse as a fetch origin"
);
}
assert!(parse_origin("h-3fa9c2d41b7e").is_ok());
assert!(parse_origin("@catalog").is_ok());
}
#[test]
fn the_reported_priority_is_the_one_the_client_uses() {
assert_eq!(priority_name(FETCH_PRIORITY), "data-low");
}
#[test]
fn an_unparseable_key_still_names_its_holder() {
let answer = FleetAnswer {
origin: "?".to_string(),
key: "zensight/v1/h-3fa9c2d41b7e/@blob/artifact/NOPE/have".to_string(),
encoding: None,
attachment: None,
answer: Answer::Error {
name: "error/x".into(),
message: String::new(),
},
};
assert_eq!(attribute("zensight", &answer), "h-3fa9c2d41b7e");
}
}