Packages
glisten
0.2.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([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.