From 51522e1f3fe0b932cdb230a443b11850ca115e65 Mon Sep 17 00:00:00 2001 From: Francesco Lorenzo D'Amico Date: Mon, 24 Aug 2026 14:41:56 +0200 Subject: [PATCH] Replace polling with push-based Telegram notification queue --- .../Controllers/Api/V1/ContactController.cs | 5 +- RazorPagesMovie/Program.cs | 1 + .../ContactMessageNotificationService.cs | 114 ++++++++++++------ .../IContactMessageNotificationQueue.cs | 32 +++++ 4 files changed, 116 insertions(+), 36 deletions(-) create mode 100644 RazorPagesMovie/Services/IContactMessageNotificationQueue.cs diff --git a/RazorPagesMovie/Controllers/Api/V1/ContactController.cs b/RazorPagesMovie/Controllers/Api/V1/ContactController.cs index 3c463a7..0d82e9e 100644 --- a/RazorPagesMovie/Controllers/Api/V1/ContactController.cs +++ b/RazorPagesMovie/Controllers/Api/V1/ContactController.cs @@ -3,12 +3,13 @@ using Microsoft.AspNetCore.Mvc; using RazorPagesMovie.Api.V1.Contact; using RazorPagesMovie.Data; using RazorPagesMovie.Models; +using RazorPagesMovie.Services; namespace RazorPagesMovie.Controllers.Api.V1; [ApiController] [Route("api/v1/[controller]")] -public class ContactController(RazorPagesMovieContext context) : ControllerBase +public class ContactController(RazorPagesMovieContext context, IContactMessageNotificationQueue notificationQueue) : ControllerBase { [HttpPost] public async Task Post(ContactDto request) @@ -31,6 +32,8 @@ public class ContactController(RazorPagesMovieContext context) : ControllerBase context.ContactMessage.Add(contactMessage); await context.SaveChangesAsync(); + notificationQueue.Enqueue(contactMessage.Id); + return Ok(new { success = true }); } } diff --git a/RazorPagesMovie/Program.cs b/RazorPagesMovie/Program.cs index 36cb4f0..d7d9f8d 100644 --- a/RazorPagesMovie/Program.cs +++ b/RazorPagesMovie/Program.cs @@ -29,6 +29,7 @@ builder.Services.AddRazorPages(options => builder.Services.AddControllers(); builder.Services.AddHttpClient(); +builder.Services.AddSingleton(); builder.Services.AddHostedService(); var app = builder.Build(); diff --git a/RazorPagesMovie/Services/ContactMessageNotificationService.cs b/RazorPagesMovie/Services/ContactMessageNotificationService.cs index ff151a3..de1d6d3 100644 --- a/RazorPagesMovie/Services/ContactMessageNotificationService.cs +++ b/RazorPagesMovie/Services/ContactMessageNotificationService.cs @@ -6,35 +6,86 @@ namespace RazorPagesMovie.Services; public class ContactMessageNotificationService : BackgroundService { - private static readonly TimeSpan CheckInterval = TimeSpan.FromMinutes(5); + private static readonly TimeSpan FallbackSweepInterval = TimeSpan.FromMinutes(2); private readonly IServiceScopeFactory _scopeFactory; private readonly IHttpClientFactory _httpClientFactory; private readonly IConfiguration _configuration; + private readonly IContactMessageNotificationQueue _queue; private readonly ILogger _logger; public ContactMessageNotificationService( IServiceScopeFactory scopeFactory, IHttpClientFactory httpClientFactory, IConfiguration configuration, + IContactMessageNotificationQueue queue, ILogger logger) { _scopeFactory = scopeFactory; _httpClientFactory = httpClientFactory; _configuration = configuration; + _queue = queue; _logger = logger; } protected override async Task ExecuteAsync(CancellationToken stoppingToken) { - while (!stoppingToken.IsCancellationRequested) + var consumerTask = ConsumeQueueAsync(stoppingToken); + var fallbackSweepTask = RunFallbackSweepAsync(stoppingToken); + + await Task.WhenAll(consumerTask, fallbackSweepTask); + } + + private async Task ConsumeQueueAsync(CancellationToken stoppingToken) + { + await foreach (var contactMessageId in _queue.DequeueAllAsync(stoppingToken)) { - await CheckForNewContactMessagesAsync(stoppingToken); - await Task.Delay(CheckInterval, stoppingToken); + try + { + await SendNotificationAsync(contactMessageId, stoppingToken); + } + finally + { + _queue.Release(contactMessageId); + } } } - private async Task CheckForNewContactMessagesAsync(CancellationToken stoppingToken) + private async Task RunFallbackSweepAsync(CancellationToken stoppingToken) + { + while (!stoppingToken.IsCancellationRequested) + { + await EnqueueUnsentFromDatabaseAsync(stoppingToken); + + try + { + await Task.Delay(FallbackSweepInterval, stoppingToken); + } + catch (OperationCanceledException) + { + // Dienst wird beendet. + } + } + } + + private async Task EnqueueUnsentFromDatabaseAsync(CancellationToken stoppingToken) + { + using var scope = _scopeFactory.CreateScope(); + var context = scope.ServiceProvider.GetRequiredService(); + + var offeneIds = await context.ContactMessage + .Where(c => !c.NotificationSent) + .OrderBy(c => c.SentAt) + .Select(c => c.Id) + .ToListAsync(stoppingToken); + + foreach (var id in offeneIds) + { + _queue.Enqueue(id); + } + } + + private async Task SendNotificationAsync(Guid contactMessageId, CancellationToken stoppingToken) { var botToken = _configuration["Telegram:BotToken"]; var chatId = _configuration["Telegram:ChatId"]; @@ -48,48 +99,41 @@ public class ContactMessageNotificationService : BackgroundService using var scope = _scopeFactory.CreateScope(); var context = scope.ServiceProvider.GetRequiredService(); - var neueEintraege = await context.ContactMessage - .Where(c => !c.NotificationSent) - .OrderBy(c => c.SentAt) - .ToListAsync(stoppingToken); - - if (neueEintraege.Count == 0) + var eintrag = await context.ContactMessage.FindAsync(new object?[] { contactMessageId }, stoppingToken); + if (eintrag is null || eintrag.NotificationSent) { return; } + var text = $"Neue Kontaktanfrage von {eintrag.Name}\n" + + $"E-Mail: {eintrag.Email}\n" + + (string.IsNullOrWhiteSpace(eintrag.Phone) ? "" : $"Telefon: {eintrag.Phone}\n") + + (string.IsNullOrWhiteSpace(eintrag.Subject) ? "" : $"Betreff: {eintrag.Subject}\n") + + $"\n{eintrag.Message}"; + var httpClient = _httpClientFactory.CreateClient(); var sendMessageUrl = $"https://api.telegram.org/bot{botToken}/sendMessage"; - foreach (var eintrag in neueEintraege) + try { - var text = $"Neue Kontaktanfrage von {eintrag.Name}\n" - + $"E-Mail: {eintrag.Email}\n" - + (string.IsNullOrWhiteSpace(eintrag.Phone) ? "" : $"Telefon: {eintrag.Phone}\n") - + (string.IsNullOrWhiteSpace(eintrag.Subject) ? "" : $"Betreff: {eintrag.Subject}\n") - + $"\n{eintrag.Message}"; + var response = await httpClient.PostAsJsonAsync( + sendMessageUrl, + new { chat_id = chatId, text }, + stoppingToken); - try + if (response.IsSuccessStatusCode) { - var response = await httpClient.PostAsJsonAsync( - sendMessageUrl, - new { chat_id = chatId, text }, - stoppingToken); - - if (response.IsSuccessStatusCode) - { - eintrag.NotificationSent = true; - await context.SaveChangesAsync(stoppingToken); - } - else - { - _logger.LogWarning("Telegram-Versand fehlgeschlagen für ContactMessage {Id}: {StatusCode}", eintrag.Id, response.StatusCode); - } + eintrag.NotificationSent = true; + await context.SaveChangesAsync(stoppingToken); } - catch (Exception ex) + else { - _logger.LogError(ex, "Telegram-Versand fehlgeschlagen für ContactMessage {Id}", eintrag.Id); + _logger.LogWarning("Telegram-Versand fehlgeschlagen für ContactMessage {Id}: {StatusCode}", eintrag.Id, response.StatusCode); } } + catch (Exception ex) + { + _logger.LogError(ex, "Telegram-Versand fehlgeschlagen für ContactMessage {Id}", eintrag.Id); + } } -} \ No newline at end of file +} diff --git a/RazorPagesMovie/Services/IContactMessageNotificationQueue.cs b/RazorPagesMovie/Services/IContactMessageNotificationQueue.cs new file mode 100644 index 0000000..10f7784 --- /dev/null +++ b/RazorPagesMovie/Services/IContactMessageNotificationQueue.cs @@ -0,0 +1,32 @@ +using System.Collections.Concurrent; +using System.Threading.Channels; + +namespace RazorPagesMovie.Services; + +public interface IContactMessageNotificationQueue +{ + void Enqueue(Guid contactMessageId); + + IAsyncEnumerable DequeueAllAsync(CancellationToken cancellationToken); + + void Release(Guid contactMessageId); +} + +public class ContactMessageNotificationQueue : IContactMessageNotificationQueue +{ + private readonly Channel _channel = Channel.CreateUnbounded(); + private readonly ConcurrentDictionary _inFlight = new(); + + public void Enqueue(Guid contactMessageId) + { + if (_inFlight.TryAdd(contactMessageId, 0)) + { + _channel.Writer.TryWrite(contactMessageId); + } + } + + public IAsyncEnumerable DequeueAllAsync(CancellationToken cancellationToken) => + _channel.Reader.ReadAllAsync(cancellationToken); + + public void Release(Guid contactMessageId) => _inFlight.TryRemove(contactMessageId, out _); +}