Current section
Files
Jump to
Current section
Files
src/m25.erl
-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 => <<?FILEPATH/utf8>>,
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 => <<?FILEPATH/utf8>>,
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 => <<?FILEPATH/utf8>>,
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 => <<?FILEPATH/utf8>>,
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 => <<?FILEPATH/utf8>>,
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 => <<?FILEPATH/utf8>>,
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 => <<?FILEPATH/utf8>>,
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 => <<?FILEPATH/utf8>>,
module => <<"m25"/utf8>>,
function => <<"enqueue"/utf8>>,
line => 247})
end end
).