Add domestic vendor push infrastructure

This commit is contained in:
2026-07-26 01:45:59 +08:00
parent 7cca34b331
commit 0738953e6d
77 changed files with 6470 additions and 855 deletions
@@ -0,0 +1,358 @@
using MiaoJiZhang.Domain.Entities;
using MiaoJiZhang.Infrastructure.Persistence;
using Microsoft.EntityFrameworkCore;
namespace MiaoJiZhang.Api.Services;
public sealed class PushDispatchService(
IServiceScopeFactory scopeFactory,
IConfiguration configuration,
ILogger<PushDispatchService> 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<bool>("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<AppDbContext>();
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<AppDbContext>();
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<AppDbContext>();
var registry = scope.ServiceProvider.GetRequiredService<PushProviderRegistry>();
var tokenProtector = scope.ServiceProvider.GetRequiredService<PushTokenProtector>();
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<long>();
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<AppDbContext>();
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;
}
}