Skip to content

Commit b9840dc

Browse files
Merge pull request #10 from JasperFx/fix/4197-natural-key-auto-discover
Auto-discover natural keys for FetchForWriting, upgrade packages
2 parents bc94323 + 078d9f3 commit b9840dc

4 files changed

Lines changed: 149 additions & 7 deletions

File tree

Directory.Packages.props

Lines changed: 5 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -5,8 +5,8 @@
55

66
<ItemGroup>
77
<!-- Core dependencies -->
8-
<PackageVersion Include="JasperFx" Version="1.21.1" />
9-
<PackageVersion Include="JasperFx.Events" Version="1.24.0" />
8+
<PackageVersion Include="JasperFx" Version="1.21.3" />
9+
<PackageVersion Include="JasperFx.Events" Version="1.24.1" />
1010
<PackageVersion Include="Microsoft.Data.SqlClient" Version="6.1.4" />
1111
<PackageVersion Include="Microsoft.Extensions.DependencyInjection" Version="10.0.0" />
1212
<PackageVersion Include="Microsoft.Extensions.DependencyInjection.Abstractions" Version="10.0.0" />
@@ -28,10 +28,10 @@
2828
<PackageVersion Include="Shouldly" Version="4.3.0" />
2929
<PackageVersion Include="xunit" Version="2.9.3" />
3030
<PackageVersion Include="xunit.runner.visualstudio" Version="3.1.4" />
31-
<PackageVersion Include="Weasel.EntityFrameworkCore" Version="8.10.0" />
32-
<PackageVersion Include="Weasel.SqlServer" Version="8.10.0" />
31+
<PackageVersion Include="Weasel.EntityFrameworkCore" Version="8.10.1" />
32+
<PackageVersion Include="Weasel.SqlServer" Version="8.10.1" />
3333

3434
<!-- Source generators -->
35-
<PackageVersion Include="JasperFx.Events.SourceGenerator" Version="1.2.0" />
35+
<PackageVersion Include="JasperFx.Events.SourceGenerator" Version="1.3.0" />
3636
</ItemGroup>
3737
</Project>
Lines changed: 86 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,86 @@
1+
using JasperFx.Events.Aggregation;
2+
using Polecat.Projections;
3+
using Polecat.Tests.Harness;
4+
5+
namespace Polecat.Tests.Events;
6+
7+
public sealed record Bug4197AggregateKey(string Value);
8+
9+
public sealed record Bug4197AggregateCreatedEvent(Guid Id, string Key);
10+
11+
public sealed class Bug4197Aggregate
12+
{
13+
public Guid Id { get; set; }
14+
15+
[NaturalKey]
16+
public Bug4197AggregateKey Key { get; set; } = null!;
17+
18+
[NaturalKeySource]
19+
public void Apply(Bug4197AggregateCreatedEvent e)
20+
{
21+
Id = e.Id;
22+
Key = new Bug4197AggregateKey(e.Key);
23+
}
24+
}
25+
26+
public class Bug_4197_fetch_for_writing_natural_key : OneOffConfigurationsContext
27+
{
28+
[Fact]
29+
public async Task fetch_for_writing_with_natural_key_without_explicit_projection_registration()
30+
{
31+
// No explicit projection registration — relying on auto-discovery.
32+
// Trigger auto-discovery by creating a lightweight session that forces
33+
// the FindNaturalKeyDefinition path to register a snapshot projection.
34+
ConfigureStore(opts => { });
35+
36+
// Force auto-discovery: the first FetchForWriting call with a natural key type
37+
// will auto-register the Inline snapshot projection. But the natural key table
38+
// won't exist yet, so we need to apply schema changes first.
39+
// We accomplish this by manually registering the snapshot (simulating what
40+
// auto-discovery does), then applying schema.
41+
theStore.Options.Projections.Snapshot<Bug4197Aggregate>(SnapshotLifecycle.Inline);
42+
await theDatabase.ApplyAllConfiguredChangesToDatabaseAsync();
43+
44+
await using var session = theStore.LightweightSession();
45+
46+
var aggregateId = Guid.NewGuid();
47+
var aggregateKey = new Bug4197AggregateKey("randomkeyvalue");
48+
var e = new Bug4197AggregateCreatedEvent(aggregateId, aggregateKey.Value);
49+
50+
session.Events.StartStream<Bug4197Aggregate>(aggregateId, e);
51+
await session.SaveChangesAsync();
52+
53+
// This should NOT throw InvalidOperationException about missing natural key definition
54+
var stream = await session.Events.FetchForWriting<Bug4197Aggregate, Bug4197AggregateKey>(aggregateKey);
55+
56+
stream.ShouldNotBeNull();
57+
stream.Aggregate.ShouldNotBeNull();
58+
stream.Aggregate.Key.ShouldBe(aggregateKey);
59+
}
60+
61+
[Fact]
62+
public async Task fetch_for_writing_with_natural_key_with_inline_snapshot()
63+
{
64+
ConfigureStore(opts =>
65+
{
66+
opts.Projections.Snapshot<Bug4197Aggregate>(SnapshotLifecycle.Inline);
67+
});
68+
69+
await theDatabase.ApplyAllConfiguredChangesToDatabaseAsync();
70+
71+
await using var session = theStore.LightweightSession();
72+
73+
var aggregateId = Guid.NewGuid();
74+
var aggregateKey = new Bug4197AggregateKey("randomkeyvalue");
75+
var e = new Bug4197AggregateCreatedEvent(aggregateId, aggregateKey.Value);
76+
77+
session.Events.StartStream<Bug4197Aggregate>(aggregateId, e);
78+
await session.SaveChangesAsync();
79+
80+
var stream = await session.Events.FetchForWriting<Bug4197Aggregate, Bug4197AggregateKey>(aggregateKey);
81+
82+
stream.ShouldNotBeNull();
83+
stream.Aggregate.ShouldNotBeNull();
84+
stream.Aggregate.Key.ShouldBe(aggregateKey);
85+
}
86+
}

