diff --git a/AGENTS.md b/AGENTS.md index dbff989..0ae6b75 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -300,17 +300,15 @@ end ### Metrics Collection -The SDK supports legacy and canonical metric surfaces, selected by the -`WORKER_CANONICAL_METRICS` environment variable. `MetricsCollector.create` -returns the appropriate collector: +`MetricsCollector.create` returns a collector that emits the canonical +(harmonized) metric surface: ```ruby metrics = Conductor::Worker::Telemetry::MetricsCollector.create(backend: :prometheus) ``` See [docs/METRICS_AND_INTERCEPTORS.md](docs/METRICS_AND_INTERCEPTORS.md) for -the full legacy and canonical metrics catalogs, label reference, and migration -guide. +the full metrics catalog and label reference. ### Worker Configuration (3-Tier Hierarchy) @@ -407,11 +405,8 @@ lib/conductor/ │ │ ├── listeners.rb # Listener protocol │ │ └── listener_registry.rb # Registration helper │ └── telemetry/ # Metrics -│ ├── metrics_collector.rb # Factory (WORKER_CANONICAL_METRICS gate) -│ ├── legacy_metrics_collector.rb # Legacy metric set -│ ├── canonical_metrics_collector.rb # Canonical metric set -│ ├── prometheus_backend.rb # Legacy Prometheus backend -│ └── canonical_prometheus_backend.rb # Canonical Prometheus backend +│ ├── metrics_collector.rb # MetricsCollector class + NullBackend +│ └── prometheus_backend.rb # PrometheusBackend + MetricsServer └── workflow/ ├── dsl/ # Workflow DSL │ ├── workflow_builder.rb # Core DSL engine (~1000 lines) diff --git a/CHANGELOG.md b/CHANGELOG.md index 0f679eb..4595cac 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -9,7 +9,7 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 ### Added -- Canonical metrics mode: opt-in harmonized metric surface via `WORKER_CANONICAL_METRICS=true` -- [details](docs/METRICS_AND_INTERCEPTORS.md#detailed-technical-notes----unreleased) +- Canonical (harmonized) metrics as the sole metric surface -- [details](docs/METRICS_AND_INTERCEPTORS.md#detailed-technical-notes----unreleased) - Bounded `uri` label on `http_api_client_request_seconds`: uses path templates (e.g. `/workflow/{workflowId}`) instead of fully-resolved paths, preventing metric cardinality explosion - `WorkflowStatusProbe` in harness: opt-in probe (via `HARNESS_PROBE_RATE_PER_SEC`) that exercises UUID-bearing endpoints to validate template URI metrics @@ -23,12 +23,10 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 - Control flow blocks: `parallel do`, `decide expr do`, `loop_over items do` - Auto-generated task reference names - Simplified LLM task methods with hash-to-ChatMessage auto-conversion -- `MetricsCollector.new(...)` is deprecated; use `MetricsCollector.create(...)` instead. `.new` still works but logs a deprecation warning. The previous implementation is preserved as `LegacyMetricsCollector` and remains the default. -- Legacy metrics emit unchanged by default; existing dashboards and alerts continue to work without modification -- HTTP request timing (`http_api_client_request_seconds`) is now zero-overhead in legacy mode: `RestClient` only enters the timing path when a canonical collector is subscribed to `GlobalDispatcher` +- HTTP request timing (`http_api_client_request_seconds`) is zero-overhead when no collector is active: `RestClient` only enters the timing path when a `MetricsCollector` is subscribed to `GlobalDispatcher` - `thread_uncaught_exceptions_total` is no longer incremented for caught exceptions in the polling loop; the metric surface is retained but unwired, matching the Python and JavaScript SDKs -- Both `CanonicalMetricsCollector` and `LegacyMetricsCollector` now respond to `stop`; `TaskHandler#stop` calls it automatically to unsubscribe from process-wide dispatchers -- `MetricsCollector.create` accepts `measure_payload_size:` (default `true` for canonical, `false` for legacy) to opt out of `workflow_input_size_bytes` JSON serialization overhead +- `MetricsCollector` responds to `stop`; `TaskHandler#stop` calls it automatically to unsubscribe from process-wide dispatchers +- `MetricsCollector.create` accepts `measure_payload_size:` (default `true`) to opt out of `workflow_input_size_bytes` JSON serialization overhead ### Removed diff --git a/docs/METRICS_AND_INTERCEPTORS.md b/docs/METRICS_AND_INTERCEPTORS.md index ada3e8c..fd2f289 100644 --- a/docs/METRICS_AND_INTERCEPTORS.md +++ b/docs/METRICS_AND_INTERCEPTORS.md @@ -5,19 +5,15 @@ execution, task result updates, payload sizes, workflow starts, and HTTP API client latency. It also provides an event-driven interceptor system for custom logging, error tracking, and observability. -This document covers the Ruby SDK metrics emitted by `MetricsCollector.create`, -`LegacyMetricsCollector`, and `CanonicalMetricsCollector`. It does not cover -Conductor server metrics or metrics emitted by other SDKs. +This document covers the Ruby SDK metrics emitted by `MetricsCollector`. It +does not cover Conductor server metrics or metrics emitted by other SDKs. ## Table of Contents -- [Legacy and Canonical Modes](#legacy-and-canonical-modes) - [Quick Start](#quick-start) -- [Canonical Metrics Catalog](#canonical-metrics-catalog) -- [Legacy Metrics Catalog](#legacy-metrics-catalog) +- [Metrics Catalog](#metrics-catalog) - [Metrics Not Applicable to Ruby](#metrics-not-applicable-to-ruby) - [Labels](#labels) -- [Migration from Legacy to Canonical](#migration-from-legacy-to-canonical) - [Prometheus Integration](#prometheus-integration) - [Custom Metrics Backends](#custom-metrics-backends) - [Troubleshooting](#troubleshooting) @@ -30,61 +26,32 @@ Conductor server metrics or metrics emitted by other SDKs. --- -## Legacy and Canonical Modes - -The Ruby SDK currently supports two mutually exclusive metric surfaces: - -- **Legacy metrics** are the default. They preserve the original Ruby SDK names - and labels, including snake_case label keys like `task_type`. -- **Canonical metrics** are opt-in with `WORKER_CANONICAL_METRICS=true`. They - use the cross-SDK canonical names, labels, units, and Prometheus histogram - bucket boundaries. - -`MetricsCollector.create` reads `WORKER_CANONICAL_METRICS` when the collector -is created: - -| Environment variable | Values | Effect | -|---|---|---| -| `WORKER_CANONICAL_METRICS` | `true`, `1`, or `yes` (case-insensitive, surrounding whitespace ignored) | Selects `CanonicalMetricsCollector`. | -| `WORKER_CANONICAL_METRICS` | unset, blank, `false`, `0`, `no`, or any other value | Selects `LegacyMetricsCollector`. | - -Only one implementation is active at a time. The SDK does not emit legacy and -canonical metrics simultaneously. Restart workers after changing -`WORKER_CANONICAL_METRICS` so the factory creates the desired collector. - -`WORKER_LEGACY_METRICS` is reserved for a future default-flip phase and is not -currently read by the Ruby SDK factory. - ### Payload Size Metrics Recording `workflow_input_size_bytes` requires JSON-serializing the workflow -input to measure its byte size. This is enabled by default in canonical mode -and disabled in legacy mode. Override explicitly via the factory: +input to measure its byte size. This is enabled by default. Override explicitly +via the factory: ```ruby -# Canonical mode, but skip the JSON serialization for large payloads +# Skip the JSON serialization for large payloads metrics = MetricsCollector.create(backend: :prometheus, measure_payload_size: false) ``` ### Collector Lifecycle -Both `CanonicalMetricsCollector` and `LegacyMetricsCollector` respond to -`stop`. Call `stop` to unsubscribe from process-wide dispatchers (e.g. the -`GlobalDispatcher` used for HTTP metrics). `TaskHandler#stop` calls `stop` -on all registered event listeners automatically. If you manage a collector -outside of `TaskHandler`, call `stop` when the collector is no longer needed. +`MetricsCollector` responds to `stop`. Call `stop` to unsubscribe from +process-wide dispatchers (e.g. the `GlobalDispatcher` used for HTTP metrics). +`TaskHandler#stop` calls `stop` on all registered event listeners +automatically. If you manage a collector outside of `TaskHandler`, call `stop` +when the collector is no longer needed. --- ## Quick Start -### Enabling Metrics (Legacy, Default) - ```ruby require 'conductor' -# MetricsCollector.create checks WORKER_CANONICAL_METRICS and returns -# the appropriate collector. Default (unset) selects legacy metrics. metrics = Conductor::Worker::Telemetry::MetricsCollector.create(backend: :prometheus) # Start metrics HTTP server @@ -104,25 +71,16 @@ handler.join metrics_server.stop ``` -### Enabling Canonical Metrics - -Set the environment variable before the worker starts: - -```shell -WORKER_CANONICAL_METRICS=true ruby my_worker.rb -``` - -The same code above will now return a `CanonicalMetricsCollector` instead. - --- -## Canonical Metrics Catalog +## Metrics Catalog -Canonical timing values are seconds. Canonical size values are bytes. Label -names use camelCase. Metrics are created lazily and appear in `/metrics` only -after the corresponding event records them. +Timing values are seconds. Size values are bytes. Label names use camelCase. +Metric definitions are pre-registered when the Prometheus backend is +created, but specific label-combination time series only appear in +`/metrics` after the corresponding event first records them. -### Canonical Counters +### Counters | Metric | Labels | Description | |---|---|---| @@ -135,9 +93,9 @@ after the corresponding event records them. | `thread_uncaught_exceptions_total` | `exception` | Incremented when a worker thread raises an uncaught exception. | | `workflow_start_error_total` | `workflowType`, `exception` | Incremented when starting a workflow fails client-side. | -### Canonical Time Histograms +### Time Histograms -All canonical time histograms use buckets (in seconds): +All time histograms use buckets (in seconds): ```text 0.001, 0.005, 0.01, 0.025, 0.05, 0.1, 0.25, 0.5, 1, 2.5, 5, 10 @@ -158,9 +116,9 @@ task_execute_time_seconds_count{taskType="my_task",status="SUCCESS"} 50.0 task_execute_time_seconds_sum{taskType="my_task",status="SUCCESS"} 2.3 ``` -### Canonical Size Histograms +### Size Histograms -All canonical size histograms use buckets (in bytes): +All size histograms use buckets (in bytes): ```text 100, 1000, 10000, 100000, 1000000, 10000000 @@ -171,7 +129,7 @@ All canonical size histograms use buckets (in bytes): | `task_result_size_bytes` | `taskType` | Serialized task result output size. | | `workflow_input_size_bytes` | `workflowType`, `version` | Serialized workflow input size. `version` is an empty string when the workflow version is absent. | -### Canonical Gauges +### Gauges | Metric | Labels | Description | |---|---|---| @@ -179,49 +137,6 @@ All canonical size histograms use buckets (in bytes): --- -## Legacy Metrics Catalog - -Legacy mode is the default so existing dashboards and alerts continue to work. -Legacy labels use snake_case (`task_type`). Legacy histograms do not carry a -`status` label. Legacy poll failure does not record poll time -- only the error -counter is incremented. - -### Legacy Counters - -| Metric | Labels | Description | -|---|---|---| -| `task_poll_total` | `task_type` | Incremented each time polling is done. | -| `task_poll_error_total` | `task_type`, `error` | Poll failures. `error` is the exception class name. | -| `task_execute_error_total` | `task_type`, `exception`, `retryable` | Task execution errors. `retryable` is `true` or `false`. | -| `task_update_failed_total` | `task_type` | Failed task result updates (critical -- task result lost). | - -### Legacy Time Histograms - -Legacy time histograms use buckets (in seconds): - -```text -0.005, 0.01, 0.025, 0.05, 0.1, 0.25, 0.5, 1, 2.5, 5, 10 -``` - -| Metric | Labels | Description | -|---|---|---| -| `task_poll_time_seconds` | `task_type` | Poll request latency. Recorded on successful polls only. | -| `task_execute_time_seconds` | `task_type` | Worker function execution duration. | - -### Legacy Size Histograms - -| Metric | Labels | Description | -|---|---|---| -| `task_result_size_bytes` | `task_type` | Serialized task result output size. Uses the same bucket set as canonical. | - -Legacy mode does not emit `task_execution_started_total`, -`task_update_time_seconds`, `task_paused_total`, -`thread_uncaught_exceptions_total`, `workflow_start_error_total`, -`http_api_client_request_seconds`, `workflow_input_size_bytes`, or -`active_workers`. - ---- - ## Metrics Not Applicable to Ruby The cross-SDK canonical catalog defines additional metrics that are not @@ -250,11 +165,40 @@ and JavaScript also define the metric surface but do not wire it. Future Ruby SDK versions may connect it to `Thread.report_on_exception` or a similar mechanism. -### Ractor Runner Limitations - -The `RactorTaskRunner` does not currently emit `active_workers` gauge updates -because each Ractor processes tasks sequentially with no shared count. All -other canonical and legacy metrics are emitted by the Ractor runner. +### Ractor Runner Limitations (Work-in-Progress -- Untested) + +> **Warning:** The `RactorTaskRunner` is an experimental work-in-progress. +> Metrics collection from Ractor-based workers is **partially implemented +> and not yet functional**. In practice, **no metrics are delivered** from Ractor workers +> in the current implementation due to the incomplete event bridge described +> below. Do not rely on Ractor worker metrics in production. + +The event communication channel between Ractor workers and the main-thread +`SyncEventDispatcher` is not yet implemented. The `TaskHandler` method +`create_event_receiver_ractor` returns `nil` (with the comment _"Ractor +event communication needs more work"_), and the spawned event-receiver thread +exits immediately without forwarding events. Because `RactorTaskRunner` +receives `nil` as its `event_queue`, the `publish_event` guard +(`return unless @event_queue`) causes **every event to be silently +discarded** -- poll events, execution events, update events, and failure +events alike. As a result, none of the metrics documented in this guide +(counters, histograms, or gauges) will be populated for Ractor-based +workers. + +Additionally: + +- The `active_workers` gauge is **never emitted** by `RactorTaskRunner`. + Unlike `TaskRunner` and `FiberTaskRunner`, which call + `publish_active_workers` when tasks start and complete, the Ractor runner + has no equivalent tracking or publishing logic. +- The `cleanup` method is a no-op stub. +- Ractor shutdown in `TaskHandler#stop` is rudimentary (`ractor.take` with + rescued errors) because Ractors lack a clean shutdown mechanism. + +These limitations will be resolved in a future release when proper +Ractor-to-main-thread event forwarding is implemented. Until then, use +`:thread` isolation (the default) or `:fiber` execution if you need metrics +and observability. --- @@ -262,58 +206,16 @@ other canonical and legacy metrics are emitted by the Ractor runner. | Label | Used by | Values | |---|---|---| -| `task_type` | Legacy worker metrics | Task definition name. Replaced by `taskType` in canonical mode. | -| `taskType` | Canonical worker metrics | Task definition name. | -| `workflowType` | Canonical workflow metrics | Workflow definition name. | +| `taskType` | Worker metrics | Task definition name. | +| `workflowType` | Workflow metrics | Workflow definition name. | | `version` | `workflow_input_size_bytes` | Workflow version as a string. Empty string when the version is absent. | -| `status` | Canonical task time metrics | `SUCCESS` or `FAILURE`. For `http_api_client_request_seconds`, the HTTP status code as a string (e.g. `"200"`), or `"0"` on network failure. | -| `exception` | Canonical error counters | Exception class name, such as `Faraday::TimeoutError`. | -| `error` | Legacy `task_poll_error_total` | Exception class name. Renamed to `exception` in canonical mode. | -| `retryable` | Legacy `task_execute_error_total` | `true` or `false`. Dropped in canonical mode. | +| `status` | Task time metrics | `SUCCESS` or `FAILURE`. For `http_api_client_request_seconds`, the HTTP status code as a string (e.g. `"200"`), or `"0"` on network failure. | +| `exception` | Error counters | Exception class name, such as `Faraday::TimeoutError`. | | `method` | HTTP metrics | HTTP verb (`GET`, `POST`, etc.). | | `uri` | HTTP metrics | API-relative path template (e.g. `/tasks/poll/batch/{taskType}`). Dynamic path segments retain `{placeholder}` tokens so label cardinality is bounded. | --- -## Migration from Legacy to Canonical - -Switching to canonical metrics is an explicit metrics-surface cutover. Enable -`WORKER_CANONICAL_METRICS=true` in a lower environment first, then update -dashboards, recording rules, and alerts before enabling it in production. - -Key changes: - -- Legacy task labels use `task_type`; canonical task labels use `taskType`. -- Legacy poll failure only increments the error counter; canonical also records - poll time with `status=FAILURE`. -- Legacy execution errors carry an extra `retryable` label; canonical drops it. -- Legacy poll errors use the `error` label; canonical uses `exception`. -- Legacy `task_update_failed_total` becomes `task_update_error_total` with an - added `exception` label. -- Canonical time histogram buckets start at 0.001s; legacy starts at 0.005s. -- Canonical adds metrics that legacy never emits: `task_execution_started_total`, - `task_update_time_seconds`, `task_paused_total`, - `thread_uncaught_exceptions_total`, `workflow_start_error_total`, - `http_api_client_request_seconds`, `workflow_input_size_bytes`, and - `active_workers`. -- Canonical and legacy collectors are mutually exclusive. During a migration, - compare scrape output by running separate worker instances or environments - with and without `WORKER_CANONICAL_METRICS=true`. - -Legacy-to-canonical replacements: - -| Legacy metric | Canonical replacement | -|---|---| -| `task_poll_total{task_type}` | `task_poll_total{taskType}` | -| `task_poll_time_seconds{task_type}` | `task_poll_time_seconds{taskType,status}` | -| `task_poll_error_total{task_type,error}` | `task_poll_error_total{taskType,exception}` | -| `task_execute_time_seconds{task_type}` | `task_execute_time_seconds{taskType,status}` | -| `task_execute_error_total{task_type,exception,retryable}` | `task_execute_error_total{taskType,exception}` | -| `task_result_size_bytes{task_type}` | `task_result_size_bytes{taskType}` | -| `task_update_failed_total{task_type}` | `task_update_error_total{taskType,exception}` | - ---- - ## Prometheus Integration ### Setup @@ -423,11 +325,11 @@ metrics = Conductor::Worker::Telemetry::MetricsCollector.create( ### Missing HTTP or Workflow Metrics -- `http_api_client_request_seconds` requires canonical mode. Legacy mode does - not emit HTTP metrics. The canonical collector auto-subscribes to - `GlobalDispatcher` for `HttpApiRequest` events from the HTTP layer. -- `workflow_input_size_bytes` and `workflow_start_error_total` require canonical - mode and only record when the corresponding `WorkflowExecutor` events fire. +- `http_api_client_request_seconds` requires the collector to be subscribed to + `GlobalDispatcher` for `HttpApiRequest` events from the HTTP layer. This + happens automatically when `subscribe_global_http: true` (the default). +- `workflow_input_size_bytes` and `workflow_start_error_total` only record when + the corresponding `WorkflowExecutor` events fire. ### High Cardinality @@ -435,8 +337,7 @@ metrics = Conductor::Worker::Telemetry::MetricsCollector.create( (e.g. `/workflow/{workflowId}`) to keep cardinality bounded. If you see fully-resolved paths in your metrics, verify that HTTP requests are going through the SDK's `ApiClient` rather than a standalone `RestClient`. -- Prefer canonical mode for bounded `exception` labels using exception class - names instead of raw error messages. +- The `exception` label uses exception class names to keep cardinality bounded. - Avoid embedding user identifiers or unbounded values in task type, workflow type, or other label values. @@ -487,29 +388,29 @@ the corresponding `on_*` method can listen for these events. | Event | When Published | Key Attributes | |---|---|---| | `TaskExecutionStarted` | Before task execution | `task_type`, `task_id`, `worker_id`, `workflow_instance_id` | -| `TaskExecutionCompleted` | After successful execution | `task_type`, `task_id`, `duration_ms`, `output_size_bytes` | -| `TaskExecutionFailure` | When execution fails | `task_type`, `task_id`, `duration_ms`, `cause`, `is_retryable` | +| `TaskExecutionCompleted` | After successful execution | `task_type`, `task_id`, `worker_id`, `workflow_instance_id`, `duration_ms`, `output_size_bytes` | +| `TaskExecutionFailure` | When execution fails | `task_type`, `task_id`, `worker_id`, `workflow_instance_id`, `duration_ms`, `cause`, `is_retryable` | ### Update Events | Event | When Published | Key Attributes | |---|---|---| -| `TaskUpdateCompleted` | After successful result update | `task_type`, `task_id`, `duration_ms` | -| `TaskUpdateFailure` | When result update fails after all retries | `task_type`, `task_id`, `retry_count`, `task_result`, `cause` | +| `TaskUpdateCompleted` | After successful result update | `task_type`, `task_id`, `worker_id`, `workflow_instance_id`, `duration_ms` | +| `TaskUpdateFailure` | When result update fails after all retries | `task_type`, `task_id`, `worker_id`, `workflow_instance_id`, `cause`, `retry_count`, `task_result`, `duration_ms` | ### Worker State Events | Event | When Published | Key Attributes | |---|---|---| | `TaskPaused` | When a paused worker skips a poll | `task_type` | -| `ThreadUncaughtException` | When a worker thread raises an uncaught exception | `cause` | +| `ThreadUncaughtException` | When a worker thread raises an uncaught exception | `cause`, `task_type` | | `ActiveWorkersChanged` | When the active worker count changes | `task_type`, `count` | ### Workflow Events | Event | When Published | Key Attributes | |---|---|---| -| `WorkflowStartError` | When starting a workflow fails client-side | `workflow_type`, `cause` | +| `WorkflowStartError` | When starting a workflow fails client-side | `workflow_type`, `version`, `cause` | | `WorkflowInputSize` | When a workflow is started | `workflow_type`, `version`, `size_bytes` | ### HTTP Events @@ -824,12 +725,9 @@ end ### Telemetry Classes -- `Conductor::Worker::Telemetry::MetricsCollector` - Factory module; `.create` returns the appropriate collector -- `Conductor::Worker::Telemetry::LegacyMetricsCollector` - Pre-harmonization metric set with `task_type` labels -- `Conductor::Worker::Telemetry::CanonicalMetricsCollector` - Canonical metric set with `taskType` labels +- `Conductor::Worker::Telemetry::MetricsCollector` - Canonical metric collector; `.create` is a convenience factory - `Conductor::Worker::Telemetry::NullBackend` - No-op metrics backend -- `Conductor::Worker::Telemetry::PrometheusBackend` - Legacy Prometheus backend -- `Conductor::Worker::Telemetry::CanonicalPrometheusBackend` - Canonical Prometheus backend +- `Conductor::Worker::Telemetry::PrometheusBackend` - Prometheus backend with canonical label schemas - `Conductor::Worker::Telemetry::MetricsServer` - WEBrick HTTP server for `/metrics` and `/health` endpoints ### Event Classes @@ -886,42 +784,39 @@ other dynamic path segments pass through the SDK. falls back to extracting the path from the full URL when `metric_uri` is `nil` (e.g. for direct `RestClient` calls outside of `ApiClient`). 5. The `HttpApiRequest` event carries the template string as its `uri` field. - `CanonicalMetricsCollector#on_http_api_request` records it as-is into the - histogram. + `MetricsCollector#on_http_api_request` records it as-is into the histogram. This approach mirrors the Python SDK (`metric_uri` parameter), the Java SDK (`PathTemplateTag` on the OkHttp request), and the Go SDK (`WithPathTemplate` context value). The base-URL path prefix (e.g. `/api`) is never included because `metric_uri` is always the raw API-relative resource path. -### Canonical metrics factory - -`MetricsCollector.create(backend:)` selects the collector implementation based -on the `WORKER_CANONICAL_METRICS` environment variable: +### MetricsCollector factory -- **Unset / falsy** -- `LegacyMetricsCollector` (default). Emits the 0.1.0 - metric names and `snake_case` label conventions unchanged. -- **Truthy** (`true`, `1`, `yes`, case-insensitive) -- - `CanonicalMetricsCollector`. Emits the harmonized cross-SDK catalog with - `camelCase` domain labels, Prometheus histograms with explicit bucket - boundaries, and the `exception` label derived from the Ruby exception class - name. +`MetricsCollector.create(backend:)` is a convenience constructor that returns +a `MetricsCollector` instance. It emits the harmonized cross-SDK catalog with +`camelCase` domain labels, Prometheus histograms with explicit bucket +boundaries, and the `exception` label derived from the Ruby exception class +name. -`WORKER_LEGACY_METRICS` is reserved for a future phase where canonical becomes -the default. +### Event system -### Event system additions +The following event types are used by the collector: -The following event types were added for the canonical collector: - -- `HttpApiRequest` -- emitted by `RestClient` via the process-wide - `GlobalDispatcher`, but only when at least one `HttpApiRequest` listener - is subscribed (i.e. a `CanonicalMetricsCollector` is active). In legacy - mode or with no collector, `RestClient` skips all timing overhead. +- `PollStarted`, `PollCompleted`, `PollFailure`, `TaskExecutionStarted`, + `TaskExecutionCompleted`, `TaskExecutionFailure`, `TaskUpdateCompleted`, + `TaskUpdateFailure`, `TaskPaused`, `ActiveWorkersChanged` -- emitted by + `TaskRunner` (and `FiberTaskRunner`). `RactorTaskRunner` constructs these + same events internally but **does not deliver them** to the + `SyncEventDispatcher` because the Ractor-to-main-thread event bridge is + not yet implemented (see + [Ractor Runner Limitations](#ractor-runner-limitations-work-in-progress----untested)). - `WorkflowStartError`, `WorkflowInputSize` -- emitted by `WorkflowExecutor`. -- `TaskUpdateCompleted`, `TaskPaused`, `ActiveWorkersChanged` -- emitted - by `TaskRunner`. +- `HttpApiRequest` -- emitted by `RestClient` via the process-wide + `GlobalDispatcher`, but only when at least one `HttpApiRequest` listener + is subscribed (i.e. a `MetricsCollector` is active). With no collector, + `RestClient` skips all timing overhead. - `ThreadUncaughtException` -- event class and collector handler exist for API completeness but are not currently emitted by any runner (see [thread_uncaught_exceptions_total](#thread_uncaught_exceptions_total)). diff --git a/docs/design/EVENT_INTERCEPTOR_SYSTEM.md b/docs/design/EVENT_INTERCEPTOR_SYSTEM.md index d21ccb3..3ea5885 100644 --- a/docs/design/EVENT_INTERCEPTOR_SYSTEM.md +++ b/docs/design/EVENT_INTERCEPTOR_SYSTEM.md @@ -125,12 +125,10 @@ TaskRunner SyncEventDispatcher Listeners | `SyncEventDispatcher` | `events/sync_event_dispatcher.rb` | Thread-safe event router | | `TaskRunnerEventsListener` | `events/listeners.rb` | Listener protocol (duck typing) | | `ListenerRegistry` | `events/listener_registry.rb` | Bulk listener registration | -| `MetricsCollector` | `telemetry/metrics_collector.rb` | Factory (WORKER_CANONICAL_METRICS gate) | -| `LegacyMetricsCollector` | `telemetry/legacy_metrics_collector.rb` | Legacy metric set | -| `CanonicalMetricsCollector` | `telemetry/canonical_metrics_collector.rb` | Canonical metric set | -| `PrometheusBackend` | `telemetry/prometheus_backend.rb` | Legacy Prometheus backend | -| `CanonicalPrometheusBackend` | `telemetry/canonical_prometheus_backend.rb` | Canonical Prometheus backend | +| `MetricsCollector` | `telemetry/metrics_collector.rb` | Canonical metric collector | | `NullBackend` | `telemetry/metrics_collector.rb` | No-op backend | +| `PrometheusBackend` | `telemetry/prometheus_backend.rb` | Prometheus backend with canonical label schemas | +| `MetricsServer` | `telemetry/prometheus_backend.rb` | WEBrick HTTP server for `/metrics` | --- @@ -465,18 +463,15 @@ handler = Conductor::Worker::TaskHandler.new( ### MetricsCollector -The SDK supports legacy and canonical metric surfaces, selected by the -`WORKER_CANONICAL_METRICS` environment variable. `MetricsCollector.create` -returns the appropriate collector (`LegacyMetricsCollector` or -`CanonicalMetricsCollector`): +`MetricsCollector.create` returns a collector that emits the canonical +(harmonized) metric surface: ```ruby metrics = Conductor::Worker::Telemetry::MetricsCollector.create(backend: :prometheus) ``` See [docs/METRICS_AND_INTERCEPTORS.md](../METRICS_AND_INTERCEPTORS.md) for the -full legacy and canonical metrics catalogs, label reference, and migration -guide. +full metrics catalog and label reference. ### Backend Protocol @@ -522,13 +517,11 @@ end ### Prometheus Backends -The SDK ships two Prometheus backends: - -- `PrometheusBackend` -- legacy metric registrations with `task_type` labels. -- `CanonicalPrometheusBackend` -- canonical metric registrations with `taskType` labels, `status` on time histograms, and canonical bucket boundaries. - -Both implement `increment`, `observe`, and `set` and integrate with the -`prometheus-client` gem. See [docs/METRICS_AND_INTERCEPTORS.md](../METRICS_AND_INTERCEPTORS.md) +The SDK ships `PrometheusBackend` with canonical metric registrations using +`taskType` labels, `status` on time histograms, and canonical bucket +boundaries. It implements `increment`, `observe`, and `set` and integrates +with the `prometheus-client` gem. See +[docs/METRICS_AND_INTERCEPTORS.md](../METRICS_AND_INTERCEPTORS.md) for the full metric catalog emitted by each backend. ### MetricsServer @@ -871,11 +864,8 @@ lib/conductor/worker/ │ ├── listeners.rb # TaskRunnerEventsListener protocol │ └── listener_registry.rb # Bulk listener registration helper ├── telemetry/ -│ ├── metrics_collector.rb # Factory (WORKER_CANONICAL_METRICS gate) + NullBackend -│ ├── legacy_metrics_collector.rb # Legacy metric set -│ ├── canonical_metrics_collector.rb # Canonical metric set -│ ├── prometheus_backend.rb # Legacy PrometheusBackend + MetricsServer -│ └── canonical_prometheus_backend.rb # Canonical PrometheusBackend +│ ├── metrics_collector.rb # MetricsCollector class + NullBackend +│ └── prometheus_backend.rb # PrometheusBackend + MetricsServer ├── task_runner.rb # Publishes events during polling/execution └── task_handler.rb # Creates dispatcher, registers listeners @@ -887,10 +877,7 @@ spec/conductor/worker/ │ └── listener_registry_spec.rb └── telemetry/ ├── metrics_collector_spec.rb - ├── legacy_metrics_collector_spec.rb - ├── canonical_metrics_collector_spec.rb - ├── prometheus_backend_spec.rb - └── canonical_prometheus_backend_spec.rb + └── prometheus_backend_spec.rb ``` --- diff --git a/docs/design/WORKER_DESIGN.md b/docs/design/WORKER_DESIGN.md index 1f190c4..1830ba2 100644 --- a/docs/design/WORKER_DESIGN.md +++ b/docs/design/WORKER_DESIGN.md @@ -1375,10 +1375,8 @@ end ### MetricsCollector -The SDK supports legacy and canonical metric surfaces, selected by the -`WORKER_CANONICAL_METRICS` environment variable. `MetricsCollector.create` -returns the appropriate collector (`LegacyMetricsCollector` or -`CanonicalMetricsCollector`): +`MetricsCollector.create` returns a collector that emits the canonical +(harmonized) metric surface: ```ruby metrics = Conductor::Worker::Telemetry::MetricsCollector.create(backend: :prometheus) @@ -1563,12 +1561,8 @@ lib/conductor/ │ ├── listener_registry.rb # Listener registration helper │ └── listeners.rb # Listener protocol module ├── worker/telemetry/ -│ ├── metrics_collector.rb # Factory (WORKER_CANONICAL_METRICS gate) -│ ├── legacy_metrics_collector.rb # Legacy metric set -│ ├── canonical_metrics_collector.rb # Canonical metric set -│ ├── prometheus_backend.rb # Legacy Prometheus backend -│ ├── canonical_prometheus_backend.rb # Canonical Prometheus backend -│ └── null_backend.rb # No-op backend +│ ├── metrics_collector.rb # MetricsCollector class + NullBackend +│ └── prometheus_backend.rb # PrometheusBackend + MetricsServer └── exceptions.rb # Add NonRetryableError ``` @@ -1606,15 +1600,22 @@ lib/conductor/ - Integration tests against local Conductor server - All Python SDK test scenarios ported -### Phase 2: Ractor-based Runner +### Phase 2: Ractor-based Runner (Work-in-Progress) **Goal:** True parallelism for CPU-bound workers. +**Status:** Partially implemented. The `RactorTaskRunner` can poll and execute +tasks inside Ractors, but the event bridge to the main thread is not yet +wired. This means metrics, interceptors, and custom event listeners receive +**no events** from Ractor workers. See +[Ractor Runner Limitations](../METRICS_AND_INTERCEPTORS.md#ractor-runner-limitations-work-in-progress----untested) +for details. + **Components:** -1. `RactorTaskRunner` -2. Ractor-local TaskContext storage -3. Event aggregation via Ractor messaging -4. `isolation: :ractor` configuration +1. `RactorTaskRunner` -- implemented, untested end-to-end +2. Ractor-local TaskContext storage -- implemented +3. Event aggregation via Ractor messaging -- **not yet implemented** +4. `isolation: :ractor` configuration -- implemented **Constraints:** - Requires Ruby 3.1+ diff --git a/examples/metrics_example.rb b/examples/metrics_example.rb index ad6e21b..9ad63f1 100644 --- a/examples/metrics_example.rb +++ b/examples/metrics_example.rb @@ -9,14 +9,22 @@ # - Track task poll times, execution times, errors, and more # - Integrate with Prometheus monitoring # -# Metrics collected: +# Metrics collected (canonical harmonized set): # - task_poll_total: Total number of task polls -# - task_poll_time_seconds: Task poll duration -# - task_poll_error_total: Poll errors by error type -# - task_execute_time_seconds: Task execution duration -# - task_execute_error_total: Execution errors by exception and retryability +# - task_poll_time_seconds: Task poll duration (with status label) +# - task_poll_error_total: Poll errors by exception type +# - task_execution_started_total: Tasks dispatched to worker function +# - task_execute_time_seconds: Task execution duration (with status label) +# - task_execute_error_total: Execution errors by exception # - task_result_size_bytes: Task result payload size -# - task_update_failed_total: Failed task updates (CRITICAL) +# - task_update_error_total: Failed task updates (CRITICAL) +# - task_update_time_seconds: Task result update latency +# - task_paused_total: Workers paused +# - active_workers: Current active worker count (gauge) +# - workflow_start_error_total: Workflow start failures +# - workflow_input_size_bytes: Workflow input payload size +# - http_api_client_request_seconds: HTTP API client latency +# - thread_uncaught_exceptions_total: Uncaught thread exceptions # # Requirements: # gem 'prometheus-client' diff --git a/harness/main.rb b/harness/main.rb index 3c51035..00e778b 100644 --- a/harness/main.rb +++ b/harness/main.rb @@ -42,7 +42,7 @@ def self.main metrics_collector = Conductor::Worker::Telemetry::MetricsCollector.create(backend: :prometheus) metrics_server = Conductor::Worker::Telemetry::MetricsServer.new(port: metrics_port) metrics_server.start - puts "Prometheus metrics server started on port #{metrics_port} (#{metrics_collector.collector_name} metrics)" + puts "Prometheus metrics server started on port #{metrics_port}" workers = SIMULATED_WORKERS.map do |def_entry| sim = SimulatedTaskWorker.new( diff --git a/harness/manifests/deployment.yaml b/harness/manifests/deployment.yaml index 0a0dbfa..0036ec9 100644 --- a/harness/manifests/deployment.yaml +++ b/harness/manifests/deployment.yaml @@ -53,14 +53,6 @@ spec: - name: HARNESS_POLL_INTERVAL_MS value: "100" - # === METRICS IMPLEMENTATION === - # Set to "true" to use the canonical (harmonized) metric set. - # Default "false" uses legacy metrics during the deprecation period. - # In a future release, canonical will become the default and - # WORKER_LEGACY_METRICS will allow opting back into legacy. - - name: WORKER_CANONICAL_METRICS - value: "true" - ports: - name: metrics containerPort: 9991 diff --git a/lib/conductor/worker/telemetry/canonical_metrics_collector.rb b/lib/conductor/worker/telemetry/canonical_metrics_collector.rb deleted file mode 100644 index 948d6a0..0000000 --- a/lib/conductor/worker/telemetry/canonical_metrics_collector.rb +++ /dev/null @@ -1,182 +0,0 @@ -# frozen_string_literal: true - -require 'logger' -require_relative '../events/listeners' -require_relative '../events/global_dispatcher' - -module Conductor - module Worker - module Telemetry - # CanonicalMetricsCollector - Canonical SDK worker metrics from the - # harmonization spec (sdk-metrics-harmonization.md). - # - # Selected when WORKER_CANONICAL_METRICS is truthy. Uses camelCase domain - # labels (taskType, workflowType) and includes status labels on time - # histograms. - # - # Legacy-only event handlers that have no canonical equivalent are - # implemented as no-ops so this collector satisfies the full listener - # interface and can be used interchangeably with LegacyMetricsCollector. - class CanonicalMetricsCollector - include Events::TaskRunnerEventsListener - include Events::WorkflowEventsListener - include Events::HttpEventsListener - - STATUS_SUCCESS = 'SUCCESS' - STATUS_FAILURE = 'FAILURE' - - # @param backend [Symbol, Object] :null, :prometheus, or a custom backend - # @param subscribe_global_http [Boolean] Auto-subscribe to GlobalDispatcher - # for HttpApiRequest events from the HTTP layer (default true). - # @param measure_payload_size [Boolean] Record workflow_input_size_bytes - # (requires JSON serialization; default true). Set false to skip - # serialization overhead for large payloads. - def initialize(backend: :null, subscribe_global_http: true, measure_payload_size: true, logger: nil) - @backend = load_backend(backend) - @logger = logger || Logger.new(File::NULL) - @measure_payload_size = measure_payload_size - @http_listener = nil - subscribe_to_global_http_events if subscribe_global_http - end - - attr_reader :backend, :measure_payload_size - - def stop - return unless @http_listener - - Events::GlobalDispatcher.instance.unregister(Events::HttpApiRequest, @http_listener) - @http_listener = nil - rescue StandardError => e - @logger.debug { "Telemetry error (non-fatal): #{e.class}: #{e.message}" } - end - - def collector_name - 'canonical' - end - - # --- Task Runner Event Handlers --- - - def on_poll_started(event) - @backend.increment('task_poll_total', labels: { taskType: event.task_type }) - end - - def on_poll_completed(event) - observe_time('task_poll_time_seconds', event.duration_ms, - { taskType: event.task_type, status: STATUS_SUCCESS }) - end - - def on_poll_failure(event) - @backend.increment('task_poll_error_total', - labels: { taskType: event.task_type, exception: event.cause.class.name }) - observe_time('task_poll_time_seconds', event.duration_ms, - { taskType: event.task_type, status: STATUS_FAILURE }) - end - - def on_task_execution_started(event) - @backend.increment('task_execution_started_total', labels: { taskType: event.task_type }) - end - - def on_task_execution_completed(event) - observe_time('task_execute_time_seconds', event.duration_ms, - { taskType: event.task_type, status: STATUS_SUCCESS }) - - return unless event.output_size_bytes - - @backend.observe('task_result_size_bytes', event.output_size_bytes, - labels: { taskType: event.task_type }) - end - - def on_task_execution_failure(event) - @backend.increment('task_execute_error_total', - labels: { taskType: event.task_type, exception: event.cause.class.name }) - observe_time('task_execute_time_seconds', event.duration_ms, - { taskType: event.task_type, status: STATUS_FAILURE }) - end - - def on_task_update_completed(event) - observe_time('task_update_time_seconds', event.duration_ms, - { taskType: event.task_type, status: STATUS_SUCCESS }) - end - - def on_task_update_failure(event) - @backend.increment('task_update_error_total', - labels: { taskType: event.task_type, exception: event.cause.class.name }) - - return unless event.respond_to?(:duration_ms) && event.duration_ms - - observe_time('task_update_time_seconds', event.duration_ms, - { taskType: event.task_type, status: STATUS_FAILURE }) - end - - def on_task_paused(event) - @backend.increment('task_paused_total', labels: { taskType: event.task_type }) - end - - def on_thread_uncaught_exception(event) - @backend.increment('thread_uncaught_exceptions_total', - labels: { exception: event.cause.class.name }) - end - - def on_active_workers_changed(event) - @backend.set('active_workers', event.count, labels: { taskType: event.task_type }) - end - - # --- Workflow Event Handlers --- - - def on_workflow_start_error(event) - @backend.increment('workflow_start_error_total', - labels: { workflowType: event.workflow_type, - exception: event.cause.class.name }) - end - - def on_workflow_input_size(event) - return unless @measure_payload_size - - @backend.observe('workflow_input_size_bytes', event.size_bytes, - labels: { workflowType: event.workflow_type, - version: (event.version || '').to_s }) - end - - # --- HTTP Event Handlers --- - - def on_http_api_request(event) - observe_time('http_api_client_request_seconds', event.duration_ms, - { method: event.method, uri: event.uri, status: event.status }) - end - - private - - def observe_time(name, duration_ms, labels) - @backend.observe(name, duration_ms / 1000.0, labels: labels) - end - - def subscribe_to_global_http_events - @http_listener = ->(event) { on_http_api_request(event) } - Events::GlobalDispatcher.instance.register(Events::HttpApiRequest, @http_listener) - rescue StandardError => e - @logger.debug { "Telemetry error (non-fatal): #{e.class}: #{e.message}" } - end - - def load_backend(backend) - case backend - when :null, nil - NullBackend.new - when :prometheus - load_prometheus_backend - else - backend - end - end - - def load_prometheus_backend - require_relative 'canonical_prometheus_backend' - CanonicalPrometheusBackend.new - rescue LoadError - raise ConfigurationError, - "The 'prometheus-client' gem is required for Prometheus metrics. " \ - "Add `gem 'prometheus-client'` to your Gemfile." - end - end - end - end -end diff --git a/lib/conductor/worker/telemetry/canonical_prometheus_backend.rb b/lib/conductor/worker/telemetry/canonical_prometheus_backend.rb deleted file mode 100644 index a0401ef..0000000 --- a/lib/conductor/worker/telemetry/canonical_prometheus_backend.rb +++ /dev/null @@ -1,173 +0,0 @@ -# frozen_string_literal: true - -module Conductor - module Worker - module Telemetry - # CanonicalPrometheusBackend - Prometheus backend for the canonical SDK metric catalog. - # - # Pre-registers every metric from the harmonization spec with its canonical - # label set and bucket configuration. Uses camelCase domain labels (taskType, - # workflowType) per the canonical convention. - class CanonicalPrometheusBackend - TIME_BUCKETS = [0.001, 0.005, 0.01, 0.025, 0.05, 0.1, 0.25, 0.5, 1, 2.5, 5, 10].freeze - SIZE_BUCKETS = [100, 1000, 10_000, 100_000, 1_000_000, 10_000_000].freeze - - COUNTER_LABELS = { - 'task_poll_total' => %i[taskType], - 'task_execution_started_total' => %i[taskType], - 'task_poll_error_total' => %i[taskType exception], - 'task_execute_error_total' => %i[taskType exception], - 'task_update_error_total' => %i[taskType exception], - 'task_paused_total' => %i[taskType], - 'thread_uncaught_exceptions_total' => %i[exception], - 'workflow_start_error_total' => %i[workflowType exception] - }.freeze - - HISTOGRAM_LABELS = { - 'task_poll_time_seconds' => %i[taskType status], - 'task_execute_time_seconds' => %i[taskType status], - 'task_update_time_seconds' => %i[taskType status], - 'http_api_client_request_seconds' => %i[method uri status], - 'task_result_size_bytes' => %i[taskType], - 'workflow_input_size_bytes' => %i[workflowType version] - }.freeze - - GAUGE_LABELS = { - 'active_workers' => %i[taskType] - }.freeze - - HISTOGRAM_BUCKETS = { - 'task_result_size_bytes' => SIZE_BUCKETS, - 'workflow_input_size_bytes' => SIZE_BUCKETS - }.freeze - - def initialize(registry: nil) - load_prometheus_client - @registry = registry || Prometheus::Client.registry - @counters = {} - @histograms = {} - @gauges = {} - setup_metrics - end - - def increment(name, labels: {}, value: 1) - metric = get_or_create_counter(name) - metric.increment(labels: normalize_labels(name, labels, COUNTER_LABELS), by: value) - end - - def observe(name, value, labels: {}) - metric = get_or_create_histogram(name) - metric.observe(value, labels: normalize_labels(name, labels, HISTOGRAM_LABELS)) - end - - def set(name, value, labels: {}) - metric = get_or_create_gauge(name) - metric.set(value, labels: normalize_labels(name, labels, GAUGE_LABELS)) - end - - attr_reader :registry - - private - - def load_prometheus_client - require 'prometheus/client' - rescue LoadError - raise ConfigurationError, - "The 'prometheus-client' gem is required for Prometheus metrics. " \ - "Add `gem 'prometheus-client'` to your Gemfile." - end - - def setup_metrics - COUNTER_LABELS.each do |name, _| - register_counter(name, "Counter for #{name}") - end - - HISTOGRAM_LABELS.each do |name, _| - register_histogram(name, "Histogram for #{name}") - end - - GAUGE_LABELS.each do |name, _| - register_gauge(name, "Gauge for #{name}") - end - end - - def register_counter(name, docstring) - metric_name = name.to_sym - labels = COUNTER_LABELS.fetch(name, %i[taskType]) - @counters[name] = register_or_reuse(metric_name) do - Prometheus::Client::Counter.new(metric_name, docstring: docstring, labels: labels) - end - end - - def register_histogram(name, docstring) - metric_name = name.to_sym - labels = HISTOGRAM_LABELS.fetch(name, %i[taskType]) - buckets = HISTOGRAM_BUCKETS[name] || TIME_BUCKETS - @histograms[name] = register_or_reuse(metric_name) do - Prometheus::Client::Histogram.new(metric_name, docstring: docstring, - labels: labels, buckets: buckets) - end - end - - def register_gauge(name, docstring) - metric_name = name.to_sym - labels = GAUGE_LABELS.fetch(name, %i[taskType]) - @gauges[name] = register_or_reuse(metric_name) do - Prometheus::Client::Gauge.new(metric_name, docstring: docstring, labels: labels) - end - end - - def register_or_reuse(metric_name) - if @registry.exist?(metric_name) - @registry.get(metric_name) - else - metric = yield - @registry.register(metric) - metric - end - end - - def get_or_create_counter(name) - @counters[name] ||= register_or_reuse(name.to_sym) do - labels = COUNTER_LABELS.fetch(name, %i[taskType]) - Prometheus::Client::Counter.new(name.to_sym, docstring: "Counter for #{name}", labels: labels) - end - end - - def get_or_create_histogram(name) - @histograms[name] ||= register_or_reuse(name.to_sym) do - labels = HISTOGRAM_LABELS.fetch(name, %i[taskType]) - buckets = HISTOGRAM_BUCKETS[name] || TIME_BUCKETS - Prometheus::Client::Histogram.new(name.to_sym, docstring: "Histogram for #{name}", - labels: labels, buckets: buckets) - end - end - - def get_or_create_gauge(name) - @gauges[name] ||= register_or_reuse(name.to_sym) do - labels = GAUGE_LABELS.fetch(name, %i[taskType]) - Prometheus::Client::Gauge.new(name.to_sym, docstring: "Gauge for #{name}", labels: labels) - end - end - - # Align provided labels to the declared label set for the metric. - # Missing keys get empty-string defaults; unknown keys are dropped. - def normalize_labels(name, labels, schema) - symbolized = {} - labels.each do |key, value| - next if value.nil? - - symbolized[key.to_sym] = value.to_s - end - - declared = schema[name] - return symbolized unless declared - - declared.each_with_object({}) do |key, acc| - acc[key] = symbolized.key?(key) ? symbolized[key] : '' - end - end - end - end - end -end diff --git a/lib/conductor/worker/telemetry/legacy_metrics_collector.rb b/lib/conductor/worker/telemetry/legacy_metrics_collector.rb deleted file mode 100644 index b7714c5..0000000 --- a/lib/conductor/worker/telemetry/legacy_metrics_collector.rb +++ /dev/null @@ -1,116 +0,0 @@ -# frozen_string_literal: true - -require_relative '../events/listeners' - -module Conductor - module Worker - module Telemetry - # LegacyMetricsCollector - The original Ruby SDK metrics implementation. - # - # Emits the pre-harmonization metric set with snake_case labels (task_type, error). - # This is the default implementation during the deprecation period while - # WORKER_CANONICAL_METRICS defaults to false. - # - # Canonical-only event handlers (on_task_update_completed, on_task_paused, etc.) - # are implemented as no-ops so this collector satisfies the full listener interface - # and can be used interchangeably with CanonicalMetricsCollector. - class LegacyMetricsCollector - include Events::TaskRunnerEventsListener - include Events::WorkflowEventsListener - include Events::HttpEventsListener - - def initialize(backend: :null, measure_payload_size: false) - @backend = load_backend(backend) - @measure_payload_size = measure_payload_size - end - - attr_reader :backend, :measure_payload_size - - def collector_name - 'legacy' - end - - def stop; end - - # --- Real legacy metrics --- - - def on_poll_started(event) - @backend.increment('task_poll_total', labels: { task_type: event.task_type }) - end - - def on_poll_completed(event) - @backend.observe('task_poll_time_seconds', event.duration_ms / 1000.0, - labels: { task_type: event.task_type }) - end - - def on_poll_failure(event) - @backend.increment('task_poll_error_total', - labels: { - task_type: event.task_type, - error: event.cause.class.name - }) - end - - def on_task_execution_started(_event) - # No legacy metric for execution-started - end - - def on_task_execution_completed(event) - @backend.observe('task_execute_time_seconds', event.duration_ms / 1000.0, - labels: { task_type: event.task_type }) - - return unless event.output_size_bytes - - @backend.observe('task_result_size_bytes', event.output_size_bytes, - labels: { task_type: event.task_type }) - end - - def on_task_execution_failure(event) - @backend.increment('task_execute_error_total', - labels: { - task_type: event.task_type, - exception: event.cause.class.name, - retryable: event.is_retryable.to_s - }) - end - - def on_task_update_failure(event) - @backend.increment('task_update_failed_total', - labels: { task_type: event.task_type }) - end - - # --- No-op stubs for canonical-only events --- - - def on_task_update_completed(_event); end - def on_task_paused(_event); end - def on_thread_uncaught_exception(_event); end - def on_active_workers_changed(_event); end - def on_workflow_start_error(_event); end - def on_workflow_input_size(_event); end - def on_http_api_request(_event); end - - private - - def load_backend(backend) - case backend - when :null, nil - NullBackend.new - when :prometheus - load_prometheus_backend - else - backend - end - end - - def load_prometheus_backend - require_relative 'prometheus_backend' - PrometheusBackend.new - rescue LoadError - raise ConfigurationError, - "The 'prometheus-client' gem is required for Prometheus metrics. " \ - "Add `gem 'prometheus-client'` to your Gemfile." - end - end - end - end -end diff --git a/lib/conductor/worker/telemetry/metrics_collector.rb b/lib/conductor/worker/telemetry/metrics_collector.rb index cf9dbc9..e9d42e6 100644 --- a/lib/conductor/worker/telemetry/metrics_collector.rb +++ b/lib/conductor/worker/telemetry/metrics_collector.rb @@ -1,52 +1,178 @@ # frozen_string_literal: true -require_relative 'legacy_metrics_collector' -require_relative 'canonical_metrics_collector' +require 'logger' +require_relative '../events/listeners' +require_relative '../events/global_dispatcher' module Conductor module Worker module Telemetry - # MetricsCollector - Factory for creating the appropriate metrics - # collector based on environment configuration. + # MetricsCollector - Canonical SDK worker metrics from the + # harmonization spec (sdk-metrics-harmonization.md). # - # Currently checks WORKER_CANONICAL_METRICS (default false). When truthy, - # returns a CanonicalMetricsCollector; otherwise a LegacyMetricsCollector. - # - # In a future release, when canonical metrics become the default, - # WORKER_LEGACY_METRICS will be checked to allow opting back in to the - # legacy implementation. - module MetricsCollector - # Backward-compatible shim: delegates to .create so existing - # MetricsCollector.new(backend: ...) callers don't crash on upgrade. - def self.new(backend: :null, **_opts) - warn '[DEPRECATION] MetricsCollector.new is deprecated. Use MetricsCollector.create instead.' - create(backend: backend) - end - - # Create a metrics collector instance gated by environment configuration. - # - # @param backend [Symbol, Object] Backend type (:null, :prometheus) or custom backend - # @param subscribe_global_http [Boolean] Auto-subscribe to HTTP events (canonical only) - # @param measure_payload_size [Boolean, nil] Record workflow_input_size_bytes. - # Defaults to true for canonical, false for legacy. Set explicitly to override. + # Uses camelCase domain labels (taskType, workflowType) and includes + # status labels on time histograms. + class MetricsCollector + include Events::TaskRunnerEventsListener + include Events::WorkflowEventsListener + include Events::HttpEventsListener + + STATUS_SUCCESS = 'SUCCESS' + STATUS_FAILURE = 'FAILURE' + + # @param backend [Symbol, Object] :null, :prometheus, or a custom backend + # @param subscribe_global_http [Boolean] Auto-subscribe to GlobalDispatcher + # for HttpApiRequest events from the HTTP layer (default true). + # @param measure_payload_size [Boolean] Record workflow_input_size_bytes + # (requires JSON serialization; default true). Set false to skip + # serialization overhead for large payloads. # @param logger [Logger, nil] Optional logger for diagnostic output in rescue blocks - # @return [LegacyMetricsCollector, CanonicalMetricsCollector] - def self.create(backend: :null, subscribe_global_http: true, measure_payload_size: nil, logger: nil) - if canonical_metrics_enabled? - mps = measure_payload_size.nil? ? true : measure_payload_size - CanonicalMetricsCollector.new( - backend: backend, subscribe_global_http: subscribe_global_http, - measure_payload_size: mps, logger: logger - ) + # @return [MetricsCollector] + def self.create(backend: :null, subscribe_global_http: true, measure_payload_size: true, logger: nil) + new(backend: backend, subscribe_global_http: subscribe_global_http, + measure_payload_size: measure_payload_size, logger: logger) + end + + def initialize(backend: :null, subscribe_global_http: true, measure_payload_size: true, logger: nil) + @backend = load_backend(backend) + @logger = logger || Logger.new(File::NULL) + @measure_payload_size = measure_payload_size + @http_listener = nil + subscribe_to_global_http_events if subscribe_global_http + end + + attr_reader :backend, :measure_payload_size + + def stop + return unless @http_listener + + Events::GlobalDispatcher.instance.unregister(Events::HttpApiRequest, @http_listener) + @http_listener = nil + rescue StandardError => e + @logger.debug { "Telemetry error (non-fatal): #{e.class}: #{e.message}" } + end + + # --- Task Runner Event Handlers --- + + def on_poll_started(event) + @backend.increment('task_poll_total', labels: { taskType: event.task_type }) + end + + def on_poll_completed(event) + observe_time('task_poll_time_seconds', event.duration_ms, + { taskType: event.task_type, status: STATUS_SUCCESS }) + end + + def on_poll_failure(event) + @backend.increment('task_poll_error_total', + labels: { taskType: event.task_type, exception: event.cause.class.name }) + observe_time('task_poll_time_seconds', event.duration_ms, + { taskType: event.task_type, status: STATUS_FAILURE }) + end + + def on_task_execution_started(event) + @backend.increment('task_execution_started_total', labels: { taskType: event.task_type }) + end + + def on_task_execution_completed(event) + observe_time('task_execute_time_seconds', event.duration_ms, + { taskType: event.task_type, status: STATUS_SUCCESS }) + + return unless event.output_size_bytes + + @backend.observe('task_result_size_bytes', event.output_size_bytes, + labels: { taskType: event.task_type }) + end + + def on_task_execution_failure(event) + @backend.increment('task_execute_error_total', + labels: { taskType: event.task_type, exception: event.cause.class.name }) + observe_time('task_execute_time_seconds', event.duration_ms, + { taskType: event.task_type, status: STATUS_FAILURE }) + end + + def on_task_update_completed(event) + observe_time('task_update_time_seconds', event.duration_ms, + { taskType: event.task_type, status: STATUS_SUCCESS }) + end + + def on_task_update_failure(event) + @backend.increment('task_update_error_total', + labels: { taskType: event.task_type, exception: event.cause.class.name }) + + return unless event.respond_to?(:duration_ms) && event.duration_ms + + observe_time('task_update_time_seconds', event.duration_ms, + { taskType: event.task_type, status: STATUS_FAILURE }) + end + + def on_task_paused(event) + @backend.increment('task_paused_total', labels: { taskType: event.task_type }) + end + + def on_thread_uncaught_exception(event) + @backend.increment('thread_uncaught_exceptions_total', + labels: { exception: event.cause.class.name }) + end + + def on_active_workers_changed(event) + @backend.set('active_workers', event.count, labels: { taskType: event.task_type }) + end + + # --- Workflow Event Handlers --- + + def on_workflow_start_error(event) + @backend.increment('workflow_start_error_total', + labels: { workflowType: event.workflow_type, + exception: event.cause.class.name }) + end + + def on_workflow_input_size(event) + return unless @measure_payload_size + + @backend.observe('workflow_input_size_bytes', event.size_bytes, + labels: { workflowType: event.workflow_type, + version: (event.version || '').to_s }) + end + + # --- HTTP Event Handlers --- + + def on_http_api_request(event) + observe_time('http_api_client_request_seconds', event.duration_ms, + { method: event.method, uri: event.uri, status: event.status }) + end + + private + + def observe_time(name, duration_ms, labels) + @backend.observe(name, duration_ms / 1000.0, labels: labels) + end + + def subscribe_to_global_http_events + @http_listener = ->(event) { on_http_api_request(event) } + Events::GlobalDispatcher.instance.register(Events::HttpApiRequest, @http_listener) + rescue StandardError => e + @logger.debug { "Telemetry error (non-fatal): #{e.class}: #{e.message}" } + end + + def load_backend(backend) + case backend + when :null, nil + NullBackend.new + when :prometheus + load_prometheus_backend else - mps = measure_payload_size.nil? ? false : measure_payload_size - LegacyMetricsCollector.new(backend: backend, measure_payload_size: mps) + backend end end - # @return [Boolean] true when the canonical metric set is selected - def self.canonical_metrics_enabled? - %w[true 1 yes].include?(ENV.fetch('WORKER_CANONICAL_METRICS', 'false').downcase.strip) + def load_prometheus_backend + require_relative 'prometheus_backend' + PrometheusBackend.new + rescue LoadError + raise ConfigurationError, + "The 'prometheus-client' gem is required for Prometheus metrics. " \ + "Add `gem 'prometheus-client'` to your Gemfile." end end diff --git a/lib/conductor/worker/telemetry/prometheus_backend.rb b/lib/conductor/worker/telemetry/prometheus_backend.rb index 3c8bc75..43649ed 100644 --- a/lib/conductor/worker/telemetry/prometheus_backend.rb +++ b/lib/conductor/worker/telemetry/prometheus_backend.rb @@ -3,68 +3,72 @@ module Conductor module Worker module Telemetry - # PrometheusBackend - Prometheus metrics backend - # Uses the prometheus-client gem for metric collection + # PrometheusBackend - Prometheus backend for the canonical SDK metric catalog. # - # Metrics exposed: - # - task_poll_total (Counter) - # - task_poll_time_seconds (Histogram) - # - task_poll_error_total (Counter) - # - task_execute_time_seconds (Histogram) - # - task_execute_error_total (Counter) - # - task_result_size_bytes (Histogram) - # - task_update_failed_total (Counter) - # - # @example - # collector = MetricsCollector.create(backend: :prometheus) - # # Metrics available at default prometheus registry + # Pre-registers every metric from the harmonization spec with its canonical + # label set and bucket configuration. Uses camelCase domain labels (taskType, + # workflowType) per the canonical convention. class PrometheusBackend - # Default histogram buckets for time measurements (in seconds) - TIME_BUCKETS = [0.005, 0.01, 0.025, 0.05, 0.1, 0.25, 0.5, 1, 2.5, 5, 10].freeze - - # Default histogram buckets for size measurements (in bytes) + TIME_BUCKETS = [0.001, 0.005, 0.01, 0.025, 0.05, 0.1, 0.25, 0.5, 1, 2.5, 5, 10].freeze SIZE_BUCKETS = [100, 1000, 10_000, 100_000, 1_000_000, 10_000_000].freeze + COUNTER_LABELS = { + 'task_poll_total' => %i[taskType], + 'task_execution_started_total' => %i[taskType], + 'task_poll_error_total' => %i[taskType exception], + 'task_execute_error_total' => %i[taskType exception], + 'task_update_error_total' => %i[taskType exception], + 'task_paused_total' => %i[taskType], + 'thread_uncaught_exceptions_total' => %i[exception], + 'workflow_start_error_total' => %i[workflowType exception] + }.freeze + + HISTOGRAM_LABELS = { + 'task_poll_time_seconds' => %i[taskType status], + 'task_execute_time_seconds' => %i[taskType status], + 'task_update_time_seconds' => %i[taskType status], + 'http_api_client_request_seconds' => %i[method uri status], + 'task_result_size_bytes' => %i[taskType], + 'workflow_input_size_bytes' => %i[workflowType version] + }.freeze + + GAUGE_LABELS = { + 'active_workers' => %i[taskType] + }.freeze + + HISTOGRAM_BUCKETS = { + 'task_result_size_bytes' => SIZE_BUCKETS, + 'workflow_input_size_bytes' => SIZE_BUCKETS + }.freeze + def initialize(registry: nil) load_prometheus_client @registry = registry || Prometheus::Client.registry + @counters = {} + @histograms = {} + @gauges = {} setup_metrics end - # Increment a counter - # @param name [String] Metric name - # @param labels [Hash] Metric labels - # @param value [Integer] Value to increment by (default: 1) def increment(name, labels: {}, value: 1) metric = get_or_create_counter(name) - metric.increment(labels: normalize_labels(labels), by: value) + metric.increment(labels: normalize_labels(name, labels, COUNTER_LABELS), by: value) end - # Observe a value in a histogram - # @param name [String] Metric name - # @param value [Numeric] Value to observe - # @param labels [Hash] Metric labels def observe(name, value, labels: {}) metric = get_or_create_histogram(name) - metric.observe(value, labels: normalize_labels(labels)) + metric.observe(value, labels: normalize_labels(name, labels, HISTOGRAM_LABELS)) end - # Set a gauge value - # @param name [String] Metric name - # @param value [Numeric] Value to set - # @param labels [Hash] Metric labels def set(name, value, labels: {}) metric = get_or_create_gauge(name) - metric.set(value, labels: normalize_labels(labels)) + metric.set(value, labels: normalize_labels(name, labels, GAUGE_LABELS)) end - # Get the prometheus registry - # @return [Prometheus::Client::Registry] attr_reader :registry private - # Load prometheus-client gem def load_prometheus_client require 'prometheus/client' rescue LoadError @@ -73,159 +77,95 @@ def load_prometheus_client "Add `gem 'prometheus-client'` to your Gemfile." end - # Each counter declares only the labels it actually receives - COUNTER_LABELS = { - 'task_poll_total' => %i[task_type], - 'task_poll_error_total' => %i[task_type error], - 'task_execute_error_total' => %i[task_type exception retryable], - 'task_update_failed_total' => %i[task_type] - }.freeze - - # Setup predefined metrics def setup_metrics - @counters = {} - @histograms = {} - @gauges = {} + COUNTER_LABELS.each do |name, _| + register_counter(name, "Counter for #{name}") + end - register_counter('task_poll_total', 'Total number of task polls', COUNTER_LABELS['task_poll_total']) - register_counter('task_poll_error_total', 'Total number of poll errors', COUNTER_LABELS['task_poll_error_total']) - register_counter('task_execute_error_total', 'Total number of execution errors', COUNTER_LABELS['task_execute_error_total']) - register_counter('task_update_failed_total', 'Total number of failed task updates (CRITICAL)', COUNTER_LABELS['task_update_failed_total']) + HISTOGRAM_LABELS.each do |name, _| + register_histogram(name, "Histogram for #{name}") + end - register_histogram('task_poll_time_seconds', 'Task poll duration in seconds', TIME_BUCKETS) - register_histogram('task_execute_time_seconds', 'Task execution duration in seconds', TIME_BUCKETS) - register_histogram('task_result_size_bytes', 'Task result size in bytes', SIZE_BUCKETS) + GAUGE_LABELS.each do |name, _| + register_gauge(name, "Gauge for #{name}") + end end - # Register a counter metric - # @param name [String] Metric name - # @param docstring [String] Metric description - # @param labels [Array] Label keys for this counter - def register_counter(name, docstring, labels = %i[task_type]) + def register_counter(name, docstring) metric_name = name.to_sym - return if @registry.exist?(metric_name) - - counter = Prometheus::Client::Counter.new( - metric_name, - docstring: docstring, - labels: labels - ) - @registry.register(counter) - @counters[name] = counter + labels = COUNTER_LABELS.fetch(name, %i[taskType]) + @counters[name] = register_or_reuse(metric_name) do + Prometheus::Client::Counter.new(metric_name, docstring: docstring, labels: labels) + end end - # Register a histogram metric - # @param name [String] Metric name - # @param docstring [String] Metric description - # @param buckets [Array] Histogram buckets - def register_histogram(name, docstring, buckets) + def register_histogram(name, docstring) metric_name = name.to_sym - return if @registry.exist?(metric_name) - - histogram = Prometheus::Client::Histogram.new( - metric_name, - docstring: docstring, - labels: [:task_type], - buckets: buckets - ) - @registry.register(histogram) - @histograms[name] = histogram + labels = HISTOGRAM_LABELS.fetch(name, %i[taskType]) + buckets = HISTOGRAM_BUCKETS[name] || TIME_BUCKETS + @histograms[name] = register_or_reuse(metric_name) do + Prometheus::Client::Histogram.new(metric_name, docstring: docstring, + labels: labels, buckets: buckets) + end end - # Register a gauge metric - # @param name [String] Metric name - # @param docstring [String] Metric description def register_gauge(name, docstring) metric_name = name.to_sym - return if @registry.exist?(metric_name) + labels = GAUGE_LABELS.fetch(name, %i[taskType]) + @gauges[name] = register_or_reuse(metric_name) do + Prometheus::Client::Gauge.new(metric_name, docstring: docstring, labels: labels) + end + end - gauge = Prometheus::Client::Gauge.new( - metric_name, - docstring: docstring, - labels: [:task_type] - ) - @registry.register(gauge) - @gauges[name] = gauge + def register_or_reuse(metric_name) + if @registry.exist?(metric_name) + @registry.get(metric_name) + else + metric = yield + @registry.register(metric) + metric + end end - # Get or create a counter metric - # @param name [String] Metric name - # @return [Prometheus::Client::Counter] def get_or_create_counter(name) - @counters[name] ||= begin - metric_name = name.to_sym - if @registry.exist?(metric_name) - @registry.get(metric_name) - else - labels = COUNTER_LABELS.fetch(name, %i[task_type]) - counter = Prometheus::Client::Counter.new( - metric_name, - docstring: "Counter for #{name}", - labels: labels - ) - @registry.register(counter) - counter - end + @counters[name] ||= register_or_reuse(name.to_sym) do + labels = COUNTER_LABELS.fetch(name, %i[taskType]) + Prometheus::Client::Counter.new(name.to_sym, docstring: "Counter for #{name}", labels: labels) end end - # Get or create a histogram metric - # @param name [String] Metric name - # @return [Prometheus::Client::Histogram] def get_or_create_histogram(name) - @histograms[name] ||= begin - metric_name = name.to_sym - if @registry.exist?(metric_name) - @registry.get(metric_name) - else - buckets = name.include?('bytes') ? SIZE_BUCKETS : TIME_BUCKETS - histogram = Prometheus::Client::Histogram.new( - metric_name, - docstring: "Histogram for #{name}", - labels: [:task_type], - buckets: buckets - ) - @registry.register(histogram) - histogram - end + @histograms[name] ||= register_or_reuse(name.to_sym) do + labels = HISTOGRAM_LABELS.fetch(name, %i[taskType]) + buckets = HISTOGRAM_BUCKETS[name] || TIME_BUCKETS + Prometheus::Client::Histogram.new(name.to_sym, docstring: "Histogram for #{name}", + labels: labels, buckets: buckets) end end - # Get or create a gauge metric - # @param name [String] Metric name - # @return [Prometheus::Client::Gauge] def get_or_create_gauge(name) - @gauges[name] ||= begin - metric_name = name.to_sym - if @registry.exist?(metric_name) - @registry.get(metric_name) - else - gauge = Prometheus::Client::Gauge.new( - metric_name, - docstring: "Gauge for #{name}", - labels: [:task_type] - ) - @registry.register(gauge) - gauge - end + @gauges[name] ||= register_or_reuse(name.to_sym) do + labels = GAUGE_LABELS.fetch(name, %i[taskType]) + Prometheus::Client::Gauge.new(name.to_sym, docstring: "Gauge for #{name}", labels: labels) end end - # Normalize labels - convert keys to symbols and filter out nil/empty values - # @param labels [Hash] Input labels - # @return [Hash] Normalized labels - def normalize_labels(labels) - result = {} + # Align provided labels to the declared label set for the metric. + # Missing keys get empty-string defaults; unknown keys are dropped. + def normalize_labels(name, labels, schema) + symbolized = {} labels.each do |key, value| next if value.nil? - sym_key = key.to_sym - result[sym_key] = value.to_s + symbolized[key.to_sym] = value.to_s end - # Ensure required labels have default values - result[:task_type] ||= 'unknown' - result + declared = schema[name] + return symbolized unless declared + + declared.each_with_object({}) do |key, acc| + acc[key] = symbolized.key?(key) ? symbolized[key] : '' + end end end @@ -234,7 +174,6 @@ def normalize_labels(labels) class MetricsServer DEFAULT_PORT = 9090 - # Initialize metrics server # @param port [Integer] Port to listen on (default: 9090) # @param registry [Prometheus::Client::Registry] Prometheus registry def initialize(port: DEFAULT_PORT, registry: nil) diff --git a/spec/conductor/worker/canonical_metrics_collector_spec.rb b/spec/conductor/worker/canonical_metrics_collector_spec.rb deleted file mode 100644 index 8ba0ef2..0000000 --- a/spec/conductor/worker/canonical_metrics_collector_spec.rb +++ /dev/null @@ -1,271 +0,0 @@ -# frozen_string_literal: true - -require 'spec_helper' -require_relative '../../../lib/conductor/worker/telemetry/metrics_collector' - -RSpec.describe Conductor::Worker::Telemetry::CanonicalMetricsCollector do - let(:backend) { double('backend') } - let(:collector) { described_class.new(backend: backend, subscribe_global_http: false) } - - before do - allow(backend).to receive(:increment) - allow(backend).to receive(:observe) - allow(backend).to receive(:set) - end - - describe '#initialize' do - it 'uses NullBackend by default' do - c = described_class.new(subscribe_global_http: false) - expect(c.backend).to be_a(Conductor::Worker::Telemetry::NullBackend) - end - end - - describe '#collector_name' do - it 'returns "canonical"' do - expect(collector.collector_name).to eq('canonical') - end - end - - # --- Task Runner Events --- - - describe '#on_poll_started' do - it 'increments task_poll_total with camelCase taskType label' do - event = Conductor::Worker::Events::PollStarted.new( - task_type: 'my_task', worker_id: 'w1', poll_count: 1 - ) - collector.on_poll_started(event) - expect(backend).to have_received(:increment).with( - 'task_poll_total', labels: { taskType: 'my_task' } - ) - end - end - - describe '#on_poll_completed' do - it 'observes task_poll_time_seconds with status=SUCCESS' do - event = Conductor::Worker::Events::PollCompleted.new( - task_type: 'my_task', duration_ms: 250.0, tasks_received: 2 - ) - collector.on_poll_completed(event) - expect(backend).to have_received(:observe).with( - 'task_poll_time_seconds', 0.25, - labels: { taskType: 'my_task', status: 'SUCCESS' } - ) - end - end - - describe '#on_poll_failure' do - it 'increments task_poll_error_total with exception label' do - error = Timeout::Error.new('timed out') - event = Conductor::Worker::Events::PollFailure.new( - task_type: 'my_task', duration_ms: 100.0, cause: error - ) - collector.on_poll_failure(event) - expect(backend).to have_received(:increment).with( - 'task_poll_error_total', - labels: { taskType: 'my_task', exception: 'Timeout::Error' } - ) - end - - it 'observes task_poll_time_seconds with status=FAILURE' do - error = StandardError.new('nope') - event = Conductor::Worker::Events::PollFailure.new( - task_type: 'my_task', duration_ms: 80.0, cause: error - ) - collector.on_poll_failure(event) - expect(backend).to have_received(:observe).with( - 'task_poll_time_seconds', 0.08, - labels: { taskType: 'my_task', status: 'FAILURE' } - ) - end - end - - describe '#on_task_execution_started' do - it 'increments task_execution_started_total' do - event = Conductor::Worker::Events::TaskExecutionStarted.new( - task_type: 'my_task', task_id: 't1', worker_id: 'w1', workflow_instance_id: 'wf1' - ) - collector.on_task_execution_started(event) - expect(backend).to have_received(:increment).with( - 'task_execution_started_total', labels: { taskType: 'my_task' } - ) - end - end - - describe '#on_task_execution_completed' do - it 'observes task_execute_time_seconds with status=SUCCESS' do - event = Conductor::Worker::Events::TaskExecutionCompleted.new( - task_type: 'my_task', task_id: 't1', worker_id: 'w1', - workflow_instance_id: 'wf1', duration_ms: 1200.0, output_size_bytes: 4096 - ) - collector.on_task_execution_completed(event) - expect(backend).to have_received(:observe).with( - 'task_execute_time_seconds', 1.2, - labels: { taskType: 'my_task', status: 'SUCCESS' } - ) - end - - it 'observes task_result_size_bytes as histogram with taskType label' do - event = Conductor::Worker::Events::TaskExecutionCompleted.new( - task_type: 'my_task', task_id: 't1', worker_id: 'w1', - workflow_instance_id: 'wf1', duration_ms: 500.0, output_size_bytes: 4096 - ) - collector.on_task_execution_completed(event) - expect(backend).to have_received(:observe).with( - 'task_result_size_bytes', 4096, labels: { taskType: 'my_task' } - ) - end - - it 'skips task_result_size_bytes when output_size_bytes is nil' do - event = Conductor::Worker::Events::TaskExecutionCompleted.new( - task_type: 'my_task', task_id: 't1', worker_id: 'w1', - workflow_instance_id: 'wf1', duration_ms: 500.0 - ) - collector.on_task_execution_completed(event) - expect(backend).to have_received(:observe).once - end - end - - describe '#on_task_execution_failure' do - it 'increments task_execute_error_total and observes time with FAILURE' do - error = ArgumentError.new('bad input') - event = Conductor::Worker::Events::TaskExecutionFailure.new( - task_type: 'my_task', task_id: 't1', worker_id: 'w1', - workflow_instance_id: 'wf1', duration_ms: 300.0, cause: error, is_retryable: true - ) - collector.on_task_execution_failure(event) - - expect(backend).to have_received(:increment).with( - 'task_execute_error_total', - labels: { taskType: 'my_task', exception: 'ArgumentError' } - ) - expect(backend).to have_received(:observe).with( - 'task_execute_time_seconds', 0.3, - labels: { taskType: 'my_task', status: 'FAILURE' } - ) - end - end - - describe '#on_task_update_completed' do - it 'observes task_update_time_seconds with status=SUCCESS' do - event = Conductor::Worker::Events::TaskUpdateCompleted.new( - task_type: 'my_task', task_id: 't1', worker_id: 'w1', - workflow_instance_id: 'wf1', duration_ms: 75.0 - ) - collector.on_task_update_completed(event) - expect(backend).to have_received(:observe).with( - 'task_update_time_seconds', 0.075, - labels: { taskType: 'my_task', status: 'SUCCESS' } - ) - end - end - - describe '#on_task_update_failure' do - it 'increments task_update_error_total and observes time with FAILURE' do - error = StandardError.new('net err') - task_result = Conductor::Http::Models::TaskResult.complete - event = Conductor::Worker::Events::TaskUpdateFailure.new( - task_type: 'my_task', task_id: 't1', worker_id: 'w1', - workflow_instance_id: 'wf1', cause: error, retry_count: 4, - task_result: task_result, duration_ms: 120.0 - ) - collector.on_task_update_failure(event) - - expect(backend).to have_received(:increment).with( - 'task_update_error_total', - labels: { taskType: 'my_task', exception: 'StandardError' } - ) - expect(backend).to have_received(:observe).with( - 'task_update_time_seconds', 0.12, - labels: { taskType: 'my_task', status: 'FAILURE' } - ) - end - - it 'skips time observation when duration_ms is nil' do - error = StandardError.new('err') - task_result = Conductor::Http::Models::TaskResult.complete - event = Conductor::Worker::Events::TaskUpdateFailure.new( - task_type: 'my_task', task_id: 't1', worker_id: 'w1', - workflow_instance_id: 'wf1', cause: error, retry_count: 4, task_result: task_result - ) - collector.on_task_update_failure(event) - expect(backend).not_to have_received(:observe) - end - end - - describe '#on_task_paused' do - it 'increments task_paused_total' do - event = Conductor::Worker::Events::TaskPaused.new(task_type: 'my_task') - collector.on_task_paused(event) - expect(backend).to have_received(:increment).with( - 'task_paused_total', labels: { taskType: 'my_task' } - ) - end - end - - describe '#on_thread_uncaught_exception' do - it 'increments thread_uncaught_exceptions_total with exception label' do - event = Conductor::Worker::Events::ThreadUncaughtException.new( - cause: RuntimeError.new('boom') - ) - collector.on_thread_uncaught_exception(event) - expect(backend).to have_received(:increment).with( - 'thread_uncaught_exceptions_total', labels: { exception: 'RuntimeError' } - ) - end - end - - describe '#on_active_workers_changed' do - it 'sets active_workers gauge' do - event = Conductor::Worker::Events::ActiveWorkersChanged.new( - task_type: 'my_task', count: 7 - ) - collector.on_active_workers_changed(event) - expect(backend).to have_received(:set).with( - 'active_workers', 7, labels: { taskType: 'my_task' } - ) - end - end - - # --- Workflow Events --- - - describe '#on_workflow_start_error' do - it 'increments workflow_start_error_total' do - event = Conductor::Worker::Events::WorkflowStartError.new( - workflow_type: 'my_wf', cause: RuntimeError.new('fail') - ) - collector.on_workflow_start_error(event) - expect(backend).to have_received(:increment).with( - 'workflow_start_error_total', - labels: { workflowType: 'my_wf', exception: 'RuntimeError' } - ) - end - end - - describe '#on_workflow_input_size' do - it 'observes workflow_input_size_bytes' do - event = Conductor::Worker::Events::WorkflowInputSize.new( - workflow_type: 'my_wf', size_bytes: 8192, version: 2 - ) - collector.on_workflow_input_size(event) - expect(backend).to have_received(:observe).with( - 'workflow_input_size_bytes', 8192, - labels: { workflowType: 'my_wf', version: '2' } - ) - end - end - - # --- HTTP Events --- - - describe '#on_http_api_request' do - it 'observes http_api_client_request_seconds' do - event = Conductor::Worker::Events::HttpApiRequest.new( - method: 'POST', uri: '/api/tasks/poll/batch/my_task', status: '200', duration_ms: 45.0 - ) - collector.on_http_api_request(event) - expect(backend).to have_received(:observe).with( - 'http_api_client_request_seconds', 0.045, - labels: { method: 'POST', uri: '/api/tasks/poll/batch/my_task', status: '200' } - ) - end - end -end diff --git a/spec/conductor/worker/canonical_prometheus_backend_spec.rb b/spec/conductor/worker/canonical_prometheus_backend_spec.rb deleted file mode 100644 index 668d0ac..0000000 --- a/spec/conductor/worker/canonical_prometheus_backend_spec.rb +++ /dev/null @@ -1,104 +0,0 @@ -# frozen_string_literal: true - -require 'spec_helper' - -CANONICAL_PROMETHEUS_AVAILABLE = begin - require 'prometheus/client' - true -rescue LoadError - false -end - -CANONICAL_BACKEND_LOADED = begin - if CANONICAL_PROMETHEUS_AVAILABLE - require_relative '../../../lib/conductor/worker/telemetry/canonical_prometheus_backend' - true - else - false - end -rescue Conductor::ConfigurationError - false -end - -if CANONICAL_BACKEND_LOADED - RSpec.describe Conductor::Worker::Telemetry::CanonicalPrometheusBackend do - let(:registry) { Prometheus::Client::Registry.new } - let(:backend) { described_class.new(registry: registry) } - - describe '#initialize' do - it 'registers canonical counters' do - backend # force lazy initialization - %i[task_poll_total task_execution_started_total task_poll_error_total - task_execute_error_total task_update_error_total task_paused_total - thread_uncaught_exceptions_total workflow_start_error_total].each do |name| - expect(registry.exist?(name)).to be(true), "Expected counter #{name} to be registered" - end - end - - it 'registers canonical histograms' do - backend - %i[task_poll_time_seconds task_execute_time_seconds task_update_time_seconds - http_api_client_request_seconds task_result_size_bytes - workflow_input_size_bytes].each do |name| - expect(registry.exist?(name)).to be(true), "Expected histogram #{name} to be registered" - end - end - - it 'registers canonical gauges' do - backend - expect(registry.exist?(:active_workers)).to be true - end - end - - describe '#increment' do - it 'increments a counter with camelCase labels' do - expect do - backend.increment('task_poll_total', labels: { taskType: 'my_task' }) - end.not_to raise_error - end - end - - describe '#observe' do - it 'observes a time histogram with status label' do - expect do - backend.observe('task_poll_time_seconds', 0.25, - labels: { taskType: 'my_task', status: 'SUCCESS' }) - end.not_to raise_error - end - - it 'observes a size histogram' do - expect do - backend.observe('task_result_size_bytes', 5000, labels: { taskType: 'my_task' }) - end.not_to raise_error - end - end - - describe '#set' do - it 'sets a gauge value' do - expect do - backend.set('active_workers', 3, labels: { taskType: 'my_task' }) - end.not_to raise_error - end - end - - describe 'label normalization' do - it 'fills missing declared labels with empty strings' do - expect do - backend.increment('task_poll_error_total', labels: { taskType: 'my_task' }) - end.not_to raise_error - end - - it 'drops undeclared labels' do - expect do - backend.increment('task_poll_total', labels: { taskType: 'my_task', extra: 'nope' }) - end.not_to raise_error - end - end - end -else - RSpec.describe 'CanonicalPrometheusBackend (prometheus-client gem unavailable)' do - it 'documents that prometheus-client gem is not installed' do - expect(CANONICAL_PROMETHEUS_AVAILABLE).to be false - end - end -end diff --git a/spec/conductor/worker/legacy_metrics_collector_spec.rb b/spec/conductor/worker/legacy_metrics_collector_spec.rb deleted file mode 100644 index 418fdc9..0000000 --- a/spec/conductor/worker/legacy_metrics_collector_spec.rb +++ /dev/null @@ -1,192 +0,0 @@ -# frozen_string_literal: true - -require 'spec_helper' -require_relative '../../../lib/conductor/worker/telemetry/metrics_collector' - -RSpec.describe Conductor::Worker::Telemetry::LegacyMetricsCollector do - let(:backend) { double('backend') } - let(:collector) { described_class.new(backend: backend) } - - before do - allow(backend).to receive(:increment) - allow(backend).to receive(:observe) - allow(backend).to receive(:set) - end - - describe '#initialize' do - it 'uses NullBackend by default' do - collector = described_class.new - expect(collector.backend).to be_a(Conductor::Worker::Telemetry::NullBackend) - end - - it 'accepts a custom backend instance' do - custom = Object.new - collector = described_class.new(backend: custom) - expect(collector.backend).to eq(custom) - end - end - - describe '#collector_name' do - it 'returns "legacy"' do - expect(collector.collector_name).to eq('legacy') - end - end - - describe '#on_poll_started' do - it 'increments task_poll_total with snake_case task_type' do - event = Conductor::Worker::Events::PollStarted.new( - task_type: 'my_task', worker_id: 'w1', poll_count: 1 - ) - collector.on_poll_started(event) - expect(backend).to have_received(:increment).with( - 'task_poll_total', labels: { task_type: 'my_task' } - ) - end - end - - describe '#on_poll_completed' do - it 'observes task_poll_time_seconds without status label' do - event = Conductor::Worker::Events::PollCompleted.new( - task_type: 'my_task', duration_ms: 150.0, tasks_received: 3 - ) - collector.on_poll_completed(event) - expect(backend).to have_received(:observe).with( - 'task_poll_time_seconds', 0.15, labels: { task_type: 'my_task' } - ) - end - end - - describe '#on_poll_failure' do - it 'increments task_poll_error_total with error label' do - error = StandardError.new('timeout') - event = Conductor::Worker::Events::PollFailure.new( - task_type: 'my_task', duration_ms: 100.0, cause: error - ) - collector.on_poll_failure(event) - expect(backend).to have_received(:increment).with( - 'task_poll_error_total', labels: { task_type: 'my_task', error: 'StandardError' } - ) - end - end - - describe '#on_task_execution_started' do - it 'is a no-op (no legacy metric)' do - event = Conductor::Worker::Events::TaskExecutionStarted.new( - task_type: 'my_task', task_id: 't1', worker_id: 'w1', workflow_instance_id: 'wf1' - ) - expect { collector.on_task_execution_started(event) }.not_to raise_error - expect(backend).not_to have_received(:increment) - end - end - - describe '#on_task_execution_completed' do - it 'observes task_execute_time_seconds' do - event = Conductor::Worker::Events::TaskExecutionCompleted.new( - task_type: 'my_task', task_id: 't1', worker_id: 'w1', - workflow_instance_id: 'wf1', duration_ms: 500.0, output_size_bytes: 2048 - ) - collector.on_task_execution_completed(event) - expect(backend).to have_received(:observe).with( - 'task_execute_time_seconds', 0.5, labels: { task_type: 'my_task' } - ) - end - - it 'observes task_result_size_bytes when output_size_bytes is present' do - event = Conductor::Worker::Events::TaskExecutionCompleted.new( - task_type: 'my_task', task_id: 't1', worker_id: 'w1', - workflow_instance_id: 'wf1', duration_ms: 500.0, output_size_bytes: 2048 - ) - collector.on_task_execution_completed(event) - expect(backend).to have_received(:observe).with( - 'task_result_size_bytes', 2048, labels: { task_type: 'my_task' } - ) - end - - it 'skips task_result_size_bytes when output_size_bytes is nil' do - event = Conductor::Worker::Events::TaskExecutionCompleted.new( - task_type: 'my_task', task_id: 't1', worker_id: 'w1', - workflow_instance_id: 'wf1', duration_ms: 500.0, output_size_bytes: nil - ) - collector.on_task_execution_completed(event) - expect(backend).to have_received(:observe).once - end - end - - describe '#on_task_execution_failure' do - it 'increments task_execute_error_total with exception and retryable labels' do - error = ArgumentError.new('bad') - event = Conductor::Worker::Events::TaskExecutionFailure.new( - task_type: 'my_task', task_id: 't1', worker_id: 'w1', - workflow_instance_id: 'wf1', duration_ms: 100.0, cause: error, is_retryable: true - ) - collector.on_task_execution_failure(event) - expect(backend).to have_received(:increment).with( - 'task_execute_error_total', - labels: { task_type: 'my_task', exception: 'ArgumentError', retryable: 'true' } - ) - end - end - - describe '#on_task_update_failure' do - it 'increments task_update_failed_total' do - error = StandardError.new('net error') - task_result = Conductor::Http::Models::TaskResult.complete - event = Conductor::Worker::Events::TaskUpdateFailure.new( - task_type: 'my_task', task_id: 't1', worker_id: 'w1', - workflow_instance_id: 'wf1', cause: error, retry_count: 4, task_result: task_result - ) - collector.on_task_update_failure(event) - expect(backend).to have_received(:increment).with( - 'task_update_failed_total', labels: { task_type: 'my_task' } - ) - end - end - - # Canonical-only events should be no-ops - describe 'canonical-only stubs' do - it 'does not raise on on_task_update_completed' do - event = Conductor::Worker::Events::TaskUpdateCompleted.new( - task_type: 'my_task', task_id: 't1', worker_id: 'w1', - workflow_instance_id: 'wf1', duration_ms: 50.0 - ) - expect { collector.on_task_update_completed(event) }.not_to raise_error - expect(backend).not_to have_received(:observe) - end - - it 'does not raise on on_task_paused' do - event = Conductor::Worker::Events::TaskPaused.new(task_type: 'my_task') - expect { collector.on_task_paused(event) }.not_to raise_error - end - - it 'does not raise on on_thread_uncaught_exception' do - event = Conductor::Worker::Events::ThreadUncaughtException.new(cause: RuntimeError.new('boom')) - expect { collector.on_thread_uncaught_exception(event) }.not_to raise_error - end - - it 'does not raise on on_active_workers_changed' do - event = Conductor::Worker::Events::ActiveWorkersChanged.new(task_type: 'my_task', count: 3) - expect { collector.on_active_workers_changed(event) }.not_to raise_error - end - - it 'does not raise on on_workflow_start_error' do - event = Conductor::Worker::Events::WorkflowStartError.new( - workflow_type: 'my_wf', cause: RuntimeError.new('fail') - ) - expect { collector.on_workflow_start_error(event) }.not_to raise_error - end - - it 'does not raise on on_workflow_input_size' do - event = Conductor::Worker::Events::WorkflowInputSize.new( - workflow_type: 'my_wf', size_bytes: 1024 - ) - expect { collector.on_workflow_input_size(event) }.not_to raise_error - end - - it 'does not raise on on_http_api_request' do - event = Conductor::Worker::Events::HttpApiRequest.new( - method: 'GET', uri: '/api/tasks', status: '200', duration_ms: 50.0 - ) - expect { collector.on_http_api_request(event) }.not_to raise_error - end - end -end diff --git a/spec/conductor/worker/metrics_collector_spec.rb b/spec/conductor/worker/metrics_collector_spec.rb index 37b8f06..42150f7 100644 --- a/spec/conductor/worker/metrics_collector_spec.rb +++ b/spec/conductor/worker/metrics_collector_spec.rb @@ -4,98 +4,274 @@ require_relative '../../../lib/conductor/worker/telemetry/metrics_collector' RSpec.describe Conductor::Worker::Telemetry::MetricsCollector do + let(:backend) { double('backend') } + let(:collector) { described_class.new(backend: backend, subscribe_global_http: false) } + + before do + allow(backend).to receive(:increment) + allow(backend).to receive(:observe) + allow(backend).to receive(:set) + end + describe '.create' do - around do |example| - old_val = ENV.fetch('WORKER_CANONICAL_METRICS', nil) - example.run - ensure - if old_val.nil? - ENV.delete('WORKER_CANONICAL_METRICS') - else - ENV['WORKER_CANONICAL_METRICS'] = old_val - end + it 'returns a MetricsCollector' do + collector = described_class.create(subscribe_global_http: false) + expect(collector).to be_a(described_class) end - it 'returns a LegacyMetricsCollector by default' do - ENV.delete('WORKER_CANONICAL_METRICS') - collector = described_class.create - expect(collector).to be_a(Conductor::Worker::Telemetry::LegacyMetricsCollector) + it 'passes backend option through' do + collector = described_class.create(backend: :null, subscribe_global_http: false) + expect(collector.backend).to be_a(Conductor::Worker::Telemetry::NullBackend) end + end - it 'returns a LegacyMetricsCollector when WORKER_CANONICAL_METRICS is false' do - ENV['WORKER_CANONICAL_METRICS'] = 'false' - collector = described_class.create - expect(collector).to be_a(Conductor::Worker::Telemetry::LegacyMetricsCollector) + describe '#initialize' do + it 'uses NullBackend by default' do + c = described_class.new(subscribe_global_http: false) + expect(c.backend).to be_a(Conductor::Worker::Telemetry::NullBackend) end + end - it 'returns a CanonicalMetricsCollector when WORKER_CANONICAL_METRICS is true' do - ENV['WORKER_CANONICAL_METRICS'] = 'true' - collector = described_class.create(subscribe_global_http: false) - expect(collector).to be_a(Conductor::Worker::Telemetry::CanonicalMetricsCollector) + # --- Task Runner Events --- + + describe '#on_poll_started' do + it 'increments task_poll_total with camelCase taskType label' do + event = Conductor::Worker::Events::PollStarted.new( + task_type: 'my_task', worker_id: 'w1', poll_count: 1 + ) + collector.on_poll_started(event) + expect(backend).to have_received(:increment).with( + 'task_poll_total', labels: { taskType: 'my_task' } + ) end + end - it 'accepts "1" as truthy for WORKER_CANONICAL_METRICS' do - ENV['WORKER_CANONICAL_METRICS'] = '1' - collector = described_class.create(subscribe_global_http: false) - expect(collector).to be_a(Conductor::Worker::Telemetry::CanonicalMetricsCollector) + describe '#on_poll_completed' do + it 'observes task_poll_time_seconds with status=SUCCESS' do + event = Conductor::Worker::Events::PollCompleted.new( + task_type: 'my_task', duration_ms: 250.0, tasks_received: 2 + ) + collector.on_poll_completed(event) + expect(backend).to have_received(:observe).with( + 'task_poll_time_seconds', 0.25, + labels: { taskType: 'my_task', status: 'SUCCESS' } + ) end + end - it 'accepts "yes" as truthy for WORKER_CANONICAL_METRICS' do - ENV['WORKER_CANONICAL_METRICS'] = 'yes' - collector = described_class.create(subscribe_global_http: false) - expect(collector).to be_a(Conductor::Worker::Telemetry::CanonicalMetricsCollector) + describe '#on_poll_failure' do + it 'increments task_poll_error_total with exception label' do + error = Timeout::Error.new('timed out') + event = Conductor::Worker::Events::PollFailure.new( + task_type: 'my_task', duration_ms: 100.0, cause: error + ) + collector.on_poll_failure(event) + expect(backend).to have_received(:increment).with( + 'task_poll_error_total', + labels: { taskType: 'my_task', exception: 'Timeout::Error' } + ) end - it 'is case-insensitive for WORKER_CANONICAL_METRICS' do - ENV['WORKER_CANONICAL_METRICS'] = 'TRUE' - collector = described_class.create(subscribe_global_http: false) - expect(collector).to be_a(Conductor::Worker::Telemetry::CanonicalMetricsCollector) + it 'observes task_poll_time_seconds with status=FAILURE' do + error = StandardError.new('nope') + event = Conductor::Worker::Events::PollFailure.new( + task_type: 'my_task', duration_ms: 80.0, cause: error + ) + collector.on_poll_failure(event) + expect(backend).to have_received(:observe).with( + 'task_poll_time_seconds', 0.08, + labels: { taskType: 'my_task', status: 'FAILURE' } + ) end + end - it 'passes backend option through to the collector' do - ENV.delete('WORKER_CANONICAL_METRICS') - collector = described_class.create(backend: :null) - expect(collector.backend).to be_a(Conductor::Worker::Telemetry::NullBackend) + describe '#on_task_execution_started' do + it 'increments task_execution_started_total' do + event = Conductor::Worker::Events::TaskExecutionStarted.new( + task_type: 'my_task', task_id: 't1', worker_id: 'w1', workflow_instance_id: 'wf1' + ) + collector.on_task_execution_started(event) + expect(backend).to have_received(:increment).with( + 'task_execution_started_total', labels: { taskType: 'my_task' } + ) end + end - it 'legacy collector returns "legacy" from collector_name' do - ENV.delete('WORKER_CANONICAL_METRICS') - collector = described_class.create - expect(collector.collector_name).to eq('legacy') + describe '#on_task_execution_completed' do + it 'observes task_execute_time_seconds with status=SUCCESS' do + event = Conductor::Worker::Events::TaskExecutionCompleted.new( + task_type: 'my_task', task_id: 't1', worker_id: 'w1', + workflow_instance_id: 'wf1', duration_ms: 1200.0, output_size_bytes: 4096 + ) + collector.on_task_execution_completed(event) + expect(backend).to have_received(:observe).with( + 'task_execute_time_seconds', 1.2, + labels: { taskType: 'my_task', status: 'SUCCESS' } + ) end - it 'canonical collector returns "canonical" from collector_name' do - ENV['WORKER_CANONICAL_METRICS'] = 'true' - collector = described_class.create(subscribe_global_http: false) - expect(collector.collector_name).to eq('canonical') + it 'observes task_result_size_bytes as histogram with taskType label' do + event = Conductor::Worker::Events::TaskExecutionCompleted.new( + task_type: 'my_task', task_id: 't1', worker_id: 'w1', + workflow_instance_id: 'wf1', duration_ms: 500.0, output_size_bytes: 4096 + ) + collector.on_task_execution_completed(event) + expect(backend).to have_received(:observe).with( + 'task_result_size_bytes', 4096, labels: { taskType: 'my_task' } + ) + end + + it 'skips task_result_size_bytes when output_size_bytes is nil' do + event = Conductor::Worker::Events::TaskExecutionCompleted.new( + task_type: 'my_task', task_id: 't1', worker_id: 'w1', + workflow_instance_id: 'wf1', duration_ms: 500.0 + ) + collector.on_task_execution_completed(event) + expect(backend).to have_received(:observe).once + end + end + + describe '#on_task_execution_failure' do + it 'increments task_execute_error_total and observes time with FAILURE' do + error = ArgumentError.new('bad input') + event = Conductor::Worker::Events::TaskExecutionFailure.new( + task_type: 'my_task', task_id: 't1', worker_id: 'w1', + workflow_instance_id: 'wf1', duration_ms: 300.0, cause: error, is_retryable: true + ) + collector.on_task_execution_failure(event) + + expect(backend).to have_received(:increment).with( + 'task_execute_error_total', + labels: { taskType: 'my_task', exception: 'ArgumentError' } + ) + expect(backend).to have_received(:observe).with( + 'task_execute_time_seconds', 0.3, + labels: { taskType: 'my_task', status: 'FAILURE' } + ) + end + end + + describe '#on_task_update_completed' do + it 'observes task_update_time_seconds with status=SUCCESS' do + event = Conductor::Worker::Events::TaskUpdateCompleted.new( + task_type: 'my_task', task_id: 't1', worker_id: 'w1', + workflow_instance_id: 'wf1', duration_ms: 75.0 + ) + collector.on_task_update_completed(event) + expect(backend).to have_received(:observe).with( + 'task_update_time_seconds', 0.075, + labels: { taskType: 'my_task', status: 'SUCCESS' } + ) + end + end + + describe '#on_task_update_failure' do + it 'increments task_update_error_total and observes time with FAILURE' do + error = StandardError.new('net err') + task_result = Conductor::Http::Models::TaskResult.complete + event = Conductor::Worker::Events::TaskUpdateFailure.new( + task_type: 'my_task', task_id: 't1', worker_id: 'w1', + workflow_instance_id: 'wf1', cause: error, retry_count: 4, + task_result: task_result, duration_ms: 120.0 + ) + collector.on_task_update_failure(event) + + expect(backend).to have_received(:increment).with( + 'task_update_error_total', + labels: { taskType: 'my_task', exception: 'StandardError' } + ) + expect(backend).to have_received(:observe).with( + 'task_update_time_seconds', 0.12, + labels: { taskType: 'my_task', status: 'FAILURE' } + ) + end + + it 'skips time observation when duration_ms is nil' do + error = StandardError.new('err') + task_result = Conductor::Http::Models::TaskResult.complete + event = Conductor::Worker::Events::TaskUpdateFailure.new( + task_type: 'my_task', task_id: 't1', worker_id: 'w1', + workflow_instance_id: 'wf1', cause: error, retry_count: 4, task_result: task_result + ) + collector.on_task_update_failure(event) + expect(backend).not_to have_received(:observe) end end - describe '.canonical_metrics_enabled?' do - around do |example| - old_val = ENV.fetch('WORKER_CANONICAL_METRICS', nil) - example.run - ensure - if old_val.nil? - ENV.delete('WORKER_CANONICAL_METRICS') - else - ENV['WORKER_CANONICAL_METRICS'] = old_val - end + describe '#on_task_paused' do + it 'increments task_paused_total' do + event = Conductor::Worker::Events::TaskPaused.new(task_type: 'my_task') + collector.on_task_paused(event) + expect(backend).to have_received(:increment).with( + 'task_paused_total', labels: { taskType: 'my_task' } + ) end + end + + describe '#on_thread_uncaught_exception' do + it 'increments thread_uncaught_exceptions_total with exception label' do + event = Conductor::Worker::Events::ThreadUncaughtException.new( + cause: RuntimeError.new('boom') + ) + collector.on_thread_uncaught_exception(event) + expect(backend).to have_received(:increment).with( + 'thread_uncaught_exceptions_total', labels: { exception: 'RuntimeError' } + ) + end + end - it 'returns false by default' do - ENV.delete('WORKER_CANONICAL_METRICS') - expect(described_class.canonical_metrics_enabled?).to be false + describe '#on_active_workers_changed' do + it 'sets active_workers gauge' do + event = Conductor::Worker::Events::ActiveWorkersChanged.new( + task_type: 'my_task', count: 7 + ) + collector.on_active_workers_changed(event) + expect(backend).to have_received(:set).with( + 'active_workers', 7, labels: { taskType: 'my_task' } + ) end + end + + # --- Workflow Events --- - it 'returns true when set to "true"' do - ENV['WORKER_CANONICAL_METRICS'] = 'true' - expect(described_class.canonical_metrics_enabled?).to be true + describe '#on_workflow_start_error' do + it 'increments workflow_start_error_total' do + event = Conductor::Worker::Events::WorkflowStartError.new( + workflow_type: 'my_wf', cause: RuntimeError.new('fail') + ) + collector.on_workflow_start_error(event) + expect(backend).to have_received(:increment).with( + 'workflow_start_error_total', + labels: { workflowType: 'my_wf', exception: 'RuntimeError' } + ) end + end + + describe '#on_workflow_input_size' do + it 'observes workflow_input_size_bytes' do + event = Conductor::Worker::Events::WorkflowInputSize.new( + workflow_type: 'my_wf', size_bytes: 8192, version: 2 + ) + collector.on_workflow_input_size(event) + expect(backend).to have_received(:observe).with( + 'workflow_input_size_bytes', 8192, + labels: { workflowType: 'my_wf', version: '2' } + ) + end + end + + # --- HTTP Events --- - it 'returns false for arbitrary strings' do - ENV['WORKER_CANONICAL_METRICS'] = 'maybe' - expect(described_class.canonical_metrics_enabled?).to be false + describe '#on_http_api_request' do + it 'observes http_api_client_request_seconds' do + event = Conductor::Worker::Events::HttpApiRequest.new( + method: 'POST', uri: '/api/tasks/poll/batch/my_task', status: '200', duration_ms: 45.0 + ) + collector.on_http_api_request(event) + expect(backend).to have_received(:observe).with( + 'http_api_client_request_seconds', 0.045, + labels: { method: 'POST', uri: '/api/tasks/poll/batch/my_task', status: '200' } + ) end end end diff --git a/spec/conductor/worker/prometheus_backend_spec.rb b/spec/conductor/worker/prometheus_backend_spec.rb index 83ca12e..f34d0b5 100644 --- a/spec/conductor/worker/prometheus_backend_spec.rb +++ b/spec/conductor/worker/prometheus_backend_spec.rb @@ -2,7 +2,6 @@ require 'spec_helper' -# Check if prometheus-client gem is available PROMETHEUS_AVAILABLE = begin require 'prometheus/client' true @@ -10,7 +9,6 @@ false end -# Conditionally load the prometheus backend PROMETHEUS_BACKEND_LOADED = begin if PROMETHEUS_AVAILABLE require_relative '../../../lib/conductor/worker/telemetry/prometheus_backend' @@ -24,50 +22,53 @@ if PROMETHEUS_BACKEND_LOADED RSpec.describe Conductor::Worker::Telemetry::PrometheusBackend do - # Use a fresh registry for each test to avoid metric conflicts let(:registry) { Prometheus::Client::Registry.new } let(:backend) { described_class.new(registry: registry) } describe '#initialize' do - it 'creates a backend with the provided registry' do - expect(backend.registry).to eq(registry) + it 'registers canonical counters' do + backend + %i[task_poll_total task_execution_started_total task_poll_error_total + task_execute_error_total task_update_error_total task_paused_total + thread_uncaught_exceptions_total workflow_start_error_total].each do |name| + expect(registry.exist?(name)).to be(true), "Expected counter #{name} to be registered" + end end - it 'registers common metrics on initialization' do - expect(backend.registry.exist?(:task_poll_total)).to be true - expect(backend.registry.exist?(:task_poll_error_total)).to be true - expect(backend.registry.exist?(:task_execute_error_total)).to be true - expect(backend.registry.exist?(:task_update_failed_total)).to be true - expect(backend.registry.exist?(:task_poll_time_seconds)).to be true - expect(backend.registry.exist?(:task_execute_time_seconds)).to be true - expect(backend.registry.exist?(:task_result_size_bytes)).to be true + it 'registers canonical histograms' do + backend + %i[task_poll_time_seconds task_execute_time_seconds task_update_time_seconds + http_api_client_request_seconds task_result_size_bytes + workflow_input_size_bytes].each do |name| + expect(registry.exist?(name)).to be(true), "Expected histogram #{name} to be registered" + end end - end - describe '#increment' do - it 'increments a counter' do - expect do - backend.increment('task_poll_total', labels: { task_type: 'my_task' }) - end.not_to raise_error + it 'registers canonical gauges' do + backend + expect(registry.exist?(:active_workers)).to be true end + end - it 'increments by a custom value' do + describe '#increment' do + it 'increments a counter with camelCase labels' do expect do - backend.increment('task_poll_total', labels: { task_type: 'my_task' }, value: 5) + backend.increment('task_poll_total', labels: { taskType: 'my_task' }) end.not_to raise_error end end describe '#observe' do - it 'observes a histogram value' do + it 'observes a time histogram with status label' do expect do - backend.observe('task_poll_time_seconds', 0.5, labels: { task_type: 'my_task' }) + backend.observe('task_poll_time_seconds', 0.25, + labels: { taskType: 'my_task', status: 'SUCCESS' }) end.not_to raise_error end - it 'observes size metrics' do + it 'observes a size histogram' do expect do - backend.observe('task_result_size_bytes', 1024, labels: { task_type: 'my_task' }) + backend.observe('task_result_size_bytes', 5000, labels: { taskType: 'my_task' }) end.not_to raise_error end end @@ -75,7 +76,21 @@ describe '#set' do it 'sets a gauge value' do expect do - backend.set('active_workers', 5, labels: { task_type: 'my_task' }) + backend.set('active_workers', 3, labels: { taskType: 'my_task' }) + end.not_to raise_error + end + end + + describe 'label normalization' do + it 'fills missing declared labels with empty strings' do + expect do + backend.increment('task_poll_error_total', labels: { taskType: 'my_task' }) + end.not_to raise_error + end + + it 'drops undeclared labels' do + expect do + backend.increment('task_poll_total', labels: { taskType: 'my_task', extra: 'nope' }) end.not_to raise_error end end @@ -91,16 +106,10 @@ server = described_class.new(port: 9091) expect(server.port).to eq(9091) end - - # NOTE: Actually starting/stopping the server in tests can be flaky - # due to port binding issues. These are integration tests. end else - # Test when prometheus is not available RSpec.describe 'PrometheusBackend (prometheus-client gem unavailable)' do it 'documents that prometheus-client gem is not installed' do - # This test documents that the prometheus-client gem is not available - # The actual functionality cannot be tested without the gem expect(PROMETHEUS_AVAILABLE).to be false end end