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 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 _logger; public InboxReceiver(ILocalActorService localActors, IRemoteActorService remoteActors, IJobQueue queue, IDomainBlocks domainBlocks, ILogger logger) { _localActors = localActors; _remoteActors = remoteActors; _queue = queue; _domainBlocks = domainBlocks; _logger = logger; } public async Task 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 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; } } }