Three shapes Funkwhale 2.0 sends, each of which lost something: - its Accept is named after the Follow it answers, on our origin (`…#follows/<uuid>/accept`), and was refused as off its actor's origin, so no follow of a channel completed: an Accept or Reject whose id extends the id of the activity it answers, on our origin, is now taken; - a channel deletes its uploads in one Delete without an id, their ids in a list as the object's id: each is deleted; - its NodeInfo discovery names the document under its swagger schema's URL: a link whose path names NodeInfo is taken when no rel is known. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01LsXgEaXee4GCU1hwYgPJXw
252 lines
12 KiB
C#
252 lines
12 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);
|
|
|
|
// ForwardedBy: the server that passed the activity on, signing with its own key (Forwarded)
|
|
public sealed record InboxPayload(string ActorURI, string Activity, string Inbox = default, string KeyId = default, string Algorithm = default,
|
|
string[] SignedHeaders = default, DateTime? ReceivedAt = default, string ForwardedBy = 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 = RequestSignature.Of(request)?.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.Refused
|
|
: result.Reason == "forwarded-ignored" ? Interactions.Dropped : Interactions.Queued,
|
|
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");
|
|
|
|
// draft-cavage's Signature, or RFC 9421's Signature-Input and Signature
|
|
var signature = RequestSignature.Of(request);
|
|
if (signature == default)
|
|
return new(StatusCodes.Status401Unauthorized, "missing or unreadable Signature header", Reason: "no-signature");
|
|
receipt.Signature = signature.Scheme;
|
|
receipt.ClaimedHost ??= HostOf(signature.KeyId);
|
|
if (_domainBlocks.IsSuspended(HostOf(signature.KeyId)) || _domainBlocks.IsSuspended(HostOf(actorUri)))
|
|
{
|
|
receipt.Blocked = true;
|
|
receipt.ClaimedHost = _domainBlocks.IsSuspended(HostOf(actorUri)) ? HostOf(actorUri) : HostOf(signature.KeyId);
|
|
return new(StatusCodes.Status202Accepted, Reason: "suspended");
|
|
}
|
|
|
|
var requestProblem = signature.Problem(request, body);
|
|
if (requestProblem != default)
|
|
return new(StatusCodes.Status401Unauthorized, requestProblem, Reason: HttpSignatures.ProblemCode(requestProblem));
|
|
|
|
var signed = signature.Signed(request, _localActors.BaseAddress);
|
|
var keyOwner = await _remoteActors.GetActorByKeyId(signature.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 || !signature.VerifiedBy(keyOwner.PublicKey, signed))
|
|
{
|
|
keyOwner = await _remoteActors.GetActorByKeyId(signature.KeyId, refresh: true, token);
|
|
if (keyOwner == default && _remoteActors.KeyTemporarilyUnavailable(signature.KeyId))
|
|
return new(StatusCodes.Status503ServiceUnavailable, "the signing key could not be fetched; try again later", KeyRetrySeconds, "key-unavailable");
|
|
if (keyOwner == default || !signature.VerifiedBy(keyOwner.PublicKey, signed))
|
|
return new(StatusCodes.Status401Unauthorized, "the signature does not verify", Reason: "signature-invalid");
|
|
}
|
|
|
|
// signed by someone else: passed on by a server that holds the thread, its signature good (a 401 would tell it
|
|
// otherwise, and make a double-knocking sender try its other scheme)
|
|
receipt.VerifiedActor = keyOwner.ActorURI;
|
|
var forwardedBy = string.Equals(keyOwner.ActorURI, actorUri, StringComparison.Ordinal) ? default : keyOwner.ActorURI;
|
|
// (ours come back too, Friendica forwarding our comments in its threads: we know them already)
|
|
if (forwardedBy != default && (!Forwarded.Takeable(type, activity, actorUri) || Origin.Same(actorUri, _localActors.BaseAddress)))
|
|
return new(StatusCodes.Status202Accepted, Reason: "forwarded-ignored");
|
|
|
|
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", signature.KeyId,
|
|
signature.Algorithm, signature.Covered, DateTime.UtcNow, forwardedBy);
|
|
// a forwarded copy never stands in for the author's own delivery, which may carry what the copy could not
|
|
// (a followers-only reply our instance actor cannot read): each is kept once, apart
|
|
var dedupe = activityId == default ? default : (forwardedBy == default ? "inbox|" : "inbox|forwarded|") + activityId;
|
|
var queuedPayload = JsonSerializer.Serialize(payload);
|
|
var host = new Uri(actorUri).Host.ToLowerInvariant();
|
|
var queued = await _queue.Enqueue(JobKind.ProcessInbox, queuedPayload, host, dedupe, token);
|
|
// one id its sender gave two activities is no duplicate (Friendica's ids are uniqid(), a prefix and the
|
|
// microsecond, which two of its processes answering at once share): what each carries tells them apart, and a
|
|
// true copy carries the same
|
|
if (!queued && dedupe != default && CarriesOther(await _queue.Payload(dedupe, token), activity))
|
|
queued = await _queue.Enqueue(JobKind.ProcessInbox, queuedPayload, host, $"{dedupe}|{Digest(activity)}", token);
|
|
_logger.LogInformation("Inbox {Recipient}: {Type} from {Actor} queued", recipient?.Handle ?? "shared", type, actorUri);
|
|
return new(StatusCodes.Status202Accepted, Reason: !queued ? "duplicate" : forwardedBy != default ? "forwarded" : "queued");
|
|
}
|
|
|
|
static string HostOf(string uri) => Uri.TryCreate(uri, UriKind.Absolute, out var parsed) ? parsed.Host.ToLowerInvariant() : default;
|
|
|
|
static bool CarriesOther(string earlierPayload, JsonNode activity)
|
|
{
|
|
if (earlierPayload == default)
|
|
return false;
|
|
var earlier = JsonNode.Parse(JsonSerializer.Deserialize<InboxPayload>(earlierPayload).Activity);
|
|
return Value(earlier, "type") != Value(activity, "type") || Id(earlier["actor"]) != Id(activity["actor"])
|
|
|| earlier["object"]?.ToJsonString() != activity["object"]?.ToJsonString();
|
|
}
|
|
|
|
// what an activity carries, short: its type, actor and object
|
|
static string Digest(JsonNode activity) =>
|
|
Convert.ToHexString(System.Security.Cryptography.SHA256.HashData(System.Text.Encoding.UTF8.GetBytes(
|
|
$"{Value(activity, "type")}\n{Id(activity["actor"])}\n{activity["object"]?.ToJsonString()}")))[..16].ToLowerInvariant();
|
|
|
|
async Task<InboxResult> ShapeProblem(string type, JsonNode activity, string actorUri, CancellationToken token)
|
|
{
|
|
var activityId = Id(activity);
|
|
// (Funkwhale names its answer after the activity it answers: `<our Follow's id>/accept`)
|
|
if (activityId != default && !Origin.Same(activityId, actorUri)
|
|
&& !(type is "Accept" or "Reject" && Id(activity["object"]) is { } answered && Origin.Same(answered, _localActors.BaseAddress)
|
|
&& activityId.StartsWith(answered + "/", StringComparison.Ordinal)))
|
|
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");
|
|
// (attributed to another account of the actor's own server, as Mobilizon's organiser creates a group's
|
|
// event: the handlers read it from that server, Origin.Same)
|
|
case "Create" or "Update" when inner is JsonObject && Origin.Same(Id(inner), actorUri)
|
|
&& inner["attributedTo"] != default && !Origin.Same(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;
|
|
}
|
|
}
|
|
}
|