diff --git a/documentation/docs/install-pmm/install-HA-clustered.md b/documentation/docs/install-pmm/install-HA-clustered.md index 704fe512ed2..c4cd8e052f5 100644 --- a/documentation/docs/install-pmm/install-HA-clustered.md +++ b/documentation/docs/install-pmm/install-HA-clustered.md @@ -829,6 +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, 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: + + - 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 + + To see what was skipped: + + ```sh + kubectl exec -n pmm -c pmm-ha -- grep -i "stale HA node" /srv/logs/pmm-managed.log + ``` To scale PMM server replicas: @@ -1206,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? 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/database.go b/managed/models/database.go index fa839fa3be0..8305e712602 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(ctx, db, params) + return db, nil } @@ -1517,6 +1519,44 @@ 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(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.WithContext(ctx), 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.InTransactionContext(ctx, nil, 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.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 + default: + nodeL.WithError(err).Warn("Failed to remove a stale HA node, keeping it.") + } + } +} + type agentConfig struct { ID string `yaml:"id"` } 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) + }) + } +} diff --git a/managed/models/node_helpers.go b/managed/models/node_helpers.go index 84c2e8e551d..21fc1a5bea7 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" @@ -90,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 { @@ -252,13 +267,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 == defaultPMMServerNodeID || (!allowPMMServerNode && n.IsPMMServerNode) { return status.Error(codes.PermissionDenied, "PMM Server node can't be removed.") } @@ -334,3 +355,160 @@ func RemoveNode(q *reform.Querier, id string, mode RemoveMode) error { //nolint: } return nil } + +// 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, nil + } + + l := logrus.WithFields(logrus.Fields{"component": "ha", "ha_node_id": haNodeID}) + + 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. + 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, 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, nil + } + + // 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 { + // 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. + if node.NodeID == defaultPMMServerNodeID { + continue + } + if _, ok := expected[node.NodeName]; ok { + continue + } + + 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 { + 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. " + + "Re-add them from a running replica; the next restart removes the node.") + continue + } + + stale = append(stale, node) + } + + return stale, nil +} + +// 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. +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) +} + +// 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 IPv4 or IPv6 addresses. +func haPeerNodeName(peer string) (string, bool) { + 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, "[" 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 + } + return label, true +} + +// 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. + agents, err := FindAgentsOnNode(q, nodeID) + if err != nil { + return nil, err + } + + var serviceIDs, pmmAgentIDs []string + for _, agent := range agents { + if agent.ServiceID != nil { + serviceIDs = append(serviceIDs, *agent.ServiceID) + } + if agent.AgentType == PMMAgentType { + pmmAgentIDs = append(pmmAgentIDs, agent.AgentID) + } + } + + 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 1bad35df9eb..c419937606f 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) @@ -267,3 +285,452 @@ func TestNodeHelpers(t *testing.T) { require.Len(t, nodes, 2) // PMM Server + HA PMM Server node }) } + +// 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) + require.NoError(t, err) +} + +func TestStaleHANodes(t *testing.T) { + sqlDB := testdb.Open(t, models.SetupFixtures, nil) + t.Cleanup(func() { + require.NoError(t, sqlDB.Close()) + }) + + 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) + insertHAFixtures(t, tx.Querier) + + teardown := func(t *testing.T) { + t.Helper() + require.NoError(t, tx.Rollback()) + } + return tx.Querier, teardown + } + + assertStale := func(t *testing.T, nodes []*models.Node, nodeIDs ...string) { + t.Helper() + actual := make([]string, 0, len(nodes)) + for _, node := range nodes { + actual = append(actual, node.NodeID) + } + assert.ElementsMatch(t, nodeIDs, actual) + } + + 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 "} + stale, err := models.StaleHANodes(q, "pmm-ha-1", peers) + require.NoError(t, err) + + // neither the live replica, the monitored nodes nor the pre-HA pmm-server Node are reported + 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("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) + + // 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"} + stale, err := models.StaleHANodes(q, "pmm-ha-1", peers) + require.NoError(t, err) + + assertStale(t, stale) + }) + + 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"} + stale, err := models.StaleHANodes(q, "pmm-ha-1", peers) + require.NoError(t, err) + + 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) + + // 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"} + stale, err := models.StaleHANodes(q, "pmm-ha-1", peers) + require.NoError(t, err) + + assertStale(t, stale, "ha-node-2") + }) + + 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"} + stale, err := models.StaleHANodes(q, "pmm-ha-1", peers) + require.NoError(t, err) + + assertStale(t, stale, "ha-node-2") + }) + + t.Run("ReportsNothingWhenPeersCantBeTrusted", 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 + {"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) + 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 + + 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) + } + + // 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) { + 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) + } + }) + + 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) + }) + + 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 + + // 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) + }) +} 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. //