diff --git a/CLAUDE.md b/CLAUDE.md index 4238e24..53a71e7 100644 --- a/CLAUDE.md +++ b/CLAUDE.md @@ -175,6 +175,8 @@ group www-data and reaches the private mongod; `sudo -u www-data` works too. actor's origin now, a Delete of a public or unlisted copy once that origin answers 404 or 410 (`RemoteActorService.IsGone`), anything else let go. Forwarded copies have their own dedupe key, so a copy that failed never hides the author's own delivery; + - an activity is queued once per id, but an id that comes back carrying another type, actor or object is a second + activity, not a copy (Friendica's ids are `uniqid()`, which two of its processes can share), and is queued apart; - the activity's `id`, and any object it creates, updates or deletes, is on the actor's origin; a cross-origin object is refetched from its own origin. 5. **Status codes:** diff --git a/PrivaPub.Tests/Federation/InboxGapTests.cs b/PrivaPub.Tests/Federation/InboxGapTests.cs index 84348b6..d0e512d 100644 --- a/PrivaPub.Tests/Federation/InboxGapTests.cs +++ b/PrivaPub.Tests/Federation/InboxGapTests.cs @@ -79,6 +79,37 @@ namespace PrivaPub.Tests.Federation // Pixelfed names our post by its page (/@name/, the post's url) in its Like, Announce and their Undo: the post is // found the same as by its id + // Friendica's activity ids are uniqid(), a prefix and the microsecond: two of its processes answering two follows at + // once gave both Accepts one id, and the second was dropped as a copy of the first, leaving that follow pending + [Fact] + public async Task Two_activities_a_server_gave_one_id_are_both_taken_and_a_true_copy_once() + { + var token = TestContext.Current.CancellationToken; + var (root, alice) = await _harness.Persona("alice"); + var dario = new RemoteActor(_harness.Peer, "dario"); + var jun = new RemoteActor(_harness.Peer, "jun"); + await _harness.Follows.Follow(root, new FollowForm { AvatarId = alice.Id, Target = dario.Id }, token); + await _harness.Follows.Follow(root, new FollowForm { AvatarId = alice.Id, Target = jun.Id }, token); + Task Of(RemoteActor target) => + DB.Default.Find().Match(f => f.AvatarId == alice.Id && f.TargetActorURI == target.Id).ExecuteSingleAsync(token); + async Task Accept(RemoteActor by, string id) => new() + { + ["id"] = id, + ["type"] = "Accept", + ["actor"] = by.Id, + ["object"] = new JsonObject { ["id"] = (await Of(by)).FollowActivityURI, ["type"] = "Follow", ["actor"] = alice.Uri, ["object"] = by.Id } + }; + var shared = NewId(dario, "activity"); + + var first = await _harness.Deliver(dario, "/human-centipede", await Accept(dario, shared)); + var second = await _harness.Deliver(jun, "/human-centipede", await Accept(jun, shared)); + var copy = await _harness.Deliver(jun, "/human-centipede", await Accept(jun, shared)); + + Assert.Equal(("queued", "queued", "duplicate"), (first.Reason, second.Reason, copy.Reason)); + Assert.Equal(FollowState.Accepted, (await Of(dario)).State); + Assert.Equal(FollowState.Accepted, (await Of(jun)).State); + } + [Fact] public async Task A_like_naming_a_post_by_its_page_counts_and_its_undo_too() { diff --git a/PrivaPub.Tests/Support/Harness.cs b/PrivaPub.Tests/Support/Harness.cs index 9ffd7bf..a642f06 100644 --- a/PrivaPub.Tests/Support/Harness.cs +++ b/PrivaPub.Tests/Support/Harness.cs @@ -135,10 +135,12 @@ namespace PrivaPub.Tests.Support public async Task Deliver(RemoteActor sender, string path, JsonNode activity) { var result = await Receiver.Receive(sender.Post(Host, path, activity), default, CancellationToken.None); + // the job this delivery queued: under its id, or apart (a forwarded copy; an id its sender reused) var id = activity is JsonObject ? activity["id"]?.GetValue() : default; var dedupe = (result.Reason == "forwarded" ? "inbox|forwarded|" : "inbox|") + id; - var job = await DB.Default.Find().Match(j => j.DedupeKey == dedupe).ExecuteFirstAsync(); - if (job != default) + var job = await DB.Default.Find().Match(j => (j.DedupeKey == dedupe || j.DedupeKey.StartsWith(dedupe + "|")) && j.State == JobState.Pending) + .Sort(j => j.CreatedAt, Order.Descending).ExecuteFirstAsync(); + if (job != default && result.Reason != "duplicate") Assert.Equal(JobResult.Done, (await Processor.Handle(job, CancellationToken.None)).Result); return result; } diff --git a/PrivaPub/Federation/Inbox/InboxReceiver.cs b/PrivaPub/Federation/Inbox/InboxReceiver.cs index a26a35a..30c76f5 100644 --- a/PrivaPub/Federation/Inbox/InboxReceiver.cs +++ b/PrivaPub/Federation/Inbox/InboxReceiver.cs @@ -191,14 +191,34 @@ namespace PrivaPub.Federation.Inbox // 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 queued = await _queue.Enqueue(JobKind.ProcessInbox, JsonSerializer.Serialize(payload), - new Uri(actorUri).Host.ToLowerInvariant(), dedupe, token); + 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(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 ShapeProblem(string type, JsonNode activity, string actorUri, CancellationToken token) { var activityId = Id(activity); diff --git a/PrivaPub/Infrastructure/Jobs/JobQueue.cs b/PrivaPub/Infrastructure/Jobs/JobQueue.cs index c3b13ae..0978a73 100644 --- a/PrivaPub/Infrastructure/Jobs/JobQueue.cs +++ b/PrivaPub/Infrastructure/Jobs/JobQueue.cs @@ -27,6 +27,7 @@ namespace PrivaPub.Infrastructure.Jobs public interface IJobQueue { Task Enqueue(JobKind kind, string payload, string host, string dedupeKey, CancellationToken token); + Task Payload(string dedupeKey, CancellationToken token); Task EnqueueMany(IEnumerable jobs, CancellationToken token); Task Lease(JobKind kind, IReadOnlyCollection busyHosts, CancellationToken token); Task Finish(Job job, JobOutcome outcome, int maxAttempts, CancellationToken token); @@ -51,6 +52,10 @@ namespace PrivaPub.Infrastructure.Jobs public async Task Enqueue(JobKind kind, string payload, string host, string dedupeKey, CancellationToken token) => await EnqueueMany(new[] { new Job { Kind = kind, Payload = payload, Host = host, DedupeKey = dedupeKey } }, token) == 1; + // what the job a dedupe key names carries, or null when there is none + public async Task Payload(string dedupeKey, CancellationToken token) => + (await DB.Default.Find().Match(j => j.DedupeKey == dedupeKey).ExecuteFirstAsync(token))?.Payload; + public async Task EnqueueMany(IEnumerable jobs, CancellationToken token) { var inserted = 0;