From 3e25edd3374bb44b428e1ab28afd000cadafd7f9 Mon Sep 17 00:00:00 2001 From: cabbage Date: Mon, 31 Aug 2026 10:03:24 +0000 Subject: [PATCH] feat: track scoped storage incidents --- internal/server/store/health.go | 78 ++++++ internal/server/store/health_test.go | 52 ++++ internal/server/store/incidents.go | 246 ++++++++++++++++++ test/coverage.toml | 24 ++ .../store/store_integration_test.go | 104 ++++++++ 5 files changed, 504 insertions(+) create mode 100644 internal/server/store/health.go create mode 100644 internal/server/store/health_test.go create mode 100644 internal/server/store/incidents.go diff --git a/internal/server/store/health.go b/internal/server/store/health.go new file mode 100644 index 0000000..df3f0d8 --- /dev/null +++ b/internal/server/store/health.go @@ -0,0 +1,78 @@ +package store + +import "sync" + +type ScopeState uint8 + +const ( + ScopeRecovering ScopeState = iota + 1 + ScopeReady + ScopeDirtyReadable + ScopeUnavailable +) + +type ScopeKey struct { + Kind string + ID string +} + +type ScopeHealth struct { + mu sync.RWMutex + required map[ScopeKey]bool + states map[ScopeKey]ScopeState +} + +func NewScopeHealth(required []ScopeKey) *ScopeHealth { + health := &ScopeHealth{required: make(map[ScopeKey]bool, len(required)), states: make(map[ScopeKey]ScopeState, len(required))} + for _, scope := range required { + health.required[scope] = true + health.states[scope] = ScopeRecovering + } + return health +} + +func (health *ScopeHealth) Set(scope ScopeKey, state ScopeState) bool { + if health == nil || scope.Kind == "" || state < ScopeRecovering || state > ScopeUnavailable { + return false + } + health.mu.Lock() + defer health.mu.Unlock() + health.states[scope] = state + return true +} + +func (health *ScopeHealth) State(scope ScopeKey) ScopeState { + if health == nil { + return ScopeUnavailable + } + health.mu.RLock() + defer health.mu.RUnlock() + state := health.states[scope] + if state == 0 { + return ScopeUnavailable + } + return state +} + +func (health *ScopeHealth) GloballyReady() bool { + if health == nil { + return false + } + health.mu.RLock() + defer health.mu.RUnlock() + for scope := range health.required { + if health.states[scope] != ScopeReady { + return false + } + } + return true +} + +func (health *ScopeHealth) ReadAllowed(scope ScopeKey) bool { + state := health.State(scope) + return state == ScopeReady || state == ScopeDirtyReadable +} + +func (health *ScopeHealth) MutationAllowed(scope ScopeKey) bool { + return health.State(scope) == ScopeReady +} diff --git a/internal/server/store/health_test.go b/internal/server/store/health_test.go new file mode 100644 index 0000000..3420c0d --- /dev/null +++ b/internal/server/store/health_test.go @@ -0,0 +1,52 @@ +package store + +import ( + "sync" + "testing" +) + +func TestScopeHealthBestEffortReadiness_HP_STORE_10(t *testing.T) { + t.Parallel() + + global := ScopeKey{Kind: "global", ID: "store"} + client := ScopeKey{Kind: "client", ID: "client-a"} + health := NewScopeHealth([]ScopeKey{global, client}) + if health.GloballyReady() || health.ReadAllowed(global) || health.MutationAllowed(global) { + t.Fatal("recovering scope reported ready") + } + if !health.Set(global, ScopeReady) || !health.Set(client, ScopeDirtyReadable) { + t.Fatal("valid state update failed") + } + if health.GloballyReady() || !health.ReadAllowed(client) || health.MutationAllowed(client) { + t.Fatal("dirty-readable semantics are wrong") + } + if !health.Set(client, ScopeReady) || !health.GloballyReady() { + t.Fatal("resolved required scopes did not become globally ready") + } + if health.Set(ScopeKey{}, ScopeReady) || health.Set(global, 0) { + t.Fatal("invalid scope state accepted") + } +} + +func TestScopeHealthConcurrentUpdates_RACE_STORE_03(t *testing.T) { + t.Parallel() + + scope := ScopeKey{Kind: "client", ID: "client-a"} + health := NewScopeHealth([]ScopeKey{scope}) + var wait sync.WaitGroup + for index := range 32 { + wait.Add(1) + go func(index int) { + defer wait.Done() + states := []ScopeState{ScopeRecovering, ScopeReady, ScopeDirtyReadable, ScopeUnavailable} + health.Set(scope, states[index%len(states)]) + _ = health.ReadAllowed(scope) + _ = health.MutationAllowed(scope) + _ = health.GloballyReady() + }(index) + } + wait.Wait() + if state := health.State(scope); state < ScopeRecovering || state > ScopeUnavailable { + t.Fatalf("final state = %d", state) + } +} diff --git a/internal/server/store/incidents.go b/internal/server/store/incidents.go new file mode 100644 index 0000000..2eb6eef --- /dev/null +++ b/internal/server/store/incidents.go @@ -0,0 +1,246 @@ +package store + +import ( + "bytes" + "context" + "crypto/sha256" + "crypto/subtle" + "database/sql" + "encoding/hex" + "errors" + "fmt" + "time" +) + +type IncidentKind uint8 + +const ( + IncidentUncommittedTail IncidentKind = iota + 1 + IncidentMissingCommittedBytes + IncidentChecksumMismatch + IncidentCounterMismatch + IncidentSQLiteIntegrity + IncidentFailedEviction + IncidentPermission + IncidentDiskExhaustion +) + +type IncidentState uint8 + +const ( + IncidentOpen IncidentState = iota + 1 + IncidentRepaired + IncidentAcknowledged +) + +type IncidentScope string + +const ( + IncidentScopeGlobal IncidentScope = "global" + IncidentScopeClient IncidentScope = "client" + IncidentScopeCommand IncidentScope = "command" + IncidentScopeSegment IncidentScope = "segment" + IncidentScopeAudit IncidentScope = "audit" +) + +var ( + ErrIncidentNotFound = errors.New("storage incident not found") + ErrIncidentResolution = errors.New("invalid storage incident resolution") + ErrMutationConflict = errors.New("request ID conflicts with an earlier mutation") +) + +type IncidentInput struct { + IncidentUUID [16]byte + DetectedAt time.Time + Kind IncidentKind + Scope IncidentScope + ScopeKey string + ClientID string + IssueUUID *[16]byte + Summary string + Evidence []byte + DataLoss bool + AutomaticallyRepairable bool +} + +type Incident struct { + IncidentUUID [16]byte + State IncidentState + Created bool +} + +type IncidentResolution struct { + RequestUUID [16]byte + IncidentUUID [16]byte + State IncidentState + Note string + ResolvedAt time.Time +} + +func (store *Store) RecordIncident(ctx context.Context, input IncidentInput) (Incident, error) { + var result Incident + if err := validateIncidentInput(input); err != nil { + return result, err + } + store.writeMu.Lock() + defer store.writeMu.Unlock() + database, err := store.openDatabase() + if err != nil { + return result, err + } + var issue any + if input.IssueUUID != nil { + issue = input.IssueUUID[:] + } + dataLoss, repairable := 0, 0 + if input.DataLoss { + dataLoss = 1 + } + if input.AutomaticallyRepairable { + repairable = 1 + } + _, err = database.ExecContext(ctx, `INSERT INTO storage_incidents ( +incident_uuid, detected_at, state, kind, scope, scope_key, client_id, issue_uuid, summary, evidence, +data_loss, automatically_repairable +) VALUES (?, ?, 1, ?, ?, ?, NULLIF(?, ''), ?, ?, ?, ?, ?)`, input.IncidentUUID[:], input.DetectedAt.UnixNano(), input.Kind, + input.Scope, input.ScopeKey, input.ClientID, issue, input.Summary, input.Evidence, dataLoss, repairable) + if err == nil { + return Incident{IncidentUUID: input.IncidentUUID, State: IncidentOpen, Created: true}, nil + } + insertErr := err + var existingBytes []byte + err = database.QueryRowContext(ctx, `SELECT incident_uuid FROM storage_incidents +WHERE scope = ? AND scope_key = ? AND kind = ? AND state = 1`, input.Scope, input.ScopeKey, input.Kind).Scan(&existingBytes) + if err != nil { + return result, fmt.Errorf("record incident: %w", insertErr) + } + if len(existingBytes) != 16 { + return result, ErrIntegrityCheck + } + copy(result.IncidentUUID[:], existingBytes) + result.State = IncidentOpen + return result, nil +} + +func (store *Store) ResolveIncident(ctx context.Context, resolution IncidentResolution) (Incident, error) { + var result Incident + if isZeroUUID(resolution.RequestUUID) || isZeroUUID(resolution.IncidentUUID) || resolution.ResolvedAt.IsZero() || + (resolution.State != IncidentRepaired && resolution.State != IncidentAcknowledged) || len(resolution.Note) > 4096 || + (resolution.State == IncidentAcknowledged && resolution.Note == "") { + return result, ErrIncidentResolution + } + store.writeMu.Lock() + defer store.writeMu.Unlock() + database, err := store.openDatabase() + if err != nil { + return result, err + } + method := "repair_storage_incident" + if resolution.State == IncidentAcknowledged { + method = "acknowledge_storage_incident" + } + target := hex.EncodeToString(resolution.IncidentUUID[:]) + requestHash := incidentResolutionHash(resolution) + if replay, found, err := lookupIncidentResolution(ctx, database, resolution.RequestUUID, method, target, requestHash); err != nil { + return result, err + } else if found { + return replay, nil + } + + tx, err := database.BeginTx(ctx, nil) + if err != nil { + return result, err + } + var currentState IncidentState + var repairable, dataLoss int + err = tx.QueryRowContext(ctx, `SELECT state, automatically_repairable, data_loss FROM storage_incidents WHERE incident_uuid = ?`, resolution.IncidentUUID[:]).Scan(¤tState, &repairable, &dataLoss) + if errors.Is(err, sql.ErrNoRows) { + err = ErrIncidentNotFound + } + if err == nil && currentState != IncidentOpen && currentState != resolution.State { + err = ErrIncidentResolution + } + if err == nil && resolution.State == IncidentRepaired && (repairable != 1 || dataLoss != 0) { + err = ErrIncidentResolution + } + if err == nil && currentState == IncidentOpen { + _, err = tx.ExecContext(ctx, `UPDATE storage_incidents SET state = ?, resolved_at = ?, resolution_note = ? +WHERE incident_uuid = ? AND state = 1`, resolution.State, resolution.ResolvedAt.UnixNano(), resolution.Note, resolution.IncidentUUID[:]) + } + resultPayload := []byte{byte(resolution.State)} + if err == nil { + _, err = tx.ExecContext(ctx, `INSERT INTO control_mutations ( +request_uuid, method, owner_kind, owner_id, target, immutable_sha256, result, created_at +) VALUES (?, ?, 'incident', ?, ?, ?, ?, ?)`, resolution.RequestUUID[:], method, target, target, requestHash[:], resultPayload, resolution.ResolvedAt.UnixNano()) + } + if err == nil { + auditPayload := []byte(resolution.Note) + auditHash := sha256.Sum256(auditPayload) + _, err = tx.ExecContext(ctx, `INSERT INTO audit_events ( +occurred_at, source, principal, action, outcome, compression, payload, raw_bytes, stored_bytes, sha256 +) VALUES (?, 'store', NULL, ?, 'success', 1, ?, ?, ?, ?)`, resolution.ResolvedAt.UnixNano(), method, + auditPayload, len(auditPayload), len(auditPayload), auditHash[:]) + } + if err != nil { + _ = tx.Rollback() + return result, err + } + if err := tx.Commit(); err != nil { + return result, err + } + return Incident{IncidentUUID: resolution.IncidentUUID, State: resolution.State}, nil +} + +func (store *Store) HasDirtyIncidents(ctx context.Context) (bool, error) { + store.writeMu.Lock() + defer store.writeMu.Unlock() + database, err := store.openDatabase() + if err != nil { + return false, err + } + var exists int + if err := database.QueryRowContext(ctx, `SELECT EXISTS(SELECT 1 FROM storage_incidents WHERE state = 1)`).Scan(&exists); err != nil { + return false, err + } + return exists == 1, nil +} + +func validateIncidentInput(input IncidentInput) error { + validScope := input.Scope == IncidentScopeGlobal || input.Scope == IncidentScopeClient || input.Scope == IncidentScopeCommand || input.Scope == IncidentScopeSegment || input.Scope == IncidentScopeAudit + if isZeroUUID(input.IncidentUUID) || input.DetectedAt.IsZero() || input.Kind < IncidentUncommittedTail || input.Kind > IncidentDiskExhaustion || + !validScope || input.ScopeKey == "" || len(input.ScopeKey) > 512 || input.Summary == "" || len(input.Summary) > 4096 || len(input.Evidence) > 16<<10 || + (input.DataLoss && input.AutomaticallyRepairable) { + return errors.New("invalid storage incident") + } + return nil +} + +func incidentResolutionHash(resolution IncidentResolution) [sha256.Size]byte { + buffer := bytes.NewBuffer(make([]byte, 0, 33+len(resolution.Note))) + buffer.Write(resolution.IncidentUUID[:]) + buffer.WriteByte(byte(resolution.State)) + buffer.WriteString(resolution.Note) + return sha256.Sum256(buffer.Bytes()) +} + +func lookupIncidentResolution(ctx context.Context, database *sql.DB, requestUUID [16]byte, method, target string, requestHash [32]byte) (Incident, bool, error) { + var storedMethod, storedTarget string + var storedHash, storedResult []byte + err := database.QueryRowContext(ctx, `SELECT method, target, immutable_sha256, result FROM control_mutations WHERE request_uuid = ?`, requestUUID[:]).Scan(&storedMethod, &storedTarget, &storedHash, &storedResult) + if errors.Is(err, sql.ErrNoRows) { + return Incident{}, false, nil + } + if err != nil { + return Incident{}, false, err + } + if storedMethod != method || storedTarget != target || len(storedHash) != 32 || subtle.ConstantTimeCompare(storedHash, requestHash[:]) != 1 || len(storedResult) != 1 { + return Incident{}, false, ErrMutationConflict + } + decoded, err := hex.DecodeString(target) + if err != nil || len(decoded) != 16 { + return Incident{}, false, fmt.Errorf("%w: stored target", ErrIntegrityCheck) + } + var incidentUUID [16]byte + copy(incidentUUID[:], decoded) + return Incident{IncidentUUID: incidentUUID, State: IncidentState(storedResult[0])}, true, nil +} diff --git a/test/coverage.toml b/test/coverage.toml index d379e55..03b5f71 100644 --- a/test/coverage.toml +++ b/test/coverage.toml @@ -321,3 +321,27 @@ id = "CRASH-STORE-04" layer = "integration" status = "implemented" tests = ["test/integration/store/store_integration_test.go:TestRetentionCrashStagesRollForward_CRASH_STORE_04"] + +[[requirements]] +id = "HP-STORE-10" +layer = "unit" +status = "implemented" +tests = ["internal/server/store/health_test.go:TestScopeHealthBestEffortReadiness_HP_STORE_10"] + +[[requirements]] +id = "RACE-STORE-03" +layer = "unit" +status = "implemented" +tests = ["internal/server/store/health_test.go:TestScopeHealthConcurrentUpdates_RACE_STORE_03"] + +[[requirements]] +id = "HP-STORE-11" +layer = "integration" +status = "implemented" +tests = ["test/integration/store/store_integration_test.go:TestIncidentDirtyResolutionIdempotencyAndRecurrence_HP_STORE_11"] + +[[requirements]] +id = "BH-STORE-09" +layer = "integration" +status = "implemented" +tests = ["test/integration/store/store_integration_test.go:TestIrreparableIncidentRequiresExplicitAcknowledgement_BH_STORE_09"] diff --git a/test/integration/store/store_integration_test.go b/test/integration/store/store_integration_test.go index 53e242d..c2427cc 100644 --- a/test/integration/store/store_integration_test.go +++ b/test/integration/store/store_integration_test.go @@ -779,6 +779,110 @@ func TestRetentionCrashStagesRollForward_CRASH_STORE_04(t *testing.T) { } } +func TestIncidentDirtyResolutionIdempotencyAndRecurrence_HP_STORE_11(t *testing.T) { + t.Parallel() + + opened := openStore(t, filepath.Join(t.TempDir(), "state")) + now := time.Unix(5_000_000, 0).UTC() + firstID := uuidBytes(10) + input := store.IncidentInput{ + IncidentUUID: firstID, DetectedAt: now, Kind: store.IncidentCounterMismatch, + Scope: store.IncidentScopeClient, ScopeKey: "client-a", ClientID: "client-a", + Summary: "charged totals disagree", Evidence: []byte("counter-version=1"), AutomaticallyRepairable: true, + } + created, err := opened.RecordIncident(context.Background(), input) + if err != nil || !created.Created || created.IncidentUUID != firstID { + t.Fatalf("RecordIncident = (%+v, %v)", created, err) + } + duplicateInput := input + duplicateInput.IncidentUUID = uuidBytes(30) + duplicate, err := opened.RecordIncident(context.Background(), duplicateInput) + if err != nil || duplicate.Created || duplicate.IncidentUUID != firstID { + t.Fatalf("duplicate RecordIncident = (%+v, %v)", duplicate, err) + } + dirty, err := opened.HasDirtyIncidents(context.Background()) + if err != nil || !dirty { + t.Fatalf("dirty = %v, err = %v", dirty, err) + } + requestID := uuidBytes(50) + resolution := store.IncidentResolution{ + RequestUUID: requestID, IncidentUUID: firstID, State: store.IncidentRepaired, + Note: "recomputed counters from command rows", ResolvedAt: now.Add(time.Minute), + } + resolved, err := opened.ResolveIncident(context.Background(), resolution) + if err != nil || resolved.State != store.IncidentRepaired { + t.Fatalf("ResolveIncident = (%+v, %v)", resolved, err) + } + replayed, err := opened.ResolveIncident(context.Background(), resolution) + if err != nil || replayed != resolved { + t.Fatalf("resolution replay = (%+v, %v)", replayed, err) + } + conflict := resolution + conflict.Note = "different repair claim" + if _, err := opened.ResolveIncident(context.Background(), conflict); !errors.Is(err, store.ErrMutationConflict) { + t.Fatalf("resolution conflict error = %v", err) + } + dirty, err = opened.HasDirtyIncidents(context.Background()) + if err != nil || dirty { + t.Fatalf("resolved dirty = %v, err = %v", dirty, err) + } + var auditCount int + if err := opened.DB().QueryRow(`SELECT count(*) FROM audit_events WHERE action = 'repair_storage_incident'`).Scan(&auditCount); err != nil || auditCount != 1 { + t.Fatalf("repair audit count = %d, err = %v", auditCount, err) + } + recurrence := input + recurrence.IncidentUUID = uuidBytes(70) + recurrence.DetectedAt = now.Add(2 * time.Minute) + recorded, err := opened.RecordIncident(context.Background(), recurrence) + if err != nil || !recorded.Created || recorded.IncidentUUID != recurrence.IncidentUUID { + t.Fatalf("incident recurrence = (%+v, %v)", recorded, err) + } + if err := opened.Close(); err != nil { + t.Fatal(err) + } +} + +func TestIrreparableIncidentRequiresExplicitAcknowledgement_BH_STORE_09(t *testing.T) { + t.Parallel() + + opened := openStore(t, filepath.Join(t.TempDir(), "state")) + now := time.Unix(6_000_000, 0).UTC() + incidentID := uuidBytes(90) + if _, err := opened.RecordIncident(context.Background(), store.IncidentInput{ + IncidentUUID: incidentID, DetectedAt: now, Kind: store.IncidentMissingCommittedBytes, + Scope: store.IncidentScopeSegment, ScopeKey: "segment-redacted", Summary: "committed bytes are missing", + Evidence: []byte("offset=128"), DataLoss: true, + }); err != nil { + t.Fatal(err) + } + if _, err := opened.ResolveIncident(context.Background(), store.IncidentResolution{ + RequestUUID: uuidBytes(110), IncidentUUID: incidentID, State: store.IncidentRepaired, + Note: "cannot really repair", ResolvedAt: now.Add(time.Minute), + }); !errors.Is(err, store.ErrIncidentResolution) { + t.Fatalf("unsafe repair error = %v", err) + } + if _, err := opened.ResolveIncident(context.Background(), store.IncidentResolution{ + RequestUUID: uuidBytes(130), IncidentUUID: incidentID, State: store.IncidentAcknowledged, + ResolvedAt: now.Add(2 * time.Minute), + }); !errors.Is(err, store.ErrIncidentResolution) { + t.Fatalf("empty acknowledgement error = %v", err) + } + acknowledged, err := opened.ResolveIncident(context.Background(), store.IncidentResolution{ + RequestUUID: uuidBytes(150), IncidentUUID: incidentID, State: store.IncidentAcknowledged, + Note: "operator accepts loss after external verification", ResolvedAt: now.Add(3 * time.Minute), + }) + if err != nil || acknowledged.State != store.IncidentAcknowledged { + t.Fatalf("acknowledgement = (%+v, %v)", acknowledged, err) + } + dirty, err := opened.HasDirtyIncidents(context.Background()) + if err != nil || dirty { + t.Fatalf("acknowledged dirty = %v, err = %v", dirty, err) + } + if err := opened.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})