diff --git a/README.md b/README.md index 82b7371..ad4b8ee 100644 --- a/README.md +++ b/README.md @@ -7,6 +7,5 @@ firehose. - Support for different subscriptions other than `com.atproto.sync.subscribeRepo' -- Support for multiple instances at once, each with unique consumers (for - listening to multiple subscriptions at once) - Tests +- Documentation diff --git a/examples/basic_consumer.ex b/examples/basic_consumer.ex new file mode 100644 index 0000000..ef040cd --- /dev/null +++ b/examples/basic_consumer.ex @@ -0,0 +1,26 @@ +defmodule BasicConsumer do + @behaviour Drinkup.Consumer + + def handle_event(%Drinkup.Event.Commit{} = event) do + IO.inspect(event, label: "Got commit event") + end + + def handle_event(_), do: :noop +end + +defmodule ExampleSupervisor do + use Supervisor + + def start_link(arg \\ []) do + Supervisor.start_link(__MODULE__, arg, name: __MODULE__) + end + + @impl true + def init(_) do + children = [ + {Drinkup, %{consumer: BasicConsumer}} + ] + + Supervisor.init(children, strategy: :one_for_one) + end +end diff --git a/examples/multiple_consumers.ex b/examples/multiple_consumers.ex new file mode 100644 index 0000000..d825d96 --- /dev/null +++ b/examples/multiple_consumers.ex @@ -0,0 +1,35 @@ +defmodule PostDeleteConsumer do + use Drinkup.RecordConsumer, collections: ["app.bsky.feed.post"] + + def handle_delete(record) do + IO.inspect(record, label: "update") + end +end + +defmodule IdentityConsumer do + @behaviour Drinkup.Consumer + + def handle_event(%Drinkup.Event.Identity{} = event) do + IO.inspect(event, label: "identity event") + end + + def handle_event(_), do: :noop +end + +defmodule ExampleSupervisor do + use Supervisor + + def start_link(arg \\ []) do + Supervisor.start_link(__MODULE__, arg, name: __MODULE__) + end + + @impl true + def init(_) do + children = [ + {Drinkup, %{consumer: PostDeleteConsumer}}, + {Drinkup, %{consumer: IdentityConsumer, name: :identities}} + ] + + Supervisor.init(children, strategy: :one_for_one) + end +end diff --git a/examples/record_consumer.ex b/examples/record_consumer.ex index 8fdb24b..b5ff0b5 100644 --- a/examples/record_consumer.ex +++ b/examples/record_consumer.ex @@ -17,14 +17,14 @@ end defmodule ExampleSupervisor do use Supervisor - def start_link(args \\ []) do + def start_link(arg \\ []) do Supervisor.start_link(__MODULE__, arg, name: __MODULE__) end @impl true def init(_) do children = [ - {Drinkup, %{module: ExampleRecordConsumer}} + {Drinkup, %{consumer: ExampleRecordConsumer}} ] Supervisor.init(children, strategy: :one_for_one) diff --git a/lib/application.ex b/lib/application.ex new file mode 100644 index 0000000..96b36d5 --- /dev/null +++ b/lib/application.ex @@ -0,0 +1,8 @@ +defmodule Drinkup.Application do + use Application + + def start(_type, _args) do + children = [{Registry, keys: :unique, name: Drinkup.Registry}] + Supervisor.start_link(children, strategy: :one_for_one) + end +end diff --git a/lib/drinkup.ex b/lib/drinkup.ex index 48ac516..866ea25 100644 --- a/lib/drinkup.ex +++ b/lib/drinkup.ex @@ -1,23 +1,32 @@ defmodule Drinkup do use Supervisor + alias Drinkup.Options - @type options() :: %{ - required(:consumer) => module(), - optional(:host) => String.t(), - optional(:cursor) => pos_integer() - } + @dialyzer nowarn_function: {:init, 1} + @impl true + def init({%Options{name: name} = drinkup_options, supervisor_options}) do + children = [ + {Task.Supervisor, name: {:via, Registry, {Drinkup.Registry, {name, Tasks}}}}, + {Drinkup.Socket, drinkup_options} + ] - @spec start_link(options()) :: Supervisor.on_start() - def start_link(options) do - Supervisor.start_link(__MODULE__, options, name: __MODULE__) + Supervisor.start_link( + children, + supervisor_options ++ [name: {:via, Registry, {Drinkup.Registry, {name, Supervisor}}}] + ) end - def init(options) do - children = [ - {Task.Supervisor, name: Drinkup.TaskSupervisor}, - {Drinkup.Socket, options} - ] + @spec child_spec(Options.options()) :: Supervisor.child_spec() + def child_spec(%{} = options), do: child_spec({options, [strategy: :one_for_one]}) - Supervisor.init(children, strategy: :one_for_one) + @spec child_spec({Options.options(), Keyword.t()}) :: Supervisor.child_spec() + def child_spec({drinkup_options, supervisor_options}) do + %{ + id: Map.get(drinkup_options, :name, __MODULE__), + start: {__MODULE__, :init, [{Options.from(drinkup_options), supervisor_options}]}, + type: :supervisor, + restart: :permanent, + shutdown: 500 + } end end diff --git a/lib/event.ex b/lib/event.ex index bade594..d0ed051 100644 --- a/lib/event.ex +++ b/lib/event.ex @@ -1,6 +1,6 @@ defmodule Drinkup.Event do require Logger - alias Drinkup.Event + alias Drinkup.{Event, Options} @type t() :: Event.Commit.t() @@ -23,10 +23,12 @@ defmodule Drinkup.Event do 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 + @spec dispatch(t(), Options.t()) :: :ok + def dispatch(message, %Options{consumer: consumer, name: name}) do + supervisor_name = {:via, Registry, {Drinkup.Registry, {name, Tasks}}} + {:ok, _pid} = - Task.Supervisor.start_child(Drinkup.TaskSupervisor, fn -> + Task.Supervisor.start_child(supervisor_name, fn -> try do consumer.handle_event(message) rescue diff --git a/lib/options.ex b/lib/options.ex new file mode 100644 index 0000000..5a4441a --- /dev/null +++ b/lib/options.ex @@ -0,0 +1,22 @@ +defmodule Drinkup.Options do + use TypedStruct + + @default_host "https://bsky.network" + + @type options() :: %{ + required(:consumer) => module(), + optional(:name) => atom(), + optional(:host) => String.t(), + optional(:cursor) => pos_integer() + } + + typedstruct do + field :consumer, module(), enforce: true + field :name, atom(), default: Drinkup + field :host, String.t(), default: @default_host + field :cursor, pos_integer() | nil + end + + @spec from(options()) :: t() + def from(%{consumer: _} = options), do: struct(__MODULE__, options) +end diff --git a/lib/socket.ex b/lib/socket.ex index 7650ddf..74600a3 100644 --- a/lib/socket.ex +++ b/lib/socket.ex @@ -4,10 +4,9 @@ defmodule Drinkup.Socket do """ require Logger - alias Drinkup.Event + alias Drinkup.{Event, Options} @behaviour :gen_statem - @default_host "https://bsky.network" @timeout :timer.seconds(5) # TODO: `flow` determines messages in buffer. Determine ideal value? @flow 10 @@ -30,9 +29,7 @@ defmodule Drinkup.Socket do } end - def start_link(%{consumer: _} = options, statem_opts) do - options = Map.merge(%{host: @default_host, cursor: nil}, options) - + def start_link(%Options{} = options, statem_opts) do :gen_statem.start_link(__MODULE__, options, statem_opts) end @@ -121,7 +118,7 @@ defmodule Drinkup.Socket do Logger.warning("Received unrecognised event from firehose: #{inspect({type, payload})}") message -> - Event.dispatch(options.consumer, message) + Event.dispatch(message, options) end {:keep_state, data} diff --git a/mix.exs b/mix.exs index 07facb1..7b21d68 100644 --- a/mix.exs +++ b/mix.exs @@ -14,7 +14,8 @@ defmodule Drinkup.MixProject do # Run "mix help compile.app" to learn about applications. def application do [ - extra_applications: [:logger] + extra_applications: [:logger], + mod: {Drinkup.Application, []} ] end