From f09ecbc01bd82f16dc0a98047fcd52976cde3e7e Mon Sep 17 00:00:00 2001 From: Nelson Vides Date: Sat, 21 Feb 2026 12:06:52 +0100 Subject: [PATCH 1/2] Use erlang:system_time/0 and convert time units only if required --- src/wpool_pool.erl | 7 ++++--- 1 file changed, 4 insertions(+), 3 deletions(-) diff --git a/src/wpool_pool.erl b/src/wpool_pool.erl index 2272ff6..6ee3e4e 100644 --- a/src/wpool_pool.erl +++ b/src/wpool_pool.erl @@ -52,7 +52,7 @@ workers :: tuple(), opts :: wpool:options(), qmanager :: wpool_queue_manager:queue_mgr(), - born = erlang:system_time(second) :: integer() + born = erlang:system_time() :: integer() }). -opaque wpool() :: #wpool{}. @@ -323,8 +323,9 @@ function_location(Function, Location) -> task(undefined) -> []; task({_TaskId, Started, Task}) -> - Time = erlang:system_time(second), - [{task, Task}, {runtime, Time - Started}]. + Time = erlang:system_time(), + Runtime = erlang:convert_time_unit(Time - Started, native, second), + [{task, Task}, {runtime, Runtime}]. %% @doc Set next within the worker pool record. Useful when using %% a custom strategy function. From d28a5e74f59e3fa8eca3d6f3ba360a94744c303d Mon Sep 17 00:00:00 2001 From: Nelson Vides Date: Sat, 21 Feb 2026 12:07:41 +0100 Subject: [PATCH 2/2] Optimise wpool_process loops Reduce the number of operations that the wpool_process needs to do on every loop: - Options get all its keys as mandatory, and will extract them all at once instead of on subsequent map gets in separate modules. - utils task_init and task_end are moved into wpool_process, as within the same module the JIT can optimise types, and readouts can be done all at once - task_end doesn't do erase(wpool_task) but put(wpool_task, undefined), as this way the process dictionary doesn't need any hash-table or bucket reshuffling, potentially. - notify_queue_manager chooses the function to call directly, instead of doing a dynamic call, it's more performant as it doesn't need to lookup the function pointer at runtime, instead it only needs to choose it based on a function clause of just 4 elements that are in the same module, so again, JIT can optimise more aggressively. --- src/wpool_pool.erl | 24 +++++--- src/wpool_process.erl | 97 +++++++++++++++++++++++---------- src/wpool_process_callbacks.erl | 10 ++-- src/wpool_utils.erl | 41 +------------- test/wpool_process_SUITE.erl | 37 +++++++------ 5 files changed, 111 insertions(+), 98 deletions(-) diff --git a/src/wpool_pool.erl b/src/wpool_pool.erl index 6ee3e4e..4c59689 100644 --- a/src/wpool_pool.erl +++ b/src/wpool_pool.erl @@ -393,8 +393,8 @@ init({Name, Options}) -> WorkerOpts0 = [{time_checker, TimeCheckerName}] ++ - maybe_queue_manager(Options, {queue_manager, QueueManagerName}) ++ - maybe_event_manager(Options, {event_manager, EventManagerName}), + maybe_queue_manager(Options, QueueManagerName) ++ + maybe_event_manager(Options, EventManagerName), WorkerOpts = maps:merge( maps:from_list(WorkerOpts0), Options @@ -437,8 +437,8 @@ init({Name, Options}) -> Children = [TimeCheckerSpec] ++ - maybe_queue_manager(Options, QueueManagerSpec) ++ - maybe_event_manager(Options, EventManagerSpec) ++ + maybe_queue_manager_child(Options, QueueManagerSpec) ++ + maybe_event_manager_child(Options, EventManagerSpec) ++ [ProcessSupSpec], SupIntensity = maps:get(pool_sup_intensity, Options, 5), @@ -610,11 +610,21 @@ build_wpool(Name) -> end. maybe_queue_manager(#{enable_queues := false}, _) -> - []; + [{queue_manager, undefined}]; maybe_queue_manager(_, Item) -> - [Item]. + [{queue_manager, Item}]. maybe_event_manager(#{enable_callbacks := true}, Item) -> - [Item]; + [{event_manager, Item}]; maybe_event_manager(_, _) -> + [{event_manager, undefined}]. + +maybe_queue_manager_child(#{enable_queues := false}, _) -> + []; +maybe_queue_manager_child(_, Item) -> + [Item]. + +maybe_event_manager_child(#{enable_callbacks := true}, Item) -> + [Item]; +maybe_event_manager_child(_, _) -> []. diff --git a/src/wpool_process.erl b/src/wpool_process.erl index 297b648..e32c077 100644 --- a/src/wpool_process.erl +++ b/src/wpool_process.erl @@ -51,15 +51,17 @@ name :: atom(), mod :: #callback_cache{}, state :: term(), - options :: - #{ - time_checker := atom(), - queue_manager := atom(), - overrun_warning := timeout(), - _ => _ - } + options :: opts() }). +-type opts() :: #{ + time_checker := atom(), + queue_manager := atom(), + event_manager := atom(), + overrun_warning := timeout(), + _ => _ +}. + -opaque state() :: #state{}. -export_type([state/0]). @@ -137,15 +139,16 @@ get_state(#state{state = State}) -> %%% init, terminate, code_change, info callbacks %%%=================================================================== %% @private --spec init({atom(), atom(), term(), wpool:options()}) -> +-spec init({atom(), atom(), term(), opts()}) -> {ok, state()} | {ok, state(), next_step()} | {stop, can_not_ignore} | {stop, term()}. init({Name, Mod, InitArgs, Options}) -> - wpool_process_callbacks:notify(handle_init_start, Options, [Name]), + #{event_manager := EventManager, queue_manager := QueueManager} = Options, + wpool_process_callbacks:notify(handle_init_start, EventManager, [Name]), CbCache = create_callback_cache(Mod), case Mod:init(InitArgs) of {ok, ModState} -> - ok = notify_queue_manager(new_worker, Name, Options), - wpool_process_callbacks:notify(handle_worker_creation, Options, [Name]), + ok = notify_queue_manager(new_worker, Name, QueueManager), + wpool_process_callbacks:notify(handle_worker_creation, EventManager, [Name]), {ok, #state{ name = Name, mod = CbCache, @@ -153,8 +156,8 @@ init({Name, Mod, InitArgs, Options}) -> options = Options }}; {ok, ModState, NextStep} -> - ok = notify_queue_manager(new_worker, Name, Options), - wpool_process_callbacks:notify(handle_worker_creation, Options, [Name]), + ok = notify_queue_manager(new_worker, Name, QueueManager), + wpool_process_callbacks:notify(handle_worker_creation, EventManager, [Name]), {ok, #state{ name = Name, @@ -176,11 +179,10 @@ terminate(Reason, State) -> mod = #callback_cache{module = Mod}, state = ModState, name = Name, - options = Options - } = - State, - ok = notify_queue_manager(worker_dead, Name, Options), - wpool_process_callbacks:notify(handle_worker_death, Options, [Name, Reason]), + options = #{event_manager := EventManager, queue_manager := QueueManager} + } = State, + ok = notify_queue_manager(worker_dead, Name, QueueManager), + wpool_process_callbacks:notify(handle_worker_death, EventManager, [Name, Reason]), case erlang:function_exported(Mod, terminate, 2) of true -> Mod:terminate(Reason, ModState); @@ -246,6 +248,13 @@ handle_continue(Continue, #state{mod = #callback_cache{module = Mod}} = State) - {stop, Reason, NewState} -> {stop, Reason, State#state{state = NewState}} catch + error:undef:Stacktrace -> + case erlang:function_exported(Mod, handle_continue, 2) of + false -> + {noreply, State}; + true -> + erlang:raise(error, undef, Stacktrace) + end; _:{noreply, NewState} -> {noreply, State#state{state = NewState}}; _:{noreply, NewState, NextStep} -> @@ -272,8 +281,9 @@ format_status(#{state := #state{mod = #callback_cache{module = Mod}}} = Status) {noreply, state()} | {noreply, state(), next_step()} | {stop, term(), state()}. handle_cast(Cast, #state{mod = CbCache, options = Options} = State) -> #callback_cache{handle_cast = HandleCast} = CbCache, - Task = wpool_utils:task_init({cast, Cast}, Options), - ok = notify_queue_manager(worker_busy, State#state.name, Options), + #{overrun_warning := OverrunWarning, queue_manager := QueueManager} = Options, + Task = task_init(OverrunWarning, {cast, Cast}, Options), + ok = notify_queue_manager(worker_busy, State#state.name, QueueManager), Reply = try HandleCast(Cast, State#state.state) of {noreply, NewState} -> @@ -290,8 +300,8 @@ handle_cast(Cast, #state{mod = CbCache, options = Options} = State) -> _:{stop, Reason, NewState} -> {stop, Reason, State#state{state = NewState}} end, - wpool_utils:task_end(Task), - ok = notify_queue_manager(worker_ready, State#state.name, Options), + task_end(Task), + ok = notify_queue_manager(worker_ready, State#state.name, QueueManager), Reply. %% @private @@ -304,8 +314,9 @@ handle_cast(Cast, #state{mod = CbCache, options = Options} = State) -> | {stop, term(), state()}. handle_call(Call, From, #state{mod = CbCache, options = Options} = State) -> #callback_cache{handle_call = HandleCall} = CbCache, - Task = wpool_utils:task_init({call, Call}, Options), - ok = notify_queue_manager(worker_busy, State#state.name, Options), + #{overrun_warning := OverrunWarning, queue_manager := QueueManager} = Options, + Task = task_init(OverrunWarning, {call, Call}, Options), + ok = notify_queue_manager(worker_busy, State#state.name, QueueManager), Reply = try HandleCall(Call, From, State#state.state) of {noreply, NewState} -> @@ -334,14 +345,40 @@ handle_call(Call, From, #state{mod = CbCache, options = Options} = State) -> _:{stop, Reason, Response, NewState} -> {stop, Reason, Response, State#state{state = NewState}} end, - wpool_utils:task_end(Task), - ok = notify_queue_manager(worker_ready, State#state.name, Options), + task_end(Task), + ok = notify_queue_manager(worker_ready, State#state.name, QueueManager), Reply. -notify_queue_manager(Function, Name, #{queue_manager := QueueManager}) -> - wpool_queue_manager:Function(QueueManager, Name); -notify_queue_manager(_, _, _) -> - ok. +notify_queue_manager(_, _, undefined) -> + ok; +notify_queue_manager(worker_busy, Name, QueueManager) -> + wpool_queue_manager:worker_busy(QueueManager, Name); +notify_queue_manager(worker_ready, Name, QueueManager) -> + wpool_queue_manager:worker_ready(QueueManager, Name); +notify_queue_manager(worker_dead, Name, QueueManager) -> + wpool_queue_manager:worker_dead(QueueManager, Name); +notify_queue_manager(new_worker, Name, QueueManager) -> + wpool_queue_manager:new_worker(QueueManager, Name). + +task_init(infinity, Task, _) -> + Time = erlang:system_time(), + erlang:put(wpool_task, {undefined, Time, Task}), + undefined; +task_init(OverrunTime, Task, #{time_checker := TimeChecker, max_overrun_warnings := MaxWarnings}) -> + TaskId = erlang:make_ref(), + Time = erlang:system_time(), + erlang:put(wpool_task, {TaskId, Time, Task}), + erlang:send_after( + OverrunTime, + TimeChecker, + {check, self(), TaskId, OverrunTime, MaxWarnings} + ). + +task_end(undefined) -> + erlang:put(wpool_task, undefined); +task_end(TimerRef) -> + _ = erlang:cancel_timer(TimerRef, [{async, true}, {info, false}]), + erlang:put(wpool_task, undefined). create_callback_cache(Mod) -> #callback_cache{ diff --git a/src/wpool_process_callbacks.erl b/src/wpool_process_callbacks.erl index 19f1e42..4b442b0 100644 --- a/src/wpool_process_callbacks.erl +++ b/src/wpool_process_callbacks.erl @@ -45,11 +45,11 @@ handle_call(Msg, State) -> {ok, {error, {unexpected_call, Msg}}, State}. %% @doc Sends a notification to all registered callback modules. --spec notify(event(), #{event_manager := any(), _ => _}, [any()]) -> ok. -notify(Event, #{event_manager := EventMgr}, Args) -> - gen_event:notify(EventMgr, {Event, Args}); -notify(_, _, _) -> - ok. +-spec notify(event(), undefined | atom(), [any()]) -> ok. +notify(_, undefined, _) -> + ok; +notify(Event, EventMgr, Args) -> + gen_event:notify(EventMgr, {Event, Args}). %% @doc Adds a callback module. -spec add_callback_module(wpool:name(), module()) -> ok | {error, any()}. diff --git a/src/wpool_utils.erl b/src/wpool_utils.erl index 201338c..2ad6734 100644 --- a/src/wpool_utils.erl +++ b/src/wpool_utils.erl @@ -11,47 +11,10 @@ % KIND, either express or implied. See the License for the % specific language governing permissions and limitations % under the License. -%%% @doc Common functions for wpool_process and other modules. +%%% @private -module(wpool_utils). -%% API --export([task_init/2, task_end/1, add_defaults/1]). - -%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%% -%% Api -%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%% - -%% @doc Marks Task as started in this worker --spec task_init(term(), #{overrun_warning := timeout(), _ => _}) -> - undefined | reference(). -task_init(Task, #{overrun_warning := infinity}) -> - Time = erlang:system_time(second), - erlang:put(wpool_task, {undefined, Time, Task}), - undefined; -task_init( - Task, - #{ - overrun_warning := OverrunTime, - time_checker := TimeChecker, - max_overrun_warnings := MaxWarnings - } -) -> - TaskId = erlang:make_ref(), - Time = erlang:system_time(second), - erlang:put(wpool_task, {TaskId, Time, Task}), - erlang:send_after( - OverrunTime, - TimeChecker, - {check, self(), TaskId, OverrunTime, MaxWarnings} - ). - -%% @doc Removes the current task from the worker --spec task_end(undefined | reference()) -> ok. -task_end(undefined) -> - erlang:erase(wpool_task); -task_end(TimerRef) -> - _ = erlang:cancel_timer(TimerRef, [{async, true}, {info, false}]), - erlang:erase(wpool_task). +-export([add_defaults/1]). %% @doc Adds default parameters to a pool configuration -spec add_defaults([wpool:option()] | wpool:options()) -> wpool:options(). diff --git a/test/wpool_process_SUITE.erl b/test/wpool_process_SUITE.erl index d4fe100..7d662c1 100644 --- a/test/wpool_process_SUITE.erl +++ b/test/wpool_process_SUITE.erl @@ -73,16 +73,16 @@ end_per_testcase(_TestCase, Config) -> -spec init(config()) -> {comment, []}. init(_Config) -> - {error, can_not_ignore} = wpool_process:start_link(?MODULE, echo_server, ignore, #{}), - {error, ?MODULE} = wpool_process:start_link(?MODULE, echo_server, {stop, ?MODULE}, #{}), - {ok, _Pid} = wpool_process:start_link(?MODULE, echo_server, {ok, state}, #{}), + {error, can_not_ignore} = wpool_process:start_link(?MODULE, echo_server, ignore, opts()), + {error, ?MODULE} = wpool_process:start_link(?MODULE, echo_server, {stop, ?MODULE}, opts()), + {ok, _Pid} = wpool_process:start_link(?MODULE, echo_server, {ok, state}, opts()), wpool_process:cast(?MODULE, {stop, normal, state}), {comment, []}. -spec init_timeout(config()) -> {comment, []}. init_timeout(_Config) -> - {ok, Pid} = wpool_process:start_link(?MODULE, echo_server, {ok, state, 0}, #{}), + {ok, Pid} = wpool_process:start_link(?MODULE, echo_server, {ok, state, 0}, opts()), timeout = get_state(?MODULE), Pid ! {stop, normal, state}, false = ktn_task:wait_for(fun() -> erlang:is_process_alive(Pid) end, false), @@ -91,7 +91,7 @@ init_timeout(_Config) -> -spec info(config()) -> {comment, []}. info(_Config) -> - {ok, Pid} = wpool_process:start_link(?MODULE, echo_server, {ok, state}, #{}), + {ok, Pid} = wpool_process:start_link(?MODULE, echo_server, {ok, state}, opts()), Pid ! {noreply, newstate}, newstate = get_state(?MODULE), Pid ! {noreply, newerstate, 1}, @@ -103,7 +103,7 @@ info(_Config) -> -spec cast(config()) -> {comment, []}. cast(_Config) -> - {ok, Pid} = wpool_process:start_link(?MODULE, echo_server, {ok, state}, #{}), + {ok, Pid} = wpool_process:start_link(?MODULE, echo_server, {ok, state}, opts()), wpool_process:cast(Pid, {noreply, newstate}), newstate = get_state(?MODULE), wpool_process:cast(Pid, {noreply, newerstate, 0}), @@ -115,7 +115,7 @@ cast(_Config) -> -spec send_request(config()) -> {comment, []}. send_request(_Config) -> - {ok, Pid} = wpool_process:start_link(?MODULE, echo_server, {ok, state}, #{}), + {ok, Pid} = wpool_process:start_link(?MODULE, echo_server, {ok, state}, opts()), Req1 = wpool_process:send_request(Pid, {reply, ok1, newstate}), ok1 = wait_response(Req1), Req2 = wpool_process:send_request(Pid, {reply, ok2, newerstate, 1}), @@ -143,7 +143,7 @@ continue(_Config) -> ?MODULE, echo_server, {ok, state, {continue, C(continue_state)}}, - #{} + opts() ), continue_state = get_state(Pid), @@ -193,14 +193,14 @@ continue(_Config) -> -spec handle_info_missing(config()) -> {comment, []}. handle_info_missing(_Config) -> %% sleepy_server does not implement handle_info/2 - {ok, Pid} = wpool_process:start_link(?MODULE, sleepy_server, 1, #{}), + {ok, Pid} = wpool_process:start_link(?MODULE, sleepy_server, 1, opts()), Pid ! test, {comment, []}. -spec handle_info_fails(config()) -> {comment, []}. handle_info_fails(_Config) -> %% sleepy_server does not implement handle_info/2 - {ok, Pid} = wpool_process:start_link(?MODULE, crashy_server, {ok, state}, #{}), + {ok, Pid} = wpool_process:start_link(?MODULE, crashy_server, {ok, state}, opts()), Pid ! undef, false = ktn_task:wait_for(fun() -> erlang:is_process_alive(Pid) end, false), {comment, []}. @@ -208,7 +208,7 @@ handle_info_fails(_Config) -> -spec format_status(config()) -> {comment, []}. format_status(_Config) -> %% echo_server implements format_status/1 - {ok, Pid} = wpool_process:start_link(?MODULE, echo_server, {ok, state}, #{}), + {ok, Pid} = wpool_process:start_link(?MODULE, echo_server, {ok, state}, opts()), %% therefore it returns State as its status state = get_state(Pid), {comment, []}. @@ -216,7 +216,7 @@ format_status(_Config) -> -spec no_format_status(config()) -> {comment, []}. no_format_status(_Config) -> %% crashy_server doesn't implement format_status/1 - {ok, Pid} = wpool_process:start_link(?MODULE, crashy_server, state, #{}), + {ok, Pid} = wpool_process:start_link(?MODULE, crashy_server, state, opts()), %% therefore it uses the default format for the stauts (but with the status of %% the gen_server, not wpool_process) state = get_state(Pid), @@ -224,7 +224,7 @@ no_format_status(_Config) -> -spec call(config()) -> {comment, []}. call(_Config) -> - {ok, Pid} = wpool_process:start_link(?MODULE, echo_server, {ok, state}, #{}), + {ok, Pid} = wpool_process:start_link(?MODULE, echo_server, {ok, state}, opts()), ok1 = wpool_process:call(Pid, {reply, ok1, newstate}, 5000), newstate = get_state(?MODULE), ok2 = wpool_process:call(Pid, {reply, ok2, newerstate, 1}, 5000), @@ -284,7 +284,7 @@ pool_norestart_crash(_Config) -> -spec stop(config()) -> {comment, []}. stop(_Config) -> ct:comment("cast_call with stop/reply"), - {ok, Pid1} = wpool_process:start_link(stopper, echo_server, {ok, state}, #{}), + {ok, Pid1} = wpool_process:start_link(stopper, echo_server, {ok, state}, opts()), ReqId1 = wpool_process:send_request(stopper, {stop, reason, response, state}), case gen_server:wait_response(ReqId1, 5000) of {reply, response} -> @@ -300,7 +300,7 @@ stop(_Config) -> end, ct:comment("cast_call with regular stop"), - {ok, Pid2} = wpool_process:start_link(stopper, echo_server, {ok, state}, #{}), + {ok, Pid2} = wpool_process:start_link(stopper, echo_server, {ok, state}, opts()), ReqId2 = wpool_process:send_request(stopper, {stop, reason, state}), case gen_server:wait_response(ReqId2, 500) of {error, {reason, Pid2}} -> @@ -316,7 +316,7 @@ stop(_Config) -> end, ct:comment("call with regular stop"), - {ok, Pid3} = wpool_process:start_link(stopper, echo_server, {ok, state}, #{}), + {ok, Pid3} = wpool_process:start_link(stopper, echo_server, {ok, state}, opts()), try wpool_process:call(stopper, {noreply, state}, 100) of _ -> ct:fail("unexpected response") @@ -351,7 +351,7 @@ stop(_Config) -> -spec complete_coverage(config()) -> {comment, []}. complete_coverage(_Config) -> ct:comment("Code Change"), - {ok, State} = wpool_process:init({complete_coverage, echo_server, {ok, state}, []}), + {ok, State} = wpool_process:init({complete_coverage, echo_server, {ok, state}, opts()}), {ok, _} = wpool_process:code_change("oldvsn", State, {ok, state}), {error, bad} = wpool_process:code_change("oldvsn", State, {error, bad}), @@ -378,3 +378,6 @@ get_state(Pid) -> Misc ), wpool_process:get_state(State). + +opts() -> + #{event_manager => undefined, queue_manager => undefined}.