From c36941ccb676069d876f3fa81d4012e49201e356 Mon Sep 17 00:00:00 2001 From: Khaliq Date: Thu, 13 Aug 2026 20:50:04 +0200 Subject: [PATCH 1/3] fix(engine): atomically claim legacy identities --- CHANGELOG.md | 4 + README.md | 6 + openapi.yaml | 54 +++++- packages/engine/CHANGELOG.md | 4 + .../conformance/legacyIdentityClaim.test.ts | 159 ++++++++++++++++++ packages/engine/src/engine/agent.ts | 68 +++++++- packages/engine/src/routes/agent.ts | 87 ++++++++++ packages/types/CHANGELOG.md | 6 +- packages/types/src/telemetry.ts | 2 + 9 files changed, 386 insertions(+), 4 deletions(-) create mode 100644 packages/engine/src/__tests__/conformance/legacyIdentityClaim.test.ts diff --git a/CHANGELOG.md b/CHANGELOG.md index 3a652402..d9a7c3e0 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -26,6 +26,10 @@ Packages without a separate changelog are covered by the cross-package notes bel - DM conversation and inbox reads remain available beyond 100 conversations, and the advertised DM-list limit now reaches the server instead of being ignored. +### Security + +- Legacy agent identity claims are atomic, and generic agent updates cannot write the reserved `identity_key` metadata field. + ## [8.0.0] - 2026-08-10 ### Added diff --git a/README.md b/README.md index 366d2cf9..1320a10b 100644 --- a/README.md +++ b/README.md @@ -86,6 +86,12 @@ npx tsx quickstart.ts That is the canonical onboarding loop: create workspace, register agents, connect realtime streams, and watch messages flow live. +Operator recovery for an offline agent registered before identity verifiers +were stored uses `PATCH /v1/agents/:name/legacy-identity`. The endpoint accepts +only a SHA-256 verifier (`identity_key_hash`) and atomically succeeds when the +record is still offline and has no `identity_key` field at write time. Generic +agent updates cannot write that reserved field. + Workspace names are not globally unique. Workspace creation is idempotent for the same workspace name and API key: repeating that combination returns the existing workspace instead of creating another one. If you want an explicit SDK helper that tells you whether setup returned an existing workspace or created a new one, use `ensureWorkspace()`: diff --git a/openapi.yaml b/openapi.yaml index 01a50f60..1286ed73 100644 --- a/openapi.yaml +++ b/openapi.yaml @@ -1785,7 +1785,10 @@ paths: patch: summary: Update agent - description: Update agent properties + description: >- + Update agent properties. Platform-managed metadata keys such as + `identity_key` are rejected; use the dedicated legacy identity claim + endpoint for the one supported backfill operation. tags: - Agents security: @@ -1836,6 +1839,55 @@ paths: '204': description: Agent deleted + /agents/{name}/legacy-identity: + patch: + summary: Atomically claim a legacy agent identity + description: >- + Operator recovery for one offline agent registered before identity + verifiers were stored. The update succeeds only when `identity_key` + is still absent and the durable agent status is still `offline` in + the same atomic database mutation. The request contains a SHA-256 + verifier, never the raw identity proof. + tags: + - Agents + security: + - workspaceKey: [] + parameters: + - name: name + in: path + required: true + schema: + type: string + requestBody: + required: true + content: + application/json: + schema: + type: object + required: + - identity_key_hash + properties: + identity_key_hash: + type: string + pattern: '^[a-f0-9]{64}$' + description: SHA-256 hex verifier for the caller-held identity proof + responses: + '200': + description: Legacy identity claimed + content: + application/json: + schema: + type: object + properties: + ok: + type: boolean + data: + $ref: '#/components/schemas/Agent' + '404': + description: Agent not found + '409': + description: Agent already has an identity or is not offline + /agents/spawn: post: summary: Request agent spawn diff --git a/packages/engine/CHANGELOG.md b/packages/engine/CHANGELOG.md index b25908b6..311468bd 100644 --- a/packages/engine/CHANGELOG.md +++ b/packages/engine/CHANGELOG.md @@ -13,6 +13,10 @@ and this project follows [Semantic Versioning](https://semver.org/spec/v2.0.0.ht - DM conversation and unread-inbox enrichment batches large identifier sets below D1's bound-parameter ceiling, so long-lived agents no longer lose both read paths after accumulating more than 100 conversations. +### Security + +- Atomically claim an offline legacy agent's identity and reject `identity_key` writes through generic agent updates. + ## [8.0.0] - 2026-08-10 ### Added diff --git a/packages/engine/src/__tests__/conformance/legacyIdentityClaim.test.ts b/packages/engine/src/__tests__/conformance/legacyIdentityClaim.test.ts new file mode 100644 index 00000000..480aac0f --- /dev/null +++ b/packages/engine/src/__tests__/conformance/legacyIdentityClaim.test.ts @@ -0,0 +1,159 @@ +import { afterEach, beforeEach, describe, expect, it } from 'vitest'; +import { createWorkspace, makeNodeStack, registerAgent, type TestStack } from './harness.js'; +import * as agentEngine from '../../engine/agent.js'; + +describe('legacy agent identity claim', () => { + let stack: TestStack; + + beforeEach(() => { + stack = makeNodeStack(); + }); + + afterEach(() => stack.close()); + + async function markOffline(workspaceKey: string, name: string): Promise { + const response = await stack.app.request(`/v1/agents/${name}`, { + method: 'PATCH', + headers: { + authorization: `Bearer ${workspaceKey}`, + 'content-type': 'application/json', + }, + body: JSON.stringify({ status: 'offline' }), + }); + expect(response.status).toBe(200); + } + + function claim(workspaceKey: string, name: string, identityKeyHash: string) { + return stack.app.request(`/v1/agents/${name}/legacy-identity`, { + method: 'PATCH', + headers: { + authorization: `Bearer ${workspaceKey}`, + 'content-type': 'application/json', + }, + body: JSON.stringify({ identity_key_hash: identityKeyHash }), + }); + } + + it('allows exactly one of two concurrent claims to stamp the identity', async () => { + const workspace = await createWorkspace(stack.app, 'legacy-identity-race'); + await registerAgent(stack.app, workspace.workspaceKey, 'legacy-node'); + await markOffline(workspace.workspaceKey, 'legacy-node'); + + const firstHash = 'a'.repeat(64); + const secondHash = 'b'.repeat(64); + const [first, second] = await Promise.all([ + claim(workspace.workspaceKey, 'legacy-node', firstHash), + claim(workspace.workspaceKey, 'legacy-node', secondHash), + ]); + + expect([first.status, second.status].sort()).toEqual([200, 409]); + const winner = first.status === 200 ? first : second; + const loser = first.status === 409 ? first : second; + const winnerBody = await winner.json() as { + data: { metadata: { identity_key: string } }; + }; + await expect(loser.json()).resolves.toMatchObject({ + error: { code: 'agent_identity_already_claimed' }, + }); + + const current = await stack.app.request('/v1/agents/legacy-node', { + headers: { authorization: `Bearer ${workspace.workspaceKey}` }, + }); + const currentBody = await current.json() as { + data: { metadata: { identity_key: string } }; + }; + expect(currentBody.data.metadata.identity_key).toBe(winnerBody.data.metadata.identity_key); + expect([firstHash, secondHash]).toContain(currentBody.data.metadata.identity_key); + }); + + it('fails closed when identity_key is present with a non-string value', async () => { + const workspace = await createWorkspace(stack.app, 'legacy-identity-null'); + const registered = await stack.app.request('/v1/agents', { + method: 'POST', + headers: { + authorization: `Bearer ${workspace.workspaceKey}`, + 'content-type': 'application/json', + }, + body: JSON.stringify({ + name: 'malformed-identity-node', + metadata: { identity_key: null }, + }), + }); + expect(registered.status).toBe(201); + await markOffline(workspace.workspaceKey, 'malformed-identity-node'); + + const response = await claim( + workspace.workspaceKey, + 'malformed-identity-node', + 'c'.repeat(64), + ); + + expect(response.status).toBe(409); + await expect(response.json()).resolves.toMatchObject({ + error: { code: 'agent_identity_already_claimed' }, + }); + }); + + it('rejects identity_key writes through the generic agent update endpoint', async () => { + const workspace = await createWorkspace(stack.app, 'reserved-agent-metadata'); + await registerAgent(stack.app, workspace.workspaceKey, 'protected-node'); + + const response = await stack.app.request('/v1/agents/protected-node', { + method: 'PATCH', + headers: { + authorization: `Bearer ${workspace.workspaceKey}`, + 'content-type': 'application/json', + }, + body: JSON.stringify({ metadata: { identity_key: 'attacker-planted-proof' } }), + }); + + expect(response.status).toBe(400); + await expect(response.json()).resolves.toMatchObject({ + error: { code: 'reserved_agent_metadata_key' }, + }); + const current = await stack.app.request('/v1/agents/protected-node', { + headers: { authorization: `Bearer ${workspace.workspaceKey}` }, + }); + const currentBody = await current.json() as { + data: { metadata: Record }; + }; + expect(currentBody.data.metadata).not.toHaveProperty('identity_key'); + }); + + it('preserves a winning claim against a generic metadata update built from a stale read', async () => { + const workspace = await createWorkspace(stack.app, 'legacy-identity-stale-patch'); + await registerAgent(stack.app, workspace.workspaceKey, 'legacy-node'); + await markOffline(workspace.workspaceKey, 'legacy-node'); + + // Model the generic PATCH route after it read the legacy row but before it + // writes its merged metadata. The claim lands between those two steps. + const staleMetadata = { operator_note: 'read-before-claim' }; + const claimedHash = 'e'.repeat(64); + const winner = await claim(workspace.workspaceKey, 'legacy-node', claimedHash); + expect(winner.status).toBe(200); + + const staleWrite = await agentEngine.updateAgent( + stack.runtime.deps.db, + workspace.workspaceId, + 'legacy-node', + { metadata: staleMetadata }, + ); + + expect(staleWrite?.metadata).toEqual({ + ...staleMetadata, + identity_key: claimedHash, + }); + }); + + it('requires the target agent to be offline at the atomic write', async () => { + const workspace = await createWorkspace(stack.app, 'legacy-identity-online'); + await registerAgent(stack.app, workspace.workspaceKey, 'online-node'); + + const response = await claim(workspace.workspaceKey, 'online-node', 'd'.repeat(64)); + + expect(response.status).toBe(409); + await expect(response.json()).resolves.toMatchObject({ + error: { code: 'agent_not_offline' }, + }); + }); +}); diff --git a/packages/engine/src/engine/agent.ts b/packages/engine/src/engine/agent.ts index 77de087b..e516fd2d 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, inArray } from 'drizzle-orm'; +import { eq, and, gt, lt, ne, inArray, sql } from 'drizzle-orm'; import type { getDb } from '../db/index.js'; import { agents, agentNodeBindings, channels, channelMembers, actions, deliveries, nodes } from '../db/schema.js'; import { randomHex, sha256Hex } from '../lib/crypto.js'; @@ -372,7 +372,22 @@ export async function updateAgent( const setClause: Record = {}; if (updates.status !== undefined) setClause.status = updates.status; if (updates.persona !== undefined) setClause.persona = updates.persona; - if (updates.metadata !== undefined) setClause.metadata = updates.metadata; + if (updates.metadata !== undefined) { + const nextMetadata = JSON.stringify(updates.metadata); + // `identity_key` is platform-managed. Preserve whatever value exists at + // write time, even when this update was built from a stale pre-claim read; + // otherwise an overlapping generic metadata PATCH could erase a winning + // legacy claim immediately after its atomic UPDATE. + setClause.metadata = sql`CASE + WHEN json_type(COALESCE(${agents.metadata}, '{}'), '$.identity_key') IS NULL + THEN json(${nextMetadata}) + ELSE json_set( + json(${nextMetadata}), + '$.identity_key', + json_extract(COALESCE(${agents.metadata}, '{}'), '$.identity_key') + ) + END`; + } if (updates.capabilities !== undefined) setClause.capabilities = updates.capabilities; if (Object.keys(setClause).length === 0) { @@ -406,6 +421,55 @@ export async function updateAgent( }; } +/** + * Atomically stamp the one-time identity verifier for a pre-gate agent. + * + * `json_type(..., '$.identity_key') IS NULL` matches an absent key only; + * explicit JSON null returns the string `null` and therefore fails closed. + * Keeping the metadata mutation and both preconditions in one UPDATE prevents + * two operators from both observing an eligible row and overwriting each + * other's claim. + */ +export async function claimLegacyAgentIdentity( + db: Db, + workspaceId: string, + name: string, + identityKeyHash: string, +) { + // Bring a genuinely stale lease to `offline` before the atomic claim. A + // concurrent heartbeat that wins after this sweep changes status first and + // makes the UPDATE predicate fail. + await sweepStaleAgents(db, workspaceId); + + const [updated] = await db + .update(agents) + .set({ + metadata: sql`json_set(COALESCE(${agents.metadata}, '{}'), '$.identity_key', ${identityKeyHash})`, + }) + .where(and( + eq(agents.workspaceId, workspaceId), + eq(agents.name, name), + eq(agents.status, 'offline'), + sql`json_type(COALESCE(${agents.metadata}, '{}'), '$.identity_key') IS NULL`, + )) + .returning(); + + if (!updated) return null; + + return { + id: updated.id, + name: updated.name, + handle: `@${updated.name}`, + type: updated.type, + status: updated.status, + persona: updated.persona, + capabilities: updated.capabilities ?? null, + created_at: updated.createdAt.toISOString(), + last_seen: updated.lastSeen.toISOString(), + metadata: updated.metadata, + }; +} + export async function deleteAgent(db: Db, workspaceId: string, name: string) { const [agent] = await db .select() diff --git a/packages/engine/src/routes/agent.ts b/packages/engine/src/routes/agent.ts index afd60708..ed69337d 100644 --- a/packages/engine/src/routes/agent.ts +++ b/packages/engine/src/routes/agent.ts @@ -62,6 +62,12 @@ const updateAgentSchema = z.object({ capabilities: capabilitiesSchema.nullable().optional(), }); +const legacyIdentityClaimSchema = z.object({ + identity_key_hash: z.string().regex(/^[a-f0-9]{64}$/), +}); + +const RESERVED_AGENT_METADATA_KEYS = new Set(['identity_key']); + const sessionEventSchema = z.object({ type: z.string().min(1), payload: z.record(z.string(), z.unknown()).default({}), @@ -307,6 +313,17 @@ agentRoutes.patch( } const body = parsed.data; + const reservedMetadataKey = body.metadata + ? Object.keys(body.metadata).find((key) => RESERVED_AGENT_METADATA_KEYS.has(key)) + : undefined; + if (reservedMetadataKey) { + return jsonError( + c, + 'reserved_agent_metadata_key', + `Agent metadata key "${reservedMetadataKey}" is managed by Relaycast and cannot be updated through the generic agent endpoint`, + 400, + ); + } const nextMetadata = body.metadata !== undefined || body.skills !== undefined ? { ...(existing.metadata || {}), @@ -347,6 +364,76 @@ agentRoutes.patch( }, ); +// PATCH /v1/agents/:name/legacy-identity - atomically claim a pre-gate identity +agentRoutes.patch( + '/agents/:name/legacy-identity', + requireWorkspaceKey, + rateLimit, + async (c) => { + try { + const db = c.get('db'); + const workspace = c.get('workspace'); + const name = c.req.param('name'); + const parsed = await parseJsonBody( + c, + legacyIdentityClaimSchema, + 'invalid legacy identity claim body', + ); + if (!parsed.ok) { + return parsed.response; + } + + const updated = await agentEngine.claimLegacyAgentIdentity( + db, + workspace.id, + name, + parsed.data.identity_key_hash, + ); + if (!updated) { + const existing = await agentEngine.getAgentByName(db, workspace.id, name); + if (!existing) { + return agentNotFound(c, name); + } + if (Object.prototype.hasOwnProperty.call(existing.metadata ?? {}, 'identity_key')) { + return jsonError( + c, + 'agent_identity_already_claimed', + `Agent "${name}" already has an identity_key and cannot be reclaimed through the legacy path`, + 409, + ); + } + if (existing.status !== 'offline') { + return jsonError( + c, + 'agent_not_offline', + `Agent "${name}" is not offline and cannot be reclaimed through the legacy path`, + 409, + ); + } + return jsonError( + c, + 'agent_identity_claim_conflict', + `Agent "${name}" changed while its legacy identity was being claimed; retry only after verifying its current state`, + 409, + ); + } + + await directoryEngine.syncSourceAgentDirectoryEntry(db, workspace.id, { + id: updated.id, + name: updated.name, + status: updated.status, + metadata: updated.metadata ?? {}, + }); + emitServerEvent(c, workspace.id, 'relaycast_server_legacy_identity_claimed', { + agent_name: name, + }); + return jsonOk(c, updated); + } catch (err: unknown) { + return errorResponse(c, err); + } + }, +); + // DELETE /v1/agents/:name - delete agent agentRoutes.delete( '/agents/:name', diff --git a/packages/types/CHANGELOG.md b/packages/types/CHANGELOG.md index 30f2f9e6..086a4fc7 100644 --- a/packages/types/CHANGELOG.md +++ b/packages/types/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] + +### Added + +- Declare the server audit event emitted after a legacy agent identity is claimed. ## [8.0.0] - 2026-08-10 diff --git a/packages/types/src/telemetry.ts b/packages/types/src/telemetry.ts index e2f91f7f..57e4be22 100644 --- a/packages/types/src/telemetry.ts +++ b/packages/types/src/telemetry.ts @@ -37,6 +37,7 @@ export const SERVER_TELEMETRY_EVENTS = [ 'relaycast_server_workspace_stream_updated', 'relaycast_server_agent_registered', 'relaycast_server_agent_updated', + 'relaycast_server_legacy_identity_claimed', 'relaycast_server_agent_deleted', 'relaycast_server_agent_token_rotated', 'relaycast_server_channel_created', @@ -145,6 +146,7 @@ const REQUIRED_SERVER_EVENT_PROPS: Record Date: Thu, 13 Aug 2026 21:18:34 +0200 Subject: [PATCH 2/3] fix(engine): harden identity metadata invariants --- CHANGELOG.md | 2 +- README.md | 5 +- openapi.yaml | 19 +++- packages/engine/CHANGELOG.md | 4 +- .../conformance/legacyIdentityClaim.test.ts | 104 ++++++++++++++++-- packages/engine/src/engine/a2a.ts | 21 ++-- packages/engine/src/engine/agent.ts | 32 +++++- packages/engine/src/routes/agent.ts | 17 ++- packages/types/CHANGELOG.md | 2 +- .../src/__tests__/sdk-openapi-sync.test.ts | 3 + 10 files changed, 170 insertions(+), 39 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index d9a7c3e0..59d46749 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -16,7 +16,7 @@ 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 - Patch] +## [Unreleased - Minor] ### Changed diff --git a/README.md b/README.md index 1320a10b..572ca82f 100644 --- a/README.md +++ b/README.md @@ -90,7 +90,10 @@ Operator recovery for an offline agent registered before identity verifiers were stored uses `PATCH /v1/agents/:name/legacy-identity`. The endpoint accepts only a SHA-256 verifier (`identity_key_hash`) and atomically succeeds when the record is still offline and has no `identity_key` field at write time. Generic -agent updates cannot write that reserved field. +agent updates cannot write that reserved field. Initial registration may set a +valid lowercase SHA-256 verifier, but malformed registration values and claims +return 400; already-claimed, non-offline, or concurrently changed records return +409. Workspace names are not globally unique. Workspace creation is idempotent for the same workspace name and API key: repeating that combination returns the existing workspace instead of creating another one. diff --git a/openapi.yaml b/openapi.yaml index 1286ed73..3778c3b8 100644 --- a/openapi.yaml +++ b/openapi.yaml @@ -1682,7 +1682,11 @@ paths: /agents: post: summary: Register agent - description: Register a new agent in the workspace. Agent names are unique per workspace; registering a duplicate name returns `409 agent_already_exists`. + description: >- + Register a new agent in the workspace. Agent names are unique per + workspace; registering a duplicate name returns 409 + agent_already_exists. Initial registration may include identity_key in + metadata only as a 64-character lowercase SHA-256 verifier. tags: - Agents security: @@ -1717,6 +1721,8 @@ paths: type: boolean data: $ref: '#/components/schemas/Agent' + '400': + description: Malformed registration identity_key (invalid_agent_identity_key) '409': description: Agent name already exists in this workspace @@ -1821,6 +1827,10 @@ paths: type: boolean data: $ref: '#/components/schemas/Agent' + '400': + description: >- + Invalid update, including reserved_agent_metadata_key when + metadata contains the platform-managed identity_key field delete: summary: Delete agent @@ -1883,10 +1893,15 @@ paths: type: boolean data: $ref: '#/components/schemas/Agent' + '400': + description: Malformed identity_key_hash (invalid_request) '404': description: Agent not found '409': - description: Agent already has an identity or is not offline + description: >- + Agent already has an identity (agent_identity_already_claimed), is + not offline (agent_not_offline), or changed during the claim + (agent_identity_claim_conflict) /agents/spawn: post: diff --git a/packages/engine/CHANGELOG.md b/packages/engine/CHANGELOG.md index 311468bd..3abdb169 100644 --- a/packages/engine/CHANGELOG.md +++ b/packages/engine/CHANGELOG.md @@ -7,7 +7,7 @@ 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 - Patch] +## [Unreleased - Minor] ### Fixed @@ -15,7 +15,7 @@ and this project follows [Semantic Versioning](https://semver.org/spec/v2.0.0.ht ### Security -- Atomically claim an offline legacy agent's identity and reject `identity_key` writes through generic agent updates. +- `PATCH /v1/agents/:name/legacy-identity` atomically claims an offline legacy agent's identity, while generic agent updates cannot overwrite the verifier. ## [8.0.0] - 2026-08-10 diff --git a/packages/engine/src/__tests__/conformance/legacyIdentityClaim.test.ts b/packages/engine/src/__tests__/conformance/legacyIdentityClaim.test.ts index 480aac0f..f1512a70 100644 --- a/packages/engine/src/__tests__/conformance/legacyIdentityClaim.test.ts +++ b/packages/engine/src/__tests__/conformance/legacyIdentityClaim.test.ts @@ -1,6 +1,9 @@ import { afterEach, beforeEach, describe, expect, it } from 'vitest'; +import { and, eq } from 'drizzle-orm'; import { createWorkspace, makeNodeStack, registerAgent, type TestStack } from './harness.js'; +import * as a2aEngine from '../../engine/a2a.js'; import * as agentEngine from '../../engine/agent.js'; +import { a2aAgents, agents } from '../../db/schema.js'; describe('legacy agent identity claim', () => { let stack: TestStack; @@ -68,19 +71,19 @@ describe('legacy agent identity claim', () => { it('fails closed when identity_key is present with a non-string value', async () => { const workspace = await createWorkspace(stack.app, 'legacy-identity-null'); - const registered = await stack.app.request('/v1/agents', { - method: 'POST', - headers: { - authorization: `Bearer ${workspace.workspaceKey}`, - 'content-type': 'application/json', - }, - body: JSON.stringify({ - name: 'malformed-identity-node', - metadata: { identity_key: null }, - }), - }); - expect(registered.status).toBe(201); + await registerAgent(stack.app, workspace.workspaceKey, 'malformed-identity-node'); await markOffline(workspace.workspaceKey, 'malformed-identity-node'); + // Seed a malformed historical row below the API boundary. Registration + // must continue accepting a valid string verifier because that is how new + // brokers establish ownership; this fixture specifically models old or + // corrupted data that the claim must reject by key presence. + await stack.runtime.deps.db + .update(agents) + .set({ metadata: { identity_key: null } }) + .where(and( + eq(agents.workspaceId, workspace.workspaceId), + eq(agents.name, 'malformed-identity-node'), + )); const response = await claim( workspace.workspaceKey, @@ -94,6 +97,45 @@ describe('legacy agent identity claim', () => { }); }); + it('rejects malformed registration verifiers while retaining valid ownership bootstrap', async () => { + const workspace = await createWorkspace(stack.app, 'registration-identity-verifier'); + const invalid = await stack.app.request('/v1/agents', { + method: 'POST', + headers: { + authorization: `Bearer ${workspace.workspaceKey}`, + 'content-type': 'application/json', + }, + body: JSON.stringify({ + name: 'malformed-verifier', + metadata: { identity_key: null }, + }), + }); + expect(invalid.status).toBe(400); + await expect(invalid.json()).resolves.toMatchObject({ + error: { code: 'invalid_agent_identity_key' }, + }); + + const validHash = 'f'.repeat(64); + const valid = await stack.app.request('/v1/agents', { + method: 'POST', + headers: { + authorization: `Bearer ${workspace.workspaceKey}`, + 'content-type': 'application/json', + }, + body: JSON.stringify({ + name: 'valid-verifier', + metadata: { identity_key: validHash }, + }), + }); + expect(valid.status).toBe(201); + const registered = await agentEngine.getAgentByName( + stack.runtime.deps.db, + workspace.workspaceId, + 'valid-verifier', + ); + expect(registered?.metadata).toMatchObject({ identity_key: validHash }); + }); + it('rejects identity_key writes through the generic agent update endpoint', async () => { const workspace = await createWorkspace(stack.app, 'reserved-agent-metadata'); await registerAgent(stack.app, workspace.workspaceKey, 'protected-node'); @@ -145,6 +187,44 @@ describe('legacy agent identity claim', () => { }); }); + it('preserves a winning claim when an A2A proxy is removed', async () => { + const workspace = await createWorkspace(stack.app, 'legacy-identity-a2a-remove'); + const proxy = await registerAgent(stack.app, workspace.workspaceKey, 'legacy-a2a-proxy'); + await stack.runtime.deps.db.insert(a2aAgents).values({ + id: 'a2a_legacy_identity_test', + workspaceId: workspace.workspaceId, + relayAgentId: proxy.agentId, + agentCard: { + name: 'legacy-a2a-proxy', + url: 'https://example.com/a2a/rpc', + version: '1.0.0', + skills: [{ id: 'echo', name: 'echo' }], + }, + externalUrl: 'https://example.com/a2a/rpc', + }); + await markOffline(workspace.workspaceKey, 'legacy-a2a-proxy'); + + const claimedHash = '9'.repeat(64); + const winner = await claim(workspace.workspaceKey, 'legacy-a2a-proxy', claimedHash); + expect(winner.status).toBe(200); + + await expect(a2aEngine.removeA2aAgent( + stack.runtime.deps.db, + workspace.workspaceId, + 'legacy-a2a-proxy', + )).resolves.toBe(true); + const current = await agentEngine.getAgentByName( + stack.runtime.deps.db, + workspace.workspaceId, + 'legacy-a2a-proxy', + ); + expect(current?.metadata).toMatchObject({ + a2a: true, + a2a_active: false, + identity_key: claimedHash, + }); + }); + it('requires the target agent to be offline at the atomic write', async () => { const workspace = await createWorkspace(stack.app, 'legacy-identity-online'); await registerAgent(stack.app, workspace.workspaceKey, 'online-node'); diff --git a/packages/engine/src/engine/a2a.ts b/packages/engine/src/engine/a2a.ts index 4cbaf677..76c714f3 100644 --- a/packages/engine/src/engine/a2a.ts +++ b/packages/engine/src/engine/a2a.ts @@ -508,19 +508,16 @@ export async function removeA2aAgent(db: Db, workspaceId: string, relayName: str if (!agentRecord) return false; await db.delete(a2aAgents).where(eq(a2aAgents.id, agentRecord.id)); - await db - .update(agents) - .set({ - status: 'offline', - metadata: { - ...(agentRecord.relay_metadata ?? {}), - a2a: true, - a2a_active: false, - }, - }) - .where(eq(agents.id, agentRecord.relay_agent_id)); + const updated = await updateAgent(db, workspaceId, relayName, { + status: 'offline', + metadata: { + ...(agentRecord.relay_metadata ?? {}), + a2a: true, + a2a_active: false, + }, + }); - return true; + return updated !== null; } export function translateRelayToA2a(message: RelayDM): A2aJsonRpcRequest { diff --git a/packages/engine/src/engine/agent.ts b/packages/engine/src/engine/agent.ts index e516fd2d..41364bb1 100644 --- a/packages/engine/src/engine/agent.ts +++ b/packages/engine/src/engine/agent.ts @@ -9,6 +9,21 @@ import { runAtomicWrites, type AtomicWrite } from '../ports/database.js'; type Db = ReturnType; +/** Metadata verifier used by brokers to reclaim their own registration. */ +export const AGENT_IDENTITY_METADATA_KEY = 'identity_key'; +export const AGENT_IDENTITY_HASH_PATTERN = /^[a-f0-9]{64}$/; +const AGENT_IDENTITY_METADATA_JSON_PATH = `$.${AGENT_IDENTITY_METADATA_KEY}`; + +export function hasValidRegistrationIdentity( + metadata: Record | undefined, +): boolean { + if (!metadata || !Object.prototype.hasOwnProperty.call(metadata, AGENT_IDENTITY_METADATA_KEY)) { + return true; + } + const verifier = metadata[AGENT_IDENTITY_METADATA_KEY]; + return typeof verifier === 'string' && AGENT_IDENTITY_HASH_PATTERN.test(verifier); +} + /** How long an authenticated agent can be silent before it is no longer present. */ export const AGENT_LIVENESS_TTL_MS = 5 * 60 * 1000; @@ -134,6 +149,13 @@ export async function registerAgent( }, ) { assertRegistrableAgentName(data.name); + if (!hasValidRegistrationIdentity(data.metadata)) { + throw codedError( + 'Agent registration identity_key must be a lowercase SHA-256 verifier', + 'invalid_agent_identity_key', + 400, + ); + } const agentId = generateId(); const token = `at_live_${randomHex(16)}`; const tokenHash = await sha256Hex(token); @@ -379,12 +401,12 @@ export async function updateAgent( // otherwise an overlapping generic metadata PATCH could erase a winning // legacy claim immediately after its atomic UPDATE. setClause.metadata = sql`CASE - WHEN json_type(COALESCE(${agents.metadata}, '{}'), '$.identity_key') IS NULL + WHEN json_type(COALESCE(${agents.metadata}, '{}'), ${AGENT_IDENTITY_METADATA_JSON_PATH}) IS NULL THEN json(${nextMetadata}) ELSE json_set( json(${nextMetadata}), - '$.identity_key', - json_extract(COALESCE(${agents.metadata}, '{}'), '$.identity_key') + ${AGENT_IDENTITY_METADATA_JSON_PATH}, + json_extract(COALESCE(${agents.metadata}, '{}'), ${AGENT_IDENTITY_METADATA_JSON_PATH}) ) END`; } @@ -444,13 +466,13 @@ export async function claimLegacyAgentIdentity( const [updated] = await db .update(agents) .set({ - metadata: sql`json_set(COALESCE(${agents.metadata}, '{}'), '$.identity_key', ${identityKeyHash})`, + metadata: sql`json_set(COALESCE(${agents.metadata}, '{}'), ${AGENT_IDENTITY_METADATA_JSON_PATH}, ${identityKeyHash})`, }) .where(and( eq(agents.workspaceId, workspaceId), eq(agents.name, name), eq(agents.status, 'offline'), - sql`json_type(COALESCE(${agents.metadata}, '{}'), '$.identity_key') IS NULL`, + sql`json_type(COALESCE(${agents.metadata}, '{}'), ${AGENT_IDENTITY_METADATA_JSON_PATH}) IS NULL`, )) .returning(); diff --git a/packages/engine/src/routes/agent.ts b/packages/engine/src/routes/agent.ts index ed69337d..0014bd0e 100644 --- a/packages/engine/src/routes/agent.ts +++ b/packages/engine/src/routes/agent.ts @@ -63,10 +63,10 @@ const updateAgentSchema = z.object({ }); const legacyIdentityClaimSchema = z.object({ - identity_key_hash: z.string().regex(/^[a-f0-9]{64}$/), + identity_key_hash: z.string().regex(agentEngine.AGENT_IDENTITY_HASH_PATTERN), }); -const RESERVED_AGENT_METADATA_KEYS = new Set(['identity_key']); +const RESERVED_AGENT_METADATA_KEYS = new Set([agentEngine.AGENT_IDENTITY_METADATA_KEY]); const sessionEventSchema = z.object({ type: z.string().min(1), @@ -217,6 +217,14 @@ agentRoutes.post( return parsed.response; } const { name, type, persona, metadata, skills, capabilities } = parsed.data; + if (!agentEngine.hasValidRegistrationIdentity(metadata)) { + return jsonError( + c, + 'invalid_agent_identity_key', + 'Agent registration identity_key must be a lowercase SHA-256 verifier', + 400, + ); + } const nextMetadata = { ...(metadata || {}), ...(skills ? { skills } : {}), @@ -394,7 +402,10 @@ agentRoutes.patch( if (!existing) { return agentNotFound(c, name); } - if (Object.prototype.hasOwnProperty.call(existing.metadata ?? {}, 'identity_key')) { + if (Object.prototype.hasOwnProperty.call( + existing.metadata ?? {}, + agentEngine.AGENT_IDENTITY_METADATA_KEY, + )) { return jsonError( c, 'agent_identity_already_claimed', diff --git a/packages/types/CHANGELOG.md b/packages/types/CHANGELOG.md index 086a4fc7..7c2a788f 100644 --- a/packages/types/CHANGELOG.md +++ b/packages/types/CHANGELOG.md @@ -7,7 +7,7 @@ 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 - Patch] +## [Unreleased - Minor] ### Added diff --git a/packages/types/src/__tests__/sdk-openapi-sync.test.ts b/packages/types/src/__tests__/sdk-openapi-sync.test.ts index 08807ab7..348adeae 100644 --- a/packages/types/src/__tests__/sdk-openapi-sync.test.ts +++ b/packages/types/src/__tests__/sdk-openapi-sync.test.ts @@ -48,6 +48,9 @@ const NON_SDK_OPENAPI_PATHS = new Set([ // node-providers work; the engine surface ships first. '/v1/nodes/{param}/actions/{param}/invoke', '/v1/nodes/{param}/providers/{param}', + // Operator-only, one-time recovery for legacy registrations. Agent SDKs do + // not expose it because normal registration establishes the verifier. + '/v1/agents/{param}/legacy-identity', ]); const CORE_SDK_PATHS = new Set([ From 6c6e5a8e9e939e91869af3b646a45e3dd87f4d42 Mon Sep 17 00:00:00 2001 From: Khaliq Date: Thu, 13 Aug 2026 21:34:14 +0200 Subject: [PATCH 3/3] fix(engine): scope cleanup updates by agent id --- .../conformance/legacyIdentityClaim.test.ts | 32 ++++++++++++ packages/engine/src/engine/a2a.ts | 4 +- packages/engine/src/engine/agent.ts | 51 ++++++++++++++----- packages/engine/src/routes/agent.ts | 8 --- 4 files changed, 73 insertions(+), 22 deletions(-) diff --git a/packages/engine/src/__tests__/conformance/legacyIdentityClaim.test.ts b/packages/engine/src/__tests__/conformance/legacyIdentityClaim.test.ts index f1512a70..f327bd7c 100644 --- a/packages/engine/src/__tests__/conformance/legacyIdentityClaim.test.ts +++ b/packages/engine/src/__tests__/conformance/legacyIdentityClaim.test.ts @@ -225,6 +225,38 @@ describe('legacy agent identity claim', () => { }); }); + it('does not redirect an id-scoped cleanup update to a same-name replacement', async () => { + const workspace = await createWorkspace(stack.app, 'legacy-identity-id-scoped-update'); + const oldProxy = await registerAgent(stack.app, workspace.workspaceKey, 'reused-proxy-name'); + await agentEngine.deleteAgent( + stack.runtime.deps.db, + workspace.workspaceId, + 'reused-proxy-name', + ); + const replacement = await registerAgent( + stack.app, + workspace.workspaceKey, + 'reused-proxy-name', + ); + + const staleUpdate = await agentEngine.updateAgentById( + stack.runtime.deps.db, + workspace.workspaceId, + oldProxy.agentId, + { status: 'offline', metadata: { a2a: true, a2a_active: false } }, + ); + + expect(staleUpdate).toBeNull(); + const current = await agentEngine.getAgentByName( + stack.runtime.deps.db, + workspace.workspaceId, + 'reused-proxy-name', + ); + expect(current?.id).toBe(replacement.agentId); + expect(current?.status).toBe('active'); + expect(current?.metadata).not.toHaveProperty('a2a_active'); + }); + it('requires the target agent to be offline at the atomic write', async () => { const workspace = await createWorkspace(stack.app, 'legacy-identity-online'); await registerAgent(stack.app, workspace.workspaceKey, 'online-node'); diff --git a/packages/engine/src/engine/a2a.ts b/packages/engine/src/engine/a2a.ts index 76c714f3..dfb2e295 100644 --- a/packages/engine/src/engine/a2a.ts +++ b/packages/engine/src/engine/a2a.ts @@ -25,7 +25,7 @@ import { z } from 'zod'; import type { FileAttachment } from '@relaycast/types'; import type { getDb } from '../db/index.js'; import { a2aAgents, agents } from '../db/schema.js'; -import { registerAgent, getAgentByName, updateAgent } from './agent.js'; +import { registerAgent, getAgentByName, updateAgent, updateAgentById } from './agent.js'; import { rotateAgentToken } from './tokenRotate.js'; import { createAndRunCertification } from './certify.js'; import { isSafeExternalUrl } from '../lib/ssrf.js'; @@ -508,7 +508,7 @@ export async function removeA2aAgent(db: Db, workspaceId: string, relayName: str if (!agentRecord) return false; await db.delete(a2aAgents).where(eq(a2aAgents.id, agentRecord.id)); - const updated = await updateAgent(db, workspaceId, relayName, { + const updated = await updateAgentById(db, workspaceId, agentRecord.relay_agent_id, { status: 'offline', metadata: { ...(agentRecord.relay_metadata ?? {}), diff --git a/packages/engine/src/engine/agent.ts b/packages/engine/src/engine/agent.ts index 41364bb1..781f3926 100644 --- a/packages/engine/src/engine/agent.ts +++ b/packages/engine/src/engine/agent.ts @@ -14,7 +14,7 @@ export const AGENT_IDENTITY_METADATA_KEY = 'identity_key'; export const AGENT_IDENTITY_HASH_PATTERN = /^[a-f0-9]{64}$/; const AGENT_IDENTITY_METADATA_JSON_PATH = `$.${AGENT_IDENTITY_METADATA_KEY}`; -export function hasValidRegistrationIdentity( +function hasValidRegistrationIdentity( metadata: Record | undefined, ): boolean { if (!metadata || !Object.prototype.hasOwnProperty.call(metadata, AGENT_IDENTITY_METADATA_KEY)) { @@ -390,6 +390,40 @@ export async function updateAgent( metadata?: Record; capabilities?: Record | null; }, +) { + if ( + updates.status === undefined + && updates.persona === undefined + && updates.metadata === undefined + && updates.capabilities === undefined + ) { + return getAgentByName(db, workspaceId, name); + } + + const [agent] = await db + .select({ id: agents.id }) + .from(agents) + .where(and(eq(agents.workspaceId, workspaceId), eq(agents.name, name))); + + if (!agent) return null; + return updateAgentById(db, workspaceId, agent.id, updates); +} + +/** + * Update one exact agent row while preserving its identity verifier at write + * time. Callers that already resolved an agent must use the id form so a + * delete-and-recreate under the same name cannot redirect their mutation. + */ +export async function updateAgentById( + db: Db, + workspaceId: string, + agentId: string, + updates: { + status?: string; + persona?: string | null; + metadata?: Record; + capabilities?: Record | null; + }, ) { const setClause: Record = {}; if (updates.status !== undefined) setClause.status = updates.status; @@ -412,23 +446,16 @@ export async function updateAgent( } if (updates.capabilities !== undefined) setClause.capabilities = updates.capabilities; - if (Object.keys(setClause).length === 0) { - return getAgentByName(db, workspaceId, name); - } - - const [agent] = await db - .select() - .from(agents) - .where(and(eq(agents.workspaceId, workspaceId), eq(agents.name, name))); - - if (!agent) return null; + if (Object.keys(setClause).length === 0) return null; const [updated] = await db .update(agents) .set(setClause) - .where(eq(agents.id, agent.id)) + .where(and(eq(agents.workspaceId, workspaceId), eq(agents.id, agentId))) .returning(); + if (!updated) return null; + return { id: updated.id, name: updated.name, diff --git a/packages/engine/src/routes/agent.ts b/packages/engine/src/routes/agent.ts index 0014bd0e..ca6b24e5 100644 --- a/packages/engine/src/routes/agent.ts +++ b/packages/engine/src/routes/agent.ts @@ -217,14 +217,6 @@ agentRoutes.post( return parsed.response; } const { name, type, persona, metadata, skills, capabilities } = parsed.data; - if (!agentEngine.hasValidRegistrationIdentity(metadata)) { - return jsonError( - c, - 'invalid_agent_identity_key', - 'Agent registration identity_key must be a lowercase SHA-256 verifier', - 400, - ); - } const nextMetadata = { ...(metadata || {}), ...(skills ? { skills } : {}),