Packages
phoenix
1.0.1
1.8.9
1.8.8
1.8.7
1.8.6
1.8.5
1.8.4
1.8.3
1.8.2
1.8.1
1.8.0
1.8.0-rc.4
1.8.0-rc.3
1.8.0-rc.2
1.8.0-rc.1
1.8.0-rc.0
1.7.24
1.7.23
1.7.22
1.7.21
1.7.20
1.7.19
1.7.18
1.7.17
1.7.16
1.7.15
1.7.14
1.7.13
1.7.12
1.7.11
1.7.10
1.7.9
1.7.8
1.7.7
1.7.6
1.7.5
1.7.4
1.7.3
1.7.2
1.7.1
1.7.0
1.7.0-rc.3
1.7.0-rc.2
1.7.0-rc.1
1.7.0-rc.0
1.6.17
1.6.16
1.6.15
1.6.14
1.6.13
1.6.12
1.6.11
1.6.10
1.6.9
1.6.8
1.6.7
1.6.6
1.6.5
1.6.4
1.6.3
1.6.2
1.6.1
1.6.0
1.6.0-rc.1
1.6.0-rc.0
1.5.15
1.5.14
1.5.13
1.5.12
1.5.11
1.5.10
1.5.9
1.5.8
1.5.7
1.5.6
1.5.5
1.5.4
1.5.3
1.5.2
1.5.1
1.5.0
1.5.0-rc.0
1.4.18
1.4.17
1.4.16
1.4.15
1.4.14
1.4.13
1.4.12
1.4.11
1.4.10
1.4.9
1.4.8
1.4.7
1.4.6
1.4.5
1.4.4
1.4.3
1.4.2
1.4.1
1.4.0
1.4.0-rc.3
1.4.0-rc.2
1.4.0-rc.1
1.4.0-rc.0
1.3.5
1.3.4
1.3.3
1.3.2
1.3.1
1.3.0
1.3.0-rc.3
1.3.0-rc.2
1.3.0-rc.1
1.3.0-rc.0
1.2.5
1.2.4
1.2.3
1.2.2
1.2.1
1.2.0
1.2.0-rc.1
1.2.0-rc.0
1.1.9
1.1.8
1.1.7
1.1.6
1.1.5
1.1.4
1.1.3
1.1.2
1.1.1
1.1.0
1.0.6
1.0.5
1.0.4
1.0.3
1.0.2
1.0.1
1.0.0
0.17.1
0.17.0
0.16.1
0.16.0
0.15.0
0.14.0
0.13.1
0.13.0
0.12.0
0.11.0
0.10.0
0.9.0
0.8.0
0.7.2
0.7.1
0.7.0
0.6.2
0.6.1
0.6.0
0.5.0
0.4.1
0.4.0
0.3.1
0.3.0
0.2.11
0.2.10
0.2.9
0.2.8
0.2.7
0.2.6
0.2.5
0.2.4
0.2.3
0.2.2
0.2.1
0.2.0
0.1.0
Productive. Reliable. Fast. A productive web framework that does not compromise speed or maintainability.
Security advisory:
This version has known vulnerabilities.
View advisories
Current section
Files
Jump to
Current section
Files
lib/phoenix/pubsub.ex
defmodule Phoenix.PubSub do
@moduledoc """
Front-end to Phoenix pubsub layer.
Used internally by Channels for pubsub broadcast but
also provides an API for direct usage.
## Adapters
Phoenix pubsub was designed to be flexible and support
multiple backends. We currently ship with two backends:
* `Phoenix.PubSub.PG2` - uses Distributed Elixir,
directly exchanging notifications between servers
* `Phoenix.PubSub.Redis` - uses Redis to exchange
data between servers
Pubsub adapters are often configured in your endpoint:
config :my_app, MyApp.Endpoint,
pubsub: [adapter: Phoenix.PubSub.PG2]
The configuration above takes care of starting the
pubsub backend and exposing its functions via the
endpoint module.
## Direct usage
It is also possible to use `Phoenix.PubSub` directly
or even run your own pubsub backends outside of an
Endpoint.
The first step is to start the adapter of choice in your
supervision tree:
supervisor(Phoenix.PubSub.Redis, [:my_redis_pubsub, host: "192.168.100.1"])
The configuration above will start a Redis pubsub and
register it with name `:my_redis_pubsub`.
You can know use the functions in this module to subscribe
and broadcast messages:
iex> PubSub.subscribe MyApp.PubSub, self, "user:123"
:ok
iex> Process.info(self)[:messages]
[]
iex> PubSub.broadcast MyApp.PubSub, "user:123", {:user_update, %{id: 123, name: "Shane"}}
:ok
iex> Process.info(self)[:messages]
{:user_update, %{id: 123, name: "Shane"}}
## Implementing your own adapter
PubSub adapters run inside their own supervision tree.
If you are interested in providing your own adapter, let's
call it `Phoenix.PubSub.MyQueue`, the first step is to provide
a supervisor module that receives the server name and a bunch
of options on `start_link/2`:
defmodule Phoenix.PubSub.MyQueue do
def start_link(name, options) do
Supervisor.start_link(__MODULE__, {name, options},
name: Module.concat(name, Supervisor))
end
def init({name, options}) do
...
end
end
On `init/1`, you will define the supervision tree and use the given
`name` to register the main pubsub process locally. This process must
be able to handle the following GenServer calls:
* `subscribe` - subscribes the given pid to the given topic
sends: `{:subscribe, pid, topic, opts}`
respond with: `:ok | {:error, reason} | {:perform, {m, f, a}}`
* `unsubscribe` - unsubscribes the given pid from the given topic
sends: `{:unsubscribe, pid, topic}`
respond with: `:ok | {:error, reason} | {:perform, {m, f, a}}`
* `broadcast` - broadcasts a message on the given topic
sends: `{:broadcast, :none | pid, topic, message}`
respond with: `:ok | {:error, reason} | {:perform, {m, f, a}}`
### Offloading work to clients via MFA response
The `Phoenix.PubSub` API allows any of its functions to handle a
response from the adapter matching `{:perform, {m, f, a}}`. The PubSub
client will recursively invoke all MFA responses until a result is
returned. This is useful for offloading work to clients without blocking
your PubSub adapter. See `Phoenix.PubSub.PG2` implementation for examples.
"""
defmodule BroadcastError do
defexception [:message]
def exception(msg) do
%BroadcastError{message: "broadcast failed with #{inspect msg}"}
end
end
@doc """
Subscribes the pid to the PubSub adapter's topic.
* `server` - The Pid registered name of the server
* `pid` - The subscriber pid to receive pubsub messages
* `topic` - The topic to subscribe to, ie: `"users:123"`
* `opts` - The optional list of options. See below.
## Options
* `:link` - links the subscriber to the pubsub adapter
* `:fastlane` - Provides a fastlane path for the broadcasts for
`%Phoenix.Socket.Broadcast{}` events. The fastlane process is
notified of a cached message instead of the normal subscriber.
Fastlane handlers must implement `fastlane/1` callbacks which accepts
a `Phoenix.Socket.Broadcast` structs and returns a fastlaned format
for the handler. For example:
PubSub.subscribe(MyApp.PubSub, self(), "topic1",
fastlane: {fast_pid, Phoenix.Transports.WebSocketSerializer, ["event1"]})
"""
@spec subscribe(atom, pid, binary, Keyword.t) :: :ok | {:error, term}
def subscribe(server, pid, topic, opts \\ []) when is_atom(server),
do: call(server, :subscribe, [pid, topic, opts])
@doc """
Unsubscribes the pid from the PubSub adapter's topic.
"""
@spec unsubscribe(atom, pid, binary) :: :ok | {:error, term}
def unsubscribe(server, pid, topic) when is_atom(server),
do: call(server, :unsubscribe, [pid, topic])
@doc """
Broadcasts message on given topic.
"""
@spec broadcast(atom, binary, term) :: :ok | {:error, term}
def broadcast(server, topic, message) when is_atom(server),
do: call(server, :broadcast, [:none, topic, message])
@doc """
Broadcasts message on given topic.
Raises `Phoenix.PubSub.BroadcastError` if broadcast fails.
"""
@spec broadcast!(atom, binary, term) :: :ok | no_return
def broadcast!(server, topic, message) do
case broadcast(server, topic, message) do
:ok -> :ok
{:error, reason} -> raise BroadcastError, message: reason
end
end
@doc """
Broadcasts message to all but `from_pid` on given topic.
"""
@spec broadcast_from(atom, pid, binary, term) :: :ok | {:error, term}
def broadcast_from(server, from_pid, topic, message) when is_atom(server) and is_pid(from_pid),
do: call(server, :broadcast, [from_pid, topic, message])
@doc """
Broadcasts message to all but `from_pid` on given topic.
Raises `Phoenix.PubSub.BroadcastError` if broadcast fails.
"""
@spec broadcast_from(atom, pid, binary, term) :: :ok | no_return
def broadcast_from!(server, from_pid, topic, message) when is_atom(server) and is_pid(from_pid) do
case broadcast_from(server, from_pid, topic, message) do
:ok -> :ok
{:error, reason} -> raise BroadcastError, message: reason
end
end
defp call(server, kind, args) do
[{^kind, module, head}] = :ets.lookup(server, kind)
apply(module, kind, head ++ args)
end
end