-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} -> <<<>/binary, Reason/binary>>; {unexpected_row_count, Reason@1} -> <<<>/binary, Reason@1/binary>>; not_found -> <>; {connection_error, Reason@2} -> <<<>/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 = <<<<<<<>/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 ).