diff --git a/.env.example b/.env.example index 83a0bb70..1eb18584 100644 --- a/.env.example +++ b/.env.example @@ -78,3 +78,14 @@ LOG_DIR=./data/logs # rotating log ring for support bundles # SENTRY_DSN= # One-time SQLite import on first boot (1.x upgrade path). Leave unset normally. # SQLITE_MIGRATE_PATH=/data/snapotter.db + +# --- OpenTelemetry Distributed Tracing (enterprise only) --- +# Requires a valid enterprise license with distributed_tracing feature. +# Set OTEL_EXPORTER_OTLP_ENDPOINT to enable. All other vars are optional. +# Docs: https://opentelemetry.io/docs/specs/otel/configuration/sdk-environment-variables/ +# OTEL_EXPORTER_OTLP_ENDPOINT=http://localhost:4318 +# OTEL_EXPORTER_OTLP_PROTOCOL=http/protobuf +# OTEL_EXPORTER_OTLP_HEADERS= +# OTEL_SERVICE_NAME=snapotter-api +# OTEL_TRACES_SAMPLER=parentbased_always_on +# OTEL_TRACES_EXPORTER=otlp diff --git a/apps/api/package.json b/apps/api/package.json index e18cf14d..98b10fd1 100644 --- a/apps/api/package.json +++ b/apps/api/package.json @@ -4,9 +4,9 @@ "private": true, "type": "module", "scripts": { - "dev": "PORT=13490 tsx watch src/index.ts", + "dev": "PORT=13490 tsx watch --import ./src/tracing.ts src/index.ts", "build": "tsc", - "start": "tsx src/index.ts", + "start": "tsx --import ./src/tracing.ts src/index.ts", "lint": "biome check src/", "typecheck": "tsc --noEmit", "clean": "rm -rf dist", @@ -20,6 +20,17 @@ "@fastify/static": "^9.1.3", "@neplex/vectorizer": "^0.1.0", "@node-saml/node-saml": "^5.1.0", + "@opentelemetry/api": "^1.9.1", + "@opentelemetry/exporter-trace-otlp-http": "^0.219.0", + "@opentelemetry/instrumentation-aws-sdk": "^0.74.0", + "@opentelemetry/instrumentation-fastify": "^0.57.0", + "@opentelemetry/instrumentation-http": "^0.219.0", + "@opentelemetry/instrumentation-ioredis": "^0.67.0", + "@opentelemetry/instrumentation-pg": "^0.71.0", + "@opentelemetry/resources": "^2.8.0", + "@opentelemetry/sdk-node": "^0.219.0", + "@opentelemetry/sdk-trace-base": "^2.8.0", + "@opentelemetry/semantic-conventions": "^1.41.1", "@scalar/fastify-api-reference": "^1.57.5", "@sentry/node": "^10.55.0", "@snapotter/ai": "workspace:*", @@ -49,6 +60,7 @@ "papaparse": "^5.5.3", "pdfkit": "^0.18.0", "pg": "^8.21.0", + "pino": "^10.3.1", "pino-roll": "^4.0.0", "playwright": "^1.60.0", "posthog-node": "^5.35.9", diff --git a/apps/api/src/index.ts b/apps/api/src/index.ts index a9ab86a6..5c647777 100644 --- a/apps/api/src/index.ts +++ b/apps/api/src/index.ts @@ -1,9 +1,9 @@ import { randomUUID } from "node:crypto"; import { statfs } from "node:fs/promises"; -import { join } from "node:path"; import cookie from "@fastify/cookie"; import cors from "@fastify/cors"; import rateLimit from "@fastify/rate-limit"; +import { trace } from "@opentelemetry/api"; import { getDispatcherStatus, initDispatcher, isGpuAvailable } from "@snapotter/ai"; import { APP_VERSION } from "@snapotter/shared"; import { eq, sql } from "drizzle-orm"; @@ -21,7 +21,7 @@ import { captureException, initAnalytics, shutdownAnalytics } from "./lib/analyt import { shouldRunStartupCleanup } from "./lib/cleanup.js"; import { buildCsp } from "./lib/csp.js"; import { ensureAiDirs, recoverInterruptedInstalls } from "./lib/feature-status.js"; - +import { logger } from "./lib/logger.js"; import { requestDuration } from "./lib/metrics.js"; import { getSettingString } from "./lib/settings-helpers.js"; import { requirePermission } from "./permissions.js"; @@ -31,6 +31,7 @@ import { ensureAnonymousUser, ensureBuiltinRoles, ensureDefaultAdmin, + getAuthUser, } from "./plugins/auth.js"; import { registerMfa } from "./plugins/mfa.js"; import { oidcRoutes } from "./plugins/oidc.js"; @@ -57,6 +58,7 @@ import { settingsRoutes } from "./routes/settings.js"; import { teamsRoutes } from "./routes/teams.js"; import { registerToolRoutes } from "./routes/tools/index.js"; import { userFileRoutes } from "./routes/user-files.js"; +import { shutdownTracing } from "./tracing.js"; // Run before anything else try { @@ -189,26 +191,7 @@ function parseTrustProxy(value: string): boolean | number | string { const app = Fastify({ genReqId: (req) => (req.headers["x-request-id"] as string) ?? randomUUID(), - logger: { - level: env.LOG_LEVEL, - transport: { - targets: [ - { target: "pino/file", options: { destination: 1 } }, - { - // Rotate at 10 MB, keep 5 files - target: "pino-roll", - options: { - file: join(env.LOG_DIR, "snapotter"), - extension: ".log", - size: "10m", - limit: { count: 5 }, - mkdir: true, - }, - }, - ], - }, - redact: ["req.headers.authorization", "req.headers.cookie"], - }, + loggerInstance: logger, bodyLimit: env.MAX_UPLOAD_SIZE_MB > 0 ? env.MAX_UPLOAD_SIZE_MB * 1024 * 1024 : 1073741824, trustProxy: parseTrustProxy(env.TRUST_PROXY), routerOptions: { maxParamLength: 500 }, @@ -331,6 +314,18 @@ await authMiddleware(app); // Per-user rate limiting (after auth so request.user is populated) await registerPerUserRateLimit(app); +// Enrich active OTel span with tool_id and user_id when available +app.addHook("preHandler", (request, _reply, done) => { + const span = trace.getActiveSpan(); + if (span) { + const params = request.params as Record | undefined; + if (params?.toolId) span.setAttribute("snapotter.tool_id", params.toolId); + const user = getAuthUser(request); + if (user) span.setAttribute("snapotter.user_id", user.id); + } + done(); +}); + // Auth routes await authRoutes(app); @@ -672,6 +667,17 @@ async function shutdown(signal: string) { // Close BullMQ resources before database (workers first so no new jobs start) try { await closeWorkers(); + } catch (err) { + console.error("Error closing workers:", err); + } + + try { + await shutdownTracing(); + } catch { + // tracing shutdown is best-effort + } + + try { await closeFlowProducer(); await closeQueueEvents(); await closeQueues(); diff --git a/apps/api/src/jobs/enqueue.ts b/apps/api/src/jobs/enqueue.ts index bead0848..7d3261fe 100644 --- a/apps/api/src/jobs/enqueue.ts +++ b/apps/api/src/jobs/enqueue.ts @@ -5,6 +5,7 @@ * the appropriate BullMQ queue. waitForJob() blocks the HTTP request * until the worker produces a result or the sync-wait window expires. */ +import { context, propagation } from "@opentelemetry/api"; import { FlowProducer, type Job, QueueEvents } from "bullmq"; import { eq } from "drizzle-orm"; import { env } from "../config.js"; @@ -54,6 +55,24 @@ export async function closeFlowProducer(): Promise { } } +// ── Trace context injection ───────────────────────────────────── + +/** + * Inject the active OpenTelemetry trace context into a ToolJobData object. + * Called from enqueueToolJob (single jobs) and from pipeline/batch routes + * that build FlowProducer trees bypassing enqueueToolJob. + */ +export function injectTraceContext(data: ToolJobData): void { + const carrier: Record = {}; + propagation.inject(context.active(), carrier); + if (carrier.traceparent) { + data._otel = { + traceparent: carrier.traceparent, + tracestate: carrier.tracestate, + }; + } +} + // ── Enqueue + wait ────────────────────────────────────────────── /** @@ -82,6 +101,8 @@ export async function enqueueToolJob(data: ToolJobData): Promise {}); } + injectTraceContext(data); + const queue = getQueue(data.pool); const job = await queue.add(data.toolId, { ...data, jobId: data.jobId }, { jobId: data.jobId }); return job; diff --git a/apps/api/src/jobs/types.ts b/apps/api/src/jobs/types.ts index 1434af3f..d8689f75 100644 --- a/apps/api/src/jobs/types.ts +++ b/apps/api/src/jobs/types.ts @@ -46,6 +46,7 @@ export interface ToolJobData { parentId?: string; totalFiles?: number; fileIndex?: number; + _otel?: { traceparent: string; tracestate?: string }; } /** Result returned by a completed BullMQ job. */ diff --git a/apps/api/src/jobs/worker.ts b/apps/api/src/jobs/worker.ts index 3689a43a..e931fc66 100644 --- a/apps/api/src/jobs/worker.ts +++ b/apps/api/src/jobs/worker.ts @@ -24,12 +24,14 @@ import { mkdir, readFile, rm } from "node:fs/promises"; import { tmpdir } from "node:os"; import { join } from "node:path"; +import { context, propagation, ROOT_CONTEXT, SpanStatusCode, trace } from "@opentelemetry/api"; import { type Job, UnrecoverableError, Worker } from "bullmq"; import { eq } from "drizzle-orm"; import { env } from "../config.js"; import { db, schema } from "../db/index.js"; import { resolveConcurrency } from "../lib/env.js"; import { stripInternalPaths } from "../lib/errors.js"; +import { logger } from "../lib/logger.js"; import { jobDuration, jobsTotal } from "../lib/metrics.js"; import { getObjectBuffer, putObject } from "../lib/object-storage.js"; import { publishEphemeral, updateSingleFileProgress } from "../routes/progress.js"; @@ -104,251 +106,300 @@ async function processToolJob(job: Job): Promise { const { jobId } = data; const startTime = Date.now(); - // Register for cooperative cancellation - const ac = registerCancelable(jobId); - const signal = ac.signal; + // Extract OTel trace context if present (no-op without SDK) + const otel = data._otel; + const parentCtx = otel?.traceparent ? propagation.extract(ROOT_CONTEXT, otel) : ROOT_CONTEXT; + const tracer = trace.getTracer("snapotter-worker"); + const span = otel?.traceparent + ? tracer.startSpan( + "job.process", + { + attributes: { + "snapotter.job_id": jobId, + "snapotter.tool_id": data.toolId, + "snapotter.pool": data.pool, + "snapotter.attempt_number": job.attemptsMade + 1, + }, + }, + parentCtx, + ) + : null; - // Timeout guard (0 means unlimited; only arm when positive) - const timeoutMs = timeoutMsFor(data.pool); - const timeoutHandle = - timeoutMs > 0 ? setTimeout(() => ac.abort("timeout"), timeoutMs) : undefined; + const runBody = async (): Promise => { + if (span) span.addEvent("job.active"); - // Per-job scratch directory - const scratchDir = join(scratchRoot(), jobId); + // Register for cooperative cancellation + const ac = registerCancelable(jobId); + const signal = ac.signal; - try { - await mkdir(scratchDir, { recursive: true }); + // Timeout guard (0 means unlimited; only arm when positive) + const timeoutMs = timeoutMsFor(data.pool); + const timeoutHandle = + timeoutMs > 0 ? setTimeout(() => ac.abort("timeout"), timeoutMs) : undefined; - // Mark job as processing in the durable row - await db - .update(schema.jobs) - .set({ - status: "processing", - startedAt: new Date(), - attempts: job.attemptsMade + 1, - }) - .where(eq(schema.jobs.id, jobId)); + // Per-job scratch directory + const scratchDir = join(scratchRoot(), jobId); - // Load all input refs from object storage. The primary input keeps - // the client-facing filename; secondary inputs derive filenames from - // their ref basenames. - const inputs: ToolProcessInputV2[] = await Promise.all( - data.inputRefs.map(async (ref) => ({ - ref, - buffer: await getObjectBuffer(ref), - filename: ref.split("/").slice(2).join("/") || data.filename, - })), - ); - inputs[0].filename = data.filename; // primary keeps the client-facing name - const inputBuffer = inputs[0].buffer; // existing metrics/size/preview paths - - // Progress reporter: emits both Redis pub/sub and BullMQ job progress - const progressJobId = data.clientJobId ?? jobId; - const report = (percent: number, stage?: string) => { - updateSingleFileProgress({ - jobId: progressJobId, - phase: "processing", - percent, - stage, - }); - void job.updateProgress({ percent, stage }); - }; - - // Check for cancellation before dispatching - if (signal.aborted) throw new Error("Canceled"); - - // Build the process context - const ctx: ToolProcessCtx = { signal, scratchDir, report }; - - // Dispatch: AI handler or standard tool registry - let resultBuffer: Buffer; - let resultFilename: string; - let resultContentType: string; - let resultPayload: Record | undefined; - let extraOutputs: Array<{ name: string; buffer: Buffer; contentType: string }> | undefined; - - if (hasAiJobHandler(data.toolId)) { - const aiResult = await runAiToolJob(data, inputBuffer, ctx); - resultBuffer = aiResult.buffer; - resultFilename = aiResult.filename; - resultContentType = aiResult.contentType; - resultPayload = aiResult.resultPayload; - extraOutputs = aiResult.extraOutputs; - } else { - const config = getToolConfig(data.toolId); - if (!config) throw new Error(`No tool config for ${data.toolId}`); - - // Use the resolved v2 process function (adapter or native) - if (!config.processV2) throw new Error(`No processV2 for ${data.toolId}`); - const result = await config.processV2({ - inputs, - settings: data.settings, - scratchDir, - signal, - report, - }); - - // Resolve buffer OR scratchPath for the primary output - if (result.buffer) { - resultBuffer = result.buffer; - } else if (result.scratchPath) { - resultBuffer = await readFile(result.scratchPath); - } else { - throw new Error(`Tool ${data.toolId} returned neither buffer nor scratchPath`); - } - resultFilename = result.filename; - resultContentType = result.contentType; - resultPayload = result.resultPayload; - - // Resolve extra outputs with the same buffer/scratchPath duality - if (result.extraOutputs) { - extraOutputs = await Promise.all( - result.extraOutputs.map(async (extra) => { - let buf: Buffer; - if (extra.buffer) { - buf = extra.buffer; - } else if (extra.scratchPath) { - buf = await readFile(extra.scratchPath); - } else { - throw new Error(`Extra output "${extra.name}" has neither buffer nor scratchPath`); - } - return { name: extra.name, buffer: buf, contentType: extra.contentType }; - }), - ); - } - } - - // Build output name with tool suffix and extension fixup - const outName = buildOutputName(resultFilename, data.filename, data.toolId, resultContentType); - - // Write primary output to object storage - const primaryKey = `outputs/${jobId}/${outName}`; - await putObject(primaryKey, resultBuffer); - const outputRefs: string[] = [primaryKey]; - - // Write extra outputs (AI tools may produce multiple files) - if (extraOutputs) { - for (const extra of extraOutputs) { - const extraKey = `outputs/${jobId}/${extra.name}`; - await putObject(extraKey, extra.buffer); - outputRefs.push(extraKey); - } - } - - // Generate preview for non-browser-previewable formats - const previewRef = await generatePreview(resultBuffer, resultContentType, jobId, inputBuffer); - - // No auto-save -- users save to library explicitly via the UI - const savedFileId: string | undefined = undefined; - - const durationMs = Date.now() - startTime; - - // Build the result - const jobResult: ToolJobResult = { - outputRefs, - filename: outName, - contentType: resultContentType, - originalSize: inputBuffer.length, - processedSize: resultBuffer.length, - previewRef, - savedFileId, - resultPayload, - }; - - // Update durable row to completed - await db - .update(schema.jobs) - .set({ - status: "completed", - completedAt: new Date(), - durationMs, - bytesIn: inputBuffer.length, - bytesOut: resultBuffer.length, - outputRefs, - progress: { percent: 100, stage: "complete" }, - }) - .where(eq(schema.jobs.id, jobId)); - - // Record Prometheus metrics - jobsTotal.inc({ pool: data.pool, status: "completed" }); - jobDuration.observe({ pool: data.pool }, durationMs / 1000); - - // Emit terminal progress event with legacy result payload - const legacyResult = buildLegacyResultPayload(jobResult, jobId); - updateSingleFileProgress({ - jobId: progressJobId, - phase: "complete", - percent: 100, - stage: "complete", - result: legacyResult, - }); - - return jobResult; - } catch (err) { - const durationMs = Date.now() - startTime; - const isTimeout = signal.aborted && signal.reason === "timeout"; - const isCanceled = signal.aborted && !isTimeout; - const errorMessage = err instanceof Error ? err.message : String(err); - const finalError = isCanceled - ? "Canceled" - : isTimeout - ? `Timed out after ${Math.round(timeoutMs / 1000)}s` - : errorMessage; - - const maxAttempts = job.opts.attempts ?? 1; - const willRetry = !isCanceled && job.attemptsMade + 1 < maxAttempts; - - const progressJobId = data.clientJobId ?? jobId; - - // When the job will be retried, do NOT write a terminal DB row or - // emit a terminal SSE frame. The row stays "processing" and the - // next attempt overwrites startedAt/attempts as usual. - if (!willRetry) { - // Record Prometheus metrics on final attempt only - jobsTotal.inc({ pool: data.pool, status: isCanceled ? "canceled" : "failed" }); - jobDuration.observe({ pool: data.pool }, durationMs / 1000); + try { + await mkdir(scratchDir, { recursive: true }); + // Mark job as processing in the durable row await db .update(schema.jobs) .set({ - status: isCanceled ? "canceled" : "failed", - completedAt: new Date(), - durationMs, - error: { message: finalError }, + status: "processing", + startedAt: new Date(), + attempts: job.attemptsMade + 1, }) - .where(eq(schema.jobs.id, jobId)) - .catch(() => {}); + .where(eq(schema.jobs.id, jobId)); - if (isCanceled) { - // Ephemeral terminal event for live SSE clients. Uses - // publishEphemeral so the replay key is set without - // overwriting the DB row (which stays "canceled"). - publishEphemeral({ - jobId: progressJobId, - type: "single", - phase: "failed", - percent: 0, - error: "Canceled", - }); - } else { + // Load all input refs from object storage. The primary input keeps + // the client-facing filename; secondary inputs derive filenames from + // their ref basenames. + const inputs: ToolProcessInputV2[] = await Promise.all( + data.inputRefs.map(async (ref) => ({ + ref, + buffer: await getObjectBuffer(ref), + filename: ref.split("/").slice(2).join("/") || data.filename, + })), + ); + inputs[0].filename = data.filename; // primary keeps the client-facing name + const inputBuffer = inputs[0].buffer; // existing metrics/size/preview paths + + // Progress reporter: emits both Redis pub/sub and BullMQ job progress + const progressJobId = data.clientJobId ?? jobId; + const report = (percent: number, stage?: string) => { updateSingleFileProgress({ jobId: progressJobId, - phase: "failed", - percent: 0, - error: stripInternalPaths(finalError), + phase: "processing", + percent, + stage, }); - } - } + void job.updateProgress({ percent, stage }); + }; - if (isCanceled) throw new UnrecoverableError("Canceled"); - if (isTimeout) throw new Error(finalError); - throw err; - } finally { - clearTimeout(timeoutHandle); - unregisterCancelable(jobId); - // Clean up scratch directory - await rm(scratchDir, { recursive: true, force: true }).catch(() => {}); + // Check for cancellation before dispatching + if (signal.aborted) throw new Error("Canceled"); + + // Build the process context + const ctx: ToolProcessCtx = { signal, scratchDir, report }; + + // Dispatch: AI handler or standard tool registry + let resultBuffer: Buffer; + let resultFilename: string; + let resultContentType: string; + let resultPayload: Record | undefined; + let extraOutputs: Array<{ name: string; buffer: Buffer; contentType: string }> | undefined; + + if (hasAiJobHandler(data.toolId)) { + const aiResult = await runAiToolJob(data, inputBuffer, ctx); + resultBuffer = aiResult.buffer; + resultFilename = aiResult.filename; + resultContentType = aiResult.contentType; + resultPayload = aiResult.resultPayload; + extraOutputs = aiResult.extraOutputs; + } else { + const config = getToolConfig(data.toolId); + if (!config) throw new Error(`No tool config for ${data.toolId}`); + + // Use the resolved v2 process function (adapter or native) + if (!config.processV2) throw new Error(`No processV2 for ${data.toolId}`); + const result = await config.processV2({ + inputs, + settings: data.settings, + scratchDir, + signal, + report, + }); + + // Resolve buffer OR scratchPath for the primary output + if (result.buffer) { + resultBuffer = result.buffer; + } else if (result.scratchPath) { + resultBuffer = await readFile(result.scratchPath); + } else { + throw new Error(`Tool ${data.toolId} returned neither buffer nor scratchPath`); + } + resultFilename = result.filename; + resultContentType = result.contentType; + resultPayload = result.resultPayload; + + // Resolve extra outputs with the same buffer/scratchPath duality + if (result.extraOutputs) { + extraOutputs = await Promise.all( + result.extraOutputs.map(async (extra) => { + let buf: Buffer; + if (extra.buffer) { + buf = extra.buffer; + } else if (extra.scratchPath) { + buf = await readFile(extra.scratchPath); + } else { + throw new Error(`Extra output "${extra.name}" has neither buffer nor scratchPath`); + } + return { name: extra.name, buffer: buf, contentType: extra.contentType }; + }), + ); + } + } + + // Build output name with tool suffix and extension fixup + const outName = buildOutputName( + resultFilename, + data.filename, + data.toolId, + resultContentType, + ); + + // Write primary output to object storage + const primaryKey = `outputs/${jobId}/${outName}`; + await putObject(primaryKey, resultBuffer); + const outputRefs: string[] = [primaryKey]; + + // Write extra outputs (AI tools may produce multiple files) + if (extraOutputs) { + for (const extra of extraOutputs) { + const extraKey = `outputs/${jobId}/${extra.name}`; + await putObject(extraKey, extra.buffer); + outputRefs.push(extraKey); + } + } + + // Generate preview for non-browser-previewable formats + const previewRef = await generatePreview(resultBuffer, resultContentType, jobId, inputBuffer); + + // No auto-save -- users save to library explicitly via the UI + const savedFileId: string | undefined = undefined; + + const durationMs = Date.now() - startTime; + + // Build the result + const jobResult: ToolJobResult = { + outputRefs, + filename: outName, + contentType: resultContentType, + originalSize: inputBuffer.length, + processedSize: resultBuffer.length, + previewRef, + savedFileId, + resultPayload, + }; + + // Update durable row to completed + await db + .update(schema.jobs) + .set({ + status: "completed", + completedAt: new Date(), + durationMs, + bytesIn: inputBuffer.length, + bytesOut: resultBuffer.length, + outputRefs, + progress: { percent: 100, stage: "complete" }, + }) + .where(eq(schema.jobs.id, jobId)); + + // Record Prometheus metrics + jobsTotal.inc({ pool: data.pool, status: "completed" }); + jobDuration.observe({ pool: data.pool }, durationMs / 1000); + + // Emit terminal progress event with legacy result payload + const legacyResult = buildLegacyResultPayload(jobResult, jobId); + updateSingleFileProgress({ + jobId: progressJobId, + phase: "complete", + percent: 100, + stage: "complete", + result: legacyResult, + }); + + // Record queue wait time and completion on the OTel span + if (span && job.processedOn) { + span.setAttribute("snapotter.queue.wait_ms", job.processedOn - job.timestamp); + } + if (span) span.addEvent("job.completed"); + + return jobResult; + } catch (err) { + const durationMs = Date.now() - startTime; + const isTimeout = signal.aborted && signal.reason === "timeout"; + const isCanceled = signal.aborted && !isTimeout; + const errorMessage = err instanceof Error ? err.message : String(err); + const finalError = isCanceled + ? "Canceled" + : isTimeout + ? `Timed out after ${Math.round(timeoutMs / 1000)}s` + : errorMessage; + + // Record error on the OTel span + if (span) { + span.setStatus({ code: SpanStatusCode.ERROR, message: finalError }); + span.recordException(err instanceof Error ? err : new Error(String(err))); + span.addEvent("job.failed"); + } + + const maxAttempts = job.opts.attempts ?? 1; + const willRetry = !isCanceled && job.attemptsMade + 1 < maxAttempts; + + const progressJobId = data.clientJobId ?? jobId; + + // When the job will be retried, do NOT write a terminal DB row or + // emit a terminal SSE frame. The row stays "processing" and the + // next attempt overwrites startedAt/attempts as usual. + if (!willRetry) { + // Record Prometheus metrics on final attempt only + jobsTotal.inc({ pool: data.pool, status: isCanceled ? "canceled" : "failed" }); + jobDuration.observe({ pool: data.pool }, durationMs / 1000); + + await db + .update(schema.jobs) + .set({ + status: isCanceled ? "canceled" : "failed", + completedAt: new Date(), + durationMs, + error: { message: finalError }, + }) + .where(eq(schema.jobs.id, jobId)) + .catch(() => {}); + + if (isCanceled) { + // Ephemeral terminal event for live SSE clients. Uses + // publishEphemeral so the replay key is set without + // overwriting the DB row (which stays "canceled"). + publishEphemeral({ + jobId: progressJobId, + type: "single", + phase: "failed", + percent: 0, + error: "Canceled", + }); + } else { + updateSingleFileProgress({ + jobId: progressJobId, + phase: "failed", + percent: 0, + error: stripInternalPaths(finalError), + }); + } + } + + if (isCanceled) throw new UnrecoverableError("Canceled"); + if (isTimeout) throw new Error(finalError); + throw err; + } finally { + if (span) span.end(); + clearTimeout(timeoutHandle); + unregisterCancelable(jobId); + // Clean up scratch directory + await rm(scratchDir, { recursive: true, force: true }).catch(() => {}); + } + }; + + // Execute with or without active span context + if (span) { + const activeCtx = trace.setSpan(parentCtx, span); + return context.with(activeCtx, runBody); } + return runBody(); } // ── Pipeline step handler ───────────────────────────────────── @@ -687,7 +738,7 @@ export function startWorkers(): void { }); worker.on("error", (err) => { - console.error(`Worker error [${pool}]:`, err); + logger.error({ err, pool }, "Worker error"); }); workers.push(worker); @@ -715,7 +766,7 @@ export function startWorkers(): void { workers.push(worker); } - console.log( + logger.info( `Workers started: ${POOLS.map((p) => `${p}(${p === "system" || p === "ai" ? 1 : concurrency})`).join(", ")}`, ); } diff --git a/apps/api/src/lib/log-trace-mixin.ts b/apps/api/src/lib/log-trace-mixin.ts new file mode 100644 index 00000000..2cba8963 --- /dev/null +++ b/apps/api/src/lib/log-trace-mixin.ts @@ -0,0 +1,8 @@ +import { trace } from "@opentelemetry/api"; + +export function traceMixin(): Record { + const span = trace.getActiveSpan(); + if (!span) return {}; + const { traceId, spanId, traceFlags } = span.spanContext(); + return { traceId, spanId, traceFlags }; +} diff --git a/apps/api/src/lib/logger.ts b/apps/api/src/lib/logger.ts new file mode 100644 index 00000000..9892177d --- /dev/null +++ b/apps/api/src/lib/logger.ts @@ -0,0 +1,26 @@ +import { join } from "node:path"; +import type { FastifyBaseLogger } from "fastify"; +import pino from "pino"; +import { env } from "../config.js"; +import { traceMixin } from "./log-trace-mixin.js"; + +export const logger: FastifyBaseLogger = pino({ + level: env.LOG_LEVEL, + mixin: traceMixin, + transport: { + targets: [ + { target: "pino/file", options: { destination: 1 } }, + { + target: "pino-roll", + options: { + file: join(env.LOG_DIR, "snapotter"), + extension: ".log", + size: "10m", + limit: { count: 5 }, + mkdir: true, + }, + }, + ], + }, + redact: ["req.headers.authorization", "req.headers.cookie"], +}); diff --git a/apps/api/src/routes/batch.ts b/apps/api/src/routes/batch.ts index 18b2e5a0..dc95ed5f 100644 --- a/apps/api/src/routes/batch.ts +++ b/apps/api/src/routes/batch.ts @@ -18,7 +18,7 @@ import sharp from "sharp"; import { env } from "../config.js"; import { db, schema } from "../db/index.js"; import { recordChildOutcome } from "../jobs/batch-progress.js"; -import { getFlowProducer, waitForJob } from "../jobs/enqueue.js"; +import { getFlowProducer, injectTraceContext, waitForJob } from "../jobs/enqueue.js"; import { type Pool, queueName, type ToolJobData } from "../jobs/types.js"; import { autoOrient } from "../lib/auto-orient.js"; import { getSecurityHeaders } from "../lib/csp.js"; @@ -39,6 +39,16 @@ interface ParsedFile { filename: string; } +/** Recursively inject OTel trace context into every node of a FlowJob tree. */ +function injectTraceContextIntoFlow(node: FlowJob): void { + injectTraceContext(node.data as ToolJobData); + if (node.children) { + for (const child of node.children) { + injectTraceContextIntoFlow(child); + } + } +} + export async function registerBatchRoutes(app: FastifyInstance): Promise { app.post( "/api/v1/tools/:toolId/batch", @@ -297,6 +307,9 @@ export async function registerBatchRoutes(app: FastifyInstance): Promise { .set({ settings: { flowChildCount: flowChildren.length } }) .where(eq(schema.jobs.id, parentId)); + // Inject OTel trace context into every node of the batch flow tree + injectTraceContextIntoFlow(batchTree); + await getFlowProducer().add(batchTree); // ── Wait for completion and stream ZIP ───────────────────────── diff --git a/apps/api/src/routes/pipeline.ts b/apps/api/src/routes/pipeline.ts index 1e222778..061e662f 100644 --- a/apps/api/src/routes/pipeline.ts +++ b/apps/api/src/routes/pipeline.ts @@ -17,7 +17,7 @@ import { z } from "zod"; import { env } from "../config.js"; import { db, schema } from "../db/index.js"; import { recordChildOutcome } from "../jobs/batch-progress.js"; -import { getFlowProducer, waitForJob } from "../jobs/enqueue.js"; +import { getFlowProducer, injectTraceContext, waitForJob } from "../jobs/enqueue.js"; import { type Pool, queueName, type ToolJobData } from "../jobs/types.js"; import { trackEvent } from "../lib/analytics.js"; import { autoOrient } from "../lib/auto-orient.js"; @@ -78,6 +78,16 @@ interface ParsedStep { */ const PASSWORD_TOOLS = new Set(["protect-pdf", "unlock-pdf"]); +/** Recursively inject OTel trace context into every node of a FlowJob tree. */ +function injectTraceContextIntoFlow(node: FlowJob): void { + injectTraceContext(node.data as ToolJobData); + if (node.children) { + for (const child of node.children) { + injectTraceContextIntoFlow(child); + } + } +} + /** * Build a FlowJob tree for a single-file pipeline. * @@ -405,6 +415,9 @@ export async function registerPipelineRoutes(app: FastifyInstance): Promise } | null = null; + +export function isTracingActive(): boolean { + return _active; +} + +export async function initTracing(options: { exporter?: SpanExporter } = {}): Promise { + if (_active) return; + + const { NodeSDK } = await import("@opentelemetry/sdk-node"); + const { BatchSpanProcessor } = await import("@opentelemetry/sdk-trace-base"); + const { OTLPTraceExporter } = await import("@opentelemetry/exporter-trace-otlp-http"); + const { defaultResource, resourceFromAttributes } = await import("@opentelemetry/resources"); + const { ATTR_SERVICE_NAME, ATTR_SERVICE_VERSION } = await import( + "@opentelemetry/semantic-conventions" + ); + const { HttpInstrumentation } = await import("@opentelemetry/instrumentation-http"); + const { FastifyInstrumentation } = await import("@opentelemetry/instrumentation-fastify"); + const { PgInstrumentation } = await import("@opentelemetry/instrumentation-pg"); + const { IORedisInstrumentation } = await import("@opentelemetry/instrumentation-ioredis"); + const { AwsInstrumentation } = await import("@opentelemetry/instrumentation-aws-sdk"); + + const require = createRequire(import.meta.url); + const { version } = require("../package.json"); + + const exporter = options.exporter ?? new OTLPTraceExporter(); + + const resource = defaultResource().merge( + resourceFromAttributes({ + [ATTR_SERVICE_NAME]: "snapotter-api", + [ATTR_SERVICE_VERSION]: version, + }), + ); + + const sdk = new NodeSDK({ + resource, + spanProcessors: [new BatchSpanProcessor(exporter)], + instrumentations: [ + new HttpInstrumentation({ + ignoreOutgoingRequestHook: (req) => { + const host = req.hostname || req.host || ""; + return host.includes("posthog") || host.includes("sentry"); + }, + }), + new FastifyInstrumentation(), + new PgInstrumentation(), + new IORedisInstrumentation({ + requireParentSpan: true, + dbStatementSerializer: (cmd, args) => { + if (cmd === "evalsha" || cmd === "eval") return `${cmd} `; + return `${cmd} ${(args ?? []).slice(0, 2).join(" ")}`; + }, + }), + new AwsInstrumentation(), + ], + }); + + sdk.start(); + _sdk = sdk; + _active = true; +} + +export async function shutdownTracing(): Promise { + if (_sdk) { + await _sdk.shutdown().catch(() => {}); + _sdk = null; + _active = false; + } +} + +// -- Preload entry point -- +// When loaded via --import, this top-level await runs before the app. + +const endpoint = process.env.OTEL_EXPORTER_OTLP_ENDPOINT; +if (endpoint) { + try { + const enterprise = await import("@snapotter/enterprise"); + const licenseKey = process.env.LICENSE_KEY ?? ""; + if (licenseKey) { + enterprise.initEnterprise(licenseKey); + } + if (enterprise.isFeatureEnabled("distributed_tracing")) { + await initTracing(); + console.log("[tracing] OpenTelemetry initialized, exporting to", endpoint); + } + } catch { + // Enterprise package not available or license invalid -- tracing stays off + } +} diff --git a/packages/ai/package.json b/packages/ai/package.json index 534b94c0..d8f1c75d 100644 --- a/packages/ai/package.json +++ b/packages/ai/package.json @@ -10,6 +10,7 @@ "clean": "rm -rf dist" }, "dependencies": { + "@opentelemetry/api": "^1.9.1", "@snapotter/shared": "workspace:*", "sharp": "^0.34.5" }, diff --git a/packages/ai/python/dispatcher.py b/packages/ai/python/dispatcher.py index 2f62a52b..7856e4c4 100644 --- a/packages/ai/python/dispatcher.py +++ b/packages/ai/python/dispatcher.py @@ -28,6 +28,32 @@ import os import traceback +# ── Optional OpenTelemetry tracing (enterprise only) ───────────── +_tracer = None +_tracer_provider = None + +def _init_tracing(): + """Initialize OTel tracing if OTEL_EXPORTER_OTLP_ENDPOINT is set.""" + global _tracer, _tracer_provider + endpoint = os.environ.get("OTEL_EXPORTER_OTLP_ENDPOINT") + if not endpoint: + return + try: + from opentelemetry import trace as otel_trace + from opentelemetry.sdk.trace import TracerProvider + from opentelemetry.sdk.trace.export import BatchSpanProcessor + from opentelemetry.exporter.otlp.proto.http.trace_exporter import OTLPSpanExporter + from opentelemetry.sdk.resources import Resource + + resource = Resource.create({"service.name": "snapotter-sidecar"}) + _tracer_provider = TracerProvider(resource=resource) + _tracer_provider.add_span_processor(BatchSpanProcessor(OTLPSpanExporter())) + otel_trace.set_tracer_provider(_tracer_provider) + _tracer = otel_trace.get_tracer("snapotter-sidecar") + except ImportError: + pass + + # ── Script allowlist ─────────────────────────────────────────────────── # Only these script names (without .py) may be dispatched. This is the # primary security gate -- no path traversal, no arbitrary file execution. @@ -319,6 +345,8 @@ def main(): print(json.dumps({"ready": True, "gpu": gpu}), file=sys.stderr, flush=True) print(f"[dispatcher] Ready. GPU: {gpu}. Max requests: {MAX_REQUESTS}. Modules: {list(available_modules.keys())}", file=sys.stderr, flush=True) + _init_tracing() + request_count = 0 for line in sys.stdin: @@ -335,8 +363,38 @@ def main(): script_name = request.get("script", "") args = request.get("args", []) + otel_data = request.pop("_otel", None) + otel_ctx = None + if otel_data and _tracer: + from opentelemetry.trace.propagation.tracecontext import TraceContextTextMapPropagator + from opentelemetry import context as otel_context + propagator = TraceContextTextMapPropagator() + otel_ctx = propagator.extract(carrier=otel_data) + try: - stdout_output, exit_code = _run_script_main(script_name, args) + if otel_ctx and _tracer: + from opentelemetry import trace as otel_trace + from opentelemetry.trace import StatusCode + from opentelemetry import context as otel_context + token = otel_context.attach(otel_ctx) + span = _tracer.start_span(f"sidecar.{script_name}", context=otel_ctx) + try: + stdout_output, exit_code = _run_script_main(script_name, args) + if exit_code != 0: + span.set_status(StatusCode.ERROR, f"exit code {exit_code}") + except Exception as exc: + span.set_status(StatusCode.ERROR, str(exc)) + span.record_exception(exc) + raise + finally: + span.end() + otel_context.detach(token) + try: + _tracer_provider.force_flush() + except Exception: + pass + else: + stdout_output, exit_code = _run_script_main(script_name, args) response = { "id": request_id, "stdout": stdout_output, @@ -361,6 +419,12 @@ def main(): file=sys.stderr, flush=True) break + if _tracer_provider: + try: + _tracer_provider.shutdown() + except Exception: + pass + if __name__ == "__main__": main() diff --git a/packages/ai/python/tests/__init__.py b/packages/ai/python/tests/__init__.py new file mode 100644 index 00000000..e69de29b diff --git a/packages/ai/python/tests/test_tracing.py b/packages/ai/python/tests/test_tracing.py new file mode 100644 index 00000000..3d7e6b9e --- /dev/null +++ b/packages/ai/python/tests/test_tracing.py @@ -0,0 +1,43 @@ +import json +import os +import pytest + + +def test_otel_stripped_from_request(): + """_otel should be popped from the request before reaching scripts.""" + request = { + "id": "test-1", + "script": "remove_bg", + "args": ["input.png"], + "_otel": {"traceparent": "00-abc123-def456-01"}, + } + otel_data = request.pop("_otel", None) + assert otel_data is not None + assert otel_data["traceparent"] == "00-abc123-def456-01" + assert "_otel" not in request + assert request["args"] == ["input.png"] + + +def test_no_otel_in_request(): + """Requests without _otel should work normally.""" + request = { + "id": "test-2", + "script": "remove_bg", + "args": ["input.png"], + } + otel_data = request.pop("_otel", None) + assert otel_data is None + assert request["args"] == ["input.png"] + + +def test_tracing_init_without_endpoint(monkeypatch): + """_init_tracing should be a no-op without OTEL_EXPORTER_OTLP_ENDPOINT.""" + monkeypatch.delenv("OTEL_EXPORTER_OTLP_ENDPOINT", raising=False) + import sys + dispatcher_dir = os.path.join(os.path.dirname(__file__), "..") + if dispatcher_dir not in sys.path: + sys.path.insert(0, dispatcher_dir) + from dispatcher import _init_tracing, _tracer + _init_tracing() + from dispatcher import _tracer as tracer_after + assert tracer_after is None diff --git a/packages/ai/src/bridge.ts b/packages/ai/src/bridge.ts index 27704983..b111ec1b 100644 --- a/packages/ai/src/bridge.ts +++ b/packages/ai/src/bridge.ts @@ -2,6 +2,7 @@ import { type ChildProcess, spawn } from "node:child_process"; import { randomUUID } from "node:crypto"; import { dirname, resolve } from "node:path"; import { fileURLToPath } from "node:url"; +import { context, propagation, SpanStatusCode, trace } from "@opentelemetry/api"; const __dirname = dirname(fileURLToPath(import.meta.url)); const PYTHON_DIR = resolve(__dirname, "../python"); @@ -32,6 +33,9 @@ function buildMinimalEnv(): Record { "DISPATCHER_MAX_REQUESTS", "PYTHON_VENV_PATH", "SNAPOTTER_GPU", + "OTEL_EXPORTER_OTLP_ENDPOINT", + "OTEL_EXPORTER_OTLP_PROTOCOL", + "OTEL_EXPORTER_OTLP_HEADERS", ]; for (const key of passthrough) { if (process.env[key] !== undefined) { @@ -364,7 +368,16 @@ export class PythonDispatcher { stderrLines: [], }); - const request = JSON.stringify({ id, script: scriptName.replace(".py", ""), args }); + const msg: Record = { id, script: scriptName.replace(".py", ""), args }; + const otelCarrier: Record = {}; + propagation.inject(context.active(), otelCarrier); + if (otelCarrier.traceparent) { + msg._otel = { + traceparent: otelCarrier.traceparent, + tracestate: otelCarrier.tracestate, + }; + } + const request = JSON.stringify(msg); try { proc.stdin!.write(request + "\n"); } catch { @@ -567,28 +580,58 @@ export class PythonDispatcher { timeout?: number; } = {}, ): Promise<{ stdout: string; stderr: string }> { - // Try persistent dispatcher first - const dispatcherPromise = this.dispatcherRun(scriptName, args, options); - if (dispatcherPromise) { - return dispatcherPromise.catch((err: Error) => { - if ( - err.message === "Python dispatcher exited unexpectedly" || - err.message === "Python dispatcher stdin closed unexpectedly" - ) { - console.warn( - `[bridge] Dispatcher crashed during ${scriptName}, retrying with per-request process`, - ); - return this.runPerRequest(scriptName, args, options).then((result) => ({ - ...result, - stderr: `${result.stderr}\n[bridge] retried after dispatcher crash`, - })); - } - throw err; - }); - } + const tracer = trace.getTracer("snapotter-sidecar"); + const span = trace.getActiveSpan() + ? tracer.startSpan("sidecar.execute", { + attributes: { + "sidecar.script": scriptName.replace(".py", ""), + "sidecar.profile": this.profile, + }, + }) + : null; - // Fall back to per-request spawning - return this.runPerRequest(scriptName, args, options); + const doRun = (): Promise<{ stdout: string; stderr: string }> => { + // Try persistent dispatcher first + const dispatcherPromise = this.dispatcherRun(scriptName, args, options); + if (dispatcherPromise) { + return dispatcherPromise.catch((err: Error) => { + if ( + err.message === "Python dispatcher exited unexpectedly" || + err.message === "Python dispatcher stdin closed unexpectedly" + ) { + console.warn( + `[bridge] Dispatcher crashed during ${scriptName}, retrying with per-request process`, + ); + return this.runPerRequest(scriptName, args, options).then((result) => ({ + ...result, + stderr: `${result.stderr}\n[bridge] retried after dispatcher crash`, + })); + } + throw err; + }); + } + + // Fall back to per-request spawning + return this.runPerRequest(scriptName, args, options); + }; + + if (!span) return doRun(); + + return doRun().then( + (result) => { + span.end(); + return result; + }, + (err) => { + span.setStatus({ + code: SpanStatusCode.ERROR, + message: err instanceof Error ? err.message : String(err), + }); + span.recordException(err instanceof Error ? err : new Error(String(err))); + span.end(); + throw err; + }, + ); } } diff --git a/packages/enterprise/src/license.ts b/packages/enterprise/src/license.ts index 343a2d4c..1d52e2ad 100644 --- a/packages/enterprise/src/license.ts +++ b/packages/enterprise/src/license.ts @@ -23,6 +23,7 @@ export const ENTERPRISE_FEATURES = [ "config_export_import", "upgrade_management", "admin_alerts", + "distributed_tracing", ] as const; export type EnterpriseFeature = (typeof ENTERPRISE_FEATURES)[number]; diff --git a/pnpm-lock.yaml b/pnpm-lock.yaml index 5e8e59e7..01b9855a 100644 --- a/pnpm-lock.yaml +++ b/pnpm-lock.yaml @@ -140,12 +140,45 @@ importers: '@node-saml/node-saml': specifier: ^5.1.0 version: 5.1.0 + '@opentelemetry/api': + specifier: ^1.9.1 + version: 1.9.1 + '@opentelemetry/exporter-trace-otlp-http': + specifier: ^0.219.0 + version: 0.219.0(@opentelemetry/api@1.9.1) + '@opentelemetry/instrumentation-aws-sdk': + specifier: ^0.74.0 + version: 0.74.0(@opentelemetry/api@1.9.1) + '@opentelemetry/instrumentation-fastify': + specifier: ^0.57.0 + version: 0.57.0(@opentelemetry/api@1.9.1) + '@opentelemetry/instrumentation-http': + specifier: ^0.219.0 + version: 0.219.0(@opentelemetry/api@1.9.1) + '@opentelemetry/instrumentation-ioredis': + specifier: ^0.67.0 + version: 0.67.0(@opentelemetry/api@1.9.1) + '@opentelemetry/instrumentation-pg': + specifier: ^0.71.0 + version: 0.71.0(@opentelemetry/api@1.9.1) + '@opentelemetry/resources': + specifier: ^2.8.0 + version: 2.8.0(@opentelemetry/api@1.9.1) + '@opentelemetry/sdk-node': + specifier: ^0.219.0 + version: 0.219.0(@opentelemetry/api@1.9.1) + '@opentelemetry/sdk-trace-base': + specifier: ^2.8.0 + version: 2.8.0(@opentelemetry/api@1.9.1) + '@opentelemetry/semantic-conventions': + specifier: ^1.41.1 + version: 1.41.1 '@scalar/fastify-api-reference': specifier: ^1.57.5 version: 1.58.0 '@sentry/node': specifier: ^10.55.0 - version: 10.56.0 + version: 10.56.0(@opentelemetry/exporter-trace-otlp-http@0.219.0(@opentelemetry/api@1.9.1)) '@snapotter/ai': specifier: workspace:* version: link:../../packages/ai @@ -227,6 +260,9 @@ importers: pg: specifier: ^8.21.0 version: 8.21.0 + pino: + specifier: ^10.3.1 + version: 10.3.1 pino-roll: specifier: ^4.0.0 version: 4.0.0 @@ -498,6 +534,9 @@ importers: packages/ai: dependencies: + '@opentelemetry/api': + specifier: ^1.9.1 + version: 1.9.1 '@snapotter/shared': specifier: workspace:* version: link:../shared @@ -2866,42 +2905,241 @@ packages: '@octokit/types@16.0.0': resolution: {integrity: sha512-sKq+9r1Mm4efXW1FCk7hFSeJo4QKreL/tTbR0rz/qx/r1Oa2VV83LTA/H/MuCOX7uCIJmQVRKBcbmWoySjAnSg==} + '@opentelemetry/api-logs@0.213.0': + resolution: {integrity: sha512-zRM5/Qj6G84Ej3F1yt33xBVY/3tnMxtL1fiDIxYbDWYaZ/eudVw3/PBiZ8G7JwUxXxjW8gU4g6LnOyfGKYHYgw==} + engines: {node: '>=8.0.0'} + '@opentelemetry/api-logs@0.214.0': resolution: {integrity: sha512-40lSJeqYO8Uz2Yj7u94/SJWE/wONa7rmMKjI1ZcIjgf3MHNHv1OZUCrCETGuaRF62d5pQD1wKIW+L4lmSMTzZA==} engines: {node: '>=8.0.0'} + '@opentelemetry/api-logs@0.219.0': + resolution: {integrity: sha512-FFx7YnaYJlIjqWW/AG/yAZ0L/NEY724PipXXXQLdtZPbLwBGbUMTGL1i/esI56TWfTUXxhLfpgrnWJCG8aUJyg==} + engines: {node: '>=8.0.0'} + '@opentelemetry/api@1.9.1': resolution: {integrity: sha512-gLyJlPHPZYdAk1JENA9LeHejZe1Ti77/pTeFm/nMXmQH/HFZlcS/O2XJB+L8fkbrNSqhdtlvjBVjxwUYanNH5Q==} engines: {node: '>=8.0.0'} + '@opentelemetry/configuration@0.219.0': + resolution: {integrity: sha512-wXZUYv4ngu43nA4WEhuXNacm46LW+17LRM8nKyIhBzroRA24PBYjMnakwzR/w777nFUB5xlgsYTTeuXxumZM1Q==} + engines: {node: ^18.19.0 || >=20.6.0} + peerDependencies: + '@opentelemetry/api': ^1.9.0 + + '@opentelemetry/context-async-hooks@2.8.0': + resolution: {integrity: sha512-/3FIraneMcng67SUJCxvyInk/oxzwsxyadufk0wwfOBLf5wqtAGX4MoQASwSbndBPeARzBryUM9Azr5kHIdWLw==} + engines: {node: ^18.19.0 || >=20.6.0} + peerDependencies: + '@opentelemetry/api': '>=1.0.0 <1.10.0' + '@opentelemetry/core@2.7.0': resolution: {integrity: sha512-DT12SXVwV2eoJrGf4nnsvZojxxeQo+LlNAsoYGRRObPWTeN6APiqZ2+nqDCQDvQX40eLi1AePONS0onoASp3yQ==} engines: {node: ^18.19.0 || >=20.6.0} peerDependencies: '@opentelemetry/api': '>=1.0.0 <1.10.0' + '@opentelemetry/core@2.8.0': + resolution: {integrity: sha512-hd1Lfh8p545nNz+jq1Ejfz+Mn1hyLuxYn1YzTfFNrxr8urEWMNQLPf1Th8kjOH+HxwawCrtgBp8JpBUR4ZSgww==} + engines: {node: ^18.19.0 || >=20.6.0} + peerDependencies: + '@opentelemetry/api': '>=1.0.0 <1.10.0' + + '@opentelemetry/exporter-logs-otlp-grpc@0.219.0': + resolution: {integrity: sha512-7SvzDCIclHWAcCwZ1MTOLcwn4BVNPGI3QxS/DJraPNe1TTL+4TvUBq5zeQV8tsnYvtDN7wKW2qocVmaCP2l7sQ==} + engines: {node: ^18.19.0 || >=20.6.0} + peerDependencies: + '@opentelemetry/api': ^1.3.0 + + '@opentelemetry/exporter-logs-otlp-http@0.219.0': + resolution: {integrity: sha512-mhl2HL6GmZI8b8PwPfqMws/5ovJfbRTxwc9Y5agVVHiQ+e5SL1btsFr/kJDgt7YCexDtsUn5HAreHQO9szFS0A==} + engines: {node: ^18.19.0 || >=20.6.0} + peerDependencies: + '@opentelemetry/api': ^1.3.0 + + '@opentelemetry/exporter-logs-otlp-proto@0.219.0': + resolution: {integrity: sha512-Ayw4Gf71PS9jhBVaYywa4WsajnqfDehMkTdVH3TSAVHqPcsAv/AhH/wTNRYNt99szeYr6Gbd/D6RjZD77wAxHg==} + engines: {node: ^18.19.0 || >=20.6.0} + peerDependencies: + '@opentelemetry/api': ^1.3.0 + + '@opentelemetry/exporter-metrics-otlp-grpc@0.219.0': + resolution: {integrity: sha512-6LaaSrPxK5L55bXevWajvOMxGOpNm0n12tG53TeZaUeNzXwLPg6d2KCC1zAlGsojan+xRG71mA4Qqs9K2VVrKQ==} + engines: {node: ^18.19.0 || >=20.6.0} + peerDependencies: + '@opentelemetry/api': ^1.3.0 + + '@opentelemetry/exporter-metrics-otlp-http@0.219.0': + resolution: {integrity: sha512-6CaDRbMVHZSDWzNXwrR8y/H4B/Z1eMNnkHiPQlTx3Ojz2OHY4X/aff/UC4P/3pHUQSuTfi3oh2UsPPZppw+Vrg==} + engines: {node: ^18.19.0 || >=20.6.0} + peerDependencies: + '@opentelemetry/api': ^1.3.0 + + '@opentelemetry/exporter-metrics-otlp-proto@0.219.0': + resolution: {integrity: sha512-DUS7XyIiEnoeccQUvuKy0G2/YqeKhpN8FVIrGbrLNIVMj10yeIFLRzRv0tibCI2kXXvlTTABVexGAk78wHk2ug==} + engines: {node: ^18.19.0 || >=20.6.0} + peerDependencies: + '@opentelemetry/api': ^1.3.0 + + '@opentelemetry/exporter-prometheus@0.219.0': + resolution: {integrity: sha512-TxOnJ85eWJY5JyOJsNMXiRTYlkDcOv0u3KbXEzWCc+tUS9sjL/BC6BcdxZ0B9r2OFVqsrZFXUzSD2sZUy42Ucw==} + engines: {node: ^18.19.0 || >=20.6.0} + peerDependencies: + '@opentelemetry/api': ^1.3.0 + + '@opentelemetry/exporter-trace-otlp-grpc@0.219.0': + resolution: {integrity: sha512-BkDNv1UD6BscW19MxbAxVmSYSSFuyeqR6buV2/HTYqA7GrR0EbTFzqG6h86T3PtXmpdbsWjMGLDdjG2rikG27Q==} + engines: {node: ^18.19.0 || >=20.6.0} + peerDependencies: + '@opentelemetry/api': ^1.3.0 + + '@opentelemetry/exporter-trace-otlp-http@0.219.0': + resolution: {integrity: sha512-9t6SvBXXBEjOBcIzgozvBbd3jWrv3Gt3ngGhl1fhdZ/zRc7oZDVOFEqbi2zlBpW9BXhgDMKv422J0DL/3iQWfw==} + engines: {node: ^18.19.0 || >=20.6.0} + peerDependencies: + '@opentelemetry/api': ^1.3.0 + + '@opentelemetry/exporter-trace-otlp-proto@0.219.0': + resolution: {integrity: sha512-lF/LUBfhOFmxJa+SQsLN7ziV4MHa2pyKgOM6JNehSOfU+npjM4gwm9oIKEJrzrWcexMcqydiyoFy0XCb1Ql3wQ==} + engines: {node: ^18.19.0 || >=20.6.0} + peerDependencies: + '@opentelemetry/api': ^1.3.0 + + '@opentelemetry/exporter-zipkin@2.8.0': + resolution: {integrity: sha512-Mj84UkEa17BK2o903VTXW3wM8CrSZexGs4tRGVZVIMM9ni1T6TuGx5IrRfoWKAbshx42D5/kc7YV+axypLPYyA==} + engines: {node: ^18.19.0 || >=20.6.0} + peerDependencies: + '@opentelemetry/api': ^1.0.0 + + '@opentelemetry/instrumentation-aws-sdk@0.74.0': + resolution: {integrity: sha512-EMLUGgx2wJSXdwMEFdwd3IaW+mkUF8PENdzDFQ1FRdztzMG1d1XN76ORIjAMsXcKQISlRRcz93AWPQeBPn4EKA==} + engines: {node: ^18.19.0 || >=20.6.0} + peerDependencies: + '@opentelemetry/api': ^1.3.0 + + '@opentelemetry/instrumentation-fastify@0.57.0': + resolution: {integrity: sha512-D+rwRtbiOediYocpKGvY/RQTpuLsLdCVwaOREyqWViwItJGibWI7O/wgd9xIV63pMP0D9IdSy27wnARfUaotKg==} + engines: {node: ^18.19.0 || >=20.6.0} + deprecated: Deprecated in favor of @fastify/otel, maintained by the Fastify authors. + peerDependencies: + '@opentelemetry/api': ^1.3.0 + + '@opentelemetry/instrumentation-http@0.219.0': + resolution: {integrity: sha512-nNt1fqpyah/OKjNHdEOu8xLwISppRU2qJuF8aR+fCcftVwdFkPgtworBLA+TI1HU2iF508jcQBF2gerWczJAXg==} + engines: {node: ^18.19.0 || >=20.6.0} + peerDependencies: + '@opentelemetry/api': ^1.3.0 + + '@opentelemetry/instrumentation-ioredis@0.67.0': + resolution: {integrity: sha512-dv64vQ4aXbJvRMMAFrMUSzDeJrNv/uQMLjfaav4LHAOar7Xn08W3pkoYoYEnzq/n2+fGgG96rp9S0gQ+VnLDiw==} + engines: {node: ^18.19.0 || >=20.6.0} + peerDependencies: + '@opentelemetry/api': ^1.3.0 + + '@opentelemetry/instrumentation-pg@0.71.0': + resolution: {integrity: sha512-jAhfyZeOkEKh3cQ5nm1tNWqHg7HFARyAe+p4BSoDHnB79c1woyEvDKqS11Hj/DjtceP+vrurIfcDs7Fqiy11mQ==} + engines: {node: ^18.19.0 || >=20.6.0} + peerDependencies: + '@opentelemetry/api': ^1.3.0 + + '@opentelemetry/instrumentation@0.213.0': + resolution: {integrity: sha512-3i9NdkET/KvQomeh7UaR/F4r9P25Rx6ooALlWXPIjypcEOUxksCmVu0zA70NBJWlrMW1rPr/LRidFAflLI+s/w==} + engines: {node: ^18.19.0 || >=20.6.0} + peerDependencies: + '@opentelemetry/api': ^1.3.0 + '@opentelemetry/instrumentation@0.214.0': resolution: {integrity: sha512-MHqEX5Dk59cqVah5LiARMACku7jXSVk9iVDWOea4x3cr7VfdByeDCURK6o1lntT1JS/Tsovw01UJrBhN3/uC5w==} engines: {node: ^18.19.0 || >=20.6.0} peerDependencies: '@opentelemetry/api': ^1.3.0 - '@opentelemetry/resources@2.7.0': - resolution: {integrity: sha512-K+oi0hNMv94EpZbnW3eyu2X6SGVpD3O5DhG2NIp65Hc7lhAj9brRXTAVzh3wB82+q3ThakEf7Zd7RsFUqcTc7A==} + '@opentelemetry/instrumentation@0.219.0': + resolution: {integrity: sha512-X5t7I8GyIO9rmGHwoedZLREpQqrF1WW2nxzNNym6HOKpFiE+rvqV3ngC0xcZVO2YwIGf3KKmRdWrYwdwz3H9RQ==} + engines: {node: ^18.19.0 || >=20.6.0} + peerDependencies: + '@opentelemetry/api': ^1.3.0 + + '@opentelemetry/otlp-exporter-base@0.219.0': + resolution: {integrity: sha512-zvIxQX/AZUVKDU+hCuYx+7UkiP7GRdnk1ZbFQRYzHvYp47cAWR4j3IhoPhV9KaeXEv2xdGq3IA6PnpzDmLcmSA==} + engines: {node: ^18.19.0 || >=20.6.0} + peerDependencies: + '@opentelemetry/api': ^1.3.0 + + '@opentelemetry/otlp-grpc-exporter-base@0.219.0': + resolution: {integrity: sha512-iIk/s8QQu39zpTrRRmsW/Eg3SE2+Hg8tLWepr2FLRgmwUpNd0IpCTLJEHJ77hpt4hgIS8MAh44UYI4xQPZwWlw==} + engines: {node: ^18.19.0 || >=20.6.0} + peerDependencies: + '@opentelemetry/api': ^1.3.0 + + '@opentelemetry/otlp-transformer@0.219.0': + resolution: {integrity: sha512-aaYKAyXhw9VchKZVGOopD3Gw/kPsyrX2c6IQ0AW32mTjqmZOh5Y6Gf5OYqTNqVktAeBjmFinhyFaCwW6GYK9YQ==} + engines: {node: ^18.19.0 || >=20.6.0} + peerDependencies: + '@opentelemetry/api': ^1.3.0 + + '@opentelemetry/propagator-b3@2.8.0': + resolution: {integrity: sha512-SazlvuSKi5533rPHTW2TwBwdMakhjZST4SYs0YauuvfGDkT13KbG1gJS75hV0uWVeevhtVP9sAIlaZLTHdSbMg==} + engines: {node: ^18.19.0 || >=20.6.0} + peerDependencies: + '@opentelemetry/api': '>=1.0.0 <1.10.0' + + '@opentelemetry/propagator-jaeger@2.8.0': + resolution: {integrity: sha512-Xnz9zZvvQzUw+9DrOn0MomR7BxFCkA2pcfXBQuHC28ndJpSbjLs7knzYb05kw5SyCjSsEWombkZMgGcJSk8JVg==} + engines: {node: ^18.19.0 || >=20.6.0} + peerDependencies: + '@opentelemetry/api': '>=1.0.0 <1.10.0' + + '@opentelemetry/redis-common@0.38.3': + resolution: {integrity: sha512-VCghU1JYs/4gP6Gqf/xro9MEsZ7LrMv2uONVsaESKL38ZOB9BqnI98FfS23wjMnHlpuE+TTaWSoAVNpTwYXzjw==} + engines: {node: ^18.19.0 || >=20.6.0} + + '@opentelemetry/resources@2.8.0': + resolution: {integrity: sha512-qmXQ27ilDbUK/vGMqwL8D4/rhn76C+sherM4wTbjlfknR8Nvfc/hCxjRJPhkzZzUsPiNg16SA31NxMabwttRjg==} engines: {node: ^18.19.0 || >=20.6.0} peerDependencies: '@opentelemetry/api': '>=1.3.0 <1.10.0' - '@opentelemetry/sdk-trace-base@2.7.0': - resolution: {integrity: sha512-Yg9zEXJB50DLVLpsKPk7NmNqlPlS+OvqhJGh0A8oawIOTPOwlm4eXs9BMJV7L79lvEwI+dWtAj+YjTyddV336A==} + '@opentelemetry/sdk-logs@0.219.0': + resolution: {integrity: sha512-s6lTKRakaPClvKoWHRChxnXjDMkM/TQ30ff78jN6EBGf7MI7VzANE5PU3f4z9qDUudWjvZjOLHG0rBnBKYvoXA==} + engines: {node: ^18.19.0 || >=20.6.0} + peerDependencies: + '@opentelemetry/api': '>=1.4.0 <1.10.0' + + '@opentelemetry/sdk-metrics@2.8.0': + resolution: {integrity: sha512-UDBGaj6W0Rgy5rTTaoxs8gVGF/aGkAKyjurJv7se6wjRxJu7FoquTLT/vt54DZfo4crbprYfhX/SOK9+BPw1qg==} + engines: {node: ^18.19.0 || >=20.6.0} + peerDependencies: + '@opentelemetry/api': '>=1.9.0 <1.10.0' + + '@opentelemetry/sdk-node@0.219.0': + resolution: {integrity: sha512-NWLpWLEb8gV3+JBHYoIrktbM385wyHpRJoh3J/4Q52d4PR+AlPMNGJT3DzBUrDSUEVbKAXoHR+EDAPxtiNcj8g==} engines: {node: ^18.19.0 || >=20.6.0} peerDependencies: '@opentelemetry/api': '>=1.3.0 <1.10.0' - '@opentelemetry/semantic-conventions@1.40.0': - resolution: {integrity: sha512-cifvXDhcqMwwTlTK04GBNeIe7yyo28Mfby85QXFe1Yk8nmi36Ab/5UQwptOx84SsoGNRg+EVSjwzfSZMy6pmlw==} + '@opentelemetry/sdk-trace-base@2.8.0': + resolution: {integrity: sha512-mhU4jp+vW0mGbFRd+GeXHvmfA4aDqWjBjLC3pE5XMpLs0IE2ryYb019Ts2AQrOq67gaTF25D91+fgvEHDZEnuQ==} + engines: {node: ^18.19.0 || >=20.6.0} + peerDependencies: + '@opentelemetry/api': '>=1.3.0 <1.10.0' + + '@opentelemetry/sdk-trace-node@2.8.0': + resolution: {integrity: sha512-nZt9OGufioAc3AfoLTqA9bsAeaMJAictYDdI2VcNQ+PmT+3rfKjAZDZvgPfd8VPX0O5Bw1hdQF6kDK8VSpZiWg==} + engines: {node: ^18.19.0 || >=20.6.0} + peerDependencies: + '@opentelemetry/api': '>=1.0.0 <1.10.0' + + '@opentelemetry/semantic-conventions@1.41.1': + resolution: {integrity: sha512-/UhIkaZgPutTFmQ7RnIJGgDXZmtEJ7Dvi86xNTFWcnRxVRNk/aotsqDJYeEvDP+FSMB2SdW+pQzNMcWP0rwuNA==} engines: {node: '>=14'} + '@opentelemetry/sql-common@0.42.0': + resolution: {integrity: sha512-nwUwUU+8O8a4bnLqk6CodWeegGMEANgC94KTAhXcpGWLrW/2/hek/0ajNbjXnSOoNuCX+nteUPs46HFHhou9Xw==} + engines: {node: ^18.19.0 || >=20.6.0} + peerDependencies: + '@opentelemetry/api': ^1.1.0 + '@oslojs/encoding@1.1.0': resolution: {integrity: sha512-70wQhgYmndg4GCPxPPxPGevRKqTIJ2Nh4OkiMWmDAVYsTQ+Ta7Sq+rPevXyXGdzr30/qZBnyOalCszoMxlyldQ==} @@ -3591,6 +3829,12 @@ packages: '@types/pdfkit@0.17.6': resolution: {integrity: sha512-tIwzxk2uWKp0Cq9JIluQXJid77lYhF52EsIOwhsMF4iWLA6YneoBR1xVKYYdAysHuepUB0OX4tdwMiUDdGKmig==} + '@types/pg-pool@2.0.7': + resolution: {integrity: sha512-U4CwmGVQcbEuqpyju8/ptOKg6gEC+Tqsvj2xS9o1g71bUh8twxnC6ZL5rZKCsGN0iyH0CwgUyc9VR5owNQF9Ng==} + + '@types/pg@8.15.6': + resolution: {integrity: sha512-NoaMtzhxOrubeL/7UZuNTrejB4MPAJ0RpxZqXQf2qXuVlTPuG6Y8p4u9dKRaue4yjmC7ZhzVO2/Yyyn25znrPQ==} + '@types/pg@8.20.0': resolution: {integrity: sha512-bEPFOaMAHTEP1EzpvHTbmwR8UsFyHSKsRisLIHVMXnpNefSbGA1bD6CVy+qKjGSqmZqNqBDV2azOBo8TgkcVow==} @@ -5131,6 +5375,9 @@ packages: resolution: {integrity: sha512-wzsgA6WOq+09wrU1tsJ09udeR/YZRaeArL9e1wPbFg3GG2yDnC2ldKpxs4xunpFF9DgqCqOIra3bc1HWrJ37Ww==} engines: {node: '>=0.4.x'} + forwarded-parse@2.1.2: + resolution: {integrity: sha512-alTFZZQDKMporBH77856pXgzhEzaUVmLCDk+egLgIgHst3Tpndzz8MnKe+GzRJRfvVdn69HhpW7cmXzvtLvJAw==} + fs-constants@1.0.0: resolution: {integrity: sha512-y6OAwoSIf7FyjMIv94u+b5rdheZEjzR63GTyZJm5qh4Bi+2YgwLCcI/fPFZkL5PSixOt6ZNKm+w+Hfp/Bciwow==} @@ -10789,16 +11036,203 @@ snapshots: dependencies: '@octokit/openapi-types': 27.0.0 + '@opentelemetry/api-logs@0.213.0': + dependencies: + '@opentelemetry/api': 1.9.1 + '@opentelemetry/api-logs@0.214.0': dependencies: '@opentelemetry/api': 1.9.1 + '@opentelemetry/api-logs@0.219.0': + dependencies: + '@opentelemetry/api': 1.9.1 + '@opentelemetry/api@1.9.1': {} + '@opentelemetry/configuration@0.219.0(@opentelemetry/api@1.9.1)': + dependencies: + '@opentelemetry/api': 1.9.1 + '@opentelemetry/core': 2.8.0(@opentelemetry/api@1.9.1) + yaml: 2.8.3 + + '@opentelemetry/context-async-hooks@2.8.0(@opentelemetry/api@1.9.1)': + dependencies: + '@opentelemetry/api': 1.9.1 + '@opentelemetry/core@2.7.0(@opentelemetry/api@1.9.1)': dependencies: '@opentelemetry/api': 1.9.1 - '@opentelemetry/semantic-conventions': 1.40.0 + '@opentelemetry/semantic-conventions': 1.41.1 + + '@opentelemetry/core@2.8.0(@opentelemetry/api@1.9.1)': + dependencies: + '@opentelemetry/api': 1.9.1 + '@opentelemetry/semantic-conventions': 1.41.1 + + '@opentelemetry/exporter-logs-otlp-grpc@0.219.0(@opentelemetry/api@1.9.1)': + dependencies: + '@grpc/grpc-js': 1.14.4 + '@opentelemetry/api': 1.9.1 + '@opentelemetry/core': 2.8.0(@opentelemetry/api@1.9.1) + '@opentelemetry/otlp-exporter-base': 0.219.0(@opentelemetry/api@1.9.1) + '@opentelemetry/otlp-grpc-exporter-base': 0.219.0(@opentelemetry/api@1.9.1) + '@opentelemetry/otlp-transformer': 0.219.0(@opentelemetry/api@1.9.1) + '@opentelemetry/sdk-logs': 0.219.0(@opentelemetry/api@1.9.1) + + '@opentelemetry/exporter-logs-otlp-http@0.219.0(@opentelemetry/api@1.9.1)': + dependencies: + '@opentelemetry/api': 1.9.1 + '@opentelemetry/api-logs': 0.219.0 + '@opentelemetry/core': 2.8.0(@opentelemetry/api@1.9.1) + '@opentelemetry/otlp-exporter-base': 0.219.0(@opentelemetry/api@1.9.1) + '@opentelemetry/otlp-transformer': 0.219.0(@opentelemetry/api@1.9.1) + '@opentelemetry/sdk-logs': 0.219.0(@opentelemetry/api@1.9.1) + + '@opentelemetry/exporter-logs-otlp-proto@0.219.0(@opentelemetry/api@1.9.1)': + dependencies: + '@opentelemetry/api': 1.9.1 + '@opentelemetry/api-logs': 0.219.0 + '@opentelemetry/core': 2.8.0(@opentelemetry/api@1.9.1) + '@opentelemetry/otlp-exporter-base': 0.219.0(@opentelemetry/api@1.9.1) + '@opentelemetry/otlp-transformer': 0.219.0(@opentelemetry/api@1.9.1) + '@opentelemetry/resources': 2.8.0(@opentelemetry/api@1.9.1) + '@opentelemetry/sdk-logs': 0.219.0(@opentelemetry/api@1.9.1) + '@opentelemetry/sdk-trace-base': 2.8.0(@opentelemetry/api@1.9.1) + + '@opentelemetry/exporter-metrics-otlp-grpc@0.219.0(@opentelemetry/api@1.9.1)': + dependencies: + '@grpc/grpc-js': 1.14.4 + '@opentelemetry/api': 1.9.1 + '@opentelemetry/core': 2.8.0(@opentelemetry/api@1.9.1) + '@opentelemetry/exporter-metrics-otlp-http': 0.219.0(@opentelemetry/api@1.9.1) + '@opentelemetry/otlp-exporter-base': 0.219.0(@opentelemetry/api@1.9.1) + '@opentelemetry/otlp-grpc-exporter-base': 0.219.0(@opentelemetry/api@1.9.1) + '@opentelemetry/otlp-transformer': 0.219.0(@opentelemetry/api@1.9.1) + '@opentelemetry/resources': 2.8.0(@opentelemetry/api@1.9.1) + '@opentelemetry/sdk-metrics': 2.8.0(@opentelemetry/api@1.9.1) + + '@opentelemetry/exporter-metrics-otlp-http@0.219.0(@opentelemetry/api@1.9.1)': + dependencies: + '@opentelemetry/api': 1.9.1 + '@opentelemetry/core': 2.8.0(@opentelemetry/api@1.9.1) + '@opentelemetry/otlp-exporter-base': 0.219.0(@opentelemetry/api@1.9.1) + '@opentelemetry/otlp-transformer': 0.219.0(@opentelemetry/api@1.9.1) + '@opentelemetry/resources': 2.8.0(@opentelemetry/api@1.9.1) + '@opentelemetry/sdk-metrics': 2.8.0(@opentelemetry/api@1.9.1) + + '@opentelemetry/exporter-metrics-otlp-proto@0.219.0(@opentelemetry/api@1.9.1)': + dependencies: + '@opentelemetry/api': 1.9.1 + '@opentelemetry/core': 2.8.0(@opentelemetry/api@1.9.1) + '@opentelemetry/exporter-metrics-otlp-http': 0.219.0(@opentelemetry/api@1.9.1) + '@opentelemetry/otlp-exporter-base': 0.219.0(@opentelemetry/api@1.9.1) + '@opentelemetry/otlp-transformer': 0.219.0(@opentelemetry/api@1.9.1) + '@opentelemetry/resources': 2.8.0(@opentelemetry/api@1.9.1) + '@opentelemetry/sdk-metrics': 2.8.0(@opentelemetry/api@1.9.1) + + '@opentelemetry/exporter-prometheus@0.219.0(@opentelemetry/api@1.9.1)': + dependencies: + '@opentelemetry/api': 1.9.1 + '@opentelemetry/core': 2.8.0(@opentelemetry/api@1.9.1) + '@opentelemetry/resources': 2.8.0(@opentelemetry/api@1.9.1) + '@opentelemetry/sdk-metrics': 2.8.0(@opentelemetry/api@1.9.1) + '@opentelemetry/semantic-conventions': 1.41.1 + + '@opentelemetry/exporter-trace-otlp-grpc@0.219.0(@opentelemetry/api@1.9.1)': + dependencies: + '@grpc/grpc-js': 1.14.4 + '@opentelemetry/api': 1.9.1 + '@opentelemetry/core': 2.8.0(@opentelemetry/api@1.9.1) + '@opentelemetry/otlp-exporter-base': 0.219.0(@opentelemetry/api@1.9.1) + '@opentelemetry/otlp-grpc-exporter-base': 0.219.0(@opentelemetry/api@1.9.1) + '@opentelemetry/otlp-transformer': 0.219.0(@opentelemetry/api@1.9.1) + '@opentelemetry/resources': 2.8.0(@opentelemetry/api@1.9.1) + '@opentelemetry/sdk-trace-base': 2.8.0(@opentelemetry/api@1.9.1) + + '@opentelemetry/exporter-trace-otlp-http@0.219.0(@opentelemetry/api@1.9.1)': + dependencies: + '@opentelemetry/api': 1.9.1 + '@opentelemetry/core': 2.8.0(@opentelemetry/api@1.9.1) + '@opentelemetry/otlp-exporter-base': 0.219.0(@opentelemetry/api@1.9.1) + '@opentelemetry/otlp-transformer': 0.219.0(@opentelemetry/api@1.9.1) + '@opentelemetry/resources': 2.8.0(@opentelemetry/api@1.9.1) + '@opentelemetry/sdk-trace-base': 2.8.0(@opentelemetry/api@1.9.1) + + '@opentelemetry/exporter-trace-otlp-proto@0.219.0(@opentelemetry/api@1.9.1)': + dependencies: + '@opentelemetry/api': 1.9.1 + '@opentelemetry/core': 2.8.0(@opentelemetry/api@1.9.1) + '@opentelemetry/otlp-exporter-base': 0.219.0(@opentelemetry/api@1.9.1) + '@opentelemetry/otlp-transformer': 0.219.0(@opentelemetry/api@1.9.1) + '@opentelemetry/resources': 2.8.0(@opentelemetry/api@1.9.1) + '@opentelemetry/sdk-trace-base': 2.8.0(@opentelemetry/api@1.9.1) + + '@opentelemetry/exporter-zipkin@2.8.0(@opentelemetry/api@1.9.1)': + dependencies: + '@opentelemetry/api': 1.9.1 + '@opentelemetry/core': 2.8.0(@opentelemetry/api@1.9.1) + '@opentelemetry/resources': 2.8.0(@opentelemetry/api@1.9.1) + '@opentelemetry/sdk-trace-base': 2.8.0(@opentelemetry/api@1.9.1) + '@opentelemetry/semantic-conventions': 1.41.1 + + '@opentelemetry/instrumentation-aws-sdk@0.74.0(@opentelemetry/api@1.9.1)': + dependencies: + '@opentelemetry/api': 1.9.1 + '@opentelemetry/core': 2.7.0(@opentelemetry/api@1.9.1) + '@opentelemetry/instrumentation': 0.219.0(@opentelemetry/api@1.9.1) + '@opentelemetry/semantic-conventions': 1.41.1 + transitivePeerDependencies: + - supports-color + + '@opentelemetry/instrumentation-fastify@0.57.0(@opentelemetry/api@1.9.1)': + dependencies: + '@opentelemetry/api': 1.9.1 + '@opentelemetry/core': 2.7.0(@opentelemetry/api@1.9.1) + '@opentelemetry/instrumentation': 0.213.0(@opentelemetry/api@1.9.1) + '@opentelemetry/semantic-conventions': 1.41.1 + transitivePeerDependencies: + - supports-color + + '@opentelemetry/instrumentation-http@0.219.0(@opentelemetry/api@1.9.1)': + dependencies: + '@opentelemetry/api': 1.9.1 + '@opentelemetry/core': 2.8.0(@opentelemetry/api@1.9.1) + '@opentelemetry/instrumentation': 0.219.0(@opentelemetry/api@1.9.1) + '@opentelemetry/semantic-conventions': 1.41.1 + forwarded-parse: 2.1.2 + transitivePeerDependencies: + - supports-color + + '@opentelemetry/instrumentation-ioredis@0.67.0(@opentelemetry/api@1.9.1)': + dependencies: + '@opentelemetry/api': 1.9.1 + '@opentelemetry/instrumentation': 0.219.0(@opentelemetry/api@1.9.1) + '@opentelemetry/redis-common': 0.38.3 + '@opentelemetry/semantic-conventions': 1.41.1 + transitivePeerDependencies: + - supports-color + + '@opentelemetry/instrumentation-pg@0.71.0(@opentelemetry/api@1.9.1)': + dependencies: + '@opentelemetry/api': 1.9.1 + '@opentelemetry/core': 2.7.0(@opentelemetry/api@1.9.1) + '@opentelemetry/instrumentation': 0.219.0(@opentelemetry/api@1.9.1) + '@opentelemetry/semantic-conventions': 1.41.1 + '@opentelemetry/sql-common': 0.42.0(@opentelemetry/api@1.9.1) + '@types/pg': 8.15.6 + '@types/pg-pool': 2.0.7 + transitivePeerDependencies: + - supports-color + + '@opentelemetry/instrumentation@0.213.0(@opentelemetry/api@1.9.1)': + dependencies: + '@opentelemetry/api': 1.9.1 + '@opentelemetry/api-logs': 0.213.0 + import-in-the-middle: 3.0.1 + require-in-the-middle: 8.0.1 + transitivePeerDependencies: + - supports-color '@opentelemetry/instrumentation@0.214.0(@opentelemetry/api@1.9.1)': dependencies: @@ -10809,20 +11243,123 @@ snapshots: transitivePeerDependencies: - supports-color - '@opentelemetry/resources@2.7.0(@opentelemetry/api@1.9.1)': + '@opentelemetry/instrumentation@0.219.0(@opentelemetry/api@1.9.1)': + dependencies: + '@opentelemetry/api': 1.9.1 + '@opentelemetry/api-logs': 0.219.0 + import-in-the-middle: 3.0.1 + require-in-the-middle: 8.0.1 + transitivePeerDependencies: + - supports-color + + '@opentelemetry/otlp-exporter-base@0.219.0(@opentelemetry/api@1.9.1)': + dependencies: + '@opentelemetry/api': 1.9.1 + '@opentelemetry/core': 2.8.0(@opentelemetry/api@1.9.1) + '@opentelemetry/otlp-transformer': 0.219.0(@opentelemetry/api@1.9.1) + + '@opentelemetry/otlp-grpc-exporter-base@0.219.0(@opentelemetry/api@1.9.1)': + dependencies: + '@grpc/grpc-js': 1.14.4 + '@opentelemetry/api': 1.9.1 + '@opentelemetry/core': 2.8.0(@opentelemetry/api@1.9.1) + '@opentelemetry/otlp-exporter-base': 0.219.0(@opentelemetry/api@1.9.1) + '@opentelemetry/otlp-transformer': 0.219.0(@opentelemetry/api@1.9.1) + + '@opentelemetry/otlp-transformer@0.219.0(@opentelemetry/api@1.9.1)': + dependencies: + '@opentelemetry/api': 1.9.1 + '@opentelemetry/api-logs': 0.219.0 + '@opentelemetry/core': 2.8.0(@opentelemetry/api@1.9.1) + '@opentelemetry/resources': 2.8.0(@opentelemetry/api@1.9.1) + '@opentelemetry/sdk-logs': 0.219.0(@opentelemetry/api@1.9.1) + '@opentelemetry/sdk-metrics': 2.8.0(@opentelemetry/api@1.9.1) + '@opentelemetry/sdk-trace-base': 2.8.0(@opentelemetry/api@1.9.1) + + '@opentelemetry/propagator-b3@2.8.0(@opentelemetry/api@1.9.1)': + dependencies: + '@opentelemetry/api': 1.9.1 + '@opentelemetry/core': 2.8.0(@opentelemetry/api@1.9.1) + + '@opentelemetry/propagator-jaeger@2.8.0(@opentelemetry/api@1.9.1)': + dependencies: + '@opentelemetry/api': 1.9.1 + '@opentelemetry/core': 2.8.0(@opentelemetry/api@1.9.1) + + '@opentelemetry/redis-common@0.38.3': {} + + '@opentelemetry/resources@2.8.0(@opentelemetry/api@1.9.1)': + dependencies: + '@opentelemetry/api': 1.9.1 + '@opentelemetry/core': 2.8.0(@opentelemetry/api@1.9.1) + '@opentelemetry/semantic-conventions': 1.41.1 + + '@opentelemetry/sdk-logs@0.219.0(@opentelemetry/api@1.9.1)': + dependencies: + '@opentelemetry/api': 1.9.1 + '@opentelemetry/api-logs': 0.219.0 + '@opentelemetry/core': 2.8.0(@opentelemetry/api@1.9.1) + '@opentelemetry/resources': 2.8.0(@opentelemetry/api@1.9.1) + '@opentelemetry/semantic-conventions': 1.41.1 + + '@opentelemetry/sdk-metrics@2.8.0(@opentelemetry/api@1.9.1)': + dependencies: + '@opentelemetry/api': 1.9.1 + '@opentelemetry/core': 2.8.0(@opentelemetry/api@1.9.1) + '@opentelemetry/resources': 2.8.0(@opentelemetry/api@1.9.1) + + '@opentelemetry/sdk-node@0.219.0(@opentelemetry/api@1.9.1)': + dependencies: + '@opentelemetry/api': 1.9.1 + '@opentelemetry/api-logs': 0.219.0 + '@opentelemetry/configuration': 0.219.0(@opentelemetry/api@1.9.1) + '@opentelemetry/context-async-hooks': 2.8.0(@opentelemetry/api@1.9.1) + '@opentelemetry/core': 2.8.0(@opentelemetry/api@1.9.1) + '@opentelemetry/exporter-logs-otlp-grpc': 0.219.0(@opentelemetry/api@1.9.1) + '@opentelemetry/exporter-logs-otlp-http': 0.219.0(@opentelemetry/api@1.9.1) + '@opentelemetry/exporter-logs-otlp-proto': 0.219.0(@opentelemetry/api@1.9.1) + '@opentelemetry/exporter-metrics-otlp-grpc': 0.219.0(@opentelemetry/api@1.9.1) + '@opentelemetry/exporter-metrics-otlp-http': 0.219.0(@opentelemetry/api@1.9.1) + '@opentelemetry/exporter-metrics-otlp-proto': 0.219.0(@opentelemetry/api@1.9.1) + '@opentelemetry/exporter-prometheus': 0.219.0(@opentelemetry/api@1.9.1) + '@opentelemetry/exporter-trace-otlp-grpc': 0.219.0(@opentelemetry/api@1.9.1) + '@opentelemetry/exporter-trace-otlp-http': 0.219.0(@opentelemetry/api@1.9.1) + '@opentelemetry/exporter-trace-otlp-proto': 0.219.0(@opentelemetry/api@1.9.1) + '@opentelemetry/exporter-zipkin': 2.8.0(@opentelemetry/api@1.9.1) + '@opentelemetry/instrumentation': 0.219.0(@opentelemetry/api@1.9.1) + '@opentelemetry/otlp-exporter-base': 0.219.0(@opentelemetry/api@1.9.1) + '@opentelemetry/otlp-grpc-exporter-base': 0.219.0(@opentelemetry/api@1.9.1) + '@opentelemetry/propagator-b3': 2.8.0(@opentelemetry/api@1.9.1) + '@opentelemetry/propagator-jaeger': 2.8.0(@opentelemetry/api@1.9.1) + '@opentelemetry/resources': 2.8.0(@opentelemetry/api@1.9.1) + '@opentelemetry/sdk-logs': 0.219.0(@opentelemetry/api@1.9.1) + '@opentelemetry/sdk-metrics': 2.8.0(@opentelemetry/api@1.9.1) + '@opentelemetry/sdk-trace-base': 2.8.0(@opentelemetry/api@1.9.1) + '@opentelemetry/sdk-trace-node': 2.8.0(@opentelemetry/api@1.9.1) + '@opentelemetry/semantic-conventions': 1.41.1 + transitivePeerDependencies: + - supports-color + + '@opentelemetry/sdk-trace-base@2.8.0(@opentelemetry/api@1.9.1)': + dependencies: + '@opentelemetry/api': 1.9.1 + '@opentelemetry/core': 2.8.0(@opentelemetry/api@1.9.1) + '@opentelemetry/resources': 2.8.0(@opentelemetry/api@1.9.1) + '@opentelemetry/semantic-conventions': 1.41.1 + + '@opentelemetry/sdk-trace-node@2.8.0(@opentelemetry/api@1.9.1)': + dependencies: + '@opentelemetry/api': 1.9.1 + '@opentelemetry/context-async-hooks': 2.8.0(@opentelemetry/api@1.9.1) + '@opentelemetry/core': 2.8.0(@opentelemetry/api@1.9.1) + '@opentelemetry/sdk-trace-base': 2.8.0(@opentelemetry/api@1.9.1) + + '@opentelemetry/semantic-conventions@1.41.1': {} + + '@opentelemetry/sql-common@0.42.0(@opentelemetry/api@1.9.1)': dependencies: '@opentelemetry/api': 1.9.1 '@opentelemetry/core': 2.7.0(@opentelemetry/api@1.9.1) - '@opentelemetry/semantic-conventions': 1.40.0 - - '@opentelemetry/sdk-trace-base@2.7.0(@opentelemetry/api@1.9.1)': - dependencies: - '@opentelemetry/api': 1.9.1 - '@opentelemetry/core': 2.7.0(@opentelemetry/api@1.9.1) - '@opentelemetry/resources': 2.7.0(@opentelemetry/api@1.9.1) - '@opentelemetry/semantic-conventions': 1.40.0 - - '@opentelemetry/semantic-conventions@1.40.0': {} '@oslojs/encoding@1.1.0': {} @@ -11114,40 +11651,41 @@ snapshots: '@sentry/core@10.56.0': {} - '@sentry/node-core@10.56.0(@opentelemetry/api@1.9.1)(@opentelemetry/core@2.7.0(@opentelemetry/api@1.9.1))(@opentelemetry/instrumentation@0.214.0(@opentelemetry/api@1.9.1))(@opentelemetry/sdk-trace-base@2.7.0(@opentelemetry/api@1.9.1))(@opentelemetry/semantic-conventions@1.40.0)': + '@sentry/node-core@10.56.0(@opentelemetry/api@1.9.1)(@opentelemetry/core@2.7.0(@opentelemetry/api@1.9.1))(@opentelemetry/exporter-trace-otlp-http@0.219.0(@opentelemetry/api@1.9.1))(@opentelemetry/instrumentation@0.214.0(@opentelemetry/api@1.9.1))(@opentelemetry/sdk-trace-base@2.8.0(@opentelemetry/api@1.9.1))(@opentelemetry/semantic-conventions@1.41.1)': dependencies: '@sentry/core': 10.56.0 - '@sentry/opentelemetry': 10.56.0(@opentelemetry/api@1.9.1)(@opentelemetry/core@2.7.0(@opentelemetry/api@1.9.1))(@opentelemetry/sdk-trace-base@2.7.0(@opentelemetry/api@1.9.1))(@opentelemetry/semantic-conventions@1.40.0) + '@sentry/opentelemetry': 10.56.0(@opentelemetry/api@1.9.1)(@opentelemetry/core@2.7.0(@opentelemetry/api@1.9.1))(@opentelemetry/sdk-trace-base@2.8.0(@opentelemetry/api@1.9.1))(@opentelemetry/semantic-conventions@1.41.1) import-in-the-middle: 3.0.1 optionalDependencies: '@opentelemetry/api': 1.9.1 '@opentelemetry/core': 2.7.0(@opentelemetry/api@1.9.1) + '@opentelemetry/exporter-trace-otlp-http': 0.219.0(@opentelemetry/api@1.9.1) '@opentelemetry/instrumentation': 0.214.0(@opentelemetry/api@1.9.1) - '@opentelemetry/sdk-trace-base': 2.7.0(@opentelemetry/api@1.9.1) - '@opentelemetry/semantic-conventions': 1.40.0 + '@opentelemetry/sdk-trace-base': 2.8.0(@opentelemetry/api@1.9.1) + '@opentelemetry/semantic-conventions': 1.41.1 - '@sentry/node@10.56.0': + '@sentry/node@10.56.0(@opentelemetry/exporter-trace-otlp-http@0.219.0(@opentelemetry/api@1.9.1))': dependencies: '@opentelemetry/api': 1.9.1 '@opentelemetry/core': 2.7.0(@opentelemetry/api@1.9.1) '@opentelemetry/instrumentation': 0.214.0(@opentelemetry/api@1.9.1) - '@opentelemetry/sdk-trace-base': 2.7.0(@opentelemetry/api@1.9.1) - '@opentelemetry/semantic-conventions': 1.40.0 + '@opentelemetry/sdk-trace-base': 2.8.0(@opentelemetry/api@1.9.1) + '@opentelemetry/semantic-conventions': 1.41.1 '@sentry-internal/server-utils': 10.56.0 '@sentry/core': 10.56.0 - '@sentry/node-core': 10.56.0(@opentelemetry/api@1.9.1)(@opentelemetry/core@2.7.0(@opentelemetry/api@1.9.1))(@opentelemetry/instrumentation@0.214.0(@opentelemetry/api@1.9.1))(@opentelemetry/sdk-trace-base@2.7.0(@opentelemetry/api@1.9.1))(@opentelemetry/semantic-conventions@1.40.0) - '@sentry/opentelemetry': 10.56.0(@opentelemetry/api@1.9.1)(@opentelemetry/core@2.7.0(@opentelemetry/api@1.9.1))(@opentelemetry/sdk-trace-base@2.7.0(@opentelemetry/api@1.9.1))(@opentelemetry/semantic-conventions@1.40.0) + '@sentry/node-core': 10.56.0(@opentelemetry/api@1.9.1)(@opentelemetry/core@2.7.0(@opentelemetry/api@1.9.1))(@opentelemetry/exporter-trace-otlp-http@0.219.0(@opentelemetry/api@1.9.1))(@opentelemetry/instrumentation@0.214.0(@opentelemetry/api@1.9.1))(@opentelemetry/sdk-trace-base@2.8.0(@opentelemetry/api@1.9.1))(@opentelemetry/semantic-conventions@1.41.1) + '@sentry/opentelemetry': 10.56.0(@opentelemetry/api@1.9.1)(@opentelemetry/core@2.7.0(@opentelemetry/api@1.9.1))(@opentelemetry/sdk-trace-base@2.8.0(@opentelemetry/api@1.9.1))(@opentelemetry/semantic-conventions@1.41.1) import-in-the-middle: 3.0.1 transitivePeerDependencies: - '@opentelemetry/exporter-trace-otlp-http' - supports-color - '@sentry/opentelemetry@10.56.0(@opentelemetry/api@1.9.1)(@opentelemetry/core@2.7.0(@opentelemetry/api@1.9.1))(@opentelemetry/sdk-trace-base@2.7.0(@opentelemetry/api@1.9.1))(@opentelemetry/semantic-conventions@1.40.0)': + '@sentry/opentelemetry@10.56.0(@opentelemetry/api@1.9.1)(@opentelemetry/core@2.7.0(@opentelemetry/api@1.9.1))(@opentelemetry/sdk-trace-base@2.8.0(@opentelemetry/api@1.9.1))(@opentelemetry/semantic-conventions@1.41.1)': dependencies: '@opentelemetry/api': 1.9.1 '@opentelemetry/core': 2.7.0(@opentelemetry/api@1.9.1) - '@opentelemetry/sdk-trace-base': 2.7.0(@opentelemetry/api@1.9.1) - '@opentelemetry/semantic-conventions': 1.40.0 + '@opentelemetry/sdk-trace-base': 2.8.0(@opentelemetry/api@1.9.1) + '@opentelemetry/semantic-conventions': 1.41.1 '@sentry/core': 10.56.0 '@sentry/react@10.56.0(react@19.2.7)': @@ -11629,6 +12167,16 @@ snapshots: dependencies: '@types/node': 22.19.19 + '@types/pg-pool@2.0.7': + dependencies: + '@types/pg': 8.20.0 + + '@types/pg@8.15.6': + dependencies: + '@types/node': 22.19.19 + pg-protocol: 1.14.0 + pg-types: 2.2.0 + '@types/pg@8.20.0': dependencies: '@types/node': 22.19.19 @@ -13356,6 +13904,8 @@ snapshots: format@0.2.2: {} + forwarded-parse@2.1.2: {} + fs-constants@1.0.0: {} fs-extra@11.3.4: diff --git a/tests/integration/tracing-lifecycle.test.ts b/tests/integration/tracing-lifecycle.test.ts new file mode 100644 index 00000000..d8de8421 --- /dev/null +++ b/tests/integration/tracing-lifecycle.test.ts @@ -0,0 +1,237 @@ +import { AsyncLocalStorage } from "node:async_hooks"; +import { + type Context, + type ContextManager, + context, + propagation, + ROOT_CONTEXT, + SpanStatusCode, + trace, +} from "@opentelemetry/api"; +import { W3CTraceContextPropagator } from "@opentelemetry/core"; +import { afterEach, describe, expect, it } from "vitest"; + +/** + * Minimal AsyncLocalStorage-based context manager for tests. + * OTel v2 BasicTracerProvider no longer registers one automatically, + * and @opentelemetry/context-async-hooks is not a direct dependency. + */ +class TestContextManager implements ContextManager { + private _als = new AsyncLocalStorage(); + + active(): Context { + return this._als.getStore() ?? ROOT_CONTEXT; + } + + with ReturnType>( + ctx: Context, + fn: F, + thisArg?: ThisParameterType, + ...args: A + ): ReturnType { + return this._als.run(ctx, () => fn.call(thisArg, ...args)); + } + + bind(_ctx: Context, target: T): T { + return target; + } + + enable(): this { + return this; + } + + disable(): this { + this._als.disable(); + return this; + } +} + +describe("tracing lifecycle", () => { + afterEach(() => { + propagation.disable(); + context.disable(); + trace.disable(); + }); + + it("propagates trace context through inject/extract cycle", async () => { + // Create HTTP span, inject to carrier, extract in "worker", create child + // Assert same traceId, different spanId, correct parent-child + const { InMemorySpanExporter, SimpleSpanProcessor, BasicTracerProvider } = await import( + "@opentelemetry/sdk-trace-base" + ); + + context.setGlobalContextManager(new TestContextManager()); + propagation.setGlobalPropagator(new W3CTraceContextPropagator()); + + const exporter = new InMemorySpanExporter(); + const provider = new BasicTracerProvider({ + spanProcessors: [new SimpleSpanProcessor(exporter)], + }); + trace.setGlobalTracerProvider(provider); + + const tracer = trace.getTracer("test"); + const httpSpan = tracer.startSpan("HTTP POST /api/v1/tools/resize"); + const httpCtx = trace.setSpan(ROOT_CONTEXT, httpSpan); + const httpTraceId = httpSpan.spanContext().traceId; + + const carrier: Record = {}; + context.with(httpCtx, () => { + propagation.inject(context.active(), carrier); + }); + + expect(carrier.traceparent).toBeDefined(); + + const workerCtx = propagation.extract(ROOT_CONTEXT, carrier); + const jobSpan = tracer.startSpan("job.process", {}, workerCtx); + + expect(jobSpan.spanContext().traceId).toBe(httpTraceId); + + const toolCtx = trace.setSpan(workerCtx, jobSpan); + const toolSpan = tracer.startSpan("tool.process", {}, toolCtx); + + expect(toolSpan.spanContext().traceId).toBe(httpTraceId); + + toolSpan.end(); + jobSpan.end(); + httpSpan.end(); + + const spans = exporter.getFinishedSpans(); + const names = spans.map((s) => s.name); + expect(names).toContain("HTTP POST /api/v1/tools/resize"); + expect(names).toContain("job.process"); + expect(names).toContain("tool.process"); + + // OTel SDK v2: parent reference is parentSpanContext (SpanContext object), + // not parentSpanId (string). + const jobFinished = spans.find((s) => s.name === "job.process")!; + expect(jobFinished.parentSpanContext?.spanId).toBe(httpSpan.spanContext().spanId); + + const toolFinished = spans.find((s) => s.name === "tool.process")!; + expect(toolFinished.parentSpanContext?.spanId).toBe(jobSpan.spanContext().spanId); + + await provider.shutdown(); + }); + + it("records error status on failed spans", async () => { + const { InMemorySpanExporter, SimpleSpanProcessor, BasicTracerProvider } = await import( + "@opentelemetry/sdk-trace-base" + ); + + context.setGlobalContextManager(new TestContextManager()); + propagation.setGlobalPropagator(new W3CTraceContextPropagator()); + + const exporter = new InMemorySpanExporter(); + const provider = new BasicTracerProvider({ + spanProcessors: [new SimpleSpanProcessor(exporter)], + }); + trace.setGlobalTracerProvider(provider); + + const tracer = trace.getTracer("test"); + const span = tracer.startSpan("job.process"); + + const error = new Error("Tool processing failed"); + span.setStatus({ code: SpanStatusCode.ERROR, message: error.message }); + span.recordException(error); + span.addEvent("job.failed"); + span.end(); + + const spans = exporter.getFinishedSpans(); + const finished = spans.find((s) => s.name === "job.process")!; + + expect(finished.status.code).toBe(SpanStatusCode.ERROR); + expect(finished.status.message).toBe("Tool processing failed"); + expect(finished.events.some((e) => e.name === "job.failed")).toBe(true); + expect(finished.events.some((e) => e.name === "exception")).toBe(true); + + await provider.shutdown(); + }); + + it("handles absent _otel gracefully", () => { + const extractedCtx = propagation.extract(ROOT_CONTEXT, {}); + const span = trace.getSpan(extractedCtx); + expect(span).toBeUndefined(); + }); + + it("tracing failures do not block request processing", async () => { + const { InMemorySpanExporter, SimpleSpanProcessor, BasicTracerProvider } = await import( + "@opentelemetry/sdk-trace-base" + ); + + context.setGlobalContextManager(new TestContextManager()); + propagation.setGlobalPropagator(new W3CTraceContextPropagator()); + + const exporter = new InMemorySpanExporter(); + const provider = new BasicTracerProvider({ + spanProcessors: [new SimpleSpanProcessor(exporter)], + }); + trace.setGlobalTracerProvider(provider); + + const tracer = trace.getTracer("test"); + const span = tracer.startSpan("request-with-bad-collector"); + expect(span).toBeDefined(); + expect(span.spanContext().traceId).toMatch(/^[0-9a-f]{32}$/); + + const carrier: Record = {}; + const ctx = trace.setSpan(ROOT_CONTEXT, span); + context.with(ctx, () => propagation.inject(context.active(), carrier)); + expect(carrier.traceparent).toBeDefined(); + + span.end(); + await provider.shutdown(); + }); + + it("supports retry spans under the same trace", async () => { + const { InMemorySpanExporter, SimpleSpanProcessor, BasicTracerProvider } = await import( + "@opentelemetry/sdk-trace-base" + ); + + context.setGlobalContextManager(new TestContextManager()); + propagation.setGlobalPropagator(new W3CTraceContextPropagator()); + + const exporter = new InMemorySpanExporter(); + const provider = new BasicTracerProvider({ + spanProcessors: [new SimpleSpanProcessor(exporter)], + }); + trace.setGlobalTracerProvider(provider); + + const tracer = trace.getTracer("test"); + const httpSpan = tracer.startSpan("HTTP POST"); + const httpCtx = trace.setSpan(ROOT_CONTEXT, httpSpan); + const traceId = httpSpan.spanContext().traceId; + + const carrier: Record = {}; + context.with(httpCtx, () => propagation.inject(context.active(), carrier)); + + const ctx1 = propagation.extract(ROOT_CONTEXT, carrier); + const attempt1 = tracer.startSpan( + "job.attempt", + { + attributes: { "snapotter.attempt_number": 1 }, + }, + ctx1, + ); + attempt1.setStatus({ code: SpanStatusCode.ERROR, message: "timeout" }); + attempt1.end(); + + const ctx2 = propagation.extract(ROOT_CONTEXT, carrier); + const attempt2 = tracer.startSpan( + "job.attempt", + { + attributes: { "snapotter.attempt_number": 2 }, + }, + ctx2, + ); + attempt2.end(); + + const spans = exporter.getFinishedSpans(); + const attempts = spans.filter((s) => s.name === "job.attempt"); + expect(attempts).toHaveLength(2); + expect(attempts[0].spanContext().traceId).toBe(traceId); + expect(attempts[1].spanContext().traceId).toBe(traceId); + expect(attempts[0].attributes["snapotter.attempt_number"]).toBe(1); + expect(attempts[1].attributes["snapotter.attempt_number"]).toBe(2); + + httpSpan.end(); + await provider.shutdown(); + }); +}); diff --git a/tests/unit/api/bullmq-context-propagation.test.ts b/tests/unit/api/bullmq-context-propagation.test.ts new file mode 100644 index 00000000..84d04bfe --- /dev/null +++ b/tests/unit/api/bullmq-context-propagation.test.ts @@ -0,0 +1,145 @@ +import { AsyncLocalStorage } from "node:async_hooks"; +import { + type Context, + type ContextManager, + context, + propagation, + ROOT_CONTEXT, + trace, +} from "@opentelemetry/api"; +import { W3CTraceContextPropagator } from "@opentelemetry/core"; +import { afterEach, describe, expect, it } from "vitest"; + +/** + * Minimal AsyncLocalStorage-based context manager for tests. + * OTel v2 BasicTracerProvider no longer registers one automatically. + */ +class TestContextManager implements ContextManager { + private _als = new AsyncLocalStorage(); + + active(): Context { + return this._als.getStore() ?? ROOT_CONTEXT; + } + + with ReturnType>( + ctx: Context, + fn: F, + thisArg?: ThisParameterType, + ...args: A + ): ReturnType { + return this._als.run(ctx, () => fn.call(thisArg, ...args)); + } + + bind(_ctx: Context, target: T): T { + return target; + } + + enable(): this { + return this; + } + + disable(): this { + this._als.disable(); + return this; + } +} + +describe("BullMQ trace context injection", () => { + afterEach(() => { + propagation.disable(); + context.disable(); + trace.disable(); + }); + + it("injects _otel with traceparent when a span is active", async () => { + const { InMemorySpanExporter, SimpleSpanProcessor, BasicTracerProvider } = await import( + "@opentelemetry/sdk-trace-base" + ); + + context.setGlobalContextManager(new TestContextManager()); + propagation.setGlobalPropagator(new W3CTraceContextPropagator()); + + const exporter = new InMemorySpanExporter(); + const provider = new BasicTracerProvider({ + spanProcessors: [new SimpleSpanProcessor(exporter)], + }); + trace.setGlobalTracerProvider(provider); + + const tracer = trace.getTracer("test"); + const span = tracer.startSpan("http-request"); + const ctx = trace.setSpan(ROOT_CONTEXT, span); + + const carrier: Record = {}; + context.with(ctx, () => { + propagation.inject(context.active(), carrier); + }); + + expect(carrier.traceparent).toBeDefined(); + expect(carrier.traceparent).toMatch(/^00-[0-9a-f]{32}-[0-9a-f]{16}-0[01]$/); + + // Verify traceparent contains the correct traceId and spanId + const parts = carrier.traceparent!.split("-"); + expect(parts[1]).toBe(span.spanContext().traceId); + expect(parts[2]).toBe(span.spanContext().spanId); + + span.end(); + await provider.shutdown(); + }); + + it("injects nothing when no SDK is registered", () => { + const carrier: Record = {}; + propagation.inject(context.active(), carrier); + expect(carrier.traceparent).toBeUndefined(); + }); +}); + +describe("BullMQ trace context extraction", () => { + afterEach(() => { + propagation.disable(); + context.disable(); + trace.disable(); + }); + + it("extracts parent context from _otel carrier", async () => { + const { InMemorySpanExporter, SimpleSpanProcessor, BasicTracerProvider } = await import( + "@opentelemetry/sdk-trace-base" + ); + + context.setGlobalContextManager(new TestContextManager()); + propagation.setGlobalPropagator(new W3CTraceContextPropagator()); + + const exporter = new InMemorySpanExporter(); + const provider = new BasicTracerProvider({ + spanProcessors: [new SimpleSpanProcessor(exporter)], + }); + trace.setGlobalTracerProvider(provider); + + // Create a parent span and inject its context + const tracer = trace.getTracer("test"); + const parentSpan = tracer.startSpan("parent-request"); + const parentCtx = trace.setSpan(ROOT_CONTEXT, parentSpan); + + const carrier: Record = {}; + context.with(parentCtx, () => { + propagation.inject(context.active(), carrier); + }); + + // Extract context from the carrier (simulating worker side) + const extractedCtx = propagation.extract(ROOT_CONTEXT, carrier); + const childSpan = tracer.startSpan("worker-process", undefined, extractedCtx); + + // Child should share traceId but have a different spanId + expect(childSpan.spanContext().traceId).toBe(parentSpan.spanContext().traceId); + expect(childSpan.spanContext().spanId).not.toBe(parentSpan.spanContext().spanId); + + parentSpan.end(); + childSpan.end(); + await provider.shutdown(); + }); + + it("creates root span when _otel is absent", () => { + const extractedCtx = propagation.extract(ROOT_CONTEXT, {}); + const span = trace.getSpan(extractedCtx); + expect(span).toBeUndefined(); + }); +}); diff --git a/tests/unit/api/enterprise-flags.test.ts b/tests/unit/api/enterprise-flags.test.ts index 9aed611a..f47c8061 100644 --- a/tests/unit/api/enterprise-flags.test.ts +++ b/tests/unit/api/enterprise-flags.test.ts @@ -32,7 +32,13 @@ describe("enterprise feature flags", () => { } }); - it("has exactly 18 features total", () => { - expect(ENTERPRISE_FEATURES).toHaveLength(18); + it("has exactly 19 features total", () => { + expect(ENTERPRISE_FEATURES).toHaveLength(19); + }); + + it("includes distributed_tracing in enterprise plan only", () => { + expect(ENTERPRISE_FEATURES).toContain("distributed_tracing"); + expect(PLAN_FEATURES.enterprise).toContain("distributed_tracing"); + expect(PLAN_FEATURES.team).not.toContain("distributed_tracing"); }); }); diff --git a/tests/unit/api/log-trace-mixin.test.ts b/tests/unit/api/log-trace-mixin.test.ts new file mode 100644 index 00000000..0f636068 --- /dev/null +++ b/tests/unit/api/log-trace-mixin.test.ts @@ -0,0 +1,84 @@ +import { AsyncLocalStorage } from "node:async_hooks"; +import { + type Context, + type ContextManager, + context, + ROOT_CONTEXT, + trace, +} from "@opentelemetry/api"; +import { afterEach, describe, expect, it } from "vitest"; +import { traceMixin } from "../../../apps/api/src/lib/log-trace-mixin.js"; + +/** + * Minimal AsyncLocalStorage-based context manager for tests. + * OTel v2 BasicTracerProvider no longer registers one automatically. + */ +class TestContextManager implements ContextManager { + private _als = new AsyncLocalStorage(); + + active(): Context { + return this._als.getStore() ?? ROOT_CONTEXT; + } + + with ReturnType>( + ctx: Context, + fn: F, + thisArg?: ThisParameterType, + ...args: A + ): ReturnType { + return this._als.run(ctx, () => fn.call(thisArg, ...args)); + } + + bind(_ctx: Context, target: T): T { + return target; + } + + enable(): this { + return this; + } + + disable(): this { + this._als.disable(); + return this; + } +} + +describe("traceMixin", () => { + afterEach(() => { + context.disable(); + trace.disable(); + }); + + it("returns empty object when no span is active", () => { + const result = traceMixin(); + expect(result).toEqual({}); + }); + + it("returns traceId, spanId, traceFlags when span is active", async () => { + const { InMemorySpanExporter, SimpleSpanProcessor, BasicTracerProvider } = await import( + "@opentelemetry/sdk-trace-base" + ); + + // Register a real context manager so context.with() propagates + context.setGlobalContextManager(new TestContextManager()); + + const exporter = new InMemorySpanExporter(); + const provider = new BasicTracerProvider({ + spanProcessors: [new SimpleSpanProcessor(exporter)], + }); + trace.setGlobalTracerProvider(provider); + + const tracer = trace.getTracer("test"); + const span = tracer.startSpan("test-span"); + const ctx = trace.setSpan(ROOT_CONTEXT, span); + + const result = context.with(ctx, () => traceMixin()); + + expect(result.traceId).toBe(span.spanContext().traceId); + expect(result.spanId).toBe(span.spanContext().spanId); + expect(result.traceFlags).toBe(span.spanContext().traceFlags); + + span.end(); + await provider.shutdown(); + }); +}); diff --git a/tests/unit/api/sidecar-context-propagation.test.ts b/tests/unit/api/sidecar-context-propagation.test.ts new file mode 100644 index 00000000..cb70c1f5 --- /dev/null +++ b/tests/unit/api/sidecar-context-propagation.test.ts @@ -0,0 +1,121 @@ +import { AsyncLocalStorage } from "node:async_hooks"; +import { + type Context, + type ContextManager, + context, + propagation, + ROOT_CONTEXT, + trace, +} from "@opentelemetry/api"; +import { W3CTraceContextPropagator } from "@opentelemetry/core"; +import { afterEach, describe, expect, it } from "vitest"; + +/** + * Minimal AsyncLocalStorage-based context manager for tests. + * OTel v2 BasicTracerProvider no longer registers one automatically. + */ +class TestContextManager implements ContextManager { + private _als = new AsyncLocalStorage(); + + active(): Context { + return this._als.getStore() ?? ROOT_CONTEXT; + } + + with ReturnType>( + ctx: Context, + fn: F, + thisArg?: ThisParameterType, + ...args: A + ): ReturnType { + return this._als.run(ctx, () => fn.call(thisArg, ...args)); + } + + bind(_ctx: Context, target: T): T { + return target; + } + + enable(): this { + return this; + } + + disable(): this { + this._als.disable(); + return this; + } +} + +describe("sidecar trace context injection", () => { + afterEach(() => { + propagation.disable(); + context.disable(); + trace.disable(); + }); + + it("produces _otel field when span is active", async () => { + const { InMemorySpanExporter, SimpleSpanProcessor, BasicTracerProvider } = await import( + "@opentelemetry/sdk-trace-base" + ); + + context.setGlobalContextManager(new TestContextManager()); + propagation.setGlobalPropagator(new W3CTraceContextPropagator()); + + const exporter = new InMemorySpanExporter(); + const provider = new BasicTracerProvider({ + spanProcessors: [new SimpleSpanProcessor(exporter)], + }); + trace.setGlobalTracerProvider(provider); + + const tracer = trace.getTracer("test"); + const span = tracer.startSpan("sidecar-call"); + const ctx = trace.setSpan(ROOT_CONTEXT, span); + + const message: Record = context.with(ctx, () => { + const carrier: Record = {}; + propagation.inject(context.active(), carrier); + + const msg: Record = { + id: "test-123", + script: "remove_bg", + args: ["input.png"], + }; + if (carrier.traceparent) { + msg._otel = { + traceparent: carrier.traceparent, + tracestate: carrier.tracestate, + }; + } + return msg; + }); + + // _otel should be present with a valid traceparent + expect(message._otel).toBeDefined(); + const otel = message._otel as { traceparent: string; tracestate?: string }; + expect(otel.traceparent).toMatch(/^00-[0-9a-f]{32}-[0-9a-f]{16}-0[01]$/); + + // traceparent should contain the correct traceId and spanId + const parts = otel.traceparent.split("-"); + expect(parts[1]).toBe(span.spanContext().traceId); + expect(parts[2]).toBe(span.spanContext().spanId); + + // args array is untouched + expect(message.args).toEqual(["input.png"]); + + span.end(); + await provider.shutdown(); + }); + + it("omits _otel when no SDK is registered", () => { + const carrier: Record = {}; + propagation.inject(context.active(), carrier); + + const message: Record = { + id: "test", + script: "remove_bg", + args: ["input.png"], + }; + if (carrier.traceparent) { + message._otel = { traceparent: carrier.traceparent }; + } + expect(message._otel).toBeUndefined(); + }); +}); diff --git a/tests/unit/api/tracing-bootstrap.test.ts b/tests/unit/api/tracing-bootstrap.test.ts new file mode 100644 index 00000000..4caac31e --- /dev/null +++ b/tests/unit/api/tracing-bootstrap.test.ts @@ -0,0 +1,46 @@ +import { afterEach, beforeEach, describe, expect, it, vi } from "vitest"; + +describe("tracing bootstrap", () => { + beforeEach(() => { + vi.unstubAllEnvs(); + }); + + afterEach(async () => { + const { shutdownTracing } = await import("../../../apps/api/src/tracing.js"); + await shutdownTracing(); + vi.resetModules(); + }); + + it("is inactive when OTEL_EXPORTER_OTLP_ENDPOINT is unset", async () => { + vi.stubEnv("OTEL_EXPORTER_OTLP_ENDPOINT", ""); + const { isTracingActive } = await import("../../../apps/api/src/tracing.js"); + expect(isTracingActive()).toBe(false); + }); + + it("is inactive when enterprise package is unavailable", async () => { + vi.stubEnv("OTEL_EXPORTER_OTLP_ENDPOINT", "http://localhost:4318"); + vi.doMock("@snapotter/enterprise", () => { + throw new Error("Cannot find module '@snapotter/enterprise'"); + }); + const { isTracingActive } = await import("../../../apps/api/src/tracing.js"); + expect(isTracingActive()).toBe(false); + }); + + it("initializes with valid config and enterprise license", async () => { + vi.stubEnv("OTEL_EXPORTER_OTLP_ENDPOINT", ""); + vi.doMock("@snapotter/enterprise", () => ({ + isFeatureEnabled: (f: string) => f === "distributed_tracing", + initEnterprise: () => true, + getActiveLicense: () => ({ + plan: "enterprise", + features: ["distributed_tracing"], + }), + })); + + const { InMemorySpanExporter } = await import("@opentelemetry/sdk-trace-base"); + const exporter = new InMemorySpanExporter(); + const { initTracing, isTracingActive } = await import("../../../apps/api/src/tracing.js"); + await initTracing({ exporter }); + expect(isTracingActive()).toBe(true); + }); +}); diff --git a/vitest.config.ts b/vitest.config.ts index e19c9286..ea614402 100644 --- a/vitest.config.ts +++ b/vitest.config.ts @@ -1,3 +1,4 @@ +import { readdirSync } from "node:fs"; import os from "node:os"; import path from "node:path"; import { defineConfig } from "vitest/config"; @@ -11,6 +12,20 @@ const webNodeModules = path.resolve(__dirname, "apps/web/node_modules"); // Resolve landing-workspace packages. const landingNodeModules = path.resolve(__dirname, "apps/landing/node_modules"); +// @opentelemetry/core is a transitive dep (via sdk-node) not hoisted by pnpm. +// Find it in the pnpm store so tests can import W3CTraceContextPropagator. +function findPnpmPackage(scope: string, name: string): string { + const pnpmDir = path.resolve(__dirname, "node_modules/.pnpm"); + const prefix = `${scope}+${name}@`; + const entries = readdirSync(pnpmDir) + .filter((e) => e.startsWith(prefix)) + .sort(); + if (entries.length === 0) { + throw new Error(`${scope}/${name} not found in pnpm store`); + } + return path.join(pnpmDir, entries[entries.length - 1], "node_modules", scope, name); +} + export default defineConfig({ esbuild: { jsx: "automatic", @@ -119,6 +134,7 @@ export default defineConfig({ qrcode: path.join(apiNodeModules, "qrcode"), jsqr: path.join(apiNodeModules, "jsqr"), pdfkit: path.join(apiNodeModules, "pdfkit"), + pino: path.join(apiNodeModules, "pino"), sharp: path.join(apiNodeModules, "sharp"), ioredis: path.join(apiNodeModules, "ioredis"), bullmq: path.join(apiNodeModules, "bullmq"), @@ -137,6 +153,39 @@ export default defineConfig({ zustand: path.join(webNodeModules, "zustand"), "posthog-js": path.join(webNodeModules, "posthog-js"), "@sentry/react": path.join(webNodeModules, "@sentry/react"), + "@opentelemetry/api": path.join(apiNodeModules, "@opentelemetry/api"), + "@opentelemetry/core": findPnpmPackage("@opentelemetry", "core"), + "@opentelemetry/sdk-node": path.join(apiNodeModules, "@opentelemetry/sdk-node"), + "@opentelemetry/sdk-trace-base": path.join(apiNodeModules, "@opentelemetry/sdk-trace-base"), + "@opentelemetry/exporter-trace-otlp-http": path.join( + apiNodeModules, + "@opentelemetry/exporter-trace-otlp-http", + ), + "@opentelemetry/resources": path.join(apiNodeModules, "@opentelemetry/resources"), + "@opentelemetry/semantic-conventions": path.join( + apiNodeModules, + "@opentelemetry/semantic-conventions", + ), + "@opentelemetry/instrumentation-http": path.join( + apiNodeModules, + "@opentelemetry/instrumentation-http", + ), + "@opentelemetry/instrumentation-fastify": path.join( + apiNodeModules, + "@opentelemetry/instrumentation-fastify", + ), + "@opentelemetry/instrumentation-pg": path.join( + apiNodeModules, + "@opentelemetry/instrumentation-pg", + ), + "@opentelemetry/instrumentation-ioredis": path.join( + apiNodeModules, + "@opentelemetry/instrumentation-ioredis", + ), + "@opentelemetry/instrumentation-aws-sdk": path.join( + apiNodeModules, + "@opentelemetry/instrumentation-aws-sdk", + ), }, }, });