feat: observe gateway connection lifecycles
This commit is contained in:
parent
bda9cc03df
commit
26572da97e
@ -41,8 +41,10 @@ Proxy Pool 用 Controller 协调这些变化,并让 Gateway 数据面只消费
|
|||||||
PostgreSQL 迁移与 Redis 活动池,并支持联动优雅停机。
|
PostgreSQL 迁移与 Redis 活动池,并支持联动优雅停机。
|
||||||
- **Checker 指标**:Metrics 启用时暴露 Checker 任务下发与 Observation 接受/拒绝计数;
|
- **Checker 指标**:Metrics 启用时暴露 Checker 任务下发与 Observation 接受/拒绝计数;
|
||||||
标签仅使用固定检查级别和结果,不记录 Proxy、IP、URL 或凭据。
|
标签仅使用固定检查级别和结果,不记录 Proxy、IP、URL 或凭据。
|
||||||
- **Gateway 指标**:Metrics 启用时暴露代理尝试的固定阶段成功/失败计数,以及本地
|
- **Gateway 指标**:Metrics 启用时暴露代理尝试的固定阶段成功/失败计数、本地
|
||||||
Outcome 队列满后的丢弃计数;不记录 Proxy、路由、目标、客户端或凭据。
|
Outcome 队列满后的丢弃计数、HTTP/CONNECT 已接收请求数与在途数,以及活跃
|
||||||
|
CONNECT 隧道数;协议标签仅有 HTTP 与 CONNECT,不记录 Proxy、路由、目标、
|
||||||
|
客户端或凭据。
|
||||||
- **Drain 指标**:Controller 暴露 `proxy_pool_controller_drain_candidates_total` 与
|
- **Drain 指标**:Controller 暴露 `proxy_pool_controller_drain_candidates_total` 与
|
||||||
`proxy_pool_controller_drains_started_total`;`reason` 仅有 `unhealthy` 与
|
`proxy_pool_controller_drains_started_total`;`reason` 仅有 `unhealthy` 与
|
||||||
`upstream_disabled`,不包含 Proxy、Worker、Upstream、会话或地址。
|
`upstream_disabled`,不包含 Proxy、Worker、Upstream、会话或地址。
|
||||||
@ -127,7 +129,8 @@ flowchart LR
|
|||||||
Observation 状态归并;Provider 连续空结果的代次化自动 Sequential 切换、禁用候选过滤、
|
Observation 状态归并;Provider 连续空结果的代次化自动 Sequential 切换、禁用候选过滤、
|
||||||
末端 `stop` 的 CAS 路由停用和 Snapshot 即时刷新。
|
末端 `stop` 的 CAS 路由停用和 Snapshot 即时刷新。
|
||||||
- **部分完成**:Docker Compose/Kubernetes 运行时 mTLS Overlay。
|
- **部分完成**:Docker Compose/Kubernetes 运行时 mTLS Overlay。
|
||||||
- **待完成**:CONNECT 长连接/Extract 压测场景、故障演练和代表性集群压测。
|
- **待完成**:故障演练和代表性集群压测;现有 HTTP、CONNECT 长连接和 Extract
|
||||||
|
场景只提供可复现的负载工具,不构成容量验证结论。
|
||||||
|
|
||||||
检查项数量不等于生产就绪度。静态部署清单与 protobuf descriptor 验证也不代表
|
检查项数量不等于生产就绪度。静态部署清单与 protobuf descriptor 验证也不代表
|
||||||
端到端拓扑已经完成;`100,000 QPS` 仍只是待验证的集群设计目标。
|
端到端拓扑已经完成;`100,000 QPS` 仍只是待验证的集群设计目标。
|
||||||
|
|||||||
@ -693,6 +693,12 @@ Client、路由、目标 URL 或凭据标签。Capacity 从既有 Provider 库
|
|||||||
结构化日志输出,包含组件和稳定错误类型,不输出错误原文;敏感属性、URL 用户信息和
|
结构化日志输出,包含组件和稳定错误类型,不输出错误原文;敏感属性、URL 用户信息和
|
||||||
查询 Secret 在写出前统一替换为 `[REDACTED]`。
|
查询 Secret 在写出前统一替换为 `[REDACTED]`。
|
||||||
|
|
||||||
|
Gateway 还按每个 Worker 暴露 HTTP/CONNECT 请求总数、在途请求数和活跃 CONNECT
|
||||||
|
隧道数;协议标签固定为 HTTP、CONNECT。请求数在 Handler 接受请求时增加,在途数在
|
||||||
|
全部拒绝、转发或隧道关闭后归零;活跃隧道只覆盖成功建立并开始 relay 的连接。这些
|
||||||
|
指标是连接池、入口并发、文件描述符和长连接排空的本地观测,不携带 Proxy、路由、
|
||||||
|
目标、Client 或凭据。
|
||||||
|
|
||||||
## 10. 启动前校验清单
|
## 10. 启动前校验清单
|
||||||
|
|
||||||
1. `version` 必须为 `1`,未知字段拒绝。
|
1. `version` 必须为 `1`,未知字段拒绝。
|
||||||
|
|||||||
@ -138,6 +138,12 @@ Proxy ID、IP、Checker ID、目标 URL、Client ID 或凭据加入指标标签
|
|||||||
|
|
||||||
Gateway 指标同样使用固定标签集:
|
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_outcomes_total{stage,result}`:每次代理尝试在最远完成阶段的成功/失败数量。
|
||||||
- `proxy_pool_gateway_outcome_queue_dropped_total`:Outcome 本地有界队列已满后丢弃的观测数量。
|
- `proxy_pool_gateway_outcome_queue_dropped_total`:Outcome 本地有界队列已满后丢弃的观测数量。
|
||||||
|
|
||||||
|
|||||||
@ -113,7 +113,8 @@ Controller/Gateway 入口,完整 mTLS 运行时拓扑仍只有静态验证。
|
|||||||
Extraction 与 Capacity 的低基数业务指标已闭环;Gateway Capacity 的动态降容保留
|
Extraction 与 Capacity 的低基数业务指标已闭环;Gateway Capacity 的动态降容保留
|
||||||
已有连接并在其排空前停止新增预留,非法 Reservation 终结汇总为固定枚举指标;Controller、Gateway、Checker 的
|
已有连接并在其排空前停止新增预留,非法 Reservation 终结汇总为固定枚举指标;Controller、Gateway、Checker 的
|
||||||
进程级致命错误现使用统一 JSON 脱敏日志出口,不输出错误原文。
|
进程级致命错误现使用统一 JSON 脱敏日志出口,不输出错误原文。
|
||||||
2. Gateway 的生产连接池调优与代表性流量压测。
|
2. Gateway 已支持连接池、每 Host 连接上限、握手/空闲超时和隧道缓冲的本地配置;
|
||||||
|
仍需代表性环境的流量压测。
|
||||||
3. Provider 分布式 singleflight/Leader、长期凭据回收和累计额度执行器。
|
3. Provider 分布式 singleflight/Leader、长期凭据回收和累计额度执行器。
|
||||||
4. Controller 的 PostgreSQL 连接池、迁移和 pgx Adapter 启动装配已完成;
|
4. Controller 的 PostgreSQL 连接池、迁移和 pgx Adapter 启动装配已完成;
|
||||||
公用 bootstrap 已通过 PostgreSQL 18 + Redis 8.2 双存储集成,Controller
|
公用 bootstrap 已通过 PostgreSQL 18 + Redis 8.2 双存储集成,Controller
|
||||||
|
|||||||
@ -146,6 +146,7 @@ func newRuntime(ctx context.Context, configuration *config.Config, options Optio
|
|||||||
return nil, fmt.Errorf("build gateway target policy: %w", err)
|
return nil, fmt.Errorf("build gateway target policy: %w", err)
|
||||||
}
|
}
|
||||||
var outcomeMetrics outcomeDomain.MetricsObserver
|
var outcomeMetrics outcomeDomain.MetricsObserver
|
||||||
|
var requestMetrics server.RequestMetricsObserver
|
||||||
storeOptions := snapshot.StoreOptions{}
|
storeOptions := snapshot.StoreOptions{}
|
||||||
if configuration.Metrics.Enabled {
|
if configuration.Metrics.Enabled {
|
||||||
collector, err := platformMetrics.NewGatewayCollector(prometheus.DefaultRegisterer)
|
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)
|
return nil, fmt.Errorf("build gateway outcome metrics: %w", err)
|
||||||
}
|
}
|
||||||
outcomeMetrics = collector
|
outcomeMetrics = collector
|
||||||
|
requestMetrics = collector
|
||||||
storeOptions.CapacityInvariantObserver = collector
|
storeOptions.CapacityInvariantObserver = collector
|
||||||
}
|
}
|
||||||
store := snapshot.NewStoreWithOptions(options.ClusterID, options.WorkerID, storeOptions)
|
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),
|
Dispatcher: dispatch.New(store),
|
||||||
Transport: proxyTransport,
|
Transport: proxyTransport,
|
||||||
Outcomes: outcomes,
|
Outcomes: outcomes,
|
||||||
|
Metrics: requestMetrics,
|
||||||
})
|
})
|
||||||
if err != nil {
|
if err != nil {
|
||||||
proxyTransport.CloseIdleConnections()
|
proxyTransport.CloseIdleConnections()
|
||||||
|
|||||||
@ -86,6 +86,15 @@ type OutcomeRecorder interface {
|
|||||||
Record(outcomeDomain.Event)
|
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 {
|
type waitingDispatcher interface {
|
||||||
AcquireWait(context.Context, dispatch.Request, time.Duration) (*dispatch.Lease, error)
|
AcquireWait(context.Context, dispatch.Request, time.Duration) (*dispatch.Lease, error)
|
||||||
}
|
}
|
||||||
@ -99,6 +108,7 @@ type Dependencies struct {
|
|||||||
Dispatcher Dispatcher
|
Dispatcher Dispatcher
|
||||||
Transport ProxyTransport
|
Transport ProxyTransport
|
||||||
Outcomes OutcomeRecorder
|
Outcomes OutcomeRecorder
|
||||||
|
Metrics RequestMetricsObserver
|
||||||
}
|
}
|
||||||
|
|
||||||
type Handler struct {
|
type Handler struct {
|
||||||
@ -109,6 +119,7 @@ type Handler struct {
|
|||||||
dispatcher Dispatcher
|
dispatcher Dispatcher
|
||||||
transport ProxyTransport
|
transport ProxyTransport
|
||||||
outcomes OutcomeRecorder
|
outcomes OutcomeRecorder
|
||||||
|
metrics RequestMetricsObserver
|
||||||
sticky *stickySession
|
sticky *stickySession
|
||||||
buffers sync.Pool
|
buffers sync.Pool
|
||||||
inFlight chan struct{}
|
inFlight chan struct{}
|
||||||
@ -152,6 +163,7 @@ func New(config Config, dependencies Dependencies) (*Handler, error) {
|
|||||||
dispatcher: dependencies.Dispatcher,
|
dispatcher: dependencies.Dispatcher,
|
||||||
transport: dependencies.Transport,
|
transport: dependencies.Transport,
|
||||||
outcomes: dependencies.Outcomes,
|
outcomes: dependencies.Outcomes,
|
||||||
|
metrics: dependencies.Metrics,
|
||||||
sticky: sticky,
|
sticky: sticky,
|
||||||
tunnels: make(map[*activeTunnel]struct{}),
|
tunnels: make(map[*activeTunnel]struct{}),
|
||||||
shutdownDone: make(chan struct{}),
|
shutdownDone: make(chan struct{}),
|
||||||
@ -169,6 +181,9 @@ func (handler *Handler) ServeHTTP(writer http.ResponseWriter, request *http.Requ
|
|||||||
return
|
return
|
||||||
}
|
}
|
||||||
defer handler.finishRequest()
|
defer handler.finishRequest()
|
||||||
|
protocol := requestMetricProtocol(request)
|
||||||
|
handler.observeRequestStarted(protocol)
|
||||||
|
defer handler.observeRequestFinished(protocol)
|
||||||
if handler.inFlight != nil {
|
if handler.inFlight != nil {
|
||||||
select {
|
select {
|
||||||
case handler.inFlight <- struct{}{}:
|
case handler.inFlight <- struct{}{}:
|
||||||
@ -438,13 +453,41 @@ func (handler *Handler) registerTunnel(tunnel *activeTunnel) bool {
|
|||||||
return false
|
return false
|
||||||
}
|
}
|
||||||
handler.tunnels[tunnel] = struct{}{}
|
handler.tunnels[tunnel] = struct{}{}
|
||||||
|
if handler.metrics != nil {
|
||||||
|
handler.metrics.ObserveTunnelOpened()
|
||||||
|
}
|
||||||
return true
|
return true
|
||||||
}
|
}
|
||||||
|
|
||||||
func (handler *Handler) unregisterTunnel(tunnel *activeTunnel) {
|
func (handler *Handler) unregisterTunnel(tunnel *activeTunnel) {
|
||||||
handler.tunnelMu.Lock()
|
handler.tunnelMu.Lock()
|
||||||
|
_, exists := handler.tunnels[tunnel]
|
||||||
|
if exists {
|
||||||
delete(handler.tunnels, tunnel)
|
delete(handler.tunnels, tunnel)
|
||||||
|
}
|
||||||
handler.tunnelMu.Unlock()
|
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(
|
func (handler *Handler) forwardHTTP(
|
||||||
|
|||||||
@ -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) {
|
func TestHandlerRejectsGatewayRoutingOutsideCredentialPolicy(t *testing.T) {
|
||||||
t.Parallel()
|
t.Parallel()
|
||||||
|
|
||||||
@ -795,6 +842,56 @@ type fakeTransport struct {
|
|||||||
relay func(context.Context, net.Conn, net.Conn) error
|
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 {
|
type outcomeRecorder struct {
|
||||||
mu sync.Mutex
|
mu sync.Mutex
|
||||||
events []outcomeDomain.Event
|
events []outcomeDomain.Event
|
||||||
|
|||||||
@ -16,6 +16,11 @@ type GatewayCollector struct {
|
|||||||
outcomes *prometheus.CounterVec
|
outcomes *prometheus.CounterVec
|
||||||
dropped prometheus.Counter
|
dropped prometheus.Counter
|
||||||
invariants *prometheus.CounterVec
|
invariants *prometheus.CounterVec
|
||||||
|
httpRequests prometheus.Counter
|
||||||
|
connectRequests prometheus.Counter
|
||||||
|
httpInFlight prometheus.Gauge
|
||||||
|
connectInFlight prometheus.Gauge
|
||||||
|
activeTunnels prometheus.Gauge
|
||||||
}
|
}
|
||||||
|
|
||||||
var (
|
var (
|
||||||
@ -48,7 +53,33 @@ func NewGatewayCollector(registerer prometheus.Registerer) (*GatewayCollector, e
|
|||||||
if err != nil {
|
if err != nil {
|
||||||
return nil, err
|
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) {
|
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()
|
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) {
|
func registerCounter(registerer prometheus.Registerer, candidate prometheus.Counter) (prometheus.Counter, error) {
|
||||||
if err := registerer.Register(candidate); err == nil {
|
if err := registerer.Register(candidate); err == nil {
|
||||||
return candidate, 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 {
|
func validGatewayStage(stage outcomeDomain.Stage) bool {
|
||||||
switch stage {
|
switch stage {
|
||||||
case outcomeDomain.StageDial, outcomeDomain.StageProxyHandshake,
|
case outcomeDomain.StageDial, outcomeDomain.StageProxyHandshake,
|
||||||
|
|||||||
@ -23,6 +23,13 @@ func TestGatewayCollectorRecordsOnlyFixedDimensions(t *testing.T) {
|
|||||||
Violation: proxyDomain.CapacityInvariantReleaseFinished,
|
Violation: proxyDomain.CapacityInvariantReleaseFinished,
|
||||||
})
|
})
|
||||||
collector.ObserveCapacityInvariant(proxyDomain.CapacityInvariant{Violation: "unknown"})
|
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{
|
assertMetricValue(t, registry, "proxy_pool_gateway_outcomes_total", map[string]string{
|
||||||
"stage": "DIAL", "result": "success",
|
"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{
|
assertMetricValue(t, registry, "proxy_pool_gateway_capacity_invariant_violations_total", map[string]string{
|
||||||
"operation": "release_finished",
|
"operation": "release_finished",
|
||||||
}, 1)
|
}, 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) {
|
func TestNewGatewayCollectorReusesRegisteredCollectors(t *testing.T) {
|
||||||
|
|||||||
Loading…
Reference in New Issue
Block a user