-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, radish@resp:value()} | {error, radish@error:error()}), integer()}. -spec 'receive'( mug:socket(), gleam@erlang@process:selector({ok, bitstring()} | {error, mug:error()}), bitstring(), {integer(), integer(), integer()}, integer() ) -> {ok, 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} -> Now = erlang:now(), case timer:now_diff(Now, Start_time) >= (Timeout * 1000) 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 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(), Start_time = erlang:now(), case 'receive'(Socket, Selector, <<>>, Start_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; 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@result:then( begin _pipe = radish@tcp:connect(Host, Port, Timeout), gleam@result:replace_error( _pipe, {init_failed, {abnormal, <<"Unable to connect to Redis server"/utf8>>}} ) end, fun(Socket) -> gleam@result:then( gleam@otp@actor:start(Socket, fun handle_message/2), fun(Client) -> Client_pid = gleam@erlang@process:subject_owner(Client), gen_tcp:controlling_process(Socket, Client_pid), {ok, Client} end ) end ).