diff --git a/CHANGELOG.md b/CHANGELOG.md index eb6314aa..bc008259 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -21,6 +21,8 @@ Packages without a separate changelog are covered by the cross-package notes bel ### Fixed - `GET /v1/nodes` no longer fetches a workspace's entire node history to serve a default listing: `name`, `capability`, and a new `status` liveness selector are now pushed into SQL, and `status=online` returns only fresh-heartbeat live nodes. An explicit `history=true` mode adds bounded, non-truncating cursor pagination for reading full history (e.g. `--all`). `active_agents` is now flagged with `active_agents_stale` once a node is offline, so a frozen historical count is never presented as current occupancy. (Fixes [#422](https://github.com/AgentWorkforce/relaycast/issues/422)) +- Keyed agent session events now survive lost responses without duplicate events, while conflicting key reuse returns a typed error. A `status.*` event's status mutation and durable completion marker are applied atomically so retries cannot leave a stale agent row behind a successful response. +- Migration `0056_session_event_status_completion.sql` reconciles pre-existing keyed status events without fabricating completion, recovering only rows proven older than the event and conservatively preserving rows touched by later status or liveness writers. ## [8.8.0] - 2026-09-10 diff --git a/README.md b/README.md index 10e01db7..af30b487 100644 --- a/README.md +++ b/README.md @@ -976,7 +976,7 @@ POST /v1/actions Register an action (agent-to-agent RPC) POST /v1/actions/:name/invoke Invoke an action (workspace-scoped / global alias) POST /v1/nodes/:node/actions/:name/invoke Invoke a node-addressed action DELETE /v1/nodes/:node/providers/:name Remove a node provider -POST /v1/agents/:name/events Emit an agent session event +POST /v1/agents/:name/events Emit an agent session event (optional Idempotency-Key replays identical retries; status.changed requires payload.status) POST /v1/directory/agents Publish an agent to the directory GET /v1/directory/search Search the agent directory POST /v1/route Skill-based agent routing diff --git a/openapi.yaml b/openapi.yaml index 80228705..1fe219ce 100644 --- a/openapi.yaml +++ b/openapi.yaml @@ -5531,8 +5531,10 @@ paths: post: summary: Emit agent session event description: > - Record a session event for an agent (e.g. `status.active`, `status.idle`). An agent token - may only post events for itself; a workspace key may post for any agent. + Record a session event for an agent (e.g. `status.active`, `status.idle`). For + `status.changed`, `payload.status` is required and accepts `active`, `idle`, `blocked`, + `waiting`, `offline`, or the legacy alias `online`. An agent token may only post events + for itself; a workspace key may post for any agent. tags: - Agents security: @@ -5544,6 +5546,16 @@ paths: required: true schema: type: string + - name: Idempotency-Key + in: header + required: false + description: >- + Stable caller-generated key for retries. Within one workspace and agent, an identical + request replays the original event (with `Idempotency-Replayed: true`) instead of + appending another event. Reusing a key with a different payload returns 409. + schema: + type: string + maxLength: 255 requestBody: required: true content: @@ -5559,11 +5571,38 @@ paths: type: object responses: '201': - description: Event recorded + description: Event recorded, or idempotently replayed for a repeated `Idempotency-Key` + headers: + Idempotency-Replayed: + description: Set to `true` when the original event is replayed. + schema: + type: string + enum: ['true'] content: application/json: schema: $ref: '#/components/schemas/SuccessResponse' + '400': + description: Invalid Idempotency-Key + content: + application/json: + schema: + $ref: '#/components/schemas/ErrorResponse' + '409': + description: Idempotency-Key was reused with a different payload + content: + application/json: + schema: + $ref: '#/components/schemas/ErrorResponse' + '503': + description: > + `idempotency_unavailable` — the durable event identity claim could not be + confirmed after a storage conflict. Retry the same logical request with the + same Idempotency-Key; do not mint a replacement key. + content: + application/json: + schema: + $ref: '#/components/schemas/ErrorResponse' get: summary: List agent session events description: List recorded session events for an agent. Observer tokens require `activity:read`; `agent_ids` and `created_after` filters apply. diff --git a/packages/engine/CHANGELOG.md b/packages/engine/CHANGELOG.md index dff8145e..3718bf45 100644 --- a/packages/engine/CHANGELOG.md +++ b/packages/engine/CHANGELOG.md @@ -12,6 +12,8 @@ and this project follows [Semantic Versioning](https://semver.org/spec/v2.0.0.ht ### Fixed - `listNodes`/`GET /v1/nodes` push `name`, `capability`, and a new `status` (`online`/`offline`) liveness selector into SQL instead of fetching the whole workspace roster and filtering in JS. `history=true` adds a bounded, non-truncating cursor pagination contract (`{ nodes, next_cursor }`, paged via `cursor`/`limit`, capped at 500/page) for explicit full-history reads. Without `history`, the response stays the legacy bare array. Roster entries add `active_agents_stale` (true once a node is offline) so `active_agents` is never read back as authoritative current occupancy. Added `idx_nodes_status_heartbeat` to keep the live-selection query indexed. An observer token's authorized `active_agents` count is computed with bounded, JSON-array-bound SQL joins instead of one `inArray`/`IN (...)` bind per visible node id, keeping every roster query under D1's 100-parameter limit regardless of page size or history/legacy path. +- `POST /v1/agents/:name/events` durably replays identical `Idempotency-Key` retries and rejects conflicting payload reuse. Status mutations and their completion markers are applied atomically, so an interrupted status retry can finish without returning success against a stale agent row. +- Migration `0056_session_event_status_completion.sql` reconciles pre-existing keyed status events without fabricating completion, recovering only rows proven older than the event and conservatively preserving rows touched by later status or liveness writers. ## [8.8.0] - 2026-09-10 diff --git a/packages/engine/src/__tests__/atomicity.test.ts b/packages/engine/src/__tests__/atomicity.test.ts index eaa14ff8..de63b930 100644 --- a/packages/engine/src/__tests__/atomicity.test.ts +++ b/packages/engine/src/__tests__/atomicity.test.ts @@ -1,4 +1,5 @@ import { afterEach, beforeEach, describe, expect, it } from 'vitest'; +import { readFileSync } from 'node:fs'; import { and, eq } from 'drizzle-orm'; import { makeNodeStack, @@ -17,6 +18,9 @@ import { messageLogs, messages, readReceipts, + pendingEvents, + sessionEvents, + workspaceEvents, } from '../db/schema.js'; import { postMessage } from '../engine/message.js'; import { sendDm } from '../engine/dm.js'; @@ -24,6 +28,10 @@ import { createGroupDm, postGroupMessage } from '../engine/groupDm.js'; import { postReply } from '../engine/thread.js'; import { markRead } from '../engine/receipt.js'; import { rotateAgentIdentity } from '../engine/agentIdentity.js'; +import { + applyStatusEventEffect, + recordSessionEventWithIdempotency, +} from '../engine/sessionEvent.js'; import type { AtomicWrite, EngineDb, TransactionCapability } from '../ports/database.js'; /** @@ -520,4 +528,486 @@ describe('atomic write paths', () => { expect(await db.select().from(deliveries)).toHaveLength(0); }); }); + + /** + * relaycast#425: the keyed status.* event mutation must be one atomic unit + * with its completion marker, or a crash between the durable event claim + * and the agent-row status write leaves a replay permanently unable to + * tell "never applied" apart from "applied", stranding a stale agent row + * behind a 201 response forever. + */ + describe('session event status effect (relaycast#425)', () => { + async function seedAgent() { + const ws = await createWorkspace(stack.app, 'status-atomicity-ws'); + const alice = await registerAgent(stack.app, ws.workspaceKey, 'alice'); + const db = stack.runtime.handle.db as unknown as EngineDb; + return { ws, alice, db }; + } + + function markLegacyStatusEvents(db: EngineDb) { + // The test schema already has 0056's additive columns and trigger. Run + // the migration's real backfill/marker statements while omitting only + // DDL that would otherwise try to add the same objects twice. + const migration = readFileSync( + new URL('../db/migrations/0056_session_event_status_completion.sql', import.meta.url), + 'utf8', + ) + .replace(/^ALTER TABLE (?:agents|session_events) ADD COLUMN .*;\n/gm, '') + .replace(/CREATE TRIGGER agents_status_reconciliation_timestamp[\s\S]*?END;\n\n/, ''); + stack.runtime.handle.sqlite.exec(migration); + } + + it('rolls back the status write and completion marker together when the agent update fails', async () => { + const { ws, alice, db } = await seedAgent(); + const { event } = await recordSessionEventWithIdempotency( + db, + ws.workspaceId, + alice.agentId, + { type: 'status.blocked', payload: {} }, + 'status-failure-1', + ); + + const restore = injectUpdateFailure(db, agents, 'injected agent status failure'); + await expect( + applyStatusEventEffect(db, ws.workspaceId, alice.agentId, event.id, 'blocked'), + ).rejects.toThrow('injected agent status failure'); + restore(); + + const [agentRow] = await db.select({ status: agents.status }).from(agents).where(eq(agents.id, alice.agentId)); + expect(agentRow!.status).not.toBe('blocked'); + const [eventRow] = await db.select({ statusAppliedAt: sessionEvents.statusAppliedAt }).from(sessionEvents).where(eq(sessionEvents.id, event.id)); + expect(eventRow!.statusAppliedAt).toBeNull(); + }); + + it('replays a keyed status event whose mutation never completed and finishes it exactly once', async () => { + const { ws, alice, db } = await seedAgent(); + const idempotencyKey = 'status-crash-replay-1'; + + // First attempt: the durable event claim commits, but the process + // crashes before the agent status mutation runs — modeled directly + // since the route always calls both in sequence. + const first = await recordSessionEventWithIdempotency( + db, ws.workspaceId, alice.agentId, { type: 'status.blocked', payload: {} }, idempotencyKey, + ); + expect(first.replayed).toBe(false); + expect(first.pendingStatusApplication).toBe(true); + + // Retry after the "crash": the event is replayed, but the completion + // marker is still NULL, so the interrupted mutation must be redone. + const replay = await recordSessionEventWithIdempotency( + db, ws.workspaceId, alice.agentId, { type: 'status.blocked', payload: {} }, idempotencyKey, + ); + expect(replay.replayed).toBe(true); + expect(replay.pendingStatusApplication).toBe(true); + expect(replay.event.id).toBe(first.event.id); + + await applyStatusEventEffect(db, ws.workspaceId, alice.agentId, replay.event.id, 'blocked'); + + const [agentRow] = await db.select({ status: agents.status }).from(agents).where(eq(agents.id, alice.agentId)); + expect(agentRow!.status).toBe('blocked'); + + // A further replay now sees the completion marker set and must not + // report a pending mutation again. + const secondReplay = await recordSessionEventWithIdempotency( + db, ws.workspaceId, alice.agentId, { type: 'status.blocked', payload: {} }, idempotencyKey, + ); + expect(secondReplay.replayed).toBe(true); + expect(secondReplay.pendingStatusApplication).toBe(false); + + expect(await db.select().from(sessionEvents)).toHaveLength(1); + }); + + it('reconciles an already-applied pre-marker status without refanout', async () => { + const { ws, alice, db } = await seedAgent(); + const key = 'status-legacy-backfill-1'; + const first = await recordSessionEventWithIdempotency( + db, ws.workspaceId, alice.agentId, { type: 'status.blocked', payload: {} }, key, + ); + + // Model the historical route: the status write completed before 0056, + // but no completion marker existed yet. Apply the actual migration's + // UPDATE against this Node database while leaving its ALTER statements + // out because the current test schema already has both columns. + await db.update(agents).set({ status: 'blocked' }).where(eq(agents.id, alice.agentId)); + await db.update(agents).set({ + statusUpdatedAt: new Date(new Date(first.event.created_at).getTime() + 1_000), + }).where(eq(agents.id, alice.agentId)); + markLegacyStatusEvents(db); + + const replay = await recordSessionEventWithIdempotency( + db, ws.workspaceId, alice.agentId, { type: 'status.blocked', payload: {} }, key, + ); + expect(replay.event.id).toBe(first.event.id); + expect(replay.pendingStatusApplication).toBe(true); + + const batches = attachFakeBatch(db); + await expect( + applyStatusEventEffect(db, ws.workspaceId, alice.agentId, replay.event.id, 'blocked'), + ).resolves.toEqual({ claimed: true, mutated: false }); + expect(batches).toHaveLength(1); + const [agentRow] = await db.select({ status: agents.status }).from(agents).where(eq(agents.id, alice.agentId)); + expect(agentRow!.status).toBe('blocked'); + const [eventRow] = await db.select({ statusAppliedAt: sessionEvents.statusAppliedAt }) + .from(sessionEvents).where(eq(sessionEvents.id, first.event.id)); + expect(eventRow!.statusAppliedAt).not.toBeNull(); + }); + + it('replays an interrupted pre-marker status when current state differs', async () => { + const { ws, alice, db } = await seedAgent(); + const first = await recordSessionEventWithIdempotency( + db, ws.workspaceId, alice.agentId, { type: 'status.blocked', payload: {} }, 'status-legacy-pending-1', + ); + await db.update(agents).set({ + statusUpdatedAt: new Date(new Date(first.event.created_at).getTime() - 1_000), + }).where(eq(agents.id, alice.agentId)); + markLegacyStatusEvents(db); + + const replay = await recordSessionEventWithIdempotency( + db, ws.workspaceId, alice.agentId, { type: 'status.blocked', payload: {} }, 'status-legacy-pending-1', + ); + expect(replay.pendingStatusApplication).toBe(true); + const batches = attachFakeBatch(db); + await expect( + applyStatusEventEffect(db, ws.workspaceId, alice.agentId, replay.event.id, 'blocked'), + ).resolves.toEqual({ claimed: true, mutated: true }); + expect(batches).toHaveLength(1); + const [agentRow] = await db.select({ status: agents.status }).from(agents).where(eq(agents.id, alice.agentId)); + expect(agentRow!.status).toBe('blocked'); + const [eventRow] = await db.select({ statusAppliedAt: sessionEvents.statusAppliedAt }) + .from(sessionEvents).where(eq(sessionEvents.id, first.event.id)); + expect(eventRow!.statusAppliedAt).not.toBeNull(); + }); + + it('does not clobber a later non-event status writer during legacy replay', async () => { + const { ws, alice, db } = await seedAgent(); + const key = 'status-legacy-later-writer-1'; + const first = await recordSessionEventWithIdempotency( + db, ws.workspaceId, alice.agentId, { type: 'status.blocked', payload: {} }, key, + ); + const eventAt = new Date(first.event.created_at); + await db.update(agents).set({ + status: 'waiting', + statusUpdatedAt: new Date(eventAt.getTime() + 1_000), + }).where(eq(agents.id, alice.agentId)); + markLegacyStatusEvents(db); + + const replay = await recordSessionEventWithIdempotency( + db, ws.workspaceId, alice.agentId, { type: 'status.blocked', payload: {} }, key, + ); + await expect( + applyStatusEventEffect(db, ws.workspaceId, alice.agentId, replay.event.id, 'blocked'), + ).resolves.toEqual({ claimed: true, mutated: false }); + const [agentRow] = await db.select({ status: agents.status }).from(agents) + .where(eq(agents.id, alice.agentId)); + expect(agentRow!.status).toBe('waiting'); + }); + + it('does not clobber a later heartbeat/liveness write during legacy replay', async () => { + const { ws, alice, db } = await seedAgent(); + const key = 'status-legacy-heartbeat-1'; + const first = await recordSessionEventWithIdempotency( + db, ws.workspaceId, alice.agentId, { type: 'status.blocked', payload: {} }, key, + ); + // Put the legacy event in the past, then use the real last_seen writer; + // 0056's trigger records that liveness write in status_updated_at. + await db.update(sessionEvents).set({ + createdAt: new Date(Date.now() - 10_000), + }).where(eq(sessionEvents.id, first.event.id)); + await db.update(agents).set({ lastSeen: new Date() }).where(eq(agents.id, alice.agentId)); + markLegacyStatusEvents(db); + + const replay = await recordSessionEventWithIdempotency( + db, ws.workspaceId, alice.agentId, { type: 'status.blocked', payload: {} }, key, + ); + await expect( + applyStatusEventEffect(db, ws.workspaceId, alice.agentId, replay.event.id, 'blocked'), + ).resolves.toEqual({ claimed: true, mutated: false }); + const [agentRow] = await db.select({ status: agents.status }).from(agents) + .where(eq(agents.id, alice.agentId)); + expect(agentRow!.status).toBe('active'); + }); + + it.each([ + { name: 'same status', status: 'blocked' }, + { name: 'different status', status: 'waiting' }, + ])('treats an equal legacy write timestamp conservatively ($name)', async ({ status }) => { + const { ws, alice, db } = await seedAgent(); + const key = `status-legacy-equal-${status}`; + const first = await recordSessionEventWithIdempotency( + db, ws.workspaceId, alice.agentId, { type: 'status.blocked', payload: {} }, key, + ); + const eventAt = new Date(first.event.created_at); + await db.update(agents).set({ + status, + statusUpdatedAt: eventAt, + }).where(eq(agents.id, alice.agentId)); + markLegacyStatusEvents(db); + + const replay = await recordSessionEventWithIdempotency( + db, ws.workspaceId, alice.agentId, { type: 'status.blocked', payload: {} }, key, + ); + await expect( + applyStatusEventEffect(db, ws.workspaceId, alice.agentId, replay.event.id, 'blocked'), + ).resolves.toEqual({ claimed: true, mutated: false }); + const [agentRow] = await db.select({ status: agents.status }).from(agents) + .where(eq(agents.id, alice.agentId)); + expect(agentRow!.status).toBe(status); + }); + + it('allows only one concurrent legacy replay to recover an older row', async () => { + const { ws, alice, db } = await seedAgent(); + const key = 'status-legacy-race-1'; + const first = await recordSessionEventWithIdempotency( + db, ws.workspaceId, alice.agentId, { type: 'status.blocked', payload: {} }, key, + ); + await db.update(agents).set({ + statusUpdatedAt: new Date(new Date(first.event.created_at).getTime() - 1_000), + }).where(eq(agents.id, alice.agentId)); + markLegacyStatusEvents(db); + + const [left, right] = await Promise.all([ + applyStatusEventEffect(db, ws.workspaceId, alice.agentId, first.event.id, 'blocked'), + applyStatusEventEffect(db, ws.workspaceId, alice.agentId, first.event.id, 'blocked'), + ]); + expect([left.claimed, right.claimed].sort()).toEqual([false, true]); + expect([left.mutated, right.mutated].sort()).toEqual([false, true]); + const [agentRow] = await db.select({ status: agents.status }).from(agents) + .where(eq(agents.id, alice.agentId)); + expect(agentRow!.status).toBe('blocked'); + }); + + it('applies the status write and completion marker in a single D1-style batch', async () => { + const { ws, alice, db } = await seedAgent(); + const { event } = await recordSessionEventWithIdempotency( + db, ws.workspaceId, alice.agentId, { type: 'status.idle', payload: {} }, 'status-batch-1', + ); + const batches = attachFakeBatch(db); + + await applyStatusEventEffect(db, ws.workspaceId, alice.agentId, event.id, 'idle'); + + expect(batches).toHaveLength(1); + expect(batches[0]).toHaveLength(2); + expectStatementOn(batches[0], 'update', 'agents'); + expectStatementOn(batches[0], 'update', 'session_events'); + + const [agentRow] = await db.select({ status: agents.status }).from(agents).where(eq(agents.id, alice.agentId)); + expect(agentRow!.status).toBe('idle'); + const [eventRow] = await db.select({ statusAppliedAt: sessionEvents.statusAppliedAt }).from(sessionEvents).where(eq(sessionEvents.id, event.id)); + expect(eventRow!.statusAppliedAt).not.toBeNull(); + }); + + it('terminalizes an older pending replay without clobbering a newer applied status', async () => { + const { ws, alice, db } = await seedAgent(); + const old = await recordSessionEventWithIdempotency( + db, ws.workspaceId, alice.agentId, { type: 'status.blocked', payload: {} }, 'status-order-old', + ); + const newer = await recordSessionEventWithIdempotency( + db, ws.workspaceId, alice.agentId, { type: 'status.active', payload: {} }, 'status-order-new', + ); + + const newerEffect = await applyStatusEventEffect( + db, ws.workspaceId, alice.agentId, newer.event.id, 'active', + ); + expect(newerEffect).toEqual({ claimed: true, mutated: true }); + + const oldReplay = await recordSessionEventWithIdempotency( + db, ws.workspaceId, alice.agentId, { type: 'status.blocked', payload: {} }, 'status-order-old', + ); + expect(oldReplay.pendingStatusApplication).toBe(true); + const oldEffect = await applyStatusEventEffect( + db, ws.workspaceId, alice.agentId, oldReplay.event.id, 'blocked', + ); + expect(oldEffect).toEqual({ claimed: true, mutated: false }); + + const [agentRow] = await db.select({ status: agents.status }).from(agents).where(eq(agents.id, alice.agentId)); + expect(agentRow!.status).toBe('active'); + const oldRow = await db.select({ statusAppliedAt: sessionEvents.statusAppliedAt }) + .from(sessionEvents).where(eq(sessionEvents.id, old.event.id)); + expect(oldRow[0]!.statusAppliedAt).not.toBeNull(); + }); + + it('preserves newer status ordering when old and new effects race', async () => { + const { ws, alice, db } = await seedAgent(); + const old = await recordSessionEventWithIdempotency( + db, ws.workspaceId, alice.agentId, { type: 'status.blocked', payload: {} }, 'status-order-race-old', + ); + const newer = await recordSessionEventWithIdempotency( + db, ws.workspaceId, alice.agentId, { type: 'status.active', payload: {} }, 'status-order-race-new', + ); + + const [oldEffect, newerEffect] = await Promise.all([ + applyStatusEventEffect(db, ws.workspaceId, alice.agentId, old.event.id, 'blocked'), + applyStatusEventEffect(db, ws.workspaceId, alice.agentId, newer.event.id, 'active'), + ]); + expect(oldEffect.claimed).toBe(true); + expect(newerEffect.claimed).toBe(true); + const [agentRow] = await db.select({ status: agents.status }).from(agents).where(eq(agents.id, alice.agentId)); + expect(agentRow!.status).toBe('active'); + const rows = await db.select({ statusAppliedAt: sessionEvents.statusAppliedAt }) + .from(sessionEvents).where(eq(sessionEvents.agentId, alice.agentId)); + expect(rows.every((row) => row.statusAppliedAt !== null)).toBe(true); + }); + + it('uses the same ordering fence in a D1-style atomic batch', async () => { + const { ws, alice, db } = await seedAgent(); + const old = await recordSessionEventWithIdempotency( + db, ws.workspaceId, alice.agentId, { type: 'status.blocked', payload: {} }, 'status-order-d1-old', + ); + const newer = await recordSessionEventWithIdempotency( + db, ws.workspaceId, alice.agentId, { type: 'status.active', payload: {} }, 'status-order-d1-new', + ); + const batches = attachFakeBatch(db); + + await expect( + applyStatusEventEffect(db, ws.workspaceId, alice.agentId, newer.event.id, 'active'), + ).resolves.toEqual({ claimed: true, mutated: true }); + await expect( + applyStatusEventEffect(db, ws.workspaceId, alice.agentId, old.event.id, 'blocked'), + ).resolves.toEqual({ claimed: true, mutated: false }); + expect(batches).toHaveLength(2); + const [agentRow] = await db.select({ status: agents.status }).from(agents).where(eq(agents.id, alice.agentId)); + expect(agentRow!.status).toBe('active'); + }); + + it('marks a released agent event complete without emitting status side effects', async () => { + const { ws, alice, db } = await seedAgent(); + await db.update(agents).set({ status: 'released' }).where(eq(agents.id, alice.agentId)); + const { event } = await recordSessionEventWithIdempotency( + db, ws.workspaceId, alice.agentId, { type: 'status.active', payload: {} }, 'status-released-1', + ); + + const effect = await applyStatusEventEffect(db, ws.workspaceId, alice.agentId, event.id, 'active'); + expect(effect).toEqual({ claimed: true, mutated: false }); + const [agentRow] = await db.select({ status: agents.status }).from(agents).where(eq(agents.id, alice.agentId)); + expect(agentRow!.status).toBe('released'); + const [eventRow] = await db.select({ statusAppliedAt: sessionEvents.statusAppliedAt }) + .from(sessionEvents).where(eq(sessionEvents.id, event.id)); + expect(eventRow!.statusAppliedAt).not.toBeNull(); + expect(await db.select().from(workspaceEvents).where(and( + eq(workspaceEvents.workspaceId, ws.workspaceId), + eq(workspaceEvents.type, 'agent.status.active'), + ))).toHaveLength(0); + expect(await db.select().from(pendingEvents).where(and( + eq(pendingEvents.workspaceId, ws.workspaceId), + eq(pendingEvents.eventType, 'agent.status.active'), + ))).toHaveLength(0); + }); + + it('does not fan out or enqueue a released agent status through the HTTP route', async () => { + const { ws, alice, db } = await seedAgent(); + await db.update(agents).set({ status: 'released' }).where(eq(agents.id, alice.agentId)); + + const response = await stack.app.request('/v1/agents/alice/events', { + method: 'POST', + headers: { + 'content-type': 'application/json', + authorization: `Bearer ${ws.workspaceKey}`, + 'Idempotency-Key': 'status-released-route-1', + }, + body: JSON.stringify({ type: 'status.active', payload: {} }), + }); + expect(response.status).toBe(201); + await stack.settle(); + + expect(await db.select().from(workspaceEvents).where(and( + eq(workspaceEvents.workspaceId, ws.workspaceId), + eq(workspaceEvents.type, 'agent.status.active'), + ))).toHaveLength(0); + expect(await db.select().from(pendingEvents).where(and( + eq(pendingEvents.workspaceId, ws.workspaceId), + eq(pendingEvents.eventType, 'agent.status.active'), + ))).toHaveLength(0); + }); + + it('rejects a bare handle with neither atomicity capability rather than silently applying only half the write', async () => { + const { ws, alice, db } = await seedAgent(); + const { event } = await recordSessionEventWithIdempotency( + db, ws.workspaceId, alice.agentId, { type: 'status.blocked', payload: {} }, 'status-bare-1', + ); + stripCapability(db); + + await expect( + applyStatusEventEffect(db, ws.workspaceId, alice.agentId, event.id, 'blocked'), + ).rejects.toThrow('Atomic write capability required'); + + const [agentRow] = await db.select({ status: agents.status }).from(agents).where(eq(agents.id, alice.agentId)); + expect(agentRow!.status).not.toBe('blocked'); + }); + + it('resolves genuinely concurrent keyed status posts to one applied status without interleaving', async () => { + const ws = await createWorkspace(stack.app, 'status-race-ws'); + const runner = await registerAgent(stack.app, ws.workspaceKey, 'runner'); + const headers = { + 'content-type': 'application/json', + authorization: `Bearer ${ws.workspaceKey}`, + 'Idempotency-Key': 'status-race-1', + }; + const body = JSON.stringify({ type: 'status.waiting', payload: {} }); + const postEvent = () => stack.app.request('/v1/agents/runner/events', { method: 'POST', headers, body }); + + const [first, second] = await Promise.all([postEvent(), postEvent()]); + expect([first.status, second.status]).toEqual([201, 201]); + + const db = stack.runtime.handle.db as unknown as EngineDb; + const [agentRow] = await db.select({ status: agents.status }).from(agents).where(eq(agents.id, runner.agentId)); + expect(agentRow!.status).toBe('waiting'); + expect(await db.select().from(sessionEvents).where(eq(sessionEvents.agentId, runner.agentId))).toHaveLength(1); + }); + + it('lets only the atomic completion winner emit status side effects', async () => { + const ws = await createWorkspace(stack.app, 'status-side-effect-race-ws'); + const runner = await registerAgent(stack.app, ws.workspaceKey, 'runner'); + const subscription = await stack.app.request('/v1/subscriptions', { + method: 'POST', + headers: { 'content-type': 'application/json', authorization: `Bearer ${ws.workspaceKey}` }, + body: JSON.stringify({ events: ['*'], url: 'http://127.0.0.1:1/hook' }), + }); + expect(subscription.status).toBe(201); + + const db = stack.runtime.handle.db as unknown as EngineDb; + const handle = db as EngineDb & TransactionCapability; + const originalWithTransaction = handle.withTransaction.bind(db); + let releaseFirstTransaction!: () => void; + const firstTransactionDelayed = new Promise((resolve) => { releaseFirstTransaction = resolve; }); + let delayed = false; + handle.withTransaction = async (fn) => { + if (!delayed) { + delayed = true; + await firstTransactionDelayed; + } + return originalWithTransaction(fn); + }; + + const headers = { + 'content-type': 'application/json', + authorization: `Bearer ${ws.workspaceKey}`, + 'Idempotency-Key': 'status-side-effect-race-1', + }; + const body = JSON.stringify({ type: 'status.waiting', payload: {} }); + const postEvent = () => stack.app.request('/v1/agents/runner/events', { method: 'POST', headers, body }); + const firstPost = postEvent(); + // Ensure the second request reaches its pending-status path while the + // first writer is paused, rather than relying on scheduler luck. + await new Promise((resolve) => setTimeout(resolve, 10)); + const secondPost = postEvent(); + await new Promise((resolve) => setTimeout(resolve, 10)); + releaseFirstTransaction(); + const [first, second] = await Promise.all([firstPost, secondPost]); + expect([first.status, second.status]).toEqual([201, 201]); + await stack.settle(); + handle.withTransaction = originalWithTransaction; + + const matchingWorkspaceEvents = await db.select().from(workspaceEvents).where(and( + eq(workspaceEvents.workspaceId, ws.workspaceId), + eq(workspaceEvents.type, 'agent.status.waiting'), + )); + const matchingOutboxRows = await db.select().from(pendingEvents).where(and( + eq(pendingEvents.workspaceId, ws.workspaceId), + eq(pendingEvents.eventType, 'agent.status.waiting'), + )); + expect(matchingWorkspaceEvents).toHaveLength(1); + expect(matchingOutboxRows).toHaveLength(1); + expect(await db.select().from(sessionEvents).where(eq(sessionEvents.agentId, runner.agentId))).toHaveLength(1); + }); + }); }); diff --git a/packages/engine/src/__tests__/conformance/sdk-contract.test.ts b/packages/engine/src/__tests__/conformance/sdk-contract.test.ts index a6d2fa43..5be836ee 100644 --- a/packages/engine/src/__tests__/conformance/sdk-contract.test.ts +++ b/packages/engine/src/__tests__/conformance/sdk-contract.test.ts @@ -1,7 +1,7 @@ import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest'; import { and, eq } from 'drizzle-orm'; import { deliverEvent } from '../../engine/eventDelivery.js'; -import { actionInvocations, webhooks } from '../../db/schema.js'; +import { actionInvocations, sessionEvents, webhooks } from '../../db/schema.js'; import { makeNodeStack, createWorkspace, @@ -864,6 +864,247 @@ describe('SDK v8 service contract', () => { }); }); + it('replays keyed harness events after a lost response and rejects payload changes', async () => { + const ws = await createWorkspace(stack.app, 'sdk-event-idempotency-ws'); + const runner = await registerAgent(stack.app, ws.workspaceKey, 'runner'); + const headers = { + 'content-type': 'application/json', + authorization: `Bearer ${ws.workspaceKey}`, + 'Idempotency-Key': 'worker-exit-generation-1', + }; + const body = JSON.stringify({ type: 'error', payload: { code: 'worker_exit', generation: 1 } }); + + // Treat the first response as lost: the durable event is still the source + // of truth for the subsequent caller retry. + const first = await stack.app.request('/v1/agents/runner/events', { + method: 'POST', headers, body, + }); + expect(first.status).toBe(201); + const firstBody = await first.json() as { data: { id: string; sequence: number } }; + + const replay = await stack.app.request('/v1/agents/runner/events', { + method: 'POST', headers, body, + }); + expect(replay.status).toBe(201); + expect(replay.headers.get('Idempotency-Replayed')).toBe('true'); + const replayBody = await replay.json() as { data: { id: string; sequence: number } }; + expect(replayBody).toEqual(firstBody); + + const conflict = await stack.app.request('/v1/agents/runner/events', { + method: 'POST', + headers, + body: JSON.stringify({ type: 'error', payload: { code: 'different' } }), + }); + expect(conflict.status).toBe(409); + await expect(conflict.json()).resolves.toMatchObject({ + ok: false, + error: { code: 'idempotency_key_reused' }, + }); + + const unkeyedHeaders = { + 'content-type': 'application/json', + authorization: `Bearer ${ws.workspaceKey}`, + }; + const unkeyedBody = JSON.stringify({ type: 'log', payload: { message: 'legacy' } }); + const unkeyed = await Promise.all([ + stack.app.request('/v1/agents/runner/events', { method: 'POST', headers: unkeyedHeaders, body: unkeyedBody }), + stack.app.request('/v1/agents/runner/events', { method: 'POST', headers: unkeyedHeaders, body: unkeyedBody }), + ]); + expect(unkeyed.map((response) => response.status)).toEqual([201, 201]); + + const runnerEvents = await stack.runtime.deps.db + .select({ id: sessionEvents.id, idempotencyKeyHash: sessionEvents.idempotencyKeyHash }) + .from(sessionEvents) + .where(and( + eq(sessionEvents.workspaceId, ws.workspaceId), + eq(sessionEvents.agentId, runner.agentId), + )); + expect(runnerEvents).toHaveLength(3); + expect(runnerEvents.filter((event) => event.idempotencyKeyHash !== null)).toHaveLength(1); + }); + + it('resolves truly concurrent same-key same-payload event posts to one persisted row', async () => { + const ws = await createWorkspace(stack.app, 'sdk-event-idempotency-race-ws'); + const runner = await registerAgent(stack.app, ws.workspaceKey, 'runner'); + const idempotencyKey = 'worker-exit-race-1'; + const headers = { + 'content-type': 'application/json', + authorization: `Bearer ${ws.workspaceKey}`, + 'Idempotency-Key': idempotencyKey, + }; + const body = JSON.stringify({ type: 'error', payload: { code: 'worker_exit', generation: 1 } }); + const postEvent = () => stack.app.request('/v1/agents/runner/events', { + method: 'POST', headers, body, + }); + + // Fire both requests without awaiting either first, so both reach the + // durable insert concurrently and genuinely race on the unique claim + // rather than serializing through a caller-side await. + const [first, second] = await Promise.all([postEvent(), postEvent()]); + expect([first.status, second.status]).toEqual([201, 201]); + + const replayedFlags = [first, second] + .map((response) => response.headers.get('Idempotency-Replayed')) + .sort(); + // Exactly one request wins the durable insert (fresh); the other reads + // back the winner's row and replays it. + expect(replayedFlags).toEqual([null, 'true']); + + const [firstBody, secondBody] = await Promise.all([ + first.json() as Promise<{ data: { id: string; sequence: number }; replayed: boolean }>, + second.json() as Promise<{ data: { id: string; sequence: number }; replayed: boolean }>, + ]); + // Both responses must describe the identical persisted event, regardless + // of which request happened to win the race. + expect(secondBody.data).toEqual(firstBody.data); + + const runnerEvents = await stack.runtime.deps.db + .select({ id: sessionEvents.id, idempotencyKeyHash: sessionEvents.idempotencyKeyHash }) + .from(sessionEvents) + .where(and( + eq(sessionEvents.workspaceId, ws.workspaceId), + eq(sessionEvents.agentId, runner.agentId), + )); + // Only one row was ever persisted for the shared key, no matter which + // request's insert physically won. + expect(runnerEvents).toHaveLength(1); + expect(runnerEvents[0]!.id).toBe(firstBody.data.id); + }); + + it('resolves truly concurrent same-key different-payload event posts deterministically', async () => { + const ws = await createWorkspace(stack.app, 'sdk-event-idempotency-conflict-race-ws'); + const runner = await registerAgent(stack.app, ws.workspaceKey, 'runner'); + const idempotencyKey = 'worker-exit-conflict-race-1'; + const headers = { + 'content-type': 'application/json', + authorization: `Bearer ${ws.workspaceKey}`, + 'Idempotency-Key': idempotencyKey, + }; + const postEvent = (code: string) => stack.app.request('/v1/agents/runner/events', { + method: 'POST', + headers, + body: JSON.stringify({ type: 'error', payload: { code } }), + }); + + // Same key, deliberately different payloads, fired concurrently. The + // unique claim on (workspace, agent, key) guarantees exactly one insert + // wins regardless of scheduling order; the loser's payload never + // matches the persisted digest, so it must fail closed rather than + // silently returning the winner's data as if it were its own. + const [a, b] = await Promise.all([postEvent('worker_exit_a'), postEvent('worker_exit_b')]); + const statuses = [a.status, b.status].sort(); + expect(statuses).toEqual([201, 409]); + + const conflictResponse = a.status === 409 ? a : b; + await expect(conflictResponse.json()).resolves.toMatchObject({ + ok: false, + error: { code: 'idempotency_key_reused' }, + }); + + const runnerEvents = await stack.runtime.deps.db + .select({ id: sessionEvents.id, idempotencyKeyHash: sessionEvents.idempotencyKeyHash }) + .from(sessionEvents) + .where(and( + eq(sessionEvents.workspaceId, ws.workspaceId), + eq(sessionEvents.agentId, runner.agentId), + )); + // Exactly one payload variant is durably persisted; the conflicting + // concurrent write never allocates a second row. + expect(runnerEvents).toHaveLength(1); + }); + + it('relaycast#425: a replay finishes an agent status mutation interrupted after the event was durably claimed', async () => { + const ws = await createWorkspace(stack.app, 'sdk-event-status-crash-ws'); + const runner = await registerAgent(stack.app, ws.workspaceKey, 'runner'); + const headers = { + 'content-type': 'application/json', + authorization: `Bearer ${ws.workspaceKey}`, + 'Idempotency-Key': 'status-crash-replay-http-1', + }; + const body = JSON.stringify({ type: 'status.blocked', payload: {} }); + + // Simulate a crash between the durable event/idempotency claim commit + // and the agent status write: inject a failure into the status mutation + // only, so the caller sees an error while the event row is already + // durable — exactly the interrupted window relaycast#425 is about. + const { agents } = await import('../../db/schema.js'); + + // Recursively proxy the builder chain (`.update(...).set(...).where(...)`) + // so the failure surfaces only when the final statement is executed + // (awaited), not when it is merely built — matching the real crash + // window (the process dies mid-statement, not before it starts). + function failOnExecute(target: T, message: string): T { + return new Proxy(target, { + get(obj, prop) { + if (prop === 'then') { + return (onFulfilled?: (v: unknown) => unknown, onRejected?: (e: unknown) => unknown) => + Promise.reject(new Error(message)).then(onFulfilled, onRejected); + } + const value = Reflect.get(obj, prop) as unknown; + if (typeof value === 'function') { + return (...args: unknown[]) => { + const result = (value as (...a: unknown[]) => unknown).apply(obj, args); + return result && typeof result === 'object' ? failOnExecute(result as object, message) : result; + }; + } + return value; + }, + }); + } + + const db = stack.runtime.deps.db as unknown as { update: (t: unknown) => object }; + const realUpdate = db.update.bind(db); + let failNext = true; + db.update = (table: unknown) => { + const builder = realUpdate(table); + if (table !== agents || !failNext) return builder; + failNext = false; + return failOnExecute(builder, 'injected crash before status write'); + }; + + const first = await stack.app.request('/v1/agents/runner/events', { method: 'POST', headers, body }); + expect(first.status).toBe(500); + db.update = realUpdate; + + // The event claim committed despite the crash; the agent row must not + // have moved yet, and the completion marker must still be NULL. + const [afterCrash] = await stack.runtime.deps.db + .select({ status: agents.status }) + .from(agents) + .where(eq(agents.id, runner.agentId)); + expect(afterCrash!.status).not.toBe('blocked'); + const [eventAfterCrash] = await stack.runtime.deps.db + .select({ statusAppliedAt: sessionEvents.statusAppliedAt }) + .from(sessionEvents) + .where(and(eq(sessionEvents.workspaceId, ws.workspaceId), eq(sessionEvents.agentId, runner.agentId))); + expect(eventAfterCrash!.statusAppliedAt).toBeNull(); + + // The retry with the same key replays the durable event, but must + // still finish the interrupted status mutation rather than returning + // 201 against a stale agent row. + const replay = await stack.app.request('/v1/agents/runner/events', { method: 'POST', headers, body }); + expect(replay.status).toBe(201); + expect(replay.headers.get('Idempotency-Replayed')).toBe('true'); + + const [afterReplay] = await stack.runtime.deps.db + .select({ status: agents.status }) + .from(agents) + .where(eq(agents.id, runner.agentId)); + expect(afterReplay!.status).toBe('blocked'); + const [eventAfterReplay] = await stack.runtime.deps.db + .select({ statusAppliedAt: sessionEvents.statusAppliedAt }) + .from(sessionEvents) + .where(and(eq(sessionEvents.workspaceId, ws.workspaceId), eq(sessionEvents.agentId, runner.agentId))); + expect(eventAfterReplay!.statusAppliedAt).not.toBeNull(); + + // No orphan second event row was ever allocated for the shared key. + expect(await stack.runtime.deps.db + .select() + .from(sessionEvents) + .where(and(eq(sessionEvents.workspaceId, ws.workspaceId), eq(sessionEvents.agentId, runner.agentId))), + ).toHaveLength(1); + }); + it('emits canonical message.reacted events for reactions', async () => { const ws = await createWorkspace(stack.app, 'sdk-reaction-ws'); const alice = await registerAgent(stack.app, ws.workspaceKey, 'alice'); diff --git a/packages/engine/src/adapters/node/__tests__/database.test.ts b/packages/engine/src/adapters/node/__tests__/database.test.ts index d49c6b23..3b210450 100644 --- a/packages/engine/src/adapters/node/__tests__/database.test.ts +++ b/packages/engine/src/adapters/node/__tests__/database.test.ts @@ -88,6 +88,72 @@ describe('delivery sequence high-water migration', () => { .toEqual({ seq: 8 }); }); }); + +describe('session event status completion migration', () => { + it('marks only historical keyed status events for ambiguous replay reconciliation', () => { + const sqlite = new Database(':memory:'); + handles.push(sqlite); + sqlite.exec(` + CREATE TABLE agents ( + id TEXT PRIMARY KEY, + created_at INTEGER NOT NULL, + last_seen INTEGER NOT NULL, + status TEXT NOT NULL + ); + CREATE TABLE session_events ( + id TEXT PRIMARY KEY, + workspace_id TEXT NOT NULL, + agent_id TEXT NOT NULL, + type TEXT NOT NULL, + payload TEXT NOT NULL DEFAULT '{}', + created_at INTEGER NOT NULL, + idempotency_key_hash TEXT, + request_digest TEXT + ); + INSERT INTO agents (id, created_at, last_seen, status) + VALUES ('agent_1', 1600000000, 1700000000, 'active'); + INSERT INTO session_events + (id, workspace_id, agent_id, type, created_at, idempotency_key_hash) + VALUES + ('keyed_status', 'ws_1', 'agent_1', 'status.blocked', 1700000001, 'hash-1'), + ('keyed_changed', 'ws_1', 'agent_1', 'status.changed', 1700000002, 'hash-2'), + ('keyed_non_status', 'ws_1', 'agent_1', 'tool.called', 1700000003, 'hash-3'), + ('unkeyed_status', 'ws_1', 'agent_1', 'status.active', 1700000004, NULL); + `); + + const migration = readFileSync( + new URL('../../../db/migrations/0056_session_event_status_completion.sql', import.meta.url), + 'utf8', + ); + sqlite.exec(migration); + sqlite.prepare(` + INSERT INTO session_events + (id, workspace_id, agent_id, type, created_at, idempotency_key_hash) + VALUES ('post_migration_status', 'ws_1', 'agent_1', 'status.idle', 1700000005, 'hash-5') + `).run(); + + expect(sqlite.prepare(` + SELECT id, status_applied_at, status_legacy_pending + FROM session_events + ORDER BY id + `).all()).toEqual([ + { id: 'keyed_changed', status_applied_at: null, status_legacy_pending: 1 }, + { id: 'keyed_non_status', status_applied_at: null, status_legacy_pending: 0 }, + { id: 'keyed_status', status_applied_at: null, status_legacy_pending: 1 }, + { id: 'post_migration_status', status_applied_at: null, status_legacy_pending: 0 }, + { id: 'unkeyed_status', status_applied_at: null, status_legacy_pending: 0 }, + ]); + expect(sqlite.prepare(`SELECT status_updated_at FROM agents WHERE id = 'agent_1'`).get()) + .toEqual({ status_updated_at: 1700000000 }); + + sqlite.prepare(`UPDATE agents SET last_seen = 1700000010 WHERE id = 'agent_1'`).run(); + const touched = sqlite.prepare(`SELECT status_updated_at FROM agents WHERE id = 'agent_1'`).get() as { + status_updated_at: number; + }; + expect(touched.status_updated_at).toBeGreaterThanOrEqual(1700000000); + expect(touched.status_updated_at).toBeLessThanOrEqual(Math.floor(Date.now() / 1_000) + 1); + }); +}); describe('action invocation provider migration', () => { it('backfills action-owned and legacy node dispatches without claiming undispatched work', () => { const sqlite = new Database(':memory:'); diff --git a/packages/engine/src/db/__tests__/compactMigrations.test.ts b/packages/engine/src/db/__tests__/compactMigrations.test.ts index 1a52967b..2dcce6a9 100644 --- a/packages/engine/src/db/__tests__/compactMigrations.test.ts +++ b/packages/engine/src/db/__tests__/compactMigrations.test.ts @@ -46,7 +46,16 @@ function fixture(count = 20) { function snapshot(handle: SqliteDbHandle) { const tables = handle.sqlite.prepare("SELECT name FROM sqlite_master WHERE type='table' AND name NOT LIKE 'sqlite_%' AND name NOT IN ('_engine_migrations','maintenance_cursors') ORDER BY name").all() as { name: string }[]; return tables.map(({ name }) => { - const rows = handle.sqlite.prepare(`SELECT * FROM "${name}"`).all().map(row => JSON.stringify(row)).sort(); + const rows = handle.sqlite.prepare(`SELECT * FROM "${name}"`).all().map(row => { + // 0056 adds a conservative reconciliation witness and backfills it from + // the existing last_seen value. Exclude that additive bookkeeping field + // so this regression continues to compare all pre-existing user data. + if (name === 'agents') { + const { status_updated_at: _statusUpdatedAt, ...existing } = row as Record; + return JSON.stringify(existing); + } + return JSON.stringify(row); + }).sort(); return [name, rows.length, createHash('sha256').update(JSON.stringify(rows)).digest('hex')]; }); } @@ -71,10 +80,30 @@ function expectConstraintsPreserved(before: ReturnType, afte // 0050 may add these redundant named lookup indexes, but must not alter any // original constraint (including a lookup that already existed via 0049). expect(after.filter(table => table.name !== 'maintenance_cursors' || before.some(original => original.name === table.name)) - .map(table => ({ ...table, uniqueIndexes: table.uniqueIndexes.filter(index => - !['idx_deliveries_id_lookup', 'idx_read_receipts_retention'].includes(index.name) - || before.find(original => original.name === table.name)!.uniqueIndexes.some(original => original.name === index.name), - ) }))).toEqual(before); + .map(table => { + const original = before.find(candidate => candidate.name === table.name); + // 0055 adds optional event identity columns and their index after the + // compact-maintenance path. Keep this regression guard focused on the + // pre-existing constraints while still checking that none was dropped. + if (table.name === 'session_events' && original) { + return { + ...table, + sql: original.sql, + uniqueIndexes: table.uniqueIndexes.filter(index => + original.uniqueIndexes.some(previous => previous.name === index.name)), + }; + } + if (table.name === 'agents' && original) { + return { ...table, sql: original.sql }; + } + return { + ...table, + uniqueIndexes: table.uniqueIndexes.filter(index => + !['idx_deliveries_id_lookup', 'idx_read_receipts_retention'].includes(index.name) + || original?.uniqueIndexes.some(previous => previous.name === index.name), + ), + }; + })).toEqual(before); } describe('compact maintenance migration path', () => { diff --git a/packages/engine/src/db/migrations/0055_session_event_idempotency.sql b/packages/engine/src/db/migrations/0055_session_event_idempotency.sql new file mode 100644 index 00000000..3cd4feef --- /dev/null +++ b/packages/engine/src/db/migrations/0055_session_event_idempotency.sql @@ -0,0 +1,8 @@ +-- relaycast#423: make optional Idempotency-Key event publication durable. +-- NULL identity columns preserve the historical append-only behavior for +-- callers that do not send a key. +ALTER TABLE session_events ADD COLUMN idempotency_key_hash TEXT; +ALTER TABLE session_events ADD COLUMN request_digest TEXT; + +CREATE UNIQUE INDEX session_events_agent_idempotency_unique + ON session_events(workspace_id, agent_id, idempotency_key_hash); diff --git a/packages/engine/src/db/migrations/0056_session_event_status_completion.sql b/packages/engine/src/db/migrations/0056_session_event_status_completion.sql new file mode 100644 index 00000000..4c9f509e --- /dev/null +++ b/packages/engine/src/db/migrations/0056_session_event_status_completion.sql @@ -0,0 +1,38 @@ +-- relaycast#425: track whether a status.* session event's agent-row mutation +-- durably completed. NULL means "not yet applied" so a replay after a crash +-- between the event insert and the agent status write can finish the +-- interrupted mutation instead of silently skipping it. Legacy keyed rows +-- are marked for reconciliation below because their old route wrote the +-- event and agent row separately, so completion is otherwise ambiguous. +-- `status_updated_at` is a conservative write-time witness for that +-- reconciliation. It is initialized from the agent's last known liveness +-- write because older rows have no historical status-write clock. The trigger +-- advances it for every later status or liveness write, so a legacy replay +-- may mutate only when the row is demonstrably older than the event itself. +ALTER TABLE agents ADD COLUMN status_updated_at INTEGER; +ALTER TABLE session_events ADD COLUMN status_applied_at INTEGER; +ALTER TABLE session_events ADD COLUMN status_legacy_pending INTEGER NOT NULL DEFAULT 0; + +UPDATE agents +SET status_updated_at = last_seen +WHERE status_updated_at IS NULL; + +CREATE TRIGGER agents_status_reconciliation_timestamp +AFTER UPDATE OF status, last_seen ON agents +FOR EACH ROW +BEGIN + UPDATE agents + SET status_updated_at = unixepoch() + WHERE id = NEW.id; +END; + +-- Before these markers existed, a keyed status event was inserted before the +-- agent-row update. Keep completion NULL so an interrupted mutation can be +-- recovered only when the agent row is demonstrably older than the event. An +-- equal or newer timestamp is conservatively treated as already handled (or +-- changed by another writer), so replay cannot clobber a later status or +-- heartbeat when the old route's outcome is ambiguous. +UPDATE session_events +SET status_legacy_pending = 1 +WHERE idempotency_key_hash IS NOT NULL + AND type LIKE 'status.%'; diff --git a/packages/engine/src/db/schema.ts b/packages/engine/src/db/schema.ts index 859a8113..82479644 100644 --- a/packages/engine/src/db/schema.ts +++ b/packages/engine/src/db/schema.ts @@ -134,6 +134,10 @@ export const agents = sqliteTable( deliverySeq: integer('delivery_seq').notNull().default(0), createdAt: integer('created_at', { mode: 'timestamp' }).notNull().default(sql`(unixepoch())`), lastSeen: integer('last_seen', { mode: 'timestamp' }).notNull().default(sql`(unixepoch())`), + // Conservative witness used to reconcile pre-0056 keyed status events. + // Migration 0056 backfills historical rows from last_seen and its SQLite + // trigger advances this value for every later status/liveness write. + statusUpdatedAt: integer('status_updated_at', { mode: 'timestamp' }), }, (table) => [ uniqueIndex('agents_workspace_name_unique').on(table.workspaceId, table.name), @@ -1148,11 +1152,27 @@ export const sessionEvents = sqliteTable( .references(() => agents.id, { onDelete: 'cascade' }), type: text('type').notNull(), payload: text('payload', { mode: 'json' }).notNull().default({}), + // Optional stable identity for callers that may retry after losing the + // response. NULL keeps the legacy append-only contract for unkeyed calls. + idempotencyKeyHash: text('idempotency_key_hash'), + requestDigest: text('request_digest'), sequence: integer('sequence').notNull().default(0), createdAt: integer('created_at', { mode: 'timestamp' }).notNull().default(sql`(unixepoch())`), + // Durable completion marker for a `status.*` event's agent-row mutation. + // NULL means "not yet applied" — set atomically with the `agents` row + // write (see `applyStatusEventEffect`) so a crash between the durable + // event insert and the status update leaves this NULL, letting a replay + // finish the interrupted mutation instead of silently skipping it forever. + statusAppliedAt: integer('status_applied_at', { mode: 'timestamp' }), + // Migration 0056 marks keyed status events created before this marker + // existed. Their old event insert and agent update were separate writes, + // so replay reconciles only when the agent write-time witness proves the + // row predates the event; otherwise it claims without clobbering state. + statusLegacyPending: integer('status_legacy_pending', { mode: 'boolean' }).notNull().default(false), }, (table) => [ uniqueIndex('session_events_agent_sequence_unique').on(table.agentId, table.sequence), + uniqueIndex('session_events_agent_idempotency_unique').on(table.workspaceId, table.agentId, table.idempotencyKeyHash), index('idx_session_events_agent').on(table.agentId, table.createdAt), index('idx_session_events_workspace').on(table.workspaceId, table.createdAt), index('idx_session_events_type').on(table.workspaceId, table.type, table.createdAt), diff --git a/packages/engine/src/engine/sessionEvent.ts b/packages/engine/src/engine/sessionEvent.ts index 7396071f..9ddf0bd1 100644 --- a/packages/engine/src/engine/sessionEvent.ts +++ b/packages/engine/src/engine/sessionEvent.ts @@ -1,7 +1,11 @@ -import { eq, and, desc, sql } from 'drizzle-orm'; +import { eq, and, ne, desc, isNull, sql } from 'drizzle-orm'; import type { getDb } from '../db/index.js'; -import { sessionEvents } from '../db/schema.js'; +import { agents, sessionEvents } from '../db/schema.js'; import { generateId } from './snowflake.js'; +import { sha256Hex } from '../lib/crypto.js'; +import { codedError } from '../lib/httpError.js'; +import { runAtomicWrites } from '../ports/database.js'; +import { RELEASED_AGENT_STATUS } from './agent.js'; type Db = ReturnType; @@ -61,6 +65,18 @@ export function isValidEventType(type: string): type is SessionEventType { return VALID_EVENT_TYPES.has(type); } +/** `status.*` events are the only ones with a side-effecting agent-row mutation to track. */ +export function isStatusEventType(type: string): boolean { + return type.startsWith('status.'); +} + +export interface StatusEventEffectResult { + /** Whether this request won the durable completion-marker claim. */ + claimed: boolean; + /** Whether the current agent row was actually changed by this event. */ + mutated: boolean; +} + export async function recordSessionEvent( db: Db, workspaceId: string, @@ -69,7 +85,7 @@ export async function recordSessionEvent( type: SessionEventType; payload: Record; }, -) { +): Promise<{ event: ReturnType; replayed: boolean; pendingStatusApplication: boolean }> { const id = `evt_${generateId()}`; // Sequence is assigned atomically via a scalar subquery — read and write in @@ -86,6 +102,235 @@ export async function recordSessionEvent( }) .returning(); + return { + event: toPublicEvent(event, agentId), + replayed: false, + // Unkeyed events have no durable claim to replay against — the caller + // always applies the status mutation immediately after this call, so + // there is nothing to recover across a retry. + pendingStatusApplication: isStatusEventType(data.type), + }; +} + +/** + * Record an event with a durable caller-provided identity. + * + * The key is hashed before persistence and scoped by workspace + agent in the + * unique index. The request digest is retained beside the event so reusing a + * key with a different event cannot accidentally replay the first event. + */ +export async function recordSessionEventWithIdempotency( + db: Db, + workspaceId: string, + agentId: string, + data: { + type: SessionEventType; + payload: Record; + }, + idempotencyKey: string, +): Promise<{ event: ReturnType; replayed: boolean; pendingStatusApplication: boolean }> { + const [idempotencyKeyHash, requestDigest] = await Promise.all([ + sha256Hex(`session-event-key-v1\0${idempotencyKey}`), + sha256Hex(`session-event-payload-v1\0${canonicalJson({ type: data.type, payload: data.payload })}`), + ]); + const id = `evt_${generateId()}`; + + // The insert and the unique identity claim are one durable operation. A + // conflict is followed by a scoped lookup so concurrent retries return the + // winner's exact event rather than allocating a second sequence number. + const [created] = await db + .insert(sessionEvents) + .values({ + id, + workspaceId, + agentId, + type: data.type, + payload: data.payload, + idempotencyKeyHash, + requestDigest, + sequence: sql`(SELECT COALESCE(MAX(sequence), 0) + 1 FROM session_events WHERE agent_id = ${agentId})`, + }) + .onConflictDoNothing() + .returning(); + + if (created) { + return { + event: toPublicEvent(created, agentId), + replayed: false, + pendingStatusApplication: isStatusEventType(data.type), + }; + } + + const [existing] = await db + .select() + .from(sessionEvents) + .where(and( + eq(sessionEvents.workspaceId, workspaceId), + eq(sessionEvents.agentId, agentId), + eq(sessionEvents.idempotencyKeyHash, idempotencyKeyHash), + )); + if (!existing) { + throw codedError( + 'The event idempotency claim could not be read after a storage conflict; retry with the same Idempotency-Key', + 'idempotency_unavailable', + 503, + ); + } + if (existing.requestDigest !== requestDigest) { + throw codedError( + 'Idempotency-Key was reused with a different request payload', + 'idempotency_key_reused', + 409, + ); + } + return { + event: toPublicEvent(existing, agentId), + replayed: true, + // `status_applied_at` is set atomically with the agent-row status write + // (see `applyStatusEventEffect`). NULL here means either the mutation + // never ran, or it ran but the process crashed before marking it durable + // — both cases are indistinguishable from "not yet applied" and safe to + // retry, because the write that sets this column is the same atomic unit + // as the status mutation itself. A replay that is still pending finishes + // the interrupted work instead of returning 201 with a stale agent row. + pendingStatusApplication: isStatusEventType(existing.type) && existing.statusAppliedAt == null, + }; +} + +/** + * Apply a `status.*` event's agent-row mutation and claim its completion + * durably, as one atomic unit. + * + * This is the fix for the crash window between "the event/idempotency claim + * committed" and "the agent's status row is updated": if either statement + * here failed independently, a retry could see `replayed: true` and return + * 201 while the agent row stayed stale forever. Batching both writes through + * `runAtomicWrites` means a failure here rolls back *both* the status change + * and the completion marker, so `pendingStatusApplication` stays true and the + * next replay retries the whole mutation — it can never observe "claimed but + * never applied" as a terminal state. + * + * The completion update is the single-winner claim for side effects: a replay + * that loses a concurrent claim gets `claimed: false` and must not fan out or + * enqueue another webhook. The agent update also carries a durable ordering + * fence: an older pending event may be terminalized, but it cannot overwrite + * a newer status event that was recorded while the older event was pending. + * A released agent's row is intentionally not updated (matching + * `updateAgentById`), but the event is still marked applied: there is no + * agent row left to reconcile, and retrying forever would just repeat the + * same no-op. + */ +export async function applyStatusEventEffect( + db: Db, + workspaceId: string, + agentId: string, + eventId: string, + status: string, +): Promise { + // The agent mutation runs first in the atomic unit. This ordering is + // important for D1 batches: the completion marker must still be NULL when + // the conditional status update is evaluated. The following marker update + // then claims the event; Node transactions and D1 batches serialize this + // pair, so an identical concurrent retry cannot mutate after the winner has + // claimed it. + const [mutationResult, claimResult] = await runAtomicWrites(db, (tx) => [ + tx.update(agents) + .set({ status }) + .where(and( + eq(agents.workspaceId, workspaceId), + eq(agents.id, agentId), + ne(agents.status, RELEASED_AGENT_STATUS), + // Do not let an interrupted older event roll a newer status back. + // The sequence is allocated per agent and is durable across retries; + // an event that is newer than this one wins the status row. A newer + // event that is itself pending will be responsible for applying its + // own status when it replays. + sql`EXISTS ( + SELECT 1 + FROM session_events AS current_event + WHERE current_event.id = ${eventId} + AND current_event.workspace_id = ${workspaceId} + AND current_event.agent_id = ${agentId} + AND current_event.status_applied_at IS NULL + )`, + sql`NOT EXISTS ( + SELECT 1 + FROM session_events AS newer_event + WHERE newer_event.workspace_id = ${workspaceId} + AND newer_event.agent_id = ${agentId} + AND newer_event.type LIKE 'status.%' + AND newer_event.sequence > ( + SELECT current_event.sequence + FROM session_events AS current_event + WHERE current_event.id = ${eventId} + AND current_event.workspace_id = ${workspaceId} + AND current_event.agent_id = ${agentId} + ) + )`, + // Migration 0056 cannot know whether the pre-marker route reached + // this agent write: it inserted the event and updated the agent in + // separate operations. Its write-time witness is therefore the only + // safe evidence available. Reapply a legacy event only when the agent + // row demonstrably predates the event; an equal or newer witness may + // be the original status write, a heartbeat, or another status writer. + // In those ambiguous cases claim the event but never clobber the row, + // regardless of whether its current status matches the event. + sql`( + NOT EXISTS ( + SELECT 1 + FROM session_events AS legacy_event + WHERE legacy_event.id = ${eventId} + AND legacy_event.workspace_id = ${workspaceId} + AND legacy_event.agent_id = ${agentId} + AND legacy_event.status_legacy_pending = 1 + ) + OR ( + ${agents.status} <> ${status} + AND + ${agents.statusUpdatedAt} IS NOT NULL + AND ${agents.statusUpdatedAt} < ( + SELECT legacy_event.created_at + FROM session_events AS legacy_event + WHERE legacy_event.id = ${eventId} + AND legacy_event.workspace_id = ${workspaceId} + AND legacy_event.agent_id = ${agentId} + ) + ) + )`, + )) + .returning({ id: agents.id }), + tx.update(sessionEvents) + .set({ statusAppliedAt: sql`(unixepoch())` }) + .where(and( + eq(sessionEvents.id, eventId), + eq(sessionEvents.workspaceId, workspaceId), + eq(sessionEvents.agentId, agentId), + isNull(sessionEvents.statusAppliedAt), + )) + .returning({ id: sessionEvents.id }), + ], { requireAtomic: true }); + + const claimed = Array.isArray(claimResult) && claimResult.length > 0; + const mutated = Array.isArray(mutationResult) && mutationResult.length > 0; + // Only the caller that both changed the current row and claimed the marker + // may emit external status side effects. The claim may intentionally win + // with `mutated === false` when a newer event superseded this one or the + // agent was released between route validation and this atomic write. + return { claimed, mutated: claimed && mutated }; +} + +function canonicalJson(value: unknown): string { + if (Array.isArray(value)) return `[${value.map(canonicalJson).join(',')}]`; + if (value && typeof value === 'object') { + return `{${Object.keys(value as Record) + .sort() + .map((key) => `${JSON.stringify(key)}:${canonicalJson((value as Record)[key])}`) + .join(',')}}`; + } + return JSON.stringify(value); +} + +function toPublicEvent(event: typeof sessionEvents.$inferSelect, agentId: string) { return { id: event.id, agent_id: agentId, diff --git a/packages/engine/src/routes/agent.ts b/packages/engine/src/routes/agent.ts index 069d2470..407b48e0 100644 --- a/packages/engine/src/routes/agent.ts +++ b/packages/engine/src/routes/agent.ts @@ -847,6 +847,9 @@ agentRoutes.post( const { type, payload } = parsed.data; + const { key: idempotencyKey, error: idempotencyError } = parseIdempotencyKey(c.req.header('Idempotency-Key')); + if (idempotencyError) return jsonError(c, 'invalid_idempotency_key', idempotencyError, 400); + if (!sessionEventEngine.isValidEventType(type)) { return jsonError(c, 'invalid_event_type', `Unknown event type: ${type}`, 400); } @@ -877,16 +880,39 @@ agentRoutes.post( } } - const event = await sessionEventEngine.recordSessionEvent(db, workspace.id, agentRecord.id, { - type, - payload, - }); - - // Update agent status after the event is durably written - if (type.startsWith('status.')) { + const recorded = idempotencyKey + ? await sessionEventEngine.recordSessionEventWithIdempotency( + db, + workspace.id, + agentRecord.id, + { type, payload }, + idempotencyKey, + ) + : await sessionEventEngine.recordSessionEvent(db, workspace.id, agentRecord.id, { type, payload }); + const { event, replayed, pendingStatusApplication } = recorded; + + // Apply the agent status mutation and durably mark it complete as one + // atomic unit (see `applyStatusEventEffect`). `pendingStatusApplication` + // is true both for a fresh event and for a replay whose status write + // never completed (crash between the durable event claim and the agent + // update) — either way this finishes the interrupted mutation instead + // of returning 201 against a stale agent row. + let statusApplied = false; + if (pendingStatusApplication && type.startsWith('status.')) { + const resolved = sessionEventEngine.resolveStatusFromEvent(type); + const newStatus = resolved ?? (payload.status as string); + const effect = await sessionEventEngine.applyStatusEventEffect( + db, + workspace.id, + agentRecord.id, + event.id, + newStatus, + ); + statusApplied = effect.mutated; + } + if (statusApplied) { const resolved = sessionEventEngine.resolveStatusFromEvent(type); const newStatus = resolved ?? (payload.status as string); - await agentEngine.updateAgent(db, workspace.id, name, { status: newStatus }); const eventType = type === 'status.changed' ? 'agent.status.changed' : `agent.status.${canonicalStatus(newStatus) ?? newStatus}`; const eventData = { agent_id: agentRecord.id, @@ -903,7 +929,7 @@ agentRoutes.post( }); } - if (!type.startsWith('status.')) { + if (!replayed && !type.startsWith('status.')) { const { type: _sessionEventType, ...eventWithoutType } = event; const eventData = { agent_name: name, ...eventWithoutType }; runInBackground(c, fanoutToWorkspace(c, `harness.${type}`, eventData), `fanout harness.${type}`); @@ -914,7 +940,7 @@ agentRoutes.post( }); } - return jsonCreated(c, event); + return jsonIdempotentOk(c, { status: 201, data: event, replayed }); } catch (err: unknown) { return errorResponse(c, err); } diff --git a/packages/sdk-rust/CHANGELOG.md b/packages/sdk-rust/CHANGELOG.md index 66044c16..122023ef 100644 --- a/packages/sdk-rust/CHANGELOG.md +++ b/packages/sdk-rust/CHANGELOG.md @@ -26,6 +26,8 @@ The format is based on Keep a Changelog, and this project follows Semantic Versi ### Added +- `RelayCast::emit_agent_event_with_idempotency_key` preserves a stable event key across retries. + - `RelayCast::create_workspace` now requires explicit provenance, preventing CLI bootstrap workspaces from being mislabeled as SDK-created. - `WorkspaceBootstrapOptions` and `RelayCast::create_workspace_with_options()` support crash-safe anonymous keyed workspace creation without a deployment-wide secret, plus opt-in self-host proof via `with_bootstrap_secret(...)`; the SDK never sends that proof to hosted Relaycast. - `RelayCast::release_agent_if_token_hash` performs generation-safe cleanup without changing the existing `ReleaseAgentRequest` struct. diff --git a/packages/sdk-rust/src/relay.rs b/packages/sdk-rust/src/relay.rs index 12459951..ff1f490f 100644 --- a/packages/sdk-rust/src/relay.rs +++ b/packages/sdk-rust/src/relay.rs @@ -849,6 +849,23 @@ impl RelayCast { .await } + /// Emit a session event with a durable identity. Reuse the same key when + /// retrying after an ambiguous response so Relaycast replays one event. + pub async fn emit_agent_event_with_idempotency_key( + &self, + name: &str, + request: EmitSessionEventRequest, + idempotency_key: impl Into, + ) -> Result { + self.client + .post( + &format!("/v1/agents/{}/events", urlencoding::encode(name)), + Some(request), + Some(RequestOptions::with_idempotency_key(idempotency_key)), + ) + .await + } + /// List recorded session events for an agent. pub async fn list_agent_events( &self, diff --git a/packages/sdk-rust/tests/parity.rs b/packages/sdk-rust/tests/parity.rs index a10e9ddd..5f1155d5 100644 --- a/packages/sdk-rust/tests/parity.rs +++ b/packages/sdk-rust/tests/parity.rs @@ -1895,6 +1895,46 @@ async fn agent_session_events_use_expected_endpoints() { assert!(events.is_empty()); } +#[tokio::test] +async fn keyed_agent_session_events_preserve_idempotency_key() { + let server = MockServer::start().await; + let relay = RelayCast::new(RelayCastOptions::new("rk_live_test").with_base_url(server.uri())) + .expect("failed to create relay client"); + + Mock::given(method("POST")) + .and(path("/v1/agents/Worker/events")) + .and(header("Idempotency-Key", "worker-exit-generation-1")) + .respond_with(ResponseTemplate::new(201).set_body_json(json!({ + "ok": true, + "data": { + "id": "evt_1", + "agent_id": "agent_1", + "type": "error", + "payload": {"code": "worker_exit"}, + "created_at": "2026-01-01T00:00:00.000Z" + } + }))) + .expect(1) + .mount(&server) + .await; + + let event = relay + .emit_agent_event_with_idempotency_key( + "Worker", + EmitSessionEventRequest { + event_type: "error".to_string(), + payload: Some(serde_json::Map::from_iter([( + "code".to_string(), + json!("worker_exit"), + )])), + }, + "worker-exit-generation-1", + ) + .await + .expect("keyed emit_agent_event failed"); + assert_eq!(event.event_type, "error"); +} + #[test] fn deserializes_action_ws_events() { let invoked: WsEvent = serde_json::from_value(json!({