From a6aae910d0c2fac5c92c7f80555ac903ff6c07c7 Mon Sep 17 00:00:00 2001 From: Neko Ayaka Date: Sat, 27 Dec 2025 17:30:04 +0800 Subject: [PATCH] feat(stage-ui): subscribe to onChatTurnComplete, emit to server channel as output:gen-ai:chat:message --- packages/stage-ui/src/stores/chat.ts | 52 ++++++++++++++++--- .../src/stores/mods/api/context-bridge.ts | 12 +++++ packages/stage-ui/src/utils/queue.ts | 1 + 3 files changed, 57 insertions(+), 8 deletions(-) diff --git a/packages/stage-ui/src/stores/chat.ts b/packages/stage-ui/src/stores/chat.ts index 3e242f4fa..15390689f 100644 --- a/packages/stage-ui/src/stores/chat.ts +++ b/packages/stage-ui/src/stores/chat.ts @@ -98,7 +98,8 @@ export const useChatStore = defineStore('chat', () => { const onTokenSpecialHooks = ref Promise>>([]) const onStreamEndHooks = ref Promise>>([]) const onAssistantResponseEndHooks = ref Promise>>([]) - const onContextPublishHooks = ref Promise | void>>([]) + const onAssistantMessageHooks = ref Promise>>([]) + const onChatTurnCompleteHooks = ref, composedMessage: Message[], output: StreamingAssistantMessage }) => Promise>>([]) function onBeforeMessageComposed(cb: (message: string) => Promise) { onBeforeMessageComposedHooks.value.push(cb) @@ -140,6 +141,16 @@ export const useChatStore = defineStore('chat', () => { return () => onAssistantResponseEndHooks.value = onAssistantResponseEndHooks.value.filter(hook => hook !== cb) // return remove listener callback } + function onAssistantMessage(cb: (message: StreamingAssistantMessage) => Promise) { + onAssistantMessageHooks.value.push(cb) + return () => onAssistantMessageHooks.value = onAssistantMessageHooks.value.filter(hook => hook !== cb) // return remove listener callback + } + + function onChatTurnComplete(cb: (chat: { input: ChatHistoryItem, contexts: Record, composedMessage: Message[], output: StreamingAssistantMessage }) => Promise) { + onChatTurnCompleteHooks.value.push(cb) + return () => onChatTurnCompleteHooks.value = onChatTurnCompleteHooks.value.filter(hook => hook !== cb) // return remove listener callback + } + function clearHooks() { onBeforeMessageComposedHooks.value = [] onAfterMessageComposedHooks.value = [] @@ -149,7 +160,8 @@ export const useChatStore = defineStore('chat', () => { onTokenSpecialHooks.value = [] onStreamEndHooks.value = [] onAssistantResponseEndHooks.value = [] - onContextPublishHooks.value = [] + onAssistantMessageHooks.value = [] + onChatTurnCompleteHooks.value = [] } async function emitBeforeMessageComposedHooks(message: string) { @@ -192,6 +204,16 @@ export const useChatStore = defineStore('chat', () => { await hook(message) } + async function emitAssistantMessageHooks(message: StreamingAssistantMessage) { + for (const hook of onAssistantMessageHooks.value) + await hook(message) + } + + async function emitChatTurnCompleteHooks(chat: { input: ChatHistoryItem, contexts: Record, composedMessage: Message[], output: StreamingAssistantMessage }) { + for (const hook of onChatTurnCompleteHooks.value) + await hook(chat) + } + // ----- Session state helpers ----- // I know this nu uh, better than loading all language on rehypeShiki const codeBlockSystemPrompt = '- For any programming code block, always specify the programming language that supported on @shikijs/rehype on the rendered markdown, eg. ```python ... ```\n' @@ -325,6 +347,7 @@ export const useChatStore = defineStore('chat', () => { ensureSession(sessionId) + const sendingCreatedAt = Date.now() const isStaleGeneration = () => getSessionGeneration(sessionId) !== generation const shouldAbort = () => isStaleGeneration() if (shouldAbort()) @@ -463,6 +486,7 @@ export const useChatStore = defineStore('chat', () => { type: 'tool-call', toolCall: event, }) + break case 'tool-result': toolCallQueue.enqueue({ @@ -470,6 +494,7 @@ export const useChatStore = defineStore('chat', () => { id: event.toolCallId, result: event.result, }) + break case 'text-delta': fullText += event.text @@ -483,6 +508,7 @@ export const useChatStore = defineStore('chat', () => { } }, }) + // Finalize the parsing of the actual message content await parser.end() @@ -491,9 +517,6 @@ export const useChatStore = defineStore('chat', () => { sessionMessagesForSend.push(toRaw(streamingMessage.value)) } - // Reset the streaming message for the next turn - streamingMessage.value = { role: 'assistant', content: '', slices: [], tool_results: [] } - // Instruct the TTS pipeline to flush by calling hooks directly const flushSignal = `${TTS_FLUSH_INSTRUCTION}${TTS_FLUSH_INSTRUCTION}` await emitTokenLiteralHooks(flushSignal) @@ -504,10 +527,17 @@ export const useChatStore = defineStore('chat', () => { // Call the end-of-response hooks with the full text await emitAssistantResponseEndHooks(fullText) - // eslint-disable-next-line no-console - console.debug('LLM output:', fullText) - await emitAfterSendHooks(sendingMessage) + await emitAssistantMessageHooks({ ...streamingMessage.value }) + await emitChatTurnCompleteHooks({ + input: { role: 'user', content: sendingMessage, createdAt: sendingCreatedAt }, + contexts: { ...activeContexts.value }, + composedMessage: newMessages as Message[], + output: { ...streamingMessage.value }, + }) + + // Reset the streaming message for the next turn + streamingMessage.value = { role: 'assistant', content: '', slices: [], tool_results: [] } } catch (error) { console.error('Error sending message:', error) @@ -551,7 +581,9 @@ export const useChatStore = defineStore('chat', () => { getAllSessions, replaceSessions, resetAllSessions, + clearHooks, + emitBeforeMessageComposedHooks, emitAfterMessageComposedHooks, emitBeforeSendHooks, @@ -560,6 +592,8 @@ export const useChatStore = defineStore('chat', () => { emitTokenSpecialHooks, emitStreamEndHooks, emitAssistantResponseEndHooks, + emitAssistantMessageHooks, + emitChatTurnCompleteHooks, onBeforeMessageComposed, onAfterMessageComposed, @@ -569,5 +603,7 @@ export const useChatStore = defineStore('chat', () => { onTokenSpecial, onStreamEnd, onAssistantResponseEnd, + onAssistantMessage, + onChatTurnComplete, } }) diff --git a/packages/stage-ui/src/stores/mods/api/context-bridge.ts b/packages/stage-ui/src/stores/mods/api/context-bridge.ts index 8e774d7b8..660339692 100644 --- a/packages/stage-ui/src/stores/mods/api/context-bridge.ts +++ b/packages/stage-ui/src/stores/mods/api/context-bridge.ts @@ -40,45 +40,57 @@ export const useContextBridgeStore = defineStore('mods:api:context-bridge', () = chatStore.onBeforeMessageComposed(async (message) => { if (isProcessingRemoteStream) return + broadcastStreamEvent({ type: 'before-compose', message, sessionId: chatStore.activeSessionId }) }), chatStore.onAfterMessageComposed(async (message) => { if (isProcessingRemoteStream) return + broadcastStreamEvent({ type: 'after-compose', message, sessionId: chatStore.activeSessionId }) }), chatStore.onBeforeSend(async (message) => { if (isProcessingRemoteStream) return + broadcastStreamEvent({ type: 'before-send', message, sessionId: chatStore.activeSessionId }) }), chatStore.onAfterSend(async (message) => { if (isProcessingRemoteStream) return + broadcastStreamEvent({ type: 'after-send', message, sessionId: chatStore.activeSessionId }) }), chatStore.onTokenLiteral(async (literal) => { if (isProcessingRemoteStream) return + broadcastStreamEvent({ type: 'token-literal', literal, sessionId: chatStore.activeSessionId }) }), chatStore.onTokenSpecial(async (special) => { if (isProcessingRemoteStream) return + broadcastStreamEvent({ type: 'token-special', special, sessionId: chatStore.activeSessionId }) }), chatStore.onStreamEnd(async () => { if (isProcessingRemoteStream) return + broadcastStreamEvent({ type: 'stream-end', sessionId: chatStore.activeSessionId }) }), chatStore.onAssistantResponseEnd(async (message) => { if (isProcessingRemoteStream) return + broadcastStreamEvent({ type: 'assistant-end', message, sessionId: chatStore.activeSessionId }) }), ) + disposeHookFns.value.push(chatStore.onChatTurnComplete(async (chat) => { + serverChannelStore.send({ type: 'output:gen-ai:chat:message', data: { messages: [chat.output] } }) + })) + const { stop: stopIncomingStreamWatch } = watch(incomingStreamEvent, async (event) => { if (!event) return diff --git a/packages/stage-ui/src/utils/queue.ts b/packages/stage-ui/src/utils/queue.ts index 618885a0e..71ffd1c08 100644 --- a/packages/stage-ui/src/utils/queue.ts +++ b/packages/stage-ui/src/utils/queue.ts @@ -82,6 +82,7 @@ export function createQueue(options: { } } } + emit('drain') drainTask = undefined }