Packages

A Gleam SDK for the Absurd durable workflow system, with type-safe database access via Parrot and OTP worker actors

Current section

Files

Jump to
gabsurd src gabsurd@context.erl
Raw

src/gabsurd@context.erl

-module(gabsurd@context).
-compile([no_auto_import, nowarn_unused_vars, nowarn_unused_function, nowarn_nomatch, inline]).
-define(FILEPATH, "src/gabsurd/context.gleam").
-export([task_id/1, run_id/1, params/1, task_name/1, attempt/1, claim_timeout/1, step/4, get_checkpoint/2, set_checkpoint/3, heartbeat/1, await_event/3]).
-export_type([context/0, event_result/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(
" Execution context for a running task.\n"
"\n"
" Constructed by the worker and passed to your handler. Encapsulates\n"
" the database connection, queue, claim details, and claim timeout so\n"
" you don't have to pass them around. Provides high-level operations\n"
" for idempotent steps, checkpoints, heartbeats, and event coordination.\n"
"\n"
" ## Example\n"
"\n"
" ```gleam\n"
" let handler = Handler(\n"
" task_name: \"process_order\",\n"
" execute: fn(ctx) {\n"
" case order_workflow(ctx) {\n"
" Ok(Nil) -> Complete(json.object([#(\"status\", json.string(\"done\"))]))\n"
" Error(e) -> Fail(encode_error(e))\n"
" }\n"
" },\n"
" on_error: option.None,\n"
" )\n"
"\n"
" fn order_workflow(ctx) -> Result(Nil, GabsurdError) {\n"
" use _ <- result.try(ctx.step(\"charge\", decode.success(Nil), fn() {\n"
" charge_card(decode_params(ctx.params(ctx)))\n"
" json.null()\n"
" }))\n"
" use _ <- result.try(ctx.step(\"reserve\", decode.success(Nil), fn() {\n"
" reserve_inventory(decode_params(ctx.params(ctx)))\n"
" json.null()\n"
" }))\n"
" Ok(Nil)\n"
" }\n"
" ```\n"
).
-type context() :: {context,
gabsurd@client:db(),
binary(),
gabsurd@task:claim(),
integer()}.
-type event_result() :: {received, binary()} | suspended.
-file("src/gabsurd/context.gleam", 70).
?DOC(" The task's unique identifier.\n").
-spec task_id(context()) -> bitstring().
task_id(Ctx) ->
erlang:element(3, erlang:element(4, Ctx)).
-file("src/gabsurd/context.gleam", 75).
?DOC(" The current run's unique identifier.\n").
-spec run_id(context()) -> bitstring().
run_id(Ctx) ->
erlang:element(2, erlang:element(4, Ctx)).
-file("src/gabsurd/context.gleam", 80).
?DOC(" The task parameters as a raw JSON string.\n").
-spec params(context()) -> binary().
params(Ctx) ->
erlang:element(6, erlang:element(4, Ctx)).
-file("src/gabsurd/context.gleam", 85).
?DOC(" The task name.\n").
-spec task_name(context()) -> binary().
task_name(Ctx) ->
erlang:element(5, erlang:element(4, Ctx)).
-file("src/gabsurd/context.gleam", 90).
?DOC(" The current attempt number (1-based).\n").
-spec attempt(context()) -> integer().
attempt(Ctx) ->
erlang:element(4, erlang:element(4, Ctx)).
-file("src/gabsurd/context.gleam", 95).
?DOC(" The claim timeout in seconds.\n").
-spec claim_timeout(context()) -> integer().
claim_timeout(Ctx) ->
erlang:element(5, Ctx).
-file("src/gabsurd/context.gleam", 132).
?DOC(
" Run an idempotent step identified by name.\n"
"\n"
" If the checkpoint already exists (from a previous attempt), the stored\n"
" value is decoded with `decoder` and returned without re-running `run`.\n"
" If not, `run` is executed, the result is persisted as a checkpoint, and\n"
" the claim lease is extended by `claim_timeout` seconds.\n"
"\n"
" The `decoder` parameter is required because Gleam's `json.Json` type is\n"
" write-only — values loaded from the database must be parsed with an\n"
" explicit decoder. For steps that don't need a return value, use\n"
" `decode.success(Nil)`.\n"
"\n"
" ## Example\n"
"\n"
" ```gleam\n"
" // Step that returns a value:\n"
" use charge_id <- result.try(\n"
" ctx.step(\"charge\", decode.field(\"charge_id\", decode.string), fn() {\n"
" let result = charge_card(...)\n"
" json.object([#(\"charge_id\", json.string(result.id))])\n"
" }),\n"
" )\n"
"\n"
" // Step that doesn't return a value:\n"
" use _ <- result.try(ctx.step(\"notify\", decode.success(Nil), fn() {\n"
" send_email(...)\n"
" json.null()\n"
" }))\n"
" ```\n"
).
-spec step(
context(),
binary(),
gleam@dynamic@decode:decoder(OCW),
fun(() -> gleam@json:json())
) -> {ok, OCW} | {error, gabsurd@client:gabsurd_error()}.
step(Ctx, Name, Decoder, Run) ->
case gabsurd@checkpoint:get(
erlang:element(2, Ctx),
erlang:element(3, Ctx),
erlang:element(3, erlang:element(4, Ctx)),
Name,
false
) of
{ok, {some, Cp}} ->
Value@1 = case gleam@json:parse(erlang:element(3, Cp), Decoder) of
{ok, Value} -> Value;
_assert_fail ->
erlang:error(#{gleam_error => let_assert,
message => <<"Pattern match failed, no pattern matched the value."/utf8>>,
file => <<?FILEPATH/utf8>>,
module => <<"gabsurd/context"/utf8>>,
function => <<"step"/utf8>>,
line => 141,
value => _assert_fail,
start => 4212,
'end' => 4264,
pattern_start => 4223,
pattern_end => 4232})
end,
{ok, Value@1};
{ok, none} ->
Json_value = Run(),
Json_string = gleam@json:to_string(Json_value),
case gabsurd@checkpoint:set(
erlang:element(2, Ctx),
erlang:element(3, Ctx),
erlang:element(3, erlang:element(4, Ctx)),
Name,
Json_value,
erlang:element(2, erlang:element(4, Ctx)),
erlang:element(5, Ctx)
) of
{ok, nil} ->
Value@3 = case gleam@json:parse(Json_string, Decoder) of
{ok, Value@2} -> Value@2;
_assert_fail@1 ->
erlang:error(#{gleam_error => let_assert,
message => <<"Pattern match failed, no pattern matched the value."/utf8>>,
file => <<?FILEPATH/utf8>>,
module => <<"gabsurd/context"/utf8>>,
function => <<"step"/utf8>>,
line => 161,
value => _assert_fail@1,
start => 4777,
'end' => 4832,
pattern_start => 4788,
pattern_end => 4797})
end,
{ok, Value@3};
{error, E} ->
{error, E}
end;
{error, E@1} ->
{error, E@1}
end.
-file("src/gabsurd/context.gleam", 177).
?DOC(
" Get a checkpoint's raw JSON state string.\n"
" Returns `Ok(Some(json_string))` if found, `Ok(None)` if not found.\n"
).
-spec get_checkpoint(context(), binary()) -> {ok, gleam@option:option(binary())} |
{error, gabsurd@client:gabsurd_error()}.
get_checkpoint(Ctx, Name) ->
case gabsurd@checkpoint:get(
erlang:element(2, Ctx),
erlang:element(3, Ctx),
erlang:element(3, erlang:element(4, Ctx)),
Name,
false
) of
{ok, {some, Cp}} ->
{ok, {some, erlang:element(3, Cp)}};
{ok, none} ->
{ok, none};
{error, E} ->
{error, E}
end.
-file("src/gabsurd/context.gleam", 189).
?DOC(" Set a checkpoint and extend the claim lease.\n").
-spec set_checkpoint(context(), binary(), gleam@json:json()) -> {ok, nil} |
{error, gabsurd@client:gabsurd_error()}.
set_checkpoint(Ctx, Name, State) ->
gabsurd@checkpoint:set(
erlang:element(2, Ctx),
erlang:element(3, Ctx),
erlang:element(3, erlang:element(4, Ctx)),
Name,
State,
erlang:element(2, erlang:element(4, Ctx)),
erlang:element(5, Ctx)
).
-file("src/gabsurd/context.gleam", 210).
?DOC(" Extend the claim lease by `claim_timeout` seconds.\n").
-spec heartbeat(context()) -> {ok, nil} |
{error, gabsurd@client:gabsurd_error()}.
heartbeat(Ctx) ->
gabsurd@task:extend_claim(
erlang:element(2, Ctx),
erlang:element(3, Ctx),
erlang:element(2, erlang:element(4, Ctx)),
erlang:element(5, Ctx)
).
-file("src/gabsurd/context.gleam", 228).
?DOC(
" Await an external event. If the event is already available, returns\n"
" `Received(payload)`. If not, the task is put to sleep and returns\n"
" `Suspended` — your handler should return `Suspend` in this case.\n"
"\n"
" `timeout` is in seconds. Set to `0` for no timeout.\n"
).
-spec await_event(context(), binary(), integer()) -> {ok, event_result()} |
{error, gabsurd@client:gabsurd_error()}.
await_event(Ctx, Event_name, Timeout) ->
case gabsurd@event:await(
erlang:element(2, Ctx),
erlang:element(3, Ctx),
erlang:element(3, erlang:element(4, Ctx)),
erlang:element(2, erlang:element(4, Ctx)),
<<"$await:"/utf8, Event_name/binary>>,
Event_name,
Timeout
) of
{ok, Result} ->
case erlang:element(2, Result) of
true ->
{ok, suspended};
false ->
{ok, {received, erlang:element(3, Result)}}
end;
{error, E} ->
{error, E}
end.