diff --git a/CHANGELOG.md b/CHANGELOG.md index 91fbb0ba..02f8de09 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -16,7 +16,12 @@ 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/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 8dd6cbf4..8d72e03f 100644 --- a/README.md +++ b/README.md @@ -624,6 +624,19 @@ joined or DMs the agent participates in. Activity feed channel-message items include `channel_id` and `channel_name`; DM items include `conversation_id`. +`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; 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 +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 atomically before any conversation state is created, so a derivation that would name another diff --git a/openapi.yaml b/openapi.yaml index e94e91f5..387ce77d 100644 --- a/openapi.yaml +++ b/openapi.yaml @@ -320,6 +320,15 @@ components: locally) can record which workspace it registered into. name: type: string + address: + type: string + 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 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 enum: [agent, human, system] @@ -804,6 +813,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 with `address` on `POST /dm` text: type: string injection_mode: @@ -863,6 +875,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 with `address` on `POST /dm` text: type: string injection_mode: @@ -1059,6 +1074,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: @@ -3840,11 +3858,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: @@ -3861,12 +3891,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: @@ -3898,6 +3930,24 @@ paths: type: boolean data: $ref: '#/components/schemas/DmSendResponse' + '400': + description: | + 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. + content: + application/json: + schema: + $ref: '#/components/schemas/ErrorResponse' + '404': + 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: | `dm_conversation_id_collision` — the deterministic 1:1 conversation diff --git a/packages/engine/CHANGELOG.md b/packages/engine/CHANGELOG.md index 990817aa..0120b1cd 100644 --- a/packages/engine/CHANGELOG.md +++ b/packages/engine/CHANGELOG.md @@ -7,7 +7,12 @@ 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/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, `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 new file mode 100644 index 00000000..7a581000 --- /dev/null +++ b/packages/engine/src/__tests__/conformance/addressedSend.test.ts @@ -0,0 +1,417 @@ +import { afterEach, beforeEach, describe, expect, it } from 'vitest'; +import { eq } from 'drizzle-orm'; +import { attachFakeBatch, makeNodeStack, createWorkspace, registerAgent, FakeSocket, type TestStack } from './harness.js'; +import { agents, messages, nodes } from '../../db/schema.js'; +import { addressSplits, SENDER_ADDRESS_METADATA_KEY } from '../../engine/address.js'; + +type Json = Record; + +/** + * 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; + beforeEach(() => { stack = makeNodeStack({ ttlMs: 60_000 }); }); + afterEach(async () => { + await stack.close(); + }); + + 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, ...(machineId ? { machine_id: machineId } : {}), tags, + 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, 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, + })); + 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'); + const bob = await registerOnNode(laptop, 'bob'); + return { ws, alice, bob, laptop }; + } + + 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({ address, ...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('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('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'); + 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 () => { + 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}` }); + 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(); + expect(deliveredDms(laptop.sock).map((message) => message.text)) + .toEqual(['hi via bob@laptop', 'hi via bob@mach-123']); + }); + + it('requires exactly one of to or address', async () => { + const { alice } = await seed(); + 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 () => { + 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 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('rejects malformed addresses and missing text', async () => { + const { alice } = await seed(); + 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('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(); + 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 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); + + 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'); + + 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 () => { + 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 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[] }; + 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 new file mode 100644 index 00000000..973b106b --- /dev/null +++ b/packages/engine/src/engine/address.ts @@ -0,0 +1,106 @@ +import { nodes } from '../db/schema.js'; +import { codedError } from '../lib/httpError.js'; + +/** + * 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'; + +/** 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; + +/** + * 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 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. + * `direct` is reserved, so a broker whose name is `direct` is addressed by its + * 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 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; +} + +/** + * 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 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 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 splits; +} + +/** + * 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 selectAddressedRecipient( + address: string, + 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]; +} + +export function addressNotFound(address: string) { + return codedError(`No agent at address "${address}"`, 'address_not_found', 404); +} + +/** `{ 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..33eb2ab9 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 { addressNodeSelection, 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: addressNodeSelection }) .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: addressNodeSelection }) .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, @@ -587,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, @@ -636,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/delivery.ts b/packages/engine/src/engine/delivery.ts index b65cfdc8..904d0818 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'; @@ -143,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, }) @@ -162,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(), @@ -595,6 +598,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..db4779dd 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,15 @@ import { buildDmReceivedEventData } from './deliveryWire.js'; import { buildWorkspaceEventWrite } from './workspaceEvents.js'; import { transformForClient } from './wsTransform.js'; import { codedError } from '../lib/httpError.js'; +import { + addressNodeSelection, + addressNotFound, + addressSplits, + formatAgentAddress, + selectAddressedRecipient, + 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 +72,12 @@ interface SendDmOptions { mailbox?: MailboxConfig; /** Server-resolved workspace growth policy; absent => no workspace guard. */ workspaceDeliveryPolicy?: WorkspaceDeliveryPolicy; + /** + * `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; } /** @@ -269,6 +285,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 | null | 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 +322,11 @@ function buildDmMessageWrites( messageId: string, createdAt = new Date(), inboundRegistration?: { id: string; tokenHash?: string }, + senderAddress?: string | null, + addressedRecipient?: { agentId: string; nodeId: 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 @@ -314,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, @@ -356,6 +395,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 +544,37 @@ 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 recipientQuery = db + .select({ agent: agents, node: addressNodeSelection }) + .from(agents) + .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), + )); + 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); } 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,8 +628,8 @@ export async function sendDm( const publicResult = buildDmResult({ id: messageId, agentId: fromAgentId, body: data.text, createdAt, - metadata: { ...sanitizeUserMessageMetadata(data.data), injection_mode: data.mode ?? 'wait' }, - }, conv, fromAgent, data, attachments); + metadata: dmMessageMetadata(data, senderAddress), + }, 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(), @@ -589,7 +638,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, addressedRecipient); if (egressId && a2aTarget && egressPayload) { // First statement owns the request identity; a competing attempt rolls @@ -655,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); } @@ -878,6 +931,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/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 fa5096d1..50d9aa98 100644 --- a/packages/engine/src/routes/dm.ts +++ b/packages/engine/src/routes/dm.ts @@ -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 { requireAgentAddress } from '../engine/address.js'; import { resolveMailboxConfig } from '../engine/mailboxConfig.js'; import { resolveWorkspaceDeliveryPolicyFor } from '../engine/workspaceDeliveryPolicy.js'; import { publishWorkspaceEvent } from './fanout.js'; @@ -20,12 +21,17 @@ 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({ @@ -46,7 +52,7 @@ dmRoutes.post( 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'; @@ -54,7 +60,11 @@ dmRoutes.post( if (!parsed.ok) { return parsed.response; } - const { to, text, attachments, data, mode } = parsed.data; + 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 @@ -63,8 +73,11 @@ dmRoutes.post( // 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)) } : {}), @@ -80,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({ @@ -105,7 +119,7 @@ dmRoutes.post( attachments: normalizedAttachments, data, mode, - }, { mailbox, resolveWorkspaceDeliveryPolicy: () => 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 }, diff --git a/packages/sdk-typescript/CHANGELOG.md b/packages/sdk-typescript/CHANGELOG.md index eb3ca66a..9d4ce98a 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 (`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 24c0e863..73dfb3a9 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 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/dm'); + expect(init.method).toBe('POST'); + expect(init.body).toBe(JSON.stringify({ address: 'Worker-1@laptop', 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..0e03d338 100644 --- a/packages/sdk-typescript/src/agent.ts +++ b/packages/sdk-typescript/src/agent.ts @@ -616,6 +616,26 @@ 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: SendDmRequest = { + address, + text, + ...(opts?.attachments ? { attachments: opts.attachments } : {}), + ...(opts?.data !== undefined ? { data: opts.data } : {}), + mode: opts?.mode ?? 'wait', + }; + return this.client.post('/v1/dm', body, idempotencyHeaders(opts)); + } + dms = { conversations: (opts?: Pick): Promise => { const query: Record = {}; 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..80334b9b 100644 --- a/packages/types/src/agent.ts +++ b/packages/types/src/agent.ts @@ -13,6 +13,11 @@ export type AgentStatus = z.infer; export const AgentSchema = z.object({ id: z.string(), name: z.string(), + /** + * `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(), type: AgentTypeSchema, status: AgentStatusSchema, persona: z.string().nullable(), 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(), 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(), 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(),