-module(m25). -compile([no_auto_import, nowarn_unused_vars, nowarn_unused_function, nowarn_nomatch]). -define(FILEPATH, "src/m25.gleam"). -export([main/0, new/1, add_queue/2, new_job/1, schedule/2, timeout/2, retry/3, unique_key/2, get_job/3, cancel_job/3, start/2, supervised/2, enqueue/3]). -export_type([queue/3, m25/0, job/1, job_cancel_error/0, job_status/0, failure_reason/0, job_id/0, job_record/3, job_record_decode_error/0, job_record_fetch_error/0, job_update_error/0, queue_manager_msg/3, queue_manager_state/3, process_jobs_error/0, job_executor_message/2, job_executor_state/3]). -if(?OTP_RELEASE >= 27). -define(MODULEDOC(Str), -moduledoc(Str)). -define(DOC(Str), -doc(Str)). -else. -define(MODULEDOC(Str), -compile([])). -define(DOC(Str), -compile([])). -endif. -type queue(YYO, YYP, YYQ) :: {queue, binary(), integer(), fun((YYO) -> gleam@json:json()), gleam@dynamic@decode:decoder(YYO), fun((YYP) -> gleam@json:json()), gleam@dynamic@decode:decoder(YYP), fun((YYQ) -> gleam@json:json()), gleam@dynamic@decode:decoder(YYQ), fun((YYO) -> {ok, YYP} | {error, YYQ}), gleam@time@duration:duration(), integer(), integer(), integer(), integer(), integer()}. -opaque m25() :: {m25, pog:connection(), list(queue(gleam@dynamic:dynamic_(), gleam@dynamic:dynamic_(), gleam@dynamic:dynamic_()))}. -opaque job(YYR) :: {job, YYR, gleam@option:option(gleam@time@timestamp:timestamp()), gleam@option:option(gleam@time@duration:duration()), integer(), gleam@option:option(gleam@time@duration:duration()), gleam@option:option(binary())}. -type job_cancel_error() :: {job_cancel_fetch_error, job_record_fetch_error()} | {invalid_state, job_status()}. -type job_status() :: pending | reserved | executing | succeeded | failed | cancelled. -type failure_reason() :: errored | heartbeat_timeout | job_timeout | crash. -type job_id() :: {job_id, youid@uuid:uuid()}. -type job_record(YYS, YYT, YYU) :: {job_record, job_id(), binary(), gleam@time@timestamp:timestamp(), gleam@option:option(gleam@time@timestamp:timestamp()), YYS, gleam@option:option(gleam@time@timestamp:timestamp()), gleam@option:option(gleam@time@timestamp:timestamp()), gleam@option:option(gleam@time@timestamp:timestamp()), gleam@option:option(gleam@time@timestamp:timestamp()), gleam@time@duration:duration(), job_status(), gleam@option:option(YYT), gleam@option:option(gleam@time@timestamp:timestamp()), gleam@option:option(gleam@time@timestamp:timestamp()), gleam@option:option(failure_reason()), gleam@option:option(YYU), integer(), integer(), gleam@option:option(job_id()), gleam@option:option(job_id()), gleam@time@duration:duration(), gleam@option:option(binary())}. -type job_record_decode_error() :: {job_record_fetch_invalid_status_error, job_id(), binary()} | {job_record_fetch_invalid_failure_reason, job_id(), binary()} | {job_record_fetch_json_decode_error, job_id(), gleam@json:decode_error()}. -type job_record_fetch_error() :: {job_record_fetch_query_error, pog:query_error()} | {job_record_fetch_decode_errors, list(job_record_decode_error())} | no_job_record_found. -type job_update_error() :: {job_not_found, job_id()} | {too_many_jobs_returned, list(job_id())} | {job_update_query_error, pog:query_error()}. -type queue_manager_msg(YYV, YYW, YYX) :: process_jobs | {work_succeeded, job_id(), YYW} | {work_failed, job_id(), YYX} | {job_executor_down, gleam@erlang@process:down()} | {job_worker_down, job_id()} | shutdown | {gleam_phantom, YYV}. -type queue_manager_state(YYY, YYZ, YZA) :: {queue_manager_state, gleam@erlang@process:subject(queue_manager_msg(YYY, YYZ, YZA)), queue(YYY, YYZ, YZA), pog:connection(), m25@internal@bimap:bimap(job_id(), gleam@erlang@process:pid_())}. -type process_jobs_error() :: {process_jobs_query_error, binary(), pog:query_error()} | {process_jobs_fetch_error, job_record_fetch_error()}. -type job_executor_message(YZB, YZC) :: heartbeat | start_work | {execution_succeeded, YZB} | {execution_failed, YZC} | {worker_down, gleam@erlang@process:exit_message()} | {manager_down, gleam@erlang@process:down()}. -type job_executor_state(YZD, YZE, YZF) :: {job_executor_state, gleam@erlang@process:subject(job_executor_message(YZE, YZF)), pog:connection(), queue(YZD, YZE, YZF), job_id(), gleam@option:option(gleam@erlang@process:pid_()), gleam@erlang@process:subject(queue_manager_msg(YZD, YZE, YZF)), fun((YZD) -> {ok, YZE} | {error, YZF}), YZD}. -file("src/m25.gleam", 26). ?DOC(false). -spec main() -> nil. main() -> m25@internal@cli:run_cli(). -file("src/m25.gleam", 101). ?DOC( " Create a new M25 instance. It's recommended that you use a supervised `pog`\n" " connection.\n" "\n" " ```gleam\n" " let conn_name = process.new_name(\"db_connection\")\n" "\n" " let conn_child =\n" " pog.default_config(conn_name)\n" " |> pog.host(\"localhost\")\n" " |> pog.database(\"my_database\")\n" " |> pog.pool_size(15)\n" " |> pog.supervised\n" "\n" " // Create a connection that can be accessed by our queue handlers\n" " let conn = pog.named_connection(conn_name)\n" "\n" " let m25 = m25.new(conn)\n" " ```\n" ). -spec new(pog:connection()) -> m25(). new(Conn) -> {m25, Conn, []}. -file("src/m25.gleam", 120). ?DOC( " Register a queue to be used by M25. All of the input, output and error values must\n" " be serialisable to JSON so that they may be inserted into the database.\n" "\n" " Returns `Error(Nil)` if a queue with the same name has already been registered.\n" "\n" " ```gleam\n" " pub fn main() {\n" " let assert Ok(m25) = m25.new(conn)\n" " |> m25.add_queue(queue1)\n" " |> result.try(m25.add_queue(_, queue2))\n" " |> result.try(m25.add_queue(_, queue3))\n" "\n" " let assert Ok(_) = m25.start(m25)\n" " }\n" " ```\n" ). -spec add_queue(m25(), queue(any(), any(), any())) -> {ok, m25()} | {error, nil}. add_queue(M25, Queue) -> case gleam@list:find( erlang:element(3, M25), fun(Existing_queue) -> erlang:element(2, Queue) =:= erlang:element(2, Existing_queue) end ) of {ok, _} -> {error, nil}; {error, _} -> {ok, begin _record = M25, {m25, erlang:element(2, _record), [m25_ffi:coerce(Queue) | erlang:element(3, M25)]} end} end. -file("src/m25.gleam", 186). ?DOC( " Create a new job with default values and the given input. The input must match the\n" " input type of the queue you'll be enqueuing it to.\n" ). -spec new_job(AAAD) -> job(AAAD). new_job(Input) -> {job, Input, none, none, 1, none, none}. -file("src/m25.gleam", 199). ?DOC( " Schedule a job to be executed at a specific time. If that time is in the past, the\n" " job will be executed immediately.\n" ). -spec schedule(job(AAJG), gleam@time@timestamp:timestamp()) -> job(AAJG). schedule(Job, Scheduled_at) -> _record = Job, {job, erlang:element(2, _record), {some, Scheduled_at}, erlang:element(4, _record), erlang:element(5, _record), erlang:element(6, _record), erlang:element(7, _record)}. -file("src/m25.gleam", 204). ?DOC(" Add a timeout to a job that overrides the queue's default timeout.\n"). -spec timeout(job(AAJO), gleam@time@duration:duration()) -> job(AAJO). timeout(Job, Timeout) -> _record = Job, {job, erlang:element(2, _record), erlang:element(3, _record), {some, Timeout}, erlang:element(5, _record), erlang:element(6, _record), erlang:element(7, _record)}. -file("src/m25.gleam", 210). ?DOC( " Configure retry behavior for a job. If no retry delay is provided, the job will be\n" " retried immediately.\n" ). -spec retry( job(AAJW), integer(), gleam@option:option(gleam@time@duration:duration()) ) -> job(AAJW). retry(Job, Max_attempts, Retry_delay) -> _record = Job, {job, erlang:element(2, _record), erlang:element(3, _record), erlang:element(4, _record), Max_attempts, Retry_delay, erlang:element(7, _record)}. -file("src/m25.gleam", 217). ?DOC( " Set a unique key for a job. This will prevent the job being enqueued if it already\n" " exists in a non-errored state. If the only matching attempts have failed or\n" " crashed, the job can still be enqueued.\n" ). -spec unique_key(job(AAKC), binary()) -> job(AAKC). unique_key(Job, Unique_key) -> _record = Job, {job, erlang:element(2, _record), erlang:element(3, _record), erlang:element(4, _record), erlang:element(5, _record), erlang:element(6, _record), {some, Unique_key}}. -file("src/m25.gleam", 335). -spec job_status_from_string(binary()) -> {ok, job_status()} | {error, nil}. job_status_from_string(Maybe_status) -> case Maybe_status of <<"pending"/utf8>> -> {ok, pending}; <<"reserved"/utf8>> -> {ok, reserved}; <<"executing"/utf8>> -> {ok, executing}; <<"succeeded"/utf8>> -> {ok, succeeded}; <<"failed"/utf8>> -> {ok, failed}; <<"cancelled"/utf8>> -> {ok, cancelled}; _ -> {error, nil} end. -file("src/m25.gleam", 359). -spec failure_reason_to_string(failure_reason()) -> binary(). failure_reason_to_string(Failure_reason) -> case Failure_reason of errored -> <<"error"/utf8>>; crash -> <<"crash"/utf8>>; heartbeat_timeout -> <<"heartbeat_timeout"/utf8>>; job_timeout -> <<"job_timeout"/utf8>> end. -file("src/m25.gleam", 368). -spec failure_reason_from_string(binary()) -> {ok, failure_reason()} | {error, nil}. failure_reason_from_string(Reason) -> case Reason of <<"error"/utf8>> -> {ok, errored}; <<"crash"/utf8>> -> {ok, crash}; <<"heartbeat_timeout"/utf8>> -> {ok, heartbeat_timeout}; <<"job_timeout"/utf8>> -> {ok, job_timeout}; _ -> {error, nil} end. -file("src/m25.gleam", 427). -spec job_record_decode_error_to_string(job_record_decode_error()) -> binary(). job_record_decode_error_to_string(Error) -> case Error of {job_record_fetch_invalid_status_error, Job_id, Status} -> <<<<<<"Job "/utf8, (youid@uuid:to_string(erlang:element(2, Job_id)))/binary>>/binary, " has invalid status "/utf8>>/binary, Status/binary>>; {job_record_fetch_invalid_failure_reason, Job_id@1, Reason} -> <<<<<<"Invalid failure reason for job "/utf8, (youid@uuid:to_string(erlang:element(2, Job_id@1)))/binary>>/binary, ": "/utf8>>/binary, Reason/binary>>; {job_record_fetch_json_decode_error, Job_id@2, Error@1} -> <<<<<<"Failed to decode job "/utf8, (youid@uuid:to_string(erlang:element(2, Job_id@2)))/binary>>/binary, ": "/utf8>>/binary, (gleam@string:inspect(Error@1))/binary>> end. -file("src/m25.gleam", 444). -spec succeed_job(pog:connection(), job_id(), gleam@json:json()) -> {ok, pog:returned(nil)} | {error, pog:query_error()}. succeed_job(Conn, Job_id, Output) -> m25@internal@sql:succeed_job(Conn, erlang:element(2, Job_id), Output). -file("src/m25.gleam", 448). -spec error_job(pog:connection(), job_id(), gleam@json:json()) -> {ok, pog:returned(m25@internal@sql:error_job_row())} | {error, pog:query_error()}. error_job(Conn, Job_id, Error) -> m25@internal@sql:error_job(Conn, erlang:element(2, Job_id), Error). -file("src/m25.gleam", 452). -spec fail_job(pog:connection(), job_id(), failure_reason()) -> {ok, pog:returned(m25@internal@sql:fail_job_row())} | {error, pog:query_error()}. fail_job(Conn, Job_id, Reason) -> m25@internal@sql:fail_job( Conn, erlang:element(2, Job_id), failure_reason_to_string(Reason) ). -file("src/m25.gleam", 468). -spec decode_job_record_row( queue(AACN, AACO, AACP), m25@internal@sql_ext:job_record_row() ) -> {ok, job_record(AACN, AACO, AACP)} | {error, job_record_decode_error()}. decode_job_record_row(Queue, Job) -> Job_id = {job_id, erlang:element(2, Job)}, gleam@result:'try'( begin _pipe = job_status_from_string(erlang:element(12, Job)), gleam@result:replace_error( _pipe, {job_record_fetch_invalid_status_error, Job_id, erlang:element(12, Job)} ) end, fun(Status) -> gleam@result:'try'( begin _pipe@1 = gleam@json:parse( erlang:element(6, Job), erlang:element(5, Queue) ), gleam@result:map_error( _pipe@1, fun(_capture) -> {job_record_fetch_json_decode_error, Job_id, _capture} end ) end, fun(Input) -> Output_result = case erlang:element(13, Job) of none -> {ok, none}; {some, Output} -> _pipe@2 = gleam@json:parse( Output, erlang:element(7, Queue) ), _pipe@3 = gleam@result:map( _pipe@2, fun(Field@0) -> {some, Field@0} end ), gleam@result:map_error( _pipe@3, fun(_capture@1) -> {job_record_fetch_json_decode_error, Job_id, _capture@1} end ) end, gleam@result:'try'( Output_result, fun(Output@1) -> Error_result = case erlang:element(17, Job) of none -> {ok, none}; {some, Error} -> _pipe@4 = gleam@json:parse( Error, erlang:element(9, Queue) ), _pipe@5 = gleam@result:map( _pipe@4, fun(Field@0) -> {some, Field@0} end ), gleam@result:map_error( _pipe@5, fun(_capture@2) -> {job_record_fetch_json_decode_error, Job_id, _capture@2} end ) end, gleam@result:'try'( Error_result, fun(Error_data) -> Failure_reason_result = case erlang:element( 16, Job ) of none -> {ok, none}; {some, Reason} -> _pipe@6 = failure_reason_from_string( Reason ), _pipe@7 = gleam@result:map( _pipe@6, fun(Field@0) -> {some, Field@0} end ), gleam@result:replace_error( _pipe@7, {job_record_fetch_invalid_failure_reason, Job_id, Reason} ) end, gleam@result:'try'( Failure_reason_result, fun(Failure_reason) -> _pipe@8 = {job_record, Job_id, erlang:element(3, Job), erlang:element(4, Job), erlang:element(5, Job), Input, erlang:element(7, Job), erlang:element(8, Job), erlang:element(9, Job), erlang:element(10, Job), gleam@time@duration:seconds( erlang:element(11, Job) ), Status, Output@1, erlang:element(14, Job), erlang:element(15, Job), Failure_reason, Error_data, erlang:element(18, Job), erlang:element(19, Job), gleam@option:map( erlang:element(20, Job), fun(Field@0) -> {job_id, Field@0} end ), gleam@option:map( erlang:element(21, Job), fun(Field@0) -> {job_id, Field@0} end ), gleam@time@duration:seconds( erlang:element(22, Job) ), erlang:element(23, Job)}, {ok, _pipe@8} end ) end ) end ) end ) end ). -file("src/m25.gleam", 252). ?DOC(" Get a job from the database by its ID.\n"). -spec get_job(pog:connection(), queue(AABB, AABC, AABD), job_id()) -> {ok, job_record(AABB, AABC, AABD)} | {error, job_record_fetch_error()}. get_job(Conn, Queue, Id) -> gleam@result:'try'( begin _pipe = m25@internal@sql_ext:get_job(Conn, erlang:element(2, Id)), gleam@result:map_error( _pipe, fun(Field@0) -> {job_record_fetch_query_error, Field@0} end ) end, fun(Job) -> case erlang:element(3, Job) of [] -> {error, no_job_record_found}; [Row] -> _pipe@1 = decode_job_record_row(Queue, Row), gleam@result:map_error( _pipe@1, fun(Err) -> {job_record_fetch_decode_errors, [Err]} end ); _ -> erlang:error(#{gleam_error => panic, message => <<"unreachable"/utf8>>, file => <>, module => <<"m25"/utf8>>, function => <<"get_job"/utf8>>, line => 266}) end end ). -file("src/m25.gleam", 280). ?DOC( " Cancel a job from the database by its ID. You can only cancel a job that is in the\n" " `Pending` state.\n" ). -spec cancel_job(pog:connection(), queue(AABM, AABN, AABO), job_id()) -> {ok, job_record(AABM, AABN, AABO)} | {error, job_cancel_error()}. cancel_job(Conn, Queue, Id) -> gleam@result:'try'( begin _pipe = m25@internal@sql_ext:cancel_job(Conn, erlang:element(2, Id)), gleam@result:map_error( _pipe, fun(Err) -> {job_cancel_fetch_error, {job_record_fetch_query_error, Err}} end ) end, fun(Job) -> case erlang:element(3, Job) of [] -> {error, {job_cancel_fetch_error, no_job_record_found}}; [{Row, Cancel_outcome}] -> gleam@result:'try'( begin _pipe@1 = decode_job_record_row(Queue, Row), gleam@result:map_error( _pipe@1, fun(Err@1) -> {job_cancel_fetch_error, {job_record_fetch_decode_errors, [Err@1]}} end ) end, fun(Job_record) -> case Cancel_outcome of not_pending -> {error, {invalid_state, erlang:element(12, Job_record)}}; cancelled -> {ok, Job_record} end end ); _ -> erlang:error(#{gleam_error => panic, message => <<"unreachable"/utf8>>, file => <>, module => <<"m25"/utf8>>, function => <<"cancel_job"/utf8>>, line => 306}) end end ). -file("src/m25.gleam", 539). -spec decode_multiple_job_record_rows( queue(AACU, AACV, AACW), list(m25@internal@sql_ext:job_record_row()) ) -> {ok, list(job_record(AACU, AACV, AACW))} | {error, job_record_fetch_error()}. decode_multiple_job_record_rows(Queue, Jobs) -> {Valid_jobs, Invalid_jobs} = begin _pipe = Jobs, _pipe@1 = gleam@list:map( _pipe, fun(_capture) -> decode_job_record_row(Queue, _capture) end ), gleam@result:partition(_pipe@1) end, case Invalid_jobs of [] -> {ok, Valid_jobs}; Invalid -> {error, {job_record_fetch_decode_errors, Invalid}} end. -file("src/m25.gleam", 456). -spec reserve_jobs(pog:connection(), queue(AACG, AACH, AACI), integer()) -> {ok, list(job_record(AACG, AACH, AACI))} | {error, job_record_fetch_error()}. reserve_jobs(Conn, Queue, Limit) -> gleam@result:'try'( begin _pipe = m25@internal@sql_ext:reserve_jobs( Conn, erlang:element(2, Queue), Limit ), gleam@result:map_error( _pipe, fun(Field@0) -> {job_record_fetch_query_error, Field@0} end ) end, fun(Jobs) -> decode_multiple_job_record_rows(Queue, erlang:element(3, Jobs)) end ). -file("src/m25.gleam", 554). -spec finalise_job_reservations( pog:connection(), list(job_id()), list(job_id()) ) -> {ok, pog:returned(m25@internal@sql:finalise_job_reservations_row())} | {error, pog:query_error()}. finalise_job_reservations(Conn, Successful_job_ids, Failed_job_ids) -> m25@internal@sql:finalise_job_reservations( Conn, gleam@list:map(Successful_job_ids, fun(Id) -> erlang:element(2, Id) end), gleam@list:map(Failed_job_ids, fun(Id@1) -> erlang:element(2, Id@1) end) ). -file("src/m25.gleam", 566). -spec cleanup_stuck_reservations(pog:connection(), binary(), integer()) -> {ok, pog:returned(m25@internal@sql:cleanup_stuck_reservations_row())} | {error, pog:query_error()}. cleanup_stuck_reservations(Conn, Queue_name, Timeout) -> m25@internal@sql:cleanup_stuck_reservations( Conn, Queue_name, erlang:float(Timeout) / 1000.0 ). -file("src/m25.gleam", 607). -spec time_out_jobs(pog:connection(), binary()) -> {ok, pog:returned(m25@internal@sql:time_out_jobs_row())} | {error, pog:query_error()}. time_out_jobs(Conn, Queue_name) -> m25@internal@sql:time_out_jobs(Conn, Queue_name). -file("src/m25.gleam", 612). ?DOC(" Returns a boolean representing whether the job has hit a heartbeat timeout\n"). -spec execute_job_heartbeat(pog:connection(), job_id(), integer(), integer()) -> {ok, pog:returned(m25@internal@sql:heartbeat_row())} | {error, pog:query_error()}. execute_job_heartbeat(Conn, Job_id, Allowed_misses, Heartbeat_interval) -> m25@internal@sql:heartbeat( Conn, erlang:element(2, Job_id), Allowed_misses, erlang:float(Heartbeat_interval) / 1000.0 ). -file("src/m25.gleam", 626). -spec timestamp_to_unix_seconds_float(gleam@time@timestamp:timestamp()) -> float(). timestamp_to_unix_seconds_float(Timestamp) -> {Seconds, Nanoseconds} = gleam@time@timestamp:to_unix_seconds_and_nanoseconds( Timestamp ), erlang:float(Seconds) + (erlang:float(Nanoseconds) / 1000000000.0). -file("src/m25.gleam", 578). -spec insert_job( pog:connection(), binary(), gleam@option:option(gleam@time@timestamp:timestamp()), gleam@json:json(), float(), integer(), integer(), gleam@option:option(youid@uuid:uuid()), gleam@option:option(youid@uuid:uuid()), float(), gleam@option:option(binary()) ) -> {ok, pog:returned(m25@internal@sql_ext:job_record_row())} | {error, pog:query_error()}. insert_job( Conn, Queue_name, Scheduled_at, Input, Timeout, Attempt, Max_attempts, Original_attempt_id, Previous_attempt_id, Retry_delay, Unique_key ) -> m25@internal@sql_ext:insert_job( Conn, youid@uuid:v7(), Queue_name, gleam@option:map(Scheduled_at, fun timestamp_to_unix_seconds_float/1), Input, Timeout, Attempt, Max_attempts, Original_attempt_id, Previous_attempt_id, Retry_delay, Unique_key ). -file("src/m25.gleam", 681). -spec retry_jobs_if_needed(pog:connection(), list(youid@uuid:uuid())) -> {ok, pog:returned(nil)} | {error, pog:query_error()}. retry_jobs_if_needed(Conn, Job_ids) -> m25@internal@sql:retry_if_needed(Conn, Job_ids). -file("src/m25.gleam", 640). -spec handle_errored_job( pog:connection(), queue(any(), any(), AADS), job_id(), AADS ) -> {ok, pog:returned(nil)} | {error, pog:transaction_error(job_update_error())}. handle_errored_job(Conn, Queue, Job_id, Error) -> pog:transaction( Conn, fun(Conn@1) -> gleam@result:'try'( begin _pipe = error_job( Conn@1, Job_id, (erlang:element(8, Queue))(Error) ), gleam@result:map_error( _pipe, fun(Field@0) -> {job_update_query_error, Field@0} end ) end, fun(Failed_job) -> case erlang:element(3, Failed_job) of [] -> {error, {job_not_found, Job_id}}; [Row] -> _pipe@1 = retry_jobs_if_needed( Conn@1, [erlang:element(2, Row)] ), gleam@result:map_error( _pipe@1, fun(Field@0) -> {job_update_query_error, Field@0} end ); Rows -> {error, {too_many_jobs_returned, gleam@list:map( Rows, fun(Row@1) -> {job_id, erlang:element(2, Row@1)} end )}} end end ) end ). -file("src/m25.gleam", 661). -spec handle_failed_job(pog:connection(), job_id(), failure_reason()) -> {ok, pog:returned(nil)} | {error, pog:transaction_error(job_update_error())}. handle_failed_job(Conn, Job_id, Failure_reason) -> pog:transaction( Conn, fun(Conn@1) -> gleam@result:'try'( begin _pipe = fail_job(Conn@1, Job_id, Failure_reason), gleam@result:map_error( _pipe, fun(Field@0) -> {job_update_query_error, Field@0} end ) end, fun(Crashed_job) -> case erlang:element(3, Crashed_job) of [] -> {error, {job_not_found, Job_id}}; [Row] -> _pipe@1 = retry_jobs_if_needed( Conn@1, [erlang:element(2, Row)] ), gleam@result:map_error( _pipe@1, fun(Field@0) -> {job_update_query_error, Field@0} end ); Rows -> {error, {too_many_jobs_returned, gleam@list:map( Rows, fun(Row@1) -> {job_id, erlang:element(2, Row@1)} end )}} end end ) end ). -file("src/m25.gleam", 1244). -spec do_retry_exponential( integer(), integer(), fun(() -> {ok, AAGA} | {error, AAGB}) ) -> {ok, AAGA} | {error, AAGB}. do_retry_exponential(Attempt, Max_attempts, Func) -> case Func() of {ok, Val} -> {ok, Val}; {error, Error} -> case Attempt > Max_attempts of true -> {error, Error}; false -> Multiplier@1 = case gleam@int:power( 2, erlang:float(Attempt) ) of {ok, Multiplier} -> Multiplier; _assert_fail -> erlang:error(#{gleam_error => let_assert, message => <<"Pattern match failed, no pattern matched the value."/utf8>>, file => <>, module => <<"m25"/utf8>>, function => <<"do_retry_exponential"/utf8>>, line => 1256, value => _assert_fail, start => 37453, 'end' => 37516, pattern_start => 37464, pattern_end => 37478}) end, gleam_erlang_ffi:sleep(1000 * erlang:round(Multiplier@1)), do_retry_exponential(Attempt + 1, Max_attempts, Func) end end. -file("src/m25.gleam", 1240). -spec retry_exponential(integer(), fun(() -> {ok, AAFV} | {error, AAFW})) -> {ok, AAFV} | {error, AAFW}. retry_exponential(Total_attempts, Func) -> do_retry_exponential(1, Total_attempts, Func). -file("src/m25.gleam", 1110). -spec handle_job_executor_message( job_executor_state(AAFM, AAFN, AAFO), job_executor_message(AAFN, AAFO) ) -> gleam@otp@actor:next(job_executor_state(AAFM, AAFN, AAFO), any()). handle_job_executor_message(State, Message) -> case Message of start_work -> Worker_function = fun() -> Message@1 = case (erlang:element(8, State))( erlang:element(9, State) ) of {ok, Output} -> {execution_succeeded, Output}; {error, Error} -> {execution_failed, Error} end, gleam@erlang@process:send(erlang:element(2, State), Message@1) end, Worker_pid = proc_lib:spawn_link(Worker_function), gleam@erlang@process:send_after( erlang:element(2, State), erlang:element(13, erlang:element(4, State)), heartbeat ), gleam@otp@actor:continue( begin _record = State, {job_executor_state, erlang:element(2, _record), erlang:element(3, _record), erlang:element(4, _record), erlang:element(5, _record), {some, Worker_pid}, erlang:element(7, _record), erlang:element(8, _record), erlang:element(9, _record)} end ); heartbeat -> case retry_exponential( 3, fun() -> execute_job_heartbeat( erlang:element(3, State), erlang:element(5, State), erlang:element(14, erlang:element(4, State)), erlang:element(13, erlang:element(4, State)) ) end ) of {ok, Timed_out} -> case erlang:element(3, Timed_out) of [{heartbeat_row, _, true}] -> gleam@option:map( erlang:element(6, State), fun gleam@erlang@process:kill/1 ), gleam@otp@actor:stop(); [{heartbeat_row, false, _}] -> gleam@erlang@process:send_after( erlang:element(2, State), erlang:element(13, erlang:element(4, State)), heartbeat ), gleam@otp@actor:continue(State); [{heartbeat_row, true, _}] -> gleam@option:map( erlang:element(6, State), fun gleam@erlang@process:kill/1 ), case retry_exponential( 3, fun() -> fail_job( erlang:element(3, State), erlang:element(5, State), heartbeat_timeout ) end ) of {ok, _} -> nil; {error, Query_error} -> logging:log( error, <<<<<<"Query error when failing job due to heartbeat timeout for job "/utf8, (youid@uuid:to_string( erlang:element( 2, erlang:element( 5, State ) ) ))/binary>>/binary, ": "/utf8>>/binary, (gleam@string:inspect(Query_error))/binary>> ) end, gleam@otp@actor:stop(); [] -> gleam@otp@actor:continue(State); _ -> erlang:error(#{gleam_error => panic, message => <<"This should never return more than one row!"/utf8>>, file => <>, module => <<"m25"/utf8>>, function => <<"handle_job_executor_message"/utf8>>, line => 1183}) end; {error, Query_error@1} -> logging:log( error, <<<<<<"Query error when checking heartbeat for job "/utf8, (youid@uuid:to_string( erlang:element( 2, erlang:element(5, State) ) ))/binary>>/binary, ": "/utf8>>/binary, (gleam@string:inspect(Query_error@1))/binary>> ), gleam@otp@actor:stop() end; {execution_succeeded, Output@1} -> gleam@erlang@process:send( erlang:element(7, State), {work_succeeded, erlang:element(5, State), Output@1} ), gleam@otp@actor:stop(); {execution_failed, Error@1} -> gleam@erlang@process:send( erlang:element(7, State), {work_failed, erlang:element(5, State), Error@1} ), gleam@otp@actor:stop(); {worker_down, Exit_message} -> case erlang:element(3, Exit_message) of normal -> gleam@otp@actor:stop(); killed -> gleam@otp@actor:stop(); {abnormal, _} -> gleam@erlang@process:send( erlang:element(7, State), {job_worker_down, erlang:element(5, State)} ), gleam@otp@actor:stop() end; {manager_down, _} -> gleam@option:map( erlang:element(6, State), fun gleam@erlang@process:kill/1 ), case retry_exponential( 3, fun() -> handle_failed_job( erlang:element(3, State), erlang:element(5, State), crash ) end ) of {ok, _} -> nil; {error, Query_error@2} -> logging:log( error, <<<<<<"Query error when failing job due to downed manager for job "/utf8, (youid@uuid:to_string( erlang:element( 2, erlang:element(5, State) ) ))/binary>>/binary, ": "/utf8>>/binary, (gleam@string:inspect(Query_error@2))/binary>> ) end, gleam@otp@actor:stop() end. -file("src/m25.gleam", 1062). -spec job_executor_spec( pog:connection(), queue(AAEY, AAEZ, AAFA), job_id(), fun((AAEY) -> {ok, AAEZ} | {error, AAFA}), AAEY, gleam@erlang@process:subject(queue_manager_msg(AAEY, AAEZ, AAFA)) ) -> gleam@otp@actor:builder(job_executor_state(AAEY, AAEZ, AAFA), job_executor_message(AAEZ, AAFA), {gleam@erlang@process:subject(job_executor_message(AAEZ, AAFA)), gleam@erlang@process:pid_()}). job_executor_spec(Conn, Queue, Job_id, Work_func, Input, Manager_subject) -> _pipe@8 = gleam@otp@actor:new_with_initialiser( erlang:element(15, Queue), fun(Self) -> gleam_erlang_ffi:trap_exits(true), gleam@result:'try'( begin _pipe = gleam@erlang@process:subject_owner(Manager_subject), gleam@result:replace_error( _pipe, <<"Failed to get queue manager PID for job ID: "/utf8, (youid@uuid:format( erlang:element(2, Job_id), string ))/binary>> ) end, fun(Manager_pid) -> Manager_monitor = gleam@erlang@process:monitor(Manager_pid), Selector = begin _pipe@1 = gleam_erlang_ffi:new_selector(), _pipe@2 = gleam@erlang@process:select(_pipe@1, Self), _pipe@3 = gleam@erlang@process:select_trapped_exits( _pipe@2, fun(Field@0) -> {worker_down, Field@0} end ), gleam@erlang@process:select_specific_monitor( _pipe@3, Manager_monitor, fun(Field@0) -> {manager_down, Field@0} end ) end, Self_pid@1 = case gleam@erlang@process:subject_owner(Self) of {ok, Self_pid} -> Self_pid; _assert_fail -> erlang:error(#{gleam_error => let_assert, message => <<"Pattern match failed, no pattern matched the value."/utf8>>, file => <>, module => <<"m25"/utf8>>, function => <<"job_executor_spec"/utf8>>, line => 1090, value => _assert_fail, start => 32586, 'end' => 32639, pattern_start => 32597, pattern_end => 32609}) end, _pipe@4 = {job_executor_state, Self, Conn, Queue, Job_id, none, Manager_subject, Work_func, Input}, _pipe@5 = gleam@otp@actor:initialised(_pipe@4), _pipe@6 = gleam@otp@actor:selecting(_pipe@5, Selector), _pipe@7 = gleam@otp@actor:returning( _pipe@6, {Self, Self_pid@1} ), {ok, _pipe@7} end ) end ), gleam@otp@actor:on_message(_pipe@8, fun handle_job_executor_message/2). -file("src/m25.gleam", 733). -spec handle_queue_message( queue_manager_state(AAEH, AAEI, AAEJ), queue_manager_msg(AAEH, AAEI, AAEJ) ) -> gleam@otp@actor:next(queue_manager_state(AAEH, AAEI, AAEJ), queue_manager_msg(AAEH, AAEI, AAEJ)). handle_queue_message(State, Message) -> case Message of process_jobs -> Reserved_jobs_result = pog:transaction( erlang:element(4, State), fun(Conn) -> gleam@result:'try'( begin _pipe = cleanup_stuck_reservations( Conn, erlang:element(2, erlang:element(3, State)), erlang:element(16, erlang:element(3, State)) ), gleam@result:map_error( _pipe, fun(Err) -> {process_jobs_query_error, <<"cleaning up stuck reservations"/utf8>>, Err} end ) end, fun(_) -> gleam@result:'try'( begin _pipe@1 = time_out_jobs( Conn, erlang:element( 2, erlang:element(3, State) ) ), gleam@result:map_error( _pipe@1, fun(Err@1) -> {process_jobs_query_error, <<"timing out jobs"/utf8>>, Err@1} end ) end, fun(Timed_out) -> gleam@result:'try'( begin _pipe@2 = erlang:element( 3, Timed_out ), _pipe@3 = gleam@list:map( _pipe@2, fun(Row) -> erlang:element(2, Row) end ), _pipe@4 = retry_jobs_if_needed( Conn, _pipe@3 ), gleam@result:map_error( _pipe@4, fun(Err@2) -> {process_jobs_query_error, <<"retrying jobs"/utf8>>, Err@2} end ) end, fun(_) -> Limit = gleam@int:max( erlang:element( 3, erlang:element(3, State) ) - m25@internal@bimap:size( erlang:element(5, State) ), 0 ), gleam@bool:guard( Limit =:= 0, {ok, []}, fun() -> gleam@result:'try'( begin _pipe@5 = reserve_jobs( Conn, erlang:element( 3, State ), Limit ), gleam@result:map_error( _pipe@5, fun(Field@0) -> {process_jobs_fetch_error, Field@0} end ) end, fun(Reserved_jobs) -> {ok, Reserved_jobs} end ) end ) end ) end ) end ) end ), case Reserved_jobs_result of {ok, Reserved_jobs@1} -> {Started@1, Start_errors} = begin _pipe@9 = gleam@list:map( Reserved_jobs@1, fun(Job) -> _pipe@6 = job_executor_spec( erlang:element(4, State), erlang:element(3, State), erlang:element(2, Job), erlang:element(10, erlang:element(3, State)), erlang:element(6, Job), erlang:element(2, State) ), _pipe@7 = gleam@otp@actor:start(_pipe@6), _pipe@8 = gleam@result:map( _pipe@7, fun(Started) -> {erlang:element(2, Job), erlang:element(3, Started)} end ), gleam@result:map_error( _pipe@8, fun(Err@3) -> {erlang:element(2, Job), Err@3} end ) end ), gleam@result:partition(_pipe@9) end, Successful_ids = gleam@list:map( Started@1, fun(S) -> erlang:element(1, S) end ), Failed_ids = gleam@list:map( Start_errors, fun(E) -> erlang:element(1, E) end ), Finalize_result = pog:transaction( erlang:element(4, State), fun(Conn@1) -> _pipe@10 = finalise_job_reservations( Conn@1, Successful_ids, Failed_ids ), gleam@result:map_error( _pipe@10, fun(Err@4) -> {process_jobs_query_error, <<"finalizing job reservations"/utf8>>, Err@4} end ) end ), case Finalize_result of {ok, _} -> Running_jobs = gleam@list:fold( Started@1, erlang:element(5, State), fun(Running, Job_data) -> gleam@erlang@process:monitor( erlang:element( 2, erlang:element(2, Job_data) ) ), gleam@erlang@process:send( erlang:element( 1, erlang:element(2, Job_data) ), start_work ), m25@internal@bimap:insert( Running, erlang:element(1, Job_data), erlang:element( 2, erlang:element(2, Job_data) ) ) end ), case Start_errors of [] -> nil; Errors -> Failure_reason = begin _pipe@11 = gleam@list:map( Errors, fun(Error) -> Error_string = case erlang:element( 2, Error ) of {init_exited, Reason} -> <<"Init exited: "/utf8, (gleam@string:inspect( Reason ))/binary>>; {init_failed, Reason@1} -> <<"Init failed: "/utf8, Reason@1/binary>>; init_timeout -> <<"Init timeout"/utf8>> end, <<<<(youid@uuid:to_string( erlang:element( 2, (erlang:element( 1, Error )) ) ))/binary, ": "/utf8>>/binary, Error_string/binary>> end ), gleam@string:join( _pipe@11, <<"\n"/utf8>> ) end, logging:log( warning, <<<<<<"Some actors failed to start (jobs reset to pending) for queue "/utf8, (erlang:element( 2, erlang:element(3, State) ))/binary>>/binary, ": \n"/utf8>>/binary, Failure_reason/binary>> ) end, gleam@erlang@process:send_after( erlang:element(2, State), erlang:element(12, erlang:element(3, State)), process_jobs ), gleam@otp@actor:continue( begin _record = State, {queue_manager_state, erlang:element(2, _record), erlang:element(3, _record), erlang:element(4, _record), Running_jobs} end ); {error, Finalize_error} -> logging:log( error, <<<<<<"Failed to finalize job reservations for queue "/utf8, (erlang:element( 2, erlang:element(3, State) ))/binary>>/binary, ": "/utf8>>/binary, (gleam@string:inspect(Finalize_error))/binary>> ), gleam@list:each( Started@1, fun(Job_data@1) -> gleam@erlang@process:kill( erlang:element( 2, erlang:element(2, Job_data@1) ) ) end ), gleam@erlang@process:send_after( erlang:element(2, State), erlang:element(12, erlang:element(3, State)), process_jobs ), gleam@otp@actor:continue(State) end; {error, Reserve_error} -> case Reserve_error of {transaction_query_error, Query_error} -> logging:log( error, <<<<<<"Failed to reserve jobs for queue "/utf8, (erlang:element( 2, erlang:element(3, State) ))/binary>>/binary, ": "/utf8>>/binary, (gleam@string:inspect(Query_error))/binary>> ); {transaction_rolled_back, {process_jobs_query_error, When, Error@1}} -> logging:log( error, <<<<<<<<<<"Transaction rolled back for queue "/utf8, (erlang:element( 2, erlang:element(3, State) ))/binary>>/binary, " when "/utf8>>/binary, When/binary>>/binary, ": "/utf8>>/binary, (gleam@string:inspect(Error@1))/binary>> ); {transaction_rolled_back, {process_jobs_fetch_error, Fetch_error}} -> case Fetch_error of {job_record_fetch_query_error, Query_error@1} -> logging:log( error, <<<<<<"Query error when reserving jobs for queue "/utf8, (erlang:element( 2, erlang:element(3, State) ))/binary>>/binary, ": "/utf8>>/binary, (gleam@string:inspect(Query_error@1))/binary>> ); {job_record_fetch_decode_errors, Decode_errors} -> logging:log( error, <<"Invalid data for multiple jobs: \n"/utf8, (begin _pipe@12 = gleam@list:map( Decode_errors, fun job_record_decode_error_to_string/1 ), gleam@string:join( _pipe@12, <<"\n"/utf8>> ) end)/binary>> ); no_job_record_found -> erlang:error(#{gleam_error => panic, message => <<"unreachable"/utf8>>, file => <>, module => <<"m25"/utf8>>, function => <<"handle_queue_message"/utf8>>, line => 927}) end end, gleam@erlang@process:send_after( erlang:element(2, State), erlang:element(12, erlang:element(3, State)), process_jobs ), gleam@otp@actor:continue(State) end; {work_succeeded, Job_id, Output} -> case retry_exponential( 3, fun() -> succeed_job( erlang:element(4, State), Job_id, (erlang:element(6, erlang:element(3, State)))(Output) ) end ) of {ok, _} -> nil; {error, Query_error@2} -> logging:log( error, <<<<<<"Query error when succeeding job"/utf8, (youid@uuid:to_string( erlang:element(2, Job_id) ))/binary>>/binary, ": "/utf8>>/binary, (gleam@string:inspect(Query_error@2))/binary>> ) end, Running_jobs@1 = m25@internal@bimap:delete_by_key( erlang:element(5, State), Job_id ), gleam@otp@actor:continue( begin _record@1 = State, {queue_manager_state, erlang:element(2, _record@1), erlang:element(3, _record@1), erlang:element(4, _record@1), Running_jobs@1} end ); {work_failed, Job_id@1, Error@2} -> case retry_exponential( 3, fun() -> handle_errored_job( erlang:element(4, State), erlang:element(3, State), Job_id@1, Error@2 ) end ) of {ok, _} -> nil; {error, Query_error@3} -> logging:log( error, <<<<<<"Query error when failing job due to failed work for job: "/utf8, (youid@uuid:to_string( erlang:element(2, Job_id@1) ))/binary>>/binary, " error: "/utf8>>/binary, (gleam@string:inspect(Query_error@3))/binary>> ) end, Running_jobs@2 = m25@internal@bimap:delete_by_key( erlang:element(5, State), Job_id@1 ), gleam@otp@actor:continue( begin _record@2 = State, {queue_manager_state, erlang:element(2, _record@2), erlang:element(3, _record@2), erlang:element(4, _record@2), Running_jobs@2} end ); {job_executor_down, Down} -> {Pid@1, Reason@3} = case Down of {process_down, _, Pid, Reason@2} -> {Pid, Reason@2}; _assert_fail -> erlang:error(#{gleam_error => let_assert, message => <<"Pattern match failed, no pattern matched the value."/utf8>>, file => <>, module => <<"m25"/utf8>>, function => <<"handle_queue_message"/utf8>>, line => 978, value => _assert_fail, start => 29378, 'end' => 29434, pattern_start => 29389, pattern_end => 29427}) end, gleam@bool:guard( Reason@3 =:= normal, gleam@otp@actor:continue(State), fun() -> case m25@internal@bimap:get_by_value( erlang:element(5, State), Pid@1 ) of {error, _} -> gleam@otp@actor:continue(State); {ok, Job_id@2} -> case retry_exponential( 3, fun() -> handle_failed_job( erlang:element(4, State), Job_id@2, crash ) end ) of {ok, _} -> nil; {error, Query_error@4} -> logging:log( error, <<<<<<"Query error when failing job due to downed executor for job "/utf8, (youid@uuid:to_string( erlang:element( 2, Job_id@2 ) ))/binary>>/binary, ": "/utf8>>/binary, (gleam@string:inspect(Query_error@4))/binary>> ) end, Running_jobs@3 = m25@internal@bimap:delete_by_key( erlang:element(5, State), Job_id@2 ), gleam@otp@actor:continue( begin _record@3 = State, {queue_manager_state, erlang:element(2, _record@3), erlang:element(3, _record@3), erlang:element(4, _record@3), Running_jobs@3} end ) end end ); {job_worker_down, Job_id@3} -> case retry_exponential( 3, fun() -> handle_failed_job(erlang:element(4, State), Job_id@3, crash) end ) of {ok, _} -> nil; {error, Query_error@5} -> logging:log( error, <<<<<<"Query error when failing job due to downed worker for job "/utf8, (youid@uuid:to_string( erlang:element(2, Job_id@3) ))/binary>>/binary, ": "/utf8>>/binary, (gleam@string:inspect(Query_error@5))/binary>> ) end, Running_jobs@4 = m25@internal@bimap:delete_by_key( erlang:element(5, State), Job_id@3 ), gleam@otp@actor:continue( begin _record@4 = State, {queue_manager_state, erlang:element(2, _record@4), erlang:element(3, _record@4), erlang:element(4, _record@4), Running_jobs@4} end ); shutdown -> _pipe@13 = m25@internal@bimap:to_list(erlang:element(5, State)), gleam@list:each( _pipe@13, fun(Job@1) -> gleam@erlang@process:kill(erlang:element(2, Job@1)) end ), gleam@otp@actor:stop() end. -file("src/m25.gleam", 706). -spec queue_manager_spec(queue(AAEA, AAEB, AAEC), pog:connection(), integer()) -> gleam@otp@actor:builder(queue_manager_state(AAEA, AAEB, AAEC), queue_manager_msg(AAEA, AAEB, AAEC), gleam@erlang@process:subject(queue_manager_msg(AAEA, AAEB, AAEC))). queue_manager_spec(Queue, Conn, Queue_init_timeout) -> _pipe@6 = gleam@otp@actor:new_with_initialiser( Queue_init_timeout, fun(Self) -> gleam@erlang@process:send(Self, process_jobs), Selector = begin _pipe = gleam_erlang_ffi:new_selector(), _pipe@1 = gleam@erlang@process:select(_pipe, Self), gleam@erlang@process:select_monitors( _pipe@1, fun(Field@0) -> {job_executor_down, Field@0} end ) end, _pipe@2 = {queue_manager_state, Self, Queue, Conn, m25@internal@bimap:new()}, _pipe@3 = gleam@otp@actor:initialised(_pipe@2), _pipe@4 = gleam@otp@actor:selecting(_pipe@3, Selector), _pipe@5 = gleam@otp@actor:returning(_pipe@4, Self), {ok, _pipe@5} end ), gleam@otp@actor:on_message(_pipe@6, fun handle_queue_message/2). -file("src/m25.gleam", 154). -spec supervisor_spec(m25(), integer()) -> gleam@otp@static_supervisor:builder(). supervisor_spec(M25, Queue_init_timeout) -> Supervisor = gleam@otp@static_supervisor:new(one_for_one), _pipe = erlang:element(3, M25), gleam@list:fold( _pipe, Supervisor, fun(Supervisor@1, Queue) -> gleam@otp@static_supervisor:add( Supervisor@1, begin _pipe@2 = gleam@otp@supervision:worker( fun() -> _pipe@1 = queue_manager_spec( Queue, erlang:element(2, M25), Queue_init_timeout ), gleam@otp@actor:start(_pipe@1) end ), gleam@otp@supervision:restart(_pipe@2, transient) end ) end ). -file("src/m25.gleam", 137). ?DOC( " Start M25 in an unsupervised fashion. This is not recommended. You should prefer\n" " using [`supervised`](#supervised) to start M25 as part of your supervision tree.\n" ). -spec start(m25(), integer()) -> {ok, gleam@otp@actor:started(gleam@otp@static_supervisor:supervisor())} | {error, gleam@otp@actor:start_error()}. start(M25, Queue_init_timeout) -> _pipe = supervisor_spec(M25, Queue_init_timeout), gleam@otp@static_supervisor:start(_pipe). -file("src/m25.gleam", 147). ?DOC( " Create a child spec for the M25 process, allowing it to be run as part of a\n" " supervision tree.\n" ). -spec supervised(m25(), integer()) -> gleam@otp@supervision:child_specification(gleam@otp@static_supervisor:supervisor()). supervised(M25, Queue_init_timeout) -> gleam@otp@supervision:supervisor( fun() -> start(M25, Queue_init_timeout) end ). -file("src/m25.gleam", 1265). -spec duration_to_seconds_float(gleam@time@duration:duration()) -> float(). duration_to_seconds_float(Duration) -> {Seconds, Nanoseconds} = gleam@time@duration:to_seconds_and_nanoseconds( Duration ), erlang:float(Seconds) + (case erlang:float(1000000000) of +0.0 -> +0.0; -0.0 -> -0.0; Gleam@denominator -> erlang:float(Nanoseconds) / Gleam@denominator end). -file("src/m25.gleam", 222). ?DOC(" Enqueue a job to be executed as soon as a worker is available.\n"). -spec enqueue(pog:connection(), queue(AAAT, AAAU, AAAV), job(AAAT)) -> {ok, job_record(AAAT, AAAU, AAAV)} | {error, job_record_fetch_error()}. enqueue(Conn, Queue, Job) -> gleam@result:'try'( begin _pipe@2 = insert_job( Conn, erlang:element(2, Queue), erlang:element(3, Job), (erlang:element(4, Queue))(erlang:element(2, Job)), begin _pipe = gleam@option:unwrap( erlang:element(4, Job), erlang:element(11, Queue) ), duration_to_seconds_float(_pipe) end, 1, erlang:element(5, Job), none, none, begin _pipe@1 = gleam@option:map( erlang:element(6, Job), fun duration_to_seconds_float/1 ), gleam@option:unwrap(_pipe@1, +0.0) end, erlang:element(7, Job) ), gleam@result:map_error( _pipe@2, fun(Field@0) -> {job_record_fetch_query_error, Field@0} end ) end, fun(Job@1) -> case erlang:element(3, Job@1) of [] -> {error, no_job_record_found}; [Row] -> _pipe@3 = decode_job_record_row(Queue, Row), gleam@result:map_error( _pipe@3, fun(Err) -> {job_record_fetch_decode_errors, [Err]} end ); _ -> erlang:error(#{gleam_error => panic, message => <<"unreachable"/utf8>>, file => <>, module => <<"m25"/utf8>>, function => <<"enqueue"/utf8>>, line => 247}) end end ).