From ec97fdf91e9c1dfe74451c0cce7af4ee5acc19fa Mon Sep 17 00:00:00 2001 From: clawbot <35+clawbot@noreply.example.org> Date: Sun, 4 Oct 2026 02:28:54 +0000 Subject: [PATCH] Remove idle host semaphores and delete .meta with its variant (closes #87) Each upstream host's semaphore now counts the fetches holding or waiting for one of its slots, and is removed from hostSems when the last of them gives its slot back or stops waiting, so a long-running pixad no longer keeps one semaphore per host it ever fetched from. The semLen test helper reads hostSems directly, as getHostSemaphore now counts its caller. VariantStorage.Delete removes the variant's .meta file too, a missing one not being an error; DeleteWithMeta, which eviction called for that, is gone. Model: opus-5-5 --- TODO.md | 7 +++ internal/httpfetcher/fetch_internal_test.go | 10 +++- internal/httpfetcher/httpfetcher.go | 52 +++++++++++++++++---- internal/imgcache/eviction.go | 2 +- internal/imgcache/storage.go | 16 ++----- 5 files changed, 64 insertions(+), 23 deletions(-) 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) {