From 653a28192ab14ace07a3cacacccdb7ef6f1a7c38 Mon Sep 17 00:00:00 2001 From: "pulse-triage[bot]" <249995291+pulse-triage[bot]@users.noreply.github.com> Date: Mon, 7 Sep 2026 22:20:18 +0100 Subject: [PATCH] fix(discovery): preserve current records during suggestion backfill Backfill could save a stale List snapshot after manual discovery repaired a service, restoring unknown identity and dropping its URL and engine version. Derive and persist missing suggestions from the current record under the store lock instead, without holding it across monitor reads. Add a deterministic SetReadState/manual-refresh interleaving and encrypted restart assertions, plus coverage for current identity, dismissed proposals, deletion and persistence failure. The discovery package passes twenty race-enabled repetitions. Change-source: pulse-maintainer --- .../v6/internal/subsystems/monitoring.md | 17 ++++ internal/servicediscovery/service.go | 13 +-- internal/servicediscovery/service_test.go | 93 +++++++++++++++++++ internal/servicediscovery/store.go | 32 +++++++ internal/servicediscovery/store_test.go | 92 ++++++++++++++++++ 5 files changed, 239 insertions(+), 8 deletions(-) diff --git a/docs/release-control/v6/internal/subsystems/monitoring.md b/docs/release-control/v6/internal/subsystems/monitoring.md index 4511892bf..71a991037 100644 --- a/docs/release-control/v6/internal/subsystems/monitoring.md +++ b/docs/release-control/v6/internal/subsystems/monitoring.md @@ -17,6 +17,23 @@ ## Purpose +**Availability backfill preserves concurrent discovery changes (7 September 2026)** + +The backfill List snapshot is a work list, not an authoritative record to save. +Read the current discovery, derive a missing suggestion from its current identity, +and persist under one store write lock. State-provider reads stay outside that +lock. Preserve newer identity, URL, engine version, user notes and existing +(including dismissed) proposals; never resurrect a discovery deleted after List. +A failed persistence attempt remains an error without installing the proposed +change in cache. `TestService_BackfillPreservesConcurrentManualRepair` in `service_test.go` +deterministically pauses SetReadState +backfill after List, completes manual ESPHome repair, then resumes backfill and +checks both cache and encrypted restart. +`TestStore_BackfillAvailabilitySuggestionUsesCurrentRecord` in `store_test.go` +covers current-identity +inference, dismissal, deletion, unsupported identity and persistence failure. +This is discovery state-integrity proof, not incident or notification acceptance. + **Correlated VM memory fallback — issue #1962 (7 September 2026)** The next Proxmox guest poll indexes live agent memory embedded in VM read views, diff --git a/internal/servicediscovery/service.go b/internal/servicediscovery/service.go index 6761bcc67..fd3b7c0ef 100644 --- a/internal/servicediscovery/service.go +++ b/internal/servicediscovery/service.go @@ -3514,14 +3514,11 @@ func (s *Service) backfillAvailabilitySuggestions(ctx context.Context) { } externalIP := s.getResourceExternalIP(req) - suggestion := SuggestAvailabilityProbe(d, externalIP) - if suggestion != nil { - d.SuggestedAvailabilityProbe = suggestion - if err := s.store.Save(d); err != nil { - log.Warn().Err(err).Str("id", d.ID).Msg("Failed to save backfilled availability suggestion") - } else { - updated++ - } + changed, err := s.store.backfillAvailabilitySuggestion(d.ID, externalIP) + if err != nil { + log.Warn().Err(err).Str("id", d.ID).Msg("Failed to save backfilled availability suggestion") + } else if changed { + updated++ } } diff --git a/internal/servicediscovery/service_test.go b/internal/servicediscovery/service_test.go index 90a12b1f6..ed9074e50 100644 --- a/internal/servicediscovery/service_test.go +++ b/internal/servicediscovery/service_test.go @@ -8,6 +8,7 @@ import ( "path/filepath" "strings" "sync" + "sync/atomic" "testing" "time" @@ -2667,3 +2668,95 @@ func TestService_ListDiscoveriesByTarget_DoesNotBridgeSharedHostnames(t *testing t.Fatalf("expected estate B to keep finding its own record, got %d", len(own)) } } + +// Pause the first snapshot read, which backfill performs after Store.List. +// Other reads remain available to the concurrent manual discovery. +type backfillBarrierState struct { + unifiedresources.ReadState + first atomic.Bool + listed chan struct{} + resume chan struct{} +} + +func (s *backfillBarrierState) VMs() []*unifiedresources.VMView { + if s.first.CompareAndSwap(false, true) { + close(s.listed) + <-s.resume + } + return s.ReadState.VMs() +} + +func TestService_BackfillPreservesConcurrentManualRepair(t *testing.T) { + store, err := NewStore(t.TempDir()) + if err != nil { + t.Fatal(err) + } + service := NewService(store, nil, DefaultConfig()) + rs := readStateFromSnapshot(StateSnapshot{Containers: []Container{ + {VMID: 102, Name: "esphome", Node: "pve1", Status: "running"}, + }}) + // Set up the fixture before starting the asynchronous backfill. + service.readState = rs + service.SetCommandScanningEnabled(true) + service.SetAIAnalyzer(&stubAnalyzer{response: `{}`}) + service.collectFingerprints(context.Background()) + id := MakeResourceID(ResourceTypeSystemContainer, "pve1", "102") + fp, err := store.GetFingerprint(id) + if err != nil || fp == nil { + t.Fatalf("fingerprint: %v, %v", fp, err) + } + if err := store.Save(&ResourceDiscovery{ + ID: id, ResourceType: ResourceTypeSystemContainer, TargetID: "pve1", ResourceID: "102", + Hostname: "esphome", ServiceType: "unknown", ServiceName: "Unknown Service", Category: CategoryUnknown, + Fingerprint: fp.Hash, FingerprintedAt: fp.GeneratedAt, FingerprintSchemaVersion: fp.SchemaVersion, + CLIAccessVersion: CLIAccessVersion, + }); err != nil { + t.Fatal(err) + } + barrier := &backfillBarrierState{ReadState: rs, listed: make(chan struct{}), resume: make(chan struct{})} + var release sync.Once + unblock := func() { release.Do(func() { close(barrier.resume) }) } + service.SetReadState(barrier) + t.Cleanup(func() { unblock(); service.backfillCancel(); <-service.backfillDone }) + select { + case <-barrier.listed: + case <-time.After(5 * time.Second): + t.Fatal("backfill did not reach snapshot barrier") + } + summary, err := service.RunManualDiscoveryRefresh(context.Background()) + if err != nil { + t.Fatal(err) + } + if summary.DiscoveredCount != 1 || summary.FailedCount != 0 { + t.Fatalf("refresh: %+v", summary) + } + repaired, err := store.Get(id) + if err != nil { + t.Fatal(err) + } + if repaired.ServiceType != "esphome" { + t.Fatalf("manual repair failed: %q", repaired.ServiceType) + } + unblock() + select { + case <-service.backfillDone: + case <-time.After(5 * time.Second): + t.Fatal("backfill did not finish") + } + // Read through a new encrypted store as well as the live cache. + restarted, err := NewStore(filepath.Dir(store.dataDir)) + if err != nil { + t.Fatal(err) + } + for name, reader := range map[string]*Store{"cache": store, "disk": restarted} { + got, err := reader.Get(id) + if err != nil || got == nil { + t.Fatalf("%s read: %v", name, err) + } + if got.ServiceType != repaired.ServiceType || got.ServiceName != repaired.ServiceName || + got.SuggestedURL != repaired.SuggestedURL || got.DiscoveryEngineVersion != repaired.DiscoveryEngineVersion { + t.Errorf("%s: backfill overwrote manual repair: type=%q name=%q URL=%q engine=%d", name, + got.ServiceType, got.ServiceName, got.SuggestedURL, got.DiscoveryEngineVersion) + } + } +} diff --git a/internal/servicediscovery/store.go b/internal/servicediscovery/store.go index ebf1b9b0d..29c532f3d 100644 --- a/internal/servicediscovery/store.go +++ b/internal/servicediscovery/store.go @@ -381,7 +381,11 @@ func (s *Store) marshalDiscoveryForStorage(discovery *ResourceDiscovery) ([]byte func (s *Store) Save(d *ResourceDiscovery) error { s.mu.Lock() defer s.mu.Unlock() + return s.saveLocked(d) +} +// saveLocked requires s.mu to be held for writing. +func (s *Store) saveLocked(d *ResourceDiscovery) error { if d.ID == "" { return fmt.Errorf("discovery ID is required") } @@ -433,7 +437,11 @@ func (s *Store) Get(id string) (*ResourceDiscovery, error) { s.mu.Lock() defer s.mu.Unlock() + return s.getLocked(id) +} +// getLocked requires s.mu to be held for writing (loads may migrate files). +func (s *Store) getLocked(id string) (*ResourceDiscovery, error) { filePath := s.getFilePath(id) activePath := filePath data, migratedPlaintext, err := s.loadDiscoveryFileData(filePath, maxDiscoveryFileReadBytes) @@ -481,6 +489,30 @@ func (s *Store) Get(id string) (*ResourceDiscovery, error) { return cloneResourceDiscovery(&discovery), nil } +// backfillAvailabilitySuggestion reads, derives and persists under one lock. +// A List snapshot is only a work list: never write its stale identity or resurrect +// a deleted discovery. State/monitor reads must happen before taking this lock. +func (s *Store) backfillAvailabilitySuggestion(id, externalIP string) (bool, error) { + s.mu.Lock() + defer s.mu.Unlock() + current, err := s.getLocked(id) + if err != nil || current == nil { + return false, err + } + if current.SuggestedAvailabilityProbe != nil { + return false, nil + } + suggestion := SuggestAvailabilityProbe(current, externalIP) + if suggestion == nil { + return false, nil + } + current.SuggestedAvailabilityProbe = suggestion + if err := s.saveLocked(current); err != nil { + return false, err + } + return true, nil +} + // GetByResource retrieves a discovery by resource type and ID. func (s *Store) GetByResource(resourceType ResourceType, targetID, resourceID string) (*ResourceDiscovery, error) { id := MakeResourceID(resourceType, targetID, resourceID) diff --git a/internal/servicediscovery/store_test.go b/internal/servicediscovery/store_test.go index 67fc64bdc..043ee43d3 100644 --- a/internal/servicediscovery/store_test.go +++ b/internal/servicediscovery/store_test.go @@ -1234,3 +1234,95 @@ func TestStore_GetStaleResources(t *testing.T) { t.Fatalf("expected GetStaleResources to return list error") } } + +func TestStore_BackfillAvailabilitySuggestionUsesCurrentRecord(t *testing.T) { + for _, scenario := range []string{"new identity", "existing dismissed proposal", "deleted", "unsupported", "write failure"} { + t.Run(scenario, func(t *testing.T) { + dir := t.TempDir() + store, err := NewStore(dir) + if err != nil { + t.Fatal(err) + } + id := MakeResourceID(ResourceTypeSystemContainer, "pve1", "102") + old := &ResourceDiscovery{ID: id, Hostname: "esphome", ServiceType: "unknown"} + if err := store.Save(old); err != nil { + t.Fatal(err) + } + work, err := store.List() + if err != nil || len(work) != 1 { + t.Fatalf("List: %v, %v", work, err) + } + current := cloneResourceDiscovery(old) + current.ServiceType = "redis" + current.ServiceName = "New Redis" + current.UserNotes = "operator note after list" + if scenario == "existing dismissed proposal" { + current.SuggestedAvailabilityProbe = SuggestAvailabilityProbe(current, "192.0.2.20") + current.DismissedAvailabilityProbeFingerprint = current.SuggestedAvailabilityProbe.EvidenceFingerprint + } + if scenario == "unsupported" { + current.ServiceType = "unknown" + current.Hostname = "unidentified" + } + if err := store.Save(current); err != nil { + t.Fatal(err) + } + before, err := store.Get(id) + if err != nil { + t.Fatal(err) + } + if scenario == "deleted" { + if err := store.Delete(id); err != nil { + t.Fatal(err) + } + } + if scenario == "write failure" { + // A directory at the temporary file path causes a real persistence error. + if err := os.Mkdir(store.getFilePath(id)+".tmp", 0700); err != nil { + t.Fatal(err) + } + } + changed, err := store.backfillAvailabilitySuggestion(work[0].ID, "192.0.2.10") + if scenario == "write failure" { + if err == nil || changed { + t.Fatalf("expected failed write, got changed=%v err=%v", changed, err) + } + } else if err != nil { + t.Fatal(err) + } + if changed != (scenario == "new identity") { + t.Fatalf("unexpected changed=%v", changed) + } + restarted, err := NewStore(dir) + if err != nil { + t.Fatal(err) + } + for name, reader := range map[string]*Store{"cache": store, "disk": restarted} { + got, err := reader.Get(id) + if err != nil { + t.Fatal(err) + } + if scenario == "deleted" { + if got != nil { + t.Errorf("%s: deleted discovery resurrected", name) + } + continue + } + want := cloneResourceDiscovery(before) + if scenario == "new identity" { + want.SuggestedAvailabilityProbe = SuggestAvailabilityProbe(before, "192.0.2.10") + want.UpdatedAt = got.UpdatedAt + if got.SuggestedAvailabilityProbe == nil || got.SuggestedAvailabilityProbe.Port != 6379 { + t.Fatalf("%s: suggestion not derived from current Redis identity", name) + } + } + // JSON comparison ignores time.Time's process-local monotonic clock. + gotJSON, _ := json.Marshal(got) + wantJSON, _ := json.Marshal(want) + if string(gotJSON) != string(wantJSON) { + t.Errorf("%s: backfill changed unrelated/current fields", name) + } + } + }) + } +}