-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 => <>, 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 => <>, 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.