Files
SocialPub/PrivaPub/Federation/Inbox/InboxReceiver.cs
T
thepraandClaude Opus 5.5 4fa53f63bd M2: every inbox answer is recorded
InboxReceiver records each answer once, in a finally, with a reason: too-large, not-json,
not-activity, missing-type-or-actor, no-signature, the signature check's own codes
(headers-unsigned, digest-mismatch, header-unreadable, date-skew, expired,
algorithm-unsupported), signature-invalid, actor-not-key-owner, key-unavailable,
id-cross-origin, undo-foreign, misattributed, unknown-recipient, and for 202s queued,
duplicate, suspended or self-delete-unknown-key. Each event carries the activity and object
type, the inbox, the signature scheme, the bytes and the time taken. A 404 for an unknown
/mouth and a rate-limited inbox (429, from OnRejected) are recorded too.

Until the signature verifies, the host is only claimed, so it is kept only if the server is
already known. A suspended server is recorded under its own name.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01ELjqpznMFMNrJoJUj6K5p2
2026-10-03 10:59:55 +02:00

215 lines
9.2 KiB
C#

using PrivaPub.Federation.Actors;
using PrivaPub.Federation.Moderation;
using PrivaPub.Federation.Objects;
using PrivaPub.Federation.Signing;
using PrivaPub.Infrastructure.Jobs;
using PrivaPub.Infrastructure.Statistics;
using PrivaPub.Models.Federation;
using PrivaPub.Models.Jobs;
using PrivaPub.Models.Statistics;
using System.Text.Json;
using System.Text.Json.Nodes;
using static PrivaPub.Federation.Objects.ActivityJson;
namespace PrivaPub.Federation.Inbox
{
public sealed record InboxResult(int StatusCode, string Error = default, int? RetryAfterSeconds = default, string Reason = default);
public sealed record InboxPayload(string ActorURI, string Activity, string Inbox = default, string KeyId = default, string Algorithm = default,
string[] SignedHeaders = default, DateTime? ReceivedAt = default);
public interface IInboxReceiver
{
Task<InboxResult> Receive(HttpRequest request, LocalActor recipient, CancellationToken token);
InboxResult NoSuchRecipient(HttpRequest request);
}
public class InboxReceiver : IInboxReceiver
{
const int KeyRetrySeconds = 300;
const int MaxBodyBytes = 1024 * 1024;
readonly ILocalActorService _localActors;
readonly IRemoteActorService _remoteActors;
readonly IJobQueue _queue;
readonly IDomainBlocks _domainBlocks;
readonly ILogger<InboxReceiver> _logger;
readonly IInteractionLedger _ledger;
public InboxReceiver(ILocalActorService localActors, IRemoteActorService remoteActors, IJobQueue queue, IDomainBlocks domainBlocks,
ILogger<InboxReceiver> logger, IInteractionLedger ledger = default)
{
_localActors = localActors;
_remoteActors = remoteActors;
_queue = queue;
_domainBlocks = domainBlocks;
_logger = logger;
_ledger = ledger;
}
sealed class Receipt
{
public string Type;
public string ObjectType;
public string ClaimedHost;
public string VerifiedActor;
public string Inbox;
public string Signature = "none";
public long? Bytes;
public bool Blocked;
}
public InboxResult NoSuchRecipient(HttpRequest request)
{
var result = new InboxResult(StatusCodes.Status404NotFound, "no such local actor", Reason: "unknown-recipient");
var keyId = HttpSignatures.Parse(request.Headers["Signature"].ToString())?.KeyId;
Record(new Receipt { ClaimedHost = HostOf(keyId), Inbox = "personal", Bytes = request.ContentLength }, result, 0);
return result;
}
public async Task<InboxResult> Receive(HttpRequest request, LocalActor recipient, CancellationToken token)
{
var started = System.Diagnostics.Stopwatch.GetTimestamp();
var receipt = new Receipt { Inbox = recipient == default ? "shared" : "personal", Bytes = request.ContentLength };
InboxResult result = default;
try
{
result = await Receive(request, recipient, receipt, token);
return result;
}
finally
{
Record(receipt, result, (int)System.Diagnostics.Stopwatch.GetElapsedTime(started).TotalMilliseconds);
}
}
void Record(Receipt receipt, InboxResult result, int latencyMs)
{
if (_ledger == default || result == default)
return;
var verified = receipt.VerifiedActor != default;
_ledger.Record(new InteractionEvent
{
Channel = Interactions.Receive,
Host = verified ? Interactions.HostOf(receipt.VerifiedActor) : receipt.ClaimedHost,
Activity = receipt.Type,
Object = receipt.ObjectType,
Status = result.StatusCode,
Outcome = result.StatusCode == StatusCodes.Status202Accepted ? Interactions.Queued : Interactions.Refused,
Reason = result.Reason,
LatencyMs = latencyMs,
Bytes = receipt.Bytes,
Inbox = receipt.Inbox,
Signature = receipt.Signature
}, verified ? receipt.VerifiedActor : default, hostClaimed: !verified && !receipt.Blocked);
}
async Task<InboxResult> Receive(HttpRequest request, LocalActor recipient, Receipt receipt, CancellationToken token)
{
if (request.ContentLength > MaxBodyBytes)
return new(StatusCodes.Status413PayloadTooLarge, Reason: "too-large");
using var buffer = new MemoryStream();
await request.Body.CopyToAsync(buffer, token);
if (buffer.Length > MaxBodyBytes)
return new(StatusCodes.Status413PayloadTooLarge, Reason: "too-large");
var body = buffer.ToArray();
receipt.Bytes = body.Length;
JsonNode activity;
try
{
activity = JsonNode.Parse(body);
}
catch (JsonException)
{
return new(StatusCodes.Status400BadRequest, "the body is not JSON", Reason: "not-json");
}
if (activity is not JsonObject)
return new(StatusCodes.Status400BadRequest, "the body is not an activity", Reason: "not-activity");
var type = Value(activity, "type");
var actorUri = Id(activity["actor"]);
receipt.Type = type;
receipt.ObjectType = activity["object"] is JsonObject embedded ? Value(embedded, "type") : default;
receipt.ClaimedHost = HostOf(actorUri);
if (string.IsNullOrEmpty(type) || string.IsNullOrEmpty(actorUri))
return new(StatusCodes.Status400BadRequest, "type and actor are required", Reason: "missing-type-or-actor");
var parameters = HttpSignatures.Parse(request.Headers["Signature"].ToString());
if (parameters == default)
return new(StatusCodes.Status401Unauthorized, "missing or unreadable Signature header", Reason: "no-signature");
receipt.Signature = "cavage:" + parameters.Algorithm;
receipt.ClaimedHost ??= HostOf(parameters.KeyId);
if (_domainBlocks.IsSuspended(HostOf(parameters.KeyId)) || _domainBlocks.IsSuspended(HostOf(actorUri)))
{
receipt.Blocked = true;
receipt.ClaimedHost = _domainBlocks.IsSuspended(HostOf(actorUri)) ? HostOf(actorUri) : HostOf(parameters.KeyId);
return new(StatusCodes.Status202Accepted, Reason: "suspended");
}
var requestProblem = HttpSignatures.CheckRequest(request, parameters, body);
if (requestProblem != default)
return new(StatusCodes.Status401Unauthorized, requestProblem, Reason: HttpSignatures.ProblemCode(requestProblem));
var signingString = HttpSignatures.SigningString(request, parameters);
var keyOwner = await _remoteActors.GetActorByKeyId(parameters.KeyId, refresh: false, token);
if (keyOwner == default && type == "Delete" && Id(activity["object"]) == actorUri)
return new(StatusCodes.Status202Accepted, Reason: "self-delete-unknown-key");
if (keyOwner == default || !HttpSignatures.Verify(keyOwner.PublicKey, signingString, parameters.Signature))
{
keyOwner = await _remoteActors.GetActorByKeyId(parameters.KeyId, refresh: true, token);
if (keyOwner == default && _remoteActors.KeyTemporarilyUnavailable(parameters.KeyId))
return new(StatusCodes.Status503ServiceUnavailable, "the signing key could not be fetched; try again later", KeyRetrySeconds, "key-unavailable");
if (keyOwner == default || !HttpSignatures.Verify(keyOwner.PublicKey, signingString, parameters.Signature))
return new(StatusCodes.Status401Unauthorized, "the signature does not verify", Reason: "signature-invalid");
}
if (!string.Equals(keyOwner.ActorURI, actorUri, StringComparison.Ordinal))
return new(StatusCodes.Status401Unauthorized, "the activity's actor is not the key's owner", Reason: "actor-not-key-owner");
receipt.VerifiedActor = actorUri;
var shapeProblem = await ShapeProblem(type, activity, actorUri, token);
if (shapeProblem != default)
return shapeProblem;
var activityId = Id(activity);
var payload = new InboxPayload(actorUri, activity.ToJsonString(), recipient == default ? "shared" : "personal", parameters.KeyId,
parameters.Algorithm, parameters.Headers, DateTime.UtcNow);
var queued = await _queue.Enqueue(JobKind.ProcessInbox, JsonSerializer.Serialize(payload),
new Uri(actorUri).Host.ToLowerInvariant(), activityId == default ? default : "inbox|" + activityId, token);
_logger.LogInformation("Inbox {Recipient}: {Type} from {Actor} queued", recipient?.Handle ?? "shared", type, actorUri);
return new(StatusCodes.Status202Accepted, Reason: queued ? "queued" : "duplicate");
}
static string HostOf(string uri) => Uri.TryCreate(uri, UriKind.Absolute, out var parsed) ? parsed.Host.ToLowerInvariant() : default;
async Task<InboxResult> ShapeProblem(string type, JsonNode activity, string actorUri, CancellationToken token)
{
var activityId = Id(activity);
if (activityId != default && !Origin.Same(activityId, actorUri))
return new(StatusCodes.Status400BadRequest, "the activity's id is not on its actor's origin", Reason: "id-cross-origin");
var inner = activity["object"];
switch (type)
{
case "Follow":
var target = await _localActors.FindByUri(Id(inner), token);
if (target is not { IsFederated: true } || target.Kind == LocalActorKind.Application)
return new(StatusCodes.Status404NotFound, "no such local actor", Reason: "unknown-recipient");
break;
case "Undo" when inner is JsonObject && Id(inner["actor"]) != actorUri:
return new(StatusCodes.Status400BadRequest, "an actor can only undo its own activities", Reason: "undo-foreign");
case "Create" or "Update" when inner is JsonObject && Origin.Same(Id(inner), actorUri)
&& inner["attributedTo"] != default && Id(inner["attributedTo"]) != actorUri
&& Value(inner, "type") is not ("Person" or "Service" or "Application" or "Group" or "Organization"):
return new(StatusCodes.Status400BadRequest, "the object is not attributed to the actor", Reason: "misattributed");
}
return default;
}
}
}