Sitelet https://github.com/fsprojects/FSharp.Control.R3/pull/38/files
Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1 change: 1 addition & 0 deletions .github/copilot-instructions.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
3 changes: 2 additions & 1 deletion CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -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

Expand Down
53 changes: 44 additions & 9 deletions src/FSharp.Control.R3/AsyncObservable.fs
Original file line number Diff line number Diff line change
Expand Up @@ -136,17 +136,52 @@ module Observable =
}

/// <summary>
/// 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 <paramref name="options"/>.
/// <para>
/// Depending on the options, not every element reaches the action: <see cref="P:FSharp.Control.R3.AwaitOperationConfiguration.AwaitDrop"/>
/// and <see cref="P:FSharp.Control.R3.AwaitOperationConfiguration.AwaitThrottleFirstLast"/> skip elements.
/// The computation completes when the source and the running actions complete; with
/// <see cref="P:FSharp.Control.R3.ProcessingOptions.CancelOnCompleted"/> it completes as soon as the source completes,
/// the running actions are cancelled without being awaited, and the queued elements never reach the action.
/// </para>
/// <para>
/// The first exception raised by the action, including an <see cref="T:System.OperationCanceledException"/> 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.
/// </para>
/// </summary>
/// <remarks>
/// 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.
/// </remarks>
/// <exception cref="T:System.ArgumentOutOfRangeException">Thrown when the concurrency limit of the options is 0 or below -1.</exception>
let iterAsync options (action : 't -> Async<unit>) 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<unit>) (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 ()
}

[<AutoOpen>]
module Extensions =
Expand Down
1 change: 1 addition & 0 deletions src/FSharp.Control.R3/FSharp.Control.R3.fsproj
Original file line number Diff line number Diff line change
Expand Up @@ -21,6 +21,7 @@
<Compile Include="Observable.fs" />
<Compile Include="ObservableOption.fs" />
<Compile Include="ObservableFactories.fs" />
<Compile Include="IterationGuard.fs" />
<Compile Include="AsyncObservable.fs" />
<Compile Include="TaskObservable.fs" />
</ItemGroup>
Expand Down
53 changes: 53 additions & 0 deletions src/FSharp.Control.R3/IterationGuard.fs
Original file line number Diff line number Diff line change
@@ -0,0 +1,53 @@
namespace FSharp.Control.R3

open System
open System.Runtime.ExceptionServices
open System.Threading

/// <summary>
/// Tracks one iteration of an asynchronous action over an observable sequence: the
/// <see cref="M:FSharp.Control.R3.Task.Observable.iterAsync``1(System.Threading.CancellationToken,FSharp.Control.R3.ProcessingOptions,Microsoft.FSharp.Core.FSharpFunc{System.Threading.CancellationToken,Microsoft.FSharp.Core.FSharpFunc{``0,System.Threading.Tasks.Task{Microsoft.FSharp.Core.Unit}}},R3.Observable{``0})"/>
/// and <see cref="M:FSharp.Control.R3.Async.Observable.iterAsync``1(FSharp.Control.R3.ProcessingOptions,Microsoft.FSharp.Core.FSharpFunc{``0,Microsoft.FSharp.Control.FSharpAsync{Microsoft.FSharp.Core.Unit}},R3.Observable{``0})"/>
/// of both flavours.
/// <para>
/// 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 <see cref="M:R3.Observable`1.Subscribe(R3.Observer{`0})"/> 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 <see cref="T:System.OperationCanceledException"/>, which stops the sequential modes for good without completing them.
/// </para>
/// <para>
/// Instead the guard records the first failure, skips the remaining actions and cancels
/// <see cref="P:FSharp.Control.R3.IterationGuard.StopToken"/>, which completes the iteration through
/// <see cref="M:R3.ObservableExtensions.TakeUntil``1(R3.Observable{``0},System.Threading.CancellationToken)"/>;
/// the iteration then fails with the recorded exception.
/// </para>
/// </summary>
[<Sealed>]
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 ()

[<DefaultValue>]
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()
62 changes: 50 additions & 12 deletions src/FSharp.Control.R3/TaskObservable.fs
Original file line number Diff line number Diff line change
@@ -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

/// <remarks>Caution! All functions returning <see cref="Task"/>/<see cref="Task`1"/> are blocking and may never return if awaited</remarks>
Expand Down Expand Up @@ -52,19 +53,56 @@ module Observable =
)

/// <summary>
/// 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 <paramref name="options"/>.
/// <para>
/// Depending on the options, not every element reaches the action: <see cref="P:FSharp.Control.R3.AwaitOperationConfiguration.AwaitDrop"/>
/// and <see cref="P:FSharp.Control.R3.AwaitOperationConfiguration.AwaitThrottleFirstLast"/> skip elements.
/// The task completes when the source and the running actions complete; with
/// <see cref="P:FSharp.Control.R3.ProcessingOptions.CancelOnCompleted"/> it completes as soon as the source completes,
/// the running actions are cancelled without being awaited, and the queued elements never reach the action.
/// </para>
/// <para>
/// The first exception of the action, including an <see cref="T:System.OperationCanceledException"/> 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.
/// </para>
/// </summary>
/// <remarks>
/// 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.
/// </remarks>
/// <exception cref="T:System.ArgumentOutOfRangeException">Thrown when the concurrency limit of the options is 0 or below -1.</exception>
let iterAsync cancellationToken options (action : CancellationToken -> 't -> Task<unit>) 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<unit>)
(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<unit> = 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 ()
}

[<AutoOpen>]
module Extensions =
Expand Down
Loading