Compare commits
1
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
213117e2ee |
@@ -71,17 +71,13 @@ prevent abuse, and allowlisted source hosts for open access.
|
|||||||
### Storage
|
### Storage
|
||||||
|
|
||||||
- **Source content**:
|
- **Source content**:
|
||||||
`<state_dir>/cache/sources/<ab>/<cd>/<sha256 of source content>`
|
`<statedir>/cache/src-content/<ab>/<cd>/<sha256 of source content>`
|
||||||
- **Source metadata**:
|
- **Source metadata**:
|
||||||
`<state_dir>/cache/metadata/<hostname>/<sha256 of path and query>.json`
|
`<statedir>/cache/src-metadata/<hostname>/<sha256 of path>.json`
|
||||||
(host, path and query, content hash, upstream status and headers, fetch time)
|
(fetch time, original headers, request, content hash)
|
||||||
- **Database**: `<state_dir>/state.sqlite3` (SQLite)
|
- **Database**: `<statedir>/state.sqlite3` (SQLite)
|
||||||
- **Transformed images**:
|
- **Output documents**:
|
||||||
`<state_dir>/cache/variants/<ab>/<cd>/<sha256 of host, path, query, size, format, quality and fit>`,
|
`<statedir>/cache/dst-content/<ab>/<cd>/<sha256 of output content>`
|
||||||
each with a `.meta` file beside it holding its content type
|
|
||||||
|
|
||||||
`<ab>` and `<cd>` are the first and second pairs of characters of the file's
|
|
||||||
name.
|
|
||||||
|
|
||||||
Multiple source paths may reference the same content blob; the
|
Multiple source paths may reference the same content blob; the
|
||||||
database tracks references rather than using filesystem refcounting.
|
database tracks references rather than using filesystem refcounting.
|
||||||
@@ -96,15 +92,12 @@ the metadata file stored beside it.
|
|||||||
/v1/image/<host>/<path>/<size>.<format>?sig=<signature>&exp=<expiration>
|
/v1/image/<host>/<path>/<size>.<format>?sig=<signature>&exp=<expiration>
|
||||||
```
|
```
|
||||||
|
|
||||||
Images are only fetched from origins using TLS with valid certificates, unless
|
Images are only fetched from origins using TLS with valid certificates.
|
||||||
`allow_http` is set: then pixa fetches every image over plain HTTP, which is for
|
|
||||||
testing only.
|
|
||||||
|
|
||||||
A request whose query string cannot be decoded, or gives any parameter more
|
A request whose query string cannot be decoded, or gives any parameter more
|
||||||
than once, is refused with 400.
|
than once, is refused with 400.
|
||||||
|
|
||||||
- `<format>`: one of `orig` (or `original`), `jpeg` (or `jpg`), `png`, `webp`,
|
- `<format>`: one of `orig`, `png`, `jpeg`, `webp`
|
||||||
`avif`, `gif`
|
|
||||||
- `<size>`: `orig` or `<width>x<height>` (e.g. `800x600`)
|
- `<size>`: `orig` or `<width>x<height>` (e.g. `800x600`)
|
||||||
|
|
||||||
An image is served with `Cache-Control: public, max-age=<seconds>, immutable`.
|
An image is served with `Cache-Control: public, max-age=<seconds>, immutable`.
|
||||||
@@ -180,8 +173,7 @@ Where:
|
|||||||
- `query` — source query string, empty string if none
|
- `query` — source query string, empty string if none
|
||||||
- `width` — requested width in pixels, `0` for original
|
- `width` — requested width in pixels, `0` for original
|
||||||
- `height` — requested height in pixels, `0` for original
|
- `height` — requested height in pixels, `0` for original
|
||||||
- `format` — output format, one of those listed under Routes, with `original`
|
- `format` — output format (jpeg, png, webp, avif, gif, orig)
|
||||||
signed as `orig` and `jpg` as `jpeg`
|
|
||||||
- `expiration` — the URL's `exp` query parameter, the Unix timestamp when
|
- `expiration` — the URL's `exp` query parameter, the Unix timestamp when
|
||||||
the signature expires; a request whose `exp` is not a whole number, an
|
the signature expires; a request whose `exp` is not a whole number, an
|
||||||
empty `exp=` included, is refused with 400
|
empty `exp=` included, is refused with 400
|
||||||
@@ -338,10 +330,7 @@ See `config.example.yml` for all options with defaults.
|
|||||||
- **Image processing**: govips (CGO wrapper for libvips)
|
- **Image processing**: govips (CGO wrapper for libvips)
|
||||||
- **Database**: SQLite via modernc.org/sqlite
|
- **Database**: SQLite via modernc.org/sqlite
|
||||||
- **Static assets**: embedded via `//go:embed`
|
- **Static assets**: embedded via `//go:embed`
|
||||||
- **Metrics**: Prometheus, at `/metrics`: generic HTTP request metrics
|
- **Metrics**: Prometheus
|
||||||
(duration, response size, requests in flight) and the Go runtime and process
|
|
||||||
metrics; requests are measured and `/metrics` is served only when
|
|
||||||
`metrics.username` and `metrics.password` are set
|
|
||||||
- **Logging**: stdlib slog
|
- **Logging**: stdlib slog
|
||||||
|
|
||||||
## Entrypoints
|
## Entrypoints
|
||||||
|
|||||||
@@ -37,21 +37,6 @@ P2: security: referer blacklist
|
|||||||
goroutine exits, stops waiting and returns its error; the handlers' stop hook
|
goroutine exits, stops waiting and returns its error; the handlers' stop hook
|
||||||
passes fx's stop context, so an eviction still running when fx's stop
|
passes fx's stop context, so an eviction still running when fx's stop
|
||||||
deadline ends fails the stop and makes the exit code 1.
|
deadline ends fails the stop and makes the exit code 1.
|
||||||
- 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
|
|
||||||
`001_schema.sql` name the same paths; the routes and the signature section
|
|
||||||
list the same output formats, `jpg` and `original` included; the TLS
|
|
||||||
sentence names `allow_http` as its exception; "Metrics" says only generic
|
|
||||||
HTTP and Go runtime metrics exist, measured and served only when the metrics
|
|
||||||
username and password are set.
|
|
||||||
- 2026-10-03 shutdown sets the exit code and waits for image processing
|
- 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
|
(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
|
signal handler is gone; fx's `Run` in `cmd/pixad` exits with the shutdown's
|
||||||
|
|||||||
@@ -2,7 +2,7 @@
|
|||||||
-- Creates all tables for the pixa caching image proxy
|
-- Creates all tables for the pixa caching image proxy
|
||||||
|
|
||||||
-- Source content blobs
|
-- Source content blobs
|
||||||
-- Files stored at: cache/sources/<ab>/<cd>/<sha256>
|
-- Files stored at: cache/src-content/<ab>/<cd>/<sha256>
|
||||||
-- last_accessed_at is NULL until the first LRU touch; eviction falls
|
-- last_accessed_at is NULL until the first LRU touch; eviction falls
|
||||||
-- back to fetched_at for rows that have never been touched.
|
-- back to fetched_at for rows that have never been touched.
|
||||||
CREATE TABLE IF NOT EXISTS source_content (
|
CREATE TABLE IF NOT EXISTS source_content (
|
||||||
@@ -16,7 +16,7 @@ CREATE INDEX IF NOT EXISTS idx_source_content_last_accessed
|
|||||||
ON source_content(last_accessed_at);
|
ON source_content(last_accessed_at);
|
||||||
|
|
||||||
-- Source URL metadata - maps URLs to content hashes
|
-- Source URL metadata - maps URLs to content hashes
|
||||||
-- JSON stored at: cache/metadata/<hostname>/<path_hash>.json
|
-- JSON stored at: cache/src-metadata/<hostname>/<path_hash>.json
|
||||||
CREATE TABLE IF NOT EXISTS source_metadata (
|
CREATE TABLE IF NOT EXISTS source_metadata (
|
||||||
id INTEGER PRIMARY KEY AUTOINCREMENT,
|
id INTEGER PRIMARY KEY AUTOINCREMENT,
|
||||||
source_host TEXT NOT NULL,
|
source_host TEXT NOT NULL,
|
||||||
@@ -56,8 +56,7 @@ CREATE INDEX IF NOT EXISTS idx_variant_content_last_accessed
|
|||||||
ON variant_content(last_accessed_at);
|
ON variant_content(last_accessed_at);
|
||||||
|
|
||||||
-- Output/transformed content blobs
|
-- Output/transformed content blobs
|
||||||
-- Not written: transformed images are stored in cache/variants and
|
-- Files stored at: cache/dst-content/<ab>/<cd>/<sha256>
|
||||||
-- tracked in variant_content above.
|
|
||||||
CREATE TABLE IF NOT EXISTS output_content (
|
CREATE TABLE IF NOT EXISTS output_content (
|
||||||
content_hash TEXT PRIMARY KEY,
|
content_hash TEXT PRIMARY KEY,
|
||||||
content_type TEXT NOT NULL,
|
content_type TEXT NOT NULL,
|
||||||
|
|||||||
@@ -191,23 +191,7 @@ func fetchBody(t *testing.T, f *HTTPFetcher, path string) string {
|
|||||||
|
|
||||||
// semLen reports how many per-host semaphore slots are currently held.
|
// semLen reports how many per-host semaphore slots are currently held.
|
||||||
func semLen(f *HTTPFetcher, host string) int {
|
func semLen(f *HTTPFetcher, host string) int {
|
||||||
f.hostSemMu.Lock()
|
return len(f.getHostSemaphore(host))
|
||||||
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.
|
|
||||||
func hostSemCount(f *HTTPFetcher) int {
|
|
||||||
f.hostSemMu.Lock()
|
|
||||||
defer f.hostSemMu.Unlock()
|
|
||||||
|
|
||||||
return len(f.hostSems)
|
|
||||||
}
|
}
|
||||||
|
|
||||||
func TestFetchRedirectToPrivateIPBlocked(t *testing.T) {
|
func TestFetchRedirectToPrivateIPBlocked(t *testing.T) {
|
||||||
|
|||||||
@@ -160,13 +160,10 @@ func DefaultConfig() *Config {
|
|||||||
// HTTPFetcher implements Fetcher with SSRF protection and connection limits
|
// HTTPFetcher implements Fetcher with SSRF protection and connection limits
|
||||||
// per host and for all hosts together.
|
// per host and for all hosts together.
|
||||||
type HTTPFetcher struct {
|
type HTTPFetcher struct {
|
||||||
client *http.Client
|
client *http.Client
|
||||||
config *Config
|
config *Config
|
||||||
// hostSems holds the semaphore of each host with a fetch holding or
|
hostSems map[string]chan struct{} // per-host semaphores
|
||||||
// waiting for one of its slots; the entry is removed when the host's
|
hostSemMu sync.Mutex // protects hostSems map
|
||||||
// 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
|
// allHostsSemaphore has one slot per connection allowed to all hosts
|
||||||
// together (config.MaxConnections).
|
// together (config.MaxConnections).
|
||||||
allHostsSemaphore chan struct{}
|
allHostsSemaphore chan struct{}
|
||||||
@@ -174,14 +171,6 @@ type HTTPFetcher struct {
|
|||||||
connectionWaitTimeout time.Duration
|
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.
|
// New creates a new HTTPFetcher with SSRF protection.
|
||||||
func New(config *Config) *HTTPFetcher {
|
func New(config *Config) *HTTPFetcher {
|
||||||
if config == nil {
|
if config == nil {
|
||||||
@@ -222,7 +211,7 @@ func New(config *Config) *HTTPFetcher {
|
|||||||
return &HTTPFetcher{
|
return &HTTPFetcher{
|
||||||
client: client,
|
client: client,
|
||||||
config: config,
|
config: config,
|
||||||
hostSems: make(map[string]*hostSemaphore),
|
hostSems: make(map[string]chan struct{}),
|
||||||
allHostsSemaphore: make(chan struct{}, config.MaxConnections),
|
allHostsSemaphore: make(chan struct{}, config.MaxConnections),
|
||||||
connectionWaitTimeout: ConnectionWaitTimeout,
|
connectionWaitTimeout: ConnectionWaitTimeout,
|
||||||
}
|
}
|
||||||
@@ -318,8 +307,6 @@ func (f *HTTPFetcher) acquireConnection(
|
|||||||
select {
|
select {
|
||||||
case hostSem <- struct{}{}:
|
case hostSem <- struct{}{}:
|
||||||
case <-ctx.Done():
|
case <-ctx.Done():
|
||||||
f.putHostSemaphore(host)
|
|
||||||
|
|
||||||
return nil, ctx.Err()
|
return nil, ctx.Err()
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -327,55 +314,32 @@ func (f *HTTPFetcher) acquireConnection(
|
|||||||
case f.allHostsSemaphore <- struct{}{}:
|
case f.allHostsSemaphore <- struct{}{}:
|
||||||
case <-time.After(f.connectionWaitTimeout):
|
case <-time.After(f.connectionWaitTimeout):
|
||||||
<-hostSem
|
<-hostSem
|
||||||
f.putHostSemaphore(host)
|
|
||||||
|
|
||||||
return nil, ErrTooManyConnections
|
return nil, ErrTooManyConnections
|
||||||
case <-ctx.Done():
|
case <-ctx.Done():
|
||||||
<-hostSem
|
<-hostSem
|
||||||
f.putHostSemaphore(host)
|
|
||||||
|
|
||||||
return nil, ctx.Err()
|
return nil, ctx.Err()
|
||||||
}
|
}
|
||||||
|
|
||||||
return func() {
|
return func() {
|
||||||
<-hostSem
|
<-hostSem
|
||||||
f.putHostSemaphore(host)
|
|
||||||
<-f.allHostsSemaphore
|
<-f.allHostsSemaphore
|
||||||
}, nil
|
}, nil
|
||||||
}
|
}
|
||||||
|
|
||||||
// getHostSemaphore returns the semaphore for a host, creating it if
|
// getHostSemaphore returns the semaphore for a host, creating it if necessary.
|
||||||
// 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{} {
|
func (f *HTTPFetcher) getHostSemaphore(host string) chan struct{} {
|
||||||
f.hostSemMu.Lock()
|
f.hostSemMu.Lock()
|
||||||
defer f.hostSemMu.Unlock()
|
defer f.hostSemMu.Unlock()
|
||||||
|
|
||||||
sem, ok := f.hostSems[host]
|
sem, ok := f.hostSems[host]
|
||||||
if !ok {
|
if !ok {
|
||||||
sem = &hostSemaphore{
|
sem = make(chan struct{}, f.config.MaxConnectionsPerHost)
|
||||||
slots: make(chan struct{}, f.config.MaxConnectionsPerHost),
|
|
||||||
}
|
|
||||||
f.hostSems[host] = sem
|
f.hostSems[host] = sem
|
||||||
}
|
}
|
||||||
|
|
||||||
sem.count++
|
return sem
|
||||||
|
|
||||||
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
|
// buildResult validates the upstream response and assembles a FetchResult
|
||||||
|
|||||||
@@ -5,7 +5,6 @@ import (
|
|||||||
"errors"
|
"errors"
|
||||||
"net"
|
"net"
|
||||||
"strconv"
|
"strconv"
|
||||||
"sync"
|
|
||||||
"testing"
|
"testing"
|
||||||
"time"
|
"time"
|
||||||
)
|
)
|
||||||
@@ -120,98 +119,6 @@ func TestFetchFreesHostSlotWhenContextEndsWaitingForConnection(t *testing.T) {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
// TestFetchRemovesIdleHostSemaphores checks that a host's semaphore is
|
|
||||||
// removed once no fetch holds or waits for one of its slots: after 100
|
|
||||||
// concurrent fetches from 50 hosts have all finished, no semaphore is left.
|
|
||||||
func TestFetchRemovesIdleHostSemaphores(t *testing.T) {
|
|
||||||
t.Parallel()
|
|
||||||
|
|
||||||
srv := startUpstream(t)
|
|
||||||
f, _ := newServerFetcher(t, srv, nil)
|
|
||||||
ctx := testContext(t)
|
|
||||||
|
|
||||||
var wg sync.WaitGroup
|
|
||||||
|
|
||||||
for i := range 100 {
|
|
||||||
wg.Go(func() {
|
|
||||||
res, err := f.Fetch(ctx, imageURLOnPort(1+i%50))
|
|
||||||
if err != nil {
|
|
||||||
t.Errorf("Fetch() error = %v", err)
|
|
||||||
|
|
||||||
return
|
|
||||||
}
|
|
||||||
|
|
||||||
_ = res.Content.Close()
|
|
||||||
})
|
|
||||||
}
|
|
||||||
|
|
||||||
wg.Wait()
|
|
||||||
|
|
||||||
if n := hostSemCount(f); n != 0 {
|
|
||||||
t.Errorf("%d host semaphores left after every fetch finished, want 0", n)
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
// TestFetchRemovesHostSemaphoreWhenNoConnection checks that a fetch that
|
|
||||||
// ends without a connection leaves no semaphore behind: when it is refused
|
|
||||||
// after waiting for a connection shared by all hosts, when its context ends
|
|
||||||
// while it waits for its host's slot, and when its context ends while it
|
|
||||||
// waits for a connection shared by all hosts, long before the 10 second
|
|
||||||
// wait timeout.
|
|
||||||
func TestFetchRemovesHostSemaphoreWhenNoConnection(t *testing.T) {
|
|
||||||
t.Parallel()
|
|
||||||
|
|
||||||
srv := startUpstream(t)
|
|
||||||
|
|
||||||
cfg := DefaultConfig()
|
|
||||||
cfg.MaxConnections = 1
|
|
||||||
cfg.MaxConnectionsPerHost = 1
|
|
||||||
|
|
||||||
f, _ := newServerFetcher(t, srv, cfg)
|
|
||||||
f.connectionWaitTimeout = 100 * time.Millisecond
|
|
||||||
|
|
||||||
open, err := f.Fetch(testContext(t), imageURLOnPort(81))
|
|
||||||
if err != nil {
|
|
||||||
t.Fatalf("first Fetch() error = %v", err)
|
|
||||||
}
|
|
||||||
|
|
||||||
_, err = f.Fetch(testContext(t), imageURLOnPort(82))
|
|
||||||
if !errors.Is(err, ErrTooManyConnections) {
|
|
||||||
t.Fatalf("Fetch() from another host: error = %v, "+
|
|
||||||
"want ErrTooManyConnections", err)
|
|
||||||
}
|
|
||||||
|
|
||||||
ctx, cancel := context.WithTimeout(t.Context(), 100*time.Millisecond)
|
|
||||||
defer cancel()
|
|
||||||
|
|
||||||
_, err = f.Fetch(ctx, imageURLOnPort(81))
|
|
||||||
if !errors.Is(err, context.DeadlineExceeded) {
|
|
||||||
t.Fatalf("Fetch() from the busy host: error = %v, "+
|
|
||||||
"want context.DeadlineExceeded", err)
|
|
||||||
}
|
|
||||||
|
|
||||||
// Back to the 10 second wait, so the next fetch's context ends first.
|
|
||||||
f.connectionWaitTimeout = ConnectionWaitTimeout
|
|
||||||
|
|
||||||
ctx, cancel = context.WithTimeout(t.Context(), 100*time.Millisecond)
|
|
||||||
defer cancel()
|
|
||||||
|
|
||||||
_, err = f.Fetch(ctx, imageURLOnPort(83))
|
|
||||||
if !errors.Is(err, context.DeadlineExceeded) {
|
|
||||||
t.Fatalf("Fetch() from a host with nothing open: error = %v, "+
|
|
||||||
"want context.DeadlineExceeded", err)
|
|
||||||
}
|
|
||||||
|
|
||||||
err = open.Content.Close()
|
|
||||||
if err != nil {
|
|
||||||
t.Fatalf("close first body: %v", err)
|
|
||||||
}
|
|
||||||
|
|
||||||
if n := hostSemCount(f); n != 0 {
|
|
||||||
t.Errorf("%d host semaphores left after every fetch finished, want 0", n)
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
// TestFetchReleasesConnectionOnError checks that a fetch that fails after
|
// TestFetchReleasesConnectionOnError checks that a fetch that fails after
|
||||||
// taking its connection gives it back: with MaxConnections at 1, the slot
|
// taking its connection gives it back: with MaxConnections at 1, the slot
|
||||||
// must be free after the failure and the next fetch must succeed.
|
// must be free after the failure and the next fetch must succeed.
|
||||||
|
|||||||
@@ -288,7 +288,7 @@ func (c *Cache) evictVariant(ctx context.Context, cacheKey VariantKey) error {
|
|||||||
|
|
||||||
c.metaCache.Remove(cacheKey)
|
c.metaCache.Remove(cacheKey)
|
||||||
|
|
||||||
err = c.variants.Delete(cacheKey)
|
err = c.variants.DeleteWithMeta(cacheKey)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -6,10 +6,8 @@ import (
|
|||||||
"database/sql"
|
"database/sql"
|
||||||
"errors"
|
"errors"
|
||||||
"io/fs"
|
"io/fs"
|
||||||
"log/slog"
|
|
||||||
"os"
|
"os"
|
||||||
"path/filepath"
|
"path/filepath"
|
||||||
"strings"
|
|
||||||
"testing"
|
"testing"
|
||||||
"time"
|
"time"
|
||||||
|
|
||||||
@@ -868,17 +866,12 @@ func TestPeriodicReconciliationAdoptsFileThatAppearsAfterStartup(t *testing.T) {
|
|||||||
// TestStopEvictionInterruptsPassInProgress holds the test database's
|
// TestStopEvictionInterruptsPassInProgress holds the test database's
|
||||||
// only connection, so the startup reconciliation pass waits for it, and
|
// only connection, so the startup reconciliation pass waits for it, and
|
||||||
// checks that StopEviction stops that pass instead of waiting for the
|
// checks that StopEviction stops that pass instead of waiting for the
|
||||||
// connection to come free, and that the stop logs one warning: the
|
// connection to come free.
|
||||||
// interrupted reconciliation's, with no eviction pass started after it.
|
|
||||||
func TestStopEvictionInterruptsPassInProgress(t *testing.T) {
|
func TestStopEvictionInterruptsPassInProgress(t *testing.T) {
|
||||||
t.Parallel()
|
t.Parallel()
|
||||||
|
|
||||||
cache, _ := newEvictionTestCache(t, 1<<30)
|
cache, _ := newEvictionTestCache(t, 1<<30)
|
||||||
|
|
||||||
var logBuf bytes.Buffer
|
|
||||||
|
|
||||||
cache.log = slog.New(slog.NewJSONHandler(&logBuf, nil))
|
|
||||||
|
|
||||||
conn, err := cache.db.Conn(t.Context())
|
conn, err := cache.db.Conn(t.Context())
|
||||||
if err != nil {
|
if err != nil {
|
||||||
t.Fatalf("failed to take the database connection: %v", err)
|
t.Fatalf("failed to take the database connection: %v", err)
|
||||||
@@ -909,72 +902,6 @@ func TestStopEvictionInterruptsPassInProgress(t *testing.T) {
|
|||||||
t.Fatalf("StopEviction() error = %v, want nil: the pass waiting for "+
|
t.Fatalf("StopEviction() error = %v, want nil: the pass waiting for "+
|
||||||
"the database did not stop", err)
|
"the database did not stop", err)
|
||||||
}
|
}
|
||||||
|
|
||||||
t.Logf("log output: %s", logBuf.String())
|
|
||||||
|
|
||||||
warnings := strings.Count(logBuf.String(), `"level":"WARN"`)
|
|
||||||
if warnings != 1 {
|
|
||||||
t.Errorf("the stop logged %d warnings, want 1", warnings)
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
// TestEvictToLimitStopsAtNextCandidateOnceCancelled cancels the context
|
|
||||||
// while the oldest of three source blobs is being evicted, and checks that
|
|
||||||
// EvictToLimit then returns context.Canceled without evicting the other
|
|
||||||
// two or logging a warning for either of them.
|
|
||||||
func TestEvictToLimitStopsAtNextCandidateOnceCancelled(t *testing.T) {
|
|
||||||
t.Parallel()
|
|
||||||
|
|
||||||
cache, _ := newEvictionTestCache(t, 1)
|
|
||||||
|
|
||||||
var logBuf bytes.Buffer
|
|
||||||
|
|
||||||
cache.log = slog.New(slog.NewJSONHandler(&logBuf, nil))
|
|
||||||
|
|
||||||
hashes := []ContentHash{
|
|
||||||
storeEvictionTestSource(t, cache, "cancel.example.com", "/a.jpg",
|
|
||||||
bytes.Repeat([]byte{0x61}, 1000)),
|
|
||||||
storeEvictionTestSource(t, cache, "cancel.example.com", "/b.jpg",
|
|
||||||
bytes.Repeat([]byte{0x62}, 1000)),
|
|
||||||
storeEvictionTestSource(t, cache, "cancel.example.com", "/c.jpg",
|
|
||||||
bytes.Repeat([]byte{0x63}, 1000)),
|
|
||||||
}
|
|
||||||
|
|
||||||
base := time.Now().Add(-time.Hour)
|
|
||||||
|
|
||||||
for i, hash := range hashes {
|
|
||||||
setSourceLastAccessed(t, cache, hash, base.Add(time.Duration(i)*time.Minute))
|
|
||||||
}
|
|
||||||
|
|
||||||
ctx, cancel := context.WithCancel(t.Context())
|
|
||||||
defer cancel()
|
|
||||||
|
|
||||||
cache.evictSourceBlobTestHook = func(ContentHash) { cancel() }
|
|
||||||
|
|
||||||
err := cache.EvictToLimit(ctx)
|
|
||||||
t.Logf("EvictToLimit() error = %v", err)
|
|
||||||
t.Logf("log output: %s", logBuf.String())
|
|
||||||
|
|
||||||
if !errors.Is(err, context.Canceled) {
|
|
||||||
t.Errorf("EvictToLimit() error = %v, want context.Canceled", err)
|
|
||||||
}
|
|
||||||
|
|
||||||
if cache.srcContent.Exists(hashes[0]) {
|
|
||||||
t.Errorf("source blob %s, evicted when the context was cancelled, "+
|
|
||||||
"is still on disk", hashes[0])
|
|
||||||
}
|
|
||||||
|
|
||||||
for _, hash := range hashes[1:] {
|
|
||||||
if !cache.srcContent.Exists(hash) {
|
|
||||||
t.Errorf("source blob %s was evicted after the context was cancelled", hash)
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
if strings.Contains(logBuf.String(), `"level":"WARN"`) {
|
|
||||||
t.Errorf("EvictToLimit logged a warning after the context was cancelled")
|
|
||||||
}
|
|
||||||
|
|
||||||
assertNoDanglingReferences(t, cache)
|
|
||||||
}
|
}
|
||||||
|
|
||||||
// TestStopEvictionReturnsWhenItsContextEnds pauses an eviction pass where
|
// TestStopEvictionReturnsWhenItsContextEnds pauses an eviction pass where
|
||||||
|
|||||||
@@ -564,8 +564,7 @@ func (s *VariantStorage) Exists(key VariantKey) bool {
|
|||||||
return err == nil
|
return err == nil
|
||||||
}
|
}
|
||||||
|
|
||||||
// Delete removes the content at the given key together with its .meta
|
// Delete removes content at the given key.
|
||||||
// sidecar file. A missing file is not an error.
|
|
||||||
func (s *VariantStorage) Delete(key VariantKey) error {
|
func (s *VariantStorage) Delete(key VariantKey) error {
|
||||||
path := s.keyToPath(key)
|
path := s.keyToPath(key)
|
||||||
|
|
||||||
@@ -574,7 +573,18 @@ func (s *VariantStorage) Delete(key VariantKey) error {
|
|||||||
return fmt.Errorf("failed to delete content: %w", err)
|
return fmt.Errorf("failed to delete content: %w", err)
|
||||||
}
|
}
|
||||||
|
|
||||||
metaPath := path + ".meta"
|
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"
|
||||||
|
|
||||||
err = os.Remove(metaPath)
|
err = os.Remove(metaPath)
|
||||||
if err != nil && !os.IsNotExist(err) {
|
if err != nil && !os.IsNotExist(err) {
|
||||||
|
|||||||
@@ -438,73 +438,3 @@ func TestVariantStorage_StoreLogsFailedMetaWrite(t *testing.T) {
|
|||||||
t.Errorf("log missing %s; got %q", want, logBuf.String())
|
t.Errorf("log missing %s; got %q", want, logBuf.String())
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
// storeTestVariant stores one variant, with its .meta file, in a new
|
|
||||||
// VariantStorage and returns the storage and the variant's key.
|
|
||||||
func storeTestVariant(t *testing.T) (*VariantStorage, VariantKey) {
|
|
||||||
t.Helper()
|
|
||||||
|
|
||||||
storage, err := NewVariantStorage(t.TempDir(), slog.New(slog.DiscardHandler))
|
|
||||||
if err != nil {
|
|
||||||
t.Fatalf("NewVariantStorage() error = %v", err)
|
|
||||||
}
|
|
||||||
|
|
||||||
key := CacheKey(&ImageRequest{SourceHost: testHostCDN, SourcePath: testPathCat})
|
|
||||||
|
|
||||||
_, err = storage.Store(key, bytes.NewReader([]byte("variant data")), "image/webp")
|
|
||||||
if err != nil {
|
|
||||||
t.Fatalf("Store() error = %v", err)
|
|
||||||
}
|
|
||||||
|
|
||||||
return storage, key
|
|
||||||
}
|
|
||||||
|
|
||||||
// TestVariantStorage_DeleteRemovesMeta verifies that Delete removes the
|
|
||||||
// variant's .meta file along with the variant file.
|
|
||||||
func TestVariantStorage_DeleteRemovesMeta(t *testing.T) {
|
|
||||||
t.Parallel()
|
|
||||||
|
|
||||||
storage, key := storeTestVariant(t)
|
|
||||||
metaPath := storage.keyToPath(key) + ".meta"
|
|
||||||
|
|
||||||
_, err := os.Stat(metaPath)
|
|
||||||
if err != nil {
|
|
||||||
t.Fatalf("Store() wrote no .meta file: %v", err)
|
|
||||||
}
|
|
||||||
|
|
||||||
err = storage.Delete(key)
|
|
||||||
if err != nil {
|
|
||||||
t.Fatalf("Delete() error = %v", err)
|
|
||||||
}
|
|
||||||
|
|
||||||
if storage.Exists(key) {
|
|
||||||
t.Error("Exists() = true after delete, want false")
|
|
||||||
}
|
|
||||||
|
|
||||||
_, err = os.Stat(metaPath)
|
|
||||||
if !os.IsNotExist(err) {
|
|
||||||
t.Errorf(".meta file left after Delete() (stat err=%v)", err)
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
// TestVariantStorage_DeleteWithoutMeta verifies that Delete succeeds for a
|
|
||||||
// variant whose .meta file is missing.
|
|
||||||
func TestVariantStorage_DeleteWithoutMeta(t *testing.T) {
|
|
||||||
t.Parallel()
|
|
||||||
|
|
||||||
storage, key := storeTestVariant(t)
|
|
||||||
|
|
||||||
err := os.Remove(storage.keyToPath(key) + ".meta")
|
|
||||||
if err != nil {
|
|
||||||
t.Fatalf("removing .meta file: %v", err)
|
|
||||||
}
|
|
||||||
|
|
||||||
err = storage.Delete(key)
|
|
||||||
if err != nil {
|
|
||||||
t.Fatalf("Delete() error = %v, want nil", err)
|
|
||||||
}
|
|
||||||
|
|
||||||
if storage.Exists(key) {
|
|
||||||
t.Error("Exists() = true after delete, want false")
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|||||||
Reference in New Issue
Block a user