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
4 changes: 4 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -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`
Expand Down
66 changes: 48 additions & 18 deletions src/FSharp.Control.R3/Observable.fs
Original file line number Diff line number Diff line change
Expand Up @@ -113,11 +113,29 @@ module ValueOptionExtensions =
[<CompilationRepresentation(CompilationRepresentationFlags.ModuleSuffix)>]
module Builders =

open System
open System.Threading

/// <summary>
/// A reactive query builder.
/// See http://mnajder.blogspot.com/2011/09/when-reactive-framework-meets-f-30.html
type RxQueryBuilder () =
/// <para>
/// 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
/// <see cref="P:FSharp.Control.R3.Observable.BuildersModule.RxQueryBuilder.CancellationToken"/>:
/// cancelling it disposes the subscription and cancels the task.
/// </para>
/// <para>See http://mnajder.blogspot.com/2011/09/when-reactive-framework-meets-f-30.html</para>
/// </summary>
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)
[<CustomOperation("select", AllowIntoPattern = true)>]
member _.Select (s : Observable<_>, [<ProjectionParameter>] selector : _ -> _) = s.Select (selector)
Expand All @@ -136,34 +154,38 @@ module Builders =
member _.Zero () : Observable<'T> = Observable.Empty<'T>()
member _.Yield (value : 'T) = Observable.Return<'T> value
[<CustomOperation("count")>]
member _.Count (s : Observable<_>) = ObservableExtensions.CountAsync (s)
member _.Count (s : Observable<_>) = ObservableExtensions.CountAsync (s, cancellationToken)
[<CustomOperation("all")>]
member _.All (s : Observable<_>, [<ProjectionParameter>] predicate : _ -> bool) = s.AllAsync (new Func<_, bool> (predicate))
member _.All (s : Observable<_>, [<ProjectionParameter>] predicate : _ -> bool) =
s.AllAsync (new Func<_, bool> (predicate), cancellationToken)
[<CustomOperation("contains")>]
member _.Contains (s : Observable<_>, key) = s.ContainsAsync (key)
member _.Contains (s : Observable<_>, key) = s.ContainsAsync (key, cancellationToken)
[<CustomOperation("distinct", MaintainsVariableSpace = true, AllowIntoPattern = true)>]
member _.Distinct (s : Observable<_>) = s.Distinct ()
[<CustomOperation("exactlyOne")>]
member _.ExactlyOne (s : Observable<_>) = s.SingleAsync ()
member _.ExactlyOne (s : Observable<_>) = s.SingleAsync (cancellationToken)
[<CustomOperation("exactlyOneOrDefault")>]
member _.ExactlyOneOrDefault (s : Observable<_>) = s.SingleOrDefaultAsync ()
member _.ExactlyOneOrDefault (s : Observable<_>) = s.SingleOrDefaultAsync (cancellationToken = cancellationToken)
[<CustomOperation("find")>]
member _.Find (s : Observable<_>, [<ProjectionParameter>] predicate : _ -> bool) = s.FirstAsync (new Func<_, bool> (predicate))
member _.Find (s : Observable<_>, [<ProjectionParameter>] predicate : _ -> bool) =
s.FirstAsync (new Func<_, bool> (predicate), cancellationToken)
[<CustomOperation("head")>]
member _.Head (s : Observable<_>) = s.FirstAsync ()
member _.Head (s : Observable<_>) = s.FirstAsync (cancellationToken)
[<CustomOperation("headOrDefault")>]
member _.HeadOrDefault (s : Observable<_>) = s.FirstOrDefaultAsync ()
member _.HeadOrDefault (s : Observable<_>) = s.FirstOrDefaultAsync (cancellationToken = cancellationToken)
[<CustomOperation("last")>]
member _.Last (s : Observable<_>) = s.LastAsync ()
member _.Last (s : Observable<_>) = s.LastAsync (cancellationToken)
[<CustomOperation("lastOrDefault")>]
member _.LastOrDefault (s : Observable<_>) = s.LastOrDefaultAsync ()
member _.LastOrDefault (s : Observable<_>) = s.LastOrDefaultAsync (cancellationToken = cancellationToken)
[<CustomOperation("maxBy")>]
member _.MaxBy (s : Observable<'a>, [<ProjectionParameter>] valueSelector : 'a -> 'b) = s.MaxByAsync (new Func<'a, 'b> (valueSelector))
member _.MaxBy (s : Observable<'a>, [<ProjectionParameter>] valueSelector : 'a -> 'b) =
s.MaxByAsync (new Func<'a, 'b> (valueSelector), cancellationToken)
[<CustomOperation("minBy")>]
member _.MinBy (s : Observable<'a>, [<ProjectionParameter>] valueSelector : 'a -> 'b) = s.MinByAsync (new Func<'a, 'b> (valueSelector))
member _.MinBy (s : Observable<'a>, [<ProjectionParameter>] valueSelector : 'a -> 'b) =
s.MinByAsync (new Func<'a, 'b> (valueSelector), cancellationToken)

[<CustomOperation("sumBy")>]
member inline _.SumBy (s : Observable<_>, [<ProjectionParameter>] valueSelector : _ -> 'Value) =
member inline this.SumBy (s : Observable<_>, [<ProjectionParameter>] 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
Expand All @@ -175,14 +197,22 @@ 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
)

[<CustomOperation("zip", IsLikeZip = true)>]
member _.Zip (s1 : Observable<_>, s2 : Observable<_>, [<ProjectionParameter>] resultSelector : _ -> _) =
s1.Zip (s2, new Func<_, _, _> (resultSelector))

[<CustomOperation("iter")>]
member _.Iter (s : Observable<_>, [<ProjectionParameter>] selector : _ -> _) = s.ForEachAsync (selector)
member _.Iter (s : Observable<_>, [<ProjectionParameter>] selector : _ -> _) = s.ForEachAsync (new Action<_> (selector), cancellationToken)

/// A reactive query builder whose query operators that return a task cannot be cancelled.
let rxquery = RxQueryBuilder ()

/// <summary>
/// A reactive query builder whose query operators that return a task observe <paramref name="cancellationToken"/>.
/// <para>Parenthesize the application in front of the query: <c>(rxqueryWith cancellationToken) { for x in source do head }</c>.</para>
/// </summary>
let rxqueryWith (cancellationToken : CancellationToken) = RxQueryBuilder cancellationToken
Loading