refactor(stage-ui,stage-web,stage-tamagotchi): use new server-sdk for context bridging, forming module context correctly

This commit is contained in:
Neko Ayaka
2025-12-27 03:50:59 +08:00
parent 3fe0e4c091
commit 7b5f209e22
13 changed files with 264 additions and 485 deletions
+7 -5
View File
@@ -1,9 +1,10 @@
<script setup lang="ts">
import { defineInvoke, defineInvokeHandler } from '@moeru/eventa'
import { useContextBridge } from '@proj-airi/stage-ui/composables'
import { useDisplayModelsStore } from '@proj-airi/stage-ui/stores/display-models'
import { useModsServerChannelStore } from '@proj-airi/stage-ui/stores/mods/api/channel-server'
import { useAiriCardStore } from '@proj-airi/stage-ui/stores/modules/airi-card'
import { useOnboardingStore } from '@proj-airi/stage-ui/stores/onboarding'
import { installChatContextBridge } from '@proj-airi/stage-ui/stores/plugins/chat-context-bridge'
import { useSettings } from '@proj-airi/stage-ui/stores/settings'
import { useTheme } from '@proj-airi/ui'
import { storeToRefs } from 'pinia'
@@ -17,6 +18,7 @@ import { themeColorFromValue, useThemeColor } from './composables/theme-color'
const { isDark: dark } = useTheme()
const i18n = useI18n()
const contextBridge = useContextBridge()
const displayModelsStore = useDisplayModelsStore()
const settingsStore = useSettings()
const { language, themeColorsHue, themeColorsHueDynamic } = storeToRefs(settingsStore)
@@ -24,7 +26,7 @@ const onboardingStore = useOnboardingStore()
const router = useRouter()
const route = useRoute()
const cardStore = useAiriCardStore()
let disposeChatBridge: (() => void) | undefined
const serverChannelStore = useModsServerChannelStore()
watch(language, () => {
i18n.locale.value = language.value
@@ -43,8 +45,8 @@ onMounted(async () => {
await displayModelsStore.loadDisplayModelsFromIndexedDB()
await settingsStore.initializeStageModel()
const bridge = installChatContextBridge()
disposeChatBridge = bridge.dispose
await serverChannelStore.initialize({ possibleEvents: ['ui:configure'] }).catch((err) => { console.error('Failed to initialize Mods Server Channel in App.vue:', err) })
await contextBridge.initialize()
const context = useElectronEventaContext()
const startTrackingCursorPoint = defineInvoke(context.value, electronStartTrackMousePosition)
@@ -62,7 +64,7 @@ watch(themeColorsHueDynamic, () => {
document.documentElement.classList.toggle('dynamic-hue', themeColorsHueDynamic.value)
}, { immediate: true })
onUnmounted(() => disposeChatBridge?.())
onUnmounted(() => contextBridge.dispose())
</script>
<template>
+9 -5
View File
@@ -1,9 +1,10 @@
<script setup lang="ts">
import { OnboardingDialog, ToasterRoot } from '@proj-airi/stage-ui/components'
import { useContextBridge } from '@proj-airi/stage-ui/composables'
import { useDisplayModelsStore } from '@proj-airi/stage-ui/stores/display-models'
import { useModsServerChannelStore } from '@proj-airi/stage-ui/stores/mods/api/channel-server'
import { useAiriCardStore } from '@proj-airi/stage-ui/stores/modules/airi-card'
import { useOnboardingStore } from '@proj-airi/stage-ui/stores/onboarding'
import { installChatContextBridge } from '@proj-airi/stage-ui/stores/plugins/chat-context-bridge'
import { useSettings } from '@proj-airi/stage-ui/stores/settings'
import { useTheme } from '@proj-airi/ui'
import { StageTransitionGroup } from '@proj-airi/ui-transitions'
@@ -18,15 +19,17 @@ import { usePWAStore } from './stores/pwa'
import 'vue-sonner/style.css'
usePWAStore()
const contextBridge = useContextBridge()
const i18n = useI18n()
const displayModelsStore = useDisplayModelsStore()
const settingsStore = useSettings()
const settings = storeToRefs(settingsStore)
const onboardingStore = useOnboardingStore()
const serverChannelStore = useModsServerChannelStore()
const { shouldShowSetup } = storeToRefs(onboardingStore)
const { isDark } = useTheme()
const cardStore = useAiriCardStore()
let disposeChatBridge: (() => void) | undefined
const primaryColor = computed(() => {
return isDark.value
@@ -67,15 +70,16 @@ onMounted(async () => {
cardStore.initialize()
onboardingStore.initializeSetupCheck()
const bridge = installChatContextBridge()
disposeChatBridge = bridge.dispose
await serverChannelStore.initialize({ possibleEvents: ['ui:configure'] }).catch((err) => { console.error('Failed to initialize Mods Server Channel in App.vue:', err) })
await contextBridge.initialize()
await displayModelsStore.loadDisplayModelsFromIndexedDB()
await settingsStore.initializeStageModel()
})
onUnmounted(() => {
disposeChatBridge?.()
contextBridge.dispose()
})
// Handle first-time setup events
@@ -1,12 +1,12 @@
<script setup lang="ts">
import type { ChatErrorMessage } from './types'
import type { ErrorMessage } from '../../../types/chat'
import { computed } from 'vue'
import MarkdownRenderer from '../../markdown/MarkdownRenderer.vue'
const props = withDefaults(defineProps<{
message: ChatErrorMessage
message: ErrorMessage
label: string
showPlaceholder?: boolean
variant?: 'desktop' | 'mobile'
@@ -1,6 +1,5 @@
<script setup lang="ts">
import type { ChatAssistantMessage } from '../../../types/chat'
import type { ChatHistoryMessage } from './types'
import type { ChatAssistantMessage, ChatHistoryItem, ContextMessage } from '../../../types/chat'
import { computed, onMounted, ref, watch } from 'vue'
import { useI18n } from 'vue-i18n'
@@ -10,8 +9,8 @@ import ChatErrorItem from './ChatErrorItem.vue'
import ChatUserItem from './ChatUserItem.vue'
const props = withDefaults(defineProps<{
messages: ChatHistoryMessage[]
streamingMessage?: ChatAssistantMessage & { context?: { ts?: number } }
messages: ChatHistoryItem[]
streamingMessage?: ChatAssistantMessage & { createdAt?: number }
sending?: boolean
assistantLabel?: string
userLabel?: string
@@ -46,10 +45,10 @@ watch([() => props.messages, () => props.streamingMessage], scrollToBottom, { de
watch(() => props.sending, scrollToBottom, { flush: 'post' })
onMounted(scrollToBottom)
const streaming = computed<ChatAssistantMessage & { context?: { ts?: number } }>(() => props.streamingMessage ?? { role: 'assistant', content: '', slices: [], tool_results: [] })
const streaming = computed<ChatAssistantMessage & { context?: ContextMessage } & { createdAt?: number }>(() => props.streamingMessage ?? { role: 'assistant', content: '', slices: [], tool_results: [], createdAt: Date.now() })
const showStreamingPlaceholder = computed(() => (streaming.value.slices?.length ?? 0) === 0 && !streaming.value.content)
const streamingTs = computed(() => streaming.value.context?.ts)
const renderMessages = computed<ChatHistoryMessage[]>(() => {
const streamingTs = computed(() => streaming.value?.createdAt)
const renderMessages = computed<ChatHistoryItem[]>(() => {
if (!props.sending)
return props.messages
@@ -57,7 +56,7 @@ const renderMessages = computed<ChatHistoryMessage[]>(() => {
if (!streamTs)
return props.messages
const hasStreamAlready = streamTs && props.messages.some(msg => msg.context?.ts === streamTs)
const hasStreamAlready = streamTs && props.messages.some(msg => msg?.createdAt === streamTs)
if (hasStreamAlready)
return props.messages
@@ -67,7 +66,7 @@ const renderMessages = computed<ChatHistoryMessage[]>(() => {
<template>
<div ref="chatHistoryRef" v-auto-animate flex="~ col" relative h-full w-full overflow-y-auto rounded-xl px="<sm:2" py="<sm:2" :class="variant === 'mobile' ? 'gap-1' : 'gap-2'">
<template v-for="(message, index) in renderMessages" :key="message.context?.ts ?? index">
<template v-for="(message, index) in renderMessages" :key="message?.createdAt ?? index">
<div v-if="message.role === 'error'">
<ChatErrorItem
:message="message"
@@ -81,7 +80,7 @@ const renderMessages = computed<ChatHistoryMessage[]>(() => {
<ChatAssistantItem
:message="message"
:label="labels.assistant"
:show-placeholder="message.context?.ts === streamingTs ? showStreamingPlaceholder : false"
:show-placeholder="message.context?.createdAt === streamingTs ? showStreamingPlaceholder : false"
:variant="variant"
/>
</div>
@@ -2,5 +2,3 @@ export { default as ChatAssistantItem } from './ChatAssistantItem.vue'
export { default as ChatErrorItem } from './ChatErrorItem.vue'
export { default as ChatHistory } from './ChatHistory.vue'
export { default as ChatUserItem } from './ChatUserItem.vue'
export type { ChatErrorMessage, ChatHistoryMessage } from './types'
@@ -1,14 +0,0 @@
import type { ChatMessage, ChatSlices } from '../../../types/chat'
export interface ChatErrorMessage {
role: 'error'
content: string
}
export type ChatHistoryMessage = (ChatMessage | ChatErrorMessage) & {
slices?: ChatSlices[]
context?: {
ts?: number
[key: string]: unknown
}
}
@@ -4,4 +4,5 @@ export * from './llmmarkerParser'
export * from './markdown'
export * from './micvad'
export * from './queues'
export * from './use-context-bridge'
export * from './whisper'
@@ -0,0 +1,144 @@
import type { ChatStreamEvent, ContextMessage } from '../types/chat'
import { useBroadcastChannel } from '@vueuse/core'
import { Mutex } from 'es-toolkit'
import { watch } from 'vue'
import { CHAT_STREAM_CHANNEL_NAME, CONTEXT_CHANNEL_NAME, useChatStore } from '../stores/chat'
import { useModsServerChannelStore } from '../stores/mods/api/channel-server'
const mutex = new Mutex()
export function useContextBridge() {
const chatStore = useChatStore()
const serverChannelStore = useModsServerChannelStore()
const { post: broadcastContext, data: incomingContext } = useBroadcastChannel<ContextMessage, ContextMessage>({ name: CONTEXT_CHANNEL_NAME })
const { post: broadcastStreamEvent, data: incomingStreamEvent } = useBroadcastChannel<ChatStreamEvent, ChatStreamEvent>({ name: CHAT_STREAM_CHANNEL_NAME })
let disposeHookFns = [] as Array<() => void>
return {
initialize: async () => {
await mutex.acquire()
try {
let isProcessingRemoteStream = false
const { stop } = watch(incomingContext, (event) => {
if (event)
chatStore.ingestContextMessage(event)
})
disposeHookFns.push(stop)
disposeHookFns.push(serverChannelStore.onContextUpdate((event) => {
chatStore.ingestContextMessage({ source: event.source, createdAt: Date.now(), ...event.data })
broadcastContext(event.data as ContextMessage)
}))
disposeHookFns.push(...[
chatStore.onBeforeMessageComposed(async (message) => {
if (isProcessingRemoteStream)
return
broadcastStreamEvent({ type: 'before-compose', message, sessionId: chatStore.activeSessionId })
}),
chatStore.onAfterMessageComposed(async (message) => {
if (isProcessingRemoteStream)
return
broadcastStreamEvent({ type: 'after-compose', message, sessionId: chatStore.activeSessionId })
}),
chatStore.onBeforeSend(async (message) => {
if (isProcessingRemoteStream)
return
broadcastStreamEvent({ type: 'before-send', message, sessionId: chatStore.activeSessionId })
}),
chatStore.onAfterSend(async (message) => {
if (isProcessingRemoteStream)
return
broadcastStreamEvent({ type: 'after-send', message, sessionId: chatStore.activeSessionId })
}),
chatStore.onTokenLiteral(async (literal) => {
if (isProcessingRemoteStream)
return
broadcastStreamEvent({ type: 'token-literal', literal, sessionId: chatStore.activeSessionId })
}),
chatStore.onTokenSpecial(async (special) => {
if (isProcessingRemoteStream)
return
broadcastStreamEvent({ type: 'token-special', special, sessionId: chatStore.activeSessionId })
}),
chatStore.onStreamEnd(async () => {
if (isProcessingRemoteStream)
return
broadcastStreamEvent({ type: 'stream-end', sessionId: chatStore.activeSessionId })
}),
chatStore.onAssistantResponseEnd(async (message) => {
if (isProcessingRemoteStream)
return
broadcastStreamEvent({ type: 'assistant-end', message, sessionId: chatStore.activeSessionId })
}),
])
const { stop: stopIncomingStreamWatch } = watch(incomingStreamEvent, async (event) => {
if (!event)
return
isProcessingRemoteStream = true
try {
if (event.sessionId && chatStore.activeSessionId !== event.sessionId)
chatStore.setActiveSession(event.sessionId)
switch (event.type) {
case 'before-compose':
await chatStore.emitBeforeMessageComposedHooks(event.message)
break
case 'after-compose':
await chatStore.emitAfterMessageComposedHooks(event.message)
break
case 'before-send':
await chatStore.emitBeforeSendHooks(event.message)
break
case 'after-send':
await chatStore.emitAfterSendHooks(event.message)
break
case 'token-literal':
await chatStore.emitTokenLiteralHooks(event.literal)
break
case 'token-special':
await chatStore.emitTokenSpecialHooks(event.special)
break
case 'stream-end':
await chatStore.emitStreamEndHooks()
break
case 'assistant-end':
await chatStore.emitAssistantResponseEndHooks(event.message)
break
}
}
finally {
isProcessingRemoteStream = false
}
})
disposeHookFns.push(stopIncomingStreamWatch)
}
finally {
mutex.release()
}
},
dispose: async () => {
await mutex.acquire()
try {
for (const fn of disposeHookFns) {
fn()
}
}
finally {
mutex.release()
}
disposeHookFns = []
},
}
}
-123
View File
@@ -1,123 +0,0 @@
import type { ContextMessage } from '@proj-airi/server-sdk'
import type { ContextPayload } from './chat'
import { createPinia, setActivePinia } from 'pinia'
import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest'
import { ref } from 'vue'
import { useChatStore } from './chat'
import { installChatContextBridge } from './plugins/chat-context-bridge'
const mockSendContextUpdate = vi.fn()
const mockInitialize = vi.fn().mockResolvedValue(undefined)
let contextUpdateHandler: ((event: { type: 'context:update', data: ContextMessage }) => void | Promise<void>) | null = null
let broadcastPosts: unknown[] = []
let bridge: { dispose: () => void } | null = null
const localStorageMap = new Map<string, unknown>()
vi.mock('@vueuse/core', () => {
return {
useLocalStorage: <T>(key: string, defaultValue: T) => {
if (!localStorageMap.has(key))
localStorageMap.set(key, ref(defaultValue))
return localStorageMap.get(key) as ReturnType<typeof ref<T>>
},
useBroadcastChannel: <T>() => {
const data = ref<T | undefined>()
const post = (value: T) => {
broadcastPosts.push(value)
data.value = value
}
return { data, post }
},
}
})
vi.mock('./llm', () => ({
useLLM: () => ({
stream: vi.fn(),
discoverToolsCompatibility: vi.fn(),
}),
}))
vi.mock('./modules', () => ({
useAiriCardStore: () => ({
systemPrompt: ref(''),
}),
}))
vi.mock('./mods/api/channel-server', () => ({
useModsServerChannelStore: () => ({
connected: ref(true),
initialize: mockInitialize,
onContextUpdate: (cb: typeof contextUpdateHandler) => {
contextUpdateHandler = cb
return () => {
contextUpdateHandler = null
}
},
sendContextUpdate: mockSendContextUpdate,
}),
}))
describe('chat store', () => {
beforeEach(() => {
setActivePinia(createPinia())
broadcastPosts = []
localStorageMap.clear()
mockSendContextUpdate.mockClear()
mockInitialize.mockClear()
contextUpdateHandler = null
bridge = null
})
afterEach(() => {
bridge?.dispose()
bridge = null
vi.clearAllMocks()
})
it('ingests assistant context updates from the channel server', async () => {
const store = useChatStore()
bridge = installChatContextBridge()
expect(mockInitialize).toHaveBeenCalled()
expect(contextUpdateHandler).toBeTruthy()
const envelope: ContextMessage = {
sessionId: 'session-ctx',
ts: 123,
role: 'assistant',
source: 'llm',
payload: { content: 'hello from server' },
}
await contextUpdateHandler?.({ type: 'context:update', data: envelope })
store.setActiveSession('session-ctx')
const last = store.messages.at(-1)
expect(last?.role).toBe('assistant')
expect(last?.content).toBe('hello from server')
})
it('publishes local context updates through the shared channel server store', () => {
const store = useChatStore()
bridge = installChatContextBridge()
const envelope: ContextMessage = {
sessionId: 'session-local',
ts: 456,
role: 'assistant',
source: 'system',
payload: { content: 'local broadcast' },
}
store.publishContextMessage(envelope as ContextMessage<ContextPayload, Record<string, unknown>>, 'local')
expect(mockSendContextUpdate).toHaveBeenCalledWith(envelope)
expect(broadcastPosts).toContain(envelope)
})
})
+56 -163
View File
@@ -1,10 +1,10 @@
import type { ContextMessage, ContextSource } from '@proj-airi/server-sdk'
import type { ChatProvider } from '@xsai-ext/shared-providers'
import type { CommonContentPart, Message, SystemMessage } from '@xsai/shared-chat'
import type { StreamEvent, StreamOptions } from '../stores/llm'
import type { ChatAssistantMessage, ChatMessage, ChatSlices } from '../types/chat'
import type { ChatAssistantMessage, ChatHistoryItem, ChatSlices, ContextMessage, StreamingAssistantMessage } from '../types/chat'
import { ContextUpdateStrategy } from '@proj-airi/server-sdk'
import { useLocalStorage } from '@vueuse/core'
import { defineStore, storeToRefs } from 'pinia'
import { computed, ref, toRaw, watch } from 'vue'
@@ -15,55 +15,24 @@ import { createQueue } from '../utils/queue'
import { TTS_FLUSH_INSTRUCTION } from '../utils/tts'
import { useAiriCardStore } from './modules'
export interface ErrorMessage {
role: 'error'
content: string
}
interface MessageContext {
sessionId: string
source: ContextSource
ts: number
meta?: Record<string, unknown>
}
export type ChatEntry = (ChatMessage | ErrorMessage) & { context?: MessageContext }
export interface ContextPayload {
content?: unknown
slices?: ChatSlices[]
tool_results?: ChatAssistantMessage['tool_results']
text?: string
}
export type ChatStreamEvent
= | { type: 'before-compose', message: string, sessionId: string }
| { type: 'after-compose', message: string, sessionId: string }
| { type: 'before-send', message: string, sessionId: string }
| { type: 'after-send', message: string, sessionId: string }
| { type: 'token-literal', literal: string, sessionId: string }
| { type: 'token-special', special: string, sessionId: string }
| { type: 'stream-end', sessionId: string }
| { type: 'assistant-end', message: string, sessionId: string }
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'
type StreamingAssistantMessage = ChatAssistantMessage & { context?: MessageContext }
export const useChatStore = defineStore('chat', () => {
const { stream, discoverToolsCompatibility } = useLLM()
const { systemPrompt } = storeToRefs(useAiriCardStore())
const activeSessionId = useLocalStorage<string>(ACTIVE_SESSION_STORAGE_KEY, 'default')
const sessionMessages = useLocalStorage<Record<string, ChatEntry[]>>(CHAT_STORAGE_KEY, {})
const sessionMessages = useLocalStorage<Record<string, ChatHistoryItem[]>>(CHAT_STORAGE_KEY, {})
const sending = ref(false)
const streamingMessage = ref<StreamingAssistantMessage>({ role: 'assistant', content: '', slices: [], tool_results: [] })
const streamingMessage = ref<StreamingAssistantMessage>({ role: 'assistant', content: '', slices: [], tool_results: [], createdAt: Date.now() })
const sessionGenerations = ref<Record<string, number>>({})
const activeContexts = ref<Record<string, ContextMessage[]>>({})
interface SendOptions {
model: string
chatProvider: ChatProvider
@@ -127,7 +96,7 @@ export const useChatStore = defineStore('chat', () => {
const onTokenSpecialHooks = ref<Array<(special: string) => Promise<void>>>([])
const onStreamEndHooks = ref<Array<() => Promise<void>>>([])
const onAssistantResponseEndHooks = ref<Array<(message: string) => Promise<void>>>([])
const onContextPublishHooks = ref<Array<(envelope: ContextMessage<ContextPayload>, origin: 'local' | 'ws' | 'broadcast') => Promise<void> | void>>([])
const onContextPublishHooks = ref<Array<(envelope: ContextMessage, origin: 'local' | 'ws' | 'broadcast') => Promise<void> | void>>([])
function onBeforeMessageComposed(cb: (message: string) => Promise<void>) {
onBeforeMessageComposedHooks.value.push(cb)
@@ -169,14 +138,6 @@ export const useChatStore = defineStore('chat', () => {
return () => onAssistantResponseEndHooks.value = onAssistantResponseEndHooks.value.filter(hook => hook !== cb) // return remove listener callback
}
function onContextPublish(cb: (envelope: ContextMessage<ContextPayload>, origin: 'local' | 'ws' | 'broadcast') => Promise<void> | void) {
onContextPublishHooks.value.push(cb)
return () => {
onContextPublishHooks.value = onContextPublishHooks.value.filter(hook => hook !== cb)
}
}
function clearHooks() {
onBeforeMessageComposedHooks.value = []
onAfterMessageComposedHooks.value = []
@@ -252,23 +213,19 @@ export const useChatStore = defineStore('chat', () => {
function generateInitialMessage() {
// TODO: compose, replace {{ user }} tag, etc
const content = codeBlockSystemPrompt + mathSyntaxSystemPrompt + systemPrompt.value
return {
role: 'system',
content: codeBlockSystemPrompt + mathSyntaxSystemPrompt + systemPrompt.value,
content,
} satisfies SystemMessage
}
function ensureSession(sessionId: string) {
ensureSessionGeneration(sessionId)
if (!sessionMessages.value[sessionId] || sessionMessages.value[sessionId].length === 0) {
sessionMessages.value[sessionId] = [{
...generateInitialMessage(),
context: {
sessionId,
source: 'system',
ts: Date.now(),
},
}]
sessionMessages.value[sessionId] = [generateInitialMessage()]
}
}
@@ -279,7 +236,7 @@ export const useChatStore = defineStore('chat', () => {
return sessionMessages.value[sessionId]!
}
const messages = computed<ChatEntry[]>({
const messages = computed<ChatHistoryItem[]>({
get: () => {
ensureSession(activeSessionId.value)
return sessionMessages.value[activeSessionId.value]
@@ -296,14 +253,8 @@ export const useChatStore = defineStore('chat', () => {
function cleanupMessages(sessionId = activeSessionId.value) {
bumpSessionGeneration(sessionId)
sessionMessages.value[sessionId] = [{
...generateInitialMessage(),
context: {
sessionId,
source: 'system',
ts: Date.now(),
},
}]
sessionMessages.value[sessionId] = [generateInitialMessage()]
// Reject pending sends for this session so callers don't hang after cleanup
for (const queued of pendingQueuedSends.value) {
if (queued.sessionId !== sessionId)
@@ -312,16 +263,17 @@ export const useChatStore = defineStore('chat', () => {
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<string, ChatEntry[]>
return JSON.parse(JSON.stringify(toRaw(sessionMessages.value))) as Record<string, ChatHistoryItem[]>
}
function replaceSessions(sessions: Record<string, ChatEntry[]>) {
function replaceSessions(sessions: Record<string, ChatHistoryItem[]>) {
sessionMessages.value = sessions
sessionGenerations.value = Object.fromEntries(Object.keys(sessions).map(sessionId => [sessionId, 0]))
const [firstSessionId] = Object.keys(sessions)
@@ -341,74 +293,22 @@ export const useChatStore = defineStore('chat', () => {
watch(systemPrompt, () => {
for (const [sessionId, history] of Object.entries(sessionMessages.value)) {
if (history.length > 0 && history[0].role === 'system') {
sessionMessages.value[sessionId][0] = {
...generateInitialMessage(),
context: {
sessionId,
source: 'system',
ts: Date.now(),
},
}
sessionMessages.value[sessionId][0] = generateInitialMessage()
}
}
}, { immediate: true })
// ----- Context bridge (WS + BroadcastChannel) -----
function normalizePayload(payload?: ContextPayload) {
const baseContent = payload?.content ?? payload?.text ?? ''
const normalizedContent = typeof baseContent === 'string' || Array.isArray(baseContent)
? baseContent
: JSON.stringify(baseContent)
return {
content: normalizedContent,
slices: payload?.slices ?? [],
tool_results: payload?.tool_results ?? [],
}
}
function ingestContextMessage(envelope: ContextMessage<ContextPayload>) {
ensureSession(envelope.sessionId)
const { content, slices, tool_results } = normalizePayload(envelope.payload)
const context: MessageContext = {
sessionId: envelope.sessionId,
source: envelope.source,
ts: envelope.ts,
meta: envelope.meta,
function ingestContextMessage(envelope: ContextMessage) {
if (!activeContexts.value[envelope.source]) {
activeContexts.value[envelope.source] = []
}
const nextHistory = sessionMessages.value[envelope.sessionId]
if (envelope.role === 'assistant') {
nextHistory.push({
role: 'assistant',
content,
slices,
tool_results,
context,
})
if (envelope.strategy === ContextUpdateStrategy.ReplaceSelf) {
activeContexts.value[envelope.source] = [envelope]
}
else if (envelope.role === 'error') {
nextHistory.push({
role: 'error',
content: typeof content === 'string' ? content : JSON.stringify(content),
context,
})
else if (envelope.strategy === ContextUpdateStrategy.AppendSelf) {
activeContexts.value[envelope.source].push(envelope)
}
else {
nextHistory.push({
role: envelope.role,
content,
context,
} as ChatEntry)
}
}
function publishContextMessage(envelope: ContextMessage<ContextPayload>, origin: 'local' | 'ws' | 'broadcast' = 'local') {
for (const hook of onContextPublishHooks.value)
void hook(envelope, origin)
}
// ----- Send flow (user -> LLM -> assistant) -----
@@ -427,13 +327,10 @@ export const useChatStore = defineStore('chat', () => {
const shouldAbort = () => isStaleGeneration()
if (shouldAbort())
return
sending.value = true
const assistantContext: MessageContext = {
sessionId,
source: 'llm',
ts: Date.now(),
}
streamingMessage.value = { role: 'assistant', content: '', slices: [], tool_results: [], context: assistantContext }
streamingMessage.value = { role: 'assistant', content: '', slices: [], tool_results: [], createdAt: Date.now() }
try {
await emitBeforeMessageComposedHooks(sendingMessage)
@@ -459,16 +356,7 @@ export const useChatStore = defineStore('chat', () => {
return
const sessionMessagesForSend = getSessionMessagesById(sessionId)
const userContext: MessageContext = { sessionId, source: 'text', ts: Date.now() }
sessionMessagesForSend.push({ role: 'user', content: finalContent, context: userContext })
publishContextMessage({
sessionId: userContext.sessionId,
ts: userContext.ts,
role: 'user',
source: userContext.source,
payload: { content: finalContent },
}, 'local')
sessionMessagesForSend.push({ role: 'user', content: finalContent })
const parser = useLlmmarkerParser({
onLiteral: async (literal) => {
@@ -493,6 +381,7 @@ export const useChatStore = defineStore('chat', () => {
onSpecial: async (special) => {
if (shouldAbort())
return
await emitTokenSpecialHooks(special)
},
minLiteralEmitLength: 24, // Avoid emitting literals too fast. This is a magic number and can be changed later.
@@ -515,9 +404,10 @@ export const useChatStore = defineStore('chat', () => {
],
})
const newMessages = sessionMessagesForSend.map((msg) => {
let newMessages = sessionMessagesForSend.map((msg) => {
const { context: _context, ...withoutContext } = msg
const rawMessage = toRaw(withoutContext)
if (rawMessage.role === 'assistant') {
const { slices: _, tool_results, ...rest } = rawMessage as ChatAssistantMessage
return {
@@ -529,6 +419,28 @@ 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 system = newMessages.slice(0, 1)
const afterSystem = newMessages.slice(1, newMessages.length)
newMessages = [
...system,
{
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` },
],
},
...afterSystem,
]
}
await emitAfterMessageComposedHooks(sendingMessage)
await emitBeforeSendHooks(sendingMessage)
@@ -573,24 +485,7 @@ export const useChatStore = defineStore('chat', () => {
// Add the completed message to the history only if it has content
if (!isStaleGeneration() && streamingMessage.value.slices.length > 0) {
const assistantMessage: ChatEntry = {
...(toRaw(streamingMessage.value) as ChatAssistantMessage),
context: assistantContext,
}
sessionMessagesForSend.push(assistantMessage)
publishContextMessage({
sessionId: assistantContext.sessionId,
ts: assistantContext.ts,
role: 'assistant',
source: assistantContext.source,
payload: {
content: assistantMessage.content,
slices: assistantMessage.slices,
tool_results: assistantMessage.tool_results,
},
}, 'local')
sessionMessagesForSend.push(toRaw(streamingMessage.value))
}
// Reset the streaming message for the next turn
@@ -649,7 +544,6 @@ export const useChatStore = defineStore('chat', () => {
send,
setActiveSession,
ingestContextMessage,
publishContextMessage,
cleanupMessages,
getAllSessions,
replaceSessions,
@@ -672,6 +566,5 @@ export const useChatStore = defineStore('chat', () => {
onTokenSpecial,
onStreamEnd,
onAssistantResponseEnd,
onContextPublish,
}
})
@@ -1,6 +1,7 @@
import type { ContextMessage, WebSocketBaseEvent, WebSocketEvent, WebSocketEvents } from '@proj-airi/server-sdk'
import type { ContextUpdate, WebSocketBaseEvent, WebSocketEvent, WebSocketEventOptionalSource, WebSocketEvents } from '@proj-airi/server-sdk'
import { Client } from '@proj-airi/server-sdk'
import { Client, WebSocketEventSource } from '@proj-airi/server-sdk'
import { isStageTamagotchi, isStageWeb } from '@proj-airi/stage-shared'
import { defineStore } from 'pinia'
import { ref } from 'vue'
@@ -8,7 +9,6 @@ export const useModsServerChannelStore = defineStore('mods:channels:proj-airi:se
const connected = ref(false)
const client = ref<Client>()
const initializing = ref<Promise<void> | null>(null)
const pendingSend = ref<Array<WebSocketEvent>>([])
function initialize(options?: { token?: string, possibleEvents?: Array<keyof WebSocketEvents> }) {
@@ -19,13 +19,12 @@ export const useModsServerChannelStore = defineStore('mods:channels:proj-airi:se
const possibleEvents = Array.from(new Set<keyof WebSocketEvents>([
'ui:configure',
'context:update',
...(options?.possibleEvents ?? []),
]))
initializing.value = new Promise<void>((resolve, reject) => {
client.value = new Client({
name: 'proj-airi:ui:stage',
name: isStageWeb() ? WebSocketEventSource.StageWeb : isStageTamagotchi() ? WebSocketEventSource.StageTamagotchi : WebSocketEventSource.StageWeb,
url: import.meta.env.VITE_AIRI_WS_URL || 'ws://localhost:6121/ws',
token: options?.token,
possibleEvents,
@@ -47,6 +46,7 @@ export const useModsServerChannelStore = defineStore('mods:channels:proj-airi:se
flush()
initializeListeners()
resolve()
return
}
@@ -64,15 +64,15 @@ export const useModsServerChannelStore = defineStore('mods:channels:proj-airi:se
return
}
function send(data: WebSocketEvent) {
function send<C = undefined>(data: WebSocketEventOptionalSource<C>) {
if (!client.value && !initializing.value)
void initialize()
if (client.value && connected.value) {
client.value.send(data)
client.value.send(data as WebSocketEvent)
}
else {
pendingSend.value.push(data)
pendingSend.value.push(data as WebSocketEvent)
}
}
@@ -86,7 +86,7 @@ export const useModsServerChannelStore = defineStore('mods:channels:proj-airi:se
}
}
function onContextUpdate(callback: (event: WebSocketBaseEvent<'context:update', ContextMessage>) => void | Promise<void>) {
function onContextUpdate(callback: (event: WebSocketBaseEvent<'context:update', ContextUpdate>) => void | Promise<void>) {
if (!client.value && !initializing.value)
void initialize()
@@ -97,11 +97,8 @@ export const useModsServerChannelStore = defineStore('mods:channels:proj-airi:se
}
}
function sendContextUpdate(message: ContextMessage) {
send({
type: 'context:update',
data: message,
})
function sendContextUpdate(message: ContextUpdate) {
send({ type: 'context:update', data: message })
}
function dispose() {
@@ -1,147 +0,0 @@
import type { ContextMessage } from '@proj-airi/server-sdk'
import type { ChatStreamEvent, ContextPayload } from '../chat'
import { useBroadcastChannel } from '@vueuse/core'
import { watch } from 'vue'
import { CHAT_STREAM_CHANNEL_NAME, CONTEXT_CHANNEL_NAME, useChatStore } from '../chat'
import { useModsServerChannelStore } from '../mods/api/channel-server'
let installed = false
export function installChatContextBridge() {
if (installed) {
return {
dispose: () => {},
}
}
const chatStore = useChatStore()
const modsChannelServer = useModsServerChannelStore()
const { post: broadcastContext, data: incomingContext } = useBroadcastChannel<ContextMessage<ContextPayload>, ContextMessage<ContextPayload>>({
name: CONTEXT_CHANNEL_NAME,
})
const { post: broadcastStreamEvent, data: incomingStreamEvent } = useBroadcastChannel<ChatStreamEvent, ChatStreamEvent>({
name: CHAT_STREAM_CHANNEL_NAME,
})
let isProcessingRemoteStream = false
const stopStreamBroadcastHooks = [
chatStore.onBeforeMessageComposed(async (message) => {
if (isProcessingRemoteStream)
return
broadcastStreamEvent({ type: 'before-compose', message, sessionId: chatStore.activeSessionId })
}),
chatStore.onAfterMessageComposed(async (message) => {
if (isProcessingRemoteStream)
return
broadcastStreamEvent({ type: 'after-compose', message, sessionId: chatStore.activeSessionId })
}),
chatStore.onBeforeSend(async (message) => {
if (isProcessingRemoteStream)
return
broadcastStreamEvent({ type: 'before-send', message, sessionId: chatStore.activeSessionId })
}),
chatStore.onAfterSend(async (message) => {
if (isProcessingRemoteStream)
return
broadcastStreamEvent({ type: 'after-send', message, sessionId: chatStore.activeSessionId })
}),
chatStore.onTokenLiteral(async (literal) => {
if (isProcessingRemoteStream)
return
broadcastStreamEvent({ type: 'token-literal', literal, sessionId: chatStore.activeSessionId })
}),
chatStore.onTokenSpecial(async (special) => {
if (isProcessingRemoteStream)
return
broadcastStreamEvent({ type: 'token-special', special, sessionId: chatStore.activeSessionId })
}),
chatStore.onStreamEnd(async () => {
if (isProcessingRemoteStream)
return
broadcastStreamEvent({ type: 'stream-end', sessionId: chatStore.activeSessionId })
}),
chatStore.onAssistantResponseEnd(async (message) => {
if (isProcessingRemoteStream)
return
broadcastStreamEvent({ type: 'assistant-end', message, sessionId: chatStore.activeSessionId })
}),
]
const stopIncomingWatch = watch(incomingContext, (event) => {
if (event)
chatStore.ingestContextMessage(event)
})
const stopIncomingStreamWatch = watch(incomingStreamEvent, async (event) => {
if (!event)
return
isProcessingRemoteStream = true
try {
if (event.sessionId && chatStore.activeSessionId !== event.sessionId)
chatStore.setActiveSession(event.sessionId)
switch (event.type) {
case 'before-compose':
await chatStore.emitBeforeMessageComposedHooks(event.message)
break
case 'after-compose':
await chatStore.emitAfterMessageComposedHooks(event.message)
break
case 'before-send':
await chatStore.emitBeforeSendHooks(event.message)
break
case 'after-send':
await chatStore.emitAfterSendHooks(event.message)
break
case 'token-literal':
await chatStore.emitTokenLiteralHooks(event.literal)
break
case 'token-special':
await chatStore.emitTokenSpecialHooks(event.special)
break
case 'stream-end':
await chatStore.emitStreamEndHooks()
break
case 'assistant-end':
await chatStore.emitAssistantResponseEndHooks(event.message)
break
}
}
finally {
isProcessingRemoteStream = false
}
})
const offPublish = chatStore.onContextPublish((envelope, origin) => {
if (origin !== 'broadcast')
broadcastContext(envelope)
if (origin === 'local')
modsChannelServer.sendContextUpdate(envelope)
})
modsChannelServer.initialize({ possibleEvents: ['context:update'] }).catch(error => console.error('Context bridge init error:', error))
const offWs = modsChannelServer.onContextUpdate((event) => {
const envelope = event.data as ContextMessage<ContextPayload, Record<string, unknown>>
chatStore.ingestContextMessage(envelope)
broadcastContext(envelope)
})
installed = true
return {
dispose: () => {
stopIncomingWatch()
stopIncomingStreamWatch()
offPublish()
offWs?.()
stopStreamBroadcastHooks.forEach(stop => stop())
installed = false
},
}
}
+25
View File
@@ -1,3 +1,4 @@
import type { ContextUpdate, WebSocketEventSource } from '@proj-airi/server-sdk'
import type { AssistantMessage, CommonContentPart, CompletionToolCall, SystemMessage, ToolMessage, UserMessage } from '@xsai/shared-chat'
export interface ChatSlicesText {
@@ -27,3 +28,27 @@ export interface ChatAssistantMessage extends AssistantMessage {
}
export type ChatMessage = ChatAssistantMessage | SystemMessage | ToolMessage | UserMessage
export interface ErrorMessage {
role: 'error'
content: string
}
export interface ContextMessage extends ContextUpdate {
source: WebSocketEventSource | string
createdAt: number
}
export type ChatHistoryItem = (ChatMessage | ErrorMessage) & { context?: ContextMessage } & { createdAt?: number }
export type ChatStreamEvent
= | { type: 'before-compose', message: string, sessionId: string }
| { type: 'after-compose', message: string, sessionId: string }
| { type: 'before-send', message: string, sessionId: string }
| { type: 'after-send', message: string, sessionId: string }
| { type: 'token-literal', literal: string, sessionId: string }
| { type: 'token-special', special: string, sessionId: string }
| { type: 'stream-end', sessionId: string }
| { type: 'assistant-end', message: string, sessionId: string }
export type StreamingAssistantMessage = ChatAssistantMessage & { context?: ContextMessage } & { createdAt?: number }