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