Merge pull request #2225 from rcourtman/maintainer/20260924T065120Z
Some checks are pending
Build and Test / Backend tests (rest-1) (push) Blocked by required conditions
Build and Test / Secret Scan (push) Waiting to run
Build and Test / Detect changed areas (push) Waiting to run
Build and Test / Frontend (push) Blocked by required conditions
Build and Test / Backend tests (api) (push) Blocked by required conditions
Build and Test / Backend tests (rest-0) (push) Blocked by required conditions
Build and Test / Script smoke tests & backend build (push) Blocked by required conditions
Build and Test / Benchmarks (push) Blocked by required conditions
Canonical Governance / governance (push) Waiting to run
Canonical Private Governance / private-governance (push) Waiting to run
Public docs / check (push) Waiting to run
Core E2E Tests / Validate E2E tier selection (push) Waiting to run
Core E2E Tests / Offline Organization provisioning (push) Waiting to run
Core E2E Tests / Playwright Core E2E (shard 1/8) (push) Blocked by required conditions
Core E2E Tests / Playwright Core E2E (shard 2/8) (push) Blocked by required conditions
Core E2E Tests / Playwright Core E2E (shard 3/8) (push) Blocked by required conditions
Core E2E Tests / Playwright Core E2E (shard 4/8) (push) Blocked by required conditions
Core E2E Tests / Playwright Core E2E (shard 5/8) (push) Blocked by required conditions
Core E2E Tests / Playwright Core E2E (shard 6/8) (push) Blocked by required conditions
Core E2E Tests / Playwright Core E2E (shard 7/8) (push) Blocked by required conditions
Core E2E Tests / Playwright Core E2E (shard 8/8) (push) Blocked by required conditions
Core E2E Tests / Agent registration lifecycle (push) Waiting to run
Core E2E Tests / E2E verdict (push) Blocked by required conditions
Unified Agent Native Verification / Linux ARM64 (push) Waiting to run
Unified Agent Native Verification / Linux x64 (push) Waiting to run
Unified Agent Native Verification / Windows x64 (push) Waiting to run
Unified Agent Native Verification / macOS ARM64 (push) Waiting to run
Unified Agent Native Verification / macOS Intel (push) Waiting to run
Unified Agent Native Verification / FreeBSD cross-build contract (push) Waiting to run

Reduce CPU work when resource snapshots are accepted
This commit is contained in:
pulse-triage[bot] 2026-09-24 07:26:20 +00:00 • committed by GitHub
commit e2b41c379d
No known key found for this signature in database
GPG key ID: B5690EEEBB952194
10 changed files with 351 additions and 7 deletions

View file

@ -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

View file

@ -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)**

View file

@ -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))

View file

@ -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 {

View file

@ -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",

View file

@ -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()

View file

@ -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, "")
}
})
}

View file

@ -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{{

View file

@ -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

View file

@ -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)
}
}
}