diff --git a/internal/monitoring/guest_memory_sources.go b/internal/monitoring/guest_memory_sources.go index 605c83d1a..d26d7c48f 100644 --- a/internal/monitoring/guest_memory_sources.go +++ b/internal/monitoring/guest_memory_sources.go @@ -382,3 +382,49 @@ func (m *Monitor) resolveGuestStatusMemory( return memTotal, memUsed, memorySource } + +// preferLinkedAgentLXCMemory replaces a low-trust Proxmox LXC memory reading +// with the linked Pulse agent's own sample when one is available. The agent runs +// inside the container, so its sample reflects the workload's real footprint, +// while cluster/resources is a cache-inclusive cgroup fallback that can badly +// under- or over-report shared-memory workloads (#2148, mirroring #1962). +// +// An agent inside a container without lxcfs sees the host's /proc/meminfo, so a +// sample is only trusted when its total matches the container's configured +// limit. Preferred provider sources are never overridden. +func preferLinkedAgentLXCMemory( + status string, + memTotal uint64, + memUsed uint64, + memorySource string, + guestRaw *VMMemoryRaw, + agentHost models.Host, + hasAgent bool, +) (uint64, uint64, string) { + if !hasAgent || status != "running" || memTotal == 0 { + return memTotal, memUsed, memorySource + } + switch CanonicalMemorySource(memorySource) { + case "cluster-resources", "unavailable": + default: + return memTotal, memUsed, memorySource + } + if !agentHost.Memory.HasKnownUsage() || agentHost.Memory.Total <= 0 || agentHost.Memory.Used < 0 { + return memTotal, memUsed, memorySource + } + agentTotal := uint64(agentHost.Memory.Total) + agentUsed := uint64(agentHost.Memory.Used) + if agentTotal != memTotal || agentUsed > memTotal { + return memTotal, memUsed, memorySource + } + if guestRaw != nil { + guestRaw.HostAgentTotal = agentTotal + guestRaw.HostAgentUsed = agentUsed + } + log.Debug(). + Uint64("total", memTotal). + Uint64("used", agentUsed). + Str("previousSource", memorySource). + Msg("LXC memory: using linked Pulse agent memory over cluster-resources fallback") + return memTotal, agentUsed, "agent" +} diff --git a/internal/monitoring/memory_trust_characterization_test.go b/internal/monitoring/memory_trust_characterization_test.go index f671d1678..2d12f5091 100644 --- a/internal/monitoring/memory_trust_characterization_test.go +++ b/internal/monitoring/memory_trust_characterization_test.go @@ -1004,6 +1004,7 @@ func TestHandleClusterContainerResourceMemoryTrustCharacterization(t *testing.T) makeGuestID("test", "node1", tt.res.VMID), client, nil, + nil, ) if !ok { t.Fatal("handleClusterContainerResource() returned ok=false") diff --git a/internal/monitoring/monitor_additional_test.go b/internal/monitoring/monitor_additional_test.go index bba390a31..15a9875f7 100644 --- a/internal/monitoring/monitor_additional_test.go +++ b/internal/monitoring/monitor_additional_test.go @@ -1009,6 +1009,45 @@ func TestCorrelatedGuestMemoryNextPoll(t *testing.T) { } } +// Issue #2148: a correlated agent merged into a system-container resource must +// surface through the previous-guest context so the LXC builder prefers it over +// the cache-inclusive cluster/resources fallback. +func TestCorrelatedContainerMemoryNextPoll(t *testing.T) { + now := time.Now() + const gib = int64(1024 * 1024 * 1024) + guestID := makeGuestID("pve-a", "node1", 100) + store := unifiedresources.NewMemoryStore() + if err := store.AddLink(unifiedresources.ResourceLink{ResourceA: "ct-test", ResourceB: "agent-test", PrimaryID: "agent-test"}); err != nil { + t.Fatal(err) + } + registry := unifiedresources.NewRegistry(store) + total, used := int64(24)*gib, int64(11)*gib + resources := []unifiedresources.Resource{ + {ID: "ct-test", Type: unifiedresources.ResourceTypeSystemContainer, Name: "npu-ct", Status: unifiedresources.StatusOnline, LastSeen: now, + Sources: []unifiedresources.DataSource{unifiedresources.SourceProxmox}, + Proxmox: &unifiedresources.ProxmoxData{SourceID: guestID, Instance: "pve-a", NodeName: "node1", VMID: 100, RuntimeStatus: "running"}, + Metrics: &unifiedresources.ResourceMetrics{Memory: &unifiedresources.MetricValue{Total: &total, Used: &used, Percent: 45, Source: unifiedresources.SourceProxmox}}}, + {ID: "agent-test", Type: unifiedresources.ResourceTypeAgent, Name: "npu-ct-agent", Status: unifiedresources.StatusOnline, LastSeen: now, + Sources: []unifiedresources.DataSource{unifiedresources.SourceAgent}, + Agent: &unifiedresources.AgentData{AgentID: "agent-test", Memory: &unifiedresources.AgentMemoryMeta{Total: total, Used: used, Free: total - used}}}, + } + registry.IngestResources(resources) + if len(registry.Containers()) != 1 || len(registry.Hosts()) != 0 { + t.Fatalf("expected one merged container and no standalone hosts, got %d/%d", len(registry.Containers()), len(registry.Hosts())) + } + mon := &Monitor{state: models.NewState(), rateTracker: NewRateTracker(), config: &config.Config{}, resourceStore: unifiedresources.NewMonitorAdapter(registry)} + prev := mon.previousGuestContextForInstance("pve-a") + if _, ok := prev.hostAgentsByVMID[guestID]; !ok { + t.Fatalf("linked container agent missing from previous context: %+v", prev.hostAgentsByVMID) + } + container, _, source, _, ok := mon.buildContainerFromClusterResource(context.Background(), "pve-a", + proxmox.ClusterResource{Type: "lxc", Status: "running", Node: "node1", VMID: 100, Name: "npu-ct", MaxMem: uint64(total), Mem: 3459743744}, + &stubPVEClient{}, map[int]bool{}, prev.hostAgentsByVMID) + if !ok || source != "agent" || container.Memory.Used != used { + t.Fatalf("container memory not agent-backed: ok=%v source=%s mem=%+v", ok, source, container.Memory) + } +} + type correlatedMemoryClient struct{ stubPVEClient } func (*correlatedMemoryClient) GetVMStatus(context.Context, string, int) (*proxmox.VMStatus, error) { diff --git a/internal/monitoring/monitor_polling_containers.go b/internal/monitoring/monitor_polling_containers.go index c3148f7f3..f9313d81a 100644 --- a/internal/monitoring/monitor_polling_containers.go +++ b/internal/monitoring/monitor_polling_containers.go @@ -154,6 +154,16 @@ func (m *Monitor) collectContainersWithNodes(ctx context.Context, instanceName s MaxMem: container.MaxMem, Mem: container.Mem, }) + agentHost, hasAgent := prevGuests.hostAgentsByVMID[guestID] + memTotal, memUsed, memorySource = preferLinkedAgentLXCMemory( + container.Status, + memTotal, + memUsed, + memorySource, + &guestRaw, + agentHost, + hasAgent, + ) memUsed, memorySource, _ = stabilizeGuestLowTrustMemory( m.previousGuestSnapshot(instanceName, "lxc", n.Node, int(container.VMID)), container.Status, diff --git a/internal/monitoring/monitor_previous_state.go b/internal/monitoring/monitor_previous_state.go index 86abbc1ac..bfd30496c 100644 --- a/internal/monitoring/monitor_previous_state.go +++ b/internal/monitoring/monitor_previous_state.go @@ -61,6 +61,9 @@ func (m *Monitor) previousGuestContextForInstance(instanceName string) previousG guestID := makeGuestID(container.Instance, container.Node, container.VMID) if guestID != "" { ctx.containersByID[guestID] = container + if memory, ok := ct.LinkedAgentMemory(); ok { + ctx.hostAgentsByVMID[guestID] = models.Host{LinkedVMID: guestID, Status: "online", Memory: memory} + } } if container.VMID > 0 && (strings.EqualFold(strings.TrimSpace(container.Type), "oci") || container.IsOCI) { ctx.containerOCIByVMID[container.VMID] = true @@ -72,10 +75,15 @@ func (m *Monitor) previousGuestContextForInstance(instanceName string) previousG continue } modelHost := previousHostFromView(host) - if modelHost.LinkedVMID == "" || modelHost.Status != "online" { + if modelHost.Status != "online" { continue } - ctx.hostAgentsByVMID[modelHost.LinkedVMID] = modelHost + if modelHost.LinkedVMID != "" { + ctx.hostAgentsByVMID[modelHost.LinkedVMID] = modelHost + } + if modelHost.LinkedContainerID != "" { + ctx.hostAgentsByVMID[modelHost.LinkedContainerID] = modelHost + } } return ctx @@ -172,12 +180,13 @@ func previousHostFromView(host *unifiedresources.HostView) models.Host { return models.Host{} } return models.Host{ - ID: host.ID(), - Hostname: host.Hostname(), - Status: string(host.Status()), - LinkedVMID: host.LinkedVMID(), - LastSeen: host.LastSeen(), - Disks: guestDisksFromReadStateView(host.Disks()), + ID: host.ID(), + Hostname: host.Hostname(), + Status: string(host.Status()), + LinkedVMID: host.LinkedVMID(), + LinkedContainerID: host.LinkedContainerID(), + LastSeen: host.LastSeen(), + Disks: guestDisksFromReadStateView(host.Disks()), Memory: models.Memory{ Used: host.MemoryUsed(), Total: host.MemoryTotal(), diff --git a/internal/monitoring/monitor_proxmox_pool_test.go b/internal/monitoring/monitor_proxmox_pool_test.go index 5e0ccca8c..87eab2990 100644 --- a/internal/monitoring/monitor_proxmox_pool_test.go +++ b/internal/monitoring/monitor_proxmox_pool_test.go @@ -58,6 +58,7 @@ func TestBuildContainerFromClusterResource_PreservesProxmoxPool(t *testing.T) { }, nil, map[int]bool{}, + nil, ) if !ok { t.Fatal("expected container to be built") diff --git a/internal/monitoring/monitor_pve_guest_lxc.go b/internal/monitoring/monitor_pve_guest_lxc.go index 7cbb4c379..6654435e9 100644 --- a/internal/monitoring/monitor_pve_guest_lxc.go +++ b/internal/monitoring/monitor_pve_guest_lxc.go @@ -44,6 +44,7 @@ func (m *Monitor) buildContainerFromClusterResource( res proxmox.ClusterResource, client PVEClientInterface, prevContainerIsOCI map[int]bool, + vmIDToHostAgent map[string]models.Host, ) (models.Container, VMMemoryRaw, string, time.Time, bool) { // Skip templates if configured if res.Template == 1 { @@ -88,6 +89,16 @@ func (m *Monitor) buildContainerFromClusterResource( ) memTotal, memUsed, memorySource, guestRaw := m.calculateLXCMemory(res) + agentHost, hasAgent := vmIDToHostAgent[guestID] + memTotal, memUsed, memorySource = preferLinkedAgentLXCMemory( + res.Status, + memTotal, + memUsed, + memorySource, + &guestRaw, + agentHost, + hasAgent, + ) memUsed, memorySource, _ = stabilizeGuestLowTrustMemory( m.previousGuestSnapshot(instanceName, "lxc", res.Node, res.VMID), res.Status, diff --git a/internal/monitoring/monitor_pve_guest_lxc_test.go b/internal/monitoring/monitor_pve_guest_lxc_test.go index 4b1e82bb6..93fd917be 100644 --- a/internal/monitoring/monitor_pve_guest_lxc_test.go +++ b/internal/monitoring/monitor_pve_guest_lxc_test.go @@ -140,7 +140,7 @@ func TestIssue1613LXCStatusLagDoesNotEraseNewerListingDiskWrite(t *testing.T) { } if _, _, _, _, ok := monitor.buildContainerFromClusterResource( - context.Background(), "cluster-a", resource, client, map[int]bool{}, + context.Background(), "cluster-a", resource, client, map[int]bool{}, nil, ); !ok { t.Fatal("expected first LXC sample") } @@ -150,7 +150,7 @@ func TestIssue1613LXCStatusLagDoesNotEraseNewerListingDiskWrite(t *testing.T) { client.containerStatus.ObservedAt = resource.ObservedAt.Add(time.Second) container, _, _, _, ok := monitor.buildContainerFromClusterResource( - context.Background(), "cluster-a", resource, client, map[int]bool{}, + context.Background(), "cluster-a", resource, client, map[int]bool{}, nil, ) if !ok { t.Fatal("expected second LXC sample") @@ -199,6 +199,7 @@ func TestBuildContainerFromClusterResource_UsesContainerStatusCountersForRates(t resource, client, map[int]bool{}, + nil, ); !ok { t.Fatal("expected first container sample to be built") } @@ -219,6 +220,7 @@ func TestBuildContainerFromClusterResource_UsesContainerStatusCountersForRates(t resource, client, map[int]bool{}, + nil, ) if !ok { t.Fatal("expected second container sample to be built") @@ -282,6 +284,7 @@ func TestIssue1634LXCMemoryFallsBackToClusterResourcesOnRealRRDShape(t *testing. resource, client, map[int]bool{}, + nil, ) if !ok { t.Fatal("expected container sample to be built") @@ -327,6 +330,7 @@ func TestIssue1634LXCMemoryStaysUnavailableWithoutListingValue(t *testing.T) { resource, client, map[int]bool{}, + nil, ) if !ok { t.Fatal("expected container sample to be built") @@ -443,6 +447,7 @@ func TestBuildContainerFromClusterResource_AppliesAgentLXCFilesystems(t *testing resource, client, map[int]bool{}, + nil, ) if !ok { t.Fatal("expected container to be built") @@ -460,3 +465,105 @@ func TestBuildContainerFromClusterResource_AppliesAgentLXCFilesystems(t *testing t.Fatalf("expected /data mount with real usage from agent data, got %+v", container.Disks) } } + +// Issue #2148: a Proxmox LXC with a correlated, online Pulse agent must report +// the agent's own memory sample instead of the cache-inclusive +// cluster/resources fallback, which can badly under-report shared-memory +// workloads (mirrors the #1962 VM behaviour). +func TestIssue2148LXCPrefersLinkedAgentMemoryOverClusterResources(t *testing.T) { + t.Parallel() + + const gib = 1024 * 1024 * 1024 + + client := &stubPVEClientLXCRRD{} + monitor := &Monitor{rateTracker: NewRateTracker()} + resource := proxmox.ClusterResource{ + Type: "lxc", + Node: "pve-a", + Name: "npu-ct", + Status: "running", + VMID: 100, + MaxMem: 24 * gib, + Mem: 3459743744, + } + guestID := makeGuestID("cluster-a", "pve-a", 100) + agentHost := models.Host{ + LinkedVMID: guestID, + Status: "online", + Memory: models.Memory{ + Total: 24 * gib, + Used: 11025571200, + Free: 24*gib - 11025571200, + }, + } + + container, _, memorySource, _, ok := monitor.buildContainerFromClusterResource( + context.Background(), + "cluster-a", + resource, + client, + map[int]bool{}, + map[string]models.Host{guestID: agentHost}, + ) + if !ok { + t.Fatal("expected container sample to be built") + } + if memorySource != "agent" { + t.Fatalf("memory source = %q, want agent", memorySource) + } + if container.Memory.Used != agentHost.Memory.Used { + t.Fatalf("memory used = %d, want linked agent value %d", container.Memory.Used, agentHost.Memory.Used) + } + if container.Memory.UsageUnavailable || !container.Memory.HasKnownUsage() { + t.Fatalf("expected usable agent-backed memory, got %+v", container.Memory) + } +} + +// An agent inside a container without lxcfs sees the host's /proc/meminfo, so a +// sample whose total does not match the container's configured limit must not +// replace the provider reading. +func TestIssue2148LXCRejectsAgentMemoryWithMismatchedTotal(t *testing.T) { + t.Parallel() + + const gib = 1024 * 1024 * 1024 + + client := &stubPVEClientLXCRRD{} + monitor := &Monitor{rateTracker: NewRateTracker()} + resource := proxmox.ClusterResource{ + Type: "lxc", + Node: "pve-a", + Name: "host-memory-leak-ct", + Status: "running", + VMID: 101, + MaxMem: 24 * gib, + Mem: 3459743744, + } + guestID := makeGuestID("cluster-a", "pve-a", 101) + agentHost := models.Host{ + LinkedVMID: guestID, + Status: "online", + Memory: models.Memory{ + Total: 128 * gib, + Used: 40 * gib, + Free: 88 * gib, + }, + } + + container, _, memorySource, _, ok := monitor.buildContainerFromClusterResource( + context.Background(), + "cluster-a", + resource, + client, + map[int]bool{}, + map[string]models.Host{guestID: agentHost}, + ) + if !ok { + t.Fatal("expected container sample to be built") + } + if CanonicalMemorySource(memorySource) != "cluster-resources" { + t.Fatalf("memory source = %q, want cluster-resources fallback", memorySource) + } + if container.Memory.Used != int64(resource.Mem) { + t.Fatalf("memory used = %d, want provider value %d", container.Memory.Used, resource.Mem) + } +} diff --git a/internal/monitoring/monitor_pve_guest_poll.go b/internal/monitoring/monitor_pve_guest_poll.go index 879035d7c..7c6f4e6bd 100644 --- a/internal/monitoring/monitor_pve_guest_poll.go +++ b/internal/monitoring/monitor_pve_guest_poll.go @@ -147,7 +147,7 @@ func (m *Monitor) collectGuestsFromClusterResources( collectionWG.Add(1) go func() { defer collectionWG.Done() - allContainers = m.collectClusterContainerResources(ctx, instanceName, containerResources, client, prevContainerIsOCI, prevContainerByID) + allContainers = m.collectClusterContainerResources(ctx, instanceName, containerResources, client, prevContainerIsOCI, prevContainerByID, vmIDToHostAgent) }() } if len(vmResources) > 0 { @@ -253,6 +253,7 @@ func (m *Monitor) collectClusterContainerResources( client PVEClientInterface, prevContainerIsOCI map[int]bool, prevContainerByID map[string]models.Container, + vmIDToHostAgent map[string]models.Host, ) []models.Container { orderedContainers := make([]models.Container, len(resources)) orderedOK := make([]bool, len(resources)) @@ -263,13 +264,13 @@ func (m *Monitor) collectClusterContainerResources( var container models.Container var ok bool ran := m.runGuestAgentVMWork(ctx, func(workCtx context.Context) { - container, ok = m.handleClusterContainerResource(workCtx, instanceName, entry.resource, entry.guestID, client, prevContainerIsOCI) + container, ok = m.handleClusterContainerResource(workCtx, instanceName, entry.resource, entry.guestID, client, prevContainerIsOCI, vmIDToHostAgent) }) if !ran { // The enrichment budget is exhausted, but cluster/resources is still // authoritative inventory. A canceled context makes remote detail // calls fail immediately while the builder retains the base row. - container, ok = m.handleClusterContainerResource(ctx, instanceName, entry.resource, entry.guestID, client, prevContainerIsOCI) + container, ok = m.handleClusterContainerResource(ctx, instanceName, entry.resource, entry.guestID, client, prevContainerIsOCI, vmIDToHostAgent) } if ctx.Err() != nil { previous := prevContainerByID[entry.guestID] @@ -415,8 +416,9 @@ func (m *Monitor) handleClusterContainerResource( guestID string, client PVEClientInterface, prevContainerIsOCI map[int]bool, + vmIDToHostAgent map[string]models.Host, ) (models.Container, bool) { - container, guestRaw, memorySource, sampleTime, ok := m.buildContainerFromClusterResource(ctx, instanceName, res, client, prevContainerIsOCI) + container, guestRaw, memorySource, sampleTime, ok := m.buildContainerFromClusterResource(ctx, instanceName, res, client, prevContainerIsOCI, vmIDToHostAgent) if !ok { return models.Container{}, false } diff --git a/internal/unifiedresources/views.go b/internal/unifiedresources/views.go index ff0a0b79a..d6440bdff 100644 --- a/internal/unifiedresources/views.go +++ b/internal/unifiedresources/views.go @@ -78,13 +78,20 @@ func NewVMView(r *Resource) VMView { return VMView{r: r} } // LinkedAgentMemory returns the agent's own sample, not the platform-priority // merged metric. Freshness belongs to the agent source, not the VM row. func (v VMView) LinkedAgentMemory() (models.Memory, bool) { - if v.r == nil || v.r.Agent == nil || v.r.Agent.Stale || v.r.Agent.Memory == nil { + return linkedAgentMemoryFromResource(v.r) +} + +// linkedAgentMemoryFromResource reads the agent-owned memory sample attached to +// a merged guest resource. Correlation removes the standalone host row, so the +// guest view is the only place the agent sample survives (#1962, #2148). +func linkedAgentMemoryFromResource(r *Resource) (models.Memory, bool) { + if r == nil || r.Agent == nil || r.Agent.Stale || r.Agent.Memory == nil { return models.Memory{}, false } - if status, ok := v.r.SourceStatus[SourceAgent]; !ok || status.Status != "online" { + if status, ok := r.SourceStatus[SourceAgent]; !ok || status.Status != "online" { return models.Memory{}, false } - m := v.r.Agent.Memory + m := r.Agent.Memory memory := models.Memory{Total: m.Total, Used: m.Used, Free: m.Free, Cache: m.Cache, UsageUnavailable: m.UsageUnavailable} if m.Total <= 0 { return models.Memory{}, false @@ -397,6 +404,14 @@ type ContainerView struct{ r *Resource } func NewContainerView(r *Resource) ContainerView { return ContainerView{r: r} } +// LinkedAgentMemory returns the agent's own sample for a system container whose +// merged resource also carries an online Pulse agent source. Proxmox +// cluster/resources memory is a low-trust fallback that can badly misreport a +// container's real footprint, so the linked agent sample is preferred (#2148). +func (v ContainerView) LinkedAgentMemory() (models.Memory, bool) { + return linkedAgentMemoryFromResource(v.r) +} + func (v ContainerView) String() string { return fmt.Sprintf("ContainerView(%s, %q)", v.ID(), v.Name()) } func (v ContainerView) ID() string { diff --git a/internal/unifiedresources/views_test.go b/internal/unifiedresources/views_test.go index f01847ae1..2ec871204 100644 --- a/internal/unifiedresources/views_test.go +++ b/internal/unifiedresources/views_test.go @@ -2267,3 +2267,34 @@ func TestVMViewLinkedAgentMemory(t *testing.T) { t.Fatal("nil view has memory") } } + +func TestContainerViewLinkedAgentMemory(t *testing.T) { + for _, tc := range []struct { + name string + status string + stale bool + total, used int64 + unavailable bool + want bool + }{ + {name: "live", status: "online", total: 24 << 30, used: 11 << 30, want: true}, + {name: "offline", status: "offline", total: 24 << 30, used: 11 << 30}, + {name: "stale agent", status: "online", stale: true, total: 24 << 30, used: 11 << 30}, + {name: "over capacity", status: "online", total: 24 << 30, used: 25 << 30}, + {name: "unavailable", status: "online", total: 24 << 30, used: 11 << 30, unavailable: true}, + } { + t.Run(tc.name, func(t *testing.T) { + r := &Resource{Status: StatusOnline, Agent: &AgentData{Stale: tc.stale, Memory: &AgentMemoryMeta{Total: tc.total, Used: tc.used, UsageUnavailable: tc.unavailable}}, SourceStatus: map[DataSource]SourceStatus{SourceAgent: {Status: tc.status}}, Metrics: &ResourceMetrics{Memory: &MetricValue{Percent: 100}}} + memory, ok := NewContainerView(r).LinkedAgentMemory() + if ok != tc.want { + t.Fatalf("valid=%v want=%v memory=%+v", ok, tc.want, memory) + } + if ok && memory.Used != tc.used { + t.Fatalf("not agent memory: %+v", memory) + } + }) + } + if _, ok := (ContainerView{}).LinkedAgentMemory(); ok { + t.Fatal("nil view has memory") + } +}