use std::fmt;
use std::sync::Arc;
use sha1::Sha1;
use sha2::{Digest as _, Sha256, Sha512};
use crate::App;
use crate::artifacts::content::{ContentError, ContentKey};
use crate::concurrency::{Resolution, SingleFlight, Slot};
use crate::http::error::ApiError;
use crate::policy::{Candidate, Decision, Digest, Ecosystem, HashAlgorithm};
use crate::store::cache::{CachedProject, ProjectKey};
use crate::store::rows::{ReferenceId, ReferenceRow};
use crate::store::{PinOutcome, StoreError};
use crate::upstream::{ArtifactRequest, UpstreamError};
#[derive(Clone, Copy, Debug)]
pub struct VerifiedContent {
pub key: ContentKey,
pub sha256: [u8; 32],
pub sha512: [u8; 64],
pub size: u64,
}
#[derive(Clone, Copy, Debug)]
struct Computed {
sha256: [u8; 32],
sha512: [u8; 64],
sha1: [u8; 20],
size: u64,
}
impl Computed {
fn digest(&self, algorithm: HashAlgorithm) -> Digest {
let bytes: Box<[u8]> = match algorithm {
HashAlgorithm::Sha256 => Box::from(self.sha256.as_slice()),
HashAlgorithm::Sha512 => Box::from(self.sha512.as_slice()),
HashAlgorithm::Sha1 => Box::from(self.sha1.as_slice()),
};
Digest { algorithm, bytes }
}
fn pinned_digests(&self) -> Vec<Digest> {
vec![
self.digest(HashAlgorithm::Sha256),
self.digest(HashAlgorithm::Sha512),
]
}
}
#[derive(Clone)]
enum SlotOutcome {
Verified(VerifiedContent),
Failed(Arc<DownloadError>),
}
pub type RefreshOutcome = Result<Arc<CachedProject>, ApiError>;
#[derive(Default)]
pub struct DownloadCoordinator {
transfers: SingleFlight<ReferenceId, SlotOutcome>,
metadata: SingleFlight<ProjectKey, RefreshOutcome>,
}
impl DownloadCoordinator {
pub fn new() -> DownloadCoordinator {
DownloadCoordinator::default()
}
fn transfers(&self) -> &SingleFlight<ReferenceId, SlotOutcome> {
&self.transfers
}
pub fn metadata(&self) -> &SingleFlight<ProjectKey, RefreshOutcome> {
&self.metadata
}
pub fn waiting_on(&self, id: &ReferenceId) -> usize {
self.transfers.waiting_on(id)
}
pub fn publishing(&self, id: &ReferenceId) -> bool {
self.transfers.publishing(id)
}
}
pub async fn fetch(
app: &Arc<App>,
row: &ReferenceRow,
) -> Result<VerifiedContent, Arc<DownloadError>> {
let (slot, leader, _waiter) = app.downloads.transfers().join(row.id);
if leader {
let app = Arc::clone(app);
let row = row.clone();
let started = Arc::clone(&slot);
tokio::spawn(async move { run(app, row, started).await });
}
match slot.wait().await {
Some(SlotOutcome::Verified(content)) => Ok(content),
Some(SlotOutcome::Failed(err)) => Err(err),
None => Err(Arc::new(DownloadError::Cancelled)),
}
}
async fn run(app: Arc<App>, row: ReferenceRow, slot: Arc<Slot<SlotOutcome>>) {
let mut resolution = Resolution::new(
app.downloads.transfers().clone(),
row.id,
Arc::clone(&slot),
SlotOutcome::Failed(Arc::new(DownloadError::Aborted)),
row.id.to_hex(),
);
let outcome = match transfer(&app, &row, &slot).await {
Ok(content) => SlotOutcome::Verified(content),
Err(err) => SlotOutcome::Failed(Arc::new(err)),
};
resolution.answer(outcome);
}
async fn transfer(
app: &App,
row: &ReferenceRow,
slot: &Slot<SlotOutcome>,
) -> Result<VerifiedContent, DownloadError> {
let _permit = app
.limits
.download_permit()
.ok_or(DownloadError::Overloaded)?;
let cap = app.config.max_artifact_bytes.get();
let body = app
.transport
.open_artifact(ArtifactRequest {
url: row.reference.upstream_url.clone(),
max_bytes: cap,
})
.await
.map_err(DownloadError::Upstream)?;
if body.declared_length.is_some_and(|declared| declared > cap) {
return Err(DownloadError::TooLarge { limit: cap });
}
let mut temp = app
.content
.create_temp(cap)
.await
.map_err(DownloadError::Capacity)?;
let mut sha256 = Sha256::new();
let mut sha512 = Sha512::new();
let mut sha1 = Sha1::new();
let mut size = 0u64;
let mut stream = body.stream;
loop {
let next = tokio::select! {
biased;
() = slot.cancelled() => return Err(DownloadError::Cancelled),
next = futures_util::StreamExt::next(&mut stream) => next,
};
let Some(chunk) = next else { break };
let chunk = chunk.map_err(DownloadError::Upstream)?;
size = size.saturating_add(chunk.len() as u64);
if size > cap {
return Err(DownloadError::TooLarge { limit: cap });
}
sha256.update(&chunk);
sha512.update(&chunk);
sha1.update(&chunk);
temp.write_all(&chunk)
.await
.map_err(DownloadError::Capacity)?;
}
let computed = Computed {
sha256: sha256.finalize().into(),
sha512: sha512.finalize().into(),
sha1: sha1.finalize().into(),
size,
};
if let Some(declared) = body.declared_length
&& declared != computed.size
{
return Err(DownloadError::SizeMismatch {
declared,
actual: computed.size,
});
}
verify_advertised(&row.reference.expected, &computed)?;
match app
.store()
.pin_computed_digests(row.id, computed.sha256, computed.sha512, computed.size)
.await
.map_err(DownloadError::Storage)?
{
PinOutcome::Established => {
app.store().caches().projects.remove(&ProjectKey::new(
row.reference.ecosystem,
row.reference.name.as_str(),
));
}
PinOutcome::MatchedExisting => {}
PinOutcome::Conflict => {
tracing::error!(
reference = %row.id.to_hex(),
package = %row.reference.name,
version = %row.reference.version,
computed = %computed.digest(HashAlgorithm::Sha256),
"the same reference now serves different bytes; keeping the original pins \
and discarding these"
);
return Err(DownloadError::PinConflict);
}
}
let now = app.clock.now_utc_micros();
let snapshot = app.blocklist();
let pinned = computed.pinned_digests();
let decision = crate::osv::evaluate(
&app.osv,
snapshot.as_deref(),
now,
app.config.cooldown_seconds,
&Candidate {
ecosystem: row.reference.ecosystem,
name: &row.reference.name,
version: &row.reference.version,
publication: row.publication_time(),
advertised_digests: &row.reference.expected,
pinned_digests: &pinned,
},
)
.await;
if decision != Decision::Allow {
if matches!(decision, Decision::Deny(_)) {
bump_digest_generation(app, row.reference.ecosystem, &row.reference.name).await;
}
return Err(DownloadError::Refused(decision));
}
if !slot.begin_publishing() {
return Err(DownloadError::Cancelled);
}
let key = ContentKey::from_sha256(computed.sha256);
let published = app
.content
.publish(temp, &key)
.await
.map_err(DownloadError::Capacity)?;
tracing::debug!(
reference = %row.id.to_hex(),
content = %key,
size = computed.size,
reused = published.reused,
steps = ?published.steps,
"published verified artifact bytes"
);
app.store()
.publish_content(key, computed.sha512, computed.size, row.id, now)
.await
.map_err(DownloadError::Storage)?;
Ok(VerifiedContent {
key,
sha256: computed.sha256,
sha512: computed.sha512,
size: computed.size,
})
}
async fn bump_digest_generation(app: &App, ecosystem: Ecosystem, name: &str) {
if let Err(err) = app.store().bump_digest_generation(ecosystem, name).await {
tracing::warn!(
package = %name,
error = %err,
"could not invalidate the project's rendered metadata after a computed digest \
revealed a block"
);
}
}
fn verify_advertised(expected: &[Digest], computed: &Computed) -> Result<(), DownloadError> {
let strongest = [
HashAlgorithm::Sha512,
HashAlgorithm::Sha256,
HashAlgorithm::Sha1,
]
.into_iter()
.find(|algorithm| expected.iter().any(|digest| digest.algorithm == *algorithm));
let Some(algorithm) = strongest else {
return Ok(());
};
let ours = computed.digest(algorithm);
if expected
.iter()
.filter(|digest| digest.algorithm == algorithm)
.any(|digest| *digest == ours)
{
return Ok(());
}
Err(DownloadError::IntegrityMismatch {
expected: expected
.iter()
.find(|digest| digest.algorithm == algorithm)
.cloned()
.unwrap_or_else(|| ours.clone()),
computed: ours,
})
}
#[derive(Debug)]
pub enum DownloadError {
Upstream(UpstreamError),
IntegrityMismatch {
expected: Digest,
computed: Digest,
},
PinConflict,
SizeMismatch {
declared: u64,
actual: u64,
},
TooLarge {
limit: u64,
},
Capacity(ContentError),
Refused(Decision),
Storage(StoreError),
Overloaded,
Cancelled,
Aborted,
}
impl fmt::Display for DownloadError {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
match self {
DownloadError::Upstream(err) => write!(f, "{err}"),
DownloadError::IntegrityMismatch { expected, computed } => write!(
f,
"the artifact hashes to {computed}, but upstream advertised {expected}"
),
DownloadError::PinConflict => f.write_str(
"the artifact no longer matches the digests permanently pinned for this reference",
),
DownloadError::SizeMismatch { declared, actual } => write!(
f,
"upstream declared {declared} bytes and delivered {actual}"
),
DownloadError::TooLarge { limit } => {
write!(f, "the artifact exceeds the {limit}-byte cap")
}
DownloadError::Capacity(err) => write!(f, "{err}"),
DownloadError::Refused(decision) => {
write!(f, "policy refused the verified bytes: {decision:?}")
}
DownloadError::Storage(err) => write!(f, "{err}"),
DownloadError::Overloaded => f.write_str("no artifact download permit was available"),
DownloadError::Cancelled => {
f.write_str("the last waiter left before the download was published")
}
DownloadError::Aborted => f.write_str("the artifact transfer ended without completing"),
}
}
}
impl std::error::Error for DownloadError {}
#[cfg(test)]
mod tests {
use super::*;
fn computed() -> Computed {
Computed {
sha256: [1; 32],
sha512: [2; 64],
sha1: [3; 20],
size: 10,
}
}
fn digest(algorithm: HashAlgorithm, byte: u8) -> Digest {
Digest {
algorithm,
bytes: vec![byte; algorithm.digest_len()].into_boxed_slice(),
}
}
#[test]
fn nothing_advertised_leaves_nothing_to_check() {
assert!(verify_advertised(&[], &computed()).is_ok());
}
#[test]
fn a_lone_legacy_sha1_is_what_decides() {
assert!(verify_advertised(&[digest(HashAlgorithm::Sha1, 3)], &computed()).is_ok());
assert!(matches!(
verify_advertised(&[digest(HashAlgorithm::Sha1, 9)], &computed()),
Err(DownloadError::IntegrityMismatch { .. })
));
}
#[test]
fn the_strongest_advertised_algorithm_decides_alone() {
let expected = vec![
digest(HashAlgorithm::Sha1, 3),
digest(HashAlgorithm::Sha512, 99),
];
assert!(matches!(
verify_advertised(&expected, &computed()),
Err(DownloadError::IntegrityMismatch { .. })
));
let expected = vec![
digest(HashAlgorithm::Sha1, 9),
digest(HashAlgorithm::Sha512, 2),
];
assert!(
verify_advertised(&expected, &computed()).is_ok(),
"and a mismatching weak digest does not veto a matching strong one"
);
}
fn fresh_slot() -> Slot<SlotOutcome> {
Slot::new()
}
#[test]
fn publishing_and_cancelling_cannot_both_win() {
let slot = fresh_slot();
assert!(slot.begin_publishing());
assert!(
!slot.begin_cancel(),
"a waiter leaving after publication began does not cancel it"
);
assert!(!slot.begin_publishing(), "and publication begins once");
let slot = fresh_slot();
assert!(slot.begin_cancel());
assert!(
!slot.begin_publishing(),
"a cancelled transfer never starts publishing"
);
assert!(!slot.begin_cancel(), "and is cancelled once");
}
#[test]
fn waiter_accounting_returns_to_zero_and_removes_the_slot() {
let coordinator = DownloadCoordinator::new();
let id = ReferenceId::from_bytes([9; 32]);
let (_slot, leader, first) = coordinator.transfers().join(id);
assert!(leader, "the first request starts the transfer");
let (_slot, leader, second) = coordinator.transfers().join(id);
assert!(!leader, "the second joins it");
assert_eq!(coordinator.waiting_on(&id), 2);
drop(second);
assert_eq!(
coordinator.waiting_on(&id),
1,
"one leaving leaves the other"
);
drop(first);
assert_eq!(coordinator.waiting_on(&id), 0);
let (_slot, leader, _third) = coordinator.transfers().join(id);
assert!(
leader,
"and the next request starts a new transfer rather than joining a cancelled one"
);
}
#[test]
fn one_of_several_entries_of_the_strongest_algorithm_is_enough() {
let expected = vec![
digest(HashAlgorithm::Sha512, 7),
digest(HashAlgorithm::Sha512, 2),
];
assert!(verify_advertised(&expected, &computed()).is_ok());
}
}