Build / Build (push) Successful in 36s
- ObjectRecord, one per stored remote object: the raw JSON (up to 256 KB, always hashed), delivered or fetched, refetched from origin or not, the activity that brought it (or caused the fetch), shared or personal inbox, the signature's key, algorithm and signed headers, received time, published and updated, the delivering activity's @context, and up to ten later revisions from Update. - Delivery details travel from InboxReceiver through the inbox job to the handlers as Arrival.Current. - A host is described from its NodeInfo when we first hear from it, at most weekly (DescribeInstance job), never when someone opens the details view. - GET /api/privapub/v1/statuses/:id/provenance and /api/privapub/v1/instances/:host, with the extensions an object used detected from its raw form (044f quotes, interaction policies, contexts, proofs, Misskey fields, MFM, FEP-8967 links, emoji, polls, language maps, url variants, Markdown content). Checked live: a GoToSocial reply shows as delivered to the shared inbox, signed hs2019 with GoToSocial's fragment-less key id, with its interaction policy detected, and gts.test is described as gotosocial 0.22.1. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_012CzABvBkbcFqoHdmi8b9WB
138 lines
5.7 KiB
C#
138 lines
5.7 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);
|
|
|
|
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 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 || !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;
|
|
}
|
|
}
|
|
}
|