Current section

Files

Jump to
ewe src ewe@internal@stream@websocket.erl
Raw

src/ewe@internal@stream@websocket.erl

-module(ewe@internal@stream@websocket).
-compile([no_auto_import, nowarn_unused_vars, nowarn_unused_function, nowarn_nomatch, inline]).
-define(FILEPATH, "src/ewe/internal/stream/websocket.gleam").
-export([send_frame/5, send_close_frame/3, start/7]).
-export_type([websocket_connection/0, websocket_message/1, websocket_next/2, websocket_state/1, internal_message/1, resolve_state/2]).
-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(),
websocks:context()}.
-type websocket_message(MRR) :: {frame, websocks:frame()} | {user_message, MRR}.
-type websocket_next(MRS, MRT) :: {continue,
MRS,
gleam@option:option(gleam@erlang@process:selector(MRT))} |
normal_stop |
{abnormal_stop, binary()}.
-type websocket_state(MRU) :: {websocket_state, MRU, websocks:context()}.
-type internal_message(MRV) :: {packet, bitstring()} |
close |
tcp_passive |
{user, MRV} |
invalid.
-type resolve_state(MRW, MRX) :: {resolve_state,
glisten@socket:socket(),
glisten@transport:transport(),
fun((websocket_connection(), MRW, websocket_message(MRX)) -> websocket_next(MRW, MRX)),
websocket_next(MRW, internal_message(MRX))}.
-file("src/ewe/internal/stream/websocket.gleam", 124).
?DOC(false).
-spec select_valid_record(
gleam@erlang@process:selector(internal_message(MSO)),
binary()
) -> gleam@erlang@process:selector(internal_message(MSO)).
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/stream/websocket.gleam", 138).
?DOC(false).
-spec create_socket_selector() -> gleam@erlang@process:selector(internal_message(any())).
create_socket_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
),
_pipe@4 = gleam@erlang@process:select_record(
_pipe@3,
erlang:binary_to_atom(<<"ssl_closed"/utf8>>),
1,
fun(_) -> close end
),
gleam@erlang@process:select_record(
_pipe@4,
erlang:binary_to_atom(<<"tcp_passive"/utf8>>),
1,
fun(_) -> tcp_passive end
).
-file("src/ewe/internal/stream/websocket.gleam", 152).
?DOC(false).
-spec user_selector(gleam@option:option(gleam@erlang@process:selector(MSW))) -> gleam@option:option(gleam@erlang@process:selector(internal_message(MSW))).
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/stream/websocket.gleam", 468).
?DOC(false).
-spec handle_close(
fun((websocket_connection(), MUR) -> nil),
websocket_state(MUR),
websocket_connection(),
gleam@option:option(binary())
) -> gleam@otp@actor:next(websocket_state(MUR), internal_message(any())).
handle_close(On_close, State, Conn, Abnormal_reason) ->
websocks:close_context(erlang:element(3, State)),
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/stream/websocket.gleam", 428).
?DOC(false).
-spec handle_user_message(
glisten@transport:transport(),
glisten@socket:socket(),
websocket_state(MUJ),
MUL,
fun((websocket_connection(), MUJ, websocket_message(MUL)) -> websocket_next(MUJ, MUL)),
fun((websocket_connection(), MUJ) -> nil)
) -> gleam@otp@actor:next(websocket_state(MUJ), internal_message(MUL)).
handle_user_message(Transport, Socket, State, User_message, Handler, On_close) ->
Conn = {websocket_connection, Transport, Socket, erlang:element(3, State)},
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(
create_socket_selector(),
_capture
)
end
)
end,
Next = gleam@otp@actor:continue(
{websocket_state, New_user_state, erlang:element(3, 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/stream/websocket.gleam", 350).
?DOC(false).
-spec handle_frame(
resolve_state(MUC, MUD),
websocks:context(),
websocks:frame()
) -> websocks:resolve_next(resolve_state(MUC, MUD)).
handle_frame(State, Context, Frame) ->
case Frame of
{control, {ping, Payload}} ->
case erlang:byte_size(Payload) of
Size when Size > 125 ->
{stop,
{resolve_state,
erlang:element(2, State),
erlang:element(3, State),
erlang:element(4, State),
{abnormal_stop,
<<"control frames are only allowed to have payload up to and including 125 octets"/utf8>>}}};
_ ->
Sent = glisten@transport:send(
erlang:element(3, State),
erlang:element(2, State),
begin
_pipe = websocks:encode_pong_frame(Payload, none),
gleam@bytes_tree:from_bit_array(_pipe)
end
),
case Sent of
{ok, nil} ->
{continue, State};
{error, _} ->
{stop,
{resolve_state,
erlang:element(2, State),
erlang:element(3, State),
erlang:element(4, State),
{abnormal_stop,
<<"Failed to send PONG frame"/utf8>>}}}
end
end;
{control, {close, Reason}} ->
_ = glisten@transport:send(
erlang:element(3, State),
erlang:element(2, State),
begin
_pipe@1 = websocks:encode_close_frame(Reason, none),
gleam@bytes_tree:from_bit_array(_pipe@1)
end
),
{stop,
{resolve_state,
erlang:element(2, State),
erlang:element(3, State),
erlang:element(4, State),
normal_stop}};
Frame@1 ->
{User_state@1, Selector@1} = case erlang:element(5, State) of
{continue, User_state, Selector} -> {User_state, Selector};
_assert_fail ->
erlang:error(#{gleam_error => let_assert,
message => <<"Pattern match failed, no pattern matched the value."/utf8>>,
file => <<?FILEPATH/utf8>>,
module => <<"ewe/internal/stream/websocket"/utf8>>,
function => <<"handle_frame"/utf8>>,
line => 399,
value => _assert_fail,
start => 12463,
'end' => 12517,
pattern_start => 12474,
pattern_end => 12504})
end,
Conn = {websocket_connection,
erlang:element(3, State),
erlang:element(2, State),
Context},
Call = ewe_ffi:rescue(
fun() ->
(erlang:element(4, State))(
Conn,
User_state@1,
{frame, Frame@1}
)
end
),
case Call of
{ok, {continue, User_state@2, New_selector}} ->
Next_selector = begin
_pipe@2 = user_selector(New_selector),
_pipe@3 = gleam@option:'or'(_pipe@2, Selector@1),
gleam@option:map(
_pipe@3,
fun(_capture) ->
gleam_erlang_ffi:merge_selector(
create_socket_selector(),
_capture
)
end
)
end,
{continue,
{resolve_state,
erlang:element(2, State),
erlang:element(3, State),
erlang:element(4, State),
{continue, User_state@2, Next_selector}}};
{ok, normal_stop} ->
{stop,
{resolve_state,
erlang:element(2, State),
erlang:element(3, State),
erlang:element(4, State),
normal_stop}};
{ok, {abnormal_stop, Reason@1}} ->
{stop,
{resolve_state,
erlang:element(2, State),
erlang:element(3, State),
erlang:element(4, State),
{abnormal_stop, Reason@1}}};
{error, _} ->
{stop,
{resolve_state,
erlang:element(2, State),
erlang:element(3, State),
erlang:element(4, State),
{abnormal_stop,
<<"Crash in websocket handler"/utf8>>}}}
end
end.
-file("src/ewe/internal/stream/websocket.gleam", 295).
?DOC(false).
-spec handle_valid_packet(
glisten@transport:transport(),
glisten@socket:socket(),
websocket_state(MTU),
bitstring(),
fun((websocket_connection(), MTU, websocket_message(MTW)) -> websocket_next(MTU, MTW)),
fun((websocket_connection(), MTU) -> nil)
) -> gleam@otp@actor:next(websocket_state(MTU), internal_message(MTW)).
handle_valid_packet(Transport, Socket, State, Data, Handler, On_close) ->
Conn = {websocket_connection, Transport, Socket, erlang:element(3, State)},
Result = websocks:process_incoming_frames(
Data,
erlang:element(3, State),
{resolve_state,
Socket,
Transport,
Handler,
{continue, erlang:element(2, State), none}},
fun handle_frame/3
),
case Result of
{ok, {Resolved_state, Context}} ->
case erlang:element(5, Resolved_state) of
{continue, User_state, Selector} ->
Next = gleam@otp@actor:continue(
{websocket_state, User_state, Context}
),
case Selector of
{some, Selector@1} ->
gleam@otp@actor:with_selector(Next, Selector@1);
none ->
Next
end;
normal_stop ->
handle_close(On_close, State, Conn, none);
{abnormal_stop, Reason} ->
handle_close(On_close, State, Conn, {some, Reason})
end;
{error, _} ->
handle_close(
On_close,
State,
Conn,
{some, <<"Received malformed message"/utf8>>}
)
end.
-file("src/ewe/internal/stream/websocket.gleam", 235).
?DOC(false).
-spec send_frame(
fun((bitstring(), websocks:context(), gleam@option:option(bitstring())) -> bitstring()),
glisten@transport:transport(),
glisten@socket:socket(),
websocks:context(),
bitstring()
) -> {ok, nil} | {error, glisten@socket:socket_reason()}.
send_frame(Encoder, Transport, Socket, Context, Payload) ->
Frame = ewe_ffi:rescue(fun() -> _pipe = Encoder(Payload, Context, none),
_pipe@1 = gleam@bytes_tree:from_bit_array(_pipe),
glisten@transport:send(Transport, Socket, _pipe@1) 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/stream/websocket"/utf8>>,
function => <<"send_frame"/utf8>>,
line => 257})
end.
-file("src/ewe/internal/stream/websocket.gleam", 262).
?DOC(false).
-spec send_close_frame(
glisten@transport:transport(),
glisten@socket:socket(),
websocks:close_reason()
) -> websocket_next(any(), any()).
send_close_frame(Transport, Socket, Code) ->
Frame = ewe_ffi:rescue(
fun() -> _pipe = websocks:encode_close_frame(Code, none),
_pipe@1 = gleam@bytes_tree:from_bit_array(_pipe),
glisten@transport:send(Transport, Socket, _pipe@1) end
),
case Frame of
{ok, {ok, nil}} ->
normal_stop;
{ok, {error, Reason}} ->
{abnormal_stop,
<<"Failed to send close frame: "/utf8,
(gleam@string:inspect(Reason))/binary>>};
{error, Reason@1} ->
logging:log(
error,
<<"Frame should be sent from the WebSocket connection, but was sent from different process: "/utf8,
(gleam@string:inspect(Reason@1))/binary>>
),
erlang:error(#{gleam_error => panic,
message => <<"Sending WebSocket message from non-owning process"/utf8>>,
file => <<?FILEPATH/utf8>>,
module => <<"ewe/internal/stream/websocket"/utf8>>,
function => <<"send_close_frame"/utf8>>,
line => 285})
end.
-file("src/ewe/internal/stream/websocket.gleam", 490).
?DOC(false).
-spec after_start(
gleam@otp@actor:started(gleam@erlang@process:subject(internal_message(any()))),
glisten@transport:transport(),
glisten@socket:socket()
) -> gleam@otp@actor:started(nil).
after_start(Started, Transport, Socket) ->
_ = glisten@transport:set_opts(
Transport,
Socket,
[{active_mode, {count, 100}}]
),
{started, erlang:element(2, Started), nil}.
-file("src/ewe/internal/stream/websocket.gleam", 169).
?DOC(false).
-spec start(
glisten@transport:transport(),
glisten@socket:socket(),
fun((websocket_connection(), gleam@erlang@process:selector(MTD)) -> {MTC,
gleam@erlang@process:selector(MTD)}),
fun((websocket_connection(), MTC, websocket_message(MTD)) -> websocket_next(MTC, MTD)),
fun((websocket_connection(), MTC) -> nil),
list(binary()),
boolean()
) -> {ok, gleam@otp@actor:started(nil)} | {error, gleam@otp@actor:start_error()}.
start(
Transport,
Socket,
On_init,
Handler,
On_close,
Extensions,
Permessage_deflate
) ->
_pipe@6 = gleam@otp@actor:new_with_initialiser(
1000,
fun(Subject) ->
Context_takeovers = websocks:get_context_takeovers(Extensions),
Compression = case Permessage_deflate of
true ->
{some, Context_takeovers};
false ->
none
end,
Context = websocks:create_context(Compression),
{User_state, User_selector} = begin
_pipe = {websocket_connection, Transport, Socket, Context},
On_init(_pipe, gleam_erlang_ffi:new_selector())
end,
Selector = begin
_pipe@1 = gleam_erlang_ffi:map_selector(
User_selector,
fun(Field@0) -> {user, Field@0} end
),
gleam_erlang_ffi:merge_selector(
_pipe@1,
create_socket_selector()
)
end,
_pipe@2 = {websocket_state, User_state, Context},
_pipe@3 = gleam@otp@actor:initialised(_pipe@2),
_pipe@4 = gleam@otp@actor:selecting(_pipe@3, Selector),
_pipe@5 = gleam@otp@actor:returning(_pipe@4, Subject),
{ok, _pipe@5}
end
),
_pipe@7 = gleam@otp@actor:on_message(_pipe@6, fun(State, Msg) -> case Msg of
{packet, Data} ->
handle_valid_packet(
Transport,
Socket,
State,
Data,
Handler,
On_close
);
{user, User_message} ->
handle_user_message(
Transport,
Socket,
State,
User_message,
Handler,
On_close
);
close ->
Conn = {websocket_connection,
Transport,
Socket,
erlang:element(3, State)},
handle_close(On_close, State, Conn, none);
invalid ->
Conn@1 = {websocket_connection,
Transport,
Socket,
erlang:element(3, State)},
handle_close(
On_close,
State,
Conn@1,
{some, <<"Received malformed message"/utf8>>}
);
tcp_passive ->
_ = glisten@transport:set_opts(
Transport,
Socket,
[{active_mode, {count, 100}}]
),
gleam@otp@actor:continue(State)
end end),
_pipe@8 = gleam@otp@actor:start(_pipe@7),
gleam@result:map(
_pipe@8,
fun(_capture) -> after_start(_capture, Transport, Socket) end
).