diff --git a/src/crew.gleam b/src/crew.gleam index c452554..ff7b205 100644 --- a/src/crew.gleam +++ b/src/crew.gleam @@ -420,7 +420,7 @@ type State(work, result) { idle_workers: List(Worker(work, result)), active_workers: Dict(Pid, ActiveWorker(work, result)), callers: Dict(Pid, Caller(result)), - requests: Dict(Receiver(result), Request(result)), + channels: Dict(Receiver(result), Channel(result)), queue: Deque(Work(work, result)), ) } @@ -434,11 +434,11 @@ type ActiveWorker(work, result) { } type Caller(result) { - Caller(pid: Pid, monitor: Monitor, requests: Set(Receiver(result))) + Caller(pid: Pid, monitor: Monitor, channels: Set(Receiver(result))) } -type Request(result) { - Request(from: Pid, workers: Set(Pid), receive: Receiver(result)) +type Channel(result) { + Channel(from: Pid, workers: Set(Pid), receive: Receiver(result)) } fn init_pool(self: Subject(PoolMsg(work, result))) { @@ -453,7 +453,7 @@ fn init_pool(self: Subject(PoolMsg(work, result))) { idle_workers: [], active_workers: dict.new(), callers: dict.new(), - requests: dict.new(), + channels: dict.new(), queue: deque.new(), ) @@ -484,20 +484,20 @@ fn pool( // we do not have an active caller, maybe we still got the cancel just before? use caller <- try_(dict.get(state.callers, pid), state) - let requests = - set.fold(caller.requests, state.requests, fn(requests, receive) { - use request <- try_(dict.get(state.requests, receive), requests) + let channels = + set.fold(caller.channels, state.channels, fn(channels, receive) { + use channel <- try_(dict.get(state.channels, receive), channels) // we kill all workers currently working on things for this caller. // the WorkerStopped messages will then clean up the workers. - set.each(request.workers, process.kill) + set.each(channel.workers, process.kill) - dict.delete(requests, receive) + dict.delete(channels, receive) }) // the queue is cleaned up while consuming it if the caller is missing. let callers = dict.delete(state.callers, pid) - State(..state, callers:, requests:) + State(..state, callers:, channels:) } // Active worker down @@ -529,18 +529,18 @@ fn pool( let active_workers = dict.delete(state.active_workers, pid) - case dict.get(state.requests, receive) { - Ok(request) -> { - let request = - Request(..request, workers: set.delete(request.workers, pid)) - let requests = dict.insert(state.requests, receive, request) + case dict.get(state.channels, receive) { + Ok(channel) -> { + let channel = + Channel(..channel, workers: set.delete(channel.workers, pid)) + let channels = dict.insert(state.channels, receive, channel) - State(..state, requests:, active_workers:) + State(..state, channels:, active_workers:) |> try_dequeue_work(worker) } - // the requests got cancelled while we were still working. - // we sent a kill request while unsubscribing, so we drop the worker here. + // the channels got cancelled while we were still working. + // we sent a kill channel while unsubscribing, so we drop the worker here. Error(_) -> State(..state, active_workers:, worker_count: state.worker_count - 1) } @@ -556,46 +556,46 @@ fn pool( // caller already exited, do not queue their work. use pid <- try_(process.subject_owner(receive), state) - // get the request data or construct a new entry if needed. - let request = - dict.get(state.requests, receive) - |> result.unwrap(Request(from: pid, workers: set.new(), receive:)) + // get the channel data or construct a new entry if needed. + let channel = + dict.get(state.channels, receive) + |> result.unwrap(Channel(from: pid, workers: set.new(), receive:)) - // get the caller data and insert the new request, starting to monitor if needed. + // get the caller data and insert the new channel, starting to monitor if needed. let caller = case dict.get(state.callers, pid) { Ok(caller) -> - Caller(..caller, requests: set.insert(caller.requests, receive)) + Caller(..caller, channels: set.insert(caller.channels, receive)) Error(_) -> { let monitor = process.monitor(pid) - Caller(pid: pid, monitor:, requests: set.new() |> set.insert(receive)) + Caller(pid: pid, monitor:, channels: set.new() |> set.insert(receive)) } } - // enqueue_loop will insert the request into the state. + // enqueue_loop will insert the channel into the state. State(..state, callers: dict.insert(state.callers, pid, caller)) - |> enqueue_loop(request, work) + |> enqueue_loop(channel, work) } Cancel(receive:) -> { // if the receiver does not exist, we don't need to do anything - use request <- try_(dict.get(state.requests, receive), state) + use channel <- try_(dict.get(state.channels, receive), state) - let requests = dict.delete(state.requests, receive) + let channels = dict.delete(state.channels, receive) // cancel all running workers - set.each(request.workers, process.kill) + set.each(channel.workers, process.kill) use caller <- try_( - dict.get(state.callers, request.from), - State(..state, requests:), + dict.get(state.callers, channel.from), + State(..state, channels:), ) let caller = - Caller(..caller, requests: set.delete(caller.requests, receive)) + Caller(..caller, channels: set.delete(caller.channels, receive)) - // if this request was the last one for the caller, demonitor and remove it. - let callers = case set.is_empty(caller.requests) { + // if this channel was the last one for the caller, demonitor and remove it. + let callers = case set.is_empty(caller.channels) { True -> { process.demonitor_process(caller.monitor) dict.delete(state.callers, caller.pid) @@ -605,7 +605,7 @@ fn pool( } } - State(..state, requests:, callers:) + State(..state, channels:, callers:) } } @@ -614,10 +614,10 @@ fn pool( fn enqueue_loop( state: State(work, result), - request: Request(result), + channel: Channel(result), work: List(work), ) -> State(work, result) { - let Request(from: caller, receive:, workers:) = request + let Channel(from: caller, receive:, workers:) = channel case work, state.idle_workers { [work, ..rest], [worker, ..idle_workers] -> { @@ -629,10 +629,10 @@ fn enqueue_loop( let active_workers = dict.insert(state.active_workers, worker.pid, active_worker) - let request = Request(..request, workers: set.insert(workers, worker.pid)) + let channel = Channel(..channel, workers: set.insert(workers, worker.pid)) State(..state, active_workers:, idle_workers:) - |> enqueue_loop(request, rest) + |> enqueue_loop(channel, rest) } _, _ -> { @@ -641,9 +641,9 @@ fn enqueue_loop( deque.push_back(queue, Work(work:, caller:, receive:)) }) - let requests = dict.insert(state.requests, receive, request) + let channels = dict.insert(state.channels, receive, channel) - State(..state, queue:, requests:) + State(..state, queue:, channels:) } } } @@ -658,8 +658,8 @@ fn try_dequeue_work( fn(_) { State(..state, idle_workers: [worker, ..state.idle_workers]) }, ) - use request <- try(dict.get(state.requests, work.receive), fn(_) { - // this request got cancelled + use channel <- try(dict.get(state.channels, work.receive), fn(_) { + // this channel got cancelled try_dequeue_work(State(..state, queue:), worker) }) @@ -669,11 +669,11 @@ fn try_dequeue_work( let active_workers = dict.insert(state.active_workers, worker.pid, active_worker) - let request = - Request(..request, workers: set.insert(request.workers, worker.pid)) - let requests = dict.insert(state.requests, work.receive, request) + let channel = + Channel(..channel, workers: set.insert(channel.workers, worker.pid)) + let channels = dict.insert(state.channels, work.receive, channel) - State(..state, active_workers:, requests:, queue:) + State(..state, active_workers:, channels:, queue:) } // -- WORKER ------------------------------------------------------------------