diff --git a/internal/server/store/retention_store.go b/internal/server/store/retention_store.go new file mode 100644 index 0000000..c6f20ec --- /dev/null +++ b/internal/server/store/retention_store.go @@ -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") } diff --git a/internal/server/store/store.go b/internal/server/store/store.go index cf518b0..24e02b0 100644 --- a/internal/server/store/store.go +++ b/internal/server/store/store.go @@ -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 } diff --git a/test/coverage.toml b/test/coverage.toml index 7560fd1..d379e55 100644 --- a/test/coverage.toml +++ b/test/coverage.toml @@ -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"] diff --git a/test/integration/store/store_integration_test.go b/test/integration/store/store_integration_test.go index 2fb7803..53e242d 100644 --- a/test/integration/store/store_integration_test.go +++ b/test/integration/store/store_integration_test.go @@ -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) + } +}