diff --git a/deploy/deployment.yaml b/deploy/deployment.yaml index 830ec9c..ca8f085 100644 --- a/deploy/deployment.yaml +++ b/deploy/deployment.yaml @@ -25,6 +25,14 @@ spec: env: - name: CHE_MCP_TRANSPORT value: "http" + - name: CHE_MCP_DATA_DIR + value: "/data" + securityContext: + allowPrivilegeEscalation: false + runAsNonRoot: true + volumeMounts: + - name: data + mountPath: /data livenessProbe: httpGet: path: /healthz @@ -37,3 +45,7 @@ spec: port: 8080 initialDelaySeconds: 3 periodSeconds: 5 + volumes: + - name: data + persistentVolumeClaim: + claimName: che-mcp-server-data diff --git a/deploy/kustomization.yaml b/deploy/kustomization.yaml index 02111fb..9e78132 100644 --- a/deploy/kustomization.yaml +++ b/deploy/kustomization.yaml @@ -4,6 +4,7 @@ resources: - rbac.yaml - deployment.yaml - service.yaml + - pvc.yaml images: - name: quay.io/che-incubator/che-mcp-server diff --git a/deploy/pvc.yaml b/deploy/pvc.yaml new file mode 100644 index 0000000..699a371 --- /dev/null +++ b/deploy/pvc.yaml @@ -0,0 +1,11 @@ +apiVersion: v1 +kind: PersistentVolumeClaim +metadata: + name: che-mcp-server-data +spec: + accessModes: + - ReadWriteOnce + storageClassName: gp3-csi + resources: + requests: + storage: 1Gi diff --git a/src/config.ts b/src/config.ts index 3357979..833ef47 100644 --- a/src/config.ts +++ b/src/config.ts @@ -28,3 +28,5 @@ function getArgValue(argv: string[], flag: string): string | undefined { if (index === -1 || index + 1 >= argv.length) return undefined; return argv[index + 1]; } + +export const DATA_DIR = process.env.CHE_MCP_DATA_DIR ?? '/data'; diff --git a/src/kube/exec.ts b/src/kube/exec.ts index 3e041d9..297c959 100644 --- a/src/kube/exec.ts +++ b/src/kube/exec.ts @@ -1,5 +1,5 @@ import * as k8s from '@kubernetes/client-node'; -import stream from 'stream'; +import stream from 'node:stream'; import { getCoreV1Api, diff --git a/src/messaging/store.ts b/src/messaging/store.ts new file mode 100644 index 0000000..d43a908 --- /dev/null +++ b/src/messaging/store.ts @@ -0,0 +1,125 @@ +import { randomUUID } from 'node:crypto'; +import { existsSync, mkdirSync, readFileSync, renameSync, rmSync, writeFileSync } from 'node:fs'; +import { join } from 'node:path'; +import { DATA_DIR } from '../config.js'; + +const inboxes: Map = new Map(); +let dataFile = ''; +let tmpFile = ''; + +export interface Message { + message_id: string; + from: string; + to: string; + body: string; + thread_id: string; + timestamp: string; +} + +export function initStore(dataDir: string): void { + dataFile = join(dataDir, 'messages.json'); + tmpFile = join(dataDir, 'messages.json.tmp'); + mkdirSync(dataDir, { recursive: true }); + inboxes.clear(); + // CRASH RECOVERY: promote a completed write that survived a crash + if (existsSync(tmpFile)) { + try { + JSON.parse(readFileSync(tmpFile, 'utf8')); + renameSync(tmpFile, dataFile); + } catch { + rmSync(tmpFile, { force: true }); + } + } + if (existsSync(dataFile)) { + try { + const entries = JSON.parse(readFileSync(dataFile, 'utf8')); + for (const [key, msgs] of entries) { + inboxes.set(key, msgs); + } + } catch { + // Corrupted data file — start with empty state + } + } +} + +export function sendMessage( + from: string, + to: string, + body: string, + thread_id?: string, +): { message_id: string; thread_id: string } { + const message_id = randomUUID(); + const resolvedThreadId = thread_id ?? randomUUID(); + + const message: Message = { + message_id, + from, + to, + body, + thread_id: resolvedThreadId, + timestamp: new Date().toISOString(), + }; + + const inbox = inboxes.get(to) ?? []; + inbox.push(message); + inboxes.set(to, inbox); + + flushToDisk(); + + return { message_id, thread_id: resolvedThreadId }; +} + +export function receiveMessages( + sessionId: string, + threadId?: string, +): { messages: Message[] } { + const inbox = inboxes.get(sessionId); + if (!inbox || inbox.length === 0) { + return { messages: [] }; + } + + if (threadId) { + const matching = inbox.filter(m => m.thread_id === threadId); + const remaining = inbox.filter(m => m.thread_id !== threadId); + if (remaining.length === 0) { + inboxes.delete(sessionId); + } else { + inboxes.set(sessionId, remaining); + } + flushToDisk(); + return { messages: matching }; + } + + inboxes.delete(sessionId); + flushToDisk(); + return { messages: inbox }; +} + +export function getUnreadCount(sessionId: string): number { + return inboxes.get(sessionId)?.length ?? 0; +} + +export function clearAllInboxes(): void { + inboxes.clear(); + if (dataFile !== '' && existsSync(dataFile)) { + rmSync(dataFile); + } + if (tmpFile !== '' && existsSync(tmpFile)) { + rmSync(tmpFile); + } +} + +function flushToDisk(): void { + if (dataFile === '') return; + const json = JSON.stringify(Array.from(inboxes.entries())); + try { + writeFileSync(tmpFile, json, 'utf8'); + renameSync(tmpFile, dataFile); + } catch (e) { + console.error('flushToDisk failed:', e); + } +} + +if (!process.env.VITEST) { + initStore(DATA_DIR); +} diff --git a/src/tools.ts b/src/tools.ts index e17d083..3af43a3 100644 --- a/src/tools.ts +++ b/src/tools.ts @@ -1,26 +1,28 @@ import { McpServer } from '@modelcontextprotocol/sdk/server/mcp.js'; import { z } from 'zod'; -import type { ServerMode } from './types.js'; -import { listWorkspaces } from './tools/list-workspaces.js'; -import { startTerminalSession } from './tools/start-terminal-session.js'; -import { readTerminalOutput } from './tools/read-terminal-output.js'; -import { sendTerminalInput } from './tools/send-terminal-input.js'; -import { getTerminalState } from './tools/get-terminal-state.js'; -import { stopTerminalSession } from './tools/stop-terminal-session.js'; -import { execInWorkspace } from './tools/exec-in-workspace.js'; import { createWorkspace } from './tools/create-workspace.js'; -import { startWorkspace } from './tools/start-workspace.js'; -import { stopWorkspace } from './tools/stop-workspace.js'; import { deleteWorkspace } from './tools/delete-workspace.js'; -import { getWorkspaceStatus } from './tools/get-workspace-status.js'; +import { execInWorkspace } from './tools/exec-in-workspace.js'; +import { getAgentOutputTool } from './tools/get-agent-output.js'; +import { getAgentStatusTool } from './tools/get-agent-status.js'; +import { getTerminalState } from './tools/get-terminal-state.js'; import { getWorkspacePod } from './tools/get-workspace-pod.js'; +import { getWorkspaceStatus } from './tools/get-workspace-status.js'; +import { injectTool } from './tools/inject-tool.js'; import { launchCodingAgentTool } from './tools/launch-coding-agent.js'; -import { getAgentStatusTool } from './tools/get-agent-status.js'; import { listAllAgentsTool } from './tools/list-all-agents.js'; +import { listWorkspaces } from './tools/list-workspaces.js'; +import { readTerminalOutput } from './tools/read-terminal-output.js'; +import { receiveMessagesTool } from './tools/receive-messages.js'; +import { sendMessageTool } from './tools/send-message.js'; import { sendMessageToAgentTool } from './tools/send-message-to-agent.js'; -import { getAgentOutputTool } from './tools/get-agent-output.js'; +import { sendTerminalInput } from './tools/send-terminal-input.js'; +import { startTerminalSession } from './tools/start-terminal-session.js'; +import { startWorkspace } from './tools/start-workspace.js'; import { stopAgentTool } from './tools/stop-agent.js'; -import { injectTool } from './tools/inject-tool.js'; +import { stopTerminalSession } from './tools/stop-terminal-session.js'; +import { stopWorkspace } from './tools/stop-workspace.js'; +import type { ServerMode } from './types.js'; const TOOL_ENUM = z.enum([ 'claude-code', @@ -588,5 +590,48 @@ export function createMcpServer(mode: ServerMode = 'orchestration'): McpServer { }, ); + server.tool( + 'send_message', + "Send a message to another agent's inbox. Messages are stored until the recipient reads them.", + { + from: z.string().describe('Sender session_id'), + to: z.string().describe('Recipient session_id'), + body: z.string().describe('Message content'), + thread_id: z + .string() + .optional() + .describe('Thread ID for grouping related messages'), + }, + async ({ from, to, body, thread_id }) => { + try { + const result = sendMessageTool({ from, to, body, thread_id }); + return { content: [{ type: 'text', text: JSON.stringify(result) }] }; + } catch (error) { + return toolError(error); + } + }, + ); + + server.tool( + 'receive_messages', + "Read and consume messages from an agent's inbox. Messages are removed after reading.", + { + session_id: z.string().describe('Whose inbox to read'), + thread_id: z + .string() + .optional() + .describe('Filter by thread — only matching messages are consumed'), + }, + { destructiveHint: true }, + async ({ session_id, thread_id }) => { + try { + const result = receiveMessagesTool({ session_id, thread_id }); + return { content: [{ type: 'text', text: JSON.stringify(result) }] }; + } catch (error) { + return toolError(error); + } + }, + ); + return server; } diff --git a/src/tools/receive-messages.ts b/src/tools/receive-messages.ts new file mode 100644 index 0000000..b6a9098 --- /dev/null +++ b/src/tools/receive-messages.ts @@ -0,0 +1,9 @@ +import type { Message } from '../messaging/store.js'; +import { receiveMessages } from '../messaging/store.js'; + +export function receiveMessagesTool(params: { + session_id: string; + thread_id?: string; +}): { messages: Message[] } { + return receiveMessages(params.session_id, params.thread_id); +} diff --git a/src/tools/send-message.ts b/src/tools/send-message.ts new file mode 100644 index 0000000..aa7f531 --- /dev/null +++ b/src/tools/send-message.ts @@ -0,0 +1,10 @@ +import { sendMessage } from '../messaging/store.js'; + +export function sendMessageTool(params: { + from: string; + to: string; + body: string; + thread_id?: string; +}): { message_id: string; thread_id: string } { + return sendMessage(params.from, params.to, params.body, params.thread_id); +} diff --git a/tests/messaging/store-persistence.test.ts b/tests/messaging/store-persistence.test.ts new file mode 100644 index 0000000..f05c60e --- /dev/null +++ b/tests/messaging/store-persistence.test.ts @@ -0,0 +1,74 @@ +import { existsSync, mkdtempSync, rmSync } from 'node:fs'; +import { join } from 'node:path'; +import { tmpdir } from 'node:os'; +import { + clearAllInboxes, + getUnreadCount, + initStore, + receiveMessages, + sendMessage, +} from '../../src/messaging/store.js'; + +let tmpDir: string; + +beforeEach(() => { + tmpDir = mkdtempSync(join(tmpdir(), 'store-persistence-test-')); + initStore(tmpDir); +}); + +afterEach(() => { + clearAllInboxes(); + rmSync(tmpDir, { recursive: true, force: true }); +}); + +describe('MessageStore persistence', () => { + it('isolation — messages in one test do not leak to the next', () => { + sendMessage('supervisor', 'worker-1', 'test message'); + const result = receiveMessages('worker-1'); + expect(result.messages).toHaveLength(1); + expect(result.messages[0].body).toBe('test message'); + }); + + it('persistence across restart — messages survive re-initialization from same directory', () => { + sendMessage('supervisor', 'worker-1', 'task brief'); + // Simulate restart: call initStore again from SAME directory WITHOUT clearAllInboxes + initStore(tmpDir); + // Messages should be loaded back from disk + expect(getUnreadCount('worker-1')).toBe(1); + const result = receiveMessages('worker-1'); + expect(result.messages).toHaveLength(1); + expect(result.messages[0].body).toBe('task brief'); + expect(result.messages[0].from).toBe('supervisor'); + expect(result.messages[0].to).toBe('worker-1'); + }); + + it('file deletion — clearAllInboxes removes messages.json from disk', () => { + sendMessage('supervisor', 'worker-2', 'hello'); + const messagesFile = join(tmpDir, 'messages.json'); + // File should exist after sendMessage (flush occurred) + expect(existsSync(messagesFile)).toBe(true); + clearAllInboxes(); + // File should be deleted after clearAllInboxes + expect(existsSync(messagesFile)).toBe(false); + // In-memory state cleared too + const result = receiveMessages('worker-2'); + expect(result.messages).toHaveLength(0); + }); + + it('empty start — initStore on fresh directory gives no messages', () => { + const result = receiveMessages('nobody'); + expect(result.messages).toEqual([]); + expect(getUnreadCount('nobody')).toBe(0); + }); + + it('count accuracy — getUnreadCount reflects correct count after reload', () => { + sendMessage('supervisor', 'worker-3', 'message one'); + sendMessage('supervisor', 'worker-3', 'message two'); + expect(getUnreadCount('worker-3')).toBe(2); + // Simulate restart + initStore(tmpDir); + expect(getUnreadCount('worker-3')).toBe(2); + receiveMessages('worker-3'); + expect(getUnreadCount('worker-3')).toBe(0); + }); +}); diff --git a/tests/messaging/store.test.ts b/tests/messaging/store.test.ts new file mode 100644 index 0000000..649d75c --- /dev/null +++ b/tests/messaging/store.test.ts @@ -0,0 +1,95 @@ +import { beforeEach, describe, expect, it } from 'vitest'; +import { + clearAllInboxes, + getUnreadCount, + receiveMessages, + sendMessage, +} from '../../src/messaging/store.js'; + +describe('Message Store', () => { + beforeEach(() => { + clearAllInboxes(); + }); + + it('sendMessage creates a message with UUID id and thread_id', () => { + const result = sendMessage('supervisor-1', 'worker-1', 'hello'); + expect(result.message_id).toMatch( + /^[0-9a-f]{8}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{12}$/, + ); + expect(result.thread_id).toMatch( + /^[0-9a-f]{8}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{12}$/, + ); + }); + + it('sendMessage with explicit thread_id uses it', () => { + const threadId = 'custom-thread-abc'; + const result = sendMessage('supervisor-1', 'worker-1', 'hello', threadId); + expect(result.thread_id).toBe(threadId); + }); + + it('receiveMessages returns and removes messages', () => { + sendMessage('supervisor-1', 'worker-1', 'task-1'); + sendMessage('supervisor-1', 'worker-1', 'task-2'); + + const first = receiveMessages('worker-1'); + expect(first.messages).toHaveLength(2); + expect(first.messages[0].body).toBe('task-1'); + expect(first.messages[1].body).toBe('task-2'); + + const second = receiveMessages('worker-1'); + expect(second.messages).toHaveLength(0); + }); + + it('receiveMessages with thread_id filters correctly', () => { + sendMessage('supervisor-1', 'worker-1', 'msg-thread-A', 'thread-A'); + sendMessage('supervisor-1', 'worker-1', 'msg-thread-B', 'thread-B'); + sendMessage('supervisor-1', 'worker-1', 'msg-thread-A-2', 'thread-A'); + + const filtered = receiveMessages('worker-1', 'thread-A'); + expect(filtered.messages).toHaveLength(2); + expect(filtered.messages[0].body).toBe('msg-thread-A'); + expect(filtered.messages[1].body).toBe('msg-thread-A-2'); + expect(filtered.messages[0].thread_id).toBe('thread-A'); + + const remaining = receiveMessages('worker-1'); + expect(remaining.messages).toHaveLength(1); + expect(remaining.messages[0].body).toBe('msg-thread-B'); + }); + + it('receiveMessages with thread_id cleans up inbox when all messages match', () => { + sendMessage('supervisor-1', 'worker-1', 'msg-1', 'thread-A'); + sendMessage('supervisor-1', 'worker-1', 'msg-2', 'thread-A'); + + const result = receiveMessages('worker-1', 'thread-A'); + expect(result.messages).toHaveLength(2); + expect(getUnreadCount('worker-1')).toBe(0); + }); + + it('receiveMessages returns empty array when no messages', () => { + const result = receiveMessages('nonexistent-session'); + expect(result.messages).toEqual([]); + }); + + it('getUnreadCount returns count without consuming', () => { + sendMessage('supervisor-1', 'worker-1', 'hello'); + sendMessage('supervisor-1', 'worker-1', 'world'); + + expect(getUnreadCount('worker-1')).toBe(2); + expect(getUnreadCount('worker-1')).toBe(2); + + receiveMessages('worker-1'); + expect(getUnreadCount('worker-1')).toBe(0); + }); + + it('multiple senders to same recipient', () => { + sendMessage('supervisor-1', 'worker-1', 'from-sup'); + sendMessage('worker-2', 'worker-1', 'from-peer'); + + const result = receiveMessages('worker-1'); + expect(result.messages).toHaveLength(2); + expect(result.messages[0].from).toBe('supervisor-1'); + expect(result.messages[0].body).toBe('from-sup'); + expect(result.messages[1].from).toBe('worker-2'); + expect(result.messages[1].body).toBe('from-peer'); + }); +}); diff --git a/tests/tools/receive-messages.test.ts b/tests/tools/receive-messages.test.ts new file mode 100644 index 0000000..7527c69 --- /dev/null +++ b/tests/tools/receive-messages.test.ts @@ -0,0 +1,54 @@ +import { mkdtempSync, rmSync } from 'node:fs'; +import { tmpdir } from 'node:os'; +import { join } from 'node:path'; +import { afterEach, beforeEach, describe, expect, it } from 'vitest'; +import { clearAllInboxes, initStore, sendMessage } from '../../src/messaging/store.js'; +import { receiveMessagesTool } from '../../src/tools/receive-messages.js'; + +describe('receiveMessagesTool', () => { + let tmpDir: string; + + beforeEach(() => { + tmpDir = mkdtempSync(join(tmpdir(), 'receive-messages-test-')); + initStore(tmpDir); + }); + + afterEach(() => { + clearAllInboxes(); + rmSync(tmpDir, { recursive: true, force: true }); + }); + + it('returns and consumes messages for a session', () => { + sendMessage('supervisor-1', 'worker-1', 'task A'); + sendMessage('supervisor-1', 'worker-1', 'task B'); + + const result = receiveMessagesTool({ session_id: 'worker-1' }); + expect(result.messages).toHaveLength(2); + expect(result.messages[0].body).toBe('task A'); + expect(result.messages[1].body).toBe('task B'); + + const empty = receiveMessagesTool({ session_id: 'worker-1' }); + expect(empty.messages).toHaveLength(0); + }); + + it('filters by thread_id when provided', () => { + sendMessage('supervisor-1', 'worker-1', 'thread-X msg', 'thread-X'); + sendMessage('supervisor-1', 'worker-1', 'thread-Y msg', 'thread-Y'); + + const result = receiveMessagesTool({ + session_id: 'worker-1', + thread_id: 'thread-X', + }); + expect(result.messages).toHaveLength(1); + expect(result.messages[0].body).toBe('thread-X msg'); + + const remaining = receiveMessagesTool({ session_id: 'worker-1' }); + expect(remaining.messages).toHaveLength(1); + expect(remaining.messages[0].body).toBe('thread-Y msg'); + }); + + it('returns empty array for unknown session', () => { + const result = receiveMessagesTool({ session_id: 'nobody' }); + expect(result.messages).toEqual([]); + }); +}); diff --git a/tests/tools/send-message.test.ts b/tests/tools/send-message.test.ts new file mode 100644 index 0000000..cf06e31 --- /dev/null +++ b/tests/tools/send-message.test.ts @@ -0,0 +1,52 @@ +import { mkdtempSync, rmSync } from 'node:fs'; +import { tmpdir } from 'node:os'; +import { join } from 'node:path'; +import { afterEach, beforeEach, describe, expect, it } from 'vitest'; +import { clearAllInboxes, initStore, receiveMessages } from '../../src/messaging/store.js'; +import { sendMessageTool } from '../../src/tools/send-message.js'; + +describe('sendMessageTool', () => { + let tmpDir: string; + + beforeEach(() => { + tmpDir = mkdtempSync(join(tmpdir(), 'send-message-test-')); + initStore(tmpDir); + }); + + afterEach(() => { + clearAllInboxes(); + rmSync(tmpDir, { recursive: true, force: true }); + }); + + it('sends a message and returns message_id and thread_id', () => { + const result = sendMessageTool({ + from: 'supervisor-1', + to: 'worker-1', + body: 'do the thing', + }); + + expect(result.message_id).toMatch( + /^[0-9a-f]{8}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{12}$/, + ); + expect(result.thread_id).toMatch( + /^[0-9a-f]{8}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{12}$/, + ); + + const inbox = receiveMessages('worker-1'); + expect(inbox.messages).toHaveLength(1); + expect(inbox.messages[0].body).toBe('do the thing'); + expect(inbox.messages[0].from).toBe('supervisor-1'); + expect(inbox.messages[0].to).toBe('worker-1'); + }); + + it('passes through explicit thread_id', () => { + const result = sendMessageTool({ + from: 'supervisor-1', + to: 'worker-1', + body: 'continue', + thread_id: 'existing-thread', + }); + + expect(result.thread_id).toBe('existing-thread'); + }); +});