From e5a86e78ae219d7ef12f6a0cf6cb705c8df85949 Mon Sep 17 00:00:00 2001 From: Neko Ayaka Date: Sun, 11 Jan 2026 21:46:23 +0800 Subject: [PATCH] refactor(stage-ui): better chat store structure --- .../pages/devtools/performance-playground.vue | 16 +- apps/stage-pocket/src/pages/index.vue | 8 +- .../renderer/components/InteractiveArea.vue | 18 +- .../src/renderer/pages/index.vue | 8 +- .../pages/devtools/performance-playground.vue | 16 +- apps/stage-web/src/pages/index.vue | 8 +- .../components/Layouts/InteractiveArea.vue | 8 +- .../Layouts/MobileInteractiveArea.vue | 18 +- .../components/Widgets/ChatActionButtons.vue | 4 +- .../src/components/Widgets/ChatArea.vue | 12 +- .../components/context-flow-actions.vue | 34 +- .../src/pages/devtools/context-flow/index.vue | 9 +- .../pages/devtools/websocket-inspector.vue | 2 +- .../stage-ui/src/components/scenes/Stage.vue | 4 +- .../scenes/StageSubtitlesWithStrip.vue | 81 --- .../src/composables/use-data-maintenance.ts | 7 +- .../stage-ui/src/stores/character/index.ts | 1 + .../stage-ui/src/stores/character/notebook.ts | 150 +++++ .../stores/character/orchestrator/store.ts | 177 +++++- packages/stage-ui/src/stores/chat.ts | 516 ++++-------------- .../stage-ui/src/stores/chat/constants.ts | 4 + .../stage-ui/src/stores/chat/context-store.ts | 39 ++ .../stage-ui/src/stores/chat/data-store.ts | 144 +++++ packages/stage-ui/src/stores/chat/hooks.ts | 217 ++++++++ .../stage-ui/src/stores/chat/maintenance.ts | 24 + .../stage-ui/src/stores/chat/session-store.ts | 93 ++++ packages/stage-ui/src/stores/chat/state.ts | 2 + .../stage-ui/src/stores/chat/stream-store.ts | 52 ++ .../stage-ui/src/stores/markdown-stress.ts | 10 +- .../src/stores/mods/api/channel-gateway.ts | 123 +++++ .../src/stores/mods/api/context-bridge.ts | 101 ++-- .../stage-ui/src/stores/mods/api/events.ts | 54 ++ 32 files changed, 1349 insertions(+), 611 deletions(-) delete mode 100644 packages/stage-ui/src/components/scenes/StageSubtitlesWithStrip.vue create mode 100644 packages/stage-ui/src/stores/character/notebook.ts create mode 100644 packages/stage-ui/src/stores/chat/constants.ts create mode 100644 packages/stage-ui/src/stores/chat/context-store.ts create mode 100644 packages/stage-ui/src/stores/chat/data-store.ts create mode 100644 packages/stage-ui/src/stores/chat/hooks.ts create mode 100644 packages/stage-ui/src/stores/chat/maintenance.ts create mode 100644 packages/stage-ui/src/stores/chat/session-store.ts create mode 100644 packages/stage-ui/src/stores/chat/state.ts create mode 100644 packages/stage-ui/src/stores/chat/stream-store.ts create mode 100644 packages/stage-ui/src/stores/mods/api/channel-gateway.ts create mode 100644 packages/stage-ui/src/stores/mods/api/events.ts diff --git a/apps/stage-pocket/src/pages/devtools/performance-playground.vue b/apps/stage-pocket/src/pages/devtools/performance-playground.vue index 34d959c1c..68a58f46e 100644 --- a/apps/stage-pocket/src/pages/devtools/performance-playground.vue +++ b/apps/stage-pocket/src/pages/devtools/performance-playground.vue @@ -8,7 +8,9 @@ import { useDelayMessageQueue, useEmotionsMessageQueue } from '@proj-airi/stage- import { llmInferenceEndToken } from '@proj-airi/stage-ui/constants' import { EMOTION_EmotionMotionName_value, EMOTION_VRMExpressionName_value, EmotionThinkMotionName } from '@proj-airi/stage-ui/constants/emotions' import { useAudioContext, useSpeakingStore } from '@proj-airi/stage-ui/stores/audio' -import { useChatStore } from '@proj-airi/stage-ui/stores/chat' +import { useChatOrchestratorStore } from '@proj-airi/stage-ui/stores/chat' +import { useChatMaintenanceStore } from '@proj-airi/stage-ui/stores/chat/maintenance' +import { useChatSessionStore } from '@proj-airi/stage-ui/stores/chat/session-store' import { useConsciousnessStore } from '@proj-airi/stage-ui/stores/modules/consciousness' import { useSpeechStore } from '@proj-airi/stage-ui/stores/modules/speech' import { useProvidersStore } from '@proj-airi/stage-ui/stores/providers' @@ -62,9 +64,11 @@ const nowSpeaking = ref(false) const currentMotion = ref<{ group: string }>({ group: EmotionThinkMotionName }) const logLines = ref([]) const chatInput = ref('') -const chatStore = useChatStore() +const chatOrchestrator = useChatOrchestratorStore() +const chatSession = useChatSessionStore() +const chatMaintenance = useChatMaintenanceStore() const chatMessages = computed(() => { - return chatStore.messages + return chatSession.messages .filter(msg => msg.role !== 'system') .map((msg) => { const text = typeof msg.content === 'string' @@ -196,7 +200,7 @@ async function sendChat() { } try { - await chatStore.send(content, { + await chatOrchestrator.ingest(content, { model: activeChatModel.value, chatProvider: provider as ChatProvider, }) @@ -209,13 +213,13 @@ async function sendChat() { } function resetChat() { - chatStore.cleanupMessages() + chatMaintenance.cleanupMessages() chatInput.value = '' logLines.value = [] playbackManager.stopAll('reset') } -const { onBeforeMessageComposed, onBeforeSend, onTokenLiteral, onTokenSpecial, onStreamEnd, onAssistantResponseEnd } = chatStore +const { onBeforeMessageComposed, onBeforeSend, onTokenLiteral, onTokenSpecial, onStreamEnd, onAssistantResponseEnd } = chatOrchestrator const chatHookCleanups: Array<() => void> = [] let currentIntent: ReturnType | null = null diff --git a/apps/stage-pocket/src/pages/index.vue b/apps/stage-pocket/src/pages/index.vue index 66c2668d1..7a39a84a5 100644 --- a/apps/stage-pocket/src/pages/index.vue +++ b/apps/stage-pocket/src/pages/index.vue @@ -13,7 +13,7 @@ import { useBackgroundStore } from '@proj-airi/stage-layouts/stores/background' import { WidgetStage } from '@proj-airi/stage-ui/components/scenes' import { useAudioRecorder } from '@proj-airi/stage-ui/composables/audio/audio-recorder' import { useVAD } from '@proj-airi/stage-ui/stores/ai/models/vad' -import { useChatStore } from '@proj-airi/stage-ui/stores/chat' +import { useChatOrchestratorStore } from '@proj-airi/stage-ui/stores/chat' import { useLive2d } from '@proj-airi/stage-ui/stores/live2d' import { useConsciousnessStore } from '@proj-airi/stage-ui/stores/modules/consciousness' import { useHearingSpeechInputPipeline } from '@proj-airi/stage-ui/stores/modules/hearing' @@ -51,7 +51,7 @@ const { supportsStreamInput } = storeToRefs(hearingPipeline) const providersStore = useProvidersStore() const consciousnessStore = useConsciousnessStore() const { activeProvider: activeChatProvider, activeModel: activeChatModel } = storeToRefs(consciousnessStore) -const chatStore = useChatStore() +const chatStore = useChatOrchestratorStore() const shouldUseStreamInput = computed(() => supportsStreamInput.value && !!stream.value) @@ -85,7 +85,7 @@ async function startAudioInteraction() { if (!provider || !activeChatModel.value) return - await chatStore.send(text, { model: activeChatModel.value, chatProvider: provider as ChatProvider }) + await chatStore.ingest(text, { model: activeChatModel.value, chatProvider: provider as ChatProvider }) } catch (err) { console.error('Failed to send chat from voice:', err) @@ -112,7 +112,7 @@ async function handleSpeechStart() { if (!provider || !activeChatModel.value) return - await chatStore.send(finalText, { model: activeChatModel.value, chatProvider: provider as ChatProvider }) + await chatStore.ingest(finalText, { model: activeChatModel.value, chatProvider: provider as ChatProvider }) } catch (err) { console.error('Failed to send chat from voice:', err) diff --git a/apps/stage-tamagotchi/src/renderer/components/InteractiveArea.vue b/apps/stage-tamagotchi/src/renderer/components/InteractiveArea.vue index a6ab0ba6b..fa9e7f9ef 100644 --- a/apps/stage-tamagotchi/src/renderer/components/InteractiveArea.vue +++ b/apps/stage-tamagotchi/src/renderer/components/InteractiveArea.vue @@ -3,7 +3,10 @@ import type { ChatHistoryItem } from '@proj-airi/stage-ui/types/chat' import type { ChatProvider } from '@xsai-ext/providers/utils' import { ChatHistory } from '@proj-airi/stage-ui/components' -import { useChatStore } from '@proj-airi/stage-ui/stores/chat' +import { useChatOrchestratorStore } from '@proj-airi/stage-ui/stores/chat' +import { useChatMaintenanceStore } from '@proj-airi/stage-ui/stores/chat/maintenance' +import { useChatSessionStore } from '@proj-airi/stage-ui/stores/chat/session-store' +import { useChatStreamStore } from '@proj-airi/stage-ui/stores/chat/stream-store' import { useConsciousnessStore } from '@proj-airi/stage-ui/stores/modules/consciousness' import { useProvidersStore } from '@proj-airi/stage-ui/stores/providers' import { BasicTextarea } from '@proj-airi/ui' @@ -16,9 +19,14 @@ import { widgetsTools } from '../stores/tools/builtin/widgets' const messageInput = ref('') const attachments = ref<{ type: 'image', data: string, mimeType: string, url: string }[]>([]) -const chatStore = useChatStore() -const { send, onAfterMessageComposed, discoverToolsCompatibility, cleanupMessages } = chatStore -const { messages, sending, streamingMessage } = storeToRefs(chatStore) +const chatOrchestrator = useChatOrchestratorStore() +const chatSession = useChatSessionStore() +const chatStream = useChatStreamStore() +const { cleanupMessages } = useChatMaintenanceStore() +const { ingest, onAfterMessageComposed, discoverToolsCompatibility } = chatOrchestrator +const { messages } = storeToRefs(chatSession) +const { streamingMessage } = storeToRefs(chatStream) +const { sending } = storeToRefs(chatOrchestrator) const { t } = useI18n() const providersStore = useProvidersStore() const { activeModel, activeProvider } = storeToRefs(useConsciousnessStore()) @@ -42,7 +50,7 @@ async function handleSend() { try { const providerConfig = providersStore.getProviderConfig(activeProvider.value) - await send(textToSend, { + await ingest(textToSend, { model: activeModel.value, chatProvider: await providersStore.getProviderInstance(activeProvider.value), providerConfig, diff --git a/apps/stage-tamagotchi/src/renderer/pages/index.vue b/apps/stage-tamagotchi/src/renderer/pages/index.vue index 5e448b666..cef24b5c5 100644 --- a/apps/stage-tamagotchi/src/renderer/pages/index.vue +++ b/apps/stage-tamagotchi/src/renderer/pages/index.vue @@ -8,7 +8,7 @@ import { WidgetStage } from '@proj-airi/stage-ui/components/scenes' import { useAudioRecorder } from '@proj-airi/stage-ui/composables/audio/audio-recorder' import { useCanvasPixelIsTransparentAtPoint } from '@proj-airi/stage-ui/composables/canvas-alpha' import { useVAD } from '@proj-airi/stage-ui/stores/ai/models/vad' -import { useChatStore } from '@proj-airi/stage-ui/stores/chat' +import { useChatOrchestratorStore } from '@proj-airi/stage-ui/stores/chat' import { useLive2d } from '@proj-airi/stage-ui/stores/live2d' import { useConsciousnessStore } from '@proj-airi/stage-ui/stores/modules/consciousness' import { useHearingSpeechInputPipeline } from '@proj-airi/stage-ui/stores/modules/hearing' @@ -136,7 +136,7 @@ const { supportsStreamInput } = storeToRefs(hearingPipeline) const providersStore = useProvidersStore() const consciousnessStore = useConsciousnessStore() const { activeProvider: activeChatProvider, activeModel: activeChatModel } = storeToRefs(consciousnessStore) -const chatStore = useChatStore() +const chatStore = useChatOrchestratorStore() const shouldUseStreamInput = computed(() => supportsStreamInput.value && !!stream.value) const { @@ -179,7 +179,7 @@ async function handleSpeechStart() { if (!provider || !activeChatModel.value) return - await chatStore.send(finalText, { model: activeChatModel.value, chatProvider: provider as ChatProvider }) + await chatStore.ingest(finalText, { model: activeChatModel.value, chatProvider: provider as ChatProvider }) } catch (err) { console.error('Failed to send chat from voice:', err) @@ -228,7 +228,7 @@ async function startAudioInteraction() { if (!provider || !activeChatModel.value) return - await chatStore.send(text, { model: activeChatModel.value, chatProvider: provider as ChatProvider }) + await chatStore.ingest(text, { model: activeChatModel.value, chatProvider: provider as ChatProvider }) } catch (err) { console.error('Failed to send chat from voice:', err) diff --git a/apps/stage-web/src/pages/devtools/performance-playground.vue b/apps/stage-web/src/pages/devtools/performance-playground.vue index 34d959c1c..68a58f46e 100644 --- a/apps/stage-web/src/pages/devtools/performance-playground.vue +++ b/apps/stage-web/src/pages/devtools/performance-playground.vue @@ -8,7 +8,9 @@ import { useDelayMessageQueue, useEmotionsMessageQueue } from '@proj-airi/stage- import { llmInferenceEndToken } from '@proj-airi/stage-ui/constants' import { EMOTION_EmotionMotionName_value, EMOTION_VRMExpressionName_value, EmotionThinkMotionName } from '@proj-airi/stage-ui/constants/emotions' import { useAudioContext, useSpeakingStore } from '@proj-airi/stage-ui/stores/audio' -import { useChatStore } from '@proj-airi/stage-ui/stores/chat' +import { useChatOrchestratorStore } from '@proj-airi/stage-ui/stores/chat' +import { useChatMaintenanceStore } from '@proj-airi/stage-ui/stores/chat/maintenance' +import { useChatSessionStore } from '@proj-airi/stage-ui/stores/chat/session-store' import { useConsciousnessStore } from '@proj-airi/stage-ui/stores/modules/consciousness' import { useSpeechStore } from '@proj-airi/stage-ui/stores/modules/speech' import { useProvidersStore } from '@proj-airi/stage-ui/stores/providers' @@ -62,9 +64,11 @@ const nowSpeaking = ref(false) const currentMotion = ref<{ group: string }>({ group: EmotionThinkMotionName }) const logLines = ref([]) const chatInput = ref('') -const chatStore = useChatStore() +const chatOrchestrator = useChatOrchestratorStore() +const chatSession = useChatSessionStore() +const chatMaintenance = useChatMaintenanceStore() const chatMessages = computed(() => { - return chatStore.messages + return chatSession.messages .filter(msg => msg.role !== 'system') .map((msg) => { const text = typeof msg.content === 'string' @@ -196,7 +200,7 @@ async function sendChat() { } try { - await chatStore.send(content, { + await chatOrchestrator.ingest(content, { model: activeChatModel.value, chatProvider: provider as ChatProvider, }) @@ -209,13 +213,13 @@ async function sendChat() { } function resetChat() { - chatStore.cleanupMessages() + chatMaintenance.cleanupMessages() chatInput.value = '' logLines.value = [] playbackManager.stopAll('reset') } -const { onBeforeMessageComposed, onBeforeSend, onTokenLiteral, onTokenSpecial, onStreamEnd, onAssistantResponseEnd } = chatStore +const { onBeforeMessageComposed, onBeforeSend, onTokenLiteral, onTokenSpecial, onStreamEnd, onAssistantResponseEnd } = chatOrchestrator const chatHookCleanups: Array<() => void> = [] let currentIntent: ReturnType | null = null diff --git a/apps/stage-web/src/pages/index.vue b/apps/stage-web/src/pages/index.vue index 739d697a0..df94b37fe 100644 --- a/apps/stage-web/src/pages/index.vue +++ b/apps/stage-web/src/pages/index.vue @@ -14,7 +14,7 @@ import { WidgetStage } from '@proj-airi/stage-ui/components/scenes' import { useAudioRecorder } from '@proj-airi/stage-ui/composables/audio/audio-recorder' import { fetchSession } from '@proj-airi/stage-ui/libs/auth' import { useVAD } from '@proj-airi/stage-ui/stores/ai/models/vad' -import { useChatStore } from '@proj-airi/stage-ui/stores/chat' +import { useChatOrchestratorStore } from '@proj-airi/stage-ui/stores/chat' import { useLive2d } from '@proj-airi/stage-ui/stores/live2d' import { useConsciousnessStore } from '@proj-airi/stage-ui/stores/modules/consciousness' import { useHearingSpeechInputPipeline } from '@proj-airi/stage-ui/stores/modules/hearing' @@ -52,7 +52,7 @@ const { supportsStreamInput } = storeToRefs(hearingPipeline) const providersStore = useProvidersStore() const consciousnessStore = useConsciousnessStore() const { activeProvider: activeChatProvider, activeModel: activeChatModel } = storeToRefs(consciousnessStore) -const chatStore = useChatStore() +const chatStore = useChatOrchestratorStore() const shouldUseStreamInput = computed(() => supportsStreamInput.value && !!stream.value) @@ -86,7 +86,7 @@ async function startAudioInteraction() { if (!provider || !activeChatModel.value) return - await chatStore.send(text, { model: activeChatModel.value, chatProvider: provider as ChatProvider }) + await chatStore.ingest(text, { model: activeChatModel.value, chatProvider: provider as ChatProvider }) } catch (err) { console.error('Failed to send chat from voice:', err) @@ -113,7 +113,7 @@ async function handleSpeechStart() { if (!provider || !activeChatModel.value) return - await chatStore.send(finalText, { model: activeChatModel.value, chatProvider: provider as ChatProvider }) + await chatStore.ingest(finalText, { model: activeChatModel.value, chatProvider: provider as ChatProvider }) } catch (err) { console.error('Failed to send chat from voice:', err) diff --git a/packages/stage-layouts/src/components/Layouts/InteractiveArea.vue b/packages/stage-layouts/src/components/Layouts/InteractiveArea.vue index d3d585200..08a322abb 100644 --- a/packages/stage-layouts/src/components/Layouts/InteractiveArea.vue +++ b/packages/stage-layouts/src/components/Layouts/InteractiveArea.vue @@ -2,7 +2,9 @@ import type { ChatHistoryItem } from '@proj-airi/stage-ui/types/chat' import { ChatHistory } from '@proj-airi/stage-ui/components' -import { useChatStore } from '@proj-airi/stage-ui/stores/chat' +import { useChatOrchestratorStore } from '@proj-airi/stage-ui/stores/chat' +import { useChatSessionStore } from '@proj-airi/stage-ui/stores/chat/session-store' +import { useChatStreamStore } from '@proj-airi/stage-ui/stores/chat/stream-store' import { useDeferredMount } from '@proj-airi/ui' import { storeToRefs } from 'pinia' import { computed, ref } from 'vue' @@ -12,7 +14,9 @@ import ChatArea from '../Widgets/ChatArea.vue' import ChatContainer from '../Widgets/ChatContainer.vue' const { isReady } = useDeferredMount() -const { messages, sending, streamingMessage } = storeToRefs(useChatStore()) +const { sending } = storeToRefs(useChatOrchestratorStore()) +const { messages } = storeToRefs(useChatSessionStore()) +const { streamingMessage } = storeToRefs(useChatStreamStore()) const isLoading = ref(true) const historyMessages = computed(() => messages.value as unknown as ChatHistoryItem[]) diff --git a/packages/stage-layouts/src/components/Layouts/MobileInteractiveArea.vue b/packages/stage-layouts/src/components/Layouts/MobileInteractiveArea.vue index 39ac71ac0..27e630002 100644 --- a/packages/stage-layouts/src/components/Layouts/MobileInteractiveArea.vue +++ b/packages/stage-layouts/src/components/Layouts/MobileInteractiveArea.vue @@ -5,7 +5,10 @@ import type { ChatProvider } from '@xsai-ext/providers/utils' import { ChatHistory, HearingConfigDialog } from '@proj-airi/stage-ui/components' import { useAudioAnalyzer } from '@proj-airi/stage-ui/composables' import { useAudioContext } from '@proj-airi/stage-ui/stores/audio' -import { useChatStore } from '@proj-airi/stage-ui/stores/chat' +import { useChatOrchestratorStore } from '@proj-airi/stage-ui/stores/chat' +import { useChatMaintenanceStore } from '@proj-airi/stage-ui/stores/chat/maintenance' +import { useChatSessionStore } from '@proj-airi/stage-ui/stores/chat/session-store' +import { useChatStreamStore } from '@proj-airi/stage-ui/stores/chat/stream-store' import { useConsciousnessStore } from '@proj-airi/stage-ui/stores/modules/consciousness' import { useProvidersStore } from '@proj-airi/stage-ui/stores/providers' import { useSettings, useSettingsAudioDevice } from '@proj-airi/stage-ui/stores/settings' @@ -25,8 +28,13 @@ import { BackgroundDialogPicker } from '../Backgrounds' const { isDark, toggleDark } = useTheme() const hearingDialogOpen = ref(false) -const chatStore = useChatStore() -const { messages, sending, streamingMessage } = storeToRefs(chatStore) +const chatOrchestrator = useChatOrchestratorStore() +const chatSession = useChatSessionStore() +const chatStream = useChatStreamStore() +const { cleanupMessages } = useChatMaintenanceStore() +const { messages } = storeToRefs(chatSession) +const { streamingMessage } = storeToRefs(chatStream) +const { sending } = storeToRefs(chatOrchestrator) const historyMessages = computed(() => messages.value as unknown as ChatHistoryItem[]) const viewControlsActiveMode = ref<'x' | 'y' | 'z' | 'scale'>('scale') @@ -44,7 +52,7 @@ useResizeObserver(document.documentElement, () => screenSafeArea.update()) const { themeColorsHueDynamic, stageViewControlsEnabled } = storeToRefs(useSettings()) const settingsAudioDevice = useSettingsAudioDevice() const { enabled, selectedAudioInput, stream, audioInputs } = storeToRefs(settingsAudioDevice) -const { send, onAfterMessageComposed, discoverToolsCompatibility, cleanupMessages } = chatStore +const { ingest, onAfterMessageComposed, discoverToolsCompatibility } = chatOrchestrator const { t } = useI18n() const { audioContext } = useAudioContext() const { startAnalyzer, stopAnalyzer, volumeLevel } = useAudioAnalyzer() @@ -71,7 +79,7 @@ async function handleSend() { try { const providerConfig = providersStore.getProviderConfig(activeProvider.value) - await send(textToSend, { + await ingest(textToSend, { chatProvider: await providersStore.getProviderInstance(activeProvider.value) as ChatProvider, model: activeModel.value, providerConfig, diff --git a/packages/stage-layouts/src/components/Widgets/ChatActionButtons.vue b/packages/stage-layouts/src/components/Widgets/ChatActionButtons.vue index b35962681..314af9f99 100644 --- a/packages/stage-layouts/src/components/Widgets/ChatActionButtons.vue +++ b/packages/stage-layouts/src/components/Widgets/ChatActionButtons.vue @@ -1,11 +1,11 @@ - - - - diff --git a/packages/stage-ui/src/composables/use-data-maintenance.ts b/packages/stage-ui/src/composables/use-data-maintenance.ts index ae07e814d..128612745 100644 --- a/packages/stage-ui/src/composables/use-data-maintenance.ts +++ b/packages/stage-ui/src/composables/use-data-maintenance.ts @@ -3,7 +3,8 @@ import type { ChatHistoryItem } from '../types/chat' import { isStageTamagotchi } from '@proj-airi/stage-shared' import { useLive2d } from '@proj-airi/stage-ui-live2d' -import { useChatStore } from '../stores/chat' +import { useChatOrchestratorStore } from '../stores/chat' +import { useChatSessionStore } from '../stores/chat/session-store' import { useDisplayModelsStore } from '../stores/display-models' import { useMcpStore } from '../stores/mcp' import { useAiriCardStore } from '../stores/modules/airi-card' @@ -19,7 +20,8 @@ import { useProvidersStore } from '../stores/providers' import { useSettings, useSettingsAudioDevice } from '../stores/settings' export function useDataMaintenance() { - const chatStore = useChatStore() + const chatStore = useChatSessionStore() + const chatOrchestrator = useChatOrchestratorStore() const displayModelsStore = useDisplayModelsStore() const providersStore = useProvidersStore() const settingsStore = useSettings() @@ -57,6 +59,7 @@ export function useDataMaintenance() { } function deleteAllChatSessions() { + chatOrchestrator.cancelPendingSends() chatStore.resetAllSessions() } diff --git a/packages/stage-ui/src/stores/character/index.ts b/packages/stage-ui/src/stores/character/index.ts index 15abedf59..b426ad539 100644 --- a/packages/stage-ui/src/stores/character/index.ts +++ b/packages/stage-ui/src/stores/character/index.ts @@ -8,6 +8,7 @@ import { useLlmmarkerParser } from '../../composables/llm-marker-parser' import { useAiriCardStore } from '../modules' import { useSpeechRuntimeStore } from '../speech-runtime' +export * from './notebook' export * from './orchestrator' export interface CharacterSparkNotifyReaction { diff --git a/packages/stage-ui/src/stores/character/notebook.ts b/packages/stage-ui/src/stores/character/notebook.ts new file mode 100644 index 000000000..d9885475e --- /dev/null +++ b/packages/stage-ui/src/stores/character/notebook.ts @@ -0,0 +1,150 @@ +import { nanoid } from 'nanoid' +import { defineStore } from 'pinia' +import { computed, ref } from 'vue' + +export type NotebookEntryKind = 'note' | 'diary' | 'focus' + +export interface NotebookEntry { + id: string + kind: NotebookEntryKind + text: string + createdAt: number + tags?: string[] + metadata?: Record +} + +export type TaskPriority = 'low' | 'normal' | 'high' | 'critical' +export type TaskStatus = 'queued' | 'scheduled' | 'done' | 'dropped' + +export interface ScheduledTask { + id: string + title: string + details?: string + priority: TaskPriority + status: TaskStatus + dueAt?: number + createdAt: number + updatedAt: number + lastNotifiedAt?: number + nextNotifyAt?: number + metadata?: Record +} + +export const useCharacterNotebookStore = defineStore('character-notebook', () => { + const entries = ref([]) + const tasks = ref([]) + + const partitionDiary = computed(() => entries.value.filter(entry => entry.kind === 'diary')) + const partitionFocus = computed(() => entries.value.filter(entry => entry.kind === 'focus')) + + function addEntry(kind: NotebookEntryKind, text: string, options?: { tags?: string[], metadata?: Record }) { + const entry: NotebookEntry = { + id: nanoid(), + kind, + text, + createdAt: Date.now(), + tags: options?.tags, + metadata: options?.metadata, + } + + entries.value.push(entry) + return entry + } + + function addNote(text: string, options?: { tags?: string[], metadata?: Record }) { + return addEntry('note', text, options) + } + + function addDiaryEntry(text: string, options?: { tags?: string[], metadata?: Record }) { + return addEntry('diary', text, options) + } + + function addFocusEntry(text: string, options?: { tags?: string[], metadata?: Record }) { + return addEntry('focus', text, options) + } + + function scheduleTask(payload: { + title: string + details?: string + priority?: TaskPriority + dueAt?: number + metadata?: Record + }) { + const now = Date.now() + const task: ScheduledTask = { + id: nanoid(), + title: payload.title, + details: payload.details, + priority: payload.priority ?? 'normal', + status: payload.dueAt ? 'scheduled' : 'queued', + dueAt: payload.dueAt, + createdAt: now, + updatedAt: now, + metadata: payload.metadata, + } + + tasks.value.push(task) + return task + } + + function markTaskDone(taskId: string) { + const task = tasks.value.find(item => item.id === taskId) + if (!task) + return + + task.status = 'done' + task.updatedAt = Date.now() + } + + function requeueTask(taskId: string, options?: { dueAt?: number, reason?: string }) { + const task = tasks.value.find(item => item.id === taskId) + if (!task) + return + + task.status = 'queued' + task.dueAt = options?.dueAt + task.updatedAt = Date.now() + task.metadata = { + ...task.metadata, + requeueReason: options?.reason, + } + } + + function markTaskNotified(taskId: string, nextNotifyAt?: number) { + const task = tasks.value.find(item => item.id === taskId) + if (!task) + return + + task.lastNotifiedAt = Date.now() + task.nextNotifyAt = nextNotifyAt + task.updatedAt = Date.now() + } + + function getDueTasks(now: number, windowMs: number) { + return tasks.value.filter((task) => { + if (task.status === 'done' || task.status === 'dropped') + return false + const dueAt = task.dueAt ?? now + if (dueAt > now + windowMs) + return false + if (typeof task.nextNotifyAt === 'number' && task.nextNotifyAt > now) + return false + return true + }) + } + + return { + entries, + tasks, + partitionDiary, + partitionFocus, + addNote, + addDiaryEntry, + addFocusEntry, + scheduleTask, + markTaskDone, + requeueTask, + markTaskNotified, + getDueTasks, + } +}) diff --git a/packages/stage-ui/src/stores/character/orchestrator/store.ts b/packages/stage-ui/src/stores/character/orchestrator/store.ts index 8cbc8d83c..5106b18cb 100644 --- a/packages/stage-ui/src/stores/character/orchestrator/store.ts +++ b/packages/stage-ui/src/stores/character/orchestrator/store.ts @@ -3,7 +3,7 @@ import type { WebSocketBaseEvent, WebSocketEvents } from '@proj-airi/server-sdk' import { defineStore, storeToRefs } from 'pinia' import { ref } from 'vue' -import { useCharacterStore } from '../' +import { useCharacterNotebookStore, useCharacterStore } from '../' import { useLLM } from '../../llm' import { useModsServerChannelStore } from '../../mods/api/channel-server' import { useConsciousnessStore } from '../../modules/consciousness' @@ -17,11 +17,27 @@ export const useCharacterOrchestratorStore = defineStore('character-orchestrator const { activeProvider, activeModel } = storeToRefs(useConsciousnessStore()) const providersStore = useProvidersStore() const characterStore = useCharacterStore() + const notebookStore = useCharacterNotebookStore() const { systemPrompt } = storeToRefs(characterStore) const modsServerChannelStore = useModsServerChannelStore() const processing = ref(false) const pendingNotifies = ref>>([]) + const scheduledNotifies = ref + enqueuedAt: number + nextRunAt: number + attempts: number + maxAttempts: number + reason?: string + }>>([]) + const attentionConfig = ref({ + tickIntervalMs: 2_000, + taskNotifyWindowMs: 60_000, + requeueDelayMs: 30_000, + maxAttempts: 3, + }) + let tickTimer: ReturnType | undefined const sparkNotifyAgent = setupAgentSparkNotifyHandler({ stream, getActiveProvider: () => activeProvider.value, @@ -36,6 +52,146 @@ export const useCharacterOrchestratorStore = defineStore('character-orchestrator setPending: next => pendingNotifies.value = next, }) + function computeNextRunAt(event: WebSocketBaseEvent<'spark:notify', WebSocketEvents['spark:notify']>, attempts: number) { + const now = Date.now() + const baseDelay = (() => { + switch (event.data.urgency) { + case 'immediate': + return 0 + case 'soon': + return 10_000 + case 'later': + return 60_000 + default: + return 30_000 + } + })() + + return now + baseDelay + (attempts * attentionConfig.value.requeueDelayMs) + } + + function removePending(eventId: string) { + pendingNotifies.value = pendingNotifies.value.filter(item => item.data.id !== eventId) + } + + function enqueueSparkNotify(event: WebSocketBaseEvent<'spark:notify', WebSocketEvents['spark:notify']>, options?: { reason?: string, nextRunAt?: number, maxAttempts?: number }) { + if (!pendingNotifies.value.find(item => item.data.id === event.data.id)) { + pendingNotifies.value = [...pendingNotifies.value, event] + } + + scheduledNotifies.value = [...scheduledNotifies.value, { + event, + enqueuedAt: Date.now(), + nextRunAt: options?.nextRunAt ?? computeNextRunAt(event, 0), + attempts: 0, + maxAttempts: options?.maxAttempts ?? attentionConfig.value.maxAttempts, + reason: options?.reason, + }] + } + + async function processSparkNotify(event: WebSocketBaseEvent<'spark:notify', WebSocketEvents['spark:notify']>) { + const result = await sparkNotifyAgent.handle(event) + if (!result?.commands?.length) + return result + + for (const command of result.commands) { + modsServerChannelStore.send({ + type: 'spark:command', + data: command, + }) + } + + return result + } + + async function handleIncomingSparkNotify(event: WebSocketBaseEvent<'spark:notify', WebSocketEvents['spark:notify']>) { + if (event.data.urgency === 'immediate' && !processing.value) { + return await processSparkNotify(event) + } + + enqueueSparkNotify(event, { reason: 'spark:notify' }) + return undefined + } + + function enqueueDueTasks(now: number) { + const dueTasks = notebookStore.getDueTasks(now, attentionConfig.value.taskNotifyWindowMs) + if (!dueTasks.length) + return + + for (const task of dueTasks) { + const event: WebSocketBaseEvent<'spark:notify', WebSocketEvents['spark:notify']> = { + type: 'spark:notify', + source: 'character:task-scheduler', + data: { + id: `task-${task.id}`, + eventId: task.id, + kind: 'reminder', + urgency: task.priority === 'critical' ? 'immediate' : 'soon', + headline: `Task reminder: ${task.title}`, + note: task.details, + destinations: ['character'], + payload: { + taskId: task.id, + dueAt: task.dueAt, + priority: task.priority, + }, + }, + } + + enqueueSparkNotify(event, { reason: 'task:due' }) + notebookStore.markTaskNotified(task.id, now + attentionConfig.value.requeueDelayMs) + } + } + + async function tick() { + if (processing.value) + return + + const now = Date.now() + enqueueDueTasks(now) + + const nextIndex = scheduledNotifies.value.findIndex(item => item.nextRunAt <= now) + if (nextIndex < 0) + return + + const [next] = scheduledNotifies.value.splice(nextIndex, 1) + removePending(next.event.data.id) + + try { + await processSparkNotify(next.event) + } + catch (error) { + if (next.attempts + 1 < next.maxAttempts) { + scheduledNotifies.value = [...scheduledNotifies.value, { + ...next, + attempts: next.attempts + 1, + nextRunAt: computeNextRunAt(next.event, next.attempts + 1), + }] + pendingNotifies.value = [...pendingNotifies.value, next.event] + } + else { + console.warn('Dropped spark:notify after max attempts:', error) + } + } + } + + function startTicker() { + if (tickTimer) + return + + tickTimer = setInterval(() => { + void tick() + }, attentionConfig.value.tickIntervalMs) + } + + function stopTicker() { + if (!tickTimer) + return + + clearInterval(tickTimer) + tickTimer = undefined + } + async function handleSparkEmit(_: WebSocketBaseEvent<'spark:emit', WebSocketEvents['spark:emit']>) { // Currently no-op return undefined @@ -44,16 +200,7 @@ export const useCharacterOrchestratorStore = defineStore('character-orchestrator function initialize() { modsServerChannelStore.onEvent('spark:notify', async (event) => { try { - const result = await sparkNotifyAgent.handle(event) - if (!result?.commands?.length) - return - - for (const command of result.commands) { - modsServerChannelStore.send({ - type: 'spark:command', - data: command, - }) - } + await handleIncomingSparkNotify(event) } catch (error) { console.warn('Failed to handle spark:notify event:', error) @@ -68,15 +215,21 @@ export const useCharacterOrchestratorStore = defineStore('character-orchestrator console.warn('Failed to handle spark:emit event:', error) } }) + + startTicker() } return { processing, pendingNotifies, + scheduledNotifies, + attentionConfig, initialize, + startTicker, + stopTicker, - handleSparkNotify: sparkNotifyAgent.handle, + handleSparkNotify: handleIncomingSparkNotify, handleSparkEmit, } }) diff --git a/packages/stage-ui/src/stores/chat.ts b/packages/stage-ui/src/stores/chat.ts index b354bc4bc..f232ffcd9 100644 --- a/packages/stage-ui/src/stores/chat.ts +++ b/packages/stage-ui/src/stores/chat.ts @@ -1,67 +1,60 @@ import type { WebSocketEventInputs } from '@proj-airi/server-sdk' import type { ChatProvider } from '@xsai-ext/providers/utils' -import type { CommonContentPart, Message, SystemMessage, ToolMessage } from '@xsai/shared-chat' +import type { CommonContentPart, Message, ToolMessage } from '@xsai/shared-chat' -import type { StreamEvent, StreamOptions } from '../stores/llm' -import type { ChatAssistantMessage, ChatHistoryItem, ChatSlices, ChatStreamEventContext, ContextMessage, StreamingAssistantMessage } from '../types/chat' +import type { ChatAssistantMessage, ChatSlices, ChatStreamEventContext, StreamingAssistantMessage } from '../types/chat' +import type { StreamEvent, StreamOptions } from './llm' -import { ContextUpdateStrategy } from '@proj-airi/server-sdk' import { createQueue } from '@proj-airi/stream-kit' -import { useLocalStorage } from '@vueuse/core' import { defineStore, storeToRefs } from 'pinia' -import { computed, ref, toRaw, watch } from 'vue' +import { ref, toRaw } from 'vue' import { useAnalytics } from '../composables' import { useLlmmarkerParser } from '../composables/llm-marker-parser' import { categorizeResponse, createStreamingCategorizer } from '../composables/response-categoriser' -import { useLLM } from '../stores/llm' -import { getEventSourceKey } from '../utils/event-source' -import { useCharacterStore } from './character' +import { useChatContextStore } from './chat/context-store' +import { createChatHooks } from './chat/hooks' +import { useChatSessionStore } from './chat/session-store' +import { useChatStreamStore } from './chat/stream-store' +import { useLLM } from './llm' import { useConsciousnessStore } from './modules/consciousness' -const CHAT_STORAGE_KEY = 'chat/messages/v2' -const ACTIVE_SESSION_STORAGE_KEY = 'chat/active-session' -export const CONTEXT_CHANNEL_NAME = 'airi-context-update' -export const CHAT_STREAM_CHANNEL_NAME = 'airi-chat-stream' +interface SendOptions { + model: string + chatProvider: ChatProvider + providerConfig?: Record + attachments?: { type: 'image', data: string, mimeType: string }[] + tools?: StreamOptions['tools'] + input?: WebSocketEventInputs +} -export const useChatStore = defineStore('chat', () => { +interface QueuedSend { + sendingMessage: string + options: SendOptions + generation: number + sessionId: string + cancelled?: boolean + deferred: { + resolve: () => void + reject: (error: unknown) => void + } +} + +export const useChatOrchestratorStore = defineStore('chat-orchestrator', () => { const llmStore = useLLM() const consciousnessStore = useConsciousnessStore() const { activeProvider } = storeToRefs(consciousnessStore) - const { systemPrompt } = storeToRefs(useCharacterStore()) const { trackFirstMessage } = useAnalytics() - const activeSessionId = useLocalStorage(ACTIVE_SESSION_STORAGE_KEY, 'default') - const sessionMessages = useLocalStorage>(CHAT_STORAGE_KEY, {}) + const chatSession = useChatSessionStore() + const chatStream = useChatStreamStore() + const chatContext = useChatContextStore() + const { activeSessionId } = storeToRefs(chatSession) + const { streamingMessage } = storeToRefs(chatStream) const sending = ref(false) - const streamingMessage = ref({ role: 'assistant', content: '', slices: [], tool_results: [], createdAt: Date.now() }) - const sessionGenerations = ref>({}) - - const activeContexts = ref>({}) - - interface SendOptions { - model: string - chatProvider: ChatProvider - providerConfig?: Record - attachments?: { type: 'image', data: string, mimeType: string }[] - tools?: StreamOptions['tools'] - input?: WebSocketEventInputs - } - - interface QueuedSend { - sendingMessage: string - options: SendOptions - generation: number - sessionId: string - cancelled?: boolean - deferred: { - resolve: () => void - reject: (error: unknown) => void - } - } - const pendingQueuedSends = ref([]) + const hooks = createChatHooks() const sendQueue = createQueue({ handlers: [ @@ -71,7 +64,7 @@ export const useChatStore = defineStore('chat', () => { if (cancelled) return - if (getSessionGeneration(sessionId) !== generation) { + if (chatSession.getSessionGeneration(sessionId) !== generation) { deferred.reject(new Error('Chat session was reset before send could start')) return } @@ -95,259 +88,6 @@ export const useChatStore = defineStore('chat', () => { pendingQueuedSends.value = pendingQueuedSends.value.filter(item => item !== queuedSend) }) - // ----- Hooks (UI callbacks) ----- - const onBeforeMessageComposedHooks = ref) => Promise>>([]) - const onAfterMessageComposedHooks = ref Promise>>([]) - const onBeforeSendHooks = ref Promise>>([]) - const onAfterSendHooks = ref Promise>>([]) - const onTokenLiteralHooks = ref Promise>>([]) - const onTokenSpecialHooks = ref Promise>>([]) - const onStreamEndHooks = ref Promise>>([]) - const onAssistantResponseEndHooks = ref Promise>>([]) - const onAssistantMessageHooks = ref Promise>>([]) - const onChatTurnCompleteHooks = ref Promise>>([]) - - function onBeforeMessageComposed(cb: (message: string, context: Omit) => Promise) { - onBeforeMessageComposedHooks.value.push(cb) - return () => onBeforeMessageComposedHooks.value = onBeforeMessageComposedHooks.value.filter(hook => hook !== cb) // return remove listener callback - } - - function onAfterMessageComposed(cb: (message: string, context: ChatStreamEventContext) => Promise) { - onAfterMessageComposedHooks.value.push(cb) - return () => onAfterMessageComposedHooks.value = onAfterMessageComposedHooks.value.filter(hook => hook !== cb) // return remove listener callback - } - - function onBeforeSend(cb: (message: string, context: ChatStreamEventContext) => Promise) { - onBeforeSendHooks.value.push(cb) - return () => onBeforeSendHooks.value = onBeforeSendHooks.value.filter(hook => hook !== cb) // return remove listener callback - } - - function onAfterSend(cb: (message: string, context: ChatStreamEventContext) => Promise) { - onAfterSendHooks.value.push(cb) - return () => onAfterSendHooks.value = onAfterSendHooks.value.filter(hook => hook !== cb) // return remove listener callback - } - - function onTokenLiteral(cb: (literal: string, context: ChatStreamEventContext) => Promise) { - onTokenLiteralHooks.value.push(cb) - return () => onTokenLiteralHooks.value = onTokenLiteralHooks.value.filter(hook => hook !== cb) // return remove listener callback - } - - function onTokenSpecial(cb: (special: string, context: ChatStreamEventContext) => Promise) { - onTokenSpecialHooks.value.push(cb) - return () => onTokenSpecialHooks.value = onTokenSpecialHooks.value.filter(hook => hook !== cb) // return remove listener callback - } - - function onStreamEnd(cb: (context: ChatStreamEventContext) => Promise) { - onStreamEndHooks.value.push(cb) - return () => onStreamEndHooks.value = onStreamEndHooks.value.filter(hook => hook !== cb) // return remove listener callback - } - - function onAssistantResponseEnd(cb: (message: string, context: ChatStreamEventContext) => Promise) { - onAssistantResponseEndHooks.value.push(cb) - return () => onAssistantResponseEndHooks.value = onAssistantResponseEndHooks.value.filter(hook => hook !== cb) // return remove listener callback - } - - function onAssistantMessage(cb: (message: StreamingAssistantMessage, messageText: string, context: ChatStreamEventContext) => Promise) { - onAssistantMessageHooks.value.push(cb) - return () => onAssistantMessageHooks.value = onAssistantMessageHooks.value.filter(hook => hook !== cb) // return remove listener callback - } - - function onChatTurnComplete(cb: (chat: { output: StreamingAssistantMessage, outputText: string, toolCalls: ToolMessage[] }, context: ChatStreamEventContext) => Promise) { - onChatTurnCompleteHooks.value.push(cb) - return () => onChatTurnCompleteHooks.value = onChatTurnCompleteHooks.value.filter(hook => hook !== cb) // return remove listener callback - } - - function clearHooks() { - onBeforeMessageComposedHooks.value = [] - onAfterMessageComposedHooks.value = [] - onBeforeSendHooks.value = [] - onAfterSendHooks.value = [] - onTokenLiteralHooks.value = [] - onTokenSpecialHooks.value = [] - onStreamEndHooks.value = [] - onAssistantResponseEndHooks.value = [] - onAssistantMessageHooks.value = [] - onChatTurnCompleteHooks.value = [] - } - - async function emitBeforeMessageComposedHooks(message: string, context: Omit) { - for (const hook of onBeforeMessageComposedHooks.value) - await hook(message, context) - } - - async function emitAfterMessageComposedHooks(message: string, context: ChatStreamEventContext) { - for (const hook of onAfterMessageComposedHooks.value) - await hook(message, context) - } - - async function emitBeforeSendHooks(message: string, context: ChatStreamEventContext) { - for (const hook of onBeforeSendHooks.value) - await hook(message, context) - } - - async function emitAfterSendHooks(message: string, context: ChatStreamEventContext) { - for (const hook of onAfterSendHooks.value) - await hook(message, context) - } - - async function emitTokenLiteralHooks(literal: string, context: ChatStreamEventContext) { - for (const hook of onTokenLiteralHooks.value) - await hook(literal, context) - } - - async function emitTokenSpecialHooks(special: string, context: ChatStreamEventContext) { - for (const hook of onTokenSpecialHooks.value) - await hook(special, context) - } - - async function emitStreamEndHooks(context: ChatStreamEventContext) { - for (const hook of onStreamEndHooks.value) - await hook(context) - } - - async function emitAssistantResponseEndHooks(message: string, context: ChatStreamEventContext) { - for (const hook of onAssistantResponseEndHooks.value) - await hook(message, context) - } - - async function emitAssistantMessageHooks(message: StreamingAssistantMessage, messageText: string, context: ChatStreamEventContext) { - for (const hook of onAssistantMessageHooks.value) - await hook(message, messageText, context) - } - - async function emitChatTurnCompleteHooks(chat: { output: StreamingAssistantMessage, outputText: string, toolCalls: ToolMessage[] }, context: ChatStreamEventContext) { - for (const hook of onChatTurnCompleteHooks.value) - await hook(chat, context) - } - - // ----- 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' - const mathSyntaxSystemPrompt = '- For any math equation, use LaTeX format, eg: $ x^3 $, always escape dollar sign outside math equation\n' - - function ensureSessionGeneration(sessionId: string) { - if (sessionGenerations.value[sessionId] === undefined) - sessionGenerations.value = { ...sessionGenerations.value, [sessionId]: 0 } - } - - function getSessionGeneration(sessionId: string) { - ensureSessionGeneration(sessionId) - return sessionGenerations.value[sessionId] ?? 0 - } - - function bumpSessionGeneration(sessionId: string) { - const nextGeneration = getSessionGeneration(sessionId) + 1 - sessionGenerations.value = { ...sessionGenerations.value, [sessionId]: nextGeneration } - return nextGeneration - } - - function getSessionGenerationValue(sessionId = activeSessionId.value) { - return getSessionGeneration(sessionId) - } - - function generateInitialMessage() { - // TODO: compose, replace {{ user }} tag, etc - const content = codeBlockSystemPrompt + mathSyntaxSystemPrompt + systemPrompt.value - - return { - role: 'system', - content, - } satisfies SystemMessage - } - - function ensureSession(sessionId: string) { - ensureSessionGeneration(sessionId) - - if (!sessionMessages.value[sessionId] || sessionMessages.value[sessionId].length === 0) { - sessionMessages.value[sessionId] = [generateInitialMessage()] - } - } - - ensureSession(activeSessionId.value) - - function getSessionMessagesById(sessionId: string) { - ensureSession(sessionId) - return sessionMessages.value[sessionId]! - } - - const messages = computed({ - get: () => { - ensureSession(activeSessionId.value) - return sessionMessages.value[activeSessionId.value] - }, - set: (value) => { - sessionMessages.value[activeSessionId.value] = value - }, - }) - - function setActiveSession(sessionId: string) { - activeSessionId.value = sessionId - ensureSession(sessionId) - } - - function cleanupMessages(sessionId = activeSessionId.value) { - bumpSessionGeneration(sessionId) - sessionMessages.value[sessionId] = [generateInitialMessage()] - activeContexts.value = {} - - // Reject pending sends for this session so callers don't hang after cleanup - for (const queued of pendingQueuedSends.value) { - if (queued.sessionId !== sessionId) - continue - - queued.cancelled = true - queued.deferred.reject(new Error('Chat session was reset before send could start')) - } - - pendingQueuedSends.value = pendingQueuedSends.value.filter(item => item.sessionId !== sessionId) - sending.value = false - streamingMessage.value = { role: 'assistant', content: '', slices: [], tool_results: [] } - } - - function getAllSessions() { - return JSON.parse(JSON.stringify(toRaw(sessionMessages.value))) as Record - } - - function replaceSessions(sessions: Record) { - sessionMessages.value = sessions - sessionGenerations.value = Object.fromEntries(Object.keys(sessions).map(sessionId => [sessionId, 0])) - const [firstSessionId] = Object.keys(sessions) - if (!sessionMessages.value[activeSessionId.value] && firstSessionId) - activeSessionId.value = firstSessionId - - ensureSession(activeSessionId.value) - } - - function resetAllSessions() { - sessionMessages.value = {} - sessionGenerations.value = {} - activeSessionId.value = 'default' - ensureSession(activeSessionId.value) - } - - watch(systemPrompt, () => { - for (const [sessionId, history] of Object.entries(sessionMessages.value)) { - if (history.length > 0 && history[0].role === 'system') { - sessionMessages.value[sessionId][0] = generateInitialMessage() - } - } - }, { immediate: true }) - - function ingestContextMessage(envelope: ContextMessage) { - const sourceKey = getEventSourceKey(envelope) - if (!activeContexts.value[sourceKey]) { - activeContexts.value[sourceKey] = [] - } - - if (envelope.strategy === ContextUpdateStrategy.ReplaceSelf) { - activeContexts.value[sourceKey] = [envelope] - } - else if (envelope.strategy === ContextUpdateStrategy.AppendSelf) { - activeContexts.value[sourceKey].push(envelope) - } - } - - // ----- Send flow (user -> LLM -> assistant) ----- async function performSend( sendingMessage: string, options: SendOptions, @@ -357,17 +97,17 @@ export const useChatStore = defineStore('chat', () => { if (!sendingMessage && !options.attachments?.length) return - ensureSession(sessionId) + chatSession.ensureSession(sessionId) const sendingCreatedAt = Date.now() const streamingMessageContext: ChatStreamEventContext = { message: { role: 'user', content: sendingMessage, createdAt: sendingCreatedAt }, - contexts: toRaw(activeContexts.value), + contexts: chatContext.getContextsSnapshot(), composedMessage: [], input: options.input, } - const isStaleGeneration = () => getSessionGeneration(sessionId) !== generation + const isStaleGeneration = () => chatSession.getSessionGeneration(sessionId) !== generation const shouldAbort = () => isStaleGeneration() if (shouldAbort()) return @@ -376,23 +116,19 @@ export const useChatStore = defineStore('chat', () => { const isForegroundSession = () => sessionId === activeSessionId.value - // Use a local object for building the message to avoid polluting the UI for background sessions const buildingMessage: StreamingAssistantMessage = { role: 'assistant', content: '', slices: [], tool_results: [], createdAt: Date.now() } - // NOTICE: Clone into reactive state only for the foreground session - // to avoid background streams mutating UI. const updateUI = () => { if (isForegroundSession()) { streamingMessage.value = JSON.parse(JSON.stringify(buildingMessage)) } } - // Initialize UI if foreground updateUI() trackFirstMessage() try { - await emitBeforeMessageComposedHooks(sendingMessage, streamingMessageContext) + await hooks.emitBeforeMessageComposedHooks(sendingMessage, streamingMessageContext) const contentParts: CommonContentPart[] = [{ type: 'text', text: sendingMessage }] @@ -422,36 +158,27 @@ export const useChatStore = defineStore('chat', () => { if (shouldAbort()) return - const sessionMessagesForSend = getSessionMessagesById(sessionId) + const sessionMessagesForSend = chatSession.getSessionMessages(sessionId) sessionMessagesForSend.push({ role: 'user', content: finalContent }) - // Create categorizer for response categorization const categorizer = createStreamingCategorizer(activeProvider.value) - let streamPosition = 0 // Track position in stream for TTS filtering + let streamPosition = 0 const parser = useLlmmarkerParser({ onLiteral: async (literal) => { if (shouldAbort()) return - // Feed to categorizer first categorizer.consume(literal) - // Filter to only include speech parts (exclude reasoning) - // The categorizer handles incomplete tags and filters based on detected tags during streaming const speechOnly = categorizer.filterToSpeech(literal, streamPosition) streamPosition += literal.length - // Only process non-empty speech content (filter empty/whitespace-only chunks) - // Preserve spacing in chunks with content for proper word boundaries if (speechOnly.trim()) { buildingMessage.content += speechOnly - // Emit TTS only for speech parts, not reasoning (clean data, no empty chunks) - await emitTokenLiteralHooks(speechOnly, streamingMessageContext) + await hooks.emitTokenLiteralHooks(speechOnly, streamingMessageContext) - // Add speech content to slices for rendering - // merge text slices for markdown const lastSlice = buildingMessage.slices.at(-1) if (lastSlice?.type === 'text') { lastSlice.text += speechOnly @@ -469,23 +196,21 @@ export const useChatStore = defineStore('chat', () => { if (shouldAbort()) return - await emitTokenSpecialHooks(special, streamingMessageContext) + await hooks.emitTokenSpecialHooks(special, streamingMessageContext) }, onEnd: async (fullText) => { if (isStaleGeneration()) return - // Categorize the full text stream const finalCategorization = categorizeResponse(fullText, activeProvider.value) - // Always store categorization (even if empty) for consistency and memory features buildingMessage.categorization = { speech: finalCategorization.speech, reasoning: finalCategorization.reasoning, } updateUI() }, - minLiteralEmitLength: 24, // Avoid emitting literals too fast. This is a magic number and can be changed later. + minLiteralEmitLength: 24, }) const toolCallQueue = createQueue({ @@ -512,7 +237,7 @@ export const useChatStore = defineStore('chat', () => { const rawMessage = toRaw(withoutContext) if (rawMessage.role === 'assistant') { - const { slices: _, tool_results, categorization: __categorization, ...rest } = rawMessage as ChatAssistantMessage + const { slices: _slices, tool_results, categorization: _categorization, ...rest } = rawMessage as ChatAssistantMessage return { ...toRaw(rest), tool_results: toRaw(tool_results), @@ -522,9 +247,8 @@ export const useChatStore = defineStore('chat', () => { return rawMessage }) - // TODO: possible prototype pollution as key of activeContexts is from external source - // TODO: sanitize keys or use a safer structure - if (Object.keys(activeContexts.value).length > 0) { + const contextsSnapshot = chatContext.getContextsSnapshot() + if (Object.keys(contextsSnapshot).length > 0) { const system = newMessages.slice(0, 1) const afterSystem = newMessages.slice(1, newMessages.length) @@ -533,11 +257,12 @@ export const useChatStore = defineStore('chat', () => { { role: 'user', content: [ - // TODO: use prompt render & i18n system later - // TODO: Module should have description & context length management - { type: 'text', text: '' - + 'These are the contextual information retrieved or on-demand updated from other modules, you may use them as context for chat, or reference of the next action, tool call, etc.:\n' - + `${Object.entries(activeContexts.value).map(([key, value]) => `Module ${key}: ${JSON.stringify(value)}`).join('\n')}\n` }, + { + type: 'text', + text: '' + + 'These are the contextual information retrieved or on-demand updated from other modules, you may use them as context for chat, or reference of the next action, tool call, etc.:\n' + + `${Object.entries(contextsSnapshot).map(([key, value]) => `Module ${key}: ${JSON.stringify(value)}`).join('\n')}\n`, + }, ], }, ...afterSystem, @@ -546,8 +271,8 @@ export const useChatStore = defineStore('chat', () => { streamingMessageContext.composedMessage = newMessages as Message[] - await emitAfterMessageComposedHooks(sendingMessage, streamingMessageContext) - await emitBeforeSendHooks(sendingMessage, streamingMessageContext) + await hooks.emitAfterMessageComposedHooks(sendingMessage, streamingMessageContext) + await hooks.emitBeforeSendHooks(sendingMessage, streamingMessageContext) let fullText = '' const headers = (options.providerConfig?.headers || {}) as Record @@ -580,7 +305,6 @@ export const useChatStore = defineStore('chat', () => { await parser.consume(event.text) break case 'finish': - // Do nothing, resolve break case 'error': throw event.error ?? new Error('Stream error') @@ -588,30 +312,23 @@ export const useChatStore = defineStore('chat', () => { }, }) - // Finalize the parsing of the actual message content - // Categorization and filtering happens in the onEnd callback await parser.end() - // Add the completed message to the history only if it has content if (!isStaleGeneration() && buildingMessage.slices.length > 0) { sessionMessagesForSend.push(toRaw(buildingMessage)) } - // Call the end-of-stream hooks - await emitStreamEndHooks(streamingMessageContext) + await hooks.emitStreamEndHooks(streamingMessageContext) + await hooks.emitAssistantResponseEndHooks(fullText, streamingMessageContext) - // Call the end-of-response hooks with the full text - await emitAssistantResponseEndHooks(fullText, streamingMessageContext) - - await emitAfterSendHooks(sendingMessage, streamingMessageContext) - await emitAssistantMessageHooks({ ...buildingMessage }, fullText, streamingMessageContext) - await emitChatTurnCompleteHooks({ + await hooks.emitAfterSendHooks(sendingMessage, streamingMessageContext) + await hooks.emitAssistantMessageHooks({ ...buildingMessage }, fullText, streamingMessageContext) + await hooks.emitChatTurnCompleteHooks({ output: { ...buildingMessage }, outputText: fullText, toolCalls: sessionMessagesForSend.filter(msg => msg.role === 'tool') as ToolMessage[], }, streamingMessageContext) - // Reset the streaming message for the next turn if (isForegroundSession()) { streamingMessage.value = { role: 'assistant', content: '', slices: [], tool_results: [] } } @@ -625,43 +342,13 @@ export const useChatStore = defineStore('chat', () => { } } - // ----- Remote stream helpers (for broadcast/devtools) ----- - function beginRemoteStream() { - streamingMessage.value = { role: 'assistant', content: '', slices: [], tool_results: [], createdAt: Date.now() } - } - - function appendRemoteLiteral(literal: string) { - streamingMessage.value.content += literal - - const lastSlice = streamingMessage.value.slices.at(-1) - if (lastSlice?.type === 'text') { - lastSlice.text += literal - return - } - - streamingMessage.value.slices.push({ - type: 'text', - text: literal, - }) - } - - function finalizeRemoteStream(fullText?: string) { - const sessionId = activeSessionId.value - const sessionMessagesForSend = getSessionMessagesById(sessionId) - if (streamingMessage.value.slices.length > 0) - sessionMessagesForSend.push(toRaw(streamingMessage.value)) - streamingMessage.value = { role: 'assistant', content: '', slices: [], tool_results: [] } - if (fullText) - streamingMessage.value.content = fullText - } - - async function send( + async function ingest( sendingMessage: string, options: SendOptions, targetSessionId?: string, ) { const sessionId = targetSessionId || activeSessionId.value - const generation = getSessionGeneration(sessionId) + const generation = chatSession.getSessionGeneration(sessionId) return new Promise((resolve, reject) => { sendQueue.enqueue({ @@ -674,51 +361,50 @@ export const useChatStore = defineStore('chat', () => { }) } + function cancelPendingSends(sessionId?: string) { + for (const queued of pendingQueuedSends.value) { + if (sessionId && queued.sessionId !== sessionId) + continue + + queued.cancelled = true + queued.deferred.reject(new Error('Chat session was reset before send could start')) + } + + pendingQueuedSends.value = sessionId + ? pendingQueuedSends.value.filter(item => item.sessionId !== sessionId) + : [] + } + return { sending, - activeSessionId, - messages, - streamingMessage, discoverToolsCompatibility: llmStore.discoverToolsCompatibility, - send, - setActiveSession, - cleanupMessages, - getAllSessions, - replaceSessions, - resetAllSessions, + ingest, + cancelPendingSends, - ingestContextMessage, + clearHooks: hooks.clearHooks, - clearHooks, + emitBeforeMessageComposedHooks: hooks.emitBeforeMessageComposedHooks, + emitAfterMessageComposedHooks: hooks.emitAfterMessageComposedHooks, + emitBeforeSendHooks: hooks.emitBeforeSendHooks, + emitAfterSendHooks: hooks.emitAfterSendHooks, + emitTokenLiteralHooks: hooks.emitTokenLiteralHooks, + emitTokenSpecialHooks: hooks.emitTokenSpecialHooks, + emitStreamEndHooks: hooks.emitStreamEndHooks, + emitAssistantResponseEndHooks: hooks.emitAssistantResponseEndHooks, + emitAssistantMessageHooks: hooks.emitAssistantMessageHooks, + emitChatTurnCompleteHooks: hooks.emitChatTurnCompleteHooks, - emitBeforeMessageComposedHooks, - emitAfterMessageComposedHooks, - emitBeforeSendHooks, - emitAfterSendHooks, - emitTokenLiteralHooks, - emitTokenSpecialHooks, - emitStreamEndHooks, - emitAssistantResponseEndHooks, - emitAssistantMessageHooks, - emitChatTurnCompleteHooks, - - getSessionGenerationValue, - - beginRemoteStream, - appendRemoteLiteral, - finalizeRemoteStream, - - onBeforeMessageComposed, - onAfterMessageComposed, - onBeforeSend, - onAfterSend, - onTokenLiteral, - onTokenSpecial, - onStreamEnd, - onAssistantResponseEnd, - onAssistantMessage, - onChatTurnComplete, + onBeforeMessageComposed: hooks.onBeforeMessageComposed, + onAfterMessageComposed: hooks.onAfterMessageComposed, + onBeforeSend: hooks.onBeforeSend, + onAfterSend: hooks.onAfterSend, + onTokenLiteral: hooks.onTokenLiteral, + onTokenSpecial: hooks.onTokenSpecial, + onStreamEnd: hooks.onStreamEnd, + onAssistantResponseEnd: hooks.onAssistantResponseEnd, + onAssistantMessage: hooks.onAssistantMessage, + onChatTurnComplete: hooks.onChatTurnComplete, } }) diff --git a/packages/stage-ui/src/stores/chat/constants.ts b/packages/stage-ui/src/stores/chat/constants.ts new file mode 100644 index 000000000..eb19ec6bd --- /dev/null +++ b/packages/stage-ui/src/stores/chat/constants.ts @@ -0,0 +1,4 @@ +export const CHAT_STORAGE_KEY = 'chat/messages/v2' +export const ACTIVE_SESSION_STORAGE_KEY = 'chat/active-session' +export const CONTEXT_CHANNEL_NAME = 'airi-context-update' +export const CHAT_STREAM_CHANNEL_NAME = 'airi-chat-stream' diff --git a/packages/stage-ui/src/stores/chat/context-store.ts b/packages/stage-ui/src/stores/chat/context-store.ts new file mode 100644 index 000000000..93807a013 --- /dev/null +++ b/packages/stage-ui/src/stores/chat/context-store.ts @@ -0,0 +1,39 @@ +import type { ContextMessage } from '../../types/chat' + +import { ContextUpdateStrategy } from '@proj-airi/server-sdk' +import { defineStore } from 'pinia' +import { ref, toRaw } from 'vue' + +import { getEventSourceKey } from '../../utils/event-source' + +export const useChatContextStore = defineStore('chat-context', () => { + const activeContexts = ref>({}) + + function ingestContextMessage(envelope: ContextMessage) { + const sourceKey = getEventSourceKey(envelope) + if (!activeContexts.value[sourceKey]) { + activeContexts.value[sourceKey] = [] + } + + if (envelope.strategy === ContextUpdateStrategy.ReplaceSelf) { + activeContexts.value[sourceKey] = [envelope] + } + else if (envelope.strategy === ContextUpdateStrategy.AppendSelf) { + activeContexts.value[sourceKey].push(envelope) + } + } + + function resetContexts() { + activeContexts.value = {} + } + + function getContextsSnapshot() { + return toRaw(activeContexts.value) + } + + return { + ingestContextMessage, + resetContexts, + getContextsSnapshot, + } +}) diff --git a/packages/stage-ui/src/stores/chat/data-store.ts b/packages/stage-ui/src/stores/chat/data-store.ts new file mode 100644 index 000000000..1a9114013 --- /dev/null +++ b/packages/stage-ui/src/stores/chat/data-store.ts @@ -0,0 +1,144 @@ +import type { SystemMessage } from '@xsai/shared-chat' + +import type { ChatHistoryItem } from '../../types/chat' + +export interface ChatDataAccess { + getActiveSessionId: () => string + setActiveSessionId: (sessionId: string) => void + getSessions: () => Record + setSessions: (sessions: Record) => void + getGenerations: () => Record + setGenerations: (generations: Record) => void +} + +export interface ChatDataStore { + ensureSession: (sessionId: string, createInitialMessage: () => SystemMessage) => void + getSessionMessages: (sessionId: string, createInitialMessage: () => SystemMessage) => ChatHistoryItem[] + setSessionMessages: (sessionId: string, next: ChatHistoryItem[]) => void + setActiveSession: (sessionId: string, createInitialMessage: () => SystemMessage) => void + getActiveSessionId: () => string + resetSession: (sessionId: string, createInitialMessage: () => SystemMessage) => void + refreshSystemMessages: (createInitialMessage: () => SystemMessage) => void + replaceSessions: (sessions: Record, createInitialMessage: () => SystemMessage) => void + resetAllSessions: (createInitialMessage: () => SystemMessage) => void + getAllSessions: () => Record + getSessionGeneration: (sessionId: string) => number + bumpSessionGeneration: (sessionId: string) => number + getSessionGenerationValue: (sessionId?: string) => number +} + +export function createChatDataStore(access: ChatDataAccess): ChatDataStore { + function ensureGeneration(sessionId: string) { + const generations = access.getGenerations() + if (generations[sessionId] === undefined) + access.setGenerations({ ...generations, [sessionId]: 0 }) + } + + function getSessionGeneration(sessionId: string) { + ensureGeneration(sessionId) + return access.getGenerations()[sessionId] ?? 0 + } + + function bumpSessionGeneration(sessionId: string) { + const nextGeneration = getSessionGeneration(sessionId) + 1 + access.setGenerations({ ...access.getGenerations(), [sessionId]: nextGeneration }) + return nextGeneration + } + + function ensureSession(sessionId: string, createInitialMessage: () => SystemMessage) { + ensureGeneration(sessionId) + + const sessions = access.getSessions() + if (!sessions[sessionId] || sessions[sessionId].length === 0) { + access.setSessions({ + ...sessions, + [sessionId]: [createInitialMessage()], + }) + } + } + + function getSessionMessages(sessionId: string, createInitialMessage: () => SystemMessage) { + ensureSession(sessionId, createInitialMessage) + return access.getSessions()[sessionId]! + } + + function setSessionMessages(sessionId: string, next: ChatHistoryItem[]) { + access.setSessions({ + ...access.getSessions(), + [sessionId]: next, + }) + } + + function setActiveSession(sessionId: string, createInitialMessage: () => SystemMessage) { + access.setActiveSessionId(sessionId) + ensureSession(sessionId, createInitialMessage) + } + + function getActiveSessionId() { + return access.getActiveSessionId() + } + + function resetSession(sessionId: string, createInitialMessage: () => SystemMessage) { + bumpSessionGeneration(sessionId) + setSessionMessages(sessionId, [createInitialMessage()]) + } + + function refreshSystemMessages(createInitialMessage: () => SystemMessage) { + const sessions = access.getSessions() + const nextSessions: Record = {} + + for (const [sessionId, history] of Object.entries(sessions)) { + if (history.length > 0 && history[0].role === 'system') { + nextSessions[sessionId] = [createInitialMessage(), ...history.slice(1)] + } + else { + nextSessions[sessionId] = history + } + } + + access.setSessions(nextSessions) + } + + function replaceSessions(sessions: Record, createInitialMessage: () => SystemMessage) { + access.setSessions(sessions) + access.setGenerations(Object.fromEntries(Object.keys(sessions).map(sessionId => [sessionId, 0]))) + + const [firstSessionId] = Object.keys(sessions) + if (!sessions[access.getActiveSessionId()] && firstSessionId) + access.setActiveSessionId(firstSessionId) + + ensureSession(access.getActiveSessionId(), createInitialMessage) + } + + function resetAllSessions(createInitialMessage: () => SystemMessage) { + access.setSessions({}) + access.setGenerations({}) + access.setActiveSessionId('default') + ensureSession('default', createInitialMessage) + } + + function getAllSessions() { + return JSON.parse(JSON.stringify(access.getSessions())) as Record + } + + function getSessionGenerationValue(sessionId?: string) { + const targetSessionId = sessionId ?? access.getActiveSessionId() + return getSessionGeneration(targetSessionId) + } + + return { + ensureSession, + getSessionMessages, + setSessionMessages, + setActiveSession, + getActiveSessionId, + resetSession, + refreshSystemMessages, + replaceSessions, + resetAllSessions, + getAllSessions, + getSessionGeneration, + bumpSessionGeneration, + getSessionGenerationValue, + } +} diff --git a/packages/stage-ui/src/stores/chat/hooks.ts b/packages/stage-ui/src/stores/chat/hooks.ts new file mode 100644 index 000000000..cb4f78c17 --- /dev/null +++ b/packages/stage-ui/src/stores/chat/hooks.ts @@ -0,0 +1,217 @@ +import type { ToolMessage } from '@xsai/shared-chat' + +import type { ChatStreamEventContext, StreamingAssistantMessage } from '../../types/chat' + +export interface ChatHookRegistry { + onBeforeMessageComposed: (cb: (message: string, context: Omit) => Promise) => () => void + onAfterMessageComposed: (cb: (message: string, context: ChatStreamEventContext) => Promise) => () => void + onBeforeSend: (cb: (message: string, context: ChatStreamEventContext) => Promise) => () => void + onAfterSend: (cb: (message: string, context: ChatStreamEventContext) => Promise) => () => void + onTokenLiteral: (cb: (literal: string, context: ChatStreamEventContext) => Promise) => () => void + onTokenSpecial: (cb: (special: string, context: ChatStreamEventContext) => Promise) => () => void + onStreamEnd: (cb: (context: ChatStreamEventContext) => Promise) => () => void + onAssistantResponseEnd: (cb: (message: string, context: ChatStreamEventContext) => Promise) => () => void + onAssistantMessage: (cb: (message: StreamingAssistantMessage, messageText: string, context: ChatStreamEventContext) => Promise) => () => void + onChatTurnComplete: (cb: (chat: { output: StreamingAssistantMessage, outputText: string, toolCalls: ToolMessage[] }, context: ChatStreamEventContext) => Promise) => () => void + emitBeforeMessageComposedHooks: (message: string, context: Omit) => Promise + emitAfterMessageComposedHooks: (message: string, context: ChatStreamEventContext) => Promise + emitBeforeSendHooks: (message: string, context: ChatStreamEventContext) => Promise + emitAfterSendHooks: (message: string, context: ChatStreamEventContext) => Promise + emitTokenLiteralHooks: (literal: string, context: ChatStreamEventContext) => Promise + emitTokenSpecialHooks: (special: string, context: ChatStreamEventContext) => Promise + emitStreamEndHooks: (context: ChatStreamEventContext) => Promise + emitAssistantResponseEndHooks: (message: string, context: ChatStreamEventContext) => Promise + emitAssistantMessageHooks: (message: StreamingAssistantMessage, messageText: string, context: ChatStreamEventContext) => Promise + emitChatTurnCompleteHooks: (chat: { output: StreamingAssistantMessage, outputText: string, toolCalls: ToolMessage[] }, context: ChatStreamEventContext) => Promise + clearHooks: () => void +} + +export function createChatHooks(): ChatHookRegistry { + const onBeforeMessageComposedHooks: Array<(message: string, context: Omit) => Promise> = [] + const onAfterMessageComposedHooks: Array<(message: string, context: ChatStreamEventContext) => Promise> = [] + const onBeforeSendHooks: Array<(message: string, context: ChatStreamEventContext) => Promise> = [] + const onAfterSendHooks: Array<(message: string, context: ChatStreamEventContext) => Promise> = [] + const onTokenLiteralHooks: Array<(literal: string, context: ChatStreamEventContext) => Promise> = [] + const onTokenSpecialHooks: Array<(special: string, context: ChatStreamEventContext) => Promise> = [] + const onStreamEndHooks: Array<(context: ChatStreamEventContext) => Promise> = [] + const onAssistantResponseEndHooks: Array<(message: string, context: ChatStreamEventContext) => Promise> = [] + const onAssistantMessageHooks: Array<(message: StreamingAssistantMessage, messageText: string, context: ChatStreamEventContext) => Promise> = [] + const onChatTurnCompleteHooks: Array<(chat: { output: StreamingAssistantMessage, outputText: string, toolCalls: ToolMessage[] }, context: ChatStreamEventContext) => Promise> = [] + + function onBeforeMessageComposed(cb: (message: string, context: Omit) => Promise) { + onBeforeMessageComposedHooks.push(cb) + return () => { + const index = onBeforeMessageComposedHooks.indexOf(cb) + if (index >= 0) + onBeforeMessageComposedHooks.splice(index, 1) + } + } + + function onAfterMessageComposed(cb: (message: string, context: ChatStreamEventContext) => Promise) { + onAfterMessageComposedHooks.push(cb) + return () => { + const index = onAfterMessageComposedHooks.indexOf(cb) + if (index >= 0) + onAfterMessageComposedHooks.splice(index, 1) + } + } + + function onBeforeSend(cb: (message: string, context: ChatStreamEventContext) => Promise) { + onBeforeSendHooks.push(cb) + return () => { + const index = onBeforeSendHooks.indexOf(cb) + if (index >= 0) + onBeforeSendHooks.splice(index, 1) + } + } + + function onAfterSend(cb: (message: string, context: ChatStreamEventContext) => Promise) { + onAfterSendHooks.push(cb) + return () => { + const index = onAfterSendHooks.indexOf(cb) + if (index >= 0) + onAfterSendHooks.splice(index, 1) + } + } + + function onTokenLiteral(cb: (literal: string, context: ChatStreamEventContext) => Promise) { + onTokenLiteralHooks.push(cb) + return () => { + const index = onTokenLiteralHooks.indexOf(cb) + if (index >= 0) + onTokenLiteralHooks.splice(index, 1) + } + } + + function onTokenSpecial(cb: (special: string, context: ChatStreamEventContext) => Promise) { + onTokenSpecialHooks.push(cb) + return () => { + const index = onTokenSpecialHooks.indexOf(cb) + if (index >= 0) + onTokenSpecialHooks.splice(index, 1) + } + } + + function onStreamEnd(cb: (context: ChatStreamEventContext) => Promise) { + onStreamEndHooks.push(cb) + return () => { + const index = onStreamEndHooks.indexOf(cb) + if (index >= 0) + onStreamEndHooks.splice(index, 1) + } + } + + function onAssistantResponseEnd(cb: (message: string, context: ChatStreamEventContext) => Promise) { + onAssistantResponseEndHooks.push(cb) + return () => { + const index = onAssistantResponseEndHooks.indexOf(cb) + if (index >= 0) + onAssistantResponseEndHooks.splice(index, 1) + } + } + + function onAssistantMessage(cb: (message: StreamingAssistantMessage, messageText: string, context: ChatStreamEventContext) => Promise) { + onAssistantMessageHooks.push(cb) + return () => { + const index = onAssistantMessageHooks.indexOf(cb) + if (index >= 0) + onAssistantMessageHooks.splice(index, 1) + } + } + + function onChatTurnComplete(cb: (chat: { output: StreamingAssistantMessage, outputText: string, toolCalls: ToolMessage[] }, context: ChatStreamEventContext) => Promise) { + onChatTurnCompleteHooks.push(cb) + return () => { + const index = onChatTurnCompleteHooks.indexOf(cb) + if (index >= 0) + onChatTurnCompleteHooks.splice(index, 1) + } + } + + function clearHooks() { + onBeforeMessageComposedHooks.length = 0 + onAfterMessageComposedHooks.length = 0 + onBeforeSendHooks.length = 0 + onAfterSendHooks.length = 0 + onTokenLiteralHooks.length = 0 + onTokenSpecialHooks.length = 0 + onStreamEndHooks.length = 0 + onAssistantResponseEndHooks.length = 0 + onAssistantMessageHooks.length = 0 + onChatTurnCompleteHooks.length = 0 + } + + async function emitBeforeMessageComposedHooks(message: string, context: Omit) { + for (const hook of onBeforeMessageComposedHooks) + await hook(message, context) + } + + async function emitAfterMessageComposedHooks(message: string, context: ChatStreamEventContext) { + for (const hook of onAfterMessageComposedHooks) + await hook(message, context) + } + + async function emitBeforeSendHooks(message: string, context: ChatStreamEventContext) { + for (const hook of onBeforeSendHooks) + await hook(message, context) + } + + async function emitAfterSendHooks(message: string, context: ChatStreamEventContext) { + for (const hook of onAfterSendHooks) + await hook(message, context) + } + + async function emitTokenLiteralHooks(literal: string, context: ChatStreamEventContext) { + for (const hook of onTokenLiteralHooks) + await hook(literal, context) + } + + async function emitTokenSpecialHooks(special: string, context: ChatStreamEventContext) { + for (const hook of onTokenSpecialHooks) + await hook(special, context) + } + + async function emitStreamEndHooks(context: ChatStreamEventContext) { + for (const hook of onStreamEndHooks) + await hook(context) + } + + async function emitAssistantResponseEndHooks(message: string, context: ChatStreamEventContext) { + for (const hook of onAssistantResponseEndHooks) + await hook(message, context) + } + + async function emitAssistantMessageHooks(message: StreamingAssistantMessage, messageText: string, context: ChatStreamEventContext) { + for (const hook of onAssistantMessageHooks) + await hook(message, messageText, context) + } + + async function emitChatTurnCompleteHooks(chat: { output: StreamingAssistantMessage, outputText: string, toolCalls: ToolMessage[] }, context: ChatStreamEventContext) { + for (const hook of onChatTurnCompleteHooks) + await hook(chat, context) + } + + return { + onBeforeMessageComposed, + onAfterMessageComposed, + onBeforeSend, + onAfterSend, + onTokenLiteral, + onTokenSpecial, + onStreamEnd, + onAssistantResponseEnd, + onAssistantMessage, + onChatTurnComplete, + emitBeforeMessageComposedHooks, + emitAfterMessageComposedHooks, + emitBeforeSendHooks, + emitAfterSendHooks, + emitTokenLiteralHooks, + emitTokenSpecialHooks, + emitStreamEndHooks, + emitAssistantResponseEndHooks, + emitAssistantMessageHooks, + emitChatTurnCompleteHooks, + clearHooks, + } +} diff --git a/packages/stage-ui/src/stores/chat/maintenance.ts b/packages/stage-ui/src/stores/chat/maintenance.ts new file mode 100644 index 000000000..dd55f6abb --- /dev/null +++ b/packages/stage-ui/src/stores/chat/maintenance.ts @@ -0,0 +1,24 @@ +import { defineStore } from 'pinia' + +import { useChatOrchestratorStore } from '../chat' +import { useChatContextStore } from './context-store' +import { useChatSessionStore } from './session-store' +import { useChatStreamStore } from './stream-store' + +export const useChatMaintenanceStore = defineStore('chat-maintenance', () => { + const chatSession = useChatSessionStore() + const chatStream = useChatStreamStore() + const chatContext = useChatContextStore() + const chatOrchestrator = useChatOrchestratorStore() + + function cleanupMessages(sessionId = chatSession.activeSessionId) { + chatSession.cleanupMessages(sessionId) + chatContext.resetContexts() + chatOrchestrator.cancelPendingSends(sessionId) + chatStream.resetStream() + } + + return { + cleanupMessages, + } +}) diff --git a/packages/stage-ui/src/stores/chat/session-store.ts b/packages/stage-ui/src/stores/chat/session-store.ts new file mode 100644 index 000000000..25a62ebf9 --- /dev/null +++ b/packages/stage-ui/src/stores/chat/session-store.ts @@ -0,0 +1,93 @@ +import type { SystemMessage } from '@xsai/shared-chat' + +import type { ChatHistoryItem } from '../../types/chat' + +import { useLocalStorage } from '@vueuse/core' +import { defineStore, storeToRefs } from 'pinia' +import { computed, ref, watch } from 'vue' + +import { useCharacterStore } from '../character' +import { ACTIVE_SESSION_STORAGE_KEY, CHAT_STORAGE_KEY } from './constants' +import { createChatDataStore } from './data-store' + +export const useChatSessionStore = defineStore('chat-session', () => { + const { systemPrompt } = storeToRefs(useCharacterStore()) + + const activeSessionId = useLocalStorage(ACTIVE_SESSION_STORAGE_KEY, 'default') + const sessionMessages = useLocalStorage>(CHAT_STORAGE_KEY, {}) + const sessionGenerations = ref>({}) + + const dataStore = createChatDataStore({ + getActiveSessionId: () => activeSessionId.value, + setActiveSessionId: sessionId => activeSessionId.value = sessionId, + getSessions: () => sessionMessages.value, + setSessions: sessions => sessionMessages.value = sessions, + getGenerations: () => sessionGenerations.value, + setGenerations: generations => sessionGenerations.value = generations, + }) + + // 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' + const mathSyntaxSystemPrompt = '- For any math equation, use LaTeX format, eg: $ x^3 $, always escape dollar sign outside math equation\n' + + function generateInitialMessage() { + const content = codeBlockSystemPrompt + mathSyntaxSystemPrompt + systemPrompt.value + + return { + role: 'system', + content, + } satisfies SystemMessage + } + + function ensureSession(sessionId: string) { + dataStore.ensureSession(sessionId, generateInitialMessage) + } + + ensureSession(activeSessionId.value) + + const messages = computed({ + get: () => dataStore.getSessionMessages(activeSessionId.value, generateInitialMessage), + set: value => dataStore.setSessionMessages(activeSessionId.value, value), + }) + + function setActiveSession(sessionId: string) { + dataStore.setActiveSession(sessionId, generateInitialMessage) + } + + function cleanupMessages(sessionId = activeSessionId.value) { + dataStore.resetSession(sessionId, generateInitialMessage) + } + + function getAllSessions() { + return dataStore.getAllSessions() + } + + function replaceSessions(sessions: Record) { + dataStore.replaceSessions(sessions, generateInitialMessage) + } + + function resetAllSessions() { + dataStore.resetAllSessions(generateInitialMessage) + } + + watch(systemPrompt, () => { + dataStore.refreshSystemMessages(generateInitialMessage) + }, { immediate: true }) + + return { + activeSessionId, + messages, + + setActiveSession, + cleanupMessages, + getAllSessions, + replaceSessions, + resetAllSessions, + + ensureSession, + getSessionMessages: (sessionId: string) => dataStore.getSessionMessages(sessionId, generateInitialMessage), + getSessionGeneration: (sessionId: string) => dataStore.getSessionGeneration(sessionId), + bumpSessionGeneration: (sessionId: string) => dataStore.bumpSessionGeneration(sessionId), + getSessionGenerationValue: (sessionId?: string) => dataStore.getSessionGenerationValue(sessionId), + } +}) diff --git a/packages/stage-ui/src/stores/chat/state.ts b/packages/stage-ui/src/stores/chat/state.ts new file mode 100644 index 000000000..723c1506c --- /dev/null +++ b/packages/stage-ui/src/stores/chat/state.ts @@ -0,0 +1,2 @@ +export * from './session-store' +export * from './stream-store' diff --git a/packages/stage-ui/src/stores/chat/stream-store.ts b/packages/stage-ui/src/stores/chat/stream-store.ts new file mode 100644 index 000000000..030e72d3b --- /dev/null +++ b/packages/stage-ui/src/stores/chat/stream-store.ts @@ -0,0 +1,52 @@ +import type { StreamingAssistantMessage } from '../../types/chat' + +import { defineStore } from 'pinia' +import { ref } from 'vue' + +import { useChatSessionStore } from './session-store' + +export const useChatStreamStore = defineStore('chat-stream', () => { + const chatSession = useChatSessionStore() + const streamingMessage = ref({ role: 'assistant', content: '', slices: [], tool_results: [], createdAt: Date.now() }) + + function beginStream() { + streamingMessage.value = { role: 'assistant', content: '', slices: [], tool_results: [], createdAt: Date.now() } + } + + function appendStreamLiteral(literal: string) { + streamingMessage.value.content += literal + + const lastSlice = streamingMessage.value.slices.at(-1) + if (lastSlice?.type === 'text') { + lastSlice.text += literal + return + } + + streamingMessage.value.slices.push({ + type: 'text', + text: literal, + }) + } + + function finalizeStream(fullText?: string) { + const sessionId = chatSession.activeSessionId + const sessionMessagesForSend = chatSession.getSessionMessages(sessionId) + if (streamingMessage.value.slices.length > 0) + sessionMessagesForSend.push(streamingMessage.value) + streamingMessage.value = { role: 'assistant', content: '', slices: [], tool_results: [] } + if (fullText) + streamingMessage.value.content = fullText + } + + function resetStream() { + streamingMessage.value = { role: 'assistant', content: '', slices: [], tool_results: [] } + } + + return { + streamingMessage, + beginStream, + appendStreamLiteral, + finalizeStream, + resetStream, + } +}) diff --git a/packages/stage-ui/src/stores/markdown-stress.ts b/packages/stage-ui/src/stores/markdown-stress.ts index 5b170e15c..3355fbf70 100644 --- a/packages/stage-ui/src/stores/markdown-stress.ts +++ b/packages/stage-ui/src/stores/markdown-stress.ts @@ -7,7 +7,7 @@ import { defaultPerfTracer, exportCsv as exportCsvFile } from '@proj-airi/stage- import { defineStore, storeToRefs } from 'pinia' import { ref } from 'vue' -import { useChatStore } from './chat' +import { useChatOrchestratorStore } from './chat' import { useLLM } from './llm' import { useConsciousnessStore } from './modules/consciousness' import { usePerfTracerBridgeStore } from './perf-tracer-bridge' @@ -292,7 +292,7 @@ export const useMarkdownStressStore = defineStore('markdownStress', () => { } async function runOnlineScenario() { - const chatStore = useChatStore() + const chatStore = useChatOrchestratorStore() const targetScenario = ensureScenario() const provider = await providersStore.getProviderInstance(activeProvider.value) as ChatProvider | undefined @@ -309,7 +309,7 @@ export const useMarkdownStressStore = defineStore('markdownStress', () => { const delay = Math.max(0, runStart + message.atMs - performance.now()) const timer = setTimeout(async () => { try { - await chatStore.send(message.text, { + await chatStore.ingest(message.text, { model: activeModel.value!, chatProvider: provider, }) @@ -323,7 +323,7 @@ export const useMarkdownStressStore = defineStore('markdownStress', () => { } async function runMockScenario() { - const chatStore = useChatStore() + const chatStore = useChatOrchestratorStore() const llm = useLLM() const targetScenario = ensureScenario() const modelToUse = mockModelId @@ -365,7 +365,7 @@ export const useMarkdownStressStore = defineStore('markdownStress', () => { const delay = Math.max(0, runStart + message.atMs - performance.now()) const timer = setTimeout(async () => { try { - await chatStore.send(message.text, { + await chatStore.ingest(message.text, { model: modelToUse, chatProvider: mockProvider, }) diff --git a/packages/stage-ui/src/stores/mods/api/channel-gateway.ts b/packages/stage-ui/src/stores/mods/api/channel-gateway.ts new file mode 100644 index 000000000..a43805154 --- /dev/null +++ b/packages/stage-ui/src/stores/mods/api/channel-gateway.ts @@ -0,0 +1,123 @@ +import type { EventEnvelope } from './events' + +export type GatewayEvent = EventEnvelope + +export interface GatewayChannel { + name: string + in?: ReadableStream + out?: (event: GatewayEvent) => void + canHandle?: (event: GatewayEvent) => boolean +} + +export interface GatewayRoute { + match: (event: GatewayEvent) => boolean + to: string[] + mode?: 'fan-out' | 'first' | 'all' +} + +export interface DispatchOptions { + origin?: string +} + +export interface ChannelGateway { + register: (channel: GatewayChannel) => void + unregister: (name: string) => void + dispatch: (event: GatewayEvent, options?: DispatchOptions) => void + route: (rule: GatewayRoute) => void + clearRoutes: () => void +} + +export function createChannelGateway(): ChannelGateway { + const channels = new Map() + const routes: GatewayRoute[] = [] + const readers = new Map>() + + function dispatch(event: GatewayEvent, options?: DispatchOptions) { + const matchedRoutes = routes.filter(rule => rule.match(event)) + + if (matchedRoutes.length > 0) { + for (const rule of matchedRoutes) { + const targets = rule.to + .map(name => channels.get(name)) + .filter((channel): channel is GatewayChannel => !!channel) + + if (rule.mode === 'first') { + const target = targets.find(channel => channel.out) + if (target && target.name !== options?.origin) + target.out?.(event) + continue + } + + for (const target of targets) { + if (!target.out) + continue + if (target.name === options?.origin) + continue + target.out(event) + if (rule.mode === 'all') + continue + } + } + return + } + + for (const channel of channels.values()) { + if (channel.name === options?.origin) + continue + if (channel.canHandle && !channel.canHandle(event)) + continue + channel.out?.(event) + } + } + + function register(channel: GatewayChannel) { + channels.set(channel.name, channel) + + if (!channel.in) + return + + const reader = channel.in.getReader() + readers.set(channel.name, reader) + + const pump = async () => { + try { + while (true) { + const result = await reader.read() + if (result.done) + break + dispatch(result.value, { origin: channel.name }) + } + } + catch (error) { + console.warn('Channel gateway stream error:', channel.name, error) + } + } + + void pump() + } + + function unregister(name: string) { + channels.delete(name) + const reader = readers.get(name) + if (reader) { + reader.cancel().catch(() => undefined) + readers.delete(name) + } + } + + function route(rule: GatewayRoute) { + routes.push(rule) + } + + function clearRoutes() { + routes.length = 0 + } + + return { + register, + unregister, + dispatch, + route, + clearRoutes, + } +} 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 ca143111d..341af57ea 100644 --- a/packages/stage-ui/src/stores/mods/api/context-bridge.ts +++ b/packages/stage-ui/src/stores/mods/api/context-bridge.ts @@ -10,7 +10,11 @@ import { nanoid } from 'nanoid' import { defineStore, storeToRefs } from 'pinia' import { ref, toRaw, watch } from 'vue' -import { CHAT_STREAM_CHANNEL_NAME, CONTEXT_CHANNEL_NAME, useChatStore } from '../../chat' +import { useChatOrchestratorStore } from '../../chat' +import { CHAT_STREAM_CHANNEL_NAME, CONTEXT_CHANNEL_NAME } from '../../chat/constants' +import { useChatContextStore } from '../../chat/context-store' +import { useChatSessionStore } from '../../chat/session-store' +import { useChatStreamStore } from '../../chat/stream-store' import { useConsciousnessStore } from '../../modules/consciousness' import { useProvidersStore } from '../../providers' import { useModsServerChannelStore } from './channel-server' @@ -18,7 +22,10 @@ import { useModsServerChannelStore } from './channel-server' export const useContextBridgeStore = defineStore('mods:api:context-bridge', () => { const mutex = new Mutex() - const chatStore = useChatStore() + const chatOrchestrator = useChatOrchestratorStore() + const chatSession = useChatSessionStore() + const chatStream = useChatStreamStore() + const chatContext = useChatContextStore() const serverChannelStore = useModsServerChannelStore() const consciousnessStore = useConsciousnessStore() const providersStore = useProvidersStore() @@ -38,7 +45,7 @@ export const useContextBridgeStore = defineStore('mods:api:context-bridge', () = const { stop } = watch(incomingContext, (event) => { if (event) - chatStore.ingestContextMessage(event) + chatContext.ingestContextMessage(event) }) disposeHookFns.value.push(stop) @@ -48,7 +55,7 @@ export const useContextBridgeStore = defineStore('mods:api:context-bridge', () = metadata: event.metadata, createdAt: Date.now(), } - chatStore.ingestContextMessage(contextMessage) + chatContext.ingestContextMessage(contextMessage) broadcastContext(toRaw(contextMessage)) })) @@ -73,7 +80,7 @@ export const useContextBridgeStore = defineStore('mods:api:context-bridge', () = if (normalizedContextUpdates?.length) { const createdAt = Date.now() for (const update of normalizedContextUpdates) { - chatStore.ingestContextMessage({ + chatContext.ingestContextMessage({ ...update, metadata: event.metadata, createdAt, @@ -88,7 +95,7 @@ export const useContextBridgeStore = defineStore('mods:api:context-bridge', () = if (overrides?.messagePrefix) messageText = `${overrides.messagePrefix}${text}` - await chatStore.send(messageText, { + await chatOrchestrator.ingest(messageText, { model: activeModel.value, chatProvider, input: { @@ -105,56 +112,56 @@ export const useContextBridgeStore = defineStore('mods:api:context-bridge', () = })) disposeHookFns.value.push( - chatStore.onBeforeMessageComposed(async (message, context) => { + chatOrchestrator.onBeforeMessageComposed(async (message, context) => { if (isProcessingRemoteStream) return - broadcastStreamEvent({ type: 'before-compose', message, sessionId: chatStore.activeSessionId, context: structuredClone(toRaw(context)) }) + broadcastStreamEvent({ type: 'before-compose', message, sessionId: chatSession.activeSessionId, context: structuredClone(toRaw(context)) }) }), - chatStore.onAfterMessageComposed(async (message, context) => { + chatOrchestrator.onAfterMessageComposed(async (message, context) => { if (isProcessingRemoteStream) return - broadcastStreamEvent({ type: 'after-compose', message, sessionId: chatStore.activeSessionId, context: structuredClone(toRaw(context)) }) + broadcastStreamEvent({ type: 'after-compose', message, sessionId: chatSession.activeSessionId, context: structuredClone(toRaw(context)) }) }), - chatStore.onBeforeSend(async (message, context) => { + chatOrchestrator.onBeforeSend(async (message, context) => { if (isProcessingRemoteStream) return - broadcastStreamEvent({ type: 'before-send', message, sessionId: chatStore.activeSessionId, context: structuredClone(toRaw(context)) }) + broadcastStreamEvent({ type: 'before-send', message, sessionId: chatSession.activeSessionId, context: structuredClone(toRaw(context)) }) }), - chatStore.onAfterSend(async (message, context) => { + chatOrchestrator.onAfterSend(async (message, context) => { if (isProcessingRemoteStream) return - broadcastStreamEvent({ type: 'after-send', message, sessionId: chatStore.activeSessionId, context: structuredClone(toRaw(context)) }) + broadcastStreamEvent({ type: 'after-send', message, sessionId: chatSession.activeSessionId, context: structuredClone(toRaw(context)) }) }), - chatStore.onTokenLiteral(async (literal, context) => { + chatOrchestrator.onTokenLiteral(async (literal, context) => { if (isProcessingRemoteStream) return - broadcastStreamEvent({ type: 'token-literal', literal, sessionId: chatStore.activeSessionId, context: structuredClone(toRaw(context)) }) + broadcastStreamEvent({ type: 'token-literal', literal, sessionId: chatSession.activeSessionId, context: structuredClone(toRaw(context)) }) }), - chatStore.onTokenSpecial(async (special, context) => { + chatOrchestrator.onTokenSpecial(async (special, context) => { if (isProcessingRemoteStream) return - broadcastStreamEvent({ type: 'token-special', special, sessionId: chatStore.activeSessionId, context: structuredClone(toRaw(context)) }) + broadcastStreamEvent({ type: 'token-special', special, sessionId: chatSession.activeSessionId, context: structuredClone(toRaw(context)) }) }), - chatStore.onStreamEnd(async (context) => { + chatOrchestrator.onStreamEnd(async (context) => { if (isProcessingRemoteStream) return - broadcastStreamEvent({ type: 'stream-end', sessionId: chatStore.activeSessionId, context: structuredClone(toRaw(context)) }) + broadcastStreamEvent({ type: 'stream-end', sessionId: chatSession.activeSessionId, context: structuredClone(toRaw(context)) }) }), - chatStore.onAssistantResponseEnd(async (message, context) => { + chatOrchestrator.onAssistantResponseEnd(async (message, context) => { if (isProcessingRemoteStream) return - broadcastStreamEvent({ type: 'assistant-end', message, sessionId: chatStore.activeSessionId, context: structuredClone(toRaw(context)) }) + broadcastStreamEvent({ type: 'assistant-end', message, sessionId: chatSession.activeSessionId, context: structuredClone(toRaw(context)) }) }), - chatStore.onAssistantMessage(async (message, _messageText, context) => { + chatOrchestrator.onAssistantMessage(async (message, _messageText, context) => { serverChannelStore.send({ type: 'output:gen-ai:chat:message', data: { @@ -172,7 +179,7 @@ export const useContextBridgeStore = defineStore('mods:api:context-bridge', () = }) }), - chatStore.onChatTurnComplete(async (chat, context) => { + chatOrchestrator.onChatTurnComplete(async (chat, context) => { serverChannelStore.send({ type: 'output:gen-ai:chat:complete', data: { @@ -209,58 +216,58 @@ export const useContextBridgeStore = defineStore('mods:api:context-bridge', () = // Use the receiver's active session to avoid clobbering chat state when events come from other windows/devtools. switch (event.type) { case 'before-compose': - await chatStore.emitBeforeMessageComposedHooks(event.message, event.context) + await chatOrchestrator.emitBeforeMessageComposedHooks(event.message, event.context) break case 'after-compose': - await chatStore.emitAfterMessageComposedHooks(event.message, event.context) + await chatOrchestrator.emitAfterMessageComposedHooks(event.message, event.context) break case 'before-send': - await chatStore.emitBeforeSendHooks(event.message, event.context) + await chatOrchestrator.emitBeforeSendHooks(event.message, event.context) remoteStreamGuard = { - sessionId: chatStore.activeSessionId, - generation: chatStore.getSessionGenerationValue(), + sessionId: chatSession.activeSessionId, + generation: chatSession.getSessionGenerationValue(), } - chatStore.sending = true - chatStore.beginRemoteStream() + chatOrchestrator.sending = true + chatStream.beginStream() break case 'after-send': - await chatStore.emitAfterSendHooks(event.message, event.context) + await chatOrchestrator.emitAfterSendHooks(event.message, event.context) break case 'token-literal': if (!remoteStreamGuard) return - if (remoteStreamGuard.sessionId !== chatStore.activeSessionId) + if (remoteStreamGuard.sessionId !== chatSession.activeSessionId) return - if (chatStore.getSessionGenerationValue(remoteStreamGuard.sessionId) !== remoteStreamGuard.generation) + if (chatSession.getSessionGenerationValue(remoteStreamGuard.sessionId) !== remoteStreamGuard.generation) return - chatStore.appendRemoteLiteral(event.literal) - await chatStore.emitTokenLiteralHooks(event.literal, event.context) + chatStream.appendStreamLiteral(event.literal) + await chatOrchestrator.emitTokenLiteralHooks(event.literal, event.context) break case 'token-special': - await chatStore.emitTokenSpecialHooks(event.special, event.context) + await chatOrchestrator.emitTokenSpecialHooks(event.special, event.context) break case 'stream-end': if (!remoteStreamGuard) break - if (remoteStreamGuard.sessionId !== chatStore.activeSessionId) + if (remoteStreamGuard.sessionId !== chatSession.activeSessionId) break - if (chatStore.getSessionGenerationValue(remoteStreamGuard.sessionId) !== remoteStreamGuard.generation) + if (chatSession.getSessionGenerationValue(remoteStreamGuard.sessionId) !== remoteStreamGuard.generation) break - await chatStore.emitStreamEndHooks(event.context) - chatStore.finalizeRemoteStream() - chatStore.sending = false + await chatOrchestrator.emitStreamEndHooks(event.context) + chatStream.finalizeStream() + chatOrchestrator.sending = false remoteStreamGuard = null break case 'assistant-end': if (!remoteStreamGuard) break - if (remoteStreamGuard.sessionId !== chatStore.activeSessionId) + if (remoteStreamGuard.sessionId !== chatSession.activeSessionId) break - if (chatStore.getSessionGenerationValue(remoteStreamGuard.sessionId) !== remoteStreamGuard.generation) + if (chatSession.getSessionGenerationValue(remoteStreamGuard.sessionId) !== remoteStreamGuard.generation) break - await chatStore.emitAssistantResponseEndHooks(event.message, event.context) - chatStore.finalizeRemoteStream(event.message) - chatStore.sending = false + await chatOrchestrator.emitAssistantResponseEndHooks(event.message, event.context) + chatStream.finalizeStream(event.message) + chatOrchestrator.sending = false remoteStreamGuard = null break } diff --git a/packages/stage-ui/src/stores/mods/api/events.ts b/packages/stage-ui/src/stores/mods/api/events.ts new file mode 100644 index 000000000..abfccf1c2 --- /dev/null +++ b/packages/stage-ui/src/stores/mods/api/events.ts @@ -0,0 +1,54 @@ +import { nanoid } from 'nanoid' + +export type EventPriority = 'critical' | 'high' | 'normal' | 'low' + +export interface EventEnvelope { + id: string + type: TType + time: number + priority?: EventPriority + source?: string + tags?: string[] + payload: TPayload +} + +export interface EventStream { + stream: ReadableStream + emit: (event: T) => void + close: () => void +} + +export function createEvent(type: string, payload: TPayload, options?: { priority?: EventPriority, source?: string, tags?: string[], id?: string, time?: number }): EventEnvelope { + return { + id: options?.id ?? nanoid(), + type, + time: options?.time ?? Date.now(), + priority: options?.priority, + source: options?.source, + tags: options?.tags, + payload, + } +} + +export function createEventStream(): EventStream { + let controller: ReadableStreamDefaultController | undefined + const stream = new ReadableStream({ + start(ctrl) { + controller = ctrl + }, + cancel() { + controller = undefined + }, + }) + + return { + stream, + emit(event) { + controller?.enqueue(event) + }, + close() { + controller?.close() + controller = undefined + }, + } +}