From 4ceb3ad433a39d86339606c260422ff38bcfa384 Mon Sep 17 00:00:00 2001 From: rcourtman Date: Wed, 1 Apr 2026 18:38:42 +0100 Subject: [PATCH] Preserve guest-agent continuity across transient PVE status failures --- .../v6/internal/subsystems/monitoring.md | 6 + .../monitoring/canonical_guardrails_test.go | 49 ++ internal/monitoring/guest_agent_evidence.go | 44 ++ .../monitoring/guest_agent_evidence_test.go | 40 ++ internal/monitoring/guest_metadata.go | 116 +++- internal/monitoring/guest_metadata_test.go | 522 ++++++++++++++++++ .../memory_trust_characterization_test.go | 2 +- .../monitoring/monitor_extra_coverage_test.go | 386 ++++++++++++- internal/monitoring/monitor_polling_vm.go | 40 +- internal/monitoring/monitor_previous_state.go | 34 +- .../monitoring/monitor_proxmox_pool_test.go | 1 + .../monitoring/monitor_pve_guest_builders.go | 48 +- internal/monitoring/monitor_pve_guest_poll.go | 20 +- 13 files changed, 1255 insertions(+), 53 deletions(-) create mode 100644 internal/monitoring/guest_agent_evidence.go create mode 100644 internal/monitoring/guest_agent_evidence_test.go diff --git a/docs/release-control/v6/internal/subsystems/monitoring.md b/docs/release-control/v6/internal/subsystems/monitoring.md index 9c8ab585f..4374e22ce 100644 --- a/docs/release-control/v6/internal/subsystems/monitoring.md +++ b/docs/release-control/v6/internal/subsystems/monitoring.md @@ -576,3 +576,9 @@ metrics into the existing shared `agent` and `vm` history stores. Pulse must not add a VMware-only charts cache, host history model, or VM metrics store just because vSphere performance collection uses a different API family from inventory and alarm reads. +That same monitoring boundary now also owns Proxmox guest-agent continuity +when `/status` is transiently missing. Recent guest-agent evidence and the +shared guest metadata cache must keep VM network and identity metadata alive +long enough to survive short Proxmox status failures, while incomplete +guest-agent metadata stays on a short retry cadence instead of freezing +partial VM summary data for minutes. diff --git a/internal/monitoring/canonical_guardrails_test.go b/internal/monitoring/canonical_guardrails_test.go index 7aa5b05b9..1849ae472 100644 --- a/internal/monitoring/canonical_guardrails_test.go +++ b/internal/monitoring/canonical_guardrails_test.go @@ -381,6 +381,55 @@ func TestProxmoxGuestPollersCarryPoolIntoCanonicalModels(t *testing.T) { } } +func TestProxmoxGuestAgentContinuityUsesCanonicalEvidenceAndRetryPaths(t *testing.T) { + requiredSnippets := map[string][]string{ + "guest_agent_evidence.go": { + "func hasRecentGuestAgentEvidence(prev *models.VM, now time.Time) bool {", + `if prev == nil || prev.Type != "qemu" {`, + "func shouldQueryGuestAgent(vmStatus *proxmox.VMStatus, prev *models.VM, now time.Time) bool {", + "return hasRecentGuestAgentEvidence(prev, now)", + }, + "guest_metadata.go": { + "guestMetadataEmptyTTL = 30 * time.Second", + "func guestMetadataCacheEntryTTL(entry guestMetadataCacheEntry) time.Duration {", + "if guestMetadataCacheHasCompleteNetworkData(entry) {", + "return guestMetadataEmptyTTL", + "func (m *Monitor) hasRecentGuestMetadataEvidence(instanceName, nodeName string, vmid int, now time.Time) bool {", + "func (m *Monitor) scheduleGuestMetadataFetchForEntry(key string, now time.Time, entry guestMetadataCacheEntry) {", + }, + "monitor_previous_state.go": { + "vmsByID map[string]models.VM", + "ctx.vmsByID[modelVM.ID] = modelVM", + "guestID := makeGuestID(modelVM.Instance, modelVM.Node, modelVM.VMID)", + "ctx.vmsByID[guestID] = modelVM", + }, + "monitor_pve_guest_builders.go": { + "guestAgentAvailable := shouldQueryGuestAgent(state.detailedStatus, prevVM, now) ||", + "m.hasRecentGuestMetadataEvidence(instanceName, res.Node, res.VMID, now)", + "if guestAgentAvailable && state.detailedStatus == nil {", + }, + "monitor_polling_vm.go": { + "prevVMByID := prevGuests.vmsByID", + "guestAgentAvailable := vm.Status == \"running\" &&", + "m.hasRecentGuestMetadataEvidence(instanceName, n.Node, vm.VMID, now)", + "if guestAgentAvailable && diskTotal > 0 {", + }, + } + + for file, snippets := range requiredSnippets { + data, err := os.ReadFile(file) + if err != nil { + t.Fatalf("failed to read %s: %v", file, err) + } + source := string(data) + for _, snippet := range snippets { + if !strings.Contains(source, snippet) { + t.Fatalf("%s must contain %q", file, snippet) + } + } + } +} + func TestUnifiedAppContainerMetricsUseCanonicalGuestHistoryPath(t *testing.T) { data, err := os.ReadFile("monitor.go") if err != nil { diff --git a/internal/monitoring/guest_agent_evidence.go b/internal/monitoring/guest_agent_evidence.go new file mode 100644 index 000000000..61e173b19 --- /dev/null +++ b/internal/monitoring/guest_agent_evidence.go @@ -0,0 +1,44 @@ +package monitoring + +import ( + "strings" + "time" + + "github.com/rcourtman/pulse-go-rewrite/internal/models" + "github.com/rcourtman/pulse-go-rewrite/pkg/proxmox" +) + +const recentGuestAgentEvidenceMaxAge = 10 * time.Minute + +func hasRecentGuestAgentEvidence(prev *models.VM, now time.Time) bool { + // Previous unified read-state snapshots normalize guest status to resource health + // (for example "online"), so continuity depends on recent guest-agent evidence + // rather than on replaying an exact prior runtime status string. + if prev == nil || prev.Type != "qemu" { + return false + } + if prev.LastSeen.IsZero() || now.Sub(prev.LastSeen) > recentGuestAgentEvidenceMaxAge { + return false + } + + if prev.AgentVersion != "" || + len(prev.IPAddresses) > 0 || + len(prev.NetworkInterfaces) > 0 || + prev.OSName != "" || + prev.OSVersion != "" { + return true + } + + if len(prev.Disks) > 0 && !strings.HasPrefix(prev.DiskStatusReason, "prev-") { + return true + } + + return false +} + +func shouldQueryGuestAgent(vmStatus *proxmox.VMStatus, prev *models.VM, now time.Time) bool { + if vmStatus != nil { + return vmStatus.Agent.Value > 0 + } + return hasRecentGuestAgentEvidence(prev, now) +} diff --git a/internal/monitoring/guest_agent_evidence_test.go b/internal/monitoring/guest_agent_evidence_test.go new file mode 100644 index 000000000..61d59f4f7 --- /dev/null +++ b/internal/monitoring/guest_agent_evidence_test.go @@ -0,0 +1,40 @@ +package monitoring + +import ( + "testing" + "time" + + "github.com/rcourtman/pulse-go-rewrite/internal/models" +) + +func TestHasRecentGuestAgentEvidenceAcceptsUnifiedReadStateStatus(t *testing.T) { + now := time.Now() + + prev := &models.VM{ + Type: "qemu", + Status: "online", + IPAddresses: []string{"192.168.1.50"}, + AgentVersion: "8.1.0", + LastSeen: now.Add(-time.Minute), + } + + if !hasRecentGuestAgentEvidence(prev, now) { + t.Fatal("expected unified read-state VM snapshot to count as recent guest-agent evidence") + } +} + +func TestHasRecentGuestAgentEvidenceRejectsStaleSnapshots(t *testing.T) { + now := time.Now() + + prev := &models.VM{ + Type: "qemu", + Status: "online", + IPAddresses: []string{"192.168.1.50"}, + AgentVersion: "8.1.0", + LastSeen: now.Add(-recentGuestAgentEvidenceMaxAge - time.Second), + } + + if hasRecentGuestAgentEvidence(prev, now) { + t.Fatal("expected stale guest-agent evidence to be rejected") + } +} diff --git a/internal/monitoring/guest_metadata.go b/internal/monitoring/guest_metadata.go index 8c320cf9e..a4cfa23bc 100644 --- a/internal/monitoring/guest_metadata.go +++ b/internal/monitoring/guest_metadata.go @@ -18,6 +18,7 @@ import ( const ( guestMetadataCacheTTL = 5 * time.Minute + guestMetadataEmptyTTL = 30 * time.Second defaultGuestMetadataHold = 15 * time.Second // Guest agent timeout defaults (configurable via environment variables) @@ -45,6 +46,48 @@ type guestMetadataCacheEntry struct { osInfoSkip bool // Skip OS info calls after repeated failures (refs #692) } +func guestMetadataCacheHasUsefulData(entry guestMetadataCacheEntry) bool { + return len(entry.ipAddresses) > 0 || + len(entry.networkInterfaces) > 0 || + entry.osName != "" || + entry.osVersion != "" || + entry.agentVersion != "" +} + +func guestMetadataCacheHasCompleteNetworkData(entry guestMetadataCacheEntry) bool { + for _, iface := range entry.networkInterfaces { + if strings.TrimSpace(iface.Name) != "" || strings.TrimSpace(iface.MAC) != "" { + return true + } + } + return false +} + +func guestMetadataCacheEntryTTL(entry guestMetadataCacheEntry) time.Duration { + // Treat identity-only and IP-only metadata as incomplete so VMs that answered + // guest-info/version or partial network calls but not full interface inventory + // are retried soon instead of freezing incomplete VM Summary data for minutes. + if guestMetadataCacheHasCompleteNetworkData(entry) { + return guestMetadataCacheTTL + } + return guestMetadataEmptyTTL +} + +func (m *Monitor) hasRecentGuestMetadataEvidence(instanceName, nodeName string, vmid int, now time.Time) bool { + if m == nil { + return false + } + + key := guestMetadataCacheKey(instanceName, nodeName, vmid) + m.guestMetadataMu.RLock() + entry, ok := m.guestMetadataCache[key] + m.guestMetadataMu.RUnlock() + if !ok || entry.fetchedAt.IsZero() || !guestMetadataCacheHasUsefulData(entry) { + return false + } + return now.Sub(entry.fetchedAt) < guestMetadataCacheEntryTTL(entry) +} + func (m *Monitor) tryReserveGuestMetadataFetch(key string, now time.Time) bool { if m == nil { return false @@ -80,6 +123,19 @@ func (m *Monitor) scheduleNextGuestMetadataFetch(key string, now time.Time) { m.guestMetadataLimiterMu.Unlock() } +func (m *Monitor) scheduleGuestMetadataFetchForEntry(key string, now time.Time, entry guestMetadataCacheEntry) { + if m == nil { + return + } + if guestMetadataCacheEntryTTL(entry) == guestMetadataEmptyTTL { + m.guestMetadataLimiterMu.Lock() + m.guestMetadataLimiter[key] = now.Add(guestMetadataEmptyTTL) + m.guestMetadataLimiterMu.Unlock() + return + } + m.scheduleNextGuestMetadataFetch(key, now) +} + func (m *Monitor) deferGuestMetadataRetry(key string, now time.Time) { if m == nil { return @@ -164,17 +220,7 @@ func (m *Monitor) retryGuestAgentCall(ctx context.Context, timeout time.Duration return nil, lastErr } -func (m *Monitor) fetchGuestAgentMetadata(ctx context.Context, client PVEClientInterface, instanceName, nodeName, vmName string, vmid int, vmStatus *proxmox.VMStatus) ([]string, []models.GuestNetworkInterface, string, string, string) { - if vmStatus == nil || client == nil { - m.clearGuestMetadataCache(instanceName, nodeName, vmid) - return nil, nil, "", "", "" - } - - if vmStatus.Agent.Value <= 0 { - m.clearGuestMetadataCache(instanceName, nodeName, vmid) - return nil, nil, "", "", "" - } - +func (m *Monitor) fetchGuestAgentMetadata(ctx context.Context, client PVEClientInterface, instanceName, nodeName, vmName string, vmid int, vmStatus *proxmox.VMStatus, allowWithoutStatus bool) ([]string, []models.GuestNetworkInterface, string, string, string) { key := guestMetadataCacheKey(instanceName, nodeName, vmid) now := time.Now() @@ -182,11 +228,20 @@ func (m *Monitor) fetchGuestAgentMetadata(ctx context.Context, client PVEClientI cached, ok := m.guestMetadataCache[key] m.guestMetadataMu.RUnlock() - if ok && now.Sub(cached.fetchedAt) < guestMetadataCacheTTL { + agentAvailable := client != nil && ((vmStatus != nil && vmStatus.Agent.Value > 0) || allowWithoutStatus) + if !agentAvailable { + if ok && now.Sub(cached.fetchedAt) < guestMetadataCacheEntryTTL(cached) { + return cloneStringSlice(cached.ipAddresses), cloneGuestNetworkInterfaces(cached.networkInterfaces), cached.osName, cached.osVersion, cached.agentVersion + } + m.clearGuestMetadataCache(instanceName, nodeName, vmid) + return nil, nil, "", "", "" + } + + if ok && now.Sub(cached.fetchedAt) < guestMetadataCacheEntryTTL(cached) { return cloneStringSlice(cached.ipAddresses), cloneGuestNetworkInterfaces(cached.networkInterfaces), cached.osName, cached.osVersion, cached.agentVersion } - needsFetch := !ok || now.Sub(cached.fetchedAt) >= guestMetadataCacheTTL + needsFetch := !ok || now.Sub(cached.fetchedAt) >= guestMetadataCacheEntryTTL(cached) if !needsFetch { return cloneStringSlice(cached.ipAddresses), cloneGuestNetworkInterfaces(cached.networkInterfaces), cached.osName, cached.osVersion, cached.agentVersion } @@ -212,9 +267,6 @@ func (m *Monitor) fetchGuestAgentMetadata(ctx context.Context, client PVEClientI return ipAddresses, networkIfaces, osName, osVersion, agentVersion } defer m.releaseGuestMetadataSlot() - defer func() { - m.scheduleNextGuestMetadataFetch(key, time.Now()) - }() } // Network interfaces with configurable timeout and retry (refs #592) @@ -229,10 +281,24 @@ func (m *Monitor) fetchGuestAgentMetadata(ctx context.Context, client PVEClientI Err(err). Msg("Guest agent network interfaces unavailable") } else if ifaces, ok := interfaces.([]proxmox.VMNetworkInterface); ok && len(ifaces) > 0 { - ipAddresses, networkIfaces = processGuestNetworkInterfaces(ifaces) + processedIPs, processedIfaces := processGuestNetworkInterfaces(ifaces) + if len(processedIPs) > 0 || len(processedIfaces) > 0 { + ipAddresses, networkIfaces = processedIPs, processedIfaces + } else if len(cached.ipAddresses) == 0 && len(cached.networkInterfaces) == 0 { + ipAddresses = nil + networkIfaces = nil + } else { + log.Debug(). + Str("instance", instanceName). + Str("vm", vmName). + Int("vmid", vmid). + Msg("Guest agent returned empty network metadata; preserving last known interfaces") + } } else { - ipAddresses = nil - networkIfaces = nil + if len(cached.ipAddresses) == 0 && len(cached.networkInterfaces) == 0 { + ipAddresses = nil + networkIfaces = nil + } } // OS info with configurable timeout and retry (refs #592) @@ -275,10 +341,13 @@ func (m *Monitor) fetchGuestAgentMetadata(ctx context.Context, client PVEClientI } } } else if agentInfo, ok := agentInfoRaw.(map[string]interface{}); ok && len(agentInfo) > 0 { - osName, osVersion = extractGuestOSInfo(agentInfo) + extractedOSName, extractedOSVersion := extractGuestOSInfo(agentInfo) + if extractedOSName != "" || extractedOSVersion != "" { + osName, osVersion = extractedOSName, extractedOSVersion + } osInfoFailureCount = 0 // Reset on success osInfoSkip = false - } else { + } else if cached.osName == "" && cached.osVersion == "" { osName = "" osVersion = "" } @@ -304,7 +373,7 @@ func (m *Monitor) fetchGuestAgentMetadata(ctx context.Context, client PVEClientI Msg("Guest agent version unavailable") } else if version, ok := versionRaw.(string); ok && version != "" { agentVersion = version - } else { + } else if cached.agentVersion == "" { agentVersion = "" } @@ -325,6 +394,9 @@ func (m *Monitor) fetchGuestAgentMetadata(ctx context.Context, client PVEClientI } m.guestMetadataCache[key] = entry m.guestMetadataMu.Unlock() + if reserved { + m.scheduleGuestMetadataFetchForEntry(key, time.Now(), entry) + } return ipAddresses, networkIfaces, osName, osVersion, agentVersion } diff --git a/internal/monitoring/guest_metadata_test.go b/internal/monitoring/guest_metadata_test.go index 9082f56aa..e9a4f96f3 100644 --- a/internal/monitoring/guest_metadata_test.go +++ b/internal/monitoring/guest_metadata_test.go @@ -11,6 +11,120 @@ import ( "github.com/rcourtman/pulse-go-rewrite/pkg/proxmox" ) +type emptyGuestMetadataClient struct { + mockPVEClient +} + +func (emptyGuestMetadataClient) GetVMNetworkInterfaces(ctx context.Context, node string, vmid int) ([]proxmox.VMNetworkInterface, error) { + return []proxmox.VMNetworkInterface{}, nil +} + +func (emptyGuestMetadataClient) GetVMAgentInfo(ctx context.Context, node string, vmid int) (map[string]interface{}, error) { + return map[string]interface{}{}, nil +} + +func (emptyGuestMetadataClient) GetVMAgentVersion(ctx context.Context, node string, vmid int) (string, error) { + return "", nil +} + +type emptyThenPopulatedGuestMetadataClient struct { + mockPVEClient + networkCalls int +} + +func (c *emptyThenPopulatedGuestMetadataClient) GetVMNetworkInterfaces(ctx context.Context, node string, vmid int) ([]proxmox.VMNetworkInterface, error) { + c.networkCalls++ + if c.networkCalls == 1 { + return []proxmox.VMNetworkInterface{}, nil + } + return []proxmox.VMNetworkInterface{ + { + Name: "Ethernet0", + HardwareAddr: "00:11:22:33:44:55", + IPAddresses: []proxmox.VMIPAddress{ + {Address: "192.168.1.50", Prefix: 24}, + }, + }, + }, nil +} + +func (*emptyThenPopulatedGuestMetadataClient) GetVMAgentInfo(ctx context.Context, node string, vmid int) (map[string]interface{}, error) { + return map[string]interface{}{}, nil +} + +func (*emptyThenPopulatedGuestMetadataClient) GetVMAgentVersion(ctx context.Context, node string, vmid int) (string, error) { + return "", nil +} + +type identityThenNetworkGuestMetadataClient struct { + mockPVEClient + networkCalls int +} + +func (c *identityThenNetworkGuestMetadataClient) GetVMNetworkInterfaces(ctx context.Context, node string, vmid int) ([]proxmox.VMNetworkInterface, error) { + c.networkCalls++ + if c.networkCalls == 1 { + return []proxmox.VMNetworkInterface{}, nil + } + return []proxmox.VMNetworkInterface{ + { + Name: "Ethernet0", + HardwareAddr: "00:11:22:33:44:55", + IPAddresses: []proxmox.VMIPAddress{ + {Address: "192.168.1.50", Prefix: 24}, + }, + }, + }, nil +} + +func (*identityThenNetworkGuestMetadataClient) GetVMAgentInfo(ctx context.Context, node string, vmid int) (map[string]interface{}, error) { + return map[string]interface{}{ + "result": map[string]interface{}{ + "name": "Microsoft Windows", + "version": "11", + }, + }, nil +} + +func (*identityThenNetworkGuestMetadataClient) GetVMAgentVersion(ctx context.Context, node string, vmid int) (string, error) { + return "8.2.2", nil +} + +type ipOnlyThenNetworkGuestMetadataClient struct { + mockPVEClient + networkCalls int +} + +func (c *ipOnlyThenNetworkGuestMetadataClient) GetVMNetworkInterfaces(ctx context.Context, node string, vmid int) ([]proxmox.VMNetworkInterface, error) { + c.networkCalls++ + if c.networkCalls == 1 { + return []proxmox.VMNetworkInterface{ + { + IPAddresses: []proxmox.VMIPAddress{ + {Address: "192.168.1.60", Prefix: 24}, + }, + }, + }, nil + } + return []proxmox.VMNetworkInterface{ + { + Name: "Ethernet0", + HardwareAddr: "00:11:22:33:44:66", + IPAddresses: []proxmox.VMIPAddress{ + {Address: "192.168.1.60", Prefix: 24}, + }, + }, + }, nil +} + +func (*ipOnlyThenNetworkGuestMetadataClient) GetVMAgentInfo(ctx context.Context, node string, vmid int) (map[string]interface{}, error) { + return map[string]interface{}{}, nil +} + +func (*ipOnlyThenNetworkGuestMetadataClient) GetVMAgentVersion(ctx context.Context, node string, vmid int) (string, error) { + return "", nil +} + func TestGuestMetadataCacheKey(t *testing.T) { t.Parallel() @@ -771,6 +885,414 @@ func TestRetryGuestAgentCall(t *testing.T) { }) } +func TestGuestMetadataCacheEntryTTL(t *testing.T) { + t.Parallel() + + tests := []struct { + name string + entry guestMetadataCacheEntry + want time.Duration + }{ + { + name: "empty entry retries quickly", + entry: guestMetadataCacheEntry{}, + want: guestMetadataEmptyTTL, + }, + { + name: "network metadata uses full ttl", + entry: guestMetadataCacheEntry{ + ipAddresses: []string{"192.168.1.10"}, + networkInterfaces: []models.GuestNetworkInterface{ + {Name: "eth0", Addresses: []string{"192.168.1.10"}}, + }, + }, + want: guestMetadataCacheTTL, + }, + { + name: "ip-only metadata retries quickly", + entry: guestMetadataCacheEntry{ + ipAddresses: []string{"192.168.1.10"}, + }, + want: guestMetadataEmptyTTL, + }, + { + name: "identity-only metadata retries quickly", + entry: guestMetadataCacheEntry{ + osName: "Ubuntu", + }, + want: guestMetadataEmptyTTL, + }, + } + + for _, tc := range tests { + tc := tc + t.Run(tc.name, func(t *testing.T) { + t.Parallel() + if got := guestMetadataCacheEntryTTL(tc.entry); got != tc.want { + t.Fatalf("guestMetadataCacheEntryTTL() = %s, want %s", got, tc.want) + } + }) + } +} + +func TestScheduleGuestMetadataFetchForEntry(t *testing.T) { + t.Parallel() + + t.Run("incomplete metadata retries sooner", func(t *testing.T) { + t.Parallel() + + m := &Monitor{ + guestMetadataLimiter: make(map[string]time.Time), + guestMetadataMinRefresh: 5 * time.Minute, + guestMetadataRefreshJitter: 0, + } + now := time.Now() + entry := guestMetadataCacheEntry{osName: "Ubuntu"} + + m.scheduleGuestMetadataFetchForEntry("vm-1", now, entry) + + next := m.guestMetadataLimiter["vm-1"] + expectedNext := now.Add(guestMetadataEmptyTTL) + if next.Sub(expectedNext) > time.Millisecond || expectedNext.Sub(next) > time.Millisecond { + t.Fatalf("expected next time ~%v, got %v", expectedNext, next) + } + }) + + t.Run("complete network metadata uses normal refresh cadence", func(t *testing.T) { + t.Parallel() + + m := &Monitor{ + guestMetadataLimiter: make(map[string]time.Time), + guestMetadataMinRefresh: 5 * time.Minute, + guestMetadataRefreshJitter: 0, + } + now := time.Now() + entry := guestMetadataCacheEntry{ + networkInterfaces: []models.GuestNetworkInterface{ + {Name: "eth0", Addresses: []string{"192.168.1.10"}}, + }, + } + + m.scheduleGuestMetadataFetchForEntry("vm-1", now, entry) + + next := m.guestMetadataLimiter["vm-1"] + expectedNext := now.Add(5 * time.Minute) + if next.Sub(expectedNext) > time.Millisecond || expectedNext.Sub(next) > time.Millisecond { + t.Fatalf("expected next time ~%v, got %v", expectedNext, next) + } + }) +} + +func TestFetchGuestAgentMetadataPreservesCachedValuesOnEmptyResponses(t *testing.T) { + t.Parallel() + + key := guestMetadataCacheKey("pve", "node1", 100) + cachedFetchedAt := time.Now().Add(-guestMetadataCacheTTL - time.Minute) + monitor := &Monitor{ + guestMetadataCache: map[string]guestMetadataCacheEntry{ + key: { + ipAddresses: []string{"192.168.1.10"}, + networkInterfaces: []models.GuestNetworkInterface{ + {Name: "eth0", MAC: "00:11:22:33:44:55", Addresses: []string{"192.168.1.10"}}, + }, + osName: "Ubuntu", + osVersion: "24.04", + agentVersion: "8.2.0", + fetchedAt: cachedFetchedAt, + }, + }, + guestMetadataLimiter: make(map[string]time.Time), + } + + gotIPs, gotIfaces, gotOSName, gotOSVersion, gotAgentVersion := monitor.fetchGuestAgentMetadata( + context.Background(), + &emptyGuestMetadataClient{}, + "pve", + "node1", + "vm100", + 100, + &proxmox.VMStatus{Agent: proxmox.VMAgentField{Value: 1}}, + false, + ) + + if len(gotIPs) != 1 || gotIPs[0] != "192.168.1.10" { + t.Fatalf("expected cached IPs to be preserved, got %#v", gotIPs) + } + if len(gotIfaces) != 1 || gotIfaces[0].Name != "eth0" { + t.Fatalf("expected cached interfaces to be preserved, got %#v", gotIfaces) + } + if gotOSName != "Ubuntu" || gotOSVersion != "24.04" { + t.Fatalf("expected cached OS info to be preserved, got %q %q", gotOSName, gotOSVersion) + } + if gotAgentVersion != "8.2.0" { + t.Fatalf("expected cached agent version to be preserved, got %q", gotAgentVersion) + } + + updated := monitor.guestMetadataCache[key] + if len(updated.ipAddresses) != 1 || updated.ipAddresses[0] != "192.168.1.10" { + t.Fatalf("expected cache IPs to remain populated, got %#v", updated.ipAddresses) + } + if len(updated.networkInterfaces) != 1 || updated.networkInterfaces[0].Name != "eth0" { + t.Fatalf("expected cache interfaces to remain populated, got %#v", updated.networkInterfaces) + } + if updated.fetchedAt.Before(cachedFetchedAt) { + t.Fatalf("expected cache timestamp to be refreshed, old=%v new=%v", cachedFetchedAt, updated.fetchedAt) + } +} + +func TestFetchGuestAgentMetadataPreservesFreshCacheWhenAgentTemporarilyUnavailable(t *testing.T) { + t.Parallel() + + key := guestMetadataCacheKey("pve", "node1", 100) + cachedFetchedAt := time.Now().Add(-time.Minute) + + newMonitor := func() *Monitor { + return &Monitor{ + guestMetadataCache: map[string]guestMetadataCacheEntry{ + key: { + ipAddresses: []string{"192.168.1.10"}, + networkInterfaces: []models.GuestNetworkInterface{ + {Name: "eth0", MAC: "00:11:22:33:44:55", Addresses: []string{"192.168.1.10"}}, + }, + osName: "Ubuntu", + osVersion: "24.04", + agentVersion: "8.2.0", + fetchedAt: cachedFetchedAt, + }, + }, + guestMetadataLimiter: make(map[string]time.Time), + } + } + + assertCachePreserved := func(t *testing.T, monitor *Monitor, gotIPs []string, gotIfaces []models.GuestNetworkInterface, gotOSName, gotOSVersion, gotAgentVersion string) { + t.Helper() + + if len(gotIPs) != 1 || gotIPs[0] != "192.168.1.10" { + t.Fatalf("expected cached IPs to be preserved, got %#v", gotIPs) + } + if len(gotIfaces) != 1 || gotIfaces[0].Name != "eth0" { + t.Fatalf("expected cached interfaces to be preserved, got %#v", gotIfaces) + } + if gotOSName != "Ubuntu" || gotOSVersion != "24.04" { + t.Fatalf("expected cached OS info to be preserved, got %q %q", gotOSName, gotOSVersion) + } + if gotAgentVersion != "8.2.0" { + t.Fatalf("expected cached agent version to be preserved, got %q", gotAgentVersion) + } + + entry, ok := monitor.guestMetadataCache[key] + if !ok { + t.Fatal("expected guest metadata cache entry to remain populated") + } + if len(entry.networkInterfaces) != 1 || entry.networkInterfaces[0].Name != "eth0" { + t.Fatalf("expected cached interfaces to remain populated, got %#v", entry.networkInterfaces) + } + } + + t.Run("nil vm status keeps fresh cache", func(t *testing.T) { + t.Parallel() + + monitor := newMonitor() + gotIPs, gotIfaces, gotOSName, gotOSVersion, gotAgentVersion := monitor.fetchGuestAgentMetadata( + context.Background(), + &emptyGuestMetadataClient{}, + "pve", + "node1", + "vm100", + 100, + nil, + false, + ) + + assertCachePreserved(t, monitor, gotIPs, gotIfaces, gotOSName, gotOSVersion, gotAgentVersion) + }) + + t.Run("agent temporarily unavailable keeps fresh cache", func(t *testing.T) { + t.Parallel() + + monitor := newMonitor() + gotIPs, gotIfaces, gotOSName, gotOSVersion, gotAgentVersion := monitor.fetchGuestAgentMetadata( + context.Background(), + &emptyGuestMetadataClient{}, + "pve", + "node1", + "vm100", + 100, + &proxmox.VMStatus{Agent: proxmox.VMAgentField{Value: 0}}, + false, + ) + + assertCachePreserved(t, monitor, gotIPs, gotIfaces, gotOSName, gotOSVersion, gotAgentVersion) + }) +} + +func TestFetchGuestAgentMetadataRetriesIdentityOnlyCacheSooner(t *testing.T) { + t.Parallel() + + client := &identityThenNetworkGuestMetadataClient{} + monitor := &Monitor{ + guestMetadataCache: make(map[string]guestMetadataCacheEntry), + guestMetadataLimiter: make(map[string]time.Time), + } + + status := &proxmox.VMStatus{Agent: proxmox.VMAgentField{Value: 1}} + + firstIPs, firstIfaces, firstOSName, firstOSVersion, firstAgentVersion := monitor.fetchGuestAgentMetadata( + context.Background(), + client, + "pve", + "node1", + "vm100", + 100, + status, + false, + ) + if len(firstIPs) != 0 || len(firstIfaces) != 0 { + t.Fatalf("expected first fetch to be missing network metadata, got ips=%#v ifaces=%#v", firstIPs, firstIfaces) + } + if firstOSName == "" || firstOSVersion == "" || firstAgentVersion == "" { + t.Fatalf("expected identity metadata on first fetch, got os=%q version=%q agent=%q", firstOSName, firstOSVersion, firstAgentVersion) + } + + key := guestMetadataCacheKey("pve", "node1", 100) + entry := monitor.guestMetadataCache[key] + if got := guestMetadataCacheEntryTTL(entry); got != guestMetadataEmptyTTL { + t.Fatalf("guestMetadataCacheEntryTTL(identity-only) = %s, want %s", got, guestMetadataEmptyTTL) + } + entry.fetchedAt = time.Now().Add(-guestMetadataEmptyTTL - time.Second) + monitor.guestMetadataCache[key] = entry + monitor.guestMetadataLimiter[key] = time.Now().Add(-time.Second) + + secondIPs, secondIfaces, secondOSName, secondOSVersion, secondAgentVersion := monitor.fetchGuestAgentMetadata( + context.Background(), + client, + "pve", + "node1", + "vm100", + 100, + status, + false, + ) + + if len(secondIPs) == 0 || len(secondIfaces) == 0 { + t.Fatalf("expected second fetch to populate network metadata, got ips=%#v ifaces=%#v", secondIPs, secondIfaces) + } + if secondOSName != firstOSName || secondOSVersion != firstOSVersion || secondAgentVersion != firstAgentVersion { + t.Fatalf("expected identity metadata to be preserved, got os=%q/%q agent=%q want os=%q/%q agent=%q", + secondOSName, secondOSVersion, secondAgentVersion, firstOSName, firstOSVersion, firstAgentVersion) + } +} + +func TestFetchGuestAgentMetadataRetriesIPOnlyCacheSooner(t *testing.T) { + t.Parallel() + + client := &ipOnlyThenNetworkGuestMetadataClient{} + monitor := &Monitor{ + guestMetadataCache: make(map[string]guestMetadataCacheEntry), + guestMetadataLimiter: make(map[string]time.Time), + } + + status := &proxmox.VMStatus{Agent: proxmox.VMAgentField{Value: 1}} + + firstIPs, firstIfaces, _, _, _ := monitor.fetchGuestAgentMetadata( + context.Background(), + client, + "pve", + "node1", + "vm100", + 100, + status, + false, + ) + if len(firstIPs) != 1 || firstIPs[0] != "192.168.1.60" { + t.Fatalf("expected first fetch to preserve discovered IP, got %#v", firstIPs) + } + if len(firstIfaces) != 1 || firstIfaces[0].Name != "" || firstIfaces[0].MAC != "" { + t.Fatalf("expected first fetch to remain identity-incomplete, got %#v", firstIfaces) + } + + key := guestMetadataCacheKey("pve", "node1", 100) + entry := monitor.guestMetadataCache[key] + if got := guestMetadataCacheEntryTTL(entry); got != guestMetadataEmptyTTL { + t.Fatalf("guestMetadataCacheEntryTTL(ip-only) = %s, want %s", got, guestMetadataEmptyTTL) + } + entry.fetchedAt = time.Now().Add(-guestMetadataEmptyTTL - time.Second) + monitor.guestMetadataCache[key] = entry + monitor.guestMetadataLimiter[key] = time.Now().Add(-time.Second) + + secondIPs, secondIfaces, _, _, _ := monitor.fetchGuestAgentMetadata( + context.Background(), + client, + "pve", + "node1", + "vm100", + 100, + status, + false, + ) + if len(secondIPs) != 1 || secondIPs[0] != "192.168.1.60" { + t.Fatalf("expected second fetch to preserve IP, got %#v", secondIPs) + } + if len(secondIfaces) != 1 || secondIfaces[0].Name != "Ethernet0" { + t.Fatalf("expected second fetch to populate interfaces, got %#v", secondIfaces) + } + if client.networkCalls != 2 { + t.Fatalf("expected network metadata to be fetched twice, got %d calls", client.networkCalls) + } +} + +func TestFetchGuestAgentMetadataRetriesEmptyCacheSooner(t *testing.T) { + t.Parallel() + + client := &emptyThenPopulatedGuestMetadataClient{} + monitor := &Monitor{ + guestMetadataCache: make(map[string]guestMetadataCacheEntry), + guestMetadataLimiter: make(map[string]time.Time), + } + + status := &proxmox.VMStatus{Agent: proxmox.VMAgentField{Value: 1}} + + firstIPs, firstIfaces, _, _, _ := monitor.fetchGuestAgentMetadata( + context.Background(), + client, + "pve", + "node1", + "vm100", + 100, + status, + false, + ) + if len(firstIPs) != 0 || len(firstIfaces) != 0 { + t.Fatalf("expected first fetch to be empty, got ips=%#v ifaces=%#v", firstIPs, firstIfaces) + } + + key := guestMetadataCacheKey("pve", "node1", 100) + entry := monitor.guestMetadataCache[key] + entry.fetchedAt = time.Now().Add(-guestMetadataEmptyTTL - time.Second) + monitor.guestMetadataCache[key] = entry + monitor.guestMetadataLimiter[key] = time.Now().Add(-time.Second) + + secondIPs, secondIfaces, _, _, _ := monitor.fetchGuestAgentMetadata( + context.Background(), + client, + "pve", + "node1", + "vm100", + 100, + status, + false, + ) + if len(secondIPs) != 1 || secondIPs[0] != "192.168.1.50" { + t.Fatalf("expected second fetch to refresh IPs, got %#v", secondIPs) + } + if len(secondIfaces) != 1 || secondIfaces[0].Name != "Ethernet0" { + t.Fatalf("expected second fetch to refresh interfaces, got %#v", secondIfaces) + } + if client.networkCalls != 2 { + t.Fatalf("expected network metadata to be fetched twice, got %d calls", client.networkCalls) + } +} + func TestAcquireGuestMetadataSlot(t *testing.T) { t.Parallel() diff --git a/internal/monitoring/memory_trust_characterization_test.go b/internal/monitoring/memory_trust_characterization_test.go index 37a534d2a..bf6598767 100644 --- a/internal/monitoring/memory_trust_characterization_test.go +++ b/internal/monitoring/memory_trust_characterization_test.go @@ -300,7 +300,7 @@ func TestHandleClusterVMResourceMemoryTrustCharacterization(t *testing.T) { MaxCPU: 4, } - vm, ok := mon.handleClusterVMResource(context.Background(), "test", res, makeGuestID("test", "node1", 101), client, nil) + vm, ok := mon.handleClusterVMResource(context.Background(), "test", res, makeGuestID("test", "node1", 101), client, nil, nil) if !ok { t.Fatal("handleClusterVMResource() returned ok=false") } diff --git a/internal/monitoring/monitor_extra_coverage_test.go b/internal/monitoring/monitor_extra_coverage_test.go index 10759b308..8174d3ed7 100644 --- a/internal/monitoring/monitor_extra_coverage_test.go +++ b/internal/monitoring/monitor_extra_coverage_test.go @@ -473,17 +473,28 @@ func TestMonitor_LinkNodeToHostAgent_Extra(t *testing.T) { type mockPVEClientExtra struct { mockPVEClient - resources []proxmox.ClusterResource - vmStatus *proxmox.VMStatus - fsInfo []proxmox.VMFileSystem - netIfaces []proxmox.VMNetworkInterface + vms []proxmox.VM + resources []proxmox.ClusterResource + vmStatus *proxmox.VMStatus + vmStatusErr error + fsInfo []proxmox.VMFileSystem + netIfaces []proxmox.VMNetworkInterface + agentInfo map[string]interface{} + agentVersion string } func (m *mockPVEClientExtra) GetClusterResources(ctx context.Context, resourceType string) ([]proxmox.ClusterResource, error) { return m.resources, nil } +func (m *mockPVEClientExtra) GetVMs(ctx context.Context, node string) ([]proxmox.VM, error) { + return m.vms, nil +} + func (m *mockPVEClientExtra) GetVMStatus(ctx context.Context, node string, vmid int) (*proxmox.VMStatus, error) { + if m.vmStatusErr != nil { + return nil, m.vmStatusErr + } return m.vmStatus, nil } @@ -520,10 +531,16 @@ func (m *mockPVEClientExtra) GetContainerInterfaces(ctx context.Context, node st } func (m *mockPVEClientExtra) GetVMAgentInfo(ctx context.Context, node string, vmid int) (map[string]interface{}, error) { + if m.agentInfo != nil { + return m.agentInfo, nil + } return map[string]interface{}{"os": "linux"}, nil } func (m *mockPVEClientExtra) GetVMAgentVersion(ctx context.Context, node string, vmid int) (string, error) { + if m.agentVersion != "" { + return m.agentVersion, nil + } return "1.0", nil } @@ -1063,6 +1080,10 @@ func TestMonitor_PreviousGuestContextForInstance_Extra(t *testing.T) { if len(prev.vms) != 1 || prev.vms[0].VMID != 101 || prev.vms[0].Instance != "pve1" || prev.vms[0].Name != "vm1" { t.Fatalf("expected only pve1 VMs, got %#v", prev.vms) } + canonicalID := prev.vms[0].ID + if len(prev.vmsByID) != 2 || prev.vmsByID[canonicalID].VMID != 101 || prev.vmsByID[makeGuestID("pve1", "", 101)].VMID != 101 { + t.Fatalf("expected previous VM lookup to be indexed by canonical and runtime guest IDs, got %#v", prev.vmsByID) + } if len(prev.containers) != 2 { t.Fatalf("expected only pve1 containers, got %#v", prev.containers) } @@ -1077,6 +1098,362 @@ func TestMonitor_PreviousGuestContextForInstance_Extra(t *testing.T) { } } +func TestBuildVMFromClusterResource_ContinuesGuestAgentQueriesAfterTransientStatusFailure(t *testing.T) { + monitor := &Monitor{ + rateTracker: NewRateTracker(), + guestMetadataCache: make(map[string]guestMetadataCacheEntry), + guestMetadataLimiter: make(map[string]time.Time), + } + client := &mockPVEClientExtra{ + vmStatusErr: fmt.Errorf("API error 500: timeout"), + fsInfo: []proxmox.VMFileSystem{ + {Mountpoint: "/", Type: "ext4", TotalBytes: 100 * 1024 * 1024 * 1024, UsedBytes: 40 * 1024 * 1024 * 1024, Disk: "/dev/vda"}, + }, + netIfaces: []proxmox.VMNetworkInterface{ + {Name: "eth0", HardwareAddr: "00:11:22:33:44:55", IPAddresses: []proxmox.VMIPAddress{{Address: "192.168.1.50", Prefix: 24}}}, + }, + agentInfo: map[string]interface{}{ + "result": map[string]interface{}{ + "name": "Debian GNU/Linux", + "version": "12", + }, + }, + agentVersion: "8.2.0", + } + guestID := makeGuestID("cluster-a", "node-a", 101) + prevVM := &models.VM{ + ID: guestID, + Instance: "cluster-a", + Node: "node-a", + VMID: 101, + Name: "app-vm", + Type: "qemu", + Status: "running", + AgentVersion: "8.1.0", + NetworkInterfaces: []models.GuestNetworkInterface{ + {Name: "eth0", MAC: "00:11:22:33:44:55", Addresses: []string{"192.168.1.50"}}, + }, + LastSeen: time.Now(), + } + + vm, _, _, _, ok := monitor.buildVMFromClusterResource( + context.Background(), + "cluster-a", + proxmox.ClusterResource{ + Type: "qemu", + Node: "node-a", + Name: "app-vm", + Status: "running", + VMID: 101, + MaxCPU: 4, + MaxMem: 8192, + MaxDisk: 100 * 1024 * 1024 * 1024, + }, + client, + guestID, + nil, + prevVM, + ) + if !ok { + t.Fatal("expected VM to be built") + } + if vm.Disk.Usage != 40 { + t.Fatalf("expected guest-agent disk usage after status failure, got %.2f", vm.Disk.Usage) + } + if len(vm.Disks) != 1 || vm.Disks[0].Device != "/dev/vda" { + t.Fatalf("expected guest-agent disk inventory, got %#v", vm.Disks) + } + if len(vm.NetworkInterfaces) != 1 || vm.NetworkInterfaces[0].Name != "eth0" { + t.Fatalf("expected guest-agent interfaces after status failure, got %#v", vm.NetworkInterfaces) + } + if vm.AgentVersion != "8.2.0" { + t.Fatalf("expected refreshed agent version, got %q", vm.AgentVersion) + } + if vm.OSName != "Debian GNU/Linux" || vm.OSVersion != "12" { + t.Fatalf("expected refreshed OS info, got %q %q", vm.OSName, vm.OSVersion) + } +} + +func TestPollVMsAndContainersEfficient_ContinuesGuestAgentQueriesAfterTransientStatusFailure(t *testing.T) { + t.Setenv("PULSE_DATA_DIR", t.TempDir()) + + m := &Monitor{ + state: models.NewState(), + guestAgentFSInfoTimeout: time.Second, + guestAgentRetries: 1, + guestAgentNetworkTimeout: time.Second, + guestAgentOSInfoTimeout: time.Second, + guestAgentVersionTimeout: time.Second, + guestMetadataCache: make(map[string]guestMetadataCacheEntry), + guestMetadataLimiter: make(map[string]time.Time), + rateTracker: NewRateTracker(), + metricsHistory: NewMetricsHistory(100, time.Hour), + alertManager: alerts.NewManager(), + stalenessTracker: NewStalenessTracker(nil), + nodeRRDMemCache: make(map[string]rrdMemCacheEntry), + vmRRDMemCache: make(map[string]rrdMemCacheEntry), + } + defer m.alertManager.Stop() + + registry := unifiedresources.NewRegistry(nil) + guestID := makeGuestID("pve1", "node1", 100) + registry.IngestSnapshot(models.StateSnapshot{ + VMs: []models.VM{ + { + ID: guestID, + VMID: 100, + Name: "vm100", + Node: "node1", + Instance: "pve1", + Type: "qemu", + Status: "running", + IPAddresses: []string{"192.168.1.50"}, + AgentVersion: "8.1.0", + NetworkInterfaces: []models.GuestNetworkInterface{ + {Name: "eth0", MAC: "00:11:22:33:44:55", Addresses: []string{"192.168.1.50"}}, + }, + Disks: []models.Disk{ + {Total: 100 * 1024 * 1024 * 1024, Used: 40 * 1024 * 1024 * 1024, Free: 60 * 1024 * 1024 * 1024, Usage: 40, Mountpoint: "/", Type: "ext4", Device: "/dev/vda"}, + }, + LastSeen: time.Now(), + }, + }, + }) + m.resourceStore = unifiedresources.NewMonitorAdapter(registry) + + prev := m.previousGuestContextForInstance("pve1") + prevVM, ok := prev.vmsByID[guestID] + if !ok || !hasRecentGuestAgentEvidence(&prevVM, time.Now()) { + t.Fatalf("expected recent guest-agent evidence for %s, got %#v", guestID, prev.vmsByID) + } + + client := &mockPVEClientExtra{ + resources: []proxmox.ClusterResource{ + {VMID: 100, Name: "vm100", Node: "node1", Status: "running", Type: "qemu", MaxMem: 8 * 1024, Mem: 4 * 1024, MaxDisk: 100 * 1024 * 1024 * 1024}, + }, + vmStatusErr: fmt.Errorf("API error 500: timeout"), + fsInfo: []proxmox.VMFileSystem{ + {Mountpoint: "/", Type: "ext4", TotalBytes: 100 * 1024 * 1024 * 1024, UsedBytes: 40 * 1024 * 1024 * 1024, Disk: "/dev/vda"}, + }, + netIfaces: []proxmox.VMNetworkInterface{ + {Name: "eth0", HardwareAddr: "00:11:22:33:44:55", IPAddresses: []proxmox.VMIPAddress{{Address: "192.168.1.50", Prefix: 24}}}, + }, + agentInfo: map[string]interface{}{ + "result": map[string]interface{}{ + "name": "Debian GNU/Linux", + "version": "12", + }, + }, + agentVersion: "8.2.0", + } + + if ok := m.pollVMsAndContainersEfficient( + context.Background(), + "pve1", + "", + false, + client, + map[string]string{"node1": "online"}, + ); !ok { + t.Fatal("expected efficient polling path to succeed") + } + + state := m.GetState() + if len(state.VMs) != 1 { + t.Fatalf("expected 1 VM, got %d", len(state.VMs)) + } + + vm := state.VMs[0] + if vm.Disk.Usage != 40 { + t.Fatalf("expected guest-agent disk usage after status failure, got %.2f", vm.Disk.Usage) + } + if len(vm.Disks) != 1 || vm.Disks[0].Device != "/dev/vda" { + t.Fatalf("expected guest-agent disk inventory, got %#v", vm.Disks) + } + if len(vm.NetworkInterfaces) != 1 || vm.NetworkInterfaces[0].Name != "eth0" { + t.Fatalf("expected guest-agent interfaces after status failure, got %#v", vm.NetworkInterfaces) + } + if vm.AgentVersion != "8.2.0" { + t.Fatalf("expected refreshed agent version, got %q", vm.AgentVersion) + } +} + +func TestPollVMsWithNodes_PreservesCachedGuestMetadataWhenStatusUnavailable(t *testing.T) { + t.Setenv("PULSE_DATA_DIR", t.TempDir()) + + m := &Monitor{ + state: models.NewState(), + guestAgentFSInfoTimeout: time.Second, + guestAgentRetries: 1, + guestAgentNetworkTimeout: time.Second, + guestAgentOSInfoTimeout: time.Second, + guestAgentVersionTimeout: time.Second, + guestMetadataCache: make(map[string]guestMetadataCacheEntry), + guestMetadataLimiter: make(map[string]time.Time), + rateTracker: NewRateTracker(), + metricsHistory: NewMetricsHistory(100, time.Hour), + alertManager: alerts.NewManager(), + stalenessTracker: NewStalenessTracker(nil), + nodeRRDMemCache: make(map[string]rrdMemCacheEntry), + vmRRDMemCache: make(map[string]rrdMemCacheEntry), + } + defer m.alertManager.Stop() + + cacheKey := guestMetadataCacheKey("pve1", "node1", 100) + m.guestMetadataCache[cacheKey] = guestMetadataCacheEntry{ + ipAddresses: []string{"192.168.1.50"}, + networkInterfaces: []models.GuestNetworkInterface{ + {Name: "eth0", MAC: "00:11:22:33:44:55", Addresses: []string{"192.168.1.50"}}, + }, + osName: "Debian", + osVersion: "12", + agentVersion: "8.2.0", + fetchedAt: time.Now().Add(-time.Minute), + } + + client := &mockPVEClientExtra{ + vms: []proxmox.VM{ + {VMID: 100, Name: "vm100", Node: "node1", Status: "running", MaxMem: 8 * 1024, Mem: 4 * 1024}, + }, + vmStatusErr: fmt.Errorf("API error 500: timeout"), + } + + m.pollVMsWithNodes( + context.Background(), + "pve1", + "", + false, + client, + []proxmox.Node{{Node: "node1", Status: "online"}}, + map[string]string{"node1": "online"}, + ) + + state := m.GetState() + if len(state.VMs) != 1 { + t.Fatalf("expected 1 VM, got %d", len(state.VMs)) + } + + vm := state.VMs[0] + if len(vm.IPAddresses) != 1 || vm.IPAddresses[0] != "192.168.1.50" { + t.Fatalf("expected cached IP addresses to be preserved, got %#v", vm.IPAddresses) + } + if len(vm.NetworkInterfaces) != 1 || vm.NetworkInterfaces[0].Name != "eth0" { + t.Fatalf("expected cached network interfaces to be preserved, got %#v", vm.NetworkInterfaces) + } + if vm.OSName != "Debian" || vm.OSVersion != "12" { + t.Fatalf("expected cached OS info to be preserved, got %q %q", vm.OSName, vm.OSVersion) + } + if vm.AgentVersion != "8.2.0" { + t.Fatalf("expected cached agent version to be preserved, got %q", vm.AgentVersion) + } +} + +func TestPollVMsWithNodes_ContinuesGuestAgentQueriesAfterTransientStatusFailure(t *testing.T) { + t.Setenv("PULSE_DATA_DIR", t.TempDir()) + + m := &Monitor{ + state: models.NewState(), + guestAgentFSInfoTimeout: time.Second, + guestAgentRetries: 1, + guestAgentNetworkTimeout: time.Second, + guestAgentOSInfoTimeout: time.Second, + guestAgentVersionTimeout: time.Second, + guestMetadataCache: make(map[string]guestMetadataCacheEntry), + guestMetadataLimiter: make(map[string]time.Time), + rateTracker: NewRateTracker(), + metricsHistory: NewMetricsHistory(100, time.Hour), + alertManager: alerts.NewManager(), + stalenessTracker: NewStalenessTracker(nil), + nodeRRDMemCache: make(map[string]rrdMemCacheEntry), + vmRRDMemCache: make(map[string]rrdMemCacheEntry), + } + defer m.alertManager.Stop() + + registry := unifiedresources.NewRegistry(nil) + guestID := makeGuestID("pve1", "node1", 100) + registry.IngestSnapshot(models.StateSnapshot{ + VMs: []models.VM{ + { + ID: guestID, + VMID: 100, + Name: "vm100", + Node: "node1", + Instance: "pve1", + Type: "qemu", + Status: "running", + IPAddresses: []string{"192.168.1.50"}, + AgentVersion: "8.1.0", + NetworkInterfaces: []models.GuestNetworkInterface{ + {Name: "eth0", MAC: "00:11:22:33:44:55", Addresses: []string{"192.168.1.50"}}, + }, + Disks: []models.Disk{ + {Total: 100 * 1024 * 1024 * 1024, Used: 40 * 1024 * 1024 * 1024, Free: 60 * 1024 * 1024 * 1024, Usage: 40, Mountpoint: "/", Type: "ext4", Device: "/dev/vda"}, + }, + LastSeen: time.Now(), + }, + }, + }) + m.resourceStore = unifiedresources.NewMonitorAdapter(registry) + + prev := m.previousGuestContextForInstance("pve1") + prevVM, ok := prev.vmsByID[guestID] + if !ok || !hasRecentGuestAgentEvidence(&prevVM, time.Now()) { + t.Fatalf("expected recent guest-agent evidence for %s, got %#v", guestID, prev.vmsByID) + } + + client := &mockPVEClientExtra{ + vms: []proxmox.VM{ + {VMID: 100, Name: "vm100", Node: "node1", Status: "running", MaxMem: 8 * 1024, Mem: 4 * 1024, MaxDisk: 100 * 1024 * 1024 * 1024}, + }, + vmStatusErr: fmt.Errorf("API error 500: timeout"), + fsInfo: []proxmox.VMFileSystem{ + {Mountpoint: "/", Type: "ext4", TotalBytes: 100 * 1024 * 1024 * 1024, UsedBytes: 40 * 1024 * 1024 * 1024, Disk: "/dev/vda"}, + }, + netIfaces: []proxmox.VMNetworkInterface{ + {Name: "eth0", HardwareAddr: "00:11:22:33:44:55", IPAddresses: []proxmox.VMIPAddress{{Address: "192.168.1.50", Prefix: 24}}}, + }, + agentInfo: map[string]interface{}{ + "result": map[string]interface{}{ + "name": "Debian GNU/Linux", + "version": "12", + }, + }, + agentVersion: "8.2.0", + } + + m.pollVMsWithNodes( + context.Background(), + "pve1", + "", + false, + client, + []proxmox.Node{{Node: "node1", Status: "online"}}, + map[string]string{"node1": "online"}, + ) + + state := m.GetState() + if len(state.VMs) != 1 { + t.Fatalf("expected 1 VM, got %d", len(state.VMs)) + } + + vm := state.VMs[0] + if vm.Disk.Usage != 40 { + t.Fatalf("expected guest-agent disk usage after status failure, got %.2f", vm.Disk.Usage) + } + if vm.DiskStatusReason != "" { + t.Fatalf("expected empty disk status reason, got %q", vm.DiskStatusReason) + } + if len(vm.Disks) != 1 || vm.Disks[0].Device != "/dev/vda" { + t.Fatalf("expected guest-agent disk inventory, got %#v", vm.Disks) + } + if len(vm.NetworkInterfaces) != 1 || vm.NetworkInterfaces[0].Name != "eth0" { + t.Fatalf("expected guest-agent interfaces after status failure, got %#v", vm.NetworkInterfaces) + } + if vm.AgentVersion != "8.2.0" { + t.Fatalf("expected refreshed agent version, got %q", vm.AgentVersion) + } +} + func TestBuildVMFromClusterResource_UsesLinkedHostAgentDiskFallback(t *testing.T) { monitor := &Monitor{rateTracker: NewRateTracker()} client := &mockPVEClientExtra{ @@ -1112,6 +1489,7 @@ func TestBuildVMFromClusterResource_UsesLinkedHostAgentDiskFallback(t *testing.T }, }, }, + nil, ) if !ok { t.Fatal("expected VM to be built") diff --git a/internal/monitoring/monitor_polling_vm.go b/internal/monitoring/monitor_polling_vm.go index fbfb6a1cc..79823e23c 100644 --- a/internal/monitoring/monitor_polling_vm.go +++ b/internal/monitoring/monitor_polling_vm.go @@ -38,6 +38,7 @@ func (m *Monitor) pollVMsWithNodes(ctx context.Context, instanceName string, clu // Capture the previous guest context once per poll cycle so fallback behavior // is based on a consistent pre-poll snapshot. prevGuests := m.previousGuestContextForInstance(instanceName) + prevVMByID := prevGuests.vmsByID vmIDToHostAgent := prevGuests.hostAgentsByVMID log.Debug(). @@ -89,6 +90,10 @@ func (m *Monitor) pollVMsWithNodes(ctx context.Context, instanceName string, clu // Generate canonical guest ID: instance:node:vmid guestID := makeGuestID(instanceName, n.Node, vm.VMID) + var prevVM *models.VM + if prev, ok := prevVMByID[guestID]; ok { + prevVM = &prev + } guestRaw := VMMemoryRaw{ ListingMem: vm.Mem, @@ -151,8 +156,12 @@ func (m *Monitor) pollVMsWithNodes(ctx context.Context, instanceName string, clu memorySource = "status-unavailable" } - if vm.Status == "running" && vmStatus != nil { - guestIPs, guestIfaces, guestOSName, guestOSVersion, agentVersion := m.fetchGuestAgentMetadata(ctx, client, instanceName, n.Node, vm.Name, vm.VMID, vmStatus) + now := time.Now() + guestAgentAvailable := vm.Status == "running" && + (shouldQueryGuestAgent(vmStatus, prevVM, now) || + m.hasRecentGuestMetadataEvidence(instanceName, n.Node, vm.VMID, now)) + if guestAgentAvailable { + guestIPs, guestIfaces, guestOSName, guestOSVersion, agentVersion := m.fetchGuestAgentMetadata(ctx, client, instanceName, n.Node, vm.Name, vm.VMID, vmStatus, vmStatus == nil) if len(guestIPs) > 0 { ipAddresses = guestIPs } @@ -221,21 +230,26 @@ func (m *Monitor) pollVMsWithNodes(ctx context.Context, instanceName string, clu // For running VMs, ALWAYS try to get filesystem info from guest agent // The cluster/resources endpoint always returns 0 for disk usage - if vm.Status == "running" && vmStatus != nil && diskTotal > 0 { + if guestAgentAvailable && diskTotal > 0 { + agentValue := 0 + if vmStatus != nil { + agentValue = vmStatus.Agent.Value + } + // Log the initial state if logging.IsLevelEnabled(zerolog.DebugLevel) { log.Debug(). Str("instance", instanceName). Str("vm", vm.Name). Int("vmid", vm.VMID). - Int("agent", vmStatus.Agent.Value). + Int("agent", agentValue). Uint64("diskUsed", diskUsed). Uint64("diskTotal", diskTotal). Msg("VM has 0 disk usage, checking guest agent") } // Check if agent is enabled - if vmStatus.Agent.Value == 0 { + if vmStatus != nil && vmStatus.Agent.Value == 0 { diskStatusReason = "agent-disabled" if logging.IsLevelEnabled(zerolog.DebugLevel) { log.Debug(). @@ -243,7 +257,7 @@ func (m *Monitor) pollVMsWithNodes(ctx context.Context, instanceName string, clu Str("vm", vm.Name). Msg("Guest agent disabled in VM config") } - } else if vmStatus.Agent.Value > 0 || diskUsed == 0 { + } else { if logging.IsLevelEnabled(zerolog.DebugLevel) { log.Debug(). Str("instance", instanceName). @@ -453,17 +467,15 @@ func (m *Monitor) pollVMsWithNodes(ctx context.Context, instanceName string, clu Msg("Guest agent provided filesystem info but no usable filesystems found (all were special mounts)") } } - } else { - // No vmStatus available or agent disabled - show allocated disk - if diskTotal > 0 { - diskUsage = -1 // Show as allocated size - diskStatusReason = "no-agent" - } } } else if vm.Status == "running" && diskTotal > 0 { - // Running VM but no vmStatus - show allocated disk + // Running VM but no guest-agent evidence - show allocated disk diskUsage = -1 - diskStatusReason = "no-status" + if vmStatus == nil { + diskStatusReason = "no-status" + } else { + diskStatusReason = "no-agent" + } } if vm.Status == "running" && !diskFromAgent { diff --git a/internal/monitoring/monitor_previous_state.go b/internal/monitoring/monitor_previous_state.go index 0c901b72e..200dd7923 100644 --- a/internal/monitoring/monitor_previous_state.go +++ b/internal/monitoring/monitor_previous_state.go @@ -9,6 +9,7 @@ import ( type previousGuestContext struct { vms []models.VM + vmsByID map[string]models.VM containers []models.Container containerOCIByVMID map[int]bool hostAgentsByVMID map[string]models.Host @@ -17,6 +18,7 @@ type previousGuestContext struct { func (m *Monitor) previousGuestContextForInstance(instanceName string) previousGuestContext { ctx := previousGuestContext{ vms: make([]models.VM, 0), + vmsByID: make(map[string]models.VM), containers: make([]models.Container, 0), containerOCIByVMID: make(map[int]bool), hostAgentsByVMID: make(map[string]models.Host), @@ -31,7 +33,15 @@ func (m *Monitor) previousGuestContextForInstance(instanceName string) previousG if vm == nil || vm.Instance() != instanceName { continue } - ctx.vms = append(ctx.vms, previousVMFromView(vm)) + modelVM := previousVMFromView(vm) + ctx.vms = append(ctx.vms, modelVM) + if modelVM.ID != "" { + ctx.vmsByID[modelVM.ID] = modelVM + } + guestID := makeGuestID(modelVM.Instance, modelVM.Node, modelVM.VMID) + if guestID != "" { + ctx.vmsByID[guestID] = modelVM + } } for _, ct := range readState.Containers() { @@ -84,13 +94,21 @@ func previousVMFromView(vm *unifiedresources.VMView) models.VM { return models.VM{} } return models.VM{ - ID: vm.ID(), - Instance: vm.Instance(), - Node: vm.Node(), - VMID: vm.VMID(), - Name: vm.Name(), - Status: string(vm.Status()), - LastSeen: vm.LastSeen(), + ID: vm.ID(), + Instance: vm.Instance(), + Node: vm.Node(), + VMID: vm.VMID(), + Name: vm.Name(), + Type: "qemu", + Status: string(vm.Status()), + IPAddresses: vm.IPAddresses(), + OSName: vm.OSName(), + OSVersion: vm.OSVersion(), + AgentVersion: vm.AgentVersion(), + NetworkInterfaces: guestNetworkInterfacesFromReadStateView(vm.NetworkInterfaces()), + Disks: guestDisksFromReadStateView(vm.Disks()), + DiskStatusReason: vm.DiskStatusReason(), + LastSeen: vm.LastSeen(), } } diff --git a/internal/monitoring/monitor_proxmox_pool_test.go b/internal/monitoring/monitor_proxmox_pool_test.go index 7ef7abb00..c6f5918fd 100644 --- a/internal/monitoring/monitor_proxmox_pool_test.go +++ b/internal/monitoring/monitor_proxmox_pool_test.go @@ -28,6 +28,7 @@ func TestBuildVMFromClusterResource_PreservesProxmoxPool(t *testing.T) { nil, "", nil, + nil, ) if !ok { t.Fatal("expected VM to be built") diff --git a/internal/monitoring/monitor_pve_guest_builders.go b/internal/monitoring/monitor_pve_guest_builders.go index 5e660c9da..9af6563ae 100644 --- a/internal/monitoring/monitor_pve_guest_builders.go +++ b/internal/monitoring/monitor_pve_guest_builders.go @@ -72,7 +72,7 @@ func (m *Monitor) applyVMStatusDetails( state.networkOutBytes = int64(status.NetOut) // Gather guest metadata from the agent when available - guestIPs, guestIfaces, guestOSName, guestOSVersion, guestAgentVersion := m.fetchGuestAgentMetadata(ctx, client, instanceName, res.Node, res.Name, res.VMID, status) + guestIPs, guestIfaces, guestOSName, guestOSVersion, guestAgentVersion := m.fetchGuestAgentMetadata(ctx, client, instanceName, res.Node, res.Name, res.VMID, status, false) if len(guestIPs) > 0 { state.ipAddresses = guestIPs } @@ -142,6 +142,7 @@ func (m *Monitor) buildVMFromClusterResource( client PVEClientInterface, guestID string, vmIDToHostAgent map[string]models.Host, + prevVM *models.VM, ) (models.VM, VMMemoryRaw, string, time.Time, bool) { // Skip templates if configured if res.Template == 1 { @@ -196,6 +197,51 @@ func (m *Monitor) buildVMFromClusterResource( Int("vmid", res.VMID). Msg("Could not get VM status, using cluster/resources disk data") } + + now := time.Now() + guestAgentAvailable := shouldQueryGuestAgent(state.detailedStatus, prevVM, now) || + m.hasRecentGuestMetadataEvidence(instanceName, res.Node, res.VMID, now) + if guestAgentAvailable && state.detailedStatus == nil { + guestIPs, guestIfaces, guestOSName, guestOSVersion, guestAgentVersion := m.fetchGuestAgentMetadata( + ctx, + client, + instanceName, + res.Node, + res.Name, + res.VMID, + nil, + true, + ) + if len(guestIPs) > 0 { + state.ipAddresses = guestIPs + } + if len(guestIfaces) > 0 { + state.networkInterfaces = guestIfaces + } + if guestOSName != "" { + state.osName = guestOSName + } + if guestOSVersion != "" { + state.osVersion = guestOSVersion + } + if guestAgentVersion != "" { + state.agentVersion = guestAgentVersion + } + + var fsDisks []models.Disk + state.diskTotal, state.diskUsed, state.diskFree, state.diskUsage, fsDisks, state.diskFromAgent = m.updateVMDisksFromGuestAgentFSInfo( + ctx, + instanceName, + res, + client, + state.diskTotal, + state.diskUsed, + state.diskUsage, + ) + if len(fsDisks) > 0 { + state.individualDisks = fsDisks + } + } } if res.Status != "running" { diff --git a/internal/monitoring/monitor_pve_guest_poll.go b/internal/monitoring/monitor_pve_guest_poll.go index 1c12303a4..4b1bdeed2 100644 --- a/internal/monitoring/monitor_pve_guest_poll.go +++ b/internal/monitoring/monitor_pve_guest_poll.go @@ -30,7 +30,15 @@ func (m *Monitor) pollVMsAndContainersEfficient(ctx context.Context, instanceNam // behavior is based on a consistent pre-poll snapshot. prevGuests := m.previousGuestContextForInstance(instanceName) - allVMs, allContainers := m.collectGuestsFromClusterResources(ctx, instanceName, resources, client, prevGuests.containerOCIByVMID, prevGuests.hostAgentsByVMID) + allVMs, allContainers := m.collectGuestsFromClusterResources( + ctx, + instanceName, + resources, + client, + prevGuests.containerOCIByVMID, + prevGuests.vmsByID, + prevGuests.hostAgentsByVMID, + ) allVMs, allContainers = m.preserveGuestsForGracePeriod(instanceName, resources, prevGuests.vms, prevGuests.containers, nodeEffectiveStatus, allVMs, allContainers) @@ -62,6 +70,7 @@ func (m *Monitor) collectGuestsFromClusterResources( resources []proxmox.ClusterResource, client PVEClientInterface, prevContainerIsOCI map[int]bool, + prevVMByID map[string]models.VM, vmIDToHostAgent map[string]models.Host, ) ([]models.VM, []models.Container) { allVMs := make([]models.VM, 0, len(resources)) @@ -81,7 +90,11 @@ func (m *Monitor) collectGuestsFromClusterResources( switch res.Type { case "qemu": - vm, ok := m.handleClusterVMResource(ctx, instanceName, res, guestID, client, vmIDToHostAgent) + var prevVM *models.VM + if prev, ok := prevVMByID[guestID]; ok { + prevVM = &prev + } + vm, ok := m.handleClusterVMResource(ctx, instanceName, res, guestID, client, prevVM, vmIDToHostAgent) if !ok { continue } @@ -104,9 +117,10 @@ func (m *Monitor) handleClusterVMResource( res proxmox.ClusterResource, guestID string, client PVEClientInterface, + prevVM *models.VM, vmIDToHostAgent map[string]models.Host, ) (models.VM, bool) { - vm, guestRaw, memorySource, sampleTime, ok := m.buildVMFromClusterResource(ctx, instanceName, res, client, guestID, vmIDToHostAgent) + vm, guestRaw, memorySource, sampleTime, ok := m.buildVMFromClusterResource(ctx, instanceName, res, client, guestID, vmIDToHostAgent, prevVM) if !ok { return models.VM{}, false }