using System.Reflection; using System.Text.Json; namespace NucleicBroker; /// /// The broker's request dispatcher (docs/WINDOWS_PORT.md §3.3): one instance per process, /// fed NDJSON lines from hostd, answering through an . Requests /// run concurrently (an image pull must not block a container.list), responses/notifications /// serialize through the writer's single drain loop. Also the /// sink the facade raises async events into. /// public sealed class BrokerService : IBrokerEvents { private readonly IWslc wslc; private readonly IAiProvider ai; private readonly OutboundWriter outbound; private readonly Dictionary 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); } /// Broker protocol revision, bumped with any wire change; hostd degrades from /// the `hello` capabilities rather than parsing versions (docs/WINDOWS_PORT.md §2.3). public const int ProtocolVersion = 1; /// Handle one inbound NDJSON line. The returned task completes once the /// response has been ENQUEUED (tests await it; the stdin loop fires and forgets). 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 { 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) 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)); } } /// Returns the result to respond with, or null when the handler has already /// responded (it needed the response on the wire before continuing). private async Task DispatchAsync( string method, JsonElement? p, JsonElement? id, CancellationToken ct) { switch (method) { case "hello": { // 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. var caps = new List { "components", "session", "image", "container", "proc" }; caps.AddRange(wslc.Capabilities); 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(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(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(p).Ref, ct).ConfigureAwait(false); return new { }; case "image.inspect": { var image = await wslc.InspectImageAsync(Rpc.Params(p).Ref, ct) .ConfigureAwait(false); return new { image }; } case "container.create": await wslc.CreateContainerAsync(Rpc.Params(p), ct).ConfigureAwait(false); return new { }; case "container.start": await wslc.StartContainerAsync(Rpc.Params(p).Name, ct).ConfigureAwait(false); return new { }; case "container.stop": { var args = Rpc.Params(p); await wslc.StopContainerAsync(args.Name, args.Signal ?? 15, args.GraceMs ?? 5000, ct) .ConfigureAwait(false); return new { }; } case "container.delete": { var args = Rpc.Params(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(p).Name, ct) .ConfigureAwait(false); return new { state }; } case "container.stats": { var stats = await wslc.ContainerStatsAsync(Rpc.Params(p).Name, ct) .ConfigureAwait(false); return new { stats }; } case "proc.exec": { var spec = Rpc.Params(p); var procId = Interlocked.Increment(ref nextProcId); var proc = await wslc.ExecAsync(procId, spec, ct).ConfigureAwait(false); lock (procsLock) procs[procId] = proc; // 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; } case "proc.stdin": { var args = Rpc.Params(p); await Proc(args.ProcId).WriteStdinAsync(Convert.FromBase64String(args.B64), ct) .ConfigureAwait(false); return new { }; } case "proc.closeStdin": await Proc(Rpc.Params(p).ProcId).CloseStdinAsync(ct).ConfigureAwait(false); return new { }; case "proc.signal": { var args = Rpc.Params(p); await Proc(args.ProcId).SignalAsync(args.Sig, ct).ConfigureAwait(false); return new { }; } case "proc.resize": { var args = Rpc.Params(p); await Proc(args.ProcId).ResizeAsync(args.Cols, args.Rows, ct).ConfigureAwait(false); return new { }; } case "ai.generate": { var args = Rpc.Params(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 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); }