84 lines
3.5 KiB
C#
84 lines
3.5 KiB
C#
using System.Threading.Channels;
|
|
|
|
namespace NucleicBroker;
|
|
|
|
/// <summary>
|
|
/// The single writer of the broker's stdout: responses and notifications from any thread are
|
|
/// enqueued (never blocking the caller — wslc raises stdio events on WinRT threads that must
|
|
/// not stall on the pipe, docs/WINDOWS_PORT.md §3.3) and drained by one loop that also
|
|
/// coalesces bursts: consecutive stdio chunks for the same (proc, stream) queued at drain
|
|
/// time merge into one `proc.stdout`/`proc.stderr` notification, cutting per-line overhead
|
|
/// for agents' high-rate NDJSON without adding any timer or latency (only already-queued
|
|
/// items merge). Ordering is preserved end-to-end — one channel, one drain loop — so
|
|
/// `proc.exit` can never overtake the output that preceded it.
|
|
/// </summary>
|
|
public sealed class OutboundWriter
|
|
{
|
|
private abstract record Item;
|
|
private sealed record JsonLine(string Line) : Item;
|
|
private sealed record ProcChunk(long ProcId, bool Stderr, byte[] Bytes) : Item;
|
|
|
|
private readonly Channel<Item> channel = Channel.CreateUnbounded<Item>(
|
|
new UnboundedChannelOptions { SingleReader = true });
|
|
private readonly TextWriter output;
|
|
private readonly Task loop;
|
|
|
|
public OutboundWriter(TextWriter output)
|
|
{
|
|
this.output = output;
|
|
// NDJSON is LF-framed on the wire (docs/WINDOWS_PORT.md §3.3). TextWriter.NewLine
|
|
// defaults to Environment.NewLine — CRLF here — which would append a stray \r to
|
|
// every line hostd's JSONRPCConnection reads. We own this writer, so pin it.
|
|
this.output.NewLine = "\n";
|
|
loop = Task.Run(DrainAsync);
|
|
}
|
|
|
|
public void EnqueueJson(string line) => channel.Writer.TryWrite(new JsonLine(line));
|
|
|
|
public void EnqueueProcOutput(long procId, bool stderr, ReadOnlySpan<byte> chunk) =>
|
|
channel.Writer.TryWrite(new ProcChunk(procId, stderr, chunk.ToArray()));
|
|
|
|
/// <summary>Close the queue, drain everything, and flush. Idempotent.</summary>
|
|
public async Task CompleteAsync()
|
|
{
|
|
channel.Writer.TryComplete();
|
|
await loop.ConfigureAwait(false);
|
|
}
|
|
|
|
private async Task DrainAsync()
|
|
{
|
|
var reader = channel.Reader;
|
|
var batch = new List<Item>(64);
|
|
while (await reader.WaitToReadAsync().ConfigureAwait(false))
|
|
{
|
|
batch.Clear();
|
|
while (batch.Count < 256 && reader.TryRead(out var item)) batch.Add(item);
|
|
|
|
for (var i = 0; i < batch.Count; i++)
|
|
{
|
|
if (batch[i] is not ProcChunk head)
|
|
{
|
|
await output.WriteLineAsync(((JsonLine)batch[i]).Line).ConfigureAwait(false);
|
|
continue;
|
|
}
|
|
// Merge the run of chunks for the same (proc, stream) into one notification.
|
|
using var merged = new MemoryStream();
|
|
merged.Write(head.Bytes);
|
|
while (i + 1 < batch.Count
|
|
&& batch[i + 1] is ProcChunk next
|
|
&& next.ProcId == head.ProcId && next.Stderr == head.Stderr)
|
|
{
|
|
merged.Write(next.Bytes);
|
|
i++;
|
|
}
|
|
var line = Rpc.Notification(
|
|
head.Stderr ? "proc.stderr" : "proc.stdout",
|
|
new { procId = head.ProcId, b64 = Convert.ToBase64String(merged.ToArray()) });
|
|
await output.WriteLineAsync(line).ConfigureAwait(false);
|
|
}
|
|
await output.FlushAsync().ConfigureAwait(false);
|
|
}
|
|
await output.FlushAsync().ConfigureAwait(false);
|
|
}
|
|
}
|