From 5ce361db471d3df02661d423bdfa87a774199577 Mon Sep 17 00:00:00 2001 From: XperiAndri Date: Mon, 5 Oct 2026 02:12:57 +0200 Subject: [PATCH] fix: stop `iterAsync` reliably and never lose a failure of its action Composing R3's `SelectAwait` with `ForEachAsync` left four gaps in both flavours of `iterAsync`: - Over a synchronous source, R3 attaches the terminal to the mapped stage only after `Subscribe` returns, so the action kept running for every element after it had thrown. - With an already cancelled token the action still ran for every element of a synchronous source, although the iteration was cancelled. - R3 drops a failure of the action that happens after the source completed, so the iteration completed successfully. - `SelectAwait` swallows an `OperationCanceledException`: a timeout such as HttpClient's `TaskCanceledException` stopped the sequential worker for good and `iterAsync` never completed. The action's exceptions no longer travel through R3 at all. An internal `IterationGuard` records the first one, skips the remaining actions and cancels an internal token that completes the iteration through `TakeUntil`, which also unsubscribes from the source at once; the iteration then fails with the recorded exception. A cancellation that R3 itself caused for an invocation (`AwaitSwitch`, `CancelOnCompleted`, disposal) is not a failure, and a cancelled token cancels the iteration without subscribing. The `iterAsync` docs now describe which elements reach the action, when the iteration completes, also with `CancelOnCompleted`, and how it fails. Co-Authored-By: Claude Opus 5.5 --- .github/copilot-instructions.md | 1 + CHANGELOG.md | 3 +- src/FSharp.Control.R3/AsyncObservable.fs | 53 +++++++++++++--- .../FSharp.Control.R3.fsproj | 1 + src/FSharp.Control.R3/IterationGuard.fs | 53 ++++++++++++++++ src/FSharp.Control.R3/TaskObservable.fs | 62 +++++++++++++++---- 6 files changed, 151 insertions(+), 22 deletions(-) create mode 100644 src/FSharp.Control.R3/IterationGuard.fs diff --git a/.github/copilot-instructions.md b/.github/copilot-instructions.md index dab8a2b..5639b73 100644 --- a/.github/copilot-instructions.md +++ b/.github/copilot-instructions.md @@ -18,6 +18,7 @@ │ ├── Observable.fs – observable operators and the `rxquery` builder │ ├── ObservableOption.fs – `option` variants of the `Observable` functions (`choose`) │ ├── ObservableFactories.fs – factories reachable as `Observable.xxx` (`ofSeq`) +│ ├── IterationGuard.fs – internal guard that stops `iterAsync` and reports its failures │ ├── AsyncObservable.fs – async observable helpers │ └── TaskObservable.fs – task-based observable helpers ├── tests/FSharp.Control.R3.Tests/ – MSTest test project diff --git a/CHANGELOG.md b/CHANGELOG.md index 4dff414..a8b750c 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -31,9 +31,10 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 - `rxquery` emitted its elements on the thread pool, out of order and after the source had moved on; `yield` and `zero` are now synchronous - `rxquery` `sumBy` passed `null` to the `(+)` of reference types - A positional comparer or element selector passed to the `Async` `toLookup` was silently ignored +- `iterAsync` kept invoking the action after it failed when the source emitted synchronously, invoked it although the token was already cancelled, completed successfully when the action failed after the source completed, and never completed when the action threw an `OperationCanceledException` of its own, such as a timeout - `iterAsync` could fault with `OverflowException` on sources with more than `Int32.MaxValue` elements - `chunkBy` accepted a non-positive window length of `ChunkTimeSpanCount` and `ChunkMillisecondsCount`, which failed every element -- Misleading XML docs of `catch` +- Misleading XML docs of `catch` and `iterAsync` ## [0.3.1] - 2026-01-28 diff --git a/src/FSharp.Control.R3/AsyncObservable.fs b/src/FSharp.Control.R3/AsyncObservable.fs index 3f5f256..d119a22 100644 --- a/src/FSharp.Control.R3/AsyncObservable.fs +++ b/src/FSharp.Control.R3/AsyncObservable.fs @@ -136,17 +136,52 @@ module Observable = } /// - /// Invokes an asynchronous action for each element in the observable sequence, and propagates all observer - /// messages through the result sequence. + /// Subscribes to the source and invokes the asynchronous action for the elements, processing elements that arrive + /// while a previous invocation is running as defined by . + /// + /// Depending on the options, not every element reaches the action: + /// and skip elements. + /// The computation completes when the source and the running actions complete; with + /// it completes as soon as the source completes, + /// the running actions are cancelled without being awaited, and the queued elements never reach the action. + /// + /// + /// The first exception raised by the action, including an that the cancellation + /// of its computation did not cause, stops the processing at once, also over a source that emits synchronously, and is raised + /// by the computation. An error of the source is raised too. A computation started with an already cancelled token is cancelled + /// without subscribing. + /// /// - /// - /// This method can be used for debugging, logging, etc. of query behavior - /// by intercepting the message stream to run arbitrary actions for messages on the pipeline. - /// /// Thrown when the concurrency limit of the options is 0 or below -1. - let iterAsync options (action : 't -> Async) source = - // Waits through iter: waiting through length counted the elements with a checked add, which overflows on long-lived sources - source |> mapAsync options action |> iter ignore + let iterAsync (options : ProcessingOptions) (action : 't -> Async) (source : Observable<'t>) = + options.Validate (nameof options) + async { + // Binding the token cancels a computation started with an already cancelled token before it subscribes + let! cancellationToken = Async.CancellationToken + let guard = IterationGuard cancellationToken + let guardedAction value = async { + if not guard.IsStopped then + let! actionToken = Async.CancellationToken + try + do! action value + with + | :? OperationCanceledException when actionToken.IsCancellationRequested -> + // Defensive: when R3 cancels this invocation (switch, cancel on completion, disposal), FSharp.Core cancels the + // computation without running this handler at all; the branch only catches a cancellation that races with it, + // and keeps the guard aligned with the Task flavour, where such a cancellation does reach the handler + () + | error -> + // Not rethrown: the guard stops the iteration and reports the failure itself (see IterationGuard) + guard.Fail error + } + // Waits through iter: waiting through length counted the elements with a checked add, which overflows on long-lived sources + do! + source + |> mapAsync options guardedAction + |> _.TakeUntil(guard.StopToken) + |> iter ignore + guard.ThrowIfFailed () + } [] module Extensions = diff --git a/src/FSharp.Control.R3/FSharp.Control.R3.fsproj b/src/FSharp.Control.R3/FSharp.Control.R3.fsproj index d60e00f..589ed19 100644 --- a/src/FSharp.Control.R3/FSharp.Control.R3.fsproj +++ b/src/FSharp.Control.R3/FSharp.Control.R3.fsproj @@ -21,6 +21,7 @@ + diff --git a/src/FSharp.Control.R3/IterationGuard.fs b/src/FSharp.Control.R3/IterationGuard.fs new file mode 100644 index 0000000..1148d24 --- /dev/null +++ b/src/FSharp.Control.R3/IterationGuard.fs @@ -0,0 +1,53 @@ +namespace FSharp.Control.R3 + +open System +open System.Runtime.ExceptionServices +open System.Threading + +/// +/// Tracks one iteration of an asynchronous action over an observable sequence: the +/// +/// and +/// of both flavours. +/// +/// The failures of the action never travel through R3, because R3 1.3.1 mishandles them in three ways. It attaches the terminal +/// operator to the mapped stage only after returns, so over a synchronous +/// source the actions that are still queued keep running after a failure. It drops a failure that happens after the source completed. +/// And it swallows an , which stops the sequential modes for good without completing them. +/// +/// +/// Instead the guard records the first failure, skips the remaining actions and cancels +/// , which completes the iteration through +/// ; +/// the iteration then fails with the recorded exception. +/// +/// +[] +type internal IterationGuard (cancellationToken : CancellationToken) = + + // Never disposed: an action still running after the iteration completed may fail and cancel it late, + // and a token source without a timer or linked tokens holds nothing that needs disposal + let stop = new CancellationTokenSource () + + [] + val mutable private failure : exn | null + + /// Whether the remaining actions must be skipped because an action failed or the iteration was cancelled. + member _.IsStopped = + stop.IsCancellationRequested + || cancellationToken.IsCancellationRequested + + /// Cancelled when an action fails, to complete the iteration. + member _.StopToken = stop.Token + + /// Records the failure of an action and stops the iteration; only the first failure is kept. + member this.Fail (error : exn) = + match Interlocked.CompareExchange (&this.failure, error, null) with + | null -> stop.Cancel () + | _ -> () + + /// Raises the recorded failure with its original stack trace when an action failed. + member this.ThrowIfFailed () = + match Volatile.Read &this.failure with + | null -> () + | error -> ExceptionDispatchInfo.Capture(error).Throw() diff --git a/src/FSharp.Control.R3/TaskObservable.fs b/src/FSharp.Control.R3/TaskObservable.fs index b8237c3..122d55f 100644 --- a/src/FSharp.Control.R3/TaskObservable.fs +++ b/src/FSharp.Control.R3/TaskObservable.fs @@ -1,8 +1,9 @@ module FSharp.Control.R3.Task -open R3 +open System open System.Threading open System.Threading.Tasks +open R3 open FSharp.Control.R3 /// Caution! All functions returning / are blocking and may never return if awaited @@ -52,19 +53,56 @@ module Observable = ) /// - /// Invokes an asynchronous action for each element in the observable sequence, and propagates all observer - /// messages through the result sequence. + /// Subscribes to the source and invokes the asynchronous action for the elements, processing elements that arrive + /// while a previous invocation is running as defined by . + /// + /// Depending on the options, not every element reaches the action: + /// and skip elements. + /// The task completes when the source and the running actions complete; with + /// it completes as soon as the source completes, + /// the running actions are cancelled without being awaited, and the queued elements never reach the action. + /// + /// + /// The first exception of the action, including an that the token passed to it + /// did not cause, stops the processing at once, also over a source that emits synchronously, and the task fails with that exception. + /// An error of the source faults the task too. An already cancelled token cancels the task without subscribing. + /// /// - /// - /// This method can be used for debugging, logging, etc. of query behavior - /// by intercepting the message stream to run arbitrary actions for messages on the pipeline. - /// /// Thrown when the concurrency limit of the options is 0 or below -1. - let iterAsync cancellationToken options (action : CancellationToken -> 't -> Task) source = - // Waits through iter: waiting through length counted the elements with a checked add, which overflows on long-lived sources - source - |> mapAsync options action - |> iter cancellationToken ignore + let iterAsync + (cancellationToken : CancellationToken) + (options : ProcessingOptions) + (action : CancellationToken -> 't -> Task) + (source : Observable<'t>) + : Task = + options.Validate (nameof options) + if cancellationToken.IsCancellationRequested then + // The guard would already skip every action; the shortcut keeps the iteration from subscribing at all + Task.FromCanceled cancellationToken + else + let guard = IterationGuard cancellationToken + let guardedAction ct value : Task = task { + if not guard.IsStopped then + try + do! action ct value + with + | :? OperationCanceledException when ct.IsCancellationRequested -> + // R3 cancelled this invocation (switch, cancel on completion, disposal), which is not a failure of the iteration + () + | error -> + // Not rethrown: the guard stops the iteration and reports the failure itself (see IterationGuard) + guard.Fail error + } + // Waits through iter: waiting through length counted the elements with a checked add, which overflows on long-lived sources + let iteration = + source + |> mapAsync options guardedAction + |> _.TakeUntil(guard.StopToken) + |> iter cancellationToken ignore + task { + do! iteration + guard.ThrowIfFailed () + } [] module Extensions =