Skip to content

fix(sql): lock the stream row in check_stream - #598

Open
nmummau wants to merge 2 commits into
Eventuous:devfrom
nmummau:fix/597-check-stream-locks
Open

nmummau wants to merge 2 commits into
Eventuous:devfrom
nmummau:fix/597-check-stream-locks

Conversation

@nmummau

@nmummau nmummau commented Sep 30, 2026

Copy link
Copy Markdown
Contributor

Fixes #597

Adds WITH (UPDLOCK, HOLDLOCK) to the stream lookup in check_stream, so concurrent appends to the same stream run one after the other instead of reading the same version. Appends to different streams still run in parallel.

  • The first commit adds ConcurrentAppendTests, which fail on dev.
  • The second commit adds the fix, and they pass.

All 58 tests in Eventuous.Tests.SqlServer pass.

check_stream reads Streams.Version without lock hints. Two appends with
ExpectedStreamVersion.Any to the same stream can read the same version,
and the second append fails:

- Existing stream: duplicate key on UQ_StreamIdAndStreamPosition.
- New stream: "WrongExpectedVersion -2, stream already exists".

The failed insert also burns a GlobalPosition, which leaves a gap.

The tests use a gate connection that holds a lock. Both appends read the
stream, then wait before they write, so the race happens on every run.
The tests fail until check_stream takes UPDLOCK, HOLDLOCK.

StoreFixture exposes SchemaName so that the tests can query the schema.
@qodo-free-for-open-source-projects

Copy link
Copy Markdown
Contributor

PR Summary by Qodo

Serialize concurrent SQL Server appends to the same stream

🐞 Bug fix 🧪 Tests 🕐 20-40 Minutes

Grey Divider

AI Description

• Lock stream lookups so concurrent appends read the latest version, including when creating a
 stream.
• Add gated concurrency tests for successful appends and uninterrupted global positions.
Diagram

graph TD
  A["Concurrent appends"] --> B["check_stream"] --> C[("Streams")] --> D["Same-stream serialization"] --> E["append_events"] --> F[("Messages")]
  E -->|"Update version"| C
Loading
High-Level Assessment

The following are alternative approaches to this PR:

1. Retry conflicting appends
  • ➕ Avoids holding a stream lookup lock while writing events.
  • ➖ Rolled-back inserts can still consume global positions; callers may see or need to handle transient failures.
2. Application-level per-stream locks
  • ➕ Keeps serialization logic outside the stored procedure.
  • ➖ Requires coordination across application instances and other database writers.

Recommendation: Prefer the database lookup lock: it coordinates all writers at the stream row or missing-key range and prevents the conflicting inserts that can leave global-position gaps.

Files changed (3) +143 / -1

Bug fix (1) +3 / -1
3_CheckStream.sqlLock stream lookups during append transactions +3/-1

Lock stream lookups during append transactions

• Adds UPDLOCK and HOLDLOCK to the stream lookup. These hints make competing appends wait on an existing stream or its missing-key range before reading the version or creating the stream.

src/SqlServer/src/Eventuous.SqlServer/Scripts/3_CheckStream.sql

Tests (2) +140 / -0
ConcurrentAppendTests.csReproduce simultaneous appends to existing and new streams +138/-0

Reproduce simultaneous appends to existing and new streams

• Adds tests that gate two Any-version appends with a third SQL connection, then check that both succeed and their events are stored. A separate test checks that the appends do not leave a gap in Messages identity values.

src/SqlServer/test/Eventuous.Tests.SqlServer/Store/ConcurrentAppendTests.cs

StoreFixture.csExpose the test schema name +2/-0

Expose the test schema name

• Exposes the fixture's generated schema name so concurrency tests can issue schema-qualified SQL queries.

src/SqlServer/test/Eventuous.Tests.SqlServer/Store/StoreFixture.cs

@qodo-free-for-open-source-projects

qodo-free-for-open-source-projects Bot commented Sep 30, 2026 •

Copy link
Copy Markdown
Contributor

Code Review by Qodo

