Friendica's activity ids are uniqid(): a short prefix and the microsecond. Two of its processes answering two follows at once gave both Accepts one id, and PrivaPub, queueing each inbox activity once per id, dropped the second as a copy: that follow stayed pending on our side while Friendica counted the persona as a follower (seen in the town's Friendica pair). An id that comes back carrying another type, actor or object is now queued apart; a true copy is still dropped. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01LsXgEaXee4GCU1hwYgPJXw
156 lines
5.1 KiB
C#
156 lines
5.1 KiB
C#
using MongoDB.Driver;
|
|
using MongoDB.Entities;
|
|
|
|
using PrivaPub.Models.Jobs;
|
|
|
|
using System.Collections.Concurrent;
|
|
using System.Linq.Expressions;
|
|
|
|
namespace PrivaPub.Infrastructure.Jobs
|
|
{
|
|
public sealed record JobOutcome(JobResult Result, string Error = default, DateTime? RetryAt = default)
|
|
{
|
|
public static readonly JobOutcome Done = new(JobResult.Done);
|
|
public static JobOutcome Retry(string error) => new(JobResult.Retry, error);
|
|
public static JobOutcome Dead(string error) => new(JobResult.Dead, error);
|
|
public static JobOutcome Defer(DateTime until, string error) => new(JobResult.Defer, error, until);
|
|
}
|
|
|
|
public enum JobResult
|
|
{
|
|
Done,
|
|
Retry,
|
|
Dead,
|
|
Defer
|
|
}
|
|
|
|
public interface IJobQueue
|
|
{
|
|
Task<bool> Enqueue(JobKind kind, string payload, string host, string dedupeKey, CancellationToken token);
|
|
Task<string> Payload(string dedupeKey, CancellationToken token);
|
|
Task<int> EnqueueMany(IEnumerable<Job> jobs, CancellationToken token);
|
|
Task<Job> Lease(JobKind kind, IReadOnlyCollection<string> busyHosts, CancellationToken token);
|
|
Task Finish(Job job, JobOutcome outcome, int maxAttempts, CancellationToken token);
|
|
Task<long> Reap(CancellationToken token);
|
|
Task WaitForWork(JobKind kind, TimeSpan poll, CancellationToken token);
|
|
}
|
|
|
|
public class JobQueue : IJobQueue
|
|
{
|
|
public static readonly TimeSpan LeaseTime = TimeSpan.FromMinutes(2);
|
|
|
|
readonly string _owner = $"{Environment.MachineName}:{Environment.ProcessId}";
|
|
readonly ConcurrentDictionary<JobKind, SemaphoreSlim> _signals = new();
|
|
readonly FilterDefinition<Job> _scope = Builders<Job>.Filter.Empty;
|
|
|
|
public JobQueue()
|
|
{
|
|
}
|
|
|
|
public JobQueue(Expression<Func<Job, bool>> scope) => _scope = Builders<Job>.Filter.Where(scope);
|
|
|
|
public async Task<bool> Enqueue(JobKind kind, string payload, string host, string dedupeKey, CancellationToken token) =>
|
|
await EnqueueMany(new[] { new Job { Kind = kind, Payload = payload, Host = host, DedupeKey = dedupeKey } }, token) == 1;
|
|
|
|
// what the job a dedupe key names carries, or null when there is none
|
|
public async Task<string> Payload(string dedupeKey, CancellationToken token) =>
|
|
(await DB.Default.Find<Job>().Match(j => j.DedupeKey == dedupeKey).ExecuteFirstAsync(token))?.Payload;
|
|
|
|
public async Task<int> EnqueueMany(IEnumerable<Job> jobs, CancellationToken token)
|
|
{
|
|
var inserted = 0;
|
|
foreach (var job in jobs)
|
|
{
|
|
try
|
|
{
|
|
await DB.Default.SaveAsync(job, token);
|
|
inserted++;
|
|
Signal(job.Kind);
|
|
}
|
|
catch (MongoWriteException ex) when (ex.WriteError?.Category == ServerErrorCategory.DuplicateKey)
|
|
{
|
|
}
|
|
}
|
|
return inserted;
|
|
}
|
|
|
|
public async Task<Job> Lease(JobKind kind, IReadOnlyCollection<string> busyHosts, CancellationToken token)
|
|
{
|
|
var now = DateTime.UtcNow;
|
|
var busy = busyHosts?.ToList() ?? new List<string>();
|
|
return await DB.Default.UpdateAndGet<Job>()
|
|
.Match(f => f.Where(j => j.Kind == kind && j.State == JobState.Pending && j.RunAt <= now && !busy.Contains(j.Host)) & _scope)
|
|
.Modify(j => j.State, JobState.Running)
|
|
.Modify(j => j.LeasedUntil, now + LeaseTime)
|
|
.Modify(j => j.LeaseOwner, _owner)
|
|
.Modify(b => b.Inc(j => j.Attempts, 1))
|
|
.Option(o => o.Sort = Builders<Job>.Sort.Ascending(j => j.RunAt))
|
|
.ExecuteAsync(token);
|
|
}
|
|
|
|
public async Task Finish(Job job, JobOutcome outcome, int maxAttempts, CancellationToken token)
|
|
{
|
|
var now = DateTime.UtcNow;
|
|
var update = DB.Default.Update<Job>().MatchID(job.ID)
|
|
.Modify(j => j.LeasedUntil, null)
|
|
.Modify(j => j.LeaseOwner, null)
|
|
.Modify(j => j.LastError, outcome.Error);
|
|
switch (outcome.Result)
|
|
{
|
|
case JobResult.Done:
|
|
update.Modify(j => j.State, JobState.Done).Modify(j => j.FinishedAt, now);
|
|
break;
|
|
case JobResult.Retry when job.Attempts < maxAttempts:
|
|
update.Modify(j => j.State, JobState.Pending).Modify(j => j.RunAt, now + Backoff.After(job.Attempts));
|
|
break;
|
|
case JobResult.Defer:
|
|
update.Modify(j => j.State, JobState.Pending)
|
|
.Modify(j => j.RunAt, outcome.RetryAt ?? now + Backoff.After(job.Attempts))
|
|
.Modify(b => b.Inc(j => j.Attempts, -1));
|
|
break;
|
|
default:
|
|
update.Modify(j => j.State, JobState.Dead).Modify(j => j.FinishedAt, now);
|
|
break;
|
|
}
|
|
await update.ExecuteAsync(token);
|
|
}
|
|
|
|
public async Task<long> Reap(CancellationToken token)
|
|
{
|
|
var now = DateTime.UtcNow;
|
|
var result = await DB.Default.Update<Job>()
|
|
.Match(f => f.Where(j => j.State == JobState.Running && j.LeasedUntil < now) & _scope)
|
|
.Modify(j => j.State, JobState.Pending)
|
|
.Modify(j => j.RunAt, now)
|
|
.Modify(j => j.LeasedUntil, null)
|
|
.Modify(j => j.LeaseOwner, null)
|
|
.ExecuteAsync(token);
|
|
return result.ModifiedCount;
|
|
}
|
|
|
|
public async Task WaitForWork(JobKind kind, TimeSpan poll, CancellationToken token)
|
|
{
|
|
try
|
|
{
|
|
await SignalFor(kind).WaitAsync(poll, token);
|
|
}
|
|
catch (OperationCanceledException) when (token.IsCancellationRequested)
|
|
{
|
|
}
|
|
}
|
|
|
|
void Signal(JobKind kind)
|
|
{
|
|
try
|
|
{
|
|
SignalFor(kind).Release();
|
|
}
|
|
catch (SemaphoreFullException)
|
|
{
|
|
}
|
|
}
|
|
|
|
SemaphoreSlim SignalFor(JobKind kind) => _signals.GetOrAdd(kind, _ => new SemaphoreSlim(0, 64));
|
|
}
|
|
}
|