apps/web/src/lib/aiResponseMonitor.server.ts
1import { waitUntil } from '@vercel/functions'
2import type { ReviewInput } from './aiConversationReview'
3import { recordConversationReview } from './aiConversationReview.server'
4export type ReviewCapture = {
5 input?: Omit<ReviewInput, 'response' | 'toolEvidence' | 'evidenceComplete'>
6 sessionId?: string
7 toolEvidence: string[]
8 evidenceComplete: boolean
9}
10/** Observe the exact returned text, including deterministic replies, without buffering delivery. */
11export function monitorChatResponse(
12 response: Response,
13 capture: ReviewCapture,
14 record: typeof recordConversationReview = recordConversationReview,
15): Response {
16 if (
17 !response.ok ||
18 !response.body ||
19 !capture.input ||
20 !capture.sessionId ||
21 process.env.AI_CONVERSATION_REVIEW === 'off'
22 )
23 return response
24 const input = capture.input,
25 sessionId = capture.sessionId
26 const reader = response.body.getReader(),
27 decoder = new TextDecoder()
28 let text = '',
29 done = false,
30 truncated = false
31 const append = (value: string) => {
32 if (text.length + value.length > 24000) truncated = true
33 text = (text + value).slice(0, 24000)
34 }
35 const finish = (incomplete: boolean) => {
36 if (done) return
37 done = true
38 const pending = record(
39 {
40 ...input,
41 response: text,
42 toolEvidence: capture.toolEvidence,
43 evidenceComplete: capture.evidenceComplete && !truncated,
44 },
45 sessionId,
46 incomplete || truncated,
47 )
48 try {
49 waitUntil(pending)
50 } catch {
51 void pending
52 }
53 }
54 const body = new ReadableStream<Uint8Array>({
55 async pull(controller) {
56 try {
57 const item = await reader.read()
58 if (item.done) {
59 append(decoder.decode())
60 finish(false)
61 controller.close()
62 } else {
63 append(decoder.decode(item.value, { stream: true }))
64 controller.enqueue(item.value)
65 }
66 } catch (error) {
67 finish(true)
68 controller.error(error)
69 }
70 },
71 async cancel(reason) {
72 finish(true)
73 await reader.cancel(reason)
74 },
75 })
76 return new Response(body, {
77 status: response.status,
78 statusText: response.statusText,
79 headers: response.headers,
80 })
81}
82