Cache.metaCache was declared and never used, so every hit read and parsed the variant's .meta file. It is now an LRU of up to 10,000 content types (hashicorp/golang-lru/v2), filled by StoreVariant and by GetVariant after it reads a .meta file. For a variant it holds, Lookup skips the disk check and GetVariant skips the .meta read; the variant file is still opened and its size taken from it. Eviction removes the entry before deleting the files, and GetVariant removes it when the file will not open, so a missing variant is never served. The cap is a constant, not a setting. README.md describes it. Model: opus-5-5
808 lines
22 KiB
Go
808 lines
22 KiB
Go
package imgcache
|
|
|
|
import (
|
|
"context"
|
|
"database/sql"
|
|
"encoding/json"
|
|
"fmt"
|
|
"io/fs"
|
|
"os"
|
|
"path/filepath"
|
|
"strings"
|
|
"time"
|
|
)
|
|
|
|
// DefaultEvictionInterval is how often the background evictor checks
|
|
// cache usage against the configured limit, in addition to the
|
|
// write-pressure wakeups triggered by stores.
|
|
const DefaultEvictionInterval = 5 * time.Minute
|
|
|
|
// evictionBatchSize is how many LRU candidates of each class (variants
|
|
// and source blobs) one eviction pass fetches from the database.
|
|
const evictionBatchSize = 100
|
|
|
|
// staleTempFileAge is how old an orphaned temp file (left behind by a
|
|
// crashed write) must be before reconciliation removes it. Fresh temp
|
|
// files may still belong to an in-flight store.
|
|
const staleTempFileAge = time.Hour
|
|
|
|
// sqliteTimestampLayout matches SQLite's CURRENT_TIMESTAMP format, so
|
|
// timestamps written by reconciliation order correctly against ones
|
|
// written by the hot path.
|
|
const sqliteTimestampLayout = "2006-01-02 15:04:05"
|
|
|
|
// tempFilePrefix is the prefix os.CreateTemp uses for in-flight cache
|
|
// writes (".tmp-*" patterns in the storage layer).
|
|
const tempFilePrefix = ".tmp-"
|
|
|
|
// variantMetaSuffix is the sidecar suffix VariantStorage writes next
|
|
// to each variant file.
|
|
const variantMetaSuffix = ".meta"
|
|
|
|
// fallbackContentType is recorded when a reconciled variant file has
|
|
// no readable .meta sidecar.
|
|
const fallbackContentType = "application/octet-stream"
|
|
|
|
// UsageBytes returns the total number of bytes of cache content
|
|
// tracked in the database (source content blobs plus processed
|
|
// variants). It never scans the cache directories.
|
|
func (c *Cache) UsageBytes(ctx context.Context) (int64, error) {
|
|
if c.disabled {
|
|
return 0, nil
|
|
}
|
|
|
|
var total int64
|
|
|
|
err := c.db.QueryRowContext(ctx, `
|
|
SELECT (SELECT COALESCE(SUM(size_bytes), 0) FROM source_content)
|
|
+ (SELECT COALESCE(SUM(size_bytes), 0) FROM variant_content)
|
|
`).Scan(&total)
|
|
if err != nil {
|
|
return 0, fmt.Errorf("failed to compute cache usage: %w", err)
|
|
}
|
|
|
|
return total, nil
|
|
}
|
|
|
|
// evictionCandidate is one LRU eviction victim candidate: either a
|
|
// processed variant (isVariant true, identified by cacheKey) or a
|
|
// source content blob (identified by contentHash).
|
|
type evictionCandidate struct {
|
|
isVariant bool
|
|
cacheKey VariantKey
|
|
contentHash ContentHash
|
|
sizeBytes int64
|
|
lastAccessedAt string
|
|
}
|
|
|
|
// EvictToLimit evicts least-recently-used cache entries until total
|
|
// tracked usage is at or below the configured MaxBytes limit. It is a
|
|
// no-op when the cache is disabled or no limit is configured.
|
|
func (c *Cache) EvictToLimit(ctx context.Context) error {
|
|
if c.disabled || c.config.MaxBytes <= 0 {
|
|
return nil
|
|
}
|
|
|
|
for {
|
|
usage, err := c.UsageBytes(ctx)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
|
|
if usage <= c.config.MaxBytes {
|
|
return nil
|
|
}
|
|
|
|
freed, err := c.evictBatch(ctx, usage-c.config.MaxBytes)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
|
|
if freed == 0 {
|
|
c.log.Warn("cache eviction made no progress",
|
|
"usage_bytes", usage,
|
|
"cache_max_bytes", c.config.MaxBytes,
|
|
)
|
|
|
|
return nil
|
|
}
|
|
|
|
c.log.Info("evicted cache content",
|
|
"freed_bytes", freed,
|
|
"usage_bytes", usage-freed,
|
|
"cache_max_bytes", c.config.MaxBytes,
|
|
)
|
|
}
|
|
}
|
|
|
|
// evictBatch fetches one batch of LRU candidates across variants and
|
|
// source blobs and evicts them oldest-first until excessBytes are
|
|
// freed or the batch is exhausted. It returns the bytes freed.
|
|
func (c *Cache) evictBatch(ctx context.Context, excessBytes int64) (int64, error) {
|
|
candidates, err := c.evictionCandidates(ctx)
|
|
if err != nil {
|
|
return 0, err
|
|
}
|
|
|
|
var freed int64
|
|
|
|
for _, candidate := range candidates {
|
|
if freed >= excessBytes {
|
|
break
|
|
}
|
|
|
|
err := c.evictCandidate(ctx, candidate)
|
|
if err != nil {
|
|
c.log.Warn("failed to evict cache entry",
|
|
"cache_key", candidate.cacheKey,
|
|
"content_hash", candidate.contentHash,
|
|
"error", err,
|
|
)
|
|
|
|
continue
|
|
}
|
|
|
|
freed += candidate.sizeBytes
|
|
}
|
|
|
|
return freed, nil
|
|
}
|
|
|
|
// evictCandidate removes a single eviction victim.
|
|
func (c *Cache) evictCandidate(ctx context.Context, candidate evictionCandidate) error {
|
|
if candidate.isVariant {
|
|
return c.evictVariant(ctx, candidate.cacheKey)
|
|
}
|
|
|
|
return c.evictSourceBlob(ctx, candidate.contentHash)
|
|
}
|
|
|
|
// evictionCandidates returns up to evictionBatchSize variants and
|
|
// evictionBatchSize source blobs, merged into a single list ordered by
|
|
// last access time (oldest first).
|
|
func (c *Cache) evictionCandidates(ctx context.Context) ([]evictionCandidate, error) {
|
|
variants, err := c.variantCandidates(ctx)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
|
|
sources, err := c.sourceCandidates(ctx)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
|
|
// Merge the two lists, each already sorted oldest-first. SQLite
|
|
// CURRENT_TIMESTAMP strings compare correctly lexicographically.
|
|
merged := make([]evictionCandidate, 0, len(variants)+len(sources))
|
|
|
|
for len(variants) > 0 && len(sources) > 0 {
|
|
if variants[0].lastAccessedAt <= sources[0].lastAccessedAt {
|
|
merged = append(merged, variants[0])
|
|
variants = variants[1:]
|
|
} else {
|
|
merged = append(merged, sources[0])
|
|
sources = sources[1:]
|
|
}
|
|
}
|
|
|
|
merged = append(merged, variants...)
|
|
merged = append(merged, sources...)
|
|
|
|
return merged, nil
|
|
}
|
|
|
|
// variantCandidates returns the least recently used variants.
|
|
func (c *Cache) variantCandidates(ctx context.Context) ([]evictionCandidate, error) {
|
|
rows, err := c.db.QueryContext(ctx, `
|
|
SELECT cache_key, size_bytes, last_accessed_at
|
|
FROM variant_content
|
|
ORDER BY last_accessed_at ASC, cache_key ASC
|
|
LIMIT ?
|
|
`, evictionBatchSize)
|
|
if err != nil {
|
|
return nil, fmt.Errorf("failed to query variant eviction candidates: %w", err)
|
|
}
|
|
|
|
defer func() { _ = rows.Close() }()
|
|
|
|
var candidates []evictionCandidate
|
|
|
|
for rows.Next() {
|
|
candidate := evictionCandidate{isVariant: true}
|
|
|
|
var key string
|
|
|
|
err := rows.Scan(&key, &candidate.sizeBytes, &candidate.lastAccessedAt)
|
|
if err != nil {
|
|
return nil, fmt.Errorf("failed to scan variant candidate: %w", err)
|
|
}
|
|
|
|
candidate.cacheKey = VariantKey(key)
|
|
candidates = append(candidates, candidate)
|
|
}
|
|
|
|
err = rows.Err()
|
|
if err != nil {
|
|
return nil, fmt.Errorf("variant candidate iteration failed: %w", err)
|
|
}
|
|
|
|
return candidates, nil
|
|
}
|
|
|
|
// sourceCandidates returns the least recently used source blobs. Rows
|
|
// written before the LRU column existed fall back to fetched_at.
|
|
func (c *Cache) sourceCandidates(ctx context.Context) ([]evictionCandidate, error) {
|
|
rows, err := c.db.QueryContext(ctx, `
|
|
SELECT content_hash, size_bytes,
|
|
COALESCE(last_accessed_at, fetched_at, '1970-01-01 00:00:00') AS lru
|
|
FROM source_content
|
|
ORDER BY lru ASC, content_hash ASC
|
|
LIMIT ?
|
|
`, evictionBatchSize)
|
|
if err != nil {
|
|
return nil, fmt.Errorf("failed to query source eviction candidates: %w", err)
|
|
}
|
|
|
|
defer func() { _ = rows.Close() }()
|
|
|
|
var candidates []evictionCandidate
|
|
|
|
for rows.Next() {
|
|
var candidate evictionCandidate
|
|
|
|
var hash string
|
|
|
|
err := rows.Scan(&hash, &candidate.sizeBytes, &candidate.lastAccessedAt)
|
|
if err != nil {
|
|
return nil, fmt.Errorf("failed to scan source candidate: %w", err)
|
|
}
|
|
|
|
candidate.contentHash = ContentHash(hash)
|
|
candidates = append(candidates, candidate)
|
|
}
|
|
|
|
err = rows.Err()
|
|
if err != nil {
|
|
return nil, fmt.Errorf("source candidate iteration failed: %w", err)
|
|
}
|
|
|
|
return candidates, nil
|
|
}
|
|
|
|
// evictVariant removes one variant: accounting row first, then the
|
|
// content and .meta files, so the database never references a deleted
|
|
// file. The metaCache entry goes before the files; a GetVariant that
|
|
// read them just before may put it back, and the next GetVariant then
|
|
// fails to open the file and removes it again.
|
|
func (c *Cache) evictVariant(ctx context.Context, cacheKey VariantKey) error {
|
|
_, err := c.db.ExecContext(ctx,
|
|
`DELETE FROM variant_content WHERE cache_key = ?`, string(cacheKey))
|
|
if err != nil {
|
|
return fmt.Errorf("failed to delete variant accounting row: %w", err)
|
|
}
|
|
|
|
c.metaCache.Remove(cacheKey)
|
|
|
|
err = c.variants.DeleteWithMeta(cacheKey)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
|
|
return nil
|
|
}
|
|
|
|
// sourceReference identifies one source_metadata row's JSON sidecar.
|
|
type sourceReference struct {
|
|
host string
|
|
pathHash PathHash
|
|
}
|
|
|
|
// evictSourceBlob removes one source content blob. All source_metadata
|
|
// rows referencing the blob are deleted together with its
|
|
// source_content row in a single transaction BEFORE the file is
|
|
// unlinked: a blob referenced by multiple source paths is only ever
|
|
// removed together with all of its references, and database rows never
|
|
// point at deleted files. The JSON metadata sidecars for the removed
|
|
// rows are deleted afterwards.
|
|
//
|
|
// The whole operation holds the content hash's lock (the same one
|
|
// StoreSource holds for its full store), so a concurrent store of
|
|
// identical content bytes can never observe the file gone but a row
|
|
// still present, or insert a fresh row between this transaction's
|
|
// commit and the file unlink below: it either runs entirely before
|
|
// this eviction starts, or is blocked until this eviction (row
|
|
// deletion and unlink together) has fully completed.
|
|
func (c *Cache) evictSourceBlob(ctx context.Context, contentHash ContentHash) error {
|
|
unlock := c.contentLocks.Lock(string(contentHash))
|
|
defer unlock()
|
|
|
|
references, err := c.sourceReferences(ctx, contentHash)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
|
|
tx, err := c.db.BeginTx(ctx, nil)
|
|
if err != nil {
|
|
return fmt.Errorf("failed to begin eviction transaction: %w", err)
|
|
}
|
|
|
|
defer func() { _ = tx.Rollback() }()
|
|
|
|
_, err = tx.ExecContext(ctx,
|
|
`DELETE FROM source_metadata WHERE content_hash = ?`, string(contentHash))
|
|
if err != nil {
|
|
return fmt.Errorf("failed to delete source metadata rows: %w", err)
|
|
}
|
|
|
|
_, err = tx.ExecContext(ctx,
|
|
`DELETE FROM source_content WHERE content_hash = ?`, string(contentHash))
|
|
if err != nil {
|
|
return fmt.Errorf("failed to delete source content row: %w", err)
|
|
}
|
|
|
|
err = tx.Commit()
|
|
if err != nil {
|
|
return fmt.Errorf("failed to commit eviction transaction: %w", err)
|
|
}
|
|
|
|
if c.evictSourceBlobTestHook != nil {
|
|
c.evictSourceBlobTestHook(contentHash)
|
|
}
|
|
|
|
// Only after the rows are gone may the files be removed.
|
|
for _, reference := range references {
|
|
err := c.srcMetadata.Delete(reference.host, reference.pathHash)
|
|
if err != nil {
|
|
c.log.Warn("failed to delete metadata sidecar",
|
|
"host", reference.host, "path_hash", reference.pathHash, "error", err)
|
|
}
|
|
}
|
|
|
|
err = c.srcContent.Delete(contentHash)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
|
|
return nil
|
|
}
|
|
|
|
// sourceReferences lists the metadata sidecar locations of every
|
|
// source_metadata row referencing the given blob.
|
|
func (c *Cache) sourceReferences(
|
|
ctx context.Context, contentHash ContentHash,
|
|
) ([]sourceReference, error) {
|
|
rows, err := c.db.QueryContext(ctx, `
|
|
SELECT source_host, path_hash FROM source_metadata WHERE content_hash = ?
|
|
`, string(contentHash))
|
|
if err != nil {
|
|
return nil, fmt.Errorf("failed to query source references: %w", err)
|
|
}
|
|
|
|
defer func() { _ = rows.Close() }()
|
|
|
|
var references []sourceReference
|
|
|
|
for rows.Next() {
|
|
var reference sourceReference
|
|
|
|
var pathHash string
|
|
|
|
err := rows.Scan(&reference.host, &pathHash)
|
|
if err != nil {
|
|
return nil, fmt.Errorf("failed to scan source reference: %w", err)
|
|
}
|
|
|
|
reference.pathHash = PathHash(pathHash)
|
|
references = append(references, reference)
|
|
}
|
|
|
|
err = rows.Err()
|
|
if err != nil {
|
|
return nil, fmt.Errorf("source reference iteration failed: %w", err)
|
|
}
|
|
|
|
return references, nil
|
|
}
|
|
|
|
// notifyWritePressure wakes the background evictor after a store, so
|
|
// eviction under write pressure happens promptly without blocking the
|
|
// storing request. The notification channel has capacity one and drops
|
|
// when a wakeup is already pending.
|
|
func (c *Cache) notifyWritePressure() {
|
|
if c.disabled || c.config.MaxBytes <= 0 {
|
|
return
|
|
}
|
|
|
|
select {
|
|
case c.evictionPressure <- struct{}{}:
|
|
default:
|
|
}
|
|
}
|
|
|
|
// StartEviction launches the background eviction goroutine, which
|
|
// reconciles the database accounting with the cache directories at
|
|
// startup and again on every periodic tick thereafter, and evicts to
|
|
// the configured limit on the given periodic interval and on
|
|
// write-pressure notifications. It is a no-op on a disabled cache or
|
|
// when already started.
|
|
func (c *Cache) StartEviction(interval time.Duration) {
|
|
if c.disabled || c.evictionStarted {
|
|
return
|
|
}
|
|
|
|
c.evictionStarted = true
|
|
|
|
go c.evictionLoop(interval)
|
|
}
|
|
|
|
// StopEviction stops the background eviction goroutine and waits for
|
|
// it to exit. It is safe to call when eviction was never started, and
|
|
// safe to call more than once.
|
|
func (c *Cache) StopEviction() {
|
|
if !c.evictionStarted {
|
|
return
|
|
}
|
|
|
|
c.evictionStopOnce.Do(func() {
|
|
close(c.evictionStop)
|
|
<-c.evictionDone
|
|
})
|
|
}
|
|
|
|
// evictionLoop is the body of the background eviction goroutine.
|
|
func (c *Cache) evictionLoop(interval time.Duration) {
|
|
defer close(c.evictionDone)
|
|
|
|
ctx := context.Background()
|
|
|
|
c.runReconciliationPass(ctx)
|
|
c.runEvictionPass(ctx)
|
|
|
|
ticker := time.NewTicker(interval)
|
|
defer ticker.Stop()
|
|
|
|
for {
|
|
select {
|
|
case <-c.evictionStop:
|
|
return
|
|
case <-ticker.C:
|
|
// Reconciliation walks the cache directories, so it only
|
|
// runs on the periodic ticker rather than on every
|
|
// write-pressure wakeup, keeping it off the per-store hot
|
|
// path. Reusing the eviction interval itself (rather than a
|
|
// separate, longer one) is a deliberate choice: it is the
|
|
// simplest option that still bounds how long a store's
|
|
// best-effort accounting insert can stay silently
|
|
// unaccounted for to one interval, on a process that is
|
|
// already running this loop regardless.
|
|
c.runReconciliationPass(ctx)
|
|
case <-c.evictionPressure:
|
|
}
|
|
|
|
c.runEvictionPass(ctx)
|
|
}
|
|
}
|
|
|
|
// runEvictionPass runs one eviction pass, logging failures instead of
|
|
// propagating them (the loop must keep running).
|
|
func (c *Cache) runEvictionPass(ctx context.Context) {
|
|
err := c.EvictToLimit(ctx)
|
|
if err != nil {
|
|
c.log.Warn("cache eviction pass failed", "error", err)
|
|
}
|
|
}
|
|
|
|
// runReconciliationPass runs one reconciliation pass, logging failures
|
|
// instead of propagating them (the loop must keep running).
|
|
func (c *Cache) runReconciliationPass(ctx context.Context) {
|
|
err := c.reconcileAccounting(ctx)
|
|
if err != nil {
|
|
c.log.Warn("cache accounting reconciliation failed", "error", err)
|
|
}
|
|
}
|
|
|
|
// reconcileAccounting synchronizes the database size accounting with
|
|
// the actual contents of the cache directories. It runs at startup and
|
|
// again on every periodic eviction tick thereafter, off the request
|
|
// hot path: it adopts variant files that predate the accounting table
|
|
// (or whose accounting insert failed, e.g. StoreVariant's best-effort
|
|
// insert under transient DB contention), drops accounting rows whose
|
|
// files are missing, removes source blob files the database does not
|
|
// know (and rows whose files are gone), 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.
|
|
func (c *Cache) reconcileAccounting(ctx context.Context) error {
|
|
if c.disabled {
|
|
return nil
|
|
}
|
|
|
|
err := c.reconcileVariantFiles(ctx)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
|
|
err = c.reconcileVariantRows(ctx)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
|
|
err = c.reconcileSourceFiles(ctx)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
|
|
err = c.reconcileSourceRows(ctx)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
|
|
return nil
|
|
}
|
|
|
|
// reconcileVariantFiles walks the variant storage directory, adopting
|
|
// files without accounting rows and sweeping stale temp files.
|
|
func (c *Cache) reconcileVariantFiles(ctx context.Context) error {
|
|
return filepath.WalkDir(
|
|
c.variants.baseDir,
|
|
func(path string, entry fs.DirEntry, err error) error {
|
|
if err != nil || entry.IsDir() {
|
|
return err
|
|
}
|
|
|
|
name := entry.Name()
|
|
|
|
if strings.HasPrefix(name, tempFilePrefix) {
|
|
c.sweepStaleTempFile(path, entry)
|
|
|
|
return nil
|
|
}
|
|
|
|
if strings.HasSuffix(name, variantMetaSuffix) {
|
|
return nil
|
|
}
|
|
|
|
return c.adoptVariantFile(ctx, path, entry, VariantKey(name))
|
|
},
|
|
)
|
|
}
|
|
|
|
// adoptVariantFile inserts an accounting row for a variant file that
|
|
// has none, using the file's size and modification time.
|
|
func (c *Cache) adoptVariantFile(
|
|
ctx context.Context, path string, entry fs.DirEntry, cacheKey VariantKey,
|
|
) error {
|
|
var rowExists int
|
|
|
|
err := c.db.QueryRowContext(ctx,
|
|
`SELECT COUNT(*) FROM variant_content WHERE cache_key = ?`, string(cacheKey),
|
|
).Scan(&rowExists)
|
|
if err != nil {
|
|
return fmt.Errorf("failed to check variant accounting row: %w", err)
|
|
}
|
|
|
|
if rowExists > 0 {
|
|
return nil
|
|
}
|
|
|
|
info, err := entry.Info()
|
|
if err != nil {
|
|
return fmt.Errorf("failed to stat variant file: %w", err)
|
|
}
|
|
|
|
modTime := info.ModTime().UTC().Format(sqliteTimestampLayout)
|
|
contentType := c.variantContentTypeFromSidecar(path)
|
|
|
|
_, err = c.db.ExecContext(ctx, `
|
|
INSERT INTO variant_content
|
|
(cache_key, size_bytes, content_type, created_at, last_accessed_at)
|
|
VALUES (?, ?, ?, ?, ?)
|
|
`, string(cacheKey), info.Size(), contentType, modTime, modTime)
|
|
if err != nil {
|
|
return fmt.Errorf("failed to adopt variant file into accounting: %w", err)
|
|
}
|
|
|
|
c.log.Info("adopted untracked variant file into size accounting",
|
|
"cache_key", cacheKey, "size_bytes", info.Size())
|
|
|
|
return nil
|
|
}
|
|
|
|
// variantContentTypeFromSidecar reads the content type from a variant
|
|
// .meta sidecar, falling back to application/octet-stream.
|
|
func (c *Cache) variantContentTypeFromSidecar(variantPath string) string {
|
|
//nolint:gosec // path from cache walk
|
|
metaData, err := os.ReadFile(variantPath + variantMetaSuffix)
|
|
if err != nil {
|
|
return fallbackContentType
|
|
}
|
|
|
|
var meta VariantMeta
|
|
if json.Unmarshal(metaData, &meta) != nil || meta.ContentType == "" {
|
|
return fallbackContentType
|
|
}
|
|
|
|
return meta.ContentType
|
|
}
|
|
|
|
// reconcileVariantRows drops accounting rows whose variant files are
|
|
// missing, so the database never references deleted content.
|
|
func (c *Cache) reconcileVariantRows(ctx context.Context) error {
|
|
keys, err := c.allVariantKeys(ctx)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
|
|
for _, key := range keys {
|
|
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)
|
|
}
|
|
|
|
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")
|
|
}
|
|
|
|
// 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.
|
|
func queryStringColumn[T ~string](
|
|
ctx context.Context, db *sql.DB, query, plural, singular string,
|
|
) ([]T, error) {
|
|
rows, err := db.QueryContext(ctx, query)
|
|
if err != nil {
|
|
return nil, fmt.Errorf("failed to query %s: %w", plural, err)
|
|
}
|
|
|
|
defer func() { _ = rows.Close() }()
|
|
|
|
var values []T
|
|
|
|
for rows.Next() {
|
|
var value string
|
|
|
|
err := rows.Scan(&value)
|
|
if err != nil {
|
|
return nil, fmt.Errorf("failed to scan %s: %w", singular, err)
|
|
}
|
|
|
|
values = append(values, T(value))
|
|
}
|
|
|
|
err = rows.Err()
|
|
if err != nil {
|
|
return nil, fmt.Errorf("%s iteration failed: %w", singular, err)
|
|
}
|
|
|
|
return values, nil
|
|
}
|
|
|
|
// reconcileSourceFiles walks the source content directory, removing
|
|
// blob files the database does not track (they are unreachable: source
|
|
// lookups always go through source_metadata) and sweeping stale temp
|
|
// files.
|
|
func (c *Cache) reconcileSourceFiles(ctx context.Context) error {
|
|
return filepath.WalkDir(
|
|
c.srcContent.baseDir,
|
|
func(path string, entry fs.DirEntry, err error) error {
|
|
if err != nil || entry.IsDir() {
|
|
return err
|
|
}
|
|
|
|
name := entry.Name()
|
|
|
|
if strings.HasPrefix(name, tempFilePrefix) {
|
|
c.sweepStaleTempFile(path, entry)
|
|
|
|
return nil
|
|
}
|
|
|
|
return c.removeUntrackedSourceFile(ctx, path, ContentHash(name))
|
|
},
|
|
)
|
|
}
|
|
|
|
// removeUntrackedSourceFile deletes a source blob file that has no
|
|
// source_content row. Any source_metadata rows referencing the hash
|
|
// are removed first so no row ever points at a deleted file.
|
|
func (c *Cache) removeUntrackedSourceFile(
|
|
ctx context.Context, path string, contentHash ContentHash,
|
|
) error {
|
|
var rowExists int
|
|
|
|
err := c.db.QueryRowContext(ctx,
|
|
`SELECT COUNT(*) FROM source_content WHERE content_hash = ?`, string(contentHash),
|
|
).Scan(&rowExists)
|
|
if err != nil {
|
|
return fmt.Errorf("failed to check source content row: %w", err)
|
|
}
|
|
|
|
if rowExists > 0 {
|
|
return nil
|
|
}
|
|
|
|
_, err = c.db.ExecContext(ctx,
|
|
`DELETE FROM source_metadata WHERE content_hash = ?`, string(contentHash))
|
|
if err != nil {
|
|
return fmt.Errorf("failed to delete metadata rows for untracked blob: %w", err)
|
|
}
|
|
|
|
err = os.Remove(path)
|
|
if err != nil && !os.IsNotExist(err) {
|
|
return fmt.Errorf("failed to remove untracked source file: %w", err)
|
|
}
|
|
|
|
c.log.Info("removed untracked source content file", "content_hash", contentHash)
|
|
|
|
return nil
|
|
}
|
|
|
|
// reconcileSourceRows removes source_content rows (and their metadata
|
|
// references and sidecars) whose blob files are missing on disk.
|
|
func (c *Cache) reconcileSourceRows(ctx context.Context) error {
|
|
hashes, err := c.allSourceContentHashes(ctx)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
|
|
for _, hash := range hashes {
|
|
if c.srcContent.Exists(hash) {
|
|
continue
|
|
}
|
|
|
|
// The blob file is already gone; evictSourceBlob removes the
|
|
// rows and sidecars and tolerates the missing file.
|
|
err := c.evictSourceBlob(ctx, hash)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
|
|
c.log.Info("dropped rows for missing source content file", "content_hash", hash)
|
|
}
|
|
|
|
return nil
|
|
}
|
|
|
|
// allSourceContentHashes returns every tracked source content hash.
|
|
func (c *Cache) allSourceContentHashes(ctx context.Context) ([]ContentHash, error) {
|
|
return queryStringColumn[ContentHash](ctx, c.db,
|
|
`SELECT content_hash FROM source_content`,
|
|
"source content hashes", "content hash")
|
|
}
|
|
|
|
// sweepStaleTempFile removes a temp file left behind by a crashed
|
|
// write once it is old enough that no in-flight store can own it.
|
|
func (c *Cache) sweepStaleTempFile(path string, entry fs.DirEntry) {
|
|
info, err := entry.Info()
|
|
if err != nil {
|
|
return
|
|
}
|
|
|
|
if time.Since(info.ModTime()) < staleTempFileAge {
|
|
return
|
|
}
|
|
|
|
err = os.Remove(path)
|
|
if err != nil && !os.IsNotExist(err) {
|
|
c.log.Warn("failed to remove stale temp file", "path", path, "error", err)
|
|
|
|
return
|
|
}
|
|
|
|
c.log.Info("removed stale temp file", "path", path)
|
|
}
|