Mastodon answers 422 when two first contacts from one actor race to create its account (ActiveRecord::RecordInvalid on the unique uri), and 409 while another worker holds its lock. A persona that followed two Mastodon accounts at once had one Follow refused that way; the job died on its first attempt and the persona waited on "requested" forever. Both answers are now retried twice on the usual backoff before they count as refusals. Found by the town (a village of 23 accounts on seven servers). Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01LsXgEaXee4GCU1hwYgPJXw
554 lines
27 KiB
C#
554 lines
27 KiB
C#
using Microsoft.AspNetCore.Builder;
|
|
using Microsoft.AspNetCore.Hosting;
|
|
using Microsoft.AspNetCore.Http;
|
|
using Microsoft.Extensions.Caching.Memory;
|
|
using Microsoft.Extensions.DependencyInjection;
|
|
using Microsoft.Extensions.Logging.Abstractions;
|
|
|
|
using MongoDB.Entities;
|
|
|
|
using OpenIddict.Abstractions;
|
|
|
|
using PrivaPub.Api.Mastodon.Auth;
|
|
using PrivaPub.Domain.Content;
|
|
using PrivaPub.Domain.Media;
|
|
using PrivaPub.Domain.Statuses;
|
|
using PrivaPub.Federation.Actors;
|
|
using PrivaPub.Federation.Inbox;
|
|
using PrivaPub.Federation.Objects;
|
|
using PrivaPub.Federation.Outbox;
|
|
using PrivaPub.Federation.Signing;
|
|
using PrivaPub.Infrastructure.Http;
|
|
using PrivaPub.Infrastructure.Jobs;
|
|
using PrivaPub.Models.Federation;
|
|
using PrivaPub.Models.Jobs;
|
|
using PrivaPub.Models.Media;
|
|
using PrivaPub.Models.Post;
|
|
using PrivaPub.Models.Social;
|
|
using PrivaPub.Tests.Support;
|
|
using PrivaPub.Tests.Support.Host;
|
|
|
|
using System.Security.Cryptography;
|
|
using System.Text;
|
|
using System.Text.Json;
|
|
using System.Text.Json.Nodes;
|
|
|
|
using static PrivaPub.Tests.Support.FederatedSeeds;
|
|
|
|
namespace PrivaPub.Tests.Federation
|
|
{
|
|
[Trait("Category", "Integration")]
|
|
public sealed class JobHandlerTests : IAsyncLifetime
|
|
{
|
|
Harness _harness;
|
|
|
|
public async ValueTask InitializeAsync()
|
|
{
|
|
Assert.SkipUnless(MongoFixture.Enabled, MongoFixture.Skip);
|
|
_harness = await Harness.Start();
|
|
}
|
|
|
|
public async ValueTask DisposeAsync()
|
|
{
|
|
if (_harness != default)
|
|
await _harness.DisposeAsync();
|
|
}
|
|
|
|
sealed class PeerDescriber : InstanceDescriber
|
|
{
|
|
readonly string _peer;
|
|
|
|
public PeerDescriber(IFederationHttp http, string peer) : base(http) => _peer = peer;
|
|
|
|
protected override string Address(string url) => _peer + new Uri(url).PathAndQuery;
|
|
}
|
|
|
|
[Fact]
|
|
public async Task Ancestors_are_fetched_one_job_at_a_time_and_stop_at_the_depth_limit()
|
|
{
|
|
var token = TestContext.Current.CancellationToken;
|
|
var poster = new RemoteActor(_harness.Peer, "poster");
|
|
var chain = new List<JsonObject>();
|
|
for (var i = 0; i < RemotePosts.MaxDepth + 3; i++)
|
|
{
|
|
var note = PublicNote(poster, $"<p>number {i}</p>");
|
|
if (i > 0)
|
|
note["inReplyTo"] = IdOf(chain[^1]);
|
|
_harness.Peer.Serve(new Uri(IdOf(note)).AbsolutePath, note.ToJsonString());
|
|
chain.Add(note);
|
|
}
|
|
var leaf = new Post
|
|
{
|
|
ObjectURI = NewId(poster, "notes"), ActorURI = poster.Id, IsFederatedCopy = true, Visibility = PostVisibility.Public,
|
|
InReplyToURI = IdOf(chain[^1]), ContentHtml = "<p>the leaf</p>"
|
|
};
|
|
await DB.Default.SaveAsync(leaf, token);
|
|
await _harness.Queue.Enqueue(JobKind.FetchAncestors, JsonSerializer.Serialize(new AncestorsPayload(leaf.ID, 1)), "127.0.0.1", "ancestors|" + leaf.ID, token);
|
|
var handler = new AncestorsJobHandler(_harness.Db, _harness.RemotePosts);
|
|
|
|
var keys = new List<string> { "ancestors|" + leaf.ID };
|
|
var ran = 0;
|
|
while (await DB.Default.Find<Job>().Match(j => j.Kind == JobKind.FetchAncestors && j.State == JobState.Pending && keys.Contains(j.DedupeKey))
|
|
.ExecuteFirstAsync(token) is { } job)
|
|
{
|
|
Assert.Equal(JobResult.Done, (await handler.Handle(job, token)).Result);
|
|
await DB.Default.Update<Job>().MatchID(job.ID).Modify(j => j.State, JobState.Done).ExecuteAsync(token);
|
|
ran++;
|
|
var uris = chain.Select(IdOf).ToList();
|
|
keys = (await DB.Default.Find<Post>().Match(p => uris.Contains(p.ObjectURI)).ExecuteAsync(token)).Select(p => "ancestors|" + p.ID).Append("ancestors|" + leaf.ID).ToList();
|
|
Assert.True(ran <= RemotePosts.MaxDepth, "the chain did not stop");
|
|
}
|
|
|
|
var all = chain.Select(IdOf).ToList();
|
|
var stored = (await DB.Default.Find<Post>().Match(p => all.Contains(p.ObjectURI)).ExecuteAsync(token)).ToDictionary(p => p.ObjectURI);
|
|
Assert.Equal(RemotePosts.MaxDepth, ran);
|
|
Assert.Equal(all.Skip(3), stored.Keys.OrderBy(all.IndexOf));
|
|
Assert.Equal(stored[all[^1]].ID, (await DB.Default.Find<Post>().OneAsync(leaf.ID, token)).AnsweringToPostId);
|
|
for (var i = 4; i < all.Count; i++)
|
|
Assert.Equal(stored[all[i - 1]].ID, stored[all[i]].AnsweringToPostId);
|
|
Assert.Null(stored[all[3]].AnsweringToPostId);
|
|
Assert.Equal(all[2], stored[all[3]].InReplyToURI);
|
|
Assert.DoesNotContain(_harness.Peer.Requests, r => r.Path == new Uri(all[2]).AbsolutePath);
|
|
}
|
|
|
|
async Task<(LocalActor Alice, RemoteActor Mallory, Post Poll)> LocalPoll()
|
|
{
|
|
var token = TestContext.Current.CancellationToken;
|
|
var (_, alice) = await _harness.Persona("alice");
|
|
var mallory = new RemoteActor(_harness.Peer, "mallory");
|
|
await _harness.Deliver(mallory, "/human-centipede", new JsonObject
|
|
{
|
|
["id"] = NewId(mallory, "follows"), ["type"] = "Follow", ["actor"] = mallory.Id, ["object"] = alice.Uri
|
|
});
|
|
var outcome = await _harness.Statuses.Publish(alice, new StatusDraft { Text = "which?", Poll = new PollDraft(new[] { "tea", "coffee" }, 3600, false, false) }, token);
|
|
Assert.True(outcome.Ok);
|
|
return (alice, mallory, outcome.Post);
|
|
}
|
|
|
|
async Task<List<JsonObject>> Updates(RemoteActor to) =>
|
|
(await _harness.Outgoing(to.Id + "/inbox")).Where(a => a["type"]!.GetValue<string>() == "Update").ToList();
|
|
|
|
static int[] Counts(JsonNode question) =>
|
|
question["oneOf"]!.AsArray().Select(o => o!["replies"]!["totalItems"]!.GetValue<int>()).ToArray();
|
|
|
|
[Fact]
|
|
public async Task A_poll_refresh_sends_the_new_counts_to_followers_as_an_update()
|
|
{
|
|
var token = TestContext.Current.CancellationToken;
|
|
var (_, mallory, poll) = await LocalPoll();
|
|
var (_, bob) = await _harness.Persona("bob");
|
|
Assert.True((await _harness.Polls.Vote(bob, poll, new[] { 1 }, token)).Ok);
|
|
var job = await DB.Default.Find<Job>().Match(j => j.Kind == JobKind.PollRefresh && j.Payload == poll.ID).ExecuteSingleAsync(token);
|
|
Assert.True(job.RunAt > DateTime.UtcNow);
|
|
Assert.Empty(await Updates(mallory));
|
|
|
|
Assert.Equal(JobResult.Done, (await new PollRefreshJob(_harness.Db, _harness.Local, _harness.Outbox).Handle(job, token)).Result);
|
|
|
|
var question = Assert.Single(await Updates(mallory))["object"]!;
|
|
Assert.Equal("Question", question["type"]!.GetValue<string>());
|
|
Assert.Equal(new[] { 0, 1 }, Counts(question));
|
|
Assert.Null(question["closed"]);
|
|
Assert.Null((await DB.Default.Find<Post>().OneAsync(poll.ID, token)).Poll.ClosedAt);
|
|
}
|
|
|
|
[Fact]
|
|
public async Task Closing_a_local_poll_tells_its_voters_and_sends_the_closed_question()
|
|
{
|
|
var token = TestContext.Current.CancellationToken;
|
|
var (alice, mallory, poll) = await LocalPoll();
|
|
var (_, bob) = await _harness.Persona("bob");
|
|
Assert.True((await _harness.Polls.Vote(bob, poll, new[] { 0 }, token)).Ok);
|
|
var job = await DB.Default.Find<Job>().Match(j => j.Kind == JobKind.PollClose && j.Payload == poll.ID).ExecuteSingleAsync(token);
|
|
|
|
Assert.Equal(JobResult.Done, (await new PollCloseJob(_harness.Db, _harness.Local, _harness.Outbox).Handle(job, token)).Result);
|
|
|
|
var closed = (await DB.Default.Find<Post>().OneAsync(poll.ID, token)).Poll.ClosedAt;
|
|
Assert.NotNull(closed);
|
|
var question = Assert.Single(await Updates(mallory))["object"]!;
|
|
Assert.NotNull(question["closed"]);
|
|
Assert.Equal(new[] { 1, 0 }, Counts(question));
|
|
var told = await DB.Default.Find<Notification>().Match(n => n.Type == NotificationType.Poll && n.PostId == poll.ID).ExecuteAsync(token);
|
|
Assert.Equal(new[] { alice.Id, bob.Id }.Order(), told.Select(n => n.AvatarId).Order());
|
|
}
|
|
|
|
[Fact]
|
|
public async Task Closing_a_remote_poll_tells_the_local_voters_and_sends_nothing()
|
|
{
|
|
var token = TestContext.Current.CancellationToken;
|
|
var (_, alice) = await _harness.Persona("alice");
|
|
var mallory = new RemoteActor(_harness.Peer, "mallory");
|
|
var question = PublicNote(mallory, "<p>pick</p>", alice.Uri);
|
|
question["type"] = "Question";
|
|
question["tag"] = new JsonArray(new JsonObject { ["type"] = "Mention", ["href"] = alice.Uri, ["name"] = "@alice" });
|
|
question["endTime"] = DateTime.UtcNow.AddMinutes(10).ToString("O");
|
|
question["oneOf"] = new JsonArray(
|
|
new JsonObject { ["type"] = "Note", ["name"] = "red", ["replies"] = new JsonObject { ["type"] = "Collection", ["totalItems"] = 0 } },
|
|
new JsonObject { ["type"] = "Note", ["name"] = "blue", ["replies"] = new JsonObject { ["type"] = "Collection", ["totalItems"] = 0 } });
|
|
await _harness.Deliver(mallory, "/human-centipede", Create(mallory, question));
|
|
var poll = await DB.Default.Find<Post>().Match(p => p.ObjectURI == IdOf(question)).ExecuteSingleAsync(token);
|
|
Assert.True((await _harness.Polls.Vote(alice, poll, new[] { 1 }, token)).Ok);
|
|
var job = await DB.Default.Find<Job>().Match(j => j.Kind == JobKind.PollClose && j.Payload == poll.ID).ExecuteSingleAsync(token);
|
|
|
|
Assert.Equal(JobResult.Done, (await new PollCloseJob(_harness.Db, _harness.Local, _harness.Outbox).Handle(job, token)).Result);
|
|
|
|
var told = Assert.Single(await DB.Default.Find<Notification>().Match(n => n.Type == NotificationType.Poll && n.PostId == poll.ID).ExecuteAsync(token));
|
|
Assert.Equal(alice.Id, told.AvatarId);
|
|
Assert.Equal(poll.AuthorAccountId, told.FromAccountId);
|
|
Assert.Empty(await Updates(mallory));
|
|
Assert.Null((await DB.Default.Find<Post>().OneAsync(poll.ID, token)).Poll.ClosedAt);
|
|
}
|
|
|
|
[Fact]
|
|
public async Task A_server_is_described_from_its_nodeinfo_and_at_most_once_a_week()
|
|
{
|
|
var token = TestContext.Current.CancellationToken;
|
|
var host = $"describe{Guid.NewGuid():N}.example";
|
|
var peer = _harness.Peer;
|
|
peer.ServeText("/.well-known/nodeinfo", new JsonObject
|
|
{
|
|
["links"] = new JsonArray(
|
|
new JsonObject { ["rel"] = "http://nodeinfo.diaspora.software/ns/schema/2.0", ["href"] = $"https://{host}/nodeinfo/2.0" },
|
|
new JsonObject { ["rel"] = "http://nodeinfo.diaspora.software/ns/schema/2.1", ["href"] = $"https://{host}/nodeinfo/2.1" })
|
|
}.ToJsonString(), "application/json");
|
|
peer.ServeText("/nodeinfo/2.1", new JsonObject
|
|
{
|
|
["version"] = "2.1",
|
|
["software"] = new JsonObject { ["name"] = "gotosocial", ["version"] = "0.20.1" },
|
|
["protocols"] = new JsonArray("activitypub", 42),
|
|
["openRegistrations"] = false,
|
|
["usage"] = new JsonObject { ["users"] = new JsonObject { ["total"] = 3 } },
|
|
["metadata"] = new JsonObject { ["nodeName"] = "A small place" }
|
|
}.ToJsonString(), "application/json");
|
|
var describer = new PeerDescriber(Peer.Http(), peer.A);
|
|
|
|
async Task Record()
|
|
{
|
|
var note = NoteParser.Parse(new JsonObject
|
|
{
|
|
["id"] = $"https://{host}/notes/{Guid.NewGuid():N}", ["type"] = "Note", ["attributedTo"] = $"https://{host}/users/x",
|
|
["content"] = "<p>hi</p>", ["to"] = new JsonArray(Addressing.Public)
|
|
});
|
|
var post = new Post { ObjectURI = note.Id, IsFederatedCopy = true };
|
|
await DB.Default.SaveAsync(post, token);
|
|
await _harness.Records.Record(note, post, ObjectPath.Fetched, refetched: false, token);
|
|
}
|
|
Task<List<Job>> Jobs() => DB.Default.Find<Job>().Match(j => j.Kind == JobKind.DescribeInstance && j.Payload == host).ExecuteAsync(token);
|
|
|
|
await Record();
|
|
await Record();
|
|
var job = Assert.Single(await Jobs());
|
|
Assert.StartsWith($"describe|{host}|", job.DedupeKey);
|
|
|
|
Assert.Equal(JobResult.Done, (await describer.Handle(job, token)).Result);
|
|
|
|
var instance = await DB.Default.Find<RemoteInstance>().Match(i => i.Host == host).ExecuteSingleAsync(token);
|
|
Assert.Equal(("gotosocial", "0.20.1", "A small place"), (instance.Software, instance.SoftwareVersion, instance.NodeName));
|
|
Assert.Equal(new[] { "activitypub" }, instance.Protocols);
|
|
Assert.False(instance.OpenRegistrations);
|
|
Assert.Equal(3, JsonNode.Parse(instance.NodeInfo)!["usage"]!["users"]!["total"]!.GetValue<int>());
|
|
Assert.Null(instance.DescriptionError);
|
|
Assert.InRange(instance.DescribedAt!.Value, DateTime.UtcNow.AddMinutes(-1), DateTime.UtcNow.AddSeconds(1));
|
|
Assert.DoesNotContain(peer.Requests, r => r.Path == "/nodeinfo/2.0");
|
|
|
|
await DB.Default.DeleteAsync<Job>(job.ID);
|
|
await Record();
|
|
Assert.Empty(await Jobs());
|
|
|
|
await DB.Default.Update<RemoteInstance>().MatchID(instance.ID).Modify(i => i.DescribedAt, DateTime.UtcNow.AddDays(-8)).ExecuteAsync(token);
|
|
await Record();
|
|
Assert.Single(await Jobs());
|
|
}
|
|
|
|
[Fact]
|
|
public async Task A_server_without_usable_nodeinfo_is_described_as_such()
|
|
{
|
|
var token = TestContext.Current.CancellationToken;
|
|
var plain = $"plain{Guid.NewGuid():N}.example";
|
|
var describer = new PeerDescriber(Peer.Http(), _harness.Peer.A + "/" + plain);
|
|
_harness.Peer.ServeText($"/{plain}/.well-known/nodeinfo", new JsonObject
|
|
{
|
|
["links"] = new JsonArray(new JsonObject { ["rel"] = "http://nodeinfo.diaspora.software/ns/schema/2.1", ["href"] = $"http://{plain}/nodeinfo/2.1" })
|
|
}.ToJsonString(), "application/json");
|
|
var silent = $"silent{Guid.NewGuid():N}.example";
|
|
|
|
await describer.Handle(new Job { Kind = JobKind.DescribeInstance, Payload = plain }, token);
|
|
await new PeerDescriber(Peer.Http(), _harness.Peer.A + "/" + silent).Handle(new Job { Kind = JobKind.DescribeInstance, Payload = silent }, token);
|
|
|
|
Assert.Equal("no NodeInfo", (await DB.Default.Find<RemoteInstance>().Match(i => i.Host == plain).ExecuteSingleAsync(token)).DescriptionError);
|
|
Assert.Equal("no NodeInfo", (await DB.Default.Find<RemoteInstance>().Match(i => i.Host == silent).ExecuteSingleAsync(token)).DescriptionError);
|
|
}
|
|
|
|
LinkPreviews Previews() => new(_harness.Db, _harness.Local, Peer.Http(), _harness.Queue,
|
|
new StaticOptions<FederationOptions>(new FederationOptions { AllowPrivateNetworks = true, AllowPlainHttp = true }));
|
|
|
|
async Task<Post> Linking(PostVisibility visibility, string path)
|
|
{
|
|
var post = new Post
|
|
{
|
|
ObjectURI = $"https://elsewhere.example/notes/{Guid.NewGuid():N}", IsFederatedCopy = true, Visibility = visibility,
|
|
ContentHtml = $"<p>read <a href=\"{_harness.Peer.A}{path}\">this</a></p>"
|
|
};
|
|
await DB.Default.SaveAsync(post, TestContext.Current.CancellationToken);
|
|
_harness.Peer.ServeText(path, """
|
|
<html><head><title>Fallback</title>
|
|
<meta property="og:title" content="A headline">
|
|
<meta property="og:description" content="What it says">
|
|
<meta property="og:image" content="/cover.jpg">
|
|
</head><body></body></html>
|
|
""", "text/html; charset=utf-8");
|
|
return post;
|
|
}
|
|
|
|
[Fact]
|
|
public async Task A_link_preview_is_built_from_open_graph_for_a_public_post_only()
|
|
{
|
|
var token = TestContext.Current.CancellationToken;
|
|
var publicPath = $"/articles/{Guid.NewGuid():N}";
|
|
var privatePath = $"/articles/{Guid.NewGuid():N}";
|
|
var open = await Linking(PostVisibility.Public, publicPath);
|
|
var followersOnly = await Linking(PostVisibility.FollowersOnly, privatePath);
|
|
|
|
Assert.Equal(JobResult.Done, (await Previews().Handle(new Job { Kind = JobKind.FetchPreview, Payload = open.ID }, token)).Result);
|
|
Assert.Equal(JobResult.Done, (await Previews().Handle(new Job { Kind = JobKind.FetchPreview, Payload = followersOnly.ID }, token)).Result);
|
|
|
|
var card = (await DB.Default.Find<Post>().OneAsync(open.ID, token)).Link;
|
|
Assert.Equal(_harness.Peer.A + publicPath, card.Href);
|
|
Assert.Equal(("A headline", "What it says", _harness.Peer.A + "/cover.jpg"), (card.Title, card.Description, card.ImageURL));
|
|
Assert.True(await DB.Default.Find<LinkPreview>().Match(p => p.Url == _harness.Peer.A + publicPath).ExecuteAnyAsync(token));
|
|
Assert.Null((await DB.Default.Find<Post>().OneAsync(followersOnly.ID, token)).Link);
|
|
Assert.DoesNotContain(_harness.Peer.Requests, r => r.Path == privatePath);
|
|
Assert.False(await DB.Default.Find<LinkPreview>().Match(p => p.Url == _harness.Peer.A + privatePath).ExecuteAnyAsync(token));
|
|
}
|
|
|
|
[Fact]
|
|
public async Task The_janitor_removes_stale_unattached_uploads_and_trims_the_proxy_cache_oldest_first()
|
|
{
|
|
var token = TestContext.Current.CancellationToken;
|
|
var root = Path.Combine(Path.GetTempPath(), $"privapub-janitor-{Guid.NewGuid():N}");
|
|
var options = new StaticOptions<MediaOptions>(new MediaOptions { Root = root, ProxyCacheBytes = 800 });
|
|
var media = new MediaService(options, _harness.Local, default, NullLogger<MediaService>.Instance);
|
|
var seeded = new List<string>();
|
|
try
|
|
{
|
|
MediaAttachment Upload(string name, DateTime created, string postId = default)
|
|
{
|
|
Directory.CreateDirectory(Path.Combine(root, "files"));
|
|
File.WriteAllBytes(Path.Combine(root, "files", name), new byte[] { 1, 2, 3 });
|
|
File.WriteAllBytes(Path.Combine(root, "files", "small-" + name), new byte[] { 1 });
|
|
return new MediaAttachment { FilePath = Path.Combine("files", name), PreviewPath = Path.Combine("files", "small-" + name), CreatedAt = created, PostId = postId };
|
|
}
|
|
var stale = Upload("stale.webp", DateTime.UtcNow.AddDays(-2));
|
|
var fresh = Upload("fresh.webp", DateTime.UtcNow.AddHours(-1));
|
|
var attached = Upload("attached.webp", DateTime.UtcNow.AddDays(-3), "000000000000000000000001");
|
|
await DB.Default.SaveAsync(new[] { stale, fresh, attached }, token);
|
|
seeded.AddRange(new[] { stale.ID, fresh.ID, attached.ID });
|
|
Directory.CreateDirectory(media.ProxyRoot);
|
|
FileInfo Cached(string name, int bytes, int minutesAgo)
|
|
{
|
|
var path = Path.Combine(media.ProxyRoot, name);
|
|
File.WriteAllBytes(path, new byte[bytes]);
|
|
File.WriteAllText(path + ".type", "image/png");
|
|
File.SetLastWriteTimeUtc(path, DateTime.UtcNow.AddMinutes(-minutesAgo));
|
|
return new FileInfo(path);
|
|
}
|
|
var oldest = Cached("a", 600, 30);
|
|
var older = Cached("b", 400, 20);
|
|
var newest = Cached("c", 300, 10);
|
|
|
|
await new MediaJanitor(media, options, NullLogger<MediaJanitor>.Instance).Sweep(token);
|
|
|
|
Assert.False(await DB.Default.Find<MediaAttachment>().Match(m => m.ID == stale.ID).ExecuteAnyAsync(token));
|
|
Assert.False(File.Exists(Path.Combine(root, stale.FilePath)));
|
|
Assert.False(File.Exists(Path.Combine(root, stale.PreviewPath)));
|
|
Assert.True(await DB.Default.Find<MediaAttachment>().Match(m => m.ID == fresh.ID).ExecuteAnyAsync(token));
|
|
Assert.True(await DB.Default.Find<MediaAttachment>().Match(m => m.ID == attached.ID).ExecuteAnyAsync(token));
|
|
Assert.True(File.Exists(Path.Combine(root, fresh.FilePath)));
|
|
Assert.False(File.Exists(oldest.FullName));
|
|
Assert.False(File.Exists(oldest.FullName + ".type"));
|
|
Assert.True(File.Exists(older.FullName));
|
|
Assert.True(File.Exists(older.FullName + ".type"));
|
|
Assert.True(File.Exists(newest.FullName));
|
|
}
|
|
finally
|
|
{
|
|
await DB.Default.DeleteAsync<MediaAttachment>(m => seeded.Contains(m.ID));
|
|
foreach (var directory in new[] { root, media.ProxyRoot }.Where(Directory.Exists))
|
|
Directory.Delete(directory, recursive: true);
|
|
}
|
|
}
|
|
}
|
|
|
|
[Trait("Category", "Integration")]
|
|
public sealed class OAuthPrunerTests
|
|
{
|
|
[Fact]
|
|
public async Task Old_tokens_and_authorizations_that_are_no_longer_valid_are_pruned_and_the_rest_kept()
|
|
{
|
|
Assert.SkipUnless(MongoFixture.Enabled, MongoFixture.Skip);
|
|
var token = TestContext.Current.CancellationToken;
|
|
var host = await PrivaPubHost.Shared();
|
|
var now = DateTimeOffset.UtcNow;
|
|
using var scope = host.Services.CreateScope();
|
|
var tokens = scope.ServiceProvider.GetRequiredService<IOpenIddictTokenManager>();
|
|
var authorizations = scope.ServiceProvider.GetRequiredService<IOpenIddictAuthorizationManager>();
|
|
var subject = $"pruned{Guid.NewGuid():N}";
|
|
async Task<string> Token(DateTimeOffset created, DateTimeOffset expires, string status) =>
|
|
await tokens.GetIdAsync(await tokens.CreateAsync(new OpenIddictTokenDescriptor
|
|
{
|
|
CreationDate = created, ExpirationDate = expires, Status = status, Subject = subject, Type = OpenIddictConstants.TokenTypeHints.AccessToken
|
|
}, token), token);
|
|
async Task<string> Authorization(DateTimeOffset created, string status) =>
|
|
await authorizations.GetIdAsync(await authorizations.CreateAsync(new OpenIddictAuthorizationDescriptor
|
|
{
|
|
CreationDate = created, Status = status, Subject = subject, Type = OpenIddictConstants.AuthorizationTypes.Permanent
|
|
}, token), token);
|
|
var expired = await Token(now.AddDays(-30), now.AddDays(-29), OpenIddictConstants.Statuses.Valid);
|
|
var revoked = await Token(now.AddDays(-30), now.AddDays(1), OpenIddictConstants.Statuses.Revoked);
|
|
var current = await Token(now.AddMinutes(-5), now.AddHours(1), OpenIddictConstants.Statuses.Valid);
|
|
var recentlyRevoked = await Token(now.AddDays(-1), now.AddDays(-1).AddMinutes(1), OpenIddictConstants.Statuses.Revoked);
|
|
var oldRevoked = await Authorization(now.AddDays(-30), OpenIddictConstants.Statuses.Revoked);
|
|
var oldValid = await Authorization(now.AddDays(-30), OpenIddictConstants.Statuses.Valid);
|
|
|
|
var (prunedTokens, prunedAuthorizations) = await new OAuthPruner(host.Services, NullLogger<OAuthPruner>.Instance).Prune(now.AddDays(-14), token);
|
|
|
|
Assert.True(prunedTokens >= 2, $"{prunedTokens} tokens pruned");
|
|
Assert.True(prunedAuthorizations >= 1, $"{prunedAuthorizations} authorizations pruned");
|
|
//a fresh scope: the managers cache what they created
|
|
using var after = host.Services.CreateScope();
|
|
tokens = after.ServiceProvider.GetRequiredService<IOpenIddictTokenManager>();
|
|
authorizations = after.ServiceProvider.GetRequiredService<IOpenIddictAuthorizationManager>();
|
|
Assert.Null(await tokens.FindByIdAsync(expired, token));
|
|
Assert.Null(await tokens.FindByIdAsync(revoked, token));
|
|
Assert.NotNull(await tokens.FindByIdAsync(current, token));
|
|
Assert.NotNull(await tokens.FindByIdAsync(recentlyRevoked, token));
|
|
Assert.Null(await authorizations.FindByIdAsync(oldRevoked, token));
|
|
Assert.NotNull(await authorizations.FindByIdAsync(oldValid, token));
|
|
}
|
|
}
|
|
|
|
[Xunit.Collection(nameof(Exclusive))]
|
|
[Trait("Category", "Integration")]
|
|
public sealed class DeliveryOutcomeTests : IAsyncLifetime
|
|
{
|
|
static readonly string[] PeerHosts = { "localhost", "127.0.0.1" };
|
|
|
|
Harness _harness;
|
|
WebApplication _hinting;
|
|
string _hintingBase;
|
|
|
|
public async ValueTask InitializeAsync()
|
|
{
|
|
Assert.SkipUnless(MongoFixture.Enabled, MongoFixture.Skip);
|
|
_harness = await Harness.Start();
|
|
await DB.Default.DeleteAsync<RemoteInstance>(i => PeerHosts.Contains(i.Host));
|
|
var builder = WebApplication.CreateSlimBuilder();
|
|
builder.WebHost.UseUrls("http://127.0.0.1:0");
|
|
_hinting = builder.Build();
|
|
//answers /{status}/{retry-after}: the peer cannot send a Retry-After
|
|
_hinting.Run(context =>
|
|
{
|
|
var parts = context.Request.Path.Value!.Trim('/').Split('/');
|
|
context.Response.StatusCode = int.Parse(parts[0]);
|
|
if (parts.Length > 1)
|
|
context.Response.Headers.RetryAfter = parts[1] == "date"
|
|
? DateTimeOffset.UtcNow.AddMinutes(10).ToString("r")
|
|
: parts[1];
|
|
return Task.CompletedTask;
|
|
});
|
|
await _hinting.StartAsync();
|
|
_hintingBase = _hinting.Urls.First();
|
|
}
|
|
|
|
public async ValueTask DisposeAsync()
|
|
{
|
|
if (_hinting != default)
|
|
await _hinting.DisposeAsync();
|
|
if (_harness == default)
|
|
return;
|
|
await _harness.DisposeAsync();
|
|
await DB.Default.DeleteAsync<RemoteInstance>(i => PeerHosts.Contains(i.Host));
|
|
}
|
|
|
|
async Task<JobOutcome> Attempt(LocalActor signer, string inbox, int attempts = 1)
|
|
{
|
|
var token = TestContext.Current.CancellationToken;
|
|
var id = $"{Harness.Base}/a/{Guid.NewGuid():N}";
|
|
await _harness.Delivery.Enqueue(signer, new[] { inbox }, new JsonObject
|
|
{
|
|
["id"] = id, ["type"] = "Create", ["actor"] = signer.Uri, ["to"] = new JsonArray(Addressing.Public),
|
|
["object"] = new JsonObject { ["type"] = "Note", ["content"] = "hi" }
|
|
}, token);
|
|
var job = await DB.Default.Find<Job>().Match(j => j.DedupeKey == $"{id}|{inbox}").ExecuteSingleAsync(token);
|
|
job.Attempts = attempts;
|
|
var handler = new DeliveryJobHandler(_harness.Local, Peer.Http(), new HostCircuitBreaker(new MemoryCache(new MemoryCacheOptions())),
|
|
NullLogger<DeliveryJobHandler>.Instance);
|
|
return await handler.Handle(job, token);
|
|
}
|
|
|
|
[Fact]
|
|
public async Task A_delivery_carries_a_valid_digest_and_a_signature_by_the_signer()
|
|
{
|
|
var (_, alice) = await _harness.Persona("alice");
|
|
_harness.Peer.Answer("/ok/inbox", 202);
|
|
|
|
var outcome = await Attempt(alice, _harness.Peer.A + "/ok/inbox");
|
|
|
|
Assert.Equal(JobResult.Done, outcome.Result);
|
|
var received = Assert.Single(_harness.Peer.Requests, r => r.Path == "/ok/inbox");
|
|
Assert.Equal("POST", received.Method);
|
|
Assert.Equal(alice.Uri, JsonNode.Parse(received.Body)!["actor"]!.GetValue<string>());
|
|
Assert.Equal(HttpSignatures.Digest(Encoding.UTF8.GetBytes(received.Body)), received.Headers["Digest"]);
|
|
var signature = HttpSignatures.Parse(received.Signature);
|
|
Assert.Equal(alice.KeyId, signature.KeyId);
|
|
Assert.Equal(new[] { "(request-target)", "host", "date", "digest" }, signature.Headers);
|
|
var signingString = $"(request-target): post /ok/inbox\nhost: {received.Headers["Host"]}\ndate: {received.Headers["Date"]}\ndigest: {received.Headers["Digest"]}";
|
|
using var key = RSA.Create();
|
|
key.ImportFromPem(alice.PublicKeyPem);
|
|
Assert.True(key.VerifyData(Encoding.UTF8.GetBytes(signingString), signature.Signature, HashAlgorithmName.SHA256, RSASignaturePadding.Pkcs1));
|
|
}
|
|
|
|
[Fact]
|
|
public async Task Refusals_are_dead_hints_defer_and_failures_retry()
|
|
{
|
|
var (_, alice) = await _harness.Persona("alice");
|
|
_harness.Peer.Answer("/gone/inbox", 410);
|
|
_harness.Peer.Answer("/missing/inbox", 404);
|
|
_harness.Peer.Answer("/broken/inbox", 500);
|
|
|
|
var gone = await Attempt(alice, _harness.Peer.A + "/gone/inbox");
|
|
var missing = await Attempt(alice, _harness.Peer.A + "/missing/inbox");
|
|
var busy = await Attempt(alice, _hintingBase + "/429/120");
|
|
var down = await Attempt(alice, _hintingBase + "/503/date");
|
|
var unhinted = await Attempt(alice, _hintingBase + "/503");
|
|
var broken = await Attempt(alice, _harness.Peer.A + "/broken/inbox");
|
|
|
|
Assert.Equal(JobResult.Dead, gone.Result);
|
|
Assert.StartsWith("410", gone.Error);
|
|
Assert.Equal(JobResult.Dead, missing.Result);
|
|
Assert.Equal(JobResult.Defer, busy.Result);
|
|
Assert.InRange(busy.RetryAt!.Value, DateTime.UtcNow.AddSeconds(100), DateTime.UtcNow.AddSeconds(125));
|
|
Assert.Equal(JobResult.Defer, down.Result);
|
|
Assert.InRange(down.RetryAt!.Value, DateTime.UtcNow.AddMinutes(9), DateTime.UtcNow.AddMinutes(11));
|
|
Assert.Equal(JobResult.Retry, unhinted.Result);
|
|
Assert.Equal(JobResult.Retry, broken.Result);
|
|
Assert.StartsWith("500", broken.Error);
|
|
}
|
|
|
|
// Mastodon's 422 for two first contacts racing to create one account, and its 409 for a held lock, are retried a
|
|
// few times; the third refusal is final like any other
|
|
[Fact]
|
|
public async Task Conflicts_and_unprocessable_answers_are_retried_twice_then_dead()
|
|
{
|
|
var (_, alice) = await _harness.Persona("alice");
|
|
_harness.Peer.Answer("/racing/inbox", 422);
|
|
_harness.Peer.Answer("/locked/inbox", 409);
|
|
|
|
Assert.Equal(JobResult.Retry, (await Attempt(alice, _harness.Peer.A + "/racing/inbox")).Result);
|
|
Assert.Equal(JobResult.Retry, (await Attempt(alice, _harness.Peer.A + "/locked/inbox", attempts: 2)).Result);
|
|
var final = await Attempt(alice, _harness.Peer.A + "/racing/inbox", attempts: 3);
|
|
Assert.Equal(JobResult.Dead, final.Result);
|
|
Assert.StartsWith("422", final.Error);
|
|
}
|
|
}
|
|
}
|