2026-07-27 22:02:34 -07:00
|
|
|
using System.Reflection;
|
|
|
|
|
using System.Text.Json;
|
|
|
|
|
|
|
|
|
|
namespace NucleicBroker;
|
|
|
|
|
|
|
|
|
|
/// <summary>
|
|
|
|
|
/// The broker's request dispatcher (docs/WINDOWS_PORT.md §3.3): one instance per process,
|
|
|
|
|
/// fed NDJSON lines from hostd, answering through an <see cref="OutboundWriter"/>. Requests
|
|
|
|
|
/// run concurrently (an image pull must not block a container.list), responses/notifications
|
|
|
|
|
/// serialize through the writer's single drain loop. Also the <see cref="IBrokerEvents"/>
|
|
|
|
|
/// sink the facade raises async events into.
|
|
|
|
|
/// </summary>
|
|
|
|
|
public sealed class BrokerService : IBrokerEvents
|
|
|
|
|
{
|
|
|
|
|
private readonly IWslc wslc;
|
|
|
|
|
private readonly IAiProvider ai;
|
|
|
|
|
private readonly OutboundWriter outbound;
|
|
|
|
|
private readonly Dictionary<long, IWslcProcess> procs = new();
|
|
|
|
|
private readonly Lock procsLock = new();
|
|
|
|
|
private long nextProcId;
|
|
|
|
|
|
|
|
|
|
public BrokerService(IWslc wslc, IAiProvider ai, OutboundWriter outbound)
|
|
|
|
|
{
|
|
|
|
|
this.wslc = wslc;
|
|
|
|
|
this.ai = ai;
|
|
|
|
|
this.outbound = outbound;
|
|
|
|
|
wslc.SetEvents(this);
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
/// <summary>Broker protocol revision, bumped with any wire change; hostd degrades from
|
|
|
|
|
/// the `hello` capabilities rather than parsing versions (docs/WINDOWS_PORT.md §2.3).</summary>
|
|
|
|
|
public const int ProtocolVersion = 1;
|
|
|
|
|
|
|
|
|
|
/// <summary>Handle one inbound NDJSON line. The returned task completes once the
|
|
|
|
|
/// response has been ENQUEUED (tests await it; the stdin loop fires and forgets).</summary>
|
|
|
|
|
public async Task HandleLineAsync(string line, CancellationToken ct = default)
|
|
|
|
|
{
|
|
|
|
|
JsonElement? id = null;
|
|
|
|
|
string method;
|
|
|
|
|
JsonElement? @params = null;
|
|
|
|
|
try
|
|
|
|
|
{
|
|
|
|
|
using var doc = JsonDocument.Parse(line);
|
|
|
|
|
var root = doc.RootElement;
|
|
|
|
|
// Clone: the id/params outlive this JsonDocument (responses are built after
|
|
|
|
|
// arbitrary awaits).
|
|
|
|
|
if (root.TryGetProperty("id", out var rawId)) id = rawId.Clone();
|
|
|
|
|
if (root.TryGetProperty("params", out var rawParams)) @params = rawParams.Clone();
|
|
|
|
|
method = root.TryGetProperty("method", out var m) && m.ValueKind == JsonValueKind.String
|
|
|
|
|
? m.GetString()!
|
|
|
|
|
: throw new JsonException("missing method");
|
|
|
|
|
}
|
|
|
|
|
catch (JsonException e)
|
|
|
|
|
{
|
|
|
|
|
outbound.EnqueueJson(Rpc.Error(id, Rpc.ParseError, $"unparseable request: {e.Message}"));
|
|
|
|
|
return;
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
try
|
|
|
|
|
{
|
2026-07-29 17:09:35 -07:00
|
|
|
var result = await DispatchAsync(method, @params, id, ct).ConfigureAwait(false);
|
|
|
|
|
// A null result means the handler already enqueued its own response because it had to
|
|
|
|
|
// send it before doing something else (proc.exec — see below). Everything else
|
|
|
|
|
// returns a value and is responded to here.
|
|
|
|
|
if (result is not null && id is { } requestId)
|
2026-07-27 22:02:34 -07:00
|
|
|
outbound.EnqueueJson(Rpc.Response(requestId, result));
|
|
|
|
|
}
|
|
|
|
|
catch (JsonException e)
|
|
|
|
|
{
|
|
|
|
|
outbound.EnqueueJson(Rpc.Error(id, Rpc.InvalidParams, e.Message));
|
|
|
|
|
}
|
|
|
|
|
catch (WslcError e)
|
|
|
|
|
{
|
|
|
|
|
outbound.EnqueueJson(Rpc.Error(id, Rpc.FacadeError, e.Message, new { kind = e.Kind }));
|
|
|
|
|
}
|
|
|
|
|
catch (NotSupportedException)
|
|
|
|
|
{
|
|
|
|
|
outbound.EnqueueJson(Rpc.Error(id, Rpc.MethodNotFound, $"unknown method: {method}"));
|
|
|
|
|
}
|
|
|
|
|
catch (Exception e)
|
|
|
|
|
{
|
|
|
|
|
outbound.EnqueueJson(Rpc.Error(id, Rpc.InternalError, e.Message));
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
|
2026-07-29 17:09:35 -07:00
|
|
|
/// <summary>Returns the result to respond with, or null when the handler has already
|
|
|
|
|
/// responded (it needed the response on the wire before continuing).</summary>
|
|
|
|
|
private async Task<object?> DispatchAsync(
|
|
|
|
|
string method, JsonElement? p, JsonElement? id, CancellationToken ct)
|
2026-07-27 22:02:34 -07:00
|
|
|
{
|
|
|
|
|
switch (method)
|
|
|
|
|
{
|
|
|
|
|
case "hello":
|
|
|
|
|
{
|
2026-07-29 02:52:41 -07:00
|
|
|
// Two kinds of capability in one list: the RPC families this broker serves, and
|
|
|
|
|
// what the facade underneath can actually do (docs/WINDOWS_PORT.md D13). hostd
|
|
|
|
|
// needs both — "proc" says the methods exist, "tty" says proc.exec(tty:true)
|
|
|
|
|
// will work rather than failing `unsupported` at the Terminal panel.
|
2026-07-27 22:02:34 -07:00
|
|
|
var caps = new List<string> { "components", "session", "image", "container", "proc" };
|
2026-07-29 02:52:41 -07:00
|
|
|
caps.AddRange(wslc.Capabilities);
|
2026-07-27 22:02:34 -07:00
|
|
|
if (ai.IsAvailable) caps.Add("ai");
|
|
|
|
|
return new
|
|
|
|
|
{
|
|
|
|
|
protocol = ProtocolVersion,
|
|
|
|
|
brokerVersion = Assembly.GetExecutingAssembly().GetName().Version?.ToString(3) ?? "0.0.0",
|
|
|
|
|
wslcVersion = wslc.WslcVersion,
|
|
|
|
|
capabilities = caps,
|
|
|
|
|
};
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
case "components.missing":
|
|
|
|
|
return new { flags = await wslc.MissingComponentsAsync(ct).ConfigureAwait(false) };
|
|
|
|
|
case "components.install":
|
|
|
|
|
await wslc.InstallComponentsAsync(ct).ConfigureAwait(false);
|
|
|
|
|
return new { };
|
|
|
|
|
|
|
|
|
|
case "session.ensure":
|
|
|
|
|
{
|
|
|
|
|
var spec = Rpc.Params<SessionSpec>(p);
|
|
|
|
|
var gateway = await wslc.EnsureSessionAsync(spec, ct).ConfigureAwait(false);
|
|
|
|
|
return new { gateway };
|
|
|
|
|
}
|
|
|
|
|
case "session.terminate":
|
|
|
|
|
await wslc.TerminateSessionAsync(ct).ConfigureAwait(false);
|
|
|
|
|
return new { };
|
|
|
|
|
|
|
|
|
|
case "image.pull":
|
|
|
|
|
{
|
|
|
|
|
var args = Rpc.Params<ImagePullParams>(p);
|
|
|
|
|
await wslc.PullImageAsync(args.Ref, args.Auth, ct).ConfigureAwait(false);
|
|
|
|
|
return new { };
|
|
|
|
|
}
|
|
|
|
|
case "image.list":
|
|
|
|
|
return new { images = await wslc.ListImagesAsync(ct).ConfigureAwait(false) };
|
|
|
|
|
case "image.delete":
|
|
|
|
|
await wslc.DeleteImageAsync(Rpc.Params<ImageRefParams>(p).Ref, ct).ConfigureAwait(false);
|
|
|
|
|
return new { };
|
|
|
|
|
case "image.inspect":
|
|
|
|
|
{
|
|
|
|
|
var image = await wslc.InspectImageAsync(Rpc.Params<ImageRefParams>(p).Ref, ct)
|
|
|
|
|
.ConfigureAwait(false);
|
|
|
|
|
return new { image };
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
case "container.create":
|
|
|
|
|
await wslc.CreateContainerAsync(Rpc.Params<ContainerCreateSpec>(p), ct).ConfigureAwait(false);
|
|
|
|
|
return new { };
|
|
|
|
|
case "container.start":
|
|
|
|
|
await wslc.StartContainerAsync(Rpc.Params<ContainerNameParams>(p).Name, ct).ConfigureAwait(false);
|
|
|
|
|
return new { };
|
|
|
|
|
case "container.stop":
|
|
|
|
|
{
|
|
|
|
|
var args = Rpc.Params<ContainerStopParams>(p);
|
|
|
|
|
await wslc.StopContainerAsync(args.Name, args.Signal ?? 15, args.GraceMs ?? 5000, ct)
|
|
|
|
|
.ConfigureAwait(false);
|
|
|
|
|
return new { };
|
|
|
|
|
}
|
|
|
|
|
case "container.delete":
|
|
|
|
|
{
|
|
|
|
|
var args = Rpc.Params<ContainerDeleteParams>(p);
|
|
|
|
|
await wslc.DeleteContainerAsync(args.Name, args.Force ?? false, ct).ConfigureAwait(false);
|
|
|
|
|
return new { };
|
|
|
|
|
}
|
|
|
|
|
case "container.list":
|
|
|
|
|
return new { containers = await wslc.ListContainersAsync(ct).ConfigureAwait(false) };
|
|
|
|
|
case "container.state":
|
|
|
|
|
{
|
|
|
|
|
var state = await wslc.ContainerStateAsync(Rpc.Params<ContainerNameParams>(p).Name, ct)
|
|
|
|
|
.ConfigureAwait(false);
|
|
|
|
|
return new { state };
|
|
|
|
|
}
|
|
|
|
|
case "container.stats":
|
|
|
|
|
{
|
|
|
|
|
var stats = await wslc.ContainerStatsAsync(Rpc.Params<ContainerNameParams>(p).Name, ct)
|
|
|
|
|
.ConfigureAwait(false);
|
|
|
|
|
return new { stats };
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
case "proc.exec":
|
|
|
|
|
{
|
|
|
|
|
var spec = Rpc.Params<ProcSpec>(p);
|
|
|
|
|
var procId = Interlocked.Increment(ref nextProcId);
|
|
|
|
|
var proc = await wslc.ExecAsync(procId, spec, ct).ConfigureAwait(false);
|
|
|
|
|
lock (procsLock) procs[procId] = proc;
|
2026-07-29 17:09:35 -07:00
|
|
|
// Respond BEFORE running it. Output and exit are enqueued on the same ordered
|
|
|
|
|
// outbound queue as this response, so a command that finishes fast (`echo`) would
|
|
|
|
|
// otherwise put `proc.exit` on the wire ahead of the `procId` naming it — and a
|
|
|
|
|
// client that registers interest when it learns the procId then waits forever.
|
|
|
|
|
// Hit on the first live run, hardware-confirmed (docs/WINDOWS_PORT.md §13.3).
|
|
|
|
|
if (id is { } requestId)
|
|
|
|
|
outbound.EnqueueJson(Rpc.Response(requestId, new { procId }));
|
|
|
|
|
|
|
|
|
|
try
|
|
|
|
|
{
|
|
|
|
|
await proc.StartAsync(ct).ConfigureAwait(false);
|
|
|
|
|
}
|
|
|
|
|
catch (Exception e)
|
|
|
|
|
{
|
|
|
|
|
// The response is already on the wire, so this failure CANNOT be reported as a
|
|
|
|
|
// JSON-RPC error — that would put two responses under one id, which is a
|
|
|
|
|
// protocol violation and left the client waiting for an exit that never came.
|
|
|
|
|
// Report it the way the process itself would have: the reason on stderr, then
|
|
|
|
|
// an exit. 126 is the shell's "command found but not executable", which is
|
|
|
|
|
// what "could not start" means to every caller above.
|
|
|
|
|
lock (procsLock) procs.Remove(procId);
|
|
|
|
|
var reason = $"nucleic-brokerd: could not start process: {e.Message}\n";
|
|
|
|
|
ProcOutput(procId, stderr: true, System.Text.Encoding.UTF8.GetBytes(reason));
|
|
|
|
|
ProcExited(procId, 126);
|
|
|
|
|
}
|
|
|
|
|
return null;
|
2026-07-27 22:02:34 -07:00
|
|
|
}
|
|
|
|
|
case "proc.stdin":
|
|
|
|
|
{
|
|
|
|
|
var args = Rpc.Params<ProcStdinParams>(p);
|
|
|
|
|
await Proc(args.ProcId).WriteStdinAsync(Convert.FromBase64String(args.B64), ct)
|
|
|
|
|
.ConfigureAwait(false);
|
|
|
|
|
return new { };
|
|
|
|
|
}
|
|
|
|
|
case "proc.closeStdin":
|
|
|
|
|
await Proc(Rpc.Params<ProcIdParams>(p).ProcId).CloseStdinAsync(ct).ConfigureAwait(false);
|
|
|
|
|
return new { };
|
|
|
|
|
case "proc.signal":
|
|
|
|
|
{
|
|
|
|
|
var args = Rpc.Params<ProcSignalParams>(p);
|
|
|
|
|
await Proc(args.ProcId).SignalAsync(args.Sig, ct).ConfigureAwait(false);
|
|
|
|
|
return new { };
|
|
|
|
|
}
|
|
|
|
|
case "proc.resize":
|
|
|
|
|
{
|
|
|
|
|
var args = Rpc.Params<ProcResizeParams>(p);
|
|
|
|
|
await Proc(args.ProcId).ResizeAsync(args.Cols, args.Rows, ct).ConfigureAwait(false);
|
|
|
|
|
return new { };
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
case "ai.generate":
|
|
|
|
|
{
|
|
|
|
|
var args = Rpc.Params<AiGenerateParams>(p);
|
|
|
|
|
var text = await ai.GenerateAsync(args.Prompt, args.Schema, ct).ConfigureAwait(false);
|
|
|
|
|
return new { text };
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
default:
|
|
|
|
|
throw new NotSupportedException(method);
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
private IWslcProcess Proc(long procId)
|
|
|
|
|
{
|
|
|
|
|
lock (procsLock)
|
|
|
|
|
{
|
|
|
|
|
return procs.TryGetValue(procId, out var proc)
|
|
|
|
|
? proc
|
|
|
|
|
: throw new WslcError(WslcError.NotFound, $"unknown procId {procId}");
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
// MARK: IBrokerEvents — facade threads enqueue and return immediately.
|
|
|
|
|
|
|
|
|
|
public void ProcOutput(long procId, bool stderr, ReadOnlySpan<byte> chunk) =>
|
|
|
|
|
outbound.EnqueueProcOutput(procId, stderr, chunk);
|
|
|
|
|
|
|
|
|
|
public void ProcExited(long procId, int code)
|
|
|
|
|
{
|
|
|
|
|
lock (procsLock) procs.Remove(procId);
|
|
|
|
|
outbound.EnqueueJson(Rpc.Notification("proc.exit", new { procId, code }));
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
public void SessionDown(string reason) =>
|
|
|
|
|
outbound.EnqueueJson(Rpc.Notification("session.down", new { reason }));
|
|
|
|
|
|
|
|
|
|
public void PullProgress(string reference, string status, long current, long total) =>
|
|
|
|
|
outbound.EnqueueJson(Rpc.Notification(
|
|
|
|
|
"image.pullProgress", new { @ref = reference, status, current, total }));
|
|
|
|
|
|
|
|
|
|
public void InstallProgress(string status, double percent) =>
|
|
|
|
|
outbound.EnqueueJson(Rpc.Notification(
|
|
|
|
|
"components.installProgress", new { status, percent }));
|
|
|
|
|
|
|
|
|
|
// Param records for methods whose shapes aren't shared DTOs.
|
|
|
|
|
private sealed record ImagePullParams(string Ref, RegistryAuth? Auth);
|
|
|
|
|
private sealed record ImageRefParams(string Ref);
|
|
|
|
|
private sealed record ContainerNameParams(string Name);
|
|
|
|
|
private sealed record ContainerStopParams(string Name, int? Signal, int? GraceMs);
|
|
|
|
|
private sealed record ContainerDeleteParams(string Name, bool? Force);
|
|
|
|
|
private sealed record ProcIdParams(long ProcId);
|
|
|
|
|
private sealed record ProcStdinParams(long ProcId, string B64);
|
|
|
|
|
private sealed record ProcSignalParams(long ProcId, int Sig);
|
|
|
|
|
private sealed record ProcResizeParams(long ProcId, int Cols, int Rows);
|
|
|
|
|
private sealed record AiGenerateParams(string Prompt, JsonElement? Schema);
|
|
|
|
|
}
|