Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
152 changes: 152 additions & 0 deletions scripts/repro-stale-message-reconciliation.mjs
Original file line number Diff line number Diff line change
@@ -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
})
186 changes: 186 additions & 0 deletions src/main/broker.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand All @@ -18,6 +19,7 @@ type MockClient = {
onEvent: ReturnType<typeof vi.fn>
addListener: ReturnType<typeof vi.fn>
connectEvents: ReturnType<typeof vi.fn>
disconnectEvents: ReturnType<typeof vi.fn>
renewLease: ReturnType<typeof vi.fn>
shutdown: ReturnType<typeof vi.fn>
disconnect: ReturnType<typeof vi.fn>
Expand Down Expand Up @@ -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(),
Expand Down Expand Up @@ -161,6 +164,26 @@ async function startLocal(manager: BrokerManager, agents: string[] = []): Promis
return lastSpawned()
}

function createMockWindow(destroyed = false): BrowserWindow {
return {
isDestroyed: vi.fn(() => destroyed),
webContents: {
send: vi.fn()
}
} as unknown as BrowserWindow
}

async function startLocalWithWindow(
manager: BrokerManager,
win: BrowserWindow,
agents: string[] = [],
projectId = PROJECT_ID
): Promise<MockClient> {
mock.state.nextLocalAgents = agents
await manager.start(projectId, `/tmp/${projectId}`, `pear-${projectId}`, win, [])
return lastSpawned()
}

async function attachCloud(manager: BrokerManager, agents: string[] = []): Promise<MockClient> {
mock.state.nextCloudAgents = agents
await manager.attachCloudSandbox(PROJECT_ID, {
Expand Down Expand Up @@ -340,6 +363,169 @@ 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<typeof vi.fn>).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('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)
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()
})

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', () => {
Expand Down
Loading
Loading