Skip to content
Open
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
2 changes: 2 additions & 0 deletions docs/docs.logflare.com/docs/concepts/endpoints.md
Original file line number Diff line number Diff line change
Expand Up @@ -106,6 +106,8 @@ Caching is performed on a query parameter basis. As such, if there are three API

When querying a specific endpoint version with `LF-ENDPOINT-VERSION`, the version number also partitions the cache. A request for the current endpoint and a request for version `1` use separate caches even when the query parameters match.

ClickHouse scheduling classes do not partition successful result caches. Operators can apply [query-class policy](../self-hosting/clickhouse-query-policy.md) on cache misses and refreshes without changing query results. Current operator-enforced endpoint limits also apply to historical versions.

### Proactive Requerying

Logflare endpoints can be proactively requeried to ensure that the cache does not become stale throughout the cache lifetime.
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,94 @@
---
title: ClickHouse query-class policy
---

# ClickHouse query-class policy

Operators can configure scheduling and resource limits independently of ClickHouse read-cluster routing. The backend configuration key is `query_class_settings`. It is disabled by default; an empty map removes the class policy without removing endpoint-owned limits.

## Operator configuration

Use the internal operator function with a persisted admin user. The admin flag is provisioned out of band and rechecked on every save. Customer backend changesets, the backend API, and customer forms cannot introduce, change, or clear this configuration. Public backend JSON omits it. An operator must clear this policy before changing backend type. Ordinary saves merge under a row lock so stale customer snapshots cannot erase newer operator policy. No new customer-supplied ClickHouse settings are enabled.

```elixir
admin = Logflare.Users.get(admin_id)
backend = Logflare.Backends.get_backend(backend_id)

settings = %{
"default" => %{"priority" => 5},
"dashboard_logs_free" => %{"priority" => 1},
"dashboard_logs_paid" => %{"priority" => 1},
"dashboard_reports_free" => %{"priority" => 1},
"dashboard_reports_paid" => %{"priority" => 1},
"dashboard_observability" => %{"priority" => 1},
"api_free" => %{"priority" => 10, "max_threads" => 4},
"api_paid" => %{"priority" => 10, "max_threads" => 4},
"mcp" => %{"priority" => 10, "max_threads" => 4}
}

{:ok, backend} =
Logflare.Backends.configure_query_class_settings(admin, backend, settings)
```

These values are an example, not production configuration. Choose them after checking the deployed query-user profile and load-testing in staging.

A nonempty configuration must include a `default` with a positive integer `priority`. A class can inherit its priority from this default. Missing, unknown, or unconfigured labels receive the default policy. Only the classes listed above are accepted at save time.

The class is selected from the original `LF-ENDPOINT-CLICKHOUSE-READ-CLUSTER-LABEL` request header, even when routing resolves to another cluster or retries on the default cluster. Known query-class labels do not generate unconfigured-routing warnings; unknown routing labels still do.

The header is a scheduling hint, not an authorization boundary. Preserve the existing endpoint authentication, and ensure the trusted Management API classifier supplies the header rather than forwarding an untrusted client's value. Selecting a label does not grant access to another tenant's data.

## Allowed settings and precedence

| Setting | Accepted value |
| --- | --- |
| `priority` | Positive integer; 1 is highest priority, larger values are lower priority |
| `max_threads` | Positive integer |
| `max_memory_usage` | Positive integer bytes |
| `max_bytes_to_read` | Positive integer bytes |
| `max_rows_to_read` | Positive integer rows |
| `max_execution_time` | Positive integer or fractional seconds |
| `read_overflow_mode` | Only `throw`; inserted automatically with byte or row limits |
| `timeout_overflow_mode` | Only `throw`; inserted automatically with an execution-time limit |

Zero (unlimited or no priority), negative values, wrong types, unknown keys, and overflow modes returning partial results are rejected. String values are serialized only from the fixed enum allowlist. Output-format, join semantics, sandbox permissions, and external table/URL settings are not allowed.