🐞 Bugs (0) 📘 Rule violations (1) 📎 Requirement gaps (0) 📜 Skill insights (0)

Grey Divider


Action required

1. New SQL tests omit context-free awaits 📘 Rule violation ⚙ Maintainability
Description
ConcurrentAppendTests awaits SQL I/O such as OpenAsync(), ExecuteNonQueryAsync(), and
RollbackAsync() without .NoContext(). These new test helpers use the context-capturing form
across connection setup, gated appends, and database queries, unlike the repository's stated async
convention.
Code

src/SqlServer/test/Eventuous.Tests.SqlServer/Store/ConcurrentAppendTests.cs[R120-122]

+        await using var creator = new SqlConnection(_fixture.Container.GetConnectionString());
+        await creator.OpenAsync();
+        await using var creatorTransaction = (SqlTransaction)await creator.BeginTransactionAsync();
Evidence
The checklist requires applicable awaits to use .NoContext(), and the repository guidance
explicitly states this convention. The added test methods await multiple SQL I/O operations without
it.

Grey Divider

Context sources
Review mode: ⚖️ Balanced: This changes SQL Server transaction locking and duplicate-stream concurrency behavior, with meaningful correctness and deadlock risks, but the logic is localized enough for one careful review pass.

Grey Divider

Tip of the day
💡 Did you know, you can route each severity your way: inline, summary, both, or drop

More tips ↗ | Customize Qodo ↗ | Qodo docs ↗

Grey Divider

Qodo Logo

Comment thread src/SqlServer/src/Eventuous.SqlServer/Scripts/3_CheckStream.sql Outdated
Comment thread src/SqlServer/test/Eventuous.Tests.SqlServer/Store/ConcurrentAppendTests.cs Outdated
@github-actions

github-actions Bot commented Sep 30, 2026 •

Copy link
Copy Markdown

Test Results

   48 files  ±0     48 suites  ±0   15m 33s ⏱️ +18s
  621 tests +8    621 ✅ +8  0 💤 ±0  0 ❌ ±0 
1 234 runs  +8  1 234 ✅ +8  0 💤 ±0  0 ❌ ±0 

Results for commit 4950055. ± Comparison against base commit ed6cfd1.

