Current section
Files
Jump to
Current section
Files
src/radish@client.erl
-module(radish@client).
-compile([no_auto_import, nowarn_unused_vars, nowarn_unused_function, nowarn_nomatch]).
-export([start/4]).
-export_type([message/0]).
-type message() :: {command,
bitstring(),
gleam@erlang@process:subject({ok, list(radish@resp:value())} |
{error, radish@error:error()}),
integer()} |
{blocking_command,
bitstring(),
gleam@erlang@process:subject({ok, list(radish@resp:value())} |
{error, radish@error:error()}),
integer()} |
{receive_forever,
gleam@erlang@process:subject({ok, list(radish@resp:value())} |
{error, radish@error:error()}),
integer()}.
-file("/home/massivefermion/Desktop/radish/src/radish/client.gleam", 128).
-spec 'receive'(
mug:socket(),
gleam@erlang@process:selector({ok, bitstring()} | {error, mug:error()}),
bitstring(),
integer(),
integer()
) -> {ok, list(radish@resp:value())} | {error, radish@error:error()}.
'receive'(Socket, Selector, Storage, Start_time, Timeout) ->
case radish@decoder:decode(Storage) of
{ok, Value} ->
{ok, Value};
{error, Error} ->
case (erlang:monotonic_time() - Start_time) >= (Timeout * 1000000) of
true ->
{error, Error};
false ->
case radish@tcp:'receive'(Socket, Selector, Timeout) of
{error, Tcp_error} ->
{error, {tcp_error, Tcp_error}};
{ok, Packet} ->
'receive'(
Socket,
Selector,
gleam@bit_array:append(Storage, Packet),
Start_time,
Timeout
)
end
end
end.
-file("/home/massivefermion/Desktop/radish/src/radish/client.gleam", 157).
-spec receive_forever(
mug:socket(),
gleam@erlang@process:selector({ok, bitstring()} | {error, mug:error()}),
bitstring(),
integer(),
integer()
) -> {ok, list(radish@resp:value())} | {error, radish@error:error()}.
receive_forever(Socket, Selector, Storage, Start_time, Timeout) ->
case radish@decoder:decode(Storage) of
{ok, Value} ->
{ok, Value};
{error, Error} when Timeout =/= 0 ->
case (erlang:monotonic_time() - Start_time) >= (Timeout * 1000000) of
true ->
{error, Error};
false ->
case radish@tcp:receive_forever(Socket, Selector) of
{error, Tcp_error} ->
{error, {tcp_error, Tcp_error}};
{ok, Packet} ->
receive_forever(
Socket,
Selector,
gleam@bit_array:append(Storage, Packet),
Start_time,
Timeout
)
end
end;
{error, _} ->
case radish@tcp:receive_forever(Socket, Selector) of
{error, Tcp_error@1} ->
{error, {tcp_error, Tcp_error@1}};
{ok, Packet@1} ->
receive_forever(
Socket,
Selector,
gleam@bit_array:append(Storage, Packet@1),
Start_time,
Timeout
)
end
end.
-file("/home/massivefermion/Desktop/radish/src/radish/client.gleam", 54).
-spec handle_message(message(), mug:socket()) -> gleam@otp@actor:next(any(), mug:socket()).
handle_message(Msg, Socket) ->
case Msg of
{command, Cmd, Reply_with, Timeout} ->
case radish@tcp:send(Socket, Cmd) of
{ok, nil} ->
Selector = radish@tcp:new_selector(),
case 'receive'(
Socket,
Selector,
<<>>,
erlang:monotonic_time(),
Timeout
) of
{ok, Reply} ->
gleam@otp@actor:send(Reply_with, {ok, Reply}),
gleam@otp@actor:continue(Socket);
{error, Error} ->
_ = mug_ffi:shutdown(Socket),
gleam@otp@actor:send(Reply_with, {error, Error}),
{stop, {abnormal, <<"TCP Error"/utf8>>}}
end;
{error, Error@1} ->
_ = mug_ffi:shutdown(Socket),
gleam@otp@actor:send(
Reply_with,
{error, {tcp_error, Error@1}}
),
{stop, {abnormal, <<"TCP Error"/utf8>>}}
end;
{blocking_command, Cmd@1, Reply_with@1, Timeout@1} ->
case radish@tcp:send(Socket, Cmd@1) of
{ok, nil} ->
Selector@1 = radish@tcp:new_selector(),
case receive_forever(
Socket,
Selector@1,
<<>>,
erlang:monotonic_time(),
Timeout@1
) of
{ok, Reply@1} ->
gleam@otp@actor:send(Reply_with@1, {ok, Reply@1}),
gleam@otp@actor:continue(Socket);
{error, Error@2} ->
_ = mug_ffi:shutdown(Socket),
gleam@otp@actor:send(Reply_with@1, {error, Error@2}),
{stop, {abnormal, <<"TCP Error"/utf8>>}}
end;
{error, Error@3} ->
_ = mug_ffi:shutdown(Socket),
gleam@otp@actor:send(
Reply_with@1,
{error, {tcp_error, Error@3}}
),
{stop, {abnormal, <<"TCP Error"/utf8>>}}
end;
{receive_forever, Reply_with@2, Timeout@2} ->
Selector@2 = radish@tcp:new_selector(),
case receive_forever(
Socket,
Selector@2,
<<>>,
erlang:monotonic_time(),
Timeout@2
) of
{ok, Reply@2} ->
gleam@otp@actor:send(Reply_with@2, {ok, Reply@2}),
gleam@otp@actor:continue(Socket);
{error, Error@4} ->
_ = mug_ffi:shutdown(Socket),
gleam@otp@actor:send(Reply_with@2, {error, Error@4}),
{stop, {abnormal, <<"TCP Error"/utf8>>}}
end
end.
-file("/home/massivefermion/Desktop/radish/src/radish/client.gleam", 37).
-spec worker_spec(binary(), integer(), integer()) -> lifeguard:spec(mug:socket(), message()).
worker_spec(Host, Port, Timeout) ->
{spec, fun(Selector) -> case radish@tcp:connect(Host, Port, Timeout) of
{ok, Socket} ->
{ready, Socket, Selector};
{error, _} ->
{failed, <<"Unable to connect to Redis server"/utf8>>}
end end, Timeout, fun handle_message/2}.
-file("/home/massivefermion/Desktop/radish/src/radish/client.gleam", 26).
-spec start(binary(), integer(), integer(), integer()) -> {ok,
lifeguard:pool(message())} |
{error, lifeguard:start_error()}.
start(Host, Port, Timeout, Pool_size) ->
_pipe = lifeguard:new(worker_spec(Host, Port, Timeout)),
_pipe@1 = lifeguard:with_size(_pipe, Pool_size),
lifeguard:start(_pipe@1, Timeout).