-
Notifications
You must be signed in to change notification settings - Fork 0
Expand file tree
/
Copy pathResubmitExample.cs
More file actions
70 lines (58 loc) · 3.16 KB
/
Copy pathResubmitExample.cs
File metadata and controls
70 lines (58 loc) · 3.16 KB
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
using Codefinity.EventEmitter.Sample.Infrastructure;
using Codefinity.EventEmitter.Sample.Inventory;
using Codefinity.EventEmitter.Sample.Orders;
using Codefinity.EventEmitter.Sample.Shipping;
using Microsoft.Extensions.DependencyInjection;
namespace Codefinity.EventEmitter.Sample.Examples;
/// <summary>
/// A module listener fails: the publication stays incomplete with its failure recorded, other listeners are
/// unaffected, and the publication can be inspected and resubmitted with IIncompleteEventPublications.
/// </summary>
internal static class ResubmitExample
{
public static async Task RunAsync()
{
var clock = new ManualClock();
await using var app = await ExampleApp.StartAsync(services => services
.AddSingleton<TimeProvider>(clock) // before AddEventEmitter(), so the library uses it
.AddOrdersModule()
.AddInventoryModule()
.AddShippingModule());
var carrier = app.Get<CarrierGateway>();
var incomplete = app.Get<IIncompleteEventPublications>();
var completed = app.Get<ICompletedEventPublications>();
app.Say("The carrier API is down while order-1 is completed:");
carrier.IsAvailable = false;
await app.InScopeAsync(sp => sp.GetRequiredService<OrderService>().CompleteAsync("order-1", "alice"));
await WaitForShippingAttemptsAsync(app, 1);
await ShowAsync(app, "Inventory's publication completed; Shipping's is incomplete:");
app.Say("ResubmitOlderThanAsync(5 minutes) skips it, because it was published just now:");
app.Say($" resubmitted: {await incomplete.ResubmitOlderThanAsync(TimeSpan.FromMinutes(5))}");
clock.Advance(TimeSpan.FromMinutes(10));
app.Say("Ten minutes later it qualifies, but the carrier is still down, so it fails again:");
app.Say($" resubmitted: {await incomplete.ResubmitOlderThanAsync(TimeSpan.FromMinutes(5))}");
await WaitForShippingAttemptsAsync(app, 2);
app.Say("The carrier is back. Resubmitting only Shipping's publications:");
carrier.IsAvailable = true;
var count = await incomplete.ResubmitAsync(p => p.ListenerId == "shipping.book-shipment");
app.Say($" resubmitted: {count}");
await app.WaitUntilAsync(async () => (await incomplete.FindAllAsync()).Count == 0);
await ShowAsync(app, "Everything has completed; attempts records each delivery:");
app.Say($"Nothing left to resubmit: {await incomplete.ResubmitAsync(_ => true)}");
async Task ShowAsync(ExampleApp app, string heading)
{
app.Say(heading);
foreach (var p in await completed.FindAllAsync())
{
app.Say(" " + Describe.Publication(p));
}
foreach (var p in await incomplete.FindAllAsync())
{
app.Say(" " + Describe.Publication(p));
}
}
}
private static Task WaitForShippingAttemptsAsync(ExampleApp app, int attempts) =>
app.WaitUntilAsync(async () => (await app.Get<IIncompleteEventPublications>().FindAllAsync())
.Any(p => p.EventType == typeof(StockReserved) && p.Attempts == attempts));
}