#!/usr/bin/env node import { createRequire } from 'node:module'; import { McpServer, ResourceTemplate } from '@modelcontextprotocol/sdk/server/mcp.js'; import { StdioServerTransport } from '@modelcontextprotocol/sdk/server/stdio.js'; import { SubscribeRequestSchema, UnsubscribeRequestSchema, } from '@modelcontextprotocol/sdk/types.js'; import { CodexWorkerApp } from './app.js'; import { SubscriptionManager } from './services/subscription-manager.js'; import { spawnTaskSchema, waitTaskSchema, respondTaskSchema, messageTaskSchema, cancelTaskSchema, getToolDescriptions, getToolAnnotation, } from './mcp/tool-definitions.js'; import { buildStatusBanner, buildRespondBanner } from './services/tool-description-banner.js'; const require = createRequire(import.meta.url); const pkg = require('../package.json') as { version: string }; const app = new CodexWorkerApp(); // --------------------------------------------------------------------------- // McpServer — the high-level MCP SDK entry point (Server class is deprecated) // --------------------------------------------------------------------------- const server = new McpServer( { name: 'mcp-codex-worker', version: pkg.version }, { capabilities: { resources: { subscribe: true, listChanged: true }, }, instructions: [ 'Codex Worker MCP Server — orchestrate coding agents across providers.', '', 'Workflow: spawn-task → wait-task → (if input_required) respond-task → wait-task.', 'Use message-task for follow-ups, cancel-task to abort.', '', 'Resources: task:///all (scoreboard), task:///{id} (detail), task:///{id}/events (debug trace).', 'Poll task:///all every ~30s to track all tasks.', ].join('\n'), }, ); // --------------------------------------------------------------------------- // Tool registration — registerTool with Zod schemas + annotations // --------------------------------------------------------------------------- const descMap = getToolDescriptions(pkg.version); function desc(name: string): string { return descMap.get(name) ?? name; } function anno(name: string) { return getToolAnnotation(name); } /** Wrap an app.callTool call into a CallToolResult. */ async function callTool(name: string, args: unknown): Promise<{ content: Array<{ type: 'text'; text: string }>; isError?: true }> { try { const text = await app.callTool(name, args); return { content: [{ type: 'text', text }] }; } catch (error) { return { content: [{ type: 'text', text: error instanceof Error ? error.message : String(error) }], isError: true, }; } } /** * Build a progress reporter from the handler's extra parameter. * Claude Code sends a progressToken in _meta when calling tools — we can * push notifications/progress back through it during long-running calls. */ function makeProgressReporter(extra: Record) { const meta = extra._meta as Record | undefined; const token = meta?.progressToken as string | number | undefined; const sendNotification = extra.sendNotification as ((n: unknown) => Promise) | undefined; if (!token || !sendNotification) return undefined; return async (progress: number, total: number, message?: string) => { await sendNotification({ method: 'notifications/progress', params: { progressToken: token, progress, total, ...(message ? { message } : {}) }, }).catch(() => {}); }; } server.registerTool('spawn-task', { description: desc('spawn-task'), inputSchema: spawnTaskSchema, annotations: anno('spawn-task'), }, (args: unknown) => callTool('spawn-task', args)); server.registerTool('wait-task', { description: desc('wait-task'), inputSchema: waitTaskSchema, annotations: anno('wait-task'), }, async (args: unknown, extra) => { // wait-task is the long-polling tool — push progress while waiting const reportProgress = makeProgressReporter(extra); try { const text = await app.callToolWithProgress('wait-task', args, reportProgress); return { content: [{ type: 'text', text }] }; } catch (error) { return { content: [{ type: 'text', text: error instanceof Error ? error.message : String(error) }], isError: true as const, }; } }); const respondTool = server.registerTool('respond-task', { description: desc('respond-task'), inputSchema: respondTaskSchema, annotations: anno('respond-task'), }, (args: unknown) => callTool('respond-task', args)); const messageTool = server.registerTool('message-task', { description: desc('message-task'), inputSchema: messageTaskSchema, annotations: anno('message-task'), }, (args: unknown) => callTool('message-task', args)); server.registerTool('cancel-task', { description: desc('cancel-task'), inputSchema: cancelTaskSchema, annotations: anno('cancel-task'), }, (args: unknown) => callTool('cancel-task', args)); // --------------------------------------------------------------------------- // Resource registration — registerResource with static URIs + ResourceTemplates // --------------------------------------------------------------------------- server.registerResource( 'Task Scoreboard', 'task:///all', { mimeType: 'text/plain', description: 'Compact overview of all tracked tasks with status badges' }, async (uri) => { const resource = await app.readResource(uri.href); return { contents: [{ uri: uri.href, mimeType: resource.mimeType, text: resource.text }] }; }, ); const taskDetailTemplate = new ResourceTemplate('task:///{id}', { list: async () => ({ resources: app.getTaskManager().getAllTasks().map(task => ({ uri: `task:///${task.id}`, name: `Task ${task.id}`, description: `Detail: ${task.prompt.slice(0, 50)}`, })), }), }); server.registerResource( 'Task Detail', taskDetailTemplate, { mimeType: 'text/markdown', description: 'Full detail view for a specific task' }, async (uri) => { const resource = await app.readResource(uri.href); return { contents: [{ uri: uri.href, mimeType: resource.mimeType, text: resource.text }] }; }, ); const taskLogTemplate = new ResourceTemplate('task:///{id}/log', { list: async () => ({ resources: app.getTaskManager().getAllTasks().map(task => ({ uri: `task:///${task.id}/log`, name: `Task ${task.id} log`, description: 'Summary log for this task', })), }), }); server.registerResource( 'Task Summary Log', taskLogTemplate, { mimeType: 'text/plain', description: 'Last 20 output lines for a task' }, async (uri) => { const resource = await app.readResource(uri.href); return { contents: [{ uri: uri.href, mimeType: resource.mimeType, text: resource.text }] }; }, ); const taskVerboseTemplate = new ResourceTemplate('task:///{id}/log.verbose', { list: undefined, }); server.registerResource( 'Task Verbose Log', taskVerboseTemplate, { mimeType: 'text/plain', description: 'Full output log for a task' }, async (uri) => { const resource = await app.readResource(uri.href); return { contents: [{ uri: uri.href, mimeType: resource.mimeType, text: resource.text }] }; }, ); const taskEventsTemplate = new ResourceTemplate('task:///{id}/events', { list: undefined, }); server.registerResource( 'Task Event Stream', taskEventsTemplate, { mimeType: 'application/jsonl', description: 'Raw JSONL event trace — every Codex notification with timestamp' }, async (uri) => { const resource = await app.readResource(uri.href); return { contents: [{ uri: uri.href, mimeType: resource.mimeType, text: resource.text }] }; }, ); const taskEventsSummaryTemplate = new ResourceTemplate('task:///{id}/events/summary', { list: undefined, }); server.registerResource( 'Task Event Summary', taskEventsSummaryTemplate, { mimeType: 'application/jsonl', description: 'Filtered event trace — excludes delta events (reasoning, message, command output). Much smaller than full events.' }, async (uri) => { const resource = await app.readResource(uri.href); return { contents: [{ uri: uri.href, mimeType: resource.mimeType, text: resource.text }] }; }, ); const taskTimelineTemplate = new ResourceTemplate('task:///{id}/timeline', { list: undefined, }); server.registerResource( 'Task Timeline', taskTimelineTemplate, { mimeType: 'text/plain', description: 'Pre-filtered execution timeline — one line per meaningful event, tail-friendly' }, async (uri) => { const resource = await app.readResource(uri.href); return { contents: [{ uri: uri.href, mimeType: resource.mimeType, text: resource.text }] }; }, ); // --------------------------------------------------------------------------- // Subscriptions — custom tracking for targeted resource update notifications // --------------------------------------------------------------------------- const subscriptions = new SubscriptionManager(); // Subscribe/unsubscribe need the underlying Server for custom handlers that // track subscriptions. McpServer's built-in handler doesn't expose state. server.server.setRequestHandler(SubscribeRequestSchema, async (request) => { subscriptions.subscribe(request.params.uri); return {}; }); server.server.setRequestHandler(UnsubscribeRequestSchema, async (request) => { subscriptions.unsubscribe(request.params.uri); return {}; }); // --------------------------------------------------------------------------- // Main // --------------------------------------------------------------------------- async function main(): Promise { await app.initialize(); // Auto-subscribe the scoreboard so every connected client gets updates subscriptions.subscribe('task:///all'); const taskManager = app.getTaskManager(); // --- Task creation: new resource URIs appear --- taskManager.onTaskCreate(() => { server.sendResourceListChanged(); server.server.sendResourceUpdated({ uri: 'task:///all' }).catch(() => {}); }); // --- Status changes: task completed, failed, etc. --- // Debounced tool description refresh for dynamic banners. let toolListTimer: ReturnType | null = null; const scheduleToolDescriptionRefresh = (): void => { if (toolListTimer) return; toolListTimer = setTimeout(() => { toolListTimer = null; const allTasks = taskManager.getAllTasks(); const statusBanner = allTasks.length > 0 ? buildStatusBanner(allTasks) : ''; const respondBannerText = allTasks.length > 0 ? buildRespondBanner(allTasks) : ''; // Update tool descriptions with fresh banners const freshDescs = getToolDescriptions(pkg.version); respondTool.update({ description: (freshDescs.get('respond-task') ?? '') + (respondBannerText ? `\n\n${respondBannerText}` : '') }); messageTool.update({ description: (freshDescs.get('message-task') ?? '') + (statusBanner ? `\n\n${statusBanner}` : '') }); // Notify clients that tool descriptions have changed server.sendToolListChanged(); }, 1000); toolListTimer.unref?.(); }; taskManager.onStatusChange((task) => { const uris = subscriptions.getMatchingSubscriptions(task.id); for (const uri of uris) { server.server.sendResourceUpdated({ uri }).catch(() => {}); } scheduleToolDescriptionRefresh(); }); // --- Output changes: new agent messages, command output, diffs --- // Throttled: notify at most once per second per task. const outputNotifyTimers = new Map>(); taskManager.onOutput((taskId) => { if (outputNotifyTimers.has(taskId)) return; const timer = setTimeout(() => { outputNotifyTimers.delete(taskId); const uris = subscriptions.getMatchingSubscriptions(taskId); for (const uri of uris) { server.server.sendResourceUpdated({ uri }).catch(() => {}); } }, 1000); timer.unref(); outputNotifyTimers.set(taskId, timer); }); const transport = new StdioServerTransport(); await server.connect(transport); process.stdin.resume(); const shutdown = async () => { if (toolListTimer) clearTimeout(toolListTimer); await server.close().catch(() => {}); await app.shutdown().catch(() => {}); process.exit(0); }; process.on('SIGINT', () => void shutdown()); process.on('SIGTERM', () => void shutdown()); } main().catch((error) => { console.error(error); process.exit(1); });