apps/web/src/lib/aiTelemetry.ts

1import { LangfuseSpanProcessor } from '@langfuse/otel'
2import { LangfuseVercelAiSdkIntegration } from '@langfuse/vercel-ai-sdk'
3import { NodeSDK } from '@opentelemetry/sdk-node'
4import { registerTelemetry } from 'ai'
5
6/**
7 * Langfuse tracing for the AI routes (AI SDK 7 + OpenTelemetry).
8 * Once registered, every generateText/streamText/generateObject call emits
9 * spans (model, tokens, tool executions) that the LangfuseSpanProcessor ships
10 * to Langfuse. Attribution (trace name, session, user, tags) is added per
11 * request with propagateAttributes() at the call sites.
12 *
13 * No-ops with a warning when LANGFUSE_* env vars are absent.
14 */
15
16const LOG_PREFIX = '[aiTelemetry]'
17
18let processor: LangfuseSpanProcessor | undefined
19let initialized = false
20
21const EMAIL_PATTERN = /[\w.+-]+@[\w-]+\.[\w.-]+/g
22const PHONE_PATTERN = /(?<![\d/])(\+?\d[\d\s().-]{7,}\d)(?![\d/])/g
23
24export function maskPiiText(text: string): string {
25	return text
26		.replace(EMAIL_PATTERN, (email) =>
27			email.toLowerCase().endsWith('@unleash.ai') ? email : '[email redacted]',
28		)
29		.replace(PHONE_PATTERN, '[phone redacted]')
30}
31
32export function redactSpanAttributes(
33	span: Pick<Parameters<LangfuseSpanProcessor['onEnd']>[0], 'attributes'>,
34): void {
35	for (const [key, value] of Object.entries(span.attributes)) {
36		if (typeof value === 'string') span.attributes[key] = maskPiiText(value)
37		else if (Array.isArray(value)) {
38			for (const [index, entry] of value.entries())
39				if (typeof entry === 'string') value[index] = maskPiiText(entry)
40		}
41	}
42}
43
44/**
45 * PII redaction in the logging layer (SoW §5.3): non-Unleash email addresses
46 * and phone numbers are masked in every exported span. Public @unleash.ai
47 * contact addresses stay readable.
48 */
49export function maskPii(value: unknown): unknown {
50	if (typeof value === 'string') {
51		return maskPiiText(value)
52	}
53	if (Array.isArray(value)) return value.map(maskPii)
54	if (value && typeof value === 'object') {
55		return Object.fromEntries(
56			Object.entries(value as Record<string, unknown>).map(([key, entry]) => [
57				key,
58				maskPii(entry),
59			]),
60		)
61	}
62	return value
63}
64
65/** Idempotent; call at the top of AI route handlers (before model calls). */
66export function initAiTelemetry(): void {
67	if (initialized) return
68	initialized = true
69
70	if (!process.env.LANGFUSE_SECRET_KEY || !process.env.LANGFUSE_PUBLIC_KEY) {
71		console.warn(`${LOG_PREFIX} LANGFUSE keys not set; tracing disabled`)
72		return
73	}
74
75	try {
76		processor = new LangfuseSpanProcessor({
77			mask: ({ data }) => maskPii(data),
78		})
79		// SDK 7 records tool arguments/messages in gen_ai.* attributes. Langfuse's
80		// mask callback currently covers only langfuse.* input/output attributes.
81		// Redact every attribute before either schema reaches the exporter.
82		const exporter = processor
83		const sdk = new NodeSDK({
84			spanProcessors: [
85				{
86					onStart: (span, context) => exporter.onStart(span, context),
87					onEnd: (span) => {
88						redactSpanAttributes(span)
89						exporter.onEnd(span)
90					},
91					forceFlush: () => exporter.forceFlush(),
92					shutdown: () => exporter.shutdown(),
93				},
94			],
95		})
96		sdk.start()
97		registerTelemetry(new LangfuseVercelAiSdkIntegration())
98	} catch (error) {
99		console.error(`${LOG_PREFIX} init failed; tracing disabled`, error)
100		processor = undefined
101	}
102}
103
104/**
105 * Ship buffered spans before the serverless function can freeze. Call after
106 * a response finishes streaming (fire-and-forget via waitUntil).
107 */
108export async function flushAiTelemetry(): Promise<void> {
109	try {
110		await processor?.forceFlush()
111	} catch (error) {
112		console.error(`${LOG_PREFIX} flush failed`, error)
113	}
114}
115