defmodule Integrations.Postgres do @moduledoc """ Postgres monitor — checks connectivity, role, connections, cache efficiency, transaction throughput, and replication health. Collection only — see `Integrations.Postgres.Display` (same package) for the dashboard panel. `Display.BundledDefault` auto-hooks it whenever this monitor starts, same as a single-module package would; a release without `raven_web` simply never compiles the display half and runs this monitor headless. Implements the Postgres frontend/backend protocol (v3) over raw TCP — no external library required, matching `maria_db`/`mongo_db`/`redis`/`kafka`'s own hand-rolled wire-protocol approach rather than depending on Postgrex. Supports trust, cleartext, MD5, and SCRAM-SHA-256 (the modern default since PG10) authentication, using only OTP's `:crypto` module (requires OTP 24+ for `:crypto.pbkdf2_hmac/5`, same requirement `mongo_db`'s SCRAM-SHA-256 implementation already has). Works against any Postgres instance. HA-aware: when connected to a primary it shows replica count; when connected to a replica it shows replication lag. Rate metrics (e.g. `xact_commit_rate`) require two consecutive successful checks to compute — they are absent from the first sample and from any sample following a counter reset (Postgres restart or `pg_stat_reset()`). ## Params * `:host` — Hostname or k8s service DNS. Required. * `:port` — Port. Defaults to `5432`. * `:database` — Database name. Defaults to `"postgres"`. * `:username` — Username. Required. * `:password` — Password. Required. * `:timeout_ms` — Connect + query timeout in ms. Defaults to `5000`. ## Health signal * `:up` — Connected and all thresholds within bounds. * `:degraded` — Any of: replica lag > 10s, idle-in-transaction > 5, any lock waiters, connection utilization > 85%, any deadlocks in the last interval, XID age > 1.5B, buffer hit ratio < 95%. * `:down` — Connection failed or query errored. ## Example raven.toml config [[monitors]] id = "postgres-primary" name = "Postgres Primary" module = "Integrations.Postgres" interval = 30 [monitors.params] host = "postgres.example.com" database = "app" username = "app" password = "secret" """ use CodeNameRaven.Monitor @default_port 5432 @default_database "postgres" @default_timeout_ms 5_000 @impl true def params_template do %{ host: "postgres.example.com", port: "5432", database: "postgres", username: "raven_monitor", password: "" } end @impl true def params_schema do [ host: [type: :string, required: true, doc: "PostgreSQL server hostname or IP"], port: [type: :non_neg_integer, default: 5432, doc: "Port number"], database: [type: :string, default: "postgres", doc: "Database name"], username: [type: :string, required: true, doc: "Login username"], password: [type: :string, required: false, doc: "Login password"], timeout_ms: [type: :non_neg_integer, default: 5_000, doc: "Connection timeout in milliseconds"], assertions: [type: {:list, :any}, default: [], doc: "List of assertion maps for query-level checks"] ] end # --------------------------------------------------------------------------- @impl true def target_uri(params) do host = to_string(params[:host] || params["host"] || "") port = parse_int(params[:port] || params["port"], @default_port) db = to_string(params[:database] || params["database"] || @default_database) if host == "", do: :none, else: {:ok, "postgres://#{host}:#{port}/#{db}"} end # Collect # --------------------------------------------------------------------------- @impl true def collect(params, state) do host = to_string(params[:host] || params["host"] || "") port = parse_int(params[:port] || params["port"], @default_port) database = to_string(params[:database] || params["database"] || @default_database) username = to_string(params[:username] || params["username"] || "") password = to_string(params[:password] || params["password"] || "") timeout_ms = parse_int(params[:timeout_ms] || params["timeout_ms"], @default_timeout_ms) if host == "" do {:error, "missing required param :host", state} else case do_collect(host, port, database, username, password, timeout_ms, state) do {:ok, raw_result, new_state} -> alias CodeNameRaven.Monitor.Assertion assertions = Assertion.parse(params[:assertions] || []) {assertion_status, failures} = Assertion.evaluate_all(assertions, raw_result) result = Map.merge(raw_result, %{ assertions: assertions, assertion_status: assertion_status, assertion_failures: failures }) {:ok, result, new_state} other -> other end end end defp do_collect(host, port, database, username, password, timeout_ms, state) do started_at = System.monotonic_time(:millisecond) tcp_opts = [:binary, active: false, packet: :raw, send_timeout: timeout_ms] case :gen_tcp.connect(to_charlist(host), port, tcp_opts, timeout_ms) do {:ok, socket} -> result = try do case handshake(socket, database, username, password, timeout_ms) do :ok -> run_checks(socket, timeout_ms, started_at, state) {:error, reason} -> {:error, reason} end after :gen_tcp.close(socket) end case result do {:ok, data, new_state} -> {:ok, data, new_state} {:error, reason} -> {:error, reason, state} end {:error, :econnrefused} -> {:error, "connection refused on port #{port}", state} {:error, :timeout} -> {:error, "connection timed out after #{timeout_ms}ms", state} {:error, :nxdomain} -> {:error, "hostname not found: #{host}", state} {:error, reason} -> {:error, "connection failed: #{inspect(reason)}", state} end end defp run_checks(socket, timeout_ms, started_at, state) do with {:ok, activity} <- query_activity(socket, timeout_ms), {:ok, db_raw} <- query_db_stat(socket, timeout_ms), {:ok, bgw_raw} <- query_bgwriter(socket, timeout_ms) do latency_ms = System.monotonic_time(:millisecond) - started_at now = DateTime.utc_now() elapsed = elapsed_sec(state[:prev_collected_at], now) db_deltas = compute_deltas(db_raw, state[:prev_db], elapsed) bgw_deltas = compute_deltas(bgw_raw, state[:prev_bgwriter], elapsed) conn_util = if activity.max_connections > 0, do: (activity.active_connections + activity.idle_connections) / activity.max_connections, else: 0.0 result = %{ latency_ms: latency_ms, is_replica: activity.is_replica, role_metrics: %{ active_connections: activity.active_connections, idle_connections: activity.idle_connections, idle_in_transaction: activity.idle_in_transaction, waiting_on_lock: activity.waiting_on_lock, max_connections: activity.max_connections, longest_query_seconds: activity.longest_query_seconds, connection_utilization: conn_util, replica_count: activity.replica_count, lag_ms: activity.lag_ms }, db_metrics: %{ xact_commit_rate: db_deltas[:xact_commit_rate], xact_rollback_rate: db_deltas[:xact_rollback_rate], tup_inserted_rate: db_deltas[:tup_inserted_rate], tup_updated_rate: db_deltas[:tup_updated_rate], tup_deleted_rate: db_deltas[:tup_deleted_rate], tup_fetched_rate: db_deltas[:tup_fetched_rate], blks_hit_rate: db_deltas[:blks_hit_rate], blks_read_rate: db_deltas[:blks_read_rate], buffer_hit_ratio: buffer_hit_ratio(db_deltas), temp_files_rate: db_deltas[:temp_files_rate], temp_bytes_rate: db_deltas[:temp_bytes_rate], deadlocks_rate: db_deltas[:deadlocks_rate], conflicts_rate: db_deltas[:conflicts_rate], db_size_bytes: db_raw.db_size_bytes, xid_wraparound_age: db_raw.xid_wraparound_age }, bgwriter_metrics: %{ checkpoints_timed_rate: bgw_deltas[:checkpoints_timed_rate], checkpoints_req_rate: bgw_deltas[:checkpoints_req_rate], buffers_checkpoint_rate: bgw_deltas[:buffers_checkpoint_rate], buffers_clean_rate: bgw_deltas[:buffers_clean_rate], buffers_backend_rate: bgw_deltas[:buffers_backend_rate], maxwritten_clean_rate: bgw_deltas[:maxwritten_clean_rate] } } new_state = %{ prev_db: db_raw, prev_bgwriter: bgw_raw, prev_collected_at: now } {:ok, result, new_state} end end defp query_activity(socket, timeout_ms) do sql = """ SELECT pg_is_in_recovery() AS is_replica, (SELECT count(*) FROM pg_stat_replication)::integer AS replica_count, CASE WHEN pg_is_in_recovery() THEN (EXTRACT(EPOCH FROM (now() - pg_last_xact_replay_timestamp())) * 1000)::bigint ELSE 0 END AS lag_ms, count(*) FILTER (WHERE state = 'active')::integer AS active_connections, count(*) FILTER (WHERE state = 'idle')::integer AS idle_connections, count(*) FILTER (WHERE state = 'idle in transaction')::integer AS idle_in_transaction, count(*) FILTER (WHERE wait_event_type = 'Lock')::integer AS waiting_on_lock, current_setting('max_connections')::integer AS max_connections, COALESCE(max(EXTRACT(EPOCH FROM (now() - query_start))), 0)::double precision AS longest_query_seconds FROM pg_stat_activity """ case query(socket, sql, timeout_ms) do {:ok, result} -> {:ok, rows_to_map(result)} {:error, reason} -> {:error, "activity query failed: #{reason}"} end end defp query_db_stat(socket, timeout_ms) do # datfrozenxid lives in pg_database, not pg_stat_database — join to get it. sql = """ SELECT s.xact_commit, s.xact_rollback, s.tup_inserted, s.tup_updated, s.tup_deleted, s.tup_fetched, s.tup_returned, s.blks_hit, s.blks_read, s.temp_files, s.temp_bytes, s.deadlocks, s.conflicts, pg_database_size(current_database()) AS db_size_bytes, age(d.datfrozenxid)::integer AS xid_wraparound_age FROM pg_stat_database s JOIN pg_database d ON d.datname = s.datname WHERE s.datname = current_database() """ case query(socket, sql, timeout_ms) do {:ok, result} -> {:ok, rows_to_map(result)} {:error, reason} -> {:error, "db stat query failed: #{reason}"} end end defp query_bgwriter(socket, timeout_ms) do # PG 17 moved checkpoint stats from pg_stat_bgwriter into pg_stat_checkpointer # (with renamed columns). Try the PG17 form first; fall back to legacy PG ≤ 16. # buffers_backend moved to pg_stat_io in PG17 and is omitted here — its rate # will be nil and filtered from the metrics map. pg17_sql = """ SELECT n.num_timed AS checkpoints_timed, n.num_requested AS checkpoints_req, n.buffers_written AS buffers_checkpoint, b.buffers_clean, b.maxwritten_clean FROM pg_stat_checkpointer n, pg_stat_bgwriter b """ legacy_sql = """ SELECT checkpoints_timed, checkpoints_req, buffers_checkpoint, buffers_clean, buffers_backend, maxwritten_clean FROM pg_stat_bgwriter """ case query(socket, pg17_sql, timeout_ms) do {:ok, result} -> {:ok, rows_to_map(result)} {:error, _} -> case query(socket, legacy_sql, timeout_ms) do {:ok, result} -> {:ok, rows_to_map(result)} {:error, reason} -> {:error, "bgwriter query failed: #{reason}"} end end end defp rows_to_map(%{columns: _cols, rows: []}), do: %{} defp rows_to_map(%{columns: cols, rows: [row | _]}) do Enum.zip(cols, row) |> Map.new(fn {col, val} -> {String.to_atom(col), val} end) end defp elapsed_sec(nil, _now), do: 0.0 defp elapsed_sec(prev_at, now), do: DateTime.diff(now, prev_at, :millisecond) / 1000.0 defp compute_deltas(_current, nil, _elapsed), do: %{} defp compute_deltas(_current, _previous, elapsed) when elapsed <= 0, do: %{} defp compute_deltas(current, previous, elapsed) do Map.new(current, fn {key, value} -> prev = Map.get(previous, key, 0) rate_key = :"#{key}_rate" if is_number(value) and is_number(prev) and value >= prev do {rate_key, (value - prev) / elapsed} else # Counter reset (value < prev) or non-numeric — no rate this tick. {rate_key, nil} end end) end defp buffer_hit_ratio(%{blks_hit_rate: hit, blks_read_rate: read}) when is_number(hit) and is_number(read) and hit + read > 0 do hit / (hit + read) end defp buffer_hit_ratio(_), do: nil # --------------------------------------------------------------------------- # Postgres frontend/backend protocol (v3) — connection setup # --------------------------------------------------------------------------- defp handshake(socket, database, username, password, timeout) do with :ok <- send_startup(socket, database, username), :ok <- authenticate(socket, username, password, timeout), :ok <- await_ready(socket, timeout) do :ok end end defp send_startup(socket, database, username) do params = "user" <> <<0>> <> username <> <<0>> <> "database" <> <<0>> <> database <> <<0>> <> <<0>> body = <<196_608::32, params::binary>> send_raw(socket, <>, "startup") end defp authenticate(socket, username, password, timeout) do case recv_message(socket, timeout) do {:ok, ?R, <<0::32>>} -> :ok {:ok, ?R, <<3::32>>} -> with :ok <- send_password_message(socket, password <> <<0>>) do expect_auth_ok(socket, timeout) end {:ok, ?R, <<5::32, salt::binary-size(4)>>} -> with :ok <- send_password_message(socket, md5_password_hash(username, password, salt) <> <<0>>) do expect_auth_ok(socket, timeout) end {:ok, ?R, <<10::32, mechanisms::binary>>} -> if scram_sha_256_offered?(mechanisms) do scram_authenticate(socket, username, password, timeout) else {:error, "server requires unsupported SASL mechanism(s): #{inspect(mechanisms)}"} end {:ok, ?E, payload} -> {:error, parse_error_message(payload)} {:ok, _type, _payload} -> {:error, "unexpected message during authentication"} {:error, reason} -> {:error, reason} end end defp expect_auth_ok(socket, timeout) do case recv_message(socket, timeout) do {:ok, ?R, <<0::32>>} -> :ok {:ok, ?E, payload} -> {:error, parse_error_message(payload)} {:ok, _, _} -> {:error, "authentication failed (unexpected response)"} {:error, reason} -> {:error, reason} end end defp md5_password_hash(username, password, salt) do inner = Base.encode16(:crypto.hash(:md5, password <> username), case: :lower) outer = Base.encode16(:crypto.hash(:md5, inner <> salt), case: :lower) "md5" <> outer end defp await_ready(socket, timeout) do case recv_message(socket, timeout) do {:ok, ?Z, _payload} -> :ok {:ok, ?E, payload} -> {:error, parse_error_message(payload)} {:ok, _type, _payload} -> await_ready(socket, timeout) {:error, reason} -> {:error, reason} end end # --------------------------------------------------------------------------- # SCRAM-SHA-256 authentication — same math as mongo_db's, different message # framing (Postgres wraps every response, initial or continuation, in a # PasswordMessage — 'p' — rather than a command document) # --------------------------------------------------------------------------- defp scram_sha_256_offered?(mechanisms) do mechanisms |> String.split(<<0>>) |> Enum.member?("SCRAM-SHA-256") end defp scram_authenticate(socket, username, password, timeout) do cnonce = :base64.encode(:crypto.strong_rand_bytes(18)) client_first_bare = "n=#{sasl_escape(username)},r=#{cnonce}" client_first = "n,," <> client_first_bare initial_payload = "SCRAM-SHA-256" <> <<0>> <> <> <> client_first with :ok <- send_password_message(socket, initial_payload), {:ok, ?R, <<11::32, server_first::binary>>} <- recv_message(socket, timeout), {:ok, snonce, salt, iters} <- parse_scram_server_first(server_first), true <- String.starts_with?(snonce, cnonce) do finish_scram(socket, timeout, username, password, client_first_bare, server_first, snonce, salt, iters) else false -> {:error, "SCRAM: server nonce does not start with client nonce"} {:ok, ?E, payload} -> {:error, parse_error_message(payload)} {:error, reason} -> {:error, reason} _ -> {:error, "SCRAM: unexpected server-first response"} end end defp finish_scram(socket, timeout, _username, password, client_first_bare, server_first, snonce, salt, iters) do salted_password = :crypto.pbkdf2_hmac(:sha256, password, salt, iters, 32) client_key = :crypto.mac(:hmac, :sha256, salted_password, "Client Key") stored_key = :crypto.hash(:sha256, client_key) client_final_bare = "c=biws,r=#{snonce}" auth_message = client_first_bare <> "," <> server_first <> "," <> client_final_bare client_signature = :crypto.mac(:hmac, :sha256, stored_key, auth_message) client_proof = :crypto.exor(client_key, client_signature) client_final = client_final_bare <> ",p=" <> :base64.encode(client_proof) with :ok <- send_password_message(socket, client_final), {:ok, ?R, <<12::32, _server_final::binary>>} <- recv_message(socket, timeout) do expect_auth_ok(socket, timeout) else {:ok, ?E, payload} -> {:error, parse_error_message(payload)} {:error, reason} -> {:error, reason} _ -> {:error, "SCRAM: unexpected server-final response"} end end defp parse_scram_server_first(data) do parts = data |> String.split(",") |> Map.new(fn part -> case String.split(part, "=", parts: 2) do [k, v] -> {k, v} _ -> {"", ""} end end) with r when is_binary(r) <- parts["r"], s when is_binary(s) <- parts["s"], i when is_binary(i) <- parts["i"], {iters, _} <- Integer.parse(i) do {:ok, r, :base64.decode(s), iters} else _ -> {:error, "SCRAM: could not parse server-first-message"} end end defp sasl_escape(user) do user |> String.replace("=", "=3D") |> String.replace(",", "=2C") end # --------------------------------------------------------------------------- # Simple query protocol # --------------------------------------------------------------------------- defp query(socket, sql, timeout) do with :ok <- send_query(socket, sql) do recv_query_response(socket, timeout, nil, []) end end defp send_query(socket, sql) do body = sql <> <<0>> send_raw(socket, <>, "query") end defp recv_query_response(socket, timeout, cols, rows) do case recv_message(socket, timeout) do {:ok, ?T, payload} -> recv_query_response(socket, timeout, parse_row_description(payload), rows) {:ok, ?D, payload} -> recv_query_response(socket, timeout, cols, [parse_data_row(payload) | rows]) {:ok, ?C, _payload} -> recv_query_response(socket, timeout, cols, rows) {:ok, ?Z, _payload} -> {:ok, %{columns: cols || [], rows: Enum.reverse(rows)}} {:ok, ?E, payload} -> message = parse_error_message(payload) drain_to_ready(socket, timeout) {:error, message} {:ok, _other_type, _payload} -> recv_query_response(socket, timeout, cols, rows) {:error, reason} -> {:error, reason} end end defp drain_to_ready(socket, timeout) do case recv_message(socket, timeout) do {:ok, ?Z, _} -> :ok {:ok, _, _} -> drain_to_ready(socket, timeout) {:error, reason} -> {:error, reason} end end defp parse_row_description(<<_num_fields::16, rest::binary>>), do: parse_fields(rest, []) defp parse_fields(<<>>, acc), do: Enum.reverse(acc) defp parse_fields(data, acc) do [name, rest] = :binary.split(data, <<0>>) <<_table_oid::32, _col_num::16, _type_oid::32, _type_size::16, _type_mod::32, _fmt::16, rest2::binary>> = rest parse_fields(rest2, [name | acc]) end defp parse_data_row(<>), do: parse_columns(rest, num_cols, []) defp parse_columns(_data, 0, acc), do: Enum.reverse(acc) defp parse_columns(<<-1::signed-32, rest::binary>>, n, acc), do: parse_columns(rest, n - 1, [nil | acc]) defp parse_columns(<>, n, acc) do <> = rest parse_columns(rest2, n - 1, [cast_value(val) | acc]) end # Simple query protocol returns everything as text — cast heuristically # rather than tracking type OIDs from RowDescription, since every column # this monitor queries is already explicitly cast to a known SQL type in # the query itself (boolean, integer/bigint, or double precision). defp cast_value("t"), do: true defp cast_value("f"), do: false defp cast_value(bin) do cond do match?({_, ""}, Integer.parse(bin)) -> elem(Integer.parse(bin), 0) match?({_, ""}, Float.parse(bin)) -> elem(Float.parse(bin), 0) true -> bin end end # --------------------------------------------------------------------------- # Message framing # --------------------------------------------------------------------------- defp send_password_message(socket, payload) do send_raw(socket, <>, "password message") end defp send_raw(socket, msg, label) do case :gen_tcp.send(socket, msg) do :ok -> :ok {:error, reason} -> {:error, "send #{label} failed: #{inspect(reason)}"} end end defp recv_message(socket, timeout) do case :gen_tcp.recv(socket, 5, timeout) do {:ok, <>} -> payload_len = len - 4 payload_result = if payload_len > 0, do: :gen_tcp.recv(socket, payload_len, timeout), else: {:ok, <<>>} case payload_result do {:ok, payload} -> {:ok, type, payload} {:error, reason} -> {:error, "recv payload failed: #{inspect(reason)}"} end {:error, reason} -> {:error, "recv header failed: #{inspect(reason)}"} end end defp parse_error_message(payload) do payload |> parse_error_fields() |> Map.get("M", "unknown Postgres error") end defp parse_error_fields(<<0, _rest::binary>>), do: %{} defp parse_error_fields(<>) do case :binary.split(rest, <<0>>) do [value, rest2] -> Map.put(parse_error_fields(rest2), <>, value) _ -> %{} end end defp parse_error_fields(<<>>), do: %{} # --------------------------------------------------------------------------- # Health # --------------------------------------------------------------------------- @impl true def healthy?(%{assertions: [_ | _], assertion_status: status}), do: status # Replica lag def healthy?(%{role_metrics: %{lag_ms: lag}}) when is_number(lag) and lag > 10_000, do: :degraded # Idle-in-transaction sessions are holding locks and blocking vacuums def healthy?(%{role_metrics: %{idle_in_transaction: iit}}) when is_number(iit) and iit > 5, do: :degraded # Any session waiting on a lock — a spike here signals contention def healthy?(%{role_metrics: %{waiting_on_lock: wl}}) when is_number(wl) and wl > 0, do: :degraded # Connection pool saturation def healthy?(%{role_metrics: %{connection_utilization: cu}}) when is_number(cu) and cu > 0.85, do: :degraded # Any deadlock in the last interval is a signal worth acting on def healthy?(%{db_metrics: %{deadlocks_rate: dl}}) when is_number(dl) and dl > 0, do: :degraded # XID wraparound: Postgres hard limit is ~2.1B; alert at 1.5B def healthy?(%{db_metrics: %{xid_wraparound_age: age}}) when is_number(age) and age > 1_500_000_000, do: :degraded # Buffer cache hit ratio < 95% — conservative; see future work for configurable thresholds def healthy?(%{db_metrics: %{buffer_hit_ratio: bhr}}) when is_number(bhr) and bhr < 0.95, do: :degraded def healthy?(_result), do: :up # --------------------------------------------------------------------------- # Metrics # --------------------------------------------------------------------------- @impl true def metrics(%{latency_ms: lat} = result) do role = result.role_metrics db = result.db_metrics bg = result.bgwriter_metrics %{ latency_ms: lat, active_connections: role.active_connections, idle_in_transaction: role.idle_in_transaction, waiting_on_lock: role.waiting_on_lock, connection_utilization: role.connection_utilization, longest_query_sec: role.longest_query_seconds, lag_ms: role.lag_ms, xact_commit_rate: db.xact_commit_rate, xact_rollback_rate: db.xact_rollback_rate, tup_inserted_rate: db.tup_inserted_rate, tup_updated_rate: db.tup_updated_rate, tup_deleted_rate: db.tup_deleted_rate, tup_fetched_rate: db.tup_fetched_rate, blks_read_rate: db.blks_read_rate, buffer_hit_ratio: db.buffer_hit_ratio, temp_bytes_rate: db.temp_bytes_rate, deadlocks_rate: db.deadlocks_rate, db_size_bytes: db.db_size_bytes, xid_wraparound_age: db.xid_wraparound_age, checkpoints_req_rate: bg.checkpoints_req_rate, buffers_backend_rate: bg.buffers_backend_rate } |> Enum.reject(fn {_k, v} -> is_nil(v) end) |> Map.new() end # --------------------------------------------------------------------------- # Helpers # --------------------------------------------------------------------------- defp parse_int(nil, default), do: default defp parse_int(v, _default) 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