From 9e623a8f8bbca4bb3cefb23ccb33462c584717d6 Mon Sep 17 00:00:00 2001 From: Hamza Alqurneh Date: Sun, 27 Sep 2026 11:19:53 +0300 Subject: [PATCH] fix: run delayed retries, manual creates and aggregation in the subscription's work group --- SW.Bitween.Api/Resources/Xchanges/Create.cs | 2 +- SW.Bitween.Api/Services/AggregationJob.cs | 2 +- SW.Bitween.Api/Services/XchangeService.cs | 2 +- .../Tests/WorkGroupRoutingTests.cs | 176 ++++++++++++++++++ 4 files changed, 179 insertions(+), 3 deletions(-) create mode 100644 SW.Bitween.IntegrationTests/Tests/WorkGroupRoutingTests.cs diff --git a/SW.Bitween.Api/Resources/Xchanges/Create.cs b/SW.Bitween.Api/Resources/Xchanges/Create.cs index 23d5ad82..b50ba239 100644 --- a/SW.Bitween.Api/Resources/Xchanges/Create.cs +++ b/SW.Bitween.Api/Resources/Xchanges/Create.cs @@ -21,7 +21,7 @@ public async Task Handle(CreateXchange request) } else if (request.Option == CreateXchangeOption.SubscriberId) { - var subscription = await dbc.Set().FirstOrDefaultAsync(d => d.Id == request.SubscriberId); + var subscription = await dbc.Subscriptions().FirstOrDefaultAsync(d => d.Id == request.SubscriberId); if (subscription == null) throw new SWValidationException("SUBSCRIPTION_NOT_FOUND", "Subscription was not found"); await xchangeService.CreateXchange(subscription, xchangeFile); } diff --git a/SW.Bitween.Api/Services/AggregationJob.cs b/SW.Bitween.Api/Services/AggregationJob.cs index fceb6761..3c73d8f8 100644 --- a/SW.Bitween.Api/Services/AggregationJob.cs +++ b/SW.Bitween.Api/Services/AggregationJob.cs @@ -21,7 +21,7 @@ public class AggregationJob( { public async Task Execute(AggregationJobParams jobParams) { - var aggSub = await dbContext.Set() + var aggSub = await dbContext.Subscriptions() .FirstOrDefaultAsync(s => s.Id == jobParams.SubscriptionId && !s.Inactive); if (aggSub == null) return; diff --git a/SW.Bitween.Api/Services/XchangeService.cs b/SW.Bitween.Api/Services/XchangeService.cs index 9194cef0..8ecf2abd 100644 --- a/SW.Bitween.Api/Services/XchangeService.cs +++ b/SW.Bitween.Api/Services/XchangeService.cs @@ -140,7 +140,7 @@ public async Task ExecuteDelayedRetry(DelayedRetry delayedRetry) return false; } - var subscription = await dbContext.Set() + var subscription = await dbContext.Subscriptions() .FirstOrDefaultAsync(s => s.Id == xchange.SubscriptionId); if (subscription == null) { diff --git a/SW.Bitween.IntegrationTests/Tests/WorkGroupRoutingTests.cs b/SW.Bitween.IntegrationTests/Tests/WorkGroupRoutingTests.cs new file mode 100644 index 00000000..b073a645 --- /dev/null +++ b/SW.Bitween.IntegrationTests/Tests/WorkGroupRoutingTests.cs @@ -0,0 +1,176 @@ +using System; +using System.Collections.Generic; +using System.Linq; +using System.Threading.Tasks; +using Microsoft.EntityFrameworkCore; +using Microsoft.Extensions.DependencyInjection; +using Microsoft.Extensions.Logging.Abstractions; +using SW.Bitween.Domain; +using SW.Bitween.IntegrationTests.Fixtures; +using SW.Bitween.Model; +using SW.Bitween.Resources.Xchanges; +using SW.PrimitiveTypes; +using Xunit; + +namespace SW.Bitween.IntegrationTests.Tests; + +/// +/// An exchange is sent to its subscription's work group queue. Each path is run on a fresh +/// DbContext, the way it runs in production: in the context that seeded the data the work group is +/// already tracked, EF fills the subscription's navigation property on its own, and a loader that +/// forgets to include it looks correct. +/// +[Collection("Bitween")] +public class WorkGroupRoutingTests(BitweenFixture fixture) +{ + // ─── Helpers ────────────────────────────────────────────────────────────── + + /// + /// Records the queue each new exchange for is published to. + /// Read while saving, because the save publishes the events and then clears them. + /// + private static List RecordQueues(BitweenDbContext db, int subscriptionId) + { + var queues = new List(); + db.SavingChanges += (_, _) => queues.AddRange(db.ChangeTracker.Entries() + .Where(e => e.State == EntityState.Added && e.Entity.SubscriptionId == subscriptionId) + .SelectMany(e => e.Entity.Events.OfType()) + .Select(ev => ev.GetBusMessageName())); + return queues; + } + + private static async Task AddWorkGroup(BitweenDbContext db) + { + var workGroup = new WorkGroup + { + Name = "Routing " + Guid.NewGuid().ToString("N")[..8], + BusMessageName = "routing" + Guid.NewGuid().ToString("N")[..8] + }; + db.Add(workGroup); + await db.SaveChangesAsync(); + return workGroup; + } + + private static async Task AddDocument(BitweenDbContext db, string name) + { + var doc = new Document(null, name, DocumentFormat.Json); + db.Set().Add(doc); + await db.SaveChangesAsync(); + return doc; + } + + // ─── Paths that load the subscription themselves ────────────────────────── + + [Fact] + public async Task Delayed_retry_runs_in_the_subscriptions_work_group() + { + int subscriptionId; + string expected; + await using (var seed = fixture.CreateScope()) + { + var db = seed.ServiceProvider.GetRequiredService(); + var xs = seed.ServiceProvider.GetRequiredService(); + + var workGroup = await AddWorkGroup(db); + var doc = await AddDocument(db, "WG Routing Retry Doc"); + var sub = new Subscription("WG Routing Retry Sub", doc.Id) { Inactive = false, WorkGroupId = workGroup.Id }; + db.Set().Add(sub); + await db.SaveChangesAsync(); + + var original = await xs.CreateXchange(sub, new XchangeFile("{}")); + await db.SaveChangesAsync(); + db.Set().Add(new XchangeResult(original.Id, null, null, exception: "boom")); + db.Set().Add(new DelayedRetry { Id = original.Id, On = DateTime.UtcNow.AddMinutes(-1) }); + await db.SaveChangesAsync(); + + subscriptionId = sub.Id; + expected = workGroup.GetBusMessageName(); + } + fixture.App.Services.GetRequiredService().Revoke(); + + await using var scope = fixture.CreateScope(); + var runDb = scope.ServiceProvider.GetRequiredService(); + var queues = RecordQueues(runDb, subscriptionId); + + await new RetryJob(runDb, scope.ServiceProvider.GetRequiredService(), + NullLogger.Instance).Execute(); + + Assert.Equal(expected, Assert.Single(queues)); + } + + [Fact] + public async Task Manual_create_for_a_subscription_runs_in_its_work_group() + { + int subscriptionId; + string expected; + await using (var seed = fixture.CreateScope()) + { + var db = seed.ServiceProvider.GetRequiredService(); + var workGroup = await AddWorkGroup(db); + var doc = await AddDocument(db, "WG Routing Create Doc"); + var sub = new Subscription("WG Routing Create Sub", doc.Id) { Inactive = false, WorkGroupId = workGroup.Id }; + db.Set().Add(sub); + await db.SaveChangesAsync(); + + subscriptionId = sub.Id; + expected = workGroup.GetBusMessageName(); + } + fixture.App.Services.GetRequiredService().Revoke(); + + await using var scope = fixture.CreateScope(); + var runDb = scope.ServiceProvider.GetRequiredService(); + var queues = RecordQueues(runDb, subscriptionId); + + await new Create(scope.ServiceProvider.GetRequiredService(), runDb).Handle(new CreateXchange + { + Option = CreateXchangeOption.SubscriberId, + SubscriberId = subscriptionId, + Data = "{}" + }); + + Assert.Equal(expected, Assert.Single(queues)); + } + + [Fact] + public async Task Aggregation_runs_in_the_aggregation_subscriptions_work_group() + { + int aggSubId; + string expected; + await using (var seed = fixture.CreateScope()) + { + var db = seed.ServiceProvider.GetRequiredService(); + var workGroup = await AddWorkGroup(db); + var sourceDoc = await AddDocument(db, "WG Routing Agg Source Doc"); + var sourceSub = new Subscription("WG Routing Agg Source", sourceDoc.Id) { Inactive = false }; + db.Set().Add(sourceSub); + await db.SaveChangesAsync(); + + var source = new Xchange(sourceSub, new XchangeFile("{\"n\":1}")); + db.Set().Add(source); + await db.SaveChangesAsync(); + db.Set().Add(new XchangeResult(source.Id, null, null)); + await db.SaveChangesAsync(); + + var aggSub = new Subscription("WG Routing Agg", sourceSub.Id, Partner.SystemId) + { + Inactive = false, + AggregationTarget = XchangeFileType.Input, + WorkGroupId = workGroup.Id + }; + db.Set().Add(aggSub); + await db.SaveChangesAsync(); + + aggSubId = aggSub.Id; + expected = workGroup.GetBusMessageName(); + } + fixture.App.Services.GetRequiredService().Revoke(); + + await using var scope = fixture.CreateScope(); + var runDb = scope.ServiceProvider.GetRequiredService(); + var queues = RecordQueues(runDb, aggSubId); + + await scope.ServiceProvider.GetRequiredService().Execute(new AggregationJobParams(aggSubId, null)); + + Assert.Equal(expected, Assert.Single(queues)); + } +}