diff --git a/CLAUDE.md b/CLAUDE.md index ea79198..31b11f3 100644 --- a/CLAUDE.md +++ b/CLAUDE.md @@ -146,9 +146,11 @@ of the few places conditional compilation is warranted. ### Stream reading -`AsyncProcessStreamReader` reads stdout and stderr concurrently with 4096-character buffers and -performs a final read after process exit, which is what ensures short-lived processes do not lose -buffered output. +`AsyncProcessStreamReader` reads stdout and stderr concurrently with 4096-character buffers, and +keeps reading each one until it reports end of stream rather than stopping when the process exits. +A process can exit with tens of kilobytes still in the pipe, so stopping at exit loses output. Each +read is raced against the caller's cancellation token, so a descendant that keeps the pipe open +cannot hang a cancelled call. ### Process configuration diff --git a/RunCommand.Test/RunCommandTests.cs b/RunCommand.Test/RunCommandTests.cs index eadac38..c5090f7 100644 --- a/RunCommand.Test/RunCommandTests.cs +++ b/RunCommand.Test/RunCommandTests.cs @@ -923,4 +923,48 @@ await Assert.ThrowsAsync( StandardInput = StandardInputMode.Closed, })).ConfigureAwait(false); } + + [TestMethod] + public async Task ExecuteAsyncShouldDeliverOutputStillInThePipeWhenTheProcessExits() + { + // Far more than one 4096-character read and more than a Linux pipe's 64 KB, from a command + // that exits as soon as it has written it, so most of the output is still unread at exit. + const int length = 200_000; + string path = Path.GetTempFileName(); + try + { + await File.WriteAllBytesAsync(path, [.. Enumerable.Repeat((byte)'a', length)]).ConfigureAwait(false); + (string fileName, string[] arguments) = GetEmitFileBytesCommand(path); + + // The loss depends on how far the reads have got when the process exits, so one clean + // run proves little. Repeat it. + StringBuilder output = new(); + for (int attempt = 0; attempt < 10; attempt++) + { + lock (output) + { + output.Clear(); + } + + int exitCode = await RunCommand.ExecuteAsync( + fileName, + arguments, + new OutputHandler(o => + { + lock (output) + { + output.Append(o); + } + }), + new CommandOptions()).ConfigureAwait(false); + + Assert.AreEqual(0, exitCode, "Expected exit code to be 0 for successful command."); + Assert.AreEqual(length, output.Length, $"Output was truncated on attempt {attempt + 1}."); + } + } + finally + { + File.Delete(path); + } + } } diff --git a/RunCommand/AsyncProcessStreamReader.cs b/RunCommand/AsyncProcessStreamReader.cs index c04be15..9d0a2d4 100644 --- a/RunCommand/AsyncProcessStreamReader.cs +++ b/RunCommand/AsyncProcessStreamReader.cs @@ -61,48 +61,11 @@ internal async Task Start(CancellationToken cancellationToken) Task cancelled = cancellationSource.Task; - Task outputTask = Task.CompletedTask; - Task errorTask = Task.CompletedTask; - - // Continuously read until the process has exited. - do - { - // A faulted task is a completed one, so without this the checks below would replace a - // failed read with a fresh one and the failure it carries would never be observed. That - // made a decode error on a long-running command a coin toss: it surfaced only when the - // process happened to exit before the loop came back around. - if (outputTask.IsFaulted || errorTask.IsFaulted) - { - break; - } - - if (outputTask.IsCompleted) - { - outputTask = ReadAndCallback(outputStream, outputBuffer, outputHandler.HandleStandardOutputData, isStandardOutput: true); - } - - if (errorTask.IsCompleted) - { - errorTask = ReadAndCallback(errorStream, errorBuffer, outputHandler.HandleStandardErrorData, isStandardOutput: false); - } - - Task first = await Task.WhenAny(outputTask, errorTask, cancelled).ConfigureAwait(false); - - if (ReferenceEquals(first, cancelled)) - { - Abandon(outputTask, errorTask); - return; - } - } while (!process.HasExited); - - if (!await DrainOrAbandon(outputTask, errorTask, cancelled).ConfigureAwait(false)) - { - return; - } - - // Read any remaining data after process exit. - outputTask = ReadAndCallback(outputStream, outputBuffer, outputHandler.HandleStandardOutputData, isStandardOutput: true); - errorTask = ReadAndCallback(errorStream, errorBuffer, outputHandler.HandleStandardErrorData, isStandardOutput: false); + // Each stream is read until it reports end of stream, not merely until the process exits. A + // process can exit with far more output still sitting in the pipe than one read returns, and + // stopping at exit plus a final read silently dropped everything past the first few kilobytes. + Task outputTask = ReadToEnd(outputStream, outputBuffer, outputHandler.HandleStandardOutputData, isStandardOutput: true); + Task errorTask = ReadToEnd(errorStream, errorBuffer, outputHandler.HandleStandardErrorData, isStandardOutput: false); _ = await DrainOrAbandon(outputTask, errorTask, cancelled).ConfigureAwait(false); } @@ -151,26 +114,25 @@ private static void Abandon(params Task[] reads) } } - private async Task ReadAndCallback(StreamReader streamReader, char[] buffer, Action? onData, bool isStandardOutput) => - await streamReader.ReadAsync(buffer, 0, buffer.Length) - .ContinueWith(t => ReadCallback(t, buffer, onData, isStandardOutput), TaskScheduler.Current) - .ConfigureAwait(false); - - private void ReadCallback(Task readTask, char[] buffer, Action? onData, bool isStandardOutput) + /// + /// Reads a stream until it reports end of stream, handing each chunk to . + /// + /// The stream to read. + /// The buffer each read fills. + /// The handler each chunk is passed to. + /// Whether the stream is standard output. + private async Task ReadToEnd(StreamReader streamReader, char[] buffer, Action? onData, bool isStandardOutput) { - int charsRead = readTask.Result; - - if (charsRead <= 0) + int charsRead; + while ((charsRead = await streamReader.ReadAsync(buffer, 0, buffer.Length).ConfigureAwait(false)) > 0) { - return; - } - - string data = new(buffer, 0, charsRead); - data = StripLeadingByteOrderMark(data, isStandardOutput); + string data = new(buffer, 0, charsRead); + data = StripLeadingByteOrderMark(data, isStandardOutput); - if (data.Length > 0) - { - onData?.Invoke(data); + if (data.Length > 0) + { + onData?.Invoke(data); + } } }