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
20 changes: 17 additions & 3 deletions src/cli/fleet.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -127,8 +127,23 @@ const issueFile = {
* flag would put the timer in scope for tests whose whole point is that no
* timer decides the ordering (#346 review, cubic).
*/
class ControlledCompletingRemoteFleetClient extends FakeFleetClient {
class CompletingRemoteFleetBase extends FakeFleetClient {
override readonly placementLocality = 'remote' as const

override async roster() {
const roster = await super.roster()
return {
agents: roster.agents.map((agent) => ({ ...agent, node: 'sf-mini' })),
nodes: [{
name: 'sf-mini',
capabilities: ['spawn:codex' as const, 'spawn:claude' as const, 'workflow:run' as const],
live: true,
}],
}
}
}

class ControlledCompletingRemoteFleetClient extends CompletingRemoteFleetBase {
implementerName?: string
exitEmitted = false

Expand All @@ -139,8 +154,7 @@ class ControlledCompletingRemoteFleetClient extends FakeFleetClient {
}
}

class CompletingRemoteFleetClient extends FakeFleetClient {
override readonly placementLocality = 'remote' as const
class CompletingRemoteFleetClient extends CompletingRemoteFleetBase {
readonly lifecycleOrder: string[] = []

override async spawn(input: SpawnInput): Promise<SpawnResult> {
Expand Down
30 changes: 30 additions & 0 deletions src/fleet/relay-fleet-client.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -815,6 +815,36 @@ describe('RelayFleetClient', () => {
expect(messaging.agentPresenceCalls).toBe(1)
})

it('confirms registration only when presence and a live capable host agree', async () => {
const messaging = new FakeMessaging()
messaging.agentRows = [
{ name: 'ar-1-impl', status: 'online', node: 'alpha' },
{ name: 'ar-offline-host-impl', status: 'online', node: 'beta' },
{ name: 'ar-hostless-impl', status: 'online' },
]
messaging.nodeRows = [
{ name: 'alpha', status: 'online', capabilities: [{ name: 'spawn:codex' }] },
{ name: 'beta', status: 'offline', capabilities: [{ name: 'spawn:codex' }] },
]
const fleet = createClient(messaging)

await expect(fleet.isAgentRegistered({
name: 'ar-1-impl',
node: 'alpha',
capability: 'spawn:codex',
})).resolves.toBe(true)
await expect(fleet.isAgentRegistered({
name: 'ar-offline-host-impl',
node: 'beta',
Comment thread
miyaontherelay marked this conversation as resolved.
capability: 'spawn:codex',
})).resolves.toBe(false)
await expect(fleet.isAgentRegistered({
name: 'ar-hostless-impl',
node: 'alpha',
capability: 'spawn:codex',
})).resolves.toBe(false)
})

it('sends DMs and channel messages through the agent-scoped surface', async () => {
const messaging = new FakeMessaging()
const fleet = createClient(messaging)
Expand Down
13 changes: 13 additions & 0 deletions src/fleet/relay-fleet-client.ts
Original file line number Diff line number Diff line change
Expand Up @@ -574,6 +574,19 @@ export class RelayFleetClient implements FleetClient {
}
}

async isAgentRegistered(input: {
name: string
node: string
capability: Capability
}): Promise<boolean> {
const roster = await this.roster()
const agent = roster.agents.find((candidate) =>
candidate.name === input.name && candidate.node === input.node)
if (!agent) return false
return roster.nodes.some((node) =>
node.name === input.node && node.live && node.capabilities.includes(input.capability))
}

async discoverTeammates(query: TeammateQuery): Promise<TeammateAgent[]> {
this.#teammateDirectory ??= this.#createTeammateDirectory()
return await this.#teammateDirectory.discover(query)
Expand Down
184 changes: 183 additions & 1 deletion src/orchestrator/factory.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -1073,7 +1073,18 @@ class RemoteLifecycleFleetClient extends FakeFleetClient {

override async roster() {
const roster = await super.roster()
return { ...roster, agents: roster.agents.map((agent) => ({ ...agent, node: 'sf-mini' })) }
return {
agents: roster.agents.map((agent) => ({ ...agent, node: 'sf-mini' })),
nodes: [{
name: 'sf-mini',
capabilities: ['spawn:codex' as const, 'spawn:claude' as const, 'workflow:run' as const],
live: true,
}],
}
}

async isAgentRegistered(): Promise<boolean> {
return true
}

override async reconcileTrackedAgents(): Promise<void> {
Expand All @@ -1084,6 +1095,49 @@ class RemoteLifecycleFleetClient extends FakeFleetClient {
}
}

class NodeDiscoveryRemoteFleetClient extends RemoteLifecycleFleetClient {
nodes: RosterEntry['nodes'] = [
{ name: 'sf-mini', capabilities: ['spawn:codex', 'spawn:claude'], live: true },
]
readonly unregistered = new Set<string>()

override async spawn(input: SpawnInput): Promise<SpawnResult> {
const result = await super.spawn(input)
return {
...result,
node: input.node ?? result.node,
locality: 'remote',
}
}

override async roster(): Promise<RosterEntry> {
const roster = await super.roster()
return {
agents: roster.agents
.filter((agent) => !this.unregistered.has(agent.name))
.map((agent) => ({
...agent,
node: this.spawns.findLast((spawn) => spawn.name === agent.name)?.node === 'self'
? agent.node
: this.spawns.findLast((spawn) => spawn.name === agent.name)?.node ?? agent.node,
})),
nodes: structuredClone(this.nodes),
}
}

override async isAgentRegistered(input: {
name: string
node: string
capability: 'spawn:codex' | 'spawn:claude' | 'workflow:run'
}): Promise<boolean> {
if (this.unregistered.has(input.name)) return false
const roster = await this.roster()
return roster.agents.some((agent) => agent.name === input.name && agent.node === input.node) &&
roster.nodes.some((node) =>
node.name === input.node && node.live && node.capabilities.includes(input.capability))
}
}

/**
* Fires a hook from inside `#finishDurableRelease`'s agent release: after
* completion recorded the terminal role at the provider write, and before the
Expand Down Expand Up @@ -3115,6 +3169,123 @@ describe('waitForDispatchTerminal', () => {
})
})

describe('remote fleet node discovery and registration admission', () => {
it('skips an offline capable node and pins every spawn to a live capable node', async () => {
const path = issuePath(390)
const issue = issueFile(390)
const fleet = new NodeDiscoveryRemoteFleetClient()
fleet.nodes = [
{ name: 'chief-broker', capabilities: ['spawn:codex', 'spawn:claude'], live: false },
{ name: 'sf-mini', capabilities: ['spawn:codex', 'spawn:claude'], live: true },
]
const factory = createFactory(config(), {
mount: new FakeMountClient({ [path]: issue }),
fleet,
stateStore: new InMemoryStateStore({ batchSize: 2 }),
triage: new StaticTriage(),
})

try {
const decision = await factory.triageIssue(parseLinearIssue(path, issue))
await expect(factory.dispatch(decision)).resolves.toMatchObject({
agents: [
{ name: 'ar-390-impl-pear', role: 'implementer' },
{ name: 'ar-390-review', role: 'reviewer' },
],
})

expect(fleet.spawns.map(({ name, node }) => ({ name, node }))).toEqual([
{ name: 'ar-390-impl-pear', node: 'sf-mini' },
{ name: 'ar-390-review', node: 'sf-mini' },
])
} finally {
await factory.stop()
}
})

it('refuses when no live capable node exists before writing a durable lifecycle claim', async () => {
const path = issuePath(391)
const issue = issueFile(391)
const fleet = new NodeDiscoveryRemoteFleetClient()
fleet.nodes = [
{ name: 'chief-broker', capabilities: ['spawn:codex', 'spawn:claude'], live: false },
{ name: 'sf-mini', capabilities: ['spawn:codex'], live: false },
]
const stateStore = new InMemoryStateStore({ batchSize: 1 })
const factory = createFactory(config({ batchSize: 1 }), {
mount: new FakeMountClient({ [path]: issue }),
fleet,
stateStore,
triage: new StaticTriage(),
})

try {
const decision = await factory.triageIssue(parseLinearIssue(path, issue))
await expect(factory.dispatch(decision)).rejects.toThrow(
'no live fleet node advertises spawn:codex',
)

expect(fleet.spawns).toEqual([])
await expect(stateStore.listDispatchLifecycles('factory-test')).resolves.toEqual([])
} finally {
await factory.stop()
}
})

it('rolls back an unregistered remote spawn and frees the slot for the next issue', async () => {
const firstPath = issuePath(392)
const secondPath = issuePath(393)
const firstIssue = issueFile(392)
const secondIssue = issueFile(393)
const fleet = new NodeDiscoveryRemoteFleetClient()
fleet.unregistered.add('ar-392-review')
const stateStore = new InMemoryStateStore({ batchSize: 1 })
const clock = new ManualClock()
const factory = createFactory(config({ batchSize: 1 }), {
mount: new FakeMountClient({ [firstPath]: firstIssue, [secondPath]: secondIssue }),
fleet,
stateStore,
triage: new StaticTriage(),
clock,
})

try {
const first = await factory.triageIssue(parseLinearIssue(firstPath, firstIssue))
await expect(factory.dispatch(first)).rejects.toThrow(
'did not register with the fleet before the startup deadline',
)
expect(clock.value).toBe(30_000)

expect(fleet.releases).toContainEqual({
name: 'ar-392-review',
reason: 'spawn-registration-timeout',
})
expect(fleet.releases).toContainEqual({
name: 'ar-392-impl-pear',
reason: 'spawn-registration-timeout',
})
await expect(stateStore.getDispatchLifecycle(
'factory-test',
dispatchIssueIdentity(first.issue),
)).resolves.toBeUndefined()

const second = await factory.triageIssue(parseLinearIssue(secondPath, secondIssue))
await expect(factory.dispatch(second)).resolves.toMatchObject({
agents: [
{ name: 'ar-393-impl-pear', role: 'implementer' },
{ name: 'ar-393-review', role: 'reviewer' },
],
})
await expect(stateStore.getDispatchLifecycle(
'factory-test',
dispatchIssueIdentity(second.issue),
)).resolves.toMatchObject({ phase: 'running' })
} finally {
await factory.stop()
}
})
})

describe('FactoryLoop', () => {
it('sweeps preview orphans on daemon startup using durable active issue owners', async () => {
const mount = new FakeMountClient()
Expand Down Expand Up @@ -26966,6 +27137,17 @@ describe('FactoryLoop PR babysitter', () => {

it('pins a remote babysitter to the preview node and gives it the live URL', async () => {
class RemotePreviewFleetClient extends RemoteLifecycleFleetClient {
override async roster(): Promise<RosterEntry> {
const roster = await super.roster()
return {
...roster,
nodes: [
...roster.nodes,
{ name: 'preview-node', capabilities: ['spawn:claude'], live: true },
],
}
}

override async createPreview(input: PreviewStartInput): Promise<PreviewReference> {
return { ...await super.createPreview(input), node: 'preview-node' }
}
Expand Down
Loading