AsyncSeq (Polyglot)¶

In [1]:
#!import ../../lib/fsharp/Notebooks.livemd
In [2]:
#!import ../../lib/fsharp/Testing.livemd
In [3]:
#r "nuget: fsharp.control.asyncseq, 3.2.1"
#r "nuget: system.reactive, 6.0.1-preview.1"
#r "nuget: system.reactive.linq, 6.0.1-preview.1"
In [4]:
#!import ../../lib/fsharp/Common.fs
In [5]:
#!import ../../lib/fsharp/Async.fs
In [6]:
#if !INTERACTIVE
open Lib
#endif
In [7]:
open Common

subscribeEvent¶

In [8]:
let inline subscribeEvent (event: IEvent<'H, 'A>) map =
    let observable = System.Reactive.Linq.Observable.FromEventPattern<'H, 'A>(event.AddHandler, event.RemoveHandler)
    System.Reactive.Linq.Observable.Select (observable, fun event -> map event.EventArgs)
    |> FSharp.Control.AsyncSeq.ofObservableBuffered
In [9]:
///- --test

type TestEvent () as self =
    member val Calls = [] with get, set
    member val Event = Event<ErrorEventHandler, ErrorEventArgs> () with get

    member _.AddCall text =
        self.Calls <- self.Calls @ [ text ]

    member _.EventInterface =
        { new IEvent<ErrorEventHandler, ErrorEventArgs> with
            member _.AddHandler handler =
                self.AddCall "AddHandler"
                self.Event.Publish.AddHandler handler

            member _.RemoveHandler handler =
                self.AddCall "RemoveHandler"
                self.Event.Publish.RemoveHandler handler

            member _.Subscribe observer =
                self.AddCall "IObservable.Subscribe"
                let disposable = self.Event.Publish.Subscribe observer
                new_disposable (fun () ->
                    self.AddCall "IObservable.Dispose"
                    disposable.Dispose ()
                )
        }

    member _.Subscribe () =
        subscribeEvent
            self.EventInterface
            (fun args ->
                let result = args.GetException () |> SpiralSm.format_exception
                self.AddCall $"TestEvent.Subscribe({result})"
                result
            )

    member _.Iter subscription =
        subscription
        |> FSharp.Control.AsyncSeq.iteriAsync (fun i error -> async {
            self.AddCall $"TestEvent.Iter({i}: {error})"
        })

    member _.WaitCall text = async {
        while self.Calls |> List.last <> text do
            do! Async.SwitchToThreadPool ()
    }
In [10]:
///- --test

let testEvent = TestEvent ()

async {
    testEvent.AddCall "1"
    let! child = testEvent.Subscribe () |> testEvent.Iter |> Async.StartChild
    do! testEvent.WaitCall "AddHandler"
    testEvent.AddCall "2"
    do! child
    testEvent.AddCall "3"
}
|> Async.runWithTimeout 300
|> _assertEqual None

testEvent.Calls
|> Seq.toList
|> _assertEqual [ "1"; "AddHandler"; "2"; "RemoveHandler" ]
11:53:49 d #1 Async.runWithTimeoutAsync / timeout: 300
<null>

["1"; "AddHandler"; "2"; "RemoveHandler"]
In [11]:
///- --test

let testEvent = TestEvent ()

async {
    testEvent.AddCall "1"
    let! child = testEvent.Subscribe () |> testEvent.Iter |> Async.StartChild
    do! testEvent.WaitCall "AddHandler"
    testEvent.AddCall "2"
    use _ = testEvent.EventInterface.Subscribe (fun args ->
        testEvent.AddCall $"testEvent.EventInterface.Subscribe({args})"
    )
    testEvent.AddCall "3"
    do! child
    testEvent.AddCall "4"
}
|> Async.runWithTimeout 300
|> _assertEqual None

testEvent.Calls
|> _assertEqual [ "1"; "AddHandler"; "2"; "IObservable.Subscribe"; "3"; "RemoveHandler"; "IObservable.Dispose" ]
11:53:49 d #2 Async.runWithTimeoutAsync / timeout: 300
<null>

["1"; "AddHandler"; "2"; "IObservable.Subscribe"; "3"; "RemoveHandler"; "IObservable.Dispose"]
In [12]:
///- --test

let testEvent = TestEvent ()

async {
    testEvent.AddCall "1"
    let! child = testEvent.Subscribe () |> testEvent.Iter |> Async.StartChild
    do! testEvent.WaitCall "AddHandler"
    testEvent.AddCall "2"
    use _ = testEvent.EventInterface.Subscribe (fun args ->
        async {
            do! testEvent.WaitCall "TestEvent.Iter(0: System.Exception: error)"
            testEvent.AddCall $"testEvent.EventInterface.Subscribe({args.GetException () |> SpiralSm.format_exception})"
        }
        |> Async.RunSynchronously
    )
    testEvent.AddCall "3"
    testEvent.Event.Trigger (null, ErrorEventArgs (Exception "error"))
    testEvent.AddCall "4"
    do! child
    testEvent.AddCall "5"
}
|> Async.runWithTimeout 300
|> _assertEqual None

testEvent.Calls
|> _assertEqual [
    "1"
    "AddHandler"
    "2"
    "IObservable.Subscribe"
    "3"
    "TestEvent.Subscribe(System.Exception: error)"
    "TestEvent.Iter(0: System.Exception: error)"
    "testEvent.EventInterface.Subscribe(System.Exception: error)"
    "4"
    "RemoveHandler"
    "IObservable.Dispose"
]
11:53:50 d #3 Async.runWithTimeoutAsync / timeout: 300
<null>

["1"; "AddHandler"; "2"; "IObservable.Subscribe"; "3"; "TestEvent.Subscribe(System.Exception: error)";
 "TestEvent.Iter(0: System.Exception: error)"; "testEvent.EventInterface.Subscribe(System.Exception: error)"; "4";
 "RemoveHandler"; "IObservable.Dispose"]

subscribeToken¶

In [13]:
let subscribeToken (token : System.Threading.CancellationToken) =
    let tcs = new System.Threading.Tasks.TaskCompletionSource ()
    System.Action tcs.SetResult |> token.Register |> ignore
    let start = System.DateTime.Now.Ticks
    FSharp.Control.AsyncSeq.unfoldAsync
        (fun () -> async {
            do! tcs.Task |> Async.AwaitTask
            return Some (System.DateTime.Now.Ticks - start, ())
        })
        ()
In [14]:
///- --test

let cts = new System.Threading.CancellationTokenSource ()

async {
    let! child =
        cts.Token
        |> subscribeToken
        |> FSharp.Control.AsyncSeq.tryFirst
        |> Async.StartChild

    do! Async.Sleep 100
    cts.Cancel ()
    return! child
}
|> Async.RunSynchronously
|> Option.get
|> fun x -> x > 900000
|> _assertEqual true
true