Packages
Reusable structured audit logging, DB change tracking, and crash reporting for Elixir/Phoenix apps
Current section
Files
Jump to
Current section
Files
lib/audit_trail/reader.ex
defmodule AuditTrail.Reader do
@moduledoc """
Query audit logs from Loki.
## Usage
AuditTrail.Reader.get_logs(%{
action: "item:approved",
actor_id: "user-uuid",
status: "success",
resource: "payment",
operation: "update",
start_date: "2026-06-01",
end_date: "2026-06-12",
search: "item_id",
limit: 100
})
`resource` and `operation` are optional payload fields (only present on
logs where the caller opted into resource/CRUD tagging). Filtering on
them is a JSON field match on the log line, not a stream label, to avoid
adding cardinality.
"""
require Logger
def get_logs(filters \\ %{}) do
url = AuditTrail.Config.loki_read_url()
if is_nil(url) do
{:error, "loki_read_url not configured"}
else
warn_if_default_window(filters)
query = build_query(filters)
params = build_params(query, filters)
case Req.get(url, params: params, receive_timeout: 10_000) do
{:ok, %{status: 200, body: body}} ->
body
|> parse_response()
|> maybe_group(filters)
{:ok, %{status: code}} ->
Logger.error("[AuditTrail.Reader] Loki returned #{code}")
{:error, "Loki returned status #{code}"}
{:error, reason} ->
Logger.error("[AuditTrail.Reader] HTTP error: #{inspect(reason)}")
{:error, reason}
end
end
end
# `get_logs/1` with no start_date/end_date silently scopes to the last
# 24h. A query that legitimately returns 0 results looks identical to one
# that returns 0 because it never looked further back than a day — so
# make that window explicit in the logs whenever it's implicit for the
# caller.
defp warn_if_default_window(filters) do
if is_nil(filters[:start_date]) and is_nil(filters[:end_date]) do
Logger.info(
"[AuditTrail.Reader] get_logs/1 called without start_date/end_date — " <>
"defaulting to the last 24h. Pass start_date/end_date explicitly for a wider window."
)
end
end
defp build_query(filters) do
app = AuditTrail.Config.app_name()
labels = [~s(app="#{app}")]
labels = if t = filters[:type], do: [~s(type="#{t}") | labels], else: labels
labels = if s = filters[:status], do: [~s(status="#{s}") | labels], else: labels
labels = if id = filters[:actor_id], do: [~s(user="#{id}") | labels], else: labels
labels = if tn = filters[:tenant], do: [~s(tenant="#{tn}") | labels], else: labels
selector = "{#{Enum.join(labels, ", ")}}"
selector = with_line_filter(selector, filters[:search])
selector = with_json_field_filter(selector, "resource", filters[:resource])
selector = with_json_field_filter(selector, "resource_id", filters[:resource_id])
with_json_field_filter(selector, "operation", filters[:operation])
end
defp with_line_filter(selector, nil), do: selector
defp with_line_filter(selector, ""), do: selector
defp with_line_filter(selector, term), do: ~s(#{selector} |= "#{term}")
defp with_json_field_filter(selector, _field, nil), do: selector
defp with_json_field_filter(selector, _field, ""), do: selector
defp with_json_field_filter(selector, field, value),
do: ~s(#{selector} | json | #{field}="#{value}")
defp build_params(query, filters) do
[
query: query,
limit: Map.get(filters, :limit, 100),
start: to_loki_time(filters[:start_date], :start),
end: end_bound(filters)
]
end
# `before:` is the cursor returned by `AuditTrail.get_logs_page/1` — the
# nanosecond timestamp of the oldest entry in the previous page. Passing
# it back narrows `end` to just before that entry so the next call picks
# up strictly older logs instead of repeating the same page.
defp end_bound(%{before: cursor}) when is_binary(cursor) do
case Integer.parse(cursor) do
{ns, ""} -> ns - 1
_ -> to_loki_time(nil, :end)
end
end
defp end_bound(filters), do: to_loki_time(filters[:end_date], :end)
defp to_loki_time(nil, :start),
do:
DateTime.utc_now()
|> DateTime.add(-86_400, :second)
|> DateTime.to_unix(:nanosecond)
defp to_loki_time(nil, :end),
do:
DateTime.utc_now()
|> DateTime.to_unix(:nanosecond)
defp to_loki_time(str, type) when is_binary(str) do
case Date.from_iso8601(str) do
{:ok, date} -> to_loki_time(date, type)
_ -> to_loki_time(nil, type)
end
end
defp to_loki_time(%Date{} = d, :start),
do:
d
|> NaiveDateTime.new!(~T[00:00:00])
|> DateTime.from_naive!("Etc/UTC")
|> DateTime.to_unix(:nanosecond)
defp to_loki_time(%Date{} = d, :end),
do:
d
|> NaiveDateTime.new!(~T[23:59:59])
|> DateTime.from_naive!("Etc/UTC")
|> DateTime.to_unix(:nanosecond)
defp parse_response(%{"data" => %{"result" => results}}) do
results
|> Enum.flat_map(fn %{"stream" => labels, "values" => values} ->
Enum.map(values, fn [ts_ns, line] ->
payload =
case Jason.decode(line) do
{:ok, json} -> json
_ -> %{"raw" => line}
end
%{
timestamp: parse_ts(ts_ns),
labels: labels,
payload: payload
}
end)
end)
|> Enum.sort_by(& &1.timestamp, {:desc, DateTime})
end
defp parse_response(_), do: []
defp maybe_group(logs, %{view: "actors"}) do
unique =
logs
|> Enum.uniq_by(fn log -> get_in(log.payload, ["actor_id"]) end)
|> Enum.reject(fn log -> is_nil(get_in(log.payload, ["actor_id"])) end)
{:ok, unique}
end
defp maybe_group(logs, _), do: {:ok, logs}
defp parse_ts(ts_ns) when is_binary(ts_ns) do
{ns, ""} = Integer.parse(ts_ns)
DateTime.from_unix!(ns, :nanosecond)
end
end