diff --git a/.gitignore b/.gitignore index 34c44bbc..ceee91c8 100644 --- a/.gitignore +++ b/.gitignore @@ -29,3 +29,4 @@ coverage/ # Claude local settings .claude/settings.local.json localdocs/ +execplan/ diff --git a/server/src/socket/handlers/cli.ts b/server/src/socket/handlers/cli.ts index 7f2eb0f0..e6b1fd36 100644 --- a/server/src/socket/handlers/cli.ts +++ b/server/src/socket/handlers/cli.ts @@ -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 + } }) }) diff --git a/server/src/store/index.ts b/server/src/store/index.ts index d4f8f2e5..c9db9bd3 100644 --- a/server/src/store/index.ts +++ b/server/src/store/index.ts @@ -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[] diff --git a/server/src/sync/syncEngine.ts b/server/src/sync/syncEngine.ts index 19d2fcbb..0992b5ce 100644 --- a/server/src/sync/syncEngine.ts +++ b/server/src/sync/syncEngine.ts @@ -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() 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 { + async sendMessage(sessionId: string, payload: { text: string; localId?: string | null; sentFrom?: 'telegram-bot' | 'webapp' }): Promise { 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( diff --git a/server/src/telegram/bot.ts b/server/src/telegram/bot.ts index ed7c6b01..95fb25a7 100644 --- a/server/src/telegram/bot.ts +++ b/server/src/telegram/bot.ts @@ -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 diff --git a/server/src/web/routes/messages.ts b/server/src/web/routes/messages.ts index e22e66d3..0abe0c3f 100644 --- a/server/src/web/routes/messages.ts +++ b/server/src/web/routes/messages.ts @@ -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 { @@ -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 }) }) diff --git a/web/index.html b/web/index.html index 542c5b58..7c61d059 100644 --- a/web/index.html +++ b/web/index.html @@ -7,6 +7,13 @@ content="width=device-width, initial-scale=1, minimum-scale=1, maximum-scale=1, user-scalable=no, viewport-fit=cover" /> Happy Mini App + diff --git a/web/src/App.tsx b/web/src/App.tsx index 2c4061e2..68ce6fa4 100644 --- a/web/src/App.tsx +++ b/web/src/App.tsx @@ -62,72 +62,74 @@ function isUserMessage(msg: DecryptedMessage): boolean { return false } -function mergeRecentMessages(existing: DecryptedMessage[], recent: DecryptedMessage[]): DecryptedMessage[] { - if (existing.length === 0) return recent - if (recent.length === 0) return existing - - // Separate optimistic messages (those with localId and status) - const optimisticMessages = existing.filter(m => m.localId && m.status) - const nonOptimisticExisting = existing.filter(m => !m.localId || !m.status) - - const anchorId = recent[0]?.id - let merged: DecryptedMessage[] - - if (anchorId) { - const matchIndex = nonOptimisticExisting.findIndex((m) => m.id === anchorId) - if (matchIndex >= 0) { - merged = [...nonOptimisticExisting.slice(0, matchIndex), ...recent] - } else { - const anchorTime = recent[0]?.createdAt - if (typeof anchorTime === 'number') { - const timeMatchIndex = nonOptimisticExisting.findIndex((m) => m.createdAt >= anchorTime) - if (timeMatchIndex >= 0) { - merged = [...nonOptimisticExisting.slice(0, timeMatchIndex), ...recent] - } else { - merged = dedupeAndSort([...nonOptimisticExisting, ...recent]) - } - } else { - merged = dedupeAndSort([...nonOptimisticExisting, ...recent]) - } - } - } else { - merged = dedupeAndSort([...nonOptimisticExisting, ...recent]) +function compareMessages(a: DecryptedMessage, b: DecryptedMessage): number { + const aSeq = typeof a.seq === 'number' ? a.seq : null + const bSeq = typeof b.seq === 'number' ? b.seq : null + if (aSeq !== null && bSeq !== null && aSeq !== bSeq) { + return aSeq - bSeq } - - // Re-add optimistic messages that are still sending or failed - // (sent messages will be replaced by server messages) - for (const opt of optimisticMessages) { - if (opt.status === 'sent') { - // Check if server USER message with similar time exists - const hasServerUserMessage = merged.some(m => - !m.localId && - isUserMessage(m) && - Math.abs(m.createdAt - opt.createdAt) < 10000 - ) - if (hasServerUserMessage) continue - } - // Keep sending and failed messages - if (!merged.some(m => m.id === opt.id)) { - merged.push(opt) - } + if (a.createdAt !== b.createdAt) { + return a.createdAt - b.createdAt } - - merged.sort((a, b) => { - if (a.createdAt !== b.createdAt) return a.createdAt - b.createdAt - return a.id.localeCompare(b.id) - }) - - return merged + return a.id.localeCompare(b.id) } -function dedupeAndSort(messages: DecryptedMessage[]): DecryptedMessage[] { - const seen = new Set() - const result: DecryptedMessage[] = [] - for (const msg of messages) { - if (seen.has(msg.id)) continue - seen.add(msg.id) - result.push(msg) +function mergeMessages(existing: DecryptedMessage[], incoming: DecryptedMessage[]): DecryptedMessage[] { + if (existing.length === 0) { + return [...incoming].sort(compareMessages) } + if (incoming.length === 0) { + return [...existing].sort(compareMessages) + } + + const byId = new Map() + for (const msg of existing) { + byId.set(msg.id, msg) + } + for (const msg of incoming) { + byId.set(msg.id, msg) + } + + let merged = Array.from(byId.values()) + + const incomingLocalIds = new Set() + for (const msg of incoming) { + if (msg.localId) { + incomingLocalIds.add(msg.localId) + } + } + + // If we received a stored message with a localId, drop any optimistic bubble with the same localId. + if (incomingLocalIds.size > 0) { + merged = merged.filter((msg) => { + if (!msg.localId || !incomingLocalIds.has(msg.localId)) { + return true + } + return !msg.status + }) + } + + // Fallback: if an optimistic message was marked as sent but we didn't get a localId echo, + // drop it when a server user message appears close in time. + const optimisticMessages = merged.filter((m) => m.localId && m.status) + const nonOptimisticMessages = merged.filter((m) => !m.localId || !m.status) + const result: DecryptedMessage[] = [...nonOptimisticMessages] + + for (const optimistic of optimisticMessages) { + if (optimistic.status === 'sent') { + const hasServerUserMessage = nonOptimisticMessages.some((m) => + !m.status && + isUserMessage(m) && + Math.abs(m.createdAt - optimistic.createdAt) < 10_000 + ) + if (hasServerUserMessage) { + continue + } + } + result.push(optimistic) + } + + result.sort(compareMessages) return result } @@ -154,7 +156,7 @@ export function App() { const [messagesLoading, setMessagesLoading] = useState(false) const [messagesLoadingMore, setMessagesLoadingMore] = useState(false) const [messagesHasMore, setMessagesHasMore] = useState(false) - const [messagesNextBefore, setMessagesNextBefore] = useState(null) + const [messagesNextBeforeSeq, setMessagesNextBeforeSeq] = useState(null) const [messagesWarning, setMessagesWarning] = useState(null) const [machines, setMachines] = useState([]) @@ -162,7 +164,7 @@ export function App() { const [machinesError, setMachinesError] = useState(null) const [isSending, setIsSending] = useState(false) - const silentSyncInFlightRef = useRef(false) + const syncInFlightRef = useRef(false) useEffect(() => { const tg = getTelegramWebApp() @@ -277,7 +279,10 @@ export function App() { setSelectedSession(res.session) }, [api]) - const loadMessages = useCallback(async (sessionId: string, options: { before?: number; appendOlder?: boolean; refresh?: boolean }) => { + const loadMessages = useCallback(async ( + sessionId: string, + options: { beforeSeq?: number | null; appendOlder?: boolean } = {} + ) => { if (!api) return if (options.appendOlder) { @@ -289,23 +294,13 @@ export function App() { try { const res = await api.getMessages(sessionId, { limit: 50, - before: options.before, - refresh: options.refresh + beforeSeq: options.beforeSeq ?? null }) - if (options.appendOlder) { - setMessages((prev) => [...res.messages, ...prev]) - } else { - setMessages(res.messages) - } + setMessages((prev) => mergeMessages(prev, res.messages)) setMessagesHasMore(res.page.hasMore) - setMessagesNextBefore(res.page.nextBefore) - if (res.warning) { - const status = res.warning.status ?? 'error' - setMessagesWarning(`Happy Bot returned ${status} while fetching message history. Showing cached/live messages.`) - } else { - setMessagesWarning(null) - } + setMessagesNextBeforeSeq(res.page.nextBeforeSeq) + setMessagesWarning(null) } catch (e) { setMessagesWarning(e instanceof Error ? e.message : 'Failed to load messages') } finally { @@ -314,11 +309,11 @@ export function App() { } }, [api]) - const silentSyncSessionAndMessages = useCallback(async (sessionId: string) => { + const syncSessionAndMessages = useCallback(async (sessionId: string) => { if (!api) return if (messagesLoading || messagesLoadingMore) return - if (silentSyncInFlightRef.current) return - silentSyncInFlightRef.current = true + if (syncInFlightRef.current) return + syncInFlightRef.current = true try { const [sessionRes, messagesRes] = await Promise.all([ @@ -331,10 +326,12 @@ export function App() { } if (messagesRes) { - setMessages((prev) => mergeRecentMessages(prev, messagesRes.messages)) + setMessages((prev) => mergeMessages(prev, messagesRes.messages)) + setMessagesHasMore(messagesRes.page.hasMore) + setMessagesNextBeforeSeq(messagesRes.page.nextBeforeSeq) } } finally { - silentSyncInFlightRef.current = false + syncInFlightRef.current = false } }, [api, messagesLoading, messagesLoadingMore]) @@ -352,7 +349,7 @@ export function App() { ) ) - api.sendMessage(selectedSessionId, text) + api.sendMessage(selectedSessionId, text, localId) .then(() => { getTelegramWebApp()?.HapticFeedback?.notificationOccurred('success') setMessages((prev) => @@ -397,18 +394,18 @@ export function App() { setSelectedSession(null) setMessages([]) setMessagesHasMore(false) - setMessagesNextBefore(null) + setMessagesNextBeforeSeq(null) setMessagesWarning(null) return } setSelectedSession(null) setMessages([]) setMessagesHasMore(false) - setMessagesNextBefore(null) + setMessagesNextBeforeSeq(null) setMessagesWarning(null) loadSession(selectedSessionId) - loadMessages(selectedSessionId, { refresh: true }) + loadMessages(selectedSessionId) }, [api, selectedSessionId, loadSession, loadMessages]) useEffect(() => { @@ -432,6 +429,11 @@ export function App() { enabled: Boolean(api && token), token: token ?? '', subscription: socketSubscription, + onConnect: () => { + if (selectedSessionId) { + syncSessionAndMessages(selectedSessionId) + } + }, onEvent: (event: SyncEvent) => { if (event.type === 'session-added' || event.type === 'session-updated' || event.type === 'session-removed') { loadSessions() @@ -440,7 +442,7 @@ export function App() { } } if (event.type === 'message-received' && selectedSessionId && event.sessionId === selectedSessionId) { - silentSyncSessionAndMessages(selectedSessionId) + setMessages((prev) => mergeMessages(prev, [event.message])) } if (event.type === 'machine-updated' && (screen.type === 'machines' || screen.type === 'spawn')) { loadMachines() @@ -448,21 +450,6 @@ export function App() { } }) - useEffect(() => { - if (!api) return - if (screen.type !== 'session') return - if (!selectedSessionId) return - - silentSyncSessionAndMessages(selectedSessionId) - const interval = setInterval(() => { - silentSyncSessionAndMessages(selectedSessionId) - }, 3_000) - - return () => { - clearInterval(interval) - } - }, [api, screen.type, selectedSessionId, silentSyncSessionAndMessages]) - if (isAuthLoading) { return (
@@ -516,11 +503,11 @@ export function App() { onBack={goBack} onRefresh={() => { loadSession(screen.sessionId) - loadMessages(screen.sessionId, { refresh: true }) + loadMessages(screen.sessionId) }} onLoadMore={() => { - if (!messagesNextBefore) return - loadMessages(screen.sessionId, { before: messagesNextBefore, appendOlder: true }) + if (messagesNextBeforeSeq === null) return + loadMessages(screen.sessionId, { beforeSeq: messagesNextBeforeSeq, appendOlder: true }) }} onSend={(text) => { if (isSending) return @@ -529,6 +516,7 @@ export function App() { const localId = makeClientSideId('local') const optimisticMessage: DecryptedMessage = { id: localId, + seq: null, localId: localId, content: { role: 'user', content: text }, createdAt: Date.now(), @@ -537,10 +525,10 @@ export function App() { } // Immediately show message - setMessages((prev) => [...prev, optimisticMessage]) + setMessages((prev) => mergeMessages(prev, [optimisticMessage])) setIsSending(true) - api.sendMessage(screen.sessionId, text) + api.sendMessage(screen.sessionId, text, localId) .then(() => { getTelegramWebApp()?.HapticFeedback?.notificationOccurred('success') // Update status to sent diff --git a/web/src/api/client.ts b/web/src/api/client.ts index 41499b21..77ffc93d 100644 --- a/web/src/api/client.ts +++ b/web/src/api/client.ts @@ -57,21 +57,24 @@ export class ApiClient { return await this.request(`/api/sessions/${encodeURIComponent(sessionId)}`) } - async getMessages(sessionId: string, options: { before?: number; limit?: number; refresh?: boolean }): Promise { + async getMessages(sessionId: string, options: { beforeSeq?: number | null; limit?: number }): Promise { const params = new URLSearchParams() - if (options.before) params.set('before', `${options.before}`) - if (options.limit) params.set('limit', `${options.limit}`) - if (options.refresh) params.set('refresh', '1') + if (options.beforeSeq !== undefined && options.beforeSeq !== null) { + params.set('beforeSeq', `${options.beforeSeq}`) + } + if (options.limit !== undefined && options.limit !== null) { + params.set('limit', `${options.limit}`) + } const qs = params.toString() const url = `/api/sessions/${encodeURIComponent(sessionId)}/messages${qs ? `?${qs}` : ''}` return await this.request(url) } - async sendMessage(sessionId: string, text: string): Promise { + async sendMessage(sessionId: string, text: string, localId?: string | null): Promise { await this.request(`/api/sessions/${encodeURIComponent(sessionId)}/messages`, { method: 'POST', - body: JSON.stringify({ text }) + body: JSON.stringify({ text, localId: localId ?? undefined }) }) } diff --git a/web/src/hooks/useSocket.ts b/web/src/hooks/useSocket.ts index bedbfe72..499087ef 100644 --- a/web/src/hooks/useSocket.ts +++ b/web/src/hooks/useSocket.ts @@ -17,9 +17,13 @@ export function useSocket(options: { token: string subscription?: SocketSubscription onEvent: (event: SyncEvent) => void + onConnect?: () => void + onDisconnect?: (reason: string) => void onError?: (error: unknown) => void }): void { const onEventRef = useRef(options.onEvent) + const onConnectRef = useRef(options.onConnect) + const onDisconnectRef = useRef(options.onDisconnect) const onErrorRef = useRef(options.onError) const subscriptionRef = useRef(options.subscription ?? {}) const socketRef = useRef | null>(null) @@ -32,6 +36,14 @@ export function useSocket(options: { onErrorRef.current = options.onError }, [options.onError]) + useEffect(() => { + onConnectRef.current = options.onConnect + }, [options.onConnect]) + + useEffect(() => { + onDisconnectRef.current = options.onDisconnect + }, [options.onDisconnect]) + useEffect(() => { subscriptionRef.current = options.subscription ?? {} }, [options.subscription]) @@ -55,6 +67,13 @@ export function useSocket(options: { const sendSubscribe = () => { socket.emit('subscribe', subscriptionRef.current) } + const handleConnect = () => { + sendSubscribe() + onConnectRef.current?.() + } + const handleDisconnect = (reason: string) => { + onDisconnectRef.current?.(reason) + } socket.on('update', (event: unknown) => { if (!isObject(event)) return @@ -70,11 +89,13 @@ export function useSocket(options: { onErrorRef.current?.(error) }) - socket.on('connect', sendSubscribe) + socket.on('connect', handleConnect) + socket.on('disconnect', handleDisconnect) sendSubscribe() return () => { - socket.off('connect', sendSubscribe) + socket.off('connect', handleConnect) + socket.off('disconnect', handleDisconnect) socket.disconnect() if (socketRef.current === socket) { socketRef.current = null diff --git a/web/src/types/api.ts b/web/src/types/api.ts index f636bb07..be745b00 100644 --- a/web/src/types/api.ts +++ b/web/src/types/api.ts @@ -52,6 +52,7 @@ export type MessageStatus = 'sending' | 'sent' | 'failed' export type DecryptedMessage = { id: string + seq: number | null localId: string | null content: unknown createdAt: number @@ -86,14 +87,10 @@ export type MessagesResponse = { messages: DecryptedMessage[] page: { limit: number - before: number | null - nextBefore: number | null + beforeSeq: number | null + nextBeforeSeq: number | null hasMore: boolean } - warning?: { - status: number | null - error: string - } | null } export type MachinesResponse = { machines: Machine[] } @@ -106,6 +103,6 @@ export type SyncEvent = | { type: 'session-added'; sessionId: string; data?: unknown } | { type: 'session-updated'; sessionId: string; data?: unknown } | { type: 'session-removed'; sessionId: string } - | { type: 'message-received'; sessionId: string; data?: unknown } + | { type: 'message-received'; sessionId: string; message: DecryptedMessage } | { type: 'machine-updated'; machineId: string; data?: unknown } | { type: 'connection-changed'; data?: { status: string } }