diff --git a/bun.lock b/bun.lock index a61c220ed..6784a17e6 100644 --- a/bun.lock +++ b/bun.lock @@ -8,6 +8,10 @@ "@aws-sdk/client-bedrock-agentcore": "^3.1092.0", "@aws-sdk/client-bedrock-agentcore-control": "^3.1079.0", "@aws-sdk/client-iam": "^3.1080.0", + "@opentelemetry/api": "^1.9.1", + "@opentelemetry/exporter-metrics-otlp-http": "^0.221.0", + "@opentelemetry/resources": "^2.10.0", + "@opentelemetry/sdk-metrics": "^2.10.0", "@smithy/core": "3.29.3", "@tanstack/react-query": "^5.101.2", "cli-truncate": "^6.1.1", @@ -81,6 +85,28 @@ "@dabh/diagnostics": ["@dabh/diagnostics@2.0.8", "", { "dependencies": { "@so-ric/colorspace": "^1.1.6", "enabled": "2.0.x", "kuler": "^2.0.0" } }, "sha512-R4MSXTVnuMzGD7bzHdW2ZhhdPC/igELENcq5IjEverBvq5hn1SXCWcsi6eSsdWP0/Ur+SItRRjAktmdoX/8R/Q=="], + "@opentelemetry/api": ["@opentelemetry/api@1.9.1", "", {}, "sha512-gLyJlPHPZYdAk1JENA9LeHejZe1Ti77/pTeFm/nMXmQH/HFZlcS/O2XJB+L8fkbrNSqhdtlvjBVjxwUYanNH5Q=="], + + "@opentelemetry/api-logs": ["@opentelemetry/api-logs@0.221.0", "", { "dependencies": { "@opentelemetry/api": "^1.3.0" } }, "sha512-OlanaW1vv7ufTqQ3/fPLI4arGt5ZoM+P8abOMki6uEYnpRazepSWDwDnnw+la7kE26SHVC18//SMccrDvLKOXQ=="], + + "@opentelemetry/core": ["@opentelemetry/core@2.10.0", "", { "dependencies": { "@opentelemetry/semantic-conventions": "^1.29.0" }, "peerDependencies": { "@opentelemetry/api": ">=1.0.0 <1.10.0" } }, "sha512-/wNZ8twnEQQA4HoHu22+vcsdru6pWPWxW+7w+FlxT6Id7PE/WIbZmVKkte+PF72e0F2dnImFeHD2syyE1Mw6MQ=="], + + "@opentelemetry/exporter-metrics-otlp-http": ["@opentelemetry/exporter-metrics-otlp-http@0.221.0", "", { "dependencies": { "@opentelemetry/core": "2.10.0", "@opentelemetry/otlp-exporter-base": "0.221.0", "@opentelemetry/otlp-transformer": "0.221.0", "@opentelemetry/resources": "2.10.0", "@opentelemetry/sdk-metrics": "2.10.0" }, "peerDependencies": { "@opentelemetry/api": "^1.3.0" } }, "sha512-sRfCKbOzgy8xZQV2as0RzIZlnCmCseCKZGLfRcrpo2CBngJDr+rPtX0zkG0+oUCV5kfQPUoW3W3C96Ag3Y/Clg=="], + + "@opentelemetry/otlp-exporter-base": ["@opentelemetry/otlp-exporter-base@0.221.0", "", { "dependencies": { "@opentelemetry/core": "2.10.0", "@opentelemetry/otlp-transformer": "0.221.0" }, "peerDependencies": { "@opentelemetry/api": "^1.3.0" } }, "sha512-UFPIq80OH3Ns/oPFHRj14d4DTOxUo+MUFU8hUiCq5jTqFhdeJnfVSANHT+xp92409cA+oxzvlZCe6NM1wvCuBA=="], + + "@opentelemetry/otlp-transformer": ["@opentelemetry/otlp-transformer@0.221.0", "", { "dependencies": { "@opentelemetry/api-logs": "0.221.0", "@opentelemetry/core": "2.10.0", "@opentelemetry/resources": "2.10.0", "@opentelemetry/sdk-logs": "0.221.0", "@opentelemetry/sdk-metrics": "2.10.0", "@opentelemetry/sdk-trace": "2.10.0" }, "peerDependencies": { "@opentelemetry/api": "^1.3.0" } }, "sha512-lg6lkOU08Az23jVcn/0Els9HP+V8PnR4Km6p0KgpTggS0n/WuhnmY64rSh83Of9iR9nD+dpWr6adlcX8KzAwjg=="], + + "@opentelemetry/resources": ["@opentelemetry/resources@2.10.0", "", { "dependencies": { "@opentelemetry/core": "2.10.0", "@opentelemetry/semantic-conventions": "^1.29.0" }, "peerDependencies": { "@opentelemetry/api": ">=1.3.0 <1.10.0" } }, "sha512-q6MMm2zhggzsHVNbabYwut+a6nbuQQe3URUoxaojM/8K1IBfwwPzvxIjNi2/lI1TFe+fMHMW9MWhrtDLEXEnkA=="], + + "@opentelemetry/sdk-logs": ["@opentelemetry/sdk-logs@0.221.0", "", { "dependencies": { "@opentelemetry/api-logs": "0.221.0", "@opentelemetry/core": "2.10.0", "@opentelemetry/resources": "2.10.0", "@opentelemetry/semantic-conventions": "^1.29.0" }, "peerDependencies": { "@opentelemetry/api": ">=1.4.0 <1.10.0" } }, "sha512-FaDcazjyMp7TZZZAsqbo4IkovP0UegoCu0EBkiNt+qCqvUf7FPAsfcrZ3+ZEkKgXZ/jHafop+JoGPDk3A0SmLg=="], + + "@opentelemetry/sdk-metrics": ["@opentelemetry/sdk-metrics@2.10.0", "", { "dependencies": { "@opentelemetry/core": "2.10.0", "@opentelemetry/resources": "2.10.0" }, "peerDependencies": { "@opentelemetry/api": ">=1.9.0 <1.10.0" } }, "sha512-t6r1VSvXNtSDnPXU1FbZeetJb7yyovHmgu0wRSoftxtE0g2rSNhQZQUy69sRUCL+iioJpX8SN/S6wq6ZtvLySQ=="], + + "@opentelemetry/sdk-trace": ["@opentelemetry/sdk-trace@2.10.0", "", { "dependencies": { "@opentelemetry/core": "2.10.0", "@opentelemetry/resources": "2.10.0", "@opentelemetry/semantic-conventions": "^1.29.0" }, "peerDependencies": { "@opentelemetry/api": ">=1.3.0 <1.10.0" } }, "sha512-MfQGq3GRmTh5fM/y+OjaO0vj6+luCB1XO2gfXCalKCfgKw0eHL++sm75DNweC6ohlp+aFvACqeE0fYayqdRaoQ=="], + + "@opentelemetry/semantic-conventions": ["@opentelemetry/semantic-conventions@1.43.0", "", {}, "sha512-eSYWTm620tTk45EKSedaUL8MFYI8hW164hIXsgIHyxu3VobUB3fFCu5t0hQby6OoWRPsG1KkKUG2M5UadiLiVg=="], + "@oxlint/binding-android-arm-eabi": ["@oxlint/binding-android-arm-eabi@1.74.0", "", { "os": "android", "cpu": "arm" }, "sha512-+gHd12muVI9ZLBaWLPkHt3Fj7jihFjgQ1MGtBaRL8vWrWrI0P7dLUty/cHrHS0oqPYIRgQUJsPu2CExQuMcwNw=="], "@oxlint/binding-android-arm64": ["@oxlint/binding-android-arm64@1.74.0", "", { "os": "android", "cpu": "arm64" }, "sha512-xjKdoMB+H+RCOByv/7l7nfIGW9mlOisqYdcyC75UqYuQecLpReAeEYUf2CNeDEI3KtmUgxpRw/+c63y4AeF/Bw=="], diff --git a/package.json b/package.json index c2fa1da53..daa19ef10 100644 --- a/package.json +++ b/package.json @@ -52,10 +52,15 @@ "@aws-sdk/client-bedrock-agentcore": "^3.1092.0", "@aws-sdk/client-bedrock-agentcore-control": "^3.1079.0", "@aws-sdk/client-iam": "^3.1080.0", + "@opentelemetry/api": "^1.9.1", + "@opentelemetry/exporter-metrics-otlp-http": "^0.221.0", + "@opentelemetry/resources": "^2.10.0", + "@opentelemetry/sdk-metrics": "^2.10.0", "@smithy/core": "3.29.3", "@tanstack/react-query": "^5.101.2", "cli-truncate": "^6.1.1", "commander": "^15.0.0", + "handlebars": "^4.7.9", "ink": "^7.1.0", "ink-scroll-view": "^0.3.7", "lodash": "^4.18.1", @@ -65,7 +70,6 @@ "string-width": "^8.2.2", "winston": "^3.19.0", "winston-daily-rotate-file": "^5.0.0", - "handlebars": "^4.7.9", "zod": "^4.4.3" } } diff --git a/src/telemetry/client.test.tsx b/src/telemetry/client.test.tsx index 1eb1a08a9..096a23294 100644 --- a/src/telemetry/client.test.tsx +++ b/src/telemetry/client.test.tsx @@ -5,9 +5,10 @@ import os, { tmpdir } from "node:os"; import { DefaultTelemetryClient } from "./client"; import { createFileLogger, type Logger } from "../logging"; import { LOG_LEVEL } from "../logging"; -import { assertLogsMatch, TestGlobalConfigAccessor } from "../testing"; +import { assertLogsMatch, createSilentLogger, TestGlobalConfigAccessor } from "../testing"; import type { MetricSink } from "./types"; import { FileSystemSink } from "./fileSystemSink"; +import { DEFAULT_GLOBAL_CONFIG } from "../globalConfig"; import { PACKAGE_VERSION } from "../constants"; describe("DefaultTelemetryClient", () => { @@ -28,9 +29,20 @@ describe("DefaultTelemetryClient", () => { test("emits complete metrics to configured JSONL filesystem sinks", async () => { const auditFilePath = join(tempDir, "telemetry", "audit.jsonl"); + const sinkResourceAttributes = { + "service.name": "agentcore-cli" as const, + "service.version": "0.0.0", + "agentcore-cli.installation_id": "00000000-0000-0000-0000-000000000000", + "agentcore-cli.session_id": "00000000-0000-0000-0000-000000000000", + "os.type": os.type(), + "os.version": os.release(), + "host.arch": os.arch(), + "node.version": process.version, + }; const fileSystemSink = new FileSystemSink({ logger: logger.child({ module: "fileSystemSink" }), filePath: auditFilePath, + resourceAttributes: sinkResourceAttributes, }); const sessionId = "aaaaaaaa-bbbb-cccc-dddd-eeeeeeeeeeee"; const globalConfigAccessor = new TestGlobalConfigAccessor(); @@ -191,6 +203,16 @@ describe("DefaultTelemetryClient", () => { const sink = new FileSystemSink({ logger: logger.child({ module: "fileSystemSink" }), filePath: tempDir, + resourceAttributes: { + "service.name": "agentcore-cli", + "service.version": "0.0.0", + "agentcore-cli.installation_id": "00000000-0000-0000-0000-000000000000", + "agentcore-cli.session_id": "00000000-0000-0000-0000-000000000000", + "os.type": os.type(), + "os.version": os.release(), + "host.arch": os.arch(), + "node.version": process.version, + }, }); const client = new DefaultTelemetryClient({ @@ -274,3 +296,109 @@ describe("DefaultTelemetryClient", () => { ]); }); }); + +describe("OtelHistogramSink", () => { + let testCollector: ReturnType; + let receivedBodies: any[]; + + const logger = createSilentLogger(); + + beforeEach(async () => { + receivedBodies = []; + testCollector = Bun.serve({ + port: 0, + async fetch(req) { + const body = await req.json(); + receivedBodies.push(body); + return new Response("", { status: 200 }); + }, + }); + }); + + afterEach(async () => { + testCollector.stop(true); + }); + + test.each([ + { enabled: true, expectRequests: true }, + { enabled: false, expectRequests: false }, + ])( + "telemetry.enabled=$enabled → collector receives requests=$expectRequests", + async ({ enabled, expectRequests }) => { + const sessionId = "aaaaaaaa-bbbb-cccc-dddd-eeeeeeeeeeee"; + const exitReason = "success"; + const commandPath = "/agentcore"; + const metricName = "cli.command_run"; + const scopeName = "agentcore-cli"; + const serviceName = "agentcore-cli"; + const globalConfigAccessor = new TestGlobalConfigAccessor({ + initialConfigData: { + ...DEFAULT_GLOBAL_CONFIG, + telemetry: { + enabled, + audit: false, + endpoint: `http://localhost:${testCollector.port}`, + }, + }, + }); + + const client = new DefaultTelemetryClient({ + logger, + sessionId, + globalConfigAccessor, + }); + + const event = client.createMetricEvent(metricName, { + exit_reason: exitReason, + command_path: commandPath, + }); + await event.emit(100); + await client.shutdown(); + + if (expectRequests) { + expect(receivedBodies.length).toBeGreaterThan(0); + + const body = receivedBodies[0]; + expect(body).toMatchObject({ + resourceMetrics: [ + { + resource: { + attributes: expect.arrayContaining([ + { key: "service.name", value: { stringValue: serviceName } }, + { + key: "agentcore-cli.session_id", + value: { stringValue: sessionId }, + }, + { key: "os.type", value: { stringValue: os.type() } }, + { key: "host.arch", value: { stringValue: os.arch() } }, + ]), + }, + scopeMetrics: [ + { + scope: { name: scopeName }, + metrics: [ + { + name: metricName, + histogram: { + dataPoints: [ + { + attributes: expect.arrayContaining([ + { key: "exit_reason", value: { stringValue: exitReason } }, + { key: "command_path", value: { stringValue: commandPath } }, + ]), + }, + ], + }, + }, + ], + }, + ], + }, + ], + }); + } else { + expect(receivedBodies).toHaveLength(0); + } + }, + ); +}); diff --git a/src/telemetry/client.tsx b/src/telemetry/client.tsx index 2ac07c3de..af13f4560 100644 --- a/src/telemetry/client.tsx +++ b/src/telemetry/client.tsx @@ -13,6 +13,7 @@ import { import type { GlobalConfigAccessor } from "../globalConfig"; import { FileSystemSink } from "./fileSystemSink"; import path from "path"; +import { OtelHistogramSink } from "./otelSink"; import { PACKAGE_VERSION } from "../constants"; export type DefaultTelemetryClientConfig = { @@ -72,6 +73,7 @@ export class DefaultTelemetryClient implements TelemetryClient { private getMetricSinks: () => Promise = once(async () => { if (this.metricSinksOverride) return this.metricSinksOverride; + const resourceAttributes = await this.getResourceAttributes(); const metricSinks = []; @@ -82,6 +84,16 @@ export class DefaultTelemetryClient implements TelemetryClient { new FileSystemSink({ logger: this.logger.child({ module: "fileSystemSink" }), filePath: this.auditFilePath, + resourceAttributes, + }), + ); + + if (globalConfig.telemetry.enabled) + metricSinks.push( + new OtelHistogramSink({ + logger: this.logger.child({ module: "otelCollectorSink" }), + collectorEndpoint: globalConfig.telemetry.endpoint, + resourceAttributes, }), ); diff --git a/src/telemetry/fileSystemSink.tsx b/src/telemetry/fileSystemSink.tsx index 26283c819..c329d3fbb 100644 --- a/src/telemetry/fileSystemSink.tsx +++ b/src/telemetry/fileSystemSink.tsx @@ -1,4 +1,5 @@ import type { Logger } from "../logging"; +import type { ResourceAttributes } from "./shapes"; import type { MetricSink } from "./types"; import { mkdir, appendFile } from "fs/promises"; import { dirname } from "path"; @@ -6,6 +7,7 @@ import { dirname } from "path"; export type FileSystemSinkConfig = { logger: Logger; filePath: string; + resourceAttributes: ResourceAttributes; }; /** An implementation of {@link MetricSink} that sends all data to the specified file in JSONL format **/ @@ -15,6 +17,8 @@ export class FileSystemSink implements MetricSink { private readonly filePath: string; private logger: Logger; + private readonly resourceAttributes: ResourceAttributes; + /* a chain of promises describing the pending writes to the audit file */ private pendingWrite: Promise; @@ -22,6 +26,7 @@ export class FileSystemSink implements MetricSink { this.filePath = config.filePath; this.logger = config.logger.child({ fsSinkFilePath: this.filePath }); this.name = new.target.name; + this.resourceAttributes = config.resourceAttributes; this.pendingWrite = Promise.resolve(); } @@ -32,7 +37,7 @@ export class FileSystemSink implements MetricSink { attributes: Record, ): void { this.pendingWrite = this.pendingWrite.then(() => - this.appendEntry({ metricName, value, attrs: attributes }), + this.appendEntry({ metricName, value, attrs: { ...this.resourceAttributes, ...attributes } }), ); } diff --git a/src/telemetry/otelSink.tsx b/src/telemetry/otelSink.tsx new file mode 100644 index 000000000..638e929b3 --- /dev/null +++ b/src/telemetry/otelSink.tsx @@ -0,0 +1,97 @@ +import { MeterProvider, PeriodicExportingMetricReader } from "@opentelemetry/sdk-metrics"; +import { resourceFromAttributes } from "@opentelemetry/resources"; +import type { Logger } from "../logging"; +import { type MetricSink } from "./types"; +import { OTLPMetricExporter } from "@opentelemetry/exporter-metrics-otlp-http"; +import { type Histogram, type Meter } from "@opentelemetry/api"; +import type { ResourceAttributes } from "./shapes"; + +export type OtelCollectorSinkConfig = { + collectorEndpoint: string; + logger: Logger; + /** The resource attributes to attach to all metrics.**/ + resourceAttributes: ResourceAttributes; + /** The time period between export flushes **/ + exportIntervalMs?: number; + /** Describes the maximum time to wait when flushing metrics to the sink.**/ + flushTimeoutMs?: number; + /** Describes the maximum time to wait when shutting down the sink.**/ + shutdownTimeoutMs?: number; + /** + * Describes the scope to attach to the given metrics. Defaults to `agentcore-cli` + * See https://opentelemetry.io/docs/concepts/instrumentation-scope/ for more information. + **/ + instrumentationScope?: string; +}; + +/** + * An implementation of {@link MetricSink} that sends data to a collector using an otel histogram. + */ +export class OtelHistogramSink implements MetricSink { + private readonly name: string; + private readonly endpoint: string; + private meterProvider: MeterProvider; + private histograms: Map; + private logger: Logger; + + private readonly flushTimeoutMs: number; + private readonly shutdownTimeoutMs: number; + private readonly scopedMeter: Meter; + + constructor(config: OtelCollectorSinkConfig) { + this.endpoint = config.collectorEndpoint; + this.logger = config.logger.child({ telemetryEndpoint: config.collectorEndpoint }); + this.name = new.target.name; + + this.flushTimeoutMs = config.flushTimeoutMs ?? 500; + this.shutdownTimeoutMs = config.shutdownTimeoutMs ?? 500; + + this.meterProvider = new MeterProvider({ + resource: resourceFromAttributes(config.resourceAttributes), + readers: [ + new PeriodicExportingMetricReader({ + exporter: new OTLPMetricExporter({ + url: this.endpoint.endsWith("/v1/metrics") + ? this.endpoint + : `${this.endpoint}/v1/metrics`, + }), + exportIntervalMillis: config.exportIntervalMs ?? 5_000, + }), + ], + }); + this.scopedMeter = this.meterProvider.getMeter(config.instrumentationScope ?? "agentcore-cli"); + this.histograms = new Map(); + } + + private getHistogram(metricName: string): Histogram { + if (!this.histograms.has(metricName)) { + this.histograms.set(metricName, this.scopedMeter.createHistogram(metricName)); + } + return this.histograms.get(metricName)!; + } + + send(metricName: string, value: number, attributes: Record): void { + this.logger + .child({ metricName, metricValue: value, metricAttributes: attributes }) + .info(`sending telemetry metric to collector`); + + this.getHistogram(metricName).record(value, attributes); + } + + async shutdown(): Promise { + try { + await this.meterProvider.forceFlush({ timeoutMillis: this.flushTimeoutMs }); + } catch (e) { + const error = e instanceof Error ? e : new Error(String(e)); + this.logger + .child({ errorName: error.name, errorMessage: error.message }) + .warn(`failed to flush metrics to ${this.getName()}`); + // don't let flush failures prevent shutdown + } + await this.meterProvider.shutdown({ timeoutMillis: this.shutdownTimeoutMs }); + } + + getName(): string { + return this.name; + } +}