Current section
Files
Jump to
Current section
Files
src/franz@group_subscriber.erl
-module(franz@group_subscriber).
-compile([no_auto_import, nowarn_unused_vars, nowarn_unused_function, nowarn_nomatch, inline]).
-define(FILEPATH, "src/franz/group_subscriber.gleam").
-export([commit/1, ack/1, new/6, with_group_config/2, with_consumer_config/2, start/1, stop/1]).
-export_type([callback_return/0, group_builder/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.
-type callback_return() :: any().
-opaque group_builder(FHQ) :: {group_builder,
franz:client(),
binary(),
list(binary()),
franz@message_type:message_type(),
fun((franz:kafka_message(), FHQ) -> callback_return()),
FHQ,
list(franz@group_config:group_config()),
list(franz@consumer_config:consumer_config())}.
-file("src/franz/group_subscriber.gleam", 24).
?DOC(" Commit the offset of the last message that was successfully processed.\n").
-spec commit(any()) -> callback_return().
commit(Cb_state) ->
franz_ffi:commit(Cb_state).
-file("src/franz/group_subscriber.gleam", 28).
?DOC(" Acknowledge the processing of the message.\n").
-spec ack(any()) -> callback_return().
ack(Cb_state) ->
franz_ffi:ack(Cb_state).
-file("src/franz/group_subscriber.gleam", 43).
?DOC(" Create a new group subscriber builder.\n").
-spec new(
franz:client(),
binary(),
list(binary()),
franz@message_type:message_type(),
fun((franz:kafka_message(), FIA) -> callback_return()),
FIA
) -> group_builder(FIA).
new(Client, Group_id, Topics, Message_type, Callback, Init_callback_state) ->
{group_builder,
Client,
Group_id,
Topics,
Message_type,
Callback,
Init_callback_state,
[],
[]}.
-file("src/franz/group_subscriber.gleam", 65).
?DOC(" Add a group configuration to the group builder.\n").
-spec with_group_config(group_builder(FIC), franz@group_config:group_config()) -> group_builder(FIC).
with_group_config(Group_builder, Group_config) ->
{group_builder,
erlang:element(2, Group_builder),
erlang:element(3, Group_builder),
erlang:element(4, Group_builder),
erlang:element(5, Group_builder),
erlang:element(6, Group_builder),
erlang:element(7, Group_builder),
[Group_config | erlang:element(8, Group_builder)],
erlang:element(9, Group_builder)}.
-file("src/franz/group_subscriber.gleam", 76).
?DOC(" Add a consumer configuration to the group builder.\n").
-spec with_consumer_config(
group_builder(FIF),
franz@consumer_config:consumer_config()
) -> group_builder(FIF).
with_consumer_config(Group_builder, Consumer_config) ->
{group_builder,
erlang:element(2, Group_builder),
erlang:element(3, Group_builder),
erlang:element(4, Group_builder),
erlang:element(5, Group_builder),
erlang:element(6, Group_builder),
erlang:element(7, Group_builder),
erlang:element(8, Group_builder),
[Consumer_config | erlang:element(9, Group_builder)]}.
-file("src/franz/group_subscriber.gleam", 87).
?DOC(" Start a new group subscriber.\n").
-spec start(group_builder(any())) -> {ok, gleam@erlang@process:pid_()} |
{error, franz:franz_error()}.
start(Group_builder) ->
franz_ffi:start_group_subscriber(
erlang:element(2, Group_builder),
erlang:element(3, Group_builder),
erlang:element(4, Group_builder),
erlang:element(9, Group_builder),
erlang:element(8, Group_builder),
erlang:element(5, Group_builder),
erlang:element(6, Group_builder),
erlang:element(7, Group_builder)
).
-file("src/franz/group_subscriber.gleam", 103).
-spec stop(gleam@erlang@process:pid_()) -> {ok, nil} |
{error, franz:franz_error()}.
stop(Pid) ->
franz_ffi:start_group_subscriber(Pid).