Current section

Files

Jump to
glisten src glisten@tcp.erl
Raw

src/glisten@tcp.erl

-module(glisten@tcp).
-compile(no_auto_import).
-export([accept_timeout/2, accept/1, receive_timeout/3, 'receive'/2, send/2, socket_info/1, close/1, do_shutdown/2, shutdown/1, set_opts/2, listen/2, echo_loop/2, start_handler/3, start_acceptor/3, receiver_to_iterator/1, start_acceptor_pool/4, handler/1]).
-export_type([socket_mode/0, active_state/0, tcp_option/0, socket_reason/0, listen_socket/0, socket/0, acceptor/0, acceptor_error/0, handler_message/0, acceptor_state/0, loop_state/1]).
-type socket_mode() :: binary.
-type active_state() :: once | passive | {count, integer()} | active.
-type tcp_option() :: {backlog, integer()} |
{nodelay, boolean()} |
{linger, {boolean(), integer()}} |
{send_timeout, integer()} |
{send_timeout_close, boolean()} |
{reuseaddr, boolean()} |
{active_mode, active_state()} |
{mode, socket_mode()}.
-type socket_reason() :: closed | timeout.
-opaque listen_socket() :: listen_socket.
-opaque socket() :: socket.
-type acceptor() :: {accept_connection, listen_socket()}.
-type acceptor_error() :: accept_error | handler_error | control_error.
-type handler_message() :: close |
ready |
{receive_message, bitstring()} |
{tcp, gleam@otp@port:port_(), bitstring()} |
{tcp_closed, nil}.
-type acceptor_state() :: {acceptor_state,
gleam@otp@process:sender(acceptor()),
gleam@option:option(socket())}.
-type loop_state(GMK) :: {loop_state,
socket(),
gleam@otp@process:sender(handler_message()),
GMK}.
-spec accept_timeout(listen_socket(), integer()) -> {ok, socket()} |
{error, socket_reason()}.
accept_timeout(A, B) ->
gen_tcp:accept(A, B).
-spec accept(listen_socket()) -> {ok, socket()} | {error, socket_reason()}.
accept(A) ->
gen_tcp:accept(A).
-spec receive_timeout(socket(), integer(), integer()) -> {ok, bitstring()} |
{error, socket_reason()}.
receive_timeout(A, B, C) ->
gen_tcp:recv(A, B, C).
-spec 'receive'(socket(), integer()) -> {ok, bitstring()} |
{error, socket_reason()}.
'receive'(A, B) ->
gen_tcp:recv(A, B).
-spec send(socket(), gleam@bit_builder:bit_builder()) -> {ok, nil} |
{error, socket_reason()}.
send(A, B) ->
tcp_ffi:send(A, B).
-spec socket_info(socket()) -> gleam@map:map_(any(), any()).
socket_info(A) ->
socket:info(A).
-spec close(socket()) -> gleam@erlang@atom:atom_().
close(A) ->
gen_tcp:close(A).
-spec do_shutdown(socket(), gleam@erlang@atom:atom_()) -> nil.
do_shutdown(A, B) ->
gen_tcp:shutdown(A, B).
-spec shutdown(socket()) -> nil.
shutdown(Socket) ->
{ok, Write@1} = case gleam@erlang@atom:from_string(<<"write"/utf8>>) of
{ok, Write} -> {ok, Write};
_try ->
erlang:error(#{gleam_error => assert,
message => <<"Assertion pattern match failed"/utf8>>,
value => _try,
module => <<"glisten/tcp"/utf8>>,
function => <<"shutdown"/utf8>>,
line => 103})
end,
gen_tcp:shutdown(Socket, Write@1).
-spec set_opts(socket(), list(tcp_option())) -> {ok, nil} | {error, nil}.
set_opts(Socket, Opts) ->
_pipe = Opts,
_pipe@1 = opts_to_map(_pipe),
_pipe@2 = gleam@map:to_list(_pipe@1),
_pipe@3 = gleam@list:map(_pipe@2, fun gleam@dynamic:from/1),
tcp_ffi:set_opts(Socket, _pipe@3).
-spec opts_to_map(list(tcp_option())) -> gleam@map:map_(gleam@erlang@atom:atom_(), gleam@dynamic:dynamic()).
opts_to_map(Options) ->
Opt_decoder = gleam@dynamic:tuple2(
fun gleam@dynamic:dynamic/1,
fun gleam@dynamic:dynamic/1
),
_pipe = Options,
_pipe@1 = gleam@list:map(_pipe, fun(Opt) -> case Opt of
{active_mode, passive} ->
gleam@dynamic:from(
{gleam@erlang@atom:create_from_string(<<"active"/utf8>>),
false}
);
{active_mode, active} ->
gleam@dynamic:from(
{gleam@erlang@atom:create_from_string(<<"active"/utf8>>),
true}
);
{active_mode, {count, N}} ->
gleam@dynamic:from(
{gleam@erlang@atom:create_from_string(<<"active"/utf8>>),
N}
);
{active_mode, once} ->
gleam@dynamic:from(
{gleam@erlang@atom:create_from_string(<<"active"/utf8>>),
gleam@erlang@atom:create_from_string(<<"once"/utf8>>)}
);
Other ->
gleam@dynamic:from(Other)
end end),
_pipe@2 = gleam@list:filter_map(_pipe@1, Opt_decoder),
_pipe@3 = gleam@list:map(
_pipe@2,
fun(_capture) ->
gleam@pair:map_first(_capture, fun gleam@dynamic:unsafe_coerce/1)
end
),
gleam@map:from_list(_pipe@3).
-spec merge_with_default_options(list(tcp_option())) -> list(tcp_option()).
merge_with_default_options(Options) ->
Overrides = opts_to_map(Options),
_pipe = [{backlog, 1024},
{nodelay, true},
{linger, {true, 30}},
{send_timeout, 30000},
{send_timeout_close, true},
{reuseaddr, true},
{mode, binary},
{active_mode, passive}],
_pipe@1 = opts_to_map(_pipe),
_pipe@2 = gleam@map:merge(_pipe@1, Overrides),
_pipe@3 = gleam@map:to_list(_pipe@2),
_pipe@4 = gleam@list:map(_pipe@3, fun gleam@dynamic:from/1),
gleam@list:map(_pipe@4, fun gleam@dynamic:unsafe_coerce/1).
-spec listen(integer(), list(tcp_option())) -> {ok, listen_socket()} |
{error, socket_reason()}.
listen(Port, Options) ->
_pipe = Options,
_pipe@1 = merge_with_default_options(_pipe),
gen_tcp:listen(Port, _pipe@1).
-spec echo_loop(handler_message(), acceptor_state()) -> gleam@otp@actor:next(acceptor_state()).
echo_loop(Msg, State) ->
case {Msg, State} of
{{receive_message, Data}, {acceptor_state, _@1, {some, Sock}}} ->
_@2 = tcp_ffi:send(Sock, gleam@bit_builder:from_bit_string(Data)),
nil;
{_@3, _@4} ->
nil
end,
{continue, State}.
-spec start_handler(
socket(),
GOI,
fun((handler_message(), loop_state(GOI)) -> gleam@otp@actor:next(loop_state(GOI)))
) -> {ok, gleam@otp@process:sender(handler_message())} |
{error, gleam@otp@actor:start_error()}.
start_handler(Socket, Initial_data, Loop) ->
gleam@otp@actor:start_spec(
{spec,
fun() ->
{Sender, Receiver} = gleam@otp@process:new_channel(),
Socket_receiver = begin
_pipe = gleam@otp@process:bare_message_receiver(),
_pipe@1 = gleam@otp@process:map_receiver(
_pipe,
fun(Msg) -> case gleam@dynamic:unsafe_coerce(Msg) of
{tcp, _@1, Data} ->
{receive_message, Data};
Message ->
Message
end end
),
gleam@otp@process:merge_receiver(_pipe@1, Receiver)
end,
{ready,
{loop_state, Socket, Sender, Initial_data},
{some, Socket_receiver}}
end,
1000,
fun(Msg@1, State) -> case Msg@1 of
{tcp_closed, _@2} ->
{stop, normal};
ready ->
{ok, _@4} = case set_opts(Socket, [{active_mode, once}]) of
{ok, _@3} -> {ok, _@3};
_try ->
erlang:error(#{gleam_error => assert,
message => <<"Assertion pattern match failed"/utf8>>,
value => _try,
module => <<"glisten/tcp"/utf8>>,
function => <<"start_handler"/utf8>>,
line => 251})
end,
{continue, State};
Msg@2 ->
case Loop(Msg@2, State) of
{continue, Next_state} ->
{ok, nil} = case set_opts(
Socket,
[{active_mode, once}]
) of
{ok, nil} -> {ok, nil};
_try@1 ->
erlang:error(#{gleam_error => assert,
message => <<"Assertion pattern match failed"/utf8>>,
value => _try@1,
module => <<"glisten/tcp"/utf8>>,
function => <<"start_handler"/utf8>>,
line => 257})
end,
{continue, Next_state};
Msg@3 ->
Msg@3
end
end end}
).
-spec start_acceptor(
listen_socket(),
GON,
fun((handler_message(), loop_state(GON)) -> gleam@otp@actor:next(loop_state(GON)))
) -> {ok, gleam@otp@process:sender(acceptor())} |
{error, gleam@otp@actor:start_error()}.
start_acceptor(Socket, Initial_data, Loop_fn) ->
gleam@otp@actor:start_spec(
{spec,
fun() ->
{Sender, Actor_receiver} = gleam@otp@process:new_channel(),
gleam@otp@process:send(Sender, {accept_connection, Socket}),
{ready, {acceptor_state, Sender, none}, {some, Actor_receiver}}
end,
1000,
fun(Msg, State) ->
{acceptor_state, Sender@1, _@1} = State,
case Msg of
{accept_connection, Listener} ->
Res = case begin
_pipe = gen_tcp:accept(Listener),
gleam@result:replace_error(_pipe, accept_error)
end of
{error, _try} -> {error, _try};
{ok, Sock} ->
case begin
_pipe@1 = start_handler(
Sock,
Initial_data,
Loop_fn
),
gleam@result:replace_error(
_pipe@1,
handler_error
)
end of
{error, _try@1} -> {error, _try@1};
{ok, Start} ->
_pipe@2 = Sock,
_pipe@3 = tcp_ffi:controlling_process(
_pipe@2,
gleam@otp@process:pid(Start)
),
_pipe@4 = gleam@result:replace_error(
_pipe@3,
control_error
),
gleam@result:map(
_pipe@4,
fun(_) ->
gleam@otp@process:send(
Start,
ready
)
end
)
end
end,
case Res of
{error, Reason} ->
{stop, {abnormal, gleam@dynamic:from(Reason)}};
_@2 ->
gleam@otp@actor:send(
Sender@1,
{accept_connection, Listener}
),
{continue, State}
end;
Msg@1 ->
{stop, {abnormal, gleam@dynamic:from(Msg@1)}}
end
end}
).
-spec receiver_to_iterator(gleam@otp@process:receiver(GOS)) -> gleam@iterator:iterator(GOS).
receiver_to_iterator(Receiver) ->
gleam@iterator:unfold(
Receiver,
fun(Recv) ->
_pipe = Recv,
_pipe@1 = gleam@otp@process:receive_forever(_pipe),
{next, _pipe@1, Recv}
end
).
-spec start_acceptor_pool(
listen_socket(),
fun((handler_message(), loop_state(GOV)) -> gleam@otp@actor:next(loop_state(GOV))),
GOV,
integer()
) -> {ok, gleam@otp@process:sender(gleam@otp@supervisor:message())} |
{error, gleam@otp@actor:start_error()}.
start_acceptor_pool(Listener_socket, Handler, Initial_data, Pool_count) ->
gleam@otp@supervisor:start_spec(
{spec,
nil,
100,
1,
fun(Children) ->
_pipe = gleam@iterator:range(0, Pool_count),
gleam@iterator:fold(
_pipe,
Children,
fun(Children@1, _) ->
gleam@otp@supervisor:add(
Children@1,
gleam@otp@supervisor:worker(
fun(_) ->
start_acceptor(
Listener_socket,
Initial_data,
Handler
)
end
)
)
end
)
end}
).
-spec handler(
fun((bitstring(), loop_state(GPA)) -> gleam@otp@actor:next(loop_state(GPA)))
) -> fun((handler_message(), loop_state(GPA)) -> gleam@otp@actor:next(loop_state(GPA))).
handler(Func) ->
fun(Msg, State) -> case Msg of
{tcp, _@1, _@2} ->
gleam@io:debug(
{<<"Received an unexpected TCP message"/utf8>>, Msg}
),
{continue, State};
ready ->
gleam@io:debug(
{<<"Received an unexpected TCP message"/utf8>>, Msg}
),
{continue, State};
close ->
gen_tcp:close(erlang:element(2, State)),
{stop, normal};
{tcp_closed, _@1} ->
{stop, normal};
{receive_message, Data} ->
Func(Data, State)
end end.