From e55db58ca909ad729b0945bb0e06706d466ca9f7 Mon Sep 17 00:00:00 2001 From: Chris Lu Date: Fri, 8 May 2026 17:49:10 -0700 Subject: [PATCH] feat(s3/lifecycle): expose Prometheus metrics (Phase 7) (#9375) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit * feat(s3/lifecycle): expose Prometheus metrics (Phase 7) Five new gauges/counters under the s3_lifecycle subsystem so operators can see what the worker is doing without grepping logs: - dispatch_total{bucket,kind,outcome} — every LifecycleDelete RPC bumps this. Outcome is the proto enum name (DONE, NOOP_RESOLVED, RETRY_LATER, BLOCKED, …) plus a synthetic "RPC_ERROR" for transport failures classified as RETRY_LATER. - schedule_depth{shard} — pending matches in each shard's schedule, sampled on the dispatcher tick. - cursor_min_ts_ns{shard} — per-shard min cursor timestamp; lag is derived as (now - min) by the scrape side. - events_total{shard} — meta-log events the reader fed to the router. - bootstrap_dispatch_total{bucket,kind} — bootstrap-walk dispatches. Test asserts the dispatch counter increments for both DONE and RPC_ERROR paths. * fix(stats): purge lifecycle bucket label series in DeleteBucketMetrics The two new bucket-labeled lifecycle counters (S3LifecycleDispatchCounter, S3LifecycleBootstrapDispatchCounter) weren't included in DeleteBucketMetrics, so explicit bucket teardown left their label series behind — same cardinality leak the existing counters above already avoid. Tack them onto the same DeletePartialMatch chain. --- weed/s3api/s3lifecycle/bootstrap/walker.go | 2 + .../s3lifecycle/dispatcher/dispatcher.go | 10 ++++ .../s3lifecycle/dispatcher/dispatcher_test.go | 50 +++++++++++++++++++ weed/s3api/s3lifecycle/dispatcher/pipeline.go | 5 ++ weed/stats/metrics.go | 48 ++++++++++++++++++ 5 files changed, 115 insertions(+) diff --git a/weed/s3api/s3lifecycle/bootstrap/walker.go b/weed/s3api/s3lifecycle/bootstrap/walker.go index 6990a11cd..e9ee1a547 100644 --- a/weed/s3api/s3lifecycle/bootstrap/walker.go +++ b/weed/s3api/s3lifecycle/bootstrap/walker.go @@ -15,6 +15,7 @@ import ( "github.com/seaweedfs/seaweedfs/weed/glog" "github.com/seaweedfs/seaweedfs/weed/s3api/s3lifecycle" "github.com/seaweedfs/seaweedfs/weed/s3api/s3lifecycle/engine" + "github.com/seaweedfs/seaweedfs/weed/stats" ) // Entry is the routing-relevant slice of a filer entry. SuccessorModTime @@ -160,6 +161,7 @@ func walkEntry(ctx context.Context, snap *engine.Snapshot, bucket string, entry bucket, entry.Path, key.ActionKind, err) return err } + stats.S3LifecycleBootstrapDispatchCounter.WithLabelValues(bucket, key.ActionKind.String()).Inc() } return nil } diff --git a/weed/s3api/s3lifecycle/dispatcher/dispatcher.go b/weed/s3api/s3lifecycle/dispatcher/dispatcher.go index a232347f3..dd4ee81fd 100644 --- a/weed/s3api/s3lifecycle/dispatcher/dispatcher.go +++ b/weed/s3api/s3lifecycle/dispatcher/dispatcher.go @@ -4,6 +4,7 @@ import ( "context" "errors" "fmt" + "strconv" "time" "github.com/seaweedfs/seaweedfs/weed/glog" @@ -11,6 +12,7 @@ import ( "github.com/seaweedfs/seaweedfs/weed/s3api/s3lifecycle" "github.com/seaweedfs/seaweedfs/weed/s3api/s3lifecycle/reader" "github.com/seaweedfs/seaweedfs/weed/s3api/s3lifecycle/router" + "github.com/seaweedfs/seaweedfs/weed/stats" ) // LifecycleClient abstracts the LifecycleDelete RPC so the dispatcher is @@ -136,9 +138,11 @@ func (d *Dispatcher) dispatchOne(ctx context.Context, m router.Match, now time.T // Transport error: classify as RETRY_LATER. The remote handler // already classifies its own filer-side errors; the only path // that hits this branch is the RPC itself failing. + stats.S3LifecycleDispatchCounter.WithLabelValues(m.Bucket, m.Key.ActionKind.String(), "RPC_ERROR").Inc() d.handleRetryLater(ctx, m, fmt.Sprintf("RPC: %v", err), now) return } + stats.S3LifecycleDispatchCounter.WithLabelValues(m.Bucket, m.Key.ActionKind.String(), resp.Outcome.String()).Inc() switch resp.Outcome { case s3_lifecycle_pb.LifecycleDeleteOutcome_DONE, s3_lifecycle_pb.LifecycleDeleteOutcome_NOOP_RESOLVED, @@ -153,6 +157,12 @@ func (d *Dispatcher) dispatchOne(ctx context.Context, m router.Match, now time.T } } +// observeScheduleDepth publishes the current schedule depth gauge for +// this shard. Called from Pipeline on each tick. +func (d *Dispatcher) observeScheduleDepth() { + stats.S3LifecycleScheduleDepthGauge.WithLabelValues(strconv.Itoa(d.ShardID)).Set(float64(d.Schedule.Len())) +} + func (d *Dispatcher) advance(m router.Match) { delete(d.retries, keyOf(m)) d.Cursor.Advance(m.Key, m.EventTs.UnixNano()) diff --git a/weed/s3api/s3lifecycle/dispatcher/dispatcher_test.go b/weed/s3api/s3lifecycle/dispatcher/dispatcher_test.go index 9b900647a..e3e4b3d90 100644 --- a/weed/s3api/s3lifecycle/dispatcher/dispatcher_test.go +++ b/weed/s3api/s3lifecycle/dispatcher/dispatcher_test.go @@ -6,10 +6,12 @@ import ( "testing" "time" + dto "github.com/prometheus/client_model/go" "github.com/seaweedfs/seaweedfs/weed/pb/s3_lifecycle_pb" "github.com/seaweedfs/seaweedfs/weed/s3api/s3lifecycle" "github.com/seaweedfs/seaweedfs/weed/s3api/s3lifecycle/reader" "github.com/seaweedfs/seaweedfs/weed/s3api/s3lifecycle/router" + "github.com/seaweedfs/seaweedfs/weed/stats" ) type fakeClient struct { @@ -266,3 +268,51 @@ func TestDispatchRestartReFreezesNaturally(t *testing.T) { t.Fatal("re-freeze cursor not pinned at event ts") } } + +func TestDispatchEmitsOutcomeMetric(t *testing.T) { + // The Prometheus counter is the operator's signal that lifecycle + // is doing real work; it must increment for every dispatch path + // (success, retry, transport error). Use the shared label tuple + // (bucket, kind, outcome) and read the counter delta. + client := &fakeClient{ + respond: func(call int) (*s3_lifecycle_pb.LifecycleDeleteResponse, error) { + if call == 1 { + return &s3_lifecycle_pb.LifecycleDeleteResponse{ + Outcome: s3_lifecycle_pb.LifecycleDeleteOutcome_DONE, + }, nil + } + return nil, errors.New("network down") + }, + } + d, sched := newDispatcher(client) + t0 := time.Now() + + bucket := "bk" + kind := s3lifecycle.ActionKindExpirationDays.String() + + startDone := counterValue(t, bucket, kind, s3_lifecycle_pb.LifecycleDeleteOutcome_DONE.String()) + startRpc := counterValue(t, bucket, kind, "RPC_ERROR") + + sched.Add(mkMatch(t0, t0, "obj-done")) + d.Tick(context.Background(), t0) + + // Second match exercises the transport-error path. + sched.Add(mkMatch(t0, t0, "obj-fail")) + d.Tick(context.Background(), t0) + + if got := counterValue(t, bucket, kind, s3_lifecycle_pb.LifecycleDeleteOutcome_DONE.String()) - startDone; got != 1 { + t.Errorf("DONE counter delta=%v, want 1", got) + } + if got := counterValue(t, bucket, kind, "RPC_ERROR") - startRpc; got != 1 { + t.Errorf("RPC_ERROR counter delta=%v, want 1", got) + } +} + +func counterValue(t *testing.T, bucket, kind, outcome string) float64 { + t.Helper() + m := &dto.Metric{} + if err := stats.S3LifecycleDispatchCounter.WithLabelValues(bucket, kind, outcome).Write(m); err != nil { + t.Fatalf("read counter: %v", err) + } + return m.GetCounter().GetValue() +} diff --git a/weed/s3api/s3lifecycle/dispatcher/pipeline.go b/weed/s3api/s3lifecycle/dispatcher/pipeline.go index 6e7f153e0..d1cf2796f 100644 --- a/weed/s3api/s3lifecycle/dispatcher/pipeline.go +++ b/weed/s3api/s3lifecycle/dispatcher/pipeline.go @@ -4,6 +4,7 @@ import ( "context" "errors" "fmt" + "strconv" "sync" "time" @@ -13,6 +14,7 @@ import ( "github.com/seaweedfs/seaweedfs/weed/s3api/s3lifecycle/engine" "github.com/seaweedfs/seaweedfs/weed/s3api/s3lifecycle/reader" "github.com/seaweedfs/seaweedfs/weed/s3api/s3lifecycle/router" + "github.com/seaweedfs/seaweedfs/weed/stats" "google.golang.org/grpc/codes" "google.golang.org/grpc/status" ) @@ -246,6 +248,7 @@ func (p *Pipeline) Run(ctx context.Context) error { if st == nil { continue } + stats.S3LifecycleEventCounter.WithLabelValues(strconv.Itoa(ev.ShardID)).Inc() // Always re-fetch the snapshot — caching it across events // means an event arriving between dispatch ticks routes // against a stale snap. With bootstrap injection, events @@ -261,12 +264,14 @@ func (p *Pipeline) Run(ctx context.Context) error { now := time.Now() for _, st := range states { st.dispatch.Tick(runCtx, now) + st.dispatch.observeScheduleDepth() } case <-ct.C: for shardID, st := range states { if err := p.Persister.Save(runCtx, shardID, st.cursor.Snapshot()); err != nil { glog.Warningf("lifecycle cursor checkpoint: shard=%d: %v", shardID, err) } + stats.S3LifecycleCursorMinTsNs.WithLabelValues(strconv.Itoa(shardID)).Set(float64(st.cursor.MinTsNs())) } } } diff --git a/weed/stats/metrics.go b/weed/stats/metrics.go index a66c81e88..b5e7965c4 100644 --- a/weed/stats/metrics.go +++ b/weed/stats/metrics.go @@ -537,6 +537,46 @@ var ( Name: "upload_error_total", Help: "Counter of upload errors by HTTP status code. Code 0 means transport error (no response received).", }, []string{"code"}) + + S3LifecycleDispatchCounter = prometheus.NewCounterVec( + prometheus.CounterOpts{ + Namespace: Namespace, + Subsystem: "s3_lifecycle", + Name: "dispatch_total", + Help: "Counter of LifecycleDelete RPC outcomes by bucket, action kind, and outcome.", + }, []string{"bucket", "kind", "outcome"}) + + S3LifecycleScheduleDepthGauge = prometheus.NewGaugeVec( + prometheus.GaugeOpts{ + Namespace: Namespace, + Subsystem: "s3_lifecycle", + Name: "schedule_depth", + Help: "Number of pending matches in the dispatcher schedule per shard.", + }, []string{"shard"}) + + S3LifecycleCursorMinTsNs = prometheus.NewGaugeVec( + prometheus.GaugeOpts{ + Namespace: Namespace, + Subsystem: "s3_lifecycle", + Name: "cursor_min_ts_ns", + Help: "Per-shard min cursor timestamp in nanoseconds since epoch (lag = now - min).", + }, []string{"shard"}) + + S3LifecycleEventCounter = prometheus.NewCounterVec( + prometheus.CounterOpts{ + Namespace: Namespace, + Subsystem: "s3_lifecycle", + Name: "events_total", + Help: "Counter of meta-log events the reader emitted to the router, partitioned by shard.", + }, []string{"shard"}) + + S3LifecycleBootstrapDispatchCounter = prometheus.NewCounterVec( + prometheus.CounterOpts{ + Namespace: Namespace, + Subsystem: "s3_lifecycle", + Name: "bootstrap_dispatch_total", + Help: "Counter of bootstrap-walk Delete dispatches by bucket and action kind.", + }, []string{"bucket", "kind"}) ) func init() { @@ -607,6 +647,12 @@ func init() { Gather.MustRegister(S3BucketPhysicalSizeBytesGauge) Gather.MustRegister(S3BucketObjectCountGauge) + Gather.MustRegister(S3LifecycleDispatchCounter) + Gather.MustRegister(S3LifecycleScheduleDepthGauge) + Gather.MustRegister(S3LifecycleCursorMinTsNs) + Gather.MustRegister(S3LifecycleEventCounter) + Gather.MustRegister(S3LifecycleBootstrapDispatchCounter) + Gather.MustRegister(UploadErrorCounter) go bucketMetricTTLControl() @@ -679,6 +725,8 @@ func DeleteBucketMetrics(bucket string) { c += S3BucketSizeBytesGauge.DeletePartialMatch(labels) c += S3BucketPhysicalSizeBytesGauge.DeletePartialMatch(labels) c += S3BucketObjectCountGauge.DeletePartialMatch(labels) + c += S3LifecycleDispatchCounter.DeletePartialMatch(labels) + c += S3LifecycleBootstrapDispatchCounter.DeletePartialMatch(labels) glog.V(0).Infof("delete bucket metrics, %s: %d", bucket, c) }