Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
@@ -0,0 +1,125 @@
using System;
using System.Data;
using System.Threading;
using System.Threading.Tasks;
using JasperFx.Events;
using JasperFx.Events.Aggregation;
using Marten;
using Marten.Events.Projections;
using Marten.Testing.Harness;
using Shouldly;
using Xunit;

namespace EventSourcingTests.Bugs;

public class Bug_4788_natrual_key_table_not_populated_on_rebuild: OneOffConfigurationsContext
{
private const string schemaName = "bug_4788";

public sealed record OrderNumber(string Value);

public sealed record OrderPlaced(Guid OrderId, string OrderNumber);

public sealed record OrderShipped(Guid OrderId, string TrackingNumber);

public sealed record Order
{
public Guid Id { get; set; }

[NaturalKey]
public OrderNumber Number { get; set; }

public string? TrackingNumber { get; set; }

[NaturalKeySource]
public static Order Create(OrderPlaced e)
{
return new Order
{
Id = e.OrderId,
Number = new OrderNumber(e.OrderNumber)
};
}

public static Order Apply(OrderShipped e, Order order)
{
return new Order
{
TrackingNumber = e.TrackingNumber
};
}
}

private static void ConfigureStore(StoreOptions opts)
{
opts.Advanced.Migrator.NameDataLength = 100;
opts.Connection(ConnectionSource.ConnectionString);
opts.DatabaseSchemaName = schemaName;

opts.Events.StreamIdentity = StreamIdentity.AsGuid;
opts.Events.AppendMode = EventAppendMode.Quick;

opts.Projections.Snapshot<Order>(SnapshotLifecycle.Inline);
}

private async Task<object> ExecuteScalar(string sql)
{
await using var conn = theStore.Storage.Database.CreateConnection();
await conn.OpenAsync();
await using var cmd = conn.CreateCommand();
cmd.CommandText = sql;
return await cmd.ExecuteScalarAsync();
}

private async Task<DataTable> GetData(string sql)
{
await using var conn = theStore.Storage.Database.CreateConnection();
await conn.OpenAsync();
await using var cmd = conn.CreateCommand();
cmd.CommandText = sql;
await using var reader = await cmd.ExecuteReaderAsync();
var dataTable = new DataTable();
dataTable.Load(reader);
return dataTable;
}

[Fact]
public async Task rebuild_snapshot_should_populate_naturalkey_table()
{
StoreOptions(ConfigureStore);
await theStore.Storage.ApplyAllConfiguredChangesToDatabaseAsync();

var streamId = Guid.NewGuid();
await using var session = theStore.LightweightSession();
session.Events.StartStream<Order>(streamId, new OrderPlaced(streamId, "12345"));
await session.SaveChangesAsync();

var anotherStore = SeparateStore(ConfigureStore);
var anotherSession = anotherStore.LightweightSession();
var order1 = await anotherSession.Events.FetchLatest<Order, OrderNumber>(new OrderNumber("12345"));
order1.ShouldNotBeNull();

var naturalKeyTableName = $"{schemaName}.mt_natural_key_order";
var natrualKeys = await GetData($"SELECT * FROM {naturalKeyTableName}");
natrualKeys.Rows.Count.ShouldBe(1);

await ExecuteScalar($"DELETE FROM {naturalKeyTableName}");
natrualKeys = await GetData($"SELECT * FROM {naturalKeyTableName}");
natrualKeys.Rows.Count.ShouldBe(0);

await using var conn = theStore.Storage.Database.CreateConnection();
await conn.OpenAsync();

var daemon = await SeparateStore(ConfigureStore).BuildProjectionDaemonAsync();
await daemon.PrepareForRebuildsAsync();
await daemon.RebuildProjectionAsync($"{nameof(Bug_4788_natrual_key_table_not_populated_on_rebuild)}.order", CancellationToken.None);

var afterRebuildStore = SeparateStore(ConfigureStore);
var afterRebuildSession = afterRebuildStore.LightweightSession();
var order2 = await afterRebuildSession.Events.FetchLatest<Order, OrderNumber>(new OrderNumber("12345"));
order2.ShouldNotBeNull();

natrualKeys = await GetData($"SELECT * FROM {naturalKeyTableName}");
natrualKeys.Rows.Count.ShouldBe(1);
}
}
39 changes: 39 additions & 0 deletions src/Marten/DocumentStore.EventStore.cs
Original file line number Diff line number Diff line change
Expand Up @@ -323,6 +323,17 @@ private void teardownProjectionStorage(IProjectionSource<IDocumentOperations, IQ

// Rewind previous DeadLetterEvents because you're going to replay them all anyway
session.DeleteWhere<DeadLetterEvent>(x => x.ProjectionName == source.Name);

// #4788: a [NaturalKey] aggregate maintains its mt_natural_key_X lookup table via the
// auto-registered NaturalKeyProjection on the inline-append path. Teardown of the parent
// projection must also wipe the natural-key table so the rebuild path repopulates it from
// scratch (the rebuild itself re-emits the upserts via StartProjectionBatchAsync).
if (source is IAggregateProjection aggregateSource && aggregateSource.NaturalKeyDefinition != null)
{
var naturalKeyTable =
$"{Events.DatabaseSchemaName}.mt_natural_key_{aggregateSource.NaturalKeyDefinition.AggregateType.Name.ToLowerInvariant()}";
session.QueueSqlCommand($"delete from {naturalKeyTable}");
}
}