This pull request removes 9 and adds 17 tests. Note that renamed tests count towards both.
Eventuous.Tests.Azure.ServiceBus.IsSerialisableByServiceBus ‑ Passes(09/21/2026 13:40:13 +00:00)
Eventuous.Tests.Azure.ServiceBus.IsSerialisableByServiceBus ‑ Passes(09/21/2026 13:40:13)
Eventuous.Tests.Azure.ServiceBus.IsSerialisableByServiceBus ‑ Passes(9a8c03f6-9bdf-4072-a43a-772bdae7bb21)
Eventuous.Tests.Subscriptions.SequenceTests ‑ ShouldReturnFirstBefore(CommitPosition { Position: 0, Sequence: 1, Timestamp: 2026-09-21T13:35:10.6224849+00:00 }, CommitPosition { Position: 0, Sequence: 2, Timestamp: 2026-09-21T13:35:10.6224849+00:00 }, CommitPosition { Position: 0, Sequence: 4, Timestamp: 2026-09-21T13:35:10.6224849+00:00 }, CommitPosition { Position: 0, Sequence: 6, Timestamp: 2026-09-21T13:35:10.6224849+00:00 }, CommitPosition { Position: 0, Sequence: 2, Timestamp: 2026-09-21T13:35:10.6224849+00:00 })
Eventuous.Tests.Subscriptions.SequenceTests ‑ ShouldReturnFirstBefore(CommitPosition { Position: 0, Sequence: 1, Timestamp: 2026-09-21T13:35:10.6224849+00:00 }, CommitPosition { Position: 0, Sequence: 2, Timestamp: 2026-09-21T13:35:10.6224849+00:00 }, CommitPosition { Position: 0, Sequence: 6, Timestamp: 2026-09-21T13:35:10.6224849+00:00 }, CommitPosition { Position: 0, Sequence: 8, Timestamp: 2026-09-21T13:35:10.6224849+00:00 }, CommitPosition { Position: 0, Sequence: 2, Timestamp: 2026-09-21T13:35:10.6224849+00:00 })
Eventuous.Tests.Subscriptions.SequenceTests ‑ ShouldReturnFirstBefore(CommitPosition { Position: 0, Sequence: 1, Timestamp: 2026-09-21T13:35:13.7108596+00:00 }, CommitPosition { Position: 0, Sequence: 2, Timestamp: 2026-09-21T13:35:13.7108596+00:00 }, CommitPosition { Position: 0, Sequence: 4, Timestamp: 2026-09-21T13:35:13.7108596+00:00 }, CommitPosition { Position: 0, Sequence: 6, Timestamp: 2026-09-21T13:35:13.7108596+00:00 }, CommitPosition { Position: 0, Sequence: 2, Timestamp: 2026-09-21T13:35:13.7108596+00:00 })
Eventuous.Tests.Subscriptions.SequenceTests ‑ ShouldReturnFirstBefore(CommitPosition { Position: 0, Sequence: 1, Timestamp: 2026-09-21T13:35:13.7108596+00:00 }, CommitPosition { Position: 0, Sequence: 2, Timestamp: 2026-09-21T13:35:13.7108596+00:00 }, CommitPosition { Position: 0, Sequence: 6, Timestamp: 2026-09-21T13:35:13.7108596+00:00 }, CommitPosition { Position: 0, Sequence: 8, Timestamp: 2026-09-21T13:35:13.7108596+00:00 }, CommitPosition { Position: 0, Sequence: 2, Timestamp: 2026-09-21T13:35:13.7108596+00:00 })
Eventuous.Tests.Subscriptions.SequenceTests ‑ ShouldReturnFirstBefore(CommitPosition { Position: 0, Sequence: 1, Timestamp: 2026-09-21T13:35:17.3347072+00:00 }, CommitPosition { Position: 0, Sequence: 2, Timestamp: 2026-09-21T13:35:17.3347072+00:00 }, CommitPosition { Position: 0, Sequence: 4, Timestamp: 2026-09-21T13:35:17.3347072+00:00 }, CommitPosition { Position: 0, Sequence: 6, Timestamp: 2026-09-21T13:35:17.3347072+00:00 }, CommitPosition { Position: 0, Sequence: 2, Timestamp: 2026-09-21T13:35:17.3347072+00:00 })
Eventuous.Tests.Subscriptions.SequenceTests ‑ ShouldReturnFirstBefore(CommitPosition { Position: 0, Sequence: 1, Timestamp: 2026-09-21T13:35:17.3347072+00:00 }, CommitPosition { Position: 0, Sequence: 2, Timestamp: 2026-09-21T13:35:17.3347072+00:00 }, CommitPosition { Position: 0, Sequence: 6, Timestamp: 2026-09-21T13:35:17.3347072+00:00 }, CommitPosition { Position: 0, Sequence: 8, Timestamp: 2026-09-21T13:35:17.3347072+00:00 }, CommitPosition { Position: 0, Sequence: 2, Timestamp: 2026-09-21T13:35:17.3347072+00:00 })
Eventuous.Tests.Azure.ServiceBus.IsSerialisableByServiceBus ‑ Passes(021707dc-e772-4616-8014-31351a28110f)
Eventuous.Tests.Azure.ServiceBus.IsSerialisableByServiceBus ‑ Passes(09/30/2026 04:00:06 +00:00)
Eventuous.Tests.Azure.ServiceBus.IsSerialisableByServiceBus ‑ Passes(09/30/2026 04:00:06)
Eventuous.Tests.SqlServer.Store.ConcurrentAppendTests ‑ ShouldAppendWithAnyWhenTwoAppendsCreateTheSameStreamTogether
Eventuous.Tests.SqlServer.Store.ConcurrentAppendTests ‑ ShouldAppendWithAnyWhenTwoAppendsCreateTheSameStreamWithDifferentCaseTogether
Eventuous.Tests.SqlServer.Store.ConcurrentAppendTests ‑ ShouldAppendWithAnyWhenTwoAppendsReadAnExistingStreamTogether
Eventuous.Tests.SqlServer.Store.ConcurrentAppendTests ‑ ShouldFailOneAppendWithNoStreamWhenTwoAppendsCreateTheSameStreamTogether
Eventuous.Tests.SqlServer.Store.ConcurrentAppendTests ‑ ShouldFailOneAppendWithoutGlobalPositionGapWhenTwoAppendsExpectTheSameVersion
Eventuous.Tests.SqlServer.Store.ConcurrentAppendTests ‑ ShouldNotLeaveGlobalPositionGapWhenTwoAppendsReadAnExistingStreamTogether
Eventuous.Tests.SqlServer.Store.ConcurrentAppendTests ‑ ShouldNotWaitToAppendToAnExistingStreamWhileAnotherStreamIsCreated
…

