-module(stratus). -compile([no_auto_import, nowarn_unused_vars, nowarn_unused_function, nowarn_nomatch]). -export([websocket/3, with_connect_timeout/2, on_close/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(KKP) :: started | {user_message, KKP} | {err, stratus@internal@socket:socket_reason()} | {data, bitstring()} | closed | shutdown. -type message(KKQ) :: {text, binary()} | {binary, bitstring()} | {user, KKQ}. -opaque builder(KKR, KKS) :: {builder, gleam@http@request:request(binary()), integer(), fun(() -> {KKR, gleam@option:option(gleam@erlang@process:selector(KKS))}), fun((message(KKS), KKR, connection()) -> gleam@otp@actor:next(KKS, KKR)), fun((KKR) -> nil)}. -type state(KKT, KKU) :: {state, bitstring(), gleam@option:option(gramps:frame()), gleam@erlang@process:subject(internal_message(KKU)), gleam@option:option(stratus@internal@socket:socket()), KKT}. -type handshake_error() :: {sock, stratus@internal@socket:socket_reason()} | {protocol, 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(() -> {KKY, gleam@option:option(gleam@erlang@process:selector(KKZ))}), fun((message(KKZ), KKY, connection()) -> gleam@otp@actor:next(KKZ, KKY)) ) -> builder(KKY, KKZ). websocket(Req, Init, Loop) -> {builder, Req, 5000, Init, Loop, fun(_) -> nil end}. -spec with_connect_timeout(builder(KLH, KLI), integer()) -> builder(KLH, KLI). with_connect_timeout(Builder, Timeout) -> erlang:setelement(3, Builder, Timeout). -spec on_close(builder(KLN, KLO), fun((KLN) -> nil)) -> builder(KLN, KLO). on_close(Builder, On_close) -> erlang:setelement(6, Builder, On_close). -spec handle_frame( builder(KMB, KMC), stratus@internal@transport:transport(), state(KMB, KMC), connection(), gramps:frame() ) -> gleam@otp@actor:next(internal_message(KMC), state(KMB, KMC)). 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 => 321}) 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 => 324}) 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: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(KMM)), KMM) -> 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:to_text_frame(Msg, true), 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:to_binary_frame(Msg, true), 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: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: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, _} = case Pair of {_, _} -> Pair; _assert_fail -> erlang:error(#{gleam_error => let_assert, message => <<"Assertion pattern match failed"/utf8>>, value => _assert_fail, module => <<"stratus"/utf8>>, function => <<"make_upgrade"/utf8>>, line => 481}) end, (((((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} = case Pair@1 of {_, _} -> Pair@1; _assert_fail@1 -> erlang:error(#{gleam_error => let_assert, message => <<"Assertion pattern match failed"/utf8>>, value => _assert_fail@1, module => <<"stratus"/utf8>>, function => <<"make_upgrade"/utf8>>, line => 490}) end, <<<>/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, User_headers), gleam@bytes_builder:append_string(_pipe@14, <<"\r\n"/utf8>>). -spec perform_handshake( gleam@http@request:request(binary()), stratus@internal@transport:transport(), integer() ) -> {ok, {stratus@internal@socket:socket(), 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 => 538}) 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) -> case gramps:read_response(Resp) of {ok, {_, Rest}} -> {ok, {Socket, Rest}}; {error, _} -> {error, {protocol, Resp}} end end ) end ) end ). -spec initialize(builder(any(), KLU)) -> {ok, gleam@erlang@process:subject(internal_message(KLU))} | {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) -> Msg@1 = <<"Failed to connect to server: "/utf8, (gleam@string:inspect(Err))/binary>>, logging:log(error, Msg@1), {stop, {abnormal, Msg@1}} 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 => 226}) 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 => 264}) end, Conn@1 = {connection, Socket@2, Transport}, {Frames, Rest} = gramps:get_messages( gleam@bit_array:append( erlang:element(2, State), Bits ), [] ), Frames@1 = gramps: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 => 273} ) 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 => 287} ) 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} ).