diff --git a/README.md b/README.md index 2497449..d603eb5 100644 --- a/README.md +++ b/README.md @@ -27,15 +27,15 @@ pub fn main() { // Start a pool with 4 workers let pool_name = process.new_name("image_pool") let assert Ok(_) = - crew.new(pool_name) + crew.new(pool_name, process_image) |> crew.fixed_size(4) |> crew.start - // Process images concurrently, but keeping the result order in sync + // Process images concurrently let images = list.range(1, 20) |> list.map(fn(i) { "photo" <> int.to_string(i) <> ".jpg" }) - |> crew.parallel_map(pool_name, 6000, process_image) + |> crew.call_many(pool_name, 6000) echo images } diff --git a/examples/src/backpressure.gleam b/examples/src/backpressure.gleam new file mode 100644 index 0000000..1d547b5 --- /dev/null +++ b/examples/src/backpressure.gleam @@ -0,0 +1,52 @@ +import crew +import gleam/erlang/process +import gleam/io + +pub fn main() { + // we will leak all messages here. + let response = process.new_subject() + + let pool = process.new_name("pool") + let assert Ok(_) = + crew.new(pool, slow_print) + |> crew.fixed_size(1) + |> crew.max_queue_length(4) + |> crew.max_overflow_reductions(2) + |> crew.start + + // start with a simple call to make sure the queue is running + crew.call(pool, 1000, "queue started") + + // + crew.submit_all(pool, 1000, response, ["1"]) + io.println("Done submitting 1") + + // we can submit more than the queue length immediately + crew.submit_all(pool, 1000, response, ["2", "3", "4", "5"]) + io.println("Done submitting 2-5") + + // this now goes into overflow + crew.submit_all(pool, 1000, response, ["6", "7", "8", "9"]) + io.println("Done submitting 6-9") + + // this now goes into overflow + process.spawn(fn() { + crew.submit_all(pool, 5000, response, ["10", "11", "12", "13", "14", "15"]) + io.println("Done submitting 10-15") + }) + + process.spawn(fn() { + crew.submit_all(pool, 5000, response, ["16", "17", "18", "19", "20"]) + io.println("Done submitting 15-20") + }) + + crew.submit_all(pool, 5000, response, ["21", "22", "23", "24", "25"]) + io.println("Done submitting 21-25") + + process.sleep_forever() +} + +fn slow_print(msg: String) { + process.sleep(200) + io.println(msg) +} diff --git a/examples/src/download_actor.gleam b/examples/src/download_actor.gleam index 5aae77c..669593d 100644 --- a/examples/src/download_actor.gleam +++ b/examples/src/download_actor.gleam @@ -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, url) + crew.submit(state.pool, 1000, 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 8f671f7..eb6ba97 100644 --- a/examples/src/image_processing.gleam +++ b/examples/src/image_processing.gleam @@ -1,21 +1,21 @@ -import crew/task_pool +import crew 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("task_pool") + let pool_name = process.new_name("image_pool") let assert Ok(_) = - task_pool.new(pool_name) - |> task_pool.fixed_size(4) - |> task_pool.start + crew.new(pool_name, process_image) + |> crew.fixed_size(4) + |> crew.start - // Process images concurrently, but keeping the result order in sync + // Process images concurrently let images = list.range(1, 20) |> list.map(fn(i) { "photo" <> int.to_string(i) <> ".jpg" }) - |> task_pool.parallel_map(pool_name, 6000, process_image) + |> crew.call_parallel(pool_name, 6000, _) echo images } diff --git a/src/crew.gleam b/src/crew.gleam index ff7b205..e622db7 100644 --- a/src/crew.gleam +++ b/src/crew.gleam @@ -37,6 +37,7 @@ import gleam/erlang/process.{ type Monitor, type Name, type Pid, type Selector, type Subject, } import gleam/list +import gleam/option.{type Option, None, Some} import gleam/otp/actor.{type Next, type StartError} import gleam/otp/static_supervisor.{type Supervisor} import gleam/otp/supervision.{type ChildSpecification} @@ -58,6 +59,8 @@ pub opaque type Builder(state, work, result) { Builder( name: Name(PoolMsg(work, result)), size: Int, + max_queue_length: Option(Int), + max_overflow_reductions: Option(Int), init: fn() -> state, init_timeout: Int, work: fn(state, work) -> result, @@ -81,13 +84,7 @@ 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) }, - ) + new_with_state(name, Nil, fn(_state, input) { work(input) }) } /// Create a new worker pool builder with the given name and state. @@ -101,13 +98,7 @@ pub fn new_with_state( state: state, work: fn(state, work) -> result, ) -> Builder(state, work, result) { - Builder( - name:, - size: scheduler_count(), - init_timeout: 1000, - init: fn() { state }, - work:, - ) + new_with_initialiser(name, 1000, fn() { state }, work) } /// Create a new worker pool builder with the given name and initialiser. @@ -124,7 +115,16 @@ pub fn new_with_initialiser( init init: fn() -> state, run work: fn(state, work) -> result, ) -> Builder(state, work, result) { - Builder(name:, size: scheduler_count(), init_timeout:, init:, work:) + let size = scheduler_count() + Builder( + name:, + size:, + max_queue_length: None, + max_overflow_reductions: None, + init_timeout:, + init:, + work:, + ) } /// Set the number of worker processes in the pool to a fixed number. @@ -144,6 +144,47 @@ pub fn fixed_size( Builder(..builder, size:) } +/// To avoid overloading the system, the internal work queue is limited. +/// +/// If this limit is reached, callers of `enqueue` have to wait until enough +/// work has been done before continuing. +/// +/// By default, the `max_queue_size` depends on the number of workers in the pool. +/// +/// ## Example +/// +/// ```gleam +/// crew.new(pool_name, worder) +/// |> crew.max_queue_length(100) +/// ``` +pub fn max_queue_length( + builder: Builder(state, work, result), + max_queue_length: Int, +) -> Builder(state, work, result) { + Builder(..builder, max_queue_length: Some(max_queue_length)) +} + +/// The `max_overflow_reductions` setting controls how many work items a single +/// request can submit before being deprioritised and being moved back to the +/// end of the queue. +/// +/// This allows other requests to have a chance at getting completed, +/// distributing the load and making the entire system better behaved. +/// By default, `max_overflow_reductions` is equal to the queue length. +/// +/// ## Example +/// +/// ```gleam +/// crew.new(pool_name, worker) +/// |> crew.max_overflow_reductions(100) +/// ``` +pub fn max_overflow_reductions( + builder: Builder(state, work, result), + max_overflow_reductions: Int, +) -> Builder(state, work, result) { + Builder(..builder, max_overflow_reductions: Some(max_overflow_reductions)) +} + // -- START ------------------------------------------------------------------- /// Start an unsupervised worker pool from the given builder. @@ -210,7 +251,7 @@ fn start_tree( let pool_spec = { use <- supervision.worker - actor.new_with_initialiser(1000, init_pool) + actor.new_with_initialiser(1000, init_pool(builder, _)) |> actor.named(builder.name) |> actor.on_message(pool) |> actor.start @@ -252,7 +293,6 @@ fn start_tree( /// - 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 -/// ``` pub fn call( in pool: Name(PoolMsg(work, result)), timeout timeout: Int, @@ -297,7 +337,9 @@ pub fn do_call( let monitor = process.monitor(pool_pid) let receive = process.new_subject() - enqueue_many(pool, receive, work) + + // we can always cast here because we will be blocking below anyways + actor.send(process.named_subject(pool), Enqueue(receive, work, None)) let selector = process.new_selector() @@ -369,6 +411,9 @@ pub fn cancel( } /// Submit a single piece of work to the pool using a subscription channel. +/// The work will be done asynchronously and the result will be sent back +/// to the provided subject. The timeout controls how long to wait for the +/// work to be added to the queue successfully. /// /// This is a lower-level function for submitting work. It is the callers /// responsibility to handle timeouts, submission order and failures. @@ -377,15 +422,19 @@ pub fn cancel( /// 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( +pub fn submit( pool: Name(PoolMsg(work, result)), + timeout: Int, receive: Subject(Result(result, process.ExitReason)), work: work, ) -> Nil { - enqueue_many(pool, receive, [work]) + submit_all(pool, timeout, receive, [work]) } /// Submit multiple pieces of work to the pool using a subscription channel. +/// The work will be done asynchronously and the result will be sent back +/// to the provided subject. The timeout controls how long to wait for the +/// work to be added to the queue successfully. /// /// This is a lower-level function for submitting work. It is the callers /// responsibility to handle timeouts, submission order and failures. @@ -394,12 +443,26 @@ pub fn enqueue( /// 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( +pub fn submit_all( pool: Name(PoolMsg(work, result)), + timeout: Int, receive: Subject(Result(result, process.ExitReason)), work: List(work), ) -> Nil { - actor.send(process.named_subject(pool), Enqueue(receive, work)) + let counter = get_counter(pool) + let subject = process.named_subject(pool) + + case work { + [] -> Nil + [_, ..] if counter > 0 -> { + actor.send(subject, Enqueue(receive:, work:, enqueued: None)) + } + [_, ..] -> { + actor.call(subject, timeout, fn(enqueued) { + Enqueue(receive:, work:, enqueued: Some(enqueued)) + }) + } + } } // -- POOL -------------------------------------------------------------------- @@ -410,18 +473,27 @@ pub opaque type PoolMsg(work, result) { MonitoredProcessExited(reason: process.Down) // GetWorkerCount(reply_to: Subject(Int)) - Enqueue(receive: Receiver(result), work: List(work)) + Enqueue( + receive: Receiver(result), + work: List(work), + enqueued: Option(Subject(Nil)), + ) Cancel(receive: Receiver(result)) } type State(work, result) { State( + name: Name(PoolMsg(work, result)), worker_count: Int, idle_workers: List(Worker(work, result)), active_workers: Dict(Pid, ActiveWorker(work, result)), + // callers: Dict(Pid, Caller(result)), channels: Dict(Receiver(result), Channel(result)), + // queue: Deque(Work(work, result)), + overflow_queue: Deque(Request(work, result)), + max_overflow_reductions: Int, ) } @@ -441,20 +513,47 @@ type Channel(result) { Channel(from: Pid, workers: Set(Pid), receive: Receiver(result)) } -fn init_pool(self: Subject(PoolMsg(work, result))) { +type Request(work, result) { + Request( + from: Pid, + receive: Receiver(result), + work: List(work), + enqueued: Option(Subject(Nil)), + reductions: Int, + ) +} + +fn init_pool( + builder: Builder(state, work, result), + self: Subject(PoolMsg(work, result)), +) { let selector = process.new_selector() |> process.select(self) |> process.select_monitors(MonitoredProcessExited) + // it's big, but it will provide backpressure eventually. + let max_queue_length = + builder.max_queue_length + |> option.unwrap(builder.size * 100) + + let max_overflow_reductions = + builder.max_overflow_reductions + |> option.unwrap(max_queue_length) + + set_counter(builder.name, max_queue_length) + let state = State( + name: builder.name, worker_count: 0, idle_workers: [], active_workers: dict.new(), callers: dict.new(), channels: dict.new(), queue: deque.new(), + overflow_queue: deque.new(), + max_overflow_reductions:, ) actor.initialised(state) @@ -551,8 +650,12 @@ fn pool( state } - Enqueue(work: [], ..) -> state - Enqueue(receive:, work:) -> { + Enqueue(receive: _, work: [], enqueued: None) -> state + Enqueue(receive: _, work: [], enqueued: Some(enqueued)) -> { + process.send(enqueued, Nil) + state + } + Enqueue(receive:, work:, enqueued:) -> { // caller already exited, do not queue their work. use pid <- try_(process.subject_owner(receive), state) @@ -574,7 +677,7 @@ fn pool( // enqueue_loop will insert the channel into the state. State(..state, callers: dict.insert(state.callers, pid, caller)) - |> enqueue_loop(channel, work) + |> enqueue_loop(channel, work, enqueued) } Cancel(receive:) -> { @@ -616,11 +719,12 @@ fn enqueue_loop( state: State(work, result), channel: Channel(result), work: List(work), + enqueued: Option(Subject(Nil)), ) -> State(work, result) { let Channel(from: caller, receive:, workers:) = channel - case work, state.idle_workers { - [work, ..rest], [worker, ..idle_workers] -> { + case work, state.idle_workers, enqueued { + [work, ..rest], [worker, ..idle_workers], _ -> { // got a work item and a worker - it's a match! let work = Work(work:, caller:, receive:) process.send(worker.send, work) @@ -632,19 +736,69 @@ fn enqueue_loop( let channel = Channel(..channel, workers: set.insert(workers, worker.pid)) State(..state, active_workers:, idle_workers:) - |> enqueue_loop(channel, rest) + |> enqueue_loop(channel, rest, enqueued) } - _, _ -> { - let queue = - list.fold(work, state.queue, fn(queue, work) { - deque.push_back(queue, Work(work:, caller:, receive:)) - }) + [_, ..], [], _ -> { + let State(queue:, overflow_queue:, ..) = state + let capacity = get_counter(state.name) + + let #(queue, overflow_queue, new_capacity) = + enqueue_loop2( + caller, + receive, + work, + enqueued, + queue, + overflow_queue, + capacity, + ) + + decrement_counter(state.name, capacity - new_capacity) let channels = dict.insert(state.channels, receive, channel) + State(..state, queue:, overflow_queue:, channels:) + } - State(..state, queue:, channels:) + [], _, Some(enqueued) -> { + process.send(enqueued, Nil) + State(..state, channels: dict.insert(state.channels, receive, channel)) } + + [], _, None -> + State(..state, channels: dict.insert(state.channels, receive, channel)) + } +} + +fn enqueue_loop2( + caller: Pid, + receive: Receiver(result), + work: List(work), + enqueued: Option(Subject(Nil)), + queue: Deque(Work(work, result)), + overflow: Deque(Request(work, result)), + capacity: Int, +) { + case work, enqueued { + [work, ..rest], _ if capacity > 0 -> { + let queue = deque.push_back(queue, Work(work:, caller:, receive:)) + let capacity = capacity - 1 + enqueue_loop2(caller, receive, rest, enqueued, queue, overflow, capacity) + } + + [_, ..], _ -> { + let request = + Request(from: caller, receive:, work:, enqueued:, reductions: 0) + let overflow = deque.push_back(overflow, request) + #(queue, overflow, capacity) + } + + [], Some(enqueued) -> { + process.send(enqueued, Nil) + #(queue, overflow, capacity) + } + + [], None -> #(queue, overflow, capacity) } } @@ -652,15 +806,14 @@ fn try_dequeue_work( state: State(work, result), worker: Worker(work, result), ) -> State(work, result) { - use #(work, queue) <- try( - deque.pop_front(state.queue), + use #(work, queue, overflow_queue) <- try(pop_work(state), fn(_) { // the queue is empty, this worker becomes idle. - fn(_) { State(..state, idle_workers: [worker, ..state.idle_workers]) }, - ) + State(..state, idle_workers: [worker, ..state.idle_workers]) + }) use channel <- try(dict.get(state.channels, work.receive), fn(_) { - // this channel got cancelled - try_dequeue_work(State(..state, queue:), worker) + // this channel got cancelled, try again + try_dequeue_work(State(..state, queue:, overflow_queue:), worker) }) process.send(worker.send, work) @@ -673,7 +826,55 @@ fn try_dequeue_work( Channel(..channel, workers: set.insert(channel.workers, worker.pid)) let channels = dict.insert(state.channels, work.receive, channel) - State(..state, active_workers:, channels:, queue:) + State(..state, active_workers:, channels:, queue:, overflow_queue:) +} + +fn pop_work(state: State(work, result)) -> Result(_, Nil) { + use #(work, queue) <- result.try(deque.pop_front(state.queue)) + + case deque.pop_front(state.overflow_queue) { + Ok(#(Request(work: [first, ..rest], ..) as request, overflow_queue)) -> { + let Request(from: caller, receive:, enqueued:, ..) = request + + let queue = deque.push_back(queue, Work(work: first, caller:, receive:)) + let overflow_queue = case rest, enqueued { + [_, ..], _ if request.reductions + 1 < state.max_overflow_reductions -> { + let request = + Request(..request, work: rest, reductions: request.reductions + 1) + deque.push_front(overflow_queue, request) + } + + [_, ..], _ -> { + let request = Request(..request, work: rest, reductions: 0) + deque.push_back(overflow_queue, request) + } + + [], Some(enqueued) -> { + process.send(enqueued, Nil) + overflow_queue + } + + [], None -> overflow_queue + } + + Ok(#(work, queue, overflow_queue)) + } + + Ok(#(Request(work: [], enqueued: Some(enqueued), ..), overflow_queue)) -> { + // We got an empty work array - this should also not happen, but we + // can handle it just fine, be it sending the enqueued notification late, or twice. + process.send(enqueued, Nil) + pop_work(State(..state, overflow_queue:)) + } + + Ok(#(Request(work: [], enqueued: None, ..), overflow_queue)) -> + pop_work(State(..state, overflow_queue:)) + + Error(Nil) -> { + increment_counter(state.name, 1) + Ok(#(work, queue, state.overflow_queue)) + } + } } // -- WORKER ------------------------------------------------------------------ @@ -767,3 +968,15 @@ fn scheduler_count() -> Int @external(erlang, "crew_ffi", "system_time") fn system_time() -> Int + +@external(erlang, "crew_ffi", "get_counter") +fn get_counter(name: Name(a)) -> Int + +@external(erlang, "crew_ffi", "increment_counter") +fn increment_counter(name: Name(a), amount: Int) -> Nil + +@external(erlang, "crew_ffi", "decrement_counter") +fn decrement_counter(name: Name(a), amount: Int) -> Nil + +@external(erlang, "crew_ffi", "set_counter") +fn set_counter(name: Name(a), value: Int) -> Nil diff --git a/src/crew/task_pool.gleam b/src/crew/task_pool.gleam index 9f4fb23..6a0e16a 100644 --- a/src/crew/task_pool.gleam +++ b/src/crew/task_pool.gleam @@ -34,7 +34,7 @@ pub type PoolMsg /// /// ```gleam /// let pool_name = process.new_name("pool") -/// let builder = crew.new(pool_name) +/// let builder = task_pool.new(pool_name) /// ``` pub fn new(name: Name(PoolMsg)) -> Builder { Builder(crew.new(cast(name), fn(f) { f() })) @@ -47,14 +47,53 @@ pub fn new(name: Name(PoolMsg)) -> Builder { /// ## Example /// /// ```gleam -/// crew.new(pool_name) -/// |> crew.fixed_size(8) // Use 8 workers regardless of CPU count +/// task_pool.new(pool_name) +/// |> task_pool.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)) } +/// To avoid overloading the system, the internal work queue is limited. +/// +/// If this limit is reached, callers of `enqueue` have to wait until enough +/// work has been done before continuing. +/// +/// By default, the `max_queue_size` depends on the number workers in the pool. +/// ## Example +/// +/// ```gleam +/// task_pool.new(pool_name) +/// |> task_pool.max_queue_length(100) +/// ``` +pub fn max_queue_length(builder: Builder, max_queue_length: Int) -> Builder { + let Builder(builder) = builder + Builder(crew.max_queue_length(builder, max_queue_length)) +} + +/// The `max_overflow_reductions` setting controls how many work items a single +/// request can submit before being deprioritised and being moved back to the +/// end of the queue. +/// +/// This allows other requests to have a chance at getting completed, +/// distributing the load and making the entire system better behaved. +/// By default, `max_overflow_reductions` is equal to the queue length. +/// +/// ## Example +/// +/// ```gleam +/// task_pool.new(pool_name) +/// |> task_pool.max_overflow_reductions(100) +/// ``` +pub fn max_overflow_reductions( + builder: Builder, + max_overflow_reductions: Int, +) -> Builder { + let Builder(builder) = builder + Builder(crew.max_overflow_reductions(builder, max_overflow_reductions)) +} + /// Start an unsupervised worker pool from the given builder. /// /// Returns a supervisor that manages the pool and its workers. In most cases, @@ -69,8 +108,8 @@ pub fn fixed_size(builder: Builder, size: Int) -> Builder { /// /// ```gleam /// let assert Ok(pool_supervisor) = -/// crew.new(pool_name) -/// |> crew.start +/// task_pool.new(pool_name) +/// |> task_pool.start /// ``` pub fn start(builder: Builder) -> Result(Supervisor, StartError) { let Builder(builder) = builder @@ -87,9 +126,9 @@ pub fn start(builder: Builder) -> Result(Supervisor, StartError) { /// /// ```gleam /// let pool_spec = -/// crew.new(pool_name) -/// |> crew.fixed_size(4) -/// |> crew.supervised +/// task_pool.new(pool_name) +/// |> task_pool.fixed_size(4) +/// |> task_pool.supervised /// /// let assert Ok(_) = /// supervisor.new(supervisor.OneForOne) @@ -119,7 +158,7 @@ pub fn supervised(builder: Builder) -> ChildSpecification(Supervisor) { /// ## Example /// /// ```gleam -/// let response = crew.run(pool_name, 5000, fn() { +/// let response = task_pool.run(pool_name, 5000, fn() { /// httpc.send(...) /// }) /// ``` @@ -151,7 +190,7 @@ pub fn run( /// ## Example /// /// ```gleam -/// let users = crew.run_all(pool_name, 10000, [ +/// let users = task_pool.run_all(pool_name, 10000, [ /// fn() { fetch_user_data(user1) }, /// fn() { fetch_user_data(user2) }, /// fn() { fetch_user_data(user3) }, @@ -185,7 +224,7 @@ pub fn run_all( /// ## Example /// /// ```gleam -/// let assert [user1, user2, user3] = crew.run_sorted(pool_name, 10000, [ +/// let assert [user1, user2, user3] = task_pool.run_sorted(pool_name, 10000, [ /// fn() { fetch_user_data(user1) }, /// fn() { fetch_user_data(user2) }, /// fn() { fetch_user_data(user3) }, @@ -225,7 +264,7 @@ pub fn run_sorted( /// /// ```gleam /// let user_ids = [1, 2, 3, 4, 5] -/// let users = crew.parallel_map(user_ids, pool_name, 5000, fetch_user_data) +/// let users = task_pool.parallel_map(user_ids, pool_name, 5000, fetch_user_data) /// ``` pub fn parallel_map( over list: List(a), @@ -255,12 +294,13 @@ pub fn parallel_map( /// 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( +pub fn submit( pool: Name(PoolMsg), + timeout: Int, receive: Subject(Result(result, process.ExitReason)), work: fn() -> result, ) -> Nil { - enqueue_many(pool, receive, [work]) + submit_all(pool, timeout, receive, [work]) } /// Submit multiple pieces of work to the pool using a subscription channel. @@ -272,12 +312,13 @@ pub fn enqueue( /// 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( +pub fn submit_all( pool: Name(PoolMsg), + timeout: Int, receive: Subject(Result(result, process.ExitReason)), work: List(fn() -> result), ) -> Nil { - crew.enqueue_many(cast(pool), receive, work) + crew.submit_all(cast(pool), timeout, receive, work) } @external(erlang, "crew_ffi", "identity") diff --git a/src/crew_ffi.erl b/src/crew_ffi.erl index dece0cf..554e318 100644 --- a/src/crew_ffi.erl +++ b/src/crew_ffi.erl @@ -1,6 +1,26 @@ -module(crew_ffi). -export([identity/1, scheduler_count/0, system_time/0]). +-export([get_counter/1, set_counter/2, increment_counter/2, decrement_counter/2]). identity(X) -> X. scheduler_count() -> erlang:system_info(schedulers_online). system_time() -> os:system_time(millisecond). + +get_counter(Name) -> + case persistent_term:get(Name, undefined) of + undefined -> 0; + Counter -> counters:get(Counter, 1) + end. + +set_counter(Name, Value) -> + case persistent_term:get(Name, undefined) of + undefined -> + Counter = counters:new(1, []), + counters:put(Counter, 1, Value), + persistent_term:put(Name, Counter); + + Counter -> counters:put(Counter, 1, Value) + end. + +increment_counter(Name, Incr) -> counters:add(persistent_term:get(Name), 1, Incr). +decrement_counter(Name, Decr) -> counters:sub(persistent_term:get(Name), 1, Decr).