Current section

Files

Jump to
franz src franz@consumer@topic_subscriber.erl
Raw

src/franz@consumer@topic_subscriber.erl

-module(franz@consumer@topic_subscriber).
-compile([no_auto_import, nowarn_unused_vars, nowarn_unused_function, nowarn_nomatch, inline]).
-define(FILEPATH, "src/franz/consumer/topic_subscriber.gleam").
-export([ack/1, new/7, with_config/2, with_commited_offset/3, named_client/1, start/1, supervised/1]).
-export_type([message/0, topic_subscriber/0, partitions/0, builder/1, ack/0]).
-if(?OTP_RELEASE >= 27).
-define(MODULEDOC(Str), -moduledoc(Str)).
-define(DOC(Str), -doc(Str)).
-else.
-define(MODULEDOC(Str), -compile([])).
-define(DOC(Str), -compile([])).
-endif.
-type message() :: any().
-type topic_subscriber() :: {topic_subscriber,
gleam@erlang@process:name(message())}.
-type partitions() :: {partitions, list(integer())} | all.
-opaque builder(GWW) :: {builder,
gleam@erlang@process:name(message()),
franz:client(),
binary(),
partitions(),
list({integer(), integer()}),
franz@consumer@message_type:message_type(),
fun((integer(), franz:kafka_message(), GWW) -> ack()),
GWW,
list(franz@consumer@config:config())}.
-type ack() :: any().
-file("src/franz/consumer/topic_subscriber.gleam", 44).
?DOC(
" Acknowledges the processing of a message.\n"
" Use this in your callback to confirm message receipt.\n"
).
-spec ack(any()) -> ack().
ack(Cb_state) ->
franz_ffi:ack(Cb_state).
-file("src/franz/consumer/topic_subscriber.gleam", 60).
?DOC(
" Creates a new topic subscriber builder.\n"
" The callback will be called for each message received from the topic partitions.\n"
).
-spec new(
gleam@erlang@process:name(message()),
franz:client(),
binary(),
partitions(),
franz@consumer@message_type:message_type(),
fun((integer(), franz:kafka_message(), GXE) -> ack()),
GXE
) -> builder(GXE).
new(
Name,
Client,
Topic,
Partitions,
Message_type,
Callback,
Init_callback_state
) ->
{builder,
Name,
Client,
Topic,
Partitions,
[],
Message_type,
Callback,
Init_callback_state,
[]}.
-file("src/franz/consumer/topic_subscriber.gleam", 84).
?DOC(
" Adds a consumer configuration option to the topic subscriber builder.\n"
" Multiple configurations can be chained together.\n"
).
-spec with_config(builder(GXG), franz@consumer@config:config()) -> builder(GXG).
with_config(Builder, Consumer_config) ->
{builder,
erlang:element(2, Builder),
erlang:element(3, Builder),
erlang:element(4, Builder),
erlang:element(5, Builder),
erlang:element(6, Builder),
erlang:element(7, Builder),
erlang:element(8, Builder),
erlang:element(9, Builder),
[Consumer_config | erlang:element(10, Builder)]}.
-file("src/franz/consumer/topic_subscriber.gleam", 97).
?DOC(
" Adds a committed offset to the topic subscriber builder.\n"
" CommittedOffsets are the offsets for the messages that have been successfully processed (acknowledged),\n"
" not the begin-offset to start fetching from.\n"
).
-spec with_commited_offset(builder(GXJ), integer(), integer()) -> builder(GXJ).
with_commited_offset(Builder, Partition, Offset) ->
{builder,
erlang:element(2, Builder),
erlang:element(3, Builder),
erlang:element(4, Builder),
erlang:element(5, Builder),
[{Partition, Offset} | erlang:element(6, Builder)],
erlang:element(7, Builder),
erlang:element(8, Builder),
erlang:element(9, Builder),
erlang:element(10, Builder)}.
-file("src/franz/consumer/topic_subscriber.gleam", 135).
-spec named_client(gleam@erlang@process:name(message())) -> topic_subscriber().
named_client(Name) ->
{topic_subscriber, Name}.
-file("src/franz/consumer/topic_subscriber.gleam", 109).
?DOC(" Starts a new topic subscriber with the configured settings.\n").
-spec start(builder(any())) -> {ok, gleam@otp@actor:started(topic_subscriber())} |
{error, gleam@otp@actor:start_error()}.
start(Builder) ->
case franz_ffi:start_topic_subscriber(
erlang:element(3, Builder),
erlang:element(4, Builder),
erlang:element(5, Builder),
erlang:element(10, Builder),
erlang:element(6, Builder),
erlang:element(7, Builder),
erlang:element(8, Builder),
erlang:element(9, Builder)
) of
{ok, Pid} ->
{ok, {started, Pid, named_client(erlang:element(2, Builder))}};
{error, Error} ->
{error, {init_exited, {abnormal, Error}}}
end.
-file("src/franz/consumer/topic_subscriber.gleam", 131).
?DOC(
" Creates a supervised worker for the topic subscriber.\n"
" This can be used with Gleam's OTP supervision trees to ensure the subscriber is restarted on failure.\n"
).
-spec supervised(builder(any())) -> gleam@otp@supervision:child_specification(topic_subscriber()).
supervised(Builder) ->
gleam@otp@supervision:worker(fun() -> start(Builder) end).