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([start/7, send_frame/5]).
-export_type([websocket_connection/0, websocket_message/1, websocket_next/2, websocket_state/1, internal_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 websocket_connection() :: {websocket_connection,
glisten@transport:transport(),
glisten@socket:socket(),
gleam@option:option(ewe@internal@gramps@websocket@compression:context())}.
-type websocket_message(MLW) :: {websocket_frame,
ewe@internal@gramps@websocket:frame()} |
{user_message, MLW}.
-type websocket_next(MLX, MLY) :: {continue,
MLX,
gleam@option:option(gleam@erlang@process:selector(MLY))} |
normal_stop |
{abnormal_stop, binary()}.
-type websocket_state(MLZ) :: {websocket_state,
MLZ,
gleam@option:option(ewe@internal@gramps@websocket@compression:compression()),
bitstring(),
list(ewe@internal@gramps@websocket:parsed_frame())}.
-type internal_message(MMA) :: {packet, bitstring()} |
close |
{user, MMA} |
invalid.
-file("src/ewe/internal/websocket.gleam", 114).
?DOC(false).
-spec get_deflate(
gleam@option:option(ewe@internal@gramps@websocket@compression:compression())
) -> gleam@option:option(ewe@internal@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", 121).
?DOC(false).
-spec get_inflate(
gleam@option:option(ewe@internal@gramps@websocket@compression:compression())
) -> gleam@option:option(ewe@internal@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", 132).
?DOC(false).
-spec select_valid_record(
gleam@erlang@process:selector(internal_message(MMV)),
binary()
) -> gleam@erlang@process:selector(internal_message(MMV)).
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({packet, Data})
end
)
end
),
gleam@result:unwrap(_pipe, invalid)
end
).
-file("src/ewe/internal/websocket.gleam", 146).
?DOC(false).
-spec glisten_selector() -> gleam@erlang@process:selector(internal_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(_) -> close end
),
gleam@erlang@process:select_record(
_pipe@3,
erlang:binary_to_atom(<<"ssl_closed"/utf8>>),
1,
fun(_) -> close end
).
-file("src/ewe/internal/websocket.gleam", 159).
?DOC(false).
-spec user_selector(gleam@option:option(gleam@erlang@process:selector(MND))) -> gleam@option:option(gleam@erlang@process:selector(internal_message(MND))).
user_selector(Selector) ->
gleam@option:map(
Selector,
fun(Selector@1) ->
gleam_erlang_ffi:map_selector(
Selector@1,
fun(Field@0) -> {user, Field@0} end
)
end
).
-file("src/ewe/internal/websocket.gleam", 170).
?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", 390).
?DOC(false).
-spec separate_frames(
list(ewe@internal@gramps@websocket:parsed_frame()),
list(ewe@internal@gramps@websocket:parsed_frame()),
list(ewe@internal@gramps@websocket:frame())
) -> {list(ewe@internal@gramps@websocket:parsed_frame()),
list(ewe@internal@gramps@websocket:frame())}.
separate_frames(Frames, Data_frames, Control_frames) ->
case Frames of
[] ->
{lists:reverse(Data_frames), lists:reverse(Control_frames)};
[{complete, {control, Control_frame}} | Rest] ->
separate_frames(
Rest,
Data_frames,
[{control, Control_frame} | Control_frames]
);
[Data_frame | Rest@1] ->
separate_frames(Rest@1, [Data_frame | Data_frames], Control_frames)
end.
-file("src/ewe/internal/websocket.gleam", 529).
?DOC(false).
-spec handle_close(
fun((websocket_connection(), MPO) -> nil),
websocket_state(MPO),
websocket_connection(),
gleam@option:option(binary())
) -> gleam@otp@actor:next(websocket_state(MPO), internal_message(any())).
handle_close(On_close, State, Conn, Abnormal_reason) ->
gleam@option:map(
erlang:element(3, State),
fun(Compression) ->
ewe@internal@gramps@websocket@compression:close(
erlang:element(3, Compression)
),
ewe@internal@gramps@websocket@compression:close(
erlang:element(2, Compression)
)
end
),
On_close(Conn, erlang:element(2, State)),
case Abnormal_reason of
{some, Reason} ->
logging:log(
error,
<<"WebSocket connection closed abnormally: "/utf8,
Reason/binary>>
),
gleam@otp@actor:stop_abnormal(Reason);
none ->
gleam@otp@actor:stop()
end.
-file("src/ewe/internal/websocket.gleam", 555).
?DOC(false).
-spec after_start(
gleam@otp@actor:started(gleam@erlang@process:subject(internal_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 => 561,
value => _assert_fail,
start => 17347,
'end' => 17403,
pattern_start => 17358,
pattern_end => 17365})
end,
_ = glisten@transport:controlling_process(Transport, Socket, Pid@1),
set_socket_active_once(Transport, Socket),
gleam@erlang@process:select_specific_monitor(
gleam_erlang_ffi:new_selector(),
gleam@erlang@process:monitor(Pid@1),
fun gleam@function:identity/1
).
-file("src/ewe/internal/websocket.gleam", 491).
?DOC(false).
-spec handle_user_message(
websocket_state(MPG),
websocket_connection(),
MPI,
fun((websocket_connection(), MPG, websocket_message(MPI)) -> websocket_next(MPG, MPI)),
fun((websocket_connection(), MPG) -> nil)
) -> gleam@otp@actor:next(websocket_state(MPG), internal_message(MPI)).
handle_user_message(State, Conn, User_message, Handler, On_close) ->
Call = ewe_ffi:rescue(
fun() ->
Handler(
Conn,
erlang:element(2, State),
{user_message, User_message}
)
end
),
case Call of
{ok, {continue, New_user_state, New_selector}} ->
Next_selector = begin
_pipe = user_selector(New_selector),
gleam@option:map(
_pipe,
fun(_capture) ->
gleam_erlang_ffi:merge_selector(
glisten_selector(),
_capture
)
end
)
end,
Next = gleam@otp@actor:continue(
{websocket_state,
New_user_state,
erlang:element(3, State),
erlang:element(4, State),
erlang:element(5, State)}
),
case Next_selector of
{some, Selector} ->
gleam@otp@actor:with_selector(Next, Selector);
none ->
Next
end;
{ok, normal_stop} ->
handle_close(On_close, State, Conn, none);
{ok, {abnormal_stop, Reason}} ->
handle_close(On_close, State, Conn, {some, Reason});
{error, _} ->
handle_close(
On_close,
State,
Conn,
{some, <<"Crash in websocket handler"/utf8>>}
)
end.
-file("src/ewe/internal/websocket.gleam", 408).
?DOC(false).
-spec loop_by_frames(
list(ewe@internal@gramps@websocket:frame()),
websocket_connection(),
fun((websocket_connection(), MOW, websocket_message(MOX)) -> websocket_next(MOW, MOX)),
websocket_next(MOW, internal_message(MOX))
) -> websocket_next(MOW, internal_message(MOX)).
loop_by_frames(Frames, Conn, Handler, Next) ->
case {Frames, Next} of
{_, normal_stop} ->
normal_stop;
{_, {abnormal_stop, Reason}} ->
{abnormal_stop, Reason};
{[], Next@1} ->
Next@1;
{[{control, {ping_frame, Payload}} | Rest], {continue, User_state, _}} ->
case erlang:byte_size(Payload) of
Size when Size > 125 ->
{abnormal_stop,
<<"control frames are only allowed to have payload up to and including 125 octets"/utf8>>};
_ ->
Sent = glisten@transport:send(
erlang:element(2, Conn),
erlang:element(3, Conn),
ewe@internal@gramps@websocket:encode_pong_frame(
Payload,
none
)
),
case Sent of
{ok, nil} ->
loop_by_frames(
Rest,
Conn,
Handler,
{continue, User_state, none}
);
{error, _} ->
{abnormal_stop,
<<"Failed to send PONG frame"/utf8>>}
end
end;
{[{control, {close_frame, Reason@1}} | _], {continue, _, _}} ->
_ = glisten@transport:send(
erlang:element(2, Conn),
erlang:element(3, Conn),
ewe@internal@gramps@websocket:encode_close_frame(Reason@1, none)
),
normal_stop;
{[{continuation, _, _} | _], {continue, _, _}} ->
{abnormal_stop, <<"Unexpected continuation frame"/utf8>>};
{[Frame | Rest@1], {continue, User_state@1, Selector}} ->
Call = ewe_ffi:rescue(
fun() ->
Handler(Conn, User_state@1, {websocket_frame, Frame})
end
),
case Call of
{ok, {continue, User_state@2, New_selector}} ->
Next_selector = begin
_pipe = user_selector(New_selector),
_pipe@1 = gleam@option:'or'(_pipe, Selector),
gleam@option:map(
_pipe@1,
fun(_capture) ->
gleam_erlang_ffi:merge_selector(
glisten_selector(),
_capture
)
end
)
end,
loop_by_frames(
Rest@1,
Conn,
Handler,
{continue, User_state@2, Next_selector}
);
{ok, normal_stop} ->
normal_stop;
{ok, {abnormal_stop, Reason@2}} ->
{abnormal_stop, Reason@2};
{error, _} ->
{abnormal_stop, <<"Crash in websocket handler"/utf8>>}
end
end.
-file("src/ewe/internal/websocket.gleam", 307).
?DOC(false).
-spec handle_frames_processing(
websocket_state(MOI),
websocket_connection(),
list(ewe@internal@gramps@websocket:parsed_frame()),
bitstring(),
fun((websocket_connection(), MOI, websocket_message(MOL)) -> websocket_next(MOI, MOL)),
fun((websocket_connection(), MOI) -> nil)
) -> gleam@otp@actor:next(websocket_state(MOI), internal_message(MOL)).
handle_frames_processing(State, Conn, Frames, Rest, Handler, On_close) ->
Frames@1 = lists:append(erlang:element(5, State), Frames),
{Data_frames, Control_frames} = separate_frames(Frames@1, [], []),
Control_result = case Control_frames of
[] ->
{continue, erlang:element(2, State), none};
_ ->
loop_by_frames(
Control_frames,
Conn,
Handler,
{continue, erlang:element(2, State), none}
)
end,
case Control_result of
normal_stop ->
handle_close(On_close, State, Conn, none);
{abnormal_stop, Reason} ->
handle_close(On_close, State, Conn, {some, Reason});
{continue, _, _} ->
Aggregated = ewe@internal@gramps@websocket:aggregate_frames(
Data_frames,
none,
[],
get_inflate(erlang:element(3, State))
),
case Aggregated of
{ok, []} ->
set_socket_active_once(
erlang:element(2, Conn),
erlang:element(3, Conn)
),
gleam@otp@actor:continue(
{websocket_state,
erlang:element(2, State),
erlang:element(3, State),
Rest,
Data_frames}
);
{ok, Data_frames@1} ->
Next = loop_by_frames(
Data_frames@1,
Conn,
Handler,
{continue, erlang:element(2, State), none}
),
case Next of
{continue, User_state, Selector} ->
set_socket_active_once(
erlang:element(2, Conn),
erlang:element(3, Conn)
),
Next@1 = gleam@otp@actor:continue(
{websocket_state,
User_state,
erlang:element(3, State),
Rest,
[]}
),
case Selector of
{some, Selector@1} ->
gleam@otp@actor:with_selector(
Next@1,
Selector@1
);
none ->
Next@1
end;
normal_stop ->
handle_close(On_close, State, Conn, none);
{abnormal_stop, Reason@1} ->
handle_close(
On_close,
State,
Conn,
{some, Reason@1}
)
end;
{error, nil} ->
handle_close(
On_close,
State,
Conn,
{some, <<"Received malformed message"/utf8>>}
)
end
end.
-file("src/ewe/internal/websocket.gleam", 269).
?DOC(false).
-spec handle_valid_packet(
websocket_state(MOA),
websocket_connection(),
bitstring(),
fun((websocket_connection(), MOA, websocket_message(MOC)) -> websocket_next(MOA, MOC)),
fun((websocket_connection(), MOA) -> nil)
) -> gleam@otp@actor:next(websocket_state(MOA), internal_message(MOC)).
handle_valid_packet(State, Conn, Data, Handler, On_close) ->
Buffer = <<(erlang:element(4, State))/bitstring, Data/bitstring>>,
Decoded = ewe@internal@gramps@websocket:decode_many_frames_result(
Buffer,
get_inflate(erlang:element(3, State)),
[]
),
case Decoded of
{ok, {Frames, Rest}} ->
handle_frames_processing(
State,
Conn,
Frames,
Rest,
Handler,
On_close
);
{error, {need_more_data_accumulated, Parsed, Rest@1}} ->
set_socket_active_once(
erlang:element(2, Conn),
erlang:element(3, Conn)
),
gleam@otp@actor:continue(
{websocket_state,
erlang:element(2, State),
erlang:element(3, State),
Rest@1,
lists:append(erlang:element(5, State), Parsed)}
);
{error, contains_invalid_frame} ->
handle_close(
On_close,
State,
Conn,
{some, <<"Received malformed message"/utf8>>}
)
end.
-file("src/ewe/internal/websocket.gleam", 180).
?DOC(false).
-spec start(
glisten@transport:transport(),
glisten@socket:socket(),
fun((websocket_connection(), gleam@erlang@process:selector(MNK)) -> {MNJ,
gleam@erlang@process:selector(MNK)}),
fun((websocket_connection(), MNJ, websocket_message(MNK)) -> websocket_next(MNJ, MNK)),
fun((websocket_connection(), MNJ) -> nil),
list(binary()),
boolean()
) -> {ok, gleam@erlang@process:selector(gleam@erlang@process:down())} |
{error, gleam@otp@actor:start_error()}.
start(
Transport,
Socket,
On_init,
Handler,
On_close,
Extensions,
Permessage_deflate
) ->
_pipe@4 = gleam@otp@actor:new_with_initialiser(
1000,
fun(Subject) ->
Takeovers = ewe@internal@gramps@websocket:get_context_takeovers(
Extensions
),
Deflate = case Permessage_deflate of
true ->
{some,
ewe@internal@gramps@websocket@compression:init(
Takeovers
)};
false ->
none
end,
Conn = {websocket_connection,
Transport,
Socket,
get_deflate(Deflate)},
{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 = {websocket_state, User_state, Deflate, <<>>, []},
_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(3, State))},
case Msg of
{packet, Data} ->
handle_valid_packet(State, Conn@1, Data, Handler, On_close);
close ->
handle_close(On_close, State, Conn@1, none);
{user, User_message} ->
handle_user_message(
State,
Conn@1,
User_message,
Handler,
On_close
);
invalid ->
handle_close(
On_close,
State,
Conn@1,
{some, <<"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
).
-file("src/ewe/internal/websocket.gleam", 238).
?DOC(false).
-spec send_frame(
fun((MNU, gleam@option:option(ewe@internal@gramps@websocket@compression:context()), gleam@option:option(bitstring())) -> gleam@bytes_tree:bytes_tree()),
glisten@transport:transport(),
glisten@socket:socket(),
gleam@option:option(ewe@internal@gramps@websocket@compression:context()),
MNU
) -> {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, Reason} ->
logging:log(
error,
<<"Frame should be sent from the WebSocket connection, but was sent from different process: "/utf8,
(gleam@string:inspect(Reason))/binary>>
),
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 => 259})
end.