diff --git a/.agentworkforce/trajectories/completed/2026-06/traj_h4wj3doqp0ra.trace.json b/.agentworkforce/trajectories/completed/2026-06/traj_h4wj3doqp0ra.trace.json new file mode 100644 index 000000000..36c23e8f6 --- /dev/null +++ b/.agentworkforce/trajectories/completed/2026-06/traj_h4wj3doqp0ra.trace.json @@ -0,0 +1,281 @@ +{ + "version": "1.0.0", + "id": "35467489-5fde-46b4-b686-b372747b0eb9", + "timestamp": "2026-06-19T00:16:10.852Z", + "trajectory": "traj_h4wj3doqp0ra", + "files": [ + { + "path": ".agentworkforce/trajectories/active/traj_h4wj3doqp0ra/trajectory.json", + "conversations": [ + { + "contributor": { + "type": "ai" + }, + "ranges": [ + { + "start_line": 1, + "end_line": 70, + "revision": "afdc289bdcd2b7de860699df39dd0df5eec94d99" + } + ] + } + ] + }, + { + "path": "packages/cli/src/cli/agent-relay-mcp.ts", + "conversations": [ + { + "contributor": { + "type": "ai" + }, + "ranges": [ + { + "start_line": 4, + "end_line": 44, + "revision": "afdc289bdcd2b7de860699df39dd0df5eec94d99" + }, + { + "start_line": 79, + "end_line": 84, + "revision": "afdc289bdcd2b7de860699df39dd0df5eec94d99" + }, + { + "start_line": 284, + "end_line": 289, + "revision": "afdc289bdcd2b7de860699df39dd0df5eec94d99" + }, + { + "start_line": 352, + "end_line": 357, + "revision": "afdc289bdcd2b7de860699df39dd0df5eec94d99" + }, + { + "start_line": 536, + "end_line": 542, + "revision": "afdc289bdcd2b7de860699df39dd0df5eec94d99" + } + ] + } + ] + }, + { + "path": "packages/cli/src/cli/lib/broker-dashboard.ts", + "conversations": [ + { + "contributor": { + "type": "ai" + }, + "ranges": [ + { + "start_line": 1, + "end_line": 570, + "revision": "afdc289bdcd2b7de860699df39dd0df5eec94d99" + } + ] + } + ] + }, + { + "path": "packages/cli/src/cli/lib/broker-lifecycle.ts", + "conversations": [ + { + "contributor": { + "type": "ai" + }, + "ranges": [ + { + "start_line": 1, + "end_line": 4, + "revision": "afdc289bdcd2b7de860699df39dd0df5eec94d99" + }, + { + "start_line": 10, + "end_line": 24, + "revision": "afdc289bdcd2b7de860699df39dd0df5eec94d99" + }, + { + "start_line": 275, + "end_line": 280, + "revision": "afdc289bdcd2b7de860699df39dd0df5eec94d99" + }, + { + "start_line": 700, + "end_line": 705, + "revision": "afdc289bdcd2b7de860699df39dd0df5eec94d99" + } + ] + } + ] + }, + { + "path": "packages/cli/src/cli/mcp/action-schema.ts", + "conversations": [ + { + "contributor": { + "type": "ai" + }, + "ranges": [ + { + "start_line": 1, + "end_line": 157, + "revision": "afdc289bdcd2b7de860699df39dd0df5eec94d99" + } + ] + } + ] + }, + { + "path": "packages/cli/src/cli/mcp/action-tools.ts", + "conversations": [ + { + "contributor": { + "type": "ai" + }, + "ranges": [ + { + "start_line": 1, + "end_line": 151, + "revision": "afdc289bdcd2b7de860699df39dd0df5eec94d99" + } + ] + } + ] + }, + { + "path": "packages/cli/src/cli/mcp/inbox.ts", + "conversations": [ + { + "contributor": { + "type": "ai" + }, + "ranges": [ + { + "start_line": 1, + "end_line": 56, + "revision": "afdc289bdcd2b7de860699df39dd0df5eec94d99" + } + ] + } + ] + }, + { + "path": "packages/cli/src/cli/mcp/messaging-tools.ts", + "conversations": [ + { + "contributor": { + "type": "ai" + }, + "ranges": [ + { + "start_line": 1, + "end_line": 428, + "revision": "afdc289bdcd2b7de860699df39dd0df5eec94d99" + } + ] + } + ] + }, + { + "path": "packages/cli/src/cli/mcp/resources.ts", + "conversations": [ + { + "contributor": { + "type": "ai" + }, + "ranges": [ + { + "start_line": 1, + "end_line": 209, + "revision": "afdc289bdcd2b7de860699df39dd0df5eec94d99" + } + ] + } + ] + }, + { + "path": "packages/cli/src/cli/mcp/telemetry.ts", + "conversations": [ + { + "contributor": { + "type": "ai" + }, + "ranges": [ + { + "start_line": 1, + "end_line": 216, + "revision": "afdc289bdcd2b7de860699df39dd0df5eec94d99" + } + ] + } + ] + }, + { + "path": "packages/cli/src/cli/mcp/tool-results.ts", + "conversations": [ + { + "contributor": { + "type": "ai" + }, + "ranges": [ + { + "start_line": 1, + "end_line": 52, + "revision": "afdc289bdcd2b7de860699df39dd0df5eec94d99" + } + ] + } + ] + }, + { + "path": "packages/cli/src/cli/mcp/tool-shapes.ts", + "conversations": [ + { + "contributor": { + "type": "ai" + }, + "ranges": [ + { + "start_line": 1, + "end_line": 17, + "revision": "afdc289bdcd2b7de860699df39dd0df5eec94d99" + } + ] + } + ] + }, + { + "path": "packages/cli/src/cli/mcp/types.ts", + "conversations": [ + { + "contributor": { + "type": "ai" + }, + "ranges": [ + { + "start_line": 1, + "end_line": 44, + "revision": "afdc289bdcd2b7de860699df39dd0df5eec94d99" + } + ] + } + ] + }, + { + "path": "packages/cli/src/cli/mcp/workspace.ts", + "conversations": [ + { + "contributor": { + "type": "ai" + }, + "ranges": [ + { + "start_line": 1, + "end_line": 48, + "revision": "afdc289bdcd2b7de860699df39dd0df5eec94d99" + } + ] + } + ] + } + ] +} diff --git a/.agentworkforce/trajectories/completed/2026-06/traj_h4wj3doqp0ra/summary.md b/.agentworkforce/trajectories/completed/2026-06/traj_h4wj3doqp0ra/summary.md new file mode 100644 index 000000000..e77c4e63c --- /dev/null +++ b/.agentworkforce/trajectories/completed/2026-06/traj_h4wj3doqp0ra/summary.md @@ -0,0 +1,52 @@ +# Trajectory: Architecture refactor of Relay monorepo + +> **Status:** ✅ Completed +> **Confidence:** 90% +> **Started:** June 18, 2026 at 11:36 PM +> **Completed:** June 19, 2026 at 12:16 AM + +--- + +## Summary + +Decomposed the two largest Agent Relay CLI god files into cohesive single-responsibility modules: agent-relay-mcp.ts (2215->969 lines, 10 mcp/ modules) and broker-lifecycle.ts (1879->1320, dashboard extracted). Three pure-move refactors, each verified with tsc + full test suite (968 pass) + live MCP probe + independent autoreview. No behavior change. + +**Approach:** Standard approach + +--- + +## Key Decisions + +### Decomposed agent-relay-mcp.ts god file into cohesive mcp/ modules + +- **Chose:** Decomposed agent-relay-mcp.ts god file into cohesive mcp/ modules +- **Reasoning:** 2215-line file mixed 8 concerns (schema adapter, telemetry, resources, inbox, workspace, action tools); split into single-responsibility modules behind unchanged public surface. Verified pure move via dual autoreview, tsc, full test suite, and live MCP stdio probe. + +### Extracted dashboard management from broker-lifecycle.ts into broker-dashboard.ts + +- **Chose:** Extracted dashboard management from broker-lifecycle.ts into broker-dashboard.ts +- **Reasoning:** 1879-line broker-lifecycle mixed broker start/stop/status with ~550 lines of dashboard discovery/spawn/asset-refresh logic. Moved the self-contained dashboard cluster behind a 7-function boundary. Verified byte-identical move via structural function diff (69->69), dual checks, tsc, full suite, and direct module live test of all 7 exports. + +### Extracted messaging MCP tools into mcp/messaging-tools.ts + +- **Chose:** Extracted messaging MCP tools into mcp/messaging-tools.ts +- **Reasoning:** registerAgentRelayTools was a 720-line god function; split out the 20 homogeneous messaging tools (channels/messages/threads/DMs/reactions/search/inbox) which depend only on getAgentClient, plus shared zod shapes into mcp/tool-shapes.ts. Main file 2215->969 lines overall. Verified byte-identical move, all 28 tools register live, full suite 968 pass, dual reviewer APPROVE. + +--- + +## Chapters + +### 1. Work + +_Agent: default_ + +- Decomposed agent-relay-mcp.ts god file into cohesive mcp/ modules: Decomposed agent-relay-mcp.ts god file into cohesive mcp/ modules +- Extracted dashboard management from broker-lifecycle.ts into broker-dashboard.ts: Extracted dashboard management from broker-lifecycle.ts into broker-dashboard.ts +- Extracted messaging MCP tools into mcp/messaging-tools.ts: Extracted messaging MCP tools into mcp/messaging-tools.ts + +--- + +## Artifacts + +**Commits:** afdc289, f5e538c, f7715d3 +**Files changed:** 14 diff --git a/.agentworkforce/trajectories/completed/2026-06/traj_h4wj3doqp0ra/trajectory.json b/.agentworkforce/trajectories/completed/2026-06/traj_h4wj3doqp0ra/trajectory.json new file mode 100644 index 000000000..2e4067a55 --- /dev/null +++ b/.agentworkforce/trajectories/completed/2026-06/traj_h4wj3doqp0ra/trajectory.json @@ -0,0 +1,93 @@ +{ + "id": "traj_h4wj3doqp0ra", + "version": 1, + "task": { + "title": "Architecture refactor of Relay monorepo" + }, + "status": "completed", + "startedAt": "2026-06-18T23:36:50.727Z", + "completedAt": "2026-06-19T00:16:10.817Z", + "agents": [ + { + "name": "default", + "role": "lead", + "joinedAt": "2026-06-18T23:58:24.013Z" + } + ], + "chapters": [ + { + "id": "chap_beyia4dq9oyw", + "title": "Work", + "agentName": "default", + "startedAt": "2026-06-18T23:58:24.013Z", + "endedAt": "2026-06-19T00:16:10.817Z", + "events": [ + { + "ts": 1781827104014, + "type": "decision", + "content": "Decomposed agent-relay-mcp.ts god file into cohesive mcp/ modules: Decomposed agent-relay-mcp.ts god file into cohesive mcp/ modules", + "raw": { + "question": "Decomposed agent-relay-mcp.ts god file into cohesive mcp/ modules", + "chosen": "Decomposed agent-relay-mcp.ts god file into cohesive mcp/ modules", + "alternatives": [], + "reasoning": "2215-line file mixed 8 concerns (schema adapter, telemetry, resources, inbox, workspace, action tools); split into single-responsibility modules behind unchanged public surface. Verified pure move via dual autoreview, tsc, full test suite, and live MCP stdio probe." + }, + "significance": "high" + }, + { + "ts": 1781827670970, + "type": "decision", + "content": "Extracted dashboard management from broker-lifecycle.ts into broker-dashboard.ts: Extracted dashboard management from broker-lifecycle.ts into broker-dashboard.ts", + "raw": { + "question": "Extracted dashboard management from broker-lifecycle.ts into broker-dashboard.ts", + "chosen": "Extracted dashboard management from broker-lifecycle.ts into broker-dashboard.ts", + "alternatives": [], + "reasoning": "1879-line broker-lifecycle mixed broker start/stop/status with ~550 lines of dashboard discovery/spawn/asset-refresh logic. Moved the self-contained dashboard cluster behind a 7-function boundary. Verified byte-identical move via structural function diff (69->69), dual checks, tsc, full suite, and direct module live test of all 7 exports." + }, + "significance": "high" + }, + { + "ts": 1781828145275, + "type": "decision", + "content": "Extracted messaging MCP tools into mcp/messaging-tools.ts: Extracted messaging MCP tools into mcp/messaging-tools.ts", + "raw": { + "question": "Extracted messaging MCP tools into mcp/messaging-tools.ts", + "chosen": "Extracted messaging MCP tools into mcp/messaging-tools.ts", + "alternatives": [], + "reasoning": "registerAgentRelayTools was a 720-line god function; split out the 20 homogeneous messaging tools (channels/messages/threads/DMs/reactions/search/inbox) which depend only on getAgentClient, plus shared zod shapes into mcp/tool-shapes.ts. Main file 2215->969 lines overall. Verified byte-identical move, all 28 tools register live, full suite 968 pass, dual reviewer APPROVE." + }, + "significance": "high" + } + ] + } + ], + "retrospective": { + "summary": "Decomposed the two largest Agent Relay CLI god files into cohesive single-responsibility modules: agent-relay-mcp.ts (2215->969 lines, 10 mcp/ modules) and broker-lifecycle.ts (1879->1320, dashboard extracted). Three pure-move refactors, each verified with tsc + full test suite (968 pass) + live MCP probe + independent autoreview. No behavior change.", + "approach": "Standard approach", + "confidence": 0.9 + }, + "commits": ["afdc289", "f5e538c", "f7715d3"], + "filesChanged": [ + ".agentworkforce/trajectories/active/traj_h4wj3doqp0ra/trajectory.json", + "packages/cli/src/cli/agent-relay-mcp.ts", + "packages/cli/src/cli/lib/broker-dashboard.ts", + "packages/cli/src/cli/lib/broker-lifecycle.ts", + "packages/cli/src/cli/mcp/action-schema.ts", + "packages/cli/src/cli/mcp/action-tools.ts", + "packages/cli/src/cli/mcp/inbox.ts", + "packages/cli/src/cli/mcp/messaging-tools.ts", + "packages/cli/src/cli/mcp/resources.ts", + "packages/cli/src/cli/mcp/telemetry.ts", + "packages/cli/src/cli/mcp/tool-results.ts", + "packages/cli/src/cli/mcp/tool-shapes.ts", + "packages/cli/src/cli/mcp/types.ts", + "packages/cli/src/cli/mcp/workspace.ts" + ], + "projectId": "AgentWorkforce/relay", + "tags": [], + "_trace": { + "startRef": "66f972f23cafb958fa1956d3654710423e2391b5", + "endRef": "afdc289bdcd2b7de860699df39dd0df5eec94d99", + "traceId": "35467489-5fde-46b4-b686-b372747b0eb9" + } +} diff --git a/packages/cli/src/cli/agent-relay-mcp.ts b/packages/cli/src/cli/agent-relay-mcp.ts index 3a4f591f0..87254c886 100644 --- a/packages/cli/src/cli/agent-relay-mcp.ts +++ b/packages/cli/src/cli/agent-relay-mcp.ts @@ -4,39 +4,41 @@ import fs from 'node:fs'; import path from 'node:path'; import { fileURLToPath } from 'node:url'; -import { McpServer, ResourceTemplate } from '@modelcontextprotocol/sdk/server/mcp.js'; +import { McpServer } from '@modelcontextprotocol/sdk/server/mcp.js'; import { StdioServerTransport } from '@modelcontextprotocol/sdk/server/stdio.js'; import { ListToolsRequestSchema, SubscribeRequestSchema, UnsubscribeRequestSchema, } from '@modelcontextprotocol/sdk/types.js'; -import { RelayCast, SDK_VERSION, WsClient, type AgentClient } from '@relaycast/sdk'; -import { - INVALID_AGENT_TOKEN_CODE, - agentTokenRecoveryMessage, - isInvalidAgentTokenError, - isInvalidAgentTokenToolResult, -} from '@agent-relay/sdk'; +import { RelayCast, SDK_VERSION, WsClient } from '@relaycast/sdk'; import { AgentRelay } from '@agent-relay/sdk'; -import type { - ActionAuditEvent, - ActionSchema, - AgentRelayActionDescriptor, - AgentRelayActions, - JsonSchemaLiteObject, - ZodLikeSchema, -} from '@agent-relay/sdk/actions'; import { z } from 'zod'; +import { initTelemetry, shutdown as shutdownTelemetry } from './telemetry/index.js'; +import { withRelaycastTelemetry } from './lib/relaycast-telemetry.js'; +import { RealtimeResourceBridge, SubscriptionManager, registerResourceDefinitions } from './mcp/resources.js'; +import { jsonContent, jsonResult, textContent } from './mcp/tool-results.js'; import { - initTelemetry, - shutdown as shutdownTelemetry, - track, - type AgentRelayToolCallCategory, - type AgentRelayToolCallType, -} from './telemetry/index.js'; -import { relaycastWorkspaceTelemetryOptions, withRelaycastTelemetry } from './lib/relaycast-telemetry.js'; -import { errorClassName } from './lib/telemetry-helpers.js'; + createWorkspace, + extractWorkspaceKey, + extractWorkspaceName, + requireWorkspaceKey, +} from './mcp/workspace.js'; +import { enableInboxPiggyback } from './mcp/telemetry.js'; +import { registerAgentRelayActionTools } from './mcp/action-tools.js'; +import { registerMessagingTools } from './mcp/messaging-tools.js'; +import { identityOverrideInputShape, messageResult } from './mcp/tool-shapes.js'; +import type { + AgentClientLike, + AgentRelayMcpServerOptions, + AgentType, + RegisteredAgent, + RegistrationSession, + RelayCastLike, + SessionSetter, + SessionState, +} from './mcp/types.js'; +export type { AgentRelayMcpServerOptions } from './mcp/types.js'; const DEFAULT_BASE_URL = 'https://gateway.relaycast.dev'; export const AGENT_RELAY_MCP_VERSION = process.env.AGENT_RELAY_CLI_VERSION ?? SDK_VERSION ?? 'unknown'; @@ -77,55 +79,6 @@ const DEFAULT_SYSTEM_PROMPT = `You are an AI agent in a collaborative workspace - React with emoji to acknowledge messages - Keep messages concise and actionable`; -const jsonResult = z.object({}).passthrough(); -const messageResult = { - message: z.string().describe('Human-readable confirmation message'), -}; -const identityOverrideInputShape = { - as: z - .string() - .optional() - .describe('Registered agent identity to act as when multiple identities have been registered'), -}; - -type AgentType = 'agent' | 'human'; -type RelayCastLike = Pick; -type AgentClientLike = AgentClient; - -export interface AgentRelayMcpServerOptions { - workspaceKey?: string; - /** @deprecated Use workspaceKey. */ - apiKey?: string; - baseUrl?: string; - agentToken?: string; - agentName?: string; - agentType?: AgentType; - strictAgentName?: boolean; - telemetryTransport?: 'stdio' | 'http'; - skipBootstrap?: boolean; - actions?: AgentRelayActions; - onActionAuditEvent?: (event: ActionAuditEvent) => Promise | void; -} - -interface RegisteredAgent { - agentName: string; - agentToken: string; -} - -interface SessionState { - workspaceKey: string | null; - agentToken: string | null; - agentName: string | null; - agents: Map; - wsBridge: RealtimeResourceBridge | null; - subscriptions: SubscriptionManager | null; - wsInitAttempted: boolean; -} - -type RegistrationSession = Pick & { - agents?: Map; -}; -type SessionSetter = (partial: Partial) => void; type AgentResultCallbackConfig = { url: string; token: string; @@ -331,795 +284,6 @@ function createInitialSession(options: { }; } -class SubscriptionManager { - private readonly subscriptions = new Set(); - - subscribe(uri: string): void { - this.subscriptions.add(uri); - } - - unsubscribe(uri: string): void { - this.subscriptions.delete(uri); - } - - getMatchingSubscriptions(uris: string[]): string[] { - return uris.filter((uri) => this.subscriptions.has(uri)); - } - - getAll(): string[] { - return [...this.subscriptions]; - } - - clear(): void { - this.subscriptions.clear(); - } -} - -function getStringEventField(event: unknown, field: string): string | null { - if (typeof event !== 'object' || event === null) { - return null; - } - const candidate = (event as Record)[field]; - return typeof candidate === 'string' ? candidate : null; -} - -function eventToResourceUris(event: unknown): string[] { - const type = getStringEventField(event, 'type'); - switch (type) { - case 'message.created': { - const channel = getStringEventField(event, 'channel'); - return channel ? ['relay://inbox', `relay://channels/${channel}/messages`] : ['relay://inbox']; - } - case 'message.updated': { - const channel = getStringEventField(event, 'channel'); - return channel ? [`relay://channels/${channel}/messages`] : []; - } - case 'thread.reply': { - const parentId = getStringEventField(event, 'parentId'); - return parentId ? ['relay://inbox', `relay://messages/${parentId}/thread`] : ['relay://inbox']; - } - case 'dm.received': - case 'group_dm.received': { - const conversationId = getStringEventField(event, 'conversationId'); - return conversationId ? ['relay://inbox', `relay://dm/${conversationId}`] : ['relay://inbox']; - } - case 'agent.online': - case 'agent.offline': - return ['relay://agents']; - case 'channel.created': - case 'channel.updated': - case 'channel.archived': - case 'member.joined': - case 'member.left': - return ['relay://channels']; - case 'webhook.received': - case 'command.invoked': { - const channel = getStringEventField(event, 'channel'); - return channel ? [`relay://channels/${channel}/messages`] : []; - } - case 'reaction.added': - case 'reaction.removed': - return ['relay://inbox']; - default: - return []; - } -} - -class RealtimeResourceBridge { - private unsubscribeFn: (() => void) | null = null; - - constructor( - private readonly wsClient: WsClient, - private readonly subscriptions: SubscriptionManager, - private readonly notifyCallback: (uri: string) => void - ) {} - - start(): void { - this.unsubscribeFn = this.wsClient.on('*', (event) => { - const type = getStringEventField(event, 'type'); - if ( - type === 'open' || - type === 'close' || - type === 'error' || - type === 'reconnecting' || - type === 'permanently_disconnected' - ) { - return; - } - const matched = this.subscriptions.getMatchingSubscriptions(eventToResourceUris(event)); - for (const uri of matched) { - this.notifyCallback(uri); - } - }); - this.wsClient.connect(); - } - - stop(): void { - if (this.unsubscribeFn) { - this.unsubscribeFn(); - this.unsubscribeFn = null; - } - this.wsClient.disconnect(); - } -} - -function registerResourceDefinitions( - server: McpServer, - getAgentClient: (asIdentity?: string) => AgentClientLike, - getRelay: () => RelayCast -): void { - server.registerResource( - 'inbox', - 'relay://inbox', - { title: 'Inbox', description: 'Unread messages, mentions, and DMs', mimeType: 'application/json' }, - async (uri) => { - const inbox = await getAgentClient().inbox(); - return { contents: [{ uri: uri.href, text: JSON.stringify(inbox) }] }; - } - ); - - server.registerResource( - 'agents', - 'relay://agents', - { - title: 'Agents', - description: 'Online and offline agents in the workspace', - mimeType: 'application/json', - }, - async (uri) => { - const agents = await getRelay().agents.list(); - return { contents: [{ uri: uri.href, text: JSON.stringify(agents) }] }; - } - ); - - server.registerResource( - 'channels', - 'relay://channels', - { title: 'Channels', description: 'Available channels in the workspace', mimeType: 'application/json' }, - async (uri) => { - const channels = await getAgentClient().channels.list(); - return { contents: [{ uri: uri.href, text: JSON.stringify(channels) }] }; - } - ); - - server.registerResource( - 'channel-messages', - new ResourceTemplate('relay://channels/{name}/messages', { list: undefined }), - { - title: 'Channel Messages', - description: 'Messages in a specific channel', - mimeType: 'application/json', - }, - async (uri, params) => { - const messages = await getAgentClient().messages(String(params.name)); - return { contents: [{ uri: uri.href, text: JSON.stringify(messages) }] }; - } - ); - - server.registerResource( - 'message-thread', - new ResourceTemplate('relay://messages/{id}/thread', { list: undefined }), - { title: 'Message Thread', description: 'Thread replies on a message', mimeType: 'application/json' }, - async (uri, params) => { - const thread = await getAgentClient().thread(String(params.id)); - return { contents: [{ uri: uri.href, text: JSON.stringify(thread) }] }; - } - ); - - server.registerResource( - 'dm-conversation', - new ResourceTemplate('relay://dm/{conversation_id}', { list: undefined }), - { - title: 'DM Conversation', - description: 'Direct message conversation', - mimeType: 'application/json', - }, - async (uri, params) => { - const messages = await getAgentClient().dms.messages(String(params.conversation_id)); - return { contents: [{ uri: uri.href, text: JSON.stringify(messages) }] }; - } - ); -} - -function hasContentArray(value: unknown): value is { content: Array> } { - return ( - typeof value === 'object' && value !== null && Array.isArray((value as { content?: unknown }).content) - ); -} - -const SKIP_PIGGYBACK = new Set(['check_inbox', 'create_workspace', 'set_workspace_key', 'register_agent']); - -function formatInbox(inbox: any, selfName?: string | null): string { - const norm = (s: string) => s.trim().replace(/^@/, '').toLowerCase(); - const selfNorm = selfName ? norm(selfName) : null; - const isSelf = (name: string) => selfNorm != null && norm(name) === selfNorm; - const lines = ['--- Pending Messages ---']; - - if (inbox.unreadChannels?.length) { - lines.push('Unread channels:'); - for (const ch of inbox.unreadChannels) { - lines.push(` #${ch.channelName}: ${ch.unreadCount} unread`); - } - } - - const mentions = selfNorm ? inbox.mentions?.filter((m: any) => !isSelf(m.agentName)) : inbox.mentions; - if (mentions?.length) { - lines.push('Mentions:'); - for (const m of mentions) { - lines.push(` @${m.agentName} in #${m.channelName}: "${m.text}"`); - } - } - - const dms = selfNorm ? inbox.unreadDms?.filter((dm: any) => !isSelf(dm.from)) : inbox.unreadDms; - if (dms?.length) { - lines.push('Unread DMs:'); - for (const dm of dms) { - lines.push(` From ${dm.from}: ${dm.unreadCount} unread`); - } - } - - const reactions = selfNorm - ? inbox.recentReactions?.filter((reaction: any) => !isSelf(reaction.agentName)) - : inbox.recentReactions; - if (reactions?.length) { - lines.push('Reactions (informational; no response required):'); - for (const reaction of reactions) { - lines.push( - ` :${reaction.emoji}: on your message in #${reaction.channelName} by @${reaction.agentName}` - ); - } - } - - return lines.length === 1 ? '' : lines.join('\n'); -} - -function readAsIdentity(args: unknown[]): string | undefined { - const [input] = args; - if (typeof input !== 'object' || input === null) return undefined; - const as = (input as { as?: unknown }).as; - return typeof as === 'string' ? as : undefined; -} - -function invalidAgentTokenToolResult(): JsonToolResult & { isError: true } { - const text = agentTokenRecoveryMessage(); - return { - content: [{ type: 'text', text }], - structuredContent: { - error: { code: INVALID_AGENT_TOKEN_CODE, message: text }, - }, - isError: true, - }; -} - -function isErrorToolResult(value: unknown): boolean { - return Boolean(value && typeof value === 'object' && (value as { isError?: unknown }).isError === true); -} - -interface AgentRelayToolCallMetadata { - toolType: AgentRelayToolCallType; - toolCategory: AgentRelayToolCallCategory; -} - -/** - * Owned tools that delegate to the actions surface (`actions.invoke(...)`) - * rather than the agents/messaging APIs. Together with the dynamic per-action - * tools (tracked via `actionToolNames`), these intentionally skip per-tool - * telemetry so the same underlying action is not counted differently depending - * on which MCP surface the caller used (e.g. `spawn` vs `invoke_action`). - */ -const ACTION_ROUTED_TOOL_NAMES = new Set(['invoke_action', 'spawn']); - -/** - * Coarse type/category metadata for the statically-registered ("owned") MCP - * tools. Action-routed calls (see `ACTION_ROUTED_TOOL_NAMES`) and the dynamic - * per-action tools surfaced from the actions registry are intentionally - * excluded from per-tool telemetry (see the skip in `enableInboxPiggyback`), - * so they have no entry. - */ -const AGENT_RELAY_TOOL_CALL_METADATA = { - add_agent: { toolType: 'agent.create', toolCategory: 'spawn' }, - remove_agent: { toolType: 'agent.release', toolCategory: 'release' }, - list_actions: { toolType: 'action.list', toolCategory: 'action' }, - submit_result: { toolType: 'result.submit', toolCategory: 'result' }, - create_workspace: { toolType: 'workspace.create', toolCategory: 'workspace' }, - set_workspace_key: { toolType: 'workspace.set_key', toolCategory: 'workspace' }, - register_agent: { toolType: 'agent.register', toolCategory: 'agent' }, - list_agents: { toolType: 'agent.list', toolCategory: 'agent' }, - post_message: { toolType: 'message.post', toolCategory: 'message' }, - send_dm: { toolType: 'message.dm', toolCategory: 'message' }, - send_group_dm: { toolType: 'message.group_dm', toolCategory: 'message' }, - list_dms: { toolType: 'message.dm_list', toolCategory: 'message' }, - list_messages: { toolType: 'message.list', toolCategory: 'message' }, - reply_to_thread: { toolType: 'message.reply', toolCategory: 'message' }, - get_message_thread: { toolType: 'message.thread', toolCategory: 'message' }, - search_messages: { toolType: 'message.search', toolCategory: 'message' }, - create_channel: { toolType: 'channel.create', toolCategory: 'channel' }, - list_channels: { toolType: 'channel.list', toolCategory: 'channel' }, - join_channel: { toolType: 'channel.join', toolCategory: 'channel' }, - leave_channel: { toolType: 'channel.leave', toolCategory: 'channel' }, - set_channel_topic: { toolType: 'channel.set_topic', toolCategory: 'channel' }, - archive_channel: { toolType: 'channel.archive', toolCategory: 'channel' }, - invite_to_channel: { toolType: 'channel.invite', toolCategory: 'channel' }, - add_reaction: { toolType: 'reaction.add', toolCategory: 'reaction' }, - remove_reaction: { toolType: 'reaction.remove', toolCategory: 'reaction' }, - check_inbox: { toolType: 'inbox.check', toolCategory: 'inbox' }, - mark_message_read: { toolType: 'inbox.mark_read', toolCategory: 'inbox' }, - get_message_readers: { toolType: 'inbox.reader_list', toolCategory: 'inbox' }, -} satisfies Record; - -function agentRelayToolCallMetadata(name: string): AgentRelayToolCallMetadata { - const known = (AGENT_RELAY_TOOL_CALL_METADATA as Partial>)[name]; - return known ?? { toolType: name, toolCategory: 'tool' }; -} - -function trackAgentRelayToolCall(input: { - toolName: string; - toolType: AgentRelayToolCallType; - toolCategory: AgentRelayToolCallCategory; - transport?: AgentRelayMcpServerOptions['telemetryTransport']; - startedAt: number; - success: boolean; - errorClass?: string; -}): void { - track('agent_relay_tool_call', { - tool_name: input.toolName, - tool_type: input.toolType, - tool_category: input.toolCategory, - transport: input.transport ?? 'unknown', - success: input.success, - duration_ms: Date.now() - input.startedAt, - ...(input.errorClass ? { error_class: input.errorClass } : {}), - }); -} - -function enableInboxPiggyback( - mcpServer: McpServer, - getSession: () => SessionState, - getAgentClient: (asIdentity?: string) => AgentClientLike, - invalidateAgentToken: (asIdentity?: string) => void, - telemetryTransport?: AgentRelayMcpServerOptions['telemetryTransport'], - actionToolNames = new Set() -): void { - const original = mcpServer.registerTool.bind(mcpServer); - const mutableServer = mcpServer as McpServer & { - registerTool: McpServer['registerTool']; - }; - - mutableServer.registerTool = (name: string, config: any, handler: any) => { - if (!handler) { - return original(name, config, handler); - } - - const wrapped = async (...args: unknown[]) => { - const asIdentity = readAsIdentity(args); - const startedAt = Date.now(); - // Action-routed calls (`invoke_action`, `spawn`, and the dynamic - // per-action tools) run through the actions surface and deliberately skip - // per-tool telemetry; only the owned tools emit `agent_relay_tool_call`. - const toolMetadata = - !ACTION_ROUTED_TOOL_NAMES.has(name) && !actionToolNames.has(name) - ? agentRelayToolCallMetadata(name) - : undefined; - - let result: any; - try { - result = await handler(...args); - } catch (err) { - if (name !== 'register_agent' && isInvalidAgentTokenError(err)) { - invalidateAgentToken(asIdentity); - if (toolMetadata) { - trackAgentRelayToolCall({ - toolName: name, - toolType: toolMetadata.toolType, - toolCategory: toolMetadata.toolCategory, - transport: telemetryTransport, - startedAt, - success: false, - errorClass: errorClassName(err) ?? 'InvalidAgentToken', - }); - } - return invalidAgentTokenToolResult(); - } - if (toolMetadata) { - trackAgentRelayToolCall({ - toolName: name, - toolType: toolMetadata.toolType, - toolCategory: toolMetadata.toolCategory, - transport: telemetryTransport, - startedAt, - success: false, - errorClass: errorClassName(err), - }); - } - throw err; - } - - if (name !== 'register_agent' && isInvalidAgentTokenToolResult(result)) { - invalidateAgentToken(asIdentity); - if (toolMetadata) { - trackAgentRelayToolCall({ - toolName: name, - toolType: toolMetadata.toolType, - toolCategory: toolMetadata.toolCategory, - transport: telemetryTransport, - startedAt, - success: false, - errorClass: 'InvalidAgentToken', - }); - } - if (hasContentArray(result)) { - result.content.push({ type: 'text', text: agentTokenRecoveryMessage() }); - } - return result; - } - - if (!SKIP_PIGGYBACK.has(name) && getSession().agentToken && hasContentArray(result)) { - try { - const inbox = await getAgentClient(asIdentity).inbox(); - const inboxText = formatInbox(inbox, asIdentity ?? getSession().agentName); - if (inboxText) { - result.content.push({ type: 'text', text: inboxText }); - } - } catch (err) { - if (isInvalidAgentTokenError(err)) { - invalidateAgentToken(asIdentity); - } - } - } - - if (toolMetadata) { - const resultIsError = isErrorToolResult(result); - trackAgentRelayToolCall({ - toolName: name, - toolType: toolMetadata.toolType, - toolCategory: toolMetadata.toolCategory, - transport: telemetryTransport, - startedAt, - success: !resultIsError, - ...(resultIsError ? { errorClass: 'ToolResultError' } : {}), - }); - } - - return result; - }; - - return original(name, config, wrapped); - }; -} - -async function createWorkspace(name: string, baseUrl?: string): Promise> { - return (await RelayCast.createWorkspace(name, { - baseUrl, - ...relaycastWorkspaceTelemetryOptions(), - })) as Record; -} - -function extractWorkspaceKey(payload: Record): string | undefined { - const data = - payload.data && typeof payload.data === 'object' ? (payload.data as Record) : {}; - const value = - payload.workspaceKey ?? - payload.workspace_key ?? - payload.apiKey ?? - payload.api_key ?? - data.workspaceKey ?? - data.workspace_key ?? - data.apiKey ?? - data.api_key; - - return typeof value === 'string' && value.trim() ? value : undefined; -} - -function extractWorkspaceName(payload: Record, fallback: string): string { - const data = - payload.data && typeof payload.data === 'object' ? (payload.data as Record) : {}; - const value = payload.workspaceName ?? payload.workspace_name ?? payload.name ?? data.workspaceName; - return typeof value === 'string' && value.trim() ? value : fallback; -} - -function requireWorkspaceKey(session: RegistrationSession): void { - if (session.workspaceKey) { - return; - } - - throw new Error( - 'Workspace key not configured. Call "create_workspace" first, or "set_workspace_key" if someone shared a workspace key.' - ); -} - -type JsonToolResult = { - content: Array<{ type: 'text'; text: string }>; - structuredContent: Record; -}; - -function jsonContent(value: unknown): JsonToolResult { - const structuredContent = - typeof value === 'object' && value !== null && !Array.isArray(value) - ? (value as Record) - : { value }; - return { - content: [{ type: 'text', text: JSON.stringify(value, null, 2) }], - structuredContent, - }; -} - -function textContent(message: string, structuredContent: Record = { message }) { - return { - content: [{ type: 'text' as const, text: message }], - structuredContent, - }; -} - -function isSchemaObject(schema: ActionSchema | undefined): schema is JsonSchemaLiteObject { - return Boolean( - schema && - typeof schema === 'object' && - !Array.isArray(schema) && - typeof (schema as { safeParse?: unknown }).safeParse !== 'function' - ); -} - -function getSchemaDescription(schema: ActionSchema | undefined): string | undefined { - return isSchemaObject(schema) && typeof schema.description === 'string' ? schema.description : undefined; -} - -function zodFromJsonSchema(schema: ActionSchema | undefined): z.ZodTypeAny { - if (schema === false) { - return z.never(); - } - - if (!isSchemaObject(schema)) { - return z.unknown(); - } - - let zodType: z.ZodTypeAny; - const schemaType = Array.isArray(schema.type) ? schema.type[0] : schema.type; - switch (schemaType) { - case 'array': - zodType = z.array(zodFromJsonSchema(schema.items)); - break; - case 'boolean': - zodType = z.boolean(); - break; - case 'integer': - zodType = z.number().int(); - break; - case 'number': - zodType = z.number(); - break; - case 'object': - if (schema.properties) { - const required = new Set(schema.required ?? []); - const shape: Record = {}; - for (const [key, childSchema] of Object.entries(schema.properties)) { - const child = zodFromJsonSchema(childSchema); - shape[key] = required.has(key) ? child : child.optional(); - } - zodType = z.object(shape).passthrough(); - } else { - zodType = z.record(z.string(), z.unknown()); - } - break; - case 'string': - zodType = z.string(); - break; - default: - zodType = z.unknown(); - break; - } - - const description = getSchemaDescription(schema); - return description ? zodType.describe(description) : zodType; -} - -function actionToolInputSchema(schema: ActionSchema | undefined): Record { - const zodShape = zodObjectShape(schema); - if (zodShape) { - return zodShape; - } - - if (!isSchemaObject(schema) || schema.type !== 'object') { - return { - input: z.unknown().describe('Action input payload. The action registry performs final validation.'), - }; - } - - const required = new Set(schema.required ?? []); - const shape: Record = {}; - for (const [key, childSchema] of Object.entries(schema.properties ?? {})) { - const child = zodFromJsonSchema(childSchema); - shape[key] = required.has(key) ? child : child.optional(); - } - return shape; -} - -function actionInvocationInput(descriptor: AgentRelayActionDescriptor, args: unknown): unknown { - const schema = descriptor.inputSchema; - if (zodObjectShape(schema)) { - return args; - } - if (!isSchemaObject(schema) || schema.type !== 'object') { - return typeof args === 'object' && args !== null && 'input' in args - ? (args as { input?: unknown }).input - : args; - } - return args; -} - -function zodObjectShape(schema: ActionSchema | undefined): Record | undefined { - if (schema instanceof z.ZodObject) { - return schema.shape; - } - return undefined; -} - -function serializableActionDescriptor(descriptor: AgentRelayActionDescriptor): Record { - return { - name: descriptor.name, - description: descriptor.description, - visibility: descriptor.visibility, - ...(descriptor.inputSchema ? { inputSchema: serializableActionSchema(descriptor.inputSchema) } : {}), - ...(descriptor.outputSchema ? { outputSchema: serializableActionSchema(descriptor.outputSchema) } : {}), - }; -} - -function serializableActionSchema(schema: ActionSchema): unknown { - if (isSchemaObject(schema)) { - return schema; - } - if (isZodLikeSchema(schema)) { - return { - type: 'zod', - ...(schema.description ? { description: schema.description } : {}), - }; - } - return schema; -} - -function isZodLikeSchema(schema: ActionSchema | undefined): schema is ZodLikeSchema { - return Boolean( - schema && - typeof schema === 'object' && - !Array.isArray(schema) && - typeof (schema as { safeParse?: unknown }).safeParse === 'function' - ); -} - -function registerAgentRelayActionTools( - server: McpServer, - actions: AgentRelayActions | undefined, - getSession: () => SessionState, - onAuditEvent?: (event: ActionAuditEvent) => Promise | void, - getAgentClient?: (asIdentity?: string) => AgentClientLike, - actionToolNames?: Set -): void { - if (!actions) { - return; - } - - /** - * Fire-and-forget invocation through the relay: returns an immediate ack - * (with an `invocation_id`) and does NOT run the handler inline. Falls back to - * the local in-process registry when the relay action surface is unavailable. - */ - const invokeAction = async (name: string, input: unknown) => { - const relayActions = getRelayAgentActions(getAgentClient); - if (relayActions) { - try { - const ack = await relayActions.invoke(name, asInputRecord(input)); - return jsonContent({ ok: true, status: 'invoked', invocation: ack }); - } catch (error) { - return { ...jsonContent({ ok: false, error: errorMessage(error) }), isError: true }; - } - } - - const session = getSession(); - const result = await actions.invoke({ - name, - input, - context: { - caller: { name: session.agentName ?? 'mcp', type: 'agent' }, - emit: onAuditEvent, - }, - }); - return result.ok ? jsonContent(result) : { ...jsonContent(result), isError: true }; - }; - - server.registerTool( - 'list_actions', - { - title: 'List Actions', - description: 'List Agent Relay actions available to this agent.', - inputSchema: {}, - outputSchema: jsonResult, - annotations: { - readOnlyHint: true, - destructiveHint: false, - idempotentHint: true, - openWorldHint: false, - }, - }, - async () => - jsonContent({ - actions: (await actions.list({ visibility: 'agent' })).map(serializableActionDescriptor), - }) - ); - - server.registerTool( - 'invoke_action', - { - title: 'Invoke Action', - description: - 'Invoke a registered Agent Relay action by name. Fire-and-forget: returns an ack with an invocation id; the result arrives asynchronously to the action handler.', - inputSchema: { - name: z.string().describe('Registered action name'), - input: z.unknown().describe('Action input payload'), - }, - outputSchema: jsonResult, - annotations: { - readOnlyHint: false, - destructiveHint: false, - idempotentHint: false, - openWorldHint: false, - }, - }, - async ({ name, input }: { name: string; input: unknown }) => invokeAction(name, input) - ); - - void actions - .list({ visibility: 'agent' }) - .then((descriptors) => { - for (const descriptor of descriptors) { - actionToolNames?.add(descriptor.name); - server.registerTool( - descriptor.name, - { - title: descriptor.name, - description: descriptor.description, - inputSchema: actionToolInputSchema(descriptor.inputSchema), - outputSchema: jsonResult, - annotations: { - readOnlyHint: false, - destructiveHint: false, - idempotentHint: false, - openWorldHint: false, - }, - }, - async (args: unknown) => invokeAction(descriptor.name, actionInvocationInput(descriptor, args)) - ); - } - }) - .catch(() => undefined); -} - -/** The relay-backed action surface on the live agent client, when available. */ -function getRelayAgentActions( - getAgentClient?: (asIdentity?: string) => AgentClientLike -): AgentClientLike['actions'] | undefined { - if (!getAgentClient) { - return undefined; - } - try { - return getAgentClient().actions; - } catch { - return undefined; - } -} - -function asInputRecord(input: unknown): Record | undefined { - if (input === undefined || input === null) { - return undefined; - } - if (typeof input === 'object' && !Array.isArray(input)) { - return input as Record; - } - return { input }; -} - -function errorMessage(error: unknown): string { - return error instanceof Error ? error.message : String(error); -} - function createRegisteredAgent(agentName: string, agentToken: string): RegisteredAgent { return { agentName, agentToken }; } @@ -1188,22 +352,6 @@ export async function registerAgentWithRebind({ }; } -function resolveEmoji(input: string): string { - const normalized = input.trim().replace(/^:/, '').replace(/:$/, '').toLowerCase(); - const aliases: Record = { - '+1': '👍', - thumbsup: '👍', - thumbs_up: '👍', - check: '✅', - white_check_mark: '✅', - rocket: '🚀', - eyes: '👀', - heart: '❤️', - clap: '👏', - }; - return aliases[normalized] ?? input; -} - function registerAgentRelayTools( server: McpServer, getRelay: () => RelayCast, @@ -1388,401 +536,7 @@ function registerAgentRelayTools( } ); - server.registerTool( - 'create_channel', - { - title: 'Create Channel', - description: 'Create a new workspace channel.', - inputSchema: { - name: z.string().describe('Unique channel name'), - topic: z.string().optional().describe('Optional channel topic'), - ...identityOverrideInputShape, - }, - outputSchema: jsonResult, - annotations: { - readOnlyHint: false, - destructiveHint: false, - idempotentHint: false, - openWorldHint: true, - }, - }, - async ({ name, topic, as }) => jsonContent(await getAgentClient(as).channels.create({ name, topic })) - ); - - server.registerTool( - 'list_channels', - { - title: 'List Channels', - description: 'List channels available in the workspace.', - inputSchema: { - include_archived: z.boolean().optional().describe('Include archived channels'), - ...identityOverrideInputShape, - }, - outputSchema: { - channels: z.array(z.object({}).passthrough()).describe('Channels'), - }, - annotations: { readOnlyHint: true, destructiveHint: false, idempotentHint: true, openWorldHint: true }, - }, - async ({ include_archived, as }) => { - const channels = await getAgentClient(as).channels.list( - include_archived ? { includeArchived: include_archived } : undefined - ); - return jsonContent({ channels }); - } - ); - - server.registerTool( - 'join_channel', - { - title: 'Join Channel', - description: 'Join an existing channel.', - inputSchema: { - channel: z.string().describe('Channel name'), - ...identityOverrideInputShape, - }, - outputSchema: messageResult, - annotations: { readOnlyHint: false, destructiveHint: false, idempotentHint: true, openWorldHint: true }, - }, - async ({ channel, as }) => { - await getAgentClient(as).channels.join(channel); - return textContent(`Joined channel #${channel}`); - } - ); - - server.registerTool( - 'leave_channel', - { - title: 'Leave Channel', - description: 'Leave a channel.', - inputSchema: { - channel: z.string().describe('Channel name'), - ...identityOverrideInputShape, - }, - outputSchema: messageResult, - annotations: { readOnlyHint: false, destructiveHint: false, idempotentHint: true, openWorldHint: true }, - }, - async ({ channel, as }) => { - await getAgentClient(as).channels.leave(channel); - return textContent(`Left channel #${channel}`); - } - ); - - server.registerTool( - 'invite_to_channel', - { - title: 'Invite to Channel', - description: 'Invite another agent to a channel.', - inputSchema: { - channel: z.string().describe('Channel name'), - agent: z.string().describe('Agent name to invite'), - ...identityOverrideInputShape, - }, - outputSchema: messageResult, - annotations: { readOnlyHint: false, destructiveHint: false, idempotentHint: true, openWorldHint: true }, - }, - async ({ channel, agent, as }) => { - await getAgentClient(as).channels.invite(channel, agent); - return textContent(`Invited ${agent} to #${channel}`); - } - ); - - server.registerTool( - 'set_channel_topic', - { - title: 'Set Channel Topic', - description: 'Update a channel topic.', - inputSchema: { - channel: z.string().describe('Channel name'), - topic: z.string().describe('New topic'), - ...identityOverrideInputShape, - }, - outputSchema: jsonResult, - annotations: { readOnlyHint: false, destructiveHint: false, idempotentHint: true, openWorldHint: true }, - }, - async ({ channel, topic, as }) => jsonContent(await getAgentClient(as).channels.setTopic(channel, topic)) - ); - - server.registerTool( - 'archive_channel', - { - title: 'Archive Channel', - description: 'Archive a channel.', - inputSchema: { - channel: z.string().describe('Channel name'), - ...identityOverrideInputShape, - }, - outputSchema: messageResult, - annotations: { readOnlyHint: false, destructiveHint: true, idempotentHint: true, openWorldHint: true }, - }, - async ({ channel, as }) => { - await getAgentClient(as).channels.archive(channel); - return textContent(`Archived channel #${channel}`); - } - ); - - server.registerTool( - 'post_message', - { - title: 'Post Message', - description: 'Post a new message to a channel as the current agent.', - inputSchema: { - channel: z.string().describe('Channel name'), - text: z.string().describe('Message text'), - attachments: z.array(z.string()).optional().describe('File attachment IDs'), - mode: z.enum(['wait', 'steer']).optional().describe('Delivery mode'), - ...identityOverrideInputShape, - }, - outputSchema: jsonResult, - annotations: { - readOnlyHint: false, - destructiveHint: false, - idempotentHint: false, - openWorldHint: true, - }, - }, - async ({ channel, text, attachments, mode, as }) => - jsonContent(await getAgentClient(as).send(channel, text, { attachments, mode })) - ); - - server.registerTool( - 'list_messages', - { - title: 'Get Messages', - description: 'Retrieve message history from a channel.', - inputSchema: { - channel: z.string().describe('Channel name'), - limit: z.number().optional().describe('Maximum messages to return'), - before: z.string().optional().describe('Older-than cursor'), - after: z.string().optional().describe('Newer-than cursor'), - ...identityOverrideInputShape, - }, - outputSchema: { - messages: z.array(z.object({}).passthrough()).describe('Messages'), - }, - annotations: { readOnlyHint: true, destructiveHint: false, idempotentHint: true, openWorldHint: true }, - }, - async ({ channel, limit, before, after, as }) => { - const messages = await getAgentClient(as).messages(channel, { limit, before, after }); - return jsonContent({ messages }); - } - ); - - server.registerTool( - 'reply_to_thread', - { - title: 'Reply to Thread', - description: 'Reply to an existing message thread.', - inputSchema: { - message_id: z.string().describe('Parent message ID'), - text: z.string().describe('Reply text'), - ...identityOverrideInputShape, - }, - outputSchema: jsonResult, - annotations: { - readOnlyHint: false, - destructiveHint: false, - idempotentHint: false, - openWorldHint: true, - }, - }, - async ({ message_id, text, as }) => jsonContent(await getAgentClient(as).reply(message_id, text)) - ); - - server.registerTool( - 'get_message_thread', - { - title: 'Get Thread', - description: 'Retrieve a message thread.', - inputSchema: { - message_id: z.string().describe('Parent message ID'), - limit: z.number().optional().describe('Maximum replies to return'), - ...identityOverrideInputShape, - }, - outputSchema: jsonResult, - annotations: { readOnlyHint: true, destructiveHint: false, idempotentHint: true, openWorldHint: true }, - }, - async ({ message_id, limit, as }) => - jsonContent(await getAgentClient(as).thread(message_id, limit ? { limit } : undefined)) - ); - - server.registerTool( - 'send_dm', - { - title: 'Send Direct Message', - description: 'Send a private direct message to another agent.', - inputSchema: { - to: z.string().describe('Recipient agent name'), - text: z.string().describe('DM text'), - mode: z.enum(['wait', 'steer']).optional().describe('Delivery mode'), - attachments: z.array(z.string()).optional().describe('File attachment IDs'), - ...identityOverrideInputShape, - }, - outputSchema: jsonResult, - annotations: { - readOnlyHint: false, - destructiveHint: false, - idempotentHint: false, - openWorldHint: true, - }, - }, - async ({ to, text, mode, attachments, as }) => - jsonContent(await getAgentClient(as).dm(to, text, { mode, attachments })) - ); - - server.registerTool( - 'list_dms', - { - title: 'List DM Conversations', - description: 'List direct message conversations for the current agent.', - inputSchema: { - ...identityOverrideInputShape, - }, - outputSchema: { - conversations: z.array(z.object({}).passthrough()).describe('DM conversations'), - }, - annotations: { readOnlyHint: true, destructiveHint: false, idempotentHint: true, openWorldHint: true }, - }, - async ({ as }) => jsonContent({ conversations: await getAgentClient(as).dms.conversations() }) - ); - - server.registerTool( - 'send_group_dm', - { - title: 'Send Group DM', - description: 'Create a group DM and send the first message.', - inputSchema: { - participants: z.array(z.string()).describe('Participant agent names'), - name: z.string().optional().describe('Optional group name'), - text: z.string().describe('Initial message'), - ...identityOverrideInputShape, - }, - outputSchema: jsonResult, - annotations: { - readOnlyHint: false, - destructiveHint: false, - idempotentHint: false, - openWorldHint: true, - }, - }, - async ({ participants, name, text, as }) => { - const client = getAgentClient(as); - const conversation = await client.dms.createGroup({ participants, name }); - const message = await client.dms.sendMessage(conversation.id, text); - return jsonContent({ conversation, message }); - } - ); - - server.registerTool( - 'add_reaction', - { - title: 'Add Reaction', - description: 'Add an emoji reaction to a message.', - inputSchema: { - message_id: z.string().describe('Message ID'), - emoji: z.string().describe('Emoji character or shortcode'), - ...identityOverrideInputShape, - }, - outputSchema: messageResult, - annotations: { readOnlyHint: false, destructiveHint: false, idempotentHint: true, openWorldHint: true }, - }, - async ({ message_id, emoji, as }) => { - const resolved = resolveEmoji(emoji); - await getAgentClient(as).react(message_id, resolved); - return textContent(`Reacted with ${resolved}`); - } - ); - - server.registerTool( - 'remove_reaction', - { - title: 'Remove Reaction', - description: 'Remove an emoji reaction from a message.', - inputSchema: { - message_id: z.string().describe('Message ID'), - emoji: z.string().describe('Emoji character or shortcode'), - ...identityOverrideInputShape, - }, - outputSchema: messageResult, - annotations: { readOnlyHint: false, destructiveHint: false, idempotentHint: true, openWorldHint: true }, - }, - async ({ message_id, emoji, as }) => { - const resolved = resolveEmoji(emoji); - await getAgentClient(as).unreact(message_id, resolved); - return textContent(`Removed reaction ${resolved}`); - } - ); - - server.registerTool( - 'search_messages', - { - title: 'Search Messages', - description: 'Search messages across the workspace.', - inputSchema: { - query: z.string().describe('Text search query'), - channel: z.string().optional().describe('Optional channel filter'), - from: z.string().optional().describe('Optional sender filter'), - limit: z.number().optional().describe('Maximum results'), - ...identityOverrideInputShape, - }, - outputSchema: { - results: z.array(z.object({}).passthrough()).describe('Search results'), - }, - annotations: { readOnlyHint: true, destructiveHint: false, idempotentHint: true, openWorldHint: true }, - }, - async ({ query, channel, from, limit, as }) => - jsonContent({ results: await getAgentClient(as).search(query, { channel, from, limit }) }) - ); - - server.registerTool( - 'check_inbox', - { - title: 'Check Inbox', - description: 'Check unread messages, mentions, DMs, and reactions for the current agent.', - inputSchema: { - limit: z.number().optional().describe('Maximum inbox items'), - ...identityOverrideInputShape, - }, - outputSchema: jsonResult, - annotations: { readOnlyHint: true, destructiveHint: false, idempotentHint: true, openWorldHint: true }, - }, - async ({ limit, as }) => - jsonContent(await getAgentClient(as).inbox(limit != null ? { limit } : undefined)) - ); - - server.registerTool( - 'mark_message_read', - { - title: 'Mark as Read', - description: 'Mark a message as read for the current agent.', - inputSchema: { - message_id: z.string().describe('Message ID'), - ...identityOverrideInputShape, - }, - outputSchema: messageResult, - annotations: { readOnlyHint: false, destructiveHint: false, idempotentHint: true, openWorldHint: true }, - }, - async ({ message_id, as }) => { - await getAgentClient(as).markRead(message_id); - return textContent(`Marked message ${message_id} as read`); - } - ); - - server.registerTool( - 'get_message_readers', - { - title: 'Get Readers', - description: 'List agents who have read a message.', - inputSchema: { - message_id: z.string().describe('Message ID'), - ...identityOverrideInputShape, - }, - outputSchema: { - readers: z.array(z.object({}).passthrough()).describe('Readers'), - }, - annotations: { readOnlyHint: true, destructiveHint: false, idempotentHint: true, openWorldHint: true }, - }, - async ({ message_id, as }) => jsonContent({ readers: await getAgentClient(as).readers(message_id) }) - ); + registerMessagingTools(server, getAgentClient); server.registerTool( 'add_agent', diff --git a/packages/cli/src/cli/lib/broker-dashboard.ts b/packages/cli/src/cli/lib/broker-dashboard.ts new file mode 100644 index 000000000..98f8a8dc0 --- /dev/null +++ b/packages/cli/src/cli/lib/broker-dashboard.ts @@ -0,0 +1,570 @@ +import fs from 'node:fs'; +import os from 'node:os'; +import path from 'node:path'; + +import type { CoreDependencies, CoreProjectPaths, SpawnedProcess } from '../commands/core.js'; + +export async function resolveDashboardPortWithFallback( + dashboardPort: number, + dashboardPortCandidates: number, + deps: CoreDependencies +): Promise { + for (let attempt = 0; attempt < dashboardPortCandidates; attempt += 1) { + const candidatePort = dashboardPort + attempt; + const inUse = await deps.isPortInUse(candidatePort); + if (!inUse) { + if (attempt > 0) { + deps.warn(`Dashboard port ${dashboardPort} is already in use; trying ${candidatePort}`); + } + return candidatePort; + } + } + + throw new Error(`Failed to find an available dashboard port near ${dashboardPort}.`); +} + +function pickDashboardStaticDir(candidates: string[], deps: CoreDependencies): string | null { + const existingCandidates = Array.from(new Set(candidates)).filter((candidate) => + deps.fs.existsSync(candidate) + ); + if (existingCandidates.length === 0) { + return null; + } + + const pageMarkerPriority = [ + ['metrics.html', path.join('metrics', 'index.html')], + ['app.html'], + ['index.html'], + ]; + + for (const markerGroup of pageMarkerPriority) { + const withMarker = existingCandidates.find((candidate) => + markerGroup.some((marker) => deps.fs.existsSync(path.join(candidate, marker))) + ); + if (withMarker) { + return withMarker; + } + } + + return existingCandidates[0]; +} + +function getHomeDashboardRoot(deps: CoreDependencies): string { + const homeDir = deps.env.HOME || deps.env.USERPROFILE || os.homedir(); + return path.join(homeDir, '.agentworkforce/relay', 'dashboard'); +} + +function getPriorDashboardRoot(deps: CoreDependencies): string | null { + const homeDir = deps.env.HOME || deps.env.USERPROFILE || ''; + if (!homeDir) { + return null; + } + return path.join(homeDir, '.relay', 'dashboard'); +} + +function getDashboardRootFromBinary(dashboardBinary: string | null, deps: CoreDependencies): string | null { + if (!dashboardBinary || dashboardBinary.endsWith('.js') || dashboardBinary.endsWith('.ts')) { + return null; + } + + const binaryDir = path.dirname(dashboardBinary); + if (path.basename(binaryDir) !== 'bin') { + return null; + } + + const homeDir = deps.env.HOME || deps.env.USERPROFILE || ''; + const resolvedBinaryDir = path.resolve(binaryDir); + const ignoredBinDirs = [ + homeDir ? path.join(homeDir, '.local', 'bin') : null, + path.join('/usr/local', 'bin'), + ] + .filter((candidate): candidate is string => Boolean(candidate)) + .map((candidate) => path.resolve(candidate)); + if (ignoredBinDirs.includes(resolvedBinaryDir)) { + return null; + } + + return path.join(path.dirname(binaryDir), 'dashboard'); +} + +function resolveDashboardStaticDir(dashboardBinary: string | null, deps: CoreDependencies): string | null { + const explicitStaticDir = deps.env.RELAY_DASHBOARD_STATIC_DIR ?? deps.env.STATIC_DIR; + if (explicitStaticDir && explicitStaticDir.trim()) { + return explicitStaticDir; + } + + if (!dashboardBinary) { + return null; + } + + if (dashboardBinary.endsWith('.js') || dashboardBinary.endsWith('.ts')) { + const dashboardServerOutDir = path.resolve(path.dirname(dashboardBinary), '..', 'out'); + const siblingDashboardOutDir = path.resolve( + path.dirname(dashboardBinary), + '..', + '..', + 'dashboard', + 'out' + ); + return pickDashboardStaticDir([dashboardServerOutDir, siblingDashboardOutDir], deps); + } + + // Installs place UI assets under the install dir (~/.agentworkforce/relay/dashboard/out + // by default, or next to a custom install's bin/ directory). ~/.relay/dashboard/out is + // read as a fallback for installs predating that move. + const installDashboardRoot = getDashboardRootFromBinary(dashboardBinary, deps); + const priorDashboardRoot = getPriorDashboardRoot(deps); + const candidates = [ + installDashboardRoot ? path.join(installDashboardRoot, 'out') : null, + path.join(getHomeDashboardRoot(deps), 'out'), + priorDashboardRoot ? path.join(priorDashboardRoot, 'out') : null, + ].filter((candidate): candidate is string => Boolean(candidate)); + return pickDashboardStaticDir(candidates, deps); +} + +function normalizeLocalhostRelayUrl(relayUrl: string): string { + try { + const parsed = new URL(relayUrl); + if (parsed.hostname === 'localhost') { + parsed.hostname = '127.0.0.1'; + } + return parsed.toString().replace(/\/+$/, ''); + } catch { + return relayUrl; + } +} + +export function getDefaultDashboardRelayUrl(apiPort: number): string { + return normalizeLocalhostRelayUrl(`http://localhost:${apiPort}`); +} + +export function resolveDashboardRelayUrl(apiPort: number, deps: CoreDependencies): string { + const explicitRelayUrl = deps.env.RELAY_DASHBOARD_RELAY_URL; + if (explicitRelayUrl && explicitRelayUrl.trim()) { + return normalizeLocalhostRelayUrl(explicitRelayUrl.trim()); + } + + return getDefaultDashboardRelayUrl(apiPort); +} + +export function isDebugLikeLoggingEnabled(deps: CoreDependencies): boolean { + const rawLevel = String(deps.env.RUST_LOG ?? '').toLowerCase(); + return rawLevel.includes('debug') || rawLevel.includes('trace'); +} + +function getDashboardSpawnEnv( + deps: CoreDependencies, + relayUrl: string, + enableVerboseLogging: boolean, + relayApiKey?: string, + brokerApiKey?: string +): NodeJS.ProcessEnv { + const env: NodeJS.ProcessEnv = { + ...deps.env, + RELAY_URL: relayUrl, + VERBOSE: enableVerboseLogging || deps.env.VERBOSE === 'true' ? 'true' : deps.env.VERBOSE, + }; + // Pass the workspace key so the dashboard can make Agent Relay calls + // (e.g. posting thread replies) without requiring a relaycast.json file. + if (relayApiKey) { + if (!env.RELAY_WORKSPACE_KEY) { + env.RELAY_WORKSPACE_KEY = relayApiKey; + } + if (!env.RELAY_API_KEY) { + env.RELAY_API_KEY = relayApiKey; + } + } + // Pass the broker API key so the dashboard can authenticate with the + // broker's HTTP API (e.g. /api/spawn, /api/spawned). + if (brokerApiKey) { + env.RELAY_BROKER_API_KEY = brokerApiKey; + } + return env; +} + +function getDashboardSpawnArgs( + paths: CoreProjectPaths, + port: number, + apiPort: number, + dashboardBinary: string | null, + relayUrl: string, + enableVerboseLogging: boolean, + deps: CoreDependencies +): string[] { + const args = ['--port', String(port), '--data-dir', paths.dataDir]; + args.push('--relay-url', relayUrl); + const staticDir = resolveDashboardStaticDir(dashboardBinary, deps); + if (staticDir) { + args.push('--static-dir', staticDir); + } + if (enableVerboseLogging) { + args.push('--verbose'); + } + return args; +} + +export function normalizeDashboardPath(rawDashboardPath: string | undefined): string | undefined { + const trimmed = rawDashboardPath?.trim(); + if (!trimmed) return undefined; + if (trimmed.startsWith('/')) { + return trimmed; + } + return `/${trimmed}`; +} + +interface DashboardStartupProcess extends SpawnedProcess { + stdout?: { + on?: (event: string, cb: (chunk: Buffer) => void) => void; + removeListener?: (event: string, cb: (...args: unknown[]) => void) => void; + off?: (event: string, cb: (...args: unknown[]) => void) => void; + }; + stderr?: { + on?: (event: string, cb: (chunk: Buffer) => void) => void; + removeListener?: (event: string, cb: (...args: unknown[]) => void) => void; + off?: (event: string, cb: (...args: unknown[]) => void) => void; + }; +} + +function startDashboard( + paths: CoreProjectPaths, + port: number, + apiPort: number, + deps: CoreDependencies, + enableVerboseLogging: boolean, + dashboardBinaryOverride?: string | null, + relayApiKey?: string, + brokerApiKey?: string +): DashboardStartupProcess { + const dashboardBinary = + dashboardBinaryOverride === undefined ? deps.findDashboardBinary() : dashboardBinaryOverride; + const relayUrl = resolveDashboardRelayUrl(apiPort, deps); + const shouldEnableVerbose = enableVerboseLogging || isDebugLikeLoggingEnabled(deps); + const args = getDashboardSpawnArgs( + paths, + port, + apiPort, + dashboardBinary, + relayUrl, + shouldEnableVerbose, + deps + ); + const launchTarget = dashboardBinary + ? dashboardBinary.endsWith('.js') + ? `node ${dashboardBinary}` + : dashboardBinary + : 'npx --yes @agent-relay/dashboard-server@latest'; + + const spawnOpts = { + stdio: ['ignore', 'pipe', 'pipe'] as unknown, + env: getDashboardSpawnEnv(deps, relayUrl, shouldEnableVerbose, relayApiKey, brokerApiKey), + }; + if (shouldEnableVerbose) { + deps.log(`[dashboard] Starting: ${launchTarget} ${args.join(' ')}`); + } + + let child: SpawnedProcess; + if (dashboardBinary) { + // If the binary is a .js file (local dev), run it with node + if (dashboardBinary.endsWith('.js')) { + child = deps.spawnProcess('node', [dashboardBinary, ...args], spawnOpts); + } else { + child = deps.spawnProcess(dashboardBinary, args, spawnOpts); + } + } else { + child = deps.spawnProcess('npx', ['--yes', '@agent-relay/dashboard-server@latest', ...args], spawnOpts); + } + + // Capture stderr for error reporting + const childAny = child as unknown as { + stdout?: { on?: (event: string, cb: (chunk: Buffer) => void) => void }; + stderr?: { on?: (event: string, cb: (chunk: Buffer) => void) => void }; + on?: (event: string, cb: (...args: unknown[]) => void) => void; + }; + let stderrBuf = ''; + + const logChunk = (chunk: Buffer, logger: (line: string) => void, prefix: string) => { + if (!shouldEnableVerbose) { + return; + } + const text = chunk.toString(); + for (const line of text.split(/\r?\n/)) { + const trimmed = line.trim(); + if (trimmed) { + logger(`[dashboard] ${prefix}: ${trimmed}`); + } + } + }; + + childAny.stdout?.on?.('data', (chunk: Buffer) => { + logChunk(chunk, deps.log, 'stdout'); + }); + childAny.stderr?.on?.('data', (chunk: Buffer) => { + stderrBuf += chunk.toString(); + logChunk(chunk, deps.warn, 'stderr'); + }); + + // Report early crashes + childAny.on?.('exit', (...exitArgs: unknown[]) => { + const code = exitArgs[0] as number | null; + const signal = exitArgs[1] as string | null; + if (code !== null && code !== 0) { + deps.error(`Dashboard process exited with code ${code}`); + if (stderrBuf.trim()) { + deps.error(stderrBuf.trim().split('\n').slice(-5).join('\n')); + } + } else if (signal && signal !== 'SIGINT' && signal !== 'SIGTERM') { + deps.error(`Dashboard process killed by signal ${signal}`); + } + }); + + return child; +} + +async function resolveStartedDashboardPort( + process: DashboardStartupProcess, + preferredPort: number, + deps: CoreDependencies +): Promise { + return new Promise((resolve) => { + let resolved = false; + const processAny = process as DashboardStartupProcess & { + on?: (event: string, cb: (...args: unknown[]) => void) => void; + off?: (event: string, cb: (...args: unknown[]) => void) => void; + removeListener?: (event: string, cb: (...args: unknown[]) => void) => void; + }; + const detach = () => { + process.stdout?.off?.('data', extractPort); + process.stdout?.removeListener?.('data', extractPort); + process.stderr?.off?.('data', extractPort); + process.stderr?.removeListener?.('data', extractPort); + processAny.off?.('exit', handleExit); + processAny.removeListener?.('exit', handleExit); + clearTimeout(timer); + }; + const timer = setTimeout(() => { + if (resolved) return; + resolved = true; + detach(); + deps.warn(`Dashboard did not report its bound port quickly; assuming requested port ${preferredPort}`); + resolve(preferredPort); + }, 3000); + + const finalize = (port: number) => { + if (resolved) return; + resolved = true; + detach(); + resolve(port); + }; + const handleExit = (...exitArgs: unknown[]) => { + const code = exitArgs[0] as number | null; + const signal = exitArgs[1] as string | null; + if (resolved) { + return; + } + resolved = true; + detach(); + if (code !== null && code !== 0) { + deps.warn(`Dashboard exited before reporting its port (code: ${code}).`); + } else if (signal && signal !== 'SIGINT' && signal !== 'SIGTERM') { + deps.warn(`Dashboard exited before reporting its port (signal: ${signal}).`); + } else { + deps.warn('Dashboard exited before reporting its bound port.'); + } + resolve(null); + }; + + const extractPort = (...chunkArgs: unknown[]) => { + const firstChunk = chunkArgs[0]; + if (!firstChunk) { + return; + } + + const chunk = Buffer.isBuffer(firstChunk) + ? firstChunk + : typeof firstChunk === 'string' + ? Buffer.from(firstChunk) + : Buffer.from(JSON.stringify(firstChunk)); + + const match = chunk.toString().match(/Server running at http:\/\/localhost:(\d+)/i); + if (!match?.[1]) { + return; + } + const parsed = Number.parseInt(match[1], 10); + if (!Number.isNaN(parsed)) { + finalize(parsed); + } + }; + + process.stdout?.on?.('data', extractPort); + process.stderr?.on?.('data', extractPort); + processAny.on?.('exit', handleExit); + }); +} + +/** + * Check if the cached dashboard UI assets match the installed dashboard-server + * binary version. If they are stale (or missing a version marker), re-download + * the latest assets from the relay-dashboard GitHub release. + */ +async function refreshDashboardAssetsIfStale( + dashboardBinary: string | null, + deps: CoreDependencies +): Promise { + if (!dashboardBinary || dashboardBinary.endsWith('.js') || dashboardBinary.endsWith('.ts')) { + // Dev mode or npx — skip + return; + } + + // Get installed binary version (async to avoid blocking event loop) + let binaryVersion: string; + try { + const versionResult = await deps.execCommand(`${JSON.stringify(dashboardBinary)} --version`); + binaryVersion = versionResult.stdout.trim(); + } catch { + return; // Can't determine version — skip + } + + if (!binaryVersion) { + return; + } + + const targetDir = getDashboardRootFromBinary(dashboardBinary, deps) ?? getHomeDashboardRoot(deps); + const assetsDir = path.join(targetDir, 'out'); + const versionFile = path.join(targetDir, '.version'); + + // Check if assets match the binary version + try { + const cachedVersion = deps.fs.readFileSync(versionFile, 'utf-8').trim(); + if (cachedVersion === binaryVersion) { + return; // Up to date + } + } catch { + // No version file — need to download if assets exist but are unversioned, + // or if assets don't exist at all + if (deps.fs.existsSync(assetsDir)) { + // Assets exist but no version marker — they're from an old install + } else { + // No assets at all — need to download + } + } + + deps.log(`Updating dashboard UI assets (${binaryVersion})...`); + + const uiUrl = + 'https://github.com/AgentWorkforce/relay-dashboard/releases/latest/download/dashboard-ui.tar.gz'; + let tempDir: string | undefined; + let tempFile: string | undefined; + + try { + tempDir = fs.mkdtempSync(path.join(os.tmpdir(), `dashboard-ui-${deps.pid}-`)); + tempFile = path.join(tempDir, 'dashboard-ui.tar.gz'); + // Download (async to avoid blocking event loop during network I/O) + await deps.execCommand( + `curl -fsSL --max-time 30 ${JSON.stringify(uiUrl)} -o ${JSON.stringify(tempFile)}` + ); + + // Verify it's a valid gzip + const header = Buffer.alloc(2); + const fd = fs.openSync(tempFile, 'r'); + fs.readSync(fd, header, 0, 2, 0); + fs.closeSync(fd); + if (header[0] !== 0x1f || header[1] !== 0x8b) { + if (tempFile) deps.fs.unlinkSync(tempFile); + return; // Not a valid gzip file + } + + // Remove old assets and extract (async to avoid blocking event loop) + deps.fs.rmSync(assetsDir, { recursive: true, force: true }); + deps.fs.mkdirSync(targetDir, { recursive: true }); + await deps.execCommand(`tar -xzf ${JSON.stringify(tempFile)} -C ${JSON.stringify(targetDir)}`); + if (tempFile) deps.fs.unlinkSync(tempFile); + + // Write version marker only after confirming extraction succeeded + if (deps.fs.existsSync(path.join(assetsDir, 'index.html'))) { + deps.fs.writeFileSync(versionFile, binaryVersion); + deps.log(`Dashboard UI assets updated to ${binaryVersion}`); + } else { + deps.warn('Dashboard UI extraction may be incomplete — skipping version marker'); + } + } catch { + // Best-effort — don't block startup + try { + if (tempFile) deps.fs.unlinkSync(tempFile); + } catch { + /* ignore */ + } + } finally { + try { + if (tempDir) deps.fs.rmSync(tempDir, { recursive: true, force: true }); + } catch { + /* ignore */ + } + } +} + +export async function startDashboardWithFallback( + paths: CoreProjectPaths, + dashboardPort: number, + apiPort: number, + deps: CoreDependencies, + enableVerboseLogging: boolean, + relayApiKey?: string, + brokerApiKey?: string +): Promise<{ process: SpawnedProcess; port: number | null }> { + const preferredBinary = deps.findDashboardBinary(); + await refreshDashboardAssetsIfStale(preferredBinary, deps); + let process = startDashboard( + paths, + dashboardPort, + apiPort, + deps, + enableVerboseLogging, + preferredBinary, + relayApiKey, + brokerApiKey + ); + let port = await resolveStartedDashboardPort(process as DashboardStartupProcess, dashboardPort, deps); + + if (port === null && preferredBinary) { + deps.warn('Retrying dashboard startup using npx @agent-relay/dashboard-server@latest'); + process = startDashboard( + paths, + dashboardPort, + apiPort, + deps, + enableVerboseLogging, + null, + relayApiKey, + brokerApiKey + ); + port = await resolveStartedDashboardPort(process as DashboardStartupProcess, dashboardPort, deps); + } + + return { process, port }; +} + +export async function waitForDashboard( + port: number, + process: SpawnedProcess, + deps: Pick, + isShuttingDown: () => boolean +): Promise { + for (let i = 0; i < 20; i++) { + await new Promise((r) => setTimeout(r, 500)); + if (process.killed) { + if (!isShuttingDown()) { + deps.warn(`Warning: Dashboard process exited before becoming ready on port ${port}`); + } + return; + } + try { + const resp = await fetch(`http://localhost:${port}/health`); + if (resp.ok) return; // Dashboard is up + } catch { + // Not ready yet + } + } + if (!isShuttingDown()) { + deps.warn(`Warning: Dashboard not responding on port ${port} after 10s`); + } +} diff --git a/packages/cli/src/cli/lib/broker-lifecycle.ts b/packages/cli/src/cli/lib/broker-lifecycle.ts index 622c671a7..4f15a4cc0 100644 --- a/packages/cli/src/cli/lib/broker-lifecycle.ts +++ b/packages/cli/src/cli/lib/broker-lifecycle.ts @@ -1,5 +1,4 @@ import fs from 'node:fs'; -import os from 'node:os'; import path from 'node:path'; import { HarnessDriverClient } from '@agent-relay/harness-driver'; @@ -11,6 +10,15 @@ import type { SpawnedProcess, } from '../commands/core.js'; import { track } from '../telemetry/index.js'; +import { + getDefaultDashboardRelayUrl, + isDebugLikeLoggingEnabled, + normalizeDashboardPath, + resolveDashboardPortWithFallback, + resolveDashboardRelayUrl, + startDashboardWithFallback, + waitForDashboard, +} from './broker-dashboard.js'; import { buildBundledAgentRelayMcpCommand } from './agent-relay-mcp-command.js'; import { errorClassName } from './telemetry-helpers.js'; import { @@ -267,25 +275,6 @@ function startImplicitLocalFleetSidecar( }); } -async function resolveDashboardPortWithFallback( - dashboardPort: number, - dashboardPortCandidates: number, - deps: CoreDependencies -): Promise { - for (let attempt = 0; attempt < dashboardPortCandidates; attempt += 1) { - const candidatePort = dashboardPort + attempt; - const inUse = await deps.isPortInUse(candidatePort); - if (!inUse) { - if (attempt > 0) { - deps.warn(`Dashboard port ${dashboardPort} is already in use; trying ${candidatePort}`); - } - return candidatePort; - } - } - - throw new Error(`Failed to find an available dashboard port near ${dashboardPort}.`); -} - function isBrokerAlreadyRunningError(message: string): boolean { return /another broker instance is already running in this directory/i.test(message); } @@ -711,552 +700,6 @@ async function waitForBrokerReadiness( return latest; } -function pickDashboardStaticDir(candidates: string[], deps: CoreDependencies): string | null { - const existingCandidates = Array.from(new Set(candidates)).filter((candidate) => - deps.fs.existsSync(candidate) - ); - if (existingCandidates.length === 0) { - return null; - } - - const pageMarkerPriority = [ - ['metrics.html', path.join('metrics', 'index.html')], - ['app.html'], - ['index.html'], - ]; - - for (const markerGroup of pageMarkerPriority) { - const withMarker = existingCandidates.find((candidate) => - markerGroup.some((marker) => deps.fs.existsSync(path.join(candidate, marker))) - ); - if (withMarker) { - return withMarker; - } - } - - return existingCandidates[0]; -} - -function getHomeDashboardRoot(deps: CoreDependencies): string { - const homeDir = deps.env.HOME || deps.env.USERPROFILE || os.homedir(); - return path.join(homeDir, '.agentworkforce/relay', 'dashboard'); -} - -function getPriorDashboardRoot(deps: CoreDependencies): string | null { - const homeDir = deps.env.HOME || deps.env.USERPROFILE || ''; - if (!homeDir) { - return null; - } - return path.join(homeDir, '.relay', 'dashboard'); -} - -function getDashboardRootFromBinary(dashboardBinary: string | null, deps: CoreDependencies): string | null { - if (!dashboardBinary || dashboardBinary.endsWith('.js') || dashboardBinary.endsWith('.ts')) { - return null; - } - - const binaryDir = path.dirname(dashboardBinary); - if (path.basename(binaryDir) !== 'bin') { - return null; - } - - const homeDir = deps.env.HOME || deps.env.USERPROFILE || ''; - const resolvedBinaryDir = path.resolve(binaryDir); - const ignoredBinDirs = [ - homeDir ? path.join(homeDir, '.local', 'bin') : null, - path.join('/usr/local', 'bin'), - ] - .filter((candidate): candidate is string => Boolean(candidate)) - .map((candidate) => path.resolve(candidate)); - if (ignoredBinDirs.includes(resolvedBinaryDir)) { - return null; - } - - return path.join(path.dirname(binaryDir), 'dashboard'); -} - -function resolveDashboardStaticDir(dashboardBinary: string | null, deps: CoreDependencies): string | null { - const explicitStaticDir = deps.env.RELAY_DASHBOARD_STATIC_DIR ?? deps.env.STATIC_DIR; - if (explicitStaticDir && explicitStaticDir.trim()) { - return explicitStaticDir; - } - - if (!dashboardBinary) { - return null; - } - - if (dashboardBinary.endsWith('.js') || dashboardBinary.endsWith('.ts')) { - const dashboardServerOutDir = path.resolve(path.dirname(dashboardBinary), '..', 'out'); - const siblingDashboardOutDir = path.resolve( - path.dirname(dashboardBinary), - '..', - '..', - 'dashboard', - 'out' - ); - return pickDashboardStaticDir([dashboardServerOutDir, siblingDashboardOutDir], deps); - } - - // Installs place UI assets under the install dir (~/.agentworkforce/relay/dashboard/out - // by default, or next to a custom install's bin/ directory). ~/.relay/dashboard/out is - // read as a fallback for installs predating that move. - const installDashboardRoot = getDashboardRootFromBinary(dashboardBinary, deps); - const priorDashboardRoot = getPriorDashboardRoot(deps); - const candidates = [ - installDashboardRoot ? path.join(installDashboardRoot, 'out') : null, - path.join(getHomeDashboardRoot(deps), 'out'), - priorDashboardRoot ? path.join(priorDashboardRoot, 'out') : null, - ].filter((candidate): candidate is string => Boolean(candidate)); - return pickDashboardStaticDir(candidates, deps); -} - -function normalizeLocalhostRelayUrl(relayUrl: string): string { - try { - const parsed = new URL(relayUrl); - if (parsed.hostname === 'localhost') { - parsed.hostname = '127.0.0.1'; - } - return parsed.toString().replace(/\/+$/, ''); - } catch { - return relayUrl; - } -} - -function getDefaultDashboardRelayUrl(apiPort: number): string { - return normalizeLocalhostRelayUrl(`http://localhost:${apiPort}`); -} - -function resolveDashboardRelayUrl(apiPort: number, deps: CoreDependencies): string { - const explicitRelayUrl = deps.env.RELAY_DASHBOARD_RELAY_URL; - if (explicitRelayUrl && explicitRelayUrl.trim()) { - return normalizeLocalhostRelayUrl(explicitRelayUrl.trim()); - } - - return getDefaultDashboardRelayUrl(apiPort); -} - -function isDebugLikeLoggingEnabled(deps: CoreDependencies): boolean { - const rawLevel = String(deps.env.RUST_LOG ?? '').toLowerCase(); - return rawLevel.includes('debug') || rawLevel.includes('trace'); -} - -function getDashboardSpawnEnv( - deps: CoreDependencies, - relayUrl: string, - enableVerboseLogging: boolean, - relayApiKey?: string, - brokerApiKey?: string -): NodeJS.ProcessEnv { - const env: NodeJS.ProcessEnv = { - ...deps.env, - RELAY_URL: relayUrl, - VERBOSE: enableVerboseLogging || deps.env.VERBOSE === 'true' ? 'true' : deps.env.VERBOSE, - }; - // Pass the workspace key so the dashboard can make Agent Relay calls - // (e.g. posting thread replies) without requiring a relaycast.json file. - if (relayApiKey) { - if (!env.RELAY_WORKSPACE_KEY) { - env.RELAY_WORKSPACE_KEY = relayApiKey; - } - if (!env.RELAY_API_KEY) { - env.RELAY_API_KEY = relayApiKey; - } - } - // Pass the broker API key so the dashboard can authenticate with the - // broker's HTTP API (e.g. /api/spawn, /api/spawned). - if (brokerApiKey) { - env.RELAY_BROKER_API_KEY = brokerApiKey; - } - return env; -} - -function getDashboardSpawnArgs( - paths: CoreProjectPaths, - port: number, - apiPort: number, - dashboardBinary: string | null, - relayUrl: string, - enableVerboseLogging: boolean, - deps: CoreDependencies -): string[] { - const args = ['--port', String(port), '--data-dir', paths.dataDir]; - args.push('--relay-url', relayUrl); - const staticDir = resolveDashboardStaticDir(dashboardBinary, deps); - if (staticDir) { - args.push('--static-dir', staticDir); - } - if (enableVerboseLogging) { - args.push('--verbose'); - } - return args; -} - -function normalizeDashboardPath(rawDashboardPath: string | undefined): string | undefined { - const trimmed = rawDashboardPath?.trim(); - if (!trimmed) return undefined; - if (trimmed.startsWith('/')) { - return trimmed; - } - return `/${trimmed}`; -} - -interface DashboardStartupProcess extends SpawnedProcess { - stdout?: { - on?: (event: string, cb: (chunk: Buffer) => void) => void; - removeListener?: (event: string, cb: (...args: unknown[]) => void) => void; - off?: (event: string, cb: (...args: unknown[]) => void) => void; - }; - stderr?: { - on?: (event: string, cb: (chunk: Buffer) => void) => void; - removeListener?: (event: string, cb: (...args: unknown[]) => void) => void; - off?: (event: string, cb: (...args: unknown[]) => void) => void; - }; -} - -function startDashboard( - paths: CoreProjectPaths, - port: number, - apiPort: number, - deps: CoreDependencies, - enableVerboseLogging: boolean, - dashboardBinaryOverride?: string | null, - relayApiKey?: string, - brokerApiKey?: string -): DashboardStartupProcess { - const dashboardBinary = - dashboardBinaryOverride === undefined ? deps.findDashboardBinary() : dashboardBinaryOverride; - const relayUrl = resolveDashboardRelayUrl(apiPort, deps); - const shouldEnableVerbose = enableVerboseLogging || isDebugLikeLoggingEnabled(deps); - const args = getDashboardSpawnArgs( - paths, - port, - apiPort, - dashboardBinary, - relayUrl, - shouldEnableVerbose, - deps - ); - const launchTarget = dashboardBinary - ? dashboardBinary.endsWith('.js') - ? `node ${dashboardBinary}` - : dashboardBinary - : 'npx --yes @agent-relay/dashboard-server@latest'; - - const spawnOpts = { - stdio: ['ignore', 'pipe', 'pipe'] as unknown, - env: getDashboardSpawnEnv(deps, relayUrl, shouldEnableVerbose, relayApiKey, brokerApiKey), - }; - if (shouldEnableVerbose) { - deps.log(`[dashboard] Starting: ${launchTarget} ${args.join(' ')}`); - } - - let child: SpawnedProcess; - if (dashboardBinary) { - // If the binary is a .js file (local dev), run it with node - if (dashboardBinary.endsWith('.js')) { - child = deps.spawnProcess('node', [dashboardBinary, ...args], spawnOpts); - } else { - child = deps.spawnProcess(dashboardBinary, args, spawnOpts); - } - } else { - child = deps.spawnProcess('npx', ['--yes', '@agent-relay/dashboard-server@latest', ...args], spawnOpts); - } - - // Capture stderr for error reporting - const childAny = child as unknown as { - stdout?: { on?: (event: string, cb: (chunk: Buffer) => void) => void }; - stderr?: { on?: (event: string, cb: (chunk: Buffer) => void) => void }; - on?: (event: string, cb: (...args: unknown[]) => void) => void; - }; - let stderrBuf = ''; - - const logChunk = (chunk: Buffer, logger: (line: string) => void, prefix: string) => { - if (!shouldEnableVerbose) { - return; - } - const text = chunk.toString(); - for (const line of text.split(/\r?\n/)) { - const trimmed = line.trim(); - if (trimmed) { - logger(`[dashboard] ${prefix}: ${trimmed}`); - } - } - }; - - childAny.stdout?.on?.('data', (chunk: Buffer) => { - logChunk(chunk, deps.log, 'stdout'); - }); - childAny.stderr?.on?.('data', (chunk: Buffer) => { - stderrBuf += chunk.toString(); - logChunk(chunk, deps.warn, 'stderr'); - }); - - // Report early crashes - childAny.on?.('exit', (...exitArgs: unknown[]) => { - const code = exitArgs[0] as number | null; - const signal = exitArgs[1] as string | null; - if (code !== null && code !== 0) { - deps.error(`Dashboard process exited with code ${code}`); - if (stderrBuf.trim()) { - deps.error(stderrBuf.trim().split('\n').slice(-5).join('\n')); - } - } else if (signal && signal !== 'SIGINT' && signal !== 'SIGTERM') { - deps.error(`Dashboard process killed by signal ${signal}`); - } - }); - - return child; -} - -async function resolveStartedDashboardPort( - process: DashboardStartupProcess, - preferredPort: number, - deps: CoreDependencies -): Promise { - return new Promise((resolve) => { - let resolved = false; - const processAny = process as DashboardStartupProcess & { - on?: (event: string, cb: (...args: unknown[]) => void) => void; - off?: (event: string, cb: (...args: unknown[]) => void) => void; - removeListener?: (event: string, cb: (...args: unknown[]) => void) => void; - }; - const detach = () => { - process.stdout?.off?.('data', extractPort); - process.stdout?.removeListener?.('data', extractPort); - process.stderr?.off?.('data', extractPort); - process.stderr?.removeListener?.('data', extractPort); - processAny.off?.('exit', handleExit); - processAny.removeListener?.('exit', handleExit); - clearTimeout(timer); - }; - const timer = setTimeout(() => { - if (resolved) return; - resolved = true; - detach(); - deps.warn(`Dashboard did not report its bound port quickly; assuming requested port ${preferredPort}`); - resolve(preferredPort); - }, 3000); - - const finalize = (port: number) => { - if (resolved) return; - resolved = true; - detach(); - resolve(port); - }; - const handleExit = (...exitArgs: unknown[]) => { - const code = exitArgs[0] as number | null; - const signal = exitArgs[1] as string | null; - if (resolved) { - return; - } - resolved = true; - detach(); - if (code !== null && code !== 0) { - deps.warn(`Dashboard exited before reporting its port (code: ${code}).`); - } else if (signal && signal !== 'SIGINT' && signal !== 'SIGTERM') { - deps.warn(`Dashboard exited before reporting its port (signal: ${signal}).`); - } else { - deps.warn('Dashboard exited before reporting its bound port.'); - } - resolve(null); - }; - - const extractPort = (...chunkArgs: unknown[]) => { - const firstChunk = chunkArgs[0]; - if (!firstChunk) { - return; - } - - const chunk = Buffer.isBuffer(firstChunk) - ? firstChunk - : typeof firstChunk === 'string' - ? Buffer.from(firstChunk) - : Buffer.from(JSON.stringify(firstChunk)); - - const match = chunk.toString().match(/Server running at http:\/\/localhost:(\d+)/i); - if (!match?.[1]) { - return; - } - const parsed = Number.parseInt(match[1], 10); - if (!Number.isNaN(parsed)) { - finalize(parsed); - } - }; - - process.stdout?.on?.('data', extractPort); - process.stderr?.on?.('data', extractPort); - processAny.on?.('exit', handleExit); - }); -} - -/** - * Check if the cached dashboard UI assets match the installed dashboard-server - * binary version. If they are stale (or missing a version marker), re-download - * the latest assets from the relay-dashboard GitHub release. - */ -async function refreshDashboardAssetsIfStale( - dashboardBinary: string | null, - deps: CoreDependencies -): Promise { - if (!dashboardBinary || dashboardBinary.endsWith('.js') || dashboardBinary.endsWith('.ts')) { - // Dev mode or npx — skip - return; - } - - // Get installed binary version (async to avoid blocking event loop) - let binaryVersion: string; - try { - const versionResult = await deps.execCommand(`${JSON.stringify(dashboardBinary)} --version`); - binaryVersion = versionResult.stdout.trim(); - } catch { - return; // Can't determine version — skip - } - - if (!binaryVersion) { - return; - } - - const targetDir = getDashboardRootFromBinary(dashboardBinary, deps) ?? getHomeDashboardRoot(deps); - const assetsDir = path.join(targetDir, 'out'); - const versionFile = path.join(targetDir, '.version'); - - // Check if assets match the binary version - try { - const cachedVersion = deps.fs.readFileSync(versionFile, 'utf-8').trim(); - if (cachedVersion === binaryVersion) { - return; // Up to date - } - } catch { - // No version file — need to download if assets exist but are unversioned, - // or if assets don't exist at all - if (deps.fs.existsSync(assetsDir)) { - // Assets exist but no version marker — they're from an old install - } else { - // No assets at all — need to download - } - } - - deps.log(`Updating dashboard UI assets (${binaryVersion})...`); - - const uiUrl = - 'https://github.com/AgentWorkforce/relay-dashboard/releases/latest/download/dashboard-ui.tar.gz'; - let tempDir: string | undefined; - let tempFile: string | undefined; - - try { - tempDir = fs.mkdtempSync(path.join(os.tmpdir(), `dashboard-ui-${deps.pid}-`)); - tempFile = path.join(tempDir, 'dashboard-ui.tar.gz'); - // Download (async to avoid blocking event loop during network I/O) - await deps.execCommand( - `curl -fsSL --max-time 30 ${JSON.stringify(uiUrl)} -o ${JSON.stringify(tempFile)}` - ); - - // Verify it's a valid gzip - const header = Buffer.alloc(2); - const fd = fs.openSync(tempFile, 'r'); - fs.readSync(fd, header, 0, 2, 0); - fs.closeSync(fd); - if (header[0] !== 0x1f || header[1] !== 0x8b) { - if (tempFile) deps.fs.unlinkSync(tempFile); - return; // Not a valid gzip file - } - - // Remove old assets and extract (async to avoid blocking event loop) - deps.fs.rmSync(assetsDir, { recursive: true, force: true }); - deps.fs.mkdirSync(targetDir, { recursive: true }); - await deps.execCommand(`tar -xzf ${JSON.stringify(tempFile)} -C ${JSON.stringify(targetDir)}`); - if (tempFile) deps.fs.unlinkSync(tempFile); - - // Write version marker only after confirming extraction succeeded - if (deps.fs.existsSync(path.join(assetsDir, 'index.html'))) { - deps.fs.writeFileSync(versionFile, binaryVersion); - deps.log(`Dashboard UI assets updated to ${binaryVersion}`); - } else { - deps.warn('Dashboard UI extraction may be incomplete — skipping version marker'); - } - } catch { - // Best-effort — don't block startup - try { - if (tempFile) deps.fs.unlinkSync(tempFile); - } catch { - /* ignore */ - } - } finally { - try { - if (tempDir) deps.fs.rmSync(tempDir, { recursive: true, force: true }); - } catch { - /* ignore */ - } - } -} - -async function startDashboardWithFallback( - paths: CoreProjectPaths, - dashboardPort: number, - apiPort: number, - deps: CoreDependencies, - enableVerboseLogging: boolean, - relayApiKey?: string, - brokerApiKey?: string -): Promise<{ process: SpawnedProcess; port: number | null }> { - const preferredBinary = deps.findDashboardBinary(); - await refreshDashboardAssetsIfStale(preferredBinary, deps); - let process = startDashboard( - paths, - dashboardPort, - apiPort, - deps, - enableVerboseLogging, - preferredBinary, - relayApiKey, - brokerApiKey - ); - let port = await resolveStartedDashboardPort(process as DashboardStartupProcess, dashboardPort, deps); - - if (port === null && preferredBinary) { - deps.warn('Retrying dashboard startup using npx @agent-relay/dashboard-server@latest'); - process = startDashboard( - paths, - dashboardPort, - apiPort, - deps, - enableVerboseLogging, - null, - relayApiKey, - brokerApiKey - ); - port = await resolveStartedDashboardPort(process as DashboardStartupProcess, dashboardPort, deps); - } - - return { process, port }; -} - -async function waitForDashboard( - port: number, - process: SpawnedProcess, - deps: Pick, - isShuttingDown: () => boolean -): Promise { - for (let i = 0; i < 20; i++) { - await new Promise((r) => setTimeout(r, 500)); - if (process.killed) { - if (!isShuttingDown()) { - deps.warn(`Warning: Dashboard process exited before becoming ready on port ${port}`); - } - return; - } - try { - const resp = await fetch(`http://localhost:${port}/health`); - if (resp.ok) return; // Dashboard is up - } catch { - // Not ready yet - } - } - if (!isShuttingDown()) { - deps.warn(`Warning: Dashboard not responding on port ${port} after 10s`); - } -} - async function discoverExistingBrokerApiPort( preferredApiPort: number, maxAttempts: number, diff --git a/packages/cli/src/cli/mcp/action-schema.ts b/packages/cli/src/cli/mcp/action-schema.ts new file mode 100644 index 000000000..c214ba3cf --- /dev/null +++ b/packages/cli/src/cli/mcp/action-schema.ts @@ -0,0 +1,157 @@ +import type { + ActionSchema, + AgentRelayActionDescriptor, + JsonSchemaLiteObject, + ZodLikeSchema, +} from '@agent-relay/sdk/actions'; +import { z } from 'zod'; + +/** + * Adapters between the JSON-schema-lite action schema format used by the Agent + * Relay actions registry and the zod schemas the MCP SDK expects for tool input + * validation. Also handles serializing descriptors for the `list_actions` tool. + */ + +export function isSchemaObject(schema: ActionSchema | undefined): schema is JsonSchemaLiteObject { + return Boolean( + schema && + typeof schema === 'object' && + !Array.isArray(schema) && + typeof (schema as { safeParse?: unknown }).safeParse !== 'function' + ); +} + +function getSchemaDescription(schema: ActionSchema | undefined): string | undefined { + return isSchemaObject(schema) && typeof schema.description === 'string' ? schema.description : undefined; +} + +/** Convert a JSON-schema-lite node into an equivalent zod type. */ +export function zodFromJsonSchema(schema: ActionSchema | undefined): z.ZodTypeAny { + if (schema === false) { + return z.never(); + } + + if (!isSchemaObject(schema)) { + return z.unknown(); + } + + let zodType: z.ZodTypeAny; + const schemaType = Array.isArray(schema.type) ? schema.type[0] : schema.type; + switch (schemaType) { + case 'array': + zodType = z.array(zodFromJsonSchema(schema.items)); + break; + case 'boolean': + zodType = z.boolean(); + break; + case 'integer': + zodType = z.number().int(); + break; + case 'number': + zodType = z.number(); + break; + case 'object': + if (schema.properties) { + const required = new Set(schema.required ?? []); + const shape: Record = {}; + for (const [key, childSchema] of Object.entries(schema.properties)) { + const child = zodFromJsonSchema(childSchema); + shape[key] = required.has(key) ? child : child.optional(); + } + zodType = z.object(shape).passthrough(); + } else { + zodType = z.record(z.string(), z.unknown()); + } + break; + case 'string': + zodType = z.string(); + break; + default: + zodType = z.unknown(); + break; + } + + const description = getSchemaDescription(schema); + return description ? zodType.describe(description) : zodType; +} + +function zodObjectShape(schema: ActionSchema | undefined): Record | undefined { + if (schema instanceof z.ZodObject) { + return schema.shape; + } + return undefined; +} + +/** Build the MCP tool `inputSchema` (a zod shape) from an action input schema. */ +export function actionToolInputSchema(schema: ActionSchema | undefined): Record { + const zodShape = zodObjectShape(schema); + if (zodShape) { + return zodShape; + } + + if (!isSchemaObject(schema) || schema.type !== 'object') { + return { + input: z.unknown().describe('Action input payload. The action registry performs final validation.'), + }; + } + + const required = new Set(schema.required ?? []); + const shape: Record = {}; + for (const [key, childSchema] of Object.entries(schema.properties ?? {})) { + const child = zodFromJsonSchema(childSchema); + shape[key] = required.has(key) ? child : child.optional(); + } + return shape; +} + +/** + * Normalize raw MCP tool args into the shape the action handler expects: object + * schemas pass through, while scalar/opaque schemas unwrap the `{ input }` envelope. + */ +export function actionInvocationInput(descriptor: AgentRelayActionDescriptor, args: unknown): unknown { + const schema = descriptor.inputSchema; + if (zodObjectShape(schema)) { + return args; + } + if (!isSchemaObject(schema) || schema.type !== 'object') { + return typeof args === 'object' && args !== null && 'input' in args + ? (args as { input?: unknown }).input + : args; + } + return args; +} + +function isZodLikeSchema(schema: ActionSchema | undefined): schema is ZodLikeSchema { + return Boolean( + schema && + typeof schema === 'object' && + !Array.isArray(schema) && + typeof (schema as { safeParse?: unknown }).safeParse === 'function' + ); +} + +function serializableActionSchema(schema: ActionSchema): unknown { + if (isSchemaObject(schema)) { + return schema; + } + if (isZodLikeSchema(schema)) { + return { + type: 'zod', + ...(schema.description ? { description: schema.description } : {}), + }; + } + return schema; +} + +/** Project an action descriptor into a JSON-serializable form for `list_actions`. */ +export function serializableActionDescriptor( + descriptor: AgentRelayActionDescriptor +): Record { + return { + name: descriptor.name, + description: descriptor.description, + visibility: descriptor.visibility, + ...(descriptor.inputSchema ? { inputSchema: serializableActionSchema(descriptor.inputSchema) } : {}), + ...(descriptor.outputSchema ? { outputSchema: serializableActionSchema(descriptor.outputSchema) } : {}), + }; +} diff --git a/packages/cli/src/cli/mcp/action-tools.ts b/packages/cli/src/cli/mcp/action-tools.ts new file mode 100644 index 000000000..a5d5cc818 --- /dev/null +++ b/packages/cli/src/cli/mcp/action-tools.ts @@ -0,0 +1,151 @@ +import { McpServer } from '@modelcontextprotocol/sdk/server/mcp.js'; +import type { ActionAuditEvent, AgentRelayActions } from '@agent-relay/sdk/actions'; +import { z } from 'zod'; + +import { + actionInvocationInput, + actionToolInputSchema, + serializableActionDescriptor, +} from './action-schema.js'; +import { jsonContent, jsonResult } from './tool-results.js'; +import type { AgentClientLike, SessionState } from './types.js'; + +/** The relay-backed action surface on the live agent client, when available. */ +function getRelayAgentActions( + getAgentClient?: (asIdentity?: string) => AgentClientLike +): AgentClientLike['actions'] | undefined { + if (!getAgentClient) { + return undefined; + } + try { + return getAgentClient().actions; + } catch { + return undefined; + } +} + +function asInputRecord(input: unknown): Record | undefined { + if (input === undefined || input === null) { + return undefined; + } + if (typeof input === 'object' && !Array.isArray(input)) { + return input as Record; + } + return { input }; +} + +function errorMessage(error: unknown): string { + return error instanceof Error ? error.message : String(error); +} + +/** + * Register the `list_actions` and `invoke_action` tools plus a dynamic per-action + * tool for each registered action. Invocation is fire-and-forget through the + * relay action surface when available, falling back to the in-process registry. + */ +export function registerAgentRelayActionTools( + server: McpServer, + actions: AgentRelayActions | undefined, + getSession: () => SessionState, + onAuditEvent?: (event: ActionAuditEvent) => Promise | void, + getAgentClient?: (asIdentity?: string) => AgentClientLike, + actionToolNames?: Set +): void { + if (!actions) { + return; + } + + /** + * Fire-and-forget invocation through the relay: returns an immediate ack + * (with an `invocation_id`) and does NOT run the handler inline. Falls back to + * the local in-process registry when the relay action surface is unavailable. + */ + const invokeAction = async (name: string, input: unknown) => { + const relayActions = getRelayAgentActions(getAgentClient); + if (relayActions) { + try { + const ack = await relayActions.invoke(name, asInputRecord(input)); + return jsonContent({ ok: true, status: 'invoked', invocation: ack }); + } catch (error) { + return { ...jsonContent({ ok: false, error: errorMessage(error) }), isError: true }; + } + } + + const session = getSession(); + const result = await actions.invoke({ + name, + input, + context: { + caller: { name: session.agentName ?? 'mcp', type: 'agent' }, + emit: onAuditEvent, + }, + }); + return result.ok ? jsonContent(result) : { ...jsonContent(result), isError: true }; + }; + + server.registerTool( + 'list_actions', + { + title: 'List Actions', + description: 'List Agent Relay actions available to this agent.', + inputSchema: {}, + outputSchema: jsonResult, + annotations: { + readOnlyHint: true, + destructiveHint: false, + idempotentHint: true, + openWorldHint: false, + }, + }, + async () => + jsonContent({ + actions: (await actions.list({ visibility: 'agent' })).map(serializableActionDescriptor), + }) + ); + + server.registerTool( + 'invoke_action', + { + title: 'Invoke Action', + description: + 'Invoke a registered Agent Relay action by name. Fire-and-forget: returns an ack with an invocation id; the result arrives asynchronously to the action handler.', + inputSchema: { + name: z.string().describe('Registered action name'), + input: z.unknown().describe('Action input payload'), + }, + outputSchema: jsonResult, + annotations: { + readOnlyHint: false, + destructiveHint: false, + idempotentHint: false, + openWorldHint: false, + }, + }, + async ({ name, input }: { name: string; input: unknown }) => invokeAction(name, input) + ); + + void actions + .list({ visibility: 'agent' }) + .then((descriptors) => { + for (const descriptor of descriptors) { + actionToolNames?.add(descriptor.name); + server.registerTool( + descriptor.name, + { + title: descriptor.name, + description: descriptor.description, + inputSchema: actionToolInputSchema(descriptor.inputSchema), + outputSchema: jsonResult, + annotations: { + readOnlyHint: false, + destructiveHint: false, + idempotentHint: false, + openWorldHint: false, + }, + }, + async (args: unknown) => invokeAction(descriptor.name, actionInvocationInput(descriptor, args)) + ); + } + }) + .catch(() => undefined); +} diff --git a/packages/cli/src/cli/mcp/inbox.ts b/packages/cli/src/cli/mcp/inbox.ts new file mode 100644 index 000000000..acee29e64 --- /dev/null +++ b/packages/cli/src/cli/mcp/inbox.ts @@ -0,0 +1,56 @@ +/** Tools whose results should never be augmented with piggybacked inbox state. */ +export const SKIP_PIGGYBACK = new Set([ + 'check_inbox', + 'create_workspace', + 'set_workspace_key', + 'register_agent', +]); + +/** + * Render an inbox payload into a compact human-readable summary appended to tool + * results. When `selfName` is provided, the agent's own mentions, DMs, and + * reactions are filtered out. Returns an empty string when there is nothing pending. + */ +export function formatInbox(inbox: any, selfName?: string | null): string { + const norm = (s: string) => s.trim().replace(/^@/, '').toLowerCase(); + const selfNorm = selfName ? norm(selfName) : null; + const isSelf = (name: string) => selfNorm != null && norm(name) === selfNorm; + const lines = ['--- Pending Messages ---']; + + if (inbox.unreadChannels?.length) { + lines.push('Unread channels:'); + for (const ch of inbox.unreadChannels) { + lines.push(` #${ch.channelName}: ${ch.unreadCount} unread`); + } + } + + const mentions = selfNorm ? inbox.mentions?.filter((m: any) => !isSelf(m.agentName)) : inbox.mentions; + if (mentions?.length) { + lines.push('Mentions:'); + for (const m of mentions) { + lines.push(` @${m.agentName} in #${m.channelName}: "${m.text}"`); + } + } + + const dms = selfNorm ? inbox.unreadDms?.filter((dm: any) => !isSelf(dm.from)) : inbox.unreadDms; + if (dms?.length) { + lines.push('Unread DMs:'); + for (const dm of dms) { + lines.push(` From ${dm.from}: ${dm.unreadCount} unread`); + } + } + + const reactions = selfNorm + ? inbox.recentReactions?.filter((reaction: any) => !isSelf(reaction.agentName)) + : inbox.recentReactions; + if (reactions?.length) { + lines.push('Reactions (informational; no response required):'); + for (const reaction of reactions) { + lines.push( + ` :${reaction.emoji}: on your message in #${reaction.channelName} by @${reaction.agentName}` + ); + } + } + + return lines.length === 1 ? '' : lines.join('\n'); +} diff --git a/packages/cli/src/cli/mcp/messaging-tools.ts b/packages/cli/src/cli/mcp/messaging-tools.ts new file mode 100644 index 000000000..591101058 --- /dev/null +++ b/packages/cli/src/cli/mcp/messaging-tools.ts @@ -0,0 +1,428 @@ +import { McpServer } from '@modelcontextprotocol/sdk/server/mcp.js'; +import { z } from 'zod'; + +import { jsonContent, jsonResult, textContent } from './tool-results.js'; +import { identityOverrideInputShape, messageResult } from './tool-shapes.js'; +import type { AgentClientLike } from './types.js'; + +function resolveEmoji(input: string): string { + const normalized = input.trim().replace(/^:/, '').replace(/:$/, '').toLowerCase(); + const aliases: Record = { + '+1': '👍', + thumbsup: '👍', + thumbs_up: '👍', + check: '✅', + white_check_mark: '✅', + rocket: '🚀', + eyes: '👀', + heart: '❤️', + clap: '👏', + }; + return aliases[normalized] ?? input; +} + +/** + * Register the channel, message, thread, DM, reaction, search, and inbox MCP + * tools. These all act through a single agent client resolved per-call from the + * optional `as` identity override. + */ +export function registerMessagingTools( + server: McpServer, + getAgentClient: (asIdentity?: string) => AgentClientLike +): void { + server.registerTool( + 'create_channel', + { + title: 'Create Channel', + description: 'Create a new workspace channel.', + inputSchema: { + name: z.string().describe('Unique channel name'), + topic: z.string().optional().describe('Optional channel topic'), + ...identityOverrideInputShape, + }, + outputSchema: jsonResult, + annotations: { + readOnlyHint: false, + destructiveHint: false, + idempotentHint: false, + openWorldHint: true, + }, + }, + async ({ name, topic, as }) => jsonContent(await getAgentClient(as).channels.create({ name, topic })) + ); + + server.registerTool( + 'list_channels', + { + title: 'List Channels', + description: 'List channels available in the workspace.', + inputSchema: { + include_archived: z.boolean().optional().describe('Include archived channels'), + ...identityOverrideInputShape, + }, + outputSchema: { + channels: z.array(z.object({}).passthrough()).describe('Channels'), + }, + annotations: { readOnlyHint: true, destructiveHint: false, idempotentHint: true, openWorldHint: true }, + }, + async ({ include_archived, as }) => { + const channels = await getAgentClient(as).channels.list( + include_archived ? { includeArchived: include_archived } : undefined + ); + return jsonContent({ channels }); + } + ); + + server.registerTool( + 'join_channel', + { + title: 'Join Channel', + description: 'Join an existing channel.', + inputSchema: { + channel: z.string().describe('Channel name'), + ...identityOverrideInputShape, + }, + outputSchema: messageResult, + annotations: { readOnlyHint: false, destructiveHint: false, idempotentHint: true, openWorldHint: true }, + }, + async ({ channel, as }) => { + await getAgentClient(as).channels.join(channel); + return textContent(`Joined channel #${channel}`); + } + ); + + server.registerTool( + 'leave_channel', + { + title: 'Leave Channel', + description: 'Leave a channel.', + inputSchema: { + channel: z.string().describe('Channel name'), + ...identityOverrideInputShape, + }, + outputSchema: messageResult, + annotations: { readOnlyHint: false, destructiveHint: false, idempotentHint: true, openWorldHint: true }, + }, + async ({ channel, as }) => { + await getAgentClient(as).channels.leave(channel); + return textContent(`Left channel #${channel}`); + } + ); + + server.registerTool( + 'invite_to_channel', + { + title: 'Invite to Channel', + description: 'Invite another agent to a channel.', + inputSchema: { + channel: z.string().describe('Channel name'), + agent: z.string().describe('Agent name to invite'), + ...identityOverrideInputShape, + }, + outputSchema: messageResult, + annotations: { readOnlyHint: false, destructiveHint: false, idempotentHint: true, openWorldHint: true }, + }, + async ({ channel, agent, as }) => { + await getAgentClient(as).channels.invite(channel, agent); + return textContent(`Invited ${agent} to #${channel}`); + } + ); + + server.registerTool( + 'set_channel_topic', + { + title: 'Set Channel Topic', + description: 'Update a channel topic.', + inputSchema: { + channel: z.string().describe('Channel name'), + topic: z.string().describe('New topic'), + ...identityOverrideInputShape, + }, + outputSchema: jsonResult, + annotations: { readOnlyHint: false, destructiveHint: false, idempotentHint: true, openWorldHint: true }, + }, + async ({ channel, topic, as }) => jsonContent(await getAgentClient(as).channels.setTopic(channel, topic)) + ); + + server.registerTool( + 'archive_channel', + { + title: 'Archive Channel', + description: 'Archive a channel.', + inputSchema: { + channel: z.string().describe('Channel name'), + ...identityOverrideInputShape, + }, + outputSchema: messageResult, + annotations: { readOnlyHint: false, destructiveHint: true, idempotentHint: true, openWorldHint: true }, + }, + async ({ channel, as }) => { + await getAgentClient(as).channels.archive(channel); + return textContent(`Archived channel #${channel}`); + } + ); + + server.registerTool( + 'post_message', + { + title: 'Post Message', + description: 'Post a new message to a channel as the current agent.', + inputSchema: { + channel: z.string().describe('Channel name'), + text: z.string().describe('Message text'), + attachments: z.array(z.string()).optional().describe('File attachment IDs'), + mode: z.enum(['wait', 'steer']).optional().describe('Delivery mode'), + ...identityOverrideInputShape, + }, + outputSchema: jsonResult, + annotations: { + readOnlyHint: false, + destructiveHint: false, + idempotentHint: false, + openWorldHint: true, + }, + }, + async ({ channel, text, attachments, mode, as }) => + jsonContent(await getAgentClient(as).send(channel, text, { attachments, mode })) + ); + + server.registerTool( + 'list_messages', + { + title: 'Get Messages', + description: 'Retrieve message history from a channel.', + inputSchema: { + channel: z.string().describe('Channel name'), + limit: z.number().optional().describe('Maximum messages to return'), + before: z.string().optional().describe('Older-than cursor'), + after: z.string().optional().describe('Newer-than cursor'), + ...identityOverrideInputShape, + }, + outputSchema: { + messages: z.array(z.object({}).passthrough()).describe('Messages'), + }, + annotations: { readOnlyHint: true, destructiveHint: false, idempotentHint: true, openWorldHint: true }, + }, + async ({ channel, limit, before, after, as }) => { + const messages = await getAgentClient(as).messages(channel, { limit, before, after }); + return jsonContent({ messages }); + } + ); + + server.registerTool( + 'reply_to_thread', + { + title: 'Reply to Thread', + description: 'Reply to an existing message thread.', + inputSchema: { + message_id: z.string().describe('Parent message ID'), + text: z.string().describe('Reply text'), + ...identityOverrideInputShape, + }, + outputSchema: jsonResult, + annotations: { + readOnlyHint: false, + destructiveHint: false, + idempotentHint: false, + openWorldHint: true, + }, + }, + async ({ message_id, text, as }) => jsonContent(await getAgentClient(as).reply(message_id, text)) + ); + + server.registerTool( + 'get_message_thread', + { + title: 'Get Thread', + description: 'Retrieve a message thread.', + inputSchema: { + message_id: z.string().describe('Parent message ID'), + limit: z.number().optional().describe('Maximum replies to return'), + ...identityOverrideInputShape, + }, + outputSchema: jsonResult, + annotations: { readOnlyHint: true, destructiveHint: false, idempotentHint: true, openWorldHint: true }, + }, + async ({ message_id, limit, as }) => + jsonContent(await getAgentClient(as).thread(message_id, limit ? { limit } : undefined)) + ); + + server.registerTool( + 'send_dm', + { + title: 'Send Direct Message', + description: 'Send a private direct message to another agent.', + inputSchema: { + to: z.string().describe('Recipient agent name'), + text: z.string().describe('DM text'), + mode: z.enum(['wait', 'steer']).optional().describe('Delivery mode'), + attachments: z.array(z.string()).optional().describe('File attachment IDs'), + ...identityOverrideInputShape, + }, + outputSchema: jsonResult, + annotations: { + readOnlyHint: false, + destructiveHint: false, + idempotentHint: false, + openWorldHint: true, + }, + }, + async ({ to, text, mode, attachments, as }) => + jsonContent(await getAgentClient(as).dm(to, text, { mode, attachments })) + ); + + server.registerTool( + 'list_dms', + { + title: 'List DM Conversations', + description: 'List direct message conversations for the current agent.', + inputSchema: { + ...identityOverrideInputShape, + }, + outputSchema: { + conversations: z.array(z.object({}).passthrough()).describe('DM conversations'), + }, + annotations: { readOnlyHint: true, destructiveHint: false, idempotentHint: true, openWorldHint: true }, + }, + async ({ as }) => jsonContent({ conversations: await getAgentClient(as).dms.conversations() }) + ); + + server.registerTool( + 'send_group_dm', + { + title: 'Send Group DM', + description: 'Create a group DM and send the first message.', + inputSchema: { + participants: z.array(z.string()).describe('Participant agent names'), + name: z.string().optional().describe('Optional group name'), + text: z.string().describe('Initial message'), + ...identityOverrideInputShape, + }, + outputSchema: jsonResult, + annotations: { + readOnlyHint: false, + destructiveHint: false, + idempotentHint: false, + openWorldHint: true, + }, + }, + async ({ participants, name, text, as }) => { + const client = getAgentClient(as); + const conversation = await client.dms.createGroup({ participants, name }); + const message = await client.dms.sendMessage(conversation.id, text); + return jsonContent({ conversation, message }); + } + ); + + server.registerTool( + 'add_reaction', + { + title: 'Add Reaction', + description: 'Add an emoji reaction to a message.', + inputSchema: { + message_id: z.string().describe('Message ID'), + emoji: z.string().describe('Emoji character or shortcode'), + ...identityOverrideInputShape, + }, + outputSchema: messageResult, + annotations: { readOnlyHint: false, destructiveHint: false, idempotentHint: true, openWorldHint: true }, + }, + async ({ message_id, emoji, as }) => { + const resolved = resolveEmoji(emoji); + await getAgentClient(as).react(message_id, resolved); + return textContent(`Reacted with ${resolved}`); + } + ); + + server.registerTool( + 'remove_reaction', + { + title: 'Remove Reaction', + description: 'Remove an emoji reaction from a message.', + inputSchema: { + message_id: z.string().describe('Message ID'), + emoji: z.string().describe('Emoji character or shortcode'), + ...identityOverrideInputShape, + }, + outputSchema: messageResult, + annotations: { readOnlyHint: false, destructiveHint: false, idempotentHint: true, openWorldHint: true }, + }, + async ({ message_id, emoji, as }) => { + const resolved = resolveEmoji(emoji); + await getAgentClient(as).unreact(message_id, resolved); + return textContent(`Removed reaction ${resolved}`); + } + ); + + server.registerTool( + 'search_messages', + { + title: 'Search Messages', + description: 'Search messages across the workspace.', + inputSchema: { + query: z.string().describe('Text search query'), + channel: z.string().optional().describe('Optional channel filter'), + from: z.string().optional().describe('Optional sender filter'), + limit: z.number().optional().describe('Maximum results'), + ...identityOverrideInputShape, + }, + outputSchema: { + results: z.array(z.object({}).passthrough()).describe('Search results'), + }, + annotations: { readOnlyHint: true, destructiveHint: false, idempotentHint: true, openWorldHint: true }, + }, + async ({ query, channel, from, limit, as }) => + jsonContent({ results: await getAgentClient(as).search(query, { channel, from, limit }) }) + ); + + server.registerTool( + 'check_inbox', + { + title: 'Check Inbox', + description: 'Check unread messages, mentions, DMs, and reactions for the current agent.', + inputSchema: { + limit: z.number().optional().describe('Maximum inbox items'), + ...identityOverrideInputShape, + }, + outputSchema: jsonResult, + annotations: { readOnlyHint: true, destructiveHint: false, idempotentHint: true, openWorldHint: true }, + }, + async ({ limit, as }) => + jsonContent(await getAgentClient(as).inbox(limit != null ? { limit } : undefined)) + ); + + server.registerTool( + 'mark_message_read', + { + title: 'Mark as Read', + description: 'Mark a message as read for the current agent.', + inputSchema: { + message_id: z.string().describe('Message ID'), + ...identityOverrideInputShape, + }, + outputSchema: messageResult, + annotations: { readOnlyHint: false, destructiveHint: false, idempotentHint: true, openWorldHint: true }, + }, + async ({ message_id, as }) => { + await getAgentClient(as).markRead(message_id); + return textContent(`Marked message ${message_id} as read`); + } + ); + + server.registerTool( + 'get_message_readers', + { + title: 'Get Readers', + description: 'List agents who have read a message.', + inputSchema: { + message_id: z.string().describe('Message ID'), + ...identityOverrideInputShape, + }, + outputSchema: { + readers: z.array(z.object({}).passthrough()).describe('Readers'), + }, + annotations: { readOnlyHint: true, destructiveHint: false, idempotentHint: true, openWorldHint: true }, + }, + async ({ message_id, as }) => jsonContent({ readers: await getAgentClient(as).readers(message_id) }) + ); +} diff --git a/packages/cli/src/cli/mcp/resources.ts b/packages/cli/src/cli/mcp/resources.ts new file mode 100644 index 000000000..8abd2afe6 --- /dev/null +++ b/packages/cli/src/cli/mcp/resources.ts @@ -0,0 +1,209 @@ +import { McpServer, ResourceTemplate } from '@modelcontextprotocol/sdk/server/mcp.js'; +import { type AgentClient, type RelayCast, type WsClient } from '@relaycast/sdk'; + +/** + * Tracks the set of `relay://` resource URIs the MCP client has subscribed to, + * so realtime events can be filtered down to only the resources that matter. + */ +export class SubscriptionManager { + private readonly subscriptions = new Set(); + + subscribe(uri: string): void { + this.subscriptions.add(uri); + } + + unsubscribe(uri: string): void { + this.subscriptions.delete(uri); + } + + getMatchingSubscriptions(uris: string[]): string[] { + return uris.filter((uri) => this.subscriptions.has(uri)); + } + + getAll(): string[] { + return [...this.subscriptions]; + } + + clear(): void { + this.subscriptions.clear(); + } +} + +function getStringEventField(event: unknown, field: string): string | null { + if (typeof event !== 'object' || event === null) { + return null; + } + const candidate = (event as Record)[field]; + return typeof candidate === 'string' ? candidate : null; +} + +/** + * Map a realtime workspace event to the `relay://` resource URIs whose contents + * it invalidates, so subscribers can be notified to re-fetch. + */ +export function eventToResourceUris(event: unknown): string[] { + const type = getStringEventField(event, 'type'); + switch (type) { + case 'message.created': { + const channel = getStringEventField(event, 'channel'); + return channel ? ['relay://inbox', `relay://channels/${channel}/messages`] : ['relay://inbox']; + } + case 'message.updated': { + const channel = getStringEventField(event, 'channel'); + return channel ? [`relay://channels/${channel}/messages`] : []; + } + case 'thread.reply': { + const parentId = getStringEventField(event, 'parentId'); + return parentId ? ['relay://inbox', `relay://messages/${parentId}/thread`] : ['relay://inbox']; + } + case 'dm.received': + case 'group_dm.received': { + const conversationId = getStringEventField(event, 'conversationId'); + return conversationId ? ['relay://inbox', `relay://dm/${conversationId}`] : ['relay://inbox']; + } + case 'agent.online': + case 'agent.offline': + return ['relay://agents']; + case 'channel.created': + case 'channel.updated': + case 'channel.archived': + case 'member.joined': + case 'member.left': + return ['relay://channels']; + case 'webhook.received': + case 'command.invoked': { + const channel = getStringEventField(event, 'channel'); + return channel ? [`relay://channels/${channel}/messages`] : []; + } + case 'reaction.added': + case 'reaction.removed': + return ['relay://inbox']; + default: + return []; + } +} + +/** + * Bridges the realtime WebSocket event stream to MCP resource notifications: + * for each workspace event, notifies any subscribed `relay://` resource URIs + * that the event invalidates. + */ +export class RealtimeResourceBridge { + private unsubscribeFn: (() => void) | null = null; + + constructor( + private readonly wsClient: WsClient, + private readonly subscriptions: SubscriptionManager, + private readonly notifyCallback: (uri: string) => void + ) {} + + start(): void { + this.unsubscribeFn = this.wsClient.on('*', (event) => { + const type = getStringEventField(event, 'type'); + if ( + type === 'open' || + type === 'close' || + type === 'error' || + type === 'reconnecting' || + type === 'permanently_disconnected' + ) { + return; + } + const matched = this.subscriptions.getMatchingSubscriptions(eventToResourceUris(event)); + for (const uri of matched) { + this.notifyCallback(uri); + } + }); + this.wsClient.connect(); + } + + stop(): void { + if (this.unsubscribeFn) { + this.unsubscribeFn(); + this.unsubscribeFn = null; + } + this.wsClient.disconnect(); + } +} + +/** + * Register the read-only `relay://` resources (inbox, agents, channels, channel + * messages, message threads, and DM conversations) on the MCP server. + */ +export function registerResourceDefinitions( + server: McpServer, + getAgentClient: (asIdentity?: string) => AgentClient, + getRelay: () => RelayCast +): void { + server.registerResource( + 'inbox', + 'relay://inbox', + { title: 'Inbox', description: 'Unread messages, mentions, and DMs', mimeType: 'application/json' }, + async (uri) => { + const inbox = await getAgentClient().inbox(); + return { contents: [{ uri: uri.href, text: JSON.stringify(inbox) }] }; + } + ); + + server.registerResource( + 'agents', + 'relay://agents', + { + title: 'Agents', + description: 'Online and offline agents in the workspace', + mimeType: 'application/json', + }, + async (uri) => { + const agents = await getRelay().agents.list(); + return { contents: [{ uri: uri.href, text: JSON.stringify(agents) }] }; + } + ); + + server.registerResource( + 'channels', + 'relay://channels', + { title: 'Channels', description: 'Available channels in the workspace', mimeType: 'application/json' }, + async (uri) => { + const channels = await getAgentClient().channels.list(); + return { contents: [{ uri: uri.href, text: JSON.stringify(channels) }] }; + } + ); + + server.registerResource( + 'channel-messages', + new ResourceTemplate('relay://channels/{name}/messages', { list: undefined }), + { + title: 'Channel Messages', + description: 'Messages in a specific channel', + mimeType: 'application/json', + }, + async (uri, params) => { + const messages = await getAgentClient().messages(String(params.name)); + return { contents: [{ uri: uri.href, text: JSON.stringify(messages) }] }; + } + ); + + server.registerResource( + 'message-thread', + new ResourceTemplate('relay://messages/{id}/thread', { list: undefined }), + { title: 'Message Thread', description: 'Thread replies on a message', mimeType: 'application/json' }, + async (uri, params) => { + const thread = await getAgentClient().thread(String(params.id)); + return { contents: [{ uri: uri.href, text: JSON.stringify(thread) }] }; + } + ); + + server.registerResource( + 'dm-conversation', + new ResourceTemplate('relay://dm/{conversation_id}', { list: undefined }), + { + title: 'DM Conversation', + description: 'Direct message conversation', + mimeType: 'application/json', + }, + async (uri, params) => { + const messages = await getAgentClient().dms.messages(String(params.conversation_id)); + return { contents: [{ uri: uri.href, text: JSON.stringify(messages) }] }; + } + ); +} diff --git a/packages/cli/src/cli/mcp/telemetry.ts b/packages/cli/src/cli/mcp/telemetry.ts new file mode 100644 index 000000000..c197e1866 --- /dev/null +++ b/packages/cli/src/cli/mcp/telemetry.ts @@ -0,0 +1,216 @@ +import { McpServer } from '@modelcontextprotocol/sdk/server/mcp.js'; +import { + agentTokenRecoveryMessage, + isInvalidAgentTokenError, + isInvalidAgentTokenToolResult, +} from '@agent-relay/sdk'; + +import { track, type AgentRelayToolCallCategory, type AgentRelayToolCallType } from '../telemetry/index.js'; +import { errorClassName } from '../lib/telemetry-helpers.js'; +import { formatInbox, SKIP_PIGGYBACK } from './inbox.js'; +import { hasContentArray, invalidAgentTokenToolResult, isErrorToolResult } from './tool-results.js'; +import type { AgentClientLike, AgentRelayMcpServerOptions, SessionState } from './types.js'; + +interface AgentRelayToolCallMetadata { + toolType: AgentRelayToolCallType; + toolCategory: AgentRelayToolCallCategory; +} + +/** + * Owned tools that delegate to the actions surface (`actions.invoke(...)`) + * rather than the agents/messaging APIs. Together with the dynamic per-action + * tools (tracked via `actionToolNames`), these intentionally skip per-tool + * telemetry so the same underlying action is not counted differently depending + * on which MCP surface the caller used (e.g. `spawn` vs `invoke_action`). + */ +const ACTION_ROUTED_TOOL_NAMES = new Set(['invoke_action', 'spawn']); + +/** + * Coarse type/category metadata for the statically-registered ("owned") MCP + * tools. Action-routed calls (see `ACTION_ROUTED_TOOL_NAMES`) and the dynamic + * per-action tools surfaced from the actions registry are intentionally + * excluded from per-tool telemetry (see the skip in `enableInboxPiggyback`), + * so they have no entry. + */ +const AGENT_RELAY_TOOL_CALL_METADATA = { + add_agent: { toolType: 'agent.create', toolCategory: 'spawn' }, + remove_agent: { toolType: 'agent.release', toolCategory: 'release' }, + list_actions: { toolType: 'action.list', toolCategory: 'action' }, + submit_result: { toolType: 'result.submit', toolCategory: 'result' }, + create_workspace: { toolType: 'workspace.create', toolCategory: 'workspace' }, + set_workspace_key: { toolType: 'workspace.set_key', toolCategory: 'workspace' }, + register_agent: { toolType: 'agent.register', toolCategory: 'agent' }, + list_agents: { toolType: 'agent.list', toolCategory: 'agent' }, + post_message: { toolType: 'message.post', toolCategory: 'message' }, + send_dm: { toolType: 'message.dm', toolCategory: 'message' }, + send_group_dm: { toolType: 'message.group_dm', toolCategory: 'message' }, + list_dms: { toolType: 'message.dm_list', toolCategory: 'message' }, + list_messages: { toolType: 'message.list', toolCategory: 'message' }, + reply_to_thread: { toolType: 'message.reply', toolCategory: 'message' }, + get_message_thread: { toolType: 'message.thread', toolCategory: 'message' }, + search_messages: { toolType: 'message.search', toolCategory: 'message' }, + create_channel: { toolType: 'channel.create', toolCategory: 'channel' }, + list_channels: { toolType: 'channel.list', toolCategory: 'channel' }, + join_channel: { toolType: 'channel.join', toolCategory: 'channel' }, + leave_channel: { toolType: 'channel.leave', toolCategory: 'channel' }, + set_channel_topic: { toolType: 'channel.set_topic', toolCategory: 'channel' }, + archive_channel: { toolType: 'channel.archive', toolCategory: 'channel' }, + invite_to_channel: { toolType: 'channel.invite', toolCategory: 'channel' }, + add_reaction: { toolType: 'reaction.add', toolCategory: 'reaction' }, + remove_reaction: { toolType: 'reaction.remove', toolCategory: 'reaction' }, + check_inbox: { toolType: 'inbox.check', toolCategory: 'inbox' }, + mark_message_read: { toolType: 'inbox.mark_read', toolCategory: 'inbox' }, + get_message_readers: { toolType: 'inbox.reader_list', toolCategory: 'inbox' }, +} satisfies Record; + +function agentRelayToolCallMetadata(name: string): AgentRelayToolCallMetadata { + const known = (AGENT_RELAY_TOOL_CALL_METADATA as Partial>)[name]; + return known ?? { toolType: name, toolCategory: 'tool' }; +} + +function trackAgentRelayToolCall(input: { + toolName: string; + toolType: AgentRelayToolCallType; + toolCategory: AgentRelayToolCallCategory; + transport?: AgentRelayMcpServerOptions['telemetryTransport']; + startedAt: number; + success: boolean; + errorClass?: string; +}): void { + track('agent_relay_tool_call', { + tool_name: input.toolName, + tool_type: input.toolType, + tool_category: input.toolCategory, + transport: input.transport ?? 'unknown', + success: input.success, + duration_ms: Date.now() - input.startedAt, + ...(input.errorClass ? { error_class: input.errorClass } : {}), + }); +} + +function readAsIdentity(args: unknown[]): string | undefined { + const [input] = args; + if (typeof input !== 'object' || input === null) return undefined; + const as = (input as { as?: unknown }).as; + return typeof as === 'string' ? as : undefined; +} + +/** + * Wrap `server.registerTool` so every owned tool emits `agent_relay_tool_call` + * telemetry, surfaces invalid-agent-token errors as recoverable results, and + * piggybacks a compact inbox summary onto successful results. + */ +export function enableInboxPiggyback( + mcpServer: McpServer, + getSession: () => SessionState, + getAgentClient: (asIdentity?: string) => AgentClientLike, + invalidateAgentToken: (asIdentity?: string) => void, + telemetryTransport?: AgentRelayMcpServerOptions['telemetryTransport'], + actionToolNames = new Set() +): void { + const original = mcpServer.registerTool.bind(mcpServer); + const mutableServer = mcpServer as McpServer & { + registerTool: McpServer['registerTool']; + }; + + mutableServer.registerTool = (name: string, config: any, handler: any) => { + if (!handler) { + return original(name, config, handler); + } + + const wrapped = async (...args: unknown[]) => { + const asIdentity = readAsIdentity(args); + const startedAt = Date.now(); + // Action-routed calls (`invoke_action`, `spawn`, and the dynamic + // per-action tools) run through the actions surface and deliberately skip + // per-tool telemetry; only the owned tools emit `agent_relay_tool_call`. + const toolMetadata = + !ACTION_ROUTED_TOOL_NAMES.has(name) && !actionToolNames.has(name) + ? agentRelayToolCallMetadata(name) + : undefined; + + let result: any; + try { + result = await handler(...args); + } catch (err) { + if (name !== 'register_agent' && isInvalidAgentTokenError(err)) { + invalidateAgentToken(asIdentity); + if (toolMetadata) { + trackAgentRelayToolCall({ + toolName: name, + toolType: toolMetadata.toolType, + toolCategory: toolMetadata.toolCategory, + transport: telemetryTransport, + startedAt, + success: false, + errorClass: errorClassName(err) ?? 'InvalidAgentToken', + }); + } + return invalidAgentTokenToolResult(); + } + if (toolMetadata) { + trackAgentRelayToolCall({ + toolName: name, + toolType: toolMetadata.toolType, + toolCategory: toolMetadata.toolCategory, + transport: telemetryTransport, + startedAt, + success: false, + errorClass: errorClassName(err), + }); + } + throw err; + } + + if (name !== 'register_agent' && isInvalidAgentTokenToolResult(result)) { + invalidateAgentToken(asIdentity); + if (toolMetadata) { + trackAgentRelayToolCall({ + toolName: name, + toolType: toolMetadata.toolType, + toolCategory: toolMetadata.toolCategory, + transport: telemetryTransport, + startedAt, + success: false, + errorClass: 'InvalidAgentToken', + }); + } + if (hasContentArray(result)) { + result.content.push({ type: 'text', text: agentTokenRecoveryMessage() }); + } + return result; + } + + if (!SKIP_PIGGYBACK.has(name) && getSession().agentToken && hasContentArray(result)) { + try { + const inbox = await getAgentClient(asIdentity).inbox(); + const inboxText = formatInbox(inbox, asIdentity ?? getSession().agentName); + if (inboxText) { + result.content.push({ type: 'text', text: inboxText }); + } + } catch (err) { + if (isInvalidAgentTokenError(err)) { + invalidateAgentToken(asIdentity); + } + } + } + + if (toolMetadata) { + const resultIsError = isErrorToolResult(result); + trackAgentRelayToolCall({ + toolName: name, + toolType: toolMetadata.toolType, + toolCategory: toolMetadata.toolCategory, + transport: telemetryTransport, + startedAt, + success: !resultIsError, + ...(resultIsError ? { errorClass: 'ToolResultError' } : {}), + }); + } + + return result; + }; + + return original(name, config, wrapped); + }; +} diff --git a/packages/cli/src/cli/mcp/tool-results.ts b/packages/cli/src/cli/mcp/tool-results.ts new file mode 100644 index 000000000..407a7859c --- /dev/null +++ b/packages/cli/src/cli/mcp/tool-results.ts @@ -0,0 +1,52 @@ +import { INVALID_AGENT_TOKEN_CODE, agentTokenRecoveryMessage } from '@agent-relay/sdk'; +import { z } from 'zod'; + +/** Permissive output schema for tools that return arbitrary JSON objects. */ +export const jsonResult = z.object({}).passthrough(); + +export type JsonToolResult = { + content: Array<{ type: 'text'; text: string }>; + structuredContent: Record; +}; + +/** Wrap an arbitrary value as an MCP tool result with both text and structured content. */ +export function jsonContent(value: unknown): JsonToolResult { + const structuredContent = + typeof value === 'object' && value !== null && !Array.isArray(value) + ? (value as Record) + : { value }; + return { + content: [{ type: 'text', text: JSON.stringify(value, null, 2) }], + structuredContent, + }; +} + +/** Wrap a human-readable message as an MCP tool result. */ +export function textContent(message: string, structuredContent: Record = { message }) { + return { + content: [{ type: 'text' as const, text: message }], + structuredContent, + }; +} + +export function hasContentArray(value: unknown): value is { content: Array> } { + return ( + typeof value === 'object' && value !== null && Array.isArray((value as { content?: unknown }).content) + ); +} + +export function isErrorToolResult(value: unknown): boolean { + return Boolean(value && typeof value === 'object' && (value as { isError?: unknown }).isError === true); +} + +/** Standard MCP error result describing an invalid/expired agent token. */ +export function invalidAgentTokenToolResult(): JsonToolResult & { isError: true } { + const text = agentTokenRecoveryMessage(); + return { + content: [{ type: 'text', text }], + structuredContent: { + error: { code: INVALID_AGENT_TOKEN_CODE, message: text }, + }, + isError: true, + }; +} diff --git a/packages/cli/src/cli/mcp/tool-shapes.ts b/packages/cli/src/cli/mcp/tool-shapes.ts new file mode 100644 index 000000000..4886ea3ef --- /dev/null +++ b/packages/cli/src/cli/mcp/tool-shapes.ts @@ -0,0 +1,17 @@ +import { z } from 'zod'; + +/** Standard output schema for tools that return a single confirmation message. */ +export const messageResult = { + message: z.string().describe('Human-readable confirmation message'), +}; + +/** + * Optional `as` input field that lets a tool act on behalf of one of several + * registered agent identities in the same MCP session. + */ +export const identityOverrideInputShape = { + as: z + .string() + .optional() + .describe('Registered agent identity to act as when multiple identities have been registered'), +}; diff --git a/packages/cli/src/cli/mcp/types.ts b/packages/cli/src/cli/mcp/types.ts new file mode 100644 index 000000000..2cb43b758 --- /dev/null +++ b/packages/cli/src/cli/mcp/types.ts @@ -0,0 +1,44 @@ +import type { RelayCast, AgentClient } from '@relaycast/sdk'; +import type { ActionAuditEvent, AgentRelayActions } from '@agent-relay/sdk/actions'; + +import type { RealtimeResourceBridge, SubscriptionManager } from './resources.js'; + +export type AgentType = 'agent' | 'human'; +export type RelayCastLike = Pick; +export type AgentClientLike = AgentClient; + +export interface AgentRelayMcpServerOptions { + workspaceKey?: string; + /** @deprecated Use workspaceKey. */ + apiKey?: string; + baseUrl?: string; + agentToken?: string; + agentName?: string; + agentType?: AgentType; + strictAgentName?: boolean; + telemetryTransport?: 'stdio' | 'http'; + skipBootstrap?: boolean; + actions?: AgentRelayActions; + onActionAuditEvent?: (event: ActionAuditEvent) => Promise | void; +} + +export interface RegisteredAgent { + agentName: string; + agentToken: string; +} + +export interface SessionState { + workspaceKey: string | null; + agentToken: string | null; + agentName: string | null; + agents: Map; + wsBridge: RealtimeResourceBridge | null; + subscriptions: SubscriptionManager | null; + wsInitAttempted: boolean; +} + +export type RegistrationSession = Pick & { + agents?: Map; +}; + +export type SessionSetter = (partial: Partial) => void; diff --git a/packages/cli/src/cli/mcp/workspace.ts b/packages/cli/src/cli/mcp/workspace.ts new file mode 100644 index 000000000..978825a84 --- /dev/null +++ b/packages/cli/src/cli/mcp/workspace.ts @@ -0,0 +1,48 @@ +import { RelayCast } from '@relaycast/sdk'; + +import { relaycastWorkspaceTelemetryOptions } from '../lib/relaycast-telemetry.js'; +import type { RegistrationSession } from './types.js'; + +/** Create a new RelayCast workspace, returning the raw provisioning payload. */ +export async function createWorkspace(name: string, baseUrl?: string): Promise> { + return (await RelayCast.createWorkspace(name, { + baseUrl, + ...relaycastWorkspaceTelemetryOptions(), + })) as Record; +} + +/** Extract a workspace key from a provisioning payload, tolerating naming variants. */ +export function extractWorkspaceKey(payload: Record): string | undefined { + const data = + payload.data && typeof payload.data === 'object' ? (payload.data as Record) : {}; + const value = + payload.workspaceKey ?? + payload.workspace_key ?? + payload.apiKey ?? + payload.api_key ?? + data.workspaceKey ?? + data.workspace_key ?? + data.apiKey ?? + data.api_key; + + return typeof value === 'string' && value.trim() ? value : undefined; +} + +/** Extract a workspace name from a provisioning payload, falling back to `fallback`. */ +export function extractWorkspaceName(payload: Record, fallback: string): string { + const data = + payload.data && typeof payload.data === 'object' ? (payload.data as Record) : {}; + const value = payload.workspaceName ?? payload.workspace_name ?? payload.name ?? data.workspaceName; + return typeof value === 'string' && value.trim() ? value : fallback; +} + +/** Throw a descriptive error when the session has no workspace key configured. */ +export function requireWorkspaceKey(session: RegistrationSession): void { + if (session.workspaceKey) { + return; + } + + throw new Error( + 'Workspace key not configured. Call "create_workspace" first, or "set_workspace_key" if someone shared a workspace key.' + ); +}