feat: reclaim terminal commands atomically

This commit is contained in:
2026-08-31 09:56:11 +00:00
parent 611100f67c
commit af9041e1ca
4 changed files with 582 additions and 2 deletions
+413
View File
@@ -0,0 +1,413 @@
package store
import (
"context"
"crypto/subtle"
"database/sql"
"errors"
"os"
"path/filepath"
"strings"
"time"
)
const (
FaultAfterEvictionMarked = "after_eviction_marked"
FaultAfterEvictionFilesMoved = "after_eviction_files_moved"
FaultAfterEvictionMetadataGone = "after_eviction_metadata_deleted"
)
type RetentionPolicy struct {
TerminalAge time.Duration
PressureBytes uint64
TombstoneMaxEntries uint64
}
type RetentionReport struct {
CommandsEvicted uint64
BytesReleased uint64
FilesDeleted uint64
TombstonesTrimmed uint64
}
type evictionCommand struct {
issueUUID [16]byte
clientID string
issueTime int64
terminalTime int64
lifecycle uint32
chargedBytes uint64
requestHash [32]byte
segmentNames []string
}
func (store *Store) RunRetention(ctx context.Context, now time.Time, policy RetentionPolicy) (RetentionReport, error) {
var report RetentionReport
if now.IsZero() || policy.TerminalAge < 0 || policy.TombstoneMaxEntries == 0 {
return report, errors.New("invalid retention policy")
}
store.writeMu.Lock()
defer store.writeMu.Unlock()
database, err := store.openDatabase()
if err != nil {
return report, err
}
if err := store.recoverEvictionsLocked(ctx, database, policy.TombstoneMaxEntries, &report); err != nil {
return report, err
}
rows, err := database.QueryContext(ctx, `SELECT issue_uuid, issue_time, terminal_time, charged_bytes, lifecycle
FROM commands WHERE lifecycle BETWEEN 5 AND 11 AND retention_status = 1 ORDER BY issue_time, issue_uuid`)
if err != nil {
return report, err
}
var candidates []RetentionCandidate
for rows.Next() {
var issueBytes []byte
var issueTime, terminalTime int64
var charged uint64
var lifecycle uint32
if err := rows.Scan(&issueBytes, &issueTime, &terminalTime, &charged, &lifecycle); err != nil {
_ = rows.Close()
return report, err
}
if len(issueBytes) != 16 {
_ = rows.Close()
return report, ErrInvalidSegmentRecord
}
var issue [16]byte
copy(issue[:], issueBytes)
candidates = append(candidates, RetentionCandidate{
IssueUUID: issue, IssueTime: time.Unix(0, issueTime), TerminalTime: time.Unix(0, terminalTime),
ChargedBytes: charged, Terminal: lifecycle >= 5 && lifecycle <= 11,
})
}
if err := rows.Close(); err != nil {
return report, err
}
selection := SelectWholeCommandEvictions(now, policy.TerminalAge, policy.PressureBytes, candidates)
for _, issue := range selection.IssueUUIDs {
evicted, err := store.loadEvictionCommand(ctx, database, issue)
if err != nil {
return report, err
}
files, trimmed, err := store.evictCommandLocked(ctx, database, now, policy.TombstoneMaxEntries, evicted)
if err != nil {
return report, err
}
report.CommandsEvicted++
report.BytesReleased += evicted.chargedBytes
report.FilesDeleted += files
report.TombstonesTrimmed += trimmed
}
return report, nil
}
func (store *Store) RecoverEvictions(ctx context.Context, tombstoneMaxEntries uint64) (RetentionReport, error) {
var report RetentionReport
if tombstoneMaxEntries == 0 {
return report, errors.New("tombstone maximum must be positive")
}
store.writeMu.Lock()
defer store.writeMu.Unlock()
database, err := store.openDatabase()
if err != nil {
return report, err
}
err = store.recoverEvictionsLocked(ctx, database, tombstoneMaxEntries, &report)
return report, err
}
func (store *Store) recoverEvictionsLocked(ctx context.Context, database *sql.DB, tombstoneMaxEntries uint64, report *RetentionReport) error {
rows, err := database.QueryContext(ctx, `SELECT issue_uuid FROM commands WHERE retention_status = 2 ORDER BY issue_time, issue_uuid`)
if err != nil {
return err
}
var issues [][16]byte
for rows.Next() {
var encoded []byte
if err := rows.Scan(&encoded); err != nil {
_ = rows.Close()
return err
}
if len(encoded) != 16 {
_ = rows.Close()
return ErrInvalidSegmentRecord
}
var issue [16]byte
copy(issue[:], encoded)
issues = append(issues, issue)
}
if err := rows.Close(); err != nil {
return err
}
for _, issue := range issues {
command, err := store.loadEvictionCommand(ctx, database, issue)
if err != nil {
return err
}
files, trimmed, err := store.evictCommandLocked(ctx, database, time.Now().UTC(), tombstoneMaxEntries, command)
if err != nil {
return err
}
report.CommandsEvicted++
report.BytesReleased += command.chargedBytes
report.FilesDeleted += files
report.TombstonesTrimmed += trimmed
}
files, err := store.sweepDeletionDirectory()
if err != nil {
return err
}
report.FilesDeleted += files
return nil
}
func (store *Store) loadEvictionCommand(ctx context.Context, database *sql.DB, issue [16]byte) (evictionCommand, error) {
command := evictionCommand{issueUUID: issue}
var issueBytes, requestHash []byte
err := database.QueryRowContext(ctx, `SELECT issue_uuid, client_id, issue_time, terminal_time, lifecycle, charged_bytes,
immutable_request_sha256 FROM commands WHERE issue_uuid = ? AND lifecycle BETWEEN 5 AND 11`, issue[:]).Scan(
&issueBytes, &command.clientID, &command.issueTime, &command.terminalTime, &command.lifecycle, &command.chargedBytes, &requestHash)
if errors.Is(err, sql.ErrNoRows) {
return command, ErrCommandNotFound
}
if err != nil {
return command, err
}
if len(issueBytes) != 16 || len(requestHash) != 32 {
return command, ErrInvalidSegmentRecord
}
copy(command.requestHash[:], requestHash)
rows, err := database.QueryContext(ctx, `SELECT path FROM output_segments WHERE issue_uuid = ? ORDER BY ordinal`, issue[:])
if err != nil {
return command, err
}
for rows.Next() {
var name string
if err := rows.Scan(&name); err != nil {
_ = rows.Close()
return command, err
}
owner, _, err := parseSegmentName(name)
if err != nil || owner != issue {
_ = rows.Close()
return command, ErrUnsafeSegmentReference
}
command.segmentNames = append(command.segmentNames, name)
}
if err := rows.Close(); err != nil {
return command, err
}
return command, nil
}
func (store *Store) evictCommandLocked(ctx context.Context, database *sql.DB, now time.Time, tombstoneMax uint64, command evictionCommand) (uint64, uint64, error) {
tx, err := database.BeginTx(ctx, nil)
if err != nil {
return 0, 0, err
}
_, err = tx.ExecContext(ctx, `INSERT INTO command_tombstones (
issue_uuid, immutable_sha256, client_id, terminal_lifecycle, terminal_time, acknowledged_at
) VALUES (?, ?, ?, ?, ?, ?) ON CONFLICT(issue_uuid) DO NOTHING`, command.issueUUID[:], command.requestHash[:], command.clientID,
command.lifecycle, command.terminalTime, now.UnixNano())
if err == nil {
var storedHash []byte
var storedClient string
var storedLifecycle uint32
var storedTerminal int64
err = tx.QueryRowContext(ctx, `SELECT immutable_sha256, client_id, terminal_lifecycle, terminal_time
FROM command_tombstones WHERE issue_uuid = ?`, command.issueUUID[:]).Scan(&storedHash, &storedClient, &storedLifecycle, &storedTerminal)
if err == nil && (len(storedHash) != 32 || subtle.ConstantTimeCompare(storedHash, command.requestHash[:]) != 1 || storedClient != command.clientID || storedLifecycle != command.lifecycle || storedTerminal != command.terminalTime) {
err = ErrEventConflict
}
}
if err == nil {
var update sql.Result
update, err = tx.ExecContext(ctx, `UPDATE commands SET retention_status = 2 WHERE issue_uuid = ? AND lifecycle BETWEEN 5 AND 11`, command.issueUUID[:])
if err == nil {
var affected int64
affected, err = update.RowsAffected()
if err == nil && affected != 1 {
err = ErrCommandNotFound
}
}
}
if err != nil {
_ = tx.Rollback()
return 0, 0, err
}
if err := tx.Commit(); err != nil {
return 0, 0, err
}
if err := store.checkpoint(FaultAfterEvictionMarked); err != nil {
return 0, 0, err
}
store.invalidateActiveSegment(command.issueUUID)
if err := store.moveCommandFiles(command); err != nil {
return 0, 0, err
}
if err := store.checkpoint(FaultAfterEvictionFilesMoved); err != nil {
return 0, 0, err
}
tx, err = database.BeginTx(ctx, nil)
if err != nil {
return 0, 0, err
}
var deleted sql.Result
deleted, err = tx.ExecContext(ctx, `DELETE FROM commands WHERE issue_uuid = ? AND retention_status = 2 AND lifecycle BETWEEN 5 AND 11`, command.issueUUID[:])
if err == nil {
var affected int64
affected, err = deleted.RowsAffected()
if err == nil && affected != 1 {
err = ErrCommandNotFound
}
}
if err == nil {
var update sql.Result
update, err = tx.ExecContext(ctx, `UPDATE clients SET charged_bytes = charged_bytes - ? WHERE client_id = ? AND charged_bytes >= ?`, command.chargedBytes, command.clientID, command.chargedBytes)
if err == nil {
var affected int64
affected, err = update.RowsAffected()
if err == nil && affected != 1 {
err = ErrQuotaCounterMismatch
}
}
}
if err == nil {
var update sql.Result
update, err = tx.ExecContext(ctx, `UPDATE storage_counters SET command_charged_bytes = command_charged_bytes - ?
WHERE singleton = 1 AND command_charged_bytes >= ?`, command.chargedBytes, command.chargedBytes)
if err == nil {
var affected int64
affected, err = update.RowsAffected()
if err == nil && affected != 1 {
err = ErrQuotaCounterMismatch
}
}
}
var tombstoneCount uint64
if err == nil {
err = tx.QueryRowContext(ctx, `SELECT count(*) FROM command_tombstones`).Scan(&tombstoneCount)
}
trimmed := FIFOEntriesToRemove(tombstoneCount, tombstoneMax)
if err == nil && trimmed > 0 {
_, err = tx.ExecContext(ctx, `DELETE FROM command_tombstones WHERE issue_uuid IN (
SELECT issue_uuid FROM command_tombstones ORDER BY acknowledged_at, issue_uuid LIMIT ?)`, trimmed)
}
if err != nil {
_ = tx.Rollback()
return 0, 0, err
}
if err := tx.Commit(); err != nil {
return 0, 0, err
}
if err := store.checkpoint(FaultAfterEvictionMetadataGone); err != nil {
return 0, 0, err
}
files, err := store.deleteMovedFiles(command.segmentNames)
return files, trimmed, err
}
func (store *Store) moveCommandFiles(command evictionCommand) error {
changed := false
for _, name := range command.segmentNames {
source := filepath.Join(store.segmentDirectory(), name)
destination := filepath.Join(store.deletionDirectory(), name+".delete")
sourceInfo, sourceErr := os.Lstat(source)
destinationInfo, destinationErr := os.Lstat(destination)
sourceExists := sourceErr == nil
destinationExists := destinationErr == nil
if sourceErr != nil && !errors.Is(sourceErr, os.ErrNotExist) {
return sourceErr
}
if destinationErr != nil && !errors.Is(destinationErr, os.ErrNotExist) {
return destinationErr
}
if sourceExists && (!sourceInfo.Mode().IsRegular() || sourceInfo.Mode().Perm()&0o077 != 0) {
return ErrUnsafeSegmentReference
}
if destinationExists && (!destinationInfo.Mode().IsRegular() || destinationInfo.Mode().Perm()&0o077 != 0) {
return ErrUnsafeSegmentReference
}
switch {
case sourceExists && !destinationExists:
if err := os.Rename(source, destination); err != nil {
return err
}
changed = true
case !sourceExists && destinationExists:
// A prior attempt already completed this exact move.
case sourceExists && destinationExists:
return errors.New("both live and deletion segment paths exist")
default:
return ErrCommittedRangeMissing
}
}
if changed {
if err := syncDirectory(store.segmentDirectory()); err != nil {
return err
}
if err := syncDirectory(store.deletionDirectory()); err != nil {
return err
}
}
return nil
}
func (store *Store) deleteMovedFiles(names []string) (uint64, error) {
var deleted uint64
for _, name := range names {
path := filepath.Join(store.deletionDirectory(), name+".delete")
if err := os.Remove(path); err != nil {
if errors.Is(err, os.ErrNotExist) {
continue
}
return deleted, err
}
deleted++
}
if deleted > 0 {
if err := syncDirectory(store.deletionDirectory()); err != nil {
return deleted, err
}
}
return deleted, nil
}
func (store *Store) sweepDeletionDirectory() (uint64, error) {
entries, err := os.ReadDir(store.deletionDirectory())
if err != nil {
return 0, err
}
var deleted uint64
for _, entry := range entries {
if !strings.HasSuffix(entry.Name(), ".delete") {
return deleted, ErrUnsafeSegmentReference
}
name := strings.TrimSuffix(entry.Name(), ".delete")
if _, _, err := parseSegmentName(name); err != nil {
return deleted, err
}
path := filepath.Join(store.deletionDirectory(), entry.Name())
info, err := os.Lstat(path)
if err != nil {
return deleted, err
}
if !info.Mode().IsRegular() || info.Mode().Perm()&0o077 != 0 {
return deleted, ErrUnsafeSegmentReference
}
if err := os.Remove(path); err != nil {
return deleted, err
}
deleted++
}
if deleted > 0 {
if err := syncDirectory(store.deletionDirectory()); err != nil {
return deleted, err
}
}
return deleted, nil
}
func (store *Store) deletionDirectory() string { return filepath.Join(store.dataDir, "deleting") }
+1 -1
View File
@@ -68,7 +68,7 @@ func Open(ctx context.Context, options Options) (*Store, error) {
if err := ensurePrivateDirectory(options.DataDir); err != nil {
return nil, err
}
for _, child := range []string{"segments", "audit"} {
for _, child := range []string{"segments", "audit", "deleting"} {
if err := ensurePrivateDirectory(filepath.Join(options.DataDir, child)); err != nil {
return nil, err
}
+18
View File
@@ -303,3 +303,21 @@ id = "BH-STORE-08"
layer = "integration"
status = "implemented"
tests = ["test/integration/store/store_integration_test.go:TestFilesystemFloorAndCounterMismatch_BH_STORE_08"]
[[requirements]]
id = "HP-STORE-08"
layer = "integration"
status = "implemented"
tests = ["test/integration/store/store_integration_test.go:TestTerminalAgeRetentionEvictsWholeCommand_HP_STORE_08"]
[[requirements]]
id = "HP-STORE-09"
layer = "integration"
status = "implemented"
tests = ["test/integration/store/store_integration_test.go:TestTombstoneFIFOIsCappedInEvictionTransaction_HP_STORE_09"]
[[requirements]]
id = "CRASH-STORE-04"
layer = "integration"
status = "implemented"
tests = ["test/integration/store/store_integration_test.go:TestRetentionCrashStagesRollForward_CRASH_STORE_04"]
@@ -3,6 +3,7 @@
package store_test
import (
"bytes"
"context"
"crypto/sha256"
"database/sql"
@@ -27,7 +28,7 @@ func TestRealSQLiteInitializationAndRestart_HP_STORE_01(t *testing.T) {
dataDir := filepath.Join(t.TempDir(), "state")
opened := openStore(t, dataDir)
for _, directory := range []string{dataDir, filepath.Join(dataDir, "segments"), filepath.Join(dataDir, "audit")} {
for _, directory := range []string{dataDir, filepath.Join(dataDir, "segments"), filepath.Join(dataDir, "audit"), filepath.Join(dataDir, "deleting")} {
info, err := os.Stat(directory)
if err != nil {
t.Fatal(err)
@@ -637,6 +638,147 @@ func TestFilesystemFloorAndCounterMismatch_BH_STORE_08(t *testing.T) {
}
}
func TestTerminalAgeRetentionEvictsWholeCommand_HP_STORE_08(t *testing.T) {
t.Parallel()
dataDir := filepath.Join(t.TempDir(), "state")
opened := openStore(t, dataDir)
now := time.Unix(2_000_000, 0).UTC()
old := uuidBytes(120)
active := uuidBytes(140)
seedCommand(t, opened.DB(), old)
seedCommand(t, opened.DB(), active)
if _, err := opened.AppendCommandEvent(context.Background(), appendEvent(old, 1, "retained history")); err != nil {
t.Fatal(err)
}
markTerminal(t, opened.DB(), old, now.Add(-40*24*time.Hour), now.Add(-31*24*time.Hour))
if _, err := opened.DB().Exec(`UPDATE commands SET issue_time = ? WHERE issue_uuid = ?`, now.Add(-100*24*time.Hour).UnixNano(), active[:]); err != nil {
t.Fatal(err)
}
report, err := opened.RunRetention(context.Background(), now, store.RetentionPolicy{
TerminalAge: 30 * 24 * time.Hour, TombstoneMaxEntries: 1_000_000,
})
if err != nil {
t.Fatal(err)
}
if report.CommandsEvicted != 1 || report.FilesDeleted != 1 || report.BytesReleased == 0 {
t.Fatalf("retention report = %+v", report)
}
var commands, tombstones int
if err := opened.DB().QueryRow(`SELECT count(*) FROM commands WHERE issue_uuid = ?`, old[:]).Scan(&commands); err != nil || commands != 0 {
t.Fatalf("old command count = %d, err = %v", commands, err)
}
if err := opened.DB().QueryRow(`SELECT count(*) FROM commands WHERE issue_uuid = ?`, active[:]).Scan(&commands); err != nil || commands != 1 {
t.Fatalf("active command count = %d, err = %v", commands, err)
}
if err := opened.DB().QueryRow(`SELECT count(*) FROM command_tombstones WHERE issue_uuid = ?`, old[:]).Scan(&tombstones); err != nil || tombstones != 1 {
t.Fatalf("tombstone count = %d, err = %v", tombstones, err)
}
for _, directory := range []string{"segments", "deleting"} {
entries, err := os.ReadDir(filepath.Join(dataDir, directory))
if err != nil || len(entries) != 0 {
t.Fatalf("%s entries = %v, err = %v", directory, entries, err)
}
}
if _, err := opened.RecoverCommandSegments(context.Background()); err != nil {
t.Fatal(err)
}
if err := opened.Close(); err != nil {
t.Fatal(err)
}
}
func TestTombstoneFIFOIsCappedInEvictionTransaction_HP_STORE_09(t *testing.T) {
t.Parallel()
opened := openStore(t, filepath.Join(t.TempDir(), "state"))
now := time.Unix(3_000_000, 0).UTC()
first := uuidBytes(150)
second := uuidBytes(170)
seedCommand(t, opened.DB(), first)
seedCommand(t, opened.DB(), second)
markTerminal(t, opened.DB(), first, now.Add(-2*time.Hour), now.Add(-time.Hour))
markTerminal(t, opened.DB(), second, now.Add(-time.Hour), now.Add(-time.Hour))
report, err := opened.RunRetention(context.Background(), now, store.RetentionPolicy{PressureBytes: 1, TombstoneMaxEntries: 1})
if err != nil {
t.Fatal(err)
}
// Zero-charge commands cannot satisfy byte pressure, so every eligible terminal command is reclaimed.
if report.CommandsEvicted != 2 || report.TombstonesTrimmed != 1 {
t.Fatalf("retention report = %+v", report)
}
var count int
var remaining []byte
if err := opened.DB().QueryRow(`SELECT count(*), max(issue_uuid) FROM command_tombstones`).Scan(&count, &remaining); err != nil {
t.Fatal(err)
}
if count != 1 || !bytes.Equal(remaining, second[:]) {
t.Fatalf("tombstones = count %d remaining %x, want %x", count, remaining, second)
}
if err := opened.Close(); err != nil {
t.Fatal(err)
}
}
func TestRetentionCrashStagesRollForward_CRASH_STORE_04(t *testing.T) {
t.Parallel()
for index, checkpoint := range []string{
store.FaultAfterEvictionMarked,
store.FaultAfterEvictionFilesMoved,
store.FaultAfterEvictionMetadataGone,
} {
t.Run(checkpoint, func(t *testing.T) {
dataDir := filepath.Join(t.TempDir(), "state")
injected := errors.New("injected retention crash")
opened := openStoreWithOptions(t, store.Options{
DataDir: dataDir, BusyTimeout: busyTimeout,
FaultInjector: &failOnceInjector{target: checkpoint, failure: injected},
})
now := time.Unix(4_000_000+int64(index), 0).UTC()
issue := uuidBytes(byte(190 + index*20))
seedCommand(t, opened.DB(), issue)
if _, err := opened.AppendCommandEvent(context.Background(), appendEvent(issue, 1, "evict me")); err != nil {
t.Fatal(err)
}
markTerminal(t, opened.DB(), issue, now.Add(-2*time.Hour), now.Add(-time.Hour))
if _, err := opened.RunRetention(context.Background(), now, store.RetentionPolicy{PressureBytes: 1, TombstoneMaxEntries: 10}); !errors.Is(err, injected) {
t.Fatalf("RunRetention error = %v", err)
}
if err := opened.Close(); err != nil {
t.Fatal(err)
}
reopened := openStore(t, dataDir)
if _, err := reopened.RecoverEvictions(context.Background(), 10); err != nil {
t.Fatal(err)
}
var commandCount, tombstoneCount int
if err := reopened.DB().QueryRow(`SELECT count(*) FROM commands WHERE issue_uuid = ?`, issue[:]).Scan(&commandCount); err != nil {
t.Fatal(err)
}
if err := reopened.DB().QueryRow(`SELECT count(*) FROM command_tombstones WHERE issue_uuid = ?`, issue[:]).Scan(&tombstoneCount); err != nil {
t.Fatal(err)
}
if commandCount != 0 || tombstoneCount != 1 {
t.Fatalf("recovered counts = commands %d tombstones %d", commandCount, tombstoneCount)
}
for _, directory := range []string{"segments", "deleting"} {
entries, err := os.ReadDir(filepath.Join(dataDir, directory))
if err != nil || len(entries) != 0 {
t.Fatalf("%s entries = %v, err = %v", directory, entries, err)
}
}
if _, err := reopened.RecoverCommandSegments(context.Background()); err != nil {
t.Fatal(err)
}
if err := reopened.Close(); err != nil {
t.Fatal(err)
}
})
}
}
func openStore(t *testing.T, dataDir string) *store.Store {
t.Helper()
return openStoreWithOptions(t, store.Options{DataDir: dataDir, BusyTimeout: busyTimeout})
@@ -746,3 +888,10 @@ func assertCommandEventState(t *testing.T, database *sql.DB, issue [16]byte, wan
t.Fatalf("event state = (last=%d, rows=%d), want (%d, %d)", last, events, wantLast, wantEvents)
}
}
func markTerminal(t *testing.T, database *sql.DB, issue [16]byte, issueTime, terminalTime time.Time) {
t.Helper()
if _, err := database.Exec(`UPDATE commands SET issue_time = ?, lifecycle = 5, terminal_time = ? WHERE issue_uuid = ?`, issueTime.UnixNano(), terminalTime.UnixNano(), issue[:]); err != nil {
t.Fatal(err)
}
}