Skip to content

Commit fada6f8

Browse files
committed
perf(pm): execute resolver manifest provider jobs
1 parent 70806c8 commit fada6f8

3 files changed

Lines changed: 269 additions & 0 deletions

File tree

crates/ruborist/src/service/provider.rs

Lines changed: 7 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -7,6 +7,7 @@
77
use std::sync::Arc;
88

99
use async_trait::async_trait;
10+
use bytes::Bytes;
1011

1112
use super::cache::VersionsInfo;
1213
use super::manifest::MetadataFormat;
@@ -80,3 +81,9 @@ pub trait ManifestProvider: RegistryClient + Clone + Send + Sync + 'static {
8081
/// de-duplication stay in the BFS loop.
8182
async fn execute_manifest_job(&self, job: ManifestJob) -> Result<ManifestJobDone, Self::Error>;
8283
}
84+
85+
/// Raw full-manifest bytes fetched by a provider before parsing.
86+
pub(crate) enum ProviderFullManifestBytes {
87+
Fresh { bytes: Bytes, etag: Option<String> },
88+
NotModified { versions: Arc<VersionsInfo> },
89+
}

crates/ruborist/src/service/registry.rs

Lines changed: 255 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -21,6 +21,7 @@
2121
use std::sync::Arc;
2222

2323
use anyhow::anyhow;
24+
use async_trait::async_trait;
2425

2526
/// Get current timestamp in seconds since UNIX epoch.
2627
/// Works on both native and WASM targets.
@@ -43,6 +44,9 @@ use dashmap::DashSet;
4344

4445
use super::cache::{PackageCache, Versions, VersionsInfo};
4546
use super::manifest;
47+
use super::provider::{
48+
ManifestFullData, ManifestJob, ManifestJobDone, ManifestProvider, ProviderFullManifestBytes,
49+
};
4650
use super::store::{ManifestStore, NoopStore};
4751
use crate::model::manifest::{CoreVersionManifest, FullManifest, extract_core_version_off_runtime};
4852
use crate::resolver::semver::normalize_spec;
@@ -199,6 +203,78 @@ enum FullManifestResult {
199203
NotModified,
200204
}
201205

