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
18 changes: 14 additions & 4 deletions electron-builder.mcp-resources.yml
Original file line number Diff line number Diff line change
Expand Up @@ -22,9 +22,16 @@ extraResources:
- '@agent-relay/cloud/**'
- '@agent-relay/config/**'
- '@agent-relay/fleet/**'
- '@agent-relay/fleet/node_modules/@relaycast/sdk/**'
- '@agent-relay/fleet/node_modules/@relaycast/sdk/node_modules/zod/**'
- '@agent-relay/fleet/node_modules/@relaycast/types/**'
- '@agent-relay/fleet/node_modules/@relaycast/types/node_modules/zod/**'
- '@agent-relay/harness-driver/**'
- '@agent-relay/harnesses/**'
- '@agent-relay/sdk/**'
- '@agent-relay/sdk/node_modules/@relaycast/sdk/**'
- '@agent-relay/sdk/node_modules/@relaycast/types/**'
- '@agent-relay/sdk/node_modules/zod/**'
- '@agent-relay/utils/**'
- '@aws-crypto/sha1-browser/**'
- '@aws-crypto/sha1-browser/node_modules/@smithy/util-utf8/**'
Expand Down Expand Up @@ -74,10 +81,6 @@ extraResources:
- '@modelcontextprotocol/sdk/**'
- '@posthog/core/**'
- '@posthog/types/**'
- '@relaycast/sdk/**'
- '@relaycast/sdk/node_modules/zod/**'
- '@relaycast/types/**'
- '@relaycast/types/node_modules/zod/**'
- '@relayfile/client/**'
- '@relayflows/browser-primitive/**'
- '@relayflows/browser-primitive/node_modules/@agent-relay/sdk/**'
Expand Down Expand Up @@ -158,6 +161,13 @@ extraResources:
- '@xterm/headless/**'
- 'accepts/**'
- 'agent-relay/**'
- 'ai-hist-native/**'
- 'ai-hist-native-darwin-arm64/**'
- 'ai-hist-native-darwin-x64/**'
- 'ai-hist-native-linux-arm64-gnu/**'
- 'ai-hist-native-linux-arm64-musl/**'
- 'ai-hist-native-linux-x64-gnu/**'
- 'ai-hist-native-linux-x64-musl/**'
- 'ajv/**'
- 'ajv-formats/**'
- 'ansi-escapes/**'
Expand Down
583 changes: 397 additions & 186 deletions package-lock.json

Large diffs are not rendered by default.

18 changes: 9 additions & 9 deletions package.json
Original file line number Diff line number Diff line change
Expand Up @@ -35,13 +35,13 @@
"lint": "eslint ."
},
"dependencies": {
"@agent-relay/cloud": "^9.2.1",
"@agent-relay/factory": "^0.1.13",
"@agent-relay/fleet": "^9.2.1",
"@agent-relay/harness-driver": "^9.2.1",
"@agent-relay/harnesses": "^9.2.1",
"@agent-relay/integration-prompts": "^9.2.1",
"@agent-relay/sdk": "^9.2.1",
"@agent-relay/cloud": "^10.6.3",
"@agent-relay/factory": "^0.1.19",
"@agent-relay/fleet": "^10.6.3",
"@agent-relay/harness-driver": "^10.6.3",
"@agent-relay/harnesses": "^10.6.3",
"@agent-relay/integration-prompts": "^10.6.3",
"@agent-relay/sdk": "^10.6.3",
"@agentworkforce/deploy": "^4.1.16",
"@relayburn/sdk": "^4.0.0",
"@relaycast/sdk": "^5.0.7",
Expand All @@ -51,7 +51,7 @@
"@xterm/addon-web-links": "^0.11.0",
"@xterm/addon-webgl": "^0.18.0",
"@xterm/xterm": "^5.5.0",
"agent-relay": "^9.2.1",
"agent-relay": "^10.6.3",
"agentworkforce": "^4.1.16",
"ai-hist": "^0.2.3",
"allotment": "^1.0.9",
Expand All @@ -71,7 +71,7 @@
"protobufjs": "8.5.0"
},
"devDependencies": {
"@agent-relay/evals": "^9.2.1",
"@agent-relay/evals": "^10.6.3",
"@playwright/test": "^1.57.0",
"@tailwindcss/vite": "^4.0.0",
"@testing-library/react": "^16.3.2",
Expand Down
16 changes: 11 additions & 5 deletions src/main/broker.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -739,9 +739,7 @@ describe('BrokerManager local + cloud coexistence', () => {
projectId: PROJECT_ID,
cwd: '/tmp/project-1',
brokerName: 'pear-project-1',
connection: {
url: 'http://127.0.0.1:4242'
}
readBrokerSession: expect.any(Function)
}))

