P7: thread backfill reads conversation containers and asks with If-None-Match
Build / Build (push) Failing after 7m52s
Build / Build (push) Failing after 7m52s
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 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_012CzABvBkbcFqoHdmi8b9WB
This commit is contained in:
1 parent
fb4949b511
commit
62189bf11d
10 files changed
+268
-35
No files matched your search
@@ -17,6 +17,8 @@ namespace PrivaPub.Federation.Actors
|
||||
public interface IRemoteActorService
|
||||
{
|
||||
Task<FetchedJson> FetchObject(string uri, CancellationToken token);
|
||||
// the same, or NotModified while the version `etag` names is still current
|
||||
Task<FetchedJson> FetchObject(string uri, string etag, CancellationToken token) => FetchObject(uri, token);
|
||||
Task<ForeignAvatar> GetActor(string actorUri, bool refresh, CancellationToken token);
|
||||
Task<ForeignAvatar> 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<FetchedJson> FetchObject(string uri, CancellationToken token)
|
||||
public Task<FetchedJson> FetchObject(string uri, CancellationToken token) => FetchObject(uri, default, token);
|
||||
|
||||
public async Task<FetchedJson> 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))
|
||||
|
||||
@@ -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<ThreadCollection>().Match(c => c.URI == uri).ExecuteFirstAsync(token);
|
||||
var (thread, etag, unchanged) = await Collection(document[name], known?.ETag, token);
|
||||
if (unchanged)
|
||||
{
|
||||
await DB.Default.Update<ThreadCollection>().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<ThreadCollection>()
|
||||
.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<JsonObject> 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<List<PostEntity>> 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<PostEntity> Posts, bool Whole)> StoreAll(JsonObject collection, bool activities, CancellationToken token)
|
||||
{
|
||||
var posts = new List<PostEntity>();
|
||||
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<string> 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<JsonObject> Page(JsonNode node, CancellationToken token)
|
||||
|
||||
@@ -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<ObjectRecord>(token, r => r.PostId);
|
||||
await Unique<DeletedObject>(d => d.ObjectURI, Builders<DeletedObject>.Filter.Gt(d => d.ObjectURI, ""), token);
|
||||
await DB.Default.Index<DeletedObject>().Key(d => d.DeletedAt, KeyType.Ascending).Option(o => o.ExpireAfter = TombstoneLifetime).CreateAsync(token);
|
||||
await Unique<ThreadCollection>(c => c.URI, Builders<ThreadCollection>.Filter.Gt(c => c.URI, ""), token);
|
||||
await DB.Default.Index<ThreadCollection>().Key(c => c.ReadAt, KeyType.Ascending).Option(o => o.ExpireAfter = ThreadCollectionLifetime).CreateAsync(token);
|
||||
await DB.Default.Index<Downvote>().Key(d => d.PostId, KeyType.Ascending).Key(d => d.ActorURI, KeyType.Ascending).Option(o => o.Unique = true).CreateAsync(token);
|
||||
await Plain<Downvote>(token, d => d.ActivityURI);
|
||||
await DB.Default.Index<PollVote>().Key(v => v.PostId, KeyType.Ascending).Key(v => v.ActorURI, KeyType.Ascending).Key(v => v.Choice, KeyType.Ascending)
|
||||
|
||||
@@ -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<FetchedJson> GetJson(string url, string accept, Action<HttpRequestMessage> sign, CancellationToken token);
|
||||
Task<(int Status, FetchedJson Json)> GetJsonStatus(string url, string accept, Action<HttpRequestMessage> sign, CancellationToken token);
|
||||
/// <summary>The document, or one that is <see cref="FetchedJson.NotModified"/> when its origin still has the version
|
||||
/// <paramref name="etag"/> names.</summary>
|
||||
Task<FetchedJson> GetJsonIfChanged(string url, string accept, string etag, Action<HttpRequestMessage> sign, CancellationToken token);
|
||||
bool FailedTemporarily(string url);
|
||||
Task<(Uri FinalUri, string Html)> GetPage(string url, CancellationToken token);
|
||||
Task<HttpResponseMessage> 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<FetchedJson> GetJsonIfChanged(string url, string accept, string etag, Action<HttpRequestMessage> 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<FetchedJson> GetJson(string url, string accept, Action<HttpRequestMessage> sign, Exchange exchange, CancellationToken token)
|
||||
async Task<FetchedJson> GetJson(string url, string accept, Action<HttpRequestMessage> 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;
|
||||
|
||||
@@ -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
|
||||
}
|
||||
}
|
||||
Reference in new issue
Block a user