diff --git a/internal/server/store/event.go b/internal/server/store/event.go index 165f7dc..d34f74f 100644 --- a/internal/server/store/event.go +++ b/internal/server/store/event.go @@ -42,6 +42,8 @@ type EventAppend struct { RawLength uint64 Payload []byte ImmutableSHA256 [32]byte + Output bool + UseCloseout bool } type EventAppendResult struct { @@ -76,8 +78,14 @@ func (store *Store) AppendCommandEvent(ctx context.Context, event EventAppend) ( if err != nil { return result, err } - var lastSequence uint64 - err = database.QueryRowContext(ctx, `SELECT last_event_seq FROM commands WHERE issue_uuid = ?`, event.IssueUUID[:]).Scan(&lastSequence) + var lastSequence, commandOutputCharged, commandCharged, clientCharged, serverCharged, closeoutRemaining uint64 + var clientID string + err = database.QueryRowContext(ctx, `SELECT commands.last_event_seq, commands.output_charged_bytes, +commands.charged_bytes, commands.closeout_remaining_bytes, commands.client_id, clients.charged_bytes, +storage_counters.command_charged_bytes +FROM commands JOIN clients ON clients.client_id = commands.client_id +JOIN storage_counters ON storage_counters.singleton = 1 WHERE commands.issue_uuid = ?`, event.IssueUUID[:]).Scan( + &lastSequence, &commandOutputCharged, &commandCharged, &closeoutRemaining, &clientID, &clientCharged, &serverCharged) if errors.Is(err, sql.ErrNoRows) { return result, ErrCommandNotFound } @@ -97,6 +105,30 @@ func (store *Store) AppendCommandEvent(ctx context.Context, event EventAppend) ( if event.EventSeq != lastSequence+1 { return result, ErrEventSequenceGap } + rowsCharged, indexesCharged := uint64(1), uint64(1) + if store.willCreateSegment(event.IssueUUID, uint64(len(encoded))) { + rowsCharged++ + indexesCharged += 2 + } + chargedBytes, err := EstimateCharge(ChargeInput{EncodedBytes: uint64(len(encoded)), SQLiteRows: rowsCharged, IndexEntries: indexesCharged}) + if err != nil { + return result, err + } + freeBytes, err := store.freeSpaceProbe.AvailableBytes(store.dataDir) + if err != nil { + return result, fmt.Errorf("check filesystem capacity: %w", err) + } + reservation, err := CheckReservation(store.quotaLimits, ReservationState{ + CommandOutputCharged: commandOutputCharged, CommandTotalCharged: commandCharged, + ClientTotalCharged: clientCharged, ServerTotalCharged: serverCharged, + CloseoutRemaining: closeoutRemaining, FilesystemFreeBytes: freeBytes, + }, ReservationRequest{ + ChargedBytes: chargedBytes, PhysicalBytes: uint64(len(encoded)), Output: event.Output || event.Stream > 0, + UseCloseout: event.UseCloseout, + }) + if err != nil { + return result, err + } active, err := store.segmentForAppend(ctx, database, event.IssueUUID, uint64(len(encoded))) if err != nil { @@ -153,8 +185,11 @@ payload, segment_ordinal, segment_record_offset, segment_record_length, immutabl } if err == nil { var update sql.Result - update, err = tx.ExecContext(ctx, `UPDATE commands SET last_event_seq = ?, retained_compressed_bytes = retained_compressed_bytes + ? -WHERE issue_uuid = ? AND last_event_seq = ?`, event.EventSeq, len(event.Payload), event.IssueUUID[:], lastSequence) + update, err = tx.ExecContext(ctx, `UPDATE commands SET last_event_seq = ?, +retained_compressed_bytes = retained_compressed_bytes + ?, output_charged_bytes = ?, charged_bytes = ?, +closeout_remaining_bytes = ? WHERE issue_uuid = ? AND last_event_seq = ?`, event.EventSeq, len(event.Payload), + reservation.CommandOutputCharged, reservation.CommandTotalCharged, reservation.CloseoutRemaining, + event.IssueUUID[:], lastSequence) if err == nil { var affected int64 affected, err = update.RowsAffected() @@ -163,6 +198,28 @@ WHERE issue_uuid = ? AND last_event_seq = ?`, event.EventSeq, len(event.Payload) } } } + if err == nil { + var update sql.Result + update, err = tx.ExecContext(ctx, `UPDATE clients SET charged_bytes = ? WHERE client_id = ? AND charged_bytes = ?`, reservation.ClientTotalCharged, clientID, clientCharged) + if err == nil { + var affected int64 + affected, err = update.RowsAffected() + if err == nil && affected != 1 { + err = errors.New("client quota counter changed unexpectedly") + } + } + } + if err == nil { + var update sql.Result + update, err = tx.ExecContext(ctx, `UPDATE storage_counters SET command_charged_bytes = ? WHERE singleton = 1 AND command_charged_bytes = ?`, reservation.ServerTotalCharged, serverCharged) + if err == nil { + var affected int64 + affected, err = update.RowsAffected() + if err == nil && affected != 1 { + err = errors.New("server quota counter changed unexpectedly") + } + } + } if err != nil { _ = tx.Rollback() store.invalidateActiveSegment(event.IssueUUID) @@ -180,6 +237,11 @@ WHERE issue_uuid = ? AND last_event_seq = ?`, event.EventSeq, len(event.Payload) return result, nil } +func (store *Store) willCreateSegment(owner [16]byte, recordLength uint64) bool { + current := store.activeSegments[owner] + return current == nil || (current.segment.offset != 0 && current.segment.offset+recordLength > store.segmentTarget) +} + func (store *Store) segmentForAppend(ctx context.Context, database *sql.DB, owner [16]byte, recordLength uint64) (*activeSegment, error) { if current := store.activeSegments[owner]; current != nil { if current.segment.offset == 0 || current.segment.offset+recordLength <= store.segmentTarget { diff --git a/internal/server/store/free_unix.go b/internal/server/store/free_unix.go new file mode 100644 index 0000000..12c5652 --- /dev/null +++ b/internal/server/store/free_unix.go @@ -0,0 +1,22 @@ +//go:build !windows + +package store + +import ( + "errors" + "math" + "syscall" +) + +func filesystemFreeBytes(path string) (uint64, error) { + var status syscall.Statfs_t + if err := syscall.Statfs(path, &status); err != nil { + return 0, err + } + blockSize := uint64(status.Bsize) + availableBlocks := uint64(status.Bavail) + if blockSize != 0 && availableBlocks > math.MaxUint64/blockSize { + return 0, errors.New("filesystem free-space value overflowed") + } + return availableBlocks * blockSize, nil +} diff --git a/internal/server/store/free_windows.go b/internal/server/store/free_windows.go new file mode 100644 index 0000000..fe83cef --- /dev/null +++ b/internal/server/store/free_windows.go @@ -0,0 +1,25 @@ +//go:build windows + +package store + +import ( + "path/filepath" + + "golang.org/x/sys/windows" +) + +func filesystemFreeBytes(path string) (uint64, error) { + absolute, err := filepath.Abs(path) + if err != nil { + return 0, err + } + pointer, err := windows.UTF16PtrFromString(absolute) + if err != nil { + return 0, err + } + var available uint64 + if err := windows.GetDiskFreeSpaceEx(pointer, &available, nil, nil); err != nil { + return 0, err + } + return available, nil +} diff --git a/internal/server/store/migrations.go b/internal/server/store/migrations.go index 5816c0a..2305c0d 100644 --- a/internal/server/store/migrations.go +++ b/internal/server/store/migrations.go @@ -79,6 +79,10 @@ CREATE TABLE commands ( retention_status INTEGER NOT NULL DEFAULT 1 CHECK(retention_status BETWEEN 1 AND 3), last_event_seq INTEGER NOT NULL DEFAULT 0 CHECK(last_event_seq >= 0), retained_compressed_bytes INTEGER NOT NULL DEFAULT 0 CHECK(retained_compressed_bytes >= 0), + output_charged_bytes INTEGER NOT NULL DEFAULT 0 CHECK(output_charged_bytes >= 0), + charged_bytes INTEGER NOT NULL DEFAULT 0 CHECK(charged_bytes >= 0), + closeout_remaining_bytes INTEGER NOT NULL DEFAULT 65536 CHECK(closeout_remaining_bytes >= 0), + charge_version INTEGER NOT NULL DEFAULT 1 CHECK(charge_version = 1), output_truncated INTEGER NOT NULL DEFAULT 0 CHECK(output_truncated IN (0,1)), output_incomplete INTEGER NOT NULL DEFAULT 0 CHECK(output_incomplete IN (0,1)), immutable_request_sha256 BLOB NOT NULL CHECK(length(immutable_request_sha256) = 32), @@ -171,4 +175,9 @@ CREATE TABLE storage_incidents ( CHECK((state = 1 AND resolved_at IS NULL) OR (state IN (2,3) AND resolved_at IS NOT NULL)) ) STRICT; CREATE UNIQUE INDEX one_open_incident_per_scope_kind ON storage_incidents(scope, scope_key, kind) WHERE state = 1; +CREATE TABLE storage_counters ( + singleton INTEGER PRIMARY KEY CHECK(singleton = 1), command_charged_bytes INTEGER NOT NULL CHECK(command_charged_bytes >= 0), + charge_version INTEGER NOT NULL CHECK(charge_version = 1) +) STRICT; +INSERT INTO storage_counters(singleton, command_charged_bytes, charge_version) VALUES (1, 0, 1); ` diff --git a/internal/server/store/quota.go b/internal/server/store/quota.go new file mode 100644 index 0000000..9dba4a7 --- /dev/null +++ b/internal/server/store/quota.go @@ -0,0 +1,203 @@ +package store + +import ( + "errors" + "fmt" + "math" +) + +const ( + ChargeFormulaVersion uint32 = 1 + ChargeSQLiteRowBytes uint64 = 192 + ChargeIndexEntryBytes uint64 = 64 +) + +var ErrCapacityExhausted = errors.New("storage capacity exhausted") + +type FreeSpaceProbe interface { + AvailableBytes(path string) (uint64, error) +} + +type FreeSpaceProbeFunc func(path string) (uint64, error) + +func (function FreeSpaceProbeFunc) AvailableBytes(path string) (uint64, error) { return function(path) } + +type CapacityTier string + +const ( + CapacityTierHardMaximum CapacityTier = "hard_maximum" + CapacityTierCommandOut CapacityTier = "command_output" + CapacityTierCommand CapacityTier = "command_total" + CapacityTierClient CapacityTier = "client_total" + CapacityTierServer CapacityTier = "server_total" + CapacityTierFilesystem CapacityTier = "filesystem_floor" +) + +type CapacityError struct { + Tier CapacityTier + Requested uint64 + Available uint64 +} + +func (failure *CapacityError) Error() string { + return fmt.Sprintf("%s: tier=%s requested=%d available=%d", ErrCapacityExhausted, failure.Tier, failure.Requested, failure.Available) +} + +func (failure *CapacityError) Unwrap() error { return ErrCapacityExhausted } + +type ChargeInput struct { + EncodedBytes uint64 + SQLiteRows uint64 + IndexEntries uint64 +} + +func EstimateCharge(input ChargeInput) (uint64, error) { + rowCharge, overflow := multiplyChecked(input.SQLiteRows, ChargeSQLiteRowBytes) + if overflow { + return 0, &CapacityError{Tier: CapacityTierHardMaximum, Requested: math.MaxUint64} + } + indexCharge, overflow := multiplyChecked(input.IndexEntries, ChargeIndexEntryBytes) + if overflow { + return 0, &CapacityError{Tier: CapacityTierHardMaximum, Requested: math.MaxUint64} + } + result, overflow := addChecked(input.EncodedBytes, rowCharge) + if !overflow { + result, overflow = addChecked(result, indexCharge) + } + if overflow { + return 0, &CapacityError{Tier: CapacityTierHardMaximum, Requested: math.MaxUint64} + } + return result, nil +} + +type QuotaLimits struct { + HardAllocationBytes uint64 + CommandOutputBytes uint64 + CommandTotalBytes uint64 + ClientTotalBytes uint64 + ServerTotalBytes uint64 + CloseoutReserveBytes uint64 + FilesystemFloorBytes uint64 +} + +func DefaultQuotaLimits() QuotaLimits { + return QuotaLimits{ + HardAllocationBytes: DefaultSegmentLimit, CommandOutputBytes: 10 << 20, CommandTotalBytes: 32 << 20, + ClientTotalBytes: 256 << 20, ServerTotalBytes: 4 << 30, CloseoutReserveBytes: 64 << 10, + FilesystemFloorBytes: 256 << 20, + } +} + +func (limits QuotaLimits) Validate() error { + if limits.HardAllocationBytes == 0 || limits.CommandOutputBytes == 0 || limits.CommandTotalBytes <= limits.CommandOutputBytes || limits.ClientTotalBytes <= limits.CommandTotalBytes || limits.ServerTotalBytes <= limits.ClientTotalBytes || limits.CloseoutReserveBytes == 0 || limits.CloseoutReserveBytes >= limits.CommandTotalBytes || limits.FilesystemFloorBytes == 0 { + return errors.New("invalid quota limits") + } + return nil +} + +type ReservationState struct { + CommandOutputCharged uint64 + CommandTotalCharged uint64 + ClientTotalCharged uint64 + ServerTotalCharged uint64 + CloseoutRemaining uint64 + FilesystemFreeBytes uint64 +} + +type ReservationRequest struct { + ChargedBytes uint64 + PhysicalBytes uint64 + Output bool + UseCloseout bool +} + +type ReservationDecision struct { + CommandOutputCharged uint64 + CommandTotalCharged uint64 + ClientTotalCharged uint64 + ServerTotalCharged uint64 + CloseoutRemaining uint64 +} + +func CheckReservation(limits QuotaLimits, state ReservationState, request ReservationRequest) (ReservationDecision, error) { + var result ReservationDecision + if err := limits.Validate(); err != nil { + return result, err + } + if request.ChargedBytes == 0 || request.ChargedBytes > limits.HardAllocationBytes { + return result, capacityFailure(CapacityTierHardMaximum, request.ChargedBytes, limits.HardAllocationBytes) + } + output, overflow := addChecked(state.CommandOutputCharged, request.ChargedBytes) + if request.Output && (overflow || output > limits.CommandOutputBytes) { + return result, capacityFailure(CapacityTierCommandOut, request.ChargedBytes, available(limits.CommandOutputBytes, state.CommandOutputCharged)) + } + command, overflow := addChecked(state.CommandTotalCharged, request.ChargedBytes) + if overflow { + return result, capacityFailure(CapacityTierCommand, request.ChargedBytes, available(limits.CommandTotalBytes, state.CommandTotalCharged)) + } + closeoutRemaining := state.CloseoutRemaining + if request.UseCloseout { + if request.ChargedBytes >= closeoutRemaining { + closeoutRemaining = 0 + } else { + closeoutRemaining -= request.ChargedBytes + } + } else if _, overflow := addChecked(command, closeoutRemaining); overflow || command+closeoutRemaining > limits.CommandTotalBytes { + return result, capacityFailure(CapacityTierCommand, request.ChargedBytes, availableForNormalCommand(limits.CommandTotalBytes, state.CommandTotalCharged, closeoutRemaining)) + } + if command > limits.CommandTotalBytes { + return result, capacityFailure(CapacityTierCommand, request.ChargedBytes, available(limits.CommandTotalBytes, state.CommandTotalCharged)) + } + client, overflow := addChecked(state.ClientTotalCharged, request.ChargedBytes) + if overflow || client > limits.ClientTotalBytes { + return result, capacityFailure(CapacityTierClient, request.ChargedBytes, available(limits.ClientTotalBytes, state.ClientTotalCharged)) + } + server, overflow := addChecked(state.ServerTotalCharged, request.ChargedBytes) + if overflow || server > limits.ServerTotalBytes { + return result, capacityFailure(CapacityTierServer, request.ChargedBytes, available(limits.ServerTotalBytes, state.ServerTotalCharged)) + } + requiredFree, overflow := addChecked(limits.FilesystemFloorBytes, request.PhysicalBytes) + if overflow || state.FilesystemFreeBytes < requiredFree { + return result, capacityFailure(CapacityTierFilesystem, request.PhysicalBytes, available(state.FilesystemFreeBytes, limits.FilesystemFloorBytes)) + } + if !request.Output { + output = state.CommandOutputCharged + } + return ReservationDecision{ + CommandOutputCharged: output, CommandTotalCharged: command, ClientTotalCharged: client, + ServerTotalCharged: server, CloseoutRemaining: closeoutRemaining, + }, nil +} + +func capacityFailure(tier CapacityTier, requested, availableBytes uint64) error { + return &CapacityError{Tier: tier, Requested: requested, Available: availableBytes} +} + +func available(limit, used uint64) uint64 { + if used >= limit { + return 0 + } + return limit - used +} + +func availableForNormalCommand(limit, used, closeout uint64) uint64 { + remaining := available(limit, used) + if closeout >= remaining { + return 0 + } + return remaining - closeout +} + +func addChecked(left, right uint64) (uint64, bool) { + if left > math.MaxUint64-right { + return 0, true + } + return left + right, false +} + +func multiplyChecked(left, right uint64) (uint64, bool) { + if left != 0 && right > math.MaxUint64/left { + return 0, true + } + return left * right, false +} diff --git a/internal/server/store/quota_test.go b/internal/server/store/quota_test.go new file mode 100644 index 0000000..3d73fe4 --- /dev/null +++ b/internal/server/store/quota_test.go @@ -0,0 +1,109 @@ +package store + +import ( + "errors" + "math" + "testing" +) + +func TestChargeFormulaVersionAndOverflow_HP_STORE_04(t *testing.T) { + t.Parallel() + + if ChargeFormulaVersion != 1 { + t.Fatalf("charge formula version = %d", ChargeFormulaVersion) + } + charge, err := EstimateCharge(ChargeInput{EncodedBytes: 100, SQLiteRows: 2, IndexEntries: 3}) + want := uint64(100 + 2*ChargeSQLiteRowBytes + 3*ChargeIndexEntryBytes) + if err != nil || charge != want { + t.Fatalf("EstimateCharge = (%d, %v), want %d", charge, err, want) + } + for _, input := range []ChargeInput{ + {SQLiteRows: math.MaxUint64}, + {IndexEntries: math.MaxUint64}, + {EncodedBytes: math.MaxUint64, SQLiteRows: 1}, + } { + if _, err := EstimateCharge(input); !errors.Is(err, ErrCapacityExhausted) { + t.Errorf("EstimateCharge(%+v) error = %v", input, err) + } + } +} + +func TestTieredReservationOrderAndBoundaries_HP_STORE_05(t *testing.T) { + t.Parallel() + + limits := quotaTestLimits() + base := ReservationState{ + CommandOutputCharged: 10, CommandTotalCharged: 20, ClientTotalCharged: 30, + ServerTotalCharged: 40, CloseoutRemaining: 20, FilesystemFreeBytes: 1000, + } + decision, err := CheckReservation(limits, base, ReservationRequest{ChargedBytes: 10, PhysicalBytes: 5, Output: true}) + if err != nil { + t.Fatal(err) + } + if decision.CommandOutputCharged != 20 || decision.CommandTotalCharged != 30 || decision.ClientTotalCharged != 40 || decision.ServerTotalCharged != 50 || decision.CloseoutRemaining != 20 { + t.Fatalf("decision = %+v", decision) + } + + cases := []struct { + name string + state ReservationState + request ReservationRequest + tier CapacityTier + }{ + {"hard maximum first", base, ReservationRequest{ChargedBytes: 101, PhysicalBytes: 1, Output: true}, CapacityTierHardMaximum}, + {"output", withState(base, 95, 20, 30, 40, 20, 1000), ReservationRequest{ChargedBytes: 10, PhysicalBytes: 1, Output: true}, CapacityTierCommandOut}, + {"command reserve", withState(base, 10, 175, 175, 175, 20, 1000), ReservationRequest{ChargedBytes: 10, PhysicalBytes: 1}, CapacityTierCommand}, + {"client", withState(base, 10, 20, 295, 295, 20, 1000), ReservationRequest{ChargedBytes: 10, PhysicalBytes: 1}, CapacityTierClient}, + {"server", withState(base, 10, 20, 30, 395, 20, 1000), ReservationRequest{ChargedBytes: 10, PhysicalBytes: 1}, CapacityTierServer}, + {"filesystem", withState(base, 10, 20, 30, 40, 20, 54), ReservationRequest{ChargedBytes: 10, PhysicalBytes: 5}, CapacityTierFilesystem}, + } + for _, test := range cases { + t.Run(test.name, func(t *testing.T) { + _, err := CheckReservation(limits, test.state, test.request) + var capacity *CapacityError + if !errors.As(err, &capacity) || capacity.Tier != test.tier { + t.Fatalf("error = %v, want tier %s", err, test.tier) + } + }) + } +} + +func TestCloseoutReservationCanBeConsumed_BH_STORE_06(t *testing.T) { + t.Parallel() + + limits := quotaTestLimits() + state := ReservationState{ + CommandOutputCharged: 90, CommandTotalCharged: 190, ClientTotalCharged: 200, + ServerTotalCharged: 200, CloseoutRemaining: 10, FilesystemFreeBytes: 1000, + } + if _, err := CheckReservation(limits, state, ReservationRequest{ChargedBytes: 10, PhysicalBytes: 1}); err == nil { + t.Fatal("normal allocation consumed closeout headroom") + } + decision, err := CheckReservation(limits, state, ReservationRequest{ChargedBytes: 10, PhysicalBytes: 1, UseCloseout: true}) + if err != nil { + t.Fatal(err) + } + if decision.CommandTotalCharged != 200 || decision.CloseoutRemaining != 0 { + t.Fatalf("closeout decision = %+v", decision) + } + if _, err := CheckReservation(limits, state, ReservationRequest{ChargedBytes: math.MaxUint64, PhysicalBytes: 1, UseCloseout: true}); !errors.Is(err, ErrCapacityExhausted) { + t.Fatalf("overflow error = %v", err) + } +} + +func quotaTestLimits() QuotaLimits { + return QuotaLimits{ + HardAllocationBytes: 100, CommandOutputBytes: 100, CommandTotalBytes: 200, + ClientTotalBytes: 300, ServerTotalBytes: 400, CloseoutReserveBytes: 20, FilesystemFloorBytes: 50, + } +} + +func withState(base ReservationState, output, command, client, server, closeout, free uint64) ReservationState { + base.CommandOutputCharged = output + base.CommandTotalCharged = command + base.ClientTotalCharged = client + base.ServerTotalCharged = server + base.CloseoutRemaining = closeout + base.FilesystemFreeBytes = free + return base +} diff --git a/internal/server/store/recovery.go b/internal/server/store/recovery.go index 61474ce..cfd38d0 100644 --- a/internal/server/store/recovery.go +++ b/internal/server/store/recovery.go @@ -9,7 +9,10 @@ import ( "path/filepath" ) -var ErrIntegrityCheck = errors.New("SQLite integrity check failed") +var ( + ErrIntegrityCheck = errors.New("SQLite integrity check failed") + ErrQuotaCounterMismatch = errors.New("quota counter mismatch") +) type RecoveryReport struct { SegmentsChecked uint64 @@ -119,12 +122,36 @@ FROM output_segments ORDER BY issue_uuid, ordinal`) return report, err } } + if err := checkQuotaCounters(ctx, database); err != nil { + return report, err + } if _, err := database.ExecContext(ctx, `UPDATE output_segments SET sealed = 1 WHERE sealed = 0`); err != nil { return report, err } return report, nil } +func checkQuotaCounters(ctx context.Context, database *sql.DB) error { + var mismatchedClients int + if err := database.QueryRowContext(ctx, `SELECT count(*) FROM clients WHERE charged_bytes != +COALESCE((SELECT sum(commands.charged_bytes) FROM commands WHERE commands.client_id = clients.client_id), 0)`).Scan(&mismatchedClients); err != nil { + return err + } + var storedServer, computedServer uint64 + if err := database.QueryRowContext(ctx, `SELECT command_charged_bytes, +COALESCE((SELECT sum(charged_bytes) FROM commands), 0) FROM storage_counters WHERE singleton = 1`).Scan(&storedServer, &computedServer); err != nil { + return err + } + var invalidCommands int + if err := database.QueryRowContext(ctx, `SELECT count(*) FROM commands WHERE output_charged_bytes > charged_bytes OR charge_version != ?`, ChargeFormulaVersion).Scan(&invalidCommands); err != nil { + return err + } + if mismatchedClients != 0 || storedServer != computedServer || invalidCommands != 0 { + return ErrQuotaCounterMismatch + } + return nil +} + func (store *Store) quickCheckDatabase(ctx context.Context, database *sql.DB) error { rows, err := database.QueryContext(ctx, `PRAGMA quick_check(100)`) if err != nil { diff --git a/internal/server/store/retention.go b/internal/server/store/retention.go new file mode 100644 index 0000000..59cadd6 --- /dev/null +++ b/internal/server/store/retention.go @@ -0,0 +1,98 @@ +package store + +import ( + "bytes" + "sort" + "time" +) + +type RetentionCandidate struct { + IssueUUID [16]byte + IssueTime time.Time + TerminalTime time.Time + ChargedBytes uint64 + Terminal bool +} + +type RetentionSelection struct { + IssueUUIDs [][16]byte + ReleasedBytes uint64 +} + +func SelectWholeCommandEvictions(now time.Time, age time.Duration, bytesNeeded uint64, candidates []RetentionCandidate) RetentionSelection { + ordered := append([]RetentionCandidate(nil), candidates...) + sort.Slice(ordered, func(left, right int) bool { + if ordered[left].IssueTime.Equal(ordered[right].IssueTime) { + return bytes.Compare(ordered[left].IssueUUID[:], ordered[right].IssueUUID[:]) < 0 + } + return ordered[left].IssueTime.Before(ordered[right].IssueTime) + }) + selected := make(map[[16]byte]bool, len(ordered)) + var result RetentionSelection + selectCandidate := func(candidate RetentionCandidate) { + if selected[candidate.IssueUUID] { + return + } + selected[candidate.IssueUUID] = true + result.IssueUUIDs = append(result.IssueUUIDs, candidate.IssueUUID) + if total, overflow := addChecked(result.ReleasedBytes, candidate.ChargedBytes); overflow { + result.ReleasedBytes = ^uint64(0) + } else { + result.ReleasedBytes = total + } + } + if age > 0 { + cutoff := now.Add(-age) + for _, candidate := range ordered { + if candidate.Terminal && !candidate.TerminalTime.IsZero() && !candidate.TerminalTime.After(cutoff) { + selectCandidate(candidate) + } + } + } + for _, candidate := range ordered { + if result.ReleasedBytes >= bytesNeeded { + break + } + if candidate.Terminal { + selectCandidate(candidate) + } + } + return result +} + +type OutputSegmentCandidate struct { + Ordinal uint32 + ChargedBytes uint64 + Sealed bool +} + +func SelectOutputRotation(currentBytes, limit uint64, candidates []OutputSegmentCandidate) []uint32 { + if currentBytes <= limit { + return nil + } + ordered := append([]OutputSegmentCandidate(nil), candidates...) + sort.Slice(ordered, func(left, right int) bool { return ordered[left].Ordinal < ordered[right].Ordinal }) + var selected []uint32 + for _, candidate := range ordered { + if currentBytes <= limit { + break + } + if !candidate.Sealed { + continue + } + selected = append(selected, candidate.Ordinal) + if candidate.ChargedBytes >= currentBytes { + currentBytes = 0 + } else { + currentBytes -= candidate.ChargedBytes + } + } + return selected +} + +func FIFOEntriesToRemove(current, maximum uint64) uint64 { + if current <= maximum { + return 0 + } + return current - maximum +} diff --git a/internal/server/store/retention_test.go b/internal/server/store/retention_test.go new file mode 100644 index 0000000..87af1e6 --- /dev/null +++ b/internal/server/store/retention_test.go @@ -0,0 +1,57 @@ +package store + +import ( + "reflect" + "testing" + "time" +) + +func TestWholeCommandEvictionOrderingAndAge_HP_STORE_06(t *testing.T) { + t.Parallel() + + now := time.Unix(2_000_000, 0) + first := retentionUUID(1) + second := retentionUUID(2) + third := retentionUUID(3) + active := retentionUUID(4) + candidates := []RetentionCandidate{ + {IssueUUID: third, IssueTime: now.Add(-time.Hour), TerminalTime: now.Add(-31 * 24 * time.Hour), ChargedBytes: 30, Terminal: true}, + {IssueUUID: active, IssueTime: now.Add(-100 * 24 * time.Hour), ChargedBytes: 1000}, + {IssueUUID: second, IssueTime: now.Add(-2 * time.Hour), TerminalTime: now.Add(-time.Hour), ChargedBytes: 20, Terminal: true}, + {IssueUUID: first, IssueTime: now.Add(-2 * time.Hour), TerminalTime: now.Add(-time.Hour), ChargedBytes: 10, Terminal: true}, + } + selection := SelectWholeCommandEvictions(now, 30*24*time.Hour, 45, candidates) + want := [][16]byte{third, first, second} + if !reflect.DeepEqual(selection.IssueUUIDs, want) || selection.ReleasedBytes != 60 { + t.Fatalf("selection = %+v, want IDs %v and 60 bytes", selection, want) + } + selection = SelectWholeCommandEvictions(now, 0, 0, candidates) + if len(selection.IssueUUIDs) != 0 { + t.Fatalf("disabled age/no pressure selected %v", selection.IssueUUIDs) + } +} + +func TestOutputRotationProtectsActiveSegmentAndFIFO_BH_STORE_07(t *testing.T) { + t.Parallel() + + selected := SelectOutputRotation(130, 100, []OutputSegmentCandidate{ + {Ordinal: 2, ChargedBytes: 50, Sealed: true}, + {Ordinal: 0, ChargedBytes: 20, Sealed: true}, + {Ordinal: 1, ChargedBytes: 20, Sealed: false}, + }) + if !reflect.DeepEqual(selected, []uint32{0, 2}) { + t.Fatalf("rotation = %v", selected) + } + if got := FIFOEntriesToRemove(1_000_001, 1_000_000); got != 1 { + t.Fatalf("FIFO removal = %d", got) + } + if got := FIFOEntriesToRemove(10, 10); got != 0 { + t.Fatalf("FIFO boundary removal = %d", got) + } +} + +func retentionUUID(last byte) [16]byte { + var result [16]byte + result[15] = last + return result +} diff --git a/internal/server/store/store.go b/internal/server/store/store.go index 3b4727e..cf518b0 100644 --- a/internal/server/store/store.go +++ b/internal/server/store/store.go @@ -25,6 +25,8 @@ type Options struct { DataDir string BusyTimeout time.Duration SegmentTargetSize uint64 + QuotaLimits QuotaLimits + FreeSpaceProbe FreeSpaceProbe FaultInjector FaultInjector } @@ -33,6 +35,8 @@ type Store struct { unlock func() error dataDir string segmentTarget uint64 + quotaLimits QuotaLimits + freeSpaceProbe FreeSpaceProbe faultInjector FaultInjector activeSegments map[[16]byte]*activeSegment mu sync.Mutex @@ -52,6 +56,15 @@ func Open(ctx context.Context, options Options) (*Store, error) { if options.SegmentTargetSize > DefaultSegmentLimit { return nil, fmt.Errorf("segment target exceeds hard record limit") } + if options.QuotaLimits == (QuotaLimits{}) { + options.QuotaLimits = DefaultQuotaLimits() + } + if err := options.QuotaLimits.Validate(); err != nil { + return nil, err + } + if options.FreeSpaceProbe == nil { + options.FreeSpaceProbe = FreeSpaceProbeFunc(filesystemFreeBytes) + } if err := ensurePrivateDirectory(options.DataDir); err != nil { return nil, err } @@ -88,6 +101,7 @@ func Open(ctx context.Context, options Options) (*Store, error) { db.SetMaxIdleConns(1) store := &Store{ db: db, unlock: unlock, dataDir: options.DataDir, segmentTarget: options.SegmentTargetSize, + quotaLimits: options.QuotaLimits, freeSpaceProbe: options.FreeSpaceProbe, faultInjector: options.FaultInjector, activeSegments: make(map[[16]byte]*activeSegment), } if err := db.PingContext(ctx); err != nil { diff --git a/test/coverage.toml b/test/coverage.toml index bd88b8b..7560fd1 100644 --- a/test/coverage.toml +++ b/test/coverage.toml @@ -261,3 +261,45 @@ id = "CRASH-STORE-03" layer = "integration" status = "implemented" tests = ["test/integration/store/store_integration_test.go:TestEventCrashOnExistingSegmentTruncatesOnlyTail_CRASH_STORE_03"] + +[[requirements]] +id = "HP-STORE-04" +layer = "unit" +status = "implemented" +tests = ["internal/server/store/quota_test.go:TestChargeFormulaVersionAndOverflow_HP_STORE_04"] + +[[requirements]] +id = "HP-STORE-05" +layer = "unit" +status = "implemented" +tests = ["internal/server/store/quota_test.go:TestTieredReservationOrderAndBoundaries_HP_STORE_05"] + +[[requirements]] +id = "BH-STORE-06" +layer = "unit" +status = "implemented" +tests = ["internal/server/store/quota_test.go:TestCloseoutReservationCanBeConsumed_BH_STORE_06"] + +[[requirements]] +id = "HP-STORE-06" +layer = "unit" +status = "implemented" +tests = ["internal/server/store/retention_test.go:TestWholeCommandEvictionOrderingAndAge_HP_STORE_06"] + +[[requirements]] +id = "BH-STORE-07" +layer = "unit" +status = "implemented" +tests = ["internal/server/store/retention_test.go:TestOutputRotationProtectsActiveSegmentAndFIFO_BH_STORE_07"] + +[[requirements]] +id = "HP-STORE-07" +layer = "integration" +status = "implemented" +tests = ["test/integration/store/store_integration_test.go:TestEventQuotaCountersAtomicAndCapacityRejectsBeforeWrite_HP_STORE_07"] + +[[requirements]] +id = "BH-STORE-08" +layer = "integration" +status = "implemented" +tests = ["test/integration/store/store_integration_test.go:TestFilesystemFloorAndCounterMismatch_BH_STORE_08"] diff --git a/test/integration/store/store_integration_test.go b/test/integration/store/store_integration_test.go index 405df66..2fb7803 100644 --- a/test/integration/store/store_integration_test.go +++ b/test/integration/store/store_integration_test.go @@ -66,7 +66,7 @@ func TestRealSQLiteInitializationAndRestart_HP_STORE_01(t *testing.T) { if err := rows.Close(); err != nil { t.Fatal(err) } - want := []string{"audit_events", "clients", "command_events", "command_payloads", "command_tombstones", "commands", "control_mutations", "output_segments", "output_truncations", "schema_migrations", "sessions", "stdin_writes", "storage_incidents", "takeover_authorizations"} + want := []string{"audit_events", "clients", "command_events", "command_payloads", "command_tombstones", "commands", "control_mutations", "output_segments", "output_truncations", "schema_migrations", "sessions", "stdin_writes", "storage_counters", "storage_incidents", "takeover_authorizations"} sort.Strings(want) if strings.Join(names, ",") != strings.Join(want, ",") { t.Fatalf("tables = %v, want %v", names, want) @@ -531,6 +531,112 @@ func TestEventGapRejectedBeforeFileWriteAndRotation_HP_STORE_03(t *testing.T) { } } +func TestEventQuotaCountersAtomicAndCapacityRejectsBeforeWrite_HP_STORE_07(t *testing.T) { + t.Parallel() + + dataDir := filepath.Join(t.TempDir(), "state") + limits := store.QuotaLimits{ + HardAllocationBytes: 1000, CommandOutputBytes: 1000, CommandTotalBytes: 2000, + ClientTotalBytes: 3000, ServerTotalBytes: 4000, CloseoutReserveBytes: 200, FilesystemFloorBytes: 50, + } + opened := openStoreWithOptions(t, store.Options{ + DataDir: dataDir, BusyTimeout: busyTimeout, QuotaLimits: limits, + FreeSpaceProbe: store.FreeSpaceProbeFunc(func(string) (uint64, error) { return 10_000, nil }), + }) + issue := uuidBytes(90) + seedCommand(t, opened.DB(), issue) + if _, err := opened.DB().Exec(`UPDATE commands SET closeout_remaining_bytes = ? WHERE issue_uuid = ?`, limits.CloseoutReserveBytes, issue[:]); err != nil { + t.Fatal(err) + } + first := appendEvent(issue, 1, "first") + if _, err := opened.AppendCommandEvent(context.Background(), first); err != nil { + t.Fatal(err) + } + var commandOutput, commandTotal, closeout, clientTotal, serverTotal uint64 + if err := opened.DB().QueryRow(`SELECT output_charged_bytes, charged_bytes, closeout_remaining_bytes FROM commands WHERE issue_uuid = ?`, issue[:]).Scan(&commandOutput, &commandTotal, &closeout); err != nil { + t.Fatal(err) + } + if err := opened.DB().QueryRow(`SELECT charged_bytes FROM clients WHERE client_id = 'client-a'`).Scan(&clientTotal); err != nil { + t.Fatal(err) + } + if err := opened.DB().QueryRow(`SELECT command_charged_bytes FROM storage_counters WHERE singleton = 1`).Scan(&serverTotal); err != nil { + t.Fatal(err) + } + if commandOutput == 0 || commandOutput != commandTotal || commandTotal != clientTotal || clientTotal != serverTotal || closeout != limits.CloseoutReserveBytes { + t.Fatalf("quota counters = output:%d command:%d client:%d server:%d closeout:%d", commandOutput, commandTotal, clientTotal, serverTotal, closeout) + } + entries, err := os.ReadDir(filepath.Join(dataDir, "segments")) + if err != nil || len(entries) != 1 { + t.Fatalf("segment entries = %v, err = %v", entries, err) + } + segmentPath := filepath.Join(dataDir, "segments", entries[0].Name()) + before, err := os.Stat(segmentPath) + if err != nil { + t.Fatal(err) + } + if result, err := opened.AppendCommandEvent(context.Background(), first); err != nil || !result.Duplicate { + t.Fatalf("duplicate = (%+v, %v)", result, err) + } + second := appendEvent(issue, 2, "second") + if _, err := opened.AppendCommandEvent(context.Background(), second); err == nil { + t.Fatal("over-quota event was accepted") + } else { + var capacity *store.CapacityError + if !errors.As(err, &capacity) || capacity.Tier != store.CapacityTierCommandOut { + t.Fatalf("capacity error = %v", err) + } + } + after, err := os.Stat(segmentPath) + if err != nil { + t.Fatal(err) + } + if after.Size() != before.Size() { + t.Fatalf("rejected event changed segment size from %d to %d", before.Size(), after.Size()) + } + assertCommandEventState(t, opened.DB(), issue, 1, 1) + if err := opened.Close(); err != nil { + t.Fatal(err) + } +} + +func TestFilesystemFloorAndCounterMismatch_BH_STORE_08(t *testing.T) { + t.Parallel() + + dataDir := filepath.Join(t.TempDir(), "state") + opened := openStoreWithOptions(t, store.Options{ + DataDir: dataDir, BusyTimeout: busyTimeout, + FreeSpaceProbe: store.FreeSpaceProbeFunc(func(string) (uint64, error) { return 0, nil }), + }) + issue := uuidBytes(110) + seedCommand(t, opened.DB(), issue) + if _, err := opened.AppendCommandEvent(context.Background(), appendEvent(issue, 1, "floor")); err == nil { + t.Fatal("filesystem-floor event was accepted") + } else { + var capacity *store.CapacityError + if !errors.As(err, &capacity) || capacity.Tier != store.CapacityTierFilesystem { + t.Fatalf("floor error = %v", err) + } + } + entries, err := os.ReadDir(filepath.Join(dataDir, "segments")) + if err != nil || len(entries) != 0 { + t.Fatalf("segments after floor rejection = %v, err = %v", entries, err) + } + if err := opened.Close(); err != nil { + t.Fatal(err) + } + + reopened := openStore(t, dataDir) + if _, err := reopened.DB().Exec(`UPDATE clients SET charged_bytes = 1 WHERE client_id = 'client-a'`); err != nil { + t.Fatal(err) + } + if _, err := reopened.RecoverCommandSegments(context.Background()); !errors.Is(err, store.ErrQuotaCounterMismatch) { + t.Fatalf("counter mismatch recovery error = %v", 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})