feat: add Kimi K3 fleet MCP orchestrator and dashboard
This commit is contained in:
@@ -0,0 +1,4 @@
|
||||
bin/
|
||||
obj/
|
||||
.vs/
|
||||
*.user
|
||||
@@ -0,0 +1,14 @@
|
||||
<!doctype html>
|
||||
<html lang="en"><head><meta charset="utf-8"><meta name="viewport" content="width=device-width,initial-scale=1"><title>KimiFleet</title>
|
||||
<style>
|
||||
:root{color-scheme:dark;--bg:#0b1020;--panel:#121a31;--line:#253252;--text:#edf2ff;--muted:#98a7ca;--accent:#76a7ff;--ok:#50d7a1;--warn:#ffc46b;--bad:#ff7777}*{box-sizing:border-box}body{margin:0;background:radial-gradient(circle at top,#19274b,#0b1020 58%);color:var(--text);font:14px Inter,ui-sans-serif,system-ui}header{padding:28px max(6vw,24px);display:flex;justify-content:space-between;align-items:center;border-bottom:1px solid var(--line)}h1{margin:0;font-size:25px;letter-spacing:-.6px}.sub{color:var(--muted);margin-top:5px}.live{color:var(--ok);font-weight:700}.dot{width:9px;height:9px;display:inline-block;border-radius:50%;background:var(--ok);margin-right:7px;box-shadow:0 0 14px var(--ok)}main{max-width:1280px;margin:0 auto;padding:26px;display:grid;grid-template-columns:1.1fr .9fr;gap:20px}.panel{background:color-mix(in srgb,var(--panel) 92%,transparent);border:1px solid var(--line);border-radius:16px;padding:18px;box-shadow:0 16px 40px #0003}h2{margin:0 0 15px;font-size:15px}.agents{display:grid;gap:12px}.agent{padding:15px;border:1px solid var(--line);border-radius:12px;background:#0d152a}.agent-top{display:flex;justify-content:space-between;align-items:center}.name{font-weight:750;font-size:16px}.badge{padding:4px 9px;border-radius:999px;background:#1a2b4c;color:var(--accent);font-size:12px}.state-ready{color:var(--ok)}.state-working{color:var(--warn)}.state-error{color:var(--bad)}.meta{color:var(--muted);font-size:12px;margin-top:8px}.scopes{margin-top:10px;display:flex;gap:6px;flex-wrap:wrap}.scope{font:11px ui-monospace,monospace;background:#17223d;border-radius:5px;padding:3px 6px;color:#c0cdf1}.events{height:600px;overflow:auto;display:flex;flex-direction:column-reverse;gap:8px}.event{border-left:3px solid var(--accent);background:#0d152a;border-radius:0 8px 8px 0;padding:10px 11px}.event-head{font-size:12px;color:var(--muted);margin-bottom:4px}.event-body{line-height:1.35;word-break:break-word}.empty{color:var(--muted);padding:18px 3px}@media(max-width:850px){main{grid-template-columns:1fr}.events{height:360px}}
|
||||
</style></head><body><header><div><h1>KimiFleet</h1><div class="sub">MCP-controlled Kimi K3 agent fleet</div></div><div class="live"><span class="dot"></span>LIVE</div></header><main><section class="panel"><h2>Agents</h2><div id="agents" class="agents"><div class="empty">Loading agents…</div></div></section><section class="panel"><h2>Activity</h2><div id="events" class="events"><div class="empty">Waiting for events…</div></div></section></main>
|
||||
<script>
|
||||
const agents=document.querySelector('#agents'),events=document.querySelector('#events');
|
||||
const esc=s=>String(s??'').replace(/[&<>"']/g,c=>({'&':'&','<':'<','>':'>','"':'"',"'":'''}[c]));
|
||||
function drawAgents(items){agents.innerHTML=items.length?items.map(a=>`<article class="agent"><div class="agent-top"><div><span class="name">${esc(a.name)}</span> <span class="badge">${esc(a.role)}</span></div><b class="state-${esc(a.state)}">${esc(a.state)}</b></div><div class="meta">${esc(a.workspace)}</div><div class="scopes">${(a.scopes||[]).map(s=>`<span class="scope">${esc(s)}</span>`).join('')||'<span class="meta">No exclusive scopes</span>'}</div></article>`).join(''):'<div class="empty">No agents yet. Create one through the MCP tool or REST API.</div>'}
|
||||
async function refresh(){try{drawAgents(await (await fetch('/agents')).json())}catch{agents.innerHTML='<div class="empty">Dashboard cannot reach KimiFleet.</div>'}}
|
||||
function add(e){const text=pretty(e);if(text===null)return;const time=new Date(e.time).toLocaleTimeString();const row=document.createElement('article');row.className='event';row.innerHTML=`<div class="event-head">${esc(time)} · ${esc(e.agent)} · ${esc(e.kind)}</div><div class="event-body">${esc(text)}</div>`;events.prepend(row);while(events.children.length>140)events.removeChild(events.lastChild)}
|
||||
function pretty(e){if(e.kind!=='session/update')return e.message;try{const u=JSON.parse(e.message).update;if(u.sessionUpdate==='tool_call_update')return `${u.title||'Tool call'} — ${u.status||'in progress'}`;if(u.sessionUpdate==='agent_thought_chunk'||u.sessionUpdate==='agent_message_chunk')return null;return u.sessionUpdate||'session update'}catch{return e.message}}
|
||||
new EventSource('/events').onmessage=x=>{add(JSON.parse(x.data));refresh()};refresh();setInterval(refresh,5000);
|
||||
</script></body></html>
|
||||
@@ -0,0 +1,11 @@
|
||||
<Project Sdk="Microsoft.NET.Sdk">
|
||||
<PropertyGroup>
|
||||
<OutputType>Exe</OutputType>
|
||||
<TargetFramework>net9.0</TargetFramework>
|
||||
<ImplicitUsings>enable</ImplicitUsings>
|
||||
<Nullable>enable</Nullable>
|
||||
</PropertyGroup>
|
||||
<ItemGroup>
|
||||
<None Update="Dashboard.html" CopyToOutputDirectory="PreserveNewest" />
|
||||
</ItemGroup>
|
||||
</Project>
|
||||
@@ -0,0 +1,107 @@
|
||||
using System.Text.Json;
|
||||
|
||||
sealed class McpStdioServer(FleetHost fleet)
|
||||
{
|
||||
private static readonly JsonSerializerOptions Json = new() { PropertyNamingPolicy = JsonNamingPolicy.CamelCase };
|
||||
|
||||
public async Task RunAsync()
|
||||
{
|
||||
string? line;
|
||||
while ((line = await Console.In.ReadLineAsync()) is not null)
|
||||
{
|
||||
try
|
||||
{
|
||||
using var document = JsonDocument.Parse(line);
|
||||
var request = document.RootElement;
|
||||
if (!request.TryGetProperty("method", out var methodNode)) continue;
|
||||
var method = methodNode.GetString();
|
||||
if (!request.TryGetProperty("id", out var id)) continue; // notifications have no response
|
||||
object result = method switch
|
||||
{
|
||||
"initialize" => new
|
||||
{
|
||||
protocolVersion = "2025-03-26",
|
||||
capabilities = new { tools = new { } },
|
||||
serverInfo = new { name = "KimiFleet", version = "0.1.0" },
|
||||
instructions = "Use KimiFleet tools to create scoped Kimi K3 agents, assign work, and require peer review."
|
||||
},
|
||||
"tools/list" => new { tools = Tools },
|
||||
"tools/call" => await CallToolAsync(request.GetProperty("params")),
|
||||
_ => throw new InvalidOperationException($"Unsupported MCP method '{method}'.")
|
||||
};
|
||||
await SendAsync(id, result);
|
||||
}
|
||||
catch (Exception ex)
|
||||
{
|
||||
Console.Error.WriteLine($"[mcp] {ex.Message}");
|
||||
await Console.Out.WriteLineAsync(JsonSerializer.Serialize(new { jsonrpc = "2.0", error = new { code = -32603, message = ex.Message } }, Json));
|
||||
await Console.Out.FlushAsync();
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
private async Task<object> CallToolAsync(JsonElement call)
|
||||
{
|
||||
var name = call.GetProperty("name").GetString() ?? throw new InvalidOperationException("Tool call has no name.");
|
||||
var arguments = call.TryGetProperty("arguments", out var value) ? value : default;
|
||||
object payload = name switch
|
||||
{
|
||||
"fleet_list_agents" => fleet.ListAgents(),
|
||||
"fleet_start_agent" => await fleet.StartAgentAsync(new FleetHost.AgentSpec(
|
||||
Required(arguments, "name"), Required(arguments, "workspace"),
|
||||
Optional(arguments, "role"), StringArray(arguments, "scopes"))),
|
||||
"fleet_assign_task" => Assign(Required(arguments, "agent"), Required(arguments, "prompt")),
|
||||
"fleet_request_review" => Review(Required(arguments, "reviewer"), Required(arguments, "author"), Optional(arguments, "context")),
|
||||
"fleet_stop_agent" => Stop(Required(arguments, "agent")),
|
||||
_ => throw new InvalidOperationException($"Unknown KimiFleet tool '{name}'.")
|
||||
};
|
||||
return new { content = new[] { new { type = "text", text = JsonSerializer.Serialize(payload, Json) } } };
|
||||
}
|
||||
|
||||
private object Assign(string agent, string prompt)
|
||||
{
|
||||
fleet.Assign(agent, prompt);
|
||||
return new { accepted = true, agent };
|
||||
}
|
||||
|
||||
private object Review(string reviewer, string author, string? context)
|
||||
{
|
||||
fleet.Review(reviewer, author, context);
|
||||
return new { accepted = true, reviewer, author };
|
||||
}
|
||||
|
||||
private object Stop(string agent)
|
||||
{
|
||||
fleet.Stop(agent);
|
||||
return new { stopped = agent };
|
||||
}
|
||||
|
||||
private static string Required(JsonElement input, string name) =>
|
||||
input.TryGetProperty(name, out var value) && value.ValueKind == JsonValueKind.String && !string.IsNullOrWhiteSpace(value.GetString())
|
||||
? value.GetString()! : throw new InvalidOperationException($"'{name}' is required.");
|
||||
|
||||
private static string? Optional(JsonElement input, string name) =>
|
||||
input.TryGetProperty(name, out var value) && value.ValueKind == JsonValueKind.String ? value.GetString() : null;
|
||||
|
||||
private static string[]? StringArray(JsonElement input, string name) =>
|
||||
input.TryGetProperty(name, out var values) && values.ValueKind == JsonValueKind.Array
|
||||
? values.EnumerateArray().Where(x => x.ValueKind == JsonValueKind.String).Select(x => x.GetString()!).ToArray() : null;
|
||||
|
||||
private static async Task SendAsync(JsonElement id, object result)
|
||||
{
|
||||
var json = $"{{\"jsonrpc\":\"2.0\",\"id\":{id.GetRawText()},\"result\":{JsonSerializer.Serialize(result, Json)}}}";
|
||||
await Console.Out.WriteLineAsync(json);
|
||||
await Console.Out.FlushAsync();
|
||||
}
|
||||
|
||||
private static readonly object[] Tools =
|
||||
[
|
||||
Tool("fleet_list_agents", "List active Kimi K3 agents, their roles, scope locks and state.", new { type = "object", properties = new { } }),
|
||||
Tool("fleet_start_agent", "Start a yolo Kimi K3 ACP agent with exclusive source scopes.", new { type = "object", required = new[] { "name", "workspace" }, properties = new { name = new { type = "string" }, workspace = new { type = "string" }, role = new { type = "string" }, scopes = new { type = "array", items = new { type = "string" } } } }),
|
||||
Tool("fleet_assign_task", "Assign an autonomous task to an active Kimi agent.", new { type = "object", required = new[] { "agent", "prompt" }, properties = new { agent = new { type = "string" }, prompt = new { type = "string" } } }),
|
||||
Tool("fleet_request_review", "Ask a second agent for read-only peer review of an author's current diff.", new { type = "object", required = new[] { "reviewer", "author" }, properties = new { reviewer = new { type = "string" }, author = new { type = "string" }, context = new { type = "string" } } }),
|
||||
Tool("fleet_stop_agent", "Stop an agent and release its exclusive scopes.", new { type = "object", required = new[] { "agent" }, properties = new { agent = new { type = "string" } } })
|
||||
];
|
||||
|
||||
private static object Tool(string name, string description, object inputSchema) => new { name, description, inputSchema };
|
||||
}
|
||||
+459
@@ -0,0 +1,459 @@
|
||||
using System.Collections.Concurrent;
|
||||
using System.Diagnostics;
|
||||
using System.Net;
|
||||
using System.Text;
|
||||
using System.Text.Json;
|
||||
using System.Text.Json.Nodes;
|
||||
using System.Threading.Channels;
|
||||
|
||||
var options = FleetOptions.Parse(args);
|
||||
var fleet = new FleetHost(options);
|
||||
Console.CancelKeyPress += (_, e) =>
|
||||
{
|
||||
e.Cancel = true;
|
||||
fleet.Dispose();
|
||||
Environment.Exit(0);
|
||||
};
|
||||
|
||||
Console.Error.WriteLine($"KimiFleet dashboard: http://127.0.0.1:{options.Port}");
|
||||
var webTask = fleet.RunAsync();
|
||||
if (options.EnableMcp)
|
||||
await Task.WhenAny(webTask, new McpStdioServer(fleet).RunAsync());
|
||||
else
|
||||
await webTask;
|
||||
|
||||
sealed record FleetOptions(int Port, string KimiExecutable, string Model, bool EnableMcp)
|
||||
{
|
||||
public static FleetOptions Parse(string[] args)
|
||||
{
|
||||
var port = 7373;
|
||||
var kimi = "kimi";
|
||||
var model = "kimi-code/k3";
|
||||
var enableMcp = false;
|
||||
for (var i = 0; i < args.Length; i++)
|
||||
{
|
||||
if (args[i] == "--port" && i + 1 < args.Length && int.TryParse(args[++i], out var value)) port = value;
|
||||
else if (args[i] == "--kimi" && i + 1 < args.Length) kimi = args[++i];
|
||||
else if (args[i] == "--model" && i + 1 < args.Length) model = args[++i];
|
||||
else if (args[i] == "--mcp") enableMcp = true;
|
||||
}
|
||||
return new FleetOptions(port, kimi, model, enableMcp);
|
||||
}
|
||||
}
|
||||
|
||||
sealed class FleetHost : IDisposable
|
||||
{
|
||||
private readonly FleetOptions _options;
|
||||
private readonly HttpListener _listener = new();
|
||||
private readonly ConcurrentDictionary<string, KimiAgent> _agents = new(StringComparer.OrdinalIgnoreCase);
|
||||
private readonly ConcurrentDictionary<string, string> _scopeOwners = new(StringComparer.OrdinalIgnoreCase);
|
||||
private readonly ConcurrentDictionary<Guid, Channel<string>> _eventSubscribers = new();
|
||||
private readonly List<string> _recentEvents = [];
|
||||
private readonly object _eventsLock = new();
|
||||
private bool _disposed;
|
||||
|
||||
public FleetHost(FleetOptions options)
|
||||
{
|
||||
_options = options;
|
||||
_listener.Prefixes.Add($"http://127.0.0.1:{options.Port}/");
|
||||
}
|
||||
|
||||
public async Task RunAsync()
|
||||
{
|
||||
_listener.Start();
|
||||
while (!_disposed)
|
||||
{
|
||||
HttpListenerContext context;
|
||||
try { context = await _listener.GetContextAsync(); }
|
||||
catch (ObjectDisposedException) { break; }
|
||||
_ = Task.Run(() => HandleAsync(context));
|
||||
}
|
||||
}
|
||||
|
||||
private async Task HandleAsync(HttpListenerContext context)
|
||||
{
|
||||
try
|
||||
{
|
||||
var path = context.Request.Url?.AbsolutePath.Trim('/') ?? string.Empty;
|
||||
var parts = path.Split('/', StringSplitOptions.RemoveEmptyEntries);
|
||||
if (context.Request.HttpMethod == "GET" && path.Length == 0)
|
||||
await ServeDashboardAsync(context);
|
||||
else if (context.Request.HttpMethod == "GET" && path == "health")
|
||||
await RespondJsonAsync(context, new { ok = true, model = _options.Model });
|
||||
else if (context.Request.HttpMethod == "GET" && path == "agents")
|
||||
await RespondJsonAsync(context, _agents.Values.Select(a => a.Snapshot()));
|
||||
else if (context.Request.HttpMethod == "GET" && path == "events")
|
||||
await StreamEventsAsync(context);
|
||||
else if (context.Request.HttpMethod == "POST" && path == "agents")
|
||||
await StartAgentAsync(context);
|
||||
else if (context.Request.HttpMethod == "POST" && parts.Length == 3 && parts[0] == "agents" && parts[2] == "prompt")
|
||||
await PromptAgentAsync(context, parts[1]);
|
||||
else if (context.Request.HttpMethod == "POST" && parts.Length == 3 && parts[0] == "agents" && parts[2] == "review")
|
||||
await ReviewAsync(context, parts[1]);
|
||||
else if (context.Request.HttpMethod == "POST" && parts.Length == 3 && parts[0] == "agents" && parts[2] == "stop")
|
||||
await StopAgentAsync(context, parts[1]);
|
||||
else
|
||||
await RespondJsonAsync(context, new { error = "Unknown endpoint." }, 404);
|
||||
}
|
||||
catch (Exception ex)
|
||||
{
|
||||
Publish("fleet", "error", ex.Message);
|
||||
if (context.Response.OutputStream.CanWrite)
|
||||
await RespondJsonAsync(context, new { error = ex.Message }, 500);
|
||||
}
|
||||
}
|
||||
|
||||
private async Task StartAgentAsync(HttpListenerContext context)
|
||||
{
|
||||
var request = await ReadJsonAsync<StartAgentRequest>(context);
|
||||
var agent = await StartAgentAsync(new AgentSpec(request.Name, request.Workspace, request.Role, request.Scopes));
|
||||
await RespondJsonAsync(context, agent, 201);
|
||||
}
|
||||
|
||||
private async Task PromptAgentAsync(HttpListenerContext context, string name)
|
||||
{
|
||||
var request = await ReadJsonAsync<PromptRequest>(context);
|
||||
Assign(name, request.Prompt);
|
||||
await RespondJsonAsync(context, new { accepted = true, agent = name });
|
||||
}
|
||||
|
||||
private async Task ReviewAsync(HttpListenerContext context, string reviewerName)
|
||||
{
|
||||
var request = await ReadJsonAsync<ReviewRequest>(context);
|
||||
Review(reviewerName, request.Author, request.Context);
|
||||
await RespondJsonAsync(context, new { accepted = true, reviewer = reviewerName, author = request.Author });
|
||||
}
|
||||
|
||||
private async Task StopAgentAsync(HttpListenerContext context, string name)
|
||||
{
|
||||
RemoveAgent(name);
|
||||
await RespondJsonAsync(context, new { stopped = name });
|
||||
}
|
||||
|
||||
public async Task<object> StartAgentAsync(AgentSpec spec)
|
||||
{
|
||||
var name = ValidateName(spec.Name);
|
||||
var workspace = Path.GetFullPath(spec.Workspace);
|
||||
if (!Directory.Exists(workspace)) throw new InvalidOperationException($"Workspace does not exist: {workspace}");
|
||||
if (_agents.ContainsKey(name)) throw new InvalidOperationException($"Agent '{name}' already exists.");
|
||||
var scopes = spec.Scopes?.Distinct(StringComparer.OrdinalIgnoreCase).ToArray() ?? [];
|
||||
foreach (var scope in scopes)
|
||||
if (_scopeOwners.TryGetValue(scope, out var owner)) throw new InvalidOperationException($"Scope '{scope}' is owned by '{owner}'.");
|
||||
var agent = new KimiAgent(name, workspace, spec.Role ?? "worker", scopes, _options, Publish);
|
||||
if (!_agents.TryAdd(name, agent)) throw new InvalidOperationException($"Could not add '{name}'.");
|
||||
foreach (var scope in scopes) _scopeOwners[scope] = name;
|
||||
try { await agent.StartAsync(); return agent.Snapshot(); }
|
||||
catch { RemoveAgent(name); throw; }
|
||||
}
|
||||
|
||||
public IReadOnlyList<object> ListAgents() => _agents.Values.Select(a => a.Snapshot()).Cast<object>().ToArray();
|
||||
public void Assign(string name, string prompt) => _ = GetAgent(name).PromptAsync(prompt);
|
||||
public void Review(string reviewerName, string authorName, string? context)
|
||||
{
|
||||
var reviewer = GetAgent(reviewerName);
|
||||
var author = GetAgent(authorName);
|
||||
_ = reviewer.PromptAsync($"Perform peer review for agent '{author.Name}'. Review only; do not edit files. Inspect current git diff and changed files. Check correctness, resource safety, tests, and scope ownership. Report findings by severity with file/line locations. Context: {context ?? "No additional context."}");
|
||||
}
|
||||
public void Stop(string name) => RemoveAgent(name);
|
||||
|
||||
private KimiAgent GetAgent(string name) =>
|
||||
_agents.TryGetValue(name, out var agent) ? agent : throw new InvalidOperationException($"Unknown agent '{name}'.");
|
||||
|
||||
private void RemoveAgent(string name)
|
||||
{
|
||||
if (!_agents.TryRemove(name, out var agent)) return;
|
||||
foreach (var scope in agent.Scopes) _scopeOwners.TryRemove(scope, out _);
|
||||
agent.Dispose();
|
||||
Publish(name, "stopped", "Agent stopped.");
|
||||
}
|
||||
|
||||
private static string ValidateName(string? name)
|
||||
{
|
||||
if (string.IsNullOrWhiteSpace(name) || name.Any(c => !char.IsAsciiLetterOrDigit(c) && c is not '-' and not '_'))
|
||||
throw new InvalidOperationException("Agent name must contain only ASCII letters, digits, '-' or '_'.");
|
||||
return name;
|
||||
}
|
||||
|
||||
private async Task<T> ReadJsonAsync<T>(HttpListenerContext context)
|
||||
{
|
||||
using var reader = new StreamReader(context.Request.InputStream, context.Request.ContentEncoding);
|
||||
var value = await JsonSerializer.DeserializeAsync<T>(reader.BaseStream, JsonOptions);
|
||||
return value ?? throw new InvalidOperationException("Request body must be valid JSON.");
|
||||
}
|
||||
|
||||
private static async Task RespondJsonAsync(HttpListenerContext context, object body, int status = 200)
|
||||
{
|
||||
var bytes = JsonSerializer.SerializeToUtf8Bytes(body, JsonOptions);
|
||||
context.Response.StatusCode = status;
|
||||
context.Response.ContentType = "application/json; charset=utf-8";
|
||||
context.Response.ContentLength64 = bytes.Length;
|
||||
await context.Response.OutputStream.WriteAsync(bytes);
|
||||
context.Response.Close();
|
||||
}
|
||||
|
||||
private static async Task ServeDashboardAsync(HttpListenerContext context)
|
||||
{
|
||||
var html = await File.ReadAllBytesAsync(Path.Combine(AppContext.BaseDirectory, "Dashboard.html"));
|
||||
context.Response.ContentType = "text/html; charset=utf-8";
|
||||
context.Response.ContentLength64 = html.Length;
|
||||
await context.Response.OutputStream.WriteAsync(html);
|
||||
context.Response.Close();
|
||||
}
|
||||
|
||||
private async Task StreamEventsAsync(HttpListenerContext context)
|
||||
{
|
||||
context.Response.StatusCode = 200;
|
||||
context.Response.ContentType = "text/event-stream";
|
||||
context.Response.SendChunked = true;
|
||||
context.Response.Headers.Add("Cache-Control", "no-cache");
|
||||
var id = Guid.NewGuid();
|
||||
var channel = Channel.CreateUnbounded<string>();
|
||||
_eventSubscribers[id] = channel;
|
||||
try
|
||||
{
|
||||
string[] snapshot;
|
||||
lock (_eventsLock) snapshot = _recentEvents.ToArray();
|
||||
foreach (var line in snapshot) await WriteSseAsync(context.Response, line);
|
||||
await foreach (var line in channel.Reader.ReadAllAsync()) await WriteSseAsync(context.Response, line);
|
||||
}
|
||||
catch (HttpListenerException) { }
|
||||
catch (IOException) { }
|
||||
finally { _eventSubscribers.TryRemove(id, out _); context.Response.Close(); }
|
||||
}
|
||||
|
||||
private static async Task WriteSseAsync(HttpListenerResponse response, string json)
|
||||
{
|
||||
var data = Encoding.UTF8.GetBytes($"data: {json}\n\n");
|
||||
await response.OutputStream.WriteAsync(data);
|
||||
await response.OutputStream.FlushAsync();
|
||||
}
|
||||
|
||||
private void Publish(string agent, string kind, string message)
|
||||
{
|
||||
var json = JsonSerializer.Serialize(new { time = DateTimeOffset.UtcNow, agent, kind, message }, JsonOptions);
|
||||
lock (_eventsLock)
|
||||
{
|
||||
_recentEvents.Add(json);
|
||||
if (_recentEvents.Count > 250) _recentEvents.RemoveAt(0);
|
||||
}
|
||||
Console.Error.WriteLine($"[{DateTime.Now:HH:mm:ss}] [{agent}] {kind}: {DescribeEvent(kind, message)}");
|
||||
foreach (var subscriber in _eventSubscribers.Values) subscriber.Writer.TryWrite(json);
|
||||
}
|
||||
|
||||
private static string DescribeEvent(string kind, string message)
|
||||
{
|
||||
if (kind != "session/update") return message;
|
||||
try
|
||||
{
|
||||
using var document = JsonDocument.Parse(message);
|
||||
var update = document.RootElement.GetProperty("update");
|
||||
var type = update.GetProperty("sessionUpdate").GetString();
|
||||
if (type == "tool_call_update") return $"{update.GetProperty("title").GetString()} — {update.GetProperty("status").GetString()}";
|
||||
if (type == "agent_message_chunk") return "Kimi sent a response";
|
||||
return type ?? "session update";
|
||||
}
|
||||
catch (JsonException) { return "session update"; }
|
||||
}
|
||||
|
||||
public void Dispose()
|
||||
{
|
||||
if (_disposed) return;
|
||||
_disposed = true;
|
||||
_listener.Close();
|
||||
foreach (var name in _agents.Keys.ToArray()) RemoveAgent(name);
|
||||
foreach (var subscriber in _eventSubscribers.Values) subscriber.Writer.TryComplete();
|
||||
}
|
||||
|
||||
private static readonly JsonSerializerOptions JsonOptions = new() { PropertyNamingPolicy = JsonNamingPolicy.CamelCase };
|
||||
|
||||
public sealed record AgentSpec(string Name, string Workspace, string? Role, string[]? Scopes);
|
||||
private sealed record StartAgentRequest(string Name, string Workspace, string? Role, string[]? Scopes);
|
||||
private sealed record PromptRequest(string Prompt);
|
||||
private sealed record ReviewRequest(string Author, string? Context);
|
||||
}
|
||||
|
||||
sealed class KimiAgent : IDisposable
|
||||
{
|
||||
private readonly FleetOptions _options;
|
||||
private readonly Action<string, string, string> _publish;
|
||||
private readonly Process _process;
|
||||
private readonly ConcurrentDictionary<long, TaskCompletionSource<JsonElement>> _pending = new();
|
||||
private readonly SemaphoreSlim _stdinLock = new(1, 1);
|
||||
private long _requestId;
|
||||
private string? _sessionId;
|
||||
private bool _disposed;
|
||||
|
||||
public string Name { get; }
|
||||
public string Workspace { get; }
|
||||
public string Role { get; }
|
||||
public IReadOnlyList<string> Scopes { get; }
|
||||
public string State { get; private set; } = "starting";
|
||||
|
||||
public KimiAgent(string name, string workspace, string role, IReadOnlyList<string> scopes, FleetOptions options, Action<string, string, string> publish)
|
||||
{
|
||||
Name = name;
|
||||
Workspace = workspace;
|
||||
Role = role;
|
||||
Scopes = scopes;
|
||||
_options = options;
|
||||
_publish = publish;
|
||||
_process = new Process
|
||||
{
|
||||
StartInfo = new ProcessStartInfo
|
||||
{
|
||||
FileName = options.KimiExecutable,
|
||||
WorkingDirectory = workspace,
|
||||
RedirectStandardInput = true,
|
||||
RedirectStandardOutput = true,
|
||||
RedirectStandardError = true,
|
||||
UseShellExecute = false,
|
||||
CreateNoWindow = true
|
||||
},
|
||||
EnableRaisingEvents = true
|
||||
};
|
||||
_process.StartInfo.ArgumentList.Add("--yolo");
|
||||
_process.StartInfo.ArgumentList.Add("--model");
|
||||
_process.StartInfo.ArgumentList.Add(options.Model);
|
||||
_process.StartInfo.ArgumentList.Add("acp");
|
||||
}
|
||||
|
||||
public async Task StartAsync()
|
||||
{
|
||||
if (!_process.Start()) throw new InvalidOperationException($"Could not start Kimi process for '{Name}'.");
|
||||
_process.Exited += (_, _) => { State = "stopped"; _publish(Name, "exit", $"Kimi exited with {_process.ExitCode}."); };
|
||||
_ = ReadOutputAsync();
|
||||
_ = ReadErrorAsync();
|
||||
var initialized = await RequestAsync("initialize", new
|
||||
{
|
||||
protocolVersion = 1,
|
||||
clientCapabilities = new { fs = new { readTextFile = true, writeTextFile = true } },
|
||||
clientInfo = new { name = "KimiFleet", version = "0.1.0" }
|
||||
});
|
||||
_publish(Name, "initialized", initialized.GetProperty("agentInfo").GetProperty("name").GetString() ?? "Kimi Code CLI");
|
||||
var session = await RequestAsync("session/new", new { cwd = Workspace, mcpServers = Array.Empty<object>() });
|
||||
_sessionId = session.GetProperty("sessionId").GetString() ?? throw new InvalidOperationException("ACP did not return sessionId.");
|
||||
State = "ready";
|
||||
_publish(Name, "ready", $"Role={Role}; scopes={string.Join(", ", Scopes)}");
|
||||
}
|
||||
|
||||
public async Task PromptAsync(string prompt)
|
||||
{
|
||||
if (string.IsNullOrWhiteSpace(prompt)) throw new InvalidOperationException("Prompt must not be empty.");
|
||||
if (_sessionId is null) throw new InvalidOperationException($"Agent '{Name}' has no ACP session.");
|
||||
State = "working";
|
||||
_publish(Name, "prompt", prompt.Length > 160 ? prompt[..160] + "…" : prompt);
|
||||
try
|
||||
{
|
||||
await RequestAsync("session/prompt", new { sessionId = _sessionId, prompt = new[] { new { type = "text", text = prompt } } });
|
||||
State = "ready";
|
||||
_publish(Name, "completed", "Prompt completed.");
|
||||
}
|
||||
catch (Exception ex)
|
||||
{
|
||||
State = "error";
|
||||
_publish(Name, "error", ex.Message);
|
||||
}
|
||||
}
|
||||
|
||||
private async Task<JsonElement> RequestAsync(string method, object parameters)
|
||||
{
|
||||
var id = Interlocked.Increment(ref _requestId);
|
||||
var completion = new TaskCompletionSource<JsonElement>(TaskCreationOptions.RunContinuationsAsynchronously);
|
||||
if (!_pending.TryAdd(id, completion)) throw new InvalidOperationException("Could not register ACP request.");
|
||||
var envelope = JsonSerializer.Serialize(new { jsonrpc = "2.0", id, method, @params = parameters });
|
||||
await SendRawAsync(envelope);
|
||||
using var timeout = new CancellationTokenSource(TimeSpan.FromMinutes(10));
|
||||
await using var registration = timeout.Token.Register(() => completion.TrySetException(new TimeoutException($"ACP request '{method}' timed out.")));
|
||||
return await completion.Task;
|
||||
}
|
||||
|
||||
private async Task ReadOutputAsync()
|
||||
{
|
||||
while (!_process.HasExited)
|
||||
{
|
||||
var line = await _process.StandardOutput.ReadLineAsync();
|
||||
if (line is null) break;
|
||||
HandleMessage(line);
|
||||
}
|
||||
}
|
||||
|
||||
private async Task ReadErrorAsync()
|
||||
{
|
||||
while (!_process.HasExited)
|
||||
{
|
||||
var line = await _process.StandardError.ReadLineAsync();
|
||||
if (line is null) break;
|
||||
_publish(Name, "stderr", line);
|
||||
}
|
||||
}
|
||||
|
||||
private void HandleMessage(string line)
|
||||
{
|
||||
try
|
||||
{
|
||||
using var document = JsonDocument.Parse(line);
|
||||
var root = document.RootElement;
|
||||
if (root.TryGetProperty("id", out var idNode) && idNode.TryGetInt64(out var id) && _pending.TryRemove(id, out var completion))
|
||||
{
|
||||
if (root.TryGetProperty("error", out var error)) completion.TrySetException(new InvalidOperationException(error.GetRawText()));
|
||||
else if (root.TryGetProperty("result", out var result)) completion.TrySetResult(result.Clone());
|
||||
else completion.TrySetException(new InvalidOperationException("ACP response has neither result nor error."));
|
||||
return;
|
||||
}
|
||||
if (root.TryGetProperty("method", out var method))
|
||||
{
|
||||
var methodName = method.GetString() ?? "notification";
|
||||
if (methodName == "session/request_permission" && root.TryGetProperty("id", out var permissionId) && root.TryGetProperty("params", out var permission))
|
||||
{
|
||||
var optionId = SelectApproval(permission);
|
||||
_ = RespondToPermissionAsync(permissionId.GetRawText(), optionId);
|
||||
_publish(Name, "permission", $"Auto-approved '{optionId}' (yolo mode).");
|
||||
return;
|
||||
}
|
||||
_publish(Name, methodName, root.TryGetProperty("params", out var p) ? p.GetRawText() : string.Empty);
|
||||
}
|
||||
else _publish(Name, "protocol", line);
|
||||
}
|
||||
catch (JsonException) { _publish(Name, "stdout", line); }
|
||||
}
|
||||
|
||||
private static string SelectApproval(JsonElement permission)
|
||||
{
|
||||
if (!permission.TryGetProperty("options", out var options) || options.ValueKind != JsonValueKind.Array)
|
||||
throw new InvalidOperationException("ACP permission request has no options.");
|
||||
var ids = options.EnumerateArray()
|
||||
.Where(option => option.TryGetProperty("optionId", out _))
|
||||
.Select(option => option.GetProperty("optionId").GetString())
|
||||
.Where(id => !string.IsNullOrWhiteSpace(id))
|
||||
.Cast<string>()
|
||||
.ToArray();
|
||||
return ids.FirstOrDefault(id => id.Contains("approve_always", StringComparison.OrdinalIgnoreCase))
|
||||
?? ids.FirstOrDefault(id => id.Contains("approve", StringComparison.OrdinalIgnoreCase))
|
||||
?? throw new InvalidOperationException("ACP permission request has no approval option.");
|
||||
}
|
||||
|
||||
private Task RespondToPermissionAsync(string requestId, string optionId) =>
|
||||
SendRawAsync($"{{\"jsonrpc\":\"2.0\",\"id\":{requestId},\"result\":{{\"outcome\":{{\"outcome\":\"selected\",\"optionId\":{JsonSerializer.Serialize(optionId)}}}}}}}");
|
||||
|
||||
private async Task SendRawAsync(string json)
|
||||
{
|
||||
await _stdinLock.WaitAsync();
|
||||
try
|
||||
{
|
||||
await _process.StandardInput.WriteLineAsync(json);
|
||||
await _process.StandardInput.FlushAsync();
|
||||
}
|
||||
finally { _stdinLock.Release(); }
|
||||
}
|
||||
|
||||
public object Snapshot() => new { name = Name, workspace = Workspace, role = Role, scopes = Scopes, state = State, sessionId = _sessionId };
|
||||
|
||||
public void Dispose()
|
||||
{
|
||||
if (_disposed) return;
|
||||
_disposed = true;
|
||||
try { if (!_process.HasExited) _process.Kill(entireProcessTree: true); } catch (InvalidOperationException) { }
|
||||
_stdinLock.Dispose();
|
||||
_process.Dispose();
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,57 @@
|
||||
# KimiFleet
|
||||
|
||||
Visible multi-agent orchestrator and MCP server for Kimi Code K3. It controls Kimi through ACP (JSON-RPC over stdio), rather than attempting to drive its terminal UI.
|
||||
|
||||
## Start
|
||||
|
||||
```bash
|
||||
dotnet run --project KimiFleet.csproj -- --port 7373
|
||||
```
|
||||
|
||||
Open a second terminal and watch live activity:
|
||||
|
||||
```bash
|
||||
curl -N http://127.0.0.1:7373/events
|
||||
```
|
||||
|
||||
Open `http://127.0.0.1:7373` for the dashboard: it shows agents, their scoped ownership, state and a human-readable live activity feed.
|
||||
|
||||
## MCP mode
|
||||
|
||||
Start both the dashboard and an MCP stdio server:
|
||||
|
||||
```bash
|
||||
dotnet run --project KimiFleet.csproj -- --mcp --port 7373
|
||||
```
|
||||
|
||||
The server implements `fleet_list_agents`, `fleet_start_agent`, `fleet_assign_task`, `fleet_request_review`, and `fleet_stop_agent`. In MCP mode stdout is protocol-only; operational logs are written to stderr and remain visible in the dashboard.
|
||||
|
||||
## API
|
||||
|
||||
Create an agent with exclusive source scopes:
|
||||
|
||||
```bash
|
||||
curl -X POST http://127.0.0.1:7373/agents \
|
||||
-H 'Content-Type: application/json' \
|
||||
-d '{"name":"vulkan","workspace":"/home/emil/Desktop/Cortex_Engine","role":"renderer","scopes":["src/Engine.Graphics.Vulkan","src/Engine.Graphics"]}'
|
||||
```
|
||||
|
||||
Give it work:
|
||||
|
||||
```bash
|
||||
curl -X POST http://127.0.0.1:7373/agents/vulkan/prompt \
|
||||
-H 'Content-Type: application/json' \
|
||||
-d '{"prompt":"Read AGENTS.md, inspect the Vulkan backend, and propose one safe improvement."}'
|
||||
```
|
||||
|
||||
Ask another agent to review the current diff without editing:
|
||||
|
||||
```bash
|
||||
curl -X POST http://127.0.0.1:7373/agents/reviewer/review \
|
||||
-H 'Content-Type: application/json' \
|
||||
-d '{"author":"vulkan","context":"Review the Vulkan change after its tests finish."}'
|
||||
```
|
||||
|
||||
Useful endpoints: `GET /health`, `GET /agents`, `GET /events`, `POST /agents/{name}/stop`.
|
||||
|
||||
All ACP children launch as `kimi --yolo --model kimi-code/k3 acp`. Scope locks prevent two fleet agents from being assigned the same area, but they do not replace a final `git status` check before edits.
|
||||
Reference in New Issue
Block a user