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
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.
  """

  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