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