fix: close TOCTOU window between blob eviction commit and unlink
StoreSource now hashes content itself and holds the per-hash contentLock across the whole store (file write plus accounting row inserts); evictSourceBlob holds the same lock across its whole operation (row deletion transaction through file unlink). A concurrent store and eviction of identical content bytes can no longer interleave: either runs to completion before the other starts, so a fresh row can never be left pointing at a file the other side is mid-unlink on. ContentStorage gains StoreHashed for callers that need the hash before writing; Store is refactored to share the write-if-absent logic with it, with no change to its existing behavior or signature. internal/imgcache/eviction_test.go: TestEvictSourceBlobExcludesConcurrentStoreOfIdenticalContent proves it: pauses eviction (via evictSourceBlobTestHook) in the exact window between commit and unlink, asserts a concurrent StoreSource for identical content blocks rather than completing, then verifies no dangling reference and that the store's data survives once eviction releases the hash.
This commit is contained in:
@@ -2,7 +2,9 @@ package imgcache
|
|||||||
|
|
||||||
import (
|
import (
|
||||||
"context"
|
"context"
|
||||||
|
"crypto/sha256"
|
||||||
"database/sql"
|
"database/sql"
|
||||||
|
"encoding/hex"
|
||||||
"encoding/json"
|
"encoding/json"
|
||||||
"errors"
|
"errors"
|
||||||
"fmt"
|
"fmt"
|
||||||
@@ -215,8 +217,33 @@ func (c *Cache) StoreSource(
|
|||||||
return "", nil
|
return "", nil
|
||||||
}
|
}
|
||||||
|
|
||||||
// Store content
|
// Hash the content ourselves (rather than via srcContent.Store,
|
||||||
contentHash, size, err := c.srcContent.Store(content)
|
// which would hash internally) so the content hash is known before
|
||||||
|
// any file or database work happens: that lets the entire store be
|
||||||
|
// serialized, per hash, against a concurrent eviction of the same
|
||||||
|
// content below.
|
||||||
|
data, err := io.ReadAll(content)
|
||||||
|
if err != nil {
|
||||||
|
return "", fmt.Errorf("failed to read source content: %w", err)
|
||||||
|
}
|
||||||
|
|
||||||
|
sum := sha256.Sum256(data)
|
||||||
|
contentHash := ContentHash(hex.EncodeToString(sum[:]))
|
||||||
|
|
||||||
|
// Hold the content hash's lock for the whole store operation. A
|
||||||
|
// concurrent eviction of this exact hash (the real SHA-256 dedup
|
||||||
|
// case: a different source path whose bytes hash identically)
|
||||||
|
// deletes the accounting rows and unlinks the file inside the same
|
||||||
|
// lock, so the two can never interleave: either this store
|
||||||
|
// completes first (and a subsequent eviction removes it together
|
||||||
|
// with its rows and file, correctly), or eviction completes first
|
||||||
|
// (and this store finds the file already gone and recreates it
|
||||||
|
// fresh) — never a fresh row left pointing at a file eviction is
|
||||||
|
// mid-unlink on.
|
||||||
|
unlock := c.contentLocks.Lock(string(contentHash))
|
||||||
|
defer unlock()
|
||||||
|
|
||||||
|
size, err := c.srcContent.StoreHashed(contentHash, data)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return "", fmt.Errorf("failed to store source content: %w", err)
|
return "", fmt.Errorf("failed to store source content: %w", err)
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -291,7 +291,18 @@ type sourceReference struct {
|
|||||||
// removed together with all of its references, and database rows never
|
// removed together with all of its references, and database rows never
|
||||||
// point at deleted files. The JSON metadata sidecars for the removed
|
// point at deleted files. The JSON metadata sidecars for the removed
|
||||||
// rows are deleted afterwards.
|
// 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 {
|
func (c *Cache) evictSourceBlob(ctx context.Context, contentHash ContentHash) error {
|
||||||
|
unlock := c.contentLocks.Lock(string(contentHash))
|
||||||
|
defer unlock()
|
||||||
|
|
||||||
references, err := c.sourceReferences(ctx, contentHash)
|
references, err := c.sourceReferences(ctx, contentHash)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return err
|
return err
|
||||||
|
|||||||
@@ -65,24 +65,49 @@ func (s *ContentStorage) Store(r io.Reader) (hash ContentHash, size int64, err e
|
|||||||
hash = ContentHash(hex.EncodeToString(h[:]))
|
hash = ContentHash(hex.EncodeToString(h[:]))
|
||||||
size = int64(len(data))
|
size = int64(len(data))
|
||||||
|
|
||||||
|
if err := s.writeIfAbsent(hash, data); err != nil {
|
||||||
|
return "", 0, err
|
||||||
|
}
|
||||||
|
|
||||||
|
return hash, size, nil
|
||||||
|
}
|
||||||
|
|
||||||
|
// StoreHashed writes pre-hashed content to storage at the path derived
|
||||||
|
// from hash, without recomputing it. Callers that already know the
|
||||||
|
// hash before writing (e.g. because they must hold a hash-keyed lock
|
||||||
|
// across the whole store operation) use this instead of Store. Like
|
||||||
|
// Store, it is idempotent: content already on disk at that path is
|
||||||
|
// left untouched.
|
||||||
|
func (s *ContentStorage) StoreHashed(hash ContentHash, data []byte) (size int64, err error) {
|
||||||
|
if err := s.writeIfAbsent(hash, data); err != nil {
|
||||||
|
return 0, err
|
||||||
|
}
|
||||||
|
|
||||||
|
return int64(len(data)), nil
|
||||||
|
}
|
||||||
|
|
||||||
|
// writeIfAbsent writes data to the path derived from hash, unless
|
||||||
|
// content already exists there, via a temp-file-plus-rename so
|
||||||
|
// concurrent readers never observe a partial file.
|
||||||
|
func (s *ContentStorage) writeIfAbsent(hash ContentHash, data []byte) (err error) {
|
||||||
// Build path: <basedir>/<ab>/<cd>/<hash>
|
// Build path: <basedir>/<ab>/<cd>/<hash>
|
||||||
path := s.hashToPath(hash)
|
path := s.hashToPath(hash)
|
||||||
|
|
||||||
// Check if already exists
|
// Check if already exists
|
||||||
if _, err := os.Stat(path); err == nil {
|
if _, statErr := os.Stat(path); statErr == nil {
|
||||||
return hash, size, nil
|
return nil
|
||||||
}
|
}
|
||||||
|
|
||||||
// Create directory structure
|
// Create directory structure
|
||||||
dir := filepath.Dir(path)
|
dir := filepath.Dir(path)
|
||||||
if err := os.MkdirAll(dir, StorageDirPerm); err != nil {
|
if err := os.MkdirAll(dir, StorageDirPerm); err != nil {
|
||||||
return "", 0, fmt.Errorf("failed to create directory: %w", err)
|
return fmt.Errorf("failed to create directory: %w", err)
|
||||||
}
|
}
|
||||||
|
|
||||||
// Write to temp file first, then rename for atomicity
|
// Write to temp file first, then rename for atomicity
|
||||||
tmpFile, err := os.CreateTemp(dir, ".tmp-*")
|
tmpFile, err := os.CreateTemp(dir, ".tmp-*")
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return "", 0, fmt.Errorf("failed to create temp file: %w", err)
|
return fmt.Errorf("failed to create temp file: %w", err)
|
||||||
}
|
}
|
||||||
tmpPath := tmpFile.Name()
|
tmpPath := tmpFile.Name()
|
||||||
|
|
||||||
@@ -95,20 +120,20 @@ func (s *ContentStorage) Store(r io.Reader) (hash ContentHash, size int64, err e
|
|||||||
if _, err := tmpFile.Write(data); err != nil {
|
if _, err := tmpFile.Write(data); err != nil {
|
||||||
_ = tmpFile.Close()
|
_ = tmpFile.Close()
|
||||||
|
|
||||||
return "", 0, fmt.Errorf("failed to write content: %w", err)
|
return fmt.Errorf("failed to write content: %w", err)
|
||||||
}
|
}
|
||||||
|
|
||||||
if err := tmpFile.Close(); err != nil {
|
if err := tmpFile.Close(); err != nil {
|
||||||
return "", 0, fmt.Errorf("failed to close temp file: %w", err)
|
return fmt.Errorf("failed to close temp file: %w", err)
|
||||||
}
|
}
|
||||||
|
|
||||||
// Atomic rename
|
// Atomic rename
|
||||||
//nolint:gosec // G703: paths from internal SHA256 hashes
|
//nolint:gosec // G703: paths from internal SHA256 hashes
|
||||||
if err := os.Rename(filepath.Clean(tmpPath), filepath.Clean(path)); err != nil {
|
if err := os.Rename(filepath.Clean(tmpPath), filepath.Clean(path)); err != nil {
|
||||||
return "", 0, fmt.Errorf("failed to rename temp file: %w", err)
|
return fmt.Errorf("failed to rename temp file: %w", err)
|
||||||
}
|
}
|
||||||
|
|
||||||
return hash, size, nil
|
return nil
|
||||||
}
|
}
|
||||||
|
|
||||||
// Load returns a reader for the content with the given hash.
|
// Load returns a reader for the content with the given hash.
|
||||||
|
|||||||
Reference in New Issue
Block a user