Repository navigation
feat: expose job operations over the C ABI - #31
Conversation
|
Warning Review limit reached
Next review available in: 1 minute You've used all free OSS reviews for now. Wait for the free limit to reset to keep reviewing this public repository. How can I continue?After more reviews become available, a review can be triggered using the To avoid repeated limits, reduce automatic review volume by pausing incremental auto-reviews earlier, using label-based review opt-in, excluding WIP or generated PR titles, or requesting reviews manually when the PR is ready. If your team needs uninterrupted high-volume reviews, an organization admin can enable usage-based reviews. How do review limits work?CodeRabbit enforces per-developer PR review limits for each organization. Most developers receive the normal plan review availability. For paid Pro and Pro+ PR reviews, CodeRabbit uses adaptive limits for sustained high-volume activity. When a developer's recent PR review activity reaches the 95th percentile or higher among CodeRabbit users, additional reviews become available more gradually as earlier reviews age out of the rolling window. Please refer docs for additional details. Review details⚙️ Run configurationConfiguration used: Path: .coderabbit.yaml Review profile: ASSERTIVE Plan: Pro Plus Run ID: 📒 Files selected for processing (11)
WalkthroughThe engine now supports atomic durable job intents and a shared job runtime. The C ABI exposes job enqueue, claim, lease, and settlement operations. Haskell bindings provide typed codecs, pagination, ABI checks, tests, Nix packaging, and CI coverage. ChangesShared engine and job runtime
C ABI job interface
Haskell binding and delivery
Possibly related PRs
🚥 Pre-merge checks | ✅ 5✅ Passed checks (5 passed)
✨ Finishing Touches📝 Generate docstrings
🧪 Generate unit tests (beta)
Thanks for using CodeRabbit! It's free for OSS, and your support helps us grow. If you like it, consider giving us a shout-out. Comment |
9f71f3c to
b7ffede
Compare
b5e54d8 to
351ed17
Compare
b7ffede to
40dd85a
Compare
351ed17 to
1de1363
Compare
40dd85a to
527824b
Compare
1de1363 to
cfdb85a
Compare
527824b to
e94291c
Compare
cfdb85a to
a0c9f00
Compare
There was a problem hiding this comment.
Actionable comments posted: 3
🤖 Prompt for all review comments with AI agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.
Inline comments:
In `@crates/event-sorcery-ffi/src/lib.rs`:
- Around line 2820-2823: Update the invalid-instant assertions around
es_job_enqueue to use assert_eq! with the exact expected ABI error code instead
of merely checking that the result is non-success. Also decode and assert the
returned error payload’s exact variant and values in both affected test cases,
preserving the existing raw-store setup.
- Around line 883-905: Update StoreInner::issue_claim and the claim lifecycle to
expire terminal or lost claims, including claims whose next claim returns
Abandoned or Skipped, and remove stale entries during subsequent operations.
Enforce a hard upper bound on the retained active claims map, evicting expired
or otherwise eligible entries before inserting new claims while preserving valid
active claims.
- Around line 909-918: Update claim() and the associated settlement paths around
the claims map to atomically transition each token from Available to Settling
while holding the mutex before invoking the backend, rejecting tokens that are
already Settling or retired. Restore the token to Available on operational
settlement failure, preserve retirement after successful settlement, and add a
concurrent replay test covering the ABI boundary so only one caller can settle a
one-shot token.
🪄 Autofix (Beta)
Fix all unresolved CodeRabbit comments on this PR:
- Push a commit to this branch (recommended)
- Create a new PR with the fixes
ℹ️ Review info
⚙️ Run configuration
Configuration used: Path: .coderabbit.yaml
Review profile: ASSERTIVE
Plan: Pro Plus
Run ID: 74a442d9-0a65-44e8-a371-c925dc609e42
📒 Files selected for processing (2)
crates/event-sorcery-ffi/src/lib.rscrates/event-sorcery/src/job_backend.rs
📜 Review details
⏰ Context from checks skipped due to timeout. (6)
- GitHub Check: fmt
- GitHub Check: examples
- GitHub Check: clippy
- GitHub Check: check
- GitHub Check: test
- GitHub Check: hooks
🧰 Additional context used
📓 Path-based instructions (2)
**/*
📄 CodeRabbit inference engine (AGENTS.md)
**/*: Before work, read SPEC.md and docs/domain.md; read relevant supplemental documentation before implementation.
New features must be documented in SPEC.md before implementation and must follow the hierarchy SPEC.md -> issue -> plan -> tests -> implementation.
Fix all known problems immediately, complete all tasks, and do not allow warnings or errors to pass through.
Keep a granular task list and clear completed tasks from the active list.
All new or modified logic must have corresponding test coverage.
Understand relevant documentation and source code before implementation, keep diffs small, and review the approach critically.
When changing direction or making an important undocumented architectural decision, obtain confirmation; record significant decisions as ADRs under adrs/.
Before handover, review the diff, revert unjustified changes, and check for scope creep.
Each aggregate in a consuming application must use exactly one SqliteCqrs instance constructed at startup; per-request construction is forbidden.
Never read secret or credential files such as .env*, credentials.json, *.key, *.pem, *.p12, *.pfx, or sensitive database files without explicit permission.
Never bypass, disable, suppress, or obscure quality-control mechanisms without explicit permission; fix lint and test issues at their root.
Use cargo check, cargo nextest, and cargo clippy for verification; never use cargo build unless build artifacts are required.
Files:
crates/event-sorcery/src/job_backend.rscrates/event-sorcery-ffi/src/lib.rs
crates/**/*.rs
📄 CodeRabbit inference engine (AGENTS.md)
crates/**/*.rs: Organize code by business feature rather than technical layer; avoid catch-all modules such as types.rs, error.rs, models.rs, utils.rs, helpers.rs, and services.rs.
Never write directly to the events table; emit events through CqrsFramework::execute() or execute_with_metadata().
Use cqrs-es Services for side effects in handle() and follow the {Action}er -> {Domain}Service -> {Domain}Manager naming pattern.
Place command-execution logging in aggregate handle() methods rather than callers.
Model invalid states with enums, ADTs, newtypes, and typestate rather than relying on runtime validation.
Use domain newtypes at APIs and convert to SDK primitives inside the callee, except at cross-crate boundaries where conversion at the call site is necessary.
Keep visibility as restrictive as possible: private over pub(crate) over pub.
Use a three-group import order: external crates, workspace crates, then crate-internal imports; do not use function-level imports except enum variants.
Do not use unwrap() or expect() in production Rust code; they are permitted in #[cfg(test)] code.
Never create error variants containing opaque String values; prefer #[from], ?, #[source], and preserve error chains.
Log a warning or error before silent early returns such as let-else failures.
Never silently mask numeric failures with caps, fallback defaults, precision truncation, unwrap_or(), or unwrap_or_default(); use explicit checked conversions and errors.
Prefer functional patterns, pattern matching, combinators, type-driven design, and iterators over imperative loops unless complexity increases.
Use ASCII in identifiers, comments, log messages, and configuration keys; Unicode is preferred only in user-facing rendered output.
Do not use single-letter variables, arguments, closure parameters, or generic type parameters except an unambiguous lone type parameter or short unambiguous closure.
Every module must have a //! docstring and should order public API, private implementati...
Files:
crates/event-sorcery/src/job_backend.rscrates/event-sorcery-ffi/src/lib.rs
🔇 Additional comments (3)
crates/event-sorcery-ffi/src/lib.rs (2)
4-70: LGTM!Also applies to: 244-261, 408-687, 1385-1590, 1709-1814
2774-2819: LGTM!Also applies to: 2841-2942
crates/event-sorcery/src/job_backend.rs (1)
63-70: LGTM!
|
Addressed the current review in b240463: invalid job instants now assert the exact storage error code and redacted payload; retained claims carry lease expiry, stale entries are pruned, and the active map is hard-bounded without evicting live claims; settlements atomically transition Available to Settling, restore availability on operational failure, and retire on applied or fenced outcomes. Regression coverage now proves expired abandoned tokens are rejected, concurrent ABI acknowledgements admit exactly one caller, and failed settlement preserves the token. All 27 FFI tests, strict clippy, rustfmt, repository hooks, and diff checks pass. Invalid finding: retiring a token on every Skipped claim result would invalidate a still-live worker claim, because the claim planner also returns Skipped while an existing lease remains active. Skipped now retires only entries whose recorded lease is actually expired. @coderabbitai review |
|
✏️ Learnings added
✅ Action performedReview finished.
|
There was a problem hiding this comment.
Actionable comments posted: 4
🤖 Prompt for all review comments with AI agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.
Inline comments:
In `@crates/event-sorcery-ffi/src/lib.rs`:
- Around line 3149-3187: Extend the rollback failure test around
es_commit_with_job to verify the job insert is also reverted: use a unique job
kind in the rollback request, poll that job kind after the failed commit, and
assert the exact empty response. Keep the existing storage-error and version
assertions unchanged.
- Around line 600-623: The claim flow around prepare_claim, claim_job, and
issue_claim must reserve claim-map capacity before persisting a durable claim to
prevent concurrent overcommit. Add a reservation guard that is released on claim
errors and non-claimed outcomes, while converting a successful reservation into
the token capacity consumed by issue_claim; update related claim handling
consistently and add a concurrency test covering MAX_ACTIVE_CLAIMS - 1.
- Around line 683-692: Update the renewal flows around lease.inner.claim,
renew_job, and update_claim_lease to transition the token from Available to
Renewing before calling the backend. Ensure pruning skips Renewing entries, then
atomically restore the token on renewal failure, update its lease and return it
to Available on Held, or retire it when no longer held. Apply the same state
handling to the additional renewal paths at the referenced locations.
- Around line 614-630: Rename the `first_execution` variable in the claim result
construction to reflect its inverted polarity, such as `has_prior_execution`,
while preserving the existing `u8::from(!claimed.handle.is_first_execution())`
value and ABI tuple position.
🪄 Autofix (Beta)
Fix all unresolved CodeRabbit comments on this PR:
- Push a commit to this branch (recommended)
- Create a new PR with the fixes
ℹ️ Review info
⚙️ Run configuration
Configuration used: Path: .coderabbit.yaml
Review profile: ASSERTIVE
Plan: Pro Plus
Run ID: 0812f7c0-466e-4854-90f8-52b449c58409
📒 Files selected for processing (2)
crates/event-sorcery-ffi/src/lib.rscrates/event-sorcery/src/job_backend.rs
📜 Review details
⏰ Context from checks skipped due to timeout. (6)
- GitHub Check: hooks
- GitHub Check: examples
- GitHub Check: test
- GitHub Check: check
- GitHub Check: clippy
- GitHub Check: fmt
🧰 Additional context used
📓 Path-based instructions (2)
**/*
📄 CodeRabbit inference engine (AGENTS.md)
**/*: Before work, read SPEC.md and docs/domain.md; read relevant supplemental documentation before implementation.
New features must be documented in SPEC.md before implementation and must follow the hierarchy SPEC.md -> issue -> plan -> tests -> implementation.
Fix all known problems immediately, complete all tasks, and do not allow warnings or errors to pass through.
Keep a granular task list and clear completed tasks from the active list.
All new or modified logic must have corresponding test coverage.
Understand relevant documentation and source code before implementation, keep diffs small, and review the approach critically.
When changing direction or making an important undocumented architectural decision, obtain confirmation; record significant decisions as ADRs under adrs/.
Before handover, review the diff, revert unjustified changes, and check for scope creep.
Each aggregate in a consuming application must use exactly one SqliteCqrs instance constructed at startup; per-request construction is forbidden.
Never read secret or credential files such as .env*, credentials.json, *.key, *.pem, *.p12, *.pfx, or sensitive database files without explicit permission.
Never bypass, disable, suppress, or obscure quality-control mechanisms without explicit permission; fix lint and test issues at their root.
Use cargo check, cargo nextest, and cargo clippy for verification; never use cargo build unless build artifacts are required.
Files:
crates/event-sorcery/src/job_backend.rscrates/event-sorcery-ffi/src/lib.rs
crates/**/*.rs
📄 CodeRabbit inference engine (AGENTS.md)
crates/**/*.rs: Organize code by business feature rather than technical layer; avoid catch-all modules such as types.rs, error.rs, models.rs, utils.rs, helpers.rs, and services.rs.
Never write directly to the events table; emit events through CqrsFramework::execute() or execute_with_metadata().
Use cqrs-es Services for side effects in handle() and follow the {Action}er -> {Domain}Service -> {Domain}Manager naming pattern.
Place command-execution logging in aggregate handle() methods rather than callers.
Model invalid states with enums, ADTs, newtypes, and typestate rather than relying on runtime validation.
Use domain newtypes at APIs and convert to SDK primitives inside the callee, except at cross-crate boundaries where conversion at the call site is necessary.
Keep visibility as restrictive as possible: private over pub(crate) over pub.
Use a three-group import order: external crates, workspace crates, then crate-internal imports; do not use function-level imports except enum variants.
Do not use unwrap() or expect() in production Rust code; they are permitted in #[cfg(test)] code.
Never create error variants containing opaque String values; prefer #[from], ?, #[source], and preserve error chains.
Log a warning or error before silent early returns such as let-else failures.
Never silently mask numeric failures with caps, fallback defaults, precision truncation, unwrap_or(), or unwrap_or_default(); use explicit checked conversions and errors.
Prefer functional patterns, pattern matching, combinators, type-driven design, and iterators over imperative loops unless complexity increases.
Use ASCII in identifiers, comments, log messages, and configuration keys; Unicode is preferred only in user-facing rendered output.
Do not use single-letter variables, arguments, closure parameters, or generic type parameters except an unambiguous lone type parameter or short unambiguous closure.
Every module must have a //! docstring and should order public API, private implementati...
Files:
crates/event-sorcery/src/job_backend.rscrates/event-sorcery-ffi/src/lib.rs
🔇 Additional comments (6)
crates/event-sorcery/src/job_backend.rs (2)
63-70: LGTM!
1112-1112: LGTM!crates/event-sorcery-ffi/src/lib.rs (4)
26-71: LGTM!Also applies to: 245-262
393-599: LGTM!Also applies to: 633-682, 693-831
894-900: LGTM!Also applies to: 1014-1119, 1544-1588, 1658-1749, 1862-1868
1924-1944: LGTM!Also applies to: 1969-1973, 2869-3148, 3188-3273
| lease.inner.prepare_claim(&job_id, now_ms)?; | ||
| let result = lease | ||
| .inner | ||
| .runtime | ||
| .block_on(lease.inner.jobs.claim_job( | ||
| &job_id, | ||
| &worker_name, | ||
| now_ms, | ||
| lease_ms, | ||
| max_claims, | ||
| )) | ||
| .map_err(job_runtime_error)?; | ||
| let wire: JobClaimWire = match result { | ||
| JobClaimResult::Claimed(claimed) => { | ||
| let attempt = claimed.handle.attempt(); | ||
| let first_execution = u8::from(!claimed.handle.is_first_execution()); | ||
| let payload = json_payload_bytes(claimed.payload)?; | ||
| require_payload_limit(payload.len())?; | ||
| let lease_until_ms = now_ms | ||
| .checked_add(lease_ms) | ||
| .ok_or(AbiError::State("claim lease instant overflowed"))?; | ||
| let token = lease | ||
| .inner | ||
| .runtime | ||
| .block_on(lease.inner.engine.current_version(&stream)) | ||
| .map_err(AbiError::storage)?; | ||
| let actual = u64::try_from(actual).map_err(AbiError::storage)?; | ||
| Err(AbiError::Conflict { | ||
| aggregate_type, | ||
| aggregate_id, | ||
| expected, | ||
| actual, | ||
| }) | ||
| .issue_claim(claimed.handle, now_ms, lease_until_ms)?; |
There was a problem hiding this comment.
🗄️ Data Integrity & Integration | 🟠 Major | 🏗️ Heavy lift
Reserve claim-map capacity before persisting the claim.
prepare_claim releases the mutex before the durable claim, then issue_claim rechecks capacity. Two concurrent claims at MAX_ACTIVE_CLAIMS - 1 can both persist, but only one receives a token; the other returns a resource-limit error with an unreachable leased job.
Reserve a slot before calling claim_job, releasing it on errors or non-claimed outcomes. Add a capacity-bound concurrency test.
Also applies to: 901-960
🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.
In `@crates/event-sorcery-ffi/src/lib.rs` around lines 600 - 623, The claim flow
around prepare_claim, claim_job, and issue_claim must reserve claim-map capacity
before persisting a durable claim to prevent concurrent overcommit. Add a
reservation guard that is released on claim errors and non-claimed outcomes,
while converting a successful reservation into the token capacity consumed by
issue_claim; update related claim handling consistently and add a concurrency
test covering MAX_ACTIVE_CLAIMS - 1.
| let attempt = claimed.handle.attempt(); | ||
| let first_execution = u8::from(!claimed.handle.is_first_execution()); | ||
| let payload = json_payload_bytes(claimed.payload)?; | ||
| require_payload_limit(payload.len())?; | ||
| let lease_until_ms = now_ms | ||
| .checked_add(lease_ms) | ||
| .ok_or(AbiError::State("claim lease instant overflowed"))?; | ||
| let token = lease | ||
| .inner | ||
| .runtime | ||
| .block_on(lease.inner.engine.current_version(&stream)) | ||
| .map_err(AbiError::storage)?; | ||
| let actual = u64::try_from(actual).map_err(AbiError::storage)?; | ||
| Err(AbiError::Conflict { | ||
| aggregate_type, | ||
| aggregate_id, | ||
| expected, | ||
| actual, | ||
| }) | ||
| .issue_claim(claimed.handle, now_ms, lease_until_ms)?; | ||
| ( | ||
| 1, | ||
| 0, | ||
| Some(OpaqueBytes(token.to_vec())), | ||
| Some(attempt), | ||
| Some(first_execution), | ||
| Some(OpaqueBytes(payload)), |
There was a problem hiding this comment.
🎯 Functional Correctness | 🔵 Trivial | 💤 Low value
🧩 Analysis chain
🏁 Script executed:
#!/bin/bash
set -euo pipefail
for path in SPEC.md docs/domain.md; do
if [[ -f "$path" ]]; then
rg -n -i -C3 'first.execution|previous.execution|execution.risk|claim' "$path"
fi
done
rg -n -C3 'first_execution|is_first_execution|expected_execution' \
crates/event-sorcery-ffi/src/lib.rs \
crates/event-sorcery/src/job_backend.rsRepository: dataclique/event-sorcery
Length of output: 9721
🏁 Script executed:
#!/bin/bash
set -euo pipefail
# Read the relevant implementation and tests around the claim wire encoding.
sed -n '580,660p' crates/event-sorcery-ffi/src/lib.rs
printf '\n---\n'
sed -n '1930,1980p' crates/event-sorcery-ffi/src/lib.rs
printf '\n---\n'
sed -n '80,110p' crates/event-sorcery/src/job_backend.rs
printf '\n---\n'
sed -n '140,150p' SPEC.md
printf '\n---\n'
sed -n '118,126p' docs/domain.mdRepository: dataclique/event-sorcery
Length of output: 7072
🏁 Script executed:
#!/bin/bash
set -euo pipefail
# Search for the claim wire schema and any mention of the execution field.
rg -n -C2 'JobClaimWire|first_execution|execution,' crates/event-sorcery-ffi/src/lib.rs crates/event-sorcery/src/job_backend.rs SPEC.md docs/domain.mdRepository: dataclique/event-sorcery
Length of output: 4611
🏁 Script executed:
#!/bin/bash
set -euo pipefail
# Inspect the wire tuple definition and the tests that assert execution values.
sed -n '55,75p' crates/event-sorcery-ffi/src/lib.rs
printf '\n---\n'
rg -n -C3 'expected_execution|is_first_execution\(\)|execution\)' crates/event-sorcery-ffi/src/lib.rsRepository: dataclique/event-sorcery
Length of output: 2310
🏁 Script executed:
#!/bin/bash
set -euo pipefail
# Find all call sites for the helper that assert the execution byte.
rg -n -C2 'claim_test_job\(' crates/event-sorcery-ffi/src/lib.rsRepository: dataclique/event-sorcery
Length of output: 2254
🏁 Script executed:
#!/bin/bash
set -euo pipefail
# Inspect any documentation or comments that name the JobClaim wire fields.
rg -n -C2 'first execution|prior-execution|subsequent execution|execution flag|JobClaimWire|expected_execution|is_first_execution|reconcile routing' \
SPEC.md docs/domain.md crates/event-sorcery-ffi/src/lib.rs crates/event-sorcery/src/job_backend.rsRepository: dataclique/event-sorcery
Length of output: 4362
🏁 Script executed:
#!/bin/bash
set -euo pipefail
# Search the repository for any external/binding-side naming of the claim execution byte.
rg -n -C2 'JobClaimWire|first_execution|execution byte|claim.*execution|reconcile routing|first-execution' \
. \
--glob '!target/**' --glob '!dist/**' --glob '!build/**'Repository: dataclique/event-sorcery
Length of output: 8138
🏁 Script executed:
#!/bin/bash
set -euo pipefail
# Inspect the relevant wire type and a couple of adjacent tests for semantics.
sed -n '55,75p' crates/event-sorcery-ffi/src/lib.rs
printf '\n---\n'
sed -n '2848,3060p' crates/event-sorcery-ffi/src/lib.rsRepository: dataclique/event-sorcery
Length of output: 7924
Rename the flag to match its polarity. u8::from(!claimed.handle.is_first_execution()) encodes 0 for the first claim and 1 afterward, so first_execution reads backwards; a name like has_prior_execution would make the ABI slot clearer.
🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.
In `@crates/event-sorcery-ffi/src/lib.rs` around lines 614 - 630, Rename the
`first_execution` variable in the claim result construction to reflect its
inverted polarity, such as `has_prior_execution`, while preserving the existing
`u8::from(!claimed.handle.is_first_execution())` value and ABI tuple position.
| let (token, claim) = lease.inner.claim(&encoded)?; | ||
| let result = lease | ||
| .inner | ||
| .runtime | ||
| .block_on(lease.inner.jobs.renew_job(&claim, new_lease_until_ms)) | ||
| .map_err(job_runtime_error)?; | ||
| let tag = match result { | ||
| JobLeaseResult::Held => { | ||
| lease.inner.update_claim_lease(token, new_lease_until_ms); | ||
| 0 |
There was a problem hiding this comment.
🩺 Stability & Availability | 🟠 Major | 🏗️ Heavy lift
Mark tokens as renewing before invoking the backend.
Renewal merely clones an Available handle. A concurrent claim can prune the expired entry while renewal is in flight; renewal may then return Held, while update_claim_lease finds no token, leaving the caller with an unusable capability.
Add a Renewing transition/guard, exclude in-flight entries from pruning, and atomically restore, update, or retire the token.
Also applies to: 901-906, 963-1013
🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.
In `@crates/event-sorcery-ffi/src/lib.rs` around lines 683 - 692, Update the
renewal flows around lease.inner.claim, renew_job, and update_claim_lease to
transition the token from Available to Renewing before calling the backend.
Ensure pruning skips Renewing entries, then atomically restore the token on
renewal failure, update its lease and return it to Available on Held, or retire
it when no longer held. Apply the same state handling to the additional renewal
paths at the referenced locations.
| let rollback_job_id = JobId::new().to_string(); | ||
| let mut rollback_encoded = encode_request(&( | ||
| 1_u8, | ||
| "ffi-domain", | ||
| "rollback", | ||
| 0_u64, | ||
| vec![("Created", "1.0", OpaqueBytes(vec![9]))], | ||
| (rollback_job_id, "ffi-test", OpaqueBytes(vec![9]), i64::MAX), | ||
| )); | ||
| let rollback = caller_buffer(&mut rollback_encoded); | ||
| assert_eq!( | ||
| unsafe { es_commit_with_job(&raw mut store, &raw const rollback, &raw mut error) }, | ||
| ES_ERR_STORAGE | ||
| ); | ||
| assert_eq!( | ||
| decode_error(&error), | ||
| ( | ||
| 1, | ||
| ES_ERR_STORAGE, | ||
| Value::Text("storage failure".to_string()), | ||
| ) | ||
| ); | ||
| unsafe { es_buf_free(&raw mut error) }; | ||
| let mut version_encoded = encode_request(&(1_u8, "ffi-domain", "rollback")); | ||
| let version_request = caller_buffer(&mut version_encoded); | ||
| let mut version = u64::MAX; | ||
| assert_eq!( | ||
| unsafe { | ||
| es_current_version( | ||
| &raw mut store, | ||
| &raw const version_request, | ||
| &raw mut version, | ||
| &raw mut error, | ||
| ) | ||
| }, | ||
| ES_OK | ||
| ); | ||
| assert_eq!(version, 0); | ||
| assert_eq!(unsafe { es_close(&raw mut store) }, ES_OK); |
There was a problem hiding this comment.
🗄️ Data Integrity & Integration | 🟡 Minor | ⚡ Quick win
Assert that the job insert also rolls back.
The test verifies only the aggregate version. A regression that persists the job while rolling back the domain event would still pass. Use a unique job kind, poll it after failure, and assert the exact empty response.
As per coding guidelines, “All new or modified logic must have corresponding test coverage” and “Tests must assert correct behavior.”
🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.
In `@crates/event-sorcery-ffi/src/lib.rs` around lines 3149 - 3187, Extend the
rollback failure test around es_commit_with_job to verify the job insert is also
reverted: use a unique job kind in the rollback request, poll that job kind
after the failed commit, and assert the exact empty response. Keep the existing
storage-error and version assertions unchanged.
Source: Coding guidelines
b240463 to
5d17f40
Compare
|
Addressed the follow-up review in 5d17f40: rollback coverage now polls a unique failed-job kind and proves the queue remains empty; claim capacity is reserved atomically before durable claiming, including concurrent coverage at MAX_ACTIVE_CLAIMS - 1; renewals transition Available to Renewing, pruning preserves in-flight operations, and RAII restores interrupted or failed renewals; and the inverted execution flag is named has_prior_execution. Validation: 29 FFI tests, strict workspace clippy with warnings denied, rustfmt, repository hooks, and diff checks all pass. Invalid findings: none. @coderabbitai review |
b1f9be7 to
f215443
Compare
0xgleb
left a comment
There was a problem hiding this comment.
Cross-reviewed this with a full panel (three iterations, adversarially verified findings) and prepared fixes for everything that survived verification. The fixes live on review-fix/pr31 (one commit, 7da24c9, sitting directly on this PR's head) — absorb it into the stack with GitButler rather than merging the ref; the branch is offered, not submitted as its own PR.
should fix, the two that lose data or spin:
- an undecodable claimed payload used to be released with
defer(now_ms), and sinceRescheduleddeliberately resets the claim budget, a single natively-enqueued (non-opaque) job on a shared kind turned every FFI worker into a claim/defer hot loop that grows the event store forever and can never dead-letter. The fix settles it terminally asDeadReason::Undecodableand reportsES_ERR_JOB_REFUSAL, which is the class the ABI already defines as per-job terminal. - the new optimistic-lock pre-check compared
expected_versiononly againstMAX(events.sequence), and compaction is exactly the operation that shrinks that maximum below the aggregate's real version. Two writers plus a compaction in between produced a silent lost update withOkreturned to both. The conflict oracle is now bounded bymax(MAX(events.sequence), snapshots.last_sequence), which stays correct for both per-commit and every-N snapshot policies.
also folded in, briefly: the Haskell binding now decodes engine error class 3 (JOB_REFUSAL) instead of laundering it into UnknownEngineError while the catch-all comment claimed the class list was exhaustive; es_load_stream/es_current_version enforce the normative 4 KiB identifier limit like the other seven entry points; the page-truncation semantics change on es_load_stream is stated in SPEC alongside a proper version treatment; duplicate job ids report as terminal refusals rather than retryable storage failures (both the enqueue path and the atomic commit path, which previously surfaced them as an unresolvable OptimisticLock); a settlement racing an in-flight renewal gets its own retryable leaf instead of the permanent invalid-handle error; the in-process re-claim keeps the superseded token so a fenced settlement is actually reachable through the public ABI; opaque payloads stop being stored as JSON number arrays (~4x inflation on the headline CBOR path); nix/event-sorcery-haskell.nix points at the real source path so the file is usable standalone; and the EsBuf layout test now pins against the generated C header instead of asserting the hand-written instance against itself.
Gates on the fix branch: 265/265 workspace tests, clippy with all targets and features clean, fmt stable. The Haskell binding builds in the CI job the branch adds.
posted on behalf of @0xgleb by his clanka.
There was a problem hiding this comment.
Actionable comments posted: 6
🤖 Prompt for all review comments with AI agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.
Inline comments:
In `@bindings/haskell/test/EngineSpec.hs`:
- Around line 79-94: Add a dedicated test alongside the existing checkAbiVersion
cases that passes the supported major with a minor greater than minimumAbiMinor
and asserts Right (). Construct the ABI value directly so the test does not
depend on the linked engine version, while preserving the existing rejection
tests.
In `@crates/event-sorcery/src/engine.rs`:
- Around line 654-674: Update the conflict check in Engine::commit to reject any
version mismatch by changing the actual_version comparison from only
greater-than to inequality, preventing sequence gaps while preserving
snapshot-backed compaction when versions match. Replace the test covering
deletion without a snapshot with a regression test asserting that a skipped
sequence is rejected.
In `@crates/event-sorcery/src/job_backend.rs`:
- Around line 1014-1031: Bound the settlement retry loop around
classify_settlement to the current lease deadline, using the lease information
from ctx instead of retrying indefinitely. In the Transient arm, stop retrying
once the deadline is reached and apply increasing backoff between attempts;
preserve the existing Applied, Fenced, and Rejected outcomes and update the
retry log to reflect deadline-bounded behavior.
- Around line 1059-1069: Update classify_settlement to match every
AggregateError variant explicitly: handle UserError, AggregateConflict,
DatabaseConnectionError, DeserializationError, and UnexpectedError with
intentional SettlementVerdict classifications. Preserve the existing fenced and
transient behavior where applicable, and remove the catch-all Err(error) arm so
newly introduced variants cannot default to Rejected silently.
In `@nix/event-sorcery-haskell.nix`:
- Line 15: Update nix/event-sorcery-haskell.nix so the package function accepts
src and assigns it directly, pass src = ./bindings/haskell; into callPackage,
and remove overrideSrc. Apply the corresponding call-site adjustment in
flake.nix lines 44-49, ensuring both sites use the explicit source path and no
deferred override.
In `@SPEC.md`:
- Around line 267-282: Update the resource-limit table and surrounding
specification text to state that an out-of-range busy timeout is rejected at
open with MALFORMED_INPUT, not RESOURCE_LIMIT. Clarify that the 4 KiB UTF-8
identifier limit applies independently to each identifier. Define and use stable
resource names for every limit row, including busy timeout and identifier text,
so RESOURCE_LIMIT details are consistent across bindings.
🪄 Autofix
Fix all unresolved CodeRabbit comments on this PR:
- Push a commit to this branch (recommended)
- Create a new PR with the fixes
ℹ️ Review info
⚙️ Run configuration
Configuration used: Path: .coderabbit.yaml
Review profile: ASSERTIVE
Plan: Pro Plus
Run ID: 0342724a-9795-4482-9844-6f2311da5acb
⛔ Files ignored due to path filters (1)
Cargo.lockis excluded by!**/*.lock
📒 Files selected for processing (17)
.github/workflows/ci.yamlCargo.tomlSPEC.mdbindings/haskell/event-sorcery.cabalbindings/haskell/src/EventSorcery/Engine.hsbindings/haskell/src/EventSorcery/Engine/Codec.hsbindings/haskell/src/EventSorcery/Engine/Internal/Paging.hsbindings/haskell/src/EventSorcery/Engine/Protocol.hsbindings/haskell/test/CodecSpec.hsbindings/haskell/test/EngineSpec.hscrates/event-sorcery-ffi/Cargo.tomlcrates/event-sorcery-ffi/src/lib.rscrates/event-sorcery/src/engine.rscrates/event-sorcery/src/job_backend.rscrates/event-sorcery/src/lib.rsflake.nixnix/event-sorcery-haskell.nix
📜 Review details
🧰 Additional context used
📓 Path-based instructions (5)
**/*
📄 CodeRabbit inference engine (AGENTS.md)
**/*: Before work, read SPEC.md and docs/domain.md; read relevant supplemental documentation before implementation.
New features must be documented in SPEC.md before implementation and must follow the hierarchy SPEC.md -> issue -> plan -> tests -> implementation.
Fix all known problems immediately, complete all tasks, and do not allow warnings or errors to pass through.
Keep a granular task list and clear completed tasks from the active list.
All new or modified logic must have corresponding test coverage.
Understand relevant documentation and source code before implementation, keep diffs small, and review the approach critically.
When changing direction or making an important undocumented architectural decision, obtain confirmation; record significant decisions as ADRs under adrs/.
Before handover, review the diff, revert unjustified changes, and check for scope creep.
Each aggregate in a consuming application must use exactly one SqliteCqrs instance constructed at startup; per-request construction is forbidden.
Never read secret or credential files such as .env*, credentials.json, *.key, *.pem, *.p12, *.pfx, or sensitive database files without explicit permission.
Never bypass, disable, suppress, or obscure quality-control mechanisms without explicit permission; fix lint and test issues at their root.
Use cargo check, cargo nextest, and cargo clippy for verification; never use cargo build unless build artifacts are required.
Files:
nix/event-sorcery-haskell.nixCargo.tomlbindings/haskell/event-sorcery.cabalbindings/haskell/src/EventSorcery/Engine/Internal/Paging.hscrates/event-sorcery-ffi/Cargo.tomlSPEC.mdcrates/event-sorcery/src/lib.rsbindings/haskell/test/CodecSpec.hsbindings/haskell/src/EventSorcery/Engine/Codec.hsbindings/haskell/test/EngineSpec.hsflake.nixbindings/haskell/src/EventSorcery/Engine.hsbindings/haskell/src/EventSorcery/Engine/Protocol.hscrates/event-sorcery/src/engine.rscrates/event-sorcery/src/job_backend.rs
**/Cargo.toml
📄 CodeRabbit inference engine (AGENTS.md)
Never manually edit Cargo.toml to add dependencies; use cargo add, with workspace dependencies declared centrally and referenced by workspace=true.
Files:
Cargo.tomlcrates/event-sorcery-ffi/Cargo.toml
*
⚙️ CodeRabbit configuration file
Focus on providing constructive criticism. Whenever you see a suboptimal approach, suggest more idiomatic or robust alternative(s). Flag potential footguns. Suggest FP alternatives to mutable/imperative code. Point out architectural flaws like leaky abstractions, tight coupling, wrong level of abstraction, poor type modeling, over-abstraction, unclear domain boundaries. Code should generally be organized based on business concerns rather than technical aspects - suggest improvements if you find violations. Point out gaps in test coverage but suggest tests that are not too coupled to the implementation and actually test domain invariants and business logic
Files:
Cargo.tomlSPEC.mdflake.nix
**/*.md
⚙️ CodeRabbit configuration file
Focus on the contents of the docs and not on cosmetic things like markdown formatting. We use markdown files for various docs including but not limited to guidelines for AI contributors (AGENTS.md), project overview and instructions for human contributors (README.md), and topic-focused references under docs/ (cqrs.md, sqlx.md, ttdd.md). Think about the target audience of a document when deciding what comment to leave. For instructions, suggest better rules and guidelines and point out missing instructions. For topic references, suggest improvements that would make non-obvious framework behavior or pitfalls easier to discover. In all cases, flag needless bloat, prefer clear concise writing, and consider the structure of the document and order of the sections
Files:
SPEC.md
crates/**/*.rs
📄 CodeRabbit inference engine (AGENTS.md)
crates/**/*.rs: Organize code by business feature rather than technical layer; avoid catch-all modules such as types.rs, error.rs, models.rs, utils.rs, helpers.rs, and services.rs.
Never write directly to the events table; emit events through CqrsFramework::execute() or execute_with_metadata().
Use cqrs-es Services for side effects in handle() and follow the {Action}er -> {Domain}Service -> {Domain}Manager naming pattern.
Place command-execution logging in aggregate handle() methods rather than callers.
Model invalid states with enums, ADTs, newtypes, and typestate rather than relying on runtime validation.
Use domain newtypes at APIs and convert to SDK primitives inside the callee, except at cross-crate boundaries where conversion at the call site is necessary.
Keep visibility as restrictive as possible: private over pub(crate) over pub.
Use a three-group import order: external crates, workspace crates, then crate-internal imports; do not use function-level imports except enum variants.
Do not use unwrap() or expect() in production Rust code; they are permitted in #[cfg(test)] code.
Never create error variants containing opaque String values; prefer #[from], ?, #[source], and preserve error chains.
Log a warning or error before silent early returns such as let-else failures.
Never silently mask numeric failures with caps, fallback defaults, precision truncation, unwrap_or(), or unwrap_or_default(); use explicit checked conversions and errors.
Prefer functional patterns, pattern matching, combinators, type-driven design, and iterators over imperative loops unless complexity increases.
Use ASCII in identifiers, comments, log messages, and configuration keys; Unicode is preferred only in user-facing rendered output.
Do not use single-letter variables, arguments, closure parameters, or generic type parameters except an unambiguous lone type parameter or short unambiguous closure.
Every module must have a //! docstring and should order public API, private implementati...
Files:
crates/event-sorcery/src/lib.rscrates/event-sorcery/src/engine.rscrates/event-sorcery/src/job_backend.rs
🪛 zizmor (1.29.0)
.github/workflows/ci.yaml
[info] 137-137: workflow or action definition without a name (anonymous-definition): this job
(anonymous-definition)
🔇 Additional comments (23)
Cargo.toml (1)
46-46: 🎯 Functional CorrectnessThe
getrandomdependency is correctly resolved.event-sorcery-ffiusesgetrandom 0.3.4; the0.2.17entry belongs to another dependency.> Likely an incorrect or invalid review comment.bindings/haskell/src/EventSorcery/Engine.hs (1)
1-360: LGTM!bindings/haskell/src/EventSorcery/Engine/Codec.hs (1)
1-198: LGTM!bindings/haskell/src/EventSorcery/Engine/Protocol.hs (1)
1-126: LGTM!bindings/haskell/src/EventSorcery/Engine/Internal/Paging.hs (1)
1-45: LGTM!bindings/haskell/event-sorcery.cabal (1)
44-84: LGTM!bindings/haskell/test/CodecSpec.hs (1)
5-547: LGTM!bindings/haskell/test/EngineSpec.hs (1)
95-283: LGTM!.github/workflows/ci.yaml (1)
137-170: LGTM!crates/event-sorcery/src/engine.rs (7)
19-58: LGTM!
489-536: LGTM!
738-809: LGTM!
985-1040: LGTM!
1082-1128: LGTM!Also applies to: 1157-1198
1329-1792: LGTM!
1965-2042: LGTM!crates/event-sorcery/src/job_backend.rs (4)
44-50: LGTM!Also applies to: 151-163
232-292: LGTM!
514-561: LGTM!Also applies to: 636-661
1323-1535: LGTM!Also applies to: 1537-1584, 2022-2077
crates/event-sorcery/src/lib.rs (1)
934-968: LGTM!crates/event-sorcery-ffi/Cargo.toml (2)
16-16: LGTM!
26-48: 📐 Maintainability & Code QualityNo change required. All nine Clippy
allowentries match[workspace.lints.clippy]; onlyunsafe_code = "allow"differs as documented.> Likely an incorrect or invalid review comment.
| loop { | ||
| match jobs.send(&ctx.job_id, command.clone()).await { | ||
| Ok(()) => return, | ||
| Err( | ||
| AggregateError::UserError(LifecycleError::Apply(JobError::Fenced)) | ||
| | AggregateError::AggregateConflict, | ||
| ) => { | ||
| if succeeded { | ||
| error!(target: "cqrs", job_id = %ctx.job_id, "DOUBLE-EXECUTION RISK: job succeeded but a re-claimer owns the outcome"); | ||
| } else { | ||
| warn!(target: "cqrs", job_id = %ctx.job_id, "ack fenced; a re-claimer owns the outcome"); | ||
| } | ||
| return; | ||
| } | ||
| Err(AggregateError::UserError(domain_error)) => { | ||
| error!(target: "cqrs", job_id = %ctx.job_id, ?domain_error, "ack hit a non-fence lifecycle error; the aggregate is failed"); | ||
| match classify_settlement(jobs.send(&ctx.job_id, command.clone()).await) { | ||
| SettlementVerdict::Applied => return, | ||
| SettlementVerdict::Fenced => { | ||
| report_fenced_settlement(ctx.job_id, settled); | ||
| return; | ||
| } | ||
| Err(AggregateError::DatabaseConnectionError(error)) => { | ||
| SettlementVerdict::Transient(error) => { | ||
| warn!(target: "cqrs", ?error, job_id = %ctx.job_id, "ack transient error; retrying within lease"); | ||
| sleep(Duration::from_millis(50)).await; | ||
| } | ||
| Err(AggregateError::DeserializationError(error)) => { | ||
| error!(target: "cqrs", ?error, job_id = %ctx.job_id, "ack deserialization error; giving up"); | ||
| return; | ||
| } | ||
| Err(AggregateError::UnexpectedError(error)) => { | ||
| error!(target: "cqrs", ?error, job_id = %ctx.job_id, "ack unexpected error; giving up"); | ||
| SettlementVerdict::Rejected(error) => { | ||
| error!(target: "cqrs", ?error, job_id = %ctx.job_id, "ack rejected by the job aggregate; giving up"); | ||
| return; | ||
| } | ||
| } | ||
| } | ||
| } |
There was a problem hiding this comment.
🩺 Stability & Availability | 🟠 Major | ⚡ Quick win
The transient-settlement retry loop is unbounded.
The Transient arm sleeps 50 ms and loops with no attempt cap, no deadline, and no backoff. A persistent DatabaseConnectionError keeps this loop running forever. The apalis task never completes, so its concurrency slot is never released and the worker's throughput degrades to zero.
The log message states "retrying within lease", but the loop never reads the lease. When the lease expires, a re-claimer owns the job, and every further attempt can only be fenced.
Bound the loop by the lease deadline and back off between attempts.
♻️ Proposed bounded retry
let command = ack_command(config, clock, ctx, plan);
let settled = settled_outcome(&command);
+ let deadline = clock.now() + config.lease_duration;
+ let mut delay = Duration::from_millis(50);
loop {
match classify_settlement(jobs.send(&ctx.job_id, command.clone()).await) {
SettlementVerdict::Applied => return,
SettlementVerdict::Fenced => {
report_fenced_settlement(ctx.job_id, settled);
return;
}
SettlementVerdict::Transient(error) => {
+ if clock.now() >= deadline {
+ error!(target: "cqrs", ?error, job_id = %ctx.job_id, "ack transient error persisted past the lease; giving up");
+ return;
+ }
warn!(target: "cqrs", ?error, job_id = %ctx.job_id, "ack transient error; retrying within lease");
- sleep(Duration::from_millis(50)).await;
+ sleep(delay).await;
+ delay = (delay * 2).min(Duration::from_secs(5));
}
SettlementVerdict::Rejected(error) => {
error!(target: "cqrs", ?error, job_id = %ctx.job_id, "ack rejected by the job aggregate; giving up");
return;
}
}
}📝 Committable suggestion
‼️ IMPORTANT
Carefully review the code before committing. Ensure that it accurately replaces the highlighted code, contains no missing lines, and has no issues with indentation. Thoroughly test & benchmark the code to ensure it meets the requirements.
| loop { | |
| match jobs.send(&ctx.job_id, command.clone()).await { | |
| Ok(()) => return, | |
| Err( | |
| AggregateError::UserError(LifecycleError::Apply(JobError::Fenced)) | |
| | AggregateError::AggregateConflict, | |
| ) => { | |
| if succeeded { | |
| error!(target: "cqrs", job_id = %ctx.job_id, "DOUBLE-EXECUTION RISK: job succeeded but a re-claimer owns the outcome"); | |
| } else { | |
| warn!(target: "cqrs", job_id = %ctx.job_id, "ack fenced; a re-claimer owns the outcome"); | |
| } | |
| return; | |
| } | |
| Err(AggregateError::UserError(domain_error)) => { | |
| error!(target: "cqrs", job_id = %ctx.job_id, ?domain_error, "ack hit a non-fence lifecycle error; the aggregate is failed"); | |
| match classify_settlement(jobs.send(&ctx.job_id, command.clone()).await) { | |
| SettlementVerdict::Applied => return, | |
| SettlementVerdict::Fenced => { | |
| report_fenced_settlement(ctx.job_id, settled); | |
| return; | |
| } | |
| Err(AggregateError::DatabaseConnectionError(error)) => { | |
| SettlementVerdict::Transient(error) => { | |
| warn!(target: "cqrs", ?error, job_id = %ctx.job_id, "ack transient error; retrying within lease"); | |
| sleep(Duration::from_millis(50)).await; | |
| } | |
| Err(AggregateError::DeserializationError(error)) => { | |
| error!(target: "cqrs", ?error, job_id = %ctx.job_id, "ack deserialization error; giving up"); | |
| return; | |
| } | |
| Err(AggregateError::UnexpectedError(error)) => { | |
| error!(target: "cqrs", ?error, job_id = %ctx.job_id, "ack unexpected error; giving up"); | |
| SettlementVerdict::Rejected(error) => { | |
| error!(target: "cqrs", ?error, job_id = %ctx.job_id, "ack rejected by the job aggregate; giving up"); | |
| return; | |
| } | |
| } | |
| } | |
| } | |
| let command = ack_command(config, clock, ctx, plan); | |
| let settled = settled_outcome(&command); | |
| let deadline = clock.now() + config.lease_duration; | |
| let mut delay = Duration::from_millis(50); | |
| loop { | |
| match classify_settlement(jobs.send(&ctx.job_id, command.clone()).await) { | |
| SettlementVerdict::Applied => return, | |
| SettlementVerdict::Fenced => { | |
| report_fenced_settlement(ctx.job_id, settled); | |
| return; | |
| } | |
| SettlementVerdict::Transient(error) => { | |
| if clock.now() >= deadline { | |
| error!(target: "cqrs", ?error, job_id = %ctx.job_id, "ack transient error persisted past the lease; giving up"); | |
| return; | |
| } | |
| warn!(target: "cqrs", ?error, job_id = %ctx.job_id, "ack transient error; retrying within lease"); | |
| sleep(delay).await; | |
| delay = (delay * 2).min(Duration::from_secs(5)); | |
| } | |
| SettlementVerdict::Rejected(error) => { | |
| error!(target: "cqrs", ?error, job_id = %ctx.job_id, "ack rejected by the job aggregate; giving up"); | |
| return; | |
| } | |
| } | |
| } | |
| } |
🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.
In `@crates/event-sorcery/src/job_backend.rs` around lines 1014 - 1031, Bound the
settlement retry loop around classify_settlement to the current lease deadline,
using the lease information from ctx instead of retrying indefinitely. In the
Transient arm, stop retrying once the deadline is reached and apply increasing
backoff between attempts; preserve the existing Applied, Fenced, and Rejected
outcomes and update the retry log to reflect deadline-bounded behavior.
| fn classify_settlement(result: Result<(), SendError<JobState>>) -> SettlementVerdict { | ||
| match result { | ||
| Ok(()) => SettlementVerdict::Applied, | ||
| Err( | ||
| AggregateError::UserError(LifecycleError::Apply(JobError::Fenced)) | ||
| | AggregateError::AggregateConflict, | ||
| ) => SettlementVerdict::Fenced, | ||
| Err(AggregateError::DatabaseConnectionError(error)) => SettlementVerdict::Transient(error), | ||
| Err(error) => SettlementVerdict::Rejected(error), | ||
| } | ||
| } |
There was a problem hiding this comment.
🎯 Functional Correctness | 🟠 Major | ⚡ Quick win
🧩 Analysis chain
🏁 Script executed:
#!/bin/bash
# Description: List the AggregateError variants the pinned cqrs-es version defines.
set -euo pipefail
rg -n 'cqrs-es' Cargo.toml crates/*/Cargo.toml
fd -t f 'aggregate_error.rs|errors.rs' ~/.cargo/registry 2>/dev/null | head -n 20
rg -n -A 30 'pub enum AggregateError' $(fd -t d 'cqrs-es-*' ~/.cargo/registry 2>/dev/null | head -n 1) 2>/dev/null || echo "cqrs-es sources not vendored in the sandbox"Repository: dataclique/event-sorcery
Length of output: 373
🏁 Script executed:
#!/bin/bash
set -euo pipefail
printf '%s\n' '--- required documentation ---'
for file in SPEC.md docs/domain.md; do
if [ -f "$file" ]; then
printf '%s\n' "--- $file ---"
cat -n "$file"
else
printf '%s\n' "$file: not found"
fi
done
printf '%s\n' '--- dependency resolution ---'
rg -n -A 8 -B 3 'name = "cqrs-es"|cqrs-es\s*=' Cargo.toml Cargo.lock crates/*/Cargo.toml 2>/dev/null || true
printf '%s\n' '--- AggregateError references and local definitions ---'
rg -n -A 35 -B 5 'pub enum AggregateError|AggregateError::|type SendError|struct SendError' crates Cargo.toml Cargo.lock 2>/dev/null | head -n 500
printf '%s\n' '--- registry sources ---'
find "${HOME}/.cargo/registry/src" -type f \( -name '*.rs' -o -name 'Cargo.toml' \) -path '*cqrs-es-*' 2>/dev/null | head -n 40Repository: dataclique/event-sorcery
Length of output: 50382
🌐 Web query:
cqrs-es 0.5.0 AggregateError enum source
💡 Result:
In the cqrs-es crate version 0.5.0, the AggregateError enum serves as the base error type for the framework [1][2]. It is defined in the cqrs_es module and is used to handle various failure scenarios within the CQRS process, ranging from business logic errors to infrastructure issues [1][3]. The AggregateError enum typically includes the following variants [1]: UserError(T): Used when a command fails due to a business rule violation. The generic type T represents the user-defined error payload. AggregateConflict: Indicates a failure due to a conflict with another command on the same aggregate instance, commonly arising from optimistic locking in relational database systems. DatabaseConnectionError(Box<dyn Error + Send + Sync + 'static>): Represents issues related to the database connection. DeserializationError(Box<dyn Error + Send + Sync + 'static>): Triggered when invalid JSON or data format prevents deserialization. UnexpectedError(Box<dyn Error + Send + Sync + 'static>): A catch-all for errors not covered by other variants. In version 0.5.0, the framework underwent changes regarding how aggregate errors are handled; specifically, the Aggregate::handle method was modified to expect a reply of Self::Error rather than the broader AggregateError, limiting the ability to overload enum branches other than UserError at that layer [4]. The source code for the enum can be found within the official documentation on docs.rs [1].
Citations:
- 1: https://docs.rs/cqrs-es/latest/cqrs_es/enum.AggregateError.html
- 2: https://docs.rs/cqrs-es/^0.5.0/
- 3: https://docs.rs/cqrs-es/latest/cqrs_es/struct.CqrsFramework.html
- 4: https://github.com/serverlesstechnology/cqrs/blob/main/docs/versions/change_log.md
🏁 Script executed:
#!/bin/bash
set -euo pipefail
printf '%s\n' '--- settlement call graph ---'
sed -n '900,1125p' crates/event-sorcery/src/job_backend.rs
printf '%s\n' '--- error definitions and imports ---'
sed -n '1,190p' crates/event-sorcery/src/job_backend.rs
printf '%s\n' '--- cqrs-es 0.5.0 source from crates.io ---'
tmpdir="$(mktemp -d)"
trap 'rm -rf "$tmpdir"' EXIT
curl --fail --silent --show-error --location \
https://crates.io/api/v1/crates/cqrs-es/0.5.0/download \
--output "$tmpdir/cqrs-es.tar.gz"
tar -xzf "$tmpdir/cqrs-es.tar.gz" -C "$tmpdir"
rg -n -A 35 -B 5 'enum AggregateError|UnexpectedError|DeserializationError|DatabaseConnectionError' \
"$tmpdir" --glob '*.rs'Repository: dataclique/event-sorcery
Length of output: 16076
🌐 Web query:
site:docs.rs/cqrs-es/0.5.0/cqrs_es/enum.AggregateError.html "pub enum AggregateError"
💡 Result:
AggregateError<T> in cqrs-es 0.5.0 is defined with these variants:
UserError(T)AggregateConflictDatabaseConnectionError(Box<dyn Error + Send + Sync + 'static>)DeserializationError(Box<dyn Error + Send + Sync + 'static>)UnexpectedError(Box<dyn Error + Send + Sync + 'static>)
It is generic over T: std::error::Error. (docs.rs)
Citations:
Use an exhaustive AggregateError match.
Handle UserError, AggregateConflict, DatabaseConnectionError, DeserializationError, and UnexpectedError explicitly. Assign each variant an intentional settlement classification instead of silently treating new or infrastructure errors as terminal Rejected.
🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.
In `@crates/event-sorcery/src/job_backend.rs` around lines 1059 - 1069, Update
classify_settlement to match every AggregateError variant explicitly: handle
UserError, AggregateConflict, DatabaseConnectionError, DeserializationError, and
UnexpectedError with intentional SettlementVerdict classifications. Preserve the
existing fenced and transient behavior where applicable, and remove the
catch-all Err(error) arm so newly introduced variants cannot default to Rejected
silently.
Source: Coding guidelines
| mkDerivation { | ||
| pname = "event-sorcery"; | ||
| version = "0.4.0.0"; | ||
| src = bindings/haskell; |
There was a problem hiding this comment.
🎯 Functional Correctness | 🟠 Major | ⚡ Quick win
🧩 Analysis chain
🏁 Script executed:
#!/bin/sh
set -eu
# Confirm that the package function declares its source input and that flake.nix
# passes the source before evaluating the derivation.
sed -n '1,20p' nix/event-sorcery-haskell.nix
sed -n '42,51p' flake.nixRepository: dataclique/event-sorcery
Length of output: 894
🏁 Script executed:
#!/bin/sh
set -eu
printf '%s\n' '--- SPEC.md ---'
sed -n '1,240p' SPEC.md
printf '%s\n' '--- docs/domain.md ---'
sed -n '1,240p' docs/domain.md
printf '%s\n' '--- package file ---'
cat -n nix/event-sorcery-haskell.nix | sed -n '1,80p'
printf '%s\n' '--- flake excerpt ---'
cat -n flake.nix | sed -n '35,58p'
printf '%s\n' '--- related source references ---'
rg -n --glob '*.nix' 'event-sorcery-haskell|haskellBinding|overrideSrc|bindings/haskell|src\s*=' .Repository: dataclique/event-sorcery
Length of output: 25031
🏁 Script executed:
#!/bin/sh
set -eu
printf '%s\n' '--- Nix tooling ---'
command -v nix-instantiate || true
command -v nix || true
command -v nix-instantiate >/dev/null 2>&1 && {
printf '%s\n' '--- Nix expression parse ---'
nix-instantiate --parse -E 'bindings/haskell' || true
nix-instantiate --parse -E 'bindings.haskell' || true
nix-instantiate --parse -E './bindings/haskell' || true
}
printf '%s\n' '--- Nix-related documentation and package definitions ---'
rg -n --hidden --glob '!target/**' --glob '!node_modules/**' \
'callPackage|overrideSrc|mkDerivation|bindings/haskell|src\s*=' \
README.md docs nix flake.nix 2>/dev/null || trueRepository: dataclique/event-sorcery
Length of output: 622
🌐 Web query:
Nix language path literal syntax bare path bindings/haskell and nixpkgs overrideSrc callPackage evaluation
💡 Result:
In the Nix language, path literal syntax requires at least one slash (/) to be recognized as a path; otherwise, a bare identifier is interpreted as a variable lookup or attribute selection [1][2]. For example,./script.sh is a path literal, while script.sh is not [2]. Regarding package customization: 1. callPackage: This function is the standard way to compose packages in Nixpkgs [3]. It automatically supplies missing dependencies to a package's function by pulling them from the provided attribute set (typically pkgs) [3][4][5]. 2. override vs. overrideAttrs: - override: Used to override the arguments passed to the function that defines a package [6][7]. This is useful for changing parameters like configuration flags or dependencies that the package function accepts [3][5]. - overrideAttrs: Used to modify the final attributes of a derivation (e.g., src, buildPhase, postInstall) [6][7]. This is the preferred method for changing the source code (src) or build logic [8][7]. 3. Overriding src: To change the source code of a package, you typically use overrideAttrs to replace the src attribute [8][9]. When doing this, ensure you provide a valid fetcher (e.g., fetchFromGitHub) and the correct hash for the new source [8][9]. Example of overriding a package source: my-package.overrideAttrs (oldAttrs: { src = fetchFromGitHub { owner = "example"; repo = "repo"; rev = "v1.2.3"; hash = "sha256-AAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAA="; }; }) For broader, reusable customizations, overlays are the recommended approach as they allow you to modify the entire package set consistently [6][8][10].
Citations:
- 1: https://nix.dev/manual/nix/2.29/language/syntax
- 2: https://github.com/NixOS/nix/blob/master/doc/manual/source/language/syntax.md
- 3: https://nix.dev/tutorials/callpackage
- 4: https://ryantm.github.io/nixpkgs/functions/library/customisation/
- 5: https://nixos-and-flakes.thiscute.world/nixpkgs/callpackage
- 6: https://nixos.org/manual/nixpkgs/unstable/
- 7: https://ryantm.github.io/nixpkgs/using/overrides/
- 8: https://wiki.nixos.org/wiki/Overlay
- 9: https://discourse.nixos.org/t/how-to-override-and-build-certain-package-from-master-main-branch/61877
- 10: https://ryantm.github.io/nixpkgs/using/overlays/
Pass src to callPackage.
nix/event-sorcery-haskell.nix evaluates src = bindings/haskell; before overrideSrc runs. Add src to the package function arguments, set src = src;, pass src = ./bindings/haskell; to callPackage, and remove overrideSrc.
📍 Affects 2 files
nix/event-sorcery-haskell.nix#L15-L15(this comment)flake.nix#L44-L49
🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.
In `@nix/event-sorcery-haskell.nix` at line 15, Update
nix/event-sorcery-haskell.nix so the package function accepts src and assigns it
directly, pass src = ./bindings/haskell; into callPackage, and remove
overrideSrc. Apply the corresponding call-site adjustment in flake.nix lines
44-49, ensuring both sites use the explicit source path and no deferred
override.
| | Resource | Limit | | ||
| | ----------------------------------------------- | ---------------: | | ||
| | Encoded request buffer | 16 MiB | | ||
| | Encoded response buffer, including error detail | 64 MiB | | ||
| | CBOR container nesting depth | 32 | | ||
| | One opaque domain payload | 1 MiB | | ||
| | Events in one commit | 1,024 | | ||
| | Items in one page or list response | 4,096 | | ||
| | UTF-8 error-detail text | 4 KiB | | ||
| | UTF-8 identifier text | 4 KiB | | ||
| | Store busy timeout | 2,147,483,647 ms | | ||
|
|
||
| An identifier is any caller-supplied name the store persists or matches on: an | ||
| aggregate type or id, a job id or kind, and a worker name. The busy-timeout | ||
| bound is what the SQLite driver accepts; a larger value is refused at open | ||
| rather than failing later as an infrastructure error. |
There was a problem hiding this comment.
🗄️ Data Integrity & Integration | 🟡 Minor | ⚡ Quick win
State the error code for a rejected busy timeout and clarify whether the identifier limit is per identifier.
Three gaps block a binding author from implementing these two new rows.
The busy-timeout row does not behave like the other rows in this table. Line 286 states that exceeding any limit returns RESOURCE_LIMIT with no partial write. The new paragraph instead says the value "is refused at open" but names no code. State which code open returns, and say whether it is RESOURCE_LIMIT or MALFORMED_INPUT.
The identifier row does not say whether 4 KiB applies to each identifier or to all identifiers in one request. Compare "One opaque domain payload", which is explicit. Use the same wording.
The RESOURCE_LIMIT detail at line 310 is [resource, observed, limit], but no section defines the stable resource names. Two bindings will emit different strings for the same limit. Add the identifier for each row in this table.
📝 Proposed wording
-| UTF-8 identifier text | 4 KiB |
+| One UTF-8 identifier text value | 4 KiB |
| Store busy timeout | 2,147,483,647 ms |
An identifier is any caller-supplied name the store persists or matches on: an
aggregate type or id, a job id or kind, and a worker name. The busy-timeout
bound is what the SQLite driver accepts; a larger value is refused at open
-rather than failing later as an infrastructure error.
+with `MALFORMED_INPUT` rather than failing later as an infrastructure error.
+
+Each row above has a stable `resource` name used in the `RESOURCE_LIMIT`
+detail. Those names are part of ABI major version 0.As per path instructions: "For topic references, suggest improvements that would make non-obvious framework behavior or pitfalls easier to discover."
🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.
In `@SPEC.md` around lines 267 - 282, Update the resource-limit table and
surrounding specification text to state that an out-of-range busy timeout is
rejected at open with MALFORMED_INPUT, not RESOURCE_LIMIT. Clarify that the 4
KiB UTF-8 identifier limit applies independently to each identifier. Define and
use stable resource names for every limit row, including busy timeout and
identifier text, so RESOURCE_LIMIT details are consistent across bindings.
Source: Path instructions
7da24c9 to
bbd11c7
Compare
superseded by the ported branch head
There was a problem hiding this comment.
Actionable comments posted: 1
Caution
Some comments are outside the diff and can’t be posted inline due to platform limitations.
⚠️ Outside diff range comments (4)
crates/event-sorcery/src/engine.rs (4)
256-269: 📐 Maintainability & Code Quality | 🔵 Trivial | ⚡ Quick winMake
with_jobsextend instead of replace.
with_jobsassignsself.jobsand discards any job already added bywith_job.with_jobappends. The two builder methods therefore have opposite semantics on the same field. No current caller mixes them, becausecommit_opaquecallswith_jobson a fresh request. If a future caller chains both, the engine drops a durable job intent silently and the commit still succeeds.Make both methods append, so the builder never loses a job intent.
♻️ Proposed change to align the two builder methods
#[must_use] /// Attaches durable jobs that must commit atomically with the events. pub(crate) fn with_jobs(mut self, jobs: Vec<EnqueueRequest>) -> Self { - self.jobs = jobs; + self.jobs.extend(jobs); self }🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the rest with a brief reason, keep changes minimal, and validate. In `@crates/event-sorcery/src/engine.rs` around lines 256 - 269, Update EventBuilder::with_jobs to append the supplied jobs to self.jobs rather than replacing existing entries, matching EventBuilder::with_job. Preserve the existing ordering and ensure chained builder calls never discard previously added job intents.
878-899: 🗄️ Data Integrity & Integration | 🟡 Minor | ⚡ Quick winReport reused job IDs as duplicate enqueues.
es_commit_with_jobaccepts a caller-providedJobId. If its job stream already exists, the atomic path maps the job-stream conflict toEngineError::OptimisticLock, then reportsES_ERR_CONFLICTfor the domain stream. The response can reportactual == expected, so retry logic can loop indefinitely. The standalone path already hasSqliteJobError::DuplicateEnqueue. Return a dedicated duplicate-enqueue error from both paths and add coverage for the FFI case.🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the rest with a brief reason, keep changes minimal, and validate. In `@crates/event-sorcery/src/engine.rs` around lines 878 - 899, Update both atomic and standalone enqueue paths to detect an existing caller-provided JobId and return the dedicated duplicate-enqueue error, rather than mapping the job-stream conflict to EngineError::OptimisticLock or ES_ERR_CONFLICT. Preserve the existing transactional behavior, and add FFI coverage through es_commit_with_job verifying reused job IDs produce the duplicate-enqueue result without retry-style actual/expected values.
501-630: 📐 Maintainability & Code Quality | 🔵 Trivial | ⚡ Quick winCollapse the four near-identical SQL statements.
This function repeats the same projection and the same
after_sequencebranch four times: twice for the byte-accountingSUMand twice for the row load. The two branches differ only by theAND sequence > ?3predicate and one bind. Any future change to the stored columns must be applied in four places, and a missed edit makes the accounting disagree with the page it guards.Bind
after_sequenceas a nullable value and use one statement per query:WHERE aggregate_type = ?1 AND aggregate_id = ?2 AND (?3 IS NULL OR sequence > ?3) ORDER BY sequence LIMIT ?4Then extract the shared column list into one
constused by both theSUMsubquery and the row query.🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the rest with a brief reason, keep changes minimal, and validate. In `@crates/event-sorcery/src/engine.rs` around lines 501 - 630, Refactor the byte-accounting and row-loading logic in the surrounding function to use one SQL statement per query instead of branching on after_sequence. Bind after_sequence as a nullable value and apply the shared “parameter is null or sequence exceeds it” predicate, adjusting placeholders and binds so both queries use the same limit position. Extract the shared stored-event column list into a single const and reuse it in the SUM subquery and StoredEventRow query.
1759-1834: 📐 Maintainability & Code Quality | 🟠 Major | ⚡ Quick winAdd coverage for the opaque commit path.
These two tests cover
CommitRequest::with_jobwell. They assert the exact error variant, the rolled-back domain stream, the empty queue, and the empty job stream.No test covers
commit_opaque. That path is new in this PR and it is the one the C ABI uses. It combines three untested behaviors:
- sequence derivation from
expected_version(lines 793-796),- opaque payload encoding into the engine envelope (line 803),
- job forwarding through
with_jobs(line 812).A regression that dropped
.with_jobs(jobs)from line 812 would still pass every test in this file, and a foreign binding would commit an event whose job never reached the queue. That is the exact invariant theOpaqueCommitRequest::with_jobdoc comment promises.
load_opaque_events_page_boundedis also untested, including the legacy-JSON fallback branch and theOpaquePayloadVersionrejection indecode_opaque_event.Add tests that assert the domain invariants rather than the encoding internals:
commit_opaquewith oneJobSeedwrites the domain event, theJobEnqueuedevent, and exactly onejob_queuerow.commit_opaquefollowed byload_opaque_events_page_boundedreturns the exact payload bytes the caller supplied.- A stored event without engine provenance decodes through the legacy JSON branch.
commit_opaqueat a staleexpected_versionreturnsOptimisticLockand writes nothing.As per coding guidelines: "All new or modified logic must have corresponding test coverage."
🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the rest with a brief reason, keep changes minimal, and validate. In `@crates/event-sorcery/src/engine.rs` around lines 1759 - 1834, Add tests covering the opaque commit path and bounded opaque loading. Extend the engine test module around commit_opaque and load_opaque_events_page_bounded to verify job forwarding writes the domain event, JobEnqueued event, and one queue row; supplied payload bytes round-trip exactly; events without engine provenance use the legacy JSON decoding path; and stale expected_version returns OptimisticLock without persisting any events or jobs.Source: Coding guidelines
🤖 Prompt for all review comments with AI agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.
Inline comments:
In `@bindings/haskell/event-sorcery.cabal`:
- Around line 55-56: Remove EventSorcery.Engine.Internal.FFI and
EventSorcery.Engine.Internal.Paging from the public exposed-modules list and add
them under other-modules in the cabal configuration. Relocate the EsBuf layout
assertion into a test-local helper, then update EngineSpec.hs to use that helper
without importing a public internal module.
---
Outside diff comments:
In `@crates/event-sorcery/src/engine.rs`:
- Around line 256-269: Update EventBuilder::with_jobs to append the supplied
jobs to self.jobs rather than replacing existing entries, matching
EventBuilder::with_job. Preserve the existing ordering and ensure chained
builder calls never discard previously added job intents.
- Around line 878-899: Update both atomic and standalone enqueue paths to detect
an existing caller-provided JobId and return the dedicated duplicate-enqueue
error, rather than mapping the job-stream conflict to
EngineError::OptimisticLock or ES_ERR_CONFLICT. Preserve the existing
transactional behavior, and add FFI coverage through es_commit_with_job
verifying reused job IDs produce the duplicate-enqueue result without
retry-style actual/expected values.
- Around line 501-630: Refactor the byte-accounting and row-loading logic in the
surrounding function to use one SQL statement per query instead of branching on
after_sequence. Bind after_sequence as a nullable value and apply the shared
“parameter is null or sequence exceeds it” predicate, adjusting placeholders and
binds so both queries use the same limit position. Extract the shared
stored-event column list into a single const and reuse it in the SUM subquery
and StoredEventRow query.
- Around line 1759-1834: Add tests covering the opaque commit path and bounded
opaque loading. Extend the engine test module around commit_opaque and
load_opaque_events_page_bounded to verify job forwarding writes the domain
event, JobEnqueued event, and one queue row; supplied payload bytes round-trip
exactly; events without engine provenance use the legacy JSON decoding path; and
stale expected_version returns OptimisticLock without persisting any events or
jobs.
🪄 Autofix
Fix all unresolved CodeRabbit comments on this PR:
- Push a commit to this branch (recommended)
- Create a new PR with the fixes
ℹ️ Review info
⚙️ Run configuration
Configuration used: Path: .coderabbit.yaml
Review profile: ASSERTIVE
Plan: Pro Plus
Run ID: fbda48ed-8338-4859-8a9b-c33db2064c43
⛔ Files ignored due to path filters (1)
Cargo.lockis excluded by!**/*.lock
📒 Files selected for processing (8)
Cargo.tomlbindings/haskell/event-sorcery.cabalbindings/haskell/src/EventSorcery/Engine.hsbindings/haskell/test/EngineSpec.hscrates/event-sorcery-ffi/Cargo.tomlcrates/event-sorcery-ffi/src/lib.rscrates/event-sorcery/src/engine.rscrates/event-sorcery/src/lib.rs
📜 Review details
⏰ Context from checks skipped due to timeout. (7)
- GitHub Check: check
- GitHub Check: clippy
- GitHub Check: examples
- GitHub Check: fmt
- GitHub Check: test
- GitHub Check: hooks
- GitHub Check: haskell
🧰 Additional context used
📓 Path-based instructions (4)
**/*
📄 CodeRabbit inference engine (AGENTS.md)
**/*: Before work, read SPEC.md and docs/domain.md; read relevant supplemental documentation before implementation.
New features must be documented in SPEC.md before implementation and must follow the hierarchy SPEC.md -> issue -> plan -> tests -> implementation.
Fix all known problems immediately, complete all tasks, and do not allow warnings or errors to pass through.
Keep a granular task list and clear completed tasks from the active list.
All new or modified logic must have corresponding test coverage.
Understand relevant documentation and source code before implementation, keep diffs small, and review the approach critically.
When changing direction or making an important undocumented architectural decision, obtain confirmation; record significant decisions as ADRs under adrs/.
Before handover, review the diff, revert unjustified changes, and check for scope creep.
Each aggregate in a consuming application must use exactly one SqliteCqrs instance constructed at startup; per-request construction is forbidden.
Never read secret or credential files such as .env*, credentials.json, *.key, *.pem, *.p12, *.pfx, or sensitive database files without explicit permission.
Never bypass, disable, suppress, or obscure quality-control mechanisms without explicit permission; fix lint and test issues at their root.
Use cargo check, cargo nextest, and cargo clippy for verification; never use cargo build unless build artifacts are required.
Files:
Cargo.tomlcrates/event-sorcery-ffi/Cargo.tomlcrates/event-sorcery/src/lib.rsbindings/haskell/test/EngineSpec.hsbindings/haskell/event-sorcery.cabalcrates/event-sorcery/src/engine.rsbindings/haskell/src/EventSorcery/Engine.hs
**/Cargo.toml
📄 CodeRabbit inference engine (AGENTS.md)
Never manually edit Cargo.toml to add dependencies; use cargo add, with workspace dependencies declared centrally and referenced by workspace=true.
Files:
Cargo.tomlcrates/event-sorcery-ffi/Cargo.toml
*
⚙️ CodeRabbit configuration file
Focus on providing constructive criticism. Whenever you see a suboptimal approach, suggest more idiomatic or robust alternative(s). Flag potential footguns. Suggest FP alternatives to mutable/imperative code. Point out architectural flaws like leaky abstractions, tight coupling, wrong level of abstraction, poor type modeling, over-abstraction, unclear domain boundaries. Code should generally be organized based on business concerns rather than technical aspects - suggest improvements if you find violations. Point out gaps in test coverage but suggest tests that are not too coupled to the implementation and actually test domain invariants and business logic
Files:
Cargo.toml
crates/**/*.rs
📄 CodeRabbit inference engine (AGENTS.md)
crates/**/*.rs: Organize code by business feature rather than technical layer; avoid catch-all modules such as types.rs, error.rs, models.rs, utils.rs, helpers.rs, and services.rs.
Never write directly to the events table; emit events through CqrsFramework::execute() or execute_with_metadata().
Use cqrs-es Services for side effects in handle() and follow the {Action}er -> {Domain}Service -> {Domain}Manager naming pattern.
Place command-execution logging in aggregate handle() methods rather than callers.
Model invalid states with enums, ADTs, newtypes, and typestate rather than relying on runtime validation.
Use domain newtypes at APIs and convert to SDK primitives inside the callee, except at cross-crate boundaries where conversion at the call site is necessary.
Keep visibility as restrictive as possible: private over pub(crate) over pub.
Use a three-group import order: external crates, workspace crates, then crate-internal imports; do not use function-level imports except enum variants.
Do not use unwrap() or expect() in production Rust code; they are permitted in #[cfg(test)] code.
Never create error variants containing opaque String values; prefer #[from], ?, #[source], and preserve error chains.
Log a warning or error before silent early returns such as let-else failures.
Never silently mask numeric failures with caps, fallback defaults, precision truncation, unwrap_or(), or unwrap_or_default(); use explicit checked conversions and errors.
Prefer functional patterns, pattern matching, combinators, type-driven design, and iterators over imperative loops unless complexity increases.
Use ASCII in identifiers, comments, log messages, and configuration keys; Unicode is preferred only in user-facing rendered output.
Do not use single-letter variables, arguments, closure parameters, or generic type parameters except an unambiguous lone type parameter or short unambiguous closure.
Every module must have a //! docstring and should order public API, private implementati...
Files:
crates/event-sorcery/src/lib.rscrates/event-sorcery/src/engine.rs
🧠 Learnings (23)
📓 Common learnings
Learnt from: 0xgleb
Repo: dataclique/event-sorcery PR: 0
File: :0-0
Timestamp: 2026-07-16T21:10:49.512Z
Learning: In `crates/event-sorcery/src/job_backend.rs`, `JobClaimHandle` is deliberately non-serializable. It contains private fencing identity, so foreign-language bindings must retain it in trusted process memory and expose only a binding-owned opaque token; accepting caller-provided serialized handles could permit forged claim state.
Learnt from: 0xgleb
Repo: dataclique/event-sorcery PR: 0
File: :0-0
Timestamp: 2026-07-16T21:45:34.697Z
Learning: In the Rust job ABI claim-token flow, a claim planner result of `Skipped` can mean an existing worker lease remains active. Token retirement must therefore occur only when the stored claim entry's recorded lease has expired; retiring every `Skipped` result would invalidate a live worker claim capability.
📚 Learning: 2026-07-16T09:05:57.710Z
Learnt from: CR
Repo: dataclique/event-sorcery PR: 0
File: AGENTS.md:0-0
Timestamp: 2026-07-16T09:05:57.710Z
Learning: Applies to **/Cargo.toml : Never manually edit Cargo.toml to add dependencies; use cargo add, with workspace dependencies declared centrally and referenced by workspace=true.
Applied to files:
Cargo.tomlcrates/event-sorcery-ffi/Cargo.toml
📚 Learning: 2026-07-16T17:19:57.719Z
Learnt from: 0xgleb
Repo: dataclique/event-sorcery PR: 26
File: crates/event-sorcery-ffi/build.rs:0-0
Timestamp: 2026-07-16T17:19:57.719Z
Learning: For the `event-sorcery-ffi` C ABI, `event_sorcery.h` is a generated build artifact defined by `SPEC.md` and must not be hand-maintained or written into the checked-in source tree during builds. `crates/event-sorcery-ffi/build.rs` should generate it under `OUT_DIR`; publishing it for consumers is deferred to a separately defined release/package destination.
Applied to files:
crates/event-sorcery-ffi/Cargo.tomlbindings/haskell/event-sorcery.cabal
📚 Learning: 2026-07-16T09:05:57.710Z
Learnt from: CR
Repo: dataclique/event-sorcery PR: 0
File: AGENTS.md:0-0
Timestamp: 2026-07-16T09:05:57.710Z
Learning: Applies to crates/**/*.rs : Use a three-group import order: external crates, workspace crates, then crate-internal imports; do not use function-level imports except enum variants.
Applied to files:
crates/event-sorcery-ffi/Cargo.toml
📚 Learning: 2026-07-16T15:10:47.551Z
Learnt from: 0xgleb
Repo: dataclique/event-sorcery PR: 24
File: crates/event-sorcery/src/lib.rs:88-88
Timestamp: 2026-07-16T15:10:47.551Z
Learning: In `crates/event-sorcery/src/lib.rs`, the crate-root `engine` module should remain private (`mod engine;`). Rust permits descendant modules, including `crates/event-sorcery/src/job_sqlite.rs` and `crates/event-sorcery/src/sqlite_event_repository.rs`, to access it through `crate::engine`; it does not need `pub(crate)` visibility for that internal access.
Applied to files:
crates/event-sorcery-ffi/Cargo.tomlcrates/event-sorcery/src/lib.rsbindings/haskell/event-sorcery.cabalcrates/event-sorcery/src/engine.rs
📚 Learning: 2026-07-16T09:05:57.710Z
Learnt from: CR
Repo: dataclique/event-sorcery PR: 0
File: AGENTS.md:0-0
Timestamp: 2026-07-16T09:05:57.710Z
Learning: Applies to crates/**/*.rs : Do not split simple-but-long pattern matches into trivial helpers; request permission before suppressing relevant lints.
Applied to files:
crates/event-sorcery-ffi/Cargo.toml
📚 Learning: 2026-07-16T09:05:57.710Z
Learnt from: CR
Repo: dataclique/event-sorcery PR: 0
File: AGENTS.md:0-0
Timestamp: 2026-07-16T09:05:57.710Z
Learning: Applies to crates/**/*.rs : Every module must have a //! docstring and should order public API, private implementation, then tests.
Applied to files:
crates/event-sorcery-ffi/Cargo.toml
📚 Learning: 2026-07-16T09:05:57.710Z
Learnt from: CR
Repo: dataclique/event-sorcery PR: 0
File: AGENTS.md:0-0
Timestamp: 2026-07-16T09:05:57.710Z
Learning: Applies to crates/**/*.rs : Use direct struct literals and field access; avoid trivial constructors and getters, destructure newtypes instead of using .0, and inline one-expression helpers.
Applied to files:
crates/event-sorcery-ffi/Cargo.toml
📚 Learning: 2026-07-16T13:23:21.626Z
Learnt from: 0xgleb
Repo: dataclique/event-sorcery PR: 0
File: :0-0
Timestamp: 2026-07-16T13:23:21.626Z
Learning: In this Rust repository, `docs/sqlx.md` requires runtime SQLx query functions such as `sqlx::query_scalar::<_, T>(...)` for `#[cfg(test)]` code. `cargo sqlx prepare --workspace -- --all-targets` does not collect test-only macro invocations, so using SQLx compile-time query macros in tests fails under `SQLX_OFFLINE=true` without unobtainable `.sqlx` metadata.
Applied to files:
crates/event-sorcery-ffi/Cargo.toml
📚 Learning: 2026-07-16T09:05:57.710Z
Learnt from: CR
Repo: dataclique/event-sorcery PR: 0
File: AGENTS.md:0-0
Timestamp: 2026-07-16T09:05:57.710Z
Learning: Applies to crates/**/*.rs : Use in-memory SQLite pools via sqlite_es::testing::create_test_pool() for database test isolation.
Applied to files:
crates/event-sorcery-ffi/Cargo.tomlcrates/event-sorcery/src/engine.rs
📚 Learning: 2026-07-16T09:05:57.710Z
Learnt from: CR
Repo: dataclique/event-sorcery PR: 0
File: AGENTS.md:0-0
Timestamp: 2026-07-16T09:05:57.710Z
Learning: Applies to crates/**/*.rs : Do not use unwrap() or expect() in production Rust code; they are permitted in #[cfg(test)] code.
Applied to files:
crates/event-sorcery-ffi/Cargo.toml
📚 Learning: 2026-07-16T09:05:57.710Z
Learnt from: CR
Repo: dataclique/event-sorcery PR: 0
File: AGENTS.md:0-0
Timestamp: 2026-07-16T09:05:57.710Z
Learning: Applies to crates/**/*.rs : Do not use single-letter variables, arguments, closure parameters, or generic type parameters except an unambiguous lone type parameter or short unambiguous closure.
Applied to files:
crates/event-sorcery-ffi/Cargo.toml
📚 Learning: 2026-07-16T09:05:57.710Z
Learnt from: CR
Repo: dataclique/event-sorcery PR: 0
File: AGENTS.md:0-0
Timestamp: 2026-07-16T09:05:57.710Z
Learning: Applies to crates/**/*.rs : Never silently mask numeric failures with caps, fallback defaults, precision truncation, unwrap_or(), or unwrap_or_default(); use explicit checked conversions and errors.
Applied to files:
crates/event-sorcery-ffi/Cargo.toml
📚 Learning: 2026-07-16T09:05:57.710Z
Learnt from: CR
Repo: dataclique/event-sorcery PR: 0
File: AGENTS.md:0-0
Timestamp: 2026-07-16T09:05:57.710Z
Learning: Applies to crates/**/*.rs : Keep visibility as restrictive as possible: private over pub(crate) over pub.
Applied to files:
crates/event-sorcery-ffi/Cargo.toml
📚 Learning: 2026-07-16T09:05:57.710Z
Learnt from: CR
Repo: dataclique/event-sorcery PR: 0
File: AGENTS.md:0-0
Timestamp: 2026-07-16T09:05:57.710Z
Learning: Applies to crates/**/*.rs : Log a warning or error before silent early returns such as let-else failures.
Applied to files:
crates/event-sorcery-ffi/Cargo.toml
📚 Learning: 2026-07-16T09:05:57.710Z
Learnt from: CR
Repo: dataclique/event-sorcery PR: 0
File: AGENTS.md:0-0
Timestamp: 2026-07-16T09:05:57.710Z
Learning: Applies to crates/**/*.rs : Avoid deep nesting, do not nest submodules inside mod tests, and prefer early returns, let-else, and guarded pattern matches.
Applied to files:
crates/event-sorcery-ffi/Cargo.toml
📚 Learning: 2026-07-16T09:05:57.710Z
Learnt from: CR
Repo: dataclique/event-sorcery PR: 0
File: AGENTS.md:0-0
Timestamp: 2026-07-16T09:05:57.710Z
Learning: Applies to crates/**/*.rs : Prefer exhaustive match over matches! in production code; matches! is acceptable in tests.
Applied to files:
crates/event-sorcery-ffi/Cargo.toml
📚 Learning: 2026-07-16T09:05:57.710Z
Learnt from: CR
Repo: dataclique/event-sorcery PR: 0
File: AGENTS.md:0-0
Timestamp: 2026-07-16T09:05:57.710Z
Learning: Applies to crates/**/*.rs : Place command-execution logging in aggregate handle() methods rather than callers.
Applied to files:
crates/event-sorcery/src/lib.rs
📚 Learning: 2026-07-16T09:05:57.710Z
Learnt from: CR
Repo: dataclique/event-sorcery PR: 0
File: AGENTS.md:0-0
Timestamp: 2026-07-16T09:05:57.710Z
Learning: Applies to crates/**/*.rs : Tests must assert correct behavior, never document known gaps; do not delete, skip, or bypass existing tests during refactors.
Applied to files:
crates/event-sorcery/src/lib.rscrates/event-sorcery/src/engine.rs
📚 Learning: 2026-07-16T09:05:57.710Z
Learnt from: CR
Repo: dataclique/event-sorcery PR: 0
File: AGENTS.md:0-0
Timestamp: 2026-07-16T09:05:57.710Z
Learning: Applies to crates/**/*.rs : Never write directly to the events table; emit events through CqrsFramework::execute() or execute_with_metadata().
Applied to files:
crates/event-sorcery/src/lib.rscrates/event-sorcery/src/engine.rs
📚 Learning: 2026-07-16T09:05:57.710Z
Learnt from: CR
Repo: dataclique/event-sorcery PR: 0
File: AGENTS.md:0-0
Timestamp: 2026-07-16T09:05:57.710Z
Learning: Applies to crates/**/*.rs : Test business logic and meaningful edge cases rather than language features or field assignments; use proptest for clear invariants.
Applied to files:
crates/event-sorcery/src/lib.rs
📚 Learning: 2026-07-16T21:10:49.512Z
Learnt from: 0xgleb
Repo: dataclique/event-sorcery PR: 0
File: :0-0
Timestamp: 2026-07-16T21:10:49.512Z
Learning: In `crates/event-sorcery/src/job_backend.rs`, `JobClaimHandle` is deliberately non-serializable. It contains private fencing identity, so foreign-language bindings must retain it in trusted process memory and expose only a binding-owned opaque token; accepting caller-provided serialized handles could permit forged claim state.
Applied to files:
crates/event-sorcery/src/lib.rscrates/event-sorcery/src/engine.rs
📚 Learning: 2026-07-16T09:05:57.710Z
Learnt from: CR
Repo: dataclique/event-sorcery PR: 0
File: AGENTS.md:0-0
Timestamp: 2026-07-16T09:05:57.710Z
Learning: Applies to crates/**/*.rs : Serialization tests must compare output with independent literals such as json!(), not re-serialized domain values.
Applied to files:
crates/event-sorcery/src/engine.rs
🔇 Additional comments (17)
bindings/haskell/test/EngineSpec.hs (2)
81-96: Add the same-major forward-minor acceptance test.The current acceptance test depends on the linked engine version. It does not prove that
checkAbiVersionaccepts a minor version greater thanminimumAbiMinor.Add a direct case with
supportedAbiMajorandminimumAbiMinor + 1. AssertRight ().As per coding guidelines, “All new or modified logic must have corresponding test coverage.”
Source: Coding guidelines
97-286: LGTM!bindings/haskell/src/EventSorcery/Engine.hs (1)
1-360: LGTM!bindings/haskell/event-sorcery.cabal (1)
43-54: LGTM!Also applies to: 59-95
Cargo.toml (2)
46-46: 📐 Maintainability & Code QualityVerify the dependency-add procedure.
The dependency is declared in the workspace table, but the final diff cannot prove that it was added with
cargo add. Confirm that the repository procedure was followed.As per coding guidelines: “Never manually edit Cargo.toml to add dependencies; use cargo add, with workspace dependencies declared centrally and referenced by workspace=true.”
Source: Coding guidelines
59-59: LGTM!crates/event-sorcery-ffi/Cargo.toml (2)
16-16: LGTM!
27-30: 📐 Maintainability & Code QualityKeep the package lint tables. The Clippy entries match
[workspace.lints.clippy], and Cargo does not support selective inheritance with a package-specificunsafe_codeoverride.> Likely an incorrect or invalid review comment.crates/event-sorcery/src/engine.rs (7)
10-19: LGTM!
107-140: LGTM!
186-191: LGTM!Also applies to: 1087-1091
651-660: LGTM!Also applies to: 906-943
767-779: LGTM!
1291-1293: LGTM!Also applies to: 1334-1348, 1389-1418, 1437-1499, 1518-1598
1139-1158: 🚀 Performance & ScalabilityKeep the JSON byte-array envelope.
SPEC.mdretains JSON SQLite persistence and defines page limits over encoded bytes, including opaque payloads. The storage-byte accounting is therefore intentional. The C ABI already transports opaque payloads as CBOR byte strings, so base64 is not required for ABI compatibility.> Likely an incorrect or invalid review comment.crates/event-sorcery/src/lib.rs (2)
128-137: LGTM!
934-968: LGTM!
refactor: package the Haskell binding by feature
stale against the current green head
Add deterministic-CBOR C exports for atomic event-and-job commits and the existing job runtime's enqueue, poll, claim, renew, acknowledge, retry, defer, and dead-letter operations. Claim settlement stays fenced by an opaque one-shot handle and all durable transitions continue through the existing cqrs-es JobState and SQLite projection.
Includes ABI integration tests against an in-memory engine store.
Closes #67.
This is part 13 of 32 in a stack made with GitButler:
Summary by CodeRabbit