diff --git a/docs/release-control/v6/internal/subsystems/monitoring.md b/docs/release-control/v6/internal/subsystems/monitoring.md index 28ee482cc..8b3794717 100644 --- a/docs/release-control/v6/internal/subsystems/monitoring.md +++ b/docs/release-control/v6/internal/subsystems/monitoring.md @@ -17,6 +17,19 @@ ## Purpose +### Canonical registry publication during accepted ingest — issue #2199 + +The monitor adapter builds a replacement resource generation away from readers. +The previous generation remains readable while change records are classified +and persisted; classification releases registry read locks before any backing +store write can block. Writer ordering, journal classification, active-alert +copying and publication of the replacement generation remain intact. A repeat +of an unchanged snapshot must not append duplicate change records. +`TestMonitorAdapterPopulateFromSnapshotDoesNotRepeatChangeJournal` and +`TestMonitorAdapterSerializesSupplementalMutationAfterSnapshotPublication` +cover the publish and blocked-persistence boundaries. Synthetic performance +improvement alone does not establish field CPU relief. + ### Linked Pulse agent memory for Proxmox LXC — issue #2148 (22 September 2026) Both LXC memory paths (the efficient `cluster/resources` builder and the diff --git a/docs/release-control/v6/internal/subsystems/unified-resources.md b/docs/release-control/v6/internal/subsystems/unified-resources.md index b3f2ac858..071783769 100644 --- a/docs/release-control/v6/internal/subsystems/unified-resources.md +++ b/docs/release-control/v6/internal/subsystems/unified-resources.md @@ -23,6 +23,21 @@ and sort the complete canonical change table while startup and ingestion wait. ## Purpose +### Bounded generation change comparison — issue #2199 + +Accepted snapshot rebuilds compare the previous and replacement registry +generations by resource ID without deep-cloning and sorting both complete +resource lists solely for change emission. Classification reads both generations +under registry read locks, materialises change records, and releases those locks +before the backing store persists them. Discovery, removal, state, relationship +and configuration classifications retain the existing journal semantics; +unchanged telemetry and volatile relationship timestamps produce no new +journal row. `TestRecordRegistryChangesBetweenGenerationsMatchesListComparison` +compares the old and new paths, while +`TestRegistryGenerationComparisonIgnoresUnchangedTelemetry` pins the no-op +boundary. The 1,000-resource benchmark is component evidence, not installed +fleet CPU attribution or a release acceptance claim. + Canonical frontend memory withdrawal is explicit: when a resource snapshot omits the canonical memory metric and its Proxmox memory facet marks usageUnavailable, the display merge must clear any previous metric. Plain partial omission remains compatible with richer REST state, and an incoming canonical metric (including measured zero) takes precedence over unavailable raw evidence. The adapter transition tests pin withdrawal and recovery; the hybrid-memory Chromium fixture exercises the rendered table and drawer at 1280px and 390px. Workload details remain canonical-only: withdrawal shows N/A in the table and removes the Memory section rather than manufacturing a raw-facet total; measured-zero recovery restores Total and Free, with screenshot positioning above fixed navigation. This does not change agent-only or arbitrary field-deletion semantics. **VM-linked agent memory read — issue #1962 (7 September 2026)** diff --git a/internal/unifiedresources/adapter_coverage_test.go b/internal/unifiedresources/adapter_coverage_test.go index 3ca54555d..12d961af3 100644 --- a/internal/unifiedresources/adapter_coverage_test.go +++ b/internal/unifiedresources/adapter_coverage_test.go @@ -335,6 +335,29 @@ func TestMonitorAdapterPopulateFromSnapshot(t *testing.T) { } } +func TestMonitorAdapterPopulateFromSnapshotDoesNotRepeatChangeJournal(t *testing.T) { + store := NewMemoryStore() + adapter := NewMonitorAdapter(NewRegistry(store)) + observedAt := time.Date(2026, 9, 24, 6, 0, 0, 0, time.UTC) + snapshot := models.StateSnapshot{VMs: []models.VM{{ + ID: "lab:node-a:101", + VMID: 101, + Name: "database", + Node: "node-a", + Instance: "lab", + Status: "running", + LastSeen: observedAt, + }}} + adapter.PopulateFromSnapshot(snapshot) + if got := len(store.changes); got != 1 { + t.Fatalf("initial snapshot emitted %d change records, want one discovery", got) + } + adapter.PopulateFromSnapshot(snapshot) + if got := len(store.changes); got != 1 { + t.Fatalf("identical snapshot emitted %d total change records, want one", got) + } +} + func TestMonitorAdapterPopulateFromSnapshotReplacesPreviousRegistryState(t *testing.T) { now := time.Date(2026, 3, 7, 12, 0, 0, 0, time.UTC) adapter := NewMonitorAdapter(NewRegistry(nil)) diff --git a/internal/unifiedresources/change_emission.go b/internal/unifiedresources/change_emission.go index 9f4070fd0..ea3a15a96 100644 --- a/internal/unifiedresources/change_emission.go +++ b/internal/unifiedresources/change_emission.go @@ -53,10 +53,71 @@ func recordRegistryChanges(store ResourceStore, before, after []Resource, observ } } +// recordRegistryChangesBetweenGenerations compares two complete registry +// generations without making the deep, sorted List clones needed by public +// readers. The resources are read only while both registry locks are held; +// persistence happens after releasing them so a slow store cannot block +// registry readers. The caller must keep the old generation stable until this +// returns (MonitorAdapter's mutationMu does so during a replacement). +func recordRegistryChangesBetweenGenerations(before, after *ResourceRegistry, observedAt time.Time, occurredAt *time.Time, sourceType ChangeSourceType, sourceAdapterHint ChangeSourceAdapter) { + if before == nil || after == nil || before.store == nil { + return + } + + var pending []ResourceChange + before.mu.RLock() + after.mu.RLock() + beforeByID := make(map[string]*Resource, len(before.resources)) + for _, resource := range before.resources { + if resource != nil && resource.ID != "" { + beforeByID[resource.ID] = resource + } + } + afterByID := make(map[string]*Resource, len(after.resources)) + for _, resource := range after.resources { + if resource != nil && resource.ID != "" { + afterByID[resource.ID] = resource + } + } + for id, previous := range beforeByID { + current, exists := afterByID[id] + var next Resource + if exists { + next = *current + } + if change := buildResourceChange(*previous, true, next, exists, observedAt, occurredAt, sourceType, sourceAdapterHint); change != nil { + pending = append(pending, *change) + } + } + for id, current := range afterByID { + if _, exists := beforeByID[id]; exists { + continue + } + if change := buildResourceChange(Resource{}, false, *current, true, observedAt, occurredAt, sourceType, sourceAdapterHint); change != nil { + pending = append(pending, *change) + } + } + after.mu.RUnlock() + before.mu.RUnlock() + + for _, change := range pending { + if err := before.store.RecordChange(change); err != nil { + log.Printf("unifiedresources: failed to record change for %s: %v", change.ResourceID, err) + } + } +} + func buildResourceChange(before Resource, beforeOK bool, after Resource, afterOK bool, observedAt time.Time, occurredAt *time.Time, sourceType ChangeSourceType, sourceAdapterHint ChangeSourceAdapter) *ResourceChange { if !beforeOK && !afterOK { return nil } + var changedFields []string + if beforeOK && afterOK { + changedFields = resourceChangedFields(before, after) + if len(changedFields) == 0 { + return nil + } + } change := &ResourceChange{ ID: uuid.NewString(), @@ -102,11 +163,6 @@ func buildResourceChange(before Resource, beforeOK bool, after Resource, afterOK return change } - changedFields := resourceChangedFields(before, after) - if len(changedFields) == 0 { - return nil - } - change.Metadata = map[string]any{"changedFields": changedFields} switch { diff --git a/internal/unifiedresources/change_emission_test.go b/internal/unifiedresources/change_emission_test.go index 4ef1ea03d..76140455b 100644 --- a/internal/unifiedresources/change_emission_test.go +++ b/internal/unifiedresources/change_emission_test.go @@ -1,10 +1,49 @@ package unifiedresources import ( + "reflect" + "sort" "testing" "time" ) +func TestRecordRegistryChangesBetweenGenerationsMatchesListComparison(t *testing.T) { + newStore := NewMemoryStore() + referenceStore := NewMemoryStore() + before := NewRegistry(newStore) + after := NewRegistry(newStore) + base := []IngestRecord{ + {SourceID: "vm-a", Resource: Resource{Type: ResourceTypeVM, Name: "a", Status: StatusOnline}}, + {SourceID: "vm-b", Resource: Resource{Type: ResourceTypeVM, Name: "b", Status: StatusOnline}}, + {SourceID: "vm-unchanged", Resource: Resource{Type: ResourceTypeVM, Name: "unchanged", Status: StatusOnline}}, + } + before.IngestRecords(SourceProxmox, base) + after.IngestRecords(SourceProxmox, []IngestRecord{ + {SourceID: "vm-a", Resource: Resource{Type: ResourceTypeVM, Name: "a", Status: StatusOffline}}, + base[2], + {SourceID: "vm-c", Resource: Resource{Type: ResourceTypeVM, Name: "c", Status: StatusOnline}}, + }) + + observedAt := time.Date(2026, 9, 24, 6, 0, 0, 0, time.UTC) + recordRegistryChanges(referenceStore, before.List(), after.List(), observedAt, nil, SourcePulseDiff, "") + recordRegistryChangesBetweenGenerations(before, after, observedAt, nil, SourcePulseDiff, "") + + if got := len(referenceStore.changes); got != 3 { + t.Fatalf("reference changes = %d, want changed, removed and added", got) + } + normalize := func(changes []ResourceChange) []ResourceChange { + out := append([]ResourceChange(nil), changes...) + for i := range out { + out[i].ID = "" // Change IDs are generated independently. + } + sort.Slice(out, func(i, j int) bool { return out[i].ResourceID < out[j].ResourceID }) + return out + } + if got, want := normalize(newStore.changes), normalize(referenceStore.changes); !reflect.DeepEqual(got, want) { + t.Fatalf("generation comparison differs from List comparison:\n got: %+v\nwant: %+v", got, want) + } +} + func TestBuildResourceChange_ReturnsNilWhenUnchanged(t *testing.T) { before := Resource{ ID: "vm:1", diff --git a/internal/unifiedresources/monitor_adapter.go b/internal/unifiedresources/monitor_adapter.go index af73ed833..2f54e6e2d 100644 --- a/internal/unifiedresources/monitor_adapter.go +++ b/internal/unifiedresources/monitor_adapter.go @@ -314,7 +314,6 @@ func (a *MonitorAdapter) replaceRegistryLocked(snapshot models.StateSnapshot, re return } - before := registry.List() rebuilt := NewRegistry(registry.store) staleThresholds := a.currentStaleThresholds() rebuilt.IngestSnapshotWithStaleThresholds(snapshot, staleThresholds) @@ -353,7 +352,7 @@ func (a *MonitorAdapter) replaceRegistryLocked(snapshot models.StateSnapshot, re occ := snapshot.LastUpdate occurredAt = &occ } - recordRegistryChanges(registry.store, before, rebuilt.List(), rebuiltAt, occurredAt, SourcePulseDiff, "") + recordRegistryChangesBetweenGenerations(registry, rebuilt, rebuiltAt, occurredAt, SourcePulseDiff, "") rebuilt.PersistIdentityPins() a.mu.Lock() diff --git a/internal/unifiedresources/monitor_adapter_bench_test.go b/internal/unifiedresources/monitor_adapter_bench_test.go new file mode 100644 index 000000000..ba83ee9dd --- /dev/null +++ b/internal/unifiedresources/monitor_adapter_bench_test.go @@ -0,0 +1,96 @@ +package unifiedresources + +import ( + "fmt" + "testing" + "time" + + "github.com/rcourtman/pulse-go-rewrite/internal/models" +) + +// benchmarkVMState is a synthetic 1,000-guest snapshot for comparing the +// accepted-ingest rebuild with its constituent registry operations. It does +// not model an installed fleet, provider polling, or durable-store latency. +func benchmarkVMState(count int) models.StateSnapshot { + now := time.Now().UTC() + snapshot := models.StateSnapshot{VMs: make([]models.VM, count), LastUpdate: now} + for i := range snapshot.VMs { + snapshot.VMs[i] = models.VM{ + ID: fmt.Sprintf("lab:node-1:%d", i+100), + VMID: i + 100, + Name: fmt.Sprintf("vm-%d", i), + Node: "node-1", + Instance: "lab", + Status: "running", + CPUs: 4, + CPU: 0.25, + LastSeen: now, + } + } + return snapshot +} + +func BenchmarkMonitorAdapterReplaceRegistry1000VMs(b *testing.B) { + snapshot := benchmarkVMState(1000) + // A memory-backed store exercises change comparison and identity-pin + // handling without measuring disk I/O or relying on host services. + adapter := NewMonitorAdapter(NewRegistry(NewMemoryStore())) + adapter.replaceRegistry(snapshot, nil) + if got := len(adapter.GetAll()); got != len(snapshot.VMs) { + b.Fatalf("warm registry length = %d, want %d", got, len(snapshot.VMs)) + } + b.ReportAllocs() + b.ResetTimer() + for i := 0; i < b.N; i++ { + adapter.replaceRegistry(snapshot, nil) + } + b.StopTimer() + if got := len(adapter.GetAll()); got != len(snapshot.VMs) { + b.Fatalf("final registry length = %d, want %d", got, len(snapshot.VMs)) + } +} + +func BenchmarkRegistryIngestSnapshot1000VMs(b *testing.B) { + snapshot := benchmarkVMState(1000) + b.ReportAllocs() + b.ResetTimer() + for i := 0; i < b.N; i++ { + registry := NewRegistry(nil) + registry.IngestSnapshot(snapshot) + } +} + +func BenchmarkRegistryList1000VMs(b *testing.B) { + registry := NewRegistry(nil) + registry.IngestSnapshot(benchmarkVMState(1000)) + if got := len(registry.List()); got != 1000 { + b.Fatalf("registry length = %d, want 1000", got) + } + b.ReportAllocs() + b.ResetTimer() + for i := 0; i < b.N; i++ { + registry.List() + } +} + +func BenchmarkRegistryChangeEmission1000VMs(b *testing.B) { + snapshot := benchmarkVMState(1000) + store := NewMemoryStore() + before := NewRegistry(store) + after := NewRegistry(store) + before.IngestSnapshot(snapshot) + after.IngestSnapshot(snapshot) + observedAt := snapshot.LastUpdate + b.Run("list-clones", func(b *testing.B) { + b.ReportAllocs() + for i := 0; i < b.N; i++ { + recordRegistryChanges(store, before.List(), after.List(), observedAt, nil, SourcePulseDiff, "") + } + }) + b.Run("locked-generations", func(b *testing.B) { + b.ReportAllocs() + for i := 0; i < b.N; i++ { + recordRegistryChangesBetweenGenerations(before, after, observedAt, nil, SourcePulseDiff, "") + } + }) +} diff --git a/internal/unifiedresources/monitor_adapter_generation_test.go b/internal/unifiedresources/monitor_adapter_generation_test.go index 8d8f42082..06eeca18d 100644 --- a/internal/unifiedresources/monitor_adapter_generation_test.go +++ b/internal/unifiedresources/monitor_adapter_generation_test.go @@ -54,6 +54,20 @@ func TestMonitorAdapterSerializesSupplementalMutationAfterSnapshotPublication(t t.Fatal("snapshot rebuild did not reach change publication") } + // Change persistence may block, but it must not retain either registry's + // read lock. Readers still use the previously published generation. + readDone := make(chan struct{}) + go func() { + adapter.GetAll() + close(readDone) + }() + select { + case <-readDone: + case <-time.After(time.Second): + close(store.release) + t.Fatal("registry read blocked behind change persistence") + } + supplementalDone := make(chan struct{}) go func() { adapter.PopulateSupplementalRecords(SourceAgent, []IngestRecord{{ diff --git a/internal/unifiedresources/registry_test.go b/internal/unifiedresources/registry_test.go index 34fc0727e..9d05b7d72 100644 --- a/internal/unifiedresources/registry_test.go +++ b/internal/unifiedresources/registry_test.go @@ -33,6 +33,31 @@ func TestRegistry_CachedReadsUseSharedLock(t *testing.T) { } } +func TestRegistryGenerationComparisonIgnoresUnchangedTelemetry(t *testing.T) { + store := NewMemoryStore() + before := NewRegistry(store) + after := NewRegistry(store) + observedAt := time.Date(2026, 9, 24, 6, 0, 0, 0, time.UTC) + first := IngestRecord{ + SourceID: "vm-101", + Resource: Resource{ + Type: ResourceTypeVM, + Name: "vm-101", + Status: StatusOnline, + LastSeen: observedAt, + }, + } + before.IngestRecords(SourceProxmox, []IngestRecord{first}) + updated := first + updated.Resource.LastSeen = observedAt.Add(time.Minute) + after.IngestRecords(SourceProxmox, []IngestRecord{updated}) + + recordRegistryChangesBetweenGenerations(before, after, observedAt.Add(time.Minute), nil, SourcePulseDiff, "") + if got := len(store.changes); got != 0 { + t.Fatalf("telemetry-only generation emitted %d change records, want none", got) + } +} + // TestMemoryStore_RecordActionAuditAppliesRedaction is an integration check // at the registry-store boundary. The MemoryStore is the backing store the // registry uses in tests and contract examples, and operator-authored audit diff --git a/internal/websocket/state_delta_bench_test.go b/internal/websocket/state_delta_bench_test.go new file mode 100644 index 000000000..c24104fbf --- /dev/null +++ b/internal/websocket/state_delta_bench_test.go @@ -0,0 +1,64 @@ +package websocket + +import ( + "fmt" + "testing" + + "github.com/rcourtman/pulse-go-rewrite/internal/models" +) + +// benchmarkFrontendState isolates once-per-broadcast snapshot/delta work at +// the reported fleet scale. It is not a network or installed CPU benchmark. +func benchmarkFrontendState(count int) models.StateFrontend { + state := models.EmptyStateFrontend() + state.Resources = make([]models.ResourceFrontend, count) + for i := range state.Resources { + name := fmt.Sprintf("vm-%d", i) + state.Resources[i] = models.ResourceFrontend{ + ID: fmt.Sprintf("proxmox:lab:%d", i+100), + Type: "vm", + Name: name, + DisplayName: name, + PlatformType: "proxmox", + SourceType: "proxmox", + Status: "online", + CPU: &models.ResourceMetricFrontend{Current: 25}, + Memory: &models.ResourceMetricFrontend{Current: 60}, + Disk: &models.ResourceMetricFrontend{Current: 45}, + } + } + return state +} + +func BenchmarkBuildClientStateSnapshot1000Resources(b *testing.B) { + state := benchmarkFrontendState(1000) + b.ReportAllocs() + b.ResetTimer() + for i := 0; i < b.N; i++ { + if _, err := buildClientStateSnapshot(state); err != nil { + b.Fatal(err) + } + } +} + +func BenchmarkBuildClientStateDelta1000Resources(b *testing.B) { + previousState := benchmarkFrontendState(1000) + currentState := previousState + currentState.Resources = append([]models.ResourceFrontend(nil), previousState.Resources...) + currentState.Resources[0].CPU = &models.ResourceMetricFrontend{Current: 26} + previous, err := buildClientStateSnapshot(previousState) + if err != nil { + b.Fatal(err) + } + current, err := buildClientStateSnapshot(currentState) + if err != nil { + b.Fatal(err) + } + b.ReportAllocs() + b.ResetTimer() + for i := 0; i < b.N; i++ { + if _, err := buildClientStateDelta(previous, current); err != nil { + b.Fatal(err) + } + } +}