fix(streaming): release ADO.NET messages during queue handoff - #10574
Conversation
Co-authored-by: Copilot App <223556219+Copilot@users.noreply.github.com>
There was a problem hiding this comment.
Pull request overview
This PR addresses a flaky ADO.NET streaming scenario where dequeued rows can remain hidden until the visibility timeout when a pulling agent relinquishes queue ownership before confirming. It adds explicit “release” support (keyed by message id + dequeue receipt) so a successor agent can immediately redeliver those rows, and strengthens regressions/diagnostics across SQL Server, PostgreSQL, and MySQL.
Changes:
- Add receiver-side tracking of unconfirmed dequeues and release those rows during receiver shutdown using the existing confirmation stored query contract.
- Extend SQL Server/PostgreSQL/MySQL confirmation procedures to interpret negative dequeue receipts as “release for immediate redelivery” (while preventing stale releases via receipt matching).
- Add cross-provider regression tests and improve streaming timeout diagnostics with observed-vs-expected counts.
Show a summary per file
| File | Description |
|---|---|
| test/TestInfrastructure/TestExtensions/Diagnostics/StreamingDiagnosticObserver.cs | Adds richer cancellation context when waiting for item delivery counts. |
| test/Extensions/Orleans.AdoNet.Tests/Streaming/RelationalOrleansQueriesTests.cs | Adds query-level regression for releasing by dequeue receipt. |
| test/Extensions/Orleans.AdoNet.Tests/Streaming/AdoNetQueueAdapterReceiverTests.cs | Adds receiver shutdown regression ensuring unconfirmed messages are released. |
| src/AdoNet/Shared/Storage/RelationalOrleansQueries.cs | Adds ReleaseStreamMessagesAsync and shared confirmation formatting for release vs confirm. |
| src/AdoNet/Orleans.Streaming.AdoNet/SQLServer-Streaming.sql | Adds release path to ConfirmStreamMessages when dequeue receipt is negative. |
| src/AdoNet/Orleans.Streaming.AdoNet/PostgreSQL-Streaming.sql | Adds release path to confirmation function when dequeue receipt is negative. |
| src/AdoNet/Orleans.Streaming.AdoNet/MySQL-Streaming.sql | Supports negative receipts (SIGNED), matches via ABS, and adds release update path. |
| src/AdoNet/Orleans.Streaming.AdoNet/AdoNetQueueAdapterReceiver.cs | Tracks pending receipts and releases unconfirmed rows on shutdown. |
Review details
💡 Add a code-review agent skill or configure MCP servers for context-aware, tailored reviews. Learn more in the docs.
Suppressed comments (2)
src/AdoNet/Orleans.Streaming.AdoNet/AdoNetQueueAdapterReceiver.cs:162
- The confirmation-failure log is tagged with
ClusterId, but the log message template and other receiver logs use the service id. This makes cross-silo/provider troubleshooting harder because it reports the wrong identifier.
catch (Exception ex)
{
LogConfirmationFailed(ex, clusterOptions.ClusterId, providerId, queueId, items);
throw;
test/TestInfrastructure/TestExtensions/Diagnostics/StreamingDiagnosticObserver.cs:214
- Same as above:
actualCount++is not thread-safe for aDiagnosticListener-backed observable and can misreport the observed count under concurrent event delivery. Use an atomic increment.
.OfType<StreamingEvents.ItemDelivered>()
.Where(e => MatchesSubscription(e.StreamId, e.SubscriptionId, e.StreamProvider, streamId, subscriptionId, streamProvider))
.Do(_ => actualCount++)
.Take(expectedCount)
- Files reviewed: 8/8 changed files
- Comments generated: 1
- Review effort level: Lite
Co-authored-by: Copilot App <223556219+Copilot@users.noreply.github.com>
There was a problem hiding this comment.
Review details
Suppressed comments (2)
src/AdoNet/Orleans.Streaming.AdoNet/AdoNetQueueAdapterReceiver.cs:52
- Shutdown waits a captured copy of
_outstandingTask, but then re-checks the_outstandingTaskfield and can return early even when the task has completed (eg, if the async continuation which clears the field hasn't run yet). That race would skip releasing pending messages during queue handoff, which is the scenario this PR is addressing.
RemoveExpiredPendingMessages();
if (_outstandingTask is not null || _pendingMessages.Count == 0)
{
return;
}
src/AdoNet/Orleans.Streaming.AdoNet/AdoNetQueueAdapterReceiver.cs:160
- This catch block logs
clusterOptions.ClusterIdas the ServiceId argument inLogConfirmationFailed, which will produce misleading identifiers in logs when confirmations fail.
catch (Exception ex)
{
- Files reviewed: 8/8 changed files
- Comments generated: 0 new
- Review effort level: Lite
Co-authored-by: Copilot App <223556219+Copilot@users.noreply.github.com>
|
CI appears to be exposing a PR-related regression rather than an infrastructure-only flake. In the failed run, every ADO.NET provider matrix failed across PostgreSQL, MySQL/MariaDB, and SQL Server on both Examples include the stream filtering, batching, queue-adapter, client-stream, and general streaming suites. The two macOS BVT/Functional failures look unrelated, but the consistent six-job ADO.NET failure pattern aligns with the queue handoff changes and needs investigation before retrying the full run. |
Restore the ADO.NET receiver constructor used by ActivatorUtilities and prevent later queue reads from overtaking cold-stream registration. Retain the first cache pin until subscriber cursors attach, with cleanup on removal and cache pressure. Co-authored-by: Copilot App <223556219+Copilot@users.noreply.github.com>
|
Pushed follow-up repair Root regression from run 31819110261 was the receiver constructor change: Stress testing then reproduced the original delivery flake on MySQL/net10: all 10 rows were enqueued, but the consumer received Validation on the pushed head:
Remaining uncertainty: local runs cannot reproduce GitHub-hosted runner scheduling exactly, so the newly triggered six provider CI matrices remain the authoritative merge gate. |
There was a problem hiding this comment.
Review details
Suppressed comments (2)
Previously missed (2) — in code that hasn't changed since the last review.
src/AdoNet/Orleans.Streaming.AdoNet/AdoNetQueueAdapterReceiver.cs:86
- Shutdown currently returns early if waiting for in-flight operations throws (for example due to timeout). That skips the new release step entirely, which can reintroduce the original handoff problem (messages remain invisible until visibility timeout). Consider logging the wait failure but still proceeding to best-effort release of any tracked pending receipts.
catch (Exception ex)
{
LogShutdownFault(ex, clusterOptions.ServiceId, providerId, queueId);
return;
}
src/AdoNet/Shared/Storage/RelationalOrleansQueries.cs:556
- FormatStreamConfirmations can emit negative dequeue receipts when release=true, but the comment implies the format is always positive ("1:2|..."). Clarifying this helps future maintainers understand the protocol contract.
// Builds a provider-neutral receipt list in the form "1:2|3:4|5:6".
- Files reviewed: 11/11 changed files
- Comments generated: 1
- Review effort level: Lite
Co-authored-by: Copilot App <223556219+Copilot@users.noreply.github.com>
There was a problem hiding this comment.
Review details
Suppressed comments (1)
Previously missed (1) — in code that hasn't changed since the last review.
src/Orleans.Streaming/PersistentStreams/PersistentStreamPullingAgent.cs:522
- The global cold-stream registration gate returns early before cache maintenance runs. When any RegistrationTask is in-flight, ReadFromQueue exits before calling queueCache.TryPurgeFromCache()/MessagesDeliveredAsync and before applying backpressure handling, which can delay confirmations/deletes for already-delivered batches and keep the cache artificially “stuck” until registration completes.
Consider moving this gate to just before rcvr.GetQueueMessagesAsync(...) so it pauses new dequeues without skipping purge/confirm and pressure handling for already-cached items.
// Pause all queue reads so a cold stream's first batch stays pinned until registration completes.
if (pubSubCache.Values.Any(static stream => stream.RegistrationTask is { IsCompleted: false }))
{
return false;
}
- Files reviewed: 11/11 changed files
- Comments generated: 0 new
- Review effort level: Lite
|
The fresh CI run on head |
|
The fresh CI run on head A targeted rerun reproduced both Redis matrix failures in attempt 2, while the unrelated SQL reminder failure cleared. Since the failures are deterministic across both frameworks and align with this PR's pulling-agent/cache registration changes, they need to be fixed before merging rather than retried again. |
Co-authored-by: Copilot App <223556219+Copilot@users.noreply.github.com>
|
Pushed follow-up repair \c8b9a67f2\ for the deterministic Redis resume failures. The registration read gate exposed two stale-token paths: cold-stream subscription setup and post-deactivation delivery realignment could request a token which had already left the local queue cache. Both now resume from the receiver's first available batch instead of reporting a false cache miss or dropping that batch. I also removed the stream-level registration pin bookkeeping which could outlive registration and interfere with inactivity cleanup; the existing registration-scoped pin remains held through all initial subscriber handshakes. The five \RedisStreamingResumeTests\ now pass on both |
There was a problem hiding this comment.
Review details
Suppressed comments (2)
Previously missed (2) — in code that hasn't changed since the last review.
src/AdoNet/Orleans.Streaming.AdoNet/AdoNetQueueAdapterReceiver.cs:86
- If waiting for in-flight operations exceeds the shutdown timeout, Shutdown logs and returns without attempting to release already-tracked pending messages. That can leave dequeued rows invisible until the visibility timeout, undermining the intent of releasing on handoff. Consider logging the timeout but continuing with a best-effort release of currently tracked pending messages.
catch (Exception ex)
{
LogShutdownFault(ex, clusterOptions.ServiceId, providerId, queueId);
return;
}
src/AdoNet/Orleans.Streaming.AdoNet/MySQL-Streaming.sql:435
- The MySQL confirmation procedure switches to release mode if any item has a negative dequeue receipt, then updates every row in _Batch. If _Items contains a mix of positive (confirm) and negative (release) receipts, this would incorrectly "release" rows that were meant to be confirmed. Other providers only update rows whose receipt is negative and matches. Consider restricting the UPDATE to rows with negative receipts and matching dequeue counters.
IF EXISTS (SELECT 1 FROM _ItemsTable WHERE Dequeued < 0) THEN
/* negative dequeue receipts release messages for immediate redelivery */
UPDATE OrleansStreamMessage AS M
INNER JOIN _Batch AS B
ON M.ServiceId = B.ServiceId
- Files reviewed: 11/11 changed files
- Comments generated: 0 new
- Review effort level: Lite
Fixes #10458.
ADO.NET stream rows dequeued by a pulling agent remain invisible until the visibility timeout if that agent relinquishes queue ownership before confirming them. A successor can therefore register a cold stream from a later row and miss the stream's first item within the test timeout. Sorting dequeue results in #10536 fixed within-poll ordering, but it could not recover rows still hidden by the previous owner.
This change tracks each receiver's unconfirmed dequeue receipts and releases those rows during shutdown so the successor can redeliver them immediately. Release operations match both message ID and dequeue receipt, preventing a stale receiver from changing a row already claimed by a new owner. The existing confirmation query contract carries the release marker, preserving startup compatibility with older database schemas.
The pulling agent now also prevents a later queue read from overtaking in-flight cold-stream registration and keeps the first cached batch pinned until subscriber cursors attach. This closes the related race where producer registration and subscriber notification cross, causing the cursor to begin at item 2 even though item 1 was successfully enqueued. Pins are disposed when subscriptions attach, streams are removed, or cache pressure requires progress.
The follow-up repair restores the concrete
RelationalOrleansQueriesreceiver constructor expected byActivatorUtilities, while retaining the injectable query abstraction used by deterministic lifecycle tests.Cross-provider regressions cover immediate release and stale-receipt safety for SQL Server, PostgreSQL, and MySQL. Streaming timeout failures report the stream, provider, expected item count, and observed item count.
Microsoft Reviewers: Open in CodeFlow