♻️ This comment has been updated with latest results.

Comment thread src/SqlServer/src/Eventuous.SqlServer/Scripts/3_CheckStream.sql Outdated
@qodo-free-for-open-source-projects

Copy link
Copy Markdown
Contributor

Code review by qodo was updated up to the latest commit eca31ab

@nmummau

nmummau commented Sep 30, 2026

Copy link
Copy Markdown
Contributor Author

The net9 failure in CancelledMessageTests.Handler_cancelled_by_shutdown_is_not_acknowledged_and_is_redelivered is in Core subscriptions, which this PR doesn't touch. It passes 30/30 locally on net9. It looks like a timing flake in the 5 s Wait.Until under CI load

@nmummau
nmummau force-pushed the fix/597-check-stream-locks branch 8 times, most recently from f5ff6c8 to 817bf34 Compare September 30, 2026 03:45
@qodo-free-for-open-source-projects

Copy link
Copy Markdown
Contributor

Code review by qodo was updated up to the latest commit 817bf34

…e decide who creates one (Eventuous#597)

check_stream read Streams without lock hints. Two appends to the same
stream could read the same version, so an append with
ExpectedStreamVersion.Any failed, and a version conflict failed after its
insert, which left a gap in GlobalPosition.

- An existing stream is read with UPDLOCK. A second append to the same
  stream waits for the first one to commit, then reads the new version.
  A version conflict now fails before anything is inserted.
- When two appends create the same stream, the insert into Streams runs
  with XACT_ABORT off. UQ_StreamName, with the collation of the lookup,
  makes the second insert wait and then fail with a duplicate key. The
  second append then reads the stream that the first one created, and
  gets the same expected version check.
- No range or application lock is used, so creating a stream does not
  block other streams, and names that differ only in case are one stream.

Tests:
- Two appends with the same expected version: one fails, no gap.
- Creating a stream does not block creating another stream, or
  appending to the existing stream next to it in the index.
- Two appends that create one stream through names that differ only
  in case both succeed.
- Two NoStream appends that create the same stream: exactly one fails.
- The gap test asserts that both appends succeed before it checks the
  identity.

Fixes Eventuous#597
@nmummau
nmummau force-pushed the fix/597-check-stream-locks branch from 817bf34 to 4950055 Compare September 30, 2026 03:55
@qodo-free-for-open-source-projects

Copy link
Copy Markdown
Contributor

Code review by qodo was updated up to the latest commit 4950055

@nmummau

nmummau commented Sep 30, 2026

Copy link
Copy Markdown
Contributor Author

@alexeyzimarev this is ready for review when you have time. It fixes #597: check_stream locks an existing stream with UPDLOCK, and when two appends create the same stream, UQ_StreamName decides which one wins, so no range or application locks are needed. The first commit adds tests that fail on dev; the second commit fixes them. All 63 tests in Eventuous.Tests.SqlServer pass. The Qodo findings are addressed.

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

SQL Server: concurrent appends with ExpectedStreamVersion.Any fail because check_stream reads the stream without locks

1 participant