fix(rxquery): emit yield/zero synchronously to keep element order - #25
Open
xperiandri wants to merge 1 commit into
Open
xperiandri wants to merge 1 commit into
xperiandri wants to merge 1 commit into
Conversation
`Yield` and `Zero` used `Observable.Return/Empty (TimeProvider.System)`, which R3 schedules on the thread pool. Every element of a query therefore hopped to another thread, `SelectMany` merged them in arbitrary order and the results arrived after the source had moved on, so `take`, `head`, `last`, `zip` and friends were nondeterministic. Both now emit synchronously on subscription. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
This was referenced Oct 5, 2026
This was referenced Oct 5, 2026
There was a problem hiding this comment.
Copilot review overview
🟡 Changes recommended
The corrected timing and ordering behavior needs regression tests within this independently merged PR.
Review effort: Balanced
Findings: 1
What changed in this PR
Fixes nondeterministic rxquery ordering by emitting yield and zero synchronously.
Changes:
- Uses immediate R3
ReturnandEmptyfactories. - Documents the fix in the changelog.
| File | Description |
|---|---|
src/FSharp.Control.R3/Observable.fs |
Makes Zero and Yield synchronous. |
CHANGELOG.md |
Records the ordering fix. |
💡 Add a code-review agent skill or configure MCP servers for context-aware, tailored reviews. Learn more in the docs.
Comment on lines
+136
to
+137
| member _.Zero () : Observable<'T> = Observable.Empty<'T>() | ||
| member _.Yield (value : 'T) = Observable.Return<'T> value |
3 of 6 tasks
This branch has not been deployed
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment
Add this suggestion to a batch that can be applied as a single commit.This suggestion is invalid because no changes were made to the code.Suggestions cannot be applied while the pull request is closed.Suggestions cannot be applied while viewing a subset of changes.Only one suggestion per line can be applied in a batch.Add this suggestion to a batch that can be applied as a single commit.Applying suggestions on deleted lines is not supported.You must change the existing code in this line in order to create a valid suggestion.Outdated suggestions cannot be applied.This suggestion has been applied or marked resolved.Suggestions cannot be applied from pending reviews.Suggestions cannot be applied on multi-line comments.Suggestions cannot be applied while the pull request is queued to merge.Suggestion cannot be applied right now. Please check back later.

Proposed Changes
Part of splitting #22 into small pull requests that can be reviewed one at a time. It is stacked on #24 (
fix/map-async-max-concurrent), so it shows only its own change; merge #24 first.YieldandZerousedObservable.Return/Empty (TimeProvider.System), which R3 schedules on the thread pool. Every element of a query therefore hopped to another thread,SelectManymerged them in arbitrary order and the results arrived after the source had moved on, sotake,head,last,zipand friends were nondeterministic. Both now emit synchronously on subscription.Types of changes
What types of changes does your code introduce to FSharp.Control.R3?
Put an
xin the boxes that applyChecklist
Put an
xin the boxes that apply. You can also fill these out after creating the PR. If you're unsure about any of them, don't hesitate to ask. We're here to help! This is simply a reminder of what we are going to look for before merging your code.Further comments
Tests (in #22):
BuilderTests.fs– rxquery for and select project every element in order, rxquery over a synchronous source completes synchronously, rxquery if-then yield skips the other elements through Zero and every operator test that relies on the order.Stack – every pull request is based on the one before it, so each shows only its own change. Merge them in this order:
AwaitOperationConfigurationcases withAwait#23 refactor!: prefixAwaitOperationConfigurationcases withAwaitmapAsyncoptions eagerly #24 fix: validate the concurrency limit ofmapAsyncoptions eagerlyyield/zerosynchronously to keep element order #25 fix(rxquery): emityield/zerosynchronously to keep element order ← this pull requestsumBywithUnchecked.defaultof#26 fix(rxquery)!: stop seedingsumBywithUnchecked.defaultofrxqueryWithto cancel the terminal query operators #27 feat(rxquery): addrxqueryWithto cancel the terminal query operatorsofSeqreachable asObservable.ofSeq#29 feat(observable)!: makeofSeqreachable asObservable.ofSeqObservable.choosetake avoptionchooser, addObservableOption#30 feat(observable)!: makeObservable.choosetake avoptionchooser, addObservableOptionchunkByBoundariesfor boundaries of any element type #31 feat(observable): addchunkByBoundariesfor boundaries of any element typecancellationTokenfromtoLookup#35 fix(async)!: remove the ignoredcancellationTokenfromtoLookupiterAsyncthroughForEachAsyncinstead ofCountAsync#37 fix: wait foriterAsyncthroughForEachAsyncinstead ofCountAsynciterAsyncreliably and never lose a failure of its action #38 fix: stopiterAsyncreliably and never lose a failure of its actionofTaskconfigureAwaittotrue#39 fix(task)!: defaultofTaskconfigureAwaittotruetoArrayandtoListinto theTask.Observablemodule #40 feat(task)!: movetoArrayandtoListinto theTask.Observablemodulebind,catchandmapAsync#42 docs: document the R3 1.3.1 limitations ofbind,catchandmapAsync🤖 Generated with Claude Code