ClickHouse profile defaults supply settings not specified by Logflare. Logflare combines class default, selected class, and endpoint-owned settings. Resource limits use the smallest configured value across all three policies; non-limit settings such as priority use the later policy. Hard profile ceilings must be ClickHouse constraints, not merely default values. The server still enforces those constraints after merging.

Endpoint-owned policy is passed separately from consumer SQL. The adaptor checks for consumer overrides anywhere in the final AST, then injects the combined policy once. Supported `EXPLAIN SELECT` wrappers are preserved. Previews use the same policy builder:

```elixir
Logflare.Endpoints.get_transformed_query(endpoint, params, read_cluster: "api_paid")
```

Historical endpoint versions retain current operator ceilings. Cached snapshot reloads also overlay current endpoint policy. Saving endpoint-owned limits invalidates caches for all endpoint versions and stops their active refresh tasks.

Successful result caches remain class-agnostic: changing the class does not create a different result-cache key or run a query on a cache hit. Misses and background refreshes enforce policy at execution time. A refresh retains the class from the request that originally created the cache, but reads current backend policy. Backend class-policy changes invalidate backend metadata across the cluster, but do not evict reusable successful results, restart ingesters, or tear down active read connections.

## Deployment preflight

Before enabling this policy on a backend, connect using its actual query credentials, including any dedicated query user. Inspect the deployed settings and constraints:

```sql
SELECT name, type, value, min, max, readonly
FROM system.settings
WHERE name IN (
'readonly', 'priority', 'max_threads', 'max_execution_time',
'max_memory_usage', 'max_bytes_to_read', 'max_rows_to_read',
'read_overflow_mode', 'timeout_overflow_mode'
);
```

