refactor: improve message synchronization with sequence-based pagination and merging

This commit is contained in:
weishu
2025-12-16 21:07:20 +08:00
parent b4654acb92
commit 0485065b87
11 changed files with 228 additions and 174 deletions
+7 -1
View File
@@ -147,7 +147,13 @@ export function registerCliHandlers(socket: Socket, deps: CliHandlersDeps): void
onWebappEvent?.({
type: 'message-received',
sessionId: sid,
data: msg.content
message: {
id: msg.id,
seq: msg.seq,
localId: msg.localId,
content: msg.content,
createdAt: msg.createdAt
}
})
})
+1 -1
View File
@@ -496,7 +496,7 @@ export class Store {
getMessages(sessionId: string, limit: number = 200, beforeSeq?: number): StoredMessage[] {
const safeLimit = Number.isFinite(limit) ? Math.max(1, Math.min(200, limit)) : 200
const rows = (beforeSeq && Number.isFinite(beforeSeq))
const rows = (beforeSeq !== undefined && beforeSeq !== null && Number.isFinite(beforeSeq))
? this.db.prepare(
'SELECT * FROM messages WHERE session_id = ? AND seq < ? ORDER BY seq DESC LIMIT ?'
).all(sessionId, beforeSeq, safeLimit) as DbMessageRow[]
+72 -11
View File
@@ -93,6 +93,7 @@ export interface Machine {
export interface DecryptedMessage {
id: string
seq: number
localId: string | null
content: unknown
createdAt: number
@@ -115,6 +116,7 @@ export interface SyncEvent {
sessionId?: string
machineId?: string
data?: unknown
message?: DecryptedMessage
}
export type SyncEventListener = (event: SyncEvent) => void
@@ -172,11 +174,18 @@ export class SyncEngine {
}
}
const webappEvent: SyncEvent = {
type: event.type,
sessionId: event.sessionId,
machineId: event.machineId
}
const webappEvent: SyncEvent = event.type === 'message-received'
? {
type: event.type,
sessionId: event.sessionId,
machineId: event.machineId,
message: event.message
}
: {
type: event.type,
sessionId: event.sessionId,
machineId: event.machineId
}
const rooms = new Set<string>()
if (webappEvent.sessionId) {
@@ -228,6 +237,46 @@ export class SyncEngine {
return this.sessionMessages.get(sessionId) || []
}
getMessagesPage(sessionId: string, options: { limit: number; beforeSeq: number | null }): {
messages: DecryptedMessage[]
page: {
limit: number
beforeSeq: number | null
nextBeforeSeq: number | null
hasMore: boolean
}
} {
const stored = this.store.getMessages(sessionId, options.limit, options.beforeSeq ?? undefined)
const messages: DecryptedMessage[] = stored.map((m) => ({
id: m.id,
seq: m.seq,
localId: m.localId,
content: m.content,
createdAt: m.createdAt
}))
let oldestSeq: number | null = null
for (const message of messages) {
if (typeof message.seq !== 'number') continue
if (oldestSeq === null || message.seq < oldestSeq) {
oldestSeq = message.seq
}
}
const nextBeforeSeq = oldestSeq
const hasMore = nextBeforeSeq !== null && this.store.getMessages(sessionId, 1, nextBeforeSeq).length > 0
return {
messages,
page: {
limit: options.limit,
beforeSeq: options.beforeSeq,
nextBeforeSeq,
hasMore
}
}
}
handleRealtimeEvent(event: SyncEvent): void {
if (event.type === 'session-updated' && event.sessionId) {
this.refreshSession(event.sessionId)
@@ -439,6 +488,7 @@ export class SyncEngine {
const stored = this.store.getMessages(sessionId, 200)
const messages: DecryptedMessage[] = stored.map((m) => ({
id: m.id,
seq: m.seq,
localId: m.localId,
content: m.content,
createdAt: m.createdAt
@@ -450,23 +500,24 @@ export class SyncEngine {
}
}
async sendMessage(sessionId: string, text: string): Promise<void> {
async sendMessage(sessionId: string, payload: { text: string; localId?: string | null; sentFrom?: 'telegram-bot' | 'webapp' }): Promise<void> {
const session = this.sessions.get(sessionId)
const sentFrom = payload.sentFrom ?? 'webapp'
const content = {
role: 'user',
content: {
type: 'text',
text
text: payload.text
},
meta: {
sentFrom: 'telegram-bot',
sentFrom,
permissionMode: session?.permissionMode || 'default',
model: session?.modelMode === 'default' ? null : session?.modelMode ?? undefined
}
}
const msg = this.store.addMessage(sessionId, content)
const msg = this.store.addMessage(sessionId, content, payload.localId ?? undefined)
const update = {
id: msg.id,
@@ -488,10 +539,20 @@ export class SyncEngine {
// Keep a small in-memory cache for Telegram rendering.
const cached = this.sessionMessages.get(sessionId) ?? []
cached.push({ id: msg.id, localId: msg.localId, content: msg.content, createdAt: msg.createdAt })
cached.push({ id: msg.id, seq: msg.seq, localId: msg.localId, content: msg.content, createdAt: msg.createdAt })
this.sessionMessages.set(sessionId, cached.slice(-200))
this.emit({ type: 'message-received', sessionId, data: msg.content })
this.emit({
type: 'message-received',
sessionId,
message: {
id: msg.id,
seq: msg.seq,
localId: msg.localId,
content: msg.content,
createdAt: msg.createdAt
}
})
}
async approvePermission(
+1 -1
View File
@@ -259,7 +259,7 @@ export class HappyBot {
}
if (event.type === 'message-received' && event.sessionId) {
const message = event.data as any
const message = (event.message?.content ?? event.data) as any
const messageContent = message?.content
const eventType = messageContent?.type === 'event' ? messageContent?.data?.type : null
+7 -37
View File
@@ -1,17 +1,17 @@
import { Hono } from 'hono'
import { z } from 'zod'
import type { FetchMessagesResult, SyncEngine } from '../../sync/syncEngine'
import type { SyncEngine } from '../../sync/syncEngine'
import type { WebAppEnv } from '../middleware/auth'
import { requireSessionFromParam, requireSyncEngine } from './guards'
const querySchema = z.object({
limit: z.coerce.number().int().min(1).max(200).optional(),
before: z.coerce.number().int().optional(),
refresh: z.string().optional()
beforeSeq: z.coerce.number().int().min(1).optional()
})
const sendMessageBodySchema = z.object({
text: z.string().min(1)
text: z.string().min(1),
localId: z.string().min(1).optional()
})
export function createMessagesRoutes(getSyncEngine: () => SyncEngine | null): Hono<WebAppEnv> {
@@ -31,38 +31,8 @@ export function createMessagesRoutes(getSyncEngine: () => SyncEngine | null): Ho
const parsed = querySchema.safeParse(c.req.query())
const limit = parsed.success ? (parsed.data.limit ?? 50) : 50
const before = parsed.success && Number.isFinite(parsed.data.before)
? (parsed.data.before ?? Number.POSITIVE_INFINITY)
: Number.POSITIVE_INFINITY
const refreshRaw = parsed.success ? parsed.data.refresh : undefined
const refresh = refreshRaw === '1' || refreshRaw === 'true'
const existing = engine.getSessionMessages(sessionId)
let fetchResult: FetchMessagesResult | null = null
if (refresh || existing.length === 0) {
fetchResult = await engine.fetchMessages(sessionId)
}
const messages = engine.getSessionMessages(sessionId)
const eligible = messages.filter((m) => m.createdAt < before)
const slice = eligible.slice(Math.max(0, eligible.length - limit))
const nextBefore = slice.length > 0 ? slice[0].createdAt : null
const hasMore = eligible.length > slice.length
return c.json({
messages: slice,
page: {
limit,
before: Number.isFinite(before) ? before : null,
nextBefore,
hasMore
},
warning: fetchResult && !fetchResult.ok
? { status: fetchResult.status, error: fetchResult.error }
: null
})
const beforeSeq = parsed.success ? (parsed.data.beforeSeq ?? null) : null
return c.json(engine.getMessagesPage(sessionId, { limit, beforeSeq }))
})
app.post('/sessions/:id/messages', async (c) => {
@@ -83,7 +53,7 @@ export function createMessagesRoutes(getSyncEngine: () => SyncEngine | null): Ho
return c.json({ error: 'Invalid body' }, 400)
}
await engine.sendMessage(sessionId, parsed.data.text)
await engine.sendMessage(sessionId, { text: parsed.data.text, localId: parsed.data.localId, sentFrom: 'webapp' })
return c.json({ ok: true })
})