Current section
Files
Jump to
Current section
Files
src/singleflight.erl
-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).