diff --git a/.changeset/pnpr-hosted-original-digests.md b/.changeset/pnpr-hosted-original-digests.md new file mode 100644 index 0000000000..873dc5c900 --- /dev/null +++ b/.changeset/pnpr-hosted-original-digests.md @@ -0,0 +1,5 @@ +--- +"@pnpm/pnpr": minor +--- + +Hosted pnpr registries now serve newly published original artifacts from registry-scoped SHA-512 digest URLs. diff --git a/pnpr/crates/pnpr/src/error.rs b/pnpr/crates/pnpr/src/error.rs index 8dc98be853..c033fc2cfd 100644 --- a/pnpr/crates/pnpr/src/error.rs +++ b/pnpr/crates/pnpr/src/error.rs @@ -131,6 +131,17 @@ pub enum RegistryError { package: String, }, + #[display("Hosted revision digest already has the maximum of {limit} references")] + #[from(skip)] + RevisionReferenceLimit { limit: usize }, + + #[display("Hosted revision reference index for digest {digest:?} changed while writing")] + #[from(skip)] + RevisionReferenceWriteConflict { + #[error(not(source))] + digest: String, + }, + #[display( "Package {package}@{version} is listed in the local OSV database as vulnerable ({advisories})" )] @@ -259,6 +270,10 @@ impl RegistryError { RegistryError::BadRequest { .. } => "bad_request", RegistryError::VersionAlreadyPublished { .. } => "version_already_published", RegistryError::PackumentWriteConflict { .. } => "packument_write_conflict", + RegistryError::RevisionReferenceLimit { .. } => "revision_reference_limit", + RegistryError::RevisionReferenceWriteConflict { .. } => { + "revision_reference_write_conflict" + } RegistryError::OsvVulnerability { .. } => "osv_vulnerability", RegistryError::RegistrationDisabled => "registration_disabled", RegistryError::TooManyUsers { .. } => "too_many_users", @@ -333,7 +348,9 @@ impl RegistryError { | RegistryError::InvalidAttachment { .. } | RegistryError::BadRequest { .. } => StatusCode::BAD_REQUEST, RegistryError::VersionAlreadyPublished { .. } - | RegistryError::PackumentWriteConflict { .. } => StatusCode::CONFLICT, + | RegistryError::PackumentWriteConflict { .. } + | RegistryError::RevisionReferenceLimit { .. } + | RegistryError::RevisionReferenceWriteConflict { .. } => StatusCode::CONFLICT, RegistryError::Unauthenticated { .. } => StatusCode::UNAUTHORIZED, RegistryError::Forbidden { .. } => StatusCode::FORBIDDEN, RegistryError::TeamsConfigManaged { .. } => StatusCode::FORBIDDEN, diff --git a/pnpr/crates/pnpr/src/journal.rs b/pnpr/crates/pnpr/src/journal.rs index c46d046a80..bee8623eab 100644 --- a/pnpr/crates/pnpr/src/journal.rs +++ b/pnpr/crates/pnpr/src/journal.rs @@ -6,8 +6,8 @@ //! tarball, one packument write per package. A crash in the middle of //! those steps could leave some packages of a batch published and //! others not. The journal closes that window: before anything is -//! promoted, the full intent — the merged packument bytes plus the -//! locations of the staged tmp files — is persisted under +//! promoted, the full intent — the merged packument bytes, revision +//! references, and locations of the staged tmp files — is persisted under //! `.pnpr-journal//` and sealed with a single atomic rename of the //! `commit` marker. [`recover_publish_journal`] runs at startup, before //! the server accepts requests: sealed transactions are rolled forward @@ -27,17 +27,17 @@ use crate::{ config::Config, - error::Result, + error::{RegistryError, Result}, package_name::PackageName, publish::{merge_manifest, now_iso}, storage::{ - RECOVERY_PACKUMENT_WRITE_RETRIES, Storage, TarballFinalize, TarballSlot, unique_tmp_path, + HostedRevisionRefWrite, RECOVERY_PACKUMENT_WRITE_RETRIES, Storage, TarballFinalize, + TarballSlot, is_canonical_revision_ref_owner, unique_tmp_path, }, - upstream::tarball_basename, }; use serde::{Deserialize, Serialize}; use std::{ - collections::HashSet, + collections::{HashMap, HashSet}, io::{self, ErrorKind}, path::{Path, PathBuf}, sync::atomic::{AtomicU64, Ordering}, @@ -77,6 +77,8 @@ struct ManifestPackage { /// packument bytes. packument_file: String, tarballs: Vec, + #[serde(default, skip_serializing_if = "Vec::is_empty")] + revision_refs: Vec, } #[derive(Debug, Serialize, Deserialize)] @@ -87,6 +89,14 @@ struct ManifestTarball { tmp_path: PathBuf, } +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct JournaledRevisionRef { + pub filename: String, + pub digest: String, + pub ref_id: String, + pub bytes: Vec, +} + /// One package of a publish about to be committed, borrowed from the /// handler's staged state. pub struct JournaledPublish<'a> { @@ -95,6 +105,7 @@ pub struct JournaledPublish<'a> { pub org: Option<&'a str>, pub packument: &'a [u8], pub slots: &'a [TarballSlot], + pub revision_refs: &'a [JournaledRevisionRef], } /// Handle to the journal directory of one [`Storage`]. @@ -108,6 +119,7 @@ pub struct PublishJournal { /// startup recovery to (idempotently) re-apply. pub struct SealedTxn { dir: PathBuf, + revision_ref_owner: String, } impl PublishJournal { @@ -120,7 +132,8 @@ impl PublishJournal { /// committed: either the caller applies it now, or startup /// recovery does. pub async fn seal(&self, packages: &[JournaledPublish<'_>]) -> Result { - let dir = self.root.join(txn_id()); + let revision_ref_owner = txn_id(); + let dir = self.root.join(&revision_ref_owner); fs::create_dir_all(&dir).await?; let mut manifest = Manifest { packages: Vec::with_capacity(packages.len()) }; for (index, package) in packages.iter().enumerate() { @@ -138,6 +151,7 @@ impl PublishJournal { tmp_path: slot.tmp_path.clone(), }) .collect(), + revision_refs: package.revision_refs.to_vec(), }); } write_synced(&dir.join(MANIFEST_FILE), &serde_json::to_vec_pretty(&manifest)?).await?; @@ -150,7 +164,7 @@ impl PublishJournal { write_synced(&marker_tmp, b"").await?; fs::rename(&marker_tmp, &marker).await?; let _ = sync_dir(&dir).await; - Ok(SealedTxn { dir }) + Ok(SealedTxn { dir, revision_ref_owner }) } /// Roll every journal entry to a consistent state: sealed @@ -177,7 +191,8 @@ impl PublishJournal { // rollback, which would delete an already-committed publish. // Abort recovery so startup fails loudly instead. if fs::try_exists(dir.join(COMMIT_MARKER)).await? { - roll_forward(storage, &dir).await?; + let revision_ref_owner = revision_ref_owner(&dir)?; + roll_forward(storage, &dir, revision_ref_owner).await?; tracing::info!(txn = %dir.display(), "rolled publish journal entry forward"); } else { roll_back(&dir).await; @@ -189,6 +204,10 @@ impl PublishJournal { } impl SealedTxn { + pub(crate) fn revision_ref_owner(&self) -> &str { + &self.revision_ref_owner + } + /// Apply the sealed transaction now, completing any apply steps that /// have not run yet, and remove the journal entry. This is the same /// idempotent roll-forward startup recovery performs; commit calls it @@ -196,7 +215,7 @@ impl SealedTxn { /// server never leaves a sealed batch partially visible until the /// next restart. pub async fn roll_forward(self, storage: &Storage) -> Result<()> { - roll_forward(storage, &self.dir).await + roll_forward(storage, &self.dir, &self.revision_ref_owner).await } /// Remove the journal entry once the publish is fully applied. @@ -221,7 +240,7 @@ pub async fn recover_publish_journal(config: &Config) -> Result<()> { /// run before the crash: a tmp file that's gone was already promoted, /// and the packument is re-merged into the current on-disk state /// instead of overwriting it. -async fn roll_forward(storage: &Storage, dir: &Path) -> Result<()> { +async fn roll_forward(storage: &Storage, dir: &Path, revision_ref_owner: &str) -> Result<()> { let manifest: Manifest = serde_json::from_slice(&fs::read(dir.join(MANIFEST_FILE)).await?)?; let mut conflicted_tmp_paths = Vec::new(); for package in &manifest.packages { @@ -233,7 +252,7 @@ async fn roll_forward(storage: &Storage, dir: &Path) -> Result<()> { Some(org) => storage.for_hosted(org), None => storage.clone(), }; - let mut conflicted: HashSet<&str> = HashSet::new(); + let mut conflicted_versions = HashSet::new(); for tarball in &package.tarballs { // A missing tmp file was already promoted before the crash, so // skip it. But never read an I/O error as "missing": that would @@ -249,28 +268,80 @@ async fn roll_forward(storage: &Storage, dir: &Path) -> Result<()> { ); match store.finalize_tarball_slot(slot).await? { TarballFinalize::Written | TarballFinalize::AlreadyIdentical => {} - // 다른 replica가 같은 버전의 다른 tarball을 먼저 확정했다. - // winner의 bytes는 immutable이므로 덮어쓰거나 우리 integrity를 - // 노출하면 안 된다. 재시도에서도 충돌을 다시 감지하도록 임시 - // 파일은 유지하고, 아래 merge에서 제외할 filename을 기록한다. + // Another replica finalized different bytes for this version. + // Keep the tmp file so a retry detects the same conflict, and + // exclude the losing version from the merge below. TarballFinalize::Conflict => { conflicted_tmp_paths.push(tarball.tmp_path.as_path()); - conflicted.insert(tarball.filename.as_str()); + let (_, version) = name.parse_tarball_name(&tarball.filename)?; + conflicted_versions.insert(version); } } } } let mut journaled: serde_json::Value = serde_json::from_slice(&fs::read(dir.join(&package.packument_file)).await?)?; - if !conflicted.is_empty() { - drop_conflicted_versions(&mut journaled, &conflicted); + let revision_refs = package + .revision_refs + .iter() + .map(|revision_ref| { + let (_, version) = name.parse_tarball_name(&revision_ref.filename)?; + Ok((revision_ref, version)) + }) + .collect::>>()?; + let mut applied_revision_refs: HashMap> = HashMap::new(); + for (revision_ref, version) in revision_refs { + if conflicted_versions.contains(&version) { + continue; + } + match store + .write_hosted_revision_ref( + &revision_ref.digest, + &revision_ref.ref_id, + revision_ref_owner, + &revision_ref.bytes, + ) + .await + { + Ok(HostedRevisionRefWrite::Claimed | HostedRevisionRefWrite::AlreadyClaimed) => { + applied_revision_refs.entry(version).or_default().push(revision_ref); + } + Ok(HostedRevisionRefWrite::Committed) => {} + Err(crate::error::RegistryError::RevisionReferenceLimit { .. }) => { + conflicted_versions.insert(version.clone()); + if let Some(applied) = applied_revision_refs.remove(&version) { + for applied_ref in applied { + store + .remove_hosted_revision_ref( + &applied_ref.digest, + &applied_ref.ref_id, + revision_ref_owner, + ) + .await?; + } + } + } + Err(err) => return Err(err), + } + } + if !conflicted_versions.is_empty() { + drop_conflicted_versions(&mut journaled, &conflicted_versions); } write_merged_packument(&store, &name, &journaled).await?; + for revision_ref in applied_revision_refs.into_values().flatten() { + store + .commit_hosted_revision_ref( + &revision_ref.digest, + &revision_ref.ref_id, + revision_ref_owner, + ) + .await?; + } } - // journal을 먼저 제거해야 중간 실패가 충돌 상태 없는 재시도를 만들지 않는다. + // Remove the journal before cleaning conflicted tmp files so an interruption + // cannot leave a retry that has lost the evidence needed to detect conflict. fs::remove_dir_all(dir).await?; - // 부모 디렉터리까지 동기화되어 journal 삭제가 내구성을 얻은 뒤에만 충돌 tmp를 지운다. - // 동기화 실패 시 tmp를 남겨 crash 후 journal이 다시 보이더라도 충돌을 재현할 수 있게 한다. + // Only clean conflict evidence after the journal removal is durable. let journal_removal_is_durable = match dir.parent() { Some(parent) => sync_dir(parent).await.is_ok(), None => false, @@ -311,25 +382,19 @@ async fn write_merged_packument( Ok(()) } -/// Drop from a journaled manifest every version whose staged tarball lost a -/// compare-and-swap to another replica. The bytes at that (immutable) version -/// key belong to the winner, so re-merging our `dist`/integrity for the version -/// would advertise metadata that no longer matches the hosted tarball. Versions -/// are matched to `conflicted` staged filenames by their `dist.tarball` -/// basename; a version we cannot match is left in place. -fn drop_conflicted_versions(journaled: &mut serde_json::Value, conflicted: &HashSet<&str>) { +/// Drop from a journaled manifest every version that lost an immutable tarball +/// slot or could not reserve a bounded digest-reference slot. The journal keeps +/// each staged attachment's canonical filename, so callers resolve that name to +/// the version before reaching this helper instead of trusting a publisher- +/// supplied `dist.tarball` URL as the transaction identity. +fn drop_conflicted_versions(journaled: &mut serde_json::Value, conflicted: &HashSet) { let Some(versions) = journaled.get_mut("versions").and_then(serde_json::Value::as_object_mut) else { return; }; let mut removed_versions = HashSet::new(); - versions.retain(|version, manifest| { - let filename = manifest - .get("dist") - .and_then(|dist| dist.get("tarball")) - .and_then(serde_json::Value::as_str) - .and_then(tarball_basename); - let keep = filename.is_none_or(|filename| !conflicted.contains(filename)); + versions.retain(|version, _| { + let keep = !conflicted.contains(version); if !keep { removed_versions.insert(version.clone()); } @@ -372,6 +437,20 @@ fn txn_id() -> String { format!("{millis:016}-{}-{counter}", std::process::id()) } +fn revision_ref_owner(dir: &Path) -> Result<&str> { + let owner = + dir.file_name().and_then(|name| name.to_str()).ok_or_else(|| RegistryError::Internal { + reason: format!("publish journal path has no transaction id: {}", dir.display()), + })?; + if is_canonical_revision_ref_owner(owner) { + Ok(owner) + } else { + Err(RegistryError::Internal { + reason: format!("publish journal transaction id is invalid: {}", dir.display()), + }) + } +} + async fn write_synced(path: &Path, bytes: &[u8]) -> Result<()> { let mut file = fs::File::create(path).await?; file.write_all(bytes).await?; diff --git a/pnpr/crates/pnpr/src/journal/tests.rs b/pnpr/crates/pnpr/src/journal/tests.rs index 4918ce7213..f0edb2976e 100644 --- a/pnpr/crates/pnpr/src/journal/tests.rs +++ b/pnpr/crates/pnpr/src/journal/tests.rs @@ -1,12 +1,13 @@ use super::{ - JournaledPublish, MANIFEST_FILE, Manifest, cleanup_conflicted_tmp_paths, - drop_conflicted_versions, roll_forward, sync_dir, + JournaledPublish, JournaledRevisionRef, MANIFEST_FILE, Manifest, cleanup_conflicted_tmp_paths, + drop_conflicted_versions, revision_ref_owner, roll_forward, sync_dir, }; use crate::{ config::HostedStoreConfig, package_name::PackageName, - storage::{Storage, TarballFinalize}, + storage::{HostedRevisionRefWrite, Storage, TarballFinalize}, }; +use base64::{Engine as _, engine::general_purpose::URL_SAFE_NO_PAD}; use object_store::{ObjectStore, memory::InMemory}; use serde_json::json; use std::{collections::HashSet, sync::Arc}; @@ -19,11 +20,11 @@ fn drop_conflicted_versions_removes_only_the_lost_versions() { "versions": { "1.0.0": { "dist": { "tarball": "http://host/pkg/-/pkg-1.0.0.tgz" } }, "2.0.0": { "dist": { "tarball": "http://host/pkg/-/pkg-2.0.0.tgz" } }, - // A version with no resolvable tarball basename is kept as-is. + // A version outside the conflict set is kept as-is. "3.0.0": { "dist": {} }, } }); - let conflicted: HashSet<&str> = std::iter::once("pkg-1.0.0.tgz").collect(); + let conflicted: HashSet = std::iter::once("1.0.0".to_string()).collect(); drop_conflicted_versions(&mut journaled, &conflicted); @@ -36,21 +37,22 @@ fn drop_conflicted_versions_removes_only_the_lost_versions() { #[test] fn drop_conflicted_versions_tolerates_a_missing_versions_map() { let mut journaled = json!({ "name": "pkg" }); - let conflicted: HashSet<&str> = std::iter::once("pkg-1.0.0.tgz").collect(); + let conflicted: HashSet = std::iter::once("1.0.0".to_string()).collect(); drop_conflicted_versions(&mut journaled, &conflicted); assert_eq!(journaled, json!({ "name": "pkg" })); } #[test] -fn drop_conflicted_versions_uses_shared_tarball_url_semantics() { +fn drop_conflicted_versions_uses_the_canonical_attachment_version() { let mut journaled = json!({ "versions": { - "1.0.0": { "dist": { "tarball": "http://host/pkg/-/pkg-1.0.0.tgz?sig=x" } }, - "2.0.0": { "dist": { "tarball": "http://host/pkg/-/pkg-2.0.0.tgz#fragment" } }, + "1.0.0": { "dist": { "tarball": "http://host/pkg/-/publisher-chosen-name.tgz" } }, + "2.0.0": { "dist": { "tarball": "http://host/pkg/-/another-name.tgz" } }, "3.0.0": { "dist": { "tarball": "http://host/pkg/-/" } }, } }); - let conflicted: HashSet<&str> = ["pkg-1.0.0.tgz", "pkg-2.0.0.tgz"].into_iter().collect(); + let conflicted: HashSet = + ["1.0.0".to_string(), "2.0.0".to_string()].into_iter().collect(); drop_conflicted_versions(&mut journaled, &conflicted); @@ -77,7 +79,7 @@ fn drop_conflicted_versions_removes_references_to_lost_versions() { "modified": "2026-07-03T00:00:00.000Z", }, }); - let conflicted: HashSet<&str> = std::iter::once("pkg-1.0.0.tgz").collect(); + let conflicted: HashSet = std::iter::once("1.0.0".to_string()).collect(); drop_conflicted_versions(&mut journaled, &conflicted); @@ -121,6 +123,185 @@ async fn sync_dir_reports_success_for_a_directory() { sync_dir(tmp.path()).await.unwrap(); } +#[tokio::test] +async fn roll_forward_persists_revision_references() { + let tmp = tempdir().unwrap(); + let object_store: Arc = Arc::new(InMemory::new()); + let storage = Storage::new( + &HostedStoreConfig::S3 { store: object_store, prefix: String::new() }, + tmp.path().join("hosted"), + tmp.path().join("cache"), + ); + let name = PackageName::parse("pkg").unwrap(); + let packument = serde_json::to_vec(&json!({ + "name": "pkg", + "versions": {}, + })) + .unwrap(); + let digest = URL_SAFE_NO_PAD.encode([7_u8; 64]); + let record = br#"{"package":"pkg","version":"1.0.0"}"#.to_vec(); + let revision_refs = [JournaledRevisionRef { + filename: "pkg-1.0.0.tgz".to_string(), + digest: digest.clone(), + ref_id: "a".repeat(64), + bytes: record.clone(), + }]; + let entries = [JournaledPublish { + name: &name, + org: None, + packument: &packument, + slots: &[], + revision_refs: &revision_refs, + }]; + + storage.publish_journal().seal(&entries).await.unwrap().roll_forward(&storage).await.unwrap(); + + assert_eq!(storage.read_hosted_revision_refs(&digest).await.unwrap(), vec![record.clone()]); + assert_eq!( + storage + .write_hosted_revision_ref(&digest, &"a".repeat(64), "later-owner", &record) + .await + .unwrap(), + HostedRevisionRefWrite::Committed, + ); +} + +#[tokio::test] +async fn roll_forward_drops_a_version_that_cannot_reserve_a_revision_reference() { + let tmp = tempdir().unwrap(); + let storage = + Storage::new(&HostedStoreConfig::Fs, tmp.path().join("hosted"), tmp.path().join("cache")); + let digest = URL_SAFE_NO_PAD.encode([7_u8; 64]); + for index in 0..crate::storage::MAX_HOSTED_REVISION_REFS { + storage + .write_hosted_revision_ref(&digest, &format!("{index:064x}"), "existing-owner", b"{}") + .await + .unwrap(); + } + let name = PackageName::parse("pkg").unwrap(); + let packument = serde_json::to_vec(&json!({ + "name": "pkg", + "versions": { + "1.0.0": { + "version": "1.0.0", + "dist": { "tarball": "http://host/pkg/-/publisher-chosen-name.tgz" }, + }, + }, + "dist-tags": { "latest": "1.0.0" }, + "time": { "1.0.0": "2026-07-01T00:00:00.000Z" }, + })) + .unwrap(); + let revision_refs = [JournaledRevisionRef { + filename: "pkg-1.0.0.tgz".to_string(), + digest: digest.clone(), + ref_id: "f".repeat(64), + bytes: br#"{"package":"pkg","version":"1.0.0"}"#.to_vec(), + }]; + let entries = [JournaledPublish { + name: &name, + org: None, + packument: &packument, + slots: &[], + revision_refs: &revision_refs, + }]; + + storage.publish_journal().seal(&entries).await.unwrap().roll_forward(&storage).await.unwrap(); + + let hosted = storage.read_hosted_packument(&name).await.unwrap().unwrap(); + let hosted: serde_json::Value = serde_json::from_slice(&hosted).unwrap(); + assert_eq!(hosted["versions"], json!({})); + assert_eq!(hosted["dist-tags"], json!({})); + assert_eq!(hosted["time"].get("1.0.0"), None); + assert_eq!( + storage.read_hosted_revision_refs(&digest).await.unwrap().len(), + crate::storage::MAX_HOSTED_REVISION_REFS, + ); +} + +#[tokio::test] +async fn roll_forward_only_removes_transaction_owned_references_for_a_dropped_version() { + let tmp = tempdir().unwrap(); + let storage = + Storage::new(&HostedStoreConfig::Fs, tmp.path().join("hosted"), tmp.path().join("cache")); + let transaction_owned_digest = URL_SAFE_NO_PAD.encode([5_u8; 64]); + let previously_owned_digest = URL_SAFE_NO_PAD.encode([6_u8; 64]); + let full_digest = URL_SAFE_NO_PAD.encode([7_u8; 64]); + for index in 0..crate::storage::MAX_HOSTED_REVISION_REFS { + storage + .write_hosted_revision_ref( + &full_digest, + &format!("{index:064x}"), + "existing-owner", + b"{}", + ) + .await + .unwrap(); + } + let name = PackageName::parse("pkg").unwrap(); + let packument = serde_json::to_vec(&json!({ + "name": "pkg", + "versions": { "1.0.0": { "version": "1.0.0" } }, + })) + .unwrap(); + let ref_id = "f".repeat(64); + let record = br#"{"package":"pkg","version":"1.0.0"}"#.to_vec(); + let revision_refs = [ + JournaledRevisionRef { + filename: "pkg-1.0.0.tgz".to_string(), + digest: transaction_owned_digest.clone(), + ref_id: ref_id.clone(), + bytes: record.clone(), + }, + JournaledRevisionRef { + filename: "pkg-1.0.0.tgz".to_string(), + digest: previously_owned_digest.clone(), + ref_id: ref_id.clone(), + bytes: record.clone(), + }, + JournaledRevisionRef { + filename: "pkg-1.0.0.tgz".to_string(), + digest: full_digest, + ref_id: ref_id.clone(), + bytes: record.clone(), + }, + ]; + let entries = [JournaledPublish { + name: &name, + org: None, + packument: &packument, + slots: &[], + revision_refs: &revision_refs, + }]; + + let txn = storage.publish_journal().seal(&entries).await.unwrap(); + let revision_ref_owner = txn.revision_ref_owner().to_string(); + storage + .write_hosted_revision_ref(&transaction_owned_digest, &ref_id, &revision_ref_owner, &record) + .await + .unwrap(); + storage + .write_hosted_revision_ref(&previously_owned_digest, &ref_id, "previous-owner", &record) + .await + .unwrap(); + storage + .commit_hosted_revision_ref(&previously_owned_digest, &ref_id, "previous-owner") + .await + .unwrap(); + txn.roll_forward(&storage).await.unwrap(); + + assert_eq!( + storage.read_hosted_revision_refs(&transaction_owned_digest).await.unwrap(), + Vec::>::new(), + ); + assert_eq!( + storage.read_hosted_revision_refs(&previously_owned_digest).await.unwrap(), + vec![record], + ); + let hosted = storage.read_hosted_packument(&name).await.unwrap().unwrap(); + let hosted: serde_json::Value = serde_json::from_slice(&hosted).unwrap(); + assert_eq!(hosted["versions"], json!({})); +} + #[tokio::test] async fn roll_forward_preserves_tarball_conflict_across_a_later_package_failure() { let tmp = tempdir().unwrap(); @@ -148,7 +329,7 @@ async fn roll_forward_preserves_tarball_conflict_across_a_later_package_failure( "1.0.0": { "version": "1.0.0", "dist": { - "tarball": "http://host/conflicted-pkg/-/conflicted-pkg-1.0.0.tgz", + "tarball": "http://host/conflicted-pkg/-/publisher-chosen-name.tgz", "integrity": "loser", }, }, @@ -166,8 +347,15 @@ async fn roll_forward_preserves_tarball_conflict_across_a_later_package_failure( org: None, packument: &conflicted_packument, slots: &conflicted_slots, + revision_refs: &[], + }, + JournaledPublish { + name: &later_name, + org: None, + packument: b"not-json", + slots: &[], + revision_refs: &[], }, - JournaledPublish { name: &later_name, org: None, packument: b"not-json", slots: &[] }, ]; let txn = storage.publish_journal().seal(&entries).await.unwrap(); let txn_dir = txn.dir.clone(); @@ -196,7 +384,7 @@ async fn roll_forward_preserves_tarball_conflict_across_a_later_package_failure( .await .unwrap(); - roll_forward(&storage, &txn_dir).await.unwrap(); + roll_forward(&storage, &txn_dir, revision_ref_owner(&txn_dir).unwrap()).await.unwrap(); let conflicted_hosted = storage.read_hosted_packument(&conflicted_name).await.unwrap().unwrap(); let conflicted_hosted: serde_json::Value = serde_json::from_slice(&conflicted_hosted).unwrap(); diff --git a/pnpr/crates/pnpr/src/s3.rs b/pnpr/crates/pnpr/src/s3.rs index d762a95245..4c953c37dc 100644 --- a/pnpr/crates/pnpr/src/s3.rs +++ b/pnpr/crates/pnpr/src/s3.rs @@ -14,9 +14,13 @@ //! only the hosted store is pluggable. use crate::{ - error::Result, + error::{RegistryError, Result}, package_name::PackageName, - storage::{STAGED_DIR, staged_id_of_meta_object}, + storage::{ + HOSTED_REVISION_REF_INDEX_FILE, HOSTED_REVISION_REFS_DIR, HostedRevisionRefIndex, + HostedRevisionRefWrite, STAGED_DIR, staged_id_of_meta_object, + wait_after_packument_write_conflict, + }, }; use axum::body::Body; use futures_util::StreamExt; @@ -33,6 +37,7 @@ use std::{ use tokio::fs; const PACKUMENT_FILE: &str = "package.json"; +const REVISION_REF_WRITE_RETRIES: usize = 32; /// The YAML `s3:` block. Selects the object-store hosted backend. /// Credentials fall back to the standard AWS environment variables @@ -341,6 +346,164 @@ impl S3Store { Ok(names) } + pub async fn read_revision_refs(&self, digest: &str) -> Result>> { + let Some((index, _)) = self.read_revision_ref_index(digest).await? else { + return Ok(Vec::new()); + }; + Ok(index.bodies().map(<[u8]>::to_vec).collect()) + } + + pub async fn write_revision_ref( + &self, + digest: &str, + ref_id: &str, + owner: &str, + bytes: &[u8], + ) -> Result { + for attempt in 0..REVISION_REF_WRITE_RETRIES { + let current = self.read_revision_ref_index(digest).await?; + let (mut index, version) = match current { + Some((index, version)) => (index, Some(version)), + None => (HostedRevisionRefIndex::default(), None), + }; + let outcome = index.insert(ref_id, owner, bytes)?; + if outcome != HostedRevisionRefWrite::Claimed { + return Ok(outcome); + } + let mode = match version { + Some(version) => PutMode::Update(version), + None => PutMode::Create, + }; + match self + .store + .put_opts( + &self.revision_ref_index_key(digest), + PutPayload::from(index.to_bytes()), + PutOptions { mode, ..PutOptions::default() }, + ) + .await + { + Ok(_) => return Ok(HostedRevisionRefWrite::Claimed), + Err( + object_store::Error::AlreadyExists { .. } + | object_store::Error::NotFound { .. } + | object_store::Error::Precondition { .. }, + ) => { + if attempt + 1 < REVISION_REF_WRITE_RETRIES { + wait_after_packument_write_conflict(attempt).await; + } + } + Err(err) => return Err(err.into()), + } + } + let mut index = self + .read_revision_ref_index(digest) + .await? + .map_or_else(HostedRevisionRefIndex::default, |(index, _)| index); + let outcome = index.insert(ref_id, owner, bytes)?; + if outcome != HostedRevisionRefWrite::Claimed { + return Ok(outcome); + } + Err(RegistryError::RevisionReferenceWriteConflict { digest: digest.to_string() }) + } + + pub async fn remove_revision_ref(&self, digest: &str, ref_id: &str, owner: &str) -> Result<()> { + for attempt in 0..REVISION_REF_WRITE_RETRIES { + let Some((mut index, version)) = self.read_revision_ref_index(digest).await? else { + return Ok(()); + }; + if !index.remove_if_owned(ref_id, owner) { + return Ok(()); + } + match self + .store + .put_opts( + &self.revision_ref_index_key(digest), + PutPayload::from(index.to_bytes()), + PutOptions { mode: PutMode::Update(version), ..PutOptions::default() }, + ) + .await + { + Ok(_) => return Ok(()), + Err( + object_store::Error::NotFound { .. } | object_store::Error::Precondition { .. }, + ) => { + if attempt + 1 < REVISION_REF_WRITE_RETRIES { + wait_after_packument_write_conflict(attempt).await; + } + } + Err(err) => return Err(err.into()), + } + } + if self + .read_revision_ref_index(digest) + .await? + .is_none_or(|(index, _)| !index.is_owned_by(ref_id, owner)) + { + return Ok(()); + } + Err(RegistryError::RevisionReferenceWriteConflict { digest: digest.to_string() }) + } + + pub async fn commit_revision_ref(&self, digest: &str, ref_id: &str, owner: &str) -> Result<()> { + for attempt in 0..REVISION_REF_WRITE_RETRIES { + let Some((mut index, version)) = self.read_revision_ref_index(digest).await? else { + return Err(RegistryError::Internal { + reason: "hosted revision reference is missing during commit".to_string(), + }); + }; + if !index.commit_if_owned(ref_id, owner)? { + return Ok(()); + } + match self + .store + .put_opts( + &self.revision_ref_index_key(digest), + PutPayload::from(index.to_bytes()), + PutOptions { mode: PutMode::Update(version), ..PutOptions::default() }, + ) + .await + { + Ok(_) => return Ok(()), + Err( + object_store::Error::NotFound { .. } | object_store::Error::Precondition { .. }, + ) => { + if attempt + 1 < REVISION_REF_WRITE_RETRIES { + wait_after_packument_write_conflict(attempt).await; + } + } + Err(err) => return Err(err.into()), + } + } + let Some((mut index, _)) = self.read_revision_ref_index(digest).await? else { + return Err(RegistryError::Internal { + reason: "hosted revision reference is missing during commit".to_string(), + }); + }; + if !index.commit_if_owned(ref_id, owner)? { + return Ok(()); + } + Err(RegistryError::RevisionReferenceWriteConflict { digest: digest.to_string() }) + } + + async fn read_revision_ref_index( + &self, + digest: &str, + ) -> Result> { + match self.store.get(&self.revision_ref_index_key(digest)).await { + Ok(result) => { + let version = UpdateVersion { + e_tag: result.meta.e_tag.clone(), + version: result.meta.version.clone(), + }; + let index = HostedRevisionRefIndex::from_bytes(&result.bytes().await?)?; + Ok(Some((index, version))) + } + Err(object_store::Error::NotFound { .. }) => Ok(None), + Err(err) => Err(err.into()), + } + } + fn packument_key(&self, name: &PackageName) -> ObjectPath { ObjectPath::from(format!("{}{}/{PACKUMENT_FILE}", self.prefix, name.as_str())) } @@ -349,6 +512,13 @@ impl S3Store { ObjectPath::from(format!("{}{}/{filename}", self.prefix, name.as_str())) } + fn revision_ref_index_key(&self, digest: &str) -> ObjectPath { + ObjectPath::from(format!( + "{}{HOSTED_REVISION_REFS_DIR}/{digest}/{HOSTED_REVISION_REF_INDEX_FILE}", + self.prefix, + )) + } + // Staged-publish records (see `storage::Storage`'s staged section for // the layout contract shared with the fs backend). diff --git a/pnpr/crates/pnpr/src/s3/tests.rs b/pnpr/crates/pnpr/src/s3/tests.rs index 07ad136d72..ae771b0b10 100644 --- a/pnpr/crates/pnpr/src/s3/tests.rs +++ b/pnpr/crates/pnpr/src/s3/tests.rs @@ -1,6 +1,6 @@ use super::{Body, ObjectStore, S3Settings, S3Store}; -use crate::package_name::PackageName; -use object_store::memory::InMemory; +use crate::{package_name::PackageName, storage::HostedRevisionRefWrite}; +use object_store::{ObjectStoreExt, PutPayload, memory::InMemory, path::Path as ObjectPath}; use std::sync::Arc; use tempfile::tempdir; @@ -201,6 +201,151 @@ async fn lists_hosted_package_names() { } } +#[tokio::test] +async fn revision_refs_roundtrip_under_the_configured_prefix() { + for prefix in ["", "packages"] { + let (store, _staging) = store_with_prefix(prefix); + let digest = "A".repeat(86); + assert_eq!(store.read_revision_refs(&digest).await.unwrap(), Vec::>::new()); + + store.write_revision_ref(&digest, &"a".repeat(64), "owner-a", b"first").await.unwrap(); + store.write_revision_ref(&digest, &"b".repeat(64), "owner-a", b"second").await.unwrap(); + let mut refs = store.read_revision_refs(&digest).await.unwrap(); + refs.sort(); + assert_eq!(refs, vec![b"first".to_vec(), b"second".to_vec()]); + } +} + +#[tokio::test] +async fn revision_ref_removal_is_scoped_to_its_owner() { + let (store, _staging) = store_with_prefix("packages"); + let digest = "A".repeat(86); + let ref_id = "a".repeat(64); + assert_eq!( + store.write_revision_ref(&digest, &ref_id, "owner-a", b"record").await.unwrap(), + HostedRevisionRefWrite::Claimed, + ); + assert_eq!( + store.write_revision_ref(&digest, &ref_id, "owner-b", b"record").await.unwrap(), + HostedRevisionRefWrite::Claimed, + ); + + store.remove_revision_ref(&digest, &ref_id, "owner-a").await.unwrap(); + assert_eq!(store.read_revision_refs(&digest).await.unwrap(), vec![b"record".to_vec()]); + + store.remove_revision_ref(&digest, &ref_id, "owner-a").await.unwrap(); + assert_eq!(store.read_revision_refs(&digest).await.unwrap(), vec![b"record".to_vec()]); + + store.commit_revision_ref(&digest, &ref_id, "owner-b").await.unwrap(); + store.remove_revision_ref(&digest, &ref_id, "owner-b").await.unwrap(); + assert_eq!(store.read_revision_refs(&digest).await.unwrap(), vec![b"record".to_vec()]); + + assert_eq!( + store.write_revision_ref(&digest, &ref_id, "owner-a", b"record").await.unwrap(), + HostedRevisionRefWrite::Committed, + ); +} + +#[tokio::test] +async fn concurrent_revision_ref_claims_survive_other_owner_removal() { + let (store, _staging) = store_with_prefix("packages"); + let digest = "A".repeat(86); + let ref_id = "a".repeat(64); + let first = { + let store = store.clone(); + let digest = digest.clone(); + let ref_id = ref_id.clone(); + tokio::spawn(async move { + store.write_revision_ref(&digest, &ref_id, "owner-a", b"record").await + }) + }; + let second = { + let store = store.clone(); + let digest = digest.clone(); + let ref_id = ref_id.clone(); + tokio::spawn(async move { + store.write_revision_ref(&digest, &ref_id, "owner-b", b"record").await + }) + }; + + assert_eq!(first.await.unwrap().unwrap(), HostedRevisionRefWrite::Claimed); + assert_eq!(second.await.unwrap().unwrap(), HostedRevisionRefWrite::Claimed); + store.remove_revision_ref(&digest, &ref_id, "owner-a").await.unwrap(); + assert_eq!(store.read_revision_refs(&digest).await.unwrap(), vec![b"record".to_vec()]); + store.remove_revision_ref(&digest, &ref_id, "owner-b").await.unwrap(); + assert_eq!(store.read_revision_refs(&digest).await.unwrap(), Vec::>::new()); +} + +#[tokio::test] +async fn revision_ref_writes_enforce_the_read_bound() { + let (store, _staging) = store_with_prefix("packages"); + let digest = "A".repeat(86); + for index in 0..crate::storage::MAX_HOSTED_REVISION_REFS { + store + .write_revision_ref(&digest, &format!("{index:064x}"), "owner-a", b"{}") + .await + .unwrap(); + } + + let overflow = crate::storage::MAX_HOSTED_REVISION_REFS; + let err = store + .write_revision_ref(&digest, &format!("{overflow:064x}"), "owner-a", b"{}") + .await + .unwrap_err(); + assert!(matches!( + err, + crate::error::RegistryError::RevisionReferenceLimit { limit } + if limit == crate::storage::MAX_HOSTED_REVISION_REFS + )); + + assert_eq!( + store.write_revision_ref(&digest, &"0".repeat(64), "owner-a", b"{}").await.unwrap(), + HostedRevisionRefWrite::AlreadyClaimed, + ); + store + .store + .put( + &ObjectPath::from(format!("packages/.revisions/sha512/{digest}/not-a-reference.json")), + PutPayload::from_static(b"stray"), + ) + .await + .unwrap(); + let refs = store.read_revision_refs(&digest).await.unwrap(); + assert_eq!(refs.len(), crate::storage::MAX_HOSTED_REVISION_REFS); + assert!(refs.iter().all(|bytes| bytes == b"{}")); +} + +#[tokio::test] +async fn concurrent_revision_ref_writes_cannot_exceed_the_limit() { + let (store, _staging) = store_with_prefix("packages"); + let digest = "A".repeat(86); + let mut writes = Vec::new(); + for index in 0..crate::storage::MAX_HOSTED_REVISION_REFS * 2 { + let store = store.clone(); + let digest = digest.clone(); + writes.push(tokio::spawn(async move { + store.write_revision_ref(&digest, &format!("{index:064x}"), "owner-a", b"{}").await + })); + } + + let mut written = 0; + let mut rejected = 0; + for write in writes { + match write.await.unwrap() { + Ok(HostedRevisionRefWrite::Claimed) => written += 1, + Ok(outcome) => panic!("unexpected revision-reference write outcome: {outcome:?}"), + Err(crate::error::RegistryError::RevisionReferenceLimit { .. }) => rejected += 1, + Err(err) => panic!("unexpected revision-reference write error: {err}"), + } + } + assert_eq!(written, crate::storage::MAX_HOSTED_REVISION_REFS); + assert_eq!(rejected, crate::storage::MAX_HOSTED_REVISION_REFS); + assert_eq!( + store.read_revision_refs(&digest).await.unwrap().len(), + crate::storage::MAX_HOSTED_REVISION_REFS, + ); +} + #[test] fn prefix_normalizes() { let normalized = |prefix: Option<&str>| { diff --git a/pnpr/crates/pnpr/src/server.rs b/pnpr/crates/pnpr/src/server.rs index 840d50577a..7a86aae06b 100644 --- a/pnpr/crates/pnpr/src/server.rs +++ b/pnpr/crates/pnpr/src/server.rs @@ -2,7 +2,7 @@ use crate::{ auth::{AuthState, TokenRecord, UpsertOutcome, identify}, config::{Config, HostedConfig}, error::RegistryError, - journal::JournaledPublish, + journal::{JournaledPublish, JournaledRevisionRef}, package_name::PackageName, policy::{Identity, PackageRules}, publish::{ @@ -36,7 +36,10 @@ use axum::{ }; use chrono::Utc; use indexmap::IndexMap; -use pnpm_crypto_hash::integrity_addressed_tarball_integrity; +use pnpm_crypto_hash::{ + create_hex_hash, integrity_addressed_tarball_integrity, integrity_addressed_tarball_path, +}; +use pnpm_lockfile::TarballRevision; use serde_json::{Value, json}; use ssri::Integrity; use std::{ @@ -112,6 +115,12 @@ struct AppInner { osv_index: Option>, } +#[derive(serde::Serialize, serde::Deserialize)] +struct HostedOriginalRef { + package: String, + version: String, +} + /// Per-package serialization for the read-modify-write packument flows /// (publish, dist-tag changes, partial-unpublish). Without it, two /// concurrent publishes of the same package both read the old @@ -750,12 +759,7 @@ async fn get_four_segments( ) -> Response { if a == "-" && b == "tarballs" && c == "sha512" { let Some(registry) = default_registry_target(&state) else { return not_found() }; - let response = serve_revision_tarball(&state, &identity, ®istry, &d).await; - return if revision_registry_is_private(&state, ®istry) { - private_no_cache(response) - } else { - response - }; + return serve_revision_tarball(&state, &identity, ®istry, &d).await; } if a == "-" && b == "package" && d == "dist-tags" { let response = get_dist_tags(&state, &identity, None, &c).await; @@ -788,7 +792,7 @@ async fn get_four_segments( /// * `/~/-/package//dist-tags` — dist-tags through a registry /// endpoint. /// * `/~/-/org/{scope}/team` — org teams through a registry endpoint. -/// * `/~/-/tarballs/sha512/{digest}` — immutable artifact through an upstream. +/// * `/~/-/tarballs/sha512/{digest}`: immutable artifact through an addressed registry. /// /// Every other 5-segment GET is a not-found catchall (the route exists so /// DELETE/PUT can sit on the same path). @@ -802,7 +806,7 @@ async fn get_five_segments( } if let Some(registry) = a.strip_prefix('~').filter(|registry| !registry.is_empty()) { if b == "-" && c == "tarballs" && d == "sha512" { - return private_no_cache(serve_revision_tarball(&state, &identity, registry, &e).await); + return serve_revision_tarball(&state, &identity, registry, &e).await; } if b.starts_with('@') && d == "-" { let full = format!("{b}/{c}"); @@ -1873,6 +1877,25 @@ async fn serve_revision_tarball( let Some(integrity) = integrity_addressed_tarball_integrity(digest) else { return not_found(); }; + if matches!(state.inner.config.registries.get(registry), Some(Registry::Upstream { .. })) { + let response = + serve_upstream_revision_tarball(state, identity, registry, digest, &integrity).await; + return if revision_registry_is_private(state, registry) { + private_no_cache(response) + } else { + response + }; + } + serve_hosted_revision_tarball(state, identity, registry, digest, &integrity).await +} + +async fn serve_upstream_revision_tarball( + state: &AppState, + identity: &Identity, + registry: &str, + digest: &str, + integrity: &Integrity, +) -> Response { let upstream = match authorized_revision_upstream(state, identity, registry) { Ok(upstream) => upstream, Err(response) => return *response, @@ -1881,7 +1904,12 @@ async fn serve_revision_tarball( if upstream.caches() { match state.inner.storage.open_upstream_revision_tarball(&namespace, digest).await { Ok(Some((file, len))) => { - return tarball_response(streaming::stream_file(file), Some(len)); + return revision_tarball_response( + streaming::stream_file(file), + Some(len), + digest, + integrity, + ); } Ok(None) => {} Err(err) => { @@ -1903,27 +1931,242 @@ async fn serve_revision_tarball( return match streaming::download_verified_to_temp( response, write, - &integrity, + integrity, MAX_TARBALL_BYTES, ) .await { - Ok((file, len, tmp_path)) => { - tarball_response(streaming::stream_file_and_remove(file, tmp_path), Some(len)) - } + Ok((file, len, tmp_path)) => revision_tarball_response( + streaming::stream_file_and_remove(file, tmp_path), + Some(len), + digest, + integrity, + ), Err(err) => { error_response(&tarball_stream_error_for_package(err, "registry revision", digest)) } }; } - match streaming::stream_verified_to_cache(response, write, &integrity, MAX_TARBALL_BYTES) { - Ok(body) => tarball_response(body, None), + match streaming::stream_verified_to_cache(response, write, integrity, MAX_TARBALL_BYTES) { + Ok(body) => revision_tarball_response(body, None, digest, integrity), Err(err) => { error_response(&tarball_stream_error_for_package(err, "registry revision", digest)) } } } +async fn serve_hosted_revision_tarball( + state: &AppState, + identity: &Identity, + registry: &str, + digest: &str, + integrity: &Integrity, +) -> Response { + let sources = hosted_revision_sources(state, registry); + if sources.is_empty() { + return not_found(); + } + + let mut private_refs = Vec::new(); + let mut policy_error = None; + for source in sources { + let Some(hosted) = state.inner.config.hosted.get(&source) else { + continue; + }; + let storage = state.inner.storage.for_hosted(&hosted.org); + let refs = match hosted_revision_refs(&storage, digest).await { + Ok(refs) => refs, + Err(err) => return private_no_cache(error_response(&err)), + }; + for original in refs { + let package = match PackageName::parse(&original.package) { + Ok(package) => package, + Err(err) => return private_no_cache(error_response(&err)), + }; + let filename = package.tarball_name_for_version(&original.version); + if let Err(err) = package.canonicalize_tarball_name(&filename) { + return private_no_cache(error_response(&err)); + } + if !matches!( + resolve_registry_source(state, registry, package.as_str()), + RegistrySource::Hosted(resolved) if resolved == source, + ) || !matches!( + hosted_gate(state, identity, &source, package.as_str()), + HostedGate::Allowed(_), + ) { + continue; + } + match hosted_original_is_current(&storage, &package, &original.version, digest).await { + Ok(true) => {} + Ok(false) => continue, + Err(err) => return private_no_cache(error_response(&err)), + } + if let Err(err) = ensure_osv_allowed(state, &package, &original.version) { + policy_error.get_or_insert(err); + continue; + } + if matches!( + hosted_gate(state, &Identity::Anonymous, &source, package.as_str()), + HostedGate::Allowed(_), + ) { + let response = open_hosted_revision_tarball( + &storage, + &package, + &original.version, + digest, + integrity, + ) + .await; + if response.status() != StatusCode::NOT_FOUND { + return response; + } + continue; + } + private_refs.push((storage.clone(), package, original.version)); + } + } + + for (storage, package, version) in private_refs { + let response = + open_hosted_revision_tarball(&storage, &package, &version, digest, integrity).await; + if response.status() != StatusCode::NOT_FOUND { + return response; + } + } + if let Some(err) = policy_error { + return private_no_cache(error_response(&err)); + } + private_no_cache(not_found()) +} + +fn hosted_revision_sources(state: &AppState, registry: &str) -> Vec { + match state.inner.config.registries.get(registry) { + Some(Registry::Hosted { .. }) => vec![registry.to_string()], + Some(Registry::Router { sources }) => sources + .iter() + .filter(|source| { + matches!(state.inner.config.registries.get(source), Some(Registry::Hosted { .. })) + }) + .cloned() + .collect(), + Some(Registry::Upstream { .. }) | None => Vec::new(), + } +} + +async fn hosted_revision_refs( + storage: &Storage, + digest: &str, +) -> Result, RegistryError> { + storage + .read_hosted_revision_refs(digest) + .await? + .into_iter() + .map(|bytes| serde_json::from_slice(&bytes).map_err(RegistryError::Json)) + .collect() +} + +async fn hosted_original_is_current( + storage: &Storage, + package: &PackageName, + version: &str, + digest: &str, +) -> Result { + let Some(bytes) = storage.read_hosted_packument(package).await? else { + return Ok(false); + }; + let packument = serde_json::from_slice::(&bytes)?; + Ok(packument + .versions + .get(version) + .and_then(|manifest| manifest.dist.as_ref()) + .and_then(original_integrity) + .and_then(|integrity| integrity_addressed_tarball_path(&integrity)) + .is_some_and(|path| path == format!("-/tarballs/sha512/{digest}"))) +} + +async fn open_hosted_revision_tarball( + storage: &Storage, + package: &PackageName, + version: &str, + digest: &str, + integrity: &Integrity, +) -> Response { + let filename = package.tarball_name_for_version(version); + match storage.open_hosted_tarball(package, &filename).await { + Ok(Some((body, len))) => revision_tarball_response(body, len, digest, integrity), + Ok(None) => not_found(), + Err(err) => error_response(&err), + } +} + +#[derive(serde::Deserialize)] +struct HostedRevisionPackument { + #[serde(default)] + versions: IndexMap, +} + +#[derive(serde::Deserialize)] +struct HostedRevisionManifest { + #[serde(default)] + dist: Option, +} + +#[derive(serde::Deserialize)] +struct HostedRevisionDist { + #[serde(default)] + integrity: Option, + #[serde(default)] + revision: RevisionField, + #[serde(default)] + revisions: Vec, +} + +#[derive(serde::Deserialize)] +struct HostedRevisionRecord { + #[serde(default)] + revision: Value, + #[serde(default)] + integrity: Option, +} + +#[derive(Default)] +enum RevisionField { + #[default] + Missing, + Present(Value), +} + +impl<'de> serde::Deserialize<'de> for RevisionField { + fn deserialize(deserializer: Deserializer) -> Result + where + Deserializer: serde::Deserializer<'de>, + { + ::deserialize(deserializer).map(Self::Present) + } +} + +fn original_integrity(dist: &HostedRevisionDist) -> Option { + let RevisionField::Present(revision) = &dist.revision else { + return dist.integrity.as_deref()?.parse().ok(); + }; + let selected_revision = + revision.as_u64().and_then(|revision| TarballRevision::try_from(revision).ok())?.get(); + let selected: Vec<_> = dist + .revisions + .iter() + .filter(|record| record.revision.as_u64() == Some(selected_revision)) + .collect(); + if selected.len() != 1 || selected[0].integrity.as_deref() != dist.integrity.as_deref() { + return None; + } + let originals: Vec<_> = + dist.revisions.iter().filter(|record| record.revision.as_u64() == Some(0)).collect(); + if originals.len() != 1 { + return None; + } + originals[0].integrity.as_deref()?.parse().ok() +} + /// The response for a cached upstream tarball, or `None` on a cache miss. A /// cache-open fault is logged and treated as a miss so the caller falls back /// to the upstream fetch rather than failing the request. @@ -3051,6 +3294,7 @@ struct StagedPublish { merged_bytes: Vec, base_version: Option, slots: Vec, + original_refs: Vec, /// Hosted-org storage namespace this publish targets, or `None` for the /// flat (path-less) hosted store. Threaded into the commit and journal so /// the write — and any crash-recovery roll-forward — lands in the right org. @@ -3135,6 +3379,10 @@ async fn stage_publish( let existing: Option = hosted.clone(); let merged = merge_manifest(existing.as_ref(), &incoming, hosted.as_ref(), now_iso); let merged_bytes = serde_json::to_vec_pretty(&merged).map_err(RegistryError::Json)?; + let original_refs = prepared + .iter() + .filter_map(|attachment| staged_hosted_original_ref(&name, attachment)) + .collect(); // `incoming` is no longer needed; drop it so the base64 strings // inside go away as soon as `prepared` (which owns each one) is // drained below. @@ -3179,12 +3427,29 @@ async fn stage_publish( merged_bytes, base_version, slots: written_slots, + original_refs, org: org.map(str::to_string), }) } +fn staged_hosted_original_ref( + package: &PackageName, + attachment: &PreparedAttachment, +) -> Option { + let integrity: Integrity = attachment.dist.get("integrity")?.as_str()?.parse().ok()?; + let path = integrity_addressed_tarball_path(&integrity)?; + let digest = path.strip_prefix("-/tarballs/sha512/")?.to_string(); + let record = HostedOriginalRef { + package: package.as_str().to_string(), + version: attachment.version.clone(), + }; + let bytes = serde_json::to_vec(&record).expect("hosted original reference serializes"); + let ref_id = create_hex_hash(&format!("{}\0{}", record.package, record.version)); + Some(JournaledRevisionRef { filename: attachment.canonical.clone(), digest, ref_id, bytes }) +} + /// Make every staged publish visible. The full intent — merged -/// packument bytes plus the staged tmp-file locations — is sealed into +/// packument bytes, revision references, and staged tmp-file locations — is sealed into /// the commit journal first, so a crash or I/O failure mid-apply can /// never leave the batch partially published: startup recovery rolls /// a sealed transaction forward. If sealing itself fails, nothing was @@ -3205,6 +3470,7 @@ async fn commit_publishes( org: stage.org.as_deref(), packument: &stage.merged_bytes, slots: &stage.slots, + revision_refs: &stage.original_refs, }) .collect(); let sealed = journal.seal(&entries).await; @@ -3218,6 +3484,7 @@ async fn commit_publishes( return Err(err); } }; + let revision_ref_owner = txn.revision_ref_owner().to_string(); // Past the seal the transaction is committed: the apply below is pure // roll-forward, and failures must NOT clean up the staged files. If // the apply fails partway, complete it immediately via the same @@ -3245,6 +3512,16 @@ async fn commit_publishes( } } } + for original in &stage.original_refs { + store + .write_hosted_revision_ref( + &original.digest, + &original.ref_id, + &revision_ref_owner, + &original.bytes, + ) + .await?; + } match store .write_hosted_packument_if_current( &stage.name, @@ -3253,7 +3530,17 @@ async fn commit_publishes( ) .await? { - PackumentWrite::Written => {} + PackumentWrite::Written => { + for original in &stage.original_refs { + store + .commit_hosted_revision_ref( + &original.digest, + &original.ref_id, + &revision_ref_owner, + ) + .await?; + } + } // Tarballs are already promoted at this point. A conflict means // another replica advanced the packument since staging, so the // base version is stale. Surfacing it drops into the seal's @@ -3279,7 +3566,13 @@ async fn commit_publishes( } Err(apply_err) => { tracing::warn!(error = %apply_err, "publish apply failed after seal; rolling forward"); - txn.roll_forward(&state.inner.storage).await.map_err(|_| apply_err) + let report_conflict = + matches!(&apply_err, RegistryError::RevisionReferenceLimit { .. }); + match txn.roll_forward(&state.inner.storage).await { + Ok(()) if report_conflict => Err(apply_err), + Ok(()) => Ok(()), + Err(_) => Err(apply_err), + } } } } @@ -4511,6 +4804,33 @@ fn tarball_response(body: Body, content_length: Option) -> Response { builder.body(body).expect("static-shape response always builds") } +fn revision_tarball_response( + body: Body, + content_length: Option, + digest: &str, + integrity: &Integrity, +) -> Response { + let mut response = tarball_response(body, content_length); + let headers = response.headers_mut(); + headers.insert( + header::CACHE_CONTROL, + "public, max-age=31536000, immutable".parse().expect("static cache control is valid"), + ); + headers.insert( + header::ETAG, + format!(r#""{digest}""#).parse().expect("canonical base64url digest is a valid ETag"), + ); + if let [hash] = integrity.hashes.as_slice() { + headers.insert( + "content-digest", + format!("sha-512=:{}:", hash.digest) + .parse() + .expect("canonical base64 digest is a valid header value"), + ); + } + private_no_cache(response) +} + fn not_found() -> Response { (StatusCode::NOT_FOUND, "Not Found").into_response() } diff --git a/pnpr/crates/pnpr/src/server/tests.rs b/pnpr/crates/pnpr/src/server/tests.rs index 6617d36cb2..7d83545a18 100644 --- a/pnpr/crates/pnpr/src/server/tests.rs +++ b/pnpr/crates/pnpr/src/server/tests.rs @@ -1,6 +1,7 @@ use super::{ - ConnectInfo, PeerAddr, bearer_credentials, canonical_ip, cidr_contains, cidr_whitelist_allows, - is_write_method, router_with_auth, token_timestamp_millis, + ConnectInfo, HostedRevisionDist, HostedRevisionRecord, PeerAddr, RevisionField, + bearer_credentials, canonical_ip, cidr_contains, cidr_whitelist_allows, is_write_method, + original_integrity, router_with_auth, token_timestamp_millis, }; use crate::{ auth::{AuthState, TokenBackend, TokenRecord, UserStore}, @@ -26,6 +27,81 @@ fn token_timestamp_millis_saturates_before_i64_conversion() { assert_eq!(token_timestamp_millis(u64::MAX), i64::MAX / 1000 * 1000); } +#[test] +fn original_integrity_uses_current_integrity_before_any_replacement() { + let integrity = format!("sha512-{}==", "A".repeat(86)); + let dist = HostedRevisionDist { + integrity: Some(integrity.clone()), + revision: RevisionField::Missing, + revisions: Vec::new(), + }; + assert_eq!(original_integrity(&dist).unwrap().to_string(), integrity); +} + +#[test] +fn original_integrity_rejects_an_explicit_null_revision() { + let dist = HostedRevisionDist { + integrity: Some(format!("sha512-{}==", "A".repeat(86))), + revision: RevisionField::Present(serde_json::Value::Null), + revisions: Vec::new(), + }; + assert_eq!(original_integrity(&dist).map(|integrity| integrity.to_string()), None); +} + +#[test] +fn original_integrity_rejects_an_explicit_revision_zero() { + let integrity = format!("sha512-{}==", "A".repeat(86)); + let dist = HostedRevisionDist { + integrity: Some(integrity.clone()), + revision: RevisionField::Present(serde_json::json!(0)), + revisions: vec![HostedRevisionRecord { + revision: serde_json::json!(0), + integrity: Some(integrity), + }], + }; + assert_eq!(original_integrity(&dist).map(|integrity| integrity.to_string()), None); +} + +#[test] +fn original_integrity_uses_validated_revision_zero_after_replacement() { + let original = format!("sha512-{}==", "A".repeat(86)); + let replacement = format!("sha512-{}Q==", "B".repeat(85)); + let dist = HostedRevisionDist { + integrity: Some(replacement.clone()), + revision: RevisionField::Present(serde_json::json!(1)), + revisions: vec![ + HostedRevisionRecord { + revision: serde_json::json!(0), + integrity: Some(original.clone()), + }, + HostedRevisionRecord { revision: serde_json::json!(1), integrity: Some(replacement) }, + ], + }; + assert_eq!(original_integrity(&dist).unwrap().to_string(), original); +} + +#[test] +fn original_integrity_rejects_ambiguous_or_inconsistent_history() { + let original = format!("sha512-{}==", "A".repeat(86)); + let replacement = format!("sha512-{}Q==", "B".repeat(85)); + let dist = HostedRevisionDist { + integrity: Some(replacement), + revision: RevisionField::Present(serde_json::json!(1)), + revisions: vec![ + HostedRevisionRecord { + revision: serde_json::json!(0), + integrity: Some(original.clone()), + }, + HostedRevisionRecord { revision: serde_json::json!(0), integrity: Some(original) }, + HostedRevisionRecord { + revision: serde_json::json!(1), + integrity: Some(format!("sha512-{}g==", "C".repeat(85))), + }, + ], + }; + assert_eq!(original_integrity(&dist).map(|integrity| integrity.to_string()), None); +} + // --------------------------------------------------------------- // CIDR matching // --------------------------------------------------------------- diff --git a/pnpr/crates/pnpr/src/storage.rs b/pnpr/crates/pnpr/src/storage.rs index 3f4dcee578..bbcd6f35a3 100644 --- a/pnpr/crates/pnpr/src/storage.rs +++ b/pnpr/crates/pnpr/src/storage.rs @@ -8,7 +8,9 @@ use crate::{ use axum::body::Body; use object_store::UpdateVersion; use pnpm_crypto_hash::integrity_addressed_tarball_integrity; +use serde::{Deserialize, Serialize}; use std::{ + collections::HashSet, io::{ErrorKind, SeekFrom}, path::{Path, PathBuf}, sync::{ @@ -23,6 +25,138 @@ use tokio::{ }; const PACKUMENT_FILE: &str = "package.json"; +pub(crate) const HOSTED_REVISION_REFS_DIR: &str = ".revisions/sha512"; +pub(crate) const HOSTED_REVISION_REF_INDEX_FILE: &str = "index.json"; +/// Bounds both the persisted candidate set and work triggered by one digest request. +pub(crate) const MAX_HOSTED_REVISION_REFS: usize = 32; + +#[derive(Debug, Default, Serialize, Deserialize)] +pub(crate) struct HostedRevisionRefIndex { + refs: Vec, +} + +#[derive(Debug, Serialize, Deserialize)] +struct HostedRevisionRefIndexEntry { + id: String, + committed: bool, + pending_owners: Vec, + bytes: Vec, +} + +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub(crate) enum HostedRevisionRefWrite { + Claimed, + AlreadyClaimed, + Committed, +} + +impl HostedRevisionRefIndex { + pub(crate) fn from_bytes(bytes: &[u8]) -> Result { + let index: Self = serde_json::from_slice(bytes)?; + if index.refs.len() > MAX_HOSTED_REVISION_REFS { + return Err(RegistryError::RevisionReferenceLimit { limit: MAX_HOSTED_REVISION_REFS }); + } + let mut seen = HashSet::with_capacity(index.refs.len()); + if index.refs.iter().any(|entry| { + !is_canonical_revision_ref_id(&entry.id) + || (entry.committed && !entry.pending_owners.is_empty()) + || (!entry.committed && entry.pending_owners.is_empty()) + || entry.pending_owners.iter().enumerate().any(|(owner_index, owner)| { + !is_canonical_revision_ref_owner(owner) + || entry.pending_owners[..owner_index].contains(owner) + }) + || !seen.insert(&entry.id) + }) { + return Err(RegistryError::Internal { + reason: "hosted revision reference index is invalid".to_string(), + }); + } + Ok(index) + } + + pub(crate) fn bodies(&self) -> impl Iterator { + self.refs.iter().map(|entry| entry.bytes.as_slice()) + } + + pub(crate) fn insert( + &mut self, + ref_id: &str, + owner: &str, + bytes: &[u8], + ) -> Result { + if let Some(entry) = self.refs.iter_mut().find(|entry| entry.id == ref_id) { + if entry.bytes != bytes { + return Err(RegistryError::Internal { + reason: "hosted revision reference body conflicts with its id".to_string(), + }); + } + if entry.committed { + return Ok(HostedRevisionRefWrite::Committed); + } + if entry.pending_owners.iter().any(|candidate| candidate == owner) { + return Ok(HostedRevisionRefWrite::AlreadyClaimed); + } + entry.pending_owners.push(owner.to_string()); + return Ok(HostedRevisionRefWrite::Claimed); + } + if self.refs.len() == MAX_HOSTED_REVISION_REFS { + return Err(RegistryError::RevisionReferenceLimit { limit: MAX_HOSTED_REVISION_REFS }); + } + self.refs.push(HostedRevisionRefIndexEntry { + id: ref_id.to_string(), + committed: false, + pending_owners: vec![owner.to_string()], + bytes: bytes.to_vec(), + }); + Ok(HostedRevisionRefWrite::Claimed) + } + + pub(crate) fn remove_if_owned(&mut self, ref_id: &str, owner: &str) -> bool { + let Some(entry_index) = self.refs.iter().position(|entry| entry.id == ref_id) else { + return false; + }; + let Some(owner_index) = + self.refs[entry_index].pending_owners.iter().position(|candidate| candidate == owner) + else { + return false; + }; + self.refs[entry_index].pending_owners.remove(owner_index); + if self.refs[entry_index].pending_owners.is_empty() { + self.refs.remove(entry_index); + } + true + } + + pub(crate) fn is_owned_by(&self, ref_id: &str, owner: &str) -> bool { + self.refs.iter().any(|entry| { + entry.id == ref_id && entry.pending_owners.iter().any(|candidate| candidate == owner) + }) + } + + pub(crate) fn commit_if_owned(&mut self, ref_id: &str, owner: &str) -> Result { + let Some(entry) = self.refs.iter_mut().find(|entry| entry.id == ref_id) else { + return Err(RegistryError::Internal { + reason: "hosted revision reference is missing during commit".to_string(), + }); + }; + if entry.committed { + return Ok(false); + } + if !entry.pending_owners.iter().any(|candidate| candidate == owner) { + return Err(RegistryError::Internal { + reason: "hosted revision reference is not owned by its committing transaction" + .to_string(), + }); + } + entry.committed = true; + entry.pending_owners.clear(); + Ok(true) + } + + pub(crate) fn to_bytes(&self) -> Vec { + serde_json::to_vec(self).expect("hosted revision reference index serializes") + } +} /// Per-process counter feeding [`unique_tmp_path`] so two concurrent /// writes to the same path don't collide on the same temp filename. @@ -187,6 +321,9 @@ pub enum CachedPackument { /// / /// package.json /// -.tgz +/// .revisions/sha512// +/// index.json +/// .json /// ``` /// /// For scoped packages the package directory is `/@scope//`. @@ -373,6 +510,40 @@ impl HostedStore { } } + async fn read_revision_refs(&self, digest: &str) -> Result>> { + match self { + HostedStore::Fs(store) => store.read_revision_refs(digest).await, + HostedStore::S3(store) => store.read_revision_refs(digest).await, + } + } + + async fn write_revision_ref( + &self, + digest: &str, + ref_id: &str, + owner: &str, + bytes: &[u8], + ) -> Result { + match self { + HostedStore::Fs(store) => store.write_revision_ref(digest, ref_id, owner, bytes).await, + HostedStore::S3(store) => store.write_revision_ref(digest, ref_id, owner, bytes).await, + } + } + + async fn remove_revision_ref(&self, digest: &str, ref_id: &str, owner: &str) -> Result<()> { + match self { + HostedStore::Fs(store) => store.remove_revision_ref(digest, ref_id, owner).await, + HostedStore::S3(store) => store.remove_revision_ref(digest, ref_id, owner).await, + } + } + + async fn commit_revision_ref(&self, digest: &str, ref_id: &str, owner: &str) -> Result<()> { + match self { + HostedStore::Fs(store) => store.commit_revision_ref(digest, ref_id, owner).await, + HostedStore::S3(store) => store.commit_revision_ref(digest, ref_id, owner).await, + } + } + /// A view rooted under `segment`, giving a hosted registry its own /// storage namespace so two orgs hosting the same `name@version` never /// collide on disk (or on object keys). @@ -435,6 +606,48 @@ impl Storage { self.hosted.list_package_names().await } + pub(crate) async fn read_hosted_revision_refs(&self, digest: &str) -> Result>> { + validate_revision_digest(digest)?; + self.hosted.read_revision_refs(digest).await + } + + pub(crate) async fn write_hosted_revision_ref( + &self, + digest: &str, + ref_id: &str, + owner: &str, + bytes: &[u8], + ) -> Result { + validate_revision_digest(digest)?; + validate_revision_ref_id(ref_id)?; + validate_revision_ref_owner(owner)?; + self.hosted.write_revision_ref(digest, ref_id, owner, bytes).await + } + + pub(crate) async fn remove_hosted_revision_ref( + &self, + digest: &str, + ref_id: &str, + owner: &str, + ) -> Result<()> { + validate_revision_digest(digest)?; + validate_revision_ref_id(ref_id)?; + validate_revision_ref_owner(owner)?; + self.hosted.remove_revision_ref(digest, ref_id, owner).await + } + + pub(crate) async fn commit_hosted_revision_ref( + &self, + digest: &str, + ref_id: &str, + owner: &str, + ) -> Result<()> { + validate_revision_digest(digest)?; + validate_revision_ref_id(ref_id)?; + validate_revision_ref_owner(owner)?; + self.hosted.commit_revision_ref(digest, ref_id, owner).await + } + /// A view whose hosted store is namespaced under `org`, so a hosted /// registry's packages live in their own storage namespace — two orgs hosting /// the same `name@version` can't collide. The disposable proxy cache is @@ -761,18 +974,22 @@ fn validated_stage_id(stage_id: &str) -> Result<&str> { #[derive(Debug, Clone)] struct Store { root: PathBuf, + revision_ref_write_lock: Arc>, } impl Store { fn new(root: PathBuf) -> Self { - Self { root } + Self { root, revision_ref_write_lock: Arc::new(tokio::sync::Mutex::new(())) } } /// A disposable store rooted at a sub-path of this one. Used to give a /// private `/~/` route its own cache namespace so its packuments /// and tarballs never collide with the public mirror or another upstream. fn namespaced(&self, prefix: &str) -> Store { - Store::new(self.root.join(prefix)) + Store { + root: self.root.join(prefix), + revision_ref_write_lock: Arc::clone(&self.revision_ref_write_lock), + } } async fn read_packument_entry( @@ -964,6 +1181,53 @@ impl Store { Ok(names) } + async fn read_revision_refs(&self, digest: &str) -> Result>> { + let index = self.read_revision_ref_index(digest).await?; + Ok(index.bodies().map(<[u8]>::to_vec).collect()) + } + + async fn write_revision_ref( + &self, + digest: &str, + ref_id: &str, + owner: &str, + bytes: &[u8], + ) -> Result { + let _guard = self.revision_ref_write_lock.lock().await; + let mut index = self.read_revision_ref_index(digest).await?; + let outcome = index.insert(ref_id, owner, bytes)?; + if outcome == HostedRevisionRefWrite::Claimed { + write_atomic(&self.revision_ref_index_path(digest), &index.to_bytes()).await?; + } + Ok(outcome) + } + + async fn remove_revision_ref(&self, digest: &str, ref_id: &str, owner: &str) -> Result<()> { + let _guard = self.revision_ref_write_lock.lock().await; + let mut index = self.read_revision_ref_index(digest).await?; + if index.remove_if_owned(ref_id, owner) { + write_atomic(&self.revision_ref_index_path(digest), &index.to_bytes()).await?; + } + Ok(()) + } + + async fn commit_revision_ref(&self, digest: &str, ref_id: &str, owner: &str) -> Result<()> { + let _guard = self.revision_ref_write_lock.lock().await; + let mut index = self.read_revision_ref_index(digest).await?; + if index.commit_if_owned(ref_id, owner)? { + write_atomic(&self.revision_ref_index_path(digest), &index.to_bytes()).await?; + } + Ok(()) + } + + async fn read_revision_ref_index(&self, digest: &str) -> Result { + match fs::read(self.revision_ref_index_path(digest)).await { + Ok(bytes) => HostedRevisionRefIndex::from_bytes(&bytes), + Err(err) if err.kind() == ErrorKind::NotFound => Ok(HostedRevisionRefIndex::default()), + Err(err) => Err(err.into()), + } + } + fn package_dir(&self, name: &PackageName) -> PathBuf { self.root.join(name.as_str()) } @@ -980,6 +1244,14 @@ impl Store { self.root.join(".revisions").join("sha512").join(digest) } + fn revision_refs_dir(&self, digest: &str) -> PathBuf { + self.root.join(HOSTED_REVISION_REFS_DIR).join(digest) + } + + fn revision_ref_index_path(&self, digest: &str) -> PathBuf { + self.revision_refs_dir(digest).join(HOSTED_REVISION_REF_INDEX_FILE) + } + async fn read_staged(&self, object: &str) -> Result>> { match fs::read(self.root.join(STAGED_DIR).join(object)).await { Ok(bytes) => Ok(Some(bytes)), @@ -1017,6 +1289,32 @@ impl Store { } } +pub(crate) fn is_canonical_revision_ref_id(ref_id: &str) -> bool { + ref_id.len() == 64 && ref_id.bytes().all(|byte| byte.is_ascii_hexdigit()) +} + +pub(crate) fn is_canonical_revision_ref_owner(owner: &str) -> bool { + !owner.is_empty() + && owner.len() <= 64 + && owner.bytes().all(|byte| byte.is_ascii_alphanumeric() || byte == b'-') +} + +fn validate_revision_ref_id(ref_id: &str) -> Result<()> { + if is_canonical_revision_ref_id(ref_id) { + Ok(()) + } else { + Err(RegistryError::BadRequest { reason: "invalid revision reference id".to_string() }) + } +} + +fn validate_revision_ref_owner(owner: &str) -> Result<()> { + if is_canonical_revision_ref_owner(owner) { + Ok(()) + } else { + Err(RegistryError::BadRequest { reason: "invalid revision reference owner".to_string() }) + } +} + /// The stage id of a metadata object name, or `None` for anything else in /// the staged namespace (bodies, tmp files from interrupted writes). pub(crate) fn staged_id_of_meta_object(object: &str) -> Option<&str> { diff --git a/pnpr/crates/pnpr/src/storage/tests.rs b/pnpr/crates/pnpr/src/storage/tests.rs index 3cefd8490f..72265fab5c 100644 --- a/pnpr/crates/pnpr/src/storage/tests.rs +++ b/pnpr/crates/pnpr/src/storage/tests.rs @@ -1,6 +1,6 @@ use super::{ - AsyncWriteExt, ErrorKind, HostedStoreConfig, PackageName, RegistryError, Storage, TarballWrite, - create_tmp_file_with, fs, + AsyncWriteExt, ErrorKind, HostedRevisionRefWrite, HostedStoreConfig, MAX_HOSTED_REVISION_REFS, + PackageName, RegistryError, Storage, TarballWrite, create_tmp_file_with, fs, }; use tempfile::TempDir; @@ -20,6 +20,118 @@ fn packument_write_conflict_delay_caps_growth() { assert_eq!(super::packument_write_conflict_delay(32).as_millis(), 250); } +#[tokio::test] +async fn hosted_revision_refs_roundtrip_in_the_org_namespace() { + let tmp = TempDir::new().unwrap(); + let storage = storage_in(&tmp).for_hosted("acme"); + let digest = "A".repeat(86); + let ref_id = "b".repeat(64); + + assert_eq!(storage.read_hosted_revision_refs(&digest).await.unwrap(), Vec::>::new()); + storage + .write_hosted_revision_ref( + &digest, + &ref_id, + "owner-a", + br#"{"package":"foo","version":"1.0.0"}"#, + ) + .await + .unwrap(); + assert_eq!( + storage.read_hosted_revision_refs(&digest).await.unwrap(), + vec![br#"{"package":"foo","version":"1.0.0"}"#.to_vec()], + ); + assert_eq!( + storage_in(&tmp).read_hosted_revision_refs(&digest).await.unwrap(), + Vec::>::new(), + ); +} + +#[tokio::test] +async fn hosted_revision_ref_paths_reject_noncanonical_segments() { + let tmp = TempDir::new().unwrap(); + let storage = storage_in(&tmp); + let digest = "A".repeat(86); + + let invalid_digest = storage.read_hosted_revision_refs("../escape").await; + assert!(invalid_digest.is_err()); + let invalid_ref = + storage.write_hosted_revision_ref(&digest, "../escape", "owner-a", b"{}").await; + assert!(invalid_ref.is_err()); +} + +#[tokio::test] +async fn hosted_revision_ref_writes_enforce_the_read_bound() { + let tmp = TempDir::new().unwrap(); + let storage = storage_in(&tmp); + let digest = "A".repeat(86); + for index in 0..MAX_HOSTED_REVISION_REFS { + storage + .write_hosted_revision_ref(&digest, &format!("{index:064x}"), "owner-a", b"{}") + .await + .unwrap(); + } + + let overflow = MAX_HOSTED_REVISION_REFS; + let err = storage + .write_hosted_revision_ref(&digest, &format!("{overflow:064x}"), "owner-a", b"{}") + .await + .unwrap_err(); + assert_eq!(err.status_code(), axum::http::StatusCode::CONFLICT); + assert!(matches!( + err, + RegistryError::RevisionReferenceLimit { limit } if limit == MAX_HOSTED_REVISION_REFS + )); + + assert_eq!( + storage + .write_hosted_revision_ref(&digest, &"0".repeat(64), "owner-a", b"{}") + .await + .unwrap(), + HostedRevisionRefWrite::AlreadyClaimed, + ); + let stray_dir = tmp.path().join("storage/.revisions/sha512").join(&digest); + fs::write(stray_dir.join("not-a-reference.json"), b"stray").await.unwrap(); + fs::write(stray_dir.join("interrupted.tmp"), b"stray").await.unwrap(); + let refs = storage.read_hosted_revision_refs(&digest).await.unwrap(); + assert_eq!(refs.len(), MAX_HOSTED_REVISION_REFS); + assert!(refs.iter().all(|bytes| bytes == b"{}")); +} + +#[tokio::test] +async fn concurrent_hosted_revision_ref_writes_cannot_exceed_the_limit() { + let tmp = TempDir::new().unwrap(); + let storage = storage_in(&tmp); + let digest = "A".repeat(86); + let mut writes = Vec::new(); + for index in 0..MAX_HOSTED_REVISION_REFS * 2 { + let storage = storage.clone(); + let digest = digest.clone(); + writes.push(tokio::spawn(async move { + storage + .write_hosted_revision_ref(&digest, &format!("{index:064x}"), "owner-a", b"{}") + .await + })); + } + + let mut written = 0; + let mut rejected = 0; + for write in writes { + match write.await.unwrap() { + Ok(HostedRevisionRefWrite::Claimed) => written += 1, + Ok(outcome) => panic!("unexpected revision-reference write outcome: {outcome:?}"), + Err(RegistryError::RevisionReferenceLimit { .. }) => rejected += 1, + Err(err) => panic!("unexpected revision-reference write error: {err}"), + } + } + assert_eq!(written, MAX_HOSTED_REVISION_REFS); + assert_eq!(rejected, MAX_HOSTED_REVISION_REFS); + assert_eq!( + storage.read_hosted_revision_refs(&digest).await.unwrap().len(), + MAX_HOSTED_REVISION_REFS, + ); +} + #[tokio::test] async fn hosted_tarball_under_non_directory_package_path_is_an_error() { let tmp = TempDir::new().unwrap(); diff --git a/pnpr/crates/pnpr/tests/server.rs b/pnpr/crates/pnpr/tests/server.rs index ec8cf59d41..0fa3cd9467 100644 --- a/pnpr/crates/pnpr/tests/server.rs +++ b/pnpr/crates/pnpr/tests/server.rs @@ -351,6 +351,37 @@ fn sha512_integrity(bytes: &[u8]) -> String { opts.result().to_string() } +fn hosted_publish_request( + url: &str, + package: &str, + version: &str, + tarball: &[u8], + token: &str, +) -> Request { + use base64::{Engine as _, engine::general_purpose::STANDARD as BASE64}; + + let basename = package.rsplit('/').next().unwrap(); + let attachment = format!("{package}-{version}.tgz"); + let body = json!({ + "name": package, + "dist-tags": { "latest": version }, + "versions": { (version): { "name": package, "version": version, "dist": { + "tarball": format!("http://example.test/{package}/-/{basename}-{version}.tgz"), + "integrity": sha512_integrity(tarball), + } } }, + "_attachments": { (attachment): { + "content_type": "application/octet-stream", + "data": BASE64.encode(tarball), + "length": tarball.len(), + } }, + }); + Request::put(url) + .header("content-type", "application/json") + .header(header::AUTHORIZATION, format!("Bearer {token}")) + .body(Body::from(serde_json::to_vec(&body).unwrap())) + .unwrap() +} + /// The 40-char hex SHA-1 the way pre-2017 npm publishes carry it in the /// legacy `dist.shasum` field. fn sha1_hex_of(bytes: &[u8]) -> String { @@ -3350,6 +3381,193 @@ async fn publish_to_hosted_round_trips_in_its_own_namespace() { assert_eq!(rejected.status(), StatusCode::BAD_REQUEST); } +#[tokio::test] +async fn hosted_original_is_served_by_digest_after_restart_and_through_a_router() { + let tmp = TempDir::new().unwrap(); + let mut config = config_for("http://127.0.0.1:1", tmp.path().to_path_buf()); + config.hosted.insert("acme".to_string(), hosted_with_access("acme", "$all")); + let graph = vec![ + ( + "acme".to_string(), + Registry::Hosted { patterns: vec![PackagePattern::parse("@acme/*").unwrap()] }, + ), + ("main".to_string(), Registry::Router { sources: vec!["acme".to_string()] }), + ]; + config.registries = Registries::new(graph.into_iter().collect(), Some("main".to_string())); + let auth = AuthState::in_memory(); + let token = auth.tokens.issue("alice").await.unwrap(); + let tarball = b"public-hosted-original"; + let integrity_text = sha512_integrity(tarball); + let integrity = integrity_text.parse().unwrap(); + let revision_path = integrity_addressed_tarball_path(&integrity).unwrap(); + let digest = revision_path.rsplit('/').next().unwrap(); + + let publish = router_with_auth(config.clone(), auth.clone()) + .oneshot(hosted_publish_request( + "/~acme/@acme/widget", + "@acme/widget", + "1.0.0", + tarball, + &token, + )) + .await + .unwrap(); + assert_eq!(publish.status(), StatusCode::CREATED); + + let app = router_with_auth(config, auth); + for path in [ + format!("/~acme/{revision_path}"), + format!("/~main/{revision_path}"), + format!("/{revision_path}"), + ] { + let response = + app.clone().oneshot(Request::get(&path).body(Body::empty()).unwrap()).await.unwrap(); + assert_eq!(response.status(), StatusCode::OK, "{path}"); + assert_eq!(response.headers().get(header::CACHE_CONTROL).unwrap(), "private, no-store"); + assert_eq!(response.headers().get(header::VARY).unwrap(), "Authorization"); + assert_eq!( + response.headers().get(header::ETAG).unwrap().to_str().unwrap(), + format!(r#""{digest}""#), + ); + assert_eq!( + response.headers().get("content-digest").unwrap().to_str().unwrap(), + format!("sha-512=:{}:", integrity_text.strip_prefix("sha512-").unwrap()), + ); + assert_eq!(body_bytes(response.into_body()).await, tarball, "{path}"); + } + + let canonical = app + .oneshot( + Request::get("/~acme/@acme/widget/-/widget-1.0.0.tgz").body(Body::empty()).unwrap(), + ) + .await + .unwrap(); + assert_eq!(canonical.status(), StatusCode::OK); + assert_eq!(body_bytes(canonical.into_body()).await, tarball); +} + +#[tokio::test] +async fn hosted_publish_rejects_digest_reference_overflow_without_disabling_existing_refs() { + let tmp = TempDir::new().unwrap(); + let mut config = config_for("http://127.0.0.1:1", tmp.path().to_path_buf()); + config.hosted.insert("acme".to_string(), hosted_with_access("acme", "$authenticated")); + config.registries = Registries::new( + vec![("acme".to_string(), Registry::Hosted { patterns: vec![] })].into_iter().collect(), + Some("acme".to_string()), + ); + let auth = AuthState::in_memory(); + let token = auth.tokens.issue("alice").await.unwrap(); + let tarball = b"shared-hosted-original"; + let integrity = sha512_integrity(tarball).parse().unwrap(); + let revision_path = integrity_addressed_tarball_path(&integrity).unwrap(); + let app = router_with_auth(config, auth); + + for index in 0..32 { + let package = format!("shared-artifact-{index}"); + let response = app + .clone() + .oneshot(hosted_publish_request( + &format!("/~acme/{package}"), + &package, + "1.0.0", + tarball, + &token, + )) + .await + .unwrap(); + assert_eq!(response.status(), StatusCode::CREATED, "{package}"); + } + + let rejected = app + .clone() + .oneshot(hosted_publish_request( + "/~acme/shared-artifact-overflow", + "shared-artifact-overflow", + "1.0.0", + tarball, + &token, + )) + .await + .unwrap(); + assert_eq!(rejected.status(), StatusCode::CONFLICT); + + let overflow_packument = app + .clone() + .oneshot( + Request::get("/~acme/shared-artifact-overflow") + .header(header::AUTHORIZATION, format!("Bearer {token}")) + .body(Body::empty()) + .unwrap(), + ) + .await + .unwrap(); + assert_eq!(overflow_packument.status(), StatusCode::OK); + let overflow_packument = body_json(overflow_packument.into_body()).await; + assert!(overflow_packument["versions"].get("1.0.0").is_none()); + + let existing = app + .oneshot( + Request::get(format!("/~acme/{revision_path}")) + .header(header::AUTHORIZATION, format!("Bearer {token}")) + .body(Body::empty()) + .unwrap(), + ) + .await + .unwrap(); + assert_eq!(existing.status(), StatusCode::OK); + assert_eq!(body_bytes(existing.into_body()).await, tarball); +} + +#[tokio::test] +async fn hosted_digest_route_rechecks_package_access() { + let tmp = TempDir::new().unwrap(); + let mut config = config_for("http://127.0.0.1:1", tmp.path().to_path_buf()); + config.hosted.insert("corp".to_string(), hosted_with_access("corp", "alice")); + config.registries = Registries::new( + vec![("corp".to_string(), Registry::Hosted { patterns: vec![] })].into_iter().collect(), + Some("corp".to_string()), + ); + let auth = AuthState::in_memory(); + let alice = auth.tokens.issue("alice").await.unwrap(); + let bob = auth.tokens.issue("bob").await.unwrap(); + let tarball = b"private-hosted-original"; + let integrity = sha512_integrity(tarball).parse().unwrap(); + let revision_path = integrity_addressed_tarball_path(&integrity).unwrap(); + let app = router_with_auth(config, auth); + + let publish = app + .clone() + .oneshot(hosted_publish_request("/~corp/secret", "secret", "1.0.0", tarball, &alice)) + .await + .unwrap(); + assert_eq!(publish.status(), StatusCode::CREATED); + + for authorization in [None, Some(format!("Bearer {bob}"))] { + let mut request = Request::get(format!("/~corp/{revision_path}")); + if let Some(authorization) = authorization { + request = request.header(header::AUTHORIZATION, authorization); + } + let response = app.clone().oneshot(request.body(Body::empty()).unwrap()).await.unwrap(); + assert_eq!(response.status(), StatusCode::NOT_FOUND); + assert_eq!(response.headers().get(header::CACHE_CONTROL).unwrap(), "private, no-store"); + assert_eq!(response.headers().get(header::VARY).unwrap(), "Authorization"); + } + + let response = app + .oneshot( + Request::get(format!("/~corp/{revision_path}")) + .header(header::AUTHORIZATION, format!("Bearer {alice}")) + .body(Body::empty()) + .unwrap(), + ) + .await + .unwrap(); + assert_eq!(response.status(), StatusCode::OK); + assert_eq!(response.headers().get(header::CACHE_CONTROL).unwrap(), "private, no-store"); + assert_eq!(response.headers().get(header::VARY).unwrap(), "Authorization"); + assert_eq!(body_bytes(response.into_body()).await, tarball); +} + /// A hosted registry's declared `patterns:` are enforced on the registry itself, on /// every path to it: an off-pattern publish is rejected and an off-pattern /// read is a definitive 404 — through a router and at the registry's own