Packages

S3-compatible object storage monitor (AWS S3, MinIO, Ceph RGW, Cloudflare R2, etc.), with its dashboard panel bundled in the same package as a separate module (Integrations.S3.Display) — one install, both halves; a release without raven_web simply runs the monitor headless.

Current section

Files

Jump to
raven_integration_s3 lib integrations s3.ex
Raw

lib/integrations/s3.ex

defmodule Integrations.S3 do
@moduledoc """
S3-compatible object storage monitor (AWS S3, MinIO, Ceph RGW, Cloudflare R2, etc.).
Checks storage health at three optional levels of depth:
1. **Liveness** — HTTP GET `/minio/health/live` (MinIO) or a basic HEAD on
the endpoint root. Always performed.
2. **Readiness / quorum** — HTTP GET `/minio/health/ready`. MinIO returns
503 when the cluster is degraded (insufficient drives/nodes for quorum).
Skipped on non-MinIO endpoints unless `:check_ready` is `true`.
3. **Bucket** — HEAD request to `/<bucket>` authenticated with AWS Signature
Version 4. Confirms the bucket exists and the credentials are valid.
Only performed when `:bucket`, `:access_key_id`, and `:secret_access_key`
are all set.
Collection only — see `Integrations.S3.Display` (same package) for the
dashboard panel. `Display.BundledDefault` auto-hooks it whenever this
monitor starts.
## Params
* `:endpoint` — Base URL of the storage endpoint. Required.
Examples: `http://minio.local:9000`,
`https://s3.amazonaws.com`.
* `:bucket` — Bucket name to verify (optional).
* `:access_key_id` — AWS / MinIO access key (optional; required for
bucket check and AWS S3).
* `:secret_access_key` — AWS / MinIO secret key (optional; required for
bucket check and AWS S3).
* `:region` — AWS region for SigV4 signing. Defaults to
`"us-east-1"`. Ignored for MinIO.
* `:check_ready` — Check the MinIO readiness endpoint. Defaults
to `true`.
* `:verify_tls` — Verify TLS certificates. Defaults to `true`.
* `:timeout_ms` — Request timeout. Defaults to `5000`.
## Health signal
* `:up` — Liveness OK; readiness OK (if checked); bucket accessible
(if configured).
* `:degraded` — Liveness OK but readiness check returned 503 (cluster
degraded / insufficient quorum).
* `:down` — Endpoint unreachable, or bucket HEAD returned 403/404.
"""
use CodeNameRaven.Monitor
@default_region "us-east-1"
@default_timeout_ms 5_000
@service "s3"
@impl true
def params_template do
%{
endpoint: "",
bucket: "",
access_key_id: "",
secret_access_key: "",
region: "us-east-1",
timeout_ms: "5000"
}
end
@impl true
def params_schema do
[
endpoint: [type: :string, required: true, doc: "Base URL of storage endpoint (e.g. http://minio:9000)"],
bucket: [type: :string, required: false, doc: "Bucket name to verify accessibility"],
access_key_id: [type: :string, required: false, doc: "AWS / MinIO access key"],
secret_access_key: [type: :string, required: false, doc: "AWS / MinIO secret key"],
region: [type: :string, default: "us-east-1", doc: "AWS region for SigV4 signing"],
check_ready: [type: :boolean, default: true, doc: "Check MinIO readiness/quorum endpoint"],
verify_tls: [type: :boolean, default: true, doc: "Verify TLS certificates"],
timeout_ms: [type: :non_neg_integer, default: 5_000, doc: "Request timeout in milliseconds"]
]
end
@impl true
def target_uri(params) do
case get_param(params, :endpoint) do
nil -> :none
ep ->
bucket = get_param(params, :bucket)
uri = if bucket, do: "#{ep}/#{bucket}", else: ep
{:ok, uri}
end
end
@impl true
def identity_params(params) do
%{
endpoint: get_param(params, :endpoint),
bucket: get_param(params, :bucket)
}
end
# ---------------------------------------------------------------------------
# Collect
# ---------------------------------------------------------------------------
@impl true
def collect(params, state) do
endpoint = get_param(params, :endpoint)
if is_nil(endpoint) do
{:error, "missing required param :endpoint", state}
else
timeout_ms = parse_int(get_param(params, :timeout_ms), @default_timeout_ms)
verify_tls = truthy?(get_param(params, :verify_tls), default: true)
check_ready = truthy?(get_param(params, :check_ready), default: true)
bucket = get_param(params, :bucket)
access_key = get_param(params, :access_key_id)
secret_key = get_param(params, :secret_access_key)
region = get_param(params, :region) || @default_region
req_opts = base_req_opts(verify_tls, timeout_ms)
t0 = System.monotonic_time(:millisecond)
with {:ok, live_status} <- check_liveness(endpoint, req_opts) do
latency_ms = System.monotonic_time(:millisecond) - t0
ready_status =
if check_ready,
do: check_readiness(endpoint, req_opts),
else: :skipped
bucket_status =
if bucket && access_key && secret_key do
check_bucket(endpoint, bucket, access_key, secret_key, region, req_opts)
else
:skipped
end
result = %{
liveness: live_status,
readiness: ready_status,
bucket_check: bucket_status,
latency_ms: latency_ms,
endpoint: endpoint,
bucket: bucket
}
{:ok, result, state}
else
{:error, reason} -> {:error, reason, state}
end
end
end
# ---------------------------------------------------------------------------
# Healthy?
# ---------------------------------------------------------------------------
@impl true
def healthy?(result) do
cond do
result.liveness == :down -> :down
result.bucket_check == :down -> :down
result.readiness == :degraded -> :degraded
true -> :up
end
end
# ---------------------------------------------------------------------------
# Metrics
# ---------------------------------------------------------------------------
@impl true
def metrics(result) do
base = %{
liveness: if(result.liveness == :up, do: 1, else: 0),
readiness: case result.readiness do
:up -> 1
:degraded -> 0
_ -> -1
end,
latency_ms: result.latency_ms
}
if result.bucket_check != :skipped do
Map.put(base, :bucket_accessible, if(result.bucket_check == :up, do: 1, else: 0))
else
base
end
end
# ---------------------------------------------------------------------------
# Health checks
# ---------------------------------------------------------------------------
defp check_liveness(endpoint, req_opts) do
# Try MinIO health endpoint; fall back to a HEAD on the root
url = "#{endpoint}/minio/health/live"
case Req.get(url, req_opts) do
{:ok, %{status: 200}} -> {:ok, :up}
{:ok, %{status: 404}} ->
# Not MinIO — try root HEAD
case Req.head(endpoint, req_opts) do
{:ok, %{status: s}} when s < 500 -> {:ok, :up}
{:ok, _} -> {:ok, :down}
{:error, reason} -> {:error, format_error(reason)}
end
{:ok, _} -> {:ok, :down}
{:error, reason} -> {:error, format_error(reason)}
end
end
defp check_readiness(endpoint, req_opts) do
url = "#{endpoint}/minio/health/ready"
case Req.get(url, req_opts) do
{:ok, %{status: 200}} -> :up
{:ok, %{status: 503}} -> :degraded
_ -> :skipped
end
end
defp check_bucket(endpoint, bucket, access_key, secret_key, region, req_opts) do
url = "#{endpoint}/#{bucket}"
now = DateTime.utc_now()
date = Calendar.strftime(now, "%Y%m%d")
amz_date = Calendar.strftime(now, "%Y%m%dT%H%M%SZ")
host = URI.parse(endpoint).host
headers = [
{"host", host},
{"x-amz-date", amz_date},
{"x-amz-content-sha256", "e3b0c44298fc1c149afbf4c8996fb92427ae41e4649b934ca495991b7852b855"}
]
auth_header = sigv4_auth_header(
"HEAD", "/#{bucket}", "", headers,
access_key, secret_key, region, date, amz_date
)
all_headers = headers ++ [{"authorization", auth_header}]
req_with_headers = Keyword.merge(req_opts, headers: all_headers)
case Req.head(url, req_with_headers) do
{:ok, %{status: 200}} -> :up
{:ok, %{status: 403}} -> :down
{:ok, %{status: 404}} -> :down
_ -> :skipped
end
end
# ---------------------------------------------------------------------------
# AWS Signature Version 4
# ---------------------------------------------------------------------------
defp sigv4_auth_header(method, path, query, headers, access_key, secret_key, region, date, amz_date) do
# Canonical headers (sorted, lowercase names)
canonical_headers =
headers
|> Enum.sort_by(fn {k, _} -> k end)
|> Enum.map_join("\n", fn {k, v} -> "#{String.downcase(k)}:#{String.trim(v)}" end)
signed_headers =
headers
|> Enum.map(fn {k, _} -> String.downcase(k) end)
|> Enum.sort()
|> Enum.join(";")
payload_hash = "e3b0c44298fc1c149afbf4c8996fb92427ae41e4649b934ca495991b7852b855"
canonical_request =
[method, path, query, canonical_headers <> "\n", signed_headers, payload_hash]
|> Enum.join("\n")
credential_scope = "#{date}/#{region}/#{@service}/aws4_request"
string_to_sign =
["AWS4-HMAC-SHA256", amz_date, credential_scope, sha256_hex(canonical_request)]
|> Enum.join("\n")
signing_key =
("AWS4" <> secret_key)
|> hmac_sha256(date)
|> hmac_sha256(region)
|> hmac_sha256(@service)
|> hmac_sha256("aws4_request")
signature = Base.encode16(hmac_sha256(signing_key, string_to_sign), case: :lower)
"AWS4-HMAC-SHA256 Credential=#{access_key}/#{credential_scope}, " <>
"SignedHeaders=#{signed_headers}, Signature=#{signature}"
end
defp sha256_hex(data),
do: Base.encode16(:crypto.hash(:sha256, data), case: :lower)
defp hmac_sha256(key, data) when is_binary(key),
do: :crypto.mac(:hmac, :sha256, key, data)
# ---------------------------------------------------------------------------
# Helpers
# ---------------------------------------------------------------------------
defp base_req_opts(verify_tls, timeout_ms) do
ssl_opts = if verify_tls,
do: [verify: :verify_peer, cacerts: :public_key.cacerts_get()],
else: [verify: :verify_none]
[
connect_options: [transport_opts: ssl_opts],
receive_timeout: timeout_ms,
retry: false
]
end
defp format_error(%{reason: reason}), do: inspect(reason)
defp format_error(reason), do: inspect(reason)
defp truthy?(nil, opts), do: Keyword.get(opts, :default, false)
defp truthy?("true", _), do: true
defp truthy?(true, _), do: true
defp truthy?("false", _), do: false
defp truthy?(false, _), do: false
defp truthy?(_, opts), do: Keyword.get(opts, :default, false)
defp get_param(params, key) when is_atom(key) do
v = params[key] || params[to_string(key)]
if is_binary(v) and String.trim(v) == "", do: nil, else: v
end
defp parse_int(nil, default), do: default
defp parse_int(v, _) when is_integer(v), do: v
defp parse_int(v, default) when is_binary(v) do
case Integer.parse(v) do
{n, _} -> n
:error -> default
end
end
defp parse_int(_, default), do: default
end