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_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(),
gleam@option:option(gramps@websocket@compression:context())}.
-opaque internal_message(LVF) :: started |
{user_message, LVF} |
{err, stratus@internal@socket:socket_reason()} |
{data, bitstring()} |
closed |
shutdown.
-type message(LVG) :: {text, binary()} | {binary, bitstring()} | {user, LVG}.
-opaque builder(LVH, LVI) :: {builder,
gleam@http@request:request(binary()),
integer(),
fun(() -> {LVH, gleam@option:option(gleam@erlang@process:selector(LVI))}),
fun((message(LVI), LVH, connection()) -> gleam@otp@actor:next(LVI, LVH)),
fun((LVH) -> nil),
fun((gleam@http@response:response(bitstring())) -> nil)}.
-type state(LVJ, LVK) :: {state,
bitstring(),
gleam@option:option(gramps@websocket:frame()),
gleam@erlang@process:subject(internal_message(LVK)),
gleam@option:option(stratus@internal@socket:socket()),
LVJ,
gleam@option:option(gramps@websocket@compression:compression())}.
-type handshake_error() :: {sock, stratus@internal@socket:socket_reason()} |
{protocol, bitstring()} |
{upgrade_failed, gleam@http@response:response(bitstring())}.
-file("/home/alex/gleams/stratus/src/stratus.gleam", 41).
-spec from_socket_message(stratus@internal@socket:socket_message()) -> internal_message(any()).
from_socket_message(Msg) ->
case Msg of
{data, Bits} ->
{data, Bits};
{err, closed} ->
closed;
{err, Reason} ->
{err, Reason}
end.
-file("/home/alex/gleams/stratus/src/stratus.gleam", 84).
-spec websocket(
gleam@http@request:request(binary()),
fun(() -> {LVO, gleam@option:option(gleam@erlang@process:selector(LVP))}),
fun((message(LVP), LVO, connection()) -> gleam@otp@actor:next(LVP, LVO))
) -> builder(LVO, LVP).
websocket(Req, Init, Loop) ->
{builder, Req, 5000, Init, Loop, fun(_) -> nil end, fun(_) -> nil end}.
-file("/home/alex/gleams/stratus/src/stratus.gleam", 104).
-spec with_connect_timeout(builder(LVX, LVY), integer()) -> builder(LVX, LVY).
with_connect_timeout(Builder, Timeout) ->
erlang:setelement(3, Builder, Timeout).
-file("/home/alex/gleams/stratus/src/stratus.gleam", 117).
-spec on_close(builder(LWD, LWE), fun((LWD) -> nil)) -> builder(LWD, LWE).
on_close(Builder, On_close) ->
erlang:setelement(6, Builder, On_close).
-file("/home/alex/gleams/stratus/src/stratus.gleam", 126).
-spec on_handshake_error(
builder(LWJ, LWK),
fun((gleam@http@response:response(bitstring())) -> nil)
) -> builder(LWJ, LWK).
on_handshake_error(Builder, On_handshake_error) ->
erlang:setelement(7, Builder, On_handshake_error).
-file("/home/alex/gleams/stratus/src/stratus.gleam", 383).
-spec handle_frame(
builder(LWY, LWZ),
stratus@internal@transport:transport(),
state(LWY, LWZ),
connection(),
gramps@websocket:frame()
) -> gleam@otp@actor:next(internal_message(LWZ), state(LWY, LWZ)).
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 => <<"Pattern match failed, no pattern matched the value."/utf8>>,
value => _assert_fail,
module => <<"stratus"/utf8>>,
function => <<"handle_frame"/utf8>>,
line => 390})
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 => <<"Pattern match failed, no pattern matched the value."/utf8>>,
value => _assert_fail@1,
module => <<"stratus"/utf8>>,
function => <<"handle_frame"/utf8>>,
line => 393})
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 = case erlang:element(4, Conn) of
{some, Context} ->
gramps@websocket:compressed_frame_to_bytes_tree(
{control, {pong_frame, Payload, Payload_length}},
Context,
{some, <<0:4/unit:8>>}
);
none ->
gramps@websocket:frame_to_bytes_tree(
{control, {pong_frame, Payload, Payload_length}},
{some, <<0:4/unit:8>>}
)
end,
_ = 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.
-file("/home/alex/gleams/stratus/src/stratus.gleam", 503).
-spec send_message(gleam@erlang@process:subject(internal_message(LXJ)), LXJ) -> nil.
send_message(Subject, Message) ->
gleam@erlang@process:send(Subject, {user_message, Message}).
-file("/home/alex/gleams/stratus/src/stratus.gleam", 512).
-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
).
-file("/home/alex/gleams/stratus/src/stratus.gleam", 522).
-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
).
-file("/home/alex/gleams/stratus/src/stratus.gleam", 532).
-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 = case erlang:element(4, Conn) of
{some, Context} ->
gramps@websocket:compressed_frame_to_bytes_tree(
{control, {ping_frame, Size, Data}},
Context,
{some, Mask}
);
none ->
gramps@websocket:frame_to_bytes_tree(
{control, {ping_frame, Size, Data}},
{some, Mask}
)
end,
stratus@internal@transport:send(
erlang:element(3, Conn),
erlang:element(2, Conn),
Frame
).
-file("/home/alex/gleams/stratus/src/stratus.gleam", 555).
-spec close(connection()) -> {ok, nil} |
{error, stratus@internal@socket:socket_reason()}.
close(Conn) ->
Frame = case erlang:element(4, Conn) of
{some, Context} ->
gramps@websocket:compressed_frame_to_bytes_tree(
{control, {close_frame, 0, <<>>}},
Context,
{some, crypto:strong_rand_bytes(4)}
);
none ->
gramps@websocket:frame_to_bytes_tree(
{control, {close_frame, 0, <<>>}},
{some, crypto:strong_rand_bytes(4)}
)
end,
stratus@internal@transport:send(
erlang:element(3, Conn),
erlang:element(2, Conn),
Frame
).
-file("/home/alex/gleams/stratus/src/stratus.gleam", 572).
-spec make_upgrade(gleam@http@request:request(binary())) -> gleam@bytes_tree:bytes_tree().
make_upgrade(Req) ->
User_headers = case erlang:element(3, Req) of
[] ->
<<""/utf8>>;
_ ->
_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>>)
end
),
_pipe@2 = gleam@list:map(
_pipe@1,
fun(Pair@1) ->
{Key@1, Value} = Pair@1,
<<<<Key@1/binary, ": "/utf8>>/binary, Value/binary>>
end
),
_pipe@3 = gleam@string:join(_pipe@2, <<"\r\n"/utf8>>),
gleam@string:append(_pipe@3, <<"\r\n"/utf8>>)
end,
Path@1 = case erlang:element(8, Req) of
<<""/utf8>> ->
<<"/"/utf8>>;
Path ->
Path
end,
Query = begin
_pipe@4 = Req,
_pipe@5 = gleam@http@request:get_query(_pipe@4),
_pipe@6 = gleam@result:map(_pipe@5, 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@6)
end,
_pipe@7 = gleam@bytes_tree:new(),
_pipe@8 = gleam@bytes_tree:append_string(
_pipe@7,
<<<<<<"GET "/utf8, Path@1/binary>>/binary, Query/binary>>/binary,
" HTTP/1.1\r\n"/utf8>>
),
_pipe@9 = gleam@bytes_tree:append_string(
_pipe@8,
<<<<"host: "/utf8, (erlang:element(6, Req))/binary>>/binary,
"\r\n"/utf8>>
),
_pipe@10 = gleam@bytes_tree:append_string(
_pipe@9,
<<"upgrade: websocket\r\n"/utf8>>
),
_pipe@11 = gleam@bytes_tree:append_string(
_pipe@10,
<<"connection: upgrade\r\n"/utf8>>
),
_pipe@12 = gleam@bytes_tree:append_string(
_pipe@11,
<<<<"sec-websocket-key: "/utf8,
(<<"dGhlIHNhbXBsZSBub25jZQ=="/utf8>>)/binary>>/binary,
"\r\n"/utf8>>
),
_pipe@13 = gleam@bytes_tree:append_string(
_pipe@12,
<<"sec-websocket-version: 13\r\n"/utf8>>
),
_pipe@14 = gleam@bytes_tree:append_string(
_pipe@13,
<<"sec-websocket-extensions: permessage-deflate\r\n"/utf8>>
),
_pipe@15 = gleam@bytes_tree:append_string(_pipe@14, User_headers),
gleam@bytes_tree:append_string(_pipe@15, <<"\r\n"/utf8>>).
-file("/home/alex/gleams/stratus/src/stratus.gleam", 717).
-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
<<Data:Length/binary, Rest/bitstring>> ->
{ok, {Data, Rest}};
_ ->
case stratus@internal@transport:receive_timeout(
Transport,
Socket,
0,
Timeout
) of
{ok, Data@1} ->
read_body(
Transport,
Socket,
Timeout,
Length,
<<Body/bitstring, Data@1/bitstring>>
);
{error, Reason} ->
{error, Reason}
end
end.
-file("/home/alex/gleams/stratus/src/stratus.gleam", 632).
-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 => <<"Pattern match failed, no pattern matched the value."/utf8>>,
value => _assert_fail,
module => <<"stratus"/utf8>>,
function => <<"perform_handshake"/utf8>>,
line => 639})
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
),
logging:log(
debug,
<<<<<<"Making request to "/utf8, (erlang:element(6, Req))/binary>>/binary,
" at "/utf8>>/binary,
(erlang:integer_to_binary(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) ->
Upgrade_req = make_upgrade(Req),
gleam@result:'try'(
gleam@result:map_error(
stratus@internal@transport:send(
Transport,
Socket,
Upgrade_req
),
fun(Field@0) -> {sock, Field@0} end
),
fun(_) ->
logging:log(
debug,
<<"Sent upgrade request, waiting "/utf8,
(erlang:integer_to_binary(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_stdlib:parse_int/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
).
-file("/home/alex/gleams/stratus/src/stratus.gleam", 737).
-spec close_contexts(
gleam@option:option(gramps@websocket@compression:compression())
) -> nil.
close_contexts(Contexts) ->
case Contexts of
{some, Compression} ->
zlib:close(erlang:element(3, Compression)),
zlib:close(erlang:element(2, Compression)),
nil;
_ ->
nil
end.
-file("/home/alex/gleams/stratus/src/stratus.gleam", 154).
-spec initialize(builder(any(), LWR)) -> {ok,
gleam@erlang@process:subject(internal_message(LWR))} |
{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, none},
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@11 = gleam@result:map(
_pipe@7,
fun(Pair@1) ->
{Socket, Resp, 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,
Extensions = begin
_pipe@8 = Resp,
_pipe@9 = gleam@http@response:get_header(
_pipe@8,
<<"sec-websocket-extensions"/utf8>>
),
_pipe@10 = gleam@result:map(
_pipe@9,
fun(_capture) ->
gleam@string:split(
_capture,
<<";"/utf8>>
)
end
),
gleam@result:unwrap(_pipe@10, [])
end,
Context = case gramps@websocket:has_deflate(
Extensions
) of
true ->
{some,
gramps@websocket@compression:init()};
false ->
none
end,
gleam@otp@actor:continue(
erlang:setelement(
7,
erlang:setelement(
2,
erlang:setelement(
5,
State,
{some, Socket}
),
Buffer
),
Context
)
)
end
),
_pipe@12 = gleam@result:map_error(
_pipe@11,
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@1} ->
(erlang:element(7, Builder))(Resp@1),
logging:log(
error,
<<"WebSocket handshake failed with status "/utf8,
(erlang:integer_to_binary(
erlang:element(2, Resp@1)
))/binary>>
),
{stop,
{abnormal,
<<"WebSocket handshake failed"/utf8>>}}
end end
),
gleam@result:unwrap_both(_pipe@12);
{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 => <<"Pattern match failed, no pattern matched the value."/utf8>>,
value => _assert_fail,
module => <<"stratus"/utf8>>,
function => <<"initialize"/utf8>>,
line => 274})
end,
Conn = {connection,
Socket@1,
Transport,
gleam@option:map(
erlang:element(7, State),
fun(Context@1) ->
erlang:element(3, Context@1)
end
)},
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@13 = User_selector@2,
_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,
{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} ->
close_contexts(erlang:element(7, State)),
{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 => <<"Pattern match failed, no pattern matched the value."/utf8>>,
value => _assert_fail@1,
module => <<"stratus"/utf8>>,
function => <<"initialize"/utf8>>,
line => 318})
end,
Conn@1 = {connection,
Socket@2,
Transport,
gleam@option:map(
erlang:element(7, State),
fun(Context@2) ->
erlang:element(3, Context@2)
end
)},
{Frames, Rest} = gramps@websocket:get_messages(
gleam@bit_array:append(
erlang:element(2, State),
Bits
),
[],
gleam@option:map(
erlang:element(7, State),
fun(Context@3) ->
erlang:element(2, Context@3)
end
)
),
Frames@1 = gramps@websocket:aggregate_frames(
Frames,
erlang:element(3, State),
[]
),
_pipe@15 = 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 => <<"Pattern match failed, no pattern matched the value."/utf8>>,
value => _assert_fail@2,
module => <<"stratus"/utf8>>,
function => <<"initialize"/utf8>>,
line => 337}
)
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 ->
close_contexts(erlang:element(7, State)),
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 => <<"Pattern match failed, no pattern matched the value."/utf8>>,
value => _assert_fail@3,
module => <<"stratus"/utf8>>,
function => <<"initialize"/utf8>>,
line => 354}
)
end,
{continue,
erlang:setelement(2, State@1, Rest),
Selector@3}
end end)(_pipe@15);
closed ->
logging:log(debug, <<"Received closed frame"/utf8>>),
(erlang:element(6, Builder))(erlang:element(6, State)),
close_contexts(erlang:element(7, State)),
{stop, normal};
shutdown ->
logging:log(debug, <<"Received shutdown messag"/utf8>>),
close_contexts(erlang:element(7, State)),
{stop, normal}
end end}
).