Files
SocialPub/PrivaPub/Federation/Inbox/InboxService.cs
T
thepraandClaude Opus 5.5 a060204dd6 A remote actor is believed only from its own origin
S1 and S2 of the roadmap. RemoteActorService:
- FetchObject accepts a document only when its id is the address it was
  served from; a same-origin document naming another address is asked for
  at that address once (how GoToSocial serves its key URIs), anything else
  is dropped;
- GetActorByKeyId accepts a key only when the actor lists it, its owner is
  the actor and it lives on the actor's origin, whether the keyId points at
  the actor or at a key document;
- a refetch for a key or an actor happens at most once per five minutes,
  so a bad signature cannot make us hammer a host;
- the cache row is written by one atomic upsert on ActorURI;
- every fetch is signed by the instance actor, never by the persona that
  happened to receive the activity.

The inbox refuses an activity whose id is not on its actor's origin, and
an Undo of someone else's activity; a Create's object, an Update and a
Delete must be on the actor's origin too, and a cross-origin object is
refetched from its own origin before it is trusted.

Tests: a fake peer on two origins serves forged actors, foreign-owned keys,
cross-origin key documents, aliases and a GoToSocial-style key address
(integration, PRIVAPUB_TEST_MONGOD=1).

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_012CzABvBkbcFqoHdmi8b9WB
2026-10-01 10:49:58 +02:00

409 lines
16 KiB
C#

