Refactor SSE subscribers to use ConcurrentDictionary for thread safety and simplify response handling.

This commit is contained in:
2026-01-01 23:06:40 +03:30
parent 1d0b42654a
commit 365267a30a
+20 -29
View File
@@ -17,7 +17,7 @@ var app = builder.Build();
app.Urls.Clear(); app.Urls.Clear();
app.Urls.Add("http://0.0.0.0:5030"); app.Urls.Add("http://0.0.0.0:5030");
var subscribers = new ConcurrentDictionary<string, List<ChannelWriter<string>>>(); var subscribers = new ConcurrentDictionary<string, ConcurrentDictionary<ChannelWriter<string>, byte>>();
var cancellationSources = var cancellationSources =
new ConcurrentDictionary<string, CancellationTokenSource>(); // To manage cancellation per feedId new ConcurrentDictionary<string, CancellationTokenSource>(); // To manage cancellation per feedId
@@ -40,34 +40,23 @@ createApi.MapGet("/{feedId}", async (HttpContext context, string feedId) =>
{ {
var jsonQueryData = JsonSerializer.Serialize(dweet, AppJsonSerializerContext.Default.Dweet); var jsonQueryData = JsonSerializer.Serialize(dweet, AppJsonSerializerContext.Default.Dweet);
if (subscribers.TryGetValue(feedId, out var subscribersList)) if (subscribers.TryGetValue(feedId, out var feedSubs))
foreach (var writer in subscribersList) foreach (var writer in feedSubs.Keys)
await writer.WriteAsync(jsonQueryData); await writer.WriteAsync(jsonQueryData);
return Results.Ok(new AddDweetSucceededResponse
{ This = "succeeded", By = "dweeting", The = "dweet", With = dweet });
} }
catch (Exception e) catch (Exception e)
{ {
var faultResponse = new AddDweetFailedResponse return Results.Json(
{ new AddDweetFailedResponse
This = "failed", {
With = "WeMessedUp", This = "failed", With = "WeMessedUp",
Because = "IDK, we couldnt dweet it. Report it at: https://github.com/mmahdium/HoolIt/issues" Because = "IDK, we couldnt dweet it. Report it at: https://github.com/mmahdium/HoolIt/issues"
}; }, statusCode: 500);
var addFailedResponse =
JsonSerializer.Serialize(faultResponse, AppJsonSerializerContext.Default.AddDweetFailedResponse);
context.Response.StatusCode = 500;
context.Response.ContentType = "application/json";
await context.Response.WriteAsync(addFailedResponse);
await context.Response.CompleteAsync();
} }
var addSuccessResponse = new AddDweetSucceededResponse
{
This = "succeeded",
By = "dweeting",
The = "dweet",
With = dweet
};
return Results.Ok(addSuccessResponse);
}); });
var getLiveDataApi = app.MapGroup("/listen/for/dweets/from"); var getLiveDataApi = app.MapGroup("/listen/for/dweets/from");
@@ -87,8 +76,10 @@ getLiveDataApi.MapGet("/{feedId}",
var channel = Channel.CreateUnbounded<string>(); var channel = Channel.CreateUnbounded<string>();
var writer = channel.Writer; var writer = channel.Writer;
subscribers.GetOrAdd(feedId, _ => new List<ChannelWriter<string>>()).Add(writer); var feedSubs = subscribers.GetOrAdd(feedId, _ => new ConcurrentDictionary<ChannelWriter<string>, byte>());
Console.WriteLine($"Added subscriber to feed {feedId}"); feedSubs.TryAdd(writer, 0);
//Console.WriteLine($"Added subscriber to feed {feedId}");
try try
{ {
@@ -104,14 +95,14 @@ getLiveDataApi.MapGet("/{feedId}",
} }
catch (OperationCanceledException) catch (OperationCanceledException)
{ {
Console.WriteLine($"Cancellation requested for feed {feedId}"); //Console.WriteLine($"Cancellation requested for feed {feedId}");
} }
finally finally
{ {
if (subscribers.TryGetValue(feedId, out var list)) if (subscribers.TryGetValue(feedId, out var list))
{ {
list.Remove(writer); feedSubs.TryRemove(writer, out _);
Console.WriteLine($"Removed subscriber from feed {feedId}"); //Console.WriteLine($"Removed subscriber from feed {feedId}");
if (list.Count == 0) subscribers.TryRemove(feedId, out _); if (list.Count == 0) subscribers.TryRemove(feedId, out _);
} }
} }