Current section

Files

Jump to
stratus src stratus.erl
Raw

src/stratus.erl

-module(stratus).
-compile([no_auto_import, nowarn_unused_vars, nowarn_unused_function, nowarn_nomatch]).
-export([websocket/3, with_init_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(KBO) :: {user_message, KBO} |
{err, stratus@internal@socket:socket_reason()} |
{data, bitstring()} |
closed |
shutdown.
-type message(KBP) :: {text, binary()} | {binary, bitstring()} | {user, KBP}.
-opaque builder(KBQ, KBR) :: {builder,
gleam@http@request:request(binary()),
gleam@option:option(integer()),
fun(() -> {KBQ, gleam@option:option(gleam@erlang@process:selector(KBR))}),
fun((message(KBR), KBQ, connection()) -> gleam@otp@actor:next(KBR, KBQ)),
fun((KBQ) -> nil)}.
-type state(KBS) :: {state,
bitstring(),
gleam@option:option(gramps:data_frame()),
stratus@internal@socket:socket(),
KBS}.
-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(() -> {KBW, gleam@option:option(gleam@erlang@process:selector(KBX))}),
fun((message(KBX), KBW, connection()) -> gleam@otp@actor:next(KBX, KBW))
) -> builder(KBW, KBX).
websocket(Req, Init, Loop) ->
{builder, Req, none, Init, Loop, fun(_) -> nil end}.
-spec with_init_timeout(builder(KCF, KCG), integer()) -> builder(KCF, KCG).
with_init_timeout(Builder, Timeout) ->
erlang:setelement(3, Builder, {some, Timeout}).
-spec on_close(builder(KCL, KCM), fun((KCL) -> nil)) -> builder(KCL, KCM).
on_close(Builder, On_close) ->
erlang:setelement(6, Builder, On_close).
-spec send_message(gleam@erlang@process:subject(internal_message(KCZ)), KCZ) -> 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 => 370})
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 => 379})
end,
<<<<Key@1/binary, ": "/utf8>>/binary, Value/binary>>
end
),
gleam@string:join(_pipe@2, <<"\r\n"/utf8>>)
end,
_pipe@3 = gleam@bytes_builder:new(),
_pipe@4 = gleam@bytes_builder:append_string(
_pipe@3,
<<<<"GET "/utf8, (erlang:element(8, Req))/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 => 410})
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,
gleam@result:'try'(
gleam@result:map_error(
stratus@internal@transport:connect(
Transport,
unicode:characters_to_list(erlang:element(6, Req)),
Port,
Opts
),
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(_) ->
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(), KCS)) -> {ok,
gleam@erlang@process:subject(internal_message(KCS))} |
{error, gleam@otp@actor:start_error()}.
initialize(Builder) ->
Transport = case erlang:element(5, erlang:element(2, Builder)) of
https ->
ssl;
_ ->
tcp
end,
Timeout = gleam@option:unwrap(erlang:element(3, Builder), 5000),
gleam@otp@actor:start_spec(
{spec,
fun() ->
_pipe = perform_handshake(
erlang:element(2, Builder),
Transport,
Timeout
),
_pipe@1 = gleam@result:'try'(
_pipe,
fun(Socket) ->
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@4 = gleam@result:map(
_pipe@1,
fun(Socket@1) ->
{User_state, User_selector} = (erlang:element(
4,
Builder
))(),
Selector@1 = case User_selector of
{some, Selector} ->
_pipe@2 = Selector,
_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
)
);
_ ->
gleam_erlang_ffi:map_selector(
stratus@internal@socket:selector(),
fun from_socket_message/1
)
end,
{ready,
{state, <<>>, none, Socket@1, User_state},
Selector@1}
end
),
_pipe@5 = gleam@result:map_error(
_pipe@4,
fun(Err) -> {failed, gleam@string:inspect(Err)} end
),
gleam@result:unwrap_both(_pipe@5)
end,
Timeout,
fun(Msg, State) ->
Conn = {connection, erlang:element(4, State), Transport},
case Msg of
{user_message, User_message} ->
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@6 = User_selector@2,
_pipe@7 = gleam_erlang_ffi:map_selector(
_pipe@6,
fun(Field@0) -> {user_message, Field@0} end
),
gleam_erlang_ffi:merge_selector(
_pipe@7,
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} ->
_pipe@8 = gramps:frame_from_message(
gleam@bit_array:append(
erlang:element(2, State),
Bits
)
),
_pipe@13 = gleam@result:map(
_pipe@8,
fun(Data) ->
{Parsed_frame, Rest} = Data,
case Parsed_frame of
{complete, {data, {text_frame, _, Data@1}}} ->
_assert_subject = gleam@bit_array:to_string(
Data@1
),
{ok, Str} = 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 => <<"initialize"/utf8>>,
line => 215}
)
end,
case (erlang:element(5, Builder))(
{text, Str},
erlang:element(5, State),
Conn
) of
{continue,
User_state@2,
User_selector@3} ->
_assert_subject@1 = stratus@internal@transport:set_opts(
Transport,
erlang:element(4, State),
stratus@internal@socket:convert_options(
[{'receive', once}]
)
),
{ok, _} = 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 => <<"initialize"/utf8>>,
line => 219}
)
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@9 = User_selector@4,
_pipe@10 = gleam_erlang_ffi:map_selector(
_pipe@9,
fun(Field@0) -> {user_message, Field@0} end
),
gleam_erlang_ffi:merge_selector(
_pipe@10,
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
) of
{continue,
User_state@3,
User_selector@5} ->
_assert_subject@2 = stratus@internal@transport:set_opts(
Transport,
erlang:element(4, State),
stratus@internal@socket:convert_options(
[{'receive', once}]
)
),
{ok, _} = 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 => 248}
)
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@11 = User_selector@6,
_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@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),
erlang:element(2, Conn),
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 => 290});
{complete, {continuation, _, _}} ->
erlang:error(#{gleam_error => panic,
message => <<"Incomplete messages not supported right now"/utf8>>,
module => <<"stratus"/utf8>>,
function => <<"initialize"/utf8>>,
line => 290})
end
end
),
gleam@result:lazy_unwrap(
_pipe@13,
fun() ->
_assert_subject@3 = stratus@internal@transport:set_opts(
Transport,
erlang:element(4, State),
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 => 294}
)
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}
).