Compare commits

...
Author SHA1 Message Date
Chris Lu a6e07181b2 admin: scrape the metrics ports the cluster advertises
Replaces scraping each node's service port with the dedicated Prometheus
listener each node now advertises. Nodes started without -metricsPort
advertise 0 and are skipped, so nothing is ever fetched from a
client-facing port.

Endpoints are deduplicated by address because a combined "weed server"
advertises one listener for all of its components; series are attributed
by metric name, which already identifies the component.
2026-09-14 23:59:38 -07:00
Chris Lu 51831d6850 admin: derive interval rates and latency quantiles from scrapes
Raw counters and histograms are process-lifetime cumulative, so
charting them directly is meaningless. Counters now become per-second
rates and histograms become p50/p95/p99, both computed from the delta
against the previous scrape, with counter resets skipped. Quantiles use
linear interpolation within the matching bucket, as Prometheus does.
2026-09-14 23:59:38 -07:00
Chris Lu d235dd280b admin: gather the admin's own registry into the metrics store
The admin's maintenance and worker metrics live in the local
stats.Gather registry, so record them directly under the admin/local
source instead of scraping over HTTP. Adds metricsStore.match for
prefix/metric lookups.
2026-09-14 23:59:38 -07:00
Chris Lu 7b9332953d admin: add server-side SVG chart renderer
Extends the existing sparklineSVG approach into a full chart helper
with axes, gridlines, multi-series lines/areas, legends, threshold
lines, and unit formatters (bytes, bps, ms, pct). No JS chart library;
safe to inline in templ pages.
2026-09-14 23:59:38 -07:00
Chris Lu 85147522a9 admin: scrape per-server /metrics into the store
Adds a 15s scrape loop that fetches /metrics from every discovered
master, volume, filer, and S3 server (addresses come from the
existing topology + ListClusterNodes helpers) and records each
series into the in-memory store. Parses Prometheus text exposition
via prometheus/common/expfmt.
2026-09-14 23:59:38 -07:00
Chris Lu 207bc0b75a admin: add in-memory metrics series store
Bounded ring buffer (240 samples) keyed by source/metric[/labels],
reusing the dashSample ring pattern. No persistence; powers the
upcoming monitoring charts.
2026-09-14 23:59:38 -07:00
Chris Lu 2f26d5779b mini: wire the shared metrics port, and tolerate an unset one
weed mini builds MasterOptions directly and never set metricsHttpPort, so
reading it in toMasterOption dereferenced nil and crashed startup. Point
mini's master, volume and filer at its single -metricsPort listener, the
same way weed server does, and treat an unset port as disabled so a
partially initialised MasterOptions cannot panic again.

Reproduced with 'weed mini -dir=... -s3.port=...', which is what the S3
filer-group and delete-regression suites start.
2026-09-14 23:59:32 -07:00
Chris Lu 14bc6e5e4f servers: advertise the configured metricsPort to the master
Filer, S3 and broker report it on the KeepConnected registration, volume
servers in their heartbeat, and masters return their own in
GetMasterConfiguration. MasterClient gains SetMetricsPort so the eight
callers that have no metrics listener are untouched, and the volume
server takes it as a constructor argument because its heartbeat goroutine
starts there.

In combined "weed server" one metrics listener serves the whole shared
registry, so every component advertises the same port.
2026-09-14 23:32:32 -07:00
Chris Lu dfde24f3ee pb: add metrics_port so nodes can advertise their metrics listener
Each server already has a -metricsPort Prometheus listener, but nothing
advertises it, so a central scraper cannot find it. Adds metrics_port to
KeepConnectedRequest and ListClusterNodes (filer, S3, broker), Heartbeat
and DataNodeInfo (volume servers), and GetMasterConfigurationResponse
(masters).

