using PrivaPub.Federation.Actors; using PrivaPub.Federation.Objects; using PrivaPub.Infrastructure.Jobs; using PrivaPub.Infrastructure.Statistics; using PrivaPub.Models.Jobs; using PrivaPub.Models.Statistics; using PrivaPub.Models.User; using System.Diagnostics; 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 _handlers; readonly ILogger _logger; readonly IInteractionLedger _ledger; public InboxProcessor(IRemoteActorService remoteActors, IEnumerable handlers, ILogger logger, IInteractionLedger ledger = default) { _remoteActors = remoteActors; _handlers = handlers.ToDictionary(h => h.Type, StringComparer.Ordinal); _logger = logger; _ledger = ledger; } public JobKind Kind => JobKind.ProcessInbox; public int Concurrency => 2; public int MaxAttempts => 8; public int PerHostLimit => 2; public async Task Handle(Job job, CancellationToken token) { var started = Stopwatch.GetTimestamp(); var payload = JsonSerializer.Deserialize(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)) { Record(job, payload, activity, type, started, new ArrivalVerdict(), Interactions.Dropped, "unknown-type"); return JobOutcome.Done; } var actor = await _remoteActors.GetActor(payload.ActorURI, refresh: false, token); if (actor == default) { Record(job, payload, activity, type, started, new ArrivalVerdict(), Interactions.Deferred, "actor-unavailable"); return JobOutcome.Retry("the actor could not be loaded"); } var arrival = new Arrival(Id(activity), type, actor.ActorURI, payload.Inbox, payload.KeyId, payload.Algorithm, payload.SignedHeaders ?? Array.Empty(), payload.ReceivedAt ?? job.CreatedAt, activity["@context"]?.ToJsonString()); Arrival.Current = arrival; try { await handler.Handle(activity, actor, token); } catch (Exception ex) when (ex is not OperationCanceledException) { Record(job, payload, activity, type, started, arrival.Verdict, Interactions.Failed, ex.GetType().Name.ToLowerInvariant()); throw; } finally { Arrival.Current = default; } Record(job, payload, activity, type, started, arrival.Verdict, default, default); _logger.LogInformation("Processed {Type} {Id} from {Actor}", type, Id(activity), actor.ActorURI); return JobOutcome.Done; } 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, Audience = verdict.Audience ?? ActivityShape.Of(activity).Audience, 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); } } }