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
22 changes: 20 additions & 2 deletions docs/docs.logflare.com/docs/concepts/endpoints.md
Original file line number Diff line number Diff line change
Expand Up @@ -102,7 +102,9 @@ All endpoint queries by default are set to cache results for 3,600 seconds. The

To disable the cache, set the cached duration to `0`. However, it is not recommended to do so unless you absolutely need up-to-date results. To prevent stale data while keeping the cache warm, use the cache proactive requerying feature.

Caching is performed on a query parameter basis. As such, if there are three API requests sent to an endpoint, `?path=123`, `?path=223`, and `?other=value`, this will result in 3 difference caches being created.
Caching is performed on a query parameter basis. As such, if there are three API requests sent to an endpoint, `?path=123`, `?path=223`, and `?other=value`, this will result in 3 different caches being created.

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.

### Proactive Requerying

Expand All @@ -114,6 +116,8 @@ When configured, the cache will be automatically updated at the set interval, pe

Endpoints support query labeling for tracking and monitoring in the backend. Labels are configured as a comma-separated allowlist and can reference parameters (`@my_param`), values provided in the `LF-ENDPOINT-LABELS` request header, or static values.

Logflare also adds a reserved `endpoint_version` label automatically for endpoint executions that use `LF-ENDPOINT-VERSION`. For example, `LF-ENDPOINT-VERSION: 1` produces `endpoint_version=1`. This label is used for query execution logging and overrides any caller-provided label with the same key. Requests without `LF-ENDPOINT-VERSION` do not include an `endpoint_version` label.

### Configuration Format