public async ValueTask<IProjectionBatch<IDocumentOperations, IQuerySession>> StartProjectionBatchAsync(
Expand Down Expand Up @@ -361,6 +372,34 @@ public async ValueTask<IProjectionBatch<IDocumentOperations, IQuerySession>> Sta

await projectionBatch.RecordProgress(range).ConfigureAwait(false);

// #4788: when rebuilding a snapshot whose aggregate has a [NaturalKey], re-emit the
// natural-key upserts for this page's events. ApplyAsync on the inline path drives off
// newly-appended StreamActions and never fires during rebuild — without this hook the
// mt_natural_key_X table stays empty after teardown. Routing the upsert SQL through
// ProjectionBatch.SessionForTenant returns a ProjectionDocumentSession whose work-tracker
// IS the ProjectionUpdateBatch, so the operations flush alongside the rebuilt snapshots
// inside the same batch transaction.
if (mode == ShardExecutionMode.Rebuild && range.Events.Any())
{
var naturalKeySource = Options.Projections.All
.FirstOrDefault(s => s.Name.EqualsIgnoreCase(range.ShardName.Name))
as IAggregateProjection;
if (naturalKeySource?.NaturalKeyDefinition == null)
{
naturalKeySource = null;
}

if (naturalKeySource != null)
{
var naturalKeyProjection = new NaturalKeyProjection(Options.EventGraph, naturalKeySource.NaturalKeyDefinition!);
foreach (var byTenant in range.Events.GroupBy(e => e.TenantId ?? StorageConstants.DefaultTenantId))
{
var ops = projectionBatch.SessionForTenant(byTenant.Key);
naturalKeyProjection.QueueUpsertsForEvents(ops, byTenant);
}
}
}

return projectionBatch;
}

Expand Down
45 changes: 38 additions & 7 deletions src/Marten/Events/Projections/NaturalKeyProjection.cs
Original file line number Diff line number Diff line change
Expand Up @@ -52,7 +52,7 @@ public Task ApplyAsync(IDocumentOperations operations, IEnumerable<StreamAction>
var innerValue = _naturalKey.Unwrap(rawValue);
if (innerValue != null)
{
queueUpsertSql(operations, stream, innerValue);
queueUpsertSql(operations, stream.Id, stream.Key, stream.TenantId, innerValue);
}
}
}
Expand All @@ -62,10 +62,41 @@ public Task ApplyAsync(IDocumentOperations operations, IEnumerable<StreamAction>
return Task.CompletedTask;
}

