use crate::progress::ProgressSender;
use deps_core::ConcreteVersion;
use deps_core::Deprecation;
use deps_core::FetchFailure;
use deps_core::PackageName;
use deps_core::PackageVersions;
use deps_core::Registry;
use deps_core::RemovalStatus;
use deps_core::VersionReq;
use std::collections::{HashMap, HashSet};
use std::sync::Arc;
use std::time::Duration;
pub type DepSources = Vec<(PackageName, deps_core::parser::DependencySource)>;
pub fn dedup_dependencies_by_source(
parse_result: &dyn deps_core::ParseResult,
formatter: &dyn deps_core::lsp_helpers::EcosystemFormatter,
) -> (
HashMap<PackageName, deps_core::parser::DependencySource>,
HashSet<PackageName>,
) {
use std::collections::hash_map::Entry;
let mut by_name: HashMap<PackageName, deps_core::parser::DependencySource> = HashMap::new();
let mut collided: HashSet<PackageName> = HashSet::new();
for dep in parse_result
.dependencies()
.into_iter()
.filter(|dep| formatter.can_resolve_source(&dep.source()))
{
let name = dep.name().clone();
let source = dep.source();
match by_name.entry(name.clone()) {
Entry::Vacant(entry) => {
entry.insert(source);
}
Entry::Occupied(entry) => {
if *entry.get() != source && collided.insert(name.clone()) {
tracing::warn!(
package = %name.for_tracing(),
source_a = ?entry.get(),
source_b = ?source,
"dependency declared against two different resolved registries; \
skipping version resolution for all occurrences"
);
}
}
}
}
for name in &collided {
by_name.remove(name);
}
(by_name, collided)
}
#[cfg(feature = "composer")]
pub fn composer_minimum_stability(parse_result: &dyn deps_core::ParseResult) -> Option<String> {
parse_result
.as_any()
.downcast_ref::<crate::setup::ComposerParseResult>()
.and_then(|r| r.minimum_stability.clone())
}
#[cfg(not(feature = "composer"))]
pub fn composer_minimum_stability(_parse_result: &dyn deps_core::ParseResult) -> Option<String> {
None
}
#[non_exhaustive]
pub struct FetchResult {
pub versions: HashMap<PackageName, PackageVersions>,
pub yanked_versions: HashMap<PackageName, (ConcreteVersion, RemovalStatus)>,
pub deprecations: HashMap<PackageName, Deprecation>,
pub fetch_failed: HashMap<PackageName, FetchFailure>,
pub no_comparable_versions: HashSet<PackageName>,
pub failed_count: usize,
pub first_error: Option<String>,
pub licenses: HashMap<PackageName, Vec<String>>,
}
impl FetchResult {
#[must_use]
#[allow(clippy::too_many_arguments)]
pub fn new(
versions: HashMap<PackageName, PackageVersions>,
yanked_versions: HashMap<PackageName, (ConcreteVersion, RemovalStatus)>,
deprecations: HashMap<PackageName, Deprecation>,
fetch_failed: HashMap<PackageName, FetchFailure>,
no_comparable_versions: HashSet<PackageName>,
failed_count: usize,
first_error: Option<String>,
licenses: HashMap<PackageName, Vec<String>>,
) -> Self {
Self {
versions,
yanked_versions,
deprecations,
fetch_failed,
no_comparable_versions,
failed_count,
first_error,
licenses,
}
}
}
#[allow(
clippy::too_many_arguments,
reason = "internal (non-pub) call-site-controlled fetch tuning + ecosystem-context \
parameters; grouping into a config struct would only move, not reduce, the \
per-call-site churn across this module's ~15 production and test call sites"
)]
pub async fn fetch_latest_versions_parallel(
registry: Arc<dyn Registry>,
package_sources: DepSources,
in_use: &HashMap<PackageName, Vec<String>>,
progress_sender: Option<ProgressSender>,
freshness: deps_core::freshness::FreshnessSettings,
timeout_secs: u64,
max_concurrent: usize,
minimum_stability: Option<&str>,
) -> FetchResult {
use futures::stream::{self, StreamExt};
use std::time::Duration;
let fetched = Arc::new(std::sync::atomic::AtomicUsize::new(0));
let failed = Arc::new(std::sync::atomic::AtomicUsize::new(0));
let first_error: Arc<std::sync::Mutex<Option<String>>> = Arc::new(std::sync::Mutex::new(None));
let timeout = Duration::from_secs(timeout_secs);
let wildcard_req = deps_core::VersionReq::new("*");
let check_yanked = registry.reports_yanked();
let results: Vec<_> = stream::iter(package_sources)
.map(|(name, source)| {
let registry = Arc::clone(®istry);
let fetched = Arc::clone(&fetched);
let failed = Arc::clone(&failed);
let first_error = Arc::clone(&first_error);
let progress_sender = progress_sender.clone();
let wildcard_req = &wildcard_req;
let in_use_versions = in_use.get(&name).cloned().unwrap_or_default();
async move {
fetch_and_classify_package(
registry.as_ref(),
name,
source,
in_use_versions,
wildcard_req,
freshness,
timeout,
minimum_stability,
check_yanked,
&fetched,
&failed,
&first_error,
progress_sender.as_ref(),
)
.await
}
})
.buffer_unordered(max_concurrent.max(1))
.collect()
.await;
let mut versions = HashMap::with_capacity(results.len());
let mut yanked_versions = HashMap::new();
let mut fetch_failed = HashMap::new();
let mut deprecations = HashMap::new();
let mut no_comparable_versions = HashSet::new();
let mut licenses = HashMap::new();
let mut priority_error: Option<String> = None;
for (version, yanked, failed_name, deprecation, no_comparable_versions_name, license) in results
{
if let Some((name, v)) = version {
versions.insert(name, v);
}
if let Some((name, v, status)) = yanked {
yanked_versions.insert(name, (v, status));
}
if let Some((name, failure, message)) = failed_name {
fetch_failed.insert(name, failure);
if priority_error.is_none() {
priority_error = Some(message);
}
}
if let Some((name, d)) = deprecation {
deprecations.insert(name, d);
}
if let Some(name) = no_comparable_versions_name {
no_comparable_versions.insert(name);
}
if let Some((name, license)) = license {
licenses.insert(name, license);
}
}
let error_message =
priority_error.or_else(|| first_error.lock().unwrap_or_else(|p| p.into_inner()).take());
FetchResult {
versions,
yanked_versions,
fetch_failed,
deprecations,
no_comparable_versions,
failed_count: failed.load(std::sync::atomic::Ordering::Relaxed),
first_error: error_message,
licenses,
}
}
type PackageFetchOutcome = (
Option<(PackageName, PackageVersions)>,
Option<(PackageName, ConcreteVersion, RemovalStatus)>,
Option<(PackageName, FetchFailure, String)>,
Option<(PackageName, Deprecation)>,
Option<PackageName>,
Option<(PackageName, Vec<String>)>,
);
#[allow(
clippy::too_many_arguments,
reason = "mirrors the per-package async closure this was extracted from — every \
parameter is either call-site fetch tuning already threaded through \
fetch_latest_versions_parallel or a counter/sender shared across the \
whole stream; grouping into a struct would only move, not reduce, churn"
)]
async fn fetch_and_classify_package(
registry: &dyn Registry,
name: PackageName,
source: deps_core::parser::DependencySource,
in_use_versions: Vec<String>,
wildcard_req: &VersionReq,
freshness: deps_core::freshness::FreshnessSettings,
timeout: Duration,
minimum_stability: Option<&str>,
check_yanked: bool,
fetched: &std::sync::atomic::AtomicUsize,
failed: &std::sync::atomic::AtomicUsize,
first_error: &std::sync::Mutex<Option<String>>,
progress_sender: Option<&ProgressSender>,
) -> PackageFetchOutcome {
let result = tokio::time::timeout(
timeout,
registry.get_versions_from(&name, &source, freshness),
)
.await;
let mut yanked: Option<(PackageName, ConcreteVersion, RemovalStatus)> = None;
let mut failed_name: Option<(PackageName, FetchFailure, String)> = None;
let mut deprecation: Option<(PackageName, Deprecation)> = None;
let mut license: Option<(PackageName, Vec<String>)> = None;
let mut no_comparable_versions = false;
let version = match result {
Ok(Ok(versions)) => {
let available: Arc<[ConcreteVersion]> = versions
.iter()
.map(|v| v.version_string().clone())
.collect();
let yanked_list: Arc<[(ConcreteVersion, RemovalStatus)]> = if check_yanked {
versions
.iter()
.filter_map(|v| {
let status = v.removal_status();
status
.is_flagged()
.then(|| (v.version_string().clone(), status))
})
.collect()
} else {
Arc::from([])
};
let resolved = if let Some(v) = registry
.select_latest_matching_with_context(&versions, wildcard_req, minimum_stability)
.and_then(|idx| versions.get(idx))
{
let latest = v.version_string().clone();
tracing::debug!(package = %name.for_tracing(), version = %latest, "fetched");
Some((
latest,
v.removal_status(),
v.published_at(),
v.deprecation().cloned(),
v.license().to_vec(),
))
} else {
let fallback = tokio::time::timeout(
timeout,
registry.get_latest_matching_from(
&name,
&source,
wildcard_req,
minimum_stability,
),
)
.await;
match fallback {
Ok(Ok(Some(v))) => {
let latest = v.version_string().clone();
tracing::debug!(
package = %name.for_tracing(),
version = %latest,
"fetched via get_latest_matching fallback"
);
Some((
latest,
v.removal_status(),
v.published_at(),
v.deprecation().cloned(),
v.license().to_vec(),
))
}
Ok(Ok(None)) => {
tracing::debug!(package = %name.for_tracing(), "no version found");
no_comparable_versions = true;
None
}
Ok(Err(e)) => {
tracing::warn!(
package = %name.for_tracing(),
error = %e,
"fetch fallback failed"
);
failed.fetch_add(1, std::sync::atomic::Ordering::Relaxed);
let mut fe = first_error.lock().unwrap_or_else(|p| p.into_inner());
if fe.is_none() {
*fe = Some(e.to_string());
}
drop(fe);
if !e.is_not_found() {
failed_name = Some((name.clone(), e.fetch_failure(), e.to_string()));
}
None
}
Err(_) => {
tracing::warn!(
package = %name.for_tracing(),
"fetch fallback timed out ({}s)",
timeout.as_secs()
);
failed.fetch_add(1, std::sync::atomic::Ordering::Relaxed);
failed_name = Some((
name.clone(),
FetchFailure::Transient,
format!(
"{}: registry request timed out after {}s",
name.for_tracing(),
timeout.as_secs()
),
));
None
}
}
};
if check_yanked {
if let Some((latest, status, _, _, _)) = &resolved
&& status.is_flagged()
{
yanked = Some((name.clone(), latest.clone(), *status));
}
if let Some((iv, status)) = in_use_versions.iter().find_map(|iv| {
versions
.iter()
.find(|v| {
v.version_string() == iv.as_str() && v.removal_status().is_flagged()
})
.map(|v| (iv, v.removal_status()))
}) {
yanked = Some((name.clone(), iv.as_str().into(), status));
}
}
if let Some((_, _, _, dep_info, _)) = &resolved
&& let Some(dep_info) = dep_info
{
deprecation = Some((name.clone(), dep_info.clone()));
}
license = resolved
.as_ref()
.map(|(_, _, _, _, lic)| lic)
.filter(|lic| !lic.is_empty())
.map(|lic| (name.clone(), lic.clone()));
resolved.map(|(latest, _, published_at, _, _)| {
let mut versions = PackageVersions::new(latest, available).with_yanked(yanked_list);
if let Some(published_at) = published_at {
versions = versions.with_published_at(published_at);
}
(name.clone(), versions)
})
}
Ok(Err(e)) => {
if e.is_offline() {
tracing::debug!(package = %name.for_tracing(), "fetch skipped: offline");
} else {
tracing::warn!(package = %name.for_tracing(), error = %e, "fetch failed");
}
failed.fetch_add(1, std::sync::atomic::Ordering::Relaxed);
let mut fe = first_error.lock().unwrap_or_else(|p| p.into_inner());
if fe.is_none() {
*fe = Some(e.to_string());
}
drop(fe);
if !e.is_not_found() {
failed_name = Some((name.clone(), e.fetch_failure(), e.to_string()));
}
None
}
Err(_) => {
tracing::warn!(package = %name.for_tracing(), "fetch timed out ({}s)", timeout.as_secs());
failed.fetch_add(1, std::sync::atomic::Ordering::Relaxed);
failed_name = Some((
name.clone(),
FetchFailure::Transient,
format!(
"{}: registry request timed out after {}s",
name.for_tracing(),
timeout.as_secs()
),
));
None
}
};
let count = fetched.fetch_add(1, std::sync::atomic::Ordering::Relaxed) + 1;
if let Some(sender) = progress_sender {
sender.send(count);
}
let no_comparable_versions_name = no_comparable_versions.then(|| name.clone());
(
version,
yanked,
failed_name,
deprecation,
no_comparable_versions_name,
license,
)
}
pub fn apply_fetch_outcomes(
outcomes: &mut deps_core::lsp_helpers::DependencyOutcomes,
yanked_versions: HashMap<PackageName, (ConcreteVersion, RemovalStatus)>,
fetch_failed: HashMap<PackageName, FetchFailure>,
collided_names: HashSet<PackageName>,
formatter: &dyn deps_core::lsp_helpers::EcosystemFormatter,
) {
for (name, version) in yanked_versions {
outcomes.set_yanked(formatter.normalize_package_name(&name), version);
}
for (name, failure) in fetch_failed {
outcomes.set_fetch_failure(formatter.normalize_package_name(&name), failure);
}
for name in collided_names {
outcomes.set_fetch_failure_if_absent(
formatter.normalize_package_name(&name),
FetchFailure::NotAttempted,
);
}
}
#[cfg(test)]
mod tests {
use super::*;
use deps_core::parser::DependencySource;
fn with_registry_source(names: Vec<PackageName>) -> Vec<(PackageName, DependencySource)> {
names
.into_iter()
.map(|name| (name, DependencySource::Registry))
.collect()
}
#[cfg(feature = "cargo")]
mod apply_fetch_outcomes_tests {
use super::*;
use crate::setup::CargoFormatter;
use deps_core::lsp_helpers::DependencyOutcomes;
#[test]
fn apply_fetch_outcomes_sets_yanked_re_keyed_to_normalized_name() {
let mut outcomes = DependencyOutcomes::new();
let mut yanked_versions = HashMap::new();
yanked_versions.insert(
PackageName::new("time"),
("0.1.43".into(), RemovalStatus::Yanked),
);
apply_fetch_outcomes(
&mut outcomes,
yanked_versions,
HashMap::new(),
HashSet::new(),
&CargoFormatter,
);
assert_eq!(
outcomes.yanked("time"),
Some(&("0.1.43".into(), RemovalStatus::Yanked))
);
}
#[test]
fn apply_fetch_outcomes_collided_name_does_not_clobber_existing_fetch_failure() {
let mut outcomes = DependencyOutcomes::new();
let mut fetch_failed = HashMap::new();
fetch_failed.insert(PackageName::new("serde"), FetchFailure::Transient);
let mut collided_names = HashSet::new();
collided_names.insert(PackageName::new("serde"));
apply_fetch_outcomes(
&mut outcomes,
HashMap::new(),
fetch_failed,
collided_names,
&CargoFormatter,
);
assert_eq!(
outcomes.fetch_failure("serde"),
Some(&FetchFailure::Transient),
"a genuine fetch failure must survive a collided name normalizing to the same key"
);
}
#[test]
fn apply_fetch_outcomes_collided_name_alone_is_recorded_as_not_attempted() {
let mut outcomes = DependencyOutcomes::new();
let mut collided_names = HashSet::new();
collided_names.insert(PackageName::new("serde"));
apply_fetch_outcomes(
&mut outcomes,
HashMap::new(),
HashMap::new(),
collided_names,
&CargoFormatter,
);
assert_eq!(
outcomes.fetch_failure("serde"),
Some(&FetchFailure::NotAttempted)
);
}
}
mod dedup_by_source_collision_tests {
use super::*;
use deps_core::Dependency;
use deps_core::lsp_helpers::{
DiagnosticMessages, DiagnosticPolicy, OsvNaming, PackageNaming, PackageRendering,
RequirementResolution, SourcePolicy,
};
use deps_core::position::{Position, Range};
use std::any::Any;
struct AlternateAwareFormatter;
impl PackageNaming for AlternateAwareFormatter {}
impl PackageRendering for AlternateAwareFormatter {
fn format_version_for_text_edit(&self, version: &ConcreteVersion) -> String {
version.to_string()
}
fn package_url(&self, name: &PackageName) -> String {
format!("https://example.com/{}", name.as_str())
}
}
impl RequirementResolution for AlternateAwareFormatter {}
impl DiagnosticMessages for AlternateAwareFormatter {}
impl DiagnosticPolicy for AlternateAwareFormatter {}
impl SourcePolicy for AlternateAwareFormatter {
fn can_resolve_source(&self, source: &DependencySource) -> bool {
matches!(
source,
DependencySource::Registry | DependencySource::AlternateRegistry { .. }
)
}
}
impl OsvNaming for AlternateAwareFormatter {}
struct MockDep {
name: PackageName,
source: DependencySource,
addr_tag: u32,
}
impl Dependency for MockDep {
fn name(&self) -> &PackageName {
&self.name
}
fn name_range(&self) -> Range {
Range::new(
Position::new(0, self.addr_tag),
Position::new(0, self.addr_tag + 1),
)
}
fn version_requirement(&self) -> Option<&VersionReq> {
None
}
fn version_range(&self) -> Option<Range> {
None
}
fn source(&self) -> DependencySource {
self.source.clone()
}
fn as_any(&self) -> &dyn Any {
self
}
}
struct MockParseResult {
deps: Vec<MockDep>,
}
impl deps_core::ParseResult for MockParseResult {
fn dependencies(&self) -> Vec<&dyn Dependency> {
self.deps.iter().map(|d| d as &dyn Dependency).collect()
}
fn workspace_root(&self) -> Option<&std::path::Path> {
None
}
fn uri(&self) -> &url::Url {
static URI: std::sync::OnceLock<url::Url> = std::sync::OnceLock::new();
URI.get_or_init(|| deps_core::test_util::test_uri("/test/Cargo.toml"))
}
fn as_any(&self) -> &dyn Any {
self
}
}
#[test]
fn test_two_different_resolvable_sources_collide_and_are_dropped() {
let parse_result = MockParseResult {
deps: vec![
MockDep {
name: PackageName::new("shared-name"),
source: DependencySource::Registry,
addr_tag: 0,
},
MockDep {
name: PackageName::new("shared-name"),
source: DependencySource::AlternateRegistry {
index: "https://index.mycorp.dev".into(),
mirrors_crates_io: false,
},
addr_tag: 1,
},
],
};
let (sources, collided) =
dedup_dependencies_by_source(&parse_result, &AlternateAwareFormatter);
assert!(
!sources.contains_key(&PackageName::new("shared-name")),
"a colliding name must not be fetched under either source"
);
assert!(
collided.contains(&PackageName::new("shared-name")),
"the collision must be recorded so the caller can mark it fetch_failed"
);
}
#[test]
fn test_identical_sources_do_not_collide() {
let parse_result = MockParseResult {
deps: vec![
MockDep {
name: PackageName::new("shared-name"),
source: DependencySource::Registry,
addr_tag: 0,
},
MockDep {
name: PackageName::new("shared-name"),
source: DependencySource::Registry,
addr_tag: 1,
},
],
};
let (sources, collided) =
dedup_dependencies_by_source(&parse_result, &AlternateAwareFormatter);
assert!(collided.is_empty());
assert_eq!(
sources.get(&PackageName::new("shared-name")),
Some(&DependencySource::Registry)
);
}
#[test]
fn test_non_resolvable_source_is_dropped_not_fetched() {
let parse_result = MockParseResult {
deps: vec![MockDep {
name: PackageName::new("local-fork"),
source: DependencySource::Path {
path: "../local-fork".into(),
},
addr_tag: 0,
}],
};
let (sources, collided) =
dedup_dependencies_by_source(&parse_result, &AlternateAwareFormatter);
assert!(sources.is_empty());
assert!(collided.is_empty());
}
#[test]
fn test_collision_warning_redacts_credentials_in_alternate_registry_debug_output() {
let parse_result = MockParseResult {
deps: vec![
MockDep {
name: PackageName::new("shared-name"),
source: DependencySource::AlternateRegistry {
index: "https://index-a.mycorp.dev/api?api_key=SECRET_A".into(),
mirrors_crates_io: false,
},
addr_tag: 0,
},
MockDep {
name: PackageName::new("shared-name"),
source: DependencySource::AlternateRegistry {
index: "https://index-b.mycorp.dev/api?api_key=SECRET_B".into(),
mirrors_crates_io: false,
},
addr_tag: 1,
},
],
};
let log = deps_core::test_util::capture_tracing_output(|| {
let (sources, collided) =
dedup_dependencies_by_source(&parse_result, &AlternateAwareFormatter);
assert!(!sources.contains_key(&PackageName::new("shared-name")));
assert!(collided.contains(&PackageName::new("shared-name")));
});
assert!(
log.contains("two different resolved registries"),
"expected the collision WARN to fire: {log:?}"
);
assert!(
!log.contains("SECRET_A") && !log.contains("SECRET_B"),
"tracing output leaked a query-string credential: {log:?}"
);
assert!(
log.contains("index-a.mycorp.dev") && log.contains("index-b.mycorp.dev"),
"host should survive redaction: {log:?}"
);
}
}
#[tokio::test]
async fn test_fetch_latest_versions_parallel_with_timeout() {
use deps_core::{Metadata, Registry, Version};
use std::any::Any;
use std::time::Duration;
struct TimeoutRegistry;
impl Registry for TimeoutRegistry {
fn get_versions<'a>(
&'a self,
_name: &'a deps_core::PackageName,
) -> deps_core::ecosystem::BoxFuture<'a, deps_core::Result<Vec<Box<dyn Version>>>>
{
Box::pin(async move {
tokio::time::sleep(Duration::from_secs(10)).await;
Ok(vec![])
})
}
fn get_latest_matching<'a>(
&'a self,
_name: &'a deps_core::PackageName,
_req: &'a deps_core::VersionReq,
) -> deps_core::ecosystem::BoxFuture<'a, deps_core::Result<Option<Box<dyn Version>>>>
{
Box::pin(async move {
tokio::time::sleep(Duration::from_secs(10)).await;
Ok(None)
})
}
fn search_raw<'a>(
&'a self,
_query: &'a str,
_limit: usize,
) -> deps_core::ecosystem::BoxFuture<'a, deps_core::Result<Vec<Box<dyn Metadata>>>>
{
Box::pin(async move { Ok(vec![]) })
}
fn as_any(&self) -> &dyn Any {
self
}
}
let registry: Arc<dyn Registry> = Arc::new(TimeoutRegistry);
let packages = vec![PackageName::new("slow-package")];
let result = fetch_latest_versions_parallel(
registry,
with_registry_source(packages),
&HashMap::new(),
None,
deps_core::freshness::FreshnessSettings::default(),
1,
10,
None,
)
.await;
assert!(result.versions.is_empty(), "Slow package should timeout");
assert_eq!(result.failed_count, 1, "Should track 1 failed package");
assert_eq!(
result.fetch_failed,
HashMap::from([(PackageName::new("slow-package"), FetchFailure::Transient)]),
"timed-out package must be recorded in fetch_failed"
);
}
#[tokio::test]
async fn test_fetch_latest_versions_parallel_fast_packages_not_blocked() {
use deps_core::{Metadata, Registry, Version};
use std::any::Any;
use std::time::Duration;
struct MixedRegistry;
impl Registry for MixedRegistry {
fn get_versions<'a>(
&'a self,
name: &'a deps_core::PackageName,
) -> deps_core::ecosystem::BoxFuture<'a, deps_core::Result<Vec<Box<dyn Version>>>>
{
Box::pin(async move {
if name == "slow-package" {
tokio::time::sleep(Duration::from_secs(10)).await;
}
Ok(vec![])
})
}
fn get_latest_matching<'a>(
&'a self,
name: &'a deps_core::PackageName,
_req: &'a deps_core::VersionReq,
) -> deps_core::ecosystem::BoxFuture<'a, deps_core::Result<Option<Box<dyn Version>>>>
{
Box::pin(async move {
if name == "slow-package" {
tokio::time::sleep(Duration::from_secs(10)).await;
}
Ok(None)
})
}
fn search_raw<'a>(
&'a self,
_query: &'a str,
_limit: usize,
) -> deps_core::ecosystem::BoxFuture<'a, deps_core::Result<Vec<Box<dyn Metadata>>>>
{
Box::pin(async move { Ok(vec![]) })
}
fn as_any(&self) -> &dyn Any {
self
}
}
let registry: Arc<dyn Registry> = Arc::new(MixedRegistry);
let packages = vec![
PackageName::new("slow-package"),
PackageName::new("fast-package"),
];
let start = std::time::Instant::now();
let result = fetch_latest_versions_parallel(
registry,
with_registry_source(packages),
&HashMap::new(),
None,
deps_core::freshness::FreshnessSettings::default(),
1,
10,
None,
)
.await;
let elapsed = start.elapsed();
assert!(
elapsed < Duration::from_secs(3),
"Should not wait for slow package: {:?}",
elapsed
);
assert!(
result.versions.is_empty(),
"No versions returned (test registry returns empty)"
);
assert_eq!(
result.failed_count, 1,
"Slow package should be marked as failed"
);
}
#[tokio::test]
async fn test_fetch_latest_versions_parallel_concurrency_limit() {
use deps_core::{Metadata, Registry, Version};
use std::any::Any;
use std::sync::atomic::{AtomicUsize, Ordering};
use std::time::Duration;
struct ConcurrencyTrackingRegistry {
current: Arc<AtomicUsize>,
max_seen: Arc<AtomicUsize>,
}
impl Registry for ConcurrencyTrackingRegistry {
fn get_versions<'a>(
&'a self,
_name: &'a deps_core::PackageName,
) -> deps_core::ecosystem::BoxFuture<'a, deps_core::Result<Vec<Box<dyn Version>>>>
{
Box::pin(async move {
let current = self.current.fetch_add(1, Ordering::SeqCst) + 1;
self.max_seen.fetch_max(current, Ordering::SeqCst);
tokio::time::sleep(Duration::from_millis(50)).await;
self.current.fetch_sub(1, Ordering::SeqCst);
Ok(vec![])
})
}
fn get_latest_matching<'a>(
&'a self,
_name: &'a deps_core::PackageName,
_req: &'a deps_core::VersionReq,
) -> deps_core::ecosystem::BoxFuture<'a, deps_core::Result<Option<Box<dyn Version>>>>
{
Box::pin(async move {
let current = self.current.fetch_add(1, Ordering::SeqCst) + 1;
self.max_seen.fetch_max(current, Ordering::SeqCst);
tokio::time::sleep(Duration::from_millis(50)).await;
self.current.fetch_sub(1, Ordering::SeqCst);
Ok(None)
})
}
fn search_raw<'a>(
&'a self,
_query: &'a str,
_limit: usize,
) -> deps_core::ecosystem::BoxFuture<'a, deps_core::Result<Vec<Box<dyn Metadata>>>>
{
Box::pin(async move { Ok(vec![]) })
}
fn as_any(&self) -> &dyn Any {
self
}
}
let current = Arc::new(AtomicUsize::new(0));
let max_seen = Arc::new(AtomicUsize::new(0));
let registry: Arc<dyn Registry> = Arc::new(ConcurrencyTrackingRegistry {
current: Arc::clone(¤t),
max_seen: Arc::clone(&max_seen),
});
let packages: Vec<PackageName> = (0..50)
.map(|i| PackageName::new(format!("package-{}", i)))
.collect();
fetch_latest_versions_parallel(
registry,
with_registry_source(packages),
&HashMap::new(),
None,
deps_core::freshness::FreshnessSettings::default(),
5,
20,
None,
)
.await;
let max = max_seen.load(Ordering::SeqCst);
assert!(
max <= 22,
"Concurrency limit violated: {} concurrent requests (limit: 20)",
max
);
}
#[tokio::test]
async fn test_fetch_latest_versions_parallel_zero_max_concurrent_still_completes() {
use deps_core::{Metadata, Registry, Version};
use std::any::Any;
struct InstantRegistry;
impl Registry for InstantRegistry {
fn get_versions<'a>(
&'a self,
_name: &'a deps_core::PackageName,
) -> deps_core::ecosystem::BoxFuture<'a, deps_core::Result<Vec<Box<dyn Version>>>>
{
Box::pin(async move { Ok(vec![]) })
}
fn get_latest_matching<'a>(
&'a self,
_name: &'a deps_core::PackageName,
_req: &'a deps_core::VersionReq,
) -> deps_core::ecosystem::BoxFuture<'a, deps_core::Result<Option<Box<dyn Version>>>>
{
Box::pin(async move { Ok(None) })
}
fn search_raw<'a>(
&'a self,
_query: &'a str,
_limit: usize,
) -> deps_core::ecosystem::BoxFuture<'a, deps_core::Result<Vec<Box<dyn Metadata>>>>
{
Box::pin(async move { Ok(vec![]) })
}
fn as_any(&self) -> &dyn Any {
self
}
}
let registry: Arc<dyn Registry> = Arc::new(InstantRegistry);
let packages = vec![PackageName::new("some-package")];
let result = tokio::time::timeout(
std::time::Duration::from_secs(5),
fetch_latest_versions_parallel(
registry,
with_registry_source(packages),
&HashMap::new(),
None,
deps_core::freshness::FreshnessSettings::default(),
5,
0,
None,
),
)
.await
.expect("fetch with max_concurrent=0 must not hang forever");
assert!(
result
.no_comparable_versions
.contains(&PackageName::new("some-package")),
"fetch must still run to completion when max_concurrent is 0"
);
}
#[tokio::test]
async fn test_fetch_partial_success_with_mixed_outcomes() {
use deps_core::{Metadata, Registry, Version};
use std::any::Any;
use std::time::Duration;
#[derive(Debug)]
struct MockVersion {
version: ConcreteVersion,
}
impl Version for MockVersion {
fn version_string(&self) -> &ConcreteVersion {
&self.version
}
fn is_prerelease(&self) -> bool {
false
}
fn as_any(&self) -> &dyn Any {
self
}
}
struct MixedOutcomeRegistry;
impl Registry for MixedOutcomeRegistry {
fn get_versions<'a>(
&'a self,
name: &'a deps_core::PackageName,
) -> deps_core::ecosystem::BoxFuture<'a, deps_core::Result<Vec<Box<dyn Version>>>>
{
Box::pin(async move {
match name.as_str() {
"package-fast" => Ok(vec![Box::new(MockVersion {
version: "1.0.0".into(),
}) as Box<dyn Version>]),
"package-slow" => {
tokio::time::sleep(Duration::from_secs(10)).await;
Ok(vec![])
}
"package-error" => Err(deps_core::error::DepsError::CacheError(
"Mock registry error".to_string(),
)),
_ => Ok(vec![]),
}
})
}
fn get_latest_matching<'a>(
&'a self,
name: &'a deps_core::PackageName,
_req: &'a deps_core::VersionReq,
) -> deps_core::ecosystem::BoxFuture<'a, deps_core::Result<Option<Box<dyn Version>>>>
{
Box::pin(async move {
match name.as_str() {
"package-fast" => Ok(Some(Box::new(MockVersion {
version: "1.0.0".into(),
}) as Box<dyn Version>)),
"package-slow" => {
tokio::time::sleep(Duration::from_secs(10)).await;
Ok(None)
}
"package-error" => Err(deps_core::error::DepsError::CacheError(
"Mock registry error".to_string(),
)),
_ => Ok(None),
}
})
}
fn search_raw<'a>(
&'a self,
_query: &'a str,
_limit: usize,
) -> deps_core::ecosystem::BoxFuture<'a, deps_core::Result<Vec<Box<dyn Metadata>>>>
{
Box::pin(async move { Ok(vec![]) })
}
fn select_latest_matching(
&self,
versions: &[Box<dyn Version>],
_req: &deps_core::VersionReq,
) -> Option<usize> {
if versions.is_empty() { None } else { Some(0) }
}
fn as_any(&self) -> &dyn Any {
self
}
}
let registry: Arc<dyn Registry> = Arc::new(MixedOutcomeRegistry);
let packages = vec![
PackageName::new("package-fast"),
PackageName::new("package-slow"),
PackageName::new("package-error"),
];
let result = fetch_latest_versions_parallel(
registry,
with_registry_source(packages),
&HashMap::new(),
None,
deps_core::freshness::FreshnessSettings::default(),
1,
10,
None,
)
.await;
assert_eq!(
result.versions.len(),
1,
"Should have exactly 1 successful package"
);
assert_eq!(
result
.versions
.get("package-fast")
.map(|v| v.latest.as_str()),
Some("1.0.0"),
"Fast package should have correct version"
);
assert!(
!result.versions.contains_key("package-slow"),
"Slow package should not be in results (timeout)"
);
assert!(
!result.versions.contains_key("package-error"),
"Error package should not be in results"
);
}
#[tokio::test]
async fn test_fetch_latest_versions_parallel_carries_yanked_flag_into_cache() {
use deps_core::{Metadata, Registry, Version};
use std::any::Any;
#[derive(Debug)]
struct MockVersion {
version: ConcreteVersion,
yanked: bool,
}
impl Version for MockVersion {
fn version_string(&self) -> &ConcreteVersion {
&self.version
}
fn removal_status(&self) -> deps_core::RemovalStatus {
deps_core::RemovalStatus::from_yanked(self.yanked)
}
fn as_any(&self) -> &dyn Any {
self
}
}
struct YankedRegistry;
impl Registry for YankedRegistry {
fn get_versions<'a>(
&'a self,
_name: &'a deps_core::PackageName,
) -> deps_core::ecosystem::BoxFuture<'a, deps_core::Result<Vec<Box<dyn Version>>>>
{
Box::pin(async move {
Ok(vec![
Box::new(MockVersion {
version: "1.0.214".into(),
yanked: false,
}) as Box<dyn Version>,
Box::new(MockVersion {
version: "1.0.213".into(),
yanked: true,
}) as Box<dyn Version>,
])
})
}
fn get_latest_matching<'a>(
&'a self,
_name: &'a deps_core::PackageName,
_req: &'a deps_core::VersionReq,
) -> deps_core::ecosystem::BoxFuture<'a, deps_core::Result<Option<Box<dyn Version>>>>
{
Box::pin(async move { Ok(None) })
}
fn search_raw<'a>(
&'a self,
_query: &'a str,
_limit: usize,
) -> deps_core::ecosystem::BoxFuture<'a, deps_core::Result<Vec<Box<dyn Metadata>>>>
{
Box::pin(async move { Ok(vec![]) })
}
fn select_latest_matching(
&self,
versions: &[Box<dyn Version>],
_req: &deps_core::VersionReq,
) -> Option<usize> {
versions
.iter()
.position(|v| !v.removal_status().blocks_resolution())
}
fn as_any(&self) -> &dyn Any {
self
}
}
let registry: Arc<dyn Registry> = Arc::new(YankedRegistry);
let packages = vec![PackageName::new("serde")];
let result = fetch_latest_versions_parallel(
registry,
with_registry_source(packages),
&HashMap::new(),
None,
deps_core::freshness::FreshnessSettings::default(),
10,
10,
None,
)
.await;
let serde = result
.versions
.get("serde")
.expect("serde should be fetched");
assert_eq!(serde.latest, "1.0.214", "latest must skip the yanked entry");
assert_eq!(
&*serde.available,
&[
ConcreteVersion::new("1.0.214"),
ConcreteVersion::new("1.0.213")
],
"available must remain unfiltered"
);
assert_eq!(
&*serde.yanked,
&[(
ConcreteVersion::new("1.0.213"),
deps_core::RemovalStatus::Yanked
)],
"yanked must carry only the entries reported as yanked, paired with their status"
);
}
#[tokio::test]
async fn test_fetch_latest_versions_parallel_carries_published_at_for_latest_only() {
use deps_core::freshness::PublishTime;
use deps_core::{Metadata, Registry, Version};
use std::any::Any;
#[derive(Debug)]
struct MockVersion {
version: ConcreteVersion,
yanked: bool,
published_at: Option<PublishTime>,
}
impl Version for MockVersion {
fn version_string(&self) -> &ConcreteVersion {
&self.version
}
fn removal_status(&self) -> deps_core::RemovalStatus {
deps_core::RemovalStatus::from_yanked(self.yanked)
}
fn published_at(&self) -> Option<PublishTime> {
self.published_at
}
fn as_any(&self) -> &dyn Any {
self
}
}
struct DatedRegistry;
impl Registry for DatedRegistry {
fn get_versions<'a>(
&'a self,
_name: &'a deps_core::PackageName,
) -> deps_core::ecosystem::BoxFuture<'a, deps_core::Result<Vec<Box<dyn Version>>>>
{
Box::pin(async move {
Ok(vec![
Box::new(MockVersion {
version: "1.0.214".into(),
yanked: false,
published_at: Some(PublishTime::from_unix_secs(2_000)),
}) as Box<dyn Version>,
Box::new(MockVersion {
version: "1.0.213".into(),
yanked: true,
published_at: Some(PublishTime::from_unix_secs(1_000)),
}) as Box<dyn Version>,
])
})
}
fn get_latest_matching<'a>(
&'a self,
_name: &'a deps_core::PackageName,
_req: &'a deps_core::VersionReq,
) -> deps_core::ecosystem::BoxFuture<'a, deps_core::Result<Option<Box<dyn Version>>>>
{
Box::pin(async move { Ok(None) })
}
fn search_raw<'a>(
&'a self,
_query: &'a str,
_limit: usize,
) -> deps_core::ecosystem::BoxFuture<'a, deps_core::Result<Vec<Box<dyn Metadata>>>>
{
Box::pin(async move { Ok(vec![]) })
}
fn select_latest_matching(
&self,
versions: &[Box<dyn Version>],
_req: &deps_core::VersionReq,
) -> Option<usize> {
versions
.iter()
.position(|v| !v.removal_status().blocks_resolution())
}
fn as_any(&self) -> &dyn Any {
self
}
}
let registry: Arc<dyn Registry> = Arc::new(DatedRegistry);
let packages = vec![PackageName::new("serde")];
let result = fetch_latest_versions_parallel(
registry,
with_registry_source(packages),
&HashMap::new(),
None,
deps_core::freshness::FreshnessSettings::default(),
10,
10,
None,
)
.await;
let serde = result
.versions
.get("serde")
.expect("serde should be fetched");
assert_eq!(serde.latest, "1.0.214");
assert_eq!(
serde.published_at,
Some(PublishTime::from_unix_secs(2_000)),
"published_at must be 1.0.214's own timestamp, not the yanked 1.0.213 entry's"
);
}
#[tokio::test]
async fn test_fetch_latest_versions_parallel_carries_license_into_fetch_result() {
use deps_core::{Metadata, Registry, Version};
use std::any::Any;
#[derive(Debug)]
struct MockVersion {
version: ConcreteVersion,
license: Vec<String>,
}
impl Version for MockVersion {
fn version_string(&self) -> &ConcreteVersion {
&self.version
}
fn as_any(&self) -> &dyn Any {
self
}
fn license(&self) -> &[String] {
&self.license
}
}
struct LicensedRegistry;
impl Registry for LicensedRegistry {
fn get_versions<'a>(
&'a self,
name: &'a PackageName,
) -> deps_core::ecosystem::BoxFuture<'a, deps_core::Result<Vec<Box<dyn Version>>>>
{
let license = if name.as_str() == "licensed-pkg" {
vec!["MIT".to_string()]
} else {
vec![]
};
Box::pin(async move {
Ok(vec![Box::new(MockVersion {
version: "1.0.0".into(),
license,
}) as Box<dyn Version>])
})
}
fn get_latest_matching<'a>(
&'a self,
_name: &'a PackageName,
_req: &'a deps_core::VersionReq,
) -> deps_core::ecosystem::BoxFuture<'a, deps_core::Result<Option<Box<dyn Version>>>>
{
Box::pin(async move { Ok(None) })
}
fn search_raw<'a>(
&'a self,
_query: &'a str,
_limit: usize,
) -> deps_core::ecosystem::BoxFuture<'a, deps_core::Result<Vec<Box<dyn Metadata>>>>
{
Box::pin(async move { Ok(vec![]) })
}
fn select_latest_matching(
&self,
versions: &[Box<dyn Version>],
_req: &deps_core::VersionReq,
) -> Option<usize> {
(!versions.is_empty()).then_some(0)
}
fn as_any(&self) -> &dyn Any {
self
}
}
let registry: Arc<dyn Registry> = Arc::new(LicensedRegistry);
let packages = vec![
PackageName::new("licensed-pkg"),
PackageName::new("unlicensed-pkg"),
];
let result = fetch_latest_versions_parallel(
registry,
with_registry_source(packages),
&HashMap::new(),
None,
deps_core::freshness::FreshnessSettings::default(),
10,
10,
None,
)
.await;
assert_eq!(
result.licenses.get(&PackageName::new("licensed-pkg")),
Some(&vec!["MIT".to_string()])
);
assert!(
!result
.licenses
.contains_key(&PackageName::new("unlicensed-pkg")),
"an empty Version::license() must produce no entry, not an empty-vec one"
);
}
#[tokio::test]
async fn test_fetch_latest_versions_parallel_uses_get_versions_with_for_freshness() {
use deps_core::freshness::{FreshnessSettings, PublishTime};
use deps_core::{Metadata, Registry, Version};
use std::any::Any;
#[derive(Debug)]
struct MockVersion {
version: ConcreteVersion,
published_at: Option<PublishTime>,
}
impl Version for MockVersion {
fn version_string(&self) -> &ConcreteVersion {
&self.version
}
fn published_at(&self) -> Option<PublishTime> {
self.published_at
}
fn as_any(&self) -> &dyn Any {
self
}
}
struct FreshnessAwareRegistry;
impl Registry for FreshnessAwareRegistry {
fn get_versions<'a>(
&'a self,
_name: &'a deps_core::PackageName,
) -> deps_core::ecosystem::BoxFuture<'a, deps_core::Result<Vec<Box<dyn Version>>>>
{
Box::pin(async move {
Ok(vec![Box::new(MockVersion {
version: "1.0.0".into(),
published_at: None,
}) as Box<dyn Version>])
})
}
fn get_versions_with<'a>(
&'a self,
_name: &'a deps_core::PackageName,
freshness: FreshnessSettings,
) -> deps_core::ecosystem::BoxFuture<'a, deps_core::Result<Vec<Box<dyn Version>>>>
{
Box::pin(async move {
Ok(vec![Box::new(MockVersion {
version: "1.0.0".into(),
published_at: freshness
.enabled
.then(|| PublishTime::from_unix_secs(5_000)),
}) as Box<dyn Version>])
})
}
fn get_latest_matching<'a>(
&'a self,
_name: &'a deps_core::PackageName,
_req: &'a deps_core::VersionReq,
) -> deps_core::ecosystem::BoxFuture<'a, deps_core::Result<Option<Box<dyn Version>>>>
{
Box::pin(async move { Ok(None) })
}
fn search_raw<'a>(
&'a self,
_query: &'a str,
_limit: usize,
) -> deps_core::ecosystem::BoxFuture<'a, deps_core::Result<Vec<Box<dyn Metadata>>>>
{
Box::pin(async move { Ok(vec![]) })
}
fn select_latest_matching(
&self,
versions: &[Box<dyn Version>],
_req: &deps_core::VersionReq,
) -> Option<usize> {
if versions.is_empty() { None } else { Some(0) }
}
fn as_any(&self) -> &dyn Any {
self
}
}
let registry: Arc<dyn Registry> = Arc::new(FreshnessAwareRegistry);
let packages = vec![PackageName::new("widget")];
let result = fetch_latest_versions_parallel(
registry,
with_registry_source(packages),
&HashMap::new(),
None,
FreshnessSettings::default(),
10,
10,
None,
)
.await;
let widget = result
.versions
.get("widget")
.expect("widget should be fetched");
assert_eq!(
widget.published_at,
Some(PublishTime::from_unix_secs(5_000)),
"published_at must come from get_versions_with, not the freshness-blind \
get_versions (#339)"
);
}
#[tokio::test]
async fn test_fetch_latest_versions_parallel_threads_minimum_stability_into_select_latest_matching_with_context()
{
use deps_core::{Metadata, Registry, Version};
use std::any::Any;
use std::sync::Mutex;
#[derive(Debug)]
struct MockVersion {
version: ConcreteVersion,
}
impl Version for MockVersion {
fn version_string(&self) -> &ConcreteVersion {
&self.version
}
fn as_any(&self) -> &dyn Any {
self
}
}
struct ContextAwareRegistry {
seen_minimum_stability: Mutex<Vec<Option<String>>>,
}
impl Registry for ContextAwareRegistry {
fn get_versions<'a>(
&'a self,
_name: &'a deps_core::PackageName,
) -> deps_core::ecosystem::BoxFuture<'a, deps_core::Result<Vec<Box<dyn Version>>>>
{
Box::pin(async move {
Ok(vec![Box::new(MockVersion {
version: "1.0.0".into(),
}) as Box<dyn Version>])
})
}
fn get_latest_matching<'a>(
&'a self,
_name: &'a deps_core::PackageName,
_req: &'a deps_core::VersionReq,
) -> deps_core::ecosystem::BoxFuture<'a, deps_core::Result<Option<Box<dyn Version>>>>
{
Box::pin(async move { Ok(None) })
}
fn search_raw<'a>(
&'a self,
_query: &'a str,
_limit: usize,
) -> deps_core::ecosystem::BoxFuture<'a, deps_core::Result<Vec<Box<dyn Metadata>>>>
{
Box::pin(async move { Ok(vec![]) })
}
fn select_latest_matching_with_context(
&self,
versions: &[Box<dyn Version>],
_req: &deps_core::VersionReq,
minimum_stability: Option<&str>,
) -> Option<usize> {
self.seen_minimum_stability
.lock()
.unwrap_or_else(|p| p.into_inner())
.push(minimum_stability.map(str::to_string));
if versions.is_empty() { None } else { Some(0) }
}
fn as_any(&self) -> &dyn Any {
self
}
}
let registry = Arc::new(ContextAwareRegistry {
seen_minimum_stability: Mutex::new(Vec::new()),
});
let packages = vec![PackageName::new("vendor/pkg")];
let result = fetch_latest_versions_parallel(
Arc::clone(®istry) as Arc<dyn Registry>,
with_registry_source(packages),
&HashMap::new(),
None,
deps_core::freshness::FreshnessSettings::default(),
10,
10,
Some("beta"),
)
.await;
assert_eq!(
*registry
.seen_minimum_stability
.lock()
.unwrap_or_else(|p| p.into_inner()),
vec![Some("beta".to_string())],
"select_latest_matching_with_context must receive the caller's minimum_stability"
);
assert!(
result.versions.contains_key("vendor/pkg"),
"the pick must still succeed via the _with_context path"
);
}
#[tokio::test]
async fn test_fetch_latest_versions_parallel_threads_minimum_stability_into_get_latest_matching_with_context()
{
use deps_core::{Metadata, Registry, Version};
use std::any::Any;
use std::sync::Mutex;
#[derive(Debug)]
struct MockVersion {
version: ConcreteVersion,
}
impl Version for MockVersion {
fn version_string(&self) -> &ConcreteVersion {
&self.version
}
fn as_any(&self) -> &dyn Any {
self
}
}
struct FallbackContextAwareRegistry {
seen_minimum_stability: Mutex<Vec<Option<String>>>,
}
impl Registry for FallbackContextAwareRegistry {
fn get_versions<'a>(
&'a self,
_name: &'a deps_core::PackageName,
) -> deps_core::ecosystem::BoxFuture<'a, deps_core::Result<Vec<Box<dyn Version>>>>
{
Box::pin(async move { Ok(vec![]) })
}
fn get_latest_matching<'a>(
&'a self,
_name: &'a deps_core::PackageName,
_req: &'a deps_core::VersionReq,
) -> deps_core::ecosystem::BoxFuture<'a, deps_core::Result<Option<Box<dyn Version>>>>
{
Box::pin(async move { Ok(None) })
}
fn get_latest_matching_with_context<'a>(
&'a self,
_name: &'a deps_core::PackageName,
_req: &'a deps_core::VersionReq,
minimum_stability: Option<&'a str>,
) -> deps_core::ecosystem::BoxFuture<'a, deps_core::Result<Option<Box<dyn Version>>>>
{
self.seen_minimum_stability
.lock()
.unwrap_or_else(|p| p.into_inner())
.push(minimum_stability.map(str::to_string));
Box::pin(async move {
Ok(Some(Box::new(MockVersion {
version: "2.0.0-beta1".into(),
}) as Box<dyn Version>))
})
}
fn search_raw<'a>(
&'a self,
_query: &'a str,
_limit: usize,
) -> deps_core::ecosystem::BoxFuture<'a, deps_core::Result<Vec<Box<dyn Metadata>>>>
{
Box::pin(async move { Ok(vec![]) })
}
fn as_any(&self) -> &dyn Any {
self
}
}
let registry = Arc::new(FallbackContextAwareRegistry {
seen_minimum_stability: Mutex::new(Vec::new()),
});
let packages = vec![PackageName::new("vendor/pkg")];
let result = fetch_latest_versions_parallel(
Arc::clone(®istry) as Arc<dyn Registry>,
with_registry_source(packages),
&HashMap::new(),
None,
deps_core::freshness::FreshnessSettings::default(),
10,
10,
Some("beta"),
)
.await;
assert_eq!(
*registry
.seen_minimum_stability
.lock()
.unwrap_or_else(|p| p.into_inner()),
vec![Some("beta".to_string())],
"get_latest_matching_with_context must receive the caller's minimum_stability"
);
let widget = result
.versions
.get("vendor/pkg")
.expect("fallback pick should succeed");
assert_eq!(widget.latest, "2.0.0-beta1");
}
#[tokio::test]
async fn test_fetch_falls_back_to_get_latest_matching_when_list_based_pick_finds_nothing() {
use deps_core::{Metadata, Registry, Version};
use std::any::Any;
#[derive(Debug)]
struct MockVersion {
version: ConcreteVersion,
}
impl Version for MockVersion {
fn version_string(&self) -> &ConcreteVersion {
&self.version
}
fn as_any(&self) -> &dyn Any {
self
}
}
struct UntaggedModuleRegistry;
impl Registry for UntaggedModuleRegistry {
fn get_versions<'a>(
&'a self,
_name: &'a deps_core::PackageName,
) -> deps_core::ecosystem::BoxFuture<'a, deps_core::Result<Vec<Box<dyn Version>>>>
{
Box::pin(async move { Ok(vec![]) })
}
fn get_latest_matching<'a>(
&'a self,
_name: &'a deps_core::PackageName,
_req: &'a deps_core::VersionReq,
) -> deps_core::ecosystem::BoxFuture<'a, deps_core::Result<Option<Box<dyn Version>>>>
{
Box::pin(async move {
Ok(Some(Box::new(MockVersion {
version: "v0.0.0-20191109021931-daa7c04131f5".into(),
}) as Box<dyn Version>))
})
}
fn search_raw<'a>(
&'a self,
_query: &'a str,
_limit: usize,
) -> deps_core::ecosystem::BoxFuture<'a, deps_core::Result<Vec<Box<dyn Metadata>>>>
{
Box::pin(async move { Ok(vec![]) })
}
fn as_any(&self) -> &dyn Any {
self
}
}
let registry: Arc<dyn Registry> = Arc::new(UntaggedModuleRegistry);
let packages = vec![PackageName::new("golang.org/x/exp")];
let result = fetch_latest_versions_parallel(
registry,
with_registry_source(packages),
&HashMap::new(),
None,
deps_core::freshness::FreshnessSettings::default(),
5,
10,
None,
)
.await;
assert_eq!(
result
.versions
.get("golang.org/x/exp")
.map(|v| v.latest.as_str()),
Some("v0.0.0-20191109021931-daa7c04131f5"),
"must fall back to get_latest_matching instead of reporting no version found"
);
}
#[tokio::test]
async fn test_fetch_registry_error_handled() {
use deps_core::{Metadata, Registry, Version};
use std::any::Any;
struct ErrorRegistry;
impl Registry for ErrorRegistry {
fn get_versions<'a>(
&'a self,
name: &'a deps_core::PackageName,
) -> deps_core::ecosystem::BoxFuture<'a, deps_core::Result<Vec<Box<dyn Version>>>>
{
Box::pin(async move {
Err(deps_core::error::DepsError::CacheError(format!(
"Failed to fetch package: {}",
name.as_str()
)))
})
}
fn get_latest_matching<'a>(
&'a self,
name: &'a deps_core::PackageName,
_req: &'a deps_core::VersionReq,
) -> deps_core::ecosystem::BoxFuture<'a, deps_core::Result<Option<Box<dyn Version>>>>
{
Box::pin(async move {
Err(deps_core::error::DepsError::CacheError(format!(
"Failed to fetch package: {}",
name.as_str()
)))
})
}
fn search_raw<'a>(
&'a self,
_query: &'a str,
_limit: usize,
) -> deps_core::ecosystem::BoxFuture<'a, deps_core::Result<Vec<Box<dyn Metadata>>>>
{
Box::pin(async move { Ok(vec![]) })
}
fn as_any(&self) -> &dyn Any {
self
}
}
let registry: Arc<dyn Registry> = Arc::new(ErrorRegistry);
let packages = vec![
PackageName::new("package-1"),
PackageName::new("package-2"),
PackageName::new("package-3"),
];
let result = fetch_latest_versions_parallel(
registry,
with_registry_source(packages),
&HashMap::new(),
None,
deps_core::freshness::FreshnessSettings::default(),
5,
10,
None,
)
.await;
assert!(
result.versions.is_empty(),
"All packages with errors should be omitted from results"
);
assert_eq!(
result.failed_count, 3,
"All 3 packages should be marked as failed"
);
assert_eq!(
result.fetch_failed,
HashMap::from([
(PackageName::new("package-1"), FetchFailure::Transient),
(PackageName::new("package-2"), FetchFailure::Transient),
(PackageName::new("package-3"), FetchFailure::Transient),
]),
"every errored package must be recorded in fetch_failed"
);
}
#[tokio::test]
async fn test_fetch_failed_log_redacts_credential_shaped_package_name() {
use deps_core::{Metadata, Registry, Version};
use std::any::Any;
#[tracing::instrument(skip_all, fields(package = %name.for_tracing()), level = "debug")]
async fn inner_fetch(name: &PackageName) -> deps_core::Result<Vec<Box<dyn Version>>> {
tracing::debug!("mock registry fetch invoked");
Err(deps_core::error::DepsError::CacheError(
"transient backend failure".to_string(),
))
}
struct AlwaysFailsRegistry;
impl Registry for AlwaysFailsRegistry {
fn get_versions<'a>(
&'a self,
name: &'a deps_core::PackageName,
) -> deps_core::ecosystem::BoxFuture<'a, deps_core::Result<Vec<Box<dyn Version>>>>
{
Box::pin(inner_fetch(name))
}
fn get_latest_matching<'a>(
&'a self,
_name: &'a deps_core::PackageName,
_req: &'a deps_core::VersionReq,
) -> deps_core::ecosystem::BoxFuture<'a, deps_core::Result<Option<Box<dyn Version>>>>
{
Box::pin(async move {
Err(deps_core::error::DepsError::CacheError(
"transient backend failure".to_string(),
))
})
}
fn search_raw<'a>(
&'a self,
_query: &'a str,
_limit: usize,
) -> deps_core::ecosystem::BoxFuture<'a, deps_core::Result<Vec<Box<dyn Metadata>>>>
{
Box::pin(async move { Ok(vec![]) })
}
fn as_any(&self) -> &dyn Any {
self
}
}
let sentinel_name =
PackageName::new("com.example:deploy:AUDITSENTINEL0000@git.internal.corp");
let registry: Arc<dyn Registry> = Arc::new(AlwaysFailsRegistry);
let packages = vec![sentinel_name.clone()];
let log =
deps_core::test_util::capture_tracing_output_async_at(tracing::Level::DEBUG, async {
let result = fetch_latest_versions_parallel(
registry,
with_registry_source(packages),
&HashMap::new(),
None,
deps_core::freshness::FreshnessSettings::default(),
5,
10,
None,
)
.await;
assert_eq!(result.failed_count, 1);
})
.await;
assert!(
log.contains("fetch failed"),
"expected the fetch-failed WARN to fire: {log:?}"
);
assert!(
log.contains("mock registry fetch invoked"),
"expected the in-span event to fire — without it the span's fields never render, \
silently downgrading this test back to event-field-only coverage: {log:?}"
);
assert!(
!log.contains("AUDITSENTINEL0000"),
"tracing output leaked a credential-shaped package name: {log:?}"
);
assert!(
log.contains("git.internal.corp"),
"host should survive redaction: {log:?}"
);
}
#[tokio::test]
async fn test_fetch_not_found_is_not_recorded_as_fetch_failed() {
use deps_core::{Metadata, Registry, Version};
use std::any::Any;
struct NotFoundRegistry;
impl Registry for NotFoundRegistry {
fn get_versions<'a>(
&'a self,
name: &'a deps_core::PackageName,
) -> deps_core::ecosystem::BoxFuture<'a, deps_core::Result<Vec<Box<dyn Version>>>>
{
Box::pin(async move {
Err(deps_core::error::DepsError::PackageNotFound {
package: name.as_str().into(),
registry: "mock",
})
})
}
fn get_latest_matching<'a>(
&'a self,
name: &'a deps_core::PackageName,
_req: &'a deps_core::VersionReq,
) -> deps_core::ecosystem::BoxFuture<'a, deps_core::Result<Option<Box<dyn Version>>>>
{
Box::pin(async move {
Err(deps_core::error::DepsError::PackageNotFound {
package: name.as_str().into(),
registry: "mock",
})
})
}
fn search_raw<'a>(
&'a self,
_query: &'a str,
_limit: usize,
) -> deps_core::ecosystem::BoxFuture<'a, deps_core::Result<Vec<Box<dyn Metadata>>>>
{
Box::pin(async move { Ok(vec![]) })
}
fn as_any(&self) -> &dyn Any {
self
}
}
let registry: Arc<dyn Registry> = Arc::new(NotFoundRegistry);
let packages = vec![PackageName::new("typo-pkg")];
let result = fetch_latest_versions_parallel(
registry,
with_registry_source(packages),
&HashMap::new(),
None,
deps_core::freshness::FreshnessSettings::default(),
5,
10,
None,
)
.await;
assert!(result.versions.is_empty());
assert!(
result.fetch_failed.is_empty(),
"a genuine not-found must not be recorded in fetch_failed, or \
generate_diagnostics_from_cache would report it as a registry \
error instead of Unknown package"
);
}
#[tokio::test]
async fn test_fetch_success_with_zero_versions_is_recorded_as_no_comparable_versions() {
use deps_core::{Metadata, Registry, Version};
use std::any::Any;
struct EmptyButRealRegistry;
impl Registry for EmptyButRealRegistry {
fn get_versions<'a>(
&'a self,
_name: &'a deps_core::PackageName,
) -> deps_core::ecosystem::BoxFuture<'a, deps_core::Result<Vec<Box<dyn Version>>>>
{
Box::pin(async move { Ok(vec![]) })
}
fn get_latest_matching<'a>(
&'a self,
_name: &'a deps_core::PackageName,
_req: &'a deps_core::VersionReq,
) -> deps_core::ecosystem::BoxFuture<'a, deps_core::Result<Option<Box<dyn Version>>>>
{
Box::pin(async move { Ok(None) })
}
fn search_raw<'a>(
&'a self,
_query: &'a str,
_limit: usize,
) -> deps_core::ecosystem::BoxFuture<'a, deps_core::Result<Vec<Box<dyn Metadata>>>>
{
Box::pin(async move { Ok(vec![]) })
}
fn as_any(&self) -> &dyn Any {
self
}
}
let registry: Arc<dyn Registry> = Arc::new(EmptyButRealRegistry);
let packages = vec![PackageName::new("dtolnay/rust-toolchain")];
let result = fetch_latest_versions_parallel(
registry,
with_registry_source(packages),
&HashMap::new(),
None,
deps_core::freshness::FreshnessSettings::default(),
5,
10,
None,
)
.await;
assert!(result.versions.is_empty());
assert!(
result.fetch_failed.is_empty(),
"a genuine empty-but-successful fetch must not be recorded as a fetch \
failure, or generate_diagnostics_from_cache would report a registry \
error instead of nothing"
);
assert!(
result
.no_comparable_versions
.contains(&PackageName::new("dtolnay/rust-toolchain")),
"a package whose fetch succeeded with zero comparable versions must be \
recorded in no_comparable_versions, or R5 would misreport it as Unknown \
package; got: {:?}",
result.no_comparable_versions
);
}
#[tokio::test]
async fn test_fetch_http_404_is_not_recorded_as_fetch_failed() {
use deps_core::{Metadata, Registry, Version};
use std::any::Any;
struct Http404Registry;
impl Registry for Http404Registry {
fn get_versions<'a>(
&'a self,
name: &'a deps_core::PackageName,
) -> deps_core::ecosystem::BoxFuture<'a, deps_core::Result<Vec<Box<dyn Version>>>>
{
Box::pin(async move {
Err(deps_core::error::DepsError::HttpStatus {
url: format!("https://example.com/{}", name.as_str()).into(),
status: 404,
})
})
}
fn get_latest_matching<'a>(
&'a self,
name: &'a deps_core::PackageName,
_req: &'a deps_core::VersionReq,
) -> deps_core::ecosystem::BoxFuture<'a, deps_core::Result<Option<Box<dyn Version>>>>
{
Box::pin(async move {
Err(deps_core::error::DepsError::HttpStatus {
url: format!("https://example.com/{}", name.as_str()).into(),
status: 404,
})
})
}
fn search_raw<'a>(
&'a self,
_query: &'a str,
_limit: usize,
) -> deps_core::ecosystem::BoxFuture<'a, deps_core::Result<Vec<Box<dyn Metadata>>>>
{
Box::pin(async move { Ok(vec![]) })
}
fn as_any(&self) -> &dyn Any {
self
}
}
let registry: Arc<dyn Registry> = Arc::new(Http404Registry);
let packages = vec![PackageName::new("typo-pkg")];
let result = fetch_latest_versions_parallel(
registry,
with_registry_source(packages),
&HashMap::new(),
None,
deps_core::freshness::FreshnessSettings::default(),
5,
10,
None,
)
.await;
assert!(result.versions.is_empty());
assert!(
result.fetch_failed.is_empty(),
"a bare HTTP 404 must not be recorded in fetch_failed either"
);
}
#[tokio::test]
async fn test_fetch_fallback_error_recorded_as_fetch_failed_unless_not_found() {
use deps_core::{Metadata, Registry, Version};
use std::any::Any;
struct FallbackErrorRegistry;
impl Registry for FallbackErrorRegistry {
fn get_versions<'a>(
&'a self,
_name: &'a deps_core::PackageName,
) -> deps_core::ecosystem::BoxFuture<'a, deps_core::Result<Vec<Box<dyn Version>>>>
{
Box::pin(async move { Ok(vec![]) })
}
fn get_latest_matching<'a>(
&'a self,
name: &'a deps_core::PackageName,
_req: &'a deps_core::VersionReq,
) -> deps_core::ecosystem::BoxFuture<'a, deps_core::Result<Option<Box<dyn Version>>>>
{
let name = name.clone();
Box::pin(async move {
if name.as_str() == "not-found" {
Err(deps_core::error::DepsError::PackageNotFound {
package: name.as_str().into(),
registry: "mock",
})
} else {
Err(deps_core::error::DepsError::CacheError(
"mock fallback failure".to_string(),
))
}
})
}
fn search_raw<'a>(
&'a self,
_query: &'a str,
_limit: usize,
) -> deps_core::ecosystem::BoxFuture<'a, deps_core::Result<Vec<Box<dyn Metadata>>>>
{
Box::pin(async move { Ok(vec![]) })
}
fn as_any(&self) -> &dyn Any {
self
}
}
let registry: Arc<dyn Registry> = Arc::new(FallbackErrorRegistry);
let packages = vec![PackageName::new("flaky"), PackageName::new("not-found")];
let result = fetch_latest_versions_parallel(
registry,
with_registry_source(packages),
&HashMap::new(),
None,
deps_core::freshness::FreshnessSettings::default(),
5,
10,
None,
)
.await;
assert!(result.versions.is_empty());
assert_eq!(
result.fetch_failed,
HashMap::from([(PackageName::new("flaky"), FetchFailure::Transient)]),
"the fallback's own non-not-found error must be recorded in fetch_failed, \
but its not-found error must not"
);
assert_eq!(
result.failed_count, 2,
"both fallback failures count toward failed_count regardless of cause (S2)"
);
}
#[tokio::test]
async fn test_fetch_fallback_timeout_recorded_as_fetch_failed() {
use deps_core::{Metadata, Registry, Version};
use std::any::Any;
use std::time::Duration;
struct FallbackTimeoutRegistry;
impl Registry for FallbackTimeoutRegistry {
fn get_versions<'a>(
&'a self,
_name: &'a deps_core::PackageName,
) -> deps_core::ecosystem::BoxFuture<'a, deps_core::Result<Vec<Box<dyn Version>>>>
{
Box::pin(async move { Ok(vec![]) })
}
fn get_latest_matching<'a>(
&'a self,
_name: &'a deps_core::PackageName,
_req: &'a deps_core::VersionReq,
) -> deps_core::ecosystem::BoxFuture<'a, deps_core::Result<Option<Box<dyn Version>>>>
{
Box::pin(async move {
tokio::time::sleep(Duration::from_secs(10)).await;
Ok(None)
})
}
fn search_raw<'a>(
&'a self,
_query: &'a str,
_limit: usize,
) -> deps_core::ecosystem::BoxFuture<'a, deps_core::Result<Vec<Box<dyn Metadata>>>>
{
Box::pin(async move { Ok(vec![]) })
}
fn as_any(&self) -> &dyn Any {
self
}
}
let registry: Arc<dyn Registry> = Arc::new(FallbackTimeoutRegistry);
let packages = vec![PackageName::new("slow-fallback")];
let result = fetch_latest_versions_parallel(
registry,
with_registry_source(packages),
&HashMap::new(),
None,
deps_core::freshness::FreshnessSettings::default(),
1,
10,
None,
)
.await;
assert!(result.versions.is_empty());
assert_eq!(
result.fetch_failed,
HashMap::from([(PackageName::new("slow-fallback"), FetchFailure::Transient)])
);
assert_eq!(result.failed_count, 1);
}
#[tokio::test]
async fn test_first_error_prefers_actionable_error_over_not_found_regardless_of_race_order() {
use deps_core::{Metadata, Registry, Version};
use std::any::Any;
use std::time::Duration;
struct MixedErrorRegistry;
impl Registry for MixedErrorRegistry {
fn get_versions<'a>(
&'a self,
name: &'a deps_core::PackageName,
) -> deps_core::ecosystem::BoxFuture<'a, deps_core::Result<Vec<Box<dyn Version>>>>
{
Box::pin(async move {
if name.as_str() == "typo-pkg" {
Err(deps_core::error::DepsError::PackageNotFound {
package: name.as_str().into(),
registry: "mock",
})
} else {
tokio::time::sleep(Duration::from_millis(50)).await;
Err(deps_core::error::DepsError::CacheError(
"rate limit exceeded".to_string(),
))
}
})
}
fn get_latest_matching<'a>(
&'a self,
_name: &'a deps_core::PackageName,
_req: &'a deps_core::VersionReq,
) -> deps_core::ecosystem::BoxFuture<'a, deps_core::Result<Option<Box<dyn Version>>>>
{
Box::pin(async move { Ok(None) })
}
fn search_raw<'a>(
&'a self,
_query: &'a str,
_limit: usize,
) -> deps_core::ecosystem::BoxFuture<'a, deps_core::Result<Vec<Box<dyn Metadata>>>>
{
Box::pin(async move { Ok(vec![]) })
}
fn as_any(&self) -> &dyn Any {
self
}
}
let registry: Arc<dyn Registry> = Arc::new(MixedErrorRegistry);
let packages = vec![
PackageName::new("typo-pkg"),
PackageName::new("rate-limited"),
];
let result = fetch_latest_versions_parallel(
registry,
with_registry_source(packages),
&HashMap::new(),
None,
deps_core::freshness::FreshnessSettings::default(),
5,
10,
None,
)
.await;
let err = result
.first_error
.expect("an actionable failure occurred and must be reported");
assert!(
err.contains("rate limit exceeded"),
"the actionable error must win the toast over the faster-finishing not-found, \
got: {err}"
);
assert!(
!err.contains("not found"),
"a not-found error must never outrank an actionable error, got: {err}"
);
}
#[tokio::test]
async fn test_first_error_falls_back_to_not_found_when_no_actionable_error_occurred() {
use deps_core::{Metadata, Registry, Version};
use std::any::Any;
struct AllNotFoundRegistry;
impl Registry for AllNotFoundRegistry {
fn get_versions<'a>(
&'a self,
name: &'a deps_core::PackageName,
) -> deps_core::ecosystem::BoxFuture<'a, deps_core::Result<Vec<Box<dyn Version>>>>
{
Box::pin(async move {
Err(deps_core::error::DepsError::PackageNotFound {
package: name.as_str().into(),
registry: "mock",
})
})
}
fn get_latest_matching<'a>(
&'a self,
_name: &'a deps_core::PackageName,
_req: &'a deps_core::VersionReq,
) -> deps_core::ecosystem::BoxFuture<'a, deps_core::Result<Option<Box<dyn Version>>>>
{
Box::pin(async move { Ok(None) })
}
fn search_raw<'a>(
&'a self,
_query: &'a str,
_limit: usize,
) -> deps_core::ecosystem::BoxFuture<'a, deps_core::Result<Vec<Box<dyn Metadata>>>>
{
Box::pin(async move { Ok(vec![]) })
}
fn as_any(&self) -> &dyn Any {
self
}
}
let registry: Arc<dyn Registry> = Arc::new(AllNotFoundRegistry);
let packages = vec![
PackageName::new("typo-pkg-1"),
PackageName::new("typo-pkg-2"),
];
let result = fetch_latest_versions_parallel(
registry,
with_registry_source(packages),
&HashMap::new(),
None,
deps_core::freshness::FreshnessSettings::default(),
5,
10,
None,
)
.await;
assert!(
result.fetch_failed.is_empty(),
"not-found errors must never be recorded in fetch_failed"
);
let err = result
.first_error
.expect("a not-found-only batch must still fall back to reporting one via first_error");
assert!(err.contains("not found"), "got: {err}");
}
#[tokio::test]
async fn test_timeout_only_batch_reports_first_error_alongside_failed_count() {
use deps_core::{Metadata, Registry, Version};
use std::any::Any;
use std::time::Duration;
struct AlwaysTimesOutRegistry;
impl Registry for AlwaysTimesOutRegistry {
fn get_versions<'a>(
&'a self,
_name: &'a deps_core::PackageName,
) -> deps_core::ecosystem::BoxFuture<'a, deps_core::Result<Vec<Box<dyn Version>>>>
{
Box::pin(async move {
tokio::time::sleep(Duration::from_secs(10)).await;
Ok(vec![])
})
}
fn get_latest_matching<'a>(
&'a self,
_name: &'a deps_core::PackageName,
_req: &'a deps_core::VersionReq,
) -> deps_core::ecosystem::BoxFuture<'a, deps_core::Result<Option<Box<dyn Version>>>>
{
Box::pin(async move {
tokio::time::sleep(Duration::from_secs(10)).await;
Ok(None)
})
}
fn search_raw<'a>(
&'a self,
_query: &'a str,
_limit: usize,
) -> deps_core::ecosystem::BoxFuture<'a, deps_core::Result<Vec<Box<dyn Metadata>>>>
{
Box::pin(async move { Ok(vec![]) })
}
fn as_any(&self) -> &dyn Any {
self
}
}
let registry: Arc<dyn Registry> = Arc::new(AlwaysTimesOutRegistry);
let packages = vec![
PackageName::new("slow-1"),
PackageName::new("slow-2"),
PackageName::new("slow-3"),
];
let result = fetch_latest_versions_parallel(
registry,
with_registry_source(packages),
&HashMap::new(),
None,
deps_core::freshness::FreshnessSettings::default(),
1,
10,
None,
)
.await;
assert_eq!(
result.failed_count, 3,
"all 3 packages must count toward failed_count"
);
let err = result
.first_error
.expect("a timeout is actionable and must populate first_error, not just failed_count");
assert!(
err.contains("timed out"),
"first_error must be the actionable timeout message, got: {err}"
);
}
#[cfg(feature = "composer")]
mod composer_tests {
use super::*;
#[tokio::test]
async fn test_composer_minimum_stability_extracts_from_real_parse_result() {
let json = r#"{
"minimum-stability": "beta",
"require": {
"symfony/console": "^6.0"
}
}"#;
let uri = deps_core::test_util::test_uri("/test/composer.json");
let parse_result = crate::setup::parse_composer_json(json, &uri).unwrap();
assert_eq!(
composer_minimum_stability(&parse_result as &dyn deps_core::ParseResult),
Some("beta".to_string())
);
}
#[tokio::test]
async fn test_composer_minimum_stability_none_when_absent() {
let json = r#"{"require": {"symfony/console": "^6.0"}}"#;
let uri = deps_core::test_util::test_uri("/test/composer.json");
let parse_result = crate::setup::parse_composer_json(json, &uri).unwrap();
assert_eq!(
composer_minimum_stability(&parse_result as &dyn deps_core::ParseResult),
None
);
}
#[test]
fn test_composer_minimum_stability_none_for_non_composer_parse_result() {
struct OtherParseResult;
impl deps_core::ParseResult for OtherParseResult {
fn dependencies(&self) -> Vec<&dyn deps_core::Dependency> {
vec![]
}
fn workspace_root(&self) -> Option<&std::path::Path> {
None
}
fn uri(&self) -> &url::Url {
unimplemented!("not exercised by this test")
}
fn as_any(&self) -> &dyn std::any::Any {
self
}
}
assert_eq!(
composer_minimum_stability(&OtherParseResult as &dyn deps_core::ParseResult),
None
);
}
}
mod yanked_check_tests {
use super::*;
use deps_core::{Metadata, Version};
use std::any::Any;
use std::sync::atomic::{AtomicUsize, Ordering};
#[derive(Debug, Clone)]
struct MockYankVersion {
version: ConcreteVersion,
yanked: bool,
}
impl Version for MockYankVersion {
fn version_string(&self) -> &ConcreteVersion {
&self.version
}
fn removal_status(&self) -> deps_core::RemovalStatus {
deps_core::RemovalStatus::from_yanked(self.yanked)
}
fn as_any(&self) -> &dyn Any {
self
}
}
enum FetchOutcome {
Versions(Vec<(&'static str, bool)>),
Error,
Timeout,
}
struct MockRegistry {
reports_yanked: bool,
versions: HashMap<&'static str, FetchOutcome>,
latest_fallback: HashMap<&'static str, (&'static str, bool)>,
fetch_calls: Arc<AtomicUsize>,
}
impl Registry for MockRegistry {
fn get_versions<'a>(
&'a self,
name: &'a PackageName,
) -> deps_core::ecosystem::BoxFuture<'a, deps_core::Result<Vec<Box<dyn Version>>>>
{
self.fetch_calls.fetch_add(1, Ordering::Relaxed);
let outcome = self.versions.get(name.as_str());
Box::pin(async move {
match outcome {
Some(FetchOutcome::Versions(vs)) => Ok(vs
.iter()
.map(|(v, y)| {
Box::new(MockYankVersion {
version: (*v).into(),
yanked: *y,
}) as Box<dyn Version>
})
.collect()),
Some(FetchOutcome::Error) => Err(deps_core::error::DepsError::CacheError(
"mock fetch error".to_string(),
)),
Some(FetchOutcome::Timeout) => {
tokio::time::sleep(std::time::Duration::from_secs(10)).await;
Ok(vec![])
}
None => Ok(vec![]),
}
})
}
fn select_latest_matching(
&self,
versions: &[Box<dyn Version>],
_req: &VersionReq,
) -> Option<usize> {
versions
.iter()
.position(|v| !v.removal_status().blocks_resolution())
}
fn get_latest_matching<'a>(
&'a self,
name: &'a PackageName,
_req: &'a VersionReq,
) -> deps_core::ecosystem::BoxFuture<'a, deps_core::Result<Option<Box<dyn Version>>>>
{
let outcome = self.latest_fallback.get(name.as_str()).copied();
Box::pin(async move {
Ok(outcome.map(|(v, y)| {
Box::new(MockYankVersion {
version: v.into(),
yanked: y,
}) as Box<dyn Version>
}))
})
}
fn search_raw<'a>(
&'a self,
_query: &'a str,
_limit: usize,
) -> deps_core::ecosystem::BoxFuture<'a, deps_core::Result<Vec<Box<dyn Metadata>>>>
{
Box::pin(async move { Ok(vec![]) })
}
fn reports_yanked(&self) -> bool {
self.reports_yanked
}
fn as_any(&self) -> &dyn Any {
self
}
}
#[tokio::test]
async fn reports_yanked_false_never_recorded() {
let fetch_calls = Arc::new(AtomicUsize::new(0));
let registry: Arc<dyn Registry> = Arc::new(MockRegistry {
reports_yanked: false,
versions: HashMap::from([(
"pkg",
FetchOutcome::Versions(vec![("2.0.0", false), ("1.0.0", true)]),
)]),
latest_fallback: HashMap::new(),
fetch_calls: Arc::clone(&fetch_calls),
});
let mut in_use = HashMap::new();
in_use.insert(PackageName::new("pkg"), vec!["1.0.0".to_string()]);
let result = fetch_latest_versions_parallel(
registry,
vec![(PackageName::new("pkg"), DependencySource::Registry)],
&in_use,
None,
deps_core::freshness::FreshnessSettings::default(),
5,
10,
None,
)
.await;
assert_eq!(fetch_calls.load(Ordering::Relaxed), 1);
assert!(result.yanked_versions.is_empty());
}
#[tokio::test]
async fn in_use_equal_to_latest_not_yanked() {
let fetch_calls = Arc::new(AtomicUsize::new(0));
let registry: Arc<dyn Registry> = Arc::new(MockRegistry {
reports_yanked: true,
versions: HashMap::from([("pkg", FetchOutcome::Versions(vec![("1.0.0", false)]))]),
latest_fallback: HashMap::new(),
fetch_calls: Arc::clone(&fetch_calls),
});
let mut in_use = HashMap::new();
in_use.insert(PackageName::new("pkg"), vec!["1.0.0".to_string()]);
let result = fetch_latest_versions_parallel(
registry,
vec![(PackageName::new("pkg"), DependencySource::Registry)],
&in_use,
None,
deps_core::freshness::FreshnessSettings::default(),
5,
10,
None,
)
.await;
assert_eq!(fetch_calls.load(Ordering::Relaxed), 1);
assert!(result.yanked_versions.is_empty());
}
#[tokio::test]
async fn no_known_in_use_version_skips_the_check() {
let fetch_calls = Arc::new(AtomicUsize::new(0));
let registry: Arc<dyn Registry> = Arc::new(MockRegistry {
reports_yanked: true,
versions: HashMap::from([("pkg", FetchOutcome::Versions(vec![("2.0.0", false)]))]),
latest_fallback: HashMap::new(),
fetch_calls: Arc::clone(&fetch_calls),
});
let result = fetch_latest_versions_parallel(
registry,
vec![(PackageName::new("pkg"), DependencySource::Registry)],
&HashMap::new(),
None,
deps_core::freshness::FreshnessSettings::default(),
5,
10,
None,
)
.await;
assert_eq!(fetch_calls.load(Ordering::Relaxed), 1);
assert!(result.yanked_versions.is_empty());
}
#[tokio::test]
async fn in_use_differs_and_yanked_is_recorded() {
let fetch_calls = Arc::new(AtomicUsize::new(0));
let registry: Arc<dyn Registry> = Arc::new(MockRegistry {
reports_yanked: true,
versions: HashMap::from([(
"pkg",
FetchOutcome::Versions(vec![("2.0.0", false), ("1.0.0", true)]),
)]),
latest_fallback: HashMap::new(),
fetch_calls: Arc::clone(&fetch_calls),
});
let mut in_use = HashMap::new();
in_use.insert(PackageName::new("pkg"), vec!["1.0.0".to_string()]);
let result = fetch_latest_versions_parallel(
registry,
vec![(PackageName::new("pkg"), DependencySource::Registry)],
&in_use,
None,
deps_core::freshness::FreshnessSettings::default(),
5,
10,
None,
)
.await;
assert_eq!(fetch_calls.load(Ordering::Relaxed), 1);
assert_eq!(
result.yanked_versions.get(&PackageName::new("pkg")),
Some(&(ConcreteVersion::new("1.0.0"), RemovalStatus::Yanked))
);
}
#[tokio::test]
async fn in_use_differs_and_not_yanked_is_not_recorded() {
let fetch_calls = Arc::new(AtomicUsize::new(0));
let registry: Arc<dyn Registry> = Arc::new(MockRegistry {
reports_yanked: true,
versions: HashMap::from([(
"pkg",
FetchOutcome::Versions(vec![("2.0.0", false), ("1.0.0", false)]),
)]),
latest_fallback: HashMap::new(),
fetch_calls: Arc::clone(&fetch_calls),
});
let mut in_use = HashMap::new();
in_use.insert(PackageName::new("pkg"), vec!["1.0.0".to_string()]);
let result = fetch_latest_versions_parallel(
registry,
vec![(PackageName::new("pkg"), DependencySource::Registry)],
&in_use,
None,
deps_core::freshness::FreshnessSettings::default(),
5,
10,
None,
)
.await;
assert_eq!(fetch_calls.load(Ordering::Relaxed), 1);
assert!(result.yanked_versions.is_empty());
}
#[tokio::test]
async fn every_version_yanked_still_checks_in_use() {
let fetch_calls = Arc::new(AtomicUsize::new(0));
let registry: Arc<dyn Registry> = Arc::new(MockRegistry {
reports_yanked: true,
versions: HashMap::from([("pkg", FetchOutcome::Versions(vec![("1.0.0", true)]))]),
latest_fallback: HashMap::new(),
fetch_calls: Arc::clone(&fetch_calls),
});
let mut in_use = HashMap::new();
in_use.insert(PackageName::new("pkg"), vec!["1.0.0".to_string()]);
let result = fetch_latest_versions_parallel(
registry,
vec![(PackageName::new("pkg"), DependencySource::Registry)],
&in_use,
None,
deps_core::freshness::FreshnessSettings::default(),
5,
10,
None,
)
.await;
assert_eq!(
result.yanked_versions.get(&PackageName::new("pkg")),
Some(&(ConcreteVersion::new("1.0.0"), RemovalStatus::Yanked))
);
assert!(result.versions.is_empty());
}
#[tokio::test]
async fn latest_pick_needs_fallback_in_use_yanked_still_found() {
let fetch_calls = Arc::new(AtomicUsize::new(0));
let registry: Arc<dyn Registry> = Arc::new(MockRegistry {
reports_yanked: true,
versions: HashMap::from([("pkg", FetchOutcome::Versions(vec![("1.0.0", true)]))]),
latest_fallback: HashMap::from([("pkg", ("2.0.0", false))]),
fetch_calls: Arc::clone(&fetch_calls),
});
let mut in_use = HashMap::new();
in_use.insert(PackageName::new("pkg"), vec!["1.0.0".to_string()]);
let result = fetch_latest_versions_parallel(
registry,
vec![(PackageName::new("pkg"), DependencySource::Registry)],
&in_use,
None,
deps_core::freshness::FreshnessSettings::default(),
5,
10,
None,
)
.await;
assert_eq!(fetch_calls.load(Ordering::Relaxed), 1);
assert_eq!(
result
.versions
.get(&PackageName::new("pkg"))
.map(|v| v.latest.as_str()),
Some("2.0.0")
);
assert_eq!(
result.yanked_versions.get(&PackageName::new("pkg")),
Some(&(ConcreteVersion::new("1.0.0"), RemovalStatus::Yanked))
);
}
#[tokio::test]
async fn in_use_checks_every_occurrence_of_a_duplicate_name() {
let registry: Arc<dyn Registry> = Arc::new(MockRegistry {
reports_yanked: true,
versions: HashMap::from([(
"pkg",
FetchOutcome::Versions(vec![
("3.0.0", false),
("2.0.0", true),
("1.0.0", false),
]),
)]),
latest_fallback: HashMap::new(),
fetch_calls: Arc::new(AtomicUsize::new(0)),
});
let mut in_use = HashMap::new();
in_use.insert(
PackageName::new("pkg"),
vec!["1.0.0".to_string(), "2.0.0".to_string()],
);
let result = fetch_latest_versions_parallel(
registry,
vec![(PackageName::new("pkg"), DependencySource::Registry)],
&in_use,
None,
deps_core::freshness::FreshnessSettings::default(),
5,
10,
None,
)
.await;
assert_eq!(
result.yanked_versions.get(&PackageName::new("pkg")),
Some(&(ConcreteVersion::new("2.0.0"), RemovalStatus::Yanked)),
"the yanked occurrence must be found even though a name-keyed \
single-value map could have kept only the non-yanked \"1.0.0\" pin"
);
}
#[tokio::test]
async fn latest_is_yanked_recorded_as_defense_in_depth() {
let registry: Arc<dyn Registry> = Arc::new(MockRegistry {
reports_yanked: true,
versions: HashMap::from([("pkg", FetchOutcome::Versions(vec![("1.0.0", true)]))]),
latest_fallback: HashMap::from([("pkg", ("1.0.0", true))]),
fetch_calls: Arc::new(AtomicUsize::new(0)),
});
let result = fetch_latest_versions_parallel(
registry,
vec![(PackageName::new("pkg"), DependencySource::Registry)],
&HashMap::new(),
None,
deps_core::freshness::FreshnessSettings::default(),
5,
10,
None,
)
.await;
assert_eq!(
result.yanked_versions.get(&PackageName::new("pkg")),
Some(&(ConcreteVersion::new("1.0.0"), RemovalStatus::Yanked))
);
}
#[tokio::test]
async fn latest_is_yanked_not_recorded_when_reports_yanked_false() {
let registry: Arc<dyn Registry> = Arc::new(MockRegistry {
reports_yanked: false,
versions: HashMap::from([("pkg", FetchOutcome::Versions(vec![("1.0.0", true)]))]),
latest_fallback: HashMap::from([("pkg", ("1.0.0", true))]),
fetch_calls: Arc::new(AtomicUsize::new(0)),
});
let result = fetch_latest_versions_parallel(
registry,
vec![(PackageName::new("pkg"), DependencySource::Registry)],
&HashMap::new(),
None,
deps_core::freshness::FreshnessSettings::default(),
5,
10,
None,
)
.await;
assert!(
result.yanked_versions.is_empty(),
"a `reports_yanked() == false` registry's `removal_status()` must never be \
trusted, even on the zero-cost row-1 path"
);
assert!(
result
.versions
.get(&PackageName::new("pkg"))
.expect("pkg was fetched")
.yanked
.is_empty(),
"`PackageVersions::yanked` must stay empty for a `reports_yanked() == false` \
registry, even though the fetched version is itself flagged"
);
}
#[tokio::test]
async fn primary_fetch_error_counts_as_failed_no_yanked_data() {
let registry: Arc<dyn Registry> = Arc::new(MockRegistry {
reports_yanked: true,
versions: HashMap::from([("pkg", FetchOutcome::Error)]),
latest_fallback: HashMap::new(),
fetch_calls: Arc::new(AtomicUsize::new(0)),
});
let mut in_use = HashMap::new();
in_use.insert(PackageName::new("pkg"), vec!["1.0.0".to_string()]);
let result = fetch_latest_versions_parallel(
registry,
vec![(PackageName::new("pkg"), DependencySource::Registry)],
&in_use,
None,
deps_core::freshness::FreshnessSettings::default(),
5,
10,
None,
)
.await;
assert!(result.yanked_versions.is_empty());
assert_eq!(result.failed_count, 1);
assert!(result.versions.is_empty());
}
#[tokio::test]
async fn primary_fetch_timeout_counts_as_failed_no_yanked_data() {
let registry: Arc<dyn Registry> = Arc::new(MockRegistry {
reports_yanked: true,
versions: HashMap::from([("pkg", FetchOutcome::Timeout)]),
latest_fallback: HashMap::new(),
fetch_calls: Arc::new(AtomicUsize::new(0)),
});
let mut in_use = HashMap::new();
in_use.insert(PackageName::new("pkg"), vec!["1.0.0".to_string()]);
let result = fetch_latest_versions_parallel(
registry,
vec![(PackageName::new("pkg"), DependencySource::Registry)],
&in_use,
None,
deps_core::freshness::FreshnessSettings::default(),
1,
10,
None,
)
.await;
assert!(result.yanked_versions.is_empty());
assert_eq!(result.failed_count, 1);
assert!(result.versions.is_empty());
}
}
mod deprecation_derivation_tests {
use super::*;
use deps_core::{Metadata, Version};
use std::any::Any;
struct MockDeprecatedVersion {
version: ConcreteVersion,
deprecation: Option<Deprecation>,
}
impl Version for MockDeprecatedVersion {
fn version_string(&self) -> &ConcreteVersion {
&self.version
}
fn removal_status(&self) -> RemovalStatus {
RemovalStatus::from_advisory(self.deprecation.is_some())
}
fn deprecation(&self) -> Option<&Deprecation> {
self.deprecation.as_ref()
}
fn as_any(&self) -> &dyn Any {
self
}
}
struct SingleVersionRegistry {
deprecation: Option<Deprecation>,
}
impl Registry for SingleVersionRegistry {
fn get_versions<'a>(
&'a self,
_name: &'a PackageName,
) -> deps_core::ecosystem::BoxFuture<'a, deps_core::Result<Vec<Box<dyn Version>>>>
{
let deprecation = self.deprecation.clone();
Box::pin(async move {
Ok(vec![Box::new(MockDeprecatedVersion {
version: "1.0.0".into(),
deprecation,
}) as Box<dyn Version>])
})
}
fn select_latest_matching(
&self,
versions: &[Box<dyn Version>],
_req: &VersionReq,
) -> Option<usize> {
(!versions.is_empty()).then_some(0)
}
fn get_latest_matching<'a>(
&'a self,
_name: &'a PackageName,
_req: &'a VersionReq,
) -> deps_core::ecosystem::BoxFuture<'a, deps_core::Result<Option<Box<dyn Version>>>>
{
Box::pin(async move { Ok(None) })
}
fn search_raw<'a>(
&'a self,
_query: &'a str,
_limit: usize,
) -> deps_core::ecosystem::BoxFuture<'a, deps_core::Result<Vec<Box<dyn Metadata>>>>
{
Box::pin(async move { Ok(vec![]) })
}
fn as_any(&self) -> &dyn Any {
self
}
}
#[tokio::test]
async fn fetch_result_carries_deprecation_from_resolved_pick() {
let registry: Arc<dyn Registry> = Arc::new(SingleVersionRegistry {
deprecation: Some(Deprecation {
reason: Some("archived".to_string()),
replacement: Some("other/pkg".to_string()),
}),
});
let result = fetch_latest_versions_parallel(
registry,
vec![(PackageName::new("pkg"), DependencySource::Registry)],
&HashMap::new(),
None,
deps_core::freshness::FreshnessSettings::default(),
5,
10,
None,
)
.await;
assert_eq!(
result.deprecations.get(&PackageName::new("pkg")),
Some(&Deprecation {
reason: Some("archived".to_string()),
replacement: Some("other/pkg".to_string()),
})
);
}
#[tokio::test]
async fn fetch_result_has_no_deprecation_when_resolved_pick_is_clean() {
let registry: Arc<dyn Registry> = Arc::new(SingleVersionRegistry { deprecation: None });
let result = fetch_latest_versions_parallel(
registry,
vec![(PackageName::new("pkg"), DependencySource::Registry)],
&HashMap::new(),
None,
deps_core::freshness::FreshnessSettings::default(),
5,
10,
None,
)
.await;
assert!(result.deprecations.is_empty());
}
}
}