Current section

Files

Jump to
glisten src glisten@tcp.erl
Raw

src/glisten@tcp.erl

-module(glisten@tcp).
-compile(no_auto_import).
-export([do_listen_tcp/2, accept_timeout/2, accept/1, do_receive/2, send/2, socket_info/1, close/1, do_shutdown/2, shutdown/1, set_opts/2, merge_with_default_options/1, listen/2, 'receive'/1, echo_loop/2, start_handler/3, start_acceptor/3, receive_timeout/2, receiver_to_iterator/1, start_acceptor_pool/4, handler/1]).
-export_type([tcp_option/0, socket_reason/0, listen_socket/0, socket/0, acceptor/0, acceptor_error/0, handler_message/0, acceptor_state/0]).
-type tcp_option() :: {backlog, integer()} |
{nodelay, boolean()} |
{linger, {boolean(), integer()}} |
{send_timeout, integer()} |
{send_timeout_close, boolean()} |
{reuseaddr, boolean()} |
{active, gleam@dynamic:dynamic()} |
binary.
-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() :: ready |
{receive_message, bitstring()} |
{tcp, gleam@otp@port:port_(), gleam@bit_builder:bit_builder()} |
{tcp_closed, nil}.
-type acceptor_state() :: {acceptor_state,
gleam@otp@process:sender(acceptor()),
gleam@option:option(socket())}.
-spec do_listen_tcp(integer(), list(tcp_option())) -> {ok, listen_socket()} |
{error, socket_reason()}.
do_listen_tcp(A, B) ->
gen_tcp:listen(A, B).
-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 do_receive(socket(), integer()) -> {ok, bitstring()} |
{error, socket_reason()}.
do_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 => 90})
end,
gen_tcp:shutdown(Socket, Write@1).
-spec set_opts(socket(), list(tcp_option())) -> {ok, nil} | {error, nil}.
set_opts(A, B) ->
tcp_ffi:set_opts(A, B).
-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 gleam@dynamic:from/1),
_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},
binary,
{active, gleam@dynamic:from(false)}],
_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 'receive'(socket()) -> {ok, bitstring()} | {error, socket_reason()}.
'receive'(Socket) ->
gen_tcp:recv(Socket, 0).
-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(),
GOA,
fun((handler_message(), {socket(), GOA}) -> gleam@otp@actor:next({socket(),
GOA}))
) -> {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() ->
{_@1, Receiver} = gleam@otp@process:new_channel(),
Socket_receiver = begin
_pipe = gleam@otp@process:bare_message_receiver(),
_pipe@3 = gleam@otp@process:map_receiver(
_pipe,
fun(Msg) -> case gleam@dynamic:unsafe_coerce(Msg) of
{tcp, _@2, Data} ->
_pipe@1 = Data,
_pipe@2 = gleam@bit_builder:to_bit_string(
_pipe@1
),
{receive_message, _pipe@2};
Message ->
Message
end end
),
gleam@otp@process:merge_receiver(_pipe@3, Receiver)
end,
{ready, {Socket, Initial_data}, {some, Socket_receiver}}
end,
1000,
fun(Msg@1, State) ->
{Socket@1, _@3} = State,
case Msg@1 of
{tcp_closed, _@4} ->
{stop, normal};
ready ->
{ok, _@6} = case tcp_ffi:set_opts(
Socket@1,
[{active,
gleam@dynamic:from(
gleam@erlang@atom:create_from_string(
<<"once"/utf8>>
)
)}]
) of
{ok, _@5} -> {ok, _@5};
_try ->
erlang:error(#{gleam_error => assert,
message => <<"Assertion pattern match failed"/utf8>>,
value => _try,
module => <<"glisten/tcp"/utf8>>,
function => <<"start_handler"/utf8>>,
line => 216})
end,
{continue, State};
Msg@2 ->
case Loop(Msg@2, State) of
{continue, Next_state} ->
{ok, nil} = case tcp_ffi:set_opts(
Socket@1,
[{active, gleam@dynamic:from(<<"once"/utf8>>)}]
) 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 => 226})
end,
{continue, Next_state};
Msg@3 ->
Msg@3
end
end
end}
).
-spec start_acceptor(
listen_socket(),
GOF,
fun((handler_message(), {socket(), GOF}) -> gleam@otp@actor:next({socket(),
GOF}))
) -> {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 receive_timeout(socket(), integer()) -> {ok, bitstring()} |
{error, socket_reason()}.
receive_timeout(Socket, Timeout) ->
gen_tcp:recv(Socket, 0, Timeout).
-spec receiver_to_iterator(gleam@otp@process:receiver(GOM)) -> gleam@iterator:iterator(GOM).
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(), {socket(), GOP}) -> gleam@otp@actor:next({socket(),
GOP})),
GOP,
integer()
) -> {ok, nil} | {error, gleam@otp@actor:start_error()}.
start_acceptor_pool(Listener_socket, Handler, Initial_data, Pool_count) ->
_pipe@1 = 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}
),
gleam@result:replace(_pipe@1, nil).
-spec handler(
fun((bitstring(), {socket(), GOT}) -> gleam@otp@actor:next({socket(), GOT}))
) -> fun((handler_message(), {socket(), GOT}) -> gleam@otp@actor:next({socket(),
GOT})).
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};
{tcp_closed, _@1} ->
{stop, normal};
{receive_message, Data} ->
Func(Data, State)
end end.