feat: observe gateway request durations

This commit is contained in:
youfak 2026-08-07 17:12:57 +08:00
parent 7f51e333c5
commit 04a9fa89a0
8 changed files with 117 additions and 11 deletions

View File

@ -36,6 +36,6 @@ PostgreSQL Adapter 和对应执行脚本完成前,不把数据库契约记为
并完成代表性集群容量验证。 并完成代表性集群容量验证。
Grafana Overview 与 Prometheus 规则只引用代码已注册的低基数指标。它们覆盖 Gateway Grafana Overview 与 Prometheus 规则只引用代码已注册的低基数指标。它们覆盖 Gateway
请求/Outcome、Controller 容量、Provider、Extraction、Checker 和 Drain不按 Proxy、IP、 请求/端到端时延/Outcome、Controller 容量、Provider、Extraction、Checker 和 Drain不按 Proxy、IP、
Client、Upstream、Worker、Session 或完整 URL 聚合。`go test ./deploy` 会解析两类资产并拒绝 Client、Upstream、Worker、Session 或完整 URL 聚合。`go test ./deploy` 会解析两类资产并拒绝
不存在的指标与禁止标签,指标改名或新增面板时必须同步更新该契约。 不存在的指标与禁止标签,指标改名或新增面板时必须同步更新该契约。

View File

