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