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