using MiaoJiZhang.Domain.Entities; using MiaoJiZhang.Infrastructure.Persistence; using Microsoft.EntityFrameworkCore; namespace MiaoJiZhang.Api.Services; public sealed class PushDispatchService( IServiceScopeFactory scopeFactory, IConfiguration configuration, ILogger logger) : BackgroundService { private static readonly TimeSpan[] RetrySchedule = [ TimeSpan.FromMinutes(1), TimeSpan.FromMinutes(5), TimeSpan.FromMinutes(30), TimeSpan.FromHours(2), TimeSpan.FromHours(6), ]; private DateTime nextCleanupAt = DateTime.MinValue; protected override async Task ExecuteAsync(CancellationToken stoppingToken) { using var timer = new PeriodicTimer(TimeSpan.FromSeconds(5)); while (!stoppingToken.IsCancellationRequested) { if (configuration.GetValue("Push:Enabled")) { try { await RecoverAndFinalizeMessages(stoppingToken); await ExpandMessages(stoppingToken); await DispatchDeliveries(stoppingToken); if (nextCleanupAt <= DateTime.UtcNow) { await Cleanup(stoppingToken); nextCleanupAt = DateTime.UtcNow.AddHours(6); } } catch (OperationCanceledException) when (stoppingToken.IsCancellationRequested) { break; } catch (Exception exception) { logger.LogError(exception, "Push dispatch loop failed"); } } await timer.WaitForNextTickAsync(stoppingToken); } } private async Task RecoverAndFinalizeMessages(CancellationToken ct) { using var scope = scopeFactory.CreateScope(); var db = scope.ServiceProvider.GetRequiredService(); var now = DateTime.UtcNow; var staleBefore = now.AddMinutes(-3); await db.PushMessages .Where(message => message.State == PushMessageStates.Sending && message.StartedAt <= staleBefore && !db.PushDeliveries.Any(delivery => delivery.PushMessageId == message.Id)) .ExecuteUpdateAsync(setters => setters .SetProperty(message => message.State, PushMessageStates.Queued) .SetProperty(message => message.StartedAt, (DateTime?)null) .SetProperty(message => message.UpdatedAt, now), ct); var ready = await db.PushMessages .Where(message => message.State == PushMessageStates.Sending && db.PushDeliveries.Any(delivery => delivery.PushMessageId == message.Id) && !db.PushDeliveries.Any(delivery => delivery.PushMessageId == message.Id && (delivery.State == PushDeliveryStates.Queued || delivery.State == PushDeliveryStates.Sending))) .Select(message => message.Id) .Take(100) .ToListAsync(ct); foreach (var messageId in ready) await FinalizeMessage(db, messageId, ct); } private async Task ExpandMessages(CancellationToken ct) { using var scope = scopeFactory.CreateScope(); var db = scope.ServiceProvider.GetRequiredService(); var now = DateTime.UtcNow; var candidates = await db.PushMessages .Where(message => message.State == PushMessageStates.Queued || (message.State == PushMessageStates.Scheduled && message.ScheduledAt <= now)) .OrderBy(message => message.ScheduledAt ?? message.CreatedAt) .Select(message => message.Id) .Take(20) .ToListAsync(ct); foreach (var id in candidates) { var claimed = await db.PushMessages .Where(message => message.Id == id && (message.State == PushMessageStates.Queued || (message.State == PushMessageStates.Scheduled && message.ScheduledAt <= now))) .ExecuteUpdateAsync(setters => setters .SetProperty(message => message.State, PushMessageStates.Sending) .SetProperty(message => message.StartedAt, now) .SetProperty(message => message.UpdatedAt, now), ct); if (claimed != 1) continue; var message = await db.PushMessages.FirstAsync(item => item.Id == id, ct); var deviceQuery = db.PushDevices .Where(device => device.IsActive && device.NotificationsAllowed && !device.User.IsBanned && device.User.AccountClosureScheduledAt == null); if (message.IsTest) { deviceQuery = deviceQuery.Where(device => device.Id == message.TestDeviceId); } else { deviceQuery = deviceQuery.Where(device => db.UserPushPreferences.Any(preference => preference.UserId == device.UserId && preference.Category == message.Category && preference.IsEnabled)); if (message.TargetUserId.HasValue) deviceQuery = deviceQuery.Where(device => device.UserId == message.TargetUserId.Value); if (!string.IsNullOrWhiteSpace(message.Flavor)) deviceQuery = deviceQuery.Where(device => device.Flavor == message.Flavor); if (!string.IsNullOrWhiteSpace(message.ProviderFilter)) deviceQuery = deviceQuery.Where(device => device.Provider == message.ProviderFilter); if (message.MinVersionCode.HasValue) deviceQuery = deviceQuery.Where(device => device.VersionCode >= message.MinVersionCode.Value); if (message.MaxVersionCode.HasValue) deviceQuery = deviceQuery.Where(device => device.VersionCode <= message.MaxVersionCode.Value); } var devices = await deviceQuery.Select(device => new { device.Id, device.UserId, device.Provider, }).ToListAsync(ct); var existing = await db.PushDeliveries .Where(delivery => delivery.PushMessageId == id) .Select(delivery => delivery.PushDeviceId) .ToListAsync(ct); var existingIds = existing.ToHashSet(); foreach (var device in devices.Where(device => !existingIds.Contains(device.Id))) { db.PushDeliveries.Add(new PushDelivery { PushMessageId = id, PushDeviceId = device.Id, UserId = device.UserId, Provider = device.Provider, State = PushDeliveryStates.Queued, NextAttemptAt = now, CreatedAt = now, UpdatedAt = now, }); } if (devices.Count == 0) { message.State = PushMessageStates.Completed; message.CompletedAt = now; message.UpdatedAt = now; } await db.SaveChangesAsync(ct); } } private async Task DispatchDeliveries(CancellationToken ct) { using var scope = scopeFactory.CreateScope(); var db = scope.ServiceProvider.GetRequiredService(); var registry = scope.ServiceProvider.GetRequiredService(); var tokenProtector = scope.ServiceProvider.GetRequiredService(); var now = DateTime.UtcNow; var candidates = await db.PushDeliveries .Where(delivery => (delivery.State == PushDeliveryStates.Queued && delivery.NextAttemptAt <= now) || (delivery.State == PushDeliveryStates.Sending && delivery.LeaseExpiresAt <= now)) .OrderBy(delivery => delivery.NextAttemptAt) .Select(delivery => delivery.Id) .Take(50) .ToListAsync(ct); var affectedMessages = new HashSet(); foreach (var id in candidates) { var claimNow = DateTime.UtcNow; var leaseId = Guid.NewGuid().ToString(); var leaseUntil = claimNow.AddMinutes(2); var claimed = await db.PushDeliveries .Where(delivery => delivery.Id == id && ((delivery.State == PushDeliveryStates.Queued && delivery.NextAttemptAt <= claimNow) || (delivery.State == PushDeliveryStates.Sending && delivery.LeaseExpiresAt <= claimNow))) .ExecuteUpdateAsync(setters => setters .SetProperty(delivery => delivery.State, PushDeliveryStates.Sending) .SetProperty(delivery => delivery.LeaseId, leaseId) .SetProperty(delivery => delivery.LeaseExpiresAt, leaseUntil) .SetProperty(delivery => delivery.AttemptCount, delivery => delivery.AttemptCount + 1) .SetProperty(delivery => delivery.UpdatedAt, claimNow), ct); if (claimed != 1) continue; var delivery = await db.PushDeliveries .Include(item => item.PushMessage) .Include(item => item.PushDevice).ThenInclude(device => device.User) .FirstAsync(item => item.Id == id, ct); affectedMessages.Add(delivery.PushMessageId); await SendOne(db, registry, tokenProtector, delivery, ct); } foreach (var messageId in affectedMessages) await FinalizeMessage(db, messageId, ct); } private static async Task SendOne( AppDbContext db, PushProviderRegistry registry, PushTokenProtector tokenProtector, PushDelivery delivery, CancellationToken ct) { var now = DateTime.UtcNow; var message = delivery.PushMessage; var device = delivery.PushDevice; var expiresAt = (message.ScheduledAt ?? message.CreatedAt).AddSeconds(message.TtlSeconds); if (expiresAt <= now) { Skip(delivery, "message_expired", now); await db.SaveChangesAsync(ct); return; } if (!device.IsActive || !device.NotificationsAllowed || device.User.IsBanned || device.User.AccountClosureScheduledAt.HasValue) { Skip(delivery, "device_or_account_inactive", now); await db.SaveChangesAsync(ct); return; } if (!message.IsTest && !await db.UserPushPreferences.AnyAsync(preference => preference.UserId == device.UserId && preference.Category == message.Category && preference.IsEnabled, ct)) { Skip(delivery, "category_disabled", now); await db.SaveChangesAsync(ct); return; } var provider = registry.Find(device.Provider); if (provider is null) { Fail(delivery, "provider_unknown", "Unknown push provider", now); await db.SaveChangesAsync(ct); return; } var envelope = new PushEnvelope( message.PublicId, message.Title, message.Body, message.Category, message.Action, message.EntityId, Math.Max(60, (int)(expiresAt - now).TotalSeconds)); var result = await provider.SendAsync( device.Flavor, device.PackageName, tokenProtector.Unprotect(device.TokenCiphertext), envelope, ct); if (result.Accepted) { delivery.State = PushDeliveryStates.Accepted; delivery.AcceptedAt = now; delivery.ProviderMessageId = result.ProviderMessageId; delivery.ErrorCode = null; delivery.ErrorMessage = null; ClearLease(delivery, now); } else if (result.Retryable && delivery.AttemptCount <= RetrySchedule.Length && expiresAt > now) { var retry = result.RetryAfter ?? RetrySchedule[Math.Clamp(delivery.AttemptCount - 1, 0, RetrySchedule.Length - 1)]; delivery.State = PushDeliveryStates.Queued; delivery.NextAttemptAt = now.Add(retry) < expiresAt ? now.Add(retry) : expiresAt; delivery.ErrorCode = result.ErrorCode; delivery.ErrorMessage = result.ErrorMessage; ClearLease(delivery, now); } else { Fail(delivery, result.ErrorCode ?? "provider_rejected", result.ErrorMessage, now); } if (result.InvalidToken) { device.IsActive = false; device.DisabledReason = "provider_invalid_token"; device.UpdatedAt = now; } await db.SaveChangesAsync(ct); } private static async Task FinalizeMessage(AppDbContext db, long messageId, CancellationToken ct) { var pending = await db.PushDeliveries.AnyAsync(delivery => delivery.PushMessageId == messageId && (delivery.State == PushDeliveryStates.Queued || delivery.State == PushDeliveryStates.Sending), ct); if (pending) return; var failed = await db.PushDeliveries.AnyAsync(delivery => delivery.PushMessageId == messageId && delivery.State == PushDeliveryStates.Failed, ct); var now = DateTime.UtcNow; await db.PushMessages.Where(message => message.Id == messageId) .ExecuteUpdateAsync(setters => setters .SetProperty(message => message.State, failed ? PushMessageStates.PartiallyFailed : PushMessageStates.Completed) .SetProperty(message => message.CompletedAt, now) .SetProperty(message => message.UpdatedAt, now), ct); } private async Task Cleanup(CancellationToken ct) { using var scope = scopeFactory.CreateScope(); var db = scope.ServiceProvider.GetRequiredService(); var staleBefore = DateTime.UtcNow.AddDays(-90); await db.PushDevices .Where(device => device.IsActive && device.LastSeenAt < staleBefore) .ExecuteUpdateAsync(setters => setters .SetProperty(device => device.IsActive, false) .SetProperty(device => device.DisabledReason, "stale_device") .SetProperty(device => device.UpdatedAt, DateTime.UtcNow), ct); await db.PushDeliveries .Where(delivery => delivery.UpdatedAt < staleBefore && delivery.State != PushDeliveryStates.Queued && delivery.State != PushDeliveryStates.Sending) .ExecuteDeleteAsync(ct); } private static void Skip(PushDelivery delivery, string code, DateTime now) { delivery.State = PushDeliveryStates.Skipped; delivery.ErrorCode = code; delivery.ErrorMessage = null; ClearLease(delivery, now); } private static void Fail(PushDelivery delivery, string code, string? message, DateTime now) { delivery.State = PushDeliveryStates.Failed; delivery.ErrorCode = code; delivery.ErrorMessage = message is { Length: > 400 } ? message[..400] : message; ClearLease(delivery, now); } private static void ClearLease(PushDelivery delivery, DateTime now) { delivery.LeaseId = null; delivery.LeaseExpiresAt = null; delivery.UpdatedAt = now; } }