Packages
glisten
0.3.2
9.0.1
9.0.0
9.0.0-rc1
8.0.3
8.0.2
8.0.1
8.0.0
8.0.0-rc1
7.0.1
7.0.0
6.0.0
5.0.0
4.0.0
3.0.0
2.0.0
1.0.0
0.11.0
0.10.2
0.10.1
0.10.0
0.9.3
0.9.2
0.9.1
0.9.0
0.8.2
0.8.1
0.8.0
0.7.0
0.6.9
0.6.8
0.6.7
0.6.6
0.6.5
0.6.4
0.6.3
0.6.2
0.6.1
0.6.0
0.5.1
0.5.0
0.4.0
0.3.2
0.3.1
0.3.0
0.2.2
0.2.1
0.2.0
0.1.3
0.1.2
0.1.1
0.1.0
a shiny Gleam TCP/TLS server
Current section
Files
Jump to
Current section
Files
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.