diff --git a/README.md b/README.md index cc1a877..8e36dc0 100644 --- a/README.md +++ b/README.md @@ -115,6 +115,9 @@ Configured via YAML file (`--config`). Key settings: - `upstream_max_response_size` — max origin response size - `downstream_timeout` — client response timeout - `signing_key` — HMAC secret for URL signatures +- `cache_max_bytes` — disk cache size limit in bytes; `0` disables the + disk cache entirely; omitted defaults to 75% of the free space on + the filesystem containing `/cache/` (minimum 500 MiB) See `config.example.yml` for all options with defaults. diff --git a/TODO.md b/TODO.md index 5c291cc..250904e 100644 --- a/TODO.md +++ b/TODO.md @@ -12,18 +12,35 @@ pre-1.0. No git tags exist. Recent work extracted the internal/magic, internal/allowlist, internal/httpfetcher, and internal/signature -packages. The gosec findings from the 2026-07-06 survey are resolved: -the last two open findings (G124, session cookie attributes in -internal/session) are fixed as of this change, so `make check` is green -on main. +packages. The gosec findings from the 2026-07-06 survey are resolved +and `make check` is green on main. The disk cache is now size-bounded +with LRU eviction (`cache_max_bytes`), closing the unbounded disk +growth DoS vector. # Next Step -P0: implement cache size management and eviction so the disk cannot -fill up +P1: implement blocked networks configuration to extend SSRF protection # Completed Steps +- 2026-08-07 implement cache size management and eviction (closes + #51): new `cache_max_bytes` config key validated by the startup + framework (explicit values used exactly with no floor, `0` disables + the disk cache entirely, omitted defaults to max(75% of free space + on the filesystem containing `/cache/`, 500 MiB), logged + at startup); processed variants are now tracked in the database + (migration 002 adds `variant_content` and an LRU timestamp on + `source_content`) so total usage is two SUMs, never a directory scan + on the hot path; a background goroutine evicts globally + least-recently-used entries (variants and source blobs merged) to + the limit, woken by a periodic ticker and by write-pressure + notifications from stores; a source blob and ALL of its + `source_metadata` references are deleted in one transaction before + the file is unlinked, so multi-referenced blobs are never removed + while referenced and rows never point at deleted files; a startup + reconciliation pass adopts untracked variant files, drops rows for + missing files, removes unreachable source blobs, and sweeps stale + temp files - 2026-08-07 validate configuration on startup, fail fast on bad config (closes #52): a config value that is set but unparseable or invalid aborts startup naming the key and value (defaults apply only @@ -79,8 +96,6 @@ fill up # Future Steps -- P1: implement blocked networks configuration to extend SSRF - protection - P1: rate limit global concurrent upstream fetches to prevent resource exhaustion - P1: strip EXIF and other metadata from processed images (privacy) diff --git a/config.example.yml b/config.example.yml index 900a0b6..122189e 100644 --- a/config.example.yml +++ b/config.example.yml @@ -28,6 +28,13 @@ allow_http: false # Maximum concurrent connections per upstream host (default: 20) upstream_connections_per_host: 20 +# Maximum disk cache size in bytes. Explicit values are used exactly as +# given; 0 disables the disk cache entirely (every request fetches and +# processes uncached). When omitted, the default is 75% of the free +# space on the filesystem containing /cache/ at startup, +# with a minimum of 500 MiB. +# cache_max_bytes: 10737418240 + # Sentry error reporting (optional) sentry_dsn: "" diff --git a/internal/config/cache_max_bytes_test.go b/internal/config/cache_max_bytes_test.go new file mode 100644 index 0000000..ac7bf72 --- /dev/null +++ b/internal/config/cache_max_bytes_test.go @@ -0,0 +1,271 @@ +package config + +import ( + "errors" + "io" + "log/slog" + "os" + "path/filepath" + "strings" + "testing" +) + +// discardLogger returns a logger that swallows all output, for tests +// that exercise code paths which log. +func discardLogger() *slog.Logger { + return slog.New(slog.NewTextHandler(io.Discard, nil)) +} + +// TestCacheMaxBytesExplicitValueUsedWithoutFloor verifies that an +// explicitly configured cache_max_bytes value is used exactly as +// given: the 500 MiB floor applies only to the computed default, never +// to explicit values. +func TestCacheMaxBytesExplicitValueUsedWithoutFloor(t *testing.T) { + yamlContent := "signing_key: " + validTestSigningKey + "\ncache_max_bytes: 1024\n" + + c, err := configFromYAML(t, yamlContent) + if err != nil { + t.Fatalf("explicit cache_max_bytes must be accepted, got error: %v", err) + } + + if c.CacheMaxBytes != 1024 { + t.Errorf("CacheMaxBytes = %d, want 1024 (no floor for explicit values)", + c.CacheMaxBytes) + } +} + +// TestCacheMaxBytesZeroIsValidAndDisablesCache verifies that an +// explicit zero is a valid value (it disables the disk cache), not an +// error. +func TestCacheMaxBytesZeroIsValidAndDisablesCache(t *testing.T) { + yamlContent := "signing_key: " + validTestSigningKey + "\ncache_max_bytes: 0\n" + + c, err := configFromYAML(t, yamlContent) + if err != nil { + t.Fatalf("cache_max_bytes: 0 must be accepted, got error: %v", err) + } + + if c.CacheMaxBytes != 0 { + t.Errorf("CacheMaxBytes = %d, want 0", c.CacheMaxBytes) + } +} + +// TestCacheMaxBytesLargeExplicitValueParses verifies that values above +// 32-bit range parse correctly (the field is an int64 byte count). +func TestCacheMaxBytesLargeExplicitValueParses(t *testing.T) { + yamlContent := "signing_key: " + validTestSigningKey + "\ncache_max_bytes: 10737418240\n" + + c, err := configFromYAML(t, yamlContent) + if err != nil { + t.Fatalf("large cache_max_bytes must be accepted, got error: %v", err) + } + + if c.CacheMaxBytes != 10737418240 { + t.Errorf("CacheMaxBytes = %d, want 10737418240", c.CacheMaxBytes) + } +} + +// TestCacheMaxBytesInvalidValuesAbortStartup verifies that a SET but +// invalid cache_max_bytes value aborts startup naming the key and the +// offending value, per the no-silent-fallback rule: defaults apply +// only to omitted keys. +func TestCacheMaxBytesInvalidValuesAbortStartup(t *testing.T) { + signingKeyLine := "signing_key: " + validTestSigningKey + "\n" + + cases := []struct { + name string + yaml string + // wantErrSubstrings must all appear in the error message. + wantErrSubstrings []string + }{ + { + name: "negative", + yaml: signingKeyLine + "cache_max_bytes: -1024\n", + wantErrSubstrings: []string{"cache_max_bytes", "-1024"}, + }, + { + name: "float", + yaml: signingKeyLine + "cache_max_bytes: 3.5\n", + wantErrSubstrings: []string{"cache_max_bytes", "3.5"}, + }, + { + name: "non-numeric string", + yaml: signingKeyLine + "cache_max_bytes: banana\n", + wantErrSubstrings: []string{"cache_max_bytes", "banana"}, + }, + { + name: "explicit null", + yaml: signingKeyLine + "cache_max_bytes: null\n", + wantErrSubstrings: []string{"cache_max_bytes", "null"}, + }, + { + name: "bare key no value", + yaml: signingKeyLine + "cache_max_bytes:\n", + wantErrSubstrings: []string{"cache_max_bytes", "null"}, + }, + { + name: "boolean", + yaml: signingKeyLine + "cache_max_bytes: true\n", + wantErrSubstrings: []string{"cache_max_bytes", "true"}, + }, + { + name: "list", + yaml: signingKeyLine + "cache_max_bytes:\n - 1\n", + wantErrSubstrings: []string{"cache_max_bytes"}, + }, + } + + for _, tc := range cases { + t.Run(tc.name, func(t *testing.T) { + c, err := configFromYAML(t, tc.yaml) + if err == nil { + t.Fatalf("config with %s cache_max_bytes must abort startup, got config: %+v", + tc.name, c) + } + + t.Logf("got expected error: %v", err) + + for _, want := range tc.wantErrSubstrings { + if !strings.Contains(err.Error(), want) { + t.Errorf("error %q does not mention %q", err.Error(), want) + } + } + }) + } +} + +// TestComputeDefaultCacheMaxBytesUses75PercentOfFreeSpace verifies the +// computed default is 75% of the probed free space when that exceeds +// the floor. +func TestComputeDefaultCacheMaxBytesUses75PercentOfFreeSpace(t *testing.T) { + // 4 GiB free -> 3 GiB default. + probe := func(string) (uint64, error) { return 4294967296, nil } + + got, err := ComputeDefaultCacheMaxBytes(t.TempDir(), probe) + if err != nil { + t.Fatalf("ComputeDefaultCacheMaxBytes returned error: %v", err) + } + + if got != 3221225472 { + t.Errorf("ComputeDefaultCacheMaxBytes = %d, want 3221225472 (75%% of 4 GiB)", got) + } +} + +// TestComputeDefaultCacheMaxBytesAppliesFloorToComputedDefault +// verifies that when 75% of free space is below 500 MiB, the computed +// default is floored at DefaultCacheMaxBytesFloor. +func TestComputeDefaultCacheMaxBytesAppliesFloorToComputedDefault(t *testing.T) { + cases := []struct { + name string + freeBytes uint64 + }{ + {name: "100 MiB free", freeBytes: 104857600}, + {name: "zero free", freeBytes: 0}, + {name: "just below floor threshold", freeBytes: 699050665}, + } + + for _, tc := range cases { + t.Run(tc.name, func(t *testing.T) { + probe := func(string) (uint64, error) { return tc.freeBytes, nil } + + got, err := ComputeDefaultCacheMaxBytes(t.TempDir(), probe) + if err != nil { + t.Fatalf("ComputeDefaultCacheMaxBytes returned error: %v", err) + } + + if got != DefaultCacheMaxBytesFloor { + t.Errorf("ComputeDefaultCacheMaxBytes = %d, want floor %d", + got, DefaultCacheMaxBytesFloor) + } + }) + } +} + +// TestComputeDefaultCacheMaxBytesPropagatesProbeError verifies that a +// failing free-space probe produces an error naming the config key, +// instead of a silently wrong default. +func TestComputeDefaultCacheMaxBytesPropagatesProbeError(t *testing.T) { + probe := func(string) (uint64, error) { return 0, errors.New("statfs failed") } + + _, err := ComputeDefaultCacheMaxBytes(t.TempDir(), probe) + if err == nil { + t.Fatal("probe failure must produce an error, got nil") + } + + t.Logf("got expected error: %v", err) + + if !strings.Contains(err.Error(), "cache_max_bytes") { + t.Errorf("error %q does not name the config key cache_max_bytes", err.Error()) + } +} + +// TestResolveCacheMaxBytesComputesDefaultWhenOmitted verifies that an +// omitted cache_max_bytes key resolves to the computed default, that +// the probe is pointed at /cache/ (which must be created +// first so statfs measures the right filesystem), and that the result +// lands on the Config. +func TestResolveCacheMaxBytesComputesDefaultWhenOmitted(t *testing.T) { + c, err := configFromYAML(t, "signing_key: "+validTestSigningKey+"\n") + if err != nil { + t.Fatalf("minimal config should be valid, got error: %v", err) + } + + c.StateDir = t.TempDir() + wantCacheDir := filepath.Join(c.StateDir, "cache") + + var probedPath string + + // 4 GiB free -> 3 GiB default. + probe := func(path string) (uint64, error) { + probedPath = path + + return 4294967296, nil + } + + if err := c.resolveCacheMaxBytes(discardLogger(), probe); err != nil { + t.Fatalf("resolveCacheMaxBytes returned error: %v", err) + } + + if c.CacheMaxBytes != 3221225472 { + t.Errorf("CacheMaxBytes = %d, want computed default 3221225472", c.CacheMaxBytes) + } + + if probedPath != wantCacheDir { + t.Errorf("free space probed at %q, want cache directory %q", probedPath, wantCacheDir) + } + + info, err := os.Stat(wantCacheDir) + if err != nil || !info.IsDir() { + t.Errorf("cache directory %q was not created before probing: info=%v err=%v", + wantCacheDir, info, err) + } +} + +// TestResolveCacheMaxBytesDoesNotOverrideExplicitValue verifies that +// an explicitly configured value survives resolution untouched and +// that the free-space probe is never consulted for it. +func TestResolveCacheMaxBytesDoesNotOverrideExplicitValue(t *testing.T) { + yamlContent := "signing_key: " + validTestSigningKey + "\ncache_max_bytes: 1024\n" + + c, err := configFromYAML(t, yamlContent) + if err != nil { + t.Fatalf("explicit cache_max_bytes must be accepted, got error: %v", err) + } + + c.StateDir = t.TempDir() + + probe := func(string) (uint64, error) { + t.Error("free-space probe must not be consulted for explicit values") + + return 0, errors.New("probe must not be called") + } + + if err := c.resolveCacheMaxBytes(discardLogger(), probe); err != nil { + t.Fatalf("resolveCacheMaxBytes returned error: %v", err) + } + + if c.CacheMaxBytes != 1024 { + t.Errorf("CacheMaxBytes = %d, want explicit 1024 (no floor, no recompute)", + c.CacheMaxBytes) + } +} diff --git a/internal/config/cachesize.go b/internal/config/cachesize.go new file mode 100644 index 0000000..2a0f8ad --- /dev/null +++ b/internal/config/cachesize.go @@ -0,0 +1,110 @@ +package config + +import ( + "fmt" + "log/slog" + "math" + "os" + "path/filepath" + "syscall" +) + +// DefaultCacheMaxBytesFloor is the minimum computed default for the +// cache_max_bytes setting: 500 MiB. The floor applies only to the +// computed default (when the key is omitted from the configuration), +// never to explicitly configured values. +const DefaultCacheMaxBytesFloor int64 = 524288000 + +// cacheDirPerms is the permission mode for the cache directory created +// before probing free space, matching the state directory permissions. +const cacheDirPerms = 0o750 + +// freeSpaceFractionNumerator and freeSpaceFractionDenominator express +// the 75% share of free space used for the computed default limit as +// integer arithmetic (dividing before multiplying avoids overflow). +const ( + freeSpaceFractionNumerator uint64 = 3 + freeSpaceFractionDenominator uint64 = 4 +) + +// FreeSpaceProbeFunc reports the number of free bytes available on the +// filesystem containing path. It is a function type so tests can +// inject a fake probe instead of depending on the host disk. +type FreeSpaceProbeFunc func(path string) (uint64, error) + +// defaultFreeSpaceProbe reports free filesystem bytes via statfs on +// the given path, as available to unprivileged processes. +func defaultFreeSpaceProbe(path string) (uint64, error) { + var stat syscall.Statfs_t + if err := syscall.Statfs(path, &stat); err != nil { + return 0, err + } + + if stat.Bsize < 0 { + return 0, fmt.Errorf("statfs reported negative block size %d for %q", stat.Bsize, path) + } + + blockSize := uint64(stat.Bsize) //nolint:gosec // G115: negative Bsize rejected above + + return stat.Bavail * blockSize, nil +} + +// ComputeDefaultCacheMaxBytes returns the default cache size limit for +// the filesystem containing cacheDir: 75% of the free bytes reported +// by probe, with a floor of DefaultCacheMaxBytesFloor. +func ComputeDefaultCacheMaxBytes(cacheDir string, probe FreeSpaceProbeFunc) (int64, error) { + freeBytes, err := probe(cacheDir) + if err != nil { + return 0, fmt.Errorf("config key %q: cannot determine free space for %q: %w", + "cache_max_bytes", cacheDir, err) + } + + computed := freeBytes / freeSpaceFractionDenominator * freeSpaceFractionNumerator + if computed > math.MaxInt64 { + computed = math.MaxInt64 + } + + limit := int64(computed) //nolint:gosec // G115: clamped to MaxInt64 above + + if limit < DefaultCacheMaxBytesFloor { + limit = DefaultCacheMaxBytesFloor + } + + return limit, nil +} + +// resolveCacheMaxBytes finalizes CacheMaxBytes after state_dir +// validation: an explicitly configured value is kept as-is (no floor +// applies), while an omitted key receives the computed default based +// on free space in /cache/. The cache directory is created +// first so statfs measures the filesystem that will actually hold the +// cache. The effective limit is logged either way. +func (c *Config) resolveCacheMaxBytes(log *slog.Logger, probe FreeSpaceProbeFunc) error { + if !c.cacheMaxBytesExplicit { + cacheDir := filepath.Join(c.StateDir, "cache") + + if err := os.MkdirAll(cacheDir, cacheDirPerms); err != nil { + return fmt.Errorf("config key %q: cannot create cache directory %q: %w", + "cache_max_bytes", cacheDir, err) + } + + limit, err := ComputeDefaultCacheMaxBytes(cacheDir, probe) + if err != nil { + return err + } + + c.CacheMaxBytes = limit + + log.Info("computed default cache size limit from free space", + "cache_max_bytes", limit, + "cache_dir", cacheDir, + ) + } + + log.Info("effective cache size limit", + "cache_max_bytes", c.CacheMaxBytes, + "cache_disabled", c.CacheMaxBytes == 0, + ) + + return nil +} diff --git a/internal/config/config.go b/internal/config/config.go index 702f1b8..9b34b04 100644 --- a/internal/config/config.go +++ b/internal/config/config.go @@ -48,6 +48,19 @@ type Config struct { AllowlistHosts []string // Hosts that don't require signatures AllowHTTP bool // Allow non-TLS upstream (testing only) UpstreamConnectionsPerHost int // Max concurrent connections per upstream host + + // CacheMaxBytes is the disk cache size limit in bytes. Zero + // disables the disk cache entirely. When cache_max_bytes is + // omitted from the configuration, this holds the computed default + // (75% of free space on the filesystem containing + // /cache/, floored at DefaultCacheMaxBytesFloor). + CacheMaxBytes int64 + + // cacheMaxBytesExplicit records whether cache_max_bytes was + // explicitly set in the configuration file. Explicit values are + // used exactly as given; only an omitted key gets the computed + // default (and its floor) in resolveCacheMaxBytes. + cacheMaxBytesExplicit bool } // New creates a new Config instance by loading configuration from file. @@ -73,6 +86,10 @@ func New(_ fx.Lifecycle, params Params) (*Config, error) { return nil, err } + if err := c.resolveCacheMaxBytes(log, defaultFreeSpaceProbe); err != nil { + return nil, err + } + if c.Debug { params.Logger.EnableDebugLogging() } @@ -111,6 +128,16 @@ func newFromSmartConfig(sc *smartconfig.Config) (*Config, error) { AllowHTTP: loader.boolVal("allow_http", false), UpstreamConnectionsPerHost: loader.intVal( "upstream_connections_per_host", DefaultUpstreamConnectionsPerHost), + CacheMaxBytes: loader.int64Val("cache_max_bytes", 0), + } + + // The computed default for cache_max_bytes needs a validated + // state_dir, so it is resolved later (resolveCacheMaxBytes); here + // we only record whether the operator set the key explicitly. + if sc != nil { + if _, present := sc.Get("cache_max_bytes"); present { + c.cacheMaxBytesExplicit = true + } } // Build DBURL from StateDir if not explicitly set. The derived URL @@ -219,7 +246,7 @@ func isKnownConfigKey(key string) bool { switch key { case "debug", "maintenance_mode", "port", "state_dir", "sentry_dsn", "db_url", "metrics", "signing_key", "allowlist_hosts", "allow_http", - "upstream_connections_per_host", "env": + "upstream_connections_per_host", "cache_max_bytes", "env": return true } @@ -289,6 +316,13 @@ func (c *Config) validate() error { return fmt.Errorf("config key %q: value must not be empty", "state_dir") } + // Zero is valid (it disables the disk cache); only negative + // values are rejected. No floor applies to explicit values. + if c.CacheMaxBytes < 0 { + return fmt.Errorf("config key %q: value %d must not be negative", + "cache_max_bytes", c.CacheMaxBytes) + } + for _, host := range c.AllowlistHosts { if err := validateAllowlistHost(host); err != nil { return err @@ -412,6 +446,19 @@ func (l *strictLoader) intVal(key string, defaultVal int) int { return val } +func (l *strictLoader) int64Val(key string, defaultVal int64) int64 { + if l.err != nil { + return 0 + } + + val, err := getInt64(l.sc, key, defaultVal) + if err != nil { + l.err = err + } + + return val +} + func (l *strictLoader) boolVal(key string, defaultVal bool) bool { if l.err != nil { return false @@ -492,6 +539,55 @@ func getInt(sc *smartconfig.Config, key string, defaultVal int) (int, error) { } } +// getInt64 returns the 64-bit integer value for key, or defaultVal if +// the key is omitted. A present value that is not a whole number, or +// is explicitly null, is an error; fractional values are never +// truncated and out-of-range values are never clamped. +func getInt64(sc *smartconfig.Config, key string, defaultVal int64) (int64, error) { + if sc == nil { + return defaultVal, nil + } + + raw, ok := sc.Get(key) + if !ok { + return defaultVal, nil + } + + if raw == nil { + return 0, errNullConfigValue(key) + } + + switch val := raw.(type) { + case int: + return int64(val), nil + case int64: + return val, nil + case uint64: + if val > math.MaxInt64 { + return 0, fmt.Errorf("config key %q: value %d overflows a 64-bit integer", + key, val) + } + + return int64(val), nil //nolint:gosec // G115: bounds checked above + case float64: + if val != math.Trunc(val) { + return 0, fmt.Errorf("config key %q: value %v is not an integer", key, val) + } + + return int64(val), nil + case string: + parsed, err := strconv.ParseInt(strings.TrimSpace(val), 10, 64) + if err != nil { + return 0, fmt.Errorf("config key %q: value %q is not an integer", key, val) + } + + return parsed, nil + default: + return 0, fmt.Errorf("config key %q: value %v (%T) is not an integer", + key, raw, raw) + } +} + // getBool returns the boolean value for key, or defaultVal if the key // is omitted. A present value that is not a boolean (or a ParseBool-able // string), or is explicitly null, is an error; numbers are not accepted diff --git a/internal/database/schema/002_cache_eviction.sql b/internal/database/schema/002_cache_eviction.sql new file mode 100644 index 0000000..42b4787 --- /dev/null +++ b/internal/database/schema/002_cache_eviction.sql @@ -0,0 +1,25 @@ +-- Migration 002: cache size accounting and eviction +-- +-- Tracks processed variants in the database (source content blobs are +-- already tracked in source_content) so total cache usage can be +-- computed without directory scans, and adds last-access timestamps +-- for LRU eviction ordering. + +-- Processed variant blobs +-- Files stored at: cache/variants/// (plus a +-- .meta sidecar with the content type) +CREATE TABLE IF NOT EXISTS variant_content ( + cache_key TEXT PRIMARY KEY, + size_bytes INTEGER NOT NULL, + content_type TEXT NOT NULL, + created_at DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP, + last_accessed_at DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP +); +CREATE INDEX IF NOT EXISTS idx_variant_content_last_accessed + ON variant_content(last_accessed_at); + +-- LRU timestamp for source content blobs. Rows written before this +-- migration have NULL here; eviction falls back to fetched_at. +ALTER TABLE source_content ADD COLUMN last_accessed_at DATETIME; +CREATE INDEX IF NOT EXISTS idx_source_content_last_accessed + ON source_content(last_accessed_at); diff --git a/internal/handlers/handlers.go b/internal/handlers/handlers.go index 89dafbe..ee5ca55 100644 --- a/internal/handlers/handlers.go +++ b/internal/handlers/handlers.go @@ -53,6 +53,13 @@ func New(lc fx.Lifecycle, params Params) (*Handlers, error) { OnStart: func(_ context.Context) error { return s.initImageService() }, + OnStop: func(_ context.Context) error { + if s.imgCache != nil { + s.imgCache.StopEviction() + } + + return nil + }, }) return s, nil @@ -60,11 +67,15 @@ func New(lc fx.Lifecycle, params Params) (*Handlers, error) { // initImageService initializes the image cache and service. func (s *Handlers) initImageService() error { - // Create the cache + // Create the cache. cache_max_bytes: 0 disables the disk cache + // entirely; any other value is the eviction limit in bytes. cache, err := imgcache.NewCache(s.db.DB(), imgcache.CacheConfig{ - StateDir: s.config.StateDir, - CacheTTL: imgcache.DefaultCacheTTL, - NegativeTTL: imgcache.DefaultNegativeTTL, + StateDir: s.config.StateDir, + CacheTTL: imgcache.DefaultCacheTTL, + NegativeTTL: imgcache.DefaultNegativeTTL, + MaxBytes: s.config.CacheMaxBytes, + DisableDiskCache: s.config.CacheMaxBytes == 0, + Logger: s.log, }) if err != nil { return err @@ -72,6 +83,10 @@ func (s *Handlers) initImageService() error { s.imgCache = cache + // Background eviction: startup reconciliation, then periodic and + // write-pressure passes. No-op when the disk cache is disabled. + cache.StartEviction(imgcache.DefaultEvictionInterval) + // Create the fetcher config fetcherCfg := httpfetcher.DefaultConfig() fetcherCfg.AllowHTTP = s.config.AllowHTTP diff --git a/internal/imgcache/cache.go b/internal/imgcache/cache.go index 5e15873..390597a 100644 --- a/internal/imgcache/cache.go +++ b/internal/imgcache/cache.go @@ -7,7 +7,9 @@ import ( "errors" "fmt" "io" + "log/slog" "path/filepath" + "sync" "time" "sneak.berlin/go/pixa/internal/httpfetcher" @@ -27,6 +29,22 @@ type CacheConfig struct { StateDir string CacheTTL time.Duration NegativeTTL time.Duration + + // MaxBytes is the disk cache size limit in bytes that eviction + // enforces. Zero means no limit is enforced (no eviction). The + // config layer supplies the computed default when the operator + // omits cache_max_bytes. + MaxBytes int64 + + // DisableDiskCache turns the disk cache off entirely: no cache + // directories are created, lookups always miss, stores are + // no-ops, and no eviction machinery runs. The config layer sets + // this when the operator configures cache_max_bytes: 0. + DisableDiskCache bool + + // Logger receives accounting and eviction log output. A nil + // Logger means slog.Default(). + Logger *slog.Logger } // variantMeta stores content type for fast cache hits without reading .meta file. @@ -42,6 +60,19 @@ type Cache struct { variants *VariantStorage // processed variants by cache key srcMetadata *MetadataStorage // source metadata by host/path config CacheConfig + log *slog.Logger + + // disabled means the disk cache is turned off entirely: lookups + // always miss, stores are no-ops, and no eviction runs. + disabled bool + + // Eviction machinery. The channels are created in NewCache so + // stores can signal write pressure without racing StartEviction. + evictionPressure chan struct{} + evictionStop chan struct{} + evictionDone chan struct{} + evictionStarted bool + evictionStopOnce sync.Once // In-memory cache of variant metadata (content type, size) to avoid reading .meta files metaCache map[VariantKey]variantMeta @@ -49,6 +80,26 @@ type Cache struct { // NewCache creates a new cache instance. func NewCache(db *sql.DB, config CacheConfig) (*Cache, error) { + log := config.Logger + if log == nil { + log = slog.Default() + } + + c := &Cache{ + db: db, + config: config, + log: log, + disabled: config.DisableDiskCache, + evictionPressure: make(chan struct{}, 1), + evictionStop: make(chan struct{}), + evictionDone: make(chan struct{}), + metaCache: make(map[VariantKey]variantMeta), + } + + if c.disabled { + return c, nil + } + srcContent, err := NewContentStorage(filepath.Join(config.StateDir, "cache", "sources")) if err != nil { return nil, fmt.Errorf("failed to create source content storage: %w", err) @@ -64,14 +115,11 @@ func NewCache(db *sql.DB, config CacheConfig) (*Cache, error) { return nil, fmt.Errorf("failed to create source metadata storage: %w", err) } - return &Cache{ - db: db, - srcContent: srcContent, - variants: variants, - srcMetadata: srcMetadata, - config: config, - metaCache: make(map[VariantKey]variantMeta), - }, nil + c.srcContent = srcContent + c.variants = variants + c.srcMetadata = srcMetadata + + return c, nil } // LookupResult contains the result of a cache lookup. @@ -83,12 +131,15 @@ type LookupResult struct { CacheStatus CacheStatus } -// Lookup checks if a processed variant exists on disk (no DB access for hits). -func (c *Cache) Lookup(_ context.Context, req *ImageRequest) (*LookupResult, error) { +// Lookup checks if a processed variant exists on disk. Hits touch the +// variant's LRU timestamp; a disabled cache always misses. +func (c *Cache) Lookup(ctx context.Context, req *ImageRequest) (*LookupResult, error) { cacheKey := CacheKey(req) // Check variant storage directly - no DB needed for cache hits - if c.variants.Exists(cacheKey) { + if !c.disabled && c.variants.Exists(cacheKey) { + c.touchVariant(ctx, cacheKey) + return &LookupResult{ Hit: true, CacheKey: cacheKey, @@ -103,18 +154,53 @@ func (c *Cache) Lookup(_ context.Context, req *ImageRequest) (*LookupResult, err }, nil } +// touchVariant updates the LRU timestamp of a variant, best-effort: +// a failed touch only makes the entry look colder to eviction. +func (c *Cache) touchVariant(ctx context.Context, cacheKey VariantKey) { + _, err := c.db.ExecContext(ctx, ` + UPDATE variant_content SET last_accessed_at = CURRENT_TIMESTAMP + WHERE cache_key = ? + `, string(cacheKey)) + if err != nil { + c.log.Debug("failed to touch variant LRU timestamp", + "cache_key", cacheKey, "error", err) + } +} + +// touchSourceContent updates the LRU timestamp of a source content +// blob, best-effort: a failed touch only makes the blob look colder. +func (c *Cache) touchSourceContent(ctx context.Context, contentHash ContentHash) { + _, err := c.db.ExecContext(ctx, ` + UPDATE source_content SET last_accessed_at = CURRENT_TIMESTAMP + WHERE content_hash = ? + `, string(contentHash)) + if err != nil { + c.log.Debug("failed to touch source content LRU timestamp", + "content_hash", contentHash, "error", err) + } +} + // GetVariant returns a reader, size, and content type for a cached variant. func (c *Cache) GetVariant(cacheKey VariantKey) (io.ReadCloser, int64, string, error) { + if c.disabled { + return nil, 0, "", ErrNotFound + } + return c.variants.LoadWithMeta(cacheKey) } -// StoreSource stores fetched source content and metadata. +// StoreSource stores fetched source content and metadata. On a +// disabled cache it is a no-op returning an empty hash. func (c *Cache) StoreSource( ctx context.Context, req *ImageRequest, content io.Reader, result *httpfetcher.FetchResult, ) (ContentHash, error) { + if c.disabled { + return "", nil + } + // Store content contentHash, size, err := c.srcContent.Store(content) if err != nil { @@ -171,19 +257,52 @@ func (c *Cache) StoreSource( _ = err } + c.notifyWritePressure() + return contentHash, nil } -// StoreVariant stores a processed variant by its cache key. +// StoreVariant stores a processed variant by its cache key and records +// it in the size accounting. On a disabled cache it is a no-op. The +// accounting insert is best-effort (the startup reconciliation pass +// adopts any variant file that misses its accounting row). func (c *Cache) StoreVariant(cacheKey VariantKey, content io.Reader, contentType string) error { - _, err := c.variants.Store(cacheKey, content, contentType) + if c.disabled { + return nil + } - return err + size, err := c.variants.Store(cacheKey, content, contentType) + if err != nil { + return err + } + + _, err = c.db.Exec(` + INSERT INTO variant_content (cache_key, size_bytes, content_type) + VALUES (?, ?, ?) + ON CONFLICT(cache_key) DO UPDATE SET + size_bytes = excluded.size_bytes, + content_type = excluded.content_type, + last_accessed_at = CURRENT_TIMESTAMP + `, string(cacheKey), size, contentType) + if err != nil { + c.log.Warn("failed to record variant in size accounting", + "cache_key", cacheKey, "error", err) + } + + c.notifyWritePressure() + + return nil } // LookupSource checks if we have cached source content for a request. -// Returns the content hash and content type if found, or empty values if not. +// Returns the content hash and content type if found, or empty values +// if not. Hits touch the blob's LRU timestamp; a disabled cache always +// reports no cached source. func (c *Cache) LookupSource(ctx context.Context, req *ImageRequest) (ContentHash, string, error) { + if c.disabled { + return "", "", nil + } + var hashStr, contentType string err := c.db.QueryRowContext(ctx, ` @@ -206,6 +325,8 @@ func (c *Cache) LookupSource(ctx context.Context, req *ImageRequest) (ContentHas return "", "", nil } + c.touchSourceContent(ctx, contentHash) + return contentHash, contentType, nil } @@ -278,6 +399,10 @@ func (c *Cache) GetSourceMetadataID(ctx context.Context, req *ImageRequest) (int // GetSourceContent returns a reader for cached source content by its hash. func (c *Cache) GetSourceContent(contentHash ContentHash) (io.ReadCloser, error) { + if c.disabled { + return nil, ErrNotFound + } + return c.srcContent.Load(contentHash) } diff --git a/internal/imgcache/eviction.go b/internal/imgcache/eviction.go new file mode 100644 index 0000000..068d8ae --- /dev/null +++ b/internal/imgcache/eviction.go @@ -0,0 +1,741 @@ +package imgcache + +import ( + "context" + "encoding/json" + "fmt" + "io/fs" + "os" + "path/filepath" + "strings" + "time" +) + +// DefaultEvictionInterval is how often the background evictor checks +// cache usage against the configured limit, in addition to the +// write-pressure wakeups triggered by stores. +const DefaultEvictionInterval = 5 * time.Minute + +// evictionBatchSize is how many LRU candidates of each class (variants +// and source blobs) one eviction pass fetches from the database. +const evictionBatchSize = 100 + +// staleTempFileAge is how old an orphaned temp file (left behind by a +// crashed write) must be before reconciliation removes it. Fresh temp +// files may still belong to an in-flight store. +const staleTempFileAge = time.Hour + +// sqliteTimestampLayout matches SQLite's CURRENT_TIMESTAMP format, so +// timestamps written by reconciliation order correctly against ones +// written by the hot path. +const sqliteTimestampLayout = "2006-01-02 15:04:05" + +// tempFilePrefix is the prefix os.CreateTemp uses for in-flight cache +// writes (".tmp-*" patterns in the storage layer). +const tempFilePrefix = ".tmp-" + +// variantMetaSuffix is the sidecar suffix VariantStorage writes next +// to each variant file. +const variantMetaSuffix = ".meta" + +// fallbackContentType is recorded when a reconciled variant file has +// no readable .meta sidecar. +const fallbackContentType = "application/octet-stream" + +// UsageBytes returns the total number of bytes of cache content +// tracked in the database (source content blobs plus processed +// variants). It never scans the cache directories. +func (c *Cache) UsageBytes(ctx context.Context) (int64, error) { + if c.disabled { + return 0, nil + } + + var total int64 + + err := c.db.QueryRowContext(ctx, ` + SELECT (SELECT COALESCE(SUM(size_bytes), 0) FROM source_content) + + (SELECT COALESCE(SUM(size_bytes), 0) FROM variant_content) + `).Scan(&total) + if err != nil { + return 0, fmt.Errorf("failed to compute cache usage: %w", err) + } + + return total, nil +} + +// evictionCandidate is one LRU eviction victim candidate: either a +// processed variant (isVariant true, identified by cacheKey) or a +// source content blob (identified by contentHash). +type evictionCandidate struct { + isVariant bool + cacheKey VariantKey + contentHash ContentHash + sizeBytes int64 + lastAccessedAt string +} + +// EvictToLimit evicts least-recently-used cache entries until total +// tracked usage is at or below the configured MaxBytes limit. It is a +// no-op when the cache is disabled or no limit is configured. +func (c *Cache) EvictToLimit(ctx context.Context) error { + if c.disabled || c.config.MaxBytes <= 0 { + return nil + } + + for { + usage, err := c.UsageBytes(ctx) + if err != nil { + return err + } + + if usage <= c.config.MaxBytes { + return nil + } + + freed, err := c.evictBatch(ctx, usage-c.config.MaxBytes) + if err != nil { + return err + } + + if freed == 0 { + c.log.Warn("cache eviction made no progress", + "usage_bytes", usage, + "cache_max_bytes", c.config.MaxBytes, + ) + + return nil + } + + c.log.Info("evicted cache content", + "freed_bytes", freed, + "usage_bytes", usage-freed, + "cache_max_bytes", c.config.MaxBytes, + ) + } +} + +// evictBatch fetches one batch of LRU candidates across variants and +// source blobs and evicts them oldest-first until excessBytes are +// freed or the batch is exhausted. It returns the bytes freed. +func (c *Cache) evictBatch(ctx context.Context, excessBytes int64) (int64, error) { + candidates, err := c.evictionCandidates(ctx) + if err != nil { + return 0, err + } + + var freed int64 + + for _, candidate := range candidates { + if freed >= excessBytes { + break + } + + if err := c.evictCandidate(ctx, candidate); err != nil { + c.log.Warn("failed to evict cache entry", + "cache_key", candidate.cacheKey, + "content_hash", candidate.contentHash, + "error", err, + ) + + continue + } + + freed += candidate.sizeBytes + } + + return freed, nil +} + +// evictCandidate removes a single eviction victim. +func (c *Cache) evictCandidate(ctx context.Context, candidate evictionCandidate) error { + if candidate.isVariant { + return c.evictVariant(ctx, candidate.cacheKey) + } + + return c.evictSourceBlob(ctx, candidate.contentHash) +} + +// evictionCandidates returns up to evictionBatchSize variants and +// evictionBatchSize source blobs, merged into a single list ordered by +// last access time (oldest first). +func (c *Cache) evictionCandidates(ctx context.Context) ([]evictionCandidate, error) { + variants, err := c.variantCandidates(ctx) + if err != nil { + return nil, err + } + + sources, err := c.sourceCandidates(ctx) + if err != nil { + return nil, err + } + + // Merge the two lists, each already sorted oldest-first. SQLite + // CURRENT_TIMESTAMP strings compare correctly lexicographically. + merged := make([]evictionCandidate, 0, len(variants)+len(sources)) + + for len(variants) > 0 && len(sources) > 0 { + if variants[0].lastAccessedAt <= sources[0].lastAccessedAt { + merged = append(merged, variants[0]) + variants = variants[1:] + } else { + merged = append(merged, sources[0]) + sources = sources[1:] + } + } + + merged = append(merged, variants...) + merged = append(merged, sources...) + + return merged, nil +} + +// variantCandidates returns the least recently used variants. +func (c *Cache) variantCandidates(ctx context.Context) ([]evictionCandidate, error) { + rows, err := c.db.QueryContext(ctx, ` + SELECT cache_key, size_bytes, last_accessed_at + FROM variant_content + ORDER BY last_accessed_at ASC, cache_key ASC + LIMIT ? + `, evictionBatchSize) + if err != nil { + return nil, fmt.Errorf("failed to query variant eviction candidates: %w", err) + } + + defer func() { _ = rows.Close() }() + + var candidates []evictionCandidate + + for rows.Next() { + candidate := evictionCandidate{isVariant: true} + + var key string + if err := rows.Scan(&key, &candidate.sizeBytes, &candidate.lastAccessedAt); err != nil { + return nil, fmt.Errorf("failed to scan variant candidate: %w", err) + } + + candidate.cacheKey = VariantKey(key) + candidates = append(candidates, candidate) + } + + if err := rows.Err(); err != nil { + return nil, fmt.Errorf("variant candidate iteration failed: %w", err) + } + + return candidates, nil +} + +// sourceCandidates returns the least recently used source blobs. Rows +// written before the LRU column existed fall back to fetched_at. +func (c *Cache) sourceCandidates(ctx context.Context) ([]evictionCandidate, error) { + rows, err := c.db.QueryContext(ctx, ` + SELECT content_hash, size_bytes, + COALESCE(last_accessed_at, fetched_at, '1970-01-01 00:00:00') AS lru + FROM source_content + ORDER BY lru ASC, content_hash ASC + LIMIT ? + `, evictionBatchSize) + if err != nil { + return nil, fmt.Errorf("failed to query source eviction candidates: %w", err) + } + + defer func() { _ = rows.Close() }() + + var candidates []evictionCandidate + + for rows.Next() { + var candidate evictionCandidate + + var hash string + if err := rows.Scan(&hash, &candidate.sizeBytes, &candidate.lastAccessedAt); err != nil { + return nil, fmt.Errorf("failed to scan source candidate: %w", err) + } + + candidate.contentHash = ContentHash(hash) + candidates = append(candidates, candidate) + } + + if err := rows.Err(); err != nil { + return nil, fmt.Errorf("source candidate iteration failed: %w", err) + } + + return candidates, nil +} + +// evictVariant removes one variant: accounting row first, then the +// content and .meta files, so the database never references a deleted +// file. +func (c *Cache) evictVariant(ctx context.Context, cacheKey VariantKey) error { + _, err := c.db.ExecContext(ctx, + `DELETE FROM variant_content WHERE cache_key = ?`, string(cacheKey)) + if err != nil { + return fmt.Errorf("failed to delete variant accounting row: %w", err) + } + + if err := c.variants.DeleteWithMeta(cacheKey); err != nil { + return err + } + + return nil +} + +// sourceReference identifies one source_metadata row's JSON sidecar. +type sourceReference struct { + host string + pathHash PathHash +} + +// evictSourceBlob removes one source content blob. All source_metadata +// rows referencing the blob are deleted together with its +// source_content row in a single transaction BEFORE the file is +// unlinked: a blob referenced by multiple source paths is only ever +// removed together with all of its references, and database rows never +// point at deleted files. The JSON metadata sidecars for the removed +// rows are deleted afterwards. +func (c *Cache) evictSourceBlob(ctx context.Context, contentHash ContentHash) error { + references, err := c.sourceReferences(ctx, contentHash) + if err != nil { + return err + } + + tx, err := c.db.BeginTx(ctx, nil) + if err != nil { + return fmt.Errorf("failed to begin eviction transaction: %w", err) + } + + defer func() { _ = tx.Rollback() }() + + if _, err := tx.ExecContext(ctx, + `DELETE FROM source_metadata WHERE content_hash = ?`, string(contentHash)); err != nil { + return fmt.Errorf("failed to delete source metadata rows: %w", err) + } + + if _, err := tx.ExecContext(ctx, + `DELETE FROM source_content WHERE content_hash = ?`, string(contentHash)); err != nil { + return fmt.Errorf("failed to delete source content row: %w", err) + } + + if err := tx.Commit(); err != nil { + return fmt.Errorf("failed to commit eviction transaction: %w", err) + } + + // Only after the rows are gone may the files be removed. + for _, reference := range references { + if err := c.srcMetadata.Delete(reference.host, reference.pathHash); err != nil { + c.log.Warn("failed to delete metadata sidecar", + "host", reference.host, "path_hash", reference.pathHash, "error", err) + } + } + + if err := c.srcContent.Delete(contentHash); err != nil { + return err + } + + return nil +} + +// sourceReferences lists the metadata sidecar locations of every +// source_metadata row referencing the given blob. +func (c *Cache) sourceReferences( + ctx context.Context, contentHash ContentHash, +) ([]sourceReference, error) { + rows, err := c.db.QueryContext(ctx, ` + SELECT source_host, path_hash FROM source_metadata WHERE content_hash = ? + `, string(contentHash)) + if err != nil { + return nil, fmt.Errorf("failed to query source references: %w", err) + } + + defer func() { _ = rows.Close() }() + + var references []sourceReference + + for rows.Next() { + var reference sourceReference + + var pathHash string + if err := rows.Scan(&reference.host, &pathHash); err != nil { + return nil, fmt.Errorf("failed to scan source reference: %w", err) + } + + reference.pathHash = PathHash(pathHash) + references = append(references, reference) + } + + if err := rows.Err(); err != nil { + return nil, fmt.Errorf("source reference iteration failed: %w", err) + } + + return references, nil +} + +// notifyWritePressure wakes the background evictor after a store, so +// eviction under write pressure happens promptly without blocking the +// storing request. The notification channel has capacity one and drops +// when a wakeup is already pending. +func (c *Cache) notifyWritePressure() { + if c.disabled || c.config.MaxBytes <= 0 { + return + } + + select { + case c.evictionPressure <- struct{}{}: + default: + } +} + +// StartEviction launches the background eviction goroutine, which +// reconciles the database accounting with the cache directories once +// at startup and then evicts to the configured limit on the given +// periodic interval and on write-pressure notifications. It is a +// no-op on a disabled cache or when already started. +func (c *Cache) StartEviction(interval time.Duration) { + if c.disabled || c.evictionStarted { + return + } + + c.evictionStarted = true + + go c.evictionLoop(interval) +} + +// StopEviction stops the background eviction goroutine and waits for +// it to exit. It is safe to call when eviction was never started, and +// safe to call more than once. +func (c *Cache) StopEviction() { + if !c.evictionStarted { + return + } + + c.evictionStopOnce.Do(func() { + close(c.evictionStop) + <-c.evictionDone + }) +} + +// evictionLoop is the body of the background eviction goroutine. +func (c *Cache) evictionLoop(interval time.Duration) { + defer close(c.evictionDone) + + ctx := context.Background() + + if err := c.reconcileAccounting(ctx); err != nil { + c.log.Warn("cache accounting reconciliation failed", "error", err) + } + + c.runEvictionPass(ctx) + + ticker := time.NewTicker(interval) + defer ticker.Stop() + + for { + select { + case <-c.evictionStop: + return + case <-ticker.C: + case <-c.evictionPressure: + } + + c.runEvictionPass(ctx) + } +} + +// runEvictionPass runs one eviction pass, logging failures instead of +// propagating them (the loop must keep running). +func (c *Cache) runEvictionPass(ctx context.Context) { + if err := c.EvictToLimit(ctx); err != nil { + c.log.Warn("cache eviction pass failed", "error", err) + } +} + +// reconcileAccounting synchronizes the database size accounting with +// the actual contents of the cache directories. It runs once when the +// background evictor starts, off the request hot path: it adopts +// variant files that predate the accounting table, drops accounting +// rows whose files are missing, removes source blob files the database +// does not know (and rows whose files are gone), and sweeps stale temp +// files left behind by crashed writes. +func (c *Cache) reconcileAccounting(ctx context.Context) error { + if c.disabled { + return nil + } + + if err := c.reconcileVariantFiles(ctx); err != nil { + return err + } + + if err := c.reconcileVariantRows(ctx); err != nil { + return err + } + + if err := c.reconcileSourceFiles(ctx); err != nil { + return err + } + + if err := c.reconcileSourceRows(ctx); err != nil { + return err + } + + return nil +} + +// reconcileVariantFiles walks the variant storage directory, adopting +// files without accounting rows and sweeping stale temp files. +func (c *Cache) reconcileVariantFiles(ctx context.Context) error { + return filepath.WalkDir(c.variants.baseDir, func(path string, entry fs.DirEntry, err error) error { + if err != nil || entry.IsDir() { + return err + } + + name := entry.Name() + + if strings.HasPrefix(name, tempFilePrefix) { + c.sweepStaleTempFile(path, entry) + + return nil + } + + if strings.HasSuffix(name, variantMetaSuffix) { + return nil + } + + return c.adoptVariantFile(ctx, path, entry, VariantKey(name)) + }) +} + +// adoptVariantFile inserts an accounting row for a variant file that +// has none, using the file's size and modification time. +func (c *Cache) adoptVariantFile( + ctx context.Context, path string, entry fs.DirEntry, cacheKey VariantKey, +) error { + var rowExists int + + err := c.db.QueryRowContext(ctx, + `SELECT COUNT(*) FROM variant_content WHERE cache_key = ?`, string(cacheKey), + ).Scan(&rowExists) + if err != nil { + return fmt.Errorf("failed to check variant accounting row: %w", err) + } + + if rowExists > 0 { + return nil + } + + info, err := entry.Info() + if err != nil { + return fmt.Errorf("failed to stat variant file: %w", err) + } + + modTime := info.ModTime().UTC().Format(sqliteTimestampLayout) + contentType := c.variantContentTypeFromSidecar(path) + + _, err = c.db.ExecContext(ctx, ` + INSERT INTO variant_content + (cache_key, size_bytes, content_type, created_at, last_accessed_at) + VALUES (?, ?, ?, ?, ?) + `, string(cacheKey), info.Size(), contentType, modTime, modTime) + if err != nil { + return fmt.Errorf("failed to adopt variant file into accounting: %w", err) + } + + c.log.Info("adopted untracked variant file into size accounting", + "cache_key", cacheKey, "size_bytes", info.Size()) + + return nil +} + +// variantContentTypeFromSidecar reads the content type from a variant +// .meta sidecar, falling back to application/octet-stream. +func (c *Cache) variantContentTypeFromSidecar(variantPath string) string { + metaData, err := os.ReadFile(variantPath + variantMetaSuffix) //nolint:gosec // path from cache walk + if err != nil { + return fallbackContentType + } + + var meta VariantMeta + if json.Unmarshal(metaData, &meta) != nil || meta.ContentType == "" { + return fallbackContentType + } + + return meta.ContentType +} + +// reconcileVariantRows drops accounting rows whose variant files are +// missing, so the database never references deleted content. +func (c *Cache) reconcileVariantRows(ctx context.Context) error { + keys, err := c.allVariantKeys(ctx) + if err != nil { + return err + } + + for _, key := range keys { + if c.variants.Exists(key) { + continue + } + + if _, err := c.db.ExecContext(ctx, + `DELETE FROM variant_content WHERE cache_key = ?`, string(key)); err != nil { + return fmt.Errorf("failed to drop stale variant accounting row: %w", err) + } + + c.log.Info("dropped accounting row for missing variant file", "cache_key", key) + } + + return nil +} + +// allVariantKeys returns every tracked variant cache key. +func (c *Cache) allVariantKeys(ctx context.Context) ([]VariantKey, error) { + rows, err := c.db.QueryContext(ctx, `SELECT cache_key FROM variant_content`) + if err != nil { + return nil, fmt.Errorf("failed to query variant keys: %w", err) + } + + defer func() { _ = rows.Close() }() + + var keys []VariantKey + + for rows.Next() { + var key string + if err := rows.Scan(&key); err != nil { + return nil, fmt.Errorf("failed to scan variant key: %w", err) + } + + keys = append(keys, VariantKey(key)) + } + + if err := rows.Err(); err != nil { + return nil, fmt.Errorf("variant key iteration failed: %w", err) + } + + return keys, nil +} + +// reconcileSourceFiles walks the source content directory, removing +// blob files the database does not track (they are unreachable: source +// lookups always go through source_metadata) and sweeping stale temp +// files. +func (c *Cache) reconcileSourceFiles(ctx context.Context) error { + return filepath.WalkDir(c.srcContent.baseDir, func(path string, entry fs.DirEntry, err error) error { + if err != nil || entry.IsDir() { + return err + } + + name := entry.Name() + + if strings.HasPrefix(name, tempFilePrefix) { + c.sweepStaleTempFile(path, entry) + + return nil + } + + return c.removeUntrackedSourceFile(ctx, path, ContentHash(name)) + }) +} + +// removeUntrackedSourceFile deletes a source blob file that has no +// source_content row. Any source_metadata rows referencing the hash +// are removed first so no row ever points at a deleted file. +func (c *Cache) removeUntrackedSourceFile( + ctx context.Context, path string, contentHash ContentHash, +) error { + var rowExists int + + err := c.db.QueryRowContext(ctx, + `SELECT COUNT(*) FROM source_content WHERE content_hash = ?`, string(contentHash), + ).Scan(&rowExists) + if err != nil { + return fmt.Errorf("failed to check source content row: %w", err) + } + + if rowExists > 0 { + return nil + } + + if _, err := c.db.ExecContext(ctx, + `DELETE FROM source_metadata WHERE content_hash = ?`, string(contentHash)); err != nil { + return fmt.Errorf("failed to delete metadata rows for untracked blob: %w", err) + } + + //nolint:gosec // G703: path comes from walking our own cache directory + if err := os.Remove(path); err != nil && !os.IsNotExist(err) { + return fmt.Errorf("failed to remove untracked source file: %w", err) + } + + c.log.Info("removed untracked source content file", "content_hash", contentHash) + + return nil +} + +// reconcileSourceRows removes source_content rows (and their metadata +// references and sidecars) whose blob files are missing on disk. +func (c *Cache) reconcileSourceRows(ctx context.Context) error { + hashes, err := c.allSourceContentHashes(ctx) + if err != nil { + return err + } + + for _, hash := range hashes { + if c.srcContent.Exists(hash) { + continue + } + + // The blob file is already gone; evictSourceBlob removes the + // rows and sidecars and tolerates the missing file. + if err := c.evictSourceBlob(ctx, hash); err != nil { + return err + } + + c.log.Info("dropped rows for missing source content file", "content_hash", hash) + } + + return nil +} + +// allSourceContentHashes returns every tracked source content hash. +func (c *Cache) allSourceContentHashes(ctx context.Context) ([]ContentHash, error) { + rows, err := c.db.QueryContext(ctx, `SELECT content_hash FROM source_content`) + if err != nil { + return nil, fmt.Errorf("failed to query source content hashes: %w", err) + } + + defer func() { _ = rows.Close() }() + + var hashes []ContentHash + + for rows.Next() { + var hash string + if err := rows.Scan(&hash); err != nil { + return nil, fmt.Errorf("failed to scan content hash: %w", err) + } + + hashes = append(hashes, ContentHash(hash)) + } + + if err := rows.Err(); err != nil { + return nil, fmt.Errorf("content hash iteration failed: %w", err) + } + + return hashes, nil +} + +// sweepStaleTempFile removes a temp file left behind by a crashed +// write once it is old enough that no in-flight store can own it. +func (c *Cache) sweepStaleTempFile(path string, entry fs.DirEntry) { + info, err := entry.Info() + if err != nil { + return + } + + if time.Since(info.ModTime()) < staleTempFileAge { + return + } + + //nolint:gosec // G703: path comes from walking our own cache directory + if err := os.Remove(path); err != nil && !os.IsNotExist(err) { + c.log.Warn("failed to remove stale temp file", "path", path, "error", err) + + return + } + + c.log.Info("removed stale temp file", "path", path) +} diff --git a/internal/imgcache/eviction_test.go b/internal/imgcache/eviction_test.go new file mode 100644 index 0000000..0a16d05 --- /dev/null +++ b/internal/imgcache/eviction_test.go @@ -0,0 +1,669 @@ +package imgcache + +import ( + "bytes" + "context" + "database/sql" + "io/fs" + "os" + "path/filepath" + "testing" + "time" + + _ "modernc.org/sqlite" + "sneak.berlin/go/pixa/internal/database" + "sneak.berlin/go/pixa/internal/httpfetcher" +) + +// sqliteTimestampFormat matches the format SQLite's CURRENT_TIMESTAMP +// produces, so injected timestamps compare correctly against ones the +// implementation writes. +const sqliteTimestampFormat = "2006-01-02 15:04:05" + +// evictionTestDB creates an in-memory SQLite database with the real +// production schema, limited to a single connection so the background +// eviction goroutine shares the same in-memory database as the test. +func evictionTestDB(t *testing.T) *sql.DB { + t.Helper() + + db, err := sql.Open("sqlite", ":memory:") + if err != nil { + t.Fatalf("failed to open test db: %v", err) + } + + db.SetMaxOpenConns(1) + + if err := database.ApplyMigrations(context.Background(), db, nil); err != nil { + t.Fatalf("failed to apply migrations: %v", err) + } + + t.Cleanup(func() { _ = db.Close() }) + + return db +} + +// newEvictionTestCache creates a Cache backed by a temp directory and +// an in-memory database, with the given size limit. +func newEvictionTestCache(t *testing.T, maxBytes int64) (*Cache, string) { + t.Helper() + + tmpDir := t.TempDir() + db := evictionTestDB(t) + + // maxBytes zero mirrors the production mapping of + // cache_max_bytes: 0 (handlers sets DisableDiskCache); at the + // CacheConfig layer itself a zero MaxBytes means "no limit" for + // backwards compatibility with existing fixtures. + cache, err := NewCache(db, CacheConfig{ + StateDir: tmpDir, + CacheTTL: time.Hour, + NegativeTTL: 5 * time.Minute, + MaxBytes: maxBytes, + DisableDiskCache: maxBytes == 0, + }) + if err != nil { + t.Fatalf("failed to create cache: %v", err) + } + + return cache, tmpDir +} + +// storeEvictionTestSource stores content as a fetched source for +// host/path and returns the resulting content hash. +func storeEvictionTestSource( + t *testing.T, cache *Cache, host, path string, content []byte, +) ContentHash { + t.Helper() + + req := &ImageRequest{ + SourceHost: host, + SourcePath: path, + Format: FormatJPEG, + Quality: 85, + FitMode: FitCover, + } + + result := &httpfetcher.FetchResult{ + StatusCode: 200, + ContentType: "image/jpeg", + ContentLength: int64(len(content)), + Headers: map[string][]string{"Content-Type": {"image/jpeg"}}, + } + + hash, err := cache.StoreSource(context.Background(), req, bytes.NewReader(content), result) + if err != nil { + t.Fatalf("StoreSource(%s%s) failed: %v", host, path, err) + } + + return hash +} + +// storeEvictionTestVariant stores content as a processed variant under +// the given cache key. +func storeEvictionTestVariant(t *testing.T, cache *Cache, key VariantKey, content []byte) { + t.Helper() + + if err := cache.StoreVariant(key, bytes.NewReader(content), "image/webp"); err != nil { + t.Fatalf("StoreVariant(%s) failed: %v", key, err) + } +} + +// setVariantLastAccessed backdates the last access time of a tracked +// variant, to make LRU ordering deterministic in tests. +func setVariantLastAccessed(t *testing.T, cache *Cache, key VariantKey, when time.Time) { + t.Helper() + + res, err := cache.db.Exec( + `UPDATE variant_content SET last_accessed_at = ? WHERE cache_key = ?`, + when.UTC().Format(sqliteTimestampFormat), string(key), + ) + if err != nil { + t.Fatalf("failed to set variant last_accessed_at: %v", err) + } + + affected, err := res.RowsAffected() + if err != nil { + t.Fatalf("failed to read affected rows: %v", err) + } + + if affected != 1 { + t.Fatalf("variant %s has no accounting row (affected=%d); "+ + "stores must track variants in the database", key, affected) + } +} + +// setSourceLastAccessed backdates the last access time of a tracked +// source content blob. +func setSourceLastAccessed(t *testing.T, cache *Cache, hash ContentHash, when time.Time) { + t.Helper() + + res, err := cache.db.Exec( + `UPDATE source_content SET last_accessed_at = ? WHERE content_hash = ?`, + when.UTC().Format(sqliteTimestampFormat), string(hash), + ) + if err != nil { + t.Fatalf("failed to set source last_accessed_at: %v", err) + } + + affected, err := res.RowsAffected() + if err != nil { + t.Fatalf("failed to read affected rows: %v", err) + } + + if affected != 1 { + t.Fatalf("source %s has no accounting row (affected=%d)", hash, affected) + } +} + +// countRows returns the number of rows the given query yields. +func countRows(t *testing.T, cache *Cache, query string, args ...interface{}) int { + t.Helper() + + var n int + if err := cache.db.QueryRow(query, args...).Scan(&n); err != nil { + t.Fatalf("count query %q failed: %v", query, err) + } + + return n +} + +// assertNoDanglingReferences verifies the core eviction invariant: +// every database row that references cache content on disk points at a +// file that actually exists. +func assertNoDanglingReferences(t *testing.T, cache *Cache) { + t.Helper() + + rows, err := cache.db.Query( + `SELECT content_hash FROM source_metadata + WHERE content_hash IS NOT NULL AND content_hash != ''`, + ) + if err != nil { + t.Fatalf("failed to query source_metadata: %v", err) + } + + defer func() { _ = rows.Close() }() + + for rows.Next() { + var hash string + if err := rows.Scan(&hash); err != nil { + t.Fatalf("failed to scan content_hash: %v", err) + } + + if !cache.srcContent.Exists(ContentHash(hash)) { + t.Errorf("source_metadata references content %s but the file is missing", hash) + } + } + + if err := rows.Err(); err != nil { + t.Fatalf("source_metadata iteration failed: %v", err) + } + + variantRows, err := cache.db.Query(`SELECT cache_key FROM variant_content`) + if err != nil { + t.Fatalf("failed to query variant_content: %v", err) + } + + defer func() { _ = variantRows.Close() }() + + for variantRows.Next() { + var key string + if err := variantRows.Scan(&key); err != nil { + t.Fatalf("failed to scan cache_key: %v", err) + } + + if !cache.variants.Exists(VariantKey(key)) { + t.Errorf("variant_content references key %s but the file is missing", key) + } + } + + if err := variantRows.Err(); err != nil { + t.Fatalf("variant_content iteration failed: %v", err) + } +} + +// waitForUsageAtOrBelow polls UsageBytes until it reaches limit or the +// timeout expires, returning the last observed usage. +func waitForUsageAtOrBelow(t *testing.T, cache *Cache, limit int64, timeout time.Duration) int64 { + t.Helper() + + deadline := time.Now().Add(timeout) + + var usage int64 + + for time.Now().Before(deadline) { + var err error + + usage, err = cache.UsageBytes(context.Background()) + if err != nil { + t.Fatalf("UsageBytes failed: %v", err) + } + + if usage <= limit { + return usage + } + + time.Sleep(25 * time.Millisecond) + } + + return usage +} + +func TestUsageBytesAccountsSourceAndVariantBytes(t *testing.T) { + cache, _ := newEvictionTestCache(t, 1<<30) + + storeEvictionTestSource(t, cache, "src.example.com", "/a.jpg", + bytes.Repeat([]byte{0xAA}, 1000)) + storeEvictionTestSource(t, cache, "src.example.com", "/b.jpg", + bytes.Repeat([]byte{0xAB}, 2000)) + storeEvictionTestVariant(t, cache, "aabbccdd0001", bytes.Repeat([]byte{0xAC}, 500)) + storeEvictionTestVariant(t, cache, "aabbccdd0002", bytes.Repeat([]byte{0xAD}, 250)) + + usage, err := cache.UsageBytes(context.Background()) + if err != nil { + t.Fatalf("UsageBytes failed: %v", err) + } + + if usage != 3750 { + t.Errorf("UsageBytes = %d, want 3750 (1000+2000+500+250)", usage) + } +} + +func TestUsageBytesCountsMultiReferencedBlobOnce(t *testing.T) { + cache, _ := newEvictionTestCache(t, 1<<30) + + content := bytes.Repeat([]byte{0xCC}, 1200) + + hashOne := storeEvictionTestSource(t, cache, "src.example.com", "/one.jpg", content) + hashTwo := storeEvictionTestSource(t, cache, "src.example.com", "/two.jpg", content) + + if hashOne != hashTwo { + t.Fatalf("identical content produced different hashes: %s vs %s", hashOne, hashTwo) + } + + usage, err := cache.UsageBytes(context.Background()) + if err != nil { + t.Fatalf("UsageBytes failed: %v", err) + } + + if usage != 1200 { + t.Errorf("UsageBytes = %d, want 1200 (deduplicated blob counted once)", usage) + } +} + +func TestEvictToLimitEvictsLeastRecentlyUsedFirst(t *testing.T) { + const limit = 3000 + + cache, _ := newEvictionTestCache(t, limit) + + now := time.Now() + keys := []VariantKey{"aabbccdd0001", "aabbccdd0002", "aabbccdd0003", "aabbccdd0004"} + fills := []byte{0x01, 0x02, 0x03, 0x04} + ages := []time.Duration{4 * time.Hour, 3 * time.Hour, 2 * time.Hour, 1 * time.Hour} + + for i, key := range keys { + storeEvictionTestVariant(t, cache, key, bytes.Repeat([]byte{fills[i]}, 1000)) + setVariantLastAccessed(t, cache, key, now.Add(-ages[i])) + } + + if err := cache.EvictToLimit(context.Background()); err != nil { + t.Fatalf("EvictToLimit failed: %v", err) + } + + usage, err := cache.UsageBytes(context.Background()) + if err != nil { + t.Fatalf("UsageBytes failed: %v", err) + } + + if usage > limit { + t.Errorf("usage after eviction = %d, want <= %d", usage, limit) + } + + if cache.variants.Exists(keys[0]) { + t.Errorf("least recently used variant %s must be evicted", keys[0]) + } + + if n := countRows(t, cache, + `SELECT COUNT(*) FROM variant_content WHERE cache_key = ?`, string(keys[0]), + ); n != 0 { + t.Errorf("evicted variant %s still has %d accounting rows", keys[0], n) + } + + for _, key := range keys[1:] { + if !cache.variants.Exists(key) { + t.Errorf("more recently used variant %s must survive eviction", key) + } + } + + assertNoDanglingReferences(t, cache) +} + +func TestEvictionRemovesMultiReferencedBlobTogetherWithAllReferences(t *testing.T) { + const limit = 1000 + + cache, _ := newEvictionTestCache(t, limit) + + now := time.Now() + + // One 800-byte blob referenced by two source paths. + sharedContent := bytes.Repeat([]byte{0xDD}, 800) + sharedHash := storeEvictionTestSource(t, cache, "src.example.com", "/a.jpg", sharedContent) + + if h := storeEvictionTestSource(t, cache, "src.example.com", "/b.jpg", sharedContent); h != sharedHash { + t.Fatalf("identical content produced different hashes: %s vs %s", h, sharedHash) + } + + // A newer 600-byte blob referenced by one source path. + recentHash := storeEvictionTestSource(t, cache, "src.example.com", "/c.jpg", + bytes.Repeat([]byte{0xEE}, 600)) + + setSourceLastAccessed(t, cache, sharedHash, now.Add(-2*time.Hour)) + setSourceLastAccessed(t, cache, recentHash, now.Add(-time.Minute)) + + if err := cache.EvictToLimit(context.Background()); err != nil { + t.Fatalf("EvictToLimit failed: %v", err) + } + + usage, err := cache.UsageBytes(context.Background()) + if err != nil { + t.Fatalf("UsageBytes failed: %v", err) + } + + if usage > limit { + t.Errorf("usage after eviction = %d, want <= %d", usage, limit) + } + + // The multi-referenced blob must be gone from disk, from + // source_content, and from BOTH source_metadata rows: references + // are removed together with the blob, never left dangling. + if cache.srcContent.Exists(sharedHash) { + t.Errorf("evicted blob %s still exists on disk", sharedHash) + } + + if n := countRows(t, cache, + `SELECT COUNT(*) FROM source_content WHERE content_hash = ?`, string(sharedHash), + ); n != 0 { + t.Errorf("evicted blob %s still has %d source_content rows", sharedHash, n) + } + + if n := countRows(t, cache, + `SELECT COUNT(*) FROM source_metadata WHERE content_hash = ?`, string(sharedHash), + ); n != 0 { + t.Errorf("evicted blob %s still has %d source_metadata references", sharedHash, n) + } + + // The JSON metadata sidecars for both referencing paths must be + // removed along with the rows. + for _, path := range []string{"/a.jpg", "/b.jpg"} { + pathHash := HashPath(path + "?") + if cache.srcMetadata.Exists("src.example.com", pathHash) { + t.Errorf("metadata sidecar for %s must be removed with its row", path) + } + } + + // The more recently used blob survives fully intact. + if !cache.srcContent.Exists(recentHash) { + t.Errorf("recently used blob %s must survive eviction", recentHash) + } + + if n := countRows(t, cache, + `SELECT COUNT(*) FROM source_metadata WHERE content_hash = ?`, string(recentHash), + ); n != 1 { + t.Errorf("recently used blob %s has %d source_metadata rows, want 1", recentHash, n) + } + + assertNoDanglingReferences(t, cache) +} + +func TestEvictionKeepsEverythingWhenUnderLimit(t *testing.T) { + cache, _ := newEvictionTestCache(t, 1<<30) + + content := bytes.Repeat([]byte{0xDF}, 800) + hash := storeEvictionTestSource(t, cache, "src.example.com", "/a.jpg", content) + + if h := storeEvictionTestSource(t, cache, "src.example.com", "/b.jpg", content); h != hash { + t.Fatalf("identical content produced different hashes: %s vs %s", h, hash) + } + + storeEvictionTestVariant(t, cache, "aabbccdd0001", bytes.Repeat([]byte{0xE0}, 500)) + + if err := cache.EvictToLimit(context.Background()); err != nil { + t.Fatalf("EvictToLimit failed: %v", err) + } + + if !cache.srcContent.Exists(hash) { + t.Errorf("blob %s must not be evicted while usage is under the limit", hash) + } + + if n := countRows(t, cache, + `SELECT COUNT(*) FROM source_metadata WHERE content_hash = ?`, string(hash), + ); n != 2 { + t.Errorf("blob %s has %d source_metadata rows, want 2", hash, n) + } + + if !cache.variants.Exists("aabbccdd0001") { + t.Error("variant must not be evicted while usage is under the limit") + } + + assertNoDanglingReferences(t, cache) +} + +func TestZeroMaxBytesDisablesDiskCache(t *testing.T) { + cache, tmpDir := newEvictionTestCache(t, 0) + ctx := context.Background() + + req := &ImageRequest{ + SourceHost: "src.example.com", + SourcePath: "/a.jpg", + Format: FormatJPEG, + Quality: 85, + FitMode: FitCover, + } + + // Writes are no-ops that report success. + if err := cache.StoreVariant(CacheKey(req), bytes.NewReader([]byte("data")), "image/webp"); err != nil { + t.Fatalf("StoreVariant on disabled cache must be a no-op, got error: %v", err) + } + + result := &httpfetcher.FetchResult{ + StatusCode: 200, + ContentType: "image/jpeg", + ContentLength: 4, + Headers: map[string][]string{}, + } + + hash, err := cache.StoreSource(ctx, req, bytes.NewReader([]byte("data")), result) + if err != nil { + t.Fatalf("StoreSource on disabled cache must be a no-op, got error: %v", err) + } + + if hash != "" { + t.Errorf("StoreSource on disabled cache returned hash %q, want empty", hash) + } + + // Reads always miss. + lookup, err := cache.Lookup(ctx, req) + if err != nil { + t.Fatalf("Lookup on disabled cache failed: %v", err) + } + + if lookup.Hit { + t.Error("Lookup on disabled cache must always miss") + } + + srcHash, srcType, err := cache.LookupSource(ctx, req) + if err != nil { + t.Fatalf("LookupSource on disabled cache failed: %v", err) + } + + if srcHash != "" || srcType != "" { + t.Errorf("LookupSource on disabled cache = (%q, %q), want empty", srcHash, srcType) + } + + // Nothing is tracked and nothing is written to disk. + usage, err := cache.UsageBytes(ctx) + if err != nil { + t.Fatalf("UsageBytes failed: %v", err) + } + + if usage != 0 { + t.Errorf("UsageBytes on disabled cache = %d, want 0", usage) + } + + if n := countRows(t, cache, `SELECT COUNT(*) FROM source_content`); n != 0 { + t.Errorf("disabled cache wrote %d source_content rows, want 0", n) + } + + if n := countRows(t, cache, `SELECT COUNT(*) FROM source_metadata`); n != 0 { + t.Errorf("disabled cache wrote %d source_metadata rows, want 0", n) + } + + if _, err := os.Stat(filepath.Join(tmpDir, "cache")); !os.IsNotExist(err) { + t.Errorf("disabled cache must not create the cache directory tree (stat err=%v)", err) + } + + var foundFiles []string + + walkErr := filepath.WalkDir(tmpDir, func(path string, d fs.DirEntry, err error) error { + if err != nil { + return err + } + + if !d.IsDir() { + foundFiles = append(foundFiles, path) + } + + return nil + }) + if walkErr != nil { + t.Fatalf("failed to walk state dir: %v", walkErr) + } + + if len(foundFiles) != 0 { + t.Errorf("disabled cache wrote files to disk: %v", foundFiles) + } +} + +func TestEvictionRunsUnderWritePressure(t *testing.T) { + const limit = 1500 + + cache, _ := newEvictionTestCache(t, limit) + + // An interval far longer than the test ensures only write + // pressure can trigger eviction here. + cache.StartEviction(time.Hour) + defer cache.StopEviction() + + keys := []VariantKey{"aabbccdd0001", "aabbccdd0002", "aabbccdd0003"} + fills := []byte{0x11, 0x12, 0x13} + + for i, key := range keys { + storeEvictionTestVariant(t, cache, key, bytes.Repeat([]byte{fills[i]}, 1000)) + } + + usage := waitForUsageAtOrBelow(t, cache, limit, 5*time.Second) + if usage > limit { + t.Errorf("write pressure did not trigger eviction: usage = %d, want <= %d", + usage, limit) + } + + assertNoDanglingReferences(t, cache) +} + +func TestEvictionRunsOnPeriodicSchedule(t *testing.T) { + const limit = 1500 + + cache, _ := newEvictionTestCache(t, limit) + + // Start the evictor while the cache is empty, then create tracked + // over-limit state WITHOUT going through the store methods, so no + // write-pressure notification fires and only the periodic ticker + // can trigger eviction. + cache.StartEviction(100 * time.Millisecond) + defer cache.StopEviction() + + keys := []VariantKey{"aabbccdd0001", "aabbccdd0002", "aabbccdd0003"} + fills := []byte{0x21, 0x22, 0x23} + + for i, key := range keys { + content := bytes.Repeat([]byte{fills[i]}, 1000) + + if _, err := cache.variants.Store(key, bytes.NewReader(content), "image/webp"); err != nil { + t.Fatalf("failed to store variant file: %v", err) + } + + if _, err := cache.db.Exec( + `INSERT INTO variant_content (cache_key, size_bytes, content_type) + VALUES (?, ?, ?)`, + string(key), len(content), "image/webp", + ); err != nil { + t.Fatalf("failed to insert variant accounting row: %v", err) + } + } + + usage := waitForUsageAtOrBelow(t, cache, limit, 5*time.Second) + if usage > limit { + t.Errorf("periodic schedule did not trigger eviction: usage = %d, want <= %d", + usage, limit) + } + + assertNoDanglingReferences(t, cache) +} + +func TestStartEvictionReconcilesAccountingWithDisk(t *testing.T) { + cache, _ := newEvictionTestCache(t, 1<<30) + + // An untracked variant file on disk (e.g. written before this + // feature existed) must be adopted into the accounting. + untracked := bytes.Repeat([]byte{0x31}, 1000) + if _, err := cache.variants.Store("aabbccdd0001", bytes.NewReader(untracked), "image/webp"); err != nil { + t.Fatalf("failed to store untracked variant file: %v", err) + } + + // An accounting row whose file is missing must be dropped. + if _, err := cache.db.Exec( + `INSERT INTO variant_content (cache_key, size_bytes, content_type) + VALUES (?, ?, ?)`, + "deadbeef0001", 700, "image/webp", + ); err != nil { + t.Fatalf("failed to insert stale variant accounting row: %v", err) + } + + cache.StartEviction(time.Hour) + defer cache.StopEviction() + + deadline := time.Now().Add(5 * time.Second) + + var usage int64 + + for time.Now().Before(deadline) { + var err error + + usage, err = cache.UsageBytes(context.Background()) + if err != nil { + t.Fatalf("UsageBytes failed: %v", err) + } + + if usage == 1000 { + break + } + + time.Sleep(25 * time.Millisecond) + } + + if usage != 1000 { + t.Errorf("usage after reconciliation = %d, want 1000 "+ + "(untracked file adopted, stale row dropped)", usage) + } + + if n := countRows(t, cache, + `SELECT COUNT(*) FROM variant_content WHERE cache_key = ?`, "aabbccdd0001", + ); n != 1 { + t.Errorf("untracked variant file was not adopted into accounting (rows=%d)", n) + } + + if n := countRows(t, cache, + `SELECT COUNT(*) FROM variant_content WHERE cache_key = ?`, "deadbeef0001", + ); n != 0 { + t.Errorf("stale accounting row without a file was not dropped (rows=%d)", n) + } +} diff --git a/internal/imgcache/storage.go b/internal/imgcache/storage.go index 371138a..9891074 100644 --- a/internal/imgcache/storage.go +++ b/internal/imgcache/storage.go @@ -493,6 +493,24 @@ func (s *VariantStorage) Delete(key VariantKey) error { 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 { + if err := s.Delete(key); err != nil { + return err + } + + metaPath := s.keyToPath(key) + ".meta" + + //nolint:gosec // G703: path derived from cache key + err := os.Remove(metaPath) + if err != nil && !os.IsNotExist(err) { + return fmt.Errorf("failed to delete variant metadata: %w", err) + } + + return nil +} + // keyToPath converts a key to a file path: /// func (s *VariantStorage) keyToPath(key VariantKey) string { k := string(key)