Deliveries run on a Mongo job queue with leases, backoff and per-host limits
The single serial DeliveryWorker is replaced by Infrastructure/Jobs: - Job rows are leased with one FindOneAndUpdate (oldest RunAt first, a two-minute lease) and a reaper returns expired leases every 30 s; - enqueueing wakes the workers, which otherwise poll every five seconds; - delivery runs eight at a time with at most two per host, so a slow or dead server holds two slots, not the queue; - a failure waits n^4 + 15 + jitter seconds (Mastodon's curve) for up to 16 attempts; a 4xx other than 408/429 is final, a 429 honours Retry-After; - RemoteInstance is a per-host circuit breaker: ten consecutive failures quarantine a host for an hour, doubling to a week, and its jobs wait without spending attempts; - a delivery is queued once per activity and inbox (unique DedupeKey), and finished jobs expire after seven days (TTL on FinishedAt). Migration _004 moves pending Delivery rows into jobs and marks them abandoned, so a rollback to the old worker cannot send them twice. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_012CzABvBkbcFqoHdmi8b9WB
This commit is contained in:
1 parent
934b6fe687
commit
d1a91c40c4
15 files changed
+704
-99
No files matched your search
@@ -10,6 +10,10 @@ using System.Text.Json.Nodes;
|
||||
using PrivaPub.Federation.Actors;
|
||||
using PrivaPub.Federation.Signing;
|
||||
using PrivaPub.Infrastructure.Http;
|
||||
using PrivaPub.Infrastructure.Jobs;
|
||||
using PrivaPub.Models.Jobs;
|
||||
|
||||
using System.Text.Json;
|
||||
|
||||
namespace PrivaPub.Federation.Outbox
|
||||
{
|
||||
@@ -20,31 +24,38 @@ namespace PrivaPub.Federation.Outbox
|
||||
Task<IReadOnlyList<string>> FollowerInboxes(LocalActor actor, CancellationToken token);
|
||||
}
|
||||
|
||||
public sealed record DeliveryPayload(string SignerId, LocalActorKind SignerKind, string Inbox, string Body);
|
||||
|
||||
public class DeliveryService : IDeliveryService
|
||||
{
|
||||
readonly DbEntities _dbEntities;
|
||||
readonly IJobQueue _queue;
|
||||
|
||||
public DeliveryService(DbEntities dbEntities)
|
||||
public DeliveryService(DbEntities dbEntities, IJobQueue queue)
|
||||
{
|
||||
_dbEntities = dbEntities;
|
||||
_queue = queue;
|
||||
}
|
||||
|
||||
public async Task Enqueue(LocalActor signer, IEnumerable<string> inboxes, JsonObject activity, CancellationToken token)
|
||||
{
|
||||
var body = activity.ToJsonString();
|
||||
var deliveries = inboxes
|
||||
var activityId = activity["id"] is JsonValue id && id.TryGetValue<string>(out var text) ? text : default;
|
||||
var jobs = inboxes
|
||||
.Where(i => !string.IsNullOrEmpty(i) && !i.StartsWith(signer.BaseAddress + "/", StringComparison.OrdinalIgnoreCase))
|
||||
.Distinct(StringComparer.Ordinal)
|
||||
.Select(inbox => new Delivery
|
||||
.Select(inbox => Uri.TryCreate(inbox, UriKind.Absolute, out var uri) ? (inbox, uri) : default)
|
||||
.Where(target => target.uri != default)
|
||||
.Select(target => new Job
|
||||
{
|
||||
SignerId = signer.Id,
|
||||
SignerKind = signer.Kind,
|
||||
InboxURL = inbox,
|
||||
Body = body
|
||||
Kind = JobKind.Deliver,
|
||||
Host = target.uri.Host.ToLowerInvariant(),
|
||||
DedupeKey = activityId == default ? default : $"{activityId}|{target.inbox}",
|
||||
Payload = JsonSerializer.Serialize(new DeliveryPayload(signer.Id, signer.Kind, target.inbox, body))
|
||||
})
|
||||
.ToList();
|
||||
if (deliveries.Count > 0)
|
||||
await DB.Default.SaveAsync(deliveries, token);
|
||||
if (jobs.Count > 0)
|
||||
await _queue.EnqueueMany(jobs, token);
|
||||
}
|
||||
|
||||
public async Task EnqueueToFollowers(LocalActor signer, JsonObject activity, CancellationToken token, IEnumerable<string> extraInboxes = default)
|
||||
@@ -64,118 +75,81 @@ namespace PrivaPub.Federation.Outbox
|
||||
}
|
||||
}
|
||||
|
||||
public class DeliveryWorker : BackgroundService
|
||||
public class DeliveryJobHandler : IJobHandler
|
||||
{
|
||||
const int MaxAttempts = 8;
|
||||
static readonly TimeSpan Poll = TimeSpan.FromSeconds(3);
|
||||
|
||||
readonly IServiceProvider _services;
|
||||
readonly ILocalActorService _actors;
|
||||
readonly IFederationHttp _http;
|
||||
readonly ILogger<DeliveryWorker> _logger;
|
||||
readonly IHostCircuitBreaker _breaker;
|
||||
readonly ILogger<DeliveryJobHandler> _logger;
|
||||
|
||||
public DeliveryWorker(IServiceProvider services, IFederationHttp http, ILogger<DeliveryWorker> logger)
|
||||
public DeliveryJobHandler(ILocalActorService actors, IFederationHttp http, IHostCircuitBreaker breaker, ILogger<DeliveryJobHandler> logger)
|
||||
{
|
||||
_services = services;
|
||||
_actors = actors;
|
||||
_http = http;
|
||||
_breaker = breaker;
|
||||
_logger = logger;
|
||||
}
|
||||
|
||||
protected override async Task ExecuteAsync(CancellationToken stoppingToken)
|
||||
{
|
||||
while (!stoppingToken.IsCancellationRequested)
|
||||
{
|
||||
try
|
||||
{
|
||||
await DeliverDue(stoppingToken);
|
||||
}
|
||||
catch (OperationCanceledException) when (stoppingToken.IsCancellationRequested)
|
||||
{
|
||||
return;
|
||||
}
|
||||
catch (Exception ex)
|
||||
{
|
||||
_logger.LogError(ex, "{Worker} pass failed", nameof(DeliveryWorker));
|
||||
}
|
||||
await Task.Delay(Poll, stoppingToken);
|
||||
}
|
||||
}
|
||||
public JobKind Kind => JobKind.Deliver;
|
||||
public int Concurrency => 8;
|
||||
public int MaxAttempts => 16;
|
||||
public int PerHostLimit => 2;
|
||||
|
||||
async Task DeliverDue(CancellationToken token)
|
||||
public async Task<JobOutcome> Handle(Job job, CancellationToken token)
|
||||
{
|
||||
using var scope = _services.CreateScope();
|
||||
var dbEntities = scope.ServiceProvider.GetRequiredService<DbEntities>();
|
||||
var actors = scope.ServiceProvider.GetRequiredService<ILocalActorService>();
|
||||
var now = DateTime.UtcNow;
|
||||
var due = await dbEntities.Deliveries
|
||||
.Match(d => !d.DeliveredAt.HasValue && !d.AbandonedAt.HasValue && d.NextAttemptAt <= now)
|
||||
.Sort(d => d.NextAttemptAt, Order.Ascending)
|
||||
.Limit(20)
|
||||
.ExecuteAsync(token);
|
||||
var payload = JsonSerializer.Deserialize<DeliveryPayload>(job.Payload);
|
||||
if (payload == default || !Uri.TryCreate(payload.Inbox, UriKind.Absolute, out var inbox) || !_http.IsAllowed(inbox))
|
||||
return JobOutcome.Dead("not a deliverable inbox");
|
||||
|
||||
foreach (var delivery in due)
|
||||
{
|
||||
var signer = await actors.FindById(delivery.SignerKind, delivery.SignerId, token);
|
||||
if (signer == default)
|
||||
{
|
||||
delivery.AbandonedAt = DateTime.UtcNow;
|
||||
delivery.LastError = "the signing actor no longer exists";
|
||||
await DB.Default.SaveAsync(delivery, token);
|
||||
continue;
|
||||
}
|
||||
await Deliver(delivery, signer, token);
|
||||
await DB.Default.SaveAsync(delivery, token);
|
||||
}
|
||||
}
|
||||
var unavailableUntil = await _breaker.UnavailableUntil(job.Host, token);
|
||||
if (unavailableUntil.HasValue)
|
||||
return JobOutcome.Defer(unavailableUntil.Value, "the host is unavailable");
|
||||
|
||||
var signer = await _actors.FindById(payload.SignerKind, payload.SignerId, token);
|
||||
if (signer == default)
|
||||
return JobOutcome.Dead("the signing actor no longer exists");
|
||||
|
||||
var body = Encoding.UTF8.GetBytes(payload.Body);
|
||||
using var request = new HttpRequestMessage(HttpMethod.Post, inbox) { Content = new ByteArrayContent(body) };
|
||||
request.Content.Headers.ContentType = MediaTypeHeaderValue.Parse(RemoteActorService.ActivityJson);
|
||||
HttpSignatures.Sign(request, signer, body);
|
||||
|
||||
async Task Deliver(Delivery delivery, LocalActor signer, CancellationToken token)
|
||||
{
|
||||
delivery.Attempts++;
|
||||
try
|
||||
{
|
||||
if (!Uri.TryCreate(delivery.InboxURL, UriKind.Absolute, out var inbox) || !_http.IsAllowed(inbox))
|
||||
{
|
||||
delivery.AbandonedAt = DateTime.UtcNow;
|
||||
delivery.LastError = "not a deliverable inbox";
|
||||
return;
|
||||
}
|
||||
|
||||
var body = Encoding.UTF8.GetBytes(delivery.Body);
|
||||
using var request = new HttpRequestMessage(HttpMethod.Post, inbox)
|
||||
{
|
||||
Content = new ByteArrayContent(body)
|
||||
};
|
||||
request.Content.Headers.ContentType = MediaTypeHeaderValue.Parse(RemoteActorService.ActivityJson);
|
||||
HttpSignatures.Sign(request, signer, body);
|
||||
|
||||
using var response = await _http.Send(request, token);
|
||||
var status = (int)response.StatusCode;
|
||||
if (response.IsSuccessStatusCode)
|
||||
{
|
||||
delivery.DeliveredAt = DateTime.UtcNow;
|
||||
delivery.LastError = default;
|
||||
return;
|
||||
await _breaker.Succeeded(job.Host, token);
|
||||
return JobOutcome.Done;
|
||||
}
|
||||
|
||||
delivery.LastError = $"{(int)response.StatusCode} {response.ReasonPhrase}";
|
||||
if (response.StatusCode is HttpStatusCode.Gone or HttpStatusCode.NotFound or HttpStatusCode.BadRequest or HttpStatusCode.Forbidden)
|
||||
if (response.StatusCode == HttpStatusCode.TooManyRequests)
|
||||
{
|
||||
delivery.AbandonedAt = DateTime.UtcNow;
|
||||
_logger.LogWarning("Delivery {Id} to {Inbox} abandoned: {Status}", delivery.ID, delivery.InboxURL, delivery.LastError);
|
||||
return;
|
||||
var retryAfter = response.Headers.RetryAfter?.Delta ?? (response.Headers.RetryAfter?.Date - DateTimeOffset.UtcNow);
|
||||
return retryAfter > TimeSpan.Zero
|
||||
? JobOutcome.Defer(DateTime.UtcNow + Min(retryAfter.Value, TimeSpan.FromHours(6)), "429")
|
||||
: JobOutcome.Retry("429");
|
||||
}
|
||||
if (status is >= 400 and < 500 && status != 408)
|
||||
{
|
||||
await _breaker.Succeeded(job.Host, token);
|
||||
_logger.LogInformation("Delivery to {Inbox} refused with {Status}", payload.Inbox, status);
|
||||
return JobOutcome.Dead($"{status} {response.ReasonPhrase}");
|
||||
}
|
||||
await _breaker.Failed(job.Host, $"{status}", token);
|
||||
return JobOutcome.Retry($"{status} {response.ReasonPhrase}");
|
||||
}
|
||||
catch (Exception ex) when (ex is HttpRequestException or BlockedDestinationException or TaskCanceledException && !token.IsCancellationRequested)
|
||||
catch (BlockedDestinationException ex)
|
||||
{
|
||||
delivery.LastError = ex.Message;
|
||||
return JobOutcome.Dead(ex.Message);
|
||||
}
|
||||
|
||||
if (delivery.Attempts >= MaxAttempts)
|
||||
catch (Exception ex) when (ex is HttpRequestException or TaskCanceledException && !token.IsCancellationRequested)
|
||||
{
|
||||
delivery.AbandonedAt = DateTime.UtcNow;
|
||||
_logger.LogWarning("Delivery {Id} to {Inbox} abandoned after {Attempts} attempts: {Error}",
|
||||
delivery.ID, delivery.InboxURL, delivery.Attempts, delivery.LastError);
|
||||
return;
|
||||
await _breaker.Failed(job.Host, ex.GetType().Name, token);
|
||||
return JobOutcome.Retry(ex.Message);
|
||||
}
|
||||
delivery.NextAttemptAt = DateTime.UtcNow.AddMinutes(Math.Pow(2, delivery.Attempts));
|
||||
}
|
||||
|
||||
static TimeSpan Min(TimeSpan a, TimeSpan b) => a < b ? a : b;
|
||||
}
|
||||
}
|
||||
Reference in new issue
Block a user