feat: enforce tiered storage reservations

This commit is contained in:
2026-08-31 09:48:04 +00:00
parent 583b539f86
commit 611100f67c
12 changed files with 780 additions and 6 deletions
+66 -4
View File
@@ -42,6 +42,8 @@ type EventAppend struct {
RawLength uint64 RawLength uint64
Payload []byte Payload []byte
ImmutableSHA256 [32]byte ImmutableSHA256 [32]byte
Output bool
UseCloseout bool
} }
type EventAppendResult struct { type EventAppendResult struct {
@@ -76,8 +78,14 @@ func (store *Store) AppendCommandEvent(ctx context.Context, event EventAppend) (
if err != nil { if err != nil {
return result, err return result, err
} }
var lastSequence uint64 var lastSequence, commandOutputCharged, commandCharged, clientCharged, serverCharged, closeoutRemaining uint64
err = database.QueryRowContext(ctx, `SELECT last_event_seq FROM commands WHERE issue_uuid = ?`, event.IssueUUID[:]).Scan(&lastSequence) 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) { if errors.Is(err, sql.ErrNoRows) {
return result, ErrCommandNotFound return result, ErrCommandNotFound
} }
@@ -97,6 +105,30 @@ func (store *Store) AppendCommandEvent(ctx context.Context, event EventAppend) (
if event.EventSeq != lastSequence+1 { if event.EventSeq != lastSequence+1 {
return result, ErrEventSequenceGap 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))) active, err := store.segmentForAppend(ctx, database, event.IssueUUID, uint64(len(encoded)))
if err != nil { if err != nil {
@@ -153,8 +185,11 @@ payload, segment_ordinal, segment_record_offset, segment_record_length, immutabl
} }
if err == nil { if err == nil {
var update sql.Result var update sql.Result
update, err = tx.ExecContext(ctx, `UPDATE commands SET last_event_seq = ?, retained_compressed_bytes = retained_compressed_bytes + ? update, err = tx.ExecContext(ctx, `UPDATE commands SET last_event_seq = ?,
WHERE issue_uuid = ? AND last_event_seq = ?`, event.EventSeq, len(event.Payload), event.IssueUUID[:], lastSequence) 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 { if err == nil {
var affected int64 var affected int64
affected, err = update.RowsAffected() 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 { if err != nil {
_ = tx.Rollback() _ = tx.Rollback()
store.invalidateActiveSegment(event.IssueUUID) store.invalidateActiveSegment(event.IssueUUID)
@@ -180,6 +237,11 @@ WHERE issue_uuid = ? AND last_event_seq = ?`, event.EventSeq, len(event.Payload)
return result, nil 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) { 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 := store.activeSegments[owner]; current != nil {
if current.segment.offset == 0 || current.segment.offset+recordLength <= store.segmentTarget { if current.segment.offset == 0 || current.segment.offset+recordLength <= store.segmentTarget {
+22
View File
@@ -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
}
+25
View File
@@ -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
}
+9
View File
@@ -79,6 +79,10 @@ CREATE TABLE commands (
retention_status INTEGER NOT NULL DEFAULT 1 CHECK(retention_status BETWEEN 1 AND 3), 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), 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), 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_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)), 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), 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)) CHECK((state = 1 AND resolved_at IS NULL) OR (state IN (2,3) AND resolved_at IS NOT NULL))
) STRICT; ) STRICT;
CREATE UNIQUE INDEX one_open_incident_per_scope_kind ON storage_incidents(scope, scope_key, kind) WHERE state = 1; 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);
` `
+203
View File
@@ -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
}
+109
View File
@@ -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
}
+28 -1
View File
@@ -9,7 +9,10 @@ import (
"path/filepath" "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 { type RecoveryReport struct {
SegmentsChecked uint64 SegmentsChecked uint64
@@ -119,12 +122,36 @@ FROM output_segments ORDER BY issue_uuid, ordinal`)
return report, err 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 { if _, err := database.ExecContext(ctx, `UPDATE output_segments SET sealed = 1 WHERE sealed = 0`); err != nil {
return report, err return report, err
} }
return report, nil 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 { func (store *Store) quickCheckDatabase(ctx context.Context, database *sql.DB) error {
rows, err := database.QueryContext(ctx, `PRAGMA quick_check(100)`) rows, err := database.QueryContext(ctx, `PRAGMA quick_check(100)`)
if err != nil { if err != nil {
+98
View File
@@ -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
}
+57
View File
@@ -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
}
+14
View File
@@ -25,6 +25,8 @@ type Options struct {
DataDir string DataDir string
BusyTimeout time.Duration BusyTimeout time.Duration
SegmentTargetSize uint64 SegmentTargetSize uint64
QuotaLimits QuotaLimits
FreeSpaceProbe FreeSpaceProbe
FaultInjector FaultInjector FaultInjector FaultInjector
} }
@@ -33,6 +35,8 @@ type Store struct {
unlock func() error unlock func() error
dataDir string dataDir string
segmentTarget uint64 segmentTarget uint64
quotaLimits QuotaLimits
freeSpaceProbe FreeSpaceProbe
faultInjector FaultInjector faultInjector FaultInjector
activeSegments map[[16]byte]*activeSegment activeSegments map[[16]byte]*activeSegment
mu sync.Mutex mu sync.Mutex
@@ -52,6 +56,15 @@ func Open(ctx context.Context, options Options) (*Store, error) {
if options.SegmentTargetSize > DefaultSegmentLimit { if options.SegmentTargetSize > DefaultSegmentLimit {
return nil, fmt.Errorf("segment target exceeds hard record limit") 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 { if err := ensurePrivateDirectory(options.DataDir); err != nil {
return nil, err return nil, err
} }
@@ -88,6 +101,7 @@ func Open(ctx context.Context, options Options) (*Store, error) {
db.SetMaxIdleConns(1) db.SetMaxIdleConns(1)
store := &Store{ store := &Store{
db: db, unlock: unlock, dataDir: options.DataDir, segmentTarget: options.SegmentTargetSize, 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), faultInjector: options.FaultInjector, activeSegments: make(map[[16]byte]*activeSegment),
} }
if err := db.PingContext(ctx); err != nil { if err := db.PingContext(ctx); err != nil {
+42
View File
@@ -261,3 +261,45 @@ id = "CRASH-STORE-03"
layer = "integration" layer = "integration"
status = "implemented" status = "implemented"
tests = ["test/integration/store/store_integration_test.go:TestEventCrashOnExistingSegmentTruncatesOnlyTail_CRASH_STORE_03"] 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"]
@@ -66,7 +66,7 @@ func TestRealSQLiteInitializationAndRestart_HP_STORE_01(t *testing.T) {
if err := rows.Close(); err != nil { if err := rows.Close(); err != nil {
t.Fatal(err) 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) sort.Strings(want)
if strings.Join(names, ",") != strings.Join(want, ",") { if strings.Join(names, ",") != strings.Join(want, ",") {
t.Fatalf("tables = %v, want %v", names, 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 { func openStore(t *testing.T, dataDir string) *store.Store {
t.Helper() t.Helper()
return openStoreWithOptions(t, store.Options{DataDir: dataDir, BusyTimeout: busyTimeout}) return openStoreWithOptions(t, store.Options{DataDir: dataDir, BusyTimeout: busyTimeout})