- When a sender's key cannot be fetched because its server timed out or answered 5xx, the inbox answers 503 with Retry-After: 300 instead of 401, so Mastodon 4.7 retries rather than switching to RFC 9421 signatures we do not verify yet. The fetcher's failure cache now remembers whether a failure was temporary. - Follow, Like, Block, Accept, Reject and Undo ids are deliberately not dereferenceable: serving them would publish who follows, likes and blocks whom. They are always sent with their object embedded. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_012CzABvBkbcFqoHdmi8b9WB
142 lines
6.0 KiB
C#
142 lines
6.0 KiB
C#
using PrivaPub.Federation.Actors;
|
|
using PrivaPub.Federation.Moderation;
|
|
using PrivaPub.Federation.Objects;
|
|
using PrivaPub.Federation.Signing;
|
|
using PrivaPub.Infrastructure.Jobs;
|
|
using PrivaPub.Models.Federation;
|
|
using PrivaPub.Models.Jobs;
|
|
|
|
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);
|
|
|
|
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);
|
|
}
|
|
|
|
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;
|
|
|
|
public InboxReceiver(ILocalActorService localActors, IRemoteActorService remoteActors, IJobQueue queue, IDomainBlocks domainBlocks,
|
|
ILogger<InboxReceiver> logger)
|
|
{
|
|
_localActors = localActors;
|
|
_remoteActors = remoteActors;
|
|
_queue = queue;
|
|
_domainBlocks = domainBlocks;
|
|
_logger = logger;
|
|
}
|
|
|
|
public async Task<InboxResult> Receive(HttpRequest request, LocalActor recipient, CancellationToken token)
|
|
{
|
|
if (request.ContentLength > MaxBodyBytes)
|
|
return new(StatusCodes.Status413PayloadTooLarge);
|
|
|
|
using var buffer = new MemoryStream();
|
|
await request.Body.CopyToAsync(buffer, token);
|
|
if (buffer.Length > MaxBodyBytes)
|
|
return new(StatusCodes.Status413PayloadTooLarge);
|
|
var body = buffer.ToArray();
|
|
|
|
JsonNode activity;
|
|
try
|
|
{
|
|
activity = JsonNode.Parse(body);
|
|
}
|
|
catch (JsonException)
|
|
{
|
|
return new(StatusCodes.Status400BadRequest, "the body is not JSON");
|
|
}
|
|
if (activity is not JsonObject)
|
|
return new(StatusCodes.Status400BadRequest, "the body is not an activity");
|
|
|
|
var type = Value(activity, "type");
|
|
var actorUri = Id(activity["actor"]);
|
|
if (string.IsNullOrEmpty(type) || string.IsNullOrEmpty(actorUri))
|
|
return new(StatusCodes.Status400BadRequest, "type and actor are required");
|
|
|
|
var parameters = HttpSignatures.Parse(request.Headers["Signature"].ToString());
|
|
if (parameters == default)
|
|
return new(StatusCodes.Status401Unauthorized, "missing or unreadable Signature header");
|
|
if (_domainBlocks.IsSuspended(HostOf(parameters.KeyId)) || _domainBlocks.IsSuspended(HostOf(actorUri)))
|
|
return new(StatusCodes.Status202Accepted);
|
|
|
|
var requestProblem = HttpSignatures.CheckRequest(request, parameters, body);
|
|
if (requestProblem != default)
|
|
return new(StatusCodes.Status401Unauthorized, 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);
|
|
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);
|
|
if (keyOwner == default || !HttpSignatures.Verify(keyOwner.PublicKey, signingString, parameters.Signature))
|
|
return new(StatusCodes.Status401Unauthorized, "the signature does not verify");
|
|
}
|
|
|
|
if (!string.Equals(keyOwner.ActorURI, actorUri, StringComparison.Ordinal))
|
|
return new(StatusCodes.Status401Unauthorized, "the activity's actor is not the key's owner");
|
|
|
|
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);
|
|
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);
|
|
}
|
|
|
|
static string HostOf(string uri) => Uri.TryCreate(uri, UriKind.Absolute, out var parsed) ? parsed.Host : 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");
|
|
|
|
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");
|
|
break;
|
|
case "Undo" when inner is JsonObject && Id(inner["actor"]) != actorUri:
|
|
return new(StatusCodes.Status400BadRequest, "an actor can only undo its own activities");
|
|
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");
|
|
}
|
|
return default;
|
|
}
|
|
}
|
|
}
|