From 4d8a88887a0a6a1a40a0b4b75e4d6c7bae9048d1 Mon Sep 17 00:00:00 2001 From: khaliqgant Date: Thu, 24 Sep 2026 18:02:20 -0700 Subject: [PATCH 1/8] feat(engine): POST /v1/to/:address routes a DM to agent@machine Resolves agent@machine against the agent's current node (node name or machine_id) and sends through the existing DM pipeline, so delivery, idempotency, and fanout are unchanged. Stale addresses 404 instead of reaching the agent elsewhere. Adds agent.sendTo() to the TypeScript SDK. Co-Authored-By: Claude Opus 5.5 (1M context) --- CHANGELOG.md | 6 +- README.md | 6 + openapi.yaml | 75 ++++++ packages/engine/CHANGELOG.md | 6 +- .../conformance/addressedSend.test.ts | 112 ++++++++ packages/engine/src/engine/address.ts | 62 +++++ packages/engine/src/routes/dm.ts | 255 ++++++++++-------- packages/sdk-typescript/CHANGELOG.md | 6 +- .../src/__tests__/agent-messaging.test.ts | 13 + packages/sdk-typescript/src/agent.ts | 19 ++ 10 files changed, 445 insertions(+), 115 deletions(-) create mode 100644 packages/engine/src/__tests__/conformance/addressedSend.test.ts create mode 100644 packages/engine/src/engine/address.ts diff --git a/CHANGELOG.md b/CHANGELOG.md index 91fbb0ba..5f8a4ca5 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -16,7 +16,11 @@ This project follows [Semantic Versioning](https://semver.org/spec/v2.0.0.html). Packages without a separate changelog are covered by the cross-package notes below. -## [Unreleased] +## [Unreleased - Minor] + +### Added + +- `POST /v1/to/:address` sends a DM to an `agent@machine` address, delivering to the agent only while it is hosted on that machine. ## [8.12.0] - 2026-09-24 diff --git a/README.md b/README.md index 8dd6cbf4..e5e750e3 100644 --- a/README.md +++ b/README.md @@ -609,6 +609,7 @@ GET /channels/:name/messages GET /sessions/:session_ref/messages?limit=<1-500>&after= POST /messages/:id/replies POST /dm +POST /to/:address DM an `agent@machine` address (body: `{ "text": "..." }`) GET /dm/conversations?limit=<1-100> List the agent's newest DM conversations (limit optional) GET /inbox GET /search @@ -624,6 +625,11 @@ joined or DMs the agent participates in. Activity feed channel-message items include `channel_id` and `channel_name`; DM items include `conversation_id`. +`POST /to/:address` routes by address instead of bare name: `agent@machine` resolves to the +agent only while it is hosted on that machine (its node's name or `machine_id`), then delivers +like `POST /dm`. A stale address returns `404 address_not_found`; a malformed one returns +`400 invalid_address`. + `POST /dm` can return **`409 dm_conversation_id_collision`**. A 1:1 conversation id is derived deterministically from `(workspace, sorted agent pair)`, and that binding is reserved atomically before any conversation state is created, so a derivation that would name another diff --git a/openapi.yaml b/openapi.yaml index e94e91f5..6887d905 100644 --- a/openapi.yaml +++ b/openapi.yaml @@ -3913,6 +3913,81 @@ paths: schema: $ref: '#/components/schemas/ErrorResponse' + /to/{address}: + post: + summary: Send DM to agent@machine + description: | + Send a direct message to an `agent@machine` address. `machine` must be + the agent's current node, matched by node name or `machine_id`; the + message is then delivered exactly like `POST /dm`. An address that no + longer matches (the agent moved or was released) returns 404 instead + of reaching the agent elsewhere. + tags: + - Direct Messages + security: + - agentToken: [] + parameters: + - name: address + in: path + required: true + description: '`agent@machine`; split on the last `@`' + schema: + type: string + - name: Idempotency-Key + in: header + schema: + type: string + requestBody: + required: true + content: + application/json: + schema: + type: object + required: + - text + properties: + text: + type: string + attachments: + type: array + items: + type: string + description: Optional file ids to attach to this DM + data: + type: object + nullable: true + additionalProperties: true + description: Public structured message metadata + mode: + type: string + enum: [wait, steer] + default: wait + description: Injection mode for the DM + responses: + '201': + description: DM sent + content: + application/json: + schema: + type: object + properties: + ok: + type: boolean + data: + $ref: '#/components/schemas/DmSendResponse' + '400': + description: '`invalid_address` — the address is not of the form `agent@machine`' + content: + application/json: + schema: + $ref: '#/components/schemas/ErrorResponse' + '404': + description: '`address_not_found` — no agent with that name is hosted on that machine' + content: + application/json: + schema: + $ref: '#/components/schemas/ErrorResponse' + /dm/conversations: get: summary: List DM conversations diff --git a/packages/engine/CHANGELOG.md b/packages/engine/CHANGELOG.md index 990817aa..a6f94590 100644 --- a/packages/engine/CHANGELOG.md +++ b/packages/engine/CHANGELOG.md @@ -7,7 +7,11 @@ See the [root changelog](../../CHANGELOG.md) for cross-package release highlight The format is based on [Keep a Changelog](https://keepachangelog.com/en/1.1.0/), and this project follows [Semantic Versioning](https://semver.org/spec/v2.0.0.html). -## [Unreleased] +## [Unreleased - Minor] + +### Added + +- `POST /v1/to/:address` resolves `agent@machine` (node name or `machine_id`) to the agent hosted there and sends it a DM; stale addresses return `404 address_not_found`. ## [8.12.0] - 2026-09-24 diff --git a/packages/engine/src/__tests__/conformance/addressedSend.test.ts b/packages/engine/src/__tests__/conformance/addressedSend.test.ts new file mode 100644 index 00000000..468e1af3 --- /dev/null +++ b/packages/engine/src/__tests__/conformance/addressedSend.test.ts @@ -0,0 +1,112 @@ +import { afterEach, beforeEach, describe, expect, it } from 'vitest'; +import { makeNodeStack, createWorkspace, registerAgent, FakeSocket, type TestStack } from './harness.js'; +import { parseAgentAddress } from '../../engine/address.js'; + +/** + * POST /v1/to/:address — send a DM to `agent@machine`, where `machine` is the + * agent's current node, by node name or machine_id. + */ +describe('addressed send', () => { + let stack: TestStack; + beforeEach(() => { stack = makeNodeStack({ ttlMs: 60_000 }); }); + afterEach(async () => { + await stack.close(); + }); + + async function seed() { + const ws = await createWorkspace(stack.app, 'address-ws'); + const alice = await registerAgent(stack.app, ws.workspaceKey, 'alice'); + + const nodeId = 'node_laptop'; + const nodeName = 'laptop'; + const enroll = await stack.app.request('/v1/nodes', { + method: 'POST', + headers: { 'content-type': 'application/json', authorization: `Bearer ${ws.workspaceKey}` }, + body: JSON.stringify({ + node_id: nodeId, + name: nodeName, + machine_id: 'mach-123', + capabilities: ['spawn:claude'], + max_agents: 4, + version: 'test-node', + }), + }); + expect(enroll.status).toBe(201); + + const sock = new FakeSocket(); + const handle = stack.runtime.realtime.attachNodeSocket(ws.workspaceId, nodeId, sock); + await handle.handleMessage(JSON.stringify({ + v: 1, type: 'node.register', name: nodeName, node_id: nodeId, + capabilities: [{ name: 'spawn:claude', kind: 'capacity' }], + max_agents: 4, tags: [], version: 'test-node', resume_cursor: null, + })); + await handle.handleMessage(JSON.stringify({ + v: 1, type: 'node.heartbeat', load: 0, active_agents: 0, handlers_live: true, + })); + await handle.handleMessage(JSON.stringify({ + v: 1, type: 'agent.register', name: 'bob', resumable: true, session_ref: 'sess-bob', + })); + const reply = sock.ofType('reply').at(-1) as { ok: boolean }; + expect(reply?.ok).toBe(true); + + return { ws, alice, sock }; + } + + function send(token: string, address: string, body: unknown) { + return stack.app.request(`/v1/to/${encodeURIComponent(address)}`, { + method: 'POST', + headers: { 'content-type': 'application/json', authorization: `Bearer ${token}` }, + body: JSON.stringify(body), + }); + } + + it('parses agent@machine on the last @', () => { + expect(parseAgentAddress('bob@laptop')).toEqual({ agent: 'bob', machine: 'laptop' }); + expect(parseAgentAddress('a@b@laptop')).toEqual({ agent: 'a@b', machine: 'laptop' }); + expect(parseAgentAddress('bob')).toBeNull(); + expect(parseAgentAddress('@laptop')).toBeNull(); + expect(parseAgentAddress('bob@')).toBeNull(); + }); + + it('routes to the agent by node name and by machine_id', async () => { + const { alice, sock } = await seed(); + + for (const address of ['bob@laptop', 'bob@mach-123']) { + const res = await send(alice.token, address, { text: `hi via ${address}` }); + expect(res.status).toBe(201); + const body = (await res.json()) as { ok: boolean; data: { text: string } }; + expect(body.ok).toBe(true); + expect(body.data.text).toBe(`hi via ${address}`); + } + + await stack.settle(); + const delivered = sock.ofType('deliver') + .map((frame: { payload?: { type?: string; data?: { message?: { text?: string } } } }) => frame.payload) + .filter((payload) => payload?.type === 'dm.received') + .map((payload) => payload?.data?.message?.text); + expect(delivered).toEqual(['hi via bob@laptop', 'hi via bob@mach-123']); + }); + + it('rejects an address whose machine does not host the agent', async () => { + const { alice } = await seed(); + const res = await send(alice.token, 'bob@desktop', { text: 'hi' }); + expect(res.status).toBe(404); + expect(((await res.json()) as { error: { code: string } }).error.code).toBe('address_not_found'); + }); + + it('rejects unknown agents and malformed addresses', async () => { + const { alice } = await seed(); + const unknown = await send(alice.token, 'carol@laptop', { text: 'hi' }); + expect(unknown.status).toBe(404); + + const malformed = await send(alice.token, 'bob', { text: 'hi' }); + expect(malformed.status).toBe(400); + expect(((await malformed.json()) as { error: { code: string } }).error.code).toBe('invalid_address'); + }); + + it('requires text', async () => { + const { alice } = await seed(); + const res = await send(alice.token, 'bob@laptop', {}); + expect(res.status).toBe(400); + }); +}); diff --git a/packages/engine/src/engine/address.ts b/packages/engine/src/engine/address.ts new file mode 100644 index 00000000..0f975ec2 --- /dev/null +++ b/packages/engine/src/engine/address.ts @@ -0,0 +1,62 @@ +import { and, eq } from 'drizzle-orm'; +import type { getDb } from '../db/index.js'; +import { agents, nodes } from '../db/schema.js'; +import { codedError } from '../lib/httpError.js'; + +type Db = ReturnType; + +export interface AgentAddress { + agent: string; + machine: string; +} + +/** + * Parse an `agent@machine` address. The split is on the last `@` so agent + * names that themselves contain `@` still resolve. + */ +export function parseAgentAddress(address: string): AgentAddress | null { + const at = address.lastIndexOf('@'); + if (at <= 0 || at === address.length - 1) return null; + return { agent: address.slice(0, at), machine: address.slice(at + 1) }; +} + +/** + * Resolve an `agent@machine` address to the agent currently hosted there. + * `machine` matches the agent's location node by node name or `machine_id`, + * so a stale address (the agent moved or was released) fails instead of + * silently reaching the agent somewhere else. + */ +export async function resolveAgentAddress(db: Db, workspaceId: string, address: string) { + const parsed = parseAgentAddress(address); + if (!parsed) { + throw codedError('Address must be of the form "agent@machine"', 'invalid_address', 400); + } + + const [row] = await db + .select({ + agentId: agents.id, + agentName: agents.name, + status: agents.status, + nodeId: nodes.id, + nodeName: nodes.name, + machineId: nodes.machineId, + }) + .from(agents) + .leftJoin(nodes, eq(nodes.id, agents.locationNodeId)) + .where(and(eq(agents.workspaceId, workspaceId), eq(agents.name, parsed.agent))); + + if ( + !row + || row.status === 'released' + || (row.nodeName !== parsed.machine && row.machineId !== parsed.machine) + ) { + throw codedError(`No agent at address "${address}"`, 'address_not_found', 404); + } + + return { + address: `${row.agentName}@${parsed.machine}`, + agent_id: row.agentId, + agent_name: row.agentName, + node_id: row.nodeId!, + }; +} diff --git a/packages/engine/src/routes/dm.ts b/packages/engine/src/routes/dm.ts index fa5096d1..9e6cf316 100644 --- a/packages/engine/src/routes/dm.ts +++ b/packages/engine/src/routes/dm.ts @@ -1,4 +1,4 @@ -import { Hono } from 'hono'; +import { Hono, type Context } from 'hono'; import { z } from 'zod'; import type { AppEnv } from '../env.js'; import { requireAgentToken } from '../middleware/auth.js'; @@ -6,6 +6,7 @@ import { rateLimit } from '../middleware/rateLimit.js'; import { jsonIdempotentOk, parseIdempotencyKey, runIdempotent } from '../middleware/idempotency.js'; import { sha256Hex } from '../lib/crypto.js'; import * as dmEngine from '../engine/dm.js'; +import { resolveAgentAddress } from '../engine/address.js'; import { resolveMailboxConfig } from '../engine/mailboxConfig.js'; import { resolveWorkspaceDeliveryPolicyFor } from '../engine/workspaceDeliveryPolicy.js'; import { publishWorkspaceEvent } from './fanout.js'; @@ -32,6 +33,129 @@ const listDmConversationsQuerySchema = z.object({ limit: positiveIntQueryParam({ max: 100 }), }); +// Body of POST /v1/to/:address — the recipient comes from the path. +const sendAddressedSchema = sendDmSchema.omit({ to: true }); + +/** + * Send a DM from the authenticated agent. Shared by POST /v1/dm (recipient by + * name) and POST /v1/to/:address (recipient by `agent@machine`). + */ +async function sendDirectMessage(c: Context, input: z.infer) { + const db = c.get('db'); + const workspace = c.get('workspace'); + const agent = c.get('agent'); + const { to, text, attachments, data, mode } = input; + const normalizedAttachments = attachments && attachments.length > 0 ? attachments : undefined; + // `data` is digested rather than embedded. It is caller-supplied and can + // be large — a Ratify proof bundle runs to MAX_PROOF_BUNDLE_BYTES (128 + // KiB) — and the fingerprint is serialized into the stored idempotency + // record, kept for the TTL, and string-compared on every replay. Inlining + // it put ~256 KiB per DM into the KV record and made each replay compare + // the whole payload. A digest answers the only question the fingerprint + // asks — "is this the same request?" — in constant size. + const fingerprintBody = { + to, + text, + ...(normalizedAttachments ? { attachments: normalizedAttachments } : {}), + ...(data !== undefined ? { data_sha256: await sha256Hex(JSON.stringify(data)) } : {}), + }; + + const { key: idempotencyKey, error: idempotencyError } = parseIdempotencyKey(c.req.header('Idempotency-Key')); + if (idempotencyError) { + return jsonError(c, 'invalid_idempotency_key', idempotencyError, 400); + } + + const mailbox = resolveMailboxConfig(c.get('engine').config, workspace.id); + const toDmReceivedEventData = (data: Awaited>) => buildDmReceivedEventData(data, { + fromName: agent!.name, + }); + + const trackDmSent = (data: { conversation_id: string; id: string }) => emitServerEvent(c, workspace.id, 'relaycast_server_dm_sent', { + conversation_id: data.conversation_id, + message_id: data.id, + from_agent_id: agent!.id, + to_agent_name: to, + }); + + const idempotent = await runIdempotent({ + workspaceId: workspace.id, + actorId: agent!.id, + scope: 'dm:direct', + key: idempotencyKey, + status: 201, + // Backward compatibility: historical fingerprint excluded mode (equivalent to wait). + // Only include mode when explicit steer is requested. + fingerprint: mode === 'steer' + ? JSON.stringify({ ...fingerprintBody, mode }) + : JSON.stringify(fingerprintBody), + kv: c.get('engine').kv, + operation: () => dmEngine.sendDm(db, workspace.id, agent!.id, { + to, + text, + attachments: normalizedAttachments, + data, + mode, + }, { mailbox, resolveWorkspaceDeliveryPolicy: () => resolveWorkspaceDeliveryPolicyFor(c.get('engine').config, workspace), idempotencyKey, + afterAdmission: (data, event) => { + runInBackground(c, c.get('engine').realtime.publishToWorkspaceStream({ + workspaceId: workspace.id, event: { ...event.payload, seq: event.seq }, + }), 'publish admitted dm.received'); + runInBackground(c, c.get('engine').webhookQueue.send({ + type: 'dm.received', workspaceId: workspace.id, + data: event.data, outboxId: event.outboxId, + }), 'queue admitted dm.received'); + if (data._delivery) runInBackground(c, + routeDeliveryOutcomes(c, [data._delivery], 'dm.received', event.data), + 'route admitted dm delivery'); + if (data._delivery_rejections.length) runInBackground(c, + notifyDeliveryRejections(c, agent!.id, data._delivery_rejections), + 'notify admitted dm delivery rejection'); + trackDmSent(data); + }, + }), + afterOperation: async (data) => { + if (data._notifications_durable) return; + await sendWebhookEvent(c, { + type: 'dm.received', + workspaceId: workspace.id, + data: toDmReceivedEventData(data), + }); + }, + }); + + if (!idempotent.replayed && !idempotent.data._notifications_durable) { + const { + _delivery, + _delivery_rejections, + ...publicDmData + } = idempotent.data as typeof idempotent.data & { + _delivery?: Parameters[1][number] | null; + _delivery_rejections?: Parameters[2]; + }; + const eventData = toDmReceivedEventData(idempotent.data); + runInBackground(c, publishWorkspaceEvent(c, 'dm.received', eventData), 'publish dm.received'); + + if (_delivery) { + runInBackground( + c, + routeDeliveryOutcomes(c, [_delivery], 'dm.received', eventData), + 'route dm delivery', + ); + } + if (_delivery_rejections && _delivery_rejections.length > 0) { + runInBackground( + c, + notifyDeliveryRejections(c, agent!.id, _delivery_rejections), + 'fanout delivery rejected', + ); + } + + trackDmSent(publicDmData); + } + + return jsonIdempotentOk(c, idempotent); +} + // POST /v1/dm - send a DM dmRoutes.post( '/dm', @@ -39,9 +163,6 @@ dmRoutes.post( rateLimit, async (c) => { try { - const db = c.get('db'); - const workspace = c.get('workspace'); - const agent = c.get('agent'); const parsed = await parseJsonBody(c, sendDmSchema, (failure) => { const hasToIssue = failure.error.issues.some((issue) => issue.path[0] === 'to'); const hasTextIssue = failure.error.issues.some((issue) => issue.path[0] === 'text'); @@ -54,116 +175,26 @@ dmRoutes.post( if (!parsed.ok) { return parsed.response; } - const { to, text, attachments, data, mode } = parsed.data; - const normalizedAttachments = attachments && attachments.length > 0 ? attachments : undefined; - // `data` is digested rather than embedded. It is caller-supplied and can - // be large — a Ratify proof bundle runs to MAX_PROOF_BUNDLE_BYTES (128 - // KiB) — and the fingerprint is serialized into the stored idempotency - // record, kept for the TTL, and string-compared on every replay. Inlining - // it put ~256 KiB per DM into the KV record and made each replay compare - // the whole payload. A digest answers the only question the fingerprint - // asks — "is this the same request?" — in constant size. - const fingerprintBody = { - to, - text, - ...(normalizedAttachments ? { attachments: normalizedAttachments } : {}), - ...(data !== undefined ? { data_sha256: await sha256Hex(JSON.stringify(data)) } : {}), - }; - - const { key: idempotencyKey, error: idempotencyError } = parseIdempotencyKey(c.req.header('Idempotency-Key')); - if (idempotencyError) { - return jsonError(c, 'invalid_idempotency_key', idempotencyError, 400); - } - - const mailbox = resolveMailboxConfig(c.get('engine').config, workspace.id); - const toDmReceivedEventData = (data: Awaited>) => buildDmReceivedEventData(data, { - fromName: agent!.name, - }); - - const trackDmSent = (data: { conversation_id: string; id: string }) => emitServerEvent(c, workspace.id, 'relaycast_server_dm_sent', { - conversation_id: data.conversation_id, - message_id: data.id, - from_agent_id: agent!.id, - to_agent_name: to, - }); - - const idempotent = await runIdempotent({ - workspaceId: workspace.id, - actorId: agent!.id, - scope: 'dm:direct', - key: idempotencyKey, - status: 201, - // Backward compatibility: historical fingerprint excluded mode (equivalent to wait). - // Only include mode when explicit steer is requested. - fingerprint: mode === 'steer' - ? JSON.stringify({ ...fingerprintBody, mode }) - : JSON.stringify(fingerprintBody), - kv: c.get('engine').kv, - operation: () => dmEngine.sendDm(db, workspace.id, agent!.id, { - to, - text, - attachments: normalizedAttachments, - data, - mode, - }, { mailbox, resolveWorkspaceDeliveryPolicy: () => resolveWorkspaceDeliveryPolicyFor(c.get('engine').config, workspace), idempotencyKey, - afterAdmission: (data, event) => { - runInBackground(c, c.get('engine').realtime.publishToWorkspaceStream({ - workspaceId: workspace.id, event: { ...event.payload, seq: event.seq }, - }), 'publish admitted dm.received'); - runInBackground(c, c.get('engine').webhookQueue.send({ - type: 'dm.received', workspaceId: workspace.id, - data: event.data, outboxId: event.outboxId, - }), 'queue admitted dm.received'); - if (data._delivery) runInBackground(c, - routeDeliveryOutcomes(c, [data._delivery], 'dm.received', event.data), - 'route admitted dm delivery'); - if (data._delivery_rejections.length) runInBackground(c, - notifyDeliveryRejections(c, agent!.id, data._delivery_rejections), - 'notify admitted dm delivery rejection'); - trackDmSent(data); - }, - }), - afterOperation: async (data) => { - if (data._notifications_durable) return; - await sendWebhookEvent(c, { - type: 'dm.received', - workspaceId: workspace.id, - data: toDmReceivedEventData(data), - }); - }, - }); + return await sendDirectMessage(c, parsed.data); + } catch (err: unknown) { + return errorResponse(c, err); + } + }, +); - if (!idempotent.replayed && !idempotent.data._notifications_durable) { - const { - _delivery, - _delivery_rejections, - ...publicDmData - } = idempotent.data as typeof idempotent.data & { - _delivery?: Parameters[1][number] | null; - _delivery_rejections?: Parameters[2]; - }; - const eventData = toDmReceivedEventData(idempotent.data); - runInBackground(c, publishWorkspaceEvent(c, 'dm.received', eventData), 'publish dm.received'); - - if (_delivery) { - runInBackground( - c, - routeDeliveryOutcomes(c, [_delivery], 'dm.received', eventData), - 'route dm delivery', - ); - } - if (_delivery_rejections && _delivery_rejections.length > 0) { - runInBackground( - c, - notifyDeliveryRejections(c, agent!.id, _delivery_rejections), - 'fanout delivery rejected', - ); - } - - trackDmSent(publicDmData); +// POST /v1/to/:address - send a DM to an `agent@machine` address +dmRoutes.post( + '/to/:address', + requireAgentToken, + rateLimit, + async (c) => { + try { + const parsed = await parseJsonBody(c, sendAddressedSchema, 'text is required'); + if (!parsed.ok) { + return parsed.response; } - - return jsonIdempotentOk(c, idempotent); + const target = await resolveAgentAddress(c.get('db'), c.get('workspace').id, c.req.param('address')); + return await sendDirectMessage(c, { ...parsed.data, to: target.agent_name }); } catch (err: unknown) { return errorResponse(c, err); } diff --git a/packages/sdk-typescript/CHANGELOG.md b/packages/sdk-typescript/CHANGELOG.md index eb3ca66a..e7ae7865 100644 --- a/packages/sdk-typescript/CHANGELOG.md +++ b/packages/sdk-typescript/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 - Minor] + +### Added + +- `agent.sendTo(address, text)` sends a DM to an `agent@machine` address via `POST /v1/to/:address`. ## [8.12.0] - 2026-09-24 diff --git a/packages/sdk-typescript/src/__tests__/agent-messaging.test.ts b/packages/sdk-typescript/src/__tests__/agent-messaging.test.ts index 24c0e863..69e91014 100644 --- a/packages/sdk-typescript/src/__tests__/agent-messaging.test.ts +++ b/packages/sdk-typescript/src/__tests__/agent-messaging.test.ts @@ -252,6 +252,19 @@ describe('AgentClient', () => { }); }); + describe('sendTo()', () => { + it('sends DM via POST /v1/to/:address', async () => { + mockFetch.mockImplementation(() => mockResponse({ id: 'dm_1' })); + + await me.sendTo('Worker-1@laptop', 'hi'); + + const [url, init] = mockFetch.mock.calls[0]!; + expect(url).toBe('https://cast.agentrelay.com/v1/to/Worker-1%40laptop'); + expect(init.method).toBe('POST'); + expect(init.body).toBe(JSON.stringify({ text: 'hi', mode: 'wait' })); + }); + }); + describe('dm()', () => { it('sends DM via POST /v1/dm', async () => { mockFetch.mockImplementation(() => mockResponse({ id: 'dm_1' })); diff --git a/packages/sdk-typescript/src/agent.ts b/packages/sdk-typescript/src/agent.ts index a43e1b78..f671a3fb 100644 --- a/packages/sdk-typescript/src/agent.ts +++ b/packages/sdk-typescript/src/agent.ts @@ -616,6 +616,25 @@ export class AgentClient { return this.client.post('/v1/dm', body, idempotencyHeaders(opts)); } + /** DM an `agent@machine` address; fails if the agent is no longer on that machine. */ + async sendTo( + address: string, + text: string, + opts?: (IdempotencyOption & { + mode?: 'wait' | 'steer'; + attachments?: string[]; + data?: Record | null; + }), + ): Promise { + const body = { + text, + ...(opts?.attachments ? { attachments: opts.attachments } : {}), + ...(opts?.data !== undefined ? { data: opts.data } : {}), + mode: opts?.mode ?? 'wait', + }; + return this.client.post(`/v1/to/${encodeURIComponent(address)}`, body, idempotencyHeaders(opts)); + } + dms = { conversations: (opts?: Pick): Promise => { const query: Record = {}; From eb6c8cb2f73d1a43abb05993c301ca396e401005 Mon Sep 17 00:00:00 2001 From: khaliqgant Date: Thu, 24 Sep 2026 20:34:37 -0700 Subject: [PATCH 2/8] feat(engine): harden agent@machine addressing - agent@direct addresses agents not hosted on a broker - the address is checked inside the idempotent operation on the same recipient row the DM is sent to: accepted retries replay after the agent moves, and there is no resolve-then-send window - the idempotency fingerprint includes the address, so reusing a key for a different address (or /v1/dm) is a 409 - agent resources expose address; DMs persist the sender's address as server-owned metadata and expose it as message.agent_address on the send response, live and redelivered dm.received, and DM history Co-Authored-By: Claude Opus 5.5 (1M context) --- CHANGELOG.md | 1 + README.md | 9 +- openapi.yaml | 28 +- packages/engine/CHANGELOG.md | 3 +- .../conformance/addressedSend.test.ts | 242 +++++++++++++++--- packages/engine/src/engine/address.ts | 90 ++++--- packages/engine/src/engine/agent.ts | 16 +- packages/engine/src/engine/delivery.ts | 2 + packages/engine/src/engine/dm.ts | 75 ++++-- packages/engine/src/engine/wsTransform.ts | 1 + packages/engine/src/routes/dm.ts | 17 +- packages/types/CHANGELOG.md | 6 +- packages/types/src/agent.ts | 2 + packages/types/src/message.ts | 2 + 14 files changed, 377 insertions(+), 117 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index 5f8a4ca5..c1a70fc6 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -21,6 +21,7 @@ Packages without a separate changelog are covered by the cross-package notes bel ### Added - `POST /v1/to/:address` sends a DM to an `agent@machine` address, delivering to the agent only while it is hosted on that machine. +- Agents report their `address`, and received DMs carry the sender's as `message.agent_address`, so agents can reply by address. ## [8.12.0] - 2026-09-24 diff --git a/README.md b/README.md index e5e750e3..d071cd75 100644 --- a/README.md +++ b/README.md @@ -626,9 +626,12 @@ Activity feed channel-message items include `channel_id` and `channel_name`; DM `conversation_id`. `POST /to/:address` routes by address instead of bare name: `agent@machine` resolves to the -agent only while it is hosted on that machine (its node's name or `machine_id`), then delivers -like `POST /dm`. A stale address returns `404 address_not_found`; a malformed one returns -`400 invalid_address`. +agent only while it is hosted on that machine (its broker node's name or `machine_id`, or +`direct` for an agent not on a broker), then delivers like `POST /dm`. Agents expose their +address as `address` on agent resources, and each DM carries the sender's as +`message.agent_address`, so a recipient can reply on it. A stale address returns +`404 address_not_found`; a malformed one returns `400 invalid_address`. An idempotent retry +replays even if the agent has moved; reusing the key for another address is a `409`. `POST /dm` can return **`409 dm_conversation_id_collision`**. A 1:1 conversation id is derived deterministically from `(workspace, sorted agent pair)`, and that binding is reserved diff --git a/openapi.yaml b/openapi.yaml index 6887d905..b3707f3e 100644 --- a/openapi.yaml +++ b/openapi.yaml @@ -320,6 +320,11 @@ components: locally) can record which workspace it registered into. name: type: string + address: + type: string + description: >- + `agent@machine` for `POST /to/{address}`. `machine` is the broker + node's name, or `direct` when the agent is not hosted on a broker. type: type: string enum: [agent, human, system] @@ -804,6 +809,9 @@ components: type: string enum: [agent, human, system] description: Sender identity type for the DM message actor + agent_address: + type: string + description: Sender's `agent@machine` address at send time; reply via `POST /to/{address}` text: type: string injection_mode: @@ -863,6 +871,9 @@ components: type: string enum: [agent, human, system] description: Sender identity type for the message actor + agent_address: + type: string + description: Sender's `agent@machine` address at send time; reply via `POST /to/{address}` text: type: string injection_mode: @@ -3918,10 +3929,19 @@ paths: summary: Send DM to agent@machine description: | Send a direct message to an `agent@machine` address. `machine` must be - the agent's current node, matched by node name or `machine_id`; the - message is then delivered exactly like `POST /dm`. An address that no - longer matches (the agent moved or was released) returns 404 instead - of reaching the agent elsewhere. + the agent's current broker node, matched by node name or `machine_id`, + or `direct` for an agent not hosted on a broker. The message is then + delivered exactly like `POST /dm`. An address that no longer matches + (the agent moved or was released) returns 404 instead of reaching the + agent elsewhere; matching is exact and case-sensitive. + + Agents learn addresses from `address` on agent resources and from + `message.agent_address` on DMs they receive. + + With an `Idempotency-Key`, a retry of an accepted send replays the + original result even if the agent has since moved. The key is bound + to the address: reusing it for a different address or for `POST /dm` + returns `409 idempotency_key_reused`. tags: - Direct Messages security: diff --git a/packages/engine/CHANGELOG.md b/packages/engine/CHANGELOG.md index a6f94590..e2e2c4d5 100644 --- a/packages/engine/CHANGELOG.md +++ b/packages/engine/CHANGELOG.md @@ -11,7 +11,8 @@ and this project follows [Semantic Versioning](https://semver.org/spec/v2.0.0.ht ### Added -- `POST /v1/to/:address` resolves `agent@machine` (node name or `machine_id`) to the agent hosted there and sends it a DM; stale addresses return `404 address_not_found`. +- `POST /v1/to/:address` resolves `agent@machine` (broker node name or `machine_id`, or `direct` for agents not on a broker) to the agent hosted there and sends it a DM; stale addresses return `404 address_not_found`, and idempotent retries replay even after the agent moves. +- Agent resources include `address`; DM responses, `dm.received` deliveries, and DM history include the sender's `message.agent_address`. ## [8.12.0] - 2026-09-24 diff --git a/packages/engine/src/__tests__/conformance/addressedSend.test.ts b/packages/engine/src/__tests__/conformance/addressedSend.test.ts index 468e1af3..30e94ed5 100644 --- a/packages/engine/src/__tests__/conformance/addressedSend.test.ts +++ b/packages/engine/src/__tests__/conformance/addressedSend.test.ts @@ -1,10 +1,15 @@ import { afterEach, beforeEach, describe, expect, it } from 'vitest'; +import { eq } from 'drizzle-orm'; import { makeNodeStack, createWorkspace, registerAgent, FakeSocket, type TestStack } from './harness.js'; -import { parseAgentAddress } from '../../engine/address.js'; +import { agents, messages } from '../../db/schema.js'; +import { parseAgentAddress, SENDER_ADDRESS_METADATA_KEY } from '../../engine/address.js'; + +type Json = Record; /** * POST /v1/to/:address — send a DM to `agent@machine`, where `machine` is the - * agent's current node, by node name or machine_id. + * agent's current broker node (by node name or machine_id), or `direct` for an + * agent not hosted on a broker. */ describe('addressed send', () => { let stack: TestStack; @@ -13,22 +18,13 @@ describe('addressed send', () => { await stack.close(); }); - async function seed() { - const ws = await createWorkspace(stack.app, 'address-ws'); - const alice = await registerAgent(stack.app, ws.workspaceKey, 'alice'); - - const nodeId = 'node_laptop'; - const nodeName = 'laptop'; + async function enrollBroker(ws: { workspaceKey: string; workspaceId: string }, nodeId: string, name: string, machineId: string) { const enroll = await stack.app.request('/v1/nodes', { method: 'POST', headers: { 'content-type': 'application/json', authorization: `Bearer ${ws.workspaceKey}` }, body: JSON.stringify({ - node_id: nodeId, - name: nodeName, - machine_id: 'mach-123', - capabilities: ['spawn:claude'], - max_agents: 4, - version: 'test-node', + node_id: nodeId, name, machine_id: machineId, + capabilities: ['spawn:claude'], max_agents: 4, version: 'test-node', }), }); expect(enroll.status).toBe(201); @@ -36,30 +32,53 @@ describe('addressed send', () => { const sock = new FakeSocket(); const handle = stack.runtime.realtime.attachNodeSocket(ws.workspaceId, nodeId, sock); await handle.handleMessage(JSON.stringify({ - v: 1, type: 'node.register', name: nodeName, node_id: nodeId, + v: 1, type: 'node.register', name, node_id: nodeId, capabilities: [{ name: 'spawn:claude', kind: 'capacity' }], max_agents: 4, tags: [], version: 'test-node', resume_cursor: null, })); await handle.handleMessage(JSON.stringify({ v: 1, type: 'node.heartbeat', load: 0, active_agents: 0, handlers_live: true, })); - await handle.handleMessage(JSON.stringify({ + return { sock, handle }; + } + + /** alice is self-connected; bob runs on broker node `laptop` (machine_id `mach-123`). */ + async function seed() { + const ws = await createWorkspace(stack.app, 'address-ws'); + const alice = await registerAgent(stack.app, ws.workspaceKey, 'alice'); + const laptop = await enrollBroker(ws, 'node_laptop', 'laptop', 'mach-123'); + await laptop.handle.handleMessage(JSON.stringify({ v: 1, type: 'agent.register', name: 'bob', resumable: true, session_ref: 'sess-bob', })); - const reply = sock.ofType('reply').at(-1) as { ok: boolean }; + const reply = laptop.sock.ofType('reply').at(-1) as { ok: boolean; data: { agent_id: string; token: string } }; expect(reply?.ok).toBe(true); - - return { ws, alice, sock }; + const bob = { agentId: reply.data.agent_id, token: reply.data.token }; + return { ws, alice, bob, laptop }; } - function send(token: string, address: string, body: unknown) { + function send(token: string, address: string, body: unknown, headers: Record = {}) { return stack.app.request(`/v1/to/${encodeURIComponent(address)}`, { method: 'POST', - headers: { 'content-type': 'application/json', authorization: `Bearer ${token}` }, + headers: { 'content-type': 'application/json', authorization: `Bearer ${token}`, ...headers }, body: JSON.stringify(body), }); } + async function errorCode(res: Response) { + return ((await res.json()) as { error: { code: string } }).error.code; + } + + function deliveredDms(sock: FakeSocket) { + return sock.ofType('deliver') + .map((frame) => (frame as { payload?: { type?: string; data?: { message?: Json } } }).payload) + .filter((payload) => payload?.type === 'dm.received') + .map((payload) => payload!.data!.message!); + } + + async function moveAgent(agentId: string, nodeId: string) { + await stack.runtime.deps.db.update(agents).set({ locationNodeId: nodeId }).where(eq(agents.id, agentId)); + } + it('parses agent@machine on the last @', () => { expect(parseAgentAddress('bob@laptop')).toEqual({ agent: 'bob', machine: 'laptop' }); expect(parseAgentAddress('a@b@laptop')).toEqual({ agent: 'a@b', machine: 'laptop' }); @@ -68,8 +87,8 @@ describe('addressed send', () => { expect(parseAgentAddress('bob@')).toBeNull(); }); - it('routes to the agent by node name and by machine_id', async () => { - const { alice, sock } = await seed(); + it('routes to the agent by node name and by machine_id, delivering on that machine', async () => { + const { alice, laptop } = await seed(); for (const address of ['bob@laptop', 'bob@mach-123']) { const res = await send(alice.token, address, { text: `hi via ${address}` }); @@ -80,33 +99,172 @@ describe('addressed send', () => { } await stack.settle(); - const delivered = sock.ofType('deliver') - .map((frame: { payload?: { type?: string; data?: { message?: { text?: string } } } }) => frame.payload) - .filter((payload) => payload?.type === 'dm.received') - .map((payload) => payload?.data?.message?.text); - expect(delivered).toEqual(['hi via bob@laptop', 'hi via bob@mach-123']); + expect(deliveredDms(laptop.sock).map((message) => message.text)) + .toEqual(['hi via bob@laptop', 'hi via bob@mach-123']); }); - it('rejects an address whose machine does not host the agent', async () => { + it('accepts an unencoded @ in the path', async () => { const { alice } = await seed(); - const res = await send(alice.token, 'bob@desktop', { text: 'hi' }); - expect(res.status).toBe(404); - expect(((await res.json()) as { error: { code: string } }).error.code).toBe('address_not_found'); + const res = await stack.app.request('/v1/to/bob@laptop', { + method: 'POST', + headers: { 'content-type': 'application/json', authorization: `Bearer ${alice.token}` }, + body: JSON.stringify({ text: 'raw' }), + }); + expect(res.status).toBe(201); }); - it('rejects unknown agents and malformed addresses', async () => { - const { alice } = await seed(); - const unknown = await send(alice.token, 'carol@laptop', { text: 'hi' }); - expect(unknown.status).toBe(404); + it('addresses agents without a broker as agent@direct', async () => { + const { alice, bob } = await seed(); + expect((await send(bob.token, 'alice@direct', { text: 'to a self-connected agent' })).status).toBe(201); + + // `direct` never matches an agent that is on a broker. + expect(await errorCode(await send(alice.token, 'bob@direct', { text: 'x' }))).toBe('address_not_found'); + // A self-connected agent is not on any named machine. + expect(await errorCode(await send(bob.token, 'alice@laptop', { text: 'x' }))).toBe('address_not_found'); + }); + + it('rejects addresses that do not match the agent, without revealing which part failed', async () => { + const { ws, alice } = await seed(); + for (const address of ['bob@desktop', 'carol@laptop', 'Bob@laptop', 'bob@Laptop', ' bob@laptop']) { + const res = await send(alice.token, address, { text: 'x' }); + expect(res.status, address).toBe(404); + expect(await errorCode(res)).toBe('address_not_found'); + } - const malformed = await send(alice.token, 'bob', { text: 'hi' }); - expect(malformed.status).toBe(400); - expect(((await malformed.json()) as { error: { code: string } }).error.code).toBe('invalid_address'); + const deleted = await stack.app.request('/v1/agents/bob', { + method: 'DELETE', headers: { authorization: `Bearer ${ws.workspaceKey}` }, + }); + expect(deleted.status).toBe(204); + expect(await errorCode(await send(alice.token, 'bob@laptop', { text: 'x' }))).toBe('address_not_found'); }); - it('requires text', async () => { + it('rejects malformed addresses and missing text', async () => { const { alice } = await seed(); - const res = await send(alice.token, 'bob@laptop', {}); - expect(res.status).toBe(400); + for (const address of ['bob', '@laptop', 'bob@']) { + const res = await send(alice.token, address, { text: 'x' }); + expect(res.status, address).toBe(400); + expect(await errorCode(res)).toBe('invalid_address'); + } + expect((await send(alice.token, 'bob@laptop', {})).status).toBe(400); + }); + + it('follows the agent when it moves: the old address fails, the new one routes', async () => { + const { ws, alice, bob } = await seed(); + const desktop = await enrollBroker(ws, 'node_desktop', 'desktop', 'mach-456'); + await moveAgent(bob.agentId, 'node_desktop'); + + expect(await errorCode(await send(alice.token, 'bob@laptop', { text: 'stale' }))).toBe('address_not_found'); + expect((await send(alice.token, 'bob@desktop', { text: 'moved' })).status).toBe(201); + await stack.settle(); + expect(deliveredDms(desktop.sock).map((message) => message.text)).toEqual(['moved']); + }); + + it('queues for an offline machine and redelivers with the sender address on reconnect', async () => { + const { ws, alice, bob, laptop } = await seed(); + await laptop.handle.handleClose(); + expect((await send(alice.token, 'bob@laptop', { text: 'while offline' })).status).toBe(201); + + const reconnected = await enrollBroker(ws, 'node_laptop', 'laptop', 'mach-123'); + await reconnected.handle.handleMessage(JSON.stringify({ + v: 1, type: 'inventory.sync', agents: [{ agent_id: bob.agentId, name: 'bob', session_ref: 'sess-bob' }], + })); + await stack.settle(); + expect(deliveredDms(reconnected.sock)).toEqual([ + expect.objectContaining({ text: 'while offline', agent_address: 'alice@direct' }), + ]); + }); + + describe('idempotent retries', () => { + it('replays an accepted send even after the agent moves', async () => { + const { ws, alice, bob } = await seed(); + const key = { 'Idempotency-Key': 'retry-1' }; + const first = await send(alice.token, 'bob@laptop', { text: 'once' }, key); + expect(first.status).toBe(201); + const firstId = ((await first.json()) as { data: { id: string } }).data.id; + + await enrollBroker(ws, 'node_desktop', 'desktop', 'mach-456'); + await moveAgent(bob.agentId, 'node_desktop'); + + const retry = await send(alice.token, 'bob@laptop', { text: 'once' }, key); + expect(retry.status).toBe(201); + expect(((await retry.json()) as { data: { id: string } }).data.id).toBe(firstId); + + const sent = await stack.runtime.deps.db.select().from(messages).where(eq(messages.body, 'once')); + expect(sent).toHaveLength(1); + }); + + it('rejects reusing a key for a different address or for /v1/dm', async () => { + const { alice } = await seed(); + const key = { 'Idempotency-Key': 'retry-2' }; + expect((await send(alice.token, 'bob@laptop', { text: 'same' }, key)).status).toBe(201); + + const otherAddress = await send(alice.token, 'bob@mach-123', { text: 'same' }, key); + expect(otherAddress.status).toBe(409); + expect(await errorCode(otherAddress)).toBe('idempotency_key_reused'); + + const viaDm = await stack.app.request('/v1/dm', { + method: 'POST', + headers: { 'content-type': 'application/json', authorization: `Bearer ${alice.token}`, ...key }, + body: JSON.stringify({ to: 'bob', text: 'same' }), + }); + expect(viaDm.status).toBe(409); + }); + + it('does not record a failed resolution, so a later send with the key can succeed', async () => { + const { alice } = await seed(); + const key = { 'Idempotency-Key': 'retry-3' }; + expect((await send(alice.token, 'bob@desktop', { text: 'x' }, key)).status).toBe(404); + // Nothing was accepted under the key yet, so a valid address is not a reuse conflict. + expect((await send(alice.token, 'bob@laptop', { text: 'x' }, key)).status).toBe(201); + }); + }); + + describe('address discovery', () => { + it('reports each agent\'s address on agent resources', async () => { + const { ws, bob } = await seed(); + const auth = { authorization: `Bearer ${ws.workspaceKey}` }; + + const list = (await (await stack.app.request('/v1/agents', { headers: auth })).json()) as { data: Json[] }; + expect(Object.fromEntries(list.data.map((agent) => [agent.name, agent.address]))) + .toMatchObject({ alice: 'alice@direct', bob: 'bob@laptop' }); + + const one = (await (await stack.app.request('/v1/agents/bob', { headers: auth })).json()) as { data: Json }; + expect(one.data.address).toBe('bob@laptop'); + + const self = (await (await stack.app.request('/v1/agent', { + headers: { authorization: `Bearer ${bob.token}` }, + })).json()) as { data: Json }; + expect(self.data.address).toBe('bob@laptop'); + }); + + it('carries the sender address on the DM response, live delivery, and history, and it round-trips', async () => { + const { alice, bob, laptop } = await seed(); + const res = await send(alice.token, 'bob@laptop', { text: 'reply to me' }); + const sent = (await res.json()) as { data: { conversation_id: string; message: Json } }; + expect(sent.data.message.agent_address).toBe('alice@direct'); + + await stack.settle(); + const [delivered] = deliveredDms(laptop.sock); + expect(delivered.agent_address).toBe('alice@direct'); + + const history = (await (await stack.app.request(`/v1/dm/${sent.data.conversation_id}/messages`, { + headers: { authorization: `Bearer ${bob.token}` }, + })).json()) as { data: Json[] }; + expect(history.data[0].agent_address).toBe('alice@direct'); + + // The recipient can answer on the address it was handed. + expect((await send(bob.token, delivered.agent_address as string, { text: 'got it' })).status).toBe(201); + }); + + it('keeps the server-owned address out of public metadata and ignores a caller-supplied one', async () => { + const { alice } = await seed(); + const res = await send(alice.token, 'bob@laptop', { + text: 'spoof', + data: { [SENDER_ADDRESS_METADATA_KEY]: 'mallory@elsewhere', topic: 'x' }, + }); + const message = ((await res.json()) as { data: { message: Json } }).data.message; + expect(message.agent_address).toBe('alice@direct'); + expect(message.metadata).toEqual({ topic: 'x', injection_mode: 'wait' }); + }); }); }); diff --git a/packages/engine/src/engine/address.ts b/packages/engine/src/engine/address.ts index 0f975ec2..709e7848 100644 --- a/packages/engine/src/engine/address.ts +++ b/packages/engine/src/engine/address.ts @@ -1,15 +1,22 @@ -import { and, eq } from 'drizzle-orm'; -import type { getDb } from '../db/index.js'; -import { agents, nodes } from '../db/schema.js'; +import { nodes } from '../db/schema.js'; import { codedError } from '../lib/httpError.js'; -type Db = ReturnType; +/** + * Machine name for agents not hosted on a broker: self-connected agents and + * agents on an implicit direct node, whose node name is an internal id. + */ +export const DIRECT_MACHINE = 'direct'; + +/** Server-owned message metadata key holding the sender's address at send time. */ +export const SENDER_ADDRESS_METADATA_KEY = '__relaycast_sender_address'; export interface AgentAddress { agent: string; machine: string; } +type AddressNode = { name: string; role: string; machineId: string | null } | null; + /** * Parse an `agent@machine` address. The split is on the last `@` so agent * names that themselves contain `@` still resolve. @@ -20,43 +27,54 @@ export function parseAgentAddress(address: string): AgentAddress | null { return { agent: address.slice(0, at), machine: address.slice(at + 1) }; } -/** - * Resolve an `agent@machine` address to the agent currently hosted there. - * `machine` matches the agent's location node by node name or `machine_id`, - * so a stale address (the agent moved or was released) fails instead of - * silently reaching the agent somewhere else. - */ -export async function resolveAgentAddress(db: Db, workspaceId: string, address: string) { +/** Canonical address: the broker node's name, or `direct` when there is no broker. */ +export function formatAgentAddress(agentName: string, node: AddressNode): string { + const machine = node && node.role !== 'direct' ? node.name : DIRECT_MACHINE; + return `${agentName}@${machine}`; +} + +/** Whether `machine` names the node the agent is on: node name, machine_id, or `direct`. */ +function machineMatches(machine: string, node: AddressNode): boolean { + if (machine === DIRECT_MACHINE && (!node || node.role === 'direct')) return true; + return node !== null && (machine === node.name || machine === node.machineId); +} + +/** Node columns needed to format or match an address; select them via a left join on the agent's location node. */ +export const addressNodeSelection = { name: nodes.name, role: nodes.role, machineId: nodes.machineId }; + +/** Reject an address that is not `agent@machine` before any other work. */ +export function requireAgentAddress(address: string): AgentAddress { const parsed = parseAgentAddress(address); if (!parsed) { throw codedError('Address must be of the form "agent@machine"', 'invalid_address', 400); } + return parsed; +} - const [row] = await db - .select({ - agentId: agents.id, - agentName: agents.name, - status: agents.status, - nodeId: nodes.id, - nodeName: nodes.name, - machineId: nodes.machineId, - }) - .from(agents) - .leftJoin(nodes, eq(nodes.id, agents.locationNodeId)) - .where(and(eq(agents.workspaceId, workspaceId), eq(agents.name, parsed.agent))); - - if ( - !row - || row.status === 'released' - || (row.nodeName !== parsed.machine && row.machineId !== parsed.machine) - ) { - throw codedError(`No agent at address "${address}"`, 'address_not_found', 404); +/** + * Throw unless `agent` (looked up by the address's agent name) is live and + * hosted on the address's machine. A stale address (the agent moved or was + * released) fails instead of silently reaching the agent somewhere else. + */ +export function assertAgentAtAddress( + address: string, + agent: { status: string } | undefined, + node: AddressNode, +): void { + const parsed = requireAgentAddress(address); + if (!agent || agent.status === 'released' || !machineMatches(parsed.machine, node)) { + throw addressNotFound(address); } +} + +function addressNotFound(address: string) { + return codedError(`No agent at address "${address}"`, 'address_not_found', 404); +} - return { - address: `${row.agentName}@${parsed.machine}`, - agent_id: row.agentId, - agent_name: row.agentName, - node_id: row.nodeId!, - }; +/** `{ agent_address }` for a DM message whose persisted metadata carries the sender's address, else `{}`. */ +export function senderAddressField( + metadata: Record | null | undefined, +): { agent_address?: string } { + const value = metadata?.[SENDER_ADDRESS_METADATA_KEY]; + return typeof value === 'string' ? { agent_address: value } : {}; } diff --git a/packages/engine/src/engine/agent.ts b/packages/engine/src/engine/agent.ts index b3007d8b..306787ff 100644 --- a/packages/engine/src/engine/agent.ts +++ b/packages/engine/src/engine/agent.ts @@ -7,6 +7,7 @@ import { invalidateChannelCache } from './cache.js'; import { queryInChunks } from '../lib/queryChunks.js'; import { codedError } from '../lib/httpError.js'; import { directNodeIdForAgent } from './node.js'; +import { formatAgentAddress } from './address.js'; import { runAtomicWrites, type AtomicWrite } from '../ports/database.js'; import { AGENT_RECOVERY_PROOF_HASH_PATTERN } from '@relaycast/types'; @@ -386,16 +387,18 @@ export async function listAgents(db: Db, workspaceId: string, status?: string) { } const rows = await db - .select() + .select({ agent: agents, node: { name: nodes.name, role: nodes.role, machineId: nodes.machineId } }) .from(agents) + .leftJoin(nodes, eq(nodes.id, agents.locationNodeId)) // Released rows are tombstones retained only to keep history attributable; // they are not roster members, so `agent list` must not fill with them. .where(and(...conditions)); - return rows.map((a) => ({ + return rows.map(({ agent: a, node }) => ({ id: a.id, name: a.name, handle: `@${a.name}`, + address: formatAgentAddress(a.name, node), type: a.type, status: effectiveAgentStatus(a, now), persona: a.persona, @@ -407,12 +410,14 @@ export async function listAgents(db: Db, workspaceId: string, status?: string) { } export async function getAgentByName(db: Db, workspaceId: string, name: string) { - const [agent] = await db - .select() + const [row] = await db + .select({ agent: agents, node: { name: nodes.name, role: nodes.role, machineId: nodes.machineId } }) .from(agents) + .leftJoin(nodes, eq(nodes.id, agents.locationNodeId)) .where(and(eq(agents.workspaceId, workspaceId), eq(agents.name, name))); - if (!agent) return null; + if (!row) return null; + const { agent, node } = row; // Get channels, actions, and pending deliveries in parallel const [memberships, allActions, pendingDeliveryRows] = await Promise.all([ @@ -478,6 +483,7 @@ export async function getAgentByName(db: Db, workspaceId: string, name: string) workspace_id: workspaceId, name: agent.name, handle: agent.handle ?? `@${agent.name}`, + address: formatAgentAddress(agent.name, node), type: agent.type, status: effectiveAgentStatus(agent), persona: agent.persona, diff --git a/packages/engine/src/engine/delivery.ts b/packages/engine/src/engine/delivery.ts index b65cfdc8..fd26fe89 100644 --- a/packages/engine/src/engine/delivery.ts +++ b/packages/engine/src/engine/delivery.ts @@ -8,6 +8,7 @@ import type { NodeConnectionRegistry } from '../ports/realtime.js'; import { isProviderAgentDeliveryReady } from '../ports/realtime.js'; import { buildDeliverFrame, buildDeliverPayload, buildMessageCreatedEventData, buildThreadReplyEventData, buildDmReceivedEventData, buildGroupDmReceivedEventData } from './deliveryWire.js'; import { displayAgentName, publicMessageMetadata } from './messageMetadata.js'; +import { senderAddressField } from './address.js'; import { toIso } from '../lib/serialize.js'; import { readNodeRedriveCandidates } from './nodeRedriveCandidates.js'; import { fetchAttachmentsBatch, type AttachmentRow } from './attachments.js'; @@ -595,6 +596,7 @@ function buildRoutableDeliveryEvent( id: row.delivery.messageId, agent_id: row.senderAgentId, agent_name: senderName, + ...senderAddressField(row.metadata as Record | null), text: row.body, injection_mode: injectionMode, attachments, diff --git a/packages/engine/src/engine/dm.ts b/packages/engine/src/engine/dm.ts index 5bad56b3..0cf757db 100644 --- a/packages/engine/src/engine/dm.ts +++ b/packages/engine/src/engine/dm.ts @@ -16,6 +16,7 @@ import { a2aAgents, pendingEvents, messageLogs, + nodes, } from '../db/schema.js'; import { sha256Hex } from '../lib/crypto.js'; import { runAtomicWrites, databaseConstraintKind, type AtomicWrite } from '../ports/database.js'; @@ -36,6 +37,13 @@ import { buildDmReceivedEventData } from './deliveryWire.js'; import { buildWorkspaceEventWrite } from './workspaceEvents.js'; import { transformForClient } from './wsTransform.js'; import { codedError } from '../lib/httpError.js'; +import { + addressNodeSelection, + assertAgentAtAddress, + formatAgentAddress, + SENDER_ADDRESS_METADATA_KEY, + senderAddressField, +} from './address.js'; import { buildMessageSessionWrite, requireSessionRefFromMetadata } from './sessionMessages.js'; import { fetchAttachmentsBatch, resolveSendAttachments, type AttachmentRow } from './attachments.js'; import { canonicalUserMessageMetadata, publicMessageMetadata, sanitizeUserMessageMetadata } from './messageMetadata.js'; @@ -62,6 +70,12 @@ interface SendDmOptions { mailbox?: MailboxConfig; /** Server-resolved workspace growth policy; absent => no workspace guard. */ workspaceDeliveryPolicy?: WorkspaceDeliveryPolicy; + /** + * `agent@machine` the caller addressed. Checked against the same recipient + * row the DM is sent to, after durable replay lookup, so an accepted retry + * still replays after the agent moves. + */ + address?: string; } /** @@ -269,6 +283,23 @@ async function resolveConversation( return conv; } +/** + * Persisted DM metadata. Server-owned keys go after caller metadata so a + * federated peer cannot override how the local runtime is injected, and + * `sanitizeUserMessageMetadata` already drops any caller-supplied + * `__relaycast_` key such as the sender address. + */ +function dmMessageMetadata( + data: { mode?: 'wait' | 'steer'; data?: Record | null }, + senderAddress: string | undefined, +): Record { + return { + ...sanitizeUserMessageMetadata(data.data), + injection_mode: data.mode ?? 'wait', + ...(senderAddress ? { [SENDER_ADDRESS_METADATA_KEY]: senderAddress } : {}), + }; +} + /** * Build the message + attachment-junction inserts for a DM without executing * them, so the send path can run them inside one atomic unit. The message @@ -289,14 +320,10 @@ function buildDmMessageWrites( messageId: string, createdAt = new Date(), inboundRegistration?: { id: string; tokenHash?: string }, + senderAddress?: string, ): AtomicWrite[] { const hasAttachments = attachments.length > 0; - const metadata = { - // Keep the server-owned delivery mode after caller metadata so a - // federated peer cannot override how the local runtime is injected. - ...sanitizeUserMessageMetadata(data.data), - injection_mode: data.mode ?? 'wait', - }; + const metadata = dmMessageMetadata(data, senderAddress); const sessionRef = requireSessionRefFromMetadata(metadata); const writes: AtomicWrite[] = [ db @@ -356,6 +383,7 @@ function buildDmResult( id: message.id, agent_id: message.agentId, agent_name: fromAgent.name, + ...senderAddressField(message.metadata), text: message.body, injection_mode: injectionMode, attachments, @@ -504,28 +532,35 @@ export async function sendDm( const workspacePolicy = options.resolveWorkspaceDeliveryPolicy ? await options.resolveWorkspaceDeliveryPolicy() : options.workspaceDeliveryPolicy; - const [toAgent] = data.to === '@self' - ? await db - .select() - .from(agents) - .where(and(eq(agents.workspaceId, workspaceId), eq(agents.id, fromAgentId))) - : await db - .select() - .from(agents) - .where(and(eq(agents.workspaceId, workspaceId), eq(agents.name, data.to))); - + const [recipient] = await db + .select({ agent: agents, node: addressNodeSelection }) + .from(agents) + .leftJoin(nodes, eq(nodes.id, agents.locationNodeId)) + .where(and( + eq(agents.workspaceId, workspaceId), + data.to === '@self' ? eq(agents.id, fromAgentId) : eq(agents.name, data.to), + )); + + // Checked on the same row the DM is sent to: there is no separate + // resolve-then-send lookup for a concurrent move or re-register to slip into. + if (options.address !== undefined) { + assertAgentAtAddress(options.address, recipient?.agent, recipient?.node ?? null); + } + const toAgent = recipient?.agent; if (!toAgent) { throw codedError(`Agent "${data.to}" not found`, 'agent_not_found', 404); } const [fromAgent] = await db - .select({ name: agents.name }) + .select({ name: agents.name, node: addressNodeSelection }) .from(agents) + .leftJoin(nodes, eq(nodes.id, agents.locationNodeId)) .where(and(eq(agents.workspaceId, workspaceId), eq(agents.id, fromAgentId))); if (!fromAgent?.name) { throw codedError('Sender agent not found', 'internal_error', 500); } + const senderAddress = formatAgentAddress(fromAgent.name, fromAgent.node); // Resolve attachments first so invalid attachments fail before any DM // metadata (channel/conversation/participant rows) is created. @@ -579,7 +614,7 @@ export async function sendDm( const publicResult = buildDmResult({ id: messageId, agentId: fromAgentId, body: data.text, createdAt, - metadata: { ...sanitizeUserMessageMetadata(data.data), injection_mode: data.mode ?? 'wait' }, + metadata: dmMessageMetadata(data, senderAddress), }, conv, fromAgent, data, attachments); const eventData = buildDmReceivedEventData(publicResult, { fromName: fromAgent.name }); const workspacePayload = transformForClient({ @@ -589,7 +624,8 @@ export async function sendDm( // None can escape a capacity rollback, or depend on external transport success. const persist = () => runAtomicWrites(db, (writeDb) => { const writes = buildDmMessageWrites(writeDb, workspaceId, fromAgentId, conv.channelId, data, attachments, messageId, createdAt, - options.receivedA2aAgentId ? { id: options.receivedA2aAgentId, tokenHash: options.receivedA2aTokenHash } : undefined); + options.receivedA2aAgentId ? { id: options.receivedA2aAgentId, tokenHash: options.receivedA2aTokenHash } : undefined, + senderAddress); if (egressId && a2aTarget && egressPayload) { // First statement owns the request identity; a competing attempt rolls @@ -878,6 +914,7 @@ export async function getDmMessages( id: r.id, agent_id: r.agentId, agent_name: r.agentName, + ...senderAddressField(r.metadata), text: r.body, injection_mode: r.metadata?.injection_mode as 'wait' | 'steer' | undefined, metadata: publicMessageMetadata(r.metadata), diff --git a/packages/engine/src/engine/wsTransform.ts b/packages/engine/src/engine/wsTransform.ts index 0c70c8a5..7ca80ff5 100644 --- a/packages/engine/src/engine/wsTransform.ts +++ b/packages/engine/src/engine/wsTransform.ts @@ -100,6 +100,7 @@ export function transformForClient(event: WsEvent): Record { id: (msg.id ?? d.id) as string, agent_id: (msg.agent_id ?? d.from_agent_id ?? d.agent_id) as string, agent_name: (msg.agent_name ?? d.from_name) as string, + ...(typeof msg.agent_address === 'string' ? { agent_address: msg.agent_address } : {}), text: (msg.text ?? d.text) as string, ...(injectionMode ? { injection_mode: injectionMode } : {}), ...(attachments.length ? { attachments } : {}), diff --git a/packages/engine/src/routes/dm.ts b/packages/engine/src/routes/dm.ts index 9e6cf316..4a6afedf 100644 --- a/packages/engine/src/routes/dm.ts +++ b/packages/engine/src/routes/dm.ts @@ -6,7 +6,7 @@ import { rateLimit } from '../middleware/rateLimit.js'; import { jsonIdempotentOk, parseIdempotencyKey, runIdempotent } from '../middleware/idempotency.js'; import { sha256Hex } from '../lib/crypto.js'; import * as dmEngine from '../engine/dm.js'; -import { resolveAgentAddress } from '../engine/address.js'; +import { requireAgentAddress } from '../engine/address.js'; import { resolveMailboxConfig } from '../engine/mailboxConfig.js'; import { resolveWorkspaceDeliveryPolicyFor } from '../engine/workspaceDeliveryPolicy.js'; import { publishWorkspaceEvent } from './fanout.js'; @@ -38,9 +38,10 @@ const sendAddressedSchema = sendDmSchema.omit({ to: true }); /** * Send a DM from the authenticated agent. Shared by POST /v1/dm (recipient by - * name) and POST /v1/to/:address (recipient by `agent@machine`). + * name) and POST /v1/to/:address (recipient by `agent@machine`, passed as + * `address` and checked by the engine against the recipient it sends to). */ -async function sendDirectMessage(c: Context, input: z.infer) { +async function sendDirectMessage(c: Context, input: z.infer, address?: string) { const db = c.get('db'); const workspace = c.get('workspace'); const agent = c.get('agent'); @@ -53,8 +54,11 @@ async function sendDirectMessage(c: Context, input: z.infer, input: z.infer resolveWorkspaceDeliveryPolicyFor(c.get('engine').config, workspace), idempotencyKey, + }, { mailbox, resolveWorkspaceDeliveryPolicy: () => resolveWorkspaceDeliveryPolicyFor(c.get('engine').config, workspace), idempotencyKey, address, afterAdmission: (data, event) => { runInBackground(c, c.get('engine').realtime.publishToWorkspaceStream({ workspaceId: workspace.id, event: { ...event.payload, seq: event.seq }, @@ -193,8 +197,9 @@ dmRoutes.post( if (!parsed.ok) { return parsed.response; } - const target = await resolveAgentAddress(c.get('db'), c.get('workspace').id, c.req.param('address')); - return await sendDirectMessage(c, { ...parsed.data, to: target.agent_name }); + const address = c.req.param('address'); + const { agent } = requireAgentAddress(address); + return await sendDirectMessage(c, { ...parsed.data, to: agent }, address); } catch (err: unknown) { return errorResponse(c, err); } diff --git a/packages/types/CHANGELOG.md b/packages/types/CHANGELOG.md index 21ff40ef..037a81cc 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 - Minor] + +### Added + +- `AgentSchema.address` and `CoreMessagePayloadSchema.agent_address` (optional) for `agent@machine` addressing. ## [8.12.0] - 2026-09-24 diff --git a/packages/types/src/agent.ts b/packages/types/src/agent.ts index 58d393df..048b84c8 100644 --- a/packages/types/src/agent.ts +++ b/packages/types/src/agent.ts @@ -13,6 +13,8 @@ export type AgentStatus = z.infer; export const AgentSchema = z.object({ id: z.string(), name: z.string(), + /** `agent@machine` for POST /v1/to/:address; `machine` is `direct` when not on a broker. */ + address: z.string().optional(), type: AgentTypeSchema, status: AgentStatusSchema, persona: z.string().nullable(), diff --git a/packages/types/src/message.ts b/packages/types/src/message.ts index 3a39cd1b..ce311b29 100644 --- a/packages/types/src/message.ts +++ b/packages/types/src/message.ts @@ -75,6 +75,8 @@ export const CoreMessagePayloadSchema = z.object({ agent_id: z.string(), agent_name: z.string(), agent_type: AgentTypeSchema.optional(), + /** Sender's `agent@machine` address at send time; present on direct messages. */ + agent_address: z.string().optional(), text: z.string(), injection_mode: MessageInjectionModeSchema.optional(), attachments: z.array(FileAttachmentSchema).optional(), From 8711873469c2f996a7d8d7143dc614aadd391807 Mon Sep 17 00:00:00 2001 From: khaliqgant Date: Thu, 24 Sep 2026 21:03:09 -0700 Subject: [PATCH 3/8] fix(engine): unhosted agents have no address instead of agent@direct Cloud tears down a sandbox by deleting its node row, which nulls the agent's location_node_id. The null node was formatted and matched as `direct`, so a torn-down sandbox agent advertised agent@direct and accepted sends nothing could deliver. `direct` now requires a direct node; an agent with no node reports address: null and matches nothing. Adds sandbox-shaped conformance tests (fleet-ensure-* node, cloud:* tags, no machine_id, teardown by row delete). Co-Authored-By: Claude Opus 5.5 (1M context) --- README.md | 4 +- openapi.yaml | 10 ++- packages/engine/CHANGELOG.md | 2 +- .../conformance/addressedSend.test.ts | 73 ++++++++++++++++++- packages/engine/src/engine/address.ts | 21 ++++-- packages/engine/src/engine/dm.ts | 4 +- packages/types/src/agent.ts | 7 +- 7 files changed, 101 insertions(+), 20 deletions(-) diff --git a/README.md b/README.md index d071cd75..16691ef7 100644 --- a/README.md +++ b/README.md @@ -627,7 +627,9 @@ Activity feed channel-message items include `channel_id` and `channel_name`; DM `POST /to/:address` routes by address instead of bare name: `agent@machine` resolves to the agent only while it is hosted on that machine (its broker node's name or `machine_id`, or -`direct` for an agent not on a broker), then delivers like `POST /dm`. Agents expose their +`direct` for a self-connected agent), then delivers like `POST /dm`. A cloud sandbox is a +broker node, so a sandboxed agent's address uses the sandbox's node name; once the sandbox +is torn down the agent has no address (`address: null`) until it is hosted again. Agents expose their address as `address` on agent resources, and each DM carries the sender's as `message.agent_address`, so a recipient can reply on it. A stale address returns `404 address_not_found`; a malformed one returns `400 invalid_address`. An idempotent retry diff --git a/openapi.yaml b/openapi.yaml index b3707f3e..de55351f 100644 --- a/openapi.yaml +++ b/openapi.yaml @@ -322,9 +322,12 @@ components: type: string address: type: string + nullable: true description: >- `agent@machine` for `POST /to/{address}`. `machine` is the broker - node's name, or `direct` when the agent is not hosted on a broker. + node's name (for a cloud sandbox, its node name), or `direct` for a + self-connected agent. Null when the agent is not hosted anywhere: + released, finished, or its node was deleted. type: type: string enum: [agent, human, system] @@ -3929,8 +3932,9 @@ paths: summary: Send DM to agent@machine description: | Send a direct message to an `agent@machine` address. `machine` must be - the agent's current broker node, matched by node name or `machine_id`, - or `direct` for an agent not hosted on a broker. The message is then + the agent's current broker node (a cloud sandbox is one), matched by + node name or `machine_id`, or `direct` for a self-connected agent. An + agent whose node was deleted has no address. The message is then delivered exactly like `POST /dm`. An address that no longer matches (the agent moved or was released) returns 404 instead of reaching the agent elsewhere; matching is exact and case-sensitive. diff --git a/packages/engine/CHANGELOG.md b/packages/engine/CHANGELOG.md index e2e2c4d5..1144198d 100644 --- a/packages/engine/CHANGELOG.md +++ b/packages/engine/CHANGELOG.md @@ -12,7 +12,7 @@ and this project follows [Semantic Versioning](https://semver.org/spec/v2.0.0.ht ### Added - `POST /v1/to/:address` resolves `agent@machine` (broker node name or `machine_id`, or `direct` for agents not on a broker) to the agent hosted there and sends it a DM; stale addresses return `404 address_not_found`, and idempotent retries replay even after the agent moves. -- Agent resources include `address`; DM responses, `dm.received` deliveries, and DM history include the sender's `message.agent_address`. +- Agent resources include `address` (null when the agent is not hosted anywhere, e.g. after its sandbox node is deleted); DM responses, `dm.received` deliveries, and DM history include the sender's `message.agent_address`. ## [8.12.0] - 2026-09-24 diff --git a/packages/engine/src/__tests__/conformance/addressedSend.test.ts b/packages/engine/src/__tests__/conformance/addressedSend.test.ts index 30e94ed5..d8addc64 100644 --- a/packages/engine/src/__tests__/conformance/addressedSend.test.ts +++ b/packages/engine/src/__tests__/conformance/addressedSend.test.ts @@ -1,7 +1,7 @@ import { afterEach, beforeEach, describe, expect, it } from 'vitest'; import { eq } from 'drizzle-orm'; import { makeNodeStack, createWorkspace, registerAgent, FakeSocket, type TestStack } from './harness.js'; -import { agents, messages } from '../../db/schema.js'; +import { agents, messages, nodes } from '../../db/schema.js'; import { parseAgentAddress, SENDER_ADDRESS_METADATA_KEY } from '../../engine/address.js'; type Json = Record; @@ -18,12 +18,18 @@ describe('addressed send', () => { await stack.close(); }); - async function enrollBroker(ws: { workspaceKey: string; workspaceId: string }, nodeId: string, name: string, machineId: string) { + async function enrollBroker( + ws: { workspaceKey: string; workspaceId: string }, + nodeId: string, + name: string, + machineId?: string, + tags: string[] = [], + ) { const enroll = await stack.app.request('/v1/nodes', { method: 'POST', headers: { 'content-type': 'application/json', authorization: `Bearer ${ws.workspaceKey}` }, body: JSON.stringify({ - node_id: nodeId, name, machine_id: machineId, + node_id: nodeId, name, ...(machineId ? { machine_id: machineId } : {}), tags, capabilities: ['spawn:claude'], max_agents: 4, version: 'test-node', }), }); @@ -174,6 +180,67 @@ describe('addressed send', () => { ]); }); + describe('cloud sandboxes', () => { + // Cloud enrolls a sandbox as a broker node named `fleet-ensure-` with + // server-owned `cloud:*` tags and no machine_id, and tears it down by + // deleting the node row, which nulls its agents' location. + const SANDBOX = 'fleet-ensure-3f9a1c2d'; + + async function seedSandbox() { + const { ws, alice, bob } = await seed(); + const sandbox = await enrollBroker(ws, 'node_sandbox', SANDBOX, undefined, + ['cloud:sandbox-provider:daytona', 'cloud:sandbox-id:sbx_123']); + await sandbox.handle.handleMessage(JSON.stringify({ + v: 1, type: 'agent.register', name: 'worker', resumable: true, session_ref: 'sess-worker', + })); + const reply = sandbox.sock.ofType('reply').at(-1) as { ok: boolean; data: { token: string } }; + expect(reply?.ok).toBe(true); + return { ws, alice, bob, sandbox, worker: { token: reply.data.token } }; + } + + async function addressOf(ws: { workspaceKey: string }, name: string) { + const res = await stack.app.request(`/v1/agents/${name}`, { + headers: { authorization: `Bearer ${ws.workspaceKey}` }, + }); + return ((await res.json()) as { data: Json }).data.address; + } + + it('addresses a sandboxed agent by its sandbox node name, both ways', async () => { + const { ws, alice, sandbox, worker } = await seedSandbox(); + expect(await addressOf(ws, 'worker')).toBe(`worker@${SANDBOX}`); + + expect((await send(alice.token, `worker@${SANDBOX}`, { text: 'into the sandbox' })).status).toBe(201); + await stack.settle(); + expect(deliveredDms(sandbox.sock).map((message) => message.text)).toEqual(['into the sandbox']); + + // Out of the sandbox, carrying the sandbox address for the reply. + const res = await send(worker.token, 'alice@direct', { text: 'from the sandbox' }); + expect(((await res.json()) as { data: { message: Json } }).data.message.agent_address) + .toBe(`worker@${SANDBOX}`); + // Server-owned cloud tags are not machine names. + expect(await errorCode(await send(alice.token, 'worker@sbx_123', { text: 'x' }))).toBe('address_not_found'); + }); + + it('leaves a torn-down sandbox agent with no address instead of falling back to direct', async () => { + const { ws, alice } = await seedSandbox(); + await stack.runtime.deps.db.delete(nodes).where(eq(nodes.id, 'node_sandbox')); + + expect(await addressOf(ws, 'worker')).toBeNull(); + for (const address of [`worker@${SANDBOX}`, 'worker@direct']) { + expect(await errorCode(await send(alice.token, address, { text: 'x' })), address).toBe('address_not_found'); + } + }); + + it('sends from an unhosted agent without an agent_address', async () => { + const { ws, alice, bob } = await seedSandbox(); + await stack.runtime.deps.db.update(agents).set({ locationNodeId: null }).where(eq(agents.id, bob.agentId)); + const res = await send(bob.token, 'alice@direct', { text: 'from nowhere' }); + expect(res.status).toBe(201); + expect(((await res.json()) as { data: { message: Json } }).data.message).not.toHaveProperty('agent_address'); + void ws; void alice; + }); + }); + describe('idempotent retries', () => { it('replays an accepted send even after the agent moves', async () => { const { ws, alice, bob } = await seed(); diff --git a/packages/engine/src/engine/address.ts b/packages/engine/src/engine/address.ts index 709e7848..5524662f 100644 --- a/packages/engine/src/engine/address.ts +++ b/packages/engine/src/engine/address.ts @@ -2,8 +2,8 @@ import { nodes } from '../db/schema.js'; import { codedError } from '../lib/httpError.js'; /** - * Machine name for agents not hosted on a broker: self-connected agents and - * agents on an implicit direct node, whose node name is an internal id. + * Machine name for agents on a direct node (self-connected agents), whose node + * name is an internal id rather than a machine. */ export const DIRECT_MACHINE = 'direct'; @@ -27,16 +27,21 @@ export function parseAgentAddress(address: string): AgentAddress | null { return { agent: address.slice(0, at), machine: address.slice(at + 1) }; } -/** Canonical address: the broker node's name, or `direct` when there is no broker. */ -export function formatAgentAddress(agentName: string, node: AddressNode): string { - const machine = node && node.role !== 'direct' ? node.name : DIRECT_MACHINE; - return `${agentName}@${machine}`; +/** + * Canonical address: the broker node's name, or `direct` for a direct node. + * Null when the agent has no node — released, finished, or its node was + * deleted (a torn-down sandbox) — because nothing could deliver to it. + */ +export function formatAgentAddress(agentName: string, node: AddressNode): string | null { + if (!node) return null; + return `${agentName}@${node.role === 'direct' ? DIRECT_MACHINE : node.name}`; } /** Whether `machine` names the node the agent is on: node name, machine_id, or `direct`. */ function machineMatches(machine: string, node: AddressNode): boolean { - if (machine === DIRECT_MACHINE && (!node || node.role === 'direct')) return true; - return node !== null && (machine === node.name || machine === node.machineId); + if (!node) return false; + if (machine === DIRECT_MACHINE && node.role === 'direct') return true; + return machine === node.name || machine === node.machineId; } /** Node columns needed to format or match an address; select them via a left join on the agent's location node. */ diff --git a/packages/engine/src/engine/dm.ts b/packages/engine/src/engine/dm.ts index 0cf757db..96a7f056 100644 --- a/packages/engine/src/engine/dm.ts +++ b/packages/engine/src/engine/dm.ts @@ -291,7 +291,7 @@ async function resolveConversation( */ function dmMessageMetadata( data: { mode?: 'wait' | 'steer'; data?: Record | null }, - senderAddress: string | undefined, + senderAddress: string | null | undefined, ): Record { return { ...sanitizeUserMessageMetadata(data.data), @@ -320,7 +320,7 @@ function buildDmMessageWrites( messageId: string, createdAt = new Date(), inboundRegistration?: { id: string; tokenHash?: string }, - senderAddress?: string, + senderAddress?: string | null, ): AtomicWrite[] { const hasAttachments = attachments.length > 0; const metadata = dmMessageMetadata(data, senderAddress); diff --git a/packages/types/src/agent.ts b/packages/types/src/agent.ts index 048b84c8..268578cf 100644 --- a/packages/types/src/agent.ts +++ b/packages/types/src/agent.ts @@ -13,8 +13,11 @@ export type AgentStatus = z.infer; export const AgentSchema = z.object({ id: z.string(), name: z.string(), - /** `agent@machine` for POST /v1/to/:address; `machine` is `direct` when not on a broker. */ - address: z.string().optional(), + /** + * `agent@machine` for POST /v1/to/:address; `machine` is `direct` for a + * self-connected agent. Null when the agent is not hosted anywhere. + */ + address: z.string().nullable().optional(), type: AgentTypeSchema, status: AgentStatusSchema, persona: z.string().nullable(), From 8400db297ccbb8d59334971c3a7aaf152c57c85e Mon Sep 17 00:00:00 2001 From: khaliqgant Date: Thu, 24 Sep 2026 22:12:20 -0700 Subject: [PATCH 4/8] fix(engine): address review feedback on agent@machine addressing - re-check the addressed recipient inside the admission write: a move or release that races the send rolls the batch back as address_not_found - try every @ as the agent/machine separator so names containing @ are addressable; two matching readings return 400 ambiguous_address - reserve `direct`: it never matches a broker, and a broker named direct is addressed by machine_id - charge POST /v1/to/:address to the POST /v1/dm rate-limit bucket - include address in PATCH /v1/agents responses - document the 409 and body-validation 400 responses Co-Authored-By: Claude Opus 5.5 (1M context) --- README.md | 5 +- openapi.yaml | 22 ++++- packages/engine/CHANGELOG.md | 2 +- .../conformance/addressedSend.test.ts | 93 ++++++++++++++++--- .../conformance/rateLimitContract.test.ts | 22 ++++- packages/engine/src/engine/address.ts | 73 +++++++++------ packages/engine/src/engine/agent.ts | 14 ++- packages/engine/src/engine/dm.ts | 49 ++++++---- packages/engine/src/middleware/rateLimit.ts | 3 + packages/engine/src/ports/database.ts | 4 +- packages/engine/src/routes/dm.ts | 6 +- 11 files changed, 224 insertions(+), 69 deletions(-) diff --git a/README.md b/README.md index 16691ef7..e81c1b57 100644 --- a/README.md +++ b/README.md @@ -632,7 +632,10 @@ broker node, so a sandboxed agent's address uses the sandbox's node name; once t is torn down the agent has no address (`address: null`) until it is hosted again. Agents expose their address as `address` on agent resources, and each DM carries the sender's as `message.agent_address`, so a recipient can reply on it. A stale address returns -`404 address_not_found`; a malformed one returns `400 invalid_address`. An idempotent retry +`404 address_not_found`, including when the agent moves while the send is in flight; a malformed +one returns `400 invalid_address`. Names may contain `@`: each `@` is tried as the separator, and an +address that reads as two different agents returns `400 ambiguous_address`. Addressed DMs share +the `POST /dm` rate-limit bucket. An idempotent retry replays even if the agent has moved; reusing the key for another address is a `409`. `POST /dm` can return **`409 dm_conversation_id_collision`**. A 1:1 conversation id is derived diff --git a/openapi.yaml b/openapi.yaml index de55351f..d867e8e4 100644 --- a/openapi.yaml +++ b/openapi.yaml @@ -3937,7 +3937,14 @@ paths: agent whose node was deleted has no address. The message is then delivered exactly like `POST /dm`. An address that no longer matches (the agent moved or was released) returns 404 instead of reaching the - agent elsewhere; matching is exact and case-sensitive. + agent elsewhere; matching is exact and case-sensitive. The check is + repeated inside the admission write, so a move that races the send + also returns 404. + + Agent names, node names and `machine_id` values may contain `@`: every + `@` is tried as the separator, and the one reading that names an agent + on that machine is used. `direct` is reserved and never matches a + broker; a broker named `direct` is addressed by its `machine_id`. Agents learn addresses from `address` on agent resources and from `message.agent_address` on DMs they receive. @@ -4000,7 +4007,10 @@ paths: data: $ref: '#/components/schemas/DmSendResponse' '400': - description: '`invalid_address` — the address is not of the form `agent@machine`' + description: | + `invalid_address` — the address has no `agent@machine` reading; + `ambiguous_address` — more than one reading names an agent on + that machine; or the body is invalid (missing or empty `text`). content: application/json: schema: @@ -4011,6 +4021,14 @@ paths: application/json: schema: $ref: '#/components/schemas/ErrorResponse' + '409': + description: | + `idempotency_key_reused` or `dm_conversation_id_collision`, with + the same semantics as `POST /dm`. + content: + application/json: + schema: + $ref: '#/components/schemas/ErrorResponse' /dm/conversations: get: diff --git a/packages/engine/CHANGELOG.md b/packages/engine/CHANGELOG.md index 1144198d..af33730e 100644 --- a/packages/engine/CHANGELOG.md +++ b/packages/engine/CHANGELOG.md @@ -11,7 +11,7 @@ and this project follows [Semantic Versioning](https://semver.org/spec/v2.0.0.ht ### Added -- `POST /v1/to/:address` resolves `agent@machine` (broker node name or `machine_id`, or `direct` for agents not on a broker) to the agent hosted there and sends it a DM; stale addresses return `404 address_not_found`, and idempotent retries replay even after the agent moves. +- `POST /v1/to/:address` resolves `agent@machine` (broker node name or `machine_id`, or `direct` for agents not on a broker) to the agent hosted there and sends it a DM; stale addresses return `404 address_not_found`, idempotent retries replay even after the agent moves, a move that races the send is rejected atomically, names containing `@` resolve (`400 ambiguous_address` when two readings match), and addressed DMs share the `POST /v1/dm` rate-limit bucket. - Agent resources include `address` (null when the agent is not hosted anywhere, e.g. after its sandbox node is deleted); DM responses, `dm.received` deliveries, and DM history include the sender's `message.agent_address`. ## [8.12.0] - 2026-09-24 diff --git a/packages/engine/src/__tests__/conformance/addressedSend.test.ts b/packages/engine/src/__tests__/conformance/addressedSend.test.ts index d8addc64..ea60e184 100644 --- a/packages/engine/src/__tests__/conformance/addressedSend.test.ts +++ b/packages/engine/src/__tests__/conformance/addressedSend.test.ts @@ -1,8 +1,8 @@ import { afterEach, beforeEach, describe, expect, it } from 'vitest'; import { eq } from 'drizzle-orm'; -import { makeNodeStack, createWorkspace, registerAgent, FakeSocket, type TestStack } from './harness.js'; +import { attachFakeBatch, makeNodeStack, createWorkspace, registerAgent, FakeSocket, type TestStack } from './harness.js'; import { agents, messages, nodes } from '../../db/schema.js'; -import { parseAgentAddress, SENDER_ADDRESS_METADATA_KEY } from '../../engine/address.js'; +import { addressSplits, SENDER_ADDRESS_METADATA_KEY } from '../../engine/address.js'; type Json = Record; @@ -48,17 +48,21 @@ describe('addressed send', () => { return { sock, handle }; } + async function registerOnNode(node: Awaited>, name: string) { + await node.handle.handleMessage(JSON.stringify({ + v: 1, type: 'agent.register', name, resumable: true, session_ref: `sess-${name}`, + })); + const reply = node.sock.ofType('reply').at(-1) as { ok: boolean; data: { agent_id: string; token: string } }; + expect(reply?.ok).toBe(true); + return { agentId: reply.data.agent_id, token: reply.data.token }; + } + /** alice is self-connected; bob runs on broker node `laptop` (machine_id `mach-123`). */ async function seed() { const ws = await createWorkspace(stack.app, 'address-ws'); const alice = await registerAgent(stack.app, ws.workspaceKey, 'alice'); const laptop = await enrollBroker(ws, 'node_laptop', 'laptop', 'mach-123'); - await laptop.handle.handleMessage(JSON.stringify({ - v: 1, type: 'agent.register', name: 'bob', resumable: true, session_ref: 'sess-bob', - })); - const reply = laptop.sock.ofType('reply').at(-1) as { ok: boolean; data: { agent_id: string; token: string } }; - expect(reply?.ok).toBe(true); - const bob = { agentId: reply.data.agent_id, token: reply.data.token }; + const bob = await registerOnNode(laptop, 'bob'); return { ws, alice, bob, laptop }; } @@ -85,12 +89,66 @@ describe('addressed send', () => { await stack.runtime.deps.db.update(agents).set({ locationNodeId: nodeId }).where(eq(agents.id, agentId)); } - it('parses agent@machine on the last @', () => { - expect(parseAgentAddress('bob@laptop')).toEqual({ agent: 'bob', machine: 'laptop' }); - expect(parseAgentAddress('a@b@laptop')).toEqual({ agent: 'a@b', machine: 'laptop' }); - expect(parseAgentAddress('bob')).toBeNull(); - expect(parseAgentAddress('@laptop')).toBeNull(); - expect(parseAgentAddress('bob@')).toBeNull(); + it('reads every @ as a candidate agent/machine split', () => { + expect(addressSplits('bob@laptop')).toEqual([{ agent: 'bob', machine: 'laptop' }]); + expect(addressSplits('a@b@laptop')).toEqual([ + { agent: 'a', machine: 'b@laptop' }, + { agent: 'a@b', machine: 'laptop' }, + ]); + expect(addressSplits('bob')).toEqual([]); + expect(addressSplits('@laptop')).toEqual([]); + expect(addressSplits('bob@')).toEqual([]); + }); + + it('resolves @ inside an agent name or a machine name, and rejects an ambiguous address', async () => { + const { ws, alice } = await seed(); + const email = await enrollBroker(ws, 'node_email', 'ops@example.com', undefined); + await registerOnNode(email, 'carol'); + const plain = await enrollBroker(ws, 'node_example', 'example.com', undefined); + await registerOnNode(plain, 'dave@ops'); + + expect((await send(alice.token, 'carol@ops@example.com', { text: 'machine has @' })).status).toBe(201); + expect((await send(alice.token, 'dave@ops@example.com', { text: 'agent has @' })).status).toBe(201); + await stack.settle(); + expect(deliveredDms(email.sock).map((message) => message.text)).toEqual(['machine has @']); + expect(deliveredDms(plain.sock).map((message) => message.text)).toEqual(['agent has @']); + + // `x@ops@example.com` is agent `x` on `ops@example.com` or agent `x@ops` + // on `example.com`; when both exist the address names neither. + await registerOnNode(email, 'x'); + await registerOnNode(plain, 'x@ops'); + const ambiguous = await send(alice.token, 'x@ops@example.com', { text: 'which one?' }); + expect(ambiguous.status).toBe(400); + expect(await errorCode(ambiguous)).toBe('ambiguous_address'); + }); + + it('reserves `direct`: a broker named direct is addressed by machine_id, never @direct', async () => { + const { ws, alice } = await seed(); + const node = await enrollBroker(ws, 'node_named_direct', 'direct', 'mach-d'); + await registerOnNode(node, 'erin'); + + const res = await stack.app.request('/v1/agents/erin', { headers: { authorization: `Bearer ${ws.workspaceKey}` } }); + expect(((await res.json()) as { data: Json }).data.address).toBe('erin@mach-d'); + expect(await errorCode(await send(alice.token, 'erin@direct', { text: 'x' }))).toBe('address_not_found'); + expect((await send(alice.token, 'erin@mach-d', { text: 'x' })).status).toBe(201); + }); + + it('rolls back an addressed send when the agent moves between selection and commit', async () => { + const { ws, alice, bob } = await seed(); + await enrollBroker(ws, 'node_desktop', 'desktop', 'mach-456'); + const sqlite = stack.runtime.handle.sqlite; + let moved = false; + attachFakeBatch(stack, stack.runtime.deps.db, async () => { + if (moved) return; + moved = true; + sqlite.prepare('UPDATE agents SET location_node_id = ? WHERE id = ?').run('node_desktop', bob.agentId); + }); + + const res = await send(alice.token, 'bob@laptop', { text: 'raced' }); + expect(moved).toBe(true); + expect(res.status).toBe(404); + expect(await errorCode(res)).toBe('address_not_found'); + expect(await stack.runtime.deps.db.select().from(messages).where(eq(messages.body, 'raced'))).toHaveLength(0); }); it('routes to the agent by node name and by machine_id, delivering on that machine', async () => { @@ -302,6 +360,13 @@ describe('addressed send', () => { headers: { authorization: `Bearer ${bob.token}` }, })).json()) as { data: Json }; expect(self.data.address).toBe('bob@laptop'); + + const patched = (await (await stack.app.request('/v1/agents/bob', { + method: 'PATCH', + headers: { 'content-type': 'application/json', ...auth }, + body: JSON.stringify({ persona: 'reviewer' }), + })).json()) as { data: Json }; + expect(patched.data.address).toBe('bob@laptop'); }); it('carries the sender address on the DM response, live delivery, and history, and it round-trips', async () => { diff --git a/packages/engine/src/__tests__/conformance/rateLimitContract.test.ts b/packages/engine/src/__tests__/conformance/rateLimitContract.test.ts index 6ba66ad6..7b5100d7 100644 --- a/packages/engine/src/__tests__/conformance/rateLimitContract.test.ts +++ b/packages/engine/src/__tests__/conformance/rateLimitContract.test.ts @@ -1,7 +1,7 @@ import { afterEach, describe, expect, it } from 'vitest'; import type { EntitlementsProvider, PlanLimits, Workspace } from '../../ports/index.js'; import { InProcessRateLimiter } from '../../adapters/node/rate-limit.js'; -import { createWorkspace, makeNodeStack, type TestStack } from './harness.js'; +import { createWorkspace, makeNodeStack, registerAgent, type TestStack } from './harness.js'; function entitlementsWithRate(ratePerMin: number): EntitlementsProvider { return { @@ -71,6 +71,26 @@ describe('rate limit contract', () => { expect(throttled.headers.get('X-RateLimit-Remaining')).toBe('0'); }); + it('charges an addressed DM to the same bucket as POST /v1/dm', async () => { + stack = makeNodeStack({ entitlements: entitlementsWithRate(10) }); + const ws = await createWorkspace(stack.app, 'addressed-dm-bucket-ws'); + const alice = await registerAgent(stack.app, ws.workspaceKey, 'alice'); + await registerAgent(stack.app, ws.workspaceKey, 'bob'); + const post = (path: string, body: unknown) => (stack as TestStack).app.request(path, { + method: 'POST', + headers: { 'content-type': 'application/json', authorization: `Bearer ${alice.token}` }, + body: JSON.stringify(body), + }); + + const dm = await post('/v1/dm', { to: 'bob', text: 'one' }); + expect(dm.headers.get('X-RateLimit-Limit')).toBe('5'); + expect(dm.headers.get('X-RateLimit-Remaining')).toBe('4'); + const addressed = await post('/v1/to/bob%40direct', { text: 'two' }); + expect(addressed.status).toBe(201); + expect(addressed.headers.get('X-RateLimit-Limit')).toBe('5'); + expect(addressed.headers.get('X-RateLimit-Remaining')).toBe('3'); + }); + it('reports the window reset on a successful response', async () => { stack = makeNodeStack({ entitlements: entitlementsWithRate(10) }); const ws = await createWorkspace(stack.app, 'reset-header-ws'); diff --git a/packages/engine/src/engine/address.ts b/packages/engine/src/engine/address.ts index 5524662f..e9087a0e 100644 --- a/packages/engine/src/engine/address.ts +++ b/packages/engine/src/engine/address.ts @@ -18,61 +18,78 @@ export interface AgentAddress { type AddressNode = { name: string; role: string; machineId: string | null } | null; /** - * Parse an `agent@machine` address. The split is on the last `@` so agent - * names that themselves contain `@` still resolve. + * Every way to read `address` as `agent@machine`. Agent names, node names and + * machine ids may all contain `@`, so each separator is a candidate split; the + * recipient lookup keeps the split that names a real agent on that machine. */ -export function parseAgentAddress(address: string): AgentAddress | null { - const at = address.lastIndexOf('@'); - if (at <= 0 || at === address.length - 1) return null; - return { agent: address.slice(0, at), machine: address.slice(at + 1) }; +export function addressSplits(address: string): AgentAddress[] { + const splits: AgentAddress[] = []; + for (let at = address.indexOf('@'); at !== -1; at = address.indexOf('@', at + 1)) { + if (at > 0 && at < address.length - 1) { + splits.push({ agent: address.slice(0, at), machine: address.slice(at + 1) }); + } + } + return splits; } /** * Canonical address: the broker node's name, or `direct` for a direct node. - * Null when the agent has no node — released, finished, or its node was - * deleted (a torn-down sandbox) — because nothing could deliver to it. + * `direct` is reserved, so a broker whose name is `direct` is addressed by its + * machine_id instead. Null when the agent has no node — released, finished, + * or its node was deleted (a torn-down sandbox) — because nothing could + * deliver to it, or when a broker has no usable identifier. */ export function formatAgentAddress(agentName: string, node: AddressNode): string | null { if (!node) return null; - return `${agentName}@${node.role === 'direct' ? DIRECT_MACHINE : node.name}`; + if (node.role === 'direct') return `${agentName}@${DIRECT_MACHINE}`; + const machine = [node.name, node.machineId].find((id) => id && id !== DIRECT_MACHINE); + return machine ? `${agentName}@${machine}` : null; } -/** Whether `machine` names the node the agent is on: node name, machine_id, or `direct`. */ -function machineMatches(machine: string, node: AddressNode): boolean { +/** + * Whether `machine` names the node the agent is on. `direct` matches only a + * direct node; a broker matches by node name or machine_id, never `direct`. + */ +export function machineMatches(machine: string, node: AddressNode): boolean { if (!node) return false; - if (machine === DIRECT_MACHINE && node.role === 'direct') return true; + if (machine === DIRECT_MACHINE || node.role === 'direct') { + return machine === DIRECT_MACHINE && node.role === 'direct'; + } return machine === node.name || machine === node.machineId; } /** Node columns needed to format or match an address; select them via a left join on the agent's location node. */ export const addressNodeSelection = { name: nodes.name, role: nodes.role, machineId: nodes.machineId }; -/** Reject an address that is not `agent@machine` before any other work. */ -export function requireAgentAddress(address: string): AgentAddress { - const parsed = parseAgentAddress(address); - if (!parsed) { +/** Reject an address with no `agent@machine` reading before any other work. */ +export function requireAgentAddress(address: string): AgentAddress[] { + const splits = addressSplits(address); + if (splits.length === 0) { throw codedError('Address must be of the form "agent@machine"', 'invalid_address', 400); } - return parsed; + return splits; } /** - * Throw unless `agent` (looked up by the address's agent name) is live and - * hosted on the address's machine. A stale address (the agent moved or was - * released) fails instead of silently reaching the agent somewhere else. + * Pick the recipient `address` names from candidate agents (looked up by every + * split's agent name). A stale address — the agent moved or was released — + * matches nothing and fails instead of reaching the agent somewhere else. */ -export function assertAgentAtAddress( +export function selectAddressedRecipient( address: string, - agent: { status: string } | undefined, - node: AddressNode, -): void { - const parsed = requireAgentAddress(address); - if (!agent || agent.status === 'released' || !machineMatches(parsed.machine, node)) { - throw addressNotFound(address); + candidates: T[], +): T { + const splits = requireAgentAddress(address); + const matches = candidates.filter(({ agent, node }) => agent.status !== 'released' + && splits.some((split) => split.agent === agent.name && machineMatches(split.machine, node))); + if (matches.length > 1) { + throw codedError(`Address "${address}" matches more than one agent`, 'ambiguous_address', 400); } + if (!matches[0]) throw addressNotFound(address); + return matches[0]; } -function addressNotFound(address: string) { +export function addressNotFound(address: string) { return codedError(`No agent at address "${address}"`, 'address_not_found', 404); } diff --git a/packages/engine/src/engine/agent.ts b/packages/engine/src/engine/agent.ts index 306787ff..33eb2ab9 100644 --- a/packages/engine/src/engine/agent.ts +++ b/packages/engine/src/engine/agent.ts @@ -7,7 +7,7 @@ import { invalidateChannelCache } from './cache.js'; import { queryInChunks } from '../lib/queryChunks.js'; import { codedError } from '../lib/httpError.js'; import { directNodeIdForAgent } from './node.js'; -import { formatAgentAddress } from './address.js'; +import { addressNodeSelection, formatAgentAddress } from './address.js'; import { runAtomicWrites, type AtomicWrite } from '../ports/database.js'; import { AGENT_RECOVERY_PROOF_HASH_PATTERN } from '@relaycast/types'; @@ -387,7 +387,7 @@ export async function listAgents(db: Db, workspaceId: string, status?: string) { } const rows = await db - .select({ agent: agents, node: { name: nodes.name, role: nodes.role, machineId: nodes.machineId } }) + .select({ agent: agents, node: addressNodeSelection }) .from(agents) .leftJoin(nodes, eq(nodes.id, agents.locationNodeId)) // Released rows are tombstones retained only to keep history attributable; @@ -411,7 +411,7 @@ export async function listAgents(db: Db, workspaceId: string, status?: string) { export async function getAgentByName(db: Db, workspaceId: string, name: string) { const [row] = await db - .select({ agent: agents, node: { name: nodes.name, role: nodes.role, machineId: nodes.machineId } }) + .select({ agent: agents, node: addressNodeSelection }) .from(agents) .leftJoin(nodes, eq(nodes.id, agents.locationNodeId)) .where(and(eq(agents.workspaceId, workspaceId), eq(agents.name, name))); @@ -593,11 +593,15 @@ export async function updateAgentById( .returning(); if (!updated) return null; + const [node] = updated.locationNodeId + ? await db.select(addressNodeSelection).from(nodes).where(eq(nodes.id, updated.locationNodeId)) + : []; return { id: updated.id, name: updated.name, handle: `@${updated.name}`, + address: formatAgentAddress(updated.name, node ?? null), type: updated.type, status: effectiveAgentStatus(updated), persona: updated.persona, @@ -642,11 +646,15 @@ export async function claimLegacyAgentIdentity( .returning(); if (!updated) return null; + const [node] = updated.locationNodeId + ? await db.select(addressNodeSelection).from(nodes).where(eq(nodes.id, updated.locationNodeId)) + : []; return { id: updated.id, name: updated.name, handle: `@${updated.name}`, + address: formatAgentAddress(updated.name, node ?? null), type: updated.type, status: updated.status, persona: updated.persona, diff --git a/packages/engine/src/engine/dm.ts b/packages/engine/src/engine/dm.ts index 96a7f056..db4779dd 100644 --- a/packages/engine/src/engine/dm.ts +++ b/packages/engine/src/engine/dm.ts @@ -39,8 +39,10 @@ import { transformForClient } from './wsTransform.js'; import { codedError } from '../lib/httpError.js'; import { addressNodeSelection, - assertAgentAtAddress, + addressNotFound, + addressSplits, formatAgentAddress, + selectAddressedRecipient, SENDER_ADDRESS_METADATA_KEY, senderAddressField, } from './address.js'; @@ -71,9 +73,9 @@ interface SendDmOptions { /** Server-resolved workspace growth policy; absent => no workspace guard. */ workspaceDeliveryPolicy?: WorkspaceDeliveryPolicy; /** - * `agent@machine` the caller addressed. Checked against the same recipient - * row the DM is sent to, after durable replay lookup, so an accepted retry - * still replays after the agent moves. + * `agent@machine` the caller addressed; replaces the `to` lookup. Resolved + * after durable replay lookup, so an accepted retry still replays after the + * agent moves, and re-checked inside the admission write. */ address?: string; } @@ -321,6 +323,7 @@ function buildDmMessageWrites( createdAt = new Date(), inboundRegistration?: { id: string; tokenHash?: string }, senderAddress?: string | null, + addressedRecipient?: { agentId: string; nodeId: string }, ): AtomicWrite[] { const hasAttachments = attachments.length > 0; const metadata = dmMessageMetadata(data, senderAddress); @@ -341,7 +344,16 @@ function buildDmMessageWrites( AND a.id = ${fromAgentId} AND a.workspace_id = ${workspaceId} ${inboundRegistration.tokenHash === undefined ? sql`` : sql`AND a.token_hash = ${inboundRegistration.tokenHash}`} )` : fromAgentId, - body: data.text, + // NULL violates messages.body when an addressed recipient left the + // addressed node (or was released) after it was selected, so the whole + // admission rolls back instead of delivering to its new location. + body: addressedRecipient ? sql`( + SELECT ${data.text} WHERE EXISTS ( + SELECT 1 FROM agents a + WHERE a.id = ${addressedRecipient.agentId} AND a.workspace_id = ${workspaceId} + AND a.status <> 'released' AND a.location_node_id = ${addressedRecipient.nodeId} + ) + )` : data.text, hasAttachments, metadata, sessionRef, @@ -532,20 +544,22 @@ export async function sendDm( const workspacePolicy = options.resolveWorkspaceDeliveryPolicy ? await options.resolveWorkspaceDeliveryPolicy() : options.workspaceDeliveryPolicy; - const [recipient] = await db + const recipientQuery = db .select({ agent: agents, node: addressNodeSelection }) .from(agents) - .leftJoin(nodes, eq(nodes.id, agents.locationNodeId)) - .where(and( + .leftJoin(nodes, eq(nodes.id, agents.locationNodeId)); + const [recipient] = options.address !== undefined + ? [selectAddressedRecipient(options.address, await recipientQuery.where(and( + eq(agents.workspaceId, workspaceId), + inArray(agents.name, addressSplits(options.address).map((split) => split.agent)), + )))] + : await recipientQuery.where(and( eq(agents.workspaceId, workspaceId), data.to === '@self' ? eq(agents.id, fromAgentId) : eq(agents.name, data.to), )); - - // Checked on the same row the DM is sent to: there is no separate - // resolve-then-send lookup for a concurrent move or re-register to slip into. - if (options.address !== undefined) { - assertAgentAtAddress(options.address, recipient?.agent, recipient?.node ?? null); - } + const addressedRecipient = options.address !== undefined && recipient?.agent.locationNodeId + ? { agentId: recipient.agent.id, nodeId: recipient.agent.locationNodeId } + : undefined; const toAgent = recipient?.agent; if (!toAgent) { throw codedError(`Agent "${data.to}" not found`, 'agent_not_found', 404); @@ -615,7 +629,7 @@ export async function sendDm( const publicResult = buildDmResult({ id: messageId, agentId: fromAgentId, body: data.text, createdAt, metadata: dmMessageMetadata(data, senderAddress), - }, conv, fromAgent, data, attachments); + }, conv, fromAgent, options.address !== undefined ? { ...data, to: toAgent.name } : data, attachments); const eventData = buildDmReceivedEventData(publicResult, { fromName: fromAgent.name }); const workspacePayload = transformForClient({ type: 'dm.received', workspace_id: workspaceId, data: eventData, timestamp: createdAt.toISOString(), @@ -625,7 +639,7 @@ export async function sendDm( const persist = () => runAtomicWrites(db, (writeDb) => { const writes = buildDmMessageWrites(writeDb, workspaceId, fromAgentId, conv.channelId, data, attachments, messageId, createdAt, options.receivedA2aAgentId ? { id: options.receivedA2aAgentId, tokenHash: options.receivedA2aTokenHash } : undefined, - senderAddress); + senderAddress, addressedRecipient); if (egressId && a2aTarget && egressPayload) { // First statement owns the request identity; a competing attempt rolls @@ -691,6 +705,9 @@ export async function sendDm( const results = await persist(); if (egressId || options.receivedA2aAgentId || inboundId) admittedEventSeq = (results[results.length - 1] as { seq: number }[])[0].seq; } catch (error) { + if (options.address !== undefined && databaseConstraintKind(error) === 'address_changed') { + throw addressNotFound(options.address); + } if (options.receivedA2aAgentId && databaseConstraintKind(error) === 'a2a_registration_changed') { throw codedError('Authenticated A2A registration changed before admission', 'a2a_registration_changed', 401); } diff --git a/packages/engine/src/middleware/rateLimit.ts b/packages/engine/src/middleware/rateLimit.ts index f52a4d1d..ad7b9fac 100644 --- a/packages/engine/src/middleware/rateLimit.ts +++ b/packages/engine/src/middleware/rateLimit.ts @@ -34,6 +34,9 @@ function getRouteKey(method: string, path: string): string | null { // Normalize path: /v1/channels/foo/messages -> /channels/*/messages const normalized = path .replace(/^\/v1/, '') + // An addressed DM is a DM: it shares the POST:/dm bucket so the new path + // cannot add a second DM budget. + .replace(/^\/to\/[^/]+$/, '/dm') .replace(/\/[a-zA-Z0-9_-]+\/messages/, '/*/messages') .replace(/\/[a-zA-Z0-9_-]+\/reactions/, '/*/reactions') .replace(/\/[a-zA-Z0-9_-]+\/replies/, '/*/replies'); diff --git a/packages/engine/src/ports/database.ts b/packages/engine/src/ports/database.ts index fe7b70b9..67ef6373 100644 --- a/packages/engine/src/ports/database.ts +++ b/packages/engine/src/ports/database.ts @@ -187,7 +187,7 @@ export async function runAtomicWrites( /** Normalize atomic admission sentinels across SQLite adapters (D1 wraps the message). * Keep driver-specific error decoding at this port boundary; engine callers use a stable kind. */ -export type DatabaseConstraintKind = 'mailbox_capacity' | 'workspace_delivery_capacity' | 'a2a_registration_changed'; +export type DatabaseConstraintKind = 'mailbox_capacity' | 'workspace_delivery_capacity' | 'a2a_registration_changed' | 'address_changed'; export function databaseConstraintKind(error: unknown): DatabaseConstraintKind | undefined { const seen = new Set(); @@ -204,6 +204,8 @@ export function databaseConstraintKind(error: unknown): DatabaseConstraintKind | // only for a row that trips both; either is a capacity refusal. if (/NOT NULL constraint failed: deliveries\.workspace_id/i.test(message)) return 'mailbox_capacity'; if (/NOT NULL constraint failed: deliveries\.status/i.test(message)) return 'workspace_delivery_capacity'; + // Addressed DM admission: the recipient left the addressed node before commit. + if (/NOT NULL constraint failed: messages\.body/i.test(message)) return 'address_changed'; cause = record.cause; } return undefined; diff --git a/packages/engine/src/routes/dm.ts b/packages/engine/src/routes/dm.ts index 4a6afedf..9dbea3aa 100644 --- a/packages/engine/src/routes/dm.ts +++ b/packages/engine/src/routes/dm.ts @@ -198,8 +198,10 @@ dmRoutes.post( return parsed.response; } const address = c.req.param('address'); - const { agent } = requireAgentAddress(address); - return await sendDirectMessage(c, { ...parsed.data, to: agent }, address); + requireAgentAddress(address); + // The engine resolves the recipient from `address`; `to` only feeds the + // idempotency fingerprint. + return await sendDirectMessage(c, { ...parsed.data, to: address }, address); } catch (err: unknown) { return errorResponse(c, err); } From 11e72c0ca60ec45300b1170b80f3775f1fb75ecc Mon Sep 17 00:00:00 2001 From: khaliqgant Date: Thu, 24 Sep 2026 22:51:02 -0700 Subject: [PATCH 5/8] refactor(engine): fold agent@machine addressing into POST /v1/dm Replace POST /v1/to/:address with an optional `address` field on POST /v1/dm (exactly one of `to` or `address`). The address is a public name, not a credential, so a separate URL added a second DM route, rate bucket, and cloud admission classification for no benefit. Resolution, the atomic admission re-check, idempotency, and sender addresses are unchanged; the DM route diff against main is now a few lines. Co-Authored-By: Claude Opus 5.5 (1M context) --- CHANGELOG.md | 2 +- README.md | 24 +- openapi.yaml | 139 +++------ packages/engine/CHANGELOG.md | 2 +- .../conformance/addressedSend.test.ts | 30 +- .../conformance/rateLimitContract.test.ts | 22 +- packages/engine/src/middleware/rateLimit.ts | 3 - packages/engine/src/routes/dm.ts | 277 ++++++++---------- packages/sdk-typescript/CHANGELOG.md | 2 +- .../src/__tests__/agent-messaging.test.ts | 6 +- packages/sdk-typescript/src/agent.ts | 5 +- packages/types/src/agent.ts | 2 +- packages/types/src/dm.ts | 4 +- 13 files changed, 201 insertions(+), 317 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index c1a70fc6..02f8de09 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -20,7 +20,7 @@ Packages without a separate changelog are covered by the cross-package notes bel ### Added -- `POST /v1/to/:address` sends a DM to an `agent@machine` address, delivering to the agent only while it is hosted on that machine. +- `POST /v1/dm` accepts `address` (`agent@machine`) in place of `to`, delivering to the agent only while it is hosted on that machine. - Agents report their `address`, and received DMs carry the sender's as `message.agent_address`, so agents can reply by address. ## [8.12.0] - 2026-09-24 diff --git a/README.md b/README.md index e81c1b57..f01406ef 100644 --- a/README.md +++ b/README.md @@ -609,7 +609,6 @@ GET /channels/:name/messages GET /sessions/:session_ref/messages?limit=<1-500>&after= POST /messages/:id/replies POST /dm -POST /to/:address DM an `agent@machine` address (body: `{ "text": "..." }`) GET /dm/conversations?limit=<1-100> List the agent's newest DM conversations (limit optional) GET /inbox GET /search @@ -625,18 +624,17 @@ joined or DMs the agent participates in. Activity feed channel-message items include `channel_id` and `channel_name`; DM items include `conversation_id`. -`POST /to/:address` routes by address instead of bare name: `agent@machine` resolves to the -agent only while it is hosted on that machine (its broker node's name or `machine_id`, or -`direct` for a self-connected agent), then delivers like `POST /dm`. A cloud sandbox is a -broker node, so a sandboxed agent's address uses the sandbox's node name; once the sandbox -is torn down the agent has no address (`address: null`) until it is hosted again. Agents expose their -address as `address` on agent resources, and each DM carries the sender's as -`message.agent_address`, so a recipient can reply on it. A stale address returns -`404 address_not_found`, including when the agent moves while the send is in flight; a malformed -one returns `400 invalid_address`. Names may contain `@`: each `@` is tried as the separator, and an -address that reads as two different agents returns `400 ambiguous_address`. Addressed DMs share -the `POST /dm` rate-limit bucket. An idempotent retry -replays even if the agent has moved; reusing the key for another address is a `409`. +`POST /dm` accepts `address` (`agent@machine`) in place of `to`. It resolves to the agent only +while it is hosted on that machine (its broker node's name or `machine_id`, or `direct` for a +self-connected agent), then delivers like any DM. A cloud sandbox is a broker node, so a sandboxed +agent's address uses the sandbox's node name; once the sandbox is torn down the agent has no address +(`address: null`) until it is hosted again. Agents expose their address as `address` on agent +resources, and each DM carries the sender's as `message.agent_address`, so a recipient can reply on +it. A stale address returns `404 address_not_found`, including when the agent moves while the send +is in flight; a malformed one returns `400 invalid_address`. Names may contain `@`: each `@` is +tried as the separator, and an address that reads as two different agents returns +`400 ambiguous_address`. An idempotent retry replays even if the agent has moved; reusing the key +for another address, or for a send by name, is a `409`. `POST /dm` can return **`409 dm_conversation_id_collision`**. A 1:1 conversation id is derived deterministically from `(workspace, sorted agent pair)`, and that binding is reserved diff --git a/openapi.yaml b/openapi.yaml index d867e8e4..55c70507 100644 --- a/openapi.yaml +++ b/openapi.yaml @@ -324,7 +324,7 @@ components: type: string nullable: true description: >- - `agent@machine` for `POST /to/{address}`. `machine` is the broker + `agent@machine` for `address` on `POST /dm`. `machine` is the broker node's name (for a cloud sandbox, its node name), or `direct` for a self-connected agent. Null when the agent is not hosted anywhere: released, finished, or its node was deleted. @@ -814,7 +814,7 @@ components: description: Sender identity type for the DM message actor agent_address: type: string - description: Sender's `agent@machine` address at send time; reply via `POST /to/{address}` + description: Sender's `agent@machine` address at send time; reply with `address` on `POST /dm` text: type: string injection_mode: @@ -876,7 +876,7 @@ components: description: Sender identity type for the message actor agent_address: type: string - description: Sender's `agent@machine` address at send time; reply via `POST /to/{address}` + description: Sender's `agent@machine` address at send time; reply with `address` on `POST /dm` text: type: string injection_mode: @@ -3854,11 +3854,23 @@ paths: post: summary: Send DM description: | - Send a direct message to another agent. + Send a direct message to another agent. Name the recipient with + exactly one of `to` or `address`. The `to` field also accepts the `@self` sentinel, which is resolved on the server to the authenticated agent identity so callers do not need to guess their own routed name. + + `address` is `agent@machine`: `machine` must be the agent's current + broker node (by node name or `machine_id`), or `direct` for a + self-connected agent. An address that no longer matches — the agent + moved, was released, or its node was deleted — returns + `404 address_not_found` instead of reaching the agent elsewhere, and + the check is repeated inside the admission write. Names may contain + `@`: every `@` is tried as the separator. `direct` never matches a + broker; a broker named `direct` is addressed by its `machine_id`. + Agents learn addresses from `address` on agent resources and from + `message.agent_address` on DMs they receive. tags: - Direct Messages security: @@ -3875,12 +3887,14 @@ paths: schema: type: object required: - - to - text properties: to: type: string - description: Recipient agent name or `@self` + description: Recipient agent name or `@self`. Exactly one of `to` or `address`. + address: + type: string + description: Recipient `agent@machine`. Exactly one of `to` or `address`. text: type: string attachments: @@ -3900,100 +3914,6 @@ paths: enum: [wait, steer] default: wait description: Injection mode for the DM - responses: - '201': - description: DM sent - content: - application/json: - schema: - type: object - properties: - ok: - type: boolean - data: - $ref: '#/components/schemas/DmSendResponse' - '409': - description: | - `dm_conversation_id_collision` — the deterministic 1:1 conversation - identifier could not be reserved for this exact - `(workspace, sorted agent pair)` tuple. A 1:1 DM id is derived from - that tuple, and the reservation is atomic, so the request fails - closed rather than resolving to another pair's conversation. Two - cases produce it: the identifier is already bound to a different - pair, or the pair is already bound to a different identifier. It is - not retryable — the same inputs will collide again. - content: - application/json: - schema: - $ref: '#/components/schemas/ErrorResponse' - - /to/{address}: - post: - summary: Send DM to agent@machine - description: | - Send a direct message to an `agent@machine` address. `machine` must be - the agent's current broker node (a cloud sandbox is one), matched by - node name or `machine_id`, or `direct` for a self-connected agent. An - agent whose node was deleted has no address. The message is then - delivered exactly like `POST /dm`. An address that no longer matches - (the agent moved or was released) returns 404 instead of reaching the - agent elsewhere; matching is exact and case-sensitive. The check is - repeated inside the admission write, so a move that races the send - also returns 404. - - Agent names, node names and `machine_id` values may contain `@`: every - `@` is tried as the separator, and the one reading that names an agent - on that machine is used. `direct` is reserved and never matches a - broker; a broker named `direct` is addressed by its `machine_id`. - - Agents learn addresses from `address` on agent resources and from - `message.agent_address` on DMs they receive. - - With an `Idempotency-Key`, a retry of an accepted send replays the - original result even if the agent has since moved. The key is bound - to the address: reusing it for a different address or for `POST /dm` - returns `409 idempotency_key_reused`. - tags: - - Direct Messages - security: - - agentToken: [] - parameters: - - name: address - in: path - required: true - description: '`agent@machine`; split on the last `@`' - schema: - type: string - - name: Idempotency-Key - in: header - schema: - type: string - requestBody: - required: true - content: - application/json: - schema: - type: object - required: - - text - properties: - text: - type: string - attachments: - type: array - items: - type: string - description: Optional file ids to attach to this DM - data: - type: object - nullable: true - additionalProperties: true - description: Public structured message metadata - mode: - type: string - enum: [wait, steer] - default: wait - description: Injection mode for the DM responses: '201': description: DM sent @@ -4008,23 +3928,32 @@ paths: $ref: '#/components/schemas/DmSendResponse' '400': description: | - `invalid_address` — the address has no `agent@machine` reading; + Invalid body (neither or both of `to`/`address`, missing `text`); + `invalid_address` — `address` has no `agent@machine` reading; `ambiguous_address` — more than one reading names an agent on - that machine; or the body is invalid (missing or empty `text`). + that machine. content: application/json: schema: $ref: '#/components/schemas/ErrorResponse' '404': - description: '`address_not_found` — no agent with that name is hosted on that machine' + description: | + `agent_not_found` for `to`; `address_not_found` when no agent with + that name is currently hosted on that machine. content: application/json: schema: $ref: '#/components/schemas/ErrorResponse' '409': description: | - `idempotency_key_reused` or `dm_conversation_id_collision`, with - the same semantics as `POST /dm`. + `dm_conversation_id_collision` — the deterministic 1:1 conversation + identifier could not be reserved for this exact + `(workspace, sorted agent pair)` tuple. A 1:1 DM id is derived from + that tuple, and the reservation is atomic, so the request fails + closed rather than resolving to another pair's conversation. Two + cases produce it: the identifier is already bound to a different + pair, or the pair is already bound to a different identifier. It is + not retryable — the same inputs will collide again. content: application/json: schema: diff --git a/packages/engine/CHANGELOG.md b/packages/engine/CHANGELOG.md index af33730e..1dd385d6 100644 --- a/packages/engine/CHANGELOG.md +++ b/packages/engine/CHANGELOG.md @@ -11,7 +11,7 @@ and this project follows [Semantic Versioning](https://semver.org/spec/v2.0.0.ht ### Added -- `POST /v1/to/:address` resolves `agent@machine` (broker node name or `machine_id`, or `direct` for agents not on a broker) to the agent hosted there and sends it a DM; stale addresses return `404 address_not_found`, idempotent retries replay even after the agent moves, a move that races the send is rejected atomically, names containing `@` resolve (`400 ambiguous_address` when two readings match), and addressed DMs share the `POST /v1/dm` rate-limit bucket. +- `POST /v1/dm` accepts `address` (`agent@machine`) in place of `to`, delivering only while the agent is hosted on that machine (broker node name or `machine_id`, or `direct` for self-connected agents). Stale addresses return `404 address_not_found`, including a move that races the send; idempotent retries replay even after the agent moves; names containing `@` resolve, with `400 ambiguous_address` when two readings match. - Agent resources include `address` (null when the agent is not hosted anywhere, e.g. after its sandbox node is deleted); DM responses, `dm.received` deliveries, and DM history include the sender's `message.agent_address`. ## [8.12.0] - 2026-09-24 diff --git a/packages/engine/src/__tests__/conformance/addressedSend.test.ts b/packages/engine/src/__tests__/conformance/addressedSend.test.ts index ea60e184..46c7eb67 100644 --- a/packages/engine/src/__tests__/conformance/addressedSend.test.ts +++ b/packages/engine/src/__tests__/conformance/addressedSend.test.ts @@ -7,9 +7,9 @@ import { addressSplits, SENDER_ADDRESS_METADATA_KEY } from '../../engine/address type Json = Record; /** - * POST /v1/to/:address — send a DM to `agent@machine`, where `machine` is the - * agent's current broker node (by node name or machine_id), or `direct` for an - * agent not hosted on a broker. + * POST /v1/dm with `address` — send a DM to `agent@machine`, where `machine` is + * the agent's current broker node (by node name or machine_id), or `direct` for + * a self-connected agent. */ describe('addressed send', () => { let stack: TestStack; @@ -66,11 +66,11 @@ describe('addressed send', () => { return { ws, alice, bob, laptop }; } - function send(token: string, address: string, body: unknown, headers: Record = {}) { - return stack.app.request(`/v1/to/${encodeURIComponent(address)}`, { + function send(token: string, address: string, body: Json, headers: Record = {}) { + return stack.app.request('/v1/dm', { method: 'POST', headers: { 'content-type': 'application/json', authorization: `Bearer ${token}`, ...headers }, - body: JSON.stringify(body), + body: JSON.stringify({ address, ...body }), }); } @@ -167,14 +167,16 @@ describe('addressed send', () => { .toEqual(['hi via bob@laptop', 'hi via bob@mach-123']); }); - it('accepts an unencoded @ in the path', async () => { + it('requires exactly one of to or address', async () => { const { alice } = await seed(); - const res = await stack.app.request('/v1/to/bob@laptop', { - method: 'POST', - headers: { 'content-type': 'application/json', authorization: `Bearer ${alice.token}` }, - body: JSON.stringify({ text: 'raw' }), - }); - expect(res.status).toBe(201); + for (const body of [{ text: 'x' }, { to: 'bob', address: 'bob@laptop', text: 'x' }]) { + const res = await stack.app.request('/v1/dm', { + method: 'POST', + headers: { 'content-type': 'application/json', authorization: `Bearer ${alice.token}` }, + body: JSON.stringify(body), + }); + expect(res.status, JSON.stringify(body)).toBe(400); + } }); it('addresses agents without a broker as agent@direct', async () => { @@ -318,7 +320,7 @@ describe('addressed send', () => { expect(sent).toHaveLength(1); }); - it('rejects reusing a key for a different address or for /v1/dm', async () => { + it('rejects reusing a key for a different address or for a send by name', async () => { const { alice } = await seed(); const key = { 'Idempotency-Key': 'retry-2' }; expect((await send(alice.token, 'bob@laptop', { text: 'same' }, key)).status).toBe(201); diff --git a/packages/engine/src/__tests__/conformance/rateLimitContract.test.ts b/packages/engine/src/__tests__/conformance/rateLimitContract.test.ts index 7b5100d7..6ba66ad6 100644 --- a/packages/engine/src/__tests__/conformance/rateLimitContract.test.ts +++ b/packages/engine/src/__tests__/conformance/rateLimitContract.test.ts @@ -1,7 +1,7 @@ import { afterEach, describe, expect, it } from 'vitest'; import type { EntitlementsProvider, PlanLimits, Workspace } from '../../ports/index.js'; import { InProcessRateLimiter } from '../../adapters/node/rate-limit.js'; -import { createWorkspace, makeNodeStack, registerAgent, type TestStack } from './harness.js'; +import { createWorkspace, makeNodeStack, type TestStack } from './harness.js'; function entitlementsWithRate(ratePerMin: number): EntitlementsProvider { return { @@ -71,26 +71,6 @@ describe('rate limit contract', () => { expect(throttled.headers.get('X-RateLimit-Remaining')).toBe('0'); }); - it('charges an addressed DM to the same bucket as POST /v1/dm', async () => { - stack = makeNodeStack({ entitlements: entitlementsWithRate(10) }); - const ws = await createWorkspace(stack.app, 'addressed-dm-bucket-ws'); - const alice = await registerAgent(stack.app, ws.workspaceKey, 'alice'); - await registerAgent(stack.app, ws.workspaceKey, 'bob'); - const post = (path: string, body: unknown) => (stack as TestStack).app.request(path, { - method: 'POST', - headers: { 'content-type': 'application/json', authorization: `Bearer ${alice.token}` }, - body: JSON.stringify(body), - }); - - const dm = await post('/v1/dm', { to: 'bob', text: 'one' }); - expect(dm.headers.get('X-RateLimit-Limit')).toBe('5'); - expect(dm.headers.get('X-RateLimit-Remaining')).toBe('4'); - const addressed = await post('/v1/to/bob%40direct', { text: 'two' }); - expect(addressed.status).toBe(201); - expect(addressed.headers.get('X-RateLimit-Limit')).toBe('5'); - expect(addressed.headers.get('X-RateLimit-Remaining')).toBe('3'); - }); - it('reports the window reset on a successful response', async () => { stack = makeNodeStack({ entitlements: entitlementsWithRate(10) }); const ws = await createWorkspace(stack.app, 'reset-header-ws'); diff --git a/packages/engine/src/middleware/rateLimit.ts b/packages/engine/src/middleware/rateLimit.ts index ad7b9fac..f52a4d1d 100644 --- a/packages/engine/src/middleware/rateLimit.ts +++ b/packages/engine/src/middleware/rateLimit.ts @@ -34,9 +34,6 @@ function getRouteKey(method: string, path: string): string | null { // Normalize path: /v1/channels/foo/messages -> /channels/*/messages const normalized = path .replace(/^\/v1/, '') - // An addressed DM is a DM: it shares the POST:/dm bucket so the new path - // cannot add a second DM budget. - .replace(/^\/to\/[^/]+$/, '/dm') .replace(/\/[a-zA-Z0-9_-]+\/messages/, '/*/messages') .replace(/\/[a-zA-Z0-9_-]+\/reactions/, '/*/reactions') .replace(/\/[a-zA-Z0-9_-]+\/replies/, '/*/replies'); diff --git a/packages/engine/src/routes/dm.ts b/packages/engine/src/routes/dm.ts index 9dbea3aa..0fa8bca6 100644 --- a/packages/engine/src/routes/dm.ts +++ b/packages/engine/src/routes/dm.ts @@ -1,4 +1,4 @@ -import { Hono, type Context } from 'hono'; +import { Hono } from 'hono'; import { z } from 'zod'; import type { AppEnv } from '../env.js'; import { requireAgentToken } from '../middleware/auth.js'; @@ -21,145 +21,23 @@ import { parsePaginationQuery, positiveIntQueryParam } from '../lib/httpQuery.js export const dmRoutes = new Hono(); +// Exactly one recipient: `to` (agent name or `@self`), or `address` +// (`agent@machine`, which also requires the agent to be on that machine). const sendDmSchema = z.object({ - to: z.string().min(1), + to: z.string().min(1).optional(), + address: z.string().min(1).optional(), text: z.string().min(1), attachments: z.array(z.string()).optional(), data: z.record(z.string(), z.unknown()).nullable().optional(), mode: z.enum(['wait', 'steer']).default('wait'), +}).refine((body) => (body.to === undefined) !== (body.address === undefined), { + path: ['to'], }); const listDmConversationsQuerySchema = z.object({ limit: positiveIntQueryParam({ max: 100 }), }); -// Body of POST /v1/to/:address — the recipient comes from the path. -const sendAddressedSchema = sendDmSchema.omit({ to: true }); - -/** - * Send a DM from the authenticated agent. Shared by POST /v1/dm (recipient by - * name) and POST /v1/to/:address (recipient by `agent@machine`, passed as - * `address` and checked by the engine against the recipient it sends to). - */ -async function sendDirectMessage(c: Context, input: z.infer, address?: string) { - const db = c.get('db'); - const workspace = c.get('workspace'); - const agent = c.get('agent'); - const { to, text, attachments, data, mode } = input; - const normalizedAttachments = attachments && attachments.length > 0 ? attachments : undefined; - // `data` is digested rather than embedded. It is caller-supplied and can - // be large — a Ratify proof bundle runs to MAX_PROOF_BUNDLE_BYTES (128 - // KiB) — and the fingerprint is serialized into the stored idempotency - // record, kept for the TTL, and string-compared on every replay. Inlining - // it put ~256 KiB per DM into the KV record and made each replay compare - // the whole payload. A digest answers the only question the fingerprint - // asks — "is this the same request?" — in constant size. - // An addressed send fingerprints its address, so reusing a key for a - // different address is a conflict rather than a replay of the first. - const fingerprintBody = { - to, - ...(address !== undefined ? { address } : {}), - text, - ...(normalizedAttachments ? { attachments: normalizedAttachments } : {}), - ...(data !== undefined ? { data_sha256: await sha256Hex(JSON.stringify(data)) } : {}), - }; - - const { key: idempotencyKey, error: idempotencyError } = parseIdempotencyKey(c.req.header('Idempotency-Key')); - if (idempotencyError) { - return jsonError(c, 'invalid_idempotency_key', idempotencyError, 400); - } - - const mailbox = resolveMailboxConfig(c.get('engine').config, workspace.id); - const toDmReceivedEventData = (data: Awaited>) => buildDmReceivedEventData(data, { - fromName: agent!.name, - }); - - const trackDmSent = (data: { conversation_id: string; id: string }) => emitServerEvent(c, workspace.id, 'relaycast_server_dm_sent', { - conversation_id: data.conversation_id, - message_id: data.id, - from_agent_id: agent!.id, - to_agent_name: to, - }); - - const idempotent = await runIdempotent({ - workspaceId: workspace.id, - actorId: agent!.id, - scope: 'dm:direct', - key: idempotencyKey, - status: 201, - // Backward compatibility: historical fingerprint excluded mode (equivalent to wait). - // Only include mode when explicit steer is requested. - fingerprint: mode === 'steer' - ? JSON.stringify({ ...fingerprintBody, mode }) - : JSON.stringify(fingerprintBody), - kv: c.get('engine').kv, - operation: () => dmEngine.sendDm(db, workspace.id, agent!.id, { - to, - text, - attachments: normalizedAttachments, - data, - mode, - }, { mailbox, resolveWorkspaceDeliveryPolicy: () => resolveWorkspaceDeliveryPolicyFor(c.get('engine').config, workspace), idempotencyKey, address, - afterAdmission: (data, event) => { - runInBackground(c, c.get('engine').realtime.publishToWorkspaceStream({ - workspaceId: workspace.id, event: { ...event.payload, seq: event.seq }, - }), 'publish admitted dm.received'); - runInBackground(c, c.get('engine').webhookQueue.send({ - type: 'dm.received', workspaceId: workspace.id, - data: event.data, outboxId: event.outboxId, - }), 'queue admitted dm.received'); - if (data._delivery) runInBackground(c, - routeDeliveryOutcomes(c, [data._delivery], 'dm.received', event.data), - 'route admitted dm delivery'); - if (data._delivery_rejections.length) runInBackground(c, - notifyDeliveryRejections(c, agent!.id, data._delivery_rejections), - 'notify admitted dm delivery rejection'); - trackDmSent(data); - }, - }), - afterOperation: async (data) => { - if (data._notifications_durable) return; - await sendWebhookEvent(c, { - type: 'dm.received', - workspaceId: workspace.id, - data: toDmReceivedEventData(data), - }); - }, - }); - - if (!idempotent.replayed && !idempotent.data._notifications_durable) { - const { - _delivery, - _delivery_rejections, - ...publicDmData - } = idempotent.data as typeof idempotent.data & { - _delivery?: Parameters[1][number] | null; - _delivery_rejections?: Parameters[2]; - }; - const eventData = toDmReceivedEventData(idempotent.data); - runInBackground(c, publishWorkspaceEvent(c, 'dm.received', eventData), 'publish dm.received'); - - if (_delivery) { - runInBackground( - c, - routeDeliveryOutcomes(c, [_delivery], 'dm.received', eventData), - 'route dm delivery', - ); - } - if (_delivery_rejections && _delivery_rejections.length > 0) { - runInBackground( - c, - notifyDeliveryRejections(c, agent!.id, _delivery_rejections), - 'fanout delivery rejected', - ); - } - - trackDmSent(publicDmData); - } - - return jsonIdempotentOk(c, idempotent); -} - // POST /v1/dm - send a DM dmRoutes.post( '/dm', @@ -167,11 +45,14 @@ dmRoutes.post( rateLimit, async (c) => { try { + const db = c.get('db'); + const workspace = c.get('workspace'); + const agent = c.get('agent'); const parsed = await parseJsonBody(c, sendDmSchema, (failure) => { const hasToIssue = failure.error.issues.some((issue) => issue.path[0] === 'to'); const hasTextIssue = failure.error.issues.some((issue) => issue.path[0] === 'text'); return hasToIssue - ? '"to" agent name is required' + ? 'exactly one of "to" (agent name) or "address" (agent@machine) is required' : hasTextIssue ? 'text is required' : 'invalid dm body'; @@ -179,29 +60,123 @@ dmRoutes.post( if (!parsed.ok) { return parsed.response; } - return await sendDirectMessage(c, parsed.data); - } catch (err: unknown) { - return errorResponse(c, err); - } - }, -); + const { address, text, attachments, data, mode } = parsed.data; + if (address !== undefined) requireAgentAddress(address); + // The engine resolves an addressed recipient from `address`; `to` then + // only names the request in the idempotency fingerprint. + const to = parsed.data.to ?? address!; + const normalizedAttachments = attachments && attachments.length > 0 ? attachments : undefined; + // `data` is digested rather than embedded. It is caller-supplied and can + // be large — a Ratify proof bundle runs to MAX_PROOF_BUNDLE_BYTES (128 + // KiB) — and the fingerprint is serialized into the stored idempotency + // record, kept for the TTL, and string-compared on every replay. Inlining + // it put ~256 KiB per DM into the KV record and made each replay compare + // the whole payload. A digest answers the only question the fingerprint + // asks — "is this the same request?" — in constant size. + // An addressed send fingerprints its address, so reusing a key for a + // different address (or for a send by name) is a conflict, not a replay. + const fingerprintBody = { + to, + ...(address !== undefined ? { address } : {}), + text, + ...(normalizedAttachments ? { attachments: normalizedAttachments } : {}), + ...(data !== undefined ? { data_sha256: await sha256Hex(JSON.stringify(data)) } : {}), + }; + + const { key: idempotencyKey, error: idempotencyError } = parseIdempotencyKey(c.req.header('Idempotency-Key')); + if (idempotencyError) { + return jsonError(c, 'invalid_idempotency_key', idempotencyError, 400); + } -// POST /v1/to/:address - send a DM to an `agent@machine` address -dmRoutes.post( - '/to/:address', - requireAgentToken, - rateLimit, - async (c) => { - try { - const parsed = await parseJsonBody(c, sendAddressedSchema, 'text is required'); - if (!parsed.ok) { - return parsed.response; + const mailbox = resolveMailboxConfig(c.get('engine').config, workspace.id); + const toDmReceivedEventData = (data: Awaited>) => buildDmReceivedEventData(data, { + fromName: agent!.name, + }); + + const trackDmSent = (data: { conversation_id: string; id: string }) => emitServerEvent(c, workspace.id, 'relaycast_server_dm_sent', { + conversation_id: data.conversation_id, + message_id: data.id, + from_agent_id: agent!.id, + to_agent_name: to, + }); + + const idempotent = await runIdempotent({ + workspaceId: workspace.id, + actorId: agent!.id, + scope: 'dm:direct', + key: idempotencyKey, + status: 201, + // Backward compatibility: historical fingerprint excluded mode (equivalent to wait). + // Only include mode when explicit steer is requested. + fingerprint: mode === 'steer' + ? JSON.stringify({ ...fingerprintBody, mode }) + : JSON.stringify(fingerprintBody), + kv: c.get('engine').kv, + operation: () => dmEngine.sendDm(db, workspace.id, agent!.id, { + to, + text, + attachments: normalizedAttachments, + data, + mode, + }, { mailbox, resolveWorkspaceDeliveryPolicy: () => resolveWorkspaceDeliveryPolicyFor(c.get('engine').config, workspace), idempotencyKey, address, + afterAdmission: (data, event) => { + runInBackground(c, c.get('engine').realtime.publishToWorkspaceStream({ + workspaceId: workspace.id, event: { ...event.payload, seq: event.seq }, + }), 'publish admitted dm.received'); + runInBackground(c, c.get('engine').webhookQueue.send({ + type: 'dm.received', workspaceId: workspace.id, + data: event.data, outboxId: event.outboxId, + }), 'queue admitted dm.received'); + if (data._delivery) runInBackground(c, + routeDeliveryOutcomes(c, [data._delivery], 'dm.received', event.data), + 'route admitted dm delivery'); + if (data._delivery_rejections.length) runInBackground(c, + notifyDeliveryRejections(c, agent!.id, data._delivery_rejections), + 'notify admitted dm delivery rejection'); + trackDmSent(data); + }, + }), + afterOperation: async (data) => { + if (data._notifications_durable) return; + await sendWebhookEvent(c, { + type: 'dm.received', + workspaceId: workspace.id, + data: toDmReceivedEventData(data), + }); + }, + }); + + if (!idempotent.replayed && !idempotent.data._notifications_durable) { + const { + _delivery, + _delivery_rejections, + ...publicDmData + } = idempotent.data as typeof idempotent.data & { + _delivery?: Parameters[1][number] | null; + _delivery_rejections?: Parameters[2]; + }; + const eventData = toDmReceivedEventData(idempotent.data); + runInBackground(c, publishWorkspaceEvent(c, 'dm.received', eventData), 'publish dm.received'); + + if (_delivery) { + runInBackground( + c, + routeDeliveryOutcomes(c, [_delivery], 'dm.received', eventData), + 'route dm delivery', + ); + } + if (_delivery_rejections && _delivery_rejections.length > 0) { + runInBackground( + c, + notifyDeliveryRejections(c, agent!.id, _delivery_rejections), + 'fanout delivery rejected', + ); + } + + trackDmSent(publicDmData); } - const address = c.req.param('address'); - requireAgentAddress(address); - // The engine resolves the recipient from `address`; `to` only feeds the - // idempotency fingerprint. - return await sendDirectMessage(c, { ...parsed.data, to: address }, address); + + return jsonIdempotentOk(c, idempotent); } catch (err: unknown) { return errorResponse(c, err); } diff --git a/packages/sdk-typescript/CHANGELOG.md b/packages/sdk-typescript/CHANGELOG.md index e7ae7865..9d4ce98a 100644 --- a/packages/sdk-typescript/CHANGELOG.md +++ b/packages/sdk-typescript/CHANGELOG.md @@ -11,7 +11,7 @@ and this project follows [Semantic Versioning](https://semver.org/spec/v2.0.0.ht ### Added -- `agent.sendTo(address, text)` sends a DM to an `agent@machine` address via `POST /v1/to/:address`. +- `agent.sendTo(address, text)` sends a DM to an `agent@machine` address (`POST /v1/dm` with `address`). ## [8.12.0] - 2026-09-24 diff --git a/packages/sdk-typescript/src/__tests__/agent-messaging.test.ts b/packages/sdk-typescript/src/__tests__/agent-messaging.test.ts index 69e91014..73dfb3a9 100644 --- a/packages/sdk-typescript/src/__tests__/agent-messaging.test.ts +++ b/packages/sdk-typescript/src/__tests__/agent-messaging.test.ts @@ -253,15 +253,15 @@ describe('AgentClient', () => { }); describe('sendTo()', () => { - it('sends DM via POST /v1/to/:address', async () => { + it('sends an addressed DM via POST /v1/dm', async () => { mockFetch.mockImplementation(() => mockResponse({ id: 'dm_1' })); await me.sendTo('Worker-1@laptop', 'hi'); const [url, init] = mockFetch.mock.calls[0]!; - expect(url).toBe('https://cast.agentrelay.com/v1/to/Worker-1%40laptop'); + expect(url).toBe('https://cast.agentrelay.com/v1/dm'); expect(init.method).toBe('POST'); - expect(init.body).toBe(JSON.stringify({ text: 'hi', mode: 'wait' })); + expect(init.body).toBe(JSON.stringify({ address: 'Worker-1@laptop', text: 'hi', mode: 'wait' })); }); }); diff --git a/packages/sdk-typescript/src/agent.ts b/packages/sdk-typescript/src/agent.ts index f671a3fb..0e03d338 100644 --- a/packages/sdk-typescript/src/agent.ts +++ b/packages/sdk-typescript/src/agent.ts @@ -626,13 +626,14 @@ export class AgentClient { data?: Record | null; }), ): Promise { - const body = { + const body: SendDmRequest = { + address, text, ...(opts?.attachments ? { attachments: opts.attachments } : {}), ...(opts?.data !== undefined ? { data: opts.data } : {}), mode: opts?.mode ?? 'wait', }; - return this.client.post(`/v1/to/${encodeURIComponent(address)}`, body, idempotencyHeaders(opts)); + return this.client.post('/v1/dm', body, idempotencyHeaders(opts)); } dms = { diff --git a/packages/types/src/agent.ts b/packages/types/src/agent.ts index 268578cf..80334b9b 100644 --- a/packages/types/src/agent.ts +++ b/packages/types/src/agent.ts @@ -14,7 +14,7 @@ export const AgentSchema = z.object({ id: z.string(), name: z.string(), /** - * `agent@machine` for POST /v1/to/:address; `machine` is `direct` for a + * `agent@machine` for `address` on POST /v1/dm; `machine` is `direct` for a * self-connected agent. Null when the agent is not hosted anywhere. */ address: z.string().nullable().optional(), diff --git a/packages/types/src/dm.ts b/packages/types/src/dm.ts index fc863b4c..5b0ed320 100644 --- a/packages/types/src/dm.ts +++ b/packages/types/src/dm.ts @@ -25,8 +25,10 @@ export type DmParticipant = z.infer; export const DmInjectionModeSchema = MessageInjectionModeSchema; +/** Exactly one of `to` (agent name) or `address` (`agent@machine`). */ export const SendDmRequestSchema = z.object({ - to: z.string(), + to: z.string().optional(), + address: z.string().optional(), text: z.string(), attachments: z.array(z.string()).optional(), data: z.record(z.string(), z.unknown()).nullable().optional(), From d2b9447f64466c9feaf504ea734315eb63c1250e Mon Sep 17 00:00:00 2001 From: khaliqgant Date: Thu, 24 Sep 2026 23:08:12 -0700 Subject: [PATCH 6/8] feat(engine): include sender agent_address on GET /v1/deliveries items Consumers that read the delivery queue (relay-desktop's injectable sessions) need the sender's address to reply by address, like the dm.received payloads already carry. Co-Authored-By: Claude Opus 5.5 (1M context) --- openapi.yaml | 3 +++ packages/engine/CHANGELOG.md | 2 +- .../engine/src/__tests__/conformance/addressedSend.test.ts | 5 +++++ packages/engine/src/engine/delivery.ts | 2 ++ packages/types/src/delivery.ts | 2 ++ 5 files changed, 13 insertions(+), 1 deletion(-) diff --git a/openapi.yaml b/openapi.yaml index 55c70507..8822355a 100644 --- a/openapi.yaml +++ b/openapi.yaml @@ -1073,6 +1073,9 @@ components: agent_name: type: string nullable: true + agent_address: + type: string + description: Sender's `agent@machine` at send time; present on direct messages text: type: string thread_id: diff --git a/packages/engine/CHANGELOG.md b/packages/engine/CHANGELOG.md index 1dd385d6..0120b1cd 100644 --- a/packages/engine/CHANGELOG.md +++ b/packages/engine/CHANGELOG.md @@ -12,7 +12,7 @@ and this project follows [Semantic Versioning](https://semver.org/spec/v2.0.0.ht ### Added - `POST /v1/dm` accepts `address` (`agent@machine`) in place of `to`, delivering only while the agent is hosted on that machine (broker node name or `machine_id`, or `direct` for self-connected agents). Stale addresses return `404 address_not_found`, including a move that races the send; idempotent retries replay even after the agent moves; names containing `@` resolve, with `400 ambiguous_address` when two readings match. -- Agent resources include `address` (null when the agent is not hosted anywhere, e.g. after its sandbox node is deleted); DM responses, `dm.received` deliveries, and DM history include the sender's `message.agent_address`. +- Agent resources include `address` (null when the agent is not hosted anywhere, e.g. after its sandbox node is deleted); DM responses, `dm.received` deliveries, `GET /v1/deliveries` items, and DM history include the sender's `message.agent_address`. ## [8.12.0] - 2026-09-24 diff --git a/packages/engine/src/__tests__/conformance/addressedSend.test.ts b/packages/engine/src/__tests__/conformance/addressedSend.test.ts index 46c7eb67..1cf381bd 100644 --- a/packages/engine/src/__tests__/conformance/addressedSend.test.ts +++ b/packages/engine/src/__tests__/conformance/addressedSend.test.ts @@ -381,6 +381,11 @@ describe('addressed send', () => { const [delivered] = deliveredDms(laptop.sock); expect(delivered.agent_address).toBe('alice@direct'); + const queued = (await (await stack.app.request('/v1/deliveries', { + headers: { authorization: `Bearer ${bob.token}` }, + })).json()) as { data: Array<{ message: Json }> }; + expect(queued.data[0].message.agent_address).toBe('alice@direct'); + const history = (await (await stack.app.request(`/v1/dm/${sent.data.conversation_id}/messages`, { headers: { authorization: `Bearer ${bob.token}` }, })).json()) as { data: Json[] }; diff --git a/packages/engine/src/engine/delivery.ts b/packages/engine/src/engine/delivery.ts index fd26fe89..904d0818 100644 --- a/packages/engine/src/engine/delivery.ts +++ b/packages/engine/src/engine/delivery.ts @@ -144,6 +144,7 @@ export async function listDeliveries( agentId: messages.agentId, agentName: agents.name, body: messages.body, + metadata: messages.metadata, threadId: messages.threadId, createdAt: messages.createdAt, }) @@ -163,6 +164,7 @@ export async function listDeliveries( channel_id: msg.channelId, agent_id: msg.agentId ?? null, agent_name: msg.agentName ?? null, + ...senderAddressField(msg.metadata), text: msg.body, thread_id: msg.threadId ?? null, created_at: msg.createdAt.toISOString(), diff --git a/packages/types/src/delivery.ts b/packages/types/src/delivery.ts index 11341586..af3851a0 100644 --- a/packages/types/src/delivery.ts +++ b/packages/types/src/delivery.ts @@ -17,6 +17,8 @@ export const DeliveryMessageSchema = z.object({ channel_id: z.string(), agent_id: z.string().nullable(), agent_name: z.string().nullable(), + /** Sender's `agent@machine` at send time; present on direct messages. */ + agent_address: z.string().optional(), text: z.string(), thread_id: z.string().nullable(), created_at: z.string(), From dc54b1d60ac4f5b348e2c331f6e3746b0ce0fccc Mon Sep 17 00:00:00 2001 From: khaliqgant Date: Thu, 24 Sep 2026 23:39:52 -0700 Subject: [PATCH 7/8] fix(engine): prefer an @-free machine identifier in published addresses An @ in the machine name can make an address read as a different agent on a different machine; publish the node's machine_id instead when it has no @, so the address has one reading wherever the node offers one. Co-Authored-By: Claude Opus 5.5 (1M context) --- .../src/__tests__/conformance/addressedSend.test.ts | 8 ++++++++ packages/engine/src/engine/address.ts | 12 ++++++++---- 2 files changed, 16 insertions(+), 4 deletions(-) diff --git a/packages/engine/src/__tests__/conformance/addressedSend.test.ts b/packages/engine/src/__tests__/conformance/addressedSend.test.ts index 1cf381bd..7a581000 100644 --- a/packages/engine/src/__tests__/conformance/addressedSend.test.ts +++ b/packages/engine/src/__tests__/conformance/addressedSend.test.ts @@ -122,6 +122,14 @@ describe('addressed send', () => { expect(await errorCode(ambiguous)).toBe('ambiguous_address'); }); + it('publishes a machine identifier without @ when the node has one', async () => { + const { ws } = await seed(); + const node = await enrollBroker(ws, 'node_email_alias', 'ops@example.com', 'mach-ops'); + await registerOnNode(node, 'x'); + const res = await stack.app.request('/v1/agents/x', { headers: { authorization: `Bearer ${ws.workspaceKey}` } }); + expect(((await res.json()) as { data: Json }).data.address).toBe('x@mach-ops'); + }); + it('reserves `direct`: a broker named direct is addressed by machine_id, never @direct', async () => { const { ws, alice } = await seed(); const node = await enrollBroker(ws, 'node_named_direct', 'direct', 'mach-d'); diff --git a/packages/engine/src/engine/address.ts b/packages/engine/src/engine/address.ts index e9087a0e..973b106b 100644 --- a/packages/engine/src/engine/address.ts +++ b/packages/engine/src/engine/address.ts @@ -35,14 +35,18 @@ export function addressSplits(address: string): AgentAddress[] { /** * Canonical address: the broker node's name, or `direct` for a direct node. * `direct` is reserved, so a broker whose name is `direct` is addressed by its - * machine_id instead. Null when the agent has no node — released, finished, - * or its node was deleted (a torn-down sandbox) — because nothing could - * deliver to it, or when a broker has no usable identifier. + * machine_id instead. A machine identifier without `@` is preferred, so the + * address has one reading wherever the node offers one (an `@` in the machine + * can make the address read as a different agent on a different machine). + * Null when the agent has no node — released, finished, or its node was + * deleted (a torn-down sandbox) — because nothing could deliver to it, or + * when a broker has no usable identifier. */ export function formatAgentAddress(agentName: string, node: AddressNode): string | null { if (!node) return null; if (node.role === 'direct') return `${agentName}@${DIRECT_MACHINE}`; - const machine = [node.name, node.machineId].find((id) => id && id !== DIRECT_MACHINE); + const usable = [node.name, node.machineId].filter((id): id is string => !!id && id !== DIRECT_MACHINE); + const machine = usable.find((id) => !id.includes('@')) ?? usable[0]; return machine ? `${agentName}@${machine}` : null; } From 554187d1d9cc1bf239ac0cb430fe237a41817a1c Mon Sep 17 00:00:00 2001 From: khaliqgant Date: Thu, 24 Sep 2026 23:46:22 -0700 Subject: [PATCH 8/8] docs(engine): document the machine_id address fallback; report resolved DM recipient - openapi/README: a published address uses the node's machine_id when its name contains @ or is `direct`. - relaycast_server_dm_sent reports the resolved recipient name rather than the raw agent@machine for addressed sends. Co-Authored-By: Claude Opus 5.5 (1M context) --- README.md | 3 ++- openapi.yaml | 5 +++-- packages/engine/src/routes/dm.ts | 5 +++-- 3 files changed, 8 insertions(+), 5 deletions(-) diff --git a/README.md b/README.md index f01406ef..8d72e03f 100644 --- a/README.md +++ b/README.md @@ -626,7 +626,8 @@ Activity feed channel-message items include `channel_id` and `channel_name`; DM `POST /dm` accepts `address` (`agent@machine`) in place of `to`. It resolves to the agent only while it is hosted on that machine (its broker node's name or `machine_id`, or `direct` for a -self-connected agent), then delivers like any DM. A cloud sandbox is a broker node, so a sandboxed +self-connected agent; the published address uses `machine_id` when the name contains `@` or is +`direct`), then delivers like any DM. A cloud sandbox is a broker node, so a sandboxed agent's address uses the sandbox's node name; once the sandbox is torn down the agent has no address (`address: null`) until it is hosted again. Agents expose their address as `address` on agent resources, and each DM carries the sender's as `message.agent_address`, so a recipient can reply on diff --git a/openapi.yaml b/openapi.yaml index 8822355a..387ce77d 100644 --- a/openapi.yaml +++ b/openapi.yaml @@ -325,8 +325,9 @@ components: nullable: true description: >- `agent@machine` for `address` on `POST /dm`. `machine` is the broker - node's name (for a cloud sandbox, its node name), or `direct` for a - self-connected agent. Null when the agent is not hosted anywhere: + node's name (for a cloud sandbox, its node name), or its + `machine_id` when the name contains `@` or is `direct`; `direct` + for a self-connected agent. Null when the agent is not hosted anywhere: released, finished, or its node was deleted. type: type: string diff --git a/packages/engine/src/routes/dm.ts b/packages/engine/src/routes/dm.ts index 0fa8bca6..50d9aa98 100644 --- a/packages/engine/src/routes/dm.ts +++ b/packages/engine/src/routes/dm.ts @@ -93,11 +93,12 @@ dmRoutes.post( fromName: agent!.name, }); - const trackDmSent = (data: { conversation_id: string; id: string }) => emitServerEvent(c, workspace.id, 'relaycast_server_dm_sent', { + // `data.to` is the resolved recipient name; `to` may be a raw address. + const trackDmSent = (data: { conversation_id: string; id: string; to?: string }) => emitServerEvent(c, workspace.id, 'relaycast_server_dm_sent', { conversation_id: data.conversation_id, message_id: data.id, from_agent_id: agent!.id, - to_agent_name: to, + to_agent_name: data.to ?? to, }); const idempotent = await runIdempotent({