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]).
-export([start/3]).
-export_type([message/0]).
-type message() :: shutdown |
{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()}.
-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.
-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.
-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;
shutdown ->
_ = mug_ffi:shutdown(Socket),
{stop, normal}
end.
-spec start(binary(), integer(), integer()) -> {ok,
gleam@erlang@process:subject(message())} |
{error, gleam@otp@actor:start_error()}.
start(Host, Port, Timeout) ->
gleam@otp@actor:start_spec(
{spec, fun() -> case radish@tcp:connect(Host, Port, Timeout) of
{ok, Socket} ->
{ready, Socket, gleam_erlang_ffi:new_selector()};
{error, _} ->
{failed, <<"Unable to connect to Redis server"/utf8>>}
end end, Timeout, fun handle_message/2}
).