Lemmy takes a report only from a person or a service, about one post or comment, addressed to its community, and it answered PrivaPub's Flag (the instance actor's, an Application, with no `to` and the account and posts as its object) 400. A report of a post or comment in a community on a server whose NodeInfo names Lemmy now leaves from `privapub_reports`, a Service with its own key that names nobody: one Flag per post, `to` the community (its own audience, else its thread's), with the persona's words, or the category, in `summary` and `content`, sent to the community's inbox. This is the second exception to "a server's software is for display" (owner decision 2026-10-06, `ReportService.ServiceReportTakers`). Every other server keeps the instance actor's report. An account alone is not reported to Lemmy, which takes no such report, and `forwarded` now says whether anything left. The reporter is read unsigned in SecureMode and answers WebFinger like the instance actor. Nobody follows or mentions it, the Mastodon API has no account for it, and a migration reserves its name. Checked live: Lemmy 1.0 and 0.19 keep the reports of a thread and of a comment, with alice's words, from "Reports from privapub.test", and none names her (69 checks). A sweep of every scenario with this and the next commit: 876 checks pass; Ghost's Network feed listed alice's post too late once, and Ghost passes alone. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01LsXgEaXee4GCU1hwYgPJXw
273 lines
13 KiB
C#
273 lines
13 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);
|
|
// an Undo of a Follow that comes again while the follow it ends exists again is meant for that follow: Mastodon
|
|
// gives every Undo it sends after reading a roll-call (FEP-8fcf) one id, `{actor}#follows//undo`. It is kept
|
|
// once per follow, so a retry still counts as a copy
|
|
if (!queued && dedupe != default && type == "Undo" && await FollowAgain(activity, actorUri, token) is { } follow)
|
|
queued = await _queue.Enqueue(JobKind.ProcessInbox, queuedPayload, host, $"{dedupe}|{Digest(activity)}|{follow}", token);
|
|
// FEP-8fcf: what the sender says of its followers here (FollowingSynchronization)
|
|
if (queued && forwardedBy == default && FollowersSynchronization.Claim(actorUri, request.Headers[FollowersSynchronization.Header].ToString()) is { } claim)
|
|
await _queue.Enqueue(JobKind.SynchronizeFollowing, JsonSerializer.Serialize(claim), host, FollowingSynchronization.DedupeKey(claim, activityId), 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");
|
|
}
|
|
|
|
// the follow an Undo of a Follow would end, when there is one: its record's id
|
|
async Task<string> FollowAgain(JsonNode activity, string actorUri, CancellationToken token)
|
|
{
|
|
if (activity["object"] is not JsonObject inner || Value(inner, "type") != "Follow" || Id(inner["object"]) is not { } targetUri
|
|
|| await _localActors.FindByUri(targetUri, token) is not { } target)
|
|
return default;
|
|
var follower = await MongoDB.Entities.DB.Default.Find<Follower>()
|
|
.Match(f => f.ActorURI == actorUri && f.LocalActorId == target.Id && f.LocalActorKind == target.Kind)
|
|
.ExecuteFirstAsync(token);
|
|
return follower?.ID;
|
|
}
|
|
|
|
static string HostOf(string uri) => Uri.TryCreate(uri, UriKind.Absolute, out var parsed) ? parsed.Host.ToLowerInvariant() : default;
|
|
|
|
public 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
|
|
public 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":
|
|
// (by its actor's id, or the profile page Forte names instead)
|
|
var target = await _localActors.FindByAddress(Id(inner), token);
|
|
if (target is not { IsFederated: true } || target.IsServerActor)
|
|
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;
|
|
}
|
|
}
|
|
}
|