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