-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]). -if(?OTP_RELEASE >= 27). -define(MODULEDOC(Str), -moduledoc(Str)). -define(DOC(Str), -doc(Str)). -else. -define(MODULEDOC(Str), -compile([])). -define(DOC(Str), -compile([])). -endif. ?MODULEDOC( " A convenience handler that will handle a consumer subscription for you.\n" " For every message it receives it will call the provided `SubscriptionHandler(a)`\n" " function and take action depending on its return value.\n" "\n" " It will also keep the state for you of type `a`.\n" "\n" " ## Pull consumer example\n" "\n" " ```gleam\n" " import gleam/io\n" " import gleam/int\n" " import gleam/string\n" " import gleam/result\n" " import gleam/function\n" " import gleam/erlang/process\n" " import glats.{Connection, Message}\n" " import glats/jetstream/stream.{Retention, WorkQueuePolicy}\n" " import glats/jetstream/consumer.{\n" " AckExplicit, AckPolicy, BindStream, Description, With,\n" " }\n" " import glats/jetstream/handler.{Ack}\n" " \n" " pub fn main() {\n" " use conn <- result.then(glats.connect(\"localhost\", 4222, []))\n" " \n" " // Create a stream\n" " let assert Ok(stream) =\n" " stream.create(conn, \"wqstream\", [\"ticket.>\"], [Retention(WorkQueuePolicy)])\n" " \n" " // Run pull handler\n" " let assert Ok(_actor) =\n" " handler.handle_pull_consumer(\n" " conn,\n" " 0, // Initial state\n" " \"ticket.*\", // Topic\n" " 100, // Batch size\n" " pull_handler, // Handler function\n" " [\n" " // Bind to stream created above\n" " BindStream(stream.config.name),\n" " // Set description for the ephemeral consumer\n" " With(Description(\"An ephemeral consumer for subscription\")),\n" " // Set ack policy for the consumer\n" " With(AckPolicy(AckExplicit)),\n" " ],\n" " )\n" " \n" " // Run a loop that publishes a message every 100ms\n" " publish_loop(conn, 0)\n" " \n" " Ok(Nil)\n" " }\n" " \n" " // Publishes a new message every 100ms\n" " fn publish_loop(conn: Connection, counter: Int) {\n" " let assert Ok(Nil) =\n" " glats.publish(\n" " conn,\n" " \"ticket.\" <> int.to_string(counter),\n" " \"ticket body\",\n" " [],\n" " )\n" " \n" " process.sleep(100)\n" " \n" " publish_loop(conn, counter + 1)\n" " }\n" " \n" " // Handler function for the pull consumer handler\n" " pub fn pull_handler(message: Message, state) {\n" " // Increment state counter, print message and instruct\n" " // pull handler to ack the message.\n" " state + 1\n" " |> function.tap(print_message(_, message.topic, message.body))\n" " |> Ack\n" " }\n" "\n" " fn print_message(num: Int, topic: String, body: String) {\n" " \"message \" <> int.to_string(num) <> \" (\" <> topic <> \"): \" <> body\n" " |> io.println\n" " }\n" " ```\n" "\n" " Will output:\n" "\n" " ```sh\n" " message 1 (ticket.0): ticket body\n" " message 2 (ticket.1): ticket body\n" " message 3 (ticket.2): ticket body\n" " message 4 (ticket.3): ticket body\n" " message 5 (ticket.4): ticket body\n" " ...\n" " ```\n" ). -type outcome(IAD) :: {ack, IAD} | {nack, IAD} | {term, IAD} | {no_reply, IAD}. -type pull_handler_state(IAE) :: {pull_handler_state, gleam@erlang@process:subject(glats:connection_message()), glats@jetstream@consumer:subscription(), integer(), integer(), fun((glats:message(), IAE) -> outcome(IAE)), IAE}. -file("src/glats/jetstream/handler.gleam", 187). -spec request_more(pull_handler_state(IAL)) -> gleam@otp@actor:next(any(), pull_handler_state(IAL)). 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, begin _record = State, {pull_handler_state, erlang:element(2, _record), erlang:element(3, _record), erlang:element(4, _record), erlang:element(4, State), erlang:element(6, _record), erlang:element(7, _record)} end, none}; {error, Err} -> {stop, {abnormal, gleam@string:inspect(Err)}} end. -file("src/glats/jetstream/handler.gleam", 211). -spec handle_pull_message(glats:message(), pull_handler_state(IAR)) -> gleam@otp@actor:next(any(), pull_handler_state(IAR)). 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( begin _record = State, {pull_handler_state, erlang:element(2, _record), erlang:element(3, _record), erlang:element(4, _record), erlang:element(5, _record), erlang:element(6, _record), Inner@4} end ); false -> {continue, begin _record@1 = State, {pull_handler_state, erlang:element(2, _record@1), erlang:element(3, _record@1), erlang:element(4, _record@1), erlang:element(5, State) - 1, erlang:element(6, _record@1), Inner@4} end, none} end. -file("src/glats/jetstream/handler.gleam", 200). -spec pull_loop(glats:subscription_message(), pull_handler_state(IAO)) -> gleam@otp@actor:next(any(), pull_handler_state(IAO)). 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("src/glats/jetstream/handler.gleam", 139). ?DOC(" Start a pull consumer handler actor.\n"). -spec handle_pull_consumer( gleam@erlang@process:subject(glats:connection_message()), IAH, binary(), integer(), fun((glats:message(), IAH) -> outcome(IAH)), 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} ).