mirror of
https://github.com/bitechdev/ResolveSpec.git
synced 2026-09-30 20:11:59 +00:00
Compare commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
1214f69e0c | ||
|
|
89a58ab3a0 |
@@ -449,11 +449,13 @@ no per-client labels so cardinality stays bounded.
|
||||
| Metric | Type | Meaning |
|
||||
|---|---|---|
|
||||
| `clientqueue_requests_total{result}` | counter | `immediate`, `queued`, `rejected_full`, `timeout`, `canceled` |
|
||||
| `clientqueue_enqueued_total` | counter | requests ever placed in a wait queue, whatever happened next (ran, timed out, cancelled) |
|
||||
| `clientqueue_wait_seconds` | histogram | wait for a slot; 0 for requests that ran immediately |
|
||||
| `clientqueue_wait_max_seconds` | gauge | longest wait since process start |
|
||||
| `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
|
||||
|
||||
@@ -70,6 +70,16 @@ var (
|
||||
Help: "Requests currently waiting for a slot",
|
||||
})
|
||||
|
||||
queueEnqueued = promauto.NewCounter(prometheus.CounterOpts{
|
||||
Name: "clientqueue_enqueued_total",
|
||||
Help: "Requests ever placed in a wait queue, whatever their outcome (ran, timed out or cancelled)",
|
||||
})
|
||||
|
||||
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,9 +285,13 @@ 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()
|
||||
queueEnqueued.Inc()
|
||||
|
||||
timer := time.NewTimer(q.cfg.MaxWait)
|
||||
defer timer.Stop()
|
||||
@@ -302,6 +316,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 +348,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()
|
||||
|
||||
@@ -2,6 +2,7 @@ package middleware
|
||||
|
||||
import (
|
||||
"context"
|
||||
"errors"
|
||||
"net/http"
|
||||
"net/http/httptest"
|
||||
"strings"
|
||||
@@ -278,6 +279,7 @@ func TestClientQueueMetrics(t *testing.T) {
|
||||
q := newTestQueue(t, ClientQueueConfig{MaxConcurrent: 2})
|
||||
imm0 := gaugeVal(t, queueRequests.WithLabelValues("immediate"))
|
||||
que0 := gaugeVal(t, queueRequests.WithLabelValues("queued"))
|
||||
enq0 := gaugeVal(t, queueEnqueued)
|
||||
bc0, bs0 := histVal(t, queueBurst)
|
||||
wc0, _ := histVal(t, queueWait)
|
||||
|
||||
@@ -317,6 +319,9 @@ func TestClientQueueMetrics(t *testing.T) {
|
||||
if got := gaugeVal(t, queueRequests.WithLabelValues("immediate")) - imm0; got != 2 {
|
||||
t.Errorf("immediate = %v, want 2", got)
|
||||
}
|
||||
if got := gaugeVal(t, queueEnqueued) - enq0; got != 3 {
|
||||
t.Errorf("enqueued = %v, want 3", got)
|
||||
}
|
||||
if got := gaugeVal(t, queueRequests.WithLabelValues("queued")) - que0; got != 3 {
|
||||
t.Errorf("queued = %v, want 3", got)
|
||||
}
|
||||
@@ -331,3 +336,99 @@ 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")
|
||||
}
|
||||
|
||||
func TestClientQueueEnqueuedCountsAllOutcomes(t *testing.T) {
|
||||
q := newTestQueue(t, ClientQueueConfig{MaxConcurrent: 1, MaxQueue: 1, MaxWait: 50 * time.Millisecond})
|
||||
base := gaugeVal(t, queueEnqueued)
|
||||
if err := q.acquire(t.Context(), "k"); err != nil { // immediate: not enqueued
|
||||
t.Fatal(err)
|
||||
}
|
||||
if err := q.acquire(t.Context(), "k"); !errors.Is(err, errQueueWait) { // enqueued, times out
|
||||
t.Fatalf("err = %v, want timeout", err)
|
||||
}
|
||||
cctx, cancel := context.WithCancel(t.Context())
|
||||
done := make(chan error, 1)
|
||||
go func() { done <- q.acquire(cctx, "k") }() // enqueued, cancelled
|
||||
time.Sleep(10 * time.Millisecond)
|
||||
if err := q.acquire(t.Context(), "k"); !errors.Is(err, errQueueFull) { // rejected: not enqueued
|
||||
t.Fatalf("err = %v, want queue full", err)
|
||||
}
|
||||
cancel()
|
||||
<-done
|
||||
q.release("k")
|
||||
if got := gaugeVal(t, queueEnqueued) - base; got != 2 {
|
||||
t.Fatalf("enqueued = %v, want 2 (timeout + cancel; not immediate or rejected)", got)
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user