diff --git a/examples/record_consumer.ex b/examples/record_consumer.ex index f8abdb5..8fdb24b 100644 --- a/examples/record_consumer.ex +++ b/examples/record_consumer.ex @@ -21,11 +21,10 @@ defmodule ExampleSupervisor do Supervisor.start_link(__MODULE__, arg, name: __MODULE__) end - @immpl true - def init(_arg) do + @impl true + def init(_) do children = [ - Drinkup, - ExampleRecordConsumer + {Drinkup, %{module: ExampleRecordConsumer}} ] Supervisor.init(children, strategy: :one_for_one) diff --git a/lib/consumer.ex b/lib/consumer.ex index f09db45..1e8d254 100644 --- a/lib/consumer.ex +++ b/lib/consumer.ex @@ -3,56 +3,7 @@ defmodule Drinkup.Consumer do An unopinionated consumer of the Firehose. Will receive all events, not just commits. """ - alias Drinkup.{ConsumerGroup, Event} + alias Drinkup.Event @callback handle_event(Event.t()) :: any() - - defmacro __using__(_opts) do - quote location: :keep do - use GenServer - require Logger - - @behaviour Drinkup.Consumer - - def child_spec(opts) do - %{ - id: __MODULE__, - start: {__MODULE__, :start_link, [opts]}, - type: :worker, - restart: :permanent, - max_restarts: 0, - shutdown: 500 - } - end - - def start_link(opts) do - GenServer.start_link(__MODULE__, [], opts) - end - - @impl GenServer - def init(_) do - ConsumerGroup.join() - {:ok, nil} - end - - @impl GenServer - def handle_info({:event, event}, state) do - {:ok, _pid} = - Task.start(fn -> - try do - __MODULE__.handle_event(event) - rescue - e -> - Logger.error( - "Error in event handler: #{Exception.format(:error, e, __STACKTRACE__)}" - ) - end - end) - - {:noreply, state} - end - - defoverridable GenServer - end - end end diff --git a/lib/consumer_group.ex b/lib/consumer_group.ex deleted file mode 100644 index 8f57e15..0000000 --- a/lib/consumer_group.ex +++ /dev/null @@ -1,39 +0,0 @@ -defmodule Drinkup.ConsumerGroup do - @moduledoc """ - Register consumers and dispatch events to them. - """ - - alias Drinkup.Event - - @scope __MODULE__ - @group :consumers - - def start_link(_) do - :pg.start_link(@scope) - end - - def child_spec(opts) do - %{ - id: __MODULE__, - start: {__MODULE__, :start_link, [opts]}, - type: :worker, - restart: :permanent, - shutdown: 500 - } - end - - @spec join() :: :ok - def join(), do: join(self()) - - @spec join(pid()) :: :ok - def join(pid), do: :pg.join(@scope, @group, pid) - - @spec dispatch(Event.t()) :: :ok - def dispatch(event) do - @scope - |> :pg.get_members(@group) - |> Enum.each(&send(&1, {:event, event})) - end - - # TODO: read `:pg` docs on what `monitor` is used fo -end diff --git a/lib/drinkup.ex b/lib/drinkup.ex index 65a958a..48ac516 100644 --- a/lib/drinkup.ex +++ b/lib/drinkup.ex @@ -1,14 +1,21 @@ defmodule Drinkup do use Supervisor - def start_link(arg \\ []) do - Supervisor.start_link(__MODULE__, arg, name: __MODULE__) + @type options() :: %{ + required(:consumer) => module(), + optional(:host) => String.t(), + optional(:cursor) => pos_integer() + } + + @spec start_link(options()) :: Supervisor.on_start() + def start_link(options) do + Supervisor.start_link(__MODULE__, options, name: __MODULE__) end - def init(_) do + def init(options) do children = [ - Drinkup.ConsumerGroup, - Drinkup.Socket + {Task.Supervisor, name: Drinkup.TaskSupervisor}, + {Drinkup.Socket, options} ] Supervisor.init(children, strategy: :one_for_one) diff --git a/lib/event.ex b/lib/event.ex index 5ccf9e8..bade594 100644 --- a/lib/event.ex +++ b/lib/event.ex @@ -1,4 +1,5 @@ defmodule Drinkup.Event do + require Logger alias Drinkup.Event @type t() :: @@ -21,4 +22,19 @@ defmodule Drinkup.Event do def valid_seq?(last_seq, nil) when is_integer(last_seq), do: true def valid_seq?(last_seq, seq) when is_integer(last_seq) and is_integer(seq), do: seq > last_seq def valid_seq?(_last_seq, _seq), do: false + + @spec dispatch(module(), t()) :: :ok + def dispatch(consumer, message) do + {:ok, _pid} = + Task.Supervisor.start_child(Drinkup.TaskSupervisor, fn -> + try do + consumer.handle_event(message) + rescue + e -> + Logger.error("Error in event handler: #{Exception.format(:error, e, __STACKTRACE__)}") + end + end) + + :ok + end end diff --git a/lib/record_consumer.ex b/lib/record_consumer.ex index 535db3a..dff42b4 100644 --- a/lib/record_consumer.ex +++ b/lib/record_consumer.ex @@ -11,7 +11,7 @@ defmodule Drinkup.RecordConsumer do {collections, _opts} = Keyword.pop(opts, :collections, []) quote location: :keep do - use Drinkup.Consumer + @behaviour Drinkup.Consumer @behaviour Drinkup.RecordConsumer def handle_event(%Drinkup.Event.Commit{} = event) do diff --git a/lib/socket.ex b/lib/socket.ex index e24ff59..7650ddf 100644 --- a/lib/socket.ex +++ b/lib/socket.ex @@ -4,7 +4,7 @@ defmodule Drinkup.Socket do """ require Logger - alias Drinkup.{ConsumerGroup, Event} + alias Drinkup.Event @behaviour :gen_statem @default_host "https://bsky.network" @@ -15,7 +15,7 @@ defmodule Drinkup.Socket do @op_regular 1 @op_error -1 - defstruct [:host, :seq, :conn, :stream] + defstruct [:options, :seq, :conn, :stream] @impl true def callback_mode, do: [:state_functions, :state_enter] @@ -30,17 +30,15 @@ defmodule Drinkup.Socket do } end - def start_link(opts \\ [], statem_opts) do - opts = Keyword.validate!(opts, host: @default_host) - host = Keyword.get(opts, :host) - cursor = Keyword.get(opts, :cursor) + def start_link(%{consumer: _} = options, statem_opts) do + options = Map.merge(%{host: @default_host, cursor: nil}, options) - :gen_statem.start_link(__MODULE__, {host, cursor}, statem_opts) + :gen_statem.start_link(__MODULE__, options, statem_opts) end @impl true - def init({host, cursor}) do - data = %__MODULE__{host: host, seq: cursor} + def init(%{cursor: seq} = options) do + data = %__MODULE__{seq: seq, options: options} {:ok, :disconnected, data, [{:next_event, :internal, :connect}]} end @@ -54,10 +52,10 @@ defmodule Drinkup.Socket do {:next_state, :connecting_http, data} end - def connecting_http(:enter, _from, data) do + def connecting_http(:enter, _from, %{options: options} = data) do Logger.debug("Connecting to http") - %{host: host, port: port} = URI.new!(data.host) + %{host: host, port: port} = URI.new!(options.host) {:ok, conn} = :gun.open(:binary.bin_to_list(host), port, %{ @@ -107,7 +105,7 @@ defmodule Drinkup.Socket do :keep_state_and_data end - def connected(:info, {:gun_ws, conn, stream, {:binary, frame}}, data) do + def connected(:info, {:gun_ws, conn, stream, {:binary, frame}}, %{options: options} = data) do # TODO: let clients specify a handler for raw* (*decoded) packets to support any atproto subscription # Will also need support for JSON frames with {:ok, header, next} <- CAR.DagCbor.decode(frame), @@ -123,7 +121,7 @@ defmodule Drinkup.Socket do Logger.warning("Received unrecognised event from firehose: #{inspect({type, payload})}") message -> - ConsumerGroup.dispatch(message) + Event.dispatch(options.consumer, message) end {:keep_state, data}