Skip to main content

fetch_latest_versions_parallel

Function fetch_latest_versions_parallel 

Source
pub async fn fetch_latest_versions_parallel(
    registry: Arc<dyn Registry>,
    package_sources: DepSources,
    in_use: &HashMap<PackageName, Vec<ConcreteVersion>>,
    progress_sender: Option<ProgressSender>,
    freshness: FreshnessSettings,
    timeout_secs: u64,
    max_concurrent: usize,
    selection_context: &SelectionContext,
    gossip: Option<&HashMap<PackageName, GossipFindings>>,
) -> FetchResult
Expand description

Fetches latest versions for multiple packages in parallel with progress reporting.

Returns a FetchResult containing successfully fetched versions and failure count. Packages that fail to fetch are omitted from the versions map.

This function executes all registry requests concurrently with per-dependency timeout isolation, preventing slow packages from blocking others.

Alongside the primary fetch, checks whether the in-use version of a dependency has been yanked (#233), for registries that report yank data. Unlike the original design, this is not a second registry round trip: registry.get_versions below already fetches the full, unfiltered version list once per package (see PackageVersions), so the in-use-version check is a zero-cost in-memory search over a list already in hand, run for every dependency with a known in-use version rather than only when it differs from latest.

§Arguments

  • registry - Package registry to fetch from
  • package_names - List of package names to fetch
  • in_use - Raw dependency name -> the version(s) this project actually has (lockfile-resolved or a concrete pin) for every occurrence of that name in the manifest, checked against the fetched version list for yank status
  • progress - Optional progress tracker (will be updated after each fetch)
  • timeout_secs - Timeout for each individual package fetch (default: 10s)
  • max_concurrent - Maximum concurrent fetches (default: 20); clamped to >= 1 internally, since buffer_unordered(0) would hang forever (issue #833)

§Timeout Behavior

Each package fetch is wrapped in an individual timeout. If a package takes longer than timeout_secs to fetch, it fails fast with a warning and does NOT block other packages.

§Performance

With 50 dependencies and 100ms per request:

  • Sequential: 50 × 100ms = 5000ms
  • Parallel (no timeout): max(100ms) ≈ 150ms
  • Parallel (10s timeout, 1 slow package at 30s): max(10s) ≈ 10s

§Examples

use deps_core::parser::DependencySource;
use deps_core::{
    ConcreteVersion, Metadata, PackageName, Registry, SelectionContext, Version, VersionReq,
};
use deps_engine::classify::fetch::fetch_latest_versions_parallel;
use std::any::Any;
use std::collections::HashMap;
use std::sync::Arc;

struct SingleVersionRegistry;

#[derive(Clone)]
struct SimpleVersion {
    version: ConcreteVersion,
}
impl Version for SimpleVersion {
    fn version_string(&self) -> &ConcreteVersion {
        &self.version
    }
    fn as_any(&self) -> &dyn Any {
        self
    }
}

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>>>> {
        Box::pin(async move {
            Ok(vec![Box::new(SimpleVersion { version: "1.0.0".into() }) as Box<dyn Version>])
        })
    }

    // The default `select_latest_matching` always returns `None` (every real registry
    // overrides it with ecosystem-specific comparison), so `fetch_and_classify_package`
    // falls back to this method for its pick.
    fn get_latest_matching<'a>(
        &'a self,
        _name: &'a PackageName,
        _req: &'a VersionReq,
        _selection_context: &'a SelectionContext,
    ) -> deps_core::ecosystem::BoxFuture<'a, deps_core::Result<Option<Box<dyn Version>>>> {
        Box::pin(async move {
            Ok(Some(Box::new(SimpleVersion { version: "1.0.0".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
    }
}

#[tokio::main]
async fn main() {
    let sources = vec![(PackageName::new("time"), DependencySource::Registry)];

    let result = fetch_latest_versions_parallel(
        Arc::new(SingleVersionRegistry),
        sources,
        &HashMap::new(),
        None,
        deps_core::freshness::FreshnessSettings::default(),
        5,
        10,
        &SelectionContext::none(),
        None,
    )
    .await;

    assert_eq!(
        result.versions.get(&PackageName::new("time")).map(|v| v.latest.to_string()),
        Some("1.0.0".to_string())
    );
}