diff --git a/CLAUDE.md b/CLAUDE.md index b01062a..cabe2bf 100644 --- a/CLAUDE.md +++ b/CLAUDE.md @@ -316,7 +316,7 @@ cd /var/www/privapub.thepra.dev && sudo -u www-data ASPNETCORE_ENVIRONMENT=Produ stored as an edit date on every remote post. Write `(T?)null`. It has shipped four times (`MastodonParams.Bool/Int`, `NoteParser.Int`, `NoteParser.Time`). - ActivityPub output is built with `System.Text.Json.Nodes` in `ActivityPubRenderer`, not typed models. Inbound - documents are read through `InboxService.Id`/`Value`, which handle string, object and array. + documents are read through `ActivityJson.Id`/`Value`, which handle string, object and array. - Libraries chosen for the roadmap: HtmlSanitizer, Markdig (`DisableHtml`), NSign (RFC 9421 inbound), OpenIddict + OpenIddict.MongoDb, NetVips, Blurhash.Core, FFMpegCore. No ImageSharp (licence key enforced), no MassTransit. @@ -327,7 +327,16 @@ without `PRIVAPUB_TEST_MONGOD=1`. CI runs the unit tests only (the box's mongods - `Support/Peer` is an in-process HTTP server answering on two origins (`127.0.0.1` and `localhost`), so origin rules can be tested; `Support/RemoteActor` signs real deliveries with its own key. -- Inbox scenarios go through `InboxService.Receive` with a signed request, not through the private handlers. +- Inbox scenarios go through `InboxReceiver.Receive` with a signed request, not through the private handlers. +- All test classes share one database and run in parallel, so a test touches only rows it made: + - random names and GUIDs; + - `Harness.Outgoing` sees only deliveries queued since that harness started, because Peer ports are reused; + - a worker gets a scoped `new JobQueue(j => ...)` so it never leases another test's jobs; + - domain blocks are set with `DomainBlocks.Load`, never written to the database. + + A test that must change something database-wide (drop indexes, run a migration over every post, let deliveries to + `localhost` fail and trip its breaker) goes in `[Xunit.Collection(nameof(Exclusive))]`, which runs alone, and + cleans up after itself. Pure logic belongs in an unconditional unit test, not in a Mongo-gated class. Beyond the tests, verify by building, running locally, and exercising: - the client API (sign up, create an avatar, a group, a post); @@ -337,7 +346,7 @@ Interop is checked against real servers, starting with the workstation's own pas ```bash DOTNET=~/.dotnet/dotnet tools/pasture/run.sh up # podman: PrivaPub + the latest GoToSocial + Mongo, behind Caddy -tools/pasture/interop.sh # 25 checks, each side driven through its own Mastodon API +tools/pasture/interop.sh # 33 checks, each side driven through its own Mastodon API tools/pasture/run.sh down ``` diff --git a/PrivaPub.Tests/Federation/DomainBlocksTests.cs b/PrivaPub.Tests/Federation/DomainBlocksTests.cs new file mode 100644 index 0000000..87701d7 --- /dev/null +++ b/PrivaPub.Tests/Federation/DomainBlocksTests.cs @@ -0,0 +1,37 @@ +using Microsoft.Extensions.Logging.Abstractions; + +using PrivaPub.Federation.Moderation; +using PrivaPub.Models.Federation; + +namespace PrivaPub.Tests.Federation +{ + public class DomainBlocksTests + { + [Fact] + public void A_block_covers_subdomains_but_not_lookalikes() + { + var blocks = new DomainBlocks(NullLogger.Instance); + blocks.Load(new[] { new DomainBlock { Domain = "evil.example", Severity = DomainBlockSeverity.Silence } }); + + Assert.NotNull(blocks.Find("evil.example")); + Assert.NotNull(blocks.Find("A.Evil.Example.")); + Assert.Null(blocks.Find("notevil.example")); + Assert.False(blocks.IsSuspended("evil.example")); + } + + [Fact] + public void A_suspension_is_found_from_any_subdomain() + { + var blocks = new DomainBlocks(NullLogger.Instance); + blocks.Load(new[] + { + new DomainBlock { Domain = " Spam.Example. ", Severity = DomainBlockSeverity.Suspend }, + new DomainBlock { Domain = "", Severity = DomainBlockSeverity.Suspend } + }); + + Assert.True(blocks.IsSuspended("deep.inside.spam.example")); + Assert.False(blocks.IsSuspended("example")); + Assert.Null(blocks.Find(default)); + } + } +} diff --git a/PrivaPub.Tests/Federation/InboxScenarioTests.cs b/PrivaPub.Tests/Federation/InboxScenarioTests.cs index 10227a8..434d6a6 100644 --- a/PrivaPub.Tests/Federation/InboxScenarioTests.cs +++ b/PrivaPub.Tests/Federation/InboxScenarioTests.cs @@ -269,38 +269,14 @@ namespace PrivaPub.Tests.Federation var token = TestContext.Current.CancellationToken; var alice = await LocalAvatar("alice"); var bob = new RemoteActor(_peer, "bob", _peer.B); - await DB.Default.SaveAsync(new DomainBlock { Domain = "localhost", Severity = DomainBlockSeverity.Suspend }, token); - try - { - await _blocks.Reload(token); - var before = _peer.Requests.Count; + _blocks.Load(new[] { new DomainBlock { Domain = "localhost", Severity = DomainBlockSeverity.Suspend } }); + var before = _peer.Requests.Count; - var result = await Deliver(bob, $"/peasants/{alice.UserName}/mouth", DirectCreate(bob, alice.Uri)); + var result = await Deliver(bob, $"/peasants/{alice.UserName}/mouth", DirectCreate(bob, alice.Uri)); - Assert.Equal(202, result.StatusCode); - Assert.Equal(before, _peer.Requests.Count); - Assert.False(await DB.Default.Find().Match(p => p.ActorURI == bob.Id).ExecuteAnyAsync(token)); - } - finally - { - await DB.Default.DeleteAsync(b => b.Domain == "localhost"); - await _blocks.Reload(token); - } - } - - [Fact] - public void A_block_covers_subdomains_but_not_lookalikes() - { - var blocks = new DomainBlocks(NullLogger.Instance); - typeof(DomainBlocks).GetField("_blocks", System.Reflection.BindingFlags.NonPublic | System.Reflection.BindingFlags.Instance)! - .SetValue(blocks, new Dictionary { ["evil.example"] = new() { Domain = "evil.example", Severity = DomainBlockSeverity.Silence } }); - typeof(DomainBlocks).GetField("_loadedAt", System.Reflection.BindingFlags.NonPublic | System.Reflection.BindingFlags.Instance)! - .SetValue(blocks, DateTime.UtcNow); - - Assert.NotNull(blocks.Find("evil.example")); - Assert.NotNull(blocks.Find("A.Evil.Example.")); - Assert.Null(blocks.Find("notevil.example")); - Assert.False(blocks.IsSuspended("evil.example")); + Assert.Equal(202, result.StatusCode); + Assert.Equal(before, _peer.Requests.Count); + Assert.False(await DB.Default.Find().Match(p => p.ActorURI == bob.Id).ExecuteAnyAsync(token)); } [Fact] diff --git a/PrivaPub.Tests/Federation/ObjectRecordsTests.cs b/PrivaPub.Tests/Federation/ObjectRecordsTests.cs new file mode 100644 index 0000000..64fca45 --- /dev/null +++ b/PrivaPub.Tests/Federation/ObjectRecordsTests.cs @@ -0,0 +1,36 @@ +using PrivaPub.Federation.Objects; + +using System.Text.Json.Nodes; + +namespace PrivaPub.Tests.Federation +{ + public class ObjectRecordsTests + { + [Fact] + public void An_object_larger_than_the_cap_keeps_only_its_hash() + { + var big = new JsonObject { ["content"] = new string('x', ObjectRecords.MaxRawBytes + 1) }; + + var (text, hash, bytes, truncated) = ObjectRecords.Capture(big); + + Assert.Null(text); + Assert.True(truncated); + Assert.True(bytes > ObjectRecords.MaxRawBytes); + Assert.Equal(64, hash.Length); + } + + [Fact] + public void A_small_object_is_kept_whole_with_its_hash() + { + var small = new JsonObject { ["type"] = "Note", ["content"] = "hi" }; + + var (text, hash, bytes, truncated) = ObjectRecords.Capture(small); + + Assert.Equal(small.ToJsonString(), text); + Assert.False(truncated); + Assert.Equal(text.Length, bytes); + Assert.Equal(64, hash.Length); + Assert.Equal(default, ObjectRecords.Capture(default)); + } + } +} diff --git a/PrivaPub.Tests/Federation/ProvenanceTests.cs b/PrivaPub.Tests/Federation/ProvenanceTests.cs index d695a42..031ffef 100644 --- a/PrivaPub.Tests/Federation/ProvenanceTests.cs +++ b/PrivaPub.Tests/Federation/ProvenanceTests.cs @@ -68,18 +68,5 @@ namespace PrivaPub.Tests.Federation Assert.Equal("

hi again

", JsonNode.Parse(Assert.Single(record.Revisions).Raw)!["content"]!.GetValue()); Assert.True(await DB.Default.Find().Match(j => j.Kind == JobKind.DescribeInstance && j.Payload == record.Host).ExecuteAnyAsync(token)); } - - [Fact] - public void An_object_larger_than_the_cap_keeps_only_its_hash() - { - var big = new JsonObject { ["content"] = new string('x', ObjectRecords.MaxRawBytes + 1) }; - - var (text, hash, bytes, truncated) = ObjectRecords.Capture(big); - - Assert.Null(text); - Assert.True(truncated); - Assert.True(bytes > ObjectRecords.MaxRawBytes); - Assert.Equal(64, hash.Length); - } } } diff --git a/PrivaPub.Tests/Infrastructure/DeliveryWorkerTests.cs b/PrivaPub.Tests/Infrastructure/DeliveryWorkerTests.cs new file mode 100644 index 0000000..9d6b575 --- /dev/null +++ b/PrivaPub.Tests/Infrastructure/DeliveryWorkerTests.cs @@ -0,0 +1,95 @@ +using Microsoft.Extensions.Caching.Memory; +using Microsoft.Extensions.Logging.Abstractions; + +using MongoDB.Entities; + +using PrivaPub.Federation.Actors; +using PrivaPub.Federation.Outbox; +using PrivaPub.Infrastructure.Jobs; +using PrivaPub.Models; +using PrivaPub.Models.Jobs; +using PrivaPub.Models.User; +using PrivaPub.StaticServices; +using PrivaPub.Tests.Support; + +using System.Text.Json.Nodes; + +namespace PrivaPub.Tests.Infrastructure +{ + [Xunit.Collection(nameof(Exclusive))] + [Trait("Category", "Integration")] + public sealed class DeliveryWorkerTests : IAsyncLifetime + { + static readonly string[] PeerHosts = { "localhost", "127.0.0.1" }; + + Peer _peer; + + public async ValueTask InitializeAsync() + { + Assert.SkipUnless(MongoFixture.Enabled, MongoFixture.Skip); + _peer = await Peer.Start(); + await DB.Default.DeleteAsync(i => PeerHosts.Contains(i.Host)); + } + + public async ValueTask DisposeAsync() + { + if (_peer == default) + return; + await _peer.DisposeAsync(); + await DB.Default.DeleteAsync(i => PeerHosts.Contains(i.Host)); + } + + [Fact] + public async Task A_dead_host_does_not_hold_up_deliveries_to_live_ones() + { + var token = TestContext.Current.CancellationToken; + var (privateKey, publicKey) = Keys.NewKeyPair(); + var avatar = new Avatar { UserName = $"sender{Guid.NewGuid():N}"[..20], PrivateKey = privateKey, PublicKey = publicKey }; + await DB.Default.SaveAsync(avatar, token); + var local = new LocalActorService(new DbEntities(), new StaticOptions(new AppConfiguration { BackendBaseAddress = "https://privapub.test" })); + var sender = local.FromAvatar(avatar); + var run = $"https://privapub.test/a/{Guid.NewGuid():N}/"; + var queue = new JobQueue(j => j.DedupeKey.StartsWith(run)); + var delivery = new DeliveryService(new DbEntities(), queue); + _peer.Answer("/dead/inbox", 503, TimeSpan.FromMilliseconds(300)); + _peer.Answer("/live/inbox", 202); + + var deadInboxes = Enumerable.Range(0, 60).Select(i => $"{_peer.B}/dead/inbox?{i}"); + await delivery.Enqueue(sender, deadInboxes, new JsonObject { ["id"] = run + "dead", ["type"] = "Create" }, token); + await delivery.Enqueue(sender, new[] { $"{_peer.A}/live/inbox" }, new JsonObject { ["id"] = run + "live", ["type"] = "Create" }, token); + + var handler = new DeliveryJobHandler(local, Peer.Http(), new HostCircuitBreaker(new MemoryCache(new MemoryCacheOptions())), NullLogger.Instance); + using var worker = new JobWorker(queue, new IJobHandler[] { handler }, NullLogger.Instance); + await worker.StartAsync(token); + try + { + await Until(() => Task.FromResult(_peer.Requests.Any(r => r.Path == "/live/inbox")), token); + var deadSoFar = _peer.Requests.Count(r => r.Path == "/dead/inbox"); + var live = Assert.Single(_peer.Requests, r => r.Path == "/live/inbox"); + Assert.Contains($"keyId=\"{sender.KeyId}\"", live.Signature); + Assert.True(deadSoFar < 30, $"{deadSoFar} dead deliveries were tried before the live one"); + + await Until(() => DB.Default.Find().Match(j => j.DedupeKey.StartsWith(run + "dead") && j.State == JobState.Pending && j.Attempts > 0).ExecuteAnyAsync(token), token); + } + finally + { + await worker.StopAsync(CancellationToken.None); + } + + var retried = await DB.Default.Find().Match(j => j.DedupeKey.StartsWith(run + "dead") && j.Attempts > 0 && j.State == JobState.Pending).ExecuteFirstAsync(token); + Assert.Equal("503", retried.LastError?.Split(' ')[0]); + Assert.True(retried.RunAt > DateTime.UtcNow); + Assert.True((await DB.Default.Find().Match(i => i.Host == "localhost").ExecuteFirstAsync(token))?.ConsecutiveFailures > 0); + } + + static async Task Until(Func> condition, CancellationToken token) + { + var deadline = DateTime.UtcNow.AddSeconds(30); + while (!await condition()) + { + Assert.True(DateTime.UtcNow < deadline, "timed out"); + await Task.Delay(50, token); + } + } + } +} diff --git a/PrivaPub.Tests/Infrastructure/HostCircuitBreakerTests.cs b/PrivaPub.Tests/Infrastructure/HostCircuitBreakerTests.cs new file mode 100644 index 0000000..a483bec --- /dev/null +++ b/PrivaPub.Tests/Infrastructure/HostCircuitBreakerTests.cs @@ -0,0 +1,67 @@ +using Microsoft.Extensions.Caching.Memory; + +using MongoDB.Entities; + +using PrivaPub.Infrastructure.Jobs; +using PrivaPub.Models.Jobs; +using PrivaPub.Tests.Support; + +namespace PrivaPub.Tests.Infrastructure +{ + [Trait("Category", "Integration")] + public sealed class HostCircuitBreakerTests : IAsyncLifetime + { + readonly string _host = $"breaker{Guid.NewGuid():N}.example"; + + public ValueTask InitializeAsync() + { + Assert.SkipUnless(MongoFixture.Enabled, MongoFixture.Skip); + return ValueTask.CompletedTask; + } + + public async ValueTask DisposeAsync() + { + if (MongoFixture.Enabled) + await DB.Default.DeleteAsync(i => i.Host == _host); + } + + [Fact] + public async Task A_host_is_quarantined_only_past_the_threshold_and_freed_by_a_success() + { + var token = TestContext.Current.CancellationToken; + var breaker = new HostCircuitBreaker(new MemoryCache(new MemoryCacheOptions())); + + for (var i = 1; i < HostCircuitBreaker.Threshold; i++) + await breaker.Failed(_host, "503", token); + Assert.Null(await breaker.UnavailableUntil(_host, token)); + + await breaker.Failed(_host, "503", token); + var until = await breaker.UnavailableUntil(_host, token); + Assert.NotNull(until); + Assert.InRange(until.Value, DateTime.UtcNow.AddMinutes(55), DateTime.UtcNow.AddMinutes(65)); + var instance = await DB.Default.Find().Match(i => i.Host == _host).ExecuteFirstAsync(token); + Assert.Equal(HostCircuitBreaker.Threshold, instance.ConsecutiveFailures); + Assert.Equal("503", instance.LastError); + + await breaker.Succeeded(_host, token); + Assert.Null(await breaker.UnavailableUntil(_host, token)); + instance = await DB.Default.Find().Match(i => i.Host == _host).ExecuteFirstAsync(token); + Assert.Equal(0, instance.ConsecutiveFailures); + Assert.NotNull(instance.LastSuccessAt); + } + + [Fact] + public async Task A_healthy_host_is_written_at_most_once_an_hour() + { + var token = TestContext.Current.CancellationToken; + var breaker = new HostCircuitBreaker(new MemoryCache(new MemoryCacheOptions())); + + await breaker.Succeeded(_host, token); + var first = (await DB.Default.Find().Match(i => i.Host == _host).ExecuteFirstAsync(token)).LastSuccessAt; + await breaker.Succeeded(_host, token); + var second = (await DB.Default.Find().Match(i => i.Host == _host).ExecuteFirstAsync(token)).LastSuccessAt; + + Assert.Equal(first, second); + } + } +} diff --git a/PrivaPub.Tests/Infrastructure/IndexTests.cs b/PrivaPub.Tests/Infrastructure/IndexTests.cs index 4ed1741..685204d 100644 --- a/PrivaPub.Tests/Infrastructure/IndexTests.cs +++ b/PrivaPub.Tests/Infrastructure/IndexTests.cs @@ -12,6 +12,7 @@ using PrivaPub.Tests.Support; namespace PrivaPub.Tests.Infrastructure { + [Xunit.Collection(nameof(Exclusive))] [Trait("Category", "Integration")] public class IndexTests { diff --git a/PrivaPub.Tests/Infrastructure/JobQueueTests.cs b/PrivaPub.Tests/Infrastructure/JobQueueTests.cs index 4a198bf..6552cdd 100644 --- a/PrivaPub.Tests/Infrastructure/JobQueueTests.cs +++ b/PrivaPub.Tests/Infrastructure/JobQueueTests.cs @@ -1,19 +1,9 @@ -using Microsoft.Extensions.Caching.Memory; -using Microsoft.Extensions.Logging.Abstractions; - using MongoDB.Entities; -using PrivaPub.Federation.Actors; -using PrivaPub.Federation.Outbox; using PrivaPub.Infrastructure.Jobs; -using PrivaPub.Models; using PrivaPub.Models.Jobs; -using PrivaPub.Models.User; -using PrivaPub.StaticServices; using PrivaPub.Tests.Support; -using System.Text.Json.Nodes; - namespace PrivaPub.Tests.Infrastructure { public class BackoffTests @@ -38,19 +28,13 @@ namespace PrivaPub.Tests.Infrastructure [Trait("Category", "Integration")] public sealed class JobQueueTests : IAsyncLifetime { - Peer _peer; - - public async ValueTask InitializeAsync() + public ValueTask InitializeAsync() { Assert.SkipUnless(MongoFixture.Enabled, MongoFixture.Skip); - _peer = await Peer.Start(); + return ValueTask.CompletedTask; } - public async ValueTask DisposeAsync() - { - if (_peer != default) - await _peer.DisposeAsync(); - } + public ValueTask DisposeAsync() => ValueTask.CompletedTask; static Job NewJob(string host = default, string dedupe = default) => new() { @@ -102,44 +86,9 @@ namespace PrivaPub.Tests.Infrastructure var job = new Job { Kind = (JobKind)98, State = JobState.Running, LeasedUntil = DateTime.UtcNow.AddMinutes(-1) }; await DB.Default.SaveAsync(job, token); - await new JobQueue().Reap(token); + await new JobQueue(j => j.ID == job.ID).Reap(token); Assert.Equal(JobState.Pending, (await DB.Default.Find().OneAsync(job.ID, token)).State); } - - [Fact] - public async Task A_dead_host_does_not_hold_up_deliveries_to_live_ones() - { - var token = TestContext.Current.CancellationToken; - var (privateKey, publicKey) = Keys.NewKeyPair(); - var avatar = new Avatar { UserName = $"sender{Guid.NewGuid():N}"[..20], PrivateKey = privateKey, PublicKey = publicKey }; - await DB.Default.SaveAsync(avatar, token); - var local = new LocalActorService(new DbEntities(), new StaticOptions(new AppConfiguration { BackendBaseAddress = "https://privapub.test" })); - var sender = local.FromAvatar(avatar); - var queue = new JobQueue(); - var delivery = new DeliveryService(new DbEntities(), queue); - _peer.Answer("/dead/inbox", 503, TimeSpan.FromMilliseconds(300)); - _peer.Answer("/live/inbox", 202); - - var deadInboxes = Enumerable.Range(0, 60).Select(i => $"{_peer.B}/dead/inbox?{i}"); - await delivery.Enqueue(sender, deadInboxes, new JsonObject { ["id"] = $"https://privapub.test/a/{Guid.NewGuid():N}", ["type"] = "Create" }, token); - await delivery.Enqueue(sender, new[] { $"{_peer.A}/live/inbox" }, new JsonObject { ["id"] = $"https://privapub.test/a/{Guid.NewGuid():N}", ["type"] = "Create" }, token); - - var handler = new DeliveryJobHandler(local, Peer.Http(), new HostCircuitBreaker(new MemoryCache(new MemoryCacheOptions())), NullLogger.Instance); - using var worker = new JobWorker(queue, new IJobHandler[] { handler }, NullLogger.Instance); - await worker.StartAsync(token); - var deadline = DateTime.UtcNow.AddSeconds(5); - while (DateTime.UtcNow < deadline && !_peer.Requests.Any(r => r.Path == "/live/inbox")) - await Task.Delay(50, token); - var live = Assert.Single(_peer.Requests, r => r.Path == "/live/inbox"); - Assert.True(_peer.Requests.Count(r => r.Path == "/dead/inbox") < 30); - while (DateTime.UtcNow < deadline && !await DB.Default.Find().Match(i => i.Host == "localhost").ExecuteAnyAsync(token)) - await Task.Delay(50, token); - await worker.StopAsync(token); - - Assert.Contains($"keyId=\"{sender.KeyId}\"", live.Signature); - var instance = await DB.Default.Find().Match(i => i.Host == "localhost").ExecuteFirstAsync(token); - Assert.True(instance.ConsecutiveFailures > 0); - } } } diff --git a/PrivaPub.Tests/Infrastructure/MigrationTests.cs b/PrivaPub.Tests/Infrastructure/MigrationTests.cs index f38b05a..611a86c 100644 --- a/PrivaPub.Tests/Infrastructure/MigrationTests.cs +++ b/PrivaPub.Tests/Infrastructure/MigrationTests.cs @@ -6,6 +6,7 @@ using PrivaPub.Tests.Support; namespace PrivaPub.Tests.Infrastructure { + [Xunit.Collection(nameof(Exclusive))] [Trait("Category", "Integration")] public class MigrationTests { diff --git a/PrivaPub.Tests/Support/Exclusive.cs b/PrivaPub.Tests/Support/Exclusive.cs new file mode 100644 index 0000000..6881df9 --- /dev/null +++ b/PrivaPub.Tests/Support/Exclusive.cs @@ -0,0 +1,7 @@ +namespace PrivaPub.Tests.Support +{ + [CollectionDefinition(nameof(Exclusive), DisableParallelization = true)] + public sealed class Exclusive + { + } +} diff --git a/PrivaPub.Tests/Support/Harness.cs b/PrivaPub.Tests/Support/Harness.cs index 0600f69..021df01 100644 --- a/PrivaPub.Tests/Support/Harness.cs +++ b/PrivaPub.Tests/Support/Harness.cs @@ -34,6 +34,8 @@ namespace PrivaPub.Tests.Support public const string Host = "privapub.test"; public const string Base = "https://" + Host; + readonly DateTime _started = Millisecond(DateTime.UtcNow); + Harness(Peer peer) { Peer = peer; @@ -136,13 +138,15 @@ namespace PrivaPub.Tests.Support } public async Task> Outgoing(string inbox) => - (await DB.Default.Find().Match(j => j.Kind == JobKind.Deliver).ExecuteAsync()) + (await DB.Default.Find().Match(j => j.Kind == JobKind.Deliver && j.CreatedAt >= _started).ExecuteAsync()) .Select(j => JsonSerializer.Deserialize(j.Payload)) .Where(p => p.Inbox == inbox) .Select(p => JsonNode.Parse(p.Body)!.AsObject()) .ToList(); public async ValueTask DisposeAsync() => await Peer.DisposeAsync(); + + static DateTime Millisecond(DateTime time) => new(time.Ticks - time.Ticks % TimeSpan.TicksPerMillisecond, DateTimeKind.Utc); } public sealed class KeyLocalizer : IStringLocalizer diff --git a/PrivaPub/Federation/Moderation/DomainBlocks.cs b/PrivaPub/Federation/Moderation/DomainBlocks.cs index 3275aec..9c5276f 100644 --- a/PrivaPub/Federation/Moderation/DomainBlocks.cs +++ b/PrivaPub/Federation/Moderation/DomainBlocks.cs @@ -44,9 +44,11 @@ namespace PrivaPub.Federation.Moderation public bool IsSuspended(string host) => Find(host)?.Severity == DomainBlockSeverity.Suspend; - public async Task Reload(CancellationToken token) + public async Task Reload(CancellationToken token) => + Load(await DB.Default.Find().ExecuteAsync(token)); + + public void Load(IEnumerable blocks) { - var blocks = await DB.Default.Find().ExecuteAsync(token); _blocks = blocks.Where(b => !string.IsNullOrEmpty(b.Domain)) .GroupBy(b => Normalise(b.Domain)) .ToDictionary(g => g.Key, g => g.First()); diff --git a/PrivaPub/Infrastructure/Jobs/JobQueue.cs b/PrivaPub/Infrastructure/Jobs/JobQueue.cs index 9d2f655..c3b13ae 100644 --- a/PrivaPub/Infrastructure/Jobs/JobQueue.cs +++ b/PrivaPub/Infrastructure/Jobs/JobQueue.cs @@ -4,6 +4,7 @@ using MongoDB.Entities; using PrivaPub.Models.Jobs; using System.Collections.Concurrent; +using System.Linq.Expressions; namespace PrivaPub.Infrastructure.Jobs { @@ -39,6 +40,13 @@ namespace PrivaPub.Infrastructure.Jobs readonly string _owner = $"{Environment.MachineName}:{Environment.ProcessId}"; readonly ConcurrentDictionary _signals = new(); + readonly FilterDefinition _scope = Builders.Filter.Empty; + + public JobQueue() + { + } + + public JobQueue(Expression> scope) => _scope = Builders.Filter.Where(scope); 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; @@ -66,7 +74,7 @@ namespace PrivaPub.Infrastructure.Jobs var now = DateTime.UtcNow; var busy = busyHosts?.ToList() ?? new List(); return await DB.Default.UpdateAndGet() - .Match(j => j.Kind == kind && j.State == JobState.Pending && j.RunAt <= now && !busy.Contains(j.Host)) + .Match(f => f.Where(j => j.Kind == kind && j.State == JobState.Pending && j.RunAt <= now && !busy.Contains(j.Host)) & _scope) .Modify(j => j.State, JobState.Running) .Modify(j => j.LeasedUntil, now + LeaseTime) .Modify(j => j.LeaseOwner, _owner) @@ -106,7 +114,7 @@ namespace PrivaPub.Infrastructure.Jobs { var now = DateTime.UtcNow; var result = await DB.Default.Update() - .Match(j => j.State == JobState.Running && j.LeasedUntil < now) + .Match(f => f.Where(j => j.State == JobState.Running && j.LeasedUntil < now) & _scope) .Modify(j => j.State, JobState.Pending) .Modify(j => j.RunAt, now) .Modify(j => j.LeasedUntil, null)