Current section

Files

Jump to
antenna lib antenna.ex
Raw

lib/antenna.ex

defmodule Antenna do
@moduledoc """
`Antenna` is a mixture of [`Phoenix.PubSub`](https://hexdocs.pm/phoenix_pubsub/Phoenix.PubSub.html)
and [`:gen_event`](https://www.erlang.org/doc/apps/stdlib/gen_event.html) functionality
with some batteries included.
It implements back-pressure on top of `GenStage`, is fully conformant with
[OTP Design Principles](https://www.erlang.org/doc/system/events). and
is distributed out of the box.
`Antenna` supports both asynchronous _and_ synchronous events. While the most preferrable way
would be to stay fully async with `Antenna.event/3`, one still might propagate the event
synchronously with `Antenna.sync_event/3` and collect all the responses from all the handlers.
One can have as many isolated `Antenna`s as necessary, distinguished by `Antenna.t:id/0`.
The workflow looks like shown below.
## Sequence Diagram
```mermaid
sequenceDiagram
Consumer->>+Broadcaster: sync_event(channel, event)
Consumer->>+Broadcaster: event(channel, event)
Broadcaster-->>+Consumer@Node1: event
Broadcaster-->>+Consumer@Node2: event
Broadcaster-->>+Consumer@NodeN: event
Consumer@Node1-->>-NoOp: mine?
Consumer@NodeN-->>-NoOp: mine?
Consumer@Node2-->>+Matchers: event
Matchers-->>+Handlers: handle_match/2
Matchers-->>+Handlers: handle_match/2
Handlers->>-Consumer: response(to: sync_event)
```
[ASCII representation](https://cascii.app/4164d).
## Usage Example
The consumer of this library is supposed to declare one or more matchers, subscribing to one
or more channels, and then call `Antenna.event/2` to propagate the event.
```elixir
assert {:ok, _pid, "{:tag_1, a, _} when is_nil(a)"} =
Antenna.match(Antenna, {:tag_1, a, _} when is_nil(a), self(), channels: [:chan_1])
assert :ok = Antenna.event(Antenna, [:chan_1], {:tag_1, nil, 42})
assert_receive {:antenna_event, :chan_1, {:tag_1, nil, 42}}
```
"""
require Logger
alias Antenna.PubSub.Broadcaster
@typedoc "The identifier of the isolated `Antenna`"
@type id :: module()
@typedoc "The identifier of the channel, messages can be sent to, preferrably `atom()`"
@type channel :: atom() | term()
@typedoc "The event being sent to the listeners"
@type event :: term()
@typedoc """
The actual handler to be associated with an event(s). It might be either a function
or a process id, in which case the message of a following shape will be sent to it.
```elixir
{:antenna_event, channel, event}
```
"""
@type handler :: (event() -> term()) | (channel(), event() -> term()) | Antenna.Matcher.t() | pid() | GenServer.name()
@id Application.compile_env(:antenna, :id, Antenna)
@doc false
def id, do: @id
@doc false
def delivery(id), do: Module.concat(id, Delivery)
@doc false
def matchers(id), do: Module.concat(id, Matchers)
@doc false
def channels(id), do: Module.concat(id, Channels)
@doc false
def guard(id), do: Module.concat(id, Guard)
@doc false
def id_opts(%{id: id} = opts), do: {id, Map.delete(opts, :id)}
def id_opts(%{name: id} = opts), do: {id, Map.delete(opts, :name)}
def id_opts(opts) do
if Keyword.keyword?(opts),
do: Keyword.pop_lazy(opts, :id, fn -> Keyword.get(opts, :name, __MODULE__) end),
else: {opts, []}
end
# Supervision tree
use Supervisor
def start_link(init_arg \\ []) do
init_arg = if Keyword.keyword?(init_arg), do: Keyword.get_lazy(init_arg, :id, &id/0), else: init_arg
Supervisor.start_link(__MODULE__, init_arg, name: __MODULE__)
end
@impl true
def init(id) do
children = [
%{id: :pg, start: {__MODULE__, :start_pg, [Antenna.channels(id)]}},
{Antenna.Guard, id: id},
{Antenna.PubSub, id: id},
{DistributedSupervisor, name: Antenna.matchers(id), monitor_nodes: true}
]
Supervisor.init(children, strategy: :one_for_one)
end
@doc false
@spec start_pg(module()) :: {:ok, pid()} | :ignore
def start_pg(scope) do
with {:error, {:already_started, _pid}} <- :pg.start_link(scope), do: :ignore
end
@doc """
Declares a matcher for tagged events.
### Example
```elixir
Antenna.match(Antenna, %{tag: _, success: false}, fn channel, message ->
Logger.warning("The processing failed for [" <>
inspect(channel) <> "], result: " <> inspect(message))
end, channels: [:rabbit])
```
"""
@doc section: :setup
defmacro match(id \\ @id, match, handlers, opts \\ [])
defmacro match(id, match, handlers, opts) do
name = Macro.to_string(match)
quote generated: true, location: :keep do
matcher = fn
unquote(match) -> true
_ -> false
end
opts = unquote(opts)
DistributedSupervisor.start_child(
Antenna.matchers(unquote(id)),
{Antenna.Matcher,
name: unquote(name),
id: unquote(id),
match: unquote(name),
matcher: matcher,
handlers: List.wrap(unquote(handlers)),
channels: List.wrap(Keyword.get(opts, :channels)),
once?: Keyword.get(opts, :once?, false)}
)
end
end
@doc """
Undeclares a matcher for tagged events previously declared with `Antenna.match/4`.
Accepts both an original match _or_ a name returned by `Antenna.match/4`,
which is effectively `Macro.to_string(match)`.
### Example
```elixir
Antenna.unmatch(Antenna, %{tag: _, success: false})
```
"""
@doc section: :setup
defmacro unmatch(id \\ @id, match) do
quote generated: true, location: :keep do
with pid when is_pid(pid) <- Antenna.whereis(unquote(id), unquote(match)),
do: DistributedSupervisor.terminate_child(Antenna.matchers(unquote(id)), pid)
end
end
@doc false
@doc section: :setup
defmacro attach(id \\ @id, channels, match) do
quote generated: true, location: :keep do
with pid when is_pid(pid) <- Antenna.whereis(unquote(id), unquote(match)),
do: Antenna.subscribe(unquote(id), unquote(channels), pid)
end
end
@doc false
@doc section: :setup
defmacro unattach(id \\ @id, channels, match) do
quote generated: true, location: :keep do
with pid when is_pid(pid) <- Antenna.whereis(unquote(id), unquote(match)),
do: Antenna.unsubscribe(unquote(id), unquote(channels), pid)
end
end
@doc """
Subscribes a matcher process specified by `pid` to a channel(s)
"""
@doc section: :setup
@spec subscribe(id :: id(), channels :: channel() | [channel()], pid()) :: :ok
def subscribe(id \\ @id, channels, pid)
def subscribe(_, [], _), do: :ok
def subscribe(id, channels, pid) when is_pid(pid) do
channels
|> List.wrap()
|> Enum.each(&join(id, &1, pid))
end
def subscribe(id, channels, pid),
do: Logger.warning("Unexpected subscription: " <> inspect(id: id, channels: channels, pid: pid))
@doc """
Unsubscribes a previously subscribed matcher process specified by `pid` from the channel(s)
"""
@doc section: :setup
@spec unsubscribe(id :: id(), channels :: channel() | [channel()], pid()) :: :ok
def unsubscribe(id \\ @id, channels, pid)
def unsubscribe(_, [], _), do: :ok
def unsubscribe(id, channels, pid) when is_pid(pid) do
channels
|> List.wrap()
|> Enum.each(&leave(id, &1, pid))
end
@doc """
Adds a handler to the matcher process specified by `pid`
"""
@doc section: :setup
@spec handle(id :: id(), handlers :: handler() | [handler()], pid()) :: :ok
def handle(id \\ @id, handlers, pid)
def handle(_, [], _), do: :ok
def handle(id, handlers, pid) when is_pid(pid) do
handlers
|> List.wrap()
|> Enum.each(&do_handle(id, &1, pid))
end
@doc """
Removes a handler from the matcher process specified by `pid`
"""
@doc section: :setup
@spec unhandle(id :: id(), handlers :: handler() | [handler()], pid()) :: :ok
def unhandle(id \\ @id, handlers, pid)
def unhandle(_, [], _), do: :ok
def unhandle(id, handlers, pid) when is_pid(pid) do
handlers
|> List.wrap()
|> Enum.each(&do_unhandle(id, &1, pid))
end
@doc false
@doc section: :internals
defmacro whereis(id \\ @id, match)
defmacro whereis(id, match) when is_binary(match), do: do_whereis(id, match)
defmacro whereis(id, match), do: do_whereis(id, Macro.to_string(match))
defp do_whereis(id, match) do
quote generated: true, location: :keep do
case DistributedSupervisor.children(Antenna.matchers(unquote(id))) do
%{unquote(match) => {pid, _spec}} when is_pid(pid) -> pid
_ -> nil
end
end
end
@doc """
Sends an event to all the associated matchers through channels.
The special `:*` might be specified as channels, then the event
will be sent to all the registered channels.
If one wants to collect results of all the registered event handlers,
they should look at `sync_event/3` instead.
"""
@doc section: :client
@spec event(id :: id(), channels :: channel() | [channel()], event :: event()) :: :ok
def event(id \\ @id, channels, event),
do: Broadcaster.async_notify(Antenna.delivery(id), {List.wrap(channels), event})
@doc """
Sends an event to all the associated matchers through channels and collects results.
The special `:*` might be specified as channels, then the event
will be sent to all the registered channels.
Unlike `event/3`, this call is _synchronous_, which means it would block until
all the registered handlers respond; the results would have been collected
and returned as a list of no specific order.
"""
@doc section: :client
@spec sync_event(id :: id(), channels :: channel() | [channel()], event :: event()) :: [term()]
def sync_event(id \\ @id, channels, event),
do: Broadcaster.sync_notify(Antenna.delivery(id), {List.wrap(channels), event})
@doc """
Returns a map of matches to matchers
"""
@doc section: :internals
@spec registered_matchers(id :: id()) :: %{term() => {pid(), Supervisor.child_spec()}}
def registered_matchers(id),
do: id |> Antenna.matchers() |> DistributedSupervisor.children()
@spec join(id(), channel(), pid()) :: :ok | :already_joined
defp join(id, channel, pid) do
scope = channels(id)
if pid in :pg.get_members(scope, channel) do
:already_joined
else
:ok = Antenna.Guard.add_channel(id, channel, pid)
:pg.join(scope, channel, pid)
end
end
@spec leave(id(), channel(), pid()) :: :ok | :not_joined
defp leave(id, channel, pid) do
with :ok <- id |> channels() |> :pg.leave(channel, pid),
do: Antenna.Guard.remove_channel(id, channel, pid)
end
@spec do_handle(id(), handler(), pid()) :: :ok
defp do_handle(id, handler, pid) do
Antenna.Guard.add_handler(id, handler, pid)
GenServer.cast(pid, {:add_handler, handler})
end
@spec do_unhandle(id(), handler(), pid()) :: :ok
defp do_unhandle(id, handler, pid) do
Antenna.Guard.remove_handler(id, handler, pid)
GenServer.cast(pid, {:remove_handler, handler})
end
end