2 Commits
Author SHA1 Message Date
clawbot 7cbcd3957f Eviction no longer reads a whole table while requests wait on the database (closes #227)
check / check (push) Canceled after 0s
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
clawbot 54328377d0 Install libvips JPEG XL support and require it at startup (part of #222)
check / check (push) Canceled after 0s
script/bootstrap --cgo installs vips-jxl when the package manager it
uses is apk, as Alpine's vips package lacks JPEG XL support; the nix
and brew libvips include it, as does apt's from Debian 12 and Ubuntu
24.04 on. The runtime stage of the Dockerfile installs vips-jxl too.

imgcache.NewService calls the new imageprocessor.CheckJPEGXLSupport
before it builds the image processor and fails with its error, which
names the fix, so pixad does not start without the support. JPEG XL
becomes the default output later in the issue. New tests check the
support and save an image as JPEG XL with govips and load it back.

Model: opus-5-5
2026-10-08 06:19:05 +02:00
12 changed files with 725 additions and 99 deletions
+6 -3
View File
@@ -22,8 +22,9 @@ FROM golang:1.25.4-alpine@sha256:d3f0cf7723f3429e3f9ed846243970b20a2de7bae6a5b66
WORKDIR /src
# script/bootstrap --cgo installs the build dependencies (a C compiler
# and the libvips and libheif headers) and downloads the Go modules.
# script/bootstrap --cgo installs the build dependencies (a C compiler,
# the libvips and libheif headers, and libvips' JPEG XL support, which
# the tests need) and downloads the Go modules.
COPY script/ ./script/
COPY go.mod go.sum ./
RUN script/bootstrap --cgo
@@ -80,9 +81,11 @@ RUN version="${VERSION:-$(git describe --tags --always)}"; \
# alpine:3.21, 2026-02-25
FROM alpine:3.21@sha256:c3f8e73fdb79deaebaa2037150150191b9dcbfba68b4a46d70103204c53f4709
# Install runtime dependencies only
# Install runtime dependencies only. vips-jxl is libvips' JPEG XL
# support, without which pixad does not start.
RUN apk add --no-cache \
vips \
vips-jxl \
libheif \
ca-certificates \
tzdata \
+10 -9
View File
@@ -89,7 +89,10 @@ another part of pixa failed to stop. A request not finished by then is cut off.
`docker stop` waits 10 seconds before it kills the container.
Outside Docker, pixa needs libvips (the image has 8.15) and libheif to run, as
it uses libvips through CGO; building it also needs their development files,
it uses libvips through CGO. pixad does not start unless libvips has its JPEG XL
support, which on Alpine is the `vips-jxl` package and which the nix and brew
packages of libvips include, as do the apt ones from Debian 12 and Ubuntu 24.04
on. Building pixa also needs the development files of libvips and libheif,
`pkg-config` and a C compiler. `script/bootstrap --cgo` installs all of these,
as the `Dockerfile` does where it compiles pixa. Plain `script/bootstrap`, which
`script/setup` and `script/cibuild` run, installs git, make and Go, and Node,
@@ -530,12 +533,10 @@ Key settings in more detail:
- `db_url` — the SQLite database to open; omitted, it is
`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
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 `<state_dir>/cache/` and the bytes of source and
@@ -585,8 +586,8 @@ provide:
- `script/bootstrap` — install git, make, Go, Node, Yarn and prettier and
download the Go modules (idempotent); with `--cgo`, the C compiler and the
libvips and libheif libraries that compiling pixa needs instead of Node, Yarn
and prettier
libvips (with its JPEG XL support) and libheif libraries that compiling and
testing pixa need instead of Node, Yarn and prettier
- `script/setup` — make a fresh clone ready for development (bootstrap, then
install-precommit)
- `script/projectname` — output the project name ("pixa")
+16
View File
@@ -30,6 +30,22 @@ P2: security: per-IP rate limiting on the image routes
# Completed Steps
- 2026-10-08 libvips' JPEG XL support is installed and required (part of #222):
`script/bootstrap --cgo` installs `vips-jxl` when its package manager is apk,
as Alpine's `vips` package lacks the support, and the runtime stage of the
`Dockerfile` installs it too. `imgcache.NewService` fails, naming the fix,
when `imageprocessor.CheckJPEGXLSupport` finds that libvips cannot load and
save JPEG XL, so pixad does not start without it. A test saves an image as
JPEG XL with govips and loads it back. JPEG XL is not yet a format pixa
serves.
- 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
+71 -6
View File
@@ -3,14 +3,12 @@
-- Source content blobs
-- 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 (
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/<ab>/<cd>/<cache_key> (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.
+18
View File
@@ -37,6 +37,24 @@ func initVips() {
})
}
// errNoJPEGXL is returned by CheckJPEGXLSupport.
var errNoJPEGXL = errors.New("libvips lacks JPEG XL support: install " +
"vips-jxl on Alpine, or use a libvips built with libjxl")
// CheckJPEGXLSupport returns an error, naming the fix, when libvips
// cannot load and save JPEG XL.
func CheckJPEGXLSupport() error {
initVips()
// govips counts a format as supported when libvips has its loader;
// libvips builds the JPEG XL loader and saver together.
if !vips.IsTypeSupported(vips.ImageTypeJXL) {
return errNoJPEGXL
}
return nil
}
// Format represents supported output image formats.
type Format string
@@ -0,0 +1,52 @@
package imageprocessor
import (
"testing"
"github.com/davidbyttow/govips/v2/vips"
)
// TestCheckJPEGXLSupport fails when libvips lacks JPEG XL support, as on
// Alpine without the vips-jxl package.
func TestCheckJPEGXLSupport(t *testing.T) {
t.Parallel()
err := CheckJPEGXLSupport()
if err != nil {
t.Fatalf("CheckJPEGXLSupport() error = %v", err)
}
}
// TestLibvipsSavesAndLoadsJPEGXL saves an image as JPEG XL with govips and
// loads it back. It fails when libvips lacks JPEG XL support, as on Alpine
// without the vips-jxl package.
func TestLibvipsSavesAndLoadsJPEGXL(t *testing.T) {
t.Parallel()
img, err := vips.NewImageFromBuffer(createTestJPEG(t, 64, 48))
if err != nil {
t.Fatalf("failed to load test JPEG: %v", err)
}
defer img.Close()
jxl, _, err := img.ExportJxl(vips.NewJxlExportParams())
if err != nil {
t.Fatalf("ExportJxl() error = %v", err)
}
loaded, err := vips.NewImageFromBuffer(jxl)
if err != nil {
t.Fatalf("failed to load the JPEG XL image: %v", err)
}
defer loaded.Close()
if loaded.Format() != vips.ImageTypeJXL {
t.Errorf("loaded format = %s, want jxl", vips.ImageTypes[loaded.Format()])
}
if loaded.Width() != 64 || loaded.Height() != 48 {
t.Errorf("loaded size = %dx%d, want 64x48", loaded.Width(), loaded.Height())
}
}
+17 -3
View File
@@ -96,6 +96,17 @@ type Cache struct {
// deterministically pause inside that window to exercise
// concurrent stores against it; production code leaves it nil.
evictSourceBlobTestHook func(ContentHash)
// reconciliationPageSize is the most rows one read of a content table
// returns in the reconciliation pass and in Stats. newCache sets it to
// defaultReconciliationPageSize; tests set it smaller.
reconciliationPageSize int
// reconciliationReadTestHook, when set, is called after each of those
// reads with the number of rows the read covered, so tests can check
// that no read covers more than one page; production code leaves it
// nil.
reconciliationReadTestHook func(rows int)
}
// NewCache creates a new cache instance.
@@ -127,6 +138,8 @@ func newCache(
evictionDone: make(chan struct{}),
metaCache: metaCache,
contentLocks: newContentLock(),
reconciliationPageSize: defaultReconciliationPageSize,
}
if c.disabled {
@@ -471,8 +484,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 +496,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)
}
@@ -0,0 +1,265 @@
package imgcache
import (
"bytes"
"strconv"
"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 stores five source images and five
// variants, sets the page size to two rows and runs a reconciliation
// pass. No read the pass makes of a content table, to check its rows or
// to sum them, may cover more than two rows, so a request's query waits
// for one page at most. Between them the reads must still cover every
// row of both tables twice, once to check it and once to sum it, and the
// sum must put a wrong total right.
func TestReconciliationReadsAPageAtATime(t *testing.T) {
t.Parallel()
cache, _ := newEvictionTestCache(t, 1<<30)
ctx := t.Context()
const pageSize = 2
cache.reconciliationPageSize = pageSize
// Sources of 100 bytes and variants of 10, five of each.
for i := range 5 {
content := []byte(strconv.Itoa(i))
storeEvictionTestSource(t, cache, "pages.example.com",
"/"+strconv.Itoa(i)+".jpg", bytes.Repeat(content, 100))
storeEvictionTestVariant(t, cache, VariantKey("aabbccdd000"+strconv.Itoa(i)),
bytes.Repeat(content, 10))
}
_, 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)
}
var rowsPerRead []int
cache.reconciliationReadTestHook = func(rows int) {
rowsPerRead = append(rowsPerRead, rows)
}
err = cache.reconcileAccounting(ctx)
if err != nil {
t.Fatalf("reconcileAccounting failed: %v", err)
}
t.Logf("rows covered by each read, in order: %v", rowsPerRead)
coveredRows := 0
for _, rows := range rowsPerRead {
if rows > pageSize {
t.Errorf("a read covered %d rows, want at most %d", rows, pageSize)
}
coveredRows += rows
}
// Ten rows, each read once to check it and once to sum it.
if coveredRows != 20 {
t.Errorf("the reads covered %d rows in all, want 20", coveredRows)
}
usage, err := cache.UsageBytes(ctx)
if err != nil {
t.Fatalf("UsageBytes failed: %v", err)
}
if usage != 550 {
t.Errorf("UsageBytes() after reconciliation = %d, want 550 (5*100 + 5*10)",
usage)
}
}
// 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)
}
}
+1 -1
View File
@@ -72,7 +72,7 @@ func (c *Cache) computeDefaultMaxBytes(
}
// Both terms are at most math.MaxInt64, so the sum cannot overflow.
//nolint:gosec // G115: UsageBytes sums file sizes, never negative
//nolint:gosec // G115: UsageBytes returns the total cache usage, never negative
spaceBytes := min(freeBytes, math.MaxInt64) + uint64(usedBytes)
computed := spaceBytes / freeSpaceFractionDenominator * freeSpaceFractionNumerator
+246 -72
View File
@@ -21,6 +21,12 @@ const DefaultEvictionInterval = 5 * time.Minute
// 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.
@@ -45,7 +51,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 +61,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 +243,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 +542,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 +574,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 +669,63 @@ 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 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]
}
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 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 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 +820,183 @@ 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 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
}
// 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
+7
View File
@@ -100,6 +100,13 @@ func NewService(cfg *ServiceConfig) (*Service, error) {
allowHTTP = cfg.FetcherConfig.AllowHTTP
}
// JPEG XL is to become the default output format, so pixad does not
// start without it.
err := imageprocessor.CheckJPEGXLSupport()
if err != nil {
return nil, err
}
maxResponseSize := fetcherCfg.MaxResponseSize
processor := imageprocessor.New(imageprocessor.Params{
MaxInputBytes: maxResponseSize,
+16 -5
View File
@@ -14,11 +14,12 @@
# script/fmt-check: all the host needs, as
# the checks compile pixa in Docker
# script/bootstrap --cgo git, make, Go, and a C compiler and the
# CGO image libraries (pkg-config, vips,
# libheif) for the govips bindings instead
# of Node: to compile pixa, in the
# Dockerfile's test phase and build stage,
# which format nothing
# CGO image libraries (pkg-config, vips
# with its JPEG XL support, libheif) for
# the govips bindings instead of Node: to
# compile pixa, in the Dockerfile's test
# phase and build stage, which format
# nothing
set -eu
ROOT="$(cd "$(dirname "$0")/.." && pwd -P)"
@@ -150,6 +151,16 @@ ensure_cgo_deps() {
if ! pkg-config --exists vips; then
pkg_install vips libvips-dev vips vips-dev
fi
# libvips' JPEG XL loader and saver are in the nix and brew vips
# packages, and in apt's from Debian 12 and Ubuntu 24.04 on, but in
# the package vips-jxl on Alpine. detect_pkgmgr is called only where
# apk exists, as on apt it updates the package lists.
if ! missing apk; then
detect_pkgmgr
fi
if [ "$PKGMGR" = "apk" ] && ! apk info -e vips-jxl >/dev/null; then
apk add --no-cache vips-jxl
fi
if ! pkg-config --exists libheif; then
pkg_install libheif libheif-dev libheif libheif-dev
fi