Files
SocialPub/PrivaPub/Federation/Inbox/InboxProcessor.cs
T

63 lines
2.0 KiB
C#
Raw Normal View History

using PrivaPub.Federation.Actors;
using PrivaPub.Infrastructure.Jobs;
using PrivaPub.Models.Jobs;
using PrivaPub.Models.User;
using System.Text.Json;
using System.Text.Json.Nodes;
using static PrivaPub.Federation.Objects.ActivityJson;
namespace PrivaPub.Federation.Inbox
{
public interface IActivityHandler
{
string Type { get; }
Task Handle(JsonNode activity, ForeignAvatar actor, CancellationToken token);
}
public class InboxProcessor : IJobHandler
{
readonly IRemoteActorService _remoteActors;
readonly IReadOnlyDictionary<string, IActivityHandler> _handlers;
readonly ILogger<InboxProcessor> _logger;
public InboxProcessor(IRemoteActorService remoteActors, IEnumerable<IActivityHandler> handlers, ILogger<InboxProcessor> logger)
{
_remoteActors = remoteActors;
_handlers = handlers.ToDictionary(h => h.Type, StringComparer.Ordinal);
_logger = logger;
}
public JobKind Kind => JobKind.ProcessInbox;
public int Concurrency => 2;
public int MaxAttempts => 8;
public int PerHostLimit => 2;
public async Task<JobOutcome> Handle(Job job, CancellationToken token)
{
var payload = JsonSerializer.Deserialize<InboxPayload>(job.Payload);
var activity = payload == default ? default : JsonNode.Parse(payload.Activity);
var type = Value(activity, "type");
if (type == default || !_handlers.TryGetValue(type, out var handler))
return JobOutcome.Done;
var actor = await _remoteActors.GetActor(payload.ActorURI, refresh: false, token);
if (actor == default)
return JobOutcome.Retry("the actor could not be loaded");
Arrival.Current = new Arrival(Id(activity), type, actor.ActorURI, payload.Inbox, payload.KeyId, payload.Algorithm,
payload.SignedHeaders ?? Array.Empty<string>(), payload.ReceivedAt ?? job.CreatedAt, activity["@context"]?.ToJsonString());
try
{
await handler.Handle(activity, actor, token);
}
finally
{
Arrival.Current = default;
}
_logger.LogInformation("Processed {Type} {Id} from {Actor}", type, Id(activity), actor.ActorURI);
return JobOutcome.Done;
}
}
}