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