Packages
tai
0.0.47
0.0.75
0.0.74
0.0.73
0.0.72
0.0.71
0.0.70
0.0.69
0.0.68
0.0.67
0.0.66
0.0.65
0.0.64
0.0.63
0.0.62
0.0.61
0.0.60
0.0.59
0.0.58
0.0.57
0.0.56
0.0.55
0.0.54
0.0.53
0.0.52
0.0.51
0.0.50
0.0.49
0.0.48
0.0.47
0.0.46
0.0.45
0.0.44
0.0.43
0.0.42
0.0.41
0.0.40
0.0.39
0.0.38
0.0.37
0.0.36
0.0.35
0.0.34
0.0.33
0.0.32
0.0.31
0.0.30
0.0.29
0.0.28
0.0.27
0.0.26
0.0.25
0.0.24
0.0.23
0.0.22
0.0.21
0.0.20
0.0.19
0.0.18
0.0.17
0.0.16
0.0.15
0.0.14
0.0.13
0.0.12
0.0.11
0.0.10
0.0.9
0.0.8
0.0.7
0.0.6
0.0.5
0.0.4
0.0.3
0.0.2
0.0.1
A composable, real time, market data and trade execution toolkit
Current section
Files
Jump to
Current section
Files
lib/tai/events.ex
defmodule Tai.Events do
@type event :: Tai.Event.t()
@type event_type :: module
@type partitions :: pos_integer
@type level :: :debug | :info | :warn | :error
@type subscribe_error_reasons :: {:already_registered, pid} | :event_not_registered
@spec child_spec(opts :: term) :: Supervisor.child_spec()
def child_spec(opts) do
%{
id: __MODULE__,
start: {__MODULE__, :start_link, [opts]},
type: :supervisor,
restart: :permanent,
shutdown: 500
}
end
@spec start_link(partitions) :: {:ok, pid} | {:error, term}
def start_link(partitions) when partitions > 0 do
Registry.start_link(keys: :duplicate, name: __MODULE__, partitions: partitions)
end
@spec firehose_subscribe :: {:ok, pid} | {:error, subscribe_error_reasons}
def firehose_subscribe do
Registry.register(__MODULE__, :firehose, [])
end
@spec subscribe(event_type) :: {:ok, pid} | {:error, :subscribe_error_reasons}
def subscribe(event_type) when is_atom(event_type) do
Registry.register(__MODULE__, event_type, [])
end
@spec error(event) :: :ok
def error(event), do: event |> broadcast(:error)
@spec warn(event) :: :ok
def warn(event), do: event |> broadcast(:warn)
@spec info(event) :: :ok
def info(event), do: event |> broadcast(:info)
@spec debug(event) :: :ok
def debug(event), do: event |> broadcast(:debug)
@spec broadcast(event, level) :: :ok
def broadcast(event, level) do
event_type = Map.fetch!(event, :__struct__)
msg = {Tai.Event, event, level}
Registry.dispatch(__MODULE__, event_type, fn entries ->
for {pid, _} <- entries, do: send(pid, msg)
end)
Registry.dispatch(__MODULE__, :firehose, fn entries ->
for {pid, _} <- entries, do: send(pid, msg)
end)
:ok
end
end