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 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 _logger; readonly IInteractionLedger _ledger; public InboxReceiver(ILocalActorService localActors, IRemoteActorService remoteActors, IJobQueue queue, IDomainBlocks domainBlocks, ILogger 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 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 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 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; } } }