Files
SocialPub/PrivaPub/Federation/Inbox/InboxProcessor.cs
T
thepraandClaude Opus 5.5 9ab87b2779 Personas prove what goes to relays; the server says what it reads
FEP-521a and FEP-8b32. Every persona has an Ed25519 key of its own (Avatar.SigningKey; migration 014 gives the earlier
ones theirs), named in its actor's assertionMethod as a Multikey, the terms defined in the actor's own context. A
persona's activity going to a relay carries an eddsa-jcs-2022 proof (JSON canonicalised by RFC 8785, Jcs), so what
Activity-Relay forwards reaches Mastodon, which verifies it with its own code. Nothing else carries one: Mitra takes a
proof over the HTTP signature and refuses one by a key it has not read, without reading the actor again. Received: an
actor's own Multikeys are kept, and a forwarded activity whose proof one of them verifies is taken as it came instead of
being read again from its origin.

Discovery: WebFinger for the server's origin links its instance actor (FEP-d556), NodeInfo links it as the application
actor (FEP-2677), and actors name RFC 9421 under implements (FEP-844e).

Checked live: relay 16 (Activity-Relay's forward of alice's post reaches Mastodon), Mitra, GoToSocial and Mastodon
unchanged (165 in all).

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01LsXgEaXee4GCU1hwYgPJXw
2026-10-06 09:01:21 +02:00

130 lines
5.1 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;
readonly ILocalActorService _localActors;
readonly int _concurrency;
public InboxProcessor(IRemoteActorService remoteActors, IEnumerable<IActivityHandler> handlers, ILogger<InboxProcessor> logger,
IInteractionLedger ledger = default, ILocalActorService localActors = default,
Microsoft.Extensions.Options.IOptions<Infrastructure.Http.FederationOptions> federation = default)
{
_concurrency = Math.Max(1, federation?.Value.InboxConcurrency ?? 8);
_localActors = localActors;
_remoteActors = remoteActors;
_handlers = handlers.ToDictionary(h => h.Type, StringComparer.Ordinal);
_logger = logger;
_ledger = ledger;
}
public JobKind Kind => JobKind.ProcessInbox;
public int Concurrency => _concurrency;
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");
}
if (payload.ForwardedBy != default && Forwarded.NamesUnknownKey(activity, actor))
actor = await _remoteActors.GetActor(actor.ActorURI, refresh: true, token) ?? actor;
if (payload.ForwardedBy != default && !Forwarded.Proven(activity, actor))
{
var (confirmed, drop) = await Forwarded.Confirm(activity, type, actor.ActorURI, _remoteActors, token);
if (drop != default)
{
Record(job, payload, activity, type, started, new ArrivalVerdict(), Interactions.Dropped, drop);
return JobOutcome.Done;
}
activity = confirmed;
}
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(),
payload.ForwardedBy == default ? payload.Activity : default, payload.ForwardedBy);
Arrival.Current = arrival;
try
{
await LocalReferences.Canonicalise(activity, _localActors?.BaseAddress, token);
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);
}
}
}