From 62189bf11d10067f68bd7602be697e5c9c0b0418 Mon Sep 17 00:00:00 2001 From: thepra Date: Wed, 7 Oct 2026 21:33:49 +0200 Subject: [PATCH] P7: thread backfill reads conversation containers and asks with If-None-Match The FetchReplies job reads a thread's FEP-171b contextHistory before its FEP-7888 context: the posts each Add (or Forte's plain Create) brought in, read from their own servers. A thread collection read whole keeps its ETag when the document itself changes with every post (it counts them, or holds them all with no further page); the next read sends it as If-None-Match and a 304 ends the job. The ETag is kept as sent, since NodeBB's has no quotes and the typed header drops it. A context naming a post we hold (Forte's first post) is not fetched. Checked in the pasture: Mastodon 4.7.3 answers the second read 304; NodeBB 4.16's unquoted ETag is kept (23/23 in its scenario); a Forte thread completes through its replies. Co-Authored-By: Claude Opus 5.5 Claude-Session: https://claude.ai/code/session_012CzABvBkbcFqoHdmi8b9WB --- FEDERATION.md | 8 +- PrivaPub.Tests/Federation/JobHandlerTests.cs | 97 +++++++++++++++++++ PrivaPub.Tests/Support/Peer.cs | 17 ++++ .../Federation/Actors/RemoteActorService.cs | 14 ++- PrivaPub/Federation/Inbox/RemoteReplies.cs | 94 ++++++++++++++---- PrivaPub/Infrastructure/Data/Indexes.cs | 3 + .../Infrastructure/Http/FederationHttp.cs | 35 ++++++- .../Models/Federation/ThreadCollection.cs | 13 +++ docs/INTEROP.md | 15 ++- docs/ROADMAP.md | 7 +- 10 files changed, 268 insertions(+), 35 deletions(-) create mode 100644 PrivaPub/Models/Federation/ThreadCollection.cs diff --git a/FEDERATION.md b/FEDERATION.md index 83949a9..4132e64 100644 --- a/FEDERATION.md +++ b/FEDERATION.md @@ -394,9 +394,11 @@ Posts with a location (shown to nearby users of this server) never leave the ser - **Fetching.** All fetches are signed by the instance actor. They go only to public addresses, follow at most three redirects and read at most 1 MB. - **Threads.** A reply's missing parents are fetched, up to 10 levels. When someone here opens a public remote thread, - its replies are read from their servers, at most once an hour: the thread's `context` collection (FEP-7888) if it has - one, else its `replies` (PeerTube's `comments`) and theirs, two levels down; 5 pages and 100 posts at most. Only - public and unlisted replies are kept, each fetched from its own origin. + its replies are read from their servers, at most once an hour: the thread's conversation container (FEP-171b's + `contextHistory`: the posts its owner's `Add`s brought in) or its `context` collection (FEP-7888) if it has one, else + its `replies` (PeerTube's `comments`) and theirs, two levels down; 5 pages and 100 posts at most. Only public and + unlisted replies are kept, each fetched from its own origin. A thread collection read whole, which counts its posts or + holds them all, is asked again with `If-None-Match` and its last ETag; a 304 ends the read. - **Reading our documents (SecureMode).** privapub.thepra.dev answers ActivityPub GETs only when they are signed, like Mastodon's authorized fetch; the server's own actors (the instance actor `/peasants/privapub` and the reporter `/peasants/privapub_reports`) are the exception, since their keys are needed first. diff --git a/PrivaPub.Tests/Federation/JobHandlerTests.cs b/PrivaPub.Tests/Federation/JobHandlerTests.cs index dbb2486..84ff0ed 100644 --- a/PrivaPub.Tests/Federation/JobHandlerTests.cs +++ b/PrivaPub.Tests/Federation/JobHandlerTests.cs @@ -216,6 +216,103 @@ namespace PrivaPub.Tests.Federation Assert.All(_harness.Peer.Requests.Where(r => r.Path.EndsWith("/replies")), r => Assert.False(string.IsNullOrEmpty(r.Signature))); } + [Fact] + public async Task A_conversation_container_is_read_before_the_context_and_only_its_added_posts_taken() + { + var owner = new RemoteActor(_harness.Peer, "owner"); + var answerer = new RemoteActor(_harness.Peer, "answerer"); + var root = PublicNote(owner, "

a conversation

"); + var history = NewId(owner, "conversations"); + root["contextHistory"] = history; + root["context"] = NewId(owner, "never"); + Served(root); + var embedded = Served(Reply(answerer, "added whole", root)); + var referenced = Served(Reply(owner, "added by reference", embedded)); + var liked = Served(Reply(answerer, "only liked", root)); + var target = new JsonObject { ["type"] = "OrderedCollection", ["id"] = history, ["attributedTo"] = owner.Id }; + JsonObject Add(JsonNode activity) => new() + { + ["id"] = NewId(owner, "activities"), ["type"] = "Add", ["actor"] = owner.Id, ["object"] = activity, ["target"] = target.DeepClone() + }; + var create = Served(Create(owner, referenced)); + create["object"] = IdOf(referenced); + Served(create); + var byReference = Served(Add(IdOf(create))); + Served(new JsonObject + { + ["id"] = history, ["type"] = "OrderedCollection", ["attributedTo"] = owner.Id, ["collectionOf"] = "Activity", + ["orderedItems"] = new JsonArray(Create(owner, root), Add(Create(answerer, embedded)), IdOf(byReference), Add(Activity(answerer, "Like", JsonValue.Create(IdOf(liked))))) + }); + var held = await Held(root); + + Assert.Equal(1, await RunReplies(held)); + var stored = await Stored(embedded, referenced, liked); + Assert.Equal(new[] { IdOf(embedded), IdOf(referenced) }.Order(), stored.Keys.Order()); + Assert.Equal(held.ID, stored[IdOf(embedded)].AnsweringToPostId); + Assert.Equal(stored[IdOf(embedded)].ID, stored[IdOf(referenced)].AnsweringToPostId); + Assert.DoesNotContain(_harness.Peer.Requests, r => r.Path == new Uri(root["context"]!.GetValue()).AbsolutePath); + } + + [Fact] + public async Task An_unchanged_context_is_not_read_again_while_it_counts_the_thread() + { + var token = TestContext.Current.CancellationToken; + var poster = new RemoteActor(_harness.Peer, "poster"); + var answerer = new RemoteActor(_harness.Peer, "answerer"); + var root = PublicNote(poster, "

a topic

"); + var context = NewId(poster, "topics"); + var path = new Uri(context).AbsolutePath; + root["context"] = context; + Served(root); + var first = Served(Reply(answerer, "first", root)); + var later = Served(Reply(poster, "later", root)); + // NodeBB's: every post inline, counted, and an ETag without quotes + JsonObject Thread(params JsonObject[] notes) => new() + { + ["id"] = context, ["type"] = "OrderedCollection", ["totalItems"] = notes.Length, + ["orderedItems"] = new JsonArray(notes.Select(n => (JsonNode)IdOf(n)).ToArray()) + }; + _harness.Peer.ServeTagged(path, Thread(root, first).ToJsonString(), "4f930186eb"); + var held = await Held(root); + var handler = new RepliesJobHandler(_harness.Db, _harness.Remote, _harness.RemotePosts, _harness.Queue); + + await handler.Handle(RepliesJobHandler.For(held, 1), token); + Assert.Equal("4f930186eb", (await DB.Default.Find().Match(c => c.URI == context).ExecuteFirstAsync(token)).ETag); + // the same version: a 304, and nothing read + _harness.Peer.ServeTagged(path, Thread(root, first, later).ToJsonString(), "4f930186eb"); + await handler.Handle(RepliesJobHandler.For(held, 1), token); + Assert.Equal("4f930186eb", _harness.Peer.Requests.Last(r => r.Path == path).Headers["If-None-Match"]); + Assert.Single(await Stored(first, later)); + // a new one + _harness.Peer.ServeTagged(path, Thread(root, first, later).ToJsonString(), "5a0b"); + await handler.Handle(RepliesJobHandler.For(held, 1), token); + Assert.Equal(2, (await Stored(first, later)).Count); + Assert.Equal("5a0b", (await DB.Default.Find().Match(c => c.URI == context).ExecuteFirstAsync(token)).ETag); + } + + [Fact] + public async Task A_paged_contexts_etag_is_not_kept_since_its_first_page_hides_what_is_added_later() + { + var token = TestContext.Current.CancellationToken; + var poster = new RemoteActor(_harness.Peer, "poster"); + var root = PublicNote(poster, "

a long thread

"); + var context = NewId(poster, "contexts"); + root["context"] = context; + Served(root); + var answer = Served(Reply(poster, "on the second page", root)); + _harness.Peer.ServeTagged(new Uri(context).AbsolutePath, new JsonObject + { + ["id"] = context, ["type"] = "Collection", + ["first"] = new JsonObject { ["type"] = "CollectionPage", ["items"] = new JsonArray(IdOf(root)), ["next"] = context + "/2" } + }.ToJsonString(), "W/\"first\""); + Served(new JsonObject { ["id"] = context + "/2", ["type"] = "CollectionPage", ["items"] = new JsonArray(IdOf(answer)) }); + var held = await Held(root); + + await new RepliesJobHandler(_harness.Db, _harness.Remote, _harness.RemotePosts, _harness.Queue).Handle(RepliesJobHandler.For(held, 1), token); + Assert.Single(await Stored(answer)); + Assert.False(await DB.Default.Find().Match(c => c.URI == context && c.ETag != null).ExecuteAnyAsync(token)); + } + [Fact] public async Task A_threads_context_collection_brings_every_reply_at_once() { diff --git a/PrivaPub.Tests/Support/Peer.cs b/PrivaPub.Tests/Support/Peer.cs index 0596dea..8a503ff 100644 --- a/PrivaPub.Tests/Support/Peer.cs +++ b/PrivaPub.Tests/Support/Peer.cs @@ -21,6 +21,7 @@ namespace PrivaPub.Tests.Support readonly ConcurrentDictionary _documents = new(); readonly ConcurrentDictionary _answers = new(); readonly ConcurrentDictionary _files = new(); + readonly ConcurrentDictionary _etags = new(); public int Port { get; } public string A => $"http://127.0.0.1:{Port}"; @@ -64,6 +65,15 @@ namespace PrivaPub.Tests.Support context.Response.StatusCode = StatusCodes.Status404NotFound; return; } + if (peer._etags.TryGetValue(key, out var etag)) + { + if (context.Request.Headers.IfNoneMatch.ToString() == etag) + { + context.Response.StatusCode = StatusCodes.Status304NotModified; + return; + } + context.Response.Headers.ETag = etag; + } context.Response.ContentType = document.ContentType; await context.Response.WriteAsync(document.Text.Replace("{A}", peer.A).Replace("{B}", peer.B)); }); @@ -74,6 +84,13 @@ namespace PrivaPub.Tests.Support public void Serve(string path, string json) => _documents[path] = (json, "application/activity+json"); + // a document with an ETag: answered 304 to a request that names it in If-None-Match + public void ServeTagged(string path, string json, string etag) + { + _documents[path] = (json, "application/activity+json"); + _etags[path] = etag; + } + public void ServeText(string path, string text, string contentType) => _documents[path] = (text, contentType); public void ServeFile(string path, byte[] bytes, string contentType) => _files[path] = (bytes, contentType); diff --git a/PrivaPub/Federation/Actors/RemoteActorService.cs b/PrivaPub/Federation/Actors/RemoteActorService.cs index 8e2a2d1..efd412b 100644 --- a/PrivaPub/Federation/Actors/RemoteActorService.cs +++ b/PrivaPub/Federation/Actors/RemoteActorService.cs @@ -17,6 +17,8 @@ namespace PrivaPub.Federation.Actors public interface IRemoteActorService { Task FetchObject(string uri, CancellationToken token); + // the same, or NotModified while the version `etag` names is still current + Task FetchObject(string uri, string etag, CancellationToken token) => FetchObject(uri, token); Task GetActor(string actorUri, bool refresh, CancellationToken token); Task GetActorByKeyId(string keyId, bool refresh, CancellationToken token); bool KeyTemporarilyUnavailable(string keyId) => false; @@ -50,13 +52,17 @@ namespace PrivaPub.Federation.Actors _options = options; } - public async Task FetchObject(string uri, CancellationToken token) + public Task FetchObject(string uri, CancellationToken token) => FetchObject(uri, default, token); + + public async Task FetchObject(string uri, string etag, CancellationToken token) { using var scope = HttpScope.Default("object"); var signer = await _localActors.GetInstanceActor(token); - var fetched = await Get(uri, signer, token); - if (fetched == default) - return default; + var fetched = string.IsNullOrEmpty(etag) + ? await Get(uri, signer, token) + : await _http.GetJsonIfChanged(uri, Accept, etag, request => HttpSignatures.Sign(request, signer, body: null), token); + if (fetched == default || fetched.NotModified) + return fetched; var id = Text(fetched.Root, "id"); if (Origin.IsDocumentAt(id, fetched.FinalUri)) diff --git a/PrivaPub/Federation/Inbox/RemoteReplies.cs b/PrivaPub/Federation/Inbox/RemoteReplies.cs index 422c7f7..9abf5dd 100644 --- a/PrivaPub/Federation/Inbox/RemoteReplies.cs +++ b/PrivaPub/Federation/Inbox/RemoteReplies.cs @@ -20,9 +20,11 @@ namespace PrivaPub.Federation.Inbox public sealed record RepliesPayload(string PostId, int Depth); // What was said under a remote post on other servers, read when a persona opens its thread, so the thread shows more - // than what happened to reach us. The thread's own collection comes first (FEP-7888's `context`, which Mastodon 4.5+ - // serves with every reply at any depth); otherwise the post's `replies` (PeerTube's `comments`), and the replies' - // own, two levels down. Read with the instance actor's signature, never a persona's, a few pages and at most + // than what happened to reach us. The thread's own collections come first: its conversation container (FEP-171b's + // `contextHistory`, every activity its owner added to the thread), then its posts (FEP-7888's `context`, which + // Mastodon 4.5+, Mitra and NodeBB serve with every reply at any depth), sent the ETag of the last whole read so an + // unchanged thread costs one 304; otherwise the post's `replies` (PeerTube's `comments`, Forte's whole thread), and the + // replies' own, two levels down. Read with the instance actor's signature, never a persona's, a few pages and at most // MaxItems posts a job; each one is fetched from its own server and stored like any post fetched for a thread // (RemotePosts.StoreContext: public and unlisted only, its author checked against its origin). public class RepliesJobHandler : IJobHandler @@ -74,14 +76,35 @@ namespace PrivaPub.Federation.Inbox if (await Document(post, token) is not { } document) return JobOutcome.Done; - if (await Collection(document["context"], token) is { } thread) + foreach (var name in new[] { "contextHistory", "context" }) { - await StoreAll(thread, token); + var uri = Id(document[name]); + // Forte's `context` is the thread's first post + if (uri != default && await _dbEntities.Posts.Match(p => p.ObjectURI == uri).ExecuteAnyAsync(token)) + continue; + var known = uri == default ? default : await DB.Default.Find().Match(c => c.URI == uri).ExecuteFirstAsync(token); + var (thread, etag, unchanged) = await Collection(document[name], known?.ETag, token); + if (unchanged) + { + await DB.Default.Update().MatchID(known.ID).Modify(c => c.ReadAt, DateTime.UtcNow).ExecuteAsync(token); + return JobOutcome.Done; + } + if (thread == default) + continue; + var (_, read) = await StoreAll(thread, name == "contextHistory", token); + var whole = read && Counted(thread); + if (uri != default && (known != default || whole && etag != default)) + await DB.Default.Update() + .Match(c => c.URI == uri) + .Modify(c => c.ETag, whole ? etag : default) + .Modify(c => c.ReadAt, DateTime.UtcNow) + .Option(o => o.IsUpsert = true) + .ExecuteAsync(token); return JobOutcome.Done; } - if (await Collection(document["replies"] ?? document["comments"], token) is not { } replies) + if ((await Collection(document["replies"] ?? document["comments"], default, token)).Collection is not { } replies) return JobOutcome.Done; - var answers = await StoreAll(replies, token); + var (answers, _) = await StoreAll(replies, false, token); if (payload.Depth < MaxDepth) await _queue.EnqueueMany(answers.Where(a => a.IsFederatedCopy && a.AnsweringToPostId == post.ID).Take(MaxBranches) .Select(a => For(a, payload.Depth + 1)), token); @@ -100,35 +123,41 @@ namespace PrivaPub.Federation.Inbox return fetched == default ? default : JsonNode.Parse(fetched.Root.GetRawText()) as JsonObject; } - // a collection as its server serves it now: what a post carries is read again by its id, since an embedded first - // page is as old as the post; a value that is no web address (Pleroma's `context` is a tag: URI), or no collection, - // gives nothing - async Task Collection(JsonNode node, CancellationToken token) + // a collection as its server serves it now, and its ETag: what a post carries is read again by its id, since an + // embedded first page is as old as the post; a value that is no web address (Pleroma's `context` is a tag: URI), or no + // collection, gives nothing, and one whose server still has the version `etag` names is unchanged + async Task<(JsonObject Collection, string ETag, bool Unchanged)> Collection(JsonNode node, string etag, CancellationToken token) { var uri = Id(node); + string served = default; if (uri != default && Uri.TryCreate(uri, UriKind.Absolute, out var parsed) && (parsed.Scheme == Uri.UriSchemeHttps || parsed.Scheme == Uri.UriSchemeHttp)) { using var scope = HttpScope.For("replies"); - using var fetched = await _remoteActors.FetchObject(uri, token); + using var fetched = await _remoteActors.FetchObject(uri, etag, token); + if (fetched?.NotModified == true) + return (default, etag, true); node = fetched == default ? default : JsonNode.Parse(fetched.Root.GetRawText()); + served = fetched?.ETag; } - return node is JsonObject collection && Value(collection, "type") is { } type && CollectionTypes.Contains(type) ? collection : default; + return (node is JsonObject collection && Value(collection, "type") is { } type && CollectionTypes.Contains(type) ? collection : default, served, false); } - // every post the collection lists, stored when we lack it: its pages in order, at most MaxPages and MaxItems - async Task> StoreAll(JsonObject collection, CancellationToken token) + // every post the collection lists, stored when we lack it: its pages in order, at most MaxPages and MaxItems; whether + // it was read to its end. A conversation container's items are its owner's Adds of the thread's activities. + async Task<(List Posts, bool Whole)> StoreAll(JsonObject collection, bool activities, CancellationToken token) { var posts = new List(); var seen = 0; var page = collection["orderedItems"] is JsonArray || collection["items"] is JsonArray ? collection : await Page(collection["first"], token); - for (var pages = 0; page != default && pages < MaxPages && seen < MaxItems; pages++) + for (var pages = 0; page != default; pages++) { + if (pages == MaxPages) + return (posts, false); foreach (var item in Items(page)) { if (++seen > MaxItems) - break; - // a context may list the activities that made the thread (FEP-f228) rather than its posts - var id = Value(item, "type") is "Create" or "Update" ? Id(item["object"]) : Id(item); + return (posts, false); + var id = activities ? await Added(item, token) : PostId(item); if (id == default) continue; var post = await _dbEntities.Posts.Match(p => p.ObjectURI == id && !p.DeletedAt.HasValue).ExecuteFirstAsync(token) @@ -138,7 +167,32 @@ namespace PrivaPub.Federation.Inbox } page = page["next"] is { } next ? await Page(next, token) : default; } - return posts; + return (posts, true); + } + + // whether the document itself changes with every post added to the thread, so its ETag is the thread's: it counts + // them (NodeBB), or holds every one with no further page (a short Mastodon thread); a long Mastodon thread's first + // page stays as it was while its last one grows + static bool Counted(JsonObject collection) + { + if (collection["totalItems"] is JsonValue) + return true; + var page = collection["orderedItems"] is JsonArray || collection["items"] is JsonArray ? collection : collection["first"] as JsonObject; + return page != default && (page["orderedItems"] is JsonArray || page["items"] is JsonArray) && page["next"] == default; + } + + // a context may list the activities that made the thread (FEP-f228) rather than its posts + static string PostId(JsonNode item) => Value(item, "type") is "Create" or "Update" ? Id(item["object"]) : Id(item); + + // the post an activity in a conversation container brought into the thread, added by its owner or (Forte's) not: a + // Create's or an Update's object; a like, a removal or anything else adds no post. Only its id is taken, and the post + // is read from its own server. + async Task Added(JsonNode item, CancellationToken token) + { + var activity = await Page(item, token); + if (activity != default && Value(activity, "type") == "Add") + activity = await Page(activity["object"], token); + return activity != default && Value(activity, "type") is "Create" or "Update" ? Id(activity["object"]) : default; } async Task Page(JsonNode node, CancellationToken token) diff --git a/PrivaPub/Infrastructure/Data/Indexes.cs b/PrivaPub/Infrastructure/Data/Indexes.cs index 4eda5cc..226901a 100644 --- a/PrivaPub/Infrastructure/Data/Indexes.cs +++ b/PrivaPub/Infrastructure/Data/Indexes.cs @@ -15,6 +15,7 @@ namespace PrivaPub.Infrastructure.Data public static class Indexes { public static readonly TimeSpan TombstoneLifetime = TimeSpan.FromDays(90); + public static readonly TimeSpan ThreadCollectionLifetime = TimeSpan.FromDays(30); public static async Task Create(CancellationToken token = default) { @@ -166,6 +167,8 @@ namespace PrivaPub.Infrastructure.Data await Plain(token, r => r.PostId); await Unique(d => d.ObjectURI, Builders.Filter.Gt(d => d.ObjectURI, ""), token); await DB.Default.Index().Key(d => d.DeletedAt, KeyType.Ascending).Option(o => o.ExpireAfter = TombstoneLifetime).CreateAsync(token); + await Unique(c => c.URI, Builders.Filter.Gt(c => c.URI, ""), token); + await DB.Default.Index().Key(c => c.ReadAt, KeyType.Ascending).Option(o => o.ExpireAfter = ThreadCollectionLifetime).CreateAsync(token); await DB.Default.Index().Key(d => d.PostId, KeyType.Ascending).Key(d => d.ActorURI, KeyType.Ascending).Option(o => o.Unique = true).CreateAsync(token); await Plain(token, d => d.ActivityURI); await DB.Default.Index().Key(v => v.PostId, KeyType.Ascending).Key(v => v.ActorURI, KeyType.Ascending).Key(v => v.Choice, KeyType.Ascending) diff --git a/PrivaPub/Infrastructure/Http/FederationHttp.cs b/PrivaPub/Infrastructure/Http/FederationHttp.cs index da6cb1b..eb66c09 100644 --- a/PrivaPub/Infrastructure/Http/FederationHttp.cs +++ b/PrivaPub/Infrastructure/Http/FederationHttp.cs @@ -15,6 +15,9 @@ namespace PrivaPub.Infrastructure.Http public Uri FinalUri { get; init; } public JsonDocument Document { get; init; } public JsonElement Root => Document.RootElement; + public string ETag { get; init; } + // the origin answered 304 to an If-None-Match: no document, the one read before still stands + public bool NotModified { get; init; } public void Dispose() => Document?.Dispose(); } @@ -24,6 +27,9 @@ namespace PrivaPub.Infrastructure.Http bool IsAllowed(Uri target); Task GetJson(string url, string accept, Action sign, CancellationToken token); Task<(int Status, FetchedJson Json)> GetJsonStatus(string url, string accept, Action sign, CancellationToken token); + /// The document, or one that is when its origin still has the version + /// names. + Task GetJsonIfChanged(string url, string accept, string etag, Action sign, CancellationToken token); bool FailedTemporarily(string url); Task<(Uri FinalUri, string Html)> GetPage(string url, CancellationToken token); Task OpenMedia(string url, System.Net.Http.Headers.RangeHeaderValue range, CancellationToken token); @@ -291,7 +297,20 @@ namespace PrivaPub.Infrastructure.Http var exchange = new Exchange(url, HttpScope.Purpose ?? "object"); try { - return await GetJson(url, accept, sign, exchange, token); + return await GetJson(url, accept, sign, default, exchange, token); + } + finally + { + Record(exchange); + } + } + + public async Task GetJsonIfChanged(string url, string accept, string etag, Action sign, CancellationToken token) + { + var exchange = new Exchange(url, HttpScope.Purpose ?? "object"); + try + { + return await GetJson(url, accept, sign, etag, exchange, token); } finally { @@ -308,7 +327,7 @@ namespace PrivaPub.Infrastructure.Http var exchange = new Exchange(url, HttpScope.Purpose ?? "object"); try { - var json = await GetJson(url, accept, sign, exchange, token); + var json = await GetJson(url, accept, sign, default, exchange, token); return (exchange.Status ?? 0, json); } finally @@ -317,7 +336,7 @@ namespace PrivaPub.Infrastructure.Http } } - async Task GetJson(string url, string accept, Action sign, Exchange exchange, CancellationToken token) + async Task GetJson(string url, string accept, Action sign, string etag, Exchange exchange, CancellationToken token) { if (!Uri.TryCreate(url, UriKind.Absolute, out var target) || !IsAllowed(target)) { @@ -338,11 +357,15 @@ namespace PrivaPub.Infrastructure.Http for (var hop = 0; hop <= MaxRedirects; hop++) { using var request = Request(target, accept); + if (!string.IsNullOrEmpty(etag)) + request.Headers.TryAddWithoutValidation("If-None-Match", etag); sign?.Invoke(request); using var response = await _httpClientFactory.CreateClient(ClientName) .SendAsync(request, HttpCompletionOption.ResponseHeadersRead, timeout.Token); exchange.Answered(response); + if (response.StatusCode == HttpStatusCode.NotModified && !string.IsNullOrEmpty(etag)) + return new FetchedJson { FinalUri = target, ETag = etag, NotModified = true }; if (IsRedirect(response.StatusCode)) { @@ -370,7 +393,7 @@ namespace PrivaPub.Infrastructure.Http return Refuse(negativeKey, url, "a body over the size limit", exchange, "too-large"); exchange.Bytes = body.Length; - return new FetchedJson { FinalUri = target, Document = JsonDocument.Parse(body) }; + return new FetchedJson { FinalUri = target, Document = JsonDocument.Parse(body), ETag = RawETag(response) }; } return Refuse(negativeKey, url, "too many redirects", exchange, "too-many-redirects"); } @@ -658,6 +681,10 @@ namespace PrivaPub.Infrastructure.Http return buffer.ToArray(); } + // as sent, since it is only ever sent back: NodeBB's carries no quotes, which the typed header rejects + static string RawETag(HttpResponseMessage response) => + response.Headers.TryGetValues("ETag", out var values) ? values.FirstOrDefault() : default; + static bool IsRedirect(HttpStatusCode status) => status is HttpStatusCode.MovedPermanently or HttpStatusCode.Found or HttpStatusCode.SeeOther or HttpStatusCode.TemporaryRedirect or HttpStatusCode.PermanentRedirect; diff --git a/PrivaPub/Models/Federation/ThreadCollection.cs b/PrivaPub/Models/Federation/ThreadCollection.cs new file mode 100644 index 0000000..a151a82 --- /dev/null +++ b/PrivaPub/Models/Federation/ThreadCollection.cs @@ -0,0 +1,13 @@ +using MongoDB.Entities; + +namespace PrivaPub.Models.Federation +{ + // a thread's collection on another server as we last read it whole: its ETag, sent back as If-None-Match so an unchanged + // thread costs its server one 304 (NodeBB's context carries a digest of its posts) + public class ThreadCollection : Entity + { + public string URI { get; set; } + public string ETag { get; set; } + public DateTime ReadAt { get; set; } = DateTime.UtcNow;//expires after Indexes.ThreadCollectionLifetime + } +} diff --git a/docs/INTEROP.md b/docs/INTEROP.md index 4d4fcfd..d2e1e71 100644 --- a/docs/INTEROP.md +++ b/docs/INTEROP.md @@ -96,7 +96,9 @@ Priorities, used throughout: - `Move`. - `QuoteRequest`, its `Accept`/`Reject` (with `result`), and `Delete{QuoteAuthorization}`. - `FeatureRequest`/`FeatureAuthorization` (4.6). -- **Threads:** `context` is a dereferenceable collection of the thread (FEP-7888, threads started on 4.5+). +- **Threads:** `context` is a dereferenceable collection of the thread (FEP-7888, threads started on 4.5+). It + answers `If-None-Match` with a 304 (a weak ETag over the body, 4.7.3), but the body embeds only the first page, so its + ETag stands for the thread only while the thread fits on that page. - **Interaction counts:** `likes` and `shares` carry `totalItems`. - **Quotes:** `quote`, plus `quoteUri` and `_misskey_quote`, `quoteAuthorization`, and `interactionPolicy.canQuote`. - **Link previews:** from 4.7, a `{type: Link, href}` **attachment** names the link a card is made from (FEP-8967). @@ -659,7 +661,9 @@ What it showed: - **NodeBB:** - A topic's first post is an `Article` with `name`, **`summary` = an excerpt** and `preview`; replies are Notes. - Categories are `Group`s **without `followers`**. - - `context` is a paged collection with an **ETag digest**; NodeBB refetches with `If-None-Match`. + - `context` is a collection with an **ETag digest** of its posts, sent **without quotes** (not an RFC 9110 + entity-tag: .NET's typed header drops it, so PrivaPub keeps the raw value); it carries `totalItems` and its items + inline, and answers `If-None-Match` with a 304. NodeBB refetches with it too. - Since 4.15, an Announce of anything but a Create or a plain object is accepted only from Group actors. - It sends `Move`/`Remove` of a whole context (FEP-f15d) and `Add{post → context}` (FEP-11dd). - Chats are private Notes, threaded by `inReplyTo` into the room they answer. @@ -855,6 +859,8 @@ without a port, before it gives out its OAuth client. - Its Mastodon API resolves an account elsewhere only when the search is not limited to a type (`/api/v2/search?resolve=true`, no `type=accounts`); reactions go through Pleroma's route (`PUT /api/v1/pleroma/statuses/:id/reactions/:emoji`). +- **Threads:** a post's `context` is `/collections/conversations/`, a collection whose first page lists the + thread's post ids (FEP-7888), with no ETag; no `contextHistory`. - **Pasture evidence (2026-10-05, `tools/pasture/scenarios/mitra.sh`):** 28 checks pass, with no change to PrivaPub: follows both ways, posts, replies both ways, likes and reposts both ways, an emoji reaction, a poll and alice's vote, edits and deletions both ways, direct messages both ways, the unfollow, statistics. @@ -918,6 +924,11 @@ without a port, before it gives out its OAuth client. - **Everything in a thread is added to its conversation** (FEP-171b): its edits come only as `Add{Update}`, under the same activity id as the post's `Create`, and likes as `Add{Like}`. PrivaPub unwraps Create, Update and Delete, telling an Update from the Create whose id it reuses by what it carries. +- **Its conversation container is not named by its posts:** `contextHistory` (and `context`) pointing at + `…/conversation/` ride on the activities only; a Note's `context` is the thread's first post and it has no + `contextHistory`, though the FEP asks the first post for one. The container (public threads: unsigned GETs work) + lists the first post's plain `Create`, then an `Add{Create}` for every post. A post's `replies` lists every comment of + the thread, embedded, and is what PrivaPub reads to complete a Forte thread (checked live 2026-10-07). - **Its own collections name themselves wrongly:** a channel's followers, following and outbox give an `id` with the gateway path twice and a trailing `?`, which 404s, so their counts are not read. - **Running it:** no image is published; `images/forte` builds its tag with composer on Hubzilla's PHP image, behind an diff --git a/docs/ROADMAP.md b/docs/ROADMAP.md index 931985d..d89bf79 100644 --- a/docs/ROADMAP.md +++ b/docs/ROADMAP.md @@ -677,8 +677,11 @@ it, raw where it doesn't. #### P7 Threads, communities and the social graph - **Thread backfill:** - - read in order: `contextHistory`, then `context` (paged, with ETag), then `replies`. **Done (2026-10-05), but for - `contextHistory` and ETags:** a persona opening a public remote thread queues `FetchReplies` for the post and its + - read in order: `contextHistory`, then `context` (paged, with ETag), then `replies`. **Done (2026-10-05; `contextHistory` + and ETags 2026-10-07: a container's `Add`ed Creates and Updates are read for their posts; a whole read of a collection + that counts its posts or holds them all keeps its ETag, sent back as `If-None-Match`, a 304 checked live on Mastodon + 4.7.3 and NodeBB 4.16, whose ETag has no quotes; neither Forte nor Mitra puts `contextHistory` on a post, and Forte's + `replies` already holds its whole thread):** a persona opening a public remote thread queues `FetchReplies` for the post and its root, at most hourly each. The job reads the thread's FEP-7888 `context` collection (posts or, as FEP-f228 allows, the activities that made them), else the post's `replies` (PeerTube's `comments`) and theirs, two levels down. It reads at most 5 pages and 100 posts per job, signed by the instance actor, and stores what it lacks through