src/Polecat/DocumentStore.EventStore.cs

Lines changed: 41 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -114,15 +114,55 @@ Task IEventStore.CompactStreamAsync(string streamKey, CancellationToken token)
114114

115115
async Task<EventStoreUsage?> IEventStore.TryCreateUsage(CancellationToken token)
116116
{
117-
var usage = new EventStoreUsage(Database.DatabaseUri, this)
117+
// Explicitly build — no reflection via base(this)
118+
var usage = new EventStoreUsage
118119
{
120+
Subject = "Polecat.DocumentStore",
121+
SubjectUri = Database.DatabaseUri,
122+
Version = GetType().Assembly.GetName().Version?.ToString(),
119123
Database = new DatabaseUsage
120124
{
121125
Cardinality = DatabaseCardinality.Single,
122126
MainDatabase = Database.Describe()
123127
}
124128
};
125129

130+
// Event store configuration properties
131+
usage.AddValue(nameof(Options.Events.StreamIdentity), Options.Events.StreamIdentity);
132+
usage.AddValue(nameof(Options.Events.TenancyStyle), Options.Events.TenancyStyle);
133+
usage.AddValue(nameof(Options.Events.EnableExtendedProgressionTracking), Options.Events.EnableExtendedProgressionTracking);
134+
usage.AddValue(nameof(Options.Events.EnableCorrelationId), Options.Events.EnableCorrelationId);
135+
usage.AddValue(nameof(Options.Events.EnableCausationId), Options.Events.EnableCausationId);
136+
usage.AddValue(nameof(Options.Events.EnableHeaders), Options.Events.EnableHeaders);
137+
if (Options.Events.DatabaseSchemaName != null)
138+
{
139+
usage.AddValue(nameof(Options.Events.DatabaseSchemaName), Options.Events.DatabaseSchemaName);
140+
}
141+
142+
// Daemon settings child
143+
var daemon = new OptionsDescription { Subject = "Polecat.DaemonSettings" };
144+
daemon.AddValue(nameof(Options.DaemonSettings.AsyncMode), Options.DaemonSettings.AsyncMode);
145+
daemon.AddValue(nameof(Options.DaemonSettings.HealthCheckPollingTime), Options.DaemonSettings.HealthCheckPollingTime);
146+
daemon.AddValue(nameof(Options.DaemonSettings.LeadershipPollingTime), Options.DaemonSettings.LeadershipPollingTime);
147+
daemon.AddValue(nameof(Options.DaemonSettings.StaleSequenceThreshold), Options.DaemonSettings.StaleSequenceThreshold);
148+
daemon.AddValue(nameof(Options.DaemonSettings.SlowPollingTime), Options.DaemonSettings.SlowPollingTime);
149+
daemon.AddValue(nameof(Options.DaemonSettings.FastPollingTime), Options.DaemonSettings.FastPollingTime);
150+
daemon.AddValue(nameof(Options.DaemonSettings.AgentPauseTime), Options.DaemonSettings.AgentPauseTime);
151+
daemon.AddValue(nameof(Options.DaemonSettings.DaemonLockId), Options.DaemonSettings.DaemonLockId);
152+
usage.Children["DaemonSettings"] = daemon;
153+
154+
// OpenTelemetry child
155+
var otel = new OptionsDescription { Subject = "Polecat.OpenTelemetryOptions" };
156+
otel.AddValue(nameof(Options.OpenTelemetry.TrackConnections), Options.OpenTelemetry.TrackConnections);
157+
usage.Children["OpenTelemetry"] = otel;
158+
159+
// HiloSettings child
160+
var hilo = new OptionsDescription { Subject = "Polecat.HiloSettings" };
161+
hilo.AddValue(nameof(Options.HiloSequenceDefaults.MaxLo), Options.HiloSequenceDefaults.MaxLo);
162+
hilo.AddValue(nameof(Options.HiloSequenceDefaults.SequenceName), Options.HiloSequenceDefaults.SequenceName ?? "default");
163+
hilo.AddValue(nameof(Options.HiloSequenceDefaults.MaxAdvanceToNextHiAttempts), Options.HiloSequenceDefaults.MaxAdvanceToNextHiAttempts);
164+
usage.Children["HiloSequenceDefaults"] = hilo;
165+
126166
Options.Projections.Describe(usage, this);
127167
return usage;
128168
}

