Files
SocialPub/PrivaPub/Federation/Inbox/InboxProcessor.cs
T
thepraandClaude Opus 5.5 0646de22bd Replies to a persona's posts reach its followers; personas join remote events
Two owner decisions of 2026-10-05, recorded in the roadmap.

Replies passed on ("the fediverse is broken without"): a public or unlisted
reply from another server to a persona's public, unlisted or followers-only
post goes on to the persona's followers as its author's server sent it, as
Mastodon forwards it, never to the replier's own server, never for a
local-only or group post; its edit and deletion follow. Only an activity its
own actor delivered is passed on (Arrival.Raw), so nothing forwarded is
forwarded again. The town checks it as relay.reply cells (specs/relay-five:
882 checks pass); Mastodon takes a passed-on activity only with an LD
signature, which GoToSocial and Akkoma don't add, and the checker knows it.

Events: a persona joins another server's event with a Join and leaves it with
a Leave, both to the organiser only, through
POST /api/privapub/v1/statuses/:id/join|leave; the organiser's Accept or
Reject is routed by our join id and shows as privapub.event.participation.
Events by invitation or taken on another site are refused before anything
is sent. Mobilizon's scenario joins and leaves an event (28 checks) and keeps
one for decePubClient's e2e.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01LsXgEaXee4GCU1hwYgPJXw
2026-10-05 14:55:20 +02:00

124 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;
readonly ILocalActorService _localActors;
public InboxProcessor(IRemoteActorService remoteActors, IEnumerable<IActivityHandler> handlers, ILogger<InboxProcessor> logger,
IInteractionLedger ledger = default, ILocalActorService localActors = default)
{
_localActors = localActors;
_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");
}
if (payload.ForwardedBy != default)
{
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);
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);
}
}
}