feat: track scoped storage incidents
This commit is contained in:
@@ -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
|
||||||
|
}
|
||||||
@@ -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)
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -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
|
||||||
|
}
|
||||||
@@ -321,3 +321,27 @@ id = "CRASH-STORE-04"
|
|||||||
layer = "integration"
|
layer = "integration"
|
||||||
status = "implemented"
|
status = "implemented"
|
||||||
tests = ["test/integration/store/store_integration_test.go:TestRetentionCrashStagesRollForward_CRASH_STORE_04"]
|
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"]
|
||||||
|
|||||||
@@ -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 {
|
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})
|
||||||
|
|||||||
Reference in New Issue
Block a user