206+
#[cfg_attr(target_arch = "wasm32", async_trait(?Send))]
207+
#[cfg_attr(not(target_arch = "wasm32"), async_trait)]
208+
impl ManifestProvider for UnifiedRegistry {
209+
async fn execute_manifest_job(&self, job: ManifestJob) -> Result<ManifestJobDone, Self::Error> {
210+
match job {
211+
ManifestJob::Full { name, spec } => {
212+
let data = match self.fetch_full_manifest_job(&name).await? {
213+
ProviderFullManifestBytes::Fresh { bytes, etag } => {
214+
let (manifest, speculative) =
215+
manifest::parse_full_manifest_with_core_off_runtime(bytes, spec)
216+
.await?;
217+
let manifest = Arc::new(manifest);
218+
let speculative = speculative.map(|(spec, core)| {
219+
let core = Arc::new(core);
220+
self.store_version_manifest(&name, Arc::clone(&core));
221+
(spec, core)
222+
});
223+
let versions = Arc::new(VersionsInfo {
224+
versions: Versions {
225+
version_list: manifest.versions.clone(),
226+
dist_tags: manifest.dist_tags.clone(),
227+
},
228+
etag,
229+
last_updated: current_timestamp_secs(),
230+
});
231+
self.store.store_versions(&name, versions);
232+
ManifestFullData::Full {
233+
manifest,
234+
speculative,
235+
}
236+
}
237+
ProviderFullManifestBytes::NotModified { versions } => {
238+
ManifestFullData::Versions(versions)
239+
}
240+
};
241+
242+
Ok(ManifestJobDone::Full { name, data })
243+
}
244+
ManifestJob::Version {
245+
name,
246+
spec,
247+
fetch_spec,
248+
format,
249+
} => {
250+
let manifest = self
251+
.fetch_version_job_manifest(&name, &spec, &fetch_spec, format)
252+
.await?;
253+
Ok(ManifestJobDone::Version {
254+
name,
255+
spec,
256+
manifest,
257+
})
258+
}
259+
ManifestJob::ExtractVersion {
260+
name,
261+
spec,
262+
version,
263+
full,
264+
} => {
265+
let manifest = self
266+
.extract_version_job_manifest(&name, &spec, version, full)
267+
.await?;
268+
Ok(ManifestJobDone::Version {
269+
name,
270+
spec,
271+
manifest,
272+
})
273+
}
274+
}
275+
}
276+
}
277+
202278
impl UnifiedRegistry {
203279
/// Create a builder for `UnifiedRegistry`.
204280
pub fn builder() -> UnifiedRegistryBuilder {
@@ -220,6 +296,87 @@ impl UnifiedRegistry {
220296
&self.cache
221297
}
222298

299+
fn store_version_manifest(&self, name: &str, manifest: Arc<CoreVersionManifest>) {
300+
let version = manifest.version.clone();
301+
self.store.store_version_manifest(name, &version, manifest);
302+
}
303+
304+
async fn fetch_full_manifest_job(
305+
&self,
306+
name: &str,
307+
) -> Result<ProviderFullManifestBytes, RegistryError> {
308+
let store_versions = self.store.load_versions(name).await.map(Arc::new);
309+
let etag = store_versions.as_ref().and_then(|v| v.etag.clone());
310+
311+
match manifest::fetch_full_manifest_bytes(manifest::FetchManifestOptions {
312+
registry_url: &self.registry_url,
313+
name,
314+
format: manifest::MetadataFormat::Abbreviated,
315+
etag: etag.as_deref(),
316+
})
317+
.await
318+
.map_err(RegistryError)?
319+
{
320+
manifest::FetchManifestBytesResult::Ok(bytes, etag) => {
321+
Ok(ProviderFullManifestBytes::Fresh { bytes, etag })
322+
}
323+
manifest::FetchManifestBytesResult::NotModified => {
324+
let versions = store_versions.ok_or_else(|| {
325+
RegistryError(anyhow!(
326+
"304 Not Modified without cached versions for {name}"
327+
))
328+
})?;
329+
Ok(ProviderFullManifestBytes::NotModified { versions })
330+
}
331+
}
332+
}
333+
334+
async fn extract_version_job_manifest(
335+
&self,
336+
name: &str,
337+
_spec: &str,
338+
version: String,
339+
full: Arc<FullManifest>,
340+
) -> Result<Arc<CoreVersionManifest>, RegistryError> {
341+
let (resolved_version, manifest) = extract_core_version_off_runtime(full, version).await;
342+
let manifest = manifest.ok_or_else(|| {
343+
RegistryError(anyhow!(
344+
"Version {} not found in manifest for {}",
345+
resolved_version,
346+
name
347+
))
348+
})?;
349+
self.store_version_manifest(name, Arc::clone(&manifest));
350+
Ok(manifest)
351+
}
352+
353+
async fn fetch_version_job_manifest(
354+
&self,
355+
name: &str,
356+
_spec: &str,
357+
fetch_spec: &str,
358+
format: manifest::MetadataFormat,
359+
) -> Result<Arc<CoreVersionManifest>, RegistryError> {
360+
if deno_semver::Version::parse_from_npm(fetch_spec).is_ok()
361+
&& let Some(manifest) = self.store.load_version_manifest(name, fetch_spec).await
362+
{
363+
return Ok(Arc::new(manifest));
364+
}
365+
366+
let manifest = Arc::new(
367+
manifest::fetch_version_manifest(manifest::FetchVersionManifestOptions {
368+
registry_url: &self.registry_url,
369+
name,
370+
spec: fetch_spec,
371+
format,
372+
})
373+
.await
374+
.map_err(RegistryError)?,
375+
);
376+
self.store_version_manifest(name, Arc::clone(&manifest));
377+
Ok(manifest)
378+
}
379+
223380
/// Resolve full manifest through memory → store → network with ETag validation.
224381
///
225382
/// Single-flight cache flow:
@@ -509,6 +666,10 @@ impl RegistryClient for UnifiedRegistry {
509666
self.supports_semver
510667
}
511668

669+
fn registry_url(&self) -> &str {
670+
&self.registry_url
671+
}
672+
512673
fn cache_version_manifest(&self, name: &str, spec: &str, manifest: Arc<CoreVersionManifest>) {
513674
self.cache
514675
.set_version_manifest(name.to_string(), spec.to_string(), manifest);
@@ -572,7 +733,44 @@ impl RegistryClient for UnifiedRegistry {
572733

573734
#[cfg(test)]
574735
mod tests {
736+
use std::sync::Mutex;
737+
575738
use super::*;
739+
use crate::service::{ManifestJob, ManifestJobDone, ManifestProvider, ManifestStore};
740+
741+
#[derive(Default)]
742+
struct RecordingStore {
743+
stored_versions: Mutex<Vec<(String, String)>>,
744+
}
745+
746+
#[async_trait]
747+
impl ManifestStore for RecordingStore {
748+
async fn load_versions(&self, _name: &str) -> Option<VersionsInfo> {
749+
None
750+
}
751+
752+
async fn load_version_manifest(
753+
&self,
754+
_name: &str,
755+
_version: &str,
756+
) -> Option<CoreVersionManifest> {
757+
None
758+
}
759+
760+
fn store_versions(&self, _name: &str, _info: Arc<VersionsInfo>) {}
761+
762+
fn store_version_manifest(
763+
&self,
764+
name: &str,
765+
version: &str,
766+
_manifest: Arc<CoreVersionManifest>,
767+
) {
768+
self.stored_versions
769+
.lock()
770+
.unwrap()
771+
.push((name.to_string(), version.to_string()));
772+
}
773+
}
576774

577775
#[test]
578776
fn test_is_npm_registry() {
@@ -642,4 +840,61 @@ mod tests {
642840
// Both registries share the same cache
643841
assert!(Arc::ptr_eq(&registry1.cache, &registry2.cache));
644842
}
843+
844+
#[tokio::test]
845+
async fn test_unified_registry_executes_extract_manifest_provider_job() {
846+
let store = Arc::new(RecordingStore::default());
847+
let registry = UnifiedRegistry::builder()
848+
.registry("https://registry.npmmirror.com")
849+
.store(store.clone())
850+
.build();
851+
852+
let (full, _) = manifest::parse_full_manifest_with_core_off_runtime(
853+
bytes::Bytes::from_static(
854+
br#"{
855+
"name":"provider-extract-demo",
856+
"dist-tags":{"latest":"1.0.0"},
857+
"versions":{
858+
"1.0.0":{
859+
"name":"provider-extract-demo",
860+
"version":"1.0.0",
861+
"dist":{"tarball":"https://registry.example/demo-1.0.0.tgz"}
862+
}
863+
}
864+
}"#,
865+
),
866+
None,
867+
)
868+
.await
869+
.unwrap();
870+
871+
let done = ManifestProvider::execute_manifest_job(
872+
&registry,
873+
ManifestJob::ExtractVersion {
874+
name: "provider-extract-demo".to_string(),
875+
spec: "latest".to_string(),
876+
version: "1.0.0".to_string(),
877+
full: Arc::new(full),
878+
},
879+
)
880+
.await
881+
.unwrap();
882+
883+
let ManifestJobDone::Version {
884+
spec,
885+
manifest: returned,
886+
..
887+
} = done
888+
else {
889+
panic!("expected version manifest job");
890+
};
891+
892+
assert_eq!(spec, "latest");
893+
assert_eq!(returned.name, "provider-extract-demo");
894+
assert_eq!(returned.version, "1.0.0");
895+
assert_eq!(
896+
store.stored_versions.lock().unwrap().as_slice(),
897+
&[("provider-extract-demo".to_string(), "1.0.0".to_string())]
898+
);
899+
}
645900
}

crates/ruborist/src/traits/registry.rs

Lines changed: 7 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -133,6 +133,13 @@ pub trait RegistryClient {
133133
false
134134
}
135135

136+
/// Base registry URL used by schedulers that need to classify raw work.
137+
///
138+
/// Implementations without a concrete URL can keep the default.
139+
fn registry_url(&self) -> &str {
140+
""
141+
}
142+
136143
/// Fetch full package manifest from registry.
137144
///
138145
/// Returns the complete package manifest with all versions, wrapped in

0 commit comments

Comments
 (0)