diff --git a/examples/src/download_actor.gleam b/examples/src/download_actor.gleam index 0b3cf00..5aae77c 100644 --- a/examples/src/download_actor.gleam +++ b/examples/src/download_actor.gleam @@ -14,7 +14,7 @@ pub type DownloadMsg { pub type State { State( - pool: process.Name(crew.PoolMsg), + pool: process.Name(crew.PoolMsg(String, String)), results: Subject(Result(String, process.ExitReason)), pending: Int, ) @@ -25,7 +25,7 @@ pub fn main() { // NOTE: In a real application, you would supervise both the pool and the actor! let pool_name = process.new_name("download_pool") let assert Ok(_) = - crew.new(pool_name) + crew.new(pool_name, download_file) |> crew.fixed_size(4) |> crew.start @@ -75,7 +75,7 @@ fn loop(state: State, msg: DownloadMsg) { case msg { StartDownload(url) -> { // queue a download on the worker pool using the channel. - crew.enqueue(state.pool, state.results, fn() { download_file(url) }) + crew.enqueue(state.pool, state.results, url) actor.continue(State(..state, pending: state.pending + 1)) } DownloadComplete(result) -> { diff --git a/examples/src/image_processing.gleam b/examples/src/image_processing.gleam index d4d503f..8f671f7 100644 --- a/examples/src/image_processing.gleam +++ b/examples/src/image_processing.gleam @@ -1,21 +1,21 @@ -import crew +import crew/task_pool import gleam/erlang/process import gleam/int import gleam/list pub fn main() { // Start a pool with 4 workers - let pool_name = process.new_name("image_pool") + let pool_name = process.new_name("task_pool") let assert Ok(_) = - crew.new(pool_name) - |> crew.fixed_size(4) - |> crew.start + task_pool.new(pool_name) + |> task_pool.fixed_size(4) + |> task_pool.start // Process images concurrently, but keeping the result order in sync let images = list.range(1, 20) |> list.map(fn(i) { "photo" <> int.to_string(i) <> ".jpg" }) - |> crew.parallel_map(pool_name, 6000, process_image) + |> task_pool.parallel_map(pool_name, 6000, process_image) echo images } diff --git a/src/crew.gleam b/src/crew.gleam index 2df28ed..c452554 100644 --- a/src/crew.gleam +++ b/src/crew.gleam @@ -11,23 +11,20 @@ //// //// ```gleam //// import crew -//// import gleam/otp/static_supervisor as supervisor +//// import gleam/erlang/process //// //// pub fn main() { //// // Create a pool name -//// let pool_name = process.new_name("my_crew") +//// let pool_name = process.new_name("image_downloader") //// //// // Start an unsupervised pool //// let assert Ok(_) = -//// crew.new(pool_name) +//// crew.new(pool_name, download_image) //// |> crew.fixed_size(4) //// |> crew.start //// //// // Execute work on the pool -//// let result = crew.work(pool_name, 5000, fn() { -//// // Some expensive computation -//// expensive_computation() -//// }) +//// let result = crew.run(pool_name, 5000, "tasteful-ramen.jpeg") //// } //// ``` @@ -39,7 +36,6 @@ import gleam/dict.{type Dict} import gleam/erlang/process.{ type Monitor, type Name, type Pid, type Selector, type Subject, } -import gleam/int import gleam/list import gleam/otp/actor.{type Next, type StartError} import gleam/otp/static_supervisor.{type Supervisor} @@ -58,8 +54,14 @@ import gleam/string // -- BUILDER ----------------------------------------------------------------- /// A builder for configuring a worker pool before starting it. -pub opaque type Builder { - Builder(name: Name(PoolMsg), size: Int) +pub opaque type Builder(state, work, result) { + Builder( + name: Name(PoolMsg(work, result)), + size: Int, + init: fn() -> state, + init_timeout: Int, + work: fn(state, work) -> result, + ) } /// Create a new worker pool builder with the given name. @@ -73,10 +75,56 @@ pub opaque type Builder { /// /// ```gleam /// let pool_name = process.new_name("pool") -/// let builder = crew.new(pool_name) +/// let builder = crew.new(pool_name, worker) /// ``` -pub fn new(name: Name(PoolMsg)) -> Builder { - Builder(name:, size: scheduler_count()) +pub fn new( + name: Name(PoolMsg(work, result)), + work: fn(work) -> result, +) -> Builder(Nil, work, result) { + Builder( + name:, + size: scheduler_count(), + init_timeout: 1000, + init: fn() { Nil }, + work: fn(_state, input) { work(input) }, + ) +} + +/// Create a new worker pool builder with the given name and state. +/// +/// The name is used to register the pool so that work can be sent to it. +/// +/// By default, the pool will have a number of workers equal to the number of +/// scheduler threads available on the system (typically the number of CPU cores). +pub fn new_with_state( + name: Name(PoolMsg(work, result)), + state: state, + work: fn(state, work) -> result, +) -> Builder(state, work, result) { + Builder( + name:, + size: scheduler_count(), + init_timeout: 1000, + init: fn() { state }, + work:, + ) +} + +/// Create a new worker pool builder with the given name and initialiser. +/// +/// The name is used to register the pool so that work can be sent to it. +/// The initialiser will run on the worker process before it registers itself +/// as a worker. +/// +/// By default, the pool will have a number of workers equal to the number of +/// scheduler threads available on the system (typically the number of CPU cores). +pub fn new_with_initialiser( + name: Name(PoolMsg(work, result)), + timeout init_timeout: Int, + init init: fn() -> state, + run work: fn(state, work) -> result, +) -> Builder(state, work, result) { + Builder(name:, size: scheduler_count(), init_timeout:, init:, work:) } /// Set the number of worker processes in the pool to a fixed number. @@ -89,7 +137,10 @@ pub fn new(name: Name(PoolMsg)) -> Builder { /// crew.new(pool_name) /// |> crew.fixed_size(8) // Use 8 workers regardless of CPU count /// ``` -pub fn fixed_size(builder: Builder, size: Int) -> Builder { +pub fn fixed_size( + builder: Builder(state, work, result), + size: Int, +) -> Builder(state, work, result) { Builder(..builder, size:) } @@ -112,9 +163,10 @@ pub fn fixed_size(builder: Builder, size: Int) -> Builder { /// crew.new(pool_name) /// |> crew.start /// ``` -pub fn start(builder: Builder) -> Result(Supervisor, StartError) { - let Builder(name:, size:) = builder - use result <- result.map(start_tree(name, size)) +pub fn start( + builder: Builder(state, work, result), +) -> Result(Supervisor, StartError) { + use result <- result.map(start_tree(builder)) result.data } @@ -137,18 +189,18 @@ pub fn start(builder: Builder) -> Result(Supervisor, StartError) { /// |> supervisor.add(pool_spec) /// |> supervisor.start /// ``` -pub fn supervised(builder: Builder) -> ChildSpecification(_) { - let Builder(name:, size:) = builder +pub fn supervised( + builder: Builder(state, work, result), +) -> ChildSpecification(Supervisor) { use <- supervision.supervisor - start_tree(name, size) + start_tree(builder) } fn start_tree( - name: Name(PoolMsg), - size: Int, + builder: Builder(state, work, result), ) -> Result(actor.Started(_), StartError) { use <- bool.guard( - when: size <= 0, + when: builder.size <= 0, return: Error(actor.InitFailed("pool size must be greater than zero")), ) @@ -159,22 +211,21 @@ fn start_tree( use <- supervision.worker actor.new_with_initialiser(1000, init_pool) - |> actor.named(name) + |> actor.named(builder.name) |> actor.on_message(pool) |> actor.start } let worker_spec = { use <- supervision.worker - let pid = process.spawn(fn() { worker(name) }) - Ok(actor.Started(pid, Nil)) + start_worker(builder) } let worker_supervisor_spec = { use <- supervision.supervisor worker_supervisor - |> repeat(times: size, with: static_supervisor.add(_, worker_spec)) + |> repeat(times: builder.size, with: static_supervisor.add(_, worker_spec)) |> static_supervisor.start } @@ -184,40 +235,34 @@ fn start_tree( |> static_supervisor.start() } -// -- WORK -------------------------------------------------------------------- +// -- CALL -------------------------------------------------------------------- -/// Execute a single piece of work on the pool and wait for the result. +/// Send a single piece of work to one of the workers and wait for the result. /// /// This function blocks until the work is completed or the timeout is reached. -/// The work function is executed on one of the worker processes in the pool. +/// If no worker is available, the message will be queued and sent to the next +/// free worker. /// /// ## Parameters /// - `pool` - The name of the pool to execute work on /// - `timeout` - Maximum time to wait for completion in milliseconds -/// - `work` - A function containing the work to be executed +/// - `work` - The message to send to the worker function /// /// ## Panics /// - If the pool does not complete the work within the specified timeout /// - If the pool is not running /// - If the worker crashes while executing the work -/// -/// ## Example -/// -/// ```gleam -/// let response = crew.work(pool_name, 5000, fn() { -/// httpc.send(...) -/// }) /// ``` -pub fn work( - in pool: Name(PoolMsg), +pub fn call( + in pool: Name(PoolMsg(work, result)), timeout timeout: Int, - do work: fn() -> any, -) -> any { - let assert [result] = work_many(pool, timeout, [work]) + msg work: work, +) -> result { + let assert [result] = do_call(pool, timeout, [work]) result } -/// Execute multiple pieces of work concurrently on the pool, without +/// Send multiple pieces of work concurrently to the pool workers, without /// ordering guarantees. /// /// Work is distributed among the available workers and executed concurrently. @@ -227,116 +272,26 @@ pub fn work( /// ## Parameters /// - `pool` - The name of the pool to execute work on /// - `timeout` - Maximum time to wait for all work to complete in milliseconds -/// - `work` - A list of functions containing work to be executed -/// -/// ## Panics -/// - If the pool does not complete all work within the specified timeout -/// - If the pool is not running -/// - If any worker crashes while executing work -/// -/// ## Example -/// -/// ```gleam -/// let users = crew.work_many(pool_name, 10000, [ -/// fn() { fetch_user_data(user1) }, -/// fn() { fetch_user_data(user2) }, -/// fn() { fetch_user_data(user3) }, -/// ]) -/// ``` -pub fn work_many( - in pool: Name(PoolMsg), - timeout timeout: Int, - do work: List(fn() -> any), -) -> List(any) { - list.reverse(do_work(pool, timeout, work)) -} - -/// Execute multiple pieces of work concurrently on the pool and return results -/// in submission order. -/// -/// Work functions are distributed across available workers and executed -/// concurrently, but results are reordered to match the original submission -/// order before being returned. -/// -/// ## Parameters -/// - `pool` - The name of the pool to execute work on -/// - `timeout` - Maximum time to wait for all work to complete in milliseconds -/// - `work` - A list of functions containing work to be executed +/// - `work` - A list of messages describing the work to be done /// /// ## Panics /// - If the pool does not complete all work within the specified timeout /// - If the pool is not running /// - If any worker crashes while executing work -/// -/// ## Example -/// -/// ```gleam -/// let assert [user1, user2, user3] = crew.work_ordered(pool_name, 10000, [ -/// fn() { fetch_user_data(user1) }, -/// fn() { fetch_user_data(user2) }, -/// fn() { fetch_user_data(user3) }, -/// ]) -/// ``` -pub fn work_ordered( - in pool: Name(PoolMsg), +pub fn call_parallel( + in pool: Name(PoolMsg(work, result)), timeout timeout: Int, - do work: List(fn() -> any), -) -> List(any) { - parallel_map(work, pool, timeout, fn(f) { f() }) + msg work: List(work), +) -> List(result) { + list.reverse(do_call(pool, timeout, work)) } -/// Apply a function to each element of a list concurrently using the worker pool. -/// -/// This is similar to `list.map` but executes the mapping function -/// concurrently across all workers. Results are returned in the same -/// order as the input list. -/// -/// There's a bunch of extra overhead involved with spawning a work item -/// per list element and making sure the order matches. Depending on your -/// workload it might make sense to split your list into chunks first to reduce -/// work queue pressure. -/// -/// ## Parameters -/// - `list` - The list of items to map over -/// - `pool` - The name of the pool to execute work on -/// - `timeout` - Maximum time to wait for all mappings to complete in milliseconds -/// - `fun` - The function to apply to each item -/// -/// ## Panics -/// - If the pool does not complete all mappings within the specified timeout -/// - If the pool is not running -/// - If any worker crashes while executing a mapping -/// -/// ## Example -/// -/// ```gleam -/// let user_ids = [1, 2, 3, 4, 5] -/// let users = crew.parallel_map(user_ids, pool_name, 5000, fetch_user_data) -/// ``` -pub fn parallel_map( - over list: List(a), - in pool: Name(PoolMsg), - timeout timeout: Int, - with fun: fn(a) -> b, -) -> List(b) { - // we don't want to do the multiple subject trick here since that is - // potentially O(n^2) - - let unordered: List(#(Int, b)) = - list - |> list.index_map(fn(item, index) { fn() { #(index, fun(item)) } }) - |> do_work(pool, timeout, _) - - unordered - |> list.sort(fn(a, b) { int.compare(a.0, b.0) }) - |> list.map(fn(x) { x.1 }) -} - -fn do_work( - pool: Name(PoolMsg), +@internal +pub fn do_call( + pool: Name(PoolMsg(work, result)), timeout: Int, - work: List(fn() -> any), -) -> List(any) { + work: List(work), +) -> List(result) { let timeout_end = system_time() + timeout let assert Ok(pool_pid) = process.named(pool) as "Pool is not running" @@ -377,11 +332,11 @@ fn do_work( } fn receive_loop( - work: List(fn() -> any), - selector: Selector(any), + work: List(work), + selector: Selector(result), timeout_end: Int, - state: List(any), -) -> Result(List(any), Nil) { + state: List(result), +) -> Result(List(result), Nil) { case work { [] -> Ok(state) [_, ..work] -> { @@ -395,9 +350,9 @@ fn receive_loop( } } -// -- SUBSCRIBE / UNSUBSCRIBE ------------------------------------------------- +// -- ASYNC API --------------------------------------------------------------- -/// The first time `enqueue` is called from a process, the pool starts to monitor +/// The first time `enqueue*` is called from a process, the pool starts to monitor /// that process and cancels all ongoing work in case it goes down. /// /// Sometimes it is useful to manually unsubscribe and cancel all ongoing work @@ -407,25 +362,25 @@ fn receive_loop( /// Note that finished work might still arrive on this selector after /// `cancel` got called. pub fn cancel( - pool: Name(PoolMsg), - receive: Subject(Result(a, process.ExitReason)), + pool: Name(PoolMsg(work, result)), + receive: Subject(Result(result, process.ExitReason)), ) -> Nil { - actor.send(process.named_subject(pool), Cancel(cast(receive))) + actor.send(process.named_subject(pool), Cancel(receive)) } /// Submit a single piece of work to the pool using a subscription channel. /// /// This is a lower-level function for submitting work. It is the callers /// responsibility to handle timeouts, submission order and failures. -/// Most users should prefer the `work*` functions. +/// Most users should prefer `call`. /// /// The first time you `enqueue` is called from a process, the pool sets up /// a monitor making sure work is cancelled when the process no longer exists /// to receive a result. You can clean up this monitor early by using `cancel`. pub fn enqueue( - pool: Name(PoolMsg), - receive: Subject(Result(a, process.ExitReason)), - work: fn() -> a, + pool: Name(PoolMsg(work, result)), + receive: Subject(Result(result, process.ExitReason)), + work: work, ) -> Nil { enqueue_many(pool, receive, [work]) } @@ -434,60 +389,59 @@ pub fn enqueue( /// /// This is a lower-level function for submitting work. It is the callers /// responsibility to handle timeouts, submission order and failures. -/// Most users should prefer the `work*` functions. +/// Most users should prefer the `call_parallel` function. /// /// The first time you `enqueue` is called from a process, the pool sets up /// a monitor making sure work is cancelled when the process no longer exists /// to receive a result. You can clean up this monitor early by using `cancel`. pub fn enqueue_many( - pool: Name(PoolMsg), - receive: Subject(Result(a, process.ExitReason)), - work: List(fn() -> a), + pool: Name(PoolMsg(work, result)), + receive: Subject(Result(result, process.ExitReason)), + work: List(work), ) -> Nil { - actor.send(process.named_subject(pool), Enqueue(cast(receive), cast(work))) + actor.send(process.named_subject(pool), Enqueue(receive, work)) } // -- POOL -------------------------------------------------------------------- -type Receiver = - Subject(Result(Any, process.ExitReason)) - -pub opaque type PoolMsg { - WorkerStarted(pid: Pid, send: Subject(Work)) +pub opaque type PoolMsg(work, result) { + WorkerStarted(pid: Pid, send: Subject(Work(work, result))) WorkerIdle(pid: Pid) MonitoredProcessExited(reason: process.Down) // - Enqueue(receive: Receiver, work: List(fn() -> Any)) - Cancel(receive: Receiver) + GetWorkerCount(reply_to: Subject(Int)) + Enqueue(receive: Receiver(result), work: List(work)) + Cancel(receive: Receiver(result)) } -type State { +type State(work, result) { State( - idle_workers: List(Worker), - active_workers: Dict(Pid, ActiveWorker), - callers: Dict(Pid, Caller), - requests: Dict(Receiver, Request), - queue: Deque(Work), + worker_count: Int, + idle_workers: List(Worker(work, result)), + active_workers: Dict(Pid, ActiveWorker(work, result)), + callers: Dict(Pid, Caller(result)), + requests: Dict(Receiver(result), Request(result)), + queue: Deque(Work(work, result)), ) } -type Worker { - Worker(pid: Pid, send: Subject(Work), monitor: Monitor) +type Worker(work, result) { + Worker(pid: Pid, send: Subject(Work(work, result)), monitor: Monitor) } -type ActiveWorker { - ActiveWorker(worker: Worker, work: Work) +type ActiveWorker(work, result) { + ActiveWorker(worker: Worker(work, result), work: Work(work, result)) } -type Caller { - Caller(pid: Pid, monitor: Monitor, requests: Set(Receiver)) +type Caller(result) { + Caller(pid: Pid, monitor: Monitor, requests: Set(Receiver(result))) } -type Request { - Request(from: Pid, workers: Set(Pid), receive: Receiver) +type Request(result) { + Request(from: Pid, workers: Set(Pid), receive: Receiver(result)) } -fn init_pool(self: Subject(PoolMsg)) { +fn init_pool(self: Subject(PoolMsg(work, result))) { let selector = process.new_selector() |> process.select(self) @@ -495,6 +449,7 @@ fn init_pool(self: Subject(PoolMsg)) { let state = State( + worker_count: 0, idle_workers: [], active_workers: dict.new(), callers: dict.new(), @@ -508,13 +463,17 @@ fn init_pool(self: Subject(PoolMsg)) { |> Ok } -fn pool(state: State, msg: PoolMsg) -> Next(State, PoolMsg) { +fn pool( + state: State(work, result), + msg: PoolMsg(work, result), +) -> Next(State(work, result), PoolMsg(work, result)) { let next_state = case msg { WorkerStarted(pid:, send:) -> { let monitor = process.monitor(pid) let worker = Worker(pid:, send:, monitor:) - try_dequeue_work(state, worker) + State(..state, worker_count: state.worker_count + 1) + |> try_dequeue_work(worker) } MonitoredProcessExited(process.ProcessDown(pid:, ..)) -> { @@ -543,15 +502,17 @@ fn pool(state: State, msg: PoolMsg) -> Next(State, PoolMsg) { // Active worker down Error(_), Ok(_) -> { - State(..state, active_workers: dict.delete(state.active_workers, pid)) + let active_workers = dict.delete(state.active_workers, pid) + let worker_count = state.worker_count - 1 + State(..state, active_workers:, worker_count:) } // idle worker down Error(_), Error(_) -> { let idle_workers = list.filter(state.idle_workers, fn(worker) { worker.pid != pid }) - - State(..state, idle_workers:) + let worker_count = state.worker_count - 1 + State(..state, idle_workers:, worker_count:) } } } @@ -580,10 +541,16 @@ fn pool(state: State, msg: PoolMsg) -> Next(State, PoolMsg) { // the requests got cancelled while we were still working. // we sent a kill request while unsubscribing, so we drop the worker here. - Error(_) -> State(..state, active_workers:) + Error(_) -> + State(..state, active_workers:, worker_count: state.worker_count - 1) } } + GetWorkerCount(reply_to:) -> { + process.send(reply_to, state.worker_count) + state + } + Enqueue(work: [], ..) -> state Enqueue(receive:, work:) -> { // caller already exited, do not queue their work. @@ -646,10 +613,10 @@ fn pool(state: State, msg: PoolMsg) -> Next(State, PoolMsg) { } fn enqueue_loop( - state: State, - request: Request, - work: List(fn() -> Any), -) -> State { + state: State(work, result), + request: Request(result), + work: List(work), +) -> State(work, result) { let Request(from: caller, receive:, workers:) = request case work, state.idle_workers { @@ -681,7 +648,10 @@ fn enqueue_loop( } } -fn try_dequeue_work(state: State, worker: Worker) -> State { +fn try_dequeue_work( + state: State(work, result), + worker: Worker(work, result), +) -> State(work, result) { use #(work, queue) <- try( deque.pop_front(state.queue), // the queue is empty, this worker becomes idle. @@ -708,29 +678,61 @@ fn try_dequeue_work(state: State, worker: Worker) -> State { // -- WORKER ------------------------------------------------------------------ -type Any +type Receiver(result) = + Subject(Result(result, process.ExitReason)) + +type Work(work, result) { + Work(work: work, caller: Pid, receive: Receiver(result)) +} + +fn start_worker( + builder: Builder(state, work, result), +) -> Result(actor.Started(_), actor.StartError) { + let initialised = process.new_subject() + let pid = process.spawn(fn() { worker(builder, initialised) }) + + let monitor = process.monitor(pid) -type Work { - Work(work: fn() -> Any, caller: Pid, receive: Receiver) + let selector = + process.new_selector() + |> process.select_map(initialised, Ok) + |> process.select_specific_monitor(monitor, Error) + + let result = case process.selector_receive(selector, builder.init_timeout) { + Ok(Ok(Nil)) -> Ok(actor.Started(pid, Nil)) + Ok(Error(down)) -> Error(actor.InitExited(down.reason)) + Error(Nil) -> { + process.unlink(pid) + process.kill(pid) + Error(actor.InitTimeout) + } + } + + process.demonitor_process(monitor) + result } -fn worker(pool: Name(PoolMsg)) -> Nil { +fn worker(builder: Builder(state, work, result), initialised: Subject(Nil)) { let self = process.self() let subject = process.new_subject() - let pool_subject = process.named_subject(pool) + let pool_subject = process.named_subject(builder.name) + let state = builder.init() + process.send(initialised, Nil) process.send(pool_subject, WorkerStarted(self, subject)) - worker_loop(pool_subject, self, subject) + + worker_loop(pool_subject, self, subject, state, builder.work) } -fn worker_loop(pool: Subject(PoolMsg), self: Pid, subject: Subject(Work)) { +fn worker_loop(pool, self, subject, state, do_work) -> Nil { let Work(work:, receive:, caller: _) = process.receive_forever(subject) - let result = work() + let result = do_work(state, work) + process.send(pool, WorkerIdle(self)) process.send(receive, Ok(result)) - worker_loop(pool, self, subject) + worker_loop(pool, self, subject, state, do_work) } // -- HELPERS ----------------------------------------------------------------- @@ -765,6 +767,3 @@ fn scheduler_count() -> Int @external(erlang, "crew_ffi", "system_time") fn system_time() -> Int - -@external(erlang, "crew_ffi", "identity") -fn cast(value: a) -> b diff --git a/src/crew/task_pool.gleam b/src/crew/task_pool.gleam new file mode 100644 index 0000000..9f4fb23 --- /dev/null +++ b/src/crew/task_pool.gleam @@ -0,0 +1,284 @@ +//// The `task_pool` module implements a simple, generic task pool that can run +//// arbitrary functions passed to it using the generic crew pool under the hood. + +import crew +import gleam/erlang/process.{type Name, type Subject} +import gleam/int +import gleam/list +import gleam/otp/actor.{type StartError} +import gleam/otp/static_supervisor.{type Supervisor} +import gleam/otp/supervision.{type ChildSpecification} + +type Any + +/// A builder for configuring a worker pool before starting it. +pub opaque type Builder { + Builder(crew.Builder(Nil, fn() -> Any, Any)) +} + +// sooo the problem here is that there is no way to hide these type aliases +// from the docs. so we lie to the compiler in the public interface since +// we also can't use wrapper types, since we can't map process.Name... +// this module is already full of crimes so who cares. + +pub type PoolMsg + +/// Create a new worker pool builder with the given name. +/// +/// The name is used to register the pool so that work can be sent to it. +/// +/// By default, the pool will have a number of workers equal to the number of +/// scheduler threads available on the system (typically the number of CPU cores). +/// +/// ## Example +/// +/// ```gleam +/// let pool_name = process.new_name("pool") +/// let builder = crew.new(pool_name) +/// ``` +pub fn new(name: Name(PoolMsg)) -> Builder { + Builder(crew.new(cast(name), fn(f) { f() })) +} + +/// Set the number of worker processes in the pool to a fixed number. +/// +/// When set to less than 1 starting the pool will fail. +/// +/// ## Example +/// +/// ```gleam +/// crew.new(pool_name) +/// |> crew.fixed_size(8) // Use 8 workers regardless of CPU count +/// ``` +pub fn fixed_size(builder: Builder, size: Int) -> Builder { + let Builder(builder) = builder + Builder(crew.fixed_size(builder, size)) +} + +/// Start an unsupervised worker pool from the given builder. +/// +/// Returns a supervisor that manages the pool and its workers. In most cases, +/// you should use `supervised` instead to get a child specification that can +/// be added to your application's supervision tree. +/// +/// ## Panics +/// This function will exit the process if any workers fail to start, similar +/// to `static_supervisor.start`. +/// +/// ## Example +/// +/// ```gleam +/// let assert Ok(pool_supervisor) = +/// crew.new(pool_name) +/// |> crew.start +/// ``` +pub fn start(builder: Builder) -> Result(Supervisor, StartError) { + let Builder(builder) = builder + crew.start(builder) +} + +/// Create a child specification for a supervised worker pool. +/// +/// This is the recommended way to start a worker pool as part of your +/// application's supervision tree. The returned child specification can be +/// added to a supervisor using `static_supervisor.add`. +/// +/// ## Example +/// +/// ```gleam +/// let pool_spec = +/// crew.new(pool_name) +/// |> crew.fixed_size(4) +/// |> crew.supervised +/// +/// let assert Ok(_) = +/// supervisor.new(supervisor.OneForOne) +/// |> supervisor.add(pool_spec) +/// |> supervisor.start +/// ``` +pub fn supervised(builder: Builder) -> ChildSpecification(Supervisor) { + let Builder(builder) = builder + crew.supervised(builder) +} + +/// Execute a single piece of work on the pool and wait for the result. +/// +/// This function blocks until the work is completed or the timeout is reached. +/// The work function is executed on one of the worker processes in the pool. +/// +/// ## Parameters +/// - `pool` - The name of the pool to execute work on +/// - `timeout` - Maximum time to wait for completion in milliseconds +/// - `work` - A function containing the work to be executed +/// +/// ## Panics +/// - If the pool does not complete the work within the specified timeout +/// - If the pool is not running +/// - If the worker crashes while executing the work +/// +/// ## Example +/// +/// ```gleam +/// let response = crew.run(pool_name, 5000, fn() { +/// httpc.send(...) +/// }) +/// ``` +pub fn run( + in pool: Name(PoolMsg), + timeout timeout: Int, + do work: fn() -> any, +) -> any { + crew.call(cast(pool), timeout, work) +} + +/// Execute multiple pieces of work concurrently on the pool, without +/// ordering guarantees. +/// +/// Work is distributed among the available workers and executed concurrently. +/// Results are returned in the order they complete, not the order they were +/// submitted. +/// +/// ## Parameters +/// - `pool` - The name of the pool to execute work on +/// - `timeout` - Maximum time to wait for all work to complete in milliseconds +/// - `work` - A list of functions containing work to be executed +/// +/// ## Panics +/// - If the pool does not complete all work within the specified timeout +/// - If the pool is not running +/// - If any worker crashes while executing work +/// +/// ## Example +/// +/// ```gleam +/// let users = crew.run_all(pool_name, 10000, [ +/// fn() { fetch_user_data(user1) }, +/// fn() { fetch_user_data(user2) }, +/// fn() { fetch_user_data(user3) }, +/// ]) +/// ``` +pub fn run_all( + in pool: Name(PoolMsg), + timeout timeout: Int, + do work: List(fn() -> any), +) -> List(any) { + crew.call_parallel(cast(pool), timeout, work) +} + +/// Execute multiple pieces of work concurrently on the pool and return results +/// in submission order. +/// +/// Work functions are distributed across available workers and executed +/// concurrently, but results are reordered to match the original submission +/// order before being returned. +/// +/// ## Parameters +/// - `pool` - The name of the pool to execute work on +/// - `timeout` - Maximum time to wait for all work to complete in milliseconds +/// - `work` - A list of functions containing work to be executed +/// +/// ## Panics +/// - If the pool does not complete all work within the specified timeout +/// - If the pool is not running +/// - If any worker crashes while executing work +/// +/// ## Example +/// +/// ```gleam +/// let assert [user1, user2, user3] = crew.run_sorted(pool_name, 10000, [ +/// fn() { fetch_user_data(user1) }, +/// fn() { fetch_user_data(user2) }, +/// fn() { fetch_user_data(user3) }, +/// ]) +/// ``` +pub fn run_sorted( + in pool: Name(PoolMsg), + timeout timeout: Int, + do work: List(fn() -> any), +) -> List(any) { + parallel_map(work, pool, timeout, fn(f) { f() }) +} + +/// Apply a function to each element of a list concurrently using the worker pool. +/// +/// This is similar to `list.map` but executes the mapping function +/// concurrently across all workers. Results are returned in the same +/// order as the input list. +/// +/// There's a bunch of extra overhead involved with spawning a work item +/// per list element and making sure the order matches. Depending on your +/// workload it might make sense to split your list into chunks first to reduce +/// work queue pressure. +/// +/// ## Parameters +/// - `list` - The list of items to map over +/// - `pool` - The name of the pool to execute work on +/// - `timeout` - Maximum time to wait for all mappings to complete in milliseconds +/// - `fun` - The function to apply to each item +/// +/// ## Panics +/// - If the pool does not complete all mappings within the specified timeout +/// - If the pool is not running +/// - If any worker crashes while executing a mapping +/// +/// ## Example +/// +/// ```gleam +/// let user_ids = [1, 2, 3, 4, 5] +/// let users = crew.parallel_map(user_ids, pool_name, 5000, fetch_user_data) +/// ``` +pub fn parallel_map( + over list: List(a), + in pool: Name(PoolMsg), + timeout timeout: Int, + with fun: fn(a) -> b, +) -> List(b) { + // we don't want to do the multiple subject trick here since that is + // potentially O(n^2) + + let unordered: List(#(Int, b)) = + list + |> list.index_map(fn(item, index) { fn() { #(index, fun(item)) } }) + |> crew.do_call(cast(pool), timeout, _) + + unordered + |> list.sort(fn(a, b) { int.compare(a.0, b.0) }) + |> list.map(fn(x) { x.1 }) +} + +/// Submit a single piece of work to the pool using a subscription channel. +/// +/// This is a lower-level function for submitting work. It is the callers +/// responsibility to handle timeouts, submission order and failures. +/// Most users should prefer the `run*` functions. +/// +/// The first time you `enqueue` is called from a process, the pool sets up +/// a monitor making sure work is cancelled when the process no longer exists +/// to receive a result. You can clean up this monitor early by using `cancel`. +pub fn enqueue( + pool: Name(PoolMsg), + receive: Subject(Result(result, process.ExitReason)), + work: fn() -> result, +) -> Nil { + enqueue_many(pool, receive, [work]) +} + +/// Submit multiple pieces of work to the pool using a subscription channel. +/// +/// This is a lower-level function for submitting work. It is the callers +/// responsibility to handle timeouts, submission order and failures. +/// Most users should prefer the `run*` functions. +/// +/// The first time you `enqueue` is called from a process, the pool sets up +/// a monitor making sure work is cancelled when the process no longer exists +/// to receive a result. You can clean up this monitor early by using `cancel`. +pub fn enqueue_many( + pool: Name(PoolMsg), + receive: Subject(Result(result, process.ExitReason)), + work: List(fn() -> result), +) -> Nil { + crew.enqueue_many(cast(pool), receive, work) +} + +@external(erlang, "crew_ffi", "identity") +fn cast(value: a) -> b