144 lines
5.9 KiB
C#
144 lines
5.9 KiB
C#
using System.Diagnostics;
|
|
using System.Text;
|
|
using System.Text.Json;
|
|
|
|
namespace WslcSpike;
|
|
|
|
/// <summary>
|
|
/// A minimal hostd stand-in: spawns `nucleic-brokerd.exe` and speaks the §3.3 NDJSON JSON-RPC
|
|
/// surface over its stdio, exactly as `WslcBrokerClient.swift` does.
|
|
///
|
|
/// Deliberately a CLIENT of the broker rather than a second implementation of the wslc calls.
|
|
/// The point of this spike is to exercise the code that actually ships — `WslcFacade`,
|
|
/// `WslcInternal`, `BrokerService`, `OutboundWriter` — so a passing run is evidence about
|
|
/// Nucleic, not about a parallel program that happens to use the same SDK.
|
|
/// </summary>
|
|
internal sealed class BrokerClient : IAsyncDisposable
|
|
{
|
|
private readonly Process process;
|
|
private readonly Dictionary<long, TaskCompletionSource<JsonElement>> pending = [];
|
|
private readonly Lock pendingLock = new();
|
|
private long nextId;
|
|
|
|
/// <summary>Raised for every notification (proc.stdout, image.pullProgress, …), on the
|
|
/// reader thread. Handlers must not block.</summary>
|
|
internal event Action<string, JsonElement>? Notification;
|
|
|
|
private BrokerClient(Process process) => this.process = process;
|
|
|
|
internal static BrokerClient Spawn(string brokerPath)
|
|
{
|
|
var info = new ProcessStartInfo(brokerPath)
|
|
{
|
|
RedirectStandardInput = true,
|
|
RedirectStandardOutput = true,
|
|
RedirectStandardError = true,
|
|
UseShellExecute = false,
|
|
// The broker pins LF framing on stdout; read it back as UTF-8 with no BOM.
|
|
StandardOutputEncoding = new UTF8Encoding(false),
|
|
StandardErrorEncoding = new UTF8Encoding(false),
|
|
};
|
|
var process = Process.Start(info)
|
|
?? throw new InvalidOperationException($"could not start {brokerPath}");
|
|
|
|
var client = new BrokerClient(process);
|
|
_ = Task.Run(client.ReadLoopAsync);
|
|
// The broker's diagnostics go to stderr and are where WslcFacade/WslcInternal report what
|
|
// bound, what degraded and why — the most useful output in the whole run when something
|
|
// is wrong. Prefix so they cannot be mistaken for spike output.
|
|
_ = Task.Run(async () =>
|
|
{
|
|
for (string? line; (line = await process.StandardError.ReadLineAsync()) is not null;)
|
|
Console.WriteLine($" [brokerd] {line}");
|
|
});
|
|
return client;
|
|
}
|
|
|
|
private async Task ReadLoopAsync()
|
|
{
|
|
for (string? line; (line = await process.StandardOutput.ReadLineAsync()) is not null;)
|
|
{
|
|
if (line.Length == 0) continue;
|
|
JsonElement root;
|
|
try { root = JsonDocument.Parse(line).RootElement.Clone(); }
|
|
catch (JsonException) { Console.WriteLine($" [unparseable] {line}"); continue; }
|
|
|
|
if (root.TryGetProperty("id", out var id) && id.ValueKind == JsonValueKind.Number)
|
|
{
|
|
TaskCompletionSource<JsonElement>? waiter;
|
|
lock (pendingLock)
|
|
{
|
|
pending.Remove(id.GetInt64(), out waiter);
|
|
}
|
|
waiter?.TrySetResult(root);
|
|
}
|
|
else if (root.TryGetProperty("method", out var method))
|
|
{
|
|
root.TryGetProperty("params", out var args);
|
|
Notification?.Invoke(method.GetString() ?? "", args);
|
|
}
|
|
}
|
|
|
|
// stdout closed: the broker exited. Fail everything outstanding rather than hanging —
|
|
// this is the `brokerLost` condition §2.3 describes, and a spike that hangs here teaches
|
|
// nothing.
|
|
lock (pendingLock)
|
|
{
|
|
foreach (var waiter in pending.Values)
|
|
waiter.TrySetException(new InvalidOperationException("brokerd exited"));
|
|
pending.Clear();
|
|
}
|
|
}
|
|
|
|
/// <summary>Send a request and await its result, throwing on a JSON-RPC error so a failing
|
|
/// step stops the scenario at the point it failed.</summary>
|
|
internal async Task<JsonElement> CallAsync(string method, object? args = null, int timeoutSeconds = 120)
|
|
{
|
|
var id = Interlocked.Increment(ref nextId);
|
|
var waiter = new TaskCompletionSource<JsonElement>(TaskCreationOptions.RunContinuationsAsynchronously);
|
|
lock (pendingLock) pending[id] = waiter;
|
|
|
|
var request = args is null
|
|
? $$"""{"jsonrpc":"2.0","id":{{id}},"method":"{{method}}"}"""
|
|
: $$"""{"jsonrpc":"2.0","id":{{id}},"method":"{{method}}","params":{{JsonSerializer.Serialize(args)}}}""";
|
|
await process.StandardInput.WriteAsync(request + "\n");
|
|
await process.StandardInput.FlushAsync();
|
|
|
|
var response = await waiter.Task.WaitAsync(TimeSpan.FromSeconds(timeoutSeconds));
|
|
if (response.TryGetProperty("error", out var error))
|
|
{
|
|
var kind = error.TryGetProperty("data", out var data)
|
|
&& data.TryGetProperty("kind", out var k) ? k.GetString() : null;
|
|
throw new BrokerError(method, error.GetProperty("message").GetString() ?? "", kind);
|
|
}
|
|
return response.GetProperty("result");
|
|
}
|
|
|
|
/// <summary>Kill the broker without letting it shut down — the §2.3 crash this spike's
|
|
/// recovery step needs to simulate.</summary>
|
|
internal void Kill()
|
|
{
|
|
try { process.Kill(entireProcessTree: true); } catch (InvalidOperationException) { }
|
|
}
|
|
|
|
public async ValueTask DisposeAsync()
|
|
{
|
|
try
|
|
{
|
|
process.StandardInput.Close(); // the graceful path: broker exits when stdin closes
|
|
await process.WaitForExitAsync().WaitAsync(TimeSpan.FromSeconds(10));
|
|
}
|
|
catch (Exception)
|
|
{
|
|
Kill();
|
|
}
|
|
process.Dispose();
|
|
}
|
|
}
|
|
|
|
internal sealed class BrokerError(string method, string message, string? kind)
|
|
: Exception($"{method} failed: {message}" + (kind is null ? "" : $" [kind={kind}]"))
|
|
{
|
|
internal string? Kind { get; } = kind;
|
|
}
|