From 64a2a22617a942230873949cb0a11ac582741438 Mon Sep 17 00:00:00 2001 From: Nelson Vides Date: Sat, 21 Feb 2026 10:06:18 +0100 Subject: [PATCH 1/3] Reformat with the new formatter --- .github/workflows/erlang.yml | 2 +- elvis.config | 58 +++-- rebar.config | 97 +++++---- src/worker_pool.app.src | 33 +-- src/wpool.erl | 139 ++++++------ src/wpool_pool.erl | 268 +++++++++++++---------- src/wpool_process.erl | 161 ++++++++------ src/wpool_process_callbacks.erl | 22 +- src/wpool_process_sup.erl | 44 ++-- src/wpool_queue_manager.erl | 58 +++-- src/wpool_sup.erl | 17 +- src/wpool_time_checker.erl | 46 ++-- src/wpool_utils.erl | 36 ++-- src/wpool_worker.erl | 44 ++-- test/crashy_server.erl | 10 +- test/echo_server.erl | 12 +- test/echo_supervisor.erl | 22 +- test/wpool_SUITE.erl | 268 ++++++++++++++--------- test/wpool_bench.erl | 10 +- test/wpool_pool_SUITE.erl | 288 +++++++++++++++---------- test/wpool_process_SUITE.erl | 62 ++++-- test/wpool_process_callbacks_SUITE.erl | 105 +++++---- test/wpool_worker_SUITE.erl | 6 +- 23 files changed, 1089 insertions(+), 719 deletions(-) diff --git a/.github/workflows/erlang.yml b/.github/workflows/erlang.yml index bdee1f3..569c2cf 100644 --- a/.github/workflows/erlang.yml +++ b/.github/workflows/erlang.yml @@ -37,7 +37,7 @@ jobs: - name: Compile run: rebar3 compile - name: Format check - run: rebar3 format --verify + run: rebar3 fmt --check - name: Run tests and verifications run: rebar3 test - name: Upload code coverage diff --git a/elvis.config b/elvis.config index ed17396..1376f78 100644 --- a/elvis.config +++ b/elvis.config @@ -1,20 +1,38 @@ -[{elvis, - [{config, - [#{dirs => ["src"], - filter => "*.erl", - ruleset => erl_files, - rules => - [{elvis_style, invalid_dynamic_call, #{ignore => [wpool_process, wpool_time_checker]}}, - {elvis_style, god_modules, #{limit => 30}}, - {elvis_style, state_record_and_type, disable}, - {elvis_style, dont_repeat_yourself, #{ignore => [wpool_SUITE], min_complexity => 13}}]}, - #{dirs => ["test"], - filter => "*.erl", - ruleset => erl_files, - rules => - [{elvis_style, no_debug_call, disable}, - {elvis_style, state_record_and_type, disable}, - {elvis_style, dont_repeat_yourself, #{min_complexity => 13}}]}, - #{dirs => ["."], - filter => "elvis.config", - ruleset => elvis_config}]}]}]. +[ + {elvis, [ + {config, [ + #{ + dirs => ["src"], + filter => "*.erl", + ruleset => erl_files, + rules => + [ + {elvis_style, invalid_dynamic_call, #{ + ignore => [wpool_process, wpool_time_checker] + }}, + {elvis_style, god_modules, #{limit => 30}}, + {elvis_style, state_record_and_type, disable}, + {elvis_style, dont_repeat_yourself, #{ + ignore => [wpool_SUITE], min_complexity => 13 + }} + ] + }, + #{ + dirs => ["test"], + filter => "*.erl", + ruleset => erl_files, + rules => + [ + {elvis_style, no_debug_call, disable}, + {elvis_style, state_record_and_type, disable}, + {elvis_style, dont_repeat_yourself, #{min_complexity => 13}} + ] + }, + #{ + dirs => ["."], + filter => "elvis.config", + ruleset => elvis_config + } + ]} + ]} +]. diff --git a/rebar.config b/rebar.config index 4ebf319..101bfb5 100644 --- a/rebar.config +++ b/rebar.config @@ -1,69 +1,76 @@ %% == Compiler and Profiles == -{erl_opts, - [warn_unused_import, warn_export_vars, warnings_as_errors, verbose, report, debug_info]}. +{erl_opts, [warn_unused_import, warn_export_vars, warnings_as_errors, verbose, report, debug_info]}. {minimum_otp_vsn, "27"}. -{profiles, - [{test, - [{ct_extra_params, - "-no_auto_compile -dir ebin -logdir log/ct --erl_args -smp enable -boot start_sasl"}, - {cover_enabled, true}, - {cover_export_enabled, true}, - {cover_opts, [verbose, {min_coverage, 92}]}, - {ct_opts, [{verbose, true}]}, - {deps, [{katana, "1.0.0"}, {mixer, "1.2.0", {pkg, inaka_mixer}}, {meck, "1.0.0"}]}, - {dialyzer, - [{warnings, [no_return, unmatched_returns, error_handling, underspecs, unknown]}, - {plt_extra_apps, [common_test, katana, meck]}]}]}]}. - -{alias, - [{test, - [compile, - format, - hank, - lint, - xref, - dialyzer, - ct, - cover, - {covertool, "generate"}, - ex_doc]}]}. +{profiles, [ + {test, [ + {ct_extra_params, + "-no_auto_compile -dir ebin -logdir log/ct --erl_args -smp enable -boot start_sasl"}, + {cover_enabled, true}, + {cover_export_enabled, true}, + {cover_opts, [verbose, {min_coverage, 92}]}, + {ct_opts, [{verbose, true}]}, + {deps, [{katana, "1.0.0"}, {mixer, "1.2.0", {pkg, inaka_mixer}}, {meck, "1.0.0"}]}, + {dialyzer, [ + {warnings, [no_return, unmatched_returns, error_handling, underspecs, unknown]}, + {plt_extra_apps, [common_test, katana, meck]} + ]} + ]} +]}. + +{alias, [ + {test, [ + compile, + fmt, + hank, + lint, + xref, + dialyzer, + ct, + cover, + {covertool, "generate"}, + ex_doc + ]} +]}. {covertool, [{coverdata_files, ["ct.coverdata"]}]}. %% == Dependencies and plugins == -{project_plugins, - [{rebar3_hank, "~> 1.4.1"}, - {rebar3_hex, "~> 7.0.11"}, - {rebar3_format, "~> 1.3.0"}, - {rebar3_lint, "~> 4.1.0"}, - {rebar3_ex_doc, "~> 0.2.30"}, - {rebar3_depup, "~> 0.4.0"}, - {covertool, "~> 2.0.7"}]}. +{project_plugins, [ + {rebar3_hank, "~> 1.4.1"}, + {rebar3_depup, "~> 0.4.0"}, + {rebar3_hex, "~> 7.0"}, + {rebar3_ex_doc, "~> 0.2"}, + {rebar3_lint, "~> 4.2"}, + {erlfmt, "~> 1.7"}, + {covertool, "~> 2.0"} +]}. %% == Documentation == -{ex_doc, - [{source_url, <<"https://github.com/inaka/worker_pool">>}, - {extras, [<<"README.md">>, <<"LICENSE">>]}, - {main, <<"README.md">>}, - {prefix_ref_vsn_with_v, false}]}. +{ex_doc, [ + {source_url, <<"https://github.com/inaka/worker_pool">>}, + {extras, [<<"README.md">>, <<"LICENSE">>]}, + {main, <<"README.md">>}, + {prefix_ref_vsn_with_v, false} +]}. {hex, [{doc, #{provider => ex_doc}}]}. %% == Format == -{format, [{files, ["*.config", "src/*", "test/*"]}]}. +{erlfmt, [ + write, + {files, ["config/**/*.config", "src/**/*.app.src", "src/**/*.erl", "test/*.erl", "*.config"]} +]}. %% == Dialyzer + XRef == -{dialyzer, - [{warnings, [no_return, unmatched_returns, error_handling, underspecs, unknown]}]}. +{dialyzer, [{warnings, [no_return, unmatched_returns, error_handling, underspecs, unknown]}]}. -{xref_checks, - [undefined_function_calls, deprecated_function_calls, deprecated_functions]}. +{xref_checks, [undefined_function_calls, deprecated_function_calls, deprecated_functions]}. {xref_extra_paths, ["test/**"]}. diff --git a/src/worker_pool.app.src b/src/worker_pool.app.src index b752e7e..8bfafcd 100644 --- a/src/worker_pool.app.src +++ b/src/worker_pool.app.src @@ -14,19 +14,20 @@ % specific language governing permissions and limitations % under the License. -{application, - worker_pool, - [{description, "Erlang Worker Pool"}, - {vsn, git}, - {id, "worker_pool"}, - {registered, []}, - {modules, []}, - {applications, [kernel, stdlib]}, - {mod, {wpool, []}}, - {env, []}, - {licenses, ["Apache2"]}, - {links, - [{"GitHub", "https://github.com/inaka/worker_pool"}, - {"Blog Post", - "https://web.archive.org/web/20170602054156/https://inaka.net/blog/2014/09/25/worker-pool/"}]}, - {build_tools, ["rebar3"]}]}. +{application, worker_pool, [ + {description, "Erlang Worker Pool"}, + {vsn, git}, + {id, "worker_pool"}, + {registered, []}, + {modules, []}, + {applications, [kernel, stdlib]}, + {mod, {wpool, []}}, + {env, []}, + {licenses, ["Apache2"]}, + {links, [ + {"GitHub", "https://github.com/inaka/worker_pool"}, + {"Blog Post", + "https://web.archive.org/web/20170602054156/https://inaka.net/blog/2014/09/25/worker-pool/"} + ]}, + {build_tools, ["rebar3"]} +]}. diff --git a/src/wpool.erl b/src/wpool.erl index 0441468..5c2a9f6 100644 --- a/src/wpool.erl +++ b/src/wpool.erl @@ -56,7 +56,7 @@ -module(wpool). %% @todo remove this line when https://github.com/AdRoll/rebar3_format/issues/356 is fixed --format ignore. +-format(ignore). -behaviour(application). @@ -195,42 +195,44 @@ %% Name of the pool -type option() :: - {workers, workers()} | - {worker, worker()} | - {worker_opt, [worker_opt()]} | - {strategy, supervisor_strategy()} | - {worker_shutdown, worker_shutdown()} | - {overrun_handler, overrun_handler() | [overrun_handler()]} | - {overrun_warning, overrun_warning()} | - {max_overrun_warnings, max_overrun_warnings()} | - {pool_sup_intensity, pool_sup_intensity()} | - {pool_sup_shutdown, pool_sup_shutdown()} | - {pool_sup_period, pool_sup_period()} | - {queue_type, queue_type()} | - {enable_callbacks, enable_callbacks()} | - {enable_queues, enable_queues()} | - {callbacks, callbacks()}. + {workers, workers()} + | {worker, worker()} + | {worker_opt, [worker_opt()]} + | {strategy, supervisor_strategy()} + | {worker_shutdown, worker_shutdown()} + | {overrun_handler, overrun_handler() | [overrun_handler()]} + | {overrun_warning, overrun_warning()} + | {max_overrun_warnings, max_overrun_warnings()} + | {pool_sup_intensity, pool_sup_intensity()} + | {pool_sup_shutdown, pool_sup_shutdown()} + | {pool_sup_period, pool_sup_period()} + | {queue_type, queue_type()} + | {enable_callbacks, enable_callbacks()} + | {enable_queues, enable_queues()} + | {callbacks, callbacks()}. %% Options that can be provided to a new pool. %% %% `child_spec/2', `start_pool/2', `start_sup_pool/2' are the callbacks %% that take a list of these options as a parameter. --type options() :: #{workers => workers(), - worker => worker(), - worker_opt => [worker_opt()], - strategy => supervisor_strategy(), - worker_shutdown => worker_shutdown(), - overrun_handler => overrun_handler() | [overrun_handler()], - overrun_warning => overrun_warning(), - max_overrun_warnings => max_overrun_warnings(), - pool_sup_intensity => pool_sup_intensity(), - pool_sup_shutdown => pool_sup_shutdown(), - pool_sup_period => pool_sup_period(), - queue_type => queue_type(), - enable_callbacks => enable_callbacks(), - enable_queues => enable_queues(), - callbacks => callbacks(), - _ => _}. +-type options() :: #{ + workers => workers(), + worker => worker(), + worker_opt => [worker_opt()], + strategy => supervisor_strategy(), + worker_shutdown => worker_shutdown(), + overrun_handler => overrun_handler() | [overrun_handler()], + overrun_warning => overrun_warning(), + max_overrun_warnings => max_overrun_warnings(), + pool_sup_intensity => pool_sup_intensity(), + pool_sup_shutdown => pool_sup_shutdown(), + pool_sup_period => pool_sup_period(), + queue_type => queue_type(), + enable_callbacks => enable_callbacks(), + enable_queues => enable_queues(), + callbacks => callbacks(), + _ => _ +}. %% Options that can be provided to a new pool. %% %% `child_spec/2', `start_pool/2', `start_sup_pool/2' are the callbacks @@ -240,13 +242,13 @@ %% A callback that gets the pool name and returns a worker's name. -type strategy() :: - best_worker | - random_worker | - next_worker | - available_worker | - next_available_worker | - {hash_worker, term()} | - custom_strategy(). + best_worker + | random_worker + | next_worker + | available_worker + | next_available_worker + | {hash_worker, term()} + | custom_strategy(). %% Strategy to use when choosing a worker. %% %%

`best_worker'

@@ -293,24 +295,40 @@ %% Statistics about a worker in a pool. -type stats() :: - [{pool, name()} | - {supervisor, pid()} | - {options, [option()] | options()} | - {size, non_neg_integer()} | - {next_worker, pos_integer()} | - {total_message_queue_len, non_neg_integer()} | - {workers, [{pos_integer(), worker_stats()}]}]. + [ + {pool, name()} + | {supervisor, pid()} + | {options, [option()] | options()} + | {size, non_neg_integer()} + | {next_worker, pos_integer()} + | {total_message_queue_len, non_neg_integer()} + | {workers, [{pos_integer(), worker_stats()}]} + ]. %% Statistics about a given live pool. --export_type([name/0, option/0, options/0, custom_strategy/0, strategy/0, - queue_type/0, run/1, worker_stats/0, stats/0]). +-export_type([ + name/0, + option/0, + options/0, + custom_strategy/0, + strategy/0, + queue_type/0, + run/1, + worker_stats/0, + stats/0 +]). -export([start/0, start/2, stop/0, stop/1]). -export([child_spec/2, start_pool/1, start_pool/2, start_sup_pool/1, start_sup_pool/2]). -export([stop_pool/1, stop_sup_pool/1]). --export([call/2, call/3, call/4, cast/2, cast/3, - run/2, run/3, run/4, broadcall/3, broadcast/2, - send_request/2, send_request/3, send_request/4]). +-export([ + call/2, call/3, call/4, + cast/2, cast/3, + run/2, run/3, run/4, + broadcall/3, + broadcast/2, + send_request/2, send_request/3, send_request/4 +]). -export([stats/0, stats/1, get_workers/1]). -export([default_strategy/0]). @@ -359,11 +377,13 @@ start_pool(Name, Options) -> -spec child_spec(name(), [option()] | options()) -> supervisor:child_spec(). child_spec(Name, Options) -> FullOptions = wpool_utils:add_defaults(Options), - #{id => Name, - start => {wpool, start_pool, [Name, FullOptions]}, - restart => permanent, - shutdown => infinity, - type => supervisor}. + #{ + id => Name, + start => {wpool, start_pool, [Name, FullOptions]}, + restart => permanent, + shutdown => infinity, + type => supervisor + }. %% @doc Stops a pool that doesn't belong to `wpool_sup'. -spec stop_pool(name()) -> true. @@ -493,7 +513,7 @@ send_request(Sup, Call) -> %% @equiv send_request(Sup, Call, Strategy, 5000) -spec send_request(name(), term(), strategy()) -> - noproc | timeout | gen_server:request_id(). + noproc | timeout | gen_server:request_id(). send_request(Sup, Call, Strategy) -> send_request(Sup, Call, Strategy, 5000). @@ -501,7 +521,7 @@ send_request(Sup, Call, Strategy) -> %% %% Timeout applies only for the time used choosing a worker in the available_worker strategy -spec send_request(name(), term(), strategy(), timeout()) -> - noproc | timeout | gen_server:request_id(). + noproc | timeout | gen_server:request_id(). send_request(Sup, Call, available_worker, Timeout) -> wpool_pool:send_request_available_worker(Sup, Call, Timeout); send_request(Sup, Call, next_available_worker, _Timeout) -> @@ -517,7 +537,6 @@ send_request(Sup, Call, {hash_worker, HashKey}, _Timeout) -> send_request(Sup, Call, Fun, _Timeout) when is_function(Fun, 1) -> wpool_process:send_request(Fun(Sup), Call). - %% @doc Casts a message to all the workers within the given pool. %% %% NOTE: These messages don't get queued, they go straight to the worker's message queues, so @@ -532,7 +551,7 @@ broadcast(Sup, Cast) -> %% %% If one worker times out, the entire call is considered timed-out. -spec broadcall(wpool:name(), term(), timeout()) -> - {[Replies :: term()], [Errors :: term()]}. + {[Replies :: term()], [Errors :: term()]}. broadcall(Sup, Call, Timeout) -> wpool_pool:broadcall(Sup, Call, Timeout). diff --git a/src/wpool_pool.erl b/src/wpool_pool.erl index 1f846e1..2272ff6 100644 --- a/src/wpool_pool.erl +++ b/src/wpool_pool.erl @@ -27,9 +27,16 @@ %% API -export([start_link/2]). --export([best_worker/1, random_worker/1, next_worker/1, hash_worker/2, - next_available_worker/1, send_request_available_worker/3, call_available_worker/3, - run_with_available_worker/3]). +-export([ + best_worker/1, + random_worker/1, + next_worker/1, + hash_worker/2, + next_available_worker/1, + send_request_available_worker/3, + call_available_worker/3, + run_with_available_worker/3 +]). -export([cast_to_available_worker/2, broadcast/2, broadcall/3]). -export([stats/0, stats/1, get_workers/1]). -export([worker_name/2, find_wpool/1]). @@ -38,14 +45,15 @@ %% Supervisor callbacks -export([init/1]). --record(wpool, - {name :: wpool:name(), - size :: pos_integer(), - next :: atomics:atomics_ref(), - workers :: tuple(), - opts :: wpool:options(), - qmanager :: wpool_queue_manager:queue_mgr(), - born = erlang:system_time(second) :: integer()}). +-record(wpool, { + name :: wpool:name(), + size :: pos_integer(), + next :: atomics:atomics_ref(), + workers :: tuple(), + opts :: wpool:options(), + qmanager :: wpool_queue_manager:queue_mgr(), + born = erlang:system_time(second) :: integer() +}). -opaque wpool() :: #wpool{}. @@ -156,11 +164,13 @@ call_available_worker(Name, Call, Timeout) -> %% @doc Picks the first available worker and sends the request to it. %% The timeout provided considers only the time it takes to get a worker -spec send_request_available_worker(wpool:name(), any(), timeout()) -> - noproc | timeout | gen_server:request_id(). + noproc | timeout | gen_server:request_id(). send_request_available_worker(Name, Call, Timeout) -> - wpool_queue_manager:send_request_available_worker(queue_manager_name(Name), - Call, - Timeout). + wpool_queue_manager:send_request_available_worker( + queue_manager_name(Name), + Call, + Timeout + ). %% @doc Picks a worker base on a hash result. %%
phash2(Term, Range)
returns hash = integer, @@ -187,15 +197,17 @@ cast_to_available_worker(Name, Cast) -> %% @doc Casts a message to all the workers within the given pool. -spec broadcast(wpool:name(), term()) -> ok. broadcast(Name, Cast) -> - lists:foreach(fun(Worker) -> ok = wpool_process:cast(Worker, Cast) end, - all_workers(Name)). + lists:foreach( + fun(Worker) -> ok = wpool_process:cast(Worker, Cast) end, + all_workers(Name) + ). %% @doc Calls all workers in the pool in parallel %% %% Waits for responses in parallel too, and it assumes that if any response times out, %% all of them did too and therefore exits with reason timeout like a regular `gen_server' does. -spec broadcall(wpool:name(), term(), timeout()) -> - {[Replies :: term()], [Errors :: term()]}. + {[Replies :: term()], [Errors :: term()]}. broadcall(Name, Call, Timeout) -> Workers = all_workers(Name), ReqId0 = gen_server:reqids_new(), @@ -203,24 +215,26 @@ broadcall(Name, Call, Timeout) -> ReqId1 = lists:foldl(RequestFold, ReqId0, Workers), WaitFold = fun(_, {Coll, Replies, Errors}) -> - case gen_server:receive_response(Coll, Timeout, true) of - {{reply, Reply}, _, Coll1} -> - {Coll1, [Reply | Replies], Errors}; - {{error, Error}, _, Coll1} -> - {Coll1, Replies, [Error | Errors]}; - timeout -> - exit({timeout, {?MODULE, broadcall, [Name, Call, Timeout]}}) - end + case gen_server:receive_response(Coll, Timeout, true) of + {{reply, Reply}, _, Coll1} -> + {Coll1, [Reply | Replies], Errors}; + {{error, Error}, _, Coll1} -> + {Coll1, Replies, [Error | Errors]}; + timeout -> + exit({timeout, {?MODULE, broadcall, [Name, Call, Timeout]}}) + end end, {_, Replies, Errors} = lists:foldl(WaitFold, {ReqId1, [], []}, Workers), {Replies, Errors}. -spec all() -> [wpool:name()]. all() -> - [Name + [ + Name || {{?MODULE, Name}, _} <- persistent_term:get(), is_atom(Name), - find_wpool(Name) /= undefined]. + find_wpool(Name) /= undefined + ]. %% @doc Retrieves the list of worker registered names. %% This can be useful to manually inspect the workers or do custom work on them. @@ -246,38 +260,50 @@ stats(Name) -> stats(Wpool, Name) -> {Total, WorkerStats} = - lists:foldl(fun(N, {T, L}) -> - case worker_info(Wpool, - N, - [message_queue_len, - memory, - current_function, - current_location, - dictionary]) - of - undefined -> - {T, L}; - [{message_queue_len, MQL} = MQLT, - Memory, - Function, - Location, - {dictionary, Dictionary}] -> - WS = [MQLT, Memory] - ++ function_location(Function, Location) - ++ task(proplists:get_value(wpool_task, Dictionary)), - {T + MQL, [{N, WS} | L]} - end - end, - {0, []}, - lists:seq(1, Wpool#wpool.size)), + lists:foldl( + fun(N, {T, L}) -> + case + worker_info( + Wpool, + N, + [ + message_queue_len, + memory, + current_function, + current_location, + dictionary + ] + ) + of + undefined -> + {T, L}; + [ + {message_queue_len, MQL} = MQLT, + Memory, + Function, + Location, + {dictionary, Dictionary} + ] -> + WS = + [MQLT, Memory] ++ + function_location(Function, Location) ++ + task(proplists:get_value(wpool_task, Dictionary)), + {T + MQL, [{N, WS} | L]} + end + end, + {0, []}, + lists:seq(1, Wpool#wpool.size) + ), PendingTasks = wpool_queue_manager:pending_task_count(Wpool#wpool.qmanager), - [{pool, Name}, - {supervisor, erlang:whereis(Name)}, - {options, maps:to_list(Wpool#wpool.opts)}, - {size, Wpool#wpool.size}, - {next_worker, atomics:get(Wpool#wpool.next, 1)}, - {total_message_queue_len, Total + PendingTasks}, - {workers, WorkerStats}]. + [ + {pool, Name}, + {supervisor, erlang:whereis(Name)}, + {options, maps:to_list(Wpool#wpool.opts)}, + {size, Wpool#wpool.size}, + {next_worker, atomics:get(Wpool#wpool.next, 1)}, + {total_message_queue_len, Total + PendingTasks}, + {workers, WorkerStats} + ]. worker_info(Wpool, N, Info) -> case erlang:whereis(nth_worker_name(Wpool, N)) of @@ -322,8 +348,9 @@ remove_callback_module(Pool, Module) -> %% @doc Get values from the worker pool record. Useful when using a custom %% strategy function. --spec wpool_get(atom(), wpool()) -> any(); - ([atom()], wpool()) -> any(). +-spec wpool_get + (atom(), wpool()) -> any(); + ([atom()], wpool()) -> any(). wpool_get(List, Wpool) when is_list(List) -> [g(Atom, Wpool) || Atom <- List]; wpool_get(Atom, Wpool) when is_atom(Atom) -> @@ -351,7 +378,7 @@ time_checker_name(Name) -> %% =================================================================== %% @private -spec init({wpool:name(), wpool:options()}) -> - {ok, {supervisor:sup_flags(), [supervisor:child_spec()]}}. + {ok, {supervisor:sup_flags(), [supervisor:child_spec()]}}. init({Name, Options}) -> Size = maps:get(workers, Options, 100), QueueType = maps:get(queue_type, Options), @@ -364,58 +391,63 @@ init({Name, Options}) -> _Wpool = store_wpool(Name, Size, Options), WorkerOpts0 = - [{time_checker, TimeCheckerName}] - ++ maybe_queue_manager(Options, {queue_manager, QueueManagerName}) - ++ maybe_event_manager(Options, {event_manager, EventManagerName}), + [{time_checker, TimeCheckerName}] ++ + maybe_queue_manager(Options, {queue_manager, QueueManagerName}) ++ + maybe_event_manager(Options, {event_manager, EventManagerName}), WorkerOpts = maps:merge( - maps:from_list(WorkerOpts0), Options), + maps:from_list(WorkerOpts0), Options + ), TimeCheckerSpec = - #{id => TimeCheckerName, - start => {wpool_time_checker, start_link, [Name, TimeCheckerName, OverrunHandler]}, - restart => permanent, - shutdown => brutal_kill, - type => worker, - modules => [wpool_time_checker]}, + #{ + id => TimeCheckerName, + start => {wpool_time_checker, start_link, [Name, TimeCheckerName, OverrunHandler]}, + restart => permanent, + shutdown => brutal_kill, + type => worker, + modules => [wpool_time_checker] + }, QueueManagerSpec = - #{id => QueueManagerName, - start => - {wpool_queue_manager, - start_link, - [Name, QueueManagerName, [{queue_type, QueueType}]]}, - restart => permanent, - shutdown => brutal_kill, - type => worker, - modules => [wpool_queue_manager]}, + #{ + id => QueueManagerName, + start => + {wpool_queue_manager, start_link, [ + Name, QueueManagerName, [{queue_type, QueueType}] + ]}, + restart => permanent, + shutdown => brutal_kill, + type => worker, + modules => [wpool_queue_manager] + }, EventManagerSpec = - #{id => EventManagerName, - start => {gen_event, start_link, [{local, EventManagerName}]}, - restart => permanent, - shutdown => brutal_kill, - type => worker, - modules => dynamic}, + #{ + id => EventManagerName, + start => {gen_event, start_link, [{local, EventManagerName}]}, + restart => permanent, + shutdown => brutal_kill, + type => worker, + modules => dynamic + }, ProcessSupSpec = - {ProcessSupName, - {wpool_process_sup, start_link, [Name, ProcessSupName, WorkerOpts]}, - permanent, - SupShutdown, - supervisor, - [wpool_process_sup]}, + {ProcessSupName, {wpool_process_sup, start_link, [Name, ProcessSupName, WorkerOpts]}, + permanent, SupShutdown, supervisor, [wpool_process_sup]}, Children = - [TimeCheckerSpec] - ++ maybe_queue_manager(Options, QueueManagerSpec) - ++ maybe_event_manager(Options, EventManagerSpec) - ++ [ProcessSupSpec], + [TimeCheckerSpec] ++ + maybe_queue_manager(Options, QueueManagerSpec) ++ + maybe_event_manager(Options, EventManagerSpec) ++ + [ProcessSupSpec], SupIntensity = maps:get(pool_sup_intensity, Options, 5), SupPeriod = maps:get(pool_sup_period, Options, 60), SupStrategy = - #{strategy => one_for_all, - intensity => SupIntensity, - period => SupPeriod}, + #{ + strategy => one_for_all, + intensity => SupIntensity, + period => SupPeriod + }, {ok, {SupStrategy, Children}}. %% @private @@ -525,12 +557,14 @@ store_wpool(Name, Size, Options) -> atomics:put(Atomic, 1, 1), WorkerNames = list_to_tuple([worker_name(Name, I) || I <- lists:seq(1, Size)]), Wpool = - #wpool{name = Name, - size = Size, - next = Atomic, - workers = WorkerNames, - opts = Options, - qmanager = queue_manager_name(Name)}, + #wpool{ + name = Name, + size = Size, + next = Atomic, + workers = WorkerNames, + opts = Options, + qmanager = queue_manager_name(Name) + }, persistent_term:put({?MODULE, Name}, Wpool), Wpool. @@ -550,19 +584,27 @@ find_wpool(Name) -> %% @doc We use this function not to report an error if for some reason we've %% lost the record on the persistent_term table. This SHOULDN'T be called too much. build_wpool(Name) -> - logger:warning(#{what => "Building a #wpool record. Something must have failed.", - pool => Name}, - ?LOCATION), + logger:warning( + #{ + what => "Building a #wpool record. Something must have failed.", + pool => Name + }, + ?LOCATION + ), try supervisor:count_children(process_sup_name(Name)) of Children -> Size = proplists:get_value(active, Children, 0), store_wpool(Name, Size, #{}) catch _:Error -> - logger:warning(#{what => "Wpool not found", - pool => Name, - reason => Error}, - ?LOCATION), + logger:warning( + #{ + what => "Wpool not found", + pool => Name, + reason => Error + }, + ?LOCATION + ), undefined end. diff --git a/src/wpool_process.erl b/src/wpool_process.erl index d752911..b9df7f4 100644 --- a/src/wpool_process.erl +++ b/src/wpool_process.erl @@ -19,38 +19,46 @@ -behaviour(gen_server). %% Taken from gen_server OTP --record(callback_cache, - {module :: module(), - handle_call :: - fun((Request :: term(), From :: from(), State :: term()) -> - {reply, Reply :: term(), NewState :: term()} | - {reply, - Reply :: term(), - NewState :: term(), - timeout() | hibernate | {continue, term()}} | - {noreply, NewState :: term()} | - {noreply, NewState :: term(), timeout() | hibernate | {continue, term()}} | - {stop, Reason :: term(), Reply :: term(), NewState :: term()} | - {stop, Reason :: term(), NewState :: term()}), - handle_cast :: - fun((Request :: term(), State :: term()) -> - {noreply, NewState :: term()} | - {noreply, NewState :: term(), timeout() | hibernate | {continue, term()}} | - {stop, Reason :: term(), NewState :: term()}), - handle_info :: - fun((Info :: timeout | term(), State :: term()) -> - {noreply, NewState :: term()} | - {noreply, NewState :: term(), timeout() | hibernate | {continue, term()}} | - {stop, Reason :: term(), NewState :: term()})}). --record(state, - {name :: atom(), - mod :: #callback_cache{}, - state :: term(), - options :: - #{time_checker := atom(), - queue_manager := atom(), - overrun_warning := timeout(), - _ => _}}). +-record(callback_cache, { + module :: module(), + handle_call :: + fun( + (Request :: term(), From :: from(), State :: term()) -> + {reply, Reply :: term(), NewState :: term()} + | {reply, Reply :: term(), NewState :: term(), + timeout() | hibernate | {continue, term()}} + | {noreply, NewState :: term()} + | {noreply, NewState :: term(), timeout() | hibernate | {continue, term()}} + | {stop, Reason :: term(), Reply :: term(), NewState :: term()} + | {stop, Reason :: term(), NewState :: term()} + ), + handle_cast :: + fun( + (Request :: term(), State :: term()) -> + {noreply, NewState :: term()} + | {noreply, NewState :: term(), timeout() | hibernate | {continue, term()}} + | {stop, Reason :: term(), NewState :: term()} + ), + handle_info :: + fun( + (Info :: timeout | term(), State :: term()) -> + {noreply, NewState :: term()} + | {noreply, NewState :: term(), timeout() | hibernate | {continue, term()}} + | {stop, Reason :: term(), NewState :: term()} + ) +}). +-record(state, { + name :: atom(), + mod :: #callback_cache{}, + state :: term(), + options :: + #{ + time_checker := atom(), + queue_manager := atom(), + overrun_warning := timeout(), + _ => _ + } +}). -opaque state() :: #state{}. @@ -74,22 +82,32 @@ -endif. %% gen_server callbacks --export([init/1, terminate/2, code_change/3, handle_call/3, handle_cast/2, handle_info/2, - handle_continue/2, format_status/1]). +-export([ + init/1, + terminate/2, + code_change/3, + handle_call/3, + handle_cast/2, + handle_info/2, + handle_continue/2, + format_status/1 +]). %%%=================================================================== %%% API %%%=================================================================== %% @doc Starts a named process -spec start_link(wpool:name(), module(), term(), wpool:options()) -> - {ok, pid()} | ignore | {error, {already_started, pid()} | term()}. + {ok, pid()} | ignore | {error, {already_started, pid()} | term()}. start_link(Name, Module, InitArgs, Options) -> FullOpts = wpool_utils:add_defaults(Options), WorkerOpt = maps:get(worker_opt, FullOpts, []), - gen_server:start_link({local, Name}, - ?MODULE, - {Name, Module, InitArgs, FullOpts}, - WorkerOpt). + gen_server:start_link( + {local, Name}, + ?MODULE, + {Name, Module, InitArgs, FullOpts}, + WorkerOpt + ). %% @doc Runs a function that takes as a parameter the given process -spec run(wpool:name() | pid(), wpool:run(Result), timeout()) -> Result. @@ -124,7 +142,7 @@ get_state(#state{state = State}) -> %%%=================================================================== %% @private -spec init({atom(), atom(), term(), wpool:options()}) -> - {ok, state()} | {ok, state(), next_step()} | {stop, can_not_ignore} | {stop, term()}. + {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]), CbCache = create_callback_cache(Mod), @@ -132,20 +150,23 @@ init({Name, Mod, InitArgs, Options}) -> {ok, ModState} -> ok = notify_queue_manager(new_worker, Name, Options), wpool_process_callbacks:notify(handle_worker_creation, Options, [Name]), - {ok, - #state{name = Name, - mod = CbCache, - state = ModState, - options = Options}}; + {ok, #state{ + name = Name, + mod = CbCache, + state = ModState, + options = Options + }}; {ok, ModState, NextStep} -> ok = notify_queue_manager(new_worker, Name, Options), wpool_process_callbacks:notify(handle_worker_creation, Options, [Name]), {ok, - #state{name = Name, + #state{ + name = Name, mod = CbCache, state = ModState, - options = Options}, - NextStep}; + options = Options + }, + NextStep}; ignore -> {stop, can_not_ignore}; Error -> @@ -155,10 +176,12 @@ init({Name, Mod, InitArgs, Options}) -> %% @private -spec terminate(atom(), state()) -> term(). terminate(Reason, State) -> - #state{mod = #callback_cache{module = Mod}, - state = ModState, - name = Name, - options = Options} = + #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]), @@ -171,7 +194,7 @@ terminate(Reason, State) -> %% @private -spec code_change(string() | {down, string()}, state(), any()) -> - {ok, state()} | {error, term()}. + {ok, state()} | {error, term()}. code_change(OldVsn, #state{mod = #callback_cache{module = Mod}} = State, Extra) -> case erlang:function_exported(Mod, code_change, 3) of true -> @@ -187,7 +210,7 @@ code_change(OldVsn, #state{mod = #callback_cache{module = Mod}} = State, Extra) %% @private -spec handle_info(any(), state()) -> - {noreply, state()} | {noreply, state(), next_step()} | {stop, term(), state()}. + {noreply, state()} | {noreply, state(), next_step()} | {stop, term(), state()}. handle_info(Info, #state{mod = CbCache} = State) -> #callback_cache{module = Mod, handle_info = HandleInfo} = CbCache, try HandleInfo(Info, State#state.state) of @@ -215,9 +238,9 @@ handle_info(Info, #state{mod = CbCache} = State) -> %% @private -spec handle_continue(any(), state()) -> - {noreply, state()} | - {noreply, state(), next_step()} | - {stop, term(), state()}. + {noreply, state()} + | {noreply, state(), next_step()} + | {stop, term(), state()}. handle_continue(Continue, #state{mod = #callback_cache{module = Mod}} = State) -> try Mod:handle_continue(Continue, State#state.state) of {noreply, NewState} -> @@ -250,7 +273,7 @@ format_status(#{state := #state{mod = #callback_cache{module = Mod}}} = Status) %%%=================================================================== %% @private -spec handle_cast(term(), state()) -> - {noreply, state()} | {noreply, state(), next_step()} | {stop, term(), state()}. + {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), @@ -277,12 +300,12 @@ handle_cast(Cast, #state{mod = CbCache, options = Options} = State) -> %% @private -spec handle_call(term(), from(), state()) -> - {reply, term(), state()} | - {reply, term(), state(), next_step()} | - {noreply, state()} | - {noreply, state(), next_step()} | - {stop, term(), term(), state()} | - {stop, term(), state()}. + {reply, term(), state()} + | {reply, term(), state(), next_step()} + | {noreply, state()} + | {noreply, state(), next_step()} + | {stop, term(), term(), 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), @@ -325,7 +348,9 @@ notify_queue_manager(_, _, _) -> ok. create_callback_cache(Mod) -> - #callback_cache{module = Mod, - handle_call = fun Mod:handle_call/3, - handle_cast = fun Mod:handle_cast/2, - handle_info = fun Mod:handle_info/2}. + #callback_cache{ + module = Mod, + handle_call = fun Mod:handle_call/3, + handle_cast = fun Mod:handle_cast/2, + handle_info = fun Mod:handle_info/2 + }. diff --git a/src/wpool_process_callbacks.erl b/src/wpool_process_callbacks.erl index c08fff0..19f1e42 100644 --- a/src/wpool_process_callbacks.erl +++ b/src/wpool_process_callbacks.erl @@ -22,8 +22,11 @@ -callback handle_worker_creation(wpool:name()) -> any(). -callback handle_worker_death(wpool:name(), term()) -> any(). --optional_callbacks([handle_init_start/1, handle_worker_creation/1, - handle_worker_death/2]). +-optional_callbacks([ + handle_init_start/1, + handle_worker_creation/1, + handle_worker_death/2 +]). %% @private -spec init(module()) -> {ok, state()}. @@ -73,17 +76,22 @@ call(Module, Event, Args) -> end catch E:R -> - logger:warning(#{what => "Could not call callback module", - error => E, - reason => R}, - ?LOCATION) + logger:warning( + #{ + what => "Could not call callback module", + error => E, + reason => R + }, + ?LOCATION + ) end. ensure_loaded(Module) -> case code:ensure_loaded(Module) of {module, Module} -> ok; - {error, embedded} -> %% We are in embedded mode so the module was loaded if exists + %% We are in embedded mode so the module was loaded if exists + {error, embedded} -> ok; Other -> Other diff --git a/src/wpool_process_sup.erl b/src/wpool_process_sup.erl index 4917d83..bbaf228 100644 --- a/src/wpool_process_sup.erl +++ b/src/wpool_process_sup.erl @@ -30,7 +30,7 @@ start_link(Parent, Name, Options) -> %% @private -spec init({wpool:name(), wpool:options()}) -> - {ok, {supervisor:sup_flags(), [supervisor:child_spec()]}}. + {ok, {supervisor:sup_flags(), [supervisor:child_spec()]}}. init({Name, Options}) -> Workers = maps:get(workers, Options, 100), Strategy = maps:get(strategy, Options, {one_for_one, 5, 60}), @@ -38,16 +38,20 @@ init({Name, Options}) -> {Worker, InitArgs} = maps:get(worker, Options, {wpool_worker, undefined}), maybe_add_event_handler(Options), WorkerSpecs = - [#{id => wpool_pool:worker_name(Name, I), - start => - {wpool_process, - start_link, - [wpool_pool:worker_name(Name, I), Worker, InitArgs, Options]}, - restart => permanent, - shutdown => WorkerShutdown, - type => worker, - modules => [Worker]} - || I <- lists:seq(1, Workers)], + [ + #{ + id => wpool_pool:worker_name(Name, I), + start => + {wpool_process, start_link, [ + wpool_pool:worker_name(Name, I), Worker, InitArgs, Options + ]}, + restart => permanent, + shutdown => WorkerShutdown, + type => worker, + modules => [Worker] + } + || I <- lists:seq(1, Workers) + ], {ok, {Strategy, WorkerSpecs}}. maybe_add_event_handler(Options) -> @@ -55,8 +59,10 @@ maybe_add_event_handler(Options) -> undefined -> ok; EventMgr -> - lists:foreach(fun(M) -> add_initial_callback(EventMgr, M) end, - maps:get(callbacks, Options, [])) + lists:foreach( + fun(M) -> add_initial_callback(EventMgr, M) end, + maps:get(callbacks, Options, []) + ) end. add_initial_callback(EventManager, Module) -> @@ -64,8 +70,12 @@ add_initial_callback(EventManager, Module) -> ok -> ok; Other -> - logger:warning(#{what => "The callback module could not be loaded", - module => Module, - reason => Other}, - ?LOCATION) + logger:warning( + #{ + what => "The callback module could not be loaded", + module => Module, + reason => Other + }, + ?LOCATION + ) end. diff --git a/src/wpool_queue_manager.erl b/src/wpool_queue_manager.erl index 787bac9..6d5bd40 100644 --- a/src/wpool_queue_manager.erl +++ b/src/wpool_queue_manager.erl @@ -18,18 +18,27 @@ %% api -export([start_link/2, start_link/3]). --export([run_with_available_worker/3, call_available_worker/3, cast_to_available_worker/2, - new_worker/2, worker_dead/2, send_request_available_worker/3, worker_ready/2, - worker_busy/2, pending_task_count/1]). +-export([ + run_with_available_worker/3, + call_available_worker/3, + cast_to_available_worker/2, + new_worker/2, + worker_dead/2, + send_request_available_worker/3, + worker_ready/2, + worker_busy/2, + pending_task_count/1 +]). %% gen_server callbacks -export([init/1, handle_call/3, handle_cast/2, handle_info/2]). --record(state, - {wpool :: wpool:name(), - clients :: queue:queue({cast | {pid(), _}, term()}), - workers :: gb_sets:set(atom()), - monitors :: #{atom() := monitored_from()}, - queue_type :: wpool:queue_type()}). +-record(state, { + wpool :: wpool:name(), + clients :: queue:queue({cast | {pid(), _}, term()}), + workers :: gb_sets:set(atom()), + monitors :: #{atom() := monitored_from()}, + queue_type :: wpool:queue_type() +}). -opaque state() :: #state{}. @@ -73,7 +82,7 @@ start_link(WPool, Name, Options) -> %% @doc returns the first available worker in the pool -spec run_with_available_worker(queue_mgr(), wpool:run(Result), timeout()) -> - noproc | timeout | Result. + noproc | timeout | Result. run_with_available_worker(QueueManager, Call, Timeout) -> case get_available_worker(QueueManager, Call, Timeout) of {ok, Worker, TimeLeft} when TimeLeft > 0 -> @@ -108,7 +117,7 @@ cast_to_available_worker(QueueManager, Cast) -> %% @doc returns the first available worker in the pool -spec send_request_available_worker(queue_mgr(), any(), timeout()) -> - noproc | timeout | gen_server:request_id(). + noproc | timeout | gen_server:request_id(). send_request_available_worker(QueueManager, Call, Timeout) -> case get_available_worker(QueueManager, Call, Timeout) of {ok, Worker, _} -> @@ -156,12 +165,13 @@ init(Args) -> WPool = proplists:get_value(pool, Args), QueueType = proplists:get_value(queue_type, Args), put(pending_tasks, 0), - {ok, - #state{wpool = WPool, - clients = queue:new(), - workers = gb_sets:new(), - monitors = #{}, - queue_type = QueueType}}. + {ok, #state{ + wpool = WPool, + clients = queue:new(), + workers = gb_sets:new(), + monitors = #{}, + queue_type = QueueType + }}. -spec handle_cast({worker_event(), atom()}, state()) -> {noreply, state()}. handle_cast({new_worker, Worker}, State) -> @@ -172,10 +182,12 @@ handle_cast({worker_dead, Worker}, #state{workers = Workers} = State) -> handle_cast({worker_busy, Worker}, #state{workers = Workers} = State) -> {noreply, State#state{workers = gb_sets:delete_any(Worker, Workers)}}; handle_cast({worker_ready, Worker}, State0) -> - #state{workers = Workers, - clients = Clients, - monitors = Mons, - queue_type = QueueType} = + #state{ + workers = Workers, + clients = Clients, + monitors = Mons, + queue_type = QueueType + } = State0, State = case Mons of @@ -217,7 +229,7 @@ handle_cast({cast_to_available_worker, Cast}, State) -> end. -spec handle_call(call_request(), from(), state()) -> - {reply, {ok, atom()}, state()} | {noreply, state()}. + {reply, {ok, atom()}, state()} | {noreply, state()}. handle_call({available_worker, ExpiresAt}, {ClientPid, _Ref} = Client, State) -> #state{workers = Workers, clients = Clients} = State, case gb_sets:is_empty(Workers) of @@ -256,7 +268,7 @@ handle_info(_Info, State) -> %%% private %%%=================================================================== -spec get_available_worker(queue_mgr(), any(), timeout()) -> - noproc | timeout | {ok, atom(), timeout()}. + noproc | timeout | {ok, atom(), timeout()}. get_available_worker(QueueManager, Call, Timeout) -> ExpiresAt = expires(Timeout), try gen_server:call(QueueManager, {available_worker, ExpiresAt}, Timeout) of diff --git a/src/wpool_sup.erl b/src/wpool_sup.erl index 9bd07dc..9522cfe 100644 --- a/src/wpool_sup.erl +++ b/src/wpool_sup.erl @@ -39,10 +39,14 @@ start_pool(Name, Options) -> stop_pool(Name) -> case erlang:whereis(Name) of undefined -> - logger:warning(#{what => "Could not stop pool", - reason => "It was not running", - pool => Name}, - ?LOCATION), + logger:warning( + #{ + what => "Could not stop pool", + reason => "It was not running", + pool => Name + }, + ?LOCATION + ), ok; Pid -> ok = supervisor:terminate_child(?MODULE, Pid) @@ -54,5 +58,6 @@ stop_pool(Name) -> -spec init([]) -> {ok, {{simple_one_for_one, 5, 60}, [supervisor:child_spec()]}}. init([]) -> {ok, - {{simple_one_for_one, 5, 60}, - [{wpool_pool, {wpool_pool, start_link, []}, permanent, 2000, supervisor, dynamic}]}}. + {{simple_one_for_one, 5, 60}, [ + {wpool_pool, {wpool_pool, start_link, []}, permanent, 2000, supervisor, dynamic} + ]}}. diff --git a/src/wpool_time_checker.erl b/src/wpool_time_checker.erl index b8d2bee..bc775a6 100644 --- a/src/wpool_time_checker.erl +++ b/src/wpool_time_checker.erl @@ -42,7 +42,7 @@ %%%=================================================================== %% @private -spec start_link(wpool:name(), atom(), handler() | [handler()]) -> - {ok, pid()} | {error, {already_started, pid()} | term()}. + {ok, pid()} | {error, {already_started, pid()} | term()}. start_link(WPool, Name, Handlers) when is_list(Handlers) -> gen_server:start_link({local, Name}, ?MODULE, {WPool, Handlers}, []); start_link(WPool, Name, Handler) when is_tuple(Handler) -> @@ -78,13 +78,15 @@ handle_call({add_handler, Handler}, _, #state{handlers = Handlers} = State) -> handle_info({check, Pid, TaskId, Runtime, WarningsLeft}, State) -> case erlang:process_info(Pid, dictionary) of {dictionary, Values} -> - run_task(TaskId, - proplists:get_value(wpool_task, Values), - Pid, - State#state.wpool, - State#state.handlers, - Runtime, - WarningsLeft); + run_task( + TaskId, + proplists:get_value(wpool_task, Values), + Pid, + State#state.wpool, + State#state.handlers, + Runtime, + WarningsLeft + ); _ -> ok end, @@ -100,13 +102,11 @@ run_task(TaskId, {TaskId, _, Task}, Pid, Pool, Handlers, Runtime, WarningsLeft) send_reports(Handlers, overrun, Pool, Pid, Task, Runtime), case new_overrun_time(Runtime, WarningsLeft) of NewOverrunTime when NewOverrunTime =< 4294967295 -> - erlang:send_after(Runtime, - self(), - {check, - Pid, - TaskId, - NewOverrunTime, - decrease_warnings(WarningsLeft)}), + erlang:send_after( + Runtime, + self(), + {check, Pid, TaskId, NewOverrunTime, decrease_warnings(WarningsLeft)} + ), ok; _ -> ok @@ -126,13 +126,15 @@ decrease_warnings(infinity) -> decrease_warnings(N) -> N - 1. --spec send_reports([{atom(), atom()}], - atom(), - atom(), - pid(), - term(), - infinity | pos_integer()) -> - ok. +-spec send_reports( + [{atom(), atom()}], + atom(), + atom(), + pid(), + term(), + infinity | pos_integer() +) -> + ok. send_reports(Handlers, Alert, Pool, Pid, Task, Runtime) -> Args = [{alert, Alert}, {pool, Pool}, {worker, Pid}, {task, Task}, {runtime, Runtime}], _ = [catch Mod:Fun(Args) || {Mod, Fun} <- Handlers], diff --git a/src/wpool_utils.erl b/src/wpool_utils.erl index 58ebe11..201338c 100644 --- a/src/wpool_utils.erl +++ b/src/wpool_utils.erl @@ -23,21 +23,27 @@ %% @doc Marks Task as started in this worker -spec task_init(term(), #{overrun_warning := timeout(), _ => _}) -> - undefined | reference(). + 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}) -> +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}). + erlang:send_after( + OverrunTime, + TimeChecker, + {check, self(), TaskId, OverrunTime, MaxWarnings} + ). %% @doc Removes the current task from the worker -spec task_end(undefined | reference()) -> ok. @@ -55,9 +61,11 @@ add_defaults(Opts) when is_list(Opts) -> maps:merge(defaults(), maps:from_list(Opts)). defaults() -> - #{max_overrun_warnings => infinity, - overrun_handler => {logger, warning}, - overrun_warning => infinity, - queue_type => fifo, - worker_opt => [], - workers => 100}. + #{ + max_overrun_warnings => infinity, + overrun_handler => {logger, warning}, + overrun_warning => infinity, + queue_type => fifo, + worker_opt => [], + workers => 100 + }. diff --git a/src/wpool_worker.erl b/src/wpool_worker.erl index 3c6485d..32ffafe 100644 --- a/src/wpool_worker.erl +++ b/src/wpool_worker.erl @@ -77,15 +77,19 @@ handle_cast({M, F, A}, State) -> {noreply, State, hibernate} end; handle_cast(Cast, State) -> - logger:error(#{what => "Invalid cast", - cast => Cast, - worker => self()}, - ?LOCATION), + logger:error( + #{ + what => "Invalid cast", + cast => Cast, + worker => self() + }, + ?LOCATION + ), {noreply, State, hibernate}. %% @private -spec handle_call(term(), from(), state()) -> - {reply, {ok, term()} | {error, term()}, state(), hibernate}. + {reply, {ok, term()} | {error, term()}, state(), hibernate}. handle_call({M, F, A}, _From, State) -> try erlang:apply(M, F, A) of R -> @@ -96,20 +100,28 @@ handle_call({M, F, A}, _From, State) -> {reply, {error, Reason}, State, hibernate} end; handle_call(Call, From, State) -> - logger:error(#{what => "Invalid call", - call => Call, - from => From, - worker => self()}, - ?LOCATION), + logger:error( + #{ + what => "Invalid call", + call => Call, + from => From, + worker => self() + }, + ?LOCATION + ), {reply, {error, invalid_request}, State, hibernate}. %%%=================================================================== %%% not exported functions %%%=================================================================== log_error(M, F, A, Class, Reason, Stacktrace) -> - logger:error(#{what => "Reason on ~p:~p~p >> ~p Backtrace ~p", - mfa => {M, F, A}, - class => Class, - reason => Reason, - stacktrace => Stacktrace}, - ?LOCATION). + logger:error( + #{ + what => "Reason on ~p:~p~p >> ~p Backtrace ~p", + mfa => {M, F, A}, + class => Class, + reason => Reason, + stacktrace => Stacktrace + }, + ?LOCATION + ). diff --git a/test/crashy_server.erl b/test/crashy_server.erl index 90ec009..ba1b6c5 100644 --- a/test/crashy_server.erl +++ b/test/crashy_server.erl @@ -17,8 +17,14 @@ -behaviour(gen_server). %% gen_server callbacks --export([init/1, terminate/2, code_change/3, handle_call/3, handle_cast/2, - handle_info/2]). +-export([ + init/1, + terminate/2, + code_change/3, + handle_call/3, + handle_cast/2, + handle_info/2 +]). -dialyzer([no_behaviours]). diff --git a/test/echo_server.erl b/test/echo_server.erl index d3b0d18..d1484c6 100644 --- a/test/echo_server.erl +++ b/test/echo_server.erl @@ -18,8 +18,16 @@ %% gen_server callbacks -export([start_link/1]). --export([init/1, terminate/2, code_change/3, handle_call/3, handle_cast/2, handle_info/2, - handle_continue/2, format_status/1]). +-export([ + init/1, + terminate/2, + code_change/3, + handle_call/3, + handle_cast/2, + handle_info/2, + handle_continue/2, + format_status/1 +]). -dialyzer([no_behaviours]). diff --git a/test/echo_supervisor.erl b/test/echo_supervisor.erl index 66b3df7..b619917 100644 --- a/test/echo_supervisor.erl +++ b/test/echo_supervisor.erl @@ -10,14 +10,18 @@ start_link() -> init(noargs) -> Children = - #{id => undefined, - start => {echo_server, start_link, []}, - restart => transient, - shutdown => 5000, - type => worker, - modules => [echo_server]}, + #{ + id => undefined, + start => {echo_server, start_link, []}, + restart => transient, + shutdown => 5000, + type => worker, + modules => [echo_server] + }, Strategy = - #{strategy => simple_one_for_one, - intensity => 5, - period => 60}, + #{ + strategy => simple_one_for_one, + intensity => 5, + period => 60 + }, {ok, {Strategy, [Children]}}. diff --git a/test/wpool_SUITE.erl b/test/wpool_SUITE.erl index e2febe3..3633de4 100644 --- a/test/wpool_SUITE.erl +++ b/test/wpool_SUITE.erl @@ -25,11 +25,27 @@ -export([all/0]). -export([init_per_suite/1, end_per_suite/1]). --export([stats/1, stop_pool/1, non_brutal_shutdown/1, brutal_worker_shutdown/1, overrun/1, - kill_on_overrun/1, too_much_overrun/1, default_strategy/1, overrun_handler1/1, - overrun_handler2/1, default_options/1, complete_coverage/1, child_spec/1, broadcall/1, - broadcast/1, send_request/1, worker_killed_stats/1, accepts_maps_and_lists_as_opts/1, - pool_of_supervisors/1]). +-export([ + stats/1, + stop_pool/1, + non_brutal_shutdown/1, + brutal_worker_shutdown/1, + overrun/1, + kill_on_overrun/1, + too_much_overrun/1, + default_strategy/1, + overrun_handler1/1, + overrun_handler2/1, + default_options/1, + complete_coverage/1, + child_spec/1, + broadcall/1, + broadcast/1, + send_request/1, + worker_killed_stats/1, + accepts_maps_and_lists_as_opts/1, + pool_of_supervisors/1 +]). -elvis([{elvis_style, no_block_expressions, disable}]). @@ -37,23 +53,25 @@ -spec all() -> [atom()]. all() -> - [too_much_overrun, - overrun, - stop_pool, - non_brutal_shutdown, - brutal_worker_shutdown, - stats, - default_strategy, - default_options, - complete_coverage, - child_spec, - broadcast, - broadcall, - send_request, - kill_on_overrun, - worker_killed_stats, - accepts_maps_and_lists_as_opts, - pool_of_supervisors]. + [ + too_much_overrun, + overrun, + stop_pool, + non_brutal_shutdown, + brutal_worker_shutdown, + stats, + default_strategy, + default_options, + complete_coverage, + child_spec, + broadcast, + broadcall, + send_request, + kill_on_overrun, + worker_killed_stats, + accepts_maps_and_lists_as_opts, + pool_of_supervisors + ]. -spec init_per_suite(config()) -> config(). init_per_suite(Config) -> @@ -78,10 +96,14 @@ too_much_overrun(_Config) -> ct:comment("Receiving overruns here..."), true = register(overrun_handler, self()), {ok, PoolPid} = - wpool:start_sup_pool(wpool_SUITE_too_much_overrun, - [{workers, 1}, - {overrun_warning, 999}, - {overrun_handler, {?MODULE, overrun_handler1}}]), + wpool:start_sup_pool( + wpool_SUITE_too_much_overrun, + [ + {workers, 1}, + {overrun_warning, 999}, + {overrun_handler, {?MODULE, overrun_handler1}} + ] + ), %% Sadly, the function that autogenerates this name is private. CheckerName = 'wpool_pool-wpool_SUITE_too_much_overrun-time-checker', @@ -96,17 +118,18 @@ too_much_overrun(_Config) -> ok = wpool:cast(wpool_SUITE_too_much_overrun, {timer, sleep, [5000]}), TaskId = ktn_task:wait_for_success(fun() -> - {dictionary, Dict} = erlang:process_info(Worker, dictionary), - {TId, _, _} = proplists:get_value(wpool_task, Dict), - TId - end), + {dictionary, Dict} = erlang:process_info(Worker, dictionary), + {TId, _, _} = proplists:get_value(wpool_task, Dict), + TId + end), ct:comment("Simulate overrun warning..."), % huge runtime => no more overruns TCPid ! {check, Worker, TaskId, 9999999999, infinity}, ct:comment("Get overrun message..."), - _ = receive + _ = + receive {overrun1, Message1} -> overrun = proplists:get_value(alert, Message1), wpool_SUITE_too_much_overrun = proplists:get_value(pool, Message1), @@ -118,7 +141,8 @@ too_much_overrun(_Config) -> end, ct:comment("Get overrun message..."), - _ = receive + _ = + receive {overrun2, Message2} -> overrun = proplists:get_value(alert, Message2), wpool_SUITE_too_much_overrun = proplists:get_value(pool, Message2), @@ -130,7 +154,8 @@ too_much_overrun(_Config) -> end, ct:comment("No more overruns..."), - _ = case get_messages(100) of + _ = + case get_messages(100) of [] -> ok; Msgs1 -> @@ -141,7 +166,8 @@ too_much_overrun(_Config) -> exit(Worker, kill), ct:comment("Simulate overrun warning..."), - TCPid ! {check, Worker, TaskId, 100}, % tiny runtime, to check + % tiny runtime, to check + TCPid ! {check, Worker, TaskId, 100}, ct:comment("Nothing happens..."), ok = no_messages(), @@ -155,12 +181,17 @@ too_much_overrun(_Config) -> overrun(_Config) -> true = register(overrun_handler, self()), {ok, _Pid} = - wpool:start_sup_pool(wpool_SUITE_overrun_pool, - [{workers, 1}, - {overrun_warning, 1000}, - {overrun_handler, {?MODULE, overrun_handler1}}]), + wpool:start_sup_pool( + wpool_SUITE_overrun_pool, + [ + {workers, 1}, + {overrun_warning, 1000}, + {overrun_handler, {?MODULE, overrun_handler1}} + ] + ), ok = wpool:cast(wpool_SUITE_overrun_pool, {timer, sleep, [1500]}), - _ = receive + _ = + receive {overrun1, Message} -> overrun = proplists:get_value(alert, Message), wpool_SUITE_overrun_pool = proplists:get_value(pool, Message), @@ -183,16 +214,22 @@ overrun(_Config) -> kill_on_overrun(_Config) -> true = register(overrun_handler, self()), {ok, _Pid} = - wpool:start_sup_pool(wpool_SUITE_kill_on_overrun_pool, - [{workers, 1}, - {overrun_warning, 500}, - {max_overrun_warnings, - 2}, %% The worker must be killed after 2 overrun - %% warnings, which is after 1 secs with this - %% configuration - {overrun_handler, {?MODULE, overrun_handler1}}]), + wpool:start_sup_pool( + wpool_SUITE_kill_on_overrun_pool, + [ + {workers, 1}, + {overrun_warning, 500}, + {max_overrun_warnings, + %% The worker must be killed after 2 overrun + 2}, + %% warnings, which is after 1 secs with this + %% configuration + {overrun_handler, {?MODULE, overrun_handler1}} + ] + ), ok = wpool:cast(wpool_SUITE_kill_on_overrun_pool, {timer, sleep, [2000]}), - _ = receive + _ = + receive {overrun1, Message} -> overrun = proplists:get_value(alert, Message), wpool_SUITE_kill_on_overrun_pool = proplists:get_value(pool, Message), @@ -203,7 +240,8 @@ kill_on_overrun(_Config) -> ct:fail(no_overrun) end, - _ = receive + _ = + receive {overrun1, Message2} -> max_overrun_limit = proplists:get_value(alert, Message2), wpool_SUITE_kill_on_overrun_pool = proplists:get_value(pool, Message2), @@ -232,8 +270,10 @@ stop_pool(_Config) -> -spec non_brutal_shutdown(config()) -> {comment, []}. non_brutal_shutdown(_Config) -> {ok, PoolPid} = - wpool:start_sup_pool(wpool_SUITE_non_brutal_shutdown, - [{workers, 1}, {pool_sup_shutdown, 100}]), + wpool:start_sup_pool( + wpool_SUITE_non_brutal_shutdown, + [{workers, 1}, {pool_sup_shutdown, 100}] + ), true = erlang:is_process_alive(PoolPid), Stats = wpool:stats(wpool_SUITE_non_brutal_shutdown), {workers, [{WorkerId, _}]} = lists:keyfind(workers, 1, Stats), @@ -252,10 +292,14 @@ non_brutal_shutdown(_Config) -> -spec brutal_worker_shutdown(config()) -> {comment, []}. brutal_worker_shutdown(_Config) -> {ok, PoolPid} = - wpool:start_sup_pool(wpool_SUITE_non_brutal_shutdown, - [{workers, 1}, - {pool_sup_shutdown, 100}, - {worker_shutdown, brutal_kill}]), + wpool:start_sup_pool( + wpool_SUITE_non_brutal_shutdown, + [ + {workers, 1}, + {pool_sup_shutdown, 100}, + {worker_shutdown, brutal_kill} + ] + ), true = erlang:is_process_alive(PoolPid), Stats = wpool:stats(wpool_SUITE_non_brutal_shutdown), {workers, [{WorkerId, _}]} = lists:keyfind(workers, 1, Stats), @@ -300,12 +344,14 @@ stats(_Config) -> 1 = Get(next_worker, InitStats), InitWorkers = Get(workers, InitStats), 10 = length(InitWorkers), - _ = [begin - WorkerStats = Get(I, InitWorkers), - 0 = Get(message_queue_len, WorkerStats), - [] = lists:keydelete(message_queue_len, 1, lists:keydelete(memory, 1, WorkerStats)) - end - || I <- lists:seq(1, 10)], + _ = [ + begin + WorkerStats = Get(I, InitWorkers), + 0 = Get(message_queue_len, WorkerStats), + [] = lists:keydelete(message_queue_len, 1, lists:keydelete(memory, 1, WorkerStats)) + end + || I <- lists:seq(1, 10) + ], % Start a long task on every worker Sleep = {timer, sleep, [2000]}, @@ -313,40 +359,44 @@ stats(_Config) -> ok = ktn_task:wait_for_success(fun() -> - WorkingStats = wpool:stats(wpool_SUITE_stats_pool), - wpool_SUITE_stats_pool = Get(pool, WorkingStats), - PoolPid = Get(supervisor, WorkingStats), - Options = Get(options, WorkingStats), - 10 = Get(size, WorkingStats), - 1 = Get(next_worker, WorkingStats), - WorkingWorkers = Get(workers, WorkingStats), - 10 = length(WorkingWorkers), - [begin - WorkerStats = Get(I, WorkingWorkers), - 0 = Get(message_queue_len, WorkerStats), - {timer, sleep, 1} = Get(current_function, WorkerStats), - {timer, sleep, 1, _} = Get(current_location, WorkerStats), - {cast, Sleep} = Get(task, WorkerStats), - true = is_number(Get(runtime, WorkerStats)) - end - || I <- lists:seq(1, 10)], - ok - end), + WorkingStats = wpool:stats(wpool_SUITE_stats_pool), + wpool_SUITE_stats_pool = Get(pool, WorkingStats), + PoolPid = Get(supervisor, WorkingStats), + Options = Get(options, WorkingStats), + 10 = Get(size, WorkingStats), + 1 = Get(next_worker, WorkingStats), + WorkingWorkers = Get(workers, WorkingStats), + 10 = length(WorkingWorkers), + [ + begin + WorkerStats = Get(I, WorkingWorkers), + 0 = Get(message_queue_len, WorkerStats), + {timer, sleep, 1} = Get(current_function, WorkerStats), + {timer, sleep, 1, _} = Get(current_location, WorkerStats), + {cast, Sleep} = Get(task, WorkerStats), + true = is_number(Get(runtime, WorkerStats)) + end + || I <- lists:seq(1, 10) + ], + ok + end), wpool:stop_sup_pool(wpool_SUITE_stats_pool), no_workers = - ktn_task:wait_for(fun() -> - try - wpool:stats(wpool_SUITE_stats_pool) - catch - _:E -> - E - end - end, - no_workers, - 100, - 50), + ktn_task:wait_for( + fun() -> + try + wpool:stats(wpool_SUITE_stats_pool) + catch + _:E -> + E + end + end, + no_workers, + 100, + 50 + ), {comment, []}. @@ -479,8 +529,10 @@ send_request(_Config) -> worker_killed_stats(_Config) -> %% Each server will take 100ms to start, but the start_sup_pool/2 call is synchronous anyway {ok, PoolPid} = - wpool:start_sup_pool(wpool_SUITE_worker_killed_stats, - [{workers, 3}, {worker, {sleepy_server, 500}}]), + wpool:start_sup_pool( + wpool_SUITE_worker_killed_stats, + [{workers, 3}, {worker, {sleepy_server, 500}}] + ), true = erlang:is_process_alive(PoolPid), Workers = @@ -504,14 +556,18 @@ worker_killed_stats(_Config) -> accepts_maps_and_lists_as_opts(_Config) -> %% Each server will take 100ms to start, but the start_sup_pool/2 call is synchronous anyway {ok, PoolPidList} = - wpool:start_sup_pool(accepts_maps_and_lists_as_opts_list, - [{workers, 3}, {worker, {sleepy_server, 500}}]), + wpool:start_sup_pool( + accepts_maps_and_lists_as_opts_list, + [{workers, 3}, {worker, {sleepy_server, 500}}] + ), true = erlang:is_process_alive(PoolPidList), ct:comment("accepts lists as opts"), {ok, PoolPidMap} = - wpool:start_sup_pool(accepts_maps_and_lists_as_opts_map, - #{workers => 3, worker => {sleepy_server, 500}}), + wpool:start_sup_pool( + accepts_maps_and_lists_as_opts_map, + #{workers => 3, worker => {sleepy_server, 500}} + ), true = erlang:is_process_alive(PoolPidMap), ct:comment("accepts lists as opts"), @@ -520,9 +576,11 @@ accepts_maps_and_lists_as_opts(_Config) -> -spec pool_of_supervisors(config()) -> {comment, string()}. pool_of_supervisors(_Config) -> Opts = - #{workers => 3, - worker_shutdown => infinity, - worker => {supervisor, {echo_supervisor, echo_supervisor, noargs}}}, + #{ + workers => 3, + worker_shutdown => infinity, + worker => {supervisor, {echo_supervisor, echo_supervisor, noargs}} + }, {ok, Pid} = wpool:start_sup_pool(pool_of_supervisors, Opts), true = erlang:is_process_alive(Pid), @@ -530,14 +588,16 @@ pool_of_supervisors(_Config) -> Run = fun(Sup, _) -> supervisor:start_child(Sup, [{ok, #{}}]) end, ForEach = fun(_) -> - {ok, EchoServer} = wpool:run(pool_of_supervisors, Run, next_worker), - true = erlang:is_process_alive(EchoServer) + {ok, EchoServer} = wpool:run(pool_of_supervisors, Run, next_worker), + true = erlang:is_process_alive(EchoServer) end, lists:foreach(ForEach, lists:seq(1, 9)), Supervisors = wpool:get_workers(pool_of_supervisors), - [3 = proplists:get_value(active, supervisor:count_children(Supervisor)) - || Supervisor <- Supervisors], + [ + 3 = proplists:get_value(active, supervisor:count_children(Supervisor)) + || Supervisor <- Supervisors + ], {comment, "Nicely load-balanced childrens across supervisors"}. diff --git a/test/wpool_bench.erl b/test/wpool_bench.erl index 4c53da6..4b19c37 100644 --- a/test/wpool_bench.erl +++ b/test/wpool_bench.erl @@ -3,10 +3,12 @@ -export([run_tasks/3]). %% @doc Returns the average time involved in processing the small tasks --spec run_tasks([{small | large, pos_integer()}, ...], - wpool:strategy(), - [wpool:option()]) -> - float(). +-spec run_tasks( + [{small | large, pos_integer()}, ...], + wpool:strategy(), + [wpool:option()] +) -> + float(). run_tasks(TaskGroups, Strategy, Options) -> Tasks = lists:flatten([lists:duplicate(N, Type) || {Type, N} <- TaskGroups]), {ok, _Pool} = wpool:start_sup_pool(?MODULE, Options), diff --git a/test/wpool_pool_SUITE.erl b/test/wpool_pool_SUITE.erl index 48e701b..09a68e5 100644 --- a/test/wpool_pool_SUITE.erl +++ b/test/wpool_pool_SUITE.erl @@ -25,9 +25,21 @@ -export([all/0]). -export([init_per_suite/1, end_per_suite/1, init_per_testcase/2, end_per_testcase/2]). --export([stop_worker/1, best_worker/1, next_worker/1, random_worker/1, available_worker/1, - hash_worker/1, custom_worker/1, next_available_worker/1, wpool_record/1, - queue_type_fifo/1, queue_type_lifo/1, get_workers/1, no_queue_manager/1]). +-export([ + stop_worker/1, + best_worker/1, + next_worker/1, + random_worker/1, + available_worker/1, + hash_worker/1, + custom_worker/1, + next_available_worker/1, + wpool_record/1, + queue_type_fifo/1, + queue_type_lifo/1, + get_workers/1, + no_queue_manager/1 +]). -export([manager_crash/1, super_fast/1, mess_up_with_store/1]). -elvis([{elvis_style, no_block_expressions, disable}]). @@ -35,9 +47,11 @@ -spec all() -> [atom()]. all() -> - [Fun + [ + Fun || {Fun, 1} <- module_info(exports), - not lists:member(Fun, [init_per_suite, end_per_suite, module_info])]. + not lists:member(Fun, [init_per_suite, end_per_suite, module_info]) + ]. -spec init_per_suite(config()) -> config(). init_per_suite(Config) -> @@ -104,20 +118,26 @@ available_worker(_Config) -> [0] = ktn_task:wait_for(fun() -> worker_msg_queue_lengths(Pool) end, [0]), - ct:log("Now send another round of messages, - the workers queues should still be empty"), + ct:log( + "Now send another round of messages,\n" + " the workers queues should still be empty" + ), [wpool:cast(Pool, {timer, sleep, [100 * I]}) || I <- lists:seq(1, ?WORKERS)], % Check that we have ?WORKERS pending tasks ?WORKERS = - ktn_task:wait_for(fun() -> - Stats1 = wpool:stats(Pool), - [0] = - lists:usort([proplists:get_value(message_queue_len, WS) - || {_, WS} <- proplists:get_value(workers, Stats1)]), - proplists:get_value(total_message_queue_len, Stats1) - end, - ?WORKERS), + ktn_task:wait_for( + fun() -> + Stats1 = wpool:stats(Pool), + [0] = + lists:usort([ + proplists:get_value(message_queue_len, WS) + || {_, WS} <- proplists:get_value(workers, Stats1) + ]), + proplists:get_value(total_message_queue_len, Stats1) + end, + ?WORKERS + ), ct:log("If we can't wait we get no workers"), try wpool:call(Pool, {erlang, self, []}, available_worker, 100) of @@ -149,18 +169,24 @@ available_worker(_Config) -> % Check we have no pending tasks 0 = - ktn_task:wait_for(fun() -> proplists:get_value(total_message_queue_len, wpool:stats(Pool)) - end, - 0), - - ct:log("We run tons of calls, and none is blocked, - because all of them are handled by different workers"), + ktn_task:wait_for( + fun() -> proplists:get_value(total_message_queue_len, wpool:stats(Pool)) end, + 0 + ), + + ct:log( + "We run tons of calls, and none is blocked,\n" + " because all of them are handled by different workers" + ), Workers = - [wpool:call(Pool, {erlang, self, []}, available_worker, 5000) - || _ <- lists:seq(1, 20 * ?WORKERS)], + [ + wpool:call(Pool, {erlang, self, []}, available_worker, 5000) + || _ <- lists:seq(1, 20 * ?WORKERS) + ], UniqueWorkers = sets:to_list( - sets:from_list(Workers)), + sets:from_list(Workers) + ), {?WORKERS, UniqueWorkers, true} = {?WORKERS, UniqueWorkers, ?WORKERS / 2 >= length(UniqueWorkers)}, @@ -214,14 +240,18 @@ next_available_worker(_Config) -> {ok, _} = wpool:run(Pool, Run, next_available_worker), ct:log("Put them all to work..."), - [wpool:cast(Pool, {timer, sleep, [1500 + I]}, next_available_worker) - || I <- lists:seq(0, (?WORKERS - 1) * 60000, 60000)], + [ + wpool:cast(Pool, {timer, sleep, [1500 + I]}, next_available_worker) + || I <- lists:seq(0, (?WORKERS - 1) * 60000, 60000) + ], AvailableWorkers = fun() -> - length([a_worker - || {_, WS} <- proplists:get_value(workers, wpool:stats(Pool)), - proplists:get_value(task, WS) == undefined]) + length([ + a_worker + || {_, WS} <- proplists:get_value(workers, wpool:stats(Pool)), + proplists:get_value(task, WS) == undefined + ]) end, ct:log("All busy..."), @@ -268,23 +298,28 @@ next_worker(_Config) -> end, Res0 = - [begin - Stats = wpool:stats(Pool), - I = proplists:get_value(next_worker, Stats), - wpool:call(Pool, {erlang, self, []}, next_worker, infinity) - end - || I <- lists:seq(1, ?WORKERS)], + [ + begin + Stats = wpool:stats(Pool), + I = proplists:get_value(next_worker, Stats), + wpool:call(Pool, {erlang, self, []}, next_worker, infinity) + end + || I <- lists:seq(1, ?WORKERS) + ], ?WORKERS = sets:size( - sets:from_list(Res0)), + sets:from_list(Res0) + ), Res0 = - [begin - Stats = wpool:stats(Pool), - I = proplists:get_value(next_worker, Stats), - wpool:call(Pool, {erlang, self, []}, next_worker) - end - || I <- lists:seq(1, ?WORKERS)], + [ + begin + Stats = wpool:stats(Pool), + I = proplists:get_value(next_worker, Stats), + wpool:call(Pool, {erlang, self, []}, next_worker) + end + || I <- lists:seq(1, ?WORKERS) + ], Req = wpool:send_request(Pool, {erlang, self, []}, next_worker), {reply, {ok, _}} = gen_server:wait_response(Req, 5000), @@ -315,19 +350,24 @@ random_worker(_Config) -> [wpool:call(Pool, {erlang, self, []}, random_worker) || _ <- lists:seq(1, 20 * ?WORKERS)], ?WORKERS = sets:size( - sets:from_list(Serial)), + sets:from_list(Serial) + ), %% Randomly ask a lot of workers to send ourselves the atom true - [wpool:cast(Pool, {erlang, send, [self(), true]}, random_worker) - || _ <- lists:seq(1, 20 * ?WORKERS)], + [ + wpool:cast(Pool, {erlang, send, [self(), true]}, random_worker) + || _ <- lists:seq(1, 20 * ?WORKERS) + ], Results = - [receive - true -> - true - after 5000 -> - ct:fail("Didn't receive 'true' in time") - end - || _ <- lists:seq(1, 20 * ?WORKERS)], + [ + receive + true -> + true + after 5000 -> + ct:fail("Didn't receive 'true' in time") + end + || _ <- lists:seq(1, 20 * ?WORKERS) + ], true = lists:all(fun(Value) -> Value end, Results), %% do a gen_server:send_request/3 @@ -337,15 +377,18 @@ random_worker(_Config) -> %% Now do the same with a freshly spawned process for each request to ensure %% randomness isn't reset with each spawn of the process_dictionary Self = self(), - _ = [spawn(fun() -> - WorkerId = wpool:call(Pool, {erlang, self, []}, random_worker), - Self ! {worker, WorkerId} - end) - || _ <- lists:seq(1, 20 * ?WORKERS)], + _ = [ + spawn(fun() -> + WorkerId = wpool:call(Pool, {erlang, self, []}, random_worker), + Self ! {worker, WorkerId} + end) + || _ <- lists:seq(1, 20 * ?WORKERS) + ], Concurrent = collect_results(20 * ?WORKERS, []), ?WORKERS = sets:size( - sets:from_list(Concurrent)), + sets:from_list(Concurrent) + ), {comment, []}. @@ -364,26 +407,34 @@ hash_worker(_Config) -> %% Use two hash keys that have different values (0, 1) to target only %% two workers. Other workers should be missing. Targeted = - [wpool:call(Pool, {erlang, self, []}, {hash_worker, I rem 2}) - || I <- lists:seq(1, 20 * ?WORKERS)], + [ + wpool:call(Pool, {erlang, self, []}, {hash_worker, I rem 2}) + || I <- lists:seq(1, 20 * ?WORKERS) + ], 2 = sets:size( - sets:from_list(Targeted)), + sets:from_list(Targeted) + ), %% Now use many different hash keys. All workers should be hit. Spread = - [wpool:call(Pool, {erlang, self, []}, {hash_worker, I}) - || I <- lists:seq(1, 20 * ?WORKERS)], + [ + wpool:call(Pool, {erlang, self, []}, {hash_worker, I}) + || I <- lists:seq(1, 20 * ?WORKERS) + ], ?WORKERS = sets:size( - sets:from_list(Spread)), + sets:from_list(Spread) + ), Run = fun(Worker, Timeout) -> gen_server:call(Worker, {erlang, self, []}, Timeout) end, [{ok, _} = wpool:run(Pool, Run, {hash_worker, I}) || I <- lists:seq(1, 20 * ?WORKERS)], %% Fill up their message queues... - [wpool:cast(Pool, {timer, sleep, [60000]}, {hash_worker, I}) - || I <- lists:seq(1, 20 * ?WORKERS)], + [ + wpool:cast(Pool, {timer, sleep, [60000]}, {hash_worker, I}) + || I <- lists:seq(1, 20 * ?WORKERS) + ], false = ktn_task:wait_for(fun() -> lists:member(0, worker_msg_queue_lengths(Pool)) end, false), @@ -404,30 +455,37 @@ custom_worker(_Config) -> ok end, - _ = [begin - Stats = wpool:stats(Pool), - I = proplists:get_value(next_worker, Stats), - wpool:cast(Pool, {io, format, ["ok!"]}, Strategy) - end - || I <- lists:seq(1, ?WORKERS)], + _ = [ + begin + Stats = wpool:stats(Pool), + I = proplists:get_value(next_worker, Stats), + wpool:cast(Pool, {io, format, ["ok!"]}, Strategy) + end + || I <- lists:seq(1, ?WORKERS) + ], Res0 = - [begin - Stats = wpool:stats(Pool), - I = proplists:get_value(next_worker, Stats), - wpool:call(Pool, {erlang, self, []}, Strategy, infinity) - end - || I <- lists:seq(1, ?WORKERS)], + [ + begin + Stats = wpool:stats(Pool), + I = proplists:get_value(next_worker, Stats), + wpool:call(Pool, {erlang, self, []}, Strategy, infinity) + end + || I <- lists:seq(1, ?WORKERS) + ], ?WORKERS = sets:size( - sets:from_list(Res0)), + sets:from_list(Res0) + ), Res0 = - [begin - Stats = wpool:stats(Pool), - I = proplists:get_value(next_worker, Stats), - wpool:call(Pool, {erlang, self, []}, Strategy) - end - || I <- lists:seq(1, ?WORKERS)], + [ + begin + Stats = wpool:stats(Pool), + I = proplists:get_value(next_worker, Stats), + wpool:call(Pool, {erlang, self, []}, Strategy) + end + || I <- lists:seq(1, ?WORKERS) + ], Req = wpool:send_request(Pool, {erlang, self, []}, Strategy), {reply, {ok, _}} = gen_server:wait_response(Req, 5000), @@ -451,8 +509,10 @@ manager_crash(_Config) -> exit(whereis(QueueManager), kill), false = - ktn_task:wait_for(fun() -> lists:member(whereis(QueueManager), [OldPid, undefined]) end, - false), + ktn_task:wait_for( + fun() -> lists:member(whereis(QueueManager), [OldPid, undefined]) end, + false + ), ct:log("Check that the pool is working again"), {ok, ok} = send_io_format(Pool), @@ -613,16 +673,18 @@ mess_up_with_store(_Config) -> Flag = process_flag(trap_exit, true), exit(whereis(Pool), kill), ok = - ktn_task:wait_for(fun() -> - try wpool:call(Pool, {io, format, ["1!~n"]}, random_worker) of - X -> - {unexpected, X} - catch - _:no_workers -> - ok - end - end, - ok), + ktn_task:wait_for( + fun() -> + try wpool:call(Pool, {io, format, ["1!~n"]}, random_worker) of + X -> + {unexpected, X} + catch + _:no_workers -> + ok + end + end, + ok + ), true = process_flag(trap_exit, Flag), @@ -636,19 +698,23 @@ mess_up_with_store(_Config) -> {comment, []}. cast_tasks(Pool, TasksNumber, ReplyTo) -> - lists:foreach(fun(N) -> wpool:cast(Pool, {erlang, send, [ReplyTo, {task, N}]}) end, - lists:seq(1, TasksNumber)). + lists:foreach( + fun(N) -> wpool:cast(Pool, {erlang, send, [ReplyTo, {task, N}]}) end, + lists:seq(1, TasksNumber) + ). collect_tasks(TasksNumber) -> - lists:map(fun(_) -> - receive - {task, N} -> - N - after 5000 -> - ct:fail("Didn't receive {'task', N} in time") - end - end, - lists:seq(1, TasksNumber)). + lists:map( + fun(_) -> + receive + {task, N} -> + N + after 5000 -> + ct:fail("Didn't receive {'task', N} in time") + end + end, + lists:seq(1, TasksNumber) + ). collect_results(0, Results) -> Results; @@ -664,8 +730,10 @@ send_io_format(Pool) -> {ok, ok} = wpool:call(Pool, {io, format, ["ok!~n"]}, available_worker). worker_msg_queue_lengths(Pool) -> - lists:usort([proplists:get_value(message_queue_len, WS) - || {_, WS} <- proplists:get_value(workers, wpool:stats(Pool))]). + lists:usort([ + proplists:get_value(message_queue_len, WS) + || {_, WS} <- proplists:get_value(workers, wpool:stats(Pool)) + ]). store_mess_up(Pool) -> true = persistent_term:erase({wpool_pool, Pool}). diff --git a/test/wpool_process_SUITE.erl b/test/wpool_process_SUITE.erl index 0f9d1f0..d4fe100 100644 --- a/test/wpool_process_SUITE.erl +++ b/test/wpool_process_SUITE.erl @@ -23,15 +23,29 @@ -export([all/0]). -export([init_per_suite/1, end_per_suite/1, init_per_testcase/2, end_per_testcase/2]). --export([init/1, init_timeout/1, info/1, cast/1, send_request/1, call/1, continue/1, - handle_info_missing/1, handle_info_fails/1, format_status/1, no_format_status/1, stop/1]). +-export([ + init/1, + init_timeout/1, + info/1, + cast/1, + send_request/1, + call/1, + continue/1, + handle_info_missing/1, + handle_info_fails/1, + format_status/1, + no_format_status/1, + stop/1 +]). -export([pool_restart_crash/1, pool_norestart_crash/1, complete_coverage/1]). -spec all() -> [atom()]. all() -> - [Fun + [ + Fun || {Fun, 1} <- module_info(exports), - not lists:member(Fun, [init_per_suite, end_per_suite, module_info])]. + not lists:member(Fun, [init_per_suite, end_per_suite, module_info]) + ]. -spec init_per_suite(config()) -> config(). init_per_suite(Config) -> @@ -51,7 +65,8 @@ init_per_testcase(_TestCase, Config) -> -spec end_per_testcase(atom(), config()) -> config(). end_per_testcase(_TestCase, Config) -> process_flag(trap_exit, false), - receive after 0 -> + receive + after 0 -> ok end, Config. @@ -124,10 +139,12 @@ continue(_Config) -> C = fun(ContinueState) -> {noreply, ContinueState} end, %% init/1 returns {continue, continue_state} {ok, Pid} = - wpool_process:start_link(?MODULE, - echo_server, - {ok, state, {continue, C(continue_state)}}, - #{}), + wpool_process:start_link( + ?MODULE, + echo_server, + {ok, state, {continue, C(continue_state)}}, + #{} + ), continue_state = get_state(Pid), %% handle_call/3 returns {continue, ...} @@ -243,11 +260,13 @@ pool_restart_crash(_Config) -> pool_norestart_crash(_Config) -> Pool = pool_norestart_crash, PoolOptions = - [{workers, 2}, - {worker, {crashy_server, []}}, - {strategy, {one_for_all, 0, 10}}, - {pool_sup_intensity, 0}, - {pool_sup_period, 10}], + [ + {workers, 2}, + {worker, {crashy_server, []}}, + {strategy, {one_for_all, 0, 10}}, + {pool_sup_intensity, 0}, + {pool_sup_period, 10} + ], {ok, Pid} = wpool:start_pool(Pool, PoolOptions), ct:log("Check that the pool is working"), @@ -349,10 +368,13 @@ get_state(Pid) -> {status, Pid, {module, gen_server}, [_PDict, _SysState, _Parent, _Dbg, Misc]} = sys:get_status(Pid), [State] = - lists:filtermap(fun ({data, [{"State", State}]}) -> - {true, State}; - (_) -> - false - end, - Misc), + lists:filtermap( + fun + ({data, [{"State", State}]}) -> + {true, State}; + (_) -> + false + end, + Misc + ), wpool_process:get_state(State). diff --git a/test/wpool_process_callbacks_SUITE.erl b/test/wpool_process_callbacks_SUITE.erl index 01f1eac..fc4752d 100644 --- a/test/wpool_process_callbacks_SUITE.erl +++ b/test/wpool_process_callbacks_SUITE.erl @@ -8,21 +8,26 @@ -export([all/0]). -export([init_per_suite/1, end_per_suite/1]). --export([complete_callback_passed_when_starting_pool/1, - partial_callback_passed_when_starting_pool/1, - callback_can_be_added_and_removed_after_pool_is_started/1, - crashing_callback_does_not_affect_others/1, non_existsing_module_does_not_affect_others/1, - complete_coverage/1]). +-export([ + complete_callback_passed_when_starting_pool/1, + partial_callback_passed_when_starting_pool/1, + callback_can_be_added_and_removed_after_pool_is_started/1, + crashing_callback_does_not_affect_others/1, + non_existsing_module_does_not_affect_others/1, + complete_coverage/1 +]). -dialyzer({no_underspecs, all/0}). -spec all() -> [atom()]. all() -> - [complete_callback_passed_when_starting_pool, - partial_callback_passed_when_starting_pool, - callback_can_be_added_and_removed_after_pool_is_started, - crashing_callback_does_not_affect_others, - non_existsing_module_does_not_affect_others]. + [ + complete_callback_passed_when_starting_pool, + partial_callback_passed_when_starting_pool, + callback_can_be_added_and_removed_after_pool_is_started, + crashing_callback_does_not_affect_others, + non_existsing_module_does_not_affect_others + ]. -spec init_per_suite(config()) -> config(). init_per_suite(Config) -> @@ -43,11 +48,15 @@ complete_callback_passed_when_starting_pool(_Config) -> meck:expect(callbacks, handle_worker_creation, fun(_AWorkerName) -> ok end), meck:expect(callbacks, handle_worker_death, fun(_AWName, _Reason) -> ok end), {ok, _Pid} = - wpool:start_pool(Pool, - [{workers, WorkersCount}, - {enable_callbacks, true}, - {worker, {crashy_server, []}}, - {callbacks, [callbacks]}]), + wpool:start_pool( + Pool, + [ + {workers, WorkersCount}, + {enable_callbacks, true}, + {worker, {crashy_server, []}}, + {callbacks, [callbacks]} + ] + ), WorkersCount = ktn_task:wait_for(function_calls(callbacks, handle_init_start, ['_']), WorkersCount), @@ -69,10 +78,14 @@ partial_callback_passed_when_starting_pool(_Config) -> meck:expect(callbacks, handle_worker_creation, fun(_AWorkerName) -> ok end), meck:expect(callbacks, handle_worker_death, fun(_AWName, _Reason) -> ok end), {ok, _Pid} = - wpool:start_pool(Pool, - [{workers, WorkersCount}, - {enable_callbacks, true}, - {callbacks, [callbacks]}]), + wpool:start_pool( + Pool, + [ + {workers, WorkersCount}, + {enable_callbacks, true}, + {callbacks, [callbacks]} + ] + ), WorkersCount = ktn_task:wait_for(function_calls(callbacks, handle_worker_creation, ['_']), WorkersCount), wpool:stop_pool(Pool), @@ -89,10 +102,14 @@ callback_can_be_added_and_removed_after_pool_is_started(_Config) -> meck:new(callbacks2, [non_strict]), meck:expect(callbacks2, handle_worker_death, fun(_AWName, _Reason) -> ok end), {ok, _Pid} = - wpool:start_pool(Pool, - [{workers, WorkersCount}, - {worker, {crashy_server, []}}, - {enable_callbacks, true}]), + wpool:start_pool( + Pool, + [ + {workers, WorkersCount}, + {worker, {crashy_server, []}}, + {enable_callbacks, true} + ] + ), %% Now we are adding 2 callback modules _ = wpool_pool:add_callback_module(Pool, callbacks), _ = wpool_pool:add_callback_module(Pool, callbacks2), @@ -125,21 +142,29 @@ crashing_callback_does_not_affect_others(_Config) -> meck:new(callbacks, [non_strict]), meck:expect(callbacks, handle_worker_creation, fun(_AWorkerName) -> ok end), meck:new(callbacks2, [non_strict]), - meck:expect(callbacks2, - handle_worker_creation, - fun(AWorkerName) -> {not_going_to_work} = AWorkerName end), + meck:expect( + callbacks2, + handle_worker_creation, + fun(AWorkerName) -> {not_going_to_work} = AWorkerName end + ), {ok, _Pid} = - wpool:start_pool(Pool, - [{workers, WorkersCount}, - {worker, {crashy_server, []}}, - {enable_callbacks, true}, - {callbacks, [callbacks, callbacks2]}]), + wpool:start_pool( + Pool, + [ + {workers, WorkersCount}, + {worker, {crashy_server, []}}, + {enable_callbacks, true}, + {callbacks, [callbacks, callbacks2]} + ] + ), WorkersCount = ktn_task:wait_for(function_calls(callbacks, handle_worker_creation, ['_']), WorkersCount), WorkersCount = - ktn_task:wait_for(function_calls(callbacks2, handle_worker_creation, ['_']), - WorkersCount), + ktn_task:wait_for( + function_calls(callbacks2, handle_worker_creation, ['_']), + WorkersCount + ), wpool:stop_pool(Pool), meck:unload(callbacks), @@ -154,11 +179,15 @@ non_existsing_module_does_not_affect_others(_Config) -> meck:new(callbacks, [non_strict]), meck:expect(callbacks, handle_worker_creation, fun(_AWorkerName) -> ok end), {ok, _Pid} = - wpool:start_pool(Pool, - [{workers, WorkersCount}, - {worker, {crashy_server, []}}, - {enable_callbacks, true}, - {callbacks, [callbacks, non_existing_m]}]), + wpool:start_pool( + Pool, + [ + {workers, WorkersCount}, + {worker, {crashy_server, []}}, + {enable_callbacks, true}, + {callbacks, [callbacks, non_existing_m]} + ] + ), {error, nofile} = wpool_pool:add_callback_module(Pool, non_existing_m2), diff --git a/test/wpool_worker_SUITE.erl b/test/wpool_worker_SUITE.erl index 20a1d16..39ef0c6 100644 --- a/test/wpool_worker_SUITE.erl +++ b/test/wpool_worker_SUITE.erl @@ -28,9 +28,11 @@ -spec all() -> [atom()]. all() -> - [Fun + [ + Fun || {Fun, 1} <- module_info(exports), - not lists:member(Fun, [init_per_suite, end_per_suite, module_info])]. + not lists:member(Fun, [init_per_suite, end_per_suite, module_info]) + ]. -spec init_per_suite(config()) -> config(). init_per_suite(Config) -> From 5ae823391bcfc2b4bacaf145df9d78120bdfdc17 Mon Sep 17 00:00:00 2001 From: Nelson Vides Date: Sat, 21 Feb 2026 10:26:43 +0100 Subject: [PATCH 2/3] Fix latest linter --- elvis.config | 11 +++-------- src/wpool.erl | 3 +-- src/wpool_process.erl | 8 ++------ src/wpool_queue_manager.erl | 8 ++------ src/wpool_time_checker.erl | 7 ++----- src/wpool_worker.erl | 6 +----- test/crashy_server.erl | 6 +----- test/echo_server.erl | 6 +----- test/sleepy_server.erl | 6 +----- 9 files changed, 14 insertions(+), 47 deletions(-) diff --git a/elvis.config b/elvis.config index 1376f78..51831c8 100644 --- a/elvis.config +++ b/elvis.config @@ -2,15 +2,15 @@ {elvis, [ {config, [ #{ - dirs => ["src"], + dirs => ["src/**"], filter => "*.erl", ruleset => erl_files, rules => [ - {elvis_style, invalid_dynamic_call, #{ + {elvis_style, no_invalid_dynamic_calls, #{ ignore => [wpool_process, wpool_time_checker] }}, - {elvis_style, god_modules, #{limit => 30}}, + {elvis_style, no_god_modules, #{limit => 30}}, {elvis_style, state_record_and_type, disable}, {elvis_style, dont_repeat_yourself, #{ ignore => [wpool_SUITE], min_complexity => 13 @@ -27,11 +27,6 @@ {elvis_style, state_record_and_type, disable}, {elvis_style, dont_repeat_yourself, #{min_complexity => 13}} ] - }, - #{ - dirs => ["."], - filter => "elvis.config", - ruleset => elvis_config } ]} ]} diff --git a/src/wpool.erl b/src/wpool.erl index 5c2a9f6..289f7b9 100644 --- a/src/wpool.erl +++ b/src/wpool.erl @@ -55,8 +55,7 @@ %%% `wpool:stats/1'. -module(wpool). -%% @todo remove this line when https://github.com/AdRoll/rebar3_format/issues/356 is fixed --format(ignore). +-elvis([{elvis_style, private_data_types, disable}]). -behaviour(application). diff --git a/src/wpool_process.erl b/src/wpool_process.erl index b9df7f4..297b648 100644 --- a/src/wpool_process.erl +++ b/src/wpool_process.erl @@ -23,7 +23,7 @@ module :: module(), handle_call :: fun( - (Request :: term(), From :: from(), State :: term()) -> + (Request :: term(), From :: gen_server:from(), State :: term()) -> {reply, Reply :: term(), NewState :: term()} | {reply, Reply :: term(), NewState :: term(), timeout() | hibernate | {continue, term()}} @@ -64,10 +64,6 @@ -export_type([state/0]). --type from() :: {pid(), reference()}. - --export_type([from/0]). - -type next_step() :: timeout() | hibernate | {continue, term()}. -export_type([next_step/0]). @@ -299,7 +295,7 @@ handle_cast(Cast, #state{mod = CbCache, options = Options} = State) -> Reply. %% @private --spec handle_call(term(), from(), state()) -> +-spec handle_call(term(), gen_server:from(), state()) -> {reply, term(), state()} | {reply, term(), state(), next_step()} | {noreply, state()} diff --git a/src/wpool_queue_manager.erl b/src/wpool_queue_manager.erl index 6d5bd40..a7526ac 100644 --- a/src/wpool_queue_manager.erl +++ b/src/wpool_queue_manager.erl @@ -44,11 +44,7 @@ -export_type([state/0]). --type from() :: {pid(), gen_server:reply_tag()}. - --export_type([from/0]). - --type monitored_from() :: {reference(), from()}. +-type monitored_from() :: {reference(), gen_server:from()}. -type options() :: [{option(), term()}]. -export_type([options/0]). @@ -228,7 +224,7 @@ handle_cast({cast_to_available_worker, Cast}, State) -> {noreply, State#state{workers = NewWorkers}} end. --spec handle_call(call_request(), from(), state()) -> +-spec handle_call(call_request(), gen_server:from(), state()) -> {reply, {ok, atom()}, state()} | {noreply, state()}. handle_call({available_worker, ExpiresAt}, {ClientPid, _Ref} = Client, State) -> #state{workers = Workers, clients = Clients} = State, diff --git a/src/wpool_time_checker.erl b/src/wpool_time_checker.erl index bc775a6..50b1b0a 100644 --- a/src/wpool_time_checker.erl +++ b/src/wpool_time_checker.erl @@ -26,16 +26,13 @@ -export_type([state/0]). --type from() :: {pid(), reference()}. - --export_type([from/0]). - %% api -export([start_link/3, add_handler/2]). %% gen_server callbacks -export([init/1, handle_call/3, handle_cast/2, handle_info/2]). -elvis([{elvis_style, no_catch_expressions, disable}]). +-elvis([{elvis_style, private_data_types, disable}]). %%%=================================================================== %%% API @@ -66,7 +63,7 @@ init({WPool, Handlers}) -> handle_cast(_Cast, State) -> {noreply, State}. --spec handle_call({add_handler, handler()}, from(), state()) -> {reply, ok, state()}. +-spec handle_call({add_handler, handler()}, gen_server:from(), state()) -> {reply, ok, state()}. handle_call({add_handler, Handler}, _, #state{handlers = Handlers} = State) -> {reply, ok, State#state{handlers = [Handler | Handlers]}}. diff --git a/src/wpool_worker.erl b/src/wpool_worker.erl index 32ffafe..02b16df 100644 --- a/src/wpool_worker.erl +++ b/src/wpool_worker.erl @@ -31,10 +31,6 @@ -export_type([state/0]). --type from() :: {pid(), reference()}. - --export_type([from/0]). - %%%=================================================================== %%% API %%%=================================================================== @@ -88,7 +84,7 @@ handle_cast(Cast, State) -> {noreply, State, hibernate}. %% @private --spec handle_call(term(), from(), state()) -> +-spec handle_call(term(), gen_server:from(), state()) -> {reply, {ok, term()} | {error, term()}, state(), hibernate}. handle_call({M, F, A}, _From, State) -> try erlang:apply(M, F, A) of diff --git a/test/crashy_server.erl b/test/crashy_server.erl index ba1b6c5..e0bf9a3 100644 --- a/test/crashy_server.erl +++ b/test/crashy_server.erl @@ -28,10 +28,6 @@ -dialyzer([no_behaviours]). --type from() :: {pid(), reference()}. - --export_type([from/0]). - %%%=================================================================== %%% callbacks %%%=================================================================== @@ -61,7 +57,7 @@ handle_cast(crash, _State) -> handle_cast(Cast, _State) -> Cast. --spec handle_call(state | Call, from(), State) -> {reply, State, State} | Call. +-spec handle_call(state | Call, gen_server:from(), State) -> {reply, State, State} | Call. handle_call(state, _From, State) -> {reply, State, State}; handle_call(crash, _From, _State) -> diff --git a/test/echo_server.erl b/test/echo_server.erl index d1484c6..2965374 100644 --- a/test/echo_server.erl +++ b/test/echo_server.erl @@ -31,10 +31,6 @@ -dialyzer([no_behaviours]). --type from() :: {pid(), reference()}. - --export_type([from/0]). - -spec start_link(term()) -> gen_server:start_ret(). start_link(Something) -> gen_server:start_link(?MODULE, Something, []). @@ -64,7 +60,7 @@ handle_info(Info, _State) -> handle_cast(Cast, _State) -> Cast. --spec handle_call(Call, from(), term()) -> Call. +-spec handle_call(Call, gen_server:from(), term()) -> Call. handle_call(Call, _From, _State) -> Call. diff --git a/test/sleepy_server.erl b/test/sleepy_server.erl index 8a5993e..2602710 100644 --- a/test/sleepy_server.erl +++ b/test/sleepy_server.erl @@ -21,10 +21,6 @@ -dialyzer([no_behaviours]). --type from() :: {pid(), reference()}. - --export_type([from/0]). - %%%=================================================================== %%% callbacks %%%=================================================================== @@ -40,7 +36,7 @@ handle_cast(TimeToSleep, State) -> _ = timer:sleep(TimeToSleep), {noreply, State}. --spec handle_call(pos_integer(), from(), State) -> {reply, ok, State}. +-spec handle_call(pos_integer(), gen_server:from(), State) -> {reply, ok, State}. handle_call(TimeToSleep, _From, State) -> _ = timer:sleep(TimeToSleep), {reply, ok, State}. From 4bd187f48bdfb3a6d7710d55abf88da76527364c Mon Sep 17 00:00:00 2001 From: Nelson Vides Date: Sat, 21 Feb 2026 10:38:41 +0100 Subject: [PATCH 3/3] Ignore reformatting commit revision --- .git-blame-ignore-revs | 1 + 1 file changed, 1 insertion(+) create mode 100644 .git-blame-ignore-revs diff --git a/.git-blame-ignore-revs b/.git-blame-ignore-revs new file mode 100644 index 0000000..f26c045 --- /dev/null +++ b/.git-blame-ignore-revs @@ -0,0 +1 @@ +64a2a22617a942230873949cb0a11ac582741438