From 0c3d3e960d8ebfc219956eecebe4797103ec75b1 Mon Sep 17 00:00:00 2001 From: Ante Gulin Date: Thu, 6 Aug 2026 09:50:41 +0200 Subject: [PATCH 01/23] PMM-15227 Remove scaled-down HA replicas from Inventory --- .../docs/install-pmm/install-HA-clustered.md | 1 + managed/models/database.go | 6 + managed/models/node_helpers.go | 104 ++++++++++++ managed/models/node_helpers_test.go | 155 ++++++++++++++++++ 4 files changed, 266 insertions(+) diff --git a/documentation/docs/install-pmm/install-HA-clustered.md b/documentation/docs/install-pmm/install-HA-clustered.md index 704fe512ed2..f4d701be749 100644 --- a/documentation/docs/install-pmm/install-HA-clustered.md +++ b/documentation/docs/install-pmm/install-HA-clustered.md @@ -829,6 +829,7 @@ When you scale PMM HA up or down, **all PMM pods will be recreated**. This happe - HAProxy continues routing to available pods during rollout - No data loss (distributed storage) - Rolling update strategy minimizes downtime + - The Nodes of removed replicas disappear from **Inventory > Nodes** once the remaining pods restart To scale PMM server replicas: diff --git a/managed/models/database.go b/managed/models/database.go index fa839fa3be0..3725ccad7f9 100644 --- a/managed/models/database.go +++ b/managed/models/database.go @@ -1528,6 +1528,12 @@ func setupPMMServerHAAgents(q *reform.Querier, params SetupDBParams) error { // create PMM Server Node and associated Agents in HA mode logrus.Infof("Setting up PMM Server agents in HA mode, Node ID: %s", params.HANodeID) + // Before the "agent already exists" early return, so restarted replicas still clean up. + err := RemoveStaleHANodes(q, params.HANodeID, params.HAPeers) + if err != nil { + return err + } + file, err := os.Open(AgentConfigFilePath) if err != nil { return err diff --git a/managed/models/node_helpers.go b/managed/models/node_helpers.go index a8c4cacff2c..521c7adb704 100644 --- a/managed/models/node_helpers.go +++ b/managed/models/node_helpers.go @@ -18,10 +18,12 @@ package models import ( "errors" "fmt" + "net" "strings" "github.com/AlekSi/pointer" "github.com/google/uuid" + "github.com/sirupsen/logrus" "google.golang.org/grpc/codes" "google.golang.org/grpc/status" "gopkg.in/reform.v1" @@ -334,3 +336,105 @@ func RemoveNode(q *reform.Querier, id string, mode RemoveMode) error { //nolint: } return nil } + +// RemoveStaleHANodes removes the PMM Server Nodes of HA replicas that are no longer configured peers, +// e.g. after a scale-down. Peers are the source of truth because they are regenerated from the replica +// count and restart every replica, while a missing memberlist member may just be restarting. +func RemoveStaleHANodes(q *reform.Querier, haNodeID string, haPeers []string) error { + if len(haPeers) == 0 { + return nil + } + + expected := make(map[string]struct{}, len(haPeers)) + for _, peer := range haPeers { + name, ok := haPeerNodeName(peer) + if !ok { + // Trusting the rest would treat a partial list as the whole cluster and remove live replicas. + logrus.Warnf("Can't read a node name from PMM_HA_PEERS entry %q, skipping the removal of stale HA nodes.", peer) + return nil + } + expected[name] = struct{}{} + } + + if _, ok := expected[haNodeID]; !ok { + logrus.Warnf("PMM_HA_PEERS %v doesn't list this node (PMM_HA_NODE_ID %q), skipping the removal of stale HA nodes.", haPeers, haNodeID) + return nil + } + + nodes, err := FindNodes(q, NodeFilters{}) + if err != nil { + return fmt.Errorf("failed to list Nodes for stale HA node cleanup: %w", err) + } + + for _, node := range nodes { + // Only HA replicas set this flag; every other Node is one the user monitors. + if !node.IsPMMServerNode { + continue + } + if _, ok := expected[node.NodeName]; ok { + continue + } + + monitored, err := haNodeMonitoredServices(q, node.NodeID) + if err != nil { + return err + } + if len(monitored) != 0 { + logrus.Warnf("Keeping stale HA node %q (%s): it still monitors services %v, which would be removed with it. "+ + "Re-add them from a running replica and remove the node from Inventory.", node.NodeName, node.NodeID, monitored) + continue + } + + err = RemoveNode(q, node.NodeID, RemoveCascade) + switch { + case err == nil: + logrus.Infof("Removed stale HA node %q (%s), it is not a part of the cluster anymore.", node.NodeName, node.NodeID) + case errors.Is(err, reform.ErrNoRows), status.Code(err) == codes.NotFound: + logrus.Infof("Stale HA node %q (%s) was already removed by another replica.", node.NodeName, node.NodeID) + default: + return fmt.Errorf("failed to remove stale HA node %q: %w", node.NodeName, err) + } + } + + return nil +} + +// haPeerNodeName maps a PMM_HA_PEERS entry ("pmm-ha-0.pmm-ha.pmm.svc.cluster.local:9761") to a Node +// name: the first label is the pod's PMM_HA_NODE_ID. Reports false for entries with no name, like bare IPs. +func haPeerNodeName(peer string) (string, bool) { + host, _, _ := strings.Cut(strings.TrimSpace(peer), ":") + if net.ParseIP(host) != nil { + return "", false + } + // "/" is memberlist's "name/address" form, "[" an IPv6 literal; neither starts with a node name. + label, _, _ := strings.Cut(host, ".") + if label == "" || strings.ContainsAny(label, "/[") { + return "", false + } + return label, true +} + +// haNodeMonitoredServices returns the IDs of Services whose exporters run under a replica's pmm-agent. +// Remote instances bind theirs to the replica that added them (see management.RDSService), so removing +// that replica's Node takes them with it. +func haNodeMonitoredServices(q *reform.Querier, nodeID string) ([]string, error) { + pmmAgents, err := FindPMMAgentsRunningOnNode(q, nodeID) + if err != nil { + return nil, err + } + + var serviceIDs []string + for _, pmmAgent := range pmmAgents { + agents, err := FindAgents(q, AgentFilters{PMMAgentID: pmmAgent.AgentID}) + if err != nil { + return nil, err + } + for _, agent := range agents { + if agent.ServiceID != nil { + serviceIDs = append(serviceIDs, *agent.ServiceID) + } + } + } + + return serviceIDs, nil +} diff --git a/managed/models/node_helpers_test.go b/managed/models/node_helpers_test.go index c95212ab38d..2a09c6219d1 100644 --- a/managed/models/node_helpers_test.go +++ b/managed/models/node_helpers_test.go @@ -256,3 +256,158 @@ func TestNodeHelpers(t *testing.T) { require.Len(t, nodes, 1) // PMM Server }) } + +func TestRemoveStaleHANodes(t *testing.T) { + sqlDB := testdb.Open(t, models.SetupFixtures, nil) + t.Cleanup(func() { + require.NoError(t, sqlDB.Close()) + }) + + // Two HA replica Nodes, one with a node_exporter, plus an unrelated monitored Node. + setup := func(t *testing.T) (*reform.Querier, func(t *testing.T)) { + t.Helper() + db := reform.NewDB(sqlDB, postgresql.Dialect, reform.NewPrintfLogger(t.Logf)) + tx, err := db.Begin() + require.NoError(t, err) + q := tx.Querier + + for _, str := range []reform.Struct{ + &models.Node{ + NodeID: "ha-node-1", + NodeType: models.GenericNodeType, + NodeName: "pmm-ha-1", + Address: models.LocalhostAddr, + IsPMMServerNode: true, + }, + &models.Agent{ + AgentID: "ha-agent-1", + AgentType: models.PMMAgentType, + RunsOnNodeID: new("ha-node-1"), + }, + &models.Node{ + NodeID: "ha-node-2", + NodeType: models.GenericNodeType, + NodeName: "pmm-ha-2", + Address: models.LocalhostAddr, + IsPMMServerNode: true, + }, + &models.Agent{ + AgentID: "ha-agent-2", + AgentType: models.PMMAgentType, + RunsOnNodeID: new("ha-node-2"), + }, + &models.Agent{ + AgentID: "ha-node-exporter-2", + AgentType: models.NodeExporterType, + PMMAgentID: new("ha-agent-2"), + NodeID: new("ha-node-2"), + }, + &models.Node{ + NodeID: "monitored-node", + NodeType: models.GenericNodeType, + NodeName: "Monitored Node", + }, + } { + require.NoError(t, q.Insert(str), "failed to INSERT %+v", str) + } + + teardown := func(t *testing.T) { + t.Helper() + require.NoError(t, tx.Rollback()) + } + return q, teardown + } + + assertNodeExists := func(t *testing.T, q *reform.Querier, nodeID string) { + t.Helper() + _, err := models.FindNodeByID(q, nodeID) + assert.NoError(t, err) + } + + t.Run("RemovesScaledDownReplicaWithItsAgents", func(t *testing.T) { + q, teardown := setup(t) + defer teardown(t) + + peers := []string{"pmm-ha-0.pmm-ha.pmm.svc.cluster.local:9761", " pmm-ha-1.pmm-ha.pmm.svc.cluster.local "} + require.NoError(t, models.RemoveStaleHANodes(q, "pmm-ha-1", peers)) + + assertNodeExists(t, q, "ha-node-1") + _, err := models.FindAgentByID(q, "ha-agent-1") + require.NoError(t, err) + + _, err = models.FindNodeByID(q, "ha-node-2") + tests.AssertGRPCErrorCode(t, codes.NotFound, err) + + // the removal cascades to the agents of the stale node + for _, agentID := range []string{"ha-agent-2", "ha-node-exporter-2"} { + _, err := models.FindAgentByID(q, agentID) + tests.AssertGRPCErrorCode(t, codes.NotFound, err) + } + + // neither monitored nodes nor the pre-HA pmm-server Node are touched + assertNodeExists(t, q, "monitored-node") + assertNodeExists(t, q, models.PMMServerNodeID) + }) + + t.Run("KeepsAllReplicasWhenNothingWasScaledDown", func(t *testing.T) { + q, teardown := setup(t) + defer teardown(t) + + // a dotless host with a port is what a hand-written PMM_HA_PEERS looks like + peers := []string{"pmm-ha-1.pmm-ha:9761", "pmm-ha-2:9761"} + require.NoError(t, models.RemoveStaleHANodes(q, "pmm-ha-1", peers)) + + assertNodeExists(t, q, "ha-node-1") + assertNodeExists(t, q, "ha-node-2") + }) + + t.Run("KeepsScaledDownReplicaThatStillMonitorsServices", func(t *testing.T) { + q, teardown := setup(t) + defer teardown(t) + + // an exporter for a remote instance, bound to the scaled-down replica's pmm-agent + for _, str := range []reform.Struct{ + &models.Service{ + ServiceID: "rds-service", + ServiceType: models.MySQLServiceType, + ServiceName: "RDS instance", + NodeID: "monitored-node", + Address: new("rds.example.com"), + Port: new(uint16(3306)), + }, + &models.Agent{ + AgentID: "rds-exporter", + AgentType: models.MySQLdExporterType, + PMMAgentID: new("ha-agent-2"), + ServiceID: new("rds-service"), + }, + } { + require.NoError(t, q.Insert(str), "failed to INSERT %+v", str) + } + + peers := []string{"pmm-ha-0.pmm-ha:9761", "pmm-ha-1.pmm-ha:9761"} + require.NoError(t, models.RemoveStaleHANodes(q, "pmm-ha-1", peers)) + + assertNodeExists(t, q, "ha-node-2") + _, err := models.FindAgentByID(q, "rds-exporter") + require.NoError(t, err) + }) + + t.Run("DoesNothingWhenPeersCantBeTrusted", func(t *testing.T) { + q, teardown := setup(t) + defer teardown(t) + + for _, peers := range [][]string{ + {"pmm-ha-2.pmm-ha:9761"}, // lists only the other replica + {"10.244.1.7:9761", "10.244.2.8:9761"}, // no node names to read + {"pmm-ha-1.pmm-ha:9761", "10.244.2.8:9761"}, // mixed: one entry hides a live replica + {"pmm-ha-1.pmm-ha:9761", "pmm-ha-2/10.0.0.2"}, // memberlist "name/address" form + nil, + } { + require.NoError(t, models.RemoveStaleHANodes(q, "pmm-ha-1", peers)) + + assertNodeExists(t, q, "ha-node-1") + assertNodeExists(t, q, "ha-node-2") + } + }) +} From b82c6579452f68eaf535d1c8749d856553ea7571 Mon Sep 17 00:00:00 2001 From: Ante Gulin Date: Thu, 6 Aug 2026 10:19:45 +0200 Subject: [PATCH 02/23] PMM-15227 Let the HA cleanup remove stale PMM Server Nodes --- managed/models/node_helpers.go | 12 +++++++++--- 1 file changed, 9 insertions(+), 3 deletions(-) diff --git a/managed/models/node_helpers.go b/managed/models/node_helpers.go index 810adac913f..88aeda84528 100644 --- a/managed/models/node_helpers.go +++ b/managed/models/node_helpers.go @@ -254,13 +254,19 @@ func CreateNode(q *reform.Querier, nodeType NodeType, params *CreateNodeParams) } // RemoveNode removes single Node. -func RemoveNode(q *reform.Querier, id string, mode RemoveMode) error { //nolint:gocognit +func RemoveNode(q *reform.Querier, id string, mode RemoveMode) error { + return removeNode(q, id, mode, false) +} + +// removeNode removes a single Node. The allowPMMServerNode flag lifts the ban on Nodes flagged as PMM +// Server Nodes; only the HA cleanup sets it, to reap replicas that are no longer part of the cluster. +func removeNode(q *reform.Querier, id string, mode RemoveMode, allowPMMServerNode bool) error { //nolint:gocognit n, err := FindNodeByID(q, id) if err != nil { return err } - if n.IsPMMServerNode || id == PMMServerNodeID { + if id == PMMServerNodeID || (!allowPMMServerNode && n.IsPMMServerNode) { return status.Error(codes.PermissionDenied, "PMM Server node can't be removed.") } @@ -385,7 +391,7 @@ func RemoveStaleHANodes(q *reform.Querier, haNodeID string, haPeers []string) er continue } - err = RemoveNode(q, node.NodeID, RemoveCascade) + err = removeNode(q, node.NodeID, RemoveCascade, true) switch { case err == nil: logrus.Infof("Removed stale HA node %q (%s), it is not a part of the cluster anymore.", node.NodeName, node.NodeID) From 1a1a3c5d25d93435de079fc91709028fabff5284 Mon Sep 17 00:00:00 2001 From: Ante Gulin Date: Thu, 6 Aug 2026 12:35:03 +0200 Subject: [PATCH 03/23] PMM-15227 Log HA cleanup with structured fields --- managed/models/node_helpers.go | 16 ++++++++++------ 1 file changed, 10 insertions(+), 6 deletions(-) diff --git a/managed/models/node_helpers.go b/managed/models/node_helpers.go index 88aeda84528..7dd5d4dfbf6 100644 --- a/managed/models/node_helpers.go +++ b/managed/models/node_helpers.go @@ -351,19 +351,21 @@ func RemoveStaleHANodes(q *reform.Querier, haNodeID string, haPeers []string) er return nil } + l := logrus.WithFields(logrus.Fields{"component": "ha", "ha_node_id": haNodeID}) + expected := make(map[string]struct{}, len(haPeers)) for _, peer := range haPeers { name, ok := haPeerNodeName(peer) if !ok { // Trusting the rest would treat a partial list as the whole cluster and remove live replicas. - logrus.Warnf("Can't read a node name from PMM_HA_PEERS entry %q, skipping the removal of stale HA nodes.", peer) + l.WithField("peer", peer).Warn("Can't read a node name from a PMM_HA_PEERS entry, skipping the removal of stale HA nodes.") return nil } expected[name] = struct{}{} } if _, ok := expected[haNodeID]; !ok { - logrus.Warnf("PMM_HA_PEERS %v doesn't list this node (PMM_HA_NODE_ID %q), skipping the removal of stale HA nodes.", haPeers, haNodeID) + l.WithField("ha_peers", haPeers).Warn("PMM_HA_PEERS doesn't list this node, skipping the removal of stale HA nodes.") return nil } @@ -381,22 +383,24 @@ func RemoveStaleHANodes(q *reform.Querier, haNodeID string, haPeers []string) er continue } + nodeL := l.WithFields(logrus.Fields{"node_id": node.NodeID, "node_name": node.NodeName}) + monitored, err := haNodeMonitoredServices(q, node.NodeID) if err != nil { return err } if len(monitored) != 0 { - logrus.Warnf("Keeping stale HA node %q (%s): it still monitors services %v, which would be removed with it. "+ - "Re-add them from a running replica and remove the node from Inventory.", node.NodeName, node.NodeID, monitored) + nodeL.WithField("service_ids", monitored).Warn("Keeping stale HA node: it still monitors services, which would be removed with it. " + + "Re-add them from a running replica and remove the node from Inventory.") continue } err = removeNode(q, node.NodeID, RemoveCascade, true) switch { case err == nil: - logrus.Infof("Removed stale HA node %q (%s), it is not a part of the cluster anymore.", node.NodeName, node.NodeID) + nodeL.Info("Removed stale HA node, it is not a part of the cluster anymore.") case errors.Is(err, reform.ErrNoRows), status.Code(err) == codes.NotFound: - logrus.Infof("Stale HA node %q (%s) was already removed by another replica.", node.NodeName, node.NodeID) + nodeL.Info("Stale HA node was already removed by another replica.") default: return fmt.Errorf("failed to remove stale HA node %q: %w", node.NodeName, err) } From d52d51a155eb74585d12aac6dc852d75cf4925a9 Mon Sep 17 00:00:00 2001 From: Ante Gulin Date: Thu, 6 Aug 2026 14:44:17 +0200 Subject: [PATCH 04/23] PMM-15227 Correct the IsPMMServerNode comment --- managed/models/node_helpers.go | 3 ++- 1 file changed, 2 insertions(+), 1 deletion(-) diff --git a/managed/models/node_helpers.go b/managed/models/node_helpers.go index 7dd5d4dfbf6..0ba5ef4ca59 100644 --- a/managed/models/node_helpers.go +++ b/managed/models/node_helpers.go @@ -375,7 +375,8 @@ func RemoveStaleHANodes(q *reform.Querier, haNodeID string, haPeers []string) er } for _, node := range nodes { - // Only HA replicas set this flag; every other Node is one the user monitors. + // Set by HA replicas, and by the PMM Server Node of a non-HA deployment; every other + // Node is one the user monitors. if !node.IsPMMServerNode { continue } From c62e178ccc379a2dbeee7467d610f78acdaae060 Mon Sep 17 00:00:00 2001 From: Ante Gulin Date: Thu, 6 Aug 2026 16:05:51 +0200 Subject: [PATCH 05/23] PMM-15227 Reject unbracketed IPv6 HA peers --- managed/models/node_helpers.go | 15 ++++++++++++--- managed/models/node_helpers_test.go | 2 ++ 2 files changed, 14 insertions(+), 3 deletions(-) diff --git a/managed/models/node_helpers.go b/managed/models/node_helpers.go index 0ba5ef4ca59..ca527745339 100644 --- a/managed/models/node_helpers.go +++ b/managed/models/node_helpers.go @@ -411,13 +411,22 @@ func RemoveStaleHANodes(q *reform.Querier, haNodeID string, haPeers []string) er } // haPeerNodeName maps a PMM_HA_PEERS entry ("pmm-ha-0.pmm-ha.pmm.svc.cluster.local:9761") to a Node -// name: the first label is the pod's PMM_HA_NODE_ID. Reports false for entries with no name, like bare IPs. +// name: the first label is the pod's PMM_HA_NODE_ID. Reports false for entries with no name, like +// bare IPv4 or IPv6 addresses. func haPeerNodeName(peer string) (string, bool) { - host, _, _ := strings.Cut(strings.TrimSpace(peer), ":") + peer = strings.TrimSpace(peer) + // Test the whole entry before cutting at ":": an unbracketed IPv6 literal would otherwise be cut + // into its first group, and the "2001" of "2001:db8::7" reads like a node name. Only IPv6 entries + // hold more than one colon, bracketed or not, and none of them starts with a name. + if strings.Count(peer, ":") > 1 || net.ParseIP(peer) != nil { + return "", false + } + host, _, _ := strings.Cut(peer, ":") if net.ParseIP(host) != nil { return "", false } - // "/" is memberlist's "name/address" form, "[" an IPv6 literal; neither starts with a node name. + // "/" is memberlist's "name/address" form, "[" a bracketed address; such a label mixes a name + // with an address instead of being one. label, _, _ := strings.Cut(host, ".") if label == "" || strings.ContainsAny(label, "/[") { return "", false diff --git a/managed/models/node_helpers_test.go b/managed/models/node_helpers_test.go index dca597c8cca..221024c4f33 100644 --- a/managed/models/node_helpers_test.go +++ b/managed/models/node_helpers_test.go @@ -413,6 +413,8 @@ func TestRemoveStaleHANodes(t *testing.T) { {"10.244.1.7:9761", "10.244.2.8:9761"}, // no node names to read {"pmm-ha-1.pmm-ha:9761", "10.244.2.8:9761"}, // mixed: one entry hides a live replica {"pmm-ha-1.pmm-ha:9761", "pmm-ha-2/10.0.0.2"}, // memberlist "name/address" form + {"pmm-ha-1.pmm-ha:9761", "2001:db8::7"}, // an unbracketed IPv6 entry hides a live replica + {"pmm-ha-1.pmm-ha:9761", "[2001:db8::7]:9761"}, nil, } { require.NoError(t, models.RemoveStaleHANodes(q, "pmm-ha-1", peers)) From eb6abb1de8248449990b83db9dd4c036af04eed2 Mon Sep 17 00:00:00 2001 From: Ante Gulin Date: Fri, 7 Aug 2026 13:00:03 +0200 Subject: [PATCH 06/23] PMM-15227 Improve wording --- documentation/docs/install-pmm/install-HA-clustered.md | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/documentation/docs/install-pmm/install-HA-clustered.md b/documentation/docs/install-pmm/install-HA-clustered.md index f4d701be749..35a178b60d6 100644 --- a/documentation/docs/install-pmm/install-HA-clustered.md +++ b/documentation/docs/install-pmm/install-HA-clustered.md @@ -829,7 +829,7 @@ When you scale PMM HA up or down, **all PMM pods will be recreated**. This happe - HAProxy continues routing to available pods during rollout - No data loss (distributed storage) - Rolling update strategy minimizes downtime - - The Nodes of removed replicas disappear from **Inventory > Nodes** once the remaining pods restart + - The Nodes of removed replicas get removed from **Inventory > Nodes** once the remaining pods restart To scale PMM server replicas: From 433b52413e301e4342dee4c963edf66f3b648ee5 Mon Sep 17 00:00:00 2001 From: Ante Gulin Date: Mon, 10 Aug 2026 09:18:19 +0200 Subject: [PATCH 07/23] PMM-15227 Keep the pre-HA PMM Server Node --- managed/models/node_helpers.go | 5 +++++ managed/models/node_helpers_test.go | 16 ++++++++++++++++ 2 files changed, 21 insertions(+) diff --git a/managed/models/node_helpers.go b/managed/models/node_helpers.go index ca527745339..3c7dbb3d4b5 100644 --- a/managed/models/node_helpers.go +++ b/managed/models/node_helpers.go @@ -380,6 +380,11 @@ func RemoveStaleHANodes(q *reform.Querier, haNodeID string, haPeers []string) er if !node.IsPMMServerNode { continue } + // The PMM Server Node of a deployment converted from non-HA: HA replicas always get a + // generated Node ID, and removeNode bans this one outright. + if node.NodeID == PMMServerNodeID { + continue + } if _, ok := expected[node.NodeName]; ok { continue } diff --git a/managed/models/node_helpers_test.go b/managed/models/node_helpers_test.go index 221024c4f33..0e7c9a5294b 100644 --- a/managed/models/node_helpers_test.go +++ b/managed/models/node_helpers_test.go @@ -404,6 +404,22 @@ func TestRemoveStaleHANodes(t *testing.T) { require.NoError(t, err) }) + t.Run("KeepsPreHAPMMServerNode", func(t *testing.T) { + q, teardown := setup(t) + defer teardown(t) + + // A deployment converted from non-HA keeps a Node with the literal "pmm-server" ID; without + // the internal PostgreSQL Service nothing marks it as monitoring, so only its ID keeps it. + service, err := models.FindServiceByName(q, models.PMMServerPostgreSQLServiceName) + require.NoError(t, err) + require.NoError(t, models.RemoveService(q, service.ServiceID, models.RemoveCascade)) + + peers := []string{"pmm-ha-0.pmm-ha:9761", "pmm-ha-1.pmm-ha:9761"} + require.NoError(t, models.RemoveStaleHANodes(q, "pmm-ha-1", peers)) + + assertNodeExists(t, q, models.PMMServerNodeID) + }) + t.Run("DoesNothingWhenPeersCantBeTrusted", func(t *testing.T) { q, teardown := setup(t) defer teardown(t) From ac9e496cc5968d4dabe7254071529002a2bc4b96 Mon Sep 17 00:00:00 2001 From: Ante Gulin Date: Mon, 10 Aug 2026 09:53:26 +0200 Subject: [PATCH 08/23] PMM-15227 Pin PMM Server Node guards to a const removeNode's ban and the stale-HA-node skip both compared against PMMServerNodeID, which setupPMMServerHAAgents reassigns to the replica's generated Node ID. pmm-managed retries SetupDB in a loop, so a commit that failed after registration reopened both guards on the next attempt and would cascade-delete a pre-HA "pmm-server" Node. Compare against defaultPMMServerNodeID instead; PMMServerNodeID keeps its runtime meaning for the callers that need it. --- managed/models/node_helpers.go | 7 ++++--- managed/models/node_helpers_test.go | 20 ++++++++++++++++++++ managed/models/node_model.go | 5 ++++- 3 files changed, 28 insertions(+), 4 deletions(-) diff --git a/managed/models/node_helpers.go b/managed/models/node_helpers.go index 3c7dbb3d4b5..e62d6e3b268 100644 --- a/managed/models/node_helpers.go +++ b/managed/models/node_helpers.go @@ -266,7 +266,7 @@ func removeNode(q *reform.Querier, id string, mode RemoveMode, allowPMMServerNod return err } - if id == PMMServerNodeID || (!allowPMMServerNode && n.IsPMMServerNode) { + if id == defaultPMMServerNodeID || (!allowPMMServerNode && n.IsPMMServerNode) { return status.Error(codes.PermissionDenied, "PMM Server node can't be removed.") } @@ -381,8 +381,9 @@ func RemoveStaleHANodes(q *reform.Querier, haNodeID string, haPeers []string) er continue } // The PMM Server Node of a deployment converted from non-HA: HA replicas always get a - // generated Node ID, and removeNode bans this one outright. - if node.NodeID == PMMServerNodeID { + // generated Node ID, and removeNode bans this one outright. Compared against the const, + // not PMMServerNodeID: setupPMMServerHAAgents reassigns that var, and SetupDB is retried. + if node.NodeID == defaultPMMServerNodeID { continue } if _, ok := expected[node.NodeName]; ok { diff --git a/managed/models/node_helpers_test.go b/managed/models/node_helpers_test.go index 0e7c9a5294b..93c5535d059 100644 --- a/managed/models/node_helpers_test.go +++ b/managed/models/node_helpers_test.go @@ -420,6 +420,26 @@ func TestRemoveStaleHANodes(t *testing.T) { assertNodeExists(t, q, models.PMMServerNodeID) }) + t.Run("KeepsPreHAPMMServerNodeAfterSetupRetry", func(t *testing.T) { + q, teardown := setup(t) + defer teardown(t) + + // setupPMMServerHAAgents points PMMServerNodeID at the replica's generated Node ID, and a + // failed commit makes pmm-managed retry the whole setup; the guards can't key off that var. + restore := models.PMMServerNodeID + models.PMMServerNodeID = "ha-node-1" + t.Cleanup(func() { models.PMMServerNodeID = restore }) + + service, err := models.FindServiceByName(q, models.PMMServerPostgreSQLServiceName) + require.NoError(t, err) + require.NoError(t, models.RemoveService(q, service.ServiceID, models.RemoveCascade)) + + peers := []string{"pmm-ha-0.pmm-ha:9761", "pmm-ha-1.pmm-ha:9761"} + require.NoError(t, models.RemoveStaleHANodes(q, "pmm-ha-1", peers)) + + assertNodeExists(t, q, "pmm-server") + }) + t.Run("DoesNothingWhenPeersCantBeTrusted", func(t *testing.T) { q, teardown := setup(t) defer teardown(t) diff --git a/managed/models/node_model.go b/managed/models/node_model.go index ab8ea76073a..c405a506eac 100644 --- a/managed/models/node_model.go +++ b/managed/models/node_model.go @@ -38,10 +38,13 @@ const ( RemoteAzureDatabaseNodeType NodeType = "remote_azure_database" ) +// defaultPMMServerNodeID is the PMM Server Node ID that never changes at runtime. +const defaultPMMServerNodeID = "pmm-server" + // PMMServerNodeID is a special Node ID representing PMM Server Node. // It takes the value of "pmm-server" in regular non-HA setups and in Active/Passive HA setups, // while in Active/Active HA setups it is set to a dynamically generated UUID. -var PMMServerNodeID = string("pmm-server") +var PMMServerNodeID = defaultPMMServerNodeID // Node represents Node as stored in database. // From dafc5903024a4c3a1b61230e14f2a3846bc079ef Mon Sep 17 00:00:00 2001 From: Ante Gulin Date: Mon, 10 Aug 2026 13:59:50 +0200 Subject: [PATCH 09/23] PMM-15227 Never let HA node cleanup block startup RemoveStaleHANodes ran inside the schema-migration transaction, so any error it returned aborted the migration and pmm-managed eventually gave up with "Could not migrate DB: timeout". A NotFound from the monitored- services pre-check is enough to trigger that whenever replicas restart together. Split it into StaleHANodes, which picks the Nodes to drop, and RemoveStaleHANode, which drops one. migrateDB runs the sweep after the migration commits, gives each Node a transaction of its own, and logs failures instead of returning them: a Node that can't be removed is rolled back whole and skipped, and the rest of the sweep continues. Extract setupFixtures to keep migrateDB under the gocognit limit. --- managed/models/database.go | 97 ++++++++---- managed/models/node_helpers.go | 42 +++--- managed/models/node_helpers_test.go | 221 ++++++++++++++++++---------- 3 files changed, 231 insertions(+), 129 deletions(-) diff --git a/managed/models/database.go b/managed/models/database.go index 3725ccad7f9..339cff5fc2d 100644 --- a/managed/models/database.go +++ b/managed/models/database.go @@ -1468,7 +1468,7 @@ func migrateDB(db *reform.DB, params SetupDBParams) error { } // rollback all migrations if one of them fails; PostgreSQL supports DDL transactions - return db.InTransaction(func(tx *reform.TX) error { + err := db.InTransaction(func(tx *reform.TX) error { for version := currentVersion + 1; version <= latestVersion; version++ { if params.Logf != nil { params.Logf("Migrating database to schema version %d ...", version) @@ -1489,32 +1489,75 @@ func migrateDB(db *reform.DB, params SetupDBParams) error { return nil } - err := EncryptDB(tx, params.Name, DefaultAgentEncryptionColumnsV3) - if err != nil { - return err - } + return setupFixtures(tx, params) + }) + if err != nil { + return err + } - // fill settings with defaults - s, err := GetSettings(tx) - if err != nil { - return err - } - err = SaveSettings(tx, s) - if err != nil { - return err - } + removeStaleHANodes(db, params) - if params.HANodeID != "" { - err = setupPMMServerHAAgents(tx.Querier, params) - } else { - err = setupPMMServerAgents(tx.Querier, params) - } - if err != nil { - return err - } + return nil +} - return nil - }) +// setupFixtures adds the initial data of a fresh PMM Server: encryption, default settings, and the +// PMM Server Node with its Agents. +func setupFixtures(tx *reform.TX, params SetupDBParams) error { + err := EncryptDB(tx, params.Name, DefaultAgentEncryptionColumnsV3) + if err != nil { + return err + } + + // fill settings with defaults + s, err := GetSettings(tx) + if err != nil { + return err + } + err = SaveSettings(tx, s) + if err != nil { + return err + } + + if params.HANodeID != "" { + return setupPMMServerHAAgents(tx.Querier, params) + } + + return setupPMMServerAgents(tx.Querier, params) +} + +// removeStaleHANodes drops the Inventory Nodes of HA replicas that were scaled away. Those rows are +// cosmetic, so this runs outside the migration transaction and only logs failures: tidying them up +// must never keep a replica from starting. +func removeStaleHANodes(db *reform.DB, params SetupDBParams) { + if params.HANodeID == "" || params.SetupFixtures == SkipFixtures { + return + } + + l := logrus.WithFields(logrus.Fields{"component": "ha", "ha_node_id": params.HANodeID}) + + nodes, err := StaleHANodes(db.Querier, params.HANodeID, params.HAPeers) + if err != nil { + l.WithError(err).Warn("Failed to look for stale HA nodes.") + return + } + + for _, node := range nodes { + nodeL := l.WithFields(logrus.Fields{"node_id": node.NodeID, "node_name": node.NodeName}) + + // A transaction per Node: a failure rolls that Node back whole instead of leaving it + // half-removed, and leaves the Nodes this sweep hasn't reached yet alone. + err := db.InTransaction(func(tx *reform.TX) error { + return RemoveStaleHANode(tx.Querier, node.NodeID) + }) + switch { + case err == nil: + nodeL.Info("Removed stale HA node, it is not a part of the cluster anymore.") + case errors.Is(err, reform.ErrNoRows), status.Code(err) == codes.NotFound: + nodeL.Info("Stale HA node was already removed by another replica.") + default: + nodeL.WithError(err).Warn("Failed to remove a stale HA node, keeping it.") + } + } } type agentConfig struct { @@ -1528,12 +1571,6 @@ func setupPMMServerHAAgents(q *reform.Querier, params SetupDBParams) error { // create PMM Server Node and associated Agents in HA mode logrus.Infof("Setting up PMM Server agents in HA mode, Node ID: %s", params.HANodeID) - // Before the "agent already exists" early return, so restarted replicas still clean up. - err := RemoveStaleHANodes(q, params.HANodeID, params.HAPeers) - if err != nil { - return err - } - file, err := os.Open(AgentConfigFilePath) if err != nil { return err diff --git a/managed/models/node_helpers.go b/managed/models/node_helpers.go index e62d6e3b268..8c67cd5f604 100644 --- a/managed/models/node_helpers.go +++ b/managed/models/node_helpers.go @@ -343,12 +343,13 @@ func removeNode(q *reform.Querier, id string, mode RemoveMode, allowPMMServerNod return nil } -// RemoveStaleHANodes removes the PMM Server Nodes of HA replicas that are no longer configured peers, -// e.g. after a scale-down. Peers are the source of truth because they are regenerated from the replica -// count and restart every replica, while a missing memberlist member may just be restarting. -func RemoveStaleHANodes(q *reform.Querier, haNodeID string, haPeers []string) error { +// StaleHANodes returns the PMM Server Nodes of HA replicas that are no longer configured peers, +// e.g. after a scale-down, and that can be removed without taking a user's monitoring with them. +// Peers are the source of truth because they are regenerated from the replica count and restart +// every replica, while a missing memberlist member may just be restarting. +func StaleHANodes(q *reform.Querier, haNodeID string, haPeers []string) ([]*Node, error) { if len(haPeers) == 0 { - return nil + return nil, nil } l := logrus.WithFields(logrus.Fields{"component": "ha", "ha_node_id": haNodeID}) @@ -359,21 +360,22 @@ func RemoveStaleHANodes(q *reform.Querier, haNodeID string, haPeers []string) er if !ok { // Trusting the rest would treat a partial list as the whole cluster and remove live replicas. l.WithField("peer", peer).Warn("Can't read a node name from a PMM_HA_PEERS entry, skipping the removal of stale HA nodes.") - return nil + return nil, nil } expected[name] = struct{}{} } if _, ok := expected[haNodeID]; !ok { l.WithField("ha_peers", haPeers).Warn("PMM_HA_PEERS doesn't list this node, skipping the removal of stale HA nodes.") - return nil + return nil, nil } nodes, err := FindNodes(q, NodeFilters{}) if err != nil { - return fmt.Errorf("failed to list Nodes for stale HA node cleanup: %w", err) + return nil, fmt.Errorf("failed to list Nodes for stale HA node cleanup: %w", err) } + var stale []*Node for _, node := range nodes { // Set by HA replicas, and by the PMM Server Node of a non-HA deployment; every other // Node is one the user monitors. @@ -392,9 +394,12 @@ func RemoveStaleHANodes(q *reform.Querier, haNodeID string, haPeers []string) er nodeL := l.WithFields(logrus.Fields{"node_id": node.NodeID, "node_name": node.NodeName}) + // Keeping a Node is per Node: another replica removing the same rows concurrently must not + // hide the Nodes this pass hasn't reached yet. monitored, err := haNodeMonitoredServices(q, node.NodeID) if err != nil { - return err + nodeL.WithError(err).Warn("Can't tell whether a stale HA node monitors services, keeping it.") + continue } if len(monitored) != 0 { nodeL.WithField("service_ids", monitored).Warn("Keeping stale HA node: it still monitors services, which would be removed with it. " + @@ -402,18 +407,17 @@ func RemoveStaleHANodes(q *reform.Querier, haNodeID string, haPeers []string) er continue } - err = removeNode(q, node.NodeID, RemoveCascade, true) - switch { - case err == nil: - nodeL.Info("Removed stale HA node, it is not a part of the cluster anymore.") - case errors.Is(err, reform.ErrNoRows), status.Code(err) == codes.NotFound: - nodeL.Info("Stale HA node was already removed by another replica.") - default: - return fmt.Errorf("failed to remove stale HA node %q: %w", node.NodeName, err) - } + stale = append(stale, node) } - return nil + return stale, nil +} + +// RemoveStaleHANode removes a Node returned by StaleHANodes, along with the Agents and Services that +// belong to it. Call it in a transaction of its own: it deletes those before the Node itself, so a +// failure half-way through would otherwise leave the Node partially removed. +func RemoveStaleHANode(q *reform.Querier, nodeID string) error { + return removeNode(q, nodeID, RemoveCascade, true) } // haPeerNodeName maps a PMM_HA_PEERS entry ("pmm-ha-0.pmm-ha.pmm.svc.cluster.local:9761") to a Node diff --git a/managed/models/node_helpers_test.go b/managed/models/node_helpers_test.go index 93c5535d059..5d55e46da80 100644 --- a/managed/models/node_helpers_test.go +++ b/managed/models/node_helpers_test.go @@ -268,108 +268,109 @@ func TestNodeHelpers(t *testing.T) { }) } -func TestRemoveStaleHANodes(t *testing.T) { +// insertHAFixtures adds two HA replica Nodes, one with a node_exporter, plus an unrelated +// monitored Node, on top of the PMM Server fixtures. +func insertHAFixtures(t *testing.T, q *reform.Querier) { + t.Helper() + + for _, str := range []reform.Struct{ + &models.Node{ + NodeID: "ha-node-1", + NodeType: models.GenericNodeType, + NodeName: "pmm-ha-1", + Address: models.LocalhostAddr, + IsPMMServerNode: true, + }, + &models.Agent{ + AgentID: "ha-agent-1", + AgentType: models.PMMAgentType, + RunsOnNodeID: new("ha-node-1"), + }, + &models.Node{ + NodeID: "ha-node-2", + NodeType: models.GenericNodeType, + NodeName: "pmm-ha-2", + Address: models.LocalhostAddr, + IsPMMServerNode: true, + }, + &models.Agent{ + AgentID: "ha-agent-2", + AgentType: models.PMMAgentType, + RunsOnNodeID: new("ha-node-2"), + }, + &models.Agent{ + AgentID: "ha-node-exporter-2", + AgentType: models.NodeExporterType, + PMMAgentID: new("ha-agent-2"), + NodeID: new("ha-node-2"), + }, + &models.Node{ + NodeID: "monitored-node", + NodeType: models.GenericNodeType, + NodeName: "Monitored Node", + }, + } { + require.NoError(t, q.Insert(str), "failed to INSERT %+v", str) + } +} + +func assertNodeExists(t *testing.T, q *reform.Querier, nodeID string) { + t.Helper() + _, err := models.FindNodeByID(q, nodeID) + assert.NoError(t, err) +} + +func TestStaleHANodes(t *testing.T) { sqlDB := testdb.Open(t, models.SetupFixtures, nil) t.Cleanup(func() { require.NoError(t, sqlDB.Close()) }) - // Two HA replica Nodes, one with a node_exporter, plus an unrelated monitored Node. setup := func(t *testing.T) (*reform.Querier, func(t *testing.T)) { t.Helper() db := reform.NewDB(sqlDB, postgresql.Dialect, reform.NewPrintfLogger(t.Logf)) tx, err := db.Begin() require.NoError(t, err) - q := tx.Querier - - for _, str := range []reform.Struct{ - &models.Node{ - NodeID: "ha-node-1", - NodeType: models.GenericNodeType, - NodeName: "pmm-ha-1", - Address: models.LocalhostAddr, - IsPMMServerNode: true, - }, - &models.Agent{ - AgentID: "ha-agent-1", - AgentType: models.PMMAgentType, - RunsOnNodeID: new("ha-node-1"), - }, - &models.Node{ - NodeID: "ha-node-2", - NodeType: models.GenericNodeType, - NodeName: "pmm-ha-2", - Address: models.LocalhostAddr, - IsPMMServerNode: true, - }, - &models.Agent{ - AgentID: "ha-agent-2", - AgentType: models.PMMAgentType, - RunsOnNodeID: new("ha-node-2"), - }, - &models.Agent{ - AgentID: "ha-node-exporter-2", - AgentType: models.NodeExporterType, - PMMAgentID: new("ha-agent-2"), - NodeID: new("ha-node-2"), - }, - &models.Node{ - NodeID: "monitored-node", - NodeType: models.GenericNodeType, - NodeName: "Monitored Node", - }, - } { - require.NoError(t, q.Insert(str), "failed to INSERT %+v", str) - } + insertHAFixtures(t, tx.Querier) teardown := func(t *testing.T) { t.Helper() require.NoError(t, tx.Rollback()) } - return q, teardown + return tx.Querier, teardown } - assertNodeExists := func(t *testing.T, q *reform.Querier, nodeID string) { + assertStale := func(t *testing.T, nodes []*models.Node, nodeIDs ...string) { t.Helper() - _, err := models.FindNodeByID(q, nodeID) - assert.NoError(t, err) + actual := make([]string, 0, len(nodes)) + for _, node := range nodes { + actual = append(actual, node.NodeID) + } + assert.ElementsMatch(t, nodeIDs, actual) } - t.Run("RemovesScaledDownReplicaWithItsAgents", func(t *testing.T) { + t.Run("ReportsScaledDownReplica", func(t *testing.T) { q, teardown := setup(t) defer teardown(t) peers := []string{"pmm-ha-0.pmm-ha.pmm.svc.cluster.local:9761", " pmm-ha-1.pmm-ha.pmm.svc.cluster.local "} - require.NoError(t, models.RemoveStaleHANodes(q, "pmm-ha-1", peers)) - - assertNodeExists(t, q, "ha-node-1") - _, err := models.FindAgentByID(q, "ha-agent-1") + stale, err := models.StaleHANodes(q, "pmm-ha-1", peers) require.NoError(t, err) - _, err = models.FindNodeByID(q, "ha-node-2") - tests.AssertGRPCErrorCode(t, codes.NotFound, err) - - // the removal cascades to the agents of the stale node - for _, agentID := range []string{"ha-agent-2", "ha-node-exporter-2"} { - _, err := models.FindAgentByID(q, agentID) - tests.AssertGRPCErrorCode(t, codes.NotFound, err) - } - - // neither monitored nodes nor the pre-HA pmm-server Node are touched - assertNodeExists(t, q, "monitored-node") - assertNodeExists(t, q, models.PMMServerNodeID) + // neither the live replica, the monitored nodes nor the pre-HA pmm-server Node are reported + assertStale(t, stale, "ha-node-2") }) - t.Run("KeepsAllReplicasWhenNothingWasScaledDown", func(t *testing.T) { + t.Run("ReportsNothingWhenNothingWasScaledDown", func(t *testing.T) { q, teardown := setup(t) defer teardown(t) // a dotless host with a port is what a hand-written PMM_HA_PEERS looks like peers := []string{"pmm-ha-1.pmm-ha:9761", "pmm-ha-2:9761"} - require.NoError(t, models.RemoveStaleHANodes(q, "pmm-ha-1", peers)) + stale, err := models.StaleHANodes(q, "pmm-ha-1", peers) + require.NoError(t, err) - assertNodeExists(t, q, "ha-node-1") - assertNodeExists(t, q, "ha-node-2") + assertStale(t, stale) }) t.Run("KeepsScaledDownReplicaThatStillMonitorsServices", func(t *testing.T) { @@ -397,11 +398,10 @@ func TestRemoveStaleHANodes(t *testing.T) { } peers := []string{"pmm-ha-0.pmm-ha:9761", "pmm-ha-1.pmm-ha:9761"} - require.NoError(t, models.RemoveStaleHANodes(q, "pmm-ha-1", peers)) - - assertNodeExists(t, q, "ha-node-2") - _, err := models.FindAgentByID(q, "rds-exporter") + stale, err := models.StaleHANodes(q, "pmm-ha-1", peers) require.NoError(t, err) + + assertStale(t, stale) }) t.Run("KeepsPreHAPMMServerNode", func(t *testing.T) { @@ -415,9 +415,10 @@ func TestRemoveStaleHANodes(t *testing.T) { require.NoError(t, models.RemoveService(q, service.ServiceID, models.RemoveCascade)) peers := []string{"pmm-ha-0.pmm-ha:9761", "pmm-ha-1.pmm-ha:9761"} - require.NoError(t, models.RemoveStaleHANodes(q, "pmm-ha-1", peers)) + stale, err := models.StaleHANodes(q, "pmm-ha-1", peers) + require.NoError(t, err) - assertNodeExists(t, q, models.PMMServerNodeID) + assertStale(t, stale, "ha-node-2") }) t.Run("KeepsPreHAPMMServerNodeAfterSetupRetry", func(t *testing.T) { @@ -435,12 +436,13 @@ func TestRemoveStaleHANodes(t *testing.T) { require.NoError(t, models.RemoveService(q, service.ServiceID, models.RemoveCascade)) peers := []string{"pmm-ha-0.pmm-ha:9761", "pmm-ha-1.pmm-ha:9761"} - require.NoError(t, models.RemoveStaleHANodes(q, "pmm-ha-1", peers)) + stale, err := models.StaleHANodes(q, "pmm-ha-1", peers) + require.NoError(t, err) - assertNodeExists(t, q, "pmm-server") + assertStale(t, stale, "ha-node-2") }) - t.Run("DoesNothingWhenPeersCantBeTrusted", func(t *testing.T) { + t.Run("ReportsNothingWhenPeersCantBeTrusted", func(t *testing.T) { q, teardown := setup(t) defer teardown(t) @@ -453,10 +455,69 @@ func TestRemoveStaleHANodes(t *testing.T) { {"pmm-ha-1.pmm-ha:9761", "[2001:db8::7]:9761"}, nil, } { - require.NoError(t, models.RemoveStaleHANodes(q, "pmm-ha-1", peers)) + stale, err := models.StaleHANodes(q, "pmm-ha-1", peers) + require.NoError(t, err, "peers: %v", peers) + + assertStale(t, stale) + } + }) +} + +func TestRemoveStaleHANode(t *testing.T) { + // The removal is meant to run in a transaction of its own, which wouldn't see fixtures held in + // an uncommitted one - so every subtest gets a database of its own instead. + setup := func(t *testing.T) *reform.DB { + t.Helper() + sqlDB := testdb.Open(t, models.SetupFixtures, nil) + db := reform.NewDB(sqlDB, postgresql.Dialect, reform.NewPrintfLogger(t.Logf)) + insertHAFixtures(t, db.Querier) + + return db + } + + t.Run("RemovesTheNodeWithItsAgents", func(t *testing.T) { + db := setup(t) + q := db.Querier - assertNodeExists(t, q, "ha-node-1") - assertNodeExists(t, q, "ha-node-2") + require.NoError(t, db.InTransaction(func(tx *reform.TX) error { + return models.RemoveStaleHANode(tx.Querier, "ha-node-2") + })) + + _, err := models.FindNodeByID(q, "ha-node-2") + tests.AssertGRPCErrorCode(t, codes.NotFound, err) + + // the removal cascades to the agents of the stale node + for _, agentID := range []string{"ha-agent-2", "ha-node-exporter-2"} { + _, err := models.FindAgentByID(q, agentID) + tests.AssertGRPCErrorCode(t, codes.NotFound, err) + } + + for _, nodeID := range []string{"ha-node-1", "monitored-node", models.PMMServerNodeID} { + assertNodeExists(t, q, nodeID) + } + }) + + t.Run("LeavesTheNodeWholeWhenRemovalFails", func(t *testing.T) { + db := setup(t) + q := db.Querier + + // RemoveAgent refuses to delete PMMServerAgentID, which HA setup points at the local + // replica's own agent; pointing it at ha-agent-2 makes removing ha-node-2 fail after its + // node_exporter, which removeNode deletes first, is already gone. + restore := models.PMMServerAgentID + models.PMMServerAgentID = "ha-agent-2" + t.Cleanup(func() { models.PMMServerAgentID = restore }) + + err := db.InTransaction(func(tx *reform.TX) error { + return models.RemoveStaleHANode(tx.Querier, "ha-node-2") + }) + tests.AssertGRPCErrorCode(t, codes.PermissionDenied, err) + + // the transaction rolled the whole Node back, node_exporter included + assertNodeExists(t, q, "ha-node-2") + for _, agentID := range []string{"ha-agent-2", "ha-node-exporter-2"} { + _, err := models.FindAgentByID(q, agentID) + require.NoError(t, err) } }) } From 82eefa467b7d549e730de2ec35870d56991fac53 Mon Sep 17 00:00:00 2001 From: Ante Gulin Date: Mon, 10 Aug 2026 15:44:03 +0200 Subject: [PATCH 10/23] PMM-15227 Tighten the HA node cleanup Call the sweep from SetupDB instead of migrateDB, which had no business rewriting Inventory and only grew a setupFixtures extraction to stay under the gocognit limit; migrateDB is back to its original form. Splitting selection from removal left the monitored-services check in a different transaction than the delete, so a service bound to the replica in between would be cascade-removed, and left RemoveStaleHANode willing to remove any Node handed to it with the PMM Server ban lifted. Re-check inside the removal, which closes both, at a negligible cost of three queries per removed Node during startup. --- .../docs/install-pmm/install-HA-clustered.md | 2 +- managed/models/database.go | 58 +++++++++---------- managed/models/node_helpers.go | 11 ++++ managed/models/node_helpers_test.go | 34 +++++++++++ 4 files changed, 72 insertions(+), 33 deletions(-) diff --git a/documentation/docs/install-pmm/install-HA-clustered.md b/documentation/docs/install-pmm/install-HA-clustered.md index 35a178b60d6..f3650c784b3 100644 --- a/documentation/docs/install-pmm/install-HA-clustered.md +++ b/documentation/docs/install-pmm/install-HA-clustered.md @@ -829,7 +829,7 @@ When you scale PMM HA up or down, **all PMM pods will be recreated**. This happe - HAProxy continues routing to available pods during rollout - No data loss (distributed storage) - Rolling update strategy minimizes downtime - - The Nodes of removed replicas get removed from **Inventory > Nodes** once the remaining pods restart + - The Nodes of removed replicas get removed from **Inventory > Nodes** once the remaining pods restart. A Node is kept (and a warning logged) if it still monitors services, for example a remote instance that was added from that replica. Re-add those services from a running client, then remove the Node from **Inventory > Nodes** To scale PMM server replicas: diff --git a/managed/models/database.go b/managed/models/database.go index 339cff5fc2d..410dd241a68 100644 --- a/managed/models/database.go +++ b/managed/models/database.go @@ -1297,6 +1297,8 @@ func SetupDB(ctx context.Context, sqlDB *sql.DB, params SetupDBParams) (*reform. return nil, err } + removeStaleHANodes(db, params) + return db, nil } @@ -1468,7 +1470,7 @@ func migrateDB(db *reform.DB, params SetupDBParams) error { } // rollback all migrations if one of them fails; PostgreSQL supports DDL transactions - err := db.InTransaction(func(tx *reform.TX) error { + return db.InTransaction(func(tx *reform.TX) error { for version := currentVersion + 1; version <= latestVersion; version++ { if params.Logf != nil { params.Logf("Migrating database to schema version %d ...", version) @@ -1489,40 +1491,32 @@ func migrateDB(db *reform.DB, params SetupDBParams) error { return nil } - return setupFixtures(tx, params) - }) - if err != nil { - return err - } - - removeStaleHANodes(db, params) - - return nil -} - -// setupFixtures adds the initial data of a fresh PMM Server: encryption, default settings, and the -// PMM Server Node with its Agents. -func setupFixtures(tx *reform.TX, params SetupDBParams) error { - err := EncryptDB(tx, params.Name, DefaultAgentEncryptionColumnsV3) - if err != nil { - return err - } + err := EncryptDB(tx, params.Name, DefaultAgentEncryptionColumnsV3) + if err != nil { + return err + } - // fill settings with defaults - s, err := GetSettings(tx) - if err != nil { - return err - } - err = SaveSettings(tx, s) - if err != nil { - return err - } + // fill settings with defaults + s, err := GetSettings(tx) + if err != nil { + return err + } + err = SaveSettings(tx, s) + if err != nil { + return err + } - if params.HANodeID != "" { - return setupPMMServerHAAgents(tx.Querier, params) - } + if params.HANodeID != "" { + err = setupPMMServerHAAgents(tx.Querier, params) + } else { + err = setupPMMServerAgents(tx.Querier, params) + } + if err != nil { + return err + } - return setupPMMServerAgents(tx.Querier, params) + return nil + }) } // removeStaleHANodes drops the Inventory Nodes of HA replicas that were scaled away. Those rows are diff --git a/managed/models/node_helpers.go b/managed/models/node_helpers.go index 8c67cd5f604..e8255336c06 100644 --- a/managed/models/node_helpers.go +++ b/managed/models/node_helpers.go @@ -416,7 +416,18 @@ func StaleHANodes(q *reform.Querier, haNodeID string, haPeers []string) ([]*Node // RemoveStaleHANode removes a Node returned by StaleHANodes, along with the Agents and Services that // belong to it. Call it in a transaction of its own: it deletes those before the Node itself, so a // failure half-way through would otherwise leave the Node partially removed. +// +// It lifts the ban on removing PMM Server Nodes, so it re-reads the Services the Node monitors and +// refuses to take a Node that has any. func RemoveStaleHANode(q *reform.Querier, nodeID string) error { + monitored, err := haNodeMonitoredServices(q, nodeID) + if err != nil { + return err + } + if len(monitored) != 0 { + return status.Errorf(codes.FailedPrecondition, "HA Node with ID %q still monitors services.", nodeID) + } + return removeNode(q, nodeID, RemoveCascade, true) } diff --git a/managed/models/node_helpers_test.go b/managed/models/node_helpers_test.go index 5d55e46da80..fc0649e46da 100644 --- a/managed/models/node_helpers_test.go +++ b/managed/models/node_helpers_test.go @@ -520,4 +520,38 @@ func TestRemoveStaleHANode(t *testing.T) { require.NoError(t, err) } }) + + t.Run("RefusesANodeThatStillMonitorsServices", func(t *testing.T) { + db := setup(t) + q := db.Querier + + // a Service bound to the stale replica's pmm-agent between StaleHANodes and the removal + for _, str := range []reform.Struct{ + &models.Service{ + ServiceID: "rds-service", + ServiceType: models.MySQLServiceType, + ServiceName: "RDS instance", + NodeID: "monitored-node", + Address: new("rds.example.com"), + Port: new(uint16(3306)), + }, + &models.Agent{ + AgentID: "rds-exporter", + AgentType: models.MySQLdExporterType, + PMMAgentID: new("ha-agent-2"), + ServiceID: new("rds-service"), + }, + } { + require.NoError(t, q.Insert(str), "failed to INSERT %+v", str) + } + + err := db.InTransaction(func(tx *reform.TX) error { + return models.RemoveStaleHANode(tx.Querier, "ha-node-2") + }) + tests.AssertGRPCErrorCode(t, codes.FailedPrecondition, err) + + assertNodeExists(t, q, "ha-node-2") + _, err = models.FindAgentByID(q, "rds-exporter") + require.NoError(t, err) + }) } From a31b5485c1552991a615228e96475ec69df038b1 Mon Sep 17 00:00:00 2001 From: Ante Gulin Date: Mon, 10 Aug 2026 16:37:29 +0200 Subject: [PATCH 11/23] PMM-15227 Fail fast in the HA node test helper assertNodeExists runs before further assertions on the same Node's agents, so with assert.NoError a missing Node let the subtest carry on and fail a second time for a derived reason. require stops at the real failure. --- managed/models/node_helpers_test.go | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/managed/models/node_helpers_test.go b/managed/models/node_helpers_test.go index fc0649e46da..a259f3831eb 100644 --- a/managed/models/node_helpers_test.go +++ b/managed/models/node_helpers_test.go @@ -317,7 +317,7 @@ func insertHAFixtures(t *testing.T, q *reform.Querier) { func assertNodeExists(t *testing.T, q *reform.Querier, nodeID string) { t.Helper() _, err := models.FindNodeByID(q, nodeID) - assert.NoError(t, err) + require.NoError(t, err) } func TestStaleHANodes(t *testing.T) { From 4b09b464b861fd05e94e37a674a3e3371168cdf7 Mon Sep 17 00:00:00 2001 From: Ante Gulin Date: Mon, 10 Aug 2026 16:45:03 +0200 Subject: [PATCH 12/23] PMM-15227 Cover the Node ID ban under a lifted flag RemoveStaleHANode is the only caller that passes allowPMMServerNode, so the id == defaultPMMServerNodeID clause is all that keeps the pre-HA "pmm-server" Node from being reaped by the sweep's own primitive, and nothing exercised it. --- managed/models/node_helpers_test.go | 20 ++++++++++++++++++++ 1 file changed, 20 insertions(+) diff --git a/managed/models/node_helpers_test.go b/managed/models/node_helpers_test.go index a259f3831eb..2a65d9e781c 100644 --- a/managed/models/node_helpers_test.go +++ b/managed/models/node_helpers_test.go @@ -554,4 +554,24 @@ func TestRemoveStaleHANode(t *testing.T) { _, err = models.FindAgentByID(q, "rds-exporter") require.NoError(t, err) }) + + t.Run("RefusesThePreHAPMMServerNode", func(t *testing.T) { + db := setup(t) + q := db.Querier + + // With the internal PostgreSQL Service gone, nothing marks it as monitoring, so only the + // literal "pmm-server" ID stands between the lifted PMM Server ban and the Node. + service, err := models.FindServiceByName(q, models.PMMServerPostgreSQLServiceName) + require.NoError(t, err) + require.NoError(t, models.RemoveService(q, service.ServiceID, models.RemoveCascade)) + + err = db.InTransaction(func(tx *reform.TX) error { + return models.RemoveStaleHANode(tx.Querier, models.PMMServerNodeID) + }) + // the Node ban, not RemoveAgent's ban on the PMM Server pmm-agent, which also answers + // PermissionDenied once the removal gets that far + tests.AssertGRPCError(t, status.New(codes.PermissionDenied, `PMM Server node can't be removed.`), err) + + assertNodeExists(t, q, models.PMMServerNodeID) + }) } From 56a34df7f8a50cabf15bc6e541d74297829c83f4 Mon Sep 17 00:00:00 2001 From: Ante Gulin Date: Tue, 11 Aug 2026 07:56:12 +0200 Subject: [PATCH 13/23] PMM-15227 Widen the stale HA Node safety check --- managed/models/agent_helpers.go | 43 +++++++++++++++ managed/models/agent_helpers_test.go | 49 +++++++++++++++++ managed/models/node_helpers.go | 46 +++++++++++----- managed/models/node_helpers_test.go | 78 ++++++++++++++++++++++++++++ 4 files changed, 202 insertions(+), 14 deletions(-) diff --git a/managed/models/agent_helpers.go b/managed/models/agent_helpers.go index 20da2a24104..d6457d213cc 100644 --- a/managed/models/agent_helpers.go +++ b/managed/models/agent_helpers.go @@ -427,6 +427,49 @@ func FindPMMAgentsRunningOnNode(q *reform.Querier, nodeID string) ([]*Agent, err return res, nil } +// FindAgentsOnNode returns Agents attached to or running on the Node: node-level exporters, the +// pmm-agents themselves, and external exporters in pull mode. +func FindAgentsOnNode(q *reform.Querier, nodeID string) ([]*Agent, error) { + structs, err := q.SelectAllFrom(AgentTable, "WHERE runs_on_node_id = $1 OR node_id = $1 ORDER BY agent_id", nodeID) + if err != nil { + return nil, fmt.Errorf("failed to select Agents on Node %q: %w", nodeID, err) + } + + res := make([]*Agent, len(structs)) + for i, str := range structs { + decryptedAgent := DecryptAgent(*str.(*Agent)) //nolint:forcetypeassert + res[i] = &decryptedAgent + } + + return res, nil +} + +// FindAgentsByPMMAgentIDs returns Agents started by any of the given pmm-agents. +func FindAgentsByPMMAgentIDs(q *reform.Querier, pmmAgentIDs []string) ([]*Agent, error) { + if len(pmmAgentIDs) == 0 { + return []*Agent{}, nil + } + + p := strings.Join(q.Placeholders(1, len(pmmAgentIDs)), ", ") + tail := fmt.Sprintf("WHERE pmm_agent_id IN (%s) ORDER BY agent_id", p) + args := make([]any, len(pmmAgentIDs)) + for i, id := range pmmAgentIDs { + args[i] = id + } + structs, err := q.SelectAllFrom(AgentTable, tail, args...) + if err != nil { + return nil, fmt.Errorf("failed to select Agents started by pmm-agents: %w", err) + } + + res := make([]*Agent, len(structs)) + for i, str := range structs { + decryptedAgent := DecryptAgent(*str.(*Agent)) //nolint:forcetypeassert + res[i] = &decryptedAgent + } + + return res, nil +} + // FindPMMAgentsForService gets pmm-agents for service. func FindPMMAgentsForService(q *reform.Querier, serviceID string) ([]*Agent, error) { _, err := q.SelectOneFrom(ServiceTable, "WHERE service_id = $1", serviceID) diff --git a/managed/models/agent_helpers_test.go b/managed/models/agent_helpers_test.go index 462b0420790..0531a2145fc 100644 --- a/managed/models/agent_helpers_test.go +++ b/managed/models/agent_helpers_test.go @@ -464,6 +464,55 @@ func TestAgentHelpers(t *testing.T) { assert.Empty(t, agents) }) + t.Run("FindAgentsOnNode", func(t *testing.T) { + q, teardown := setup(t) + defer teardown(t) + + agents, err := models.FindAgentsOnNode(q, "N1") + require.NoError(t, err) + agentIDs := make([]string, len(agents)) + for i, agent := range agents { + agentIDs[i] = agent.AgentID + assert.True(t, + pointer.GetString(agent.RunsOnNodeID) == "N1" || pointer.GetString(agent.NodeID) == "N1", + "%s is on neither runs_on_node_id nor node_id of N1", agent.AgentID) + } + assert.Contains(t, agentIDs, "A1") // a pmm-agent running on the Node + assert.Contains(t, agentIDs, "A3") // an exporter attached to the Node + assert.Contains(t, agentIDs, "A7") // attached to the Node, but started by a pmm-agent on N2 + assert.NotContains(t, agentIDs, "A2") // bound to a Service, not to the Node + + // find with non existing node. + agents, err = models.FindAgentsOnNode(q, "X1") + require.NoError(t, err) + assert.Empty(t, agents) + }) + + t.Run("FindAgentsByPMMAgentIDs", func(t *testing.T) { + q, teardown := setup(t) + defer teardown(t) + + agents, err := models.FindAgentsByPMMAgentIDs(q, []string{"A1"}) + require.NoError(t, err) + agentIDs := make([]string, len(agents)) + for i, agent := range agents { + agentIDs[i] = agent.AgentID + } + assert.Equal(t, []string{"A2", "A3"}, agentIDs) + + agents, err = models.FindAgentsByPMMAgentIDs(q, []string{"A1", "A4"}) + require.NoError(t, err) + agentIDs = make([]string, len(agents)) + for i, agent := range agents { + agentIDs[i] = agent.AgentID + } + assert.Equal(t, []string{"A2", "A3", "A5", "A6", "A7"}, agentIDs) + + agents, err = models.FindAgentsByPMMAgentIDs(q, nil) + require.NoError(t, err) + assert.Empty(t, agents) + }) + t.Run("FindPMMAgentsForServicesOnNode", func(t *testing.T) { q, teardown := setup(t) defer teardown(t) diff --git a/managed/models/node_helpers.go b/managed/models/node_helpers.go index e8255336c06..751e6822812 100644 --- a/managed/models/node_helpers.go +++ b/managed/models/node_helpers.go @@ -455,27 +455,45 @@ func haPeerNodeName(peer string) (string, bool) { return label, true } -// haNodeMonitoredServices returns the IDs of Services whose exporters run under a replica's pmm-agent. -// Remote instances bind theirs to the replica that added them (see management.RDSService), so removing -// that replica's Node takes them with it. +// haNodeMonitoredServices returns the IDs of the Services that removing the Node would take with it: +// those attached to the Node, those whose exporters run on it (an external exporter in pull mode), and +// those whose exporters run under its pmm-agent (remote instances bind theirs to the replica that +// added them. All three are cascaded away by removeNode. func haNodeMonitoredServices(q *reform.Querier, nodeID string) ([]string, error) { - pmmAgents, err := FindPMMAgentsRunningOnNode(q, nodeID) + // An external exporter carries a service_id itself, and the pmm-agents are the parents of the + // exporters read next. + agents, err := FindAgentsOnNode(q, nodeID) if err != nil { return nil, err } - var serviceIDs []string - for _, pmmAgent := range pmmAgents { - agents, err := FindAgents(q, AgentFilters{PMMAgentID: pmmAgent.AgentID}) - if err != nil { - return nil, err + var serviceIDs, pmmAgentIDs []string + for _, agent := range agents { + if agent.ServiceID != nil { + serviceIDs = append(serviceIDs, *agent.ServiceID) } - for _, agent := range agents { - if agent.ServiceID != nil { - serviceIDs = append(serviceIDs, *agent.ServiceID) - } + if agent.AgentType == PMMAgentType { + pmmAgentIDs = append(pmmAgentIDs, agent.AgentID) } } - return serviceIDs, nil + started, err := FindAgentsByPMMAgentIDs(q, pmmAgentIDs) + if err != nil { + return nil, err + } + for _, agent := range started { + if agent.ServiceID != nil { + serviceIDs = append(serviceIDs, *agent.ServiceID) + } + } + + services, err := FindServices(q, ServiceFilters{NodeID: nodeID}) + if err != nil { + return nil, err + } + for _, service := range services { + serviceIDs = append(serviceIDs, service.ServiceID) + } + + return deduplicateStrings(serviceIDs), nil } diff --git a/managed/models/node_helpers_test.go b/managed/models/node_helpers_test.go index 2a65d9e781c..da52eb35c73 100644 --- a/managed/models/node_helpers_test.go +++ b/managed/models/node_helpers_test.go @@ -404,6 +404,59 @@ func TestStaleHANodes(t *testing.T) { assertStale(t, stale) }) + t.Run("KeepsScaledDownReplicaWithAServiceOnIt", func(t *testing.T) { + q, teardown := setup(t) + defer teardown(t) + + // a Service registered against the replica's own Node, monitored from somewhere else + require.NoError(t, q.Insert(&models.Service{ + ServiceID: "service-on-replica", + ServiceType: models.MySQLServiceType, + ServiceName: "MySQL on the replica", + NodeID: "ha-node-2", + Address: new("mysql.example.com"), + Port: new(uint16(3306)), + })) + + peers := []string{"pmm-ha-0.pmm-ha:9761", "pmm-ha-1.pmm-ha:9761"} + stale, err := models.StaleHANodes(q, "pmm-ha-1", peers) + require.NoError(t, err) + + assertStale(t, stale) + }) + + t.Run("KeepsScaledDownReplicaRunningAnExternalExporter", func(t *testing.T) { + q, teardown := setup(t) + defer teardown(t) + + // an external exporter in pull mode: no pmm-agent owns it, it just runs on the replica + for _, str := range []reform.Struct{ + &models.Service{ + ServiceID: "external-service", + ServiceType: models.ExternalServiceType, + ServiceName: "External instance", + NodeID: "monitored-node", + Address: new("external.example.com"), + Port: new(uint16(9100)), + ExternalGroup: "external", + }, + &models.Agent{ + AgentID: "external-exporter", + AgentType: models.ExternalExporterType, + RunsOnNodeID: new("ha-node-2"), + ServiceID: new("external-service"), + }, + } { + require.NoError(t, q.Insert(str), "failed to INSERT %+v", str) + } + + peers := []string{"pmm-ha-0.pmm-ha:9761", "pmm-ha-1.pmm-ha:9761"} + stale, err := models.StaleHANodes(q, "pmm-ha-1", peers) + require.NoError(t, err) + + assertStale(t, stale) + }) + t.Run("KeepsPreHAPMMServerNode", func(t *testing.T) { q, teardown := setup(t) defer teardown(t) @@ -555,6 +608,31 @@ func TestRemoveStaleHANode(t *testing.T) { require.NoError(t, err) }) + t.Run("RefusesANodeWithAServiceOnIt", func(t *testing.T) { + db := setup(t) + q := db.Querier + + // the re-check is the last line of defence, so it has to cover the same ground as the + // selection: a Service attached to the Node would be cascaded away with it + require.NoError(t, q.Insert(&models.Service{ + ServiceID: "service-on-replica", + ServiceType: models.MySQLServiceType, + ServiceName: "MySQL on the replica", + NodeID: "ha-node-2", + Address: new("mysql.example.com"), + Port: new(uint16(3306)), + })) + + err := db.InTransaction(func(tx *reform.TX) error { + return models.RemoveStaleHANode(tx.Querier, "ha-node-2") + }) + tests.AssertGRPCErrorCode(t, codes.FailedPrecondition, err) + + assertNodeExists(t, q, "ha-node-2") + _, err = models.FindServiceByID(q, "service-on-replica") + require.NoError(t, err) + }) + t.Run("RefusesThePreHAPMMServerNode", func(t *testing.T) { db := setup(t) q := db.Querier From dcaf91eb8c65991f4e30f8b9e133de69c8768660 Mon Sep 17 00:00:00 2001 From: Ante Gulin Date: Tue, 11 Aug 2026 08:15:07 +0200 Subject: [PATCH 14/23] PMM-15227 Skip blank PMM_HA_PEERS entries --- managed/models/node_helpers.go | 6 ++++++ managed/models/node_helpers_test.go | 13 +++++++++++++ 2 files changed, 19 insertions(+) diff --git a/managed/models/node_helpers.go b/managed/models/node_helpers.go index 751e6822812..9a61631dbed 100644 --- a/managed/models/node_helpers.go +++ b/managed/models/node_helpers.go @@ -356,6 +356,12 @@ func StaleHANodes(q *reform.Querier, haNodeID string, haPeers []string) ([]*Node expected := make(map[string]struct{}, len(haPeers)) for _, peer := range haPeers { + // A trailing comma in PMM_HA_PEERS, or a blank element in the list the chart joins, yields an + // empty entry. It names no replica, so unlike an unreadable one it hides nothing. + if strings.TrimSpace(peer) == "" { + continue + } + name, ok := haPeerNodeName(peer) if !ok { // Trusting the rest would treat a partial list as the whole cluster and remove live replicas. diff --git a/managed/models/node_helpers_test.go b/managed/models/node_helpers_test.go index da52eb35c73..ecfdf85c06c 100644 --- a/managed/models/node_helpers_test.go +++ b/managed/models/node_helpers_test.go @@ -361,6 +361,18 @@ func TestStaleHANodes(t *testing.T) { assertStale(t, stale, "ha-node-2") }) + t.Run("ReportsScaledDownReplicaDespiteBlankPeers", func(t *testing.T) { + q, teardown := setup(t) + defer teardown(t) + + // a trailing comma in PMM_HA_PEERS, or a blank element in the list the chart joins + peers := []string{"pmm-ha-0.pmm-ha:9761", "pmm-ha-1.pmm-ha:9761", "", " "} + stale, err := models.StaleHANodes(q, "pmm-ha-1", peers) + require.NoError(t, err) + + assertStale(t, stale, "ha-node-2") + }) + t.Run("ReportsNothingWhenNothingWasScaledDown", func(t *testing.T) { q, teardown := setup(t) defer teardown(t) @@ -506,6 +518,7 @@ func TestStaleHANodes(t *testing.T) { {"pmm-ha-1.pmm-ha:9761", "pmm-ha-2/10.0.0.2"}, // memberlist "name/address" form {"pmm-ha-1.pmm-ha:9761", "2001:db8::7"}, // an unbracketed IPv6 entry hides a live replica {"pmm-ha-1.pmm-ha:9761", "[2001:db8::7]:9761"}, + {"", " "}, // only blank entries, so nothing is left to compare against nil, } { stale, err := models.StaleHANodes(q, "pmm-ha-1", peers) From 18a6c9b54b485dd2c4d9db3e112fca67f75b3bed Mon Sep 17 00:00:00 2001 From: Ante Gulin Date: Tue, 11 Aug 2026 08:46:13 +0200 Subject: [PATCH 15/23] PMM-15227 Select only PMM Server nodes instead of all --- managed/models/node_helpers.go | 25 +++++++++++++++++-------- managed/models/node_helpers_test.go | 18 ++++++++++++++++++ 2 files changed, 35 insertions(+), 8 deletions(-) diff --git a/managed/models/node_helpers.go b/managed/models/node_helpers.go index 9a61631dbed..2caf20237df 100644 --- a/managed/models/node_helpers.go +++ b/managed/models/node_helpers.go @@ -92,15 +92,28 @@ func CheckUniqueNodeAddressRegion(q *reform.Querier, address string, region *str type NodeFilters struct { // Return Nodes with provided type. NodeType *NodeType + // Return only Nodes that are (or are not) PMM Server Nodes. + IsPMMServerNode *bool } // FindNodes returns Nodes by filters. func FindNodes(q *reform.Querier, filters NodeFilters) ([]*Node, error) { - var whereClause string + var conditions []string var args []any + idx := 1 if filters.NodeType != nil { - whereClause = "WHERE node_type = $1" + conditions = append(conditions, "node_type = "+q.Placeholder(idx)) args = append(args, *filters.NodeType) + idx++ + } + if filters.IsPMMServerNode != nil { + conditions = append(conditions, "is_pmm_server_node = "+q.Placeholder(idx)) + args = append(args, *filters.IsPMMServerNode) + // idx++ + } + var whereClause string + if len(conditions) != 0 { + whereClause = "WHERE " + strings.Join(conditions, " AND ") } structs, err := q.SelectAllFrom(NodeTable, whereClause+" ORDER BY node_id", args...) if err != nil { @@ -376,18 +389,14 @@ func StaleHANodes(q *reform.Querier, haNodeID string, haPeers []string) ([]*Node return nil, nil } - nodes, err := FindNodes(q, NodeFilters{}) + // Only PMM Server Nodes can be stale replicas; the rest are Nodes the user monitors. + nodes, err := FindNodes(q, NodeFilters{IsPMMServerNode: new(true)}) if err != nil { return nil, fmt.Errorf("failed to list Nodes for stale HA node cleanup: %w", err) } var stale []*Node for _, node := range nodes { - // Set by HA replicas, and by the PMM Server Node of a non-HA deployment; every other - // Node is one the user monitors. - if !node.IsPMMServerNode { - continue - } // The PMM Server Node of a deployment converted from non-HA: HA replicas always get a // generated Node ID, and removeNode bans this one outright. Compared against the const, // not PMMServerNodeID: setupPMMServerHAAgents reassigns that var, and SetupDB is retried. diff --git a/managed/models/node_helpers_test.go b/managed/models/node_helpers_test.go index ecfdf85c06c..1829f38799e 100644 --- a/managed/models/node_helpers_test.go +++ b/managed/models/node_helpers_test.go @@ -221,6 +221,24 @@ func TestNodeHelpers(t *testing.T) { require.Equal(t, expected, nodes) }) + t.Run("FindNodesByIsPMMServerNode", func(t *testing.T) { + q, teardown := setup(t) + defer teardown(t) + + nodes, err := models.FindNodes(q, models.NodeFilters{IsPMMServerNode: new(true)}) + require.NoError(t, err) + require.Len(t, nodes, 1) + assert.Equal(t, models.PMMServerNodeID, nodes[0].NodeID) + + nodes, err = models.FindNodes(q, models.NodeFilters{IsPMMServerNode: new(false)}) + require.NoError(t, err) + nodeIDs := make([]string, len(nodes)) + for i, node := range nodes { + nodeIDs[i] = node.NodeID + } + assert.Equal(t, []string{"EmptyNode", "GenericNode", "MySQLNode", "NodeWithPMMAgent"}, nodeIDs) + }) + t.Run("RemoveNode", func(t *testing.T) { q, teardown := setup(t) defer teardown(t) From 41682365e9147d8c5b5029ac1448632c04e9ddbf Mon Sep 17 00:00:00 2001 From: Ante Gulin Date: Tue, 11 Aug 2026 08:48:17 +0200 Subject: [PATCH 16/23] PMM-15227 Document when a stale Node is kept --- .../docs/install-pmm/install-HA-clustered.md | 15 ++++++++++++++- 1 file changed, 14 insertions(+), 1 deletion(-) diff --git a/documentation/docs/install-pmm/install-HA-clustered.md b/documentation/docs/install-pmm/install-HA-clustered.md index f3650c784b3..1be65a05e43 100644 --- a/documentation/docs/install-pmm/install-HA-clustered.md +++ b/documentation/docs/install-pmm/install-HA-clustered.md @@ -829,7 +829,20 @@ When you scale PMM HA up or down, **all PMM pods will be recreated**. This happe - HAProxy continues routing to available pods during rollout - No data loss (distributed storage) - Rolling update strategy minimizes downtime - - The Nodes of removed replicas get removed from **Inventory > Nodes** once the remaining pods restart. A Node is kept (and a warning logged) if it still monitors services, for example a remote instance that was added from that replica. Re-add those services from a running client, then remove the Node from **Inventory > Nodes** + - The Nodes of removed replicas get removed from **Inventory > Nodes** once the remaining pods restart + +!!! info "When PMM keeps a stale Node" + PMM logs a warning (`component=ha`) and keeps the Node when: + + - the Node still monitors services, for example a remote instance that was added from that replica. Re-add those services from a running client, then remove the Node from **Inventory > Nodes** + - `PMM_HA_PEERS` carries no readable node names, for example bare IP addresses + - `PMM_HA_PEERS` does not list the pod that is doing the cleanup + + To see what was skipped: + + ```sh + kubectl exec -n pmm -c pmm-ha -- grep "stale HA node" /srv/logs/pmm-managed.log + ``` To scale PMM server replicas: From 326f45ecf7370dc4ec2c4b4229741332d93bf369 Mon Sep 17 00:00:00 2001 From: Ante Gulin Date: Tue, 11 Aug 2026 08:51:13 +0200 Subject: [PATCH 17/23] PMM-15227 Fix typo --- managed/models/node_helpers.go | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/managed/models/node_helpers.go b/managed/models/node_helpers.go index 2caf20237df..ddd69cf9a85 100644 --- a/managed/models/node_helpers.go +++ b/managed/models/node_helpers.go @@ -473,7 +473,7 @@ func haPeerNodeName(peer string) (string, bool) { // haNodeMonitoredServices returns the IDs of the Services that removing the Node would take with it: // those attached to the Node, those whose exporters run on it (an external exporter in pull mode), and // those whose exporters run under its pmm-agent (remote instances bind theirs to the replica that -// added them. All three are cascaded away by removeNode. +// added them). All three are cascaded away by removeNode. func haNodeMonitoredServices(q *reform.Querier, nodeID string) ([]string, error) { // An external exporter carries a service_id itself, and the pmm-agents are the parents of the // exporters read next. From b7b543ed1f1b4ad4cd129a6861e3e0e7c10b3f2c Mon Sep 17 00:00:00 2001 From: Ante Gulin Date: Tue, 11 Aug 2026 09:07:01 +0200 Subject: [PATCH 18/23] PMM-15227 Improve docs for stale node removal --- documentation/docs/install-pmm/install-HA-clustered.md | 6 +++--- 1 file changed, 3 insertions(+), 3 deletions(-) diff --git a/documentation/docs/install-pmm/install-HA-clustered.md b/documentation/docs/install-pmm/install-HA-clustered.md index 1be65a05e43..d6d31f1775e 100644 --- a/documentation/docs/install-pmm/install-HA-clustered.md +++ b/documentation/docs/install-pmm/install-HA-clustered.md @@ -829,7 +829,7 @@ When you scale PMM HA up or down, **all PMM pods will be recreated**. This happe - HAProxy continues routing to available pods during rollout - No data loss (distributed storage) - Rolling update strategy minimizes downtime - - The Nodes of removed replicas get removed from **Inventory > Nodes** once the remaining pods restart + - The Nodes of removed replicas get removed from **Inventory > Nodes** once the remaining pods restart, unless one of the conditions in the note below applies !!! info "When PMM keeps a stale Node" PMM logs a warning (`component=ha`) and keeps the Node when: @@ -841,7 +841,7 @@ When you scale PMM HA up or down, **all PMM pods will be recreated**. This happe To see what was skipped: ```sh - kubectl exec -n pmm -c pmm-ha -- grep "stale HA node" /srv/logs/pmm-managed.log + kubectl exec -n pmm -c pmm-ha -- grep -i "stale HA node" /srv/logs/pmm-managed.log ``` To scale PMM server replicas: @@ -1220,4 +1220,4 @@ This Tech Preview release is designed to gather community feedback before GA. Yo - What works well in your environment? - What's challenging or confusing? - What features are you missing? -- How does performance compare to single-instance deployments? \ No newline at end of file +- How does performance compare to single-instance deployments? From 25e9b85187a9bd0a071086e5cef12e75bf9cd3e8 Mon Sep 17 00:00:00 2001 From: Ante Gulin Date: Tue, 11 Aug 2026 09:12:00 +0200 Subject: [PATCH 19/23] PMM-15227 Add test for scale-to-one-replica sweep scenario --- managed/models/node_helpers_test.go | 23 +++++++++++++++++++++++ 1 file changed, 23 insertions(+) diff --git a/managed/models/node_helpers_test.go b/managed/models/node_helpers_test.go index 1829f38799e..4deae83b666 100644 --- a/managed/models/node_helpers_test.go +++ b/managed/models/node_helpers_test.go @@ -391,6 +391,29 @@ func TestStaleHANodes(t *testing.T) { assertStale(t, stale, "ha-node-2") }) + t.Run("ReportsEveryDepartedReplicaWhenScaledToOne", func(t *testing.T) { + q, teardown := setup(t) + defer teardown(t) + + // The surviving replica's own Node. It goes here rather than into insertHAFixtures, where + // ReportsNothingWhenNothingWasScaledDown would then see it as stale. + require.NoError(t, q.Insert(&models.Node{ + NodeID: "ha-node-0", + NodeType: models.GenericNodeType, + NodeName: "pmm-ha-0", + Address: models.LocalhostAddr, + IsPMMServerNode: true, + })) + + // what the chart renders at replicas: 1 - a single entry, and it is this pod + peers := []string{"pmm-ha-0.monitoring-service.pmm.svc.cluster.local"} + stale, err := models.StaleHANodes(q, "pmm-ha-0", peers) + require.NoError(t, err) + + // both departed replicas in one sweep, the survivor's own Node untouched + assertStale(t, stale, "ha-node-1", "ha-node-2") + }) + t.Run("ReportsNothingWhenNothingWasScaledDown", func(t *testing.T) { q, teardown := setup(t) defer teardown(t) From 79324be60535e6fec6dd714b6f53da1055d81997 Mon Sep 17 00:00:00 2001 From: Ante Gulin Date: Tue, 11 Aug 2026 10:34:14 +0200 Subject: [PATCH 20/23] PMM-15227 Bind the HA node cleanup to the startup ctx --- managed/models/database.go | 11 +++++++---- 1 file changed, 7 insertions(+), 4 deletions(-) diff --git a/managed/models/database.go b/managed/models/database.go index 410dd241a68..ba6488891f2 100644 --- a/managed/models/database.go +++ b/managed/models/database.go @@ -1297,7 +1297,7 @@ func SetupDB(ctx context.Context, sqlDB *sql.DB, params SetupDBParams) (*reform. return nil, err } - removeStaleHANodes(db, params) + removeStaleHANodes(ctx, db, params) return db, nil } @@ -1522,14 +1522,14 @@ func migrateDB(db *reform.DB, params SetupDBParams) error { // removeStaleHANodes drops the Inventory Nodes of HA replicas that were scaled away. Those rows are // cosmetic, so this runs outside the migration transaction and only logs failures: tidying them up // must never keep a replica from starting. -func removeStaleHANodes(db *reform.DB, params SetupDBParams) { +func removeStaleHANodes(ctx context.Context, db *reform.DB, params SetupDBParams) { if params.HANodeID == "" || params.SetupFixtures == SkipFixtures { return } l := logrus.WithFields(logrus.Fields{"component": "ha", "ha_node_id": params.HANodeID}) - nodes, err := StaleHANodes(db.Querier, params.HANodeID, params.HAPeers) + nodes, err := StaleHANodes(db.WithContext(ctx), params.HANodeID, params.HAPeers) if err != nil { l.WithError(err).Warn("Failed to look for stale HA nodes.") return @@ -1540,7 +1540,7 @@ func removeStaleHANodes(db *reform.DB, params SetupDBParams) { // A transaction per Node: a failure rolls that Node back whole instead of leaving it // half-removed, and leaves the Nodes this sweep hasn't reached yet alone. - err := db.InTransaction(func(tx *reform.TX) error { + err := db.InTransactionContext(ctx, nil, func(tx *reform.TX) error { return RemoveStaleHANode(tx.Querier, node.NodeID) }) switch { @@ -1548,6 +1548,9 @@ func removeStaleHANodes(db *reform.DB, params SetupDBParams) { nodeL.Info("Removed stale HA node, it is not a part of the cluster anymore.") case errors.Is(err, reform.ErrNoRows), status.Code(err) == codes.NotFound: nodeL.Info("Stale HA node was already removed by another replica.") + case errors.Is(err, context.Canceled), errors.Is(err, context.DeadlineExceeded): + nodeL.WithError(err).Warn("Startup was cancelled, stopping the removal of stale HA nodes.") + return default: nodeL.WithError(err).Warn("Failed to remove a stale HA node, keeping it.") } From 940d5d60a79360d533b222ec5fe7dfb4938a6e4b Mon Sep 17 00:00:00 2001 From: Ante Gulin Date: Tue, 11 Aug 2026 11:02:15 +0200 Subject: [PATCH 21/23] PMM-15227 Fix the advice for a Node we keep The warning and the HA scaling doc both told the operator to re-add the services and then delete the Node from Inventory. Every user-facing delete path goes through RemoveNode, which refuses PMM Server Nodes - that refusal is what this ticket works around - so the advice ended at the wall the feature exists to remove. Point at the next restart's sweep instead, and say the same thing in both places. Two doc comments overstated the cascade: removeNode deletes Services attached to the Node, but for Services whose exporters merely run on it only the exporters go, and RemoveStaleHANode can never delete a Service because its own re-check refuses such a Node. Keep the error on the "already removed by another replica" branch: a NotFound from a child agent mid-cascade lands there too, and that case rolls back and leaves the Node in Inventory. --- .../docs/install-pmm/install-HA-clustered.md | 2 +- managed/models/database.go | 2 +- managed/models/node_helpers.go | 16 ++++++++-------- 3 files changed, 10 insertions(+), 10 deletions(-) diff --git a/documentation/docs/install-pmm/install-HA-clustered.md b/documentation/docs/install-pmm/install-HA-clustered.md index d6d31f1775e..c4cd8e052f5 100644 --- a/documentation/docs/install-pmm/install-HA-clustered.md +++ b/documentation/docs/install-pmm/install-HA-clustered.md @@ -834,7 +834,7 @@ When you scale PMM HA up or down, **all PMM pods will be recreated**. This happe !!! info "When PMM keeps a stale Node" PMM logs a warning (`component=ha`) and keeps the Node when: - - the Node still monitors services, for example a remote instance that was added from that replica. Re-add those services from a running client, then remove the Node from **Inventory > Nodes** + - the Node still monitors services, for example a remote instance that was added from that replica. Re-add those services from a running replica; the next restart removes the Node - `PMM_HA_PEERS` carries no readable node names, for example bare IP addresses - `PMM_HA_PEERS` does not list the pod that is doing the cleanup diff --git a/managed/models/database.go b/managed/models/database.go index ba6488891f2..8305e712602 100644 --- a/managed/models/database.go +++ b/managed/models/database.go @@ -1547,7 +1547,7 @@ func removeStaleHANodes(ctx context.Context, db *reform.DB, params SetupDBParams case err == nil: nodeL.Info("Removed stale HA node, it is not a part of the cluster anymore.") case errors.Is(err, reform.ErrNoRows), status.Code(err) == codes.NotFound: - nodeL.Info("Stale HA node was already removed by another replica.") + nodeL.WithError(err).Info("Stale HA node was already removed by another replica.") case errors.Is(err, context.Canceled), errors.Is(err, context.DeadlineExceeded): nodeL.WithError(err).Warn("Startup was cancelled, stopping the removal of stale HA nodes.") return diff --git a/managed/models/node_helpers.go b/managed/models/node_helpers.go index ddd69cf9a85..21fc1a5bea7 100644 --- a/managed/models/node_helpers.go +++ b/managed/models/node_helpers.go @@ -418,7 +418,7 @@ func StaleHANodes(q *reform.Querier, haNodeID string, haPeers []string) ([]*Node } if len(monitored) != 0 { nodeL.WithField("service_ids", monitored).Warn("Keeping stale HA node: it still monitors services, which would be removed with it. " + - "Re-add them from a running replica and remove the node from Inventory.") + "Re-add them from a running replica; the next restart removes the node.") continue } @@ -428,9 +428,9 @@ func StaleHANodes(q *reform.Querier, haNodeID string, haPeers []string) ([]*Node return stale, nil } -// RemoveStaleHANode removes a Node returned by StaleHANodes, along with the Agents and Services that -// belong to it. Call it in a transaction of its own: it deletes those before the Node itself, so a -// failure half-way through would otherwise leave the Node partially removed. +// RemoveStaleHANode removes a Node returned by StaleHANodes together with its Agents. Call it in a +// transaction of its own: it deletes those before the Node itself, so a failure half-way through +// would otherwise leave the Node partially removed. // // It lifts the ban on removing PMM Server Nodes, so it re-reads the Services the Node monitors and // refuses to take a Node that has any. @@ -470,10 +470,10 @@ func haPeerNodeName(peer string) (string, bool) { return label, true } -// haNodeMonitoredServices returns the IDs of the Services that removing the Node would take with it: -// those attached to the Node, those whose exporters run on it (an external exporter in pull mode), and -// those whose exporters run under its pmm-agent (remote instances bind theirs to the replica that -// added them). All three are cascaded away by removeNode. +// haNodeMonitoredServices returns the IDs of the Services that removing the Node would damage: those +// attached to the Node, which removeNode deletes outright, and those whose exporters run on it (an +// external exporter in pull mode) or under its pmm-agent (remote instances bind theirs to the replica +// that added them), which survive but lose their monitoring. A Node with any of them is kept. func haNodeMonitoredServices(q *reform.Querier, nodeID string) ([]string, error) { // An external exporter carries a service_id itself, and the pmm-agents are the parents of the // exporters read next. From a8e4a511d39241834ee8ceec4eeb7bdf6dcbb79a Mon Sep 17 00:00:00 2001 From: Ante Gulin Date: Tue, 11 Aug 2026 11:08:59 +0200 Subject: [PATCH 22/23] PMM-15227 Close three gaps in the HA cleanup tests --- managed/models/node_helpers_test.go | 27 +++++++++++++++++++++++++++ 1 file changed, 27 insertions(+) diff --git a/managed/models/node_helpers_test.go b/managed/models/node_helpers_test.go index 4deae83b666..c419937606f 100644 --- a/managed/models/node_helpers_test.go +++ b/managed/models/node_helpers_test.go @@ -602,6 +602,33 @@ func TestRemoveStaleHANode(t *testing.T) { for _, nodeID := range []string{"ha-node-1", "monitored-node", models.PMMServerNodeID} { assertNodeExists(t, q, nodeID) } + + // the cascade must not reach past the stale Node: the live replica's pmm-agent, and PMM + // Server's own Service and pmm-agent, are untouched + _, err = models.FindAgentByID(q, "ha-agent-1") + require.NoError(t, err) + _, err = models.FindServiceByName(q, models.PMMServerPostgreSQLServiceName) + require.NoError(t, err) + _, err = models.FindAgentByID(q, models.PMMServerAgentID) + require.NoError(t, err) + }) + + t.Run("ReportsNotFoundWhenAlreadyRemoved", func(t *testing.T) { + db := setup(t) + q := db.Querier + + require.NoError(t, db.InTransaction(func(tx *reform.TX) error { + return models.RemoveStaleHANode(tx.Querier, "ha-node-2") + })) + + // what a replica sees when another one won the race; removeStaleHANodes reads this as + // "already removed by another replica" rather than as a failure + err := db.InTransaction(func(tx *reform.TX) error { + return models.RemoveStaleHANode(tx.Querier, "ha-node-2") + }) + tests.AssertGRPCErrorCode(t, codes.NotFound, err) + + assertNodeExists(t, q, "ha-node-1") }) t.Run("LeavesTheNodeWholeWhenRemovalFails", func(t *testing.T) { From 0a3f7b408e5f29e45e744b294312a2d258e81004 Mon Sep 17 00:00:00 2001 From: Ante Gulin Date: Tue, 11 Aug 2026 11:20:56 +0200 Subject: [PATCH 23/23] PMM-15227 Cover the HA peer parser directly haPeerNodeName decides the whole sweep: an entry it reads no name from stops the cleanup and keeps every Node, while a name it does read is trusted as a live replica. Until now it was only reachable through DB-backed subtests asserting on the sweep's result, where a failure says "elements differ" rather than naming the input. The table goes in models_test.go, the package's existing internal test, so no database is involved and each case gets its own subtest. Three shapes the DB-backed table never reached are covered: bare IPv4 without a port, a bracketed IPv6 without a port, and memberlist's name/address form carrying an IPv6 address. --- managed/models/models_test.go | 33 +++++++++++++++++++++++++++++++++ 1 file changed, 33 insertions(+) diff --git a/managed/models/models_test.go b/managed/models/models_test.go index ae1764eb7ce..52b350e7aff 100644 --- a/managed/models/models_test.go +++ b/managed/models/models_test.go @@ -69,3 +69,36 @@ func TestLabels(t *testing.T) { tests.AssertGRPCError(t, status.New(codes.InvalidArgument, `Invalid label name "__1".`), err) }) } + +// The two outcomes are not symmetric: an entry that yields no name stops the whole sweep, while one +// that yields a name is trusted as naming a live replica. Reading a name out of an entry that carries +// none would turn "keep every Node" into "remove every Node this entry didn't name". +func TestHAPeerNodeName(t *testing.T) { + for _, tc := range []struct { + peer string + name string + ok bool + }{ + {peer: "pmm-ha-0.monitoring-service.pmm.svc.cluster.local", name: "pmm-ha-0", ok: true}, // what the chart renders + {peer: "pmm-ha-0.pmm-ha:9761", name: "pmm-ha-0", ok: true}, + {peer: " pmm-ha-1.pmm-ha.pmm.svc.cluster.local ", name: "pmm-ha-1", ok: true}, // trimmed + {peer: "pmm-ha-2:9761", name: "pmm-ha-2", ok: true}, // a dotless host with a port + {peer: "pmm-ha-2", name: "pmm-ha-2", ok: true}, + {peer: "10.244.1.7"}, // bare IPv4, with and without a port + {peer: "10.244.1.7:9761"}, + {peer: "2001:db8::7"}, // IPv6, unbracketed and bracketed + {peer: "[2001:db8::7]:9761"}, + {peer: "[2001:db8::7]"}, + {peer: "pmm-ha-2/10.0.0.2"}, // memberlist's "name/address" form + {peer: "pmm-ha-2/[2001:db8::7]:9761"}, + {peer: ":9761"}, + {peer: ""}, + {peer: " "}, + } { + t.Run(tc.peer, func(t *testing.T) { + name, ok := haPeerNodeName(tc.peer) + assert.Equal(t, tc.ok, ok) + assert.Equal(t, tc.name, name) + }) + } +}