diff --git a/weed/admin/dash/admin_data.go b/weed/admin/dash/admin_data.go index c57b10a56..2d1f1f914 100644 --- a/weed/admin/dash/admin_data.go +++ b/weed/admin/dash/admin_data.go @@ -146,6 +146,8 @@ type FilerNode struct { DataCenter string `json:"datacenter"` Rack string `json:"rack"` LastUpdated time.Time `json:"last_updated"` + // MetricsPort is the node's advertised Prometheus port, 0 when disabled. + MetricsPort uint32 `json:"metrics_port"` } type MessageBrokerNode struct { @@ -159,6 +161,8 @@ type S3Node struct { Address string `json:"address"` DataCenter string `json:"datacenter"` LastUpdated time.Time `json:"last_updated"` + // MetricsPort is the node's advertised Prometheus port, 0 when disabled. + MetricsPort uint32 `json:"metrics_port"` } // GetAdminData retrieves admin data as a struct (for reuse by both JSON and HTML handlers) @@ -356,6 +360,7 @@ func (s *AdminServer) getFilerNodesStatus() []FilerNode { DataCenter: node.DataCenter, Rack: node.Rack, LastUpdated: time.Now(), + MetricsPort: node.MetricsPort, }) } @@ -433,6 +438,7 @@ func (s *AdminServer) getS3NodesStatus() []S3Node { Address: pb.ServerAddress(node.Address).ToHttpAddress(), DataCenter: node.DataCenter, LastUpdated: time.Now(), + MetricsPort: node.MetricsPort, }) } diff --git a/weed/admin/dash/cluster_topology.go b/weed/admin/dash/cluster_topology.go index 21759dd34..fd443e29b 100644 --- a/weed/admin/dash/cluster_topology.go +++ b/weed/admin/dash/cluster_topology.go @@ -207,6 +207,7 @@ func (s *AdminServer) getTopologyViaGRPC(topology *ClusterTopology) error { DiskCapacity: diskCapacity, LastHeartbeat: time.Now(), RemoteSize: remoteSize, + MetricsPort: node.MetricsPort, } rackObj.Nodes = append(rackObj.Nodes, vs) diff --git a/weed/admin/dash/metrics_derive_test.go b/weed/admin/dash/metrics_derive_test.go index fb93c159c..c826d40ad 100644 --- a/weed/admin/dash/metrics_derive_test.go +++ b/weed/admin/dash/metrics_derive_test.go @@ -90,6 +90,25 @@ func TestDeriveCounterResetSkipped(t *testing.T) { } } +func TestMetricsEndpoint(t *testing.T) { + for _, tc := range []struct { + node string + port uint32 + want string + }{ + // The advertised port replaces the node's service port. + {"127.0.0.1:8080", 9327, "127.0.0.1:9327"}, + {"127.0.0.1:8888.18888", 9327, "127.0.0.1:9327"}, + {"[::1]:8080", 9327, "[::1]:9327"}, + // A node without -metricsPort must never be scraped. + {"127.0.0.1:8080", 0, ""}, + } { + if got := metricsEndpoint(tc.node, tc.port); got != tc.want { + t.Errorf("metricsEndpoint(%q, %d) = %q, want %q", tc.node, tc.port, got, tc.want) + } + } +} + func TestStoreRingIsBounded(t *testing.T) { s := newMetricsSeries() for i := 0; i < metricsMaxSamples+50; i++ { diff --git a/weed/admin/dash/metrics_scraper.go b/weed/admin/dash/metrics_scraper.go index 4d57ac7a1..f3c37d689 100644 --- a/weed/admin/dash/metrics_scraper.go +++ b/weed/admin/dash/metrics_scraper.go @@ -4,7 +4,10 @@ import ( "context" "fmt" "io" + "net" "net/http" + "sort" + "strconv" "strings" "time" @@ -12,6 +15,8 @@ import ( "github.com/prometheus/common/expfmt" "github.com/seaweedfs/seaweedfs/weed/glog" + "github.com/seaweedfs/seaweedfs/weed/pb" + "github.com/seaweedfs/seaweedfs/weed/pb/master_pb" stats_collect "github.com/seaweedfs/seaweedfs/weed/stats" util_http "github.com/seaweedfs/seaweedfs/weed/util/http" ) @@ -153,26 +158,95 @@ func (s *AdminServer) scrapeAllServers(ctx context.Context) { } } +// scrapeTarget is one Prometheus endpoint. source is the endpoint address, not +// a component name: a combined "weed server" advertises one listener for +// master, volume, filer and S3 alike, and metric names already identify the +// component. nodes records which cluster members advertised this endpoint, so +// the UI can label it. type scrapeTarget struct { source string address string + nodes []string } +// scrapeTargets lists the distinct metrics endpoints advertised by the cluster. +// Nodes started without -metricsPort advertise 0 and are skipped, so nothing is +// scraped from a client-facing service port. func (s *AdminServer) scrapeTargets() []scrapeTarget { - var out []scrapeTarget - if topo, err := s.GetClusterTopology(); err == nil && topo != nil { - for _, m := range topo.Masters { - out = append(out, scrapeTarget{source: "master/" + m.Address, address: m.Address}) + byAddress := map[string][]string{} + add := func(nodeAddress string, metricsPort uint32) { + endpoint := metricsEndpoint(nodeAddress, metricsPort) + if endpoint == "" { + return } + byAddress[endpoint] = append(byAddress[endpoint], nodeAddress) + } + + for _, m := range s.mastersWithMetricsPort() { + add(m.address, m.metricsPort) + } + if topo, err := s.GetClusterTopology(); err == nil && topo != nil { for _, vs := range topo.VolumeServers { - out = append(out, scrapeTarget{source: "volume/" + vs.Address, address: vs.Address}) + add(vs.Address, vs.MetricsPort) } } for _, f := range s.getFilerNodesStatus() { - out = append(out, scrapeTarget{source: "filer/" + f.Address, address: f.Address}) + add(f.Address, f.MetricsPort) } for _, n := range s.getS3NodesStatus() { - out = append(out, scrapeTarget{source: "s3/" + n.Address, address: n.Address}) + add(n.Address, n.MetricsPort) + } + + out := make([]scrapeTarget, 0, len(byAddress)) + for endpoint, nodes := range byAddress { + sort.Strings(nodes) + out = append(out, scrapeTarget{source: endpoint, address: endpoint, nodes: nodes}) + } + sort.Slice(out, func(i, j int) bool { return out[i].source < out[j].source }) + return out +} + +// metricsEndpoint combines a node's host with its advertised metrics port. +// Returns "" when the node does not run a metrics listener. +func metricsEndpoint(nodeAddress string, metricsPort uint32) string { + if metricsPort == 0 { + return "" + } + host, _, err := net.SplitHostPort(nodeAddress) + if err != nil { + host = nodeAddress + } + return net.JoinHostPort(host, strconv.Itoa(int(metricsPort))) +} + +type masterMetricsTarget struct { + address string + metricsPort uint32 +} + +// mastersWithMetricsPort asks each master for its own metrics port. +// GetMasterConfiguration reports the configuration of the master that answers, +// so it is called per address rather than once via the leader. +func (s *AdminServer) mastersWithMetricsPort() []masterMetricsTarget { + md, err := s.GetClusterMasters() + if err != nil || md == nil { + return nil + } + var out []masterMetricsTarget + for _, m := range md.Masters { + address := m.Address + err := pb.WithMasterClient(context.Background(), false, pb.ServerAddress(address), s.grpcDialOption, false, + func(client master_pb.SeaweedClient) error { + resp, err := client.GetMasterConfiguration(context.Background(), &master_pb.GetMasterConfigurationRequest{}) + if err != nil { + return err + } + out = append(out, masterMetricsTarget{address: address, metricsPort: resp.MetricsPort}) + return nil + }) + if err != nil { + glog.V(1).Infof("master %s configuration: %v", address, err) + } } return out } diff --git a/weed/admin/dash/types.go b/weed/admin/dash/types.go index b30da7961..a268c7d9c 100644 --- a/weed/admin/dash/types.go +++ b/weed/admin/dash/types.go @@ -50,6 +50,8 @@ type VolumeServer struct { DiskUsage int64 `json:"disk_usage"` DiskCapacity int64 `json:"disk_capacity"` LastHeartbeat time.Time `json:"last_heartbeat"` + // MetricsPort is the node's advertised Prometheus port, 0 when disabled. + MetricsPort uint32 `json:"metrics_port"` // EC shard information EcVolumes int `json:"ec_volumes"` // Number of EC volumes this server has shards for