Fix bridge history, profile models, and Windows gateway handling (#845)
* feat: support profile-aware group chat bridge flows * feat: route cron jobs through hermes cli * Fix group chat routing and isolate bridge tests * Add Grok image-to-video media skill * Default Grok videos to media directory * Fix bridge profile fallback and cron repeat clearing * Refine bridge chat and gateway platform handling * Filter bridge tool-call text deltas * Preserve structured bridge chat history * Prepare beta release build artifacts * Fix Windows run profile resolution * Fix Windows path compatibility checks * Fix profile-scoped model page display * Hide Windows subprocess windows for jobs and updates * Hide Windows file backend subprocess windows * Avoid Windows gateway restart lock conflicts * Treat Windows gateway lock as running on startup * Force release Windows gateway lock on restart * Tighten Windows gateway lock cleanup * Update chat e2e source expectation * Bump package version to 0.5.30 --------- Co-authored-by: Codex <codex@openai.com>
This commit is contained in:
@@ -6,6 +6,8 @@ import { getDb } from '../../../db'
|
||||
import { AgentClients } from './agent-clients'
|
||||
import { ContextEngine } from '../context-engine/compressor'
|
||||
import { SessionDeleter } from '../session-deleter'
|
||||
import { countTokens, SUMMARY_PREFIX } from '../../../lib/context-compressor'
|
||||
import { AgentBridgeClient } from '../agent-bridge'
|
||||
|
||||
// ─── Types ────────────────────────────────────────────────────
|
||||
|
||||
@@ -16,6 +18,43 @@ interface ChatMessage {
|
||||
senderName: string
|
||||
content: string
|
||||
timestamp: number
|
||||
role?: string
|
||||
tool_call_id?: string | null
|
||||
tool_calls?: any[] | null
|
||||
tool_name?: string | null
|
||||
finish_reason?: string | null
|
||||
reasoning?: string | null
|
||||
reasoning_details?: string | null
|
||||
reasoning_content?: string | null
|
||||
mentionDepth?: number
|
||||
}
|
||||
|
||||
function contentToStorageString(content: unknown): string {
|
||||
if (typeof content === 'string') return content
|
||||
return JSON.stringify(content ?? '')
|
||||
}
|
||||
|
||||
function contentToText(content: unknown): string {
|
||||
if (typeof content === 'string') {
|
||||
const trimmed = content.trim()
|
||||
if (trimmed.startsWith('[') && trimmed.endsWith(']')) {
|
||||
try {
|
||||
return contentToText(JSON.parse(trimmed))
|
||||
} catch {
|
||||
return content
|
||||
}
|
||||
}
|
||||
return content
|
||||
}
|
||||
if (Array.isArray(content)) {
|
||||
return content.map((block: any) => {
|
||||
if (block?.type === 'text') return block.text || ''
|
||||
if (block?.type === 'image') return `[Image: ${block.name || block.path || ''}]`
|
||||
if (block?.type === 'file') return `[File: ${block.name || block.path || ''}]`
|
||||
return ''
|
||||
}).filter(Boolean).join('\n')
|
||||
}
|
||||
return content == null ? '' : String(content)
|
||||
}
|
||||
|
||||
interface RoomAgent {
|
||||
@@ -64,6 +103,64 @@ export interface PendingSessionDeleteDrainResult {
|
||||
failed: Array<{ sessionId: string; error: string }>
|
||||
}
|
||||
|
||||
function parseJsonArray(value: unknown): any[] | null {
|
||||
if (value == null || value === '') return null
|
||||
if (Array.isArray(value)) return value
|
||||
if (typeof value !== 'string') return null
|
||||
try {
|
||||
const parsed = JSON.parse(value)
|
||||
return Array.isArray(parsed) ? parsed : null
|
||||
} catch {
|
||||
return null
|
||||
}
|
||||
}
|
||||
|
||||
function normalizeMessageRole(role: unknown): string {
|
||||
const value = String(role || '').trim()
|
||||
return ['user', 'assistant', 'tool', 'command'].includes(value) ? value : 'user'
|
||||
}
|
||||
|
||||
function normalizeMentionDepth(depth: unknown): number {
|
||||
const value = Number(depth)
|
||||
return Number.isFinite(value) && value > 0 ? Math.floor(value) : 0
|
||||
}
|
||||
|
||||
function groupRunOrder(id: string): { baseId: string; phase: number } {
|
||||
const value = String(id || '')
|
||||
const partMatch = value.match(/^(.*)_part_(\d+)(?:_(toolcall|toolresult)_.+)?$/)
|
||||
if (partMatch) {
|
||||
const part = Number(partMatch[2] || 0)
|
||||
const kind = partMatch[3] || 'assistant'
|
||||
const offset = kind === 'toolcall' ? 1 : kind === 'toolresult' ? 2 : 0
|
||||
return { baseId: partMatch[1], phase: part * 3 + offset }
|
||||
}
|
||||
const toolIdx = value.indexOf('_toolcall_')
|
||||
if (toolIdx >= 0) return { baseId: value.slice(0, toolIdx), phase: 0 }
|
||||
const resultIdx = value.indexOf('_toolresult_')
|
||||
if (resultIdx >= 0) return { baseId: value.slice(0, resultIdx), phase: 1 }
|
||||
return { baseId: value, phase: 2 }
|
||||
}
|
||||
|
||||
function sortGroupMessages<T extends { id: string; timestamp: number }>(messages: T[]): T[] {
|
||||
const baseMinTimestamp = new Map<string, number>()
|
||||
for (const msg of messages) {
|
||||
const { baseId } = groupRunOrder(msg.id)
|
||||
const existing = baseMinTimestamp.get(baseId)
|
||||
if (existing == null || msg.timestamp < existing) baseMinTimestamp.set(baseId, msg.timestamp)
|
||||
}
|
||||
return [...messages].sort((a, b) => {
|
||||
const ao = groupRunOrder(a.id)
|
||||
const bo = groupRunOrder(b.id)
|
||||
const at = baseMinTimestamp.get(ao.baseId) ?? a.timestamp
|
||||
const bt = baseMinTimestamp.get(bo.baseId) ?? b.timestamp
|
||||
if (at !== bt) return at - bt
|
||||
if (ao.baseId !== bo.baseId) return ao.baseId.localeCompare(bo.baseId)
|
||||
if (ao.phase !== bo.phase) return ao.phase - bo.phase
|
||||
if (a.timestamp !== b.timestamp) return a.timestamp - b.timestamp
|
||||
return a.id.localeCompare(b.id)
|
||||
})
|
||||
}
|
||||
|
||||
class ChatStorage {
|
||||
private db() { return getDb() }
|
||||
|
||||
@@ -175,16 +272,16 @@ class ChatStorage {
|
||||
|
||||
// ─── Rooms ────────────────────────────────────────────────
|
||||
|
||||
getRoom(roomId: string): { id: string; name: string; inviteCode: string | null; triggerTokens: number; maxHistoryTokens: number; tailMessageCount: number; totalTokens: number } | undefined {
|
||||
return this.db()?.prepare('SELECT id, name, inviteCode, triggerTokens, maxHistoryTokens, tailMessageCount, totalTokens FROM gc_rooms WHERE id = ?').get(roomId) as any
|
||||
getRoom(roomId: string): { id: string; name: string; inviteCode: string | null; triggerTokens: number; maxHistoryTokens: number; tailMessageCount: number; totalTokens: number; sessionSeed: string } | undefined {
|
||||
return this.db()?.prepare('SELECT id, name, inviteCode, triggerTokens, maxHistoryTokens, tailMessageCount, totalTokens, sessionSeed FROM gc_rooms WHERE id = ?').get(roomId) as any
|
||||
}
|
||||
|
||||
getRoomByInviteCode(code: string): { id: string; name: string; inviteCode: string | null; triggerTokens: number; maxHistoryTokens: number; tailMessageCount: number; totalTokens: number } | undefined {
|
||||
return this.db()?.prepare('SELECT id, name, inviteCode, triggerTokens, maxHistoryTokens, tailMessageCount, totalTokens FROM gc_rooms WHERE inviteCode = ?').get(code) as any
|
||||
getRoomByInviteCode(code: string): { id: string; name: string; inviteCode: string | null; triggerTokens: number; maxHistoryTokens: number; tailMessageCount: number; totalTokens: number; sessionSeed: string } | undefined {
|
||||
return this.db()?.prepare('SELECT id, name, inviteCode, triggerTokens, maxHistoryTokens, tailMessageCount, totalTokens, sessionSeed FROM gc_rooms WHERE inviteCode = ?').get(code) as any
|
||||
}
|
||||
|
||||
getAllRooms(): { id: string; name: string; inviteCode: string | null; triggerTokens: number; maxHistoryTokens: number; tailMessageCount: number; totalTokens: number }[] {
|
||||
return (this.db()?.prepare('SELECT id, name, inviteCode, triggerTokens, maxHistoryTokens, tailMessageCount, totalTokens FROM gc_rooms ORDER BY id').all() || []) as any[]
|
||||
getAllRooms(): { id: string; name: string; inviteCode: string | null; triggerTokens: number; maxHistoryTokens: number; tailMessageCount: number; totalTokens: number; sessionSeed: string }[] {
|
||||
return (this.db()?.prepare('SELECT id, name, inviteCode, triggerTokens, maxHistoryTokens, tailMessageCount, totalTokens, sessionSeed FROM gc_rooms ORDER BY id').all() || []) as any[]
|
||||
}
|
||||
|
||||
saveRoom(id: string, name: string, inviteCode?: string, config?: { triggerTokens?: number; maxHistoryTokens?: number; tailMessageCount?: number }): void {
|
||||
@@ -212,25 +309,132 @@ class ChatStorage {
|
||||
this.db()?.prepare('UPDATE gc_rooms SET totalTokens = ? WHERE id = ?').run(tokens, roomId)
|
||||
}
|
||||
|
||||
rotateRoomSessionSeed(roomId: string): string {
|
||||
const seed = `${Date.now().toString(36)}${Math.random().toString(36).slice(2, 8)}`
|
||||
this.db()?.prepare('UPDATE gc_rooms SET sessionSeed = ? WHERE id = ?').run(seed, roomId)
|
||||
return seed
|
||||
}
|
||||
|
||||
estimateTokens(text: string): number {
|
||||
const cjk = (text.match(/[\u2e80-\u9fff\uac00-\ud7af\u3000-\u303f\uff00-\uffef]/g) || []).length
|
||||
const other = text.length - cjk
|
||||
return Math.ceil(cjk * 1.5 + other / 4)
|
||||
}
|
||||
|
||||
private contentToUsageText(content: unknown): string {
|
||||
if (typeof content === 'string') return content
|
||||
if (!content) return ''
|
||||
if (Array.isArray(content)) {
|
||||
return content.map((block: any) => {
|
||||
if (typeof block?.text === 'string') return block.text
|
||||
if (typeof block?.type === 'string') return `[${block.type}]`
|
||||
return String(block || '')
|
||||
}).join('\n')
|
||||
}
|
||||
return String(content)
|
||||
}
|
||||
|
||||
private estimateUsageTokensFromMessages(messages: ChatMessage[]): { inputTokens: number; outputTokens: number } {
|
||||
const inputTokens = messages
|
||||
.filter(m => (m.role || 'user') === 'user')
|
||||
.reduce((sum, m) => sum + countTokens(this.contentToUsageText(m.content)), 0)
|
||||
const outputTokens = messages
|
||||
.filter(m => m.role === 'assistant' || m.role === 'tool')
|
||||
.reduce((sum, m) => sum + countTokens(this.contentToUsageText(m.content)) + countTokens(String(m.tool_calls || '')), 0)
|
||||
return { inputTokens, outputTokens }
|
||||
}
|
||||
|
||||
private estimateRoomTotalTokens(roomId: string, messages: ChatMessage[]): number {
|
||||
const snapshot = this.getContextSnapshot(roomId)
|
||||
if (snapshot && messages.length) {
|
||||
const snapshotIdx = messages.findIndex(m => m.id === snapshot.lastMessageId)
|
||||
const newMessages = snapshotIdx >= 0
|
||||
? messages.slice(snapshotIdx + 1)
|
||||
: messages.filter(m => m.timestamp > snapshot.lastMessageTimestamp)
|
||||
const newUsage = this.estimateUsageTokensFromMessages(newMessages)
|
||||
return countTokens(SUMMARY_PREFIX + snapshot.summary) + newUsage.inputTokens + newUsage.outputTokens
|
||||
}
|
||||
const usage = this.estimateUsageTokensFromMessages(messages)
|
||||
return usage.inputTokens + usage.outputTokens
|
||||
}
|
||||
|
||||
// ─── Messages ─────────────────────────────────────────────
|
||||
|
||||
getMessages(roomId: string, limit = 500): ChatMessage[] {
|
||||
const rows = (this.db()?.prepare(
|
||||
'SELECT id, roomId, senderId, senderName, content, timestamp FROM gc_messages WHERE roomId = ? ORDER BY timestamp DESC LIMIT ?'
|
||||
'SELECT id, roomId, senderId, senderName, content, timestamp, role, tool_call_id, tool_calls, tool_name, finish_reason, reasoning, reasoning_details, reasoning_content FROM gc_messages WHERE roomId = ? ORDER BY timestamp DESC LIMIT ?'
|
||||
).all(roomId, limit) || []) as any[]
|
||||
return rows.reverse()
|
||||
return sortGroupMessages(rows.map(row => ({
|
||||
...row,
|
||||
tool_calls: parseJsonArray(row.tool_calls),
|
||||
})))
|
||||
}
|
||||
|
||||
getMessage(messageId: string): ChatMessage | null {
|
||||
const row = this.db()?.prepare(
|
||||
'SELECT id, roomId, senderId, senderName, content, timestamp, role, tool_call_id, tool_calls, tool_name, finish_reason, reasoning, reasoning_details, reasoning_content FROM gc_messages WHERE id = ?'
|
||||
).get(messageId) as any
|
||||
if (!row) return null
|
||||
return {
|
||||
...row,
|
||||
tool_calls: parseJsonArray(row.tool_calls),
|
||||
}
|
||||
}
|
||||
|
||||
addMessage(msg: ChatMessage): void {
|
||||
this.upsertMessage(msg)
|
||||
}
|
||||
|
||||
upsertMessage(msg: ChatMessage): void {
|
||||
const toolCallsJson = msg.tool_calls ? JSON.stringify(msg.tool_calls) : null
|
||||
this.db()?.prepare(
|
||||
'INSERT INTO gc_messages (id, roomId, senderId, senderName, content, timestamp) VALUES (?, ?, ?, ?, ?, ?)'
|
||||
).run(msg.id, msg.roomId, msg.senderId, msg.senderName, msg.content, msg.timestamp)
|
||||
`INSERT INTO gc_messages (id, roomId, senderId, senderName, content, timestamp, role, tool_call_id, tool_calls, tool_name, finish_reason, reasoning, reasoning_details, reasoning_content)
|
||||
VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)`
|
||||
+ ` ON CONFLICT(id) DO UPDATE SET
|
||||
roomId = excluded.roomId,
|
||||
senderId = excluded.senderId,
|
||||
senderName = excluded.senderName,
|
||||
content = excluded.content,
|
||||
timestamp = excluded.timestamp,
|
||||
role = excluded.role,
|
||||
tool_call_id = excluded.tool_call_id,
|
||||
tool_calls = excluded.tool_calls,
|
||||
tool_name = excluded.tool_name,
|
||||
finish_reason = excluded.finish_reason,
|
||||
reasoning = excluded.reasoning,
|
||||
reasoning_details = excluded.reasoning_details,
|
||||
reasoning_content = excluded.reasoning_content`
|
||||
).run(
|
||||
msg.id, msg.roomId, msg.senderId, msg.senderName, msg.content, msg.timestamp,
|
||||
msg.role || 'user',
|
||||
msg.tool_call_id ?? null,
|
||||
toolCallsJson,
|
||||
msg.tool_name ?? null,
|
||||
msg.finish_reason ?? null,
|
||||
msg.reasoning ?? null,
|
||||
msg.reasoning_details ?? null,
|
||||
msg.reasoning_content ?? null,
|
||||
)
|
||||
}
|
||||
|
||||
saveMessageAndRefreshRoom(msg: ChatMessage, options: { preserveExistingTimestamp?: boolean } = {}): { message: ChatMessage; totalTokens: number } {
|
||||
const db = this.db()
|
||||
if (!db) return { message: msg, totalTokens: 0 }
|
||||
db.exec('BEGIN IMMEDIATE')
|
||||
try {
|
||||
const existing = this.getMessage(msg.id)
|
||||
const message = existing && options.preserveExistingTimestamp ? { ...msg, timestamp: existing.timestamp } : msg
|
||||
this.upsertMessage(message)
|
||||
this.pruneMessages(msg.roomId)
|
||||
const messages = this.getMessages(msg.roomId)
|
||||
const totalTokens = this.estimateRoomTotalTokens(msg.roomId, messages)
|
||||
this.updateRoomTotalTokens(msg.roomId, totalTokens)
|
||||
db.exec('COMMIT')
|
||||
return { message, totalTokens }
|
||||
} catch (err) {
|
||||
try { db.exec('ROLLBACK') } catch { /* ignore */ }
|
||||
throw err
|
||||
}
|
||||
}
|
||||
|
||||
clearRoomContext(roomId: string): void {
|
||||
@@ -238,7 +442,7 @@ class ChatStorage {
|
||||
if (!db) return
|
||||
db.prepare('DELETE FROM gc_messages WHERE roomId = ?').run(roomId)
|
||||
db.prepare('DELETE FROM gc_context_snapshots WHERE roomId = ?').run(roomId)
|
||||
db.prepare('UPDATE gc_rooms SET totalTokens = 0 WHERE id = ?').run(roomId)
|
||||
db.prepare('UPDATE gc_rooms SET totalTokens = 0, sessionSeed = ? WHERE id = ?').run(`${Date.now().toString(36)}${Math.random().toString(36).slice(2, 8)}`, roomId)
|
||||
}
|
||||
|
||||
pruneMessages(roomId: string, keep = 500): void {
|
||||
@@ -419,13 +623,6 @@ export class GroupChatServer {
|
||||
/** roomId -> (agentName -> { agentName, status }) */
|
||||
private contextStatusState = new Map<string, Map<string, { agentName: string; status: string }>>()
|
||||
|
||||
setGatewayManager(manager: any): void {
|
||||
this.agentClients.setGatewayManager(manager)
|
||||
if (this._contextEngine && manager) {
|
||||
this._contextEngine.setUpstream(manager.getUpstream(''), manager.getApiKey(''))
|
||||
}
|
||||
}
|
||||
|
||||
constructor(httpServers: HttpServer | HttpServer[]) {
|
||||
this.storage = new ChatStorage()
|
||||
this.storage.init()
|
||||
@@ -569,10 +766,18 @@ export class GroupChatServer {
|
||||
logger.debug(`[GroupChat] Connected: ${userName} (socket=${socket.id}, user=${userId})`)
|
||||
|
||||
socket.on('join', (data: { roomId?: string; name?: string }, ack?: (response?: unknown) => void) => this.handleJoin(socket, data, ack))
|
||||
socket.on('message', (data: { roomId?: string; content: string }, ack?: (response?: unknown) => void) => this.handleMessage(socket, data, ack))
|
||||
socket.on('message', (data: Partial<ChatMessage> & { roomId?: string; content: string | Array<Record<string, unknown>>; id?: string; mentionDepth?: number }, ack?: (response?: unknown) => void) => this.handleMessage(socket, data, ack))
|
||||
socket.on('message_stream_start', (data: { roomId?: string; id?: string; senderId?: string; senderName?: string; timestamp?: number }) => this.handleMessageStreamStart(socket, data))
|
||||
socket.on('message_stream_delta', (data: { roomId?: string; id?: string; delta?: string }) => this.handleMessageStreamDelta(socket, data))
|
||||
socket.on('message_reasoning_delta', (data: { roomId?: string; id?: string; delta?: string }) => this.handleMessageReasoningDelta(socket, data))
|
||||
socket.on('message_stream_end', (data: { roomId?: string; id?: string }) => this.handleMessageStreamEnd(socket, data))
|
||||
socket.on('typing', (data: { roomId?: string }) => this.handleTyping(socket, data))
|
||||
socket.on('stop_typing', (data: { roomId?: string }) => this.handleStopTyping(socket, data))
|
||||
socket.on('context_status', (data: { roomId?: string; agentName?: string; status?: string }) => this.handleContextStatus(socket, data))
|
||||
socket.on('interrupt_agent', (data: { roomId?: string; agentName?: string }, ack?: (response?: unknown) => void) => this.handleInterruptAgent(socket, data, ack))
|
||||
socket.on('approval.requested', (data: { roomId?: string; agentName?: string; approval_id?: string; command?: string; description?: string; choices?: string[]; allow_permanent?: boolean }) => this.handleApprovalRequested(socket, data))
|
||||
socket.on('approval.resolved', (data: { roomId?: string; agentName?: string; approval_id?: string; choice?: string }) => this.handleApprovalResolved(socket, data))
|
||||
socket.on('approval.respond', (data: { roomId?: string; approval_id?: string; choice?: string }, ack?: (response?: unknown) => void) => this.handleApprovalRespond(socket, data, ack))
|
||||
socket.on('disconnect', () => this.handleDisconnect(socket))
|
||||
}
|
||||
|
||||
@@ -581,14 +786,18 @@ export class GroupChatServer {
|
||||
private handleJoin(socket: Socket, data: { roomId?: string; name?: string; description?: string }, ack?: (res: any) => void): void {
|
||||
const socketId = socket.id
|
||||
const userId = this.socketUserMap.get(socketId) || socketId
|
||||
const userInfo = this.userInfoMap.get(userId) || { name: `User-${userId.slice(0, 6)}`, description: '' }
|
||||
const userName = data.name || userInfo.name
|
||||
const description = data.description || userInfo.description
|
||||
const roomId = data.roomId || 'general'
|
||||
const existingMember = this.storage.getMemberByUserId(roomId, userId)
|
||||
const userInfo = this.userInfoMap.get(userId) || {
|
||||
name: existingMember?.name || `User-${userId.slice(0, 6)}`,
|
||||
description: existingMember?.description || '',
|
||||
}
|
||||
const userName = data.name || existingMember?.name || userInfo.name
|
||||
const description = data.description || existingMember?.description || userInfo.description
|
||||
|
||||
// Update stored user info
|
||||
this.userInfoMap.set(userId, { name: userName, description })
|
||||
|
||||
const roomId = data.roomId || 'general'
|
||||
let room = this.rooms.get(roomId)
|
||||
if (!room) {
|
||||
room = new ChatRoom(roomId)
|
||||
@@ -628,7 +837,7 @@ export class GroupChatServer {
|
||||
logger.debug(`[GroupChat] ${userName} (user=${userId}) joined room: ${roomId}`)
|
||||
}
|
||||
|
||||
private handleMessage(socket: Socket, data: { roomId?: string; content: string }, ack?: (res: any) => void): void {
|
||||
private handleMessage(socket: Socket, data: Partial<ChatMessage> & { roomId?: string; content: string | Array<Record<string, unknown>>; id?: string; mentionDepth?: number }, ack?: (res: any) => void): void {
|
||||
const socketId = socket.id
|
||||
const roomId = data.roomId || 'general'
|
||||
const room = this.rooms.get(roomId)
|
||||
@@ -643,37 +852,105 @@ export class GroupChatServer {
|
||||
const userName = member?.name || `User-${socketId.slice(0, 6)}`
|
||||
|
||||
const msg: ChatMessage = {
|
||||
id: this.generateId(),
|
||||
id: this.normalizeClientMessageId(data.id) || this.generateId(),
|
||||
roomId,
|
||||
senderId: userId,
|
||||
senderName: userName,
|
||||
content: data.content,
|
||||
timestamp: Date.now(),
|
||||
content: contentToStorageString(data.content),
|
||||
timestamp: this.normalizeMessageTimestamp(data.timestamp, data.role),
|
||||
role: normalizeMessageRole(data.role),
|
||||
tool_call_id: data.tool_call_id ?? null,
|
||||
tool_calls: Array.isArray(data.tool_calls) ? data.tool_calls : null,
|
||||
tool_name: data.tool_name ?? null,
|
||||
finish_reason: data.finish_reason ?? null,
|
||||
reasoning: data.reasoning ?? null,
|
||||
reasoning_details: data.reasoning_details ?? null,
|
||||
reasoning_content: data.reasoning_content ?? null,
|
||||
}
|
||||
|
||||
this.storage.addMessage(msg)
|
||||
this.storage.pruneMessages(roomId)
|
||||
const saved = this.storage.saveMessageAndRefreshRoom(msg)
|
||||
const savedMsg = saved.message
|
||||
const totalTokens = saved.totalTokens
|
||||
|
||||
// Recalculate total tokens for the room
|
||||
const messages = this.storage.getMessages(roomId)
|
||||
const totalTokens = this.storage.estimateTokens(messages.map(m => m.content + m.senderName).join(''))
|
||||
this.storage.updateRoomTotalTokens(roomId, totalTokens)
|
||||
|
||||
this.nsp.to(roomId).emit('message', msg)
|
||||
this.nsp.to(roomId).emit('message', savedMsg)
|
||||
this.nsp.to(roomId).emit('room_updated', { roomId, totalTokens })
|
||||
ack?.({ id: msg.id })
|
||||
ack?.({ id: savedMsg.id })
|
||||
|
||||
// Server-side @mention routing — parse mentions and invoke agents directly
|
||||
this.agentClients.processMentions(roomId, {
|
||||
content: msg.content,
|
||||
senderName: msg.senderName,
|
||||
senderId: msg.senderId,
|
||||
timestamp: msg.timestamp,
|
||||
}).catch((err) => {
|
||||
logger.error(`[GroupChat] processMentions error: ${err.message}`)
|
||||
const mentionDepth = normalizeMentionDepth(data.mentionDepth)
|
||||
const shouldRouteMentions =
|
||||
savedMsg.role === 'user' ||
|
||||
(savedMsg.role === 'assistant' && mentionDepth < 2)
|
||||
|
||||
if (shouldRouteMentions) {
|
||||
// Server-side @mention routing — parse user mentions and invoke agents directly.
|
||||
this.agentClients.processMentions(roomId, {
|
||||
content: contentToText(savedMsg.content),
|
||||
input: Array.isArray(data.content) ? data.content : undefined,
|
||||
senderName: savedMsg.senderName,
|
||||
senderId: savedMsg.senderId,
|
||||
timestamp: savedMsg.timestamp,
|
||||
mentionDepth,
|
||||
}).catch((err) => {
|
||||
logger.error(`[GroupChat] processMentions error: ${err.message}`)
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
private handleMessageStreamStart(socket: Socket, data: { roomId?: string; id?: string; senderId?: string; senderName?: string; timestamp?: number }): void {
|
||||
const roomId = data.roomId || 'general'
|
||||
const room = this.rooms.get(roomId)
|
||||
if (!room || !room.hasOnlineMember(socket.id)) return
|
||||
const id = this.normalizeClientMessageId(data.id)
|
||||
if (!id) return
|
||||
|
||||
const member = room.getOnlineMemberBySocketId(socket.id)
|
||||
this.nsp.to(roomId).emit('message_stream_start', {
|
||||
id,
|
||||
roomId,
|
||||
senderId: data.senderId || member?.userId || socket.id,
|
||||
senderName: data.senderName || member?.name || `User-${socket.id.slice(0, 6)}`,
|
||||
content: '',
|
||||
timestamp: data.timestamp || Date.now(),
|
||||
role: 'assistant',
|
||||
finish_reason: 'streaming',
|
||||
})
|
||||
}
|
||||
|
||||
private handleMessageStreamDelta(socket: Socket, data: { roomId?: string; id?: string; delta?: string }): void {
|
||||
const roomId = data.roomId || 'general'
|
||||
const room = this.rooms.get(roomId)
|
||||
if (!room || !room.hasOnlineMember(socket.id)) return
|
||||
const id = this.normalizeClientMessageId(data.id)
|
||||
if (!id || !data.delta) return
|
||||
this.nsp.to(roomId).emit('message_stream_delta', {
|
||||
roomId,
|
||||
id,
|
||||
delta: String(data.delta),
|
||||
})
|
||||
}
|
||||
|
||||
private handleMessageReasoningDelta(socket: Socket, data: { roomId?: string; id?: string; delta?: string }): void {
|
||||
const roomId = data.roomId || 'general'
|
||||
const room = this.rooms.get(roomId)
|
||||
if (!room || !room.hasOnlineMember(socket.id)) return
|
||||
const id = this.normalizeClientMessageId(data.id)
|
||||
if (!id || !data.delta) return
|
||||
this.nsp.to(roomId).emit('message_reasoning_delta', {
|
||||
roomId,
|
||||
id,
|
||||
delta: String(data.delta),
|
||||
})
|
||||
}
|
||||
|
||||
private handleMessageStreamEnd(socket: Socket, data: { roomId?: string; id?: string }): void {
|
||||
const roomId = data.roomId || 'general'
|
||||
const room = this.rooms.get(roomId)
|
||||
if (!room || !room.hasOnlineMember(socket.id)) return
|
||||
const id = this.normalizeClientMessageId(data.id)
|
||||
if (!id) return
|
||||
this.nsp.to(roomId).emit('message_stream_end', { roomId, id })
|
||||
}
|
||||
|
||||
private handleTyping(socket: Socket, data: { roomId?: string }): void {
|
||||
const roomId = data.roomId || 'general'
|
||||
const userId = this.socketUserMap.get(socket.id) || socket.id
|
||||
@@ -749,6 +1026,75 @@ export class GroupChatServer {
|
||||
})
|
||||
}
|
||||
|
||||
private async handleInterruptAgent(socket: Socket, data: { roomId?: string; agentName?: string }, ack?: (response?: unknown) => void): Promise<void> {
|
||||
const roomId = data.roomId
|
||||
const agentName = data.agentName
|
||||
if (!roomId || !agentName) {
|
||||
ack?.({ error: 'roomId and agentName are required' })
|
||||
return
|
||||
}
|
||||
const room = this.rooms.get(roomId)
|
||||
if (!room?.hasOnlineMember(socket.id)) {
|
||||
ack?.({ error: 'Not in room' })
|
||||
return
|
||||
}
|
||||
try {
|
||||
await this.agentClients.interruptAgent(roomId, agentName)
|
||||
this.nsp.to(roomId).emit('context_status', { roomId, agentName, status: 'ready' })
|
||||
ack?.({ ok: true })
|
||||
} catch (err: any) {
|
||||
logger.warn(`[GroupChat] failed to interrupt agent ${agentName} in room ${roomId}: ${err.message}`)
|
||||
ack?.({ error: err.message || 'interrupt failed' })
|
||||
}
|
||||
}
|
||||
|
||||
private handleApprovalRequested(socket: Socket, data: { roomId?: string; agentName?: string; approval_id?: string; command?: string; description?: string; choices?: string[]; allow_permanent?: boolean }): void {
|
||||
const roomId = data.roomId
|
||||
if (!roomId || !data.approval_id) return
|
||||
this.nsp.to(roomId).emit('approval.requested', {
|
||||
event: 'approval.requested',
|
||||
roomId,
|
||||
agentName: data.agentName || '',
|
||||
approval_id: data.approval_id,
|
||||
command: data.command || '',
|
||||
description: data.description || '',
|
||||
choices: Array.isArray(data.choices) ? data.choices : ['once', 'session', 'deny'],
|
||||
allow_permanent: Boolean(data.allow_permanent),
|
||||
})
|
||||
}
|
||||
|
||||
private handleApprovalResolved(socket: Socket, data: { roomId?: string; agentName?: string; approval_id?: string; choice?: string }): void {
|
||||
const roomId = data.roomId
|
||||
if (!roomId || !data.approval_id) return
|
||||
this.nsp.to(roomId).emit('approval.resolved', {
|
||||
event: 'approval.resolved',
|
||||
roomId,
|
||||
agentName: data.agentName || '',
|
||||
approval_id: data.approval_id,
|
||||
choice: data.choice || '',
|
||||
})
|
||||
}
|
||||
|
||||
private async handleApprovalRespond(socket: Socket, data: { roomId?: string; approval_id?: string; choice?: string }, ack?: (response?: unknown) => void): Promise<void> {
|
||||
const roomId = data.roomId
|
||||
if (!roomId || !data.approval_id) {
|
||||
ack?.({ error: 'roomId and approval_id are required' })
|
||||
return
|
||||
}
|
||||
const room = this.rooms.get(roomId)
|
||||
if (!room?.hasOnlineMember(socket.id)) {
|
||||
ack?.({ error: 'Not in room' })
|
||||
return
|
||||
}
|
||||
try {
|
||||
const result = await new AgentBridgeClient().approvalRespond(data.approval_id, data.choice || 'deny')
|
||||
ack?.({ ok: true, resolved: Boolean((result as any)?.resolved) })
|
||||
} catch (err: any) {
|
||||
logger.warn(`[GroupChat] failed to respond approval ${data.approval_id}: ${err.message}`)
|
||||
ack?.({ error: err.message || 'approval response failed' })
|
||||
}
|
||||
}
|
||||
|
||||
private handleDisconnect(socket: Socket): void {
|
||||
const socketId = socket.id
|
||||
const userId = this.socketUserMap.get(socketId)
|
||||
@@ -804,4 +1150,19 @@ export class GroupChatServer {
|
||||
private generateId(): string {
|
||||
return Date.now().toString(36) + Math.random().toString(36).slice(2, 8)
|
||||
}
|
||||
|
||||
private normalizeClientMessageId(id?: string): string | null {
|
||||
const cleaned = String(id || '').trim()
|
||||
if (!cleaned || cleaned.length > 160) return null
|
||||
return /^[a-zA-Z0-9_-]+$/.test(cleaned) ? cleaned : null
|
||||
}
|
||||
|
||||
private normalizeMessageTimestamp(timestamp?: unknown, role?: unknown): number {
|
||||
const normalizedRole = normalizeMessageRole(role)
|
||||
if (normalizedRole !== 'user') {
|
||||
const value = Number(timestamp)
|
||||
if (Number.isFinite(value) && value > 0) return value
|
||||
}
|
||||
return Date.now()
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user