Arrival carries a verdict that handlers set with one line before an existing return (Arrival.Drop, Reject, Accept, About), so no handler signature changes. InboxProcessor records it as an 'in' event, with the time taken, the wait since the inbox accepted it, the attempt, the audience (from the post's visibility, or else the activity's addressing), the local actor kind, the object's age for updates, deletes and reactions, and the FEP features of a delivered object. It also records types nothing handles (unknown-type), actors that cannot be loaded (deferred) and handler failures. Drop reasons: fetch-failed, unparseable, misattributed, duplicate, deleted, not-addressed, not-followed, not-public, not-visible, not-deleted, unknown-object, unknown-recipient, cross-origin, unsupported. Accepted sub-reasons: stored, poll-vote, edit, refresh, actor-refresh, actor-delete, removed, tombstone-only, auto-accepted, pending, follow-answer, quote-answer, quote-granted, undone, reaction, reported. Rejected: blocked, ignored, quote-refused. A circle's traffic is private and kindless, and a stranger posting into one is just not-addressed, so the event store cannot reveal a circle. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01ELjqpznMFMNrJoJUj6K5p2
129 lines
4.7 KiB
C#
129 lines
4.7 KiB
C#
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<string, IActivityHandler> _handlers;
|
|
readonly ILogger<InboxProcessor> _logger;
|
|
readonly IInteractionLedger _ledger;
|
|
|
|
public InboxProcessor(IRemoteActorService remoteActors, IEnumerable<IActivityHandler> handlers, ILogger<InboxProcessor> 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<JobOutcome> Handle(Job job, CancellationToken token)
|
|
{
|
|
var started = Stopwatch.GetTimestamp();
|
|
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))
|
|
{
|
|
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<string>(), 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 ?? Audience(activity),
|
|
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);
|
|
}
|
|
|
|
static string Audience(JsonNode activity)
|
|
{
|
|
var source = activity?["to"] != default || activity?["cc"] != default ? activity : activity?["object"] as JsonObject;
|
|
var to = Addresses(source?["to"]);
|
|
var cc = Addresses(source?["cc"]);
|
|
if (to.Any(Addressing.IsPublic))
|
|
return Interactions.Public;
|
|
if (cc.Any(Addressing.IsPublic))
|
|
return Interactions.Unlisted;
|
|
return to.Count + cc.Count == 0 ? Interactions.None : Interactions.Private;
|
|
}
|
|
|
|
static List<string> Addresses(JsonNode node) => node switch
|
|
{
|
|
JsonArray array => array.Select(Id).Where(id => id != default).ToList(),
|
|
JsonNode single when Id(single) is { } id => new List<string> { id },
|
|
_ => new List<string>()
|
|
};
|
|
}
|
|
}
|