From 284b58299f0d327c6023529b0d08f7151a213e0f Mon Sep 17 00:00:00 2001 From: Mohammad Mahdi Date: Sun, 15 Feb 2026 11:01:54 +0330 Subject: [PATCH] Refactor SSE implementation to improve efficiency, replace `ConcurrentDictionary` structure, and optimize response handling. --- HoolIt.sln.DotSettings.user | 3 +- HoolIt/Program.cs | 142 ++++++++++++++++++------------------ 2 files changed, 72 insertions(+), 73 deletions(-) diff --git a/HoolIt.sln.DotSettings.user b/HoolIt.sln.DotSettings.user index 1a152eb..2c2b2b7 100644 --- a/HoolIt.sln.DotSettings.user +++ b/HoolIt.sln.DotSettings.user @@ -1,3 +1,4 @@  ForceIncluded - ForceIncluded \ No newline at end of file + ForceIncluded + ForceIncluded \ No newline at end of file diff --git a/HoolIt/Program.cs b/HoolIt/Program.cs index 0bcf408..019d7ad 100644 --- a/HoolIt/Program.cs +++ b/HoolIt/Program.cs @@ -1,6 +1,5 @@ +using System.Buffers; using System.Collections.Concurrent; -using System.Runtime.CompilerServices; -using System.Text; using System.Text.Json; using System.Text.Json.Serialization; using System.Threading.Channels; @@ -17,97 +16,96 @@ var app = builder.Build(); app.Urls.Clear(); app.Urls.Add("http://0.0.0.0:5030"); -var subscribers = new ConcurrentDictionary, byte>>(); -var cancellationSources = - new ConcurrentDictionary(); // To manage cancellation per feedId +var topicSubscribers = new ConcurrentDictionary>>(); app.MapGet("/", () => Results.Redirect("https://github.com/mmahdium/HoolIt")); -// HAPI! -// https://github.com/jheising/HAPI var createApi = app.MapGroup("/dweet/for"); -createApi.MapGet("/{feedId}", async (HttpContext context, string feedId) => +createApi.MapGet("/{feedId}", (HttpContext context, string feedId) => { - var queryDataDic = context.Request.Query.ToDictionary(k => k.Key, v => v.Value[0]); var dweet = new Dweet { - Content = queryDataDic, + Content = context.Request.Query.ToDictionary(k => k.Key, v => v.Value[0])!, Created = DateTime.UtcNow, Thing = feedId }; - try + var utf8Bytes = JsonSerializer.SerializeToUtf8Bytes(dweet, AppJsonSerializerContext.Default.Dweet); + + if (topicSubscribers.TryGetValue(feedId, out var subscribers)) { - var jsonQueryData = JsonSerializer.Serialize(dweet, AppJsonSerializerContext.Default.Dweet); + var writers = subscribers.Values.ToArray(); - if (subscribers.TryGetValue(feedId, out var feedSubs)) - foreach (var writer in feedSubs.Keys) - await writer.WriteAsync(jsonQueryData); - - - return Results.Ok(new AddDweetSucceededResponse - { This = "succeeded", By = "dweeting", The = "dweet", With = dweet }); + foreach (var w in writers) w.TryWrite(utf8Bytes); } - catch (Exception e) + + return Results.Ok(new AddDweetSucceededResponse { - return Results.Json( - new AddDweetFailedResponse - { - This = "failed", With = "WeMessedUp", - Because = "IDK, we couldnt dweet it. Report it at: https://github.com/mmahdium/HoolIt/issues" - }, statusCode: 500); - } + This = "succeeded", + By = "dweeting", + The = "dweet", + With = dweet + }); }); +// Subscribe endpoint var getLiveDataApi = app.MapGroup("/listen/for/dweets/from"); -getLiveDataApi.MapGet("/{feedId}", - async (HttpContext context, string feedId, IHostApplicationLifetime appLifetime, - CancellationToken reqCancellationToken) => +getLiveDataApi.MapGet("/{feedId}", async (HttpContext context, string feedId, CancellationToken reqCancellationToken) => +{ + context.Response.StatusCode = 200; + context.Response.Headers.ContentType = "text/plain; charset=utf-8"; + context.Response.Headers.CacheControl = "no-cache"; + context.Response.Headers["X-Content-Type-Options"] = "nosniff"; + + context.Features.Get()? + .DisableBuffering(); + + var channel = Channel.CreateBounded(new BoundedChannelOptions(512) // tune buffer count { - context.Response.StatusCode = 200; - context.Response.Headers.ContentType = "text/plain"; - context.Response.Headers.CacheControl = "no-cache"; - context.Response.Headers["X-Content-Type-Options"] = "nosniff"; - - // Disable response buffering - context.Features.Get()? - .DisableBuffering(); - - var channel = Channel.CreateUnbounded(); - var writer = channel.Writer; - - var feedSubs = subscribers.GetOrAdd(feedId, _ => new ConcurrentDictionary, byte>()); - feedSubs.TryAdd(writer, 0); - - //Console.WriteLine($"Added subscriber to feed {feedId}"); - - try - { - var reader = channel.Reader; - - while (!reqCancellationToken.IsCancellationRequested && - await reader.WaitToReadAsync(reqCancellationToken)) - while (reader.TryRead(out var msg)) - { - await context.Response.WriteAsync(msg + "\n", reqCancellationToken); - await context.Response.Body.FlushAsync(reqCancellationToken); - } - } - catch (OperationCanceledException) - { - //Console.WriteLine($"Cancellation requested for feed {feedId}"); - } - finally - { - if (subscribers.TryGetValue(feedId, out var list)) - { - feedSubs.TryRemove(writer, out _); - //Console.WriteLine($"Removed subscriber from feed {feedId}"); - if (list.Count == 0) subscribers.TryRemove(feedId, out _); - } - } + SingleReader = true, + SingleWriter = true, + FullMode = BoundedChannelFullMode.DropOldest }); + var writer = channel.Writer; + + var subscribers = topicSubscribers.GetOrAdd(feedId, _ => new ConcurrentDictionary>()); + var subscriberId = Guid.CreateVersion7(DateTimeOffset.UtcNow); + subscribers.TryAdd(subscriberId, writer); + + try + { + var reader = channel.Reader; + + var bodyWriter = context.Response.BodyWriter; + + while (await reader.WaitToReadAsync(reqCancellationToken)) + while (reader.TryRead(out var msgBytes)) + { + var memory = bodyWriter.GetMemory(msgBytes.Length + 1); // +1 for newline + msgBytes.CopyTo(memory); + memory.Span[msgBytes.Length] = (byte)'\n'; + bodyWriter.Advance(msgBytes.Length + 1); + + await bodyWriter.FlushAsync(reqCancellationToken); + } + } + catch (OperationCanceledException) + { + // client disconnected + } + finally + { + if (topicSubscribers.TryGetValue(feedId, out var existing)) + { + existing.TryRemove(subscriberId, out var removedWriter); + removedWriter?.TryComplete(); + + if (existing.IsEmpty) topicSubscribers.TryRemove(feedId, out _); + } + } +}); + app.Run(); [JsonSerializable(typeof(Dweet))]