From 32b274d7bf38c01b9e23d0d885db872b31ae55b6 Mon Sep 17 00:00:00 2001 From: kjgbot Date: Sun, 20 Sep 2026 10:00:38 -0700 Subject: [PATCH 1/2] fix(engine): drain deliveries when agents are released --- .../2026-09/traj_2t44mh839vla/summary.md | 32 ++++++++ .../2026-09/traj_2t44mh839vla/trajectory.json | 73 +++++++++++++++++++ CHANGELOG.md | 6 +- packages/engine/CHANGELOG.md | 6 +- .../conformance/agentLifecycle.test.ts | 21 ++++++ .../conformance/deleteAgentRoute.test.ts | 32 +++++++- .../conformance/nodeCompletedRelease.test.ts | 29 +++++++- .../__tests__/workspaceDeliveryDepth.test.ts | 26 +++++++ packages/engine/src/engine/action.ts | 30 +++++++- packages/engine/src/engine/agent.ts | 52 ++++++++++++- 10 files changed, 301 insertions(+), 6 deletions(-) create mode 100644 .agentworkforce/trajectories/completed/2026-09/traj_2t44mh839vla/summary.md create mode 100644 .agentworkforce/trajectories/completed/2026-09/traj_2t44mh839vla/trajectory.json diff --git a/.agentworkforce/trajectories/completed/2026-09/traj_2t44mh839vla/summary.md b/.agentworkforce/trajectories/completed/2026-09/traj_2t44mh839vla/summary.md new file mode 100644 index 00000000..c618cc5b --- /dev/null +++ b/.agentworkforce/trajectories/completed/2026-09/traj_2t44mh839vla/summary.md @@ -0,0 +1,32 @@ +# Trajectory: Drain workspace delivery capacity when agents are released + +> **Status:** ✅ Completed +> **Confidence:** 93% +> **Started:** September 20, 2026 at 09:53 AM +> **Completed:** September 20, 2026 at 10:00 AM + +--- + +## Summary + +Dead-lettered active delivery rows during irreversible agent release across direct, local, guarded node, and legacy node lifecycle paths; added release-path and expired-unswept workspace-cap regressions; full engine suite, typecheck, and lint pass. + +**Approach:** Standard approach + +--- + +## Key Decisions + +### Fix capacity recovery in the Relaycast engine release lifecycle +- **Chose:** Fix capacity recovery in the Relaycast engine release lifecycle +- **Reasoning:** Release stops future fan-out but queued/delivered rows for the tombstoned recipient remain active and continue consuming the workspace cap until TTL. The engine owns both the lifecycle mutation and active-depth accounting, so it can atomically dead-letter those rows on irreversible delete_agent releases across every path. + +--- + +## Chapters + +### 1. Work +*Agent: default* + +- Fix capacity recovery in the Relaycast engine release lifecycle: Fix capacity recovery in the Relaycast engine release lifecycle +- Release-time settlement now covers direct deletion, local reaping, guarded node completion, and legacy node completion. Full engine suite passed (1,154 tests), and a workspace-cap regression proves expired unswept rows are already excluded from admission. diff --git a/.agentworkforce/trajectories/completed/2026-09/traj_2t44mh839vla/trajectory.json b/.agentworkforce/trajectories/completed/2026-09/traj_2t44mh839vla/trajectory.json new file mode 100644 index 00000000..925b4765 --- /dev/null +++ b/.agentworkforce/trajectories/completed/2026-09/traj_2t44mh839vla/trajectory.json @@ -0,0 +1,73 @@ +{ + "id": "traj_2t44mh839vla", + "version": 1, + "task": { + "title": "Drain workspace delivery capacity when agents are released" + }, + "status": "completed", + "startedAt": "2026-09-20T16:53:13.768Z", + "completedAt": "2026-09-20T17:00:33.145Z", + "agents": [ + { + "name": "default", + "role": "lead", + "joinedAt": "2026-09-20T16:54:06.278Z" + } + ], + "chapters": [ + { + "id": "chap_xydn9gw6mt5u", + "title": "Work", + "agentName": "default", + "startedAt": "2026-09-20T16:54:06.278Z", + "endedAt": "2026-09-20T17:00:33.145Z", + "events": [ + { + "ts": 1789923246278, + "type": "decision", + "content": "Fix capacity recovery in the Relaycast engine release lifecycle: Fix capacity recovery in the Relaycast engine release lifecycle", + "raw": { + "question": "Fix capacity recovery in the Relaycast engine release lifecycle", + "chosen": "Fix capacity recovery in the Relaycast engine release lifecycle", + "alternatives": [], + "reasoning": "Release stops future fan-out but queued/delivered rows for the tombstoned recipient remain active and continue consuming the workspace cap until TTL. The engine owns both the lifecycle mutation and active-depth accounting, so it can atomically dead-letter those rows on irreversible delete_agent releases across every path." + }, + "significance": "high" + }, + { + "ts": 1789923595490, + "type": "reflection", + "content": "Release-time settlement now covers direct deletion, local reaping, guarded node completion, and legacy node completion. Full engine suite passed (1,154 tests), and a workspace-cap regression proves expired unswept rows are already excluded from admission.", + "raw": { + "focalPoints": [ + "release atomicity", + "capacity accounting", + "regression coverage" + ], + "confidence": 0.9 + }, + "significance": "high", + "tags": [ + "focal:release atomicity", + "focal:capacity accounting", + "focal:regression coverage", + "confidence:0.9" + ] + } + ] + } + ], + "retrospective": { + "summary": "Dead-lettered active delivery rows during irreversible agent release across direct, local, guarded node, and legacy node lifecycle paths; added release-path and expired-unswept workspace-cap regressions; full engine suite, typecheck, and lint pass.", + "approach": "Standard approach", + "confidence": 0.93 + }, + "commits": [], + "filesChanged": [], + "projectId": "AgentWorkforce/relaycast", + "tags": [], + "_trace": { + "startRef": "6414a59e619459fa64640937f4c4634c207191f2", + "endRef": "6414a59e619459fa64640937f4c4634c207191f2" + } +} \ No newline at end of file diff --git a/CHANGELOG.md b/CHANGELOG.md index c3ad436a..4f811ae6 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 + +- Permanently releasing an agent now dead-letters its queued deliveries immediately, restoring workspace messaging capacity instead of waiting for delivery TTL expiry. ## [8.11.5] - 2026-09-20 diff --git a/packages/engine/CHANGELOG.md b/packages/engine/CHANGELOG.md index bf5c11e1..99aff483 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 + +- Irreversible agent release paths atomically dead-letter that recipient's active deliveries with `recipient agent released`, so tombstoned identities cannot hold workspace delivery-depth capacity until TTL. ## [8.11.5] - 2026-09-20 diff --git a/packages/engine/src/__tests__/conformance/agentLifecycle.test.ts b/packages/engine/src/__tests__/conformance/agentLifecycle.test.ts index 774e34fb..140cca92 100644 --- a/packages/engine/src/__tests__/conformance/agentLifecycle.test.ts +++ b/packages/engine/src/__tests__/conformance/agentLifecycle.test.ts @@ -694,6 +694,7 @@ describe('agent presence and release lifecycle', () => { it('reaps a hostless agent that has already spoken', async () => { const ws = await createWorkspace(stack.app, 'hostless-agent-delete-with-history'); const target = await registerAgent(stack.app, ws.workspaceKey, 'talkative-agent'); + const sender = await registerAgent(stack.app, ws.workspaceKey, 'hostless-release-sender'); const nodeId = `node_direct_${target.agentId}`; // Every agent worth reaping has history. Four FKs reference agents.id @@ -710,6 +711,15 @@ describe('agent presence and release lifecycle', () => { body: JSON.stringify({ text: 'i have said something' }), }); expect(posted.status).toBe(201); + const queued = await stack.app.request('/v1/channels/general/messages', { + method: 'POST', + headers: { + 'content-type': 'application/json', + authorization: `Bearer ${sender.token}`, + }, + body: JSON.stringify({ text: 'this delivery must be settled by local release' }), + }); + expect(queued.status).toBe(201); await stack.runtime.deps.db .update(agents) @@ -745,6 +755,17 @@ describe('agent presence and release lifecycle', () => { eq(actionInvocations.actionName, 'release'), )); expect(invocation.status).toBe('completed'); + expect( + await stack.runtime.deps.db + .select({ status: deliveries.status, error: deliveries.error }) + .from(deliveries) + .where(eq(deliveries.agentId, target.agentId)), + ).toEqual([ + expect.objectContaining({ + status: 'dead_lettered', + error: 'recipient agent released', + }), + ]); }); it('refuses to register into the reserved released-agent namespace', async () => { diff --git a/packages/engine/src/__tests__/conformance/deleteAgentRoute.test.ts b/packages/engine/src/__tests__/conformance/deleteAgentRoute.test.ts index 86240e9b..f735c155 100644 --- a/packages/engine/src/__tests__/conformance/deleteAgentRoute.test.ts +++ b/packages/engine/src/__tests__/conformance/deleteAgentRoute.test.ts @@ -1,7 +1,7 @@ import { afterEach, beforeEach, describe, expect, it } from 'vitest'; import { and, eq } from 'drizzle-orm'; import { createWorkspace, makeNodeStack, registerAgent, type TestStack } from './harness.js'; -import { agents, channelMembers, messages } from '../../db/schema.js'; +import { agents, channelMembers, deliveries, messages } from '../../db/schema.js'; /** * `DELETE /v1/agents/:name` -> `agentEngine.deleteAgent` was the last release @@ -48,7 +48,16 @@ describe('DELETE /v1/agents/:name preserves attributed history', () => { it('removes an agent that has authored messages, keeping the attribution', async () => { const ws = await createWorkspace(stack.app, 'route-delete-with-history'); const target = await registerAgent(stack.app, ws.workspaceKey, 'talkative'); + const sender = await registerAgent(stack.app, ws.workspaceKey, 'sender'); await post(target.token, 'this message must keep its author'); + await post(sender.token, 'this queued delivery must stop consuming capacity'); + + expect( + await stack.runtime.deps.db + .select({ status: deliveries.status }) + .from(deliveries) + .where(eq(deliveries.agentId, target.agentId)), + ).toContainEqual({ status: 'queued' }); const res = await removeAgent(ws.workspaceKey, target.name); expect(res.status).toBeLessThan(300); @@ -61,6 +70,27 @@ describe('DELETE /v1/agents/:name preserves attributed history', () => { .where(and(eq(agents.workspaceId, ws.workspaceId), eq(agents.name, target.name))), ).toHaveLength(0); + // A tombstoned recipient can never ACK its old queue. Settle those rows + // immediately so they stop consuming the workspace delivery-depth cap. + expect( + await stack.runtime.deps.db + .select({ + status: deliveries.status, + error: deliveries.error, + retryable: deliveries.retryable, + deadLetteredAt: deliveries.deadLetteredAt, + }) + .from(deliveries) + .where(eq(deliveries.agentId, target.agentId)), + ).toEqual([ + expect.objectContaining({ + status: 'dead_lettered', + error: 'recipient agent released', + retryable: false, + deadLetteredAt: expect.any(Date), + }), + ]); + // The row survives as a tombstone so history keeps its author. const [tombstone] = await stack.runtime.deps.db .select({ name: agents.name, status: agents.status }) diff --git a/packages/engine/src/__tests__/conformance/nodeCompletedRelease.test.ts b/packages/engine/src/__tests__/conformance/nodeCompletedRelease.test.ts index a65d8285..3eb89e9b 100644 --- a/packages/engine/src/__tests__/conformance/nodeCompletedRelease.test.ts +++ b/packages/engine/src/__tests__/conformance/nodeCompletedRelease.test.ts @@ -1,7 +1,7 @@ import { afterEach, beforeEach, describe, expect, it } from 'vitest'; import { and, eq } from 'drizzle-orm'; import { attachDirectNodeSocket, createWorkspace, makeNodeStack, registerAgent, type TestStack } from './harness.js'; -import { actionInvocations, agentNodeBindings, agents, messages, nodes } from '../../db/schema.js'; +import { actionInvocations, agentNodeBindings, agents, deliveries, messages, nodes } from '../../db/schema.js'; import { sha256Hex } from '../../lib/crypto.js'; /** @@ -61,7 +61,9 @@ describe('node-completed release preserves attributed history', () => { it('tombstones an agent that has spoken when a NODE completes the release', async () => { const ws = await createWorkspace(stack.app, 'node-release-with-history'); const target = await registerAgent(stack.app, ws.workspaceKey, 'spoke-then-released'); + const sender = await registerAgent(stack.app, ws.workspaceKey, 'release-sender'); await post(target.token, 'this message must keep its author'); + await post(sender.token, 'this delivery must be settled by guarded release'); // A LIVE node binding is what routes the release through // `applyReleaseCompletionEffect` instead of the local tombstone path. const { sock, handle } = await attachDirectNodeSocket(stack, ws.workspaceId, target); @@ -138,6 +140,18 @@ describe('node-completed release preserves attributed history', () => { }, }); + expect( + await stack.runtime.deps.db + .select({ status: deliveries.status, error: deliveries.error }) + .from(deliveries) + .where(eq(deliveries.agentId, target.agentId)), + ).toEqual([ + expect.objectContaining({ + status: 'dead_lettered', + error: 'recipient agent released', + }), + ]); + // Attribution intact, and the old credential is dead. expect( await stack.runtime.deps.db.select().from(messages).where(eq(messages.agentId, target.agentId)), @@ -159,6 +173,8 @@ describe('node-completed release preserves attributed history', () => { it('releases an agent that never spoke through the same node-completed path', async () => { const ws = await createWorkspace(stack.app, 'node-release-no-history'); const target = await registerAgent(stack.app, ws.workspaceKey, 'never-spoke'); + const sender = await registerAgent(stack.app, ws.workspaceKey, 'legacy-release-sender'); + await post(sender.token, 'this delivery must be settled by legacy release'); const { handle } = await attachDirectNodeSocket(stack, ws.workspaceId, target); const { data } = await release(ws.workspaceKey, target.name); @@ -181,6 +197,17 @@ describe('node-completed release preserves attributed history', () => { .from(agents) .where(and(eq(agents.workspaceId, ws.workspaceId), eq(agents.name, target.name))), ).toHaveLength(0); + expect( + await stack.runtime.deps.db + .select({ status: deliveries.status, error: deliveries.error }) + .from(deliveries) + .where(eq(deliveries.agentId, target.agentId)), + ).toEqual([ + expect.objectContaining({ + status: 'dead_lettered', + error: 'recipient agent released', + }), + ]); }); it('does not apply a guarded completion to a same-id takeover generation', async () => { diff --git a/packages/engine/src/engine/__tests__/workspaceDeliveryDepth.test.ts b/packages/engine/src/engine/__tests__/workspaceDeliveryDepth.test.ts index 258c7bcf..30233a87 100644 --- a/packages/engine/src/engine/__tests__/workspaceDeliveryDepth.test.ts +++ b/packages/engine/src/engine/__tests__/workspaceDeliveryDepth.test.ts @@ -143,6 +143,32 @@ describe('workspace delivery growth guard (channel broadcast)', () => { expect(await currentWorkspaceDepth(db, ws)).toBe(3); }); + it('does not charge expired unswept rows against workspace capacity', async () => { + stack = makeNodeStack(); + const db = stack.runtime.deps.db; + const { ws, channel } = await seedWorkspace(db, 1); + const policy = { cap: 1 }; + + await send(db, ws, channel, 'msg_expired', policy); + await db + .update(deliveries) + .set({ expiresAt: new Date(Date.now() - 60_000) }) + .where(eq(deliveries.workspaceId, ws)); + + // The row deliberately remains stored as `queued`; capacity admission + // must use effective (unexpired) depth rather than depend on a maintenance + // sweep having rewritten its durable status first. + expect(await db + .select({ status: deliveries.status }) + .from(deliveries) + .where(eq(deliveries.messageId, 'msg_expired'))) + .toEqual([{ status: 'queued' }]); + expect(await currentWorkspaceDepth(db, ws)).toBe(0); + + await expect(send(db, ws, channel, 'msg_after_expiry', policy)).resolves.toBeUndefined(); + expect(await activeDepth(db, ws)).toBe(1); + }); + it('resolves the dynamic host policy through the async resolver, clamped within cap', async () => { const config = { workspaceDelivery: { diff --git a/packages/engine/src/engine/action.ts b/packages/engine/src/engine/action.ts index d482a7cb..eb9bb843 100644 --- a/packages/engine/src/engine/action.ts +++ b/packages/engine/src/engine/action.ts @@ -3,7 +3,12 @@ import type { getDb } from '../db/index.js'; import { actions, actionInvocations, agents, agentNodeBindings, channelMembers, dmParticipants, nodes } from '../db/schema.js'; import { waitForPendingInvocationRetry } from './invocationRetry.js'; import { generateId } from './snowflake.js'; -import { assertRegistrableAgentName, RELEASED_AGENT_STATUS, releasedAgentName } from './agent.js'; +import { + assertRegistrableAgentName, + buildDeadLetterReleasedAgentDeliveriesWrite, + RELEASED_AGENT_STATUS, + releasedAgentName, +} from './agent.js'; import { randomHex, sha256Hex } from '../lib/crypto.js'; import { codedError } from '../lib/httpError.js'; import { D1_SAFE_IN_QUERY_CHUNK_SIZE } from '../lib/queryChunks.js'; @@ -1337,6 +1342,13 @@ async function dispatchRelease(args: { invocationCompleted, generationStillCurrent, ))); + writes.push(buildDeadLetterReleasedAgentDeliveriesWrite( + writeDb, + args.workspaceId, + agent.id, + completedAt, + and(invocationCompleted, generationStillCurrent), + )); writes.push(writeDb .delete(nodes) .where(and( @@ -2422,6 +2434,16 @@ async function completeGuardedReleaseNodeInvocation( invocationCompleted, ); if (input.delete_agent === true) { + // Settle while the release generation is still current. The tombstone + // update below intentionally rotates its token and name, after which the + // same generation predicate must no longer match. + writes.push(buildDeadLetterReleasedAgentDeliveriesWrite( + writeDb, + workspaceId, + agent.id, + completedAt, + and(invocationCompleted, generationStillCurrent), + )); writes.push(writeDb .update(agents) .set({ @@ -2595,6 +2617,12 @@ async function applyReleaseCompletionEffect( // explicitly or the released agent stays a delivery target. await db.delete(channelMembers).where(eq(channelMembers.agentId, agent.id)); await db.delete(dmParticipants).where(eq(dmParticipants.agentId, agent.id)); + await buildDeadLetterReleasedAgentDeliveriesWrite( + db, + workspaceId, + agent.id, + new Date(), + ); const implicitNodeId = `node_direct_${agent.id}`; await db.delete(nodes).where(and(eq(nodes.workspaceId, workspaceId), eq(nodes.id, implicitNodeId))); } else { diff --git a/packages/engine/src/engine/agent.ts b/packages/engine/src/engine/agent.ts index 2b7e8e75..e40190ab 100644 --- a/packages/engine/src/engine/agent.ts +++ b/packages/engine/src/engine/agent.ts @@ -1,4 +1,4 @@ -import { eq, and, gt, lt, ne, sql, inArray } from 'drizzle-orm'; +import { eq, and, gt, lt, ne, sql, inArray, type SQL } from 'drizzle-orm'; import type { getDb } from '../db/index.js'; import { agents, agentNodeBindings, agentRecoveryCredentials, channels, channelMembers, dmParticipants, actions, deliveries, nodes } from '../db/schema.js'; import { randomHex, sha256Hex } from '../lib/crypto.js'; @@ -67,6 +67,50 @@ type AgentPresenceRow = Pick; */ export const RELEASED_AGENT_STATUS = 'released'; +/** Terminal reason applied to deliveries whose recipient is permanently released. */ +export const RELEASED_AGENT_DELIVERY_ERROR = 'recipient agent released'; + +/** + * Build the terminal delivery transition that accompanies an irreversible + * agent release. + * + * Removing channel/DM membership prevents future fan-out, but it does not + * touch already queued rows. Those rows count against the workspace delivery + * cap until they are acknowledged, failed, dead-lettered, or expire. A + * tombstoned agent can never acknowledge them, so leaving them active turns a + * clean fleet teardown into a workspace-wide messaging outage until TTL. + * + * Callers include this statement in the same atomic release unit whenever + * possible. `releaseGuard` binds the transition to the caller's own CAS (for + * example a successfully completed release invocation), so a losing release + * race cannot discard deliveries owned by the surviving generation. + */ +export function buildDeadLetterReleasedAgentDeliveriesWrite( + db: Db, + workspaceId: string, + agentId: string, + releasedAt: Date, + releaseGuard?: SQL, +): AtomicWrite { + return db + .update(deliveries) + .set({ + status: 'dead_lettered', + error: RELEASED_AGENT_DELIVERY_ERROR, + retryable: false, + nextAttemptAt: null, + deadLetteredAt: releasedAt, + updatedAt: releasedAt, + }) + .where(and( + eq(deliveries.workspaceId, workspaceId), + eq(deliveries.agentId, agentId), + // Keep the active predicate literal so SQLite can use the partial index. + sql`${deliveries.status} IN ('queued', 'delivered')`, + releaseGuard, + )); +} + /** * Marker separating a released agent's original name from its tombstone * suffix. Reserved: registration rejects it, which is what makes @@ -668,6 +712,12 @@ export async function deleteAgent(db: Db, workspaceId: string, name: string) { // a tombstone, matching the release paths. writes.push(writeDb.delete(agentNodeBindings).where(eq(agentNodeBindings.agentId, agent.id))); writes.push(writeDb.delete(nodes).where(eq(nodes.id, directNodeIdForAgent(agent.id)))); + writes.push(buildDeadLetterReleasedAgentDeliveriesWrite( + writeDb, + workspaceId, + agent.id, + releasedAt, + )); return writes; }); // Capture memberships at deletion, not a preflight read that can race joins. From bd76719bc6409daca48cc90e666d41a37a7312ae Mon Sep 17 00:00:00 2001 From: agentrelaybot Date: Sun, 20 Sep 2026 10:28:11 -0700 Subject: [PATCH 2/2] fix(engine): require atomic agent release and deletion Session-Id: 01a0bfc9-024f-7503-b140-bcb525e973c0 --- .../compact_whqlsrrkhceb_2026-09-20.json | 52 ++++++ .../compact_whqlsrrkhceb_2026-09-20.md | 23 +++ CHANGELOG.md | 2 +- README.md | 6 +- openapi.yaml | 10 +- packages/engine/CHANGELOG.md | 2 +- .../conformance/nodeCompletedRelease.test.ts | 2 +- .../conformance/releaseAtomicity.test.ts | 109 ++++++++++++ packages/engine/src/engine/action.ts | 162 +----------------- packages/engine/src/engine/agent.ts | 8 +- 10 files changed, 212 insertions(+), 164 deletions(-) create mode 100644 .agentworkforce/trajectories/compacted/compact_whqlsrrkhceb_2026-09-20.json create mode 100644 .agentworkforce/trajectories/compacted/compact_whqlsrrkhceb_2026-09-20.md create mode 100644 packages/engine/src/__tests__/conformance/releaseAtomicity.test.ts diff --git a/.agentworkforce/trajectories/compacted/compact_whqlsrrkhceb_2026-09-20.json b/.agentworkforce/trajectories/compacted/compact_whqlsrrkhceb_2026-09-20.json new file mode 100644 index 00000000..c19ae6ed --- /dev/null +++ b/.agentworkforce/trajectories/compacted/compact_whqlsrrkhceb_2026-09-20.json @@ -0,0 +1,52 @@ +{ + "id": "compact_x03t9bwpqklx", + "version": 1, + "type": "compacted", + "compactedAt": "2026-09-20T17:28:10.496Z", + "sourceTrajectories": [ + "traj_pcfd4ebl1pam" + ], + "dateRange": { + "start": "2026-09-20T17:13:27.676Z", + "end": "2026-09-20T17:27:21.837Z" + }, + "summary": { + "totalDecisions": 2, + "totalEvents": 4, + "uniqueAgents": [ + "default" + ] + }, + "decisionGroups": [ + { + "category": "other", + "decisions": [ + { + "question": "Require atomic writes for every irreversible release and share node completion implementation", + "chosen": "Require atomic writes for every irreversible release and share node completion implementation", + "reasoning": "Sequential fallback can permanently settle deliveries before a later tombstone failure; legacy node completion also committed lifecycle writes separately even on capable adapters.", + "fromTrajectory": "traj_pcfd4ebl1pam" + } + ] + }, + { + "category": "api", + "decisions": [ + { + "question": "Address all three PR threads", + "chosen": "Address all three PR threads", + "reasoning": "Devin requires fail-closed local releases; CodeRabbit independently identifies sequential legacy node completion and requests naming both active delivery states in release notes.", + "fromTrajectory": "traj_pcfd4ebl1pam" + } + ] + } + ], + "keyLearnings": [], + "keyFindings": [ + "Engine typecheck and lint passed under Node 22.23.2.", + "Focused release and capacity tests: 101 passed, including 21 new refusal/rollback regressions.", + "Full engine suite: 1176 tests across 100 files passed with --maxWorkers=2. Initial default-parallel run had one unrelated retention CLI timeout; no timeout or dependency changes were committed." + ], + "filesAffected": [], + "commits": [] +} diff --git a/.agentworkforce/trajectories/compacted/compact_whqlsrrkhceb_2026-09-20.md b/.agentworkforce/trajectories/compacted/compact_whqlsrrkhceb_2026-09-20.md new file mode 100644 index 00000000..eed38bb9 --- /dev/null +++ b/.agentworkforce/trajectories/compacted/compact_whqlsrrkhceb_2026-09-20.md @@ -0,0 +1,23 @@ +# Trajectory Compaction: Sep 20, 2026 - Sep 20, 2026 + +## Summary +- Sessions: 1 +- Decisions: 2 +- Events: 4 +- Agents: default +- Files: 0 +- Commits: 0 + +## Other +- Require atomic writes for every irreversible release and share node completion implementation -> Require atomic writes for every irreversible release and share node completion implementation (traj_pcfd4ebl1pam) + +## Api +- Address all three PR threads -> Address all three PR threads (traj_pcfd4ebl1pam) + +## Key Learnings +- None + +## Key Findings +- Engine typecheck and lint passed under Node 22.23.2. +- Focused release and capacity tests: 101 passed, including 21 new refusal/rollback regressions. +- Full engine suite: 1176 tests across 100 files passed with `--maxWorkers=2`. Initial default-parallel run had one unrelated retention CLI timeout; no timeout or dependency changes were committed. diff --git a/CHANGELOG.md b/CHANGELOG.md index 4f811ae6..e73e7b0e 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -20,7 +20,7 @@ Packages without a separate changelog are covered by the cross-package notes bel ### Fixed -- Permanently releasing an agent now dead-letters its queued deliveries immediately, restoring workspace messaging capacity instead of waiting for delivery TTL expiry. +- Permanently releasing an agent now dead-letters its active `queued` and `delivered` deliveries immediately, restoring workspace messaging capacity instead of waiting for delivery TTL expiry. ## [8.11.5] - 2026-09-20 diff --git a/README.md b/README.md index f9c26e89..d3fb92ec 100644 --- a/README.md +++ b/README.md @@ -701,7 +701,11 @@ live host when one exists. If the host is absent or offline, a normal release fails explicitly with `agent_host_unavailable` instead of creating an ownerless pending invocation. A `delete_agent` request can be completed locally in that case: Relaycast deactivates bindings, frees the live name, and removes its -implicit direct node. Cleanup callers that retain the issued agent token can +implicit direct node. Irreversible releases and agent deletion require a database +transaction or atomic batch; adapters without either capability are refused +before identity, membership, or queued deliveries change. Successful cleanup +dead-letters the released agent's active deliveries in the same atomic write. +Cleanup callers that retain the issued agent token can send its SHA-256 hash as `expected_token_hash`; Relaycast then rejects a stale release with `agent_release_generation_conflict` before dispatch or completion, so a same-name takeover is left untouched. diff --git a/openapi.yaml b/openapi.yaml index ce2edd21..907df993 100644 --- a/openapi.yaml +++ b/openapi.yaml @@ -2584,7 +2584,10 @@ paths: delete: summary: Delete agent - description: Delete an agent from the workspace + description: >- + Tombstone an agent, remove memberships, and dead-letter its active deliveries + in one atomic write. Database adapters without transaction or atomic batch + support are refused before these changes. tags: - Agents security: @@ -2926,7 +2929,10 @@ paths: release fails explicitly with `503 agent_host_unavailable`; it never creates an ownerless pending invocation. With `delete_agent`, the engine can reap the database record directly and returns a completed invocation, - deleting the agent and any implicit direct node. + tombstoning the agent and deleting any implicit direct node. Irreversible + release removes memberships and dead-letters active deliveries in the + same atomic write. Database adapters without transaction or atomic batch + support are refused before these changes. tags: - Agents security: diff --git a/packages/engine/CHANGELOG.md b/packages/engine/CHANGELOG.md index 99aff483..3c5a3540 100644 --- a/packages/engine/CHANGELOG.md +++ b/packages/engine/CHANGELOG.md @@ -11,7 +11,7 @@ and this project follows [Semantic Versioning](https://semver.org/spec/v2.0.0.ht ### Fixed -- Irreversible agent release paths atomically dead-letter that recipient's active deliveries with `recipient agent released`, so tombstoned identities cannot hold workspace delivery-depth capacity until TTL. +- Irreversible agent release paths atomically dead-letter that recipient's active deliveries with `recipient agent released`, so tombstoned identities cannot hold workspace delivery-depth capacity until TTL. Release and deletion refuse adapters without atomic writes before changing identity, membership, or deliveries. ## [8.11.5] - 2026-09-20 diff --git a/packages/engine/src/__tests__/conformance/nodeCompletedRelease.test.ts b/packages/engine/src/__tests__/conformance/nodeCompletedRelease.test.ts index 3eb89e9b..7c0efb8b 100644 --- a/packages/engine/src/__tests__/conformance/nodeCompletedRelease.test.ts +++ b/packages/engine/src/__tests__/conformance/nodeCompletedRelease.test.ts @@ -65,7 +65,7 @@ describe('node-completed release preserves attributed history', () => { await post(target.token, 'this message must keep its author'); await post(sender.token, 'this delivery must be settled by guarded release'); // A LIVE node binding is what routes the release through - // `applyReleaseCompletionEffect` instead of the local tombstone path. + // `completeReleaseNodeInvocation` instead of the local tombstone path. const { sock, handle } = await attachDirectNodeSocket(stack, ws.workspaceId, target); const { data } = await release( diff --git a/packages/engine/src/__tests__/conformance/releaseAtomicity.test.ts b/packages/engine/src/__tests__/conformance/releaseAtomicity.test.ts new file mode 100644 index 00000000..7c7d79d9 --- /dev/null +++ b/packages/engine/src/__tests__/conformance/releaseAtomicity.test.ts @@ -0,0 +1,109 @@ +import { afterEach, beforeEach, describe, expect, it } from 'vitest'; +import { eq } from 'drizzle-orm'; +import { + attachDirectNodeSocket, attachFakeBatch, createWorkspace, makeNodeStack, + registerAgent, stripTransactionCapability, type TestStack, +} from './harness.js'; +import { + actionInvocations, agentNodeBindings, agents, channelMembers, deliveries, + dmParticipants, nodes, +} from '../../db/schema.js'; +import { completeNodeInvocation, dispatchAgentRelease } from '../../engine/action.js'; +import { deleteAgent } from '../../engine/agent.js'; +import { sha256Hex } from '../../lib/crypto.js'; + +const paths = ['delete', 'local', 'local-token', 'local-exact', 'node', 'node-token', 'node-exact'] as const; + +describe('irreversible release atomicity', () => { + let stack: TestStack; + beforeEach(() => { stack = makeNodeStack(); }); + afterEach(() => stack.close()); + + for (const path of paths) { + for (const adapter of ['sequential', 'transaction', 'batch'] as const) { + it(`${path} preserves identity, membership, and deliveries on ${adapter} failure`, async () => { + const ws = await createWorkspace(stack.app, `release-${path}-${adapter}`); + const target = await registerAgent(stack.app, ws.workspaceKey, 'target'); + const sender = await registerAgent(stack.app, ws.workspaceKey, 'sender'); + const db = stack.runtime.deps.db; + const input = { + name: target.name, + delete_agent: true, + ...(path.endsWith('-token') ? { expected_token_hash: await sha256Hex(target.token) } : {}), + ...(path.endsWith('-exact') ? { expected_agent_id: target.agentId } : {}), + }; + let complete: (() => Promise) | undefined; + let invocationId: string | undefined; + if (path.startsWith('node')) { + const { nodeId } = await attachDirectNodeSocket(stack, ws.workspaceId, target); + const ack = await dispatchAgentRelease(db, ws.workspaceId, { input }, { + nodeConnections: stack.runtime.deps.nodeConnections, + }); + expect(ack.status).toBe('dispatched'); + invocationId = ack.invocation_id; + const [invocation] = await db.select().from(actionInvocations) + .where(eq(actionInvocations.id, invocationId)); + complete = () => completeNodeInvocation( + db, stack.runtime.deps.nodeConnections, ws.workspaceId, nodeId, + invocation.dispatchedProvider!, invocation.id, { output: { released: true } }, + ); + } + for (const text of ['must remain replayable', 'must remain acknowledgeable']) { + const dm = await stack.app.request('/v1/dm', { + method: 'POST', + headers: { 'content-type': 'application/json', authorization: `Bearer ${sender.token}` }, + body: JSON.stringify({ to: target.name, text }), + }); + expect(dm.status).toBe(201); + } + await stack.settle(); + // Exercise both queued and delivered active rows without socket timing + // changing the snapshot while the release is under test. + await db.update(deliveries).set({ status: 'queued' }) + .where(eq(deliveries.agentId, target.agentId)); + const [queued] = await db.select().from(deliveries) + .where(eq(deliveries.agentId, target.agentId)); + expect(queued).toBeDefined(); + await db.update(deliveries).set({ status: 'delivered' }) + .where(eq(deliveries.id, queued.id)); + + const snapshot = async () => ({ + agent: await db.select().from(agents).where(eq(agents.id, target.agentId)), + channels: await db.select().from(channelMembers).where(eq(channelMembers.agentId, target.agentId)), + dms: await db.select().from(dmParticipants).where(eq(dmParticipants.agentId, target.agentId)), + bindings: await db.select().from(agentNodeBindings).where(eq(agentNodeBindings.agentId, target.agentId)), + nodes: await db.select().from(nodes).where(eq(nodes.workspaceId, ws.workspaceId)), + deliveries: await db.select().from(deliveries).where(eq(deliveries.agentId, target.agentId)), + }); + const before = await snapshot(); + expect(before.channels.length).toBeGreaterThan(0); + expect(before.dms.length).toBeGreaterThan(0); + expect(before.deliveries.map(row => row.status).sort()).toEqual(['delivered', 'queued']); + if (adapter === 'sequential') stripTransactionCapability(db); + if (adapter === 'batch') attachFakeBatch(stack, db); + // A late tombstone failure reproduces the review's lost-message case. + // Atomic adapters must roll back; sequential adapters must refuse first. + stack.runtime.handle.sqlite.exec(` + CREATE TRIGGER refuse_release_tombstone BEFORE UPDATE ON agents + WHEN NEW.status = 'released' + BEGIN SELECT RAISE(ABORT, 'forced tombstone failure'); END + `); + const release = complete ?? (() => path === 'delete' + ? deleteAgent(db, ws.workspaceId, target.name) + : dispatchAgentRelease(db, ws.workspaceId, { input })); + await expect(release()).rejects.toThrow(adapter === 'sequential' + ? /Atomic write capability required/ : /forced tombstone failure/); + expect(await snapshot()).toEqual(before); + if (invocationId) { + const [invocation] = await db.select().from(actionInvocations) + .where(eq(actionInvocations.id, invocationId)); + expect(invocation.status).toBe('dispatched'); + } else if (path !== 'delete') { + const [invocation] = await db.select().from(actionInvocations) + .where(eq(actionInvocations.workspaceId, ws.workspaceId)); + expect(invocation.status).toBe('pending'); + } + }); + } + } +}); diff --git a/packages/engine/src/engine/action.ts b/packages/engine/src/engine/action.ts index eb9bb843..cbe0e6df 100644 --- a/packages/engine/src/engine/action.ts +++ b/packages/engine/src/engine/action.ts @@ -1387,7 +1387,7 @@ async function dispatchRelease(args: { ))); return writes; - }, expectedTokenHash || expectedAgentId ? { requireAtomic: true } : undefined); + }, { requireAtomic: true }); const generationConflict = results[0] as Array<{ id: string }>; const completed = results[1] as Array<{ id: string; handlerNodeId: string | null }>; const completedExitNodeId = completed[0]?.handlerNodeId ?? null; @@ -2239,7 +2239,7 @@ function publicInvocation(row: InvocationRow) { }; } -async function completeGuardedReleaseNodeInvocation( +async function completeReleaseNodeInvocation( db: Db, workspaceId: string, nodeId: string, @@ -2256,7 +2256,7 @@ async function completeGuardedReleaseNodeInvocation( ? persistedTokenHash : null; const expectedAgentId = releaseExpectedAgentId(input); - if (!name || (!expectedTokenHash && !expectedAgentId)) { + if (!name || (persistedTokenHash !== undefined && !expectedTokenHash && !expectedAgentId)) { const [failed] = await db .update(actionInvocations) .set({ @@ -2342,10 +2342,9 @@ async function completeGuardedReleaseNodeInvocation( const results = await runAtomicWrites(db, (writeDb) => { const writes: AtomicWrite[] = []; - // The conflict settlement comes first. In a sequential SQLite batch it - // closes the invocation before every mutation below if either the token - // generation or node binding changed; on transactional adapters the same - // statement ordering and rollback guarantee apply. + // Settle conflicts before lifecycle mutations. Every statement belongs + // to one required transaction or atomic batch, including legacy releases + // without an explicit identity or token guard. writes.push(writeDb .update(actionInvocations) .set({ @@ -2524,142 +2523,6 @@ async function completeGuardedReleaseNodeInvocation( return publicInvocation(settled); } -async function applyReleaseCompletionEffect( - db: Db, - workspaceId: string, - nodeId: string | null, - invocation: Pick, - data: { error?: string }, - deps?: InvocationCompletionDeps, - options: { allowMissingBinding?: boolean; expectedAgentId?: string } = {}, -): Promise { - if (!isBuiltinReleaseInvocation(invocation) || data.error) return false; - - const input = recordInput(invocation.input); - const name = typeof input.name === 'string' ? input.name : null; - if (!name) return false; - - const [agent] = await db - .select() - .from(agents) - .where(and( - eq(agents.workspaceId, workspaceId), - eq(agents.name, name), - ...(options.expectedAgentId ? [eq(agents.id, options.expectedAgentId)] : []), - ...(!options.allowMissingBinding && nodeId ? [ - eq(agents.locationType, 'via_node'), - eq(agents.locationNodeId, nodeId), - ] : []), - )); - if (!agent) return false; - - // Only proceed if an active binding actually flipped to inactive. This guards - // against a second release (e.g. a retry) double-decrementing activeAgents for - // an agent that was already released from this node. - const deactivatedBindings = await db - .update(agentNodeBindings) - .set({ status: 'inactive', updatedAt: new Date() }) - .where(and( - eq(agentNodeBindings.workspaceId, workspaceId), - eq(agentNodeBindings.agentId, agent.id), - eq(agentNodeBindings.status, 'active'), - ...(nodeId ? [eq(agentNodeBindings.nodeId, nodeId)] : []), - )) - .returning({ nodeId: agentNodeBindings.nodeId }); - if (deactivatedBindings.length === 0 && !options.allowMissingBinding) return false; - - // Capture exit correlation BEFORE the mutation deletes the row or strips the - // spawn/cli metadata, so a durable agent.exited can still be emitted. - const exited = { agentId: agent.id, agentName: agent.name, invocationId: fleetInvocationId(agent.metadata) }; - - const deactivatedNodeIds = Array.from(new Set(deactivatedBindings.map((binding) => binding.nodeId))); - if (deactivatedNodeIds.length > 0) { - await db - .update(nodes) - .set({ - activeAgents: sql`CASE WHEN ${nodes.activeAgents} > 0 THEN ${nodes.activeAgents} - 1 ELSE 0 END`, - }) - .where(and(eq(nodes.workspaceId, workspaceId), inArray(nodes.id, deactivatedNodeIds))); - } - - if (input.delete_agent === true) { - // Tombstone rather than DELETE, matching `dispatchRelease`'s - // `completeLocally`. Four FKs reference `agents.id` without an ON DELETE - // action (`messages.agent_id`, `channels.created_by`, `files.uploaded_by`, - // `webhooks.created_by`), so a bare delete is refused for any agent that - // has ever spoken — and this runs inside the completion's atomic unit, so - // that refusal aborts the invocation completion too. The seat and the name - // then stay claimed forever and the caller only ever sees `dispatched`. - // Renaming frees the unique `(workspace_id, name)` immediately while every - // FK target stays valid and every message keeps its sender. - const releasedName = releasedAgentName(agent.name, agent.id); - // The row survives, so its credential must not. `token_hash` is NOT NULL - // UNIQUE and cannot be cleared, so rotate it to a value nobody holds. - const releasedTokenHash = await sha256Hex(`released:${agent.id}:${randomHex(16)}`); - await db - .update(agents) - .set({ - name: releasedName, - handle: `@${releasedName}`, - status: RELEASED_AGENT_STATUS, - tokenHash: releasedTokenHash, - // Same reason as the other release paths — the grace slot survives - // `token_hash` rewrites unless we clear it. See 0035_agent_token_grace. - previousTokenHash: null, - previousTokenExpiresAt: null, - locationType: 'self_connected', - locationNodeId: null, - lastSeen: new Date(), - }) - .where(and(eq(agents.workspaceId, workspaceId), eq(agents.id, agent.id))); - // `channel_members` and `dm_participants` reference `agents.id` ON DELETE - // CASCADE; an UPDATE does not fire that cascade, so drop the memberships - // explicitly or the released agent stays a delivery target. - await db.delete(channelMembers).where(eq(channelMembers.agentId, agent.id)); - await db.delete(dmParticipants).where(eq(dmParticipants.agentId, agent.id)); - await buildDeadLetterReleasedAgentDeliveriesWrite( - db, - workspaceId, - agent.id, - new Date(), - ); - const implicitNodeId = `node_direct_${agent.id}`; - await db.delete(nodes).where(and(eq(nodes.workspaceId, workspaceId), eq(nodes.id, implicitNodeId))); - } else { - const existingMetadata = agent.metadata ?? {}; - const { spawn: _spawn, cli: _cli, ...restMetadata } = existingMetadata; - await db - .update(agents) - .set({ - status: 'offline', - // Clear the node location so the agent is no longer routable to the released - // node and a repeat release can't re-decrement the node's active count. - locationType: 'self_connected', - locationNodeId: null, - lastSeen: new Date(), - metadata: { - ...restMetadata, - release: { - reason: typeof input.reason === 'string' ? input.reason : null, - released_at: new Date().toISOString(), - }, - }, - }) - .where(and(eq(agents.workspaceId, workspaceId), eq(agents.id, agent.id))); - } - - if (deps && nodeId) { - await emitAgentExitedEffects(deps, workspaceId, { - agentId: exited.agentId, - agentName: exited.agentName, - nodeId, - invocationId: exited.invocationId, - reason: 'released', - }); - } - return true; -} - async function dispatchNodeAttempt( db: Db, workspaceId: string, @@ -3835,13 +3698,8 @@ export async function completeNodeInvocation( } } - const existingInput = recordInput(existing.input); - if ( - !data.error - && isBuiltinReleaseInvocation(existing) - && (existingInput.expected_token_hash !== undefined || existingInput.expected_agent_id !== undefined) - ) { - return completeGuardedReleaseNodeInvocation( + if (!data.error && isBuiltinReleaseInvocation(existing)) { + return completeReleaseNodeInvocation( db, workspaceId, nodeId, @@ -3877,10 +3735,6 @@ export async function completeNodeInvocation( await releaseNodeCapacity(db, workspaceId, updated.dispatchedNodeId); } - if (updated) { - await applyReleaseCompletionEffect(db, workspaceId, nodeId, existing, data, deps); - } - return updated ? publicInvocation(updated) : null; } diff --git a/packages/engine/src/engine/agent.ts b/packages/engine/src/engine/agent.ts index e40190ab..b3007d8b 100644 --- a/packages/engine/src/engine/agent.ts +++ b/packages/engine/src/engine/agent.ts @@ -80,8 +80,8 @@ export const RELEASED_AGENT_DELIVERY_ERROR = 'recipient agent released'; * tombstoned agent can never acknowledge them, so leaving them active turns a * clean fleet teardown into a workspace-wide messaging outage until TTL. * - * Callers include this statement in the same atomic release unit whenever - * possible. `releaseGuard` binds the transition to the caller's own CAS (for + * Callers include this statement in the same required atomic release unit. + * `releaseGuard` binds the transition to the caller's own CAS (for * example a successfully completed release invocation), so a losing release * race cannot discard deliveries owned by the surviving generation. */ @@ -660,7 +660,7 @@ export async function deleteAgent(db: Db, workspaceId: string, name: string) { if (!agent) return false; // Tombstone rather than DELETE, matching both release paths - // (`dispatchRelease` -> `completeLocally`, and `applyReleaseCompletionEffect`). + // (`dispatchRelease` -> `completeLocally`, and `completeReleaseNodeInvocation`). // Four FKs reference `agents.id` with no ON DELETE action — // `messages.agent_id`, `channels.created_by`, `files.uploaded_by`, // `webhooks.created_by` — so a bare DELETE is refused for any agent that has @@ -719,7 +719,7 @@ export async function deleteAgent(db: Db, workspaceId: string, name: string) { releasedAt, )); return writes; - }); + }, { requireAtomic: true }); // Capture memberships at deletion, not a preflight read that can race joins. const removed = releaseResults[1] as Array<{ channelId: string }>; const joinedChannels = await queryInChunks(removed.map(row => row.channelId), ids => db