Skip to content

Commit 9ba3198

Browse files
author
Daniel Robaszkiewicz
committed
RM: projections
1 parent 6ab5e46 commit 9ba3198

7 files changed

Lines changed: 71 additions & 24 deletions

File tree

src/MicroPlumberd.Services.ProcessManager/ProcessManagerClient.cs

Lines changed: 8 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -103,8 +103,14 @@ public async Task<IAsyncDisposable> SubscribeProcessManager<TProcessManager>() w
103103
c += await Plumber.SubscribeEventHandlerPersistently(sender, $"{typeof(TProcessManager).Name}Outbox", ensureOutputStreamProjection: true);
104104
c += await Plumber.SubscribeEventHandlerPersistently(executor, $"{typeof(TProcessManager).Name}Inbox", ensureOutputStreamProjection: true);
105105

106-
await Plumber.ProjectionManagementClient.EnsureLookupProjection(Plumber.Client, Plumber.ProjectionRegister,
107-
typeof(TProcessManager).Name, "RecipientId", $"{typeof(TProcessManager).Name}Lookup");
106+
var lookupStream = $"{typeof(TProcessManager).Name}Lookup";
107+
var changed = await Plumber.ProjectionManagementClient.EnsureLookupProjection(Plumber.Client, Plumber.ProjectionRegister,
108+
typeof(TProcessManager).Name, "RecipientId", lookupStream);
109+
var logger = _serviceProvider.GetService<ILogger<ProcessManagerClient>>();
110+
if (changed)
111+
logger?.LogInformation("Lookup projection {Projection} was created or updated.", lookupStream);
112+
else
113+
logger?.LogDebug("Lookup projection {Projection} is up-to-date; skipped recreate.", lookupStream);
108114

109115
return c;
110116
}

src/MicroPlumberd.Tests/Integration/ReadModelTests.cs

Lines changed: 14 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -89,6 +89,20 @@ public async Task SubscribeModelWithEventStoreRestart()
8989
fooModel.AssertionDb.Index.Should().HaveCount(2);
9090

