From 2af504fd495656655f01200eccd2ec7392ec9d47 Mon Sep 17 00:00:00 2001 From: youfak Date: Sun, 2 Aug 2026 09:50:34 +0800 Subject: [PATCH] feat: expose gateway outcome metrics --- README.md | 2 + docs/configuration/reference.md | 7 +- docs/development/implementation-plan.md | 6 +- docs/operations/runbook.md | 8 ++ docs/requirements/completion-audit.md | 3 +- docs/requirements/traceability.md | 2 +- internal/domain/outcome/metrics.go | 9 ++ internal/gateway/bootstrap/bootstrap.go | 11 +++ internal/gateway/bootstrap/bootstrap_test.go | 13 +++ internal/gateway/outcome/queue.go | 14 ++- internal/gateway/outcome/queue_test.go | 32 +++++++ internal/platform/metrics/checker.go | 4 +- internal/platform/metrics/gateway.go | 99 ++++++++++++++++++++ internal/platform/metrics/gateway_test.go | 44 +++++++++ 14 files changed, 245 insertions(+), 9 deletions(-) create mode 100644 internal/domain/outcome/metrics.go create mode 100644 internal/platform/metrics/gateway.go create mode 100644 internal/platform/metrics/gateway_test.go diff --git a/README.md b/README.md index c432d74..d78616d 100644 --- a/README.md +++ b/README.md @@ -41,6 +41,8 @@ Proxy Pool 用 Controller 协调这些变化,并让 Gateway 数据面只消费 PostgreSQL 迁移与 Redis 活动池,并支持联动优雅停机。 - **Checker 指标**:Metrics 启用时暴露 Checker 任务下发与 Observation 接受/拒绝计数; 标签仅使用固定检查级别和结果,不记录 Proxy、IP、URL 或凭据。 +- **Gateway 指标**:Metrics 启用时暴露代理尝试的固定阶段成功/失败计数,以及本地 + Outcome 队列满后的丢弃计数;不记录 Proxy、路由、目标、客户端或凭据。 - **PostgreSQL 管理面**:持久化配置版本、Upstream/Routing 管理状态、Admin 审计与 Outbox;不保存 Proxy 明细或逐次提取记录。 - **Gateway 组件**:HTTP 正向代理、HTTPS CONNECT、双向 Tunnel、重试、超时、 diff --git a/docs/configuration/reference.md b/docs/configuration/reference.md index 2852dda..fdf642d 100644 --- a/docs/configuration/reference.md +++ b/docs/configuration/reference.md @@ -513,8 +513,11 @@ Metrics 启用时 `listen` 必须是合法 `host:port`。该入口固定提供 ` 网络策略限制。当前 `/metrics` 包含 Go/进程基础指标,以及 Checker 的 `proxy_pool_checker_tasks_dispatched_total{level}` 和 `proxy_pool_checker_observations_total{level,result}`。`level` 固定为 BASIC、EGRESS、TARGET, -`result` 固定为 accepted、rejected;绝不包含 Proxy、IP、Checker、Client、URL 或凭据标签。 -Provider、提取和容量等业务指标仍在后续实施范围。Distribution 启用时 `/readyz` 只以 Redis 活动池为 +`result` 固定为 accepted、rejected。Gateway 另暴露 +`proxy_pool_gateway_outcomes_total{stage,result}` 与 +`proxy_pool_gateway_outcome_queue_dropped_total`;`stage` 固定为 DIAL、PROXY_HANDSHAKE、 +RESPONSE_HEADERS、TUNNEL,`result` 固定为 success、failure。绝不包含 Proxy、IP、Checker、 +Client、路由、目标 URL 或凭据标签。Provider、提取和容量等业务指标仍在后续实施范围。Distribution 启用时 `/readyz` 只以 Redis 活动池为 服务流量门槛,PostgreSQL 故障由 Admin 接口独立报告。Metrics 开关或监听地址 变更需要重启 Controller。 diff --git a/docs/development/implementation-plan.md b/docs/development/implementation-plan.md index 090ecb6..d402b92 100644 --- a/docs/development/implementation-plan.md +++ b/docs/development/implementation-plan.md @@ -336,8 +336,10 @@ Profile 在启用 Routing 与 Upstream 的组合上才进入调度。 当前进度(2026-08-02):已接入 Checker 任务流指标 `proxy_pool_checker_tasks_dispatched_total{level}` 与 `proxy_pool_checker_observations_total{level,result}`;标签值仅允许固定的 -`BASIC`、`EGRESS`、`TARGET` 级别及 `accepted`、`rejected` 结果。Provider、提取、 -容量指标与密钥安全的结构化日志仍待补齐。 +`BASIC`、`EGRESS`、`TARGET` 级别及 `accepted`、`rejected` 结果。Gateway 还暴露 +`proxy_pool_gateway_outcomes_total{stage,result}` 和 +`proxy_pool_gateway_outcome_queue_dropped_total`,其中阶段和结果均为固定枚举。Provider、 +提取、容量指标与密钥安全的结构化日志仍待补齐。 ## Task 14: Documentation, Examples, and Diagrams diff --git a/docs/operations/runbook.md b/docs/operations/runbook.md index 8889192..99ae977 100644 --- a/docs/operations/runbook.md +++ b/docs/operations/runbook.md @@ -136,6 +136,14 @@ Checker 指标使用固定标签集: `level` 仅为 BASIC、EGRESS、TARGET,`result` 仅为 accepted、rejected。不得将 Proxy ID、IP、Checker ID、目标 URL、Client ID 或凭据加入指标标签。 +Gateway 指标同样使用固定标签集: + +- `proxy_pool_gateway_outcomes_total{stage,result}`:每次代理尝试在最远完成阶段的成功/失败数量。 +- `proxy_pool_gateway_outcome_queue_dropped_total`:Outcome 本地有界队列已满后丢弃的观测数量。 + +`stage` 仅为 DIAL、PROXY_HANDSHAKE、RESPONSE_HEADERS、TUNNEL,`result` 仅为 success、failure。 +前者在写入本地队列前计数,后者用于识别观测背压;二者均不包含 Proxy ID、路由、目标、客户端或凭据。 + ### 3.4 发布顺序与兼容性 1. 先做向前兼容数据库迁移。 diff --git a/docs/requirements/completion-audit.md b/docs/requirements/completion-audit.md index b7b64b1..42117b6 100644 --- a/docs/requirements/completion-audit.md +++ b/docs/requirements/completion-audit.md @@ -104,7 +104,8 @@ Controller/Gateway 入口,完整 mTLS 运行时拓扑仍只有静态验证。 响应校验场景;`proxy-checker` 的 BASIC 任务进程已经完成, `proxy-controller` 已完成 Admin/Distribution/Metrics 与 PostgreSQL/Redis 启动装配,`proxy-gateway` 已完成 - HTTP/Metrics 与控制面 Session 装配,但 Provider 和业务指标链未闭环。 + HTTP/Metrics 与控制面 Session 装配;Checker 与 Gateway Outcome 的低基数业务指标 + 已接入,但 Provider、提取和容量指标链未闭环。 2. Gateway 的生产连接池调优与代表性流量压测。 3. Provider 分布式 singleflight/Leader、长期凭据回收和累计额度执行器。 4. Controller 的 PostgreSQL 连接池、迁移和 pgx Adapter 启动装配已完成; diff --git a/docs/requirements/traceability.md b/docs/requirements/traceability.md index 8167c31..999af15 100644 --- a/docs/requirements/traceability.md +++ b/docs/requirements/traceability.md @@ -87,5 +87,5 @@ | OPS-001 | 配置校验后构建不可变快照并原子替换 | 8959-8999 | 100k 索引、版本/epoch 与并发 Apply/Acquire 测试 | | OPS-002 | 优雅停机停止新请求/Fetch,等待现有流量后超时关闭 | 8981-9000 | Provider Run 收敛与 `Handler.Shutdown` HTTP 排空、Hijacked CONNECT 超时关闭测试 | | OPS-003 | PostgreSQL 只保存管理修订、Upstream/Routing 状态、Admin 审计与 Outbox | 当前会话 | ADR-006、`adminstate` 公用契约和六表 Schema 边界测试;真实 PostgreSQL 契约待完成 | -| OBS-001 | 指标禁止 Proxy IP、session、Client、完整 URL 高基数标签 | 9001-9029 | Controller Prometheus/探针模块已实现;Checker Collector 仅暴露 `proxy_pool_checker_tasks_dispatched_total{level}` 与 `proxy_pool_checker_observations_total{level,result}`,并以注册表测试固定标签集。Provider、提取和容量业务指标待实现 | +| OBS-001 | 指标禁止 Proxy IP、session、Client、完整 URL 高基数标签 | 9001-9029 | Controller Prometheus/探针模块已实现;Checker Collector 暴露固定等级和结果,Gateway Collector 暴露 `proxy_pool_gateway_outcomes_total{stage,result}` 与 `proxy_pool_gateway_outcome_queue_dropped_total`,均以注册表测试固定标签集。Provider、提取和容量业务指标待实现 | | TEST-001 | 覆盖对话中列出的 11 个关键并发与故障场景 | 9030-9082 | 测试清单;Redis 活动池由 Memory/Redis 公用契约覆盖,跨进程故障场景仍按清单推进 | diff --git a/internal/domain/outcome/metrics.go b/internal/domain/outcome/metrics.go new file mode 100644 index 0000000..35dea1f --- /dev/null +++ b/internal/domain/outcome/metrics.go @@ -0,0 +1,9 @@ +package outcome + +// MetricsObserver receives aggregate Gateway outcome signals on the local +// request path. Implementations must only perform bounded in-process work and +// must never block request forwarding. +type MetricsObserver interface { + Observe(Event) + ObserveDropped() +} diff --git a/internal/gateway/bootstrap/bootstrap.go b/internal/gateway/bootstrap/bootstrap.go index 63a070a..210d6bc 100644 --- a/internal/gateway/bootstrap/bootstrap.go +++ b/internal/gateway/bootstrap/bootstrap.go @@ -19,6 +19,7 @@ import ( controlplanev1 "proxy-pool/gen/controlplane/v1" "proxy-pool/internal/config" "proxy-pool/internal/controlplane/clienttransport" + outcomeDomain "proxy-pool/internal/domain/outcome" proxyDomain "proxy-pool/internal/domain/proxy" "proxy-pool/internal/domain/workerruntime" "proxy-pool/internal/gateway/controlplane" @@ -146,8 +147,18 @@ func newRuntime(ctx context.Context, configuration *config.Config, options Optio return nil, fmt.Errorf("build gateway target policy: %w", err) } proxyTransport := transport.New(transport.Config{}, snapshotCredentialResolver{store: store}) + var outcomeMetrics outcomeDomain.MetricsObserver + if configuration.Metrics.Enabled { + collector, err := platformMetrics.NewGatewayCollector(prometheus.DefaultRegisterer) + if err != nil { + proxyTransport.CloseIdleConnections() + return nil, fmt.Errorf("build gateway outcome metrics: %w", err) + } + outcomeMetrics = collector + } outcomes, err := gatewayOutcome.NewQueue(gatewayOutcome.QueueOptions{ Capacity: defaultOutcomeQueueCapacity, MaxBatch: min(defaultOutcomeBatchSize, configuration.ControlPlane.MaxRuntimeCounters), + Metrics: outcomeMetrics, }) if err != nil { proxyTransport.CloseIdleConnections() diff --git a/internal/gateway/bootstrap/bootstrap_test.go b/internal/gateway/bootstrap/bootstrap_test.go index 04ebee3..bf34614 100644 --- a/internal/gateway/bootstrap/bootstrap_test.go +++ b/internal/gateway/bootstrap/bootstrap_test.go @@ -11,6 +11,7 @@ import ( "net/url" "os" "path/filepath" + "strings" "testing" "time" @@ -79,6 +80,18 @@ func TestRunServesGatewayOnlyAfterApplyingControlPlaneSnapshot(t *testing.T) { t.Fatalf("gateway did not become ready: %v", err) } } + metricsResponse, err := http.Get("http://" + metricsListener.Addr().String() + "/metrics") + if err != nil { + t.Fatalf("GET gateway metrics: %v", err) + } + metricsBody, err := io.ReadAll(metricsResponse.Body) + _ = metricsResponse.Body.Close() + if err != nil { + t.Fatalf("ReadAll(gateway metrics) = %v", err) + } + if metricsResponse.StatusCode != http.StatusOK || !strings.Contains(string(metricsBody), "proxy_pool_gateway_outcome_queue_dropped_total") { + t.Fatalf("gateway metrics = (%d, %q), want outcome queue metric", metricsResponse.StatusCode, metricsBody) + } proxyURL, err := url.Parse("http://" + proxyListener.Addr().String()) if err != nil { t.Fatalf("Parse(proxy URL): %v", err) diff --git a/internal/gateway/outcome/queue.go b/internal/gateway/outcome/queue.go index 474967e..44154f7 100644 --- a/internal/gateway/outcome/queue.go +++ b/internal/gateway/outcome/queue.go @@ -15,6 +15,7 @@ var ErrInvalidQueueOptions = errors.New("invalid gateway outcome queue options") type QueueOptions struct { Capacity int MaxBatch int + Metrics domain.MetricsObserver } // Queue accepts observations from the request path without blocking. Events @@ -23,6 +24,7 @@ type QueueOptions struct { type Queue struct { events chan domain.Event maxBatch int + metrics domain.MetricsObserver dropped atomic.Uint64 } @@ -30,17 +32,27 @@ func NewQueue(options QueueOptions) (*Queue, error) { if options.Capacity <= 0 || options.MaxBatch <= 0 || options.MaxBatch > options.Capacity { return nil, ErrInvalidQueueOptions } - return &Queue{events: make(chan domain.Event, options.Capacity), maxBatch: options.MaxBatch}, nil + return &Queue{ + events: make(chan domain.Event, options.Capacity), + maxBatch: options.MaxBatch, + metrics: options.Metrics, + }, nil } func (queue *Queue) Record(event domain.Event) { if queue == nil { return } + if queue.metrics != nil { + queue.metrics.Observe(event) + } select { case queue.events <- event: default: queue.dropped.Add(1) + if queue.metrics != nil { + queue.metrics.ObserveDropped() + } } } diff --git a/internal/gateway/outcome/queue_test.go b/internal/gateway/outcome/queue_test.go index 0255e82..a2ce260 100644 --- a/internal/gateway/outcome/queue_test.go +++ b/internal/gateway/outcome/queue_test.go @@ -8,6 +8,19 @@ import ( domain "proxy-pool/internal/domain/outcome" ) +type metricsRecorder struct { + events []domain.Event + dropped int +} + +func (recorder *metricsRecorder) Observe(event domain.Event) { + recorder.events = append(recorder.events, event) +} + +func (recorder *metricsRecorder) ObserveDropped() { + recorder.dropped++ +} + func TestQueueBatchesWithoutBlockingAndCountsOverflow(t *testing.T) { queue, err := NewQueue(QueueOptions{Capacity: 2, MaxBatch: 2}) if err != nil { @@ -39,3 +52,22 @@ func TestQueueHonorsCanceledContext(t *testing.T) { t.Fatalf("Next(canceled) = %v, want context canceled", err) } } + +func TestQueueReportsAttemptedOutcomesAndOverflow(t *testing.T) { + metrics := &metricsRecorder{} + queue, err := NewQueue(QueueOptions{Capacity: 1, MaxBatch: 1, Metrics: metrics}) + if err != nil { + t.Fatalf("NewQueue() = %v", err) + } + first := domain.Event{ProxyID: "proxy-a", Stage: domain.StageDial, Success: true} + second := domain.Event{ProxyID: "proxy-b", Stage: domain.StageTunnel, Success: false} + queue.Record(first) + queue.Record(second) + + if len(metrics.events) != 2 || metrics.events[0] != first || metrics.events[1] != second { + t.Fatalf("observed events = %+v, want both attempts", metrics.events) + } + if metrics.dropped != 1 { + t.Fatalf("observed dropped = %d, want 1", metrics.dropped) + } +} diff --git a/internal/platform/metrics/checker.go b/internal/platform/metrics/checker.go index 67fee3a..e083838 100644 --- a/internal/platform/metrics/checker.go +++ b/internal/platform/metrics/checker.go @@ -60,11 +60,11 @@ func registerCounterVec(registerer prometheus.Registerer, candidate *prometheus. } else { var registered prometheus.AlreadyRegisteredError if !errors.As(err, ®istered) { - return nil, fmt.Errorf("register checker metrics: %w", err) + return nil, fmt.Errorf("register counter vector: %w", err) } existing, ok := registered.ExistingCollector.(*prometheus.CounterVec) if !ok { - return nil, fmt.Errorf("register checker metrics: existing collector has unexpected type") + return nil, fmt.Errorf("register counter vector: existing collector has unexpected type") } return existing, nil } diff --git a/internal/platform/metrics/gateway.go b/internal/platform/metrics/gateway.go new file mode 100644 index 0000000..4148efa --- /dev/null +++ b/internal/platform/metrics/gateway.go @@ -0,0 +1,99 @@ +package metrics + +import ( + "errors" + "fmt" + + "github.com/prometheus/client_golang/prometheus" + + outcomeDomain "proxy-pool/internal/domain/outcome" +) + +// GatewayCollector exposes fixed-cardinality request-path outcome metrics. +// Proxy, route, destination, client and credential values are never labels. +type GatewayCollector struct { + outcomes *prometheus.CounterVec + dropped prometheus.Counter +} + +var _ outcomeDomain.MetricsObserver = (*GatewayCollector)(nil) + +func NewGatewayCollector(registerer prometheus.Registerer) (*GatewayCollector, error) { + if registerer == nil { + return nil, ErrInvalidDependencies + } + outcomes, err := registerCounterVec(registerer, prometheus.NewCounterVec(prometheus.CounterOpts{ + Namespace: "proxy_pool", Subsystem: "gateway", Name: "outcomes_total", + Help: "Number of Gateway proxy attempts by furthest completed stage and result.", + }, []string{"stage", "result"})) + if err != nil { + return nil, err + } + dropped, err := registerCounter(registerer, prometheus.NewCounter(prometheus.CounterOpts{ + Namespace: "proxy_pool", Subsystem: "gateway", Name: "outcome_queue_dropped_total", + Help: "Number of Gateway outcome observations dropped after the local queue was full.", + })) + if err != nil { + return nil, err + } + return &GatewayCollector{outcomes: outcomes, dropped: dropped}, nil +} + +func (collector *GatewayCollector) Observe(event outcomeDomain.Event) { + if collector == nil || collector.outcomes == nil || !validGatewayStage(event.Stage) { + return + } + result := "failure" + if event.Success { + result = "success" + } + collector.outcomes.WithLabelValues(gatewayStageLabel(event.Stage), result).Inc() +} + +func (collector *GatewayCollector) ObserveDropped() { + if collector == nil || collector.dropped == nil { + return + } + collector.dropped.Inc() +} + +func registerCounter(registerer prometheus.Registerer, candidate prometheus.Counter) (prometheus.Counter, error) { + if err := registerer.Register(candidate); err == nil { + return candidate, nil + } else { + var registered prometheus.AlreadyRegisteredError + if !errors.As(err, ®istered) { + return nil, fmt.Errorf("register counter: %w", err) + } + existing, ok := registered.ExistingCollector.(prometheus.Counter) + if !ok { + return nil, fmt.Errorf("register counter: existing collector has unexpected type") + } + return existing, nil + } +} + +func validGatewayStage(stage outcomeDomain.Stage) bool { + switch stage { + case outcomeDomain.StageDial, outcomeDomain.StageProxyHandshake, + outcomeDomain.StageResponseHeaders, outcomeDomain.StageTunnel: + return true + default: + return false + } +} + +func gatewayStageLabel(stage outcomeDomain.Stage) string { + switch stage { + case outcomeDomain.StageDial: + return "DIAL" + case outcomeDomain.StageProxyHandshake: + return "PROXY_HANDSHAKE" + case outcomeDomain.StageResponseHeaders: + return "RESPONSE_HEADERS" + case outcomeDomain.StageTunnel: + return "TUNNEL" + default: + return "" + } +} diff --git a/internal/platform/metrics/gateway_test.go b/internal/platform/metrics/gateway_test.go new file mode 100644 index 0000000..a5a5b0c --- /dev/null +++ b/internal/platform/metrics/gateway_test.go @@ -0,0 +1,44 @@ +package metrics + +import ( + "testing" + + "github.com/prometheus/client_golang/prometheus" + + outcomeDomain "proxy-pool/internal/domain/outcome" +) + +func TestGatewayCollectorRecordsOnlyFixedDimensions(t *testing.T) { + registry := prometheus.NewRegistry() + collector, err := NewGatewayCollector(registry) + if err != nil { + t.Fatalf("NewGatewayCollector() = %v", err) + } + collector.Observe(outcomeDomain.Event{Stage: outcomeDomain.StageDial, Success: true}) + collector.Observe(outcomeDomain.Event{Stage: outcomeDomain.StageTunnel, Success: false}) + collector.Observe(outcomeDomain.Event{Stage: outcomeDomain.StageUnspecified, Success: true}) + collector.ObserveDropped() + + assertMetricValue(t, registry, "proxy_pool_gateway_outcomes_total", map[string]string{ + "stage": "DIAL", "result": "success", + }, 1) + assertMetricValue(t, registry, "proxy_pool_gateway_outcomes_total", map[string]string{ + "stage": "TUNNEL", "result": "failure", + }, 1) + assertMetricValue(t, registry, "proxy_pool_gateway_outcome_queue_dropped_total", nil, 1) +} + +func TestNewGatewayCollectorReusesRegisteredCollectors(t *testing.T) { + registry := prometheus.NewRegistry() + first, err := NewGatewayCollector(registry) + if err != nil { + t.Fatalf("first NewGatewayCollector() = %v", err) + } + second, err := NewGatewayCollector(registry) + if err != nil { + t.Fatalf("second NewGatewayCollector() = %v", err) + } + first.ObserveDropped() + second.ObserveDropped() + assertMetricValue(t, registry, "proxy_pool_gateway_outcome_queue_dropped_total", nil, 2) +}