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
213 lines
9.7 KiB
C#
213 lines
9.7 KiB
C#
using System.Text.Json;
|
|
using System.Text.Json.Nodes;
|
|
|
|
using MongoDB.Entities;
|
|
|
|
using PrivaPub.Federation.Actors;
|
|
using PrivaPub.Infrastructure.Http;
|
|
using PrivaPub.Infrastructure.Jobs;
|
|
using PrivaPub.Models.Federation;
|
|
using PrivaPub.Models.Jobs;
|
|
using PrivaPub.Models.Post;
|
|
using PrivaPub.StaticServices;
|
|
|
|
using static PrivaPub.Federation.Objects.ActivityJson;
|
|
|
|
using PostEntity = PrivaPub.Models.Post.Post;
|
|
|
|
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 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
|
|
{
|
|
public const int MaxPages = 5;
|
|
public const int MaxItems = 100;
|
|
public const int MaxDepth = 2;
|
|
public const int MaxBranches = 20;
|
|
|
|
static readonly HashSet<string> CollectionTypes = new(StringComparer.Ordinal)
|
|
{
|
|
"Collection", "OrderedCollection", "CollectionPage", "OrderedCollectionPage"
|
|
};
|
|
|
|
readonly DbEntities _dbEntities;
|
|
readonly IRemoteActorService _remoteActors;
|
|
readonly IRemotePosts _remotePosts;
|
|
readonly IJobQueue _queue;
|
|
|
|
public RepliesJobHandler(DbEntities dbEntities, IRemoteActorService remoteActors, IRemotePosts remotePosts, IJobQueue queue)
|
|
{
|
|
_dbEntities = dbEntities;
|
|
_remoteActors = remoteActors;
|
|
_remotePosts = remotePosts;
|
|
_queue = queue;
|
|
}
|
|
|
|
public JobKind Kind => JobKind.FetchReplies;
|
|
public int Concurrency => 2;
|
|
public int MaxAttempts => 2;
|
|
public int PerHostLimit => 1;
|
|
|
|
// a thread is read again at most once an hour
|
|
public static Job For(PostEntity post, int depth) => new()
|
|
{
|
|
Kind = JobKind.FetchReplies,
|
|
Payload = JsonSerializer.Serialize(new RepliesPayload(post.ID, depth)),
|
|
Host = new Uri(post.ObjectURI).Host,
|
|
DedupeKey = $"replies|{post.ID}|{DateTime.UtcNow:yyyyMMddHH}"
|
|
};
|
|
|
|
public async Task<JobOutcome> Handle(Job job, CancellationToken token)
|
|
{
|
|
var payload = JsonSerializer.Deserialize<RepliesPayload>(job.Payload);
|
|
var post = await _dbEntities.Posts.MatchID(payload.PostId).ExecuteFirstAsync(token);
|
|
if (post is not { IsFederatedCopy: true, Visibility: PostVisibility.Public or PostVisibility.Unlisted } || post.DeletedAt.HasValue
|
|
|| string.IsNullOrEmpty(post.ObjectURI))
|
|
return JobOutcome.Done;
|
|
if (await Document(post, token) is not { } document)
|
|
return JobOutcome.Done;
|
|
|
|
foreach (var name in new[] { "contextHistory", "context" })
|
|
{
|
|
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"], default, token)).Collection is not { } replies)
|
|
return JobOutcome.Done;
|
|
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);
|
|
return JobOutcome.Done;
|
|
}
|
|
|
|
// the post as it reached us, or as its server serves it now when we kept no copy
|
|
async Task<JsonObject> Document(PostEntity post, CancellationToken token)
|
|
{
|
|
var record = await DB.Default.Find<ObjectRecord>().Match(r => r.ObjectURI == post.ObjectURI && r.Raw != null && !r.RawTruncated)
|
|
.Project(r => new ObjectRecord { Raw = r.Raw }).ExecuteFirstAsync(token);
|
|
if (record != default)
|
|
return JsonNode.Parse(record.Raw) as JsonObject;
|
|
using var scope = HttpScope.For("replies");
|
|
using var fetched = await _remoteActors.FetchObject(post.ObjectURI, token);
|
|
return fetched == default ? default : JsonNode.Parse(fetched.Root.GetRawText()) as JsonObject;
|
|
}
|
|
|
|
// 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, 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, served, false);
|
|
}
|
|
|
|
// 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++)
|
|
{
|
|
if (pages == MaxPages)
|
|
return (posts, false);
|
|
foreach (var item in Items(page))
|
|
{
|
|
if (++seen > MaxItems)
|
|
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)
|
|
?? await _remotePosts.StoreContext(id, 0, token);
|
|
if (post != default)
|
|
posts.Add(post);
|
|
}
|
|
page = page["next"] is { } next ? await Page(next, token) : default;
|
|
}
|
|
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)
|
|
{
|
|
if (node is JsonValue && Id(node) is { } uri)
|
|
{
|
|
using var scope = HttpScope.For("replies");
|
|
using var fetched = await _remoteActors.FetchObject(uri, token);
|
|
node = fetched == default ? default : JsonNode.Parse(fetched.Root.GetRawText());
|
|
}
|
|
return node as JsonObject;
|
|
}
|
|
|
|
static IEnumerable<JsonNode> Items(JsonNode page) =>
|
|
(page["orderedItems"] as JsonArray ?? page["items"] as JsonArray ?? new JsonArray()).Where(i => i != default);
|
|
}
|
|
}
|