Skip to content

Commit 188ed9c

Browse files
authored
fix(cdn): bound goroutine fan-out and add per-request timeout (#37)
The CDN resolver created a fresh semaphore per recursion level, allowing 10^depth concurrent goroutines for deep dependency trees. Share a single semaphore across all depths, acquiring it only for HTTP work and releasing before spawning children to prevent deadlock. Add WithRequestTimeout for per-request context deadlines on registry and CDN fetch calls, and WithConcurrency to configure the shared limit. Extract clone() to eliminate field-list duplication across With* builders, using slices.Clone for slice fields to prevent aliasing. Closes #34, closes #35 Assisted-By: Claude Opus 4.6 (1M context) <noreply@anthropic.com>
1 parent a79a8a2 commit 188ed9c

20 files changed

Lines changed: 351 additions & 119 deletions

File tree

resolve/cdn/cdn.go

Lines changed: 108 additions & 117 deletions
Original file line numberDiff line numberDiff line change
@@ -24,6 +24,7 @@ import (
2424
"slices"
2525
"strings"
2626
"sync"
27+
"time"
2728

2829
mappacdn "bennypowers.dev/mappa/cdn"
2930
"bennypowers.dev/mappa/importmap"
@@ -42,8 +43,10 @@ type Resolver struct {
4243
conditions []string
4344
includeDev bool
4445
excludePackages []string
45-
maxDepth int // Maximum dependency depth (0 = unlimited)
46-
resolveScope bool // Whether to resolve transitive dependencies as scopes
46+
maxDepth int // Maximum dependency depth (0 = unlimited)
47+
resolveScope bool // Whether to resolve transitive dependencies as scopes
48+
requestTimeout time.Duration // Per-request timeout (0 = no timeout)
49+
concurrency int // Max concurrent goroutines across all depths
4750
}
4851

4952
// New creates a new CDN resolver with default settings.
@@ -56,25 +59,17 @@ func New(fetcher mappacdn.Fetcher) *Resolver {
5659
template: tmpl,
5760
cache: mappacdn.NewPackageCache(100),
5861
resolveScope: true,
62+
concurrency: 10,
5963
}
6064
}
6165

6266
// WithProvider returns a new Resolver using the specified CDN provider.
6367
func (r *Resolver) WithProvider(provider mappacdn.Provider) *Resolver {
6468
tmpl, _ := resolve.ParseTemplate(provider.ModuleTemplate)
65-
return &Resolver{
66-
fetcher: r.fetcher,
67-
provider: provider,
68-
registry: r.registry,
69-
template: tmpl,
70-
cache: r.cache,
71-
logger: r.logger,
72-
conditions: r.conditions,
73-
includeDev: r.includeDev,
74-
excludePackages: r.excludePackages,
75-
maxDepth: r.maxDepth,
76-
resolveScope: r.resolveScope,
77-
}
69+
c := r.clone()
70+
c.provider = provider
71+
c.template = tmpl
72+
return c
7873
}
7974

8075
// WithTemplate returns a new Resolver using a custom URL template.
@@ -83,122 +78,88 @@ func (r *Resolver) WithTemplate(pattern string) (*Resolver, error) {
8378
if err != nil {
8479
return nil, err
8580
}
86-
return &Resolver{
87-
fetcher: r.fetcher,
88-
provider: r.provider,
89-
registry: r.registry,
90-
template: tmpl,
91-
cache: r.cache,
92-
logger: r.logger,
93-
conditions: r.conditions,
94-
includeDev: r.includeDev,
95-
excludePackages: r.excludePackages,
96-
maxDepth: r.maxDepth,
97-
resolveScope: r.resolveScope,
98-
}, nil
81+
c := r.clone()
82+
c.template = tmpl
83+
return c, nil
9984
}
10085

10186
// WithLogger returns a new Resolver with the specified logger.
10287
func (r *Resolver) WithLogger(logger resolve.Logger) *Resolver {
103-
return &Resolver{
104-
fetcher: r.fetcher,
105-
provider: r.provider,
106-
registry: r.registry,
107-
template: r.template,
108-
cache: r.cache,
109-
logger: logger,
110-
conditions: r.conditions,
111-
includeDev: r.includeDev,
112-
excludePackages: r.excludePackages,
113-
maxDepth: r.maxDepth,
114-
resolveScope: r.resolveScope,
115-
}
88+
c := r.clone()
89+
c.logger = logger
90+
return c
11691
}
11792

11893
// WithConditions returns a new Resolver with the specified export conditions.
11994
func (r *Resolver) WithConditions(conditions []string) *Resolver {
120-
return &Resolver{
121-
fetcher: r.fetcher,
122-
provider: r.provider,
123-
registry: r.registry,
124-
template: r.template,
125-
cache: r.cache,
126-
logger: r.logger,
127-
conditions: conditions,
128-
includeDev: r.includeDev,
129-
excludePackages: r.excludePackages,
130-
maxDepth: r.maxDepth,
131-
resolveScope: r.resolveScope,
132-
}
95+
c := r.clone()
96+
c.conditions = conditions
97+
return c
13398
}
13499

135100
// WithIncludeDev returns a new Resolver that includes devDependencies.
136101
func (r *Resolver) WithIncludeDev(include bool) *Resolver {
137-
return &Resolver{
138-
fetcher: r.fetcher,
139-
provider: r.provider,
140-
registry: r.registry,
141-
template: r.template,
142-
cache: r.cache,
143-
logger: r.logger,
144-
conditions: r.conditions,
145-
includeDev: include,
146-
excludePackages: r.excludePackages,
147-
maxDepth: r.maxDepth,
148-
resolveScope: r.resolveScope,
149-
}
102+
c := r.clone()
103+
c.includeDev = include
104+
return c
150105
}
151106

152107
// WithMaxDepth returns a new Resolver with a maximum dependency depth.
153108
// 0 means unlimited (default), 1 means direct dependencies only.
154109
func (r *Resolver) WithMaxDepth(depth int) *Resolver {
155-
return &Resolver{
156-
fetcher: r.fetcher,
157-
provider: r.provider,
158-
registry: r.registry,
159-
template: r.template,
160-
cache: r.cache,
161-
logger: r.logger,
162-
conditions: r.conditions,
163-
includeDev: r.includeDev,
164-
excludePackages: r.excludePackages,
165-
maxDepth: depth,
166-
resolveScope: r.resolveScope,
167-
}
110+
c := r.clone()
111+
c.maxDepth = depth
112+
return c
168113
}
169114

170115
// WithResolveScope controls whether to generate scopes for transitive dependencies.
171116
func (r *Resolver) WithResolveScope(resolveScope bool) *Resolver {
172-
return &Resolver{
173-
fetcher: r.fetcher,
174-
provider: r.provider,
175-
registry: r.registry,
176-
template: r.template,
177-
cache: r.cache,
178-
logger: r.logger,
179-
conditions: r.conditions,
180-
includeDev: r.includeDev,
181-
excludePackages: r.excludePackages,
182-
maxDepth: r.maxDepth,
183-
resolveScope: resolveScope,
184-
}
117+
c := r.clone()
118+
c.resolveScope = resolveScope
119+
return c
185120
}
186121

187122
// WithExclude returns a new Resolver that excludes the specified packages
188123
// from the generated import map, including as transitive dependencies.
189124
func (r *Resolver) WithExclude(packages []string) *Resolver {
125+
c := r.clone()
126+
c.excludePackages = packages
127+
return c
128+
}
129+
130+
// WithRequestTimeout returns a new Resolver with a per-request timeout.
131+
// 0 means no timeout (default). Applied to each HTTP fetch and registry call.
132+
func (r *Resolver) WithRequestTimeout(d time.Duration) *Resolver {
133+
c := r.clone()
134+
c.requestTimeout = d
135+
return c
136+
}
137+
138+
// WithConcurrency returns a new Resolver with the specified max concurrent goroutines.
139+
// Shared across all recursion depths to prevent unbounded fan-out.
140+
func (r *Resolver) WithConcurrency(n int) *Resolver {
141+
c := r.clone()
142+
if n > 0 {
143+
c.concurrency = n
144+
}
145+
return c
146+
}
147+
148+
func (r *Resolver) clone() *Resolver {
190149
return &Resolver{
191150
fetcher: r.fetcher,
192151
provider: r.provider,
193152
registry: r.registry,
194153
template: r.template,
195154
cache: r.cache,
196155
logger: r.logger,
197-
conditions: r.conditions,
156+
conditions: slices.Clone(r.conditions),
198157
includeDev: r.includeDev,
199-
excludePackages: packages,
158+
excludePackages: slices.Clone(r.excludePackages),
200159
maxDepth: r.maxDepth,
201160
resolveScope: r.resolveScope,
161+
requestTimeout: r.requestTimeout,
162+
concurrency: r.concurrency,
202163
}
203164
}
204165

@@ -237,17 +198,14 @@ func (r *Resolver) ResolvePackageJSON(ctx context.Context, pkg *packagejson.Pack
237198
// Resolve each dependency
238199
var wg sync.WaitGroup
239200
var mu sync.Mutex
240-
sem := make(chan struct{}, 10) // Limit concurrency
201+
sem := make(chan struct{}, r.concurrency)
241202
visited := sync.Map{}
242203

243204
for name, versionRange := range deps {
244205
wg.Add(1)
245206
go func(pkgName, verRange string) {
246207
defer wg.Done()
247-
sem <- struct{}{}
248-
defer func() { <-sem }()
249-
250-
if err := r.resolvePackage(ctx, result, &mu, &visited, pkgName, verRange, 0); err != nil {
208+
if err := r.resolvePackage(ctx, result, &mu, &visited, sem, pkgName, verRange, 0); err != nil {
251209
if r.logger != nil {
252210
r.logger.Warning("Failed to resolve %s@%s: %v", pkgName, verRange, err)
253211
}
@@ -265,41 +223,47 @@ func (r *Resolver) ResolvePackageJSON(ctx context.Context, pkg *packagejson.Pack
265223
}
266224

267225
// resolvePackage resolves a single package and its dependencies.
226+
// sem is shared across all recursion depths to bound total concurrency.
227+
// Acquires sem only for HTTP work, releases before spawning children
228+
// to prevent deadlock when parents hold slots while waiting on children.
268229
func (r *Resolver) resolvePackage(
269230
ctx context.Context,
270231
im *importmap.ImportMap,
271232
mu *sync.Mutex,
272233
visited *sync.Map,
234+
sem chan struct{},
273235
pkgName, versionRange string,
274236
depth int,
275237
) error {
276-
// Check max depth
277238
if r.maxDepth > 0 && depth >= r.maxDepth {
278239
return nil
279240
}
280241

281-
// Resolve version
282-
version, err := r.registry.ResolveVersion(ctx, pkgName, versionRange)
242+
if err := r.acquireSem(ctx, sem); err != nil {
243+
return err
244+
}
245+
246+
version, err := r.resolveVersion(ctx, pkgName, versionRange)
283247
if err != nil {
248+
<-sem
284249
return err
285250
}
286251

287-
// Check if already visited at this or higher version
288252
cacheKey := pkgName + "@" + version
289253
if _, loaded := visited.LoadOrStore(cacheKey, true); loaded {
254+
<-sem
290255
return nil
291256
}
292257

293-
// Fetch package.json from CDN
294258
pkg, err := r.fetchPackageJSON(ctx, pkgName, version)
259+
// Release sem -- HTTP work done, recursive work doesn't need it
260+
<-sem
295261
if err != nil {
296262
return err
297263
}
298264

299-
// Add to imports
300265
r.addPackageImports(im, mu, pkgName, version, pkg)
301266

302-
// Resolve transitive dependencies if enabled
303267
if r.resolveScope && (r.maxDepth == 0 || depth < r.maxDepth) && len(pkg.Dependencies) > 0 {
304268
scopeKey := r.template.Expand(pkgName, version, "")
305269
if !strings.HasSuffix(scopeKey, "/") {
@@ -309,7 +273,6 @@ func (r *Resolver) resolvePackage(
309273
scopeEntries := make(map[string]string)
310274
var wg sync.WaitGroup
311275
var scopeMu sync.Mutex
312-
sem := make(chan struct{}, 10)
313276

314277
for depName, depVer := range pkg.Dependencies {
315278
if slices.Contains(r.excludePackages, depName) {
@@ -318,35 +281,35 @@ func (r *Resolver) resolvePackage(
318281
wg.Add(1)
319282
go func(name, ver string) {
320283
defer wg.Done()
321-
sem <- struct{}{}
322-
defer func() { <-sem }()
323284

324-
// Resolve transitive dependency version
325-
resolvedVer, err := r.registry.ResolveVersion(ctx, name, ver)
285+
if err := r.acquireSem(ctx, sem); err != nil {
286+
return
287+
}
288+
289+
resolvedVer, err := r.resolveVersion(ctx, name, ver)
326290
if err != nil {
291+
<-sem
327292
if r.logger != nil {
328293
r.logger.Warning("Failed to resolve transitive dep %s@%s: %v", name, ver, err)
329294
}
330295
return
331296
}
332297

333-
// Fetch package.json
334298
depPkg, err := r.fetchPackageJSON(ctx, name, resolvedVer)
299+
<-sem
335300
if err != nil {
336301
if r.logger != nil {
337302
r.logger.Warning("Failed to fetch %s@%s: %v", name, resolvedVer, err)
338303
}
339304
return
340305
}
341306

342-
// Build scope entries
343307
entries := r.buildPackageImports(name, resolvedVer, depPkg)
344308
scopeMu.Lock()
345309
maps.Copy(scopeEntries, entries)
346310
scopeMu.Unlock()
347311

348-
// Recursively resolve deeper dependencies using resolved version
349-
if err := r.resolvePackage(ctx, im, mu, visited, name, resolvedVer, depth+1); err != nil {
312+
if err := r.resolvePackage(ctx, im, mu, visited, sem, name, resolvedVer, depth+1); err != nil {
350313
if r.logger != nil {
351314
r.logger.Warning("Failed to resolve transitive dep %s: %v", name, err)
352315
}
@@ -371,18 +334,46 @@ func (r *Resolver) resolvePackage(
371334
return nil
372335
}
373336

337+
// acquireSem blocks until a semaphore slot is available or the context is cancelled.
338+
func (r *Resolver) acquireSem(ctx context.Context, sem chan struct{}) error {
339+
select {
340+
case sem <- struct{}{}:
341+
return nil
342+
case <-ctx.Done():
343+
return ctx.Err()
344+
}
345+
}
346+
347+
// resolveVersion wraps registry.ResolveVersion with an optional per-request timeout.
348+
func (r *Resolver) resolveVersion(ctx context.Context, pkgName, versionRange string) (string, error) {
349+
reqCtx, cancel := r.withTimeout(ctx)
350+
defer cancel()
351+
return r.registry.ResolveVersion(reqCtx, pkgName, versionRange)
352+
}
353+
374354
// fetchPackageJSON fetches and parses a package.json from the CDN.
375355
func (r *Resolver) fetchPackageJSON(ctx context.Context, pkgName, version string) (*packagejson.PackageJSON, error) {
376356
return r.cache.GetOrLoad(pkgName, version, func() (*packagejson.PackageJSON, error) {
377357
url := r.buildPackageJSONURL(pkgName, version)
378-
data, err := r.fetcher.Fetch(ctx, url)
358+
reqCtx, cancel := r.withTimeout(ctx)
359+
defer cancel()
360+
data, err := r.fetcher.Fetch(reqCtx, url)
379361
if err != nil {
380362
return nil, err
381363
}
382364
return packagejson.Parse(data)
383365
})
384366
}
385367

368+
// withTimeout derives a context with the configured request timeout.
369+
// Returns the original context and a no-op cancel if no timeout is set.
370+
func (r *Resolver) withTimeout(ctx context.Context) (context.Context, context.CancelFunc) {
371+
if r.requestTimeout > 0 {
372+
return context.WithTimeout(ctx, r.requestTimeout)
373+
}
374+
return ctx, func() {}
375+
}
376+
386377
// buildPackageJSONURL builds the URL for a package.json file.
387378
func (r *Resolver) buildPackageJSONURL(pkgName, version string) string {
388379
url := r.provider.PackageJSONTemplate

0 commit comments

Comments
 (0)