-module(gleam@otp@task). -compile(no_auto_import). -export([async/1, try_await/2, await/2, try_await_forever/1, await_forever/1]). -export_type([task/1, await_error/0, message/1]). -opaque task(EQP) :: {task, gleam@otp@process:pid_(), gleam@otp@process:pid_(), gleam@otp@process:receiver(message(EQP))}. -type await_error() :: timeout | {exit, gleam@dynamic:dynamic()}. -type message(EQQ) :: {mon, gleam@otp@process:process_down()} | {chan, EQQ}. -spec async(fun(() -> EQV)) -> task(EQV). async(Work) -> Owner = gleam@otp@process:self(), {Sender, Receiver} = gleam@otp@process:new_channel(), Pid = gleam@otp@process:start( fun() -> gleam@otp@process:send(Sender, Work()) end ), Receiver@1 = begin _pipe = Pid, _pipe@1 = gleam@otp@process:monitor_process(_pipe), _pipe@2 = gleam@otp@process:map_receiver( _pipe@1, fun(A) -> {mon, A} end ), gleam@otp@process:merge_receiver( _pipe@2, begin _pipe@3 = Receiver, gleam@otp@process:map_receiver(_pipe@3, fun(A) -> {chan, A} end) end ) end, {task, Owner, Pid, Receiver@1}. -spec assert_owner(task(any())) -> nil. assert_owner(Task) -> Self = gleam@otp@process:self(), case erlang:element(2, Task) =:= Self of true -> nil; false -> gleam@otp@process:send_exit( Self, <<"awaited on a task that does not belong to this process"/utf8>> ) end. -spec try_await(task(EQZ), integer()) -> {ok, EQZ} | {error, await_error()}. try_await(Task, Timeout) -> assert_owner(Task), case gleam@otp@process:'receive'(erlang:element(4, Task), Timeout) of {ok, {chan, X}} -> gleam@otp@process:close_channels(erlang:element(4, Task)), {ok, X}; {ok, {mon, {process_down, _@1, Reason}}} -> gleam@otp@process:close_channels(erlang:element(4, Task)), {error, {exit, Reason}}; {error, nil} -> {error, timeout} end. -spec await(task(ERD), integer()) -> ERD. await(Task, Timeout) -> {ok, Value@1} = case try_await(Task, Timeout) of {ok, Value} -> {ok, Value}; _try -> erlang:error(#{gleam_error => assert, message => <<"Assertion pattern match failed"/utf8>>, value => _try, module => <<"gleam/otp/task"/utf8>>, function => <<"await"/utf8>>, line => 116}) end, Value@1. -spec try_await_forever(task(ERF)) -> {ok, ERF} | {error, await_error()}. try_await_forever(Task) -> assert_owner(Task), case gleam@otp@process:receive_forever(erlang:element(4, Task)) of {chan, X} -> gleam@otp@process:close_channels(erlang:element(4, Task)), {ok, X}; {mon, {process_down, _@1, Reason}} -> gleam@otp@process:close_channels(erlang:element(4, Task)), {error, {exit, Reason}} end. -spec await_forever(task(ERJ)) -> ERJ. await_forever(Task) -> {ok, Value@1} = case try_await_forever(Task) of {ok, Value} -> {ok, Value}; _try -> erlang:error(#{gleam_error => assert, message => <<"Assertion pattern match failed"/utf8>>, value => _try, module => <<"gleam/otp/task"/utf8>>, function => <<"await_forever"/utf8>>, line => 151}) end, Value@1.