Packages
phoenix
1.7.21
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/transports/long_poll_server.ex
defmodule Phoenix.Transports.LongPoll.Server do
@moduledoc false
use GenServer, restart: :temporary
alias Phoenix.PubSub
def start_link(arg) do
GenServer.start_link(__MODULE__, arg)
end
def init({endpoint, handler, options, params, priv_topic, connect_info}) do
config = %{
endpoint: endpoint,
transport: :longpoll,
options: options,
params: params,
connect_info: connect_info
}
window_ms = Keyword.fetch!(options, :window_ms)
case handler.connect(config) do
{:ok, handler_state} ->
{:ok, handler_state} = handler.init(handler_state)
state = %{
buffer: [],
handler: {handler, handler_state},
window_ms: trunc(window_ms * 1.5),
pubsub_server: endpoint.config(:pubsub_server),
priv_topic: priv_topic,
last_client_poll: now_ms(),
client_ref: nil
}
:ok = PubSub.subscribe(state.pubsub_server, priv_topic, link: true)
schedule_inactive_shutdown(state.window_ms)
{:ok, state}
:error ->
:ignore
{:error, _reason} ->
:ignore
end
end
def handle_info({:dispatch, client_ref, {body, opcode}, ref}, state) do
%{handler: {handler, handler_state}} = state
case handler.handle_in({body, opcode: opcode}, handler_state) do
{:reply, status, {_, reply}, handler_state} ->
state = %{state | handler: {handler, handler_state}}
status = if status == :ok, do: :ok, else: :error
broadcast_from!(state, client_ref, {status, ref})
publish_reply(state, reply)
{:ok, handler_state} ->
state = %{state | handler: {handler, handler_state}}
broadcast_from!(state, client_ref, {:ok, ref})
{:noreply, state}
{:stop, reason, handler_state} ->
state = %{state | handler: {handler, handler_state}}
broadcast_from!(state, client_ref, {:error, ref})
{:stop, reason, state}
end
end
def handle_info({:subscribe, client_ref, ref}, state) do
broadcast_from!(state, client_ref, {:subscribe, ref})
{:noreply, state}
end
def handle_info({:flush, client_ref, ref}, state) do
case state.buffer do
[] ->
{:noreply, %{state | client_ref: {client_ref, ref}, last_client_poll: now_ms()}}
buffer ->
broadcast_from!(state, client_ref, {:messages, Enum.reverse(buffer), ref})
{:noreply, %{state | client_ref: nil, last_client_poll: now_ms(), buffer: []}}
end
end
def handle_info({:expired, client_ref, ref}, state) do
case state.client_ref do
{^client_ref, ^ref} ->
{:noreply, %{state | client_ref: nil}}
_ ->
{:noreply, state}
end
end
def handle_info(:shutdown_if_inactive, state) do
if now_ms() - state.last_client_poll > state.window_ms do
{:stop, {:shutdown, :inactive}, state}
else
schedule_inactive_shutdown(state.window_ms)
{:noreply, state}
end
end
def handle_info(message, state) do
%{handler: {handler, handler_state}} = state
case handler.handle_info(message, handler_state) do
{:push, {_, reply}, handler_state} ->
state = %{state | handler: {handler, handler_state}}
publish_reply(state, reply)
{:ok, handler_state} ->
state = %{state | handler: {handler, handler_state}}
{:noreply, state}
{:stop, reason, handler_state} ->
state = %{state | handler: {handler, handler_state}}
{:stop, reason, state}
end
end
def terminate(reason, state) do
%{handler: {handler, handler_state}} = state
handler.terminate(reason, handler_state)
:ok
end
defp broadcast_from!(state, client_ref, msg) when is_binary(client_ref),
do: PubSub.broadcast_from!(state.pubsub_server, self(), client_ref, msg)
defp broadcast_from!(_state, client_ref, msg) when is_pid(client_ref),
do: send(client_ref, msg)
defp publish_reply(state, reply) when is_map(reply) do
IO.warn(
"Returning a map from the LongPolling serializer is deprecated. " <>
"Please return JSON encoded data instead (see Phoenix.Socket.Serializer)"
)
publish_reply(state, Phoenix.json_library().encode_to_iodata!(reply))
end
defp publish_reply(state, reply) do
notify_client_now_available(state)
{:noreply, update_in(state.buffer, &[IO.iodata_to_binary(reply) | &1])}
end
defp notify_client_now_available(state) do
case state.client_ref do
{client_ref, ref} -> broadcast_from!(state, client_ref, {:now_available, ref})
nil -> :ok
end
end
defp now_ms, do: System.system_time(:millisecond)
defp schedule_inactive_shutdown(window_ms) do
Process.send_after(self(), :shutdown_if_inactive, window_ms)
end
end