diff --git a/TODO.md b/TODO.md index c70f11b..404abd1 100644 --- a/TODO.md +++ b/TODO.md @@ -29,6 +29,13 @@ P2: security: referer blacklist # Completed Steps +- 2026-10-04 upstream host semaphores and variant `.meta` files no longer + outlive their use (closes #87): the fetcher counts the fetches holding or + waiting for a slot of each upstream host's semaphore and removes the host's + semaphore once none is left, so fetches from many hosts no longer leave one + semaphore each until restart; `VariantStorage.Delete` removes the variant's + `.meta` file along with it, a missing `.meta` file not being an error, and + `DeleteWithMeta`, which eviction called for that, is gone. - 2026-10-04 `README.md` matches the code (closes #74): "Storage" names the cache directories pixa uses (`cache/sources`, `cache/metadata`, `cache/variants`) and how files are named in each, and the comments in diff --git a/internal/httpfetcher/fetch_internal_test.go b/internal/httpfetcher/fetch_internal_test.go index 619b85a..83c78c0 100644 --- a/internal/httpfetcher/fetch_internal_test.go +++ b/internal/httpfetcher/fetch_internal_test.go @@ -191,7 +191,15 @@ func fetchBody(t *testing.T, f *HTTPFetcher, path string) string { // semLen reports how many per-host semaphore slots are currently held. func semLen(f *HTTPFetcher, host string) int { - return len(f.getHostSemaphore(host)) + f.hostSemMu.Lock() + defer f.hostSemMu.Unlock() + + sem, ok := f.hostSems[host] + if !ok { + return 0 + } + + return len(sem.slots) } // hostSemCount reports how many hosts have a semaphore in hostSems. diff --git a/internal/httpfetcher/httpfetcher.go b/internal/httpfetcher/httpfetcher.go index d65f2da..3cfee1c 100644 --- a/internal/httpfetcher/httpfetcher.go +++ b/internal/httpfetcher/httpfetcher.go @@ -160,10 +160,13 @@ func DefaultConfig() *Config { // HTTPFetcher implements Fetcher with SSRF protection and connection limits // per host and for all hosts together. type HTTPFetcher struct { - client *http.Client - config *Config - hostSems map[string]chan struct{} // per-host semaphores - hostSemMu sync.Mutex // protects hostSems map + client *http.Client + config *Config + // hostSems holds the semaphore of each host with a fetch holding or + // waiting for one of its slots; the entry is removed when the host's + // last such fetch gives its slot back or stops waiting. + hostSems map[string]*hostSemaphore + hostSemMu sync.Mutex // protects hostSems and each entry's count // allHostsSemaphore has one slot per connection allowed to all hosts // together (config.MaxConnections). allHostsSemaphore chan struct{} @@ -171,6 +174,14 @@ type HTTPFetcher struct { connectionWaitTimeout time.Duration } +// hostSemaphore is one host's connection slots +// (config.MaxConnectionsPerHost) and the number of fetches holding or +// waiting for one of them. +type hostSemaphore struct { + slots chan struct{} + count int +} + // New creates a new HTTPFetcher with SSRF protection. func New(config *Config) *HTTPFetcher { if config == nil { @@ -211,7 +222,7 @@ func New(config *Config) *HTTPFetcher { return &HTTPFetcher{ client: client, config: config, - hostSems: make(map[string]chan struct{}), + hostSems: make(map[string]*hostSemaphore), allHostsSemaphore: make(chan struct{}, config.MaxConnections), connectionWaitTimeout: ConnectionWaitTimeout, } @@ -307,6 +318,8 @@ func (f *HTTPFetcher) acquireConnection( select { case hostSem <- struct{}{}: case <-ctx.Done(): + f.putHostSemaphore(host) + return nil, ctx.Err() } @@ -314,32 +327,55 @@ func (f *HTTPFetcher) acquireConnection( case f.allHostsSemaphore <- struct{}{}: case <-time.After(f.connectionWaitTimeout): <-hostSem + f.putHostSemaphore(host) return nil, ErrTooManyConnections case <-ctx.Done(): <-hostSem + f.putHostSemaphore(host) return nil, ctx.Err() } return func() { <-hostSem + f.putHostSemaphore(host) <-f.allHostsSemaphore }, nil } -// getHostSemaphore returns the semaphore for a host, creating it if necessary. +// getHostSemaphore returns the semaphore for a host, creating it if +// necessary, and counts the caller among the fetches using it. The caller +// calls putHostSemaphore once it holds no slot and waits for none. func (f *HTTPFetcher) getHostSemaphore(host string) chan struct{} { f.hostSemMu.Lock() defer f.hostSemMu.Unlock() sem, ok := f.hostSems[host] if !ok { - sem = make(chan struct{}, f.config.MaxConnectionsPerHost) + sem = &hostSemaphore{ + slots: make(chan struct{}, f.config.MaxConnectionsPerHost), + } f.hostSems[host] = sem } - return sem + sem.count++ + + return sem.slots +} + +// putHostSemaphore stops counting the caller among the fetches using the +// host's semaphore, and removes the semaphore when no fetch uses it. +func (f *HTTPFetcher) putHostSemaphore(host string) { + f.hostSemMu.Lock() + defer f.hostSemMu.Unlock() + + sem := f.hostSems[host] + + sem.count-- + if sem.count == 0 { + delete(f.hostSems, host) + } } // buildResult validates the upstream response and assembles a FetchResult diff --git a/internal/imgcache/eviction.go b/internal/imgcache/eviction.go index 5d185ea..9f90b8f 100644 --- a/internal/imgcache/eviction.go +++ b/internal/imgcache/eviction.go @@ -283,7 +283,7 @@ func (c *Cache) evictVariant(ctx context.Context, cacheKey VariantKey) error { c.metaCache.Remove(cacheKey) - err = c.variants.DeleteWithMeta(cacheKey) + err = c.variants.Delete(cacheKey) if err != nil { return err } diff --git a/internal/imgcache/storage.go b/internal/imgcache/storage.go index 7e1cb38..395178d 100644 --- a/internal/imgcache/storage.go +++ b/internal/imgcache/storage.go @@ -564,7 +564,8 @@ func (s *VariantStorage) Exists(key VariantKey) bool { return err == nil } -// Delete removes content at the given key. +// Delete removes the content at the given key together with its .meta +// sidecar file. A missing file is not an error. func (s *VariantStorage) Delete(key VariantKey) error { path := s.keyToPath(key) @@ -573,18 +574,7 @@ func (s *VariantStorage) Delete(key VariantKey) error { return fmt.Errorf("failed to delete content: %w", err) } - return nil -} - -// DeleteWithMeta removes the content at the given key together with -// its .meta sidecar file. A missing file is not an error. -func (s *VariantStorage) DeleteWithMeta(key VariantKey) error { - err := s.Delete(key) - if err != nil { - return err - } - - metaPath := s.keyToPath(key) + ".meta" + metaPath := path + ".meta" err = os.Remove(metaPath) if err != nil && !os.IsNotExist(err) {