Skip to content
Draft
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
9 changes: 9 additions & 0 deletions .github/workflows/elixir-ci.yml
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
9 changes: 9 additions & 0 deletions .github/workflows/elixir-migration-check.yml
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
2 changes: 1 addition & 1 deletion DEVELOPMENT.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
44 changes: 39 additions & 5 deletions docs/docs.logflare.com/docs/backends/victoria-metrics.mdx
Original file line number Diff line number Diff line change
Expand Up @@ -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

Expand All @@ -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

Expand All @@ -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_<name>`, 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_<name>`, 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.
4 changes: 3 additions & 1 deletion docs/docs.logflare.com/docs/concepts/querying.md
Original file line number Diff line number Diff line change
Expand Up @@ -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.
:::

Expand All @@ -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.
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
Original file line number Diff line number Diff line change
@@ -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
Expand Down
156 changes: 117 additions & 39 deletions lib/logflare/backends/adaptor/victoria_metrics_adaptor.ex
Original file line number Diff line number Diff line change
Expand Up @@ -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_<name>`.
"""

Expand All @@ -37,15 +41,20 @@ 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
alias Logflare.Sources
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,
Expand Down Expand Up @@ -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

Expand All @@ -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
Expand All @@ -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}")

Expand All @@ -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

Expand All @@ -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
Expand Down Expand Up @@ -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

Expand Down Expand Up @@ -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
Expand Down
Loading
Loading