pub mod document;
pub mod name;
pub mod render;
pub mod tags;
use std::collections::{BTreeSet, HashMap, HashSet};
use std::sync::{Arc, OnceLock};
use std::time::{Duration, Instant};
use serde_json::Value;
use crate::App;
use crate::concurrency::Resolution;
use crate::http::error::ApiError;
use crate::http::logging;
use crate::npm::document::{PackageDocument, VersionEntry};
use crate::policy::{
BlocklistSnapshot, Candidate, Decision, DenyReason, Ecosystem, PublicationTime,
};
use crate::store::cache::{
AbsentMark, CachedProject, ProjectKey, RenderKey, RenderedResponse, Representation,
};
use crate::store::rows::{Generation, ProjectRefresh, ProjectRow, ReferenceId, ReferenceUpsert};
use crate::store::{self, StoreError};
use crate::upstream::{MetadataRequest, MetadataResponse, OriginKind, UpstreamError};
pub use name::{InvalidPackageName, PackageName};
const NPM_ACCEPT: &str = "application/json";
pub struct Rendered {
pub body: Arc<[u8]>,
pub content_type: &'static str,
}
pub async fn serve(
app: &App,
name: &PackageName,
representation: Representation,
) -> Result<Rendered, ApiError> {
let now = app.clock.now_utc_micros();
let monotonic = app.clock.now_monotonic();
let snapshot = app.blocklist().ok_or(ApiError::PolicyUnavailable)?;
if !snapshot.is_valid_at(now) {
return Err(ApiError::PolicyUnavailable);
}
let project = ProjectKey::new(Ecosystem::Npm, name.as_str());
let render_key = RenderKey {
project: project.clone(),
representation: representation.clone(),
};
let caches = app.store().caches();
if let Some(rendered) = caches.rendered.get(&render_key)
&& let Some(cached) = caches.projects.get(&project)
&& rendered.is_reusable(
cached.row.generation,
snapshot.revision,
cached.row.digest_generation,
now,
monotonic,
)
{
logging::record_cache(logging::CacheStatus::Hit);
return Ok(Rendered {
body: Arc::clone(&rendered.body),
content_type: rendered.content_type,
});
}
logging::record_cache(logging::CacheStatus::Miss);
let resolved = resolve(app, name, &snapshot, now, monotonic).await?;
let (body, content_type) = render(app, &resolved, &representation, now)?;
let body: Arc<[u8]> = Arc::from(body);
if resolved.deadline_utc_micros > now {
caches.rendered.insert(
render_key,
Arc::new(RenderedResponse {
body: Arc::clone(&body),
content_type,
project_generation: resolved.generation,
blocklist_revision: snapshot.revision,
digest_generation: resolved.digest_generation,
deadline_utc_micros: resolved.deadline_utc_micros,
deadline_monotonic: resolved.deadline_monotonic,
}),
body.len() as u64 + 256,
);
}
Ok(Rendered { body, content_type })
}
struct Resolved {
document: PackageDocument,
entries: Vec<VersionEntry>,
decisions: Vec<Decision>,
eligible: BTreeSet<String>,
tags: serde_json::Map<String, Value>,
generation: Generation,
digest_generation: u64,
deadline_utc_micros: i64,
deadline_monotonic: Instant,
}
async fn current_project(
app: &App,
name: &PackageName,
now: i64,
monotonic: Instant,
) -> Result<Arc<CachedProject>, ApiError> {
let key = ProjectKey::new(Ecosystem::Npm, name.as_str());
let caches = app.store().caches();
let ttl = Duration::from_secs(app.config.metadata_ttl_seconds);
let ttl_micros = i64::try_from(ttl.as_micros()).unwrap_or(i64::MAX);
if let Some(mark) = caches.absent.get(&key) {
if monotonic.duration_since(mark.observed_monotonic) < ttl {
return Err(ApiError::NotFound);
}
caches.absent.remove(&key);
}
let (seen, cached) = caches.projects.get_with_seen(&key);
let cached = match cached {
Some(cached) => Some(cached),
None => load_project(app, name).await?.map(Arc::new),
};
let over_age = cached.as_ref().is_some_and(|cached| {
store::is_over_age(
&key,
app.config.metadata_max_age_seconds,
cached.row.fetched_at_micros,
now,
)
});
let cached = match cached {
Some(cached)
if !over_age && now < cached.row.validated_at_micros.saturating_add(ttl_micros) =>
{
cached
}
stale => refresh_coalesced(app, name, &key, now, monotonic, over_age, stale).await?,
};
caches
.projects
.insert_if_current(key, Arc::clone(&cached), cached.approximate_bytes(), seen);
Ok(cached)
}
pub async fn ensure_fresh_project(
app: &App,
name: &str,
) -> Result<Arc<HashSet<ReferenceId>>, ApiError> {
let name = PackageName::parse_route(name).map_err(|err| {
tracing::warn!(error = %err, "a stored reference names an unusable package");
ApiError::NotFound
})?;
let cached = current_project(
app,
&name,
app.clock.now_utc_micros(),
app.clock.now_monotonic(),
)
.await?;
if let Some(advertised) = cached.advertised.get() {
return Ok(Arc::clone(advertised));
}
let document = parse_stored(app, &name, &cached.row.payload)?;
let entries = document
.entries(app.config.max_references_per_project.get())
.map_err(|err| reference_cap_error(&name, err))?;
let advertised: Arc<HashSet<ReferenceId>> =
Arc::new(entries.iter().map(|entry| entry.id).collect());
let _ = cached.advertised.set(Arc::clone(&advertised));
Ok(advertised)
}
fn parse_stored(
app: &App,
name: &PackageName,
payload: &[u8],
) -> Result<PackageDocument, ApiError> {
app.store().caches().note_stored_parse();
PackageDocument::parse(payload).map_err(|err| {
tracing::error!(package = %name, error = %err, "the stored project payload is unusable");
ApiError::UpstreamInvalid
})
}
async fn resolve(
app: &App,
name: &PackageName,
snapshot: &BlocklistSnapshot,
now: i64,
monotonic: Instant,
) -> Result<Resolved, ApiError> {
let cached = current_project(app, name, now, monotonic).await?;
let ttl_micros =
i64::try_from(Duration::from_secs(app.config.metadata_ttl_seconds).as_micros())
.unwrap_or(i64::MAX);
let document = parse_stored(app, name, &cached.row.payload)?;
let entries = document
.entries(app.config.max_references_per_project.get())
.map_err(|err| reference_cap_error(name, err))?;
let mut decisions = Vec::with_capacity(entries.len());
let mut eligible = BTreeSet::new();
let mut next_release: Option<i64> = None;
for entry in &entries {
let candidate = Candidate {
ecosystem: Ecosystem::Npm,
name: name.as_str(),
version: entry.version(),
publication: publication_of(entry, &cached.first_seen),
advertised_digests: &entry.reference.expected,
pinned_digests: cached.pins_of(&entry.id),
};
let decision = crate::osv::evaluate(
&app.osv,
Some(snapshot),
now,
app.config.cooldown_seconds,
&candidate,
)
.await;
match decision {
Decision::Allow => {
eligible.insert(entry.version().to_owned());
}
Decision::Hold { eligible_at_micros } => {
next_release = Some(
next_release
.map_or(eligible_at_micros, |held: i64| held.min(eligible_at_micros)),
);
}
Decision::Deny(_) | Decision::Unavailable => {}
}
decisions.push(decision);
}
let tags = tags::resolve(&document.dist_tags, &eligible);
let ceiling = store::effective_max_age_micros(
&ProjectKey::new(Ecosystem::Npm, name.as_str()),
app.config.metadata_max_age_seconds,
)
.map_or(i64::MAX, |ceiling| {
cached.row.fetched_at_micros.saturating_add(ceiling)
});
let deadline_utc_micros = cached
.row
.validated_at_micros
.saturating_add(ttl_micros)
.min(snapshot.expires_at_micros)
.min(next_release.unwrap_or(i64::MAX))
.min(ceiling);
let deadline_monotonic =
monotonic + Duration::from_micros(deadline_utc_micros.saturating_sub(now).max(0) as u64);
Ok(Resolved {
generation: cached.row.generation,
digest_generation: cached.row.digest_generation,
document,
entries,
decisions,
eligible,
tags,
deadline_utc_micros,
deadline_monotonic,
})
}
fn publication_of(entry: &VersionEntry, first_seen: &HashMap<ReferenceId, i64>) -> PublicationTime {
match entry.publication {
PublicationTime::Upstream(micros) => PublicationTime::Upstream(micros),
PublicationTime::Malformed => PublicationTime::Malformed,
PublicationTime::FirstSeen(_) | PublicationTime::Unknown => first_seen
.get(&entry.id)
.map_or(PublicationTime::Unknown, |micros| {
PublicationTime::FirstSeen(*micros)
}),
}
}
async fn load_project(app: &App, name: &PackageName) -> Result<Option<CachedProject>, ApiError> {
let row = app
.store()
.get_project(Ecosystem::Npm, name.as_str())
.await
.map_err(storage_error)?;
let Some(row) = row else {
return Ok(None);
};
let first_seen = app
.store()
.list_project_first_seen(Ecosystem::Npm, name.as_str())
.await
.map_err(storage_error)?;
let pins = app
.store()
.list_project_pins(Ecosystem::Npm, name.as_str())
.await
.map_err(storage_error)?;
Ok(Some(CachedProject {
row,
first_seen,
pins,
advertised: OnceLock::new(),
}))
}
async fn refresh_coalesced(
app: &App,
name: &PackageName,
key: &ProjectKey,
now: i64,
monotonic: Instant,
over_age: bool,
stale: Option<Arc<CachedProject>>,
) -> Result<Arc<CachedProject>, ApiError> {
let refreshes = app.downloads.metadata();
let (slot, leader, _waiter) = refreshes.join(key.clone());
if !leader {
return slot.wait().await.unwrap_or(Err(ApiError::InternalFailure));
}
let mut resolution = Resolution::new(
refreshes.clone(),
key.clone(),
Arc::clone(&slot),
Err(ApiError::InternalFailure),
name.as_str().to_owned(),
);
let refreshed = refresh(app, name, key, now, monotonic, over_age, stale).await;
resolution.answer(refreshed.clone());
refreshed
}
async fn refresh(
app: &App,
name: &PackageName,
key: &ProjectKey,
now: i64,
monotonic: Instant,
over_age: bool,
stale: Option<Arc<CachedProject>>,
) -> Result<Arc<CachedProject>, ApiError> {
let url = app
.origins
.url_for(OriginKind::NpmMetadata, &[name.upstream_segment()])
.map_err(|rejection| {
tracing::warn!(package = %name, %rejection, "refused to build an upstream URL");
ApiError::NotFound
})?;
let fetched = app
.transport
.fetch_metadata(MetadataRequest {
url,
accept: NPM_ACCEPT,
validators: (!over_age)
.then(|| {
stale
.as_ref()
.map(|cached| cached.row.validators.clone())
.filter(|validators| !validators.is_empty())
})
.flatten(),
max_bytes: app.config.max_metadata_bytes.get(),
})
.await;
let (payload, validators, fetched_at_micros): (Arc<[u8]>, _, i64) = match fetched {
Ok(MetadataResponse::Fresh { body, validators }) => {
(Arc::from(body.as_ref()), validators, now)
}
Ok(MetadataResponse::NotModified { validators }) => match &stale {
Some(cached) if !over_age => {
let validators = if validators.is_empty() {
cached.row.validators.clone()
} else {
validators
};
(
Arc::clone(&cached.row.payload),
validators,
cached.row.fetched_at_micros,
)
}
_ => {
tracing::warn!(package = %name, "upstream answered 304 to an unconditional request");
return Err(ApiError::UpstreamInvalid);
}
},
Ok(MetadataResponse::Missing) => {
app.store().caches().absent.insert(
key.clone(),
Arc::new(AbsentMark {
observed_monotonic: monotonic,
}),
128,
);
return Err(ApiError::NotFound);
}
Err(UpstreamError::Timeout) => {
tracing::warn!(package = %name, "the upstream registry did not answer in time");
return Err(ApiError::UpstreamTimeout);
}
Err(err) => {
tracing::warn!(package = %name, error = %err, "the upstream fetch failed");
return Err(ApiError::UpstreamFailure);
}
};
let document = PackageDocument::parse(&payload).map_err(|err| {
tracing::warn!(package = %name, error = %err, "the upstream document is unusable");
ApiError::UpstreamInvalid
})?;
let entries = document
.entries(app.config.max_references_per_project.get())
.map_err(|err| reference_cap_error(name, err))?;
let references = entries
.iter()
.map(|entry| ReferenceUpsert {
id: entry.id,
reference: entry.reference.clone(),
publication_micros: match entry.publication {
PublicationTime::Upstream(micros) => Some(micros),
_ => None,
},
first_seen_micros: matches!(entry.publication, PublicationTime::Unknown).then_some(now),
})
.collect();
let pins = app
.store()
.list_project_pins(Ecosystem::Npm, name.as_str())
.await
.map_err(storage_error)?;
let committed = app
.store()
.commit_project_refresh(ProjectRefresh {
ecosystem: Ecosystem::Npm,
name: name.as_str().to_owned(),
payload: Arc::clone(&payload),
validators: validators.clone(),
validated_at_micros: now,
fetched_at_micros,
references,
})
.await
.map_err(storage_error)?;
Ok(Arc::new(CachedProject {
row: ProjectRow {
ecosystem: Ecosystem::Npm,
name: name.as_str().to_owned(),
payload,
validators,
validated_at_micros: now,
fetched_at_micros,
generation: committed.generation,
digest_generation: committed.digest_generation,
},
first_seen: committed.first_seen,
pins,
advertised: OnceLock::new(),
}))
}
fn render(
app: &App,
resolved: &Resolved,
representation: &Representation,
now: i64,
) -> Result<(Vec<u8>, &'static str), ApiError> {
let public_url = &app.config.public_url;
if let Representation::NpmVersion(spelling) = representation {
return single_version(resolved, spelling, public_url, now);
}
if resolved.eligible.is_empty() {
return Err(denial(resolved, now));
}
let kept: Vec<&VersionEntry> = resolved
.entries
.iter()
.filter(|entry| resolved.eligible.contains(entry.version()))
.collect();
Ok(match representation {
Representation::NpmAbbreviated => (
render::abbreviated(&resolved.document, &kept, &resolved.tags, public_url),
render::ABBREVIATED_CONTENT_TYPE,
),
_ => (
render::full(&resolved.document, &kept, &resolved.tags, public_url),
render::FULL_CONTENT_TYPE,
),
})
}
fn single_version(
resolved: &Resolved,
spelling: &str,
public_url: &url::Url,
now: i64,
) -> Result<(Vec<u8>, &'static str), ApiError> {
let target = if resolved.document.dist_tags.contains_key(spelling) {
match resolved.tags.get(spelling).and_then(Value::as_str) {
Some(target) => target.to_owned(),
None => return Err(ApiError::NotFound),
}
} else {
spelling.to_owned()
};
let Some(index) = resolved
.entries
.iter()
.position(|entry| entry.version() == target)
else {
return Err(ApiError::NotFound);
};
match resolved.decisions[index] {
Decision::Allow => Ok((
render::single_version(&resolved.entries[index], public_url),
render::FULL_CONTENT_TYPE,
)),
Decision::Hold { eligible_at_micros } => Err(ApiError::held(eligible_at_micros, now)),
Decision::Deny(reason) => Err(ApiError::Blocked {
reason: deny_reason(reason),
}),
Decision::Unavailable => Err(ApiError::PolicyUnavailable),
}
}
fn denial(resolved: &Resolved, now: i64) -> ApiError {
for decision in &resolved.decisions {
if let Decision::Deny(reason) = decision {
return ApiError::Blocked {
reason: deny_reason(*reason),
};
}
}
let earliest = resolved
.decisions
.iter()
.filter_map(|decision| match decision {
Decision::Hold { eligible_at_micros } => Some(*eligible_at_micros),
_ => None,
})
.min();
match earliest {
Some(eligible_at_micros) => ApiError::held(eligible_at_micros, now),
None => ApiError::NotFound,
}
}
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 version is blocked by a known OSV malicious-package advisory"
}
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn deny_reason_maps_blocked_by_osv_in_npm_mod() {
assert_eq!(
deny_reason(DenyReason::BlockedByOsv),
"the version is blocked by a known OSV malicious-package advisory"
);
}
}
fn reference_cap_error(name: &PackageName, err: document::DocumentError) -> ApiError {
if let document::DocumentError::TooManyReferences { count, limit } = &err {
tracing::error!(
package = %name,
references = count,
limit = limit,
"refusing a project whose reference count exceeds max_references_per_project"
);
}
ApiError::UpstreamInvalid
}
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
}
}
}