-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 => <>, 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 => <>, module => <<"ewe/internal/websocket"/utf8>>, function => <<"send_frame"/utf8>>, line => 259}) end.