diff --git a/CLAUDE.md b/CLAUDE.md index 8af9316..2289cde 100644 --- a/CLAUDE.md +++ b/CLAUDE.md @@ -226,7 +226,9 @@ group www-data and reaches the private mongod; `sudo -u www-data` works too. marks (`Domain/Statuses/ConversationStates.cs`; writing in a conversation reads it). 10. **Nothing slow happens inside a request.** Deliveries and inbox processing are `Job`s (`Infrastructure/Jobs`): leased, retried on Mastodon's curve, at most two per host, paused per host by `RemoteInstance`. The inbox answers - 202 once it has verified and queued; a handler must be idempotent (unique `ObjectURI`, job `DedupeKey`). + 202 once it has verified and queued; a handler must be idempotent (unique `ObjectURI`, job `DedupeKey`). A lease + (2 minutes, stamped with its own owner) is renewed every third of its length while the handler runs, so a long job + runs once; a lease found taken cancels the handler, and `Finish` only counts for the lease it was given. 11. **Every number about posts and users goes through `Domain/Privacy/Counted`** (owner decision 2026-10-04: one answer everywhere): a persona's `statuses_count`, its outbox `totalItems`, NodeInfo `localPosts` and the instance `status_count` all count Mastodon's way (not deleted, not a DM, boosts included, circle and located posts too); diff --git a/PrivaPub.Tests/Infrastructure/JobQueueTests.cs b/PrivaPub.Tests/Infrastructure/JobQueueTests.cs index 6552cdd..e243d76 100644 --- a/PrivaPub.Tests/Infrastructure/JobQueueTests.cs +++ b/PrivaPub.Tests/Infrastructure/JobQueueTests.cs @@ -1,3 +1,5 @@ +using Microsoft.Extensions.Logging.Abstractions; + using MongoDB.Entities; using PrivaPub.Infrastructure.Jobs; @@ -74,6 +76,8 @@ namespace PrivaPub.Tests.Infrastructure Assert.Equal(JobState.Pending, after.State); Assert.True(after.RunAt > DateTime.UtcNow.AddSeconds(10)); + // leased again for its third attempt (a finish only counts for the lease it was given) + await DB.Default.Update().MatchID(leased.ID).Modify(j => j.State, JobState.Running).Modify(j => j.LeaseOwner, leased.LeaseOwner).ExecuteAsync(token); leased.Attempts = 3; await queue.Finish(leased, JobOutcome.Retry("boom"), maxAttempts: 3, token); Assert.Equal(JobState.Dead, (await DB.Default.Find().OneAsync(leased.ID, token)).State); @@ -90,5 +94,100 @@ namespace PrivaPub.Tests.Infrastructure Assert.Equal(JobState.Pending, (await DB.Default.Find().OneAsync(job.ID, token)).State); } + + // a handler running three times its lease keeps it: the reaper finds nothing to take back and it runs once + [Fact] + public async Task A_long_job_keeps_its_lease_and_runs_once() + { + var token = TestContext.Current.CancellationToken; + var kind = (JobKind)(2000 + Random.Shared.Next(100000)); + var queue = new JobQueue(TimeSpan.FromSeconds(1.5), j => j.Kind == kind); + var handler = new SlowHandler(kind, TimeSpan.FromSeconds(4.5)); + await queue.Enqueue(kind, "{}", default, default, token); + + using var worker = new JobWorker(queue, new IJobHandler[] { handler }, NullLogger.Instance); + await worker.StartAsync(token); + try + { + var deadline = DateTime.UtcNow.AddSeconds(20); + while (!await DB.Default.Find().Match(j => j.Kind == kind && j.State == JobState.Done).ExecuteAnyAsync(token)) + { + Assert.True(DateTime.UtcNow < deadline, "the job never finished"); + Assert.Equal(0, await queue.Reap(token)); + await Task.Delay(250, token); + } + } + finally + { + await worker.StopAsync(CancellationToken.None); + } + + Assert.Equal(1, handler.Started); + } + + // a lease taken by another worker (reaped, then leased again) stops the handler, and its outcome is not recorded + [Fact] + public async Task A_lost_lease_stops_the_handler_and_drops_its_outcome() + { + var token = TestContext.Current.CancellationToken; + var kind = (JobKind)(2000 + Random.Shared.Next(100000)); + var queue = new JobQueue(TimeSpan.FromSeconds(1.5), j => j.Kind == kind); + var handler = new SlowHandler(kind, TimeSpan.FromSeconds(30)); + await queue.Enqueue(kind, "{}", default, default, token); + + using var worker = new JobWorker(queue, new IJobHandler[] { handler }, NullLogger.Instance); + await worker.StartAsync(token); + try + { + var deadline = DateTime.UtcNow.AddSeconds(20); + while (handler.Started == 0) + { + Assert.True(DateTime.UtcNow < deadline, "the job never started"); + await Task.Delay(100, token); + } + await DB.Default.Update().Match(j => j.Kind == kind).Modify(j => j.LeaseOwner, "another worker").ExecuteAsync(token); + while (!handler.Cancelled) + { + Assert.True(DateTime.UtcNow < deadline, "the handler was never stopped"); + await Task.Delay(100, token); + } + await Task.Delay(300, token); + } + finally + { + await worker.StopAsync(CancellationToken.None); + } + + var job = await DB.Default.Find().Match(j => j.Kind == kind).ExecuteFirstAsync(token); + Assert.Equal(JobState.Running, job.State); + Assert.Equal("another worker", job.LeaseOwner); + } + + sealed class SlowHandler(JobKind kind, TimeSpan takes) : IJobHandler + { + int _started; + + public int Started => _started; + public bool Cancelled { get; private set; } + public JobKind Kind => kind; + public int Concurrency => 1; + public int MaxAttempts => 3; + public int PerHostLimit => 1; + + public async Task Handle(Job job, CancellationToken token) + { + Interlocked.Increment(ref _started); + try + { + await Task.Delay(takes, token); + } + catch (OperationCanceledException) + { + Cancelled = true; + throw; + } + return JobOutcome.Done; + } + } } } diff --git a/PrivaPub/Infrastructure/Jobs/JobQueue.cs b/PrivaPub/Infrastructure/Jobs/JobQueue.cs index 0978a73..f124124 100644 --- a/PrivaPub/Infrastructure/Jobs/JobQueue.cs +++ b/PrivaPub/Infrastructure/Jobs/JobQueue.cs @@ -1,3 +1,4 @@ +using MongoDB.Bson; using MongoDB.Driver; using MongoDB.Entities; @@ -29,8 +30,13 @@ namespace PrivaPub.Infrastructure.Jobs Task Enqueue(JobKind kind, string payload, string host, string dedupeKey, CancellationToken token); Task Payload(string dedupeKey, CancellationToken token); Task EnqueueMany(IEnumerable jobs, CancellationToken token); + /// How long a lease lasts unless it is renewed. + TimeSpan LeaseFor { get; } Task Lease(JobKind kind, IReadOnlyCollection busyHosts, CancellationToken token); - Task Finish(Job job, JobOutcome outcome, int maxAttempts, CancellationToken token); + /// Extends a running job's lease; false once it is no longer this lease's (reaped and leased again). + Task Renew(Job job, CancellationToken token); + /// Records the outcome; false when the lease was lost meanwhile, and then nothing changes. + Task Finish(Job job, JobOutcome outcome, int maxAttempts, CancellationToken token); Task Reap(CancellationToken token); Task WaitForWork(JobKind kind, TimeSpan poll, CancellationToken token); } @@ -39,6 +45,8 @@ namespace PrivaPub.Infrastructure.Jobs { public static readonly TimeSpan LeaseTime = TimeSpan.FromMinutes(2); + // every lease is stamped with its own owner (the process and a fresh id): two workers of one process, the second + // leasing a job the reaper took back from the first, must not renew nor finish each other's lease readonly string _owner = $"{Environment.MachineName}:{Environment.ProcessId}"; readonly ConcurrentDictionary _signals = new(); readonly FilterDefinition _scope = Builders.Filter.Empty; @@ -49,6 +57,15 @@ namespace PrivaPub.Infrastructure.Jobs public JobQueue(Expression> scope) => _scope = Builders.Filter.Where(scope); + public JobQueue(TimeSpan leaseFor, Expression> scope = default) + { + LeaseFor = leaseFor; + if (scope != default) + _scope = Builders.Filter.Where(scope); + } + + public TimeSpan LeaseFor { get; } = LeaseTime; + public async Task 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; @@ -81,17 +98,26 @@ namespace PrivaPub.Infrastructure.Jobs return await DB.Default.UpdateAndGet() .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(j => j.LeasedUntil, now + LeaseFor) + .Modify(j => j.LeaseOwner, $"{_owner}/{ObjectId.GenerateNewId()}") .Modify(b => b.Inc(j => j.Attempts, 1)) .Option(o => o.Sort = Builders.Sort.Ascending(j => j.RunAt)) .ExecuteAsync(token); } - public async Task Finish(Job job, JobOutcome outcome, int maxAttempts, CancellationToken token) + public async Task Renew(Job job, CancellationToken token) + { + var result = await DB.Default.Update() + .Match(j => j.ID == job.ID && j.State == JobState.Running && j.LeaseOwner == job.LeaseOwner) + .Modify(j => j.LeasedUntil, DateTime.UtcNow + LeaseFor) + .ExecuteAsync(token); + return result.MatchedCount == 1; + } + + public async Task Finish(Job job, JobOutcome outcome, int maxAttempts, CancellationToken token) { var now = DateTime.UtcNow; - var update = DB.Default.Update().MatchID(job.ID) + var update = DB.Default.Update().Match(j => j.ID == job.ID && j.LeaseOwner == job.LeaseOwner) .Modify(j => j.LeasedUntil, null) .Modify(j => j.LeaseOwner, null) .Modify(j => j.LastError, outcome.Error); @@ -112,7 +138,7 @@ namespace PrivaPub.Infrastructure.Jobs update.Modify(j => j.State, JobState.Dead).Modify(j => j.FinishedAt, now); break; } - await update.ExecuteAsync(token); + return (await update.ExecuteAsync(token)).MatchedCount == 1; } public async Task Reap(CancellationToken token) diff --git a/PrivaPub/Infrastructure/Jobs/JobWorker.cs b/PrivaPub/Infrastructure/Jobs/JobWorker.cs index b362307..a078609 100644 --- a/PrivaPub/Infrastructure/Jobs/JobWorker.cs +++ b/PrivaPub/Infrastructure/Jobs/JobWorker.cs @@ -69,18 +69,27 @@ namespace PrivaPub.Infrastructure.Jobs var host = job.Host ?? string.Empty; inFlight.AddOrUpdate(host, 1, (_, count) => count + 1); + // the lease is renewed while the handler runs; once it is lost (another worker leased the job again) the + // handler is cancelled and its outcome dropped + using var running = CancellationTokenSource.CreateLinkedTokenSource(stoppingToken); + var renewing = KeepLease(job, running); try { JobOutcome outcome; try { using var scope = HttpScope.Triggered(handler.Kind.ToString().ToLowerInvariant()); - outcome = await handler.Handle(job, stoppingToken); + outcome = await handler.Handle(job, running.Token); } catch (OperationCanceledException) when (stoppingToken.IsCancellationRequested) { return; } + catch (OperationCanceledException) when (running.IsCancellationRequested) + { + _logger.LogWarning("{Kind} job {Id} lost its lease and was stopped", handler.Kind, job.ID); + continue; + } catch (Exception ex) { _logger.LogError(ex, "{Kind} job {Id} threw", handler.Kind, job.ID); @@ -90,7 +99,8 @@ namespace PrivaPub.Infrastructure.Jobs 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); + if (!await _queue.Finish(job, outcome, handler.MaxAttempts, CancellationToken.None)) + _logger.LogWarning("{Kind} job {Id} lost its lease before it finished; its outcome is dropped", handler.Kind, job.ID); } catch (Exception ex) { @@ -98,11 +108,43 @@ namespace PrivaPub.Infrastructure.Jobs } finally { + await running.CancelAsync(); + await renewing; inFlight.AddOrUpdate(host, 0, (_, count) => count - 1); } } } + // renews the lease every third of its length until the handler ends; a lease found lost cancels the handler + async Task KeepLease(Job job, CancellationTokenSource running) + { + var every = _queue.LeaseFor / 3; + while (!running.IsCancellationRequested) + { + try + { + await Task.Delay(every, running.Token); + } + catch (OperationCanceledException) + { + return; + } + + try + { + if (await _queue.Renew(job, CancellationToken.None)) + continue; + await running.CancelAsync(); + return; + } + catch (Exception ex) + { + // a passing database error: the lease has two more thirds to go, so try again at the next turn + _logger.LogWarning(ex, "{Worker} could not renew the lease of job {Id}", nameof(JobWorker), job.ID); + } + } + } + async Task Reap(CancellationToken stoppingToken) { while (!stoppingToken.IsCancellationRequested)