From 540ee52d80fb937cffb8aba7f9fa6ac21baa4a0f Mon Sep 17 00:00:00 2001 From: Khaliq Date: Mon, 17 Aug 2026 11:21:33 +0200 Subject: [PATCH 1/7] fix(factory): fail closed on local relay placement --- .agentworkforce/features/manifest.yaml | 4 +-- src/fleet/relay-fleet-client.test.ts | 45 ++++++++++++++++++++++++++ src/fleet/relay-fleet-client.ts | 11 +++++++ 3 files changed, 58 insertions(+), 2 deletions(-) diff --git a/.agentworkforce/features/manifest.yaml b/.agentworkforce/features/manifest.yaml index d659e62..1ce0249 100644 --- a/.agentworkforce/features/manifest.yaml +++ b/.agentworkforce/features/manifest.yaml @@ -279,7 +279,7 @@ categories: - id: fleet-target-node name: Target Fleet Node cli: factory fleet spawn --node - description: Constrain a spawn or resume request to the local node or a named hosted node + description: Constrain internal placement to the local node or Relay placement to a named hosted node; Relay treats self as no preference and fails closed unless the result proves a named remote node location: src/cli/fleet.ts, src/fleet/relay-fleet-client.ts verify_tier: 4 @@ -1238,7 +1238,7 @@ categories: - id: fleet-relay-backend name: Hosted Relay Fleet api: RelayFleetClient - description: Invoke hosted spawn and release actions, place by capability or node, and expose a unified roster + description: Invoke hosted spawn and release actions, place by capability or node, reject results that do not prove a named remote node, and expose a unified roster location: src/fleet/relay-fleet-client.ts verify_tier: 4 diff --git a/src/fleet/relay-fleet-client.test.ts b/src/fleet/relay-fleet-client.test.ts index d14ac7c..9097af4 100644 --- a/src/fleet/relay-fleet-client.test.ts +++ b/src/fleet/relay-fleet-client.test.ts @@ -1,6 +1,7 @@ import { describe, expect, it, vi } from 'vitest' import { RelayFleetClient, type RelayClientFactoryOptions, type RelayClientLike } from './relay-fleet-client' +import { runFleetCli } from '../cli/fleet' import type { RelayActionInvocation, @@ -306,6 +307,50 @@ describe('RelayFleetClient', () => { }) }) + it.each([ + ['self', 'self'], + ['a missing node', ''], + ])('fails closed when placement resolves to %s', async (_label, node) => { + const messaging = new FakeMessaging() + messaging.placementAck = { placement: { node } } + const fleet = createClient(messaging) + + await expect(fleet.spawn({ + name: 'ar-1-impl', + capability: 'spawn:codex', + node: 'self', + repo: 'AgentWorkforce/factory', + task: 'do work', + })).rejects.toThrow('Relay placement did not prove a named remote node') + + expect(fleet.trackedAgents().size).toBe(0) + }) + + it('surfaces a self-placement refusal as a non-zero CLI result', async () => { + const messaging = new FakeMessaging() + messaging.placementAck = { placement: { node: 'self' } } + const fleet = createClient(messaging) + const stderr: string[] = [] + + const code = await runFleetCli([ + 'fleet', + 'spawn', + 'spawn:codex', + '--name', + 'ar-1-impl', + ], { + fleet, + stdout: { write: () => true } as never, + stderr: { write: (chunk: string | Uint8Array) => { + stderr.push(String(chunk)) + return true + } } as never, + }) + + expect(code).toBe(1) + expect(stderr.join('')).toContain('Relay placement did not prove a named remote node') + }) + it('sweeps previews on every live preview-capable node', async () => { const messaging = new FakeMessaging() messaging.nodeRows = [{ diff --git a/src/fleet/relay-fleet-client.ts b/src/fleet/relay-fleet-client.ts index 38ff4c5..e481043 100644 --- a/src/fleet/relay-fleet-client.ts +++ b/src/fleet/relay-fleet-client.ts @@ -187,6 +187,7 @@ export class RelayFleetClient implements FleetClient { }) const invocation = await this.#awaitInvocation(ack.actionName || 'spawn', ack) const result = spawnResultFromInvocation(input.name, input.sessionRef, invocation, ack) + assertNamedRemotePlacement(result) this.#track(result.name, ack) return result } @@ -807,6 +808,16 @@ function spawnResultFromInvocation( } } +function assertNamedRemotePlacement(result: SpawnResult): void { + const node = result.node?.trim() + if (!node || node === 'self') { + throw new Error( + `Relay placement did not prove a named remote node for ${result.name}; ` + + `refusing to accept the spawn result`, + ) + } +} + function lifecycleSignalFromInvocation( input: Record | undefined, callerName: string | null | undefined, From 39f8138ee0add24e259a8fbdc51fffacb91be319 Mon Sep 17 00:00:00 2001 From: Khaliq Date: Mon, 17 Aug 2026 11:58:08 +0200 Subject: [PATCH 2/7] fix(factory): clean up rejected relay placements --- src/fleet/relay-fleet-client.test.ts | 41 ++++++++++++++++++++++++++++ src/fleet/relay-fleet-client.ts | 24 ++++++++++++++-- 2 files changed, 62 insertions(+), 3 deletions(-) diff --git a/src/fleet/relay-fleet-client.test.ts b/src/fleet/relay-fleet-client.test.ts index 9097af4..03b3d08 100644 --- a/src/fleet/relay-fleet-client.test.ts +++ b/src/fleet/relay-fleet-client.test.ts @@ -324,6 +324,47 @@ describe('RelayFleetClient', () => { })).rejects.toThrow('Relay placement did not prove a named remote node') expect(fleet.trackedAgents().size).toBe(0) + expect(messaging.invokes).toContainEqual({ + name: 'release', + input: { + name: 'ar-1-impl', + agent: 'ar-1-impl', + reason: 'unverified-placement', + }, + }) + }) + + it('rejects the acknowledgement node even when action output synthesizes a named node', async () => { + const messaging = new FakeMessaging() + messaging.placementAck = { + invocationId: 'self-placement', + status: 'pending', + placement: { node: 'self' }, + } + messaging.invocations.set('self-placement', [{ + invocationId: 'self-placement', + actionName: 'spawn', + status: 'completed', + output: { name: 'ar-1-impl', node: 'mac-mini' }, + }]) + const fleet = createClient(messaging) + + await expect(fleet.spawn({ + name: 'ar-1-impl', + capability: 'spawn:codex', + node: 'self', + repo: 'AgentWorkforce/factory', + task: 'do work', + })).rejects.toThrow('Relay placement did not prove a named remote node') + + expect(messaging.invokes).toContainEqual({ + name: 'release', + input: expect.objectContaining({ + name: 'ar-1-impl', + reason: 'unverified-placement', + }), + }) + expect(fleet.trackedAgents().size).toBe(0) }) it('surfaces a self-placement refusal as a non-zero CLI result', async () => { diff --git a/src/fleet/relay-fleet-client.ts b/src/fleet/relay-fleet-client.ts index e481043..3bb4cea 100644 --- a/src/fleet/relay-fleet-client.ts +++ b/src/fleet/relay-fleet-client.ts @@ -187,7 +187,22 @@ export class RelayFleetClient implements FleetClient { }) const invocation = await this.#awaitInvocation(ack.actionName || 'spawn', ack) const result = spawnResultFromInvocation(input.name, input.sessionRef, invocation, ack) - assertNamedRemotePlacement(result) + try { + assertNamedRemotePlacement(result, ack.placement?.node) + } catch (error) { + // A completed placement invocation has already launched the worker. If + // Relay cannot prove that it ran on a named remote node, tear that worker + // down before refusing the result so a rejected spawn cannot keep acting + // outside Factory's lifecycle tracking. + try { + await this.release(result.name, 'unverified-placement') + } catch (releaseError) { + this.#log( + `Failed to release ${result.name} after unverified Relay placement: ${errorMessage(releaseError)}`, + ) + } + throw error + } this.#track(result.name, ack) return result } @@ -808,8 +823,11 @@ function spawnResultFromInvocation( } } -function assertNamedRemotePlacement(result: SpawnResult): void { - const node = result.node?.trim() +function assertNamedRemotePlacement(result: SpawnResult, acknowledgedNode: string | undefined): void { + // The placement acknowledgement is authoritative. Action output may name + // the node executing the spawn handler even when Relay acknowledged `self`, + // so accepting the synthesized SpawnResult would let self-placement pass. + const node = acknowledgedNode?.trim() if (!node || node === 'self') { throw new Error( `Relay placement did not prove a named remote node for ${result.name}; ` + From 9cebfe04d4f83ec200b6b81ee27285e7f6ad4ea1 Mon Sep 17 00:00:00 2001 From: Khaliq Date: Mon, 17 Aug 2026 12:04:54 +0200 Subject: [PATCH 3/7] test(factory): trust normalized placement acknowledgements --- src/fleet/relay-fleet-client.test.ts | 43 ++++++++++++++++++++++++---- src/fleet/relay-fleet-client.ts | 17 ++++++----- 2 files changed, 47 insertions(+), 13 deletions(-) diff --git a/src/fleet/relay-fleet-client.test.ts b/src/fleet/relay-fleet-client.test.ts index 03b3d08..262fc77 100644 --- a/src/fleet/relay-fleet-client.test.ts +++ b/src/fleet/relay-fleet-client.test.ts @@ -102,15 +102,18 @@ class FakeMessaging { spawn: async (input: RelaySpawnPlacementInput) => { this.placements.push(input) const invocationId = this.placementAck.invocationId ?? `inv-${++this.nextInvocationId}` + const acknowledgedNode = Object.prototype.hasOwnProperty.call(this.placementAck, 'placement') + ? this.placementAck.placement?.node + : 'node-a' return { invocationId, actionName: 'spawn', status: this.placementAck.status ?? (this.invocations.has(invocationId) ? 'pending' : 'completed'), dispatchedNodeId: this.placementAck.dispatchedNodeId, - node: { name: this.placementAck.placement?.node ?? 'node-a' } as RelayNode, + node: { name: acknowledgedNode } as RelayNode, placement: { capability: input.capability, - node: this.placementAck.placement?.node ?? 'node-a', + node: acknowledgedNode, attempts: 1, queued: false, }, @@ -210,6 +213,7 @@ describe('RelayFleetClient', () => { exit_after_task: true, }) expect(fleet.trackedAgents().get('ar-1-impl')).toMatchObject({ invocationId: 'inv-1', node: 'mac-mini' }) + expect(messaging.invokes).not.toContainEqual(expect.objectContaining({ name: 'release' })) }) it('passes explicit node targets through to placement', async () => { @@ -308,11 +312,12 @@ describe('RelayFleetClient', () => { }) it.each([ - ['self', 'self'], - ['a missing node', ''], - ])('fails closed when placement resolves to %s', async (_label, node) => { + ['self', { node: 'self' }], + ['an empty node', { node: '' }], + ['an absent node', {}], + ])('fails closed when placement resolves to %s', async (_label, placement) => { const messaging = new FakeMessaging() - messaging.placementAck = { placement: { node } } + messaging.placementAck = { placement } const fleet = createClient(messaging) await expect(fleet.spawn({ @@ -367,6 +372,32 @@ describe('RelayFleetClient', () => { expect(fleet.trackedAgents().size).toBe(0) }) + it('returns and tracks the normalized acknowledgement node instead of action output', async () => { + const messaging = new FakeMessaging() + messaging.placementAck = { + invocationId: 'remote-placement', + status: 'pending', + placement: { node: ' mac-mini ' }, + } + messaging.invocations.set('remote-placement', [{ + invocationId: 'remote-placement', + actionName: 'spawn', + status: 'completed', + output: { name: 'ar-1-impl', node: 'wrong-node' }, + }]) + const fleet = createClient(messaging) + + await expect(fleet.spawn({ + name: 'ar-1-impl', + capability: 'spawn:codex', + node: 'self', + repo: 'AgentWorkforce/factory', + })).resolves.toMatchObject({ node: 'mac-mini' }) + + expect(fleet.trackedAgents().get('ar-1-impl')).toMatchObject({ node: 'mac-mini' }) + expect(messaging.invokes).not.toContainEqual(expect.objectContaining({ name: 'release' })) + }) + it('surfaces a self-placement refusal as a non-zero CLI result', async () => { const messaging = new FakeMessaging() messaging.placementAck = { placement: { node: 'self' } } diff --git a/src/fleet/relay-fleet-client.ts b/src/fleet/relay-fleet-client.ts index 3bb4cea..db2afb1 100644 --- a/src/fleet/relay-fleet-client.ts +++ b/src/fleet/relay-fleet-client.ts @@ -187,8 +187,9 @@ export class RelayFleetClient implements FleetClient { }) const invocation = await this.#awaitInvocation(ack.actionName || 'spawn', ack) const result = spawnResultFromInvocation(input.name, input.sessionRef, invocation, ack) + let acknowledgedNode: string try { - assertNamedRemotePlacement(result, ack.placement?.node) + acknowledgedNode = assertNamedRemotePlacement(result, ack.placement?.node) } catch (error) { // A completed placement invocation has already launched the worker. If // Relay cannot prove that it ran on a named remote node, tear that worker @@ -203,8 +204,9 @@ export class RelayFleetClient implements FleetClient { } throw error } - this.#track(result.name, ack) - return result + const trustedResult = { ...result, node: acknowledgedNode } + this.#track(trustedResult.name, { invocationId: ack.invocationId, node: acknowledgedNode }) + return trustedResult } async resume(input: { @@ -539,10 +541,10 @@ export class RelayFleetClient implements FleetClient { return invocation } - #track(name: string, ack: { invocationId: string; dispatchedNodeId?: string | null; placement?: { node?: string } }): void { + #track(name: string, placement: { invocationId: string; node: string }): void { this.#tracked.set(name, { - invocationId: ack.invocationId, - node: ack.placement?.node ?? ack.dispatchedNodeId ?? undefined, + invocationId: placement.invocationId, + node: placement.node, spawnedAtMs: this.#now(), }) this.#syncExitWatcher() @@ -823,7 +825,7 @@ function spawnResultFromInvocation( } } -function assertNamedRemotePlacement(result: SpawnResult, acknowledgedNode: string | undefined): void { +function assertNamedRemotePlacement(result: SpawnResult, acknowledgedNode: string | undefined): string { // The placement acknowledgement is authoritative. Action output may name // the node executing the spawn handler even when Relay acknowledged `self`, // so accepting the synthesized SpawnResult would let self-placement pass. @@ -834,6 +836,7 @@ function assertNamedRemotePlacement(result: SpawnResult, acknowledgedNode: strin `refusing to accept the spawn result`, ) } + return node } function lifecycleSignalFromInvocation( From 35ab71f2b06911a01953e2e72de46005d0d386d9 Mon Sep 17 00:00:00 2001 From: Khaliq Date: Mon, 17 Aug 2026 12:22:00 +0200 Subject: [PATCH 4/7] fix: retain unverified workers until release succeeds --- src/fleet/relay-fleet-client.test.ts | 47 ++++++++++++++++++++++++ src/fleet/relay-fleet-client.ts | 53 ++++++++++++++++++++-------- 2 files changed, 86 insertions(+), 14 deletions(-) diff --git a/src/fleet/relay-fleet-client.test.ts b/src/fleet/relay-fleet-client.test.ts index 262fc77..cd0cf6e 100644 --- a/src/fleet/relay-fleet-client.test.ts +++ b/src/fleet/relay-fleet-client.test.ts @@ -372,6 +372,53 @@ describe('RelayFleetClient', () => { expect(fleet.trackedAgents().size).toBe(0) }) + it('retains and retries cleanup when release of an unverified placement fails', async () => { + const messaging = new FakeMessaging() + messaging.placementAck = { + invocationId: 'self-placement', + status: 'pending', + placement: { node: 'self' }, + } + messaging.invocations.set('self-placement', [{ + invocationId: 'self-placement', + actionName: 'spawn', + status: 'completed', + output: { name: 'ar-1-impl', node: 'untrusted-output-node' }, + }]) + messaging.invocations.set('inv-1', [{ + invocationId: 'inv-1', + actionName: 'release', + status: 'failed', + error: 'temporary cleanup failure', + }]) + messaging.agentRows = [{ name: 'ar-1-impl', status: 'online' }] + const fleet = createClient(messaging) + + await expect(fleet.spawn({ + name: 'ar-1-impl', + capability: 'spawn:codex', + node: 'self', + repo: 'AgentWorkforce/factory', + })).rejects.toThrow('Relay placement did not prove a named remote node') + + expect(fleet.trackedAgents().get('ar-1-impl')).toMatchObject({ + invocationId: 'self-placement', + pendingReleaseReason: 'unverified-placement', + }) + expect(fleet.trackedAgents().get('ar-1-impl')).not.toHaveProperty('node') + + messaging.invocations.set('inv-2', [{ + invocationId: 'inv-2', + actionName: 'release', + status: 'completed', + output: {}, + }]) + await fleet.reconcileTrackedAgents() + + expect(messaging.invokes.filter((invoke) => invoke.name === 'release')).toHaveLength(2) + expect(fleet.trackedAgents().size).toBe(0) + }) + it('returns and tracks the normalized acknowledgement node instead of action output', async () => { const messaging = new FakeMessaging() messaging.placementAck = { diff --git a/src/fleet/relay-fleet-client.ts b/src/fleet/relay-fleet-client.ts index db2afb1..68fdeb9 100644 --- a/src/fleet/relay-fleet-client.ts +++ b/src/fleet/relay-fleet-client.ts @@ -23,6 +23,8 @@ export interface TrackedAgent { node?: string spawnedAtMs: number nodeOfflineSinceMs?: number + /** A rejected spawn whose compensating release must be retried. */ + pendingReleaseReason?: string } export interface RelayClientLike { @@ -198,6 +200,14 @@ export class RelayFleetClient implements FleetClient { try { await this.release(result.name, 'unverified-placement') } catch (releaseError) { + // Do not forget a live worker merely because its compensating release + // failed. Retain it outside the accepted placement result and let the + // reconciliation loop retry the idempotent release until Relay + // confirms cleanup (or roster evidence proves the worker exited). + this.#track(result.name, { + invocationId: ack.invocationId, + pendingReleaseReason: 'unverified-placement', + }) this.#log( `Failed to release ${result.name} after unverified Relay placement: ${errorMessage(releaseError)}`, ) @@ -234,17 +244,14 @@ export class RelayFleetClient implements FleetClient { async release(name: string, reason?: string): Promise { const messaging = await this.#ensureMessaging() - try { - const ack = await messaging.commands.invoke('release', { - name, - agent: name, - ...(reason ? { reason } : {}), - }) - await this.#awaitInvocation(ack.actionName || 'release', ack) - } finally { - this.#tracked.delete(name) - this.#syncExitWatcher() - } + const ack = await messaging.commands.invoke('release', { + name, + agent: name, + ...(reason ? { reason } : {}), + }) + await this.#awaitInvocation(ack.actionName || 'release', ack) + this.#tracked.delete(name) + this.#syncExitWatcher() } async createPreview(input: PreviewStartInput): Promise { @@ -541,17 +548,27 @@ export class RelayFleetClient implements FleetClient { return invocation } - #track(name: string, placement: { invocationId: string; node: string }): void { + #track(name: string, placement: { + invocationId: string + node?: string + pendingReleaseReason?: string + }): void { this.#tracked.set(name, { invocationId: placement.invocationId, - node: placement.node, + ...(placement.node ? { node: placement.node } : {}), + ...(placement.pendingReleaseReason + ? { pendingReleaseReason: placement.pendingReleaseReason } + : {}), spawnedAtMs: this.#now(), }) this.#syncExitWatcher() } #syncExitWatcher(): void { - const shouldRun = !this.#disposed && this.#tracked.size > 0 && this.#agentExitListeners.size > 0 + const hasPendingRelease = [...this.#tracked.values()].some((agent) => agent.pendingReleaseReason) + const shouldRun = !this.#disposed && this.#tracked.size > 0 && ( + this.#agentExitListeners.size > 0 || hasPendingRelease + ) if (shouldRun && !this.#watchTimer) { this.#watchTimer = setInterval(() => { void this.reconcileTrackedAgents().catch((error) => { @@ -579,6 +596,14 @@ export class RelayFleetClient implements FleetClient { const nowMs = this.#now() for (const [name, entry] of [...this.#tracked]) { + if (entry.pendingReleaseReason) { + try { + await this.release(name, entry.pendingReleaseReason) + continue + } catch (error) { + this.#log(`Pending release retry failed for ${name}: ${errorMessage(error)}`) + } + } if (!onlineAgentNames.has(name)) { // A just-spawned agent may not have registered with the engine yet; // absence only counts as an exit once the grace window has passed. From a8b8e829f4cc198ad8c4cc5a11030eca1d6f7113 Mon Sep 17 00:00:00 2001 From: Khaliq Date: Mon, 17 Aug 2026 12:29:26 +0200 Subject: [PATCH 5/7] fix(fleet): drain rejected placement cleanup on dispose --- src/fleet/relay-fleet-client.test.ts | 82 ++++++++++++++++++++++++++++ src/fleet/relay-fleet-client.ts | 19 +++++++ 2 files changed, 101 insertions(+) diff --git a/src/fleet/relay-fleet-client.test.ts b/src/fleet/relay-fleet-client.test.ts index cd0cf6e..0320b0b 100644 --- a/src/fleet/relay-fleet-client.test.ts +++ b/src/fleet/relay-fleet-client.test.ts @@ -419,6 +419,88 @@ describe('RelayFleetClient', () => { expect(fleet.trackedAgents().size).toBe(0) }) + it('drains a pending rejected-placement release before disposal', async () => { + const messaging = new FakeMessaging() + messaging.placementAck = { + invocationId: 'self-placement', + status: 'pending', + placement: { node: 'self' }, + } + messaging.invocations.set('self-placement', [{ + invocationId: 'self-placement', + actionName: 'spawn', + status: 'completed', + output: { name: 'ar-1-impl' }, + }]) + messaging.invocations.set('inv-1', [{ + invocationId: 'inv-1', + actionName: 'release', + status: 'failed', + error: 'temporary cleanup failure', + }]) + messaging.agentRows = [{ name: 'ar-1-impl', status: 'online' }] + const fleet = createClient(messaging) + + await expect(fleet.spawn({ + name: 'ar-1-impl', + capability: 'spawn:codex', + node: 'self', + })).rejects.toThrow('Relay placement did not prove a named remote node') + + messaging.invocations.set('inv-2', [{ + invocationId: 'inv-2', + actionName: 'release', + status: 'completed', + output: {}, + }]) + await expect(fleet.dispose()).resolves.toBeUndefined() + + expect(messaging.invokes.filter((invoke) => invoke.name === 'release')).toHaveLength(2) + expect(fleet.trackedAgents().size).toBe(0) + }) + + it('refuses disposal without erasing an unconfirmed rejected-placement release', async () => { + const messaging = new FakeMessaging() + messaging.placementAck = { + invocationId: 'self-placement', + status: 'pending', + placement: { node: 'self' }, + } + messaging.invocations.set('self-placement', [{ + invocationId: 'self-placement', + actionName: 'spawn', + status: 'completed', + output: { name: 'ar-1-impl' }, + }]) + messaging.invocations.set('inv-1', [{ + invocationId: 'inv-1', + actionName: 'release', + status: 'failed', + error: 'cleanup unavailable', + }]) + messaging.invocations.set('inv-2', [{ + invocationId: 'inv-2', + actionName: 'release', + status: 'failed', + error: 'cleanup still unavailable', + }]) + messaging.agentRows = [{ name: 'ar-1-impl', status: 'online' }] + const fleet = createClient(messaging) + + await expect(fleet.spawn({ + name: 'ar-1-impl', + capability: 'spawn:codex', + node: 'self', + })).rejects.toThrow('Relay placement did not prove a named remote node') + + await expect(fleet.dispose()).rejects.toThrow( + 'Refusing to dispose Relay fleet client with unconfirmed worker cleanup: ar-1-impl', + ) + expect(fleet.trackedAgents().get('ar-1-impl')).toMatchObject({ + pendingReleaseReason: 'unverified-placement', + }) + }) + it('returns and tracks the normalized acknowledgement node instead of action output', async () => { const messaging = new FakeMessaging() messaging.placementAck = { diff --git a/src/fleet/relay-fleet-client.ts b/src/fleet/relay-fleet-client.ts index 68fdeb9..e9dd95d 100644 --- a/src/fleet/relay-fleet-client.ts +++ b/src/fleet/relay-fleet-client.ts @@ -430,6 +430,25 @@ export class RelayFleetClient implements FleetClient { async dispose(): Promise { if (this.#disposed) return + + // A rejected placement may already have launched a worker. If its first + // compensating release failed, disposal must not erase the only ownership + // record before giving cleanup one final synchronous retry. Refuse to + // dispose while any such release remains unconfirmed; callers then get a + // hard failure containing the deterministic worker name instead of a + // successful shutdown that silently strands it. + if ([...this.#tracked.values()].some((agent) => agent.pendingReleaseReason)) { + await this.reconcileTrackedAgents() + const pendingNames = [...this.#tracked] + .filter(([, agent]) => agent.pendingReleaseReason) + .map(([name]) => name) + if (pendingNames.length > 0) { + throw new Error( + `Refusing to dispose Relay fleet client with unconfirmed worker cleanup: ${pendingNames.join(', ')}`, + ) + } + } + this.#disposed = true if (this.#watchTimer) { clearInterval(this.#watchTimer) From 99e22687b01cb7480a0d83785200846dc506fe43 Mon Sep 17 00:00:00 2001 From: Khaliq Date: Mon, 17 Aug 2026 12:38:45 +0200 Subject: [PATCH 6/7] fix(fleet): retry cleanup before roster checks --- src/fleet/relay-fleet-client.test.ts | 42 ++++++++++++++++++++++++++++ src/fleet/relay-fleet-client.ts | 27 ++++++++++++------ 2 files changed, 60 insertions(+), 9 deletions(-) diff --git a/src/fleet/relay-fleet-client.test.ts b/src/fleet/relay-fleet-client.test.ts index 0320b0b..99a8fca 100644 --- a/src/fleet/relay-fleet-client.test.ts +++ b/src/fleet/relay-fleet-client.test.ts @@ -459,6 +459,48 @@ describe('RelayFleetClient', () => { expect(fleet.trackedAgents().size).toBe(0) }) + it('retries pending cleanup during disposal without depending on roster health', async () => { + const messaging = new FakeMessaging() + messaging.placementAck = { + invocationId: 'self-placement', + status: 'pending', + placement: { node: 'self' }, + } + messaging.invocations.set('self-placement', [{ + invocationId: 'self-placement', + actionName: 'spawn', + status: 'completed', + output: { name: 'ar-1-impl' }, + }]) + messaging.invocations.set('inv-1', [{ + invocationId: 'inv-1', + actionName: 'release', + status: 'failed', + error: 'temporary cleanup failure', + }]) + const presence = vi.spyOn(messaging.agents, 'presence') + .mockRejectedValue(new Error('roster unavailable')) + const fleet = createClient(messaging) + + await expect(fleet.spawn({ + name: 'ar-1-impl', + capability: 'spawn:codex', + node: 'self', + })).rejects.toThrow('Relay placement did not prove a named remote node') + + messaging.invocations.set('inv-2', [{ + invocationId: 'inv-2', + actionName: 'release', + status: 'completed', + output: {}, + }]) + await expect(fleet.dispose()).resolves.toBeUndefined() + + expect(presence).not.toHaveBeenCalled() + expect(messaging.invokes.filter((invoke) => invoke.name === 'release')).toHaveLength(2) + expect(fleet.trackedAgents().size).toBe(0) + }) + it('refuses disposal without erasing an unconfirmed rejected-placement release', async () => { const messaging = new FakeMessaging() messaging.placementAck = { diff --git a/src/fleet/relay-fleet-client.ts b/src/fleet/relay-fleet-client.ts index e9dd95d..9246ddb 100644 --- a/src/fleet/relay-fleet-client.ts +++ b/src/fleet/relay-fleet-client.ts @@ -438,7 +438,7 @@ export class RelayFleetClient implements FleetClient { // hard failure containing the deterministic worker name instead of a // successful shutdown that silently strands it. if ([...this.#tracked.values()].some((agent) => agent.pendingReleaseReason)) { - await this.reconcileTrackedAgents() + await this.#retryPendingReleases() const pendingNames = [...this.#tracked] .filter(([, agent]) => agent.pendingReleaseReason) .map(([name]) => name) @@ -603,6 +603,12 @@ export class RelayFleetClient implements FleetClient { async #reconcileTracked(): Promise { if (this.#tracked.size === 0) return + // Cleanup is more important than roster metadata and does not depend on + // it. Retry rejected-placement releases first so a presence/node outage + // cannot prevent the compensating action. + await this.#retryPendingReleases() + if (this.#tracked.size === 0) return + const messaging = await this.#ensureMessaging() const [presence, nodes] = await Promise.all([ messaging.agents.presence(), @@ -615,14 +621,6 @@ export class RelayFleetClient implements FleetClient { const nowMs = this.#now() for (const [name, entry] of [...this.#tracked]) { - if (entry.pendingReleaseReason) { - try { - await this.release(name, entry.pendingReleaseReason) - continue - } catch (error) { - this.#log(`Pending release retry failed for ${name}: ${errorMessage(error)}`) - } - } if (!onlineAgentNames.has(name)) { // A just-spawned agent may not have registered with the engine yet; // absence only counts as an exit once the grace window has passed. @@ -642,6 +640,17 @@ export class RelayFleetClient implements FleetClient { } } + async #retryPendingReleases(): Promise { + for (const [name, entry] of [...this.#tracked]) { + if (!entry.pendingReleaseReason) continue + try { + await this.release(name, entry.pendingReleaseReason) + } catch (error) { + this.#log(`Pending release retry failed for ${name}: ${errorMessage(error)}`) + } + } + } + #emitExit(name: string, reason: string): void { this.#tracked.delete(name) this.#syncExitWatcher() From 7c59720951e3cdbf79d1cc60222342b1dfb53139 Mon Sep 17 00:00:00 2001 From: Khaliq Date: Mon, 17 Aug 2026 12:55:45 +0200 Subject: [PATCH 7/7] fix: serialize rejected worker cleanup --- src/fleet/relay-fleet-client.test.ts | 52 ++++++++++++++++++++++++++++ src/fleet/relay-fleet-client.ts | 10 +++++- 2 files changed, 61 insertions(+), 1 deletion(-) diff --git a/src/fleet/relay-fleet-client.test.ts b/src/fleet/relay-fleet-client.test.ts index 99a8fca..76f936d 100644 --- a/src/fleet/relay-fleet-client.test.ts +++ b/src/fleet/relay-fleet-client.test.ts @@ -459,6 +459,58 @@ describe('RelayFleetClient', () => { expect(fleet.trackedAgents().size).toBe(0) }) + it('shares a pending release retry between reconciliation and disposal', async () => { + const messaging = new FakeMessaging() + messaging.placementAck = { + invocationId: 'self-placement', + status: 'pending', + placement: { node: 'self' }, + } + messaging.invocations.set('self-placement', [{ + invocationId: 'self-placement', + actionName: 'spawn', + status: 'completed', + output: { name: 'ar-1-impl' }, + }]) + messaging.invocations.set('inv-1', [{ + invocationId: 'inv-1', + actionName: 'release', + status: 'failed', + error: 'temporary cleanup failure', + }]) + messaging.agentRows = [{ name: 'ar-1-impl', status: 'online' }] + const fleet = createClient(messaging) + + await expect(fleet.spawn({ + name: 'ar-1-impl', + capability: 'spawn:codex', + node: 'self', + })).rejects.toThrow('Relay placement did not prove a named remote node') + + let unblockRelease!: () => void + const releaseBlocked = new Promise((resolve) => { + unblockRelease = resolve + }) + const invoke = messaging.commands.invoke.bind(messaging.commands) + vi.spyOn(messaging.commands, 'invoke').mockImplementation(async (name, input) => { + if (name !== 'release') return await invoke(name, input) + messaging.invokes.push({ name, input }) + await releaseBlocked + return { invocationId: 'shared-release', actionName: name, status: 'completed' } + }) + + const reconciliation = fleet.reconcileTrackedAgents() + await flush() + const disposal = fleet.dispose() + await flush() + + expect(messaging.invokes.filter((candidate) => candidate.name === 'release')).toHaveLength(2) + unblockRelease() + await expect(Promise.all([reconciliation, disposal])).resolves.toEqual([undefined, undefined]) + expect(messaging.invokes.filter((candidate) => candidate.name === 'release')).toHaveLength(2) + expect(fleet.trackedAgents().size).toBe(0) + }) + it('retries pending cleanup during disposal without depending on roster health', async () => { const messaging = new FakeMessaging() messaging.placementAck = { diff --git a/src/fleet/relay-fleet-client.ts b/src/fleet/relay-fleet-client.ts index 9246ddb..6d75526 100644 --- a/src/fleet/relay-fleet-client.ts +++ b/src/fleet/relay-fleet-client.ts @@ -117,6 +117,7 @@ export class RelayFleetClient implements FleetClient { #disposed = false #watchTimer: ReturnType | undefined #reconciling: Promise | undefined + #pendingReleaseRetry: Promise | undefined constructor(options: RelayFleetClientOptions = {}) { this.#options = options @@ -640,7 +641,14 @@ export class RelayFleetClient implements FleetClient { } } - async #retryPendingReleases(): Promise { + #retryPendingReleases(): Promise { + this.#pendingReleaseRetry ??= this.#runPendingReleaseRetries().finally(() => { + this.#pendingReleaseRetry = undefined + }) + return this.#pendingReleaseRetry + } + + async #runPendingReleaseRetries(): Promise { for (const [name, entry] of [...this.#tracked]) { if (!entry.pendingReleaseReason) continue try {