From 4691b9150ca6fb090b627909f92d6870b63a34d5 Mon Sep 17 00:00:00 2001 From: youfak Date: Sun, 2 Aug 2026 14:08:45 +0800 Subject: [PATCH] feat: expose extraction metrics --- README.md | 6 ++ docs/development/implementation-plan.md | 7 +- docs/requirements/traceability.md | 2 +- findings.md | 5 ++ internal/controller/bootstrap/bootstrap.go | 77 ++++++++++------- internal/controller/extraction/metrics.go | 80 +++++++++++++++++ internal/controller/extraction/service.go | 30 ++++++- .../controller/extraction/service_test.go | 86 +++++++++++++++++++ internal/platform/metrics/extraction.go | 70 +++++++++++++++ internal/platform/metrics/extraction_test.go | 58 +++++++++++++ progress.md | 5 ++ 11 files changed, 385 insertions(+), 41 deletions(-) create mode 100644 internal/controller/extraction/metrics.go create mode 100644 internal/platform/metrics/extraction.go create mode 100644 internal/platform/metrics/extraction_test.go diff --git a/README.md b/README.md index cf08038..22e9136 100644 --- a/README.md +++ b/README.md @@ -51,6 +51,12 @@ Proxy Pool 用 Controller 协调这些变化,并让 Gateway 数据面只消费 `proxy_pool_controller_provider_valid_candidates_total` 与 `proxy_pool_controller_provider_new_proxies_total`;`class` 仅有 `valid`、 `empty`、`duplicate_only`、`error`,不包含 Upstream、Proxy 或错误文本标签。 +- **提取指标**:Controller 暴露 + `proxy_pool_controller_extraction_requests_total{result}`、 + `proxy_pool_controller_extraction_requested_proxies_total` 与 + `proxy_pool_controller_extraction_returned_proxies_total`;`result` 仅有完成、 + 部分、空、库存不足、幂等冲突、限流、不可用、无效与内部错误等固定枚举, + 不包含 Client、请求、过滤条件、Upstream、Proxy 或错误文本标签。 - **PostgreSQL 管理面**:持久化配置版本、Upstream/Routing 管理状态、Admin 审计与 Outbox;不保存 Proxy 明细或逐次提取记录。 - **Gateway 组件**:HTTP 正向代理、HTTPS CONNECT、双向 Tunnel、重试、超时、 diff --git a/docs/development/implementation-plan.md b/docs/development/implementation-plan.md index 4f09a80..8cfb2c4 100644 --- a/docs/development/implementation-plan.md +++ b/docs/development/implementation-plan.md @@ -205,8 +205,8 @@ Distribution/Admin 服务构造、错误合并和资源关闭。生产 Provider 停机;Admin disable 会取消 Runtime,reload 在提交前预检并在发布后替换运行实例。 组合 fixture 已验证隔离 Redis namespace 下的选主、Provider HTTP 调用、模板解析 和活动池写入。Controller Metrics 独立入口现已提供 `/livez`、`/readyz` 与基础 -Prometheus 运行时指标,三监听器隔离已通过测试;Checker、Gateway、Drain 与 -Provider 业务指标已接入,提取、容量指标和密钥安全的结构化日志仍待实现。双存储 bootstrap 已通过 PostgreSQL 18 + Redis 8.2 组合 fixture,覆盖 +Prometheus 运行时指标,三监听器隔离已通过测试;Checker、Gateway、Drain、 +Provider 与 Extraction 业务指标已接入,容量指标和密钥安全的结构化日志仍待实现。双存储 bootstrap 已通过 PostgreSQL 18 + Redis 8.2 组合 fixture,覆盖 迁移、启动配置提交、Readiness、Admin Status 和 Metrics 探针。 WorkerControlPlane 现已接入 Controller 生命周期:Register、ACK 和 Runtime @@ -354,7 +354,8 @@ Supervisor 也改为同时服从静态配置与管理态,消除两条启停消 拉取已暴露 `proxy_pool_controller_provider_fetch_results_total{class}`、 `proxy_pool_controller_provider_valid_candidates_total` 与 `proxy_pool_controller_provider_new_proxies_total`;`class` 仅允许 `valid`、`empty`、 -`duplicate_only`、`error`。Provider 指标已完成,提取、容量指标与密钥安全的结构化日志仍待补齐。 +`duplicate_only`、`error`。Extraction 还暴露固定 `result` 的请求次数、请求数与响应 +交付数,幂等重放按响应交付统计。容量指标与密钥安全的结构化日志仍待补齐。 ## Task 14: Documentation, Examples, and Diagrams diff --git a/docs/requirements/traceability.md b/docs/requirements/traceability.md index a162bdb..9cfc7ae 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 暴露固定等级和结果,Gateway Collector 暴露 `proxy_pool_gateway_outcomes_total{stage,result}` 与 `proxy_pool_gateway_outcome_queue_dropped_total`,Drain Collector 暴露 `proxy_pool_controller_drain_candidates_total{reason}` 与 `proxy_pool_controller_drains_started_total{reason}`,其中 reason 固定为 unhealthy/upstream_disabled;Provider Collector 暴露 `proxy_pool_controller_provider_fetch_results_total{class}`、`proxy_pool_controller_provider_valid_candidates_total` 与 `proxy_pool_controller_provider_new_proxies_total`,其中 class 固定为 valid/empty/duplicate_only/error;注册表测试锁定标签集。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`,Drain Collector 暴露 `proxy_pool_controller_drain_candidates_total{reason}` 与 `proxy_pool_controller_drains_started_total{reason}`,其中 reason 固定为 unhealthy/upstream_disabled;Provider Collector 暴露固定 class 的拉取次数与候选计数;Extraction Collector 暴露 `proxy_pool_controller_extraction_requests_total{result}`、请求数与交付数,result 固定为 complete/partial/empty/insufficient/idempotency_conflict/rate_limited/unavailable/invalid/error。注册表测试锁定标签集;容量指标和密钥安全的结构化日志仍待实现 | | TEST-001 | 覆盖对话中列出的 11 个关键并发与故障场景 | 9030-9082 | 测试清单;Redis 活动池由 Memory/Redis 公用契约覆盖,跨进程故障场景仍按清单推进 | diff --git a/findings.md b/findings.md index bdc4631..6cef519 100644 --- a/findings.md +++ b/findings.md @@ -56,6 +56,11 @@ Extract 请求使用哪些 Upstream。 文本或凭据拆分。 - Provider 指标已完成;提取、容量指标和密钥安全的结构化日志仍待实现。 +- Extraction 只按固定 `result` 聚合请求:`complete`、`partial`、`empty`、 + `insufficient`、`idempotency_conflict`、`rate_limited`、`unavailable`、`invalid`、 + `error`。计数只记录请求数量和响应交付数量;幂等重放不被误记为新的 Proxy 消费。 +- Extraction 指标已完成;容量指标和密钥安全的结构化日志仍待实现。 + ## 集群与性能 - 用户补充:高峰可能达到 100,000 请求/秒。 diff --git a/internal/controller/bootstrap/bootstrap.go b/internal/controller/bootstrap/bootstrap.go index d563edf..b032401 100644 --- a/internal/controller/bootstrap/bootstrap.go +++ b/internal/controller/bootstrap/bootstrap.go @@ -179,11 +179,54 @@ func runWithWorkerFactory( snapshotRefresh := worker.NewSnapshotRefreshBroker() dependencies := controllerRuntime.Dependencies{} + var checkerMetrics healthDomain.TaskMetricsObserver + var drainMetrics healthDomain.DrainMetricsObserver + var extractionMetrics extraction.MetricsObserver + 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 + drainCollector, drainCollectorErr := platformMetrics.NewDrainCollector(prometheus.DefaultRegisterer) + if drainCollectorErr != nil { + return fmt.Errorf("%w: build Drain metrics: %w", ErrStartup, drainCollectorErr) + } + drainMetrics = drainCollector + providerCollector, providerCollectorErr := platformMetrics.NewProviderCollector(prometheus.DefaultRegisterer) + if providerCollectorErr != nil { + return fmt.Errorf("%w: build Provider metrics: %w", ErrStartup, providerCollectorErr) + } + if registrar, ok := opened.providerResults.(provider.ResultObserverRegistrar); ok && !nilInterface(registrar) { + registrar.AddResultObserver(providerCollector) + } + extractionCollector, extractionCollectorErr := platformMetrics.NewExtractionCollector(prometheus.DefaultRegisterer) + if extractionCollectorErr != nil { + return fmt.Errorf("%w: build extraction metrics: %w", ErrStartup, extractionCollectorErr) + } + extractionMetrics = extractionCollector + handler, handlerErr := platformMetrics.NewHandler(platformMetrics.Dependencies{ + Gatherer: prometheus.DefaultGatherer, Readiness: opened.metricsReadiness, + }) + if handlerErr != nil { + return fmt.Errorf("%w: build metrics handler: %w", ErrStartup, handlerErr) + } + dependencies.MetricsHandler = handler + } if loaded.Value.Distribution.Enabled { if nilInterface(opened.activity) || nilInterface(opened.readiness) || nilInterface(opened.admission) { return errors.Join(ErrStartup, ErrInvalidOptions) } - service, serviceErr := extraction.NewService(opened.activity, extractionPolicy(loaded.Value), opened.admission, options.Now) + service, serviceErr := extraction.NewService( + opened.activity, + extractionPolicy(loaded.Value), + opened.admission, + options.Now, + extractionMetrics, + ) if serviceErr != nil { return fmt.Errorf("%w: build extraction service: %w", ErrStartup, serviceErr) } @@ -221,38 +264,6 @@ func runWithWorkerFactory( } dependencies.AdminService = service } - var checkerMetrics healthDomain.TaskMetricsObserver - var drainMetrics healthDomain.DrainMetricsObserver - 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 - drainCollector, drainCollectorErr := platformMetrics.NewDrainCollector(prometheus.DefaultRegisterer) - if drainCollectorErr != nil { - return fmt.Errorf("%w: build Drain metrics: %w", ErrStartup, drainCollectorErr) - } - drainMetrics = drainCollector - providerCollector, providerCollectorErr := platformMetrics.NewProviderCollector(prometheus.DefaultRegisterer) - if providerCollectorErr != nil { - return fmt.Errorf("%w: build Provider metrics: %w", ErrStartup, providerCollectorErr) - } - if registrar, ok := opened.providerResults.(provider.ResultObserverRegistrar); ok && !nilInterface(registrar) { - registrar.AddResultObserver(providerCollector) - } - handler, handlerErr := platformMetrics.NewHandler(platformMetrics.Dependencies{ - Gatherer: prometheus.DefaultGatherer, Readiness: opened.metricsReadiness, - }) - if handlerErr != nil { - return fmt.Errorf("%w: build metrics handler: %w", ErrStartup, handlerErr) - } - dependencies.MetricsHandler = handler - } - runners := make([]lifecycle.Runner, 0, 3+len(loaded.Value.Upstreams)) if hasHTTPRuntime(loaded.Value) { runner, err := factory.New(configurationStore.Current(), dependencies, controllerRuntime.Options{HTTP: options.HTTP}) diff --git a/internal/controller/extraction/metrics.go b/internal/controller/extraction/metrics.go new file mode 100644 index 0000000..2e29e64 --- /dev/null +++ b/internal/controller/extraction/metrics.go @@ -0,0 +1,80 @@ +package extraction + +import ( + "errors" + + domain "proxy-pool/internal/domain/extraction" +) + +// MetricResult is a fixed, low-cardinality classification of one extraction +// request after application-level validation, admission, and storage handling. +type MetricResult string + +const ( + MetricComplete MetricResult = "complete" + MetricPartial MetricResult = "partial" + MetricEmpty MetricResult = "empty" + MetricInsufficient MetricResult = "insufficient" + MetricIdempotencyConflict MetricResult = "idempotency_conflict" + MetricRateLimited MetricResult = "rate_limited" + MetricUnavailable MetricResult = "unavailable" + MetricInvalid MetricResult = "invalid" + MetricError MetricResult = "error" +) + +// MetricsEvent contains only aggregate request and response counts. It must +// never contain a Client, Proxy, Upstream, filter, identifier, or error text. +type MetricsEvent struct { + Result MetricResult + Requested int + Returned int +} + +// MetricsObserver receives one completed extraction request. Implementations +// must return promptly and must not perform network or persistent storage I/O. +type MetricsObserver interface { + ObserveExtraction(MetricsEvent) +} + +func classifyMetricsEvent(response Response, resultErr error) MetricsEvent { + event := MetricsEvent{Requested: response.Requested, Returned: response.Returned} + if resultErr == nil { + switch { + case response.Returned == response.Requested: + event.Result = MetricComplete + case response.Returned > 0 && response.Returned < response.Requested: + event.Result = MetricPartial + case response.Returned == 0: + event.Result = MetricEmpty + default: + event.Result = MetricError + } + return event + } + switch { + case errors.Is(resultErr, domain.ErrInsufficientProxies): + event.Result = MetricInsufficient + case errors.Is(resultErr, domain.ErrIdempotencyConflict): + event.Result = MetricIdempotencyConflict + case errors.Is(resultErr, ErrAdmissionRejected): + event.Result = MetricRateLimited + case errors.Is(resultErr, ErrUnavailable): + event.Result = MetricUnavailable + case errors.Is(resultErr, ErrInvalidRequest), errors.Is(resultErr, ErrCountExceeded), + errors.Is(resultErr, ErrInvalidFulfillment), errors.Is(resultErr, domain.ErrInvalidCommand): + event.Result = MetricInvalid + default: + event.Result = MetricError + } + return event +} + +func validMetricResult(result MetricResult) bool { + switch result { + case MetricComplete, MetricPartial, MetricEmpty, MetricInsufficient, MetricIdempotencyConflict, + MetricRateLimited, MetricUnavailable, MetricInvalid, MetricError: + return true + default: + return false + } +} diff --git a/internal/controller/extraction/service.go b/internal/controller/extraction/service.go index 75626b8..5a04360 100644 --- a/internal/controller/extraction/service.go +++ b/internal/controller/extraction/service.go @@ -73,9 +73,16 @@ type Service struct { policy Policy admission admission.Admitter now func() time.Time + metrics MetricsObserver } -func NewService(store domain.Store, policy Policy, admitter admission.Admitter, now func() time.Time) (*Service, error) { +func NewService( + store domain.Store, + policy Policy, + admitter admission.Admitter, + now func() time.Time, + metrics ...MetricsObserver, +) (*Service, error) { if admitter == nil { return nil, fmt.Errorf("%w: admission is required", ErrInvalidServicePolicy) } @@ -92,11 +99,19 @@ func NewService(store domain.Store, policy Policy, admitter admission.Admitter, if now == nil { now = time.Now } - return &Service{store: store, policy: policy, admission: admitter, now: now}, nil + if len(metrics) > 1 { + return nil, ErrInvalidServicePolicy + } + var observer MetricsObserver + if len(metrics) == 1 { + observer = metrics[0] + } + return &Service{store: store, policy: policy, admission: admitter, now: now, metrics: observer}, nil } -func (s *Service) Extract(ctx context.Context, request Request) (Response, error) { - response := Response{RequestID: request.RequestID, Requested: request.Count} +func (s *Service) Extract(ctx context.Context, request Request) (response Response, resultErr error) { + response = Response{RequestID: request.RequestID, Requested: request.Count} + defer func() { s.observeMetrics(classifyMetricsEvent(response, resultErr)) }() if request.RequestID == "" || (request.ClientID == "" && request.SourceIP == "") || request.Count <= 0 { return response, ErrInvalidRequest } @@ -176,6 +191,13 @@ func (s *Service) Extract(ctx context.Context, request Request) (Response, error return response, nil } +func (s *Service) observeMetrics(event MetricsEvent) { + if s == nil || s.metrics == nil || !validMetricResult(event.Result) { + return + } + s.metrics.ObserveExtraction(event) +} + func admissionKey(request Request) string { if request.ClientID != "" { return "client:" + request.ClientID diff --git a/internal/controller/extraction/service_test.go b/internal/controller/extraction/service_test.go index a2b5b89..3d52078 100644 --- a/internal/controller/extraction/service_test.go +++ b/internal/controller/extraction/service_test.go @@ -244,6 +244,86 @@ func TestServiceIdempotentReplayKeepsOriginalExtractionTime(t *testing.T) { } } +func TestServiceEmitsFixedMetricsForTerminalResults(t *testing.T) { + observer := &recordingMetricsObserver{} + policy := Policy{MaxCountPerRequest: 2, DefaultFulfillment: domain.Partial} + tests := []struct { + name string + request Request + result domain.Result + storeErr error + admitErr error + wantEvent MetricsEvent + }{ + { + name: "complete", + request: Request{RequestID: "req-complete", ClientID: "client-a", Count: 2}, + result: domain.Result{Requested: 2, Returned: 2, Items: []domain.Candidate{{}, {}}}, + wantEvent: MetricsEvent{Result: MetricComplete, Requested: 2, Returned: 2}, + }, + { + name: "partial", + request: Request{RequestID: "req-partial", ClientID: "client-a", Count: 2}, + result: domain.Result{Requested: 2, Returned: 1, Items: []domain.Candidate{{}}}, + wantEvent: MetricsEvent{Result: MetricPartial, Requested: 2, Returned: 1}, + }, + { + name: "empty", + request: Request{RequestID: "req-empty", ClientID: "client-a", Count: 1}, + result: domain.Result{Requested: 1}, + wantEvent: MetricsEvent{Result: MetricEmpty, Requested: 1, Returned: 0}, + }, + { + name: "unexpected over delivery", + request: Request{RequestID: "req-over", ClientID: "client-a", Count: 1}, + result: domain.Result{Requested: 1, Returned: 2, Items: []domain.Candidate{{}, {}}}, + wantEvent: MetricsEvent{Result: MetricError, Requested: 1, Returned: 2}, + }, + { + name: "insufficient", + request: Request{RequestID: "req-insufficient", ClientID: "client-a", Count: 1}, + storeErr: domain.ErrInsufficientProxies, + wantEvent: MetricsEvent{Result: MetricInsufficient, Requested: 1}, + }, + { + name: "idempotency conflict", + request: Request{RequestID: "req-conflict", ClientID: "client-a", Count: 1}, + storeErr: domain.ErrIdempotencyConflict, + wantEvent: MetricsEvent{Result: MetricIdempotencyConflict, Requested: 1}, + }, + { + name: "rate limited", + request: Request{RequestID: "req-rate", ClientID: "client-a", Count: 1}, + admitErr: errors.New("rate limited"), + wantEvent: MetricsEvent{Result: MetricRateLimited, Requested: 1}, + }, + { + name: "unavailable", + request: Request{RequestID: "req-unavailable", ClientID: "client-a", Count: 1}, + storeErr: domain.ErrStoreUnavailable, + wantEvent: MetricsEvent{Result: MetricUnavailable, Requested: 1}, + }, + { + name: "invalid", + request: Request{RequestID: "req-invalid", Count: 1}, + wantEvent: MetricsEvent{Result: MetricInvalid, Requested: 1}, + }, + } + for _, test := range tests { + t.Run(test.name, func(t *testing.T) { + service, err := NewService(&recordingStore{result: test.result, err: test.storeErr}, policy, + &recordingAdmission{err: test.admitErr}, time.Now, observer) + if err != nil { + t.Fatalf("NewService() error = %v", err) + } + _, _ = service.Extract(context.Background(), test.request) + if got := observer.events[len(observer.events)-1]; got != test.wantEvent { + t.Fatalf("metrics event = %+v, want %+v", got, test.wantEvent) + } + }) + } +} + type recordingStore struct { command domain.Command result domain.Result @@ -257,6 +337,8 @@ type recordingAdmission struct { calls int } +type recordingMetricsObserver struct{ events []MetricsEvent } + type allowAllAdmission struct{} func (allowAllAdmission) Admit(context.Context, string) error { return nil } @@ -267,6 +349,10 @@ func (a *recordingAdmission) Admit(_ context.Context, key string) error { return a.err } +func (observer *recordingMetricsObserver) ObserveExtraction(event MetricsEvent) { + observer.events = append(observer.events, event) +} + func (s *recordingStore) Extract(_ context.Context, command domain.Command) (domain.Result, error) { s.calls++ s.command = command diff --git a/internal/platform/metrics/extraction.go b/internal/platform/metrics/extraction.go new file mode 100644 index 0000000..51a23fc --- /dev/null +++ b/internal/platform/metrics/extraction.go @@ -0,0 +1,70 @@ +package metrics + +import ( + "github.com/prometheus/client_golang/prometheus" + + controllerExtraction "proxy-pool/internal/controller/extraction" +) + +// ExtractionCollector exposes bounded Distribution business metrics. Client, +// proxy, upstream, filter, request, and error values are never labels. +type ExtractionCollector struct { + requests *prometheus.CounterVec + requested prometheus.Counter + returned prometheus.Counter +} + +var _ controllerExtraction.MetricsObserver = (*ExtractionCollector)(nil) + +func NewExtractionCollector(registerer prometheus.Registerer) (*ExtractionCollector, error) { + if registerer == nil { + return nil, ErrInvalidDependencies + } + requests, err := registerCounterVec(registerer, prometheus.NewCounterVec(prometheus.CounterOpts{ + Namespace: "proxy_pool", Subsystem: "controller", Name: "extraction_requests_total", + Help: "Number of completed Distribution extraction requests by fixed result.", + }, []string{"result"})) + if err != nil { + return nil, err + } + requested, err := registerCounter(registerer, prometheus.NewCounter(prometheus.CounterOpts{ + Namespace: "proxy_pool", Subsystem: "controller", Name: "extraction_requested_proxies_total", + Help: "Number of proxies requested from completed Distribution extraction requests.", + })) + if err != nil { + return nil, err + } + returned, err := registerCounter(registerer, prometheus.NewCounter(prometheus.CounterOpts{ + Namespace: "proxy_pool", Subsystem: "controller", Name: "extraction_returned_proxies_total", + Help: "Number of proxies returned in Distribution extraction responses, including idempotent replays.", + })) + if err != nil { + return nil, err + } + return &ExtractionCollector{requests: requests, requested: requested, returned: returned}, nil +} + +func (collector *ExtractionCollector) ObserveExtraction(event controllerExtraction.MetricsEvent) { + if collector == nil || collector.requests == nil || !validExtractionMetricResult(event.Result) { + return + } + collector.requests.WithLabelValues(string(event.Result)).Inc() + if event.Requested > 0 && collector.requested != nil { + collector.requested.Add(float64(event.Requested)) + } + if event.Returned > 0 && collector.returned != nil { + collector.returned.Add(float64(event.Returned)) + } +} + +func validExtractionMetricResult(result controllerExtraction.MetricResult) bool { + switch result { + case controllerExtraction.MetricComplete, controllerExtraction.MetricPartial, controllerExtraction.MetricEmpty, + controllerExtraction.MetricInsufficient, controllerExtraction.MetricIdempotencyConflict, + controllerExtraction.MetricRateLimited, controllerExtraction.MetricUnavailable, + controllerExtraction.MetricInvalid, controllerExtraction.MetricError: + return true + default: + return false + } +} diff --git a/internal/platform/metrics/extraction_test.go b/internal/platform/metrics/extraction_test.go new file mode 100644 index 0000000..04636d6 --- /dev/null +++ b/internal/platform/metrics/extraction_test.go @@ -0,0 +1,58 @@ +package metrics + +import ( + "testing" + + "github.com/prometheus/client_golang/prometheus" + + controllerExtraction "proxy-pool/internal/controller/extraction" +) + +func TestExtractionCollectorRecordsOnlyFixedResults(t *testing.T) { + registry := prometheus.NewRegistry() + collector, err := NewExtractionCollector(registry) + if err != nil { + t.Fatalf("NewExtractionCollector() error = %v", err) + } + collector.ObserveExtraction(controllerExtraction.MetricsEvent{ + Result: controllerExtraction.MetricComplete, Requested: 2, Returned: 2, + }) + collector.ObserveExtraction(controllerExtraction.MetricsEvent{ + Result: controllerExtraction.MetricPartial, Requested: 3, Returned: 1, + }) + collector.ObserveExtraction(controllerExtraction.MetricsEvent{ + Result: controllerExtraction.MetricUnavailable, Requested: 1, + }) + collector.ObserveExtraction(controllerExtraction.MetricsEvent{ + Result: "client-a", Requested: 10, Returned: 10, + }) + + assertMetricValue(t, registry, "proxy_pool_controller_extraction_requests_total", map[string]string{"result": "complete"}, 1) + assertMetricValue(t, registry, "proxy_pool_controller_extraction_requests_total", map[string]string{"result": "partial"}, 1) + assertMetricValue(t, registry, "proxy_pool_controller_extraction_requests_total", map[string]string{"result": "unavailable"}, 1) + assertMetricValue(t, registry, "proxy_pool_controller_extraction_requested_proxies_total", nil, 6) + assertMetricValue(t, registry, "proxy_pool_controller_extraction_returned_proxies_total", nil, 3) +} + +func TestNewExtractionCollectorReusesRegisteredCollectors(t *testing.T) { + registry := prometheus.NewRegistry() + first, err := NewExtractionCollector(registry) + if err != nil { + t.Fatalf("first NewExtractionCollector() error = %v", err) + } + second, err := NewExtractionCollector(registry) + if err != nil { + t.Fatalf("second NewExtractionCollector() error = %v", err) + } + first.ObserveExtraction(controllerExtraction.MetricsEvent{ + Result: controllerExtraction.MetricEmpty, Requested: 1, + }) + second.ObserveExtraction(controllerExtraction.MetricsEvent{ + Result: controllerExtraction.MetricComplete, Requested: 1, Returned: 1, + }) + + assertMetricValue(t, registry, "proxy_pool_controller_extraction_requests_total", map[string]string{"result": "empty"}, 1) + assertMetricValue(t, registry, "proxy_pool_controller_extraction_requests_total", map[string]string{"result": "complete"}, 1) + assertMetricValue(t, registry, "proxy_pool_controller_extraction_requested_proxies_total", nil, 2) + assertMetricValue(t, registry, "proxy_pool_controller_extraction_returned_proxies_total", nil, 1) +} diff --git a/progress.md b/progress.md index 121d10b..6d30a50 100644 --- a/progress.md +++ b/progress.md @@ -2,6 +2,11 @@ ## 2026-08-02 +- Extraction 可观测性已接入公用 `MetricsObserver`:Controller 暴露固定 `result` + 的请求计数、请求 Proxy 总数和响应交付 Proxy 总数。事件在服务单一出口分类,覆盖 + complete/partial/empty/insufficient/idempotency_conflict/rate_limited/unavailable/ + invalid/error;不带 Client、Request、Filter、Upstream、Proxy 或错误文本标签。 + 幂等重放统计响应交付而不误计为新的 Redis 消费。容量指标和密钥安全结构化日志仍待实现。 - Provider 拉取可观测性已接入公用 `ResultObserver`:Controller 暴露 `proxy_pool_controller_provider_fetch_results_total{class}`、 `proxy_pool_controller_provider_valid_candidates_total` 与