```text
Expand Down Expand Up @@ -149,6 +153,21 @@ LF-ENDPOINT-LABELS: session_id=abc123,ignored=xyz

Only allowlisted labels are processed. Query parameters override header values for the same key.

## Querying a Specific Endpoint Version

Endpoints can be queried at a previous version by passing the `LF-ENDPOINT-VERSION` request header with the version number.

```bash
curl "https://api.logflare.app/api/endpoints/query/my-endpoint" \
-H 'X-API-KEY: YOUR-ACCESS-TOKEN' \
-H 'LF-ENDPOINT-VERSION: 1' \
-H 'Content-Type: application/json; charset=utf-8'
```

If the header is omitted the current endpoint definition is used. If the requested version does not exist for the resolved endpoint the response contains `version not found`.

Versioned endpoint requests are tagged in query execution logs. For example, `LF-ENDPOINT-VERSION: 1` produces `endpoint_version=1`. See [Query tagging with labels](#query-tagging-with-labels) for more details on labels.

## Subquery Expansion with Other Endpoints

Logflare endpoints support subquery expansion, allowing you to reference and query data from other endpoints within your SQL queries. This enables powerful data composition and cross-endpoint analysis.
Expand Down Expand Up @@ -241,4 +260,3 @@ The main Logflare service account must be granted the following IAM roles on the

- **Project IAM Admin** (`roles/resourcemanager.projectIamAdmin`) — required to set IAM policies for managed service accounts on the additional project
- **BigQuery Admin** (`roles/bigquery.admin`) — required for managed service accounts to execute queries against reservations in the additional project

61 changes: 59 additions & 2 deletions lib/logflare/endpoints.ex
Original file line number Diff line number Diff line change
Expand Up @@ -12,8 +12,8 @@ defmodule Logflare.Endpoints do
alias Logflare.Backends.Adaptor.QueryResult
alias Logflare.Backends.Backend
alias Logflare.Backends.QueryError
alias Logflare.Endpoints.PiiRedactor
alias Logflare.Endpoints.EndpointQuery
alias Logflare.Endpoints.PiiRedactor
alias Logflare.Endpoints.Resolver
alias Logflare.Endpoints.ResultsCache
alias Logflare.Lql
Expand All @@ -30,6 +30,7 @@ defmodule Logflare.Endpoints do
alias Logflare.Utils
alias PaperTrail.Version

@endpoint_version_label "endpoint_version"
@valid_sql_languages ~w(bq_sql ch_sql pg_sql)a

@typep language :: :bq_sql | :ch_sql | :pg_sql | :lql
Expand Down Expand Up @@ -100,6 +101,31 @@ defmodule Logflare.Endpoints do
def get_endpoint_query(query_id) when is_integer_or_string(query_id),
do: Repo.get(EndpointQuery, query_id)

@spec get_endpoint_query_at_version(EndpointQuery.t() | integer() | String.t(), integer()) ::
{:ok, EndpointQuery.t()} | {:error, :invalid_version | :version_not_found}
def get_endpoint_query_at_version(%EndpointQuery{} = endpoint, version_number)
when is_integer(version_number) and version_number > 0 do
case get_endpoint_query_version(endpoint.id, version_number) do
%Version{} = version ->
{:ok, endpoint_from_version(endpoint, version, version_number)}

nil ->
{:error, :version_not_found}
end
end

def get_endpoint_query_at_version(%EndpointQuery{}, _version_number),
do: {:error, :invalid_version}

def get_endpoint_query_at_version(nil, _version_number), do: {:error, :version_not_found}

def get_endpoint_query_at_version(query_id, version_number)
when is_integer_or_string(query_id) do
query_id
|> get_endpoint_query()
|> get_endpoint_query_at_version(version_number)
end

@doc """
Retrieves a mapped endpoint `EndpointQuery` by token
"""
Expand Down Expand Up @@ -312,11 +338,28 @@ defmodule Logflare.Endpoints do
from(version in Version,
where:
version.item_type == "EndpointQuery" and version.item_id == ^endpoint_id and
fragment("(?->>'version_number')::integer = ?", version.meta, ^version_number)
fragment("?->>'version_number' = ?", version.meta, ^Integer.to_string(version_number))
)
|> Repo.one()
end

@spec endpoint_from_version(EndpointQuery.t(), Version.t(), pos_integer()) :: EndpointQuery.t()
defp endpoint_from_version(%EndpointQuery{} = endpoint, %Version{meta: meta}, version_number)
when is_map(meta) do
snapshot = Map.get(meta, "endpoint_snapshot", %{})

%EndpointQuery{}
|> Ecto.Changeset.cast(snapshot, EndpointQuery.version_snapshot_fields())
|> Ecto.Changeset.apply_changes()
|> Map.merge(%{
id: endpoint.id,
user_id: endpoint.user_id,
token: endpoint.token,
version_number: version_number
})
|> map_query_sources()
end

@spec delete_query(EndpointQuery.t(), origin()) :: {:ok, EndpointQuery.t()} | {:error, any()}
def delete_query(%EndpointQuery{} = query, origin) do
Repo.transact(fn ->
Expand Down Expand Up @@ -443,6 +486,8 @@ defmodule Logflare.Endpoints do
opts \\ []
)
when is_map(params) and is_list(opts) do
endpoint_query = put_endpoint_version_label(endpoint_query)

%EndpointQuery{query: query_string, user_id: user_id, sandboxable: sandboxable} =
endpoint_query

Expand Down Expand Up @@ -494,6 +539,18 @@ defmodule Logflare.Endpoints do
end
end

@spec put_endpoint_version_label(EndpointQuery.t()) :: EndpointQuery.t()
defp put_endpoint_version_label(%EndpointQuery{version_number: version_number} = endpoint_query)
when is_integer(version_number) do
parsed_labels =
(endpoint_query.parsed_labels || %{})
|> Map.put(@endpoint_version_label, Integer.to_string(version_number))

%{endpoint_query | parsed_labels: parsed_labels}
end

defp put_endpoint_version_label(%EndpointQuery{} = endpoint_query), do: endpoint_query

@spec emit_query_telemetry(result :: run_query_return(), endpoint_query :: EndpointQuery.t()) ::
{run_query_return(), map()}
defp emit_query_telemetry({:ok, data} = result, endpoint_query) do
Expand Down
4 changes: 4 additions & 0 deletions lib/logflare/endpoints/cache.ex
Original file line number Diff line number Diff line change
Expand Up @@ -29,6 +29,10 @@ defmodule Logflare.Endpoints.Cache do
end

def get_endpoint_query(kw), do: apply_repo_fun(:get_endpoint_query, [kw])

def get_endpoint_query_at_version(query_id, version_number),
do: apply_repo_fun(:get_endpoint_query_at_version, [query_id, version_number])

def get_by(kw), do: apply_repo_fun(:get_by, [kw])
def get_mapped_query_by_token(token), do: apply_repo_fun(:get_mapped_query_by_token, [token])

Expand Down
1 change: 1 addition & 0 deletions lib/logflare/endpoints/endpoint_query.ex
Original file line number Diff line number Diff line change
Expand Up @@ -71,6 +71,7 @@ defmodule Logflare.Endpoints.EndpointQuery do
field(:labels, :string)
field(:parsed_labels, :map, virtual: true)
field(:metrics, :map, virtual: true)
field(:version_number, :integer, virtual: true)

belongs_to(:user, User)
belongs_to(:backend, Backend)
Expand Down
12 changes: 7 additions & 5 deletions lib/logflare/endpoints/resolver.ex
Original file line number Diff line number Diff line change
Expand Up @@ -10,10 +10,11 @@ defmodule Logflare.Endpoints.Resolver do
@doc """
Lists all caches for an endpoint across all paritions
"""
def list_caches(%Logflare.Endpoints.EndpointQuery{id: id}) do
endpoints_partition = ResultsCache.endpoints_part(id)
def list_caches(%Logflare.Endpoints.EndpointQuery{} = query) do
{query_id, version_key} = ResultsCache.endpoint_cache_key(query)
endpoints_partition = ResultsCache.endpoints_part(query_id, version_key)

:syn.members(endpoints_partition, id)
:syn.members(endpoints_partition, {query_id, version_key})
|> Enum.map(fn {pid, _} -> pid end)
end

Expand All @@ -29,7 +30,7 @@ defmodule Logflare.Endpoints.Resolver do
"endpoint.user_id" => query.user_id
}

ResultsCache.name(query.id, params)
ResultsCache.name(query, params)
|> GenServer.whereis()
|> case do
pid when is_pid(pid) ->
Expand All @@ -45,7 +46,8 @@ defmodule Logflare.Endpoints.Resolver do

via =
{:via, PartitionSupervisor,
{Logflare.Endpoints.ResultsCache.PartitionSupervisor, {id, params, opts}}}
{Logflare.Endpoints.ResultsCache.PartitionSupervisor,
ResultsCache.cache_partition_key(query, params, opts)}}

case DynamicSupervisor.start_child(via, spec) do
{:ok, pid} ->
Expand Down
85 changes: 66 additions & 19 deletions lib/logflare/endpoints/results_cache.ex
Original file line number Diff line number Diff line change
Expand Up @@ -6,11 +6,14 @@ defmodule Logflare.Endpoints.ResultsCache do
require Logger

alias Logflare.Endpoints
alias Logflare.Endpoints.EndpointQuery
alias Logflare.Utils
alias Logflare.Utils.Tasks

use GenServer, restart: :temporary

@latest_version :latest

defstruct query_tasks: [],
params: %{},
opts: [],
Expand All @@ -19,22 +22,24 @@ defmodule Logflare.Endpoints.ResultsCache do
refresh_timer: nil,
endpoint_query_id: nil,
endpoint_query_token: nil,
endpoint_version_number: nil,
parsed_labels: %{}

@type t :: %__MODULE__{
endpoint_query_id: integer(),
endpoint_query_token: String.t(),
endpoint_version_number: integer() | nil,
query_tasks: list(%Task{}),
params: map(),
opts: list(),
cached_result: binary(),
shutdown_timer: reference(),
refresh_timer: reference(),
cached_result: map() | nil,
shutdown_timer: reference() | nil,
refresh_timer: reference() | nil,
parsed_labels: map()
}

def start_link({query, params, opts}) do
name = name(query.id, params)
name = name(query, params)
GenServer.start_link(__MODULE__, {query, params, opts}, name: name, hibernate_after: 5_000)
end

Expand Down Expand Up @@ -90,15 +95,17 @@ defmodule Logflare.Endpoints.ResultsCache do
end

def init({query, params, opts}) do
endpoints = endpoints_part(query.id)
:syn.join(endpoints, query.id, self())
endpoint_cache_key = {query_id, version_key} = endpoint_cache_key(query)
endpoints = endpoints_part(query_id, version_key)
:syn.join(endpoints, endpoint_cache_key, self())

timer = query |> cache_duration_ms() |> shutdown()

state =
%__MODULE__{
endpoint_query_id: query.id,
endpoint_query_token: query.token,
endpoint_version_number: query.version_number,
params: params,
opts: opts,
shutdown_timer: timer,
Expand Down Expand Up @@ -188,28 +195,68 @@ defmodule Logflare.Endpoints.ResultsCache do
end

def do_query(state) do
query =
Endpoints.Cache.get_mapped_query_by_token(state.endpoint_query_token)
|> Map.put(:parsed_labels, state.parsed_labels)
with %EndpointQuery{} = query <- endpoint_query(state) do
query = Map.put(query, :parsed_labels, state.parsed_labels)

Logflare.Endpoints.run_query(query, state.params, state.opts)
|> Utils.append_to_tuple(query)
else
{:error, error} ->
{:error, error, nil}

Logflare.Endpoints.run_query(query, state.params, state.opts)
|> Utils.append_to_tuple(query)
nil ->
{:error, :not_found, nil}
end
end

def endpoints_part(query_id, params) do
part = :erlang.phash2({query_id, params}, System.schedulers_online())
"endpoints_#{part}" |> String.to_existing_atom()
def endpoints_part(query_id, version_key) do
endpoints_part({query_id, version_key})
end

def endpoints_part(query_id) do
part = :erlang.phash2(query_id, System.schedulers_online())
def endpoints_part(query_id, version_key, params) do
endpoints_part({query_id, version_key, params})
end

defp endpoints_part(partition_key) do
part = :erlang.phash2(partition_key, System.schedulers_online())
"endpoints_#{part}" |> String.to_existing_atom()
end

def name(query_id, params) do
partition = endpoints_part(query_id, params)
def name(%EndpointQuery{} = query, params) do
param_hash = :erlang.phash2(params)
{:via, :syn, {partition, {query_id, param_hash}}}
{query_id, version_key} = endpoint_cache_key(query)

key = {query_id, version_key, param_hash}

{:via, :syn, {endpoints_part(query_id, version_key, params), key}}
end

@spec cache_partition_key(EndpointQuery.t(), map(), Keyword.t()) :: tuple()
def cache_partition_key(%EndpointQuery{} = query, params, opts) do
{query_id, version_key} = endpoint_cache_key(query)

{query_id, version_key, params, opts}
end

def endpoint_cache_key(%EndpointQuery{id: id, version_number: version_number}),
do: {id, version_cache_key(version_number)}

defp version_cache_key(version_number) when is_integer(version_number),
do: {:version, version_number}

defp version_cache_key(_version_number), do: {:version, @latest_version}

@spec endpoint_query(t()) :: EndpointQuery.t() | {:error, atom()} | nil
defp endpoint_query(%__MODULE__{endpoint_version_number: version_number} = state)
when is_integer(version_number) do
case Endpoints.Cache.get_endpoint_query_at_version(state.endpoint_query_id, version_number) do
{:ok, query} -> query
{:error, error} -> {:error, error}
end
end

defp endpoint_query(%__MODULE__{} = state) do
Endpoints.Cache.get_mapped_query_by_token(state.endpoint_query_token)
end

defp refresh(every) do
Expand Down
6 changes: 6 additions & 0 deletions lib/logflare_web/controllers/api/fallback_controller.ex
Original file line number Diff line number Diff line change
Expand Up @@ -48,6 +48,12 @@ defmodule LogflareWeb.Api.FallbackController do
|> json(%{error: "Not Found"})
end

def call(conn, {:error, :not_found, message}) when is_binary(message) do
conn
|> put_status(:not_found)
|> json(%{error: message})
end

def call(conn, {:error, %QueryError{}}) do
conn
|> put_status(400)
Expand Down
Loading
Loading