Eviction no longer reads a whole table while requests wait on the database (closes #227)
check / check (push) Waiting to run
check / check (push) Waiting to run
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
This commit is contained in:
@@ -530,12 +530,10 @@ Key settings in more detail:
|
|||||||
- `db_url` — the SQLite database to open; omitted, it is
|
- `db_url` — the SQLite database to open; omitted, it is
|
||||||
`file:<state_dir>/state.sqlite3?_pragma=journal_mode(WAL)`, which keeps the
|
`file:<state_dir>/state.sqlite3?_pragma=journal_mode(WAL)`, which keeps the
|
||||||
database in WAL mode. pixa opens one connection to it, so its own reads and
|
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
|
writes run one at a time. pixa adds `_pragma=busy_timeout(5000)` to any
|
||||||
queries, some of which read a whole table. pixa adds
|
`db_url`, so a write that finds another program writing to the file waits up
|
||||||
`_pragma=busy_timeout(5000)` to any `db_url`, so a write that finds another
|
to five seconds for it instead of failing. WAL mode comes only from the URL:
|
||||||
program writing to the file waits up to five seconds for it instead of
|
keep `_pragma=journal_mode(WAL)` in one you set
|
||||||
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_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
|
cache entirely; omitted defaults to 75% of the sum of the free space on the
|
||||||
filesystem containing `<state_dir>/cache/` and the bytes of source and
|
filesystem containing `<state_dir>/cache/` and the bytes of source and
|
||||||
|
|||||||
@@ -30,6 +30,14 @@ P2: security: per-IP rate limiting on the image routes
|
|||||||
|
|
||||||
# Completed Steps
|
# 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
|
- 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
|
(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
|
pixa's own reads and writes run on it one at a time instead of competing for
|
||||||
|
|||||||
@@ -3,14 +3,12 @@
|
|||||||
|
|
||||||
-- Source content blobs
|
-- Source content blobs
|
||||||
-- Files stored at: cache/sources/<ab>/<cd>/<sha256>
|
-- Files stored at: cache/sources/<ab>/<cd>/<sha256>
|
||||||
-- 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 (
|
CREATE TABLE IF NOT EXISTS source_content (
|
||||||
content_hash TEXT PRIMARY KEY,
|
content_hash TEXT PRIMARY KEY,
|
||||||
content_type TEXT NOT NULL,
|
content_type TEXT NOT NULL,
|
||||||
size_bytes INTEGER NOT NULL,
|
size_bytes INTEGER NOT NULL,
|
||||||
fetched_at DATETIME DEFAULT CURRENT_TIMESTAMP,
|
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
|
CREATE INDEX IF NOT EXISTS idx_source_content_last_accessed
|
||||||
ON source_content(last_accessed_at);
|
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
|
-- Processed variant blobs
|
||||||
-- Files stored at: cache/variants/<ab>/<cd>/<cache_key> (plus a .meta
|
-- Files stored at: cache/variants/<ab>/<cd>/<cache_key> (plus a .meta
|
||||||
-- sidecar with the content type). Tracked here (like source content
|
-- sidecar with the content type). Tracked here (like source content
|
||||||
-- blobs above) so total cache usage can be computed with a SUM query,
|
-- blobs above) so total cache usage is known without a directory scan,
|
||||||
-- never a directory scan, and so LRU eviction has a timestamp to order
|
-- and so LRU eviction has a timestamp to order on.
|
||||||
-- on.
|
|
||||||
CREATE TABLE IF NOT EXISTS variant_content (
|
CREATE TABLE IF NOT EXISTS variant_content (
|
||||||
cache_key TEXT PRIMARY KEY,
|
cache_key TEXT PRIMARY KEY,
|
||||||
size_bytes INTEGER NOT NULL,
|
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
|
CREATE INDEX IF NOT EXISTS idx_variant_content_last_accessed
|
||||||
ON variant_content(last_accessed_at);
|
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
|
-- Output/transformed content blobs
|
||||||
-- Not written: transformed images are stored in cache/variants and
|
-- Not written: transformed images are stored in cache/variants and
|
||||||
-- tracked in variant_content above.
|
-- tracked in variant_content above.
|
||||||
|
|||||||
@@ -471,8 +471,9 @@ func (c *Cache) Stats(ctx context.Context) (*CacheStats, error) {
|
|||||||
return nil, fmt.Errorf("failed to get cache stats: %w", err)
|
return nil, fmt.Errorf("failed to get cache stats: %w", err)
|
||||||
}
|
}
|
||||||
|
|
||||||
// Count and size the cached source images and processed variants. A
|
// Count and size the cached source images and processed variants from
|
||||||
// disabled cache holds none, whatever rows an earlier run left.
|
// their tables. A disabled cache holds none, whatever rows an earlier
|
||||||
|
// run left.
|
||||||
if !c.disabled {
|
if !c.disabled {
|
||||||
err = c.db.QueryRowContext(ctx, `
|
err = c.db.QueryRowContext(ctx, `
|
||||||
SELECT (SELECT COUNT(*) FROM source_content)
|
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)
|
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 {
|
if err != nil {
|
||||||
c.log.Warn("failed to sum cache size for stats", "error", err)
|
c.log.Warn("failed to sum cache size for stats", "error", err)
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -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)
|
||||||
|
}
|
||||||
|
}
|
||||||
+194
-39
@@ -21,6 +21,11 @@ const DefaultEvictionInterval = 5 * time.Minute
|
|||||||
// and source blobs) one eviction pass fetches from the database.
|
// and source blobs) one eviction pass fetches from the database.
|
||||||
const evictionBatchSize = 100
|
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
|
// staleTempFileAge is how old an orphaned temp file (left behind by a
|
||||||
// crashed write) must be before reconciliation removes it. Fresh temp
|
// crashed write) must be before reconciliation removes it. Fresh temp
|
||||||
// files may still belong to an in-flight store.
|
// 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
|
// UsageBytes returns the total number of bytes of cache content
|
||||||
// tracked in the database (source content blobs plus processed
|
// 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) {
|
func (c *Cache) UsageBytes(ctx context.Context) (int64, error) {
|
||||||
if c.disabled {
|
if c.disabled {
|
||||||
return 0, nil
|
return 0, nil
|
||||||
@@ -53,12 +60,11 @@ func (c *Cache) UsageBytes(ctx context.Context) (int64, error) {
|
|||||||
|
|
||||||
var total int64
|
var total int64
|
||||||
|
|
||||||
err := c.db.QueryRowContext(ctx, `
|
err := c.db.QueryRowContext(ctx,
|
||||||
SELECT (SELECT COALESCE(SUM(size_bytes), 0) FROM source_content)
|
`SELECT total_size_bytes FROM cache_usage WHERE id = 1`,
|
||||||
+ (SELECT COALESCE(SUM(size_bytes), 0) FROM variant_content)
|
).Scan(&total)
|
||||||
`).Scan(&total)
|
|
||||||
if err != nil {
|
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
|
return total, nil
|
||||||
@@ -236,16 +242,19 @@ func (c *Cache) variantCandidates(ctx context.Context) ([]evictionCandidate, err
|
|||||||
return candidates, nil
|
return candidates, nil
|
||||||
}
|
}
|
||||||
|
|
||||||
// sourceCandidates returns the least recently used source blobs. Rows
|
// sourceCandidatesQuery selects the least recently used source blobs. It
|
||||||
// written before the LRU column existed fall back to fetched_at.
|
// orders by the last_accessed_at column itself, not by an expression, so
|
||||||
func (c *Cache) sourceCandidates(ctx context.Context) ([]evictionCandidate, error) {
|
// SQLite reads the rows in order from that column's index instead of
|
||||||
rows, err := c.db.QueryContext(ctx, `
|
// sorting the whole table.
|
||||||
SELECT content_hash, size_bytes,
|
const sourceCandidatesQuery = `
|
||||||
COALESCE(last_accessed_at, fetched_at, '1970-01-01 00:00:00') AS lru
|
SELECT content_hash, size_bytes, last_accessed_at
|
||||||
FROM source_content
|
FROM source_content
|
||||||
ORDER BY lru ASC, content_hash ASC
|
ORDER BY last_accessed_at ASC, content_hash ASC
|
||||||
LIMIT ?
|
LIMIT ?`
|
||||||
`, evictionBatchSize)
|
|
||||||
|
// 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 {
|
if err != nil {
|
||||||
return nil, fmt.Errorf("failed to query source eviction candidates: %w", err)
|
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
|
// (or whose accounting insert failed, e.g. StoreVariant's best-effort
|
||||||
// insert under transient DB contention), drops accounting rows whose
|
// insert under transient DB contention), drops accounting rows whose
|
||||||
// files are missing, removes source blob files the database does not
|
// files are missing, removes source blob files the database does not
|
||||||
// know (and rows whose files are gone), and sweeps stale temp files
|
// know (and rows whose files are gone), sweeps stale temp files left
|
||||||
// left behind by crashed writes. Running it periodically, not just
|
// behind by crashed writes, and last checks the total cache usage
|
||||||
// once, bounds how long such drift can accumulate unaccounted for on a
|
// against the tables. It reads the tables a page at a time. Running it
|
||||||
// long-running process to one eviction interval. Once ctx is cancelled,
|
// periodically, not just once, bounds how long such drift can
|
||||||
// it stops at the next file or row and returns ctx's error.
|
// 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 {
|
func (c *Cache) reconcileAccounting(ctx context.Context) error {
|
||||||
if c.disabled {
|
if c.disabled {
|
||||||
return nil
|
return nil
|
||||||
@@ -562,7 +573,7 @@ func (c *Cache) reconcileAccounting(ctx context.Context) error {
|
|||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
|
|
||||||
return nil
|
return c.reconcileUsageTotal(ctx)
|
||||||
}
|
}
|
||||||
|
|
||||||
// reconcileVariantFiles walks the variant storage directory, adopting
|
// reconcileVariantFiles walks the variant storage directory, adopting
|
||||||
@@ -657,11 +668,18 @@ func (c *Cache) variantContentTypeFromSidecar(variantPath string) string {
|
|||||||
// reconcileVariantRows drops accounting rows whose variant files are
|
// reconcileVariantRows drops accounting rows whose variant files are
|
||||||
// missing, so the database never references deleted content.
|
// missing, so the database never references deleted content.
|
||||||
func (c *Cache) reconcileVariantRows(ctx context.Context) error {
|
func (c *Cache) reconcileVariantRows(ctx context.Context) error {
|
||||||
keys, err := c.allVariantKeys(ctx)
|
var after VariantKey
|
||||||
|
|
||||||
|
for {
|
||||||
|
keys, err := c.variantKeysAfter(ctx, after)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
|
|
||||||
|
if len(keys) == 0 {
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
|
||||||
for _, key := range keys {
|
for _, key := range keys {
|
||||||
if ctx.Err() != nil {
|
if ctx.Err() != nil {
|
||||||
return ctx.Err()
|
return ctx.Err()
|
||||||
@@ -680,23 +698,29 @@ func (c *Cache) reconcileVariantRows(ctx context.Context) error {
|
|||||||
c.log.Info("dropped accounting row for missing variant file", "cache_key", key)
|
c.log.Info("dropped accounting row for missing variant file", "cache_key", key)
|
||||||
}
|
}
|
||||||
|
|
||||||
return nil
|
after = keys[len(keys)-1]
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
// allVariantKeys returns every tracked variant cache key.
|
// variantKeysAfter returns, in order, up to reconciliationPageSize
|
||||||
func (c *Cache) allVariantKeys(ctx context.Context) ([]VariantKey, error) {
|
// tracked variant cache keys that sort after the given one.
|
||||||
return queryStringColumn[VariantKey](ctx, c.db,
|
func (c *Cache) variantKeysAfter(
|
||||||
`SELECT cache_key FROM variant_content`, "variant keys", "variant key")
|
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
|
// queryStringColumn runs a single-column query with args and returns the
|
||||||
// values as T. plural names the set for the query and scan failure
|
// column values as T. plural names the set for the query and scan
|
||||||
// messages; singular names one row for the scan and iteration failure
|
// failure messages; singular names one row for the scan and iteration
|
||||||
// messages.
|
// failure messages.
|
||||||
func queryStringColumn[T ~string](
|
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) {
|
) ([]T, error) {
|
||||||
rows, err := db.QueryContext(ctx, query)
|
rows, err := db.QueryContext(ctx, query, args...)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return nil, fmt.Errorf("failed to query %s: %w", plural, err)
|
return nil, fmt.Errorf("failed to query %s: %w", plural, err)
|
||||||
}
|
}
|
||||||
@@ -791,11 +815,18 @@ func (c *Cache) removeUntrackedSourceFile(
|
|||||||
// reconcileSourceRows removes source_content rows (and their metadata
|
// reconcileSourceRows removes source_content rows (and their metadata
|
||||||
// references and sidecars) whose blob files are missing on disk.
|
// references and sidecars) whose blob files are missing on disk.
|
||||||
func (c *Cache) reconcileSourceRows(ctx context.Context) error {
|
func (c *Cache) reconcileSourceRows(ctx context.Context) error {
|
||||||
hashes, err := c.allSourceContentHashes(ctx)
|
var after ContentHash
|
||||||
|
|
||||||
|
for {
|
||||||
|
hashes, err := c.sourceContentHashesAfter(ctx, after)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
|
|
||||||
|
if len(hashes) == 0 {
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
|
||||||
for _, hash := range hashes {
|
for _, hash := range hashes {
|
||||||
if ctx.Err() != nil {
|
if ctx.Err() != nil {
|
||||||
return ctx.Err()
|
return ctx.Err()
|
||||||
@@ -815,14 +846,138 @@ func (c *Cache) reconcileSourceRows(ctx context.Context) error {
|
|||||||
c.log.Info("dropped rows for missing source content file", "content_hash", hash)
|
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
|
return nil
|
||||||
}
|
}
|
||||||
|
|
||||||
// allSourceContentHashes returns every tracked source content hash.
|
// correctUsageTotal sets the total cache usage to sumBytes, a sum of the
|
||||||
func (c *Cache) allSourceContentHashes(ctx context.Context) ([]ContentHash, error) {
|
// content tables taken when the change count was changeCount, and
|
||||||
return queryStringColumn[ContentHash](ctx, c.db,
|
// reports whether it did. If a row was added, removed or resized since,
|
||||||
`SELECT content_hash FROM source_content`,
|
// the count has moved and the total is left alone: the sum may have
|
||||||
"source content hashes", "content hash")
|
// 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
|
// sweepStaleTempFile removes a temp file left behind by a crashed
|
||||||
|
|||||||
Reference in New Issue
Block a user