diff --git a/.github/workflows/elixir-ci.yml b/.github/workflows/elixir-ci.yml index 79e4525ec9..54748c7c6e 100644 --- a/.github/workflows/elixir-ci.yml +++ b/.github/workflows/elixir-ci.yml @@ -333,6 +333,15 @@ jobs: --cpus="1.0" --shm-size="256m" --tmpfs /var/lib/clickhouse:rw,size=1G + vm: + image: victoriametrics/victoria-metrics:v1.106.1 + ports: + - 8428:8428 + options: >- + --health-cmd "wget -q -O /dev/null http://127.0.0.1:8428/health || exit 1" + --health-interval 10s + --health-timeout 5s + --health-retries 10 env: MIX_ENV: test SHELL: /bin/bash diff --git a/.github/workflows/elixir-migration-check.yml b/.github/workflows/elixir-migration-check.yml index 968a676946..665fdbe62c 100644 --- a/.github/workflows/elixir-migration-check.yml +++ b/.github/workflows/elixir-migration-check.yml @@ -69,6 +69,15 @@ jobs: --cpus="1.0" --shm-size="256m" --tmpfs /var/lib/clickhouse:rw,size=1G + vm: + image: victoriametrics/victoria-metrics:v1.106.1 + ports: + - 8428:8428 + options: >- + --health-cmd "wget -q -O /dev/null http://127.0.0.1:8428/health || exit 1" + --health-interval 10s + --health-timeout 5s + --health-retries 10 env: MIX_ENV: test SHELL: /bin/bash diff --git a/DEVELOPMENT.md b/DEVELOPMENT.md index 6756c96d86..21ee4b58cd 100644 --- a/DEVELOPMENT.md +++ b/DEVELOPMENT.md @@ -48,7 +48,7 @@ make start.pink ```bash # start local databases -docker-compose up -d db clickhouse telegraf +docker-compose up -d db clickhouse telegraf vm # install dependencies make setup diff --git a/docs/docs.logflare.com/docs/backends/victoria-metrics.mdx b/docs/docs.logflare.com/docs/backends/victoria-metrics.mdx index 4749ddaed4..7788b8b51f 100644 --- a/docs/docs.logflare.com/docs/backends/victoria-metrics.mdx +++ b/docs/docs.logflare.com/docs/backends/victoria-metrics.mdx @@ -4,7 +4,7 @@ sidebar_position: 14 # VictoriaMetrics -The VictoriaMetrics backend is **ingest-only**, and sends metric events to [VictoriaMetrics](https://victoriametrics.com) using the Prometheus remote write protocol (v1). +The VictoriaMetrics backend sends metric events to [VictoriaMetrics](https://victoriametrics.com) using the Prometheus remote write protocol (v1). Configure a read URL to also run raw PromQL queries through Logflare's management API. ## Behaviour and configurations @@ -13,11 +13,41 @@ The VictoriaMetrics backend is **ingest-only**, and sends metric events to [Vict The backend can be configured with the following options: - `url` (`string`, required) - The remote write endpoint, such as `https://victoriametrics.example.com/api/v1/write`. Private and reserved addresses are rejected. -- `username` and `password` (`string`, optional) - Basic auth credentials. Both must be provided together. -- `headers` (`map`, optional) - Extra headers attached to each request. Credential headers such as `authorization` are redacted when the backend is read back. +- `query_url` (`string`, optional) - The HTTP(S) read **base URL**, such as `https://victoriametrics.example.com` for a single-node deployment or `https://vmselect.example.com/select/0/prometheus` for a cluster. Include any tenant or proxy path, but omit `/api/v1/query` and `/api/v1/query_range`; Logflare appends the appropriate path. Embedded credentials, query parameters, fragments, and private or reserved addresses are rejected. Leave this unset or clear it to disable querying. +- `username` and `password` (`string`, optional) - Basic auth credentials shared by writes and queries. Both must be provided together. Queries support basic authentication only. +- `headers` (`map`, optional) - Extra headers attached to metric writes only. Queries do not use these headers. Credential headers such as `authorization` are redacted when the backend is read back. - `labels` (`map`, optional) - Static labels added to every time series. These take precedence over event attributes with the same name. Labels can be set through the management API. -Stored credentials are only sent to the host they were entered for. Changing the URL to another host clears the password and credential headers, so they need to be entered again. +Changing the scheme, host, or port of the write URL clears the password and credential headers, so they need to be entered again. Adding a read URL on a different origin, or changing its origin, requires entering the shared password again; write-only headers are preserved. Changing only a URL's path preserves credentials. + +### Querying + +Use an access token with the `private` scope and the ID of a VictoriaMetrics backend that has a `query_url`: + +```bash +curl --get 'https://api.logflare.app/api/query' \ + -H 'X-API-KEY: YOUR-ACCESS-TOKEN' \ + --data-urlencode 'backend_id=123' \ + --data-urlencode 'promql=sum(rate(http_requests_total[5m]))' +``` + +An instant query uses VictoriaMetrics' current time by default. Pass `time` to evaluate it at a specific timestamp. For a range query, supply `start`, `end`, and `step` together: + +```bash +curl --get 'https://api.logflare.app/api/query' \ + -H 'X-API-KEY: YOUR-ACCESS-TOKEN' \ + --data-urlencode 'backend_id=123' \ + --data-urlencode 'promql=sum(rate(http_requests_total[5m]))' \ + --data-urlencode 'start=2026-01-01T00:00:00Z' \ + --data-urlencode 'end=2026-01-01T01:00:00Z' \ + --data-urlencode 'step=60s' +``` + +Do not combine `time` with range parameters. Timestamp and duration values are forwarded to VictoriaMetrics for validation. Use the optional `timeout` parameter to set a query evaluation timeout. The query connection has a 30-second receive timeout. + +Queries return the native Prometheus response envelope, including `status`, `data.resultType`, and `data.result`, or the upstream error envelope. Logflare passes the expression through without SQL conversion or source filtering. Queries can access the data permitted by the configured VictoriaMetrics credentials and tenant URL. + +PromQL is available through this API; the dashboard query editor, saved Endpoints, and Alerts do not support it. The backend's **Test connection** action checks the write endpoint; execute a query to verify the read endpoint. ### Implementation Details @@ -34,6 +64,10 @@ Implementation is based on the [webhook backend](/backends/webhook). - `source` is set to the source name. Renaming a source starts new series under the new name. - `job` is set from the `service.name` resource attribute, prefixed with `service.namespace/` when present, and `instance` from `service.instance.id`. +- `otel_scope_name`, `otel_scope_version`, and `otel_scope_schema_url` contain the corresponding instrumentation scope fields when present in the normalized event. +- `logflare_resource_id` and `logflare_scope_id` identify non-empty resource and scope contexts. These deterministic SHA-256 IDs keep otherwise identical metric series separate when their contexts differ, including nested attributes, array order, and value types. Other resource and scope attributes are not promoted to individual labels. - Data point attributes with non-empty string, number or boolean values become labels; list and map values are skipped. Logflare normalizes attribute keys at ingest, so an attribute such as `http.route` is sent as the `_http_route` label. - Metric and label names are converted into valid Prometheus identifiers, so a metric named `http.server.duration` becomes `http_server_duration`. -- The labels the backend sets (`__name__`, `le`, `source`, `job` and `instance`) always win. An attribute or configured label with one of those names is kept as `exported_`, for example `exported_job`. +- The labels the backend sets (`__name__`, `le`, `source`, `job`, `instance`, the three scope fields, and the two context IDs) always win. An attribute or configured label with one of those names is kept as `exported_`, for example `exported_logflare_resource_id`; it cannot replace the generated identity. + +The context IDs cover the data available after Logflare's ingestion normalization. Ingestion may already normalize attribute keys or remove empty values; the current OTEL conversion also omits schema URLs. The IDs cannot recover distinctions removed before the backend receives an event. Adding these labels starts new series for existing metrics that have resource or scope contexts. diff --git a/docs/docs.logflare.com/docs/concepts/querying.md b/docs/docs.logflare.com/docs/concepts/querying.md index 0b2309fa24..519c58a3a3 100644 --- a/docs/docs.logflare.com/docs/concepts/querying.md +++ b/docs/docs.logflare.com/docs/concepts/querying.md @@ -11,7 +11,7 @@ To run adhoc queries for exploratory analysis, use the Querying or Search functi :::info -You will need to use an access token with the `management` scope to query the management API. +You will need to use an access token with the `private` scope to query the management API. An `ingest` or `query` scoped token **cannot** be used to for this querying API. ::: @@ -20,6 +20,8 @@ An `ingest` or `query` scoped token **cannot** be used to for this querying API. Sources can be queried through SQL using our management API. +[VictoriaMetrics backends support raw PromQL queries](/backends/victoria-metrics#querying) through the same API using `promql` and `backend_id`. The SQL parameters and caveats below do not apply to PromQL. + The following query parameters are available: - `?sql=` (string): the SQL query string. diff --git a/lib/logflare/backends/adaptor/clickhouse_adaptor/ingester.ex b/lib/logflare/backends/adaptor/clickhouse_adaptor/ingester.ex index d81d21688e..facadc6e12 100644 --- a/lib/logflare/backends/adaptor/clickhouse_adaptor/ingester.ex +++ b/lib/logflare/backends/adaptor/clickhouse_adaptor/ingester.ex @@ -6,9 +6,9 @@ defmodule Logflare.Backends.Adaptor.ClickHouseAdaptor.Ingester do import Logflare.Utils.Guards alias Logflare.Backends.Adaptor.ClickHouseAdaptor.EndpointUtils - alias Logflare.Backends.Adaptor.ClickHouseAdaptor.FinchPoolTimeoutNormalizer alias Logflare.Backends.Adaptor.ClickHouseAdaptor.QueryTemplates alias Logflare.Backends.Adaptor.ClickHouseAdaptor.RowBinaryEncoder + alias Logflare.Backends.Adaptor.HttpBased.FinchPoolTimeoutNormalizer alias Logflare.Backends.Backend alias Logflare.LogEvent alias Logflare.LogEvent.TypeDetection diff --git a/lib/logflare/backends/adaptor/clickhouse_adaptor/finch_pool_timeout_normalizer.ex b/lib/logflare/backends/adaptor/http_based/finch_pool_timeout_normalizer.ex similarity index 65% rename from lib/logflare/backends/adaptor/clickhouse_adaptor/finch_pool_timeout_normalizer.ex rename to lib/logflare/backends/adaptor/http_based/finch_pool_timeout_normalizer.ex index 6269c3560e..1f976635af 100644 --- a/lib/logflare/backends/adaptor/clickhouse_adaptor/finch_pool_timeout_normalizer.ex +++ b/lib/logflare/backends/adaptor/http_based/finch_pool_timeout_normalizer.ex @@ -1,16 +1,10 @@ -defmodule Logflare.Backends.Adaptor.ClickHouseAdaptor.FinchPoolTimeoutNormalizer do +defmodule Logflare.Backends.Adaptor.HttpBased.FinchPoolTimeoutNormalizer do @moduledoc """ Tesla middleware that turns a Finch HTTP/1 pool checkout timeout into `{:error, :pool_timeout}`. - Finch only returns `%Finch.Error{reason: :pool_timeout}` for HTTP/2 pools. On an - HTTP/1 pool — which both ClickHouse ingest pools are — a checkout timeout surfaces - as a `NimblePool` exit that `Finch.HTTP1.Pool` catches and re-raises as a - `RuntimeError`. An exception is invisible to `Tesla.Middleware.Retry`, which only - inspects the return value of the stack below it, so a saturated pool would neither - be retried nor counted as an insert failure. - - Must sit *after* `Tesla.Middleware.Retry` in the middleware list so it runs inside - the retry loop and its `{:error, :pool_timeout}` is visible to `should_retry`. + Finch raises a `RuntimeError` when an HTTP/1 connection cannot be checked out in + time. Place this middleware before the adapter and after any retry middleware + so callers can handle the timeout as an ordinary transport error. Finch gives no structured way to tell this exception apart from any other `RuntimeError`, so it is matched on a stable fragment of the message and anything diff --git a/lib/logflare/backends/adaptor/victoria_metrics_adaptor.ex b/lib/logflare/backends/adaptor/victoria_metrics_adaptor.ex index 7db546bedd..af909cddc6 100644 --- a/lib/logflare/backends/adaptor/victoria_metrics_adaptor.ex +++ b/lib/logflare/backends/adaptor/victoria_metrics_adaptor.ex @@ -23,12 +23,16 @@ defmodule Logflare.Backends.Adaptor.VictoriaMetricsAdaptor do * `source` - the source name * `job` and `instance` - from the `service.namespace`/`service.name` and `service.instance.id` resource attributes, when present + * `logflare_resource_id` and `logflare_scope_id` - stable IDs for the normalized + resource and scope maps, including nested attributes + * `otel_scope_name`, `otel_scope_version` and `otel_scope_schema_url` - readable + scope metadata, when present * data point attributes with non-empty string, number or boolean values; list and map values are skipped. Names are as stored by Logflare, which normalizes keys at ingest (e.g. `http.route` becomes `_http_route`) * the optional `labels` config map, which wins over attributes on collision - Labels the adaptor sets (`__name__`, `le`, `source`, `job`, `instance`) always win. + Labels the adaptor sets, including `__name__` and `le`, always win. An attribute or config label with one of those names is kept as `exported_`. """ @@ -37,7 +41,9 @@ defmodule Logflare.Backends.Adaptor.VictoriaMetricsAdaptor do require Logger alias Logflare.Backends.Adaptor.HttpBased.Headers + alias Logflare.Backends.Adaptor.VictoriaMetricsAdaptor.Query alias Logflare.Backends.Adaptor.VictoriaMetricsAdaptor.RemoteWrite + alias Logflare.Backends.Adaptor.VictoriaMetricsAdaptor.SeriesIdentity alias Logflare.Backends.Adaptor.WebhookAdaptor alias Logflare.Backends.Backend alias Logflare.LogEvent @@ -45,7 +51,10 @@ defmodule Logflare.Backends.Adaptor.VictoriaMetricsAdaptor do alias Logflare.Sources.Source alias Logflare.Utils - @reserved_labels ["__name__", "le", "source", "job", "instance"] + @reserved_labels ~w( + __name__ le source job instance logflare_resource_id logflare_scope_id + otel_scope_name otel_scope_version otel_scope_schema_url + ) @redacted_value Headers.redacted_value() @max_float 1.7_976_931_348_623_157e308 # Snappy encoding of an empty WriteRequest. snappyer returns "" for empty input, @@ -123,17 +132,31 @@ defmodule Logflare.Backends.Adaptor.VictoriaMetricsAdaptor do def cast_config(params, existing_config \\ %{}) do changeset = {existing_config, - %{url: :string, headers: :map, username: :string, password: :string, labels: :map}} - |> Ecto.Changeset.cast(params, [:url, :headers, :username, :password, :labels]) + %{ + url: :string, + query_url: :string, + headers: :map, + username: :string, + password: :string, + labels: :map + }} + |> Ecto.Changeset.cast(params, [:url, :query_url, :headers, :username, :password, :labels]) + + cond do + destination_changed?(changeset, existing_config) -> + changeset + |> WebhookAdaptor.unredact_credentials(Map.drop(existing_config, [:headers, "headers"])) + |> require_new_credentials(existing_config) - if destination_changed?(changeset, existing_config) do - changeset - |> WebhookAdaptor.unredact_credentials(Map.drop(existing_config, [:headers, "headers"])) - |> require_new_credentials(existing_config) - else - changeset - |> WebhookAdaptor.unredact_credentials(existing_config) - |> unredact_password() + query_destination_changed?(changeset, existing_config) -> + changeset + |> WebhookAdaptor.unredact_credentials(existing_config) + |> require_new_password() + + true -> + changeset + |> WebhookAdaptor.unredact_credentials(existing_config) + |> unredact_password() end end @@ -143,6 +166,7 @@ defmodule Logflare.Backends.Adaptor.VictoriaMetricsAdaptor do changeset |> Ecto.Changeset.validate_required([:url]) |> Ecto.Changeset.validate_format(:url, ~r/https?\:\/\/.+/) + |> Ecto.Changeset.validate_change(:query_url, &validate_query_url/2) |> validate_user_pass() |> WebhookAdaptor.validate_no_ssrf() end @@ -162,6 +186,18 @@ defmodule Logflare.Backends.Adaptor.VictoriaMetricsAdaptor do WebhookAdaptor.test_connection(backend, @empty_write_request) end + @spec execute_promql(Backend.t(), String.t(), map()) :: + {:ok, map()} | {:error, {pos_integer(), map()}} + defdelegate execute_promql(backend, query, params), to: Query, as: :execute + + @spec validate_query_url(:query_url, String.t()) :: keyword(String.t()) + defp validate_query_url(:query_url, url) do + case Query.validate_url(url) do + :ok -> [] + {:error, message} -> [query_url: message] + end + end + defp put_basic_auth(headers, nil), do: headers defp put_basic_auth(headers, encoded), do: Map.put(headers, "authorization", "Basic #{encoded}") @@ -177,23 +213,54 @@ defmodule Logflare.Backends.Adaptor.VictoriaMetricsAdaptor do # Stored credentials only ever go to the destination they were entered for. When the # URL moves to another origin, the password and credential headers must be entered # again rather than restored from the redacted form or kept from storage. + @spec require_new_credentials(Ecto.Changeset.t(), map()) :: Ecto.Changeset.t() defp require_new_credentials(changeset, existing_config) do - changeset = - if Ecto.Changeset.get_change(changeset, :password) in [nil, @redacted_value], - do: Ecto.Changeset.put_change(changeset, :password, nil), - else: changeset + headers = + case Map.fetch(changeset.params, "headers") do + {:ok, submitted} when is_map(submitted) -> + for {key, value} <- submitted, + value != @redacted_value, + into: %{}, + do: {Headers.normalize_key(key), value} + + {:ok, nil} -> + nil + + _ -> + stored = + Map.get(existing_config, :headers) || Map.get(existing_config, "headers") || %{} + + for {key, value} <- stored, + not Headers.sensitive?(key), + into: %{}, + do: {Headers.normalize_key(key), value} + end - case Ecto.Changeset.get_change(changeset, :headers) do - nil -> - stored = Map.get(existing_config, :headers) || Map.get(existing_config, "headers") || %{} + changeset + |> require_new_password() + |> Ecto.Changeset.put_change(:headers, headers) + end - kept = - for {key, value} <- stored, not Headers.sensitive?(key), into: %{}, do: {key, value} + @spec require_new_password(Ecto.Changeset.t()) :: Ecto.Changeset.t() + defp require_new_password(changeset) do + case changeset.params["password"] do + password when is_binary(password) and password not in ["", @redacted_value] -> + Ecto.Changeset.put_change(changeset, :password, password) - Ecto.Changeset.put_change(changeset, :headers, kept) + _ -> + Ecto.Changeset.put_change(changeset, :password, nil) + end + end - _submitted -> - changeset + @spec query_destination_changed?(Ecto.Changeset.t(), map()) :: boolean() + defp query_destination_changed?(changeset, existing_config) do + previous_url = + existing_config[:query_url] || existing_config["query_url"] || + existing_config[:url] || existing_config["url"] + + case Ecto.Changeset.get_change(changeset, :query_url) do + url when is_binary(url) and is_binary(previous_url) -> origin(url) != origin(previous_url) + _ -> false end end @@ -206,9 +273,13 @@ defmodule Logflare.Backends.Adaptor.VictoriaMetricsAdaptor do end end + @spec origin(String.t()) :: + {String.t() | nil, String.t() | nil, non_neg_integer() | nil} | :invalid defp origin(url) do - uri = URI.parse(url) - {uri.scheme, uri.host && String.downcase(uri.host), uri.port} + case URI.new(url) do + {:ok, uri} -> {uri.scheme, uri.host && String.downcase(uri.host), uri.port} + {:error, _reason} -> :invalid + end end defp encode_write_request(series) do @@ -239,9 +310,13 @@ defmodule Logflare.Backends.Adaptor.VictoriaMetricsAdaptor do defp labeled_samples(%{event_type: :metric, body: body} = event, context) do with {:ok, kind} <- series_kind(body), {:ok, name} <- metric_name(body["event_message"]), - {:ok, point} <- data_point(kind, body) do - labels = series_labels(body, Map.get(context.source_names, event.source_id), context) + {:ok, point} <- data_point(kind, body), + labels when is_list(labels) <- + series_labels(body, Map.get(context.source_names, event.source_id), context) do {:ok, samples(point, name, labels, micro_to_ms(body["timestamp"]))} + else + :error -> {:drop, :invalid} + {:drop, _reason} = drop -> drop end end @@ -362,17 +437,20 @@ defmodule Logflare.Backends.Adaptor.VictoriaMetricsAdaptor do defp integer_value(_value), do: :negative_infinity defp series_labels(body, source_name, context) do - adaptor_labels = - body["resource"] - |> resource_labels() - |> Map.put("source", source_name || "unknown") - - body["attributes"] - |> flat_labels() - |> Map.merge(context.static_labels) - |> export_reserved() - |> Map.merge(adaptor_labels) - |> Enum.sort() + with %{} = identity_labels <- SeriesIdentity.labels(body) do + adaptor_labels = + body["resource"] + |> resource_labels() + |> Map.merge(identity_labels) + |> Map.put("source", source_name || "unknown") + + body["attributes"] + |> flat_labels() + |> Map.merge(context.static_labels) + |> export_reserved() + |> Map.merge(adaptor_labels) + |> Enum.sort() + end end # LogEvent.make/2 stores keys in BigQuery column form, so `service.name` arrives as diff --git a/lib/logflare/backends/adaptor/victoria_metrics_adaptor/query.ex b/lib/logflare/backends/adaptor/victoria_metrics_adaptor/query.ex new file mode 100644 index 0000000000..95cfccb33b --- /dev/null +++ b/lib/logflare/backends/adaptor/victoria_metrics_adaptor/query.ex @@ -0,0 +1,176 @@ +defmodule Logflare.Backends.Adaptor.VictoriaMetricsAdaptor.Query do + @moduledoc """ + Executes raw PromQL against a VictoriaMetrics read endpoint. + + Query expressions and Prometheus response data are forwarded without SQL + parsing or result conversion. Only stored basic authentication is used. + """ + + alias Logflare.Backends.Adaptor.HttpBased.FinchPoolTimeoutNormalizer + alias Logflare.Backends.Adaptor.HttpBased.SSRFProtection + alias Logflare.Backends.Backend + alias Logflare.Utils + alias Logflare.Utils.SSRF + + @parameters ~w(time start end step timeout) + @range_parameters ~w(start end step) + @receive_timeout 30_000 + + @type result :: {:ok, map()} | {:error, {pos_integer(), map()}} + + @spec execute(Backend.t(), String.t(), map()) :: result() + def execute(%Backend{config: config}, query, params) do + with :ok <- validate_query(query), + {:ok, params} <- query_parameters(params), + :ok <- validate_url(config[:query_url]) do + url = String.trim_trailing(config.query_url, "/") <> query_path(params) + + client = + Tesla.client( + [Tesla.Middleware.Telemetry, SSRFProtection, FinchPoolTimeoutNormalizer], + {Tesla.Adapter.Finch, + name: Logflare.FinchDefaultHttp1, receive_timeout: @receive_timeout} + ) + + client + |> Tesla.post(url, URI.encode_query(Map.put(params, "query", query)), + headers: headers(config) + ) + |> decode_response() + |> response() + else + {:error, message} -> error(400, "bad_data", message) + end + end + + @spec validate_url(term()) :: :ok | {:error, String.t()} + def validate_url(url) when is_binary(url) do + case URI.new(url) do + {:ok, + %URI{scheme: scheme, host: host, port: port, userinfo: nil, query: nil, fragment: nil}} + when scheme in ["http", "https"] and is_binary(host) and host != "" and + port in 1..65_535 -> + case SSRF.safe_resolve(host) do + {:ok, _address} -> :ok + {:error, reason} -> {:error, reason} + end + + _ -> + {:error, + "query_url must be an HTTP(S) base URL without credentials, query parameters, or a fragment"} + end + end + + def validate_url(_), do: {:error, "Configure query_url on the VictoriaMetrics backend first"} + + @spec validate_query(term()) :: :ok | {:error, String.t()} + defp validate_query(query) when is_binary(query) do + if String.valid?(query) and String.trim(query) != "", + do: :ok, + else: {:error, "promql must be a non-empty UTF-8 string"} + end + + defp validate_query(_), do: {:error, "promql must be a non-empty UTF-8 string"} + + @spec query_parameters(term()) :: {:ok, map()} | {:error, String.t()} + defp query_parameters(params) when is_map(params) do + params = Map.take(params, @parameters) + + with :ok <- validate_parameter_values(params), + :ok <- validate_range(params) do + {:ok, params} + end + end + + defp query_parameters(_), do: {:error, "Query parameters must be an object"} + + @spec validate_parameter_values(map()) :: :ok | {:error, String.t()} + defp validate_parameter_values(params) do + case Enum.find(params, fn {_key, value} -> not parameter_value?(value) end) do + nil -> :ok + {key, _value} -> {:error, "#{key} must be a non-empty string or a number"} + end + end + + @spec parameter_value?(term()) :: boolean() + defp parameter_value?(value) when is_binary(value), + do: String.valid?(value) and String.trim(value) != "" + + defp parameter_value?(value), do: is_number(value) + + @spec validate_range(map()) :: :ok | {:error, String.t()} + defp validate_range(params) do + range? = Enum.any?(@range_parameters, &Map.has_key?(params, &1)) + + cond do + range? and not Enum.all?(@range_parameters, &Map.has_key?(params, &1)) -> + {:error, "Range queries require start, end, and step together"} + + range? and Map.has_key?(params, "time") -> + {:error, "time cannot be combined with start, end, or step"} + + true -> + :ok + end + end + + @spec query_path(map()) :: String.t() + defp query_path(%{"start" => _start}), do: "/api/v1/query_range" + defp query_path(_params), do: "/api/v1/query" + + @spec headers(map()) :: [{String.t(), String.t()}] + defp headers(config) do + headers = [ + {"accept", "application/json"}, + {"content-type", "application/x-www-form-urlencoded"} + ] + + case Utils.encode_basic_auth(config) do + nil -> headers + encoded -> [{"authorization", "Basic " <> encoded} | headers] + end + end + + @spec decode_response(Tesla.Env.result()) :: Tesla.Env.result() + defp decode_response({:ok, %Tesla.Env{body: body} = env}) when is_binary(body) do + case Jason.decode(body) do + {:ok, decoded} -> {:ok, %{env | body: decoded}} + {:error, _reason} -> {:ok, env} + end + end + + defp decode_response(result), do: result + + @spec response(Tesla.Env.result()) :: result() + defp response( + {:ok, %Tesla.Env{status: 200, body: %{"status" => "success", "data" => data} = body}} + ) + when is_map(data), + do: {:ok, body} + + defp response({:ok, %Tesla.Env{status: status, body: %{"status" => "error"} = body}}) + when status in 400..599, + do: {:error, {status, body}} + + defp response({:ok, %Tesla.Env{status: status}}) when status in [401, 403], + do: error(status, "unauthorized", "VictoriaMetrics rejected the backend credentials") + + defp response({:ok, %Tesla.Env{status: status}}) when status in 400..599, + do: error(status, "upstream_error", "VictoriaMetrics returned HTTP #{status}") + + defp response({:error, :pool_timeout}), + do: error(503, "unavailable", "VictoriaMetrics query capacity is temporarily unavailable") + + defp response({:error, :timeout}), + do: error(504, "timeout", "VictoriaMetrics query timed out") + + defp response({:error, _reason}), + do: error(502, "unavailable", "Unable to query VictoriaMetrics") + + defp response(_response), + do: error(502, "bad_response", "VictoriaMetrics returned an invalid query response") + + @spec error(pos_integer(), String.t(), String.t()) :: result() + defp error(status, type, message), + do: {:error, {status, %{"status" => "error", "errorType" => type, "error" => message}}} +end diff --git a/lib/logflare/backends/adaptor/victoria_metrics_adaptor/series_identity.ex b/lib/logflare/backends/adaptor/victoria_metrics_adaptor/series_identity.ex new file mode 100644 index 0000000000..cea7323815 --- /dev/null +++ b/lib/logflare/backends/adaptor/victoria_metrics_adaptor/series_identity.ex @@ -0,0 +1,90 @@ +defmodule Logflare.Backends.Adaptor.VictoriaMetricsAdaptor.SeriesIdentity do + @moduledoc """ + Preserves the resource and scope identity available in a normalized LogEvent. + + Context IDs hash a typed JSON representation with sorted map entries, ordered + lists, base64 binaries, and IEEE 754 float bytes. This keeps identities stable + across map insertion order and runtime upgrades without flattening attributes. + """ + + @scope_fields ~w(name version schema_url) + + @spec labels(map()) :: %{String.t() => String.t()} | :error + def labels(body) do + if valid_context?(body["resource"]) and valid_context?(body["scope"]) do + body + |> Map.get("scope") + |> scope_labels() + |> put_context_id("logflare_resource_id", body["resource"]) + |> put_context_id("logflare_scope_id", body["scope"]) + else + :error + end + end + + @spec valid_context?(term()) :: boolean() + defp valid_context?(context) when is_map(context), do: valid_value?(context) + defp valid_context?(_context), do: true + + @spec valid_value?(term()) :: boolean() + defp valid_value?(value) when is_map(value) do + value + |> Map.to_list() + |> Enum.all?(fn {key, value} -> valid_value?(key) and valid_value?(value) end) + end + + defp valid_value?(value) when is_list(value), do: valid_list?(value) + defp valid_value?(value) when is_tuple(value), do: value |> Tuple.to_list() |> valid_list?() + + defp valid_value?(value) + when is_binary(value) or is_integer(value) or is_float(value) or is_atom(value), + do: true + + defp valid_value?(_value), do: false + + @spec valid_list?(term()) :: boolean() + defp valid_list?([]), do: true + defp valid_list?([head | tail]), do: valid_value?(head) and valid_list?(tail) + defp valid_list?(_value), do: false + + @spec scope_labels(term()) :: map() + defp scope_labels(scope) when is_map(scope) do + for field <- @scope_fields, + value = Map.get(scope, field), + is_binary(value) and value != "", + into: %{}, + do: {"otel_scope_" <> field, value} + end + + defp scope_labels(_scope), do: %{} + + @spec put_context_id(map(), String.t(), term()) :: map() + defp put_context_id(labels, name, context) when is_map(context) and map_size(context) > 0 do + encoded = context |> canonical() |> Jason.encode!() + id = :sha256 |> :crypto.hash(encoded) |> Base.encode16(case: :lower) + Map.put(labels, name, id) + end + + defp put_context_id(labels, _name, _context), do: labels + + @spec canonical(term()) :: list() + defp canonical(value) when is_map(value) do + entries = + value + |> Map.to_list() + |> Enum.map(fn {key, value} -> [canonical(key), canonical(value)] end) + |> Enum.sort() + + ["map", entries] + end + + defp canonical(value) when is_list(value), do: ["list", Enum.map(value, &canonical/1)] + + defp canonical(value) when is_tuple(value), + do: ["tuple", value |> Tuple.to_list() |> Enum.map(&canonical/1)] + + defp canonical(value) when is_binary(value), do: ["binary", Base.encode64(value)] + defp canonical(value) when is_integer(value), do: ["integer", Integer.to_string(value)] + defp canonical(value) when is_float(value), do: ["float", Base.encode16(<>)] + defp canonical(value) when is_atom(value), do: ["atom", Atom.to_string(value)] +end diff --git a/lib/logflare_web/controllers/api/query_controller.ex b/lib/logflare_web/controllers/api/query_controller.ex index 1cfc929e9a..79fc2bff2c 100644 --- a/lib/logflare_web/controllers/api/query_controller.ex +++ b/lib/logflare_web/controllers/api/query_controller.ex @@ -4,6 +4,7 @@ defmodule LogflareWeb.Api.QueryController do alias Logflare.Alerting alias Logflare.Backends + alias Logflare.Backends.Adaptor.VictoriaMetricsAdaptor alias Logflare.Backends.Backend alias Logflare.Endpoints alias Logflare.Endpoints.EndpointQuery @@ -13,8 +14,10 @@ defmodule LogflareWeb.Api.QueryController do alias LogflareWeb.OpenApi.BadRequest alias LogflareWeb.OpenApi.One alias LogflareWeb.OpenApi.Unauthorized + alias LogflareWeb.OpenApiSchemas.PromQLQueryResponse alias LogflareWeb.OpenApiSchemas.QueryParseResult alias LogflareWeb.OpenApiSchemas.QueryResult + alias OpenApiSpex.Schema action_fallback(LogflareWeb.Api.FallbackController) @@ -84,6 +87,14 @@ defmodule LogflareWeb.Api.QueryController do operation(:query, summary: "Execute a query", parameters: [ + promql: [ + in: :query, + description: + "Raw PromQL expression. Requires a VictoriaMetrics backend_id and cannot be combined with SQL parameters.", + type: :string, + required: false, + example: "sum(rate(http_requests_total[5m]))" + ], sql: [ in: :query, description: @@ -115,18 +126,82 @@ defmodule LogflareWeb.Api.QueryController do backend_id: [ in: :query, description: - "Backend ID to execute the query against. The backend type determines the SQL language.", + "Backend ID to execute the query against. Required for PromQL; determines the language for SQL.", type: :integer, required: false + ], + time: [ + in: :query, + description: + "PromQL instant query evaluation time as a Unix timestamp or RFC3339 string.", + type: :string, + required: false + ], + start: [ + in: :query, + description: + "PromQL range start as a Unix timestamp or RFC3339 string. Requires end and step.", + type: :string, + required: false + ], + end: [ + in: :query, + description: + "PromQL range end as a Unix timestamp or RFC3339 string. Requires start and step.", + type: :string, + required: false + ], + step: [ + in: :query, + description: + "PromQL range resolution as seconds or a duration. Cannot be combined with time.", + type: :string, + required: false + ], + timeout: [ + in: :query, + description: "PromQL query evaluation timeout as a duration.", + type: :string, + required: false ] ], responses: %{ - 200 => One.response(QueryResult), - 400 => BadRequest.response(), - 401 => Unauthorized.response() + 200 => + {"SQL rows or a native Prometheus response", "application/json", + %Schema{oneOf: [QueryResult, PromQLQueryResponse]}}, + 400 => + {"Invalid query", "application/json", + %Schema{anyOf: [BadRequest.schema(), PromQLQueryResponse]}}, + 401 => + {"Unauthorized", "application/json", + %Schema{anyOf: [Unauthorized.schema(), PromQLQueryResponse]}}, + "4XX" => {"PromQL request error", "application/json", PromQLQueryResponse}, + "5XX" => {"PromQL backend error", "application/json", PromQLQueryResponse}, + :default => {"PromQL error", "application/json", PromQLQueryResponse} } ) + def query(%{assigns: %{user: user}} = conn, %{"promql" => query} = params) do + with :ok <- validate_promql_params(params), + {:ok, backend} <- fetch_backend(user, params, :promql), + {:ok, response} <- + VictoriaMetricsAdaptor.execute_promql( + backend, + query, + Map.take(params, ~w(time start end step timeout)) + ) do + json(conn, response) + else + {:error, {status, response}} -> + conn |> put_status(status) |> json(response) + + {:error, message} -> + conn + |> put_status(400) + |> json(%{status: "error", errorType: "bad_data", error: message}) + end + end + def query(%{assigns: %{user: user}} = conn, params) do with {:ok, requested_language, sql} <- extract_query(params), {:ok, backend} <- fetch_backend(user, params), @@ -137,6 +212,15 @@ defmodule LogflareWeb.Api.QueryController do end end + @spec validate_promql_params(map()) :: :ok | {:error, String.t()} + defp validate_promql_params(params) do + if Enum.any?(~w(sql bq_sql ch_sql pg_sql), &Map.has_key?(params, &1)) do + {:error, "promql cannot be combined with SQL parameters"} + else + :ok + end + end + @spec extract_query(map()) :: {:ok, :infer | :bq_sql | :ch_sql | :pg_sql, String.t()} | {:error, String.t()} defp extract_query(%{"sql" => sql}), do: {:ok, :infer, sql} @@ -156,32 +240,62 @@ defmodule LogflareWeb.Api.QueryController do defp resolve_language(language, _backend), do: language - @spec fetch_backend(User.t(), map()) :: {:ok, Backend.t() | nil} | {:error, String.t()} - defp fetch_backend(_user, %{"backend_id" => backend_id}) when backend_id in [nil, ""], + @spec fetch_backend(User.t(), map(), :sql | :promql) :: + {:ok, Backend.t() | nil} | {:error, String.t()} + defp fetch_backend(user, params, query_type \\ :sql) + + defp fetch_backend(_user, %{"backend_id" => backend_id}, :sql) when backend_id in [nil, ""], do: {:ok, nil} - defp fetch_backend(user, %{"backend_id" => backend_id}) when is_binary(backend_id) do + defp fetch_backend(_user, %{"backend_id" => backend_id}, :promql) when backend_id in [nil, ""], + do: {:error, "backend_id is required for PromQL queries"} + + defp fetch_backend(user, %{"backend_id" => backend_id}, query_type) + when is_binary(backend_id) do case Integer.parse(backend_id) do - {id, ""} -> fetch_backend(user, %{"backend_id" => id}) + {id, ""} -> fetch_backend(user, %{"backend_id" => id}, query_type) _ -> {:error, "Invalid backend_id: must be an integer"} end end - defp fetch_backend(user, %{"backend_id" => backend_id}) when is_integer(backend_id) do + defp fetch_backend(_user, %{"backend_id" => backend_id}, :promql) + when is_integer(backend_id) and backend_id not in 1..9_223_372_036_854_775_807, + do: {:error, "Invalid backend_id: must be a positive 64-bit integer"} + + defp fetch_backend(user, %{"backend_id" => backend_id}, query_type) + when is_integer(backend_id) do case Backends.get_backend(backend_id) do %Backend{user_id: user_id} = backend when user_id == user.id -> - if Backends.Adaptor.can_query?(backend) do - {:ok, backend} - else - {:error, "Backend does not support querying"} - end + validate_backend_type(backend, query_type) _ -> {:error, "Backend not found"} end end - defp fetch_backend(_user, _params), do: {:ok, nil} + defp fetch_backend(_user, %{"backend_id" => _backend_id}, :promql), + do: {:error, "Invalid backend_id: must be an integer"} + + defp fetch_backend(_user, _params, :promql), + do: {:error, "backend_id is required for PromQL queries"} + + defp fetch_backend(_user, _params, :sql), do: {:ok, nil} + + @spec validate_backend_type(Backend.t(), :sql | :promql) :: + {:ok, Backend.t()} | {:error, String.t()} + defp validate_backend_type(%Backend{type: :victoria_metrics} = backend, :promql), + do: {:ok, backend} + + defp validate_backend_type(_backend, :promql), + do: {:error, "Backend does not support PromQL queries"} + + defp validate_backend_type(backend, :sql) do + if Backends.Adaptor.can_query?(backend) do + {:ok, backend} + else + {:error, "Backend does not support querying"} + end + end @spec build_query_opts(Backend.t() | nil) :: keyword() defp build_query_opts(nil), do: [] diff --git a/lib/logflare_web/live/backends/backends_live.ex b/lib/logflare_web/live/backends/backends_live.ex index 3194259b7b..0266a77ec4 100644 --- a/lib/logflare_web/live/backends/backends_live.ex +++ b/lib/logflare_web/live/backends/backends_live.ex @@ -8,6 +8,7 @@ defmodule LogflareWeb.BackendsLive do alias Logflare.Backends alias Logflare.Backends.Adaptor.HttpBased.Headers + alias Logflare.Backends.Adaptor.VictoriaMetricsAdaptor alias Logflare.Backends.Adaptor.WebhookAdaptor alias Logflare.Backends.Backend alias Logflare.Rules @@ -78,7 +79,8 @@ defmodule LogflareWeb.BackendsLive do %{"backend" => params}, %{assigns: %{live_action: :edit}} = socket ) do - with {:ok, params} <- transform_params(params, existing_headers(socket.assigns.backend)) do + with {:ok, params} <- + transform_params(params, existing_headers(socket.assigns.backend, params)) do socket = case Backends.update_backend(socket.assigns.backend, params) do {:ok, backend} -> @@ -488,12 +490,25 @@ defmodule LogflareWeb.BackendsLive do |> assemble_read_clusters() end - @spec existing_headers(Backend.t() | nil) :: map() - defp existing_headers(%Backend{config: config}) when is_map(config) do + @spec existing_headers(Backend.t() | nil, map()) :: map() + defp existing_headers( + %Backend{type: :victoria_metrics, config: config}, + %{"config" => submitted_config} + ) + when is_map(config) and is_map(submitted_config) do + headers = + submitted_config + |> VictoriaMetricsAdaptor.cast_config(config) + |> Ecto.Changeset.get_field(:headers) + + Headers.normalize_keys(headers || %{}) + end + + defp existing_headers(%Backend{config: config}, _params) when is_map(config) do Headers.normalize_keys(Map.get(config, :headers) || %{}) end - defp existing_headers(_backend), do: %{} + defp existing_headers(_backend, _params), do: %{} @spec header_form_keys(map()) :: [String.t()] defp header_form_keys(config) when is_map(config) do diff --git a/lib/logflare_web/live/backends/components/backend_form.heex b/lib/logflare_web/live/backends/components/backend_form.heex index b55d334ece..05b54c6e71 100644 --- a/lib/logflare_web/live/backends/components/backend_form.heex +++ b/lib/logflare_web/live/backends/components/backend_form.heex @@ -532,6 +532,14 @@ The Prometheus remote write endpoint, e.g. https://victoriametrics.example.com/api/v1/write. Only metric events are sent +
+ {label(f_config, :query_url, "Query Base URL")} + {text_input(f_config, :query_url, class: "form-control")} + + Optional read URL for the query API, e.g. https://victoriametrics.example.com or https://vmselect.example.com/select/0/prometheus. + Include any tenant or proxy path, but omit /api/v1/query. Leave blank to disable querying. + +
{label(f_config, :username, "Basic Auth Username")} {text_input(f_config, :username, class: "form-control")} @@ -542,7 +550,7 @@ value: Logflare.Backends.Adaptor.HttpBased.Headers.mask_value(input_value(f_config, :password)) )} - For basic auth, both username and password must be provided. Changing the URL to another host requires entering the password again. + Both fields are required for basic authentication and are shared by writes and queries. Changing the scheme, host or port of either URL requires entering the password again. Custom headers apply only to writes.
<.header_inputs form={f_config} /> diff --git a/lib/logflare_web/open_api_schemas.ex b/lib/logflare_web/open_api_schemas.ex index cbc0411762..837f86a205 100644 --- a/lib/logflare_web/open_api_schemas.ex +++ b/lib/logflare_web/open_api_schemas.ex @@ -507,7 +507,13 @@ defmodule LogflareWeb.OpenApiSchemas do defmodule VictoriaMetricsConfigSchema do @properties %{ url: %Schema{type: :string}, - headers: %Schema{type: :object}, + query_url: %Schema{ + type: :string, + nullable: true, + description: + "Optional HTTP(S) read base URL, including tenant or proxy paths, without /api/v1/query. Enables raw PromQL queries." + }, + headers: %Schema{type: :object, description: "Additional headers for metric writes only."}, username: %Schema{type: :string, nullable: true}, password: %Schema{type: :string, nullable: true}, labels: %Schema{type: :object, nullable: true} diff --git a/lib/logflare_web/open_api_schemas/promql_query_response.ex b/lib/logflare_web/open_api_schemas/promql_query_response.ex new file mode 100644 index 0000000000..f560a7a245 --- /dev/null +++ b/lib/logflare_web/open_api_schemas/promql_query_response.ex @@ -0,0 +1,34 @@ +defmodule LogflareWeb.OpenApiSchemas.PromQLQueryResponse do + @moduledoc """ + Native Prometheus query response, including vector, matrix, scalar, and string results. + """ + + require OpenApiSpex + + alias OpenApiSpex.Schema + + OpenApiSpex.schema(%{ + type: :object, + properties: %{ + status: %Schema{type: :string, enum: ["success", "error"]}, + data: %Schema{ + type: :object, + properties: %{ + resultType: %Schema{type: :string, enum: ["vector", "matrix", "scalar", "string"]}, + result: %Schema{ + type: :array, + items: %Schema{}, + description: + "Native Prometheus values; sample values remain strings, including NaN and infinities." + } + }, + required: [:resultType, :result] + }, + errorType: %Schema{type: :string}, + error: %Schema{type: :string}, + warnings: %Schema{type: :array, items: %Schema{type: :string}}, + infos: %Schema{type: :array, items: %Schema{type: :string}} + }, + required: [:status] + }) +end diff --git a/test/logflare/backends/adaptor/clickhouse_adaptor/finch_pool_timeout_normalizer_test.exs b/test/logflare/backends/adaptor/http_based/finch_pool_timeout_normalizer_test.exs similarity index 95% rename from test/logflare/backends/adaptor/clickhouse_adaptor/finch_pool_timeout_normalizer_test.exs rename to test/logflare/backends/adaptor/http_based/finch_pool_timeout_normalizer_test.exs index e9d1204791..ccbf4c9ceb 100644 --- a/test/logflare/backends/adaptor/clickhouse_adaptor/finch_pool_timeout_normalizer_test.exs +++ b/test/logflare/backends/adaptor/http_based/finch_pool_timeout_normalizer_test.exs @@ -1,7 +1,7 @@ -defmodule Logflare.Backends.Adaptor.ClickHouseAdaptor.FinchPoolTimeoutNormalizerTest do +defmodule Logflare.Backends.Adaptor.HttpBased.FinchPoolTimeoutNormalizerTest do use ExUnit.Case, async: false - alias Logflare.Backends.Adaptor.ClickHouseAdaptor.FinchPoolTimeoutNormalizer + alias Logflare.Backends.Adaptor.HttpBased.FinchPoolTimeoutNormalizer @pool __MODULE__.Pool diff --git a/test/logflare/backends/adaptor/victoria_metrics_adaptor/query_test.exs b/test/logflare/backends/adaptor/victoria_metrics_adaptor/query_test.exs new file mode 100644 index 0000000000..4e3cf2c87d --- /dev/null +++ b/test/logflare/backends/adaptor/victoria_metrics_adaptor/query_test.exs @@ -0,0 +1,273 @@ +defmodule Logflare.Backends.Adaptor.VictoriaMetricsAdaptor.QueryTest do + use ExUnit.Case, async: true + use Mimic + + alias Logflare.Backends.Adaptor.VictoriaMetricsAdaptor.Query + alias Logflare.Backends.Backend + alias Logflare.Utils.SSRF + + @pool __MODULE__.Pool + @url "https://8.8.8.8/select/42/prometheus" + + setup :verify_on_exit! + + test "sends raw PromQL to the explicit read endpoint with basic auth only" do + query = ~S'sum(rate(requests_total{route="/a+b",env=~"prod|staging"}[5m]))' + result = success("vector", []) + + expect(Tesla.Adapter.Finch, :call, fn env, options -> + assert env.method == :post + assert env.url == @url <> "/api/v1/query" + + assert URI.decode_query(env.body) == %{ + "query" => query, + "time" => "2026-10-08T12:00:00Z", + "timeout" => "10s" + } + + assert Map.new(env.headers)["authorization"] == + "Basic " <> Base.encode64("user:secret") + + assert Map.new(env.headers)["content-type"] == "application/x-www-form-urlencoded" + refute Map.has_key?(Map.new(env.headers), "x-write-only") + refute Map.has_key?(Map.new(env.headers), "content-encoding") + assert options[:receive_timeout] == 30_000 + assert options[:name] == Logflare.FinchDefaultHttp1 + reply(env, 200, result) + end) + + backend = + backend(%{ + username: "user", + password: "secret", + headers: %{"authorization" => "Bearer ignored", "x-write-only" => "ignored"} + }) + + assert {:ok, ^result} = + Query.execute(backend, query, %{ + "time" => "2026-10-08T12:00:00Z", + "timeout" => "10s", + "backend_id" => 42, + "query_url" => "https://other.example.com" + }) + end + + test "range requests preserve proxy paths and numeric timestamps" do + result = success("matrix", [%{"metric" => %{"job" => "api"}, "values" => [[10, "2"]]}]) + + expect(Tesla.Adapter.Finch, :call, fn env, _options -> + assert env.url == @url <> "/api/v1/query_range" + + assert URI.decode_query(env.body) == %{ + "query" => "up", + "start" => "10", + "end" => "20.5", + "step" => "5s" + } + + refute Map.has_key?(Map.new(env.headers), "authorization") + reply(env, 200, result) + end) + + assert {:ok, ^result} = + Query.execute(backend(%{query_url: @url <> "/"}), "up", %{ + "start" => 10, + "end" => 20.5, + "step" => "5s" + }) + end + + test "preserves result types, special values, warnings, and informational messages" do + for {type, data} <- [ + {"vector", [%{"metric" => %{"env" => "prod"}, "value" => [123.5, "NaN"]}]}, + {"matrix", [%{"metric" => %{}, "values" => [[123.5, "+Inf"], [124.5, "-Inf"]]}]}, + {"scalar", [123.5, "42"]}, + {"string", [123.5, "hello"]} + ] do + result = Map.merge(success(type, data), %{"warnings" => ["warning"], "infos" => ["info"]}) + expect(Tesla.Adapter.Finch, :call, fn env, _options -> reply(env, 200, result) end) + assert {:ok, ^result} = Query.execute(backend(), "up", %{}) + end + end + + test "preserves native backend query errors and their HTTP status" do + for status <- [400, 422, 503] do + body = %{"status" => "error", "errorType" => "bad_data", "error" => "invalid expression"} + expect(Tesla.Adapter.Finch, :call, fn env, _options -> reply(env, status, body) end) + assert {:error, {^status, ^body}} = Query.execute(backend(), "invalid(", %{}) + end + end + + test "normalizes non-JSON auth failures without returning response bodies" do + for status <- [401, 403], headers <- [[], [{"content-type", "application/json"}]] do + expect(Tesla.Adapter.Finch, :call, fn env, _ -> + {:ok, %{env | status: status, body: "private upstream details", headers: headers}} + end) + + assert {:error, {^status, %{"errorType" => "unauthorized"} = error}} = + Query.execute(backend(), "up", %{}) + + refute inspect(error) =~ "private upstream details" + end + end + + test "does not follow redirects or treat unexpected bodies as successful queries" do + for {status, body} <- [{302, "redirect"}, {200, %{"unexpected" => "body"}}, {204, ""}] do + expect(Tesla.Adapter.Finch, :call, fn env, _ -> + {:ok, + %{env | status: status, body: body, headers: [{"location", "https://other.example.com"}]}} + end) + + assert {:error, {502, %{"errorType" => "bad_response"}}} = + Query.execute(backend(), "up", %{}) + end + end + + test "maps timeouts and transport failures to gateway errors without leaking credentials" do + expect(Tesla.Adapter.Finch, :call, fn _env, _ -> {:error, :timeout} end) + assert {:error, {504, %{"errorType" => "timeout"}}} = Query.execute(backend(), "up", %{}) + + expect(Tesla.Adapter.Finch, :call, fn _env, _ -> {:error, {:connection_failed, "secret"}} end) + assert {:error, {502, body}} = Query.execute(backend(), "up", %{}) + refute inspect(body) =~ "secret" + end + + test "returns a native capacity error when the HTTP/1 connection pool is exhausted" do + start_supervised!( + {Finch, name: @pool, pools: %{default: [protocols: [:http1], size: 1, count: 1]}} + ) + + {:ok, listener} = :gen_tcp.listen(0, [:binary, active: false, reuseaddr: true]) + {:ok, port} = :inet.port(listener) + parent = self() + + server = + Task.async(fn -> + {:ok, socket} = :gen_tcp.accept(listener) + {:ok, _request} = :gen_tcp.recv(socket, 0, 5_000) + send(parent, :pool_busy) + + receive do + :release -> :gen_tcp.close(socket) + end + end) + + holder = + Task.async(fn -> + request = Mimic.call_original(Finch, :build, [:get, "http://127.0.0.1:#{port}/", [], nil]) + Mimic.call_original(Finch, :request, [request, @pool, [receive_timeout: 5_000]]) + end) + + try do + assert_receive :pool_busy, 1_000 + + stub(SSRF, :safe_resolve, fn "victoriametrics.test" -> {:ok, {127, 0, 0, 1}} end) + + stub(Finch, :build, fn method, url, headers, body -> + Mimic.call_original(Finch, :build, [method, url, headers, body]) + end) + + stub(Finch, :request, fn request, pool, options -> + Mimic.call_original(Finch, :request, [request, pool, options]) + end) + + expect(Tesla.Adapter.Finch, :call, fn env, options -> + options = Keyword.merge(options, name: @pool, pool_timeout: 50) + Mimic.call_original(Tesla.Adapter.Finch, :call, [env, options]) + end) + + assert {:error, + {503, + %{ + "status" => "error", + "errorType" => "unavailable", + "error" => "VictoriaMetrics query capacity is temporarily unavailable" + }}} = + Query.execute( + backend(%{query_url: "http://victoriametrics.test:#{port}"}), + "up", + %{} + ) + after + send(server.pid, :release) + Task.shutdown(holder, :brutal_kill) + Task.shutdown(server, :brutal_kill) + :gen_tcp.close(listener) + end + end + + test "does not normalize unrelated runtime errors" do + expect(Tesla.Adapter.Finch, :call, fn _env, _options -> raise "unrelated boom" end) + + assert_raise RuntimeError, "unrelated boom", fn -> + Query.execute(backend(), "up", %{}) + end + end + + test "requires an explicitly configured safe read URL before sending requests" do + reject(Tesla.Adapter.Finch, :call, 2) + + for url <- [ + nil, + "", + "ftp://8.8.8.8", + "https://user:secret@8.8.8.8", + "https://8.8.8.8?token=secret", + "https://8.8.8.8/#fragment", + "https://8.8.8.8:99999", + "http://127.0.0.1:8428", + "http://169.254.169.254" + ] do + assert {:error, {400, %{"status" => "error"}}} = + Query.execute(backend(%{query_url: url}), "up", %{}) + end + end + + test "rechecks the resolved destination before the transport call" do + expect(SSRF, :safe_resolve, fn "8.8.8.8" -> {:ok, {8, 8, 8, 8}} end) + expect(SSRF, :safe_resolve, fn "8.8.8.8" -> {:error, "private destination"} end) + reject(Tesla.Adapter.Finch, :call, 2) + assert {:error, {502, _body}} = Query.execute(backend(), "up", %{}) + end + + test "rejects malformed and ambiguous request parameters before sending requests" do + reject(Tesla.Adapter.Finch, :call, 2) + + for query <- [nil, "", " \n ", ["up"], 1, <<255>>] do + assert {:error, {400, _body}} = Query.execute(backend(), query, %{}) + end + + for params <- [ + nil, + %{"time" => []}, + %{"time" => nil}, + %{"timeout" => ""}, + %{"start" => 1, "end" => 2}, + %{"step" => "5s"}, + %{"time" => 1, "start" => 1, "end" => 2, "step" => 1} + ] do + assert {:error, {400, _body}} = Query.execute(backend(), "up", params) + end + end + + @spec backend(map()) :: Backend.t() + defp backend(config \\ %{}), + do: %Backend{ + config: Map.merge(%{url: "https://1.1.1.1/api/v1/write", query_url: @url}, config) + } + + @spec success(String.t(), term()) :: map() + defp success(type, result), + do: %{"status" => "success", "data" => %{"resultType" => type, "result" => result}} + + @spec reply(Tesla.Env.t(), pos_integer(), map()) :: Tesla.Env.result() + defp reply(env, status, body), + do: + {:ok, + %{ + env + | status: status, + body: Jason.encode!(body), + headers: [{"content-type", "application/json"}] + }} +end diff --git a/test/logflare/backends/adaptor/victoria_metrics_adaptor/series_identity_test.exs b/test/logflare/backends/adaptor/victoria_metrics_adaptor/series_identity_test.exs new file mode 100644 index 0000000000..dcc40c805e --- /dev/null +++ b/test/logflare/backends/adaptor/victoria_metrics_adaptor/series_identity_test.exs @@ -0,0 +1,184 @@ +defmodule Logflare.Backends.Adaptor.VictoriaMetricsAdaptor.SeriesIdentityTest do + use ExUnit.Case, async: true + + alias Logflare.Backends.Adaptor.VictoriaMetricsAdaptor + alias Logflare.Backends.Adaptor.VictoriaMetricsAdaptor.SeriesIdentity + alias Logflare.LogEvent + alias Logflare.TestUtils + alias Prometheus.WriteRequest + + test "keeps resource and scope variants in separate remote write series" do + body = metric_body() + + bodies = [ + body, + put_in(body, ["resource", "deployment_environment"], "staging"), + put_in(body, ["scope", "name"], "library.b"), + put_in(body, ["scope", "version"], "2"), + put_in(body, ["scope", "attributes", "region"], "eu") + ] + + series = + bodies + |> Enum.with_index(1) + |> Enum.map(fn {body, value} -> event(Map.put(body, "value", value)) end) + |> decode() + + assert length(series) == 5 + assert Enum.sort(for %{samples: [%{value: value}]} <- series, do: value) == [1, 2, 3, 4, 5] + + labels = Enum.map(series, &label_map/1) + assert Enum.all?(labels, &(&1["job"] == "api" and &1["instance"] == "one")) + assert labels |> Enum.map(& &1["logflare_resource_id"]) |> Enum.uniq() |> length() == 2 + assert labels |> Enum.map(& &1["logflare_scope_id"]) |> Enum.uniq() |> length() == 4 + end + + test "identical context still batches samples into one series" do + body = metric_body() + later = body |> Map.put("value", 2.0) |> Map.update!("timestamp", &(&1 + 1_000_000)) + + assert [%{samples: [first, second]}] = decode([event(later), event(body)]) + assert first.value == 1.0 + assert second.value == 2.0 + assert second.timestamp - first.timestamp == 1_000 + end + + test "context IDs preserve composite values, scalar types, and array order" do + values = [1, 1.0, "1", true, "true", nil, [], %{}, [1, 2], [2, 1], %{"a" => 1}, [["a", 1]]] + + for {context, label} <- [{"resource", "logflare_resource_id"}, {"scope", "logflare_scope_id"}] do + ids = + Enum.map(values, fn value -> + SeriesIdentity.labels(%{context => %{"attribute" => value}})[label] + end) + + assert MapSet.size(MapSet.new(ids)) == length(values) + assert Enum.all?(ids, &Regex.match?(~r/\A[a-f0-9]{64}\z/, &1)) + end + end + + test "context IDs are independent of map insertion order, including nested maps" do + entries = for n <- 1..40, do: {"key_#{n}", n} + first = Map.new(entries) |> Map.put("nested", %{"x" => 1, "y" => [2, 3]}) + + second = + entries |> Enum.reverse() |> Map.new() |> Map.put("nested", %{"y" => [2, 3], "x" => 1}) + + assert SeriesIdentity.labels(%{"resource" => first, "scope" => first}) == + SeriesIdentity.labels(%{"resource" => second, "scope" => second}) + end + + test "context IDs use a stable canonical encoding" do + assert SeriesIdentity.labels(%{"resource" => %{"env" => "prod"}}) == %{ + "logflare_resource_id" => + "636dae545c02c33adfe04be8199fb1ef0f76ccfe546c757333c6ab37925b22b6" + } + end + + test "context IDs preserve distinctions lost by label sanitization and stringification" do + contexts = [ + %{"a.b" => "x"}, + %{"a_b" => "x"}, + %{"name" => "library", "attributes" => %{"name" => "first"}}, + %{"name" => "library", "attributes" => %{"name" => "second"}}, + %{"value" => :nan}, + %{"value" => "nan"}, + %{"value" => <<255>>} + ] + + ids = Enum.map(contexts, &SeriesIdentity.labels(%{"scope" => &1})["logflare_scope_id"]) + assert MapSet.size(MapSet.new(ids)) == length(contexts) + end + + test "generated identity and scope metadata cannot be overridden by configured labels" do + body = put_in(metric_body(), ["scope", "schema_url"], "https://opentelemetry.io/schemas/1.0") + generated = SeriesIdentity.labels(body) + configured = Map.new(generated, fn {name, _value} -> {name, "configured"} end) + attributes = Map.new(generated, fn {name, _value} -> {name, "point"} end) + attributes = Map.put(attributes, "exported_logflare_scope_id", "already_exported") + + assert [series] = + decode([event(Map.put(body, "attributes", attributes))], %{labels: configured}) + + labels = label_map(series) + assert Map.take(labels, Map.keys(generated)) == generated + assert labels["otel_scope_name"] == "library.a" + assert labels["otel_scope_version"] == "1" + assert labels["otel_scope_schema_url"] == "https://opentelemetry.io/schemas/1.0" + assert labels["exported_logflare_scope_id"] == "already_exported" + assert labels["exported_exported_logflare_scope_id"] == "configured" + assert labels["exported_otel_scope_name"] == "configured" + assert labels["exported_logflare_resource_id"] == "configured" + end + + test "absent or malformed outer contexts do not prevent a mixed batch from exporting" do + for context <- [nil, %{}, [], "invalid", false, 1] do + assert SeriesIdentity.labels(%{"resource" => context, "scope" => context}) == %{} + end + + malformed = metric_body() |> Map.put("resource", "invalid") |> Map.put("scope", ["invalid"]) + assert length(decode([event(malformed), event(metric_body())])) == 2 + end + + test "rejects unsupported nested values and map keys without encoding them" do + unsupported = [self(), make_ref(), fn -> :value end, <<1::1>>, [1 | 2]] + + for context <- ["resource", "scope"], value <- unsupported do + assert :error = + SeriesIdentity.labels(%{context => %{"attributes" => %{"nested" => [value]}}}) + + assert :error = SeriesIdentity.labels(%{context => %{"attributes" => %{value => "value"}}}) + assert :error = SeriesIdentity.labels(%{context => %{"attributes" => {"nested", value}}}) + end + end + + @tag capture_log: true + test "drops malformed context events without losing valid samples in the batch" do + telemetry_event = [:logflare, :backends, :victoria_metrics, :drop] + TestUtils.attach_forwarder(telemetry_event) + backend_id = System.unique_integer([:positive]) + valid = metric_body() + + invalid_resource = put_in(valid, ["resource", "nested"], %{"value" => self()}) + invalid_scope = put_in(valid, ["scope", "attributes", "nested"], [1 | 2]) + + assert [%{samples: [%{value: 1.0}]}] = + decode( + [event(invalid_resource), event(valid), event(invalid_scope)], + %{backend_id: backend_id} + ) + + assert_received {:telemetry_event, ^telemetry_event, %{count: 2}, + %{reason: :invalid, backend_id: ^backend_id}} + end + + @spec metric_body() :: map() + defp metric_body do + %{ + "event_message" => "requests", + "metric_type" => "gauge", + "value" => 1.0, + "timestamp" => 1_700_000_000_000_000, + "attributes" => %{}, + "resource" => %{ + "service.name" => "api", + "service.instance.id" => "one", + "deployment_environment" => "production" + }, + "scope" => %{"name" => "library.a", "version" => "1", "attributes" => %{"region" => "us"}} + } + end + + @spec event(map()) :: LogEvent.t() + defp event(body), do: %LogEvent{source_id: nil, event_type: :metric, body: body} + + @spec decode([LogEvent.t()], map()) :: [Prometheus.TimeSeries.t()] + defp decode(events, config \\ %{}) do + payload = VictoriaMetricsAdaptor.format_batch(events, config) + {:ok, protobuf} = :snappyer.decompress(payload) + WriteRequest.decode(protobuf).timeseries + end + + @spec label_map(Prometheus.TimeSeries.t()) :: map() + defp label_map(series), do: Map.new(series.labels, &{&1.name, &1.value}) +end diff --git a/test/logflare/backends/adaptor/victoria_metrics_adaptor_test.exs b/test/logflare/backends/adaptor/victoria_metrics_adaptor_test.exs index 61962c0a42..580740e8d4 100644 --- a/test/logflare/backends/adaptor/victoria_metrics_adaptor_test.exs +++ b/test/logflare/backends/adaptor/victoria_metrics_adaptor_test.exs @@ -65,6 +65,123 @@ defmodule Logflare.Backends.Adaptor.VictoriaMetricsAdaptorTest do }).valid? end + test "accepts an optional query base URL and allows clearing it" do + existing = %{ + url: @vm_remote_write_url, + query_url: "https://vmselect.example.com/select/0/prometheus" + } + + assert Adaptor.cast_and_validate_config(@subject, existing).valid? + + for query_url <- [nil, ""] do + changeset = Adaptor.cast_and_validate_config(@subject, %{query_url: query_url}, existing) + assert changeset.valid? + assert Ecto.Changeset.get_field(changeset, :query_url) == nil + end + end + + test "rejects query base URLs with invalid transport or embedded request options" do + for query_url <- [ + "not-a-url", + "ftp://vm.example.com", + "https:///api", + "https://vm.example.com:0", + "https://vm.example.com:65536", + "https://user:password@vm.example.com", + "https://vm.example.com?query=up", + "https://vm.example.com#fragment" + ] do + changeset = + Adaptor.cast_and_validate_config(@subject, %{ + url: @vm_remote_write_url, + query_url: query_url + }) + + refute changeset.valid?, query_url + assert changeset.errors[:query_url] + end + end + + test "returns a validation error when editing the query URL to a malformed port" do + existing = %{ + url: "https://8.8.8.8/api/v1/write", + query_url: "https://8.8.8.8" + } + + changeset = + Adaptor.cast_and_validate_config( + @subject, + %{query_url: "https://8.8.8.8:notaport"}, + existing + ) + + refute changeset.valid? + assert changeset.errors[:query_url] + end + + test "requires password reentry for a new read origin while keeping write-only headers" do + for previous_query_url <- [nil, "https://reader.example.com/select/0/prometheus"] do + existing = %{ + url: "https://vm.example.com/api/v1/write", + query_url: previous_query_url, + username: "user", + password: "pass", + headers: %{"authorization" => "Bearer write-token"} + } + + params = + existing + |> @subject.redact_config() + |> Map.put(:query_url, "https://other-reader.example.com/select/0/prometheus") + + changeset = Adaptor.cast_and_validate_config(@subject, params, existing) + refute changeset.valid? + assert changeset.errors[:password] + assert Ecto.Changeset.get_field(changeset, :password) == nil + assert Ecto.Changeset.get_field(changeset, :headers) == existing.headers + + reentered = + Adaptor.cast_and_validate_config(@subject, Map.put(params, :password, "pass"), existing) + + assert reentered.valid? + assert Ecto.Changeset.get_field(reentered, :password) == "pass" + end + end + + test "keeps shared credentials for same-origin read paths and when disabling reads" do + existing = %{ + url: "https://vm.example.com/api/v1/write", + username: "user", + password: "pass" + } + + for old_url <- [nil, "https://vm.example.com/select/0/prometheus"], + new_url <- [nil, "https://vm.example.com/select/1/prometheus"] do + stored = Map.put(existing, :query_url, old_url) + params = stored |> @subject.redact_config() |> Map.put(:query_url, new_url) + changeset = Adaptor.cast_and_validate_config(@subject, params, stored) + + assert changeset.valid? + assert Ecto.Changeset.get_field(changeset, :password) == "pass" + assert Ecto.Changeset.get_field(changeset, :query_url) == new_url + end + end + + test "treats query URL scheme and port changes as a new destination" do + existing = %{ + url: "https://vm.example.com/api/v1/write", + query_url: "https://reader.example.com/select/0/prometheus", + username: "user", + password: "pass" + } + + for query_url <- ["http://reader.example.com", "https://reader.example.com:8443"] do + changeset = Adaptor.cast_and_validate_config(@subject, %{query_url: query_url}, existing) + refute changeset.valid? + assert changeset.errors[:password] + end + end + test "keeps stored credentials submitted back in redacted form" do existing = %{ url: "https://user:secret@vm.example.com/api/v1/write", @@ -132,6 +249,49 @@ defmodule Logflare.Backends.Adaptor.VictoriaMetricsAdaptorTest do assert Ecto.Changeset.apply_changes(reentered).password == "new" end + test "keeps explicitly reentered identical credentials when moving the write destination" do + existing = %{ + url: "https://vm.example.com/api/v1/write", + headers: %{"authorization" => "Bearer token", "x-tenant" => "acme"} + } + + changeset = + Adaptor.cast_and_validate_config( + @subject, + %{existing | url: "https://other.example.com/api/v1/write"}, + existing + ) + + assert changeset.valid? + assert Ecto.Changeset.get_field(changeset, :headers) == existing.headers + end + + test "drops masked credentials in mixed header edits only when the write origin changes" do + existing = %{ + url: "https://vm.example.com/api/v1/write", + headers: %{"authorization" => "Bearer token", "x-tenant" => "acme"} + } + + for {url, expected_headers} <- [ + {"https://other.example.com/api/v1/write", %{"x-tenant" => "updated"}}, + {"https://vm.example.com/api/v1/write", + %{"authorization" => "Bearer token", "x-tenant" => "updated"}} + ] do + changeset = + Adaptor.cast_and_validate_config( + @subject, + %{ + url: url, + headers: %{"Authorization" => "REDACTED", "x-tenant" => "updated"} + }, + existing + ) + + assert changeset.valid? + assert Ecto.Changeset.get_field(changeset, :headers) == expected_headers + end + end + test "keeps credentials across a stored round trip through the backend API" do user = insert(:user) @@ -169,6 +329,23 @@ defmodule Logflare.Backends.Adaptor.VictoriaMetricsAdaptorTest do refute changeset.valid? assert {_msg, [validation: :ssrf]} = changeset.errors[:url] end + + test "rejects a private read destination independently of the write URL" do + stub(Logflare.Utils.SSRF, :safe_resolve, fn + "vm.example.com" -> {:ok, {1, 2, 3, 4}} + "127.0.0.1" -> {:error, "private address is not allowed"} + end) + + changeset = + Adaptor.cast_and_validate_config(@subject, %{ + url: "https://vm.example.com/api/v1/write", + query_url: "http://127.0.0.1:8428" + }) + + refute changeset.valid? + assert changeset.errors[:query_url] + refute changeset.errors[:url] + end end describe "redact_config/1" do @@ -694,22 +871,28 @@ defmodule Logflare.Backends.Adaptor.VictoriaMetricsAdaptorTest do end end - # End-to-end tests against the docker-compose `vm` service. - # - # Requires `docker compose up -d vm`. Excluded by default via the - # :integration tag (see test/test_helper.exs). - # - # Run with: - # mix test test/logflare/backends/adaptor/victoria_metrics_adaptor_test.exs --include integration describe "victoriametrics e2e" do - @describetag :integration - # The vm service is on loopback, which SSRFProtection blocks. The pipeline sends # from its own processes, so the stub has to be global. setup :set_mimic_global setup do - stub(Logflare.Utils.SSRF, :safe_resolve, fn _ -> {:ok, {127, 0, 0, 1}} end) + stub(Logflare.Utils.SSRF, :safe_resolve, fn + "localhost" -> {:ok, {127, 0, 0, 1}} + host -> Mimic.call_original(Logflare.Utils.SSRF, :safe_resolve, [host]) + end) + + stub(Finch, :build, fn method, url, headers, body -> + Mimic.call_original(Finch, :build, [method, url, headers, body]) + end) + + stub(Finch, :request, fn + %Finch.Request{host: "127.0.0.1", port: 8428} = request, pool, opts -> + Mimic.call_original(Finch, :request, [request, pool, opts]) + + _request, _pool, _opts -> + {:error, :unexpected_external_request} + end) insert(:plan) user = insert(:user) diff --git a/test/logflare/backends/adaptor/victoria_metrics_query_integration_test.exs b/test/logflare/backends/adaptor/victoria_metrics_query_integration_test.exs new file mode 100644 index 0000000000..ca8cf3755a --- /dev/null +++ b/test/logflare/backends/adaptor/victoria_metrics_query_integration_test.exs @@ -0,0 +1,310 @@ +defmodule Logflare.Backends.Adaptor.VictoriaMetricsQueryIntegrationTest do + use LogflareWeb.ConnCase, async: false + + alias Logflare.Backends + alias Logflare.Backends.Adaptor.VictoriaMetricsAdaptor + alias Logflare.Backends.Adaptor.VictoriaMetricsAdaptor.SeriesIdentity + alias Logflare.Backends.AdaptorSupervisor + alias Logflare.LogEvent + alias Logflare.Logs.OtelMetric + alias Logflare.Sources.Source + alias Logflare.SystemMetrics.AllLogsLogged + alias Logflare.Utils.SSRF + alias Opentelemetry.Proto.Common.V1.AnyValue + alias Opentelemetry.Proto.Common.V1.ArrayValue + alias Opentelemetry.Proto.Common.V1.InstrumentationScope + alias Opentelemetry.Proto.Common.V1.KeyValue + alias Opentelemetry.Proto.Metrics.V1.Gauge + alias Opentelemetry.Proto.Metrics.V1.Metric + alias Opentelemetry.Proto.Metrics.V1.NumberDataPoint + alias Opentelemetry.Proto.Metrics.V1.ResourceMetrics + alias Opentelemetry.Proto.Metrics.V1.ScopeMetrics + alias Opentelemetry.Proto.Resource.V1.Resource + + @receiver_url "http://victoriametrics.test:8428" + @retry [sleep: 250, duration: 30_000] + + setup do + stub(SSRF, :safe_resolve, fn + "victoriametrics.test" -> {:ok, {127, 0, 0, 1}} + host -> Mimic.call_original(SSRF, :safe_resolve, [host]) + end) + + stub(Finch, :build, fn method, url, headers, body -> + Mimic.call_original(Finch, :build, [method, url, headers, body]) + end) + + stub(Finch, :request, fn + %Finch.Request{host: "127.0.0.1", port: 8428} = request, pool, opts -> + Mimic.call_original(Finch, :request, [request, pool, opts]) + + _request, _pool, _opts -> + {:error, :unexpected_external_request} + end) + + start_supervised!(AllLogsLogged) + insert(:plan) + user = insert(:user) + source = insert(:source, user: user) + + backend = + insert(:backend, + type: :victoria_metrics, + user: user, + sources: [source], + config: %{ + url: @receiver_url <> "/api/v1/write", + query_url: @receiver_url, + labels: %{"env" => "integration"} + } + ) + + start_supervised!({AdaptorSupervisor, {source, backend}}) + + prefix = "vm_query_" <> String.replace(Ecto.UUID.generate(), "-", "") + %{source: source, backend: backend, prefix: prefix, user: user} + end + + test "queries ingested OTEL metrics as native instant and range results", %{ + source: source, + backend: backend, + prefix: prefix, + user: user, + conn: conn + } do + start = div(System.system_time(:second), 10) * 10 - 120 + finish = start + 20 + samples = [{start, 10.0}, {start + 10, 20.0}, {finish, 30.0}] + metric_name = prefix <> ".temperature" + query_name = prefix <> "_temperature" + + events = metric_events(source, metric_name, samples) + + assert Enum.map(events, & &1.body["timestamp"]) == + Enum.map(samples, &(elem(&1, 0) * 1_000_000)) + + assert {:ok, _count} = Backends.ingest_logs(events, source) + + expected_labels = + events + |> hd() + |> Map.fetch!(:body) + |> SeriesIdentity.labels() + |> Map.take(["logflare_resource_id", "logflare_scope_id"]) + |> Map.merge(%{ + "__name__" => query_name, + "source" => source.name, + "job" => "integration/metrics-api", + "instance" => "replica-a", + "env" => "integration", + "_http_route" => "/metrics", + "otel_scope_name" => "vm.query.integration", + "otel_scope_version" => "1.0" + }) + + TestUtils.retry_assert(@retry, fn -> + assert {:ok, + %{ + "status" => "success", + "data" => %{"resultType" => "vector", "result" => [result]} + }} = VictoriaMetricsAdaptor.execute_promql(backend, query_name, %{"time" => finish}) + + assert result["metric"] == expected_labels + assert result["value"] == [finish, "30"] + end) + + conn = add_access_token(conn, user, ~w(private)) + + instant = + conn + |> get(~p"/api/query", %{promql: query_name, backend_id: backend.id, time: finish}) + |> json_response(200) + + assert %{ + "status" => "success", + "data" => %{"resultType" => "vector", "result" => [instant_result]} + } = instant + + assert instant_result == %{"metric" => expected_labels, "value" => [finish, "30"]} + + range = + conn + |> get(~p"/api/query", %{ + promql: query_name, + backend_id: backend.id, + start: start, + end: finish, + step: "10s" + }) + |> json_response(200) + + assert %{ + "status" => "success", + "data" => %{"resultType" => "matrix", "result" => [range_result]} + } = range + + assert range_result == %{ + "metric" => expected_labels, + "values" => [[start, "10"], [start + 10, "20"], [finish, "30"]] + } + end + + test "keeps resource and scope variants separate at the same timestamp", %{ + source: source, + backend: backend, + prefix: prefix + } do + timestamp = System.system_time(:second) - 120 + metric_name = prefix <> "_identity" + + context = [ + resource_attributes: [attribute("host.name", "host-a")], + scope_attributes: [attribute("build.flags", ["alpha", "beta"])] + ] + + variants = [ + {10.0, []}, + {20.0, resource_attributes: [attribute("host.name", "host-b")]}, + {30.0, scope_name: "vm.query.other"}, + {40.0, scope_version: "2.0"}, + {50.0, scope_attributes: [attribute("build.flags", ["alpha", "gamma"])]}, + {60.0, scope_attributes: [attribute("build.flags", ["beta", "alpha"])]} + ] + + events = + Enum.flat_map(variants, fn {value, overrides} -> + metric_events( + source, + metric_name, + [{timestamp, value}], + Keyword.merge(context, overrides) + ) + end) + + assert {:ok, _count} = Backends.ingest_logs(events, source) + + TestUtils.retry_assert(@retry, fn -> + assert {:ok, + %{ + "status" => "success", + "data" => %{"resultType" => "vector", "result" => results} + }} = + VictoriaMetricsAdaptor.execute_promql(backend, metric_name, %{"time" => timestamp}) + + assert length(results) == length(variants) + + by_value = + Map.new(results, fn %{"metric" => labels, "value" => [sample_time, value]} -> + assert sample_time == timestamp + assert labels["__name__"] == metric_name + assert labels["source"] == source.name + assert labels["job"] == "integration/metrics-api" + assert labels["instance"] == "replica-a" + assert labels["env"] == "integration" + assert labels["_http_route"] == "/metrics" + assert labels["logflare_resource_id"] =~ ~r/\A[0-9a-f]{64}\z/ + assert labels["logflare_scope_id"] =~ ~r/\A[0-9a-f]{64}\z/ + {value, labels} + end) + + assert Enum.sort(Map.keys(by_value)) == ~w(10 20 30 40 50 60) + baseline = by_value["10"] + + assert by_value["20"]["logflare_resource_id"] != baseline["logflare_resource_id"] + assert by_value["20"]["logflare_scope_id"] == baseline["logflare_scope_id"] + assert by_value["30"]["otel_scope_name"] == "vm.query.other" + assert by_value["40"]["otel_scope_version"] == "2.0" + + scope_variants = Enum.map(~w(10 30 40 50 60), &Map.fetch!(by_value, &1)) + + assert scope_variants |> Enum.map(& &1["logflare_scope_id"]) |> Enum.uniq() |> length() == 5 + + assert Enum.all?(scope_variants, fn labels -> + labels["logflare_resource_id"] == baseline["logflare_resource_id"] + end) + + for value <- ~w(10 20 50 60) do + assert by_value[value]["otel_scope_name"] == "vm.query.integration" + assert by_value[value]["otel_scope_version"] == "1.0" + end + end) + end + + test "returns native errors from the real PromQL parser", %{ + backend: backend, + user: user, + conn: conn + } do + response = + conn + |> add_access_token(user, ~w(private)) + |> get(~p"/api/query", %{promql: "{", backend_id: backend.id}) + |> json_response(422) + + assert %{"status" => "error", "errorType" => error_type, "error" => message} = response + assert is_binary(error_type) + assert is_binary(message) + assert message != "" + end + + test "requires API authentication before querying the receiver", %{backend: backend, conn: conn} do + conn = get(conn, ~p"/api/query", %{promql: "up", backend_id: backend.id}) + assert conn.status == 401 + end + + @spec metric_events(Source.t(), String.t(), [{integer(), float()}], keyword()) :: [LogEvent.t()] + defp metric_events(source, name, samples, context \\ []) do + resource_metrics = %ResourceMetrics{ + resource: %Resource{ + attributes: + [ + attribute("service.name", "metrics-api"), + attribute("service.namespace", "integration"), + attribute("service.instance.id", "replica-a") + ] ++ Keyword.get(context, :resource_attributes, []) + }, + scope_metrics: [ + %ScopeMetrics{ + scope: %InstrumentationScope{ + name: Keyword.get(context, :scope_name, "vm.query.integration"), + version: Keyword.get(context, :scope_version, "1.0"), + attributes: Keyword.get(context, :scope_attributes, []) + }, + metrics: [ + %Metric{ + name: name, + data: + {:gauge, + %Gauge{ + data_points: + for {timestamp, value} <- samples do + %NumberDataPoint{ + time_unix_nano: timestamp * 1_000_000_000, + value: {:as_double, value}, + attributes: [ + attribute("env", "event"), + attribute("http.route", "/metrics") + ] + } + end + }} + } + ] + } + ] + } + + [resource_metrics] + |> OtelMetric.handle_batch(source) + |> Enum.map(&LogEvent.make(&1, %{source: source})) + end + + @spec attribute(String.t(), String.t() | [String.t()]) :: KeyValue.t() + defp attribute(key, values) when is_list(values) do + array = %ArrayValue{values: Enum.map(values, &%AnyValue{value: {:string_value, &1}})} + %KeyValue{key: key, value: %AnyValue{value: {:array_value, array}}} + end + + defp attribute(key, value) when is_binary(value), + do: %KeyValue{key: key, value: %AnyValue{value: {:string_value, value}}} +end diff --git a/test/logflare_web/controllers/api/query_controller_test.exs b/test/logflare_web/controllers/api/query_controller_test.exs index c3d8339e8a..d9be4f7fec 100644 --- a/test/logflare_web/controllers/api/query_controller_test.exs +++ b/test/logflare_web/controllers/api/query_controller_test.exs @@ -1,6 +1,7 @@ defmodule LogflareWeb.Api.QueryControllerTest do use LogflareWeb.ConnCase + alias Logflare.Backends.Adaptor alias Logflare.Backends.Adaptor.ClickHouseAdaptor alias Logflare.Backends.Adaptor.PostgresAdaptor alias Logflare.DataCase @@ -351,4 +352,281 @@ defmodule LogflareWeb.Api.QueryControllerTest do assert %{"error" => "Backend does not support querying"} = response end end + + describe "query with promql" do + setup %{user: user} do + backend = + insert(:backend, + user: user, + type: :victoria_metrics, + config: %{ + url: "https://8.8.8.8/api/v1/write", + query_url: "https://8.8.8.8", + username: "query-user", + password: "query:password" + } + ) + + stub(Tesla.Adapter.Finch, :call, fn _env, _opts -> + flunk("Unexpected VictoriaMetrics request") + end) + + %{backend: backend} + end + + test "forwards raw expressions and instant options with stored basic authentication", %{ + conn: conn, + user: user, + backend: backend + } do + query = ~S'sum(rate(http_requests_total{path=~"/a|b",service="café"}[5m])) by (job)' + data = %{"resultType" => "vector", "result" => []} + expected = %{"status" => "success", "data" => data, "warnings" => ["partial data"]} + + expect(Tesla.Adapter.Finch, :call, fn env, _opts -> + assert env.method == :post + assert env.url == "https://8.8.8.8/api/v1/query" + + assert Tesla.get_header(env, "authorization") == + "Basic " <> Base.encode64("query-user:query:password") + + assert URI.decode_query(env.body) == %{ + "query" => query, + "time" => "2026-10-08T01:02:03Z", + "timeout" => "10s" + } + + tesla_response(env, 200, expected) + end) + + assert promql_response(conn, user, %{ + promql: query, + backend_id: backend.id, + time: "2026-10-08T01:02:03Z", + timeout: "10s", + query: "ignored", + username: "ignored", + password: "ignored" + }) == expected + end + + test "range queries preserve timestamps, labels, and special sample values", %{ + conn: conn, + user: user, + backend: backend + } do + expected = %{ + "status" => "success", + "data" => %{ + "resultType" => "matrix", + "result" => [ + %{ + "metric" => %{"__name__" => "temperature", "instance" => "host-1"}, + "values" => [[1_791_421_200.25, "NaN"], [1_791_421_215.25, "+Inf"]] + } + ] + }, + "infos" => ["sample information"] + } + + expect(Tesla.Adapter.Finch, :call, fn env, _opts -> + assert env.url == "https://8.8.8.8/api/v1/query_range" + + assert URI.decode_query(env.body) == %{ + "query" => "temperature", + "start" => "1791421200.25", + "end" => "1791421260.25", + "step" => "15s" + } + + tesla_response(env, 200, expected) + end) + + assert promql_response(conn, user, %{ + promql: "temperature", + backend_id: backend.id, + start: "1791421200.25", + end: "1791421260.25", + step: "15s" + }) == expected + end + + test "preserves vector, scalar, and string response shapes", %{ + conn: conn, + user: user, + backend: backend + } do + for {type, result} <- [ + {"vector", [%{"metric" => %{"job" => "api"}, "value" => [123.5, "-Inf"]}]}, + {"scalar", [123.5, "NaN"]}, + {"string", [123.5, "a string result"]} + ] do + expected = %{"status" => "success", "data" => %{"resultType" => type, "result" => result}} + + expect(Tesla.Adapter.Finch, :call, fn env, _opts -> + tesla_response(env, 200, expected) + end) + + assert promql_response(conn, user, %{promql: "up", backend_id: backend.id}) == expected + end + end + + test "preserves native backend errors and HTTP statuses", %{ + conn: conn, + user: user, + backend: backend + } do + for {status, type} <- [{400, "bad_data"}, {422, "execution"}, {503, "timeout"}] do + expected = %{ + "status" => "error", + "errorType" => type, + "error" => "query failed", + "warnings" => ["backend warning"] + } + + expect(Tesla.Adapter.Finch, :call, fn env, _opts -> + tesla_response(env, status, expected) + end) + + assert promql_response(conn, user, %{promql: "up", backend_id: backend.id}, status) == + expected + end + end + + test "returns a native capacity error when Finch cannot check out a connection", %{ + conn: conn, + user: user, + backend: backend + } do + for options <- [%{}, %{start: "1", end: "2", step: "1s"}] do + expect(Tesla.Adapter.Finch, :call, fn _env, _options -> + raise "Finch was unable to provide a connection within the timeout due to excess queuing for connections" + end) + + params = Map.merge(%{promql: "up", backend_id: backend.id}, options) + + assert promql_response(conn, user, params, 503) == %{ + "status" => "error", + "errorType" => "unavailable", + "error" => "VictoriaMetrics query capacity is temporarily unavailable" + } + end + end + + test "requires an explicit valid backend ID", %{conn: conn, user: user} do + for params <- [%{promql: "up"}, %{promql: "up", backend_id: ""}] do + assert %{ + "status" => "error", + "errorType" => "bad_data", + "error" => "backend_id is required for PromQL queries" + } = + promql_response(conn, user, params, 400) + end + + for backend_id <- ["invalid", "1.5", %{id: "1"}] do + assert %{"error" => "Invalid backend_id: must be an integer"} = + promql_response(conn, user, %{promql: "up", backend_id: backend_id}, 400) + end + + for backend_id <- [0, -1, "9223372036854775808"] do + assert %{"error" => "Invalid backend_id: must be a positive 64-bit integer"} = + promql_response(conn, user, %{promql: "up", backend_id: backend_id}, 400) + end + + assert %{"error" => "Backend not found"} = + promql_response(conn, user, %{promql: "up", backend_id: 999_999}, 400) + end + + test "rejects another user's backend and non-VictoriaMetrics backends", %{ + conn: conn, + user: user + } do + other_backend = insert(:backend, user: insert(:user), type: :victoria_metrics) + sql_backend = insert(:backend, user: user, type: :clickhouse) + + assert %{"error" => "Backend not found"} = + promql_response(conn, user, %{promql: "up", backend_id: other_backend.id}, 400) + + assert %{"error" => "Backend does not support PromQL queries"} = + promql_response(conn, user, %{promql: "up", backend_id: sql_backend.id}, 400) + end + + test "rejects every SQL parameter mixed with PromQL", %{ + conn: conn, + user: user, + backend: backend + } do + for sql_key <- [:sql, :bq_sql, :ch_sql, :pg_sql] do + params = Map.put(%{promql: "up", backend_id: backend.id}, sql_key, "SELECT 1") + + assert %{"error" => "promql cannot be combined with SQL parameters"} = + promql_response(conn, user, params, 400) + end + end + + test "rejects malformed expressions and options before sending a request", %{ + conn: conn, + user: user, + backend: backend + } do + for invalid <- [ + %{promql: " "}, + %{promql: %{query: "up"}}, + %{time: %{invalid: "value"}}, + %{start: "1"}, + %{start: "1", end: "2", step: "1s", time: "1"} + ] do + params = Map.merge(%{promql: "up", backend_id: backend.id}, invalid) + + assert %{"status" => "error", "errorType" => "bad_data", "error" => error} = + promql_response(conn, user, params, 400) + + assert is_binary(error) + end + end + + test "VictoriaMetrics remains unavailable to SQL execution and parsing", %{ + conn: conn, + user: user, + backend: backend + } do + refute Adaptor.can_query?(backend) + params = %{sql: "SELECT 1", backend_id: backend.id} + + assert %{"error" => "Backend does not support querying"} = + promql_response(conn, user, params, 400) + + response = + conn + |> add_access_token(user, ~w(private)) + |> get(~p"/api/query/parse", params) + |> json_response(400) + + assert %{"error" => "Backend does not support querying"} = response + end + + test "requires management authentication", %{conn: conn, backend: backend} do + conn = get(conn, ~p"/api/query", %{promql: "up", backend_id: backend.id}) + assert %{"error" => _error} = json_response(conn, 401) + end + end + + @spec promql_response(Plug.Conn.t(), Logflare.User.t(), map(), pos_integer()) :: map() + defp promql_response(conn, user, params, status \\ 200) do + conn + |> add_access_token(user, ~w(private)) + |> get(~p"/api/query", params) + |> json_response(status) + end + + @spec tesla_response(Tesla.Env.t(), pos_integer(), map()) :: Tesla.Env.result() + defp tesla_response(%Tesla.Env{} = env, status, body) do + {:ok, + %Tesla.Env{ + env + | status: status, + headers: [{"content-type", "application/json"}], + body: Jason.encode!(body) + }} + end end diff --git a/test/logflare_web/live/backends/backends_live_test.exs b/test/logflare_web/live/backends/backends_live_test.exs index c60014434b..112916a366 100644 --- a/test/logflare_web/live/backends/backends_live_test.exs +++ b/test/logflare_web/live/backends/backends_live_test.exs @@ -604,6 +604,7 @@ defmodule LogflareWeb.BackendsLiveTest do type: "victoria_metrics", config: %{ url: "https://example.com/api/v1/write", + query_url: "https://example.com/select/0/prometheus", username: "user", password: "pass" } @@ -616,6 +617,7 @@ defmodule LogflareWeb.BackendsLiveTest do [backend] = Backends.list_backends_by_user_access(user, type: :victoria_metrics) assert backend.config.url == "https://example.com/api/v1/write" + assert backend.config.query_url == "https://example.com/select/0/prometheus" assert backend.config.username == "user" end @@ -839,6 +841,115 @@ defmodule LogflareWeb.BackendsLiveTest do assert html =~ "some description" end + test "victoria_metrics edit requires password reentry for a new query destination", %{ + conn: conn, + source: source, + user: user + } do + backend = + insert(:backend, + sources: [source], + user: user, + type: :victoria_metrics, + config: %{ + url: "https://example.com/api/v1/write", + query_url: "https://example.com/select/0/prometheus", + username: "user", + password: "vm-secret", + headers: %{"authorization" => "Bearer write-token"} + } + ) + + {:ok, view, html} = live_with_redirect(conn, ~p"/backends/#{backend.id}/edit") + + refute html =~ "vm-secret" + refute html =~ "write-token" + + assert view + |> element("input[name='backend[config][query_url]']") + |> render() =~ "https://example.com/select/0/prometheus" + + html = + view + |> form("form", %{ + backend: %{config: %{query_url: "https://example.org/select/0/prometheus"}} + }) + |> render_submit() + + assert html =~ "Both username and password must be provided for basic auth" + assert Backends.get_backend(backend.id).config.query_url == backend.config.query_url + + view + |> form("form", %{ + backend: %{ + config: %{ + query_url: "https://example.org/select/0/prometheus", + password: "vm-secret" + } + } + }) + |> render_submit() + + updated = Backends.get_backend(backend.id) + assert updated.config.query_url == "https://example.org/select/0/prometheus" + assert updated.config.password == "vm-secret" + assert updated.config.headers == %{"authorization" => "Bearer write-token"} + end + + test "victoria_metrics mixed header edits require token reentry for a new write origin", %{ + conn: conn, + source: source, + user: user + } do + for {url, header_key, token, expected_headers} <- [ + {"https://example.org/api/v1/write", "authorization", "REDACTED", + %{"x-tenant" => "updated"}}, + {"http://example.com/api/v1/write", "authorization", "REDACTED", + %{"x-tenant" => "updated"}}, + {"https://example.com:8443/api/v1/write", "authorization", "REDACTED", + %{"x-tenant" => "updated"}}, + {"https://example.org/api/v1/write", "x-api-key", "REDACTED", + %{"x-tenant" => "updated"}}, + {"https://example.org/api/v1/write", "authorization", "Bearer write-token", + %{"authorization" => "Bearer write-token", "x-tenant" => "updated"}}, + {"https://example.com/api/v1/write", "authorization", "REDACTED", + %{"authorization" => "Bearer write-token", "x-tenant" => "updated"}} + ] do + backend = + insert(:backend, + sources: [source], + user: user, + type: :victoria_metrics, + config: %{ + url: "https://example.com/api/v1/write", + headers: %{"authorization" => "Bearer write-token", "x-tenant" => "acme"} + } + ) + + {:ok, view, html} = live_with_redirect(conn, ~p"/backends/#{backend.id}/edit") + refute html =~ "write-token" + + html = + view + |> form("form", %{ + backend: %{ + config: %{ + url: url, + header1_key: header_key, + header1_value: token, + header2_value: "updated" + } + } + }) + |> render_submit() + + assert html =~ "Successfully updated backend" + updated = Backends.get_backend(backend.id) + assert updated.config.url == url + assert updated.config.headers == expected_headers + end + end + test "webhook edit renders stored headers and preserves them on submit", %{ conn: conn, source: source, diff --git a/test/logflare_web/open_api_test.exs b/test/logflare_web/open_api_test.exs index a7aa44def5..cb11c3aa3a 100644 --- a/test/logflare_web/open_api_test.exs +++ b/test/logflare_web/open_api_test.exs @@ -7,6 +7,7 @@ defmodule LogflareWeb.OpenApiTest do alias LogflareWeb.ApiSpec alias LogflareWeb.OpenApiSchemas.AccessToken alias LogflareWeb.OpenApiSchemas.ClickhouseConfigSchema + alias LogflareWeb.OpenApiSchemas.PromQLQueryResponse alias LogflareWeb.OpenApiSchemas.QueryResult alias OpenApiSpex.MediaType alias OpenApiSpex.Response @@ -20,14 +21,18 @@ defmodule LogflareWeb.OpenApiTest do } @string_error_schema %Schema{type: :string} - test "Management API query success is documented as an object containing result rows" do + test "Management API query success documents SQL rows and native PromQL responses" do response = QueryController.open_api_operation(:query) |> Map.fetch!(:responses) |> Map.fetch!(200) assert %Response{ - content: %{"application/json" => %MediaType{schema: QueryResult}} + content: %{ + "application/json" => %MediaType{ + schema: %Schema{oneOf: [QueryResult, PromQLQueryResponse]} + } + } } = response assert %Schema{ @@ -35,6 +40,22 @@ defmodule LogflareWeb.OpenApiTest do properties: %{result: %Schema{type: :array, items: %Schema{type: :object}}}, required: [:result] } = QueryResult.schema() + + assert %Schema{ + type: :object, + properties: %{ + status: %Schema{type: :string, enum: ["success", "error"]}, + data: %Schema{ + type: :object, + properties: %{ + resultType: %Schema{enum: ["vector", "matrix", "scalar", "string"]}, + result: %Schema{type: :array} + }, + required: [:resultType, :result] + } + }, + required: [:status] + } = PromQLQueryResponse.schema() end test "Management API access token timestamps are documented as RFC3339" do @@ -76,7 +97,11 @@ defmodule LogflareWeb.OpenApiTest do QueryController.open_api_operation(action) |> Map.fetch!(:responses) |> Map.fetch!(400) - |> assert_json_error_response("BadRequestResponse", @bad_request_error_schema) + |> assert_json_error_response( + "BadRequestResponse", + @bad_request_error_schema, + action == :query + ) end end @@ -90,7 +115,11 @@ defmodule LogflareWeb.OpenApiTest do QueryController.open_api_operation(action) |> Map.fetch!(:responses) |> Map.fetch!(401) - |> assert_json_error_response("UnauthorizedResponse", @string_error_schema) + |> assert_json_error_response( + "UnauthorizedResponse", + @string_error_schema, + action == :query + ) end end @@ -101,18 +130,27 @@ defmodule LogflareWeb.OpenApiTest do |> assert_json_error_response("NotFoundResponse", @string_error_schema) end - defp assert_json_error_response(response, schema_title, error_schema) do + @spec assert_json_error_response(Response.t(), String.t(), Schema.t(), boolean()) :: Schema.t() + defp assert_json_error_response(response, schema_title, error_schema, promql? \\ false) do assert %Response{ content: %{ - "application/json" => %MediaType{ - schema: %Schema{ - title: ^schema_title, - type: :object, - properties: %{error: ^error_schema}, - required: [:error] - } - } + "application/json" => %MediaType{schema: schema} } } = response + + schema = + if promql? do + assert %Schema{anyOf: [json_error_schema, PromQLQueryResponse]} = schema + json_error_schema + else + schema + end + + assert %Schema{ + title: ^schema_title, + type: :object, + properties: %{error: ^error_schema}, + required: [:error] + } = schema end end