Current section

Files

Jump to
pulsar_elixir lib pulsar consumer.ex~
Raw

lib/pulsar/consumer.ex~

defmodule Pulsar.Consumer do
use GenServer
# 1. Topic Discovery
# 2. Partition Discovery
# 3. Connect to Broker
# 4. Subscribe
# 5. Flow
# 6. Consume, Forward, ACK and NACK
defstruct connection: nil,
module: nil,
topic: "",
subscription_type: "",
subscription_name: "",
consumer_id: 0
def start_link(connection, module, topic, subscription_type, subscription_name) do
GenServer.start_link(__MODULE__, [connection, module, topic, subscription_type, subscription_name])
end
@impl true
def init(args) do
[connection, module, topic, subscription_type, subscription_name] = args
consumer_id = System.unique_integer([:monotonic, :positive])
state = %__MODULE__{
connection: connection,
module: module,
topic: topic,
subscription_type: subscription_type,
subscription_name: subscription_name,
consumer_id: consumer_id
}
{:ok, state, {:continue, :subscribe}}
end
@impl true
def handle_continue(:subscribe, state) do
%__MODULE__{
connection: conn,
topic: topic,
subscription_type: subscription_type,
subscription_name: subscription_name,
consumer_id: consumer_id
} = state
Pulsar.Connection.subscribe(conn, consumer_id, topic, subscription_type, subscription_name)
Process.sleep(5_000)
Pulsar.Connection.flow(conn, consumer_id, 5)
{:noreply, state}
end
@impl true
def handle_call() do
end
@impl true
def handle_cast() do
end
end