Changes: - Add Oban configuration to config/dev.exs (was only in runtime.exs) - Replace Exq telemetry metrics with Oban equivalents - Update publish_exq_stats to publish_oban_stats using Oban.Job queries - Track queue sizes, executing jobs, and available jobs This fixes the startup error where Oban config was not available in dev.
324 lines
9.3 KiB
Elixir
324 lines
9.3 KiB
Elixir
defmodule ToweropsWeb.Telemetry do
|
|
@moduledoc false
|
|
use Supervisor
|
|
|
|
import Telemetry.Metrics
|
|
|
|
def start_link(arg) do
|
|
Supervisor.start_link(__MODULE__, arg, name: __MODULE__)
|
|
end
|
|
|
|
@impl true
|
|
def init(_arg) do
|
|
# Attach telemetry handlers for logging request failures
|
|
_ =
|
|
:telemetry.attach(
|
|
"towerops-router-exception",
|
|
[:phoenix, :router_dispatch, :exception],
|
|
&__MODULE__.handle_router_exception/4,
|
|
nil
|
|
)
|
|
|
|
_ =
|
|
:telemetry.attach(
|
|
"towerops-endpoint-stop",
|
|
[:phoenix, :endpoint, :stop],
|
|
&__MODULE__.handle_endpoint_stop/4,
|
|
nil
|
|
)
|
|
|
|
children = [
|
|
# Telemetry poller will execute the given period measurements
|
|
# every 10_000ms. Learn more here: https://hexdocs.pm/telemetry_metrics
|
|
{:telemetry_poller, measurements: periodic_measurements(), period: 10_000}
|
|
# Add reporters as children of your supervision tree.
|
|
# {Telemetry.Metrics.ConsoleReporter, metrics: metrics()}
|
|
]
|
|
|
|
Supervisor.init(children, strategy: :one_for_one)
|
|
end
|
|
|
|
def metrics do
|
|
[
|
|
# Phoenix Metrics
|
|
summary("phoenix.endpoint.start.system_time",
|
|
unit: {:native, :millisecond}
|
|
),
|
|
summary("phoenix.endpoint.stop.duration",
|
|
unit: {:native, :millisecond}
|
|
),
|
|
summary("phoenix.router_dispatch.start.system_time",
|
|
tags: [:route],
|
|
unit: {:native, :millisecond}
|
|
),
|
|
summary("phoenix.router_dispatch.exception.duration",
|
|
tags: [:route],
|
|
unit: {:native, :millisecond}
|
|
),
|
|
summary("phoenix.router_dispatch.stop.duration",
|
|
tags: [:route],
|
|
unit: {:native, :millisecond}
|
|
),
|
|
summary("phoenix.socket_connected.duration",
|
|
unit: {:native, :millisecond}
|
|
),
|
|
sum("phoenix.socket_drain.count"),
|
|
summary("phoenix.channel_joined.duration",
|
|
unit: {:native, :millisecond}
|
|
),
|
|
summary("phoenix.channel_handled_in.duration",
|
|
tags: [:event],
|
|
unit: {:native, :millisecond}
|
|
),
|
|
|
|
# Database Metrics
|
|
summary("towerops.repo.query.total_time",
|
|
unit: {:native, :millisecond},
|
|
description: "The sum of the other measurements"
|
|
),
|
|
summary("towerops.repo.query.decode_time",
|
|
unit: {:native, :millisecond},
|
|
description: "The time spent decoding the data received from the database"
|
|
),
|
|
summary("towerops.repo.query.query_time",
|
|
unit: {:native, :millisecond},
|
|
description: "The time spent executing the query"
|
|
),
|
|
summary("towerops.repo.query.queue_time",
|
|
unit: {:native, :millisecond},
|
|
description: "The time spent waiting for a database connection"
|
|
),
|
|
summary("towerops.repo.query.idle_time",
|
|
unit: {:native, :millisecond},
|
|
description: "The time the connection spent waiting before being checked out for the query"
|
|
),
|
|
|
|
# VM Metrics
|
|
summary("vm.memory.total", unit: {:byte, :kilobyte}),
|
|
summary("vm.total_run_queue_lengths.total"),
|
|
summary("vm.total_run_queue_lengths.cpu"),
|
|
summary("vm.total_run_queue_lengths.io"),
|
|
|
|
# Oban/Redis Metrics
|
|
last_value("towerops.oban.queue.size",
|
|
tags: [:queue],
|
|
description: "Number of jobs in each Oban queue"
|
|
),
|
|
last_value("towerops.oban.jobs.executing",
|
|
description: "Number of currently executing Oban jobs"
|
|
),
|
|
last_value("towerops.oban.jobs.available",
|
|
description: "Number of available Oban jobs"
|
|
),
|
|
last_value("towerops.redis.connected_clients",
|
|
description: "Number of Redis/Valkey connected clients"
|
|
),
|
|
last_value("towerops.redis.used_memory",
|
|
unit: {:byte, :megabyte},
|
|
description: "Redis/Valkey memory usage"
|
|
),
|
|
last_value("towerops.redis.commands_processed",
|
|
description: "Total commands processed by Redis/Valkey"
|
|
)
|
|
]
|
|
end
|
|
|
|
defp periodic_measurements do
|
|
[
|
|
# Measure Oban and Redis stats every 10 seconds
|
|
{__MODULE__, :publish_oban_stats, []},
|
|
{__MODULE__, :publish_redis_stats, []}
|
|
]
|
|
end
|
|
|
|
@doc """
|
|
Publishes Oban queue and job statistics.
|
|
"""
|
|
def publish_oban_stats do
|
|
if Application.get_env(:towerops, :env) == :test do
|
|
:ok
|
|
else
|
|
try do
|
|
import Ecto.Query
|
|
|
|
queues = ["default", "discovery", "pollers", "monitors", "maintenance"]
|
|
|
|
# Query queue sizes
|
|
for queue <- queues do
|
|
try do
|
|
size =
|
|
Towerops.Repo.one(
|
|
from j in Oban.Job,
|
|
where: j.queue == ^queue and j.state in ["available", "scheduled"],
|
|
select: count(j.id)
|
|
)
|
|
|
|
:telemetry.execute(
|
|
[:towerops, :oban, :queue, :size],
|
|
%{value: size},
|
|
%{queue: queue}
|
|
)
|
|
rescue
|
|
_ -> :ok
|
|
end
|
|
end
|
|
|
|
# Query executing jobs count
|
|
try do
|
|
executing_count =
|
|
Towerops.Repo.one(
|
|
from j in Oban.Job,
|
|
where: j.state == "executing",
|
|
select: count(j.id)
|
|
)
|
|
|
|
:telemetry.execute(
|
|
[:towerops, :oban, :jobs, :executing],
|
|
%{value: executing_count},
|
|
%{}
|
|
)
|
|
rescue
|
|
_ -> :ok
|
|
end
|
|
|
|
# Query available jobs count
|
|
try do
|
|
available_count =
|
|
Towerops.Repo.one(
|
|
from j in Oban.Job,
|
|
where: j.state == "available",
|
|
select: count(j.id)
|
|
)
|
|
|
|
:telemetry.execute(
|
|
[:towerops, :oban, :jobs, :available],
|
|
%{value: available_count},
|
|
%{}
|
|
)
|
|
rescue
|
|
_ -> :ok
|
|
end
|
|
|
|
:ok
|
|
rescue
|
|
_ -> :ok
|
|
catch
|
|
_, _ -> :ok
|
|
end
|
|
end
|
|
end
|
|
|
|
@doc """
|
|
Publishes Redis/Valkey statistics.
|
|
"""
|
|
def publish_redis_stats do
|
|
# Only run if Redis is configured (not in test env)
|
|
if Application.get_env(:towerops, :env) == :test do
|
|
:ok
|
|
|
|
# Connect to Redis and get INFO
|
|
|
|
# Parse INFO response
|
|
|
|
# Publish metrics
|
|
else
|
|
try do
|
|
redis_config = Application.get_env(:towerops, :redis, [])
|
|
host = Keyword.get(redis_config, :host, "localhost")
|
|
port = Keyword.get(redis_config, :port, 6379)
|
|
|
|
with {:ok, conn} <- Redix.start_link(host: host, port: port),
|
|
{:ok, info} <- Redix.command(conn, ["INFO", "stats"]),
|
|
{:ok, memory_info} <- Redix.command(conn, ["INFO", "memory"]),
|
|
{:ok, clients_info} <- Redix.command(conn, ["INFO", "clients"]) do
|
|
stats = parse_redis_info(info)
|
|
memory_stats = parse_redis_info(memory_info)
|
|
client_stats = parse_redis_info(clients_info)
|
|
|
|
if total_commands = stats["total_commands_processed"] do
|
|
:telemetry.execute(
|
|
[:towerops, :redis, :commands_processed],
|
|
%{value: String.to_integer(total_commands)},
|
|
%{}
|
|
)
|
|
end
|
|
|
|
if used_memory = memory_stats["used_memory"] do
|
|
:telemetry.execute(
|
|
[:towerops, :redis, :used_memory],
|
|
%{value: String.to_integer(used_memory)},
|
|
%{}
|
|
)
|
|
end
|
|
|
|
if connected_clients = client_stats["connected_clients"] do
|
|
:telemetry.execute(
|
|
[:towerops, :redis, :connected_clients],
|
|
%{value: String.to_integer(connected_clients)},
|
|
%{}
|
|
)
|
|
end
|
|
|
|
Redix.stop(conn)
|
|
:ok
|
|
end
|
|
rescue
|
|
_ -> :ok
|
|
end
|
|
end
|
|
end
|
|
|
|
@doc """
|
|
Parse Redis INFO command output into a map.
|
|
"""
|
|
def parse_redis_info(info_string) do
|
|
info_string
|
|
|> String.split(~r/\r?\n/)
|
|
|> Enum.map(&String.trim/1)
|
|
|> Enum.reject(&(String.starts_with?(&1, "#") or &1 == ""))
|
|
|> Enum.map(&String.split(&1, ":", parts: 2))
|
|
|> Enum.filter(&(length(&1) == 2))
|
|
|> Map.new(fn [key, value] -> {key, value} end)
|
|
end
|
|
|
|
# Telemetry handler for router exceptions
|
|
def handle_router_exception(_event, _measurements, metadata, _config) do
|
|
require Logger
|
|
|
|
Logger.error(
|
|
"Router exception on #{metadata.plug} #{metadata.conn.method} #{metadata.conn.request_path}",
|
|
kind: metadata.kind,
|
|
reason: metadata.reason,
|
|
stacktrace: metadata.stacktrace,
|
|
request_id: metadata.conn.assigns[:request_id]
|
|
)
|
|
end
|
|
|
|
# Telemetry handler for endpoint stop events (log slow requests and errors)
|
|
def handle_endpoint_stop(_event, measurements, metadata, _config) do
|
|
require Logger
|
|
|
|
duration_ms = System.convert_time_unit(measurements.duration, :native, :millisecond)
|
|
|
|
# Log slow requests (over 5 seconds)
|
|
if duration_ms > 5_000 do
|
|
Logger.warning(
|
|
"Slow request: #{metadata.conn.method} #{metadata.conn.request_path} took #{duration_ms}ms",
|
|
request_id: metadata.conn.assigns[:request_id],
|
|
duration_ms: duration_ms
|
|
)
|
|
end
|
|
|
|
# Log requests with non-2xx status codes
|
|
status = metadata.conn.status
|
|
|
|
if status >= 500 do
|
|
Logger.error(
|
|
"Server error: #{metadata.conn.method} #{metadata.conn.request_path} returned #{status}",
|
|
request_id: metadata.conn.assigns[:request_id],
|
|
status: status,
|
|
duration_ms: duration_ms
|
|
)
|
|
end
|
|
end
|
|
end
|