From c1db220e3a2f99afe56c7bf5db2a028aa954000d Mon Sep 17 00:00:00 2001 From: rcourtman Date: Wed, 13 May 2026 14:18:16 +0100 Subject: [PATCH] Make maintenance evidence writes atomic --- internal/api/resources_operator_state.go | 55 +---- internal/api/resources_operator_state_test.go | 102 ++++++++- .../unifiedresources/loop_reports_store.go | 27 ++- .../loop_reports_store_test.go | 30 +++ .../resource_operator_state.go | 25 +++ internal/unifiedresources/store.go | 193 +++++++++++++++++- 6 files changed, 361 insertions(+), 71 deletions(-) diff --git a/internal/api/resources_operator_state.go b/internal/api/resources_operator_state.go index 5e376277f..1a624158c 100644 --- a/internal/api/resources_operator_state.go +++ b/internal/api/resources_operator_state.go @@ -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) -} diff --git a/internal/api/resources_operator_state_test.go b/internal/api/resources_operator_state_test.go index 89833f87d..2d7f87ca9 100644 --- a/internal/api/resources_operator_state_test.go +++ b/internal/api/resources_operator_state_test.go @@ -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) diff --git a/internal/unifiedresources/loop_reports_store.go b/internal/unifiedresources/loop_reports_store.go index 056e39148..9249bada4 100644 --- a/internal/unifiedresources/loop_reports_store.go +++ b/internal/unifiedresources/loop_reports_store.go @@ -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 } diff --git a/internal/unifiedresources/loop_reports_store_test.go b/internal/unifiedresources/loop_reports_store_test.go index 1aa499f67..2c3ebf23c 100644 --- a/internal/unifiedresources/loop_reports_store_test.go +++ b/internal/unifiedresources/loop_reports_store_test.go @@ -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") diff --git a/internal/unifiedresources/resource_operator_state.go b/internal/unifiedresources/resource_operator_state.go index 642cc05fe..06685238f 100644 --- a/internal/unifiedresources/resource_operator_state.go +++ b/internal/unifiedresources/resource_operator_state.go @@ -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" diff --git a/internal/unifiedresources/store.go b/internal/unifiedresources/store.go index 9a28b4aa4..3c7f56f59 100644 --- a/internal/unifiedresources/store.go +++ b/internal/unifiedresources/store.go @@ -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)