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