private void queueUpsertSql(IDocumentOperations operations, StreamAction stream, object innerValue)
/// <summary>
/// #4788: rebuild-time counterpart to <see cref="ApplyAsync"/>. The async-daemon rebuild path
/// replays already-persisted events without appending streams, so <c>ApplyAsync</c>'s
/// stream-driven dispatch never fires and the <c>mt_natural_key_X</c> table stays empty after
/// teardown. This entry-point feeds raw <see cref="IEvent"/>s straight through the same
/// upsert SQL builder, pulling stream id/key + tenant id off the event itself (events written
/// to <c>mt_events</c> always carry these). Called from <c>StartProjectionBatchAsync</c> per
/// rebuild page, with <paramref name="operations"/> routed through
/// <see cref="ProjectionBatch.SessionForTenant"/> so the SQL flushes into the projection batch
/// rather than the bare session's unit-of-work.
/// </summary>
internal void QueueUpsertsForEvents(IDocumentOperations operations, IEnumerable<IEvent> events)
{
foreach (var @event in events)
{
foreach (var mapping in _naturalKey.EventMappings)
{
if (mapping.EventType.IsAssignableFrom(@event.Data.GetType()))
{
var rawValue = mapping.Extractor(@event.Data);
var innerValue = _naturalKey.Unwrap(rawValue);
if (innerValue != null)
{
queueUpsertSql(operations, @event.StreamId, @event.StreamKey, @event.TenantId, innerValue);
}
}
}
}
}

private void queueUpsertSql(IDocumentOperations operations, Guid streamId, string? streamKey,
string tenantId, object innerValue)
{
var streamCol = _isGuid ? "stream_id" : "stream_key";
object streamId = _isGuid ? (object)stream.Id : stream.Key!;
object streamIdValue = _isGuid ? (object)streamId : streamKey!;

// When UseArchivedStreamPartitioning is on, is_archived is part of the PK
// and must be included in the ON CONFLICT clause
Expand All @@ -74,28 +105,28 @@ private void queueUpsertSql(IDocumentOperations operations, StreamAction stream,
var sql = $"INSERT INTO {_tableName} (natural_key_value, {streamCol}, tenant_id, is_archived) " +
$"VALUES (?, ?, ?, false) " +
$"ON CONFLICT (natural_key_value, tenant_id, is_archived) DO UPDATE SET {streamCol} = ?";
operations.QueueSqlCommand(sql, innerValue, streamId, stream.TenantId, streamId);
operations.QueueSqlCommand(sql, innerValue, streamIdValue, tenantId, streamIdValue);
}
else if (_isConjoined)
{
var sql = $"INSERT INTO {_tableName} (natural_key_value, {streamCol}, tenant_id, is_archived) " +
$"VALUES (?, ?, ?, false) " +
$"ON CONFLICT (natural_key_value, tenant_id) DO UPDATE SET {streamCol} = ?, is_archived = false";
operations.QueueSqlCommand(sql, innerValue, streamId, stream.TenantId, streamId);
operations.QueueSqlCommand(sql, innerValue, streamIdValue, tenantId, streamIdValue);
}
else if (_useArchivedPartitioning)
{
var sql = $"INSERT INTO {_tableName} (natural_key_value, {streamCol}, is_archived) " +
$"VALUES (?, ?, false) " +
$"ON CONFLICT (natural_key_value, is_archived) DO UPDATE SET {streamCol} = ?";
operations.QueueSqlCommand(sql, innerValue, streamId, streamId);
operations.QueueSqlCommand(sql, innerValue, streamIdValue, streamIdValue);
}
else
{
var sql = $"INSERT INTO {_tableName} (natural_key_value, {streamCol}, is_archived) " +
$"VALUES (?, ?, false) " +
$"ON CONFLICT (natural_key_value) DO UPDATE SET {streamCol} = ?, is_archived = false";
operations.QueueSqlCommand(sql, innerValue, streamId, streamId);
operations.QueueSqlCommand(sql, innerValue, streamIdValue, streamIdValue);
}
}

Expand Down
Loading