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 CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -10,6 +10,7 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0
### Changed

- **Breaking:** `AwaitOperationConfiguration` cases are prefixed with `Await` (`AwaitSequential`, `AwaitParallel 4`, ...) so they no longer collide with `System.Threading.Tasks.Parallel`
- `mapAsync` and `iterAsync` validate the concurrency limit of the options when they are called

## [0.3.1] - 2026-01-28

Expand Down
5 changes: 4 additions & 1 deletion src/FSharp.Control.R3/AsyncObservable.fs
Original file line number Diff line number Diff line change
Expand Up @@ -65,8 +65,10 @@ module Observable =
|> Async.AwaitTask
}

/// Maps the given observable with the given asynchronous function
/// <summary>Maps the given observable with the given asynchronous function</summary>
/// <exception cref="T:System.ArgumentOutOfRangeException">Thrown when the concurrency limit of the options is 0 or below -1.</exception>
let mapAsync (options : ProcessingOptions) (f : 't -> Async<'r>) source =
options.Validate (nameof options)
let selector x ct = ValueTask<'r>(Async.StartImmediateAsTask (f x, ct))
ObservableExtensions.SelectAwait (
source,
Expand Down Expand Up @@ -107,6 +109,7 @@ module Observable =
/// 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 = source |> mapAsync options action |> length |> Async.Ignore

[<AutoOpen>]
Expand Down
18 changes: 16 additions & 2 deletions src/FSharp.Control.R3/ProcessingOptions.fs
Original file line number Diff line number Diff line change
Expand Up @@ -21,11 +21,11 @@ type AwaitOperationConfiguration =
| AwaitSwitch
/// <summary>All values are sent immediately to the asynchronous method.</summary>
| AwaitParallel of
/// If set to -1, there is no limit.
/// Maximum number of concurrent invocations; -1 means no limit, otherwise it must be greater than 0.
MaxConcurrent : int
/// <summary>All values are sent immediately to the asynchronous method, but the results are queued and passed to the next operator in order.</summary>
| AwaitSequentialParallel of
/// If set to -1, there is no limit.
/// Maximum number of concurrent invocations; -1 means no limit, otherwise it must be greater than 0.
MaxConcurrent : int
/// <summary>Send the first value and the last value while the asynchronous method is running.</summary>
| AwaitThrottleFirstLast
Expand Down Expand Up @@ -69,6 +69,20 @@ type ProcessingOptions = {
| AwaitOperationConfiguration.AwaitSequentialParallel _ -> AwaitOperation.SequentialParallel
| AwaitOperationConfiguration.AwaitThrottleFirstLast -> AwaitOperation.ThrottleFirstLast

/// <summary>
/// Throws <see cref="T:System.ArgumentOutOfRangeException"/> when a parallel configuration has an invalid concurrency limit.
/// <para>
/// R3 validates the limit only when the mapped sequence is subscribed, far away from the code that built the options,
/// so the operators that accept the options validate them eagerly.
/// </para>
/// </summary>
member internal this.Validate (paramName : string) =
match this.AwaitOperationConfiguration with
| AwaitOperationConfiguration.AwaitParallel maxConcurrent
| AwaitOperationConfiguration.AwaitSequentialParallel maxConcurrent when maxConcurrent = 0 || maxConcurrent < -1 ->
raise (ArgumentOutOfRangeException (paramName, maxConcurrent, "MaxConcurrent must be -1 (no limit) or greater than 0."))
| _ -> ()

type ChunkConfiguration<'T> =
| ChunkCount of WindowLength : int
| ChunkTimeSpan of WindowTime : TimeSpan * TimeProvider : TimeProvider
Expand Down
5 changes: 4 additions & 1 deletion src/FSharp.Control.R3/TaskObservable.fs
Original file line number Diff line number Diff line change
Expand Up @@ -37,8 +37,10 @@ module Observable =
/// Returns the length of the observable sequence till its completion or cancellation
let length cancellationToken source = ObservableExtensions.CountAsync (source, cancellationToken)

/// Maps the given observable with the given asynchronous function
/// <summary>Maps the given observable with the given asynchronous function</summary>
/// <exception cref="T:System.ArgumentOutOfRangeException">Thrown when the concurrency limit of the options is 0 or below -1.</exception>
let mapAsync (options : ProcessingOptions) (f : CancellationToken -> 'T -> Task<'R>) source =
options.Validate (nameof options)
let selector x ct = ValueTask<'R>(f ct x)
ObservableExtensions.SelectAwait (
source,
Expand All @@ -57,6 +59,7 @@ module Observable =
/// 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 =
source
|> mapAsync options action
Expand Down
Loading