2026-10-01 11:12:34 +02:00
|
|
|
using PrivaPub.Federation.Actors;
|
2026-10-03 11:04:42 +02:00
|
|
|
using PrivaPub.Federation.Objects;
|
2026-10-01 11:12:34 +02:00
|
|
|
using PrivaPub.Infrastructure.Jobs;
|
2026-10-03 11:04:42 +02:00
|
|
|
using PrivaPub.Infrastructure.Statistics;
|
2026-10-01 11:12:34 +02:00
|
|
|
using PrivaPub.Models.Jobs;
|
2026-10-03 11:04:42 +02:00
|
|
|
using PrivaPub.Models.Statistics;
|
2026-10-01 11:12:34 +02:00
|
|
|
using PrivaPub.Models.User;
|
|
|
|
|
|
2026-10-03 11:04:42 +02:00
|
|
|
using System.Diagnostics;
|
2026-10-01 11:12:34 +02:00
|
|
|
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;
|
2026-10-03 11:04:42 +02:00
|
|
|
readonly IInteractionLedger _ledger;
|
2026-10-01 11:12:34 +02:00
|
|
|
|
2026-10-03 11:04:42 +02:00
|
|
|
public InboxProcessor(IRemoteActorService remoteActors, IEnumerable<IActivityHandler> handlers, ILogger<InboxProcessor> logger,
|
|
|
|
|
IInteractionLedger ledger = default)
|
2026-10-01 11:12:34 +02:00
|
|
|
{
|
|
|
|
|
_remoteActors = remoteActors;
|
|
|
|
|
_handlers = handlers.ToDictionary(h => h.Type, StringComparer.Ordinal);
|
|
|
|
|
_logger = logger;
|
2026-10-03 11:04:42 +02:00
|
|
|
_ledger = ledger;
|
2026-10-01 11:12:34 +02:00
|
|
|
}
|
|
|
|
|
|
|
|
|
|
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)
|
|
|
|
|
{
|
2026-10-03 11:04:42 +02:00
|
|
|
var started = Stopwatch.GetTimestamp();
|
2026-10-01 11:12:34 +02:00
|
|
|
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))
|
2026-10-03 11:04:42 +02:00
|
|
|
{
|
|
|
|
|
Record(job, payload, activity, type, started, new ArrivalVerdict(), Interactions.Dropped, "unknown-type");
|
2026-10-01 11:12:34 +02:00
|
|
|
return JobOutcome.Done;
|
2026-10-03 11:04:42 +02:00
|
|
|
}
|
2026-10-01 11:12:34 +02:00
|
|
|
|
|
|
|
|
var actor = await _remoteActors.GetActor(payload.ActorURI, refresh: false, token);
|
|
|
|
|
if (actor == default)
|
2026-10-03 11:04:42 +02:00
|
|
|
{
|
|
|
|
|
Record(job, payload, activity, type, started, new ArrivalVerdict(), Interactions.Deferred, "actor-unavailable");
|
2026-10-01 11:12:34 +02:00
|
|
|
return JobOutcome.Retry("the actor could not be loaded");
|
2026-10-03 11:04:42 +02:00
|
|
|
}
|
2026-10-01 11:12:34 +02:00
|
|
|
|
2026-10-03 11:04:42 +02:00
|
|
|
var arrival = new Arrival(Id(activity), type, actor.ActorURI, payload.Inbox, payload.KeyId, payload.Algorithm,
|
2026-10-01 17:49:22 +02:00
|
|
|
payload.SignedHeaders ?? Array.Empty<string>(), payload.ReceivedAt ?? job.CreatedAt, activity["@context"]?.ToJsonString());
|
2026-10-03 11:04:42 +02:00
|
|
|
Arrival.Current = arrival;
|
2026-10-01 17:49:22 +02:00
|
|
|
try
|
|
|
|
|
{
|
|
|
|
|
await handler.Handle(activity, actor, token);
|
|
|
|
|
}
|
2026-10-03 11:04:42 +02:00
|
|
|
catch (Exception ex) when (ex is not OperationCanceledException)
|
|
|
|
|
{
|
|
|
|
|
Record(job, payload, activity, type, started, arrival.Verdict, Interactions.Failed, ex.GetType().Name.ToLowerInvariant());
|
|
|
|
|
throw;
|
|
|
|
|
}
|
2026-10-01 17:49:22 +02:00
|
|
|
finally
|
|
|
|
|
{
|
|
|
|
|
Arrival.Current = default;
|
|
|
|
|
}
|
2026-10-03 11:04:42 +02:00
|
|
|
Record(job, payload, activity, type, started, arrival.Verdict, default, default);
|
2026-10-01 11:12:34 +02:00
|
|
|
_logger.LogInformation("Processed {Type} {Id} from {Actor}", type, Id(activity), actor.ActorURI);
|
|
|
|
|
return JobOutcome.Done;
|
|
|
|
|
}
|
2026-10-03 11:04:42 +02:00
|
|
|
|
|
|
|
|
void Record(Job job, InboxPayload payload, JsonNode activity, string type, long started, ArrivalVerdict verdict, string outcome, string reason)
|
|
|
|
|
{
|
|
|
|
|
if (_ledger == default || payload == default)
|
|
|
|
|
return;
|
|
|
|
|
var embedded = activity?["object"] as JsonObject;
|
|
|
|
|
_ledger.Record(new InteractionEvent
|
|
|
|
|
{
|
|
|
|
|
Channel = Interactions.In,
|
|
|
|
|
Host = Interactions.HostOf(payload.ActorURI),
|
|
|
|
|
Activity = type,
|
|
|
|
|
Object = verdict.ObjectType ?? Value(embedded, "type"),
|
|
|
|
|
Outcome = outcome ?? verdict.Outcome ?? Interactions.Accepted,
|
|
|
|
|
Reason = reason ?? verdict.Reason,
|
|
|
|
|
LatencyMs = (int)Stopwatch.GetElapsedTime(started).TotalMilliseconds,
|
|
|
|
|
WaitMs = (int)Math.Min(int.MaxValue, Math.Max(0, (DateTime.UtcNow - (payload.ReceivedAt ?? job.CreatedAt)).TotalMilliseconds)),
|
|
|
|
|
Attempt = job.Attempts,
|
2026-10-03 11:06:24 +02:00
|
|
|
Audience = verdict.Audience ?? ActivityShape.Of(activity).Audience,
|
2026-10-03 11:04:42 +02:00
|
|
|
LocalKind = verdict.LocalKind,
|
|
|
|
|
AgeSeconds = verdict.AgeSeconds,
|
|
|
|
|
Inbox = payload.Inbox,
|
|
|
|
|
Signature = payload.Algorithm == default ? default : "cavage:" + payload.Algorithm,
|
|
|
|
|
Features = type is "Create" or "Update" && embedded != default ? ObjectFeatures.Detect(embedded) : default
|
|
|
|
|
}, payload.ActorURI);
|
|
|
|
|
}
|
2026-10-01 11:12:34 +02:00
|
|
|
}
|
|
|
|
|
}
|