From 70cda1fa222ffd9891e941dc8665b1909f1ff4bb Mon Sep 17 00:00:00 2001 From: Steven Pelech Date: Sun, 20 Sep 2026 16:57:17 -0500 Subject: [PATCH] fix(components): serialize SSE progress streaming via Channel and fix premature termination race condition --- Endpoints/ComponentEndpoints.cs | 67 ++++++++++--------- .../ComponentManagerAndEndpointsTests.cs | 4 ++ 2 files changed, 40 insertions(+), 31 deletions(-) diff --git a/Endpoints/ComponentEndpoints.cs b/Endpoints/ComponentEndpoints.cs index 0ff3163..df5e8ea 100644 --- a/Endpoints/ComponentEndpoints.cs +++ b/Endpoints/ComponentEndpoints.cs @@ -26,58 +26,63 @@ public static void MapComponentEndpoints(this WebApplication app) httpContext.Response.Headers.CacheControl = "no-cache"; httpContext.Response.Headers.Connection = "keep-alive"; - var syncLock = new SemaphoreSlim(1, 1); + var channel = System.Threading.Channels.Channel.CreateUnbounded(new System.Threading.Channels.UnboundedChannelOptions + { + SingleWriter = true, + SingleReader = true + }); + var progress = new Progress(percent => { - _ = Task.Run(async () => + channel.Writer.TryWrite(percent); + }); + + var writerTask = Task.Run(async () => + { + try { - await syncLock.WaitAsync(); - try + while (await channel.Reader.WaitToReadAsync(httpContext.RequestAborted)) { - var eventData = JsonSerializer.Serialize(new { progress = percent, status = "installing" }); - await httpContext.Response.WriteAsync($"data: {eventData}\n\n"); - await httpContext.Response.Body.FlushAsync(); + while (channel.Reader.TryRead(out var percent)) + { + var eventData = JsonSerializer.Serialize(new { progress = percent, status = "installing" }); + await httpContext.Response.WriteAsync($"data: {eventData}\n\n", httpContext.RequestAborted); + await httpContext.Response.Body.FlushAsync(httpContext.RequestAborted); + } } - catch { } - finally - { - syncLock.Release(); - } - }); + } + catch (OperationCanceledException) { } + catch { } }); try { var result = await componentService.InstallComponentAsync(componentId, progress, httpContext.RequestAborted); - await syncLock.WaitAsync(); - try - { - var finalData = JsonSerializer.Serialize(new { progress = 100.0, status = result ? "completed" : "failed", success = result }); - await httpContext.Response.WriteAsync($"data: {finalData}\n\n"); - await httpContext.Response.Body.FlushAsync(); - } - finally - { - syncLock.Release(); - } + channel.Writer.TryComplete(); + await writerTask; + + var finalData = JsonSerializer.Serialize(new { progress = 100.0, status = result ? "completed" : "failed", success = result }); + await httpContext.Response.WriteAsync($"data: {finalData}\n\n", httpContext.RequestAborted); + await httpContext.Response.Body.FlushAsync(httpContext.RequestAborted); } catch (OperationCanceledException) { + channel.Writer.TryComplete(); + await writerTask; // Request canceled } catch (Exception ex) { - await syncLock.WaitAsync(); + channel.Writer.TryComplete(); + await writerTask; + try { var errData = JsonSerializer.Serialize(new { progress = 0.0, status = "error", message = ex.Message }); - await httpContext.Response.WriteAsync($"data: {errData}\n\n"); - await httpContext.Response.Body.FlushAsync(); - } - finally - { - syncLock.Release(); + await httpContext.Response.WriteAsync($"data: {errData}\n\n", httpContext.RequestAborted); + await httpContext.Response.Body.FlushAsync(httpContext.RequestAborted); } + catch { } } return Results.Empty; diff --git a/LocalLLMServerManager.Tests/ComponentManagerAndEndpointsTests.cs b/LocalLLMServerManager.Tests/ComponentManagerAndEndpointsTests.cs index e1b1be1..b078fd6 100644 --- a/LocalLLMServerManager.Tests/ComponentManagerAndEndpointsTests.cs +++ b/LocalLLMServerManager.Tests/ComponentManagerAndEndpointsTests.cs @@ -59,6 +59,8 @@ public async Task InstallAndUninstallComponent_Succeeds() var installReq = new ComponentInstallRequest { ComponentId = "audio-tts" }; var installResponse = await client.PostAsJsonAsync("/api/components/install", installReq); Assert.Equal(HttpStatusCode.OK, installResponse.StatusCode); + var installContent = await installResponse.Content.ReadAsStringAsync(); + Assert.Contains("completed", installContent); // Uninstall request var uninstallReq = new ComponentInstallRequest { ComponentId = "audio-tts" }; @@ -69,6 +71,8 @@ public async Task InstallAndUninstallComponent_Succeeds() var aiInstallReq = new ComponentInstallRequest { ComponentId = "ai-assistant" }; var aiInstallResponse = await client.PostAsJsonAsync("/api/components/install", aiInstallReq); Assert.Equal(HttpStatusCode.OK, aiInstallResponse.StatusCode); + var aiInstallContent = await aiInstallResponse.Content.ReadAsStringAsync(); + Assert.Contains("completed", aiInstallContent); var aiUninstallReq = new ComponentInstallRequest { ComponentId = "ai-assistant" }; var aiUninstallResponse = await client.PostAsJsonAsync("/api/components/uninstall", aiUninstallReq);