From 64dfe5dd1e1e195a0868a2fec241e1060a3da9d0 Mon Sep 17 00:00:00 2001 From: Sagnik Ghosh Date: Tue, 6 Oct 2026 00:58:21 +0530 Subject: [PATCH] feat: add operation logging and release Runtime 1.4.0 --- README.md | 1 + docs/ARCHITECTURE.md | 1 + docs/GETTING-STARTED.md | 8 +- docs/OPERATION-LOGGING.md | 114 ++++ package-lock.json | 8 +- packages/otlp-ingester/package.json | 2 +- packages/otlp-ingester/src/clickhouse.ts | 14 + packages/otlp-ingester/src/context.ts | 124 +++++ .../src/continuous-detection.test.ts | 2 +- packages/otlp-ingester/src/logs.test.ts | 153 ++++++ packages/otlp-ingester/src/logs.ts | 200 +++++++ packages/otlp-ingester/src/migrations.ts | 2 + packages/otlp-ingester/src/normalize-otlp.ts | 22 +- packages/otlp-ingester/src/otlp-proto.ts | 23 + .../otlp-ingester/src/server.logs.test.ts | 132 +++++ packages/otlp-ingester/src/server.ts | 16 +- packages/runtime-next/README.md | 11 + packages/runtime-next/package.json | 4 +- packages/runtime-next/src/server.ts | 12 + packages/runtime-node/README.md | 9 +- packages/runtime-node/package.json | 4 +- packages/runtime-node/src/index.ts | 13 + packages/runtime-node/src/logger.ts | 511 ++++++++++++++++++ packages/runtime-node/src/server.ts | 11 + .../test/logger-ingester.test.mjs | 100 ++++ packages/runtime-node/test/logger.test.mjs | 168 ++++++ 26 files changed, 1639 insertions(+), 26 deletions(-) create mode 100644 docs/OPERATION-LOGGING.md create mode 100644 packages/otlp-ingester/src/context.ts create mode 100644 packages/otlp-ingester/src/logs.test.ts create mode 100644 packages/otlp-ingester/src/logs.ts create mode 100644 packages/otlp-ingester/src/server.logs.test.ts create mode 100644 packages/runtime-node/src/logger.ts create mode 100644 packages/runtime-node/test/logger-ingester.test.mjs create mode 100644 packages/runtime-node/test/logger.test.mjs diff --git a/README.md b/README.md index 92a18f7..bc86342 100644 --- a/README.md +++ b/README.md @@ -143,6 +143,7 @@ all, and ad-blockers can't tell it apart from your own API traffic. - [Stack integrations](docs/INTEGRATIONS.md) — React, Node, Next.js, Go, Rust, generic OTel - [Using Autter Runtime **without npm**](docs/WITHOUT-NPM.md) — any OTel SDK, an OTel Collector, or plain HTTP from any language - [Architecture & data model](docs/ARCHITECTURE.md) +- [Operation logging and diagnostic context](docs/OPERATION-LOGGING.md) - [Continuous detection, profiles, and outcomes](docs/CONTINUOUS-DETECTION.md) - [Roadmap](docs/PLAN.md) · [Releasing](docs/RELEASING.md) diff --git a/docs/ARCHITECTURE.md b/docs/ARCHITECTURE.md index a42c000..c038fb7 100644 --- a/docs/ARCHITECTURE.md +++ b/docs/ARCHITECTURE.md @@ -43,6 +43,7 @@ Two key scopes separate frontend and backend credentials: | --- | --- | --- | --- | | `runtime_error_occurrences` | MergeTree | `(org_id, repository_id, fingerprint, occurred_at)` | 14 d | | `runtime_spans` | MergeTree | `(org_id, repository_id, trace_id, started_at)` | 7 d | +| `runtime_logs` | ReplacingMergeTree | `(org_id, repository_id, occurred_at, event_id)` | 14 d | | `runtime_metrics_1m` | SummingMergeTree | `(org_id, repository_id, service, environment, release, route, bucket_at)` | 90 d | | `runtime_llm_calls` | MergeTree | `(org_id, repository_id, started_at)` | 90 d | | `runtime_profile_samples` | MergeTree | `(org_id, repository_id, service, environment, release, observed_at, profile_id)` | 7 d | diff --git a/docs/GETTING-STARTED.md b/docs/GETTING-STARTED.md index 6298e12..b19d872 100644 --- a/docs/GETTING-STARTED.md +++ b/docs/GETTING-STARTED.md @@ -1,3 +1,9 @@ +# Operation logging + +For structured messages, completed operation summaries, and explicit outcomes, +see [Operation logging and diagnostic context](OPERATION-LOGGING.md). These APIs +require the updated Node/Next.js SDK and ingester described in that guide. + # Getting started Already collecting logs in an external provider? See [External sources](EXTERNAL-SOURCES.md) for repository connectors. The SDK/OTel setup below applies when you instrument application code. @@ -335,4 +341,4 @@ clickhouse-client --password dev`. - [ ] Direct browser ingest: your CSP includes `connect-src https://your-ingester…`, and you accept that ad-blockers may drop some events (the relay avoids this). -- [ ] The ingester's `/healthz` is wired to your load-balancer health check. +- [ ] The ingester's `/healthz` is wired to your load-balancer health check. \ No newline at end of file diff --git a/docs/OPERATION-LOGGING.md b/docs/OPERATION-LOGGING.md new file mode 100644 index 0000000..3f3ab4b --- /dev/null +++ b/docs/OPERATION-LOGGING.md @@ -0,0 +1,114 @@ +# Operation logging and diagnostic context + +Operation logging starts in `@autter/runtime-node` and `@autter/runtime-next` +version **1.4.0**. Update the ingester to **1.4.0** before upgrading SDKs: it creates +`runtime_logs` through migration `0011-runtime-logs` and accepts `/v1/logs`. + +## Node and Next.js + +Initialize `initAutterServer` once, before the application starts. For Next.js, +use `registerAutter` in `instrumentation.ts` and import the APIs below from +`@autter/runtime-next/server`. Use the Node runtime; these APIs use Node's async +context and do not run in the browser or an edge worker. + +```ts +import { + initAutterServer, withRuntimeOperation, runtimeLogger, +} from "@autter/runtime-node"; + +const runtime = initAutterServer({ + apiKey: process.env.AUTTER_RUNTIME_KEY!, + service: "payments-api", + release: process.env.GIT_SHA, + logging: { console: false, minLevel: "info" }, +}); + +await withRuntimeOperation("checkout", async (operation) => { + operation.setContext({ "payment.provider": "stripe", "cart.item_count": 3 }); + await operation.step("reserve_inventory", () => reserveInventory()); + const payment = await operation.step("confirm_payment", () => confirmPayment()); + runtimeLogger.info("Payment confirmation returned", { "payment.attempts": 2 }); + if (!payment.confirmed) { + operation.outcome("failed", "Payment was not confirmed; no order created"); + return; + } + await operation.step("create_order", () => createOrder(payment)); +}); + +// Flush at the end of a short-lived invocation; shutdown before process exit. +await runtime.shutdown(); +``` + +`withRuntimeOperation(name, fn, attributes?)` runs the callback in an always-recorded +process span and emits one completed operation summary. It returns the callback's +result and rethrows its original error. An unhandled callback error records an +exception in the trace; the log summary is related evidence, not a second issue. + +`operation.setContext(attributes)` merges redacted nested context (arrays are +replaced) and adds attributes to the operation and +its active span. `operation.step(name, fn)` records the step result and elapsed +time; it rethrows failures. A caught step error can be recovered by application +code: the final outcome follows the callback's result or an explicit outcome. +Up to 64 steps are recorded. Log names and context should describe the operation, +not contain request bodies, payment data, or personal information. +`RuntimeLogContext` supports nested objects, arrays and scalar values; standard +trace attributes encode structured values as JSON for OTLP compatibility. + +`operation.outcome(status, message?)` accepts `succeeded`, `failed`, `degraded`, +`cancelled`, and `pending`. The default is `succeeded` when the callback returns. +Returning normally is not proof of business success: declare an outcome when +the intended result was not achieved. An explicit `failed` outcome also emits +the existing `autter.outcome` trace event, including the reporting call site. +`pending` means the operation has not confirmed its downstream result. + +`createRuntimeLogger(attributes?)` creates a logger with `debug`, `info`, `warn`, +and `error` methods. `runtimeLogger` is the default instance. `error` records +diagnostic context; use `captureException` for exception grouping or declare a +failed operation outcome for business-failure grouping. Logs inherit the +active operation ID, name, parent operation ID, and application attributes, plus +the active OTel trace/span IDs. Concurrent operations keep separate contexts. +Async context does not cross a queue or process: propagate an opaque workflow +identifier in the job payload and pass it as a custom attribute in the consumer. + +## Collection and privacy + +- Logs are OTLP/HTTP log records. Both OTLP JSON and protobuf are accepted. + evlog and other OTLP log exporters can use `/v1/logs` with a server ingest key. +- Source attributes are bounded and scrubbed again at ingestion. Tenant identity + comes from the validated key, never the supplied event. Client/browser keys + cannot send server logs. +- Redaction applies before console output and export. Never intentionally send + secrets or personal data; pattern-based redaction cannot identify every secret. +- `logging.minLevel` filters ordinary messages; completed operation summaries + are retained independently. `logging.console` defaults to true and can be + disabled without disabling export. +- Logs batch in memory (up to 1,000 records or 4 MiB, whichever comes first). + Requests contain up to 50 records or 512 KiB. Context has depth, field, step + and string limits; truncated context is marked. Each flush has a 10-second + budget and makes at most three attempts per batch. Failed batches remain + buffered for a later flush, with automatic retry intervals up to 30 seconds. Buffer overflow and failed shutdown delivery are + reported. `flushRuntimeLogs()` rejects if delivery fails, and + `runtimeLogStats()` reports buffered and dropped counts. This is best-effort + telemetry, not durable delivery or an audit log. +- Records expire after 14 days. Successful operation counts represent captured + summaries, not a guarantee that every application operation was observed. + +## Investigation and fixes + +Autter's repository Runtime Logs view displays captured messages, operation +summaries, outcomes, steps, and trace IDs. Investigation readers retrieve related +logs and spans using captured IDs and compare bounded successful operations. +They disclose unavailable sources and truncation. Event text is untrusted +diagnostic evidence, never instructions to an agent. Business context improves +an investigation but does not establish a cause by itself. + +The observations feed the existing root-cause analysis and draft-fix pipelines, +which retain their repository settings, source checks and validation rules. A reported outcome call site can locate the reporting code; it is not +proof that this location caused the failure. The fix agent must confirm the +cause, reproduce applicable conditions, and validate the intended result. + +Log and trace exporters are independent. Operation evidence may arrive after an +initial analysis; refresh the evidence panel or rerun analysis to incorporate +later arrivals. Shutdown attempts log export before closing tracers. Regular Runtime flushing +exports logs and traces independently so a failed log export does not block +exception telemetry. diff --git a/package-lock.json b/package-lock.json index cd0c9b3..8d5c550 100644 --- a/package-lock.json +++ b/package-lock.json @@ -4667,7 +4667,7 @@ }, "packages/otlp-ingester": { "name": "@autter/otlp-ingester", - "version": "1.3.4", + "version": "1.4.0", "license": "MIT", "dependencies": { "@clickhouse/client": "^1.12.0", @@ -4789,11 +4789,11 @@ }, "packages/runtime-next": { "name": "@autter/runtime-next", - "version": "1.3.4", + "version": "1.4.0", "license": "MIT", "dependencies": { "@autter/runtime-browser": "^1.3.3", - "@autter/runtime-node": "^1.3.4" + "@autter/runtime-node": "^1.4.0" }, "devDependencies": { "@types/react": "^18.3.0", @@ -4809,7 +4809,7 @@ }, "packages/runtime-node": { "name": "@autter/runtime-node", - "version": "1.3.4", + "version": "1.4.0", "license": "MIT", "dependencies": { "@opentelemetry/api": "^1.9.0", diff --git a/packages/otlp-ingester/package.json b/packages/otlp-ingester/package.json index 3da4b97..24aa3a3 100644 --- a/packages/otlp-ingester/package.json +++ b/packages/otlp-ingester/package.json @@ -1,6 +1,6 @@ { "name": "@autter/otlp-ingester", - "version": "1.3.4", + "version": "1.4.0", "description": "Self-hostable OTLP + browser-error ingest service for Autter Runtime: normalises telemetry into a per-repo ClickHouse data model", "license": "MIT", "type": "module", diff --git a/packages/otlp-ingester/src/clickhouse.ts b/packages/otlp-ingester/src/clickhouse.ts index c4a0d96..2a37371 100644 --- a/packages/otlp-ingester/src/clickhouse.ts +++ b/packages/otlp-ingester/src/clickhouse.ts @@ -8,6 +8,7 @@ import { latencyTableDDL, type LatencyHistogram } from "./latency.js"; import { profileTableDDL, type ProfileSample } from "./profiles.js"; import { memoryTableDDL, platformEventTableDDL, type MemorySample, type PlatformEvent } from "./memory.js"; import { sourceMapTableDDL } from "./source-maps.js"; +import { logTableDDL, type RuntimeLogRecord } from "./logs.js"; import { MIGRATIONS, migrationsTableDDL, @@ -75,6 +76,7 @@ export class ClickHouseStore { memoryTableDDL(db), platformEventTableDDL(db), sourceMapTableDDL(db), + logTableDDL(db), `CREATE TABLE IF NOT EXISTS ${db}.runtime_error_occurrences ( org_id String, repository_id String, @@ -294,6 +296,18 @@ export class ClickHouseStore { }); } + async insertLogs(ctx: IngestContext, logs: RuntimeLogRecord[]): Promise { + if (!logs.length) return; + if (!this.configured) throw new Error("CLICKHOUSE_URL is not configured"); + await this.ensureSchema(); + await this.getClient().insert({ table: this.table("runtime_logs"), format: "JSONEachRow", clickhouse_settings: INSERT_SETTINGS, + values: logs.map((row) => ({ org_id: ctx.orgId, repository_id: ctx.repositoryId, event_id: row.id, + service: row.service, environment: row.environment, release: row.release, trace_id: row.traceId, + span_id: row.spanId, operation_id: row.operationId, operation: row.operation, event_type: row.type, + severity: row.severity, message: row.message, outcome: row.outcome, duration_ms: row.durationMs, + attributes: JSON.stringify(row.attributes), occurred_at: row.occurredAt.toISOString() })) }); + } + async insertProfileSamples(ctx: IngestContext, samples: ProfileSample[]): Promise { if (!samples.length) return; if (!this.configured) throw new Error("CLICKHOUSE_URL is not configured"); diff --git a/packages/otlp-ingester/src/context.ts b/packages/otlp-ingester/src/context.ts new file mode 100644 index 0000000..ff5e30d --- /dev/null +++ b/packages/otlp-ingester/src/context.ts @@ -0,0 +1,124 @@ +/** Bounded, defensive privacy boundary for custom telemetry from any OTLP SDK. */ +export function sanitizeRuntimeContext( + input: unknown, +): Record { + let budget = 512; + const seen = new WeakSet(); + const scrub = (value: string) => + value + .replace(/[A-Z0-9._%+-]+@[A-Z0-9.-]+\.[A-Z]{2,}/gi, "[redacted]") + .replace( + /\b(?:bearer\s+[A-Za-z0-9._~+\/=~-]{10,}|eyJ[A-Za-z0-9_-]+\.[A-Za-z0-9_-]+\.[A-Za-z0-9_-]+|gh[pousr]_[A-Za-z0-9]{20,}|sk-[A-Za-z0-9_-]{20,}|xox[baprs]-[A-Za-z0-9-]{10,}|autter_(?:rt|pat)_[A-Za-z0-9_-]{10,}|(?:AKIA|ASIA)[A-Z0-9]{16})\b/gi, + "[redacted]", + ) + .replace( + /-----BEGIN [A-Z ]*PRIVATE KEY-----[\s\S]*?(?:-----END [A-Z ]*PRIVATE KEY-----|$)/g, + "[redacted]", + ) + .replace(/([a-z][a-z0-9+.-]*:\/\/)[^\s/:@]+:[^\s/@]+@/gi, "$1[redacted]@") + .replace(/(https?:\/\/[^\s?#]+)[?#][^\s]*/gi, "$1"); + const visit = (value: unknown, key: string, depth: number): unknown => { + if (--budget < 0 || depth > 6) return "[truncated]"; + if (/__proto__|constructor|prototype/i.test(key)) return undefined; + const usage = + /(?:^|\.)(?:input|output|total|prompt|completion)_?tokens$/i.test(key) && + typeof value === "number" && + Number.isFinite(value) && + value >= 0; + if ( + !usage && + /password|passwd|secret|token|credential|authorization|cookie|email|phone|ssn|card[._-]?number|connection[._-]?string|api[._-]?key|private[._-]?key|request[._-]?body|response[._-]?body|headers|url[._-]?query/i.test( + key, + ) + ) + return "[redacted]"; + if (typeof value === "string") { + // Serialized custom context crosses the same privacy boundary as nested values. + if (/^\s*[\[{]/.test(value)) { + try { + return visit(JSON.parse(value), key, depth + 1); + } catch { + /* plain text */ + } + } + const text = /(?:url|path|route|target)$/i.test(key) + ? value.split(/[?#]/)[0]! + : value; + return scrub(text).slice(0, /stack/i.test(key) ? 32000 : 2048); + } + if (typeof value === "number") + return Number.isFinite(value) ? value : undefined; + if (typeof value === "boolean" || value === null) return value; + if (!value || typeof value !== "object") return undefined; + if (seen.has(value)) return "[circular]"; + seen.add(value); + if (Array.isArray(value)) + return value.slice(0, 64).map((v) => visit(v, key, depth + 1)); + const result: Record = {}; + for (const [k, v] of Object.entries(value).slice(0, 128)) { + const safe = visit(v, k, depth + 1); + if (safe !== undefined) result[k.slice(0, 200)] = safe; + } + return result; + }; + const result = visit(input, "", 0); + return result && typeof result === "object" && !Array.isArray(result) + ? (result as Record) + : {}; +} + +export interface OtlpValue { + stringValue?: string; + intValue?: string | number; + doubleValue?: number; + boolValue?: boolean; + arrayValue?: { values?: OtlpValue[] }; + kvlistValue?: { values?: OtlpAttribute[] }; +} +export interface OtlpAttribute { + key?: string; + value?: OtlpValue; +} + +export function decodeOtlpAttributes( + attributes: OtlpAttribute[] | undefined, +): Record { + const decode = (v: OtlpValue | undefined, depth = 0): unknown => { + if (!v || depth > 6) return undefined; + if (v.stringValue !== undefined) return v.stringValue; + if (v.intValue !== undefined) return Number(v.intValue); + if (v.doubleValue !== undefined) return v.doubleValue; + if (v.boolValue !== undefined) return v.boolValue; + if (v.arrayValue) + return (v.arrayValue.values ?? []) + .slice(0, 64) + .map((value) => decode(value, depth + 1)); + if (v.kvlistValue) + return Object.fromEntries( + (v.kvlistValue.values ?? []) + .slice(0, 128) + .filter( + (a) => + a.key && !/^(?:__proto__|constructor|prototype)$/.test(a.key), + ) + .map((a) => [a.key!, decode(a.value, depth + 1)]), + ); + return undefined; + }; + return sanitizeRuntimeContext( + Object.fromEntries( + (attributes ?? []) + .slice() + .sort( + (a, b) => + Number(/^autter\.(?:operation|event)\./.test(b.key ?? "")) - + Number(/^autter\.(?:operation|event)\./.test(a.key ?? "")), + ) + .slice(0, 128) + .filter( + (a) => a.key && !/^(?:__proto__|constructor|prototype)$/.test(a.key), + ) + .map((a) => [a.key!, decode(a.value)]), + ), + ); +} diff --git a/packages/otlp-ingester/src/continuous-detection.test.ts b/packages/otlp-ingester/src/continuous-detection.test.ts index 19cadb2..9beb73c 100644 --- a/packages/otlp-ingester/src/continuous-detection.test.ts +++ b/packages/otlp-ingester/src/continuous-detection.test.ts @@ -50,7 +50,7 @@ test("sampled handled exception markers survive OTLP normalization", () => { { key: "autter.sampled", value: { boolValue: true } }, ] }], }] }] }] }); - assert.deepEqual(result.occurrences[0]?.attributes, { "autter.handled": true, "autter.sampled": true }); + assert.deepEqual(result.occurrences[0]?.attributes, { "exception.type": "ValueError", "autter.handled": true, "autter.sampled": true }); }); test("browser timing is aggregated; failed outcome becomes an issue", () => { diff --git a/packages/otlp-ingester/src/logs.test.ts b/packages/otlp-ingester/src/logs.test.ts new file mode 100644 index 0000000..fdd637f --- /dev/null +++ b/packages/otlp-ingester/src/logs.test.ts @@ -0,0 +1,153 @@ +import assert from "node:assert/strict"; +import { test } from "node:test"; +import { normalizeLogs, logTableDDL } from "./logs.js"; +import { normalizeTraces } from "./normalize-otlp.js"; +import { decodeLogsRequest } from "./otlp-proto.js"; + +const attr = (key: string, value: string) => ({ + key, + value: { stringValue: value }, +}); +test("OTLP context survives ingestion while sensitive fields are scrubbed", () => { + const result = normalizeTraces({ + resourceSpans: [ + { + resource: { attributes: [attr("service.name", "checkout")] }, + scopeSpans: [ + { + spans: [ + { + traceId: "a".repeat(32), + spanId: "b".repeat(16), + name: "checkout", + startTimeUnixNano: "1760000000000000000", + endTimeUnixNano: "1760000000001000000", + attributes: [ + attr("autter.operation.id", "op-1"), + attr("payment.provider", "stripe"), + attr("token", "secret"), + ], + events: [ + { + name: "autter.outcome", + attributes: [ + attr("autter.outcome.status", "error"), + attr("autter.outcome.name", "checkout"), + attr( + "autter.outcome.stack", + "Error\n at checkout (/app/checkout.ts:20:1)", + ), + ], + }, + ], + }, + ], + }, + ], + }, + ], + }); + assert.equal(result.spans[0]!.attributes!["payment.provider"], "stripe"); + assert.equal( + result.occurrences[0]!.attributes!["autter.operation.id"], + "op-1", + ); + assert.equal(result.occurrences[0]!.attributes!.token, "[redacted]"); + assert.match(result.occurrences[0]!.stack!, /checkout.ts/); +}); +test("wide event logs preserve nested context, stable IDs and native trace IDs", () => { + const payload = { + resourceLogs: [ + { + resource: { attributes: [attr("service.name", "api")] }, + scopeLogs: [ + { + logRecords: [ + { + timeUnixNano: "1760000000000000000", + traceId: "a".repeat(32), + body: { + stringValue: JSON.stringify({ + payment: { attempts: 3, token: "private" }, + message: "person@example.com", + }), + }, + attributes: [ + attr("autter.event.type", "operation"), + attr("autter.operation.outcome", "failed"), + ], + }, + ], + }, + ], + }, + ], + }; + const [row] = normalizeLogs(payload); + assert.deepEqual(row!.attributes.payment, { + attempts: 3, + token: "[redacted]", + }); + assert.equal(row!.message, "[redacted]"); + assert.equal(row!.traceId, "a".repeat(32)); + assert.equal(row!.id, normalizeLogs(payload)[0]!.id); + assert.match(logTableDDL("test"), /org_id String, repository_id String/); + assert.throws(() => normalizeLogs({ resourceLogs: "bad" } as never)); +}); +test("OTLP protobuf logs are decoded with the standard field numbers", () => { + // resourceLogs(1) -> scopeLogs(2) -> logRecords(2) -> body(5) -> stringValue(1) + const decoded = decodeLogsRequest( + Buffer.from([10, 10, 18, 8, 18, 6, 42, 4, 10, 2, 111, 107]), + ); + assert.equal(normalizeLogs(decoded)[0]!.message, "ok"); +}); + +test("serialized context and path queries are scrubbed and excess batches fail explicitly", () => { + const payload = { + resourceLogs: [ + { + scopeLogs: [ + { + logRecords: [ + { + body: { + kvlistValue: { + values: [ + attr("message", "failed"), + attr("details", '{"password":"unsafe","attempts":2}'), + attr("path", "/checkout?token=unsafe"), + ], + }, + }, + }, + ], + }, + ], + }, + ], + }; + const [row] = normalizeLogs(payload); + assert.equal(row!.message, "failed"); + assert.deepEqual(row!.attributes.details, { + password: "[redacted]", + attempts: 2, + }); + assert.equal(row!.attributes.path, "/checkout"); + assert.throws( + () => + normalizeLogs({ + resourceLogs: [ + { + scopeLogs: [ + { + logRecords: Array.from({ length: 2001 }, () => ({ + body: { stringValue: "x" }, + })), + }, + ], + }, + ], + }), + /too many log records/, + ); +}); diff --git a/packages/otlp-ingester/src/logs.ts b/packages/otlp-ingester/src/logs.ts new file mode 100644 index 0000000..ec0ffaf --- /dev/null +++ b/packages/otlp-ingester/src/logs.ts @@ -0,0 +1,200 @@ +import { createHash } from "node:crypto"; +import { + decodeOtlpAttributes, + sanitizeRuntimeContext, + type OtlpAttribute, + type OtlpValue, +} from "./context.js"; +import { normalizeRoute } from "./fingerprint.js"; + +export interface RuntimeLogRecord { + id: string; + service: string; + environment: string; + release: string; + traceId: string; + spanId: string; + operationId: string; + operation: string; + type: "log" | "operation"; + severity: string; + message: string; + outcome: string; + durationMs: number; + attributes: Record; + occurredAt: Date; +} +export interface OtlpLogsRequest { + resourceLogs?: Array<{ + resource?: { attributes?: OtlpAttribute[] }; + scopeLogs?: Array<{ + scope?: { name?: string }; + logRecords?: Array<{ + timeUnixNano?: string | number; + observedTimeUnixNano?: string | number; + severityNumber?: number; + severityText?: string; + traceId?: string; + spanId?: string; + body?: OtlpValue; + attributes?: OtlpAttribute[]; + }>; + }>; + }>; +} + +export function normalizeLogs(request: OtlpLogsRequest): RuntimeLogRecord[] { + if (!request || !Array.isArray(request.resourceLogs)) + throw new Error("resourceLogs must be an array"); + const rows: RuntimeLogRecord[] = []; + if (request.resourceLogs.length > 128) + throw new Error("too many log resources"); + for (const resource of request.resourceLogs) { + if ( + !resource || + (resource.scopeLogs && !Array.isArray(resource.scopeLogs)) || + (resource.scopeLogs?.length ?? 0) > 128 + ) + throw new Error("invalid log scopes"); + const info = decodeOtlpAttributes(resource.resource?.attributes); + for (const scope of resource.scopeLogs ?? []) { + if (!scope || (scope.logRecords && !Array.isArray(scope.logRecords))) + throw new Error("invalid log records"); + for (const record of scope.logRecords ?? []) { + if (!record || typeof record !== "object") + throw new Error("invalid log record"); + if (rows.length >= 2000) throw new Error("too many log records"); + const attrs = decodeOtlpAttributes(record.attributes); + let event: Record = {}; + const decodedBody = decodeOtlpAttributes([ + { key: "event", value: record.body }, + ]).event; + const body = + typeof decodedBody === "string" + ? decodedBody + : decodedBody == null + ? "" + : JSON.stringify(decodedBody); + if ( + decodedBody && + typeof decodedBody === "object" && + !Array.isArray(decodedBody) + ) + event = sanitizeRuntimeContext(decodedBody); + const metadata = Object.fromEntries( + Object.entries(attrs).filter(([key]) => + /^autter\.(?:operation|event)\./.test(key), + ), + ); + const attributes = sanitizeRuntimeContext({ + ...metadata, + ...event, + ...attrs, + }); + const nanos = record.timeUnixNano ?? record.observedTimeUnixNano; + const occurredAt = + nanos === undefined + ? new Date() + : new Date(Number(BigInt(String(nanos)) / 1_000_000n)); + if (!Number.isFinite(occurredAt.getTime())) + throw new Error("invalid log timestamp"); + const number = Number(record.severityNumber ?? 9); + const text = String( + record.severityText ?? event.level ?? "", + ).toLowerCase(); + const severity = + number >= 21 || text === "fatal" + ? "fatal" + : number >= 17 || text === "error" + ? "error" + : number >= 13 || /warn/.test(text) + ? "warning" + : number < 9 || text === "debug" + ? "debug" + : "info"; + const errorMessage = + event.error && typeof event.error === "object" + ? (event.error as Record).message + : event.error; + const message = String( + sanitizeRuntimeContext({ + message: event.message ?? errorMessage ?? body, + }).message ?? "", + ).slice(0, 4000); + const operationId = String( + attributes["autter.operation.id"] ?? "", + ).slice(0, 128); + const traceId = /^[a-f0-9]{32}$/i.test(record.traceId ?? "") + ? record.traceId!.toLowerCase().replace(/^0+$/, "") + : ""; + const spanId = /^[a-f0-9]{16}$/i.test(record.spanId ?? "") + ? record.spanId!.toLowerCase().replace(/^0+$/, "") + : ""; + if (typeof attributes.path === "string") + attributes.path = normalizeRoute(attributes.path.split(/[?#]/)[0]!); + rows.push({ + id: createHash("sha256") + .update( + JSON.stringify([ + info, + nanos, + traceId, + spanId, + severity, + attributes, + message, + ]), + ) + .digest("hex") + .slice(0, 32), + service: String(info["service.name"] ?? "unknown").slice(0, 200), + environment: String( + info["deployment.environment.name"] ?? + info["deployment.environment"] ?? + "production", + ).slice(0, 100), + release: String(info["service.version"] ?? "").slice(0, 200), + traceId, + spanId, + operationId, + operation: String(attributes["autter.operation.name"] ?? "").slice( + 0, + 200, + ), + type: + attributes["autter.event.type"] === "operation" + ? "operation" + : "log", + severity, + message, + outcome: String(attributes["autter.operation.outcome"] ?? "").slice( + 0, + 30, + ), + durationMs: Math.max( + 0, + Number(attributes["autter.operation.duration_ms"]) || 0, + ), + attributes, + occurredAt, + }); + } + } + } + return rows; +} + +export function logTableDDL(db: string): string { + return `CREATE TABLE IF NOT EXISTS ${db}.runtime_logs ( + org_id String, repository_id String, event_id String, + service LowCardinality(String), environment LowCardinality(String), release String DEFAULT '', + trace_id String DEFAULT '', span_id String DEFAULT '', operation_id String DEFAULT '', + operation String DEFAULT '', event_type LowCardinality(String) DEFAULT 'log', + severity LowCardinality(String) DEFAULT 'info', message String CODEC(ZSTD(1)), + outcome LowCardinality(String) DEFAULT '', duration_ms Float64 DEFAULT 0, + attributes String DEFAULT '{}' CODEC(ZSTD(1)), occurred_at DateTime64(3, 'UTC'), + ingested_at DateTime64(3, 'UTC') DEFAULT now64(3) + ) ENGINE = ReplacingMergeTree PARTITION BY toDate(occurred_at) + ORDER BY (org_id, repository_id, occurred_at, event_id) + TTL toDateTime(occurred_at) + INTERVAL 14 DAY`; +} diff --git a/packages/otlp-ingester/src/migrations.ts b/packages/otlp-ingester/src/migrations.ts index ccfda37..d1f0a9f 100644 --- a/packages/otlp-ingester/src/migrations.ts +++ b/packages/otlp-ingester/src/migrations.ts @@ -36,6 +36,7 @@ import { latencyTableDDL } from "./latency.js"; import { profileTableDDL } from "./profiles.js"; import { sourceMapTableDDL } from "./source-maps.js"; +import { logTableDDL } from "./logs.js"; import { memoryTableDDL, platformEventTableDDL } from "./memory.js"; export interface Migration { @@ -135,6 +136,7 @@ export const MIGRATIONS: Migration[] = [ { id: "0010-memory-temporality", statements: [ `ALTER TABLE {db}.runtime_memory_samples ADD COLUMN IF NOT EXISTS temporality LowCardinality(String) DEFAULT 'gauge'`, ] }, + { id: "0011-runtime-logs", statements: [logTableDDL("{db}")] }, ]; /** The tracking table itself — created by the runner before anything else. */ diff --git a/packages/otlp-ingester/src/normalize-otlp.ts b/packages/otlp-ingester/src/normalize-otlp.ts index 526a71c..e939916 100644 --- a/packages/otlp-ingester/src/normalize-otlp.ts +++ b/packages/otlp-ingester/src/normalize-otlp.ts @@ -1,3 +1,4 @@ +import { decodeOtlpAttributes, sanitizeRuntimeContext, type OtlpValue } from "./context.js"; import { normalizeRoute } from "./fingerprint.js"; import { extractLlmCall } from "./llm.js"; import { @@ -23,12 +24,7 @@ const EMAIL_VALUE_RE = /[A-Z0-9._%+-]+@[A-Z0-9.-]+\.[A-Z]{2,}/gi; interface OtlpKeyValue { key?: string; - value?: { - stringValue?: string; - intValue?: string | number; - doubleValue?: number; - boolValue?: boolean; - }; + value?: OtlpValue; } interface OtlpEvent { @@ -261,6 +257,7 @@ export function normalizeTraces(request: OtlpTraceRequest): NormalizedTraces { statusCode, durationMs, attributes: { + ...decodeOtlpAttributes(span.attributes), "http.request.method": (attrs.get("http.request.method") ?? attrs.get("http.method") ?? "").slice(0, 20), }, startedAt, @@ -294,10 +291,8 @@ export function normalizeTraces(request: OtlpTraceRequest): NormalizedTraces { statusCode, traceId: span.traceId ?? null, sessionId: null, - attributes: handled ? { - "autter.handled": true, - "autter.sampled": eventAttrs.get("autter.sampled") === "true", - } : null, + attributes: sanitizeRuntimeContext({ ...decodeOtlpAttributes(span.attributes), ...decodeOtlpAttributes(event.attributes), + ...(handled ? { "autter.handled": true, "autter.sampled": eventAttrs.get("autter.sampled") === "true" } : {}) }), occurredAt: nanosToDate(event.timeUnixNano ?? span.startTimeUnixNano), }); } @@ -314,8 +309,9 @@ export function normalizeTraces(request: OtlpTraceRequest): NormalizedTraces { environment: resource.environment, release: resource.release, errorType: "OutcomeFailure", message: `${name}: ${(outcome.get("autter.outcome.message") ?? "failed").replace(EMAIL_VALUE_RE, "[redacted]").slice(0, 1000)}`, - stack: null, route, method: methodOf(attrs), statusCode, - traceId: span.traceId ?? null, sessionId: null, attributes: null, + stack: String(decodeOtlpAttributes(event.attributes)["autter.outcome.stack"] ?? "") || null, route, method: methodOf(attrs), statusCode, + traceId: span.traceId ?? null, sessionId: null, + attributes: sanitizeRuntimeContext({ ...decodeOtlpAttributes(span.attributes), ...decodeOtlpAttributes(event.attributes) }), occurredAt: nanosToDate(event.timeUnixNano ?? span.startTimeUnixNano), }); } @@ -337,7 +333,7 @@ export function normalizeTraces(request: OtlpTraceRequest): NormalizedTraces { statusCode, traceId: span.traceId ?? null, sessionId: null, - attributes: null, + attributes: decodeOtlpAttributes(span.attributes), occurredAt: startedAt, }); } diff --git a/packages/otlp-ingester/src/otlp-proto.ts b/packages/otlp-ingester/src/otlp-proto.ts index 5131765..9fec7f0 100644 --- a/packages/otlp-ingester/src/otlp-proto.ts +++ b/packages/otlp-ingester/src/otlp-proto.ts @@ -1,5 +1,6 @@ import protobuf from "protobufjs"; import type { OtlpMetricsRequest, OtlpTraceRequest } from "./normalize-otlp.js"; +import type { OtlpLogsRequest } from "./logs.js"; /** * OTLP/HTTP protobuf decode (`content-type: application/x-protobuf`) — @@ -24,10 +25,27 @@ message AnyValue { bool bool_value = 2; int64 int_value = 3; double double_value = 4; + ArrayValue array_value = 5; + KeyValueList kvlist_value = 6; } } +message ArrayValue { repeated AnyValue values = 1; } +message KeyValueList { repeated KeyValue values = 1; } message KeyValue { string key = 1; AnyValue value = 2; } message Resource { repeated KeyValue attributes = 1; } +message ExportLogsServiceRequest { repeated ResourceLogs resource_logs = 1; } +message ResourceLogs { Resource resource = 1; repeated ScopeLogs scope_logs = 2; } +message ScopeLogs { repeated LogRecord log_records = 2; } +message LogRecord { + fixed64 time_unix_nano = 1; + int32 severity_number = 2; + string severity_text = 3; + AnyValue body = 5; + repeated KeyValue attributes = 6; + bytes trace_id = 9; + bytes span_id = 10; + fixed64 observed_time_unix_nano = 11; +} message ExportTraceServiceRequest { repeated ResourceSpans resource_spans = 1; } message ResourceSpans { Resource resource = 1; repeated ScopeSpans scope_spans = 2; } @@ -84,6 +102,7 @@ message HistogramDataPoint { const root = protobuf.parse(PROTO).root; const TraceRequest = root.lookupType("otlp.ExportTraceServiceRequest"); const MetricsRequest = root.lookupType("otlp.ExportMetricsServiceRequest"); +const LogsRequest = root.lookupType("otlp.ExportLogsServiceRequest"); const TO_OBJECT_OPTIONS: protobuf.IConversionOptions = { longs: String, // 64-bit ints → strings (matches OTLP/JSON) @@ -122,3 +141,7 @@ export function decodeMetricsRequest(body: Buffer): OtlpMetricsRequest { MetricsRequest.toObject(message, TO_OBJECT_OPTIONS), ) as OtlpMetricsRequest; } + +export function decodeLogsRequest(body: Buffer): OtlpLogsRequest { + return hexifyIds(LogsRequest.toObject(LogsRequest.decode(body), TO_OBJECT_OPTIONS)) as OtlpLogsRequest; +} diff --git a/packages/otlp-ingester/src/server.logs.test.ts b/packages/otlp-ingester/src/server.logs.test.ts new file mode 100644 index 0000000..d6fe642 --- /dev/null +++ b/packages/otlp-ingester/src/server.logs.test.ts @@ -0,0 +1,132 @@ +import assert from "node:assert/strict"; +import { createServer, type Server } from "node:http"; +import type { AddressInfo } from "node:net"; +import { test } from "node:test"; +import type { IngesterConfig } from "./config.js"; +import { createIngesterApp } from "./server.js"; + +const listen = async (server: Server) => { + await new Promise((resolve) => server.listen(0, "127.0.0.1", resolve)); + return `http://127.0.0.1:${(server.address() as AddressInfo).port}`; +}; +const close = (server: Server) => + new Promise((resolve) => { + server.close(resolve); + server.closeAllConnections(); + }); + +test("log ingestion authenticates, preserves tenant mapping and refuses malformed or undelivered records", async () => { + const inserts: Array> = []; + let fail = false; + const ch = createServer((req, res) => { + let body = ""; + req.on("data", (chunk) => (body += chunk)); + req.on("end", () => { + const query = + new URL(req.url!, "http://localhost").searchParams.get("query") ?? ""; + if (query.includes("INSERT INTO") && query.includes("runtime_logs")) { + if (fail) { + res.writeHead(400).end("storage rejected"); + return; + } + inserts.push( + ...body + .trim() + .split("\n") + .map((line) => JSON.parse(line)), + ); + } + res.end(); + }); + }); + const chUrl = await listen(ch); + const config: IngesterConfig = { + port: 0, + clickhouseUrl: chUrl, + clickhouseUser: "default", + clickhousePassword: "", + clickhouseDatabase: "autter_runtime", + ingestKeys: [ + { key: "server", orgId: "org-owner", repositoryId: "repo-owner" }, + { + key: "client", + orgId: "org-owner", + repositoryId: "repo-owner", + scope: "client", + }, + ], + keyValidatorUrl: null, + keyValidatorToken: null, + sinkUrl: null, + sinkToken: null, + sinkMaxAttempts: 1, + sinkMaxBufferedBatches: 10, + sinkMaxBufferedMb: 2, + maxBodyBytes: 1024 * 1024, + rateLimitPerMinute: 300, + clientRateLimitPerMinute: 120, + occurrenceTtlDays: 14, + spanTtlDays: 7, + metricsTtlDays: 90, + llmCallTtlDays: 90, + }; + const app = createIngesterApp(config).app.listen(); + const url = await new Promise((resolve) => + app.once("listening", () => + resolve(`http://127.0.0.1:${(app.address() as AddressInfo).port}`), + ), + ); + const payload = { + resourceLogs: [ + { + resource: { + attributes: [{ key: "org_id", value: { stringValue: "other-org" } }], + }, + scopeLogs: [ + { + logRecords: [ + { + timeUnixNano: "1760000000000000000", + body: { stringValue: "checkout started" }, + }, + ], + }, + ], + }, + ], + }; + const post = (key: string, body: unknown) => + fetch(`${url}/v1/logs`, { + method: "POST", + headers: { + "content-type": "application/json", + authorization: `Bearer ${key}`, + }, + body: JSON.stringify(body), + }); + try { + assert.equal((await post("unknown", payload)).status, 401); + assert.equal((await post("client", payload)).status, 403); + assert.equal((await post("server", { resourceLogs: "bad" })).status, 400); + assert.equal((await post("server", payload)).status, 200); + assert.equal(inserts.length, 1); + assert.equal(inserts[0]!.org_id, "org-owner"); + assert.equal(inserts[0]!.repository_id, "repo-owner"); + const proto = await fetch(`${url}/v1/logs`, { + method: "POST", + headers: { + "content-type": "application/x-protobuf", + authorization: "Bearer server", + }, + body: Buffer.from([10, 10, 18, 8, 18, 6, 42, 4, 10, 2, 111, 107]), + }); + assert.equal(proto.status, 200); + assert.equal(inserts.length, 2); + assert.equal(inserts[1]!.message, "ok"); + fail = true; + assert.equal((await post("server", payload)).status, 503); + } finally { + await close(app); + await close(ch); + } +}); diff --git a/packages/otlp-ingester/src/server.ts b/packages/otlp-ingester/src/server.ts index 5e262de..24d94c4 100644 --- a/packages/otlp-ingester/src/server.ts +++ b/packages/otlp-ingester/src/server.ts @@ -22,7 +22,8 @@ import { type OtlpMetricsRequest, type OtlpTraceRequest, } from "./normalize-otlp.js"; -import { decodeMetricsRequest, decodeTraceRequest } from "./otlp-proto.js"; +import { decodeMetricsRequest, decodeTraceRequest, decodeLogsRequest } from "./otlp-proto.js"; +import { normalizeLogs, type OtlpLogsRequest } from "./logs.js"; import { decodeProfile } from "./profiles.js"; import { normalizeMemoryMetrics, normalizePlatformEvent, platformEventSchema } from "./memory.js"; import { validateSourceMap } from "./source-maps.js"; @@ -291,6 +292,19 @@ export function createIngesterApp(config: IngesterConfig): IngesterApp { otlpSuccess(req, res); }); + app.post("/v1/logs", async (req, res) => { + const ctx = await authenticate(req, res, "otlp"); + if (!ctx) return; + let logs; + try { logs = normalizeLogs(req.is("application/x-protobuf") ? decodeLogsRequest(req.body as Buffer) : req.body as OtlpLogsRequest); } + catch { res.status(400).json({ error: "invalid OTLP logs payload" }); return; } + try { await store.insertLogs(ctx, logs); } + catch (err) { storageError(res, err); return; } + // Logs are diagnostic evidence. Exceptions and failed outcomes use the trace sink, + // so one operation does not create duplicate issues through two export paths. + otlpSuccess(req, res); + }); + app.post("/v1/metrics", async (req, res) => { const ctx = await authenticate(req, res, "otlp"); if (!ctx) return; diff --git a/packages/runtime-next/README.md b/packages/runtime-next/README.md index d3bff4b..8cec646 100644 --- a/packages/runtime-next/README.md +++ b/packages/runtime-next/README.md @@ -92,3 +92,14 @@ calls with `experimental_telemetry: { isEnabled: true }` (and anything wrapped in `withLlmCall` or an `instrumentLlmClient` client) are recorded at 100% — every call, with model, tokens, latency, and cost. See the [`@autter/runtime-node` README](../runtime-node) for the API. +## Operation logging (1.4.0+) + +Import `withRuntimeOperation`, `runtimeLogger`, `createRuntimeLogger`, +`flushRuntimeLogs`, and `runtimeLogStats` from `@autter/runtime-next/server`. +Initialize through the existing `registerAutter` call in `instrumentation.ts`. +These APIs require the Node runtime, a 1.4.0+ ingester, and matching platform +evidence support. Client components keep `@autter/runtime-next/client`. +See the [operation logging guide](../../docs/OPERATION-LOGGING.md) for measured +steps, explicit business outcomes, privacy and flushing. Ordinary error logs +are diagnostics; captured exceptions and declared failed outcomes use tracing +for issue grouping. diff --git a/packages/runtime-next/package.json b/packages/runtime-next/package.json index b66ca77..67205e9 100644 --- a/packages/runtime-next/package.json +++ b/packages/runtime-next/package.json @@ -1,6 +1,6 @@ { "name": "@autter/runtime-next", - "version": "1.3.4", + "version": "1.4.0", "description": "One-command Autter Runtime for Next.js: server OTel, browser tracker, relay route, and React error boundary", "license": "MIT", "type": "module", @@ -37,7 +37,7 @@ }, "dependencies": { "@autter/runtime-browser": "^1.3.3", - "@autter/runtime-node": "^1.3.4" + "@autter/runtime-node": "^1.4.0" }, "peerDependencies": { "react": ">=18" diff --git a/packages/runtime-next/src/server.ts b/packages/runtime-next/src/server.ts index ed4ae56..4f8955a 100644 --- a/packages/runtime-next/src/server.ts +++ b/packages/runtime-next/src/server.ts @@ -34,6 +34,17 @@ import { } from "@autter/runtime-node"; export { + createRuntimeLogger, + runtimeLogger, + withRuntimeOperation, + flushRuntimeLogs, + runtimeLogStats, + type RuntimeOperation, + type RuntimeOutcome, + type RuntimeLogContext, + type RuntimeLogContextValue, + type RuntimeLogLevel, + type RuntimeLoggingOptions, captureException as captureServerException, captureMessage as captureServerMessage, reportOutcome as reportServerOutcome, @@ -46,6 +57,7 @@ export { installAutterAutoFlush, redactAttributes, } from "@autter/runtime-node"; + export type { LlmCallHandle, LlmCallInfo, diff --git a/packages/runtime-node/README.md b/packages/runtime-node/README.md index c90a7af..1403b01 100644 --- a/packages/runtime-node/README.md +++ b/packages/runtime-node/README.md @@ -1,3 +1,10 @@ +# Operation logging + +Use `withRuntimeOperation`, `runtimeLogger`, and `createRuntimeLogger` to attach +application context to logs and completed operations. See the full +[operation logging guide](../../docs/OPERATION-LOGGING.md) for initialization, +steps, outcomes, redaction, export lifecycle, and release requirements. + # @autter/runtime-node The server tracker exports per-instance RSS, heap, memory limit, and GC @@ -287,4 +294,4 @@ Note on the slow-process monitor: Autter flags HTTP routes from the unsampled request metrics, so route detection works out of the box. Non-HTTP work is only visible where a span exists — relying on 1%-sampled regular traces there would undercount ~100×, which is why these spans skip head -sampling. +sampling. \ No newline at end of file diff --git a/packages/runtime-node/package.json b/packages/runtime-node/package.json index df51c8b..b94a8aa 100644 --- a/packages/runtime-node/package.json +++ b/packages/runtime-node/package.json @@ -1,6 +1,6 @@ { "name": "@autter/runtime-node", - "version": "1.3.4", + "version": "1.4.0", "description": "Autter Runtime for Node.js: same-origin browser relay handler + curated OpenTelemetry server tracker", "license": "MIT", "type": "module", @@ -24,7 +24,7 @@ }, "scripts": { "build": "tsup src/index.ts --format esm,cjs --dts --target node20 --clean", - "test": "npm run build && node --test test/*.test.mjs" + "test": "npm run build && npm run build -w @autter/otlp-ingester && node --test test/*.test.mjs" }, "dependencies": { "@opentelemetry/api": "^1.9.0", diff --git a/packages/runtime-node/src/index.ts b/packages/runtime-node/src/index.ts index b7296f9..bc3603f 100644 --- a/packages/runtime-node/src/index.ts +++ b/packages/runtime-node/src/index.ts @@ -38,3 +38,16 @@ export { type InstrumentLlmOptions, } from "./llm-instrument.js"; export { startCaughtExceptionSampler } from "./caught-exceptions.js"; +export { + createRuntimeLogger, + runtimeLogger, + withRuntimeOperation, + flushRuntimeLogs, + runtimeLogStats, + type RuntimeLogContext, + type RuntimeLogContextValue, + type RuntimeLogLevel, + type RuntimeOutcome, + type RuntimeOperation, + type RuntimeLoggingOptions, +} from "./logger.js"; diff --git a/packages/runtime-node/src/logger.ts b/packages/runtime-node/src/logger.ts new file mode 100644 index 0000000..9abf427 --- /dev/null +++ b/packages/runtime-node/src/logger.ts @@ -0,0 +1,511 @@ +import { AsyncLocalStorage } from "node:async_hooks"; +import { randomUUID } from "node:crypto"; +import { trace, type Attributes } from "@opentelemetry/api"; +import { redactAttributes } from "./redact.js"; + +export type RuntimeLogContextValue = + | string + | number + | boolean + | null + | undefined + | RuntimeLogContextValue[] + | { [key: string]: RuntimeLogContextValue }; +export type RuntimeLogContext = Record; +export type RuntimeLogLevel = "debug" | "info" | "warning" | "error"; +export type RuntimeOutcome = + | "succeeded" + | "failed" + | "degraded" + | "cancelled" + | "pending"; +export interface RuntimeLoggingOptions { + console?: boolean; + minLevel?: RuntimeLogLevel; +} +export interface RuntimeOperation { + readonly id: string; + setContext(attributes: RuntimeLogContext): void; + outcome(status: RuntimeOutcome, message?: string): void; + step(name: string, fn: () => T | Promise): Promise; +} +interface OperationState { + id: string; + name: string; + parentId?: string; + attributes: RuntimeLogContext; + startedAt: number; + outcome: RuntimeOutcome; + message?: string; + reportingStack?: string; + sealed?: boolean; + steps: Array<{ name: string; status: string; durationMs: number }>; +} +interface LoggerConfig { + endpoint: string; + apiKey: string; + service: string; + environment: string; + release?: string; + options?: RuntimeLoggingOptions; + redact(attributes: Attributes): Attributes; + run( + name: string, + fn: () => Promise, + attributes: Attributes, + ): Promise; + reportOutcome(name: string, message: string, attributes: Attributes): void; +} +const local = new AsyncLocalStorage(); +let config: LoggerConfig | null = null; +interface QueuedLog { + record: Record; + bytes: number; +} +let queue: QueuedLog[] = []; +let queueBytes = 0; +let inFlightBytes = 0; +let timer: ReturnType | undefined; +let flushing: Promise | null = null; +let dropped = 0; +let inFlight = 0; +let stopping = false; +let failureStreak = 0; +const severity = { debug: 5, info: 9, warning: 13, error: 17 }; + +function safe(attributes: RuntimeLogContext = {}): RuntimeLogContext { + const redacted = + config?.redact(attributes as Attributes) ?? + redactAttributes(attributes as Attributes); + let budget = 512; + let characters = 16384; + let truncated = false; + const visit = (input: unknown, depth: number): unknown => { + if (--budget < 0 || depth > 6) { + truncated = true; + return "[truncated]"; + } + if (typeof input === "string") { + const text = input.replace(/(https?:\/\/[^\s?#]+)[?#][^\s]*/gi, "$1"); + const kept = text.slice(0, Math.max(0, characters)); + truncated ||= kept.length < text.length; + characters -= kept.length; + return kept; + } + if (typeof input === "number") + return Number.isFinite(input) ? input : undefined; + if (typeof input === "boolean" || input === null) return input; + if (Array.isArray(input)) { + truncated ||= input.length > 64; + return input.slice(0, 64).map((item) => visit(item, depth + 1)); + } + if (input && typeof input === "object") { + const entries = Object.entries(input); + if (depth === 0) { + const priority = (key: string) => + /^autter\.(?:operation|event)\./.test(key) + ? 2 + : /^exception\./.test(key) + ? 1 + : 0; + entries.sort(([a], [b]) => priority(b) - priority(a)); + } + truncated ||= entries.length > 100; + return Object.fromEntries( + entries + .slice(0, 100) + .filter(([key]) => !/^(?:__proto__|constructor|prototype)$/.test(key)) + .map(([key, item]) => [key.slice(0, 200), visit(item, depth + 1)]), + ); + } + return undefined; + }; + const result = visit(redacted, 0) as RuntimeLogContext; + if (truncated) result["autter.context.truncated"] = true; + return result; +} +function operationAttributes(state?: OperationState): RuntimeLogContext { + return state + ? { + "autter.operation.id": state.id, + "autter.operation.name": state.name, + ...(state.parentId + ? { "autter.operation.parent_id": state.parentId } + : {}), + ...state.attributes, + } + : {}; +} +function userAttributes(attributes: RuntimeLogContext = {}): RuntimeLogContext { + return Object.fromEntries( + Object.entries(safe(attributes)).filter( + ([key]) => !/^autter\.(?:operation|event)\./.test(key), + ), + ); +} +function mergeContext( + current: RuntimeLogContext, + next: RuntimeLogContext = {}, +): RuntimeLogContext { + const merge = ( + left: RuntimeLogContext, + right: RuntimeLogContext, + depth: number, + ): RuntimeLogContext => { + const result = { ...left }; + for (const [key, value] of Object.entries(right)) { + const prior = result[key]; + result[key] = + depth < 6 && + value && + typeof value === "object" && + !Array.isArray(value) && + prior && + typeof prior === "object" && + !Array.isArray(prior) + ? merge(prior, value, depth + 1) + : value; + } + return result; + }; + return userAttributes(merge(current, userAttributes(next), 0)); +} +function traceAttributes(state: OperationState): Attributes { + return Object.fromEntries( + Object.entries(operationAttributes(state)) + .filter(([, item]) => item !== null && item !== undefined) + .map(([key, item]) => [ + key, + typeof item === "object" ? JSON.stringify(item) : item, + ]), + ); +} +function scheduleFlush(): void { + if (!timer && config && !stopping && queue.length) { + timer = setTimeout( + () => { + timer = undefined; + void flushRuntimeLogs().catch(() => {}); + }, + Math.min(30000, 2000 * 2 ** Math.min(failureStreak, 4)), + ); + timer.unref(); + } +} +function value(input: unknown): Record { + if (typeof input === "number") return { doubleValue: input }; + if (typeof input === "boolean") return { boolValue: input }; + if (Array.isArray(input)) + return { arrayValue: { values: input.slice(0, 64).map(value) } }; + if (input && typeof input === "object") + return { + kvlistValue: { + values: Object.entries(input) + .slice(0, 128) + .map(([key, item]) => ({ key, value: value(item) })), + }, + }; + return { stringValue: String(input ?? "").slice(0, 32000) }; +} +function emit( + level: RuntimeLogLevel, + message: string, + attributes: RuntimeLogContext, + extra: Record = {}, +): void { + const state = local.getStore(); + if (state?.sealed) return; // work continuing after completion cannot mutate the emitted operation + if ( + extra["autter.event.type"] !== "operation" && + severity[level] < severity[config?.options?.minLevel ?? "debug"] + ) + return; + const span = trace.getActiveSpan()?.spanContext(); + const user = mergeContext(state?.attributes ?? {}, attributes); + const exceptions = Object.fromEntries( + Object.entries(user).filter(([key]) => key.startsWith("exception.")), + ); + const attrs = safe({ + "autter.event.id": randomUUID(), + ...extra, + ...exceptions, + ...operationAttributes(state), + ...user, + } as RuntimeLogContext); + const record = { + timeUnixNano: String(BigInt(Date.now()) * 1_000_000n), + severityNumber: severity[level], + severityText: level.toUpperCase(), + body: value(String(safe({ message }).message ?? "").slice(0, 4000)), + attributes: Object.entries(attrs).map(([key, item]) => ({ + key, + value: value(item), + })), + ...(span && + /^[a-f0-9]{32}$/.test(span.traceId) && + !/^0+$/.test(span.traceId) + ? { traceId: span.traceId, spanId: span.spanId } + : {}), + }; + if (config?.options?.console !== false) + console.log( + JSON.stringify({ ...attrs, level, message: record.body.stringValue }), + ); + if (!config) return; + if (stopping) { + dropped++; + return; + } + const bytes = Buffer.byteLength(JSON.stringify(record)); + if ( + bytes > 256 * 1024 || + queue.length + inFlight >= 1000 || + queueBytes + inFlightBytes + bytes > 4 * 1024 * 1024 + ) { + dropped++; + console.warn( + "[autter-runtime] log buffer or record limit reached; record dropped", + ); + return; + } + queue.push({ record, bytes }); + queueBytes += bytes; + scheduleFlush(); +} + +export function configureRuntimeLogger(next: LoggerConfig): void { + config = next; + stopping = false; + failureStreak = 0; +} +export function runtimeLogStats(): { buffered: number; dropped: number } { + return { buffered: queue.length + inFlight, dropped }; +} +export async function flushRuntimeLogs(): Promise { + if (flushing) return flushing; + if (!config || !queue.length) return; + if (timer) { + clearTimeout(timer); + timer = undefined; + } + const activeConfig = config; + const deadline = Date.now() + 10000; + flushing = (async () => { + const count = queue.length; + for (let remaining = count; remaining > 0; ) { + const batch: QueuedLog[] = []; + inFlightBytes = 0; + while (batch.length < Math.min(50, remaining) && queue.length) { + const next = queue[0]!; + if (batch.length && inFlightBytes + next.bytes > 512 * 1024) break; + batch.push(queue.shift()!); + inFlightBytes += next.bytes; + queueBytes -= next.bytes; + } + remaining -= batch.length; + inFlight = batch.length; + const body = JSON.stringify({ + resourceLogs: [ + { + resource: { + attributes: Object.entries({ + "service.name": activeConfig.service, + "deployment.environment.name": activeConfig.environment, + ...(activeConfig.release + ? { "service.version": activeConfig.release } + : {}), + }).map(([key, item]) => ({ key, value: value(item) })), + }, + scopeLogs: [ + { + scope: { name: "autter-runtime" }, + logRecords: batch.map((entry) => entry.record), + }, + ], + }, + ], + }); + let failure: unknown; + for (let attempt = 0; attempt < 3; attempt++) { + try { + if (Date.now() >= deadline) + throw new Error("Runtime log flush exceeded 10 seconds"); + const response = await fetch(`${activeConfig.endpoint}/v1/logs`, { + method: "POST", + headers: { + "content-type": "application/json", + authorization: `Bearer ${activeConfig.apiKey}`, + }, + body, + signal: AbortSignal.timeout( + Math.max(1, Math.min(3000, deadline - Date.now())), + ), + }); + if (!response.ok) + throw new Error(`Runtime log export failed (${response.status})`); + failure = undefined; + break; + } catch (error) { + failure = error; + } + } + inFlight = 0; + if (failure) { + failureStreak++; + queue = [...batch, ...queue]; + queueBytes += inFlightBytes; + inFlightBytes = 0; + throw failure; + } + inFlightBytes = 0; + failureStreak = 0; + } + })().finally(() => { + flushing = null; + scheduleFlush(); + }); + return flushing; +} +export async function shutdownRuntimeLogger(): Promise { + stopping = true; + try { + await flushRuntimeLogs(); + } finally { + if (queue.length) { + dropped += queue.length; + console.warn( + `[autter-runtime] ${queue.length} log records could not be delivered before shutdown`, + ); + } + queue = []; + queueBytes = 0; + config = null; + if (timer) clearTimeout(timer); + timer = undefined; + } +} + +export function createRuntimeLogger(attributes: RuntimeLogContext = {}) { + const attrs = userAttributes(attributes); + return { + debug: (message: string, extra?: RuntimeLogContext) => + emit("debug", message, mergeContext(attrs, extra)), + info: (message: string, extra?: RuntimeLogContext) => + emit("info", message, mergeContext(attrs, extra)), + warn: (message: string, extra?: RuntimeLogContext) => + emit("warning", message, mergeContext(attrs, extra)), + error: (error: unknown, extra?: RuntimeLogContext) => + emit("error", error instanceof Error ? error.message : String(error), { + ...attrs, + ...extra, + ...(error instanceof Error + ? { + "exception.type": error.name, + "exception.stacktrace": error.stack ?? "", + } + : {}), + }), + }; +} +export const runtimeLogger = createRuntimeLogger(); + +/** Context stays local to this operation; queues/processes require explicit propagation. */ +export function withRuntimeOperation( + name: string, + fn: (operation: RuntimeOperation) => T | Promise, + attributes: RuntimeLogContext = {}, +): Promise { + const parent = local.getStore(); + const state: OperationState = { + id: randomUUID(), + name: String(safe({ name }).name).slice(0, 200), + parentId: parent?.sealed ? undefined : parent?.id, + attributes: mergeContext( + parent?.sealed ? {} : (parent?.attributes ?? {}), + attributes, + ), + startedAt: Date.now(), + outcome: "succeeded", + steps: [], + }; + const operation: RuntimeOperation = { + id: state.id, + setContext: (attrs) => { + if (!state.sealed) { + state.attributes = mergeContext(state.attributes, attrs); + trace.getActiveSpan()?.setAttributes(traceAttributes(state)); + } + }, + outcome: (status, message) => { + if (state.sealed) return; + state.outcome = status; + state.message = message; + if (status === "failed") + state.reportingStack = + new Error("Operation outcome reported here").stack ?? ""; + }, + step: async (stepName, stepFn) => { + const started = Date.now(); + let status = "succeeded"; + try { + return await stepFn(); + } catch (error) { + status = "failed"; + throw error; + } finally { + if (!state.sealed && state.steps.length < 64) + state.steps.push({ + name: String(safe({ name: stepName }).name).slice(0, 200), + status, + durationMs: Date.now() - started, + }); + } + }, + }; + const run = () => + local.run(state, async () => { + let thrown: unknown; + try { + return await fn(operation); + } catch (error) { + state.outcome = "failed"; + thrown = error; + throw error; + } finally { + trace.getActiveSpan()?.setAttributes({ + ...traceAttributes(state), + "autter.operation.outcome": state.outcome, + }); + if (state.outcome === "failed" && !thrown) + config?.reportOutcome( + state.name, + state.message ?? "Operation failed", + { + ...traceAttributes(state), + "autter.outcome.stack": state.reportingStack ?? "", + }, + ); + emit( + state.outcome === "failed" ? "error" : "info", + state.message ?? `${state.name}: ${state.outcome}`, + { + ...operationAttributes(state), + ...(thrown instanceof Error + ? { + "exception.type": thrown.name, + "exception.stacktrace": thrown.stack ?? "", + } + : {}), + }, + { + "autter.event.type": "operation", + "autter.operation.outcome": state.outcome, + "autter.operation.duration_ms": Date.now() - state.startedAt, + "autter.operation.steps": state.steps, + }, + ); + state.sealed = true; + } + }); + return config ? config.run(state.name, run, traceAttributes(state)) : run(); +} diff --git a/packages/runtime-node/src/server.ts b/packages/runtime-node/src/server.ts index 3667385..883b1a4 100644 --- a/packages/runtime-node/src/server.ts +++ b/packages/runtime-node/src/server.ts @@ -57,6 +57,7 @@ import { type RedactOptions, } from "./redact.js"; import { startMemoryMetrics } from "./memory.js"; +import { configureRuntimeLogger, flushRuntimeLogs, shutdownRuntimeLogger, type RuntimeLoggingOptions } from "./logger.js"; /** * Curated OpenTelemetry setup for Autter Runtime — deliberately NOT the @@ -84,6 +85,8 @@ import { startMemoryMetrics } from "./memory.js"; */ export interface AutterServerOptions { + /** Structured logs and completed operation summaries. */ + logging?: RuntimeLoggingOptions; /** Private ingest key (autter_rt_…). */ apiKey: string; /** Ingester base URL. Default: https://otlp.autter.dev */ @@ -957,6 +960,11 @@ export function initAutterServer(options: AutterServerOptions): AutterServer { span.setStatus({ code: SpanStatusCode.ERROR, message: safeMessage }); if (!reuseActive) span.end(); } + configureRuntimeLogger({ endpoint, apiKey: options.apiKey, service: options.service, environment, + release: options.release, options: options.logging, redact: activeRedactor, + run: (name, fn, attributes) => runWithSpan(processTracer, name, fn, activeRedactor(attributes)), + reportOutcome, + }); if (options.captureGlobalErrors !== false) { // `uncaughtExceptionMonitor` observes crashes WITHOUT changing the @@ -977,6 +985,7 @@ export function initAutterServer(options: AutterServerOptions): AutterServer { const flushTarget: FlushTarget = { forceFlush: async () => { const results = await Promise.allSettled([ + flushRuntimeLogs(), Promise.resolve().then(() => alwaysOnProvider.forceFlush()), Promise.resolve().then(() => mainSpanProcessor.forceFlush()), ...(errorTraceBuffer @@ -1010,6 +1019,7 @@ export function initAutterServer(options: AutterServerOptions): AutterServer { withLlmCall: (info, fn) => runLlmSpan(llmTracer, info, fn), trackLlmCall: (call) => recordLlmCall(llmTracer, call), shutdown: async () => { + const logShutdown = await Promise.allSettled([shutdownRuntimeLogger()]); stopMemoryMetrics?.(); active = null; activeAlwaysOnProvider = null; @@ -1018,6 +1028,7 @@ export function initAutterServer(options: AutterServerOptions): AutterServer { unregisterFlushTargets(); telemetryStats.markAllFlushed(); await Promise.allSettled([alwaysOnProvider.shutdown(), sdk.shutdown()]); + if (logShutdown[0]?.status === "rejected") throw logShutdown[0].reason; }, }; active = server; diff --git a/packages/runtime-node/test/logger-ingester.test.mjs b/packages/runtime-node/test/logger-ingester.test.mjs new file mode 100644 index 0000000..79ed3d1 --- /dev/null +++ b/packages/runtime-node/test/logger-ingester.test.mjs @@ -0,0 +1,100 @@ +import assert from "node:assert/strict"; +import { createServer } from "node:http"; +import { test } from "node:test"; +import { + initAutterServer, + withRuntimeOperation, + runtimeLogger, +} from "../dist/index.js"; +import { createIngesterApp } from "../../otlp-ingester/dist/server.js"; +const close = (server) => + new Promise((resolve) => { + server.close(resolve); + server.closeAllConnections(); + }); + +test("customer SDK operations survive the ingester and ClickHouse write boundary", async () => { + const inserts = []; + const ch = createServer((req, res) => { + let body = ""; + req.on("data", (chunk) => (body += chunk)); + req.on("end", () => { + const query = + new URL(req.url, "http://localhost").searchParams.get("query") ?? ""; + if (query.includes("INSERT INTO") && query.includes("runtime_logs")) + inserts.push(...body.trim().split("\n").map(JSON.parse)); + res.end(); + }); + }); + await new Promise((resolve) => ch.listen(0, "127.0.0.1", resolve)); + const { app } = createIngesterApp({ + port: 0, + clickhouseUrl: `http://127.0.0.1:${ch.address().port}`, + clickhouseUser: "default", + clickhousePassword: "", + clickhouseDatabase: "autter_runtime", + ingestKeys: [ + { key: "test-server", orgId: "org-sdk", repositoryId: "repo-sdk" }, + ], + keyValidatorUrl: null, + keyValidatorToken: null, + sinkUrl: null, + sinkToken: null, + maxBodyBytes: 1048576, + rateLimitPerMinute: 300, + clientRateLimitPerMinute: 120, + occurrenceTtlDays: 14, + spanTtlDays: 7, + metricsTtlDays: 90, + llmCallTtlDays: 90, + }); + const ingester = app.listen(0); + await new Promise((resolve) => ingester.once("listening", resolve)); + const runtime = initAutterServer({ + service: "checkout-api", + environment: "test", + release: "abc123", + apiKey: "test-server", + endpoint: `http://127.0.0.1:${ingester.address().port}`, + captureGlobalErrors: false, + autoFlush: false, + logging: { console: false, minLevel: "error" }, + }); + try { + const attrs = Object.fromEntries( + Array.from({ length: 200 }, (_, i) => [`custom.${i}`, i]), + ); + await withRuntimeOperation("checkout", async (operation) => { + operation.setContext({ + payment: { provider: "stripe" }, + ...attrs, + token: "secret-token", + "autter.operation.id": "spoofed", + }); + operation.setContext({ payment: { attempts: 2 } }); + await operation.step("reserve", async () => true); + runtimeLogger.info("filtered message"); + operation.outcome("failed", "Payment not confirmed"); + }); + await runtime.shutdown(); + assert.equal( + inserts.length, + 1, + "operation summary survives minLevel filtering", + ); + const summary = inserts[0]; + const context = JSON.parse(summary.attributes); + assert.equal(summary.org_id, "org-sdk"); + assert.equal(summary.repository_id, "repo-sdk"); + assert.equal(summary.operation, "checkout"); + assert.equal(summary.outcome, "failed"); + assert.deepEqual(context.payment, { provider: "stripe", attempts: 2 }); + assert.notEqual(summary.operation_id, "spoofed"); + assert.match(summary.trace_id, /^[a-f0-9]{32}$/); + assert.equal(context["autter.operation.steps"][0].name, "reserve"); + assert.equal(JSON.stringify(summary).includes("secret-token"), false); + } finally { + await close(ingester); + await close(ch); + } +}); diff --git a/packages/runtime-node/test/logger.test.mjs b/packages/runtime-node/test/logger.test.mjs new file mode 100644 index 0000000..eebf079 --- /dev/null +++ b/packages/runtime-node/test/logger.test.mjs @@ -0,0 +1,168 @@ +import assert from "node:assert/strict"; +import { createServer } from "node:http"; +import { test } from "node:test"; +import { + initAutterServer, + withRuntimeOperation, + runtimeLogger, + flushRuntimeLogs, +} from "../dist/index.js"; + +test("operation context isolates concurrent requests and exports outcomes, steps and redacted logs", async () => { + const requests = []; + const collector = createServer((req, res) => { + let body = ""; + req.on("data", (chunk) => (body += chunk)); + req.on("end", () => { + requests.push({ path: req.url, body: JSON.parse(body) }); + res.end("{}"); + }); + }); + await new Promise((resolve) => collector.listen(0, "127.0.0.1", resolve)); + const server = initAutterServer({ + service: "checkout", + apiKey: "test-key", + endpoint: `http://127.0.0.1:${collector.address().port}`, + captureGlobalErrors: false, + autoFlush: false, + logging: { console: false }, + }); + try { + await Promise.all( + ["a", "b"].map((id) => + withRuntimeOperation("checkout", async (op) => { + op.setContext({ + "checkout.id": id, + token: "secret-value", + note: "person@example.com", + }); + await op.step( + "reserve", + async () => + new Promise((resolve) => + setTimeout(resolve, id === "a" ? 15 : 1), + ), + ); + runtimeLogger.info("reservation complete"); + if (id === "b") op.outcome("failed", "payment timeout"); + }), + ), + ); + await assert.rejects( + withRuntimeOperation("throwing", async () => { + throw new Error("declined"); + }), + /declined/, + ); + await flushRuntimeLogs(); + const records = requests + .filter((r) => r.path === "/v1/logs") + .flatMap((r) => r.body.resourceLogs[0].scopeLogs[0].logRecords); + const attributes = (r) => + Object.fromEntries( + r.attributes.map((a) => [ + a.key, + a.value.stringValue ?? + a.value.doubleValue ?? + a.value.kvlistValue ?? + a.value.arrayValue, + ]), + ); + const summaries = records.filter( + (r) => attributes(r)["autter.event.type"] === "operation", + ); + assert.equal(summaries.length, 3); + const checkouts = summaries.filter( + (r) => attributes(r)["autter.operation.name"] === "checkout", + ); + assert.deepEqual( + checkouts.map((r) => attributes(r)["checkout.id"]).sort(), + ["a", "b"], + ); + assert.equal( + new Set(checkouts.map((r) => attributes(r)["autter.operation.id"])).size, + 2, + ); + assert.equal( + attributes(checkouts.find((r) => attributes(r)["checkout.id"] === "b"))[ + "autter.operation.outcome" + ], + "failed", + ); + assert.ok(checkouts.every((r) => /^[a-f0-9]{32}$/.test(r.traceId))); + assert.ok( + checkouts.every( + (r) => attributes(r)["autter.operation.steps"].values.length === 1, + ), + ); + assert.equal(JSON.stringify(records).includes("secret-value"), false); + assert.equal(JSON.stringify(records).includes("person@example.com"), false); + } finally { + await server.shutdown(); + await new Promise((resolve) => collector.close(resolve)); + } +}); + +test("export retries retain buffered records, enforce limits and report undelivered shutdown", async () => { + const { runtimeLogStats } = await import("../dist/index.js"); + let fail = true; + let attempts = 0; + const records = []; + const collector = createServer((req, res) => { + let body = ""; + req.on("data", (chunk) => (body += chunk)); + req.on("end", () => { + if (req.url === "/v1/logs") { + attempts++; + if (fail) { + res.writeHead(503).end(); + return; + } + records.push( + ...JSON.parse(body).resourceLogs[0].scopeLogs[0].logRecords, + ); + } + res.end("{}"); + }); + }); + await new Promise((resolve) => collector.listen(0, "127.0.0.1", resolve)); + const runtime = initAutterServer({ + service: "limits", + apiKey: "test-key", + endpoint: `http://127.0.0.1:${collector.address().port}`, + captureGlobalErrors: false, + autoFlush: false, + logging: { console: false }, + }); + try { + const before = runtimeLogStats().dropped; + for (let i = 0; i < 1001; i++) runtimeLogger.info(`record-${i}`); + assert.equal(runtimeLogStats().buffered, 1000); + assert.equal(runtimeLogStats().dropped, before + 1); + await assert.rejects(flushRuntimeLogs(), /Runtime log export failed/); + assert.equal(attempts, 3); + assert.equal(runtimeLogStats().buffered, 1000); + fail = false; + await flushRuntimeLogs(); + assert.equal(runtimeLogStats().buffered, 0); + assert.equal(records.length, 1000); + runtimeLogger.error(new Error("failure"), { huge: "x".repeat(50000) }); + await flushRuntimeLogs(); + const oversized = records.at(-1); + assert.ok(JSON.stringify(oversized).length < 40000); + assert.ok( + oversized.attributes.some( + (attribute) => + attribute.key === "autter.context.truncated" && + attribute.value.boolValue, + ), + ); + fail = true; + runtimeLogger.info("shutdown failure"); + await assert.rejects(runtime.shutdown(), /Runtime log export failed/); + assert.equal(runtimeLogStats().buffered, 0); + assert.equal(runtimeLogStats().dropped, before + 2); + } finally { + await new Promise((resolve) => collector.close(resolve)); + } +});