-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, 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(KHH) :: started | {user_message, KHH} | {err, stratus@internal@socket:socket_reason()} | {data, bitstring()} | closed | shutdown. -type message(KHI) :: {text, binary()} | {binary, bitstring()} | {user, KHI}. -opaque builder(KHJ, KHK) :: {builder, gleam@http@request:request(binary()), integer(), fun(() -> {KHJ, gleam@option:option(gleam@erlang@process:selector(KHK))}), fun((message(KHK), KHJ, connection()) -> gleam@otp@actor:next(KHK, KHJ)), fun((KHJ) -> nil)}. -type state(KHL, KHM) :: {state, bitstring(), gleam@option:option(gramps:frame()), gleam@erlang@process:subject(internal_message(KHM)), gleam@option:option(stratus@internal@socket:socket()), KHL}. -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(() -> {KHQ, gleam@option:option(gleam@erlang@process:selector(KHR))}), fun((message(KHR), KHQ, connection()) -> gleam@otp@actor:next(KHR, KHQ)) ) -> builder(KHQ, KHR). websocket(Req, Init, Loop) -> {builder, Req, 5000, Init, Loop, fun(_) -> nil end}. -spec with_connect_timeout(builder(KHZ, KIA), integer()) -> builder(KHZ, KIA). with_connect_timeout(Builder, Timeout) -> erlang:setelement(3, Builder, Timeout). -spec on_close(builder(KIF, KIG), fun((KIF) -> nil)) -> builder(KIF, KIG). on_close(Builder, On_close) -> erlang:setelement(6, Builder, On_close). -spec handle_frame( builder(KIT, KIU), stratus@internal@transport:transport(), state(KIT, KIU), connection(), gramps:frame() ) -> gleam@otp@actor:next(internal_message(KIU), state(KIT, KIU)). 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 => 308}) 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 => 311}) end, case (erlang:element(5, Builder))( {text, Str}, erlang:element(6, State), Conn ) of {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; {stop, Reason} -> {stop, Reason} end; {data, {binary_frame, _, Data@1}} -> case (erlang:element(5, Builder))( {binary, Data@1}, erlang:element(6, State), Conn ) of {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; {stop, Reason@1} -> {stop, Reason@1} 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(info, 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(KJE)), KJE) -> 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 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(32)} ), 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 => 436}) 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 => 445}) 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, _pipe@3 = gleam@bytes_builder:new(), _pipe@4 = gleam@bytes_builder:append_string( _pipe@3, <<<<"GET "/utf8, Path@1/binary>>/binary, " HTTP/1.1\r\n"/utf8>> ), _pipe@5 = gleam@bytes_builder:append_string( _pipe@4, <<<<"Host: "/utf8, (erlang:element(6, Req))/binary>>/binary, "\r\n"/utf8>> ), _pipe@6 = gleam@bytes_builder:append_string( _pipe@5, <<"Upgrade: websocket\r\n"/utf8>> ), _pipe@7 = gleam@bytes_builder:append_string( _pipe@6, <<"Connection: Upgrade\r\n"/utf8>> ), _pipe@8 = gleam@bytes_builder:append_string( _pipe@7, <<<<"Sec-WebSocket-Key: "/utf8, (<<"dGhlIHNhbXBsZSBub25jZQ=="/utf8>>)/binary>>/binary, "\r\n"/utf8>> ), _pipe@9 = gleam@bytes_builder:append_string( _pipe@8, <<"Sec-WebSocket-Version: 13\r\n"/utf8>> ), _pipe@10 = gleam@bytes_builder:append_string( _pipe@9, <<<<"Origin: "/utf8, Origin/binary>>/binary, "\r\n"/utf8>> ), _pipe@11 = gleam@bytes_builder:append_string(_pipe@10, User_headers), gleam@bytes_builder:append_string(_pipe@11, <<"\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 => 481}) end, [{cacerts, public_key:cacerts_get()}, stratus_ffi:custom_sni_matcher()]; http -> [] end, Opts = stratus@internal@socket:convert_options( gleam@list: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( info, <<<<<<"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( info, <<"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(), KIM)) -> {ok, gleam@erlang@process:subject(internal_message(KIM))} | {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(info, <<"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) -> logging:log( info, <<"got a message: "/utf8, (gleam@string:inspect(Msg))/binary>> ), case Msg of started -> logging:log( info, <<"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( info, <<"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( info, <<"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}, case (erlang:element(5, Builder))( {user, User_message}, erlang:element(6, State), Conn ) of {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; {stop, Reason} -> {stop, Reason} end; {err, Reason@1} -> {stop, {abnormal, gleam@string:inspect(Reason@1)}}; {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 => 253}) 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 => 262} ) 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 => 276} ) end, {continue, erlang:setelement(2, State@1, Rest), Selector@3} end end)(_pipe@12); closed -> (erlang:element(6, Builder))(erlang:element(6, State)), {stop, normal}; shutdown -> {stop, normal} end end} ).