Packages

phoenix_kit

2.60.2
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 bucket_log.ex
Raw

lib/modules/storage/bucket_log.ex

defmodule PhoenixKit.Modules.Storage.BucketLog do
  @moduledoc """
  What went wrong with a site storage bucket, and what its probes said (V208).

  `Manager` reports a write, read or delete that failed (`record_failure/3`),
  and `Storage.probe_bucket/1` reports every probe of a saved bucket
  (`record/4`). The bucket's own page reads it back (`recent/2`, `summary/1`).

    * **Failures only.** A successful write or read is not recorded — on a busy
      site that would be a row per request. Latency over time therefore comes
      from the probes.
    * **A miss is not a failure.** A bucket that does not hold an object is the
      normal case during failover; a "not found" reply is not logged.
    * **A repeat is one row.** The same failure of the same kind on the same
      bucket within a minute from its first occurrence bumps `count` and `last_at`
      instead of adding a row,
      so a bucket that is down does not write a row per request.
    * **Site buckets only.** A user's own bucket (V206) is private to them and is
      never logged here; neither is a bucket that has no uuid yet (a form's
      unsaved test).
    * **Never in the way.** Writing is best effort: `record_failure/3` runs in
      the background (a dedicated supervisor capped at four tasks), swallows
      every error, and
      never queries the database in the read or write that reported it. If the
      supervisor is absent or full, the diagnostic is dropped.

  Rows are pruned daily to `bucket_log_retention_days` (default 30)
  by `Workers.BucketLogPruneWorker`.
  """

  import Ecto.Query

  require Logger

  alias PhoenixKit.Modules.Storage.BucketLogEntry
  alias PhoenixKit.Settings

  @kinds ~w(probe write read delete)
  @merge_window_seconds 60
  @message_limit 500
  @default_retention_days 30
  @latency_points 20

  @doc "The kinds of entry."
  @spec kinds() :: [String.t()]
  def kinds, do: @kinds

  @doc """
  Reports that `kind` (`"write"`, `"read"` or `"delete"`) failed on `bucket`,
  with the provider's `reason`. Returns `:ok` at once; the row is written in the
  background. A missing object is ignored for reads and deletes, but never for writes.
  """
  @spec record_failure(map() | nil, String.t(), term()) :: :ok
  def record_failure(bucket, kind, reason) when kind in @kinds do
    if loggable?(bucket) and not (kind in ~w(read delete) and not_found?(reason)) do
      message = safe_message(reason)
      # Keep credentials out of the background task as well.
      bucket = %{uuid: bucket.uuid, owner_uuid: nil}
      run(fn -> record(bucket, kind, false, message: message) end)
    end

    :ok
  rescue
    _ -> :ok
  catch
    :exit, _ -> :ok
  end

  def record_failure(_bucket, _kind, _reason), do: :ok

  @doc """
  Writes an entry now: `:ok`, or `:skipped` for a bucket that is not logged or a
  write that did not happen (the table is unreachable). Options: `:message`,
  `:latency_ms`. A repeat of a failure that is not a probe, within a minute,
  bumps the earlier row instead.
  """
  @spec record(map() | nil, String.t(), boolean(), keyword()) :: :ok | :skipped
  def record(bucket, kind, ok?, opts \\ []) when kind in @kinds and is_boolean(ok?) do
    if loggable?(bucket), do: write(bucket.uuid, kind, ok?, opts), else: :skipped
  rescue
    error ->
      Logger.debug("BucketLog: not written (#{inspect(error.__struct__)})")
      :skipped
  catch
    :exit, _ -> :skipped
  end

  defp write(bucket_uuid, kind, ok?, opts) do
    now = NaiveDateTime.utc_now() |> NaiveDateTime.truncate(:second)
    message = opts |> Keyword.get(:message) |> safe_message()

    merged? = not ok? and kind != "probe" and merge(bucket_uuid, kind, message, now) > 0

    unless merged? do
      repo().insert!(%BucketLogEntry{
        bucket_uuid: bucket_uuid,
        kind: kind,
        ok: ok?,
        latency_ms: opts[:latency_ms],
        message: message,
        count: 1,
        inserted_at: now,
        last_at: now
      })
    end

    :ok
  end

  # Bumps the newest row of the same failure seen within the window.
  defp merge(bucket_uuid, kind, message, now) do
    since = NaiveDateTime.add(now, -@merge_window_seconds, :second)
    text = message || ""

    newest =
      from(e in BucketLogEntry,
        where:
          e.bucket_uuid == ^bucket_uuid and e.kind == ^kind and e.ok == false and
            coalesce(e.message, "") == ^text and e.inserted_at >= ^since and e.last_at >= ^since,
        order_by: [desc: e.last_at],
        limit: 1,
        select: e.uuid
      )

    {count, _} =
      repo().update_all(
        from(e in BucketLogEntry, where: e.uuid in subquery(newest)),
        inc: [count: 1],
        set: [last_at: now]
      )

    count
  end

  @doc """
  The bucket's entries, newest first: `%{entries, page, total_pages, total}`.
  Options: `:page` (1), `:per_page` (10), `:filter` (`:all`, or `:failures`).
  """
  @spec recent(term(), keyword()) :: map()
  def recent(bucket_uuid, opts \\ []) do
    page = max(Keyword.get(opts, :page, 1), 1)
    per_page = Keyword.get(opts, :per_page, 10)

    base = from(e in BucketLogEntry, where: e.bucket_uuid == ^bucket_uuid)
    base = if opts[:filter] == :failures, do: where(base, [e], e.ok == false), else: base

    total = repo().aggregate(base, :count)
    total_pages = max(ceil(total / per_page), 1)
    page = min(page, total_pages)

    entries =
      base
      |> order_by([e], desc: e.last_at, desc: e.uuid)
      |> limit(^per_page)
      |> offset(^((page - 1) * per_page))
      |> repo().all()

    %{entries: entries, page: page, total_pages: total_pages, total: total}
  end

  @doc """
  The bucket's log at a glance: `%{failures_24h, last_failure, last_probe,
  probes}`. `failures_24h` counts every failure seen in the last day (a row
  that merged repeats counts them all, with at most one minute of boundary
  approximation); `probes` are the latest probes, oldest
  first, for the latency chart.
  """
  @spec summary(term()) :: map()
  def summary(bucket_uuid) do
    since = NaiveDateTime.add(NaiveDateTime.utc_now(), -86_400, :second)
    mine = from(e in BucketLogEntry, where: e.bucket_uuid == ^bucket_uuid)

    failures =
      repo().one(
        from(e in mine,
          where: e.ok == false and e.last_at >= ^since,
          select: sum(e.count)
        )
      )

    last_failure =
      repo().one(
        from(e in mine, where: e.ok == false, order_by: [desc: e.last_at, desc: e.uuid], limit: 1)
      )

    probes =
      from(e in mine,
        where: e.kind == "probe",
        order_by: [desc: e.last_at, desc: e.uuid],
        limit: @latency_points
      )
      |> repo().all()

    %{
      failures_24h: to_int(failures),
      last_failure: last_failure,
      last_probe: List.first(probes),
      probes: Enum.reverse(probes)
    }
  end

  @doc "Removes every entry of a bucket (it was deleted). Returns how many."
  @spec delete_for_bucket(term()) :: non_neg_integer()
  def delete_for_bucket(bucket_uuid) do
    {count, _} =
      repo().delete_all(from(e in BucketLogEntry, where: e.bucket_uuid == ^bucket_uuid))

    count
  rescue
    _ -> 0
  catch
    :exit, _ -> 0
  end

  @doc "How long entries are kept, in days (`bucket_log_retention_days`, default 30)."
  @spec retention_days() :: pos_integer()
  def retention_days do
    case Integer.parse(
           Settings.get_setting("bucket_log_retention_days", "#{@default_retention_days}")
         ) do
      {days, _} when days > 0 -> days
      _ -> @default_retention_days
    end
  end

  @doc "Deletes the entries last seen before the retention. Returns how many."
  @spec prune() :: non_neg_integer()
  def prune do
    cutoff = NaiveDateTime.add(NaiveDateTime.utc_now(), -retention_days() * 86_400, :second)
    {count, _} = repo().delete_all(from(e in BucketLogEntry, where: e.last_at < ^cutoff))
    count
  end

  # ---- internals ----

  # A site bucket that is saved. A user's own bucket (owner_uuid) is theirs.
  defp loggable?(%{uuid: uuid, owner_uuid: nil}) when not is_nil(uuid), do: true
  defp loggable?(_bucket), do: false

  @doc false
  def task_supervisor, do: __MODULE__.TaskSupervisor

  defp run(fun) do
    # No synchronous fallback: even a rescued SQL error can poison a caller's
    # transaction. max_children bounds tasks waiting on a dead DB pool.
    case Process.whereis(task_supervisor()) do
      nil -> :ok
      supervisor -> Task.Supervisor.start_child(supervisor, fun)
    end

    :ok
  end

  # Only reads/deletes treat a missing object as normal. Writes reporting a
  # missing source, directory or bucket are real failures.
  defp not_found?({:http_error, 404, body}), do: not (reason_text(body) =~ ~r/\bNoSuchBucket\b/)
  defp not_found?({:http_error, _status, _}), do: false
  defp not_found?(:enoent), do: true
  defp not_found?("NoSuchKey"), do: true

  defp not_found?(reason) do
    text = reason_text(reason)

    case Regex.run(~r/\{:http_error,\s*([45]\d{2})\s*,/, text) do
      [_, "404"] -> not (text =~ ~r/\bNoSuchBucket\b/)
      [_, _status] -> false
      nil -> text =~ ~r/(?<![a-z0-9_]):enoent\b/i
    end
  end

  @doc """
  A controlled diagnostic for persistence and display. Provider response
  bodies, headers, URLs, keys and arbitrary exception text are never copied.
  Unknown errors deliberately lose detail; server-side provider diagnostics
  remain the place to investigate them.
  """
  @spec safe_message(term()) :: String.t() | nil
  def safe_message(nil), do: nil
  def safe_message({:error, reason}), do: safe_message(reason)
  def safe_message({:http_error, status, _}) when status in 400..599, do: "HTTP #{status}"

  def safe_message(reason) do
    text = reason_text(reason)

    cond do
      # These are static messages, never a substring of the provider response.
      text in ["Disk full", "timeout", "Timed out", "Access denied", "denied"] ->
        text

      text =~ ~r/\{:http_error,\s*([45]\d{2})\s*,/ ->
        [_, status] = Regex.run(~r/\{:http_error,\s*([45]\d{2})\s*,/, text)
        "HTTP #{status}"

      text =~ ~r/\bHTTP [45]\d{2}\b/ ->
        Regex.run(~r/\bHTTP [45]\d{2}\b/, text) |> hd()

      text =~
          ~r/\b(?:AccessDenied|InvalidAccessKeyId|SignatureDoesNotMatch|ExpiredToken|NoSuchBucket|NoSuchKey)\b/ ->
        Regex.run(
          ~r/\b(?:AccessDenied|InvalidAccessKeyId|SignatureDoesNotMatch|ExpiredToken|NoSuchBucket|NoSuchKey)\b/,
          text
        )
        |> hd()

      text =~
          ~r/(?<![a-z0-9_]):(?:enoent|enospc|eacces|eperm|enotdir|eisdir|eio|erofs|etimedout|econnrefused)\b/i ->
        Regex.run(
          ~r/(?<![a-z0-9_]):(?:enoent|enospc|eacces|eperm|enotdir|eisdir|eio|erofs|etimedout|econnrefused)\b/i,
          text
        )
        |> hd()
        |> String.downcase()

      true ->
        "Storage operation failed"
    end
  end

  defp reason_text(reason) when is_binary(reason), do: String.slice(reason, 0, 4096)
  defp reason_text(%{__exception__: true} = error), do: Exception.message(error)
  defp reason_text(reason), do: inspect(reason, limit: 10, printable_limit: @message_limit)

  defp to_int(nil), do: 0
  defp to_int(%Decimal{} = value), do: Decimal.to_integer(value)
  defp to_int(value) when is_integer(value), do: value

  defp repo, do: PhoenixKit.RepoHelper.repo()
end