using System.Threading.Channels; namespace NucleicBroker; /// /// 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. /// 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 channel = Channel.CreateUnbounded( new UnboundedChannelOptions { SingleReader = true }); private readonly TextWriter output; private readonly Task loop; public OutboundWriter(TextWriter output) { this.output = output; loop = Task.Run(DrainAsync); } public void EnqueueJson(string line) => channel.Writer.TryWrite(new JsonLine(line)); public void EnqueueProcOutput(long procId, bool stderr, ReadOnlySpan chunk) => channel.Writer.TryWrite(new ProcChunk(procId, stderr, chunk.ToArray())); /// Close the queue, drain everything, and flush. Idempotent. public async Task CompleteAsync() { channel.Writer.TryComplete(); await loop.ConfigureAwait(false); } private async Task DrainAsync() { var reader = channel.Reader; var batch = new List(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); } }