Current section

Files

Jump to
franz src franz@group_subscriber.erl
Raw

src/franz@group_subscriber.erl

-module(franz@group_subscriber).
-compile([no_auto_import, nowarn_unused_vars, nowarn_unused_function, nowarn_nomatch]).
-export([commit/1, ack/1, new/6, with_group_config/2, with_consumer_config/2, start/1]).
-export_type([callback_return/0, group_builder/1]).
-type callback_return() :: any().
-type group_builder(FPI) :: {group_builder,
franz:franz_client(),
binary(),
list(binary()),
franz@message_type:message_type(),
fun((franz:kafka_message(), FPI) -> callback_return()),
FPI,
list(franz@group_config:group_config()),
list(franz@consumer_config:consumer_config())}.
-file("src/franz/group_subscriber.gleam", 23).
-spec commit(any()) -> callback_return().
commit(Cb_state) ->
franz_ffi:commit(Cb_state).
-file("src/franz/group_subscriber.gleam", 26).
-spec ack(any()) -> callback_return().
ack(Cb_state) ->
franz_ffi:ack(Cb_state).
-file("src/franz/group_subscriber.gleam", 40).
-spec new(
franz:franz_client(),
binary(),
list(binary()),
franz@message_type:message_type(),
fun((franz:kafka_message(), FPS) -> callback_return()),
FPS
) -> group_builder(FPS).
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", 60).
-spec with_group_config(group_builder(FPU), franz@group_config:group_config()) -> group_builder(FPU).
with_group_config(Group_builder, Group_config) ->
_record = Group_builder,
{group_builder,
erlang:element(2, _record),
erlang:element(3, _record),
erlang:element(4, _record),
erlang:element(5, _record),
erlang:element(6, _record),
erlang:element(7, _record),
[Group_config | erlang:element(8, Group_builder)],
erlang:element(9, _record)}.
-file("src/franz/group_subscriber.gleam", 70).
-spec with_consumer_config(
group_builder(FPX),
franz@consumer_config:consumer_config()
) -> group_builder(FPX).
with_consumer_config(Group_builder, Consumer_config) ->
_record = Group_builder,
{group_builder,
erlang:element(2, _record),
erlang:element(3, _record),
erlang:element(4, _record),
erlang:element(5, _record),
erlang:element(6, _record),
erlang:element(7, _record),
erlang:element(8, _record),
[Consumer_config | erlang:element(9, Group_builder)]}.
-file("src/franz/group_subscriber.gleam", 80).
-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)
).