From b187d00c10abf1aa9d810f40f7d072b89b0cdf74 Mon Sep 17 00:00:00 2001 From: "Jeremy D. Miller" Date: Tue, 7 Jul 2026 05:59:57 -0500 Subject: [PATCH 1/7] Consume Marten 9.13.0-alpha.3 + JasperFx 2.21.0 (epic jasperfx#486 WS0/WS1) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Marten 9.13.0-alpha.3 carries the per-tenant distribution work this PR's e2e tests exercise: marten#4864 (broadened DistributesAgentsPerTenant gate + single-DB descriptor TenantIds), marten#4865 (ExternallyManaged honored), marten#4869 (per-tenant high-water advancement — unblocks the failover floor-continuity test), and marten#4870 (native per-tenant distribution). JasperFx/JasperFx.Events pins move 2.20.0 -> 2.21.0 to match what Marten alpha.3 already pulls transitively — keeps the explicit pins from drifting below the resolved graph. No envelope-operation contract drift from marten#4872 (Weasel 9.12 metadata binders): Wolverine.Marten's StoreIncomingEnvelope/StoreOutgoingEnvelope compile unchanged. Co-Authored-By: Claude Fable 5 --- Directory.Packages.props | 6 +++--- 1 file changed, 3 insertions(+), 3 deletions(-) diff --git a/Directory.Packages.props b/Directory.Packages.props index f20df00cc..f1d92f9c8 100644 --- a/Directory.Packages.props +++ b/Directory.Packages.props @@ -35,15 +35,15 @@ - - + + - + From 3ca9067e70a965be437e8e022d5232181bdd2f93 Mon Sep 17 00:00:00 2001 From: "Jeremy D. Miller" Date: Mon, 6 Jul 2026 20:01:12 -0500 Subject: [PATCH 2/7] E2E: runtime tenant churn under Wolverine-managed distribution (jasperfx#486 WS1) Sharded + UseTenantPartitionedEvents, single node: with tenant-a's per-tenant agent already projecting, add tenant-b to the same shard database at RUNTIME via Advanced.AddTenantToShardAsync. Within the assignment-evaluation window (CheckAssignmentPeriod = 1s) the leader re-enumerates the store usage, fans out the ":all/tenant-b" agent, tenant-b's fresh seq-1.. events project, its HighWaterMark-independent PartitionedCounter:All:tenant-b progression row lands in the shard database, and no store-global (tenant-less) agent ever runs alongside the per-tenant agents (the #3328 stale-agent-retirement edge). Tenant REMOVAL is deliberately not covered end-to-end: Marten 9.13.0-alpha.2 exposes RemoveTenantAsync/DisableTenantAsync (ShardedTenancy.cs:378/:470), but a running store's usage never shrinks - MartenDatabase.TenantIds is add-only (Fill at ShardedTenancy.cs:211/:364) and DescribeDatabasesAsync copies that in-memory list, so Wolverine keeps enumerating the removed tenant's agent until a process restart. The Wolverine half of the removal transition stays pinned by the running_per_tenant_agent_is_stopped_when_its_tenant_is_removed unit test. Co-Authored-By: Claude Fable 5 --- ...tenant_churn_under_managed_distribution.cs | 272 ++++++++++++++++++ 1 file changed, 272 insertions(+) create mode 100644 src/Persistence/MartenTests/Distribution/runtime_tenant_churn_under_managed_distribution.cs diff --git a/src/Persistence/MartenTests/Distribution/runtime_tenant_churn_under_managed_distribution.cs b/src/Persistence/MartenTests/Distribution/runtime_tenant_churn_under_managed_distribution.cs new file mode 100644 index 000000000..1530837f2 --- /dev/null +++ b/src/Persistence/MartenTests/Distribution/runtime_tenant_churn_under_managed_distribution.cs @@ -0,0 +1,272 @@ +using IntegrationTests; +using JasperFx.CodeGeneration; +using JasperFx.Core; +using JasperFx.Events; +using JasperFx.Events.Projections; +using JasperFx.MultiTenancy; +using Marten; +using Marten.Events; +using MartenTests.Distribution.Support; +using Microsoft.Extensions.DependencyInjection; +using Microsoft.Extensions.Hosting; +using Microsoft.Extensions.Logging; +using Npgsql; +using Shouldly; +using Weasel.Core; +using Weasel.Postgresql; +using Weasel.Postgresql.Migrations; +using Wolverine; +using Wolverine.Marten; +using Wolverine.Runtime.Agents; +using Wolverine.Tracking; +using Xunit; +using Xunit.Abstractions; + +namespace MartenTests.Distribution; + +// JasperFx/jasperfx#486 WS1 — runtime tenant churn under Wolverine-managed distribution. The sibling +// sharded_tenant_partitioned_distribution test provisions its co-located tenants BEFORE the host starts, +// so the first agent enumeration already fans out per tenant. This test covers the RUNTIME transition: +// with the host up and tenant-a's per-tenant agent already projecting, a NEW tenant is added to the same +// shard database via Advanced.AddTenantToShardAsync. Within the assignment-evaluation window +// (Durability.CheckAssignmentPeriod) the leader must re-enumerate the store's usage, see the new tenant, +// fan out its ":all/tenant-b" agent, and start projecting the new tenant's events from its OWN per-tenant +// sequence 1.. — all WITHOUT ever running a store-global (tenant-less) agent for the shard alongside the +// per-tenant agents, which would double-process the same events (the stale-agent-retirement edge of +// wolverine#3328, EventSubscriptionAgentFamily.RetireSupersededAgentsAsync). +// +// TENANT REMOVAL (the inverse transition) IS DELIBERATELY NOT COVERED END-TO-END, because no runtime +// removal path exists end-to-end in published Marten 9.13.0-alpha.2: +// - Marten DOES expose removal APIs on ShardedTenancy (src/Marten/Storage/ShardedTenancy.cs): +// RemoveTenantAsync(tenantId, ct) (:378) hard-deletes the pool assignment row, and +// DisableTenantAsync(tenantId) (:470) soft-deletes (disabled = true) — both via the store-agnostic +// IDynamicTenantSource surface from marten#4607. +// - BUT the running store's usage never shrinks: MartenDatabase.TenantIds is only ever ADDED to — +// BuildDatabases() re-reads the assignment table on every call but only Fill()s ids in +// (ShardedTenancy.cs:211), as do AssignTenantAsync (:364) and AddTenantToShardAsync (:660) — and +// nothing ever removes an id from that in-memory list. DescribeDatabasesAsync (:735-751) copies +// exactly that list onto the DatabaseDescriptor that IEventStore.TryCreateUsage returns. +// - Wolverine's retirement only stops an agent that the current SupportedAgentsAsync enumeration no +// longer lists (EventSubscriptionAgentFamily.RetireSupersededAgentsAsync), and EventStoreAgents +// .SupportedAgentsAsync fans out from DatabaseDescriptor.TenantIds — which still contains the +// removed tenant for the life of the store instance. So the removed tenant's agent keeps running +// until a process restart (a FRESH store's BuildDatabases() no longer reads the deleted/disabled +// row, so only then does the enumeration shrink and retirement fire). +// - The Wolverine half of the removal transition IS pinned by the unit test +// event_subscription_family_stale_agent_retirement.running_per_tenant_agent_is_stopped_when_its_tenant_is_removed +// (CoreTests/Runtime/Agents/event_subscription_family_cardinality_assignment.cs), which stubs the +// shrunken usage directly. The missing piece is Marten-side (TenantIds shrink on remove/disable), +// not Wolverine-side. +public class runtime_tenant_churn_under_managed_distribution(ITestOutputHelper output) + : PostgresqlContext, IAsyncLifetime +{ + private const string TenantA = "tenant-a"; + private const string TenantB = "tenant-b"; + + private IHost theHost = null!; + private IDocumentStore theStore = null!; + private string theShardConnectionString = null!; + + // Unique names per run: managed partitions/databases from one run otherwise break the next run's + // resource-setup DDL reconciliation. + private readonly string theSuffix = Guid.NewGuid().ToString("N")[..8]; + private string ShardDatabase => $"w486_churn_{theSuffix}"; + private string MasterSchema => $"csp_churn_{theSuffix}"; + + public async Task InitializeAsync() + { + await using (var conn = new NpgsqlConnection(Servers.PostgresConnectionString)) + { + await conn.OpenAsync(); + if (!await conn.DatabaseExists(ShardDatabase)) + { + await new DatabaseSpecification().BuildDatabase(conn, ShardDatabase); + } + } + + theShardConnectionString = new NpgsqlConnectionStringBuilder(Servers.PostgresConnectionString) + { + Database = ShardDatabase + }.ConnectionString; + + void ConfigureShardedStore(StoreOptions m) + { + m.DisableNpgsqlLogging = true; + m.MultiTenantedWithShardedDatabases(x => + { + x.ConnectionString = Servers.PostgresConnectionString; + x.SchemaName = MasterSchema; + x.PartitionSchemaName = MasterSchema + "_tenants"; + x.AddDatabase("shard-1", theShardConnectionString); + }); + + m.Events.StreamIdentity = StreamIdentity.AsString; + m.Events.TenancyStyle = TenancyStyle.Conjoined; + m.Events.AppendMode = EventAppendMode.Quick; + m.Events.UseTenantPartitionedEvents = true; + m.Events.UseIdentityMapForAggregates = false; + + m.Schema.For().MultiTenanted(); + m.Projections.Snapshot(SnapshotLifecycle.Async); + } + + // ONLY tenant-a is provisioned before the host starts. tenant-b arrives at runtime, in the test + // body — that transition is the point of this test. Because the shard database already has a + // tenant, the first enumeration fans out per tenant from the start and no store-global agent + // ever legitimately exists for this shard. + await using (var provisioning = DocumentStore.For(ConfigureShardedStore)) + { + await provisioning.Storage.ApplyAllConfiguredChangesToDatabaseAsync(); + await provisioning.Advanced.AddTenantToShardAsync(TenantA, "shard-1", default); + } + + theHost = await Host.CreateDefaultBuilder() + .UseWolverine(opts => + { + opts.Durability.HealthCheckPollingTime = 1.Seconds(); + opts.Durability.CheckAssignmentPeriod = 1.Seconds(); + + opts.Services.AddMarten(ConfigureShardedStore) + .IntegrateWithWolverine(m => + { + m.UseWolverineManagedEventSubscriptionDistribution = true; + // Sharded tenancy has no single "main" database, so Wolverine's message store needs + // an explicit master (mirrors a real deployment pointing this at a references DB). + m.MainDatabaseConnectionString = Servers.PostgresConnectionString; + }); + + opts.Services.AddSingleton(new OutputLoggerProvider(output)); + opts.Discovery.DisableConventionalDiscovery(); + opts.CodeGeneration.TypeLoadMode = TypeLoadMode.Auto; + }).StartAsync(); + + theStore = theHost.Services.GetRequiredService(); + } + + public async Task DisposeAsync() + { + theHost.GetRuntime().Agents.DisableHealthChecks(); + await theHost.StopAsync(); + theHost.Dispose(); + } + + [Fact] + public async Task adding_a_tenant_at_runtime_fans_out_its_agent_within_the_assignment_window() + { + await theHost.WaitUntilAssumesLeadershipAsync(10.Seconds()); + + // Phase 1: steady state with the pre-provisioned tenant only — exactly ONE per-tenant agent. + await theHost.WaitUntilAssignmentsChangeTo(w => + { + w.AgentScheme = EventSubscriptionAgentFamily.SchemeName; + w.ExpectRunningAgents(theHost, 1); + }, 60.Seconds()); + + var running = RunningSubscriptionUris(); + running.Length.ShouldBe(1); + running[0].ShouldEndWith("/" + TenantA); + + var idA = "pc-" + Guid.NewGuid().ToString("N"); + await AppendAsync(TenantA, idA, 5); + (await WaitForProjectionAsync(TenantA, idA, 5, 60.Seconds()))!.Total.ShouldBe(5); + + // Phase 2: WHILE the host is running, add tenant-b to the same shard database through the + // running store — the same provisioning path the sibling tests use before startup. + await theStore.Advanced.AddTenantToShardAsync(TenantB, "shard-1", default); + + // Within the assignment-evaluation window (CheckAssignmentPeriod = 1s; generous ceiling for CI) + // the leader re-enumerates the usage and starts tenant-b's agent. + await theHost.WaitUntilAssignmentsChangeTo(w => + { + w.AgentScheme = EventSubscriptionAgentFamily.SchemeName; + w.ExpectRunningAgents(theHost, 2); + }, 60.Seconds()); + + running = RunningSubscriptionUris(); + running.Length.ShouldBe(2); + running.ShouldContain(u => u.EndsWith("/" + TenantA, StringComparison.OrdinalIgnoreCase)); + running.ShouldContain(u => u.EndsWith("/" + TenantB, StringComparison.OrdinalIgnoreCase)); + + // The wolverine#3328 transition edge: no store-global (tenant-less) agent may run alongside the + // per-tenant agents — a lingering "…/partitionedcounter/all" agent would track the same events + // as the per-tenant agents and double-process. A tenant-less agent URI ends with the bare "/all" + // shard-key segment; every per-tenant URI carries a trailing tenant segment. + running.ShouldAllBe(u => !u.TrimEnd('/').EndsWith("/all", StringComparison.OrdinalIgnoreCase)); + + // Phase 3: tenant-b's fresh events — its OWN per-tenant sequence, from 1 — actually project… + var idB = "pc-" + Guid.NewGuid().ToString("N"); + await AppendAsync(TenantB, idB, 7); + (await WaitForProjectionAsync(TenantB, idB, 7, 60.Seconds()))!.Total.ShouldBe(7); + + // …and the per-tenant progression row exists in the shard database under the tenant-scoped shard + // identity (`{Projection}:{ShardKey}:{tenant}` — see JasperFx ShardName.Identity and the comment + // in Marten's EventProgressionTable). The progression write commits atomically with the projected + // document above, so no extra polling is needed here. + var progressionB = await TenantProgressionAsync($"PartitionedCounter:All:{TenantB}"); + progressionB.ShouldNotBeNull(); + progressionB.Value.ShouldBeGreaterThanOrEqualTo(7); + + // And tenant-a was undisturbed by the churn. + (await LoadAsync(TenantA, idA))!.Total.ShouldBe(5); + + // Still no store-global agent after processing settled. + RunningSubscriptionUris() + .ShouldAllBe(u => !u.TrimEnd('/').EndsWith("/all", StringComparison.OrdinalIgnoreCase)); + } + + private string[] RunningSubscriptionUris() + { + return theHost.GetRuntime().Agents.AllRunningAgentUris() + .Where(x => x.Scheme == EventSubscriptionAgentFamily.SchemeName) + .Select(x => x.AbsoluteUri.TrimEnd('/')) + .ToArray(); + } + + private async Task AppendAsync(string tenant, string id, int events) + { + await using var session = theStore.LightweightSession(tenant); + session.Events.StartStream(id, + Enumerable.Range(0, events).Select(_ => (object)new CounterChanged(1)).ToArray()); + await session.SaveChangesAsync(); + } + + private async Task LoadAsync(string tenant, string id) + { + await using var session = theStore.QuerySession(tenant); + return await session.LoadAsync(id); + } + + private async Task WaitForProjectionAsync(string tenant, string id, int expectedTotal, + TimeSpan timeout) + { + var deadline = DateTimeOffset.UtcNow + timeout; + PartitionedCounter? doc = null; + while (DateTimeOffset.UtcNow < deadline) + { + doc = await LoadAsync(tenant, id); + if (doc?.Total >= expectedTotal) + { + return doc; + } + + await Task.Delay(250.Milliseconds()); + } + + return doc; + } + + // Events (and therefore mt_event_progression) live in the store's default schema ("public" — the + // sharded configuration's SchemaName only scopes the pool/assignment tables in the MASTER database) + // of the shard database. + private async Task TenantProgressionAsync(string shardIdentity) + { + await using var conn = new NpgsqlConnection(theShardConnectionString); + await conn.OpenAsync(); + var result = await conn + .CreateCommand("select last_seq_id from public.mt_event_progression where name = :name") + .With("name", shardIdentity) + .ExecuteScalarAsync(); + + return result is long value ? value : null; + } +} From bccc2d6ccdc714b104628b047420de506bd8e4b0 Mon Sep 17 00:00:00 2001 From: "Jeremy D. Miller" Date: Mon, 6 Jul 2026 20:01:12 -0500 Subject: [PATCH 3/7] E2E (pinned/skipped): failover preserves per-tenant progression floors (jasperfx#486 WS1) Two nodes, sharded + UseTenantPartitionedEvents, two tenants co-located on one shard database: append N events per tenant, record each tenant's progression floor, kill the node owning the affine agent group, assert the survivor picks both agents up and resumes tenant-b exactly from its floor (doc == N+M, no skips, no rewind) after M more post-failover events. SKIPPED, pinning a Marten 9.13.0-alpha.2 product bug found by this test: the failover mechanics work (survivor restarts both agents, floors never rewind, tenant-a intact), but per-tenant high-water detection FREEZES at the first persisted HighWaterMark: row - HighWaterDetector.loadPerTenantStatistics (src/Marten/Events/Daemon/HighWater/HighWaterDetector.cs:260): var currentMark = lastSeqId > 0 ? lastSeqId : lastValue; Once the row exists, CurrentMark is pinned to it and max(seq_id) is never consulted again, so ANY second batch of events per tenant never projects (not failover-specific; sibling tests all append one batch per tenant and never see it). Empirically confirmed: tenant-b at max seq_id 11 committed, store-global HighWaterMark advanced to 11 (global detection + per-tenant poll trigger both firing), HighWaterMark:tenant-b frozen at 6. Unskip on the Marten-side fix. Co-Authored-By: Claude Fable 5 --- ...preserves_per_tenant_progression_floors.cs | 371 ++++++++++++++++++ 1 file changed, 371 insertions(+) create mode 100644 src/Persistence/MartenTests/Distribution/failover_preserves_per_tenant_progression_floors.cs diff --git a/src/Persistence/MartenTests/Distribution/failover_preserves_per_tenant_progression_floors.cs b/src/Persistence/MartenTests/Distribution/failover_preserves_per_tenant_progression_floors.cs new file mode 100644 index 000000000..563a5adf9 --- /dev/null +++ b/src/Persistence/MartenTests/Distribution/failover_preserves_per_tenant_progression_floors.cs @@ -0,0 +1,371 @@ +using IntegrationTests; +using JasperFx.CodeGeneration; +using JasperFx.Core; +using JasperFx.Events; +using JasperFx.Events.Projections; +using JasperFx.MultiTenancy; +using Marten; +using Marten.Events; +using MartenTests.Distribution.Support; +using Microsoft.Extensions.DependencyInjection; +using Microsoft.Extensions.Hosting; +using Microsoft.Extensions.Logging; +using Npgsql; +using Shouldly; +using Weasel.Core; +using Weasel.Postgresql; +using Weasel.Postgresql.Migrations; +using Wolverine; +using Wolverine.Marten; +using Wolverine.Runtime.Agents; +using Wolverine.Tracking; +using Xunit; +using Xunit.Abstractions; + +namespace MartenTests.Distribution; + +// JasperFx/jasperfx#486 WS1 — node failover must preserve each tenant's progression floor. Under +// MultiTenantedWithShardedDatabases + UseTenantPartitionedEvents two tenants are co-located in one shard +// database, each with its own agent and its own per-tenant progression row +// (`{Projection}:All:{tenant}` in mt_event_progression). The database-affine distribution keeps both +// agents together on ONE node of the two-node cluster; when that node leaves, the surviving node must +// (1) pick both agents up, (2) resume each tenant EXACTLY where its progression row left off — events +// appended after the failover project on top of the pre-failover state with nothing skipped — and +// (3) never rewind a progression row below its pre-failover floor (a rewind would mean reprocessing, +// which the additive Apply() below would surface as an inflated Total). +// +// NOTE: the group of agents may land whole on either node (whichever way the leader balances the +// affine group against its own baseline agents), and it can be rebalanced a beat after the second node +// joins — the test resolves the owner dynamically and only after the placement holds stable across +// several assignment-evaluation periods. +// +// ── SKIPPED: pinned Marten product bug (frozen per-tenant high-water) ───────────────────────────────── +// The failover mechanics themselves work (verified from the run logs + diagnostics of this test before +// it was skipped): the survivor picks up both per-tenant agents, tenant-a's doc and both progression +// floors come through intact, and nothing ever rewinds. What fails is the LAST leg — the M post-failover +// tenant-b events never project, because per-tenant high-water detection FREEZES at the first persisted +// per-tenant high-water row in published Marten 9.13.0-alpha.2: +// +// Marten src/Marten/Events/Daemon/HighWater/HighWaterDetector.cs, loadPerTenantStatistics (:260): +// var currentMark = lastSeqId > 0 ? lastSeqId : lastValue; +// where lastSeqId = the persisted `HighWaterMark:` mt_event_progression row and +// lastValue = max(seq_id) of the tenant's events (per #4847). The FIRST poll after a tenant's first +// events works (no row yet → lastSeqId = 0 → CurrentMark = max(seq_id)) and then persists the row via +// TenantedHighWaterCoordinator.PollAndRouteAsync → MarkHighWaterForTenantAsync. Every subsequent poll +// reads the row back as CurrentMark and re-persists the same value — max(seq_id) is never consulted +// again (the "per-tenant gap detection not yet wired in (a Phase 3 refinement)" comment right above +// that line is the gap). So ANY second batch of events for a tenant stalls forever — not just across +// failover; the sibling tests never see it because they all append exactly one batch per tenant. +// +// Empirically pinned here (diagnostics from the failing run): tenant-b had 11 committed events +// (max seq_id = 11 after a probe append), the store-global HighWaterMark row advanced to 11 (global +// detection + the JasperFx per-tenant poll trigger both fired — JasperFx.Events +// JasperFxAsyncDaemon.OnNext reuses the global high-water cadence), yet HighWaterMark:tenant-b and +// PartitionedCounter:All:tenant-b stayed frozen at 6 = the first persisted mark. +// +// Unskip (and delete this block) once Marten advances per-tenant CurrentMark past the persisted row +// (e.g. max(lastSeqId, lastValue) or real per-tenant gap detection). +public class failover_preserves_per_tenant_progression_floors(ITestOutputHelper output) + : PostgresqlContext, IAsyncLifetime +{ + private const string TenantA = "tenant-a"; + private const string TenantB = "tenant-b"; + + private readonly List _hosts = []; + private string theShardConnectionString = null!; + + // Unique names per run: managed partitions/databases from one run otherwise break the next run's + // resource-setup DDL reconciliation. + private readonly string theSuffix = Guid.NewGuid().ToString("N")[..8]; + private string ShardDatabase => $"w486_fo_{theSuffix}"; + private string MasterSchema => $"csp_fo_{theSuffix}"; + + private void ConfigureShardedStore(StoreOptions m) + { + m.DisableNpgsqlLogging = true; + m.MultiTenantedWithShardedDatabases(x => + { + x.ConnectionString = Servers.PostgresConnectionString; + x.SchemaName = MasterSchema; + x.PartitionSchemaName = MasterSchema + "_tenants"; + x.AddDatabase("shard-1", theShardConnectionString); + }); + + m.Events.StreamIdentity = StreamIdentity.AsString; + m.Events.TenancyStyle = TenancyStyle.Conjoined; + m.Events.AppendMode = EventAppendMode.Quick; + m.Events.UseTenantPartitionedEvents = true; + m.Events.UseIdentityMapForAggregates = false; + + m.Schema.For().MultiTenanted(); + m.Projections.Snapshot(SnapshotLifecycle.Async); + } + + public async Task InitializeAsync() + { + await using (var conn = new NpgsqlConnection(Servers.PostgresConnectionString)) + { + await conn.OpenAsync(); + if (!await conn.DatabaseExists(ShardDatabase)) + { + await new DatabaseSpecification().BuildDatabase(conn, ShardDatabase); + } + } + + theShardConnectionString = new NpgsqlConnectionStringBuilder(Servers.PostgresConnectionString) + { + Database = ShardDatabase + }.ConnectionString; + + // Two co-located tenants on the one shard database, provisioned before any host starts so the + // first agent enumeration already fans out per tenant. + await using (var provisioning = DocumentStore.For(ConfigureShardedStore)) + { + await provisioning.Storage.ApplyAllConfiguredChangesToDatabaseAsync(); + await provisioning.Advanced.AddTenantToShardAsync(TenantA, "shard-1", default); + await provisioning.Advanced.AddTenantToShardAsync(TenantB, "shard-1", default); + } + } + + public async Task DisposeAsync() + { + foreach (var host in _hosts) + { + host.GetRuntime().Agents.DisableHealthChecks(); + await host.StopAsync(); + host.Dispose(); + } + } + + private async Task StartHostAsync() + { + var host = await Host.CreateDefaultBuilder() + .UseWolverine(opts => + { + opts.Durability.HealthCheckPollingTime = 1.Seconds(); + opts.Durability.CheckAssignmentPeriod = 1.Seconds(); + + opts.Services.AddMarten(ConfigureShardedStore) + .IntegrateWithWolverine(m => + { + m.UseWolverineManagedEventSubscriptionDistribution = true; + // Sharded tenancy has no single "main" database, so Wolverine's message store needs + // an explicit master. + m.MainDatabaseConnectionString = Servers.PostgresConnectionString; + }); + + opts.Services.AddSingleton(new OutputLoggerProvider(output)); + opts.Discovery.DisableConventionalDiscovery(); + opts.CodeGeneration.TypeLoadMode = TypeLoadMode.Auto; + }).StartAsync(); + + _hosts.Add(host); + return host; + } + + [Fact(Skip = "Pinned Marten 9.13.0-alpha.2 bug: per-tenant high-water freezes at the first persisted " + + "HighWaterMark: row (HighWaterDetector.loadPerTenantStatistics, " + + "src/Marten/Events/Daemon/HighWater/HighWaterDetector.cs:260 — " + + "`lastSeqId > 0 ? lastSeqId : lastValue`), so any second batch of events per tenant " + + "never projects. See the class comment for the full mechanism; unskip on the Marten fix.")] + public async Task surviving_node_resumes_each_tenant_from_its_progression_floor() + { + const int n = 6; // events per tenant before the failover + const int m = 4; // tenant-b events appended after the node dies + + var first = await StartHostAsync(); + await first.WaitUntilAssumesLeadershipAsync(10.Seconds()); + + var second = await StartHostAsync(); + + // One shard database × two tenants = one affine group of two agents, whole on exactly one node. + var owner = await WaitForSingleOwnerAsync(60.Seconds()); + var survivor = ReferenceEquals(owner, first) ? second : first; + + // Pre-failover: N events per tenant, both projected. + var store = owner.Services.GetRequiredService(); + var idA = "pc-" + Guid.NewGuid().ToString("N"); + var idB = "pc-" + Guid.NewGuid().ToString("N"); + await AppendAsync(store, TenantA, idA, n); + await AppendAsync(store, TenantB, idB, n); + + var survivorStore = survivor.Services.GetRequiredService(); + (await WaitForTotalAsync(survivorStore, TenantA, idA, n, 60.Seconds()))!.Total.ShouldBe(n); + (await WaitForTotalAsync(survivorStore, TenantB, idB, n, 60.Seconds()))!.Total.ShouldBe(n); + + // Record each tenant's progression floor (its own per-tenant sequence, so exactly N here). + var floorA = (await TenantProgressionAsync($"PartitionedCounter:All:{TenantA}")).ShouldNotBeNull(); + var floorB = (await TenantProgressionAsync($"PartitionedCounter:All:{TenantB}")).ShouldNotBeNull(); + floorA.ShouldBeGreaterThanOrEqualTo(n); + floorB.ShouldBeGreaterThanOrEqualTo(n); + + // Kill the node that owns tenant-b's agent (both agents, by affinity). Same stop pattern as + // tenant_partitioned_distribution_multinode.agents_fail_over_to_the_surviving_node_when_a_node_leaves. + owner.GetRuntime().Agents.DisableHealthChecks(); + await owner.StopAsync(); + + // If the owner was the leader (the usual case — see class comment), the survivor must first + // assume leadership before it can evaluate assignments. Returns immediately when the survivor + // already leads. + await survivor.WaitUntilAssumesLeadershipAsync(60.Seconds()); + + await survivor.WaitUntilAssignmentsChangeTo(w => + { + w.AgentScheme = EventSubscriptionAgentFamily.SchemeName; + w.ExpectRunningAgents(survivor, 2); + }, 60.Seconds()); + + // Post-failover: M more events for tenant-b, appended through the surviving node's store. + await AppendAsync(survivorStore, TenantB, idB, m); + + // The projected doc must reflect EXACTLY N+M: a skipped event (resuming above the floor) would + // leave Total below, a reprocessed one (progression rewound below the floor) would push it above + // because Apply() is additive. While polling, also assert the progression row NEVER dips below + // the pre-failover floor. + var deadline = DateTimeOffset.UtcNow + 60.Seconds(); + PartitionedCounter? docB = null; + while (DateTimeOffset.UtcNow < deadline) + { + var progression = await TenantProgressionAsync($"PartitionedCounter:All:{TenantB}"); + if (progression.HasValue) + { + progression.Value.ShouldBeGreaterThanOrEqualTo(floorB, + "tenant-b's progression must never rewind below its pre-failover floor"); + } + + docB = await LoadAsync(survivorStore, TenantB, idB); + if (docB?.Total >= n + m) + { + break; + } + + await Task.Delay(250.Milliseconds()); + } + + docB.ShouldNotBeNull(); + docB.Total.ShouldBe(n + m); + + // Progression advanced monotonically past the floor by the M new events. + var finalB = (await TenantProgressionAsync($"PartitionedCounter:All:{TenantB}")).ShouldNotBeNull(); + finalB.ShouldBeGreaterThanOrEqualTo(floorB + m); + + // And the co-located tenant-a came through the failover untouched: same doc, no reprocessing, + // floor intact. + (await LoadAsync(survivorStore, TenantA, idA))!.Total.ShouldBe(n); + var finalA = (await TenantProgressionAsync($"PartitionedCounter:All:{TenantA}")).ShouldNotBeNull(); + finalA.ShouldBeGreaterThanOrEqualTo(floorA); + } + + /// + /// Wait until exactly two event-subscription agents run cluster-wide AND both run on the same node + /// (the database-affine grouping), returning that owning host. The owner must hold STABLY across + /// several consecutive samples spanning multiple assignment evaluations (CheckAssignmentPeriod = 1s): + /// right after the second node joins, the leader may still be running both agents from its + /// single-node phase and rebalance the whole group a beat later — sampling once would race that move + /// and hand back the wrong owner. + /// + private async Task WaitForSingleOwnerAsync(TimeSpan timeout) + { + const int requiredStableSamples = 6; // 6 × 500ms = 3s > 2 assignment evaluation periods + var deadline = DateTimeOffset.UtcNow + timeout; + IHost? candidate = null; + var stableSamples = 0; + + while (DateTimeOffset.UtcNow < deadline) + { + var byHost = _hosts + .Select(host => (Host: host, Agents: host.GetRuntime().Agents.AllRunningAgentUris() + .Where(x => x.Scheme == EventSubscriptionAgentFamily.SchemeName) + .ToArray())) + .Where(x => x.Agents.Length > 0) + .ToArray(); + + if (byHost.Length == 1 && byHost[0].Agents.Length == 2) + { + if (ReferenceEquals(byHost[0].Host, candidate)) + { + if (++stableSamples >= requiredStableSamples) + { + return candidate!; + } + } + else + { + candidate = byHost[0].Host; + stableSamples = 1; + } + } + else + { + candidate = null; + stableSamples = 0; + } + + await Task.Delay(500.Milliseconds()); + } + + throw new TimeoutException( + "The two per-tenant agents never settled together on a single node. Running agents: " + + _hosts.Select(h => h.GetRuntime().Agents.AllRunningAgentUris() + .Where(x => x.Scheme == EventSubscriptionAgentFamily.SchemeName) + .Select(x => x.ToString()) + .Join(", ")) + .Join(" | ")); + } + + private static async Task AppendAsync(IDocumentStore store, string tenant, string id, int events) + { + await using var session = store.LightweightSession(tenant); + var payload = Enumerable.Range(0, events).Select(_ => (object)new CounterChanged(1)).ToArray(); + if (await session.Events.FetchStreamStateAsync(id) == null) + { + session.Events.StartStream(id, payload); + } + else + { + session.Events.Append(id, payload); + } + + await session.SaveChangesAsync(); + } + + private static async Task LoadAsync(IDocumentStore store, string tenant, string id) + { + await using var session = store.QuerySession(tenant); + return await session.LoadAsync(id); + } + + private static async Task WaitForTotalAsync(IDocumentStore store, string tenant, + string id, int expectedTotal, TimeSpan timeout) + { + var deadline = DateTimeOffset.UtcNow + timeout; + PartitionedCounter? doc = null; + while (DateTimeOffset.UtcNow < deadline) + { + doc = await LoadAsync(store, tenant, id); + if (doc?.Total >= expectedTotal) + { + return doc; + } + + await Task.Delay(250.Milliseconds()); + } + + return doc; + } + + // Events (and therefore mt_event_progression) live in the store's default schema ("public" — the + // sharded configuration's SchemaName only scopes the pool/assignment tables in the MASTER database) + // of the shard database. + private async Task TenantProgressionAsync(string shardIdentity) + { + await using var conn = new NpgsqlConnection(theShardConnectionString); + await conn.OpenAsync(); + var result = await conn + .CreateCommand("select last_seq_id from public.mt_event_progression where name = :name") + .With("name", shardIdentity) + .ExecuteScalarAsync(); + + return result is long value ? value : null; + } +} From eccf5d68365a97ef9c04a9fce49667456b10c3da Mon Sep 17 00:00:00 2001 From: "Jeremy D. Miller" Date: Tue, 7 Jul 2026 06:03:23 -0500 Subject: [PATCH 4/7] =?UTF-8?q?Unskip=20the=20failover=20floor-continuity?= =?UTF-8?q?=20test=20=E2=80=94=20fixed=20by=20marten#4869=20in=209.13.0-al?= =?UTF-8?q?pha.3?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit The skip pinned marten#4867: per-tenant high-water detection froze at the first persisted HighWaterMark: row, so the M post-failover tenant-b events never projected. Marten 9.13.0-alpha.3 advances per-tenant CurrentMark past the persisted row (marten#4869), which makes the final leg pass. The class comment keeps a condensed history of the pinned bug. Co-Authored-By: Claude Fable 5 --- ...preserves_per_tenant_progression_floors.cs | 37 +++---------------- 1 file changed, 6 insertions(+), 31 deletions(-) diff --git a/src/Persistence/MartenTests/Distribution/failover_preserves_per_tenant_progression_floors.cs b/src/Persistence/MartenTests/Distribution/failover_preserves_per_tenant_progression_floors.cs index 563a5adf9..6d2057bda 100644 --- a/src/Persistence/MartenTests/Distribution/failover_preserves_per_tenant_progression_floors.cs +++ b/src/Persistence/MartenTests/Distribution/failover_preserves_per_tenant_progression_floors.cs @@ -39,32 +39,11 @@ namespace MartenTests.Distribution; // joins — the test resolves the owner dynamically and only after the placement holds stable across // several assignment-evaluation periods. // -// ── SKIPPED: pinned Marten product bug (frozen per-tenant high-water) ───────────────────────────────── -// The failover mechanics themselves work (verified from the run logs + diagnostics of this test before -// it was skipped): the survivor picks up both per-tenant agents, tenant-a's doc and both progression -// floors come through intact, and nothing ever rewinds. What fails is the LAST leg — the M post-failover -// tenant-b events never project, because per-tenant high-water detection FREEZES at the first persisted -// per-tenant high-water row in published Marten 9.13.0-alpha.2: -// -// Marten src/Marten/Events/Daemon/HighWater/HighWaterDetector.cs, loadPerTenantStatistics (:260): -// var currentMark = lastSeqId > 0 ? lastSeqId : lastValue; -// where lastSeqId = the persisted `HighWaterMark:` mt_event_progression row and -// lastValue = max(seq_id) of the tenant's events (per #4847). The FIRST poll after a tenant's first -// events works (no row yet → lastSeqId = 0 → CurrentMark = max(seq_id)) and then persists the row via -// TenantedHighWaterCoordinator.PollAndRouteAsync → MarkHighWaterForTenantAsync. Every subsequent poll -// reads the row back as CurrentMark and re-persists the same value — max(seq_id) is never consulted -// again (the "per-tenant gap detection not yet wired in (a Phase 3 refinement)" comment right above -// that line is the gap). So ANY second batch of events for a tenant stalls forever — not just across -// failover; the sibling tests never see it because they all append exactly one batch per tenant. -// -// Empirically pinned here (diagnostics from the failing run): tenant-b had 11 committed events -// (max seq_id = 11 after a probe append), the store-global HighWaterMark row advanced to 11 (global -// detection + the JasperFx per-tenant poll trigger both fired — JasperFx.Events -// JasperFxAsyncDaemon.OnNext reuses the global high-water cadence), yet HighWaterMark:tenant-b and -// PartitionedCounter:All:tenant-b stayed frozen at 6 = the first persisted mark. -// -// Unskip (and delete this block) once Marten advances per-tenant CurrentMark past the persisted row -// (e.g. max(lastSeqId, lastValue) or real per-tenant gap detection). +// This test was originally skipped against Marten 9.13.0-alpha.2, which froze per-tenant high-water +// detection at the first persisted `HighWaterMark:` row (marten#4867) — any SECOND batch of +// events per tenant never projected, so the M post-failover tenant-b events below stalled forever. +// Marten 9.13.0-alpha.3 fixes it (marten#4869: per-tenant CurrentMark advances past the persisted row), +// which is what unskipped this test. public class failover_preserves_per_tenant_progression_floors(ITestOutputHelper output) : PostgresqlContext, IAsyncLifetime { @@ -163,11 +142,7 @@ private async Task StartHostAsync() return host; } - [Fact(Skip = "Pinned Marten 9.13.0-alpha.2 bug: per-tenant high-water freezes at the first persisted " + - "HighWaterMark: row (HighWaterDetector.loadPerTenantStatistics, " + - "src/Marten/Events/Daemon/HighWater/HighWaterDetector.cs:260 — " + - "`lastSeqId > 0 ? lastSeqId : lastValue`), so any second batch of events per tenant " + - "never projects. See the class comment for the full mechanism; unskip on the Marten fix.")] + [Fact] public async Task surviving_node_resumes_each_tenant_from_its_progression_floor() { const int n = 6; // events per tenant before the failover From 79098490bd62e79046101eb6de3e06eab4266af1 Mon Sep 17 00:00:00 2001 From: "Jeremy D. Miller" Date: Tue, 7 Jul 2026 06:03:23 -0500 Subject: [PATCH 5/7] Regression test for marten#4862: single-DB tenant-partitioned per-tenant distribution (jasperfx#486 WS1) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Promotes the WS1 single-database fact-finding probe (found the bug 2026-07-06) into the marten#4862 regression test. Against Marten 9.13.0-alpha.2 a single-database store reported DistributesAgentsPerTenant = false, so managed distribution ran ONE store-global agent whose progression floor (driven to the leading tenant's max per-tenant seq) permanently skipped a lagging tenant's later events at per-tenant seq 1..N — the wolverine#3280 failure mode on one database. Marten 9.13.0-alpha.3 (marten#4864) broadens the gate to any tenant-partitioned store, so the startup expectation flips to per-tenant fan-out (2 tenants -> 2 agents) and the lagging-tenant projection assertions now pass unchanged. Co-Authored-By: Claude Fable 5 --- ...ngle_db_tenant_partitioned_distribution.cs | 214 ++++++++++++++++++ 1 file changed, 214 insertions(+) create mode 100644 src/Persistence/MartenTests/Distribution/single_db_tenant_partitioned_distribution.cs diff --git a/src/Persistence/MartenTests/Distribution/single_db_tenant_partitioned_distribution.cs b/src/Persistence/MartenTests/Distribution/single_db_tenant_partitioned_distribution.cs new file mode 100644 index 000000000..48b5fde85 --- /dev/null +++ b/src/Persistence/MartenTests/Distribution/single_db_tenant_partitioned_distribution.cs @@ -0,0 +1,214 @@ +using IntegrationTests; +using JasperFx.CodeGeneration; +using JasperFx.Core; +using JasperFx.Events; +using JasperFx.Events.Projections; +using JasperFx.MultiTenancy; +using Marten; +using Marten.Events; +using MartenTests.Distribution.Support; +using Microsoft.Extensions.DependencyInjection; +using Microsoft.Extensions.Hosting; +using Microsoft.Extensions.Logging; +using Npgsql; +using Shouldly; +using Weasel.Postgresql; +using Wolverine; +using Wolverine.Marten; +using Wolverine.Runtime.Agents; +using Wolverine.Tracking; +using Xunit; +using Xunit.Abstractions; + +namespace MartenTests.Distribution; + +// Regression test for marten#4862: SINGLE-database × Wolverine-managed cell of the distribution +// matrix (JasperFx/jasperfx#486 WS1). +// +// The bug was found by the WS1 fact-finding probe this test grew out of +// (single_db_tenant_partitioned_lagging_tenant_probe, 2026-07-06): under UseTenantPartitionedEvents +// every tenant draws seq_ids from its own mt_events_sequence_{tenant} starting at 1, so sequences +// OVERLAP inside the one database. Tenant A appends many events first (a store-global agent's +// progression floor advances to A's max), then tenant B appends its first few events at per-tenant +// seq 1..N — far BELOW that floor. Against Marten 9.13.0-alpha.2, where +// IEventStore.DistributesAgentsPerTenant required ShardedTenancy, Wolverine started only the single +// store-global "PartitionedCounter:All" agent, and the lagging tenant's events fell below the shared +// floor and were skipped forever — the wolverine#3280 failure mode, on a single database. +// +// Marten 9.13.0-alpha.3 fixes this (marten#4864: DistributesAgentsPerTenant broadened to ANY +// tenant-partitioned store, and single-DB store descriptors carry their TenantIds), so managed +// distribution now fans out one agent per tenant on a single database, each with its own per-tenant +// progression row — the lag shape is representable and nothing is skipped. This test pins BOTH the +// per-tenant fan-out (2 tenants -> 2 agents) and the lagging-tenant projection. +public class single_db_tenant_partitioned_distribution(ITestOutputHelper output) + : PostgresqlContext, IAsyncLifetime +{ + // Letters/digits/underscores only — the tenant id doubles as the partition suffix. + private const string TenantA = "tenant1"; + private const string TenantB = "tenant2"; + + private IHost theHost = null!; + private IDocumentStore theStore = null!; + + // Unique schema per run: managed partitions created in one run otherwise break the next run's + // resource-setup DDL reconciliation against the same schema. + private readonly string theSchema = "csp_lag_" + Guid.NewGuid().ToString("N")[..8]; + + public async Task InitializeAsync() + { + theHost = await Host.CreateDefaultBuilder() + .UseWolverine(opts => + { + opts.Durability.HealthCheckPollingTime = 1.Seconds(); + opts.Durability.CheckAssignmentPeriod = 1.Seconds(); + + opts.Services.AddMarten(m => + { + m.DisableNpgsqlLogging = true; + m.Connection(Servers.PostgresConnectionString); + m.DatabaseSchemaName = theSchema; + + m.Events.StreamIdentity = StreamIdentity.AsString; + m.Events.TenancyStyle = TenancyStyle.Conjoined; + m.Events.AppendMode = EventAppendMode.Quick; + m.Events.UseTenantPartitionedEvents = true; + m.Events.UseIdentityMapForAggregates = false; + + m.Schema.For().MultiTenanted(); + m.Projections.Snapshot(SnapshotLifecycle.Async); + }) + .IntegrateWithWolverine(m => m.UseWolverineManagedEventSubscriptionDistribution = true); + + opts.Services.AddSingleton(new OutputLoggerProvider(output)); + opts.Discovery.DisableConventionalDiscovery(); + opts.CodeGeneration.TypeLoadMode = TypeLoadMode.Auto; + }).StartAsync(); + + theStore = theHost.Services.GetRequiredService(); + await theStore.Advanced.AddMartenManagedTenantsAsync(default, new Dictionary + { + [TenantA] = TenantA, + [TenantB] = TenantB + }); + } + + public async Task DisposeAsync() + { + theHost.GetRuntime().Agents.DisableHealthChecks(); + await theHost.StopAsync(); + theHost.Dispose(); + } + + [Fact] + public async Task lagging_tenant_events_are_projected_under_managed_distribution() + { + await theHost.WaitUntilAssumesLeadershipAsync(10.Seconds()); + + // marten#4864 (Marten 9.13.0-alpha.3): a single-database tenant-partitioned store reports + // DistributesAgentsPerTenant = true, so managed distribution fans out one agent PER TENANT — + // two tenants => two agents, not one store-global agent. + await theHost.WaitUntilAssignmentsChangeTo(w => + { + w.AgentScheme = EventSubscriptionAgentFamily.SchemeName; + w.ExpectRunningAgents(theHost, 2); + }, 60.Seconds()); + + output.WriteLine("Agent URIs after startup:"); + foreach (var uri in await GetAgentUrisAsync(theHost)) + { + output.WriteLine(" " + uri); + } + + // Phase 1 — the LEADING tenant: 20 events over several commits, per-tenant seq 1..20. + var id = "pc-" + Guid.NewGuid().ToString("N"); + await StartStreamAsync(TenantA, id); + for (var i = 0; i < 4; i++) + { + await AppendAsync(TenantA, id, 5); + } + + // Wait until the agent has demonstrably processed ALL of tenant A's events — the store-global + // progression floor is now at/above A's max seq (20). + var docA = await WaitForProjectionAsync(TenantA, id, x => x.Total >= 21, 60.Seconds()); + docA.ShouldNotBeNull("The leading tenant's projection never caught up — the store-global agent is not processing at all"); + docA.Total.ShouldBe(21); + + await DumpProgressionAsync("after tenant A caught up"); + + // Phase 2 — the LAGGING tenant appends its FIRST events now, at per-tenant seq 1..3: far below + // the store-global floor of ~21. + await StartStreamAsync(TenantB, id); + await AppendAsync(TenantB, id, 2); + + // Observe: does the lagging tenant's projection ever appear? + var docB = await WaitForProjectionAsync(TenantB, id, x => x.Total >= 3, 45.Seconds()); + + await DumpProgressionAsync("after waiting for tenant B"); + output.WriteLine($"lagging tenant '{TenantB}' projection doc: {(docB == null ? "MISSING" : $"Total={docB.Total}")}"); + + docB.ShouldNotBeNull( + "wolverine#3280 on a single database: the lagging tenant's events (per-tenant seq 1..3) fell " + + "below the store-global progression floor and were skipped by the single store-global agent"); + docB.Total.ShouldBe(3); + } + + private async Task StartStreamAsync(string tenant, string id) + { + await using var session = theStore.LightweightSession(tenant); + session.Events.StartStream(id, new CounterChanged(1)); + await session.SaveChangesAsync(); + } + + private async Task AppendAsync(string tenant, string id, int count) + { + await using var session = theStore.LightweightSession(tenant); + for (var i = 0; i < count; i++) + { + session.Events.Append(id, new CounterChanged(1)); + } + + await session.SaveChangesAsync(); + } + + private async Task WaitForProjectionAsync( + string tenant, string id, Func predicate, TimeSpan timeout) + { + var deadline = DateTimeOffset.UtcNow + timeout; + PartitionedCounter? last = null; + while (DateTimeOffset.UtcNow < deadline) + { + await using var session = theStore.QuerySession(tenant); + last = await session.LoadAsync(id); + if (last != null && predicate(last)) + { + return last; + } + + await Task.Delay(250.Milliseconds()); + } + + return last; + } + + private async Task DumpProgressionAsync(string label) + { + output.WriteLine($"mt_event_progression ({label}):"); + await using var conn = new NpgsqlConnection(Servers.PostgresConnectionString); + await conn.OpenAsync(); + await using var reader = await conn + .CreateCommand($"select name, last_seq_id from {theSchema}.mt_event_progression order by name") + .ExecuteReaderAsync(); + while (await reader.ReadAsync()) + { + output.WriteLine($" {reader.GetString(0)} => {reader.GetInt64(1)}"); + } + } + + private static async Task GetAgentUrisAsync(IHost host) + { + var family = host.Services.GetServices() + .OfType().Single(); + var agents = await family.AllKnownAgentsAsync(); + return [.. agents.Select(x => x.AbsoluteUri)]; + } +} From adbb1d2f5ab15436acc93a1130ef03813c8296d2 Mon Sep 17 00:00:00 2001 From: "Jeremy D. Miller" Date: Tue, 7 Jul 2026 06:33:05 -0500 Subject: [PATCH 6/7] Flip single-DB distribution granularity tests to per-tenant fan-out (marten#4862/#4864) With Marten 9.13.0-alpha.3, DistributesAgentsPerTenant is true for ANY tenant-partitioned store, so one async projection over a single tenant-partitioned database now fans out one subscription agent per managed tenant instead of a single store-global "/all" agent whose lone progression floor could silently skip a lagging tenant. - tenant_partitioned_distribution_granularity: 1 agent -> 3 (one per registered tenant, default tenant included), agent URIs carry the tenant in the shard path - tenant_partitioned_distribution_multinode: 2 store-global agents -> 6 per-tenant agents (2 projections x 3 tenants), spread 3+3 across the cluster and all reassigned to the survivor on node loss - AssignmentWaiter.WritePersistedActualsAsync no longer throws KeyNotFoundException when the database holds a node record the waiter isn't tracking (ghost rows from a hard-stopped prior run); it reports the untracked node instead so the TimeoutException diagnostic survives Part of jasperfx#486 WS1. Co-Authored-By: Claude Fable 5 --- ...ngle_db_tenant_partitioned_distribution.cs | 28 +++++++++-- ...nt_partitioned_distribution_granularity.cs | 48 ++++++++++-------- ...nant_partitioned_distribution_multinode.cs | 50 ++++++++++++------- src/Wolverine/TestingExtensions.cs | 28 +++++++---- 4 files changed, 101 insertions(+), 53 deletions(-) diff --git a/src/Persistence/MartenTests/Distribution/single_db_tenant_partitioned_distribution.cs b/src/Persistence/MartenTests/Distribution/single_db_tenant_partitioned_distribution.cs index 48b5fde85..0228ecdb9 100644 --- a/src/Persistence/MartenTests/Distribution/single_db_tenant_partitioned_distribution.cs +++ b/src/Persistence/MartenTests/Distribution/single_db_tenant_partitioned_distribution.cs @@ -39,7 +39,9 @@ namespace MartenTests.Distribution; // tenant-partitioned store, and single-DB store descriptors carry their TenantIds), so managed // distribution now fans out one agent per tenant on a single database, each with its own per-tenant // progression row — the lag shape is representable and nothing is skipped. This test pins BOTH the -// per-tenant fan-out (2 tenants -> 2 agents) and the lagging-tenant projection. +// per-tenant fan-out (2 tenants -> 2 agents) and the lagging-tenant projection. (One residual +// upstream caveat: the per-tenant high-water POLL still rides the global high-water cadence — +// jasperfx#492 — so Phase 2 nudges the global mark to trigger it; see the comment there.) public class single_db_tenant_partitioned_distribution(ITestOutputHelper output) : PostgresqlContext, IAsyncLifetime { @@ -140,8 +142,28 @@ await theHost.WaitUntilAssignmentsChangeTo(w => await StartStreamAsync(TenantB, id); await AppendAsync(TenantB, id, 2); - // Observe: does the lagging tenant's projection ever appear? - var docB = await WaitForProjectionAsync(TenantB, id, x => x.Total >= 3, 45.Seconds()); + // KNOWN CADENCE GAP (jasperfx#492 — same mitigation family as Marten's own + // dynamic_tenant_lifecycle_during_continuous_daemon test): the vectorized per-tenant + // high-water poll rides the GLOBAL high-water cadence, and tenant2's events at per-tenant + // seq 1..3 do not move the store-global max(seq_id) (still 21 from tenant1), so nothing + // triggers the poll that would advance HighWaterMark:tenant2 — the agent idles at ceiling 0. + // Mitigation: NUDGE the global mark forward with tenant1 events while waiting; each nudge + // (global seq 22, 23, ...) fires the per-tenant poll, which reads tenant2's OWN max(seq_id) + // and routes its lagging events. The lag shape this test pins is untouched: under the OLD + // single store-global agent (pre-marten#4864) the nudges cannot help — the store-global floor + // only climbs FURTHER above tenant2's seq 1..3 — so the regression discrimination stands. + // Drop the nudging once jasperfx#492 puts per-tenant polls on the timer instead. + PartitionedCounter? docB = null; + var deadline = DateTimeOffset.UtcNow + 45.Seconds(); + while (DateTimeOffset.UtcNow < deadline) + { + await AppendAsync(TenantA, id, 1); // move the global mark -> trigger the per-tenant poll + docB = await WaitForProjectionAsync(TenantB, id, x => x.Total >= 3, 3.Seconds()); + if (docB?.Total >= 3) + { + break; + } + } await DumpProgressionAsync("after waiting for tenant B"); output.WriteLine($"lagging tenant '{TenantB}' projection doc: {(docB == null ? "MISSING" : $"Total={docB.Total}")}"); diff --git a/src/Persistence/MartenTests/Distribution/tenant_partitioned_distribution_granularity.cs b/src/Persistence/MartenTests/Distribution/tenant_partitioned_distribution_granularity.cs index a40388a4a..47f72bf59 100644 --- a/src/Persistence/MartenTests/Distribution/tenant_partitioned_distribution_granularity.cs +++ b/src/Persistence/MartenTests/Distribution/tenant_partitioned_distribution_granularity.cs @@ -26,12 +26,15 @@ namespace MartenTests.Distribution; // Phase 2 of the per-tenant-partitioned-events matrix (#3021): projection/subscription distribution // granularity under a SINGLE database with Conjoined + Quick + UseTenantPartitionedEvents. // -// The load-bearing caveat this pins: Wolverine-managed event-subscription agents are distributed per -// (shard × DATABASE), NOT per tenant. With per-tenant *partitioning* the tenants all live in one -// database, so a single async projection yields exactly ONE agent ("…/all") spanning every tenant -// partition — not one agent per tenant. ("One agent per tenant" is only true under sharded databases, -// which is Phase 3.) The single shard still processes every tenant's events into that tenant's own -// projection document. +// FLIPPED for Marten 9.13.0-alpha.3 (marten#4862 / marten#4864): this test originally pinned the old +// behavior — agents distributed per (shard × DATABASE), one store-global "…/all" agent spanning every +// tenant partition, with "one agent per tenant" true only under sharded databases. That store-global +// agent's single progression floor silently skipped a lagging tenant's low per-tenant seq_ids (see +// single_db_tenant_partitioned_distribution for the regression shape). marten#4864 broadens +// IEventStore.DistributesAgentsPerTenant to ANY tenant-partitioned store and surfaces the managed +// tenants on the single-DB usage descriptor's TenantIds, so one async projection now fans out one +// agent PER TENANT ("…/all/{tenant}") even on a single database, each tracking its own per-tenant +// progression. public class tenant_partitioned_distribution_granularity(ITestOutputHelper output) : PostgresqlContext, IAsyncLifetime { private IHost theHost = null!; @@ -88,36 +91,41 @@ public async Task DisposeAsync() } [Fact] - public async Task one_async_projection_yields_one_agent_per_shard_not_per_tenant() + public async Task one_async_projection_fans_out_one_agent_per_tenant() { await theHost.WaitUntilAssumesLeadershipAsync(10.Seconds()); - // One async projection over a single database => exactly ONE subscription agent, even though - // three tenants (tenant1/tenant2/*DEFAULT*) are registered. Per-tenant would be three. + // marten#4862/#4864 (Marten 9.13.0-alpha.3): one async projection over a single + // tenant-partitioned database fans out per tenant — three registered tenants + // (tenant1/tenant2/*DEFAULT*) => exactly THREE subscription agents. This used to be ONE + // store-global agent with no tenant component. await theHost.WaitUntilAssignmentsChangeTo(w => { w.AgentScheme = EventSubscriptionAgentFamily.SchemeName; - w.ExpectRunningAgents(theHost, 1); + w.ExpectRunningAgents(theHost, 3); }, 60.Seconds()); var uris = await GetAgentUrisAsync(theHost); - uris.Length.ShouldBe(1); - - // The agent Uri is keyed by (store/database/shard) — no tenant component. - var uri = uris.Single(); - uri.ShouldContain("/all"); - uri.ShouldNotContain("tenant1"); - uri.ShouldNotContain("tenant2"); + uris.Length.ShouldBe(3); + + // Every agent Uri now carries the tenant in the shard path: "…/all/{tenant}". + uris.ShouldAllBe(u => u.Contains("/all")); + uris.ShouldContain(u => u.TrimEnd('/').EndsWith("/tenant1", StringComparison.OrdinalIgnoreCase)); + uris.ShouldContain(u => u.TrimEnd('/').EndsWith("/tenant2", StringComparison.OrdinalIgnoreCase)); + // ...and the default tenant gets its own agent too (it is a registered managed tenant). + uris.ShouldContain(u => !u.TrimEnd('/').EndsWith("/tenant1", StringComparison.OrdinalIgnoreCase) + && !u.TrimEnd('/').EndsWith("/tenant2", StringComparison.OrdinalIgnoreCase)); } [Fact] - public async Task the_single_shard_projects_every_tenant_partition() + public async Task every_tenant_partition_projects_through_its_own_agent() { await theHost.WaitUntilAssumesLeadershipAsync(10.Seconds()); + // marten#4864: per-tenant fan-out — three registered tenants => three agents. await theHost.WaitUntilAssignmentsChangeTo(w => { w.AgentScheme = EventSubscriptionAgentFamily.SchemeName; - w.ExpectRunningAgents(theHost, 1); + w.ExpectRunningAgents(theHost, 3); }, 60.Seconds()); var id = "pc-" + Guid.NewGuid().ToString("N"); @@ -126,7 +134,7 @@ await theHost.WaitUntilAssignmentsChangeTo(w => await theStore.WaitForNonStaleProjectionDataAsync(60.Seconds()); - // The one shard processed both tenant partitions into each tenant's own projection doc. + // Each tenant's own agent processed its partition into that tenant's own projection doc. (await LoadAsync("tenant1", id))!.Total.ShouldBe(5); (await LoadAsync("tenant2", id))!.Total.ShouldBe(11); } diff --git a/src/Persistence/MartenTests/Distribution/tenant_partitioned_distribution_multinode.cs b/src/Persistence/MartenTests/Distribution/tenant_partitioned_distribution_multinode.cs index 02ce3f683..9500acb55 100644 --- a/src/Persistence/MartenTests/Distribution/tenant_partitioned_distribution_multinode.cs +++ b/src/Persistence/MartenTests/Distribution/tenant_partitioned_distribution_multinode.cs @@ -25,11 +25,16 @@ namespace MartenTests.Distribution; -// Phase 2 of #3021 (multi-node slice): the clustered case of the granularity caveat. Two async +// Phase 2 of #3021 (multi-node slice): the clustered case of the distribution granularity. Two async // projections on a SINGLE database with Conjoined + Quick + UseTenantPartitionedEvents and three -// managed tenants. Across a two-node cluster the subscription agents distribute per (shard × DATABASE) -// — two agents total, one per node — NOT per tenant (which would be 2 × 3 = 6). Confirms Wolverine- -// managed distribution behaves identically under per-tenant partitioning as for a non-partitioned store. +// managed tenants. +// +// FLIPPED for Marten 9.13.0-alpha.3 (marten#4862 / marten#4864): this test originally pinned the old +// per-(shard × DATABASE) granularity — two agents total, one per node, never per tenant. marten#4864 +// broadens DistributesAgentsPerTenant to any tenant-partitioned store (the single store-global agent +// skipped lagging tenants — see single_db_tenant_partitioned_distribution), so the cluster now runs +// 2 projections × 3 tenants = SIX per-tenant agents, spread evenly (single-database stores keep the +// even blue/green distribution, not database-affine grouping) — three per node on a two-node cluster. // // This was previously blocked by GH-3037 (a 2nd node's resource-setup threw 42P07 "relation // mt_streams_default already exists"), root-caused to Weasel's name-based default-partition @@ -106,30 +111,36 @@ private async Task StartHostAsync() } [Fact] - public async Task agents_spread_one_per_node_across_the_cluster_and_stay_per_shard() + public async Task agents_spread_evenly_across_the_cluster_per_tenant() { await theOriginalHost.WaitUntilAssumesLeadershipAsync(10.Seconds()); - // Single node: both shards run here (2 agents, not 2 × 3 tenants). + // marten#4862/#4864 (Marten 9.13.0-alpha.3): per-tenant fan-out on a single tenant-partitioned + // database. Single node: 2 projections × 3 tenants = SIX agents (previously 2 store-global). await theOriginalHost.WaitUntilAssignmentsChangeTo(w => { w.AgentScheme = EventSubscriptionAgentFamily.SchemeName; - w.ExpectRunningAgents(theOriginalHost, 2); + w.ExpectRunningAgents(theOriginalHost, 6); }, 60.Seconds()); - // Add a second node — the two shards rebalance one-per-node. + // Add a second node — single-database stores keep the even blue/green distribution + // (database-affine grouping only applies to multi-database stores), so the six agents + // rebalance three-per-node. var second = await StartHostAsync(); await theOriginalHost.WaitUntilAssignmentsChangeTo(w => { w.AgentScheme = EventSubscriptionAgentFamily.SchemeName; - w.ExpectRunningAgents(theOriginalHost, 1); - w.ExpectRunningAgents(second, 1); + w.ExpectRunningAgents(theOriginalHost, 3); + w.ExpectRunningAgents(second, 3); }, 60.Seconds()); - // Exactly two agents cluster-wide, keyed by (store/database/shard) — no tenant component. + // Exactly six agents cluster-wide, every one carrying its tenant in the shard path — one per + // (projection × tenant): two per non-default tenant (NodeCounterA + NodeCounterB), plus the + // default tenant's pair. var uris = await GetAgentUrisAsync(theOriginalHost); - uris.Length.ShouldBe(2); - uris.ShouldAllBe(u => !u.Contains("tenant1") && !u.Contains("tenant2")); + uris.Length.ShouldBe(6); + uris.Count(u => u.TrimEnd('/').EndsWith("/tenant1", StringComparison.OrdinalIgnoreCase)).ShouldBe(2); + uris.Count(u => u.TrimEnd('/').EndsWith("/tenant2", StringComparison.OrdinalIgnoreCase)).ShouldBe(2); } [Fact] @@ -138,25 +149,26 @@ public async Task agents_fail_over_to_the_surviving_node_when_a_node_leaves() await theOriginalHost.WaitUntilAssumesLeadershipAsync(10.Seconds()); var second = await StartHostAsync(); + // marten#4864 per-tenant fan-out: 2 projections × 3 tenants = 6 agents, spread 3 + 3. await theOriginalHost.WaitUntilAssignmentsChangeTo(w => { w.AgentScheme = EventSubscriptionAgentFamily.SchemeName; - w.ExpectRunningAgents(theOriginalHost, 1); - w.ExpectRunningAgents(second, 1); + w.ExpectRunningAgents(theOriginalHost, 3); + w.ExpectRunningAgents(second, 3); }, 60.Seconds()); - // The second node leaves the cluster — its subscription agent must reassign to the survivor, - // and stay per-shard (still two agents total, both now on the original node). + // The second node leaves the cluster — its subscription agents must reassign to the survivor + // (all six per-tenant agents now on the original node). second.GetRuntime().Agents.DisableHealthChecks(); await second.StopAsync(); await theOriginalHost.WaitUntilAssignmentsChangeTo(w => { w.AgentScheme = EventSubscriptionAgentFamily.SchemeName; - w.ExpectRunningAgents(theOriginalHost, 2); + w.ExpectRunningAgents(theOriginalHost, 6); }, 60.Seconds()); - (await GetAgentUrisAsync(theOriginalHost)).Length.ShouldBe(2); + (await GetAgentUrisAsync(theOriginalHost)).Length.ShouldBe(6); } private static async Task GetAgentUrisAsync(IHost host) diff --git a/src/Wolverine/TestingExtensions.cs b/src/Wolverine/TestingExtensions.cs index b5a88ea40..7f3a07f65 100644 --- a/src/Wolverine/TestingExtensions.cs +++ b/src/Wolverine/TestingExtensions.cs @@ -239,10 +239,10 @@ public async ValueTask StartAsync(TimeSpan timeout) } catch (TaskCanceledException) { - } - - if (HasReached()) return true; - + } + + if (HasReached()) return true; + throw new TimeoutException("Did not assume the leadership in the time allowed"); } @@ -291,12 +291,12 @@ public async ValueTask StartAsync(TimeSpan timeout) } catch (TaskCanceledException) { - } - - if (HasReached()) return true; - - var builder = await WritePersistedActualsAsync(); - + } + + if (HasReached()) return true; + + var builder = await WritePersistedActualsAsync(); + throw new TimeoutException(builder.ToString()); } @@ -323,7 +323,13 @@ private async Task WritePersistedActualsAsync() foreach (var node in nodes.OrderBy(x => x.AssignedNodeNumber)) { await writer.WriteLineAsync($"Node {node.AssignedNodeNumber} is running:"); - var runtime = _runtimes[node.NodeId]; + if (!_runtimes.TryGetValue(node.NodeId, out var runtime)) + { + // A persisted node this waiter isn't tracking — typically a ghost record + // left behind by a previous test run's hard-stopped host + await writer.WriteLineAsync($"(node {node.NodeId} is not tracked by this waiter)"); + continue; + } foreach (var uri in runtime.AllRunningAgentUris().Select(x => x.ToString()).OrderBy(x => x)) { From 0ca54625346380133995151a77c9ccc15adee3f4 Mon Sep 17 00:00:00 2001 From: "Jeremy D. Miller" Date: Tue, 7 Jul 2026 07:07:41 -0500 Subject: [PATCH 7/7] Ancillary stores: dispatch message-store shape on database cardinality, not tenancy type name MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Marten 9.13.0-alpha.3 (marten#4864) gives a single-database store with UseTenantPartitionedEvents a TenantPartitionedSingleDatabaseTenancy — a DefaultTenancy subclass, chosen precisely so 'is DefaultTenancy' dispatch keeps working. The ancillary IntegrateWithWolverine seam compared Tenancy.GetType().Name == "DefaultTenancy" instead, so such a store was misrouted to the multi-tenanted message-database path and host startup failed with "no configured connectivity for the required master PostgreSQL message database" after 20 broker retries (the CIMarten failure on this PR). Both the Marten and Polecat ancillary seams now dispatch on Tenancy.Cardinality == DatabaseCardinality.Single, mirroring the main-store integration (WolverineOptionsMartenExtensions), which is why only ancillary stores were affected. Co-Authored-By: Claude Fable 5 --- .../AncillaryWolverineOptionsMartenExtensions.cs | 8 ++++++-- .../AncillaryWolverineOptionsPolecatExtensions.cs | 6 +++++- 2 files changed, 11 insertions(+), 3 deletions(-) diff --git a/src/Persistence/Wolverine.Marten/AncillaryWolverineOptionsMartenExtensions.cs b/src/Persistence/Wolverine.Marten/AncillaryWolverineOptionsMartenExtensions.cs index 27589865b..7cf4a9270 100644 --- a/src/Persistence/Wolverine.Marten/AncillaryWolverineOptionsMartenExtensions.cs +++ b/src/Persistence/Wolverine.Marten/AncillaryWolverineOptionsMartenExtensions.cs @@ -94,8 +94,12 @@ public static MartenServiceCollectionExtensions.MartenStoreExpression Integra integration.SchemaName ??= runtime.Options.Durability.MessageStorageSchemaName ?? store.Options.DatabaseSchemaName; - // TODO -- hacky. Need a way to expose this in Marten - if (store.Tenancy.GetType().Name == "DefaultTenancy") + // Dispatch on database cardinality, mirroring the main-store integration. The old + // exact-type-name check ("DefaultTenancy") misrouted Marten's + // TenantPartitionedSingleDatabaseTenancy (marten#4864, a DefaultTenancy subclass for + // single-database stores with UseTenantPartitionedEvents) to the multi-tenanted + // message-database path, which then failed for lack of a master database. + if (store.Tenancy.Cardinality == JasperFx.Descriptors.DatabaseCardinality.Single) { return BuildSinglePostgresqlMessageStore(integration.SchemaName, integration.AutoCreate, store, runtime, logger); } diff --git a/src/Persistence/Wolverine.Polecat/AncillaryWolverineOptionsPolecatExtensions.cs b/src/Persistence/Wolverine.Polecat/AncillaryWolverineOptionsPolecatExtensions.cs index f861da515..cde4abaad 100644 --- a/src/Persistence/Wolverine.Polecat/AncillaryWolverineOptionsPolecatExtensions.cs +++ b/src/Persistence/Wolverine.Polecat/AncillaryWolverineOptionsPolecatExtensions.cs @@ -77,8 +77,12 @@ public static PolecatStoreExpression IntegrateWithWolverine( runtime.Options.Durability.MessageStorageSchemaName ?? "wolverine"; + // Dispatch on database cardinality rather than the tenancy's exact type name — the + // name check breaks for DefaultTenancy subclasses (see the Marten twin of this seam, + // where marten#4864's TenantPartitionedSingleDatabaseTenancy was misrouted to the + // multi-tenanted path and startup failed for lack of a master database). if (store.Options.Tenancy != null && - store.Options.Tenancy.GetType().Name != "DefaultTenancy") + store.Options.Tenancy.Cardinality != JasperFx.Descriptors.DatabaseCardinality.Single) { return BuildMultiTenantedMessageDatabase(schemaName, integration.AutoCreate, integration.MainConnectionString, store, runtime);