From 863efa20631c58ef58cd96896760422eaa661a2d Mon Sep 17 00:00:00 2001 From: Miya Date: Fri, 5 Jun 2026 10:54:02 +0200 Subject: [PATCH 1/4] fix: reconcile stale broker messages --- .../repro-stale-message-reconciliation.mjs | 152 +++++++++ src/main/broker.test.ts | 82 +++++ src/main/broker.ts | 299 ++++++++++++++++- src/main/ipc-handlers.ts | 10 + src/preload/index.ts | 12 + src/renderer/src/App.tsx | 2 + .../hooks/use-message-reconciliation.test.ts | 175 ++++++++++ .../src/hooks/use-message-reconciliation.ts | 302 ++++++++++++++++++ src/renderer/src/stores/agent-store.ts | 71 ++++ src/shared/types/ipc.ts | 49 +++ tsconfig.web.json | 3 +- vitest.config.mjs | 10 +- 12 files changed, 1156 insertions(+), 11 deletions(-) create mode 100644 scripts/repro-stale-message-reconciliation.mjs create mode 100644 src/renderer/src/hooks/use-message-reconciliation.test.ts create mode 100644 src/renderer/src/hooks/use-message-reconciliation.ts diff --git a/scripts/repro-stale-message-reconciliation.mjs b/scripts/repro-stale-message-reconciliation.mjs new file mode 100644 index 00000000..9a340e85 --- /dev/null +++ b/scripts/repro-stale-message-reconciliation.mjs @@ -0,0 +1,152 @@ +#!/usr/bin/env node +import { spawn } from 'node:child_process' +import { existsSync, readFileSync } from 'node:fs' +import { join, resolve } from 'node:path' +import WebSocket from 'ws' + +function readArg(name, fallback) { + const index = process.argv.indexOf(name) + if (index === -1 || index + 1 >= process.argv.length) return fallback + return process.argv[index + 1] +} + +function readNumberArg(name, fallback) { + const value = readArg(name, undefined) + if (value === undefined) return fallback + const parsed = Number(value) + return Number.isFinite(parsed) && parsed >= 0 ? parsed : fallback +} + +function runAgentRelay(args, cwd) { + return new Promise((resolvePromise) => { + const child = spawn('agent-relay', args, { + cwd, + stdio: ['ignore', 'pipe', 'pipe'] + }) + let stdout = '' + let stderr = '' + child.stdout.on('data', (chunk) => { + stdout += chunk.toString() + }) + child.stderr.on('data', (chunk) => { + stderr += chunk.toString() + }) + child.on('error', (error) => { + resolvePromise({ ok: false, stdout, stderr: error.message }) + }) + child.on('close', (code) => { + resolvePromise({ ok: code === 0, stdout, stderr }) + }) + }) +} + +function parseJsonArray(value) { + try { + const parsed = JSON.parse(value) + return Array.isArray(parsed) ? parsed : [] + } catch { + return [] + } +} + +async function main() { + const cwd = resolve(readArg('--cwd', process.cwd())) + const channel = readArg('--channel', 'general') + const idleMs = readNumberArg('--idle-ms', 0) + const timeoutMs = readNumberArg('--timeout-ms', 10_000) + const connectionPath = resolve( + readArg('--connection', join(cwd, '.agentworkforce', 'relay', 'connection.json')) + ) + const text = readArg('--text', `pear-stale-reconcile-probe ${new Date().toISOString()}`) + + if (!existsSync(connectionPath)) { + throw new Error(`Broker connection file not found: ${connectionPath}`) + } + + const connection = JSON.parse(readFileSync(connectionPath, 'utf8')) + if (!connection.url || !connection.api_key) { + throw new Error(`Connection file is missing url/api_key: ${connectionPath}`) + } + + const wsUrl = `${String(connection.url).replace(/^http/, 'ws')}/ws?sinceSeq=0` + const result = { + cwd, + channel, + text, + connectionPath, + wsUrl, + idleMs, + postOk: false, + cliListed: false, + brokerWsReceived: false, + brokerEvent: null, + postError: null + } + + const socket = new WebSocket(wsUrl, { + headers: { 'X-API-Key': connection.api_key } + }) + + const completion = new Promise((resolvePromise) => { + const timeout = setTimeout(() => resolvePromise('timeout'), timeoutMs) + + socket.on('message', (data) => { + try { + const event = JSON.parse(data.toString()) + if (event?.kind === 'relay_inbound' && event.body === text) { + result.brokerWsReceived = true + result.brokerEvent = { + kind: event.kind, + from: event.from, + target: event.target, + event_id: event.event_id, + seq: event.seq + } + clearTimeout(timeout) + resolvePromise('event') + } + } catch { + // Ignore non-JSON frames. + } + }) + + socket.on('error', (error) => { + result.postError = error.message + clearTimeout(timeout) + resolvePromise('error') + }) + }) + + await new Promise((resolvePromise, reject) => { + socket.once('open', resolvePromise) + socket.once('error', reject) + }) + + if (idleMs > 0) { + await new Promise((resolvePromise) => setTimeout(resolvePromise, idleMs)) + } + + const post = await runAgentRelay(['message', 'post', channel, text], cwd) + result.postOk = post.ok + if (!post.ok) { + result.postError = post.stderr || post.stdout || 'agent-relay message post failed' + } + + await completion + + const list = await runAgentRelay(['message', 'list', channel, '--limit', '10'], cwd) + const messages = parseJsonArray(list.stdout) + result.cliListed = messages.some((message) => message?.text === text) + + socket.terminate() + console.log(JSON.stringify(result, null, 2)) + + if (!result.postOk || !result.cliListed || !result.brokerWsReceived) { + process.exitCode = 1 + } +} + +main().catch((error) => { + console.error(error instanceof Error ? error.message : String(error)) + process.exitCode = 1 +}) diff --git a/src/main/broker.test.ts b/src/main/broker.test.ts index d43131b2..f3a1a6ef 100644 --- a/src/main/broker.test.ts +++ b/src/main/broker.test.ts @@ -2,6 +2,7 @@ import { chmod, mkdir, mkdtemp, rm, writeFile } from 'node:fs/promises' import { tmpdir } from 'node:os' import { dirname, join } from 'node:path' import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest' +import type { BrowserWindow } from 'electron' // Covers the multi-session BrokerManager: a project's local broker and cloud // sandbox broker coexist instead of clobbering each other in the sessions map. @@ -18,6 +19,7 @@ type MockClient = { onEvent: ReturnType addListener: ReturnType connectEvents: ReturnType + disconnectEvents: ReturnType renewLease: ReturnType shutdown: ReturnType disconnect: ReturnType @@ -46,6 +48,7 @@ const mock = vi.hoisted(() => { onEvent: vi.fn(() => () => undefined), addListener: vi.fn(() => () => undefined), connectEvents: vi.fn(), + disconnectEvents: vi.fn(), renewLease: vi.fn(async () => undefined), shutdown: vi.fn(async () => undefined), disconnect: vi.fn(), @@ -161,6 +164,25 @@ async function startLocal(manager: BrokerManager, agents: string[] = []): Promis return lastSpawned() } +function createMockWindow(): BrowserWindow { + return { + isDestroyed: vi.fn(() => false), + webContents: { + send: vi.fn() + } + } as unknown as BrowserWindow +} + +async function startLocalWithWindow( + manager: BrokerManager, + win: BrowserWindow, + agents: string[] = [] +): Promise { + mock.state.nextLocalAgents = agents + await manager.start(PROJECT_ID, '/tmp/project-1', 'pear-project-1', win, []) + return lastSpawned() +} + async function attachCloud(manager: BrokerManager, agents: string[] = []): Promise { mock.state.nextCloudAgents = agents await manager.attachCloudSandbox(PROJECT_ID, { @@ -340,6 +362,66 @@ describe('BrokerManager local + cloud coexistence', () => { await manager.shutdown() }) + + it('routes broker events to the current window after an existing session swaps windows', async () => { + const manager = new BrokerManager() + const firstWindow = createMockWindow() + const secondWindow = createMockWindow() + const local = await startLocalWithWindow(manager, firstWindow) + + await manager.start(PROJECT_ID, '/tmp/project-1', 'pear-project-1', secondWindow, []) + + const listener = local.onEvent.mock.calls.at(-1)?.[0] + expect(listener).toBeTypeOf('function') + listener?.({ + kind: 'relay_inbound', + from: 'codex-2', + target: '#general', + body: 'window swap proof', + event_id: 'evt-window-swap', + seq: 12 + }) + + expect((firstWindow.webContents.send as ReturnType).mock.calls + .some(([channel, payload]) => + channel === 'broker:event' && + (payload as { body?: string }).body === 'window swap proof' + )).toBe(false) + expect(secondWindow.webContents.send).toHaveBeenCalledWith( + 'broker:event', + expect.objectContaining({ + kind: 'relay_inbound', + body: 'window swap proof', + projectId: PROJECT_ID + }) + ) + + await manager.shutdown() + }) + + it('refreshEventStream rebinds the harness stream from the last seen sequence', async () => { + const manager = new BrokerManager() + const local = await startLocal(manager) + const listener = local.onEvent.mock.calls.at(-1)?.[0] + expect(listener).toBeTypeOf('function') + + listener?.({ + kind: 'relay_inbound', + from: 'codex-2', + target: '#general', + body: 'seq proof', + event_id: 'evt-seq-proof', + seq: 477 + }) + + await manager.refreshEventStream(PROJECT_ID, 'test-rebind') + + expect(local.disconnectEvents).toHaveBeenCalledTimes(1) + expect(local.onEvent).toHaveBeenCalledTimes(2) + expect(local.connectEvents).toHaveBeenLastCalledWith(477) + + await manager.shutdown() + }) }) describe('BrokerManager spawnAgent CLI preflight', () => { diff --git a/src/main/broker.ts b/src/main/broker.ts index 96b27432..16be98fe 100644 --- a/src/main/broker.ts +++ b/src/main/broker.ts @@ -17,6 +17,7 @@ import { type PendingRelayMessage, type PtyInputStream } from '@agent-relay/harness-driver' +import { AgentRelay, type RelayMessage } from '@agent-relay/sdk' import { getAccessToken, getApiUrl } from './auth' import { assertDirectory } from './path-utils' import { createPearBurnSpawnListener, stampPearBurnSpawnedAgent } from './burn-spawn-hook' @@ -27,7 +28,12 @@ import { type GeneratedCommitDraft } from './schemas' import { compactBrokerEvent, normalizeEventTimestamp } from '../shared/lib/broker-events' -import type { WorkforcePersona } from '../shared/types/ipc' +import type { + BrokerEventStreamDiagnostic, + BrokerReconciledChatMessage, + BrokerReconcileMessagesInput, + WorkforcePersona +} from '../shared/types/ipc' import { canExecute, resolveAgentRelayMcpCommand as resolveAgentRelayMcpCommandForOptions, @@ -334,6 +340,9 @@ const COMMIT_DRAFT_MAX_DIFF_CHARS = 80_000 const COMMIT_DRAFT_TIMEOUT_MS = 180_000 const MAX_BROKER_EVENT_HISTORY = 3_000 const BROKER_EVENT_HISTORY_TTL_MS = 12 * 60 * 60 * 1_000 +const DEFAULT_RECONCILE_MESSAGE_LIMIT = 50 +const MAX_RECONCILE_MESSAGE_LIMIT = 100 +const EVENT_STREAM_REBIND_COOLDOWN_MS = 5_000 // After this many consecutive failures to open a PTY input stream, give up on // the WS fast path for that agent briefly and send over HTTP while it cools down. const MAX_INPUT_STREAM_OPEN_FAILURES = 3 @@ -475,6 +484,101 @@ function toErrorMessage(err: unknown): string { type BrokerEventRecordPayload = Record & { kind: string } +function isBrokerDebugEnabled(): boolean { + return process.env.PEAR_BROKER_DEBUG === '1' || process.env.PEAR_BROKER_DEBUG === 'true' +} + +function normalizeReconcileLimit(limit: number | undefined): number { + if (!Number.isFinite(limit) || !limit || limit <= 0) return DEFAULT_RECONCILE_MESSAGE_LIMIT + return Math.min(Math.floor(limit), MAX_RECONCILE_MESSAGE_LIMIT) +} + +function normalizeChatChannelTarget(channelName: string): string { + const normalized = channelName.trim().replace(/^#/, '') + return normalized ? `#${normalized}` : '#general' +} + +function normalizeRelayTimestamp(value: string | undefined): number { + if (!value) return Date.now() + const parsed = Date.parse(value) + return Number.isFinite(parsed) ? parsed : Date.now() +} + +function isHumanSenderName(sender: string): boolean { + return sender.trim().toLowerCase() === 'human' +} + +function senderNameFromRelayMessage(message: RelayMessage): string { + return message.from.name || message.from.id || 'unknown' +} + +function directMessageTargetFromRelayMessage( + message: RelayMessage, + participants: string[] | undefined +): string { + const from = senderNameFromRelayMessage(message) + const normalizedParticipants = (participants || []) + .map((participant) => participant.trim().replace(/^@+/, '')) + .filter(Boolean) + const target = message.target + + if (target?.kind === 'agent' && target.agentName) return target.agentName + if (target?.kind === 'channel' && target.channelName) return normalizeChatChannelTarget(target.channelName) + + if (normalizedParticipants.length > 0) { + const otherParticipants = normalizedParticipants.filter((participant) => + participant.toLowerCase() !== from.toLowerCase() + ) + if (otherParticipants.length > 0) return otherParticipants.join(', ') + } + + const targetConversationId = target && 'conversationId' in target && typeof target.conversationId === 'string' + ? target.conversationId + : undefined + return message.conversationId || targetConversationId || 'direct-message' +} + +function normalizeRelayMessageForChat( + message: RelayMessage, + input: BrokerReconcileMessagesInput +): BrokerReconciledChatMessage | null { + const id = message.id || message.messageId + const body = message.text + if (!id || !body) return null + + const from = senderNameFromRelayMessage(message) + const to = input.kind === 'channel' + ? normalizeChatChannelTarget(input.channelName || message.channel?.name || 'general') + : directMessageTargetFromRelayMessage(message, input.dmParticipants) + + return { + id, + kind: 'message', + from, + to, + body, + timestamp: normalizeRelayTimestamp(message.createdAt || message.updatedAt), + isHuman: isHumanSenderName(from), + projectId: input.projectId, + ...(message.conversationId ? { conversationId: message.conversationId } : {}), + reactions: message.reactions?.map((reaction) => ({ + emoji: reaction.emoji, + count: reaction.count, + reactedByHuman: reaction.agents.some(isHumanSenderName) + })) + } +} + +function brokerEventSeq(event: BrokerEvent): number | undefined { + const seq = (event as Record).seq + return typeof seq === 'number' && Number.isFinite(seq) ? seq : undefined +} + +function brokerEventId(event: BrokerEvent): string | undefined { + const eventId = (event as Record).event_id + return typeof eventId === 'string' && eventId.trim() ? eventId : undefined +} + function getErrorStatus(err: unknown): unknown { if (typeof err !== 'object' || err === null || !('status' in err)) return undefined return (err as { status?: unknown }).status @@ -855,6 +959,11 @@ interface BrokerSession { channels: string[] cloudSandboxId: string | null pearLineage: Map + lastEventSeq?: number + lastEventAt?: number + lastEventId?: string + lastEventStreamRebindAt?: number + eventStreamRebinds?: number // For attach-to-remote-broker sessions (cloud sandboxes), the SDK doesn't // auto-renew the owner lease the way .spawn() does. The remote broker // auto-shuts-down after 120s without a lease renewal, so we own the timer @@ -1029,6 +1138,7 @@ export class BrokerManager { existing.cwd = cwd existing.name = name await this.syncChannels(normalizedProjectId, nextChannels) + await this.refreshEventStream(normalizedProjectId, 'existing-session-start', win) this.sendStatus(normalizedProjectId, 'connected') return } @@ -1071,7 +1181,7 @@ export class BrokerManager { existingClient.connectEvents() await this.syncChannels(normalizedProjectId, nextChannels) - this.publishBrokerEvent(normalizedProjectId, win, { + this.publishBrokerEvent(normalizedProjectId, normalizedProjectId, win, { kind: 'broker_initialized', name, cwd, @@ -1121,7 +1231,7 @@ export class BrokerManager { pearLineage: new Map() }) - this.publishBrokerEvent(normalizedProjectId, win, { + this.publishBrokerEvent(normalizedProjectId, normalizedProjectId, win, { kind: 'broker_initialized', name, cwd, @@ -1320,7 +1430,7 @@ export class BrokerManager { }) client.connectEvents() - this.publishBrokerEvent(normalizedProjectId, win, { + this.publishBrokerEvent(sessionKey, normalizedProjectId, win, { kind: 'broker_initialized', name: `cloud-${normalizedProjectId}`, url: getClientBaseUrl(client), @@ -1564,6 +1674,173 @@ export class BrokerManager { } } + private windowForSession(sessionKey: string, fallback?: BrowserWindow): BrowserWindow | undefined { + const win = this.sessions.get(sessionKey)?.window || fallback + if (!win || win.isDestroyed()) return undefined + return win + } + + private emitEventStreamDiagnostic( + sessionKey: string, + fallback: BrowserWindow | undefined, + diagnostic: BrokerEventStreamDiagnostic + ): void { + if (diagnostic.status === 'received' && !isBrokerDebugEnabled()) return + + if (isBrokerDebugEnabled()) { + console.info('[broker:event-stream]', diagnostic) + } + + const win = this.windowForSession(sessionKey, fallback) + if (win) { + win.webContents.send('broker:event-stream-diagnostic', diagnostic) + } + } + + private noteBrokerEventReceipt(sessionKey: string, event: BrokerEvent): void { + const session = this.sessions.get(sessionKey) + const seq = brokerEventSeq(event) + const eventId = brokerEventId(event) + if (session) { + session.lastEventAt = Date.now() + if (seq !== undefined) session.lastEventSeq = seq + if (eventId) session.lastEventId = eventId + } + + this.emitEventStreamDiagnostic(sessionKey, session?.window, { + projectId: projectIdFromSessionKey(sessionKey), + status: 'received', + at: Date.now(), + eventKind: event.kind, + ...(eventId ? { eventId } : {}), + ...(seq !== undefined ? { seq } : {}) + }) + } + + async refreshEventStream(projectId?: string, reason = 'manual', win?: BrowserWindow): Promise { + const sessions = projectId + ? await this.getOrAwaitSessionsForProject(projectId) + : Array.from(this.sessions.values()) + + for (const session of sessions) { + await this.rebindSessionEventStream(sessionKeyFor(session), session, reason, win) + } + } + + private async rebindSessionEventStream( + sessionKey: string, + session: BrokerSession, + reason: string, + win?: BrowserWindow + ): Promise { + const now = Date.now() + if (win && !win.isDestroyed()) { + session.window = win + } + + if ( + session.lastEventStreamRebindAt && + now - session.lastEventStreamRebindAt < EVENT_STREAM_REBIND_COOLDOWN_MS + ) { + this.emitEventStreamDiagnostic(sessionKey, session.window, { + projectId: session.projectId, + status: 'rebind-skipped', + reason, + at: now, + reconnects: session.eventStreamRebinds || 0, + ...(session.lastEventSeq !== undefined ? { seq: session.lastEventSeq } : {}), + ...(session.lastEventId ? { eventId: session.lastEventId } : {}) + }) + return + } + + this.emitEventStreamDiagnostic(sessionKey, session.window, { + projectId: session.projectId, + status: 'rebind-started', + reason, + at: now, + reconnects: session.eventStreamRebinds || 0, + ...(session.lastEventSeq !== undefined ? { seq: session.lastEventSeq } : {}), + ...(session.lastEventId ? { eventId: session.lastEventId } : {}) + }) + + try { + session.unsubEvent() + session.client.disconnectEvents() + session.unsubEvent = this.attachClient(sessionKey, session.client, session.window) + session.client.connectEvents(session.lastEventSeq) + session.lastEventStreamRebindAt = now + session.eventStreamRebinds = (session.eventStreamRebinds || 0) + 1 + + this.emitEventStreamDiagnostic(sessionKey, session.window, { + projectId: session.projectId, + status: 'rebound', + reason, + at: Date.now(), + reconnects: session.eventStreamRebinds, + ...(session.lastEventSeq !== undefined ? { seq: session.lastEventSeq } : {}), + ...(session.lastEventId ? { eventId: session.lastEventId } : {}) + }) + } catch (err) { + this.emitEventStreamDiagnostic(sessionKey, session.window, { + projectId: session.projectId, + status: 'rebind-error', + reason, + at: Date.now(), + error: toErrorMessage(err), + reconnects: session.eventStreamRebinds || 0 + }) + throw err + } + } + + async reconcileMessages(input: BrokerReconcileMessagesInput): Promise { + const projectId = input.projectId.trim() + if (!projectId) throw new Error('Project id is required') + if (input.kind === 'channel' && !input.channelName?.trim()) { + throw new Error('Channel name is required for channel reconciliation') + } + if (input.kind === 'dm' && !input.conversationId?.trim()) { + throw new Error('Conversation id is required for DM reconciliation') + } + + const [session] = await this.getOrAwaitSessionsForProject(projectId) + if (!session) return [] + + const metadata = await session.client.getSession() + const workspaceKey = metadata.workspace_key + if (!workspaceKey) { + throw new Error('Broker session does not expose a Relay workspace key') + } + + const relay = new AgentRelay({ + workspaceKey, + ...(process.env.RELAY_BASE_URL ? { baseUrl: process.env.RELAY_BASE_URL } : {}) + }) + const limit = normalizeReconcileLimit(input.limit) + const messages = input.kind === 'channel' + ? await relay.messages.list(input.channelName!, { limit }) + : await relay.messages.listDirect({ conversationId: input.conversationId!, limit }) + + const normalized = messages + .map((message) => normalizeRelayMessageForChat(message, input)) + .filter((message): message is BrokerReconciledChatMessage => message !== null) + .sort((left, right) => left.timestamp - right.timestamp) + + if (isBrokerDebugEnabled()) { + console.info('[broker:reconcile]', { + projectId, + kind: input.kind, + channelName: input.channelName, + conversationId: input.conversationId, + requested: limit, + returned: normalized.length + }) + } + + return normalized + } + private attachClient(sessionKey: string, client: AgentRelayClient, win?: BrowserWindow): () => void { const projectId = projectIdFromSessionKey(sessionKey) // Stamp every spawn this session does so burn can attribute Pear sessions @@ -1584,6 +1861,7 @@ export class BrokerManager { const unsubBurn = client.addListener('beforeAgentSpawn', burnHandler) const unsubEvent = client.onEvent((event: BrokerEvent) => { + this.noteBrokerEventReceipt(sessionKey, event) // Fast path for PTY chunks: ship just (projectId, name, chunk) over a // dedicated channel so typing latency doesn't pay for compactBrokerEvent, // the broker:event metadata spread, or pushing into eventHistory per @@ -1594,8 +1872,9 @@ export class BrokerManager { 'name' in event && typeof event.name === 'string' && 'chunk' in event && typeof event.chunk === 'string' ) { - if (win && !win.isDestroyed()) { - win.webContents.send('broker:pty-chunk', projectId, event.name, event.chunk) + const targetWindow = this.windowForSession(sessionKey, win) + if (targetWindow && !targetWindow.isDestroyed()) { + targetWindow.webContents.send('broker:pty-chunk', projectId, event.name, event.chunk) } this.rememberAgentSession(event.name, sessionKey) if (this.sessions.get(sessionKey)?.cloudSandboxId) { @@ -1606,7 +1885,7 @@ export class BrokerManager { return } - this.publishBrokerEvent(projectId, win, event as unknown as BrokerEventRecordPayload) + this.publishBrokerEvent(sessionKey, projectId, win, event as unknown as BrokerEventRecordPayload) if (event.kind === 'agent_spawned' && event.name) { this.rememberAgentSession(event.name, sessionKey) @@ -1646,13 +1925,15 @@ export class BrokerManager { } private publishBrokerEvent( + sessionKey: string, projectId: string, win: BrowserWindow | undefined, event: BrokerEventRecordPayload ): BrokerEventRecord { const record = this.recordBrokerEvent(projectId, event) - if (win && !win.isDestroyed()) { - win.webContents.send('broker:event', { + const targetWindow = this.windowForSession(sessionKey, win) + if (targetWindow && !targetWindow.isDestroyed()) { + targetWindow.webContents.send('broker:event', { ...event, projectId, observedAt: record.timestamp, diff --git a/src/main/ipc-handlers.ts b/src/main/ipc-handlers.ts index 4ce14e5d..de82f1f1 100644 --- a/src/main/ipc-handlers.ts +++ b/src/main/ipc-handlers.ts @@ -27,6 +27,7 @@ import { burnManager, type BurnAgentInput, type BurnProjectInput, type BurnSessi import { resetRelayWorkspaceManager } from './relay-workspace' import { assertDirectory, isDirectory } from './path-utils' import { findProjectForPath } from './cli' +import type { BrokerReconcileMessagesInput } from '../shared/types/ipc' import type { ProactiveAgentDraft } from './proactive-agent.types' function getProjectIdForPath(targetPath: string): string | null { @@ -258,6 +259,15 @@ export function registerIpcHandlers(): void { await brokerManager.sendMessage(projectId, input) }) + ipcMain.handle('broker:reconcile-messages', async (_, input: BrokerReconcileMessagesInput) => { + return brokerManager.reconcileMessages(input) + }) + + ipcMain.handle('broker:refresh-event-stream', async (event, projectId?: string, reason?: string) => { + const win = BrowserWindow.fromWebContents(event.sender) + await brokerManager.refreshEventStream(projectId, reason || 'renderer-request', win || undefined) + }) + ipcMain.handle('broker:subscribe-agent-channel', async (_, projectId: string | undefined, name: string, channel: string) => { await brokerManager.subscribeAgentChannel(projectId, name, channel) }) diff --git a/src/preload/index.ts b/src/preload/index.ts index 945ebbb4..ce84e821 100644 --- a/src/preload/index.ts +++ b/src/preload/index.ts @@ -12,7 +12,10 @@ import type { BrokerAttachTerminalResult, BrokerDetails, BrokerEventRecord, + BrokerEventStreamDiagnostic, BrokerListAgent, + BrokerReconciledChatMessage, + BrokerReconcileMessagesInput, BrokerSendMessageInput, BrokerSetTerminalModeResult, BrokerSpawnAgentInput, @@ -83,7 +86,10 @@ export type { BrokerAttachTerminalResult, BrokerDetails, BrokerEventRecord, + BrokerEventStreamDiagnostic, BrokerListAgent, + BrokerReconciledChatMessage, + BrokerReconcileMessagesInput, BrokerSendMessageInput, BrokerSetTerminalModeResult, BrokerSpawnAgentInput, @@ -241,6 +247,10 @@ const api = { invoke('broker:input-srtt', projectId, name), sendMessage: (projectId: string | undefined, input: BrokerSendMessageInput) => invoke('broker:send-message', projectId, input), + reconcileMessages: (input: BrokerReconcileMessagesInput) => + invoke('broker:reconcile-messages', input), + refreshEventStream: (projectId?: string, reason?: string) => + invoke('broker:refresh-event-stream', projectId, reason), subscribeAgentChannel: (projectId: string | undefined, name: string, channel: string) => invoke('broker:subscribe-agent-channel', projectId, name, channel), unsubscribeAgentChannel: (projectId: string | undefined, name: string, channel: string) => @@ -253,6 +263,8 @@ const api = { listEvents: () => invoke('broker:list-events'), shutdown: () => invoke('broker:shutdown'), onEvent: (callback: (event: unknown) => void) => subscribe('broker:event', callback), + onEventStreamDiagnostic: (callback: (event: BrokerEventStreamDiagnostic) => void) => + subscribe('broker:event-stream-diagnostic', callback), onPtyChunk: (callback: (projectId: string, name: string, chunk: string) => void) => { const handler = (_: unknown, projectId: string, name: string, chunk: string): void => callback(projectId, name, chunk) diff --git a/src/renderer/src/App.tsx b/src/renderer/src/App.tsx index 70665cf3..d0c6bc54 100644 --- a/src/renderer/src/App.tsx +++ b/src/renderer/src/App.tsx @@ -24,6 +24,7 @@ import { useUIStore } from '@/stores/ui-store' import { useBrokerEvents } from '@/hooks/use-broker-events' import { useCloudAgentEvents } from '@/hooks/use-cloud-agent' import { useGitStatus } from '@/hooks/use-git-status' +import { useMessageReconciliation } from '@/hooks/use-message-reconciliation' import { useAgentStore } from '@/stores/agent-store' import { initTypingTrace } from '@/lib/typing-trace' @@ -45,6 +46,7 @@ export default function App(): React.ReactNode { useBrokerEvents() useCloudAgentEvents() useGitStatus() + useMessageReconciliation() useEffect(() => { load() diff --git a/src/renderer/src/hooks/use-message-reconciliation.test.ts b/src/renderer/src/hooks/use-message-reconciliation.test.ts new file mode 100644 index 00000000..783708a9 --- /dev/null +++ b/src/renderer/src/hooks/use-message-reconciliation.test.ts @@ -0,0 +1,175 @@ +import { afterEach, beforeAll, describe, expect, it, vi } from 'vitest' +import type { MessageReconciliationRequest } from './use-message-reconciliation' +import type { ChatMessage } from '@/stores/agent-store' + +vi.mock('@/lib/ipc', () => ({ + pear: { + broker: {}, + project: {} + } +})) + +let hooks: typeof import('./use-message-reconciliation') +let agentStore: typeof import('@/stores/agent-store') + +beforeAll(async () => { + Object.defineProperty(globalThis, 'localStorage', { + configurable: true, + value: { + getItem: vi.fn(() => null), + setItem: vi.fn(), + removeItem: vi.fn() + } + }) + + hooks = await import('./use-message-reconciliation') + agentStore = await import('@/stores/agent-store') +}) + +const channelMessage: ChatMessage = { + id: 'msg-1', + from: 'codex-1', + to: '#general', + body: 'missed while idle', + timestamp: 1_717_000_000_000, + isHuman: false, + projectId: 'project-1' +} + +describe('getActiveMessageReconciliationRequest', () => { + it('builds a bounded channel reconciliation request for the active tab', () => { + expect(hooks.getActiveMessageReconciliationRequest({ + activeProjectId: 'fallback-project', + activeTab: { + id: 'channel:project-1:general', + kind: 'channel', + title: 'general', + projectId: 'project-1', + channelName: '#general' + } + })).toEqual({ + projectId: 'project-1', + kind: 'channel', + channelName: 'general', + limit: 50 + }) + }) + + it('builds a DM request from the active tab participants', () => { + expect(hooks.getActiveMessageReconciliationRequest({ + activeProjectId: 'project-1', + activeTab: { + id: 'dm:project-1:human|worker', + kind: 'dm', + title: 'Worker', + dmParticipants: ['Worker', 'human'] + }, + limit: 25 + })).toEqual({ + projectId: 'project-1', + kind: 'dm', + conversationId: 'human|worker', + dmParticipants: ['Worker', 'human'], + limit: 25 + }) + }) + + it('does not reconcile non-chat tabs', () => { + expect(hooks.getActiveMessageReconciliationRequest({ + activeProjectId: 'project-1', + activeTab: { + id: 'agents', + kind: 'agents', + title: 'Agents' + } + })).toBeNull() + }) +}) + +describe('createMessageReconciler', () => { + afterEach(() => { + vi.useRealTimers() + agentStore.useAgentStore.getState().clearAll() + }) + + it('debounces triggers, fetches canonical messages, and merges the result', async () => { + vi.useFakeTimers() + const request: MessageReconciliationRequest = { + projectId: 'project-1', + kind: 'channel', + channelName: 'general', + limit: 50 + } + const reconcileMessages = vi.fn(async () => [channelMessage]) + const mergeMessages = vi.fn() + const debug = vi.fn() + const reconciler = hooks.createMessageReconciler({ + getRequest: () => request, + reconcileMessages, + mergeMessages, + setTimeout: setTimeout as unknown as typeof window.setTimeout, + clearTimeout: clearTimeout as unknown as typeof window.clearTimeout, + debounceMs: 250, + now: () => 123, + debug + }) + + reconciler.schedule('broker-status') + reconciler.schedule('window-focus') + await vi.advanceTimersByTimeAsync(249) + + expect(reconcileMessages).not.toHaveBeenCalled() + + await vi.advanceTimersByTimeAsync(1) + + expect(reconcileMessages).toHaveBeenCalledTimes(1) + expect(reconcileMessages).toHaveBeenCalledWith(request) + expect(mergeMessages).toHaveBeenCalledTimes(1) + expect(mergeMessages).toHaveBeenCalledWith([channelMessage]) + expect(debug).toHaveBeenCalledWith(expect.objectContaining({ + kind: 'merged', + reason: 'window-focus', + messageCount: 1, + timestamp: 123 + })) + }) + + it('skips fetches when no active chat room can be reconciled', async () => { + vi.useFakeTimers() + const reconcileMessages = vi.fn(async () => [channelMessage]) + const mergeMessages = vi.fn() + const reconciler = hooks.createMessageReconciler({ + getRequest: () => null, + reconcileMessages, + mergeMessages, + setTimeout: setTimeout as unknown as typeof window.setTimeout, + clearTimeout: clearTimeout as unknown as typeof window.clearTimeout, + debounceMs: 1 + }) + + reconciler.schedule('active-room') + await vi.advanceTimersByTimeAsync(1) + + expect(reconcileMessages).not.toHaveBeenCalled() + expect(mergeMessages).not.toHaveBeenCalled() + }) + + it('merges reconciled messages into the Zustand store and dedups by id', () => { + const store = agentStore.useAgentStore.getState() + + store.reconcileMessages([channelMessage]) + store.reconcileMessages([{ + ...channelMessage, + body: 'updated canonical body', + timestamp: channelMessage.timestamp + 1 + }]) + + const messages = agentStore.useAgentStore.getState().messages + expect(messages).toHaveLength(1) + expect(messages[0]).toMatchObject({ + id: channelMessage.id, + body: 'updated canonical body', + timestamp: channelMessage.timestamp + 1 + }) + }) +}) diff --git a/src/renderer/src/hooks/use-message-reconciliation.ts b/src/renderer/src/hooks/use-message-reconciliation.ts new file mode 100644 index 00000000..c950e25b --- /dev/null +++ b/src/renderer/src/hooks/use-message-reconciliation.ts @@ -0,0 +1,302 @@ +import { useEffect, useMemo } from 'react' +import { getDirectMessageRoomId } from '@/lib/direct-messages' +import { useAgentStore, type ChatMessage } from '@/stores/agent-store' +import { useProjectStore } from '@/stores/project-store' +import { type AppTab, useUIStore } from '@/stores/ui-store' +import type { + BrokerEventStreamDiagnostic, + BrokerReconciledChatMessage, + BrokerReconcileMessagesInput, + PearAPI +} from '@shared/types/ipc' + +const DEFAULT_RECONCILE_LIMIT = 50 +const DEFAULT_RECONCILE_DEBOUNCE_MS = 750 +const BROKER_CONNECTED_STATUSES = new Set([ + 'connected', + 'event_stream_connected', + 'event_stream_reconnected', + 'event-stream-connected', + 'event-stream-reconnected' +]) + +const EVENT_STREAM_RECONCILED_STATUSES = new Set([ + 'rebound', + 'received' +]) + +export type MessageReconciliationRequest = BrokerReconcileMessagesInput + +interface BrokerWithMessageReconciliation { + reconcileMessages: (input: MessageReconciliationRequest) => Promise + refreshEventStream?: (projectId?: string, reason?: string) => Promise + onEventStreamDiagnostic?: (callback: (event: BrokerEventStreamDiagnostic) => void) => () => void +} + +interface StoreWithMessageReconciliation { + reconcileMessages?: (messages: ChatMessage[]) => void + handleBrokerEvent: (event: Record & { kind: string }) => void +} + +interface MessageReconcilerDeps { + getRequest: () => MessageReconciliationRequest | null + reconcileMessages: (input: MessageReconciliationRequest) => Promise + mergeMessages: (messages: ChatMessage[]) => void + setTimeout: (handler: () => void, timeout: number) => number + clearTimeout: (handle: number) => void + debounceMs?: number + now?: () => number + debug?: (event: MessageReconciliationDebugEvent) => void +} + +export interface MessageReconciliationDebugEvent { + kind: 'scheduled' | 'started' | 'skipped' | 'merged' | 'failed' + reason: string + timestamp: number + messageCount?: number + error?: string +} + +export interface MessageReconciler { + schedule: (reason: string) => void + runNow: (reason: string) => Promise + dispose: () => void +} + +function normalizeChannelName(value: string | undefined): string | null { + const normalized = value?.trim().replace(/^#/, '') + return normalized || null +} + +function getActiveTab(tabs: AppTab[], activeTabId: string): AppTab | undefined { + return tabs.find((tab) => tab.id === activeTabId) +} + +export function getActiveMessageReconciliationRequest(input: { + activeProjectId: string | null + activeTab?: AppTab + limit?: number +}): MessageReconciliationRequest | null { + const projectId = input.activeTab?.projectId || input.activeProjectId + if (!projectId) return null + + if (input.activeTab?.kind === 'channel') { + const channelName = normalizeChannelName(input.activeTab.channelName) + if (!channelName) return null + return { + projectId, + kind: 'channel', + channelName, + limit: input.limit ?? DEFAULT_RECONCILE_LIMIT + } + } + + if (input.activeTab?.kind === 'dm') { + const conversationId = getDirectMessageRoomId(input.activeTab.dmParticipants || []) + if (!conversationId) return null + return { + projectId, + kind: 'dm', + conversationId, + dmParticipants: input.activeTab.dmParticipants || [], + limit: input.limit ?? DEFAULT_RECONCILE_LIMIT + } + } + + return null +} + +export function createMessageReconciler(deps: MessageReconcilerDeps): MessageReconciler { + const debounceMs = deps.debounceMs ?? DEFAULT_RECONCILE_DEBOUNCE_MS + const now = deps.now ?? (() => Date.now()) + let timer: number | null = null + let disposed = false + let inFlight: Promise | null = null + + const debug = (event: Omit): void => { + deps.debug?.({ ...event, timestamp: now() }) + } + + const runNow = async (reason: string): Promise => { + if (disposed) return + if (inFlight) { + debug({ kind: 'skipped', reason }) + return inFlight + } + + const request = deps.getRequest() + if (!request) { + debug({ kind: 'skipped', reason }) + return + } + + inFlight = (async () => { + debug({ kind: 'started', reason }) + try { + const messages = await deps.reconcileMessages(request) + if (disposed) return + if (messages.length > 0) { + deps.mergeMessages(messages) + } + debug({ kind: 'merged', reason, messageCount: messages.length }) + } catch (err) { + debug({ + kind: 'failed', + reason, + error: err instanceof Error ? err.message : String(err) + }) + } + })() + + try { + await inFlight + } finally { + inFlight = null + } + } + + const schedule = (reason: string): void => { + if (disposed) return + if (timer) { + deps.clearTimeout(timer) + } + debug({ kind: 'scheduled', reason }) + timer = deps.setTimeout(() => { + timer = null + void runNow(reason) + }, debounceMs) + } + + return { + schedule, + runNow, + dispose: () => { + disposed = true + if (timer) { + deps.clearTimeout(timer) + timer = null + } + } + } +} + +function mergeReconciledMessages(messages: ChatMessage[]): void { + const state = useAgentStore.getState() as unknown as StoreWithMessageReconciliation + if (state.reconcileMessages) { + state.reconcileMessages(messages) + return + } + + // Compatibility while the store merge action and IPC surface land together. + for (const message of messages) { + state.handleBrokerEvent({ + kind: 'relay_inbound', + event_id: message.id, + from: message.from, + target: message.to, + body: message.body, + projectId: message.projectId + }) + } +} + +function debugReconciliation(event: MessageReconciliationDebugEvent): void { + if (typeof localStorage === 'undefined') return + if (localStorage.getItem('pear:debug-message-reconciliation') !== '1') return + console.debug('[broker:message-reconciliation]', event) +} + +function refreshEventStream(reason: string): void { + const projectId = useProjectStore.getState().activeProjectId || undefined + const broker = window.pear.broker as PearAPI['broker'] & BrokerWithMessageReconciliation + void broker.refreshEventStream?.(projectId, reason).catch(() => undefined) +} + +export function useMessageReconciliation(): void { + const activeProjectId = useProjectStore((s) => s.activeProjectId) + const activeTabId = useUIStore((s) => s.activeTabId) + const tabs = useUIStore((s) => s.tabs) + const brokerStatus = useAgentStore((s) => s.brokerStatus) + const activeTab = getActiveTab(tabs, activeTabId) + const activeRoomKey = activeTab?.kind === 'channel' + ? `channel:${activeTab.projectId || activeProjectId || ''}:${normalizeChannelName(activeTab.channelName) || ''}` + : activeTab?.kind === 'dm' + ? `dm:${activeTab.projectId || activeProjectId || ''}:${getDirectMessageRoomId(activeTab.dmParticipants || [])}` + : 'none' + + const reconciler = useMemo(() => createMessageReconciler({ + getRequest: () => { + const ui = useUIStore.getState() + const project = useProjectStore.getState() + return getActiveMessageReconciliationRequest({ + activeProjectId: project.activeProjectId, + activeTab: getActiveTab(ui.tabs, ui.activeTabId) + }) + }, + reconcileMessages: (input) => + (window.pear.broker as PearAPI['broker'] & BrokerWithMessageReconciliation).reconcileMessages(input), + mergeMessages: mergeReconciledMessages, + setTimeout: window.setTimeout.bind(window), + clearTimeout: window.clearTimeout.bind(window), + debug: debugReconciliation + }), []) + + useEffect(() => () => reconciler.dispose(), [reconciler]) + + useEffect(() => { + reconciler.schedule('active-room') + }, [activeRoomKey, reconciler]) + + useEffect(() => { + if (brokerStatus === 'connected') { + refreshEventStream('broker-status') + reconciler.schedule('broker-status') + } + }, [brokerStatus, reconciler]) + + useEffect(() => { + const scheduleAfterRefresh = (reason: string): void => { + refreshEventStream(reason) + reconciler.schedule(reason) + } + const onFocus = (): void => scheduleAfterRefresh('window-focus') + const onVisibility = (): void => { + if (document.visibilityState === 'visible') { + scheduleAfterRefresh('document-visible') + } + } + const onPageShow = (): void => scheduleAfterRefresh('pageshow') + const onOnline = (): void => scheduleAfterRefresh('online') + + window.addEventListener('focus', onFocus) + document.addEventListener('visibilitychange', onVisibility) + window.addEventListener('pageshow', onPageShow) + window.addEventListener('online', onOnline) + + return () => { + window.removeEventListener('focus', onFocus) + document.removeEventListener('visibilitychange', onVisibility) + window.removeEventListener('pageshow', onPageShow) + window.removeEventListener('online', onOnline) + } + }, [reconciler]) + + useEffect(() => { + return window.pear.broker.onStatus((status) => { + if (BROKER_CONNECTED_STATUSES.has(status.status)) { + refreshEventStream(`broker:${status.status}`) + reconciler.schedule(`broker:${status.status}`) + } + }) + }, [reconciler]) + + useEffect(() => { + const broker = window.pear.broker as PearAPI['broker'] & BrokerWithMessageReconciliation + if (!broker.onEventStreamDiagnostic) return + return broker.onEventStreamDiagnostic((event) => { + if (EVENT_STREAM_RECONCILED_STATUSES.has(event.status)) { + reconciler.schedule(`event-stream:${event.status}`) + } + }) + }, [reconciler]) +} diff --git a/src/renderer/src/stores/agent-store.ts b/src/renderer/src/stores/agent-store.ts index 003739ef..745d1396 100644 --- a/src/renderer/src/stores/agent-store.ts +++ b/src/renderer/src/stores/agent-store.ts @@ -5,6 +5,7 @@ import type { BrokerDetails, BrokerEventRecord, BrokerListAgent, + BrokerReconciledChatMessage, InboundDeliveryMode, TerminalAttachMode } from '@/lib/ipc' @@ -41,6 +42,7 @@ export interface ChatMessage { timestamp: number isHuman: boolean projectId?: string + conversationId?: string reactions?: ChatReaction[] threadReplies?: ChatThreadReply[] } @@ -320,6 +322,49 @@ function appendJoinNotices(messages: ChatMessage[], notices: ChatMessage[]): Cha return capByCount(nextMessages, MAX_CHAT_MESSAGES) } +function isBrokerDebugEnabled(): boolean { + if (typeof localStorage === 'undefined') return false + return localStorage.getItem('pear-broker-debug') === '1' || + localStorage.getItem('pear-broker-debug') === 'true' +} + +function reconcileChatMessages( + existingMessages: ChatMessage[], + incomingMessages: BrokerReconciledChatMessage[] +): ChatMessage[] { + if (incomingMessages.length === 0) return existingMessages + + const byId = new Map(existingMessages.map((message) => [message.id, message])) + let changed = false + + for (const incoming of incomingMessages) { + const next: ChatMessage = { + ...incoming, + kind: incoming.kind || 'message' + } + const previous = byId.get(next.id) + if (previous) { + byId.set(next.id, { + ...previous, + ...next, + threadReplies: next.threadReplies || previous.threadReplies, + reactions: next.reactions || previous.reactions + }) + changed = true + continue + } + byId.set(next.id, next) + changed = true + } + + if (!changed) return existingMessages + + return capByCount( + Array.from(byId.values()).sort((left, right) => left.timestamp - right.timestamp), + MAX_CHAT_MESSAGES + ) +} + interface AgentState { agents: Agent[] activeAgentKey: string | null @@ -345,6 +390,7 @@ interface AgentState { syncBrokerDetailsStatus: (details: Pick[]) => void hydrateBrokerEvents: (events: BrokerEventRecord[]) => void recordBrokerEvent: (event: BrokerEvent) => void + reconcileMessages: (messages: BrokerReconciledChatMessage[]) => void handleBrokerEvent: (event: BrokerEvent) => void handleBrokerStatus: (status: { projectId?: string; status: string; error?: string }) => void addHumanMessage: (to: string, body: string, projectId?: string) => void @@ -583,6 +629,21 @@ export const useAgentStore = create()(subscribeWithSelector((set, ge }) }, + reconcileMessages: (messages) => { + set((state) => { + const nextMessages = reconcileChatMessages(state.messages, messages) + if (nextMessages === state.messages) return {} + if (isBrokerDebugEnabled()) { + console.info('[broker:renderer-reconcile]', { + incoming: messages.length, + before: state.messages.length, + after: nextMessages.length + }) + } + return { messages: nextMessages } + }) + }, + handleBrokerEvent: (event) => { const { kind } = event @@ -801,6 +862,16 @@ export const useAgentStore = create()(subscribeWithSelector((set, ge ? state.messages : capByCount([...state.messages, msg], MAX_CHAT_MESSAGES) + if (messages !== state.messages && isBrokerDebugEnabled()) { + console.info('[broker:renderer-receipt]', { + projectId, + eventId: msg.id, + kind, + from: msg.from, + to: msg.to + }) + } + return { agents: state.agents.map((a) => { const nextAgent = matchesAgent(a, projectId, eventFrom) ? clearPendingDeliveries(a) : a diff --git a/src/shared/types/ipc.ts b/src/shared/types/ipc.ts index 81ed3383..5b4a4570 100644 --- a/src/shared/types/ipc.ts +++ b/src/shared/types/ipc.ts @@ -201,6 +201,52 @@ export interface PendingRelayMessage { event_id?: string } +export interface BrokerReconciledChatMessage { + id: string + kind?: 'message' | 'notice' + from: string + to: string + body: string + timestamp: number + isHuman: boolean + projectId?: string + conversationId?: string + reactions?: Array<{ + emoji: string + count: number + reactedByHuman: boolean + }> + threadReplies?: Array<{ + id: string + from: string + body: string + timestamp: number + isHuman: boolean + projectId?: string + }> +} + +export interface BrokerReconcileMessagesInput { + projectId: string + kind: 'channel' | 'dm' + channelName?: string + conversationId?: string + dmParticipants?: string[] + limit?: number +} + +export interface BrokerEventStreamDiagnostic { + projectId: string + status: 'received' | 'rebind-started' | 'rebound' | 'rebind-skipped' | 'rebind-error' + reason?: string + at: number + eventKind?: string + eventId?: string + seq?: number + reconnects?: number + error?: string +} + export interface BrokerListAgent { name: string projectId: string @@ -759,6 +805,8 @@ export interface PearAPI { resizePty: (projectId: string | undefined, name: string, rows: number, cols: number) => Promise inputSrtt: (projectId: string | undefined, name: string) => Promise sendMessage: (projectId: string | undefined, input: BrokerSendMessageInput) => Promise + reconcileMessages: (input: BrokerReconcileMessagesInput) => Promise + refreshEventStream: (projectId?: string, reason?: string) => Promise subscribeAgentChannel: (projectId: string | undefined, name: string, channel: string) => Promise unsubscribeAgentChannel: (projectId: string | undefined, name: string, channel: string) => Promise releaseAgent: (projectId: string | undefined, name: string) => Promise @@ -767,6 +815,7 @@ export interface PearAPI { listEvents: () => Promise shutdown: () => Promise onEvent: (callback: (event: unknown) => void) => () => void + onEventStreamDiagnostic: (callback: (event: BrokerEventStreamDiagnostic) => void) => () => void onPtyChunk: (callback: (projectId: string, name: string, chunk: string) => void) => () => void onStatus: (callback: (status: BrokerStatusEvent) => void) => () => void } diff --git a/tsconfig.web.json b/tsconfig.web.json index 4b9842ae..504201b8 100644 --- a/tsconfig.web.json +++ b/tsconfig.web.json @@ -21,5 +21,6 @@ "@shared/*": ["./src/shared/*"] } }, - "include": ["src/renderer/src/**/*", "src/shared/**/*"] + "include": ["src/renderer/src/**/*", "src/shared/**/*"], + "exclude": ["src/**/*.test.ts", "src/**/*.test.tsx"] } diff --git a/vitest.config.mjs b/vitest.config.mjs index 1590b581..af2980ee 100644 --- a/vitest.config.mjs +++ b/vitest.config.mjs @@ -1,6 +1,14 @@ +import { resolve } from 'node:path' + export default { + resolve: { + alias: { + '@': resolve('src/renderer/src'), + '@shared': resolve('src/shared') + } + }, test: { - include: ['src/main/**/*.test.ts', 'packages/**/*.test.ts'], + include: ['src/main/**/*.test.ts', 'src/renderer/src/**/*.test.ts', 'packages/**/*.test.ts'], exclude: ['**/node_modules/**', '**/dist/**', '**/out/**', 'src/main/__tests__/**'] } } From 2e94a15f423eaa5cf4910d6b946f96a963dbcf65 Mon Sep 17 00:00:00 2001 From: Miya Date: Fri, 5 Jun 2026 11:09:20 +0200 Subject: [PATCH 2/4] fix: harden stale message reconciliation --- src/main/broker.test.ts | 112 +++++++++++++++++- src/main/broker.ts | 24 ++-- src/main/ipc-handlers.ts | 4 +- .../hooks/use-message-reconciliation.test.ts | 62 ++++++++++ .../src/hooks/use-message-reconciliation.ts | 59 ++++----- src/renderer/src/stores/agent-store.ts | 49 +++++++- 6 files changed, 260 insertions(+), 50 deletions(-) diff --git a/src/main/broker.test.ts b/src/main/broker.test.ts index f3a1a6ef..eabcb8fc 100644 --- a/src/main/broker.test.ts +++ b/src/main/broker.test.ts @@ -164,9 +164,9 @@ async function startLocal(manager: BrokerManager, agents: string[] = []): Promis return lastSpawned() } -function createMockWindow(): BrowserWindow { +function createMockWindow(destroyed = false): BrowserWindow { return { - isDestroyed: vi.fn(() => false), + isDestroyed: vi.fn(() => destroyed), webContents: { send: vi.fn() } @@ -176,10 +176,11 @@ function createMockWindow(): BrowserWindow { async function startLocalWithWindow( manager: BrokerManager, win: BrowserWindow, - agents: string[] = [] + agents: string[] = [], + projectId = PROJECT_ID ): Promise { mock.state.nextLocalAgents = agents - await manager.start(PROJECT_ID, '/tmp/project-1', 'pear-project-1', win, []) + await manager.start(projectId, `/tmp/${projectId}`, `pear-${projectId}`, win, []) return lastSpawned() } @@ -399,6 +400,32 @@ describe('BrokerManager local + cloud coexistence', () => { await manager.shutdown() }) + it('does not publish broker events to a destroyed captured window', async () => { + const manager = new BrokerManager() + const destroyedWindow = createMockWindow(true) + const local = await startLocalWithWindow(manager, destroyedWindow) + const listener = local.onEvent.mock.calls.at(-1)?.[0] + expect(listener).toBeTypeOf('function') + + listener?.({ + kind: 'relay_inbound', + from: 'codex-2', + target: '#general', + body: 'destroyed window proof', + event_id: 'evt-destroyed-window', + seq: 13 + }) + + expect(destroyedWindow.webContents.send).not.toHaveBeenCalledWith( + 'broker:event', + expect.objectContaining({ + body: 'destroyed window proof' + }) + ) + + await manager.shutdown() + }) + it('refreshEventStream rebinds the harness stream from the last seen sequence', async () => { const manager = new BrokerManager() const local = await startLocal(manager) @@ -422,6 +449,83 @@ describe('BrokerManager local + cloud coexistence', () => { await manager.shutdown() }) + + it('keeps a replacement event listener when reconnect throws during refreshEventStream', async () => { + const manager = new BrokerManager() + const win = createMockWindow() + const local = await startLocalWithWindow(manager, win) + + local.connectEvents.mockImplementationOnce(() => { + throw new Error('connect failed') + }) + + await manager.refreshEventStream(PROJECT_ID, 'test-rebind-failure') + + expect(local.disconnectEvents).toHaveBeenCalledTimes(1) + expect(local.onEvent).toHaveBeenCalledTimes(2) + expect(local.connectEvents).toHaveBeenLastCalledWith(undefined) + + const replacementListener = local.onEvent.mock.calls.at(-1)?.[0] + expect(replacementListener).toBeTypeOf('function') + replacementListener?.({ + kind: 'relay_inbound', + from: 'codex-2', + target: '#general', + body: 'rebind failure still subscribed', + event_id: 'evt-rebind-failure', + seq: 478 + }) + + expect(win.webContents.send).toHaveBeenCalledWith( + 'broker:event', + expect.objectContaining({ + kind: 'relay_inbound', + body: 'rebind failure still subscribed', + projectId: PROJECT_ID + }) + ) + + await manager.shutdown() + }) + + it('does not let a global event-stream refresh overwrite session windows', async () => { + const manager = new BrokerManager() + const firstWindow = createMockWindow() + const secondWindow = createMockWindow() + const intruderWindow = createMockWindow() + await startLocalWithWindow(manager, firstWindow, [], PROJECT_ID) + const secondClient = await startLocalWithWindow(manager, secondWindow, [], 'project-2') + + await manager.refreshEventStream(undefined, 'global-refresh', intruderWindow) + + const secondListener = secondClient.onEvent.mock.calls.at(-1)?.[0] + expect(secondListener).toBeTypeOf('function') + secondListener?.({ + kind: 'relay_inbound', + from: 'codex-2', + target: '#general', + body: 'global refresh proof', + event_id: 'evt-global-refresh', + seq: 479 + }) + + expect(intruderWindow.webContents.send).not.toHaveBeenCalledWith( + 'broker:event', + expect.objectContaining({ + body: 'global refresh proof' + }) + ) + expect(secondWindow.webContents.send).toHaveBeenCalledWith( + 'broker:event', + expect.objectContaining({ + kind: 'relay_inbound', + body: 'global refresh proof', + projectId: 'project-2' + }) + ) + + await manager.shutdown() + }) }) describe('BrokerManager spawnAgent CLI preflight', () => { diff --git a/src/main/broker.ts b/src/main/broker.ts index 16be98fe..4e4029b9 100644 --- a/src/main/broker.ts +++ b/src/main/broker.ts @@ -509,7 +509,7 @@ function isHumanSenderName(sender: string): boolean { } function senderNameFromRelayMessage(message: RelayMessage): string { - return message.from.name || message.from.id || 'unknown' + return message.from?.name || message.from?.id || 'unknown' } function directMessageTargetFromRelayMessage( @@ -564,7 +564,7 @@ function normalizeRelayMessageForChat( reactions: message.reactions?.map((reaction) => ({ emoji: reaction.emoji, count: reaction.count, - reactedByHuman: reaction.agents.some(isHumanSenderName) + reactedByHuman: Array.isArray(reaction.agents) && reaction.agents.some(isHumanSenderName) })) } } @@ -1718,12 +1718,14 @@ export class BrokerManager { } async refreshEventStream(projectId?: string, reason = 'manual', win?: BrowserWindow): Promise { - const sessions = projectId - ? await this.getOrAwaitSessionsForProject(projectId) + const normalizedProjectId = projectId?.trim() + const sessions = normalizedProjectId + ? await this.getOrAwaitSessionsForProject(normalizedProjectId) : Array.from(this.sessions.values()) + const rebindWindow = normalizedProjectId ? win : undefined for (const session of sessions) { - await this.rebindSessionEventStream(sessionKeyFor(session), session, reason, win) + await this.rebindSessionEventStream(sessionKeyFor(session), session, reason, rebindWindow) } } @@ -1764,10 +1766,14 @@ export class BrokerManager { ...(session.lastEventId ? { eventId: session.lastEventId } : {}) }) + const previousUnsubEvent = session.unsubEvent + let nextUnsubEvent: (() => void) | undefined + try { - session.unsubEvent() session.client.disconnectEvents() - session.unsubEvent = this.attachClient(sessionKey, session.client, session.window) + nextUnsubEvent = this.attachClient(sessionKey, session.client, session.window) + session.unsubEvent = nextUnsubEvent + previousUnsubEvent() session.client.connectEvents(session.lastEventSeq) session.lastEventStreamRebindAt = now session.eventStreamRebinds = (session.eventStreamRebinds || 0) + 1 @@ -1790,7 +1796,9 @@ export class BrokerManager { error: toErrorMessage(err), reconnects: session.eventStreamRebinds || 0 }) - throw err + if (!nextUnsubEvent) { + session.unsubEvent = previousUnsubEvent + } } } diff --git a/src/main/ipc-handlers.ts b/src/main/ipc-handlers.ts index de82f1f1..93fb5295 100644 --- a/src/main/ipc-handlers.ts +++ b/src/main/ipc-handlers.ts @@ -264,8 +264,10 @@ export function registerIpcHandlers(): void { }) ipcMain.handle('broker:refresh-event-stream', async (event, projectId?: string, reason?: string) => { + const normalizedProjectId = projectId?.trim() + if (!normalizedProjectId) return const win = BrowserWindow.fromWebContents(event.sender) - await brokerManager.refreshEventStream(projectId, reason || 'renderer-request', win || undefined) + await brokerManager.refreshEventStream(normalizedProjectId, reason || 'renderer-request', win || undefined) }) ipcMain.handle('broker:subscribe-agent-channel', async (_, projectId: string | undefined, name: string, channel: string) => { diff --git a/src/renderer/src/hooks/use-message-reconciliation.test.ts b/src/renderer/src/hooks/use-message-reconciliation.test.ts index 783708a9..557737e2 100644 --- a/src/renderer/src/hooks/use-message-reconciliation.test.ts +++ b/src/renderer/src/hooks/use-message-reconciliation.test.ts @@ -36,6 +36,17 @@ const channelMessage: ChatMessage = { projectId: 'project-1' } +function deferredMessages(): { + promise: Promise + resolve: (messages: ChatMessage[]) => void +} { + let resolve!: (messages: ChatMessage[]) => void + const promise = new Promise((done) => { + resolve = done + }) + return { promise, resolve } +} + describe('getActiveMessageReconciliationRequest', () => { it('builds a bounded channel reconciliation request for the active tab', () => { expect(hooks.getActiveMessageReconciliationRequest({ @@ -134,6 +145,53 @@ describe('createMessageReconciler', () => { })) }) + it('queues a rerun when a trigger fires during an in-flight reconciliation', async () => { + vi.useFakeTimers() + const request: MessageReconciliationRequest = { + projectId: 'project-1', + kind: 'channel', + channelName: 'general', + limit: 50 + } + const firstFetch = deferredMessages() + const secondFetch = deferredMessages() + const secondMessage = { ...channelMessage, id: 'msg-2', body: 'arrived during first fetch' } + const reconcileMessages = vi.fn() + .mockReturnValueOnce(firstFetch.promise) + .mockReturnValueOnce(secondFetch.promise) + const mergeMessages = vi.fn() + const reconciler = hooks.createMessageReconciler({ + getRequest: () => request, + reconcileMessages, + mergeMessages, + setTimeout: setTimeout as unknown as typeof window.setTimeout, + clearTimeout: clearTimeout as unknown as typeof window.clearTimeout, + debounceMs: 1 + }) + + reconciler.schedule('initial') + await vi.advanceTimersByTimeAsync(1) + expect(reconcileMessages).toHaveBeenCalledTimes(1) + + reconciler.schedule('window-focus') + await vi.advanceTimersByTimeAsync(1) + expect(reconcileMessages).toHaveBeenCalledTimes(1) + + firstFetch.resolve([channelMessage]) + await Promise.resolve() + await Promise.resolve() + + expect(reconcileMessages).toHaveBeenCalledTimes(2) + expect(reconcileMessages).toHaveBeenNthCalledWith(2, request) + + secondFetch.resolve([secondMessage]) + await Promise.resolve() + await Promise.resolve() + + expect(mergeMessages).toHaveBeenNthCalledWith(1, [channelMessage]) + expect(mergeMessages).toHaveBeenNthCalledWith(2, [secondMessage]) + }) + it('skips fetches when no active chat room can be reconciled', async () => { vi.useFakeTimers() const reconcileMessages = vi.fn(async () => [channelMessage]) @@ -171,5 +229,9 @@ describe('createMessageReconciler', () => { body: 'updated canonical body', timestamp: channelMessage.timestamp + 1 }) + + const messagesAfterUpdate = agentStore.useAgentStore.getState().messages + store.reconcileMessages([messagesAfterUpdate[0]]) + expect(agentStore.useAgentStore.getState().messages).toBe(messagesAfterUpdate) }) }) diff --git a/src/renderer/src/hooks/use-message-reconciliation.ts b/src/renderer/src/hooks/use-message-reconciliation.ts index c950e25b..a42aa9ed 100644 --- a/src/renderer/src/hooks/use-message-reconciliation.ts +++ b/src/renderer/src/hooks/use-message-reconciliation.ts @@ -28,16 +28,11 @@ const EVENT_STREAM_RECONCILED_STATUSES = new Set Promise + reconcileMessages?: (input: MessageReconciliationRequest) => Promise refreshEventStream?: (projectId?: string, reason?: string) => Promise onEventStreamDiagnostic?: (callback: (event: BrokerEventStreamDiagnostic) => void) => () => void } -interface StoreWithMessageReconciliation { - reconcileMessages?: (messages: ChatMessage[]) => void - handleBrokerEvent: (event: Record & { kind: string }) => void -} - interface MessageReconcilerDeps { getRequest: () => MessageReconciliationRequest | null reconcileMessages: (input: MessageReconciliationRequest) => Promise @@ -112,6 +107,7 @@ export function createMessageReconciler(deps: MessageReconcilerDeps): MessageRec let timer: number | null = null let disposed = false let inFlight: Promise | null = null + let pendingRerun: string | null = null const debug = (event: Omit): void => { deps.debug?.({ ...event, timestamp: now() }) @@ -120,7 +116,8 @@ export function createMessageReconciler(deps: MessageReconcilerDeps): MessageRec const runNow = async (reason: string): Promise => { if (disposed) return if (inFlight) { - debug({ kind: 'skipped', reason }) + pendingRerun = reason + debug({ kind: 'scheduled', reason }) return inFlight } @@ -130,7 +127,7 @@ export function createMessageReconciler(deps: MessageReconcilerDeps): MessageRec return } - inFlight = (async () => { + const currentRun = (async () => { debug({ kind: 'started', reason }) try { const messages = await deps.reconcileMessages(request) @@ -147,11 +144,19 @@ export function createMessageReconciler(deps: MessageReconcilerDeps): MessageRec }) } })() + inFlight = currentRun try { - await inFlight + await currentRun } finally { - inFlight = null + if (inFlight === currentRun) { + inFlight = null + } + const rerunReason = pendingRerun + pendingRerun = null + if (rerunReason && !disposed) { + await runNow(rerunReason) + } } } @@ -181,23 +186,7 @@ export function createMessageReconciler(deps: MessageReconcilerDeps): MessageRec } function mergeReconciledMessages(messages: ChatMessage[]): void { - const state = useAgentStore.getState() as unknown as StoreWithMessageReconciliation - if (state.reconcileMessages) { - state.reconcileMessages(messages) - return - } - - // Compatibility while the store merge action and IPC surface land together. - for (const message of messages) { - state.handleBrokerEvent({ - kind: 'relay_inbound', - event_id: message.id, - from: message.from, - target: message.to, - body: message.body, - projectId: message.projectId - }) - } + useAgentStore.getState().reconcileMessages(messages) } function debugReconciliation(event: MessageReconciliationDebugEvent): void { @@ -208,8 +197,8 @@ function debugReconciliation(event: MessageReconciliationDebugEvent): void { function refreshEventStream(reason: string): void { const projectId = useProjectStore.getState().activeProjectId || undefined - const broker = window.pear.broker as PearAPI['broker'] & BrokerWithMessageReconciliation - void broker.refreshEventStream?.(projectId, reason).catch(() => undefined) + const broker = window.pear?.broker as (PearAPI['broker'] & BrokerWithMessageReconciliation) | undefined + void broker?.refreshEventStream?.(projectId, reason)?.catch(() => undefined) } export function useMessageReconciliation(): void { @@ -233,8 +222,10 @@ export function useMessageReconciliation(): void { activeTab: getActiveTab(ui.tabs, ui.activeTabId) }) }, - reconcileMessages: (input) => - (window.pear.broker as PearAPI['broker'] & BrokerWithMessageReconciliation).reconcileMessages(input), + reconcileMessages: (input) => { + const broker = window.pear?.broker as (PearAPI['broker'] & BrokerWithMessageReconciliation) | undefined + return broker?.reconcileMessages?.(input) ?? Promise.resolve([]) + }, mergeMessages: mergeReconciledMessages, setTimeout: window.setTimeout.bind(window), clearTimeout: window.clearTimeout.bind(window), @@ -282,7 +273,7 @@ export function useMessageReconciliation(): void { }, [reconciler]) useEffect(() => { - return window.pear.broker.onStatus((status) => { + return window.pear?.broker?.onStatus?.((status) => { if (BROKER_CONNECTED_STATUSES.has(status.status)) { refreshEventStream(`broker:${status.status}`) reconciler.schedule(`broker:${status.status}`) @@ -291,8 +282,8 @@ export function useMessageReconciliation(): void { }, [reconciler]) useEffect(() => { - const broker = window.pear.broker as PearAPI['broker'] & BrokerWithMessageReconciliation - if (!broker.onEventStreamDiagnostic) return + const broker = window.pear?.broker as (PearAPI['broker'] & BrokerWithMessageReconciliation) | undefined + if (!broker?.onEventStreamDiagnostic) return return broker.onEventStreamDiagnostic((event) => { if (EVENT_STREAM_RECONCILED_STATUSES.has(event.status)) { reconciler.schedule(`event-stream:${event.status}`) diff --git a/src/renderer/src/stores/agent-store.ts b/src/renderer/src/stores/agent-store.ts index 745d1396..c74d3e63 100644 --- a/src/renderer/src/stores/agent-store.ts +++ b/src/renderer/src/stores/agent-store.ts @@ -328,6 +328,46 @@ function isBrokerDebugEnabled(): boolean { localStorage.getItem('pear-broker-debug') === 'true' } +function reactionsEqual(left: ChatReaction[] | undefined, right: ChatReaction[] | undefined): boolean { + if (left === right) return true + if (!left || !right || left.length !== right.length) return false + return left.every((reaction, index) => { + const candidate = right[index] + return candidate && + reaction.emoji === candidate.emoji && + reaction.count === candidate.count && + reaction.reactedByHuman === candidate.reactedByHuman + }) +} + +function threadRepliesEqual(left: ChatThreadReply[] | undefined, right: ChatThreadReply[] | undefined): boolean { + if (left === right) return true + if (!left || !right || left.length !== right.length) return false + return left.every((reply, index) => { + const candidate = right[index] + return candidate && + reply.id === candidate.id && + reply.from === candidate.from && + reply.body === candidate.body && + reply.timestamp === candidate.timestamp && + reply.isHuman === candidate.isHuman && + reply.projectId === candidate.projectId + }) +} + +function chatMessagesEqual(left: ChatMessage, right: ChatMessage): boolean { + return left.kind === right.kind && + left.from === right.from && + left.to === right.to && + left.body === right.body && + left.timestamp === right.timestamp && + left.isHuman === right.isHuman && + left.projectId === right.projectId && + left.conversationId === right.conversationId && + reactionsEqual(left.reactions, right.reactions) && + threadRepliesEqual(left.threadReplies, right.threadReplies) +} + function reconcileChatMessages( existingMessages: ChatMessage[], incomingMessages: BrokerReconciledChatMessage[] @@ -344,13 +384,16 @@ function reconcileChatMessages( } const previous = byId.get(next.id) if (previous) { - byId.set(next.id, { + const merged = { ...previous, ...next, threadReplies: next.threadReplies || previous.threadReplies, reactions: next.reactions || previous.reactions - }) - changed = true + } + if (!chatMessagesEqual(previous, merged)) { + byId.set(next.id, merged) + changed = true + } continue } byId.set(next.id, next) From bf9e51b1d3aefe1a8434d109fb39cec10f560bc2 Mon Sep 17 00:00:00 2001 From: Miya Date: Fri, 5 Jun 2026 11:18:46 +0200 Subject: [PATCH 3/4] fix: reconcile after human message sends --- .../hooks/use-message-reconciliation.test.ts | 40 +++++++++++++++++++ .../src/hooks/use-message-reconciliation.ts | 14 +++++++ src/renderer/src/stores/agent-store.ts | 8 +++- 3 files changed, 60 insertions(+), 2 deletions(-) diff --git a/src/renderer/src/hooks/use-message-reconciliation.test.ts b/src/renderer/src/hooks/use-message-reconciliation.test.ts index 557737e2..56fa6302 100644 --- a/src/renderer/src/hooks/use-message-reconciliation.test.ts +++ b/src/renderer/src/hooks/use-message-reconciliation.test.ts @@ -95,10 +95,27 @@ describe('getActiveMessageReconciliationRequest', () => { } })).toBeNull() }) + + it('schedules reconciliation after a human message send timestamp changes', () => { + const reconciler = { schedule: vi.fn() } + + hooks.scheduleHumanMessageSentReconciliation({ + lastHumanMessageSentAt: 0, + reconciler + }) + hooks.scheduleHumanMessageSentReconciliation({ + lastHumanMessageSentAt: 1_717_000_000_000, + reconciler + }) + + expect(reconciler.schedule).toHaveBeenCalledTimes(1) + expect(reconciler.schedule).toHaveBeenCalledWith('human-message-sent') + }) }) describe('createMessageReconciler', () => { afterEach(() => { + vi.restoreAllMocks() vi.useRealTimers() agentStore.useAgentStore.getState().clearAll() }) @@ -234,4 +251,27 @@ describe('createMessageReconciler', () => { store.reconcileMessages([messagesAfterUpdate[0]]) expect(agentStore.useAgentStore.getState().messages).toBe(messagesAfterUpdate) }) + + it('tracks optimistic human sends for reconciliation triggers', () => { + const now = 1_717_000_123_000 + vi.spyOn(Date, 'now').mockReturnValue(now) + + agentStore.useAgentStore.getState().addHumanMessage('#general', 'hello from human', 'project-1') + + const state = agentStore.useAgentStore.getState() + expect(state.lastHumanMessageSentAt).toBe(now) + expect(state.messages[0]).toMatchObject({ + from: 'human', + to: '#general', + body: 'hello from human', + projectId: 'project-1' + }) + }) + + it('wires human message sends into the reconciliation hook', () => { + const source = hooks.useMessageReconciliation.toString() + expect(source).toContain('s.lastHumanMessageSentAt') + expect(source).toMatch(/scheduleHumanMessageSentReconciliation\([\s\S]*lastHumanMessageSentAt[\s\S]*reconciler/) + expect(source).toMatch(/\[lastHumanMessageSentAt,\s*reconciler\]/) + }) }) diff --git a/src/renderer/src/hooks/use-message-reconciliation.ts b/src/renderer/src/hooks/use-message-reconciliation.ts index a42aa9ed..a136ceca 100644 --- a/src/renderer/src/hooks/use-message-reconciliation.ts +++ b/src/renderer/src/hooks/use-message-reconciliation.ts @@ -58,6 +58,15 @@ export interface MessageReconciler { dispose: () => void } +export function scheduleHumanMessageSentReconciliation(input: { + lastHumanMessageSentAt: number + reconciler: Pick +}): void { + if (input.lastHumanMessageSentAt > 0) { + input.reconciler.schedule('human-message-sent') + } +} + function normalizeChannelName(value: string | undefined): string | null { const normalized = value?.trim().replace(/^#/, '') return normalized || null @@ -206,6 +215,7 @@ export function useMessageReconciliation(): void { const activeTabId = useUIStore((s) => s.activeTabId) const tabs = useUIStore((s) => s.tabs) const brokerStatus = useAgentStore((s) => s.brokerStatus) + const lastHumanMessageSentAt = useAgentStore((s) => s.lastHumanMessageSentAt) const activeTab = getActiveTab(tabs, activeTabId) const activeRoomKey = activeTab?.kind === 'channel' ? `channel:${activeTab.projectId || activeProjectId || ''}:${normalizeChannelName(activeTab.channelName) || ''}` @@ -245,6 +255,10 @@ export function useMessageReconciliation(): void { } }, [brokerStatus, reconciler]) + useEffect(() => { + scheduleHumanMessageSentReconciliation({ lastHumanMessageSentAt, reconciler }) + }, [lastHumanMessageSentAt, reconciler]) + useEffect(() => { const scheduleAfterRefresh = (reason: string): void => { refreshEventStream(reason) diff --git a/src/renderer/src/stores/agent-store.ts b/src/renderer/src/stores/agent-store.ts index c74d3e63..70e68e8b 100644 --- a/src/renderer/src/stores/agent-store.ts +++ b/src/renderer/src/stores/agent-store.ts @@ -417,6 +417,7 @@ interface AgentState { brokerError: string | null brokerErrors: BrokerErrorEntry[] brokerEvents: BrokerEventRecord[] + lastHumanMessageSentAt: number setActiveAgentKey: (key: string | null) => void markAgentActive: (projectId: string | undefined, name: string) => void @@ -460,6 +461,7 @@ export const useAgentStore = create()(subscribeWithSelector((set, ge brokerError: null, brokerErrors: [], brokerEvents: [], + lastHumanMessageSentAt: 0, setActiveAgentKey: (key) => set({ activeAgentKey: key }), @@ -1002,7 +1004,8 @@ export const useAgentStore = create()(subscribeWithSelector((set, ge set((state) => ({ messages: isDuplicateHumanEcho(state.messages, msg) ? state.messages - : capByCount([...state.messages, msg], MAX_CHAT_MESSAGES) + : capByCount([...state.messages, msg], MAX_CHAT_MESSAGES), + lastHumanMessageSentAt: timestamp })) }, @@ -1147,7 +1150,8 @@ export const useAgentStore = create()(subscribeWithSelector((set, ge brokerStatus: 'disconnected', brokerError: null, brokerErrors: [], - brokerEvents: [] + brokerEvents: [], + lastHumanMessageSentAt: 0 }), getAgentBuffer: (projectId, name) => getPtyChunks(getAgentKey(projectId, name)) From 51c6e6051244e9a537754b6d77b81669f10ecdd7 Mon Sep 17 00:00:00 2001 From: Miya Date: Fri, 5 Jun 2026 11:22:58 +0200 Subject: [PATCH 4/4] Revert "fix: reconcile after human message sends" This reverts commit bf9e51b1d3aefe1a8434d109fb39cec10f560bc2. --- .../hooks/use-message-reconciliation.test.ts | 40 ------------------- .../src/hooks/use-message-reconciliation.ts | 14 ------- src/renderer/src/stores/agent-store.ts | 8 +--- 3 files changed, 2 insertions(+), 60 deletions(-) diff --git a/src/renderer/src/hooks/use-message-reconciliation.test.ts b/src/renderer/src/hooks/use-message-reconciliation.test.ts index 56fa6302..557737e2 100644 --- a/src/renderer/src/hooks/use-message-reconciliation.test.ts +++ b/src/renderer/src/hooks/use-message-reconciliation.test.ts @@ -95,27 +95,10 @@ describe('getActiveMessageReconciliationRequest', () => { } })).toBeNull() }) - - it('schedules reconciliation after a human message send timestamp changes', () => { - const reconciler = { schedule: vi.fn() } - - hooks.scheduleHumanMessageSentReconciliation({ - lastHumanMessageSentAt: 0, - reconciler - }) - hooks.scheduleHumanMessageSentReconciliation({ - lastHumanMessageSentAt: 1_717_000_000_000, - reconciler - }) - - expect(reconciler.schedule).toHaveBeenCalledTimes(1) - expect(reconciler.schedule).toHaveBeenCalledWith('human-message-sent') - }) }) describe('createMessageReconciler', () => { afterEach(() => { - vi.restoreAllMocks() vi.useRealTimers() agentStore.useAgentStore.getState().clearAll() }) @@ -251,27 +234,4 @@ describe('createMessageReconciler', () => { store.reconcileMessages([messagesAfterUpdate[0]]) expect(agentStore.useAgentStore.getState().messages).toBe(messagesAfterUpdate) }) - - it('tracks optimistic human sends for reconciliation triggers', () => { - const now = 1_717_000_123_000 - vi.spyOn(Date, 'now').mockReturnValue(now) - - agentStore.useAgentStore.getState().addHumanMessage('#general', 'hello from human', 'project-1') - - const state = agentStore.useAgentStore.getState() - expect(state.lastHumanMessageSentAt).toBe(now) - expect(state.messages[0]).toMatchObject({ - from: 'human', - to: '#general', - body: 'hello from human', - projectId: 'project-1' - }) - }) - - it('wires human message sends into the reconciliation hook', () => { - const source = hooks.useMessageReconciliation.toString() - expect(source).toContain('s.lastHumanMessageSentAt') - expect(source).toMatch(/scheduleHumanMessageSentReconciliation\([\s\S]*lastHumanMessageSentAt[\s\S]*reconciler/) - expect(source).toMatch(/\[lastHumanMessageSentAt,\s*reconciler\]/) - }) }) diff --git a/src/renderer/src/hooks/use-message-reconciliation.ts b/src/renderer/src/hooks/use-message-reconciliation.ts index a136ceca..a42aa9ed 100644 --- a/src/renderer/src/hooks/use-message-reconciliation.ts +++ b/src/renderer/src/hooks/use-message-reconciliation.ts @@ -58,15 +58,6 @@ export interface MessageReconciler { dispose: () => void } -export function scheduleHumanMessageSentReconciliation(input: { - lastHumanMessageSentAt: number - reconciler: Pick -}): void { - if (input.lastHumanMessageSentAt > 0) { - input.reconciler.schedule('human-message-sent') - } -} - function normalizeChannelName(value: string | undefined): string | null { const normalized = value?.trim().replace(/^#/, '') return normalized || null @@ -215,7 +206,6 @@ export function useMessageReconciliation(): void { const activeTabId = useUIStore((s) => s.activeTabId) const tabs = useUIStore((s) => s.tabs) const brokerStatus = useAgentStore((s) => s.brokerStatus) - const lastHumanMessageSentAt = useAgentStore((s) => s.lastHumanMessageSentAt) const activeTab = getActiveTab(tabs, activeTabId) const activeRoomKey = activeTab?.kind === 'channel' ? `channel:${activeTab.projectId || activeProjectId || ''}:${normalizeChannelName(activeTab.channelName) || ''}` @@ -255,10 +245,6 @@ export function useMessageReconciliation(): void { } }, [brokerStatus, reconciler]) - useEffect(() => { - scheduleHumanMessageSentReconciliation({ lastHumanMessageSentAt, reconciler }) - }, [lastHumanMessageSentAt, reconciler]) - useEffect(() => { const scheduleAfterRefresh = (reason: string): void => { refreshEventStream(reason) diff --git a/src/renderer/src/stores/agent-store.ts b/src/renderer/src/stores/agent-store.ts index 70e68e8b..c74d3e63 100644 --- a/src/renderer/src/stores/agent-store.ts +++ b/src/renderer/src/stores/agent-store.ts @@ -417,7 +417,6 @@ interface AgentState { brokerError: string | null brokerErrors: BrokerErrorEntry[] brokerEvents: BrokerEventRecord[] - lastHumanMessageSentAt: number setActiveAgentKey: (key: string | null) => void markAgentActive: (projectId: string | undefined, name: string) => void @@ -461,7 +460,6 @@ export const useAgentStore = create()(subscribeWithSelector((set, ge brokerError: null, brokerErrors: [], brokerEvents: [], - lastHumanMessageSentAt: 0, setActiveAgentKey: (key) => set({ activeAgentKey: key }), @@ -1004,8 +1002,7 @@ export const useAgentStore = create()(subscribeWithSelector((set, ge set((state) => ({ messages: isDuplicateHumanEcho(state.messages, msg) ? state.messages - : capByCount([...state.messages, msg], MAX_CHAT_MESSAGES), - lastHumanMessageSentAt: timestamp + : capByCount([...state.messages, msg], MAX_CHAT_MESSAGES) })) }, @@ -1150,8 +1147,7 @@ export const useAgentStore = create()(subscribeWithSelector((set, ge brokerStatus: 'disconnected', brokerError: null, brokerErrors: [], - brokerEvents: [], - lastHumanMessageSentAt: 0 + brokerEvents: [] }), getAgentBuffer: (projectId, name) => getPtyChunks(getAgentKey(projectId, name))