| 1 | using System.Text; |
| 2 | using System.Threading.Channels; |
| 3 | |
| 4 | namespace CascadeIDE.Services; |
| 5 | |
| 6 | /// <summary> |
| 7 | /// Чтение кусков лога из канала (stdout/stderr процесса) с батчированием обратных вызовов |
| 8 | /// <paramref name="append"/>, чтобы не вызывать UI на каждый мелкий read (ADR 0094). |
| 9 | /// </summary> |
| 10 | public static class BuildLogIngestion |
| 11 | { |
| 12 | /// <summary>Дефолт: очередь «кусков» строки; при <see cref="BoundedChannelFullMode.Wait"/> — backpressure к продюсеру (помпы процесса).</summary> |
| 13 | public const int DefaultChannelChunkCapacity = 32; |
| 14 | |
| 15 | /// <summary>Один писатель, один читатель, <see cref="BoundedChannelFullMode.Wait"/>: наполненная очередь тормозит <c>WriteAsync</c> до снятия с хвоста.</summary> |
| 16 | public static Channel<string> CreateBuildLogChannel(int capacity = DefaultChannelChunkCapacity) => |
| 17 | Channel.CreateBounded<string>(new BoundedChannelOptions(Math.Max(1, capacity)) |
| 18 | { |
| 19 | SingleReader = true, |
| 20 | SingleWriter = true, |
| 21 | FullMode = BoundedChannelFullMode.Wait |
| 22 | }); |
| 23 | |
| 24 | /// <summary> |
| 25 | /// Читает все элементы из <paramref name="reader"/>; при накоплении не менее |
| 26 | /// <paramref name="maxBatchChars"/> символов сбрасывает пачку в <paramref name="append"/>. |
| 27 | /// <paramref name="onEachDequeuedChunk"/> (если задан) вызывается на каждом снятом с канала куске |
| 28 | /// — для накопления полного лога (MCP, парсеры) параллельно батчу в панель. |
| 29 | /// </summary> |
| 30 | public static async Task DrainToAppendAsync( |
| 31 | ChannelReader<string> reader, |
| 32 | Action<string> append, |
| 33 | int maxBatchChars = 8192, |
| 34 | Action<string>? onEachDequeuedChunk = null, |
| 35 | CancellationToken cancellationToken = default) |
| 36 | { |
| 37 | ArgumentNullException.ThrowIfNull(reader); |
| 38 | ArgumentNullException.ThrowIfNull(append); |
| 39 | ArgumentOutOfRangeException.ThrowIfLessThan(maxBatchChars, 1); |
| 40 | |
| 41 | var batch = new StringBuilder(); |
| 42 | var size = 0; |
| 43 | |
| 44 | await foreach (var chunk in reader.ReadAllAsync(cancellationToken).ConfigureAwait(false)) |
| 45 | { |
| 46 | if (chunk.Length == 0) |
| 47 | continue; |
| 48 | |
| 49 | onEachDequeuedChunk?.Invoke(chunk); |
| 50 | batch.Append(chunk); |
| 51 | size += chunk.Length; |
| 52 | if (size < maxBatchChars) |
| 53 | continue; |
| 54 | |
| 55 | append(batch.ToString()); |
| 56 | batch.Clear(); |
| 57 | size = 0; |
| 58 | } |
| 59 | |
| 60 | if (batch.Length > 0) |
| 61 | append(batch.ToString()); |
| 62 | } |
| 63 | } |
| 64 | |