diff --git a/src/crew.gleam b/src/crew.gleam index 45cfa10..c458509 100644 --- a/src/crew.gleam +++ b/src/crew.gleam @@ -466,12 +466,16 @@ pub fn select_map( /// This function should be called when you're done using a channel obtained /// from `subscribe`. It stops any in-progress work for this channel and /// removes the channel's handlers from the selector. +/// +/// Any finished work that is still coming in is dropped. pub fn unsubscribe( channel: Channel(a), selector: Selector(msg), ) -> Selector(msg) { - actor.send(channel.pool, Unsubscribe(channel.receive)) + let unsubscribed = process.new_subject() + actor.send(channel.pool, Unsubscribe(channel.receive, unsubscribed:)) + wait_for_unsubscribe(channel, unsubscribed) process.demonitor_process(channel.monitor) selector @@ -479,6 +483,23 @@ pub fn unsubscribe( |> process.deselect(channel.receive) } +fn wait_for_unsubscribe(channel: Channel(a), unsubscribed: Subject(b)) -> Nil { + let unsubscribe_selector = + process.new_selector() + |> process.select_map(unsubscribed, fn(_) { True }) + |> process.select_specific_monitor(channel.monitor, fn(_) { True }) + |> process.select_map(channel.receive, fn(_) { False }) + + drop_messages_loop(unsubscribe_selector) +} + +fn drop_messages_loop(selector: Selector(Bool)) -> Nil { + case process.selector_receive_forever(selector) { + True -> Nil + False -> drop_messages_loop(selector) + } +} + /// Submit a single piece of work to the pool using a subscription channel. /// /// This is a lower-level function for submitting work. Results must be handled @@ -505,7 +526,7 @@ pub opaque type PoolMsg { MonitoredProcessExited(reason: process.Down) // Subscribe(receive: Subject(WorkResult)) - Unsubscribe(receive: Subject(WorkResult)) + Unsubscribe(receive: Subject(WorkResult), unsubscribed: Subject(Nil)) Enqueue(receive: Subject(WorkResult), work: List(fn() -> Any)) } @@ -655,7 +676,7 @@ fn pool(state: State, msg: PoolMsg) -> Next(State, PoolMsg) { State(..state, callers:, requests:) } - Unsubscribe(receive:) -> { + Unsubscribe(receive:, unsubscribed:) -> { use request <- try_(dict.get(state.requests, receive), state) let requests = dict.delete(state.requests, receive) @@ -679,6 +700,8 @@ fn pool(state: State, msg: PoolMsg) -> Next(State, PoolMsg) { } } + process.send(unsubscribed, Nil) + State(..state, requests:, callers:) } }