@ -33,6 +33,7 @@ var observableMetricNames = map[string]struct{}{
"proxy_pool_gateway_capacity_invariant_violations_total": {}, "proxy_pool_gateway_capacity_invariant_violations_total": {},
"proxy_pool_gateway_outcome_queue_dropped_total": {}, "proxy_pool_gateway_outcome_queue_dropped_total": {},
"proxy_pool_gateway_outcomes_total": {}, "proxy_pool_gateway_outcomes_total": {},
"proxy_pool_gateway_request_duration_seconds": {},
"proxy_pool_gateway_requests_in_flight": {}, "proxy_pool_gateway_requests_in_flight": {},
"proxy_pool_gateway_requests_total": {}, "proxy_pool_gateway_requests_total": {},
} }
@ -305,7 +306,8 @@ func TestObservabilityAssetsUseRegisteredLowCardinalityMetrics(t *testing.T) {
forbiddenLabel := regexp.MustCompile(`(?:by\s*\([^)]*\b(?:upstream|worker)\b|\{[^}]*\b(?:upstream|worker)\s*=)`) forbiddenLabel := regexp.MustCompile(`(?:by\s*\([^)]*\b(?:upstream|worker)\b|\{[^}]*\b(?:upstream|worker)\s*=)`)
for _, expression := range expressions { for _, expression := range expressions {
for _, name := range metricPattern.FindAllString(expression, -1) { for _, name := range metricPattern.FindAllString(expression, -1) {
if _, exists := observableMetricNames[name]; !exists { baseName := strings.TrimSuffix(strings.TrimSuffix(strings.TrimSuffix(name, "_bucket"), "_count"), "_sum")
if _, exists := observableMetricNames[baseName]; !exists {
t.Errorf("observability expression references unregistered metric %q: %s", name, expression) t.Errorf("observability expression references unregistered metric %q: %s", name, expression)
} }
} }
@ -314,7 +316,7 @@ func TestObservabilityAssetsUseRegisteredLowCardinalityMetrics(t *testing.T) {
} }
} }
for _, required := range []string{ for _, required := range []string{
"proxy_pool_gateway_requests_total", "proxy_pool_controller_capacity_available_slots", "proxy_pool_gateway_requests_total", "proxy_pool_gateway_request_duration_seconds", "proxy_pool_controller_capacity_available_slots",
"proxy_pool_controller_provider_fetch_results_total", "proxy_pool_controller_extraction_requests_total", "proxy_pool_controller_provider_fetch_results_total", "proxy_pool_controller_extraction_requests_total",
"proxy_pool_checker_observations_total", "proxy_pool_checker_observations_total",
} { } {

View File

@ -4,8 +4,8 @@
"graphTooltip": 1, "graphTooltip": 1,
"panels": [ "panels": [
{"type":"timeseries","title":"Gateway QPS","gridPos":{"h":8,"w":8,"x":0,"y":0},"targets":[{"expr":"sum by (protocol) (rate(proxy_pool_gateway_requests_total[1m]))","legendFormat":"{{protocol}}"}]}, {"type":"timeseries","title":"Gateway QPS","gridPos":{"h":8,"w":8,"x":0,"y":0},"targets":[{"expr":"sum by (protocol) (rate(proxy_pool_gateway_requests_total[1m]))","legendFormat":"{{protocol}}"}]},
{"type":"timeseries","title":"Gateway In Flight","gridPos":{"h":8,"w":8,"x":8,"y":0},"targets":[{"expr":"sum by (protocol) (proxy_pool_gateway_requests_in_flight)","legendFormat":"{{protocol}}"}]}, {"type":"timeseries","title":"Gateway p99","gridPos":{"h":8,"w":8,"x":8,"y":0},"targets":[{"expr":"histogram_quantile(0.99, sum by (le, protocol) (rate(proxy_pool_gateway_request_duration_seconds_bucket[5m])))","legendFormat":"{{protocol}}"}]},
{"type":"timeseries","title":"Active CONNECT Tunnels","gridPos":{"h":8,"w":8,"x":16,"y":0},"targets":[{"expr":"sum(proxy_pool_gateway_active_tunnels)","legendFormat":"tunnels"}]}, {"type":"timeseries","title":"Gateway In Flight and Tunnels","gridPos":{"h":8,"w":8,"x":16,"y":0},"targets":[{"expr":"sum by (protocol) (proxy_pool_gateway_requests_in_flight)","legendFormat":"in-flight {{protocol}}"},{"expr":"sum(proxy_pool_gateway_active_tunnels)","legendFormat":"tunnels"}]},
{"type":"timeseries","title":"Gateway Outcome Failures","gridPos":{"h":8,"w":12,"x":0,"y":8},"targets":[{"expr":"sum by (stage) (rate(proxy_pool_gateway_outcomes_total{result=\"failure\"}[5m]))","legendFormat":"{{stage}}"}]}, {"type":"timeseries","title":"Gateway Outcome Failures","gridPos":{"h":8,"w":12,"x":0,"y":8},"targets":[{"expr":"sum by (stage) (rate(proxy_pool_gateway_outcomes_total{result=\"failure\"}[5m]))","legendFormat":"{{stage}}"}]},
{"type":"timeseries","title":"Gateway Outcome Queue Drops","gridPos":{"h":8,"w":12,"x":12,"y":8},"targets":[{"expr":"sum(rate(proxy_pool_gateway_outcome_queue_dropped_total[5m]))","legendFormat":"drops/s"}]}, {"type":"timeseries","title":"Gateway Outcome Queue Drops","gridPos":{"h":8,"w":12,"x":12,"y":8},"targets":[{"expr":"sum(rate(proxy_pool_gateway_outcome_queue_dropped_total[5m]))","legendFormat":"drops/s"}]},
{"type":"timeseries","title":"Controller Capacity","gridPos":{"h":8,"w":12,"x":0,"y":16},"targets":[{"expr":"sum(proxy_pool_controller_capacity_available_slots)","legendFormat":"available slots"},{"expr":"sum(proxy_pool_controller_capacity_effective_slots)","legendFormat":"effective slots"},{"expr":"sum(proxy_pool_controller_capacity_pending_expected_proxies)","legendFormat":"pending proxies"}]}, {"type":"timeseries","title":"Controller Capacity","gridPos":{"h":8,"w":12,"x":0,"y":16},"targets":[{"expr":"sum(proxy_pool_controller_capacity_available_slots)","legendFormat":"available slots"},{"expr":"sum(proxy_pool_controller_capacity_effective_slots)","legendFormat":"effective slots"},{"expr":"sum(proxy_pool_controller_capacity_pending_expected_proxies)","legendFormat":"pending proxies"}]},

View File

@ -91,6 +91,7 @@ type OutcomeRecorder interface {
type RequestMetricsObserver interface { type RequestMetricsObserver interface {
ObserveRequestStarted(protocol string) ObserveRequestStarted(protocol string)
ObserveRequestFinished(protocol string) ObserveRequestFinished(protocol string)
ObserveRequestDuration(protocol string, duration time.Duration)
ObserveTunnelOpened() ObserveTunnelOpened()
ObserveTunnelClosed() ObserveTunnelClosed()
} }
@ -182,8 +183,12 @@ func (handler *Handler) ServeHTTP(writer http.ResponseWriter, request *http.Requ
} }
defer handler.finishRequest() defer handler.finishRequest()
protocol := requestMetricProtocol(request) protocol := requestMetricProtocol(request)
started := time.Now()
handler.observeRequestStarted(protocol) handler.observeRequestStarted(protocol)
defer handler.observeRequestFinished(protocol) defer func() {
handler.observeRequestFinished(protocol)
handler.observeRequestDuration(protocol, time.Since(started))
}()
if handler.inFlight != nil { if handler.inFlight != nil {
select { select {
case handler.inFlight <- struct{}{}: case handler.inFlight <- struct{}{}:
@ -495,6 +500,12 @@ func (handler *Handler) observeRequestFinished(protocol string) {
} }
} }
func (handler *Handler) observeRequestDuration(protocol string, duration time.Duration) {
if handler.metrics != nil {
handler.metrics.ObserveRequestDuration(protocol, duration)
}
}
func requestMetricProtocol(request *http.Request) string { func requestMetricProtocol(request *http.Request) string {
if request != nil && request.Method == http.MethodConnect { if request != nil && request.Method == http.MethodConnect {
return "CONNECT" return "CONNECT"

View File

@ -107,12 +107,18 @@ func TestHandlerReportsLocalRequestAndTunnelLifecycles(t *testing.T) {
if got := metrics.finished("HTTP"); got != 1 { if got := metrics.finished("HTTP"); got != 1 {
t.Fatalf("HTTP request finishes = %d, want 1", got) t.Fatalf("HTTP request finishes = %d, want 1", got)
} }
if got := metrics.durationCount("HTTP"); got != 1 {
t.Fatalf("HTTP request durations = %d, want 1", got)
}
if got := metrics.started("CONNECT"); got != 1 { if got := metrics.started("CONNECT"); got != 1 {
t.Fatalf("CONNECT request starts = %d, want 1", got) t.Fatalf("CONNECT request starts = %d, want 1", got)
} }
if got := metrics.finished("CONNECT"); got != 1 { if got := metrics.finished("CONNECT"); got != 1 {
t.Fatalf("CONNECT request finishes = %d, want 1", got) t.Fatalf("CONNECT request finishes = %d, want 1", got)
} }
if got := metrics.durationCount("CONNECT"); got != 1 {
t.Fatalf("CONNECT request durations = %d, want 1", got)
}
if metrics.opened != 1 || metrics.closed != 1 { if metrics.opened != 1 || metrics.closed != 1 {
t.Fatalf("tunnel lifecycle = opened:%d closed:%d, want 1:1", metrics.opened, metrics.closed) t.Fatalf("tunnel lifecycle = opened:%d closed:%d, want 1:1", metrics.opened, metrics.closed)
} }
@ -963,6 +969,7 @@ type recordingRequestMetrics struct {
mu sync.Mutex mu sync.Mutex
starts map[string]int starts map[string]int
finishes map[string]int finishes map[string]int
durations map[string][]time.Duration
opened int opened int
closed int closed int
} }
@ -985,6 +992,15 @@ func (metrics *recordingRequestMetrics) ObserveRequestFinished(protocol string)
metrics.finishes[protocol]++ metrics.finishes[protocol]++
} }
func (metrics *recordingRequestMetrics) ObserveRequestDuration(protocol string, duration time.Duration) {
metrics.mu.Lock()
defer metrics.mu.Unlock()
if metrics.durations == nil {
metrics.durations = make(map[string][]time.Duration)
}
metrics.durations[protocol] = append(metrics.durations[protocol], duration)
}
func (metrics *recordingRequestMetrics) ObserveTunnelOpened() { func (metrics *recordingRequestMetrics) ObserveTunnelOpened() {
metrics.mu.Lock() metrics.mu.Lock()
defer metrics.mu.Unlock() defer metrics.mu.Unlock()
@ -1009,6 +1025,12 @@ func (metrics *recordingRequestMetrics) finished(protocol string) int {
return metrics.finishes[protocol] return metrics.finishes[protocol]
} }
func (metrics *recordingRequestMetrics) durationCount(protocol string) int {
metrics.mu.Lock()
defer metrics.mu.Unlock()
return len(metrics.durations[protocol])
}
type outcomeRecorder struct { type outcomeRecorder struct {
mu sync.Mutex mu sync.Mutex
events []outcomeDomain.Event events []outcomeDomain.Event

View File

@ -3,6 +3,7 @@ package metrics
import ( import (
"errors" "errors"
"fmt" "fmt"
"time"
"github.com/prometheus/client_golang/prometheus" "github.com/prometheus/client_golang/prometheus"
@ -20,6 +21,8 @@ type GatewayCollector struct {
connectRequests prometheus.Counter connectRequests prometheus.Counter
httpInFlight prometheus.Gauge httpInFlight prometheus.Gauge
connectInFlight prometheus.Gauge connectInFlight prometheus.Gauge
httpDuration prometheus.Observer
connectDuration prometheus.Observer
activeTunnels prometheus.Gauge activeTunnels prometheus.Gauge
} }
@ -67,6 +70,14 @@ func NewGatewayCollector(registerer prometheus.Registerer) (*GatewayCollector, e
if err != nil { if err != nil {
return nil, err return nil, err
} }
durations, err := registerHistogramVec(registerer, prometheus.NewHistogramVec(prometheus.HistogramOpts{
Namespace: "proxy_pool", Subsystem: "gateway", Name: "request_duration_seconds",
Help: "Gateway request duration from admission to completion by fixed proxy protocol.",
Buckets: []float64{0.001, 0.005, 0.01, 0.025, 0.05, 0.1, 0.25, 0.5, 1, 2.5, 5, 10, 30, 60, 300},
}, []string{"protocol"}))
if err != nil {
return nil, err
}
activeTunnels, err := registerGauge(registerer, prometheus.NewGauge(prometheus.GaugeOpts{ activeTunnels, err := registerGauge(registerer, prometheus.NewGauge(prometheus.GaugeOpts{
Namespace: "proxy_pool", Subsystem: "gateway", Name: "active_tunnels", Namespace: "proxy_pool", Subsystem: "gateway", Name: "active_tunnels",
Help: "Number of established CONNECT tunnels currently relaying through this Gateway Worker.", Help: "Number of established CONNECT tunnels currently relaying through this Gateway Worker.",
@ -78,6 +89,7 @@ func NewGatewayCollector(registerer prometheus.Registerer) (*GatewayCollector, e
outcomes: outcomes, dropped: dropped, invariants: invariants, outcomes: outcomes, dropped: dropped, invariants: invariants,
httpRequests: requests.WithLabelValues("HTTP"), connectRequests: requests.WithLabelValues("CONNECT"), httpRequests: requests.WithLabelValues("HTTP"), connectRequests: requests.WithLabelValues("CONNECT"),
httpInFlight: inFlight.WithLabelValues("HTTP"), connectInFlight: inFlight.WithLabelValues("CONNECT"), httpInFlight: inFlight.WithLabelValues("HTTP"), connectInFlight: inFlight.WithLabelValues("CONNECT"),
httpDuration: durations.WithLabelValues("HTTP"), connectDuration: durations.WithLabelValues("CONNECT"),
activeTunnels: activeTunnels, activeTunnels: activeTunnels,
}, nil }, nil
} }
@ -145,6 +157,22 @@ func (collector *GatewayCollector) ObserveRequestFinished(protocol string) {
} }
} }
func (collector *GatewayCollector) ObserveRequestDuration(protocol string, duration time.Duration) {
if collector == nil || duration < 0 {
return
}
switch protocol {
case "HTTP":
if collector.httpDuration != nil {
collector.httpDuration.Observe(duration.Seconds())
}
case "CONNECT":
if collector.connectDuration != nil {
collector.connectDuration.Observe(duration.Seconds())
}
}
}
func (collector *GatewayCollector) ObserveTunnelOpened() { func (collector *GatewayCollector) ObserveTunnelOpened() {
if collector != nil && collector.activeTunnels != nil { if collector != nil && collector.activeTunnels != nil {
collector.activeTunnels.Inc() collector.activeTunnels.Inc()
@ -205,6 +233,22 @@ func registerGaugeVec(registerer prometheus.Registerer, candidate *prometheus.Ga
} }
} }
func registerHistogramVec(registerer prometheus.Registerer, candidate *prometheus.HistogramVec) (*prometheus.HistogramVec, error) {
if err := registerer.Register(candidate); err == nil {
return candidate, nil
} else {
var registered prometheus.AlreadyRegisteredError
if !errors.As(err, &registered) {
return nil, fmt.Errorf("register histogram vector: %w", err)
}
existing, ok := registered.ExistingCollector.(*prometheus.HistogramVec)
if !ok {
return nil, fmt.Errorf("register histogram vector: existing collector has unexpected type")
}
return existing, nil
}
}
func validGatewayStage(stage outcomeDomain.Stage) bool { func validGatewayStage(stage outcomeDomain.Stage) bool {
switch stage { switch stage {
case outcomeDomain.StageDial, outcomeDomain.StageProxyHandshake, case outcomeDomain.StageDial, outcomeDomain.StageProxyHandshake,

View File

@ -2,6 +2,7 @@ package metrics
import ( import (
"testing" "testing"
"time"
"github.com/prometheus/client_golang/prometheus" "github.com/prometheus/client_golang/prometheus"
@ -27,6 +28,9 @@ func TestGatewayCollectorRecordsOnlyFixedDimensions(t *testing.T) {
collector.ObserveRequestStarted("CONNECT") collector.ObserveRequestStarted("CONNECT")
collector.ObserveRequestFinished("HTTP") collector.ObserveRequestFinished("HTTP")
collector.ObserveRequestFinished("CONNECT") collector.ObserveRequestFinished("CONNECT")
collector.ObserveRequestDuration("HTTP", 25*time.Millisecond)
collector.ObserveRequestDuration("CONNECT", 2*time.Second)
collector.ObserveRequestDuration("unknown", time.Second)
collector.ObserveTunnelOpened() collector.ObserveTunnelOpened()
collector.ObserveTunnelClosed() collector.ObserveTunnelClosed()
collector.ObserveRequestStarted("unknown") collector.ObserveRequestStarted("unknown")
@ -54,6 +58,8 @@ func TestGatewayCollectorRecordsOnlyFixedDimensions(t *testing.T) {
"protocol": "CONNECT", "protocol": "CONNECT",
}, 0) }, 0)
assertMetricValue(t, registry, "proxy_pool_gateway_active_tunnels", nil, 0) assertMetricValue(t, registry, "proxy_pool_gateway_active_tunnels", nil, 0)
assertHistogramSampleCount(t, registry, "proxy_pool_gateway_request_duration_seconds", map[string]string{"protocol": "HTTP"}, 1)
assertHistogramSampleCount(t, registry, "proxy_pool_gateway_request_duration_seconds", map[string]string{"protocol": "CONNECT"}, 1)
} }
func TestNewGatewayCollectorReusesRegisteredCollectors(t *testing.T) { func TestNewGatewayCollectorReusesRegisteredCollectors(t *testing.T) {
@ -70,3 +76,22 @@ func TestNewGatewayCollectorReusesRegisteredCollectors(t *testing.T) {
second.ObserveDropped() second.ObserveDropped()
assertMetricValue(t, registry, "proxy_pool_gateway_outcome_queue_dropped_total", nil, 2) assertMetricValue(t, registry, "proxy_pool_gateway_outcome_queue_dropped_total", nil, 2)
} }
func assertHistogramSampleCount(t *testing.T, registry *prometheus.Registry, name string, labels map[string]string, want uint64) {
t.Helper()
metrics, err := registry.Gather()
if err != nil {
t.Fatalf("Gather() = %v", err)
}
for _, family := range metrics {
if family.GetName() != name {
continue
}
for _, metric := range family.GetMetric() {
if metricLabelsMatch(metric.GetLabel(), labels) && metric.GetHistogram().GetSampleCount() == want {
return
}
}
}
t.Fatalf("histogram %s labels=%v count=%d was not found", name, labels, want)
}

View File

@ -20,6 +20,8 @@
- Grafana Overview 与 Prometheus 告警已从早期失效指标迁移到当前代码注册的低基数指标,覆盖 - Grafana Overview 与 Prometheus 告警已从早期失效指标迁移到当前代码注册的低基数指标,覆盖
Gateway、Controller 容量、Provider、Extraction、Checker 与 Drain。部署测试会解析仪表盘和 Gateway、Controller 容量、Provider、Extraction、Checker 与 Drain。部署测试会解析仪表盘和
规则,拒绝未注册指标以及 `upstream`/`worker` 聚合或 selector 标签,防止观测资产再次漂移。 规则,拒绝未注册指标以及 `upstream`/`worker` 聚合或 selector 标签,防止观测资产再次漂移。
- Gateway 新增按固定 `HTTP`/`CONNECT` protocol 标签聚合的请求总耗时 Histogram覆盖从准入到
HTTP 响应或 CONNECT 隧道结束的完整生命周期Grafana Overview 已恢复该真实指标的 p99 面板。
- 全仓 `go test -count=1 -timeout 60s ./...`、`go vet ./...`、`go build ./...`、 - 全仓 `go test -count=1 -timeout 60s ./...`、`go vet ./...`、`go build ./...`、
Protobuf descriptor、Kustomize Base 渲染及开发证书 SAN/SPIFFE 校验均通过。Compose Protobuf descriptor、Kustomize Base 渲染及开发证书 SAN/SPIFFE 校验均通过。Compose
容器端到端启动在拉取 Dockerfile 前端与监控镜像时受 Docker Desktop HTTPS 代理缺失阻断, 容器端到端启动在拉取 Dockerfile 前端与监控镜像时受 Docker Desktop HTTPS 代理缺失阻断,