diff --git a/docs/docs.logflare.com/docs/concepts/endpoints.md b/docs/docs.logflare.com/docs/concepts/endpoints.md index 4c6de23c2d..a54f90fba9 100644 --- a/docs/docs.logflare.com/docs/concepts/endpoints.md +++ b/docs/docs.logflare.com/docs/concepts/endpoints.md @@ -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 @@ -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 @@ -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. @@ -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 - diff --git a/lib/logflare/endpoints.ex b/lib/logflare/endpoints.ex index 681334eabc..e0a1db25c5 100644 --- a/lib/logflare/endpoints.ex +++ b/lib/logflare/endpoints.ex @@ -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 @@ -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 @@ -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 """ @@ -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 -> @@ -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 @@ -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 diff --git a/lib/logflare/endpoints/cache.ex b/lib/logflare/endpoints/cache.ex index b12c609ff1..4a653589f4 100644 --- a/lib/logflare/endpoints/cache.ex +++ b/lib/logflare/endpoints/cache.ex @@ -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]) diff --git a/lib/logflare/endpoints/endpoint_query.ex b/lib/logflare/endpoints/endpoint_query.ex index 548d97a78b..43f9eeb830 100644 --- a/lib/logflare/endpoints/endpoint_query.ex +++ b/lib/logflare/endpoints/endpoint_query.ex @@ -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) diff --git a/lib/logflare/endpoints/resolver.ex b/lib/logflare/endpoints/resolver.ex index 1b1103cf3f..66e37b80be 100644 --- a/lib/logflare/endpoints/resolver.ex +++ b/lib/logflare/endpoints/resolver.ex @@ -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 @@ -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) -> @@ -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} -> diff --git a/lib/logflare/endpoints/results_cache.ex b/lib/logflare/endpoints/results_cache.ex index d8eaa74621..6b7e9df935 100644 --- a/lib/logflare/endpoints/results_cache.ex +++ b/lib/logflare/endpoints/results_cache.ex @@ -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: [], @@ -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 @@ -90,8 +95,9 @@ 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() @@ -99,6 +105,7 @@ defmodule Logflare.Endpoints.ResultsCache do %__MODULE__{ endpoint_query_id: query.id, endpoint_query_token: query.token, + endpoint_version_number: query.version_number, params: params, opts: opts, shutdown_timer: timer, @@ -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 diff --git a/lib/logflare_web/controllers/api/fallback_controller.ex b/lib/logflare_web/controllers/api/fallback_controller.ex index ad648a2c78..f5fe776120 100644 --- a/lib/logflare_web/controllers/api/fallback_controller.ex +++ b/lib/logflare_web/controllers/api/fallback_controller.ex @@ -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) diff --git a/lib/logflare_web/controllers/endpoints_controller.ex b/lib/logflare_web/controllers/endpoints_controller.ex index 3a9598d50d..d8d2091d72 100644 --- a/lib/logflare_web/controllers/endpoints_controller.ex +++ b/lib/logflare_web/controllers/endpoints_controller.ex @@ -6,9 +6,12 @@ defmodule LogflareWeb.EndpointsController do alias Logflare.Backends.QueryError alias Logflare.Endpoints + alias LogflareWeb.Api.FallbackController alias LogflareWeb.JsonParser - alias LogflareWeb.OpenApi.Unauthorized + alias LogflareWeb.OpenApi.BadRequest + alias LogflareWeb.OpenApi.NotFound alias LogflareWeb.OpenApi.ServerError + alias LogflareWeb.OpenApi.Unauthorized alias LogflareWeb.OpenApiSchemas.EndpointQuery alias LogflareWeb.QueryErrorHelpers @@ -32,7 +35,8 @@ defmodule LogflareWeb.EndpointsController do "X-API-Key", "LF-ENDPOINT-LABELS", "LF-ENDPOINT-REDACT-PII", - "LF-ENDPOINT-BIGQUERY-RESERVATION" + "LF-ENDPOINT-BIGQUERY-RESERVATION", + "LF-ENDPOINT-VERSION" ], methods: ["GET", "POST", "OPTIONS"], send_preflight_response?: true @@ -48,20 +52,28 @@ defmodule LogflareWeb.EndpointsController do type: :string, example: "a040ae88-3e27-448b-9ee6-622278b23193", required: true + ], + lf_endpoint_version: [ + in: :header, + name: "LF-ENDPOINT-VERSION", + description: "Endpoint version number to execute", + type: :integer, + required: false ] ], responses: %{ 200 => EndpointQuery.response(), + 400 => BadRequest.response(), 401 => Unauthorized.response(), + 404 => NotFound.response(), 500 => ServerError.response() } ) + plug(:load_endpoint_query) plug(:parse_get_body) - def query(%{assigns: %{endpoint: endpoint}} = conn, params) do - endpoint_query = Endpoints.map_query_sources(endpoint) - + def query(%{assigns: %{endpoint_query: endpoint_query}} = conn, params) do header_str = get_req_header(conn, "lf-endpoint-labels") |> case do @@ -98,13 +110,44 @@ defmodule LogflareWeb.EndpointsController do end end - # only parse body for get when `?sql=` and `?lql=` are empty and it is sandboxable + @spec load_endpoint_query(Plug.Conn.t(), any()) :: Plug.Conn.t() + defp load_endpoint_query(%{assigns: %{endpoint: endpoint}} = conn, _opts) do + with [version_str | _] <- get_req_header(conn, "lf-endpoint-version"), + {:ok, version_number} <- parse_endpoint_version(version_str), + {:ok, endpoint_query} <- + Endpoints.get_endpoint_query_at_version(endpoint, version_number) do + assign(conn, :endpoint_query, endpoint_query) + else + [] -> + assign(conn, :endpoint_query, Endpoints.map_query_sources(endpoint)) + + {:error, :version_not_found} -> + conn + |> FallbackController.call({:error, :not_found, "version not found"}) + |> halt() + + _ -> + conn + |> FallbackController.call({:error, "invalid lf-endpoint-version"}) + |> halt() + end + end + + # only parse body for get when `?sql=` and `?lql=` are empty # passthrough for all other cases defp parse_get_body( - %{method: "GET", assigns: %{endpoint: %_{sandboxable: true}}, query_params: qp} = conn, + %{method: "GET", assigns: %{endpoint_query: %{sandboxable: true}}, query_params: qp} = + conn, _opts ) when not is_map_key(qp, "sql") and not is_map_key(qp, "lql") do + parse_get_body(conn) + end + + defp parse_get_body(conn, _opts), do: conn + + @spec parse_get_body(Plug.Conn.t()) :: Plug.Conn.t() + defp parse_get_body(conn) do conn # Plug.Parsers only supports POST/PUT/PATCH |> Map.put(:method, "POST") @@ -113,5 +156,14 @@ defmodule LogflareWeb.EndpointsController do |> Map.put(:method, "GET") end - defp parse_get_body(conn, _opts), do: conn + @spec parse_endpoint_version(String.t()) :: {:ok, pos_integer()} | {:error, :invalid_version} + defp parse_endpoint_version(version_str) do + case Integer.parse(version_str) do + {version_number, ""} when version_number > 0 -> + {:ok, version_number} + + _ -> + {:error, :invalid_version} + end + end end diff --git a/lib/logflare_web/live/endpoints/endpoints_versions_live.ex b/lib/logflare_web/live/endpoints/endpoints_versions_live.ex index b0635e1c1e..be2e186918 100644 --- a/lib/logflare_web/live/endpoints/endpoints_versions_live.ex +++ b/lib/logflare_web/live/endpoints/endpoints_versions_live.ex @@ -7,7 +7,6 @@ defmodule LogflareWeb.EndpointsVersionsLive do import LogflareWeb.Utils, only: [time_ago: 1] alias Logflare.Endpoints - alias Logflare.Endpoints.EndpointQuery alias Logflare.Repo alias LogflareWeb.Endpoints.Components alias LogflareWeb.Endpoints.SnapshotModalComponent @@ -161,10 +160,12 @@ defmodule LogflareWeb.EndpointsVersionsLive do socket = with {version_number, ""} <- Integer.parse(version_number), selected_version when is_struct(selected_version) <- - Endpoints.get_endpoint_query_version(endpoint.id, version_number) do + Endpoints.get_endpoint_query_version(endpoint.id, version_number), + {:ok, endpoint_snapshot} <- + Endpoints.get_endpoint_query_at_version(endpoint, version_number) do socket |> assign(:selected_version, selected_version) - |> assign(:endpoint_snapshot, snapshot_to_endpoint(selected_version)) + |> assign(:endpoint_snapshot, endpoint_snapshot) else _ -> socket @@ -313,15 +314,6 @@ defmodule LogflareWeb.EndpointsVersionsLive do defp version_number(%Version{meta: %{"version_number" => version_number}}), do: version_number defp version_number(_version), do: nil - @spec snapshot_to_endpoint(Version.t()) :: EndpointQuery.t() - defp snapshot_to_endpoint(version) do - snapshot = Map.get(version.meta, "endpoint_snapshot", %{}) - - %EndpointQuery{} - |> Ecto.Changeset.cast(snapshot, EndpointQuery.version_snapshot_fields()) - |> Ecto.Changeset.apply_changes() - end - defp maybe_assign_team_context(socket, %{"t" => _team_id}, _endpoint), do: socket defp maybe_assign_team_context(socket, _params, endpoint) do diff --git a/test/logflare/endpoints/cache_test.exs b/test/logflare/endpoints/cache_test.exs index 7c0a3e860d..2091328f12 100644 --- a/test/logflare/endpoints/cache_test.exs +++ b/test/logflare/endpoints/cache_test.exs @@ -1,6 +1,7 @@ defmodule Logflare.Endpoints.CacheTest do use Logflare.DataCase + alias Logflare.Backends.Adaptor.ClickHouseAdaptor alias Logflare.Backends.QueryError alias Logflare.Endpoints @@ -29,6 +30,17 @@ defmodule Logflare.Endpoints.CacheTest do %{user: user, endpoint: endpoint, endpoint_2: endpoint_2} end + setup context do + if context[:clickhouse_cache] do + {_source, backend} = setup_clickhouse_test(user: context.user) + start_supervised!({ClickHouseAdaptor, backend}) + + %{clickhouse_backend: backend} + else + :ok + end + end + test "cache starts and serves cached results", %{endpoint: endpoint} do # Mock response by setting up test backend test_response = [%{"testing" => "123"}] @@ -49,6 +61,151 @@ defmodule Logflare.Endpoints.CacheTest do assert {:ok, %{rows: [%{"testing" => "123"}]}} = Endpoints.run_cached_query(endpoint) end + @tag :clickhouse_cache + test "cache separates current and versioned endpoint results", %{ + user: user, + clickhouse_backend: backend + } do + endpoint = + insert(:endpoint, + user: user, + backend: backend, + language: :ch_sql, + query: "SELECT 'current' AS testing", + cache_duration_seconds: 60, + proactive_requerying_seconds: 60 + ) + + insert(:endpoint_version, + endpoint: endpoint, + version_number: 1, + snapshot_overrides: %{"query" => "SELECT 'historical' AS testing"} + ) + + assert {:ok, versioned_endpoint} = Endpoints.get_endpoint_query_at_version(endpoint, 1) + + assert {:ok, %{rows: [%{"testing" => "current"}]}} = Endpoints.run_cached_query(endpoint) + + assert {:ok, %{rows: [%{"testing" => "historical"}]}} = + Endpoints.run_cached_query(versioned_endpoint) + + assert {:ok, %{rows: [%{"testing" => "current"}]}} = Endpoints.run_cached_query(endpoint) + + assert {:ok, %{rows: [%{"testing" => "historical"}]}} = + Endpoints.run_cached_query(versioned_endpoint) + end + + @tag :clickhouse_cache + test "versioned cache refresh keeps running the selected snapshot", %{ + user: user, + clickhouse_backend: backend + } do + endpoint = + insert(:endpoint, + user: user, + backend: backend, + language: :ch_sql, + query: "SELECT 'current' AS testing", + cache_duration_seconds: 60, + proactive_requerying_seconds: 1 + ) + + insert(:endpoint_version, + endpoint: endpoint, + version_number: 1, + snapshot_overrides: %{ + "query" => "SELECT concat('historical-', toString(generateUUIDv4())) AS testing", + "cache_duration_seconds" => 60, + "proactive_requerying_seconds" => 1 + } + ) + + assert {:ok, versioned_endpoint} = Endpoints.get_endpoint_query_at_version(endpoint, 1) + + assert {:ok, %{rows: [%{"testing" => first_value}]}} = + Endpoints.run_cached_query(versioned_endpoint) + + assert String.starts_with?(first_value, "historical-") + + endpoint_id = endpoint.id + + Logflare.Repo.update_all( + from(endpoint_query in Endpoints.EndpointQuery, where: endpoint_query.id == ^endpoint_id), + set: [query: "SELECT 'current updated' AS testing"] + ) + + Process.sleep(versioned_endpoint.proactive_requerying_seconds * 1000 + 100) + + TestUtils.retry_assert(fn -> + assert {:ok, %{rows: [%{"testing" => second_value}]}} = + Endpoints.run_cached_query(versioned_endpoint) + + assert String.starts_with?(second_value, "historical-") + assert second_value != first_value + end) + + versioned_endpoint + |> Endpoints.ResultsCache.name(%{}) + |> GenServer.whereis() + |> Endpoints.ResultsCache.invalidate() + end + + @tag :clickhouse_cache + test "endpoint updates invalidate latest caches without touching versioned caches", %{ + user: user, + clickhouse_backend: backend + } do + endpoint = + insert(:endpoint, + user: user, + backend: backend, + language: :ch_sql, + query: "SELECT 'current' AS testing", + cache_duration_seconds: 60, + proactive_requerying_seconds: 60 + ) + + insert(:endpoint_version, + endpoint: endpoint, + version_number: 1, + snapshot_overrides: %{ + "query" => "SELECT 'historical' AS testing", + "cache_duration_seconds" => 60, + "proactive_requerying_seconds" => 60 + } + ) + + assert {:ok, versioned_endpoint} = Endpoints.get_endpoint_query_at_version(endpoint, 1) + + assert {:ok, %{rows: [%{"testing" => "current"}]}} = Endpoints.run_cached_query(endpoint) + + assert {:ok, %{rows: [%{"testing" => "historical"}]}} = + Endpoints.run_cached_query(versioned_endpoint) + + latest_cache_pid = + endpoint + |> Endpoints.ResultsCache.name(%{}) + |> GenServer.whereis() + + versioned_cache_pid = + versioned_endpoint + |> Endpoints.ResultsCache.name(%{}) + |> GenServer.whereis() + + assert is_pid(latest_cache_pid) + assert is_pid(versioned_cache_pid) + + assert {:ok, _updated_endpoint} = + Endpoints.update_query(endpoint, %{query: "select 'updated' as testing"}, user) + + TestUtils.retry_assert(fn -> + refute Process.alive?(latest_cache_pid) + end) + + assert Process.alive?(versioned_cache_pid) + Endpoints.ResultsCache.invalidate(versioned_cache_pid) + end + test "cache dies on timeout error from query", %{endpoint: endpoint} do GoogleApi.BigQuery.V2.Api.Jobs |> expect(:bigquery_jobs_query, 1, fn _conn, _proj_id, _opts -> diff --git a/test/logflare/endpoints_test.exs b/test/logflare/endpoints_test.exs index d51e7448f1..702f485313 100644 --- a/test/logflare/endpoints_test.exs +++ b/test/logflare/endpoints_test.exs @@ -363,6 +363,68 @@ defmodule Logflare.EndpointsTest do assert nil == Endpoints.get_endpoint_query_version(endpoint.id, 3) assert nil == Endpoints.get_endpoint_query_version(endpoint.id + 1, 1) end + + test "get_endpoint_query_at_version/2 returns the endpoint snapshot", %{ + user: user + } do + source = insert(:source, user: user, name: "old_table") + + assert {:ok, endpoint} = + Endpoints.create_query( + user, + %{ + name: "versioned-runnable", + query: "select a from old_table", + language: :bq_sql, + labels: "environment", + redact_pii: true, + enable_dynamic_reservation: true, + cache_duration_seconds: 123, + proactive_requerying_seconds: 45 + }, + user + ) + + assert {:ok, updated_endpoint} = + Endpoints.update_query(endpoint, %{query: "select b from old_table"}, user) + + source + |> Ecto.Changeset.change(name: "new_table") + |> Logflare.Repo.update!() + + endpoint_id = updated_endpoint.id + token = updated_endpoint.token + user_id = user.id + + assert {:ok, + %EndpointQuery{ + id: ^endpoint_id, + user_id: ^user_id, + token: ^token, + version_number: 1, + labels: "environment", + redact_pii: true, + enable_dynamic_reservation: true, + cache_duration_seconds: 123, + proactive_requerying_seconds: 45, + query: query + }} = + Endpoints.get_endpoint_query_at_version(updated_endpoint, 1) + + assert String.downcase(query) == "select a from new_table" + end + + test "get_endpoint_query_at_version/2 rejects invalid version numbers", %{ + user: user, + endpoint_params: endpoint_params + } do + assert {:ok, endpoint} = Endpoints.create_query(user, endpoint_params, user) + + for version <- [0, -1, "1", nil] do + assert {:error, :invalid_version} = + Endpoints.get_endpoint_query_at_version(endpoint, version) + end + end end test "parse_query_string/1" do @@ -520,9 +582,9 @@ defmodule Logflare.EndpointsTest do endpoint_id_label_value = Integer.to_string(endpoint_id) - assert_received %{ - "endpoint_id" => ^endpoint_id_label_value - } + assert_received labels + assert labels["endpoint_id"] == endpoint_id_label_value + refute Map.has_key?(labels, "endpoint_version") end test "run an endpoint query with query composition" do @@ -563,7 +625,9 @@ defmodule Logflare.EndpointsTest do ) assert {:ok, %{rows: [%{"testing" => _}]}} = Endpoints.run_query(endpoint) - assert_received %{"my_label" => "my_value"} + assert_received labels + assert labels["my_label"] == "my_value" + refute Map.has_key?(labels, "endpoint_version") end test "run_query_string/3" do diff --git a/test/logflare_web/controllers/endpoints_controller_test.exs b/test/logflare_web/controllers/endpoints_controller_test.exs index 59192f146b..d8a8dbf5f1 100644 --- a/test/logflare_web/controllers/endpoints_controller_test.exs +++ b/test/logflare_web/controllers/endpoints_controller_test.exs @@ -5,6 +5,7 @@ defmodule LogflareWeb.EndpointsControllerTest do alias Logflare.Backends alias Logflare.Backends.Adaptor.PostgresAdaptor.PgRepo alias Logflare.Backends.Adaptor.PostgresAdaptor.SharedRepo + alias Logflare.Endpoints alias Logflare.Google.BigQuery.GenUtils alias Logflare.SingleTenant alias Logflare.Sources @@ -394,6 +395,156 @@ defmodule LogflareWeb.EndpointsControllerTest do end end + describe "lf-endpoint-version" do + setup do + insert(:plan, name: "Free") + user = insert(:user) + {_source, backend} = Logflare.DataCase.setup_clickhouse_test(user: user) + + assert {:ok, endpoint} = + Endpoints.create_query( + user, + %{ + name: "versioned-clickhouse-endpoint", + query: "select 'historical' as versioned_value", + backend_id: backend.id, + cache_duration_seconds: 0, + enable_auth: false, + sandboxable: false, + labels: "endpoint_version=caller" + }, + user + ) + + assert {:ok, endpoint} = + Endpoints.update_query( + endpoint, + %{ + query: "select 'current' as versioned_value", + sandboxable: true + }, + user + ) + + {:ok, user: user, endpoint: endpoint} + end + + test "runs the requested endpoint version", %{ + conn: init_conn, + endpoint: endpoint, + user: user + } do + conn = + init_conn + |> put_req_header("lf-endpoint-version", "1") + |> get(~p"/endpoints/query/#{endpoint.token}") + + assert %{"result" => [%{"versioned_value" => "historical"}]} = json_response(conn, 200) + + conn = + init_conn + |> put_req_header("x-api-key", user.api_key) + |> put_req_header("lf-endpoint-version", "1") + |> get(~p"/api/endpoints/query/#{endpoint.name}") + + assert %{"result" => [%{"versioned_value" => "historical"}]} = json_response(conn, 200) + + conn = + init_conn + |> put_req_header("x-api-key", user.api_key) + |> put_req_header("lf-endpoint-version", "1") + |> post(~p"/api/endpoints/query/#{endpoint.name}", %{}) + + assert %{"result" => [%{"versioned_value" => "historical"}]} = json_response(conn, 200) + end + + test "uses version sandboxability when deciding whether to parse body", %{ + conn: init_conn, + endpoint: endpoint + } do + conn = + init_conn + |> put_req_header("content-type", "application/json") + |> put_req_header("lf-endpoint-version", "1") + |> get( + ~p"/endpoints/query/#{endpoint.token}", + Jason.encode!(%{sql: "select 'override' as versioned_value"}) + ) + + assert %{"result" => [%{"versioned_value" => "historical"}]} = json_response(conn, 200) + refute conn.halted + end + + test "returns version not found for a missing version", %{ + conn: init_conn, + endpoint: endpoint + } do + conn = + init_conn + |> put_req_header("lf-endpoint-version", "99") + |> get(~p"/endpoints/query/#{endpoint.token}") + + assert %{"error" => "version not found"} = json_response(conn, 404) + + conn = + init_conn + |> put_req_header("lf-endpoint-version", "2147483648") + |> get(~p"/endpoints/query/#{endpoint.token}") + + assert %{"error" => "version not found"} = json_response(conn, 404) + end + + test "returns a client error for an invalid version", %{ + conn: init_conn, + endpoint: endpoint + } do + for version <- ["latest", "0", "-1"] do + conn = + init_conn + |> put_req_header("lf-endpoint-version", version) + |> get(~p"/endpoints/query/#{endpoint.token}") + + assert %{"error" => "invalid lf-endpoint-version"} = json_response(conn, 400) + end + end + + test "adds the requested endpoint version to query telemetry labels", %{ + conn: init_conn, + endpoint: endpoint, + user: user + } do + test_pid = self() + handler_id = "test-endpoint-version-label-#{inspect(self())}" + + :telemetry.attach( + handler_id, + [:logflare, :endpoints, :query], + fn _event, _measurements, metadata, _config -> + send(test_pid, {:endpoint_query_metadata, metadata}) + end, + nil + ) + + on_exit(fn -> :telemetry.detach(handler_id) end) + + conn = + init_conn + |> put_req_header("x-api-key", user.api_key) + |> put_req_header("content-type", "application/json") + |> put_req_header("lf-endpoint-version", "1") + |> get(~p"/api/endpoints/query/#{endpoint.name}") + + assert [_] = json_response(conn, 200)["result"] + assert conn.halted == false + + assert_received {:endpoint_query_metadata, + %{ + "endpoint_id" => _endpoint_id, + "endpoint_version" => "1" + }} + end + end + describe "bigquery with labels" do setup do _plan = insert(:plan, name: "Free")