-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} ).