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 was merged in pull request #228.
This commit is contained in:
+246
-72
@@ -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
|
||||
|
||||
Reference in New Issue
Block a user