Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
15 commits
Select commit Hold shift + click to select a range
a62e0a5
MicroPlumberd.Migration: offline raw event-store rewrite tool with pa…
rafalmaciag Jul 31, 2026
397fce3
MicroPlumberd.Migration: user-defined-index read-path (deterministic …
rafalmaciag Aug 2, 2026
1b6d354
MicroPlumberd: feasibility spike — index-backed merge streams (Kurren…
rafalmaciag Aug 2, 2026
8e8a8a1
MicroPlumberd: index-backed merge streams — feasibility final + design
rafalmaciag Aug 2, 2026
3340d95
MicroPlumberd: index-backed merge design — v1 FINAL (527 lines)
rafalmaciag Aug 2, 2026
3b38e5c
MicroPlumberd: index-backed merge design — v2 $ce filter caveat-free …
rafalmaciag Aug 2, 2026
fe87609
MicroPlumberd: index-backed merge streams v1 (shape-A, opt-in) — revi…
rafalmaciag Aug 3, 2026
4bf87fc
Merge remote-tracking branch 'origin/master' into feature/mp-index-ba…
rafalmaciag Sep 7, 2026
2dab5c4
Migration.Tests: skip SPIKE-2a — recorded negative finding, not a liv…
rafalmaciag Sep 7, 2026
d2bcdab
feat(migration): generic Transform op + RawEvent identity fields + Mi…
rafalmaciag Sep 7, 2026
c71c334
feat(rewrite): mp-rewrite tool — docker orchestration, swap, rollback…
rafalmaciag Sep 7, 2026
0d02ef6
ci+docs(rewrite): rewrite-e2e.yml, tool pack + single-file release as…
rafalmaciag Sep 7, 2026
5710652
fix(rewrite): review BLOCKERs #1-#3 — pure copy records history, inde…
rafalmaciag Sep 7, 2026
c140508
fix(rewrite): review items #4-#22, M1 — re-guard before stop, source …
rafalmaciag Sep 7, 2026
8e8ecf6
fix(rewrite): gRPC readiness probe; --status anchored on the producti…
rafalmaciag Sep 7, 2026
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
67 changes: 67 additions & 0 deletions .github/workflows/publish-nuget.yml
Original file line number Diff line number Diff line change
Expand Up @@ -25,6 +25,9 @@ env:
jobs:
publish:
runs-on: ubuntu-latest
outputs:
VERSION: ${{ steps.get_version.outputs.VERSION }}
CONFIGURATION: ${{ steps.get_version.outputs.CONFIGURATION }}

steps:
- name: Checkout code
Expand Down Expand Up @@ -187,6 +190,26 @@ jobs:
--output ./artifacts \
/p:PackageVersion=${{ steps.get_version.outputs.VERSION }}

dotnet pack src/MicroPlumberd.Migration/MicroPlumberd.Migration.csproj \
--configuration ${{ steps.get_version.outputs.CONFIGURATION }} \
--no-build \
--output ./artifacts \
/p:PackageVersion=${{ steps.get_version.outputs.VERSION }}

dotnet pack src/MicroPlumberd.Migration.Scripting/MicroPlumberd.Migration.Scripting.csproj \
--configuration ${{ steps.get_version.outputs.CONFIGURATION }} \
--no-build \
--output ./artifacts \
/p:PackageVersion=${{ steps.get_version.outputs.VERSION }}

# mp-rewrite ships as a dotnet tool. The package BUNDLES its dependencies (Jint, the migration
# libraries) under tools/net10.0/any, so it installs without them being on the feed first.
dotnet pack src/MicroPlumberd.Rewrite/MicroPlumberd.Rewrite.csproj \
--configuration ${{ steps.get_version.outputs.CONFIGURATION }} \
--no-build \
--output ./artifacts \
/p:PackageVersion=${{ steps.get_version.outputs.VERSION }}

- name: List artifacts
run: ls -lh ./artifacts

Expand All @@ -211,3 +234,47 @@ jobs:
name: nuget-packages-${{ steps.get_version.outputs.VERSION }}
path: ./artifacts/*.nupkg
retention-days: 30

# Neurons have no SDK, so the tool also ships as one self-contained file per architecture, attached to the
# release for the tag. Runs only for a tag push: a workflow_dispatch has no tag to release against.
rewrite-binaries:
needs: publish
if: startsWith(github.ref, 'refs/tags/')
runs-on: ubuntu-latest
permissions:
contents: write

steps:
- name: Checkout code
uses: actions/checkout@v4

- name: Setup .NET
uses: actions/setup-dotnet@v4
with:
dotnet-version: '10.0.x'

- name: Publish single-file binaries
run: |
for RID in linux-x64 linux-arm64; do
# Always Release: a preview TAG is about the package version, not about shipping an unoptimised
# binary to a neuron that has no SDK to rebuild it with.
dotnet publish src/MicroPlumberd.Rewrite/MicroPlumberd.Rewrite.csproj \
--configuration Release \
--runtime "$RID" \
--self-contained true \
-p:PublishSingleFile=true \
-p:IncludeNativeLibrariesForSelfExtract=true \
-p:EnableCompressionInSingleFile=true \
-p:Version=${{ needs.publish.outputs.VERSION }} \
--output "./binaries/$RID"
# One file is the whole point — an operator scps it onto a neuron and runs it.
mv "./binaries/$RID/mp-rewrite" "./binaries/mp-rewrite-${{ needs.publish.outputs.VERSION }}-$RID"
done
ls -lh ./binaries/mp-rewrite-*

- name: Attach binaries to the release
uses: softprops/action-gh-release@v2
with:
tag_name: ${{ github.ref_name }}
files: ./binaries/mp-rewrite-*
fail_on_unmatched_files: true
103 changes: 103 additions & 0 deletions .github/workflows/rewrite-e2e.yml
Original file line number Diff line number Diff line change
@@ -0,0 +1,103 @@
# End-to-end suite for mp-rewrite (epic-095).
#
# These tests are the acceptance bar for the tool: they start REAL KurrentDB containers, bind-mount a data
# directory, rewrite it, swap it in and read the result back through the original container. They cannot be
# replaced by unit tests, because the whole value of the tool is the orchestration.
name: mp-rewrite end-to-end

on:
pull_request:
branches: [master]
push:
branches: [master, 'feature/**']
workflow_dispatch:

# One run per ref. These scenarios start real containers on a shared docker host and the tool refuses when a
# sibling of the same compose project is up, so two concurrent runs would refuse each other.
concurrency:
group: rewrite-e2e-${{ github.ref }}
cancel-in-progress: true

env:
DOTNET_VERSION: '10.0.x'
DOTNET_NOLOGO: 'true'
DOTNET_CLI_TELEMETRY_OPTOUT: 'true'

jobs:
e2e:
runs-on: ubuntu-latest
# The suite starts two containers per scenario and runs strictly sequentially; ~7 minutes on saturn.
# 30 minutes leaves room for a cold image pull and a slower runner without letting a hang burn an hour.
timeout-minutes: 30

steps:
- name: Checkout
uses: actions/checkout@v4

- name: Setup .NET
uses: actions/setup-dotnet@v4
with:
dotnet-version: ${{ env.DOTNET_VERSION }}

- name: Restore
run: dotnet restore src/MicroPlumberd.Rewrite.Tests/MicroPlumberd.Rewrite.Tests.csproj

- name: Build
run: >-
dotnet build src/MicroPlumberd.Rewrite.Tests/MicroPlumberd.Rewrite.Tests.csproj
--configuration Release --no-restore

# The fixture pulls this itself if it is missing, but doing it here gives the pull its own step, its own
# log and its own failure — rather than surfacing as a container-start timeout inside the first scenario.
- name: Pull KurrentDB image
run: docker pull docker.kurrent.io/kurrent-latest/kurrentdb:latest

# Sequential is not optional: two fixtures up at once are siblings of each other's compose project and
# the tool would correctly refuse. xunit.runner.json pins it too, so a local run behaves the same.
- name: Run end-to-end suite
run: >-
dotnet test src/MicroPlumberd.Rewrite.Tests/MicroPlumberd.Rewrite.Tests.csproj
--configuration Release --no-build
--logger "trx;LogFileName=rewrite-e2e.trx"
--results-directory ${{ github.workspace }}/test-results
-- xUnit.ParallelizeTestCollections=false
Comment thread
rafalmaciag marked this conversation as resolved.

# A discovery break runs ZERO tests and dotnet test exits 0 — green, with nothing to notice it in. For a
# suite that is the acceptance bar for a tool whose failure mode is losing a production store, "no tests
# ran" must be a failure, and the floor must be a number someone has to consciously lower.
- name: Require the suite to have actually run
if: always()
run: |
TRX=$(ls ${{ github.workspace }}/test-results/*.trx 2>/dev/null | head -1)
if [ -z "$TRX" ]; then echo "::error::No TRX produced — the suite did not run."; exit 1; fi
EXECUTED=$(grep -o 'executed="[0-9]*"' "$TRX" | head -1 | grep -o '[0-9]*')
PASSED=$(grep -o 'passed="[0-9]*"' "$TRX" | head -1 | grep -o '[0-9]*')
OUTCOME=$(grep -o 'outcome="[A-Za-z]*"' "$TRX" | tail -1 | sed 's/.*="\(.*\)"/\1/')
echo "executed=$EXECUTED passed=$PASSED outcome=$OUTCOME"
# Measured on saturn 2026-09-07: a run that ABORTED after 71 of 76 tests still printed "Passed!"
# for the 71 that had run. A dying run looks green, so the count is the only thing that catches it —
# and a floor far below the real total would have let that very run through. Raise this
# deliberately when tests are added; never lower it to make a red build green.
MIN_TESTS=83
if [ -z "$EXECUTED" ] || [ "$EXECUTED" -lt "$MIN_TESTS" ]; then
echo "::error::Only ${EXECUTED:-0} test(s) ran; expected at least $MIN_TESTS. A run that dies partway still reports success."
exit 1
fi
if [ "$OUTCOME" != "Completed" ]; then
echo "::error::Test run outcome was '$OUTCOME', not 'Completed'."
exit 1
fi

- name: Upload test results
if: always()
uses: actions/upload-artifact@v4
with:
name: rewrite-e2e-trx
path: ${{ github.workspace }}/test-results/*.trx
retention-days: 14

# A scenario that fails midway can leave a container behind; the runner is thrown away, but leaving them
# named in the log makes a CI-only failure diagnosable. Only ever by the exact name prefix this suite owns.
- name: Report leftover containers
if: always()
run: docker ps -a --filter "name=mp-rewrite-" --format '{{.Names}}\t{{.Status}}' || true
129 changes: 129 additions & 0 deletions dev-log-index-backed-merge.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,129 @@
# Dev-log: Index-backed merge streams (v1, shape-A) — implementation

Branch: `feature/mp-index-backed-merge`. Implements v1 (slices S1–S4) of
`docs/design-index-backed-merge-streams.md` exactly. Shape-B (S5) is out of scope (v2). No git actions performed
by the engineer — reviewed branch only.

Build with `dotnet.exe`. Every slice ends with a green integration test against a REAL KurrentDB 26.1
(`MicroPlumberd.Testing` `EventStoreServer.StartInDocker`, NEVER Testcontainers, no `ClearAllPools`).

## What shipped, by slice

### S1 — index-backed subscription primitive + shared subscription seam (core)
- **Relocated the index primitives into core** as `MicroPlumberd.UserDefinedIndex` and made `KurrentHttpEndpoint`
a public core utility (added `FromSettings(KurrentDBClientSettings)` alongside the existing
`Parse(connectionString)`). `MicroPlumberd.Migration` now has a `ProjectReference` → core and its
`UserDefinedIndexSource` DELEGATES create/name/filter/stream to the core type — ONE implementation
(CLAUDE.md "NEVER DUPLICATE"), dependency arrow Migration → core (design "Assembly layering", team-lead
APPROVED). `UserDefinedIndexSource` keeps its offline-only members (`CountMatchingAsync`,
`WaitUntilReadyAsync` count-gate, `ReadAsync`, `CreateWaitReadAsync`) unchanged.
- **`ISubscriptionState` seam** driving ONE `SubscriptionRunner` loop (no fork): `SubscriptionRunnerState`
(existing stream path via `SubscribeToStream`, resume `FromStream.After(rev)`) and the new
`IndexSubscriptionState` (filtered `SubscribeToAll(FromAll, resolveLinkTos:true,
StreamFilter.Prefix($idx-user-<name>))`, resume `FromAll.After(pos)`). `SubscriptionRunner.WithHandler`
changed only its two source-specific lines (`Subscribe()` / `Advance(e)`); dispatch, `ICaughtUpHandler`,
5s-resubscribe-backoff, `FailFastException` are literally reused.
- **Loud guard** (SPIKE-2a/SPIKE-9 footgun): `SubscriptionRunnerState.Subscribe()` throws
`InvalidOperationException` if the stream name starts with `$idx-user-` (a direct `SubscribeToStream` there
silently delivers zero). Index streams are read ONLY via filtered `$all`.

### S2 — opt-in wiring + PROJECTION-vs-INDEX PARITY (the acceptance bar)
- **`MergeSource` enum** (`MicroPlumberd.Services`, default `Projection`). `AddEventHandler<T>` /
`AddSingletonEventHandler<T>` / `AddScopedEventHandler<T>` gained `MergeSource mergeSource = Projection`;
`EventHandlerStarter<T>.Configure` carries it and `Start()` routes to `SubscribeEventHandlerViaIndex<T>`.
- **`PlumberEngine.SubscribeEventHandlerViaIndex<T>`** (public): convention output stream + `GetEventNamesFor<T>`
→ reconcile (ensure index + return the name) → subscribe via `IndexSubscriptionState` in a `SubscriptionRunner`,
same dispatch/converter as the projection path. Catch-up only. Logs handler → index → stream → types + start.
- **Registration guards (fail fast):** `UserDefinedIndex + persistently` → throw (SPIKE-7, permanent);
`UserDefinedIndex + specific-revision start` → throw (only Start/End meaningful on filtered `$all`).

### S3 — create-new-and-swap reconciler (`UserDefinedIndexReconciler`)
- Filter-hashed name `mpidx-<normalized-output>-<hash8-of-filter>` (`UserDefinedIndex.IndexNameFor`,
`hash8` = first 8 lower-hex of `SHA-256(BuildFilter(types))`, ordinal-sorted ⇒ stable per set). Changed
event-type set ⇒ new name ⇒ read model rebuilds from `FromAll.Start`; unchanged ⇒ same name ⇒ 409 reuse ⇒
no rebuild. Best-effort DELETE orphan cleanup of superseded `mpidx-<output>-*` via `ListNamesAsync` +
`DeleteAsync` (correctness never depends on cleanup running).

### S4 — coexistence + permanent persistent guard + proliferation
- Projection-backed and index-backed handlers coexist in one engine; the projection default creates NO index;
the persistent+index registration guard is permanent; N index-backed merges all build (SPIKE-6 headroom).

## Key decisions / deviations
- **Relocation over copy** (design-mandated): required adding `Migration → MicroPlumberd` `ProjectReference`.
This surfaced a namespace collision — `MicroPlumberd.StreamMetadata` began shadowing
`KurrentDB.Client.StreamMetadata` at an unqualified use site in `ProjectionCopier.cs`; fixed by fully
qualifying that one `new KurrentDB.Client.StreamMetadata(...)`. No behavior change (the 6 Migration
integration tests remain green = the S1-T5 relocation-parity guard).
- **CaughtUp semantics on the index path (empirical nuance, not a defect — DOCUMENTED at the opt-in surface):**
an index-backed subscription's history→live boundary tracks the `$all` position, so while the index is still
backfilling, `CaughtUp` can fire BEFORE the backfilled links arrive as "live". The guaranteed properties are
CaughtUp-fires + no-loss + commit-order — NOT the projection output-stream path's strict "all history
processed before CaughtUp" timeline. S2-T3 asserts the guaranteed properties.
**Consumer guidance (surfaced at the decision point):** because this is a per-handler opt-in, the trade-off is
documented in XML docs on `MergeSource.UserDefinedIndex`, on the `mergeSource:` parameter of
`AddEventHandler`/`AddSingletonEventHandler`/`AddScopedEventHandler`, and on
`PlumberEngine.SubscribeEventHandlerViaIndex<T>` — a read model that uses `ICaughtUpHandler.CaughtUp()` as a
"fully caught up / now authoritative" readiness signal should NOT opt into index-backing (or must tolerate the
weaker guarantee); eventual delivery of all history + commit ordering are still guaranteed.
- **`HttpEndpoint` from settings:** core derives the management base + basic-auth from
`KurrentDBClientSettings.ConnectivitySettings.Address` (or first gossip seed) + `DefaultCredentials` — the
same node the gRPC client talks to; mirrors the existing `WaitUntilReady` pattern.
- Removed the dead `SubscriptionRunnerState.Handler` property (only ever assigned, never read).

## Test results (all against real KurrentDB 26.1, image with `/v2/indexes`)
- **S1:** 4/4 — catch-up→live in commit order; resume from recorded `$all` position; footgun guard throws;
stream-backed path no-regression. Plus **6/6** pre-existing `UserDefinedIndexIntegrationTests` (Migration
relocation is byte-identical — S1-T5).
- **S2:** 4/4 — **PARITY (S2-T1, the acceptance bar): the same read model fed the same interleaved events via
the projection path and the index path reached IDENTICAL final state in the IDENTICAL delivery order**;
live-after-boot; ICaughtUpHandler fires; registration guards reject persistent + revision-start.
- **S3:** 3/3 — definition change → new filter-hashed index with widened content; no-op when unchanged;
superseded index DELETEd (orphan cleanup verified end-to-end, incl. `ListNamesAsync` parse + DELETE).
- **S4:** 4/4 — projection + index coexist; default creates no index; persistent+index rejected (projection
persistent untouched); 6 concurrent index-backed merges all build.
- **Regression:** whole solution builds 0 errors; existing subscription/read-model suites pass. One
timing-sensitive from-End test (`ReadModelTests.SubscribeModelFromEnd`) flaked once under parallel Docker
load but passes 3/3 in isolation — pre-existing flakiness, the refactored path is behavior-identical.

## Review follow-ups closed (post-approval, before v1 lands)
- **Acceptance-path integration test (was compile-checked only).** Added a full-DI test
(`IndexBackedMergeReviewFollowupsTests.DiAcceptancePath`): a consumer registers
`AddSingletonEventHandler<T>(mergeSource: MergeSource.UserDefinedIndex)`, boots a real `TestAppHost`
(`Host.StartAsync` → `EventHandlerService` → `EventHandlerStarter.Start` →
`SubscribeEventHandlerViaIndex<T>(eh: null)`), and the DI-resolved singleton read model folds index-delivered
events to the expected state + tails a live append. Exercises the `eh == null` DI-resolution branch
(`SubscriptionRunner.WithHandler<T>(func)`) + starter routing the other tests bypassed by passing an instance.
- **Cross-output orphan-delete collision guard (data-destructive edge).** Two distinct output streams whose names
NORMALIZE to the same managed base (e.g. `"Foo.1"` and `"Foo-1"` → `"foo-1"`) would share the
`mpidx-<base>-*` prefix, so reconciling one could DELETE the other's live index. Added
`UserDefinedIndex.RegisterManagedBase(outputStream)` — a process-wide owner registry keyed on the normalized
base — called first in `UserDefinedIndexReconciler.ReconcileAsync` (before any DELETE). A collision throws a
clear `InvalidOperationException` at reconcile; re-registering the SAME output stream is idempotent. Unit test
`ManagedBase_collision_is_rejected_same_name_is_idempotent` (no server).
- **Deferred (team-lead logged as follow-ups):** a mid-backfill CaughtUp reorder-window integration test; the
`NormalizeName` doc-charset nit.

## Files
### Added (core `src/MicroPlumberd/`)
- `KurrentHttpEndpoint.cs` — relocated, public, `+FromSettings`.
- `UserDefinedIndex.cs` — relocated primitives + `EnsureAsync`/`DeleteAsync`/`GetFilterAsync`/`ListNamesAsync`/
`IndexNameFor`/`ManagedNamePrefixFor`.
- `ISubscriptionState.cs` — the seam + `IndexSubscriptionState`.
- `UserDefinedIndexReconciler.cs` — `IIndexDefinitionReconciler` + `UserDefinedIndexReconciler`.
### Added (`src/MicroPlumberd.Services/`)
- `MergeSource.cs`.
### Added (tests `src/MicroPlumberd.Tests/Integration/`)
- `IndexBackedMergeS1Tests.cs`, `IndexBackedMergeS2Tests.cs`, `IndexBackedMergeS3Tests.cs`,
`IndexBackedMergeS4Tests.cs`, `IndexBackedMergeReviewFollowupsTests.cs` (DI acceptance path + collision guard).
### Changed
- `src/MicroPlumberd/SubscriptionRunner.cs` — `SubscriptionRunnerState : ISubscriptionState` (+guard, +Advance),
runner ctor `ISubscriptionState`, loop uses `Advance(e)`.
- `src/MicroPlumberd/PlumberEngine.cs` — store settings; lazy `UserDefinedIndex`/reconciler;
`SubscribeEventHandlerViaIndex<T>` + `IsTailOnlyStart`.
- `src/MicroPlumberd.Services/EventHandlerStarter.cs` — `Configure(..., MergeSource)` + `Start` routing.
- `src/MicroPlumberd.Services/ContainerExtensions.cs` — `mergeSource` params + `ValidateMergeSource` guards.
- `src/MicroPlumberd.Migration/MicroPlumberd.Migration.csproj` — `ProjectReference` → core.
- `src/MicroPlumberd.Migration/UserDefinedIndexSource.cs` — delegates create/name/filter to core.
- `src/MicroPlumberd.Migration/ProjectionCopier.cs` — qualify `KurrentDB.Client.StreamMetadata`.
### Removed
- `src/MicroPlumberd.Migration/KurrentHttpEndpoint.cs` — relocated to core.
Loading
Loading