-
Notifications
You must be signed in to change notification settings - Fork 0
Expand file tree
/
Copy pathDurableStorageExample.cs
More file actions
86 lines (75 loc) · 3.5 KB
/
Copy pathDurableStorageExample.cs
File metadata and controls
86 lines (75 loc) · 3.5 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
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
using Codefinity.EventEmitter.Sample.Infrastructure;
using Codefinity.EventEmitter.Sample.Orders;
using Microsoft.Extensions.DependencyInjection;
using Microsoft.Extensions.Logging;
namespace Codefinity.EventEmitter.Sample.Examples;
/// <summary>
/// A custom, durable publication repository (a JSON file) together with ShutdownTimeout and
/// RepublishOutstandingEventsOnStartup: an event that couldn't be handled before the app stopped is
/// delivered when the app starts again.
/// </summary>
internal static class DurableStorageExample
{
private sealed record BillingSettings(TimeSpan InvoiceDuration);
private sealed class InvoiceListener(BillingSettings settings, ILogger<InvoiceListener> logger)
{
// An explicit id: publications stored under it must still match after a redeploy.
[ApplicationModuleListener(Id = "billing.create-invoice")]
public async Task On(OrderCompleted evt, CancellationToken cancellationToken)
{
logger.LogInformation("Creating the invoice for {OrderId}...", evt.OrderId);
await Task.Delay(settings.InvoiceDuration, cancellationToken);
logger.LogInformation("Invoice for {OrderId} created", evt.OrderId);
}
}
public static async Task RunAsync()
{
var path = Path.Combine(Path.GetTempPath(), $"event-publications-{Guid.NewGuid():N}.json");
try
{
await FirstRunAsync(path);
await SecondRunAsync(path);
}
finally
{
File.Delete(path);
}
}
private static async Task FirstRunAsync(string path)
{
var app = await ExampleApp.StartAsync(services =>
{
services.AddSingleton(new BillingSettings(InvoiceDuration: TimeSpan.FromSeconds(10)));
services
.AddEventEmitter(o => o.ShutdownTimeout = TimeSpan.FromMilliseconds(300))
.UsePublicationRepository(new JsonFileEventPublicationRepository(path))
.AddListener<InvoiceListener>();
});
app.Say($"Run 1: publications are stored in {Path.GetFileName(path)}. Invoices are slow today.");
await app.PublishAsync(new OrderCompleted("order-1", "alice"));
await Task.Delay(100);
app.Say("Stopping the app mid-invoice (a deploy). ShutdownTimeout is 300 ms, then the listener is cancelled:");
await app.StopAsync();
await app.DisposeAsync();
var stored = await new JsonFileEventPublicationRepository(path).FindIncompleteAsync();
Console.WriteLine($" Left in the file: {string.Join("; ", stored.Select(Describe.Publication))}");
}
private static async Task SecondRunAsync(string path)
{
Console.WriteLine(" Run 2: a new process starts with RepublishOutstandingEventsOnStartup = true, and invoices are fast again:");
await using var app = await ExampleApp.StartAsync(services =>
{
services.AddSingleton(new BillingSettings(InvoiceDuration: TimeSpan.FromMilliseconds(100)));
services
.AddEventEmitter(o => o.RepublishOutstandingEventsOnStartup = true)
.UsePublicationRepository(new JsonFileEventPublicationRepository(path))
.AddListener<InvoiceListener>();
});
await app.WaitForCompletedAsync(1);
app.Say("The publication from run 1 is now complete:");
foreach (var p in await app.Get<ICompletedEventPublications>().FindAllAsync())
{
app.Say(" " + Describe.Publication(p));
}
}
}