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
4 changes: 3 additions & 1 deletion src/main/broker.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -203,7 +203,9 @@ vi.mock('@agent-relay/harness-driver', () => ({
}))

vi.mock('./pear-fleet-node', () => ({
startPearFleetSidecar: fleetNodeMock.startPearFleetSidecar
startPearFleetSidecar: fleetNodeMock.startPearFleetSidecar,
pearFleetProviderName: (options: { brokerName?: string; projectId?: string }) =>
`${options.brokerName || options.projectId || 'pear'}-local-fleet`
}))

vi.mock('./auth', () => ({
Expand Down
151 changes: 149 additions & 2 deletions src/main/broker.ts
Original file line number Diff line number Diff line change
Expand Up @@ -16,7 +16,7 @@ import {
type InboundDeliveryMode,
type PendingRelayMessage
} from '@agent-relay/harness-driver'
import { AgentRelay, type RelayMessage } from '@agent-relay/sdk'
import { AgentRelay, RelayPlacementError, type RelayMessage } from '@agent-relay/sdk'
import { getAccessToken, getApiUrl } from './auth'
import { assertDirectory } from './path-utils'
import { toErrorMessage } from './errors'
Expand Down Expand Up @@ -71,13 +71,26 @@ import {
resolveCommandOnPath,
resolvePackageBin
} from './mcp-command'
import { startPearFleetSidecar, type RunningPearFleetSidecar } from './pear-fleet-node'
import { startPearFleetSidecar, pearFleetProviderName, type RunningPearFleetSidecar } from './pear-fleet-node'
import {
isObserverStreamEnabled,
ObserverStreamManager,
ObserverStreamUnsupportedError
} from './observer-stream'
import { getObserverStreamCursor, setObserverStreamCursor } from './store'
import {
BrokerPlacementError,
buildPlacementMessage,
placementRequesterName,
toBrokerNodeSummary
} from './placement'
import type {
BrokerPlaceAgentInput,
BrokerPlaceAgentResult,
BrokerNodeSummary
} from '../shared/types/ipc'

export { BrokerPlacementError } from './placement'

function isShellLikeCommand(cli: string): boolean {
const normalized = basename(cli).toLowerCase()
Expand Down Expand Up @@ -1089,6 +1102,11 @@ interface BrokerSession {
leaseTimer?: ReturnType<typeof setInterval>
fleetSidecar?: RunningPearFleetSidecar
fleetSidecarCwd?: string
// Agent-scoped relay client used to invoke placement (#411). Lazily created on
// first placeAgent/listNodes: workspace key from the broker session + a
// dedicated `pear-requester-<projectId>` agent identity (placement.spawn ->
// commands.invoke requires an agent-scoped connection). Cleared in dropSession.
placementRelay?: AgentRelay
operationQueue: BrokerOperationQueue
}

Expand Down Expand Up @@ -2544,6 +2562,131 @@ export class BrokerManager {
return normalized
}

// Lazily build (and cache on the session) an agent-scoped relay client for
// placement. placement.spawn -> commands.invoke requires an agent-scoped
// connection, so a workspace-key-only client (as reconcileMessages uses) is
// insufficient: register/rotate a dedicated `pear-requester-<projectId>`
// identity and construct the client with its agent token.
private async getPlacementRelay(session: BrokerSession): Promise<AgentRelay> {
if (session.placementRelay) return session.placementRelay
Comment on lines +2570 to +2571

Copy link
Copy Markdown

Choose a reason for hiding this comment

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

P2 Badge Coalesce placement requester registration

If two listNodes/placeAgent calls arrive before the first cache fill—such as the doubled roster effect under the existing React StrictMode or calls from multiple windows—both pass this check and invoke registerOrRotate for the same identity. The second rotation invalidates the first token, so one caller can immediately issue its roster or placement request with a revoked token and fail intermittently. Cache a keyed in-flight promise and clear it safely with the session, as required by the repository's duplicate-event lifecycle guidance.

Useful? React with 👍 / 👎.


const meta = await session.client.getSession()
const workspaceKey = meta.workspace_key
if (!workspaceKey) {
throw new Error('Broker session does not expose a Relay workspace key')
}
const baseUrl = normalizeRelaycastBaseUrl(process.env.RELAYCAST_BASE_URL || process.env.RELAY_BASE_URL)
const workspaceRelay = new AgentRelay({ workspaceKey, ...(baseUrl ? { baseUrl } : {}) })
const requesterName = placementRequesterName(session.projectId)
// registerOrRotate adopts an existing `pear-requester-<projectId>` identity
// (rotating its token) so restarts don't strand orphan agents; fall back to
// register on backends without rotation support.
const register = workspaceRelay.agents.registerOrRotate ?? workspaceRelay.agents.register
const registration = await register.call(workspaceRelay.agents, {
name: requesterName,
type: 'agent',
metadata: { pearRequester: true, projectId: session.projectId }
})
const relay = new AgentRelay({
workspaceKey,
agentToken: registration.token,
...(baseUrl ? { baseUrl } : {})
})
session.placementRelay = relay
return relay
}

private localFleetNodeName(session: BrokerSession): string {
return pearFleetProviderName({
projectId: session.projectId,
cwd: session.cwd,
brokerName: session.name
})
}

/**
* Placement requester (#411). Dispatch a spawn onto an eligible fleet node via
* the relay placement engine. `input.node` omitted → any eligible least-loaded
* node; set → that exact node ('self' = this machine). Returns the node the
* agent landed on and whether this machine owns its PTY (`local`). A remote
* placement is reachable over relay chat only — its raw terminal has no relay
* transport yet (acceptance #2b, upstream-blocked).
*/
async placeAgent(projectId: string, input: BrokerPlaceAgentInput): Promise<BrokerPlaceAgentResult> {
const normalizedProjectId = projectId.trim()
if (!normalizedProjectId) throw new Error('Project id is required')
const cli = spawnCliLabel(input.cli)
if (!cli) throw new Error('cli is required for placement')

const session = this.getSessionForProject(normalizedProjectId)
const relay = await this.getPlacementRelay(session)
const capability = `spawn:${cli}`
const selfNodeName = this.localFleetNodeName(session)
const requestedNode = input.node?.trim() || undefined

// Dedupe the placed name against the shared workspace roster — a remote
// landing shares the workspace, so a name collision there is workspace-wide
// even though a single broker would accept it locally.
const existingNames = new Set(
(await relay.agents.list().catch(() => [])).map((agent) => agent.name)
)
const name = getAvailableAgentName(input.name?.trim() || `${cli}-1`, existingNames)
// Any-node placement fails fast (a clear message, never a silent hang);
// an explicit target queues up to the TTL so a briefly-offline node recovers.
const failFast = input.failFast ?? !requestedNode

try {
const ack = await relay.messaging.placement.spawn({
capability,
...(requestedNode ? { node: requestedNode } : {}),
...(requestedNode === 'self' ? { selfNodeName } : {}),
...(input.repo?.trim() ? { repo: input.repo.trim() } : {}),
input: {
name,
...(input.task?.trim() ? { task: input.task.trim() } : {}),
...(input.model?.trim() ? { model: input.model.trim() } : {})
},
failFast
})
const landedNode = ack.node.name
const nodeId = ack.node.nodeId ?? ack.node.id
return {
name,
node: landedNode,
...(nodeId ? { nodeId } : {}),
invocationId: ack.invocationId,
queued: ack.placement.queued,
local: landedNode === selfNodeName
}
} catch (err) {
if (err instanceof RelayPlacementError) {
throw new BrokerPlacementError(err.code, buildPlacementMessage(err), {
capability: err.capability,
node: err.node,
repo: err.repo
})
}
throw err
}
}

/**
* Fleet node roster for the spawn/node-picker UI (#411). Lists nodes visible
* in the project's relay workspace, optionally filtered to those advertising a
* capability (e.g. `spawn:claude`), flagging this machine's own node.
*/
async listNodes(projectId: string, capability?: string): Promise<BrokerNodeSummary[]> {
const normalizedProjectId = projectId.trim()
if (!normalizedProjectId) throw new Error('Project id is required')
const session = this.getSessionForProject(normalizedProjectId)
const relay = await this.getPlacementRelay(session)
const selfNodeName = this.localFleetNodeName(session)
const nodes = await relay.nodes.list(
capability?.trim() ? { capability: capability.trim() } : undefined
)
return nodes.map((node) => toBrokerNodeSummary(node, selfNodeName))
}

private attachClient(
sessionKey: string,
client: AgentRelayClient,
Expand Down Expand Up @@ -4342,6 +4485,10 @@ export class BrokerManager {

session.unsubEvent()
if (session.leaseTimer) clearInterval(session.leaseTimer)
// Drop the cached placement requester client; a session reconnecting to a
// different workspace would otherwise reuse an agent token scoped to the old
// workspace (mirrors the observer-token invalidation above).
session.placementRelay = undefined
void this.stopSessionFleetSidecar(session)
if (options.disconnectOnly) {
const disconnect = (session.client as { disconnect?: () => void }).disconnect
Expand Down
32 changes: 31 additions & 1 deletion src/main/ipc-handlers.ts
Original file line number Diff line number Diff line change
Expand Up @@ -21,7 +21,7 @@ import {
addProjectIntegration,
removeProjectIntegration
} from './store'
import { brokerManager, isCommandAvailableWithAugmentedPath } from './broker'
import { brokerManager, isCommandAvailableWithAugmentedPath, BrokerPlacementError } from './broker'
import * as git from './git'
import * as filesystem from './filesystem'
import * as auth from './auth'
Expand All @@ -37,6 +37,9 @@ import { findProjectForPath, projectContainsPath } from './cli'
import type {
BrokerReconcileMessagesInput,
BrokerSpawnAgentResult,
BrokerPlaceAgentInput,
BrokerPlaceAgentOutcome,
BrokerNodeSummary,
FactoryAgentStatus,
FactoryConfigReadResult,
FactoryIssueStatus,
Expand Down Expand Up @@ -762,6 +765,33 @@ export function registerIpcHandlers(): void {
return toBrokerSpawnAgentResult(result)
})

// Placement requester (#411): dispatch a spawn onto an eligible fleet node.
// Placement errors (no eligible node, capability mismatch, queue full, unmapped
// repo) are returned as structured outcomes so the UI can show a clear message
// and never hang; unexpected failures still reject.
ipcMain.handle('broker:place-agent', async (_, projectId: string, input: BrokerPlaceAgentInput): Promise<BrokerPlaceAgentOutcome> => {
try {
const result = await brokerManager.placeAgent(projectId, input)
integrationEventBridge.invalidateProjectAgentCache(projectId)
return { status: 'placed', result }
} catch (err) {
if (err instanceof BrokerPlacementError) {
return {
status: 'error',
code: err.code,
message: err.message,
...(err.node ? { node: err.node } : {}),
...(err.repo ? { repo: err.repo } : {})
}
}
throw err
}
})

ipcMain.handle('broker:list-nodes', async (_, projectId: string, capability?: string): Promise<BrokerNodeSummary[]> => {
return brokerManager.listNodes(projectId, capability)
})

ipcMain.handle('broker:list-personas', async (_, projectId: string, cwd?: string) => {
const normalizedProjectId = projectId.trim()
const personaCwd = cwd?.trim()
Expand Down
2 changes: 1 addition & 1 deletion src/main/pear-fleet-node.ts
Original file line number Diff line number Diff line change
Expand Up @@ -377,7 +377,7 @@ export function startPearFleetSidecar(options: PearFleetSidecarOptions): Running
}
}

function pearFleetProviderName(options: PearFleetNodeOptions): string {
export function pearFleetProviderName(options: PearFleetNodeOptions): string {
const rawName = `${options.brokerName || options.projectId || 'pear'}-local-fleet`
return rawName.replace(/[^\w.-]+/gu, '-').replace(/^-+|-+$/gu, '') || 'pear-local-fleet'
}
Expand Down
97 changes: 97 additions & 0 deletions src/main/placement-helpers.test.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,97 @@
import { describe, it, expect } from 'vitest'
import { RelayPlacementError, type RelayNode } from '@agent-relay/sdk'
import {
buildPlacementMessage,
placementRequesterName,
toBrokerNodeSummary
} from './placement'

function placementError(code: RelayPlacementError['code'], ctx: { capability?: string; node?: string; repo?: string } = {}): RelayPlacementError {
return new RelayPlacementError(code, `raw ${code}`, {
capability: ctx.capability ?? 'spawn:claude',
node: ctx.node,
repo: ctx.repo,
attempts: 1
})
}

describe('buildPlacementMessage', () => {
it('names the offending node and cli for capability_mismatch', () => {
const message = buildPlacementMessage(placementError('capability_mismatch', { node: 'gpu-box-1' }))
expect(message).toBe('Node "gpu-box-1" can\'t run claude.')
})

it('falls back to a generic capability message when no node is named', () => {
const message = buildPlacementMessage(placementError('capability_mismatch'))
expect(message).toBe('No node can run claude.')
})

it('reports queue saturation for placement_queue_full', () => {
expect(buildPlacementMessage(placementError('placement_queue_full'))).toMatch(/try again/i)
})

it('reports no eligible node for placement_ttl_expired (never a hang)', () => {
expect(buildPlacementMessage(placementError('placement_ttl_expired'))).toBe(
'No node advertises claude right now.'
)
})

it('names the repo for unmapped_repo', () => {
expect(buildPlacementMessage(placementError('unmapped_repo', { repo: 'pear' }))).toBe(
'No node has repo "pear" checked out.'
)
})

it('strips the spawn: prefix from the capability in messages', () => {
const message = buildPlacementMessage(placementError('placement_ttl_expired', { capability: 'spawn:codex' }))
expect(message).toContain('codex')
expect(message).not.toContain('spawn:')
})
})

describe('placementRequesterName', () => {
it('derives a sanitized, workspace-safe requester identity per project', () => {
expect(placementRequesterName('project-1')).toBe('pear-requester-project-1')
})

it('replaces path/space characters that a workspace name cannot carry', () => {
expect(placementRequesterName('a/b c:d')).toBe('pear-requester-a-b-c-d')
})

it('never returns an empty name', () => {
expect(placementRequesterName('')).toBe('pear-requester')
})
})

describe('toBrokerNodeSummary', () => {
const baseNode: RelayNode = {
name: 'other-node',
status: 'online',
live: true,
load: 2,
activeAgents: 1,
maxAgents: 4,
capabilities: [{ name: 'spawn:claude' }, { name: 'spawn:codex' }],
repoKeys: ['pear'],
tags: ['pear', 'local']
} as RelayNode

it('flattens capabilities and preserves liveness/load for the picker', () => {
const summary = toBrokerNodeSummary(baseNode, 'my-self-node')
expect(summary.capabilities).toEqual(['spawn:claude', 'spawn:codex'])
expect(summary.live).toBe(true)
expect(summary.load).toBe(2)
expect(summary.activeAgents).toBe(1)
expect(summary.isSelf).toBe(false)
})

it('flags this machine when the node name matches the local fleet node', () => {
const summary = toBrokerNodeSummary({ ...baseNode, name: 'my-self-node' }, 'my-self-node')
expect(summary.isSelf).toBe(true)
})

it('treats an absent live flag as offline', () => {
const summary = toBrokerNodeSummary({ ...baseNode, live: undefined }, 'my-self-node')
expect(summary.live).toBe(false)
})
})
Loading
Loading