-
Notifications
You must be signed in to change notification settings - Fork 1
fix(engine): accept a release result after the node deregistered the agent #455
New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
base: main
Are you sure you want to change the base?
Changes from all commits
File filter
Filter by extension
Conversations
Jump to
Diff view
Diff view
There are no files selected for viewing
| Original file line number | Diff line number | Diff line change |
|---|---|---|
| @@ -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<string, unknown>; | ||
|
|
||
| /** | ||
| * 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<ReturnType<typeof enrollBroker>>, 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'); | ||
| }); | ||
| }); |
| Original file line number | Diff line number | Diff line change |
|---|---|---|
|
|
@@ -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}`} | ||
|
Comment on lines
+2353
to
+2355
Contributor
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. 🔴 Plain release strands direct-node capacity After broker deregistration, Learn moreBroker deregistration re-homes the agent using ensureDirectNodeForAgent, creating an active direct binding and setting the direct node's Example: Node A deregisters worker, creating an active Recommended fix: In Was this helpful? React with 👍 or 👎 to provide feedback.
Comment on lines
+2353
to
+2355
Contributor
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. 🔴 Broker result releases a new direct session If an agent reconnects directly after broker deregistration, Learn morederegisterAgentViaNode re-homes the agent to its implicit direct node, which has an active binding even while its node is offline. The direct connection can subsequently become live before the original broker reports its release result. This predicate treats that active binding exactly like the offline fallback. A plain release lacks an Example: A dispatches a plain release for worker. A deregisters worker; worker connects through Recommended fix: Distinguish a direct fallback left offline by deregistration from a subsequent active direct connection. Fence the completion against a newer direct session using an atomic agent/binding state or generation check; do not allow the reporting broker to release the newer session. Was this helpful? React with 👍 or 👎 to provide feedback.
Comment on lines
+2353
to
+2355
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more.
If the agent opens its direct-node connection after Useful? React with 👍 / 👎. |
||
| ) | ||
| )`; | ||
| const releaseCanApply = sql`(${generationStillCurrent}) AND (${boundOnlyHere})`; | ||
|
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more.
In the broker ordering this change is designed to accept, Useful? React with 👍 / 👎.
Contributor
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Release leaves implicit binding activeLow Severity Accepting a release after Additional Locations (1)Reviewed by Cursor Bugbot for commit 140b1b5. Configure here. |
||
| const releasedName = releasedAgentName(agent.name, agent.id); | ||
| const releasedTokenHash = await sha256Hex(`released:${agent.id}:${randomHex(16)}`); | ||
| let successIndex = -1; | ||
|
|
||


There was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.
For the ordinary non-deleting release covered by the new test,
agent.deregistercreates an active implicit-direct binding and sets that direct node'sactiveAgentsto 1. This predicate now accepts that state, but the completion batch only deactivates/decrements the dispatched broker node before setting the agent's location to null, leaving the fallback binding active and its node capacity occupied. Node-agent listings and later binding selection consequently continue to report the released agent as hosted; the accepted fallback binding and its capacity need to be retired as part of completion.Useful? React with 👍 / 👎.