Note the existing metrics_address fields are unrelated: they carry the
Prometheus push gateway that servers push to, not a scrape target.
2026-09-14 23:23:25 -07:00
31 changed files with 1059 additions and 66 deletions
+1 -1
View File
@@ -60,7 +60,7 @@ require (
github.com/pquerna/cachecontrol v0.2.0
github.com/prometheus/client_golang v1.24.1
github.com/prometheus/client_model v0.6.3
github.com/prometheus/common v0.70.1 // indirect
github.com/prometheus/common v0.70.1
github.com/prometheus/procfs v0.22.0
github.com/rcrowley/go-metrics v0.0.0-20201227073835-cf1acfcdf475 // indirect
github.com/remyoudompheng/bigfft v0.0.0-20230129092748-24d4a6f8daec // indirect
+6
View File
@@ -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,
})
}
+10
View File
@@ -126,6 +126,13 @@ type AdminServer struct {
dashSamples []dashSample
dashSamplesMu sync.Mutex
// metricsStore holds scraped per-server Prometheus series for the
// monitoring pages. Filled by the scrape loop in startMetricsScraper.
// metricsDeriver turns raw counters and histograms into interval rates
// and latency quantiles before they are stored.
metricsStore *metricsStore
metricsDeriver *metricsDeriver
// Filer discovery and caching
cachedFilers []string
lastFilerUpdate time.Time
@@ -210,6 +217,8 @@ func NewAdminServer(masters string, filerGroup string, templateFS http.FileSyste
pluginLock: lockManager,
adminPresenceLock: presenceLock,
bgCancel: bgCancel,
metricsStore: newMetricsStore(),
metricsDeriver: newMetricsDeriver(),
}
// Initialize topic retention purger
@@ -321,6 +330,7 @@ func NewAdminServer(masters string, filerGroup string, templateFS http.FileSyste
}
go server.publishMaintenanceMetrics(bgCtx)
go server.startMetricsScraper(bgCtx)
return server
}
+174
View File
@@ -0,0 +1,174 @@
package dash
import (
"fmt"
"strings"
)
type chartSeries struct {
name string
color string
data []float64
area bool
}
type chartOptions struct {
unit string
threshold *float64
labels []string
}
func renderChartSVG(series []chartSeries, opts chartOptions) string {
const w, h = 560.0, 190.0
const l, r, t, b = 46.0, 8.0, 10.0, 22.0
if len(series) == 0 || len(series[0].data) < 2 {
return flatChart(w, h, l, r, t, b)
}
n := len(series[0].data)
min, max := computeRange(series, opts.threshold)
if min > 0 {
min = 0
}
if max == min {
max = min + 1
}
span := max - min
px := func(i int) float64 { return l + float64(i)*(w-l-r)/float64(n-1) }
py := func(v float64) float64 { return t + (h-t-b)*(1-(v-min)/span) }
var s strings.Builder
for g := 0; g <= 4; g++ {
v := min + span*float64(g)/4
y := py(v)
fmt.Fprintf(&s, `<line x1="%g" y1="%g" x2="%g" y2="%g" stroke="#e3e6f0" stroke-width="1"/>`, l, y, w-r, y)
fmt.Fprintf(&s, `<text x="%g" y="%g" text-anchor="end" font-size="9" fill="#858796">%s</text>`, l-6, y+3, formatChartValue(v, opts.unit))
}
for _, i := range []int{0, n / 2, n - 1} {
fmt.Fprintf(&s, `<text x="%g" y="%g" text-anchor="middle" font-size="9" fill="#858796">%s</text>`, px(i), h-6, xLabel(opts, i, n))
}
if opts.threshold != nil {
y := py(*opts.threshold)
fmt.Fprintf(&s, `<line x1="%g" y1="%g" x2="%g" y2="%g" stroke="#a5615c" stroke-width="1" stroke-dasharray="4 3"/>`, l, y, w-r, y)
}
for _, se := range series {
var d strings.Builder
for i, v := range se.data {
if i == 0 {
fmt.Fprintf(&d, "M%.1f %.1f", px(i), py(v))
} else {
fmt.Fprintf(&d, " L%.1f %.1f", px(i), py(v))
}
}
if se.area {
fmt.Fprintf(&s, `<path d="%s L%.1f %.1f L%.1f %.1f Z" fill="%s" opacity="0.12"/>`, d.String(), px(n-1), py(0), px(0), py(0), se.color)
}
fmt.Fprintf(&s, `<path d="%s" fill="none" stroke="%s" stroke-width="1.8" stroke-linejoin="round"/>`, d.String(), se.color)
}
return fmt.Sprintf(`<svg viewBox="0 0 %g %g" preserveAspectRatio="xMidYMid meet" style="width:100%%;height:auto">%s</svg>`, w, h, s.String())
}
func flatChart(w, h, l, r, t, b float64) string {
return fmt.Sprintf(`<svg viewBox="0 0 %g %g" preserveAspectRatio="xMidYMid meet" style="width:100%%;height:auto"><line x1="%g" y1="%g" x2="%g" y2="%g" stroke="#e3e6f0" stroke-width="1"/></svg>`, w, h, l, (t + (h-t-b)/2), w-r, (t + (h-t-b)/2))
}
func computeRange(series []chartSeries, threshold *float64) (float64, float64) {
min, max := 1e9, -1e9
for _, se := range series {
for _, v := range se.data {
if v < min {
min = v
}
if v > max {
max = v
}
}
}
if threshold != nil && *threshold > max {
max = *threshold * 1.15
}
return min, max
}
func xLabel(opts chartOptions, i, n int) string {
if i < len(opts.labels) {
return opts.labels[i]
}
return fmt.Sprintf("-%dm", n-1-i)
}
func formatChartValue(v float64, unit string) string {
switch unit {
case "bytes":
return chartFormatBytes(int64(v))
case "bps":
return chartFormatBytes(int64(v)) + "/s"
case "ms":
return fmt.Sprintf("%.0f ms", v)
case "pct":
return fmt.Sprintf("%.0f%%", v)
default:
if v >= 1000 {
return fmt.Sprintf("%.1fk", v/1000)
}
if v == float64(int64(v)) {
return fmt.Sprintf("%d", int64(v))
}
return fmt.Sprintf("%.1f", v)
}
}
func chartFormatBytes(b int64) string {
const u = 1024
if b < u {
return fmt.Sprintf("%d B", b)
}
div, exp := int64(u), 0
for n := b / u; n >= u; n /= u {
div *= u
exp++
}
return fmt.Sprintf("%.1f %cB", float64(b)/float64(div), "KMGTPE"[exp])
}
func renderLegend(series []chartSeries) string {
var b strings.Builder
b.WriteString(`<div class="legend">`)
for _, se := range series {
fmt.Fprintf(&b, `<span class="key"><span class="dot" style="background:%s"></span>%s</span>`, se.color, se.name)
}
b.WriteString(`</div>`)
return b.String()
}
func sampleTimes(n int) []string {
out := make([]string, n)
for i := 0; i < n; i++ {
m := -(n - 1 - i)
if m == 0 {
out[i] = "now"
} else {
out[i] = fmt.Sprintf("%dm", m)
}
}
return out
}
func seriesFromSamples(samples []metricsSample) []float64 {
out := make([]float64, len(samples))
for i, s := range samples {
if v, ok := s.values[""]; ok {
out[i] = v
}
}
return out
}
func timeLabels(samples []metricsSample) []string {
out := make([]string, len(samples))
for i, s := range samples {
out[i] = s.t.Format("15:04")
}
return out
}
+1
View File
@@ -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)
+123
View File
@@ -0,0 +1,123 @@
package dash
import (
"math"
"sync"
"time"
)
// Suffixes for series derived from raw scrapes. Counters become per-second
// rates and histograms become latency quantiles, both computed from the delta
// against the previous scrape so the values reflect the last interval rather
// than process lifetime totals.
const (
suffixRate = ":rate"
suffixP50 = ":p50"
suffixP95 = ":p95"
suffixP99 = ":p99"
)
type prevScrape struct {
t time.Time
value float64
buckets []histogramBucket
}
type metricsDeriver struct {
mu sync.Mutex
prev map[string]prevScrape
}
func newMetricsDeriver() *metricsDeriver {
return &metricsDeriver{prev: make(map[string]prevScrape)}
}
// record stores the raw value and, for counters and histograms, the derived
// rate/quantile series for this interval.
func (d *metricsDeriver) record(store *metricsStore, source string, m scrapedMetric, now time.Time) {
key := source + "/" + m.name + "/" + labelKey(m.labels)
d.mu.Lock()
prev, hadPrev := d.prev[key]
d.prev[key] = prevScrape{t: now, value: m.value, buckets: m.buckets}
d.mu.Unlock()
switch m.kind {
case kindGauge:
store.recordLabeled(source, m.name, m.labels, m.value, now)
case kindCounter:
if !hadPrev {
return
}
dt := now.Sub(prev.t).Seconds()
if dt <= 0 {
return
}
delta := m.value - prev.value
if delta < 0 {
// Counter reset (process restart); skip this interval.
return
}
store.recordLabeled(source, m.name+suffixRate, m.labels, delta/dt, now)
case kindHistogram:
if !hadPrev {
return
}
delta := bucketDelta(prev.buckets, m.buckets)
if len(delta) == 0 {
return
}
for suffix, q := range map[string]float64{suffixP50: 0.5, suffixP95: 0.95, suffixP99: 0.99} {
store.recordLabeled(source, m.name+suffix, m.labels, histogramQuantile(delta, q), now)
}
}
}
// bucketDelta subtracts cumulative bucket counts, yielding the distribution
// observed during the interval. Returns nil on a reset or bucket mismatch.
func bucketDelta(prev, cur []histogramBucket) []histogramBucket {
if len(prev) != len(cur) {
return nil
}
out := make([]histogramBucket, len(cur))
for i := range cur {
if cur[i].upperBound != prev[i].upperBound {
return nil
}
c := cur[i].count - prev[i].count
if c < 0 {
return nil
}
out[i] = histogramBucket{upperBound: cur[i].upperBound, count: c}
}
if out[len(out)-1].count == 0 {
return nil
}
return out
}
// histogramQuantile estimates a quantile from cumulative buckets by linear
// interpolation within the matching bucket, matching Prometheus' approach.
func histogramQuantile(buckets []histogramBucket, q float64) float64 {
total := buckets[len(buckets)-1].count
if total == 0 {
return 0
}
rank := q * total
prevCount, prevBound := 0.0, 0.0
for _, b := range buckets {
if b.count < rank {
prevCount, prevBound = b.count, b.upperBound
continue
}
if math.IsInf(b.upperBound, 1) {
return prevBound
}
span := b.count - prevCount
if span <= 0 {
return b.upperBound
}
return prevBound + (b.upperBound-prevBound)*(rank-prevCount)/span
}
return buckets[len(buckets)-1].upperBound
}
+124
View File
@@ -0,0 +1,124 @@
package dash
import (
"math"
"testing"
"time"
)
func TestHistogramQuantile(t *testing.T) {
// 100 observations spread evenly across 0-1s in 10 buckets.
buckets := []histogramBucket{
{0.1, 10}, {0.2, 20}, {0.3, 30}, {0.4, 40}, {0.5, 50},
{0.6, 60}, {0.7, 70}, {0.8, 80}, {0.9, 90}, {1.0, 100},
{math.Inf(1), 100},
}
for _, tc := range []struct {
q float64
want float64
}{
{0.5, 0.5},
{0.95, 0.95},
{0.99, 0.99},
} {
got := histogramQuantile(buckets, tc.q)
if math.Abs(got-tc.want) > 1e-9 {
t.Errorf("quantile(%v) = %v, want %v", tc.q, got, tc.want)
}
}
}
func TestHistogramQuantileEmpty(t *testing.T) {
if got := histogramQuantile([]histogramBucket{{math.Inf(1), 0}}, 0.99); got != 0 {
t.Errorf("empty histogram quantile = %v, want 0", got)
}
}
func TestBucketDeltaResetAndMismatch(t *testing.T) {
prev := []histogramBucket{{0.1, 5}, {math.Inf(1), 10}}
if got := bucketDelta(prev, []histogramBucket{{0.1, 1}, {math.Inf(1), 2}}); got != nil {
t.Errorf("counter reset should yield nil, got %v", got)
}
if got := bucketDelta(prev, []histogramBucket{{0.2, 5}, {math.Inf(1), 10}}); got != nil {
t.Errorf("bound mismatch should yield nil, got %v", got)
}
if got := bucketDelta(prev, prev); got != nil {
t.Errorf("no new observations should yield nil, got %v", got)
}
got := bucketDelta(prev, []histogramBucket{{0.1, 7}, {math.Inf(1), 14}})
if len(got) != 2 || got[0].count != 2 || got[1].count != 4 {
t.Errorf("unexpected delta %v", got)
}
}
func TestDeriveCounterRate(t *testing.T) {
store := newMetricsStore()
d := newMetricsDeriver()
t0 := time.Now()
m := scrapedMetric{name: "reqs", kind: kindCounter, value: 100}
d.record(store, "volume/a", m, t0)
if got := store.match("volume", "reqs"+suffixRate); len(got) != 0 {
t.Fatalf("first scrape should not emit a rate, got %d series", len(got))
}
m.value = 250
d.record(store, "volume/a", m, t0.Add(15*time.Second))
series := store.match("volume", "reqs"+suffixRate)
if len(series) != 1 {
t.Fatalf("expected 1 rate series, got %d", len(series))
}
samples := series[0].snapshot()
if len(samples) != 1 {
t.Fatalf("expected 1 sample, got %d", len(samples))
}
if want := 10.0; math.Abs(samples[0].values[""]-want) > 1e-9 {
t.Errorf("rate = %v, want %v", samples[0].values[""], want)
}
}
func TestDeriveCounterResetSkipped(t *testing.T) {
store := newMetricsStore()
d := newMetricsDeriver()
t0 := time.Now()
m := scrapedMetric{name: "reqs", kind: kindCounter, value: 100}
d.record(store, "volume/a", m, t0)
m.value = 5
d.record(store, "volume/a", m, t0.Add(15*time.Second))
if got := store.match("volume", "reqs"+suffixRate); len(got) != 0 {
t.Errorf("counter reset should emit no rate, got %d series", len(got))
}
}
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++ {
s.record(time.Now(), map[string]float64{"": float64(i)})
}
got := s.snapshot()
if len(got) != metricsMaxSamples {
t.Fatalf("len = %d, want %d", len(got), metricsMaxSamples)
}
if got[len(got)-1].values[""] != float64(metricsMaxSamples+49) {
t.Errorf("newest sample not retained: %v", got[len(got)-1].values[""])
}
}
+267
View File
@@ -0,0 +1,267 @@
package dash
import (
"context"
"fmt"
"io"
"net"
"net/http"
"sort"
"strconv"
"strings"
"time"
dto "github.com/prometheus/client_model/go"
"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"
)
type metricKind int
const (
kindGauge metricKind = iota
kindCounter
kindHistogram
)
type histogramBucket struct {
upperBound float64
count float64
}
type scrapedMetric struct {
name string
labels map[string]string
kind metricKind
value float64
buckets []histogramBucket
}
func scrapeMetrics(ctx context.Context, target string) ([]scrapedMetric, error) {
if !strings.HasPrefix(target, "http://") && !strings.HasPrefix(target, "https://") {
target = "http://" + target
}
target = strings.TrimRight(target, "/") + "/metrics"
req, err := http.NewRequestWithContext(ctx, http.MethodGet, target, nil)
if err != nil {
return nil, err
}
req.Header.Set("Accept", string(expfmt.TextVersion))
client := util_http.GetGlobalHttpClient()
resp, err := client.Do(req)
if err != nil {
return nil, err
}
defer resp.Body.Close()
if resp.StatusCode != http.StatusOK {
return nil, fmt.Errorf("scrape %s: status %d", target, resp.StatusCode)
}
return parsePrometheusText(resp.Body)
}
func parsePrometheusText(r io.Reader) ([]scrapedMetric, error) {
dec := expfmt.NewDecoder(r, expfmt.NewFormat(expfmt.TypeTextPlain))
var out []scrapedMetric
for {
var fam dto.MetricFamily
if err := dec.Decode(&fam); err != nil {
if err == io.EOF {
break
}
return nil, err
}
for _, m := range fam.Metric {
out = append(out, toScrapedMetric(fam.GetName(), m))
}
}
return out, nil
}
func toScrapedMetric(name string, m *dto.Metric) scrapedMetric {
labels := map[string]string{}
for _, l := range m.Label {
labels[l.GetName()] = l.GetValue()
}
sm := scrapedMetric{name: name, labels: labels}
switch {
case m.Counter != nil:
sm.kind, sm.value = kindCounter, m.Counter.GetValue()
case m.Histogram != nil:
sm.kind = kindHistogram
for _, b := range m.Histogram.Bucket {
sm.buckets = append(sm.buckets, histogramBucket{upperBound: b.GetUpperBound(), count: float64(b.GetCumulativeCount())})
}
case m.Summary != nil:
sm.kind, sm.value = kindCounter, m.Summary.GetSampleSum()
case m.Gauge != nil:
sm.value = m.Gauge.GetValue()
case m.Untyped != nil:
sm.value = m.Untyped.GetValue()
}
return sm
}
// gatherLocalMetrics records the admin's own registry (maintenance tasks,
// worker slots) without a network round trip.
func (s *AdminServer) gatherLocalMetrics(now time.Time) {
families, err := stats_collect.Gather.Gather()
if err != nil {
glog.V(1).Infof("gather admin metrics: %v", err)
return
}
for _, fam := range families {
for _, m := range fam.Metric {
s.metricsDeriver.record(s.metricsStore, "admin/local", toScrapedMetric(fam.GetName(), m), now)
}
}
}
func (s *AdminServer) scrapeAllServers(ctx context.Context) {
now := time.Now()
s.gatherLocalMetrics(now)
targets := s.scrapeTargets()
if len(targets) == 0 {
return
}
type result struct {
source string
metrics []scrapedMetric
err error
}
results := make(chan result, len(targets))
scrapeCtx, cancel := context.WithTimeout(ctx, 5*time.Second)
defer cancel()
for _, t := range targets {
go func(t scrapeTarget) {
ms, err := scrapeMetrics(scrapeCtx, t.address)
results <- result{source: t.source, metrics: ms, err: err}
}(t)
}
for i := 0; i < len(targets); i++ {
r := <-results
if r.err != nil {
glog.V(1).Infof("metrics scrape %s: %v", r.source, r.err)
continue
}
for _, m := range r.metrics {
s.metricsDeriver.record(s.metricsStore, r.source, m, now)
}
}
}
// 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 {
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 {
add(vs.Address, vs.MetricsPort)
}
}
for _, f := range s.getFilerNodesStatus() {
add(f.Address, f.MetricsPort)
}
for _, n := range s.getS3NodesStatus() {
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
}
func (s *AdminServer) startMetricsScraper(ctx context.Context) {
const interval = 15 * time.Second
ticker := time.NewTicker(interval)
defer ticker.Stop()
s.scrapeAllServers(ctx)
for {
select {
case <-ctx.Done():
return
case <-ticker.C:
s.scrapeAllServers(ctx)
}
}
}
+148
View File
@@ -0,0 +1,148 @@
package dash
import (
"strings"
"sync"
"time"
)
const metricsMaxSamples = 240
type metricsSample struct {
t time.Time
values map[string]float64
}
type metricsSeries struct {
mu sync.Mutex
samples []metricsSample
}
func newMetricsSeries() *metricsSeries {
return &metricsSeries{samples: make([]metricsSample, 0, metricsMaxSamples)}
}
func (s *metricsSeries) record(t time.Time, values map[string]float64) {
s.mu.Lock()
defer s.mu.Unlock()
s.samples = append(s.samples, metricsSample{t: t, values: values})
if len(s.samples) > metricsMaxSamples {
s.samples = s.samples[len(s.samples)-metricsMaxSamples:]
}
}
func (s *metricsSeries) snapshot() []metricsSample {
s.mu.Lock()
defer s.mu.Unlock()
out := make([]metricsSample, len(s.samples))
copy(out, s.samples)
return out
}
type metricsStore struct {
mu sync.Mutex
series map[string]*metricsSeries
}
func newMetricsStore() *metricsStore {
return &metricsStore{series: make(map[string]*metricsSeries)}
}
func (s *metricsStore) record(source, name string, value float64, t time.Time) {
s.mu.Lock()
key := source + "/" + name
ser, ok := s.series[key]
if !ok {
ser = newMetricsSeries()
s.series[key] = ser
}
s.mu.Unlock()
ser.record(t, map[string]float64{"": value})
}
func (s *metricsStore) recordLabeled(source, name string, labels map[string]string, value float64, t time.Time) {
s.mu.Lock()
key := source + "/" + name + "/" + labelKey(labels)
ser, ok := s.series[key]
if !ok {
ser = newMetricsSeries()
s.series[key] = ser
}
s.mu.Unlock()
ser.record(t, map[string]float64{"": value})
}
func (s *metricsStore) get(source, name string) []metricsSample {
s.mu.Lock()
key := source + "/" + name
ser, ok := s.series[key]
s.mu.Unlock()
if !ok {
return nil
}
return ser.snapshot()
}
func (s *metricsStore) getLabeled(source, name string, labels map[string]string) []metricsSample {
s.mu.Lock()
key := source + "/" + name + "/" + labelKey(labels)
ser, ok := s.series[key]
s.mu.Unlock()
if !ok {
return nil
}
return ser.snapshot()
}
// match returns every series whose source has the given prefix and whose
// metric name matches exactly.
func (s *metricsStore) match(sourcePrefix, name string) []*metricsSeries {
s.mu.Lock()
defer s.mu.Unlock()
var out []*metricsSeries
for k, ser := range s.series {
if !strings.HasPrefix(k, sourcePrefix) {
continue
}
rest := k[len(sourcePrefix):]
if !strings.HasPrefix(rest, "/") {
continue
}
rest = rest[1:]
// rest is either "<addr>/<metric>[/<labels>]" or "<metric>[/<labels>]".
if rest == name || strings.HasPrefix(rest, name+"/") {
out = append(out, ser)
continue
}
if i := strings.Index(rest, "/"); i >= 0 {
tail := rest[i+1:]
if tail == name || strings.HasPrefix(tail, name+"/") {
out = append(out, ser)
}
}
}
return out
}
func labelKey(labels map[string]string) string {
if len(labels) == 0 {
return ""
}
keys := make([]string, 0, len(labels))
for k := range labels {
keys = append(keys, k)
}
for i := 1; i < len(keys); i++ {
for j := i; j > 0 && keys[j] < keys[j-1]; j-- {
keys[j], keys[j-1] = keys[j-1], keys[j]
}
}
out := ""
for i, k := range keys {
if i > 0 {
out += ","
}
out += k + "=" + labels[k]
}
return out
}
+2
View File
@@ -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
+9 -6
View File
@@ -27,6 +27,9 @@ type ClusterNode struct {
CreatedTs time.Time
DataCenter DataCenter
Rack Rack
// MetricsPort is the node's Prometheus /metrics port, or 0 when the node
// does not run a metrics listener.
MetricsPort uint32
}
type ClusterNodeGroups struct {
@@ -53,11 +56,11 @@ func (g *ClusterNodeGroups) getGroupMembers(filerGroup FilerGroupName, createIfN
return members
}
func (g *ClusterNodeGroups) AddClusterNode(filerGroup FilerGroupName, nodeType string, dataCenter DataCenter, rack Rack, address pb.ServerAddress, version string) []*master_pb.KeepConnectedResponse {
func (g *ClusterNodeGroups) AddClusterNode(filerGroup FilerGroupName, nodeType string, dataCenter DataCenter, rack Rack, address pb.ServerAddress, version string, metricsPort uint32) []*master_pb.KeepConnectedResponse {
g.Lock()
defer g.Unlock()
m := g.getGroupMembers(filerGroup, true)
if t := m.addMember(dataCenter, rack, address, version); t != nil {
if t := m.addMember(dataCenter, rack, address, version, metricsPort); t != nil {
return buildClusterNodeUpdateMessage(true, filerGroup, nodeType, address)
}
return nil
@@ -95,15 +98,15 @@ func NewCluster() *Cluster {
}
}
func (cluster *Cluster) AddClusterNode(ns, nodeType string, dataCenter DataCenter, rack Rack, address pb.ServerAddress, version string) []*master_pb.KeepConnectedResponse {
func (cluster *Cluster) AddClusterNode(ns, nodeType string, dataCenter DataCenter, rack Rack, address pb.ServerAddress, version string, metricsPort uint32) []*master_pb.KeepConnectedResponse {
filerGroup := FilerGroupName(ns)
switch nodeType {
case FilerType:
return cluster.filerGroups.AddClusterNode(filerGroup, nodeType, dataCenter, rack, address, version)
return cluster.filerGroups.AddClusterNode(filerGroup, nodeType, dataCenter, rack, address, version, metricsPort)
case BrokerType:
return cluster.brokerGroups.AddClusterNode(filerGroup, nodeType, dataCenter, rack, address, version)
return cluster.brokerGroups.AddClusterNode(filerGroup, nodeType, dataCenter, rack, address, version, metricsPort)
case S3Type:
return cluster.s3Groups.AddClusterNode(filerGroup, nodeType, dataCenter, rack, address, version)
return cluster.s3Groups.AddClusterNode(filerGroup, nodeType, dataCenter, rack, address, version, metricsPort)
case MasterType:
return buildClusterNodeUpdateMessage(true, filerGroup, nodeType, address)
}
+4 -4
View File
@@ -16,7 +16,7 @@ func TestConcurrentAddRemoveNodes(t *testing.T) {
go func(i int) {
defer wg.Done()
address := strconv.Itoa(i)
c.AddClusterNode("", "filer", "", "", pb.ServerAddress(address), "23.45")
c.AddClusterNode("", "filer", "", "", pb.ServerAddress(address), "23.45", 0)
}(i)
}
wg.Wait()
@@ -43,8 +43,8 @@ func TestConcurrentAddRemoveNodes(t *testing.T) {
func TestListClusterNodeUpdates(t *testing.T) {
c := NewCluster()
filer := pb.ServerAddress("10.0.0.20:8888")
c.AddClusterNode("group", FilerType, "dc1", "rack1", filer, "test")
c.AddClusterNode("group", BrokerType, "dc1", "rack1", pb.ServerAddress("10.0.0.20:17777"), "test")
c.AddClusterNode("group", FilerType, "dc1", "rack1", filer, "test", 0)
c.AddClusterNode("group", BrokerType, "dc1", "rack1", pb.ServerAddress("10.0.0.20:17777"), "test", 0)
updates := c.ListClusterNodeUpdates("group", FilerType)
if len(updates) != 1 {
@@ -64,7 +64,7 @@ func TestListClusterNodeUpdates(t *testing.T) {
func TestIsKnownNode(t *testing.T) {
c := NewCluster()
filer := pb.ServerAddress("10.0.0.20:8888")
c.AddClusterNode("", FilerType, "dc1", "rack1", filer, "test")
c.AddClusterNode("", FilerType, "dc1", "rack1", filer, "test", 0)
if !c.IsKnownNode(FilerType, filer) {
t.Fatalf("registered filer %s should be known", filer)
+10 -7
View File
@@ -16,18 +16,21 @@ func newGroupMembers() *GroupMembers {
}
}
func (m *GroupMembers) addMember(dataCenter DataCenter, rack Rack, address pb.ServerAddress, version string) *ClusterNode {
func (m *GroupMembers) addMember(dataCenter DataCenter, rack Rack, address pb.ServerAddress, version string, metricsPort uint32) *ClusterNode {
if existingNode, found := m.members[address]; found {
existingNode.counter++
// A restarted node may have gained or lost its metrics listener.
existingNode.MetricsPort = metricsPort
return nil
}
t := &ClusterNode{
Address: address,
Version: version,
counter: 1,
CreatedTs: time.Now(),
DataCenter: dataCenter,
Rack: rack,
Address: address,
Version: version,
counter: 1,
CreatedTs: time.Now(),
DataCenter: dataCenter,
Rack: rack,
MetricsPort: metricsPort,
}
m.members[address] = t
return t
+1
View File
@@ -380,6 +380,7 @@ func (fo *FilerOptions) startFiler() {
fs, nfs_err := weed_server.NewFilerServer(defaultMux, publicVolumeMux, &weed_server.FilerOption{
Masters: fo.masters,
MetricsPort: uint32(*fo.metricsHttpPort),
FilerGroup: *fo.filerGroup,
Collection: *fo.collection,
DefaultReplication: *fo.defaultReplicaPlacement,
+7
View File
@@ -465,6 +465,12 @@ func peerIndex(self pb.ServerAddress, peers []pb.ServerAddress) int {
}
func (m *MasterOptions) toMasterOption(whiteList []string) *weed_server.MasterOption {
// Not every caller wires up every flag, so treat an unset metrics port as
// disabled rather than dereferencing nil.
metricsPort := 0
if m.metricsHttpPort != nil {
metricsPort = *m.metricsHttpPort
}
masterAddress := pb.NewServerAddress(*m.ip, *m.port, *m.portGrpc)
return &weed_server.MasterOption{
Master: masterAddress,
@@ -479,6 +485,7 @@ func (m *MasterOptions) toMasterOption(whiteList []string) *weed_server.MasterOp
WhiteList: whiteList,
DisableHttp: *m.disableHttp,
MetricsAddress: *m.metricsAddress,
MetricsPort: metricsPort,
MetricsIntervalSec: *m.metricsIntervalSec,
TelemetryUrl: *m.telemetryUrl,
TelemetryEnabled: *m.telemetryEnabled,
+7
View File
@@ -1308,6 +1308,13 @@ func runMini(cmd *Command, args []string) bool {
}
pb.RegisterLocalGrpcSocket(*miniIp, *miniAdminOptions.grpcPort, fmt.Sprintf("/tmp/seaweedfs-admin-grpc-%d.sock", *miniAdminOptions.grpcPort))
// One process, one shared Prometheus registry, so this single listener
// serves every component's series. Point each component at it so they all
// advertise the same port to the master.
miniMasterOptions.metricsHttpPort = miniMetricsHttpPort
miniOptions.v.metricsHttpPort = miniMetricsHttpPort
miniFilerOptions.metricsHttpPort = miniMetricsHttpPort
go stats_collect.StartMetricsServer(*miniMetricsHttpIp, *miniMetricsHttpPort)
if *miniMasterOptions.volumeSizeLimitMB > util.MaxVolumeSizeLimitMB {
+1
View File
@@ -358,6 +358,7 @@ func (s3opt *S3Options) startS3Server() bool {
Filers: filerAddresses,
Masters: masterAddresses,
Port: *s3opt.port,
MetricsPort: uint32(*s3opt.metricsHttpPort),
Config: *s3opt.config,
DomainName: *s3opt.domainName,
AllowedOrigins: strings.Split(*s3opt.allowedOrigins, ","),
+9
View File
@@ -343,6 +343,15 @@ func runServer(cmd *Command, args []string) bool {
webdavOptions.filer = &filerAddress
mqBrokerOptions.filerGroup = filerOptions.filerGroup
// One process, one shared Prometheus registry, so a single metrics listener
// serves master, volume, filer and S3 series. Every component advertises
// that same port; the admin server deduplicates targets by address and
// attributes series by metric name.
masterOptions.metricsHttpPort = serverMetricsHttpPort
serverOptions.v.metricsHttpPort = serverMetricsHttpPort
filerOptions.metricsHttpPort = serverMetricsHttpPort
s3Options.metricsHttpPort = serverMetricsHttpPort
go stats_collect.StartMetricsServer(*serverMetricsHttpIp, *serverMetricsHttpPort)
*volumeDataFolders = util.ResolveCommaSeparatedPaths(*volumeDataFolders)
+1 -1
View File
@@ -438,7 +438,7 @@ func (v VolumeServerOptions) startVolumeServer(volumeFolders, maxVolumeCounts, v
RecoveryCoef: *v.diskRecoveryCoef,
}
volumeServer := weed_server.NewVolumeServer(volumeMux, publicVolumeMux,
*v.ip, *v.port, *v.portGrpc, *v.publicUrl, volumeServerId,
*v.ip, *v.port, *v.portGrpc, *v.metricsHttpPort, *v.publicUrl, volumeServerId,
v.folders, v.folderMaxLimits, minFreeSpaces, diskTypes, folderTags,
util.ResolvePath(*v.idxFolder),
volumeNeedleMapKind,
+14
View File
@@ -98,6 +98,9 @@ message Heartbeat {
map<string, uint32> max_volume_counts = 4;
uint32 grpc_port = 20;
// Port of this volume server's Prometheus /metrics listener (-metricsPort),
// or 0 when it is not enabled.
uint32 metrics_port = 29;
repeated string location_uuids = 21;
string id = 22; // volume server id, independent of ip:port for stable identification
@@ -213,6 +216,10 @@ message KeepConnectedRequest {
string filer_group = 5;
string data_center = 6;
string rack = 7;
// Port of this node's Prometheus /metrics listener (-metricsPort), or 0 when
// it is not enabled. Advertised so the admin server can scrape it without
// exposing metrics on the client-facing service port.
uint32 metrics_port = 8;
}
message VolumeLocation {
@@ -400,6 +407,7 @@ message DataNodeInfo {
map<string, DiskInfo> diskInfos = 2;
uint32 grpc_port = 3;
string address = 4; // ip:port for connecting to the volume server
uint32 metrics_port = 5; // Prometheus /metrics port, or 0 when not enabled
}
message RackInfo {
string id = 1;
@@ -514,6 +522,10 @@ message GetMasterConfigurationResponse {
// MIGRATION: fields 8-9 help migrate master.toml [master.maintenance] to admin script plugin. Remove after March 2027.
string maintenance_scripts = 8;
uint32 maintenance_sleep_minutes = 9;
// Port of this master's Prometheus /metrics listener (-metricsPort), or 0
// when it is not enabled. metrics_address above is unrelated: it is the
// Prometheus push gateway that servers push to.
uint32 metrics_port = 10;
}
message ListClusterNodesRequest {
@@ -528,6 +540,8 @@ message ListClusterNodesResponse {
int64 created_at_ns = 4;
string data_center = 5;
string rack = 6;
// Port of this node's Prometheus /metrics listener, or 0 when not enabled.
uint32 metrics_port = 7;
}
repeated ClusterNode cluster_nodes = 1;
}
+77 -22
View File
@@ -113,8 +113,11 @@ type Heartbeat struct {
HasNoEcShards bool `protobuf:"varint,19,opt,name=has_no_ec_shards,json=hasNoEcShards,proto3" json:"has_no_ec_shards,omitempty"`
MaxVolumeCounts map[string]uint32 `protobuf:"bytes,4,rep,name=max_volume_counts,json=maxVolumeCounts,proto3" json:"max_volume_counts,omitempty" protobuf_key:"bytes,1,opt,name=key" protobuf_val:"varint,2,opt,name=value"`
GrpcPort uint32 `protobuf:"varint,20,opt,name=grpc_port,json=grpcPort,proto3" json:"grpc_port,omitempty"`
LocationUuids []string `protobuf:"bytes,21,rep,name=location_uuids,json=locationUuids,proto3" json:"location_uuids,omitempty"`
Id string `protobuf:"bytes,22,opt,name=id,proto3" json:"id,omitempty"` // volume server id, independent of ip:port for stable identification
// Port of this volume server's Prometheus /metrics listener (-metricsPort),
// or 0 when it is not enabled.
MetricsPort uint32 `protobuf:"varint,29,opt,name=metrics_port,json=metricsPort,proto3" json:"metrics_port,omitempty"`
LocationUuids []string `protobuf:"bytes,21,rep,name=location_uuids,json=locationUuids,proto3" json:"location_uuids,omitempty"`
Id string `protobuf:"bytes,22,opt,name=id,proto3" json:"id,omitempty"` // volume server id, independent of ip:port for stable identification
// state flags
State *volume_server_pb.VolumeServerState `protobuf:"bytes,23,opt,name=state,proto3" json:"state,omitempty"`
DiskTags []*DiskTag `protobuf:"bytes,24,rep,name=disk_tags,json=diskTags,proto3" json:"disk_tags,omitempty"`
@@ -283,6 +286,13 @@ func (x *Heartbeat) GetGrpcPort() uint32 {
return 0
}
func (x *Heartbeat) GetMetricsPort() uint32 {
if x != nil {
return x.MetricsPort
}
return 0
}
func (x *Heartbeat) GetLocationUuids() []string {
if x != nil {
return x.LocationUuids
@@ -998,6 +1008,10 @@ type KeepConnectedRequest struct {
FilerGroup string `protobuf:"bytes,5,opt,name=filer_group,json=filerGroup,proto3" json:"filer_group,omitempty"`
DataCenter string `protobuf:"bytes,6,opt,name=data_center,json=dataCenter,proto3" json:"data_center,omitempty"`
Rack string `protobuf:"bytes,7,opt,name=rack,proto3" json:"rack,omitempty"`
// Port of this node's Prometheus /metrics listener (-metricsPort), or 0 when
// it is not enabled. Advertised so the admin server can scrape it without
// exposing metrics on the client-facing service port.
MetricsPort uint32 `protobuf:"varint,8,opt,name=metrics_port,json=metricsPort,proto3" json:"metrics_port,omitempty"`
unknownFields protoimpl.UnknownFields
sizeCache protoimpl.SizeCache
}
@@ -1074,6 +1088,13 @@ func (x *KeepConnectedRequest) GetRack() string {
return ""
}
func (x *KeepConnectedRequest) GetMetricsPort() uint32 {
if x != nil {
return x.MetricsPort
}
return 0
}
type VolumeLocation struct {
state protoimpl.MessageState `protogen:"open.v1"`
Url string `protobuf:"bytes,1,opt,name=url,proto3" json:"url,omitempty"`
@@ -2608,7 +2629,8 @@ type DataNodeInfo struct {
Id string `protobuf:"bytes,1,opt,name=id,proto3" json:"id,omitempty"`
DiskInfos map[string]*DiskInfo `protobuf:"bytes,2,rep,name=diskInfos,proto3" json:"diskInfos,omitempty" protobuf_key:"bytes,1,opt,name=key" protobuf_val:"bytes,2,opt,name=value"`
GrpcPort uint32 `protobuf:"varint,3,opt,name=grpc_port,json=grpcPort,proto3" json:"grpc_port,omitempty"`
Address string `protobuf:"bytes,4,opt,name=address,proto3" json:"address,omitempty"` // ip:port for connecting to the volume server
Address string `protobuf:"bytes,4,opt,name=address,proto3" json:"address,omitempty"` // ip:port for connecting to the volume server
MetricsPort uint32 `protobuf:"varint,5,opt,name=metrics_port,json=metricsPort,proto3" json:"metrics_port,omitempty"` // Prometheus /metrics port, or 0 when not enabled
unknownFields protoimpl.UnknownFields
sizeCache protoimpl.SizeCache
}
@@ -2671,6 +2693,13 @@ func (x *DataNodeInfo) GetAddress() string {
return ""
}
func (x *DataNodeInfo) GetMetricsPort() uint32 {
if x != nil {
return x.MetricsPort
}
return 0
}
type RackInfo struct {
state protoimpl.MessageState `protogen:"open.v1"`
Id string `protobuf:"bytes,1,opt,name=id,proto3" json:"id,omitempty"`
@@ -3639,8 +3668,12 @@ type GetMasterConfigurationResponse struct {
// MIGRATION: fields 8-9 help migrate master.toml [master.maintenance] to admin script plugin. Remove after March 2027.
MaintenanceScripts string `protobuf:"bytes,8,opt,name=maintenance_scripts,json=maintenanceScripts,proto3" json:"maintenance_scripts,omitempty"`
MaintenanceSleepMinutes uint32 `protobuf:"varint,9,opt,name=maintenance_sleep_minutes,json=maintenanceSleepMinutes,proto3" json:"maintenance_sleep_minutes,omitempty"`
unknownFields protoimpl.UnknownFields
sizeCache protoimpl.SizeCache
// Port of this master's Prometheus /metrics listener (-metricsPort), or 0
// when it is not enabled. metrics_address above is unrelated: it is the
// Prometheus push gateway that servers push to.
MetricsPort uint32 `protobuf:"varint,10,opt,name=metrics_port,json=metricsPort,proto3" json:"metrics_port,omitempty"`
unknownFields protoimpl.UnknownFields
sizeCache protoimpl.SizeCache
}
func (x *GetMasterConfigurationResponse) Reset() {
@@ -3736,6 +3769,13 @@ func (x *GetMasterConfigurationResponse) GetMaintenanceSleepMinutes() uint32 {
return 0
}
func (x *GetMasterConfigurationResponse) GetMetricsPort() uint32 {
if x != nil {
return x.MetricsPort
}
return 0
}
type ListClusterNodesRequest struct {
state protoimpl.MessageState `protogen:"open.v1"`
ClientType string `protobuf:"bytes,1,opt,name=client_type,json=clientType,proto3" json:"client_type,omitempty"`
@@ -4865,12 +4905,14 @@ func (x *LookupEcVolumeResponse_EcShardIdLocation) GetLocations() []*Location {
}
type ListClusterNodesResponse_ClusterNode struct {
state protoimpl.MessageState `protogen:"open.v1"`
Address string `protobuf:"bytes,1,opt,name=address,proto3" json:"address,omitempty"`
Version string `protobuf:"bytes,2,opt,name=version,proto3" json:"version,omitempty"`
CreatedAtNs int64 `protobuf:"varint,4,opt,name=created_at_ns,json=createdAtNs,proto3" json:"created_at_ns,omitempty"`
DataCenter string `protobuf:"bytes,5,opt,name=data_center,json=dataCenter,proto3" json:"data_center,omitempty"`
Rack string `protobuf:"bytes,6,opt,name=rack,proto3" json:"rack,omitempty"`
state protoimpl.MessageState `protogen:"open.v1"`
Address string `protobuf:"bytes,1,opt,name=address,proto3" json:"address,omitempty"`
Version string `protobuf:"bytes,2,opt,name=version,proto3" json:"version,omitempty"`
CreatedAtNs int64 `protobuf:"varint,4,opt,name=created_at_ns,json=createdAtNs,proto3" json:"created_at_ns,omitempty"`
DataCenter string `protobuf:"bytes,5,opt,name=data_center,json=dataCenter,proto3" json:"data_center,omitempty"`
Rack string `protobuf:"bytes,6,opt,name=rack,proto3" json:"rack,omitempty"`
// Port of this node's Prometheus /metrics listener, or 0 when not enabled.
MetricsPort uint32 `protobuf:"varint,7,opt,name=metrics_port,json=metricsPort,proto3" json:"metrics_port,omitempty"`
unknownFields protoimpl.UnknownFields
sizeCache protoimpl.SizeCache
}
@@ -4940,6 +4982,13 @@ func (x *ListClusterNodesResponse_ClusterNode) GetRack() string {
return ""
}
func (x *ListClusterNodesResponse_ClusterNode) GetMetricsPort() uint32 {
if x != nil {
return x.MetricsPort
}
return 0
}
type RaftListClusterServersResponse_ClusterServers struct {
state protoimpl.MessageState `protogen:"open.v1"`
Id string `protobuf:"bytes,1,opt,name=id,proto3" json:"id,omitempty"`
@@ -5017,7 +5066,7 @@ const file_master_proto_rawDesc = "" +
"\adisk_id\x18\x01 \x01(\rR\x06diskId\x12\x12\n" +
"\x04tags\x18\x02 \x03(\tR\x04tags\x12\x12\n" +
"\x04type\x18\x03 \x01(\tR\x04type\x12(\n" +
"\x10max_volume_count\x18\x04 \x01(\x03R\x0emaxVolumeCount\"\xf0\v\n" +
"\x10max_volume_count\x18\x04 \x01(\x03R\x0emaxVolumeCount\"\x93\f\n" +
"\tHeartbeat\x12\x0e\n" +
"\x02ip\x18\x01 \x01(\tR\x02ip\x12\x12\n" +
"\x04port\x18\x02 \x01(\rR\x04port\x12\x1d\n" +
@@ -5041,7 +5090,8 @@ const file_master_proto_rawDesc = "" +
"\x11deleted_ec_shards\x18\x12 \x03(\v2*.master_pb.VolumeEcShardInformationMessageR\x0fdeletedEcShards\x12'\n" +
"\x10has_no_ec_shards\x18\x13 \x01(\bR\rhasNoEcShards\x12U\n" +
"\x11max_volume_counts\x18\x04 \x03(\v2).master_pb.Heartbeat.MaxVolumeCountsEntryR\x0fmaxVolumeCounts\x12\x1b\n" +
"\tgrpc_port\x18\x14 \x01(\rR\bgrpcPort\x12%\n" +
"\tgrpc_port\x18\x14 \x01(\rR\bgrpcPort\x12!\n" +
"\fmetrics_port\x18\x1d \x01(\rR\vmetricsPort\x12%\n" +
"\x0elocation_uuids\x18\x15 \x03(\tR\rlocationUuids\x12\x0e\n" +
"\x02id\x18\x16 \x01(\tR\x02id\x129\n" +
"\x05state\x18\x17 \x01(\v2#.volume_server_pb.VolumeServerStateR\x05state\x12/\n" +
@@ -5137,7 +5187,7 @@ const file_master_proto_rawDesc = "" +
"\x04data\x18\x01 \x01(\rR\x04data\x12\x16\n" +
"\x06parity\x18\x02 \x01(\rR\x06parity\x12\x1d\n" +
"\n" +
"volume_ids\x18\x03 \x03(\rR\tvolumeIds\"\xce\x01\n" +
"volume_ids\x18\x03 \x03(\rR\tvolumeIds\"\xf1\x01\n" +
"\x14KeepConnectedRequest\x12\x1f\n" +
"\vclient_type\x18\x01 \x01(\tR\n" +
"clientType\x12%\n" +
@@ -5147,7 +5197,8 @@ const file_master_proto_rawDesc = "" +
"filerGroup\x12\x1f\n" +
"\vdata_center\x18\x06 \x01(\tR\n" +
"dataCenter\x12\x12\n" +
"\x04rack\x18\a \x01(\tR\x04rack\"\x9e\x03\n" +
"\x04rack\x18\a \x01(\tR\x04rack\x12!\n" +
"\fmetrics_port\x18\b \x01(\rR\vmetricsPort\"\x9e\x03\n" +
"\x0eVolumeLocation\x12\x10\n" +
"\x03url\x18\x01 \x01(\tR\x03url\x12\x1d\n" +
"\n" +
@@ -5296,12 +5347,13 @@ const file_master_proto_rawDesc = "" +
"\x0fdisk_free_bytes\x18\r \x01(\x04R\rdiskFreeBytes\x1aG\n" +
"\x19MaxVolumeCountByDiskEntry\x12\x10\n" +
"\x03key\x18\x01 \x01(\rR\x03key\x12\x14\n" +
"\x05value\x18\x02 \x01(\x03R\x05value:\x028\x01\"\xee\x01\n" +
"\x05value\x18\x02 \x01(\x03R\x05value:\x028\x01\"\x91\x02\n" +
"\fDataNodeInfo\x12\x0e\n" +
"\x02id\x18\x01 \x01(\tR\x02id\x12D\n" +
"\tdiskInfos\x18\x02 \x03(\v2&.master_pb.DataNodeInfo.DiskInfosEntryR\tdiskInfos\x12\x1b\n" +
"\tgrpc_port\x18\x03 \x01(\rR\bgrpcPort\x12\x18\n" +
"\aaddress\x18\x04 \x01(\tR\aaddress\x1aQ\n" +
"\aaddress\x18\x04 \x01(\tR\aaddress\x12!\n" +
"\fmetrics_port\x18\x05 \x01(\rR\vmetricsPort\x1aQ\n" +
"\x0eDiskInfosEntry\x12\x10\n" +
"\x03key\x18\x01 \x01(\tR\x03key\x12)\n" +
"\x05value\x18\x02 \x01(\v2\x13.master_pb.DiskInfoR\x05value:\x028\x01\"\xf0\x01\n" +
@@ -5385,7 +5437,7 @@ const file_master_proto_rawDesc = "" +
" \x01(\bR\n" +
"isReadonly\"\x1c\n" +
"\x1aVolumeMarkReadonlyResponse\"\x1f\n" +
"\x1dGetMasterConfigurationRequest\"\xe0\x03\n" +
"\x1dGetMasterConfigurationRequest\"\x83\x04\n" +
"\x1eGetMasterConfigurationResponse\x12'\n" +
"\x0fmetrics_address\x18\x01 \x01(\tR\x0emetricsAddress\x128\n" +
"\x18metrics_interval_seconds\x18\x02 \x01(\rR\x16metricsIntervalSeconds\x12D\n" +
@@ -5395,22 +5447,25 @@ const file_master_proto_rawDesc = "" +
"\x15volume_size_limit_m_b\x18\x06 \x01(\rR\x11volumeSizeLimitMB\x12-\n" +
"\x12volume_preallocate\x18\a \x01(\bR\x11volumePreallocate\x12/\n" +
"\x13maintenance_scripts\x18\b \x01(\tR\x12maintenanceScripts\x12:\n" +
"\x19maintenance_sleep_minutes\x18\t \x01(\rR\x17maintenanceSleepMinutes\"q\n" +
"\x19maintenance_sleep_minutes\x18\t \x01(\rR\x17maintenanceSleepMinutes\x12!\n" +
"\fmetrics_port\x18\n" +
" \x01(\rR\vmetricsPort\"q\n" +
"\x17ListClusterNodesRequest\x12\x1f\n" +
"\vclient_type\x18\x01 \x01(\tR\n" +
"clientType\x12\x1f\n" +
"\vfiler_group\x18\x02 \x01(\tR\n" +
"filerGroup\x12\x14\n" +
"\x05limit\x18\x04 \x01(\x05R\x05limit\"\x8d\x02\n" +
"\x05limit\x18\x04 \x01(\x05R\x05limit\"\xb0\x02\n" +
"\x18ListClusterNodesResponse\x12T\n" +
"\rcluster_nodes\x18\x01 \x03(\v2/.master_pb.ListClusterNodesResponse.ClusterNodeR\fclusterNodes\x1a\x9a\x01\n" +
"\rcluster_nodes\x18\x01 \x03(\v2/.master_pb.ListClusterNodesResponse.ClusterNodeR\fclusterNodes\x1a\xbd\x01\n" +
"\vClusterNode\x12\x18\n" +
"\aaddress\x18\x01 \x01(\tR\aaddress\x12\x18\n" +
"\aversion\x18\x02 \x01(\tR\aversion\x12\"\n" +
"\rcreated_at_ns\x18\x04 \x01(\x03R\vcreatedAtNs\x12\x1f\n" +
"\vdata_center\x18\x05 \x01(\tR\n" +
"dataCenter\x12\x12\n" +
"\x04rack\x18\x06 \x01(\tR\x04rack\"\xc5\x01\n" +
"\x04rack\x18\x06 \x01(\tR\x04rack\x12!\n" +
"\fmetrics_port\x18\a \x01(\rR\vmetricsPort\"\xc5\x01\n" +
"\x16LeaseAdminTokenRequest\x12%\n" +
"\x0eprevious_token\x18\x01 \x01(\x03R\rpreviousToken\x12,\n" +
"\x12previous_lock_time\x18\x02 \x01(\x03R\x10previousLockTime\x12\x1b\n" +
+10 -4
View File
@@ -65,10 +65,15 @@ type S3ApiServerOption struct {
Ip string // address advertised to the cluster; empty falls back to BindIp
BindIp string
GrpcPort int
ExternalUrl string // external URL clients use, tried first during signature verification behind a reverse proxy
DefaultFileMode uint32 // default file permission mode for S3 uploads (e.g. 0660, 0644)
CacheSizeMB int64 // in-memory chunk cache capacity in MB for the shared ReaderCache; 0 disables
MaxMB int32 // filer's -maxMB, read from the filer configuration at startup
// MetricsPort is the Prometheus /metrics port this S3 server serves,
// advertised to the master so the admin server can scrape it. 0 when
// disabled. S3 metrics carry bucket labels, so they are deliberately not
// served on the client-facing S3 port.
MetricsPort uint32
ExternalUrl string // external URL clients use, tried first during signature verification behind a reverse proxy
DefaultFileMode uint32 // default file permission mode for S3 uploads (e.g. 0660, 0644)
CacheSizeMB int64 // in-memory chunk cache capacity in MB for the shared ReaderCache; 0 disables
MaxMB int32 // filer's -maxMB, read from the filer configuration at startup
// AllowUntrustedRemoteEndpoints lets a read of a remote-only object dial a
// mounted endpoint that resolves to a loopback / private / metadata host.
AllowUntrustedRemoteEndpoints bool
@@ -209,6 +214,7 @@ func NewS3ApiServerWithStore(router *mux.Router, option *S3ApiServerOption, expl
}
clientHost := option.advertisedHost()
masterClient = wdclient.NewMasterClient(option.GrpcDialOption, option.FilerGroup, cluster.S3Type, pb.ServerAddress(util.JoinHostPort(clientHost, option.GrpcPort)), option.DataCenter, "", *pb.NewServiceDiscoveryFromMap(masterMap))
masterClient.SetMetricsPort(option.MetricsPort)
// Build the object-write lock client and subscribe to the master's
// lock-ring updates BEFORE starting the master loop, so the initial
// LockRingUpdate sent on connect isn't dropped (the master only delivers
+5 -1
View File
@@ -88,7 +88,10 @@ type FilerOption struct {
TusMaxSize int64
TusSessionExpiry time.Duration
S3ConfigFile string // optional path to static S3 identity config file
CredentialManager *credential.CredentialManager
// MetricsPort is the Prometheus /metrics port this filer serves, advertised
// to the master so the admin server can scrape it. 0 when disabled.
MetricsPort uint32
CredentialManager *credential.CredentialManager
// AllowUntrustedRemoteEndpoints lets a read of a remote-only entry dial a
// mounted endpoint that resolves to a loopback / private / metadata host.
AllowUntrustedRemoteEndpoints bool
@@ -244,6 +247,7 @@ func NewFilerServer(defaultMux, readonlyMux *http.ServeMux, option *FilerOption)
fs.checkWithMaster()
go stats.LoopPushingMetric("filer", string(fs.option.Host), fs.metricsAddress, fs.metricsIntervalSec)
fs.filer.MasterClient.SetMetricsPort(option.MetricsPort)
go fs.filer.MasterClient.KeepConnectedToMaster(context.Background())
fs.option.recursiveDelete = v.GetBool("filer.options.recursive_delete")
+3 -1
View File
@@ -188,6 +188,7 @@ func (ms *MasterServer) SendHeartbeat(stream master_pb.Seaweed_SendHeartbeatServ
dc := ms.Topo.GetOrCreateDataCenter(dcName)
rack := dc.GetOrCreateRack(rackName)
dn = rack.GetOrCreateDataNode(heartbeat.Ip, int(heartbeat.Port), int(heartbeat.GrpcPort), heartbeat.PublicUrl, heartbeat.Id, heartbeat.MaxVolumeCounts)
dn.MetricsPort = int(heartbeat.MetricsPort)
glog.V(0).Infof("added volume server %d: %v (id=%s, ip=%v:%d) %v", dn.Counter, dn.Id(), heartbeat.Id, heartbeat.GetIp(), heartbeat.GetPort(), heartbeat.LocationUuids)
uuidlist, err := ms.RegisterUuids(heartbeat)
if err != nil {
@@ -418,7 +419,7 @@ func (ms *MasterServer) KeepConnected(stream master_pb.Seaweed_KeepConnectedServ
stopChan := make(chan bool, 1)
clientName, messageChan := ms.addClient(req.FilerGroup, req.ClientType, peerAddress)
for _, update := range ms.Cluster.AddClusterNode(req.FilerGroup, req.ClientType, cluster.DataCenter(req.DataCenter), cluster.Rack(req.Rack), peerAddress, req.Version) {
for _, update := range ms.Cluster.AddClusterNode(req.FilerGroup, req.ClientType, cluster.DataCenter(req.DataCenter), cluster.Rack(req.Rack), peerAddress, req.Version, req.MetricsPort) {
glog.V(1).Infof("Cluster: %s node %s added to group '%s'", req.ClientType, peerAddress, req.FilerGroup)
ms.broadcastToClients(update)
}
@@ -662,6 +663,7 @@ func (ms *MasterServer) GetMasterConfiguration(ctx context.Context, req *master_
resp := &master_pb.GetMasterConfigurationResponse{
MetricsAddress: ms.option.MetricsAddress,
MetricsIntervalSeconds: uint32(ms.option.MetricsIntervalSec),
MetricsPort: uint32(ms.option.MetricsPort),
StorageBackends: backend.ToPbStorageBackends(),
DefaultReplication: ms.option.DefaultReplicaPlacement,
VolumeSizeLimitMB: uint32(ms.option.VolumeSizeLimitMB),
@@ -27,7 +27,7 @@ func TestMasterIsKnownPingTarget(t *testing.T) {
c := cluster.NewCluster()
filerAddr := pb.ServerAddress("10.0.0.20:8888")
c.AddClusterNode("", cluster.FilerType, "dc1", "rack1", filerAddr, "test")
c.AddClusterNode("", cluster.FilerType, "dc1", "rack1", filerAddr, "test", 0)
ms := &MasterServer{
option: &MasterOption{Master: pb.ServerAddress("10.0.0.1:9333")},
@@ -22,6 +22,7 @@ func (ms *MasterServer) ListClusterNodes(ctx context.Context, req *master_pb.Lis
CreatedAtNs: node.CreatedTs.UnixNano(),
DataCenter: string(node.DataCenter),
Rack: string(node.Rack),
MetricsPort: node.MetricsPort,
})
}
return resp, nil
+8 -4
View File
@@ -58,10 +58,14 @@ type MasterOption struct {
DisableHttp bool
MetricsAddress string
MetricsIntervalSec int
IsFollower bool
TelemetryUrl string
TelemetryEnabled bool
VolumeGrowthDisabled bool
// MetricsPort is this master's Prometheus /metrics port (-metricsPort),
// reported in GetMasterConfiguration so the admin server can scrape it.
// Unrelated to MetricsAddress, which is the Prometheus push gateway.
MetricsPort int
IsFollower bool
TelemetryUrl string
TelemetryEnabled bool
VolumeGrowthDisabled bool
}
type MasterServer struct {
+4 -1
View File
@@ -63,7 +63,7 @@ type VolumeServer struct {
}
func NewVolumeServer(adminMux, publicMux *http.ServeMux, ip string,
port int, grpcPort int, publicUrl string, id string,
port int, grpcPort int, metricsPort int, publicUrl string, id string,
folders []string, maxCounts []int32, minFreeSpaces []util.MinFreeSpace, diskTypes []types.DiskType, diskTags [][]string,
idxFolder string,
needleMapKind storage.NeedleMapKind,
@@ -136,6 +136,9 @@ func NewVolumeServer(adminMux, publicMux *http.ServeMux, ip string,
vs.checkWithMaster()
vs.store = storage.NewStore(vs.grpcDialOption, ip, port, grpcPort, publicUrl, id, folders, maxCounts, minFreeSpaces, idxFolder, vs.needleMapKind, diskTypes, diskTags, ldbTimeout, diskProbeConfig)
// Set before the heartbeat goroutine starts below, since the heartbeat is
// what advertises this port to the master.
vs.store.MetricsPort = metricsPort
vs.guard = security.NewGuard(whiteList, signingKey, expiresAfterSec, readSigningKey, readExpiresAfterSec)
handleStaticResources(adminMux)
+12 -7
View File
@@ -60,13 +60,17 @@ type ReadOption struct {
* A VolumeServer contains one Store
*/
type Store struct {
MasterAddress pb.ServerAddress
grpcDialOption grpc.DialOption
volumeSizeLimit uint64 // read from the master
preallocate atomic.Bool // read from the master
Ip string
Port int
GrpcPort int
MasterAddress pb.ServerAddress
grpcDialOption grpc.DialOption
volumeSizeLimit uint64 // read from the master
preallocate atomic.Bool // read from the master
Ip string
Port int
GrpcPort int
// MetricsPort is the Prometheus /metrics port this volume server serves,
// advertised to the master so the admin server can scrape it. 0 when
// disabled.
MetricsPort int
PublicUrl string
Id string // volume server id, independent of ip:port for stable identification
Locations []*DiskLocation
@@ -684,6 +688,7 @@ func (s *Store) CollectHeartbeat() *master_pb.Heartbeat {
Ip: s.Ip,
Port: uint32(s.Port),
GrpcPort: uint32(s.GrpcPort),
MetricsPort: uint32(s.MetricsPort),
PublicUrl: s.PublicUrl,
Id: s.Id,
MaxVolumeCounts: maxVolumeCounts,
+10 -6
View File
@@ -15,9 +15,12 @@ import (
type DataNode struct {
NodeImpl
Ip string
Port int
GrpcPort int
Ip string
Port int
GrpcPort int
// MetricsPort is the volume server's Prometheus /metrics port as reported in
// its heartbeat, or 0 when it does not run a metrics listener.
MetricsPort int
PublicUrl string
LastSeen int64 // unix time in seconds
Counter int // in race condition, the previous dataNode was not dead
@@ -404,9 +407,10 @@ func (dn *DataNode) ToDataNodeInfo(filter VolumeFilter) *master_pb.DataNodeInfo
Id: string(dn.Id()),
// Start from disk usage counters so empty disks are still represented
// even when there are no volumes/EC shards on this data node yet.
DiskInfos: dn.diskUsages.ToDiskInfo(),
GrpcPort: uint32(dn.GrpcPort),
Address: dn.Url(), // ip:port for connecting to the volume server
DiskInfos: dn.diskUsages.ToDiskInfo(),
GrpcPort: uint32(dn.GrpcPort),
MetricsPort: uint32(dn.MetricsPort),
Address: dn.Url(), // ip:port for connecting to the volume server
}
if m.DiskInfos == nil {
m.DiskInfos = make(map[string]*master_pb.DiskInfo)
+9
View File
@@ -149,6 +149,7 @@ type MasterClient struct {
clientType string
clientHost pb.ServerAddress
rack string
metricsPort uint32
currentMaster pb.ServerAddress
currentMasterLock sync.RWMutex
masters pb.ServerDiscovery
@@ -180,6 +181,13 @@ func NewMasterClient(grpcDialOption grpc.DialOption, filerGroup string, clientTy
return mc
}
// SetMetricsPort advertises this node's Prometheus /metrics port to the master,
// so a central scraper can discover it. Must be called before
// KeepConnectedToMaster starts, which is what publishes the value.
func (mc *MasterClient) SetMetricsPort(port uint32) {
mc.metricsPort = port
}
func (mc *MasterClient) SetOnPeerUpdateFn(onPeerUpdate func(update *master_pb.ClusterNodeUpdate, startFrom time.Time)) {
mc.OnPeerUpdateLock.Lock()
mc.OnPeerUpdate = onPeerUpdate
@@ -246,6 +254,7 @@ func (mc *MasterClient) tryConnectToMaster(ctx context.Context, master pb.Server
ClientType: mc.clientType,
ClientAddress: string(mc.clientHost),
Version: version.Version(),
MetricsPort: mc.metricsPort,
}); err != nil {
glog.V(0).Infof("%s.%s masterClient failed to send to %s: %v", mc.FilerGroup, mc.clientType, master, err)
stats.MasterClientConnectCounter.WithLabelValues(stats.FailedToSend).Inc()