mod common;
use std::sync::Arc;
use common::{
FakeAnswer, FakeRegistry, TestClock, TestServer, npm_upstream_path, parse_rfc3339,
sample_config,
};
use probation::policy::Ecosystem;
use probation::store::cache::MemoryCaches;
use probation::store::rows::{
ArtifactReference, Generation, ProjectRefresh, ReferenceId, ReferenceUpsert,
};
use probation::store::{StoreHandle, startup};
use probation::upstream::UpstreamValidators;
use serde_json::Value;
use tempfile::TempDir;
use url::Url;
const NAME: &str = "fixture-untimed";
const T0: &str = "2026-04-06T12:00:00Z";
const ONE_DAY: u64 = 86_400;
async fn store(dir: &TempDir) -> (StoreHandle, tokio::task::JoinHandle<()>) {
let opened = startup::open_and_recover(dir.path())
.await
.expect("the data directory opens");
let connection = opened.connection.expect("a fresh database recovers");
probation::store::spawn(
connection,
opened.lock,
Arc::new(MemoryCaches::new(1 << 20)),
)
}
fn reference(version: &str, filename: &str) -> ArtifactReference {
ArtifactReference {
ecosystem: Ecosystem::Npm,
name: "widget".to_owned(),
version: version.to_owned(),
filename: filename.to_owned(),
upstream_url: Url::parse(&format!(
"https://npm.invalid/widget/-/widget-{version}.tgz"
))
.expect("a test URL"),
expected: Vec::new(),
}
}
fn upsert(version: &str, filename: &str, first_seen: Option<i64>) -> ReferenceUpsert {
let reference = reference(version, filename);
ReferenceUpsert {
id: ReferenceId::compute(&reference),
reference,
publication_micros: None,
first_seen_micros: first_seen,
}
}
fn refresh(validated_at: i64, references: Vec<ReferenceUpsert>) -> ProjectRefresh {
ProjectRefresh {
ecosystem: Ecosystem::Npm,
name: "widget".to_owned(),
payload: Arc::from(br#"{"name":"widget"}"#.as_slice()),
validators: UpstreamValidators::default(),
validated_at_micros: validated_at,
fetched_at_micros: validated_at,
references,
}
}
#[tokio::test]
async fn project_and_references_commit_atomically() {
let dir = tempfile::tempdir().expect("a temporary data directory");
let (store, task) = store(&dir).await;
let references = vec![
upsert("1.0.0", "widget-1.0.0.tgz", Some(1_000)),
upsert("1.1.0", "widget-1.1.0.tgz", Some(2_000)),
upsert("2.0.0", "widget-2.0.0.tgz", None),
];
let committed = store
.commit_project_refresh(refresh(5_000, references.clone()))
.await
.expect("the refresh commits");
assert_eq!(committed.generation, Generation(1));
let project = store
.get_project(Ecosystem::Npm, "widget")
.await
.expect("a query")
.expect("the project landed");
assert_eq!(project.generation, Generation(1));
assert_eq!(project.validated_at_micros, 5_000);
for upsert in &references {
let row = store
.get_reference(upsert.id)
.await
.expect("a query")
.unwrap_or_else(|| {
panic!(
"{} landed in the same transaction as its project",
upsert.reference.version
)
});
assert_eq!(row.reference, upsert.reference);
assert_eq!(row.first_seen_micros, upsert.first_seen_micros);
}
assert_eq!(
committed.first_seen.len(),
2,
"the commit reports the first-seen values in force, and only the two that have one"
);
let again = store
.commit_project_refresh(refresh(6_000, references))
.await
.expect("the second refresh commits");
assert_eq!(again.generation, Generation(2));
drop(store);
task.await.expect("the storage task stops");
}
#[tokio::test]
async fn transaction_rollback_leaves_no_partial_project() {
let dir = tempfile::tempdir().expect("a temporary data directory");
let (store, task) = store(&dir).await;
let good = upsert("1.0.0", "widget-1.0.0.tgz", Some(1_000));
let bad = upsert("1.1.0", "", None);
let outcome = store
.commit_project_refresh(refresh(5_000, vec![good.clone(), bad]))
.await;
assert!(
outcome.is_err(),
"the commit must fail rather than write what it can"
);
assert!(
store
.get_project(Ecosystem::Npm, "widget")
.await
.expect("a query")
.is_none(),
"the project row was written before the failure and must not have survived it"
);
assert!(
store
.get_reference(good.id)
.await
.expect("a query")
.is_none(),
"nor the reference that was written before it"
);
store
.commit_project_refresh(refresh(6_000, vec![good.clone()]))
.await
.expect("a later refresh still commits");
assert!(
store
.get_reference(good.id)
.await
.expect("a query")
.is_some()
);
drop(store);
task.await.expect("the storage task stops");
}
fn untimed_document() -> String {
format!(
r#"{{
"name": "{NAME}",
"dist-tags": {{"latest": "1.0.0"}},
"versions": {{
"1.0.0": {{"name": "{NAME}", "version": "1.0.0",
"dist": {{"tarball": "https://npm.invalid/{NAME}/-/{NAME}-1.0.0.tgz"}}}}
}},
"time": {{}}
}}"#
)
}
fn timed_document() -> String {
untimed_document().replace(
r#""time": {}"#,
r#""time": {"1.0.0": "2026-01-01T00:00:00.000Z"}"#,
)
}
fn untimed_reference_id() -> ReferenceId {
ReferenceId::compute(&ArtifactReference {
ecosystem: Ecosystem::Npm,
name: NAME.to_owned(),
version: "1.0.0".to_owned(),
filename: format!("{NAME}-1.0.0.tgz"),
upstream_url: Url::parse(&format!("https://npm.invalid/{NAME}/-/{NAME}-1.0.0.tgz"))
.expect("a test URL"),
expected: Vec::new(),
})
}
fn registry_answering(document: String) -> Arc<FakeRegistry> {
let registry = FakeRegistry::new();
registry.answer(&npm_upstream_path(NAME), FakeAnswer::Body(document));
registry
}
fn config_in(dir: &TempDir, cooldown: u64) -> probation::config::Config {
let blocklist_file = dir.path().join("blocklist.json");
if !blocklist_file.exists() {
std::fs::write(
&blocklist_file,
common::snapshot(1, "2020-01-01T00:00:00Z", "2099-01-01T00:00:00Z", ""),
)
.expect("the blocklist is written");
}
let mut config = sample_config();
config.blocklist_file = blocklist_file;
config.cooldown_seconds = cooldown;
config
}
#[tokio::test]
async fn first_seen_persisted_before_it_grants_eligibility() {
let t0 = parse_rfc3339(T0);
let dir = tempfile::tempdir().expect("a temporary data directory");
let server = TestServer::start_in_with_registry(
dir.path(),
config_in(&dir, 0),
TestClock::at_rfc3339(T0).shared(),
registry_answering(untimed_document()),
)
.await;
let document: Value = server.json(&format!("/npm/{NAME}")).await;
assert!(
document["versions"].get("1.0.0").is_some(),
"with a zero cooldown an untimed release is eligible as soon as its first-seen \
time exists"
);
let row = server
.running()
.app()
.store()
.get_reference(untimed_reference_id())
.await
.expect("a query")
.expect("the first-seen time is already committed when the response is produced");
assert_eq!(row.first_seen_micros, Some(t0));
assert_eq!(
row.publication_micros, None,
"upstream supplied no time, so none was invented"
);
server.shutdown().await;
let dir = tempfile::tempdir().expect("a temporary data directory");
let server = TestServer::start_in_with_registry(
dir.path(),
config_in(&dir, ONE_DAY),
TestClock::at_rfc3339(T0).shared(),
registry_answering(untimed_document()),
)
.await;
let response = server.get(&format!("/npm/{NAME}")).await;
assert_eq!(response.status().as_u16(), 403);
let body: Value = serde_json::from_str(&response.text().await.expect("a body"))
.expect("the error body is JSON");
assert_eq!(
body["eligible_at"],
Value::from("2026-04-07T12:00:00Z"),
"the first-seen time plus the cooldown"
);
server.shutdown().await;
let dir = tempfile::tempdir().expect("a temporary data directory");
std::fs::create_dir_all(dir.path().join("state")).expect("the state directory");
std::fs::write(dir.path().join("state").join("firewall.db-wal"), [7u8; 64])
.expect("an unrecoverable write-ahead log with no database beside it");
let server = TestServer::start_in_with_registry(
dir.path(),
config_in(&dir, 0),
TestClock::at_rfc3339(T0).shared(),
registry_answering(untimed_document()),
)
.await;
assert_eq!(
server.status(&format!("/npm/{NAME}")).await,
503,
"nothing that could not be committed is served"
);
server.shutdown().await;
}
#[tokio::test]
async fn first_seen_survives_restart() {
let dir = tempfile::tempdir().expect("a temporary data directory");
let t0 = parse_rfc3339(T0);
let server = TestServer::start_in_with_registry(
dir.path(),
config_in(&dir, ONE_DAY),
TestClock::at_rfc3339(T0).shared(),
registry_answering(untimed_document()),
)
.await;
assert_eq!(
server.status(&format!("/npm/{NAME}")).await,
403,
"the release is untimed, so its cooldown starts the moment it is first seen"
);
server.shutdown().await;
let restarted = TestServer::start_in_with_registry(
dir.path(),
config_in(&dir, ONE_DAY),
TestClock::at_rfc3339("2026-04-08T12:00:00Z").shared(),
registry_answering(untimed_document()),
)
.await;
let document: Value = restarted.json(&format!("/npm/{NAME}")).await;
assert!(
document["versions"].get("1.0.0").is_some(),
"a restart does not reset the clock: the cooldown is counted from the original \
first-seen time, so the release is eligible by now"
);
assert_eq!(
restarted
.running()
.app()
.store()
.get_reference(untimed_reference_id())
.await
.expect("a query")
.expect("the reference survived")
.first_seen_micros,
Some(t0),
"and the stored value is the original one, not a fresh one from this start"
);
restarted.shutdown().await;
}
#[tokio::test]
async fn upstream_timestamp_supersedes_first_seen() {
let dir = tempfile::tempdir().expect("a temporary data directory");
let t0 = parse_rfc3339(T0);
let server = TestServer::start_in_with_registry(
dir.path(),
config_in(&dir, ONE_DAY),
TestClock::at_rfc3339(T0).shared(),
registry_answering(untimed_document()),
)
.await;
assert_eq!(server.status(&format!("/npm/{NAME}")).await, 403);
server.shutdown().await;
let restarted = TestServer::start_in_with_registry(
dir.path(),
config_in(&dir, ONE_DAY),
TestClock::at_rfc3339("2026-04-06T12:10:00Z").shared(),
registry_answering(timed_document()),
)
.await;
let document: Value = restarted.json(&format!("/npm/{NAME}")).await;
assert!(
document["versions"].get("1.0.0").is_some(),
"the upstream time is three months old, so the release is eligible at once"
);
let row = restarted
.running()
.app()
.store()
.get_reference(untimed_reference_id())
.await
.expect("a query")
.expect("the reference is still there");
assert_eq!(
row.publication_micros,
Some(parse_rfc3339("2026-01-01T00:00:00Z"))
);
assert_eq!(
row.first_seen_micros,
Some(t0),
"the first-seen time is kept rather than discarded; upstream supersedes it, \
which is not the same as erasing it"
);
restarted.shutdown().await;
}
#[tokio::test]
async fn a_full_fetch_time_survives_a_restart() {
let dir = tempfile::tempdir().expect("a temporary data directory");
let (first, task) = store(&dir).await;
first
.commit_project_refresh(ProjectRefresh {
ecosystem: Ecosystem::Npm,
name: "widget".to_owned(),
payload: Arc::from(br#"{"name":"widget"}"#.as_slice()),
validators: UpstreamValidators::default(),
validated_at_micros: 9_000,
fetched_at_micros: 5_000,
references: Vec::new(),
})
.await
.expect("the refresh commits");
drop(first);
task.await.expect("the storage task stops");
let (restarted, task) = store(&dir).await;
let row = restarted
.get_project(Ecosystem::Npm, "widget")
.await
.expect("a query")
.expect("the project survived the restart");
assert_eq!(
row.fetched_at_micros, 5_000,
"the last FULL fetch is read back from disk, not re-derived from this start"
);
assert_eq!(
row.validated_at_micros, 9_000,
"and it is a column of its own rather than a second name for the validation time"
);
drop(restarted);
task.await.expect("the storage task stops");
}
#[tokio::test]
async fn a_304_after_a_restart_still_reads_the_stored_full_fetch_time() {
let dir = tempfile::tempdir().expect("a temporary data directory");
let t0 = parse_rfc3339(T0);
let path = format!("/npm/{NAME}");
let registry = FakeRegistry::new();
registry.answer(
&npm_upstream_path(NAME),
FakeAnswer::Validated {
etag: "\"v1\"".to_owned(),
body: untimed_document(),
},
);
let server = TestServer::start_in_with_registry(
dir.path(),
config_in(&dir, ONE_DAY),
TestClock::at_rfc3339(T0).shared(),
Arc::clone(®istry),
)
.await;
assert_eq!(server.status(&path).await, 403);
assert_eq!(
stored_full_fetch(&server).await,
t0,
"the first fetch records itself as a full fetch"
);
server.shutdown().await;
let later = "2026-04-06T13:00:00Z";
let restarted = TestServer::start_in_with_registry(
dir.path(),
config_in(&dir, ONE_DAY),
TestClock::at_rfc3339(later).shared(),
Arc::clone(®istry),
)
.await;
assert_eq!(restarted.status(&path).await, 403);
assert_eq!(
registry.not_modified_answers(),
1,
"the reloaded snapshot really was revalidated conditionally and answered 304"
);
let row = restarted
.running()
.app()
.store()
.get_project(Ecosystem::Npm, NAME)
.await
.expect("a query")
.expect("the project survived the restart");
assert_eq!(
row.fetched_at_micros, t0,
"the 304 carried the original full-fetch time forward across the restart"
);
assert_eq!(
row.validated_at_micros,
parse_rfc3339(later),
"while the validation time is this process's, which is what a 304 renews"
);
restarted.shutdown().await;
}
async fn stored_full_fetch(server: &TestServer) -> i64 {
server
.running()
.app()
.store()
.get_project(Ecosystem::Npm, NAME)
.await
.expect("a query")
.expect("the project is stored")
.fetched_at_micros
}