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 `/` 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