add a clide-hosted team coordination broker over the MCP control channel
Claude's tmux team mode let teammates message each other and share a task
list; that mode is undocumented and unavailable headless. clide rebuilds
the same behavior over its own managed sessions, as the broker.
Verified live against claude 2.1.150 that a spawner can host an in-process
("SDK") MCP server entirely over the stream-json control channel — no
subprocess, no --mcp-config, no socket: declare the server name in the
initialize handshake's sdkMcpServers, answer the mcp_message JSON-RPC
round-trips (initialize / tools/list / tools/call) under
response.response.mcp_response. SDK tool calls are permission-gated through
the existing can_use_tool path. Documented in the 2.1.150 spike §6.
StreamJsonSession gains an McpServer hosting seam; TeamBroker + TeamMcpServer
expose send_message / broadcast / list_teammates / inbox / claim_task /
task_status, all routed through one shared broker. The orchestrator owns the
broker, registers each team session, delivers a message into the target's
next turn on its stdin, and injects roster + role via --append-system-prompt.
Solo sessions are unchanged (no MCP server, no initialize handshake).
T-170, D-77.
Co-Authored-By: Claude Opus 4.7 (1M context) <noreply@anthropic.com>
This commit is contained in:
@@ -6,11 +6,12 @@ import 'package:flutter_test/flutter_test.dart';
|
||||
|
||||
class _FakeProc implements StreamJsonProcess {
|
||||
final _ctl = StreamController<String>.broadcast();
|
||||
final List<String> writes = [];
|
||||
bool killed = false;
|
||||
@override
|
||||
Stream<String> get lines => _ctl.stream;
|
||||
@override
|
||||
void writeLine(String line) {}
|
||||
void writeLine(String line) => writes.add(line);
|
||||
@override
|
||||
Future<void> kill() async => killed = true;
|
||||
}
|
||||
@@ -89,4 +90,36 @@ void main() {
|
||||
await Future<void>.delayed(Duration.zero);
|
||||
expect(created.every((p) => p.killed), isTrue);
|
||||
});
|
||||
|
||||
group('team broker wiring (T-170)', () {
|
||||
SpawnSpec teamSpec(String id, String name, String role) => SpawnSpec(id: id, role: role, sessionId: '$id-uuid', cwd: '/repo', team: true, memberName: name);
|
||||
|
||||
test('team sessions register in the broker; solo sessions do not', () async {
|
||||
await orch.spawn(spec('solo'));
|
||||
expect(orch.broker.members, isEmpty);
|
||||
await orch.spawn(teamSpec('primary', 'lead', 'lead'));
|
||||
expect(orch.broker.members.map((m) => m.name), ['lead']);
|
||||
});
|
||||
|
||||
test('a message between team members is delivered into the target session stdin', () async {
|
||||
await orch.spawn(teamSpec('primary', 'lead', 'lead'));
|
||||
await orch.spawn(teamSpec('teammate:tyre', 'tyre', 'teammate'));
|
||||
orch.broker.sendMessage('primary', 'tyre', 'pick up T-9');
|
||||
await Future<void>.delayed(Duration.zero);
|
||||
final tyreProc = created[1];
|
||||
expect(tyreProc.writes.any((w) => w.contains('[team] lead: pick up T-9')), isTrue);
|
||||
});
|
||||
|
||||
test('a team session declares the clide-team MCP server in its init handshake', () async {
|
||||
await orch.spawn(teamSpec('primary', 'lead', 'lead'));
|
||||
expect(created.single.writes.any((w) => w.contains('"sdkMcpServers":["clide-team"]')), isTrue);
|
||||
});
|
||||
|
||||
test('closing a team member removes it from the broker roster', () async {
|
||||
await orch.spawn(teamSpec('primary', 'lead', 'lead'));
|
||||
await orch.spawn(teamSpec('teammate:tyre', 'tyre', 'teammate'));
|
||||
await orch.close('teammate:tyre');
|
||||
expect(orch.broker.members.map((m) => m.name), ['lead']);
|
||||
});
|
||||
});
|
||||
}
|
||||
|
||||
@@ -20,6 +20,38 @@ class _FakeProc implements StreamJsonProcess {
|
||||
void emit(String line) => _ctl.add(line);
|
||||
}
|
||||
|
||||
class _FakeMcpServer implements McpServer {
|
||||
@override
|
||||
String get name => 'clide-team';
|
||||
@override
|
||||
String get version => '9.9.9';
|
||||
final List<String> calls = [];
|
||||
@override
|
||||
List<Map<String, dynamic>> get tools => [
|
||||
{
|
||||
'name': 'ping',
|
||||
'description': 'p',
|
||||
'inputSchema': {'type': 'object', 'properties': <String, dynamic>{}},
|
||||
},
|
||||
];
|
||||
@override
|
||||
Future<Map<String, dynamic>> callTool(String name, Map<String, dynamic> arguments) async {
|
||||
calls.add(name);
|
||||
return {
|
||||
'content': [
|
||||
{'type': 'text', 'text': 'pong'},
|
||||
],
|
||||
'isError': false,
|
||||
};
|
||||
}
|
||||
}
|
||||
|
||||
String mcpMessage(String rid, Map<String, dynamic> message, {String server = 'clide-team'}) => jsonEncode({
|
||||
'type': 'control_request',
|
||||
'request_id': rid,
|
||||
'request': {'subtype': 'mcp_message', 'server_name': server, 'message': message},
|
||||
});
|
||||
|
||||
String assistantText(String text) => jsonEncode({
|
||||
'type': 'assistant',
|
||||
'uuid': 'a1',
|
||||
@@ -361,4 +393,72 @@ void main() {
|
||||
await session.dispose();
|
||||
expect(proc.killed, isTrue);
|
||||
});
|
||||
|
||||
group('MCP server hosting (T-170)', () {
|
||||
late _FakeProc mproc;
|
||||
late StreamJsonSession msession;
|
||||
late _FakeMcpServer server;
|
||||
|
||||
setUp(() {
|
||||
mproc = _FakeProc();
|
||||
server = _FakeMcpServer();
|
||||
msession = StreamJsonSession(mproc, mcpServers: [server]);
|
||||
msession.start();
|
||||
});
|
||||
|
||||
tearDown(() => msession.dispose());
|
||||
|
||||
Map<String, dynamic> mcpResponseOf(String write) {
|
||||
final resp = jsonDecode(write) as Map<String, dynamic>;
|
||||
return ((resp['response'] as Map)['response'] as Map)['mcp_response'] as Map<String, dynamic>;
|
||||
}
|
||||
|
||||
test('declares its sdkMcpServers in the initialize handshake', () {
|
||||
final init = mproc.writes.map((w) => jsonDecode(w) as Map).firstWhere(
|
||||
(m) => (m['request'] as Map?)?['subtype'] == 'initialize',
|
||||
);
|
||||
expect((init['request'] as Map)['sdkMcpServers'], ['clide-team']);
|
||||
});
|
||||
|
||||
test('answers mcp initialize with our serverInfo', () async {
|
||||
mproc.emit(mcpMessage('m1', {
|
||||
'method': 'initialize',
|
||||
'params': {'protocolVersion': '2025-11-25'},
|
||||
'jsonrpc': '2.0',
|
||||
'id': 0,
|
||||
}));
|
||||
await Future<void>.delayed(Duration.zero);
|
||||
final r = mcpResponseOf(mproc.writes.last);
|
||||
expect((r['result'] as Map)['serverInfo'], {'name': 'clide-team', 'version': '9.9.9'});
|
||||
});
|
||||
|
||||
test('answers tools/list with the server tools', () async {
|
||||
mproc.emit(mcpMessage('m2', {'method': 'tools/list', 'jsonrpc': '2.0', 'id': 1}));
|
||||
await Future<void>.delayed(Duration.zero);
|
||||
final r = mcpResponseOf(mproc.writes.last);
|
||||
final tools = (r['result'] as Map)['tools'] as List;
|
||||
expect(tools.single['name'], 'ping');
|
||||
});
|
||||
|
||||
test('routes tools/call to the server and returns its result', () async {
|
||||
mproc.emit(mcpMessage('m3', {
|
||||
'method': 'tools/call',
|
||||
'params': {'name': 'ping', 'arguments': <String, dynamic>{}},
|
||||
'jsonrpc': '2.0',
|
||||
'id': 2,
|
||||
}));
|
||||
await Future<void>.delayed(Duration.zero);
|
||||
expect(server.calls, ['ping']);
|
||||
final r = mcpResponseOf(mproc.writes.last);
|
||||
final content = (r['result'] as Map)['content'] as List;
|
||||
expect(content.single['text'], 'pong');
|
||||
});
|
||||
|
||||
test('an mcp_message for an unknown server is answered with an error', () async {
|
||||
mproc.emit(mcpMessage('m4', {'method': 'tools/list', 'jsonrpc': '2.0', 'id': 3}, server: 'nope'));
|
||||
await Future<void>.delayed(Duration.zero);
|
||||
final r = mcpResponseOf(mproc.writes.last);
|
||||
expect(r['error'], isNotNull);
|
||||
});
|
||||
});
|
||||
}
|
||||
|
||||
@@ -0,0 +1,128 @@
|
||||
import 'dart:convert';
|
||||
|
||||
import 'package:clide/builtin/claude/src/team_broker.dart';
|
||||
import 'package:test/test.dart';
|
||||
|
||||
/// Decode a [TeamMcpServer] tool result's single text block back into the
|
||||
/// structured map the broker returned.
|
||||
Map<String, dynamic> decode(Map<String, dynamic> mcpResult) {
|
||||
final text = ((mcpResult['content'] as List).single as Map)['text'] as String;
|
||||
return jsonDecode(text) as Map<String, dynamic>;
|
||||
}
|
||||
|
||||
void main() {
|
||||
late TeamBroker broker;
|
||||
late List<(String, String)> delivered; // (toMemberId, text)
|
||||
late TeamMcpServer lead;
|
||||
late TeamMcpServer tyre;
|
||||
|
||||
setUp(() {
|
||||
delivered = [];
|
||||
broker = TeamBroker(deliver: (to, text) => delivered.add((to, text)));
|
||||
broker.addMember(const TeamMemberRef(id: 'primary', name: 'lead', role: 'lead'));
|
||||
broker.addMember(const TeamMemberRef(id: 'teammate:tyre', name: 'tyre', role: 'teammate'));
|
||||
lead = TeamMcpServer(broker: broker, memberId: 'primary');
|
||||
tyre = TeamMcpServer(broker: broker, memberId: 'teammate:tyre');
|
||||
});
|
||||
|
||||
test('send_message delivers into the named teammate next turn', () async {
|
||||
final r = decode(await lead.callTool('send_message', {'to': 'tyre', 'text': 'pick up T-9'}));
|
||||
expect(r['ok'], isTrue);
|
||||
expect(r['to'], 'tyre');
|
||||
expect(delivered, [('teammate:tyre', '[team] lead: pick up T-9')]);
|
||||
});
|
||||
|
||||
test('addressing an unknown teammate fails with a helpful error', () async {
|
||||
final result = await lead.callTool('send_message', {'to': 'ghost', 'text': 'hi'});
|
||||
expect(result['isError'], isTrue);
|
||||
expect(decode(result)['ok'], isFalse);
|
||||
expect(delivered, isEmpty);
|
||||
});
|
||||
|
||||
test('the recipient can read the message from its inbox', () async {
|
||||
await lead.callTool('send_message', {'to': 'tyre', 'text': 'hello'});
|
||||
final box = decode(await tyre.callTool('inbox', {}));
|
||||
final msgs = box['messages'] as List;
|
||||
expect(msgs.single['from'], 'lead');
|
||||
expect(msgs.single['text'], 'hello');
|
||||
// Draining: a second read is empty.
|
||||
expect((decode(await tyre.callTool('inbox', {}))['messages'] as List), isEmpty);
|
||||
});
|
||||
|
||||
test('broadcast reaches every other member but not the sender', () async {
|
||||
broker.addMember(const TeamMemberRef(id: 'teammate:qatux', name: 'qatux', role: 'teammate'));
|
||||
final r = decode(await lead.callTool('broadcast', {'text': 'standup'}));
|
||||
expect((r['recipients'] as List).toSet(), {'tyre', 'qatux'});
|
||||
expect(delivered.map((d) => d.$1).toSet(), {'teammate:tyre', 'teammate:qatux'});
|
||||
});
|
||||
|
||||
test('list_teammates returns the other members with roles', () async {
|
||||
final r = decode(await lead.callTool('list_teammates', {}));
|
||||
final mates = r['teammates'] as List;
|
||||
expect(mates.single, {'name': 'tyre', 'role': 'teammate'});
|
||||
});
|
||||
|
||||
test('a claimed task is visible to every member as shared state', () async {
|
||||
final claimed = decode(await tyre.callTool('claim_task', {'title': 'wire the broker'}));
|
||||
final taskId = (claimed['task'] as Map)['id'] as String;
|
||||
expect((claimed['task'] as Map)['owner'], 'tyre');
|
||||
|
||||
final seenByLead = decode(await lead.callTool('task_status', {}));
|
||||
final tasks = seenByLead['tasks'] as List;
|
||||
expect(tasks.single['id'], taskId);
|
||||
expect(tasks.single['status'], 'claimed');
|
||||
|
||||
final done = decode(await lead.callTool('task_status', {'id': taskId, 'status': 'done'}));
|
||||
expect((done['task'] as Map)['status'], 'done');
|
||||
});
|
||||
|
||||
test('removing a member releases its claimed tasks', () async {
|
||||
final claimed = decode(await tyre.callTool('claim_task', {'title': 'temp'}));
|
||||
final taskId = (claimed['task'] as Map)['id'] as String;
|
||||
broker.removeMember('teammate:tyre');
|
||||
final tasks = decode(await lead.callTool('task_status', {}))['tasks'] as List;
|
||||
final t = tasks.firstWhere((t) => t['id'] == taskId);
|
||||
expect(t['status'], 'open');
|
||||
expect(t.containsKey('owner'), isFalse);
|
||||
});
|
||||
|
||||
test('claim_task by id claims an existing open task', () async {
|
||||
final created = decode(await lead.callTool('task_status', {'title': 'open work'}));
|
||||
final id = (created['task'] as Map)['id'] as String;
|
||||
expect((created['task'] as Map)['status'], 'open');
|
||||
|
||||
final claimed = decode(await tyre.callTool('claim_task', {'id': id}));
|
||||
expect((claimed['task'] as Map)['status'], 'claimed');
|
||||
expect((claimed['task'] as Map)['owner'], 'tyre');
|
||||
});
|
||||
|
||||
test('claim_task with neither id nor title is an error', () async {
|
||||
final r = await lead.callTool('claim_task', {});
|
||||
expect(r['isError'], isTrue);
|
||||
});
|
||||
|
||||
test('claiming an unknown task id is an error', () async {
|
||||
final r = await lead.callTool('claim_task', {'id': 'task-999'});
|
||||
expect(r['isError'], isTrue);
|
||||
});
|
||||
|
||||
test('task_status on an unknown id is an error', () async {
|
||||
final r = await lead.callTool('task_status', {'id': 'task-999', 'status': 'done'});
|
||||
expect(r['isError'], isTrue);
|
||||
});
|
||||
|
||||
test('an unknown team tool is an error', () async {
|
||||
final r = await lead.callTool('nope', {});
|
||||
expect(r['isError'], isTrue);
|
||||
});
|
||||
|
||||
test('removing an unknown member is a no-op', () {
|
||||
broker.removeMember('teammate:ghost');
|
||||
expect(broker.members.map((m) => m.name).toSet(), {'lead', 'tyre'});
|
||||
});
|
||||
|
||||
test('the MCP tool surface lists all six team tools', () {
|
||||
final names = lead.tools.map((t) => t['name']).toSet();
|
||||
expect(names, {'send_message', 'broadcast', 'list_teammates', 'inbox', 'claim_task', 'task_status'});
|
||||
});
|
||||
}
|
||||
Reference in New Issue
Block a user