mirror of
https://github.com/windmill-labs/windmill.git
synced 2026-10-04 08:02:23 +00:00
* feat: render an AI agent result as its answer, not as raw JSON Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> * fix: sanitize agent markdown through the shared plugin chain Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> * refactor: fold agent stream events incrementally per poll Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> * fix: separate the agent meta line from the result toggle group Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> * feat: add a transcript view of an agent run's conversation Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> * refactor: retire AIAgentLogViewer in favour of the transcript Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> * feat: show the partial transcript a max-iterations failure carries Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> * fix: move the agent meta line and system prompt below the conversation Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> * refactor: show an agent run as what it did, not as a conversation Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> * refactor: name the agent run breakdown a trace Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> * fix(ai-agent): label the save-as-agent form fields per the guidelines Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> * fix(ai-agent): reuse the resource form's path and description fields Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> * fix(ai-agent): keep the action tags on a max-iterations failure Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> * fix(ai-agent): keep offline replay inert and the streamed answer to one turn Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> * fix(ai-agent): reset the streamed answer on providers that skip tool_call Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> * fix: carry the event type narrowing through the stream parser Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> * fix: survive a malformed message rather than take the viewer down Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> * refactor: coerce agent messages once at the parse boundary Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> * feat: show an agent run as one scroll ending in its output Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> * fix: keep every streamed turn instead of dropping the narration Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> * fix: share the chat divider and drop the unsafe run auto-scroll Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> * fix: give the labelled divider a border colour and the standard pretty icon Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> * fix: pad the agent run below its badges as well as above Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> * feat: follow a streaming run's pane without moving the page Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> * refactor: share the chat's stick-to-bottom mechanics with the agent run Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> * fix: keep the answer's citations and end a turn's reasoning with the turn Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> * refactor: read the agent stream through the windmill-chat sdk parser Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> * fix: show an agent run's thinking instead of falling back to raw json Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> * fix: stop a stream the fold cannot use from claiming the run pane Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> * fix: identify an agent run by more than the job id the replay withholds Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> * refactor: drop the tests and comment lines that were not earning their place Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> * fix: stop a streamed turn's text shifting when a tool call closes it Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> --------- Co-authored-by: Claude Opus 5 (1M context) <noreply@anthropic.com>
342 lines
12 KiB
TypeScript
342 lines
12 KiB
TypeScript
import type { FlowStatusModule } from '$lib/gen'
|
|
import { parseStreamEvents, type AgentStreamEvent } from 'windmill-chat'
|
|
|
|
/** The `agent_action` tag the worker puts on every message it records. */
|
|
export type AgentAction = NonNullable<FlowStatusModule['agent_actions']>[number]
|
|
|
|
export type AgentTokenUsage = {
|
|
input_tokens?: number
|
|
output_tokens?: number
|
|
total_tokens?: number
|
|
cache_read_input_tokens?: number
|
|
cache_write_input_tokens?: number
|
|
}
|
|
|
|
export type AgentMessage = {
|
|
role: string
|
|
content?: unknown
|
|
tool_calls?: Array<{
|
|
id?: string
|
|
type?: string
|
|
function?: { name?: string; arguments?: string }
|
|
}>
|
|
tool_call_id?: string
|
|
agent_action?: AgentAction
|
|
annotations?: Array<{ url: string; title?: string; start_index?: number; end_index?: number }>
|
|
}
|
|
|
|
/** The envelope every AI agent step returns, built by `AIAgentResult`. */
|
|
export type AgentResult = {
|
|
output: unknown
|
|
messages: AgentMessage[]
|
|
usage?: AgentTokenUsage
|
|
wm_stream?: string
|
|
/** The model's thinking across every iteration, blank-line separated. */
|
|
reasoning?: string
|
|
}
|
|
|
|
/**
|
|
* Every key `AIAgentResult` can serialize; the optional ones are skipped when
|
|
* empty. The signature below is closed, so a key the worker gains and this list
|
|
* does not makes every run carrying it fall back to the raw JSON.
|
|
*/
|
|
const ENVELOPE_KEYS = ['output', 'messages', 'usage', 'wm_stream', 'reasoning']
|
|
|
|
function isRecord(value: unknown): value is Record<string, unknown> {
|
|
return typeof value === 'object' && value !== null && !Array.isArray(value)
|
|
}
|
|
|
|
function hasRole(message: unknown): boolean {
|
|
return isRecord(message) && typeof message.role === 'string'
|
|
}
|
|
|
|
/**
|
|
* A job result is whatever its script returned, so a message that passed the shape
|
|
* check still has anything underneath. Coerced once here rather than guarded at
|
|
* each reader: a wrong type reaching one throws mid-render and takes the whole
|
|
* result viewer down, including the plain error it usually rides on.
|
|
*/
|
|
function toAgentMessage(raw: Record<string, unknown>): AgentMessage {
|
|
const toolCalls = Array.isArray(raw.tool_calls)
|
|
? raw.tool_calls.filter(isRecord).map((call) => ({
|
|
id: typeof call.id === 'string' ? call.id : undefined,
|
|
type: typeof call.type === 'string' ? call.type : undefined,
|
|
function: isRecord(call.function)
|
|
? {
|
|
name: typeof call.function.name === 'string' ? call.function.name : undefined,
|
|
arguments:
|
|
typeof call.function.arguments === 'string' ? call.function.arguments : undefined
|
|
}
|
|
: undefined
|
|
}))
|
|
: undefined
|
|
const annotations = Array.isArray(raw.annotations)
|
|
? raw.annotations.filter(
|
|
(a): a is Record<string, unknown> => isRecord(a) && typeof a.url === 'string'
|
|
)
|
|
: undefined
|
|
return {
|
|
role: raw.role as string,
|
|
content: raw.content,
|
|
tool_calls: toolCalls,
|
|
tool_call_id: typeof raw.tool_call_id === 'string' ? raw.tool_call_id : undefined,
|
|
// The union is discriminated on `type`; an action without a string one
|
|
// matches no branch and is treated as untagged.
|
|
agent_action:
|
|
isRecord(raw.agent_action) && typeof raw.agent_action.type === 'string'
|
|
? (raw.agent_action as unknown as AgentAction)
|
|
: undefined,
|
|
annotations: annotations as AgentMessage['annotations']
|
|
}
|
|
}
|
|
|
|
function toAgentMessages(raw: unknown[]): AgentMessage[] {
|
|
return raw.map((message) => toAgentMessage(message as Record<string, unknown>))
|
|
}
|
|
|
|
/**
|
|
* Recognised by shape, not a marker key: sniffing works on completed runs, and an
|
|
* added key would travel into a parent agent's conversation. Deliberately not also
|
|
* requiring a tagged `agent_action` — a run answering through a structured-output
|
|
* tool tags nothing, and hiding a real answer costs more than claiming a lookalike.
|
|
*/
|
|
export function parseAgentResult(result: unknown): AgentResult | undefined {
|
|
if (!isRecord(result)) {
|
|
return undefined
|
|
}
|
|
const keys = Object.keys(result)
|
|
if (!keys.every((key) => ENVELOPE_KEYS.includes(key))) {
|
|
return undefined
|
|
}
|
|
if (!('output' in result) || !Array.isArray(result.messages)) {
|
|
return undefined
|
|
}
|
|
// An agent always records at least the message it was asked, so an empty list
|
|
// is someone else's result rather than a run that did nothing.
|
|
if (result.messages.length === 0 || !result.messages.every(hasRole)) {
|
|
return undefined
|
|
}
|
|
return {
|
|
output: result.output,
|
|
messages: toAgentMessages(result.messages),
|
|
usage: isRecord(result.usage) ? (result.usage as AgentTokenUsage) : undefined,
|
|
wm_stream: typeof result.wm_stream === 'string' ? result.wm_stream : undefined,
|
|
reasoning: typeof result.reasoning === 'string' ? result.reasoning : undefined
|
|
}
|
|
}
|
|
|
|
/**
|
|
* A run stopped by `max_iterations` fails, so it returns an error rather than an
|
|
* envelope — but the worker attaches the conversation so far to it. That partial
|
|
* transcript is the whole reason to look at a run that hit the cap.
|
|
*/
|
|
export function parseAgentErrorMessages(result: unknown): AgentMessage[] | undefined {
|
|
if (!isRecord(result) || !isRecord(result.error)) {
|
|
return undefined
|
|
}
|
|
const inner = result.error.result
|
|
if (!isRecord(inner) || !Array.isArray(inner.messages)) {
|
|
return undefined
|
|
}
|
|
if (inner.messages.length === 0 || !inner.messages.every(hasRole)) {
|
|
return undefined
|
|
}
|
|
return toAgentMessages(inner.messages)
|
|
}
|
|
|
|
export type AgentResultSummary = {
|
|
toolCalls: number
|
|
webSearches: number
|
|
tokens: number | undefined
|
|
cachedTokens: number | undefined
|
|
}
|
|
|
|
function actionType(message: AgentMessage): string | undefined {
|
|
return message.agent_action?.type
|
|
}
|
|
|
|
export function summarizeAgentResult(result: AgentResult): AgentResultSummary {
|
|
let toolCalls = 0
|
|
let webSearches = 0
|
|
for (const message of result.messages) {
|
|
const type = actionType(message)
|
|
if (type === 'tool_call' || type === 'mcp_tool_call') {
|
|
toolCalls++
|
|
} else if (type === 'web_search') {
|
|
webSearches++
|
|
}
|
|
}
|
|
const usage = result.usage
|
|
// `total_tokens` is what providers report when they report anything; fall back
|
|
// to the parts so a provider that only sends the split still shows a count.
|
|
const tokens =
|
|
usage?.total_tokens ??
|
|
(usage?.input_tokens !== undefined || usage?.output_tokens !== undefined
|
|
? (usage?.input_tokens ?? 0) + (usage?.output_tokens ?? 0)
|
|
: undefined)
|
|
return {
|
|
toolCalls,
|
|
webSearches,
|
|
tokens,
|
|
cachedTokens: usage?.cache_read_input_tokens
|
|
}
|
|
}
|
|
|
|
export type AgentStreamEntry =
|
|
| { kind: 'assistant'; content: string }
|
|
| { kind: 'tool'; callId: string; name: string; running: boolean; success?: boolean }
|
|
|
|
export type AgentStream = {
|
|
/**
|
|
* What the run has finished doing, in order — the same rows the completed
|
|
* trace will show, so nothing on screen moves when the result lands.
|
|
*/
|
|
entries: AgentStreamEntry[]
|
|
/**
|
|
* The text of the turn being written. It is not yet the output: a turn that
|
|
* goes on to call a tool was narration, and only the run ending decides which
|
|
* this was. So it stays unlabelled here and becomes one or the other.
|
|
*/
|
|
current: string
|
|
reasoning: string
|
|
}
|
|
|
|
/** How much of the stream has been folded in, so the next poll starts there. */
|
|
export type AgentStreamProgress = { consumed: number; stream: AgentStream }
|
|
|
|
export function emptyAgentStreamProgress(): AgentStreamProgress {
|
|
return { consumed: 0, stream: { entries: [], current: '', reasoning: '' } }
|
|
}
|
|
|
|
/**
|
|
* Any of these means a tool call is beginning, and which one arrives depends on
|
|
* the provider: the SSE parsers announce `tool_call`, while Bedrock streams only
|
|
* the arguments and the worker follows with `tool_execution`. Resetting on all
|
|
* three is idempotent and keeps the rule provider-independent.
|
|
*/
|
|
const TOOL_TURN_STARTED: AgentStreamEvent['type'][] = [
|
|
'tool_call',
|
|
'tool_call_arguments',
|
|
'tool_execution'
|
|
]
|
|
|
|
/**
|
|
* Whether an event carries what its type declares: a delta its text, a tool event
|
|
* the `call_id` its row is keyed on and the name that labels it. `result_stream` is
|
|
* whatever the job wrote, so detection and the fold ask the same question — were
|
|
* they to differ, a stream would be claimed and then render nothing.
|
|
*/
|
|
function isWellFormedEvent(event: AgentStreamEvent): boolean {
|
|
if (event.type === 'token_delta' || event.type === 'reasoning_token_delta') {
|
|
return typeof event.content === 'string'
|
|
}
|
|
return (
|
|
typeof event.call_id === 'string' &&
|
|
event.call_id !== '' &&
|
|
typeof event.function_name === 'string' &&
|
|
event.function_name !== ''
|
|
)
|
|
}
|
|
|
|
/**
|
|
* Whether `result_stream` is an agent's event stream rather than something a
|
|
* script printed. Reads only the first complete line, because it runs on every
|
|
* poll of a running job.
|
|
*/
|
|
export function isAgentStream(raw: string): boolean {
|
|
let start = 0
|
|
while (start < raw.length) {
|
|
const end = raw.indexOf('\n', start)
|
|
if (end === -1) {
|
|
// Only a partial first line so far; wait for the poll that completes it.
|
|
return false
|
|
}
|
|
const line = raw.slice(start, end)
|
|
if (line.trim() !== '') {
|
|
return parseStreamEvents(line).some(isWellFormedEvent)
|
|
}
|
|
start = end + 1
|
|
}
|
|
return false
|
|
}
|
|
|
|
/**
|
|
* Fold the events that arrived since `previous` into the answer so far. Incremental
|
|
* rather than a parse of the whole buffer: the stream only grows, a poll can arrive
|
|
* every 50ms, and a `tool_result` carries the tool's entire output, so re-reading it
|
|
* all each tick is quadratic with a large constant.
|
|
*/
|
|
export function advanceAgentStream(
|
|
raw: string,
|
|
previous: AgentStreamProgress
|
|
): AgentStreamProgress {
|
|
// A trailing line with no newline yet is still being written, so it stays
|
|
// unconsumed until the poll that completes it.
|
|
const complete = raw.lastIndexOf('\n') + 1
|
|
if (complete <= previous.consumed) {
|
|
return previous
|
|
}
|
|
const stream: AgentStream = { ...previous.stream, entries: [...previous.stream.entries] }
|
|
for (const event of parseStreamEvents(raw.slice(previous.consumed, complete))) {
|
|
if (!isWellFormedEvent(event)) {
|
|
continue
|
|
}
|
|
if (event.type === 'token_delta') {
|
|
stream.current += event.content
|
|
continue
|
|
}
|
|
if (event.type === 'reasoning_token_delta') {
|
|
stream.reasoning += event.content
|
|
continue
|
|
}
|
|
if (TOOL_TURN_STARTED.includes(event.type)) {
|
|
if (stream.current !== '') {
|
|
// A model can narrate and request a tool in the same turn. The call
|
|
// settles what that text was: narration, not the output. It becomes a
|
|
// row rather than being dropped, so nothing vanishes from the screen
|
|
// only to reappear when the result lands.
|
|
stream.entries.push({ kind: 'assistant', content: stream.current })
|
|
stream.current = ''
|
|
}
|
|
// Thinking belongs to the turn that produced it, and a turn can think
|
|
// without narrating — extended thinking before a tool call is exactly
|
|
// that shape. So this clears on the boundary itself, not with the
|
|
// narration, or one turn's thoughts run into the next turn's.
|
|
stream.reasoning = ''
|
|
}
|
|
// The same call is announced, then argued, then executed, then answered.
|
|
// Keyed on `call_id` so those four events are one row rather than four.
|
|
const existing = stream.entries.find(
|
|
(e): e is Extract<AgentStreamEntry, { kind: 'tool' }> =>
|
|
e.kind === 'tool' && e.callId === event.call_id
|
|
)
|
|
const settled = event.type === 'tool_result'
|
|
if (existing) {
|
|
existing.running = !settled
|
|
existing.success = settled ? event.success === true : existing.success
|
|
} else {
|
|
stream.entries.push({
|
|
kind: 'tool',
|
|
callId: event.call_id,
|
|
name: event.function_name,
|
|
running: !settled,
|
|
success: settled ? event.success === true : undefined
|
|
})
|
|
}
|
|
}
|
|
return { consumed: complete, stream }
|
|
}
|
|
|
|
/**
|
|
* Token counts run to five and six figures, where the exact digit is noise. The
|
|
* millions branch is not decoration: usage accumulates over every loop
|
|
* iteration, and each one re-sends the whole context.
|
|
*/
|
|
export function formatTokenCount(count: number): string {
|
|
if (count < 1000) {
|
|
return String(count)
|
|
}
|
|
const [scaled, unit] = count < 1_000_000 ? [count / 1000, 'k'] : [count / 1_000_000, 'M']
|
|
return `${scaled < 10 ? scaled.toFixed(1) : Math.round(scaled)}${unit}`
|
|
}
|