-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/1, handshake_error/0]). -opaque connection() :: {connection, stratus@internal@socket:socket(), stratus@internal@transport:transport()}. -opaque internal_message(KBY) :: started | {user_message, KBY} | {err, stratus@internal@socket:socket_reason()} | {data, bitstring()} | closed | shutdown. -type message(KBZ) :: {text, binary()} | {binary, bitstring()} | {user, KBZ}. -opaque builder(KCA, KCB) :: {builder, gleam@http@request:request(binary()), integer(), fun(() -> {KCA, gleam@option:option(gleam@erlang@process:selector(KCB))}), fun((message(KCB), KCA, connection()) -> gleam@otp@actor:next(KCB, KCA)), fun((KCA) -> nil)}. -type state(KCC) :: {state, bitstring(), gleam@option:option(gramps:data_frame()), gleam@option:option(stratus@internal@socket:socket()), KCC}. -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(() -> {KCG, gleam@option:option(gleam@erlang@process:selector(KCH))}), fun((message(KCH), KCG, connection()) -> gleam@otp@actor:next(KCH, KCG)) ) -> builder(KCG, KCH). websocket(Req, Init, Loop) -> {builder, Req, 5000, Init, Loop, fun(_) -> nil end}. -spec with_connect_timeout(builder(KCP, KCQ), integer()) -> builder(KCP, KCQ). with_connect_timeout(Builder, Timeout) -> erlang:setelement(3, Builder, Timeout). -spec on_close(builder(KCV, KCW), fun((KCV) -> nil)) -> builder(KCV, KCW). on_close(Builder, On_close) -> erlang:setelement(6, Builder, On_close). -spec send_message(gleam@erlang@process:subject(internal_message(KDJ)), KDJ) -> 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, <<>>} ), 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 => 408}) 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 => 417}) 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()} | {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 => 453}) end, [{cacerts, public_key:cacerts_get()}]; http -> [] end, Opts = stratus@internal@socket:convert_options( gleam@list:append( [{'receive', once}, {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>>/binary, " with opts: "/utf8>>/binary, (gleam@string:inspect(Opts))/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 Resp of <<"HTTP/1.1 101 Switching Protocols"/utf8, _/bitstring>> -> {ok, Socket}; _ -> {error, {protocol, Resp}} end end ) end ) end ). -spec initialize(builder(any(), KDC)) -> {ok, gleam@erlang@process:subject(internal_message(KDC))} | {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, 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@5 = gleam@result:'try'( _pipe@4, fun(Socket) -> logging:log( debug, <<"Handshake successful"/utf8>> ), case stratus@internal@transport:set_opts( Transport, Socket, stratus@internal@socket:convert_options( [{'receive', once}] ) ) of {ok, _} -> {ok, Socket}; {error, Reason} -> {error, {sock, Reason}} end end ), _pipe@6 = gleam@result:map( _pipe@5, fun(Socket@1) -> logging:log( debug, <<"WebSocket process ready to start receiving"/utf8>> ), gleam@otp@actor:continue( erlang:setelement( 4, State, {some, Socket@1} ) ) end ), _pipe@7 = gleam@result:map_error( _pipe@6, 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@7); {user_message, User_message} -> _assert_subject = erlang:element(4, State), {some, Socket@2} = 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 => 218}) end, Conn = {connection, Socket@2, Transport}, case (erlang:element(5, Builder))( {user, User_message}, erlang:element(5, State), Conn ) of {continue, User_state@1, User_selector@1} -> New_state = erlang:setelement( 5, State, User_state@1 ), case User_selector@1 of {some, User_selector@2} -> Selector@2 = begin _pipe@8 = User_selector@2, _pipe@9 = gleam_erlang_ffi:map_selector( _pipe@8, fun(Field@0) -> {user_message, Field@0} end ), gleam_erlang_ffi:merge_selector( _pipe@9, 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@1} -> {stop, Reason@1} end; {err, Reason@2} -> {stop, {abnormal, gleam@string:inspect(Reason@2)}}; {data, Bits} -> _assert_subject@1 = erlang:element(4, State), {some, Socket@3} = 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 => 245}) end, Conn@1 = {connection, Socket@3, Transport}, _pipe@10 = gramps:frame_from_message( gleam@bit_array:append( erlang:element(2, State), Bits ) ), _pipe@15 = gleam@result:map( _pipe@10, fun(Data) -> {Parsed_frame, Rest} = Data, case Parsed_frame of {complete, {data, {text_frame, _, Data@1}}} -> _assert_subject@2 = gleam@bit_array:to_string( Data@1 ), {ok, Str} = case _assert_subject@2 of {ok, _} -> _assert_subject@2; _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 => 253} ) end, case (erlang:element(5, Builder))( {text, Str}, erlang:element(5, State), Conn@1 ) of {continue, User_state@2, User_selector@3} -> _assert_subject@3 = stratus@internal@transport:set_opts( Transport, Socket@3, stratus@internal@socket:convert_options( [{'receive', once}] ) ), {ok, _} = case _assert_subject@3 of {ok, _} -> _assert_subject@3; _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 => 257} ) end, New_state@1 = erlang:setelement( 2, erlang:setelement( 5, State, User_state@2 ), Rest ), case User_selector@3 of {some, User_selector@4} -> Selector@3 = begin _pipe@11 = User_selector@4, _pipe@12 = gleam_erlang_ffi:map_selector( _pipe@11, fun(Field@0) -> {user_message, Field@0} end ), gleam_erlang_ffi:merge_selector( _pipe@12, gleam_erlang_ffi:map_selector( stratus@internal@socket:selector( ), fun from_socket_message/1 ) ) end, {continue, New_state@1, {some, Selector@3}}; _ -> gleam@otp@actor:continue( New_state@1 ) end; {stop, Reason@3} -> {stop, Reason@3} end; {complete, {data, {binary_frame, _, Data@2}}} -> case (erlang:element(5, Builder))( {binary, Data@2}, erlang:element(5, State), Conn@1 ) of {continue, User_state@3, User_selector@5} -> _assert_subject@4 = stratus@internal@transport:set_opts( Transport, Socket@3, stratus@internal@socket:convert_options( [{'receive', once}] ) ), {ok, _} = case _assert_subject@4 of {ok, _} -> _assert_subject@4; _assert_fail@4 -> erlang:error( #{gleam_error => let_assert, message => <<"Assertion pattern match failed"/utf8>>, value => _assert_fail@4, module => <<"stratus"/utf8>>, function => <<"initialize"/utf8>>, line => 286} ) end, New_state@2 = erlang:setelement( 2, erlang:setelement( 5, State, User_state@3 ), Rest ), case User_selector@5 of {some, User_selector@6} -> Selector@4 = begin _pipe@13 = User_selector@6, _pipe@14 = gleam_erlang_ffi:map_selector( _pipe@13, fun(Field@0) -> {user_message, Field@0} end ), gleam_erlang_ffi:merge_selector( _pipe@14, gleam_erlang_ffi:map_selector( stratus@internal@socket:selector( ), fun from_socket_message/1 ) ) end, {continue, New_state@2, {some, Selector@4}}; _ -> gleam@otp@actor:continue( New_state@2 ) end; {stop, Reason@4} -> {stop, Reason@4} end; {complete, {control, {ping_frame, Payload, Payload_length}}} -> Frame = gramps:frame_to_bytes_builder( {control, {pong_frame, Payload, Payload_length}}, {some, <<>>} ), _ = stratus@internal@transport:send( erlang:element(3, Conn@1), erlang:element(2, Conn@1), Frame ), gleam@otp@actor:continue(State); {complete, {control, {pong_frame, _, _}}} -> gleam@otp@actor:continue(State); {complete, {control, {close_frame, _, _}}} -> (erlang:element(6, Builder))( erlang:element(5, State) ), {stop, normal}; {incomplete, _} -> erlang:error(#{gleam_error => panic, message => <<"Incomplete messages not supported right now"/utf8>>, module => <<"stratus"/utf8>>, function => <<"initialize"/utf8>>, line => 328}); {complete, {continuation, _, _}} -> erlang:error(#{gleam_error => panic, message => <<"Incomplete messages not supported right now"/utf8>>, module => <<"stratus"/utf8>>, function => <<"initialize"/utf8>>, line => 328}) end end ), gleam@result:lazy_unwrap( _pipe@15, fun() -> _assert_subject@5 = stratus@internal@transport:set_opts( Transport, Socket@3, stratus@internal@socket:convert_options( [{'receive', once}] ) ), {ok, _} = case _assert_subject@5 of {ok, _} -> _assert_subject@5; _assert_fail@5 -> erlang:error( #{gleam_error => let_assert, message => <<"Assertion pattern match failed"/utf8>>, value => _assert_fail@5, module => <<"stratus"/utf8>>, function => <<"initialize"/utf8>>, line => 332} ) end, gleam@otp@actor:continue( erlang:setelement( 2, State, gleam@bit_array:append( erlang:element(2, State), Bits ) ) ) end ); closed -> (erlang:element(6, Builder))(erlang:element(5, State)), {stop, normal}; shutdown -> {stop, normal} end end} ).