diff --git a/CHANGELOG.md b/CHANGELOG.md
index 6583e86..4cc6deb 100644
--- a/CHANGELOG.md
+++ b/CHANGELOG.md
@@ -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
diff --git a/src/FSharp.Control.R3/AsyncObservable.fs b/src/FSharp.Control.R3/AsyncObservable.fs
index 81751a5..cb25a62 100644
--- a/src/FSharp.Control.R3/AsyncObservable.fs
+++ b/src/FSharp.Control.R3/AsyncObservable.fs
@@ -65,8 +65,10 @@ module Observable =
|> Async.AwaitTask
}
- /// Maps the given observable with the given asynchronous function
+ /// Maps the given observable with the given asynchronous function
+ /// Thrown when the concurrency limit of the options is 0 or below -1.
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,
@@ -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.
///
+ /// Thrown when the concurrency limit of the options is 0 or below -1.
let iterAsync options (action : 't -> Async) source = source |> mapAsync options action |> length |> Async.Ignore
[]
diff --git a/src/FSharp.Control.R3/ProcessingOptions.fs b/src/FSharp.Control.R3/ProcessingOptions.fs
index a868b18..cd4ca69 100644
--- a/src/FSharp.Control.R3/ProcessingOptions.fs
+++ b/src/FSharp.Control.R3/ProcessingOptions.fs
@@ -21,11 +21,11 @@ type AwaitOperationConfiguration =
| AwaitSwitch
/// All values are sent immediately to the asynchronous method.
| 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
/// All values are sent immediately to the asynchronous method, but the results are queued and passed to the next operator in order.
| 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
/// Send the first value and the last value while the asynchronous method is running.
| AwaitThrottleFirstLast
@@ -69,6 +69,20 @@ type ProcessingOptions = {
| AwaitOperationConfiguration.AwaitSequentialParallel _ -> AwaitOperation.SequentialParallel
| AwaitOperationConfiguration.AwaitThrottleFirstLast -> AwaitOperation.ThrottleFirstLast
+ ///
+ /// Throws when a parallel configuration has an invalid concurrency limit.
+ ///
+ /// 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.
+ ///
+ ///
+ 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
diff --git a/src/FSharp.Control.R3/TaskObservable.fs b/src/FSharp.Control.R3/TaskObservable.fs
index fe506a7..60a28b3 100644
--- a/src/FSharp.Control.R3/TaskObservable.fs
+++ b/src/FSharp.Control.R3/TaskObservable.fs
@@ -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
+ /// Maps the given observable with the given asynchronous function
+ /// Thrown when the concurrency limit of the options is 0 or below -1.
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,
@@ -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.
///
+ /// Thrown when the concurrency limit of the options is 0 or below -1.
let iterAsync cancellationToken options (action : CancellationToken -> 't -> Task) source =
source
|> mapAsync options action