From b35092ccdd149cc246b5d9302e4d5a93e007c9a6 Mon Sep 17 00:00:00 2001 From: youfak Date: Sun, 2 Aug 2026 09:24:46 +0800 Subject: [PATCH] feat: expose checker task metrics --- README.md | 2 + docs/configuration/reference.md | 7 +- docs/development/implementation-plan.md | 6 ++ docs/operations/runbook.md | 8 ++ internal/controller/bootstrap/bootstrap.go | 8 ++ internal/controller/health/grpc_handler.go | 10 +++ .../controller/health/grpc_handler_test.go | 45 ++++++++++- internal/domain/health/metrics.go | 16 ++++ internal/platform/metrics/checker.go | 80 +++++++++++++++++++ internal/platform/metrics/checker_test.go | 79 ++++++++++++++++++ 10 files changed, 258 insertions(+), 3 deletions(-) create mode 100644 internal/domain/health/metrics.go create mode 100644 internal/platform/metrics/checker.go create mode 100644 internal/platform/metrics/checker_test.go diff --git a/README.md b/README.md index 7d0505d..85f10a1 100644 --- a/README.md +++ b/README.md @@ -39,6 +39,8 @@ Proxy Pool 用 Controller 协调这些变化,并让 Gateway 数据面只消费 库存预留,以及 `AVAILABLE -> EXTRACTED` 的一次性独占提取。 - **Controller 入口**:`proxy-controller` 已装配 Distribution、Admin、Metrics、 PostgreSQL 迁移与 Redis 活动池,并支持联动优雅停机。 +- **Checker 指标**:Metrics 启用时暴露 Checker 任务下发与 Observation 接受/拒绝计数; + 标签仅使用固定检查级别和结果,不记录 Proxy、IP、URL 或凭据。 - **PostgreSQL 管理面**:持久化配置版本、Upstream/Routing 管理状态、Admin 审计与 Outbox;不保存 Proxy 明细或逐次提取记录。 - **Gateway 组件**:HTTP 正向代理、HTTPS CONNECT、双向 Tunnel、重试、超时、 diff --git a/docs/configuration/reference.md b/docs/configuration/reference.md index 95a5a38..2852dda 100644 --- a/docs/configuration/reference.md +++ b/docs/configuration/reference.md @@ -510,8 +510,11 @@ Runtime;指纹已更新但本地源尚未同步时停止旧 Provider,待源 Metrics 启用时 `listen` 必须是合法 `host:port`。该入口固定提供 `/livez`、 `/readyz` 和 `/metrics`,不复用 Distribution/Admin 的认证边界;外部访问必须由 -网络策略限制。当前 `/metrics` 已包含 Go/进程基础指标,Provider、提取和容量等 -业务指标仍在后续实施范围。Distribution 启用时 `/readyz` 只以 Redis 活动池为 +网络策略限制。当前 `/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 活动池为 服务流量门槛,PostgreSQL 故障由 Admin 接口独立报告。Metrics 开关或监听地址 变更需要重启 Controller。 diff --git a/docs/development/implementation-plan.md b/docs/development/implementation-plan.md index b36aa8f..5eb8e54 100644 --- a/docs/development/implementation-plan.md +++ b/docs/development/implementation-plan.md @@ -332,6 +332,12 @@ Profile 在启用 Routing 与 Upstream 的组合上才进入调度。 - [ ] Document backup, recovery, rollout, rollback, capacity, kernel, file descriptor, NAT/conntrack, and incident runbooks. +当前进度(2026-08-02):已接入 Checker 任务流指标 +`proxy_pool_checker_tasks_dispatched_total{level}` 与 +`proxy_pool_checker_observations_total{level,result}`;标签值仅允许固定的 +`BASIC`、`EGRESS`、`TARGET` 级别及 `accepted`、`rejected` 结果。Provider、提取、 +容量指标与密钥安全的结构化日志仍待补齐。 + ## Task 14: Documentation, Examples, and Diagrams **Files:** `docs/**`, `examples/**`, `diagrams/**` diff --git a/docs/operations/runbook.md b/docs/operations/runbook.md index bb963ca..8889192 100644 --- a/docs/operations/runbook.md +++ b/docs/operations/runbook.md @@ -128,6 +128,14 @@ kubectl -n proxy-pool rollout status deployment/proxy-gateway --timeout=10m 探针不得执行 Provider 请求或完整数据库扫描。 +Checker 指标使用固定标签集: + +- `proxy_pool_checker_tasks_dispatched_total{level}`:成功写入 Checker 任务流后的任务数。 +- `proxy_pool_checker_observations_total{level,result}`:已解码的 Observation 被 Controller 接受或拒绝的数量。 + +`level` 仅为 BASIC、EGRESS、TARGET,`result` 仅为 accepted、rejected。不得将 +Proxy ID、IP、Checker ID、目标 URL、Client ID 或凭据加入指标标签。 + ### 3.4 发布顺序与兼容性 1. 先做向前兼容数据库迁移。 diff --git a/internal/controller/bootstrap/bootstrap.go b/internal/controller/bootstrap/bootstrap.go index 9aede1a..4d58aa9 100644 --- a/internal/controller/bootstrap/bootstrap.go +++ b/internal/controller/bootstrap/bootstrap.go @@ -23,6 +23,7 @@ import ( "proxy-pool/internal/controller/worker" "proxy-pool/internal/domain/activitypool" extractionDomain "proxy-pool/internal/domain/extraction" + healthDomain "proxy-pool/internal/domain/health" ownershipDomain "proxy-pool/internal/domain/ownership" "proxy-pool/internal/domain/workerruntime" "proxy-pool/internal/platform/admission" @@ -219,10 +220,16 @@ func runWithWorkerFactory( } dependencies.AdminService = service } + var checkerMetrics healthDomain.TaskMetricsObserver if loaded.Value.Metrics.Enabled { if nilInterface(opened.metricsReadiness) { return errors.Join(ErrStartup, ErrInvalidOptions) } + collector, collectorErr := platformMetrics.NewCheckerCollector(prometheus.DefaultRegisterer) + if collectorErr != nil { + return fmt.Errorf("%w: build Checker metrics: %w", ErrStartup, collectorErr) + } + checkerMetrics = collector handler, handlerErr := platformMetrics.NewHandler(platformMetrics.Dependencies{ Gatherer: prometheus.DefaultGatherer, Readiness: opened.metricsReadiness, }) @@ -289,6 +296,7 @@ func runWithWorkerFactory( return fmt.Errorf("%w: build Checker identity authorizer: %w", ErrStartup, identityErr) } checkerOptions := controllerHealth.DefaultGRPCHandlerOptions() + checkerOptions.Metrics = checkerMetrics if tasks, ok := opened.activity.(healthTaskRuntime); ok && !nilInterface(tasks) { checkerOptions.TaskBroker = tasks } diff --git a/internal/controller/health/grpc_handler.go b/internal/controller/health/grpc_handler.go index c079366..36af9e8 100644 --- a/internal/controller/health/grpc_handler.go +++ b/internal/controller/health/grpc_handler.go @@ -31,6 +31,7 @@ type GRPCHandlerOptions struct { MaxObservationsPerBatch int MaxTasksPerClaim int TaskBroker TaskBroker + Metrics healthDomain.TaskMetricsObserver Now func() time.Time } @@ -122,6 +123,9 @@ func (handler *GRPCHandler) StreamCheckTasks( if sendErr := stream.Send(wire); sendErr != nil { return sendErr } + if handler.options.Metrics != nil { + handler.options.Metrics.ObserveTaskDispatch(task.Level, 1) + } } return nil } @@ -160,10 +164,16 @@ func (handler *GRPCHandler) ReportObservations( } if err == nil { response.Accepted++ + if handler.options.Metrics != nil { + handler.options.Metrics.ObserveObservation(observation.Level, healthDomain.ObservationMetricAccepted) + } continue } if rejectedObservationError(err) { response.Rejected++ + if handler.options.Metrics != nil && observation.Level != "" { + handler.options.Metrics.ObserveObservation(observation.Level, healthDomain.ObservationMetricRejected) + } continue } return nil, healthGRPCError(err) diff --git a/internal/controller/health/grpc_handler_test.go b/internal/controller/health/grpc_handler_test.go index f4caa9f..aa75819 100644 --- a/internal/controller/health/grpc_handler_test.go +++ b/internal/controller/health/grpc_handler_test.go @@ -5,6 +5,7 @@ import ( "errors" "io" "net" + "sync" "testing" "time" @@ -127,8 +128,9 @@ func TestGRPCHandlerStreamsLeasedTasksAndFencesReportedFacts(t *testing.T) { if err != nil { t.Fatalf("NewReducer(): %v", err) } + metrics := &recordingTaskMetrics{} handler, err := NewGRPCHandler(reducer, &recordingCheckerIdentity{}, GRPCHandlerOptions{ - MaxObservationsPerBatch: 2, MaxTasksPerClaim: 2, TaskBroker: broker, Now: func() time.Time { return now }, + MaxObservationsPerBatch: 2, MaxTasksPerClaim: 2, TaskBroker: broker, Metrics: metrics, Now: func() time.Time { return now }, }) if err != nil { t.Fatalf("NewGRPCHandler(): %v", err) @@ -171,6 +173,11 @@ func TestGRPCHandlerStreamsLeasedTasksAndFencesReportedFacts(t *testing.T) { if err != nil || rejected.GetRejected() != 1 || len(global.commands) != 1 { t.Fatalf("ReportObservations(other checker) = (%+v, %v), global=%d", rejected, err, len(global.commands)) } + if metrics.dispatches(healthDomain.LevelBasic) != 1 || + metrics.observations(healthDomain.LevelBasic, healthDomain.ObservationMetricAccepted) != 1 || + metrics.observations(healthDomain.LevelBasic, healthDomain.ObservationMetricRejected) != 1 { + t.Fatalf("Checker metrics = %+v", metrics) + } } func TestGRPCHandlerStreamsEgressProbeURLAndAcceptsGlobalFact(t *testing.T) { @@ -265,6 +272,42 @@ type recordingCheckerIdentity struct { err error } +type recordingTaskMetrics struct { + mu sync.Mutex + dispatch map[healthDomain.Level]int + observation map[string]int +} + +func (metrics *recordingTaskMetrics) ObserveTaskDispatch(level healthDomain.Level, count int) { + metrics.mu.Lock() + defer metrics.mu.Unlock() + if metrics.dispatch == nil { + metrics.dispatch = make(map[healthDomain.Level]int) + } + metrics.dispatch[level] += count +} + +func (metrics *recordingTaskMetrics) ObserveObservation(level healthDomain.Level, result healthDomain.ObservationMetricResult) { + metrics.mu.Lock() + defer metrics.mu.Unlock() + if metrics.observation == nil { + metrics.observation = make(map[string]int) + } + metrics.observation[string(level)+"\x00"+string(result)]++ +} + +func (metrics *recordingTaskMetrics) dispatches(level healthDomain.Level) int { + metrics.mu.Lock() + defer metrics.mu.Unlock() + return metrics.dispatch[level] +} + +func (metrics *recordingTaskMetrics) observations(level healthDomain.Level, result healthDomain.ObservationMetricResult) int { + metrics.mu.Lock() + defer metrics.mu.Unlock() + return metrics.observation[string(level)+"\x00"+string(result)] +} + func (identity *recordingCheckerIdentity) AuthorizeChecker(_ context.Context, checkerID string) error { identity.checkerID = checkerID return identity.err diff --git a/internal/domain/health/metrics.go b/internal/domain/health/metrics.go new file mode 100644 index 0000000..746c10b --- /dev/null +++ b/internal/domain/health/metrics.go @@ -0,0 +1,16 @@ +package health + +// TaskMetricsObserver accepts only the low-cardinality dimensions emitted by +// Checker task processing. Implementations must not add task, proxy, checker, +// endpoint, target URL, or credential values as metric labels. +type TaskMetricsObserver interface { + ObserveTaskDispatch(Level, int) + ObserveObservation(Level, ObservationMetricResult) +} + +type ObservationMetricResult string + +const ( + ObservationMetricAccepted ObservationMetricResult = "accepted" + ObservationMetricRejected ObservationMetricResult = "rejected" +) diff --git a/internal/platform/metrics/checker.go b/internal/platform/metrics/checker.go new file mode 100644 index 0000000..67fee3a --- /dev/null +++ b/internal/platform/metrics/checker.go @@ -0,0 +1,80 @@ +package metrics + +import ( + "errors" + "fmt" + + "github.com/prometheus/client_golang/prometheus" + + healthDomain "proxy-pool/internal/domain/health" +) + +// CheckerCollector exposes bounded Checker task-flow counters. The only +// labels are the three fixed health levels and accepted/rejected outcome. +type CheckerCollector struct { + tasksDispatched *prometheus.CounterVec + observations *prometheus.CounterVec +} + +var _ healthDomain.TaskMetricsObserver = (*CheckerCollector)(nil) + +func NewCheckerCollector(registerer prometheus.Registerer) (*CheckerCollector, error) { + if registerer == nil { + return nil, ErrInvalidDependencies + } + tasks, err := registerCounterVec(registerer, prometheus.NewCounterVec(prometheus.CounterOpts{ + Namespace: "proxy_pool", Subsystem: "checker", Name: "tasks_dispatched_total", + Help: "Number of checker tasks dispatched after a successful task-stream send.", + }, []string{"level"})) + if err != nil { + return nil, err + } + observations, err := registerCounterVec(registerer, prometheus.NewCounterVec(prometheus.CounterOpts{ + Namespace: "proxy_pool", Subsystem: "checker", Name: "observations_total", + Help: "Number of checker observations accepted or rejected after decoding a known level.", + }, []string{"level", "result"})) + if err != nil { + return nil, err + } + return &CheckerCollector{tasksDispatched: tasks, observations: observations}, nil +} + +func (collector *CheckerCollector) ObserveTaskDispatch(level healthDomain.Level, count int) { + if collector == nil || collector.tasksDispatched == nil || count <= 0 || !validCheckerLevel(level) { + return + } + collector.tasksDispatched.WithLabelValues(string(level)).Add(float64(count)) +} + +func (collector *CheckerCollector) ObserveObservation(level healthDomain.Level, result healthDomain.ObservationMetricResult) { + if collector == nil || collector.observations == nil || !validCheckerLevel(level) || + (result != healthDomain.ObservationMetricAccepted && result != healthDomain.ObservationMetricRejected) { + return + } + collector.observations.WithLabelValues(string(level), string(result)).Inc() +} + +func registerCounterVec(registerer prometheus.Registerer, candidate *prometheus.CounterVec) (*prometheus.CounterVec, 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 checker metrics: %w", err) + } + existing, ok := registered.ExistingCollector.(*prometheus.CounterVec) + if !ok { + return nil, fmt.Errorf("register checker metrics: existing collector has unexpected type") + } + return existing, nil + } +} + +func validCheckerLevel(level healthDomain.Level) bool { + switch level { + case healthDomain.LevelBasic, healthDomain.LevelEgress, healthDomain.LevelTarget: + return true + default: + return false + } +} diff --git a/internal/platform/metrics/checker_test.go b/internal/platform/metrics/checker_test.go new file mode 100644 index 0000000..0b7380c --- /dev/null +++ b/internal/platform/metrics/checker_test.go @@ -0,0 +1,79 @@ +package metrics + +import ( + "testing" + + "github.com/prometheus/client_golang/prometheus" + dto "github.com/prometheus/client_model/go" + + healthDomain "proxy-pool/internal/domain/health" +) + +func TestCheckerCollectorRecordsOnlyFixedDimensions(t *testing.T) { + registry := prometheus.NewRegistry() + collector, err := NewCheckerCollector(registry) + if err != nil { + t.Fatalf("NewCheckerCollector() = %v", err) + } + collector.ObserveTaskDispatch(healthDomain.LevelBasic, 2) + collector.ObserveTaskDispatch(healthDomain.LevelTarget, 1) + collector.ObserveTaskDispatch("invalid", 99) + collector.ObserveObservation(healthDomain.LevelBasic, healthDomain.ObservationMetricAccepted) + collector.ObserveObservation(healthDomain.LevelTarget, healthDomain.ObservationMetricRejected) + collector.ObserveObservation("invalid", healthDomain.ObservationMetricAccepted) + + assertMetricValue(t, registry, "proxy_pool_checker_tasks_dispatched_total", map[string]string{"level": "BASIC"}, 2) + assertMetricValue(t, registry, "proxy_pool_checker_tasks_dispatched_total", map[string]string{"level": "TARGET"}, 1) + assertMetricValue(t, registry, "proxy_pool_checker_observations_total", map[string]string{ + "level": "BASIC", "result": "accepted", + }, 1) + assertMetricValue(t, registry, "proxy_pool_checker_observations_total", map[string]string{ + "level": "TARGET", "result": "rejected", + }, 1) +} + +func TestNewCheckerCollectorReusesRegisteredCollectors(t *testing.T) { + registry := prometheus.NewRegistry() + first, err := NewCheckerCollector(registry) + if err != nil { + t.Fatalf("first NewCheckerCollector() = %v", err) + } + second, err := NewCheckerCollector(registry) + if err != nil { + t.Fatalf("second NewCheckerCollector() = %v", err) + } + first.ObserveTaskDispatch(healthDomain.LevelEgress, 1) + second.ObserveTaskDispatch(healthDomain.LevelEgress, 1) + assertMetricValue(t, registry, "proxy_pool_checker_tasks_dispatched_total", map[string]string{"level": "EGRESS"}, 2) +} + +func assertMetricValue(t *testing.T, registry *prometheus.Registry, name string, labels map[string]string, want float64) { + 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.GetCounter().GetValue() == want { + return + } + } + } + t.Fatalf("metric %s labels=%v value=%v was not found", name, labels, want) +} + +func metricLabelsMatch(labels []*dto.LabelPair, want map[string]string) bool { + if len(labels) != len(want) { + return false + } + for _, label := range labels { + if want[label.GetName()] != label.GetValue() { + return false + } + } + return true +}