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") }