Files
SocialPub/PrivaPub/Federation/Inbox/InboxReceiver.cs
T
thepraandClaude Opus 5.5 d37999970c M5: outbound requests, previews, the media proxy and served traffic
- HttpScope (AsyncLocal) tags each outbound request with a purpose and a trigger:
  - purpose is set by the caller: actor, key, object, webfinger, context, nodeinfo;
  - trigger is set by the job kind, by "verify" during inbox verification, or defaults
    to "request".
- FederationHttp records every JSON, media and stream fetch: status, time, bytes, hops,
  and an outcome of ok, refused or failed, with a reason: disallowed, remembered,
  bad-redirect, too-many-redirects, content-type, too-large, bad-json, private-address,
  timeout, network, or the status. A fetch a reader caused (trigger "request") is only
  counted per server per day.
- Link previews record a 'preview' event: card, no-card or failed.
- The media proxy counts cache hits.
- TrafficMeter counts the client API per endpoint group, method and status class. It
  counts our served documents (actor, outbox, collection, object, activity, licence,
  webfinger, nodeinfo) by kind, status and whether signed, per day and never per server,
  and never names a circle's collections.

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

217 lines
9.3 KiB
C#

using PrivaPub.Federation.Actors;
using PrivaPub.Federation.Moderation;
using PrivaPub.Federation.Objects;
using PrivaPub.Federation.Signing;
using PrivaPub.Infrastructure.Http;
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)
{
using var scope = HttpScope.Triggered("verify");
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;
}
}
}