Add external log exporters and fix missing external trace exporters in deployed tasks (#2038)

* Add external log exporters and fix missing external trace exporters in deployed tasks

* Generate the external traceID correctly and exporter 3rd party logs with the external traceID as well
This commit is contained in:
Eric Allam
2025-05-13 16:54:24 +01:00
committed by GitHub
parent 84a0336bdd
commit a815633712
10 changed files with 143 additions and 7 deletions
+5
View File
@@ -0,0 +1,5 @@
---
"trigger.dev": patch
---
Add external log exporters and fix missing external trace exporters in deployed tasks
@@ -164,6 +164,7 @@ async function bootstrap() {
url: env.OTEL_EXPORTER_OTLP_ENDPOINT ?? "http://0.0.0.0:4318",
instrumentations: config.telemetry?.instrumentations ?? config.instrumentations ?? [],
exporters: config.telemetry?.exporters ?? [],
logExporters: config.telemetry?.logExporters ?? [],
diagLogLevel: (env.OTEL_LOG_LEVEL as TracingDiagnosticLogLevel) ?? "none",
forceFlushTimeoutMillis: 30_000,
});
@@ -174,7 +175,8 @@ async function bootstrap() {
const tracer = new TriggerTracer({ tracer: otelTracer, logger: otelLogger });
const consoleInterceptor = new ConsoleInterceptor(
otelLogger,
typeof config.enableConsoleLogging === "boolean" ? config.enableConsoleLogging : true
typeof config.enableConsoleLogging === "boolean" ? config.enableConsoleLogging : true,
typeof config.disableConsoleInterceptor === "boolean" ? config.disableConsoleInterceptor : false
);
const configLogLevel = triggerLogLevel ?? config.logLevel ?? "info";
@@ -163,6 +163,8 @@ async function bootstrap() {
instrumentations: config.instrumentations ?? [],
diagLogLevel: (env.OTEL_LOG_LEVEL as TracingDiagnosticLogLevel) ?? "none",
forceFlushTimeoutMillis: 30_000,
exporters: config.telemetry?.exporters ?? [],
logExporters: config.telemetry?.logExporters ?? [],
});
const otelTracer: Tracer = tracingSDK.getTracer("trigger-dev-worker", VERSION);
@@ -171,7 +173,8 @@ async function bootstrap() {
const tracer = new TriggerTracer({ tracer: otelTracer, logger: otelLogger });
const consoleInterceptor = new ConsoleInterceptor(
otelLogger,
typeof config.enableConsoleLogging === "boolean" ? config.enableConsoleLogging : true
typeof config.enableConsoleLogging === "boolean" ? config.enableConsoleLogging : true,
typeof config.disableConsoleInterceptor === "boolean" ? config.disableConsoleInterceptor : false
);
const configLogLevel = triggerLogLevel ?? config.logLevel ?? "info";
+13
View File
@@ -11,6 +11,7 @@ import type {
} from "./index.js";
import type { LogLevel } from "./logger/taskLogger.js";
import type { MachinePresetName } from "./schemas/common.js";
import { LogRecordExporter } from "@opentelemetry/sdk-logs";
export type CompatibilityFlag = "run_engine_v2";
@@ -80,6 +81,13 @@ export type TriggerConfig = {
* @see https://trigger.dev/docs/config/config-file#exporters
*/
exporters?: Array<SpanExporter>;
/**
* Log exporters to use for OpenTelemetry. This is useful if you want to add custom log exporters to your tasks.
*
* @see https://trigger.dev/docs/config/config-file#exporters
*/
logExporters?: Array<LogRecordExporter>;
};
/**
@@ -131,6 +139,11 @@ export type TriggerConfig = {
*/
enableConsoleLogging?: boolean;
/**
* Disable the console interceptor. This will prevent logs from being sent to the trigger.dev backend.
*/
disableConsoleInterceptor?: boolean;
build?: {
/**
* Add custom conditions to the esbuild build. For example, if you are importing `ai/rsc`, you'll need to add "react-server" condition.
+6 -1
View File
@@ -10,12 +10,17 @@ import { clock } from "./clock-api.js";
export class ConsoleInterceptor {
constructor(
private readonly logger: logsAPI.Logger,
private readonly sendToStdIO: boolean
private readonly sendToStdIO: boolean,
private readonly interceptingDisabled: boolean
) {}
// Intercept the console and send logs to the OpenTelemetry logger
// during the execution of the callback
async intercept<T>(console: Console, callback: () => Promise<T>): Promise<T> {
if (this.interceptingDisabled) {
return await callback();
}
// Save the original console methods
const originalConsole = {
log: console.log,
+76 -1
View File
@@ -1,4 +1,5 @@
import { DiagConsoleLogger, DiagLogLevel, TracerProvider, diag } from "@opentelemetry/api";
import { RandomIdGenerator } from "@opentelemetry/sdk-trace-base";
import { logs } from "@opentelemetry/api-logs";
import { OTLPLogExporter } from "@opentelemetry/exporter-logs-otlp-http";
import { OTLPTraceExporter } from "@opentelemetry/exporter-trace-otlp-http";
@@ -15,6 +16,8 @@ import {
import {
BatchLogRecordProcessor,
LoggerProvider,
LogRecordExporter,
ReadableLogRecord,
SimpleLogRecordProcessor,
} from "@opentelemetry/sdk-logs";
import {
@@ -87,9 +90,12 @@ export type TracingSDKConfig = {
resource?: IResource;
instrumentations?: Instrumentation[];
exporters?: SpanExporter[];
logExporters?: LogRecordExporter[];
diagLogLevel?: TracingDiagnosticLogLevel;
};
const idGenerator = new RandomIdGenerator();
export class TracingSDK {
public readonly asyncResourceDetector = new AsyncResourceDetector();
private readonly _logProvider: LoggerProvider;
@@ -158,7 +164,7 @@ export class TracingSDK {
)
);
const externalTraceId = crypto.randomUUID();
const externalTraceId = idGenerator.generateTraceId();
for (const exporter of config.exporters ?? []) {
traceProvider.addSpanProcessor(
@@ -210,6 +216,28 @@ export class TracingSDK {
)
);
for (const externalLogExporter of config.logExporters ?? []) {
loggerProvider.addLogRecordProcessor(
getEnvVar("OTEL_BATCH_PROCESSING_ENABLED") === "1"
? new BatchLogRecordProcessor(
new ExternalLogRecordExporterWrapper(externalLogExporter, externalTraceId),
{
maxExportBatchSize: parseInt(getEnvVar("OTEL_LOG_MAX_EXPORT_BATCH_SIZE") ?? "64"),
scheduledDelayMillis: parseInt(
getEnvVar("OTEL_LOG_SCHEDULED_DELAY_MILLIS") ?? "200"
),
exportTimeoutMillis: parseInt(
getEnvVar("OTEL_LOG_EXPORT_TIMEOUT_MILLIS") ?? "30000"
),
maxQueueSize: parseInt(getEnvVar("OTEL_LOG_MAX_QUEUE_SIZE") ?? "512"),
}
)
: new SimpleLogRecordProcessor(
new ExternalLogRecordExporterWrapper(externalLogExporter, externalTraceId)
)
);
}
this._logProvider = loggerProvider;
this._spanExporter = spanExporter;
this._traceProvider = traceProvider;
@@ -306,3 +334,50 @@ class ExternalSpanExporterWrapper {
: Promise.resolve();
}
}
class ExternalLogRecordExporterWrapper {
constructor(
private underlyingExporter: LogRecordExporter,
private externalTraceId: string
) {}
export(logs: any[], resultCallback: (result: any) => void): void {
const modifiedLogs = logs.map(this.transformLogRecord.bind(this));
this.underlyingExporter.export(modifiedLogs, resultCallback);
}
shutdown(): Promise<void> {
return this.underlyingExporter.shutdown();
}
transformLogRecord(logRecord: ReadableLogRecord): ReadableLogRecord {
// If there's no spanContext, or if the externalTraceId is not set, return the original logRecord.
if (!logRecord.spanContext || !this.externalTraceId) {
return logRecord;
}
// Capture externalTraceId for use within the proxy's scope.
const { externalTraceId } = this;
return new Proxy(logRecord, {
get(target, prop, receiver) {
if (prop === "spanContext") {
// Intercept access to spanContext.
const originalSpanContext = target.spanContext;
// Ensure originalSpanContext exists (it should, due to the check above, but good for safety).
if (originalSpanContext) {
return {
...originalSpanContext,
traceId: externalTraceId, // Override traceId.
};
}
// Fallback if, for some reason, originalSpanContext is undefined here.
return undefined;
}
// For all other properties, defer to the original object.
return Reflect.get(target, prop, receiver);
},
});
}
}
+10 -3
View File
@@ -1828,6 +1828,12 @@ importers:
'@e2b/code-interpreter':
specifier: ^1.1.0
version: 1.1.0
'@opentelemetry/exporter-logs-otlp-http':
specifier: 0.52.1
version: 0.52.1(@opentelemetry/api@1.9.0)
'@opentelemetry/exporter-trace-otlp-http':
specifier: 0.52.1
version: 0.52.1(@opentelemetry/api@1.9.0)
'@radix-ui/react-avatar':
specifier: ^1.1.3
version: 1.1.3(@types/react-dom@19.0.4)(@types/react@19.0.12)(react-dom@19.0.0)(react@19.0.0)
@@ -1866,7 +1872,7 @@ importers:
version: 5.1.5
next:
specifier: 15.2.4
version: 15.2.4(@playwright/test@1.37.0)(react-dom@19.0.0)(react@19.0.0)
version: 15.2.4(@opentelemetry/api@1.9.0)(@playwright/test@1.37.0)(react-dom@19.0.0)(react@19.0.0)
react:
specifier: ^19.0.0
version: 19.0.0
@@ -1954,7 +1960,7 @@ importers:
version: 5.1.5
next:
specifier: 15.2.4
version: 15.2.4(@playwright/test@1.37.0)(react-dom@19.0.0)(react@19.0.0)
version: 15.2.4(@opentelemetry/api@1.9.0)(@playwright/test@1.37.0)(react-dom@19.0.0)(react@19.0.0)
react:
specifier: ^19.0.0
version: 19.0.0
@@ -28612,7 +28618,7 @@ packages:
- babel-plugin-macros
dev: false
/next@15.2.4(@playwright/test@1.37.0)(react-dom@19.0.0)(react@19.0.0):
/next@15.2.4(@opentelemetry/api@1.9.0)(@playwright/test@1.37.0)(react-dom@19.0.0)(react@19.0.0):
resolution: {integrity: sha512-VwL+LAaPSxEkd3lU2xWbgEOtrM8oedmyhBqaVNmgKB+GvZlCy9rgaEc+y2on0wv+l0oSFqLtYD6dcC1eAedUaQ==}
engines: {node: ^18.18.0 || ^19.8.0 || >= 20.0.0}
hasBin: true
@@ -28634,6 +28640,7 @@ packages:
optional: true
dependencies:
'@next/env': 15.2.4
'@opentelemetry/api': 1.9.0
'@playwright/test': 1.37.0
'@swc/counter': 0.1.3
'@swc/helpers': 0.5.15
+2
View File
@@ -27,6 +27,8 @@
"@trigger.dev/python": "workspace:*",
"@trigger.dev/react-hooks": "workspace:*",
"@trigger.dev/sdk": "workspace:*",
"@opentelemetry/exporter-logs-otlp-http": "0.52.1",
"@opentelemetry/exporter-trace-otlp-http": "0.52.1",
"@vercel/postgres": "^0.10.0",
"ai": "4.2.5",
"class-variance-authority": "^0.7.1",
+2
View File
@@ -227,6 +227,8 @@ export const interruptibleChat = schemaTask({
prompt: z.string().describe("The prompt to chat with the AI"),
}),
run: async ({ prompt }, { signal }) => {
logger.info("interruptible-chat: starting", { prompt });
const chunks: TextStreamPart<{}>[] = [];
// 👇 This is a global onCancel hook, but it's inside of the run function
+22
View File
@@ -1,10 +1,32 @@
import { defineConfig } from "@trigger.dev/sdk";
import { pythonExtension } from "@trigger.dev/python/extension";
import { installPlaywrightChromium } from "./src/extensions/playwright";
import { OTLPLogExporter } from "@opentelemetry/exporter-logs-otlp-http";
import { OTLPTraceExporter } from "@opentelemetry/exporter-trace-otlp-http";
export default defineConfig({
project: "proj_cdmymsrobxmcgjqzhdkq",
dirs: ["./src/trigger"],
telemetry: {
logExporters: [
new OTLPLogExporter({
url: "https://api.axiom.co/v1/logs",
headers: {
Authorization: `Bearer ${process.env.AXIOM_TOKEN}`,
"X-Axiom-Dataset": "d3-chat-tester",
},
}),
],
exporters: [
new OTLPTraceExporter({
url: "https://api.axiom.co/v1/traces",
headers: {
Authorization: `Bearer ${process.env.AXIOM_TOKEN}`,
"X-Axiom-Dataset": "d3-chat-tester",
},
}),
],
},
maxDuration: 3600,
build: {
extensions: [