Current section

Files

Jump to
off_broadway_pulsar lib off_broadway pulsar acknowledger.ex
Raw

lib/off_broadway/pulsar/acknowledger.ex

defmodule OffBroadway.Pulsar.Acknowledger do
@moduledoc false
@behaviour Broadway.Acknowledger
@impl Broadway.Acknowledger
def ack(%{consumer: consumer}, successful, failed) do
ack_ids =
Enum.flat_map(successful, fn %{acknowledger: {_, _, message_id}} ->
List.wrap(message_id)
end)
:ok = Pulsar.ack(consumer, ack_ids)
nack_ids =
Enum.flat_map(failed, fn %{acknowledger: {_, _, message_id}} ->
List.wrap(message_id)
end)
:ok = Pulsar.nack(consumer, nack_ids)
end
end