From 113d330a2273e7482b6ba0fa19aaa47b53d9f7a0 Mon Sep 17 00:00:00 2001 From: Mike Shearer Date: Thu, 6 Aug 2026 22:29:28 -0600 Subject: [PATCH 1/5] build: ship an env template and a dedicated postgres service The service-backed tests read three connection strings and nothing in the repo said so, so a fresh checkout panics with "File .env or Env Vars not found" and no way to learn the variable names. .env.example documents all three and the README says how to run the suite. Postgres had no compose service at all: the test helper falls back to a hardcoded connection string pointing at a different project's database, which only works on a machine that happens to be running it. The new service publishes on 54321 to stay clear of both a system postgres and that container, and the backend migrates its own schema onto it. The obsolete compose version attribute goes too; it only ever printed a warning. The hardcoded fallback in the test helper is left alone: it is live on this machine and changing it is a behavior call, not a packaging one. --- .env.example | 28 ++++++++++++++++++++++++++++ README.md | 23 +++++++++++++++++++++++ docker-compose.yml | 17 ++++++++++++++++- 3 files changed, 67 insertions(+), 1 deletion(-) create mode 100644 .env.example diff --git a/.env.example b/.env.example new file mode 100644 index 0000000..3603b4b --- /dev/null +++ b/.env.example @@ -0,0 +1,28 @@ +# Copy this file to .env before running the integration tests: +# +# cp .env.example .env +# docker compose up -d +# cargo test +# +# The backends read these at test time via dotenv, and the suite panics +# with "File .env or Env Vars not found" when the file is missing. + +# EventStoreDB. Matches the eventstore.db service in docker-compose.yml. +# Used by the esdb feature, which is on by default. +ESDB_CONNECTION_STRING=esdb://admin:changeit@localhost:2113?tls=false + +# Redis with the RedisJSON module. Matches the redis service in +# docker-compose.yml. Used by the redis feature, which is on by default. +REDIS_CONNECTION_STRING=redis://localhost:6379 + +# Postgres. Matches the postgres service in docker-compose.yml, which +# publishes on 54321 to stay clear of any system postgres on 5432. +# +# The postgres feature is NOT in the default set, so its tests only run +# when you ask for them: +# +# cargo test --features postgres +# +# Leaving this unset falls back to a hardcoded connection string in the +# test helper that points at a different project's database, so set it. +EPOCH_PG_TEST_URL=postgres://epoch:epoch@localhost:54321/epoch diff --git a/README.md b/README.md index f657575..81d85e2 100644 --- a/README.md +++ b/README.md @@ -5,6 +5,29 @@ Event Sourcing + CQRS Framework This project is a collection of event sourcing and cqrs types to support some small personal projects heavily inluenced by [Thalo](https://github.com/thalo-rs/thalo) but borrowing (or will be borrowing in the future) ideas from Haskell [Eventful](https://github.com/jdreaver/eventful), F# [Equinox](https://github.com/jet/equinox), and Kotlin [f(model)](https://github.com/fraktalio/fmodel). +### Running the tests + +The in-memory backend needs nothing. Every other backend talks to a real +service, so bring the services up and give the suite their connection +strings: + +```sh +cp .env.example .env +docker compose up -d +cargo test +``` + +`cargo test` covers the default features (in-memory, EventStoreDB, +Redis). The postgres backend is behind a non-default feature and +migrates its own schema on first use: + +```sh +cargo test --features postgres +``` + +Without a `.env` the service-backed tests panic with `File .env or Env +Vars not found`; `.env.example` documents every variable they read. + ### Example ```rust use epoch::{event_store::ESDBEventStore, EventEnvelope, EventStore, EventContext}; diff --git a/docker-compose.yml b/docker-compose.yml index 76d98a4..9cf24c8 100644 --- a/docker-compose.yml +++ b/docker-compose.yml @@ -1,4 +1,3 @@ -version: "3.9" services: redis: image: "redislabs/rejson:latest" @@ -8,6 +7,20 @@ services: command: redis-server --save 120 1 --appendonly yes --loglevel warning --loadmodule /usr/lib/redis/modules/rejson.so volumes: - redis:/data + postgres: + image: postgres:16-alpine + restart: always + environment: + - POSTGRES_USER=epoch + - POSTGRES_PASSWORD=epoch + - POSTGRES_DB=epoch + # Published on 54321 rather than 5432 so a system postgres, or + # another project's container, does not collide with it. The + # postgres backend migrates its own schema on first use. + ports: + - '54321:5432' + volumes: + - postgres:/var/lib/postgresql/data eventstore.db: image: eventstore/eventstore:20.10.2-buster-slim environment: @@ -33,6 +46,8 @@ services: volumes: redis: driver: local + postgres: + driver: local eventstore-volume-data: driver: local eventstore-volume-logs: From 6ff5d97d58f2a1e8e35adb2e7c3e8ee529270ec3 Mon Sep 17 00:00:00 2001 From: Mike Shearer Date: Thu, 6 Aug 2026 22:37:32 -0600 Subject: [PATCH 2/5] fix: honor .env in the postgres test helper, harden the service Review of the previous commit found the env template documented a variable that did nothing. The postgres helper reads EPOCH_PG_TEST_URL through std::env::var while only the redis and esdb helpers load .env, so a connection string set in .env was ignored and the hardcoded fallback ran instead - silently pointing the suite at an unrelated project's database. Proven before and after: with .env aimed at a dead port the tests passed, and now they fail. The service also bound to every interface with throwaway credentials and carried a restart policy it has no business having as a test fixture. It binds to loopback now and stays down until asked for. --- docker-compose.yml | 10 ++++++---- src/repository/postgres/mod.rs | 5 +++++ 2 files changed, 11 insertions(+), 4 deletions(-) diff --git a/docker-compose.yml b/docker-compose.yml index 9cf24c8..1a08ff4 100644 --- a/docker-compose.yml +++ b/docker-compose.yml @@ -9,16 +9,18 @@ services: - redis:/data postgres: image: postgres:16-alpine - restart: always environment: - POSTGRES_USER=epoch - POSTGRES_PASSWORD=epoch - POSTGRES_DB=epoch # Published on 54321 rather than 5432 so a system postgres, or - # another project's container, does not collide with it. The - # postgres backend migrates its own schema on first use. + # another project's container, does not collide with it, and bound + # to the loopback address because the credentials are throwaway and + # this database is only ever a test fixture. No restart policy for + # the same reason: it comes up with `docker compose up` and should + # not survive a reboot on its own. ports: - - '54321:5432' + - '127.0.0.1:54321:5432' volumes: - postgres:/var/lib/postgresql/data eventstore.db: diff --git a/src/repository/postgres/mod.rs b/src/repository/postgres/mod.rs index bc94e04..936c0ff 100644 --- a/src/repository/postgres/mod.rs +++ b/src/repository/postgres/mod.rs @@ -308,6 +308,11 @@ mod tests { const BASE_STREAM: u32 = const_random!(u32); async fn repo_from_environment(stream_type: &str) -> PgEventRepository { + // Load .env the way the other backends do. Without this the + // variable below is only ever read from the ambient environment, + // so a connection string set in .env is silently ignored and the + // fallback runs instead - against whatever database that names. + let _ = dotenv::dotenv(); let conn_str = std::env::var("EPOCH_PG_TEST_URL") .unwrap_or_else(|_| "postgres://vikunja:devpass@localhost:54320/vikunja".to_string()); let pool = PgEventRepository::::pool_from_conn_str(&conn_str) From b4254274c79b7b3817c1e9a1cc5fcbbfe616c149 Mon Sep 17 00:00:00 2001 From: Mike Shearer Date: Thu, 6 Aug 2026 22:47:25 -0600 Subject: [PATCH 3/5] fix: require an explicit postgres test url, drop the fallback The helper fell back to a hardcoded connection string naming another project's database. These tests migrate schema and write events into whatever they connect to, so an unset or misspelled variable was not a harmless default - it silently pointed the suite at someone else's data. There is no safe database to guess, so the variable is required and the run stops without it. Verified both ways: the suite passes driven by .env, and with .env removed every postgres test fails with the message naming .env.example rather than connecting anywhere. --- .env.example | 4 ++-- src/repository/postgres/mod.rs | 15 +++++++++------ 2 files changed, 11 insertions(+), 8 deletions(-) diff --git a/.env.example b/.env.example index 3603b4b..aaccee0 100644 --- a/.env.example +++ b/.env.example @@ -23,6 +23,6 @@ REDIS_CONNECTION_STRING=redis://localhost:6379 # # cargo test --features postgres # -# Leaving this unset falls back to a hardcoded connection string in the -# test helper that points at a different project's database, so set it. +# Required: the postgres tests refuse to run without it rather than +# guessing a database, since they migrate schema into whatever they hit. EPOCH_PG_TEST_URL=postgres://epoch:epoch@localhost:54321/epoch diff --git a/src/repository/postgres/mod.rs b/src/repository/postgres/mod.rs index 936c0ff..3b3f172 100644 --- a/src/repository/postgres/mod.rs +++ b/src/repository/postgres/mod.rs @@ -308,13 +308,16 @@ mod tests { const BASE_STREAM: u32 = const_random!(u32); async fn repo_from_environment(stream_type: &str) -> PgEventRepository { - // Load .env the way the other backends do. Without this the - // variable below is only ever read from the ambient environment, - // so a connection string set in .env is silently ignored and the - // fallback runs instead - against whatever database that names. + // Load .env the way the other backends do, then require the + // variable. These tests migrate schema and write events into + // whatever database they are pointed at, so there is no safe + // default to fall back on: an unset or misspelled variable must + // stop the run rather than silently pick a database. let _ = dotenv::dotenv(); - let conn_str = std::env::var("EPOCH_PG_TEST_URL") - .unwrap_or_else(|_| "postgres://vikunja:devpass@localhost:54320/vikunja".to_string()); + let conn_str = std::env::var("EPOCH_PG_TEST_URL").expect( + "EPOCH_PG_TEST_URL must be set (see .env.example; \ + `cp .env.example .env && docker compose up -d`)", + ); let pool = PgEventRepository::::pool_from_conn_str(&conn_str) .await .expect("pg pool from EPOCH_PG_TEST_URL"); From 34bdd74ecf005257e2d43d593edf836b9a582516 Mon Sep 17 00:00:00 2001 From: Mike Shearer Date: Fri, 2 Oct 2026 12:08:52 -0600 Subject: [PATCH 4/5] feat(release): prepare epoch-journal 0.1.0 Package the current Decider API for its first public release. Update the backend dependencies and test checks to support the release gates. Card: E17 Card: E23 --- .claude/context.md | 33 ++- .claude/session-start.md | 6 +- .github/workflows/ci.yml | 49 ++++ AGENTS.md | 15 ++ CHANGELOG.md | 42 ++- CLAUDE.md | 1 + Cargo.toml | 20 +- README.md | 246 ++++++++---------- TODO.md | 17 +- deny.toml | 20 ++ examples/counter.rs | 46 ++++ src/decider.rs | 2 +- src/lib.rs | 3 + src/repository/esdb/mod.rs | 30 +-- src/repository/event.rs | 6 +- src/repository/in_memory/simple.rs | 6 +- src/repository/in_memory/state/versioned.rs | 5 +- .../in_memory/versioned_with_streams/mod.rs | 10 +- src/repository/postgres/mod.rs | 14 +- src/repository/redis/mod.rs | 35 +-- src/repository/redis/versioned_event.rs | 13 +- .../redis/versioned_stream_snapshot.rs | 15 +- src/strategies/mod.rs | 37 +-- src/test_helpers/deciders.rs | 39 ++- src/test_helpers/redis.rs | 37 +-- src/test_helpers/repository.rs | 35 +-- 26 files changed, 445 insertions(+), 337 deletions(-) create mode 100644 .github/workflows/ci.yml create mode 100644 AGENTS.md create mode 100644 CLAUDE.md create mode 100644 deny.toml create mode 100644 examples/counter.rs diff --git a/.claude/context.md b/.claude/context.md index 9adea40..31e2cb1 100644 --- a/.claude/context.md +++ b/.claude/context.md @@ -1,37 +1,34 @@ # Epoch Development Context > **Automatically loaded in Claude Code sessions** -> This file provides essential context about the Epoch repository for LLM-assisted development. +> Start here for Epoch's development conventions and source map. ## Quick Reference - **Repository**: Event Sourcing + CQRS Framework (Rust) -- **Version**: 1.0.0-alpha.18 +- **Release candidate**: 0.1.0 (`epoch-journal` package, `epoch` library) - **Core Pattern**: Decider Pattern (pure functional event sourcing) - **Rust Edition**: 2021 +- **Minimum Rust version**: 1.88 ## Session Start Protocol **At the start of each session, review**: -1. **[TODO.md](../TODO.md)** - Current work items and priorities - - Check "Current Sprint" section for active work - - Review "High Priority" for next tasks - - Note any blockers in "Known Issues" +1. **[TODO.md](../TODO.md)** - Archived 2025 planning notes, not the + current release backlog. Use the README and changelog for the + shipped API and version. 2. **[CHANGELOG.md](../CHANGELOG.md)** - Recent changes - Review "Unreleased" section for latest updates - Understand what changed since last session -3. **[Session Start Checklist](.claude/session-start.md)** - Detailed session setup +3. **[Session Start Checklist](session-start.md)** - Detailed session setup - Verify development environment - Review core principles - Set session goals -**Throughout the session**: -- Update TODO.md when starting/completing work -- Add to CHANGELOG.md for user-facing changes -- Keep both files synchronized with work progress +**Throughout the session**: Update CHANGELOG.md for user-facing changes. ## Essential Reading @@ -220,9 +217,8 @@ pub enum RepositoryVersion { ## Known Issues & TODOs -1. **README Outdated**: References old `EventContext` API instead of Decider pattern -2. **Retry Logic**: Uses `thread::sleep` instead of `tokio::sleep` (alpha pragmatism) -3. **Missing Feature**: `LoadDecideAppendWithSnapshot` not yet implemented +1. **Retry Logic**: Uses `thread::sleep` instead of `tokio::sleep` +2. **Missing Feature**: `LoadDecideAppendWithSnapshot` not yet implemented ## File Locations @@ -238,7 +234,8 @@ pub enum RepositoryVersion { │ │ ├── state.rs # State repository traits │ │ ├── in_memory/ # In-memory backend │ │ ├── esdb/ # EventStoreDB backend -│ │ └── redis/ # Redis backend +│ │ ├── redis/ # Redis backend +│ │ └── postgres/ # PostgreSQL backend │ └── test_helpers/ │ ├── deciders.rs # Example UserDecider │ └── repository.rs # Generic spec tests @@ -259,10 +256,10 @@ Following guidelines from [claude-skills](https://github.com/Shearerbeard/claude ### Internal vs. External Documentation -- **Internal** (`docs/internal/`): Architecture decisions, coding patterns, LLM context +- **Internal** (`docs/internal/`): Design decisions and coding conventions for maintainers - Target audience: Developers and LLM assistants - Focus: WHY decisions were made, HOW patterns work - - Examples: Architecture docs, style guides, planning documents + - Examples: Architecture and style guides, plus implementation plans - **External** (README, doc comments): User-facing API documentation - Target audience: Library users @@ -277,7 +274,7 @@ Files in `.claude/` and `docs/internal/` are designed to provide LLM assistants - Common tasks and their implementations - Known issues and limitations -This allows for consistent, context-aware assistance across sessions. +The source tree and README remain authoritative when these notes fall behind. ## Quick Start Commands diff --git a/.claude/session-start.md b/.claude/session-start.md index 2b7f0db..cbb47bf 100644 --- a/.claude/session-start.md +++ b/.claude/session-start.md @@ -1,7 +1,7 @@ # Session Start Checklist -> **Automatic Reference for Claude Code Sessions** -> This file provides a quick checklist and context for starting new development sessions. +> Archived 2025 checklist. Use the root README for current setup, +> feature, and test commands. `TODO.md` is also a historical snapshot. **Date**: {SESSION_DATE} @@ -99,7 +99,7 @@ cargo fmt --check ### Success Criteria By end of session: -- [ ] All tests passing +- [ ] Run the feature-specific suite documented in the README - [ ] Code formatted and linted - [ ] TODO.md updated - [ ] CHANGELOG.md updated (if applicable) diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml new file mode 100644 index 0000000..918b016 --- /dev/null +++ b/.github/workflows/ci.yml @@ -0,0 +1,49 @@ +name: CI + +on: + push: + pull_request: + +jobs: + rust: + runs-on: ubuntu-latest + steps: + - uses: actions/checkout@v4 + - name: Install supported Rust toolchain + run: rustup toolchain install 1.88.0 --profile minimal --component clippy,rustfmt + - name: Format + run: cargo +1.88.0 fmt --check + - name: Lint all targets and features + run: cargo +1.88.0 clippy --all-targets --all-features -- -D warnings + - name: Test without backends + run: cargo +1.88.0 test --no-default-features + - name: Check README doctest and example + run: cargo +1.88.0 test --doc && cargo +1.88.0 run --no-default-features --example counter + + backends: + runs-on: ubuntu-latest + steps: + - uses: actions/checkout@v4 + - name: Install supported Rust toolchain + run: rustup toolchain install 1.88.0 --profile minimal + - name: Start test services + run: | + cp .env.example .env + docker compose up -d + for attempt in $(seq 1 60); do + if ( /dev/null && + ( /dev/null && + ( /dev/null; then + exit 0 + fi + sleep 2 + done + docker compose ps + exit 1 + - name: Test default backends + run: cargo +1.88.0 test + - name: Test every backend + run: cargo +1.88.0 test --all-features + - name: Stop test services + if: always() + run: docker compose down diff --git a/AGENTS.md b/AGENTS.md new file mode 100644 index 0000000..c65a535 --- /dev/null +++ b/AGENTS.md @@ -0,0 +1,15 @@ +# Epoch contributor context + +Start with [README.md](README.md). It owns the install, example, +feature, build, and test instructions. [CHANGELOG.md](CHANGELOG.md) +records changes. The `TODO.md` file is an archived 2025 snapshot, not +the current release plan. + +`src/decider.rs` defines the domain interfaces. Persistence interfaces +and backends live in `src/repository/`; `src/strategies/` combines +deciders with repositories. The public counter example is in +`examples/counter.rs` and appears verbatim in the README. + +Treat the README's Rust block as executable: it is included in crate +docs and runs under `cargo test --doc`. Run the relevant feature tests +and the checks named in the README when changing behavior or docs. diff --git a/CHANGELOG.md b/CHANGELOG.md index d292f32..f3e289e 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -1,6 +1,9 @@ # Changelog + + All notable changes to this project will be documented in this file. + The format is based on [Keep a Changelog](https://keepachangelog.com/en/1.0.0/), and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0.html). @@ -8,7 +11,13 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 ## [Unreleased] ### Added -- Comprehensive internal documentation structure +- A compiled counter example that exercises the current Decider and + Evolver traits without a backend. +- An accessor for the events in `CommandResponse`. +- Pull-request checks for the supported Rust floor and service-backed + repository tests. +- PostgreSQL repository backend behind the opt-in `postgres` feature. +- Internal documentation structure - Architecture and philosophy documentation (docs/internal/planning/epoch-architecture-philosophy.md) - Coding style guide with Railway-Oriented Programming patterns (docs/internal/planning/coding-style-guide.md) - Documentation guidelines for internal vs external docs (docs/internal/documentation-guidelines.md) @@ -30,7 +39,7 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 - LLM-assisted development patterns - Templates for TODO items and CHANGELOG entries - PostgreSQL repository implementation planning document (docs/internal/planning/postgres-repository-implementation.md) - - Comprehensive 4-phase implementation plan (8-12 hours total) + - 4-phase implementation plan (8-12 hours total) - Database schema design with JSONB event storage - Trait implementation patterns following ESDB and Redis - Connection pooling strategy with bb8 @@ -41,6 +50,15 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 - Performance considerations and indexing strategy ### Changed +- Declare Rust 1.88 as the minimum supported version after checking + an unlocked dependency resolution with every feature enabled. +- Prepare the first public package as `epoch-journal` while retaining + `epoch` as the library import. Git consumers must update their + dependency declaration to name the new package. +- Update the EventStoreDB client to 4.0 and replace the test-only + `dotenv` dependency with `dotenvy`. +- Take event slices in repository append methods instead of requiring + `Vec` references. - Enhanced coding style guide with trucker_buddy_rs patterns - Added Railway-Oriented Programming section with visual diagrams - Added "Making Illegal States Unrepresentable" principle @@ -64,12 +82,17 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 - None ### Fixed -- None +- Replace the removed EventContext README example with a runnable + Decider example. +- Use one consistent ordering implementation for Redis stream versions. ### Security - None -## [1.0.0-alpha.18] - Prior to Documentation +## Repository history (unpublished) + +The old `1.0.0-alpha.18` version was used in the repository, not +published to crates.io. The notes below describe that code line. ### Added - Decider pattern traits (Evolver, Decider, DeciderWithContext) @@ -101,9 +124,9 @@ This project uses [Semantic Versioning](https://semver.org/): - **MAJOR** version: Incompatible API changes - **MINOR** version: Add functionality in a backwards compatible manner - **PATCH** version: Backwards compatible bug fixes -- **Alpha/Beta** suffix: Pre-release versions (current: alpha) +- **Alpha/Beta** suffix: Pre-release versions -### Alpha Status (1.0.0-alpha.x) +### Earlier alpha status (1.0.0-alpha.x) During alpha: - API may change without notice @@ -116,8 +139,8 @@ During alpha: To move from alpha to 1.0.0 stable: - [ ] API is stable and documented - [ ] README examples use current Decider pattern -- [ ] All three backends (in-memory, ESDB, Redis) fully tested -- [ ] Comprehensive documentation complete +- [ ] All four backends (in-memory, ESDB, Redis, PostgreSQL) fully tested +- [ ] Public traits and backend limits documented - [ ] LoadDecideAppendWithSnapshot implemented - [ ] Migration guide from alpha to 1.0 - [ ] Performance benchmarks established @@ -193,6 +216,3 @@ To move from alpha to 1.0.0 stable: - [Conventional Commits](https://www.conventionalcommits.org/) --- - -[Unreleased]: https://github.com/Shearerbeard/Epoch/compare/v1.0.0-alpha.18...HEAD -[1.0.0-alpha.18]: https://github.com/Shearerbeard/Epoch/releases/tag/v1.0.0-alpha.18 diff --git a/CLAUDE.md b/CLAUDE.md new file mode 100644 index 0000000..43c994c --- /dev/null +++ b/CLAUDE.md @@ -0,0 +1 @@ +@AGENTS.md diff --git a/Cargo.toml b/Cargo.toml index 05f907d..e401255 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -1,9 +1,18 @@ [package] -name = "epoch" -version = "1.0.0-alpha.18" +name = "epoch-journal" +version = "0.1.0" edition = "2021" +rust-version = "1.88" +description = "Event sourcing with deciders and pluggable event repositories" +license = "Apache-2.0" +repository = "https://github.com/Shearerbeard/Epoch" +readme = "README.md" +keywords = ["event-sourcing", "cqrs", "event-store", "decider"] +categories = ["database-implementations", "asynchronous"] +exclude = ["/.claude/", "/.github/", "/AGENTS.md", "/CLAUDE.md", "/TODO.md", "/deny.toml", "/docs/internal/", "/docs/design/", "/docs/research/", "/docs/README.md"] -# See more keys and their definitions at https://doc.rust-lang.org/cargo/reference/manifest.html +[lib] +name = "epoch" [features] default = ["in_memory", "esdb", "redis"] @@ -14,7 +23,7 @@ postgres = ["dep:tokio-postgres", "dep:bb8", "dep:bb8-postgres", "dep:serde_json [dependencies] async-trait = "0.1.53" -eventstore = { version = "2.2.0", optional = true } +eventstore = { version = "4.0.0", optional = true } redis-om = { version = "0.1.0", features = ["json"], optional = true} rusty_ulid = "2.0.0" serde = { version = "1.0.136", features = ["derive"] } @@ -30,5 +39,6 @@ actix-rt = "2.7.0" assert_matches = "1.5.0" const-random = "0.1.15" autoincrement = "1" -dotenv = "0.15.0" +dotenvy = "0.15" futures = "0.3.25" +tokio = { version = "1", features = ["time"] } diff --git a/README.md b/README.md index 81d85e2..7ddfa8a 100644 --- a/README.md +++ b/README.md @@ -1,163 +1,143 @@ # Epoch -Event Sourcing + CQRS Framework -### Inspiration -This project is a collection of event sourcing and cqrs types to support some small personal projects heavily inluenced by [Thalo](https://github.com/thalo-rs/thalo) but borrowing (or will be borrowing in the future) ideas from Haskell [Eventful](https://github.com/jdreaver/eventful), F# [Equinox](https://github.com/jet/equinox), and Kotlin [f(model)](https://github.com/fraktalio/fmodel). +Epoch is a Rust event-sourcing library built around deciders and event +repositories. It separates the decision to emit events from the work of +storing them. `epoch-journal` is the package name on crates.io; the Rust +import remains `epoch`. +The library grew out of small personal projects, heavily influenced by +[Thalo](https://github.com/thalo-rs/thalo). Ideas from Haskell +[Eventful](https://github.com/jdreaver/eventful), F# +[Equinox](https://github.com/jet/equinox), and Kotlin +[f(model)](https://github.com/fraktalio/fmodel) have also informed its +direction. The API is still changing before 1.0. -### Running the tests +## Start here -The in-memory backend needs nothing. Every other backend talks to a real -service, so bring the services up and give the suite their connection -strings: +Add the package under its Rust import name: -```sh -cp .env.example .env -docker compose up -d -cargo test +```toml +[dependencies] +epoch = { package = "epoch-journal", version = "0.1.0", default-features = false } ``` -`cargo test` covers the default features (in-memory, EventStoreDB, -Redis). The postgres backend is behind a non-default feature and -migrates its own schema on first use: +The example below uses the pure `Decider` and `Evolver` traits. It needs +no database, runtime, or feature flag. Put it in `src/main.rs` and run +`cargo run`; it prints `1`. The same program lives in +[`examples/counter.rs`](examples/counter.rs), where the repository +build checks it as an example target. -```sh -cargo test --features postgres -``` - -Without a `.env` the service-backed tests panic with `File .env or Env -Vars not found`; `.env.example` documents every variable they read. - -### Example ```rust -use epoch::{event_store::ESDBEventStore, EventEnvelope, EventStore, EventContext}; -use serde::{Deserialize, Serialize}; -use thiserror::Error; -use async_trait::async_trait; - -use example::domain::user::{UserId, User, UserName}, +use epoch::decider::{Decider, Event, Evolver}; #[derive(Debug)] -pub struct Users {} - -/// IMPL EventContext to definte relationships between Event, Cmd, CmdErr, and some sort of reified State (This will be reverbed to a Decider pattern with fn decide() and fn evolve()) -#[async_trait] -impl EventContext for Users { - type Id = String; - type Command = UserCommand; - type Event = UserEvent; - type Err = UserError; - type Services = (); - type State = UserState; - -// Event Context requires we implement string names for each context - this is used at the event persistance layer - fn event_context() -> String { - "Users".to_string() - } +enum CounterEvent { + Incremented, +} - async fn handle( - state: Self::State, - cmd: Self::Command, - _services: Self::Services, - ) -> Result, Self::Err> { - match cmd { - UserCommand::AddUser(AddUserCommand { name }) => { - todo!() - } - UserCommand::UpdateUser(UpdateUserCommand { user_id, name }) => { - todo!() - } - } - } +impl Event for CounterEvent { + type EntityId = (); - fn apply(mut state: UserState, event: &EventEnvelope) -> UserState { - match event.data.clone() { - UserEvent::UserAdded { user_id, name } => { - todo!() - } - UserEvent::UserUpdated { user_id, name } => { - todo!() - } - }; - state + fn event_type(&self) -> String { + "Incremented".to_owned() } -} -#[derive(Default, Debug, Clone)] -pub struct UserState { - users: HashMap, + fn get_id(&self) -> Self::EntityId {} } -#[derive(Error, Debug, PartialEq, Eq, Clone)] -pub enum UserError { - // Error ADT -} +struct Counter; -// Arbitrary service calling EventContext::execute() -pub struct UserService { - event_store: ESDBEventStore -} +impl Evolver for Counter { + type State = u64; + type Evt = CounterEvent; -impl UserService { - pub async fn cmd_add_user( - &self, - cmd: AddUserCommand, - ) -> Result< - ( - EventEnvelope, - ::Position, - ), - UserServiceError, - > { - // Execute knows how to handle a command and apply state automatically in the context of an event store (soon to be re-verbed for a decider pattern - execute knows how to decide and evolve state) - self.event_store - .execute::(UserCommand::AddUser(cmd), (), None) - .await - .map_err(UserServiceError::EventStoreError) + fn evolve(state: u64, event: &CounterEvent) -> u64 { + match event { + CounterEvent::Incremented => state + 1, + } } } -#[derive(Error, Debug)] -pub enum UserServiceError { - // Error ADT -} +impl Decider for Counter { + type Cmd = (); + type Err = std::convert::Infallible; -// EventSourcing types and primatives can be whatever you want - usually Event and Cmd are parameteratized enums as ADT -pub enum UserCommand { - AddUser(AddUserCommand), - UpdateUser(UpdateUserCommand), + fn decide(_state: &u64, _cmd: &()) -> Result, Self::Err> { + Ok(vec![CounterEvent::Incremented]) + } } -#[derive(Serialize, Deserialize, Debug)] -pub struct AddUserCommand { - pub name: String, +fn main() { + let state = 0; + let events = Counter::decide(&state, &()).unwrap(); + let next = events.iter().fold(state, Counter::evolve); + assert_eq!(next, 1); + println!("{next}"); } +``` -#[derive(Serialize, Deserialize, Debug)] -pub struct UpdateUserCommand { - pub user_id: String, - pub name: Option, -} +`decide` produces events from a command and current state; `evolve` +folds those events into a new state. A repository stores the resulting +events. The example stops at the domain boundary so it can run without +choosing a backend. -#[derive(Debug, Clone, Serialize, Deserialize)] -pub enum UserEvent { - UserAdded { - user_id: UserId, - name: UserName, - }, - UserUpdated { - user_id: UserId, - name: Option, - }, -} +## Repositories and features -// Event Context requires we implement string names for event types - this are used at the event persistance layer -impl Event for UserEvent { - fn event_type(&self) -> String { - match self { - UserEvent::UserAdded { .. } => "UserAdded".to_string(), - UserEvent::UserUpdated { .. } => "UserUpdated".to_string(), - } - } -} -``` \ No newline at end of file +`src/decider.rs` defines the domain traits. `src/repository/` holds +repository interfaces and backend implementations; `src/strategies/` +composes repository operations with deciders. The `in_memory` feature +needs no service. The default feature set enables `in_memory`, `esdb` +(EventStoreDB), and `redis` (RedisJSON). `postgres` is opt-in and +provides a PostgreSQL repository. Enable only the backend you use, for +example: + +```toml +[dependencies] +epoch = { package = "epoch-journal", version = "0.1.0", default-features = false, features = ["postgres"] } +``` + +This release has no `streams`, atomic-batch, feed, saga, or outbox API. +Backend version semantics differ; consult the repository trait and +backend implementation when moving stored events between backends. +The PostgreSQL repository applies its schema with +`PgEventRepository::migrate` before use. Deciders do not persist +anything on their own. + +## Build and test + +Rust 1.88 or later is required. For the example and the pure domain +surface, `cargo test --no-default-features` needs no external services. +The default suite exercises EventStoreDB and Redis; the PostgreSQL +tests run only with its feature enabled. From a clone of the repository: + +```sh +cp .env.example .env +docker compose up -d +cargo test +cargo test --features postgres +cargo test --doc +cargo fmt --check +cargo clippy --all-targets --all-features -- -D warnings +``` + +Docker Compose supplies EventStoreDB, RedisJSON, and PostgreSQL. +`.env.example` names every connection string used by the integration +tests. These services and credentials are local test fixtures; use your +own connection settings in an application. The PostgreSQL tests call +`migrate` and require `EPOCH_PG_TEST_URL` rather than choosing a +database implicitly. Run `docker compose down` when finished. + +## Versions and contributions + +`0.1.0` is the first crates.io release of `epoch-journal`. Earlier +`1.0.0-alpha.*` versions were repository versions, not publications +under this package name. Existing git consumers that declare +`epoch = { git = "https://github.com/Shearerbeard/Epoch" }` must update +their Cargo manifest to use the package declaration above; keeping +the dependency key `epoch` keeps source imports the same. + +See [CHANGELOG.md](CHANGELOG.md) for changes and earlier repository +history. [LICENSE](LICENSE) contains the Apache-2.0 terms. To contribute, +open a pull request with the relevant `cargo test` feature set, +`cargo fmt --check`, and `cargo clippy --all-targets --all-features -- +-D warnings` results. CI runs the build checks on pull requests. diff --git a/TODO.md b/TODO.md index 5bbc845..7a4b322 100644 --- a/TODO.md +++ b/TODO.md @@ -1,7 +1,8 @@ # Epoch TODO -> **Project Work Items and Planning** -> This file tracks current work items, planned features, and known issues for the Epoch project. +> **Archived planning snapshot from 2025-11-21.** The work items and +> statuses below describe that session, not the current release. For +> the shipped API and version, start with README.md and CHANGELOG.md. **Last Updated**: 2025-11-21 @@ -42,7 +43,7 @@ 2. [ ] Create examples directory with compilable examples - Location: Create examples/ directory - Reason: Helps users understand patterns in practice - - Contents: User domain, Truck domain, Expense domain examples + - Contents: User and Truck examples, with an Expense example later - Effort: ~3 hours 3. [ ] Replace thread::sleep with tokio::sleep in retry logic @@ -85,7 +86,7 @@ - Reason: Provide relational database option for event sourcing - Dependencies: tokio-postgres, bb8, bb8-postgres - Effort: ~8-12 hours (4 phases) - - Includes: Connection pooling, optimistic concurrency, generic spec tests + - Includes: A connection pool and version checks, covered by the repository spec suite - [ ] Implement LoadDecideAppendWithSnapshot strategy - [ ] Add projection pattern for read models - [ ] Add event upcasting support for schema evolution @@ -179,7 +180,7 @@ ## Completed ### Recently Completed (2025-11-21 Session) -- [x] Create comprehensive internal documentation structure +- [x] Create internal documentation structure - Added docs/internal/planning/ directory - Created epoch-architecture-philosophy.md - Created coding-style-guide.md @@ -207,7 +208,7 @@ - Integration with git workflow - [x] Plan PostgreSQL repository implementation - - Created comprehensive planning document + - Created PostgreSQL planning document - Researched ESDB and Redis patterns - Reviewed Thalo PostgreSQL implementation - Defined schema, trait implementation, testing strategy @@ -229,7 +230,7 @@ - **Ideas / Research**: Needs investigation - **Backlog**: Deferred for later -2. **Use clear, actionable descriptions**: +2. **Describe the work plainly**: ```markdown - [ ] Add validation helpers for Railway-Oriented Programming ``` @@ -276,7 +277,7 @@ When completing TODO items: ## Priority Definitions -- **High**: Blocks other work, affects users, or critical for next release +- **High**: Blocks the next release or an existing user - **Medium**: Important features or improvements, but not blocking - **Low**: Nice to have, quality of life improvements - **Backlog**: Future ideas, not currently planned diff --git a/deny.toml b/deny.toml new file mode 100644 index 0000000..f48c2cf --- /dev/null +++ b/deny.toml @@ -0,0 +1,20 @@ +[advisories] +ignore = [ + { id = "RUSTSEC-2025-0134", reason = "User-approved for 0.1.0: EventStoreDB 4.0 pulls tonic 0.12.3, which still depends on unmaintained rustls-pemfile 2.2.0; no vulnerability was reported in the release audit" }, +] +unused-ignored-advisory = "deny" + +[licenses] +allow = [ + "Apache-2.0", + "BSD-2-Clause", + "BSD-3-Clause", + "CDLA-Permissive-2.0", + "ISC", + "MIT", + "Unicode-3.0", +] + +[sources] +unknown-registry = "deny" +unknown-git = "deny" diff --git a/examples/counter.rs b/examples/counter.rs new file mode 100644 index 0000000..83b6fe4 --- /dev/null +++ b/examples/counter.rs @@ -0,0 +1,46 @@ +use epoch::decider::{Decider, Event, Evolver}; + +#[derive(Debug)] +enum CounterEvent { + Incremented, +} + +impl Event for CounterEvent { + type EntityId = (); + + fn event_type(&self) -> String { + "Incremented".to_owned() + } + + fn get_id(&self) -> Self::EntityId {} +} + +struct Counter; + +impl Evolver for Counter { + type State = u64; + type Evt = CounterEvent; + + fn evolve(state: u64, event: &CounterEvent) -> u64 { + match event { + CounterEvent::Incremented => state + 1, + } + } +} + +impl Decider for Counter { + type Cmd = (); + type Err = std::convert::Infallible; + + fn decide(_state: &u64, _cmd: &()) -> Result, Self::Err> { + Ok(vec![CounterEvent::Incremented]) + } +} + +fn main() { + let state = 0; + let events = Counter::decide(&state, &()).unwrap(); + let next = events.iter().fold(state, Counter::evolve); + assert_eq!(next, 1); + println!("{next}"); +} diff --git a/src/decider.rs b/src/decider.rs index b7468f4..5063e8c 100644 --- a/src/decider.rs +++ b/src/decider.rs @@ -30,7 +30,7 @@ pub trait Evolver { fn evolve(state: Self::State, event: &Self::Evt) -> Self::State; } -#[cfg(test)] +#[cfg(all(test, feature = "in_memory"))] mod tests { use assert_matches::assert_matches; diff --git a/src/lib.rs b/src/lib.rs index 6f34de6..7c014e9 100644 --- a/src/lib.rs +++ b/src/lib.rs @@ -1,3 +1,6 @@ +#![doc = include_str!("../README.md")] +#![allow(clippy::needless_doctest_main)] + pub mod decider; pub mod repository; pub mod strategies; diff --git a/src/repository/esdb/mod.rs b/src/repository/esdb/mod.rs index 75457b7..d00f8a8 100644 --- a/src/repository/esdb/mod.rs +++ b/src/repository/esdb/mod.rs @@ -33,7 +33,7 @@ impl ESDBEventRepository { Self { client: client.to_owned(), stream_name: stream_name.to_owned(), - _hidden: PhantomData::default(), + _hidden: PhantomData, } } @@ -136,7 +136,7 @@ where &mut self, version: &RepositoryVersion, stream: &Self::StreamId, - events: &Vec, + events: &[E], ) -> Result<(Vec, RepositoryVersion), VersionedRepositoryError> where 'a: 'async_trait, @@ -174,7 +174,7 @@ where })?; Ok(( - events.to_owned(), + events.to_vec(), RepositoryVersion::Exact(res.next_expected_version.try_into().unwrap()), )) } @@ -182,8 +182,6 @@ where #[cfg(test)] mod tests { - use const_random::const_random; - use eventstore::DeleteStreamOptions; use super::*; @@ -196,11 +194,9 @@ mod tests { }, }; - const BASE_STREAM: u32 = const_random!(u32); - async fn store_from_environment(base_stream: &str, ids: Vec) -> eventstore::Client { - let _ = dotenv::dotenv().expect("File .env or Env Vars not found"); - let settings = dotenv::var("ESDB_CONNECTION_STRING") + let _ = dotenvy::dotenv().expect("File .env or Env Vars not found"); + let settings = dotenvy::var("ESDB_CONNECTION_STRING") .expect("ESDB to be set in env") .parse() .expect("ESDB connection string to parse"); @@ -210,7 +206,7 @@ mod tests { for id in ids { let _ = client .delete_stream( - format!("{}-{}", base_stream, id), + format!("{base_stream}-{id}"), &DeleteStreamOptions::default(), ) .await; @@ -221,20 +217,18 @@ mod tests { #[actix_rt::test] async fn repository_spec_tests() { - let base_stream = BASE_STREAM; - let client = store_from_environment(&base_stream.to_string(), vec![1, 2]).await; - let event_repository = - ESDBEventRepository::::new(&client, &base_stream.to_string()); + let base_stream = uuid::Uuid::new_v4().simple().to_string(); + let client = store_from_environment(&base_stream, vec![1, 2]).await; + let event_repository = ESDBEventRepository::::new(&client, &base_stream); let _ = versioned_event_repository_with_streams_spec(event_repository).await; } #[actix_rt::test] async fn repository_with_occ_spec_test() { - let base_stream = format!("{}_with_occ", BASE_STREAM); - let client = store_from_environment(&base_stream.to_string(), vec![1]).await; - let event_repository = - ESDBEventRepository::::new(&client, &base_stream.to_string()); + let base_stream = format!("{}_with_occ", uuid::Uuid::new_v4().simple()); + let client = store_from_environment(&base_stream, vec![1]).await; + let event_repository = ESDBEventRepository::::new(&client, &base_stream); let _ = versioned_event_repository_with_streams_occ_spec(event_repository).await; } diff --git a/src/repository/event.rs b/src/repository/event.rs index 7a944a4..a5885f0 100644 --- a/src/repository/event.rs +++ b/src/repository/event.rs @@ -11,7 +11,7 @@ where E: Event + Sync + Send, { async fn load(&self) -> Result, Err>; - async fn append(&mut self, events: &Vec) -> Result, Err>; + async fn append(&mut self, events: &[E]) -> Result, Err>; } #[async_trait] @@ -26,7 +26,7 @@ where async fn append( &mut self, version: &Self::Version, - events: &Vec, + events: &[E], ) -> Result<(Vec, Self::Version), Err>; } @@ -65,7 +65,7 @@ where &mut self, version: &RepositoryVersion, stream: &Self::StreamId, - events: &Vec, + events: &[E], ) -> Result< (Vec, RepositoryVersion), VersionedRepositoryError, diff --git a/src/repository/in_memory/simple.rs b/src/repository/in_memory/simple.rs index 5dca007..a6c8e08 100644 --- a/src/repository/in_memory/simple.rs +++ b/src/repository/in_memory/simple.rs @@ -41,12 +41,12 @@ where Ok(lock.events.clone()) } - async fn append(&mut self, events: &Vec) -> Result, ()> { + async fn append(&mut self, events: &[E]) -> Result, ()> { let mut lock = self.state.lock().unwrap(); - lock.events.extend(events.to_owned()); + lock.events.extend_from_slice(events); lock.position = lock.events.len(); - Ok(events.clone()) + Ok(events.to_vec()) } } diff --git a/src/repository/in_memory/state/versioned.rs b/src/repository/in_memory/state/versioned.rs index bd3c1ef..029e769 100644 --- a/src/repository/in_memory/state/versioned.rs +++ b/src/repository/in_memory/state/versioned.rs @@ -124,10 +124,7 @@ pub enum Error { #[cfg(test)] mod tests { use crate::test_helpers::{ - deciders::user::UserDeciderState, - repository::{ - versioned_event_repository_with_streams_occ_spec, vesioned_state_repository_spec, - }, + deciders::user::UserDeciderState, repository::vesioned_state_repository_spec, }; use super::*; diff --git a/src/repository/in_memory/versioned_with_streams/mod.rs b/src/repository/in_memory/versioned_with_streams/mod.rs index ab0b913..dcc213a 100644 --- a/src/repository/in_memory/versioned_with_streams/mod.rs +++ b/src/repository/in_memory/versioned_with_streams/mod.rs @@ -52,7 +52,7 @@ where } fn get_stream_or_new(&mut self, key: &str) -> &Arc>> { - if self.state.get(key).is_none() { + if !self.state.contains_key(key) { self.state.insert( key.to_owned(), Arc::new(Mutex::new(InMemoryEventRepositoryState::new())), @@ -113,7 +113,7 @@ where &mut self, version: &RepositoryVersion, stream: &Self::StreamId, - events: &Vec, + events: &[E], ) -> Result<(Vec, RepositoryVersion), VersionedRepositoryError> where 'a: 'async_trait, @@ -124,7 +124,7 @@ where let mut stream = self.get_stream_or_new(&stream_key).lock().unwrap(); if stream.position == Self::index_from_version(version) { - stream.events.extend(events.clone()); + stream.events.extend_from_slice(events); let position = stream.events.len() - 1; stream.position = position; @@ -134,11 +134,11 @@ where .get_stream_or_new(&self.get_base_stream_key()) .lock() .unwrap(); - sub_stream.events.extend(events.clone()); + sub_stream.events.extend_from_slice(events); let sub_position = sub_stream.events.len() - 1; sub_stream.position = sub_position; - Ok((events.to_owned(), RepositoryVersion::Exact(position))) + Ok((events.to_vec(), RepositoryVersion::Exact(position))) } else { Err(Error::VersionConflict(VersionDiff::new( *version, diff --git a/src/repository/postgres/mod.rs b/src/repository/postgres/mod.rs index 3b3f172..069aa58 100644 --- a/src/repository/postgres/mod.rs +++ b/src/repository/postgres/mod.rs @@ -183,7 +183,7 @@ where &mut self, version: &RepositoryVersion, stream: &Self::StreamId, - events: &Vec, + events: &[E], ) -> RepoResult<(Vec, RepositoryVersion)> where 'a: 'async_trait, @@ -285,7 +285,7 @@ where .map_err(VersionedRepositoryError::RepoErr)?; Ok(( - events.to_owned(), + events.to_vec(), RepositoryVersion::Exact(PgVersion::from(next_sequence)), )) } @@ -313,7 +313,7 @@ mod tests { // whatever database they are pointed at, so there is no safe // default to fall back on: an unset or misspelled variable must // stop the run rather than silently pick a database. - let _ = dotenv::dotenv(); + let _ = dotenvy::dotenv(); let conn_str = std::env::var("EPOCH_PG_TEST_URL").expect( "EPOCH_PG_TEST_URL must be set (see .env.example; \ `cp .env.example .env && docker compose up -d`)", @@ -337,21 +337,21 @@ mod tests { #[actix_rt::test] async fn versioned_event_repository_with_streams_spec_postgres() { - let stream_type = format!("spec-{}", BASE_STREAM); + let stream_type = format!("spec-{BASE_STREAM}"); let repo = repo_from_environment(&stream_type).await; versioned_event_repository_with_streams_spec(repo).await; } #[actix_rt::test] async fn versioned_event_repository_with_streams_occ_spec_postgres() { - let stream_type = format!("occ-spec-{}", BASE_STREAM); + let stream_type = format!("occ-spec-{BASE_STREAM}"); let repo = repo_from_environment(&stream_type).await; versioned_event_repository_with_streams_occ_spec(repo).await; } #[actix_rt::test] async fn load_missing_stream_reports_no_stream() { - let stream_type = format!("missing-{}", BASE_STREAM); + let stream_type = format!("missing-{BASE_STREAM}"); let repo = repo_from_environment(&stream_type).await; let res = repo.load(Some(&"never-written".to_string())).await; assert!(matches!(res, Ok((v, RepositoryVersion::NoStream)) if v.is_empty())); @@ -361,7 +361,7 @@ mod tests { async fn concurrent_appends_conflict_deterministically() { use crate::test_helpers::deciders::user::{User, UserName}; - let stream_type = format!("race-{}", BASE_STREAM); + let stream_type = format!("race-{BASE_STREAM}"); let repo_a = repo_from_environment(&stream_type).await; let repo_b = repo_a.clone(); let stream_id = "shared".to_string(); diff --git a/src/repository/redis/mod.rs b/src/repository/redis/mod.rs index a94e092..f1c55f2 100644 --- a/src/repository/redis/mod.rs +++ b/src/repository/redis/mod.rs @@ -1,4 +1,7 @@ -use std::{error::Error, fmt::Debug}; +use std::{ + error::Error, + fmt::{self, Debug}, +}; use redis_om::RedisError; @@ -30,34 +33,12 @@ pub enum RedisVersionError { ParseVersion(String), } -#[derive(Debug, Eq, PartialEq, Copy, Clone, Serialize, Deserialize)] +#[derive(Debug, Eq, PartialEq, Ord, PartialOrd, Copy, Clone, Serialize, Deserialize)] pub struct RedisVersion { timestamp: usize, version: usize, } -impl Ord for RedisVersion { - fn cmp(&self, other: &Self) -> std::cmp::Ordering { - if self > other { - std::cmp::Ordering::Greater - } else if self < other { - std::cmp::Ordering::Less - } else { - std::cmp::Ordering::Equal - } - } -} - -impl PartialOrd for RedisVersion { - fn partial_cmp(&self, other: &Self) -> Option { - match self.timestamp.partial_cmp(&other.timestamp) { - Some(core::cmp::Ordering::Equal) => {} - ord => return ord, - } - self.version.partial_cmp(&other.version) - } -} - impl TryFrom<&str> for RedisVersion { type Error = RedisVersionError; @@ -79,9 +60,9 @@ impl TryFrom<&str> for RedisVersion { } } -impl ToString for RedisVersion { - fn to_string(&self) -> String { - format!("{}-{}", self.timestamp, self.version) +impl fmt::Display for RedisVersion { + fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result { + write!(f, "{}-{}", self.timestamp, self.version) } } diff --git a/src/repository/redis/versioned_event.rs b/src/repository/redis/versioned_event.rs index 1be9818..3568a2f 100644 --- a/src/repository/redis/versioned_event.rs +++ b/src/repository/redis/versioned_event.rs @@ -41,7 +41,7 @@ where pub fn new(client: &Client) -> Self { Self { client: client.to_owned(), - _sm: PhantomData::default(), + _sm: PhantomData, } } @@ -179,7 +179,7 @@ where &mut self, version: &RepositoryVersion, _stream: &Self::StreamId, - events: &Vec, + events: &[E], ) -> Result< (Vec, RepositoryVersion), VersionedRepositoryError, RedisVersion>, @@ -257,7 +257,7 @@ where RepositoryVersion::Exact(RedisVersion::try_from(version_str.as_str()).unwrap()); } - Ok((events.to_owned(), version)) + Ok((events.to_vec(), version)) } } @@ -267,8 +267,7 @@ mod tests { use super::*; use crate::test_helpers::{ - deciders::user::UserEvent, - redis::{TestUserDTOErr, TestUserEventDTO, TestUserEventDTOManager}, + redis::{TestUserEventDTO, TestUserEventDTOManager}, repository::{ versioned_event_repository_with_streams_occ_spec, versioned_event_repository_with_streams_spec, @@ -276,9 +275,9 @@ mod tests { }; async fn client_from_environment() -> Client { - let _ = dotenv::dotenv().expect("File .env or Env Vars not found"); + let _ = dotenvy::dotenv().expect("File .env or Env Vars not found"); - let settings: String = dotenv::var("REDIS_CONNECTION_STRING") + let settings: String = dotenvy::var("REDIS_CONNECTION_STRING") .expect("Redis to be set in env") .parse() .expect("Redis connection string to parse"); diff --git a/src/repository/redis/versioned_stream_snapshot.rs b/src/repository/redis/versioned_stream_snapshot.rs index 371d9f6..5bb74a8 100644 --- a/src/repository/redis/versioned_stream_snapshot.rs +++ b/src/repository/redis/versioned_stream_snapshot.rs @@ -78,15 +78,14 @@ where Self { client: client.clone(), snapshot_expiry: None, - _st: PhantomData::default(), - _jm: PhantomData::default(), + _st: PhantomData, + _jm: PhantomData, } } } #[async_trait] -impl<'a, State, JM> VersionedStreamSnapshotRepository - for RedisJSONSnapshotRepository +impl VersionedStreamSnapshotRepository for RedisJSONSnapshotRepository where State: Send + Sync @@ -131,7 +130,7 @@ where let mut conn = self.get_connection().await?; state - .to_dto(version.clone()) + .to_dto(*version) .save(&mut conn) .await .map_err(RedisRepositoryError::SaveError) @@ -187,7 +186,7 @@ mod tests { } fn version(&self) -> RedisVersion { - self.version.clone() + self.version } fn data(&self) -> TestModel { @@ -200,9 +199,9 @@ mod tests { } async fn client_from_environment() -> Client { - let _ = dotenv::dotenv().expect("File .env or Env Vars not found"); + let _ = dotenvy::dotenv().expect("File .env or Env Vars not found"); - let settings: String = dotenv::var("REDIS_CONNECTION_STRING") + let settings: String = dotenvy::var("REDIS_CONNECTION_STRING") .expect("Redis to be set in env") .parse() .expect("Redis connection string to parse"); diff --git a/src/strategies/mod.rs b/src/strategies/mod.rs index f3131f3..36e1ecd 100644 --- a/src/strategies/mod.rs +++ b/src/strategies/mod.rs @@ -136,11 +136,11 @@ where match event_repository.append(&version, &stream, &new_evts).await { Ok((appended_evts, _)) => return Ok(appended_evts), Err(VersionedRepositoryError::RepoErr(e)) => { - println!("Max Retries for {:?}!!", &cmd); + println!("Max Retries for {cmd:?}!!"); return Err(LoadDecideAppendError::RepositoryErr(e)); } Err(VersionedRepositoryError::VersionConflict(_)) => { - println!("RETRY #{} for {:?}!!", &r, &cmd); + println!("RETRY #{r} for {cmd:?}!!"); thread::sleep(time::Duration::new(0, 100000000 * r)); let (mut catchup_evts, new_version) = event_repository .load_from_version(&version, Some(&stream)) @@ -214,7 +214,7 @@ where return Err(ReifyDecideSaveError::RepositoryErr(e)) } Err(VersionedRepositoryError::VersionConflict(_)) => { - println!("Retry #{} for {:?} - Reload State", &r, &cmd); + println!("Retry #{r} for {cmd:?} - Reload State"); (state, version) = state_repository .reify() .await @@ -234,6 +234,12 @@ pub struct CommandResponse::State, ); +impl> CommandResponse { + pub fn events(&self) -> &[E] { + &self.1 + } +} + #[async_trait] pub trait DecideEvolveWithCommandResponse where @@ -286,7 +292,7 @@ pub enum ReifyDecideSaveError { RepositoryErr(RepoErr), } -#[cfg(test)] +#[cfg(all(test, feature = "in_memory"))] mod tests { use std::collections::HashMap; @@ -332,7 +338,7 @@ mod tests { assert_matches!( evts.first().expect("one event"), - UserEvent::UserAdded(User { id, name, .. }) if (&first_id == id) && (name.value() == "Mike".to_string()) + UserEvent::UserAdded(User { id, name, .. }) if (&first_id == id) && (name.value() == "Mike") ); let state = UserDeciderState::load_by_id( @@ -345,7 +351,7 @@ mod tests { assert_matches!( state, - UserDeciderState { users } if users == HashMap::from([(first_id.clone(), User::new(first_id, UserName::try_from("Mike".to_string()).unwrap()))]) + UserDeciderState { users } if users == HashMap::from([(first_id, User::new(first_id, UserName::try_from("Mike".to_string()).unwrap()))]) ); let cmd2 = UserCommand::AddUser("Dmitiry".to_string()); @@ -364,7 +370,7 @@ mod tests { assert_matches!( evts.first().expect("one event"), - UserEvent::UserAdded(User { id, name, .. }) if (&second_id == id) && (name.value() == "Dmitiry".to_string()) + UserEvent::UserAdded(User { id, name, .. }) if (&second_id == id) && (name.value() == "Dmitiry") ); let state = UserDeciderState::load_by_id( @@ -377,10 +383,10 @@ mod tests { assert_matches!( state, - UserDeciderState { users } if users == HashMap::from([(second_id.clone(), User::new(second_id, UserName::try_from("Dmitiry".to_string()).unwrap()))]) + UserDeciderState { users } if users == HashMap::from([(second_id, User::new(second_id, UserName::try_from("Dmitiry".to_string()).unwrap()))]) ); - let cmd3 = UserCommand::UpdateUserName(second_id.clone(), "Dmitiry2".to_string()); + let cmd3 = UserCommand::UpdateUserName(second_id, "Dmitiry2".to_string()); let evts = UserDecider::execute( UserDeciderState::default(), &mut event_repository, @@ -407,11 +413,10 @@ mod tests { assert_matches!( state, - UserDeciderState { users } if users == HashMap::from([(second_id.clone(), User::new(second_id, UserName::try_from("Dmitiry2".to_string()).unwrap()))]) + UserDeciderState { users } if users == HashMap::from([(second_id, User::new(second_id, UserName::try_from("Dmitiry2".to_string()).unwrap()))]) ); - let cmd4 = - UserCommand::UpdateUserName(second_id.clone(), "DmitiryWayToLongToSucceed".to_string()); + let cmd4 = UserCommand::UpdateUserName(second_id, "DmitiryWayToLongToSucceed".to_string()); let res = UserDecider::execute( UserDeciderState::default(), @@ -425,7 +430,7 @@ mod tests { assert_matches!( res, - Err(LoadDecideAppendError::DecideErr(UserDeciderError::UserField(UserFieldError::NameToLong(n)))) if n == "DmitiryWayToLongToSucceed".to_string() + Err(LoadDecideAppendError::DecideErr(UserDeciderError::UserField(UserFieldError::NameToLong(n)))) if n == "DmitiryWayToLongToSucceed" ); let state = UserDeciderState::load_by_id( @@ -438,7 +443,7 @@ mod tests { assert_matches!( state, - UserDeciderState { users } if users == HashMap::from([(second_id.clone(), User::new(second_id, UserName::try_from("Dmitiry2".to_string()).unwrap()))]) + UserDeciderState { users } if users == HashMap::from([(second_id, User::new(second_id, UserName::try_from("Dmitiry2".to_string()).unwrap()))]) ); let state = UserDeciderState::load(UserDeciderState::default(), &event_repository) @@ -448,8 +453,8 @@ mod tests { assert_matches!( state, UserDeciderState { users } if users == HashMap::from([ - (first_id.clone(), User::new(first_id, UserName::try_from("Mike".to_string()).unwrap())), - (second_id.clone(), User::new(second_id, UserName::try_from("Dmitiry2".to_string()).unwrap())) + (first_id, User::new(first_id, UserName::try_from("Mike".to_string()).unwrap())), + (second_id, User::new(second_id, UserName::try_from("Dmitiry2".to_string()).unwrap())) ]) ); } diff --git a/src/test_helpers/deciders.rs b/src/test_helpers/deciders.rs index b9f0ac6..112767e 100644 --- a/src/test_helpers/deciders.rs +++ b/src/test_helpers/deciders.rs @@ -120,8 +120,8 @@ pub(crate) mod user { ) -> Result, UserDeciderError> { match cmd { UserCommand::AddUser(user_name) => { - let name = UserName::try_from(user_name) - .map_err(|e| UserDeciderError::UserField(e))?; + let name = + UserName::try_from(user_name).map_err(UserDeciderError::UserField)?; Ok(vec![UserEvent::UserAdded(User { id: 1, @@ -130,18 +130,18 @@ pub(crate) mod user { })]) } UserCommand::UpdateUserName(user_id, user_name) => { - let name = UserName::try_from(user_name) - .map_err(|e| UserDeciderError::UserField(e))?; + let name = + UserName::try_from(user_name).map_err(UserDeciderError::UserField)?; Ok(vec![UserEvent::UserNameUpdated(user_id.to_owned(), name)]) } UserCommand::AddGuitar(user_id, guitar) => { - println!("ADD GUITAR STATE: {:?}", &state); + println!("ADD GUITAR STATE: {state:?}"); let user = state .users - .get(&user_id) + .get(user_id) .ok_or(UserDeciderError::NotFound(*user_id))?; - if user.guitars.contains(&guitar) { + if user.guitars.contains(guitar) { Err(UserDeciderError::AlreadyHasGuitar(guitar.to_owned())) } else { Ok(vec![UserEvent::UserGuitarAdded( @@ -166,13 +166,13 @@ pub(crate) mod user { state } UserEvent::UserNameUpdated(user_id, user_name) => { - state.users.get_mut(&user_id).unwrap().name = user_name.to_owned(); + state.users.get_mut(user_id).unwrap().name = user_name.to_owned(); state } UserEvent::UserGuitarAdded(user_id, guitar) => { state .users - .get_mut(&user_id) + .get_mut(user_id) .unwrap() .guitars .insert(guitar.to_owned()); @@ -199,8 +199,8 @@ pub(crate) mod user { let seq = ctx.id_sequence.lock().unwrap(); let IdGen(id) = seq.pull(); - let name = UserName::try_from(user_name) - .map_err(|e| UserDeciderError::UserField(e))?; + let name = + UserName::try_from(user_name).map_err(UserDeciderError::UserField)?; Ok(vec![UserEvent::UserAdded(User { id, @@ -209,17 +209,17 @@ pub(crate) mod user { })]) } UserCommand::UpdateUserName(user_id, user_name) => { - let name = UserName::try_from(user_name) - .map_err(|e| UserDeciderError::UserField(e))?; + let name = + UserName::try_from(user_name).map_err(UserDeciderError::UserField)?; Ok(vec![UserEvent::UserNameUpdated(user_id.to_owned(), name)]) } UserCommand::AddGuitar(user_id, guitar) => { let user = state .users - .get(&user_id) + .get(user_id) .ok_or(UserDeciderError::NotFound(*user_id))?; - if user.guitars.contains(&guitar) { + if user.guitars.contains(guitar) { Err(UserDeciderError::AlreadyHasGuitar(guitar.to_owned())) } else { Ok(vec![UserEvent::UserGuitarAdded( @@ -255,7 +255,7 @@ pub(crate) mod user { } pub fn set_users(&self, users: HashMap) -> Self { - Self { users, ..*self } + Self { users } } } @@ -278,11 +278,6 @@ pub(crate) mod user { id_sequence: Arc::new(Mutex::new(IdGen::init())), } } - - pub fn current(&self) -> usize { - let IdGen(id) = self.id_sequence.lock().unwrap().current(); - id - } } impl StateFromEventRepository for UserDeciderState { @@ -296,6 +291,8 @@ pub(crate) mod user { AddGuitar(UserId, Guitar), } + // The event names are persisted by the repository fixtures. + #[allow(clippy::enum_variant_names)] #[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)] pub(crate) enum UserEvent { UserAdded(User), diff --git a/src/test_helpers/redis.rs b/src/test_helpers/redis.rs index f629f02..615383b 100644 --- a/src/test_helpers/redis.rs +++ b/src/test_helpers/redis.rs @@ -3,10 +3,7 @@ use std::collections::HashSet; use redis_om::{redis, RedisTransportValue, StreamModel}; use thiserror::Error; -use crate::repository::{ - redis::{versioned_event::StreamModelDTO, RedisRepositoryError}, - WithFineGrainedStreamId, -}; +use crate::repository::{redis::versioned_event::StreamModelDTO, WithFineGrainedStreamId}; use super::{ deciders::user::{Guitar, User, UserEvent, UserId, UserName}, @@ -22,6 +19,8 @@ pub struct TestUserEventDTO { guitar: Option, } +// These names are part of the stored Redis fixture format. +#[allow(clippy::enum_variant_names)] #[derive(RedisTransportValue, Debug, Clone, Copy)] pub(crate) enum UserEventTypeDTO { UserAdded, @@ -93,21 +92,21 @@ impl StreamModelDTO for UserEvent { fn into_dto(self) -> ::Data { match self { UserEvent::UserAdded(user) => TestUserEventDTO { - user_id: user.id.into(), + user_id: user.id, event_type: UserEventTypeDTO::UserAdded, user: Some(user.into()), guitar: None, user_name: None, }, UserEvent::UserNameUpdated(user_id, user_name) => TestUserEventDTO { - user_id: user_id.into(), + user_id, event_type: UserEventTypeDTO::UserNameUpdated, user: None, user_name: Some(user_name.value()), guitar: None, }, UserEvent::UserGuitarAdded(user_id, guitar) => TestUserEventDTO { - user_id: user_id.into(), + user_id, event_type: UserEventTypeDTO::UserGuitarAdded, user: None, user_name: None, @@ -125,34 +124,26 @@ impl StreamModelDTO for UserEvent { match model.event_type { UserEventTypeDTO::UserAdded => match model.user { Some(user) => Ok(UserEvent::UserAdded(user.into())), - None => Err(TestUserDTOErr( - format!( - "Redis UserEventDTO invalid: missing Some(User), {:?}", - model - ) - .into(), - )), + None => Err(TestUserDTOErr(format!( + "Redis UserEventDTO invalid: missing Some(User), {model:?}" + ))), }, UserEventTypeDTO::UserNameUpdated => match model.user_name { Some(user_name) => Ok(UserEvent::UserNameUpdated( model.user_id, UserName::try_from(user_name).map_err(|e| { - TestUserDTOErr(format!("Redis UserEventDTO invalid: {:?}", e).into()) + TestUserDTOErr(format!("Redis UserEventDTO invalid: {e:?}")) })?, )), None => Err(TestUserDTOErr(format!( - "Redis UserEventDTO invalid: missing Some(UserName), {:?}", - model - )) - .into()), + "Redis UserEventDTO invalid: missing Some(UserName), {model:?}" + ))), }, UserEventTypeDTO::UserGuitarAdded => match model.guitar { Some(guitar) => Ok(UserEvent::UserGuitarAdded(model.user_id, guitar.into())), None => Err(TestUserDTOErr(format!( - "Redis UserEventDTO invalid: missing Some(Guitar), {:?}", - model - )) - .into()), + "Redis UserEventDTO invalid: missing Some(Guitar), {model:?}" + ))), }, } } diff --git a/src/test_helpers/repository.rs b/src/test_helpers/repository.rs index 1e3e5d1..59afcb7 100644 --- a/src/test_helpers/repository.rs +++ b/src/test_helpers/repository.rs @@ -1,8 +1,7 @@ -use core::time; use std::{ collections::{HashMap, HashSet}, fmt::Debug, - thread, + time::{Duration, Instant}, }; use assert_matches::assert_matches; @@ -81,18 +80,24 @@ pub(crate) async fn versioned_event_repository_with_streams_spec< .await .expect("Successful append"); - // Crude but we need to wait for ESDB to catch up its "Categories" auto projection - thread::sleep(time::Duration::from_secs(1)); - let res = event_repository.load(Some(&id_1)).await; assert_matches!(res, Ok((v, RepositoryVersion::Exact(_))) if v == events1); let res = event_repository.load(Some(&id_2)).await; assert_matches!(res, Ok((v, RepositoryVersion::Exact(_))) if v == events2); - let res = event_repository.load(None).await; - - let events_combined: Vec = events1.into_iter().chain(events2.into_iter()).collect(); + let events_combined: Vec = events1.into_iter().chain(events2).collect(); + // EventStoreDB's category projection trails the stream writes. + let deadline = Instant::now() + Duration::from_secs(10); + let res = loop { + let res = event_repository.load(None).await; + if matches!(&res, Ok((events, _)) if *events == events_combined) + || Instant::now() >= deadline + { + break res; + } + tokio::time::sleep(Duration::from_millis(100)).await; + }; assert_matches!(res, Ok((v, RepositoryVersion::Exact(_))) if v == events_combined); let res = event_repository.load(Some(&id_1)).await; @@ -127,7 +132,7 @@ pub(crate) async fn vesioned_state_repository_spec<'a, Err: Debug + Send + Sync> )])); let version = RepositoryVersion::Exact(0); - println!("Saving: state={:?}, version={:?}", &new_state, &version); + println!("Saving: state={new_state:?}, version={version:?}"); let _ = state_repository .save(&version, &new_state) .await @@ -172,7 +177,7 @@ pub(crate) async fn versioned_event_repository_with_streams_occ_spec< assert_matches!( evts.first().expect("one event"), - UserEvent::UserAdded(User { id, name, .. }) if (&first_id == id) && (name.value() == "Mike".to_string()) + UserEvent::UserAdded(User { id, name, .. }) if (&first_id == id) && (name.value() == "Mike") ); let state = UserDeciderState::load_by_id( @@ -185,7 +190,7 @@ pub(crate) async fn versioned_event_repository_with_streams_occ_spec< assert_matches!( state, - UserDeciderState { users } if users == HashMap::from([(first_id.clone(), User::new(first_id, UserName::try_from("Mike".to_string()).unwrap()))]) + UserDeciderState { users } if users == HashMap::from([(first_id, User::new(first_id, UserName::try_from("Mike".to_string()).unwrap()))]) ); let guitars = vec![ @@ -221,13 +226,11 @@ pub(crate) async fn versioned_event_repository_with_streams_occ_spec< let futures = guitars .iter() .cloned() - .map(|g| add_guitar(event_repository.clone(), first_id.clone(), g).boxed()) + .map(|g| add_guitar(event_repository.clone(), first_id, g).boxed()) .collect::>>(); future::join_all(futures).await; - thread::sleep(time::Duration::from_secs(1)); - let state = UserDeciderState::load_by_id( UserDeciderState::default(), &event_repository, @@ -255,7 +258,7 @@ async fn add_guitar< ) { let ctx = UserDeciderCtx::new(); - println!("Adding Guitar {:?} for user {}", &guitar.brand, &user_id); + println!("Adding Guitar {:?} for user {}", guitar.brand, user_id); let cmd = UserCommand::AddGuitar(user_id, guitar.to_owned()); @@ -271,6 +274,6 @@ async fn add_guitar< println!( "Result for Guitar {:?} for user {}: {:?}", - &guitar.brand, &user_id, res + guitar.brand, user_id, res ); } From 3e20bac8bc8b7587f80f05328e09580129edfa3f Mon Sep 17 00:00:00 2001 From: Mike Shearer Date: Fri, 2 Oct 2026 12:58:24 -0600 Subject: [PATCH 5/5] fix(test): limit category polling to esdb Wait for EventStoreDB's asynchronous category projection only. Other backends keep an immediate category assertion, and read errors fail without retrying. Card: E23 --- src/repository/esdb/mod.rs | 6 +++- .../in_memory/versioned_with_streams/mod.rs | 2 +- src/repository/postgres/mod.rs | 2 +- src/repository/redis/versioned_event.rs | 2 +- src/test_helpers/repository.rs | 30 ++++++++++++------- 5 files changed, 27 insertions(+), 15 deletions(-) diff --git a/src/repository/esdb/mod.rs b/src/repository/esdb/mod.rs index d00f8a8..5a5dd95 100644 --- a/src/repository/esdb/mod.rs +++ b/src/repository/esdb/mod.rs @@ -221,7 +221,11 @@ mod tests { let client = store_from_environment(&base_stream, vec![1, 2]).await; let event_repository = ESDBEventRepository::::new(&client, &base_stream); - let _ = versioned_event_repository_with_streams_spec(event_repository).await; + let _ = versioned_event_repository_with_streams_spec( + event_repository, + Some(std::time::Duration::from_secs(10)), + ) + .await; } #[actix_rt::test] diff --git a/src/repository/in_memory/versioned_with_streams/mod.rs b/src/repository/in_memory/versioned_with_streams/mod.rs index dcc213a..40288e0 100644 --- a/src/repository/in_memory/versioned_with_streams/mod.rs +++ b/src/repository/in_memory/versioned_with_streams/mod.rs @@ -170,6 +170,6 @@ mod tests { #[actix_rt::test] async fn repository_spec_test() { let event_repository = InMemoryEventRepository::::new(BASE_STREAM); - let _ = versioned_event_repository_with_streams_spec(event_repository).await; + let _ = versioned_event_repository_with_streams_spec(event_repository, None).await; } } diff --git a/src/repository/postgres/mod.rs b/src/repository/postgres/mod.rs index 069aa58..d72e0b2 100644 --- a/src/repository/postgres/mod.rs +++ b/src/repository/postgres/mod.rs @@ -339,7 +339,7 @@ mod tests { async fn versioned_event_repository_with_streams_spec_postgres() { let stream_type = format!("spec-{BASE_STREAM}"); let repo = repo_from_environment(&stream_type).await; - versioned_event_repository_with_streams_spec(repo).await; + versioned_event_repository_with_streams_spec(repo, None).await; } #[actix_rt::test] diff --git a/src/repository/redis/versioned_event.rs b/src/repository/redis/versioned_event.rs index 3568a2f..d8ee62f 100644 --- a/src/repository/redis/versioned_event.rs +++ b/src/repository/redis/versioned_event.rs @@ -299,7 +299,7 @@ mod tests { RedisStreamsEventRepository::::new(&client); // Run both tests in sequence because we cannot specify a stream identifier per test in redis - let _ = versioned_event_repository_with_streams_spec(event_repository.clone()).await; + let _ = versioned_event_repository_with_streams_spec(event_repository.clone(), None).await; let _ = versioned_event_repository_with_streams_occ_spec(event_repository).await; } } diff --git a/src/test_helpers/repository.rs b/src/test_helpers/repository.rs index 59afcb7..d6dbfc3 100644 --- a/src/test_helpers/repository.rs +++ b/src/test_helpers/repository.rs @@ -1,7 +1,7 @@ use std::{ collections::{HashMap, HashSet}, fmt::Debug, - time::{Duration, Instant}, + time::Duration, }; use assert_matches::assert_matches; @@ -39,6 +39,7 @@ pub(crate) async fn versioned_event_repository_with_streams_spec< Version = V, StreamId = String, >, + category_projection_timeout: Option, ) { println!("RUNNING UNIVERSAL SPEC TEST FOR VersionedEventRepositoryWithStreams"); let id_1 = "1".to_string(); @@ -87,16 +88,23 @@ pub(crate) async fn versioned_event_repository_with_streams_spec< assert_matches!(res, Ok((v, RepositoryVersion::Exact(_))) if v == events2); let events_combined: Vec = events1.into_iter().chain(events2).collect(); - // EventStoreDB's category projection trails the stream writes. - let deadline = Instant::now() + Duration::from_secs(10); - let res = loop { - let res = event_repository.load(None).await; - if matches!(&res, Ok((events, _)) if *events == events_combined) - || Instant::now() >= deadline - { - break res; - } - tokio::time::sleep(Duration::from_millis(100)).await; + let res = if let Some(timeout) = category_projection_timeout { + // EventStoreDB's category projection trails the stream writes. + tokio::time::timeout(timeout, async { + loop { + let res = event_repository.load(None).await; + match &res { + Ok((events, _)) if events != &events_combined => { + tokio::time::sleep(Duration::from_millis(100)).await; + } + _ => break res, + } + } + }) + .await + .expect("category projection caught up before the deadline") + } else { + event_repository.load(None).await }; assert_matches!(res, Ok((v, RepositoryVersion::Exact(_))) if v == events_combined);