-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(JBW) :: {ack, JBW} | {nack, JBW} | {term, JBW} | {no_reply, JBW}. -type pull_handler_state(JBX) :: {pull_handler_state, gleam@erlang@process:subject(glats:connection_message()), glats@jetstream@consumer:subscription(), integer(), integer(), fun((glats:message(), JBX) -> outcome(JBX)), JBX}. -spec request_more(pull_handler_state(JCE)) -> gleam@otp@actor:next(any(), pull_handler_state(JCE)). 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. -spec handle_pull_message(glats:message(), pull_handler_state(JCK)) -> gleam@otp@actor:next(any(), pull_handler_state(JCK)). 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. -spec pull_loop(glats:subscription_message(), pull_handler_state(JCH)) -> gleam@otp@actor:next(any(), pull_handler_state(JCH)). 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. -spec handle_pull_consumer( gleam@erlang@process:subject(glats:connection_message()), JCA, binary(), integer(), fun((glats:message(), JCA) -> outcome(JCA)), 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} ).