pub mod content;
pub mod download;
pub mod reference;
pub mod stream;
use std::sync::Arc;
use axum::http::Method;
use crate::App;
use crate::artifacts::content::PinnedFile;
use crate::artifacts::stream::{ArtifactResponse, Authorized, RangeRequest};
use crate::http::error::ApiError;
use crate::http::logging;
use crate::policy::{BlocklistSnapshot, Candidate, Decision, DenyReason, Digest, Ecosystem};
use crate::store::StoreError;
use crate::store::rows::{ReferenceId, ReferenceRow};
use crate::{npm, osv, pypi};
pub use download::{DownloadError, VerifiedContent};
pub use reference::{ArtifactReference, InvalidReferenceId};
pub async fn serve_artifact(
app: &Arc<App>,
ecosystem: Ecosystem,
id: ReferenceId,
filename: &str,
method: &Method,
range_header: Option<&str>,
) -> Result<ArtifactResponse, ApiError> {
let snapshot = app.blocklist().ok_or(ApiError::PolicyUnavailable)?;
let now = app.clock.now_utc_micros();
if !snapshot.is_valid_at(now) {
return Err(ApiError::PolicyUnavailable);
}
let row = load_reference(app, id).await?;
logging::record_ecosystem(row.reference.ecosystem);
if row.reference.ecosystem != ecosystem {
return Err(ApiError::NotFound);
}
if row.reference.filename != filename {
return Err(ApiError::NotFound);
}
let pinned = row.pinned_digests();
match evaluate(app, &snapshot, now, &row, &pinned).await {
Decision::Allow => {}
other => return Err(decision_error(other, now)),
}
ensure_membership(app, &row).await?;
match evaluate(app, &snapshot, now, &row, &pinned).await {
Decision::Allow => {}
other => return Err(decision_error(other, now)),
}
let (file, pinned, downloaded) = open_or_download(app, &row, pinned).await?;
logging::record_cache(if downloaded {
logging::CacheStatus::Miss
} else {
logging::CacheStatus::Hit
});
let total_length = file.size();
if downloaded {
ensure_membership(app, &row).await?;
}
let snapshot = app.blocklist().ok_or(ApiError::PolicyUnavailable)?;
let now = app.clock.now_utc_micros();
let decision = evaluate(app, &snapshot, now, &row, &pinned).await;
let Some(auth) = Authorized::from_decision(decision, Some(snapshot.revision)) else {
return Err(decision_error(decision, now));
};
if method == Method::HEAD {
return Ok(ArtifactResponse::head_only(auth, total_length));
}
Ok(match stream::parse_range(range_header, total_length) {
RangeRequest::Whole => ArtifactResponse::with_body(auth, file, None, total_length),
RangeRequest::One(range) => {
ArtifactResponse::with_body(auth, file, Some(range), total_length)
}
RangeRequest::Unsatisfiable => ArtifactResponse::range_not_satisfiable(auth, total_length),
})
}
async fn ensure_membership(app: &App, row: &ReferenceRow) -> Result<(), ApiError> {
let advertised = match row.reference.ecosystem {
Ecosystem::Npm => npm::ensure_fresh_project(app, &row.reference.name).await?,
Ecosystem::PyPi => pypi::ensure_fresh_project(app, &row.reference.name).await?,
};
if advertised.contains(&row.id) {
return Ok(());
}
tracing::info!(
reference = %row.id.to_hex(),
ecosystem = row.reference.ecosystem.as_tag(),
package = %row.reference.name,
version = %row.reference.version,
"refusing a reference the current upstream snapshot no longer advertises"
);
Err(ApiError::NotFound)
}
async fn evaluate(
app: &App,
snapshot: &BlocklistSnapshot,
now: i64,
row: &ReferenceRow,
pinned: &[Digest],
) -> Decision {
osv::evaluate(
&app.osv,
Some(snapshot),
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
}
async fn load_reference(app: &App, id: ReferenceId) -> Result<ReferenceRow, ApiError> {
let caches = app.store().caches();
if let Some(row) = caches.references.get(&id) {
return Ok((*row).clone());
}
let row = app
.store()
.get_reference(id)
.await
.map_err(|err| storage_error(&err))?
.ok_or(ApiError::NotFound)?;
caches
.references
.insert(id, Arc::new(row.clone()), approximate_bytes(&row));
Ok(row)
}
fn approximate_bytes(row: &ReferenceRow) -> u64 {
(row.reference.name.len()
+ row.reference.version.len()
+ row.reference.filename.len()
+ row.reference.upstream_url.as_str().len()) as u64
+ 512
}
async fn open_or_download(
app: &Arc<App>,
row: &ReferenceRow,
pinned: Vec<Digest>,
) -> Result<(PinnedFile, Vec<Digest>, bool), ApiError> {
if let (Some(key), Some(size)) = (row.content_key, row.pinned_size) {
match app.content.open_verified(&key, size).await {
Ok(file) => return Ok((file, pinned, false)),
Err(err) => {
tracing::warn!(
reference = %row.id.to_hex(),
content = %key,
error = %err,
"the cached file is unusable; dropping the mapping and refetching"
);
if let Err(err) = app.store().clear_content_key(key).await {
tracing::warn!(error = %err, "the stale content mapping could not be cleared");
}
app.store().caches().references.remove(&row.id);
}
}
}
let verified = download::fetch(app, row)
.await
.map_err(|err| download_error(&err, app.clock.now_utc_micros()))?;
app.store().caches().references.remove(&row.id);
let file = app
.content
.open_verified(&verified.key, verified.size)
.await
.map_err(|err| {
tracing::error!(
reference = %row.id.to_hex(),
error = %err,
"the artifact was published and then could not be opened"
);
ApiError::CapacityExhausted
})?;
let mut pinned = pinned;
if pinned.is_empty() {
pinned = vec![
Digest {
algorithm: crate::policy::HashAlgorithm::Sha256,
bytes: Box::from(verified.sha256.as_slice()),
},
Digest {
algorithm: crate::policy::HashAlgorithm::Sha512,
bytes: Box::from(verified.sha512.as_slice()),
},
];
}
Ok((file, pinned, true))
}
fn decision_error(decision: Decision, now: i64) -> ApiError {
match decision {
Decision::Allow => ApiError::Blocked {
reason: "the artifact is not authorized",
},
Decision::Hold { eligible_at_micros } => ApiError::held(eligible_at_micros, now),
Decision::Deny(reason) => ApiError::Blocked {
reason: deny_reason(reason),
},
Decision::Unavailable => ApiError::PolicyUnavailable,
}
}
const fn deny_reason(reason: DenyReason) -> &'static str {
match reason {
DenyReason::BlockedPackage => "the package is blocked",
DenyReason::BlockedVersion => "the version is blocked",
DenyReason::BlockedDigest => "the artifact digest is blocked",
DenyReason::MalformedTimestamp => "the upstream publication time is malformed",
DenyReason::FutureTimestamp => "the upstream publication time is in the future",
DenyReason::NoTimestamp => "no publication time has been established",
DenyReason::BlockedByOsv => {
"the artifact is blocked by a known OSV malicious-package advisory"
}
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn deny_reason_maps_blocked_by_osv_in_artifacts_mod() {
assert_eq!(
deny_reason(DenyReason::BlockedByOsv),
"the artifact is blocked by a known OSV malicious-package advisory"
);
}
}
fn download_error(err: &DownloadError, now: i64) -> ApiError {
match err {
DownloadError::Refused(decision) => {
tracing::info!(error = %err, "policy refused verified artifact bytes");
return decision_error(*decision, now);
}
DownloadError::Capacity(_) => {
tracing::error!(error = %err, "local storage could not take the artifact");
return ApiError::CapacityExhausted;
}
_ => tracing::warn!(error = %err, "an artifact download failed"),
}
match err {
DownloadError::Upstream(crate::upstream::UpstreamError::Timeout) => {
ApiError::UpstreamTimeout
}
DownloadError::Upstream(_) => ApiError::UpstreamFailure,
DownloadError::IntegrityMismatch { .. }
| DownloadError::PinConflict
| DownloadError::SizeMismatch { .. } => ApiError::IntegrityMismatch,
DownloadError::TooLarge { .. } => ApiError::UpstreamInvalid,
DownloadError::Storage(err) => storage_error(err),
DownloadError::Overloaded => ApiError::Overloaded,
DownloadError::Cancelled => ApiError::Overloaded,
DownloadError::Aborted => ApiError::InternalFailure,
DownloadError::Capacity(_) => ApiError::CapacityExhausted,
DownloadError::Refused(_) => ApiError::Blocked {
reason: "the artifact is blocked",
},
}
}
fn storage_error(err: &StoreError) -> ApiError {
tracing::warn!(error = %err, "a storage command failed; refusing the request");
match err {
StoreError::Busy => ApiError::Overloaded,
StoreError::Closed | StoreError::Database(_) | StoreError::Corrupt(_) => {
ApiError::StorageUnusable
}
}
}