Current section
Files
Jump to
Current section
Files
src/gabsurd@worker.erl
-module(gabsurd@worker).
-compile([no_auto_import, nowarn_unused_vars, nowarn_unused_function, nowarn_nomatch, inline]).
-define(FILEPATH, "src/gabsurd/worker.gleam").
-export([new/3, with_worker_id/2, with_poll_interval/2, with_claim_timeout/2, with_batch_size/2, with_max_backoff/2, start/1, stop/1, child_spec/2, pool_child_specs/3]).
-export_type([handler_result/0, handler/0, config/0, worker/0, message/0, worker_state/0]).
-if(?OTP_RELEASE >= 27).
-define(MODULEDOC(Str), -moduledoc(Str)).
-define(DOC(Str), -doc(Str)).
-else.
-define(MODULEDOC(Str), -compile([])).
-define(DOC(Str), -compile([])).
-endif.
?MODULEDOC(
" OTP Worker actor for the Absurd durable workflow system.\n"
" Polls a queue for tasks, dispatches to registered handlers,\n"
" and completes/fails tasks based on handler results.\n"
"\n"
" ## Distributed systems behaviour\n"
"\n"
" - **Claim extension**: The primary lease extension mechanism is\n"
" `checkpoint.set`, which passes `claim_timeout` as `extend_claim_by`\n"
" to `set_task_checkpoint_state` — every checkpoint write extends the\n"
" lease. For handlers doing a single long operation without checkpoints,\n"
" the claim timeout is the safety net: if the handler takes too long,\n"
" the claim expires and another worker picks up the task.\n"
" - **Error backoff**: On transient claim errors, backs off exponentially\n"
" up to `max_backoff` (default 60s), resets on success.\n"
" - **Unknown task deferral**: Tasks with no registered handler are deferred\n"
" (rescheduled with a delay) rather than failed. This supports rolling\n"
" deployments where a new task type may arrive before its handler is\n"
" deployed.\n"
" - **Terminal state tolerance**: complete/fail errors from already-\n"
" completed or already-failed runs are silently ignored (matching the\n"
" official Absurd SDK behaviour).\n"
).
-type handler_result() :: {complete, gleam@json:json()} |
{fail, gleam@json:json()} |
suspend.
-type handler() :: {handler,
binary(),
fun((gabsurd@context:context()) -> handler_result()),
gleam@option:option(fun((gabsurd@context:context(), gleam@json:json()) -> nil))}.
-type config() :: {config,
gabsurd@client:db(),
binary(),
binary(),
integer(),
integer(),
integer(),
integer(),
list(handler())}.
-type worker() :: {worker, gleam@erlang@process:subject(message())}.
-type message() :: poll | shutdown.
-type worker_state() :: {worker_state,
config(),
gleam@dict:dict(binary(), handler()),
gleam@erlang@process:subject(message()),
integer()}.
-file("src/gabsurd/worker.gleam", 132).
?DOC(" Create a new worker config with defaults.\n").
-spec new(gabsurd@client:db(), binary(), list(handler())) -> config().
new(Db, Queue_name, Handlers) ->
{config,
Db,
Queue_name,
<<"gabsurd_worker"/utf8>>,
5000,
30,
1,
60000,
Handlers}.
-file("src/gabsurd/worker.gleam", 146).
?DOC(" Set the worker ID (used as `worker_id` in claim_task).\n").
-spec with_worker_id(config(), binary()) -> config().
with_worker_id(Config, Worker_id) ->
{config,
erlang:element(2, Config),
erlang:element(3, Config),
Worker_id,
erlang:element(5, Config),
erlang:element(6, Config),
erlang:element(7, Config),
erlang:element(8, Config),
erlang:element(9, Config)}.
-file("src/gabsurd/worker.gleam", 151).
?DOC(" Set the poll interval in milliseconds.\n").
-spec with_poll_interval(config(), integer()) -> config().
with_poll_interval(Config, Interval_ms) ->
{config,
erlang:element(2, Config),
erlang:element(3, Config),
erlang:element(4, Config),
Interval_ms,
erlang:element(6, Config),
erlang:element(7, Config),
erlang:element(8, Config),
erlang:element(9, Config)}.
-file("src/gabsurd/worker.gleam", 156).
?DOC(" Set the claim timeout in seconds.\n").
-spec with_claim_timeout(config(), integer()) -> config().
with_claim_timeout(Config, Timeout_secs) ->
{config,
erlang:element(2, Config),
erlang:element(3, Config),
erlang:element(4, Config),
erlang:element(5, Config),
Timeout_secs,
erlang:element(7, Config),
erlang:element(8, Config),
erlang:element(9, Config)}.
-file("src/gabsurd/worker.gleam", 161).
?DOC(" Set the batch size (tasks claimed per poll).\n").
-spec with_batch_size(config(), integer()) -> config().
with_batch_size(Config, Size) ->
{config,
erlang:element(2, Config),
erlang:element(3, Config),
erlang:element(4, Config),
erlang:element(5, Config),
erlang:element(6, Config),
Size,
erlang:element(8, Config),
erlang:element(9, Config)}.
-file("src/gabsurd/worker.gleam", 167).
?DOC(
" Set the maximum backoff in milliseconds for retrying after claim errors.\n"
" Default: 60000 (60 seconds).\n"
).
-spec with_max_backoff(config(), integer()) -> config().
with_max_backoff(Config, Max_backoff_ms) ->
{config,
erlang:element(2, Config),
erlang:element(3, Config),
erlang:element(4, Config),
erlang:element(5, Config),
erlang:element(6, Config),
erlang:element(7, Config),
Max_backoff_ms,
erlang:element(9, Config)}.
-file("src/gabsurd/worker.gleam", 413).
-spec log_error(binary(), gabsurd@client:gabsurd_error()) -> nil.
log_error(Context, Error) ->
Msg = case Error of
{query_error, Reason} ->
<<<<Context/binary, ": query error: "/utf8>>/binary, Reason/binary>>;
{unexpected_row_count, Reason@1} ->
<<<<Context/binary, ": "/utf8>>/binary, Reason@1/binary>>;
not_found ->
<<Context/binary, ": not found"/utf8>>;
{connection_error, Reason@2} ->
<<<<Context/binary, ": connection error: "/utf8>>/binary,
Reason@2/binary>>
end,
gleam_stdlib:println_error(<<"gabsurd worker: "/utf8, Msg/binary>>).
-file("src/gabsurd/worker.gleam", 406).
-spec power_of_2(integer()) -> integer().
power_of_2(N) ->
case N of
0 ->
1;
_ ->
2 * power_of_2(N - 1)
end.
-file("src/gabsurd/worker.gleam", 399).
?DOC(" Calculate exponential backoff: base * 2^errors.\n").
-spec exponential_backoff(integer(), integer()) -> integer().
exponential_backoff(Base, Errors) ->
case Errors of
0 ->
Base;
_ ->
Base * power_of_2(Errors)
end.
-file("src/gabsurd/worker.gleam", 386).
-spec string_starts_with(binary(), binary()) -> boolean().
string_starts_with(Haystack, Prefix) ->
Hay_len = string:length(Haystack),
Pre_len = string:length(Prefix),
case Hay_len < Pre_len of
true ->
false;
false ->
Prefix_slice = gleam@string:slice(Haystack, 0, Pre_len),
Prefix_slice =:= Prefix
end.
-file("src/gabsurd/worker.gleam", 369).
?DOC(
" Handle errors from complete/fail calls.\n"
"\n"
" The Absurd schema raises SQLSTATE AB001 (cancelled) and AB002 (already\n"
" failed) when you try to complete or fail a run that is already in a\n"
" terminal state. The official SDKs silently swallow these errors. We\n"
" log unexpected errors but silently ignore terminal-state conflicts.\n"
).
-spec handle_completion_error(binary(), gabsurd@client:gabsurd_error()) -> nil.
handle_completion_error(Context, Error) ->
case Error of
{query_error, Reason} when Reason =:= <<"AB002"/utf8>> ->
nil;
{query_error, Reason@1} ->
case string_starts_with(Reason@1, <<"AB0"/utf8>>) of
true ->
nil;
false ->
log_error(Context, Error)
end;
_ ->
log_error(Context, Error)
end.
-file("src/gabsurd/worker.gleam", 293).
-spec execute_task(worker_state(), gabsurd@task:claim()) -> nil.
execute_task(State, Claim) ->
case gleam_stdlib:map_get(
erlang:element(3, State),
erlang:element(5, Claim)
) of
{ok, Handler} ->
Ctx = {context,
erlang:element(2, erlang:element(2, State)),
erlang:element(3, erlang:element(2, State)),
Claim,
erlang:element(6, erlang:element(2, State))},
case (erlang:element(3, Handler))(Ctx) of
{complete, Result_json} ->
case gabsurd@task:complete(
erlang:element(2, erlang:element(2, State)),
erlang:element(3, erlang:element(2, State)),
erlang:element(2, Claim),
Result_json
) of
{ok, _} ->
nil;
{error, Error} ->
handle_completion_error(<<"complete"/utf8>>, Error)
end;
{fail, Error_json} ->
case erlang:element(4, Handler) of
{some, Hook} ->
Hook(Ctx, Error_json);
none ->
nil
end,
case gabsurd@task:fail(
erlang:element(2, erlang:element(2, State)),
erlang:element(3, erlang:element(2, State)),
erlang:element(2, Claim),
Error_json
) of
{ok, _} ->
nil;
{error, Error@1} ->
handle_completion_error(<<"fail"/utf8>>, Error@1)
end;
suspend ->
nil
end;
{error, _} ->
_ = gabsurd@task:schedule_run(
erlang:element(2, erlang:element(2, State)),
erlang:element(3, erlang:element(2, State)),
erlang:element(2, Claim),
60
),
gleam_stdlib:println_error(
<<<<"gabsurd worker: deferred unknown task \""/utf8,
(erlang:element(5, Claim))/binary>>/binary,
"\" (no handler registered)"/utf8>>
)
end.
-file("src/gabsurd/worker.gleam", 239).
-spec handle_message(worker_state(), message()) -> gleam@otp@actor:next(worker_state(), message()).
handle_message(State, Message) ->
case Message of
poll ->
Result = gabsurd@task:claim(
erlang:element(2, erlang:element(2, State)),
erlang:element(3, erlang:element(2, State)),
erlang:element(4, erlang:element(2, State)),
erlang:element(6, erlang:element(2, State)),
erlang:element(7, erlang:element(2, State))
),
case Result of
{ok, Claims} ->
State@1 = {worker_state,
erlang:element(2, State),
erlang:element(3, State),
erlang:element(4, State),
0},
gleam@list:each(
Claims,
fun(Claim) -> execute_task(State@1, Claim) end
),
_ = gleam@erlang@process:send_after(
erlang:element(4, State@1),
erlang:element(5, erlang:element(2, State@1)),
poll
),
gleam@otp@actor:continue(State@1);
{error, Error} ->
Errors = erlang:element(5, State) + 1,
State@2 = {worker_state,
erlang:element(2, State),
erlang:element(3, State),
erlang:element(4, State),
Errors},
Backoff = gleam@int:min(
exponential_backoff(
erlang:element(5, erlang:element(2, State@2)),
Errors
),
erlang:element(8, erlang:element(2, State@2))
),
log_error(<<"claim error"/utf8>>, Error),
_ = gleam@erlang@process:send_after(
erlang:element(4, State@2),
Backoff,
poll
),
gleam@otp@actor:continue(State@2)
end;
shutdown ->
gleam@otp@actor:stop()
end.
-file("src/gabsurd/worker.gleam", 176).
?DOC(" Start a worker actor.\n").
-spec start(config()) -> {ok, gleam@otp@actor:started(worker())} |
{error, gleam@otp@actor:start_error()}.
start(Config) ->
Handler_map = begin
_pipe = erlang:element(9, Config),
_pipe@1 = gleam@list:map(_pipe, fun(H) -> {erlang:element(2, H), H} end),
maps:from_list(_pipe@1)
end,
Init = fun(Subject) ->
State = {worker_state, Config, Handler_map, Subject, 0},
_ = gleam@erlang@process:send_after(
Subject,
erlang:element(5, Config),
poll
),
_pipe@2 = gleam@otp@actor:initialised(State),
_pipe@3 = gleam@otp@actor:returning(_pipe@2, {worker, Subject}),
{ok, _pipe@3}
end,
_pipe@4 = gleam@otp@actor:new_with_initialiser(5000, Init),
_pipe@5 = gleam@otp@actor:on_message(_pipe@4, fun handle_message/2),
gleam@otp@actor:start(_pipe@5).
-file("src/gabsurd/worker.gleam", 203).
?DOC(" Stop a worker gracefully.\n").
-spec stop(worker()) -> nil.
stop(Worker) ->
gleam@erlang@process:send(erlang:element(2, Worker), shutdown).
-file("src/gabsurd/worker.gleam", 208).
?DOC(" Create a child spec for adding to a static_supervisor.\n").
-spec child_spec(binary(), config()) -> gleam@otp@supervision:child_specification(worker()).
child_spec(_, Config) ->
gleam@otp@supervision:worker(fun() -> start(Config) end).
-file("src/gabsurd/worker.gleam", 218).
?DOC(
" Create a list of child specs for a worker pool (N workers).\n"
" Each worker gets a unique worker_id incorporating a unique integer\n"
" to avoid collisions between pools.\n"
).
-spec pool_child_specs(binary(), config(), integer()) -> list(gleam@otp@supervision:child_specification(worker())).
pool_child_specs(Name, Config, Count) ->
Unique = erlang:unique_integer(),
gleam@list:index_fold(
gleam@list:repeat(nil, Count),
[],
fun(Acc, _, I) ->
I@1 = I + 1,
Wid = <<<<<<<<Name/binary, "_"/utf8>>/binary,
(erlang:integer_to_binary(Unique))/binary>>/binary,
"_"/utf8>>/binary,
(erlang:integer_to_binary(I@1))/binary>>,
Config@1 = {config,
erlang:element(2, Config),
erlang:element(3, Config),
Wid,
erlang:element(5, Config),
erlang:element(6, Config),
erlang:element(7, Config),
erlang:element(8, Config),
erlang:element(9, Config)},
[gleam@otp@supervision:worker(fun() -> start(Config@1) end) | Acc]
end
).