-module(stratus). -compile([no_auto_import, nowarn_unused_vars, nowarn_unused_function, nowarn_nomatch]). -export([websocket/3, with_connect_timeout/2, on_close/2, on_handshake_error/2, send_message/2, send_text_message/2, send_binary_message/2, send_ping/2, close/1, initialize/1]). -export_type([connection/0, internal_message/1, message/1, builder/2, state/2, handshake_error/0]). -opaque connection() :: {connection, stratus@internal@socket:socket(), stratus@internal@transport:transport()}. -opaque internal_message(KMG) :: started | {user_message, KMG} | {err, stratus@internal@socket:socket_reason()} | {data, bitstring()} | closed | shutdown. -type message(KMH) :: {text, binary()} | {binary, bitstring()} | {user, KMH}. -opaque builder(KMI, KMJ) :: {builder, gleam@http@request:request(binary()), integer(), fun(() -> {KMI, gleam@option:option(gleam@erlang@process:selector(KMJ))}), fun((message(KMJ), KMI, connection()) -> gleam@otp@actor:next(KMJ, KMI)), fun((KMI) -> nil), fun((gleam@http@response:response(bitstring())) -> nil)}. -type state(KMK, KML) :: {state, bitstring(), gleam@option:option(gramps@websocket:frame()), gleam@erlang@process:subject(internal_message(KML)), gleam@option:option(stratus@internal@socket:socket()), KMK}. -type handshake_error() :: {sock, stratus@internal@socket:socket_reason()} | {protocol, bitstring()} | {upgrade_failed, gleam@http@response:response(bitstring())}. -spec from_socket_message(stratus@internal@socket:socket_message()) -> internal_message(any()). from_socket_message(Msg) -> case Msg of {data, Bits} -> {data, Bits}; closed -> closed; {err, Reason} -> {err, Reason} end. -spec websocket( gleam@http@request:request(binary()), fun(() -> {KMP, gleam@option:option(gleam@erlang@process:selector(KMQ))}), fun((message(KMQ), KMP, connection()) -> gleam@otp@actor:next(KMQ, KMP)) ) -> builder(KMP, KMQ). websocket(Req, Init, Loop) -> {builder, Req, 5000, Init, Loop, fun(_) -> nil end, fun(_) -> nil end}. -spec with_connect_timeout(builder(KMY, KMZ), integer()) -> builder(KMY, KMZ). with_connect_timeout(Builder, Timeout) -> erlang:setelement(3, Builder, Timeout). -spec on_close(builder(KNE, KNF), fun((KNE) -> nil)) -> builder(KNE, KNF). on_close(Builder, On_close) -> erlang:setelement(6, Builder, On_close). -spec on_handshake_error( builder(KNK, KNL), fun((gleam@http@response:response(bitstring())) -> nil) ) -> builder(KNK, KNL). on_handshake_error(Builder, On_handshake_error) -> erlang:setelement(7, Builder, On_handshake_error). -spec handle_frame( builder(KNZ, KOA), stratus@internal@transport:transport(), state(KNZ, KOA), connection(), gramps@websocket:frame() ) -> gleam@otp@actor:next(internal_message(KOA), state(KNZ, KOA)). handle_frame(Builder, Transport, State, Conn, Frame) -> _assert_subject = erlang:element(5, State), {some, Socket} = case _assert_subject of {some, _} -> _assert_subject; _assert_fail -> erlang:error(#{gleam_error => let_assert, message => <<"Assertion pattern match failed"/utf8>>, value => _assert_fail, module => <<"stratus"/utf8>>, function => <<"handle_frame"/utf8>>, line => 354}) end, case Frame of {data, {text_frame, _, Data}} -> _assert_subject@1 = gleam@bit_array:to_string(Data), {ok, Str} = case _assert_subject@1 of {ok, _} -> _assert_subject@1; _assert_fail@1 -> erlang:error(#{gleam_error => let_assert, message => <<"Assertion pattern match failed"/utf8>>, value => _assert_fail@1, module => <<"stratus"/utf8>>, function => <<"handle_frame"/utf8>>, line => 357}) end, Res = gleam_erlang_ffi:rescue( fun() -> (erlang:element(5, Builder))( {text, Str}, erlang:element(6, State), Conn ) end ), case Res of {ok, {continue, User_state, User_selector}} -> New_state = erlang:setelement(6, State, User_state), case User_selector of {some, User_selector@1} -> Selector = begin _pipe = User_selector@1, _pipe@1 = gleam_erlang_ffi:map_selector( _pipe, fun(Field@0) -> {user_message, Field@0} end ), gleam_erlang_ffi:merge_selector( _pipe@1, gleam_erlang_ffi:map_selector( stratus@internal@socket:selector(), fun from_socket_message/1 ) ) end, {continue, New_state, {some, Selector}}; _ -> gleam@otp@actor:continue(New_state) end; {ok, {stop, Reason}} -> {stop, Reason}; {error, Reason@1} -> logging:log( error, <<"Caught error in user handler: "/utf8, (gleam@string:inspect(Reason@1))/binary>> ), gleam@otp@actor:continue(State) end; {data, {binary_frame, _, Data@1}} -> Res@1 = gleam_erlang_ffi:rescue( fun() -> (erlang:element(5, Builder))( {binary, Data@1}, erlang:element(6, State), Conn ) end ), case Res@1 of {ok, {continue, User_state@1, User_selector@2}} -> New_state@1 = erlang:setelement(6, State, User_state@1), case User_selector@2 of {some, User_selector@3} -> Selector@1 = begin _pipe@2 = User_selector@3, _pipe@3 = gleam_erlang_ffi:map_selector( _pipe@2, fun(Field@0) -> {user_message, Field@0} end ), gleam_erlang_ffi:merge_selector( _pipe@3, gleam_erlang_ffi:map_selector( stratus@internal@socket:selector(), fun from_socket_message/1 ) ) end, {continue, New_state@1, {some, Selector@1}}; _ -> gleam@otp@actor:continue(New_state@1) end; {ok, {stop, Reason@2}} -> {stop, Reason@2}; {error, Reason@3} -> logging:log( error, <<"Caught error in user handler: "/utf8, (gleam@string:inspect(Reason@3))/binary>> ), gleam@otp@actor:continue(State) end; {control, {ping_frame, Payload, Payload_length}} -> Frame@1 = gramps@websocket:frame_to_bytes_builder( {control, {pong_frame, Payload, Payload_length}}, {some, <<0:4/unit:8>>} ), _ = stratus@internal@transport:send( erlang:element(3, Conn), erlang:element(2, Conn), Frame@1 ), gleam@otp@actor:continue(State); {control, {pong_frame, _, _}} -> gleam@otp@actor:continue(State); {control, {close_frame, Length, Payload@1}} -> Size = Length - 2, case Payload@1 of <<_:2/integer-unit:8, Message:Size/binary>> -> Msg = <<"WebSocket closing: "/utf8, (gleam@string:inspect(Message))/binary>>, logging:log(debug, Msg); _ -> nil end, (erlang:element(6, Builder))(erlang:element(6, State)), {stop, normal}; {continuation, _, _} -> gleam@otp@actor:continue(State) end. -spec send_message(gleam@erlang@process:subject(internal_message(KOK)), KOK) -> nil. send_message(Subject, Message) -> gleam@erlang@process:send(Subject, {user_message, Message}). -spec send_text_message(connection(), binary()) -> {ok, nil} | {error, stratus@internal@socket:socket_reason()}. send_text_message(Conn, Msg) -> Frame = gramps@websocket:to_text_frame( Msg, none, {some, crypto:strong_rand_bytes(4)} ), stratus@internal@transport:send( erlang:element(3, Conn), erlang:element(2, Conn), Frame ). -spec send_binary_message(connection(), bitstring()) -> {ok, nil} | {error, stratus@internal@socket:socket_reason()}. send_binary_message(Conn, Msg) -> Frame = gramps@websocket:to_binary_frame( Msg, none, {some, crypto:strong_rand_bytes(4)} ), stratus@internal@transport:send( erlang:element(3, Conn), erlang:element(2, Conn), Frame ). -spec send_ping(connection(), bitstring()) -> {ok, nil} | {error, stratus@internal@socket:socket_reason()}. send_ping(Conn, Data) -> Size = erlang:byte_size(Data), Mask = case Size of 0 -> <<0:4>>; _ -> crypto:strong_rand_bytes(4) end, Frame = gramps@websocket:frame_to_bytes_builder( {control, {ping_frame, Size, Data}}, {some, Mask} ), stratus@internal@transport:send( erlang:element(3, Conn), erlang:element(2, Conn), Frame ). -spec close(connection()) -> {ok, nil} | {error, stratus@internal@socket:socket_reason()}. close(Conn) -> Frame = gramps@websocket:frame_to_bytes_builder( {control, {close_frame, 0, <<>>}}, {some, crypto:strong_rand_bytes(4)} ), stratus@internal@transport:send( erlang:element(3, Conn), erlang:element(2, Conn), Frame ). -spec make_upgrade(gleam@http@request:request(binary()), binary()) -> gleam@bytes_builder:bytes_builder(). make_upgrade(Req, Origin) -> User_headers = begin _pipe = erlang:element(3, Req), _pipe@1 = gleam@list:filter( _pipe, fun(Pair) -> {Key, _} = Pair, (((((Key /= <<"host"/utf8>>) andalso (Key /= <<"upgrade"/utf8>>)) andalso (Key /= <<"connection"/utf8>>)) andalso (Key /= <<"sec-websocket-key"/utf8>>)) andalso (Key /= <<"sec-websocket-version"/utf8>>)) andalso (Key /= <<"origin"/utf8>>) end ), _pipe@2 = gleam@list:map( _pipe@1, fun(Pair@1) -> {Key@1, Value} = Pair@1, <<<>/binary, Value/binary>> end ), gleam@string:join(_pipe@2, <<"\r\n"/utf8>>) end, Path@1 = case erlang:element(8, Req) of <<""/utf8>> -> <<"/"/utf8>>; Path -> Path end, Query = begin _pipe@3 = Req, _pipe@4 = gleam@http@request:get_query(_pipe@3), _pipe@5 = gleam@result:map(_pipe@4, fun gleam@uri:query_to_string/1), (fun(Str) -> case Str of {ok, <<""/utf8>>} -> <<""/utf8>>; {ok, Str@1} -> <<"?"/utf8, Str@1/binary>>; _ -> <<""/utf8>> end end)(_pipe@5) end, _pipe@6 = gleam@bytes_builder:new(), _pipe@7 = gleam@bytes_builder:append_string( _pipe@6, <<<<<<"GET "/utf8, Path@1/binary>>/binary, Query/binary>>/binary, " HTTP/1.1\r\n"/utf8>> ), _pipe@8 = gleam@bytes_builder:append_string( _pipe@7, <<<<"host: "/utf8, (erlang:element(6, Req))/binary>>/binary, "\r\n"/utf8>> ), _pipe@9 = gleam@bytes_builder:append_string( _pipe@8, <<"upgrade: websocket\r\n"/utf8>> ), _pipe@10 = gleam@bytes_builder:append_string( _pipe@9, <<"connection: Upgrade\r\n"/utf8>> ), _pipe@11 = gleam@bytes_builder:append_string( _pipe@10, <<<<"sec-websocket-key: "/utf8, (<<"dGhlIHNhbXBsZSBub25jZQ=="/utf8>>)/binary>>/binary, "\r\n"/utf8>> ), _pipe@12 = gleam@bytes_builder:append_string( _pipe@11, <<"sec-websocket-version: 13\r\n"/utf8>> ), _pipe@13 = gleam@bytes_builder:append_string( _pipe@12, <<<<"origin: "/utf8, Origin/binary>>/binary, "\r\n"/utf8>> ), _pipe@14 = gleam@bytes_builder:append_string( _pipe@13, <<"sec-websocket-extensions: permessage-deflate\r\n"/utf8>> ), _pipe@15 = gleam@bytes_builder:append_string(_pipe@14, User_headers), gleam@bytes_builder:append_string(_pipe@15, <<"\r\n"/utf8>>). -spec read_body( stratus@internal@transport:transport(), stratus@internal@socket:socket(), integer(), integer(), bitstring() ) -> {ok, {bitstring(), bitstring()}} | {error, stratus@internal@socket:socket_reason()}. read_body(Transport, Socket, Timeout, Length, Body) -> case Body of <> -> {ok, {Data, Rest}}; _ -> case stratus@internal@transport:receive_timeout( Transport, Socket, 0, Timeout ) of {ok, Data@1} -> read_body( Transport, Socket, Timeout, Length, <> ); {error, Reason} -> {error, Reason} end end. -spec perform_handshake( gleam@http@request:request(binary()), stratus@internal@transport:transport(), integer() ) -> {ok, {stratus@internal@socket:socket(), gleam@http@response:response(bitstring()), bitstring()}} | {error, handshake_error()}. perform_handshake(Req, Transport, Timeout) -> Certs = case erlang:element(5, Req) of https -> _assert_subject = stratus_ffi:ssl_start(), {ok, _} = case _assert_subject of {ok, _} -> _assert_subject; _assert_fail -> erlang:error(#{gleam_error => let_assert, message => <<"Assertion pattern match failed"/utf8>>, value => _assert_fail, module => <<"stratus"/utf8>>, function => <<"perform_handshake"/utf8>>, line => 577}) end, [{cacerts, public_key:cacerts_get()}, stratus_ffi:custom_sni_matcher()]; http -> [] end, Opts = stratus@internal@socket:convert_options( lists:append( [{packets_of, binary}, {send_timeout, 30000}, {send_timeout_close, true}, {reuseaddr, true}, {nodelay, true}], [{'receive', pull} | Certs] ) ), Port = gleam@option:lazy_unwrap( erlang:element(7, Req), fun() -> case Transport of ssl -> 443; tcp -> 80 end end ), Origin = case {erlang:element(5, Req), Port} of {https, 443} -> <<"https://"/utf8, (erlang:element(6, Req))/binary>>; {http, 80} -> <<"http://"/utf8, (erlang:element(6, Req))/binary>>; {https, _} -> <<<<<<"https://"/utf8, (erlang:element(6, Req))/binary>>/binary, ":"/utf8>>/binary, (gleam@int:to_string(Port))/binary>>; {_, _} -> <<<<<<"http://"/utf8, (erlang:element(6, Req))/binary>>/binary, ":"/utf8>>/binary, (gleam@int:to_string(Port))/binary>> end, logging:log( debug, <<<<<<"Making request to "/utf8, (erlang:element(6, Req))/binary>>/binary, " at "/utf8>>/binary, (gleam@int:to_string(Port))/binary>> ), gleam@result:'try'( gleam@result:map_error( stratus@internal@transport:connect( Transport, unicode:characters_to_list(erlang:element(6, Req)), Port, Opts, Timeout ), fun(Field@0) -> {sock, Field@0} end ), fun(Socket) -> gleam@result:'try'( gleam@result:map_error( stratus@internal@transport:send( Transport, Socket, make_upgrade(Req, Origin) ), fun(Field@0) -> {sock, Field@0} end ), fun(_) -> logging:log( debug, <<"Sent upgrade request, waiting "/utf8, (gleam@int:to_string(Timeout))/binary>> ), gleam@result:'try'( gleam@result:map_error( stratus@internal@transport:receive_timeout( Transport, Socket, 0, Timeout ), fun(Field@0) -> {sock, Field@0} end ), fun(Resp) -> _pipe = Resp, _pipe@1 = gramps@http:read_response(_pipe), _pipe@2 = gleam@result:map_error( _pipe@1, fun(_) -> {protocol, Resp} end ), _pipe@6 = gleam@result:then( _pipe@2, fun(Pair) -> {Resp@1, Body} = Pair, Body_size = begin _pipe@3 = erlang:element(3, Resp@1), _pipe@4 = gleam@list:key_find( _pipe@3, <<"content-length"/utf8>> ), _pipe@5 = gleam@result:then( _pipe@4, fun gleam@int:parse/1 ), gleam@result:unwrap(_pipe@5, 0) end, case read_body( Transport, Socket, Timeout, Body_size, Body ) of {ok, {Body@1, Rest}} -> {ok, {gleam@http@response:set_body( Resp@1, Body@1 ), Rest}}; {error, Reason} -> {error, {sock, Reason}} end end ), gleam@result:then( _pipe@6, fun(Pair@1) -> {Resp@2, Rest@1} = Pair@1, case erlang:element(2, Resp@2) of 101 -> {ok, {Socket, Resp@2, Rest@1}}; _ -> {error, {upgrade_failed, Resp@2}} end end ) end ) end ) end ). -spec initialize(builder(any(), KNS)) -> {ok, gleam@erlang@process:subject(internal_message(KNS))} | {error, gleam@otp@actor:start_error()}. initialize(Builder) -> Transport = case erlang:element(5, erlang:element(2, Builder)) of https -> ssl; _ -> tcp end, gleam@otp@actor:start_spec( {spec, fun() -> Subj = gleam@erlang@process:new_subject(), Started_selector = gleam@erlang@process:selecting( gleam_erlang_ffi:new_selector(), Subj, fun gleam@function:identity/1 ), logging:log(debug, <<"Calling user initializer"/utf8>>), {User_state, User_selector} = (erlang:element(4, Builder))(), Selector@1 = case User_selector of {some, Selector} -> _pipe = Selector, _pipe@1 = gleam_erlang_ffi:map_selector( _pipe, fun(Field@0) -> {user_message, Field@0} end ), _pipe@2 = gleam_erlang_ffi:merge_selector( _pipe@1, Started_selector ), gleam_erlang_ffi:merge_selector( _pipe@2, gleam_erlang_ffi:map_selector( stratus@internal@socket:selector(), fun from_socket_message/1 ) ); _ -> _pipe@3 = Started_selector, gleam_erlang_ffi:merge_selector( _pipe@3, gleam_erlang_ffi:map_selector( stratus@internal@socket:selector(), fun from_socket_message/1 ) ) end, gleam@erlang@process:send(Subj, started), {ready, {state, <<>>, none, Subj, none, User_state}, Selector@1} end, 1000, fun(Msg, State) -> case Msg of started -> logging:log( debug, <<"Attempting handshake to "/utf8, (gleam@uri:to_string( gleam@http@request:to_uri( erlang:element(2, Builder) ) ))/binary>> ), _pipe@4 = perform_handshake( erlang:element(2, Builder), Transport, erlang:element(3, Builder) ), _pipe@7 = gleam@result:then( _pipe@4, fun(Pair) -> logging:log( debug, <<"Handshake successful"/utf8>> ), _pipe@5 = stratus@internal@transport:set_opts( Transport, erlang:element(1, Pair), stratus@internal@socket:convert_options( [{'receive', once}] ) ), _pipe@6 = gleam@result:replace(_pipe@5, Pair), gleam@result:map_error( _pipe@6, fun(Field@0) -> {sock, Field@0} end ) end ), _pipe@8 = gleam@result:map( _pipe@7, fun(Pair@1) -> {Socket, _, Buffer} = Pair@1, logging:log( debug, <<"WebSocket process ready to start receiving"/utf8>> ), _ = case Buffer of <<>> -> nil; Data -> gleam@erlang@process:send( erlang:element(4, State), {data, Data} ) end, gleam@otp@actor:continue( erlang:setelement( 2, erlang:setelement( 5, State, {some, Socket} ), Buffer ) ) end ), _pipe@9 = gleam@result:map_error( _pipe@8, fun(Err) -> case Err of {protocol, _} -> Msg@1 = <<"Failed to connect to server: "/utf8, (gleam@string:inspect(Err))/binary>>, logging:log(error, Msg@1), {stop, {abnormal, Msg@1}}; {sock, _} -> Msg@1 = <<"Failed to connect to server: "/utf8, (gleam@string:inspect(Err))/binary>>, logging:log(error, Msg@1), {stop, {abnormal, Msg@1}}; {upgrade_failed, Resp} -> (erlang:element(7, Builder))(Resp), logging:log( error, <<"WebSocket handshake failed with status "/utf8, (gleam@int:to_string( erlang:element(2, Resp) ))/binary>> ), {stop, {abnormal, <<"WebSocket handshake failed"/utf8>>}} end end ), gleam@result:unwrap_both(_pipe@9); {user_message, User_message} -> _assert_subject = erlang:element(5, State), {some, Socket@1} = case _assert_subject of {some, _} -> _assert_subject; _assert_fail -> erlang:error(#{gleam_error => let_assert, message => <<"Assertion pattern match failed"/utf8>>, value => _assert_fail, module => <<"stratus"/utf8>>, function => <<"initialize"/utf8>>, line => 254}) end, Conn = {connection, Socket@1, Transport}, Res = gleam_erlang_ffi:rescue( fun() -> (erlang:element(5, Builder))( {user, User_message}, erlang:element(6, State), Conn ) end ), case Res of {ok, {continue, User_state@1, User_selector@1}} -> New_state = erlang:setelement( 6, State, User_state@1 ), case User_selector@1 of {some, User_selector@2} -> Selector@2 = begin _pipe@10 = User_selector@2, _pipe@11 = gleam_erlang_ffi:map_selector( _pipe@10, fun(Field@0) -> {user_message, Field@0} end ), gleam_erlang_ffi:merge_selector( _pipe@11, gleam_erlang_ffi:map_selector( stratus@internal@socket:selector( ), fun from_socket_message/1 ) ) end, {continue, New_state, {some, Selector@2}}; _ -> gleam@otp@actor:continue(New_state) end; {ok, {stop, Reason}} -> {stop, Reason}; {error, Reason@1} -> logging:log( error, <<"Caught error in user handler: "/utf8, (gleam@string:inspect(Reason@1))/binary>> ), gleam@otp@actor:continue(State) end; {err, Reason@2} -> {stop, {abnormal, gleam@string:inspect(Reason@2)}}; {data, Bits} -> _assert_subject@1 = erlang:element(5, State), {some, Socket@2} = case _assert_subject@1 of {some, _} -> _assert_subject@1; _assert_fail@1 -> erlang:error(#{gleam_error => let_assert, message => <<"Assertion pattern match failed"/utf8>>, value => _assert_fail@1, module => <<"stratus"/utf8>>, function => <<"initialize"/utf8>>, line => 292}) end, Conn@1 = {connection, Socket@2, Transport}, {Frames, Rest} = gramps@websocket:get_messages( gleam@bit_array:append( erlang:element(2, State), Bits ), [], none ), Frames@1 = gramps@websocket:aggregate_frames( Frames, erlang:element(3, State), [] ), _pipe@12 = case Frames@1 of {error, nil} -> gleam@otp@actor:continue(State); {ok, Frames@2} -> gleam@list:fold_until( Frames@2, gleam@otp@actor:continue(State), fun(Acc, Frame) -> {continue, Prev_state, _} = case Acc of {continue, _, _} -> Acc; _assert_fail@2 -> erlang:error( #{gleam_error => let_assert, message => <<"Assertion pattern match failed"/utf8>>, value => _assert_fail@2, module => <<"stratus"/utf8>>, function => <<"initialize"/utf8>>, line => 306} ) end, case handle_frame( Builder, Transport, Prev_state, Conn@1, Frame ) of {continue, _, _} = Next -> {continue, Next}; {stop, _} = Err@1 -> {stop, Err@1} end end ) end, (fun(Next@1) -> case Next@1 of {stop, _} = Stop -> Stop; {continue, State@1, Selector@3} -> _assert_subject@2 = stratus@internal@transport:set_opts( Transport, Socket@2, stratus@internal@socket:convert_options( [{'receive', once}] ) ), {ok, _} = case _assert_subject@2 of {ok, _} -> _assert_subject@2; _assert_fail@3 -> erlang:error( #{gleam_error => let_assert, message => <<"Assertion pattern match failed"/utf8>>, value => _assert_fail@3, module => <<"stratus"/utf8>>, function => <<"initialize"/utf8>>, line => 320} ) end, {continue, erlang:setelement(2, State@1, Rest), Selector@3} end end)(_pipe@12); closed -> logging:log(debug, <<"Received closed frame"/utf8>>), (erlang:element(6, Builder))(erlang:element(6, State)), {stop, normal}; shutdown -> logging:log(debug, <<"Received shutdown messag"/utf8>>), {stop, normal} end end} ).