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