-module(gleam@otp@task). -compile(no_auto_import). -export([async/1, try_await/2, await/2]). -export_type([task/1, await_error/0, message/1]). -opaque task(FR) :: {task, gleam@otp@process:pid_(), gleam@otp@process:pid_(), gleam@otp@process:receiver(message(FR))}. -type await_error() :: timeout | {exit, gleam@dynamic:dynamic()}. -type message(FS) :: {mon, gleam@otp@process:process_down()} | {chan, FS}. -spec async(fun(() -> FX)) -> task(FX). 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(GB), integer()) -> {ok, GB} | {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(GF), integer()) -> GF. 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.