mirror of
https://github.com/langgenius/dify.git
synced 2026-07-30 08:49:31 +08:00
1651 lines
61 KiB
TypeScript
1651 lines
61 KiB
TypeScript
import type { ChatConfig, ChatItem, ChatItemInTree, Inputs } from '../types'
|
|
import type { InputForm, ThoughtItem } from './type'
|
|
import type AudioPlayer from '@/app/components/base/audio-btn/audio'
|
|
import type { FileEntity } from '@/app/components/base/file-uploader/types'
|
|
import type { Annotation } from '@/models/log'
|
|
import type { IOnDataMoreInfo, IOtherOptions } from '@/service/base'
|
|
import type { VisionFile } from '@/types/app'
|
|
import type { FileResponse, ReasoningChunkResponse } from '@/types/workflow'
|
|
import { toast } from '@langgenius/dify-ui/toast'
|
|
import { uniqBy } from 'es-toolkit/compat'
|
|
import { noop } from 'es-toolkit/function'
|
|
import { produce, setAutoFreeze } from 'immer'
|
|
import { useCallback, useEffect, useMemo, useRef, useState } from 'react'
|
|
import { useTranslation } from 'react-i18next'
|
|
import { v4 as uuidV4 } from 'uuid'
|
|
import { AudioPlayerManager } from '@/app/components/base/audio-btn/audio.player.manager'
|
|
import { enrichSubmittedHumanInputFormData } from '@/app/components/base/chat/chat/answer/human-input-content/submitted-utils'
|
|
import {
|
|
getProcessedFiles,
|
|
getProcessedFilesFromResponse,
|
|
} from '@/app/components/base/file-uploader/utils'
|
|
import { isInstalledAppPath } from '@/app/components/explore/installed-app/routes'
|
|
import { addFileInfos, sortAgentSorts } from '@/app/components/tools/utils'
|
|
import { NodeRunningStatus, WorkflowRunningStatus } from '@/app/components/workflow/types'
|
|
import useTimestamp from '@/hooks/use-timestamp'
|
|
import { useParams, usePathname } from '@/next/navigation'
|
|
import { sseGet, ssePost } from '@/service/base'
|
|
import { TransferMethod } from '@/types/app'
|
|
import { getThreadMessages } from '../utils'
|
|
import { getProcessedInputs, processOpeningStatement } from './utils'
|
|
|
|
type GetAbortController = (abortController: AbortController) => void
|
|
type HistoryMessageFile = Partial<FileResponse> & {
|
|
id?: string
|
|
belongs_to?: string
|
|
}
|
|
type HistoryConversationMessage = {
|
|
id: string
|
|
answer?: string
|
|
message?: { role: string; text: string; files?: FileEntity[] }[]
|
|
message_files?: HistoryMessageFile[]
|
|
agent_thoughts?: ThoughtItem[] | null
|
|
retriever_resources?: ChatItem['citation']
|
|
metadata?: {
|
|
reasoning?: ChatItem['reasoningContent']
|
|
}
|
|
created_at?: number
|
|
answer_tokens?: number
|
|
message_tokens?: number
|
|
provider_response_latency?: number
|
|
workflow_run_id?: string
|
|
feedback?: ChatItem['feedback']
|
|
inputs?: unknown
|
|
query?: string
|
|
}
|
|
type ConversationMessagesResponse = {
|
|
data: HistoryConversationMessage[]
|
|
}
|
|
type SendCallback = {
|
|
onGetConversationMessages?: (
|
|
conversationId: string,
|
|
getAbortController: GetAbortController,
|
|
) => Promise<unknown>
|
|
onGetSuggestedQuestions?: (
|
|
responseItemId: string,
|
|
getAbortController: GetAbortController,
|
|
) => Promise<unknown>
|
|
onConversationComplete?: (conversationId: string, workflowRunId?: string) => void
|
|
onUnhandledEvent?: IOtherOptions['onUnhandledEvent']
|
|
onSendSettled?: (hasError?: boolean) => void
|
|
isPublicAPI?: boolean
|
|
}
|
|
|
|
type UseChatOptions = {
|
|
isNewAgent?: boolean
|
|
timezone?: string
|
|
}
|
|
|
|
function mergeStreamingThought(currentThought: ThoughtItem, nextThought: ThoughtItem): ThoughtItem {
|
|
return {
|
|
...nextThought,
|
|
message_files: nextThought.message_files?.length
|
|
? nextThought.message_files
|
|
: currentThought.message_files,
|
|
}
|
|
}
|
|
|
|
function appendAgentResponseMessagePart(responseItem: ChatItemInTree, message: string) {
|
|
if (!responseItem.agent_response_parts) responseItem.agent_response_parts = []
|
|
|
|
const lastPart = responseItem.agent_response_parts.at(-1)
|
|
if (lastPart?.type === 'message') {
|
|
lastPart.content += message
|
|
} else {
|
|
responseItem.agent_response_parts.push({
|
|
type: 'message',
|
|
content: message,
|
|
})
|
|
}
|
|
}
|
|
|
|
function upsertAgentResponseThoughtPart(responseItem: ChatItemInTree, thought: ThoughtItem) {
|
|
if (!responseItem.agent_response_parts) responseItem.agent_response_parts = []
|
|
|
|
const partIndex = responseItem.agent_response_parts.findIndex(
|
|
(part) => part.type === 'thought' && part.thought.id === thought.id,
|
|
)
|
|
if (partIndex > -1) {
|
|
responseItem.agent_response_parts[partIndex] = {
|
|
type: 'thought',
|
|
thought,
|
|
}
|
|
return
|
|
}
|
|
|
|
responseItem.agent_response_parts.push({
|
|
type: 'thought',
|
|
thought,
|
|
})
|
|
}
|
|
|
|
function getHistoryAgentThoughts(responseItem: HistoryConversationMessage) {
|
|
if (!Array.isArray(responseItem.agent_thoughts)) return []
|
|
|
|
const messageFiles: VisionFile[] =
|
|
responseItem.message_files?.map((file) => ({
|
|
id: file.id,
|
|
type: file.type || '',
|
|
transfer_method: file.transfer_method || TransferMethod.remote_url,
|
|
url: file.url || '',
|
|
upload_file_id: file.upload_file_id || '',
|
|
belongs_to: file.belongs_to,
|
|
})) || []
|
|
|
|
return addFileInfos(sortAgentSorts(responseItem.agent_thoughts), messageFiles)
|
|
}
|
|
|
|
function toHistoryFileResponse(file: HistoryMessageFile): FileResponse {
|
|
return {
|
|
related_id: file.related_id || file.id || '',
|
|
extension: file.extension || '',
|
|
filename: file.filename || '',
|
|
size: file.size || 0,
|
|
mime_type: file.mime_type || file.type || '',
|
|
transfer_method: file.transfer_method || TransferMethod.remote_url,
|
|
type: file.type || '',
|
|
url: file.url || '',
|
|
upload_file_id: file.upload_file_id || '',
|
|
remote_url: file.remote_url || '',
|
|
}
|
|
}
|
|
|
|
function getHistoryAnswerFiles(responseItem: HistoryConversationMessage) {
|
|
const answerFiles =
|
|
responseItem.message_files?.filter((file) => file.belongs_to === 'assistant') || []
|
|
|
|
return getProcessedFilesFromResponse(
|
|
answerFiles.map((file) =>
|
|
toHistoryFileResponse({
|
|
...file,
|
|
related_id: file.related_id || file.id || '',
|
|
upload_file_id: file.upload_file_id || '',
|
|
}),
|
|
),
|
|
)
|
|
}
|
|
|
|
function isHistoryConversationMessage(value: unknown): value is HistoryConversationMessage {
|
|
return (
|
|
typeof value === 'object' &&
|
|
value !== null &&
|
|
typeof (value as { id?: unknown }).id === 'string'
|
|
)
|
|
}
|
|
|
|
function getConversationMessagesData(response: unknown): ConversationMessagesResponse['data'] {
|
|
if (typeof response !== 'object' || response === null) return []
|
|
|
|
const data = (response as { data?: unknown }).data
|
|
return Array.isArray(data) ? data.filter(isHistoryConversationMessage) : []
|
|
}
|
|
|
|
export const useChat = (
|
|
config?: ChatConfig,
|
|
formSettings?: {
|
|
inputs: Inputs
|
|
inputsForm: InputForm[]
|
|
},
|
|
prevChatTree?: ChatItemInTree[],
|
|
stopChat?: (taskId: string) => void,
|
|
clearChatList?: boolean,
|
|
clearChatListCallback?: (state: boolean) => void,
|
|
initialConversationId?: string,
|
|
options: UseChatOptions = {},
|
|
) => {
|
|
const { t } = useTranslation()
|
|
const { formatTime } = useTimestamp({ timezone: options.timezone })
|
|
const conversationIdRef = useRef(initialConversationId ?? '')
|
|
const initialConversationIdRef = useRef(initialConversationId ?? '')
|
|
const hasStopRespondedRef = useRef(false)
|
|
const [isResponding, setIsResponding] = useState(false)
|
|
const isRespondingRef = useRef(false)
|
|
const taskIdRef = useRef('')
|
|
const pausedStateRef = useRef(false)
|
|
const [suggestedQuestions, setSuggestedQuestions] = useState<string[]>([])
|
|
const conversationMessagesAbortControllerRef = useRef<AbortController | null>(null)
|
|
const suggestedQuestionsAbortControllerRef = useRef<AbortController | null>(null)
|
|
const workflowEventsAbortControllerRef = useRef<AbortController | null>(null)
|
|
const params = useParams()
|
|
const pathname = usePathname()
|
|
|
|
const [chatTree, setChatTree] = useState<ChatItemInTree[]>(prevChatTree || [])
|
|
const chatTreeRef = useRef<ChatItemInTree[]>(chatTree)
|
|
const [targetMessageId, setTargetMessageId] = useState<string>()
|
|
const threadMessages = useMemo(
|
|
() => getThreadMessages(chatTree, targetMessageId),
|
|
[chatTree, targetMessageId],
|
|
)
|
|
|
|
const getIntroduction = useCallback(
|
|
(str: string) => {
|
|
return processOpeningStatement(
|
|
str,
|
|
formSettings?.inputs || {},
|
|
formSettings?.inputsForm || [],
|
|
)
|
|
},
|
|
[formSettings?.inputs, formSettings?.inputsForm],
|
|
)
|
|
|
|
const processedOpeningContent = config?.opening_statement
|
|
? getIntroduction(config.opening_statement)
|
|
: undefined
|
|
const processedSuggestionsKey = config?.suggested_questions
|
|
? JSON.stringify(config.suggested_questions.map((q) => getIntroduction(q)))
|
|
: undefined
|
|
|
|
const openingStatementItem = useMemo<ChatItemInTree | null>(() => {
|
|
if (!processedOpeningContent) return null
|
|
return {
|
|
id: 'opening-statement',
|
|
content: processedOpeningContent,
|
|
isAnswer: true,
|
|
isOpeningStatement: true,
|
|
suggestedQuestions: processedSuggestionsKey
|
|
? (JSON.parse(processedSuggestionsKey) as string[])
|
|
: undefined,
|
|
}
|
|
}, [processedOpeningContent, processedSuggestionsKey])
|
|
|
|
const threadOpener = useMemo(
|
|
() => threadMessages.find((item) => item.isOpeningStatement) ?? null,
|
|
[threadMessages],
|
|
)
|
|
|
|
const mergedOpeningItem = useMemo<ChatItemInTree | null>(() => {
|
|
if (!threadOpener || !openingStatementItem) return null
|
|
return {
|
|
...threadOpener,
|
|
content: openingStatementItem.content,
|
|
suggestedQuestions: openingStatementItem.suggestedQuestions,
|
|
}
|
|
}, [threadOpener, openingStatementItem])
|
|
|
|
/** Final chat list that will be rendered */
|
|
const chatList = useMemo(() => {
|
|
const ret = [...threadMessages]
|
|
if (openingStatementItem) {
|
|
const index = threadMessages.findIndex((item) => item.isOpeningStatement)
|
|
if (index > -1 && mergedOpeningItem) ret[index] = mergedOpeningItem
|
|
else if (index === -1) ret.unshift(openingStatementItem)
|
|
}
|
|
return ret
|
|
}, [threadMessages, openingStatementItem, mergedOpeningItem])
|
|
|
|
useEffect(() => {
|
|
setAutoFreeze(false)
|
|
return () => {
|
|
setAutoFreeze(true)
|
|
}
|
|
}, [])
|
|
|
|
useEffect(() => {
|
|
const nextConversationId = initialConversationId ?? ''
|
|
initialConversationIdRef.current = nextConversationId
|
|
conversationIdRef.current = nextConversationId
|
|
}, [initialConversationId])
|
|
|
|
/** Find the target node by bfs and then operate on it */
|
|
const produceChatTreeNode = useCallback(
|
|
(targetId: string, operation: (node: ChatItemInTree) => void) => {
|
|
return produce(chatTreeRef.current, (draft) => {
|
|
const queue: ChatItemInTree[] = [...draft]
|
|
while (queue.length > 0) {
|
|
const current = queue.shift()!
|
|
if (current.id === targetId) {
|
|
operation(current)
|
|
break
|
|
}
|
|
if (current.children) queue.push(...current.children)
|
|
}
|
|
})
|
|
},
|
|
[],
|
|
)
|
|
|
|
type UpdateChatTreeNode = {
|
|
(id: string, fields: Partial<ChatItemInTree>): void
|
|
(id: string, update: (node: ChatItemInTree) => void): void
|
|
}
|
|
|
|
const updateChatTreeNode: UpdateChatTreeNode = useCallback(
|
|
(id: string, fieldsOrUpdate: Partial<ChatItemInTree> | ((node: ChatItemInTree) => void)) => {
|
|
const nextState = produceChatTreeNode(id, (node) => {
|
|
if (typeof fieldsOrUpdate === 'function') {
|
|
fieldsOrUpdate(node)
|
|
} else {
|
|
Object.keys(fieldsOrUpdate).forEach((key) => {
|
|
;(node as any)[key] = (fieldsOrUpdate as any)[key]
|
|
})
|
|
}
|
|
})
|
|
setChatTree(nextState)
|
|
chatTreeRef.current = nextState
|
|
},
|
|
[produceChatTreeNode],
|
|
)
|
|
|
|
const handleResponding = useCallback((isResponding: boolean) => {
|
|
setIsResponding(isResponding)
|
|
isRespondingRef.current = isResponding
|
|
}, [])
|
|
|
|
const handleStop = useCallback(() => {
|
|
hasStopRespondedRef.current = true
|
|
handleResponding(false)
|
|
if (stopChat && taskIdRef.current && !pausedStateRef.current) stopChat(taskIdRef.current)
|
|
if (conversationMessagesAbortControllerRef.current)
|
|
conversationMessagesAbortControllerRef.current.abort()
|
|
if (suggestedQuestionsAbortControllerRef.current)
|
|
suggestedQuestionsAbortControllerRef.current.abort()
|
|
if (workflowEventsAbortControllerRef.current) workflowEventsAbortControllerRef.current.abort()
|
|
}, [stopChat, handleResponding])
|
|
|
|
const handleRestart = useCallback(
|
|
(cb?: any) => {
|
|
conversationIdRef.current = initialConversationIdRef.current
|
|
taskIdRef.current = ''
|
|
handleStop()
|
|
setChatTree([])
|
|
setSuggestedQuestions([])
|
|
cb?.()
|
|
},
|
|
[handleStop],
|
|
)
|
|
|
|
const createAudioPlayerManager = useCallback(() => {
|
|
let ttsUrl = ''
|
|
let ttsIsPublic = false
|
|
if (params.token) {
|
|
ttsUrl = '/text-to-audio'
|
|
ttsIsPublic = true
|
|
} else if (params.appId) {
|
|
if (isInstalledAppPath(pathname)) ttsUrl = `/installed-apps/${params.appId}/text-to-audio`
|
|
else ttsUrl = `/apps/${params.appId}/text-to-audio`
|
|
}
|
|
|
|
let player: AudioPlayer | null = null
|
|
const getOrCreatePlayer = () => {
|
|
if (!player)
|
|
player = AudioPlayerManager.getInstance().getAudioPlayer(
|
|
ttsUrl,
|
|
ttsIsPublic,
|
|
uuidV4(),
|
|
'none',
|
|
'none',
|
|
noop,
|
|
)
|
|
|
|
return player
|
|
}
|
|
|
|
return getOrCreatePlayer
|
|
}, [params.token, params.appId, pathname])
|
|
|
|
const handleResume = useCallback(
|
|
async (
|
|
messageId: string,
|
|
workflowRunId: string,
|
|
{ onGetSuggestedQuestions, onConversationComplete, onSendSettled, isPublicAPI }: SendCallback,
|
|
) => {
|
|
const getOrCreatePlayer = createAudioPlayerManager()
|
|
let hasSettled = false
|
|
const settleSend = (hasError?: boolean) => {
|
|
if (hasSettled) return
|
|
|
|
hasSettled = true
|
|
onSendSettled?.(hasError)
|
|
}
|
|
// Re-subscribe to workflow events for the specific message
|
|
const url = `/workflow/${workflowRunId}/events?include_state_snapshot=true`
|
|
|
|
const otherOptions: IOtherOptions = {
|
|
isPublicAPI,
|
|
getAbortController: (abortController) => {
|
|
workflowEventsAbortControllerRef.current = abortController
|
|
},
|
|
onData: (
|
|
message: string,
|
|
isFirstMessage: boolean,
|
|
{ event, conversationId: newConversationId, messageId, taskId }: IOnDataMoreInfo,
|
|
) => {
|
|
updateChatTreeNode(messageId, (responseItem) => {
|
|
const agentThoughts = responseItem.agent_thoughts ?? []
|
|
const isNewAgentMessage =
|
|
options.isNewAgent && (event === 'agent_message' || event === 'message')
|
|
if (isNewAgentMessage) {
|
|
appendAgentResponseMessagePart(responseItem, message)
|
|
} else if (!agentThoughts.length || options.isNewAgent) {
|
|
responseItem.content = responseItem.content + message
|
|
} else {
|
|
const lastThought = agentThoughts[agentThoughts.length - 1]
|
|
if (lastThought) lastThought.thought = lastThought.thought + message
|
|
}
|
|
if (messageId) responseItem.id = messageId
|
|
})
|
|
|
|
if (isFirstMessage && newConversationId) conversationIdRef.current = newConversationId
|
|
|
|
if (taskId) taskIdRef.current = taskId
|
|
},
|
|
onReasoning: ({ data: reasoningData }: ReasoningChunkResponse) => {
|
|
const { message_id, reasoning, node_id, is_final } = reasoningData
|
|
updateChatTreeNode(message_id, (responseItem) => {
|
|
const reasoningContent =
|
|
responseItem.reasoningContent || (responseItem.reasoningContent = {})
|
|
const key = node_id || '_'
|
|
if (reasoning) reasoningContent[key] = (reasoningContent[key] || '') + reasoning
|
|
if (is_final) responseItem.reasoningFinished = true
|
|
})
|
|
},
|
|
async onCompleted(hasError?: boolean) {
|
|
handleResponding(false)
|
|
|
|
try {
|
|
if (hasError) return
|
|
|
|
if (onConversationComplete)
|
|
onConversationComplete(conversationIdRef.current, workflowRunId)
|
|
|
|
if (
|
|
config?.suggested_questions_after_answer?.enabled &&
|
|
!hasStopRespondedRef.current &&
|
|
onGetSuggestedQuestions
|
|
) {
|
|
try {
|
|
const { data }: any = await onGetSuggestedQuestions(
|
|
messageId,
|
|
(newAbortController) =>
|
|
(suggestedQuestionsAbortControllerRef.current = newAbortController),
|
|
)
|
|
setSuggestedQuestions(data)
|
|
} catch {
|
|
setSuggestedQuestions([])
|
|
}
|
|
}
|
|
} finally {
|
|
settleSend(hasError)
|
|
}
|
|
},
|
|
onFile(file) {
|
|
// Convert simple file type to MIME type for non-agent mode
|
|
// Backend sends: { id, type: "image", belongs_to, url }
|
|
// Frontend expects: { id, type: "image/png", transferMethod, url, uploadedId, supportFileType, name, size }
|
|
|
|
// Determine file type for MIME conversion
|
|
const fileType = (file as { type?: string }).type || 'image'
|
|
|
|
// If file already has transferMethod, use it as base and ensure all required fields exist
|
|
// Otherwise, create a new complete file object
|
|
const baseFile = 'transferMethod' in file ? (file as Partial<FileEntity>) : null
|
|
|
|
const convertedFile: FileEntity = {
|
|
id: baseFile?.id || (file as { id: string }).id,
|
|
type:
|
|
baseFile?.type ||
|
|
(fileType === 'image'
|
|
? 'image/png'
|
|
: fileType === 'video'
|
|
? 'video/mp4'
|
|
: fileType === 'audio'
|
|
? 'audio/mpeg'
|
|
: 'application/octet-stream'),
|
|
transferMethod:
|
|
(baseFile?.transferMethod as FileEntity['transferMethod']) ||
|
|
(fileType === 'image' ? 'remote_url' : 'local_file'),
|
|
uploadedId: baseFile?.uploadedId || (file as { id: string }).id,
|
|
supportFileType:
|
|
baseFile?.supportFileType ||
|
|
(fileType === 'image'
|
|
? 'image'
|
|
: fileType === 'video'
|
|
? 'video'
|
|
: fileType === 'audio'
|
|
? 'audio'
|
|
: 'document'),
|
|
progress: baseFile?.progress ?? 100,
|
|
name:
|
|
baseFile?.name ||
|
|
`generated_${fileType}.${fileType === 'image' ? 'png' : fileType === 'video' ? 'mp4' : fileType === 'audio' ? 'mp3' : 'bin'}`,
|
|
url: baseFile?.url || (file as { url?: string }).url,
|
|
size: baseFile?.size ?? 0, // Generated files don't have a known size
|
|
}
|
|
updateChatTreeNode(messageId, (responseItem) => {
|
|
const lastThought =
|
|
responseItem.agent_thoughts?.[responseItem.agent_thoughts?.length - 1]
|
|
if (lastThought) {
|
|
responseItem.agent_thoughts!.at(-1)!.message_files = [
|
|
...(lastThought as any).message_files,
|
|
convertedFile,
|
|
]
|
|
} else {
|
|
const currentFiles = (responseItem.message_files as FileEntity[] | undefined) ?? []
|
|
responseItem.message_files = [...currentFiles, convertedFile]
|
|
}
|
|
})
|
|
},
|
|
onThought(thought) {
|
|
updateChatTreeNode(messageId, (responseItem) => {
|
|
if (thought.message_id) responseItem.id = thought.message_id
|
|
if (thought.conversation_id) responseItem.conversationId = thought.conversation_id
|
|
|
|
if (!responseItem.agent_thoughts) responseItem.agent_thoughts = []
|
|
|
|
if (responseItem.agent_thoughts.length === 0) {
|
|
responseItem.agent_thoughts.push(thought)
|
|
} else {
|
|
const lastThought = responseItem.agent_thoughts.at(-1)
|
|
if (lastThought?.id === thought.id) {
|
|
responseItem.agent_thoughts[responseItem.agent_thoughts.length - 1] =
|
|
mergeStreamingThought(lastThought, thought)
|
|
} else {
|
|
responseItem.agent_thoughts.push(thought)
|
|
}
|
|
}
|
|
if (options.isNewAgent) {
|
|
const currentThought =
|
|
responseItem.agent_thoughts.find((item) => item.id === thought.id) ?? thought
|
|
upsertAgentResponseThoughtPart(responseItem, currentThought)
|
|
}
|
|
})
|
|
},
|
|
onMessageEnd: (messageEnd) => {
|
|
if (options.isNewAgent && messageEnd.conversation_id)
|
|
conversationIdRef.current = messageEnd.conversation_id
|
|
updateChatTreeNode(messageId, (responseItem) => {
|
|
if (messageEnd.metadata?.annotation_reply) {
|
|
responseItem.annotation = {
|
|
id: messageEnd.metadata.annotation_reply.id,
|
|
authorName: messageEnd.metadata.annotation_reply.account.name,
|
|
}
|
|
return
|
|
}
|
|
responseItem.citation = messageEnd.metadata?.retriever_resources || []
|
|
const processedFilesFromResponse = getProcessedFilesFromResponse(messageEnd.files || [])
|
|
responseItem.allFiles = uniqBy(
|
|
[...(responseItem.allFiles || []), ...(processedFilesFromResponse || [])],
|
|
'id',
|
|
)
|
|
})
|
|
},
|
|
onMessageReplace: (messageReplace) => {
|
|
updateChatTreeNode(messageId, (responseItem) => {
|
|
responseItem.content = messageReplace.answer
|
|
})
|
|
},
|
|
onError() {
|
|
handleResponding(false)
|
|
settleSend(true)
|
|
},
|
|
onWorkflowStarted: ({ workflow_run_id, task_id }) => {
|
|
handleResponding(true)
|
|
hasStopRespondedRef.current = false
|
|
updateChatTreeNode(messageId, (responseItem) => {
|
|
if (responseItem.workflowProcess && responseItem.workflowProcess.tracing.length > 0) {
|
|
responseItem.workflowProcess = {
|
|
...responseItem.workflowProcess,
|
|
status: WorkflowRunningStatus.Running,
|
|
error: undefined,
|
|
}
|
|
} else {
|
|
taskIdRef.current = task_id
|
|
responseItem.workflow_run_id = workflow_run_id
|
|
responseItem.workflowProcess = {
|
|
status: WorkflowRunningStatus.Running,
|
|
tracing: [],
|
|
}
|
|
}
|
|
})
|
|
},
|
|
onWorkflowFinished: ({ data: workflowFinishedData }) => {
|
|
updateChatTreeNode(messageId, (responseItem) => {
|
|
if (responseItem.workflowProcess) {
|
|
responseItem.workflowProcess = {
|
|
...responseItem.workflowProcess,
|
|
status: workflowFinishedData.status as WorkflowRunningStatus,
|
|
error: workflowFinishedData.error,
|
|
}
|
|
}
|
|
})
|
|
},
|
|
onIterationStart: ({ data: iterationStartedData }) => {
|
|
updateChatTreeNode(messageId, (responseItem) => {
|
|
if (!responseItem.workflowProcess) return
|
|
if (!responseItem.workflowProcess.tracing) responseItem.workflowProcess.tracing = []
|
|
responseItem.workflowProcess.tracing.push({
|
|
...iterationStartedData,
|
|
status: WorkflowRunningStatus.Running,
|
|
})
|
|
})
|
|
},
|
|
onIterationFinish: ({ data: iterationFinishedData }) => {
|
|
updateChatTreeNode(messageId, (responseItem) => {
|
|
if (!responseItem.workflowProcess?.tracing) return
|
|
const tracing = responseItem.workflowProcess.tracing
|
|
const iterationIndex = tracing.findIndex(
|
|
(item) =>
|
|
item.node_id === iterationFinishedData.node_id &&
|
|
(item.execution_metadata?.parallel_id ===
|
|
iterationFinishedData.execution_metadata?.parallel_id ||
|
|
item.parallel_id === iterationFinishedData.execution_metadata?.parallel_id),
|
|
)!
|
|
if (iterationIndex > -1) {
|
|
tracing[iterationIndex] = {
|
|
...tracing[iterationIndex],
|
|
...iterationFinishedData,
|
|
status: WorkflowRunningStatus.Succeeded,
|
|
}
|
|
}
|
|
})
|
|
},
|
|
onNodeStarted: ({ data: nodeStartedData }) => {
|
|
updateChatTreeNode(messageId, (responseItem) => {
|
|
if (!responseItem.workflowProcess) return
|
|
if (!responseItem.workflowProcess.tracing) responseItem.workflowProcess.tracing = []
|
|
|
|
const currentIndex = responseItem.workflowProcess.tracing.findIndex(
|
|
(item) => item.node_id === nodeStartedData.node_id,
|
|
)
|
|
// if the node is already started, update the node
|
|
if (currentIndex > -1) {
|
|
responseItem.workflowProcess.tracing[currentIndex] = {
|
|
...nodeStartedData,
|
|
status: NodeRunningStatus.Running,
|
|
}
|
|
} else {
|
|
if (nodeStartedData.iteration_id) return
|
|
|
|
responseItem.workflowProcess.tracing.push({
|
|
...nodeStartedData,
|
|
status: WorkflowRunningStatus.Running,
|
|
})
|
|
}
|
|
})
|
|
},
|
|
onNodeFinished: ({ data: nodeFinishedData }) => {
|
|
updateChatTreeNode(messageId, (responseItem) => {
|
|
if (!responseItem.workflowProcess?.tracing) return
|
|
|
|
if (nodeFinishedData.iteration_id) return
|
|
|
|
const currentIndex = responseItem.workflowProcess.tracing.findIndex((item) => {
|
|
if (!item.execution_metadata?.parallel_id) return item.id === nodeFinishedData.id
|
|
|
|
return (
|
|
item.id === nodeFinishedData.id &&
|
|
item.execution_metadata?.parallel_id ===
|
|
nodeFinishedData.execution_metadata?.parallel_id
|
|
)
|
|
})
|
|
if (currentIndex > -1)
|
|
responseItem.workflowProcess.tracing[currentIndex] = nodeFinishedData as any
|
|
})
|
|
},
|
|
onTTSChunk: (messageId: string, audio: string) => {
|
|
if (!audio || audio === '') return
|
|
const audioPlayer = getOrCreatePlayer()
|
|
if (audioPlayer) {
|
|
audioPlayer.playAudioWithAudio(audio, true)
|
|
AudioPlayerManager.getInstance().resetMsgId(messageId)
|
|
}
|
|
},
|
|
onTTSEnd: (messageId: string, audio: string) => {
|
|
const audioPlayer = getOrCreatePlayer()
|
|
if (audioPlayer) audioPlayer.playAudioWithAudio(audio, false)
|
|
},
|
|
onLoopStart: ({ data: loopStartedData }) => {
|
|
updateChatTreeNode(messageId, (responseItem) => {
|
|
if (!responseItem.workflowProcess) return
|
|
if (!responseItem.workflowProcess.tracing) responseItem.workflowProcess.tracing = []
|
|
responseItem.workflowProcess.tracing.push({
|
|
...loopStartedData,
|
|
status: WorkflowRunningStatus.Running,
|
|
})
|
|
})
|
|
},
|
|
onLoopFinish: ({ data: loopFinishedData }) => {
|
|
updateChatTreeNode(messageId, (responseItem) => {
|
|
if (!responseItem.workflowProcess?.tracing) return
|
|
const tracing = responseItem.workflowProcess.tracing
|
|
const loopIndex = tracing.findIndex(
|
|
(item) =>
|
|
item.node_id === loopFinishedData.node_id &&
|
|
(item.execution_metadata?.parallel_id ===
|
|
loopFinishedData.execution_metadata?.parallel_id ||
|
|
item.parallel_id === loopFinishedData.execution_metadata?.parallel_id),
|
|
)!
|
|
if (loopIndex > -1) {
|
|
tracing[loopIndex] = {
|
|
...tracing[loopIndex],
|
|
...loopFinishedData,
|
|
status: WorkflowRunningStatus.Succeeded,
|
|
}
|
|
}
|
|
})
|
|
},
|
|
onHumanInputRequired: ({ data: humanInputRequiredData }) => {
|
|
updateChatTreeNode(messageId, (responseItem) => {
|
|
if (!responseItem.humanInputFormDataList) {
|
|
responseItem.humanInputFormDataList = [humanInputRequiredData]
|
|
} else {
|
|
const currentFormIndex = responseItem.humanInputFormDataList.findIndex(
|
|
(item) => item.node_id === humanInputRequiredData.node_id,
|
|
)
|
|
if (currentFormIndex > -1) {
|
|
responseItem.humanInputFormDataList[currentFormIndex] = humanInputRequiredData
|
|
} else {
|
|
responseItem.humanInputFormDataList.push(humanInputRequiredData)
|
|
}
|
|
}
|
|
if (responseItem.workflowProcess?.tracing) {
|
|
const currentTracingIndex = responseItem.workflowProcess.tracing.findIndex(
|
|
(item) => item.node_id === humanInputRequiredData.node_id,
|
|
)
|
|
if (currentTracingIndex > -1)
|
|
responseItem.workflowProcess.tracing[currentTracingIndex]!.status =
|
|
NodeRunningStatus.Paused
|
|
}
|
|
})
|
|
},
|
|
onHumanInputFormFilled: ({ data: humanInputFilledFormData }) => {
|
|
updateChatTreeNode(messageId, (responseItem) => {
|
|
let requiredFormData:
|
|
| NonNullable<ChatItem['humanInputFormDataList']>[number]
|
|
| undefined
|
|
if (responseItem.humanInputFormDataList?.length) {
|
|
const currentFormIndex = responseItem.humanInputFormDataList.findIndex(
|
|
(item) => item.node_id === humanInputFilledFormData.node_id,
|
|
)
|
|
if (currentFormIndex > -1) {
|
|
requiredFormData = responseItem.humanInputFormDataList[currentFormIndex]
|
|
responseItem.humanInputFormDataList.splice(currentFormIndex, 1)
|
|
}
|
|
}
|
|
const enrichedHumanInputFilledFormData = enrichSubmittedHumanInputFormData(
|
|
humanInputFilledFormData,
|
|
requiredFormData,
|
|
)
|
|
if (!responseItem.humanInputFilledFormDataList) {
|
|
responseItem.humanInputFilledFormDataList = [enrichedHumanInputFilledFormData]
|
|
} else {
|
|
responseItem.humanInputFilledFormDataList.push(enrichedHumanInputFilledFormData)
|
|
}
|
|
})
|
|
},
|
|
onHumanInputFormTimeout: ({ data: humanInputFormTimeoutData }) => {
|
|
updateChatTreeNode(messageId, (responseItem) => {
|
|
if (responseItem.humanInputFormDataList?.length) {
|
|
const currentFormIndex = responseItem.humanInputFormDataList.findIndex(
|
|
(item) => item.node_id === humanInputFormTimeoutData.node_id,
|
|
)
|
|
responseItem.humanInputFormDataList[currentFormIndex]!.expiration_time =
|
|
humanInputFormTimeoutData.expiration_time
|
|
}
|
|
})
|
|
},
|
|
onWorkflowPaused: ({ data: workflowPausedData }) => {
|
|
const resumeUrl = `/workflow/${workflowPausedData.workflow_run_id}/events`
|
|
pausedStateRef.current = true
|
|
sseGet(resumeUrl, {}, otherOptions)
|
|
updateChatTreeNode(messageId, (responseItem) => {
|
|
responseItem.workflowProcess!.status = WorkflowRunningStatus.Paused
|
|
})
|
|
},
|
|
}
|
|
|
|
if (workflowEventsAbortControllerRef.current) workflowEventsAbortControllerRef.current.abort()
|
|
|
|
sseGet(url, {}, otherOptions)
|
|
},
|
|
[
|
|
updateChatTreeNode,
|
|
handleResponding,
|
|
createAudioPlayerManager,
|
|
config?.suggested_questions_after_answer,
|
|
options.isNewAgent,
|
|
],
|
|
)
|
|
|
|
const updateCurrentQAOnTree = useCallback(
|
|
({
|
|
parentId,
|
|
responseItem,
|
|
placeholderQuestionId,
|
|
questionItem,
|
|
}: {
|
|
parentId?: string
|
|
responseItem: ChatItem
|
|
placeholderQuestionId: string
|
|
questionItem: ChatItem
|
|
}) => {
|
|
let nextState: ChatItemInTree[]
|
|
const currentQA = { ...questionItem, children: [{ ...responseItem, children: [] }] }
|
|
if (
|
|
!parentId &&
|
|
!chatTree.some((item) => [placeholderQuestionId, questionItem.id].includes(item.id))
|
|
) {
|
|
// QA whose parent is not provided is considered as a first message of the conversation,
|
|
// and it should be a root node of the chat tree
|
|
nextState = produce(chatTree, (draft) => {
|
|
draft.push(currentQA)
|
|
})
|
|
} else {
|
|
// find the target QA in the tree and update it; if not found, insert it to its parent node
|
|
nextState = produceChatTreeNode(parentId!, (parentNode) => {
|
|
const questionNodeIndex = parentNode.children!.findIndex((item) =>
|
|
[placeholderQuestionId, questionItem.id].includes(item.id),
|
|
)
|
|
if (questionNodeIndex === -1) parentNode.children!.push(currentQA)
|
|
else parentNode.children![questionNodeIndex] = currentQA
|
|
})
|
|
}
|
|
setChatTree(nextState)
|
|
chatTreeRef.current = nextState
|
|
},
|
|
[chatTree, produceChatTreeNode],
|
|
)
|
|
|
|
const handleSend = useCallback(
|
|
async (
|
|
url: string,
|
|
data: {
|
|
query: string
|
|
files?: FileEntity[]
|
|
parent_message_id?: string
|
|
[key: string]: any
|
|
},
|
|
{
|
|
onGetConversationMessages,
|
|
onGetSuggestedQuestions,
|
|
onConversationComplete,
|
|
onUnhandledEvent,
|
|
onSendSettled,
|
|
isPublicAPI,
|
|
}: SendCallback,
|
|
) => {
|
|
setSuggestedQuestions([])
|
|
|
|
if (isRespondingRef.current) {
|
|
toast.info(t(($) => $['errorMessage.waitForResponse'], { ns: 'appDebug' }))
|
|
return false
|
|
}
|
|
|
|
const parentMessage = threadMessages.find((item) => item.id === data.parent_message_id)
|
|
|
|
const placeholderQuestionId = `question-${Date.now()}`
|
|
const questionItem = {
|
|
id: placeholderQuestionId,
|
|
content: data.query,
|
|
isAnswer: false,
|
|
message_files: data.files,
|
|
parentMessageId: data.parent_message_id,
|
|
}
|
|
|
|
const placeholderAnswerId = `answer-placeholder-${Date.now()}`
|
|
const placeholderAnswerItem = {
|
|
id: placeholderAnswerId,
|
|
content: '',
|
|
isAnswer: true,
|
|
parentMessageId: questionItem.id,
|
|
siblingIndex: parentMessage?.children?.length ?? chatTree.length,
|
|
}
|
|
|
|
setTargetMessageId(parentMessage?.id)
|
|
updateCurrentQAOnTree({
|
|
parentId: data.parent_message_id,
|
|
responseItem: placeholderAnswerItem,
|
|
placeholderQuestionId,
|
|
questionItem,
|
|
})
|
|
|
|
// answer
|
|
const responseItem: ChatItemInTree = {
|
|
id: placeholderAnswerId,
|
|
content: '',
|
|
agent_thoughts: [],
|
|
message_files: [],
|
|
isAnswer: true,
|
|
parentMessageId: questionItem.id,
|
|
siblingIndex: parentMessage?.children?.length ?? chatTree.length,
|
|
}
|
|
|
|
handleResponding(true)
|
|
hasStopRespondedRef.current = false
|
|
|
|
const { query, files, inputs, overrideInputsForm, ...restData } = data
|
|
const requestInputsForm = overrideInputsForm ?? formSettings?.inputsForm ?? []
|
|
const bodyParams = {
|
|
response_mode: 'streaming',
|
|
conversation_id: conversationIdRef.current,
|
|
files: getProcessedFiles(files || []),
|
|
query,
|
|
inputs: getProcessedInputs(inputs || {}, requestInputsForm),
|
|
...restData,
|
|
}
|
|
if (bodyParams?.files?.length) {
|
|
bodyParams.files = bodyParams.files.map((item) => {
|
|
if (item.transfer_method === TransferMethod.local_file) {
|
|
return {
|
|
...item,
|
|
url: '',
|
|
}
|
|
}
|
|
return item
|
|
})
|
|
}
|
|
|
|
let isAgentMode = false
|
|
let hasSetResponseId = false
|
|
let hasSettled = false
|
|
const settleSend = (hasError?: boolean) => {
|
|
if (hasSettled) return
|
|
|
|
hasSettled = true
|
|
onSendSettled?.(hasError)
|
|
}
|
|
|
|
const getOrCreatePlayer = createAudioPlayerManager()
|
|
|
|
const otherOptions: IOtherOptions = {
|
|
isPublicAPI,
|
|
onUnhandledEvent,
|
|
getAbortController: (abortController) => {
|
|
workflowEventsAbortControllerRef.current = abortController
|
|
},
|
|
onData: (
|
|
message: string,
|
|
isFirstMessage: boolean,
|
|
{ event, conversationId: newConversationId, messageId, taskId }: any,
|
|
) => {
|
|
const isNewAgentMessage =
|
|
options.isNewAgent && (event === 'agent_message' || event === 'message')
|
|
if (isNewAgentMessage) {
|
|
appendAgentResponseMessagePart(responseItem, message)
|
|
} else if (!isAgentMode || options.isNewAgent) {
|
|
responseItem.content = responseItem.content + message
|
|
} else {
|
|
const lastThought =
|
|
responseItem.agent_thoughts?.[responseItem.agent_thoughts.length - 1]
|
|
if (lastThought) lastThought.thought = lastThought.thought + message // need immer setAutoFreeze
|
|
}
|
|
|
|
if (messageId && !hasSetResponseId) {
|
|
questionItem.id = `question-${messageId}`
|
|
responseItem.id = messageId
|
|
responseItem.parentMessageId = questionItem.id
|
|
hasSetResponseId = true
|
|
}
|
|
|
|
if (isFirstMessage && newConversationId) conversationIdRef.current = newConversationId
|
|
|
|
taskIdRef.current = taskId
|
|
if (messageId) responseItem.id = messageId
|
|
|
|
updateCurrentQAOnTree({
|
|
placeholderQuestionId,
|
|
questionItem,
|
|
responseItem,
|
|
parentId: data.parent_message_id,
|
|
})
|
|
},
|
|
onReasoning: ({ data: reasoningData }: ReasoningChunkResponse) => {
|
|
const { reasoning, node_id, is_final } = reasoningData
|
|
const reasoningContent =
|
|
responseItem.reasoningContent || (responseItem.reasoningContent = {})
|
|
const key = node_id || '_'
|
|
if (reasoning) reasoningContent[key] = (reasoningContent[key] || '') + reasoning
|
|
if (is_final) responseItem.reasoningFinished = true
|
|
|
|
updateCurrentQAOnTree({
|
|
placeholderQuestionId,
|
|
questionItem,
|
|
responseItem,
|
|
parentId: data.parent_message_id,
|
|
})
|
|
},
|
|
async onCompleted(hasError?: boolean) {
|
|
handleResponding(false)
|
|
|
|
try {
|
|
if (hasError) return
|
|
|
|
let completedWorkflowRunId = responseItem.workflow_run_id
|
|
|
|
if (
|
|
conversationIdRef.current &&
|
|
!hasStopRespondedRef.current &&
|
|
onGetConversationMessages
|
|
) {
|
|
const conversationMessagesResponse = await onGetConversationMessages(
|
|
conversationIdRef.current,
|
|
(newAbortController) =>
|
|
(conversationMessagesAbortControllerRef.current = newAbortController),
|
|
)
|
|
const data = getConversationMessagesData(conversationMessagesResponse)
|
|
const newResponseItem = data.find((item) => item.id === responseItem.id)
|
|
completedWorkflowRunId = newResponseItem?.workflow_run_id ?? completedWorkflowRunId
|
|
if (!newResponseItem)
|
|
return onConversationComplete?.(conversationIdRef.current, completedWorkflowRunId)
|
|
|
|
const historyAgentThoughts = getHistoryAgentThoughts(newResponseItem)
|
|
const lastHistoryAgentThought = historyAgentThoughts.at(-1)
|
|
const historyAnswer = newResponseItem.answer || ''
|
|
const isUseAgentThought =
|
|
!options.isNewAgent && lastHistoryAgentThought?.thought === historyAnswer
|
|
const messageLog = Array.isArray(newResponseItem.message)
|
|
? newResponseItem.message
|
|
: []
|
|
const answerTokens = newResponseItem.answer_tokens ?? 0
|
|
const messageTokens = newResponseItem.message_tokens ?? 0
|
|
const providerResponseLatency = newResponseItem.provider_response_latency ?? 0
|
|
const historyAnswerFiles = getHistoryAnswerFiles(newResponseItem)
|
|
updateChatTreeNode(responseItem.id, {
|
|
content: isUseAgentThought ? '' : historyAnswer,
|
|
agent_thoughts: historyAgentThoughts,
|
|
agent_response_parts: undefined,
|
|
citation: newResponseItem.retriever_resources,
|
|
reasoningContent: newResponseItem.metadata?.reasoning,
|
|
reasoningFinished: true,
|
|
message_files: historyAnswerFiles,
|
|
allFiles: undefined,
|
|
workflowProcess: undefined,
|
|
workflow_run_id: newResponseItem.workflow_run_id ?? completedWorkflowRunId,
|
|
feedback: newResponseItem.feedback,
|
|
log: [
|
|
...messageLog,
|
|
...(messageLog.at(-1)?.role !== 'assistant'
|
|
? [
|
|
{
|
|
role: 'assistant',
|
|
text: historyAnswer,
|
|
files: historyAnswerFiles,
|
|
},
|
|
]
|
|
: []),
|
|
],
|
|
more: {
|
|
time: formatTime(newResponseItem.created_at ?? Date.now(), 'hh:mm A'),
|
|
tokens: answerTokens + messageTokens,
|
|
latency: providerResponseLatency.toFixed(2),
|
|
tokens_per_second:
|
|
providerResponseLatency > 0
|
|
? (answerTokens / providerResponseLatency).toFixed(2)
|
|
: undefined,
|
|
},
|
|
// for agent log
|
|
conversationId: conversationIdRef.current,
|
|
input: {
|
|
inputs: newResponseItem.inputs,
|
|
query: newResponseItem.query,
|
|
},
|
|
})
|
|
}
|
|
|
|
onConversationComplete?.(conversationIdRef.current, completedWorkflowRunId)
|
|
|
|
if (
|
|
config?.suggested_questions_after_answer?.enabled &&
|
|
!hasStopRespondedRef.current &&
|
|
onGetSuggestedQuestions
|
|
) {
|
|
try {
|
|
const { data }: any = await onGetSuggestedQuestions(
|
|
responseItem.id,
|
|
(newAbortController) =>
|
|
(suggestedQuestionsAbortControllerRef.current = newAbortController),
|
|
)
|
|
setSuggestedQuestions(data)
|
|
} catch {
|
|
setSuggestedQuestions([])
|
|
}
|
|
}
|
|
} finally {
|
|
settleSend(hasError)
|
|
}
|
|
},
|
|
onFile(file) {
|
|
// Convert simple file type to MIME type for non-agent mode
|
|
// Backend sends: { id, type: "image", belongs_to, url }
|
|
// Frontend expects: { id, type: "image/png", transferMethod, url, uploadedId, supportFileType, name, size }
|
|
|
|
// Determine file type for MIME conversion
|
|
const fileType = (file as { type?: string }).type || 'image'
|
|
|
|
// If file already has transferMethod, use it as base and ensure all required fields exist
|
|
// Otherwise, create a new complete file object
|
|
const baseFile = 'transferMethod' in file ? (file as Partial<FileEntity>) : null
|
|
|
|
const convertedFile: FileEntity = {
|
|
id: baseFile?.id || (file as { id: string }).id,
|
|
type:
|
|
baseFile?.type ||
|
|
(fileType === 'image'
|
|
? 'image/png'
|
|
: fileType === 'video'
|
|
? 'video/mp4'
|
|
: fileType === 'audio'
|
|
? 'audio/mpeg'
|
|
: 'application/octet-stream'),
|
|
transferMethod:
|
|
(baseFile?.transferMethod as FileEntity['transferMethod']) ||
|
|
(fileType === 'image' ? 'remote_url' : 'local_file'),
|
|
uploadedId: baseFile?.uploadedId || (file as { id: string }).id,
|
|
supportFileType:
|
|
baseFile?.supportFileType ||
|
|
(fileType === 'image'
|
|
? 'image'
|
|
: fileType === 'video'
|
|
? 'video'
|
|
: fileType === 'audio'
|
|
? 'audio'
|
|
: 'document'),
|
|
progress: baseFile?.progress ?? 100,
|
|
name:
|
|
baseFile?.name ||
|
|
`generated_${fileType}.${fileType === 'image' ? 'png' : fileType === 'video' ? 'mp4' : fileType === 'audio' ? 'mp3' : 'bin'}`,
|
|
url: baseFile?.url || (file as { url?: string }).url,
|
|
size: baseFile?.size ?? 0, // Generated files don't have a known size
|
|
}
|
|
|
|
// For agent mode, add files to the last thought
|
|
const lastThought = responseItem.agent_thoughts?.[responseItem.agent_thoughts?.length - 1]
|
|
if (lastThought) {
|
|
const thought = lastThought as { message_files?: FileEntity[] }
|
|
responseItem.agent_thoughts!.at(-1)!.message_files = [
|
|
...(thought.message_files ?? []),
|
|
convertedFile,
|
|
]
|
|
}
|
|
// For non-agent mode, add files directly to responseItem.message_files
|
|
else {
|
|
const currentFiles = (responseItem.message_files as FileEntity[] | undefined) ?? []
|
|
responseItem.message_files = [...currentFiles, convertedFile]
|
|
}
|
|
|
|
updateCurrentQAOnTree({
|
|
placeholderQuestionId,
|
|
questionItem,
|
|
responseItem,
|
|
parentId: data.parent_message_id,
|
|
})
|
|
},
|
|
onThought(thought) {
|
|
isAgentMode = true
|
|
const response = responseItem as any
|
|
if (thought.message_id && !hasSetResponseId) response.id = thought.message_id
|
|
if (thought.conversation_id) response.conversationId = thought.conversation_id
|
|
|
|
if (response.agent_thoughts.length === 0) {
|
|
response.agent_thoughts.push(thought)
|
|
} else {
|
|
const lastThought = response.agent_thoughts.at(-1)
|
|
// thought changed but still the same thought, so update.
|
|
if (lastThought.id === thought.id) {
|
|
responseItem.agent_thoughts![response.agent_thoughts.length - 1] =
|
|
mergeStreamingThought(lastThought, thought)
|
|
} else {
|
|
responseItem.agent_thoughts!.push(thought)
|
|
}
|
|
}
|
|
if (options.isNewAgent) {
|
|
const currentThought =
|
|
responseItem.agent_thoughts?.find((item) => item.id === thought.id) ?? thought
|
|
upsertAgentResponseThoughtPart(responseItem, currentThought)
|
|
}
|
|
updateCurrentQAOnTree({
|
|
placeholderQuestionId,
|
|
questionItem,
|
|
responseItem,
|
|
parentId: data.parent_message_id,
|
|
})
|
|
},
|
|
onMessageEnd: (messageEnd) => {
|
|
if (options.isNewAgent && messageEnd.conversation_id)
|
|
conversationIdRef.current = messageEnd.conversation_id
|
|
if (messageEnd.metadata?.annotation_reply) {
|
|
responseItem.id = messageEnd.id
|
|
responseItem.annotation = {
|
|
id: messageEnd.metadata.annotation_reply.id,
|
|
authorName: messageEnd.metadata.annotation_reply.account.name,
|
|
}
|
|
updateCurrentQAOnTree({
|
|
placeholderQuestionId,
|
|
questionItem,
|
|
responseItem,
|
|
parentId: data.parent_message_id,
|
|
})
|
|
handleResponding(false)
|
|
return
|
|
}
|
|
responseItem.citation = messageEnd.metadata?.retriever_resources || []
|
|
const processedFilesFromResponse = getProcessedFilesFromResponse(messageEnd.files || [])
|
|
responseItem.allFiles = uniqBy(
|
|
[...(responseItem.allFiles || []), ...(processedFilesFromResponse || [])],
|
|
'id',
|
|
)
|
|
|
|
updateCurrentQAOnTree({
|
|
placeholderQuestionId,
|
|
questionItem,
|
|
responseItem,
|
|
parentId: data.parent_message_id,
|
|
})
|
|
},
|
|
onMessageReplace: (messageReplace) => {
|
|
responseItem.content = messageReplace.answer
|
|
},
|
|
onError() {
|
|
handleResponding(false)
|
|
settleSend(true)
|
|
updateCurrentQAOnTree({
|
|
placeholderQuestionId,
|
|
questionItem,
|
|
responseItem,
|
|
parentId: data.parent_message_id,
|
|
})
|
|
},
|
|
onWorkflowStarted: ({ workflow_run_id, task_id, conversation_id, message_id }) => {
|
|
handleResponding(true)
|
|
// If there are no streaming messages, we still need to set the conversation_id to avoid create a new conversation when regeneration in chat-flow.
|
|
if (conversation_id) {
|
|
conversationIdRef.current = conversation_id
|
|
}
|
|
if (message_id && !hasSetResponseId) {
|
|
questionItem.id = `question-${message_id}`
|
|
responseItem.id = message_id
|
|
responseItem.parentMessageId = questionItem.id
|
|
hasSetResponseId = true
|
|
}
|
|
|
|
if (responseItem.workflowProcess && responseItem.workflowProcess.tracing.length > 0) {
|
|
responseItem.workflowProcess = {
|
|
...responseItem.workflowProcess,
|
|
status: WorkflowRunningStatus.Running,
|
|
error: undefined,
|
|
}
|
|
} else {
|
|
taskIdRef.current = task_id
|
|
responseItem.workflow_run_id = workflow_run_id
|
|
responseItem.workflowProcess = {
|
|
status: WorkflowRunningStatus.Running,
|
|
tracing: [],
|
|
}
|
|
}
|
|
updateCurrentQAOnTree({
|
|
placeholderQuestionId,
|
|
questionItem,
|
|
responseItem,
|
|
parentId: data.parent_message_id,
|
|
})
|
|
},
|
|
onWorkflowFinished: ({ data: workflowFinishedData }) => {
|
|
if (pausedStateRef.current) pausedStateRef.current = false
|
|
responseItem.workflowProcess = {
|
|
...responseItem.workflowProcess!,
|
|
status: workflowFinishedData.status as WorkflowRunningStatus,
|
|
error: workflowFinishedData.error,
|
|
}
|
|
updateCurrentQAOnTree({
|
|
placeholderQuestionId,
|
|
questionItem,
|
|
responseItem,
|
|
parentId: data.parent_message_id,
|
|
})
|
|
},
|
|
onIterationStart: ({ data: iterationStartedData }) => {
|
|
responseItem.workflowProcess!.tracing!.push({
|
|
...iterationStartedData,
|
|
status: WorkflowRunningStatus.Running,
|
|
})
|
|
updateCurrentQAOnTree({
|
|
placeholderQuestionId,
|
|
questionItem,
|
|
responseItem,
|
|
parentId: data.parent_message_id,
|
|
})
|
|
},
|
|
onIterationFinish: ({ data: iterationFinishedData }) => {
|
|
const tracing = responseItem.workflowProcess!.tracing!
|
|
const iterationIndex = tracing.findIndex(
|
|
(item) =>
|
|
item.node_id === iterationFinishedData.node_id &&
|
|
(item.execution_metadata?.parallel_id ===
|
|
iterationFinishedData.execution_metadata?.parallel_id ||
|
|
item.parallel_id === iterationFinishedData.execution_metadata?.parallel_id),
|
|
)!
|
|
tracing[iterationIndex] = {
|
|
...tracing[iterationIndex],
|
|
...iterationFinishedData,
|
|
status: WorkflowRunningStatus.Succeeded,
|
|
}
|
|
|
|
updateCurrentQAOnTree({
|
|
placeholderQuestionId,
|
|
questionItem,
|
|
responseItem,
|
|
parentId: data.parent_message_id,
|
|
})
|
|
},
|
|
onNodeStarted: ({ data: nodeStartedData }) => {
|
|
if (!responseItem.workflowProcess) return
|
|
if (!responseItem.workflowProcess.tracing) responseItem.workflowProcess.tracing = []
|
|
|
|
const currentIndex = responseItem.workflowProcess.tracing.findIndex(
|
|
(item) => item.node_id === nodeStartedData.node_id,
|
|
)
|
|
if (currentIndex > -1) {
|
|
responseItem.workflowProcess.tracing[currentIndex] = {
|
|
...nodeStartedData,
|
|
status: NodeRunningStatus.Running,
|
|
}
|
|
} else {
|
|
if (nodeStartedData.iteration_id) return
|
|
|
|
if (data.loop_id) return
|
|
|
|
responseItem.workflowProcess.tracing.push({
|
|
...nodeStartedData,
|
|
status: WorkflowRunningStatus.Running,
|
|
})
|
|
}
|
|
updateCurrentQAOnTree({
|
|
placeholderQuestionId,
|
|
questionItem,
|
|
responseItem,
|
|
parentId: data.parent_message_id,
|
|
})
|
|
},
|
|
onNodeFinished: ({ data: nodeFinishedData }) => {
|
|
if (nodeFinishedData.iteration_id) return
|
|
|
|
if (data.loop_id) return
|
|
|
|
const currentIndex = responseItem.workflowProcess!.tracing!.findIndex((item) => {
|
|
if (!item.execution_metadata?.parallel_id) return item.id === nodeFinishedData.id
|
|
|
|
return (
|
|
item.id === nodeFinishedData.id &&
|
|
item.execution_metadata?.parallel_id ===
|
|
nodeFinishedData.execution_metadata?.parallel_id
|
|
)
|
|
})
|
|
responseItem.workflowProcess!.tracing[currentIndex] = nodeFinishedData as any
|
|
|
|
updateCurrentQAOnTree({
|
|
placeholderQuestionId,
|
|
questionItem,
|
|
responseItem,
|
|
parentId: data.parent_message_id,
|
|
})
|
|
},
|
|
onTTSChunk: (messageId: string, audio: string) => {
|
|
if (!audio || audio === '') return
|
|
const audioPlayer = getOrCreatePlayer()
|
|
if (audioPlayer) {
|
|
audioPlayer.playAudioWithAudio(audio, true)
|
|
AudioPlayerManager.getInstance().resetMsgId(messageId)
|
|
}
|
|
},
|
|
onTTSEnd: (messageId: string, audio: string) => {
|
|
const audioPlayer = getOrCreatePlayer()
|
|
if (audioPlayer) audioPlayer.playAudioWithAudio(audio, false)
|
|
},
|
|
onLoopStart: ({ data: loopStartedData }) => {
|
|
responseItem.workflowProcess!.tracing!.push({
|
|
...loopStartedData,
|
|
status: WorkflowRunningStatus.Running,
|
|
})
|
|
updateCurrentQAOnTree({
|
|
placeholderQuestionId,
|
|
questionItem,
|
|
responseItem,
|
|
parentId: data.parent_message_id,
|
|
})
|
|
},
|
|
onLoopFinish: ({ data: loopFinishedData }) => {
|
|
const tracing = responseItem.workflowProcess!.tracing!
|
|
const loopIndex = tracing.findIndex(
|
|
(item) =>
|
|
item.node_id === loopFinishedData.node_id &&
|
|
(item.execution_metadata?.parallel_id ===
|
|
loopFinishedData.execution_metadata?.parallel_id ||
|
|
item.parallel_id === loopFinishedData.execution_metadata?.parallel_id),
|
|
)!
|
|
tracing[loopIndex] = {
|
|
...tracing[loopIndex],
|
|
...loopFinishedData,
|
|
status: WorkflowRunningStatus.Succeeded,
|
|
}
|
|
|
|
updateCurrentQAOnTree({
|
|
placeholderQuestionId,
|
|
questionItem,
|
|
responseItem,
|
|
parentId: data.parent_message_id,
|
|
})
|
|
},
|
|
onHumanInputRequired: ({ data: humanInputRequiredData }) => {
|
|
if (!responseItem.humanInputFormDataList) {
|
|
responseItem.humanInputFormDataList = [humanInputRequiredData]
|
|
} else {
|
|
const currentFormIndex = responseItem.humanInputFormDataList!.findIndex(
|
|
(item) => item.node_id === humanInputRequiredData.node_id,
|
|
)
|
|
if (currentFormIndex > -1) {
|
|
responseItem.humanInputFormDataList[currentFormIndex] = humanInputRequiredData
|
|
} else {
|
|
responseItem.humanInputFormDataList.push(humanInputRequiredData)
|
|
}
|
|
}
|
|
const currentTracingIndex = responseItem.workflowProcess!.tracing!.findIndex(
|
|
(item) => item.node_id === humanInputRequiredData.node_id,
|
|
)
|
|
if (currentTracingIndex > -1) {
|
|
responseItem.workflowProcess!.tracing[currentTracingIndex]!.status =
|
|
NodeRunningStatus.Paused
|
|
updateCurrentQAOnTree({
|
|
placeholderQuestionId,
|
|
questionItem,
|
|
responseItem,
|
|
parentId: data.parent_message_id,
|
|
})
|
|
}
|
|
},
|
|
onHumanInputFormFilled: ({ data: humanInputFilledFormData }) => {
|
|
let requiredFormData: NonNullable<ChatItem['humanInputFormDataList']>[number] | undefined
|
|
if (responseItem.humanInputFormDataList?.length) {
|
|
const currentFormIndex = responseItem.humanInputFormDataList!.findIndex(
|
|
(item) => item.node_id === humanInputFilledFormData.node_id,
|
|
)
|
|
if (currentFormIndex > -1) {
|
|
requiredFormData = responseItem.humanInputFormDataList[currentFormIndex]
|
|
responseItem.humanInputFormDataList.splice(currentFormIndex, 1)
|
|
}
|
|
}
|
|
const enrichedHumanInputFilledFormData = enrichSubmittedHumanInputFormData(
|
|
humanInputFilledFormData,
|
|
requiredFormData,
|
|
)
|
|
if (!responseItem.humanInputFilledFormDataList) {
|
|
responseItem.humanInputFilledFormDataList = [enrichedHumanInputFilledFormData]
|
|
} else {
|
|
responseItem.humanInputFilledFormDataList.push(enrichedHumanInputFilledFormData)
|
|
}
|
|
updateCurrentQAOnTree({
|
|
placeholderQuestionId,
|
|
questionItem,
|
|
responseItem,
|
|
parentId: data.parent_message_id,
|
|
})
|
|
},
|
|
onHumanInputFormTimeout: ({ data: humanInputFormTimeoutData }) => {
|
|
if (responseItem.humanInputFormDataList?.length) {
|
|
const currentFormIndex = responseItem.humanInputFormDataList!.findIndex(
|
|
(item) => item.node_id === humanInputFormTimeoutData.node_id,
|
|
)
|
|
responseItem.humanInputFormDataList[currentFormIndex]!.expiration_time =
|
|
humanInputFormTimeoutData.expiration_time
|
|
}
|
|
updateCurrentQAOnTree({
|
|
placeholderQuestionId,
|
|
questionItem,
|
|
responseItem,
|
|
parentId: data.parent_message_id,
|
|
})
|
|
},
|
|
onWorkflowPaused: ({ data: workflowPausedData }) => {
|
|
const url = `/workflow/${workflowPausedData.workflow_run_id}/events`
|
|
pausedStateRef.current = true
|
|
sseGet(url, {}, otherOptions)
|
|
responseItem.workflowProcess!.status = WorkflowRunningStatus.Paused
|
|
updateCurrentQAOnTree({
|
|
placeholderQuestionId,
|
|
questionItem,
|
|
responseItem,
|
|
parentId: data.parent_message_id,
|
|
})
|
|
},
|
|
}
|
|
|
|
// Abort the previous workflow events SSE request
|
|
if (workflowEventsAbortControllerRef.current) workflowEventsAbortControllerRef.current.abort()
|
|
|
|
ssePost(
|
|
url,
|
|
{
|
|
body: bodyParams,
|
|
},
|
|
otherOptions,
|
|
)
|
|
return true
|
|
},
|
|
[
|
|
t,
|
|
chatTree.length,
|
|
threadMessages,
|
|
config?.suggested_questions_after_answer,
|
|
updateCurrentQAOnTree,
|
|
updateChatTreeNode,
|
|
handleResponding,
|
|
formatTime,
|
|
createAudioPlayerManager,
|
|
formSettings,
|
|
options.isNewAgent,
|
|
],
|
|
)
|
|
|
|
const handleAnnotationEdited = useCallback(
|
|
(query: string, answer: string, index: number) => {
|
|
const targetQuestionId = chatList[index - 1]!.id
|
|
const targetAnswerId = chatList[index]!.id
|
|
|
|
updateChatTreeNode(targetQuestionId, {
|
|
content: query,
|
|
})
|
|
updateChatTreeNode(targetAnswerId, {
|
|
content: answer,
|
|
annotation: {
|
|
...chatList[index]!.annotation,
|
|
logAnnotation: undefined,
|
|
} as any,
|
|
})
|
|
},
|
|
[chatList, updateChatTreeNode],
|
|
)
|
|
|
|
const handleAnnotationAdded = useCallback(
|
|
(annotationId: string, authorName: string, query: string, answer: string, index: number) => {
|
|
const targetQuestionId = chatList[index - 1]!.id
|
|
const targetAnswerId = chatList[index]!.id
|
|
|
|
updateChatTreeNode(targetQuestionId, {
|
|
content: query,
|
|
})
|
|
|
|
updateChatTreeNode(targetAnswerId, {
|
|
content: chatList[index]!.content,
|
|
annotation: {
|
|
id: annotationId,
|
|
authorName,
|
|
logAnnotation: {
|
|
content: answer,
|
|
account: {
|
|
id: '',
|
|
name: authorName,
|
|
email: '',
|
|
},
|
|
},
|
|
} as Annotation,
|
|
})
|
|
},
|
|
[chatList, updateChatTreeNode],
|
|
)
|
|
|
|
const handleAnnotationRemoved = useCallback(
|
|
(index: number) => {
|
|
const targetAnswerId = chatList[index]!.id
|
|
|
|
updateChatTreeNode(targetAnswerId, {
|
|
content: chatList[index]!.content,
|
|
annotation: {
|
|
...chatList[index]!.annotation,
|
|
id: '',
|
|
} as Annotation,
|
|
})
|
|
},
|
|
[chatList, updateChatTreeNode],
|
|
)
|
|
|
|
const handleSwitchSibling = useCallback(
|
|
(siblingMessageId: string, callbacks: SendCallback) => {
|
|
setTargetMessageId(siblingMessageId)
|
|
|
|
// Helper to find message in tree
|
|
const findMessageInTree = (
|
|
nodes: ChatItemInTree[],
|
|
targetId: string,
|
|
): ChatItemInTree | undefined => {
|
|
for (const node of nodes) {
|
|
if (node.id === targetId) return node
|
|
if (node.children) {
|
|
const found = findMessageInTree(node.children, targetId)
|
|
if (found) return found
|
|
}
|
|
}
|
|
return undefined
|
|
}
|
|
|
|
const targetMessage = findMessageInTree(chatTreeRef.current, siblingMessageId)
|
|
if (
|
|
targetMessage?.workflow_run_id &&
|
|
targetMessage.humanInputFormDataList &&
|
|
targetMessage.humanInputFormDataList.length > 0
|
|
) {
|
|
handleResume(targetMessage.id, targetMessage.workflow_run_id, callbacks)
|
|
}
|
|
},
|
|
[setTargetMessageId, handleResume],
|
|
)
|
|
|
|
useEffect(() => {
|
|
if (clearChatList) handleRestart(() => clearChatListCallback?.(false))
|
|
}, [clearChatList, clearChatListCallback, handleRestart])
|
|
|
|
return {
|
|
chatList,
|
|
setTargetMessageId,
|
|
isResponding,
|
|
setIsResponding,
|
|
handleSend,
|
|
handleResume,
|
|
handleSwitchSibling,
|
|
suggestedQuestions,
|
|
handleRestart,
|
|
handleStop,
|
|
handleAnnotationEdited,
|
|
handleAnnotationAdded,
|
|
handleAnnotationRemoved,
|
|
}
|
|
}
|