What is a good way to wrap a MailboxProcessor into and IObservable in F#?

Viewed 92

Suppose I have a MailboxProcessor that takes AsyncReplyChannel messages an fulfills them asynchronously.

Is there an easy way to build an IObservable from a MailboxProcessor like this?

let actor = MailboxProcessor.Start (fun inbox ->
  async {
    let mutable x = 0

    while true do
      let! ch : AsyncReplyChannel<int> = inbox.Receive ()

      ch.Reply x

      x <- x + 1
  })

let obs : IObservable<int> = Observable.ofActor actor 

actor.PostAndReply (fun ch -> ch) // Fires obs too

I suppose there are some decisions to be made around hot / cold observables etc.

This is touched on here: Does MailboxProcessor just duplicate IObservable?

1 Answers

If you want to get an existing MailboxProcessor and obtain an IObservable that will get triggered whenever the mailbox processor responds to a message, then I do not think there is a way of doing this - there is no hook that would let you detect this.

The way to do this would be to define some kind of wrapper over MailboxProcessor. For example you could define NotifyingMailboxProcessor<'Msg, 'Evt> that triggers an event of type 'Evt each time a message of type 'Msg is processed. You could start with something like this:

type NotifyingMailboxProcessor<'Msg, 'Evt>(mbox:MailboxProcessor<'T>) = 
  let evt = Event<'E>()
  member x.OnPostAndReply = evt.Publish
  member x.PostAndReply<'R>(f:AsyncReplyChannel<'R> -> 'T, g:'T -> 'R -> 'E) = 
    let mutable msg = Unchecked.defaultof<_>
    let res = mbox.PostAndReply(fun ch -> msg <- f(ch); msg)
    evt.Trigger(g msg res)
    res

type NotifyingMailboxProcessor =
  static member Start<'Msg, 'Evt>(f) = 
    NotifyingMailboxProcessor<'Msg, 'Evt>(MailboxProcessor.Start(f))

This emulates the standard interface, so you can create one using NotifyingMailboxProcessor.Start. A subtle issue is that to wrap PostAndReply, you need to give it another function that constructs the event 'Evt from a pair consisting of the message sent to the mailbox and the reply (these can have different types, so you have to wrap the result into a type 'Evt, much like you have to wrap multiple message types into a single discriminated union, typically).

An example inspired by your motivating one would then be:

let actor = NotifyingMailboxProcessor.Start<AsyncReplyChannel<int>, int>(fun inbox ->
  async {
    let mutable x = 0
    while true do
      let! (ch : AsyncReplyChannel<int>) = inbox.Receive ()
      ch.Reply x
      x <- x + 1
  })

actor.OnPostAndReply.Add(printfn "Message: %A")
actor.PostAndReply((fun ch -> ch), (fun msg ans -> ans))
Related