-
Notifications
You must be signed in to change notification settings - Fork 1.1k
Expand file tree
/
Copy pathtool-stream-parser.ts
More file actions
195 lines (177 loc) · 5.04 KB
/
Copy pathtool-stream-parser.ts
File metadata and controls
195 lines (177 loc) · 5.04 KB
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
import { AnalyticsEvent } from '@codebuff/common/constants/analytics-events'
import {
createStreamParserState,
parseStreamChunk,
} from './util/stream-xml-parser'
import type { StreamParserState } from './util/stream-xml-parser'
import type { Model } from '@codebuff/common/old-constants'
import type { TrackEventFn } from '@codebuff/common/types/contracts/analytics'
import type { StreamChunk } from '@codebuff/common/types/contracts/llm'
import type { Logger } from '@codebuff/common/types/contracts/logger'
import type {
PrintModeError,
PrintModeText,
} from '@codebuff/common/types/print-mode'
import type { PromptResult } from '@codebuff/common/util/error'
export async function* processStreamWithTools(params: {
stream: AsyncGenerator<StreamChunk, PromptResult<string | null>>
processors: Record<
string,
{
onTagStart: (
tagName: string,
attributes: Record<string, string>,
) => void | Promise<void>
onTagEnd: (
tagName: string,
params: Record<string, any>,
) => void | Promise<void>
}
>
defaultProcessor: (toolName: string) => {
onTagStart: (
tagName: string,
attributes: Record<string, string>,
) => void | Promise<void>
onTagEnd: (
tagName: string,
params: Record<string, any>,
) => void | Promise<void>
}
onResponseChunk: (chunk: PrintModeText | PrintModeError) => void
logger: Logger
loggerOptions?: {
userId?: string
model?: Model
agentName?: string
}
trackEvent: TrackEventFn
executeXmlToolCall: (params: {
toolName: string
input: Record<string, unknown>
}) => Promise<void>
}): AsyncGenerator<StreamChunk, PromptResult<string | null>> {
const {
stream,
processors,
defaultProcessor,
onResponseChunk,
logger,
loggerOptions,
trackEvent,
executeXmlToolCall,
} = params
let streamCompleted = false
let buffer = ''
let autocompleted = false
// State for parsing XML tool calls from text stream
const xmlParserState: StreamParserState = createStreamParserState()
async function processToolCallObject(params: {
toolName: string
input: any
contents?: string
}): Promise<void> {
const { toolName, contents } = params
let { input } = params
// AI SDK sometimes emits tool-call chunks with a raw JSON string as `input`
// when its repair pass can't produce a parsed object. Try to parse; if it
// fails, leave as string — the executor surfaces a clear error.
if (typeof input === 'string') {
try {
input = JSON.parse(input)
} catch {}
}
const processor = processors[toolName] ?? defaultProcessor(toolName)
trackEvent({
event: AnalyticsEvent.TOOL_USE,
userId: loggerOptions?.userId ?? '',
properties: {
toolName,
contents,
parsedParams: input,
autocompleted,
model: loggerOptions?.model,
agent: loggerOptions?.agentName,
},
logger,
})
await processor.onTagStart(toolName, {})
await processor.onTagEnd(toolName, input)
}
function flush() {
if (buffer) {
onResponseChunk({
type: 'text',
text: buffer,
})
}
buffer = ''
}
async function* processChunk(
chunk: StreamChunk | undefined,
): AsyncGenerator<StreamChunk> {
if (chunk === undefined) {
flush()
streamCompleted = true
return
}
if (chunk.type === 'text') {
// Parse XML tool calls from the text stream
const { filteredText, toolCalls } = parseStreamChunk(
chunk.text,
xmlParserState,
)
if (filteredText) {
buffer += filteredText
yield {
type: 'text',
text: filteredText,
}
}
// Flush buffer before yielding tool calls so text event is sent first
if (toolCalls.length > 0) {
flush()
}
// Then process and yield any XML tool calls found
for (const toolCall of toolCalls) {
// Execute the tool immediately if callback provided, pausing the stream
// The callback handles emitting tool_call and tool_result events
await executeXmlToolCall({
toolName: toolCall.toolName,
input: toolCall.input,
})
}
return
} else {
flush()
}
if (chunk.type === 'tool-call') {
await processToolCallObject(chunk)
}
yield chunk
}
let result: PromptResult<string | null> = { aborted: false, value: null }
try {
while (true) {
const { value, done } = await stream.next()
if (done) {
result = value
break
}
if (streamCompleted) {
break
}
yield* processChunk(value)
}
if (!streamCompleted) {
// After the stream ends, try parsing one last time in case there's leftover text
yield* processChunk(undefined)
}
} finally {
// Flush any remaining buffered text so it reaches onResponseChunk even on
// abort. Without this, text streamed after the last tool call would be lost
// from the message history.
flush()
}
return result
}