diff --git a/api/descriptor.bin b/api/descriptor.bin index e6481fc092b..67db13f5863 100644 Binary files a/api/descriptor.bin and b/api/descriptor.bin differ diff --git a/api/management/v1/azure.pb.go b/api/management/v1/azure.pb.go index 04a05c47ec1..0266389ebc3 100644 --- a/api/management/v1/azure.pb.go +++ b/api/management/v1/azure.pb.go @@ -381,8 +381,10 @@ type AddAzureDatabaseRequest struct { Type DiscoverAzureDatabaseType `protobuf:"varint,25,opt,name=type,proto3,enum=management.v1.DiscoverAzureDatabaseType" json:"type,omitempty"` // Connection timeout for exporter (if set). ConnectionTimeout *durationpb.Duration `protobuf:"bytes,26,opt,name=connection_timeout,json=connectionTimeout,proto3" json:"connection_timeout,omitempty"` - unknownFields protoimpl.UnknownFields - sizeCache protoimpl.SizeCache + // The pmm-agent identifier which should run agents. Defaults to the PMM Server's own pmm-agent. + PmmAgentId string `protobuf:"bytes,27,opt,name=pmm_agent_id,json=pmmAgentId,proto3" json:"pmm_agent_id,omitempty"` + unknownFields protoimpl.UnknownFields + sizeCache protoimpl.SizeCache } func (x *AddAzureDatabaseRequest) Reset() { @@ -597,6 +599,13 @@ func (x *AddAzureDatabaseRequest) GetConnectionTimeout() *durationpb.Duration { return nil } +func (x *AddAzureDatabaseRequest) GetPmmAgentId() string { + if x != nil { + return x.PmmAgentId + } + return "" +} + type AddAzureDatabaseResponse struct { state protoimpl.MessageState `protogen:"open.v1"` unknownFields protoimpl.UnknownFields @@ -658,7 +667,8 @@ const file_management_v1_azure_proto_rawDesc = "" + "node_model\x18\n" + " \x01(\tR\tnodeModel\"\x85\x01\n" + "\x1dDiscoverAzureDatabaseResponse\x12d\n" + - "\x17azure_database_instance\x18\x01 \x03(\v2,.management.v1.DiscoverAzureDatabaseInstanceR\x15azureDatabaseInstance\"\xfc\t\n" + + "\x17azure_database_instance\x18\x01 \x03(\v2,.management.v1.DiscoverAzureDatabaseInstanceR\x15azureDatabaseInstance\"\x9e\n" + + "\n" + "\x17AddAzureDatabaseRequest\x12\x1f\n" + "\x06region\x18\x01 \x01(\tB\a\xfaB\x04r\x02\x10\x01R\x06region\x12\x0e\n" + "\x02az\x18\x02 \x01(\tR\x02az\x12(\n" + @@ -688,7 +698,9 @@ const file_management_v1_azure_proto_rawDesc = "" + "\x16disable_query_examples\x18\x17 \x01(\bR\x14disableQueryExamples\x12?\n" + "\x1ctablestats_group_table_limit\x18\x18 \x01(\x05R\x19tablestatsGroupTableLimit\x12<\n" + "\x04type\x18\x19 \x01(\x0e2(.management.v1.DiscoverAzureDatabaseTypeR\x04type\x12R\n" + - "\x12connection_timeout\x18\x1a \x01(\v2\x19.google.protobuf.DurationB\b\xfaB\x05\xaa\x01\x022\x00R\x11connectionTimeout\x1a?\n" + + "\x12connection_timeout\x18\x1a \x01(\v2\x19.google.protobuf.DurationB\b\xfaB\x05\xaa\x01\x022\x00R\x11connectionTimeout\x12 \n" + + "\fpmm_agent_id\x18\x1b \x01(\tR\n" + + "pmmAgentId\x1a?\n" + "\x11CustomLabelsEntry\x12\x10\n" + "\x03key\x18\x01 \x01(\tR\x03key\x12\x14\n" + "\x05value\x18\x02 \x01(\tR\x05value:\x028\x01\"\x1a\n" + diff --git a/api/management/v1/azure.pb.validate.go b/api/management/v1/azure.pb.validate.go index be225936500..c66ee420f6e 100644 --- a/api/management/v1/azure.pb.validate.go +++ b/api/management/v1/azure.pb.validate.go @@ -637,6 +637,8 @@ func (m *AddAzureDatabaseRequest) validate(all bool) error { } } + // no validation rules for PmmAgentId + if len(errors) > 0 { return AddAzureDatabaseRequestMultiError(errors) } diff --git a/api/management/v1/azure.proto b/api/management/v1/azure.proto index 2d6eb7b5622..8f7d7aca3be 100644 --- a/api/management/v1/azure.proto +++ b/api/management/v1/azure.proto @@ -130,6 +130,8 @@ message AddAzureDatabaseRequest { google.protobuf.Duration connection_timeout = 26 [(validate.rules).duration = { gte: {seconds: 0} }]; + // The pmm-agent identifier which should run agents. Defaults to the PMM Server's own pmm-agent. + string pmm_agent_id = 27; } message AddAzureDatabaseResponse {} diff --git a/api/management/v1/json/client/management_service/add_azure_database_responses.go b/api/management/v1/json/client/management_service/add_azure_database_responses.go index 44e631ba826..9363d4f9dfb 100644 --- a/api/management/v1/json/client/management_service/add_azure_database_responses.go +++ b/api/management/v1/json/client/management_service/add_azure_database_responses.go @@ -272,6 +272,9 @@ type AddAzureDatabaseBody struct { // Connection timeout for exporter (if set). ConnectionTimeout string `json:"connection_timeout,omitempty"` + + // The pmm-agent identifier which should run agents. Defaults to the PMM Server's own pmm-agent. + PMMAgentID string `json:"pmm_agent_id,omitempty"` } // Validate validates this add azure database body diff --git a/api/management/v1/json/client/management_service/get_node_responses.go b/api/management/v1/json/client/management_service/get_node_responses.go index 9c9faa6b669..bf4972126b2 100644 --- a/api/management/v1/json/client/management_service/get_node_responses.go +++ b/api/management/v1/json/client/management_service/get_node_responses.go @@ -584,6 +584,10 @@ type GetNodeOKBodyNode struct { // True if this node is a PMM Server node (HA mode). IsPMMServerNode bool `json:"is_pmm_server_node,omitempty"` + + // True if this node belongs to the internal infrastructure of a PMM deployment + // (e.g. the HA persistence layer) and must not host user monitoring workloads. + IsPMMInternalNode bool `json:"is_pmm_internal_node,omitempty"` } // Validate validates this get node OK body node diff --git a/api/management/v1/json/client/management_service/list_nodes_responses.go b/api/management/v1/json/client/management_service/list_nodes_responses.go index 1f3aba69c3a..2bf6834181b 100644 --- a/api/management/v1/json/client/management_service/list_nodes_responses.go +++ b/api/management/v1/json/client/management_service/list_nodes_responses.go @@ -593,6 +593,10 @@ type ListNodesOKBodyNodesItems0 struct { // True if this node is a PMM Server node (HA mode). IsPMMServerNode bool `json:"is_pmm_server_node,omitempty"` + + // True if this node belongs to the internal infrastructure of a PMM deployment + // (e.g. the HA persistence layer) and must not host user monitoring workloads. + IsPMMInternalNode bool `json:"is_pmm_internal_node,omitempty"` } // Validate validates this list nodes OK body nodes items0 diff --git a/api/management/v1/json/v1.json b/api/management/v1/json/v1.json index 1b83e01532d..664c306093f 100644 --- a/api/management/v1/json/v1.json +++ b/api/management/v1/json/v1.json @@ -803,6 +803,11 @@ "description": "True if this node is a PMM Server node (HA mode).", "type": "boolean", "x-order": 18 + }, + "is_pmm_internal_node": { + "description": "True if this node belongs to the internal infrastructure of a PMM deployment\n(e.g. the HA persistence layer) and must not host user monitoring workloads.", + "type": "boolean", + "x-order": 19 } } }, @@ -1360,6 +1365,11 @@ "description": "True if this node is a PMM Server node (HA mode).", "type": "boolean", "x-order": 18 + }, + "is_pmm_internal_node": { + "description": "True if this node belongs to the internal infrastructure of a PMM deployment\n(e.g. the HA persistence layer) and must not host user monitoring workloads.", + "type": "boolean", + "x-order": 19 } }, "x-order": 0 @@ -7131,6 +7141,11 @@ "description": "Connection timeout for exporter (if set).", "type": "string", "x-order": 25 + }, + "pmm_agent_id": { + "description": "The pmm-agent identifier which should run agents. Defaults to the PMM Server's own pmm-agent.", + "type": "string", + "x-order": 26 } } } diff --git a/api/management/v1/node.pb.go b/api/management/v1/node.pb.go index 149fb30ad94..7402fee8c0b 100644 --- a/api/management/v1/node.pb.go +++ b/api/management/v1/node.pb.go @@ -619,8 +619,11 @@ type UniversalNode struct { InstanceId string `protobuf:"bytes,18,opt,name=instance_id,json=instanceId,proto3" json:"instance_id,omitempty"` // True if this node is a PMM Server node (HA mode). IsPmmServerNode bool `protobuf:"varint,19,opt,name=is_pmm_server_node,json=isPmmServerNode,proto3" json:"is_pmm_server_node,omitempty"` - unknownFields protoimpl.UnknownFields - sizeCache protoimpl.SizeCache + // True if this node belongs to the internal infrastructure of a PMM deployment + // (e.g. the HA persistence layer) and must not host user monitoring workloads. + IsPmmInternalNode bool `protobuf:"varint,20,opt,name=is_pmm_internal_node,json=isPmmInternalNode,proto3" json:"is_pmm_internal_node,omitempty"` + unknownFields protoimpl.UnknownFields + sizeCache protoimpl.SizeCache } func (x *UniversalNode) Reset() { @@ -786,6 +789,13 @@ func (x *UniversalNode) GetIsPmmServerNode() bool { return false } +func (x *UniversalNode) GetIsPmmInternalNode() bool { + if x != nil { + return x.IsPmmInternalNode + } + return false +} + type ListNodesRequest struct { state protoimpl.MessageState `protogen:"open.v1"` // Node type to be filtered out. @@ -1159,7 +1169,7 @@ const file_management_v1_node_proto_rawDesc = "" + "\anode_id\x18\x01 \x01(\tB\a\xfaB\x04r\x02\x10\x01R\x06nodeId\x12\x14\n" + "\x05force\x18\x02 \x01(\bR\x05force\"2\n" + "\x16UnregisterNodeResponse\x12\x18\n" + - "\awarning\x18\x01 \x01(\tR\awarning\"\x9d\t\n" + + "\awarning\x18\x01 \x01(\tR\awarning\"\xce\t\n" + "\rUniversalNode\x12\x17\n" + "\anode_id\x18\x01 \x01(\tR\x06nodeId\x12\x1b\n" + "\tnode_type\x18\x02 \x01(\tR\bnodeType\x12\x1b\n" + @@ -1185,7 +1195,8 @@ const file_management_v1_node_proto_rawDesc = "" + "\bservices\x18\x11 \x03(\v2$.management.v1.UniversalNode.ServiceR\bservices\x12\x1f\n" + "\vinstance_id\x18\x12 \x01(\tR\n" + "instanceId\x12+\n" + - "\x12is_pmm_server_node\x18\x13 \x01(\bR\x0fisPmmServerNode\x1an\n" + + "\x12is_pmm_server_node\x18\x13 \x01(\bR\x0fisPmmServerNode\x12/\n" + + "\x14is_pmm_internal_node\x18\x14 \x01(\bR\x11isPmmInternalNode\x1an\n" + "\aService\x12\x1d\n" + "\n" + "service_id\x18\x01 \x01(\tR\tserviceId\x12!\n" + diff --git a/api/management/v1/node.pb.validate.go b/api/management/v1/node.pb.validate.go index 34411b3f938..0e606deeec4 100644 --- a/api/management/v1/node.pb.validate.go +++ b/api/management/v1/node.pb.validate.go @@ -906,6 +906,8 @@ func (m *UniversalNode) validate(all bool) error { // no validation rules for IsPmmServerNode + // no validation rules for IsPmmInternalNode + if len(errors) > 0 { return UniversalNodeMultiError(errors) } diff --git a/api/management/v1/node.proto b/api/management/v1/node.proto index d5e26d3bc57..e4cdfa5c0b0 100644 --- a/api/management/v1/node.proto +++ b/api/management/v1/node.proto @@ -165,6 +165,9 @@ message UniversalNode { string instance_id = 18; // True if this node is a PMM Server node (HA mode). bool is_pmm_server_node = 19; + // True if this node belongs to the internal infrastructure of a PMM deployment + // (e.g. the HA persistence layer) and must not host user monitoring workloads. + bool is_pmm_internal_node = 20; } message ListNodesRequest { diff --git a/api/swagger/swagger-dev.json b/api/swagger/swagger-dev.json index 29fa37229f6..e7da061958b 100644 --- a/api/swagger/swagger-dev.json +++ b/api/swagger/swagger-dev.json @@ -22326,6 +22326,11 @@ "description": "True if this node is a PMM Server node (HA mode).", "type": "boolean", "x-order": 18 + }, + "is_pmm_internal_node": { + "description": "True if this node belongs to the internal infrastructure of a PMM deployment\n(e.g. the HA persistence layer) and must not host user monitoring workloads.", + "type": "boolean", + "x-order": 19 } } }, @@ -22883,6 +22888,11 @@ "description": "True if this node is a PMM Server node (HA mode).", "type": "boolean", "x-order": 18 + }, + "is_pmm_internal_node": { + "description": "True if this node belongs to the internal infrastructure of a PMM deployment\n(e.g. the HA persistence layer) and must not host user monitoring workloads.", + "type": "boolean", + "x-order": 19 } }, "x-order": 0 @@ -28654,6 +28664,11 @@ "description": "Connection timeout for exporter (if set).", "type": "string", "x-order": 25 + }, + "pmm_agent_id": { + "description": "The pmm-agent identifier which should run agents. Defaults to the PMM Server's own pmm-agent.", + "type": "string", + "x-order": 26 } } } diff --git a/api/swagger/swagger.json b/api/swagger/swagger.json index 56c5ccb51d7..e085c0e9477 100644 --- a/api/swagger/swagger.json +++ b/api/swagger/swagger.json @@ -21353,6 +21353,11 @@ "description": "True if this node is a PMM Server node (HA mode).", "type": "boolean", "x-order": 18 + }, + "is_pmm_internal_node": { + "description": "True if this node belongs to the internal infrastructure of a PMM deployment\n(e.g. the HA persistence layer) and must not host user monitoring workloads.", + "type": "boolean", + "x-order": 19 } } }, @@ -21910,6 +21915,11 @@ "description": "True if this node is a PMM Server node (HA mode).", "type": "boolean", "x-order": 18 + }, + "is_pmm_internal_node": { + "description": "True if this node belongs to the internal infrastructure of a PMM deployment\n(e.g. the HA persistence layer) and must not host user monitoring workloads.", + "type": "boolean", + "x-order": 19 } }, "x-order": 0 @@ -27681,6 +27691,11 @@ "description": "Connection timeout for exporter (if set).", "type": "string", "x-order": 25 + }, + "pmm_agent_id": { + "description": "The pmm-agent identifier which should run agents. Defaults to the PMM Server's own pmm-agent.", + "type": "string", + "x-order": 26 } } } diff --git a/managed/cmd/pmm-managed/main.go b/managed/cmd/pmm-managed/main.go index 227f86921c3..79a3bd38d07 100644 --- a/managed/cmd/pmm-managed/main.go +++ b/managed/cmd/pmm-managed/main.go @@ -234,6 +234,20 @@ type gRPCServerDeps struct { versionCache *versioncache.Service vmdb *victoriametrics.Service vmalert *vmalert.Service + internalNodePrefixes []string +} + +// parseNodeNamePrefixes splits a comma-separated list of Node name prefixes. +func parseNodeNamePrefixes(value string) []string { + var prefixes []string + for p := range strings.SplitSeq(value, ",") { + p = strings.TrimSpace(p) + if p != "" { + prefixes = append(prefixes, p) + } + } + + return prefixes } // runGRPCServer runs gRPC server until context is canceled, then gracefully stops it. @@ -304,6 +318,8 @@ func runGRPCServer(ctx context.Context, deps *gRPCServerDeps) { deps.db, deps.agentsRegistry, deps.agentsStateUpdater, deps.connectionCheck, deps.serviceInfoBroker, deps.vmdb, deps.versionCache, deps.grafanaClient, v1.NewAPI(*deps.vmClient), + deps.internalNodePrefixes, + deps.ha.Params().Enabled, ) managementv1.RegisterManagementServiceServer(gRPCServer, managementSvc) @@ -743,6 +759,11 @@ func main() { //nolint:gocognit,maintidx,cyclop Default("9762"). Int() + internalNodePrefixesF := kingpin.Flag("internal-node-name-prefixes", + "Comma-separated list of Node name prefixes reserved for the internal infrastructure of this PMM deployment"). + Envar("PMM_INTERNAL_NODE_NAME_PREFIXES"). + String() + supervisordConfigDirF := kingpin.Flag("supervisord-config-dir", "Supervisord configuration directory").Required().String() logLevelF := kingpin.Flag("log-level", "Set logging level").Envar("PMM_LOG_LEVEL").Default("info").Enum("trace", "debug", "info", "warn", "error", "fatal") @@ -1192,6 +1213,7 @@ func main() { //nolint:gocognit,maintidx,cyclop grafanaClient: grafanaClient, handler: agentsHandler, ha: haService, + internalNodePrefixes: parseNodeNamePrefixes(*internalNodePrefixesF), jobsService: jobsService, minioClient: minioClient, pbmPITRService: pbmPITRService, diff --git a/managed/cmd/pmm-managed/main_test.go b/managed/cmd/pmm-managed/main_test.go index 17e536726f0..6b309ed1337 100644 --- a/managed/cmd/pmm-managed/main_test.go +++ b/managed/cmd/pmm-managed/main_test.go @@ -206,3 +206,17 @@ func formatPkgName(t *testing.T, name string) string { return name } + +func TestParseNodeNamePrefixes(t *testing.T) { + for _, tc := range []struct { + value string + expected []string + }{ + {value: "", expected: nil}, + {value: ",,", expected: nil}, + {value: "pmm-pmm-ha-pg-db-", expected: []string{"pmm-pmm-ha-pg-db-"}}, + {value: " pmm-pmm-ha-pg-db- , pmm-pmm-ha-ch- ", expected: []string{"pmm-pmm-ha-pg-db-", "pmm-pmm-ha-ch-"}}, + } { + assert.Equal(t, tc.expected, parseNodeNamePrefixes(tc.value), tc.value) + } +} diff --git a/managed/services/management/add_service_exporter_timeout_test.go b/managed/services/management/add_service_exporter_timeout_test.go index 73b5ada35b1..9a6123a405f 100644 --- a/managed/services/management/add_service_exporter_timeout_test.go +++ b/managed/services/management/add_service_exporter_timeout_test.go @@ -75,7 +75,7 @@ func TestAddServiceExporterTimeout(t *testing.T) { vmClient.AssertExpectations(t) }) - s := NewManagementService(db, ar, state, cc, sib, vmdb, vc, grafanaClient, vmClient) + s := NewManagementService(db, ar, state, cc, sib, vmdb, vc, grafanaClient, vmClient, nil, false) want := durationpb.New(17 * time.Second) t.Run("MySQL", func(t *testing.T) { diff --git a/managed/services/management/agent_test.go b/managed/services/management/agent_test.go index 676ccaf1b2f..e3e352c7232 100644 --- a/managed/services/management/agent_test.go +++ b/managed/services/management/agent_test.go @@ -97,7 +97,7 @@ func setup(t *testing.T) (context.Context, *ManagementService, func(t *testing.T vmClient.AssertExpectations(t) } - s := NewManagementService(db, ar, state, cc, sib, vmdb, vc, grafanaClient, vmClient) + s := NewManagementService(db, ar, state, cc, sib, vmdb, vc, grafanaClient, vmClient, nil, false) return ctx, s, teardown } diff --git a/managed/services/management/annotation_test.go b/managed/services/management/annotation_test.go index 88398ce1dbb..203b109b48e 100644 --- a/managed/services/management/annotation_test.go +++ b/managed/services/management/annotation_test.go @@ -66,7 +66,7 @@ func TestAnnotations(t *testing.T) { vmClient := &mockVictoriaMetricsClient{} vmClient.Test(t) - s := NewManagementService(db, ar, state, cc, sib, vmdb, vc, grafanaClient, vmClient) + s := NewManagementService(db, ar, state, cc, sib, vmdb, vc, grafanaClient, vmClient, nil, false) teardown := func(t *testing.T) { t.Helper() diff --git a/managed/services/management/azure_database.go b/managed/services/management/azure_database.go index 71c61928b4c..eca8a43b4b3 100644 --- a/managed/services/management/azure_database.go +++ b/managed/services/management/azure_database.go @@ -195,6 +195,16 @@ func (s *ManagementService) AddAzureDatabase(ctx context.Context, req *managemen } l := logger.Get(ctx).WithField("component", "discover/azureDatabase") + + pmmAgentID := models.PMMServerAgentID + if req.GetPmmAgentId() != "" { + pmmAgentID = req.GetPmmAgentId() + } + err := s.checkNodeIsEligible(ctx, pmmAgentID, req.Address) + if err != nil { + return nil, err + } + // tweak according to API docs if req.NodeName == "" { req.NodeName = req.InstanceId @@ -260,7 +270,7 @@ func (s *ManagementService) AddAzureDatabase(ctx context.Context, req *managemen if req.AzureDatabaseExporter { azureDatabaseExporter, err := models.CreateAgent(tx.Querier, models.AzureDatabaseExporterType, &models.CreateAgentParams{ - PMMAgentID: models.PMMServerAgentID, + PMMAgentID: pmmAgentID, ServiceID: service.ServiceID, AzureOptions: models.AzureOptionsFromRequest(req), }) @@ -271,7 +281,7 @@ func (s *ManagementService) AddAzureDatabase(ctx context.Context, req *managemen } metricsExporter, err := models.CreateAgent(tx.Querier, exporterType, &models.CreateAgentParams{ - PMMAgentID: models.PMMServerAgentID, + PMMAgentID: pmmAgentID, ServiceID: service.ServiceID, Username: req.Username, Password: req.Password, @@ -302,7 +312,7 @@ func (s *ManagementService) AddAzureDatabase(ctx context.Context, req *managemen if req.Qan { qanAgent, err := models.CreateAgent(tx.Querier, qanAgentType, &models.CreateAgentParams{ - PMMAgentID: models.PMMServerAgentID, + PMMAgentID: pmmAgentID, ServiceID: service.ServiceID, Username: req.Username, Password: req.Password, @@ -324,6 +334,6 @@ func (s *ManagementService) AddAzureDatabase(ctx context.Context, req *managemen return nil, e } - s.state.RequestStateUpdate(ctx, models.PMMServerAgentID) + s.state.RequestStateUpdate(ctx, pmmAgentID) return &managementv1.AddAzureDatabaseResponse{}, nil } diff --git a/managed/services/management/azure_database_test.go b/managed/services/management/azure_database_test.go new file mode 100644 index 00000000000..87eb313e7b3 --- /dev/null +++ b/managed/services/management/azure_database_test.go @@ -0,0 +1,93 @@ +// Copyright (C) 2023 Percona LLC +// +// This program is free software: you can redistribute it and/or modify +// it under the terms of the GNU Affero General Public License as published by +// the Free Software Foundation, either version 3 of the License, or +// (at your option) any later version. +// +// This program is distributed in the hope that it will be useful, +// but WITHOUT ANY WARRANTY; without even the implied warranty of +// MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the +// GNU Affero General Public License for more details. +// +// You should have received a copy of the GNU Affero General Public License +// along with this program. If not, see . + +package management + +import ( + "testing" + + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" + "gopkg.in/reform.v1" + "gopkg.in/reform.v1/dialects/postgresql" + + managementv1 "github.com/percona/pmm/api/management/v1" + "github.com/percona/pmm/managed/models" + "github.com/percona/pmm/managed/utils/testdb" + "github.com/percona/pmm/utils/logger" +) + +// TestAddAzureDatabaseRunsOnRequestedAgent covers the Agent an Azure Database is delegated to. +// The Agents are created under the requested pmm-agent, so it is the one which has to be told +// about them, otherwise it does not start monitoring until its next state update. +func TestAddAzureDatabaseRunsOnRequestedAgent(t *testing.T) { + ctx := logger.Set(t.Context(), t.Name()) + + sqlDB := testdb.Open(t, models.SetupFixtures, nil) + t.Cleanup(func() { + require.NoError(t, sqlDB.Close()) + }) + db := reform.NewDB(sqlDB, postgresql.Dialect, reform.NewPrintfLogger(t.Logf)) + + _, err := models.UpdateSettings(sqlDB, &models.ChangeSettingsParams{ + EnableAzurediscover: new(true), + }) + require.NoError(t, err) + + node, err := models.CreateNode(db.Querier, models.GenericNodeType, &models.CreateNodeParams{ + NodeName: "azure-delegate", + Address: "10.2.3.4", + }) + require.NoError(t, err) + agent, err := models.CreatePMMAgent(db.Querier, node.NodeID, nil) + require.NoError(t, err) + + state := &mockAgentsStateUpdater{} + state.Test(t) + state.On("RequestStateUpdate", ctx, agent.AgentID).Once() + t.Cleanup(func() { + state.AssertExpectations(t) + }) + + s := NewManagementService(db, nil, state, nil, nil, nil, nil, nil, nil, nil, false) + + res, err := s.AddAzureDatabase(ctx, &managementv1.AddAzureDatabaseRequest{ + PmmAgentId: agent.AgentID, + Region: "westeurope", + InstanceId: "azure-mysql-instance", + Address: "test.mysql.database.azure.com", + Port: 3306, + Username: "azure-user", + Type: managementv1.DiscoverAzureDatabaseType_DISCOVER_AZURE_DATABASE_TYPE_MYSQL, + AzureDatabaseExporter: true, + Qan: true, + SkipConnectionCheck: true, + }) + require.NoError(t, err) + require.NotNil(t, res) + + agents, err := models.FindAgents(db.Querier, models.AgentFilters{PMMAgentID: agent.AgentID}) + require.NoError(t, err) + + agentTypes := make([]models.AgentType, 0, len(agents)) + for _, a := range agents { + agentTypes = append(agentTypes, a.AgentType) + } + assert.ElementsMatch(t, []models.AgentType{ + models.AzureDatabaseExporterType, + models.MySQLdExporterType, + models.QANMySQLPerfSchemaAgentType, + }, agentTypes) +} diff --git a/managed/services/management/internal_node_test.go b/managed/services/management/internal_node_test.go new file mode 100644 index 00000000000..e51d53d8eb1 --- /dev/null +++ b/managed/services/management/internal_node_test.go @@ -0,0 +1,325 @@ +// Copyright (C) 2023 Percona LLC +// +// This program is free software: you can redistribute it and/or modify +// it under the terms of the GNU Affero General Public License as published by +// the Free Software Foundation, either version 3 of the License, or +// (at your option) any later version. +// +// This program is distributed in the hope that it will be useful, +// but WITHOUT ANY WARRANTY; without even the implied warranty of +// MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the +// GNU Affero General Public License for more details. +// +// You should have received a copy of the GNU Affero General Public License +// along with this program. If not, see . + +package management + +import ( + "fmt" + "testing" + + "github.com/prometheus/common/model" + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/mock" + "github.com/stretchr/testify/require" + "google.golang.org/grpc/codes" + "google.golang.org/grpc/status" + "gopkg.in/reform.v1" + "gopkg.in/reform.v1/dialects/postgresql" + + managementv1 "github.com/percona/pmm/api/management/v1" + "github.com/percona/pmm/managed/models" + "github.com/percona/pmm/managed/utils/testdb" + "github.com/percona/pmm/managed/utils/tests" + "github.com/percona/pmm/utils/logger" +) + +// internalNodePrefix mimics what the PMM HA Helm chart reports for its PostgreSQL cluster: +// the Nodes are named "-" by the PostgreSQL operator. +const internalNodePrefix = "pmm-pmm-ha-pg-db-" + +func TestIsInternalNode(t *testing.T) { + s := &ManagementService{internalNodePrefixes: []string{internalNodePrefix, "pmm-pmm-ha-ch-"}} + + for nodeName, expected := range map[string]bool{ + "pmm-pmm-ha-pg-db-instance1-qjjl-0": true, + "pmm-pmm-ha-ch-0": true, + "pmm-ha-0": false, + "pmm-server": false, + "": false, + } { + assert.Equal(t, expected, s.isInternalNode(&models.Node{NodeName: nodeName}), nodeName) + } + + t.Run("no prefixes configured", func(t *testing.T) { + s := &ManagementService{} + assert.False(t, s.isInternalNode(&models.Node{NodeName: internalNodePrefix + "instance1-qjjl-0"})) + }) + + // The PMM Server Nodes of an HA deployment serve PMM itself. A single-node deployment has no + // other Node to delegate monitoring to, so its Node stays available. + t.Run("PMM Server Nodes", func(t *testing.T) { + serverNode := &models.Node{NodeName: "pmm-ha-0", IsPMMServerNode: true} + clientNode := &models.Node{NodeName: "pmm-pmm-ha-client-0"} + + ha := &ManagementService{haEnabled: true} + assert.True(t, ha.isInternalNode(serverNode)) + assert.False(t, ha.isInternalNode(clientNode)) + + singleNode := &ManagementService{} + assert.False(t, singleNode.isInternalNode(&models.Node{NodeName: "pmm-server", IsPMMServerNode: true})) + }) +} + +func TestAddServiceTarget(t *testing.T) { + const ( + agentID = "00000000-0000-4000-8000-000000000005" + address = "mysql.example.com" + ) + + for _, tc := range []struct { + name string + req *managementv1.AddServiceRequest + expectedAgentID string + expectedAddress string + }{ + { + name: "MySQL", + req: &managementv1.AddServiceRequest{Service: &managementv1.AddServiceRequest_Mysql{ + Mysql: &managementv1.AddMySQLServiceParams{PmmAgentId: agentID, Address: address}, + }}, + expectedAgentID: agentID, + expectedAddress: address, + }, + { + name: "MongoDB", + req: &managementv1.AddServiceRequest{Service: &managementv1.AddServiceRequest_Mongodb{ + Mongodb: &managementv1.AddMongoDBServiceParams{PmmAgentId: agentID, Address: address}, + }}, + expectedAgentID: agentID, + expectedAddress: address, + }, + { + name: "PostgreSQL", + req: &managementv1.AddServiceRequest{Service: &managementv1.AddServiceRequest_Postgresql{ + Postgresql: &managementv1.AddPostgreSQLServiceParams{PmmAgentId: agentID, Address: address}, + }}, + expectedAgentID: agentID, + expectedAddress: address, + }, + { + name: "ProxySQL", + req: &managementv1.AddServiceRequest{Service: &managementv1.AddServiceRequest_Proxysql{ + Proxysql: &managementv1.AddProxySQLServiceParams{PmmAgentId: agentID, Address: address}, + }}, + expectedAgentID: agentID, + expectedAddress: address, + }, + { + name: "Valkey", + req: &managementv1.AddServiceRequest{Service: &managementv1.AddServiceRequest_Valkey{ + Valkey: &managementv1.AddValkeyServiceParams{PmmAgentId: agentID, Address: address}, + }}, + expectedAgentID: agentID, + expectedAddress: address, + }, + { + name: "RDS", + req: &managementv1.AddServiceRequest{Service: &managementv1.AddServiceRequest_Rds{ + Rds: &managementv1.AddRDSServiceParams{PmmAgentId: agentID, Address: address}, + }}, + expectedAgentID: agentID, + expectedAddress: address, + }, + { + name: "External Services are scraped on the Node they run on", + req: &managementv1.AddServiceRequest{Service: &managementv1.AddServiceRequest_External{ + External: &managementv1.AddExternalServiceParams{RunsOnNodeId: "00000000-0000-4000-8000-000000000006"}, + }}, + }, + { + name: "HAProxy Services are scraped on the Node they run on", + req: &managementv1.AddServiceRequest{Service: &managementv1.AddServiceRequest_Haproxy{ + Haproxy: &managementv1.AddHAProxyServiceParams{NodeId: "00000000-0000-4000-8000-000000000006"}, + }}, + }, + } { + t.Run(tc.name, func(t *testing.T) { + pmmAgentID, address := addServiceTarget(tc.req) + assert.Equal(t, tc.expectedAgentID, pmmAgentID) + assert.Equal(t, tc.expectedAddress, address) + }) + } +} + +func TestListNodesMarksInternalNodes(t *testing.T) { + ctx := logger.Set(t.Context(), t.Name()) + + sqlDB := testdb.Open(t, models.SetupFixtures, nil) + t.Cleanup(func() { + require.NoError(t, sqlDB.Close()) + }) + db := reform.NewDB(sqlDB, postgresql.Dialect, reform.NewPrintfLogger(t.Logf)) + + node, err := models.CreateNode(db.Querier, models.GenericNodeType, &models.CreateNodeParams{ + NodeName: internalNodePrefix + "instance1-qjjl-0", + Address: "10.1.2.3", + }) + require.NoError(t, err) + + ar := &mockAgentsRegistry{} + ar.Test(t) + ar.On("IsConnected", mock.Anything).Return(false) + + vmdb := &mockPrometheusService{} + vmdb.Test(t) + + vmClient := &mockVictoriaMetricsClient{} + vmClient.Test(t) + vmClient.On("Query", ctx, mock.Anything, mock.Anything).Return(model.Vector{}, nil, nil) + + s := NewManagementService(db, ar, nil, nil, nil, vmdb, nil, nil, vmClient, []string{internalNodePrefix}, false) + + res, err := s.ListNodes(ctx, &managementv1.ListNodesRequest{}) + require.NoError(t, err) + + isInternal := make(map[string]bool, len(res.Nodes)) + for _, n := range res.Nodes { + isInternal[n.NodeName] = n.IsPmmInternalNode + } + assert.True(t, isInternal[node.NodeName], node.NodeName) + assert.False(t, isInternal["pmm-server"]) +} + +// TestListNodesMarksPMMServerNodesInternalInHA covers the PMM Server Nodes of an HA deployment, +// which serve PMM itself and are reported as internal so that the UI does not offer them. +// The fixtures register the PMM Server Node the same way a deployment does. +func TestListNodesMarksPMMServerNodesInternalInHA(t *testing.T) { + ctx := logger.Set(t.Context(), t.Name()) + + sqlDB := testdb.Open(t, models.SetupFixtures, nil) + t.Cleanup(func() { + require.NoError(t, sqlDB.Close()) + }) + db := reform.NewDB(sqlDB, postgresql.Dialect, reform.NewPrintfLogger(t.Logf)) + + client, err := models.CreateNode(db.Querier, models.GenericNodeType, &models.CreateNodeParams{ + NodeName: "pmm-pmm-ha-client-0", + Address: "10.1.2.4", + }) + require.NoError(t, err) + + newService := func(haEnabled bool) *ManagementService { + ar := &mockAgentsRegistry{} + ar.Test(t) + ar.On("IsConnected", mock.Anything).Return(false) + + vmdb := &mockPrometheusService{} + vmdb.Test(t) + + vmClient := &mockVictoriaMetricsClient{} + vmClient.Test(t) + vmClient.On("Query", ctx, mock.Anything, mock.Anything).Return(model.Vector{}, nil, nil) + + return NewManagementService(db, ar, nil, nil, nil, vmdb, nil, nil, vmClient, nil, haEnabled) + } + + listNodes := func(t *testing.T, haEnabled bool) map[string]bool { + t.Helper() + + res, err := newService(haEnabled).ListNodes(ctx, &managementv1.ListNodesRequest{}) + require.NoError(t, err) + + isInternal := make(map[string]bool, len(res.Nodes)) + for _, n := range res.Nodes { + isInternal[n.NodeName] = n.IsPmmInternalNode + } + + return isInternal + } + + t.Run("HA hides the PMM Server Node", func(t *testing.T) { + isInternal := listNodes(t, true) + assert.True(t, isInternal["pmm-server"]) + assert.False(t, isInternal[client.NodeName], client.NodeName) + }) + + t.Run("a single-node deployment keeps it", func(t *testing.T) { + isInternal := listNodes(t, false) + assert.False(t, isInternal["pmm-server"]) + assert.False(t, isInternal[client.NodeName], client.NodeName) + }) +} + +func TestCheckNodeIsEligible(t *testing.T) { + ctx := logger.Set(t.Context(), t.Name()) + + sqlDB := testdb.Open(t, models.SetupFixtures, nil) + t.Cleanup(func() { + require.NoError(t, sqlDB.Close()) + }) + db := reform.NewDB(sqlDB, postgresql.Dialect, reform.NewPrintfLogger(t.Logf)) + + node, err := models.CreateNode(db.Querier, models.GenericNodeType, &models.CreateNodeParams{ + NodeName: internalNodePrefix + "instance1-qjjl-0", + Address: "10.1.2.3", + }) + require.NoError(t, err) + agent, err := models.CreatePMMAgent(db.Querier, node.NodeID, nil) + require.NoError(t, err) + + s := NewManagementService(db, nil, nil, nil, nil, nil, nil, nil, nil, []string{internalNodePrefix}, false) + expectedErr := status.New(codes.FailedPrecondition, fmt.Sprintf( + "Node '%s' is a part of the internal infrastructure of this PMM deployment and cannot monitor other services.", node.NodeName, + )) + + t.Run("a remote address on an internal Node is rejected", func(t *testing.T) { + err := s.checkNodeIsEligible(ctx, agent.AgentID, "mysql.example.com") + tests.AssertGRPCError(t, expectedErr, err) + }) + + t.Run("local addresses on an internal Node are allowed", func(t *testing.T) { + for _, address := range []string{"", "localhost", "127.0.0.1", "::1"} { + assert.NoError(t, s.checkNodeIsEligible(ctx, agent.AgentID, address), address) + } + }) + + t.Run("a remote address on a regular Node is allowed", func(t *testing.T) { + assert.NoError(t, s.checkNodeIsEligible(ctx, models.PMMServerAgentID, "mysql.example.com")) + }) + + t.Run("no prefixes configured", func(t *testing.T) { + s := NewManagementService(db, nil, nil, nil, nil, nil, nil, nil, nil, nil, false) + assert.NoError(t, s.checkNodeIsEligible(ctx, agent.AgentID, "mysql.example.com")) + }) + + t.Run("AddService rejects an internal Node", func(t *testing.T) { + res, err := s.AddService(ctx, &managementv1.AddServiceRequest{Service: &managementv1.AddServiceRequest_Mysql{ + Mysql: &managementv1.AddMySQLServiceParams{ + PmmAgentId: agent.AgentID, + ServiceName: "test-mysql", + Address: "mysql.example.com", + Port: 3306, + }, + }}) + assert.Nil(t, res) + tests.AssertGRPCError(t, expectedErr, err) + }) + + t.Run("AddAzureDatabase rejects an internal Node", func(t *testing.T) { + _, err := models.UpdateSettings(sqlDB, &models.ChangeSettingsParams{ + EnableAzurediscover: new(true), + }) + require.NoError(t, err) + + res, err := s.AddAzureDatabase(ctx, &managementv1.AddAzureDatabaseRequest{ + PmmAgentId: agent.AgentID, + InstanceId: "test-azure", + Address: "test.mysql.database.azure.com", + Port: 3306, + }) + assert.Nil(t, res) + tests.AssertGRPCError(t, expectedErr, err) + }) +} diff --git a/managed/services/management/node.go b/managed/services/management/node.go index 5e8e13bc399..172bd660710 100644 --- a/managed/services/management/node.go +++ b/managed/services/management/node.go @@ -316,22 +316,23 @@ func (s *ManagementService) ListNodes(ctx context.Context, req *managementv1.Lis } uNode := &managementv1.UniversalNode{ - Address: node.Address, - CustomLabels: labels, - NodeId: node.NodeID, - NodeName: node.NodeName, - NodeType: string(node.NodeType), - Az: node.AZ, - CreatedAt: timestamppb.New(node.CreatedAt), - ContainerId: pointer.GetString(node.ContainerID), - ContainerName: pointer.GetString(node.ContainerName), - Distro: node.Distro, - MachineId: pointer.GetString(node.MachineID), - NodeModel: node.NodeModel, - Region: pointer.GetString(node.Region), - UpdatedAt: timestamppb.New(node.UpdatedAt), - InstanceId: node.InstanceID, - IsPmmServerNode: node.IsPMMServerNode, + Address: node.Address, + CustomLabels: labels, + NodeId: node.NodeID, + NodeName: node.NodeName, + NodeType: string(node.NodeType), + Az: node.AZ, + CreatedAt: timestamppb.New(node.CreatedAt), + ContainerId: pointer.GetString(node.ContainerID), + ContainerName: pointer.GetString(node.ContainerName), + Distro: node.Distro, + MachineId: pointer.GetString(node.MachineID), + NodeModel: node.NodeModel, + Region: pointer.GetString(node.Region), + UpdatedAt: timestamppb.New(node.UpdatedAt), + InstanceId: node.InstanceID, + IsPmmServerNode: node.IsPMMServerNode, + IsPmmInternalNode: s.isInternalNode(node), } freshUp, hasFresh := metrics[node.NodeID] @@ -397,21 +398,22 @@ func (s *ManagementService) GetNode(ctx context.Context, req *managementv1.GetNo } uNode := &managementv1.UniversalNode{ - Address: node.Address, - Az: node.AZ, - CreatedAt: timestamppb.New(node.CreatedAt), - ContainerId: pointer.GetString(node.ContainerID), - ContainerName: pointer.GetString(node.ContainerName), - CustomLabels: labels, - Distro: node.Distro, - MachineId: pointer.GetString(node.MachineID), - NodeId: node.NodeID, - NodeName: node.NodeName, - NodeType: string(node.NodeType), - NodeModel: node.NodeModel, - Region: pointer.GetString(node.Region), - UpdatedAt: timestamppb.New(node.UpdatedAt), - IsPmmServerNode: node.IsPMMServerNode, + Address: node.Address, + Az: node.AZ, + CreatedAt: timestamppb.New(node.CreatedAt), + ContainerId: pointer.GetString(node.ContainerID), + ContainerName: pointer.GetString(node.ContainerName), + CustomLabels: labels, + Distro: node.Distro, + MachineId: pointer.GetString(node.MachineID), + NodeId: node.NodeID, + NodeName: node.NodeName, + NodeType: string(node.NodeType), + NodeModel: node.NodeModel, + Region: pointer.GetString(node.Region), + UpdatedAt: timestamppb.New(node.UpdatedAt), + IsPmmServerNode: node.IsPMMServerNode, + IsPmmInternalNode: s.isInternalNode(node), } freshUp, hasFresh := metrics[node.NodeID] diff --git a/managed/services/management/node_test.go b/managed/services/management/node_test.go index 9b83eecdc58..559b403815e 100644 --- a/managed/services/management/node_test.go +++ b/managed/services/management/node_test.go @@ -88,7 +88,7 @@ func TestNodeService(t *testing.T) { vmClient.AssertExpectations(t) } - s := NewManagementService(db, r, state, nil, nil, vmdb, nil, authProvider, vmClient) + s := NewManagementService(db, r, state, nil, nil, vmdb, nil, authProvider, vmClient, nil, false) return ctx, s, teardown } @@ -272,7 +272,7 @@ func TestNodeService(t *testing.T) { grafanaClient := &mockGrafanaClient{} grafanaClient.Test(t) - s := NewManagementService(db, ar, state, cc, sib, vmdb, vc, grafanaClient, vmClient) + s := NewManagementService(db, ar, state, cc, sib, vmdb, vc, grafanaClient, vmClient, nil, false) teardown := func(t *testing.T) { t.Helper() @@ -549,7 +549,7 @@ func TestNodeService(t *testing.T) { vmClient := &mockVictoriaMetricsClient{} vmClient.Test(t) - s := NewManagementService(db, ar, state, cc, sib, vmdb, vc, grafanaClient, vmClient) + s := NewManagementService(db, ar, state, cc, sib, vmdb, vc, grafanaClient, vmClient, nil, false) teardown := func(t *testing.T) { t.Helper() diff --git a/managed/services/management/rds_test.go b/managed/services/management/rds_test.go index e92220e97fc..707fa17fa9e 100644 --- a/managed/services/management/rds_test.go +++ b/managed/services/management/rds_test.go @@ -81,7 +81,7 @@ func TestRDSService(t *testing.T) { vmClient.AssertExpectations(t) }() - s := NewManagementService(db, ar, state, cc, sib, vmdb, vc, grafanaClient, vmClient) + s := NewManagementService(db, ar, state, cc, sib, vmdb, vc, grafanaClient, vmClient, nil, false) t.Run("DiscoverRDS", func(t *testing.T) { t.Run("ListRegions", func(t *testing.T) { diff --git a/managed/services/management/service.go b/managed/services/management/service.go index c92ef207839..5d948de7012 100644 --- a/managed/services/management/service.go +++ b/managed/services/management/service.go @@ -49,6 +49,12 @@ type ManagementService struct { //nolint:revive grafanaClient grafanaClient vmClient victoriaMetricsClient l *logrus.Entry + + // internalNodePrefixes holds the Node name prefixes reserved for the internal + // infrastructure of this PMM deployment, e.g. its HA persistence layer. + internalNodePrefixes []string + // haEnabled indicates whether this PMM Server is a node of an HA cluster. + haEnabled bool } // upMetricSelectors match the per-service-type "up" metrics that back the service status. @@ -80,19 +86,43 @@ func NewManagementService( vc versionCache, grafanaClient grafanaClient, vmClient victoriaMetricsClient, + internalNodePrefixes []string, + haEnabled bool, ) *ManagementService { return &ManagementService{ - db: db, - r: r, - state: state, - cc: cc, - sib: sib, - vmdb: vmdb, - vc: vc, - grafanaClient: grafanaClient, - vmClient: vmClient, - l: logrus.WithField("service", "management"), + db: db, + r: r, + state: state, + cc: cc, + sib: sib, + vmdb: vmdb, + vc: vc, + grafanaClient: grafanaClient, + vmClient: vmClient, + l: logrus.WithField("service", "management"), + internalNodePrefixes: internalNodePrefixes, + haEnabled: haEnabled, + } +} + +// isInternalNode reports whether the Node belongs to the internal infrastructure of this +// PMM deployment and therefore must not host user monitoring workloads. +// +// In an HA deployment that covers the PMM Server Nodes themselves: they are expected to spend +// their resources on serving PMM, and a Client is pre-provisioned to carry the monitoring instead. +// A single-node deployment keeps its Node available, as it is the only one there is. +func (s *ManagementService) isInternalNode(node *models.Node) bool { + if s.haEnabled && node.IsPMMServerNode { + return true } + + for _, prefix := range s.internalNodePrefixes { + if strings.HasPrefix(node.NodeName, prefix) { + return true + } + } + + return false } // A map to check if the service is supported. @@ -108,8 +138,74 @@ var supportedServices = map[string]inventoryv1.ServiceType{ string(models.HAProxyServiceType): inventoryv1.ServiceType_SERVICE_TYPE_HAPROXY_SERVICE, } +// localAddresses resolve to the Node an Agent runs on. Monitoring them delegates no work to +// that Node beyond what it already does for itself. +var localAddresses = map[string]struct{}{"": {}, "localhost": {}, "127.0.0.1": {}, "::1": {}} + +// checkNodeIsEligible rejects requests which delegate monitoring of a remote address to an +// Agent running on a Node reserved for the internal infrastructure of this PMM deployment. +// Those Nodes still monitor the services inside their own pod, hence the address check. +func (s *ManagementService) checkNodeIsEligible(ctx context.Context, pmmAgentID, address string) error { + if pmmAgentID == "" || (len(s.internalNodePrefixes) == 0 && !s.haEnabled) { + return nil + } + _, isLocal := localAddresses[address] + if isLocal { + return nil + } + + agent, err := models.FindAgentByID(s.db.WithContext(ctx), pmmAgentID) + if err != nil { + return err + } + nodeID := pointer.GetString(agent.RunsOnNodeID) + if nodeID == "" { + return nil + } + + node, err := models.FindNodeByID(s.db.WithContext(ctx), nodeID) + if err != nil { + return err + } + if !s.isInternalNode(node) { + return nil + } + + return status.Errorf(codes.FailedPrecondition, + "Node '%s' is a part of the internal infrastructure of this PMM deployment and cannot monitor other services.", node.NodeName) +} + +// addServiceTarget returns the pmm-agent which is to run the Service's Agents together with the +// address it is to monitor, for the Service types which delegate monitoring to an existing Node. +func addServiceTarget(req *managementv1.AddServiceRequest) (string, string) { + switch req.Service.(type) { + case *managementv1.AddServiceRequest_Mysql: + return req.GetMysql().GetPmmAgentId(), req.GetMysql().GetAddress() + case *managementv1.AddServiceRequest_Mongodb: + return req.GetMongodb().GetPmmAgentId(), req.GetMongodb().GetAddress() + case *managementv1.AddServiceRequest_Postgresql: + return req.GetPostgresql().GetPmmAgentId(), req.GetPostgresql().GetAddress() + case *managementv1.AddServiceRequest_Proxysql: + return req.GetProxysql().GetPmmAgentId(), req.GetProxysql().GetAddress() + case *managementv1.AddServiceRequest_Valkey: + return req.GetValkey().GetPmmAgentId(), req.GetValkey().GetAddress() + case *managementv1.AddServiceRequest_Rds: + return req.GetRds().GetPmmAgentId(), req.GetRds().GetAddress() + default: + // External and HAProxy Services are scraped on the Node their exporter runs on, + // so they cannot offload work onto it. + return "", "" + } +} + // AddService add a Service and its Agents. func (s *ManagementService) AddService(ctx context.Context, req *managementv1.AddServiceRequest) (*managementv1.AddServiceResponse, error) { + pmmAgentID, address := addServiceTarget(req) + err := s.checkNodeIsEligible(ctx, pmmAgentID, address) + if err != nil { + return nil, err + } + switch req.Service.(type) { case *managementv1.AddServiceRequest_Mysql: return s.addMySQL(ctx, req.GetMysql()) diff --git a/managed/services/management/service_test.go b/managed/services/management/service_test.go index f1e93d63695..70fb98af4c9 100644 --- a/managed/services/management/service_test.go +++ b/managed/services/management/service_test.go @@ -89,7 +89,7 @@ func TestServiceService(t *testing.T) { vmClient.AssertExpectations(t) } - s := NewManagementService(db, ar, state, cc, sib, vmdb, vc, grafanaClient, vmClient) + s := NewManagementService(db, ar, state, cc, sib, vmdb, vc, grafanaClient, vmClient, nil, false) return ctx, s, teardown } @@ -325,7 +325,7 @@ func TestServiceService(t *testing.T) { vmClient.AssertExpectations(t) } - s := NewManagementService(db, ar, state, cc, sib, vmdb, vc, grafanaClient, vmClient) + s := NewManagementService(db, ar, state, cc, sib, vmdb, vc, grafanaClient, vmClient, nil, false) return ctx, s, teardown, vmdb } diff --git a/managed/utils/envvars/parser.go b/managed/utils/envvars/parser.go index 953f6ee5f82..ee5377f1a38 100644 --- a/managed/utils/envvars/parser.go +++ b/managed/utils/envvars/parser.go @@ -127,6 +127,9 @@ func ParseEnvVars(envs []string) (*models.ChangeSettingsParams, []error, []strin case "PERCONA_TELEMETRY_DISABLE": // skip the Pillars telemetry environment variable continue + case "PMM_INTERNAL_NODE_NAME_PREFIXES": + // skip the env variable that is already handled by kingpin + continue case "PMM_ENABLE_UPDATES": b, err := strconv.ParseBool(v) if err != nil { diff --git a/managed/utils/envvars/parser_test.go b/managed/utils/envvars/parser_test.go index 9881de244d7..e0a40c42212 100644 --- a/managed/utils/envvars/parser_test.go +++ b/managed/utils/envvars/parser_test.go @@ -150,6 +150,18 @@ func TestEnvVarValidator(t *testing.T) { assert.Nil(t, gotWarns) }) + t.Run("Skipped internal node name prefixes env var", func(t *testing.T) { + t.Parallel() + + envs := []string{"PMM_INTERNAL_NODE_NAME_PREFIXES=pmm-pmm-ha-pg-db-"} + expectedEnvVars := &models.ChangeSettingsParams{} + + gotEnvVars, gotErrs, gotWarns := ParseEnvVars(envs) + assert.Equal(t, expectedEnvVars, gotEnvVars) + assert.Nil(t, gotErrs) + assert.Nil(t, gotWarns) + }) + t.Run("Invalid env variables values", func(t *testing.T) { t.Parallel()