feat(core-agent): extract pure runtime logic from stage-ui into core-agent (#1524)

Depends on #1487 

## Summary

- Introduce `packages/core-agent` — a new pure-runtime package that
extracts zero-Vue/Pinia algorithmic logic from `@proj-airi/stage-ui`,
following a "compatible facade + core sinking" strategy
- Migrate shared chat types, LLM streaming types, chat hook registry,
session message merge logic, context registry algorithm, and LLM service
utilities into `core-agent`
- `stage-ui` files are preserved as re-exports / thin wrappers — all
external imports, store IDs, and public APIs remain unchanged

## Motivation

`stage-ui` currently mixes pure agent runtime logic (algorithms, type
definitions, stateless utilities) with Vue/Pinia state management and
browser-specific adapters. This coupling makes it hard to:
- Test agent logic in isolation
- Reuse agent algorithms outside of Vue contexts (e.g., server-side,
CLI, other frameworks)
- Reason about the boundary between "what the agent does" vs "how the UI
manages state"

---------

Co-authored-by: autofix-ci[bot] <114827586+autofix-ci[bot]@users.noreply.github.com>
Co-authored-by-agent: Unknown <unknown@example.com>
This commit is contained in:
Nashchennc
2026-04-16 16:09:22 +08:00
committed by GitHub
co-authored by autofix-ci[bot]
parent 11ac511dd9
commit bf41655dbb
22 changed files with 1399 additions and 659 deletions
+46
View File
@@ -0,0 +1,46 @@
{
"name": "@proj-airi/core-agent",
"type": "module",
"version": "0.9.0-alpha.35",
"private": true,
"description": "Core agent runtime orchestration for AIRI",
"author": {
"name": "Moeru AI Project AIRI Team",
"email": "airi@moeru.ai",
"url": "https://github.com/moeru-ai"
},
"license": "MIT",
"repository": {
"type": "git",
"url": "https://github.com/moeru-ai/airi.git",
"directory": "packages/core-agent"
},
"exports": {
".": {
"types": "./dist/index.d.mts",
"default": "./dist/index.mjs"
}
},
"main": "./dist/index.mjs",
"types": "./dist/index.d.mts",
"files": [
"README.md",
"dist",
"package.json"
],
"scripts": {
"build": "tsdown",
"typecheck": "tsc --noEmit"
},
"dependencies": {
"@proj-airi/server-shared": "workspace:^",
"@xsai-ext/providers": "catalog:",
"@xsai/model": "catalog:",
"@xsai/shared-chat": "catalog:",
"@xsai/stream-text": "catalog:"
},
"devDependencies": {
"tsdown": "catalog:",
"typescript": "^5.9.2"
}
}
@@ -0,0 +1,7 @@
import type { ContextMessage } from '../types/chat'
export interface AgentContextPort {
ingest: (envelope: ContextMessage) => void
snapshot: () => Record<string, ContextMessage[]>
reset: () => void
}
@@ -0,0 +1,55 @@
import type { ToolMessage } from '@xsai/shared-chat'
import type { ChatStreamEventContext, StreamingAssistantMessage } from '../types/chat'
export interface ChatHookRegistry {
onBeforeMessageComposed: (cb: (message: string, context: Omit<ChatStreamEventContext, 'composedMessage'>) => Promise<void>) => () => void
onAfterMessageComposed: (cb: (message: string, context: ChatStreamEventContext) => Promise<void>) => () => void
onBeforeSend: (cb: (message: string, context: ChatStreamEventContext) => Promise<void>) => () => void
onAfterSend: (cb: (message: string, context: ChatStreamEventContext) => Promise<void>) => () => void
onTokenLiteral: (cb: (literal: string, context: ChatStreamEventContext) => Promise<void>) => () => void
onTokenSpecial: (cb: (special: string, context: ChatStreamEventContext) => Promise<void>) => () => void
onStreamEnd: (cb: (context: ChatStreamEventContext) => Promise<void>) => () => void
onAssistantResponseEnd: (cb: (message: string, context: ChatStreamEventContext) => Promise<void>) => () => void
onAssistantMessage: (cb: (message: StreamingAssistantMessage, messageText: string, context: ChatStreamEventContext) => Promise<void>) => () => void
onChatTurnComplete: (cb: (chat: { output: StreamingAssistantMessage, outputText: string, toolCalls: ToolMessage[] }, context: ChatStreamEventContext) => Promise<void>) => () => void
emitBeforeMessageComposedHooks: (message: string, context: Omit<ChatStreamEventContext, 'composedMessage'>) => Promise<void>
emitAfterMessageComposedHooks: (message: string, context: ChatStreamEventContext) => Promise<void>
emitBeforeSendHooks: (message: string, context: ChatStreamEventContext) => Promise<void>
emitAfterSendHooks: (message: string, context: ChatStreamEventContext) => Promise<void>
emitTokenLiteralHooks: (literal: string, context: ChatStreamEventContext) => Promise<void>
emitTokenSpecialHooks: (special: string, context: ChatStreamEventContext) => Promise<void>
emitStreamEndHooks: (context: ChatStreamEventContext) => Promise<void>
emitAssistantResponseEndHooks: (message: string, context: ChatStreamEventContext) => Promise<void>
emitAssistantMessageHooks: (message: StreamingAssistantMessage, messageText: string, context: ChatStreamEventContext) => Promise<void>
emitChatTurnCompleteHooks: (chat: { output: StreamingAssistantMessage, outputText: string, toolCalls: ToolMessage[] }, context: ChatStreamEventContext) => Promise<void>
clearHooks: () => void
}
export interface HookUnsubscribe {
(): void
}
export interface AgentHookRegistry<TContext, TAssistantMessage, TToolCall> {
onBeforeMessageComposed: (cb: (message: string, context: Omit<TContext, 'composedMessage'>) => Promise<void>) => HookUnsubscribe
onAfterMessageComposed: (cb: (message: string, context: TContext) => Promise<void>) => HookUnsubscribe
onBeforeSend: (cb: (message: string, context: TContext) => Promise<void>) => HookUnsubscribe
onAfterSend: (cb: (message: string, context: TContext) => Promise<void>) => HookUnsubscribe
onTokenLiteral: (cb: (literal: string, context: TContext) => Promise<void>) => HookUnsubscribe
onTokenSpecial: (cb: (special: string, context: TContext) => Promise<void>) => HookUnsubscribe
onStreamEnd: (cb: (context: TContext) => Promise<void>) => HookUnsubscribe
onAssistantResponseEnd: (cb: (message: string, context: TContext) => Promise<void>) => HookUnsubscribe
onAssistantMessage: (cb: (message: TAssistantMessage, messageText: string, context: TContext) => Promise<void>) => HookUnsubscribe
onChatTurnComplete: (cb: (chat: { output: TAssistantMessage, outputText: string, toolCalls: TToolCall[] }, context: TContext) => Promise<void>) => HookUnsubscribe
emitBeforeMessageComposedHooks: (message: string, context: Omit<TContext, 'composedMessage'>) => Promise<void>
emitAfterMessageComposedHooks: (message: string, context: TContext) => Promise<void>
emitBeforeSendHooks: (message: string, context: TContext) => Promise<void>
emitAfterSendHooks: (message: string, context: TContext) => Promise<void>
emitTokenLiteralHooks: (literal: string, context: TContext) => Promise<void>
emitTokenSpecialHooks: (special: string, context: TContext) => Promise<void>
emitStreamEndHooks: (context: TContext) => Promise<void>
emitAssistantResponseEndHooks: (message: string, context: TContext) => Promise<void>
emitAssistantMessageHooks: (message: TAssistantMessage, messageText: string, context: TContext) => Promise<void>
emitChatTurnCompleteHooks: (chat: { output: TAssistantMessage, outputText: string, toolCalls: TToolCall[] }, context: TContext) => Promise<void>
clearHooks: () => void
}
@@ -0,0 +1,8 @@
import type { ChatProvider } from '@xsai-ext/providers/utils'
import type { Message } from '@xsai/shared-chat'
import type { StreamOptions } from '../types/llm'
export interface AgentLLMPort {
stream: (model: string, chatProvider: ChatProvider, messages: Message[], options?: StreamOptions) => Promise<void>
}
@@ -0,0 +1,8 @@
import type { ChatHistoryItem } from '../types/chat'
export interface AgentSessionPort {
ensureSession: (sessionId: string) => void
getSessionMessages: (sessionId: string) => ChatHistoryItem[]
appendSessionMessage: (sessionId: string, message: ChatHistoryItem) => void
getSessionGeneration: (sessionId: string) => number
}
@@ -0,0 +1,6 @@
import type { StreamingAssistantMessage } from '../types/chat'
export interface AgentForegroundStreamPort {
patch: (message: StreamingAssistantMessage) => void
reset: () => void
}
+40
View File
@@ -0,0 +1,40 @@
export type { AgentContextPort } from './contracts/context-port'
export type { ChatHookRegistry } from './contracts/hook-types'
export type { AgentLLMPort } from './contracts/llm-port'
export type { AgentSessionPort } from './contracts/session-port'
export type { AgentForegroundStreamPort } from './contracts/stream-port'
export { createChatHooks } from './runtime/agent-hooks'
export type { ContextHistoryEntry, ContextRegistry } from './runtime/context-registry'
export { createContextRegistry } from './runtime/context-registry'
export {
isToolRelatedError,
modelKey,
sanitizeMessages,
streamFrom,
streamOptionsToolsCompatibilityOk,
} from './runtime/llm-service'
export { mergeLoadedSessionMessages } from './session/merge-loaded-session-messages'
export type {
ChatAssistantMessage,
ChatHistoryItem,
ChatMessage,
ChatSlices,
ChatSlicesText,
ChatSlicesToolCall,
ChatSlicesToolCallResult,
ChatStreamEvent,
ChatStreamEventContext,
ContextMessage,
ErrorMessage,
StreamingAssistantMessage,
} from './types/chat'
export type {
BuiltinToolsResolver,
StreamEvent,
StreamFromOptions,
StreamOptions,
} from './types/llm'
@@ -0,0 +1,259 @@
import type { ToolMessage } from '@xsai/shared-chat'
import type { AgentHookRegistry, ChatHookRegistry } from '../contracts/hook-types'
import type { ChatStreamEventContext, StreamingAssistantMessage } from '../types/chat'
export function createChatHooks(): ChatHookRegistry {
const onBeforeMessageComposedHooks: Array<(message: string, context: Omit<ChatStreamEventContext, 'composedMessage'>) => Promise<void>> = []
const onAfterMessageComposedHooks: Array<(message: string, context: ChatStreamEventContext) => Promise<void>> = []
const onBeforeSendHooks: Array<(message: string, context: ChatStreamEventContext) => Promise<void>> = []
const onAfterSendHooks: Array<(message: string, context: ChatStreamEventContext) => Promise<void>> = []
const onTokenLiteralHooks: Array<(literal: string, context: ChatStreamEventContext) => Promise<void>> = []
const onTokenSpecialHooks: Array<(special: string, context: ChatStreamEventContext) => Promise<void>> = []
const onStreamEndHooks: Array<(context: ChatStreamEventContext) => Promise<void>> = []
const onAssistantResponseEndHooks: Array<(message: string, context: ChatStreamEventContext) => Promise<void>> = []
const onAssistantMessageHooks: Array<(message: StreamingAssistantMessage, messageText: string, context: ChatStreamEventContext) => Promise<void>> = []
const onChatTurnCompleteHooks: Array<(chat: { output: StreamingAssistantMessage, outputText: string, toolCalls: ToolMessage[] }, context: ChatStreamEventContext) => Promise<void>> = []
function onBeforeMessageComposed(cb: (message: string, context: Omit<ChatStreamEventContext, 'composedMessage'>) => Promise<void>) {
onBeforeMessageComposedHooks.push(cb)
return () => {
const index = onBeforeMessageComposedHooks.indexOf(cb)
if (index >= 0)
onBeforeMessageComposedHooks.splice(index, 1)
}
}
function onAfterMessageComposed(cb: (message: string, context: ChatStreamEventContext) => Promise<void>) {
onAfterMessageComposedHooks.push(cb)
return () => {
const index = onAfterMessageComposedHooks.indexOf(cb)
if (index >= 0)
onAfterMessageComposedHooks.splice(index, 1)
}
}
function onBeforeSend(cb: (message: string, context: ChatStreamEventContext) => Promise<void>) {
onBeforeSendHooks.push(cb)
return () => {
const index = onBeforeSendHooks.indexOf(cb)
if (index >= 0)
onBeforeSendHooks.splice(index, 1)
}
}
function onAfterSend(cb: (message: string, context: ChatStreamEventContext) => Promise<void>) {
onAfterSendHooks.push(cb)
return () => {
const index = onAfterSendHooks.indexOf(cb)
if (index >= 0)
onAfterSendHooks.splice(index, 1)
}
}
function onTokenLiteral(cb: (literal: string, context: ChatStreamEventContext) => Promise<void>) {
onTokenLiteralHooks.push(cb)
return () => {
const index = onTokenLiteralHooks.indexOf(cb)
if (index >= 0)
onTokenLiteralHooks.splice(index, 1)
}
}
function onTokenSpecial(cb: (special: string, context: ChatStreamEventContext) => Promise<void>) {
onTokenSpecialHooks.push(cb)
return () => {
const index = onTokenSpecialHooks.indexOf(cb)
if (index >= 0)
onTokenSpecialHooks.splice(index, 1)
}
}
function onStreamEnd(cb: (context: ChatStreamEventContext) => Promise<void>) {
onStreamEndHooks.push(cb)
return () => {
const index = onStreamEndHooks.indexOf(cb)
if (index >= 0)
onStreamEndHooks.splice(index, 1)
}
}
function onAssistantResponseEnd(cb: (message: string, context: ChatStreamEventContext) => Promise<void>) {
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<void>) {
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<void>) {
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<ChatStreamEventContext, 'composedMessage'>) {
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,
}
}
export function createAgentHooks<TContext, TAssistantMessage, TToolCall>(): AgentHookRegistry<TContext, TAssistantMessage, TToolCall> {
const onBeforeMessageComposedHooks: Array<(message: string, context: Omit<TContext, 'composedMessage'>) => Promise<void>> = []
const onAfterMessageComposedHooks: Array<(message: string, context: TContext) => Promise<void>> = []
const onBeforeSendHooks: Array<(message: string, context: TContext) => Promise<void>> = []
const onAfterSendHooks: Array<(message: string, context: TContext) => Promise<void>> = []
const onTokenLiteralHooks: Array<(literal: string, context: TContext) => Promise<void>> = []
const onTokenSpecialHooks: Array<(special: string, context: TContext) => Promise<void>> = []
const onStreamEndHooks: Array<(context: TContext) => Promise<void>> = []
const onAssistantResponseEndHooks: Array<(message: string, context: TContext) => Promise<void>> = []
const onAssistantMessageHooks: Array<(message: TAssistantMessage, messageText: string, context: TContext) => Promise<void>> = []
const onChatTurnCompleteHooks: Array<(chat: { output: TAssistantMessage, outputText: string, toolCalls: TToolCall[] }, context: TContext) => Promise<void>> = []
function createSubscribe<T>(bucket: T[], cb: T) {
bucket.push(cb)
return () => {
const index = bucket.indexOf(cb)
if (index >= 0)
bucket.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 emitHooks<T extends any[]>(hooks: Array<(...args: T) => Promise<void>>, ...args: T) {
for (const hook of hooks)
await hook(...args)
}
return {
onBeforeMessageComposed: cb => createSubscribe(onBeforeMessageComposedHooks, cb),
onAfterMessageComposed: cb => createSubscribe(onAfterMessageComposedHooks, cb),
onBeforeSend: cb => createSubscribe(onBeforeSendHooks, cb),
onAfterSend: cb => createSubscribe(onAfterSendHooks, cb),
onTokenLiteral: cb => createSubscribe(onTokenLiteralHooks, cb),
onTokenSpecial: cb => createSubscribe(onTokenSpecialHooks, cb),
onStreamEnd: cb => createSubscribe(onStreamEndHooks, cb),
onAssistantResponseEnd: cb => createSubscribe(onAssistantResponseEndHooks, cb),
onAssistantMessage: cb => createSubscribe(onAssistantMessageHooks, cb),
onChatTurnComplete: cb => createSubscribe(onChatTurnCompleteHooks, cb),
emitBeforeMessageComposedHooks: (message, context) => emitHooks(onBeforeMessageComposedHooks, message, context),
emitAfterMessageComposedHooks: (message, context) => emitHooks(onAfterMessageComposedHooks, message, context),
emitBeforeSendHooks: (message, context) => emitHooks(onBeforeSendHooks, message, context),
emitAfterSendHooks: (message, context) => emitHooks(onAfterSendHooks, message, context),
emitTokenLiteralHooks: (literal, context) => emitHooks(onTokenLiteralHooks, literal, context),
emitTokenSpecialHooks: (special, context) => emitHooks(onTokenSpecialHooks, special, context),
emitStreamEndHooks: context => emitHooks(onStreamEndHooks, context),
emitAssistantResponseEndHooks: (message, context) => emitHooks(onAssistantResponseEndHooks, message, context),
emitAssistantMessageHooks: (message, messageText, context) => emitHooks(onAssistantMessageHooks, message, messageText, context),
emitChatTurnCompleteHooks: (chat, context) => emitHooks(onChatTurnCompleteHooks, chat, context),
clearHooks,
}
}
@@ -0,0 +1,95 @@
import type { MetadataEventSource } from '@proj-airi/server-shared/types'
import type { ContextMessage } from '../types/chat'
const CONTEXT_UPDATE_REPLACE_SELF = 'replace-self'
const CONTEXT_UPDATE_APPEND_SELF = 'append-self'
interface EventSourcePayload {
source?: string
metadata?: { source?: MetadataEventSource }
}
export interface ContextHistoryEntry extends ContextMessage {
sourceKey: string
}
export interface ContextRegistry {
ingest: (envelope: ContextMessage) => void
reset: () => void
snapshot: () => Record<string, ContextMessage[]>
activeContexts: () => Record<string, ContextMessage[]>
contextHistory: () => ContextHistoryEntry[]
}
interface CreateContextRegistryOptions {
historyLimit?: number
getSourceKey?: (event: EventSourcePayload, fallback?: string) => string
}
function formatMetadataSource(source?: MetadataEventSource) {
if (!source?.plugin)
return undefined
const pluginId = source.plugin.id
const instanceId = source.id
return instanceId ? `${pluginId}:${instanceId}` : pluginId
}
function defaultGetSourceKey(event: EventSourcePayload, fallback = 'unknown') {
return (
formatMetadataSource(event.metadata?.source)
?? event.source
?? fallback
)
}
export function createContextRegistry(options: CreateContextRegistryOptions = {}): ContextRegistry {
const historyLimit = options.historyLimit ?? 400
const getSourceKey = options.getSourceKey ?? defaultGetSourceKey
let currentActiveContexts: Record<string, ContextMessage[]> = {}
let currentContextHistory: ContextHistoryEntry[] = []
function ingest(envelope: ContextMessage) {
const sourceKey = getSourceKey(envelope)
if (!currentActiveContexts[sourceKey]) {
currentActiveContexts[sourceKey] = []
}
const safeEnvelopeToStore = structuredClone(envelope)
if (envelope.strategy === CONTEXT_UPDATE_REPLACE_SELF) {
currentActiveContexts[sourceKey] = [safeEnvelopeToStore]
}
else if (envelope.strategy === CONTEXT_UPDATE_APPEND_SELF) {
currentActiveContexts[sourceKey].push(safeEnvelopeToStore)
}
currentContextHistory = [
...currentContextHistory,
{
...safeEnvelopeToStore,
sourceKey,
},
].slice(-historyLimit)
}
function reset() {
currentActiveContexts = {}
currentContextHistory = []
}
function snapshot() {
return structuredClone(currentActiveContexts)
}
return {
ingest,
reset,
snapshot,
activeContexts: () => structuredClone(currentActiveContexts),
contextHistory: () => [...currentContextHistory],
}
}
@@ -0,0 +1,145 @@
import type { ChatProvider } from '@xsai-ext/providers/utils'
import type { Message } from '@xsai/shared-chat'
import type { StreamFromOptions, StreamOptions } from '../types/llm'
import { stepCountAtLeast } from '@xsai/shared-chat'
import { streamText } from '@xsai/stream-text'
export function sanitizeMessages(messages: unknown[]): Message[] {
return messages.map((message: any) => {
if (message && message.role === 'error') {
return {
role: 'user',
content: `User encountered error: ${String(message.content ?? '')}`,
} as Message
}
// NOTICE: Flatten array content for providers (e.g. DeepSeek) that expect string,
// not content-part arrays. Skipped when image_url parts are present.
if (message && Array.isArray(message.content)) {
const contentParts = message.content as { type?: string, text?: string }[]
if (!contentParts.some(part => part?.type === 'image_url')) {
return { ...message, content: contentParts.map(part => part?.text ?? '').join('') } as Message
}
}
return message as Message
})
}
export function modelKey(model: string, chatProvider: ChatProvider): string {
return `${chatProvider.chat(model).baseURL}-${model}`
}
export function streamOptionsToolsCompatibilityOk(model: string, chatProvider: ChatProvider, options?: StreamOptions): boolean {
if (options?.supportsTools)
return true
const key = modelKey(model, chatProvider)
return options?.toolsCompatibility?.get(key) !== false
}
async function resolveTools(options?: StreamOptions) {
const tools = typeof options?.tools === 'function'
? await options.tools()
: options?.tools
return tools ?? []
}
export async function streamFrom({
model,
chatProvider,
messages,
options,
builtinToolsResolver,
}: StreamFromOptions) {
const chatConfig = chatProvider.chat(model)
const sanitized = sanitizeMessages(messages as unknown[])
const supportedTools = streamOptionsToolsCompatibilityOk(model, chatProvider, options)
const builtinTools = supportedTools
? await (builtinToolsResolver?.(model, chatProvider) ?? Promise.resolve([]))
: []
const customTools = supportedTools ? await resolveTools(options) : []
const mergedTools = supportedTools ? [...builtinTools, ...customTools] : []
const tools = mergedTools.length > 0 ? mergedTools : undefined
return new Promise<void>((resolve, reject) => {
let settled = false
const resolveOnce = () => {
if (settled)
return
settled = true
resolve()
}
const rejectOnce = (error: unknown) => {
if (settled)
return
settled = true
reject(error)
}
const onEvent = async (event: unknown) => {
try {
await options?.onStreamEvent?.(event as any)
if (event && (event as any).type === 'finish') {
const finishReason = (event as any).finishReason
const waitingForToolRound = finishReason === 'tool_calls' || finishReason === 'tool-calls'
if (!waitingForToolRound || !options?.waitForTools)
resolveOnce()
}
else if (event && (event as any).type === 'error') {
rejectOnce((event as any).error ?? new Error('Stream error'))
}
}
catch (error) {
rejectOnce(error)
}
}
try {
const streamResult = streamText({
...chatConfig,
abortSignal: options?.abortSignal,
messages: sanitized,
headers: options?.headers,
stopWhen: stepCountAtLeast(10),
tools,
captureToolErrors: true,
onEvent,
})
// NOTICE: Consume underlying promises to prevent unhandled rejections from
// @xsai/stream-text's SSE parser surfacing as faulted app state.
void streamResult.steps.catch((error) => {
rejectOnce(error)
console.error('Stream steps error:', error)
})
void streamResult.messages.catch(error => console.error('Stream messages error:', error))
void streamResult.usage.catch(error => console.error('Stream usage error:', error))
void streamResult.totalUsage.catch(error => console.error('Stream totalUsage error:', error))
}
catch (error) {
rejectOnce(error)
}
})
}
// Runtime auto-degrade: patterns that indicate the model/provider does not support tool calling.
const TOOLS_RELATED_ERROR_PATTERNS: RegExp[] = [
/does not support tools/i, // Ollama
/no endpoints found that support tool use/i, // OpenRouter
/invalid schema for function/i, // OpenAI-compatible
/invalid.?function.?parameters/i, // OpenAI-compatible
/functions are not supported/i, // Azure AI Foundry
/unrecognized request argument.+tools/i, // Azure AI Foundry
/tool use with function calling is unsupported/i, // Google Generative AI
/tool_use_failed/i, // Groq
/does not support function.?calling/i, // Anthropic
/tools?\s+(is|are)\s+not\s+supported/i, // Cloudflare Workers AI
]
export function isToolRelatedError(error: unknown): boolean {
const message = String(error)
return TOOLS_RELATED_ERROR_PATTERNS.some(pattern => pattern.test(message))
}
@@ -0,0 +1,59 @@
import type { ChatHistoryItem } from '../types/chat'
function extractMessageContent(message: ChatHistoryItem) {
if (typeof message.content === 'string')
return message.content
if (Array.isArray(message.content)) {
return message.content.map((part) => {
if (typeof part === 'string')
return part
if (part && typeof part === 'object' && 'text' in part)
return String(part.text ?? '')
return ''
}).join('')
}
return ''
}
function getMessageFingerprint(message: ChatHistoryItem) {
return [
message.id ?? '',
message.role,
message.createdAt ?? '',
extractMessageContent(message),
].join('\u001F')
}
export function mergeLoadedSessionMessages(storedMessages: ChatHistoryItem[], currentMessages: ChatHistoryItem[]) {
if (currentMessages.length === 0)
return storedMessages
const currentNonSystemMessages = currentMessages.filter((message, index) => index !== 0 || message.role !== 'system')
if (currentNonSystemMessages.length === 0)
return storedMessages
const seen = new Set(storedMessages.map(getMessageFingerprint))
const extraMessages = currentNonSystemMessages.filter((message) => {
const fingerprint = getMessageFingerprint(message)
if (seen.has(fingerprint))
return false
seen.add(fingerprint)
return true
})
if (extraMessages.length === 0)
return storedMessages
const systemMessage = storedMessages[0]?.role === 'system'
? storedMessages[0]
: currentMessages[0]?.role === 'system'
? currentMessages[0]
: undefined
if (storedMessages.length === 0 && systemMessage)
return [systemMessage, ...extraMessages]
return [...storedMessages, ...extraMessages]
}
+70
View File
@@ -0,0 +1,70 @@
import type { ContextUpdate, MetadataEventSource, WebSocketEventInputs } from '@proj-airi/server-shared/types'
import type { AssistantMessage, CommonContentPart, CompletionToolCall, Message, SystemMessage, ToolMessage, UserMessage } from '@xsai/shared-chat'
export interface ChatSlicesText {
type: 'text'
text: string
}
export interface ChatSlicesToolCall {
type: 'tool-call'
toolCall: CompletionToolCall
}
export interface ChatSlicesToolCallResult {
type: 'tool-call-result'
id: string
isError?: boolean
result?: string | CommonContentPart[]
}
export type ChatSlices = ChatSlicesText | ChatSlicesToolCall | ChatSlicesToolCallResult
export interface ChatAssistantMessage extends AssistantMessage {
slices: ChatSlices[]
tool_results: {
id: string
isError?: boolean
result?: string | CommonContentPart[]
}[]
categorization?: {
speech: string
reasoning: string
}
}
export type ChatMessage = ChatAssistantMessage | SystemMessage | ToolMessage | UserMessage
export interface ErrorMessage {
role: 'error'
content: string
}
export interface ContextMessage extends ContextUpdate<Record<string, unknown>, unknown> {
metadata?: {
source: MetadataEventSource
}
createdAt: number
}
export type ChatHistoryItem = (ChatMessage | ErrorMessage) & { context?: ContextMessage } & { createdAt?: number, id?: string }
export interface ChatStreamEventContext {
message: ChatHistoryItem
contexts: Record<string, ContextMessage[]>
composedMessage: Array<Message>
input?: WebSocketEventInputs
}
export type ChatStreamEvent
= | { type: 'before-compose', message: string, sessionId: string, context: Omit<ChatStreamEventContext, 'composedMessage'> }
| { type: 'after-compose', message: string, sessionId: string, context: ChatStreamEventContext }
| { type: 'before-send', message: string, sessionId: string, context: ChatStreamEventContext }
| { type: 'after-send', message: string, sessionId: string, context: ChatStreamEventContext }
| { type: 'token-literal', literal: string, sessionId: string, context: ChatStreamEventContext }
| { type: 'token-special', special: string, sessionId: string, context: ChatStreamEventContext }
| { type: 'stream-end', sessionId: string, context: ChatStreamEventContext }
| { type: 'assistant-end', message: string, sessionId: string, context: ChatStreamEventContext }
| { type: 'assistant-message', message: ChatAssistantMessage, sessionId: string, messageText: string, context: ChatStreamEventContext }
export type StreamingAssistantMessage = ChatAssistantMessage & { context?: ContextMessage } & { createdAt?: number, id?: string }
+30
View File
@@ -0,0 +1,30 @@
import type { ChatProvider } from '@xsai-ext/providers/utils'
import type { CommonContentPart, CompletionToolCall, CompletionToolResult, Message, Tool } from '@xsai/shared-chat'
export type StreamEvent
= | { type: 'text-delta', text: string }
| ({ type: 'finish' } & any)
| ({ type: 'tool-call' } & CompletionToolCall)
| (CompletionToolResult & { type: 'tool-error' })
| { type: 'tool-result', toolCallId: string, result?: string | CommonContentPart[] }
| { type: 'error', error: any }
export interface StreamOptions {
abortSignal?: AbortSignal
headers?: Record<string, string>
onStreamEvent?: (event: StreamEvent) => void | Promise<void>
toolsCompatibility?: Map<string, boolean>
supportsTools?: boolean
waitForTools?: boolean
tools?: Tool[] | (() => Promise<Tool[] | undefined>)
}
export type BuiltinToolsResolver = (model: string, chatProvider: ChatProvider) => Promise<Tool[]>
export interface StreamFromOptions {
model: string
chatProvider: ChatProvider
messages: Message[]
options?: StreamOptions
builtinToolsResolver?: BuiltinToolsResolver
}
+14
View File
@@ -0,0 +1,14 @@
{
"compilerOptions": {
"target": "ESNext",
"lib": ["ESNext", "DOM"],
"module": "ESNext",
"moduleResolution": "bundler",
"esModuleInterop": true,
"forceConsistentCasingInFileNames": true,
"isolatedModules": true,
"verbatimModuleSyntax": true,
"skipLibCheck": true
},
"include": ["src/**/*.ts"]
}
+8
View File
@@ -0,0 +1,8 @@
import { defineConfig } from 'tsdown'
export default defineConfig({
entry: [
'src/index.ts',
],
dts: true,
})
+1
View File
@@ -69,6 +69,7 @@
"@proj-airi/audio": "workspace:^",
"@proj-airi/ccc": "workspace:^",
"@proj-airi/chromatic": "^1.1.1",
"@proj-airi/core-agent": "workspace:^",
"@proj-airi/core-character": "workspace:^",
"@proj-airi/drizzle-duckdb-wasm": "catalog:",
"@proj-airi/font-chillroundm": "workspace:^",
+1 -217
View File
@@ -1,217 +1 @@
import type { ToolMessage } from '@xsai/shared-chat'
import type { ChatStreamEventContext, StreamingAssistantMessage } from '../../types/chat'
export interface ChatHookRegistry {
onBeforeMessageComposed: (cb: (message: string, context: Omit<ChatStreamEventContext, 'composedMessage'>) => Promise<void>) => () => void
onAfterMessageComposed: (cb: (message: string, context: ChatStreamEventContext) => Promise<void>) => () => void
onBeforeSend: (cb: (message: string, context: ChatStreamEventContext) => Promise<void>) => () => void
onAfterSend: (cb: (message: string, context: ChatStreamEventContext) => Promise<void>) => () => void
onTokenLiteral: (cb: (literal: string, context: ChatStreamEventContext) => Promise<void>) => () => void
onTokenSpecial: (cb: (special: string, context: ChatStreamEventContext) => Promise<void>) => () => void
onStreamEnd: (cb: (context: ChatStreamEventContext) => Promise<void>) => () => void
onAssistantResponseEnd: (cb: (message: string, context: ChatStreamEventContext) => Promise<void>) => () => void
onAssistantMessage: (cb: (message: StreamingAssistantMessage, messageText: string, context: ChatStreamEventContext) => Promise<void>) => () => void
onChatTurnComplete: (cb: (chat: { output: StreamingAssistantMessage, outputText: string, toolCalls: ToolMessage[] }, context: ChatStreamEventContext) => Promise<void>) => () => void
emitBeforeMessageComposedHooks: (message: string, context: Omit<ChatStreamEventContext, 'composedMessage'>) => Promise<void>
emitAfterMessageComposedHooks: (message: string, context: ChatStreamEventContext) => Promise<void>
emitBeforeSendHooks: (message: string, context: ChatStreamEventContext) => Promise<void>
emitAfterSendHooks: (message: string, context: ChatStreamEventContext) => Promise<void>
emitTokenLiteralHooks: (literal: string, context: ChatStreamEventContext) => Promise<void>
emitTokenSpecialHooks: (special: string, context: ChatStreamEventContext) => Promise<void>
emitStreamEndHooks: (context: ChatStreamEventContext) => Promise<void>
emitAssistantResponseEndHooks: (message: string, context: ChatStreamEventContext) => Promise<void>
emitAssistantMessageHooks: (message: StreamingAssistantMessage, messageText: string, context: ChatStreamEventContext) => Promise<void>
emitChatTurnCompleteHooks: (chat: { output: StreamingAssistantMessage, outputText: string, toolCalls: ToolMessage[] }, context: ChatStreamEventContext) => Promise<void>
clearHooks: () => void
}
export function createChatHooks(): ChatHookRegistry {
const onBeforeMessageComposedHooks: Array<(message: string, context: Omit<ChatStreamEventContext, 'composedMessage'>) => Promise<void>> = []
const onAfterMessageComposedHooks: Array<(message: string, context: ChatStreamEventContext) => Promise<void>> = []
const onBeforeSendHooks: Array<(message: string, context: ChatStreamEventContext) => Promise<void>> = []
const onAfterSendHooks: Array<(message: string, context: ChatStreamEventContext) => Promise<void>> = []
const onTokenLiteralHooks: Array<(literal: string, context: ChatStreamEventContext) => Promise<void>> = []
const onTokenSpecialHooks: Array<(special: string, context: ChatStreamEventContext) => Promise<void>> = []
const onStreamEndHooks: Array<(context: ChatStreamEventContext) => Promise<void>> = []
const onAssistantResponseEndHooks: Array<(message: string, context: ChatStreamEventContext) => Promise<void>> = []
const onAssistantMessageHooks: Array<(message: StreamingAssistantMessage, messageText: string, context: ChatStreamEventContext) => Promise<void>> = []
const onChatTurnCompleteHooks: Array<(chat: { output: StreamingAssistantMessage, outputText: string, toolCalls: ToolMessage[] }, context: ChatStreamEventContext) => Promise<void>> = []
function onBeforeMessageComposed(cb: (message: string, context: Omit<ChatStreamEventContext, 'composedMessage'>) => Promise<void>) {
onBeforeMessageComposedHooks.push(cb)
return () => {
const index = onBeforeMessageComposedHooks.indexOf(cb)
if (index >= 0)
onBeforeMessageComposedHooks.splice(index, 1)
}
}
function onAfterMessageComposed(cb: (message: string, context: ChatStreamEventContext) => Promise<void>) {
onAfterMessageComposedHooks.push(cb)
return () => {
const index = onAfterMessageComposedHooks.indexOf(cb)
if (index >= 0)
onAfterMessageComposedHooks.splice(index, 1)
}
}
function onBeforeSend(cb: (message: string, context: ChatStreamEventContext) => Promise<void>) {
onBeforeSendHooks.push(cb)
return () => {
const index = onBeforeSendHooks.indexOf(cb)
if (index >= 0)
onBeforeSendHooks.splice(index, 1)
}
}
function onAfterSend(cb: (message: string, context: ChatStreamEventContext) => Promise<void>) {
onAfterSendHooks.push(cb)
return () => {
const index = onAfterSendHooks.indexOf(cb)
if (index >= 0)
onAfterSendHooks.splice(index, 1)
}
}
function onTokenLiteral(cb: (literal: string, context: ChatStreamEventContext) => Promise<void>) {
onTokenLiteralHooks.push(cb)
return () => {
const index = onTokenLiteralHooks.indexOf(cb)
if (index >= 0)
onTokenLiteralHooks.splice(index, 1)
}
}
function onTokenSpecial(cb: (special: string, context: ChatStreamEventContext) => Promise<void>) {
onTokenSpecialHooks.push(cb)
return () => {
const index = onTokenSpecialHooks.indexOf(cb)
if (index >= 0)
onTokenSpecialHooks.splice(index, 1)
}
}
function onStreamEnd(cb: (context: ChatStreamEventContext) => Promise<void>) {
onStreamEndHooks.push(cb)
return () => {
const index = onStreamEndHooks.indexOf(cb)
if (index >= 0)
onStreamEndHooks.splice(index, 1)
}
}
function onAssistantResponseEnd(cb: (message: string, context: ChatStreamEventContext) => Promise<void>) {
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<void>) {
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<void>) {
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<ChatStreamEventContext, 'composedMessage'>) {
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,
}
}
export { type ChatHookRegistry, createChatHooks } from '@proj-airi/core-agent'
@@ -1,59 +1 @@
import type { ChatHistoryItem } from '../../types/chat'
function extractMessageContent(message: ChatHistoryItem) {
if (typeof message.content === 'string')
return message.content
if (Array.isArray(message.content)) {
return message.content.map((part) => {
if (typeof part === 'string')
return part
if (part && typeof part === 'object' && 'text' in part)
return String(part.text ?? '')
return ''
}).join('')
}
return ''
}
function getMessageFingerprint(message: ChatHistoryItem) {
return [
message.id ?? '',
message.role,
message.createdAt ?? '',
extractMessageContent(message),
].join('\u001F')
}
export function mergeLoadedSessionMessages(storedMessages: ChatHistoryItem[], currentMessages: ChatHistoryItem[]) {
if (currentMessages.length === 0)
return storedMessages
const currentNonSystemMessages = currentMessages.filter((message, index) => index !== 0 || message.role !== 'system')
if (currentNonSystemMessages.length === 0)
return storedMessages
const seen = new Set(storedMessages.map(getMessageFingerprint))
const extraMessages = currentNonSystemMessages.filter((message) => {
const fingerprint = getMessageFingerprint(message)
if (seen.has(fingerprint))
return false
seen.add(fingerprint)
return true
})
if (extraMessages.length === 0)
return storedMessages
const systemMessage = storedMessages[0]?.role === 'system'
? storedMessages[0]
: currentMessages[0]?.role === 'system'
? currentMessages[0]
: undefined
if (storedMessages.length === 0 && systemMessage)
return [systemMessage, ...extraMessages]
return [...storedMessages, ...extraMessages]
}
export { mergeLoadedSessionMessages } from '@proj-airi/core-agent'
@@ -31,6 +31,8 @@ describe('stage-ui exports contract', () => {
'./constants/*',
'./libs',
'./libs/*',
'./libs/inference',
'./libs/inference/adapters/*',
'./stores',
'./stores/*',
'./stores/analytics',
+29 -174
View File
@@ -1,198 +1,53 @@
import type { StreamOptions } from '@proj-airi/core-agent'
import type { WebSocketEvents } from '@proj-airi/server-sdk'
import type { ChatProvider } from '@xsai-ext/providers/utils'
import type { CommonContentPart, CompletionToolCall, CompletionToolResult, Message, Tool } from '@xsai/shared-chat'
import type { Message } from '@xsai/shared-chat'
import { streamFrom as coreStreamFrom, isToolRelatedError, modelKey } from '@proj-airi/core-agent'
import { listModels } from '@xsai/model'
import {
stepCountAtLeast,
} from '@xsai/shared-chat'
import { streamText } from '@xsai/stream-text'
import { defineStore } from 'pinia'
import { ref } from 'vue'
import { createSparkCommandTool, debug, mcp } from '../tools'
import { useModsServerChannelStore } from './mods/api/channel-server'
export type StreamEvent
= | { type: 'text-delta', text: string }
| ({ type: 'finish' } & any)
| ({ type: 'tool-call' } & CompletionToolCall)
| (CompletionToolResult & { type: 'tool-error' })
| { type: 'tool-result', toolCallId: string, result?: string | CommonContentPart[] }
| { type: 'error', error: any }
export interface StreamOptions {
abortSignal?: AbortSignal
headers?: Record<string, string>
onStreamEvent?: (event: StreamEvent) => void | Promise<void>
toolsCompatibility?: Map<string, boolean>
supportsTools?: boolean
waitForTools?: boolean
tools?: Tool[] | (() => Promise<Tool[] | undefined>)
}
function sanitizeMessages(messages: unknown[]): Message[] {
return messages.map((m: any) => {
if (m && m.role === 'error') {
return {
role: 'user',
content: `User encountered error: ${String(m.content ?? '')}`,
} as Message
}
// NOTICE: Flatten array content for providers (e.g. DeepSeek) that expect string,
// not content-part arrays. Skipped when image_url parts are present.
if (m && Array.isArray(m.content)) {
const contentParts = m.content as { type?: string, text?: string }[]
if (!contentParts.some(p => p?.type === 'image_url')) {
return { ...m, content: contentParts.map(p => p?.text ?? '').join('') } as Message
}
}
return m as Message
})
}
function streamOptionsToolsCompatibilityOk(model: string, chatProvider: ChatProvider, _: Message[], options?: StreamOptions): boolean {
if (options?.supportsTools)
return true
const key = `${chatProvider.chat(model).baseURL}-${model}`
return options?.toolsCompatibility?.get(key) !== false
}
async function streamFrom(model: string, chatProvider: ChatProvider, messages: Message[], sendSparkCommand: (command: WebSocketEvents['spark:command']) => void, options?: StreamOptions) {
const chatConfig = chatProvider.chat(model)
const sanitized = sanitizeMessages(messages as unknown[])
const resolveTools = async () => {
const tools = typeof options?.tools === 'function'
? await options.tools()
: options?.tools
return tools ?? []
}
const supportedTools = streamOptionsToolsCompatibilityOk(model, chatProvider, messages, options)
const tools = supportedTools
? [
...await mcp(),
...await debug(),
...await resolveTools(),
await createSparkCommandTool({ sendSparkCommand }),
]
: undefined
return new Promise<void>((resolve, reject) => {
let settled = false
const resolveOnce = () => {
if (settled)
return
settled = true
resolve()
}
const rejectOnce = (err: unknown) => {
if (settled)
return
settled = true
reject(err)
}
const onEvent = async (event: unknown) => {
try {
await options?.onStreamEvent?.(event as StreamEvent)
if (event && (event as StreamEvent).type === 'finish') {
const finishReason = (event as any).finishReason
const waitingForToolRound = finishReason === 'tool_calls' || finishReason === 'tool-calls'
if (!waitingForToolRound || !options?.waitForTools)
resolveOnce()
}
else if (event && (event as StreamEvent).type === 'error') {
rejectOnce((event as any).error ?? new Error('Stream error'))
}
}
catch (err) {
rejectOnce(err)
}
}
try {
const streamResult = streamText({
...chatConfig,
abortSignal: options?.abortSignal,
messages: sanitized,
headers: options?.headers,
stopWhen: stepCountAtLeast(10),
tools,
captureToolErrors: true,
onEvent,
})
// NOTICE: Consume underlying promises to prevent unhandled rejections from
// @xsai/stream-text's SSE parser surfacing as faulted app state.
void streamResult.steps.catch((err) => {
rejectOnce(err)
console.error('Stream steps error:', err)
})
void streamResult.messages.catch(err => console.error('Stream messages error:', err))
void streamResult.usage.catch(err => console.error('Stream usage error:', err))
void streamResult.totalUsage.catch(err => console.error('Stream totalUsage error:', err))
}
catch (err) {
rejectOnce(err)
}
})
}
// Runtime auto-degrade: patterns that indicate the model/provider does not support tool calling.
const TOOLS_RELATED_ERROR_PATTERNS: RegExp[] = [
/does not support tools/i, // Ollama
/no endpoints found that support tool use/i, // OpenRouter
/invalid schema for function/i, // OpenAI-compatible
/invalid.?function.?parameters/i, // OpenAI-compatible
/functions are not supported/i, // Azure AI Foundry
/unrecognized request argument.+tools/i, // Azure AI Foundry
/tool use with function calling is unsupported/i, // Google Generative AI
/tool_use_failed/i, // Groq
/does not support function.?calling/i, // Anthropic
/tools?\s+(is|are)\s+not\s+supported/i, // Cloudflare Workers AI
]
export function isToolRelatedError(err: unknown): boolean {
const msg = String(err)
return TOOLS_RELATED_ERROR_PATTERNS.some(p => p.test(msg))
}
export type { StreamEvent, StreamOptions } from '@proj-airi/core-agent'
export { isToolRelatedError } from '@proj-airi/core-agent'
export const useLLM = defineStore('llm', () => {
const toolsCompatibility = ref<Map<string, boolean>>(new Map())
const modsServerChannelStore = useModsServerChannelStore()
function modelKey(model: string, chatProvider: ChatProvider): string {
return `${chatProvider.chat(model).baseURL}-${model}`
}
async function stream(model: string, chatProvider: ChatProvider, messages: Message[], options?: StreamOptions) {
const key = modelKey(model, chatProvider)
try {
await streamFrom(
// TODO(@nekomeowww,@shinohara-rin): we should not register the command callback on every stream anyway...
const sendSparkCommand = (command: WebSocketEvents['spark:command']) => {
// TODO(@nekomeowww): instruct the LLM to understand what destination is.
// Currently without skill like prompt injection, many issues occur.
// destination mostly are wrong or hallucinated, we need to find a way to make it more reliable.
//
// For now, since destinations as array will always broadcast to all connected modules/agents, we can set it to
// empty array to avoid wrong routing.
command.destinations = []
modsServerChannelStore.send({
type: 'spark:command',
data: command,
})
}
await coreStreamFrom({
model,
chatProvider,
messages,
// TODO(@nekomeowww,@shinohara-rin): we should not register the command callback on every stream anyway...
(command) => {
// TODO(@nekomeowww): instruct the LLM to understand what destination is.
// Currently without skill like prompt injection, many issues occur.
// destination mostly are wrong or hallucinated, we need to find a way to make it more reliable.
//
// For now, since destinations as array will always broadcast to all connected modules/agents, we can set it to
// empty array to avoid wrong routing.
command.destinations = []
modsServerChannelStore.send({
type: 'spark:command',
data: command,
})
},
{ ...options, toolsCompatibility: toolsCompatibility.value },
)
options: { ...options, toolsCompatibility: toolsCompatibility.value },
builtinToolsResolver: async () => [
...await mcp(),
...await debug(),
await createSparkCommandTool({ sendSparkCommand }),
],
})
}
catch (err) {
if (isToolRelatedError(err)) {
+14 -70
View File
@@ -1,70 +1,14 @@
import type { ContextUpdate, MetadataEventSource, WebSocketEventInputs } from '@proj-airi/server-sdk'
import type { AssistantMessage, CommonContentPart, CompletionToolCall, Message, SystemMessage, ToolMessage, UserMessage } from '@xsai/shared-chat'
export interface ChatSlicesText {
type: 'text'
text: string
}
export interface ChatSlicesToolCall {
type: 'tool-call'
toolCall: CompletionToolCall
}
export interface ChatSlicesToolCallResult {
type: 'tool-call-result'
id: string
isError?: boolean
result?: string | CommonContentPart[]
}
export type ChatSlices = ChatSlicesText | ChatSlicesToolCall | ChatSlicesToolCallResult
export interface ChatAssistantMessage extends AssistantMessage {
slices: ChatSlices[]
tool_results: {
id: string
isError?: boolean
result?: string | CommonContentPart[]
}[]
categorization?: {
speech: string
reasoning: string
}
}
export type ChatMessage = ChatAssistantMessage | SystemMessage | ToolMessage | UserMessage
export interface ErrorMessage {
role: 'error'
content: string
}
export interface ContextMessage extends ContextUpdate<Record<string, unknown>, unknown> {
metadata?: {
source: MetadataEventSource
}
createdAt: number
}
export type ChatHistoryItem = (ChatMessage | ErrorMessage) & { context?: ContextMessage } & { createdAt?: number, id?: string }
export interface ChatStreamEventContext {
message: ChatHistoryItem
contexts: Record<string, ContextMessage[]>
composedMessage: Array<Message>
input?: WebSocketEventInputs
}
export type ChatStreamEvent
= | { type: 'before-compose', message: string, sessionId: string, context: Omit<ChatStreamEventContext, 'composedMessage'> }
| { type: 'after-compose', message: string, sessionId: string, context: ChatStreamEventContext }
| { type: 'before-send', message: string, sessionId: string, context: ChatStreamEventContext }
| { type: 'after-send', message: string, sessionId: string, context: ChatStreamEventContext }
| { type: 'token-literal', literal: string, sessionId: string, context: ChatStreamEventContext }
| { type: 'token-special', special: string, sessionId: string, context: ChatStreamEventContext }
| { type: 'stream-end', sessionId: string, context: ChatStreamEventContext }
| { type: 'assistant-end', message: string, sessionId: string, context: ChatStreamEventContext }
| { type: 'assistant-message', message: ChatAssistantMessage, sessionId: string, messageText: string, context: ChatStreamEventContext }
export type StreamingAssistantMessage = ChatAssistantMessage & { context?: ContextMessage } & { createdAt?: number, id?: string }
export type {
ChatAssistantMessage,
ChatHistoryItem,
ChatMessage,
ChatSlices,
ChatSlicesText,
ChatSlicesToolCall,
ChatSlicesToolCallResult,
ChatStreamEvent,
ChatStreamEventContext,
ContextMessage,
ErrorMessage,
StreamingAssistantMessage,
} from '@proj-airi/core-agent'
+501 -139
View File
File diff suppressed because it is too large Load Diff