Merge nucleic/lucid-river-toad-6efj into dev
This commit is contained in:
@@ -0,0 +1,143 @@
|
||||
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;
|
||||
}
|
||||
Reference in New Issue
Block a user