From 89a58ab3a0c3f705d5f245620daea3ff9bfa3f27 Mon Sep 17 00:00:00 2001 From: Hein Date: Wed, 30 Sep 2026 16:14:27 +0200 Subject: [PATCH] feat(middleware): export clientqueue_waiting_clients gauge Counts clients with at least one queued request, updated whenever a client's waiter count crosses zero (enqueue, grant, timeout, cancel). --- pkg/middleware/README.md | 1 + pkg/middleware/clientqueue.go | 14 ++++++ pkg/middleware/clientqueue_test.go | 73 ++++++++++++++++++++++++++++++ 3 files changed, 88 insertions(+) diff --git a/pkg/middleware/README.md b/pkg/middleware/README.md index 5bbfb8e..01b7811 100644 --- a/pkg/middleware/README.md +++ b/pkg/middleware/README.md @@ -454,6 +454,7 @@ no per-client labels so cardinality stays bounded. | `clientqueue_burst_size` | histogram | peak outstanding requests per client busy period | | `clientqueue_burst_max` | gauge | largest burst since process start | | `clientqueue_active` / `clientqueue_queue_depth` | gauge | running / waiting now | +| `clientqueue_waiting_clients` | gauge | clients with at least one request waiting (`queue_depth` counts requests, this counts clients) | | `clientqueue_clients` | gauge | clients currently tracked | A **burst** is one client's busy period: the peak number of its requests running plus waiting between diff --git a/pkg/middleware/clientqueue.go b/pkg/middleware/clientqueue.go index fc19f80..bc647e2 100644 --- a/pkg/middleware/clientqueue.go +++ b/pkg/middleware/clientqueue.go @@ -70,6 +70,11 @@ var ( Help: "Requests currently waiting for a slot", }) + queueWaitingClients = promauto.NewGauge(prometheus.GaugeOpts{ + Name: "clientqueue_waiting_clients", + Help: "Clients currently with at least one request waiting for a slot", + }) + queueClients = promauto.NewGauge(prometheus.GaugeOpts{ Name: "clientqueue_clients", Help: "Clients currently tracked by the queue", @@ -275,6 +280,9 @@ func (q *ClientQueue) acquire(ctx context.Context, key string) error { } w := &queueWaiter{ready: make(chan struct{})} w.elem = c.waiters.PushBack(w) + if c.waiters.Len() == 1 { + queueWaitingClients.Inc() + } c.enter() q.mu.Unlock() queueDepth.Inc() @@ -302,6 +310,9 @@ func (q *ClientQueue) acquire(ctx context.Context, key string) error { q.releaseLocked(key) } else { c.waiters.Remove(w.elem) + if c.waiters.Len() == 0 { + queueWaitingClients.Dec() + } c.leave() queueDepth.Dec() } @@ -331,6 +342,9 @@ func (q *ClientQueue) releaseLocked(key string) { queueActive.Dec() if front := c.waiters.Front(); front != nil { w := c.waiters.Remove(front).(*queueWaiter) + if c.waiters.Len() == 0 { + queueWaitingClients.Dec() + } w.granted = true queueDepth.Dec() queueActive.Inc() diff --git a/pkg/middleware/clientqueue_test.go b/pkg/middleware/clientqueue_test.go index 974cd36..53f59aa 100644 --- a/pkg/middleware/clientqueue_test.go +++ b/pkg/middleware/clientqueue_test.go @@ -2,6 +2,7 @@ package middleware import ( "context" + "errors" "net/http" "net/http/httptest" "strings" @@ -331,3 +332,75 @@ func TestClientQueueMetrics(t *testing.T) { t.Errorf("burst max = %v, want >= 5", got) } } + +func TestClientQueueWaitingClientsGauge(t *testing.T) { + q := newTestQueue(t, ClientQueueConfig{MaxConcurrent: 1, MaxWait: 100 * time.Millisecond}) + base := gaugeVal(t, queueWaitingClients) + waiting := func() float64 { return gaugeVal(t, queueWaitingClients) - base } + waitFor := func(want float64) { + t.Helper() + deadline := time.Now().Add(2 * time.Second) + for waiting() != want { + if time.Now().After(deadline) { + t.Fatalf("waiting clients = %v, want %v", waiting(), want) + } + time.Sleep(time.Millisecond) + } + } + + // Two clients each hold their only slot. + for _, k := range []string{"a", "b"} { + if err := q.acquire(t.Context(), k); err != nil { + t.Fatal(err) + } + } + if waiting() != 0 { + t.Fatalf("waiting = %v with nobody queued", waiting()) + } + + // Two waiters for "a" count as one waiting client, not two. + bg := func(ctx context.Context, k string) chan error { + ch := make(chan error, 1) + go func() { ch <- q.acquire(ctx, k) }() + return ch + } + a1 := bg(t.Context(), "a") + waitFor(1) + a2 := bg(t.Context(), "a") + time.Sleep(10 * time.Millisecond) + waitFor(1) + + // A waiter for "b" adds a second waiting client, and cancelling it drops it. + cctx, cancel := context.WithCancel(t.Context()) + b1 := bg(cctx, "b") + waitFor(2) + cancel() + if err := <-b1; err == nil { + t.Fatal("expected cancellation") + } + waitFor(1) + + // Draining "a": the first release hands over to a1 (one still waits), + // the second hands over to a2 and "a" stops waiting. + q.release("a") + if err := <-a1; err != nil { + t.Fatal(err) + } + waitFor(1) + q.release("a") + if err := <-a2; err != nil { + t.Fatal(err) + } + waitFor(0) + q.release("a") + q.release("a") + + // Timeout: "b" still holds its slot, so a new waiter times out. + b2 := bg(t.Context(), "b") + waitFor(1) + if err := <-b2; !errors.Is(err, errQueueWait) { + t.Fatalf("err = %v, want timeout", err) + } + waitFor(0) + q.release("b") +}