apps/web/src/lib/aiConversationReview.server.ts

1import { createGateway } from '@ai-sdk/gateway'
2import { propagateAttributes, startObservation } from '@langfuse/tracing'
3import { generateObject, jsonSchema } from 'ai'
4import {
5	type ConversationReview,
6	hasMalformedInternalLink,
7	type ReviewInput,
8	reviewFlags,
9	reviewNeedsAttention,
10	reviewRubric,
11} from './aiConversationReview'
12import { flushAiTelemetry } from './aiTelemetry'
13
14const schema = jsonSchema<ConversationReview>({
15	type: 'object',
16	additionalProperties: false,
17	required: ['topic', 'status', 'flags', 'reason'],
18	properties: {
19		topic: {
20			type: 'string',
21			enum: [
22				'tickets',
23				'spex',
24				'agenda',
25				'hotels',
26				'content',
27				'support',
28				'other',
29				'unknown',
30			],
31		},
32		status: {
33			type: 'string',
34			enum: ['no_detected_issue', 'needs_review', 'insufficient_evidence'],
35		},
36		flags: {
37			type: 'array',
38			items: { type: 'string', enum: [...reviewFlags] },
39			maxItems: 3,
40		},
41		reason: { type: 'string', maxLength: 700 },
42	},
43})
44export async function assessConversation(
45	input: ReviewInput,
46): Promise<ConversationReview> {
47	const apiKey = process.env.AI_GATEWAY_API_KEY || process.env.AI_GATEWAY
48	if (!apiKey) throw new Error('Review model unavailable')
49	const result = await propagateAttributes(
50		{
51			traceName: 'conversation-review-judge',
52			tags: ['evaluation', 'conversation-monitor'],
53			metadata: { purpose: 'quality-review' },
54		},
55		() =>
56			generateObject({
57				model: createGateway({ apiKey })(
58					process.env.AI_REVIEW_MODEL || 'openai/gpt-5.4',
59				),
60				schema,
61				system: reviewRubric,
62				prompt: JSON.stringify({
63					current_user_message: input.message,
64					saved_profile_data: input.profileContext,
65					assistant_response_to_evaluate: input.response,
66					earlier_conversation_for_context_only: input.history,
67					page_default_only: input.location,
68					tool_evidence_for_this_turn: input.toolEvidence,
69					evidence_complete: input.evidenceComplete,
70				}),
71				maxRetries: 0,
72				abortSignal: AbortSignal.timeout(30000),
73				providerOptions: {
74					gateway: { zeroDataRetention: true },
75					openai: { reasoningEffort: 'medium' },
76				},
77			}),
78	)
79	const review = result.object
80	if (review.status === 'needs_review' && review.flags.length === 0)
81		return { ...review, status: 'insufficient_evidence' }
82	if (review.flags.length > 0) return { ...review, status: 'needs_review' }
83	return review
84}
85export async function publishReviewScores(
86	target: { traceId: string; observationId: string },
87	review: ConversationReview | null,
88	malformed: boolean,
89	incomplete = false,
90) {
91	const key = process.env.LANGFUSE_PUBLIC_KEY,
92		secret = process.env.LANGFUSE_SECRET_KEY
93	if (!key || !secret) return
94	const status = incomplete
95		? 'incomplete'
96		: review
97			? reviewNeedsAttention(review, malformed)
98				? 'needs_review'
99				: review.status
100			: malformed
101				? 'needs_review'
102				: 'evaluator_unavailable'
103	const scores = [
104		{ name: 'conversation-review-status', value: status },
105		{ name: 'conversation-topic', value: review?.topic ?? 'unknown' },
106		{
107			name: 'conversation-confirmed-defect',
108			value: malformed ? 'malformed-internal-link' : 'none_detected',
109		},
110	]
111	const reason =
112		review?.reason.replace(/[\w.+-]+@[\w-]+\.[\w.-]+/g, '[email redacted]') ??
113		(incomplete
114			? 'Response was not completed; no quality verdict.'
115			: 'Semantic review unavailable; no quality verdict.')
116	const batch = scores.map((s) => ({
117		id: crypto.randomUUID(),
118		timestamp: new Date().toISOString(),
119		type: 'score-create',
120		body: {
121			id: `review-v1-${target.observationId}-${s.name}`,
122			traceId: target.traceId,
123			observationId: target.observationId,
124			name: s.name,
125			dataType: 'CATEGORICAL',
126			value: s.value,
127			comment:
128				s.name === 'conversation-confirmed-defect'
129					? malformed
130						? 'Protocol-relative internal path points to an invalid host.'
131						: 'No deterministic defect detected.'
132					: malformed
133						? `Malformed internal link: protocol-relative /events URL is rejected by the chat renderer. ${reason}`
134						: reason,
135			metadata: {
136				reviewVersion: '1',
137				flags: review?.flags ?? [],
138				semanticVerdict: review?.status ?? 'unavailable',
139			},
140		},
141	}))
142	const response = await fetch(
143		`${process.env.LANGFUSE_BASE_URL || 'https://cloud.langfuse.com'}/api/public/ingestion`,
144		{
145			method: 'POST',
146			headers: {
147				Authorization: `Basic ${Buffer.from(`${key}:${secret}`).toString('base64')}`,
148				'Content-Type': 'application/json',
149			},
150			body: JSON.stringify({ batch }),
151			signal: AbortSignal.timeout(5000),
152		},
153	)
154	if (!response.ok)
155		throw new Error(`Review score upload failed (${response.status})`)
156	const result = await response.json()
157	if (result.errors?.length) throw new Error('Review score ingestion rejected')
158}
159export async function recordConversationReview(
160	input: ReviewInput,
161	sessionId: string,
162	incomplete = false,
163) {
164	if (!process.env.LANGFUSE_PUBLIC_KEY || !process.env.LANGFUSE_SECRET_KEY)
165		return
166	const span = propagateAttributes(
167		{
168			sessionId,
169			traceName: 'conversation-turn',
170			tags: ['conversation-monitor'],
171			metadata: { reviewVersion: '1' },
172		},
173		() =>
174			startObservation(
175				'conversation-turn',
176				{ input: { ...input, response: undefined }, output: input.response },
177				{ asType: 'span' },
178			),
179	)
180	try {
181		const malformed = hasMalformedInternalLink(input.response)
182		let review: ConversationReview | null = null
183		if (!incomplete) {
184			try {
185				review = await assessConversation(input)
186			} catch {
187				span.update({
188					level: 'WARNING',
189					statusMessage: 'Semantic review unavailable',
190				})
191			}
192		}
193		span.update({
194			...(malformed || (review && reviewNeedsAttention(review, false))
195				? { level: 'WARNING' as const, statusMessage: 'Conversation needs review' }
196				: {}),
197			metadata: {
198				review: review ?? {
199					status: incomplete ? 'incomplete' : 'evaluator_unavailable',
200				},
201				malformedInternalLink: malformed,
202			},
203		})
204		await publishReviewScores(
205			{ traceId: span.traceId, observationId: span.id },
206			review,
207			malformed,
208			incomplete,
209		)
210	} catch {
211		console.warn('[ai-review] Review scoring failed')
212	} finally {
213		span.end()
214		await flushAiTelemetry()
215	}
216}
217