Packages
tai
0.0.64
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/system_bus.ex
defmodule Tai.SystemBus do
@type partitions :: pos_integer
@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) do
Registry.start_link(
keys: :duplicate,
name: __MODULE__,
partitions: partitions
)
end
def subscribe([]), do: :ok
def subscribe([topic | tail]) do
Registry.register(__MODULE__, topic, [])
subscribe(tail)
end
def subscribe(topic) do
topic
|> List.wrap()
|> subscribe
end
def unsubscribe([]), do: :ok
def unsubscribe([topic | tail]) do
Registry.unregister(__MODULE__, topic)
unsubscribe(tail)
end
def unsubscribe(topic) do
topic
|> List.wrap()
|> unsubscribe
end
def broadcast(topic, message) do
Registry.dispatch(__MODULE__, topic, fn entries ->
for {pid, _} <- entries, do: send(pid, message)
end)
end
end