feat(mongodb): implement the MongoDB document persistence platform
Implements the mongodb-superpowers-package design: Stable Tasks 1-50 and Advanced Tasks 1-15. The design assumes 19 Stable + 12 Advanced Gradle projects under modules/mongodb*. This repository's fail-closed registry declares exactly 19 leaf identities, so those modules become package boundaries inside the registered leaf :adapter:outbound:persistence-mongo, with the design's module dependency table enforced by ten ArchUnit rules. The mapping and every deviation are recorded in docs/mongodb/repository-adaptation.md. Contract highlights, all enforced by tests rather than convention: - Transaction body retry and commit retry are separate loops. A new session per body attempt; commit-only retry on an unknown commit. The body is never replayed after a commit ambiguity, so a failover cannot become a duplicate. - MongoExecutionOutcome keeps both ambiguous outcomes distinct from success and failure, and MongoFailureContext records only the design-permitted fields. - Failure classification reads server error labels before numeric codes. - BSON representations come from a pinned manifest, never a library default, and a golden type-signature gate fails on any drift. - Index and validator changes go through the manifest and the admin plane; metadata ownership gates every drop. - Every Advanced capability refuses construction unless its flag is enabled. Verified against real servers, not only unit tests. Running the lanes for the first time exposed four defects that a green `check` had hidden: - Four release lanes passed while executing zero tests; the gate now counts executed tests per lane and fails on zero. - The "single replica set" fixture was a standalone, because Testcontainers 2.x needs withReplicaSet(); its test only asserted a connection string. - The three-node fixture was three independent clusters, so no election could occur, and awaitNewPrimary() compared against the post-stop primary. - The migration lease checked modifiedCount, so a same-millisecond refresh read as a lost lease. scripts/verify-mongodb-platform.sh now reports: 9 lanes, 0 skipped, 0 failed, every evidence category produced. scripts/verify-mongodb-advanced.sh reports NOT PROMOTABLE: actual-topology evidence (real sharded cluster, real KMS, real target deployment) is unobtainable here, so it is named rather than assumed. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
This commit is contained in:
co-authored by
Claude Opus 5
parent
3b5aee50e3
commit
d57d2f62a0
@@ -0,0 +1,105 @@
|
||||
# Change Stream Guide
|
||||
|
||||
Design §20, decision D-12. A change stream is an **at-least-once projector**, not an event bus.
|
||||
|
||||
## 1. What a change stream is not
|
||||
|
||||
D-12 is explicit: a physical change event is not a business integration event. The two differ in
|
||||
ways that matter to every consumer:
|
||||
|
||||
| Change event | Integration event |
|
||||
|---|---|
|
||||
| Emitted per document write | Emitted per business fact |
|
||||
| Shape follows the storage schema | Shape is a published contract |
|
||||
| A refactor of the document changes it | A refactor of the document does not change it |
|
||||
| Replayed on resume, duplicated on retry | Versioned and deliberately evolved |
|
||||
|
||||
Publishing raw change events externally makes your storage schema a public API, and the first time
|
||||
someone renames a field the downstream consumers break. If you need to bridge to messaging, use the
|
||||
Advanced bridge, which maps to an owned envelope
|
||||
([advanced/multi-tenancy.md](advanced/multi-tenancy.md) is separate;
|
||||
the bridge is described in §7 below).
|
||||
|
||||
## 2. Subscription and resume
|
||||
|
||||
`MongoChangeStreamSubscription` declares the collection, pipeline and consistency. `MongoResumePosition`
|
||||
is either a resume token or a cluster time; `MongoResumeCheckpoint` is what gets persisted and
|
||||
`MongoResumeCheckpointStore` persists it.
|
||||
|
||||
The checkpoint stores the token as a Base64 `encodedToken` string rather than a byte array — a record
|
||||
with an array component has broken equality, and a checkpoint that does not compare correctly is a
|
||||
checkpoint that silently fails its own dedup test.
|
||||
|
||||
## 3. Checkpoint after processing, not after receiving
|
||||
|
||||
The ordering rule that makes at-least-once actually hold:
|
||||
|
||||
```
|
||||
receive event
|
||||
→ process it (idempotently)
|
||||
→ persist the checkpoint
|
||||
```
|
||||
|
||||
Checkpointing on receipt turns the delivery guarantee into at-most-once, and the events lost are
|
||||
exactly the ones the process died while handling.
|
||||
|
||||
## 4. Idempotency
|
||||
|
||||
`MongoChangeEventIdentity` is the dedup key: `(resumeToken, documentKey, clusterTime, operationType)`.
|
||||
`MongoChangeDeduplicationStore` records what has been applied. Duplicates are not an edge case — every
|
||||
resume after any interruption replays at least one event, so a projector that is not idempotent is
|
||||
wrong on its first restart, not on some rare day.
|
||||
|
||||
`MongoChangeProjector` returns a `MongoChangeProjectionResult` so the runner can distinguish applied
|
||||
from skipped-as-duplicate, and the skip count is worth a metric: a sudden rise means something is
|
||||
looping.
|
||||
|
||||
## 5. States and recovery
|
||||
|
||||
`MongoChangeStreamState`: `STARTING`, `RUNNING`, `RESUMING`, `STOPPED`, `HISTORY_LOST`.
|
||||
|
||||
`MongoChangeStreamRecoveryPolicy` returns a `MongoChangeStreamRecoveryDecision`, which is either
|
||||
`resume()` (auto-resume from the checkpoint) or `halt(state, runbook)`. A halting decision **must**
|
||||
name a runbook — a decision that only says "stopped" leaves the on-call engineer to work out from
|
||||
scratch whether the projection can be rebuilt and from what.
|
||||
|
||||
| Situation | Decision |
|
||||
|---|---|
|
||||
| Transient network error, token still valid | `resume()` |
|
||||
| Primary failover | `resume()` — the token survives an election |
|
||||
| `invalidate` (collection dropped/renamed) | `halt(STOPPED, …)` → `MongoInvalidateRecovery` |
|
||||
| Token no longer in the oplog | `halt(HISTORY_LOST, "history-lost")` → `MongoChangeHistoryLostException` |
|
||||
|
||||
## 6. History lost
|
||||
|
||||
`MongoChangeHistoryLostException` is raised when the resume token predates the oldest oplog entry.
|
||||
The stream **cannot** be resumed: the events between the checkpoint and now are gone from the server,
|
||||
and no amount of retrying brings them back.
|
||||
|
||||
What the platform will not do is silently restart from "now". That looks like a recovery and is
|
||||
actually a silent gap in the projection — the worst possible outcome, because nothing reports it. The
|
||||
runner halts and requires an operator decision. See
|
||||
[runbooks/history-lost.md](runbooks/history-lost.md).
|
||||
|
||||
## 7. Bridging to messaging (Advanced)
|
||||
|
||||
`MongoChangeMessagingBridge` is opt-in behind `MongoCapability.CHANGE_STREAM` plus the bridge's own
|
||||
flag. It maps a change event to a platform-owned `MongoIntegrationEventEnvelope` through
|
||||
`MongoChangeToIntegrationEventMapper` and hands it to a `MongoIntegrationEventPublisher` port.
|
||||
|
||||
The port is defined in the bridge package rather than imported from the messaging adapter because the
|
||||
architecture registry forbids adapter-to-adapter dependencies; the composition root supplies the
|
||||
implementation.
|
||||
|
||||
`MongoBridgeOutboxPolicy` and `MongoBridgeCheckpointPolicy` state the delivery contract: publish then
|
||||
checkpoint, at-least-once, consumers must dedup on the envelope's event id.
|
||||
|
||||
## 8. Operating notes
|
||||
|
||||
- Change streams require a replica set. `MongoStartupValidator` refuses a change-stream profile on
|
||||
`STANDALONE`.
|
||||
- The change-stream principal is its own role (`MongoPrincipalRole.CHANGE_STREAM`) with
|
||||
`changeStream` and `find` — not the application write credential.
|
||||
- Oplog window is the recovery budget. If the oplog holds four hours, a consumer that is down for five
|
||||
hours needs a rebuild, not a resume. Alert on consumer lag against the oplog window, not against
|
||||
wall-clock.
|
||||
Reference in New Issue
Block a user