For every newly used key, issue a small SELECT with the intended setting using those credentials. Verify both the accepted values and rejection of values outside profile constraints. `readonly = 1` can prohibit changes unless the setting is specifically changeable in readonly mode; `readonly = 2` allows settings changes subject to constraints. Do not weaken `readonly`, grant DDL, or make `readonly` itself changeable to make these probes pass. See [ClickHouse setting constraints](https://clickhouse.com/docs/operations/settings/constraints-on-settings).

Local regression tests exercise every allowlisted key with a temporary `readonly = 2` user and profile ceilings. They do not establish that the production `logflare_prod` profile permits these values. Production verification is an operator rollout requirement.

Deploy Logflare support on every instance and drain older versions before enabling the policy. Apply the backend class configuration and a matching `priority = 5` baseline in the ClickHouse `logflare` profile in the same coordinated rollout window, **not the profile change before Logflare class policy**. Consider `priority MIN 1` to prevent bypass paths from opting out with zero. Ensure configured limits fit the profile's constraints; use server-side MAX constraints where cluster-wide ceilings are required.

## Verification and rollout monitoring

Ensure query-setting logging is enabled for the verification queries. For actual API/MCP, dashboard, missing-label, and unknown-label requests, inspect `system.query_log.Settings` on Logflare Reads and verify priority and limits, including requests routed to the default cluster. Use the query ID to identify the request; a small local regression test verifies the recorded API values.

In staging, compare dashboard p95 latency and timeout rates under competing API load. During rollout, monitor query errors, delayed concurrency slots, paused queries, memory, and connections. Lower-priority queries retain memory and connections while paused and may see more execution-time failures. Priority is not tenant isolation or aggregate admission control.

Aggregate Management API throttling (O11Y-2662) and moving API/MCP classes to separately sized compute remain separate mitigations. Roll back class policy through the admin function with `%{}`; coordinate any profile rollback after accounting for requests and instances still using class policy.
92 changes: 74 additions & 18 deletions lib/logflare/backends.ex
Original file line number Diff line number Diff line change
Expand Up @@ -8,6 +8,7 @@ defmodule Logflare.Backends do

alias Ecto.Changeset
alias Logflare.Backends.Adaptor
alias Logflare.Backends.Adaptor.ClickHouseAdaptor.QueryClassSettings
alias Logflare.Backends.Backend
alias Logflare.Backends.BackendRegistry
alias Logflare.Backends.ConsolidatedSup
Expand Down Expand Up @@ -318,7 +319,27 @@ defmodule Logflare.Backends do
Updates the config of a Backend.
"""
@spec update_backend(Backend.t(), map()) :: {:ok, Backend.t()} | {:error, Changeset.t()}
def update_backend(%Backend{} = backend, attrs) do
def update_backend(%Backend{} = backend, attrs), do: update_backend_config(backend, attrs, [])

@spec configure_query_class_settings(User.t(), Backend.t(), map()) ::
{:ok, Backend.t()} | {:error, term()}
def configure_query_class_settings(%User{id: user_id}, %Backend{id: id}, settings) do
with %User{admin: true} <- Repo.get(User, user_id),
%Backend{type: :clickhouse} = backend <- get_backend(id),
{:ok, settings} <- QueryClassSettings.normalize(settings) do
update_backend_config(backend, %{config: %{query_class_settings: settings}},
allow_query_class_settings: true,
query_policy_only: true
)
else
nil -> {:error, :not_found}
%User{} -> {:error, :forbidden}
%Backend{} -> {:error, :not_clickhouse_backend}
error -> error
end
end

defp update_backend_config(%Backend{} = backend, attrs, opts) do
alerts_modified = Map.has_key?(attrs, :alert_queries)

default_ingest_modified? =
Expand All @@ -328,20 +349,29 @@ defmodule Logflare.Backends do

source_id = Map.get(attrs, "source_id") || Map.get(attrs, :source_id)

changeset =
backend
|> Backend.changeset(attrs)
|> validate_default_ingest_source(source_id)
|> then(fn changeset ->
if alerts_modified do
Changeset.put_assoc(changeset, :alert_queries, Map.get(attrs, :alert_queries))
else
changeset
result =
Repo.transact(fn ->
current =
from(b in Backend, where: b.id == ^backend.id, lock: "FOR UPDATE")
|> Repo.one!()
|> typecast_config_string_map_to_atom_map()

current = if alerts_modified, do: Repo.preload(current, :alert_queries), else: current

changeset =
current
|> Backend.changeset(attrs, opts)
|> validate_default_ingest_source(source_id)
|> maybe_put_backend_alerts(attrs, alerts_modified)

with :ok <- validate_query_policy_backend(current, opts),
{:ok, updated} <- Repo.update(changeset) do
{:ok, {current, updated}}
end
end)

case Repo.update(changeset) do
{:ok, updated} ->
case result do
{:ok, {backend, updated}} ->
updated = preload_sources(updated)

updated =
Expand All @@ -357,11 +387,15 @@ defmodule Logflare.Backends do
updated
end

Enum.each(updated.sources, &restart_source_sup(&1))

apply_consolidated_pipeline_change(backend, updated, default_ingest_modified?)

if config_modified?, do: Adaptor.on_backend_config_changed(updated)
if Keyword.get(opts, :query_policy_only, false) do
keys = [{__MODULE__, updated.id}]
Logflare.ContextCache.bust_keys(keys)
Cluster.Utils.rpc_multicast(Logflare.ContextCache, :bust_keys, [keys])
else
Enum.each(updated.sources, &restart_source_sup(&1))
apply_consolidated_pipeline_change(backend, updated, default_ingest_modified?)
notify_backend_config_change(updated, config_modified?)
end

{:ok, typecast_config_string_map_to_atom_map(updated)}

Expand All @@ -370,11 +404,33 @@ defmodule Logflare.Backends do
end
end

@spec maybe_put_backend_alerts(Changeset.t(), map(), boolean()) :: Changeset.t()
defp maybe_put_backend_alerts(changeset, attrs, true),
do: Changeset.put_assoc(changeset, :alert_queries, Map.get(attrs, :alert_queries))

defp maybe_put_backend_alerts(changeset, _attrs, false), do: changeset

@spec notify_backend_config_change(Backend.t(), boolean()) :: term()
defp notify_backend_config_change(_backend, false), do: :ok

defp notify_backend_config_change(backend, true),
do: Adaptor.on_backend_config_changed(backend)

@spec validate_query_policy_backend(Backend.t(), Keyword.t()) ::
:ok | {:error, :not_clickhouse_backend}
defp validate_query_policy_backend(%Backend{type: :clickhouse}, _opts), do: :ok

defp validate_query_policy_backend(_backend, opts) do
if Keyword.get(opts, :query_policy_only, false),
do: {:error, :not_clickhouse_backend},
else: :ok
end

@spec validate_default_ingest_source(Changeset.t(), String.t() | integer() | nil) ::
Changeset.t()
defp validate_default_ingest_source(%{changes: %{default_ingest?: true}} = changeset, source_id)
when is_non_empty_binary(source_id) or is_integer(source_id) do
case Sources.get(source_id) do
case Repo.get(Source, source_id) do
%Source{default_ingest_backend_enabled?: true} ->
changeset

Expand Down
86 changes: 69 additions & 17 deletions lib/logflare/backends/adaptor/clickhouse_adaptor.ex
Original file line number Diff line number Diff line change
Expand Up @@ -21,6 +21,7 @@ defmodule Logflare.Backends.Adaptor.ClickHouseAdaptor do
alias __MODULE__.Ingester
alias __MODULE__.Pipeline
alias __MODULE__.Provisioner
alias __MODULE__.QueryClassSettings
alias __MODULE__.QueryConnectionSup
alias __MODULE__.QueryErrorNormalizer
alias __MODULE__.QueryTemplates
Expand All @@ -34,6 +35,7 @@ defmodule Logflare.Backends.Adaptor.ClickHouseAdaptor do
alias Logflare.Backends.IngestEventQueue
alias Logflare.Backends.Adaptor.QueryResult
alias Logflare.Backends.QueryError
alias Logflare.Endpoints.ClickHouseSettings
alias Logflare.LogEvent
alias Logflare.LogEvent.TypeDetection
alias Logflare.Sources.Source
Expand All @@ -57,6 +59,7 @@ defmodule Logflare.Backends.Adaptor.ClickHouseAdaptor do
:labeled_read_pool_size,
:read_only_urls,
:default_read_cluster,
:query_class_settings,
:query_user,
:query_password,
:replica_routing_param
Expand Down Expand Up @@ -144,6 +147,7 @@ defmodule Logflare.Backends.Adaptor.ClickHouseAdaptor do
def redact_config(config) do
config
|> Map.put(:password, "REDACTED")
|> Map.delete(:query_class_settings)
|> redact_query_password()
end

Expand Down Expand Up @@ -198,6 +202,13 @@ defmodule Logflare.Backends.Adaptor.ClickHouseAdaptor do
)
when is_non_empty_binary(query_string) and is_list(declared_params) and is_map(input_params) and
is_list(opts) do
endpoint_settings =
if is_map(endpoint_query),
do: Map.get(endpoint_query, :enforced_clickhouse_settings, %{}),
else: %{}

opts = Keyword.put(opts, :enforced_clickhouse_settings, endpoint_settings)

with {:ok, {limited_query, max_rows}} <- limit_endpoint_query(query_string, endpoint_query) do
execute_query_with_params(
backend,
Expand Down Expand Up @@ -314,6 +325,7 @@ defmodule Logflare.Backends.Adaptor.ClickHouseAdaptor do
labeled_read_pool_size: :integer,
read_only_urls: {:map, :string},
default_read_cluster: :string,
query_class_settings: :map,
use_async_inserts_for_small_batches: :boolean,
use_async_inserts_only: :boolean,
async_insert_cluster_url: :string,
Expand All @@ -335,6 +347,7 @@ defmodule Logflare.Backends.Adaptor.ClickHouseAdaptor do
:labeled_read_pool_size,
:read_only_urls,
:default_read_cluster,
:query_class_settings,
:use_async_inserts_for_small_batches,
:use_async_inserts_only,
:async_insert_cluster_url,
Expand Down Expand Up @@ -413,6 +426,7 @@ defmodule Logflare.Backends.Adaptor.ClickHouseAdaptor do
)
|> validate_read_only_urls()
|> validate_default_read_cluster()
|> validate_query_class_settings()
|> validate_user_pass()
|> validate_query_user_pass()
|> validate_number(:read_pool_size,
Expand Down Expand Up @@ -745,27 +759,62 @@ defmodule Logflare.Backends.Adaptor.ClickHouseAdaptor do

def execute_ch_query(%Backend{} = backend, statement, params, opts)
when is_list_or_map(params) and is_list(opts) do
requested = Keyword.get(opts, :read_cluster)
label = resolve_read_cluster_label(backend, requested)
with {:ok, statement} <- prepare_query(backend, statement, opts) do
requested = Keyword.get(opts, :read_cluster)
label = resolve_read_cluster_label(backend, requested)

warn_on_unconfigured_read_cluster(backend, requested, label)
headers = Keyword.get(opts, :headers, [])
warn_on_unconfigured_read_cluster(backend, requested, label)
headers = Keyword.get(opts, :headers, [])

{result, queried_label} =
case do_ch_query_on_label(backend, statement, params, label, headers) do
{:error, %QueryError{kind: :connection_error}} = error ->
maybe_retry_on_default_cluster(backend, statement, params, label, headers, error)
{result, queried_label} =
case do_ch_query_on_label(backend, statement, params, label, headers) do
{:error, %QueryError{kind: :connection_error}} = error ->
maybe_retry_on_default_cluster(backend, statement, params, label, headers, error)

result ->
{result, label}
end
result ->
{result, label}
end

emit_query_error_telemetry(result, %{
backend_id: backend.id,
read_cluster: read_cluster_tag(queried_label)
})
emit_query_error_telemetry(result, %{
backend_id: backend.id,
read_cluster: read_cluster_tag(queried_label)
})

result
result
else
{:error, reason} ->
{:error, %{query_error(:invalid_query, reason) | description: reason}}
end
end

@spec prepare_query(Backend.t(), iodata(), Keyword.t()) ::
{:ok, String.t()} | {:error, String.t()}
def prepare_query(%Backend{config: config}, statement, opts \\ []) do
classes = Map.get(config, :query_class_settings) || %{}
endpoint_settings = Keyword.get(opts, :enforced_clickhouse_settings, %{})

with {:ok, settings} <-
QueryClassSettings.resolve(
classes,
Keyword.get(opts, :read_cluster),
endpoint_settings
) do
ClickHouseSettings.enforce(IO.iodata_to_binary(statement), settings)
end
end

@spec validate_query_class_settings(Changeset.t()) :: Changeset.t()
defp validate_query_class_settings(changeset) do
case Changeset.get_field(changeset, :query_class_settings) do
nil ->
changeset

settings ->
case QueryClassSettings.normalize(settings) do
{:ok, settings} -> Changeset.put_change(changeset, :query_class_settings, settings)
{:error, reason} -> Changeset.add_error(changeset, :query_class_settings, reason)
end
end
end

@spec do_ch_query_on_label(
Expand Down Expand Up @@ -816,7 +865,10 @@ defmodule Logflare.Backends.Adaptor.ClickHouseAdaptor do
@spec warn_on_unconfigured_read_cluster(Backend.t(), String.t() | nil, String.t() | nil) :: :ok
defp warn_on_unconfigured_read_cluster(%Backend{} = backend, requested, label)
when is_non_empty_binary(requested) and requested != label do
Logger.warning(
level = if QueryClassSettings.class_label?(requested), do: :debug, else: :warning

Logger.log(
level,
"ClickHouse read cluster not configured, falling back to resolved read cluster",
user_id: backend.user_id,
backend_id: backend.id,
Expand Down
Loading
Loading