Skip to content

Commit 9ef4c7e

Browse files
committed
fix Langfuse trace correlation ids
1 parent 0c5f98e commit 9ef4c7e

9 files changed

Lines changed: 49 additions & 18 deletions

apps/daemon/src/langfuse-trace.ts

Lines changed: 7 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -55,6 +55,10 @@ const DEFAULT_FETCH_RETRIES = 1;
5555
const PROMPT_STACK_BLAME_MAX_SECTIONS = 8;
5656
let missingTelemetrySinkWarned = false;
5757

58+
export function langfuseTraceIdForRun(runId: string): string {
59+
return createHash('sha256').update(runId).digest('hex').slice(0, 32);
60+
}
61+
5862
export interface LangfuseConfig {
5963
authHeader: string;
6064
baseUrl: string;
@@ -1540,7 +1544,7 @@ export function buildTracePayload(ctx: ReportContext): unknown[] {
15401544
const performanceDiagnostics = buildPerformanceDiagnostics(ctx);
15411545

15421546
const success = ctx.run.status === 'succeeded';
1543-
const traceId = ctx.run.runId;
1547+
const traceId = langfuseTraceIdForRun(ctx.run.runId);
15441548
const langfuseDelivery =
15451549
ctx.langfuse ??
15461550
deriveLangfuseDeliveryState(ctx.prefs, readRunTelemetrySinkConfig());
@@ -1577,6 +1581,7 @@ export function buildTracePayload(ctx: ReportContext): unknown[] {
15771581
status: ctx.run.status,
15781582
error: safeRunError,
15791583
error_code: ctx.run.errorCode,
1584+
run_id: ctx.run.runId,
15801585
langfuse_trace_id: traceId,
15811586
...langfuseDelivery,
15821587
...(ctx.run.failure ?? {}),
@@ -2478,7 +2483,7 @@ export async function reportRunCompleted(
24782483
// thread `removedReasonCodes` through and emit overwriting "cleared"
24792484
// scores for them; not done here to keep this PR scoped to the bridge.
24802485
export function buildFeedbackPayload(ctx: FeedbackReportContext): unknown[] {
2481-
const traceId = ctx.runId;
2486+
const traceId = langfuseTraceIdForRun(ctx.runId);
24822487
const nowIso = new Date().toISOString();
24832488
const batch: unknown[] = [];
24842489

apps/daemon/src/routes/runs.ts

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -45,6 +45,7 @@ import { readVelaLoginStatus } from '../integrations/vela.js';
4545
import { getDetectedRuntimeVersions } from '../runtimes/detection.js';
4646
import {
4747
deriveLangfuseDeliveryState,
48+
langfuseTraceIdForRun,
4849
readTelemetrySinkConfig,
4950
} from '../langfuse-trace.js';
5051
import { parseMediaExecutionPolicyInput } from '../media/policy.js';
@@ -1979,7 +1980,7 @@ export function registerRunRoutes(app: Express, ctx: RegisterRunRoutesDeps) {
19791980
...(typeof run.lastAgentActivityAt === 'number'
19801981
? { last_progress_age_ms: Math.max(0, analyticsCapturedAt - run.lastAgentActivityAt) }
19811982
: {}),
1982-
langfuse_trace_id: run.id,
1983+
langfuse_trace_id: langfuseTraceIdForRun(run.id),
19831984
...langfuseDeliveryForAnalytics,
19841985
...(errorCode ? { error_code: errorCode } : {}),
19851986
...(failure ?? {}),

apps/daemon/src/runtimes/run-terminal-reconciliation.ts

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -5,6 +5,7 @@ import { randomUUID } from 'node:crypto';
55
import type Database from 'better-sqlite3';
66

77
import { appendMessageStatusEvent } from '../db.js';
8+
import { langfuseTraceIdForRun } from '../langfuse-trace.js';
89
import { classifyRunFailure } from '../run-failure-classification.js';
910
import { deriveRunErrorCode, runResultFromStatus } from '../run-result.js';
1011
import { runAskedUserQuestion } from './run-artifacts.js';
@@ -281,7 +282,7 @@ export async function reconcileDurableRunTerminals(
281282
artifact_count: state.artifactCount ?? 0,
282283
asked_user_question: runAskedUserQuestion(events),
283284
total_duration_ms: Math.max(0, state.updatedAt - state.createdAt),
284-
langfuse_trace_id: state.id,
285+
langfuse_trace_id: langfuseTraceIdForRun(state.id),
285286
terminal_reconciled: true,
286287
terminal_recovery_reason: recoveryReason,
287288
...(errorCode ? { error_code: errorCode } : {}),

apps/daemon/src/server.ts

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -411,6 +411,7 @@ import {
411411
snapshotAiHtmlVersionsForRun,
412412
} from './run-html-version-snapshots.js';
413413
import { reportRunCompletedFromDaemon } from './langfuse-bridge.js';
414+
import { langfuseTraceIdForRun } from './langfuse-trace.js';
414415
import { reconcileDurableRunTerminals } from './runtimes/run-terminal-reconciliation.js';
415416
import { buildPromptStackTelemetry } from './prompt-telemetry.js';
416417
import { readAnalyticsContext } from './analytics.js';
@@ -1540,7 +1541,7 @@ export function createFinalizedMessageTelemetryReporter({
15401541
project_id: run?.projectId ?? projectId ?? null,
15411542
conversation_id: run?.conversationId ?? conversationId ?? null,
15421543
run_id: runId,
1543-
langfuse_trace_id: runId,
1544+
langfuse_trace_id: langfuseTraceIdForRun(runId),
15441545
langfuse_expected: delivery.langfuse_expected,
15451546
langfuse_delivery_status: delivery.langfuse_delivery_status,
15461547
...(delivery.langfuse_drop_reason

apps/daemon/tests/langfuse-trace.test.ts

Lines changed: 27 additions & 9 deletions
Original file line numberDiff line numberDiff line change
@@ -3,6 +3,7 @@ import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest';
33
import {
44
buildFeedbackPayload,
55
buildTracePayload,
6+
langfuseTraceIdForRun,
67
deriveLangfuseDeliveryState,
78
isContentToolName,
89
isPartialRedactToolName,
@@ -371,9 +372,18 @@ describe('shouldFullyRedactToolPayload (fail-closed)', () => {
371372
});
372373
});
373374

375+
describe('langfuseTraceIdForRun', () => {
376+
it('derives the Langfuse-compliant trace id deterministically from the Open Design run id', () => {
377+
expect(langfuseTraceIdForRun('run-1')).toBe(
378+
'4e65d3fbe8ad6535681b021b30785b12',
379+
);
380+
});
381+
});
382+
374383
describe('buildTracePayload', () => {
375384
it('emits a trace with nested agent + generation observations', () => {
376385
const batch = buildTracePayload(makeCtx());
386+
const traceId = langfuseTraceIdForRun('run-1');
377387
const types = (batch as Array<{ type: string }>).map((e) => e.type);
378388
expect(types).toEqual([
379389
'trace-create',
@@ -386,9 +396,13 @@ describe('buildTracePayload', () => {
386396
const gen = bodyOf(batch, 'generation-create', 'llm');
387397
const bash = bodyOf(batch, 'span-create', 'tool:Bash');
388398
const write = bodyOf(batch, 'span-create', 'tool:Write');
399+
const trace = (batch[0] as any).body;
400+
expect(trace.id).toBe(traceId);
401+
expect(trace.metadata.run_id).toBe('run-1');
402+
expect(trace.metadata.langfuse_trace_id).toBe(traceId);
389403
expect(span.id).toBe('run-1-agent');
390-
expect(span.traceId).toBe('run-1');
391-
expect(gen.traceId).toBe('run-1');
404+
expect(span.traceId).toBe(traceId);
405+
expect(gen.traceId).toBe(traceId);
392406
expect(gen.parentObservationId).toBe('run-1-agent');
393407
expect(bash.parentObservationId).toBe('run-1-agent');
394408
expect(bash.input).toBeUndefined();
@@ -1278,7 +1292,10 @@ describe('buildTracePayload', () => {
12781292
);
12791293
const metadata = (batch[0] as any).body.metadata;
12801294
expect(metadata.error_code).toBe('RATE_LIMITED');
1281-
expect(metadata.langfuse_trace_id).toBe('run-rate-limit');
1295+
expect(metadata.run_id).toBe('run-rate-limit');
1296+
expect(metadata.langfuse_trace_id).toBe(
1297+
langfuseTraceIdForRun('run-rate-limit'),
1298+
);
12821299
expect(metadata.langfuse_expected).toBe(false);
12831300
expect(metadata.langfuse_delivery_status).toBe('not_expected');
12841301
expect(metadata.langfuse_drop_reason).toBe('content_consent_off');
@@ -2528,8 +2545,9 @@ describe('buildFeedbackPayload', () => {
25282545
) as Array<Record<string, any>>;
25292546
expect(batch).toHaveLength(3);
25302547
const ratingScore = batch[0]!;
2548+
const traceId = langfuseTraceIdForRun('run-feedback-1');
25312549
expect(ratingScore.type).toBe('score-create');
2532-
expect(ratingScore.body.traceId).toBe('run-feedback-1');
2550+
expect(ratingScore.body.traceId).toBe(traceId);
25332551
expect(ratingScore.body.name).toBe('user_rating');
25342552
expect(ratingScore.body.value).toBe(-1);
25352553
expect(ratingScore.body.dataType).toBe('NUMERIC');
@@ -2543,7 +2561,7 @@ describe('buildFeedbackPayload', () => {
25432561
expect(reasonScore.body.name).toBe('user_rating_reason');
25442562
expect(reasonScore.body.dataType).toBe('CATEGORICAL');
25452563
expect(reasonScore.body.comment).toBe('negative');
2546-
expect(reasonScore.body.traceId).toBe('run-feedback-1');
2564+
expect(reasonScore.body.traceId).toBe(traceId);
25472565
}
25482566
expect(batch[1]!.body.value).toBe('missed_request');
25492567
expect(batch[2]!.body.value).toBe('weak_visual');
@@ -2611,14 +2629,14 @@ describe('reportRunFeedback', () => {
26112629
'score',
26122630
]);
26132631
expect(envelope.events[0].data).toMatchObject({
2614-
id: 'run-feedback-1-rating',
2615-
traceId: 'run-feedback-1',
2632+
id: '908d98293a13a7b7811f44a271aa3b5b-rating',
2633+
traceId: '908d98293a13a7b7811f44a271aa3b5b',
26162634
name: 'user_rating',
26172635
value: 1,
26182636
});
26192637
expect(envelope.events[1].data).toMatchObject({
2620-
id: 'run-feedback-1-reason-matched_request',
2621-
traceId: 'run-feedback-1',
2638+
id: '908d98293a13a7b7811f44a271aa3b5b-reason-matched_request',
2639+
traceId: '908d98293a13a7b7811f44a271aa3b5b',
26222640
name: 'user_rating_reason',
26232641
value: 'matched_request',
26242642
});

apps/daemon/tests/runtimes/run-terminal-reconciliation.test.ts

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -106,6 +106,7 @@ describe('durable run terminal reconciliation', () => {
106106
failure_stage: 'finalize',
107107
retryable: true,
108108
user_action: 'retry',
109+
langfuse_trace_id: '109f933ded33d54c42a3faba33818f1e',
109110
terminal_reconciled: true,
110111
terminal_recovery_reason: 'daemon_restart',
111112
}),

apps/daemon/tests/telemetry-message-finalization.test.ts

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -453,7 +453,7 @@ describe('Langfuse message finalization gate', () => {
453453
insertId: 'run-accepted-langfuse-report-terminal_fallback-accepted',
454454
properties: expect.objectContaining({
455455
run_id: 'run-accepted',
456-
langfuse_trace_id: 'run-accepted',
456+
langfuse_trace_id: '7d63aa20a8824b5b21fda43b3a7cfbc9',
457457
langfuse_expected: true,
458458
langfuse_delivery_status: 'accepted',
459459
langfuse_report_result: 'accepted',

specs/current/ai-native-observability-loop.md

Lines changed: 3 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -17,7 +17,9 @@ Issue: [#3713](https://github.com/nexu-io/open-design/issues/3713)
1717
Open Design already reports completed runs to Langfuse and PostHog. The current
1818
implementation captures useful operational facts:
1919

20-
- trace identity: `run_id == langfuse_trace_id == traceId`;
20+
- trace identity: `run_id` is the Open Design external identifier, while
21+
`langfuse_trace_id == traceId == sha256(run_id).slice(0, 32)`; Langfuse
22+
metadata retains `run_id` for external correlation;
2123
- run status, error code, failure category, failure detail, retryability, and
2224
user action;
2325
- timing fields for queue, prompt build, spawn, first token, generation, tool

specs/current/run-reliability-optimization-plan.md

Lines changed: 4 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -55,8 +55,10 @@ observability slice:
5555
`failure_category`, `failure_detail`, `failure_stage`, `retryable`, and
5656
`user_action`;
5757
- failed runs preserve the non-empty `error_code` invariant;
58-
- `langfuse_trace_id` makes PostHog and Langfuse correlate by
59-
`run_id == langfuse_trace_id == traceId`;
58+
- `langfuse_trace_id` makes PostHog and Langfuse correlate by storing the
59+
actual Langfuse trace ID, deterministically derived as
60+
`sha256(run_id).slice(0, 32) == traceId`; Langfuse metadata retains `run_id`
61+
for the reverse join;
6062
- Langfuse trace metadata mirrors failure, timing, and token/cache fields;
6163
- `run_finished` emits main-path timing fields for queue, spawn, first token,
6264
generation, tool aggregate, finalization, and total duration;

0 commit comments

Comments
 (0)