Packages
phoenix_live_view
0.20.13
1.2.7
1.2.6
1.2.5
1.2.4
1.2.3
1.2.2
1.2.1
1.2.0
1.2.0-rc.3
1.2.0-rc.2
1.2.0-rc.1
1.2.0-rc.0
1.1.32
1.1.31
1.1.30
1.1.29
1.1.28
1.1.27
1.1.26
1.1.25
1.1.24
1.1.23
1.1.22
1.1.21
1.1.20
1.1.19
1.1.18
1.1.17
1.1.16
1.1.15
1.1.14
1.1.13
1.1.12
1.1.11
1.1.10
1.1.9
1.1.8
1.1.7
1.1.6
retired
1.1.5
1.1.4
1.1.3
1.1.2
1.1.1
1.1.0
1.1.0-rc.4
1.1.0-rc.3
1.1.0-rc.2
1.1.0-rc.1
1.1.0-rc.0
1.0.18
1.0.17
1.0.16
1.0.15
1.0.14
1.0.13
1.0.12
1.0.11
1.0.10
1.0.9
1.0.8
retired
1.0.7
1.0.6
retired
1.0.5
1.0.4
1.0.3
1.0.2
1.0.1
1.0.0
1.0.0-rc.9
1.0.0-rc.8
1.0.0-rc.7
1.0.0-rc.6
1.0.0-rc.5
1.0.0-rc.4
1.0.0-rc.3
1.0.0-rc.2
1.0.0-rc.1
1.0.0-rc.0
0.20.17
0.20.16
0.20.15
0.20.14
0.20.13
0.20.12
0.20.11
0.20.10
0.20.9
0.20.8
0.20.7
0.20.6
0.20.5
0.20.4
0.20.3
0.20.2
0.20.1
0.20.0
0.19.5
0.19.4
0.19.3
0.19.2
0.19.1
0.19.0
0.18.18
0.18.17
0.18.16
0.18.15
0.18.14
0.18.13
0.18.12
0.18.11
0.18.10
0.18.9
0.18.8
0.18.7
0.18.6
0.18.5
0.18.4
0.18.3
0.18.2
0.18.1
0.18.0
0.17.14
0.17.13
0.17.12
0.17.11
0.17.10
0.17.9
0.17.8
0.17.7
0.17.6
0.17.5
0.17.4
0.17.3
0.17.2
0.17.1
0.17.0
0.16.4
0.16.3
0.16.2
0.16.1
0.16.0
0.15.7
0.15.6
0.15.5
0.15.4
0.15.3
0.15.2
0.15.1
0.15.0
0.14.8
0.14.7
0.14.6
0.14.5
0.14.4
0.14.3
0.14.2
0.14.1
0.14.0
0.13.3
0.13.2
0.13.1
0.13.0
0.12.1
0.12.0
0.11.1
0.11.0
0.10.0
0.9.0
0.8.1
0.8.0
0.7.1
0.7.0
0.6.0
0.6.0-dev
0.5.2
0.5.1
0.5.0
0.4.1
0.4.0
0.3.1
0.3.0
0.2.1
0.2.0
0.1.1
0.1.0
Rich, real-time user experiences with server-rendered HTML
Current section
Files
Jump to
Current section
Files
lib/phoenix_live_view/upload_channel.ex
defmodule Phoenix.LiveView.UploadChannel do
@moduledoc false
use Phoenix.Channel, log_handle_in: false
@timeout :infinity
require Logger
alias Phoenix.LiveView.{Static, Channel}
def cancel(pid) do
GenServer.call(pid, :cancel, @timeout)
end
def consume(pid, entry, func) when is_function(func, 1) or is_function(func, 2) do
case GenServer.call(pid, :consume_start, @timeout) do
{:ok, file_meta} ->
try do
result =
cond do
is_function(func, 1) -> func.(file_meta)
is_function(func, 2) -> func.(file_meta, entry)
end
case result do
{:ok, return} ->
GenServer.call(pid, :consume_done, @timeout)
return
{:postpone, return} ->
return
return ->
IO.warn("""
consuming uploads requires a return signature matching:
{:ok, value} | {:postpone, value}
got:
#{inspect(return)}
""")
GenServer.call(pid, :consume_done, @timeout)
return
end
rescue
exception ->
GenServer.call(pid, :consume_done, @timeout)
reraise(exception, __STACKTRACE__)
end
{:error, :in_progress} ->
raise RuntimeError, "cannot consume uploaded file that is still in progress"
end
end
@impl true
def join(_topic, auth_payload, socket) do
%{"token" => token} = auth_payload
with {:ok, %{pid: pid, ref: ref, cid: cid}} <- Static.verify_token(socket.endpoint, token),
{:ok, config} <- Channel.register_upload(pid, ref, cid),
%{max_file_size: max_file_size, chunk_timeout: chunk_timeout} = config,
{writer, writer_opts} <- config.writer,
{:ok, writer_state} <- writer.init(writer_opts) do
Process.monitor(pid)
Process.flag(:trap_exit, true)
socket =
assign(socket, %{
writer: writer,
writer_state: writer_state,
live_view_pid: pid,
max_file_size: max_file_size,
chunk_timeout: chunk_timeout,
chunk_timer: nil,
writer_closed?: false,
done?: false,
uploaded_size: 0
})
{:ok, socket}
else
{:error, reason} when reason in [:expired, :invalid] ->
{:error, %{reason: :invalid_token}}
{:error, reason} when reason in [:already_registered, :disallowed] ->
{:error, %{reason: reason}}
# writer init error
{:error, _reason} ->
{:error, %{reason: :writer_error}}
end
end
@impl true
def handle_in("chunk", {:binary, payload}, socket) do
%{uploaded_size: uploaded_size, max_file_size: max_file_size} = socket.assigns
socket = reschedule_chunk_timer(socket)
if !socket.assigns.writer_closed? and byte_size(payload) + uploaded_size <= max_file_size do
case write_bytes(socket, payload) do
{:ok, new_socket} ->
{:reply, :ok, new_socket}
{:error, reason, new_socket} ->
new_socket =
case close_writer(new_socket, {:error, reason}) do
{:ok, new_socket} -> new_socket
{:error, _reason, new_socket} -> new_socket
end
Channel.report_writer_error(socket.assigns.live_view_pid, reason)
{:reply, {:error, %{reason: :writer_error}}, new_socket}
end
else
reply = %{reason: :file_size_limit_exceeded, limit: max_file_size}
{:stop, {:shutdown, :closed}, {:error, reply}, socket}
end
end
@impl true
def handle_info({:EXIT, _pid, reason}, socket) do
{:stop, reason, socket}
end
def handle_info(
{:DOWN, _, _, live_view_pid, reason},
%{assigns: %{live_view_pid: live_view_pid}} = socket
) do
reason = if reason == :normal, do: {:shutdown, :closed}, else: reason
{:stop, reason, maybe_cancel_writer(socket)}
end
def handle_info(:chunk_timeout, socket) do
{:stop, {:shutdown, :closed}, socket}
end
@impl true
def handle_call(:consume_start, _from, socket) do
if socket.assigns.done? do
{:reply, {:ok, file_meta(socket)}, socket}
else
{:reply, {:error, :in_progress}, socket}
end
end
@impl true
def handle_call(:consume_done, from, socket) do
GenServer.reply(from, :ok)
{:stop, {:shutdown, :closed}, socket}
end
def handle_call(:cancel, from, socket) do
if socket.assigns.writer_closed? do
GenServer.reply(from, :ok)
{:stop, {:shutdown, :closed}, socket}
else
case close_writer(socket, :cancel) do
{:ok, new_socket} ->
GenServer.reply(from, :ok)
{:stop, {:shutdown, :closed}, new_socket}
{:error, reason, new_socket} ->
GenServer.reply(from, {:error, reason})
{:stop, {:shutdown, :closed}, new_socket}
end
end
end
@impl true
def terminate(_reason, socket) do
_ = maybe_cancel_writer(socket)
:ok
end
defp reschedule_chunk_timer(socket) do
cancel_timer(socket.assigns.chunk_timer, :chunk_timeout)
new_timer = Process.send_after(self(), :chunk_timeout, socket.assigns.chunk_timeout)
assign(socket, :chunk_timer, new_timer)
end
defp cancel_timer(nil = _timer, _msg), do: :ok
defp cancel_timer(timer, msg) do
if Process.cancel_timer(timer) do
:ok
else
receive do
^msg -> :ok
after
0 -> :ok
end
end
end
defp write_bytes(socket, payload) do
case socket.assigns.writer.write_chunk(payload, socket.assigns.writer_state) do
{:ok, writer_state} ->
socket
|> assign(:uploaded_size, socket.assigns.uploaded_size + byte_size(payload))
|> assign(:writer_state, writer_state)
|> maybe_close_completed_file()
{:error, reason, writer_state} ->
cancel_timer(socket.assigns.chunk_timer, :chunk_timeout)
{:error, reason, assign(socket, writer_state: writer_state, chunk_timer: nil)}
end
end
defp maybe_close_completed_file(socket) do
if socket.assigns.uploaded_size == socket.assigns.max_file_size do
case close_writer(socket, :done) do
{:ok, socket} -> {:ok, assign(socket, done?: true)}
{:error, reason, new_socket} -> {:error, reason, new_socket}
end
else
{:ok, socket}
end
end
# we need to handle the case where socket assigns aren't set yet because
# we are trapping exits and may enter terminate before joining is complete
defp maybe_cancel_writer(socket) do
case socket.assigns do
%{writer_closed?: false} ->
case close_writer(socket, :cancel) do
{:ok, new_socket} -> new_socket
{:error, _reason, new_socket} -> new_socket
end
%{} ->
socket
end
end
defp close_writer(socket, reason) do
cancel_timer(socket.assigns.chunk_timer, :chunk_timeout)
socket = assign(socket, chunk_timer: nil, writer_closed?: true)
case socket.assigns.writer.close(socket.assigns.writer_state, reason) do
{:ok, writer_state} ->
{:ok,
socket
|> assign(writer_state: writer_state)
|> garbage_collect()}
{:error, reason} ->
{:error, reason, socket}
end
end
defp garbage_collect(socket) do
send(socket.transport_pid, :garbage_collect)
:erlang.garbage_collect(self())
socket
end
defp file_meta(socket), do: socket.assigns.writer.meta(socket.assigns.writer_state)
end