Packages

HTTP telemetry dashboard for Elixir/Phoenix — monitor outbound & inbound requests with LiveView

Current section

Files

Jump to
monitorex lib monitorex event_handler.ex
Raw

lib/monitorex/event_handler.ex

defmodule Monitorex.EventHandler do
@moduledoc """
Handles telemetry events from Tesla, Finch, Req, and Phoenix, transforming them
into `Monitorex.Event` structs.
Each handler function follows the `Telemetry` handler callback signature:
`(event_name, measurements, metadata, config)`.
"""
alias Monitorex.ConsumerIdentifier
alias Monitorex.Event
alias Monitorex.HeaderRedactor
alias Monitorex.UrlNormalizer
alias Monitorex.URLRedactor
@doc """
Handles a Tesla telemetry event (`[:tesla, :request, :stop]`).
Parses the metadata and measurements into a `Monitorex.Event` struct with
source `:tesla` and direction `:outbound`.
The URL is normalised via `UrlNormalizer.normalize/1` and sensitive query
parameters are redacted via `URLRedactor.redact/1`.
## Telemetry metadata shape
%{
url: %URI{},
method: :get | :post | …,
status: 200,
pid: pid,
monotonic_time: integer
}
The `dedup_key` is set to `{pid, monotonic_time}` for deduplication.
"""
@spec handle_tesla_event(
event_name :: [atom()],
measurements :: map(),
metadata :: map(),
config :: keyword()
) :: Event.t() | nil
def handle_tesla_event([:tesla, :request, :stop], measurements, metadata, _config) do
# Support both legacy Tesla telemetry (flat keys) and modern Tesla (env struct)
{url_str, method, status, req_headers, resp_headers, ts, pid} =
case metadata do
%{env: %{url: url, method: m, status: s} = env} ->
url_str = url_to_string(url)
{url_str, Event.normalize_method(m), s, redact_headers_from_metadata(env.headers || []),
redact_headers_from_metadata(env.headers || []),
metadata[:monotonic_time] || measurements[:monotonic_time] || System.monotonic_time(),
metadata[:pid] || self()}
%{url: url, method: m, status: s} ->
url_str = url_to_string(url)
{url_str, Event.normalize_method(m), s,
redact_headers_from_metadata(metadata[:req_headers]),
redact_headers_from_metadata(metadata[:resp_headers]),
metadata[:monotonic_time] || measurements[:monotonic_time] || System.monotonic_time(),
metadata[:pid] || self()}
end
normalized_url = UrlNormalizer.normalize(url_str)
redacted_url = URLRedactor.redact(normalized_url)
%URI{host: host, path: path} = URI.parse(normalized_url)
tag_slow_request(
%Event{
source: :tesla,
direction: :outbound,
method: method,
host: host,
path: path,
full_url: redacted_url,
status: status,
status_class: Event.classify_status(status || 0),
duration_ms: Event.duration_ms(measurements.duration),
timestamp: System.system_time(:microsecond),
dedup_key: {pid, ts},
request_headers: req_headers,
response_headers: resp_headers,
request_body: maybe_store_body(metadata[:request_body], :request),
response_body: maybe_store_body(metadata[:response_body], :response)
},
metadata
)
end
def handle_tesla_event([:tesla, :request, :exception], measurements, metadata, _config) do
{url_str, method, req_headers, ts, pid} =
case metadata do
%{env: %{url: url, method: m} = env} ->
{url_to_string(url), Event.normalize_method(m),
redact_headers_from_metadata(env.headers || []),
metadata[:monotonic_time] || measurements[:monotonic_time] || System.monotonic_time(),
metadata[:pid] || self()}
%{url: url, method: m} ->
{url_to_string(url), Event.normalize_method(m),
redact_headers_from_metadata(metadata[:req_headers]),
metadata[:monotonic_time] || measurements[:monotonic_time] || System.monotonic_time(),
metadata[:pid] || self()}
_ ->
{"unknown", "UNKNOWN", [], System.monotonic_time(), self()}
end
normalized_url = if url_str != "unknown", do: UrlNormalizer.normalize(url_str), else: url_str
redacted_url = URLRedactor.redact(normalized_url)
tag_slow_request(
%Event{
source: :tesla,
direction: :outbound,
method: method,
host: URI.parse(normalized_url).host,
path: URI.parse(normalized_url).path,
full_url: redacted_url,
status: nil,
status_class: :server_error,
duration_ms: Event.duration_ms(measurements.duration),
timestamp: System.system_time(:microsecond),
dedup_key: {pid, ts},
error: inspect(metadata[:reason] || metadata[:kind] || "Tesla exception"),
request_headers: req_headers,
response_headers: nil
},
metadata
)
end
# Catch-all for unexpected Tesla telemetry events
def handle_tesla_event(_event_name, _measurements, _metadata, _config), do: nil
@doc """
Handles a Finch telemetry event (`[:finch, :request, :stop]`).
Parses the metadata and measurements into a `Monitorex.Event` struct with
source `:finch` and direction `:outbound`.
## Telemetry metadata shape
%{
url: %URI{} | String.t(),
method: :get | "GET" | …,
status: 200,
pid: pid,
monotonic_time: integer
}
The `url` field may be a `URI.t()` struct or a string; both are handled.
The `method` field may be an atom or string; both are normalised.
"""
@spec handle_finch_event(
event_name :: [atom()],
measurements :: map(),
metadata :: map(),
config :: keyword()
) :: Event.t() | nil
def handle_finch_event([:finch, :request, :stop], measurements, metadata, _config) do
{url_str, method, host, path, req_headers, status} =
case metadata do
%{request: %{method: m} = req} ->
# New Finch telemetry format (Finch.Request struct)
url = build_finch_url(req)
{url, Event.normalize_method(m), req.host, URI.parse(url).path,
redact_headers_from_metadata(req.headers || []), extract_finch_status(metadata)}
%{url: url, method: m, status: s} ->
# Legacy Finch telemetry format
url_str = url_to_string(url)
uri = URI.parse(url_str)
{url_str, Event.normalize_method(m), Event.extract_host(url), uri.path,
redact_headers_from_metadata(metadata[:req_headers] || []), s}
end
resp_headers = redact_headers_from_metadata(metadata[:resp_headers] || [])
ts = metadata[:monotonic_time] || measurements[:monotonic_time] || System.monotonic_time()
pid = metadata[:pid] || self()
tag_slow_request(
%Event{
source: :finch,
direction: :outbound,
method: method,
host: host,
path: path,
full_url: URLRedactor.redact(url_str),
status: status,
status_class: Event.classify_status(status || 0),
duration_ms: Event.duration_ms(measurements.duration),
timestamp: System.system_time(:microsecond),
dedup_key: {pid, ts},
request_headers: req_headers,
response_headers: resp_headers,
request_body: maybe_store_body(metadata[:request_body], :request),
response_body: maybe_store_body(metadata[:response_body], :response)
},
metadata
)
end
def handle_finch_event([:finch, :request, :exception], measurements, metadata, _config) do
ts = metadata[:monotonic_time] || measurements[:monotonic_time] || System.monotonic_time()
pid = metadata[:pid] || self()
case metadata do
%{request: request} ->
url_str = build_finch_url(request)
tag_slow_request(
%Event{
source: :finch,
direction: :outbound,
method: Event.normalize_method(request.method),
host: request.host,
path: URI.parse(url_str).path,
full_url: URLRedactor.redact(url_str),
status: nil,
status_class: :server_error,
duration_ms: Event.duration_ms(measurements.duration),
timestamp: System.system_time(:microsecond),
dedup_key: {pid, ts},
error: inspect(metadata[:reason] || metadata[:result] || "Finch exception"),
request_headers: redact_headers_from_metadata(request.headers || []),
response_headers: nil
},
metadata
)
_ ->
%Event{
source: :finch,
direction: :outbound,
method: nil,
host: nil,
path: nil,
full_url: "unknown",
status: nil,
status_class: :server_error,
duration_ms: Event.duration_ms(measurements.duration),
timestamp: System.system_time(:microsecond),
dedup_key: {pid, ts},
error: "Finch exception",
request_headers: [],
response_headers: nil,
slow: false
}
end
end
# Catch-all for unexpected Finch telemetry events
def handle_finch_event(_event_name, _measurements, _metadata, _config), do: nil
defp build_finch_url(%{scheme: scheme, host: host, port: port, path: path} = req) do
query = if req.query && req.query != "", do: "?#{req.query}", else: ""
"#{scheme}://#{host}#{if port != 443 && port != 80, do: ":#{port}"}#{path}#{query}"
end
defp build_finch_url(_), do: nil
defp extract_finch_status(metadata) do
case metadata[:response] do
%{status: status} ->
status
_ ->
case metadata[:result] do
{:ok, %{status: status}} -> status
_ -> metadata[:status]
end
end
end
@doc """
Handles a Req telemetry event from `req_telemetry` (`[:req, :request, :pipeline, :stop]`).
Parses the metadata and measurements into a `Monitorex.Event` struct with
source `:req` and direction `:outbound`.
## Telemetry metadata shape (req_telemetry pipeline events)
%{
url: %URI{},
method: :get,
status: 200,
resp_headers: [{name, value}],
error: %RuntimeError{},
ref: term(),
metadata: %{}
}
Duration is in `measurements.duration` (native time units).
"""
@spec handle_req_event(
event_name :: [atom()],
measurements :: map(),
metadata :: map(),
config :: keyword()
) :: Event.t() | nil
def handle_req_event([:req, :request, :pipeline, :stop], measurements, metadata, _config) do
url_str = url_to_string(metadata.url)
method = Event.normalize_method(metadata.method)
uri = URI.parse(url_str)
host = uri.host
path = uri.path
resp_headers = redact_headers_from_metadata(metadata[:resp_headers] || [])
normalized_url = UrlNormalizer.normalize(url_str)
redacted_url = URLRedactor.redact(normalized_url)
tag_slow_request(
%Event{
source: :req,
direction: :outbound,
method: method,
host: host,
path: path,
full_url: redacted_url,
status: metadata.status,
status_class: Event.classify_status(metadata.status || 0),
duration_ms: Event.duration_ms(measurements.duration),
timestamp: System.system_time(:microsecond),
response_headers: resp_headers
},
metadata
)
end
def handle_req_event([:req, :request, :pipeline, :error], measurements, metadata, _config) do
url_str = url_to_string(metadata.url)
method = Event.normalize_method(metadata.method)
uri = URI.parse(url_str)
tag_slow_request(
%Event{
source: :req,
direction: :outbound,
method: method,
host: uri.host,
path: uri.path,
full_url: URLRedactor.redact(url_str),
status: nil,
status_class: :server_error,
duration_ms: Event.duration_ms(measurements.duration),
timestamp: System.system_time(:microsecond),
request_headers: redact_headers_from_metadata(metadata[:headers] || []),
response_headers: nil,
error: inspect(metadata[:error] || "Req exception")
},
metadata
)
end
# Catch-all for unexpected Req telemetry events
def handle_req_event(_event_name, _measurements, _metadata, _config), do: nil
@doc """
Handles a Phoenix telemetry event (`[:phoenix, :router_dispatch, :stop]`).
Parses the metadata and measurements into a `Monitorex.Event` struct with
source `:phoenix` and direction `:inbound`.
The consumer is extracted via `ConsumerIdentifier.identify/1`.
If the application config key `:inbound_path_prefixes` is set to a list of
path prefixes, only requests whose path starts with one of the prefixes
produce an event. Returns `nil` (i.e. the event is filtered) when no
prefix matches.
## Telemetry metadata shape
%{
conn: %Plug.Conn{}
}
"""
@spec handle_phoenix_event(
event_name :: [atom()],
measurements :: map(),
metadata :: map(),
config :: keyword()
) :: Event.t() | nil
def handle_phoenix_event([:phoenix, :router_dispatch, :stop], measurements, metadata, _config) do
conn = metadata.conn
path = conn.request_path
inbound_path_prefixes = Application.get_env(:monitorex, :inbound_path_prefixes, nil)
if accepts_path?(path, inbound_path_prefixes) do
req_headers = redact_headers_from_metadata(conn.req_headers)
resp_headers = redact_headers_from_metadata(conn.resp_headers)
tag_slow_request(
%Event{
source: :phoenix,
direction: :inbound,
method: conn.method,
host: conn.host,
path: path,
full_url: conn.request_path,
status: conn.status,
status_class: Event.classify_status(conn.status),
duration_ms: Event.duration_ms(measurements.duration),
consumer: ConsumerIdentifier.identify(conn),
timestamp: System.system_time(:microsecond),
dedup_key: {self(), System.monotonic_time()},
request_headers: req_headers,
response_headers: resp_headers,
request_body: maybe_store_body(metadata[:request_body], :request),
response_body: maybe_store_body(metadata[:response_body], :response)
},
metadata
)
end
end
# Catch-all for unexpected Phoenix telemetry events
def handle_phoenix_event(_event_name, _measurements, _metadata, _config), do: nil
# ── private helpers ──
defp redact_headers_from_metadata(nil), do: nil
defp redact_headers_from_metadata([]), do: []
defp redact_headers_from_metadata(headers) when is_list(headers) do
HeaderRedactor.redact_headers(headers)
end
defp redact_headers_from_metadata(headers) when is_map(headers) do
headers
|> Enum.flat_map(fn {name, values} when is_list(values) ->
Enum.map(values, &{name, &1})
end)
|> HeaderRedactor.redact_headers()
end
defp maybe_store_body(body, _direction) when is_nil(body) or not is_binary(body), do: nil
defp maybe_store_body(body, :request) do
if Application.get_env(:monitorex, :store_request_body, false), do: body, else: nil
end
defp maybe_store_body(body, :response) do
if Application.get_env(:monitorex, :store_response_body, false), do: body, else: nil
end
# Converts a URI struct or string URL to a string
defp url_to_string(%URI{} = uri), do: URI.to_string(uri)
defp url_to_string(url) when is_binary(url), do: url
defp accepts_path?(_path, nil), do: true
defp accepts_path?(path, prefixes) when is_list(prefixes) do
Enum.any?(prefixes, &String.starts_with?(path, &1))
end
@doc """
Post-processes an Event to tag it as slow if its duration exceeds the threshold.
When a request is slow, full request/response bodies are captured even if
body storage is globally disabled. This provides debugging data for slow
requests without the memory overhead of storing all bodies.
The threshold is configured via:
config :monitorex, :slow_request_threshold_ms, 2_000
Set to `nil` or `0` to disable slow request tracing.
## Returns
Modified `%Event{}` with `:slow` set and bodies populated (if applicable).
"""
@spec tag_slow_request(Monitorex.Event.t() | nil, any()) :: Monitorex.Event.t() | nil
def tag_slow_request(event, metadata \\ %{})
def tag_slow_request(nil, _metadata), do: nil
def tag_slow_request(%Monitorex.Event{} = event, metadata) when not is_map(metadata) do
tag_slow_request(event, %{})
end
def tag_slow_request(%Monitorex.Event{} = event, metadata) do
threshold = Application.get_env(:monitorex, :slow_request_threshold_ms, 2_000)
if threshold && threshold > 0 && event.duration_ms && event.duration_ms >= threshold do
req_body = metadata[:request_body] || metadata[:req_body]
resp_body = metadata[:response_body] || metadata[:resp_body]
max_bytes = Application.get_env(:monitorex, :max_body_bytes, 10_000)
%{
event
| slow: true,
request_body: truncate_body(req_body, max_bytes),
response_body: truncate_body(resp_body, max_bytes)
}
else
%{event | slow: false}
end
end
defp truncate_body(nil, _max), do: nil
defp truncate_body(body, max) when is_binary(body) and byte_size(body) <= max, do: body
defp truncate_body(body, max) when is_binary(body) do
<<trimmed::binary-size(^max), _rest::binary>> = body
trimmed <> "\n[truncated]"
end
defp truncate_body(_, _), do: nil
end