mirror of
https://github.com/rcourtman/Pulse.git
synced 2026-09-11 14:00:29 +00:00
Make maintenance evidence writes atomic
This commit is contained in:
@@ -96,11 +96,6 @@ func (h *ResourceHandlers) HandleResourceOperatorState(w http.ResponseWriter, r
|
||||
http.Error(w, "Invalid JSON", http.StatusBadRequest)
|
||||
return
|
||||
}
|
||||
previous, previousFound, err := store.GetResourceOperatorState(resourceID)
|
||||
if err != nil {
|
||||
http.Error(w, sanitizeErrorForClient(err, "Internal server error"), http.StatusInternalServerError)
|
||||
return
|
||||
}
|
||||
// canonical_id from URL wins over body to prevent scope confusion
|
||||
// (operator wrote vm:101 in the URL but vm:102 in the body).
|
||||
state := unified.ResourceOperatorState{
|
||||
@@ -117,7 +112,8 @@ func (h *ResourceHandlers) HandleResourceOperatorState(w http.ResponseWriter, r
|
||||
SetAt: time.Now().UTC(),
|
||||
SetBy: getUserID(r),
|
||||
}
|
||||
if err := store.SetResourceOperatorState(state); err != nil {
|
||||
persisted, err := unified.SetResourceOperatorStateWithMaintenanceLifecycle(store, state)
|
||||
if err != nil {
|
||||
if errors.Is(err, unified.ErrResourceOperatorStateInvalid) {
|
||||
writeJSONError(w, http.StatusBadRequest, "operator_state_invalid", err.Error())
|
||||
return
|
||||
@@ -125,41 +121,12 @@ func (h *ResourceHandlers) HandleResourceOperatorState(w http.ResponseWriter, r
|
||||
http.Error(w, sanitizeErrorForClient(err, "Internal server error"), http.StatusInternalServerError)
|
||||
return
|
||||
}
|
||||
// Read-after-write: return the persisted state so the caller can
|
||||
// see exactly what the server stored, including the
|
||||
// server-populated attribution fields.
|
||||
persisted, _, err := store.GetResourceOperatorState(resourceID)
|
||||
if err != nil {
|
||||
http.Error(w, sanitizeErrorForClient(err, "Internal server error"), http.StatusInternalServerError)
|
||||
return
|
||||
}
|
||||
if err := recordMaintenanceWindowLifecycleChange(store, previous, previousFound, persisted, true, persisted.SetAt, persisted.SetBy); err != nil {
|
||||
http.Error(w, sanitizeErrorForClient(err, "Internal server error"), http.StatusInternalServerError)
|
||||
return
|
||||
}
|
||||
writeJSON(w, http.StatusOK, toResourceOperatorStateAPI(persisted))
|
||||
|
||||
case http.MethodDelete:
|
||||
previous, previousFound, err := store.GetResourceOperatorState(resourceID)
|
||||
if err != nil {
|
||||
http.Error(w, sanitizeErrorForClient(err, "Internal server error"), http.StatusInternalServerError)
|
||||
return
|
||||
}
|
||||
observedAt := time.Now().UTC()
|
||||
actor := getUserID(r)
|
||||
if err := store.ClearResourceOperatorState(resourceID); err != nil {
|
||||
http.Error(w, sanitizeErrorForClient(err, "Internal server error"), http.StatusInternalServerError)
|
||||
return
|
||||
}
|
||||
if err := recordMaintenanceWindowLifecycleChange(
|
||||
store,
|
||||
previous,
|
||||
previousFound,
|
||||
unified.ResourceOperatorState{CanonicalID: resourceID},
|
||||
false,
|
||||
observedAt,
|
||||
actor,
|
||||
); err != nil {
|
||||
if err := unified.ClearResourceOperatorStateWithMaintenanceLifecycle(store, resourceID, observedAt, actor); err != nil {
|
||||
http.Error(w, sanitizeErrorForClient(err, "Internal server error"), http.StatusInternalServerError)
|
||||
return
|
||||
}
|
||||
@@ -182,19 +149,3 @@ func extractOperatorStateResourceID(path string) string {
|
||||
trimmed = strings.TrimSuffix(trimmed, "/")
|
||||
return unified.CanonicalResourceID(trimmed)
|
||||
}
|
||||
|
||||
func recordMaintenanceWindowLifecycleChange(
|
||||
store unified.ResourceStore,
|
||||
previous unified.ResourceOperatorState,
|
||||
previousFound bool,
|
||||
current unified.ResourceOperatorState,
|
||||
currentFound bool,
|
||||
observedAt time.Time,
|
||||
actor string,
|
||||
) error {
|
||||
change, ok := unified.BuildMaintenanceWindowLifecycleChange(previous, previousFound, current, currentFound, observedAt, actor)
|
||||
if !ok {
|
||||
return nil
|
||||
}
|
||||
return store.RecordChange(change)
|
||||
}
|
||||
|
||||
@@ -2,9 +2,11 @@ package api
|
||||
|
||||
import (
|
||||
"bytes"
|
||||
"database/sql"
|
||||
"encoding/json"
|
||||
"net/http"
|
||||
"net/http/httptest"
|
||||
"path/filepath"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
@@ -17,8 +19,27 @@ import (
|
||||
// existing resources_test.go fixtures.
|
||||
func newOperatorStateHandlers(t *testing.T) *ResourceHandlers {
|
||||
t.Helper()
|
||||
cfg := &config.Config{DataPath: t.TempDir()}
|
||||
return NewResourceHandlers(cfg)
|
||||
h, _ := newOperatorStateHandlersWithDataDir(t)
|
||||
return h
|
||||
}
|
||||
|
||||
func newOperatorStateHandlersWithDataDir(t *testing.T) (*ResourceHandlers, string) {
|
||||
t.Helper()
|
||||
dataDir := t.TempDir()
|
||||
cfg := &config.Config{DataPath: dataDir}
|
||||
return NewResourceHandlers(cfg), dataDir
|
||||
}
|
||||
|
||||
func dropOperatorStateTimelineTable(t *testing.T, dataDir string) {
|
||||
t.Helper()
|
||||
db, err := sql.Open("sqlite", filepath.Join(dataDir, "resources", "unified_resources.db"))
|
||||
if err != nil {
|
||||
t.Fatalf("open resource db for failure injection: %v", err)
|
||||
}
|
||||
defer db.Close()
|
||||
if _, err := db.Exec(`DROP TABLE resource_changes`); err != nil {
|
||||
t.Fatalf("drop resource_changes for failure injection: %v", err)
|
||||
}
|
||||
}
|
||||
|
||||
func TestHandleResourceOperatorState_GetReturns404WhenUnset(t *testing.T) {
|
||||
@@ -109,7 +130,7 @@ func TestHandleResourceOperatorState_PutPersistsAndGetReturns200(t *testing.T) {
|
||||
}
|
||||
}
|
||||
|
||||
func TestHandleResourceOperatorState_RecordsMaintenanceWindowLifecycleChanges(t *testing.T) {
|
||||
func TestResourceOperatorState_RecordsMaintenanceWindowLifecycleChanges(t *testing.T) {
|
||||
h := newOperatorStateHandlers(t)
|
||||
store, err := h.getStore("default")
|
||||
if err != nil {
|
||||
@@ -178,6 +199,81 @@ func TestHandleResourceOperatorState_RecordsMaintenanceWindowLifecycleChanges(t
|
||||
}
|
||||
}
|
||||
|
||||
func TestResourceOperatorState_PutRollsBackWhenTimelineProjectionFails(t *testing.T) {
|
||||
h, dataDir := newOperatorStateHandlersWithDataDir(t)
|
||||
store, err := h.getStore("default")
|
||||
if err != nil {
|
||||
t.Fatalf("get store: %v", err)
|
||||
}
|
||||
dropOperatorStateTimelineTable(t, dataDir)
|
||||
|
||||
start := time.Date(2026, 5, 9, 12, 0, 0, 0, time.UTC)
|
||||
end := time.Date(2026, 5, 9, 14, 0, 0, 0, time.UTC)
|
||||
body, _ := json.Marshal(map[string]any{
|
||||
"maintenanceStartAt": start.Format(time.RFC3339),
|
||||
"maintenanceEndAt": end.Format(time.RFC3339),
|
||||
"maintenanceReason": "storage controller patch",
|
||||
})
|
||||
|
||||
rec := httptest.NewRecorder()
|
||||
req := httptest.NewRequest(http.MethodPut, "/api/resources/vm:101/operator-state", bytes.NewReader(body))
|
||||
h.HandleResourceOperatorState(rec, req)
|
||||
|
||||
if rec.Code != http.StatusInternalServerError {
|
||||
t.Fatalf("expected 500 on timeline projection failure; got %d body=%s", rec.Code, rec.Body.String())
|
||||
}
|
||||
if _, found, err := store.GetResourceOperatorState("vm:101"); err != nil {
|
||||
t.Fatalf("get operator state after rollback: %v", err)
|
||||
} else if found {
|
||||
t.Fatal("operator state source row persisted despite projection failure")
|
||||
}
|
||||
}
|
||||
|
||||
func TestResourceOperatorState_DeleteRollsBackWhenTimelineProjectionFails(t *testing.T) {
|
||||
h, dataDir := newOperatorStateHandlersWithDataDir(t)
|
||||
store, err := h.getStore("default")
|
||||
if err != nil {
|
||||
t.Fatalf("get store: %v", err)
|
||||
}
|
||||
|
||||
start := time.Date(2026, 5, 9, 12, 0, 0, 0, time.UTC)
|
||||
end := time.Date(2026, 5, 9, 14, 0, 0, 0, time.UTC)
|
||||
seed := unified.ResourceOperatorState{
|
||||
CanonicalID: "vm:101",
|
||||
MaintenanceStartAt: &start,
|
||||
MaintenanceEndAt: &end,
|
||||
MaintenanceReason: "storage controller patch",
|
||||
IntentionallyOffline: true,
|
||||
SetAt: start.Add(-time.Hour),
|
||||
SetBy: "operator",
|
||||
}
|
||||
if err := store.SetResourceOperatorState(seed); err != nil {
|
||||
t.Fatalf("seed operator state: %v", err)
|
||||
}
|
||||
dropOperatorStateTimelineTable(t, dataDir)
|
||||
|
||||
rec := httptest.NewRecorder()
|
||||
req := httptest.NewRequest(http.MethodDelete, "/api/resources/vm:101/operator-state", nil)
|
||||
h.HandleResourceOperatorState(rec, req)
|
||||
|
||||
if rec.Code != http.StatusInternalServerError {
|
||||
t.Fatalf("expected 500 on timeline projection failure; got %d body=%s", rec.Code, rec.Body.String())
|
||||
}
|
||||
got, found, err := store.GetResourceOperatorState("vm:101")
|
||||
if err != nil {
|
||||
t.Fatalf("get operator state after rollback: %v", err)
|
||||
}
|
||||
if !found {
|
||||
t.Fatal("operator state source row was deleted despite projection failure")
|
||||
}
|
||||
if got.MaintenanceStartAt == nil || !got.MaintenanceStartAt.Equal(start) {
|
||||
t.Fatalf("maintenance start after rollback = %v want %v", got.MaintenanceStartAt, start)
|
||||
}
|
||||
if got.MaintenanceEndAt == nil || !got.MaintenanceEndAt.Equal(end) {
|
||||
t.Fatalf("maintenance end after rollback = %v want %v", got.MaintenanceEndAt, end)
|
||||
}
|
||||
}
|
||||
|
||||
func TestHandleResourceOperatorState_PutRejectsInvalidWith400(t *testing.T) {
|
||||
h := newOperatorStateHandlers(t)
|
||||
|
||||
|
||||
@@ -137,7 +137,20 @@ func (s *SQLiteResourceStore) RecordLoopReport(report LoopReport) error {
|
||||
}
|
||||
|
||||
s.mu.Lock()
|
||||
_, err = s.db.Exec(`
|
||||
defer s.mu.Unlock()
|
||||
|
||||
tx, err := s.db.Begin()
|
||||
if err != nil {
|
||||
return fmt.Errorf("begin loop report transaction: %w", err)
|
||||
}
|
||||
committed := false
|
||||
defer func() {
|
||||
if !committed {
|
||||
_ = tx.Rollback()
|
||||
}
|
||||
}()
|
||||
|
||||
_, err = tx.Exec(`
|
||||
INSERT INTO loop_reports (
|
||||
id, report_type, scope, trigger, goal, status, started_at, completed_at,
|
||||
window_started_at, window_ended_at, evidence_json,
|
||||
@@ -166,15 +179,18 @@ func (s *SQLiteResourceStore) RecordLoopReport(report LoopReport) error {
|
||||
report.ReviewedBy,
|
||||
report.ReviewNote,
|
||||
)
|
||||
s.mu.Unlock()
|
||||
if err != nil {
|
||||
return fmt.Errorf("insert loop report: %w", err)
|
||||
}
|
||||
if change, ok := BuildLoopReportResourceChange(report); ok {
|
||||
if err := s.RecordChange(change); err != nil {
|
||||
if err := recordChangeSQL(tx, change, s.resourceChangesHasTimestamp); err != nil {
|
||||
return fmt.Errorf("record loop report resource change: %w", err)
|
||||
}
|
||||
}
|
||||
if err := tx.Commit(); err != nil {
|
||||
return fmt.Errorf("commit loop report transaction: %w", err)
|
||||
}
|
||||
committed = true
|
||||
return nil
|
||||
}
|
||||
|
||||
@@ -422,12 +438,13 @@ func (m *MemoryStore) RecordLoopReport(report LoopReport) error {
|
||||
return fmt.Errorf("%w: id %q already exists", ErrLoopReportInvalid, report.ID)
|
||||
}
|
||||
m.loopReports[report.ID] = report
|
||||
m.mu.Unlock()
|
||||
if change, ok := BuildLoopReportResourceChange(report); ok {
|
||||
if err := m.RecordChange(change); err != nil {
|
||||
if err := m.recordChangeLocked(change); err != nil {
|
||||
m.mu.Unlock()
|
||||
return fmt.Errorf("record loop report resource change: %w", err)
|
||||
}
|
||||
}
|
||||
m.mu.Unlock()
|
||||
return nil
|
||||
}
|
||||
|
||||
|
||||
@@ -1,6 +1,7 @@
|
||||
package unifiedresources
|
||||
|
||||
import (
|
||||
"strings"
|
||||
"testing"
|
||||
"time"
|
||||
)
|
||||
@@ -97,6 +98,35 @@ func TestLoopReport_RecordMaintenanceVerificationResourceChange(t *testing.T) {
|
||||
}
|
||||
}
|
||||
|
||||
func TestSQLiteRecordLoopReport_RollsBackWhenTimelineProjectionFails(t *testing.T) {
|
||||
dataDir := t.TempDir()
|
||||
store, err := NewSQLiteResourceStore(dataDir, "default")
|
||||
if err != nil {
|
||||
t.Fatalf("new store: %v", err)
|
||||
}
|
||||
defer store.Close()
|
||||
|
||||
if _, err := store.db.Exec(`DROP TABLE resource_changes`); err != nil {
|
||||
t.Fatalf("drop resource_changes: %v", err)
|
||||
}
|
||||
|
||||
windowEnd := time.Date(2026, 5, 12, 12, 0, 0, 0, time.UTC)
|
||||
report := newLoopReport("mv-rollback", "vm:101", windowEnd, LoopReportStatusHealthy)
|
||||
err = store.RecordLoopReport(report)
|
||||
if err == nil {
|
||||
t.Fatal("expected timeline projection failure")
|
||||
}
|
||||
if !strings.Contains(err.Error(), "record loop report resource change") {
|
||||
t.Fatalf("error = %v, want resource change failure", err)
|
||||
}
|
||||
|
||||
if _, found, err := store.GetLoopReport(report.ID); err != nil {
|
||||
t.Fatalf("get loop report after rollback: %v", err)
|
||||
} else if found {
|
||||
t.Fatal("loop report source row persisted despite projection failure")
|
||||
}
|
||||
}
|
||||
|
||||
func TestSQLiteRecordLoopReport_AllowsRerunForSameWindow(t *testing.T) {
|
||||
dataDir := t.TempDir()
|
||||
store, err := NewSQLiteResourceStore(dataDir, "default")
|
||||
|
||||
@@ -176,6 +176,31 @@ func NormalizeResourceOperatorState(state ResourceOperatorState) ResourceOperato
|
||||
return state
|
||||
}
|
||||
|
||||
type resourceOperatorStateLifecycleStore interface {
|
||||
SetResourceOperatorStateWithMaintenanceLifecycle(state ResourceOperatorState) (ResourceOperatorState, error)
|
||||
ClearResourceOperatorStateWithMaintenanceLifecycle(canonicalID string, observedAt time.Time, actor string) error
|
||||
}
|
||||
|
||||
// SetResourceOperatorStateWithMaintenanceLifecycle persists operator state and
|
||||
// any derived maintenance-window timeline change through one store-owned write.
|
||||
func SetResourceOperatorStateWithMaintenanceLifecycle(store ResourceStore, state ResourceOperatorState) (ResourceOperatorState, error) {
|
||||
lifecycleStore, ok := store.(resourceOperatorStateLifecycleStore)
|
||||
if !ok {
|
||||
return ResourceOperatorState{}, errors.New("resource operator state maintenance lifecycle projection requires atomic store support")
|
||||
}
|
||||
return lifecycleStore.SetResourceOperatorStateWithMaintenanceLifecycle(state)
|
||||
}
|
||||
|
||||
// ClearResourceOperatorStateWithMaintenanceLifecycle clears operator state and
|
||||
// any derived maintenance-window timeline change through one store-owned write.
|
||||
func ClearResourceOperatorStateWithMaintenanceLifecycle(store ResourceStore, canonicalID string, observedAt time.Time, actor string) error {
|
||||
lifecycleStore, ok := store.(resourceOperatorStateLifecycleStore)
|
||||
if !ok {
|
||||
return errors.New("resource operator state maintenance lifecycle projection requires atomic store support")
|
||||
}
|
||||
return lifecycleStore.ClearResourceOperatorStateWithMaintenanceLifecycle(canonicalID, observedAt, actor)
|
||||
}
|
||||
|
||||
const (
|
||||
MaintenanceWindowLifecycleEventScheduled = "maintenance_window_scheduled"
|
||||
MaintenanceWindowLifecycleEventUpdated = "maintenance_window_updated"
|
||||
|
||||
@@ -727,12 +727,16 @@ func (s *SQLiteResourceStore) RecordChange(change ResourceChange) error {
|
||||
s.mu.Lock()
|
||||
defer s.mu.Unlock()
|
||||
|
||||
return recordChangeSQL(s.db, change, s.resourceChangesHasTimestamp)
|
||||
}
|
||||
|
||||
func recordChangeSQL(execer sqlExecutor, change ResourceChange, includeTimestamp bool) error {
|
||||
relJSON, _ := json.Marshal(change.RelatedResources)
|
||||
metaJSON, _ := json.Marshal(change.Metadata)
|
||||
|
||||
columns := []string{"id", "canonical_id", "observed_at"}
|
||||
values := []any{change.ID, CanonicalResourceID(change.ResourceID), change.ObservedAt}
|
||||
if s.resourceChangesHasTimestamp {
|
||||
if includeTimestamp {
|
||||
columns = append(columns, "timestamp")
|
||||
values = append(values, change.ObservedAt)
|
||||
}
|
||||
@@ -768,7 +772,7 @@ func (s *SQLiteResourceStore) RecordChange(change ResourceChange) error {
|
||||
placeholders[i] = "?"
|
||||
}
|
||||
|
||||
_, err := s.db.Exec(
|
||||
_, err := execer.Exec(
|
||||
`INSERT INTO resource_changes (`+strings.Join(columns, ", ")+`) VALUES (`+strings.Join(placeholders, ", ")+`) ON CONFLICT(id) DO NOTHING`,
|
||||
values...,
|
||||
)
|
||||
@@ -1480,16 +1484,38 @@ func (s *SQLiteResourceStore) GetExportAudits(since time.Time, limit int) ([]Exp
|
||||
// signal is meaningful for the API GET path which returns 404 vs
|
||||
// returning the default no-state record.
|
||||
func (s *SQLiteResourceStore) GetResourceOperatorState(canonicalID string) (ResourceOperatorState, bool, error) {
|
||||
return getResourceOperatorStateSQL(s.db, canonicalID)
|
||||
}
|
||||
|
||||
type resourceOperatorStateQueryRower interface {
|
||||
QueryRow(query string, args ...any) *sql.Row
|
||||
}
|
||||
|
||||
type resourceOperatorStateScanner interface {
|
||||
Scan(dest ...any) error
|
||||
}
|
||||
|
||||
func getResourceOperatorStateSQL(queryer resourceOperatorStateQueryRower, canonicalID string) (ResourceOperatorState, bool, error) {
|
||||
canonicalID = strings.TrimSpace(canonicalID)
|
||||
if canonicalID == "" {
|
||||
return ResourceOperatorState{}, false, nil
|
||||
}
|
||||
row := s.db.QueryRow(`
|
||||
row := queryer.QueryRow(`
|
||||
SELECT canonical_id, intentionally_offline, never_auto_remediate,
|
||||
maintenance_start_at, maintenance_end_at, maintenance_reason,
|
||||
criticality, note, set_at, set_by
|
||||
FROM resource_operator_state WHERE canonical_id = ?`, canonicalID)
|
||||
state, err := scanResourceOperatorState(row)
|
||||
if err != nil {
|
||||
if errors.Is(err, sql.ErrNoRows) {
|
||||
return ResourceOperatorState{}, false, nil
|
||||
}
|
||||
return ResourceOperatorState{}, false, fmt.Errorf("query resource operator state: %w", err)
|
||||
}
|
||||
return state, true, nil
|
||||
}
|
||||
|
||||
func scanResourceOperatorState(scanner resourceOperatorStateScanner) (ResourceOperatorState, error) {
|
||||
var state ResourceOperatorState
|
||||
var (
|
||||
intentional int
|
||||
@@ -1500,7 +1526,7 @@ func (s *SQLiteResourceStore) GetResourceOperatorState(canonicalID string) (Reso
|
||||
note sql.NullString
|
||||
setBy sql.NullString
|
||||
)
|
||||
if err := row.Scan(
|
||||
if err := scanner.Scan(
|
||||
&state.CanonicalID,
|
||||
&intentional,
|
||||
&neverRemediate,
|
||||
@@ -1512,10 +1538,7 @@ func (s *SQLiteResourceStore) GetResourceOperatorState(canonicalID string) (Reso
|
||||
&state.SetAt,
|
||||
&setBy,
|
||||
); err != nil {
|
||||
if errors.Is(err, sql.ErrNoRows) {
|
||||
return ResourceOperatorState{}, false, nil
|
||||
}
|
||||
return ResourceOperatorState{}, false, fmt.Errorf("query resource operator state: %w", err)
|
||||
return ResourceOperatorState{}, err
|
||||
}
|
||||
state.IntentionallyOffline = intentional != 0
|
||||
state.NeverAutoRemediate = neverRemediate != 0
|
||||
@@ -1539,7 +1562,7 @@ func (s *SQLiteResourceStore) GetResourceOperatorState(canonicalID string) (Reso
|
||||
if setBy.Valid {
|
||||
state.SetBy = setBy.String
|
||||
}
|
||||
return state, true, nil
|
||||
return state, nil
|
||||
}
|
||||
|
||||
// SetResourceOperatorState upserts the state row. Validates and
|
||||
@@ -1552,6 +1575,14 @@ func (s *SQLiteResourceStore) SetResourceOperatorState(state ResourceOperatorSta
|
||||
if err := ValidateResourceOperatorState(state); err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
s.mu.Lock()
|
||||
defer s.mu.Unlock()
|
||||
|
||||
return setResourceOperatorStateSQL(s.db, state)
|
||||
}
|
||||
|
||||
func setResourceOperatorStateSQL(execer sqlExecutor, state ResourceOperatorState) error {
|
||||
intentional := 0
|
||||
if state.IntentionallyOffline {
|
||||
intentional = 1
|
||||
@@ -1591,7 +1622,7 @@ func (s *SQLiteResourceStore) SetResourceOperatorState(state ResourceOperatorSta
|
||||
setBy.String = state.SetBy
|
||||
setBy.Valid = true
|
||||
}
|
||||
_, err := s.db.Exec(`
|
||||
_, err := execer.Exec(`
|
||||
INSERT INTO resource_operator_state (
|
||||
canonical_id, intentionally_offline, never_auto_remediate,
|
||||
maintenance_start_at, maintenance_end_at, maintenance_reason,
|
||||
@@ -1624,6 +1655,48 @@ func (s *SQLiteResourceStore) SetResourceOperatorState(state ResourceOperatorSta
|
||||
return nil
|
||||
}
|
||||
|
||||
// SetResourceOperatorStateWithMaintenanceLifecycle persists the operator-state
|
||||
// source row and the derived maintenance-window timeline projection in a single
|
||||
// SQLite transaction.
|
||||
func (s *SQLiteResourceStore) SetResourceOperatorStateWithMaintenanceLifecycle(state ResourceOperatorState) (ResourceOperatorState, error) {
|
||||
state = NormalizeResourceOperatorState(state)
|
||||
if err := ValidateResourceOperatorState(state); err != nil {
|
||||
return ResourceOperatorState{}, err
|
||||
}
|
||||
|
||||
s.mu.Lock()
|
||||
defer s.mu.Unlock()
|
||||
|
||||
tx, err := s.db.Begin()
|
||||
if err != nil {
|
||||
return ResourceOperatorState{}, fmt.Errorf("begin resource operator state transaction: %w", err)
|
||||
}
|
||||
committed := false
|
||||
defer func() {
|
||||
if !committed {
|
||||
_ = tx.Rollback()
|
||||
}
|
||||
}()
|
||||
|
||||
previous, previousFound, err := getResourceOperatorStateSQL(tx, state.CanonicalID)
|
||||
if err != nil {
|
||||
return ResourceOperatorState{}, err
|
||||
}
|
||||
if err := setResourceOperatorStateSQL(tx, state); err != nil {
|
||||
return ResourceOperatorState{}, err
|
||||
}
|
||||
if change, ok := BuildMaintenanceWindowLifecycleChange(previous, previousFound, state, true, state.SetAt, state.SetBy); ok {
|
||||
if err := recordChangeSQL(tx, change, s.resourceChangesHasTimestamp); err != nil {
|
||||
return ResourceOperatorState{}, fmt.Errorf("record maintenance window lifecycle change: %w", err)
|
||||
}
|
||||
}
|
||||
if err := tx.Commit(); err != nil {
|
||||
return ResourceOperatorState{}, fmt.Errorf("commit resource operator state transaction: %w", err)
|
||||
}
|
||||
committed = true
|
||||
return state, nil
|
||||
}
|
||||
|
||||
// ClearResourceOperatorState removes the row for the given canonical
|
||||
// ID. Idempotent — returns nil whether or not a row existed.
|
||||
func (s *SQLiteResourceStore) ClearResourceOperatorState(canonicalID string) error {
|
||||
@@ -1631,13 +1704,64 @@ func (s *SQLiteResourceStore) ClearResourceOperatorState(canonicalID string) err
|
||||
if canonicalID == "" {
|
||||
return nil
|
||||
}
|
||||
_, err := s.db.Exec(`DELETE FROM resource_operator_state WHERE canonical_id = ?`, canonicalID)
|
||||
|
||||
s.mu.Lock()
|
||||
defer s.mu.Unlock()
|
||||
|
||||
return clearResourceOperatorStateSQL(s.db, canonicalID)
|
||||
}
|
||||
|
||||
func clearResourceOperatorStateSQL(execer sqlExecutor, canonicalID string) error {
|
||||
_, err := execer.Exec(`DELETE FROM resource_operator_state WHERE canonical_id = ?`, canonicalID)
|
||||
if err != nil {
|
||||
return fmt.Errorf("delete resource operator state: %w", err)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
// ClearResourceOperatorStateWithMaintenanceLifecycle deletes the operator-state
|
||||
// source row and the derived maintenance-window timeline projection in a single
|
||||
// SQLite transaction.
|
||||
func (s *SQLiteResourceStore) ClearResourceOperatorStateWithMaintenanceLifecycle(canonicalID string, observedAt time.Time, actor string) error {
|
||||
canonicalID = strings.TrimSpace(canonicalID)
|
||||
if canonicalID == "" {
|
||||
return nil
|
||||
}
|
||||
|
||||
s.mu.Lock()
|
||||
defer s.mu.Unlock()
|
||||
|
||||
tx, err := s.db.Begin()
|
||||
if err != nil {
|
||||
return fmt.Errorf("begin resource operator state clear transaction: %w", err)
|
||||
}
|
||||
committed := false
|
||||
defer func() {
|
||||
if !committed {
|
||||
_ = tx.Rollback()
|
||||
}
|
||||
}()
|
||||
|
||||
previous, previousFound, err := getResourceOperatorStateSQL(tx, canonicalID)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
if err := clearResourceOperatorStateSQL(tx, canonicalID); err != nil {
|
||||
return err
|
||||
}
|
||||
current := ResourceOperatorState{CanonicalID: canonicalID}
|
||||
if change, ok := BuildMaintenanceWindowLifecycleChange(previous, previousFound, current, false, observedAt, actor); ok {
|
||||
if err := recordChangeSQL(tx, change, s.resourceChangesHasTimestamp); err != nil {
|
||||
return fmt.Errorf("record maintenance window lifecycle change: %w", err)
|
||||
}
|
||||
}
|
||||
if err := tx.Commit(); err != nil {
|
||||
return fmt.Errorf("commit resource operator state clear transaction: %w", err)
|
||||
}
|
||||
committed = true
|
||||
return nil
|
||||
}
|
||||
|
||||
// MemoryStore is an in-memory implementation for tests.
|
||||
type MemoryStore struct {
|
||||
mu sync.RWMutex
|
||||
@@ -1700,6 +1824,11 @@ func (m *MemoryStore) Close() error {
|
||||
func (m *MemoryStore) RecordChange(change ResourceChange) error {
|
||||
m.mu.Lock()
|
||||
defer m.mu.Unlock()
|
||||
|
||||
return m.recordChangeLocked(change)
|
||||
}
|
||||
|
||||
func (m *MemoryStore) recordChangeLocked(change ResourceChange) error {
|
||||
for _, existing := range m.changes {
|
||||
if existing.ID == change.ID && change.ID != "" {
|
||||
return nil
|
||||
@@ -2173,6 +2302,28 @@ func (m *MemoryStore) SetResourceOperatorState(state ResourceOperatorState) erro
|
||||
return nil
|
||||
}
|
||||
|
||||
// SetResourceOperatorStateWithMaintenanceLifecycle updates the in-memory
|
||||
// source row and derived timeline projection under one lock.
|
||||
func (m *MemoryStore) SetResourceOperatorStateWithMaintenanceLifecycle(state ResourceOperatorState) (ResourceOperatorState, error) {
|
||||
state = NormalizeResourceOperatorState(state)
|
||||
if err := ValidateResourceOperatorState(state); err != nil {
|
||||
return ResourceOperatorState{}, err
|
||||
}
|
||||
m.mu.Lock()
|
||||
defer m.mu.Unlock()
|
||||
if m.resourceOperatorState == nil {
|
||||
m.resourceOperatorState = make(map[string]ResourceOperatorState)
|
||||
}
|
||||
previous, previousFound := m.resourceOperatorState[state.CanonicalID]
|
||||
m.resourceOperatorState[state.CanonicalID] = state
|
||||
if change, ok := BuildMaintenanceWindowLifecycleChange(previous, previousFound, state, true, state.SetAt, state.SetBy); ok {
|
||||
if err := m.recordChangeLocked(change); err != nil {
|
||||
return ResourceOperatorState{}, fmt.Errorf("record maintenance window lifecycle change: %w", err)
|
||||
}
|
||||
}
|
||||
return state, nil
|
||||
}
|
||||
|
||||
// ClearResourceOperatorState removes any operator-set state for the
|
||||
// given canonical ID. Returns nil whether or not an entry was present —
|
||||
// the operation is idempotent so the API surface can issue
|
||||
@@ -2188,6 +2339,26 @@ func (m *MemoryStore) ClearResourceOperatorState(canonicalID string) error {
|
||||
return nil
|
||||
}
|
||||
|
||||
// ClearResourceOperatorStateWithMaintenanceLifecycle clears the in-memory
|
||||
// source row and derived timeline projection under one lock.
|
||||
func (m *MemoryStore) ClearResourceOperatorStateWithMaintenanceLifecycle(canonicalID string, observedAt time.Time, actor string) error {
|
||||
canonicalID = strings.TrimSpace(canonicalID)
|
||||
if canonicalID == "" {
|
||||
return nil
|
||||
}
|
||||
m.mu.Lock()
|
||||
defer m.mu.Unlock()
|
||||
previous, previousFound := m.resourceOperatorState[canonicalID]
|
||||
delete(m.resourceOperatorState, canonicalID)
|
||||
current := ResourceOperatorState{CanonicalID: canonicalID}
|
||||
if change, ok := BuildMaintenanceWindowLifecycleChange(previous, previousFound, current, false, observedAt, actor); ok {
|
||||
if err := m.recordChangeLocked(change); err != nil {
|
||||
return fmt.Errorf("record maintenance window lifecycle change: %w", err)
|
||||
}
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
func normalizePair(a, b string) (string, string) {
|
||||
a = CanonicalResourceID(a)
|
||||
b = CanonicalResourceID(b)
|
||||
|
||||
Reference in New Issue
Block a user