-module(glats@jetstream@handler). -compile([no_auto_import, nowarn_unused_vars]). -export([handle_pull_consumer/6]). -export_type([outcome/1, pull_handler_state/1]). -type outcome(IGI) :: {ack, IGI} | {nack, IGI} | {term, IGI} | {no_reply, IGI}. -type pull_handler_state(IGJ) :: {pull_handler_state, gleam@erlang@process:subject(glats:connection_message()), glats@jetstream@consumer:subscription(), integer(), integer(), fun((glats:message(), IGJ) -> outcome(IGJ)), IGJ}. -spec request_more(pull_handler_state(IGQ)) -> gleam@otp@actor:next(pull_handler_state(IGQ)). 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))}; {error, Err} -> {stop, {abnormal, gleam@string:inspect(Err)}} end. -spec handle_pull_message(glats:message(), pull_handler_state(IGW)) -> gleam@otp@actor:next(pull_handler_state(IGW)). 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 )} end. -spec pull_loop(glats:subscription_message(), pull_handler_state(IGT)) -> gleam@otp@actor:next(pull_handler_state(IGT)). 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()), IGM, binary(), integer(), fun((glats:message(), IGM) -> outcome(IGM)), 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} ).