Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
12 changes: 12 additions & 0 deletions deploy/deployment.yaml
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand All @@ -37,3 +45,7 @@ spec:
port: 8080
initialDelaySeconds: 3
periodSeconds: 5
volumes:
- name: data
persistentVolumeClaim:
claimName: che-mcp-server-data
1 change: 1 addition & 0 deletions deploy/kustomization.yaml
Original file line number Diff line number Diff line change
Expand Up @@ -4,6 +4,7 @@ resources:
- rbac.yaml
- deployment.yaml
- service.yaml
- pvc.yaml

images:
- name: quay.io/che-incubator/che-mcp-server
Expand Down
11 changes: 11 additions & 0 deletions deploy/pvc.yaml
Original file line number Diff line number Diff line change
@@ -0,0 +1,11 @@
apiVersion: v1
kind: PersistentVolumeClaim
metadata:
name: che-mcp-server-data
spec:
accessModes:
- ReadWriteOnce
storageClassName: gp3-csi
resources:
requests:
storage: 1Gi
2 changes: 2 additions & 0 deletions src/config.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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';
2 changes: 1 addition & 1 deletion src/kube/exec.ts
Original file line number Diff line number Diff line change
@@ -1,5 +1,5 @@
import * as k8s from '@kubernetes/client-node';
import stream from 'stream';
import stream from 'node:stream';

import {
getCoreV1Api,
Expand Down
125 changes: 125 additions & 0 deletions src/messaging/store.ts
Original file line number Diff line number Diff line change
@@ -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<string, Message[]> = 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
}
}
}
Comment thread
coderabbitai[bot] marked this conversation as resolved.

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);
}
73 changes: 59 additions & 14 deletions src/tools.ts
Original file line number Diff line number Diff line change
@@ -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',
Expand Down Expand Up @@ -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;
}
9 changes: 9 additions & 0 deletions src/tools/receive-messages.ts
Original file line number Diff line number Diff line change
@@ -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);
}
10 changes: 10 additions & 0 deletions src/tools/send-message.ts
Original file line number Diff line number Diff line change
@@ -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);
}
74 changes: 74 additions & 0 deletions tests/messaging/store-persistence.test.ts
Original file line number Diff line number Diff line change
@@ -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);
});
});
Loading
Loading