From 1f025e748c94f400fb2a5a641cc28744f6dd2da2 Mon Sep 17 00:00:00 2001 From: XperiAndri Date: Mon, 5 Oct 2026 02:12:28 +0200 Subject: [PATCH] feat(rxquery): add `rxqueryWith` to cancel the terminal query operators The query operators that return a task (`count`, `head`, `iter`, ...) called the R3 operators without a token, so a query over a hot source could not be cancelled and kept its subscription alive. `RxQueryBuilder` now carries a `CancellationToken` that every such operator observes; `rxquery` keeps using `CancellationToken.None` and `(rxqueryWith ct) { ... }` passes a token. The `Builders` module opens `System.Threading` for the token instead of repeating the file-level `open System`. Co-Authored-By: Claude Opus 5.5 --- CHANGELOG.md | 4 ++ src/FSharp.Control.R3/Observable.fs | 66 +++++++++++++++++++++-------- 2 files changed, 52 insertions(+), 18 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index ec0ff78..63d3ec8 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -7,6 +7,10 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 ## [Unreleased] +### Added + +- `rxqueryWith cancellationToken` and `RxQueryBuilder (cancellationToken)` to cancel the query operators that return a task + ### Changed - **Breaking:** `AwaitOperationConfiguration` cases are prefixed with `Await` (`AwaitSequential`, `AwaitParallel 4`, ...) so they no longer collide with `System.Threading.Tasks.Parallel` diff --git a/src/FSharp.Control.R3/Observable.fs b/src/FSharp.Control.R3/Observable.fs index 8fb14bd..eed71cc 100644 --- a/src/FSharp.Control.R3/Observable.fs +++ b/src/FSharp.Control.R3/Observable.fs @@ -113,11 +113,29 @@ module ValueOptionExtensions = [] module Builders = - open System + open System.Threading + /// /// A reactive query builder. - /// See http://mnajder.blogspot.com/2011/09/when-reactive-framework-meets-f-30.html - type RxQueryBuilder () = + /// + /// The query operators that return an observable sequence are lazy. The ones that return a task subscribe at once, + /// complete when the deciding element arrives or the source terminates, and observe + /// : + /// cancelling it disposes the subscription and cancels the task. + /// + /// See http://mnajder.blogspot.com/2011/09/when-reactive-framework-meets-f-30.html + /// + type RxQueryBuilder + /// Creates a builder whose query operators that return a task observe the token. + (cancellationToken : CancellationToken) + = + + /// Creates a builder whose query operators that return a task cannot be cancelled. + new () = RxQueryBuilder CancellationToken.None + + /// The token observed by the query operators that return a task. + member _.CancellationToken = cancellationToken + member _.For (s : Observable<_>, body : _ -> Observable<_>) = s.SelectMany (body) [] member _.Select (s : Observable<_>, [] selector : _ -> _) = s.Select (selector) @@ -136,34 +154,38 @@ module Builders = member _.Zero () : Observable<'T> = Observable.Empty<'T>() member _.Yield (value : 'T) = Observable.Return<'T> value [] - member _.Count (s : Observable<_>) = ObservableExtensions.CountAsync (s) + member _.Count (s : Observable<_>) = ObservableExtensions.CountAsync (s, cancellationToken) [] - member _.All (s : Observable<_>, [] predicate : _ -> bool) = s.AllAsync (new Func<_, bool> (predicate)) + member _.All (s : Observable<_>, [] predicate : _ -> bool) = + s.AllAsync (new Func<_, bool> (predicate), cancellationToken) [] - member _.Contains (s : Observable<_>, key) = s.ContainsAsync (key) + member _.Contains (s : Observable<_>, key) = s.ContainsAsync (key, cancellationToken) [] member _.Distinct (s : Observable<_>) = s.Distinct () [] - member _.ExactlyOne (s : Observable<_>) = s.SingleAsync () + member _.ExactlyOne (s : Observable<_>) = s.SingleAsync (cancellationToken) [] - member _.ExactlyOneOrDefault (s : Observable<_>) = s.SingleOrDefaultAsync () + member _.ExactlyOneOrDefault (s : Observable<_>) = s.SingleOrDefaultAsync (cancellationToken = cancellationToken) [] - member _.Find (s : Observable<_>, [] predicate : _ -> bool) = s.FirstAsync (new Func<_, bool> (predicate)) + member _.Find (s : Observable<_>, [] predicate : _ -> bool) = + s.FirstAsync (new Func<_, bool> (predicate), cancellationToken) [] - member _.Head (s : Observable<_>) = s.FirstAsync () + member _.Head (s : Observable<_>) = s.FirstAsync (cancellationToken) [] - member _.HeadOrDefault (s : Observable<_>) = s.FirstOrDefaultAsync () + member _.HeadOrDefault (s : Observable<_>) = s.FirstOrDefaultAsync (cancellationToken = cancellationToken) [] - member _.Last (s : Observable<_>) = s.LastAsync () + member _.Last (s : Observable<_>) = s.LastAsync (cancellationToken) [] - member _.LastOrDefault (s : Observable<_>) = s.LastOrDefaultAsync () + member _.LastOrDefault (s : Observable<_>) = s.LastOrDefaultAsync (cancellationToken = cancellationToken) [] - member _.MaxBy (s : Observable<'a>, [] valueSelector : 'a -> 'b) = s.MaxByAsync (new Func<'a, 'b> (valueSelector)) + member _.MaxBy (s : Observable<'a>, [] valueSelector : 'a -> 'b) = + s.MaxByAsync (new Func<'a, 'b> (valueSelector), cancellationToken) [] - member _.MinBy (s : Observable<'a>, [] valueSelector : 'a -> 'b) = s.MinByAsync (new Func<'a, 'b> (valueSelector)) + member _.MinBy (s : Observable<'a>, [] valueSelector : 'a -> 'b) = + s.MinByAsync (new Func<'a, 'b> (valueSelector), cancellationToken) [] - member inline _.SumBy (s : Observable<_>, [] valueSelector : _ -> 'Value) = + member inline this.SumBy (s : Observable<_>, [] valueSelector : _ -> 'Value) = // The first element seeds the sum: seeding with Unchecked.defaultof passed null to the (+) of reference types, // while requiring a Zero member would reject types such as TimeSpan whose zero is a field s @@ -175,7 +197,8 @@ module Builders = | ValueNone -> ValueSome value | ValueSome sum -> ValueSome (sum + value) ), - new Func<_, _> (ValueOption.defaultValue Unchecked.defaultof<'Value>) + new Func<_, _> (ValueOption.defaultValue Unchecked.defaultof<'Value>), + this.CancellationToken ) [] @@ -183,6 +206,13 @@ module Builders = s1.Zip (s2, new Func<_, _, _> (resultSelector)) [] - member _.Iter (s : Observable<_>, [] selector : _ -> _) = s.ForEachAsync (selector) + member _.Iter (s : Observable<_>, [] selector : _ -> _) = s.ForEachAsync (new Action<_> (selector), cancellationToken) + /// A reactive query builder whose query operators that return a task cannot be cancelled. let rxquery = RxQueryBuilder () + + /// + /// A reactive query builder whose query operators that return a task observe . + /// Parenthesize the application in front of the query: (rxqueryWith cancellationToken) { for x in source do head }. + /// + let rxqueryWith (cancellationToken : CancellationToken) = RxQueryBuilder cancellationToken