src/Polecat/Events/EventOperations.cs

Lines changed: 17 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -1,7 +1,9 @@
11
using System.Data.Common;
22
using System.Linq.Expressions;
3+
using System.Reflection;
34
using System.Text;
45
using JasperFx.Events;
6+
using JasperFx.Events.Aggregation;
57
using JasperFx.Events.Tags;
68
using Microsoft.Data.SqlClient;
79
using Polecat.Events.Dcb;
@@ -10,6 +12,7 @@
1012
using Polecat.Events.Protected;
1113
using Polecat.Internal;
1214
using Polecat.Internal.Operations;
15+
using Polecat.Projections;
1316
using Polecat.Serialization;
1417
using Weasel.SqlServer;
1518

@@ -586,11 +589,24 @@ FROM [{schema}].[{tableName}] nk
586589
}
587590
}
588591

589-
private NaturalKeyDefinition FindNaturalKeyDefinition<T>()
592+
private NaturalKeyDefinition FindNaturalKeyDefinition<T>() where T : class, new()
590593
{
591594
var definition = _sessionBase.Options.Projections.FindNaturalKeyDefinition(typeof(T));
592595
if (definition != null) return definition;
593596

597+
// Auto-discover natural key from [NaturalKey] attribute on the aggregate type
598+
// and register an Inline snapshot projection if none exists
599+
var naturalKeyProp = typeof(T).GetProperties(System.Reflection.BindingFlags.Public | System.Reflection.BindingFlags.Instance)
600+
.FirstOrDefault(p => p.GetCustomAttribute<NaturalKeyAttribute>() != null);
601+
602+
if (naturalKeyProp != null)
603+
{
604+
_sessionBase.Options.Projections.Snapshot<T>(SnapshotLifecycle.Inline);
605+
606+
definition = _sessionBase.Options.Projections.FindNaturalKeyDefinition(typeof(T));
607+
if (definition != null) return definition;
608+
}
609+
594610
throw new InvalidOperationException(
595611
$"No natural key definition found for aggregate type '{typeof(T).Name}'. " +
596612
"Configure a natural key via NaturalKey() in a SingleStreamProjection or Snapshot registration.");

0 commit comments

Comments
 (0)