await manager.shutdown()
Expand Down Expand Up @@ -2021,14 +2019,22 @@ exit 2
kind: 'worker_stream',
name: 'claude-1',
chunk: 'pong\n',
seq: 22
seq: 22,
offset: 101
}
listener?.(chunkEvent)
listener?.(chunkEvent)

const ptyCalls = (win.webContents.send as ReturnType<typeof vi.fn>).mock.calls
.filter(([channel]) => channel === 'broker:pty-chunk')
expect(ptyCalls).toEqual([['broker:pty-chunk', PROJECT_ID, 'claude-1', 'pong\n']])
expect(ptyCalls).toEqual([[
'broker:pty-chunk',
PROJECT_ID,
'claude-1',
'pong\n',
101,
expect.any(Number)
]])

await manager.shutdown()
})
Expand Down
31 changes: 18 additions & 13 deletions src/main/broker.ts
Original file line number Diff line number Diff line change
Expand Up @@ -363,6 +363,8 @@ export interface AttachTerminalResult {
cols: number
cursor: [number, number]
screen: string
offset?: number
generation?: number
}
}

Expand Down Expand Up @@ -1840,20 +1842,11 @@ export class BrokerManager {
await this.stopSessionFleetSidecar(session)
}

const url = getClientBaseUrl(session.client)
if (!url) {
console.warn(`[broker] Local fleet node skipped for project ${session.projectId}: broker URL unavailable`)
return
}

const sidecar = startPearFleetSidecar({
projectId: session.projectId,
cwd: session.cwd,
brokerName: session.name,
connection: {
url,
...(getClientApiKey(session.client) ? { apiKey: getClientApiKey(session.client) } : {})
},
readBrokerSession: () => session.client.getSession(),
log: (message) => console.log(`[broker] ${message}`),
warn: (message) => console.warn(`[broker] ${message}`)
})
Expand Down Expand Up @@ -2522,7 +2515,14 @@ export class BrokerManager {
}
const targetWindow = this.windowForSession(sessionKey, win)
if (targetWindow && !targetWindow.isDestroyed()) {
targetWindow.webContents.send('broker:pty-chunk', projectId, event.name, event.chunk)
const offset = brokerEventNumber(event, 'offset')
targetWindow.webContents.send(
'broker:pty-chunk',
projectId,
event.name,
event.chunk,
...(offset !== undefined ? [offset, eventStreamGeneration] : [])
)
}
this.rememberAgentSession(event.name, sessionKey)
if (this.sessions.get(sessionKey)?.cloudSandboxId) {
Expand Down Expand Up @@ -3445,7 +3445,11 @@ export class BrokerManager {
rows: snapshot.rows,
cols: snapshot.cols,
cursor: snapshot.cursor,
screen: Buffer.from(snapshot.screen, 'base64').toString('utf-8')
screen: Buffer.from(snapshot.screen, 'base64').toString('utf-8'),
...(typeof snapshot.offset === 'number' && Number.isFinite(snapshot.offset)
? { offset: snapshot.offset }
: {}),
generation: session.eventStreamGeneration
}
}
} catch (err) {
Expand Down Expand Up @@ -3682,7 +3686,8 @@ export class BrokerManager {
screen:
format === 'ansi'
? Buffer.from(snapshot.screen, 'base64').toString('utf-8')
: snapshot.screen
: snapshot.screen,
...(snapshot.offset !== undefined ? { offset: snapshot.offset } : {})
}
} catch (err) {
if (isMissingAgentError(err)) return null
Expand Down
124 changes: 95 additions & 29 deletions src/main/pear-fleet-node.test.ts
Original file line number Diff line number Diff line change
@@ -1,6 +1,11 @@
import { describe, expect, it, vi } from 'vitest'
import { invokeNodeHandler, nodeInfo, nodeManifest } from '@agent-relay/fleet'
import { createPearFleetNodeDefinition, PEAR_LOCAL_SPAWN_HARNESSES, reconnectDelayMs } from './pear-fleet-node'
import { invokeNodeHandler, nodeInfo } from '@agent-relay/fleet'
import {
createPearFleetNodeDefinition,
PEAR_LOCAL_SPAWN_HARNESSES,
resolvePearFleetConnection,
startPearFleetSidecar
} from './pear-fleet-node'

const expectedCapabilities = Object.keys(PEAR_LOCAL_SPAWN_HARNESSES).map((cli) => `spawn:${cli}`)

Expand All @@ -12,16 +17,19 @@ describe('Pear local fleet node', () => {
brokerName: 'pear-project-1'
})

const manifest = nodeManifest(definition)
expect(manifest.name).toBe('pear-project-1-local-fleet')
expect(manifest.capabilities.map((capability) => capability.name)).toEqual(expectedCapabilities)
expect(manifest.capabilities).toEqual(expectedCapabilities.map((name) => expect.objectContaining({
name,
metadata: expect.objectContaining({
pearLocalNode: true,
clonePaths: { 'project-1': '/tmp/project-1' }
})
})))
const info = nodeInfo(definition)
expect(info.name).toBe('pear-project-1-local-fleet')
expect(info.capabilities).toEqual(expectedCapabilities)
expect(expectedCapabilities.map((name) => definition.capabilities[name])).toEqual(
expectedCapabilities.map(() =>
expect.objectContaining({
metadata: expect.objectContaining({
pearLocalNode: true,
clonePaths: { 'project-1': '/tmp/project-1' }
})
})
)
)
})