using MongoDB.Entities;
using PrivaPub.Models.Federation;
using PrivaPub.Models.Group;
using PrivaPub.Models.User;
using PrivaPub.StaticServices;
using System.Text.Json;
using System.Text.Json.Nodes;
using DmPostEntity = PrivaPub.Models.Post.DmPost;
using GroupEntity = PrivaPub.Models.Group.Group;
using PostEntity = PrivaPub.Models.Post.Post;
using PrivaPub.Federation.Actors;
using PrivaPub.Federation.Signing;
using PrivaPub.Federation.Rendering;
using PrivaPub.Federation.Objects;
using PrivaPub.Federation.Outbox;
namespace PrivaPub.Federation.Inbox
{
public sealed record InboxResult(int StatusCode, string Error = default);
public interface IInboxService
{
Task<InboxResult> Receive(HttpRequest request, LocalActor recipient, CancellationToken token);
}
public class InboxService : IInboxService
{
const int MaxBodyBytes = 1024 * 1024;
readonly DbEntities _dbEntities;
readonly ILocalActorService _localActors;
readonly IRemoteActorService _remoteActors;
readonly IDeliveryService _delivery;
readonly ILogger<InboxService> _logger;
public InboxService(DbEntities dbEntities, ILocalActorService localActors, IRemoteActorService remoteActors,
IDeliveryService delivery, ILogger<InboxService> logger)
{
_dbEntities = dbEntities;
_localActors = localActors;
_remoteActors = remoteActors;
_delivery = delivery;
_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");
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 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");
_logger.LogInformation("Inbox {Recipient}: {Type} from {Actor}", recipient?.Handle ?? "shared", type, actorUri);
return type switch
{
"Follow" => await Follow(activity, keyOwner, token),
"Undo" => await Undo(activity, keyOwner, token),
"Create" => await Create(activity, keyOwner, token),
"Delete" => await Delete(activity, keyOwner, token),
"Update" => await Update(activity, keyOwner, token),
_ => new(StatusCodes.Status202Accepted)
};
}
async Task<InboxResult> Follow(JsonNode follow, ForeignAvatar follower, CancellationToken token)
{
var target = await _localActors.FindByUri(Id(follow["object"]), token);
if (target == default || target.Kind == LocalActorKind.Application)
return new(StatusCodes.Status404NotFound, "no such local actor");
var existing = await _dbEntities.Followers
.Match(f => f.LocalActorId == target.Id && f.LocalActorKind == target.Kind && f.ActorURI == follower.ActorURI)
.ExecuteFirstAsync(token);
var record = existing ?? new Follower
{
LocalActorId = target.Id,
LocalActorKind = target.Kind,
ActorURI = follower.ActorURI
};
record.InboxURL = follower.InboxURL;
record.SharedInboxURL = follower.SharedInboxURL;
record.FollowActivityURI = Id(follow);
record.IsAccepted = existing?.IsAccepted == true || !target.ManuallyApprovesFollowers;
await DB.Default.SaveAsync(record, token);
if (!record.IsAccepted)
return new(StatusCodes.Status202Accepted);
if (target.Kind == LocalActorKind.Group)
await AddForeignMember(target.Id, follower.ActorURI, token);
var accept = ActivityPubRenderer.Accept(target, follow, $"accept-{record.ID}-{DateTime.UtcNow.Ticks}");
await _delivery.Enqueue(target, new[] { follower.InboxURL }, accept, token);
return new(StatusCodes.Status202Accepted);
}
async Task<InboxResult> Undo(JsonNode undo, ForeignAvatar actor, CancellationToken token)
{
var inner = undo["object"];
var innerType = inner is JsonObject ? Value(inner, "type") : default;
var innerId = Id(inner);
if (inner is JsonObject && Id(inner["actor"]) != actor.ActorURI)
return new(StatusCodes.Status400BadRequest, "an actor can only undo its own activities");
if (innerType is not (null or "Follow"))
return new(StatusCodes.Status202Accepted);
var followers = await _dbEntities.Followers.Match(f => f.ActorURI == actor.ActorURI).ExecuteAsync(token);
var targetUri = inner is JsonObject ? Id(inner["object"]) : default;
foreach (var follower in followers)
{
var matchesActivity = innerId != default && follower.FollowActivityURI == innerId;
var target = targetUri == default ? default : await _localActors.FindByUri(targetUri, token);
var matchesTarget = target != default && target.Id == follower.LocalActorId && target.Kind == follower.LocalActorKind;
if (!matchesActivity && !matchesTarget)
continue;
await DB.Default.DeleteAsync<Follower>(follower.ID);
if (follower.LocalActorKind == LocalActorKind.Group)
await RemoveForeignMember(follower.LocalActorId, actor.ActorURI, token);
}
return new(StatusCodes.Status202Accepted);
}
async Task<InboxResult> Create(JsonNode create, ForeignAvatar author, CancellationToken token)
{
var note = create["object"];
if (note is not JsonObject || !Origin.Same(Id(note), author.ActorURI))
{
using var fetched = await _remoteActors.FetchObject(Id(note), token);
note = fetched == default ? default : JsonNode.Parse(fetched.Root.GetRawText());
}
if (note is not JsonObject || Value(note, "type") is not ("Note" or "Article" or "Page" or "Question"))
return new(StatusCodes.Status202Accepted);
var objectUri = Id(note);
if (string.IsNullOrEmpty(objectUri) || Id(note["attributedTo"]) != author.ActorURI)
return new(StatusCodes.Status400BadRequest, "the object is not attributed to the actor");
var addressed = Addresses(note).Concat(Addresses(create)).Distinct(StringComparer.Ordinal).ToList();
var isPublic = addressed.Contains(ActivityPubRenderer.Public) || addressed.Contains("as:Public") || addressed.Contains("Public");
var inReplyTo = Id(note["inReplyTo"]);
var localTargets = new List<LocalActor>();
foreach (var uri in addressed.Concat(new[] { Id(note["audience"]) }).Where(u => u != default).Distinct())
{
var local = await _localActors.FindByUri(uri, token);
if (local != default && localTargets.All(l => l.Id != local.Id))
localTargets.Add(local);
}
if (isPublic || addressed.Any(a => a == author.ActorURI + "/followers" || a.EndsWith("/followers")))
{
var group = localTargets.FirstOrDefault(t => t.Kind == LocalActorKind.Group);
if (group != default && !await IsAcceptedFollower(group, author.ActorURI, token))
group = default;
if (group == default && localTargets.All(t => t.Kind != LocalActorKind.Person))
return new(StatusCodes.Status202Accepted);
if (await _dbEntities.Posts.Match(p => p.ObjectURI == objectUri).ExecuteAnyAsync(token))
return new(StatusCodes.Status202Accepted);
var post = new PostEntity
{
ObjectURI = objectUri,
ActorURI = author.ActorURI,
GroupId = group?.Id,
Title = Value(note, "summary") ?? Value(note, "name"),
Text = Value(note, "content"),
HasContentWarning = note["sensitive"] is JsonValue sensitive && sensitive.TryGetValue<bool>(out var s) && s,
AnsweringToPostId = await LocalPostId(inReplyTo, token) ?? inReplyTo,
IsFederatedCopy = true,
CreationDate = DateTime.TryParse(Value(note, "published"), out var published) ? published.ToUniversalTime() : DateTime.UtcNow
};
await DB.Default.SaveAsync(post, token);
if (group != default)
{
var announce = ActivityPubRenderer.Announce(group, objectUri, $"announce-{post.ID}");
await _delivery.EnqueueToFollowers(group, announce, token);
}
return new(StatusCodes.Status202Accepted);
}
var recipients = localTargets.Where(t => t.Kind == LocalActorKind.Person).ToList();
if (recipients.Count == 0)
return new(StatusCodes.Status202Accepted);
if (await _dbEntities.DmPosts.Match(p => p.ObjectURI == objectUri).ExecuteAnyAsync(token))
return new(StatusCodes.Status202Accepted);
var participants = addressed.Where(a => a != ActivityPubRenderer.Public).Append(author.ActorURI).ToList();
var dmGroup = await FindOrCreateDmGroup(participants, Value(note, "context") ?? Value(note, "conversation"), token);
var dm = new DmPostEntity
{
ObjectURI = objectUri,
ActorURI = author.ActorURI,
GroupId = dmGroup.ID,
Title = Value(note, "summary"),
Text = Value(note, "content"),
HasContentWarning = note["sensitive"] is JsonValue dmSensitive && dmSensitive.TryGetValue<bool>(out var ds) && ds,
AnsweringToPostId = inReplyTo,
IsFederatedCopy = true,
CreationDate = DateTime.TryParse(Value(note, "published"), out var dmPublished) ? dmPublished.ToUniversalTime() : DateTime.UtcNow
};
await DB.Default.SaveAsync(dm, token);
await DB.Default.Update<DmGroup>().MatchID(dmGroup.ID).Modify(g => g.UpdatedAt, DateTime.UtcNow).ExecuteAsync(token);
return new(StatusCodes.Status202Accepted);
}
async Task<InboxResult> Delete(JsonNode delete, ForeignAvatar actor, CancellationToken token)
{
var objectUri = Id(delete["object"]);
if (!Origin.Same(objectUri, actor.ActorURI))
return new(StatusCodes.Status202Accepted);
if (objectUri == actor.ActorURI)
{
actor.DeletionAt = DateTime.UtcNow;
actor.AccountState = AvatarAccountState.Deleted;
await DB.Default.SaveAsync(actor, token);
var followers = await _dbEntities.Followers.Match(f => f.ActorURI == actor.ActorURI).ExecuteAsync(token);
foreach (var follower in followers)
{
await DB.Default.DeleteAsync<Follower>(follower.ID);
if (follower.LocalActorKind == LocalActorKind.Group)
await RemoveForeignMember(follower.LocalActorId, actor.ActorURI, token);
}
return new(StatusCodes.Status202Accepted);
}
await DB.Default.DeleteAsync<PostEntity>(p => p.ObjectURI == objectUri && p.ActorURI == actor.ActorURI);
await DB.Default.DeleteAsync<DmPostEntity>(p => p.ObjectURI == objectUri && p.ActorURI == actor.ActorURI);
return new(StatusCodes.Status202Accepted);
}
async Task<InboxResult> Update(JsonNode update, ForeignAvatar actor, CancellationToken token)
{
var inner = update["object"];
if (Id(inner) == actor.ActorURI)
{
await _remoteActors.GetActor(actor.ActorURI, refresh: true, token);
return new(StatusCodes.Status202Accepted);
}
if (inner is not JsonObject || Id(inner["attributedTo"]) != actor.ActorURI || !Origin.Same(Id(inner), actor.ActorURI))
return new(StatusCodes.Status202Accepted);
var objectUri = Id(inner);
var text = Value(inner, "content");
await DB.Default.Update<PostEntity>()
.Match(p => p.ObjectURI == objectUri && p.ActorURI == actor.ActorURI)
.Modify(p => p.Text, text)
.Modify(p => p.UpdateDate, DateTime.UtcNow)
.ExecuteAsync(token);
await DB.Default.Update<DmPostEntity>()
.Match(p => p.ObjectURI == objectUri && p.ActorURI == actor.ActorURI)
.Modify(p => p.Text, text)
.Modify(p => p.UpdateDate, DateTime.UtcNow)
.ExecuteAsync(token);
return new(StatusCodes.Status202Accepted);
}
async Task<bool> IsAcceptedFollower(LocalActor group, string actorUri, CancellationToken token) =>
await _dbEntities.Followers
.Match(f => f.LocalActorId == group.Id && f.LocalActorKind == LocalActorKind.Group && f.ActorURI == actorUri && f.IsAccepted)
.ExecuteAnyAsync(token);
async Task AddForeignMember(string groupId, string actorUri, CancellationToken token)
{
var group = await _dbEntities.Groups.MatchID(groupId).ExecuteFirstAsync(token);
if (group == default || group.Members.Any(m => m.IsForeign && m.AvatarId == actorUri))
return;
group.Members.Add(new GroupMember { AvatarId = actorUri, IsForeign = true });
group.UpdatedAt = DateTime.UtcNow;
await DB.Default.SaveAsync(group, token);
}
async Task RemoveForeignMember(string groupId, string actorUri, CancellationToken token)
{
var group = await _dbEntities.Groups.MatchID(groupId).ExecuteFirstAsync(token);
if (group == default || group.Members.RemoveAll(m => m.IsForeign && m.AvatarId == actorUri) == 0)
return;
group.UpdatedAt = DateTime.UtcNow;
await DB.Default.SaveAsync<GroupEntity>(group, token);
}
async Task<DmGroup> FindOrCreateDmGroup(List<string> participantUris, string context, CancellationToken token)
{
var members = new List<GroupMember>();
foreach (var uri in participantUris.Distinct(StringComparer.Ordinal))
{
var local = await _localActors.FindByUri(uri, token);
members.Add(local != default && local.Kind == LocalActorKind.Person
? new GroupMember { AvatarId = local.Id }
: new GroupMember { AvatarId = uri, IsForeign = true });
}
if (!string.IsNullOrEmpty(context))
{
var byContext = await _dbEntities.DmGroups.Match(g => g.ConversationURI == context).ExecuteFirstAsync(token);
if (byContext != default)
return byContext;
}
var keys = members.Select(m => m.AvatarId).OrderBy(k => k, StringComparer.Ordinal).ToList();
var candidates = await _dbEntities.DmGroups
.Match(g => !g.DeletionAt.HasValue && g.Members.Count == keys.Count)
.ExecuteAsync(token);
var match = candidates.FirstOrDefault(g => g.Members.Select(m => m.AvatarId).OrderBy(k => k, StringComparer.Ordinal).SequenceEqual(keys));
if (match != default)
return match;
var dmGroup = new DmGroup { Members = members, ConversationURI = context };
await DB.Default.SaveAsync(dmGroup, token);
return dmGroup;
}
async Task<string> LocalPostId(string uri, CancellationToken token)
{
if (string.IsNullOrEmpty(uri) || !uri.StartsWith(_localActors.BaseAddress + "/peasants/", StringComparison.OrdinalIgnoreCase))
return default;
var id = uri[(uri.LastIndexOf('/') + 1)..];
return await _dbEntities.Posts.MatchID(id).ExecuteAnyAsync(token) ? id : default;
}
static IEnumerable<string> Addresses(JsonNode node)
{
foreach (var field in new[] { "to", "cc", "bto", "bcc", "audience" })
{
switch (node[field])
{
case JsonArray array:
foreach (var item in array)
{
var id = Id(item);
if (id != default)
yield return id;
}
break;
case JsonNode single when Id(single) is { } id:
yield return id;
break;
}
}
}
public static string Id(JsonNode node) => node switch
{
JsonValue value when value.TryGetValue<string>(out var text) => text,
JsonObject obj => Value(obj, "id") ?? Value(obj, "href"),
JsonArray array => array.Select(Id).FirstOrDefault(i => i != default),
_ => default
};
public static string Value(JsonNode node, string property) =>
node is JsonObject obj && obj[property] is JsonValue value && value.TryGetValue<string>(out var text) ? text : default;
}
}