mod common;
use std::sync::Arc;
use std::time::Duration;
use common::{
FakeAnswer, FakeRegistry, Gate, TestClock, TestServer, artifact_path,
config_with_open_blocklist, npm_artifact_upstream_path, npm_tarball_url, npm_upstream_path,
publish_blocklist, snapshot_with,
};
use probation::store::rows::ReferenceId;
use serde_json::{Map, Value, json};
use tempfile::TempDir;
use tokio::io::{AsyncReadExt, AsyncWriteExt};
const WIDGET: &str = "fixture-widget";
const FILENAME: &str = "fixture-widget-1.0.0.tgz";
const VERSION: &str = "1.0.0";
const PUBLISHED: &str = "2026-04-01T00:00:00Z";
const NOW: &str = "2026-04-06T12:00:00Z";
const BODY_SHA256: &str = "029830248baf17af5d9a9e23d3e7054a8860882d1cdc06bbbb1549056d347acb";
const BODY_SHA512: &str = "72e64e9e9403218ba8896d15a6576c64b994690fde6701a3547793650a1745d4de04c8799a3d57c8732512a4f33e36055b7031cee2998bc71be0667640fe06ee";
const BODY_SHA1: &str = "06e87889ae32c865ca8a744b780e57253f0b0a47";
const BODY_SRI: &str = "sha512-cuZOnpQDIYuoiW0VpldsZLmUaQ/eZwGjVHeTZQoXRdTeBMh5mj1XyHMlEqTzPjYFW3AxzuKZi8cb4GZ2QP4G7g==";
const TAMPERED_SHA256: &str = "f753ce595d22fe1a96c17aa0d93e1fc1bfcb1d7bfa756b23d96f67ff31ffdc3b";
const TAMPERED_SRI: &str = "sha512-hl/BywthAR9bBzuAS/JOKoCECO9fgaHl2bce2C43qs4r89tbyU9hMv02qnBf4Vn5193+NTmwCSt1RBlrihU54Q==";
fn body() -> String {
common::fixture("artifacts/harmless-widget-1.0.0.tgz")
}
fn tampered() -> String {
common::fixture("artifacts/tampered-widget-1.0.0.tgz")
}
fn document(dist: Value) -> String {
let mut versions = Map::new();
versions.insert(
VERSION.to_owned(),
json!({"name": WIDGET, "version": VERSION, "dist": dist}),
);
let mut time = Map::new();
time.insert(VERSION.to_owned(), json!(PUBLISHED));
json!({
"name": WIDGET,
"dist-tags": {"latest": VERSION},
"versions": Value::Object(versions),
"time": Value::Object(time),
})
.to_string()
}
fn strong_integrity() -> Value {
json!({"tarball": npm_tarball_url(WIDGET, FILENAME), "integrity": BODY_SRI})
}
fn legacy_shasum_only() -> Value {
json!({"tarball": npm_tarball_url(WIDGET, FILENAME), "shasum": BODY_SHA1})
}
fn no_advertised_digest() -> Value {
json!({"tarball": npm_tarball_url(WIDGET, FILENAME)})
}
struct Harness {
server: TestServer,
registry: Arc<FakeRegistry>,
clock: Arc<TestClock>,
data_dir: std::path::PathBuf,
_dir: TempDir,
}
impl Harness {
async fn artifact_path(&self) -> String {
let document = self.server.json(&format!("/npm/{WIDGET}")).await;
artifact_path(&document, VERSION)
}
fn reference_id(&self, path: &str) -> ReferenceId {
let hex = path
.split('/')
.nth(3)
.unwrap_or_else(|| panic!("an artifact path has a reference id: {path}"));
ReferenceId::parse_hex(hex).expect("the reference id is hexadecimal")
}
async fn reference(&self, path: &str) -> probation::store::rows::ReferenceRow {
self.server
.running()
.app()
.store()
.get_reference(self.reference_id(path))
.await
.expect("the reference query")
.expect("the reference was committed before its URL was advertised")
}
fn content_files(&self) -> Vec<std::path::PathBuf> {
let mut found = Vec::new();
let objects = self.data_dir.join("content").join("objects");
let Ok(fan_out) = std::fs::read_dir(&objects) else {
return found;
};
for directory in fan_out.flatten() {
if let Ok(entries) = std::fs::read_dir(directory.path()) {
found.extend(entries.flatten().map(|entry| entry.path()));
}
}
found
}
fn temp_files(&self) -> usize {
std::fs::read_dir(self.data_dir.join("content").join("tmp"))
.map(|entries| entries.count())
.unwrap_or(0)
}
fn artifact_calls(&self) -> usize {
self.registry
.calls()
.iter()
.filter(|url| url.path() == npm_artifact_upstream_path(WIDGET, FILENAME))
.count()
}
async fn shutdown(self) {
self.server.shutdown().await;
}
}
async fn harness(dist: Value, artifact: FakeAnswer) -> Harness {
let dir = tempfile::tempdir().expect("a temporary directory");
let config = config_with_open_blocklist(dir.path());
let data_dir = dir.path().join("data");
let registry = FakeRegistry::new();
registry.answer(&npm_upstream_path(WIDGET), FakeAnswer::Body(document(dist)));
registry.answer(&npm_artifact_upstream_path(WIDGET, FILENAME), artifact);
let clock = TestClock::at_rfc3339(NOW);
let server = TestServer::start_in_with_registry(
&data_dir,
config,
clock.shared(),
Arc::clone(®istry),
)
.await;
Harness {
server,
registry,
clock,
data_dir,
_dir: dir,
}
}
fn blocklist(revision: u64, packages: &str, hashes: &str) -> String {
snapshot_with(
revision,
"2026-04-05T00:00:00Z",
"2099-01-01T00:00:00Z",
packages,
hashes,
)
}
fn blocked_digest(algorithm: &str, digest: &str) -> String {
format!(
r#"{{"algorithm":"{algorithm}","digest":"{digest}","reason":"known malicious artifact"}}"#
)
}
#[tokio::test]
async fn zero_body_bytes_before_cold_verification_completes() {
let gate = Gate::new();
let whole = body();
let (head, tail) = whole.split_at(whole.len() / 2);
let harness = harness(
strong_integrity(),
FakeAnswer::Gated {
head: head.to_owned(),
tail: tail.to_owned(),
gate: Arc::clone(&gate),
},
)
.await;
let path = harness.artifact_path().await;
let addr = harness.server.running().local_addr;
let mut socket = tokio::net::TcpStream::connect(addr)
.await
.expect("the server accepts a connection");
socket
.write_all(
format!("GET {path} HTTP/1.1\r\nHost: {addr}\r\nConnection: close\r\n\r\n").as_bytes(),
)
.await
.expect("the request is written");
gate.wait_until_reached().await;
let mut early = [0u8; 1];
let peeked = tokio::time::timeout(Duration::from_millis(300), socket.read(&mut early)).await;
assert!(
peeked.is_err(),
"half the artifact has arrived from upstream and the client already received \
{peeked:?}; nothing may cross the socket before verification completes"
);
assert!(
harness.content_files().is_empty(),
"and nothing is in the content cache yet either"
);
gate.release();
let mut response = Vec::new();
tokio::time::timeout(Duration::from_secs(5), socket.read_to_end(&mut response))
.await
.expect("the response arrives once verification completes")
.expect("the response is read");
let response = String::from_utf8_lossy(&response).into_owned();
assert!(
response.starts_with("HTTP/1.1 200 OK"),
"the verified artifact is served once, whole: {response}"
);
assert!(
response.ends_with(&whole),
"and the body is the artifact itself"
);
harness.shutdown().await;
}
#[tokio::test]
async fn truncated_download_never_becomes_a_cache_hit() {
let whole = body();
let harness = harness(
no_advertised_digest(),
FakeAnswer::Truncated {
declared_length: whole.len() as u64,
body: whole[..whole.len() / 2].to_owned(),
},
)
.await;
let path = harness.artifact_path().await;
assert_eq!(harness.server.status(&path).await, 502);
let row = harness.reference(&path).await;
assert_eq!(row.content_key, None, "no mapping was committed");
assert_eq!(
row.pinned_sha256, None,
"SPEC §9: an incomplete download never establishes a pin"
);
assert!(harness.content_files().is_empty());
assert_eq!(harness.temp_files(), 0, "and its temporary file is gone");
assert_eq!(
harness.server.status(&path).await,
502,
"a second request is refused too, rather than finding the first attempt cached"
);
assert_eq!(
harness.artifact_calls(),
2,
"both attempts really went upstream; neither was served from a cache"
);
harness.shutdown().await;
}
#[tokio::test]
async fn integrity_failure_never_becomes_a_cache_hit() {
let harness = harness(
json!({"tarball": npm_tarball_url(WIDGET, FILENAME), "integrity": TAMPERED_SRI}),
FakeAnswer::Body(body()),
)
.await;
let path = harness.artifact_path().await;
assert_eq!(harness.server.status(&path).await, 502);
let error = common::body_error(harness.server.get(&path).await).await;
assert_eq!(error, "INTEGRITY_MISMATCH");
let row = harness.reference(&path).await;
assert_eq!(row.content_key, None);
assert_eq!(
row.pinned_sha256, None,
"SPEC §9: an upstream-integrity-failing download never establishes a pin"
);
assert!(harness.content_files().is_empty());
assert_eq!(harness.temp_files(), 0);
harness.shutdown().await;
}
#[tokio::test]
async fn blocked_bytes_never_become_a_cache_hit() {
let harness = harness(strong_integrity(), FakeAnswer::Body(body())).await;
let path = harness.artifact_path().await;
publish_blocklist(
&harness.server,
&blocklist(2, "", &blocked_digest("sha256", BODY_SHA256)),
common::parse_rfc3339(NOW),
);
assert_eq!(harness.server.status(&path).await, 403);
let row = harness.reference(&path).await;
assert_eq!(
row.pinned_sha256.map(hex::encode).as_deref(),
Some(BODY_SHA256),
"SPEC §9: persist computed digests, even when policy now blocks them"
);
assert_eq!(row.content_key, None, "but the bytes are not published");
assert!(harness.content_files().is_empty());
assert_eq!(harness.temp_files(), 0);
assert_eq!(
harness.server.status(&path).await,
403,
"and the next request is refused from the pins, not served from a cache"
);
harness.shutdown().await;
}
#[tokio::test]
async fn eviction_and_refetch_preserve_pins() {
let harness = harness(strong_integrity(), FakeAnswer::Body(body())).await;
let path = harness.artifact_path().await;
assert_eq!(harness.server.status(&path).await, 200);
let pinned = harness.reference(&path).await;
assert_eq!(
pinned.pinned_sha256.map(hex::encode).as_deref(),
Some(BODY_SHA256)
);
assert_eq!(
pinned.pinned_sha512.map(hex::encode).as_deref(),
Some(BODY_SHA512)
);
assert!(pinned.content_key.is_some());
harness.clock.advance_seconds(600);
let _ = harness.server.json(&format!("/npm/{WIDGET}")).await;
let refreshed = harness.reference(&path).await;
assert_eq!(
refreshed.pinned_sha256, pinned.pinned_sha256,
"an upstream metadata refresh must not erase what this reference was proven \
to mean"
);
assert_eq!(refreshed.content_key, pinned.content_key);
for file in harness.content_files() {
std::fs::remove_file(&file).expect("the evicted file is removed");
}
let response = harness.server.get(&path).await;
assert_eq!(response.status().as_u16(), 200);
assert_eq!(
response.text().await.expect("a body"),
body(),
"the refetched bytes are served"
);
assert_eq!(
harness.artifact_calls(),
2,
"the eviction really did force a second transfer"
);
let refetched = harness.reference(&path).await;
assert_eq!(
refetched.pinned_sha256, pinned.pinned_sha256,
"the pin survived the eviction of the bytes it describes"
);
assert_eq!(refetched.pinned_sha512, pinned.pinned_sha512);
assert_eq!(refetched.pinned_size, pinned.pinned_size);
harness.shutdown().await;
}
#[tokio::test]
async fn changed_bytes_for_the_same_reference_are_refused_with_502() {
let harness = harness(no_advertised_digest(), FakeAnswer::Body(body())).await;
let path = harness.artifact_path().await;
assert_eq!(harness.server.status(&path).await, 200);
let pinned = harness.reference(&path).await;
assert_eq!(
pinned.pinned_sha256.map(hex::encode).as_deref(),
Some(BODY_SHA256)
);
assert!(
pinned.reference.expected.is_empty(),
"upstream advertised nothing, so only the pins can refuse the new bytes"
);
for file in harness.content_files() {
std::fs::remove_file(&file).expect("the evicted file is removed");
}
harness.registry.answer(
&npm_artifact_upstream_path(WIDGET, FILENAME),
FakeAnswer::Body(tampered()),
);
let response = harness.server.get(&path).await;
assert_eq!(response.status().as_u16(), 502);
assert_eq!(
common::body_error(response).await,
"INTEGRITY_MISMATCH",
"SPEC §9: a mismatch is 502 INTEGRITY_MISMATCH"
);
let after = harness.reference(&path).await;
assert_eq!(
after.pinned_sha256, pinned.pinned_sha256,
"the original pins are retained, never replaced by the new bytes"
);
assert_ne!(
after.pinned_sha256.map(hex::encode).as_deref(),
Some(TAMPERED_SHA256)
);
assert_eq!(after.content_key, None);
assert!(
harness.content_files().is_empty(),
"and the refused bytes were discarded"
);
assert_eq!(harness.temp_files(), 0);
harness.shutdown().await;
}
#[tokio::test]
async fn pin_established_without_a_strong_advertised_digest() {
let harness = harness(legacy_shasum_only(), FakeAnswer::Body(body())).await;
let path = harness.artifact_path().await;
assert_eq!(harness.server.status(&path).await, 200);
let row = harness.reference(&path).await;
assert_eq!(
row.reference.expected.len(),
1,
"upstream advertised exactly one digest, and it is the legacy SHA-1"
);
assert_eq!(row.reference.expected[0].to_hex_lowercase(), BODY_SHA1);
assert_eq!(
row.pinned_sha256.map(hex::encode).as_deref(),
Some(BODY_SHA256),
"a SHA-256 was computed and pinned although upstream never advertised one"
);
assert_eq!(
row.pinned_sha512.map(hex::encode).as_deref(),
Some(BODY_SHA512)
);
assert_eq!(row.pinned_size, Some(body().len() as u64));
harness.shutdown().await;
}
#[tokio::test]
async fn local_block_denies_with_upstream_unreachable() {
let harness = harness(strong_integrity(), FakeAnswer::Body(body())).await;
let path = harness.artifact_path().await;
publish_blocklist(
&harness.server,
&blocklist(
2,
&format!(
r#"{{"ecosystem":"npm","name":"{WIDGET}","version":null,"reason":"malware"}}"#
),
"",
),
common::parse_rfc3339(NOW),
);
harness.registry.go_offline();
let calls_before = harness.registry.calls().len();
let response = harness.server.get(&path).await;
assert_eq!(
response.status().as_u16(),
403,
"a locally conclusive block answers without upstream"
);
assert_eq!(common::body_error(response).await, "BLOCKED");
assert_eq!(
harness.registry.calls().len(),
calls_before,
"and it reached upstream not once — an unreachable registry cannot even be \
observed from this path"
);
harness.shutdown().await;
}
#[tokio::test]
async fn a_dot_dot_filename_never_serves_another_reference() {
let harness = harness(strong_integrity(), FakeAnswer::Body(body())).await;
let path = harness.artifact_path().await;
assert_eq!(harness.server.status(&path).await, 200);
let id = path.split('/').nth(3).expect("a reference id").to_owned();
for hostile in [
format!("/npm/artifacts/{id}/%2E%2E"),
format!("/npm/artifacts/{id}/..%2F{FILENAME}"),
format!("/npm/artifacts/{id}/other-name.tgz"),
] {
let status = harness.server.raw_get_status(&hostile).await;
assert_eq!(
status, 404,
"{hostile} names no reference this instance has committed"
);
}
harness.shutdown().await;
}
#[tokio::test]
async fn a_warm_artifact_request_does_not_reparse_the_project_document() {
let harness = harness(strong_integrity(), FakeAnswer::Body(body())).await;
let path = harness.artifact_path().await;
assert_eq!(harness.server.status(&path).await, 200);
let before = harness.server.app().store().caches().stored_parses();
assert_eq!(harness.server.status(&path).await, 200);
let warm = harness.server.app().store().caches().stored_parses() - before;
assert_eq!(
warm, 0,
"a warm artifact request parsed the stored project document {warm} time(s); \
the advertised reference-id set is supposed to be kept with the snapshot it \
was parsed from"
);
harness.shutdown().await;
}