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