mirror of
https://github.com/windmill-labs/windmill.git
synced 2026-09-21 00:02:30 +00:00
refactor(chat): give a flow chat turn a module of its own
A turn was not a thing. What one is made of lived in three places — a reactive status record, a non-reactive runtime map, and locals inside the function following the job — so the question every piece of a turn's work has to answer, "am I still the turn this chat is on?", had no owner. Each await site answered it again from whatever abort flag it had in scope, and each new await needed another answer; the counter this replaces was the eighth such answer in as many review rounds. `Turn` holds the run, the handle that stops it, the cursor its stream resumes from, the rows it is writing and the timers it armed, and the chat keys one per conversation. Work carries its turn and passes `isCurrent` before writing, so a frame buffered before a Stop, a poll dispatched a tick ago, or a settle that outlived its own turn all land nowhere instead of in whatever replaced it. Ending a turn is one call rather than an abort, two timers and a pair of reveals that every caller had to remember, and the sites that release one say which turn they mean. The chat itself is keyed on its flow and workspace, so changing flow builds a new panel rather than re-pointing the old one. That is what makes the guards removable: nothing survives a flow change to need them. Fixes found while doing this, each of which the old shape allowed: a two-minute polling deadline that outlived its turn and stopped the next one's reading; a run launched into a stopped turn that nothing could name; a conversation-list refresh on every poll tick of a new chat's first turn; a failed page walk that freed a chat whose run it never got to ask about; and a Stop during a launch that tore down the turn that came after it. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
This commit is contained in:
co-authored by
Claude Opus 5
parent
d3344eb3cd
commit
dda7ab2f30
+21
@@ -37,6 +37,27 @@ _Avoid_: argument field, param
|
||||
Any other place a property can be picked into: the loop iterator, skip and early-stop predicates, the retry condition, a branch predicate, timeout. Its prop picker opens in a popover from the connect button rather than taking a pane.
|
||||
_Avoid_: JS field, code input
|
||||
|
||||
### Flow chat
|
||||
|
||||
**Conversation**:
|
||||
One thread of messages against one chat-enabled flow, with its own agent memory. A flow has
|
||||
many; the chat shows one at a time and keeps the others live, so a turn keeps writing while
|
||||
the reader is in another conversation.
|
||||
_Avoid_: thread, session (that names an AI session, a different thing), chat (that names the surface)
|
||||
|
||||
**Turn**:
|
||||
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. A turn exists from the moment it takes the chat — before it has a job, while an
|
||||
attachment uploads or a reload works out whether a run is live — until it is ended.
|
||||
_Avoid_: request, exchange, message round
|
||||
|
||||
**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
|
||||
from what it holds is one to ask the server rather than to guess at.
|
||||
_Avoid_: history, messages (too easily read as "all of them")
|
||||
|
||||
### Permissions
|
||||
|
||||
**Member**:
|
||||
|
||||
@@ -1,133 +1,28 @@
|
||||
<script lang="ts">
|
||||
import { workspaceStore } from '$lib/stores'
|
||||
import { createFlowChatManager, type ConversationKind } from './FlowChatManager.svelte'
|
||||
import FlowConversationsSidebar from './FlowConversationsSidebar.svelte'
|
||||
import FlowChatInterface from './FlowChatInterface.svelte'
|
||||
import { getContext, untrack } from 'svelte'
|
||||
import { getContext } from 'svelte'
|
||||
import type { FlowEditorContext } from '../types'
|
||||
import type { FlowModule } from '$lib/gen'
|
||||
import type { FlowChatProps } from './flowChatProps'
|
||||
import FlowChatPanel from './FlowChatPanel.svelte'
|
||||
|
||||
type ChatFrame = 'boxed' | 'top' | 'none'
|
||||
|
||||
const FRAME_CLASS: Record<ChatFrame, string> = {
|
||||
boxed: 'border rounded-md',
|
||||
top: 'border-t',
|
||||
none: ''
|
||||
}
|
||||
|
||||
interface Props {
|
||||
onRunFlow: (
|
||||
userMessage: string,
|
||||
conversationId: string,
|
||||
additionalInputs?: Record<string, any>
|
||||
) => Promise<string | undefined>
|
||||
useStreaming?: boolean
|
||||
deploymentInProgress?: boolean
|
||||
path: string
|
||||
/** The flow's own description, shown where the chat has room for it: the empty
|
||||
* transcript, and the sidebar once a conversation has replaced it. */
|
||||
description?: string
|
||||
inputSchema?: Record<string, any>
|
||||
/** The flow's modules, used to find which inputs an AI agent step reads directly. */
|
||||
flowModules?: FlowModule[]
|
||||
/** Wider centered column, for the full-page chat. */
|
||||
wideLayout?: boolean
|
||||
/** What separates the chat from what sits above it. `boxed` is its own panel, corners
|
||||
* clipped so the sidebar's edge follows them; `top` a dividing line under an enclosing
|
||||
* header; `none` for a surface where the chat is the whole pane. */
|
||||
frame?: ChatFrame
|
||||
/** Which chats the sidebar lists before the reader filters it themselves. */
|
||||
conversationKind?: ConversationKind
|
||||
/** Whether a turn may run while another chat's is still going. Off where one run
|
||||
* owns the surface — the editor's panel shows it on the graph. */
|
||||
parallelTurns?: boolean
|
||||
}
|
||||
|
||||
let {
|
||||
onRunFlow,
|
||||
deploymentInProgress = false,
|
||||
useStreaming = false,
|
||||
path,
|
||||
description = undefined,
|
||||
inputSchema = undefined,
|
||||
flowModules = undefined,
|
||||
wideLayout = false,
|
||||
frame = 'top',
|
||||
conversationKind = 'deployed',
|
||||
parallelTurns = false
|
||||
}: Props = $props()
|
||||
let props: FlowChatProps = $props()
|
||||
|
||||
const flowEditorContext = getContext<FlowEditorContext>('FlowEditorContext')
|
||||
|
||||
const manager = createFlowChatManager()
|
||||
manager.operatingWorkspace = () => flowEditorContext?.opWorkspace?.()
|
||||
manager.conversationKind = conversationKind
|
||||
manager.allowsParallelTurns = parallelTurns
|
||||
// The filter moves; what this surface runs does not.
|
||||
manager.surfaceKind = conversationKind === 'test' ? 'test' : 'deployed'
|
||||
// The editor is the only surface with both kinds in play, and it is the one that opens
|
||||
// on test chats. A deployed flow lists what its users started, with no way to ask for
|
||||
// anything else.
|
||||
manager.canFilterConversationKind = conversationKind !== 'deployed'
|
||||
|
||||
// Initialize manager when component mounts
|
||||
$effect(() => {
|
||||
if ($workspaceStore) {
|
||||
manager.initialize(onRunFlow, path, useStreaming)
|
||||
// Reads the open conversation, and this effect tears down with `cleanup()`: tracked,
|
||||
// the first send of a fresh chat would select the conversation it just created and
|
||||
// so abort its own turn.
|
||||
untrack(() => manager.selectLatestConversation())
|
||||
}
|
||||
|
||||
return () => {
|
||||
manager.cleanup()
|
||||
}
|
||||
})
|
||||
|
||||
// Initialize InfiniteList when component mounts or flowPath changes
|
||||
$effect(() => {
|
||||
if ($workspaceStore && path && manager.conversationListComponent) {
|
||||
untrack(() => {
|
||||
manager.setupInfiniteList()
|
||||
})
|
||||
}
|
||||
})
|
||||
|
||||
// Everything the chat asks for beyond the message itself. `user_message` is the
|
||||
// composer: the server requires that exact argument on a chat-enabled flow and
|
||||
// stores it as the conversation's message (`handle_chat_conversation_messages`), so
|
||||
// the name is a contract rather than the author's choice.
|
||||
const additionalInputsSchema = $derived.by(() => {
|
||||
const props = inputSchema?.properties ?? {}
|
||||
const messageInput = 'user_message'
|
||||
const filtered = Object.fromEntries(Object.entries(props).filter(([k]) => k !== messageInput))
|
||||
if (Object.keys(filtered).length === 0) return undefined
|
||||
const required = inputSchema?.required
|
||||
const requiredArray: string[] = Array.isArray(required) ? required : []
|
||||
return {
|
||||
...inputSchema,
|
||||
properties: filtered,
|
||||
required: requiredArray.filter((k: string) => k !== messageInput)
|
||||
}
|
||||
})
|
||||
/**
|
||||
* A chat holds one flow's conversations, the turn running in one of them, and the rows
|
||||
* that turn is writing. Pointing the same chat at another flow would leave every one of
|
||||
* those in flight — an upload, a launch, a poll, a list request — landing in the chat
|
||||
* that replaced it, so the panel is replaced wholesale instead and nothing outlives it.
|
||||
*
|
||||
* The workspace is the one the surface operates on, not the navigated one: an embedded
|
||||
* editor acts on a session's fork while `workspaceStore` stays where the reader left it.
|
||||
*/
|
||||
const chatKey = $derived(
|
||||
`${flowEditorContext?.opWorkspace?.() ?? $workspaceStore ?? ''}:${props.path}`
|
||||
)
|
||||
</script>
|
||||
|
||||
<!-- The column's max width and side padding come from AIChatDisplay itself. -->
|
||||
<div class="flex overflow-hidden flex-1 {FRAME_CLASS[frame]}">
|
||||
<FlowConversationsSidebar {manager} {description} />
|
||||
<!-- pb-3 on the chat alone, not on the row: the transcript and composer stop short of
|
||||
the panel edge the way the session chat does, while the sidebar and the border
|
||||
dividing it from the chat still reach the bottom. -->
|
||||
<div class="flex flex-1 min-w-0 min-h-0 pb-3">
|
||||
<FlowChatInterface
|
||||
{manager}
|
||||
{deploymentInProgress}
|
||||
{additionalInputsSchema}
|
||||
{flowModules}
|
||||
{path}
|
||||
{description}
|
||||
{wideLayout}
|
||||
/>
|
||||
</div>
|
||||
</div>
|
||||
{#key chatKey}
|
||||
<FlowChatPanel {...props} />
|
||||
{/key}
|
||||
|
||||
@@ -13,17 +13,8 @@ import { base } from '$lib/base'
|
||||
import { followJob, WindmillChatApi, type AgentStreamEvent } from 'windmill-chat'
|
||||
import type { StreamEvent } from '$lib/components/chat/utils'
|
||||
import { randomUUID } from '$lib/utils/uuid'
|
||||
import {
|
||||
prefersInstantReveal,
|
||||
TypewriterReveal
|
||||
} from '$lib/components/copilot/chat/typewriterReveal'
|
||||
import {
|
||||
appendRevealed,
|
||||
applyStreamEvent,
|
||||
emptyTurnState,
|
||||
turnFailed,
|
||||
type TurnState
|
||||
} from './turnTranscript'
|
||||
import { appendRevealed, applyStreamEvent, turnFailed } from './turnTranscript'
|
||||
import { Turn, type RevealKind } from './turn.svelte'
|
||||
|
||||
export interface ChatMessage extends FlowConversationMessage {
|
||||
loading?: boolean
|
||||
@@ -57,11 +48,6 @@ type TurnStatus = {
|
||||
* the send lands in.
|
||||
*/
|
||||
isDispatchingTurn: boolean
|
||||
/** The thinking of the turn in flight, until it is attached to the answer it produced. */
|
||||
currentReasoning: string
|
||||
/** The model is reasoning: true from the first thinking token until the answer starts. */
|
||||
isReasoningActive: boolean
|
||||
jobId?: string
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -75,8 +61,6 @@ const FOLLOW_RETRY_DELAY_MS = 500
|
||||
/** How often a turn with no stream asks its job whether the run is over. */
|
||||
const SETTLE_POLL_MS = 2000
|
||||
|
||||
/** How long a conversation is polled before the reader is assumed to have left it running. */
|
||||
const POLL_MAX_MS = 2 * 60 * 1000
|
||||
/** Rows per request when a poll reads a conversation. A batch shorter than this is how the
|
||||
* endpoint says there are no more. */
|
||||
const POLL_PAGE_SIZE = 50
|
||||
@@ -84,6 +68,10 @@ const POLL_PAGE_SIZE = 50
|
||||
* stops short of is left to the next tick, which resumes from the cursor it reached. */
|
||||
const POLL_MAX_REQUESTS = 20
|
||||
|
||||
/** Pages walked back looking for the message that started the newest turn. Bounds the search
|
||||
* on a turn that wrote an implausible number of rows; the conversation's own start ends it. */
|
||||
const RESUME_MAX_PAGES = 20
|
||||
|
||||
/**
|
||||
* The agent events the worker streams, in the shape the transcript applies. The SDK names
|
||||
* the same six events after the wire protocol; this is the rest of the app's vocabulary.
|
||||
@@ -116,35 +104,11 @@ function toStreamEvent(event: AgentStreamEvent): StreamEvent {
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* The machinery one conversation's live turn runs on. Deliberately not `$state`: nothing
|
||||
* renders from it, and an abort handle or a TypewriterReveal has no business behind a proxy.
|
||||
*/
|
||||
type TurnRuntime = {
|
||||
/** Stops this turn's `followJob`. */
|
||||
follow?: AbortController
|
||||
pollingInterval?: ReturnType<typeof setInterval>
|
||||
/** Stops the interval above at `POLL_MAX_MS`, and is cleared with it: one left armed by a
|
||||
* turn that ended early would fire into the next turn on the same conversation. */
|
||||
pollingDeadline?: ReturnType<typeof setTimeout>
|
||||
// What the turn has written so far — which row is open, and the text in it. Held here
|
||||
// rather than in the stream handler's locals: the typewriter reveals on animation
|
||||
// frames, long after the chunk that delivered the text was applied.
|
||||
turn: TurnState
|
||||
// The worker's events reach us in bursts — the provider batches tokens, and the SSE
|
||||
// endpoint ships whatever accumulated — so display is paced separately from arrival,
|
||||
// exactly as the session chat does it. Answer and thinking pace independently.
|
||||
replyReveal: TypewriterReveal
|
||||
reasoningReveal: TypewriterReveal
|
||||
}
|
||||
|
||||
function emptyStatus(): TurnStatus {
|
||||
return {
|
||||
isLoading: false,
|
||||
isWaitingForResponse: false,
|
||||
isDispatchingTurn: false,
|
||||
currentReasoning: '',
|
||||
isReasoningActive: false
|
||||
isDispatchingTurn: false
|
||||
}
|
||||
}
|
||||
|
||||
@@ -167,9 +131,6 @@ export class FlowChatManager {
|
||||
*/
|
||||
allowsParallelTurns = $state(false)
|
||||
|
||||
/** Bumped when the chat is re-pointed at another flow. Work started before that must not
|
||||
* write what it fetched into the chat that replaced it. */
|
||||
#generation = 0
|
||||
/** Each conversation's rows, live ones included, so a turn keeps writing while the
|
||||
* reader is in another chat. Doubles as the load cache: rows here are never re-fetched. */
|
||||
#rowsById = $state<Record<string, ChatMessage[]>>({})
|
||||
@@ -180,7 +141,16 @@ export class FlowChatManager {
|
||||
#pagedTo = $state<Record<string, number>>({})
|
||||
#hasMoreById = $state<Record<string, boolean>>({})
|
||||
#status = $state<Record<string, TurnStatus>>({})
|
||||
#runtime = new Map<string, TurnRuntime>()
|
||||
/**
|
||||
* The turn each conversation is on, and the only record of whether it is on one. A turn
|
||||
* is in here from the moment it takes the chat until it is ended.
|
||||
*
|
||||
* `$state` for the record, not for the turns in it: Svelte proxies plain objects and
|
||||
* arrays and leaves class instances alone, which is what `#isCurrent` needs. A proxied
|
||||
* turn would fail `===` against the one its own work holds, and its `AbortController`
|
||||
* would stop working behind the proxy.
|
||||
*/
|
||||
#turns = $state<Record<string, Turn>>({})
|
||||
|
||||
/** Row ids are temp- prefixed: the sweep after a run keeps only what the server stored. */
|
||||
#newRowId = () => 'temp-' + randomUUID()
|
||||
@@ -214,36 +184,56 @@ export class FlowChatManager {
|
||||
return this.#status[conversationId]
|
||||
}
|
||||
|
||||
#liveRuntime(conversationId: string): TurnRuntime {
|
||||
const existing = this.#runtime.get(conversationId)
|
||||
if (existing) return existing
|
||||
const created: TurnRuntime = {
|
||||
turn: emptyTurnState(conversationId),
|
||||
replyReveal: new TypewriterReveal({
|
||||
onReveal: (chunk) => this.#reveal(conversationId, 'answer', chunk),
|
||||
instant: prefersInstantReveal()
|
||||
}),
|
||||
reasoningReveal: new TypewriterReveal({
|
||||
onReveal: (chunk) => this.#reveal(conversationId, 'reasoning', chunk),
|
||||
instant: prefersInstantReveal()
|
||||
})
|
||||
}
|
||||
this.#runtime.set(conversationId, created)
|
||||
return created
|
||||
/**
|
||||
* Start a turn on this conversation, ending whatever it was on.
|
||||
*
|
||||
* One turn per conversation at a time: a second would write the same agent memory, and
|
||||
* the chat is held for the length of the first precisely so there cannot be one.
|
||||
*/
|
||||
#startTurn(conversationId: string, createdConversation = false): Turn {
|
||||
this.#turns[conversationId]?.end()
|
||||
const turn: Turn = new Turn({
|
||||
conversationId,
|
||||
createdConversation,
|
||||
onReveal: (kind, chunk) => this.#reveal(turn, kind, chunk),
|
||||
onPoll: () =>
|
||||
this.pollConversationMessages(conversationId, {
|
||||
isNewConversation: turn.listPending,
|
||||
turn
|
||||
})
|
||||
})
|
||||
this.#turns[conversationId] = turn
|
||||
return turn
|
||||
}
|
||||
|
||||
#reveal(conversationId: string, kind: 'answer' | 'reasoning', chunk: string) {
|
||||
const runtime = this.#liveRuntime(conversationId)
|
||||
/** The turn this conversation is on, if it is on one. */
|
||||
#turnOf(conversationId: string): Turn | undefined {
|
||||
return this.#turns[conversationId]
|
||||
}
|
||||
|
||||
/**
|
||||
* Whether work started by a turn may still write to the chat.
|
||||
*
|
||||
* The one question every write-back asks. A stream frame already in the reader's buffer,
|
||||
* a poll dispatched a tick ago, a settle that outlived its Stop: each resolves after the
|
||||
* turn that wanted it may be gone, and each would otherwise land in whatever turn is
|
||||
* there now.
|
||||
*/
|
||||
#isCurrent(turn: Turn): boolean {
|
||||
return this.#turns[turn.conversationId] === turn && !turn.ended
|
||||
}
|
||||
|
||||
#reveal(turn: Turn, kind: RevealKind, chunk: string) {
|
||||
if (!this.#isCurrent(turn)) return
|
||||
const step = appendRevealed(
|
||||
{ rows: this.#rowsOf(conversationId), state: runtime.turn },
|
||||
{ rows: this.#rowsOf(turn.conversationId), state: turn.transcript },
|
||||
kind,
|
||||
chunk,
|
||||
this.#newRowId
|
||||
)
|
||||
this.#rowsById[conversationId] = step.rows
|
||||
runtime.turn = step.state
|
||||
if (kind === 'reasoning')
|
||||
this.#liveStatus(conversationId).currentReasoning = step.state.reasoning
|
||||
this.#rowsById[turn.conversationId] = step.rows
|
||||
turn.transcript = step.state
|
||||
if (kind === 'reasoning') turn.reasoning = step.state.reasoning
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -255,18 +245,17 @@ export class FlowChatManager {
|
||||
* a retry replay the turn it is looking at rather than run the text again with whatever
|
||||
* the composer holds later.
|
||||
*/
|
||||
#nameTurnJob(conversationId: string, userRowId: string, jobId: string) {
|
||||
this.#liveStatus(conversationId).jobId = jobId
|
||||
this.#rowsById[conversationId] = this.#rowsOf(conversationId).map((row) =>
|
||||
#nameTurnJob(turn: Turn, userRowId: string, jobId: string) {
|
||||
// The row is named whatever became of the turn: it is the row this run was started
|
||||
// for, and Retry replays that run's arguments from the job named here, as does the
|
||||
// chip lane under the message. Both have to keep working on a turn that was stopped.
|
||||
this.#rowsById[turn.conversationId] = this.#rowsOf(turn.conversationId).map((row) =>
|
||||
row.id === userRowId ? { ...row, job_id: jobId } : row
|
||||
)
|
||||
}
|
||||
|
||||
/** Reveal everything buffered now, so the row is whole before the turn moves on. */
|
||||
#flushReveals(conversationId: string) {
|
||||
const runtime = this.#runtime.get(conversationId)
|
||||
runtime?.replyReveal.flush()
|
||||
runtime?.reasoningReveal.flush()
|
||||
// The turn is passed rather than looked up: a Stop landing while the run request was
|
||||
// in flight ends this one, and the conversation may be on another by now — naming
|
||||
// that one with this run would have Stop cancel a job it does not own.
|
||||
if (this.#isCurrent(turn)) turn.jobId = jobId
|
||||
}
|
||||
|
||||
// The open conversation's turn, which is what the composer and the transcript render.
|
||||
@@ -293,17 +282,17 @@ export class FlowChatManager {
|
||||
}
|
||||
get currentReasoning(): string {
|
||||
return this.selectedConversationId
|
||||
? this.#statusOf(this.selectedConversationId).currentReasoning
|
||||
? (this.#turnOf(this.selectedConversationId)?.reasoning ?? '')
|
||||
: ''
|
||||
}
|
||||
get isReasoningActive(): boolean {
|
||||
return this.selectedConversationId
|
||||
? this.#statusOf(this.selectedConversationId).isReasoningActive
|
||||
? (this.#turnOf(this.selectedConversationId)?.reasoningActive ?? false)
|
||||
: false
|
||||
}
|
||||
get currentJobId(): string | undefined {
|
||||
return this.selectedConversationId
|
||||
? this.#statusOf(this.selectedConversationId).jobId
|
||||
? this.#turnOf(this.selectedConversationId)?.jobId
|
||||
: undefined
|
||||
}
|
||||
|
||||
@@ -365,7 +354,7 @@ export class FlowChatManager {
|
||||
*/
|
||||
#forget(conversationId: string) {
|
||||
this.endTurn(conversationId)
|
||||
this.#runtime.delete(conversationId)
|
||||
delete this.#turns[conversationId]
|
||||
delete this.#rowsById[conversationId]
|
||||
delete this.#status[conversationId]
|
||||
delete this.#lastSeenCount[conversationId]
|
||||
@@ -392,6 +381,9 @@ export class FlowChatManager {
|
||||
if (rows[i].message_type !== 'user') continue
|
||||
return turnFailed(rows, i) ? 'error' : 'idle'
|
||||
}
|
||||
// A turn that wrote more rows than the page the transcript opened on leaves no
|
||||
// message of its own in it, and every row held is then part of that one turn.
|
||||
return turnFailed(rows, -1) ? 'error' : 'idle'
|
||||
}
|
||||
return 'idle'
|
||||
}
|
||||
@@ -467,50 +459,43 @@ export class FlowChatManager {
|
||||
*/
|
||||
onTurnSettled?: (conversationId: string) => void
|
||||
|
||||
/**
|
||||
* End this turn, if it is still the one its conversation is on.
|
||||
*
|
||||
* What every caller holding a turn means by "end it". A turn that is no longer current
|
||||
* has already been ended by whatever replaced it, so there is nothing to do — and
|
||||
* ending *the conversation's* turn would end the one that replaced it, freeing a chat
|
||||
* whose run is still going.
|
||||
*/
|
||||
#endTurnIfCurrent(turn: Turn, options?: { settled?: boolean }) {
|
||||
if (this.#isCurrent(turn)) this.endTurn(turn.conversationId, options)
|
||||
}
|
||||
|
||||
/** Stop following one conversation's turn and forget what it was mid-way through. */
|
||||
endTurn(conversationId: string, options?: { settled?: boolean }) {
|
||||
const runtime = this.#runtime.get(conversationId)
|
||||
if (runtime) {
|
||||
runtime.replyReveal.reset()
|
||||
runtime.reasoningReveal.reset()
|
||||
runtime.turn = emptyTurnState(conversationId)
|
||||
runtime.follow?.abort()
|
||||
runtime.follow = undefined
|
||||
this.stopPolling(conversationId)
|
||||
}
|
||||
// Out of the record first: the turn is no longer this conversation's, so anything it
|
||||
// started that resolves from here on fails `#isCurrent` and writes nothing.
|
||||
const turn = this.#turns[conversationId]
|
||||
delete this.#turns[conversationId]
|
||||
turn?.end()
|
||||
const status = this.#liveStatus(conversationId)
|
||||
status.currentReasoning = ''
|
||||
status.isReasoningActive = false
|
||||
status.isLoading = false
|
||||
status.isWaitingForResponse = false
|
||||
// Including a hold taken before a turn had a job — an upload, or a load working out
|
||||
// whether this chat has a run in flight. Stop is the way out of those too, and it
|
||||
// reaches them only here.
|
||||
status.isDispatchingTurn = false
|
||||
status.jobId = undefined
|
||||
if (options?.settled) this.onTurnSettled?.(conversationId)
|
||||
}
|
||||
|
||||
/**
|
||||
* Stop following every turn, and forget what this chat holds about one flow.
|
||||
*
|
||||
* Called when the chat goes away — which includes being re-pointed at another flow or
|
||||
* workspace, since the route component is reused between them. Keeping the selection
|
||||
* there would open the next flow on the previous one's conversation, and send its
|
||||
* messages into that transcript and its agent's memory.
|
||||
* Stop following every turn. Called when the chat goes away, which is the only way one
|
||||
* flow's chat is left: a manager belongs to the panel built for one flow in one
|
||||
* workspace, and `FlowChat` replaces that panel rather than re-pointing it.
|
||||
*/
|
||||
cleanup() {
|
||||
this.#generation++
|
||||
for (const conversationId of Object.keys(this.#status)) this.endTurn(conversationId)
|
||||
this.selectedConversationId = undefined
|
||||
this.#rowsById = {}
|
||||
this.#pagedTo = {}
|
||||
this.#hasMoreById = {}
|
||||
this.#status = {}
|
||||
this.#runtime.clear()
|
||||
this.#lastSeenCount = {}
|
||||
this.conversations = []
|
||||
this.inputMessage = ''
|
||||
const held = new Set([...Object.keys(this.#status), ...Object.keys(this.#turns)])
|
||||
for (const conversationId of held) this.endTurn(conversationId)
|
||||
}
|
||||
|
||||
// Public methods for component to call
|
||||
@@ -598,11 +583,9 @@ export class FlowChatManager {
|
||||
*/
|
||||
async selectLatestConversation() {
|
||||
if (this.selectedConversationId || !this.#workspace() || !this.#path) return
|
||||
const startedIn = this.#generation
|
||||
const [latest] = await this.loadConversations(1, 1)
|
||||
// Re-checked after the await: a message sent meanwhile has already opened its own, or
|
||||
// the chat has been re-pointed and this is the previous flow's latest conversation.
|
||||
if (!latest || this.selectedConversationId || startedIn !== this.#generation) return
|
||||
// Re-checked after the await: a message sent meanwhile has already opened its own.
|
||||
if (!latest || this.selectedConversationId) return
|
||||
await this.selectConversation(latest.id)
|
||||
}
|
||||
|
||||
@@ -679,11 +662,12 @@ export class FlowChatManager {
|
||||
}
|
||||
|
||||
/** No job came of the send, so nothing is in flight and nothing is waiting on one. */
|
||||
#turnFailedToStart(conversationId: string) {
|
||||
const status = this.#liveStatus(conversationId)
|
||||
status.isLoading = false
|
||||
status.isWaitingForResponse = false
|
||||
status.isDispatchingTurn = false
|
||||
#turnFailedToStart(turn: Turn) {
|
||||
// The turn goes with the holds it took: nothing ran, so there is nothing for it to be
|
||||
// current for. `IfCurrent` because the run request is not tied to the turn's signal —
|
||||
// a Stop and a second send during it leave this one answering for a turn that is now
|
||||
// streaming, and tearing that down would free a chat whose run carries on.
|
||||
this.#endTurnIfCurrent(turn)
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -760,7 +744,11 @@ export class FlowChatManager {
|
||||
if (!this.#workspace() || !conversationId) {
|
||||
return
|
||||
}
|
||||
const jobId = this.#statusOf(conversationId).jobId
|
||||
// Held, not looked up again: the turn can settle while the cancel request is in
|
||||
// flight, and the failure path must not hand this job to whatever turn is running by
|
||||
// then — it would name that turn with the wrong run.
|
||||
const turn = this.#turnOf(conversationId)
|
||||
const jobId = turn?.jobId
|
||||
|
||||
try {
|
||||
if (jobId) {
|
||||
@@ -771,14 +759,18 @@ export class FlowChatManager {
|
||||
})
|
||||
sendUserToast(`Job ${jobId} cancelled`)
|
||||
}
|
||||
this.endTurn(conversationId)
|
||||
// The turn that settled during the cancel request has already gone, and the one
|
||||
// that replaced it owns a run of its own.
|
||||
if (!turn || this.#isCurrent(turn)) this.endTurn(conversationId)
|
||||
} catch (error) {
|
||||
// The run may well still be going, and freeing the chat would let the next turn
|
||||
// write the same agent memory. It is left to the job to say when it is over.
|
||||
console.error('Error cancelling job:', error)
|
||||
sendUserToast('Could not stop the run; waiting for it to finish', true)
|
||||
if (jobId) this.#followFromJob(conversationId, jobId)
|
||||
else this.endTurn(conversationId)
|
||||
// The turn's own follower is still on the job and is what will settle it; another
|
||||
// would ask the same question twice. With no turn left there is nothing following
|
||||
// it, and the chat is released instead.
|
||||
if (!turn || !this.#isCurrent(turn)) this.endTurn(conversationId)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -790,7 +782,6 @@ export class FlowChatManager {
|
||||
private async loadConversations(page: number, perPage: number) {
|
||||
if (!this.#workspace() || !this.#path) return []
|
||||
|
||||
const startedIn = this.#generation
|
||||
try {
|
||||
const response = await FlowConversationsService.listFlowConversations({
|
||||
workspace: this.#workspace()!,
|
||||
@@ -799,10 +790,6 @@ export class FlowChatManager {
|
||||
page: page,
|
||||
perPage: perPage
|
||||
})
|
||||
// The list is bound straight into the sidebar, so one fetched for the flow just
|
||||
// left would otherwise become the list of the flow that replaced it.
|
||||
if (startedIn !== this.#generation) return []
|
||||
|
||||
return response
|
||||
} catch (error) {
|
||||
console.error('Failed to load conversations:', error)
|
||||
@@ -835,7 +822,6 @@ export class FlowChatManager {
|
||||
try {
|
||||
const previousScrollHeight = this.messagesContainer?.scrollHeight || 0
|
||||
|
||||
const startedIn = this.#generation
|
||||
const response = await FlowConversationsService.listConversationMessages({
|
||||
workspace: this.#workspace()!,
|
||||
conversationId: conversationIdToUse,
|
||||
@@ -843,10 +829,6 @@ export class FlowChatManager {
|
||||
perPage: this.#perPage
|
||||
})
|
||||
|
||||
// The chat was re-pointed while this was in flight; these rows belong to a flow it
|
||||
// no longer shows, and the hold below would be released for the wrong one.
|
||||
if (startedIn !== this.#generation) return
|
||||
|
||||
if (reset) {
|
||||
this.#rowsById[conversationIdToUse] = response
|
||||
this.#pagedTo[conversationIdToUse] = 1
|
||||
@@ -926,12 +908,11 @@ export class FlowChatManager {
|
||||
// Polling
|
||||
private async pollConversationMessages(
|
||||
conversationId: string,
|
||||
options?: { isNewConversation?: boolean; removeTempMessages?: boolean }
|
||||
options?: { isNewConversation?: boolean; removeTempMessages?: boolean; turn?: Turn }
|
||||
) {
|
||||
if (!this.#workspace()) return
|
||||
|
||||
try {
|
||||
const startedIn = this.#generation
|
||||
// Paged, not one request: the endpoint answers oldest-first with a limit, so a turn
|
||||
// that wrote more rows than a page — an agent calling several tools a round — would
|
||||
// hand back its earliest and leave its answer behind, and the sweep below drops the
|
||||
@@ -943,6 +924,9 @@ export class FlowChatManager {
|
||||
let afterSeq = this.getLastPersistedMessageSeq(conversationId) ?? 0
|
||||
let readWhole = false
|
||||
for (let request = 0; request < POLL_MAX_REQUESTS; request++) {
|
||||
// A read for a turn stops with it: a page boundary is where a poll walking a long
|
||||
// conversation notices the chat it was reading for has gone.
|
||||
if (options?.turn && !this.#isCurrent(options.turn)) return
|
||||
const batch = await FlowConversationsService.listConversationMessages({
|
||||
workspace: this.#workspace()!,
|
||||
conversationId: conversationId,
|
||||
@@ -950,9 +934,6 @@ export class FlowChatManager {
|
||||
perPage: POLL_PAGE_SIZE,
|
||||
afterSeq
|
||||
})
|
||||
// An interval tick already dispatched outlives `clearInterval`, and a turn's final
|
||||
// poll outlives its abort — either would put a forgotten flow's rows back.
|
||||
if (startedIn !== this.#generation) return
|
||||
response.push(...batch)
|
||||
if (batch.length < POLL_PAGE_SIZE) {
|
||||
readWhole = true
|
||||
@@ -966,12 +947,16 @@ export class FlowChatManager {
|
||||
`Stopped reading conversation ${conversationId} at seq ${afterSeq}: a full ` +
|
||||
`page of ${POLL_PAGE_SIZE} rows did not move the cursor`
|
||||
)
|
||||
this.stopPolling(conversationId)
|
||||
this.#turnOf(conversationId)?.stopPolling()
|
||||
break
|
||||
}
|
||||
afterSeq = furthest
|
||||
}
|
||||
|
||||
// The read is done; the turn it was for may not be. Everything below writes to the
|
||||
// transcript, so it asks the same question every other write-back asks.
|
||||
if (options?.turn && !this.#isCurrent(options.turn)) return
|
||||
|
||||
if (!readWhole && options?.removeTempMessages) {
|
||||
// The last poll of a turn, and no tick comes after it to carry on from where this
|
||||
// one stopped. Its temp rows have to stay, being the only copy of what the read
|
||||
@@ -987,6 +972,8 @@ export class FlowChatManager {
|
||||
|
||||
if (options?.isNewConversation) {
|
||||
await this.refreshConversations()
|
||||
// The list has the row now, so the ticks after this one do not ask again.
|
||||
if (options.turn) options.turn.listPending = false
|
||||
}
|
||||
|
||||
// Written to this conversation's rows, not the open one's: a turn keeps landing
|
||||
@@ -1011,29 +998,6 @@ export class FlowChatManager {
|
||||
}
|
||||
}
|
||||
|
||||
private startPolling(conversationId: string, isNewConversation?: boolean) {
|
||||
const runtime = this.#liveRuntime(conversationId)
|
||||
if (runtime.pollingInterval) return
|
||||
runtime.pollingInterval = setInterval(() => {
|
||||
this.pollConversationMessages(conversationId, { isNewConversation })
|
||||
}, 500) // Poll every 0.5 seconds
|
||||
runtime.pollingDeadline = setTimeout(() => {
|
||||
this.stopPolling(conversationId)
|
||||
}, POLL_MAX_MS)
|
||||
}
|
||||
|
||||
private stopPolling(conversationId: string) {
|
||||
const runtime = this.#runtime.get(conversationId)
|
||||
if (runtime?.pollingInterval) {
|
||||
clearInterval(runtime.pollingInterval)
|
||||
runtime.pollingInterval = undefined
|
||||
}
|
||||
if (runtime?.pollingDeadline) {
|
||||
clearTimeout(runtime.pollingDeadline)
|
||||
runtime.pollingDeadline = undefined
|
||||
}
|
||||
}
|
||||
|
||||
// Message sending
|
||||
/**
|
||||
* Send `inputMessage` as a turn. `onUserRow` is called with the id of the row added
|
||||
@@ -1075,8 +1039,9 @@ export class FlowChatManager {
|
||||
|
||||
const isNewConversation = this.#rowsOf(currentConversationId).length === 0
|
||||
|
||||
// Reset state for new message
|
||||
this.stopPolling(currentConversationId)
|
||||
// The turn begins here, before the run: it is what holds the chat for the length of
|
||||
// the send, and what a Stop during it has to find.
|
||||
const turn = this.#startTurn(currentConversationId, isNewConversation)
|
||||
|
||||
const userMessage: ChatMessage = {
|
||||
id: `temp-${randomUUID()}`,
|
||||
@@ -1105,16 +1070,14 @@ export class FlowChatManager {
|
||||
if (this.#useStreaming && this.#path) {
|
||||
started = await this.handleStreamingMessage(
|
||||
messageContent,
|
||||
currentConversationId,
|
||||
isNewConversation,
|
||||
turn,
|
||||
userMessage.id,
|
||||
additionalInputs
|
||||
)
|
||||
} else {
|
||||
started = await this.handlePollingMessage(
|
||||
messageContent,
|
||||
currentConversationId,
|
||||
isNewConversation,
|
||||
turn,
|
||||
userMessage.id,
|
||||
additionalInputs
|
||||
)
|
||||
@@ -1125,7 +1088,7 @@ export class FlowChatManager {
|
||||
// A turn that never started leaves nothing to wait for. Said here as well as in
|
||||
// the finally because the streaming path keeps `isLoading` for its own stream,
|
||||
// and without this the composer and the sidebar stay locked until a reload.
|
||||
this.#turnFailedToStart(currentConversationId)
|
||||
this.#turnFailedToStart(turn)
|
||||
started = false
|
||||
} finally {
|
||||
if (!this.#useStreaming) {
|
||||
@@ -1163,42 +1126,33 @@ export class FlowChatManager {
|
||||
/** Answers whether a job was actually started. */
|
||||
private async handleStreamingMessage(
|
||||
messageContent: string,
|
||||
currentConversationId: string,
|
||||
isNewConversation: boolean,
|
||||
turn: Turn,
|
||||
userRowId: string,
|
||||
additionalInputs?: Record<string, any>
|
||||
): Promise<boolean> {
|
||||
const runtime = this.#liveRuntime(currentConversationId)
|
||||
runtime.follow?.abort()
|
||||
|
||||
// Track stream state for this message
|
||||
runtime.turn = emptyTurnState(currentConversationId)
|
||||
runtime.replyReveal.reset()
|
||||
runtime.reasoningReveal.reset()
|
||||
|
||||
const currentConversationId = turn.conversationId
|
||||
try {
|
||||
const jobId = await this.#onRunFlow?.(messageContent, currentConversationId, additionalInputs)
|
||||
if (!jobId) {
|
||||
console.error('No jobId returned from onRunFlow')
|
||||
this.#turnFailedToStart(currentConversationId)
|
||||
this.#turnFailedToStart(turn)
|
||||
return false
|
||||
}
|
||||
// What Stop cancels: the flow job, never the streaming step's own sub-job — which
|
||||
// would leave the steps after the agent running.
|
||||
this.#nameTurnJob(currentConversationId, userRowId, jobId)
|
||||
this.#nameTurnJob(turn, userRowId, jobId)
|
||||
|
||||
// start polling
|
||||
this.startPolling(currentConversationId, isNewConversation)
|
||||
turn.startPolling()
|
||||
|
||||
// Runs for the length of the turn, and reports its own failures.
|
||||
void this.#followJob(currentConversationId, jobId)
|
||||
void this.#followJob(turn, jobId)
|
||||
} catch (error) {
|
||||
// The run request itself, which the deployed page's launcher throws from. The
|
||||
// stream is not live yet, so no turn ran.
|
||||
console.error('Stream connection error:', error)
|
||||
sendUserToast('Failed to connect to stream', true)
|
||||
this.endTurn(currentConversationId)
|
||||
this.#turnFailedToStart(currentConversationId)
|
||||
this.#turnFailedToStart(turn)
|
||||
return false
|
||||
}
|
||||
return true
|
||||
@@ -1217,21 +1171,11 @@ export class FlowChatManager {
|
||||
* recovered is handed to its job rather than ended: the flow may still be running, and
|
||||
* freeing the composer would let the next turn write the same agent memory.
|
||||
*/
|
||||
async #followJob(currentConversationId: string, jobId: string) {
|
||||
const runtime = this.#liveRuntime(currentConversationId)
|
||||
const status = this.#liveStatus(currentConversationId)
|
||||
runtime.follow?.abort()
|
||||
const controller = new AbortController()
|
||||
runtime.follow = controller
|
||||
|
||||
async #followJob(turn: Turn, jobId: string) {
|
||||
const currentConversationId = turn.conversationId
|
||||
const signal = turn.signal
|
||||
const api = this.#chatApi()
|
||||
|
||||
// Kept across attempts so a reconnect resumes after what is already on screen. A
|
||||
// retried agent step gets its own stream, which `followJob` detects and restarts from
|
||||
// — but only within one call, so an offset carried into a *new* call can index the
|
||||
// previous sub-job. Resuming too far in drops chunks the final row poll then repairs;
|
||||
// not resuming at all would duplicate the answer on screen, which nothing repairs.
|
||||
let streamOffset: number | undefined
|
||||
// Counted since the last attempt that delivered anything, so a long turn blipping
|
||||
// once an hour is not the same as an endpoint that has gone.
|
||||
let sinceProgress = 0
|
||||
@@ -1239,63 +1183,64 @@ export class FlowChatManager {
|
||||
let delivered = false
|
||||
try {
|
||||
for await (const update of followJob(api, jobId, {
|
||||
signal: controller.signal,
|
||||
streamOffset,
|
||||
onOffset: (offset) => (streamOffset = offset)
|
||||
signal,
|
||||
// Kept on the turn so a reconnect resumes after what is already on screen. A
|
||||
// retried agent step gets its own stream, which `followJob` detects and
|
||||
// restarts from. Resuming too far in drops chunks the final row poll then
|
||||
// repairs; not resuming duplicates the answer, which nothing repairs.
|
||||
streamOffset: turn.streamOffset,
|
||||
onOffset: (offset) => (turn.streamOffset = offset)
|
||||
})) {
|
||||
// Frames already buffered by the SSE reader keep arriving after an abort, and
|
||||
// `endTurn` has reset the transcript state they would be applied to.
|
||||
if (controller.signal.aborted) return
|
||||
// the turn they would be applied to is gone by then.
|
||||
if (!this.#isCurrent(turn)) return
|
||||
delivered = true
|
||||
if (update.type === 'stream') {
|
||||
// Stop polling since we are receiving last step streaming
|
||||
this.stopPolling(currentConversationId)
|
||||
turn.stopPolling()
|
||||
// One chunk can carry several events, so each is applied in turn: a
|
||||
// chunk holding a call and its result must produce both.
|
||||
for (const streamed of update.events) {
|
||||
const event = toStreamEvent(streamed)
|
||||
if (event.kind === 'reasoning') {
|
||||
status.isReasoningActive = true
|
||||
runtime.reasoningReveal.push(event.content)
|
||||
turn.reasoningActive = true
|
||||
turn.pushReasoning(event.content)
|
||||
} else if (event.kind === 'token') {
|
||||
runtime.replyReveal.push(event.content)
|
||||
turn.pushAnswer(event.content)
|
||||
} else {
|
||||
// Whatever the pacing still holds belongs to the row before the tool —
|
||||
// thinking that led straight to the call included — so it is revealed
|
||||
// before the event that closes that row.
|
||||
if (event.kind === 'tool_call' || event.kind === 'tool_execution') {
|
||||
this.#flushReveals(currentConversationId)
|
||||
status.currentReasoning = ''
|
||||
status.isReasoningActive = false
|
||||
turn.flushReveals()
|
||||
turn.reasoning = ''
|
||||
turn.reasoningActive = false
|
||||
}
|
||||
const step = applyStreamEvent(
|
||||
{ rows: this.#rowsOf(currentConversationId), state: runtime.turn },
|
||||
{ rows: this.#rowsOf(currentConversationId), state: turn.transcript },
|
||||
event,
|
||||
this.#newRowId
|
||||
)
|
||||
this.#rowsById[currentConversationId] = step.rows
|
||||
runtime.turn = step.state
|
||||
turn.transcript = step.state
|
||||
}
|
||||
}
|
||||
continue
|
||||
}
|
||||
// Anything still buffered would be dropped by the temp-row sweep below.
|
||||
this.#flushReveals(currentConversationId)
|
||||
turn.flushReveals()
|
||||
// Do a final poll to get all messages from database
|
||||
const startedIn = this.#generation
|
||||
await this.pollConversationMessages(currentConversationId, {
|
||||
removeTempMessages: true
|
||||
removeTempMessages: true,
|
||||
turn
|
||||
})
|
||||
// Settling releases this conversation's queue, and after a re-point that
|
||||
// would run the message against the flow now loaded.
|
||||
if (controller.signal.aborted || startedIn !== this.#generation) return
|
||||
this.endTurn(currentConversationId, { settled: true })
|
||||
this.#endTurnIfCurrent(turn, { settled: true })
|
||||
}
|
||||
return
|
||||
} catch (error) {
|
||||
// A Stop, a conversation's turn ending, or the chat going away: whoever aborted
|
||||
// it has already torn the turn down.
|
||||
if (controller.signal.aborted) return
|
||||
// A Stop, a conversation's turn ending, or the chat going away: whoever ended
|
||||
// the turn has already torn it down.
|
||||
if (!this.#isCurrent(turn)) return
|
||||
console.error('Error following the flow job:', error)
|
||||
if (delivered) sinceProgress = 0
|
||||
if (sinceProgress < FOLLOW_RETRIES) {
|
||||
@@ -1303,7 +1248,7 @@ export class FlowChatManager {
|
||||
setTimeout(resolve, FOLLOW_RETRY_DELAY_MS * 2 ** sinceProgress)
|
||||
)
|
||||
sinceProgress++
|
||||
if (controller.signal.aborted) return
|
||||
if (!this.#isCurrent(turn)) return
|
||||
continue
|
||||
}
|
||||
// Out of attempts, and the run is the only thing that knows whether it is over.
|
||||
@@ -1312,7 +1257,7 @@ export class FlowChatManager {
|
||||
`Lost the live answer for this turn; it will land when the run finishes. ${reason}`,
|
||||
true
|
||||
)
|
||||
void this.#settleFromJob(currentConversationId, jobId, api, controller.signal)
|
||||
void this.#settleFromJob(turn, jobId)
|
||||
return
|
||||
}
|
||||
}
|
||||
@@ -1329,74 +1274,151 @@ export class FlowChatManager {
|
||||
* The chat is held from the moment the question is asked, not from the answer: a send
|
||||
* accepted during that round trip is the very thing this exists to prevent.
|
||||
*
|
||||
* Only page one is loaded, so a turn with more rows than a page — a tool-heavy agent —
|
||||
* has no user row here to name its job, and is not picked up. Finding it needs the
|
||||
* server to answer which turn is running rather than this inferring it from a page.
|
||||
* Which turn is running is inferred from the rows rather than asked of the server, so a
|
||||
* turn whose own rows fill the page the transcript opened on takes a walk back through
|
||||
* older pages to find the message that started it (see `#newestUserRow`).
|
||||
*/
|
||||
async #resumeRunningTurn(conversationId: string) {
|
||||
const status = this.#liveStatus(conversationId)
|
||||
let tookOver = false
|
||||
// The turn starts here, before anything is awaited, so a chat torn down or stopped
|
||||
// while this is in flight can still stop what it would start.
|
||||
const turn = this.#startTurn(conversationId)
|
||||
try {
|
||||
tookOver = await this.#takeOverRunningTurn(conversationId, status)
|
||||
} finally {
|
||||
// Held by the load that called this; a turn taken over is busy on its own flags
|
||||
// by now, and one that was not leaves nothing for Stop to cancel.
|
||||
status.isDispatchingTurn = false
|
||||
if (!tookOver) status.jobId = undefined
|
||||
await this.#takeOverRunningTurn(turn, status)
|
||||
// The hold this attempt was given back, and only it: a turn that replaced this one
|
||||
// took a hold of its own, and releasing that would unlock a chat with a run in it.
|
||||
// A turn taken over is busy on its own flags by now.
|
||||
if (!this.#turnOf(conversationId) || this.#isCurrent(turn)) {
|
||||
status.isDispatchingTurn = false
|
||||
}
|
||||
} catch (error) {
|
||||
// Whether a run owns this conversation's agent memory is exactly what could not be
|
||||
// read, so the chat stays held rather than taking a message that would write it a
|
||||
// second time. Stop is the reader's way out.
|
||||
console.error('Could not tell whether a conversation had a run in flight:', error)
|
||||
if (!this.#turnOf(conversationId)) status.isDispatchingTurn = true
|
||||
}
|
||||
}
|
||||
|
||||
async #takeOverRunningTurn(conversationId: string, status: TurnStatus): Promise<boolean> {
|
||||
// A turn already in flight here owns the chat; only the load's own hold is set.
|
||||
if (status.isLoading || status.isWaitingForResponse) return false
|
||||
if (!this.#workspace()) return false
|
||||
// The newest turn is the only one that can still be running; the rows arrive oldest
|
||||
// first, and a user row is named with its flow job by the run that created it.
|
||||
const rows = this.#rowsOf(conversationId)
|
||||
const startedAt = rows.findLastIndex((row) => row.message_type === 'user' && row.job_id)
|
||||
const jobId = startedAt >= 0 ? rows[startedAt].job_id : undefined
|
||||
if (!jobId) return false
|
||||
// Named before the question, not after it: the chat is held from here, so Stop is on
|
||||
// offer, and Stop with no job cancels nothing while the run carries on.
|
||||
status.jobId = jobId
|
||||
/**
|
||||
* The message that started the newest turn in a conversation, which is the only turn
|
||||
* that can still be running.
|
||||
*
|
||||
* The transcript opens on the newest page, and one turn can fill it on its own: an agent
|
||||
* writes a row per round and per tool call, so a turn that used a lot of tools pushes the
|
||||
* message that started it further back than a page. Pages are walked back until it turns
|
||||
* up, since without it a reload cannot tell a finished conversation from one still
|
||||
* running, and would offer a composer that writes a second turn into the same memory.
|
||||
*/
|
||||
async #newestUserRow(
|
||||
conversationId: string,
|
||||
signal?: AbortSignal
|
||||
): Promise<ChatMessage | undefined> {
|
||||
const held = this.#rowsOf(conversationId).findLast((row) => row.message_type === 'user')
|
||||
if (held) return held
|
||||
const from = (this.#pagedTo[conversationId] ?? 1) + 1
|
||||
for (let page = from; page < from + RESUME_MAX_PAGES; page++) {
|
||||
if (signal?.aborted) return undefined
|
||||
const batch = await FlowConversationsService.listConversationMessages({
|
||||
workspace: this.#workspace()!,
|
||||
conversationId,
|
||||
page,
|
||||
perPage: this.#perPage
|
||||
})
|
||||
const found = batch.findLast((row) => row.message_type === 'user')
|
||||
if (found) return found
|
||||
// The start of the conversation, which a user row always opens — so this only runs
|
||||
// out of rows on one written by something other than a chat.
|
||||
if (batch.length < this.#perPage) return undefined
|
||||
}
|
||||
console.warn(
|
||||
`Gave up looking for the message that started the newest turn of conversation ` +
|
||||
`${conversationId} after ${RESUME_MAX_PAGES} pages`
|
||||
)
|
||||
return undefined
|
||||
}
|
||||
|
||||
const runtime = this.#liveRuntime(conversationId)
|
||||
async #takeOverRunningTurn(turn: Turn, status: TurnStatus): Promise<boolean> {
|
||||
const conversationId = turn.conversationId
|
||||
// A turn already in flight here owns the chat; only the load's own hold is set.
|
||||
if (status.isLoading || status.isWaitingForResponse) {
|
||||
this.#endTurnIfCurrent(turn)
|
||||
return false
|
||||
}
|
||||
if (!this.#workspace()) {
|
||||
this.#endTurnIfCurrent(turn)
|
||||
return false
|
||||
}
|
||||
|
||||
// The newest turn is the only one that can still be running, and the message that
|
||||
// started it is named with its flow job by the run that created it.
|
||||
let startedBy: ChatMessage | undefined
|
||||
try {
|
||||
startedBy = await this.#newestUserRow(conversationId, turn.signal)
|
||||
} catch (error) {
|
||||
// The caller decides what happens to the chat's hold; what this owns is the turn
|
||||
// it started, which answered nothing and has no run to follow.
|
||||
this.#endTurnIfCurrent(turn)
|
||||
throw error
|
||||
}
|
||||
// Torn down or stopped while the walk was in flight. The run is left alone either
|
||||
// way: leaving a chat has never cancelled one, and this cannot tell the two apart.
|
||||
// A turn that is no longer this conversation's has already been ended by whatever
|
||||
// replaced it; ending "the conversation's turn" here would end that one.
|
||||
if (!this.#isCurrent(turn)) return false
|
||||
if (!startedBy?.job_id) {
|
||||
this.endTurn(conversationId)
|
||||
return false
|
||||
}
|
||||
const jobId = startedBy.job_id
|
||||
const api = this.#chatApi()
|
||||
// Taken before the question so a chat torn down while it is in flight can still stop
|
||||
// what the answer would start.
|
||||
runtime.follow?.abort()
|
||||
const controller = new AbortController()
|
||||
runtime.follow = controller
|
||||
// Named before the question, not after it: Stop with no job cancels nothing while the
|
||||
// run carries on.
|
||||
turn.jobId = jobId
|
||||
{
|
||||
// Asked until answered, for the same reason a turn is: an API that cannot be
|
||||
// reached has not said the run is over, and a chat freed on that guess takes a
|
||||
// message that writes the same agent memory. Stop is on screen throughout.
|
||||
for (;;) {
|
||||
if (controller.signal.aborted) return false
|
||||
if (!this.#isCurrent(turn)) return false
|
||||
try {
|
||||
const { completed } = await api.getCompletedResult(jobId, controller.signal)
|
||||
if (completed) return false
|
||||
const { completed } = await api.getCompletedResult(jobId, turn.signal)
|
||||
// Nothing to take over, so the turn opened to ask the question goes with it.
|
||||
if (completed) {
|
||||
this.#endTurnIfCurrent(turn)
|
||||
return false
|
||||
}
|
||||
break
|
||||
} catch (error) {
|
||||
if (controller.signal.aborted) return false
|
||||
if (!this.#isCurrent(turn)) return false
|
||||
console.error('Could not tell whether a conversation had a run in flight:', error)
|
||||
await new Promise((resolve) => setTimeout(resolve, SETTLE_POLL_MS))
|
||||
}
|
||||
}
|
||||
}
|
||||
if (controller.signal.aborted) return false
|
||||
if (!this.#isCurrent(turn)) return false
|
||||
|
||||
status.isLoading = true
|
||||
status.isWaitingForResponse = true
|
||||
// Rows written while this turn ran are replayed by the stream it is about to
|
||||
// re-attach to — an agent persists one per round — and the transcript has no way to
|
||||
// tell a replayed round from a second one. The final poll brings them back.
|
||||
if (this.#useStreaming) this.#rowsById[conversationId] = rows.slice(0, startedAt + 1)
|
||||
this.startPolling(conversationId)
|
||||
//
|
||||
// The message that started the turn stays, and when the turn was long enough to push
|
||||
// it off the page the transcript opened on it is all that is left: polls only ever
|
||||
// add what a turn wrote, never the questions, so a transcript emptied here would come
|
||||
// back as answers with nothing asking them. It is also the cursor the next poll reads
|
||||
// forward from.
|
||||
if (this.#useStreaming) {
|
||||
void this.#followJob(conversationId, jobId)
|
||||
const held = this.#rowsOf(conversationId)
|
||||
const startedAt = held.findLastIndex((row) => row.id === startedBy.id)
|
||||
this.#rowsById[conversationId] = startedAt >= 0 ? held.slice(0, startedAt + 1) : [startedBy]
|
||||
}
|
||||
turn.startPolling()
|
||||
if (this.#useStreaming) {
|
||||
void this.#followJob(turn, jobId)
|
||||
} else {
|
||||
void this.#settleFromJob(conversationId, jobId, api, controller.signal)
|
||||
void this.#settleFromJob(turn, jobId)
|
||||
}
|
||||
return true
|
||||
}
|
||||
@@ -1413,15 +1435,6 @@ export class FlowChatManager {
|
||||
})
|
||||
}
|
||||
|
||||
/** Hand a turn to its job, under this conversation's abort handle. */
|
||||
#followFromJob(conversationId: string, jobId: string) {
|
||||
const runtime = this.#liveRuntime(conversationId)
|
||||
runtime.follow?.abort()
|
||||
const controller = new AbortController()
|
||||
runtime.follow = controller
|
||||
void this.#settleFromJob(conversationId, jobId, this.#chatApi(), controller.signal)
|
||||
}
|
||||
|
||||
/**
|
||||
* Wait out a turn with no stream to follow, on the job itself.
|
||||
*
|
||||
@@ -1433,57 +1446,53 @@ export class FlowChatManager {
|
||||
* Deliberately not `waitJob`, which stops answering without settling: on its fifth
|
||||
* failed request it cancels and returns, leaving its promise pending forever.
|
||||
*/
|
||||
async #settleFromJob(
|
||||
conversationId: string,
|
||||
jobId: string,
|
||||
api: WindmillChatApi,
|
||||
signal: AbortSignal
|
||||
) {
|
||||
while (!signal.aborted) {
|
||||
async #settleFromJob(turn: Turn, jobId: string) {
|
||||
const conversationId = turn.conversationId
|
||||
const api = this.#chatApi()
|
||||
const signal = turn.signal
|
||||
while (this.#isCurrent(turn)) {
|
||||
try {
|
||||
const { completed } = await api.getCompletedResult(jobId, signal)
|
||||
if (completed) break
|
||||
} catch (error) {
|
||||
if (signal.aborted) return
|
||||
if (!this.#isCurrent(turn)) return
|
||||
console.error('Could not read the flow job while settling a turn:', error)
|
||||
}
|
||||
await new Promise((resolve) => setTimeout(resolve, SETTLE_POLL_MS))
|
||||
}
|
||||
if (signal.aborted) return
|
||||
const startedIn = this.#generation
|
||||
if (!this.#isCurrent(turn)) return
|
||||
try {
|
||||
await this.pollConversationMessages(conversationId, { removeTempMessages: true })
|
||||
await this.pollConversationMessages(conversationId, { removeTempMessages: true, turn })
|
||||
} catch {}
|
||||
// `endTurn` would put this conversation's status back into a map the re-point emptied.
|
||||
if (signal.aborted || startedIn !== this.#generation) return
|
||||
this.endTurn(conversationId, { settled: true })
|
||||
this.#endTurnIfCurrent(turn, { settled: true })
|
||||
}
|
||||
|
||||
/** Answers whether a job was actually started. */
|
||||
private async handlePollingMessage(
|
||||
messageContent: string,
|
||||
currentConversationId: string,
|
||||
isNewConversation: boolean,
|
||||
turn: Turn,
|
||||
userRowId: string,
|
||||
additionalInputs?: Record<string, any>
|
||||
): Promise<boolean> {
|
||||
const currentConversationId = turn.conversationId
|
||||
const jobId = await this.#onRunFlow?.(messageContent, currentConversationId, additionalInputs)
|
||||
if (!jobId) {
|
||||
console.error('No jobId returned from onRunFlow')
|
||||
this.#turnFailedToStart(currentConversationId)
|
||||
this.#turnFailedToStart(turn)
|
||||
return false
|
||||
}
|
||||
|
||||
// Store the current job ID so it can be cancelled
|
||||
this.#nameTurnJob(currentConversationId, userRowId, jobId)
|
||||
this.#nameTurnJob(turn, userRowId, jobId)
|
||||
|
||||
if (isNewConversation) {
|
||||
if (turn.listPending) {
|
||||
await this.refreshConversations()
|
||||
turn.listPending = false
|
||||
}
|
||||
|
||||
// Start polling for intermediate messages in non-streaming mode too
|
||||
this.startPolling(currentConversationId)
|
||||
this.#followFromJob(currentConversationId, jobId)
|
||||
turn.startPolling()
|
||||
void this.#settleFromJob(turn, jobId)
|
||||
return true
|
||||
}
|
||||
}
|
||||
|
||||
@@ -128,50 +128,6 @@ describe('unread bookkeeping', () => {
|
||||
})
|
||||
})
|
||||
|
||||
describe('a chat re-pointed at another flow', () => {
|
||||
/**
|
||||
* The route component is reused between two flows, so the chat is re-pointed rather than
|
||||
* rebuilt. A selection carried across would open the new flow on the old flow's
|
||||
* conversation and send into its transcript and its agent's memory.
|
||||
*/
|
||||
it('forgets the flow it was pointed at when it is re-pointed', async () => {
|
||||
vi.mocked(FlowConversationsService.listConversationMessages).mockResolvedValue(
|
||||
rows('a', 2) as any
|
||||
)
|
||||
const manager = managerWithRows()
|
||||
await manager.selectConversation('a')
|
||||
expect(manager.selectedConversationId).toBe('a')
|
||||
expect(manager.liveRowIds.size).toBe(2)
|
||||
|
||||
manager.cleanup()
|
||||
|
||||
expect(manager.selectedConversationId).toBeUndefined()
|
||||
expect(manager.liveRowIds.size).toBe(0)
|
||||
})
|
||||
|
||||
/**
|
||||
* Forgetting has to hold against work already in flight: a transcript fetched for the
|
||||
* flow just left would otherwise be written into the one that replaced it.
|
||||
*/
|
||||
it('does not let a load started before the re-point write its rows back', async () => {
|
||||
let release = () => {}
|
||||
const held = new Promise<void>((resolve) => (release = resolve))
|
||||
vi.mocked(FlowConversationsService.listConversationMessages).mockImplementation((async () => {
|
||||
await held
|
||||
return rows('a', 2)
|
||||
}) as any)
|
||||
const manager = managerWithRows()
|
||||
const loading = manager.selectConversation('a')
|
||||
|
||||
manager.cleanup()
|
||||
release()
|
||||
await loading
|
||||
|
||||
expect(manager.liveRowIds.size).toBe(0)
|
||||
expect(manager.selectedConversationId).toBeUndefined()
|
||||
})
|
||||
})
|
||||
|
||||
/**
|
||||
* The messages endpoint answers oldest-first with a limit, so one request only reaches the
|
||||
* start of a turn that wrote a lot of rows — an agent calling several tools a round. Its
|
||||
@@ -188,6 +144,18 @@ describe('reading a turn longer than one page', () => {
|
||||
created_seq: from + i
|
||||
}))
|
||||
|
||||
// The sidebar reads the turn's outcome off the loaded rows, and the same long turn hides
|
||||
// the message it started from, leaving a failed chat looking like an idle one.
|
||||
it('reads a failed turn that filled the page as failed', async () => {
|
||||
const manager = managerWithRows()
|
||||
manager.selectedConversationId = 'a'
|
||||
manager.messages = assistantRows(1, 50).map((row, i) =>
|
||||
i === 49 ? ({ ...row, success: false } as any) : (row as any)
|
||||
)
|
||||
|
||||
expect(manager.conversationStatus('a')).toBe('error')
|
||||
})
|
||||
|
||||
it('keeps reading until a page comes back short', async () => {
|
||||
vi.mocked(FlowConversationsService.listConversationMessages)
|
||||
.mockReset()
|
||||
@@ -406,6 +374,30 @@ describe('an SSE timeout re-attaches instead of re-running', () => {
|
||||
* A chunk is not guaranteed to end on a line boundary. Split mid-JSON, the two halves
|
||||
* were parsed separately and both discarded, losing that token from the answer.
|
||||
*/
|
||||
// The thinking indicator reads these two off the manager while the turn runs. They are
|
||||
// written by the turn and read through the open conversation, so a turn that writes them
|
||||
// somewhere the getters do not look leaves the indicator permanently off.
|
||||
it('surfaces the thinking of the turn in flight', async () => {
|
||||
const { manager } = turnWith([
|
||||
[
|
||||
{
|
||||
type: 'update',
|
||||
new_result_stream: `${JSON.stringify({
|
||||
type: 'reasoning_token_delta',
|
||||
content: 'weighing it up'
|
||||
})}\n`
|
||||
},
|
||||
{ type: 'timeout' }
|
||||
]
|
||||
])
|
||||
|
||||
manager.inputMessage = 'ask something'
|
||||
await manager.sendMessage(undefined, undefined, 'a')
|
||||
|
||||
await vi.waitFor(() => expect(manager.isReasoningActive).toBe(true))
|
||||
await vi.waitFor(() => expect(manager.currentReasoning).toContain('weighing it up'))
|
||||
})
|
||||
|
||||
it('keeps a token whose chunk ended mid-line', async () => {
|
||||
const line = `${JSON.stringify({ type: 'token_delta', content: 'from-the-stream' })}\n`
|
||||
const cut = line.indexOf('from-the') + 3
|
||||
@@ -536,6 +528,51 @@ describe('a conversation opened while its run is still going', () => {
|
||||
await vi.waitFor(() => expect(streamCalls.some((call) => call.jobId === 'job-live')).toBe(true))
|
||||
})
|
||||
|
||||
/**
|
||||
* A turn that called a lot of tools writes more rows than the page the transcript opens
|
||||
* on, pushing the message that started it out of sight. Read from that page alone the
|
||||
* conversation looks finished, and the composer is handed back while the run still owns
|
||||
* the agent's memory.
|
||||
*/
|
||||
it('walks back past a turn that filled the page to find the message that started it', async () => {
|
||||
const answerRows = Array.from({ length: 50 }, (_, i) => ({
|
||||
id: `answer-${i}`,
|
||||
conversation_id: 'a',
|
||||
message_type: i % 2 === 0 ? 'assistant' : 'tool',
|
||||
content: `round ${i}`,
|
||||
created_at: new Date().toISOString(),
|
||||
created_seq: 100 + i
|
||||
}))
|
||||
vi.mocked(FlowConversationsService.listConversationMessages)
|
||||
.mockReset()
|
||||
// The newest page is this turn's answer, all of it.
|
||||
.mockResolvedValueOnce(answerRows as any)
|
||||
// The page before it ends with the message that started the turn.
|
||||
.mockResolvedValueOnce([
|
||||
{
|
||||
id: 'the-question',
|
||||
conversation_id: 'a',
|
||||
message_type: 'user',
|
||||
content: 'do the thing',
|
||||
created_at: new Date().toISOString(),
|
||||
created_seq: 99,
|
||||
job_id: 'job-live'
|
||||
}
|
||||
] as any)
|
||||
const manager = (live = managerWithRows())
|
||||
;(manager as any).initialize(vi.fn(), 'u/admin/flow', true)
|
||||
manager.operatingWorkspace = () => 'ws'
|
||||
|
||||
await manager.selectConversation('a')
|
||||
|
||||
await vi.waitFor(() => expect(manager.isConversationBusy('a')).toBe(true))
|
||||
await vi.waitFor(() => expect(streamCalls.some((call) => call.jobId === 'job-live')).toBe(true))
|
||||
// The rows the turn already wrote are dropped, since the stream replays them. The
|
||||
// message that started it is not: polls only ever add what a turn wrote, so a
|
||||
// transcript emptied here would come back as answers with nothing asking them.
|
||||
expect(manager.messages.map((m) => m.id)).toEqual(['the-question'])
|
||||
})
|
||||
|
||||
// The gap this closes: until the job answers, whether the chat is free is unknown, and a
|
||||
// message accepted meanwhile starts a second run against the same agent memory.
|
||||
it('holds the chat while it is asking whether a run is live', async () => {
|
||||
@@ -562,6 +599,36 @@ describe('a conversation opened while its run is still going', () => {
|
||||
* The hold a failed load leaves behind gates the whole surface — New chat and every
|
||||
* other conversation with it — so Stop has to be able to release it.
|
||||
*/
|
||||
/**
|
||||
* The walk back for the message that started the turn is a request of its own, and it can
|
||||
* fail on its own. Whether a run owns this conversation's memory is then precisely what is
|
||||
* unknown, so the chat stays held rather than taking a message that would write it twice.
|
||||
*/
|
||||
it('keeps holding the chat when it could not read far enough to ask', async () => {
|
||||
const answerRows = Array.from({ length: 50 }, (_, i) => ({
|
||||
id: `answer-${i}`,
|
||||
conversation_id: 'a',
|
||||
message_type: 'assistant',
|
||||
content: `round ${i}`,
|
||||
created_at: new Date().toISOString(),
|
||||
created_seq: 100 + i
|
||||
}))
|
||||
vi.mocked(FlowConversationsService.listConversationMessages)
|
||||
.mockReset()
|
||||
.mockResolvedValueOnce(answerRows as any)
|
||||
.mockRejectedValue(new Error('page unavailable'))
|
||||
const manager = (live = managerWithRows())
|
||||
;(manager as any).initialize(vi.fn(), 'u/admin/flow', true)
|
||||
manager.operatingWorkspace = () => 'ws'
|
||||
|
||||
await manager.selectConversation('a')
|
||||
|
||||
await vi.waitFor(() => expect(manager.isConversationBusy('a')).toBe(true))
|
||||
// And Stop is still the way out of it.
|
||||
await manager.cancelCurrentJob()
|
||||
expect(manager.isConversationBusy('a')).toBe(false)
|
||||
})
|
||||
|
||||
it('lets Stop release a chat whose rows never loaded', async () => {
|
||||
vi.mocked(FlowConversationsService.listConversationMessages).mockRejectedValue(
|
||||
new Error('transcript unavailable')
|
||||
|
||||
@@ -0,0 +1,116 @@
|
||||
<!--
|
||||
One chat, for one flow in one workspace. Mounted only by `FlowChat`, which keys it on
|
||||
that pair: the manager created here holds the conversations, the turn running in one of
|
||||
them and the rows that turn is writing, so pointing an existing panel at another flow
|
||||
would leave an upload, a launch or a poll landing in the chat that replaced it.
|
||||
-->
|
||||
<script lang="ts">
|
||||
import { workspaceStore } from '$lib/stores'
|
||||
import { createFlowChatManager } from './FlowChatManager.svelte'
|
||||
import FlowConversationsSidebar from './FlowConversationsSidebar.svelte'
|
||||
import FlowChatInterface from './FlowChatInterface.svelte'
|
||||
import { getContext, untrack } from 'svelte'
|
||||
import type { FlowEditorContext } from '../types'
|
||||
import type { ChatFrame, FlowChatProps } from './flowChatProps'
|
||||
|
||||
const FRAME_CLASS: Record<ChatFrame, string> = {
|
||||
boxed: 'border rounded-md',
|
||||
top: 'border-t',
|
||||
none: ''
|
||||
}
|
||||
|
||||
let {
|
||||
onRunFlow,
|
||||
deploymentInProgress = false,
|
||||
useStreaming = false,
|
||||
path,
|
||||
description = undefined,
|
||||
inputSchema = undefined,
|
||||
flowModules = undefined,
|
||||
wideLayout = false,
|
||||
frame = 'top',
|
||||
conversationKind = 'deployed',
|
||||
parallelTurns = false
|
||||
}: FlowChatProps = $props()
|
||||
|
||||
const flowEditorContext = getContext<FlowEditorContext>('FlowEditorContext')
|
||||
|
||||
// The flow and the workspace this panel was built for. `FlowChat` keys on the same pair,
|
||||
// so neither can change under it: what they change is which panel exists.
|
||||
const workspace = $derived(flowEditorContext?.opWorkspace?.() ?? $workspaceStore)
|
||||
|
||||
const manager = createFlowChatManager()
|
||||
manager.operatingWorkspace = () => flowEditorContext?.opWorkspace?.()
|
||||
manager.conversationKind = conversationKind
|
||||
manager.allowsParallelTurns = parallelTurns
|
||||
// The filter moves; what this surface runs does not.
|
||||
manager.surfaceKind = conversationKind === 'test' ? 'test' : 'deployed'
|
||||
// The editor is the only surface with both kinds in play, and it is the one that opens
|
||||
// on test chats. A deployed flow lists what its users started, with no way to ask for
|
||||
// anything else.
|
||||
manager.canFilterConversationKind = conversationKind !== 'deployed'
|
||||
|
||||
// The manager's inputs, kept current rather than set once: `useStreaming` follows the
|
||||
// flow's last step, which an edit can change while a turn is running.
|
||||
$effect(() => {
|
||||
manager.initialize(onRunFlow, path, useStreaming)
|
||||
})
|
||||
|
||||
// Opens the chat on the conversation last spoken in. Untracked: it reads the open
|
||||
// conversation, and tracked, the first send of a fresh chat would select the conversation
|
||||
// it just created and so abort its own turn. Waits for a workspace, which an embedded
|
||||
// editor resolves from its session rather than having on the first run.
|
||||
$effect(() => {
|
||||
if (workspace) untrack(() => manager.selectLatestConversation())
|
||||
})
|
||||
|
||||
// Ends whatever is still running when the panel goes away, which is the only thing that
|
||||
// should end it. Nothing is tracked here, so no input changing under the panel can.
|
||||
$effect(() => () => manager.cleanup())
|
||||
|
||||
// Initialize InfiniteList when component mounts
|
||||
$effect(() => {
|
||||
if (workspace && path && manager.conversationListComponent) {
|
||||
untrack(() => {
|
||||
manager.setupInfiniteList()
|
||||
})
|
||||
}
|
||||
})
|
||||
|
||||
// Everything the chat asks for beyond the message itself. `user_message` is the
|
||||
// composer: the server requires that exact argument on a chat-enabled flow and
|
||||
// stores it as the conversation's message (`handle_chat_conversation_messages`), so
|
||||
// the name is a contract rather than the author's choice.
|
||||
const additionalInputsSchema = $derived.by(() => {
|
||||
const props = inputSchema?.properties ?? {}
|
||||
const messageInput = 'user_message'
|
||||
const filtered = Object.fromEntries(Object.entries(props).filter(([k]) => k !== messageInput))
|
||||
if (Object.keys(filtered).length === 0) return undefined
|
||||
const required = inputSchema?.required
|
||||
const requiredArray: string[] = Array.isArray(required) ? required : []
|
||||
return {
|
||||
...inputSchema,
|
||||
properties: filtered,
|
||||
required: requiredArray.filter((k: string) => k !== messageInput)
|
||||
}
|
||||
})
|
||||
</script>
|
||||
|
||||
<!-- The column's max width and side padding come from AIChatDisplay itself. -->
|
||||
<div class="flex overflow-hidden flex-1 {FRAME_CLASS[frame]}">
|
||||
<FlowConversationsSidebar {manager} {description} />
|
||||
<!-- pb-3 on the chat alone, not on the row: the transcript and composer stop short of
|
||||
the panel edge the way the session chat does, while the sidebar and the border
|
||||
dividing it from the chat still reach the bottom. -->
|
||||
<div class="flex flex-1 min-w-0 min-h-0 pb-3">
|
||||
<FlowChatInterface
|
||||
{manager}
|
||||
{deploymentInProgress}
|
||||
{additionalInputsSchema}
|
||||
{flowModules}
|
||||
{path}
|
||||
{description}
|
||||
{wideLayout}
|
||||
/>
|
||||
</div>
|
||||
</div>
|
||||
@@ -0,0 +1,36 @@
|
||||
import type { FlowModule } from '$lib/gen'
|
||||
import type { ConversationKind } from './FlowChatManager.svelte'
|
||||
|
||||
/** What separates the chat from what sits above it. `boxed` is its own panel, corners
|
||||
* clipped so the sidebar's edge follows them; `top` a dividing line under an enclosing
|
||||
* header; `none` for a surface where the chat is the whole pane. */
|
||||
export type ChatFrame = 'boxed' | 'top' | 'none'
|
||||
|
||||
/**
|
||||
* Shared by `FlowChat` and the `FlowChatPanel` it keys, which passes them straight through:
|
||||
* one declaration, so a prop cannot reach the panel without the wrapper accepting it.
|
||||
*/
|
||||
export interface FlowChatProps {
|
||||
onRunFlow: (
|
||||
userMessage: string,
|
||||
conversationId: string,
|
||||
additionalInputs?: Record<string, any>
|
||||
) => Promise<string | undefined>
|
||||
useStreaming?: boolean
|
||||
deploymentInProgress?: boolean
|
||||
path: string
|
||||
/** The flow's own description, shown where the chat has room for it: the empty
|
||||
* transcript, and the sidebar once a conversation has replaced it. */
|
||||
description?: string
|
||||
inputSchema?: Record<string, any>
|
||||
/** The flow's modules, used to find which inputs an AI agent step reads directly. */
|
||||
flowModules?: FlowModule[]
|
||||
/** Wider centered column, for the full-page chat. */
|
||||
wideLayout?: boolean
|
||||
frame?: ChatFrame
|
||||
/** Which chats the sidebar lists before the reader filters it themselves. */
|
||||
conversationKind?: ConversationKind
|
||||
/** Whether a turn may run while another chat's is still going. Off where one run
|
||||
* owns the surface — the editor's panel shows it on the graph. */
|
||||
parallelTurns?: boolean
|
||||
}
|
||||
@@ -0,0 +1,153 @@
|
||||
import { emptyTurnState, type TurnState } from './turnTranscript'
|
||||
import {
|
||||
prefersInstantReveal,
|
||||
TypewriterReveal
|
||||
} from '$lib/components/copilot/chat/typewriterReveal'
|
||||
|
||||
/** How often a running turn asks the server for the rows it has written so far. */
|
||||
const POLL_MS = 500
|
||||
/** How long a turn is polled before the reader is assumed to have left it running. */
|
||||
const POLL_MAX_MS = 2 * 60 * 1000
|
||||
|
||||
export type RevealKind = 'answer' | 'reasoning'
|
||||
|
||||
export interface TurnOptions {
|
||||
conversationId: string
|
||||
/** Paced text leaving the turn, a chunk at a time. Called for the turn's whole life. */
|
||||
onReveal: (kind: RevealKind, chunk: string) => void
|
||||
/** One tick of the catch-up read while the turn runs. */
|
||||
onPoll: () => void
|
||||
/** This turn created the conversation, so the sidebar's list has a row it has not seen. */
|
||||
createdConversation?: boolean
|
||||
}
|
||||
|
||||
/**
|
||||
* One turn of one conversation: the run it started, the handle that stops it, and the
|
||||
* display state it is writing.
|
||||
*
|
||||
* It exists to be identified. A turn's work is a stream, a poll and a settle, each of which
|
||||
* outlives the moment it was started in, and each of which writes back to the chat. Holding
|
||||
* that work's turn means the chat can ask "is this still the turn I am on?" — one question,
|
||||
* asked with `===` — instead of every write-back re-deriving the answer from an abort flag
|
||||
* it happens to have in scope. A turn that has been replaced or ended answers no, and its
|
||||
* work lands nowhere.
|
||||
*
|
||||
* Everything a turn has to undo is held here, so ending one is a single call rather than
|
||||
* an abort, two timers and a pair of reveals that each caller has to remember.
|
||||
*/
|
||||
export class Turn {
|
||||
readonly conversationId: string
|
||||
readonly #controller = new AbortController()
|
||||
|
||||
/** Stops the work this turn started. Read by everything it hands to the network. */
|
||||
get signal(): AbortSignal {
|
||||
return this.#controller.signal
|
||||
}
|
||||
|
||||
/** The flow job, once the run has one. A turn holding the chat before it has a job — an
|
||||
* attachment upload, a reload working out whether a run is live — has none, and Stop
|
||||
* has nothing to cancel. */
|
||||
jobId = $state<string | undefined>(undefined)
|
||||
|
||||
/** The model's thinking, as revealed. Read by the chat while the turn is the open one. */
|
||||
reasoning = $state('')
|
||||
/** From the first thinking token until a tool call closes the row. */
|
||||
reasoningActive = $state(false)
|
||||
|
||||
/**
|
||||
* Where the live answer resumes after a dropped connection. Kept for the turn rather
|
||||
* than for one attempt, so a reconnect picks up after what is already on screen.
|
||||
*/
|
||||
streamOffset: number | undefined = undefined
|
||||
|
||||
/** Which row is open and what is in it. The rows themselves belong to the transcript. */
|
||||
transcript: TurnState
|
||||
|
||||
/**
|
||||
* The conversation this turn created is not in the sidebar's list yet. Cleared by the
|
||||
* read that picks it up, so a turn long enough to be read many times asks for the list
|
||||
* once rather than on every tick.
|
||||
*/
|
||||
listPending = $state(false)
|
||||
|
||||
#replyReveal: TypewriterReveal
|
||||
#reasoningReveal: TypewriterReveal
|
||||
#pollInterval: ReturnType<typeof setInterval> | undefined
|
||||
#pollDeadline: ReturnType<typeof setTimeout> | undefined
|
||||
#onPoll: () => void
|
||||
#ended = false
|
||||
|
||||
constructor(options: TurnOptions) {
|
||||
this.conversationId = options.conversationId
|
||||
this.transcript = emptyTurnState(options.conversationId)
|
||||
this.listPending = options.createdConversation ?? false
|
||||
this.#onPoll = options.onPoll
|
||||
const instant = prefersInstantReveal()
|
||||
this.#replyReveal = new TypewriterReveal({
|
||||
onReveal: (chunk) => options.onReveal('answer', chunk),
|
||||
instant
|
||||
})
|
||||
this.#reasoningReveal = new TypewriterReveal({
|
||||
onReveal: (chunk) => options.onReveal('reasoning', chunk),
|
||||
instant
|
||||
})
|
||||
}
|
||||
|
||||
/** Stopped, or replaced by a later turn that ended this one. */
|
||||
get ended(): boolean {
|
||||
return this.#ended || this.#controller.signal.aborted
|
||||
}
|
||||
|
||||
pushAnswer(chunk: string) {
|
||||
if (this.ended) return
|
||||
this.#replyReveal.push(chunk)
|
||||
}
|
||||
|
||||
pushReasoning(chunk: string) {
|
||||
if (this.ended) return
|
||||
this.#reasoningReveal.push(chunk)
|
||||
}
|
||||
|
||||
/** Put what the pacing still holds on screen now. What is buffered belongs to the row
|
||||
* being closed — a tool card, or the end of the turn — and would otherwise be dropped. */
|
||||
flushReveals() {
|
||||
this.#replyReveal.flush()
|
||||
this.#reasoningReveal.flush()
|
||||
}
|
||||
|
||||
/** The catch-up read, for the rows a turn writes that the stream does not carry. Safe to
|
||||
* call more than once; the deadline bounds a turn whose reader has gone. */
|
||||
startPolling() {
|
||||
if (this.#pollInterval || this.ended) return
|
||||
this.#pollInterval = setInterval(() => this.#onPoll(), POLL_MS)
|
||||
this.#pollDeadline = setTimeout(() => this.stopPolling(), POLL_MAX_MS)
|
||||
}
|
||||
|
||||
/** Both timers, always together: a deadline left armed by a turn that stopped early
|
||||
* would fire into whatever is polling by then. */
|
||||
stopPolling() {
|
||||
if (this.#pollInterval) {
|
||||
clearInterval(this.#pollInterval)
|
||||
this.#pollInterval = undefined
|
||||
}
|
||||
if (this.#pollDeadline) {
|
||||
clearTimeout(this.#pollDeadline)
|
||||
this.#pollDeadline = undefined
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Stop everything this turn started. The run itself is not cancelled — leaving a chat
|
||||
* has never stopped a flow, and a turn cannot tell being stopped from being left.
|
||||
*/
|
||||
end() {
|
||||
this.#ended = true
|
||||
this.stopPolling()
|
||||
this.#replyReveal.reset()
|
||||
this.#reasoningReveal.reset()
|
||||
this.transcript = emptyTurnState(this.conversationId)
|
||||
this.reasoning = ''
|
||||
this.reasoningActive = false
|
||||
this.#controller.abort()
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,61 @@
|
||||
import { describe, it, expect, vi, afterEach } from 'vitest'
|
||||
import { Turn } from './turn.svelte'
|
||||
|
||||
function turnWith(onPoll = vi.fn()) {
|
||||
return new Turn({ conversationId: 'a', onReveal: vi.fn(), onPoll })
|
||||
}
|
||||
|
||||
afterEach(() => {
|
||||
vi.useRealTimers()
|
||||
})
|
||||
|
||||
/**
|
||||
* A turn is the thing the chat checks its own work against, so the only property that has to
|
||||
* hold is that an ended one says so and leaves nothing of itself running.
|
||||
*/
|
||||
describe('a turn that has ended', () => {
|
||||
it('answers that it has, and aborts what it started', () => {
|
||||
const turn = turnWith()
|
||||
expect(turn.ended).toBe(false)
|
||||
expect(turn.signal.aborted).toBe(false)
|
||||
|
||||
turn.end()
|
||||
|
||||
expect(turn.ended).toBe(true)
|
||||
expect(turn.signal.aborted).toBe(true)
|
||||
})
|
||||
|
||||
it('stops polling, and cannot be started again', () => {
|
||||
vi.useFakeTimers()
|
||||
const onPoll = vi.fn()
|
||||
const turn = turnWith(onPoll)
|
||||
turn.startPolling()
|
||||
vi.advanceTimersByTime(1500)
|
||||
expect(onPoll).toHaveBeenCalled()
|
||||
const ticks = onPoll.mock.calls.length
|
||||
|
||||
turn.end()
|
||||
// Both timers go together: a deadline left armed by a turn that stopped early would
|
||||
// fire into whatever is polling by the time it came round.
|
||||
turn.startPolling()
|
||||
vi.advanceTimersByTime(10 * 60 * 1000)
|
||||
|
||||
expect(onPoll.mock.calls.length).toBe(ticks)
|
||||
})
|
||||
|
||||
it('gives the reader back the chat after two minutes of a run nobody is watching', () => {
|
||||
vi.useFakeTimers()
|
||||
const onPoll = vi.fn()
|
||||
const turn = turnWith(onPoll)
|
||||
turn.startPolling()
|
||||
|
||||
vi.advanceTimersByTime(2 * 60 * 1000)
|
||||
const ticks = onPoll.mock.calls.length
|
||||
vi.advanceTimersByTime(60 * 1000)
|
||||
|
||||
expect(onPoll.mock.calls.length).toBe(ticks)
|
||||
// The deadline stops the reading, not the turn: the run is still the one that says
|
||||
// when it is over.
|
||||
expect(turn.ended).toBe(false)
|
||||
})
|
||||
})
|
||||
Reference in New Issue
Block a user