use std::time::Duration;
use anyhow::{Result, bail};
use sha2::{Digest, Sha256};
use tracing::warn;
const MAX_ATTEMPTS: u32 = 5;
async fn with_backoff<F, Fut>(mut once: F) -> bool
where
F: FnMut(u32) -> Fut,
Fut: std::future::Future<Output = bool>,
{
let mut delay = Duration::from_millis(200);
for attempt in 1..=MAX_ATTEMPTS {
if once(attempt).await {
return true;
}
if attempt < MAX_ATTEMPTS {
tokio::time::sleep(delay).await;
delay = (delay * 2).min(Duration::from_secs(3));
}
}
false
}
fn nats_digest(hasher: Sha256) -> String {
use base64::Engine;
let b64 = base64::engine::general_purpose::URL_SAFE.encode(hasher.finalize());
format!("SHA-256={b64}")
}
pub async fn verify_http_readback(
client: &reqwest::Client,
url: &reqwest::Url,
key: &str,
expected_digest: Option<&str>,
expected_size: u64,
) -> Result<()> {
let ok = with_backoff(|attempt| async move {
match http_read_and_hash(client, url).await {
Ok((got_digest, got_size)) => {
if got_size == expected_size && expected_digest.is_none_or(|d| d == got_digest) {
return true;
}
warn!(
attempt,
?expected_digest,
got_digest = %got_digest,
expected_size,
got_size,
"publish read-back mismatch — object store not yet consistent (#277)"
);
}
Err(e) => {
warn!(attempt, error = %e, "publish read-back: download failed (transient?)");
}
}
false
})
.await;
if ok {
return Ok(());
}
bail!(
"publish read-back: {key:?} still inconsistent after {MAX_ATTEMPTS} attempts \
— JetStream race (#277). Retry the publish in a few seconds; if it persists, check broker health."
);
}
async fn http_read_and_hash(client: &reqwest::Client, url: &reqwest::Url) -> Result<(String, u64)> {
use futures::StreamExt;
let resp = client.get(url.clone()).send().await?;
if !resp.status().is_success() {
bail!("GET {url}: {}", resp.status());
}
let mut hasher = Sha256::new();
let mut total: u64 = 0;
let mut stream = resp.bytes_stream();
while let Some(chunk) = stream.next().await {
let chunk = chunk?;
hasher.update(&chunk);
total += chunk.len() as u64;
}
Ok((nats_digest(hasher), total))
}