Current section
Files
Jump to
Current section
Files
src/glats@jetstream@handler.erl
-module(glats@jetstream@handler).
-compile([no_auto_import, nowarn_unused_vars, nowarn_unused_function, nowarn_nomatch]).
-export([handle_pull_consumer/6]).
-export_type([outcome/1, pull_handler_state/1]).
-type outcome(JLD) :: {ack, JLD} | {nack, JLD} | {term, JLD} | {no_reply, JLD}.
-type pull_handler_state(JLE) :: {pull_handler_state,
gleam@erlang@process:subject(glats:connection_message()),
glats@jetstream@consumer:subscription(),
integer(),
integer(),
fun((glats:message(), JLE) -> outcome(JLE)),
JLE}.
-file("/home/arnar/Code/glats/src/glats/jetstream/handler.gleam", 187).
-spec request_more(pull_handler_state(JLL)) -> gleam@otp@actor:next(any(), pull_handler_state(JLL)).
request_more(State) ->
case glats@jetstream@consumer:request_batch(
erlang:element(3, State),
[{batch, erlang:element(4, State)}, {expires, 10000 * 1000000}]
) of
{ok, nil} ->
{continue,
erlang:setelement(5, State, erlang:element(4, State)),
none};
{error, Err} ->
{stop, {abnormal, gleam@string:inspect(Err)}}
end.
-file("/home/arnar/Code/glats/src/glats/jetstream/handler.gleam", 211).
-spec handle_pull_message(glats:message(), pull_handler_state(JLR)) -> gleam@otp@actor:next(any(), pull_handler_state(JLR)).
handle_pull_message(Message, State) ->
Inner@4 = case (erlang:element(6, State))(Message, erlang:element(7, State)) of
{ack, Inner} ->
glats@jetstream:ack(erlang:element(2, State), Message),
Inner;
{nack, Inner@1} ->
glats@jetstream:nack(erlang:element(2, State), Message),
Inner@1;
{term, Inner@2} ->
glats@jetstream:term(erlang:element(2, State), Message),
Inner@2;
{no_reply, Inner@3} ->
Inner@3
end,
case erlang:element(5, State) =< 1 of
true ->
request_more(erlang:setelement(7, State, Inner@4));
false ->
{continue,
erlang:setelement(
7,
erlang:setelement(5, State, erlang:element(5, State) - 1),
Inner@4
),
none}
end.
-file("/home/arnar/Code/glats/src/glats/jetstream/handler.gleam", 200).
-spec pull_loop(glats:subscription_message(), pull_handler_state(JLO)) -> gleam@otp@actor:next(any(), pull_handler_state(JLO)).
pull_loop(Message, State) ->
case Message of
{received_message, _, _, {some, 408}, _} ->
request_more(State);
{received_message, _, _, {some, 404}, _} ->
request_more(State);
{received_message, _, _, _, Msg} ->
handle_pull_message(Msg, State)
end.
-file("/home/arnar/Code/glats/src/glats/jetstream/handler.gleam", 139).
-spec handle_pull_consumer(
gleam@erlang@process:subject(glats:connection_message()),
JLH,
binary(),
integer(),
fun((glats:message(), JLH) -> outcome(JLH)),
list(glats@jetstream@consumer:subscription_option())
) -> {ok, gleam@erlang@process:subject(glats:subscription_message())} |
{error, gleam@otp@actor:start_error()}.
handle_pull_consumer(Conn, Initial_state, Topic, Batch_size, Handler, Opts) ->
gleam@otp@actor:start_spec(
{spec,
fun() ->
Subject = gleam@erlang@process:new_subject(),
Selector = begin
_pipe = gleam_erlang_ffi:new_selector(),
gleam@erlang@process:selecting(
_pipe,
Subject,
fun gleam@function:identity/1
)
end,
case glats@jetstream@consumer:subscribe(
Conn,
Subject,
Topic,
Opts
) of
{ok, Sub} ->
case glats@jetstream@consumer:request_batch(
Sub,
[{batch, Batch_size}, {expires, 10000 * 1000000}]
) of
{ok, nil} ->
{ready,
{pull_handler_state,
Conn,
Sub,
Batch_size,
Batch_size,
Handler,
Initial_state},
Selector};
{error, Err} ->
{failed, gleam@string:inspect(Err)}
end;
{error, Err@1} ->
{failed, gleam@string:inspect(Err@1)}
end
end,
5000,
fun pull_loop/2}
).