Packages

phoenix_kit

2.52.2
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