Repository navigation
epic-095 iter1: mp-rewrite — KurrentDB rewrite tool (+ MicroPlumberd.Migration, Scripting) #17
New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
Merged
Merged
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 397fce3
MicroPlumberd.Migration: user-defined-index read-path (deterministic …
rafalmaciag 1b6d354
MicroPlumberd: feasibility spike — index-backed merge streams (Kurren…
rafalmaciag 8e8a8a1
MicroPlumberd: index-backed merge streams — feasibility final + design
rafalmaciag 3340d95
MicroPlumberd: index-backed merge design — v1 FINAL (527 lines)
rafalmaciag 3b38e5c
MicroPlumberd: index-backed merge design — v2 $ce filter caveat-free …
rafalmaciag fe87609
MicroPlumberd: index-backed merge streams v1 (shape-A, opt-in) — revi…
rafalmaciag 4bf87fc
Merge remote-tracking branch 'origin/master' into feature/mp-index-ba…
rafalmaciag 2dab5c4
Migration.Tests: skip SPIKE-2a — recorded negative finding, not a liv…
rafalmaciag d2bcdab
feat(migration): generic Transform op + RawEvent identity fields + Mi…
rafalmaciag c71c334
feat(rewrite): mp-rewrite tool — docker orchestration, swap, rollback…
rafalmaciag 0d02ef6
ci+docs(rewrite): rewrite-e2e.yml, tool pack + single-file release as…
rafalmaciag 5710652
fix(rewrite): review BLOCKERs #1-#3 — pure copy records history, inde…
rafalmaciag c140508
fix(rewrite): review items #4-#22, M1 — re-guard before stop, source …
rafalmaciag 8e8ecf6
fix(rewrite): gRPC readiness probe; --status anchored on the producti…
rafalmaciag File filter
Filter by extension
Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
There are no files selected for viewing
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
| 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 | ||
|
|
||
| # 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 | ||
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
| 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. |
Oops, something went wrong.
Oops, something went wrong.
Add this suggestion to a batch that can be applied as a single commit.
This suggestion is invalid because no changes were made to the code.
Suggestions cannot be applied while the pull request is closed.
Suggestions cannot be applied while viewing a subset of changes.
Only one suggestion per line can be applied in a batch.
Add this suggestion to a batch that can be applied as a single commit.
Applying suggestions on deleted lines is not supported.
You must change the existing code in this line in order to create a valid suggestion.
Outdated suggestions cannot be applied.
This suggestion has been applied or marked resolved.
Suggestions cannot be applied from pending reviews.
Suggestions cannot be applied on multi-line comments.
Suggestions cannot be applied while the pull request is queued to merge.
Suggestion cannot be applied right now. Please check back later.
Uh oh!
There was an error while loading. Please reload this page.