diff --git a/CONTEXT.md b/CONTEXT.md index fbc2592cc9..054960607c 100644 --- a/CONTEXT.md +++ b/CONTEXT.md @@ -48,8 +48,26 @@ _Avoid_: thread, session (that names an AI session, a different thing), chat (th One question and the answer to it: the run the question started, the handle that stops it, and the rows it is writing. At most one per conversation, and the chat is held for its whole length — from the moment the question takes the chat, before it has a job, until it is ended. +The server holds the same rule: a question sent to a conversation whose turn is still running +is refused, whoever sends it. Several conversations of one flow can each have a turn running. _Avoid_: request, exchange, message round +**Running turn**: +A turn whose run has not finished. Which conversations have one is the server's to say, so a +chat opened after the turn started — a reload, another tab — still sees it and follows it. +_Avoid_: busy, active, in flight + +**Queued message**: +A question typed into a conversation while its turn runs, sent when that turn ends answered. +At most one per conversation; typing another adds to it. A turn that fails, or is stopped, +hands it back to the composer instead of sending it. +_Avoid_: pending message (a pending message is one already sent and not yet confirmed) + +**Unread**: +The answers that arrived in a conversation while it was not the one shown. Counted per open +chat and forgotten on reload. +_Avoid_: new messages, notifications + **Transcript**: The rows a conversation's chat holds. Not the conversation: it is the newest page plus whatever older pages the reader has scrolled back through, so a question it cannot answer diff --git a/backend/.sqlx/query-2bf039a2880555c5e6eec9a29c4a380aa832d72d2aaf5b109ee9554f1e9da6a4.json b/backend/.sqlx/query-2bf039a2880555c5e6eec9a29c4a380aa832d72d2aaf5b109ee9554f1e9da6a4.json new file mode 100644 index 0000000000..46d609ecb4 --- /dev/null +++ b/backend/.sqlx/query-2bf039a2880555c5e6eec9a29c4a380aa832d72d2aaf5b109ee9554f1e9da6a4.json @@ -0,0 +1,34 @@ +{ + "db_name": "PostgreSQL", + "query": "SELECT c.id AS \"id!\", u.job_id AS \"job_id!\", u.created_seq AS \"user_seq!\"\n FROM unnest($1::uuid[]) AS c(id)\n CROSS JOIN LATERAL (\n SELECT job_id, created_seq\n FROM flow_conversation_message\n WHERE conversation_id = c.id AND message_type = 'user'\n ORDER BY created_seq DESC\n LIMIT 1\n ) u\n WHERE u.job_id IS NOT NULL\n AND EXISTS (SELECT 1 FROM v2_job_queue q WHERE q.id = u.job_id)", + "describe": { + "columns": [ + { + "ordinal": 0, + "name": "id!", + "type_info": "Uuid" + }, + { + "ordinal": 1, + "name": "job_id!", + "type_info": "Uuid" + }, + { + "ordinal": 2, + "name": "user_seq!", + "type_info": "Int8" + } + ], + "parameters": { + "Left": [ + "UuidArray" + ] + }, + "nullable": [ + null, + true, + false + ] + }, + "hash": "2bf039a2880555c5e6eec9a29c4a380aa832d72d2aaf5b109ee9554f1e9da6a4" +} diff --git a/backend/tests/flow_conversation_running_turn.rs b/backend/tests/flow_conversation_running_turn.rs new file mode 100644 index 0000000000..392ff6f5cf --- /dev/null +++ b/backend/tests/flow_conversation_running_turn.rs @@ -0,0 +1,87 @@ +//! A conversation answers one message at a time: while the run its newest user message started +//! is queued or running, the list reports that turn and a run into the conversation is refused. +//! +//! Uses runtime `sqlx::query` (not the compile-time macros) so no offline query cache is +//! needed, matching v2_job_delete_orphans.rs. + +use sqlx::{Pool, Postgres}; +use uuid::Uuid; +use windmill_common::error::Error; +use windmill_common::flow_conversations::{get_or_create_conversation_with_id, running_turns}; +use windmill_test_utils::*; + +const WS: &str = "test-workspace"; +const CONV: Uuid = Uuid::from_u128(0x5eed); + +/// Takes the conversation for a turn, as a run does, without keeping what it writes. +async fn take_conversation(db: &Pool) -> windmill_common::error::Result { + let mut tx = db.begin().await?; + let result = + get_or_create_conversation_with_id(&mut tx, WS, "f/flow", "test-user", "t", CONV, false) + .await; + tx.rollback().await?; + result.map(|c| c.id) +} + +#[sqlx::test(fixtures("base"))] +async fn test_a_conversation_refuses_a_turn_while_one_runs( + db: Pool, +) -> anyhow::Result<()> { + initialize_tracing().await; + + let job_id = Uuid::new_v4(); + sqlx::query("INSERT INTO v2_job (id, workspace_id, kind) VALUES ($1, $2, 'flow')") + .bind(job_id) + .bind(WS) + .execute(&db) + .await?; + sqlx::query( + "INSERT INTO v2_job_queue (id, workspace_id, scheduled_for) VALUES ($1, $2, now())", + ) + .bind(job_id) + .bind(WS) + .execute(&db) + .await?; + sqlx::query( + "INSERT INTO flow_conversation (id, workspace_id, flow_path, created_by) + VALUES ($1, $2, 'f/flow', 'test-user')", + ) + .bind(CONV) + .bind(WS) + .execute(&db) + .await?; + let user_seq: i64 = sqlx::query_scalar( + "INSERT INTO flow_conversation_message (conversation_id, message_type, content, job_id) + VALUES ($1, 'user', 'hi', $2) RETURNING created_seq", + ) + .bind(CONV) + .bind(job_id) + .fetch_one(&db) + .await?; + + let running = running_turns(&db, &[CONV]).await?; + let turn = running + .get(&CONV) + .expect("the queued run is the conversation's running turn"); + assert_eq!((turn.job_id, turn.user_seq), (job_id, user_seq)); + + match take_conversation(&db).await { + Err(Error::Generic(status, body)) => { + assert_eq!(status.as_u16(), 409); + assert!( + body.contains(&job_id.to_string()), + "the refusal names the running job: {body}" + ); + } + other => panic!("expected a 409 while the first turn runs, got {other:?}"), + } + + // The run ends: it leaves the queue, and the conversation takes the next message. + sqlx::query("DELETE FROM v2_job_queue WHERE id = $1") + .bind(job_id) + .execute(&db) + .await?; + assert!(running_turns(&db, &[CONV]).await?.is_empty()); + assert_eq!(take_conversation(&db).await?, CONV); + Ok(()) +} diff --git a/backend/windmill-api-flow-conversations/src/lib.rs b/backend/windmill-api-flow-conversations/src/lib.rs index 0ec295c23f..7320721ea5 100644 --- a/backend/windmill-api-flow-conversations/src/lib.rs +++ b/backend/windmill-api-flow-conversations/src/lib.rs @@ -14,7 +14,7 @@ pub use windmill_common::flow_conversations::FlowConversation; use windmill_common::{ db::{UserDB, DB}, error::{JsonResult, Result}, - flow_conversations::MessageType, + flow_conversations::{running_turns, MessageType, RunningTurn}, utils::{not_found_if_none, paginate, truncate_with_ellipsis, Pagination}, }; @@ -81,7 +81,7 @@ async fn list_conversations( Path(w_id): Path, Query(pagination): Query, Query(query): Query, -) -> JsonResult> { +) -> JsonResult> { let (per_page, offset) = paginate(pagination); let mut tx = user_db.clone().begin(&authed).await?; @@ -123,9 +123,27 @@ async fn list_conversations( let conversations = sqlx::query_as::(&sql) .fetch_all(&mut *tx) .await?; + let ids: Vec = conversations.iter().map(|c| c.id).collect(); + let mut running = running_turns(&mut *tx, &ids).await?; tx.commit().await?; - Ok(Json(conversations)) + Ok(Json( + conversations + .into_iter() + .map(|conversation| ListedConversation { + running_turn: running.remove(&conversation.id), + conversation, + }) + .collect(), + )) +} + +#[derive(Serialize)] +pub struct ListedConversation { + #[serde(flatten)] + pub conversation: FlowConversation, + /// Lets a chat that opens after a turn started follow it, and tell which chats are busy. + pub running_turn: Option, } async fn delete_conversation( diff --git a/backend/windmill-api/openapi.yaml b/backend/windmill-api/openapi.yaml index 075ced45a9..00fb3e07de 100644 --- a/backend/windmill-api/openapi.yaml +++ b/backend/windmill-api/openapi.yaml @@ -11791,6 +11791,11 @@ paths: content: application/json: schema: {} + "409": + description: >- + Chat-enabled flow only: the conversation named by `memory_id` is still answering + a message. The body is JSON: `{ "error": string, "running_turn": { "job_id", + "user_seq" } }`. /w/{workspace}/jobs/run_wait_result/fv/{version}: post: @@ -11831,6 +11836,11 @@ paths: content: application/json: schema: {} + "409": + description: >- + Chat-enabled flow only: the conversation named by `memory_id` is still answering + a message. The body is JSON: `{ "error": string, "running_turn": { "job_id", + "user_seq" } }`. get: summary: run flow by version with GET and wait until completion @@ -11863,6 +11873,11 @@ paths: content: application/json: schema: {} + "409": + description: >- + Chat-enabled flow only: the conversation named by `memory_id` is still answering + a message. The body is JSON: `{ "error": string, "running_turn": { "job_id", + "user_seq" } }`. /w/{workspace}/jobs/run_and_stream/f/{path}: post: @@ -11904,6 +11919,11 @@ paths: text/event-stream: schema: type: string + "409": + description: >- + Chat-enabled flow only: the conversation named by `memory_id` is still answering + a message. The body is JSON: `{ "error": string, "running_turn": { "job_id", + "user_seq" } }`. get: summary: run flow by path with GET and stream updates via SSE @@ -11937,6 +11957,11 @@ paths: text/event-stream: schema: type: string + "409": + description: >- + Chat-enabled flow only: the conversation named by `memory_id` is still answering + a message. The body is JSON: `{ "error": string, "running_turn": { "job_id", + "user_seq" } }`. /w/{workspace}/jobs/run_and_stream/fv/{version}: post: @@ -11984,6 +12009,11 @@ paths: text/event-stream: schema: type: string + "409": + description: >- + Chat-enabled flow only: the conversation named by `memory_id` is still answering + a message. The body is JSON: `{ "error": string, "running_turn": { "job_id", + "user_seq" } }`. get: summary: run flow by version with GET and stream updates via SSE @@ -12023,6 +12053,11 @@ paths: text/event-stream: schema: type: string + "409": + description: >- + Chat-enabled flow only: the conversation named by `memory_id` is still answering + a message. The body is JSON: `{ "error": string, "running_turn": { "job_id", + "user_seq" } }`. /w/{workspace}/jobs/run_and_stream/p/{path}: post: @@ -15432,6 +15467,11 @@ paths: schema: type: string format: uuid + "409": + description: >- + Chat-enabled flow only: the conversation named by `memory_id` is still answering + a message. The body is JSON: `{ "error": string, "running_turn": { "job_id", + "user_seq" } }`. /w/{workspace}/jobs/run/fv/{version}: post: @@ -15489,6 +15529,11 @@ paths: schema: type: string format: uuid + "409": + description: >- + Chat-enabled flow only: the conversation named by `memory_id` is still answering + a message. The body is JSON: `{ "error": string, "running_turn": { "job_id", + "user_seq" } }`. /w/{workspace}/jobs/run/batch_rerun_jobs: post: @@ -15982,6 +16027,11 @@ paths: schema: type: string format: uuid + "409": + description: >- + Chat-enabled flow only: the conversation named by `memory_id` is still answering + a message. The body is JSON: `{ "error": string, "running_turn": { "job_id", + "user_seq" } }`. /w/{workspace}/jobs/run_wait_result/preview_flow: post: @@ -16011,6 +16061,11 @@ paths: content: application/json: schema: {} + "409": + description: >- + Chat-enabled flow only: the conversation named by `memory_id` is still answering + a message. The body is JSON: `{ "error": string, "running_turn": { "job_id", + "user_seq" } }`. /w/{workspace}/jobs/run/dynamic_select: post: @@ -28902,6 +28957,23 @@ components: is_test: type: boolean description: Started from the flow editor's test panel rather than a deployed run + running_turn: + type: object + nullable: true + description: >- + The turn the conversation is still answering, set by the list endpoint: its + newest user message, while the flow run it started is queued or running. A + run into this conversation is refused with 409 until the turn ends. + required: [job_id, user_seq] + properties: + job_id: + type: string + format: uuid + description: The flow run of the turn + user_seq: + type: integer + format: int64 + description: created_seq of the user message that started the turn FlowConversationMessage: type: object diff --git a/backend/windmill-api/src/tracing_init.rs b/backend/windmill-api/src/tracing_init.rs index d99a7e5907..2fb668a39c 100644 --- a/backend/windmill-api/src/tracing_init.rs +++ b/backend/windmill-api/src/tracing_init.rs @@ -48,7 +48,9 @@ impl OnResponse for MyOnResponse { let status = response.status().as_u16(); if response.status().is_success() || response.status().is_redirection() { tracing::info!(latency = latency, status = status, "response") - } else if response.status().as_u16() == 404 { + } else if status == 404 || status == 409 { + // A refused turn is as expected as a miss: the flow chat takes the turn that + // refused it and sends its message after it. tracing::warn!(latency = latency, status = status, "response") } else { tracing::error!(latency = latency, status = status, "response") diff --git a/backend/windmill-common/src/error.rs b/backend/windmill-common/src/error.rs index da9f5f3df2..c06defaab5 100644 --- a/backend/windmill-common/src/error.rs +++ b/backend/windmill-common/src/error.rs @@ -303,7 +303,11 @@ impl IntoResponse for Error { let e = &self; - if matches!(status, axum::http::StatusCode::NOT_FOUND) { + // A refused turn is as expected as a miss: the chat follows the turn that refused it. + if matches!( + status, + axum::http::StatusCode::NOT_FOUND | axum::http::StatusCode::CONFLICT + ) { tracing::warn!(message = e.to_string()); } else { tracing::error!(message = e.to_string(), error = ?e); diff --git a/backend/windmill-common/src/flow_conversations.rs b/backend/windmill-common/src/flow_conversations.rs index 91bfa4a39f..5ee99d0681 100644 --- a/backend/windmill-common/src/flow_conversations.rs +++ b/backend/windmill-common/src/flow_conversations.rs @@ -69,7 +69,9 @@ pub async fn get_or_create_conversation_with_id( is_test: bool, ) -> Result { if let Some(existing) = lock_conversation(tx, w_id, conversation_id).await? { - return same_kind(existing, is_test); + let existing = same_kind(existing, is_test)?; + refuse_running_turn(tx, conversation_id).await?; + return Ok(existing); } // Truncate title to 25 characters max @@ -104,7 +106,74 @@ pub async fn get_or_create_conversation_with_id( "conversation {conversation_id} belongs to another workspace" )) })?; - same_kind(existing, is_test) + let existing = same_kind(existing, is_test)?; + refuse_running_turn(tx, conversation_id).await?; + Ok(existing) +} + +/// The turn a conversation is still answering: its newest user message, while the flow run +/// that message started is still queued or running. +#[derive(Serialize, Debug, Clone, Copy)] +pub struct RunningTurn { + pub job_id: Uuid, + /// `created_seq` of the user message that started the turn. + pub user_seq: i64, +} + +/// One running turn per conversation holds its agent memory; a second run would write the +/// same memory concurrently. Checked under the conversation's row lock, so two runs sent at +/// once cannot both pass: the second waits, then sees the first's message and queued job. +async fn refuse_running_turn( + tx: &mut sqlx::Transaction<'_, sqlx::Postgres>, + conversation_id: Uuid, +) -> Result<()> { + let Some(turn) = running_turns(&mut **tx, &[conversation_id]) + .await? + .remove(&conversation_id) + else { + return Ok(()); + }; + // A JSON body, so a chat client can follow the running turn instead of failing. + Err(crate::error::Error::Generic( + axum::http::StatusCode::CONFLICT, + serde_json::json!({ + "error": "this conversation is still answering a message; wait for it to finish or stop it before sending another", + "running_turn": turn, + }) + .to_string(), + )) +} + +/// The running turn of each of `conversation_ids` that has one. +/// +/// It answers for whatever ids it is given and checks no permission of its own, so the +/// executor must be one the caller is entitled to read those conversations through: a +/// `user_db` transaction under RLS, or a transaction holding ids the caller has already +/// authorized. Handed a raw pool and ids from a request, it would report other users' jobs. +pub async fn running_turns<'e, E: sqlx::PgExecutor<'e>>( + executor: E, + conversation_ids: &[Uuid], +) -> Result> { + let rows = sqlx::query!( + r#"SELECT c.id AS "id!", u.job_id AS "job_id!", u.created_seq AS "user_seq!" + FROM unnest($1::uuid[]) AS c(id) + CROSS JOIN LATERAL ( + SELECT job_id, created_seq + FROM flow_conversation_message + WHERE conversation_id = c.id AND message_type = 'user' + ORDER BY created_seq DESC + LIMIT 1 + ) u + WHERE u.job_id IS NOT NULL + AND EXISTS (SELECT 1 FROM v2_job_queue q WHERE q.id = u.job_id)"#, + conversation_ids + ) + .fetch_all(executor) + .await?; + Ok(rows + .into_iter() + .map(|r| (r.id, RunningTurn { job_id: r.job_id, user_seq: r.user_seq })) + .collect()) } /// `memory_id` is the caller's to choose, so a preview run could name a deployed diff --git a/chat-sdk/README.md b/chat-sdk/README.md index 1d09b68265..62a5467a47 100644 --- a/chat-sdk/README.md +++ b/chat-sdk/README.md @@ -229,7 +229,8 @@ set) means the turn could not run or be followed at all, such as a refused reque Methods: `sendMessage(text, { inputs?, attachments?, attachmentsInput? })`, `stop()`, `newConversation()`, `selectConversation(id)`, `loadConversations({ page?, perPage?, kind? })`, `deleteConversation(id)`, `renameConversation(id, title)`, `loadOlderMessages()`, -`destroy()`. `kind` lists the flow editor's test chats (`'test'`), the deployed flow's +`resumeTurn(turn)`, `refreshMessages()` (reads what another tab added to the open +conversation), `destroy()`. `kind` lists the flow editor's test chats (`'test'`), the deployed flow's own (`'deployed'`, the server's default) or both (`'all'`); each `Conversation` carries `isTest`. A rename keeps the conversation's place in the list. Switching conversations stops following the current answer; the flow keeps running and, with server history, @@ -261,6 +262,27 @@ user needs read and write on `windmill_uploads/*`, which the default rules grant `job_helpers`, so a restricted token needs `job_helpers:write`; a sandboxed raw app cannot request that scope today, so attachments are not available there yet. +## One turn at a time + +A conversation answers one message at a time, and Windmill enforces it: a message sent +while its previous turn still runs (from another tab, or before a reload) is refused. +`sendMessage` then rejects with a `TurnRunningError` and shows nothing of the message. +With server history, a listed `Conversation` also carries `runningTurn` while it is +answering. Either way, `resumeTurn(turn)` follows that turn in the selected +conversation: its answer streams in from the start and the turn finishes as if it had +been sent here, after which the message can be sent again. + +```ts +// `runningTurn` comes from the listing, so read it before selecting. +await chat.loadConversations() +const running = chat.getState().conversations.find((c) => c.id === id)?.runningTurn +await chat.selectConversation(id) +if (running) await chat.resumeTurn(running) +``` + +A dropped connection to the answer is retried; when it keeps failing, the chat stops +streaming and waits for the flow's result instead, so the turn still ends. + ## History Windmill stores every conversation of a chat-mode flow, and each Windmill user sees diff --git a/chat-sdk/src/ai-sdk.ts b/chat-sdk/src/ai-sdk.ts index d93bb99e2f..0ff422562b 100644 --- a/chat-sdk/src/ai-sdk.ts +++ b/chat-sdk/src/ai-sdk.ts @@ -169,9 +169,19 @@ function chunkStream( for (const e of event.events) parts.apply(e) continue } + // A result polled after the stream failed: the round that was streaming stopped + // wherever the connection did. Sent chunks cannot be taken back, so the answer's + // rest follows them, or the whole answer when it does not continue them. + const cut = parts.openText ?? '' parts.closeOpen() failure = await failureText(api, entry.jobId, event.result, signal) - if (failure === undefined && !parts.streamedText) { + if (failure === undefined && event.streamLost) { + const answer = extractChatAnswer(event.result) + if (answer !== undefined && answer !== cut) { + parts.text(answer.startsWith(cut) ? answer.slice(cut.length) : answer) + parts.closeOpen() + } + } else if (failure === undefined && !parts.streamedText) { // No agent streamed: the flow's result is the answer. const answer = extractChatAnswer(event.result) if (answer !== undefined) { @@ -216,6 +226,8 @@ async function failureText( */ class PartWriter { streamedText = false + /** The text of the round still streaming, until a tool call or the end closes it. */ + openText: string | undefined #textId: string | undefined #reasoningId: string | undefined #started = new Set() @@ -227,8 +239,10 @@ class PartWriter { this.streamedText = true if (!this.#textId) { this.#textId = randomId() + this.openText = '' this.emit({ type: 'text-start', id: this.#textId }) } + this.openText += delta this.emit({ type: 'text-delta', id: this.#textId, delta }) } @@ -248,6 +262,7 @@ class PartWriter { if (this.#textId) { this.emit({ type: 'text-end', id: this.#textId }) this.#textId = undefined + this.openText = undefined } } diff --git a/chat-sdk/src/api.ts b/chat-sdk/src/api.ts index a14e9800a1..ed97832aae 100644 --- a/chat-sdk/src/api.ts +++ b/chat-sdk/src/api.ts @@ -13,13 +13,53 @@ export interface WindmillChatApiOptions { export class WindmillApiError extends Error { constructor( message: string, - readonly status: number + 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 @@ -30,6 +70,8 @@ export interface FlowConversation { 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 } /** @@ -111,18 +153,27 @@ export class WindmillChatApi { this.#pollDelayMs = options.pollDelayMs } - /** Starts a turn: runs the flow with `memory_id` set to the conversation id. Returns the job id. */ + /** + * 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, options: { memoryId: string; signal?: AbortSignal } ): Promise { - const res = await this.#request(`jobs/run/f/${encodePath(flowPath)}`, { - method: 'POST', - query: { memory_id: options.memoryId, skip_preprocessor: 'true' }, - body: args, - signal: options.signal - }) + 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() } @@ -146,7 +197,9 @@ export class WindmillChatApi { signal: options.signal }) if (!res.body) { - throw new WindmillApiError('The job update stream has no body', res.status) + // 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 { @@ -294,7 +347,8 @@ export class WindmillChatApi { const text = await res.text().catch(() => '') throw new WindmillApiError( `${init.method ?? 'GET'} ${path} failed (${res.status})${text ? `: ${text}` : ''}`, - res.status + res.status, + text ) } return res diff --git a/chat-sdk/src/chat.ts b/chat-sdk/src/chat.ts index 2e2c52f3df..dfef438f33 100644 --- a/chat-sdk/src/chat.ts +++ b/chat-sdk/src/chat.ts @@ -1,9 +1,11 @@ import { + TurnRunningError, WindmillApiError, WindmillChatApi, type ConversationKind, type FlowConversation, - type FlowConversationMessage + type FlowConversationMessage, + type RunningTurn } from './api' import { resolveConfig, type ResolvedConfig } from './config' import { followJob } from './follow' @@ -38,6 +40,8 @@ const PERSIST_DEBOUNCE_MS = 250 /** Messages persist from spawned tasks that can land just after the flow completes. */ const RECONCILE_ATTEMPTS = 3 const RECONCILE_DELAY_MS = 400 +/** Rows per read of what a turn wrote; a fuller page is read on from its last row. */ +const ROWS_PAGE = 100 interface Turn { controller: AbortController @@ -73,6 +77,11 @@ class ChatImpl implements Chat { /** The kind the caller last listed, so the refresh after a new turn lists the same rows. */ #conversationKind: ConversationKind | undefined #persistTimer: ReturnType | undefined + /** Settles once the selected conversation's first page has been read. */ + #selecting: Promise = Promise.resolve() + /** Bumped whenever a turn is taken or the conversation is left: a read that started before + * it describes a conversation this chat has moved past. */ + #epoch = 0 constructor(options: ChatOptions) { this.#config = resolveConfig(options) @@ -141,6 +150,7 @@ class ChatImpl implements Chat { streamedText: false } this.#turn = turn + this.#epoch++ const timestamp = now() const conversation: Conversation = this.#state.conversations.find( @@ -186,6 +196,109 @@ class ChatImpl implements Chat { turn.jobId = this.#config.run ? await this.#config.run(args, context) : await this.#api.runFlow(this.#config.flowPath, args, context) + } catch (e) { + try { + // Nothing of this message reached the server, whether the conversation was already + // answering one sent elsewhere or an attachment never uploaded: the message is taken + // back rather than shown as a failed turn, and the caller gets the reason. + if (e instanceof TurnRunningError || !turn.started) { + this.#withdrawTurn(turn) + throw e + } + // stop() and a conversation switch abort the turn and settle the state themselves. + if (turn.controller.signal.aborted || isAbortError(e)) return + this.#failTurn(turn, e) + return + } finally { + if (this.#turn === turn) this.#turn = undefined + } + } + await this.#followTurn(turn, isNew) + } + + resumeTurn = async ({ jobId, userSeq }: RunningTurn): Promise => { + const conversationId = this.#state.conversationId + if (!conversationId || this.#turn) return + if (this.#state.history !== 'server') { + if (this.#state.status === 'submitted') this.#set({ status: 'idle' }) + return + } + const turn: Turn = { + controller: new AbortController(), + conversationId, + userMessageId: '', + jobId, + // Its run was asked for elsewhere and is already going: there is nothing to withdraw, + // and the conversation it belongs to is listed. + isNew: false, + started: true, + withdrawn: false, + streamedText: false + } + this.#turn = turn + this.#epoch++ + this.#set({ status: 'submitted', error: undefined }) + try { + // The conversation's first page may still be on its way; it would land over the turn. + await this.#selecting + if (!this.#turnActive(turn)) return + // A listing is a snapshot: the turn it named can have ended and another one started + // since. The newest user message this chat holds names the turn to follow instead — + // following the one the listing named would drop the newer turn's rows and leave the + // chat idle while it runs. + const newest = this.#latestUserAfter(userSeq) + if (newest) { + if (!newest.jobId) { + // Its row is written and its run is not named yet: there is nothing to follow, so + // it is left to the next listing to name the turn that runs now. + if (this.#turn === turn) this.#turn = undefined + this.#set({ status: 'idle' }) + return + } + turn.jobId = newest.jobId + userSeq = newest.seq! + } + // The stream replays the turn from its start, so what this chat shows of it goes and + // comes back as it replays: its rows, and what an earlier follow of it streamed or + // failed with, its own unread question included. The message that started it stays: + // it is the turn's anchor, and a long turn can have pushed it off the page this chat + // opened on. An earlier turn's answer that only its flow result gave has no row and so + // no seq; it names that turn's job, and it stays wherever the fallback put it. + const held = this.#state.messages + const at = held.findIndex((m) => m.seq === userSeq) + let user = at === -1 ? undefined : held[at] + let messages = held.filter((m, i) => { + if (m.seq !== undefined) return m.seq < userSeq + if (m.pending || m.jobId === turn.jobId) return false + return m.jobId !== undefined || (at !== -1 && i < at) + }) + if (user) messages = [...messages, user] + if (!user) { + const [row] = await this.#api.listMessages(conversationId, { + afterSeq: userSeq - 1, + perPage: 1, + signal: turn.controller.signal + }) + if (!this.#turnActive(turn)) return + if (row?.created_seq !== userSeq || row.message_type !== 'user') { + throw new Error('windmill-chat: the message that started the running turn is gone') + } + user = fromRow(row) + messages = [...messages, user] + } + turn.userMessageId = user.id + this.#set({ messages }) + } catch (e) { + if (!(turn.controller.signal.aborted || isAbortError(e))) this.#failTurn(turn, e) + if (this.#turn === turn) this.#turn = undefined + return + } + await this.#followTurn(turn, false) + } + + async #followTurn(turn: Turn, isNew: boolean): Promise { + let nextTurn: RunningTurn | undefined + try { const stopPolling = this.#state.history === 'server' ? this.#startPolling(turn) : () => {} let result: unknown try { @@ -193,7 +306,7 @@ class ChatImpl implements Chat { } finally { stopPolling() } - await this.#finishTurn(turn, result, isNew) + nextTurn = await this.#finishTurn(turn, result, isNew) } catch (e) { if (!turn.started) { // Nothing ran: the message is withdrawn rather than shown as a failed turn, and the @@ -208,6 +321,7 @@ class ChatImpl implements Chat { } finally { if (this.#turn === turn) this.#turn = undefined } + if (nextTurn) await this.resumeTurn(nextTurn) } stop = async (): Promise => { @@ -249,8 +363,14 @@ class ChatImpl implements Chat { }) } - selectConversation = async (conversationId: string): Promise => { - if (conversationId === this.#state.conversationId) return + selectConversation = (conversationId: string): Promise => { + if (conversationId === this.#state.conversationId) return this.#selecting + const selecting = this.#select(conversationId) + this.#selecting = selecting.catch(() => {}) + return selecting + } + + async #select(conversationId: string): Promise { this.#leaveConversation() this.#page = 1 this.#set({ @@ -379,9 +499,19 @@ class ChatImpl implements Chat { }) if (this.#state.conversationId !== conversationId) return const known = new Set(this.#state.messages.map((m) => m.serverId ?? m.id)) + // Pages count back from the newest row, so rows written since the first page shift + // newer rows into this one; and a resumed turn drops the rows it replays. Only what + // is older than everything held belongs above it. + const oldest = this.#state.messages.reduce( + (min, m) => (m.seq !== undefined && (min === undefined || m.seq < min) ? m.seq : min), + undefined + ) + const older = rows + .map(fromRow) + .filter((m) => !known.has(m.id) && (oldest === undefined || m.seq! < oldest)) this.#page = page this.#set({ - messages: [...rows.map(fromRow).filter((m) => !known.has(m.id)), ...this.#state.messages], + messages: [...older, ...this.#state.messages], hasMoreMessages: rows.length === this.#config.pageSize }) } finally { @@ -389,6 +519,56 @@ class ChatImpl implements Chat { } } + refreshMessages = async (): Promise => { + const conversationId = this.#state.conversationId + if ( + !conversationId || + this.#turn || + this.#state.history !== 'server' || + this.#state.loadingMessages + ) { + return + } + const newestBefore = latestSeq(this.#state.messages) + const failure = lastFailureShown(this.#state.messages) + // A turn that starts and ends while these rows are read, or the conversation being left + // and opened again, leaves the chat holding messages this read knows nothing about, and + // the rows would land under them rather than in their own place. They are left to the + // next read, which sees the conversation as it is now. + const epoch = this.#epoch + await this.#syncFromServer(conversationId, () => this.#epoch === epoch) + if (this.#epoch !== epoch) return + if (!failure || this.#state.conversationId !== conversationId) return + if (!this.#state.messages.some((m) => m.id === failure.id)) return + // The turn this chat lost may have been carried to its end elsewhere, which only a row + // this read brought can say, and only one of that turn's own: an agent writes its answer + // from a task the run does not wait for, so an earlier turn's answer can commit after + // this turn's question and would otherwise settle it with someone else's answer. The + // failed turn's job is on the message it left; unknown jobs accept the row, as everywhere + // else the turn's jobs are read. + const jobs = failure.jobId ? await this.#turnJobIds(failure.jobId) : undefined + if (this.#state.conversationId !== conversationId) return + const held = this.#state.messages + if (!held.some((m) => m.id === failure.id)) return + const answered = held.some( + (m) => + m.seq !== undefined && + m.seq > newestBefore && + m.role === 'assistant' && + m.success && + (jobs === undefined || m.jobId === undefined || jobs.has(m.jobId)) + ) + if (!answered) return + // A turn that ran, or failed, while those jobs were read owns the state now: this one's + // failure still goes, since its answer is here, but that turn's outcome stands. + const settles = + this.#state.status === 'error' && lastFailureShown(held)?.id === failure.id + this.#set({ + messages: held.filter((m) => m.id !== failure.id), + ...(settles ? { status: 'idle' as const, error: undefined } : {}) + }) + } + destroy = (): void => { this.#leaveConversation() } @@ -398,7 +578,10 @@ class ChatImpl implements Chat { async #follow(turn: Turn, onStreamStart: () => void): Promise { let started = false for await (const event of followJob(this.#api, turn.jobId!, { signal: turn.controller.signal })) { - if (event.type === 'completed') return event.result + if (event.type === 'completed') { + if (event.streamLost) this.#dropCutRound(turn) + return event.result + } if (!started) { started = true // Persisted rows for the streaming step would duplicate what is streaming. @@ -409,6 +592,19 @@ class ChatImpl implements Chat { throw new Error('windmill-chat: the job stream ended before the flow completed') } + /** + * The stream failed while a round's text was arriving, so that text stops wherever the + * connection did. It goes: the persisted rows or, without them, the flow result give + * the whole answer instead. Rounds a tool call closed were complete and stay. + */ + #dropCutRound(turn: Turn): void { + if (!this.#turnActive(turn)) return + const cut = turn.assistantId + turn.assistantId = undefined + turn.streamedText = false + if (cut) this.#set({ messages: this.#state.messages.filter((m) => m.id !== cut) }) + } + #applyEvents(turn: Turn, events: AgentStreamEvent[]): void { if (events.length === 0 || !this.#turnActive(turn)) return let messages = [...this.#state.messages] @@ -506,18 +702,20 @@ class ChatImpl implements Chat { this.#set({ messages, status: 'streaming' }) } - async #finishTurn(turn: Turn, result: unknown, isNew: boolean): Promise { + async #finishTurn(turn: Turn, result: unknown, isNew: boolean): Promise { if (!this.#turnActive(turn)) return if (this.#state.history === 'server') { - turn.jobIds = await this.#turnJobIds(turn) + turn.jobIds = await this.#turnJobIds(turn.jobId!, turn.controller.signal) if (!this.#turnActive(turn)) return const reconciled = await this.#reconcileTurn(turn) if (!this.#turnActive(turn)) return if (reconciled) { - this.#set({ status: 'idle' }) + const nextTurn = this.#nextTurnAfter(turn) + this.#set({ status: nextTurn ? 'submitted' : 'idle' }) this.#config.onFinish?.({ conversationId: turn.conversationId, jobId: turn.jobId, messages: this.#state.messages }) if (isNew) await this.loadConversations().catch(() => {}) - return + if (!this.#turnActive(turn)) return + return nextTurn } // Server history just proved unreadable: the turn completes as local history. } @@ -540,9 +738,13 @@ class ChatImpl implements Chat { messages = [...messages, assistantMessage(answer, true, turn.jobId)] } } - this.#set({ messages: finalized(messages), status: 'idle' }) + const nextTurn = this.#state.history === 'server' ? this.#nextTurnAfter(turn) : undefined + this.#set({ messages: finalized(messages), status: nextTurn ? 'submitted' : 'idle' }) this.#persistLocal() this.#config.onFinish?.({ conversationId: turn.conversationId, jobId: turn.jobId, messages: this.#state.messages }) + // `onFinish` can switch conversation, and the next turn belongs to this one. + if (!this.#turnActive(turn)) return + return nextTurn } /** @@ -561,11 +763,7 @@ class ChatImpl implements Chat { for (let attempt = 1; attempt <= RECONCILE_ATTEMPTS; attempt++) { let rows: FlowConversationMessage[] try { - rows = await this.#api.listMessages(turn.conversationId, { - afterSeq: this.#lastSeq(), - perPage: 100, - signal: turn.controller.signal - }) + rows = await this.#rowsAfterLastSeq(turn.conversationId, turn.controller.signal) } catch (e) { if (isAbortError(e)) throw e if (this.#fallBackToLocal(e)) return false @@ -592,22 +790,26 @@ class ChatImpl implements Chat { * a badly delayed one can invert that order at the cost of the reconcile * retries). The content is not compared with the flow result: an image answer, a * structured one and a forwarded agent result are all persisted in a shape the - * result does not reproduce. Rows carrying a job id belong to the turn when the - * job is one of the turn's, which leaves out an earlier turn whose job outlived - * `stop()` (a token without `jobs:write` cannot cancel it); a tool row without one - * (an MCP call runs inside the agent step) belongs to whatever turn is under way. + * result does not reproduce. A row belongs to the turn when it was created after + * the user message and, when it carries a job id, the job is one of the turn's. The + * server refuses a turn while the previous run is queued, but an agent writes its + * answer row from a task the run does not wait for, so that row can still land after + * the next user message; its job says whose it is. A tool row without a job (an MCP + * call runs inside the agent step) belongs to the turn under way. Until the user + * message's own row has been read, its position in the list stands in for its seq. */ #answered(turn: Turn): boolean { const messages = this.#state.messages const from = messages.findIndex((m) => m.id === turn.userMessageId) + const userSeq = messages[from]?.seq const ownJob = (m: ChatMessage) => turn.jobIds === undefined || (m.jobId === undefined ? m.role === 'tool' : turn.jobIds.has(m.jobId)) let latest: ChatMessage | undefined - for (let i = from + 1; i < messages.length; i++) { - const m = messages[i] - if (m.seq === undefined || m.role === 'user' || !ownJob(m)) continue + messages.forEach((m, i) => { + if (m.seq === undefined || m.role === 'user' || !ownJob(m)) return + if (userSeq !== undefined ? m.seq <= userSeq : i <= from) return if (latest === undefined || m.seq > latest.seq!) latest = m - } + }) return latest?.role === 'assistant' } @@ -615,22 +817,29 @@ class ChatImpl implements Chat { * The flow job plus every step job it ran, the failure and preprocessor steps * included (a failure handler's answer is persisted under its own job), and the * jobs an agent step's tool calls ran as (a tool row is persisted under its own - * job too). Unknown when the read fails. + * job too). A failed read is retried like the rows are; unknown when it keeps + * failing or the credential may not read jobs. Unknown accepts every row after the + * question: refusing them would leave a token without job access with no turn ever + * answered, each one finished a second time from its result. */ - async #turnJobIds(turn: Turn): Promise | undefined> { - try { - const job = await this.#api.getFlowJob(turn.jobId!, turn.controller.signal) - const ids = new Set([turn.jobId!]) - const status = job.flow_status - for (const m of [...(status?.modules ?? []), status?.failure_module, status?.preprocessor_module]) { - if (m?.job) ids.add(m.job) - for (const j of m?.flow_jobs ?? []) ids.add(j) - for (const a of m?.agent_actions ?? []) if (a.job_id) ids.add(a.job_id) + async #turnJobIds(jobId: string, signal?: AbortSignal): Promise | undefined> { + for (let attempt = 1; ; attempt++) { + try { + const job = await this.#api.getFlowJob(jobId, signal) + const ids = new Set([jobId]) + const status = job.flow_status + for (const m of [...(status?.modules ?? []), status?.failure_module, status?.preprocessor_module]) { + if (m?.job) ids.add(m.job) + for (const j of m?.flow_jobs ?? []) ids.add(j) + for (const a of m?.agent_actions ?? []) if (a.job_id) ids.add(a.job_id) + } + return ids + } catch (e) { + if (isAbortError(e)) throw e + const refused = e instanceof WindmillApiError && e.status >= 400 && e.status < 500 + if (refused || attempt === RECONCILE_ATTEMPTS) return undefined + await sleep(RECONCILE_DELAY_MS, signal) } - return ids - } catch (e) { - if (isAbortError(e)) throw e - return undefined } } @@ -689,11 +898,7 @@ class ChatImpl implements Chat { } if (stopped) return try { - const rows = await this.#api.listMessages(turn.conversationId, { - afterSeq: this.#lastSeq(), - perPage: 100, - signal - }) + const rows = await this.#rowsAfterLastSeq(turn.conversationId, signal) if (!stopped && this.#turnActive(turn)) this.#mergeRows(rows) } catch { // transient; the completion reconciliation catches up @@ -706,12 +911,22 @@ class ChatImpl implements Chat { } } - async #syncFromServer(conversationId: string): Promise { - const rows = await this.#api.listMessages(conversationId, { - afterSeq: this.#lastSeq(), - perPage: 100 - }) + /** Every row created after the newest one held, however many pages that takes. */ + async #rowsAfterLastSeq(conversationId: string, signal?: AbortSignal): Promise { + const rows: FlowConversationMessage[] = [] + let afterSeq = this.#lastSeq() + while (true) { + const page = await this.#api.listMessages(conversationId, { afterSeq, perPage: ROWS_PAGE, signal }) + rows.push(...page) + if (page.length < ROWS_PAGE) return rows + afterSeq = page[page.length - 1].created_seq + } + } + + async #syncFromServer(conversationId: string, stillCurrent?: () => boolean): Promise { + const rows = await this.#rowsAfterLastSeq(conversationId) if (this.#turn || this.#state.conversationId !== conversationId) return + if (stillCurrent && !stillCurrent()) return this.#mergeRows(rows) this.#set({ messages: finalized(this.#state.messages) }) } @@ -828,6 +1043,21 @@ class ChatImpl implements Chat { return this.#turn === turn && this.#state.conversationId === turn.conversationId } + #latestUserAfter(userSeq: number): ChatMessage | undefined { + return this.#state.messages.reduce( + (found, message) => + message.role === 'user' && message.seq !== undefined && message.seq > (found?.seq ?? userSeq) ? message : found, + undefined + ) + } + + #nextTurnAfter(turn: Turn): RunningTurn | undefined { + const userSeq = this.#state.messages.find((message) => message.id === turn.userMessageId)?.seq + if (userSeq === undefined) return undefined + const next = this.#latestUserAfter(userSeq) + return next?.jobId && next.seq !== undefined ? { jobId: next.jobId, userSeq: next.seq } : undefined + } + /** Stops following the current answer; the flow itself keeps running. */ #detachTurn(): void { const turn = this.#turn @@ -842,6 +1072,7 @@ class ChatImpl implements Chat { * written out now rather than on the debounce that may never fire. */ #leaveConversation(): void { + this.#epoch++ if (this.#turn) { this.#detachTurn() this.#set({ messages: finalized(this.#state.messages), status: 'idle' }) @@ -859,6 +1090,20 @@ class ChatImpl implements Chat { } } +/** The failure this chat is showing: the message `#failTurn` left, which is never a row. */ +function lastFailureShown(messages: readonly ChatMessage[]): ChatMessage | undefined { + for (let i = messages.length - 1; i >= 0; i--) { + const m = messages[i] + if (m.seq === undefined && m.success === false) return m + } + return undefined +} + +/** The newest row seq a list holds; 0 when it holds none. */ +function latestSeq(messages: readonly ChatMessage[]): number { + return messages.reduce((newest, m) => (m.seq !== undefined && m.seq > newest ? m.seq : newest), 0) +} + function fromRow(row: FlowConversationMessage): ChatMessage { const toolName = row.message_type === 'tool' @@ -898,7 +1143,10 @@ function fromConversation(row: FlowConversation): Conversation { title: row.title ?? undefined, createdAt: row.created_at, updatedAt: row.updated_at, - isTest: row.is_test + isTest: row.is_test, + runningTurn: row.running_turn + ? { jobId: row.running_turn.job_id, userSeq: row.running_turn.user_seq } + : undefined } } diff --git a/chat-sdk/src/follow.ts b/chat-sdk/src/follow.ts index 50f101393e..6402b0d6f7 100644 --- a/chat-sdk/src/follow.ts +++ b/chat-sdk/src/follow.ts @@ -1,13 +1,18 @@ -import type { WindmillChatApi } from './api' +import { WindmillApiError, type WindmillChatApi } from './api' import { createStreamEventParser, type AgentStreamEvent } from './stream' -import { abortError, sleep } from './utils' +import { abortError, isAbortError, sleep } from './utils' const RECONNECT_DELAY_MS = 300 +const MAX_RECONNECT_DELAY_MS = 5000 +/** Consecutive failed connections before the job is polled instead. */ +const MAX_CONNECTION_FAILURES = 3 +const RESULT_POLL_MS = 2000 export type FollowEvent = /** Agent events decoded from the job's result stream; empty when a chunk ended mid-line. */ | { type: 'stream'; events: AgentStreamEvent[] } - | { type: 'completed'; result: unknown } + /** `streamLost`: the result was polled after the stream failed, so what streamed may stop short. */ + | { type: 'completed'; result: unknown; streamLost?: boolean } /** * Follows a job to completion across the server's stream timeouts: every @@ -18,6 +23,10 @@ export type FollowEvent = * The offset indexes the stream of one sub-job (`flow_stream_job_id`, the flow's * streaming step). A retried step gets a new one, so when the id changes the * offset is dropped and the connection reopened from that sub-job's start. + * + * A connection that fails (a proxy restarting, the network dropping) is retried with + * backoff. The run is still going, so after a few failures in a row the job's result is + * polled instead: the rest of the answer is not streamed, but the turn still ends. */ export async function* followJob( api: WindmillChatApi, @@ -27,45 +36,86 @@ export async function* followJob( let parser = createStreamEventParser() let offset = options.streamOffset let streamJobId: string | undefined - while (true) { + let failures = 0 + while (failures < MAX_CONNECTION_FAILURES) { let reopen = false - for await (const update of api.streamJob(jobId, { streamOffset: offset, signal: options.signal })) { - if (update.type === 'ping') continue - if (update.type === 'timeout') { - reopen = true - break - } - if (update.type === 'error') throw new Error(update.error) - if (update.type === 'notfound') throw new Error(`Job ${jobId} not found`) - if (update.flow_stream_job_id && update.flow_stream_job_id !== streamJobId) { - const switched = streamJobId !== undefined && offset !== undefined - streamJobId = update.flow_stream_job_id - if (switched) { - // This connection skipped the new sub-job's first chunks: start it over. - offset = undefined - options.onOffset?.(undefined) - parser = createStreamEventParser() + try { + for await (const update of api.streamJob(jobId, { streamOffset: offset, signal: options.signal })) { + // Opening a connection does not prove it is carrying the job. Only stream progress + // clears the count, or status snapshots followed by EOF would never reach polling. + if (update.type === 'ping') continue + if (update.type === 'timeout') { reopen = true break } + if (update.type === 'error') throw new Error(update.error) + if (update.type === 'notfound') throw new Error(`Job ${jobId} not found`) + if (update.flow_stream_job_id && update.flow_stream_job_id !== streamJobId) { + const switched = streamJobId !== undefined && offset !== undefined + streamJobId = update.flow_stream_job_id + if (switched) { + // This connection skipped the new sub-job's first chunks: start it over. + failures = 0 + offset = undefined + options.onOffset?.(undefined) + parser = createStreamEventParser() + reopen = true + break + } + } + if (update.stream_offset !== undefined) { + if (update.stream_offset !== offset) failures = 0 + offset = update.stream_offset + options.onOffset?.(offset) + } + if (update.new_result_stream) { + failures = 0 + yield { type: 'stream', events: parser.push(update.new_result_stream) } + } + if (update.completed) { + const rest = parser.flush() + if (rest.length > 0) yield { type: 'stream', events: rest } + yield { type: 'completed', result: update.only_result } + return + } } - if (update.stream_offset !== undefined) { - offset = update.stream_offset - options.onOffset?.(offset) - } - if (update.new_result_stream) { - yield { type: 'stream', events: parser.push(update.new_result_stream) } - } - if (update.completed) { - const rest = parser.flush() - if (rest.length > 0) yield { type: 'stream', events: rest } - yield { type: 'completed', result: update.only_result } - return - } + } catch (e) { + if (options.signal?.aborted || isAbortError(e) || !isConnectionFailure(e)) throw e + failures++ + if (failures >= MAX_CONNECTION_FAILURES) break + await sleep(Math.min(RECONNECT_DELAY_MS * 2 ** failures, MAX_RECONNECT_DELAY_MS), options.signal) + continue } if (options.signal?.aborted) throw abortError() - // The server closes the connection after its timeout; a dropped connection looks - // the same minus the event. Either way the offset lets the next one resume. - if (!reopen) await sleep(RECONNECT_DELAY_MS, options.signal) + // The server closes the connection after its timeout; a connection that ends without + // that event, and without the job completing, carried nothing to its end. The offset + // lets the next one resume, and it counts like a failed one so a gateway closing every + // stream this way still reaches the polling below rather than reconnecting for ever. + if (!reopen) { + failures++ + if (failures >= MAX_CONNECTION_FAILURES) break + await sleep(RECONNECT_DELAY_MS, options.signal) + } + } + while (true) { + await sleep(RESULT_POLL_MS, options.signal) + try { + const { completed, result } = await api.getCompletedResult(jobId, options.signal) + if (completed) { + yield { type: 'completed', result, streamLost: true } + return + } + } catch (e) { + if (options.signal?.aborted || isAbortError(e) || !isConnectionFailure(e)) throw e + } } } + +/** + * A failure that says nothing about the job: the request never reached Windmill, or a + * gateway in front of it answered. A 4xx from Windmill itself (not found, refused) does. + */ +function isConnectionFailure(e: unknown): boolean { + if (e instanceof WindmillApiError) return e.status >= 500 || e.status === 0 + return e instanceof TypeError +} diff --git a/chat-sdk/src/index.ts b/chat-sdk/src/index.ts index 9ed2ed0b7e..a050edce37 100644 --- a/chat-sdk/src/index.ts +++ b/chat-sdk/src/index.ts @@ -3,7 +3,10 @@ export { detectRawApp, type RawAppContext } from './config' export { WindmillChatApi, WindmillApiError, + TurnRunningError, + turnRunningError, readServerSentEvents, + type RunningTurn, type WindmillChatApiOptions, type ConversationKind, type FlowConversation, diff --git a/chat-sdk/src/react.ts b/chat-sdk/src/react.ts index 33193842f2..c32324ee18 100644 --- a/chat-sdk/src/react.ts +++ b/chat-sdk/src/react.ts @@ -6,6 +6,7 @@ export type UseWindmillChat = ChatState & Pick< Chat, | 'sendMessage' + | 'resumeTurn' | 'stop' | 'newConversation' | 'selectConversation' @@ -13,6 +14,7 @@ export type UseWindmillChat = ChatState & | 'deleteConversation' | 'renameConversation' | 'loadOlderMessages' + | 'refreshMessages' > & { chat: Chat } /** @@ -61,13 +63,15 @@ export function useWindmillChat(options: ChatOptions): UseWindmillChat { chat, sendMessage: (text, options) => chat.sendMessage(text, { ...options, inputs: { ...latest.current.inputs, ...options?.inputs } }), + resumeTurn: chat.resumeTurn, stop: chat.stop, newConversation: chat.newConversation, selectConversation: chat.selectConversation, loadConversations: chat.loadConversations, deleteConversation: chat.deleteConversation, renameConversation: chat.renameConversation, - loadOlderMessages: chat.loadOlderMessages + loadOlderMessages: chat.loadOlderMessages, + refreshMessages: chat.refreshMessages }), [state, chat] ) diff --git a/chat-sdk/src/types.ts b/chat-sdk/src/types.ts index 99e16fe304..180238e50c 100644 --- a/chat-sdk/src/types.ts +++ b/chat-sdk/src/types.ts @@ -1,3 +1,5 @@ +import type { RunningTurn } from './api' + export type ChatRole = 'user' | 'assistant' | 'tool' | 'system' /** @@ -61,6 +63,11 @@ export interface Conversation { * server has listed the conversation; unset for one only this client has seen. */ isTest?: boolean + /** + * The turn the conversation was still answering when the server listed it: started in + * another tab, or before this chat was created. Pass it to `resumeTurn` to follow it. + */ + runningTurn?: RunningTurn } export interface ChatState { @@ -164,8 +171,17 @@ export interface Chat { * Sends a message in the current conversation, starting one when there is none. Resolves * when the answer is complete. Rejects when the message could not be sent at all — a turn * already running, an attachment that failed to upload — without touching the transcript. + * A conversation still answering a message sent elsewhere rejects with `TurnRunningError`: + * follow that turn with `resumeTurn`, then send again. */ sendMessage(text: string, options?: SendMessageOptions): Promise + /** + * Follows a turn of the current conversation that this chat did not start, as named by + * `Conversation.runningTurn` or a `TurnRunningError`: its answer streams into `messages` + * from the start and the turn finishes like one sent here. Server history only; resolves + * when the answer is complete, at once when a turn is already being followed. + */ + resumeTurn(turn: RunningTurn): Promise /** Stops following the answer and asks Windmill to cancel the run. */ stop(): Promise newConversation(): void @@ -183,6 +199,11 @@ export interface Chat { /** Sets a conversation's title. The list keeps its order: only a turn moves a conversation. */ renameConversation(conversationId: string, title: string): Promise loadOlderMessages(): Promise + /** + * Reads what the current conversation gained since its newest message held here, such + * as a turn another tab ran. Server history only; does nothing while a turn is followed. + */ + refreshMessages(): Promise /** Stops background work (stream, polling) and writes local history out. The chat stays usable. */ destroy(): void } diff --git a/chat-sdk/test/ai-sdk-chat.test.ts b/chat-sdk/test/ai-sdk-chat.test.ts index a0bbb7b8a1..4cef282a82 100644 --- a/chat-sdk/test/ai-sdk-chat.test.ts +++ b/chat-sdk/test/ai-sdk-chat.test.ts @@ -59,4 +59,26 @@ describe('AI SDK Chat over the Windmill transport', () => { expect(chat.error?.message).toBe('boom') expect(memoryIds.length === 1 || calls.filter((c) => c.method === 'POST')[1].url.searchParams.get('memory_id') === memoryIds[0]).toBe(true) }) + + test('a text part cut by a lost stream is completed from the polled result', async () => { + let streams = 0 + const { fetch } = fetchMock( + (c) => (c.method === 'POST' && c.url.pathname === `/api/w/ws/jobs/run/f/${FLOW}` ? text('job-1') : undefined), + (c) => + c.url.pathname.endsWith('/getupdate_sse/job-1') + ? ++streams === 1 + ? sse([{ type: 'update', new_result_stream: ndjson({ type: 'token_delta', content: 'The ans' }), stream_offset: 1 }]) + : text('bad gateway', 502) + : undefined, + (c) => + c.url.pathname.endsWith('/get_result_maybe/job-1') + ? json({ completed: true, success: true, result: { output: 'The answer is 42', messages: [] } }) + : undefined + ) + const transport = createWindmillChatTransport({ baseUrl: 'http://wm.test', workspace: 'ws', flowPath: FLOW, token: 'tok', fetch }) + const chat = new Chat({ id: 'cut-chat', transport }) + await chat.sendMessage({ text: 'what is it?' }) + const texts = chat.messages[1].parts.filter((p) => p.type === 'text').map((p) => (p as { text: string }).text) + expect(texts.join('')).toBe('The answer is 42') + }, 15000) }) diff --git a/chat-sdk/test/chat.test.ts b/chat-sdk/test/chat.test.ts index a739abf3ad..3845c91ac7 100644 --- a/chat-sdk/test/chat.test.ts +++ b/chat-sdk/test/chat.test.ts @@ -1,4 +1,5 @@ import { describe, expect, test } from 'bun:test' +import { TurnRunningError } from '../src/api' import { createChat } from '../src/chat' import type { ChatOptions } from '../src/types' import { fetchMock, json, memoryStorage, messageRow, ndjson, sse, sseTimed, text, type Route } from './support' @@ -11,6 +12,15 @@ const run: Route = (c) => const streamPath = '/api/w/ws/jobs_u/getupdate_sse/job-1' +/** Waits for a condition the code under test must reach, and fails saying which one. */ +async function until(done: () => boolean, what: string, timeoutMs = 2000): Promise { + const deadline = Date.now() + timeoutMs + while (!done()) { + if (Date.now() > deadline) throw new Error(`timed out waiting for ${what}`) + await new Promise((resolve) => setTimeout(resolve, 0)) + } +} + function options(extra: Partial, fetch: ChatOptions['fetch']): ChatOptions { return { flowPath: FLOW, baseUrl: BASE, workspace: 'ws', fetch, storage: memoryStorage(), ...extra } } @@ -945,41 +955,752 @@ describe('createChat with server history', () => { ]) }) - test('a late answer from a stopped job is not taken as the next turn answer', async () => { - let jobs = 0 - let reads = 0 + test('a message refused because a turn is running leaves nothing behind and names that turn', async () => { const { fetch } = fetchMock( - (c) => (c.method === 'POST' && c.url.pathname.includes('/jobs/run/f/') ? text(`job-${++jobs}`) : undefined), - // job-1 never completes: the connection just ends, so the turn keeps waiting. - (c) => (c.url.pathname.endsWith('/getupdate_sse/job-1') ? sse([{ type: 'update' }]) : undefined), (c) => - c.url.pathname.endsWith('/getupdate_sse/job-2') - ? sse([{ type: 'update', completed: true, only_result: { windmill_chat_answer: 'second answer' } }]) + c.method === 'POST' && c.url.pathname === `/api/w/ws/jobs/run/f/${FLOW}` + ? text(JSON.stringify({ error: 'still answering', running_turn: { job_id: 'job-9', user_seq: 41 } }), 409) + : undefined, + (c) => (c.url.pathname.endsWith('/messages') ? json([messageRow(40, 'user', 'earlier')]) : undefined) + ) + const chat = createChat(options({}, fetch)) + await chat.selectConversation('conv') + const refused = await chat.sendMessage('again').catch((e) => e) + expect(refused).toBeInstanceOf(TurnRunningError) + expect(refused.turn).toEqual({ jobId: 'job-9', userSeq: 41 }) + const state = chat.getState() + expect(state.messages.map((m) => m.content)).toEqual(['earlier']) + expect(state.status).toBe('idle') + // The conversation is a real one, answering elsewhere: it stays listed, its message gone. + expect(state.conversations.map((c) => c.id)).toEqual(['conv']) + }) + + test('a listing that names a turn already over follows the one running now', async () => { + const { fetch, calls } = fetchMock( + (c) => + c.url.pathname === '/api/w/ws/jobs_u/getupdate_sse/job-2' + ? sse([ + { + type: 'update', + new_result_stream: ndjson({ type: 'token_delta', content: 'second answer' }), + stream_offset: 1, + completed: true, + only_result: { windmill_chat_answer: 'second answer' } + } + ]) : undefined, - // The run-only token cannot cancel: job-1 keeps running after stop(). - (c) => (c.url.pathname.includes('/queue/cancel/') ? text('forbidden', 400) : undefined), (c) => c.url.pathname.endsWith('/jobs_u/get/job-2') ? json({ flow_status: { modules: [{ job: 'step-2' }] } }) : undefined, - // Read 1 is stop()'s sync; the stopped job's answer lands after the second user row. + (c) => { + if (!c.url.pathname.endsWith('/messages')) return undefined + const after = c.url.searchParams.get('after_seq') + if (after === '52') return json([messageRow(53, 'assistant', 'second answer', { job_id: 'step-2' })]) + return json([ + messageRow(50, 'user', 'first'), + messageRow(51, 'assistant', 'first answer'), + messageRow(52, 'user', 'second', { job_id: 'job-2' }) + ]) + } + ) + const chat = createChat(options({}, fetch)) + await chat.selectConversation('conv') + // The listing named the first turn; it ended and the second one started before this select. + await chat.resumeTurn({ jobId: 'job-1', userSeq: 50 }) + expect(calls.some((c) => c.url.pathname.includes('getupdate_sse/job-1'))).toBe(false) + expect(calls.some((c) => c.url.pathname.includes('getupdate_sse/job-2'))).toBe(true) + expect(chat.getState().messages.map((m) => [m.serverId, m.content])).toEqual([ + ['row-50', 'first'], + ['row-51', 'first answer'], + ['row-52', 'second'], + ['row-53', 'second answer'] + ]) + expect(chat.getState().status).toBe('idle') + expect(chat.getState().error).toBeUndefined() + }) + + test('a newer message whose run is not named yet leaves the chat free', async () => { + const { fetch, calls } = fetchMock( (c) => c.url.pathname.endsWith('/messages') - ? json( - ++reads === 1 - ? [messageRow(71, 'user', 'first')] - : reads === 2 - ? [messageRow(72, 'user', 'second'), messageRow(73, 'assistant', 'first answer, late', { job_id: 'step-1' })] - : [messageRow(74, 'assistant', 'second answer', { job_id: 'step-2' })] - ) + ? json([messageRow(50, 'user', 'first'), messageRow(52, 'user', 'second')]) + : undefined + ) + const chat = createChat(options({}, fetch)) + await chat.selectConversation('conv') + await chat.resumeTurn({ jobId: 'job-1', userSeq: 50 }) + expect(chat.getState().status).toBe('idle') + expect(calls.some((c) => c.url.pathname.includes('getupdate_sse'))).toBe(false) + // A message sent now starts its own turn rather than being refused. + expect(chat.getState().error).toBeUndefined() + }) + + test('a turn started while a resumed turn reconciles is followed next', async () => { + const { fetch, calls } = fetchMock( + (c) => + c.url.pathname === streamPath + ? sse([{ type: 'update', completed: true, only_result: { windmill_chat_answer: 'first answer' } }]) + : undefined, + (c) => + c.url.pathname === '/api/w/ws/jobs_u/getupdate_sse/job-2' + ? sse([{ type: 'update', completed: true, only_result: { windmill_chat_answer: 'second answer' } }]) + : undefined, + (c) => + c.url.pathname.endsWith('/jobs_u/get/job-1') + ? json({ flow_status: { modules: [{ job: 'step-1' }] } }) + : undefined, + (c) => + c.url.pathname.endsWith('/jobs_u/get/job-2') + ? json({ flow_status: { modules: [{ job: 'step-2' }] } }) + : undefined, + (c) => { + if (!c.url.pathname.endsWith('/messages')) return undefined + const after = c.url.searchParams.get('after_seq') + if (after === null) return json([messageRow(50, 'user', 'first', { job_id: 'job-1' })]) + if (after === '50') { + return json([ + messageRow(51, 'assistant', 'first answer', { job_id: 'step-1' }), + messageRow(52, 'user', 'second', { job_id: 'job-2' }) + ]) + } + if (after === '52') return json([messageRow(53, 'assistant', 'second answer', { job_id: 'step-2' })]) + return json([]) + } + ) + const chat = createChat(options({}, fetch)) + const statuses: string[] = [] + chat.subscribe((state) => statuses.push(state.status)) + await chat.selectConversation('conv') + await chat.resumeTurn({ jobId: 'job-1', userSeq: 50 }) + expect(calls.some((c) => c.url.pathname === streamPath)).toBe(true) + expect(calls.some((c) => c.url.pathname.includes('getupdate_sse/job-2'))).toBe(true) + expect(chat.getState().messages.map((m) => [m.serverId, m.content])).toEqual([ + ['row-50', 'first'], + ['row-51', 'first answer'], + ['row-52', 'second'], + ['row-53', 'second answer'] + ]) + expect(chat.getState().status).toBe('idle') + // Vacuous without this: with no 'submitted' at all, both lookups are -1 and the slice empty. + expect(statuses).toContain('submitted') + expect(statuses.slice(statuses.indexOf('submitted'), statuses.lastIndexOf('submitted') + 1)).not.toContain('idle') + }) + + test('a fallback to local history during handoff does not leave the chat busy', async () => { + let messageFetches = 0 + const { fetch, calls } = fetchMock( + (c) => + c.url.pathname === streamPath + ? sse([{ type: 'update', completed: true, only_result: { windmill_chat_answer: 'first answer' } }]) + : undefined, + (c) => { + if (!c.url.pathname.endsWith('/messages')) return undefined + messageFetches++ + if (messageFetches === 1) return json([messageRow(50, 'user', 'first', { job_id: 'job-1' })]) + if (messageFetches === 2) return json([messageRow(52, 'user', 'second', { job_id: 'job-2' })]) + return text('forbidden', 403) + } + ) + const chat = createChat(options({}, fetch)) + await chat.selectConversation('conv') + await chat.resumeTurn({ jobId: 'job-1', userSeq: 50 }) + expect(chat.getState().history).toBe('local') + expect(chat.getState().status).toBe('idle') + expect(calls.some((c) => c.url.pathname.includes('getupdate_sse/job-2'))).toBe(false) + }) + + test('a conversation list fallback during handoff does not leave the chat busy', async () => { + const { fetch, calls } = fetchMock( + run, + (c) => + c.url.pathname === streamPath + ? sse([{ type: 'update', completed: true, only_result: { windmill_chat_answer: 'first answer' } }]) + : undefined, + (c) => + c.url.pathname.endsWith('/jobs_u/get/job-1') + ? json({ flow_status: { modules: [{ job: 'job-1' }] } }) + : undefined, + (c) => + c.url.pathname.endsWith('/messages') + ? json([ + messageRow(50, 'user', 'first', { job_id: 'job-1' }), + messageRow(51, 'assistant', 'first answer', { job_id: 'job-1' }), + messageRow(52, 'user', 'second', { job_id: 'job-2' }) + ]) + : undefined, + (c) => (c.url.pathname === '/api/w/ws/flow_conversations/list' ? text('forbidden', 403) : undefined) + ) + const chat = createChat(options({}, fetch)) + await chat.sendMessage('first') + expect(chat.getState().history).toBe('local') + expect(chat.getState().status).toBe('idle') + expect(calls.some((c) => c.url.pathname.includes('getupdate_sse/job-2'))).toBe(false) + }) + + test('a conversation switch during handoff does not resume the next turn elsewhere', async () => { + let releaseList!: (response: Response) => void + const listGate = new Promise((resolve) => { + releaseList = resolve + }) + const { fetch, calls } = fetchMock( + run, + (c) => + c.url.pathname === streamPath + ? sse([{ type: 'update', completed: true, only_result: { windmill_chat_answer: 'first answer' } }]) + : undefined, + (c) => + c.url.pathname.endsWith('/jobs_u/get/job-1') + ? json({ flow_status: { modules: [{ job: 'job-1' }] } }) + : undefined, + (c) => { + if (!c.url.pathname.endsWith('/messages')) return undefined + if (c.url.pathname.includes('/flow_conversations/other/messages')) return json([]) + return json([ + messageRow(50, 'user', 'first', { job_id: 'job-1' }), + messageRow(51, 'assistant', 'first answer', { job_id: 'job-1' }), + messageRow(52, 'user', 'second', { job_id: 'job-2' }) + ]) + }, + (c) => (c.url.pathname === '/api/w/ws/flow_conversations/list' ? listGate : undefined) + ) + const chat = createChat(options({}, fetch)) + const sent = chat.sendMessage('first') + await until(() => calls.some((c) => c.url.pathname === '/api/w/ws/flow_conversations/list'), 'the list to be asked for') + await chat.selectConversation('other') + releaseList(json([])) + await sent + expect(chat.getState().conversationId).toBe('other') + expect(chat.getState().status).toBe('idle') + expect(chat.getState().messages).toEqual([]) + expect(calls.some((c) => c.url.pathname.includes('getupdate_sse/job-2'))).toBe(false) + }) + + test('resuming a later turn keeps an earlier answer that only the flow result gave', async () => { + const { fetch } = fetchMock( + run, + (c) => + c.url.pathname === streamPath + ? sse([{ type: 'update', completed: true, only_result: { windmill_chat_answer: 'first answer' } }]) + : undefined, + (c) => + c.url.pathname === '/api/w/ws/jobs_u/getupdate_sse/job-2' + ? sse([{ type: 'update', completed: true, only_result: { windmill_chat_answer: 'second answer' } }]) + : undefined, + (c) => + c.url.pathname.endsWith('/jobs_u/get/job-1') + ? json({ flow_status: { modules: [{ job: 'step-1' }] } }) + : undefined, + (c) => + c.url.pathname.endsWith('/jobs_u/get/job-2') + ? json({ flow_status: { modules: [{ job: 'step-2' }] } }) + : undefined, + (c) => { + if (!c.url.pathname.endsWith('/messages')) return undefined + const after = c.url.searchParams.get('after_seq') + if (after === '51') return json([messageRow(52, 'user', 'second', { job_id: 'job-2' })]) + if (after === '52') return json([messageRow(53, 'assistant', 'second answer', { job_id: 'step-2' })]) + // The first answer never got a row: the turn finishes from the flow result. + return json([messageRow(50, 'user', 'first', { job_id: 'job-1' })]) + }, + (c) => (c.url.pathname === '/api/w/ws/flow_conversations/list' ? json([]) : undefined) + ) + const chat = createChat(options({}, fetch)) + await chat.sendMessage('first') + await chat.resumeTurn({ jobId: 'job-2', userSeq: 52 }) + expect(chat.getState().messages.map((m) => m.content)).toEqual([ + 'first', + 'first answer', + 'second', + 'second answer' + ]) + }) + + test('resuming a turn this chat lost drops what it had shown of that turn', async () => { + let streams = 0 + const { fetch } = fetchMock( + (c) => { + if (c.url.pathname !== streamPath) return undefined + return ++streams === 1 + ? sse([ + { type: 'update', new_result_stream: ndjson({ type: 'token_delta', content: 'partial' }), stream_offset: 1 }, + { type: 'error', error: 'stream broke' } + ]) + : sse([{ type: 'update', completed: true, only_result: { windmill_chat_answer: 'the answer' } }]) + }, + (c) => + c.url.pathname.endsWith('/jobs_u/get/job-1') + ? json({ flow_status: { modules: [{ job: 'step-1' }] } }) + : undefined, + (c) => { + if (!c.url.pathname.endsWith('/messages')) return undefined + if (c.url.searchParams.get('after_seq') === '50') { + return json([messageRow(51, 'assistant', 'the answer', { job_id: 'step-1' })]) + } + return json([messageRow(50, 'user', 'question', { job_id: 'job-1' })]) + } + ) + const chat = createChat(options({}, fetch)) + await chat.selectConversation('conv') + await chat.resumeTurn({ jobId: 'job-1', userSeq: 50 }) + expect(chat.getState().status).toBe('error') + // The run goes on; following it again replays the whole answer under its question. + await chat.resumeTurn({ jobId: 'job-1', userSeq: 50 }) + expect(chat.getState().messages.map((m) => m.content)).toEqual(['question', 'the answer']) + }) + + test('resuming a turn this chat sent and lost shows its question once', async () => { + let streams = 0 + const { fetch } = fetchMock( + run, + (c) => { + if (c.url.pathname !== streamPath) return undefined + return ++streams === 1 + ? sse([ + { type: 'update', new_result_stream: ndjson({ type: 'token_delta', content: 'partial' }), stream_offset: 1 }, + { type: 'error', error: 'stream broke' } + ]) + : sse([{ type: 'update', completed: true, only_result: { windmill_chat_answer: 'the answer' } }]) + }, + (c) => + c.url.pathname.endsWith('/jobs_u/get/job-1') + ? json({ flow_status: { modules: [{ job: 'step-1' }] } }) + : undefined, + (c) => { + if (!c.url.pathname.endsWith('/messages')) return undefined + const after = c.url.searchParams.get('after_seq') + if (after === '49') return json([messageRow(50, 'user', 'question', { job_id: 'job-1' })]) + if (after === '50') return json([messageRow(51, 'assistant', 'the answer', { job_id: 'step-1' })]) + // The stream fails before the question's row is read. + return json([]) + }, + (c) => (c.url.pathname === '/api/w/ws/flow_conversations/list' ? json([]) : undefined) + ) + const chat = createChat(options({}, fetch)) + await chat.sendMessage('question') + expect(chat.getState().status).toBe('error') + await chat.resumeTurn({ jobId: 'job-1', userSeq: 50 }) + expect(chat.getState().messages.map((m) => m.content)).toEqual(['question', 'the answer']) + }) + + test('a re-read that finds the lost turn answered clears the failure it was left with', async () => { + let answered = false + const { fetch } = fetchMock( + (c) => + c.url.pathname === streamPath + ? sse([{ type: 'error', error: 'stream broke' }]) + : undefined, + (c) => { + if (!c.url.pathname.endsWith('/messages')) return undefined + const rows = [messageRow(50, 'user', 'question', { job_id: 'job-1' })] + if (answered) rows.push(messageRow(51, 'assistant', 'the answer', { job_id: 'job-1' })) + return json(rows) + } + ) + const chat = createChat(options({}, fetch)) + await chat.selectConversation('conv') + await chat.resumeTurn({ jobId: 'job-1', userSeq: 50 }) + expect(chat.getState().status).toBe('error') + // Another tab carried the turn to its end and its answer is in the rows now. + answered = true + await chat.refreshMessages() + // The answer stands alone: the failure this chat showed was never a row. + expect(chat.getState().messages.map((m) => m.content)).toEqual(['question', 'the answer']) + expect(chat.getState().status).toBe('idle') + expect(chat.getState().error).toBeUndefined() + }) + + test("a re-read does not settle a failed turn with the previous turn's late answer", async () => { + const { fetch } = fetchMock( + (c) => + c.url.pathname === '/api/w/ws/jobs_u/getupdate_sse/job-b' + ? sse([{ type: 'error', error: 'stream broke' }]) + : undefined, + (c) => + c.url.pathname.endsWith('/jobs_u/get/job-b') + ? json({ flow_status: { modules: [{ job: 'step-b' }] } }) + : undefined, + (c) => { + if (!c.url.pathname.endsWith('/messages')) return undefined + const rows = [ + messageRow(50, 'user', 'qA', { job_id: 'job-a' }), + messageRow(51, 'user', 'qB', { job_id: 'job-b' }) + ] + // Turn A's answer commits from its detached task, after B's question. + if (c.url.searchParams.get('after_seq') === '51') { + return json([messageRow(52, 'assistant', "A's answer", { job_id: 'step-a' })]) + } + return json(rows) + } + ) + const chat = createChat(options({}, fetch)) + await chat.selectConversation('conv') + await chat.resumeTurn({ jobId: 'job-b', userSeq: 51 }) + expect(chat.getState().status).toBe('error') + await chat.refreshMessages() + // A's answer is not B's: B's failure stands, and it is still shown. + expect(chat.getState().status).toBe('error') + expect(chat.getState().messages.some((m) => m.success === false && m.seq === undefined)).toBe(true) + }) + + test("a re-read settles the failed turn from its own answer, not only from the newest row", async () => { + const { fetch } = fetchMock( + (c) => + c.url.pathname === streamPath ? sse([{ type: 'error', error: 'stream broke' }]) : undefined, + (c) => + c.url.pathname.endsWith('/jobs_u/get/job-1') + ? json({ flow_status: { modules: [{ job: 'step-1' }] } }) + : undefined, + (c) => { + if (!c.url.pathname.endsWith('/messages')) return undefined + if (c.url.searchParams.get('after_seq') === '50') { + // The turn's answer, and behind it the question of a turn started elsewhere. + return json([ + messageRow(51, 'assistant', 'the answer', { job_id: 'step-1' }), + messageRow(52, 'user', 'next question', { job_id: 'job-2' }) + ]) + } + return json([messageRow(50, 'user', 'question', { job_id: 'job-1' })]) + } + ) + const chat = createChat(options({}, fetch)) + await chat.selectConversation('conv') + await chat.resumeTurn({ jobId: 'job-1', userSeq: 50 }) + expect(chat.getState().status).toBe('error') + await chat.refreshMessages() + expect(chat.getState().status).toBe('idle') + expect(chat.getState().messages.map((m) => m.content)).toEqual([ + 'question', + 'the answer', + 'next question' + ]) + }) + + test('a conversation left and opened again while its rows are read drops that read', async () => { + let releaseRows!: (r: Response) => void + const rowsGate = new Promise((resolve) => (releaseRows = resolve)) + let reads = 0 + const { fetch } = fetchMock((c) => { + if (!c.url.pathname.endsWith('/messages')) return undefined + if (c.url.pathname.includes('/flow_conversations/other/')) return json([]) + reads++ + // The read the refresh made, answered only after the reader has come back. + if (reads === 2) return rowsGate + return json([messageRow(150, 'user', 'newest page', { job_id: 'job-x' })]) + }) + const chat = createChat(options({}, fetch)) + await chat.selectConversation('conv') + const refreshed = chat.refreshMessages() + await chat.selectConversation('other') + await chat.selectConversation('conv') + // Rows from before the reader left: older than the page the chat holds now. + releaseRows(json([messageRow(51, 'user', 'older page', { job_id: 'job-y' })])) + await refreshed + expect(chat.getState().messages.map((m) => m.content)).toEqual(['newest page']) + }) + + test('a turn that runs while the rows are read leaves them to the next read', async () => { + let releaseRows!: (r: Response) => void + const rowsGate = new Promise((resolve) => (releaseRows = resolve)) + let reads = 0 + const { fetch } = fetchMock( + run, + (c) => + c.url.pathname.includes('/jobs_u/getupdate_sse/') + ? sse([{ type: 'error', error: 'stream broke' }]) + : undefined, + (c) => + c.url.pathname.endsWith('/jobs_u/get/job-1') + ? json({ flow_status: { modules: [{ job: 'job-1' }] } }) + : undefined, + (c) => { + if (!c.url.pathname.endsWith('/messages')) return undefined + if (++reads === 1) return json([messageRow(50, 'user', 'qA', { job_id: 'job-a' })]) + // The read the refresh made, answered only after another turn has come and gone. + if (reads === 2) return rowsGate + return json([]) + }, + (c) => (c.url.pathname === '/api/w/ws/flow_conversations/list' ? json([]) : undefined) + ) + const chat = createChat(options({}, fetch)) + await chat.selectConversation('conv') + const refreshed = chat.refreshMessages() + // A turn starts and fails entirely while that read is in flight. + await chat.sendMessage('qB') + expect(chat.getState().status).toBe('error') + releaseRows(json([messageRow(51, 'assistant', "A's answer", { job_id: 'job-a' })])) + await refreshed + // A's answer is not appended under B: it waits for a read that sees the conversation as + // it is now. + expect(chat.getState().messages.map((m) => m.content)).toEqual([ + 'qA', + 'qB', + 'stream broke' + ]) + }) + + test('a turn that starts while the jobs are read still frees the answered failure', async () => { + let releaseJobs!: (r: Response) => void + const jobsGate = new Promise((resolve) => (releaseJobs = resolve)) + let releaseRun!: (r: Response) => void + const runGate = new Promise((resolve) => (releaseRun = resolve)) + const { fetch, calls } = fetchMock( + (c) => (c.method === 'POST' && c.url.pathname === `/api/w/ws/jobs/run/f/${FLOW}` ? runGate : undefined), + (c) => + c.url.pathname.includes('/jobs_u/getupdate_sse/') + ? sse([{ type: 'error', error: 'stream broke' }]) + : undefined, + (c) => (c.url.pathname.endsWith('/jobs_u/get/job-a') ? jobsGate : undefined), + (c) => + c.url.pathname.endsWith('/jobs_u/get/job-b') + ? json({ flow_status: { modules: [{ job: 'job-b' }] } }) + : undefined, + (c) => { + if (!c.url.pathname.endsWith('/messages')) return undefined + if (c.url.searchParams.get('after_seq') === '50') { + return json([messageRow(51, 'assistant', "A's answer", { job_id: 'job-a' })]) + } + return json([messageRow(50, 'user', 'qA', { job_id: 'job-a' })]) + }, + (c) => (c.url.pathname === '/api/w/ws/flow_conversations/list' ? json([]) : undefined) + ) + const chat = createChat(options({}, fetch)) + await chat.selectConversation('conv') + await chat.resumeTurn({ jobId: 'job-a', userSeq: 50 }) + expect(chat.getState().status).toBe('error') + const refreshed = chat.refreshMessages() + await until( + () => calls.some((c) => c.url.pathname.endsWith('/jobs_u/get/job-a')), + 'the job read to start' + ) + // A turn starts while that read is in flight, and is still running when it answers. + const sent = chat.sendMessage('qB') + await until(() => chat.getState().status === 'submitted', 'the new turn to take the chat') + releaseJobs(json({ flow_status: { modules: [{ job: 'job-a' }] } })) + await refreshed + // A's failure goes with its answer; the running turn keeps the chat busy. + expect(chat.getState().messages.some((m) => m.success === false)).toBe(false) + expect(chat.getState().status).toBe('submitted') + releaseRun(text('job-b')) + await sent + }) + + test('a turn that fails while the jobs are read keeps its own failure', async () => { + let releaseJobs!: (r: Response) => void + const jobsGate = new Promise((resolve) => (releaseJobs = resolve)) + const { fetch, calls } = fetchMock( + run, + (c) => + c.url.pathname === streamPath || c.url.pathname === '/api/w/ws/jobs_u/getupdate_sse/job-1' + ? sse([{ type: 'error', error: 'stream broke' }]) + : undefined, + (c) => (c.url.pathname.endsWith('/jobs_u/get/job-a') ? jobsGate : undefined), + (c) => + c.url.pathname.endsWith('/jobs_u/get/job-1') + ? json({ flow_status: { modules: [{ job: 'job-1' }] } }) + : undefined, + (c) => { + if (!c.url.pathname.endsWith('/messages')) return undefined + if (c.url.searchParams.get('after_seq') === '50') { + return json([messageRow(51, 'assistant', "A's answer", { job_id: 'job-a' })]) + } + return json([messageRow(50, 'user', 'qA', { job_id: 'job-a' })]) + }, + (c) => (c.url.pathname === '/api/w/ws/flow_conversations/list' ? json([]) : undefined) + ) + const chat = createChat(options({}, fetch)) + await chat.selectConversation('conv') + await chat.resumeTurn({ jobId: 'job-a', userSeq: 50 }) + expect(chat.getState().status).toBe('error') + // The refresh stops on the job read; a second turn is sent and fails while it waits. + const refreshed = chat.refreshMessages() + await until( + () => calls.some((c) => c.url.pathname.endsWith('/jobs_u/get/job-a')), + 'the job read to start' + ) + await chat.sendMessage('qB') + expect(chat.getState().status).toBe('error') + releaseJobs(json({ flow_status: { modules: [{ job: 'job-a' }] } })) + await refreshed + // A's failure goes, since its answer is here; the error the chat shows is B's. + expect(chat.getState().status).toBe('error') + const failures = chat.getState().messages.filter((m) => m.success === false) + expect(failures).toHaveLength(1) + expect(chat.getState().messages.map((m) => m.content)).toContain("A's answer") + }) + + test('a re-read that finds no answer leaves the failure standing', async () => { + const { fetch } = fetchMock( + (c) => + c.url.pathname === streamPath ? sse([{ type: 'error', error: 'stream broke' }]) : undefined, + (c) => + c.url.pathname.endsWith('/messages') + ? json([messageRow(50, 'user', 'question', { job_id: 'job-1' })]) + : undefined + ) + const chat = createChat(options({}, fetch)) + await chat.selectConversation('conv') + await chat.resumeTurn({ jobId: 'job-1', userSeq: 50 }) + expect(chat.getState().status).toBe('error') + await chat.refreshMessages() + expect(chat.getState().status).toBe('error') + expect(chat.getState().messages.map((m) => m.content)).toEqual(['question', 'stream broke']) + }) + + test('a conversation switch from onFinish does not resume the next turn elsewhere', async () => { + const { fetch, calls } = fetchMock( + run, + (c) => + c.url.pathname === streamPath + ? sse([{ type: 'update', completed: true, only_result: { windmill_chat_answer: 'first answer' } }]) + : undefined, + (c) => + c.url.pathname.endsWith('/jobs_u/get/job-1') + ? json({ flow_status: { modules: [{ job: 'job-1' }] } }) + : undefined, + (c) => { + if (!c.url.pathname.endsWith('/messages')) return undefined + if (c.url.pathname.includes('/flow_conversations/other/messages')) return json([]) + // No answer row lands, so the turn finishes from the flow result. + return json([ + messageRow(50, 'user', 'first', { job_id: 'job-1' }), + messageRow(52, 'user', 'second', { job_id: 'job-2' }) + ]) + }, + (c) => (c.url.pathname === '/api/w/ws/flow_conversations/list' ? json([]) : undefined) + ) + const chat = createChat(options({ onFinish: () => void chat.selectConversation('other') }, fetch)) + await chat.sendMessage('first') + expect(chat.getState().conversationId).toBe('other') + expect(chat.getState().messages).toEqual([]) + expect(calls.some((c) => c.url.pathname.includes('getupdate_sse/job-2'))).toBe(false) + }) + + test('stopping during handoff does not resume the next turn', async () => { + let releaseList!: (response: Response) => void + const listGate = new Promise((resolve) => { + releaseList = resolve + }) + const { fetch, calls } = fetchMock( + run, + (c) => + c.url.pathname === streamPath + ? sse([{ type: 'update', completed: true, only_result: { windmill_chat_answer: 'first answer' } }]) + : undefined, + (c) => + c.url.pathname.endsWith('/jobs_u/get/job-1') + ? json({ flow_status: { modules: [{ job: 'job-1' }] } }) + : undefined, + (c) => + c.url.pathname.endsWith('/messages') + ? json([ + messageRow(50, 'user', 'first', { job_id: 'job-1' }), + messageRow(51, 'assistant', 'first answer', { job_id: 'job-1' }), + messageRow(52, 'user', 'second', { job_id: 'job-2' }) + ]) + : undefined, + (c) => (c.url.pathname === '/api/w/ws/flow_conversations/list' ? listGate : undefined) + ) + const chat = createChat(options({}, fetch)) + const sent = chat.sendMessage('first') + await until(() => calls.some((c) => c.url.pathname === '/api/w/ws/flow_conversations/list'), 'the list to be asked for') + const stopped = chat.stop() + releaseList(json([])) + await Promise.all([sent, stopped]) + expect(chat.getState().status).toBe('idle') + expect(calls.some((c) => c.url.pathname.includes('getupdate_sse/job-2'))).toBe(false) + }) + + test('a stream that keeps ending before the job completes hands the turn to polling', async () => { + let streams = 0 + const { fetch } = fetchMock( + run, + // Every real connection starts with a status snapshot. It is not stream progress. + (c) => (c.url.pathname === streamPath ? (streams++, sse([{ type: 'update', running: true }])) : undefined), + (c) => + c.url.pathname.endsWith('/get_result_maybe/job-1') + ? json({ completed: true, success: true, result: { windmill_chat_answer: 'polled' } }) + : undefined, + (c) => + c.url.pathname.endsWith('/messages') + ? json([messageRow(81, 'user', 'hi'), messageRow(82, 'assistant', 'polled')]) : undefined, (c) => (c.url.pathname === '/api/w/ws/flow_conversations/list' ? json([]) : undefined) ) const chat = createChat(options({}, fetch)) - const first = chat.sendMessage('first') - await new Promise((r) => setTimeout(r, 50)) - await chat.stop() - await first + await chat.sendMessage('hi') + expect(streams).toBe(3) + expect(chat.getState().status).toBe('idle') + expect(chat.getState().messages.map((m) => m.content)).toEqual(['hi', 'polled']) + }, 15000) + + test('resuming a turn whose message is off the first page replays it without duplicating rows', async () => { + const { fetch } = fetchMock( + (c) => + c.url.pathname === streamPath && !c.url.searchParams.has('stream_offset') + ? sse([ + { + type: 'update', + new_result_stream: ndjson( + { type: 'tool_call', call_id: 'c1', function_name: 'lookup' }, + { type: 'tool_result', call_id: 'c1', function_name: 'lookup', result: '1', success: true }, + { type: 'token_delta', content: 'Done' } + ), + stream_offset: 3, + completed: true, + only_result: { output: 'Done', messages: [] } + } + ]) + : undefined, + (c) => { + if (!c.url.pathname.endsWith('/messages')) return undefined + const after = c.url.searchParams.get('after_seq') + // The first page holds only what the running turn wrote so far. + if (after === null) return json([messageRow(51, 'tool', 'Used lookup tool')]) + if (after === '49') return json([messageRow(50, 'user', 'hi')]) + return json([messageRow(51, 'tool', 'Used lookup tool'), messageRow(52, 'assistant', 'Done')]) + } + ) + const chat = createChat(options({}, fetch)) + void chat.selectConversation('conv') + await chat.resumeTurn({ jobId: 'job-1', userSeq: 50 }) + expect(chat.getState().messages.map((m) => [m.serverId, m.role, m.content, m.pending])).toEqual([ + ['row-50', 'user', 'hi', false], + ['row-51', 'tool', 'Used lookup tool', false], + ['row-52', 'assistant', 'Done', false] + ]) + expect(chat.getState().status).toBe('idle') + }) + + test("the previous run's answer landing after the next message is not that turn's answer", async () => { + let reads = 0 + let jobReads = 0 + const { fetch } = fetchMock( + run, + (c) => + c.url.pathname === streamPath + ? sse([{ type: 'update', completed: true, only_result: { windmill_chat_answer: 'second answer' } }]) + : undefined, + // A transient failure of the job read is retried, not taken as "any row counts". + (c) => + c.url.pathname.endsWith('/jobs_u/get/job-1') + ? ++jobReads === 1 + ? text('bad gateway', 502) + : json({ flow_status: { modules: [{ job: 'step-2' }] } }) + : undefined, + (c) => { + if (!c.url.pathname.endsWith('/messages')) return undefined + if (!c.url.searchParams.has('after_seq')) return json([messageRow(71, 'user', 'first')]) + // The earlier agent wrote its answer from a task its run did not wait for. + return json( + ++reads === 1 + ? [messageRow(72, 'user', 'second'), messageRow(73, 'assistant', 'first answer, late', { job_id: 'step-1' })] + : [messageRow(74, 'assistant', 'second answer', { job_id: 'step-2' })] + ) + } + ) + const chat = createChat(options({}, fetch)) + await chat.selectConversation('conv') await chat.sendMessage('second') expect(chat.getState().messages.map((m) => [m.role, m.content, m.serverId])).toEqual([ ['user', 'first', 'row-71'], @@ -989,6 +1710,64 @@ describe('createChat with server history', () => { ]) }) + test('an answer cut by a lost stream gives way to the polled result', async () => { + let streams = 0 + const { fetch } = fetchMock( + run, + (c) => + c.url.pathname === streamPath + ? ++streams === 1 + ? sse([{ type: 'update', new_result_stream: ndjson({ type: 'token_delta', content: 'Hel' }), stream_offset: 1 }]) + : text('bad gateway', 502) + : undefined, + (c) => + c.url.pathname.endsWith('/get_result_maybe/job-1') + ? json({ completed: true, success: true, result: { windmill_chat_answer: 'Hello, full answer' } }) + : undefined + ) + const chat = createChat(options({ history: 'none' }, fetch)) + await chat.sendMessage('hi') + expect(chat.getState().messages.map((m) => [m.role, m.content, m.pending])).toEqual([ + ['user', 'hi', false], + ['assistant', 'Hello, full answer', false] + ]) + }, 15000) + + test('a stream that keeps failing hands the turn to polling the job', async () => { + // Each connection opens, sends the server's ping, then drops. + const pingThenDrop = () => + new Response( + new ReadableStream({ + start(controller) { + controller.enqueue(new TextEncoder().encode(`data: ${JSON.stringify({ type: 'ping' })}\n\n`)) + setTimeout(() => controller.error(new TypeError('network connection was lost')), 20) + } + }), + { status: 200, headers: { 'content-type': 'text/event-stream' } } + ) + const { fetch } = fetchMock( + run, + (c) => (c.url.pathname === streamPath ? pingThenDrop() : undefined), + (c) => + c.url.pathname.endsWith('/get_result_maybe/job-1') + ? json({ completed: true, success: true, result: { windmill_chat_answer: 'polled' } }) + : undefined, + (c) => + c.url.pathname.endsWith('/messages') + ? json([messageRow(61, 'user', 'hi'), messageRow(62, 'assistant', 'polled')]) + : undefined, + (c) => (c.url.pathname === '/api/w/ws/flow_conversations/list' ? json([]) : undefined) + ) + const chat = createChat(options({}, fetch)) + await chat.sendMessage('hi') + const state = chat.getState() + expect(state.status).toBe('idle') + expect(state.messages.map((m) => [m.role, m.content, m.success])).toEqual([ + ['user', 'hi', true], + ['assistant', 'polled', true] + ]) + }, 15000) + test('a failure handler answer is attributed to the turn', async () => { let reads = 0 const { fetch } = fetchMock( diff --git a/chat-sdk/test/follow.test.ts b/chat-sdk/test/follow.test.ts index 19af6baf38..991b1a7886 100644 --- a/chat-sdk/test/follow.test.ts +++ b/chat-sdk/test/follow.test.ts @@ -23,3 +23,20 @@ test('a retried streaming step reports its offset as lost before the new sub-job } expect(offsets).toEqual([3, undefined, 1]) }) + +test('a 200 carrying no stream is reconnected like any other failed connection', async () => { + let streams = 0 + const { fetch } = fetchMock((c) => { + if (!c.url.pathname.endsWith('/jobs_u/getupdate_sse/job-1')) return undefined + // Something in front of Windmill answers 200 with nothing in it. + if (++streams === 1) return new Response(null, { status: 200 }) + return sse([{ type: 'update', completed: true, only_result: 'ok' }]) + }) + const api = new WindmillChatApi({ baseUrl: 'http://wm.test', workspace: 'ws', token: 'tok', fetch }) + let result: unknown + for await (const event of followJob(api, 'job-1')) { + if (event.type === 'completed') result = event.result + } + expect(streams).toBe(2) + expect(result).toBe('ok') +}) diff --git a/frontend/src/lib/components/FlowPreviewContent.svelte b/frontend/src/lib/components/FlowPreviewContent.svelte index 9ccac1d8c6..5489e4abb5 100644 --- a/frontend/src/lib/components/FlowPreviewContent.svelte +++ b/frontend/src/lib/components/FlowPreviewContent.svelte @@ -1,5 +1,6 @@ + +{#if count > 0} + + {count > 9 ? '9+' : count} + +{/if} + + diff --git a/frontend/src/lib/components/copilot/chat/AIChatDisplay.svelte b/frontend/src/lib/components/copilot/chat/AIChatDisplay.svelte index 2d88d1deb7..387eea15d0 100644 --- a/frontend/src/lib/components/copilot/chat/AIChatDisplay.svelte +++ b/frontend/src/lib/components/copilot/chat/AIChatDisplay.svelte @@ -247,8 +247,16 @@ function onWindowKeydownCapture(e: KeyboardEvent) { if (e.key !== 'Escape' || !chatHost.loading) return const active = document.activeElement + // Focus parked on the body answers for the panel on screen only: the flow chat keeps + // a panel per conversation mounted, several of which can be loading, and a hidden + // one's listener would otherwise swallow the press (immediate form, below) and stop + // its own turn instead. `inert` marks the panels behind, and is matched as an + // attribute so the browsers without `checkVisibility` read it too. + const hidden = + panelEl?.closest('[inert]') != null || + panelEl?.checkVisibility?.({ visibilityProperty: true }) === false const focusOnChat = - !active || active === document.body || (panelEl?.contains(active) ?? false) + !hidden && (!active || active === document.body || panelEl?.contains(active) === true) // An Escape while a run form is open must not discard what the user typed, so the action // row alone stops the turn — wherever it is mounted, since the preview panel holds the // form outside `panelEl`. Matched by call: two chats can be loading at once, and one's @@ -487,22 +495,30 @@ const flatFiles = Array.from(dt.files ?? []) const topLevelImages = flatFiles.filter(isImageFile) const imageWork: Promise[] = [] + // The composer this drop landed on, held for the whole routing: a panel destroyed + // while a file is being read clears the binding, and the file would then reach no + // composer at all. + const dropped = aiChatInput if (topLevelImages.length > 0) { - imageWork.push(aiChatInput?.addImages(topLevelImages) ?? Promise.resolve()) + imageWork.push(dropped?.addImages(topLevelImages) ?? Promise.resolve()) } // Text-file routing must await handle/entry resolution before it can call // addTextFiles — hold sending across that window (taken BEFORE the first // await) or a send mid-resolution would land the drop on the next message. - const releaseSendHold = aiChatInput?.holdSendForIngestion() + const releaseSendHold = dropped?.holdSendForIngestion() try { - await routeDroppedTextAndFolders(dt, flatFiles) + await routeDroppedTextAndFolders(dt, flatFiles, dropped) } finally { releaseSendHold?.() } await Promise.all(imageWork) } - async function routeDroppedTextAndFolders(dt: DataTransfer, flatFiles: File[]) { + async function routeDroppedTextAndFolders( + dt: DataTransfer, + flatFiles: File[], + input: typeof aiChatInput + ) { if (canUseFsAccess) { // getAsFileSystemHandle calls are kicked off synchronously inside this call. const handles = await handlesFromDataTransfer(dt) @@ -514,7 +530,10 @@ ? flatFiles : await Promise.all(handles.filter(isFileHandle).map((h) => h.getFile())) // Loose files attach to the message, like images. - await attachNonImageFiles(looseFiles.filter((f) => !isImageFile(f))) + await attachNonImageFiles( + looseFiles.filter((f) => !isImageFile(f)), + input + ) // Folders link as a live handle. const dirs = handles.filter(isDirectoryHandle) if (dirs.length > 0 && !canLinkFolders) { @@ -551,7 +570,7 @@ if (canLinkFolders) await handleAddFiles(folderEntries) else sendUserToast('Folders cannot be attached in this chat — drop individual files.', true) } - await attachNonImageFiles(topLevelText) + await attachNonImageFiles(topLevelText, input) } } @@ -563,8 +582,10 @@ input.value = '' // allow re-selecting the same file } - async function attachNonImageFiles(files: File[]) { - await aiChatInput?.addNonImageFiles(files) + /** `input` is the composer the files are for, captured before any await: read off the + * binding afterwards it would be null once the panel has been destroyed. */ + async function attachNonImageFiles(files: File[], input: typeof aiChatInput) { + await input?.addNonImageFiles(files) } async function attachPickedFiles(picked: File[]) { @@ -572,7 +593,7 @@ const others = picked.filter((f) => !isImageFile(f)) // Reserved before the other work is awaited — see onPanelDrop. const imageWork = imageFiles.length > 0 ? aiChatInput?.addImages(imageFiles) : undefined - await attachNonImageFiles(others) + await attachNonImageFiles(others, aiChatInput) await imageWork } diff --git a/frontend/src/lib/components/copilot/chat/flow/FlowAIChat.svelte b/frontend/src/lib/components/copilot/chat/flow/FlowAIChat.svelte index 9744085567..6cb398ee44 100644 --- a/frontend/src/lib/components/copilot/chat/flow/FlowAIChat.svelte +++ b/frontend/src/lib/components/copilot/chat/flow/FlowAIChat.svelte @@ -2,7 +2,7 @@ import FlowModuleSchemaMap from '$lib/components/flows/map/FlowModuleSchemaMap.svelte' import { getContext, tick, untrack } from 'svelte' import type { ExtendedOpenFlow, FlowEditorContext } from '$lib/components/flows/types' - import type { InputTransform } from '$lib/gen' + import { ApiError, type InputTransform } from '$lib/gen' import type { FlowAIChatHelpers } from './core' import { chatMemoryId } from '../global/core' import { createInlineScriptSession } from './inlineScriptsUtils' @@ -174,7 +174,19 @@ previewArgs.val = args } // Call the UI test function which opens preview panel - return await onTestFlow?.(conversationId ?? chatMemoryId(flowStore.val.value)) + try { + return await onTestFlow?.(conversationId ?? chatMemoryId(flowStore.val.value)) + } catch (e) { + // The flow chat is still answering in that conversation, and the server takes one + // turn at a time. Said plainly here: the raw refusal names a job the model has + // no use for. + if (e instanceof ApiError && e.status === 409) { + throw new Error( + 'That chat is still answering an earlier message; wait for it to finish before testing again.' + ) + } + throw e + } }, getLintErrors: async (moduleId: string): Promise => { diff --git a/frontend/src/lib/components/flows/conversations/FlowChat.svelte b/frontend/src/lib/components/flows/conversations/FlowChat.svelte index 01e2f60850..04671e061b 100644 --- a/frontend/src/lib/components/flows/conversations/FlowChat.svelte +++ b/frontend/src/lib/components/flows/conversations/FlowChat.svelte @@ -1,13 +1,22 @@
- {#if chat && chatState} + {#if listChat && listState && pool && poolState} {#if !hideSidebar} @@ -136,24 +204,45 @@ -
- - {#key chat} - - {/key} +
+ + {#each poolState.mounted as key (key)} + {@const panel = pool.get(key)} + {#if panel} + {@const shown = key === shownKey} +
+ + listState?.conversations.find((c) => c.id === id)?.isTest ?? + conversationKinds.get(id)} + {deploymentInProgress} + {additionalInputsSchema} + {flowModules} + {path} + {identity} + {workspace} + {description} + {wideLayout} + {conversationKind} + {subject} + /> +
+ {/if} + {/each}
{/if}
diff --git a/frontend/src/lib/components/flows/conversations/FlowChatInterface.svelte b/frontend/src/lib/components/flows/conversations/FlowChatInterface.svelte index cad5688d89..f502be90cb 100644 --- a/frontend/src/lib/components/flows/conversations/FlowChatInterface.svelte +++ b/frontend/src/lib/components/flows/conversations/FlowChatInterface.svelte @@ -3,14 +3,14 @@ import { Loader2, MessageCircle, Settings2 } from 'lucide-svelte' import AIChatDisplay from '$lib/components/copilot/chat/AIChatDisplay.svelte' import { setChatViewHost } from '$lib/components/copilot/chat/chatViewHost' - import { FlowChatViewHost } from './flowChatViewHost.svelte' + import type { FlowChatViewHost } from './flowChatViewHost.svelte' import Modal from '$lib/components/common/modal/Modal.svelte' import SchemaForm from '$lib/components/SchemaForm.svelte' import GfmMarkdown from '$lib/components/GfmMarkdown.svelte' import { emptyString, type DynamicInput } from '$lib/utils' - import { onDestroy, tick, untrack } from 'svelte' + import { tick, untrack } from 'svelte' import type { Chat } from 'windmill-chat' - import { chatFlowKey } from './flowChatProps' + import { saveFlowChatInputs } from './flowChatProps' import type { FlowModule } from '$lib/gen' import { useWorkspaceStorageConfigured } from '$lib/components/inputTransformEnv.svelte' import { @@ -30,6 +30,14 @@ interface Props { chat: Chat + /** The conversation's host, which outlives this panel. One panel follows one chat for + * its whole life, so a later value of either prop never reaches it. */ + chatHost: FlowChatViewHost + /** The flow inputs the reader chose, shared by every conversation in the flow: FlowChat + * holds them, since each conversation has a panel of its own mounted. */ + inputValues: Record + /** Whether a conversation is a test chat, once the list has said. */ + isTestOf?: (conversationId: string) => boolean | undefined deploymentInProgress?: boolean additionalInputsSchema?: Record /** The flow's modules, read for the AI agent inputs the composer drives: the provider wiring @@ -50,6 +58,9 @@ let { chat, + chatHost: chatHostProp, + inputValues = $bindable(), + isTestOf = undefined, deploymentInProgress = false, additionalInputsSchema, flowModules, @@ -97,12 +108,7 @@ const modelGap = $derived(agentModelGap(modelWiring, subject)) const showModelButton = $derived(showsModelButton(modelWiring)) - // LocalStorage helpers - const STORAGE_KEY_PREFIX = 'windmill_flow_chat_inputs_' - let showInputsModal = $state(false) - // Conversation settings, persisted per flow: what the reader chose, and nothing else. - let inputValues = $state>(loadInputsFromStorage() ?? {}) let modalDraft = $state>({}) /** What the flow's own form would open on. */ @@ -128,31 +134,9 @@ // value, an author's default — is made safe before it reaches the provider. const runInputs = $derived(withoutRejectedEffort(modelWiring, effectiveInputs)) - function getStorageKey(): string { - return `${STORAGE_KEY_PREFIX}${chatFlowKey({ path, identity })}` - } - - function loadInputsFromStorage(): Record | null { - try { - const stored = localStorage.getItem(getStorageKey()) - return stored ? JSON.parse(stored) : null - } catch (e) { - console.error('Failed to load inputs from localStorage:', e) - return null - } - } - - function saveInputsToStorage(values: Record) { - try { - localStorage.setItem(getStorageKey(), JSON.stringify(values)) - } catch (e) { - console.error('Failed to save inputs to localStorage:', e) - } - } - function setInputValue(name: string, value: any) { inputValues = { ...inputValues, [name]: value } - saveInputsToStorage(inputValues) + saveFlowChatInputs({ path, identity }, inputValues) } function handleModalConfirm() { @@ -167,47 +151,48 @@ ) ) inputValues = kept - saveInputsToStorage(inputValues) + saveFlowChatInputs({ path, identity }, inputValues) showInputsModal = false } function openInputsModal() { - modalDraft = { ...effectiveInputs, ...(loadInputsFromStorage() ?? inputValues) } + modalDraft = { ...effectiveInputs, ...inputValues } showInputsModal = true } - // The host follows the chat it was built on for the life of this component: FlowChat - // remounts the interface under `{#key chat}`, so a later value of the prop never reaches it. - const chatHost = new FlowChatViewHost( - untrack(() => chat), - { - additionalInputs: () => (additionalInputsSchema ? { ...runInputs } : undefined), - attachmentsTarget: () => attachmentsTarget, - attachmentsUnavailable: () => - workspaceStorage.current - ? undefined - : 'This workspace has no object storage, so files cannot be attached.', - workspace: () => workspace, - sendDisabled: () => deploymentInProgress || !!modelGap || !!wrongKindReason, - // The model controls only: a retry changes model when the reader did, but replays - // the run's own attachments rather than whatever the composer holds now. - inputsShownInComposer: () => composerOwnedInputs(modelWiring, undefined) - } - ) + // The host belongs to the conversation, not to this panel: the pool keeps it alive so a + // message queued here still goes out once the reader has moved on. What it reads is this + // panel's, set on mount; one panel shows one conversation for its whole life. + const chatHost = untrack(() => chatHostProp) + chatHost.setOptions({ + additionalInputs: () => (additionalInputsSchema ? { ...runInputs } : undefined), + attachmentsTarget: () => attachmentsTarget, + attachmentsUnavailable: () => + workspaceStorage.current + ? undefined + : 'This workspace has no object storage, so files cannot be attached.', + workspace: () => workspace, + sendDisabled: () => deploymentInProgress || !!modelGap || !!wrongKindReason, + // The model controls only: a retry changes model when the reader did, but replays + // the run's own attachments rather than whatever the composer holds now. + inputsShownInComposer: () => composerOwnedInputs(modelWiring, undefined) + }) setChatViewHost(chatHost) + // Read off the host's own conversation: a message queued here goes out after the reader + // has moved on, when the shown conversation may be of the other kind. + const isTest = $derived.by(() => { + const id = chatHost.state.conversationId + return id === undefined ? undefined : isTestOf?.(id) + }) // A chat of the other kind can be read from here but not added to: the server refuses a // preview run into a deployed conversation and the reverse, so the composer says why first. const wrongKindReason = $derived.by(() => { - const { conversationId, conversations } = chatHost.state - const open = conversations.find((c) => c.id === conversationId) - if (open?.isTest === undefined || open.isTest === (conversationKind === 'test')) - return undefined - return open.isTest + if (isTest === undefined || isTest === (conversationKind === 'test')) return undefined + return isTest ? 'This chat was run from the flow editor. Start a new chat to continue here.' : 'This chat belongs to the deployed flow. Start a new chat to test.' }) - onDestroy(() => chatHost.dispose()) // What the Configure-inputs modal asks for: every flow input the composer does not // edit itself. diff --git a/frontend/src/lib/components/flows/conversations/FlowConversationsSidebar.svelte b/frontend/src/lib/components/flows/conversations/FlowConversationsSidebar.svelte index d2ea5c5061..8c6a102fef 100644 --- a/frontend/src/lib/components/flows/conversations/FlowConversationsSidebar.svelte +++ b/frontend/src/lib/components/flows/conversations/FlowConversationsSidebar.svelte @@ -5,11 +5,13 @@ Plus, Trash2, Pen, + PencilLine, Filter, PanelLeftClose, PanelLeftOpen } from 'lucide-svelte' - import CountBadge from '$lib/components/common/badge/CountBadge.svelte' + import UnreadCountBadge from '$lib/components/common/badge/UnreadCountBadge.svelte' + import SessionStatusDot from '$lib/components/sessions/SessionStatusDot.svelte' import InfiniteList from '$lib/components/InfiniteList.svelte' import DropdownV2 from '$lib/components/DropdownV2.svelte' import Popover from '$lib/components/meltComponents/Popover.svelte' @@ -21,11 +23,16 @@ import { twMerge } from 'tailwind-merge' import { fade } from 'svelte/transition' import { tick, untrack } from 'svelte' - import type { Chat, ChatState, Conversation, ConversationKind } from 'windmill-chat' + import type { Chat, Conversation, ConversationKind } from 'windmill-chat' + import type { FlowChatPool, FlowChatPoolState } from './flowChatPool' + import type { ComposerAttachment, FlowChatViewHost } from './flowChatViewHost.svelte' interface Props { - chat: Chat - chatState: ChatState + /** Lists, renames and deletes the flow's conversations; runs no turn itself. */ + listChat: Chat + /** The conversations' own chats: which one is shown, and what each is doing. */ + pool: FlowChatPool + poolState: FlowChatPoolState /** * Which conversations the list holds at first. The editor shows its own test chats, * since testing is what happens there; a deployed flow shows the chats its users @@ -40,7 +47,13 @@ canFilterKind?: boolean } - let { chat, chatState, defaultKind = 'deployed', canFilterKind = false }: Props = $props() + let { + listChat, + pool, + poolState, + defaultKind = 'deployed', + canFilterKind = false + }: Props = $props() let expanded = $state(false) let list = $state(undefined) @@ -58,13 +71,13 @@ let renameDraft = $state('') let renameInput = $state(undefined) - const turnInFlight = $derived( - chatState.status === 'submitted' || chatState.status === 'streaming' + const totalUnread = $derived( + Object.values(poolState.unread).reduce((total, count) => total + count, 0) ) $effect(() => { const l = list - const c = chat + const c = listChat if (!l) return untrack(() => { // Every load goes through here, the first one and infinite scroll included. A @@ -72,13 +85,18 @@ // selected kind brings its own, whichever of the two lands last. l.setLoader(async (page, perPage) => { const requested = kind + // Only a listing read now says which turns run: the rows the chat holds keep + // what the last one said, and a rename or delete publishes those again. + const since = pool.listingStarted() const rows = await c.loadConversations({ page, perPage, kind: requested }) + pool.setListed(rows, since) return requested === kind ? rows : items }) l.setDeleteItemFn(async (id: string) => { deletingId = id try { await c.deleteConversation(id) + pool.forget(id) sendUserToast('Conversation deleted successfully') } catch (error) { console.error('Failed to delete conversation:', error) @@ -103,10 +121,10 @@ await list?.loadData('forceRefresh') } - const draftShown = $derived(draft && !items.some((c) => c.id === chatState.conversationId)) + const draftShown = $derived(draft && poolState.selectedId === undefined) function newChat() { - chat.newConversation() + pool.newChat() draft = true } @@ -118,17 +136,16 @@ /** * Narrow the list to one kind of chat and reload it. The open conversation goes with it - * when it is not of the new kind: the composer sends into whatever is selected, and a + * when it is not of the new kind: the composer sends into whatever is shown, and a * conversation keeps the kind it was created with, so a turn sent into one the list no - * longer shows would be stored where nothing here lists it. + * longer shows would be stored where nothing here lists it. A turn running in it goes on. */ async function setKind(next: ConversationKind) { - // A turn writes into the open conversation, which a kind that excludes it would close. - if (next === kind || turnInFlight) return + if (next === kind) return kind = next - const open = items.find((c) => c.id === chatState.conversationId) + const open = items.find((c) => c.id === poolState.selectedId) const stillListed = open === undefined || next === 'all' || (next === 'test') === open.isTest - if (!stillListed) chat.newConversation() + if (!stillListed) pool.newChat() await list?.loadData('forceRefresh') } @@ -149,11 +166,11 @@ const current = items.find((c) => c.id === id) if (!current || title === '' || title === current.title) return try { - await chat.renameConversation(id, title) + await listChat.renameConversation(id, title) // The list holds its own rows, loaded through the loader: patched rather than // reloaded, so the row keeps its place without a round trip. The title is read // back from the chat, which holds it as the server stored it (a long one is cut). - const stored = chat.getState().conversations.find((c) => c.id === id)?.title ?? title + const stored = listChat.getState().conversations.find((c) => c.id === id)?.title ?? title items = items.map((c) => (c.id === id ? { ...c, title: stored } : c)) } catch (error) { console.error('Failed to rename conversation:', error) @@ -177,8 +194,34 @@ function getConversationTitle(conversation: Conversation): string { return conversation.title || `Conversation ${conversation.createdAt.slice(0, 10)}` } + + /** The session sidebar's dot vocabulary, which knows no `queued`: that is its own mark. */ + function dotStatus(conversationId: string): 'streaming' | 'error' | 'idle' { + const activity = poolState.activity[conversationId] + return activity === 'running' ? 'streaming' : activity === 'error' ? 'error' : 'idle' + } +{#snippet statusDot(conversation: Conversation)} + + + {#snippet resting()} + + {/snippet} + +{/snippet} + {/if} @@ -299,7 +341,6 @@ onClick={(e) => { e?.stopPropagation() draft = false - chat.newConversation() }} title="Discard draft" destructive @@ -314,7 +355,7 @@
{:else} + {@const unread = poolState.unread[conversation.id] ?? 0} + {@const queued = !!poolState.queued[conversation.id]} diff --git a/frontend/src/lib/components/sessions/SessionStatusDot.svelte b/frontend/src/lib/components/sessions/SessionStatusDot.svelte index 646163bcaf..e97de79caf 100644 --- a/frontend/src/lib/components/sessions/SessionStatusDot.svelte +++ b/frontend/src/lib/components/sessions/SessionStatusDot.svelte @@ -8,8 +8,18 @@ let { status, isFork, - forkDetached = false - }: { status: SessionChatStatus; isFork: boolean; forkDetached?: boolean } = $props() + forkDetached = false, + resting, + restingTitle + }: { + status: SessionChatStatus + isFork: boolean + forkDetached?: boolean + /** What the slot shows when there is no live signal. Sessions leave it unset and get + * the workspace/fork mark below; another list passes its own resting mark. */ + resting?: import('svelte').Snippet + restingTitle?: string + } = $props() const statusTooltip: Record = { idle: 'No chat activity', @@ -39,7 +49,7 @@ : 'Root workspace session' ) - const title = $derived(liveOverride ? statusTooltip[status] : persistentTitle) + const title = $derived(liveOverride ? statusTooltip[status] : (restingTitle ?? persistentTitle)) @@ -55,6 +65,8 @@ {:else if status === 'error'} + {:else if resting} + {@render resting()} {:else if isFork} {#if forkDetached}