Packages

phoenix_kit

2.60.3
2.60.3 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
phoenix_kit lib modules storage upload_inbox.ex
Raw

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