it('spawns non-Claude/Codex harnesses through the broker', async () => {
Expand Down Expand Up @@ -110,26 +118,84 @@ describe('Pear local fleet node', () => {
})
})

describe('reconnectDelayMs', () => {
it('caps transient reconnects (already registered) at the fast ceiling', () => {
expect(reconnectDelayMs(1, true)).toBe(500)
expect(reconnectDelayMs(2, true)).toBe(1_000)
expect(reconnectDelayMs(4, true)).toBe(4_000)
// 500 * 2**4 = 8000 -> clamped to the 5s ceiling
expect(reconnectDelayMs(5, true)).toBe(5_000)
expect(reconnectDelayMs(50, true)).toBe(5_000)
describe('resolvePearFleetConnection', () => {
it('waits for the broker-minted node token and attaches to the same v10 node', async () => {
let now = 0
let reads = 0
const connection = await resolvePearFleetConnection(async () => {
reads += 1
if (reads === 1) {
return { node_id: 'node-1', node_name: 'pear-project-1' }
}
return {
node_id: 'node-1',
node_name: 'pear-project-1',
node_token: 'nt_live_test',
relay_base_url: 'https://cast.example'
}
}, new AbortController().signal, {
timeoutMs: 1_000,
pollIntervalMs: 250,
now: () => now,
sleep: async (ms) => {
now += ms
}
})

expect(reads).toBe(2)
expect(connection).toEqual({
connection: {
nodeId: 'node-1',
nodeToken: 'nt_live_test',
baseUrl: 'https://cast.example'
},
nodeName: 'pear-project-1'
})
})

it('fails clearly instead of silently disabling the provider when the token never arrives', async () => {
let now = 0
await expect(resolvePearFleetConnection(
async () => ({ node_id: 'node-1', node_name: 'pear-project-1' }),
new AbortController().signal,
{
timeoutMs: 500,
pollIntervalMs: 250,
now: () => now,
sleep: async (ms) => {
now += ms
}
}
)).rejects.toThrow('timed out waiting for a node token for node-1')
})

it('backs off much harder while the node has never registered', () => {
// Grows past the registered ceiling so a wedged broker is not stormed.
expect(reconnectDelayMs(5, false)).toBe(8_000)
expect(reconnectDelayMs(6, false)).toBe(16_000)
expect(reconnectDelayMs(8, false)).toBe(60_000)
expect(reconnectDelayMs(50, false)).toBe(60_000)
it('bounds a broker session read that never settles', async () => {
let now = 0
await expect(resolvePearFleetConnection(
() => new Promise(() => {}),
new AbortController().signal,
{
timeoutMs: 500,
now: () => now,
sleep: async (ms) => {
now += ms
}
}
)).rejects.toThrow('timed out waiting for the broker node id')
expect(now).toBe(500)
})

it('never returns a negative or sub-base delay for the first attempt', () => {
expect(reconnectDelayMs(0, false)).toBe(500)
expect(reconnectDelayMs(1, false)).toBe(500)
it('stops promptly while a broker session read is hung', async () => {
const sidecar = startPearFleetSidecar({
projectId: 'project-1',
cwd: '/tmp/project-1',
brokerName: 'pear-project-1',
readBrokerSession: () => new Promise(() => {})
})
const registered = sidecar.registered.catch((error: unknown) => error)

await sidecar.stop()

await expect(registered).resolves.toMatchObject({ name: 'AbortError' })
})
})
Loading
Loading