9191
}
92+
[Fact]
93+
public async Task TryCreateJoinProjection_SkipsRecreate_OnSecondBoot_WhenQueryUnchanged()
94+
{
95+
await _eventStore.StartInDocker();
96+
97+
var firstResult = await plumber.TryCreateJoinProjection<FooModel>();
98+
firstResult.Should().BeTrue("the projection does not exist on first boot");
99+
100+
var plumber2 = Plumber.Create(_eventStore.GetEventStoreSettings());
101+
102+
var secondResult = await plumber2.TryCreateJoinProjection<FooModel>();
103+
secondResult.Should().BeFalse("the projection's query is unchanged since first boot");
104+
}
105+
92106
[Fact]
93107
public async Task SubscribeScopedModel()
94108
{

src/MicroPlumberd/Abstractions/Api/IPlumberApi.cs

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -513,7 +513,7 @@ Task<IAsyncDisposable> SubscribeStateEventHandler<TEventHandler>(TEventHandler?
513513
/// <returns>
514514
/// A <see cref="Task"/> representing the asynchronous operation.
515515
/// </returns>
516-
Task TryCreateJoinProjection(string outputStream, IEnumerable<string> eventTypes, CancellationToken token = default);
516+
Task<bool> TryCreateJoinProjection(string outputStream, IEnumerable<string> eventTypes, CancellationToken token = default);
517517

518518
/// <summary>
519519
/// Ensures that a join projection is created for the specified event handler type.
@@ -531,7 +531,7 @@ Task<IAsyncDisposable> SubscribeStateEventHandler<TEventHandler>(TEventHandler?
531531
/// <returns>
532532
/// A <see cref="Task"/> representing the asynchronous operation.
533533
/// </returns>
534-
Task TryCreateJoinProjection<TEventHandler>(string? outputStream=null, CancellationToken token = default) where TEventHandler : class, IEventHandler, ITypeRegister;
534+
Task<bool> TryCreateJoinProjection<TEventHandler>(string? outputStream=null, CancellationToken token = default) where TEventHandler : class, IEventHandler, ITypeRegister;
535535

536536
/// <summary>
537537
/// Appends metadata to a stream derived from the specified event type and identifier.

src/MicroPlumberd/EventStoreProjectionManagementClientExtensions.cs

Lines changed: 21 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -17,22 +17,27 @@ public static class KurrentDBProjectionManagementClientExtensions
1717
/// <summary>
1818
/// Attempts to create or update a join projection in the EventStore.
1919
/// </summary>
20-
public static async Task TryCreateJoinProjection(this KurrentDBProjectionManagementClient client,
20+
/// <returns><c>true</c> if the projection was created or updated; <c>false</c> if it was already up-to-date and no change was applied.</returns>
21+
public static async Task<bool> TryCreateJoinProjection(this KurrentDBProjectionManagementClient client,
2122
KurrentDBClient esClient,
2223
string outputStream, IEnumerable<string> eventTypes)
2324
{
2425
var query = CreateQuery(outputStream, eventTypes);
2526

2627
if (await client.ListContinuousAsync().AnyAsync(x => x.Name == outputStream))
27-
await UpdateIfChanged(client, esClient, outputStream, query);
28+
return await UpdateIfChanged(client, esClient, outputStream, query);
2829
else
30+
{
2931
await CreateAndStoreHash(client, esClient, outputStream, query);
32+
return true;
33+
}
3034
}
3135

3236
/// <summary>
3337
/// Ensures the existence and proper configuration of a lookup projection in the EventStore.
3438
/// </summary>
35-
public static async Task EnsureLookupProjection(this KurrentDBProjectionManagementClient client,
39+
/// <returns><c>true</c> if the projection was created or updated; <c>false</c> if it was already up-to-date and no change was applied.</returns>
40+
public static async Task<bool> EnsureLookupProjection(this KurrentDBProjectionManagementClient client,
3641
KurrentDBClient esClient,
3742
IProjectionRegister register,
3843
string category, string eventProperty, string outputStreamCategory,
@@ -42,15 +47,19 @@ public static async Task EnsureLookupProjection(this KurrentDBProjectionManageme
4247
$"fromStreams(['$ce-{category}']).when( {{ \n $any : function(s,e) {{ \n if(e.body && e.body.{eventProperty}) {{\n linkTo('{outputStreamCategory}-' + e.body.{eventProperty}, e) \n }}\n \n }}\n}});";
4348

4449
if ((await register.Get(outputStreamCategory)) != null)
45-
await UpdateIfChanged(client, esClient, outputStreamCategory, query, token);
50+
return await UpdateIfChanged(client, esClient, outputStreamCategory, query, token);
4651
else
52+
{
4753
await CreateAndStoreHash(client, esClient, outputStreamCategory, query, token);
54+
return true;
55+
}
4856
}
4957

5058
/// <summary>
5159
/// Attempts to create or update a join projection in the EventStore.
5260
/// </summary>
53-
public static async Task TryCreateJoinProjection(this KurrentDBProjectionManagementClient client,
61+
/// <returns><c>true</c> if the projection was created or updated; <c>false</c> if it was already up-to-date and no change was applied.</returns>
62+
public static async Task<bool> TryCreateJoinProjection(this KurrentDBProjectionManagementClient client,
5463
KurrentDBClient esClient,
5564
string outputStream, IProjectionRegister register, IEnumerable<string> eventTypes,
5665
CancellationToken token = default)
@@ -62,20 +71,24 @@ public static async Task TryCreateJoinProjection(this KurrentDBProjectionManagem
6271
var query = CreateQuery(outputStream, eventTypes);
6372

6473
if ((await register.Get(outputStream)) != null)
65-
await UpdateIfChanged(client, esClient, outputStream, query, token);
74+
return await UpdateIfChanged(client, esClient, outputStream, query, token);
6675
else
76+
{
6777
await CreateAndStoreHash(client, esClient, outputStream, query, token);
78+
return true;
79+
}
6880
}
6981

70-
private static async Task UpdateIfChanged(KurrentDBProjectionManagementClient client, KurrentDBClient esClient,
82+
private static async Task<bool> UpdateIfChanged(KurrentDBProjectionManagementClient client, KurrentDBClient esClient,
7183
string outputStream, string query, CancellationToken token = default)
7284
{
7385
var newHash = ComputeQueryHash(query);
7486
var existing = await TryGetStoredQueryHash(esClient, outputStream, token);
75-
if (existing == newHash) return;
87+
if (existing == newHash) return false;
7688

7789
await UpdateWithRetry(client, outputStream, query, token);
7890
await StoreQueryHash(esClient, outputStream, newHash, token);
91+
return true;
7992
}
8093

8194
private static async Task CreateAndStoreHash(KurrentDBProjectionManagementClient client, KurrentDBClient esClient,

src/MicroPlumberd/Plumber.cs

Lines changed: 4 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -314,13 +314,13 @@ public Task<IAsyncDisposable> SubscribeStateEventHandler<TEventHandler>(IEnumera
314314
}
315315

316316
/// <inheritdoc />
317-
public Task TryCreateJoinProjection(string outputStream, IEnumerable<string> eventTypes, CancellationToken token = default)
317+
public Task<bool> TryCreateJoinProjection(string outputStream, IEnumerable<string> eventTypes, CancellationToken token = default)
318318
{
319319
return engine.TryCreateJoinProjection(outputStream, eventTypes, token);
320320
}
321321

322322
/// <inheritdoc />
323-
public Task TryCreateJoinProjection<TEventHandler>(string? outputStream = null, CancellationToken token = default) where TEventHandler : class, IEventHandler, ITypeRegister
323+
public Task<bool> TryCreateJoinProjection<TEventHandler>(string? outputStream = null, CancellationToken token = default) where TEventHandler : class, IEventHandler, ITypeRegister
324324
{
325325
return engine.TryCreateJoinProjection<TEventHandler>(outputStream, token);
326326
}
@@ -767,13 +767,13 @@ public Task<IAsyncDisposable> SubscribeStateEventHandler<TEventHandler>(IEnumera
767767
}
768768

769769
/// <inheritdoc />
770-
public Task TryCreateJoinProjection(string outputStream, IEnumerable<string> eventTypes, CancellationToken token = default)
770+
public Task<bool> TryCreateJoinProjection(string outputStream, IEnumerable<string> eventTypes, CancellationToken token = default)
771771
{
772772
return engine.TryCreateJoinProjection(outputStream, eventTypes, token);
773773
}
774774

775775
/// <inheritdoc />
776-
public Task TryCreateJoinProjection<TEventHandler>(string? outputStream = null, CancellationToken token = default) where TEventHandler : class, IEventHandler, ITypeRegister
776+
public Task<bool> TryCreateJoinProjection<TEventHandler>(string? outputStream = null, CancellationToken token = default) where TEventHandler : class, IEventHandler, ITypeRegister
777777
{
778778
return engine.TryCreateJoinProjection<TEventHandler>(outputStream, token);
779779
}

src/MicroPlumberd/PlumberEngine.cs

Lines changed: 20 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -3,6 +3,8 @@
33
using System.Text;
44
using System.Text.Json;
55
using KurrentDB.Client;
6+
using Microsoft.Extensions.DependencyInjection;
7+
using Microsoft.Extensions.Logging;
68
using MicroPlumberd.Utils;
79

810
namespace MicroPlumberd;
@@ -158,7 +160,7 @@ public Task<IAsyncDisposable> SubscribeEventHandler<TEventHandler>(TypeEventConv
158160
/// <returns>
159161
/// A <see cref="Task"/> representing the asynchronous operation.
160162
/// </returns>
161-
public Task TryCreateJoinProjection<TEventHandler>(string? outputStream=null, CancellationToken token = default) where TEventHandler : class, IEventHandler, ITypeRegister
163+
public Task<bool> TryCreateJoinProjection<TEventHandler>(string? outputStream=null, CancellationToken token = default) where TEventHandler : class, IEventHandler, ITypeRegister
162164
{
163165
return TryCreateJoinProjection(outputStream ?? Conventions.OutputStreamModelConvention(typeof(TEventHandler)), _typeHandlerRegisters.GetEventNamesFor<TEventHandler>(), token);
164166
}
@@ -178,9 +180,21 @@ public Task TryCreateJoinProjection<TEventHandler>(string? outputStream=null, Ca
178180
/// <returns>
179181
/// A <see cref="Task"/> representing the asynchronous operation.
180182
/// </returns>
181-
public async Task TryCreateJoinProjection(string outputStream, IEnumerable<string> eventTypes, CancellationToken token = default)
183+
public async Task<bool> TryCreateJoinProjection(string outputStream, IEnumerable<string> eventTypes, CancellationToken token = default)
182184
{
183-
await ProjectionManagementClient.TryCreateJoinProjection(Client, outputStream, ProjectionRegister, eventTypes, token: token);
185+
var changed = await ProjectionManagementClient.TryCreateJoinProjection(Client, outputStream, ProjectionRegister, eventTypes, token: token);
186+
LogProjectionEnsured(outputStream, changed);
187+
return changed;
188+
}
189+
190+
internal void LogProjectionEnsured(string outputStream, bool changed)
191+
{
192+
var logger = ServiceProvider?.GetService<ILogger<PlumberEngine>>();
193+
if (logger == null) return;
194+
if (changed)
195+
logger.LogInformation("Projection {Projection} was created or updated.", outputStream);
196+
else
197+
logger.LogDebug("Projection {Projection} is up-to-date; skipped recreate.", outputStream);
184198
}
185199
/// <summary>
186200
/// Subscribes an event handler to a stream using auto-discovered event types.
@@ -244,7 +258,7 @@ public async Task<IAsyncDisposable> SubscribeStateEventHandler<TEventHandler>(
244258

245259
outputStream ??= Conventions.OutputStreamModelConvention(typeof(TEventHandler));
246260
if (ensureOutputStreamProjection)
247-
await ProjectionManagementClient.TryCreateJoinProjection(Client, outputStream, ProjectionRegister, eventTypes, token: token);
261+
LogProjectionEnsured(outputStream, await ProjectionManagementClient.TryCreateJoinProjection(Client, outputStream, ProjectionRegister, eventTypes, token: token));
248262
var sub = Subscribe(outputStream, start ?? FromStream.Start, cancellationToken: token);
249263
if (eh == null)
250264
await sub.WithSnapshotHandler<TEventHandler>();
@@ -274,7 +288,7 @@ public async Task<IAsyncDisposable> SubscribeEventHandler<TEventHandler>(TypeEve
274288

275289
outputStream ??= Conventions.OutputStreamModelConvention(typeof(TEventHandler));
276290
if (ensureOutputStreamProjection)
277-
await ProjectionManagementClient.TryCreateJoinProjection(Client, outputStream, ProjectionRegister, eventTypes, token: token);
291+
LogProjectionEnsured(outputStream, await ProjectionManagementClient.TryCreateJoinProjection(Client, outputStream, ProjectionRegister, eventTypes, token: token));
278292
var sub = Subscribe(outputStream, start ?? FromStream.Start, cancellationToken:token);
279293
if (eh == null)
280294
await sub.WithHandler<TEventHandler>(mapFunc);
@@ -311,7 +325,7 @@ public async Task<IAsyncDisposable> SubscribeEventHandlerPersistently<TEventHand
311325
outputStream ??= Conventions.OutputStreamModelConvention(handlerType);
312326
groupName ??= Conventions.GroupNameModelConvention(handlerType);
313327
if (ensureOutputStreamProjection)
314-
await ProjectionManagementClient.TryCreateJoinProjection(Client, outputStream, ProjectionRegister, events, token);
328+
LogProjectionEnsured(outputStream, await ProjectionManagementClient.TryCreateJoinProjection(Client, outputStream, ProjectionRegister, events, token: token));
315329

316330
try
317331
{

src/MicroPlumberd/SubscriptionSet.cs

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -24,7 +24,7 @@ public IEngineSubscriptionSet With<TModel>(TModel model)
2424
public async Task SubscribePersistentlyAsync(OperationContext context, string outputStream, string? groupName = null)
2525
{
2626
groupName ??= outputStream;
27-
await plumber.ProjectionManagementClient.TryCreateJoinProjection(plumber.Client, outputStream, _register.Keys);
27+
plumber.LogProjectionEnsured(outputStream, await plumber.ProjectionManagementClient.TryCreateJoinProjection(plumber.Client, outputStream, _register.Keys));
2828
var subscription = plumber.PersistentSubscriptionClient.SubscribeToStream(outputStream, groupName);
2929
var state = Tuple.Create(this,context, subscription);
3030

@@ -52,7 +52,7 @@ await Task.Factory.StartNew(static async (x) =>
5252

5353
public async Task SubscribeAsync(OperationContext context, string name, FromStream start)
5454
{
55-
await plumber.ProjectionManagementClient.TryCreateJoinProjection(plumber.Client, name, _register.Keys);
55+
plumber.LogProjectionEnsured(name, await plumber.ProjectionManagementClient.TryCreateJoinProjection(plumber.Client, name, _register.Keys));
5656

5757
KurrentDBClient.StreamSubscriptionResult subscription = plumber.Client.SubscribeToStream(name, start, true);
5858
var state = Tuple.Create(this, context,subscription);

0 commit comments

Comments
 (0)