Skip to content
Open
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
57 changes: 57 additions & 0 deletions src/FSharp.Control.R3/Observable.fs
Original file line number Diff line number Diff line change
Expand Up @@ -110,6 +110,7 @@ module ValueOptionExtensions =
static member choose f = map f >> where ValueOption.isSome >> map ValueOption.get


/// The <see cref="T:FSharp.Control.R3.Observable.BuildersModule.RxQueryBuilder"/> query builder and its instances.
[<CompilationRepresentation(CompilationRepresentationFlags.ModuleSuffix)>]
module Builders =

Expand All @@ -136,54 +137,108 @@ module Builders =
/// The token observed by the query operators that return a task.
member _.CancellationToken = cancellationToken

/// Projects every element of the source to an observable sequence and merges the results.
member _.For (s : Observable<_>, body : _ -> Observable<_>) = s.SelectMany (body)

/// Projects every element of the query.
[<CustomOperation("select", AllowIntoPattern = true)>]
member _.Select (s : Observable<_>, [<ProjectionParameter>] selector : _ -> _) = s.Select (selector)

/// Keeps the elements of the query that satisfy the predicate.
[<CustomOperation("where", MaintainsVariableSpace = true, AllowIntoPattern = true)>]
member _.Where (s : Observable<_>, [<ProjectionParameter>] predicate : _ -> bool) = s.Where (predicate)

/// Keeps the elements of the query while they satisfy the predicate and completes at the first one that does not.
[<CustomOperation("takeWhile", MaintainsVariableSpace = true, AllowIntoPattern = true)>]
member _.TakeWhile (s : Observable<_>, [<ProjectionParameter>] predicate : _ -> bool) = s.TakeWhile (predicate)

/// Keeps the first count elements of the query and then completes.
[<CustomOperation("take", MaintainsVariableSpace = true, AllowIntoPattern = true)>]
member _.Take (s : Observable<_>, count : int) = s.Take (count)

/// Bypasses the leading elements of the query that satisfy the predicate.
[<CustomOperation("skipWhile", MaintainsVariableSpace = true, AllowIntoPattern = true)>]
member _.SkipWhile (s : Observable<_>, [<ProjectionParameter>] predicate : _ -> bool) = s.SkipWhile (predicate)

/// Bypasses the first count elements of the query.
[<CustomOperation("skip", MaintainsVariableSpace = true, AllowIntoPattern = true)>]
member _.Skip (s : Observable<_>, count : int) = s.Skip (count)

// Zero and Yield emit synchronously on subscription: scheduling them on TimeProvider.System hopped every element of
// a query to the thread pool, so SelectMany merged them in arbitrary order and the results arrived late

/// An empty sequence, used for the elements skipped by an if-then expression without else.
member _.Zero () : Observable<'T> = Observable.Empty<'T>()

/// A sequence of the single yielded element.
member _.Yield (value : 'T) = Observable.Return<'T> value

/// Counts the elements of the query.
[<CustomOperation("count")>]
member _.Count (s : Observable<_>) = ObservableExtensions.CountAsync (s, cancellationToken)

/// Determines whether every element of the query satisfies the predicate; completes at the first one that does not.
[<CustomOperation("all")>]
member _.All (s : Observable<_>, [<ProjectionParameter>] predicate : _ -> bool) =
s.AllAsync (new Func<_, bool> (predicate), cancellationToken)

/// Determines whether the query contains the element; completes at the first equal one.
[<CustomOperation("contains")>]
member _.Contains (s : Observable<_>, key) = s.ContainsAsync (key, cancellationToken)

/// Removes the repeated elements of the query.
[<CustomOperation("distinct", MaintainsVariableSpace = true, AllowIntoPattern = true)>]
member _.Distinct (s : Observable<_>) = s.Distinct ()

/// Returns the only element of the query; fails when the query has no element or more than one.
[<CustomOperation("exactlyOne")>]
member _.ExactlyOne (s : Observable<_>) = s.SingleAsync (cancellationToken)

/// <summary>
/// Returns the only element of the query, or the default value when the query has no element; fails on more than one.
/// <para>The default value of a reference type is <see langword="null"/>.</para>
/// </summary>
[<CustomOperation("exactlyOneOrDefault")>]
member _.ExactlyOneOrDefault (s : Observable<_>) = s.SingleOrDefaultAsync (cancellationToken = cancellationToken)

/// Returns the first element of the query that satisfies the predicate; fails when there is none.
[<CustomOperation("find")>]
member _.Find (s : Observable<_>, [<ProjectionParameter>] predicate : _ -> bool) =
s.FirstAsync (new Func<_, bool> (predicate), cancellationToken)

/// Returns the first element of the query; fails when the query has no element.
[<CustomOperation("head")>]
member _.Head (s : Observable<_>) = s.FirstAsync (cancellationToken)

/// <summary>
/// Returns the first element of the query, or the default value when the query has no element.
/// <para>The default value of a reference type is <see langword="null"/>.</para>
/// </summary>
[<CustomOperation("headOrDefault")>]
member _.HeadOrDefault (s : Observable<_>) = s.FirstOrDefaultAsync (cancellationToken = cancellationToken)

/// Returns the last element of the query; fails when the query has no element.
[<CustomOperation("last")>]
member _.Last (s : Observable<_>) = s.LastAsync (cancellationToken)

/// <summary>
/// Returns the last element of the query, or the default value when the query has no element.
/// <para>The default value of a reference type is <see langword="null"/>.</para>
/// </summary>
[<CustomOperation("lastOrDefault")>]
member _.LastOrDefault (s : Observable<_>) = s.LastOrDefaultAsync (cancellationToken = cancellationToken)

/// Returns the element of the query with the largest projected value, the first one on ties; fails when the query has no element.
[<CustomOperation("maxBy")>]
member _.MaxBy (s : Observable<'a>, [<ProjectionParameter>] valueSelector : 'a -> 'b) =
s.MaxByAsync (new Func<'a, 'b> (valueSelector), cancellationToken)

/// Returns the element of the query with the smallest projected value, the first one on ties; fails when the query has no element.
[<CustomOperation("minBy")>]
member _.MinBy (s : Observable<'a>, [<ProjectionParameter>] valueSelector : 'a -> 'b) =
s.MinByAsync (new Func<'a, 'b> (valueSelector), cancellationToken)

/// Adds up the projected values of the query; returns the default value of the projected type when the query has no element.
[<CustomOperation("sumBy")>]
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,
Expand All @@ -201,10 +256,12 @@ module Builders =
this.CancellationToken
)

/// Pairs the elements of the query with the elements of another sequence by position; completes with the shorter one.
[<CustomOperation("zip", IsLikeZip = true)>]
member _.Zip (s1 : Observable<_>, s2 : Observable<_>, [<ProjectionParameter>] resultSelector : _ -> _) =
s1.Zip (s2, new Func<_, _, _> (resultSelector))

/// Invokes the action for every element of the query; completes when the query completes.
[<CustomOperation("iter")>]
member _.Iter (s : Observable<_>, [<ProjectionParameter>] selector : _ -> _) = s.ForEachAsync (new Action<_> (selector), cancellationToken)

Expand Down
Loading