feat: Use concurrent dependency resolver

This commit is contained in:
abhisek
2025-05-15 15:45:30 +05:30
parent 9bb6bbfe6b
commit 82f21e68e9
2 changed files with 56 additions and 17 deletions
+51 -17
View File
@@ -16,6 +16,7 @@ type dependencyResolverConfig struct {
IncludeTransitiveDependencies bool IncludeTransitiveDependencies bool
TransitiveDepth int TransitiveDepth int
FailFast bool FailFast bool
MaxConcurrency int
} }
type dependencyResolver struct { type dependencyResolver struct {
@@ -25,6 +26,10 @@ type dependencyResolver struct {
} }
func newDependencyResolver(client packageregistry.Client, config dependencyResolverConfig) *dependencyResolver { func newDependencyResolver(client packageregistry.Client, config dependencyResolverConfig) *dependencyResolver {
if config.MaxConcurrency <= 0 {
config.MaxConcurrency = 10
}
return &dependencyResolver{ return &dependencyResolver{
client: client, client: client,
config: config, config: config,
@@ -44,8 +49,8 @@ func (r *dependencyResolver) resolveDependencies(ctx context.Context,
// Result collection // Result collection
dependencies := make([]*packagev1.PackageVersion, 0) dependencies := make([]*packagev1.PackageVersion, 0)
// Start recursive resolution // Start concurrent resolution
err = r.resolvePackageDependenciesRecursive(ctx, pd, packageVersion, 0, visitedPackages, &dependencies) err = r.resolvePackageDependenciesConcurrent(ctx, pd, packageVersion, 0, visitedPackages, &dependencies)
if err != nil { if err != nil {
return nil, fmt.Errorf("failed to resolve dependencies: %w", err) return nil, fmt.Errorf("failed to resolve dependencies: %w", err)
} }
@@ -53,8 +58,8 @@ func (r *dependencyResolver) resolveDependencies(ctx context.Context,
return dependencies, nil return dependencies, nil
} }
// resolvePackageDependenciesRecursive resolves dependencies for a package version recursively // resolvePackageDependenciesConcurrent resolves dependencies for a package version concurrently
func (r *dependencyResolver) resolvePackageDependenciesRecursive( func (r *dependencyResolver) resolvePackageDependenciesConcurrent(
ctx context.Context, ctx context.Context,
pd packageregistry.PackageDiscovery, pd packageregistry.PackageDiscovery,
packageVersion *packagev1.PackageVersion, packageVersion *packagev1.PackageVersion,
@@ -85,7 +90,13 @@ func (r *dependencyResolver) resolvePackageDependenciesRecursive(
// Skip if already visited // Skip if already visited
packageKey := r.packageKey(packageVersion) packageKey := r.packageKey(packageVersion)
if _, ok := visitedPackages[packageKey]; ok {
alreadyVisited := false
r.synchronize(func() {
alreadyVisited = visitedPackages[packageKey]
})
if alreadyVisited {
return nil return nil
} }
@@ -120,22 +131,45 @@ func (r *dependencyResolver) resolvePackageDependenciesRecursive(
}) })
} }
// Process transitive dependencies if enabled // Add resolved dependencies to the result
if r.config.IncludeTransitiveDependencies && depth < r.config.TransitiveDepth { r.synchronize(func() {
for _, dependency := range resolvedDependencies { for _, dependency := range resolvedDependencies {
err := r.resolvePackageDependenciesRecursive(ctx, pd, dependency, depth+1, visitedPackages, result) if !slices.Contains(*result, dependency) {
if err != nil { *result = append(*result, dependency)
return ff(fmt.Errorf("failed to resolve transitive dependency: %w", err))
} }
} }
} })
// Finally add resolved dependencies to the result // Process transitive dependencies if enabled and depth limit not reached
for _, dependency := range resolvedDependencies { if r.config.IncludeTransitiveDependencies && depth < r.config.TransitiveDepth && len(resolvedDependencies) > 0 {
if !slices.Contains(*result, dependency) { // Create worker pool using semaphore pattern
r.synchronize(func() { semaphore := make(chan struct{}, r.config.MaxConcurrency)
*result = append(*result, dependency) errCh := make(chan error, len(resolvedDependencies))
}) var wg sync.WaitGroup
for _, dependency := range resolvedDependencies {
wg.Add(1)
go func(dep *packagev1.PackageVersion) {
defer wg.Done()
semaphore <- struct{}{}
defer func() { <-semaphore }()
err := r.resolvePackageDependenciesConcurrent(ctx, pd, dep, depth+1, visitedPackages, result)
if err != nil {
errCh <- err
}
}(dependency)
}
// Wait for all goroutines to finish
wg.Wait()
close(errCh)
// Check for errors
for err := range errCh {
return ff(fmt.Errorf("failed to resolve transitive dependency: %w", err))
} }
} }
+5
View File
@@ -16,6 +16,9 @@ type NpmDependencyResolverConfig struct {
// FailFast will stop resolving dependencies after the first error // FailFast will stop resolving dependencies after the first error
FailFast bool FailFast bool
// MaxConcurrency limits the number of concurrent goroutines used for dependency resolution
MaxConcurrency int
} }
func NewDefaultNpmDependencyResolverConfig() NpmDependencyResolverConfig { func NewDefaultNpmDependencyResolverConfig() NpmDependencyResolverConfig {
@@ -24,6 +27,7 @@ func NewDefaultNpmDependencyResolverConfig() NpmDependencyResolverConfig {
IncludeTransitiveDependencies: true, IncludeTransitiveDependencies: true,
TransitiveDepth: 5, TransitiveDepth: 5,
FailFast: false, FailFast: false,
MaxConcurrency: 10,
} }
} }
@@ -73,6 +77,7 @@ func (r *npmDependencyResolver) ResolveDependencies(ctx context.Context,
IncludeTransitiveDependencies: r.config.IncludeTransitiveDependencies, IncludeTransitiveDependencies: r.config.IncludeTransitiveDependencies,
TransitiveDepth: r.config.TransitiveDepth, TransitiveDepth: r.config.TransitiveDepth,
FailFast: r.config.FailFast, FailFast: r.config.FailFast,
MaxConcurrency: r.config.MaxConcurrency,
}) })
return resolver.resolveDependencies(ctx, packageVersion) return resolver.resolveDependencies(ctx, packageVersion)