Current section
Files
Jump to
Current section
Files
src/gabsurd@task.erl
-module(gabsurd@task).
-compile([no_auto_import, nowarn_unused_vars, nowarn_unused_function, nowarn_nomatch, inline]).
-define(FILEPATH, "src/gabsurd/task.gleam").
-export([new_options/0, with_max_attempts/2, with_retry_strategy/2, with_cancellation/2, with_headers/2, with_idempotency_key/2, encode_options/1, spawn/5, claim/5, complete/4, fail/4, fail_with_retry/5, extend_claim/4, schedule_run/4, cancel/3, get_result/3, retry/4]).
-export_type([spawn_info/0, claim/0, task_result/0, spawn_options/0, retry_strategy/0, cancellation/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(
" Task lifecycle operations for the Absurd durable workflow system.\n"
" Provides high-level functions for spawning, claiming, completing,\n"
" failing, and cancelling tasks.\n"
).
-type spawn_info() :: {spawn_info,
bitstring(),
bitstring(),
integer(),
boolean()}.
-type claim() :: {claim,
bitstring(),
bitstring(),
integer(),
binary(),
binary(),
binary(),
integer(),
binary(),
binary(),
binary()}.
-type task_result() :: {task_result, binary(), binary(), binary()}.
-type spawn_options() :: {spawn_options,
gleam@option:option(integer()),
gleam@option:option(retry_strategy()),
gleam@option:option(cancellation()),
gleam@option:option(gleam@json:json()),
gleam@option:option(binary())}.
-type retry_strategy() :: {fixed_retry, integer()} |
{exponential_retry, integer(), float(), gleam@option:option(float())}.
-type cancellation() :: {cancellation, integer()}.
-file("src/gabsurd/task.gleam", 92).
?DOC(" Create empty spawn options (all fields default to absent).\n").
-spec new_options() -> spawn_options().
new_options() ->
{spawn_options, none, none, none, none, none}.
-file("src/gabsurd/task.gleam", 103).
?DOC(" Set the maximum number of attempts for this task.\n").
-spec with_max_attempts(spawn_options(), integer()) -> spawn_options().
with_max_attempts(Options, Max) ->
{spawn_options,
{some, Max},
erlang:element(3, Options),
erlang:element(4, Options),
erlang:element(5, Options),
erlang:element(6, Options)}.
-file("src/gabsurd/task.gleam", 108).
?DOC(" Set the retry strategy for failed tasks.\n").
-spec with_retry_strategy(spawn_options(), retry_strategy()) -> spawn_options().
with_retry_strategy(Options, Strategy) ->
{spawn_options,
erlang:element(2, Options),
{some, Strategy},
erlang:element(4, Options),
erlang:element(5, Options),
erlang:element(6, Options)}.
-file("src/gabsurd/task.gleam", 116).
?DOC(" Set the cancellation policy.\n").
-spec with_cancellation(spawn_options(), cancellation()) -> spawn_options().
with_cancellation(Options, Cancellation) ->
{spawn_options,
erlang:element(2, Options),
erlang:element(3, Options),
{some, Cancellation},
erlang:element(5, Options),
erlang:element(6, Options)}.
-file("src/gabsurd/task.gleam", 124).
?DOC(" Set headers (arbitrary JSON metadata).\n").
-spec with_headers(spawn_options(), gleam@json:json()) -> spawn_options().
with_headers(Options, Headers) ->
{spawn_options,
erlang:element(2, Options),
erlang:element(3, Options),
erlang:element(4, Options),
{some, Headers},
erlang:element(6, Options)}.
-file("src/gabsurd/task.gleam", 132).
?DOC(" Set an idempotency key to prevent duplicate task creation.\n").
-spec with_idempotency_key(spawn_options(), binary()) -> spawn_options().
with_idempotency_key(Options, Key) ->
{spawn_options,
erlang:element(2, Options),
erlang:element(3, Options),
erlang:element(4, Options),
erlang:element(5, Options),
{some, Key}}.
-file("src/gabsurd/task.gleam", 191).
-spec encode_cancellation(cancellation()) -> gleam@json:json().
encode_cancellation(C) ->
case C of
{cancellation, Max_duration} ->
gleam@json:object(
[{<<"max_duration"/utf8>>, gleam@json:int(Max_duration)}]
)
end.
-file("src/gabsurd/task.gleam", 168).
-spec encode_retry_strategy(retry_strategy()) -> gleam@json:json().
encode_retry_strategy(Strategy) ->
case Strategy of
{fixed_retry, Base_seconds} ->
gleam@json:object(
[{<<"kind"/utf8>>, gleam@json:string(<<"fixed"/utf8>>)},
{<<"base_seconds"/utf8>>, gleam@json:int(Base_seconds)}]
);
{exponential_retry, Base_seconds@1, Factor, Max_seconds} ->
Entries = [{<<"kind"/utf8>>,
gleam@json:string(<<"exponential"/utf8>>)},
{<<"base_seconds"/utf8>>, gleam@json:int(Base_seconds@1)},
{<<"factor"/utf8>>, gleam@json:float(Factor)}],
Entries@1 = case Max_seconds of
{some, Max} ->
[{<<"max_seconds"/utf8>>, gleam@json:float(Max)} | Entries];
none ->
Entries
end,
gleam@json:object(Entries@1)
end.
-file("src/gabsurd/task.gleam", 140).
?DOC(" Encode spawn options to a JSON string for the database.\n").
-spec encode_options(spawn_options()) -> binary().
encode_options(Options) ->
Entries = [],
Entries@1 = case erlang:element(2, Options) of
{some, Max} ->
[{<<"max_attempts"/utf8>>, gleam@json:int(Max)} | Entries];
none ->
Entries
end,
Entries@2 = case erlang:element(3, Options) of
{some, Strategy} ->
[{<<"retry_strategy"/utf8>>, encode_retry_strategy(Strategy)} |
Entries@1];
none ->
Entries@1
end,
Entries@3 = case erlang:element(4, Options) of
{some, C} ->
[{<<"cancellation"/utf8>>, encode_cancellation(C)} | Entries@2];
none ->
Entries@2
end,
Entries@4 = case erlang:element(5, Options) of
{some, H} ->
[{<<"headers"/utf8>>, H} | Entries@3];
none ->
Entries@3
end,
Entries@5 = case erlang:element(6, Options) of
{some, Key} ->
[{<<"idempotency_key"/utf8>>, gleam@json:string(Key)} | Entries@4];
none ->
Entries@4
end,
gleam@json:to_string(gleam@json:object(Entries@5)).
-file("src/gabsurd/task.gleam", 219).
?DOC(
" Spawn a new task in a queue with typed options.\n"
"\n"
" `params` is a `json.Json` value — use `json.object`, `json.string`, etc.\n"
" to build it. `options` is a `SpawnOptions` record — use `new_options()` and\n"
" `with_*` builders.\n"
"\n"
" ## Example\n"
"\n"
" ```gleam\n"
" let assert Ok(info) = task.spawn(\n"
" db,\n"
" \"emails\",\n"
" \"send_welcome\",\n"
" json.object([#(\"to\", json.string(\"user@example.com\"))]),\n"
" task.new_options() |> task.with_max_attempts(3),\n"
" )\n"
" ```\n"
).
-spec spawn(
gabsurd@client:db(),
binary(),
binary(),
gleam@json:json(),
spawn_options()
) -> {ok, spawn_info()} | {error, gabsurd@client:gabsurd_error()}.
spawn(Db, Queue_name, Task_name, Params, Options) ->
Params_str = gleam@json:to_string(Params),
Options_str = encode_options(Options),
gleam@result:'try'(
gabsurd@client:query_one(
Db,
gabsurd@sql:spawn_task(
Queue_name,
Task_name,
Params_str,
Options_str
)
),
fun(Row) ->
{ok,
{spawn_info,
erlang:element(2, Row),
erlang:element(3, Row),
erlang:element(4, Row),
erlang:element(5, Row)}}
end
).
-file("src/gabsurd/task.gleam", 243).
?DOC(" Claim available tasks from a queue for a worker.\n").
-spec claim(gabsurd@client:db(), binary(), binary(), integer(), integer()) -> {ok,
list(claim())} |
{error, gabsurd@client:gabsurd_error()}.
claim(Db, Queue_name, Worker_id, Claim_timeout, Qty) ->
gleam@result:'try'(
gabsurd@client:query_many(
Db,
gabsurd@sql:claim_task(Queue_name, Worker_id, Claim_timeout, Qty)
),
fun(Rows) ->
{ok,
gleam@list:map(
Rows,
fun(Row) ->
{claim,
erlang:element(2, Row),
erlang:element(3, Row),
erlang:element(4, Row),
erlang:element(5, Row),
erlang:element(6, Row),
erlang:element(7, Row),
erlang:element(8, Row),
erlang:element(9, Row),
erlang:element(10, Row),
erlang:element(11, Row)}
end
)}
end
).
-file("src/gabsurd/task.gleam", 275).
?DOC(" Mark a run as completed with optional result state.\n").
-spec complete(gabsurd@client:db(), binary(), bitstring(), gleam@json:json()) -> {ok,
nil} |
{error, gabsurd@client:gabsurd_error()}.
complete(Db, Queue_name, Run_id, State) ->
gabsurd@client:exec(
Db,
gabsurd@sql:complete_run(
Queue_name,
Run_id,
gleam@json:to_string(State)
)
).
-file("src/gabsurd/task.gleam", 289).
?DOC(
" Mark a run as failed with a reason.\n"
" Passes NULL for retry_at so the queue's retry policy controls retries.\n"
).
-spec fail(gabsurd@client:db(), binary(), bitstring(), gleam@json:json()) -> {ok,
nil} |
{error, gabsurd@client:gabsurd_error()}.
fail(Db, Queue_name, Run_id, Reason) ->
gabsurd@client:exec(
Db,
gabsurd@sql:fail_run(Queue_name, Run_id, gleam@json:to_string(Reason))
).
-file("src/gabsurd/task.gleam", 302).
?DOC(" Mark a run as failed and schedule a retry at a specific time.\n").
-spec fail_with_retry(
gabsurd@client:db(),
binary(),
bitstring(),
gleam@json:json(),
gleam@time@timestamp:timestamp()
) -> {ok, nil} | {error, gabsurd@client:gabsurd_error()}.
fail_with_retry(Db, Queue_name, Run_id, Reason, Retry_at) ->
gabsurd@client:exec(
Db,
gabsurd@sql:fail_run_with_retry(
Queue_name,
Run_id,
gleam@json:to_string(Reason),
Retry_at
)
).
-file("src/gabsurd/task.gleam", 328).
?DOC(
" Extend a worker's claim lease on a run by `extend_by` seconds.\n"
"\n"
" This is the manual heartbeat mechanism. The primary lease extension\n"
" mechanism is `checkpoint.set` which calls `set_task_checkpoint_state`\n"
" with `extend_claim_by` — every checkpoint write extends the lease.\n"
"\n"
" Use this function when you have long-running work between checkpoints\n"
" and need to keep the lease alive.\n"
).
-spec extend_claim(gabsurd@client:db(), binary(), bitstring(), integer()) -> {ok,
nil} |
{error, gabsurd@client:gabsurd_error()}.
extend_claim(Db, Queue_name, Run_id, Extend_by) ->
gabsurd@client:exec(
Db,
gabsurd@sql:extend_claim(Queue_name, Run_id, Extend_by)
).
-file("src/gabsurd/task.gleam", 341).
?DOC(
" Schedule a run to become available again at a future time.\n"
" Used for deferring unknown tasks during rolling deployments.\n"
"\n"
" `defer_seconds` is how many seconds from now to reschedule.\n"
).
-spec schedule_run(gabsurd@client:db(), binary(), bitstring(), integer()) -> {ok,
nil} |
{error, gabsurd@client:gabsurd_error()}.
schedule_run(Db, Queue_name, Run_id, Defer_seconds) ->
Now = gleam@time@timestamp:system_time(),
Wake_at = gleam@time@timestamp:add(
Now,
gleam@time@duration:seconds(Defer_seconds)
),
gabsurd@client:exec(
Db,
gabsurd@sql:schedule_run(Queue_name, Run_id, Wake_at)
).
-file("src/gabsurd/task.gleam", 353).
?DOC(" Cancel a task by its task_id.\n").
-spec cancel(gabsurd@client:db(), binary(), bitstring()) -> {ok, nil} |
{error, gabsurd@client:gabsurd_error()}.
cancel(Db, Queue_name, Task_id) ->
gabsurd@client:exec(Db, gabsurd@sql:cancel_task(Queue_name, Task_id)).
-file("src/gabsurd/task.gleam", 362).
?DOC(" Get the result of a completed task.\n").
-spec get_result(gabsurd@client:db(), binary(), bitstring()) -> {ok,
task_result()} |
{error, gabsurd@client:gabsurd_error()}.
get_result(Db, Queue_name, Task_id) ->
gleam@result:'try'(
gabsurd@client:query_one(
Db,
gabsurd@sql:get_task_result(Queue_name, Task_id)
),
fun(Row) ->
{ok,
{task_result,
erlang:element(3, Row),
erlang:element(4, Row),
erlang:element(5, Row)}}
end
).
-file("src/gabsurd/task.gleam", 381).
?DOC(" Retry a task with typed options.\n").
-spec retry(gabsurd@client:db(), binary(), bitstring(), spawn_options()) -> {ok,
spawn_info()} |
{error, gabsurd@client:gabsurd_error()}.
retry(Db, Queue_name, Task_id, Options) ->
Options_str = encode_options(Options),
gleam@result:'try'(
gabsurd@client:query_one(
Db,
gabsurd@sql:retry_task(Queue_name, Task_id, Options_str)
),
fun(Row) ->
{ok,
{spawn_info,
erlang:element(2, Row),
erlang:element(3, Row),
erlang:element(4, Row),
erlang:element(5, Row)}}
end
).