diff --git a/lib/logflare/backends/adaptor/bigquery_adaptor.ex b/lib/logflare/backends/adaptor/bigquery_adaptor.ex index f5fbfc595f..6371986ea1 100644 --- a/lib/logflare/backends/adaptor/bigquery_adaptor.ex +++ b/lib/logflare/backends/adaptor/bigquery_adaptor.ex @@ -11,16 +11,16 @@ defmodule Logflare.Backends.Adaptor.BigQueryAdaptor do require OpenTelemetry.Tracer alias Ecto.Changeset + alias GoogleApi.BigQuery.V2.Model alias GoogleApi.IAM.V1.Api.Projects, as: IAMProjects alias GoogleApi.IAM.V1.Model.CreateServiceAccountRequest - alias GoogleApi.BigQuery.V2.Model alias Logflare.Backends alias Logflare.Backends.Adaptor.BigQueryAdaptor.GoogleApiClient + alias Logflare.Backends.Adaptor.QueryResult alias Logflare.Backends.Backend alias Logflare.Backends.DynamicPipeline alias Logflare.Backends.Ecto.SqlUtils alias Logflare.Backends.IngestEventQueue - alias Logflare.Backends.Adaptor.QueryResult alias Logflare.Backends.QueryError alias Logflare.BigQuery.SchemaTypes alias Logflare.Billing @@ -31,7 +31,9 @@ defmodule Logflare.Backends.Adaptor.BigQueryAdaptor do alias Logflare.Google.BigQuery.GCPConfig alias Logflare.Google.BigQuery.GenUtils alias Logflare.Google.CloudResourceManager + alias Logflare.Rules alias Logflare.Sources + alias Logflare.Sources.Source alias Logflare.Sources.Source.BigQuery.Pipeline alias Logflare.Sources.Source.BigQuery.Schema alias Logflare.User @@ -45,6 +47,16 @@ defmodule Logflare.Backends.Adaptor.BigQueryAdaptor do @timeout_error_regex ~r/timed out/i @search_query_timeout_ms 60_000 @endpoint_query_timeout_ms 60_000 + @connection_test_message "Logflare BigQuery connection test. No action required." + + @type connection_test_error :: + :connection_error + | :http_client_error + | :http_server_error + | :insert_error + | :invalid_config + | :source_required + | :unknown_error @impl Logflare.Backends.Adaptor def start_link({source, backend} = source_backend) do @@ -161,8 +173,184 @@ defmodule Logflare.Backends.Adaptor.BigQueryAdaptor do |> String.replace("-", "_") end + @doc """ + Writes a labeled event through the REST ingestion path to an attached source's table. + + Direct source attachments are preferred over rule sources. + """ @impl Logflare.Backends.Adaptor - def test_connection(_), do: {:error, :not_implemented} + @spec test_connection(Backend.t()) :: :ok | {:error, connection_test_error()} + def test_connection( + %Backend{ + id: backend_id, + user_id: user_id, + config: %{project_id: project_id, dataset_id: dataset_id} + } = backend + ) + when is_integer(backend_id) and is_integer(user_id) and is_non_empty_binary(project_id) and + is_non_empty_binary(dataset_id) do + with {:ok, %Source{token: source_token}} <- connection_test_source(backend) do + Google.BigQuery.stream_batch!( + %{ + bigquery_project_id: project_id, + bigquery_dataset_id: dataset_id, + source_token: source_token + }, + [connection_test_row()] + ) + |> handle_connection_test_response(backend, source_token) + end + end + + def test_connection(%Backend{}), do: {:error, :invalid_config} + + @spec connection_test_source(Backend.t()) :: {:ok, Source.t()} | {:error, :source_required} + defp connection_test_source(%Backend{id: backend_id, user_id: user_id}) do + direct_source = + [backend_id: backend_id, user_id: user_id] + |> Sources.list_sources() + |> Enum.min_by(& &1.id, fn -> nil end) + + case direct_source do + %Source{} = source -> + {:ok, source} + + nil -> + connection_test_rule_source(backend_id, user_id) + end + end + + @spec connection_test_rule_source(pos_integer(), pos_integer()) :: + {:ok, Source.t()} | {:error, :source_required} + defp connection_test_rule_source(backend_id, user_id) do + rule = + user_id + |> Rules.list_rules_by_user_id(backend_id) + |> Enum.min_by(& &1.source_id, fn -> nil end) + + case rule do + %{source_id: source_id} -> + case Sources.fetch_source_by(id: source_id, user_id: user_id) do + {:ok, %Source{} = source} -> {:ok, source} + {:error, :not_found} -> {:error, :source_required} + end + + nil -> + {:error, :source_required} + end + end + + @spec connection_test_row() :: Model.TableDataInsertAllRequestRows.t() + defp connection_test_row do + id = Ecto.UUID.generate() + + %Model.TableDataInsertAllRequestRows{ + insertId: id, + json: %{ + "event_message" => @connection_test_message, + "id" => id, + "timestamp" => DateTime.utc_now() + } + } + end + + @spec handle_connection_test_response(term(), Backend.t(), atom()) :: + :ok | {:error, connection_test_error()} + defp handle_connection_test_response( + {:ok, %Model.TableDataInsertAllResponse{insertErrors: insert_errors}}, + _backend, + _source_token + ) + when insert_errors in [nil, []], + do: :ok + + defp handle_connection_test_response( + {:ok, %Model.TableDataInsertAllResponse{insertErrors: insert_errors}}, + backend, + source_token + ) do + connection_test_failed(backend, source_token, :insert_error, insert_errors) + end + + defp handle_connection_test_response( + {:error, %Tesla.Env{status: status, body: body}}, + backend, + source_token + ) + when status in 400..499 do + connection_test_failed(backend, source_token, :http_client_error, %{ + status: status, + body: body + }) + end + + defp handle_connection_test_response( + {:error, %Tesla.Env{status: status, body: body}}, + backend, + source_token + ) + when status in 500..599 do + connection_test_failed(backend, source_token, :http_server_error, %{ + status: status, + body: body + }) + end + + defp handle_connection_test_response( + {:error, %Tesla.Env{status: status, body: body}}, + backend, + source_token + ) do + connection_test_failed(backend, source_token, :unknown_error, %{status: status, body: body}) + end + + defp handle_connection_test_response({:error, error}, backend, source_token) + when error in [:closed, :emfile, :timeout] do + connection_test_failed(backend, source_token, :connection_error, error) + end + + defp handle_connection_test_response({:error, error}, backend, source_token) do + connection_test_failed(backend, source_token, :unknown_error, error) + end + + defp handle_connection_test_response(response, backend, source_token) do + connection_test_failed(backend, source_token, :unknown_error, response) + end + + @spec connection_test_failed( + Backend.t(), + atom(), + connection_test_error(), + term() + ) :: {:error, connection_test_error()} + defp connection_test_failed(backend, source_token, reason, error) do + Logger.warning("BigQuery connection test failed.", + backend_id: backend.id, + user_id: backend.user_id, + table_name: format_table_name(source_token), + reason: reason, + error_string: format_connection_test_error(error) + ) + + {:error, reason} + end + + @spec format_connection_test_error(term()) :: String.t() + defp format_connection_test_error(%{status: status, body: body}) do + inspect(%{status: status, body: body}, limit: 20, printable_limit: 2_000) + end + + defp format_connection_test_error(errors) when is_list(errors) do + inspect(errors, limit: 20, printable_limit: 2_000) + end + + defp format_connection_test_error(error) when is_atom(error), do: inspect(error) + + defp format_connection_test_error(%{__struct__: module}) do + "unexpected #{inspect(module)}" + end + + defp format_connection_test_error(_error), do: "unexpected response" @impl Logflare.Backends.Adaptor def ecto_to_sql(%Ecto.Query{} = query, _opts) do diff --git a/lib/logflare_web/live/backends/components.ex b/lib/logflare_web/live/backends/components.ex index 2b6cc507f8..066694dc7f 100644 --- a/lib/logflare_web/live/backends/components.ex +++ b/lib/logflare_web/live/backends/components.ex @@ -33,6 +33,11 @@ defmodule LogflareWeb.Backends.Components do LogflareWeb.QueryErrorHelpers.query_error_message(query_error) end + defp status_error_message(reason) + when reason in [:source_required, {:error, :source_required}] do + "Attach a source or add a drain rule before testing this connection." + end + defp status_error_message(_reason) do LogflareWeb.QueryErrorHelpers.generic_query_error_message() end diff --git a/test/logflare/backends/adaptor/bigquery_adaptor_test.exs b/test/logflare/backends/adaptor/bigquery_adaptor_test.exs index 891b48bee6..01eafebcef 100644 --- a/test/logflare/backends/adaptor/bigquery_adaptor_test.exs +++ b/test/logflare/backends/adaptor/bigquery_adaptor_test.exs @@ -5,10 +5,18 @@ defmodule Logflare.Backends.Adaptor.BigQueryAdaptorTest do import Ecto.Query alias GoogleApi.BigQuery.V2.Api.Jobs, as: BqJobs - alias Logflare.Backends.Backend + alias GoogleApi.BigQuery.V2.Api.Tabledata, as: BqTabledata + alias GoogleApi.BigQuery.V2.Api.Tables, as: BqTables + alias GoogleApi.BigQuery.V2.Model + alias Logflare.Backends alias Logflare.Backends.Adaptor.BigQueryAdaptor + alias Logflare.Backends.Adaptor.BigQueryAdaptor.GoogleApiClient alias Logflare.Backends.Adaptor.QueryResult + alias Logflare.Backends.Backend alias Logflare.Backends.QueryError + alias Logflare.Google + + @connection_test_message "Logflare BigQuery connection test. No action required." # Characters illegal in a BigQuery dataset identifier: SQL delimiters, # identifier-quoting characters, whitespace, and shell metacharacters. @@ -60,6 +68,205 @@ defmodule Logflare.Backends.Adaptor.BigQueryAdaptorTest do end end + describe "test_connection/1" do + setup do + insert(:plan) + user = insert(:user) + source = insert(:source, user: user, bq_storage_write_api: true) + + backend = + insert(:backend, + type: :bigquery, + user: user, + sources: [source], + config: %{project_id: "test-project", dataset_id: "test_dataset"} + ) + + [backend: Backends.get_backend(backend.id), source: source, user: user] + end + + test "writes one labeled event to the attached source table through REST", %{ + backend: backend, + source: source + } do + Mimic.reject(BqTables, :bigquery_tables_get, 4) + Mimic.reject(Google.BigQuery, :create_table, 4) + Mimic.reject(GoogleApiClient, :append_rows, 3) + + BqTabledata + |> expect(:bigquery_tabledata_insert_all, fn _conn, + project_id, + dataset_id, + table_name, + opts -> + assert project_id == "test-project" + assert dataset_id == "test_dataset" + assert table_name == BigQueryAdaptor.format_table_name(source.token) + + assert %Model.TableDataInsertAllRequest{ + ignoreUnknownValues: true, + skipInvalidRows: true, + rows: [%Model.TableDataInsertAllRequestRows{} = row] + } = opts[:body] + + assert row.insertId == row.json["id"] + assert is_binary(row.insertId) + assert %DateTime{} = row.json["timestamp"] + assert row.json["event_message"] == @connection_test_message + + {:ok, %Model.TableDataInsertAllResponse{insertErrors: nil}} + end) + + assert :ok = BigQueryAdaptor.test_connection(backend) + end + + test "uses a rule source when no source is attached directly", %{user: user} do + source = insert(:source, user: user) + + backend = + insert(:backend, + type: :bigquery, + user: user, + config: %{project_id: "rule-project", dataset_id: "rule_dataset"} + ) + + insert(:rule, backend: backend, source: source) + + Google.BigQuery + |> expect(:stream_batch!, fn context, [_row] -> + assert context == %{ + bigquery_project_id: "rule-project", + bigquery_dataset_id: "rule_dataset", + source_token: source.token + } + + {:ok, %Model.TableDataInsertAllResponse{insertErrors: nil}} + end) + + assert :ok = BigQueryAdaptor.test_connection(Backends.get_backend(backend.id)) + end + + test "prefers a direct source over an older rule source", %{user: user} do + rule_source = insert(:source, user: user) + direct_source = insert(:source, user: user) + + backend = + insert(:backend, + type: :bigquery, + user: user, + sources: [direct_source], + config: %{project_id: "test-project", dataset_id: "test_dataset"} + ) + + insert(:rule, backend: backend, source: rule_source) + + Google.BigQuery + |> expect(:stream_batch!, fn %{source_token: source_token}, [_row] -> + assert source_token == direct_source.token + {:ok, %Model.TableDataInsertAllResponse{insertErrors: nil}} + end) + + assert rule_source.id < direct_source.id + assert :ok = BigQueryAdaptor.test_connection(Backends.get_backend(backend.id)) + end + + test "selects the oldest directly attached source deterministically", %{user: user} do + first_source = insert(:source, user: user) + second_source = insert(:source, user: user) + + backend = + insert(:backend, + type: :bigquery, + user: user, + sources: [second_source, first_source], + config: %{project_id: "test-project", dataset_id: "test_dataset"} + ) + + Google.BigQuery + |> expect(:stream_batch!, fn %{source_token: source_token}, [_row] -> + assert source_token == first_source.token + {:ok, %Model.TableDataInsertAllResponse{insertErrors: nil}} + end) + + assert first_source.id < second_source.id + assert :ok = BigQueryAdaptor.test_connection(Backends.get_backend(backend.id)) + end + + test "requires an attached direct or rule source", %{user: user} do + backend = + insert(:backend, + type: :bigquery, + user: user, + config: %{project_id: "test-project", dataset_id: "test_dataset"} + ) + + Mimic.reject(BqTables, :bigquery_tables_get, 4) + Mimic.reject(Google.BigQuery, :create_table, 4) + Mimic.reject(Google.BigQuery, :stream_batch!, 2) + + assert {:error, :source_required} = + BigQueryAdaptor.test_connection(Backends.get_backend(backend.id)) + end + + test "ignores source associations owned by another user", %{user: user} do + other_user = insert(:user) + direct_source = insert(:source, user: other_user) + rule_source = insert(:source, user: other_user) + + backend = + insert(:backend, + type: :bigquery, + user: user, + sources: [direct_source], + config: %{project_id: "test-project", dataset_id: "test_dataset"} + ) + + insert(:rule, backend: backend, source: rule_source) + + Mimic.reject(Google.BigQuery, :stream_batch!, 2) + + assert {:error, :source_required} = + BigQueryAdaptor.test_connection(Backends.get_backend(backend.id)) + end + + test "normalizes REST insert responses to atom reasons", %{backend: backend} do + responses = [ + {{:ok, %Model.TableDataInsertAllResponse{insertErrors: []}}, :ok}, + {{:ok, %Model.TableDataInsertAllResponse{insertErrors: [%{index: 0}]}}, :insert_error}, + {{:error, %Tesla.Env{status: 404, body: %{"error" => "not found"}}}, :http_client_error}, + {{:error, %Tesla.Env{status: 403, body: %{"error" => "forbidden"}}}, :http_client_error}, + {{:error, %Tesla.Env{status: 503, body: %{"error" => "unavailable"}}}, + :http_server_error}, + {{:error, :timeout}, :connection_error}, + {{:error, :unexpected}, :unknown_error}, + {{:ok, :unexpected}, :unknown_error} + ] + + for {response, expected_reason} <- responses do + Google.BigQuery + |> expect(:stream_batch!, fn _context, [_row] -> response end) + + expected = if expected_reason == :ok, do: :ok, else: {:error, expected_reason} + assert BigQueryAdaptor.test_connection(backend) == expected + end + end + + test "rejects incomplete configuration without making a request", %{backend: backend} do + Mimic.reject(BqTables, :bigquery_tables_get, 4) + Mimic.reject(Google.BigQuery, :create_table, 4) + Mimic.reject(Google.BigQuery, :stream_batch!, 2) + + for config <- [ + %{}, + %{project_id: "", dataset_id: "test_dataset"}, + %{project_id: "test-project", dataset_id: nil} + ] do + assert {:error, :invalid_config} = + BigQueryAdaptor.test_connection(%{backend | config: config}) + end + end + end + describe "ecto_to_sql/2" do test "converts Ecto query to BigQuery SQL format" do query = diff --git a/test/logflare_web/live/backends/components_test.exs b/test/logflare_web/live/backends/components_test.exs index 878f23257a..d5a4710eb6 100644 --- a/test/logflare_web/live/backends/components_test.exs +++ b/test/logflare_web/live/backends/components_test.exs @@ -37,6 +37,15 @@ defmodule LogflareWeb.Backends.ComponentsTest do refute html =~ "fa-check" end + test "explains that BigQuery connection tests require a source" do + result = AsyncResult.failed(AsyncResult.loading(), :source_required) + html = render_component(&Components.status_indicator/1, %{status: result}) + + assert html =~ "fa-times" + assert html =~ "Attach a source or add a drain rule before testing this connection." + refute html =~ "Backend error!" + end + test "renders a check icon when the async result is ok" do result = AsyncResult.ok(AsyncResult.loading(), :connected) html = render_component(&Components.status_indicator/1, %{status: result})