Packages
phoenix_kit
2.59.0
2.60.2
2.60.1
2.60.0
2.59.0
2.58.0
2.57.1
2.57.0
2.56.1
2.56.0
2.55.1
2.55.0
2.54.2
2.54.1
2.54.0
2.53.0
2.52.2
2.52.1
2.52.0
2.51.0
2.50.0
2.49.1
2.49.0
2.48.0
2.47.0
2.46.0
2.45.0
2.44.0
2.43.1
2.43.0
2.42.1
2.42.0
2.41.6
2.41.4
2.41.3
2.41.2
2.41.1
2.41.0
2.40.1
2.40.0
2.39.0
2.38.1
2.38.0
2.37.5
2.37.4
2.37.3
2.37.2
2.37.1
2.37.0
2.36.1
2.36.0
2.35.0
2.34.0
2.33.0
2.32.1
2.32.0
2.31.1
2.31.0
2.30.0
2.29.1
2.29.0
2.28.2
2.28.1
2.28.0
2.27.2
2.27.1
2.27.0
2.26.1
2.26.0
2.25.0
2.24.0
2.23.3
2.23.2
2.23.1
2.23.0
2.22.24
2.22.23
2.22.22
2.22.21
2.22.20
2.22.19
2.22.18
2.22.17
2.22.16
2.22.15
2.22.14
2.22.13
2.22.12
2.22.11
2.22.10
2.22.9
2.22.8
2.22.7
2.22.6
2.22.5
2.22.4
2.22.3
2.22.2
2.22.1
2.22.0
2.21.5
2.21.4
2.21.3
2.21.2
2.21.1
2.21.0
2.20.0
2.19.0
2.18.1
2.18.0
2.17.0
2.16.0
2.15.1
2.15.0
2.14.2
2.14.1
2.14.0
2.13.19
2.13.18
2.13.17
2.13.16
2.13.15
2.13.13
2.13.12
2.13.11
2.13.10
2.13.9
2.13.8
2.13.7
2.13.6
2.13.5
2.13.4
2.13.3
2.13.2
2.13.1
2.13.0
2.12.1
2.12.0
2.11.0
2.10.0
2.9.0
2.8.1
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.1
2.0.0
1.7.236
1.7.235
1.7.234
1.7.233
1.7.232
1.7.231
1.7.230
1.7.229
1.7.228
1.7.227
1.7.226
1.7.225
1.7.224
1.7.223
1.7.222
1.7.221
1.7.220
1.7.219
1.7.218
1.7.217
1.7.216
1.7.215
1.7.214
1.7.213
1.7.212
1.7.211
1.7.210
1.7.209
1.7.208
1.7.207
1.7.206
1.7.205
1.7.204
1.7.203
1.7.202
1.7.201
1.7.200
1.7.199
1.7.198
1.7.197
1.7.196
1.7.194
1.7.193
1.7.192
1.7.191
1.7.190
1.7.189
1.7.187
1.7.186
1.7.185
1.7.184
1.7.183
1.7.182
1.7.181
1.7.180
1.7.179
1.7.178
1.7.177
1.7.176
1.7.175
1.7.174
1.7.173
1.7.172
1.7.171
1.7.170
1.7.169
1.7.168
1.7.167
1.7.166
1.7.165
1.7.164
1.7.162
1.7.161
1.7.160
1.7.159
1.7.157
1.7.156
1.7.155
1.7.154
1.7.153
1.7.152
1.7.151
1.7.150
1.7.149
1.7.146
1.7.145
1.7.144
1.7.143
1.7.138
1.7.133
1.7.132
1.7.131
1.7.130
1.7.128
1.7.126
1.7.125
1.7.121
1.7.120
1.7.119
1.7.118
1.7.117
1.7.116
1.7.115
1.7.114
1.7.113
1.7.112
1.7.111
1.7.110
1.7.109
1.7.108
1.7.107
1.7.106
1.7.105
1.7.104
1.7.103
1.7.102
1.7.101
1.7.100
1.7.99
1.7.98
1.7.97
1.7.96
1.7.95
1.7.94
1.7.93
1.7.92
1.7.91
1.7.90
1.7.89
1.7.88
1.7.87
1.7.86
1.7.85
1.7.84
1.7.83
1.7.82
1.7.81
1.7.80
1.7.79
1.7.78
1.7.77
1.7.76
1.7.75
1.7.74
1.7.71
1.7.70
1.7.69
1.7.66
1.7.65
1.7.64
1.7.63
1.7.62
1.7.61
1.7.59
1.7.58
1.7.57
1.7.56
1.7.55
1.7.54
1.7.53
1.7.52
1.7.51
1.7.49
1.7.44
1.7.43
1.7.42
1.7.41
1.7.39
1.7.38
1.7.37
1.7.36
1.7.34
1.7.33
1.7.31
1.7.30
1.7.29
1.7.28
1.7.27
1.7.26
1.7.25
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.6.20
1.6.19
1.6.18
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.5.2
1.5.1
1.5.0
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.3.2
1.3.1
1.3.0
1.2.10
1.2.9
1.2.8
1.2.7
1.2.5
1.2.4
1.2.2
1.2.1
1.2.0
1.1.0
1.0.0
A foundation for building Elixir Phoenix apps — SaaS, social networks, ERP systems, marketplaces, and more
Current section
Files
Jump to
Current section
Files
lib/modules/storage/upload_inbox.ex
defmodule PhoenixKit.Modules.Storage.UploadInbox do
@moduledoc """
Where an upload's bytes wait between "the server has them" and "they are
stored", so a failure in between is something a person can see and act on.
The media browser used to copy a finished transfer to a temp path, try to
store it, and delete the temp file either way. A failed store became one
anonymous count in a flash ("2 failed") and the bytes were gone; a raise
took the LiveView down with every queued file in it; a page refresh while
files were queued lost them without a word. Someone who dropped ten
pictures and found six had no way to learn which four, or to try again.
Now a received upload is written here first and removed only once it is
stored. One that fails stays, marked with why; one whose processing never
finished (the LiveView died, the page was refreshed) stays too, and reads
as interrupted once it is old enough that nothing can still be working on
it. `list/1` is what the browser's "Upload problems" panel shows, and it
survives a refresh because it lives on disk rather than in a socket.
One directory per user (`<root>/<user_uuid>/`), holding each upload's bytes
at `<id>` and what is known about it at `<id>.json`. The root defaults to
the system temp dir — this is a waiting room, not storage — and can be
moved with `config :phoenix_kit, :upload_inbox_dir, "/path"`. Anything
older than a week is swept on the next listing, and all users' inboxes are
swept at most once an hour when an upload arrives.
## Where an upload belongs
An item remembers the destination it was dropped on (`library_uuid` and
`folder_uuid`, recorded by the first browser that claims it) and who is
working on it (`owner`, the LiveView's pid). A retry goes to the recorded
destination, never to whichever browser the click happened in — a private
library's file must not land in Media — and a panel lists only the items of
its own library. An item whose owner is still alive is live, whatever its
age; one whose owner is gone reads as interrupted at once. Retry and
Discard from a panel that has gone stale go through `claim/4` and
`discard/2`, which refuse an item another process is working on.
## Persistence
The default root is node-local temporary storage: a multi-node deployment,
or a container with an ephemeral filesystem, loses "survives a refresh"
when the reconnect lands elsewhere or the host restarts. Point
`:upload_inbox_dir` at storage every node shares if that matters. Directories
are created `0700`. There is no per-user quota yet.
"""
require Logger
@uuid ~r/\A[0-9a-f]{8}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{12}\z/i
# A "received" item younger than this may still be in a drain queue
# somewhere; older, nothing is working on it any more.
@interrupted_after_s 120
@keep_for_s 7 * 24 * 3600
@type item :: %{
id: String.t(),
client_name: String.t(),
client_type: String.t() | nil,
client_size: non_neg_integer() | nil,
status: String.t(),
error: String.t() | nil,
received_at: integer(),
dest: %{library_uuid: String.t() | nil, folder_uuid: String.t() | nil} | nil,
owner: String.t() | nil
}
@doc "The inbox root."
def root,
do:
Application.get_env(:phoenix_kit, :upload_inbox_dir) ||
Path.join(System.tmp_dir!(), "phoenix_kit_upload_inbox")
@doc """
Moves `src` into `user_uuid`'s inbox as a received upload and returns the
item. `meta` carries the client's `:client_name`, `:client_type` and
`:client_size`. `{:error, reason}` when the user is not a uuid or the disk
refuses — the caller then falls back to its own temp file, which is still
where it was: a failure here never costs the bytes.
`copy: true` leaves `src` in place (a second browser's own item).
"""
@spec put(String.t(), Path.t(), map(), keyword()) ::
{:ok, item(), Path.t()} | {:error, term()}
def put(user_uuid, src, meta, opts \\ []) do
with {:ok, dir} <- user_dir(user_uuid),
:ok <- File.mkdir_p(dir) do
File.chmod(root(), 0o700)
File.chmod(dir, 0o700)
maybe_sweep_all()
id = Ecto.UUID.generate()
dest = Path.join(dir, id)
item = %{
id: id,
client_name: to_string(meta[:client_name] || "upload"),
client_type: meta[:client_type],
client_size: meta[:client_size],
status: "received",
error: nil,
received_at: System.system_time(:second),
dest: nil,
owner: nil
}
store = if opts[:copy], do: &File.cp/2, else: &move/2
with :ok <- store.(src, dest),
:ok <- write_meta(dir, item) do
{:ok, item, dest}
else
error ->
give_back(src, dest, opts[:copy])
error
end
else
{:error, _} = error -> error
end
end
# The bytes are in the inbox but their record could not be written: put them
# back where they came from, and only then forget the inbox copy. When that is
# not possible the inbox copy stays — a copy without a record still beats no
# copy — and the caller carries on with the error.
defp give_back(_src, dest, true), do: File.rm(dest)
defp give_back(src, dest, _copy) do
if File.exists?(src) do
File.rm(dest)
else
move(dest, src)
end
end
@doc "The inbox path of `path`'s item as `{user_uuid, id}`, or nil if `path` is not in an inbox."
@spec locate(Path.t()) :: {String.t(), String.t()} | nil
def locate(path) when is_binary(path) do
root = Path.expand(root())
expanded = Path.expand(path)
with true <- String.starts_with?(expanded, root <> "/"),
[user_uuid, id] <- Path.split(Path.relative_to(expanded, root)),
true <- Regex.match?(@uuid, user_uuid) and Regex.match?(@uuid, id) do
{user_uuid, id}
else
_ -> nil
end
end
def locate(_), do: nil
@doc "Marks an item failed, keeping its bytes for a retry."
@spec fail(String.t(), String.t(), String.t()) :: :ok | {:error, term()}
def fail(user_uuid, id, reason) do
update(user_uuid, id, &%{&1 | status: "failed", error: reason, owner: nil})
end
@doc """
Takes an item to work on, for the calling process: it reads as live while
this process is, and `dest` (`%{library_uuid: _, folder_uuid: _}`) becomes
its destination unless it already has one. `{:ok, item}` carries the item as
it now stands — a retry stores it at `item.dest`, not where the click was.
`{:error, :busy}` when another live process holds it (a stale panel retried
an item a second tab has already taken), `{:error, :not_found}` when its
bytes are gone.
"""
@spec claim(String.t(), String.t(), map() | nil, pid()) :: {:ok, item()} | {:error, term()}
def claim(user_uuid, id, dest, owner \\ self()) do
with_lock(id, fn ->
with {:ok, item} <- read(user_uuid, id),
false <- busy?(item, owner),
claimed = %{
item
| status: "received",
error: nil,
received_at: now(),
owner: encode_owner(owner),
dest: item.dest || dest
},
:ok <- write_item(user_uuid, claimed) do
{:ok, claimed}
else
true -> {:error, :busy}
{:error, _} = error -> error
end
end)
end
@doc """
An item as it stands, or `nil` when it is gone. Its `dest` is where a retry
must put it.
"""
@spec get(String.t(), String.t()) :: item() | nil
def get(user_uuid, id) do
case read(user_uuid, id) do
{:ok, item} -> item
_ -> nil
end
end
@doc """
Removes an item the user let go of. `{:error, :busy}` — and nothing removed —
while another live process is working on it: a panel that went stale must not
delete bytes a second tab is storing.
"""
@spec discard(String.t(), String.t(), pid()) :: :ok | {:error, :busy}
def discard(user_uuid, id, owner \\ self()) do
with_lock(id, fn ->
case read(user_uuid, id) do
{:ok, item} ->
if busy?(item, owner), do: {:error, :busy}, else: delete(user_uuid, id)
_ ->
:ok
end
end)
end
@doc "Marks an item as being worked on again (a retry), so it reads as live, not interrupted."
@spec touch(String.t(), String.t()) :: :ok | {:error, term()}
def touch(user_uuid, id) do
update(user_uuid, id, &%{&1 | status: "received", error: nil, received_at: now()})
end
@doc "Removes an item — it was stored, or the user discarded it."
@spec delete(String.t(), String.t()) :: :ok
def delete(user_uuid, id) do
case item_paths(user_uuid, id) do
{:ok, bytes, meta} ->
File.rm(bytes)
File.rm(meta)
:ok
_ ->
:ok
end
end
@doc "The bytes of an item, when it is still here."
@spec path(String.t(), String.t()) :: {:ok, Path.t()} | :error
def path(user_uuid, id) do
case item_paths(user_uuid, id) do
{:ok, bytes, _meta} -> if File.exists?(bytes), do: {:ok, bytes}, else: :error
_ -> :error
end
end
@doc """
The items that need the user's attention, oldest first: failed ones, and
received ones old enough that their processing was clearly interrupted
(those come back with status `"interrupted"`). Sweeps week-old items.
"""
@spec list(String.t() | nil) :: [item()]
def list(user_uuid) do
case user_dir(user_uuid) do
{:ok, dir} ->
now = now()
dir
|> Path.join("*.json")
|> Path.wildcard()
|> Enum.flat_map(&read_item(&1, dir, now))
|> Enum.sort_by(& &1.received_at)
_ ->
[]
end
rescue
error ->
Logger.warning("[UploadInbox] list failed: #{Exception.message(error)}")
[]
end
@doc """
Whether an upload is in somebody's hands right now: received, and either
claimed by a live process or too young to call abandoned. A page uses it to
keep looking while that is so (see `list/1`).
"""
@spec waiting?(String.t() | nil) :: boolean()
def waiting?(user_uuid) do
case user_dir(user_uuid) do
{:ok, dir} ->
now = now()
dir
|> Path.join("*.json")
|> Path.wildcard()
|> Enum.any?(&live_meta?(&1, now))
_ ->
false
end
end
defp live_meta?(meta_path, now) do
with {:ok, json} <- File.read(meta_path),
{:ok, map} <- Jason.decode(json),
%{status: "received"} = item <- decode(map) do
if item.owner,
do: owner_state(item, now) == :alive,
else: now - item.received_at <= @interrupted_after_s
else
_ -> false
end
end
defp read_item(meta_path, dir, now) do
with {:ok, json} <- File.read(meta_path),
{:ok, map} <- Jason.decode(json),
%{} = item <- decode(map),
true <- File.exists?(Path.join(dir, item.id)) do
cond do
now - item.received_at > @keep_for_s ->
File.rm(Path.join(dir, item.id))
File.rm(meta_path)
[]
item.status == "failed" ->
[item]
# Claimed: live while its owner is, interrupted the moment it is not —
# a refresh or a crash does not make anyone wait two minutes.
item.owner != nil ->
case owner_state(item, now) do
:alive -> []
:gone -> [%{item | status: "interrupted"}]
end
# Not claimed by anyone yet: it may be on its way to a browser.
now - item.received_at > @interrupted_after_s ->
[%{item | status: "interrupted"}]
true ->
[]
end
else
# Bytes gone (or the sidecar unreadable): nothing left to retry.
_ ->
File.rm(meta_path)
[]
end
end
# An owner on this node is checked; one on another node cannot be, so it is
# taken as alive until the item is old enough that nothing could still be
# working on it.
defp owner_state(%{owner: owner, received_at: received_at}, now) do
case decode_owner(owner) do
pid when is_pid(pid) and node(pid) == node() ->
if Process.alive?(pid), do: :alive, else: :gone
pid when is_pid(pid) ->
if now - received_at > @interrupted_after_s, do: :gone, else: :alive
_ ->
:gone
end
end
defp busy?(%{owner: nil}, _me), do: false
defp busy?(%{owner: owner} = item, me) do
case decode_owner(owner) do
^me -> false
_ -> item.status == "received" and owner_state(item, now()) == :alive
end
end
defp encode_owner(pid) when is_pid(pid), do: pid |> :erlang.pid_to_list() |> List.to_string()
defp decode_owner(owner) when is_binary(owner) do
owner |> String.to_charlist() |> :erlang.list_to_pid()
rescue
ArgumentError -> nil
end
defp decode_owner(_), do: nil
defp decode(%{"id" => id} = map) when is_binary(id) do
%{
id: id,
client_name: map["client_name"] || "upload",
client_type: map["client_type"],
client_size: map["client_size"],
status: map["status"] || "received",
error: map["error"],
received_at: map["received_at"] || 0,
dest: decode_dest(map["dest"]),
owner: if(is_binary(map["owner"]), do: map["owner"])
}
end
defp decode(_), do: nil
defp decode_dest(%{} = dest),
do: %{library_uuid: dest["library_uuid"], folder_uuid: dest["folder_uuid"]}
defp decode_dest(_), do: nil
defp read(user_uuid, id) do
with {:ok, _bytes, meta_path} <- item_paths(user_uuid, id),
{:ok, json} <- File.read(meta_path),
{:ok, map} <- Jason.decode(json),
%{} = item <- decode(map) do
{:ok, item}
else
_ -> {:error, :not_found}
end
end
defp update(user_uuid, id, fun) do
with {:ok, item} <- read(user_uuid, id) do
write_item(user_uuid, fun.(item))
end
end
defp write_item(user_uuid, item) do
with {:ok, dir} <- user_dir(user_uuid), do: write_meta(dir, item)
end
# One claim at a time per item, across nodes.
defp with_lock(id, fun), do: :global.trans({{__MODULE__, id}, self()}, fun)
defp item_paths(user_uuid, id) do
with {:ok, dir} <- user_dir(user_uuid),
true <- is_binary(id) and Regex.match?(@uuid, id) do
{:ok, Path.join(dir, id), Path.join(dir, id <> ".json")}
else
_ -> :error
end
end
# The uuid check is the path-traversal guard: nothing but a uuid ever
# becomes a directory or file name here.
defp user_dir(user_uuid) when is_binary(user_uuid) do
if Regex.match?(@uuid, user_uuid),
do: {:ok, Path.join(root(), String.downcase(user_uuid))},
else: {:error, :invalid_user}
end
defp user_dir(_), do: {:error, :invalid_user}
defp write_meta(dir, item) do
# Written beside and renamed over: a reader (`list/1`) must never see a
# half-written sidecar, which it would take for a broken one and remove.
meta = Path.join(dir, item.id <> ".json")
tmp = meta <> ".#{System.unique_integer([:positive])}.tmp"
with :ok <- File.write(tmp, Jason.encode!(item)),
:ok <- File.rename(tmp, meta) do
:ok
else
error ->
File.rm(tmp)
error
end
end
# Every user's inbox, swept at most once an hour — a user who never comes
# back has no listing of their own to do it.
defp maybe_sweep_all do
key = {__MODULE__, :last_sweep}
last = :persistent_term.get(key, 0)
if now() - last > 3600 do
:persistent_term.put(key, now())
for dir <- Path.wildcard(Path.join(root(), "*")), File.dir?(dir) do
dir |> Path.join("*.json") |> Path.wildcard() |> Enum.each(&read_item(&1, dir, now()))
sweep_orphans(dir)
end
end
rescue
_ -> :ok
end
# Bytes with no record (a record that could not be written, or an old-format
# copy): nothing can show or retry them, so they go once a week has passed.
defp sweep_orphans(dir) do
cutoff = now() - @keep_for_s
for path <- Path.wildcard(Path.join(dir, "*")),
Regex.match?(@uuid, Path.basename(path)),
not File.exists?(path <> ".json"),
{:ok, %{mtime: mtime}} <- [File.stat(path, time: :posix)],
mtime < cutoff do
File.rm(path)
end
end
# A rename when source and inbox share a filesystem (the usual case — both
# under the temp dir); a copy otherwise.
defp move(src, dest) do
case File.rename(src, dest) do
:ok ->
:ok
{:error, _} ->
with :ok <- File.cp(src, dest) do
File.rm(src)
:ok
end
end
end
defp now, do: System.system_time(:second)
end