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
This commit is contained in:
2026-10-04 03:43:09 +00:00
parent 96be48f127
commit c9867d801a
5 changed files with 64 additions and 23 deletions
+7
View File
@@ -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-03 shutdown sets the exit code and waits for image processing
(closes #86): fx alone handles SIGINT and SIGTERM, and the server's own
signal handler is gone; fx's `Run` in `cmd/pixad` exits with the shutdown's
+9 -1
View File
@@ -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.
+44 -8
View File
@@ -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
+1 -1
View File
@@ -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
}
+3 -13
View File
@@ -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) {