apps/web/src/lib/chatDelivery.ts

1import { HANDOFF_FAILURE_REPLY } from '@local/ai/components/workflows'
2import { type HandoffOutcome, verifiedHandoffReply } from './aiHandoff'
3
4type DeliveryPart = { type: string; text?: string }
5
6/** Pass generated text through immediately; tool-owned replies are emitted at tool completion. */
7export function createChatDelivery(options: {
8	handoffRequested: boolean
9	outcome: () => HandoffOutcome | undefined
10	plan: () => string | undefined
11	onComplete?: (text: string) => void
12	suffix?: (text: string) => string
13}) {
14	let text = ''
15	let authoritativeReplyEmitted = false
16	let stepBoundary = false
17	const encoder = new TextEncoder()
18	const stream = new TransformStream<DeliveryPart, Uint8Array>({
19		transform(part, controller) {
20			if (part.type === 'start-step' && text) stepBoundary = true
21			if (part.type === 'error') {
22				if (options.handoffRequested) {
23					if (!authoritativeReplyEmitted) {
24						text = HANDOFF_FAILURE_REPLY
25						controller.enqueue(encoder.encode(text))
26						authoritativeReplyEmitted = true
27					}
28				} else controller.error(new Error('AI response stream failed'))
29				return
30			}
31			const outcome = options.outcome()
32			const plan = options.plan()
33			const emit = (value: string) => {
34				text += value
35				controller.enqueue(encoder.encode(value))
36			}
37			if ((outcome || plan) && !authoritativeReplyEmitted) {
38				emit(`${text ? '\n\n' : ''}${plan || verifiedHandoffReply('', outcome)}`)
39				authoritativeReplyEmitted = true
40			}
41			if (authoritativeReplyEmitted || options.handoffRequested) return
42			if (part.type === 'text-delta' && part.text) {
43				if (stepBoundary) {
44					emit('\n\n')
45					stepBoundary = false
46				}
47				emit(part.text)
48			}
49		},
50		flush(controller) {
51			if (options.handoffRequested && !authoritativeReplyEmitted) {
52				text = HANDOFF_FAILURE_REPLY
53				controller.enqueue(encoder.encode(text))
54			}
55			if (!authoritativeReplyEmitted && !options.handoffRequested) {
56				const suffix = options.suffix?.(text) || ''
57				text += suffix
58				if (suffix) controller.enqueue(encoder.encode(suffix))
59			}
60			options.onComplete?.(text)
61		},
62	})
63	return { stream, text: () => text }
64}
65