Files
pixa/internal/imgcache/eviction.go
T
clawbot 7cbcd3957f
check / check (push) Waiting to run
Eviction no longer reads a whole table while requests wait on the database (closes #227)
UsageBytes now reads the new cache_usage row, which triggers on
source_content and variant_content keep up to date in the statement
that adds, removes or resizes a row. The reconciliation pass reads both
tables 1000 rows per query, sums them, and corrects the total when it
differs, unless a row changed while it summed. Source rows now get
last_accessed_at when added, so choosing source images to evict reads
that column's index instead of sorting the whole table. Stats still
sums the tables, now in pages: an existing test drops both tables and
expects that sum to fail.

Model: opus-5-5
2026-10-08 04:28:31 +00:00

1023 lines
28 KiB
Go

package imgcache
import (
"context"
"database/sql"
"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
// defaultReconciliationPageSize is the most rows one read of the
// reconciliation pass returns, unless a test sets
// Cache.reconciliationPageSize smaller. Each read is a query of its own,
// so a request waits for one page at most, however large the cache is.
const defaultReconciliationPageSize = 1000
// 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 the content type given to a variant file that
// has no readable .meta sidecar, when it is served or reconciled.
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 reads the total the database keeps up to date as rows
// are added and removed, so it neither scans the cache directories nor
// sums the tables.
func (c *Cache) UsageBytes(ctx context.Context) (int64, error) {
if c.disabled {
return 0, nil
}
var total int64
err := c.db.QueryRowContext(ctx,
`SELECT total_size_bytes FROM cache_usage WHERE id = 1`,
).Scan(&total)
if err != nil {
return 0, fmt.Errorf("failed to read 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. A
// candidate that fails once ctx is cancelled (every one started after
// that fails at its first database call) ends the batch with ctx's
// error, without a warning.
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
}
err := c.evictCandidate(ctx, candidate)
if err != nil {
if ctx.Err() != nil {
return freed, ctx.Err()
}
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
err := rows.Scan(&key, &candidate.sizeBytes, &candidate.lastAccessedAt)
if err != nil {
return nil, fmt.Errorf("failed to scan variant candidate: %w", err)
}
candidate.cacheKey = VariantKey(key)
candidates = append(candidates, candidate)
}
err = rows.Err()
if err != nil {
return nil, fmt.Errorf("variant candidate iteration failed: %w", err)
}
return candidates, nil
}
// sourceCandidatesQuery selects the least recently used source blobs. It
// orders by the last_accessed_at column itself, not by an expression, so
// SQLite reads the rows in order from that column's index instead of
// sorting the whole table.
const sourceCandidatesQuery = `
SELECT content_hash, size_bytes, last_accessed_at
FROM source_content
ORDER BY last_accessed_at ASC, content_hash ASC
LIMIT ?`
// sourceCandidates returns the least recently used source blobs.
func (c *Cache) sourceCandidates(ctx context.Context) ([]evictionCandidate, error) {
rows, err := c.db.QueryContext(ctx, sourceCandidatesQuery, 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
err := rows.Scan(&hash, &candidate.sizeBytes, &candidate.lastAccessedAt)
if err != nil {
return nil, fmt.Errorf("failed to scan source candidate: %w", err)
}
candidate.contentHash = ContentHash(hash)
candidates = append(candidates, candidate)
}
err = rows.Err()
if 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. The metaCache entry goes before the files; a GetVariant that
// read them just before may put it back, and the next GetVariant then
// fails to open the file and removes it again.
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)
}
c.metaCache.Remove(cacheKey)
err = c.variants.Delete(cacheKey)
if 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.
//
// The whole operation holds the content hash's lock (the same one
// StoreSource holds for its full store), so a concurrent store of
// identical content bytes can never observe the file gone but a row
// still present, or insert a fresh row between this transaction's
// commit and the file unlink below: it either runs entirely before
// this eviction starts, or is blocked until this eviction (row
// deletion and unlink together) has fully completed.
func (c *Cache) evictSourceBlob(ctx context.Context, contentHash ContentHash) error {
unlock := c.contentLocks.Lock(string(contentHash))
defer unlock()
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() }()
_, err = tx.ExecContext(ctx,
`DELETE FROM source_metadata WHERE content_hash = ?`, string(contentHash))
if err != nil {
return fmt.Errorf("failed to delete source metadata rows: %w", err)
}
_, err = tx.ExecContext(ctx,
`DELETE FROM source_content WHERE content_hash = ?`, string(contentHash))
if err != nil {
return fmt.Errorf("failed to delete source content row: %w", err)
}
err = tx.Commit()
if err != nil {
return fmt.Errorf("failed to commit eviction transaction: %w", err)
}
if c.evictSourceBlobTestHook != nil {
c.evictSourceBlobTestHook(contentHash)
}
// Only after the rows are gone may the files be removed.
for _, reference := range references {
err := c.srcMetadata.Delete(reference.host, reference.pathHash)
if err != nil {
c.log.Warn("failed to delete metadata sidecar",
"host", reference.host, "path_hash", reference.pathHash, "error", err)
}
}
err = c.srcContent.Delete(contentHash)
if 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
err := rows.Scan(&reference.host, &pathHash)
if err != nil {
return nil, fmt.Errorf("failed to scan source reference: %w", err)
}
reference.pathHash = PathHash(pathHash)
references = append(references, reference)
}
err = rows.Err()
if 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 at
// startup and again on every periodic tick thereafter, and 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. The goroutine outlives the caller, so it runs
// with its own context, which StopEviction cancels.
func (c *Cache) StartEviction(interval time.Duration) {
if c.disabled || c.evictionCancel != nil {
return
}
ctx, cancel := context.WithCancel(context.Background())
c.evictionCancel = cancel
go c.evictionLoop(ctx, interval)
}
// StopEviction cancels the background eviction goroutine, which
// interrupts a pass in progress, and waits for it to exit or for ctx to
// end, whichever comes first. In the second case it returns an error
// wrapping ctx's error. It is safe to call when eviction was never
// started, and safe to call more than once.
func (c *Cache) StopEviction(ctx context.Context) error {
if c.evictionCancel == nil {
return nil
}
c.evictionCancel()
select {
case <-c.evictionDone:
return nil
case <-ctx.Done():
return fmt.Errorf("cache eviction still running: %w", ctx.Err())
}
}
// evictionLoop is the body of the background eviction goroutine. It
// returns when ctx is cancelled, and starts no pass after that.
func (c *Cache) evictionLoop(ctx context.Context, interval time.Duration) {
defer close(c.evictionDone)
c.runReconciliationPass(ctx)
c.runEvictionPass(ctx)
ticker := time.NewTicker(interval)
defer ticker.Stop()
for {
select {
case <-ctx.Done():
return
case <-ticker.C:
// Reconciliation walks the cache directories, so it only
// runs on the periodic ticker rather than on every
// write-pressure wakeup, keeping it off the per-store hot
// path. Reusing the eviction interval itself (rather than a
// separate, longer one) is a deliberate choice: it is the
// simplest option that still bounds how long a store's
// best-effort accounting insert can stay silently
// unaccounted for to one interval, on a process that is
// already running this loop regardless.
c.runReconciliationPass(ctx)
case <-c.evictionPressure:
}
c.runEvictionPass(ctx)
}
}
// runEvictionPass runs one eviction pass, logging failures instead of
// propagating them (the loop must keep running). It does nothing once
// ctx is cancelled.
func (c *Cache) runEvictionPass(ctx context.Context) {
if ctx.Err() != nil {
return
}
err := c.EvictToLimit(ctx)
if err != nil {
c.log.Warn("cache eviction pass failed", "error", err)
}
}
// runReconciliationPass runs one reconciliation pass, logging failures
// instead of propagating them (the loop must keep running). It does
// nothing once ctx is cancelled.
func (c *Cache) runReconciliationPass(ctx context.Context) {
if ctx.Err() != nil {
return
}
err := c.reconcileAccounting(ctx)
if err != nil {
c.log.Warn("cache accounting reconciliation failed", "error", err)
}
}
// reconcileAccounting synchronizes the database size accounting with
// the actual contents of the cache directories. It runs at startup and
// again on every periodic eviction tick thereafter, off the request
// hot path: it adopts variant files that predate the accounting table
// (or whose accounting insert failed, e.g. StoreVariant's best-effort
// insert under transient DB contention), drops accounting rows whose
// files are missing, removes source blob files the database does not
// know (and rows whose files are gone), sweeps stale temp files left
// behind by crashed writes, and last checks the total cache usage
// against the tables. It reads the tables a page at a time. Running it
// periodically, not just once, bounds how long such drift can
// accumulate unaccounted for on a long-running process to one eviction
// interval. Once ctx is cancelled, it stops at the next file or row and
// returns ctx's error.
func (c *Cache) reconcileAccounting(ctx context.Context) error {
if c.disabled {
return nil
}
err := c.reconcileVariantFiles(ctx)
if err != nil {
return err
}
err = c.reconcileVariantRows(ctx)
if err != nil {
return err
}
err = c.reconcileSourceFiles(ctx)
if err != nil {
return err
}
err = c.reconcileSourceRows(ctx)
if err != nil {
return err
}
return c.reconcileUsageTotal(ctx)
}
// 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 ctx.Err() != nil {
return ctx.Err()
}
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 {
//nolint:gosec // path from cache walk
metaData, err := os.ReadFile(variantPath + variantMetaSuffix)
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 {
var after VariantKey
for {
keys, err := c.variantKeysAfter(ctx, after)
if err != nil {
return err
}
if c.reconciliationReadTestHook != nil {
c.reconciliationReadTestHook(len(keys))
}
if len(keys) == 0 {
return nil
}
for _, key := range keys {
if ctx.Err() != nil {
return ctx.Err()
}
if c.variants.Exists(key) {
continue
}
_, err := c.db.ExecContext(ctx,
`DELETE FROM variant_content WHERE cache_key = ?`, string(key))
if 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)
}
after = keys[len(keys)-1]
}
}
// variantKeysAfter returns, in order, up to c.reconciliationPageSize
// tracked variant cache keys that sort after the given one.
func (c *Cache) variantKeysAfter(
ctx context.Context, after VariantKey,
) ([]VariantKey, error) {
return queryStringColumn[VariantKey](ctx, c.db, `
SELECT cache_key FROM variant_content
WHERE cache_key > ? ORDER BY cache_key LIMIT ?
`, "variant keys", "variant key", string(after), c.reconciliationPageSize)
}
// queryStringColumn runs a single-column query with args and returns the
// column values as T. plural names the set for the query and scan
// failure messages; singular names one row for the scan and iteration
// failure messages.
func queryStringColumn[T ~string](
ctx context.Context, db *sql.DB, query, plural, singular string, args ...any,
) ([]T, error) {
rows, err := db.QueryContext(ctx, query, args...)
if err != nil {
return nil, fmt.Errorf("failed to query %s: %w", plural, err)
}
defer func() { _ = rows.Close() }()
var values []T
for rows.Next() {
var value string
err := rows.Scan(&value)
if err != nil {
return nil, fmt.Errorf("failed to scan %s: %w", singular, err)
}
values = append(values, T(value))
}
err = rows.Err()
if err != nil {
return nil, fmt.Errorf("%s iteration failed: %w", singular, err)
}
return values, 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 ctx.Err() != nil {
return ctx.Err()
}
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
}
_, err = c.db.ExecContext(ctx,
`DELETE FROM source_metadata WHERE content_hash = ?`, string(contentHash))
if err != nil {
return fmt.Errorf("failed to delete metadata rows for untracked blob: %w", err)
}
err = os.Remove(path)
if 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 {
var after ContentHash
for {
hashes, err := c.sourceContentHashesAfter(ctx, after)
if err != nil {
return err
}
if c.reconciliationReadTestHook != nil {
c.reconciliationReadTestHook(len(hashes))
}
if len(hashes) == 0 {
return nil
}
for _, hash := range hashes {
if ctx.Err() != nil {
return ctx.Err()
}
if c.srcContent.Exists(hash) {
continue
}
// The blob file is already gone; evictSourceBlob removes the
// rows and sidecars and tolerates the missing file.
err := c.evictSourceBlob(ctx, hash)
if err != nil {
return err
}
c.log.Info("dropped rows for missing source content file", "content_hash", hash)
}
after = hashes[len(hashes)-1]
}
}
// sourceContentHashesAfter returns, in order, up to
// c.reconciliationPageSize tracked source content hashes that sort after
// the given one.
func (c *Cache) sourceContentHashesAfter(
ctx context.Context, after ContentHash,
) ([]ContentHash, error) {
return queryStringColumn[ContentHash](ctx, c.db, `
SELECT content_hash FROM source_content
WHERE content_hash > ? ORDER BY content_hash LIMIT ?
`, "source content hashes", "content hash", string(after),
c.reconciliationPageSize)
}
// Each of these queries sums size_bytes over the next page of rows of
// one content table, the rows that sort after a key, and returns the
// page's last key, the sum and the number of rows in the page. Past the
// last row the page has no rows and the key is NULL.
const (
sourceSizePageQuery = `
SELECT MAX(content_hash), COALESCE(SUM(size_bytes), 0), COUNT(*)
FROM (
SELECT content_hash, size_bytes FROM source_content
WHERE content_hash > ? ORDER BY content_hash LIMIT ?
)`
variantSizePageQuery = `
SELECT MAX(cache_key), COALESCE(SUM(size_bytes), 0), COUNT(*)
FROM (
SELECT cache_key, size_bytes FROM variant_content
WHERE cache_key > ? ORDER BY cache_key LIMIT ?
)`
)
// sumContentSizeBytes sums size_bytes over both content tables, a page
// of rows per query.
func (c *Cache) sumContentSizeBytes(ctx context.Context) (int64, error) {
sourceBytes, err := c.sumSizeBytesInPages(ctx, sourceSizePageQuery)
if err != nil {
return 0, err
}
variantBytes, err := c.sumSizeBytesInPages(ctx, variantSizePageQuery)
if err != nil {
return 0, err
}
return sourceBytes + variantBytes, nil
}
// sumSizeBytesInPages runs pageQuery, one of the size page queries
// above, from the first page to the last and adds up the page sums.
func (c *Cache) sumSizeBytesInPages(
ctx context.Context, pageQuery string,
) (int64, error) {
var total int64
after := ""
for {
var lastKey sql.NullString
var pageBytes int64
var pageRows int
err := c.db.QueryRowContext(ctx, pageQuery, after, c.reconciliationPageSize).
Scan(&lastKey, &pageBytes, &pageRows)
if err != nil {
return 0, fmt.Errorf("failed to sum cache content sizes: %w", err)
}
if c.reconciliationReadTestHook != nil {
c.reconciliationReadTestHook(pageRows)
}
if pageRows == 0 {
return total, nil
}
total += pageBytes
after = lastKey.String
}
}
// reconcileUsageTotal checks the total cache usage the database keeps
// against size_bytes summed over both content tables, and corrects the
// total when they differ.
func (c *Cache) reconcileUsageTotal(ctx context.Context) error {
var totalBytes, changeCount int64
err := c.db.QueryRowContext(ctx,
`SELECT total_size_bytes, change_count FROM cache_usage WHERE id = 1`,
).Scan(&totalBytes, &changeCount)
if err != nil {
return fmt.Errorf("failed to read cache usage: %w", err)
}
sumBytes, err := c.sumContentSizeBytes(ctx)
if err != nil {
return err
}
if sumBytes == totalBytes {
return nil
}
corrected, err := c.correctUsageTotal(ctx, sumBytes, changeCount)
if err != nil || !corrected {
return err
}
c.log.Warn("corrected total cache usage to the sum of the content tables",
"previous_usage_bytes", totalBytes, "usage_bytes", sumBytes)
return nil
}
// correctUsageTotal sets the total cache usage to sumBytes, a sum of the
// content tables taken when the change count was changeCount, and
// reports whether it did. If a row was added, removed or resized since,
// the count has moved and the total is left alone: the sum may have
// missed that change, and the next pass checks again.
func (c *Cache) correctUsageTotal(
ctx context.Context, sumBytes, changeCount int64,
) (bool, error) {
result, err := c.db.ExecContext(ctx, `
UPDATE cache_usage SET total_size_bytes = ?
WHERE id = 1 AND change_count = ?
`, sumBytes, changeCount)
if err != nil {
return false, fmt.Errorf("failed to correct cache usage: %w", err)
}
affected, err := result.RowsAffected()
if err != nil {
return false, fmt.Errorf("failed to read cache usage correction: %w", err)
}
return affected > 0, 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
}
err = os.Remove(path)
if 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)
}