Current section

Files

Jump to
glats src glats@jetstream@handler.erl
Raw

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]).
-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}
).