64 lines
2.2 KiB
C#
64 lines
2.2 KiB
C#
using Microsoft.Extensions.DependencyInjection;
|
|||
|
|
|
||
|
|
using MongoDB.Entities;
|
||
|
|
|
||
|
|
using PrivaPub.Federation.Outbox;
|
||
|
|
using PrivaPub.Infrastructure.Jobs;
|
||
|
|
using PrivaPub.Models.Jobs;
|
||
|
|
|
||
|
|
using System.Linq.Expressions;
|
||
|
|
using System.Text.Json;
|
||
|
|
using System.Text.Json.Nodes;
|
||
|
|
|
||
|
|
namespace PrivaPub.Tests.Support.Host
|
||
|
|
{
|
||
|
|
public static class Jobs
|
||
|
|
{
|
||
|
|
public static async Task<int> Run(this PrivaPubHost host, Expression<Func<Job, bool>> which, CancellationToken token = default)
|
||
|
|
{
|
||
|
|
var handlers = host.Services.GetServices<IJobHandler>().ToDictionary(h => h.Kind);
|
||
|
|
var queue = host.Get<IJobQueue>();
|
||
|
|
var ran = 0;
|
||
|
|
for (var round = 0; round < 50; round++)
|
||
|
|
{
|
||
|
|
var pending = await DB.Default.Find<Job>().Match(which).Match(j => j.State == JobState.Pending).ExecuteAsync(token);
|
||
|
|
if (pending.Count == 0)
|
||
|
|
return ran;
|
||
|
|
foreach (var candidate in pending)
|
||
|
|
{
|
||
|
|
var job = await DB.Default.UpdateAndGet<Job>()
|
||
|
|
.Match(j => j.ID == candidate.ID && j.State == JobState.Pending)
|
||
|
|
.Modify(j => j.State, JobState.Running)
|
||
|
|
.Modify(j => j.LeasedUntil, DateTime.UtcNow + JobQueue.LeaseTime)
|
||
|
|
.Modify(b => b.Inc(j => j.Attempts, 1))
|
||
|
|
.ExecuteAsync(token);
|
||
|
|
if (job == default)
|
||
|
|
continue;
|
||
|
|
var handler = handlers[job.Kind];
|
||
|
|
JobOutcome outcome;
|
||
|
|
try
|
||
|
|
{
|
||
|
|
outcome = await handler.Handle(job, token);
|
||
|
|
}
|
||
|
|
catch (Exception ex)
|
||
|
|
{
|
||
|
|
outcome = JobOutcome.Dead(ex.GetType().Name + ": " + ex.Message);
|
||
|
|
}
|
||
|
|
await queue.Finish(job, outcome, handler.MaxAttempts, token);
|
||
|
|
ran++;
|
||
|
|
}
|
||
|
|
}
|
||
|
|
return ran;
|
||
|
|
}
|
||
|
|
|
||
|
|
public static Task<int> RunInbox(this PrivaPubHost host, string activityId, CancellationToken token = default) =>
|
||
|
|
host.Run(j => j.Kind == JobKind.ProcessInbox && j.DedupeKey == "inbox|" + activityId, token);
|
||
|
|
|
||
|
|
public static async Task<List<JsonObject>> Deliveries(string inbox, DateTime since, CancellationToken token = default) =>
|
||
|
|
(await DB.Default.Find<Job>().Match(j => j.Kind == JobKind.Deliver && j.CreatedAt >= since).ExecuteAsync(token))
|
||
|
|
.Select(j => JsonSerializer.Deserialize<DeliveryPayload>(j.Payload))
|
||
|
|
.Where(p => p.Inbox == inbox)
|
||
|
|
.Select(p => JsonNode.Parse(p.Body)!.AsObject())
|
||
|
|
.ToList();
|
||
|
|
}
|
||
|
|
}
|