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
124 changes: 122 additions & 2 deletions src/main/broker.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -26,7 +26,9 @@ type MockClient = {
release: ReturnType<typeof vi.fn>
subscribeChannels: ReturnType<typeof vi.fn>
unsubscribeChannels: ReturnType<typeof vi.fn>
getStatus: ReturnType<typeof vi.fn>
brokerPid?: number
baseUrl?: string
agentNames: string[]
}

Expand All @@ -45,6 +47,10 @@ const mock = vi.hoisted(() => {
snapshot: vi.fn(async () => ({ rows: 24, cols: 80, cursor: { x: 0, y: 0 }, screen: 'aGVsbG8=' })),
resizePty: vi.fn(async () => undefined),
getPending: vi.fn(async () => []),
getStatus: vi.fn(async () => ({
agents: client.agentNames.map((name) => ({ name, runtime: 'pty', channels: [] })),
pending_delivery_count: 0
})),
onEvent: vi.fn(() => () => undefined),
addListener: vi.fn(() => () => undefined),
connectEvents: vi.fn(),
Expand All @@ -63,8 +69,11 @@ const mock = vi.hoisted(() => {
const state = {
spawnedClients: [] as MockClient[],
constructedClients: [] as MockClient[],
connectedClients: [] as MockClient[],
nextLocalAgents: [] as string[],
nextCloudAgents: [] as string[]
nextCloudAgents: [] as string[],
nextConnectedAgents: [] as string[],
nextConnectedSessionErrors: [] as Error[]
}

class HarnessDriverClient {
Expand All @@ -75,7 +84,13 @@ const mock = vi.hoisted(() => {
})

static connect = vi.fn(() => {
throw new Error('not used in tests')
const client = createMockClient(state.nextConnectedAgents.splice(0))
const sessionError = state.nextConnectedSessionErrors.shift()
if (sessionError) {
client.getSession.mockRejectedValueOnce(sessionError)
}
state.connectedClients.push(client)
return client
})

constructor() {
Expand Down Expand Up @@ -314,9 +329,13 @@ describe('BrokerManager local + cloud coexistence', () => {
beforeEach(() => {
mock.state.spawnedClients.length = 0
mock.state.constructedClients.length = 0
mock.state.connectedClients.length = 0
mock.state.nextLocalAgents = []
mock.state.nextCloudAgents = []
mock.state.nextConnectedAgents = []
mock.state.nextConnectedSessionErrors = []
mock.HarnessDriverClient.spawn.mockClear()
mock.HarnessDriverClient.connect.mockClear()
})

it('keeps the local session alive when a cloud sandbox attaches', async () => {
Expand Down Expand Up @@ -354,6 +373,107 @@ describe('BrokerManager local + cloud coexistence', () => {
await manager.shutdown()
})

it('reuses current harness-driver connection files instead of spawning another broker', async () => {
const tempDir = await mkdtemp(join(tmpdir(), 'pear-current-connection-'))
const connectionPath = join(tempDir, '.agentworkforce', 'relay', 'connection.json')
await mkdir(dirname(connectionPath), { recursive: true })
await writeFile(connectionPath, JSON.stringify({
url: 'http://127.0.0.1:43210',
apiKey: 'test-key',
pid: 4242
}))

try {
const manager = new BrokerManager()
mock.state.nextConnectedAgents = ['codex-1']

const started = await manager.start(PROJECT_ID, tempDir, 'pear-project-1', undefined as never, [])

expect(started).toBe(true)
expect(mock.HarnessDriverClient.connect).toHaveBeenCalledWith({ cwd: tempDir, connectionPath })
expect(mock.HarnessDriverClient.spawn).not.toHaveBeenCalled()
expect(mock.state.connectedClients).toHaveLength(1)

await manager.shutdown()
} finally {
await rm(tempDir, { recursive: true, force: true })
}
})

it('reports the matching connection file when a stale current file and matching legacy file coexist', async () => {
const tempDir = await mkdtemp(join(tmpdir(), 'pear-connection-status-'))
const currentConnectionPath = join(tempDir, '.agentworkforce', 'relay', 'connection.json')
const legacyConnectionPath = join(tempDir, '.agent-relay', 'connection.json')
await mkdir(dirname(currentConnectionPath), { recursive: true })
await mkdir(dirname(legacyConnectionPath), { recursive: true })

try {
const manager = new BrokerManager()
mock.state.nextLocalAgents = []
await manager.start(PROJECT_ID, tempDir, 'pear-project-1', undefined as never, [])
const local = lastSpawned()
local.baseUrl = 'http://127.0.0.1:43210'
local.brokerPid = 4242

await writeFile(currentConnectionPath, JSON.stringify({
url: 'http://127.0.0.1:1',
api_key: 'stale-key',
pid: 1
}))
await writeFile(legacyConnectionPath, JSON.stringify({
url: 'http://127.0.0.1:43210/',
api_key: 'legacy-key',
pid: 4242
}))

const [details] = await manager.listBrokerDetails()

expect(details.connectionPath).toBe(legacyConnectionPath)
expect(details.connectionFileStatus).toBe('matches')
expect(details.apiKey).toBe('legacy-key')

await manager.shutdown()
} finally {
await rm(tempDir, { recursive: true, force: true })
}
})

it('disconnects failed broker connection probes before trying the next candidate', async () => {
const tempDir = await mkdtemp(join(tmpdir(), 'pear-connection-fallback-'))
const currentConnectionPath = join(tempDir, '.agentworkforce', 'relay', 'connection.json')
const legacyConnectionPath = join(tempDir, '.agent-relay', 'connection.json')
await mkdir(dirname(currentConnectionPath), { recursive: true })
await mkdir(dirname(legacyConnectionPath), { recursive: true })
await writeFile(currentConnectionPath, JSON.stringify({
url: 'http://127.0.0.1:1',
api_key: 'stale-key',
pid: 1
}))
await writeFile(legacyConnectionPath, JSON.stringify({
url: 'http://127.0.0.1:43210',
api_key: 'legacy-key',
pid: 4242
}))

try {
const manager = new BrokerManager()
mock.state.nextConnectedSessionErrors = [new Error('stale broker')]

const started = await manager.start(PROJECT_ID, tempDir, 'pear-project-1', undefined as never, [])

expect(started).toBe(true)
expect(mock.HarnessDriverClient.connect).toHaveBeenCalledTimes(2)
expect(mock.HarnessDriverClient.connect).toHaveBeenNthCalledWith(1, { cwd: tempDir, connectionPath: currentConnectionPath })
expect(mock.HarnessDriverClient.connect).toHaveBeenNthCalledWith(2, { cwd: tempDir, connectionPath: legacyConnectionPath })
expect(mock.state.connectedClients[0].disconnect).toHaveBeenCalledTimes(1)
expect(mock.HarnessDriverClient.spawn).not.toHaveBeenCalled()

await manager.shutdown()
} finally {
await rm(tempDir, { recursive: true, force: true })
}
})

it('keeps the cloud session alive when a local broker starts afterwards', async () => {
const manager = new BrokerManager()
const cloud = await attachCloud(manager, ['cloud-agent'])
Expand Down
72 changes: 56 additions & 16 deletions src/main/broker.ts
Original file line number Diff line number Diff line change
Expand Up @@ -718,7 +718,25 @@ function getBrokerConnectionFileInfo(
return { hasApiKey: false }
}

const connectionPath = join(cwd, '.agent-relay', 'connection.json')
const candidates = brokerConnectionPathCandidates(cwd)
let firstExisting: BrokerConnectionFileInfo | undefined

for (const connectionPath of candidates) {
if (!existsSync(connectionPath)) continue

const info = readBrokerConnectionFileInfo(connectionPath, baseUrl, brokerPid)
if (info.status === 'matches') return info
firstExisting ??= info
}

return firstExisting ?? { path: candidates[0], status: 'missing', hasApiKey: false }
}

function readBrokerConnectionFileInfo(
connectionPath: string,
baseUrl: string | undefined,
brokerPid: number | undefined
): BrokerConnectionFileInfo {
if (!existsSync(connectionPath)) {
return { path: connectionPath, status: 'missing', hasApiKey: false }
}
Expand All @@ -745,6 +763,18 @@ function getBrokerConnectionFileInfo(
}
}

function brokerConnectionPathCandidates(cwd: string): string[] {
return [
join(cwd, '.agentworkforce', 'relay', 'connection.json'),
join(cwd, '.agent-relay', 'connection.json')
]
}

function resolveBrokerConnectionPath(cwd: string): string {
const candidates = brokerConnectionPathCandidates(cwd)
return candidates.find((candidate) => existsSync(candidate)) ?? candidates[0]
}
Comment on lines +773 to +776

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

medium

Avoid calling brokerConnectionPathCandidates(cwd) twice to prevent redundant array allocations and path joins. Storing the candidates in a local variable is cleaner and more efficient.

function resolveBrokerConnectionPath(cwd: string): string {
  const candidates = brokerConnectionPathCandidates(cwd)
  return candidates.find((candidate) => existsSync(candidate)) ?? candidates[0]
}


function toInboundDeliveryMode(mode?: TerminalAttachMode): InboundDeliveryMode {
return mode === 'drive' ? 'manual_flush' : 'auto_inject'
}
Expand Down Expand Up @@ -845,14 +875,19 @@ function getBrokerRuntimeSafeName(name: string): string {
}

function getBrokerRuntimeCleanupPaths(cwd: string, brokerName: string): string[] {
const root = join(cwd, '.agent-relay')
const legacyRoot = join(cwd, '.agent-relay')
const currentRoot = join(cwd, '.agentworkforce', 'relay')
const safeName = getBrokerRuntimeSafeName(brokerName)

return [
join(root, 'connection.json'),
join(root, `broker-${safeName}.lock`),
join(root, `state-${safeName}.json`),
join(root, `pending-${safeName}.json`)
join(currentRoot, 'connection.json'),
join(currentRoot, `broker-${safeName}.lock`),
join(currentRoot, `state-${safeName}.json`),
join(currentRoot, `pending-${safeName}.json`),
join(legacyRoot, 'connection.json'),
join(legacyRoot, `broker-${safeName}.lock`),
join(legacyRoot, `state-${safeName}.json`),
join(legacyRoot, `pending-${safeName}.json`)
]
}

Expand Down Expand Up @@ -1268,20 +1303,25 @@ export class BrokerManager {
}

private async connectExistingBroker(projectId: string, cwd: string): Promise<AgentRelayClient | null> {
const connectionPath = join(cwd, '.agent-relay', 'connection.json')
if (!existsSync(connectionPath)) {
const connectionPaths = brokerConnectionPathCandidates(cwd).filter((candidate) => existsSync(candidate))
if (connectionPaths.length === 0) {
return null
}

try {
const client = AgentRelayClient.connect({ cwd })
await client.getSession()
console.log(`[broker] Reusing existing broker for project ${projectId}: ${connectionPath}`)
return client
} catch (err) {
console.warn(`[broker] Existing broker connection is not reusable for project ${projectId}:`, err)
return null
for (const connectionPath of connectionPaths) {
let client: AgentRelayClient | undefined
try {
client = AgentRelayClient.connect({ cwd, connectionPath })
await client.getSession()
console.log(`[broker] Reusing existing broker for project ${projectId}: ${connectionPath}`)
return client
} catch (err) {
client?.disconnect()
console.warn(`[broker] Existing broker connection is not reusable for project ${projectId}:`, err)
}
}
Comment on lines +1311 to 1322

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Suggestion: When probing multiple connection files, a client instance is created before health-checking with getSession(), but the failure path only logs and continues. If connect() allocates sockets/listeners, failed attempts are never torn down, so repeated starts can leak broker connections/resources. Explicitly disconnect/release the temporary client inside the catch path before trying the next candidate. [resource leak]

Severity Level: Major ⚠️
- ❌ Main process can leak AgentRelayClient sockets on failed reuse.
- ⚠️ Repeated broker:start calls may exhaust broker connection resources.
- ⚠️ Long-lived sessions risk instability from accumulated leaked clients.
Steps of Reproduction ✅
1. Open the Electron app and trigger a broker start for a project, which sends the IPC
call `'broker:start'` handled in `src/main/ipc-handlers.ts` (lines 170–199 from tool
output) where `brokerManager.start(projectId, cwd, name, win, channels)` is invoked.

2. Inside `BrokerManager.start` in `src/main/broker.ts` (lines 52–79 from tool output),
the `startBroker` closure is created; its first step is `const existingClient = await
this.connectExistingBroker(normalizedProjectId, cwd)`, which calls
`connectExistingBroker`.

3. Ensure the project `cwd` contains at least one stale or unusable local broker
connection file (for example, `.agentworkforce/relay/connection.json` left over from a
previous run) so that `brokerConnectionPathCandidates(cwd).filter((candidate) =>
existsSync(candidate))` in `connectExistingBroker` (file `src/main/broker.ts`, around
lines 69–75) returns one or more paths that exist but point to a dead or incompatible
broker.

4. When `connectExistingBroker` runs, for each such `connectionPath` it executes the loop
shown in `src/main/broker.ts` lines 1293–1302: it calls `const client =
AgentRelayClient.connect({ cwd, connectionPath })` and then `await client.getSession()`.
If the broker behind that file is unreachable or mismatched, `client.getSession()` throws,
execution enters the `catch (err)` block, and only a warning is logged
(`console.warn('[broker] Existing broker connection is not reusable…')`) without calling
`client.disconnect()` or any shutdown.

5. Because the failing `client` is never stored in `this.sessions` and never explicitly
disconnected in the catch block, any sockets/listeners allocated by
`AgentRelayClient.connect` remain live in the Electron main process. The loop then
continues to the next candidate (if any), or returns `null`, causing `startBroker` to
spawn a new broker via `AgentRelayClient.spawn(opts)` while the failed client instance is
effectively leaked.

6. Repeating the `'broker:start'` IPC for the same project while the stale connection file
remains (for example, reopen the project or let `reviveSession` call `this.start` again
from `src/main/broker.ts` lines 1312–1335) causes `connectExistingBroker` to allocate and
abandon additional `AgentRelayClient` instances, leading to a cumulative leak of broker
client resources over time.

Fix in Cursor | Fix in VSCode Claude

(Use Cmd/Ctrl + Click for best experience)

Prompt for AI Agent 🤖
This is a comment left during a code review.

**Path:** src/main/broker.ts
**Line:** 1293:1302
**Comment:**
	*Resource Leak: When probing multiple connection files, a client instance is created before health-checking with `getSession()`, but the failure path only logs and continues. If `connect()` allocates sockets/listeners, failed attempts are never torn down, so repeated starts can leak broker connections/resources. Explicitly disconnect/release the temporary client inside the `catch` path before trying the next candidate.

Validate the correctness of the flagged issue. If correct, How can I resolve this? If you propose a fix, implement it and please make it concise.
Once fix is implemented, also check other comments on the same PR, and ask user if the user wants to fix the rest of the comments as well. if said yes, then fetch all the comments validate the correctness and implement a minimal fix
👍 | 👎


return null
}

// Re-establish a local broker whose process has died or wedged. When this
Expand Down
Loading