diff --git a/README.md b/README.md index 23b0080..6c79510 100644 --- a/README.md +++ b/README.md @@ -41,8 +41,10 @@ Proxy Pool 用 Controller 协调这些变化,并让 Gateway 数据面只消费 PostgreSQL 迁移与 Redis 活动池,并支持联动优雅停机。 - **Checker 指标**:Metrics 启用时暴露 Checker 任务下发与 Observation 接受/拒绝计数; 标签仅使用固定检查级别和结果,不记录 Proxy、IP、URL 或凭据。 -- **Gateway 指标**:Metrics 启用时暴露代理尝试的固定阶段成功/失败计数,以及本地 - Outcome 队列满后的丢弃计数;不记录 Proxy、路由、目标、客户端或凭据。 +- **Gateway 指标**:Metrics 启用时暴露代理尝试的固定阶段成功/失败计数、本地 + Outcome 队列满后的丢弃计数、HTTP/CONNECT 已接收请求数与在途数,以及活跃 + CONNECT 隧道数;协议标签仅有 HTTP 与 CONNECT,不记录 Proxy、路由、目标、 + 客户端或凭据。 - **Drain 指标**:Controller 暴露 `proxy_pool_controller_drain_candidates_total` 与 `proxy_pool_controller_drains_started_total`;`reason` 仅有 `unhealthy` 与 `upstream_disabled`,不包含 Proxy、Worker、Upstream、会话或地址。 @@ -127,7 +129,8 @@ flowchart LR Observation 状态归并;Provider 连续空结果的代次化自动 Sequential 切换、禁用候选过滤、 末端 `stop` 的 CAS 路由停用和 Snapshot 即时刷新。 - **部分完成**:Docker Compose/Kubernetes 运行时 mTLS Overlay。 -- **待完成**:CONNECT 长连接/Extract 压测场景、故障演练和代表性集群压测。 +- **待完成**:故障演练和代表性集群压测;现有 HTTP、CONNECT 长连接和 Extract + 场景只提供可复现的负载工具,不构成容量验证结论。 检查项数量不等于生产就绪度。静态部署清单与 protobuf descriptor 验证也不代表 端到端拓扑已经完成;`100,000 QPS` 仍只是待验证的集群设计目标。 diff --git a/docs/configuration/reference.md b/docs/configuration/reference.md index 7f0e4f8..e97ff21 100644 --- a/docs/configuration/reference.md +++ b/docs/configuration/reference.md @@ -693,6 +693,12 @@ Client、路由、目标 URL 或凭据标签。Capacity 从既有 Provider 库 结构化日志输出,包含组件和稳定错误类型,不输出错误原文;敏感属性、URL 用户信息和 查询 Secret 在写出前统一替换为 `[REDACTED]`。 +Gateway 还按每个 Worker 暴露 HTTP/CONNECT 请求总数、在途请求数和活跃 CONNECT +隧道数;协议标签固定为 HTTP、CONNECT。请求数在 Handler 接受请求时增加,在途数在 +全部拒绝、转发或隧道关闭后归零;活跃隧道只覆盖成功建立并开始 relay 的连接。这些 +指标是连接池、入口并发、文件描述符和长连接排空的本地观测,不携带 Proxy、路由、 +目标、Client 或凭据。 + ## 10. 启动前校验清单 1. `version` 必须为 `1`,未知字段拒绝。 diff --git a/docs/operations/runbook.md b/docs/operations/runbook.md index 602a4fe..3f1450e 100644 --- a/docs/operations/runbook.md +++ b/docs/operations/runbook.md @@ -138,6 +138,12 @@ Proxy ID、IP、Checker ID、目标 URL、Client ID 或凭据加入指标标签 Gateway 指标同样使用固定标签集: +每个 Gateway Worker 还暴露 requests_total、requests_in_flight 与 active_tunnels +三类本地连接指标。protocol 标签只允许 HTTP、CONNECT;请求指标覆盖从 Handler +接受请求到全部转发、拒绝或隧道关闭的生命周期。活跃隧道只在成功写出 CONNECT 200 +并开始 relay 后增加,关闭后立即减少。这些指标用于核对入口准入、连接池调优、 +文件描述符预算和长连接排空,不得按 Proxy、路由、目标、Client 或凭据拆分。 + - `proxy_pool_gateway_outcomes_total{stage,result}`:每次代理尝试在最远完成阶段的成功/失败数量。 - `proxy_pool_gateway_outcome_queue_dropped_total`:Outcome 本地有界队列已满后丢弃的观测数量。 diff --git a/docs/requirements/completion-audit.md b/docs/requirements/completion-audit.md index f13ed0b..5656635 100644 --- a/docs/requirements/completion-audit.md +++ b/docs/requirements/completion-audit.md @@ -113,7 +113,8 @@ Controller/Gateway 入口,完整 mTLS 运行时拓扑仍只有静态验证。 Extraction 与 Capacity 的低基数业务指标已闭环;Gateway Capacity 的动态降容保留 已有连接并在其排空前停止新增预留,非法 Reservation 终结汇总为固定枚举指标;Controller、Gateway、Checker 的 进程级致命错误现使用统一 JSON 脱敏日志出口,不输出错误原文。 -2. Gateway 的生产连接池调优与代表性流量压测。 +2. Gateway 已支持连接池、每 Host 连接上限、握手/空闲超时和隧道缓冲的本地配置; + 仍需代表性环境的流量压测。 3. Provider 分布式 singleflight/Leader、长期凭据回收和累计额度执行器。 4. Controller 的 PostgreSQL 连接池、迁移和 pgx Adapter 启动装配已完成; 公用 bootstrap 已通过 PostgreSQL 18 + Redis 8.2 双存储集成,Controller diff --git a/internal/gateway/bootstrap/bootstrap.go b/internal/gateway/bootstrap/bootstrap.go index 5b1e7be..847dbff 100644 --- a/internal/gateway/bootstrap/bootstrap.go +++ b/internal/gateway/bootstrap/bootstrap.go @@ -146,6 +146,7 @@ func newRuntime(ctx context.Context, configuration *config.Config, options Optio return nil, fmt.Errorf("build gateway target policy: %w", err) } var outcomeMetrics outcomeDomain.MetricsObserver + var requestMetrics server.RequestMetricsObserver storeOptions := snapshot.StoreOptions{} if configuration.Metrics.Enabled { collector, err := platformMetrics.NewGatewayCollector(prometheus.DefaultRegisterer) @@ -153,6 +154,7 @@ func newRuntime(ctx context.Context, configuration *config.Config, options Optio return nil, fmt.Errorf("build gateway outcome metrics: %w", err) } outcomeMetrics = collector + requestMetrics = collector storeOptions.CapacityInvariantObserver = collector } store := snapshot.NewStoreWithOptions(options.ClusterID, options.WorkerID, storeOptions) @@ -174,6 +176,7 @@ func newRuntime(ctx context.Context, configuration *config.Config, options Optio Dispatcher: dispatch.New(store), Transport: proxyTransport, Outcomes: outcomes, + Metrics: requestMetrics, }) if err != nil { proxyTransport.CloseIdleConnections() diff --git a/internal/gateway/server/handler.go b/internal/gateway/server/handler.go index 0900a4d..a6aef82 100644 --- a/internal/gateway/server/handler.go +++ b/internal/gateway/server/handler.go @@ -86,6 +86,15 @@ type OutcomeRecorder interface { Record(outcomeDomain.Event) } +// RequestMetricsObserver records local Gateway request lifecycle signals. +// Implementations must be concurrency-safe and must not block forwarding. +type RequestMetricsObserver interface { + ObserveRequestStarted(protocol string) + ObserveRequestFinished(protocol string) + ObserveTunnelOpened() + ObserveTunnelClosed() +} + type waitingDispatcher interface { AcquireWait(context.Context, dispatch.Request, time.Duration) (*dispatch.Lease, error) } @@ -99,6 +108,7 @@ type Dependencies struct { Dispatcher Dispatcher Transport ProxyTransport Outcomes OutcomeRecorder + Metrics RequestMetricsObserver } type Handler struct { @@ -109,6 +119,7 @@ type Handler struct { dispatcher Dispatcher transport ProxyTransport outcomes OutcomeRecorder + metrics RequestMetricsObserver sticky *stickySession buffers sync.Pool inFlight chan struct{} @@ -152,6 +163,7 @@ func New(config Config, dependencies Dependencies) (*Handler, error) { dispatcher: dependencies.Dispatcher, transport: dependencies.Transport, outcomes: dependencies.Outcomes, + metrics: dependencies.Metrics, sticky: sticky, tunnels: make(map[*activeTunnel]struct{}), shutdownDone: make(chan struct{}), @@ -169,6 +181,9 @@ func (handler *Handler) ServeHTTP(writer http.ResponseWriter, request *http.Requ return } defer handler.finishRequest() + protocol := requestMetricProtocol(request) + handler.observeRequestStarted(protocol) + defer handler.observeRequestFinished(protocol) if handler.inFlight != nil { select { case handler.inFlight <- struct{}{}: @@ -438,13 +453,41 @@ func (handler *Handler) registerTunnel(tunnel *activeTunnel) bool { return false } handler.tunnels[tunnel] = struct{}{} + if handler.metrics != nil { + handler.metrics.ObserveTunnelOpened() + } return true } func (handler *Handler) unregisterTunnel(tunnel *activeTunnel) { handler.tunnelMu.Lock() - delete(handler.tunnels, tunnel) + _, exists := handler.tunnels[tunnel] + if exists { + delete(handler.tunnels, tunnel) + } handler.tunnelMu.Unlock() + if exists && handler.metrics != nil { + handler.metrics.ObserveTunnelClosed() + } +} + +func (handler *Handler) observeRequestStarted(protocol string) { + if handler.metrics != nil { + handler.metrics.ObserveRequestStarted(protocol) + } +} + +func (handler *Handler) observeRequestFinished(protocol string) { + if handler.metrics != nil { + handler.metrics.ObserveRequestFinished(protocol) + } +} + +func requestMetricProtocol(request *http.Request) string { + if request != nil && request.Method == http.MethodConnect { + return "CONNECT" + } + return "HTTP" } func (handler *Handler) forwardHTTP( diff --git a/internal/gateway/server/handler_test.go b/internal/gateway/server/handler_test.go index f424579..22201a3 100644 --- a/internal/gateway/server/handler_test.go +++ b/internal/gateway/server/handler_test.go @@ -70,6 +70,53 @@ func TestHandlerRunsProtectionAndTargetPolicyBeforeRouting(t *testing.T) { } } +func TestHandlerReportsLocalRequestAndTunnelLifecycles(t *testing.T) { + t.Parallel() + + metrics := &recordingRequestMetrics{} + handler, err := New(Config{}, Dependencies{ + Targets: fakeTargets{}, + Router: RouteFunc(func(*http.Request) (dispatch.Request, error) { + return dispatch.Request{}, errors.New("route stopped") + }), + Dispatcher: DispatcherFunc(func(dispatch.Request) (*dispatch.Lease, error) { + t.Fatal("dispatcher must not run after route error") + return nil, nil + }), + Transport: &fakeTransport{}, + Metrics: metrics, + }) + if err != nil { + t.Fatalf("New() error = %v", err) + } + + handler.ServeHTTP(httptest.NewRecorder(), httptest.NewRequest(http.MethodGet, "http://example.test/resource", nil)) + connect := httptest.NewRequest(http.MethodConnect, "http://example.test", nil) + connect.Host = "example.test:443" + handler.ServeHTTP(httptest.NewRecorder(), connect) + tunnel := &activeTunnel{} + if !handler.registerTunnel(tunnel) { + t.Fatal("registerTunnel() = false") + } + handler.unregisterTunnel(tunnel) + + if got := metrics.started("HTTP"); got != 1 { + t.Fatalf("HTTP request starts = %d, want 1", got) + } + if got := metrics.finished("HTTP"); got != 1 { + t.Fatalf("HTTP request finishes = %d, want 1", got) + } + if got := metrics.started("CONNECT"); got != 1 { + t.Fatalf("CONNECT request starts = %d, want 1", got) + } + if got := metrics.finished("CONNECT"); got != 1 { + t.Fatalf("CONNECT request finishes = %d, want 1", got) + } + if metrics.opened != 1 || metrics.closed != 1 { + t.Fatalf("tunnel lifecycle = opened:%d closed:%d, want 1:1", metrics.opened, metrics.closed) + } +} + func TestHandlerRejectsGatewayRoutingOutsideCredentialPolicy(t *testing.T) { t.Parallel() @@ -795,6 +842,56 @@ type fakeTransport struct { relay func(context.Context, net.Conn, net.Conn) error } +type recordingRequestMetrics struct { + mu sync.Mutex + starts map[string]int + finishes map[string]int + opened int + closed int +} + +func (metrics *recordingRequestMetrics) ObserveRequestStarted(protocol string) { + metrics.mu.Lock() + defer metrics.mu.Unlock() + if metrics.starts == nil { + metrics.starts = make(map[string]int) + } + metrics.starts[protocol]++ +} + +func (metrics *recordingRequestMetrics) ObserveRequestFinished(protocol string) { + metrics.mu.Lock() + defer metrics.mu.Unlock() + if metrics.finishes == nil { + metrics.finishes = make(map[string]int) + } + metrics.finishes[protocol]++ +} + +func (metrics *recordingRequestMetrics) ObserveTunnelOpened() { + metrics.mu.Lock() + defer metrics.mu.Unlock() + metrics.opened++ +} + +func (metrics *recordingRequestMetrics) ObserveTunnelClosed() { + metrics.mu.Lock() + defer metrics.mu.Unlock() + metrics.closed++ +} + +func (metrics *recordingRequestMetrics) started(protocol string) int { + metrics.mu.Lock() + defer metrics.mu.Unlock() + return metrics.starts[protocol] +} + +func (metrics *recordingRequestMetrics) finished(protocol string) int { + metrics.mu.Lock() + defer metrics.mu.Unlock() + return metrics.finishes[protocol] +} + type outcomeRecorder struct { mu sync.Mutex events []outcomeDomain.Event diff --git a/internal/platform/metrics/gateway.go b/internal/platform/metrics/gateway.go index 2e3256f..018a27c 100644 --- a/internal/platform/metrics/gateway.go +++ b/internal/platform/metrics/gateway.go @@ -13,9 +13,14 @@ import ( // 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 - invariants *prometheus.CounterVec + outcomes *prometheus.CounterVec + dropped prometheus.Counter + invariants *prometheus.CounterVec + httpRequests prometheus.Counter + connectRequests prometheus.Counter + httpInFlight prometheus.Gauge + connectInFlight prometheus.Gauge + activeTunnels prometheus.Gauge } var ( @@ -48,7 +53,33 @@ func NewGatewayCollector(registerer prometheus.Registerer) (*GatewayCollector, e if err != nil { return nil, err } - return &GatewayCollector{outcomes: outcomes, dropped: dropped, invariants: invariants}, nil + requests, err := registerCounterVec(registerer, prometheus.NewCounterVec(prometheus.CounterOpts{ + Namespace: "proxy_pool", Subsystem: "gateway", Name: "requests_total", + Help: "Number of accepted Gateway requests by fixed proxy protocol.", + }, []string{"protocol"})) + if err != nil { + return nil, err + } + inFlight, err := registerGaugeVec(registerer, prometheus.NewGaugeVec(prometheus.GaugeOpts{ + Namespace: "proxy_pool", Subsystem: "gateway", Name: "requests_in_flight", + Help: "Number of accepted Gateway requests currently executing by fixed proxy protocol.", + }, []string{"protocol"})) + if err != nil { + return nil, err + } + activeTunnels, err := registerGauge(registerer, prometheus.NewGauge(prometheus.GaugeOpts{ + Namespace: "proxy_pool", Subsystem: "gateway", Name: "active_tunnels", + Help: "Number of established CONNECT tunnels currently relaying through this Gateway Worker.", + })) + if err != nil { + return nil, err + } + return &GatewayCollector{ + outcomes: outcomes, dropped: dropped, invariants: invariants, + httpRequests: requests.WithLabelValues("HTTP"), connectRequests: requests.WithLabelValues("CONNECT"), + httpInFlight: inFlight.WithLabelValues("HTTP"), connectInFlight: inFlight.WithLabelValues("CONNECT"), + activeTunnels: activeTunnels, + }, nil } func (collector *GatewayCollector) Observe(event outcomeDomain.Event) { @@ -76,6 +107,56 @@ func (collector *GatewayCollector) ObserveCapacityInvariant(event proxyDomain.Ca collector.invariants.WithLabelValues(string(event.Violation)).Inc() } +func (collector *GatewayCollector) ObserveRequestStarted(protocol string) { + if collector == nil { + return + } + switch protocol { + case "HTTP": + if collector.httpRequests != nil { + collector.httpRequests.Inc() + } + if collector.httpInFlight != nil { + collector.httpInFlight.Inc() + } + case "CONNECT": + if collector.connectRequests != nil { + collector.connectRequests.Inc() + } + if collector.connectInFlight != nil { + collector.connectInFlight.Inc() + } + } +} + +func (collector *GatewayCollector) ObserveRequestFinished(protocol string) { + if collector == nil { + return + } + switch protocol { + case "HTTP": + if collector.httpInFlight != nil { + collector.httpInFlight.Dec() + } + case "CONNECT": + if collector.connectInFlight != nil { + collector.connectInFlight.Dec() + } + } +} + +func (collector *GatewayCollector) ObserveTunnelOpened() { + if collector != nil && collector.activeTunnels != nil { + collector.activeTunnels.Inc() + } +} + +func (collector *GatewayCollector) ObserveTunnelClosed() { + if collector != nil && collector.activeTunnels != nil { + collector.activeTunnels.Dec() + } +} + func registerCounter(registerer prometheus.Registerer, candidate prometheus.Counter) (prometheus.Counter, error) { if err := registerer.Register(candidate); err == nil { return candidate, nil @@ -92,6 +173,38 @@ func registerCounter(registerer prometheus.Registerer, candidate prometheus.Coun } } +func registerGauge(registerer prometheus.Registerer, candidate prometheus.Gauge) (prometheus.Gauge, 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 gauge: %w", err) + } + existing, ok := registered.ExistingCollector.(prometheus.Gauge) + if !ok { + return nil, fmt.Errorf("register gauge: existing collector has unexpected type") + } + return existing, nil + } +} + +func registerGaugeVec(registerer prometheus.Registerer, candidate *prometheus.GaugeVec) (*prometheus.GaugeVec, 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 gauge vector: %w", err) + } + existing, ok := registered.ExistingCollector.(*prometheus.GaugeVec) + if !ok { + return nil, fmt.Errorf("register gauge vector: existing collector has unexpected type") + } + return existing, nil + } +} + func validGatewayStage(stage outcomeDomain.Stage) bool { switch stage { case outcomeDomain.StageDial, outcomeDomain.StageProxyHandshake, diff --git a/internal/platform/metrics/gateway_test.go b/internal/platform/metrics/gateway_test.go index 8d9f84f..2719080 100644 --- a/internal/platform/metrics/gateway_test.go +++ b/internal/platform/metrics/gateway_test.go @@ -23,6 +23,13 @@ func TestGatewayCollectorRecordsOnlyFixedDimensions(t *testing.T) { Violation: proxyDomain.CapacityInvariantReleaseFinished, }) collector.ObserveCapacityInvariant(proxyDomain.CapacityInvariant{Violation: "unknown"}) + collector.ObserveRequestStarted("HTTP") + collector.ObserveRequestStarted("CONNECT") + collector.ObserveRequestFinished("HTTP") + collector.ObserveRequestFinished("CONNECT") + collector.ObserveTunnelOpened() + collector.ObserveTunnelClosed() + collector.ObserveRequestStarted("unknown") assertMetricValue(t, registry, "proxy_pool_gateway_outcomes_total", map[string]string{ "stage": "DIAL", "result": "success", @@ -34,6 +41,19 @@ func TestGatewayCollectorRecordsOnlyFixedDimensions(t *testing.T) { assertMetricValue(t, registry, "proxy_pool_gateway_capacity_invariant_violations_total", map[string]string{ "operation": "release_finished", }, 1) + assertMetricValue(t, registry, "proxy_pool_gateway_requests_total", map[string]string{ + "protocol": "HTTP", + }, 1) + assertMetricValue(t, registry, "proxy_pool_gateway_requests_total", map[string]string{ + "protocol": "CONNECT", + }, 1) + assertMetricValue(t, registry, "proxy_pool_gateway_requests_in_flight", map[string]string{ + "protocol": "HTTP", + }, 0) + assertMetricValue(t, registry, "proxy_pool_gateway_requests_in_flight", map[string]string{ + "protocol": "CONNECT", + }, 0) + assertMetricValue(t, registry, "proxy_pool_gateway_active_tunnels", nil, 0) } func TestNewGatewayCollectorReusesRegisteredCollectors(t *testing.T) {