Packages
🐑 a fluffy Gleam web server
Security advisory:
This version has known vulnerabilities.
View advisories
Current section
Files
Jump to
Current section
Files
src/ewe@internal@websocket.erl
-module(ewe@internal@websocket).
-compile([no_auto_import, nowarn_unused_vars, nowarn_unused_function, nowarn_nomatch, inline]).
-define(FILEPATH, "src/ewe/internal/websocket.gleam").
-export([send_frame/5, start/6]).
-export_type([exit_reason/0, next/1, state/1, websocket_connection/0, handler_message/1, valid_message/0, message/1]).
-if(?OTP_RELEASE >= 27).
-define(MODULEDOC(Str), -moduledoc(Str)).
-define(DOC(Str), -doc(Str)).
-else.
-define(MODULEDOC(Str), -compile([])).
-define(DOC(Str), -compile([])).
-endif.
?MODULEDOC(false).
-type exit_reason() :: normal | {abnormal, binary()}.
-type next(KOJ) :: {continue, KOJ} | {stop, exit_reason()}.
-type state(KOK) :: {state,
KOK,
bitstring(),
gleam@option:option(gramps@websocket@compression:compression())}.
-type websocket_connection() :: {websocket_connection,
glisten@transport:transport(),
glisten@socket:socket(),
gleam@option:option(gramps@websocket@compression:context())}.
-type handler_message(KOL) :: {frame, gramps@websocket:frame()} |
{user_message, KOL}.
-type valid_message() :: {packet, bitstring()} | close.
-type message(KOM) :: {valid, valid_message()} | {user, KOM} | invalid.
-file("src/ewe/internal/websocket.gleam", 61).
?DOC(false).
-spec select_valid_record(gleam@erlang@process:selector(message(KON)), binary()) -> gleam@erlang@process:selector(message(KON)).
select_valid_record(Selector, Binary_atom) ->
gleam@erlang@process:select_record(
Selector,
erlang:binary_to_atom(Binary_atom),
2,
fun(Record) ->
_pipe = gleam@dynamic@decode:run(
Record,
begin
gleam@dynamic@decode:field(
2,
{decoder, fun gleam@dynamic@decode:decode_bit_array/1},
fun(Data) ->
gleam@dynamic@decode:success(
{valid, {packet, Data}}
)
end
)
end
),
gleam@result:unwrap(_pipe, invalid)
end
).
-file("src/ewe/internal/websocket.gleam", 74).
?DOC(false).
-spec glisten_selector() -> gleam@erlang@process:selector(message(any())).
glisten_selector() ->
_pipe = gleam_erlang_ffi:new_selector(),
_pipe@1 = select_valid_record(_pipe, <<"tcp"/utf8>>),
_pipe@2 = select_valid_record(_pipe@1, <<"ssl"/utf8>>),
_pipe@3 = gleam@erlang@process:select_record(
_pipe@2,
erlang:binary_to_atom(<<"tcp_closed"/utf8>>),
1,
fun(_) -> {valid, close} end
),
gleam@erlang@process:select_record(
_pipe@3,
erlang:binary_to_atom(<<"ssl_closed"/utf8>>),
1,
fun(_) -> {valid, close} end
).
-file("src/ewe/internal/websocket.gleam", 88).
?DOC(false).
-spec get_deflate(
gleam@option:option(gramps@websocket@compression:compression())
) -> gleam@option:option(gramps@websocket@compression:context()).
get_deflate(Compression) ->
gleam@option:map(
Compression,
fun(Compression@1) -> erlang:element(3, Compression@1) end
).
-file("src/ewe/internal/websocket.gleam", 94).
?DOC(false).
-spec get_inflate(
gleam@option:option(gramps@websocket@compression:compression())
) -> gleam@option:option(gramps@websocket@compression:context()).
get_inflate(Compression) ->
gleam@option:map(
Compression,
fun(Compression@1) -> erlang:element(2, Compression@1) end
).
-file("src/ewe/internal/websocket.gleam", 100).
?DOC(false).
-spec close_compression(
gleam@option:option(gramps@websocket@compression:compression())
) -> gleam@option:option(nil).
close_compression(Compression) ->
gleam@option:map(
Compression,
fun(Compression@1) ->
gramps@websocket@compression:close(erlang:element(3, Compression@1)),
gramps@websocket@compression:close(erlang:element(2, Compression@1))
end
).
-file("src/ewe/internal/websocket.gleam", 206).
?DOC(false).
-spec loop_by_frames(
list(gramps@websocket:frame()),
websocket_connection(),
fun((websocket_connection(), KPT, handler_message(any())) -> next(KPT)),
next(KPT)
) -> next(KPT).
loop_by_frames(Frames, Conn, Handler, Next) ->
case {Frames, Next} of
{_, {stop, normal}} ->
{stop, normal};
{_, {stop, {abnormal, Reason}}} ->
{stop, {abnormal, Reason}};
{[], Next@1} ->
Next@1;
{[{control, {ping_frame, Payload}}], {continue, User_state}} ->
Sent = glisten@transport:send(
erlang:element(2, Conn),
erlang:element(3, Conn),
gramps@websocket:encode_pong_frame(Payload, none)
),
case Sent of
{ok, nil} ->
{continue, User_state};
{error, _} ->
{stop, {abnormal, <<"Failed to send PONG frame"/utf8>>}}
end;
{[{control, {close_frame, Reason@1}}], {continue, _}} ->
_ = glisten@transport:send(
erlang:element(2, Conn),
erlang:element(3, Conn),
gramps@websocket:encode_close_frame(Reason@1, none)
),
{stop, normal};
{[Frame | Rest], {continue, User_state@1}} ->
case ewe_ffi:rescue(
fun() -> Handler(Conn, User_state@1, {frame, Frame}) end
) of
{ok, {continue, New_user_state}} ->
loop_by_frames(
Rest,
Conn,
Handler,
{continue, New_user_state}
);
{ok, Stop} ->
Stop;
{error, _} ->
{stop, {abnormal, <<"Crash in websocket handler"/utf8>>}}
end
end.
-file("src/ewe/internal/websocket.gleam", 258).
?DOC(false).
-spec handle_user_message(
state(KPZ),
websocket_connection(),
KQB,
fun((websocket_connection(), KPZ, handler_message(KQB)) -> next(KPZ))
) -> gleam@otp@actor:next(state(KPZ), message(KQB)).
handle_user_message(State, Conn, User_message, Handler) ->
Call = ewe_ffi:rescue(
fun() ->
Handler(
Conn,
erlang:element(2, State),
{user_message, User_message}
)
end
),
case Call of
{ok, {continue, New_user_state}} ->
gleam@otp@actor:continue(
{state,
New_user_state,
erlang:element(3, State),
erlang:element(4, State)}
);
{ok, {stop, normal}} ->
close_compression(erlang:element(4, State)),
gleam@otp@actor:stop();
{ok, {stop, {abnormal, Reason}}} ->
close_compression(erlang:element(4, State)),
gleam@otp@actor:stop_abnormal(Reason);
{error, _} ->
close_compression(erlang:element(4, State)),
gleam@otp@actor:stop_abnormal(<<"Crash in websocket handler"/utf8>>)
end.
-file("src/ewe/internal/websocket.gleam", 309).
?DOC(false).
-spec send_frame(
fun((KQN, gleam@option:option(gramps@websocket@compression:context()), gleam@option:option(bitstring())) -> gleam@bytes_tree:bytes_tree()),
glisten@transport:transport(),
glisten@socket:socket(),
gleam@option:option(gramps@websocket@compression:context()),
KQN
) -> {ok, nil} | {error, glisten@socket:socket_reason()}.
send_frame(Encoder, Transport, Socket, Deflate, Data) ->
Frame = ewe_ffi:rescue(fun() -> _pipe = Encoder(Data, Deflate, none),
glisten@transport:send(Transport, Socket, _pipe) end),
case Frame of
{ok, Frame@1} ->
Frame@1;
{error, _} ->
erlang:error(#{gleam_error => panic,
message => <<"Sending WebSocket message from non-owning process"/utf8>>,
file => <<?FILEPATH/utf8>>,
module => <<"ewe/internal/websocket"/utf8>>,
function => <<"send_frame"/utf8>>,
line => 325})
end.
-file("src/ewe/internal/websocket.gleam", 330).
?DOC(false).
-spec set_socket_active_once(
glisten@transport:transport(),
glisten@socket:socket()
) -> nil.
set_socket_active_once(Transport, Socket) ->
_ = glisten@transport:set_opts(Transport, Socket, [{active_mode, once}]),
nil.
-file("src/ewe/internal/websocket.gleam", 288).
?DOC(false).
-spec after_start(
gleam@otp@actor:started(gleam@erlang@process:subject(message(any()))),
glisten@transport:transport(),
glisten@socket:socket()
) -> gleam@erlang@process:selector(gleam@erlang@process:down()).
after_start(Started, Transport, Socket) ->
Pid@1 = case gleam@erlang@process:subject_owner(erlang:element(3, Started)) of
{ok, Pid} -> Pid;
_assert_fail ->
erlang:error(#{gleam_error => let_assert,
message => <<"Pattern match failed, no pattern matched the value."/utf8>>,
file => <<?FILEPATH/utf8>>,
module => <<"ewe/internal/websocket"/utf8>>,
function => <<"after_start"/utf8>>,
line => 294,
value => _assert_fail,
start => 8265,
'end' => 8321,
pattern_start => 8276,
pattern_end => 8283})
end,
_ = glisten@transport:controlling_process(Transport, Socket, Pid@1),
set_socket_active_once(Transport, Socket),
Selector = gleam@erlang@process:select_specific_monitor(
gleam_erlang_ffi:new_selector(),
gleam@erlang@process:monitor(Pid@1),
fun gleam@function:identity/1
),
Selector.
-file("src/ewe/internal/websocket.gleam", 167).
?DOC(false).
-spec handle_valid_packet(
state(KPJ),
websocket_connection(),
bitstring(),
fun((websocket_connection(), KPJ, handler_message(KPL)) -> next(KPJ))
) -> gleam@otp@actor:next(state(KPJ), message(KPL)).
handle_valid_packet(State, Conn, Data, Handler) ->
Buffer = <<(erlang:element(3, State))/bitstring, Data/bitstring>>,
{Frames, Rest} = gramps@websocket:decode_many_frames(
Buffer,
get_inflate(erlang:element(4, State)),
[]
),
Next = begin
_pipe = gramps@websocket:aggregate_frames(Frames, none, []),
gleam@result:map(
_pipe,
fun(_capture) ->
loop_by_frames(
_capture,
Conn,
Handler,
{continue, erlang:element(2, State)}
)
end
)
end,
case Next of
{ok, {continue, New_user_state}} ->
set_socket_active_once(
erlang:element(2, Conn),
erlang:element(3, Conn)
),
gleam@otp@actor:continue(
{state, New_user_state, Rest, erlang:element(4, State)}
);
{ok, {stop, normal}} ->
close_compression(erlang:element(4, State)),
gleam@otp@actor:stop();
{ok, {stop, {abnormal, Reason}}} ->
close_compression(erlang:element(4, State)),
gleam@otp@actor:stop_abnormal(Reason);
{error, nil} ->
close_compression(erlang:element(4, State)),
gleam@otp@actor:stop_abnormal(<<"Received malformed message"/utf8>>)
end.
-file("src/ewe/internal/websocket.gleam", 109).
?DOC(false).
-spec start(
glisten@transport:transport(),
glisten@socket:socket(),
fun((websocket_connection(), gleam@erlang@process:selector(KOZ)) -> {KPB,
gleam@erlang@process:selector(KOZ)}),
fun((websocket_connection(), KPB, handler_message(KOZ)) -> next(KPB)),
list(binary()),
boolean()
) -> {ok, gleam@erlang@process:selector(gleam@erlang@process:down())} |
{error, gleam@otp@actor:start_error()}.
start(Transport, Socket, On_init, Handler, Extensions, Permessage_deflate) ->
_pipe@4 = gleam@otp@actor:new_with_initialiser(
1000,
fun(Subject) ->
Takeovers = gramps@websocket:get_context_takeovers(Extensions),
Compression = case Permessage_deflate of
true ->
{some, gramps@websocket@compression:init(Takeovers)};
false ->
none
end,
Conn = {websocket_connection,
Transport,
Socket,
get_deflate(Compression)},
{User_state, User_selector} = On_init(
Conn,
gleam_erlang_ffi:new_selector()
),
Selector = begin
_pipe = gleam_erlang_ffi:map_selector(
User_selector,
fun(Field@0) -> {user, Field@0} end
),
gleam_erlang_ffi:merge_selector(_pipe, glisten_selector())
end,
Ws_state = {state, User_state, <<>>, Compression},
_pipe@1 = gleam@otp@actor:initialised(Ws_state),
_pipe@2 = gleam@otp@actor:selecting(_pipe@1, Selector),
_pipe@3 = gleam@otp@actor:returning(_pipe@2, Subject),
{ok, _pipe@3}
end
),
_pipe@5 = gleam@otp@actor:on_message(
_pipe@4,
fun(State, Msg) ->
Conn@1 = {websocket_connection,
Transport,
Socket,
get_deflate(erlang:element(4, State))},
case Msg of
{valid, {packet, Data}} ->
handle_valid_packet(State, Conn@1, Data, Handler);
{valid, close} ->
close_compression(erlang:element(4, State)),
gleam@otp@actor:stop();
{user, User_message} ->
handle_user_message(State, Conn@1, User_message, Handler);
invalid ->
close_compression(erlang:element(4, State)),
gleam@otp@actor:stop_abnormal(
<<"Received malformed message"/utf8>>
)
end
end
),
_pipe@6 = gleam@otp@actor:start(_pipe@5),
gleam@result:map(
_pipe@6,
fun(_capture) -> after_start(_capture, Transport, Socket) end
).