-module(singleflight). -compile([no_auto_import, nowarn_unused_vars, nowarn_unused_function, nowarn_nomatch, inline]). -define(FILEPATH, "src/singleflight.gleam"). -export([config/2, fetch/3, start/2]). -export_type([config/0, fetch_error/0, message/2, singleflight/2, state/2]). -type config() :: {config, integer(), integer()}. -type fetch_error() :: crashed | timed_out. -type message(FLZ, FMA) :: {request, FLZ, fun((FLZ) -> FMA), gleam@erlang@process:subject({ok, FMA} | {error, fetch_error()})} | {done, gleam@erlang@process:pid_(), FLZ, {ok, FMA} | {error, fetch_error()}} | {worker_down, gleam@erlang@process:pid_()}. -opaque singleflight(FMB, FMC) :: {singleflight, gleam@erlang@process:subject(message(FMB, FMC)), integer()}. -type state(FMD, FME) :: {state, gleam@dict:dict(FMD, list(gleam@erlang@process:subject({ok, FME} | {error, fetch_error()}))), gleam@dict:dict(gleam@erlang@process:pid_(), FMD), gleam@erlang@process:subject(message(FMD, FME))}. -file("src/singleflight.gleam", 39). -spec config(integer(), integer()) -> config(). config(Initialisation_timeout_ms, Fetch_timeout_ms) -> {config, Initialisation_timeout_ms, Fetch_timeout_ms}. -file("src/singleflight.gleam", 80). -spec fetch(singleflight(FMT, FMU), FMT, fun((FMT) -> FMU)) -> {ok, FMU} | {error, fetch_error()}. fetch(Singleflight, Key, Work) -> {singleflight, Subject, Fetch_timeout_ms} = Singleflight, case gleam@erlang@process:subject_owner(Subject) of {ok, Actor_pid} -> Caller = gleam@erlang@process:new_subject(), Monitor = gleam@erlang@process:monitor(Actor_pid), gleam@erlang@process:send(Subject, {request, Key, Work, Caller}), Result = begin _pipe = gleam_erlang_ffi:new_selector(), _pipe@1 = gleam@erlang@process:select(_pipe, Caller), _pipe@2 = gleam@erlang@process:select_specific_monitor( _pipe@1, Monitor, fun(_) -> {error, crashed} end ), gleam_erlang_ffi:select(_pipe@2, Fetch_timeout_ms) end, gleam@erlang@process:demonitor_process(Monitor), case Result of {ok, Result@1} -> Result@1; {error, nil} -> {error, timed_out} end; {error, nil} -> {error, crashed} end. -file("src/singleflight.gleam", 111). -spec handle_message(state(FMY, FMZ), message(FMY, FMZ)) -> gleam@otp@actor:next(state(FMY, FMZ), message(FMY, FMZ)). handle_message(State, Message) -> case Message of {request, Key, Work, Caller} -> case gleam_stdlib:map_get(erlang:element(2, State), Key) of {ok, Waiters} -> gleam@otp@actor:continue( {state, gleam@dict:insert( erlang:element(2, State), Key, [Caller | Waiters] ), erlang:element(3, State), erlang:element(4, State)} ); {error, nil} -> Self = erlang:element(4, State), Worker_pid = proc_lib:spawn( fun() -> Result = Work(Key), gleam@otp@actor:send( Self, {done, erlang:self(), Key, {ok, Result}} ) end ), gleam@erlang@process:monitor(Worker_pid), gleam@otp@actor:continue( {state, gleam@dict:insert( erlang:element(2, State), Key, [Caller] ), gleam@dict:insert( erlang:element(3, State), Worker_pid, Key ), erlang:element(4, State)} ) end; {worker_down, Pid} -> case gleam_stdlib:map_get(erlang:element(3, State), Pid) of {ok, Key@1} -> case gleam_stdlib:map_get(erlang:element(2, State), Key@1) of {ok, Waiters@1} -> gleam@list:each( Waiters@1, fun(Waiter) -> gleam@erlang@process:send( Waiter, {error, crashed} ) end ); {error, nil} -> nil end, gleam@otp@actor:continue( {state, gleam@dict:delete(erlang:element(2, State), Key@1), gleam@dict:delete(erlang:element(3, State), Pid), erlang:element(4, State)} ); {error, nil} -> gleam@otp@actor:continue(State) end; {done, Pid@1, Key@2, Result@1} -> case gleam_stdlib:map_get(erlang:element(2, State), Key@2) of {ok, Waiters@2} -> gleam@list:each( Waiters@2, fun(Waiter@1) -> gleam@erlang@process:send(Waiter@1, Result@1) end ); {error, nil} -> nil end, gleam@otp@actor:continue( {state, gleam@dict:delete(erlang:element(2, State), Key@2), gleam@dict:delete(erlang:element(3, State), Pid@1), erlang:element(4, State)} ) end. -file("src/singleflight.gleam", 46). -spec start(config(), gleam@erlang@process:name(message(FML, FMM))) -> {ok, gleam@otp@actor:started(singleflight(FML, FMM))} | {error, gleam@otp@actor:start_error()}. start(Config, Name) -> {config, Initialisation_timeout_ms, Fetch_timeout_ms} = Config, _pipe@5 = gleam@otp@actor:new_with_initialiser( Initialisation_timeout_ms, fun(Self) -> Selector = begin _pipe = gleam_erlang_ffi:new_selector(), _pipe@1 = gleam@erlang@process:select(_pipe, Self), gleam@erlang@process:select_monitors( _pipe@1, fun(Down) -> case Down of {process_down, _, Pid, _} -> {worker_down, Pid}; {port_down, _, _, _} -> {worker_down, erlang:self()} end end ) end, _pipe@2 = gleam@otp@actor:initialised( {state, maps:new(), maps:new(), Self} ), _pipe@3 = gleam@otp@actor:selecting(_pipe@2, Selector), _pipe@4 = gleam@otp@actor:returning( _pipe@3, {singleflight, Self, Fetch_timeout_ms} ), {ok, _pipe@4} end ), _pipe@6 = gleam@otp@actor:on_message(_pipe@5, fun handle_message/2), _pipe@7 = gleam@otp@actor:named(_pipe@6, Name), gleam@otp@actor:start(_pipe@7).