mirror of
https://github.com/windmill-labs/windmill.git
synced 2026-10-04 16:02:17 +00:00
* feat: run turns in several flow chat conversations at once Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> * fix: keep finished turns finished and cached chats current in the flow chat pool Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> * fix: attribute a turn's rows by job id as well as sequence Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> * fix: count only real stream updates and retry the job-id read Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> * fix: keep a chat that holds an unsent draft, and take one back when its first turn is withdrawn Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> * fix: ignore a stale running-turn snapshot, keep a withdrawn chat's draft, poll after clean stream ends Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> * fix: follow the turn running now when the listing named one already over Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> * fix: keep replacement turns and SSE fallback moving * fix: keep replacement turn handoffs active * fix: preserve unread badge line height * fix: settle local fallback handoffs * fix: settle refused turn handoffs * fix: scope turn handoffs to conversation * fix: drop stale turn handoffs * refactor: move the queued message and 409 handling into per-conversation turns Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> * fix: address cubic's review of the parallel flow chat turns Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> * fix: clear a stale failure on refresh, and tighten the docs and test waits Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> * fix: recover running rows past the first page, and drop the failure a re-read disproves Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> * refactor: check the running-turn query at compile time, and narrow what a refresh clears Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> * fix: settle a failed turn only from an answer that turn wrote Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> * fix: settle a failed turn from its own answer, and only while it is still the failure shown Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> * fix: drop a failure whose answer arrived even when a newer turn owns the error Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> * fix: free an answered failure whatever the turn that started meanwhile is doing Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> * fix: drop a rows read that a turn outran, rather than merging it under newer messages Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> * fix: drop a rows read whose conversation was left and opened again Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> * fix: hand over a file still being read when its composer goes Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> * fix: count a drop's routing as work in flight, so its file is handed over too Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> * fix: hold the send until every file a conversation is owed has landed Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> * refactor: keep a panel mounted per conversation instead of handing its draft over Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> * fix: keep the withdrawn chat whose composer was written in, not the empty one Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> * fix: keep the chat in front of the reader when both withdrawn composers were written in Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> * fix: keep a retry's own run arguments when a turn elsewhere refuses it Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> * fix: name panels apart across pools, and read a flow's inputs when its chat is built Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> --------- Co-authored-by: Claude Opus 5 (1M context) <noreply@anthropic.com>
419 lines
14 KiB
TypeScript
419 lines
14 KiB
TypeScript
import type { ChatAttachment, FetchLike, TokenSource } from './types'
|
|
|
|
export interface WindmillChatApiOptions {
|
|
baseUrl: string
|
|
workspace: string
|
|
/** Omit to rely on the session cookie of the Windmill origin. */
|
|
token?: TokenSource
|
|
fetch?: FetchLike
|
|
/** Server poll interval for a turn's stream (Enterprise; see `ChatOptions.pollDelayMs`). */
|
|
pollDelayMs?: number
|
|
}
|
|
|
|
export class WindmillApiError extends Error {
|
|
constructor(
|
|
message: string,
|
|
readonly status: number,
|
|
/** The response body, as the server sent it. */
|
|
readonly body?: string
|
|
) {
|
|
super(message)
|
|
this.name = 'WindmillApiError'
|
|
}
|
|
}
|
|
|
|
/** The turn a conversation is still answering, as the server reports it. */
|
|
export interface RunningTurn {
|
|
jobId: string
|
|
/** `created_seq` of the user message that started the turn. */
|
|
userSeq: number
|
|
}
|
|
|
|
/**
|
|
* A message was sent to a conversation whose turn is still running. The server refuses
|
|
* it (409) so two runs never write one agent memory; `turn` is the run to follow instead.
|
|
*/
|
|
export class TurnRunningError extends Error {
|
|
constructor(
|
|
message: string,
|
|
readonly turn: RunningTurn
|
|
) {
|
|
super(message)
|
|
this.name = 'TurnRunningError'
|
|
}
|
|
}
|
|
|
|
/**
|
|
* The running turn named by a run's 409 body, or undefined for any other body. Exported
|
|
* for a custom `run` that calls Windmill through its own client: it rethrows the refusal
|
|
* as a `TurnRunningError` so the chat can follow the running turn.
|
|
*/
|
|
export function turnRunningError(body: string): TurnRunningError | undefined {
|
|
try {
|
|
const parsed = JSON.parse(body) as { error?: unknown; running_turn?: FlowConversation['running_turn'] }
|
|
const turn = parsed?.running_turn
|
|
if (!turn || typeof turn.job_id !== 'string' || typeof turn.user_seq !== 'number') return undefined
|
|
const message = typeof parsed.error === 'string' ? parsed.error : 'this conversation is still answering a message'
|
|
return new TurnRunningError(message, { jobId: turn.job_id, userSeq: turn.user_seq })
|
|
} catch {
|
|
return undefined
|
|
}
|
|
}
|
|
|
|
export interface FlowConversation {
|
|
id: string
|
|
workspace_id: string
|
|
flow_path: string
|
|
title?: string | null
|
|
created_at: string
|
|
updated_at: string
|
|
created_by: string
|
|
/** Started from the flow editor's test panel rather than a deployed run. */
|
|
is_test: boolean
|
|
/** Set by the list: the turn this conversation is still answering. */
|
|
running_turn?: { job_id: string; user_seq: number } | null
|
|
}
|
|
|
|
/**
|
|
* Which conversations a listing holds: the flow editor's test chats, the deployed flow's
|
|
* own (the server's default), or both.
|
|
*/
|
|
export type ConversationKind = 'test' | 'deployed' | 'all'
|
|
|
|
export interface FlowConversationMessage {
|
|
id: string
|
|
conversation_id: string
|
|
message_type: 'user' | 'assistant' | 'system' | 'tool'
|
|
content: string
|
|
job_id?: string | null
|
|
created_at: string
|
|
created_seq: number
|
|
step_name?: string | null
|
|
success?: boolean
|
|
/** On a tool row, the arguments the model wrote, without the inputs a step wires in; null for a web search. */
|
|
tool_arguments?: string | null
|
|
/** On a tool row, the text the model got back or what the call failed with; a web search's citations. */
|
|
tool_result?: string | null
|
|
/** On an answer, the thinking that produced it; on a tool row, the thinking that led to the call. */
|
|
reasoning?: string | null
|
|
/** The files a user message carried, as object-storage references. */
|
|
attachments?: ChatAttachment[] | null
|
|
}
|
|
|
|
export type JobUpdateEvent =
|
|
| {
|
|
type: 'update'
|
|
running?: boolean
|
|
completed?: boolean
|
|
new_result_stream?: string
|
|
stream_offset?: number
|
|
only_result?: unknown
|
|
flow_stream_job_id?: string
|
|
}
|
|
| { type: 'error'; error: string }
|
|
| { type: 'notfound' }
|
|
| { type: 'timeout' }
|
|
| { type: 'ping' }
|
|
|
|
export interface CompletedJobResult {
|
|
completed: boolean
|
|
success?: boolean
|
|
result?: unknown
|
|
}
|
|
|
|
/** The part of a flow job's status that names the jobs its steps ran as. */
|
|
export interface FlowJobStatus {
|
|
flow_status?: {
|
|
modules?: FlowStepStatus[] | null
|
|
failure_module?: FlowStepStatus | null
|
|
preprocessor_module?: FlowStepStatus | null
|
|
} | null
|
|
}
|
|
|
|
export interface FlowStepStatus {
|
|
job?: string | null
|
|
flow_jobs?: string[] | null
|
|
/** An agent step's rounds; a tool call ran as a job of its own, which its row is persisted under. */
|
|
agent_actions?: { type?: string; job_id?: string | null }[] | null
|
|
}
|
|
|
|
/** Thin client over the Windmill endpoints a chat-mode flow uses. */
|
|
export class WindmillChatApi {
|
|
readonly #baseUrl: string
|
|
readonly #workspace: string
|
|
readonly #token: TokenSource | undefined
|
|
readonly #fetch: FetchLike
|
|
readonly #pollDelayMs: number | undefined
|
|
|
|
constructor(options: WindmillChatApiOptions) {
|
|
this.#baseUrl = normalizeBaseUrl(options.baseUrl)
|
|
this.#workspace = options.workspace
|
|
this.#token = options.token
|
|
this.#fetch = options.fetch ?? ((input, init) => globalThis.fetch(input, init))
|
|
this.#pollDelayMs = options.pollDelayMs
|
|
}
|
|
|
|
/**
|
|
* Starts a turn: runs the flow with `memory_id` set to the conversation id. Returns the
|
|
* job id. Throws `TurnRunningError` when the conversation is still answering.
|
|
*/
|
|
async runFlow(
|
|
flowPath: string,
|
|
args: Record<string, unknown>,
|
|
options: { memoryId: string; signal?: AbortSignal }
|
|
): Promise<string> {
|
|
let res: Response
|
|
try {
|
|
res = await this.#request(`jobs/run/f/${encodePath(flowPath)}`, {
|
|
method: 'POST',
|
|
query: { memory_id: options.memoryId, skip_preprocessor: 'true' },
|
|
body: args,
|
|
signal: options.signal
|
|
})
|
|
} catch (e) {
|
|
if (e instanceof WindmillApiError && e.status === 409) throw turnRunningError(e.body ?? '') ?? e
|
|
throw e
|
|
}
|
|
return (await res.text()).trim()
|
|
}
|
|
|
|
/**
|
|
* One server-sent-events connection to a job's updates. The server closes it after
|
|
* `TIMEOUT_SSE_STREAM` (a `timeout` event); resume by calling again with the last
|
|
* `stream_offset`, never by re-running the flow.
|
|
*/
|
|
async *streamJob(
|
|
jobId: string,
|
|
options: { streamOffset?: number; signal?: AbortSignal } = {}
|
|
): AsyncGenerator<JobUpdateEvent> {
|
|
const query: Record<string, string> = { fast: 'true', only_result: 'true' }
|
|
if (this.#pollDelayMs !== undefined) query.poll_delay_ms = String(this.#pollDelayMs)
|
|
if (options.streamOffset !== undefined) {
|
|
query.stream_offset = String(options.streamOffset)
|
|
}
|
|
const res = await this.#request(`jobs_u/getupdate_sse/${encodeURIComponent(jobId)}`, {
|
|
query,
|
|
accept: 'text/event-stream',
|
|
signal: options.signal
|
|
})
|
|
if (!res.body) {
|
|
// Status 0, whatever the status line said: a response with no stream in it comes from
|
|
// something in front of Windmill, so it is followed by a reconnect like any other.
|
|
throw new WindmillApiError('The job update stream has no body', 0)
|
|
}
|
|
for await (const data of readServerSentEvents(res.body)) {
|
|
try {
|
|
yield JSON.parse(data) as JobUpdateEvent
|
|
} catch {
|
|
// A frame that isn't JSON carries nothing the chat can use.
|
|
}
|
|
}
|
|
}
|
|
|
|
async getCompletedResult(jobId: string, signal?: AbortSignal): Promise<CompletedJobResult> {
|
|
const res = await this.#request(
|
|
`jobs_u/completed/get_result_maybe/${encodeURIComponent(jobId)}`,
|
|
{ signal }
|
|
)
|
|
return (await res.json()) as CompletedJobResult
|
|
}
|
|
|
|
/** A flow job with its status: the step job ids are what persisted messages carry as `job_id`. */
|
|
async getFlowJob(jobId: string, signal?: AbortSignal): Promise<FlowJobStatus> {
|
|
const res = await this.#request(`jobs_u/get/${encodeURIComponent(jobId)}`, {
|
|
query: { no_logs: 'true' },
|
|
signal
|
|
})
|
|
return (await res.json()) as FlowJobStatus
|
|
}
|
|
|
|
/**
|
|
* Where a message's attachment downloads from. The endpoint authenticates like every other
|
|
* request: a consumer holding a token must fetch it with that token, not put the URL in an
|
|
* `img src`, which would send only the Windmill session cookie.
|
|
*/
|
|
attachmentUrl(attachment: ChatAttachment): string {
|
|
const query = new URLSearchParams({ file_key: attachment.s3 })
|
|
if (attachment.storage) query.set('storage', attachment.storage)
|
|
return `${this.#baseUrl}/api/w/${encodeURIComponent(this.#workspace)}/job_helpers/download_s3_file?${query}`
|
|
}
|
|
|
|
async cancelJob(jobId: string, reason = 'Stopped from the chat'): Promise<void> {
|
|
await this.#request(`jobs_u/queue/cancel/${encodeURIComponent(jobId)}`, {
|
|
method: 'POST',
|
|
body: { reason }
|
|
})
|
|
}
|
|
|
|
async listConversations(
|
|
flowPath: string,
|
|
options: { page?: number; perPage?: number; kind?: ConversationKind; signal?: AbortSignal } = {}
|
|
): Promise<FlowConversation[]> {
|
|
const extra: Record<string, string> = { flow_path: flowPath }
|
|
if (options.kind !== undefined) extra.kind = options.kind
|
|
const res = await this.#request('flow_conversations/list', {
|
|
query: pagination(options, extra),
|
|
signal: options.signal
|
|
})
|
|
return (await res.json()) as FlowConversation[]
|
|
}
|
|
|
|
/** Sets a conversation's title. Its place in the list is kept: only a turn moves one. */
|
|
async renameConversation(conversationId: string, title: string): Promise<void> {
|
|
await this.#request(`flow_conversations/update/${encodeURIComponent(conversationId)}`, {
|
|
method: 'POST',
|
|
body: { title }
|
|
})
|
|
}
|
|
|
|
/**
|
|
* Without `afterSeq`: one page counted from the newest message, returned oldest first.
|
|
* With `afterSeq`: the messages created after that cursor, oldest first.
|
|
*/
|
|
async listMessages(
|
|
conversationId: string,
|
|
options: { page?: number; perPage?: number; afterSeq?: number; signal?: AbortSignal } = {}
|
|
): Promise<FlowConversationMessage[]> {
|
|
const extra: Record<string, string> = {}
|
|
if (options.afterSeq !== undefined) extra.after_seq = String(options.afterSeq)
|
|
const res = await this.#request(
|
|
`flow_conversations/${encodeURIComponent(conversationId)}/messages`,
|
|
{ query: pagination(options, extra), signal: options.signal }
|
|
)
|
|
return (await res.json()) as FlowConversationMessage[]
|
|
}
|
|
|
|
/**
|
|
* Puts bytes in the workspace's object storage under `fileKey` and returns the key they were
|
|
* stored under (the server may rewrite it). Needs the workspace to have object storage set.
|
|
*/
|
|
async uploadFile(
|
|
fileKey: string,
|
|
body: Blob,
|
|
options: { contentType?: string; signal?: AbortSignal } = {}
|
|
): Promise<{ file_key: string }> {
|
|
const query: Record<string, string> = { file_key: fileKey }
|
|
if (options.contentType) query.content_type = options.contentType
|
|
const res = await this.#request('job_helpers/upload_s3_file', {
|
|
method: 'POST',
|
|
query,
|
|
raw: body,
|
|
contentType: options.contentType || 'application/octet-stream',
|
|
signal: options.signal
|
|
})
|
|
return (await res.json()) as { file_key: string }
|
|
}
|
|
|
|
async deleteConversation(conversationId: string): Promise<void> {
|
|
await this.#request(`flow_conversations/delete/${encodeURIComponent(conversationId)}`, {
|
|
method: 'DELETE'
|
|
})
|
|
}
|
|
|
|
async #request(
|
|
path: string,
|
|
init: {
|
|
method?: string
|
|
query?: Record<string, string>
|
|
/** JSON-encoded. */
|
|
body?: unknown
|
|
/** Sent as is, under `contentType`. */
|
|
raw?: Blob
|
|
contentType?: string
|
|
accept?: string
|
|
signal?: AbortSignal
|
|
} = {}
|
|
): Promise<Response> {
|
|
const url = new URL(`${this.#baseUrl}/api/w/${encodeURIComponent(this.#workspace)}/${path}`)
|
|
for (const [k, v] of Object.entries(init.query ?? {})) url.searchParams.set(k, v)
|
|
|
|
const headers: Record<string, string> = {}
|
|
if (init.accept) headers['Accept'] = init.accept
|
|
if (init.body !== undefined) headers['Content-Type'] = 'application/json'
|
|
else if (init.raw !== undefined) headers['Content-Type'] = init.contentType ?? 'application/octet-stream'
|
|
const token = typeof this.#token === 'function' ? await this.#token() : this.#token
|
|
if (token) headers['Authorization'] = `Bearer ${token}`
|
|
|
|
const res = await this.#fetch(url.toString(), {
|
|
method: init.method ?? 'GET',
|
|
headers,
|
|
body: init.body === undefined ? init.raw : JSON.stringify(init.body),
|
|
// A token must not be paired with ambient cookies; without one, the cookie is
|
|
// the credential and only rides same-origin requests.
|
|
credentials: token ? 'omit' : 'same-origin',
|
|
signal: init.signal
|
|
})
|
|
if (!res.ok) {
|
|
const text = await res.text().catch(() => '')
|
|
throw new WindmillApiError(
|
|
`${init.method ?? 'GET'} ${path} failed (${res.status})${text ? `: ${text}` : ''}`,
|
|
res.status,
|
|
text
|
|
)
|
|
}
|
|
return res
|
|
}
|
|
}
|
|
|
|
export function normalizeBaseUrl(baseUrl: string): string {
|
|
return baseUrl.replace(/\/+$/, '').replace(/\/api$/, '')
|
|
}
|
|
|
|
function encodePath(path: string): string {
|
|
return path.split('/').map(encodeURIComponent).join('/')
|
|
}
|
|
|
|
function pagination(
|
|
options: { page?: number; perPage?: number },
|
|
extra: Record<string, string>
|
|
): Record<string, string> {
|
|
const query = { ...extra }
|
|
if (options.page !== undefined) query.page = String(options.page)
|
|
if (options.perPage !== undefined) query.per_page = String(options.perPage)
|
|
return query
|
|
}
|
|
|
|
/** Yields the `data` payload of each event in a `text/event-stream` body. */
|
|
export async function* readServerSentEvents(
|
|
body: ReadableStream<Uint8Array>
|
|
): AsyncGenerator<string> {
|
|
const reader = body.getReader()
|
|
const decoder = new TextDecoder()
|
|
let buffer = ''
|
|
// A CR ending a chunk may be half of a CRLF; it waits for the next chunk.
|
|
let carry = ''
|
|
try {
|
|
while (true) {
|
|
const { value, done } = await reader.read()
|
|
if (done) break
|
|
let text = carry + decoder.decode(value, { stream: true })
|
|
carry = ''
|
|
if (text.endsWith('\r')) {
|
|
carry = '\r'
|
|
text = text.slice(0, -1)
|
|
}
|
|
buffer += text.replace(/\r\n?/g, '\n')
|
|
let end: number
|
|
while ((end = buffer.indexOf('\n\n')) !== -1) {
|
|
const data = eventData(buffer.slice(0, end))
|
|
buffer = buffer.slice(end + 2)
|
|
if (data !== undefined) yield data
|
|
}
|
|
}
|
|
if (carry) buffer += '\n'
|
|
const data = eventData(buffer)
|
|
if (data !== undefined) yield data
|
|
} finally {
|
|
// Closes the connection when the consumer stops early.
|
|
reader.cancel().catch(() => {})
|
|
}
|
|
}
|
|
|
|
function eventData(block: string): string | undefined {
|
|
const lines = block
|
|
.split('\n')
|
|
.filter((line) => line.startsWith('data:'))
|
|
.map((line) => line.slice(line.startsWith('data: ') ? 6 : 5))
|
|
return lines.length > 0 ? lines.join('\n') : undefined
|
|
}
|