From 140b1b5ec9729e238927ffd27c83f7a1d637e0f5 Mon Sep 17 00:00:00 2001 From: khaliqgant Date: Fri, 25 Sep 2026 06:48:18 -0700 Subject: [PATCH] fix(engine): accept a release result after the node deregistered the agent A relay broker completes a release by stopping the worker, queueing agent.deregister and then sending action.result on its one control channel. Since #448 routed plain releases through the guarded completion, that deregister left no active binding by the time the result arrived, so every broker-completed release failed with agent_release_generation_conflict. Accept the result when the agent was bound to the reporting node and has not been bound to any node but its own implicit direct node since. A release whose agent moved to a different node is still refused. Co-Authored-By: Claude Opus 5.5 (1M context) --- CHANGELOG.md | 6 +- packages/engine/CHANGELOG.md | 6 +- .../releaseAfterNodeDeregister.test.ts | 125 ++++++++++++++++++ packages/engine/src/engine/action.ts | 23 +++- 4 files changed, 157 insertions(+), 3 deletions(-) create mode 100644 packages/engine/src/__tests__/conformance/releaseAfterNodeDeregister.test.ts diff --git a/CHANGELOG.md b/CHANGELOG.md index a7781ba0..fc9c9089 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -16,7 +16,11 @@ This project follows [Semantic Versioning](https://semver.org/spec/v2.0.0.html). Packages without a separate changelog are covered by the cross-package notes below. -## [Unreleased] +## [Unreleased - Patch] + +### Fixed + +- Releasing an agent hosted by a relay broker completes again. It was failing with `agent_release_generation_conflict` because the broker deregisters the agent just before reporting the release result. ## [8.13.0] - 2026-09-25 diff --git a/packages/engine/CHANGELOG.md b/packages/engine/CHANGELOG.md index 03c3f62a..6ad7e9c2 100644 --- a/packages/engine/CHANGELOG.md +++ b/packages/engine/CHANGELOG.md @@ -7,7 +7,11 @@ See the [root changelog](../../CHANGELOG.md) for cross-package release highlight The format is based on [Keep a Changelog](https://keepachangelog.com/en/1.1.0/), and this project follows [Semantic Versioning](https://semver.org/spec/v2.0.0.html). -## [Unreleased] +## [Unreleased - Patch] + +### Fixed + +- A node-completed release is accepted when that node already deregistered the agent (a relay broker sends `agent.deregister` before `action.result`), instead of failing with `agent_release_generation_conflict`. A release whose agent has since been bound to a different node is still refused. ## [8.13.0] - 2026-09-25 diff --git a/packages/engine/src/__tests__/conformance/releaseAfterNodeDeregister.test.ts b/packages/engine/src/__tests__/conformance/releaseAfterNodeDeregister.test.ts new file mode 100644 index 00000000..2ec25a12 --- /dev/null +++ b/packages/engine/src/__tests__/conformance/releaseAfterNodeDeregister.test.ts @@ -0,0 +1,125 @@ +import { afterEach, beforeEach, describe, expect, it } from 'vitest'; +import { eq } from 'drizzle-orm'; +import { makeNodeStack, createWorkspace, FakeSocket, type TestStack } from './harness.js'; +import { actionInvocations, agentNodeBindings, agents } from '../../db/schema.js'; + +type Json = Record; + +/** + * A relay broker completes a release by stopping the worker, queueing + * `agent.deregister` and then sending `action.result`, in that order on its + * one control channel. The deregister deactivates the agent's binding on the + * node before the result arrives, so the completion must accept a binding the + * same release already removed, while still refusing an agent that has since + * been bound to a different node. + */ +describe('node release completed after the node deregistered the agent', () => { + let stack: TestStack; + beforeEach(() => { stack = makeNodeStack({ ttlMs: 60_000 }); }); + afterEach(async () => { + await stack.close(); + }); + + async function enrollBroker(ws: { workspaceKey: string; workspaceId: string }, nodeId: string, name: string) { + const enroll = await stack.app.request('/v1/nodes', { + method: 'POST', + headers: { 'content-type': 'application/json', authorization: `Bearer ${ws.workspaceKey}` }, + body: JSON.stringify({ node_id: nodeId, name, capabilities: ['spawn:claude'], max_agents: 4, version: 'test-node' }), + }); + expect(enroll.status).toBe(201); + const sock = new FakeSocket(); + const handle = stack.runtime.realtime.attachNodeSocket(ws.workspaceId, nodeId, sock); + await handle.handleMessage(JSON.stringify({ + v: 1, type: 'node.register', name, node_id: nodeId, + capabilities: [{ name: 'spawn:claude', kind: 'capacity' }, { name: 'release', kind: 'capacity' }], + max_agents: 4, tags: [], version: 'test-node', resume_cursor: null, + })); + await handle.handleMessage(JSON.stringify({ + v: 1, type: 'node.heartbeat', load: 0, active_agents: 0, handlers_live: true, + })); + return { sock, handle }; + } + + async function registerOnNode(node: Awaited>, name: string) { + await node.handle.handleMessage(JSON.stringify({ + v: 1, type: 'agent.register', name, resumable: true, session_ref: `sess-${name}`, + })); + const reply = node.sock.ofType('reply').at(-1) as { ok: boolean; data: { agent_id: string } }; + expect(reply?.ok).toBe(true); + return reply.data.agent_id; + } + + async function releaseDispatched(workspaceKey: string, name: string) { + const res = await stack.app.request('/v1/agents/release', { + method: 'POST', + headers: { 'content-type': 'application/json', authorization: `Bearer ${workspaceKey}` }, + body: JSON.stringify({ name }), + }); + expect(res.status).toBe(201); + const { data } = (await res.json()) as { data: { status: string; invocation_id: string } }; + expect(data.status).toBe('dispatched'); + return data.invocation_id; + } + + async function invocation(id: string) { + const [row] = await stack.runtime.deps.db + .select({ status: actionInvocations.status, error: actionInvocations.error }) + .from(actionInvocations) + .where(eq(actionInvocations.id, id)); + return row; + } + + it('completes a plain release whose node deregistered the agent before reporting the result', async () => { + const ws = await createWorkspace(stack.app, 'release-after-deregister'); + const node = await enrollBroker(ws, 'node_a', 'node-a'); + const agentId = await registerOnNode(node, 'worker'); + + const invocationId = await releaseDispatched(ws.workspaceKey, 'worker'); + expect(node.sock.ofType('action.invoke').map((frame) => (frame as Json).action)).toContain('release'); + + // The broker's order: deregister first, then the release result. + await node.handle.handleMessage(JSON.stringify({ v: 1, type: 'agent.deregister', agent_id: agentId, name: 'worker' })); + await node.handle.handleMessage(JSON.stringify({ + v: 1, type: 'action.result', invocation_id: invocationId, output: { released: true }, + })); + + expect(await invocation(invocationId)).toEqual({ status: 'completed', error: null }); + const [agent] = await stack.runtime.deps.db + .select({ status: agents.status, locationNodeId: agents.locationNodeId }) + .from(agents) + .where(eq(agents.id, agentId)); + expect(agent).toEqual({ status: 'offline', locationNodeId: null }); + }); + + it('still refuses the release once the agent is bound to a different node', async () => { + const ws = await createWorkspace(stack.app, 'release-after-move'); + const nodeA = await enrollBroker(ws, 'node_a', 'node-a'); + const agentId = await registerOnNode(nodeA, 'worker'); + + const invocationId = await releaseDispatched(ws.workspaceKey, 'worker'); + await nodeA.handle.handleMessage(JSON.stringify({ v: 1, type: 'agent.deregister', agent_id: agentId, name: 'worker' })); + // Before node A reports, the identity is bound to node B: the state a + // concurrent move leaves behind, written directly. + await enrollBroker(ws, 'node_b', 'node-b'); + const db = stack.runtime.deps.db; + await db.insert(agentNodeBindings).values({ + id: 'bind_moved_to_b', workspaceId: ws.workspaceId, agentId, nodeId: 'node_b', status: 'active', priority: 0, + }).onConflictDoUpdate({ + target: [agentNodeBindings.agentId, agentNodeBindings.nodeId], + set: { status: 'active' }, + }); + await db.update(agents).set({ locationType: 'via_node', locationNodeId: 'node_b', status: 'active' }).where(eq(agents.id, agentId)); + + await nodeA.handle.handleMessage(JSON.stringify({ + v: 1, type: 'action.result', invocation_id: invocationId, output: { released: true }, + })); + + expect(await invocation(invocationId)).toEqual({ status: 'failed', error: 'agent_release_generation_conflict' }); + const [agent] = await stack.runtime.deps.db + .select({ status: agents.status, locationNodeId: agents.locationNodeId }) + .from(agents) + .where(eq(agents.id, agentId)); + expect(agent.locationNodeId).toBe('node_b'); + expect(agent.status).not.toBe('offline'); + }); +}); diff --git a/packages/engine/src/engine/action.ts b/packages/engine/src/engine/action.ts index cbe0e6df..43e6d7da 100644 --- a/packages/engine/src/engine/action.ts +++ b/packages/engine/src/engine/action.ts @@ -2334,7 +2334,28 @@ async function completeReleaseNodeInvocation( AND ${actionInvocations.dispatchedProvider} = ${providerName} AND ${actionInvocations.status} = 'completed' )`; - const releaseCanApply = sql`(${generationStillCurrent}) AND (${activeBindingStillCurrent})`; + // A relay broker queues `agent.deregister` ahead of the release result on + // its one control channel, so the binding this release targets is usually + // already inactive when the result lands. Accept that, as long as the agent + // was bound to this node and has not since been bound to any other node + // than its own implicit direct node (where deregistration re-homes it). + const boundOnlyHere = sql`( + EXISTS ( + SELECT 1 FROM ${agentNodeBindings} + WHERE ${agentNodeBindings.workspaceId} = ${workspaceId} + AND ${agentNodeBindings.agentId} = ${agent.id} + AND ${agentNodeBindings.nodeId} = ${nodeId} + ) + AND NOT EXISTS ( + SELECT 1 FROM ${agentNodeBindings} + WHERE ${agentNodeBindings.workspaceId} = ${workspaceId} + AND ${agentNodeBindings.agentId} = ${agent.id} + AND ${agentNodeBindings.status} = 'active' + AND ${agentNodeBindings.nodeId} <> ${nodeId} + AND ${agentNodeBindings.nodeId} <> ${`node_direct_${agent.id}`} + ) + )`; + const releaseCanApply = sql`(${generationStillCurrent}) AND (${boundOnlyHere})`; const releasedName = releasedAgentName(agent.name, agent.id); const releasedTokenHash = await sha256Hex(`released:${agent.id}:${randomHex(16)}`); let successIndex = -1;