merge PR #5507 via eval/pr-5507: feat(telemetry): support full OTel — log and metric export
This commit is contained in:
@@ -89,9 +89,14 @@
|
||||
"@oh-my-pi/pi-wire": "catalog:",
|
||||
"@oh-my-pi/snapcompact": "catalog:",
|
||||
"@opentelemetry/api": "catalog:",
|
||||
"@opentelemetry/api-logs": "catalog:",
|
||||
"@opentelemetry/context-async-hooks": "catalog:",
|
||||
"@opentelemetry/exporter-logs-otlp-proto": "catalog:",
|
||||
"@opentelemetry/exporter-metrics-otlp-proto": "catalog:",
|
||||
"@opentelemetry/exporter-trace-otlp-proto": "catalog:",
|
||||
"@opentelemetry/resources": "catalog:",
|
||||
"@opentelemetry/sdk-logs": "catalog:",
|
||||
"@opentelemetry/sdk-metrics": "catalog:",
|
||||
"@opentelemetry/sdk-trace-base": "catalog:",
|
||||
"@opentelemetry/sdk-trace-node": "catalog:",
|
||||
"@puppeteer/browsers": "catalog:",
|
||||
@@ -377,11 +382,16 @@
|
||||
"@oh-my-pi/pi-wire": "17.0.1",
|
||||
"@oh-my-pi/snapcompact": "17.0.1",
|
||||
"@opentelemetry/api": "^1.9.1",
|
||||
"@opentelemetry/context-async-hooks": "^2.7.1",
|
||||
"@opentelemetry/api-logs": "^0.220.0",
|
||||
"@opentelemetry/context-async-hooks": "^2.9.0",
|
||||
"@opentelemetry/exporter-logs-otlp-proto": "^0.220.0",
|
||||
"@opentelemetry/exporter-metrics-otlp-proto": "^0.220.0",
|
||||
"@opentelemetry/exporter-trace-otlp-proto": "^0.220.0",
|
||||
"@opentelemetry/resources": "^2.7.1",
|
||||
"@opentelemetry/sdk-trace-base": "^2.7.1",
|
||||
"@opentelemetry/sdk-trace-node": "^2.7.1",
|
||||
"@opentelemetry/resources": "^2.9.0",
|
||||
"@opentelemetry/sdk-logs": "^0.220.0",
|
||||
"@opentelemetry/sdk-metrics": "^2.9.0",
|
||||
"@opentelemetry/sdk-trace-base": "^2.9.0",
|
||||
"@opentelemetry/sdk-trace-node": "^2.9.0",
|
||||
"@puppeteer/browsers": "^3.0.6",
|
||||
"@tailwindcss/node": "^4.3.2",
|
||||
"@tailwindcss/vite": "^4.3.2",
|
||||
@@ -800,6 +810,12 @@
|
||||
|
||||
"@opentelemetry/core": ["@opentelemetry/core@2.9.0", "", { "dependencies": { "@opentelemetry/semantic-conventions": "^1.29.0" }, "peerDependencies": { "@opentelemetry/api": ">=1.0.0 <1.10.0" } }, "sha512-m2nckMT80NnmjTYSPjJQObBJ+8dgkoajEOUbznL8AHZ3T3yHRk2P7gI1PhEBc1+lOnrYE9UWrWHqJDsmqjmNbw=="],
|
||||
|
||||
"@opentelemetry/exporter-logs-otlp-proto": ["@opentelemetry/exporter-logs-otlp-proto@0.220.0", "", { "dependencies": { "@opentelemetry/otlp-exporter-base": "0.220.0", "@opentelemetry/otlp-transformer": "0.220.0", "@opentelemetry/sdk-logs": "0.220.0" }, "peerDependencies": { "@opentelemetry/api": "^1.3.0" } }, "sha512-8LZAxdJ0ENDAFwr4j0oY35mHBltiSzvlhdQAPGiC7p9VnxtuSq4SW1gfBAdW6t6hiQG6OwUl8w7KHaOdJPKHWg=="],
|
||||
|
||||
"@opentelemetry/exporter-metrics-otlp-http": ["@opentelemetry/exporter-metrics-otlp-http@0.220.0", "", { "dependencies": { "@opentelemetry/core": "2.9.0", "@opentelemetry/otlp-exporter-base": "0.220.0", "@opentelemetry/otlp-transformer": "0.220.0", "@opentelemetry/resources": "2.9.0", "@opentelemetry/sdk-metrics": "2.9.0" }, "peerDependencies": { "@opentelemetry/api": "^1.3.0" } }, "sha512-Yqt3RBw/bRVncaE9qIIhk4WfjbAQqXuP9FgAaU+IKPndnLEp/cUqZlSC324+bpmduRz7DoTjig8Ub0PeILWXUA=="],
|
||||
|
||||
"@opentelemetry/exporter-metrics-otlp-proto": ["@opentelemetry/exporter-metrics-otlp-proto@0.220.0", "", { "dependencies": { "@opentelemetry/core": "2.9.0", "@opentelemetry/exporter-metrics-otlp-http": "0.220.0", "@opentelemetry/otlp-exporter-base": "0.220.0", "@opentelemetry/otlp-transformer": "0.220.0", "@opentelemetry/resources": "2.9.0", "@opentelemetry/sdk-metrics": "2.9.0" }, "peerDependencies": { "@opentelemetry/api": "^1.3.0" } }, "sha512-lyO+IQBdSvqHN/ZOW/OzrSWemtfD+HgWngn+HBNLhjy0YrCQQTz0OE/kSekH2Pl340dn9DWzhqHdz5Eftr+HLA=="],
|
||||
|
||||
"@opentelemetry/exporter-trace-otlp-proto": ["@opentelemetry/exporter-trace-otlp-proto@0.220.0", "", { "dependencies": { "@opentelemetry/core": "2.9.0", "@opentelemetry/otlp-exporter-base": "0.220.0", "@opentelemetry/otlp-transformer": "0.220.0", "@opentelemetry/resources": "2.9.0", "@opentelemetry/sdk-trace": "2.9.0" }, "peerDependencies": { "@opentelemetry/api": "^1.3.0" } }, "sha512-voTAD8XgJxlK7zLkXh8EzMB09zrQr3tyY/BsnDTlDiQU/UdK58MZ63A3mUjdEDrxMjCVmBHU3WQJhRmQe+Dvzg=="],
|
||||
|
||||
"@opentelemetry/otlp-exporter-base": ["@opentelemetry/otlp-exporter-base@0.220.0", "", { "dependencies": { "@opentelemetry/core": "2.9.0", "@opentelemetry/otlp-transformer": "0.220.0" }, "peerDependencies": { "@opentelemetry/api": "^1.3.0" } }, "sha512-CXYo8UD5Mn9YbgebO2EL4wejtA+gxLmLiu6HCk2KH2BR7XhFN6/6p1UlCb23DYCjeYkndevLHuejCCN1yx4+OQ=="],
|
||||
|
||||
+9
-4
@@ -39,11 +39,16 @@
|
||||
"@oh-my-pi/pi-wire": "17.0.1",
|
||||
"@oh-my-pi/snapcompact": "17.0.1",
|
||||
"@opentelemetry/api": "^1.9.1",
|
||||
"@opentelemetry/context-async-hooks": "^2.7.1",
|
||||
"@opentelemetry/api-logs": "^0.220.0",
|
||||
"@opentelemetry/context-async-hooks": "^2.9.0",
|
||||
"@opentelemetry/exporter-logs-otlp-proto": "^0.220.0",
|
||||
"@opentelemetry/exporter-metrics-otlp-proto": "^0.220.0",
|
||||
"@opentelemetry/exporter-trace-otlp-proto": "^0.220.0",
|
||||
"@opentelemetry/resources": "^2.7.1",
|
||||
"@opentelemetry/sdk-trace-base": "^2.7.1",
|
||||
"@opentelemetry/sdk-trace-node": "^2.7.1",
|
||||
"@opentelemetry/resources": "^2.9.0",
|
||||
"@opentelemetry/sdk-logs": "^0.220.0",
|
||||
"@opentelemetry/sdk-metrics": "^2.9.0",
|
||||
"@opentelemetry/sdk-trace-base": "^2.9.0",
|
||||
"@opentelemetry/sdk-trace-node": "^2.9.0",
|
||||
"@puppeteer/browsers": "^3.0.6",
|
||||
"@tailwindcss/node": "^4.3.2",
|
||||
"@tailwindcss/vite": "^4.3.2",
|
||||
|
||||
@@ -10,6 +10,7 @@
|
||||
- Added native Warp CLI-agent events for rich session status, tool approvals, and completion notifications ([#5592](https://github.com/can1357/oh-my-pi/pull/5592) by [@metaphorics](https://github.com/metaphorics)).
|
||||
- Added Codex (ChatGPT subscription) support to `generate_image`. The tool now resolves a connected `openai-codex` OAuth credential and drives OpenAI's hosted `image_generation` tool through the ChatGPT backend (`chatgpt.com/backend-api/codex/responses`, `chatgpt-account-id` header) **independent of the active chat model** — so image generation works on a ChatGPT/Codex subscription with no metered `OPENAI_API_KEY`, even when the active model is Claude/Gemini/etc. A new `providers.image: "openai-codex"` option forces it; `auto` now auto-detects a connected subscription (priority: active GPT image tool > Codex subscription > Antigravity > xAI > OpenRouter > Gemini), and the `openai` preference falls back to it when no `OPENAI_API_KEY`/active GPT model is present.
|
||||
- Added an optional `provider` parameter to `generate_image` (`auto` | `openai` | `openai-codex` | `antigravity` | `xai` | `gemini` | `openrouter`) that overrides the `providers.image` setting **for a single request** — so "generate this using gemini / codex / xai" routes per-call without changing the global setting. Absent → the `providers.image` setting applies, unchanged; the named provider uses the same resolution semantics (falls back to auto-detect if it has no credentials). File: `tools/image-gen.ts` (`imageProviderSchema`, `findImageApiKey` `preference` arg).
|
||||
- Added OpenTelemetry log and metric export alongside the existing trace export. When `OTEL_EXPORTER_OTLP_LOGS_ENDPOINT` (or the shared `OTEL_EXPORTER_OTLP_ENDPOINT`) is set, `omp` registers a `LoggerProvider` and forwards every centralized-logger event as an OTLP log record (severity + attributes + active span context for log↔trace correlation, min level via `OTEL_LOG_LEVEL`, plus a structured `agent run completed` summary event). When `OTEL_EXPORTER_OTLP_METRICS_ENDPOINT` (or the shared endpoint) is set, it registers a `MeterProvider` with a `PeriodicExportingMetricReader` and records GenAI-semconv `gen_ai.client.token.usage` plus `pi.omp.agent.*` counters/histograms (runs, steps, chat/tool calls by name+status+finish reason, latencies, estimated cost, errors) from the agent run summary and per-chat usage hooks. Each signal honors its own `OTEL_*_EXPORTER=none` kill switch, the global `OTEL_SDK_DISABLED`, and declines non-`http/protobuf` protocols independently ([#4604](https://github.com/can1357/oh-my-pi/issues/4604)).
|
||||
- `retry.fallbackChains` wildcards now support id-prefixed targets and keys: a chain entry like `"openrouter/google/*"` re-prefixes the failing model's bare id (`google-antigravity/gemini-x` → `openrouter/google/gemini-x`), a plain `"provider/*"` entry falling back *from* an aggregator strips the vendor prefix when the target provider only knows the bare id (`openrouter/google/x` → `google-vertex/x`), and an id-prefixed key (`"openrouter/google/*"`) scopes a chain to that provider's ids under the prefix.
|
||||
|
||||
### Changed
|
||||
|
||||
@@ -64,9 +64,14 @@
|
||||
"@oh-my-pi/pi-wire": "catalog:",
|
||||
"@oh-my-pi/snapcompact": "catalog:",
|
||||
"@opentelemetry/api": "catalog:",
|
||||
"@opentelemetry/api-logs": "catalog:",
|
||||
"@opentelemetry/context-async-hooks": "catalog:",
|
||||
"@opentelemetry/exporter-logs-otlp-proto": "catalog:",
|
||||
"@opentelemetry/exporter-metrics-otlp-proto": "catalog:",
|
||||
"@opentelemetry/exporter-trace-otlp-proto": "catalog:",
|
||||
"@opentelemetry/resources": "catalog:",
|
||||
"@opentelemetry/sdk-logs": "catalog:",
|
||||
"@opentelemetry/sdk-metrics": "catalog:",
|
||||
"@opentelemetry/sdk-trace-base": "catalog:",
|
||||
"@opentelemetry/sdk-trace-node": "catalog:",
|
||||
"@puppeteer/browsers": "catalog:",
|
||||
|
||||
@@ -77,7 +77,7 @@ import { executeBuiltinSlashCommand } from "./slash-commands/builtin-registry";
|
||||
import { shouldShowStartupSplash } from "./startup-splash";
|
||||
import { discoverTitleSystemPromptFile, resolvePromptInput } from "./system-prompt";
|
||||
import { createPersistedSubagentReviverFactory } from "./task/persisted-revive";
|
||||
import { initTelemetryExport, isTelemetryExportEnabled } from "./telemetry-export";
|
||||
import { createTelemetryExportConfig, initTelemetryExport, isTelemetryExportEnabled } from "./telemetry-export";
|
||||
import { concreteThinkingLevel, parseConfiguredThinkingLevel } from "./thinking";
|
||||
import type { LspStartupServerInfo } from "./tools";
|
||||
import {
|
||||
@@ -1351,15 +1351,13 @@ export async function runRootCommand(
|
||||
sessionOptions.hasUI = isInteractive || mode === "rpc-ui";
|
||||
sessionOptions.settings = settingsInstance;
|
||||
|
||||
// OTEL: register the global OTLP trace exporter when an OTLP endpoint is
|
||||
// configured via env, then switch on the agent loop's telemetry so its
|
||||
// GenAI spans (invoke_agent / chat / execute_tool) are actually emitted.
|
||||
// Both are no-ops when OTEL_EXPORTER_OTLP_ENDPOINT is unset. An empty config
|
||||
// is enough to enable telemetry — content capture is governed by the
|
||||
// standard OTEL_INSTRUMENTATION_GENAI_CAPTURE_MESSAGE_CONTENT env var.
|
||||
// OTEL: register global OTLP exporters when an endpoint is configured via
|
||||
// env, then switch on the agent loop's telemetry hooks so traces, run-level
|
||||
// metrics, and structured logs have source events to export. Content capture
|
||||
// remains governed by OTEL_INSTRUMENTATION_GENAI_CAPTURE_MESSAGE_CONTENT.
|
||||
await logger.time("initTelemetryExport", initTelemetryExport);
|
||||
if (isTelemetryExportEnabled()) {
|
||||
sessionOptions.telemetry = {};
|
||||
sessionOptions.telemetry = createTelemetryExportConfig(sessionOptions.telemetry);
|
||||
}
|
||||
|
||||
// Handle CLI --api-key as runtime override (not persisted)
|
||||
|
||||
@@ -1,144 +1,500 @@
|
||||
/**
|
||||
* OTLP trace export bootstrap.
|
||||
* OTLP telemetry export bootstrap.
|
||||
*
|
||||
* oh-my-pi's agent core (`@oh-my-pi/pi-agent-core`) emits OpenTelemetry GenAI
|
||||
* spans through the global `@opentelemetry/api` tracer, but only when a
|
||||
* TracerProvider is registered in the process — otherwise the API returns a
|
||||
* no-op tracer and the spans are silently dropped. The shipped CLI never
|
||||
* registered one, so headless / embedded hosts (e.g. an ACP harness that
|
||||
* spawns `omp` as a child process) had no way to collect omp's internal traces.
|
||||
* spans through the global `@opentelemetry/api` tracer, and exposes run-level
|
||||
* callbacks for metrics/log pipelines. This module registers the OTLP/proto
|
||||
* trace, log, and metric SDK providers when the standard `OTEL_*` endpoint env
|
||||
* vars are set so `omp` can be observed by any OTLP collector without vendor
|
||||
* coupling.
|
||||
*
|
||||
* This module registers a NodeTracerProvider with an OTLP/proto exporter when
|
||||
* the standard `OTEL_EXPORTER_OTLP_ENDPOINT` (or `..._TRACES_ENDPOINT`) env var
|
||||
* is set, following the zero-code OTEL env contract: the exporter reads its
|
||||
* endpoint, headers, and timeout from `OTEL_EXPORTER_OTLP_*` itself. The
|
||||
* consuming process configures the destination entirely through env; omp stays
|
||||
* provider-agnostic and ships no vendor coupling. Only the `http/protobuf`
|
||||
* transport is supported — an `OTEL_EXPORTER_OTLP*_PROTOCOL` of `grpc` or
|
||||
* `http/json` declines rather than misrouting spans.
|
||||
*
|
||||
* The OTLP/proto exporter on the 2.x line is used deliberately: the 1.x line
|
||||
* deadlocks under Bun — its `req.on('close')` handler fires a spurious failure
|
||||
* after the success path. `exporter-trace-otlp-proto@0.218` paired with
|
||||
* `sdk-trace-base@2.7` exports cleanly on Bun.
|
||||
* Only the `http/protobuf` transport is supported — an
|
||||
* `OTEL_EXPORTER_OTLP*_PROTOCOL` of `grpc` or `http/json` declines rather than
|
||||
* misrouting protobuf payloads. The exporter line is pinned to the 0.218/2.7
|
||||
* family validated under Bun; the 1.x OTLP line deadlocks when its
|
||||
* `req.on("close")` handler fires after a successful export.
|
||||
*/
|
||||
import type {
|
||||
AgentRunCoverage,
|
||||
AgentRunSummary,
|
||||
AgentTelemetryConfig,
|
||||
AgentTelemetryWarning,
|
||||
ChatUsageEvent,
|
||||
ToolStatus,
|
||||
} from "@oh-my-pi/pi-agent-core";
|
||||
import { logger, postmortem } from "@oh-my-pi/pi-utils";
|
||||
import type * as TraceNode from "@opentelemetry/sdk-trace-node";
|
||||
import {
|
||||
type Attributes,
|
||||
type AttributeValue,
|
||||
type Counter,
|
||||
context,
|
||||
type Histogram,
|
||||
type Meter,
|
||||
metrics,
|
||||
} from "@opentelemetry/api";
|
||||
import { type LogAttributes, logs, type Logger as OtelLogger, SeverityNumber } from "@opentelemetry/api-logs";
|
||||
import { AsyncLocalStorageContextManager } from "@opentelemetry/context-async-hooks";
|
||||
import { OTLPLogExporter } from "@opentelemetry/exporter-logs-otlp-proto";
|
||||
import { OTLPMetricExporter } from "@opentelemetry/exporter-metrics-otlp-proto";
|
||||
import { OTLPTraceExporter } from "@opentelemetry/exporter-trace-otlp-proto";
|
||||
import { resourceFromAttributes } from "@opentelemetry/resources";
|
||||
import { BatchLogRecordProcessor, LoggerProvider } from "@opentelemetry/sdk-logs";
|
||||
import { MeterProvider, PeriodicExportingMetricReader } from "@opentelemetry/sdk-metrics";
|
||||
import { BatchSpanProcessor } from "@opentelemetry/sdk-trace-base";
|
||||
import { NodeTracerProvider } from "@opentelemetry/sdk-trace-node";
|
||||
|
||||
/**
|
||||
* Periodic flush interval. A long-lived `omp` process (the ACP server is
|
||||
* spawned once and reused across many turns) would otherwise hold finished
|
||||
* spans until the batch window elapses or the process exits.
|
||||
* telemetry until a batch window elapses or the process exits.
|
||||
*/
|
||||
const FLUSH_INTERVAL_MS = 30_000;
|
||||
|
||||
let provider: TraceNode.NodeTracerProvider | undefined;
|
||||
const SERVICE_NAME = "oh-my-pi";
|
||||
|
||||
type TelemetrySignal = "trace" | "log" | "metric";
|
||||
type OtelLogLevel = "none" | logger.LogLevel;
|
||||
|
||||
interface SignalConfig {
|
||||
readonly trace: boolean;
|
||||
readonly log: boolean;
|
||||
readonly metric: boolean;
|
||||
}
|
||||
|
||||
const LOG_SEVERITY: Record<logger.LogLevel, SeverityNumber> = {
|
||||
error: SeverityNumber.ERROR,
|
||||
warn: SeverityNumber.WARN,
|
||||
info: SeverityNumber.INFO,
|
||||
debug: SeverityNumber.DEBUG,
|
||||
};
|
||||
|
||||
const LOG_LEVEL_WEIGHT: Record<logger.LogLevel, number> = {
|
||||
error: 0,
|
||||
warn: 1,
|
||||
info: 2,
|
||||
debug: 3,
|
||||
};
|
||||
|
||||
const TOOL_STATUSES = ["ok", "error", "skipped", "blocked", "timeout", "aborted"] satisfies readonly ToolStatus[];
|
||||
|
||||
let traceProvider: NodeTracerProvider | undefined;
|
||||
let logProvider: LoggerProvider | undefined;
|
||||
let meterProvider: MeterProvider | undefined;
|
||||
let metricRecorder: AgentMetricRecorder | undefined;
|
||||
let otelLogger: OtelLogger | undefined;
|
||||
let unregisterLogSink: (() => void) | undefined;
|
||||
let initPromise: Promise<void> | undefined;
|
||||
|
||||
/**
|
||||
* Whether {@link initTelemetryExport} registered a real provider. The CLI uses
|
||||
* this to decide whether to switch on the agent loop's telemetry config — there
|
||||
* is no point emitting spans into a no-op tracer.
|
||||
* Whether {@link initTelemetryExport} registered any real OTLP signal provider.
|
||||
* The CLI uses this to decide whether to switch on the agent loop's telemetry
|
||||
* hooks; metrics and structured logs need those callbacks even when traces are
|
||||
* disabled.
|
||||
*/
|
||||
export function isTelemetryExportEnabled(): boolean {
|
||||
return provider !== undefined;
|
||||
if (traceProvider) return true;
|
||||
if (logProvider) return true;
|
||||
if (meterProvider) return true;
|
||||
return false;
|
||||
}
|
||||
|
||||
/**
|
||||
* Register the global TracerProvider + OTLP exporter when an OTLP endpoint is
|
||||
* configured via env. Idempotent, and a no-op when no endpoint is set (or when
|
||||
* the OTEL kill-switches are engaged), so it is safe to call unconditionally at
|
||||
* startup.
|
||||
* Merge OTLP metrics/log hooks into an existing agent telemetry config.
|
||||
*
|
||||
* The caller still owns content-capture policy, cost estimation, and custom
|
||||
* attributes. This only appends host-level metrics/log forwarding for the
|
||||
* providers registered by {@link initTelemetryExport}.
|
||||
*/
|
||||
export function createTelemetryExportConfig(
|
||||
config: AgentTelemetryConfig | undefined,
|
||||
): AgentTelemetryConfig | undefined {
|
||||
if (!isTelemetryExportEnabled()) return config;
|
||||
return {
|
||||
...config,
|
||||
onChatUsage: async event => {
|
||||
await config?.onChatUsage?.(event);
|
||||
metricRecorder?.recordChatUsage(event);
|
||||
},
|
||||
onRunEnd: (summary, coverage) => {
|
||||
config?.onRunEnd?.(summary, coverage);
|
||||
metricRecorder?.recordRun(summary, coverage);
|
||||
emitRunSummaryLog(summary, coverage);
|
||||
},
|
||||
onTelemetryWarning: warning => {
|
||||
config?.onTelemetryWarning?.(warning);
|
||||
emitTelemetryWarningLog(warning);
|
||||
},
|
||||
};
|
||||
}
|
||||
|
||||
/**
|
||||
* Register global trace/log/meter providers when OTLP endpoints are configured
|
||||
* through env. Idempotent, and a no-op when no signal has an endpoint (or when
|
||||
* the OTEL kill-switches are engaged), so startup can call it unconditionally.
|
||||
*/
|
||||
export async function initTelemetryExport(): Promise<void> {
|
||||
if (provider) return;
|
||||
if (isTelemetryExportEnabled()) return;
|
||||
if (initPromise) return initPromise;
|
||||
|
||||
// The OTEL env contract parses booleans and enum lists case-insensitively, so
|
||||
// OTEL_SDK_DISABLED=TRUE and OTEL_TRACES_EXPORTER=None must also disable export.
|
||||
if (process.env.OTEL_SDK_DISABLED?.trim().toLowerCase() === "true") return;
|
||||
if (tracesExporterDisabled(process.env.OTEL_TRACES_EXPORTER)) return;
|
||||
|
||||
const endpoint = process.env.OTEL_EXPORTER_OTLP_TRACES_ENDPOINT ?? process.env.OTEL_EXPORTER_OTLP_ENDPOINT;
|
||||
if (!endpoint) return;
|
||||
const signalConfig = resolveSignalConfig();
|
||||
if (!signalConfig.trace && !signalConfig.log && !signalConfig.metric) return;
|
||||
|
||||
// We only ship the http/protobuf transport (the line validated on Bun). The
|
||||
// OTEL contract lets OTEL_EXPORTER_OTLP*_PROTOCOL select grpc / http/json;
|
||||
// rather than silently send protobuf-over-HTTP to a grpc :4317 port and lose
|
||||
// every span, decline when an unsupported protocol is requested.
|
||||
const protocol = (process.env.OTEL_EXPORTER_OTLP_TRACES_PROTOCOL ?? process.env.OTEL_EXPORTER_OTLP_PROTOCOL)
|
||||
?.trim()
|
||||
.toLowerCase();
|
||||
if (protocol && protocol !== "http/protobuf") {
|
||||
logger.warn(
|
||||
`OTEL trace export disabled: OTEL_EXPORTER_OTLP_PROTOCOL=${protocol} is unsupported (only http/protobuf)`,
|
||||
);
|
||||
return;
|
||||
}
|
||||
|
||||
initPromise = registerProvider();
|
||||
initPromise = registerProviders(signalConfig);
|
||||
return initPromise;
|
||||
}
|
||||
|
||||
async function registerProvider(): Promise<void> {
|
||||
const [
|
||||
{ AsyncLocalStorageContextManager },
|
||||
{ OTLPTraceExporter },
|
||||
{ resourceFromAttributes },
|
||||
{ BatchSpanProcessor },
|
||||
{ NodeTracerProvider },
|
||||
] = await Promise.all([
|
||||
import("@opentelemetry/context-async-hooks"),
|
||||
import("@opentelemetry/exporter-trace-otlp-proto"),
|
||||
import("@opentelemetry/resources"),
|
||||
import("@opentelemetry/sdk-trace-base"),
|
||||
import("@opentelemetry/sdk-trace-node"),
|
||||
]);
|
||||
|
||||
// The exporter reads endpoint/headers/timeout from OTEL_EXPORTER_OTLP_* itself,
|
||||
// so there is nothing to thread through here.
|
||||
const exporter = new OTLPTraceExporter();
|
||||
const tracerProvider = new NodeTracerProvider({
|
||||
resource: resourceFromAttributes({
|
||||
"service.name": process.env.OTEL_SERVICE_NAME ?? "oh-my-pi",
|
||||
}),
|
||||
spanProcessors: [new BatchSpanProcessor(exporter)],
|
||||
async function registerProviders(signalConfig: SignalConfig): Promise<void> {
|
||||
const resource = resourceFromAttributes({
|
||||
"service.name": process.env.OTEL_SERVICE_NAME ?? SERVICE_NAME,
|
||||
});
|
||||
// register() installs the global tracer provider and the W3C trace-context +
|
||||
// baggage propagators; the explicit AsyncLocalStorage context manager keeps
|
||||
// parent/child span linkage working under Bun.
|
||||
tracerProvider.register({ contextManager: new AsyncLocalStorageContextManager().enable() });
|
||||
provider = tracerProvider;
|
||||
|
||||
if (signalConfig.trace) {
|
||||
const exporter = new OTLPTraceExporter();
|
||||
traceProvider = new NodeTracerProvider({
|
||||
resource,
|
||||
spanProcessors: [new BatchSpanProcessor(exporter)],
|
||||
});
|
||||
traceProvider.register({ contextManager: new AsyncLocalStorageContextManager().enable() });
|
||||
}
|
||||
|
||||
if (signalConfig.metric) {
|
||||
const exporter = new OTLPMetricExporter();
|
||||
meterProvider = new MeterProvider({
|
||||
resource,
|
||||
readers: [new PeriodicExportingMetricReader({ exporter })],
|
||||
});
|
||||
metrics.setGlobalMeterProvider(meterProvider);
|
||||
metricRecorder = new AgentMetricRecorder(metrics.getMeter("@oh-my-pi/pi-coding-agent"));
|
||||
}
|
||||
|
||||
if (signalConfig.log) {
|
||||
const exporter = new OTLPLogExporter();
|
||||
logProvider = new LoggerProvider({
|
||||
resource,
|
||||
processors: [new BatchLogRecordProcessor({ exporter })],
|
||||
});
|
||||
logs.setGlobalLoggerProvider(logProvider);
|
||||
otelLogger = logProvider.getLogger("@oh-my-pi/pi-coding-agent");
|
||||
unregisterLogSink = logger.registerLogSink(event => {
|
||||
emitOtelLog(
|
||||
event.level,
|
||||
event.message,
|
||||
logAttributesFromContext(event.context),
|
||||
"pi.omp.log",
|
||||
event.timestamp,
|
||||
);
|
||||
});
|
||||
}
|
||||
|
||||
const flushTimer = setInterval(() => {
|
||||
provider?.forceFlush().catch(() => {});
|
||||
flushTelemetryExport().catch(() => {});
|
||||
}, FLUSH_INTERVAL_MS);
|
||||
flushTimer.unref();
|
||||
|
||||
// Shut down through postmortem rather than a bare signal listener. postmortem
|
||||
// owns SIGINT/SIGTERM/SIGHUP/exit and quit(), and awaits registered cleanups
|
||||
// before calling process.exit — so the batch processor's final OTLP export
|
||||
// completes instead of being cut off mid-flight on the shutdown path.
|
||||
postmortem.register("otel-trace-export", async () => {
|
||||
postmortem.register("otel-export", async () => {
|
||||
clearInterval(flushTimer);
|
||||
await provider?.shutdown();
|
||||
unregisterLogSink?.();
|
||||
unregisterLogSink = undefined;
|
||||
const shutdowns: Promise<void>[] = [];
|
||||
if (traceProvider) shutdowns.push(traceProvider.shutdown());
|
||||
if (logProvider) shutdowns.push(logProvider.shutdown());
|
||||
if (meterProvider) shutdowns.push(meterProvider.shutdown());
|
||||
await Promise.all(shutdowns);
|
||||
});
|
||||
}
|
||||
|
||||
/**
|
||||
* Parse the `OTEL_TRACES_EXPORTER` selection. The value is a case-insensitive,
|
||||
* comma-separated list; the literal `none` disables span export entirely.
|
||||
*/
|
||||
function tracesExporterDisabled(raw: string | undefined): boolean {
|
||||
if (!raw) return false;
|
||||
return raw.split(",").some(entry => entry.trim().toLowerCase() === "none");
|
||||
function resolveSignalConfig(): SignalConfig {
|
||||
const signalConfig: SignalConfig = {
|
||||
trace: signalEnabled(
|
||||
"trace",
|
||||
process.env.OTEL_EXPORTER_OTLP_TRACES_ENDPOINT ?? process.env.OTEL_EXPORTER_OTLP_ENDPOINT,
|
||||
process.env.OTEL_TRACES_EXPORTER,
|
||||
process.env.OTEL_EXPORTER_OTLP_TRACES_PROTOCOL ?? process.env.OTEL_EXPORTER_OTLP_PROTOCOL,
|
||||
),
|
||||
log: signalEnabled(
|
||||
"log",
|
||||
process.env.OTEL_EXPORTER_OTLP_LOGS_ENDPOINT ?? process.env.OTEL_EXPORTER_OTLP_ENDPOINT,
|
||||
process.env.OTEL_LOGS_EXPORTER,
|
||||
process.env.OTEL_EXPORTER_OTLP_LOGS_PROTOCOL ?? process.env.OTEL_EXPORTER_OTLP_PROTOCOL,
|
||||
),
|
||||
metric: signalEnabled(
|
||||
"metric",
|
||||
process.env.OTEL_EXPORTER_OTLP_METRICS_ENDPOINT ?? process.env.OTEL_EXPORTER_OTLP_ENDPOINT,
|
||||
process.env.OTEL_METRICS_EXPORTER,
|
||||
process.env.OTEL_EXPORTER_OTLP_METRICS_PROTOCOL ?? process.env.OTEL_EXPORTER_OTLP_PROTOCOL,
|
||||
),
|
||||
};
|
||||
return signalConfig;
|
||||
}
|
||||
|
||||
function signalEnabled(
|
||||
signal: TelemetrySignal,
|
||||
endpoint: string | undefined,
|
||||
exporterSelection: string | undefined,
|
||||
protocolSelection: string | undefined,
|
||||
): boolean {
|
||||
if (exporterSelection) {
|
||||
for (const entry of exporterSelection.split(",")) {
|
||||
if (entry.trim().toLowerCase() === "none") return false;
|
||||
}
|
||||
}
|
||||
if (!endpoint) return false;
|
||||
|
||||
const protocol = protocolSelection?.trim().toLowerCase();
|
||||
if (protocol && protocol !== "http/protobuf") {
|
||||
logger.warn(`OTEL ${signal} export disabled: OTEL_EXPORTER_OTLP_PROTOCOL=${protocol} is unsupported`, {
|
||||
supported: "http/protobuf",
|
||||
});
|
||||
return false;
|
||||
}
|
||||
return true;
|
||||
}
|
||||
|
||||
class AgentMetricRecorder {
|
||||
readonly #tokenUsage: Histogram<Attributes>;
|
||||
readonly #chatCostUsd: Counter<Attributes>;
|
||||
readonly #runs: Counter<Attributes>;
|
||||
readonly #steps: Counter<Attributes>;
|
||||
readonly #chatCalls: Counter<Attributes>;
|
||||
readonly #chatDurationMs: Histogram<Attributes>;
|
||||
readonly #toolCalls: Counter<Attributes>;
|
||||
readonly #toolDurationMs: Histogram<Attributes>;
|
||||
readonly #errors: Counter<Attributes>;
|
||||
|
||||
constructor(meter: Meter) {
|
||||
this.#tokenUsage = meter.createHistogram("gen_ai.client.token.usage", {
|
||||
description: "Token usage reported by GenAI chat calls.",
|
||||
unit: "{token}",
|
||||
});
|
||||
this.#chatCostUsd = meter.createCounter("pi.omp.agent.chat.cost.estimated_usd", {
|
||||
description: "Estimated USD cost for completed chat calls.",
|
||||
unit: "USD",
|
||||
});
|
||||
this.#runs = meter.createCounter("pi.omp.agent.runs", {
|
||||
description: "Completed agent runs.",
|
||||
unit: "{run}",
|
||||
});
|
||||
this.#steps = meter.createCounter("pi.omp.agent.steps", {
|
||||
description: "Agent loop steps completed inside a run.",
|
||||
unit: "{step}",
|
||||
});
|
||||
this.#chatCalls = meter.createCounter("pi.omp.agent.chat.calls", {
|
||||
description: "Chat calls completed inside agent runs.",
|
||||
unit: "{call}",
|
||||
});
|
||||
this.#chatDurationMs = meter.createHistogram("pi.omp.agent.chat.duration", {
|
||||
description: "Total chat latency observed in an agent run.",
|
||||
unit: "ms",
|
||||
});
|
||||
this.#toolCalls = meter.createCounter("pi.omp.agent.tool.calls", {
|
||||
description: "Tool calls completed inside agent runs.",
|
||||
unit: "{call}",
|
||||
});
|
||||
this.#toolDurationMs = meter.createHistogram("pi.omp.agent.tool.duration", {
|
||||
description: "Total tool latency observed in an agent run.",
|
||||
unit: "ms",
|
||||
});
|
||||
this.#errors = meter.createCounter("pi.omp.agent.errors", {
|
||||
description: "Errors observed in chat and tool execution.",
|
||||
unit: "{error}",
|
||||
});
|
||||
}
|
||||
|
||||
recordChatUsage(event: ChatUsageEvent): void {
|
||||
const baseAttrs = metricAttributes({
|
||||
"gen_ai.operation.name": "chat",
|
||||
"gen_ai.provider.name": event.provider,
|
||||
"gen_ai.request.model": event.model,
|
||||
"gen_ai.response.service_tier": event.serviceTier,
|
||||
"pi.gen_ai.agent.id": event.agent?.id,
|
||||
"pi.gen_ai.agent.name": event.agent?.name,
|
||||
});
|
||||
|
||||
this.#recordToken(event.usage.inputTokens, baseAttrs, "input");
|
||||
this.#recordToken(event.usage.outputTokens, baseAttrs, "output");
|
||||
this.#recordToken(event.usage.totalTokens, baseAttrs, "total");
|
||||
this.#recordToken(event.usage.cachedInputTokens, baseAttrs, "cache_read_input");
|
||||
this.#recordToken(event.usage.cacheWriteTokens, baseAttrs, "cache_write_input");
|
||||
this.#recordToken(event.usage.reasoningOutputTokens, baseAttrs, "reasoning_output");
|
||||
|
||||
if (event.cost && "usd" in event.cost && event.cost.usd > 0) {
|
||||
this.#chatCostUsd.add(event.cost.usd, baseAttrs);
|
||||
}
|
||||
}
|
||||
|
||||
recordRun(summary: AgentRunSummary, coverage: AgentRunCoverage): void {
|
||||
const runAttrs = metricAttributes({
|
||||
"pi.omp.agent.models_used.count": coverage.modelsUsed.length,
|
||||
"pi.omp.agent.providers_used.count": coverage.providersUsed.length,
|
||||
"pi.omp.agent.tools_available.count": coverage.toolsAvailable.length,
|
||||
"pi.omp.agent.tools_invoked.count": coverage.toolsInvoked.length,
|
||||
"pi.omp.agent.tools_unused.count": coverage.toolsUnused.length,
|
||||
});
|
||||
|
||||
this.#runs.add(1, runAttrs);
|
||||
if (summary.stepCount > 0) this.#steps.add(summary.stepCount, runAttrs);
|
||||
if (summary.chats.totalLatencyMs > 0) this.#chatDurationMs.record(summary.chats.totalLatencyMs, runAttrs);
|
||||
|
||||
for (const reason in summary.chats.byStopReason) {
|
||||
const count = summary.chats.byStopReason[reason];
|
||||
if (count > 0)
|
||||
this.#chatCalls.add(count, metricAttributes({ ...runAttrs, "gen_ai.response.finish_reason": reason }));
|
||||
}
|
||||
for (const toolName in summary.tools.byName) {
|
||||
const counters = summary.tools.byName[toolName];
|
||||
const toolAttrs = metricAttributes({ ...runAttrs, "gen_ai.tool.name": toolName });
|
||||
if (counters.totalLatencyMs > 0) this.#toolDurationMs.record(counters.totalLatencyMs, toolAttrs);
|
||||
for (const status of TOOL_STATUSES) {
|
||||
const count = counters[status];
|
||||
if (count > 0) this.#toolCalls.add(count, metricAttributes({ ...toolAttrs, "pi.omp.tool.status": status }));
|
||||
}
|
||||
}
|
||||
for (const errorType in summary.errors.byType) {
|
||||
const count = summary.errors.byType[errorType];
|
||||
if (count > 0) this.#errors.add(count, metricAttributes({ ...runAttrs, "error.type": errorType }));
|
||||
}
|
||||
}
|
||||
|
||||
#recordToken(value: number | undefined, baseAttrs: Attributes, tokenType: string): void {
|
||||
if (!value || value <= 0) return;
|
||||
this.#tokenUsage.record(value, metricAttributes({ ...baseAttrs, "gen_ai.token.type": tokenType }));
|
||||
}
|
||||
}
|
||||
|
||||
function metricAttributes(fields: Readonly<Record<string, unknown>>): Attributes {
|
||||
const out: Attributes = {};
|
||||
for (const key in fields) {
|
||||
const value = fields[key];
|
||||
if (value === undefined || value === null) continue;
|
||||
if (typeof value === "string" || typeof value === "number" || typeof value === "boolean") {
|
||||
out[key] = value;
|
||||
continue;
|
||||
}
|
||||
const text = String(value);
|
||||
if (text.length > 0) out[key] = text;
|
||||
}
|
||||
return out;
|
||||
}
|
||||
|
||||
function emitRunSummaryLog(summary: AgentRunSummary, coverage: AgentRunCoverage): void {
|
||||
emitOtelLog(
|
||||
"info",
|
||||
"agent run completed",
|
||||
{
|
||||
"pi.omp.agent.step_count": summary.stepCount,
|
||||
"pi.omp.agent.chats.total": summary.chats.total,
|
||||
"pi.omp.agent.chats.total_latency_ms": summary.chats.totalLatencyMs,
|
||||
"pi.omp.agent.tools.total": summary.tools.total,
|
||||
"pi.omp.agent.tools.ok": summary.tools.ok,
|
||||
"pi.omp.agent.tools.error": summary.tools.error,
|
||||
"pi.omp.agent.tools.skipped": summary.tools.skipped,
|
||||
"pi.omp.agent.tools.blocked": summary.tools.blocked,
|
||||
"pi.omp.agent.tools.timeout": summary.tools.timeout,
|
||||
"pi.omp.agent.tools.aborted": summary.tools.aborted,
|
||||
"pi.omp.agent.tools.total_latency_ms": summary.tools.totalLatencyMs,
|
||||
"pi.omp.agent.usage.input_tokens": summary.usage.inputTokens,
|
||||
"pi.omp.agent.usage.output_tokens": summary.usage.outputTokens,
|
||||
"pi.omp.agent.usage.cached_input_tokens": summary.usage.cachedInputTokens,
|
||||
"pi.omp.agent.usage.cache_write_tokens": summary.usage.cacheWriteTokens,
|
||||
"pi.omp.agent.usage.reasoning_output_tokens": summary.usage.reasoningOutputTokens,
|
||||
"pi.omp.agent.usage.total_tokens": summary.usage.totalTokens,
|
||||
"pi.omp.agent.cost.estimated_usd": summary.cost.estimatedUsd,
|
||||
"pi.omp.agent.cost.unavailable_reasons": summary.cost.unavailableReasons.join(","),
|
||||
"pi.omp.agent.errors.total": summary.errors.total,
|
||||
"pi.omp.agent.coverage.tools_available": coverage.toolsAvailable.join(","),
|
||||
"pi.omp.agent.coverage.tools_invoked": coverage.toolsInvoked.join(","),
|
||||
"pi.omp.agent.coverage.tools_unused": coverage.toolsUnused.join(","),
|
||||
"pi.omp.agent.coverage.models_used": coverage.modelsUsed.join(","),
|
||||
"pi.omp.agent.coverage.providers_used": coverage.providersUsed.join(","),
|
||||
},
|
||||
"pi.omp.agent.run.completed",
|
||||
);
|
||||
}
|
||||
|
||||
function emitTelemetryWarningLog(warning: AgentTelemetryWarning): void {
|
||||
const attrs = logAttributesFromContext({
|
||||
code: warning.code,
|
||||
error: warning.error,
|
||||
});
|
||||
emitOtelLog("warn", warning.message, attrs, "pi.omp.telemetry.warning");
|
||||
}
|
||||
|
||||
function emitOtelLog(
|
||||
level: logger.LogLevel,
|
||||
body: string,
|
||||
attributes: LogAttributes,
|
||||
eventName: string,
|
||||
timestamp = new Date(),
|
||||
): void {
|
||||
if (!otelLogger) return;
|
||||
const minLevel = parseOtelLogLevel(process.env.OTEL_LOG_LEVEL);
|
||||
if (minLevel === "none") return;
|
||||
if (LOG_LEVEL_WEIGHT[level] > LOG_LEVEL_WEIGHT[minLevel]) return;
|
||||
otelLogger.emit({
|
||||
eventName,
|
||||
timestamp,
|
||||
observedTimestamp: new Date(),
|
||||
severityNumber: LOG_SEVERITY[level],
|
||||
severityText: level.toUpperCase(),
|
||||
body,
|
||||
attributes,
|
||||
context: context.active(),
|
||||
});
|
||||
}
|
||||
|
||||
function parseOtelLogLevel(raw: string | undefined): OtelLogLevel {
|
||||
if (!raw) return "info";
|
||||
switch (raw.trim().toLowerCase()) {
|
||||
case "none":
|
||||
return "none";
|
||||
case "error":
|
||||
return "error";
|
||||
case "warn":
|
||||
case "warning":
|
||||
return "warn";
|
||||
case "debug":
|
||||
return "debug";
|
||||
default:
|
||||
return "info";
|
||||
}
|
||||
}
|
||||
|
||||
function logAttributesFromContext(input: Record<string, unknown> | undefined): LogAttributes {
|
||||
const out: LogAttributes = { "process.pid": process.pid };
|
||||
if (!input) return out;
|
||||
for (const key in input) {
|
||||
const attr = logAttributeValue(input[key]);
|
||||
if (attr !== undefined) out[key] = attr;
|
||||
}
|
||||
return out;
|
||||
}
|
||||
|
||||
function logAttributeValue(value: unknown): AttributeValue | undefined {
|
||||
if (value === undefined || value === null) return undefined;
|
||||
if (typeof value === "string" || typeof value === "number" || typeof value === "boolean") return value;
|
||||
if (value instanceof Error) {
|
||||
return `${value.name}: ${value.message}`;
|
||||
}
|
||||
try {
|
||||
const text = JSON.stringify(value);
|
||||
if (text && text.length > 0) return text;
|
||||
} catch {
|
||||
return String(value);
|
||||
}
|
||||
return String(value);
|
||||
}
|
||||
|
||||
/**
|
||||
* Flush any buffered spans to the exporter. No-op when export is disabled.
|
||||
* Flush buffered spans, log records, and metrics. No-op when export is disabled.
|
||||
* Hosts embedding the agent can call this at natural boundaries (e.g. the end
|
||||
* of a turn) so traces surface promptly rather than on the batch interval.
|
||||
* of a turn) so telemetry surfaces promptly rather than on the batch interval.
|
||||
*/
|
||||
export async function flushTelemetryExport(): Promise<void> {
|
||||
await provider?.forceFlush();
|
||||
const flushes: Promise<void>[] = [];
|
||||
if (traceProvider) flushes.push(traceProvider.forceFlush());
|
||||
if (logProvider) flushes.push(logProvider.forceFlush());
|
||||
if (meterProvider) flushes.push(meterProvider.forceFlush());
|
||||
await Promise.all(flushes);
|
||||
}
|
||||
|
||||
@@ -0,0 +1,209 @@
|
||||
/**
|
||||
* Positive-path probe for the OTLP log + metric exporters, run as a subprocess
|
||||
* by telemetry-export.test.ts. Keeping it out-of-process means the global
|
||||
* LoggerProvider / MeterProvider singletons that initTelemetryExport() registers
|
||||
* never leak into the test runner.
|
||||
*
|
||||
* Stands up a loopback OTLP/proto receiver, points the standard env vars at it,
|
||||
* registers the providers, drives a log record through the bridged
|
||||
* `@oh-my-pi/pi-utils` logger and metric instruments through the agent
|
||||
* telemetry hooks, flushes, and exits 0 only if the receiver got a non-empty
|
||||
* protobuf POST at both /v1/logs and /v1/metrics.
|
||||
*/
|
||||
|
||||
import type { AgentRunCoverage, AgentRunSummary, ChatUsageEvent } from "@oh-my-pi/pi-agent-core";
|
||||
import { emptyAgentRunCoverage, emptyAgentRunSummary } from "@oh-my-pi/pi-agent-core";
|
||||
import {
|
||||
createTelemetryExportConfig,
|
||||
flushTelemetryExport,
|
||||
initTelemetryExport,
|
||||
isTelemetryExportEnabled,
|
||||
} from "@oh-my-pi/pi-coding-agent/telemetry-export";
|
||||
import { logger } from "@oh-my-pi/pi-utils";
|
||||
|
||||
const seen = new Set<string>();
|
||||
const metricPayloads: Uint8Array[] = [];
|
||||
|
||||
interface ProtobufField {
|
||||
readonly number: number;
|
||||
readonly bytes?: Uint8Array;
|
||||
}
|
||||
|
||||
function readVarint(bytes: Uint8Array, offset: number): [number, number] {
|
||||
let value = 0;
|
||||
let shift = 0;
|
||||
while (offset < bytes.length) {
|
||||
const byte = bytes[offset++];
|
||||
value += (byte & 0x7f) * 2 ** shift;
|
||||
if ((byte & 0x80) === 0) return [value, offset];
|
||||
shift += 7;
|
||||
}
|
||||
throw new Error("Truncated protobuf varint");
|
||||
}
|
||||
|
||||
function protobufFields(bytes: Uint8Array): ProtobufField[] {
|
||||
const fields: ProtobufField[] = [];
|
||||
for (let offset = 0; offset < bytes.length; ) {
|
||||
const [tag, nextOffset] = readVarint(bytes, offset);
|
||||
offset = nextOffset;
|
||||
const wireType = tag & 7;
|
||||
const number = tag >>> 3;
|
||||
if (wireType === 0) {
|
||||
[, offset] = readVarint(bytes, offset);
|
||||
fields.push({ number });
|
||||
} else if (wireType === 1) {
|
||||
offset += 8;
|
||||
fields.push({ number });
|
||||
} else if (wireType === 2) {
|
||||
const [length, valueOffset] = readVarint(bytes, offset);
|
||||
offset = valueOffset;
|
||||
const end = offset + length;
|
||||
if (end > bytes.length) throw new Error("Truncated protobuf field");
|
||||
fields.push({ number, bytes: bytes.slice(offset, end) });
|
||||
offset = end;
|
||||
} else if (wireType === 5) {
|
||||
offset += 4;
|
||||
fields.push({ number });
|
||||
} else {
|
||||
throw new Error(`Unsupported protobuf wire type ${wireType}`);
|
||||
}
|
||||
}
|
||||
return fields;
|
||||
}
|
||||
|
||||
function pointCountForMetric(bytes: Uint8Array, metricName: string): number | undefined {
|
||||
const fields = protobufFields(bytes);
|
||||
const isMetric = fields.some(
|
||||
field => field.number === 1 && field.bytes && new TextDecoder().decode(field.bytes) === metricName,
|
||||
);
|
||||
if (isMetric) {
|
||||
const aggregation = fields.find(field => field.number === 7 || field.number === 9)?.bytes;
|
||||
if (!aggregation) return undefined;
|
||||
return protobufFields(aggregation).filter(field => field.number === 1).length;
|
||||
}
|
||||
for (const field of fields) {
|
||||
if (!field.bytes) continue;
|
||||
try {
|
||||
const count = pointCountForMetric(field.bytes, metricName);
|
||||
if (count !== undefined) return count;
|
||||
} catch {
|
||||
// This length-delimited field is a scalar string or bytes value, not a nested message.
|
||||
}
|
||||
}
|
||||
return undefined;
|
||||
}
|
||||
|
||||
function assertSingleMetricPoint(metricName: string): void {
|
||||
const counts = metricPayloads.map(payload => pointCountForMetric(payload, metricName));
|
||||
if (!counts.includes(1)) {
|
||||
throw new Error(`${metricName} expected one dimensioned point, got ${counts.join(",")}`);
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
const server = Bun.serve({
|
||||
port: 0,
|
||||
async fetch(req) {
|
||||
const path = new URL(req.url).pathname;
|
||||
if (req.method === "POST" && req.headers.get("content-type")?.startsWith("application/x-protobuf")) {
|
||||
const body = await req.arrayBuffer();
|
||||
if (path.endsWith("/v1/metrics")) metricPayloads.push(new Uint8Array(body));
|
||||
if (body.byteLength > 0) {
|
||||
if (path.endsWith("/v1/logs")) seen.add("logs");
|
||||
if (path.endsWith("/v1/metrics")) seen.add("metrics");
|
||||
}
|
||||
}
|
||||
return new Response('{"partialSuccess":{}}', {
|
||||
status: 200,
|
||||
headers: { "content-type": "application/json" },
|
||||
});
|
||||
},
|
||||
});
|
||||
|
||||
const base = `http://localhost:${server.port}`;
|
||||
process.env.OTEL_EXPORTER_OTLP_LOGS_ENDPOINT = `${base}/v1/logs`;
|
||||
process.env.OTEL_EXPORTER_OTLP_METRICS_ENDPOINT = `${base}/v1/metrics`;
|
||||
process.env.OTEL_SERVICE_NAME = "oh-my-pi-signals-probe";
|
||||
// Force a short metric export interval so the periodic reader flushes fast.
|
||||
process.env.OTEL_METRIC_EXPORT_INTERVAL = "500";
|
||||
|
||||
await initTelemetryExport();
|
||||
if (!isTelemetryExportEnabled()) {
|
||||
console.error("PROBE: providers did not register");
|
||||
await server.stop(true);
|
||||
process.exit(2);
|
||||
}
|
||||
|
||||
const config = createTelemetryExportConfig(undefined);
|
||||
if (!config) {
|
||||
console.error("PROBE: export config not produced");
|
||||
await server.stop(true);
|
||||
process.exit(2);
|
||||
}
|
||||
|
||||
// Bridged utility logger -> OTel log record.
|
||||
logger.error("probe error", { code: "probe" });
|
||||
|
||||
// Metric instruments via the agent telemetry hooks.
|
||||
const usage: ChatUsageEvent = {
|
||||
span: undefined as never,
|
||||
agent: { id: "main", name: "Main" },
|
||||
conversationId: "probe-session",
|
||||
stepNumber: 0,
|
||||
model: "claude-haiku-4-5",
|
||||
provider: "anthropic",
|
||||
serviceTier: undefined,
|
||||
usage: {
|
||||
inputTokens: 1000,
|
||||
outputTokens: 200,
|
||||
totalTokens: 1200,
|
||||
cachedInputTokens: 0,
|
||||
cacheWriteTokens: 0,
|
||||
reasoningOutputTokens: 0,
|
||||
},
|
||||
cost: { usd: 0.01 },
|
||||
attributes: undefined,
|
||||
headers: undefined,
|
||||
};
|
||||
await config.onChatUsage?.(usage);
|
||||
|
||||
const summary: AgentRunSummary = {
|
||||
...emptyAgentRunSummary(),
|
||||
chats: { total: 1, byStopReason: { end_turn: 1 }, totalLatencyMs: 1500 },
|
||||
tools: {
|
||||
total: 1,
|
||||
ok: 1,
|
||||
error: 0,
|
||||
skipped: 0,
|
||||
blocked: 0,
|
||||
timeout: 0,
|
||||
aborted: 0,
|
||||
totalLatencyMs: 42,
|
||||
byName: {
|
||||
read: { total: 1, ok: 1, error: 0, skipped: 0, blocked: 0, timeout: 0, aborted: 0, totalLatencyMs: 42 },
|
||||
},
|
||||
},
|
||||
stepCount: 1,
|
||||
};
|
||||
const coverage: AgentRunCoverage = {
|
||||
...emptyAgentRunCoverage(),
|
||||
toolsAvailable: ["read", "write"],
|
||||
toolsInvoked: ["read"],
|
||||
toolsUnused: ["write"],
|
||||
modelsUsed: ["claude-haiku-4-5"],
|
||||
providersUsed: ["anthropic"],
|
||||
};
|
||||
config.onRunEnd?.(summary, coverage);
|
||||
|
||||
await flushTelemetryExport();
|
||||
// The metric reader exports on its own interval; wait one cycle then flush.
|
||||
await Bun.sleep(700);
|
||||
await flushTelemetryExport();
|
||||
assertSingleMetricPoint("pi.omp.agent.chat.calls");
|
||||
assertSingleMetricPoint("pi.omp.agent.tool.calls");
|
||||
assertSingleMetricPoint("pi.omp.agent.tool.duration");
|
||||
await server.stop(true);
|
||||
|
||||
const ok = seen.has("logs") && seen.has("metrics");
|
||||
console.log(ok ? "PROBE: RECEIVED" : `PROBE: MISSING ${["logs", "metrics"].filter(s => !seen.has(s)).join(",")}`);
|
||||
process.exit(ok ? 0 : 1);
|
||||
@@ -11,10 +11,16 @@ import { initTelemetryExport, isTelemetryExportEnabled } from "@oh-my-pi/pi-codi
|
||||
const OTEL_KEYS = [
|
||||
"OTEL_EXPORTER_OTLP_ENDPOINT",
|
||||
"OTEL_EXPORTER_OTLP_TRACES_ENDPOINT",
|
||||
"OTEL_EXPORTER_OTLP_LOGS_ENDPOINT",
|
||||
"OTEL_EXPORTER_OTLP_METRICS_ENDPOINT",
|
||||
"OTEL_EXPORTER_OTLP_PROTOCOL",
|
||||
"OTEL_EXPORTER_OTLP_TRACES_PROTOCOL",
|
||||
"OTEL_EXPORTER_OTLP_LOGS_PROTOCOL",
|
||||
"OTEL_EXPORTER_OTLP_METRICS_PROTOCOL",
|
||||
"OTEL_SDK_DISABLED",
|
||||
"OTEL_TRACES_EXPORTER",
|
||||
"OTEL_LOGS_EXPORTER",
|
||||
"OTEL_METRICS_EXPORTER",
|
||||
] as const;
|
||||
|
||||
let saved: Record<string, string | undefined>;
|
||||
@@ -45,8 +51,8 @@ describe("initTelemetryExport gating", () => {
|
||||
expect(isTelemetryExportEnabled()).toBe(false);
|
||||
});
|
||||
|
||||
it("stays disabled when OTEL_TRACES_EXPORTER=none even with an endpoint", async () => {
|
||||
process.env.OTEL_EXPORTER_OTLP_ENDPOINT = "http://localhost:4318";
|
||||
it("stays disabled when OTEL_TRACES_EXPORTER=none and only the traces endpoint is set", async () => {
|
||||
process.env.OTEL_EXPORTER_OTLP_TRACES_ENDPOINT = "http://localhost:4318";
|
||||
process.env.OTEL_TRACES_EXPORTER = "none";
|
||||
await initTelemetryExport();
|
||||
expect(isTelemetryExportEnabled()).toBe(false);
|
||||
@@ -71,12 +77,23 @@ describe("initTelemetryExport gating", () => {
|
||||
|
||||
delete process.env.OTEL_SDK_DISABLED;
|
||||
process.env.OTEL_TRACES_EXPORTER = "otlp,None";
|
||||
process.env.OTEL_LOGS_EXPORTER = "none";
|
||||
process.env.OTEL_METRICS_EXPORTER = "none";
|
||||
await initTelemetryExport();
|
||||
expect(isTelemetryExportEnabled()).toBe(false);
|
||||
});
|
||||
|
||||
it("stays disabled when every signal exporter is set to none", async () => {
|
||||
process.env.OTEL_EXPORTER_OTLP_ENDPOINT = "http://localhost:4318";
|
||||
process.env.OTEL_TRACES_EXPORTER = "none";
|
||||
process.env.OTEL_LOGS_EXPORTER = "none";
|
||||
process.env.OTEL_METRICS_EXPORTER = "none";
|
||||
await initTelemetryExport();
|
||||
expect(isTelemetryExportEnabled()).toBe(false);
|
||||
});
|
||||
});
|
||||
|
||||
describe("initTelemetryExport export path", () => {
|
||||
describe("initTelemetryExport signals export path", () => {
|
||||
it("registers a provider and exports spans to an OTLP/proto receiver", async () => {
|
||||
// Run in a subprocess: initTelemetryExport() registers a process-global
|
||||
// provider, so exercising the positive path in-process would leak that
|
||||
@@ -88,4 +105,15 @@ describe("initTelemetryExport export path", () => {
|
||||
expect(stdout).toContain("PROBE: RECEIVED");
|
||||
expect(code).toBe(0);
|
||||
}, 20_000);
|
||||
|
||||
it("exports log records and metrics to OTLP/proto receivers", async () => {
|
||||
// Same subprocess isolation as the trace probe: the logs/metrics probe
|
||||
// drives the bridged logger and the agent telemetry metric hooks, then
|
||||
// asserts protobuf POSTs landed at both /v1/logs and /v1/metrics.
|
||||
const probe = fileURLToPath(new URL("./otel-signals-probe.ts", import.meta.url));
|
||||
const proc = Bun.spawn(["bun", probe], { stdout: "pipe", stderr: "pipe" });
|
||||
const [code, stdout] = await Promise.all([proc.exited, new Response(proc.stdout).text()]);
|
||||
expect(stdout).toContain("PROBE: RECEIVED");
|
||||
expect(code).toBe(0);
|
||||
}, 20_000);
|
||||
});
|
||||
|
||||
@@ -2,6 +2,9 @@
|
||||
|
||||
## [Unreleased]
|
||||
|
||||
### Added
|
||||
|
||||
- Added a structured log sink API to the centralized logger (`registerLogSink`, `LogEvent`, `LogLevel`) so out-of-band consumers (e.g. OpenTelemetry log export) receive every `error`/`warn`/`info`/`debug` event after the local transport path runs, without disturbing existing file/console logging ([#4604](https://github.com/can1357/oh-my-pi/issues/4604)).
|
||||
### Fixed
|
||||
|
||||
- Fixed fatal cleanup failing to reach `process.exit()` when terminal stderr is revoked, and isolated rotating log files/audit state per process to prevent concurrent OMP instances from racing compression and rotation ([#5716](https://github.com/can1357/oh-my-pi/issues/5716)).
|
||||
|
||||
@@ -17,6 +17,42 @@ import winston from "winston";
|
||||
import DailyRotateFile from "winston-daily-rotate-file";
|
||||
import { getLogsDir } from "./dirs";
|
||||
import { drainModuleLoadEvents } from "./timing-buffer";
|
||||
/** Severity names accepted by the centralized logger. */
|
||||
export type LogLevel = "error" | "warn" | "info" | "debug";
|
||||
|
||||
/** Structured log event forwarded to out-of-band sinks such as OpenTelemetry. */
|
||||
export interface LogEvent {
|
||||
readonly level: LogLevel;
|
||||
readonly message: string;
|
||||
readonly context: Record<string, unknown> | undefined;
|
||||
readonly timestamp: Date;
|
||||
}
|
||||
|
||||
/** Receives each structured log event after the local transport path runs. */
|
||||
export type LogSink = (event: LogEvent) => void;
|
||||
|
||||
const logSinks = new Set<LogSink>();
|
||||
|
||||
/** Register an out-of-band log sink and return a disposer. */
|
||||
export function registerLogSink(sink: LogSink): () => void {
|
||||
logSinks.add(sink);
|
||||
return () => {
|
||||
logSinks.delete(sink);
|
||||
};
|
||||
}
|
||||
|
||||
function emitToSinks(level: LogLevel, message: string, context: Record<string, unknown> | undefined): void {
|
||||
if (logSinks.size === 0) return;
|
||||
const event: LogEvent = { level, message, context, timestamp: new Date() };
|
||||
for (const sink of logSinks) {
|
||||
try {
|
||||
sink(event);
|
||||
} catch {
|
||||
// Sinks are side channels; they must never break local logging.
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
const PROCESS_LOG_PATTERN = /^omp\.\d{4}-\d{2}-\d{2}\.(\d+)\.log(?:\.\d+)?$/;
|
||||
const PROCESS_AUDIT_PATTERN = /^\.omp\.(\d+)-audit\.json$/;
|
||||
@@ -212,6 +248,7 @@ export function error(message: string, context?: Record<string, unknown>): void
|
||||
} catch {
|
||||
// Silently ignore logging failures
|
||||
}
|
||||
emitToSinks("error", message, context);
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -225,6 +262,7 @@ export function warn(message: string, context?: Record<string, unknown>): void {
|
||||
} catch {
|
||||
// Silently ignore logging failures
|
||||
}
|
||||
emitToSinks("warn", message, context);
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -238,6 +276,7 @@ export function info(message: string, context?: Record<string, unknown>): void {
|
||||
} catch {
|
||||
// Silently ignore logging failures
|
||||
}
|
||||
emitToSinks("info", message, context);
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -251,6 +290,7 @@ export function debug(message: string, context?: Record<string, unknown>): void
|
||||
} catch {
|
||||
// Silently ignore logging failures
|
||||
}
|
||||
emitToSinks("debug", message, context);
|
||||
}
|
||||
|
||||
/**
|
||||
|
||||
Reference in New Issue
Block a user