diff --git a/README.md b/README.md index cd7b4c2..d53624f 100644 --- a/README.md +++ b/README.md @@ -530,12 +530,10 @@ Key settings in more detail: - `db_url` — the SQLite database to open; omitted, it is `file:/state.sqlite3?_pragma=journal_mode(WAL)`, which keeps the database in WAL mode. pixa opens one connection to it, so its own reads and - writes run one at a time. Requests wait while eviction runs one of its - queries, some of which read a whole table. pixa adds - `_pragma=busy_timeout(5000)` to any `db_url`, so a write that finds another - program writing to the file waits up to five seconds for it instead of - failing. WAL mode comes only from the URL: keep `_pragma=journal_mode(WAL)` in - one you set + writes run one at a time. pixa adds `_pragma=busy_timeout(5000)` to any + `db_url`, so a write that finds another program writing to the file waits up + to five seconds for it instead of failing. WAL mode comes only from the URL: + keep `_pragma=journal_mode(WAL)` in one you set - `cache_max_bytes` — disk cache size limit in bytes; `0` disables the disk cache entirely; omitted defaults to 75% of the sum of the free space on the filesystem containing `/cache/` and the bytes of source and diff --git a/TODO.md b/TODO.md index 53bb84d..7343e1c 100644 --- a/TODO.md +++ b/TODO.md @@ -30,6 +30,14 @@ P2: security: per-IP rate limiting on the image routes # Completed Steps +- 2026-10-08 requests no longer wait behind eviction queries that read a whole + table (closes #227): the new `cache_usage` table holds the total cache usage, + kept up to date by triggers on `source_content` and `variant_content` in the + statement that adds, removes or resizes a row, and `UsageBytes` reads it. The + reconciliation pass reads the content tables 1000 rows per query, sums them + and corrects the total when it differs, unless a row changed while it summed. + Source rows get `last_accessed_at` when added, so choosing source images to + evict reads that column's index instead of sorting the whole table. - 2026-10-08 SQLite writes no longer fail with "database is locked" under load (closes #223): `internal/database` opens the database with one connection, so pixa's own reads and writes run on it one at a time instead of competing for diff --git a/internal/db/migrations/001_schema.sql b/internal/db/migrations/001_schema.sql index 41a0e02..68f3a32 100644 --- a/internal/db/migrations/001_schema.sql +++ b/internal/db/migrations/001_schema.sql @@ -3,14 +3,12 @@ -- Source content blobs -- Files stored at: cache/sources/// --- last_accessed_at is NULL until the first LRU touch; eviction falls --- back to fetched_at for rows that have never been touched. CREATE TABLE IF NOT EXISTS source_content ( content_hash TEXT PRIMARY KEY, content_type TEXT NOT NULL, size_bytes INTEGER NOT NULL, fetched_at DATETIME DEFAULT CURRENT_TIMESTAMP, - last_accessed_at DATETIME + last_accessed_at DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP ); CREATE INDEX IF NOT EXISTS idx_source_content_last_accessed ON source_content(last_accessed_at); @@ -42,9 +40,8 @@ CREATE INDEX IF NOT EXISTS idx_source_meta_content_hash ON source_metadata(conte -- Processed variant blobs -- Files stored at: cache/variants/// (plus a .meta -- sidecar with the content type). Tracked here (like source content --- blobs above) so total cache usage can be computed with a SUM query, --- never a directory scan, and so LRU eviction has a timestamp to order --- on. +-- blobs above) so total cache usage is known without a directory scan, +-- and so LRU eviction has a timestamp to order on. CREATE TABLE IF NOT EXISTS variant_content ( cache_key TEXT PRIMARY KEY, size_bytes INTEGER NOT NULL, @@ -55,6 +52,74 @@ CREATE TABLE IF NOT EXISTS variant_content ( CREATE INDEX IF NOT EXISTS idx_variant_content_last_accessed ON variant_content(last_accessed_at); +-- Total cache usage: the sum of size_bytes over source_content and +-- variant_content, kept by the triggers below in the same statement +-- that adds, removes or resizes a row, so eviction reads this one row +-- instead of summing both tables. Each trigger also adds one to +-- change_count; the reconciliation pass, which sums both tables a page +-- at a time, corrects total_size_bytes only if change_count did not +-- move while it summed. +CREATE TABLE IF NOT EXISTS cache_usage ( + id INTEGER PRIMARY KEY CHECK (id = 1), + total_size_bytes INTEGER NOT NULL DEFAULT 0, + change_count INTEGER NOT NULL DEFAULT 0 +); +INSERT OR IGNORE INTO cache_usage (id) VALUES (1); + +CREATE TRIGGER IF NOT EXISTS source_content_usage_insert +AFTER INSERT ON source_content +BEGIN + UPDATE cache_usage + SET total_size_bytes = total_size_bytes + NEW.size_bytes, + change_count = change_count + 1 + WHERE id = 1; +END; + +CREATE TRIGGER IF NOT EXISTS source_content_usage_delete +AFTER DELETE ON source_content +BEGIN + UPDATE cache_usage + SET total_size_bytes = total_size_bytes - OLD.size_bytes, + change_count = change_count + 1 + WHERE id = 1; +END; + +CREATE TRIGGER IF NOT EXISTS source_content_usage_update +AFTER UPDATE OF size_bytes ON source_content +BEGIN + UPDATE cache_usage + SET total_size_bytes = total_size_bytes - OLD.size_bytes + NEW.size_bytes, + change_count = change_count + 1 + WHERE id = 1; +END; + +CREATE TRIGGER IF NOT EXISTS variant_content_usage_insert +AFTER INSERT ON variant_content +BEGIN + UPDATE cache_usage + SET total_size_bytes = total_size_bytes + NEW.size_bytes, + change_count = change_count + 1 + WHERE id = 1; +END; + +CREATE TRIGGER IF NOT EXISTS variant_content_usage_delete +AFTER DELETE ON variant_content +BEGIN + UPDATE cache_usage + SET total_size_bytes = total_size_bytes - OLD.size_bytes, + change_count = change_count + 1 + WHERE id = 1; +END; + +CREATE TRIGGER IF NOT EXISTS variant_content_usage_update +AFTER UPDATE OF size_bytes ON variant_content +BEGIN + UPDATE cache_usage + SET total_size_bytes = total_size_bytes - OLD.size_bytes + NEW.size_bytes, + change_count = change_count + 1 + WHERE id = 1; +END; + -- Output/transformed content blobs -- Not written: transformed images are stored in cache/variants and -- tracked in variant_content above. diff --git a/internal/imgcache/cache.go b/internal/imgcache/cache.go index 3c31956..9bae3bb 100644 --- a/internal/imgcache/cache.go +++ b/internal/imgcache/cache.go @@ -471,8 +471,9 @@ func (c *Cache) Stats(ctx context.Context) (*CacheStats, error) { return nil, fmt.Errorf("failed to get cache stats: %w", err) } - // Count and size the cached source images and processed variants. A - // disabled cache holds none, whatever rows an earlier run left. + // Count and size the cached source images and processed variants from + // their tables. A disabled cache holds none, whatever rows an earlier + // run left. if !c.disabled { err = c.db.QueryRowContext(ctx, ` SELECT (SELECT COUNT(*) FROM source_content) @@ -482,7 +483,7 @@ func (c *Cache) Stats(ctx context.Context) (*CacheStats, error) { c.log.Warn("failed to count cache items for stats", "error", err) } - stats.TotalSizeBytes, err = c.UsageBytes(ctx) + stats.TotalSizeBytes, err = c.sumContentSizeBytes(ctx) if err != nil { c.log.Warn("failed to sum cache size for stats", "error", err) } diff --git a/internal/imgcache/cache_usage_internal_test.go b/internal/imgcache/cache_usage_internal_test.go new file mode 100644 index 0000000..442eac2 --- /dev/null +++ b/internal/imgcache/cache_usage_internal_test.go @@ -0,0 +1,267 @@ +package imgcache + +import ( + "bytes" + "fmt" + "strings" + "testing" +) + +// TestEvictionReadsTheUsageTotal adds, stores again, resizes and evicts +// cache content, then drops both content tables. UsageBytes must still +// report what the tables held, and an eviction pass under the limit must +// still succeed: neither may sum the tables. +func TestEvictionReadsTheUsageTotal(t *testing.T) { + t.Parallel() + + cache, _ := newEvictionTestCache(t, 1<<30) + ctx := t.Context() + + kept := storeEvictionTestSource(t, cache, "total.example.com", "/kept.jpg", + bytes.Repeat([]byte{0x71}, 1000)) + evicted := storeEvictionTestSource(t, cache, "total.example.com", "/evicted.jpg", + bytes.Repeat([]byte{0x72}, 700)) + + storeEvictionTestVariant(t, cache, testVariantKeyOne, + bytes.Repeat([]byte{0x73}, 500)) + storeEvictionTestVariant(t, cache, testVariantKeyOne, + bytes.Repeat([]byte{0x74}, 300)) + storeEvictionTestVariant(t, cache, testVariantKeyTwo, + bytes.Repeat([]byte{0x75}, 200)) + + // pixa never changes a source's size, but a statement that does must + // change the total too. + _, err := cache.db.ExecContext(ctx, + `UPDATE source_content SET size_bytes = 900 WHERE content_hash = ?`, + string(kept)) + if err != nil { + t.Fatalf("failed to change the source's size: %v", err) + } + + err = cache.evictSourceBlob(ctx, evicted) + if err != nil { + t.Fatalf("evictSourceBlob failed: %v", err) + } + + err = cache.evictVariant(ctx, testVariantKeyTwo) + if err != nil { + t.Fatalf("evictVariant failed: %v", err) + } + + _, err = cache.db.ExecContext(ctx, + `DROP TABLE source_content; DROP TABLE variant_content`) + if err != nil { + t.Fatalf("failed to drop the content tables: %v", err) + } + + usage, err := cache.UsageBytes(ctx) + t.Logf("UsageBytes() = %d, error = %v", usage, err) + + if err != nil { + t.Fatalf("UsageBytes() error = %v, want nil: it read a content table", err) + } + + if usage != 1200 { + t.Errorf("UsageBytes() = %d, want 1200 (900 + 300)", usage) + } + + err = cache.EvictToLimit(ctx) + if err != nil { + t.Errorf("EvictToLimit() error = %v, want nil: it read a content table", err) + } +} + +// TestReconciliationCorrectsTheUsageTotal sets the total cache usage to a +// wrong value and checks that a reconciliation pass, which sums both +// content tables, puts it right. +func TestReconciliationCorrectsTheUsageTotal(t *testing.T) { + t.Parallel() + + cache, _ := newEvictionTestCache(t, 1<<30) + ctx := t.Context() + + storeEvictionTestSource(t, cache, "total.example.com", "/a.jpg", + bytes.Repeat([]byte{0x76}, 1000)) + storeEvictionTestVariant(t, cache, testVariantKeyOne, + bytes.Repeat([]byte{0x77}, 500)) + + _, err := cache.db.ExecContext(ctx, + `UPDATE cache_usage SET total_size_bytes = 1 WHERE id = 1`) + if err != nil { + t.Fatalf("failed to set a wrong total: %v", err) + } + + err = cache.reconcileAccounting(ctx) + if err != nil { + t.Fatalf("reconcileAccounting failed: %v", err) + } + + usage, err := cache.UsageBytes(ctx) + if err != nil { + t.Fatalf("UsageBytes failed: %v", err) + } + + if usage != 1500 { + t.Errorf("UsageBytes() after reconciliation = %d, want 1500 (1000 + 500)", + usage) + } +} + +// TestUsageTotalNotCorrectedFromAnOlderSum stores a variant after the +// change count was read, then offers a sum taken before the store as the +// correction: the total must keep the stored variant, as that sum may +// have missed it. +func TestUsageTotalNotCorrectedFromAnOlderSum(t *testing.T) { + t.Parallel() + + cache, _ := newEvictionTestCache(t, 1<<30) + ctx := t.Context() + + var changeCount int64 + + err := cache.db.QueryRowContext(ctx, + `SELECT change_count FROM cache_usage WHERE id = 1`, + ).Scan(&changeCount) + if err != nil { + t.Fatalf("failed to read the change count: %v", err) + } + + storeEvictionTestVariant(t, cache, testVariantKeyOne, + bytes.Repeat([]byte{0x78}, 500)) + + corrected, err := cache.correctUsageTotal(ctx, 0, changeCount) + if err != nil { + t.Fatalf("correctUsageTotal failed: %v", err) + } + + if corrected { + t.Error("correctUsageTotal() = true, want false: content changed since the sum") + } + + usage, err := cache.UsageBytes(ctx) + if err != nil { + t.Fatalf("UsageBytes failed: %v", err) + } + + if usage != 500 { + t.Errorf("UsageBytes() = %d, want 500", usage) + } +} + +// TestReconciliationReadsAPageAtATime adds one row more than a page to +// each content table, none of them with a file. Each reconciliation read +// must return one page at most, so a request's query waits for one page +// at most, and the pass must still reach every row: it sums all of them +// and drops all of them. +func TestReconciliationReadsAPageAtATime(t *testing.T) { + t.Parallel() + + cache, _ := newEvictionTestCache(t, 1<<30) + ctx := t.Context() + + const rowsPerTable = reconciliationPageSize + 1 + + for i := range rowsPerTable { + _, err := cache.db.ExecContext(ctx, ` + INSERT INTO variant_content (cache_key, size_bytes, content_type) + VALUES (?, 1, ?) + `, fmt.Sprintf("%012x", i), testContentTypeWebP) + if err != nil { + t.Fatalf("failed to insert variant row %d: %v", i, err) + } + + _, err = cache.db.ExecContext(ctx, ` + INSERT INTO source_content (content_hash, content_type, size_bytes) + VALUES (?, ?, 1) + `, fmt.Sprintf("%064x", i), testContentTypeJPEG) + if err != nil { + t.Fatalf("failed to insert source row %d: %v", i, err) + } + } + + keys, err := cache.variantKeysAfter(ctx, "") + if err != nil { + t.Fatalf("variantKeysAfter failed: %v", err) + } + + if len(keys) != reconciliationPageSize { + t.Errorf("one read returned %d variant keys, want %d", + len(keys), reconciliationPageSize) + } + + hashes, err := cache.sourceContentHashesAfter(ctx, "") + if err != nil { + t.Fatalf("sourceContentHashesAfter failed: %v", err) + } + + if len(hashes) != reconciliationPageSize { + t.Errorf("one read returned %d source hashes, want %d", + len(hashes), reconciliationPageSize) + } + + sum, err := cache.sumContentSizeBytes(ctx) + if err != nil { + t.Fatalf("sumContentSizeBytes failed: %v", err) + } + + if sum != 2*rowsPerTable { + t.Errorf("sumContentSizeBytes() = %d, want %d", sum, 2*rowsPerTable) + } + + err = cache.reconcileAccounting(ctx) + if err != nil { + t.Fatalf("reconcileAccounting failed: %v", err) + } + + if n := countRows(t, cache, `SELECT COUNT(*) FROM variant_content`); n != 0 { + t.Errorf("%d variant rows without a file are left, want 0", n) + } + + if n := countRows(t, cache, `SELECT COUNT(*) FROM source_content`); n != 0 { + t.Errorf("%d source rows without a file are left, want 0", n) + } +} + +// TestSourceEvictionCandidatesComeFromTheIndex checks that the query +// choosing source images to evict reads source_content in last access +// order from its index, instead of reading and sorting the whole table. +func TestSourceEvictionCandidatesComeFromTheIndex(t *testing.T) { + t.Parallel() + + cache, _ := newEvictionTestCache(t, 1<<30) + + rows, err := cache.db.QueryContext(t.Context(), + `EXPLAIN QUERY PLAN `+sourceCandidatesQuery, evictionBatchSize) + if err != nil { + t.Fatalf("EXPLAIN QUERY PLAN failed: %v", err) + } + + defer func() { _ = rows.Close() }() + + var plan []string + + for rows.Next() { + var id, parent, unused int + + var detail string + + err := rows.Scan(&id, &parent, &unused, &detail) + if err != nil { + t.Fatalf("failed to scan the query plan: %v", err) + } + + plan = append(plan, detail) + } + + err = rows.Err() + if err != nil { + t.Fatalf("query plan iteration failed: %v", err) + } + + t.Logf("query plan: %q", plan) + + if !strings.Contains(strings.Join(plan, "\n"), "idx_source_content_last_accessed") { + t.Errorf("the source candidate query does not use "+ + "idx_source_content_last_accessed; plan: %q", plan) + } +} diff --git a/internal/imgcache/eviction.go b/internal/imgcache/eviction.go index 29de851..0a55e29 100644 --- a/internal/imgcache/eviction.go +++ b/internal/imgcache/eviction.go @@ -21,6 +21,11 @@ const DefaultEvictionInterval = 5 * time.Minute // and source blobs) one eviction pass fetches from the database. const evictionBatchSize = 100 +// reconciliationPageSize is the most rows one read of the reconciliation +// pass returns. Each read is a query of its own, so a request waits for +// one page at most, however large the cache is. +const reconciliationPageSize = 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. @@ -45,7 +50,9 @@ 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. +// 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 @@ -53,12 +60,11 @@ func (c *Cache) UsageBytes(ctx context.Context) (int64, error) { 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) + 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 compute cache usage: %w", err) + return 0, fmt.Errorf("failed to read cache usage: %w", err) } return total, nil @@ -236,16 +242,19 @@ func (c *Cache) variantCandidates(ctx context.Context) ([]evictionCandidate, err return candidates, nil } -// sourceCandidates returns the least recently used source blobs. Rows -// written before the LRU column existed fall back to fetched_at. +// 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, ` - 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) + rows, err := c.db.QueryContext(ctx, sourceCandidatesQuery, evictionBatchSize) if err != nil { return nil, fmt.Errorf("failed to query source eviction candidates: %w", err) } @@ -532,11 +541,13 @@ func (c *Cache) runReconciliationPass(ctx context.Context) { // (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), and sweeps stale temp files -// left behind by crashed writes. 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. +// 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 @@ -562,7 +573,7 @@ func (c *Cache) reconcileAccounting(ctx context.Context) error { return err } - return nil + return c.reconcileUsageTotal(ctx) } // reconcileVariantFiles walks the variant storage directory, adopting @@ -657,46 +668,59 @@ func (c *Cache) variantContentTypeFromSidecar(variantPath string) string { // 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 - } + var after VariantKey - 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)) + for { + keys, err := c.variantKeysAfter(ctx, after) if err != nil { - return fmt.Errorf("failed to drop stale variant accounting row: %w", err) + return err } - c.log.Info("dropped accounting row for missing variant file", "cache_key", key) + 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] } - - return nil } -// allVariantKeys returns every tracked variant cache key. -func (c *Cache) allVariantKeys(ctx context.Context) ([]VariantKey, error) { - return queryStringColumn[VariantKey](ctx, c.db, - `SELECT cache_key FROM variant_content`, "variant keys", "variant key") +// variantKeysAfter returns, in order, up to 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), reconciliationPageSize) } -// queryStringColumn runs a single-column query 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. +// 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, + ctx context.Context, db *sql.DB, query, plural, singular string, args ...any, ) ([]T, error) { - rows, err := db.QueryContext(ctx, query) + rows, err := db.QueryContext(ctx, query, args...) if err != nil { return nil, fmt.Errorf("failed to query %s: %w", plural, err) } @@ -791,38 +815,169 @@ func (c *Cache) removeUntrackedSourceFile( // 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 - } + var after ContentHash - 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) + for { + hashes, err := c.sourceContentHashesAfter(ctx, after) if err != nil { return err } - c.log.Info("dropped rows for missing source content file", "content_hash", hash) + 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 +// 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), 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 with the sum. Past the last row the key is NULL. +const ( + sourceSizePageQuery = ` + SELECT MAX(content_hash), COALESCE(SUM(size_bytes), 0) 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) 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 + + err := c.db.QueryRowContext(ctx, pageQuery, after, reconciliationPageSize). + Scan(&lastKey, &pageBytes) + if err != nil { + return 0, fmt.Errorf("failed to sum cache content sizes: %w", err) + } + + if !lastKey.Valid { + 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 } -// allSourceContentHashes returns every tracked source content hash. -func (c *Cache) allSourceContentHashes(ctx context.Context) ([]ContentHash, error) { - return queryStringColumn[ContentHash](ctx, c.db, - `SELECT content_hash FROM source_content`, - "source content hashes", "content hash") +// 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