diff --git a/CHANGELOG.md b/CHANGELOG.md index db3d3054..13137b8c 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 + +- Reconnecting brokers recover same-node legacy `default` workers without letting another provider or node claim their identities; adoption rechecks provider liveness and ignores future-dated default heartbeats. ## [8.11.6] - 2026-09-20 diff --git a/packages/engine/CHANGELOG.md b/packages/engine/CHANGELOG.md index 19180155..6b34d40e 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 + +- Inventory reconciliation adopts a bound legacy `default` worker into its named broker only when the agent ID and node match and no live default provider owns it; conflicting claims still fail closed; adoption rechecks provider liveness and ignores future-dated default heartbeats. ## [8.11.6] - 2026-09-20 diff --git a/packages/engine/src/__tests__/conformance/inventoryPresenceIsolation.test.ts b/packages/engine/src/__tests__/conformance/inventoryPresenceIsolation.test.ts index 8b2ac2f2..760383ad 100644 --- a/packages/engine/src/__tests__/conformance/inventoryPresenceIsolation.test.ts +++ b/packages/engine/src/__tests__/conformance/inventoryPresenceIsolation.test.ts @@ -1,4 +1,4 @@ -import { afterEach, beforeEach, describe, expect, it } from 'vitest'; +import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest'; import { and, eq, inArray } from 'drizzle-orm'; import { FLEET_DELIVERY_CURSOR_CAPABILITY } from '@relaycast/types'; import { @@ -8,7 +8,8 @@ import { registerAgent, type TestStack, } from './harness.js'; -import { agents, nodes } from '../../db/schema.js'; +import { agents, nodeProviders, nodes } from '../../db/schema.js'; +import { NODE_LIVENESS_TTL_MS } from '../../engine/placement.js'; import { AGENT_LIVENESS_TTL_MS } from '../../engine/agent.js'; type Workspace = Awaited>; @@ -37,7 +38,12 @@ describe('node inventory presence isolation', () => { afterEach(() => stack.close()); - async function connectNode(ws: Workspace, id: string, name: string): Promise { + async function connectNode( + ws: Workspace, + id: string, + name: string, + providerName = 'broker', + ): Promise { const sock = new FakeSocket(); const handle = stack.runtime.realtime.attachNodeSocket(ws.workspaceId, id, sock); await handle.handleMessage(JSON.stringify({ @@ -46,7 +52,7 @@ describe('node inventory presence isolation', () => { type: 'node.register', name, node_id: id, - provider: { name: 'broker', instance_id: `${id}-broker` }, + provider: { name: providerName, instance_id: `${id}-${providerName}` }, capabilities: [ { name: 'spawn:codex', kind: 'capacity' }, { name: FLEET_DELIVERY_CURSOR_CAPABILITY, kind: 'capacity' }, @@ -60,13 +66,13 @@ describe('node inventory presence isolation', () => { if (!registerReply) throw new Error(`node registration failed: ${JSON.stringify(sock.received)}`); expect(registerReply).toMatchObject({ ok: true, - data: { provider: { name: 'broker' } }, + data: { provider: { name: providerName } }, }); await handle.handleMessage(JSON.stringify({ v: 1, id: `heartbeat-${id}`, type: 'node.heartbeat', - provider: { name: 'broker', instance_id: `${id}-broker` }, + provider: { name: providerName, instance_id: `${id}-${providerName}` }, load: 0, active_agents: 0, handlers_live: true, @@ -173,6 +179,198 @@ describe('node inventory presence isolation', () => { })); } + async function setProvider(agent: NodeAgent, providerName: string, status: 'active' | 'offline') { + await stack.runtime.handle.db + .update(agents) + .set({ providerName, status }) + .where(eq(agents.id, agent.agentId)); + } + + it.each(['active', 'offline'] as const)( + 'adopts a %s same-node legacy default worker into the broker and restores delivery', + async (status) => { + const ws = await createWorkspace(stack.app, `inventory-legacy-default-${status}`); + const sender = await registerAgent(stack.app, ws.workspaceKey, `legacy-sender-${status}`); + let node = await attachNode(ws, `node_legacy_${status}`, `legacy-node-${status}`); + const worker = await registerViaNode(node, `legacy-worker-${status}`); + await setProvider(worker, 'default', status); + + node = await reconnectNode(ws, node); + await syncInventory(node, `legacy-renewal-${status}`, [worker]); + expect(node.sock.ofType('reply').find((frame) => frame.id === `legacy-renewal-${status}`)).toMatchObject({ + ok: true, + data: { rebound_agents: 1, rejected_agents: 0 }, + }); + + const [row] = await stack.runtime.handle.db + .select({ status: agents.status, providerName: agents.providerName, locationNodeId: agents.locationNodeId }) + .from(agents) + .where(eq(agents.id, worker.agentId)); + expect(row).toEqual({ status: 'active', providerName: 'broker', locationNodeId: node.id }); + + node.sock.received.length = 0; + await postFrom(sender.token, 'delivery after legacy adoption'); + expect(node.sock.ofType('deliver').filter((frame) => frame.agent === worker.name)).toHaveLength(1); + await ackFirstDelivery(node, worker); + expect((await readAgent(ws, worker.name)).pending_deliveries).toHaveLength(0); + }, + ); + + it('does not let a broker claim a worker while the default provider is live', async () => { + const ws = await createWorkspace(stack.app, 'inventory-live-default-provider'); + const broker = await attachNode(ws, 'node_live_default', 'live-default-node'); + const legacy = await connectNode(ws, broker.id, broker.name, 'default'); + const worker = await registerViaNode(legacy, 'live-default-worker'); + + await syncInventory(broker, 'live-default-claim', [worker]); + expect(broker.sock.ofType('error').find((frame) => frame.id === 'live-default-claim')).toMatchObject({ + code: 'agent_provider_conflict', + }); + const [row] = await stack.runtime.handle.db + .select({ providerName: agents.providerName }) + .from(agents) + .where(eq(agents.id, worker.agentId)); + expect(row.providerName).toBe('default'); + }); + + it('adopts a worker after its persisted default provider disconnects', async () => { + const ws = await createWorkspace(stack.app, 'inventory-disconnected-default-provider'); + const broker = await attachNode(ws, 'node_disconnected_default', 'disconnected-default-node'); + const legacy = await connectNode(ws, broker.id, broker.name, 'default'); + const worker = await registerViaNode(legacy, 'disconnected-default-worker'); + await legacy.handle.handleClose(); + + await syncInventory(broker, 'disconnected-default-renewal', [worker]); + expect(broker.sock.ofType('reply').find((frame) => frame.id === 'disconnected-default-renewal')).toMatchObject({ + ok: true, + data: { rebound_agents: 1, rejected_agents: 0 }, + }); + const [row] = await stack.runtime.handle.db + .select({ providerName: agents.providerName, status: agents.status }) + .from(agents) + .where(eq(agents.id, worker.agentId)); + expect(row).toEqual({ providerName: 'broker', status: 'active' }); + }); + + it('adopts a worker when the default provider heartbeat is in the future', async () => { + const ws = await createWorkspace(stack.app, 'inventory-future-default'); + const broker = await attachNode(ws, 'node_future_default', 'future-default'); + const legacy = await connectNode(ws, broker.id, broker.name, 'default'); + const worker = await registerViaNode(legacy, 'future-default-worker'); + await stack.runtime.handle.db.update(nodeProviders) + .set({ lastHeartbeatAt: new Date(Date.now() + 60_000) }) + .where(and(eq(nodeProviders.nodeId, broker.id), eq(nodeProviders.name, 'default'))); + + await syncInventory(broker, 'future-default-adoption', [worker]); + expect(broker.sock.ofType('reply').find((frame) => frame.id === 'future-default-adoption')).toMatchObject({ + ok: true, + data: { rebound_agents: 1, rejected_agents: 0 }, + }); + const [row] = await stack.runtime.handle.db.select().from(agents).where(eq(agents.id, worker.agentId)); + expect(row.providerName).toBe('broker'); + }); + + it.each(['offline', 'handlers-down', 'expired', 'future', 'missing-heartbeat', 'default-reconnected'])( + 'fences legacy adoption when provider liveness changes before the update: %s', + async (change) => { + const ws = await createWorkspace(stack.app, `inventory-race-${change}`); + const broker = await attachNode(ws, 'node_race', 'race-node'); + const legacy = await connectNode(ws, broker.id, broker.name, 'default'); + const worker = await registerViaNode(legacy, 'race-worker'); + await legacy.handle.handleClose(); + await setProvider(worker, 'default', 'offline'); + await stack.settle(); + const db = stack.runtime.handle.db; + const update = db.update.bind(db); + let changed = false; + // Change persisted liveness after validation, immediately before the + // ownership update is built. The real database evaluates the CAS. + const spy = vi.spyOn(db, 'update').mockImplementation((table) => { + if (table === agents && !changed) { + changed = true; + update(nodeProviders).set(change === 'default-reconnected' + ? { status: 'online', handlersLive: true, lastHeartbeatAt: new Date() } + : change === 'offline' ? { status: 'offline' } + : change === 'handlers-down' ? { handlersLive: false } + : { lastHeartbeatAt: change === 'missing-heartbeat' ? null + : new Date(Date.now() + (change === 'future' ? 60_000 : -NODE_LIVENESS_TTL_MS - 1_000)) }) + .where(and(eq(nodeProviders.nodeId, broker.id), + eq(nodeProviders.name, change === 'default-reconnected' ? 'default' : 'broker'))).run(); + } + return update(table); + }); + try { + await syncInventory(broker, 'race-adoption', [worker]); + expect(changed).toBe(true); + expect(broker.sock.ofType('error').find((frame) => frame.id === 'race-adoption')).toMatchObject({ + code: 'agent_provider_conflict', + }); + const [row] = await db.select().from(agents).where(eq(agents.id, worker.agentId)); + expect(row).toMatchObject({ providerName: 'default', status: 'offline' }); + } finally { + spy.mockRestore(); + } + }, + ); + + it('does not adopt a legacy worker with a wrong id or from another node', async () => { + const ws = await createWorkspace(stack.app, 'inventory-legacy-foreign-claim'); + const owner = await attachNode(ws, 'node_legacy_owner', 'legacy-owner'); + const worker = await registerViaNode(owner, 'legacy-owned-worker'); + await setProvider(worker, 'default', 'offline'); + + await syncInventory(owner, 'wrong-id-claim', [{ ...worker, agentId: 'agt_wrong_identity' }]); + expect(owner.sock.ofType('error').find((frame) => frame.id === 'wrong-id-claim')).toMatchObject({ + code: 'agent_provider_conflict', + }); + + const foreign = await attachNode(ws, 'node_legacy_foreign', 'legacy-foreign'); + await syncInventory(foreign, 'foreign-node-claim', [worker]); + expect(foreign.sock.ofType('error').find((frame) => frame.id === 'foreign-node-claim')).toMatchObject({ + code: 'agent_provider_conflict', + }); + const [row] = await stack.runtime.handle.db + .select({ providerName: agents.providerName, locationNodeId: agents.locationNodeId }) + .from(agents) + .where(eq(agents.id, worker.agentId)); + expect(row).toEqual({ providerName: 'default', locationNodeId: owner.id }); + }); + + it('does not transfer ownership from a different named provider', async () => { + const ws = await createWorkspace(stack.app, 'inventory-named-provider-claim'); + const broker = await attachNode(ws, 'node_named_provider', 'named-provider-node'); + const worker = await registerViaNode(broker, 'named-provider-worker'); + await setProvider(worker, 'other-provider', 'offline'); + + await syncInventory(broker, 'named-provider-claim', [worker]); + expect(broker.sock.ofType('error').find((frame) => frame.id === 'named-provider-claim')).toMatchObject({ + code: 'agent_provider_conflict', + }); + const [row] = await stack.runtime.handle.db + .select({ providerName: agents.providerName }) + .from(agents) + .where(eq(agents.id, worker.agentId)); + expect(row.providerName).toBe('other-provider'); + }); + + it('does not let a sibling named provider adopt the legacy default identity', async () => { + const ws = await createWorkspace(stack.app, 'inventory-sibling-provider-claim'); + const broker = await attachNode(ws, 'node_sibling_provider', 'sibling-provider-node'); + const worker = await registerViaNode(broker, 'sibling-provider-worker'); + await setProvider(worker, 'default', 'offline'); + const sibling = await connectNode(ws, broker.id, broker.name, 'sibling'); + + await syncInventory(sibling, 'sibling-provider-claim', [worker]); + expect(sibling.sock.ofType('error').find((frame) => frame.id === 'sibling-provider-claim')).toMatchObject({ + code: 'agent_provider_conflict', + }); + const [row] = await stack.runtime.handle.db + .select({ providerName: agents.providerName }) + .from(agents) + .where(eq(agents.id, worker.agentId)); + expect(row.providerName).toBe('default'); + }); + it('keeps a healthy two-agent node active and drains both deliveries after inventory renewal', async () => { const ws = await createWorkspace(stack.app, 'inventory-presence-control'); const sender = await registerAgent(stack.app, ws.workspaceKey, 'control-sender'); diff --git a/packages/engine/src/engine/node.ts b/packages/engine/src/engine/node.ts index 17b6821d..70940cff 100644 --- a/packages/engine/src/engine/node.ts +++ b/packages/engine/src/engine/node.ts @@ -1,6 +1,6 @@ import { z } from 'zod'; import { invalidateChannelCache } from './cache.js'; -import { and, asc, eq, gt, gte, inArray, isNotNull, isNull, lt, lte, ne, or, sql } from 'drizzle-orm'; +import { and, asc, eq, exists, gt, gte, inArray, isNotNull, isNull, lt, lte, ne, notExists, or, sql } from 'drizzle-orm'; import type { FleetAgentRecoverMessage, FleetAgentRegisterMessage, @@ -2085,6 +2085,8 @@ export async function reconcileInventory( const existingByName = new Map(); const acceptedInventoryAgents: FleetInventoryAgent[] = []; const rejectedInventoryErrors: Array> = []; + const legacyDefaultAdoptions = new Set(); + let canAdoptLegacyDefault: boolean | undefined; let acceptedExistingAgentCount = 0; for (const item of inventoryAgents) { const [existing] = await db @@ -2097,12 +2099,42 @@ export async function reconcileInventory( } let rejection: ReturnType | undefined; if (existing.providerName !== providerName) { - rejection = codedError( - `Agent "${item.name}" belongs to provider "${existing.providerName}"`, - 'agent_provider_conflict', - 409, - ); - } else if (existing.status === 'active') { + // A pre-provider registration can be bound to this node while no + // provider is connected, leaving its row on synthetic `default` even + // though the surviving session is in the named broker's inventory. + // Only the broker on the *same* node may adopt that exact identity, and + // only after the old default provider has ceased to be live. Other + // provider/name conflicts retain the fail-closed behavior. + const legacyDefaultClaim = existing.providerName === DEFAULT_PROVIDER_NAME + && providerName === 'broker' + && existing.locationType === 'via_node' + && existing.locationNodeId === nodeId + && existing.id === item.agent_id; + if (legacyDefaultClaim && canAdoptLegacyDefault === undefined) { + const providerRows = await db + .select() + .from(nodeProviders) + .where(and( + eq(nodeProviders.workspaceId, workspaceId), + eq(nodeProviders.nodeId, nodeId), + inArray(nodeProviders.name, [DEFAULT_PROVIDER_NAME, 'broker']), + )); + const defaultProvider = providerRows.find((row) => row.name === DEFAULT_PROVIDER_NAME); + const brokerProvider = providerRows.find((row) => row.name === 'broker'); + canAdoptLegacyDefault = !!brokerProvider && isProviderLive(brokerProvider) + && (!defaultProvider || !isProviderLive(defaultProvider)); + } + if (legacyDefaultClaim && canAdoptLegacyDefault) { + legacyDefaultAdoptions.add(item.name); + } else { + rejection = codedError( + `Agent "${item.name}" belongs to provider "${existing.providerName}"`, + 'agent_provider_conflict', + 409, + ); + } + } + if (!rejection && existing.status === 'active') { const [boundNode] = await db .select() .from(nodes) @@ -2212,6 +2244,7 @@ export async function reconcileInventory( const targetWasActive = activeNodeIds.includes(nodeId); const wasRoutableThroughProvider = existing.locationType === 'via_node' && existing.locationNodeId === nodeId + && existing.providerName === providerName && targetWasActive; let reservedTargetSlot = false; try { @@ -2223,9 +2256,39 @@ export async function reconcileInventory( }); reservedTargetSlot = true; } - await db + const adoptingLegacyDefault = legacyDefaultAdoptions.has(item.name); + // The earlier liveness check is for a useful error decision. This + // compare-and-set fences broker liveness loss, a default provider + // reconnecting, and concurrent ownership changes during the apply pass. + const adoptionNow = Date.now(); + const noLiveDefaultProvider = notExists(db + .select({ id: nodeProviders.id }) + .from(nodeProviders) + .where(and( + eq(nodeProviders.workspaceId, workspaceId), + eq(nodeProviders.nodeId, nodeId), + eq(nodeProviders.name, DEFAULT_PROVIDER_NAME), + eq(nodeProviders.status, 'online'), + eq(nodeProviders.handlersLive, true), + gte(nodeProviders.lastHeartbeatAt, new Date(adoptionNow - NODE_LIVENESS_TTL_MS)), + lte(nodeProviders.lastHeartbeatAt, new Date(adoptionNow)), + ))); + const liveBrokerProvider = exists(db + .select({ id: nodeProviders.id }) + .from(nodeProviders) + .where(and( + eq(nodeProviders.workspaceId, workspaceId), + eq(nodeProviders.nodeId, nodeId), + eq(nodeProviders.name, 'broker'), + eq(nodeProviders.status, 'online'), + eq(nodeProviders.handlersLive, true), + gte(nodeProviders.lastHeartbeatAt, new Date(adoptionNow - NODE_LIVENESS_TTL_MS)), + lte(nodeProviders.lastHeartbeatAt, new Date(adoptionNow)), + ))); + const [updated] = await db .update(agents) .set({ + ...(adoptingLegacyDefault ? { providerName } : {}), status: 'active', lastSeen: new Date(), locationType: 'via_node', @@ -2233,7 +2296,21 @@ export async function reconcileInventory( originNodeId: existing.originNodeId ?? nodeId, sessionRef: item.session_ref ?? existing.sessionRef, }) - .where(eq(agents.id, existing.id)); + .where(and( + eq(agents.workspaceId, workspaceId), + eq(agents.id, existing.id), + ...(adoptingLegacyDefault ? [ + eq(agents.providerName, DEFAULT_PROVIDER_NAME), + eq(agents.locationType, 'via_node'), + eq(agents.locationNodeId, nodeId), + noLiveDefaultProvider, + liveBrokerProvider, + ] : []), + )) + .returning({ id: agents.id }); + if (!updated) { + throw codedError(`Agent "${item.name}" changed ownership during inventory reconciliation`, 'agent_provider_conflict', 409); + } await upsertAgentNodeBinding(db, workspaceId, existing, nodeId, { sessionRef: item.session_ref ?? existing.sessionRef, deactivateExisting: true,