Packages
phoenix_kit
2.57.0
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.
"""
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()
}
@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.
"""
@spec put(String.t(), Path.t(), map()) :: {:ok, item(), Path.t()} | {:error, term()}
def put(user_uuid, src, meta) do
with {:ok, dir} <- user_dir(user_uuid),
:ok <- File.mkdir_p(dir) do
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)
}
with :ok <- move(src, dest),
:ok <- write_meta(dir, item) do
{:ok, item, dest}
else
error ->
File.rm(dest)
error
end
else
{:error, _} = error -> error
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})
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
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]
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
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
}
end
defp decode(_), do: nil
defp update(user_uuid, id, fun) 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
write_meta(Path.dirname(meta_path), fun.(item))
else
_ -> {:error, :not_found}
end
end
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 <> ".tmp"
with :ok <- File.write(tmp, Jason.encode!(item)),
:ok <- File.rename(tmp, meta) do
:ok
else
error ->
File.rm(tmp)
error
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