Files
SocialPub/PrivaPub/Infrastructure/Jobs/JobWorker.cs
T
thepraandClaude Opus 5.5 d37999970c M5: outbound requests, previews, the media proxy and served traffic
- HttpScope (AsyncLocal) tags each outbound request with a purpose and a trigger:
  - purpose is set by the caller: actor, key, object, webfinger, context, nodeinfo;
  - trigger is set by the job kind, by "verify" during inbox verification, or defaults
    to "request".
- FederationHttp records every JSON, media and stream fetch: status, time, bytes, hops,
  and an outcome of ok, refused or failed, with a reason: disallowed, remembered,
  bad-redirect, too-many-redirects, content-type, too-large, bad-json, private-address,
  timeout, network, or the status. A fetch a reader caused (trigger "request") is only
  counted per server per day.
- Link previews record a 'preview' event: card, no-card or failed.
- The media proxy counts cache hits.
- TrafficMeter counts the client API per endpoint group, method and status class. It
  counts our served documents (actor, outbox, collection, object, activity, licence,
  webfinger, nodeinfo) by kind, status and whether signed, per day and never per server,
  and never names a circle's collections.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01ELjqpznMFMNrJoJUj6K5p2
2026-10-03 11:10:42 +02:00

140 lines
3.9 KiB
C#

using PrivaPub.Infrastructure.Http;
using PrivaPub.Models.Jobs;
using System.Collections.Concurrent;
namespace PrivaPub.Infrastructure.Jobs
{
public interface IJobHandler
{
JobKind Kind { get; }
int Concurrency { get; }
int MaxAttempts { get; }
int PerHostLimit { get; }
Task<JobOutcome> Handle(Job job, CancellationToken token);
}
public class JobWorker : BackgroundService
{
static readonly TimeSpan Poll = TimeSpan.FromSeconds(5);
static readonly TimeSpan ReapInterval = TimeSpan.FromSeconds(30);
readonly IJobQueue _queue;
readonly IEnumerable<IJobHandler> _handlers;
readonly ILogger<JobWorker> _logger;
public JobWorker(IJobQueue queue, IEnumerable<IJobHandler> handlers, ILogger<JobWorker> logger)
{
_queue = queue;
_handlers = handlers;
_logger = logger;
}
protected override async Task ExecuteAsync(CancellationToken stoppingToken)
{
var loops = _handlers.SelectMany(handler =>
{
var inFlight = new ConcurrentDictionary<string, int>(StringComparer.OrdinalIgnoreCase);
return Enumerable.Range(0, handler.Concurrency).Select(_ => Run(handler, inFlight, stoppingToken));
}).Append(Reap(stoppingToken));
await Task.WhenAll(loops);
}
async Task Run(IJobHandler handler, ConcurrentDictionary<string, int> inFlight, CancellationToken stoppingToken)
{
while (!stoppingToken.IsCancellationRequested)
{
Job job;
try
{
var busy = inFlight.Where(h => h.Value >= handler.PerHostLimit).Select(h => h.Key).ToList();
job = await _queue.Lease(handler.Kind, busy, stoppingToken);
}
catch (OperationCanceledException) when (stoppingToken.IsCancellationRequested)
{
return;
}
catch (Exception ex)
{
_logger.LogError(ex, "{Worker} could not lease a {Kind} job", nameof(JobWorker), handler.Kind);
await Delay(Poll, stoppingToken);
continue;
}
if (job == default)
{
await _queue.WaitForWork(handler.Kind, Poll, stoppingToken);
continue;
}
var host = job.Host ?? string.Empty;
inFlight.AddOrUpdate(host, 1, (_, count) => count + 1);
try
{
JobOutcome outcome;
try
{
using var scope = HttpScope.Triggered(handler.Kind.ToString().ToLowerInvariant());
outcome = await handler.Handle(job, stoppingToken);
}
catch (OperationCanceledException) when (stoppingToken.IsCancellationRequested)
{
return;
}
catch (Exception ex)
{
_logger.LogError(ex, "{Kind} job {Id} threw", handler.Kind, job.ID);
outcome = JobOutcome.Retry(ex.GetType().Name);
}
if (outcome.Result != JobResult.Done && job.Attempts >= handler.MaxAttempts && outcome.Result == JobResult.Retry)
_logger.LogWarning("{Kind} job {Id} for {Host} is dead after {Attempts} attempts: {Error}",
handler.Kind, job.ID, job.Host, job.Attempts, outcome.Error);
await _queue.Finish(job, outcome, handler.MaxAttempts, CancellationToken.None);
}
catch (Exception ex)
{
_logger.LogError(ex, "{Kind} job {Id} could not be finished; its lease will expire", handler.Kind, job.ID);
}
finally
{
inFlight.AddOrUpdate(host, 0, (_, count) => count - 1);
}
}
}
async Task Reap(CancellationToken stoppingToken)
{
while (!stoppingToken.IsCancellationRequested)
{
try
{
var reaped = await _queue.Reap(stoppingToken);
if (reaped > 0)
_logger.LogWarning("{Worker} returned {Count} expired leases to the queue", nameof(JobWorker), reaped);
}
catch (OperationCanceledException) when (stoppingToken.IsCancellationRequested)
{
return;
}
catch (Exception ex)
{
_logger.LogError(ex, "{Worker} reaper failed", nameof(JobWorker));
}
await Delay(ReapInterval, stoppingToken);
}
}
static async Task Delay(TimeSpan delay, CancellationToken token)
{
try
{
await Task.Delay(delay, token);
}
catch (OperationCanceledException)
{
}
}
}
}