Skip to content
Merged
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
15 changes: 15 additions & 0 deletions lib/logflare/backends/adaptor/clickhouse_adaptor.ex
Original file line number Diff line number Diff line change
Expand Up @@ -785,6 +785,21 @@ defmodule Logflare.Backends.Adaptor.ClickHouseAdaptor do
end

@spec handle_read_pool_log(DBConnection.LogEntry.t(), pos_integer(), String.t() | nil) :: :ok
defp handle_read_pool_log(
%DBConnection.LogEntry{
result: {:error, %DBConnection.ConnectionError{reason: reason}},
connection_time: nil
},
backend_id,
label
) do
:telemetry.execute(
[:logflare, :clickhouse, :read_pool, :checkout_error],
%{count: 1},
%{backend_id: backend_id, read_cluster: read_cluster_tag(label), reason: reason}
)
end

defp handle_read_pool_log(%DBConnection.LogEntry{} = entry, backend_id, label) do
metadata = %{backend_id: backend_id, read_cluster: read_cluster_tag(label)}

Expand Down
9 changes: 8 additions & 1 deletion lib/telemetry.ex
Original file line number Diff line number Diff line change
Expand Up @@ -385,7 +385,7 @@ defmodule Logflare.Telemetry do
measurement: :count,
tags: [:backend_id, :read_cluster, :error_kind],
description:
"ClickHouse read queries that resolved to an error, excluding invalid-query (user SQL) errors, counted once per query after any failover retry and tagged with the read cluster that produced the final error"
"ClickHouse read queries that resolved to an error, excluding invalid-query (user SQL) errors, counted once per query after any failover retry and tagged with the read cluster that produced the final error. Checkout failures are additionally counted per attempt in `checkout_error`; `connection_error` here covers both failed checkouts and connections lost mid-query"
),
sum("logflare.clickhouse.read_pool.failover",
event_name: [:logflare, :clickhouse, :read_pool, :failover],
Expand All @@ -412,6 +412,13 @@ defmodule Logflare.Telemetry do
description:
"ClickHouse read pool connections lost or recycled, per backend and read cluster. Paired with `connected`, this is the pool's connection churn rate"
),
sum("logflare.clickhouse.read_pool.checkout_error",
event_name: [:logflare, :clickhouse, :read_pool, :checkout_error],
measurement: :count,
tags: [:backend_id, :read_cluster, :reason],
description:
"ClickHouse read pool checkouts that failed before a connection was obtained, counted once per checkout attempt (including attempts masked by a successful failover) and tagged with the read cluster that shed it. `reason` mirrors `DBConnection.ConnectionError.reason`: `:queue_timeout` when the request was dropped from a saturated pool's queue, `:error` for other checkout failures such as the deadline expiring while queued. A `queue_timeout` here also surfaces in `query_error` as `pool_exhausted`; do not sum the two series"
),
sum("logflare.clickhouse.insert.result.count",
event_name: [:logflare, :clickhouse, :insert, :result],
measurement: :count,
Expand Down
116 changes: 116 additions & 0 deletions test/logflare/backends/adaptor/clickhouse_adaptor_test.exs
Original file line number Diff line number Diff line change
Expand Up @@ -214,6 +214,107 @@ defmodule Logflare.Backends.Adaptor.ClickHouseAdaptorTest do
assert measurements.idle_time < minute_in_native
end

test "emits checkout error telemetry when the pool sheds a checkout", %{backend: backend} do
TestUtils.attach_forwarder([:logflare, :clickhouse, :read_pool, :checkout_error])

error = %DBConnection.ConnectionError{message: "dropped", reason: :queue_timeout}
stub_failed_checkout(error)

assert {:error, %QueryError{kind: :pool_exhausted}} =
ClickHouseAdaptor.execute_ch_query(backend, "SELECT 1 as test")

assert_receive {:telemetry_event, [:logflare, :clickhouse, :read_pool, :checkout_error],
%{count: 1}, metadata}

assert metadata.backend_id == backend.id
assert metadata.read_cluster == "(unlabeled)"
assert metadata.reason == :queue_timeout
end

test "tags non-queue checkout failures with the error reason", %{backend: backend} do
TestUtils.attach_forwarder([:logflare, :clickhouse, :read_pool, :checkout_error])

error = %DBConnection.ConnectionError{
message: "connection not available because deadline reached while in queue"
}

stub_failed_checkout(error)

assert {:error, %QueryError{kind: :connection_error}} =
ClickHouseAdaptor.execute_ch_query(backend, "SELECT 1 as test")

assert_receive {:telemetry_event, [:logflare, :clickhouse, :read_pool, :checkout_error],
%{count: 1}, %{reason: :error}}
end

test "does not record checkout latency for a failed checkout", %{backend: backend} do
TestUtils.attach_forwarder([:logflare, :clickhouse, :read_pool, :checkout])

error = %DBConnection.ConnectionError{message: "dropped", reason: :queue_timeout}
stub_failed_checkout(error)

assert {:error, %QueryError{kind: :pool_exhausted}} =
ClickHouseAdaptor.execute_ch_query(backend, "SELECT 1 as test")

refute_received {:telemetry_event, [:logflare, :clickhouse, :read_pool, :checkout], _, _}
end

test "does not warn about a slow checkout that never got a connection", %{backend: backend} do
original = Application.get_env(:logflare, ClickHouseAdaptor)

on_exit(fn ->
if original do
Application.put_env(:logflare, ClickHouseAdaptor, original)
else
Application.delete_env(:logflare, ClickHouseAdaptor)
end
end)

Application.put_env(:logflare, ClickHouseAdaptor, slow_pool_checkout_ms: 0)

error = %DBConnection.ConnectionError{message: "dropped", reason: :queue_timeout}
stub_failed_checkout(error)

log =
ExUnit.CaptureLog.capture_log(fn ->
assert {:error, %QueryError{kind: :pool_exhausted}} =
ClickHouseAdaptor.execute_ch_query(backend, "SELECT 1 as test")
end)

refute log =~ "ClickHouse slow connection checkout"
end

test "still emits query error telemetry for a connection error after checkout",
%{backend: backend} do
TestUtils.attach_forwarder([:logflare, :clickhouse, :read_pool, :checkout_error])
TestUtils.attach_forwarder([:logflare, :clickhouse, :read_pool, :checkout])

error = %DBConnection.ConnectionError{message: "socket closed"}

expect(Ch, :query, fn _pool, statement, params, opts ->
entry = %DBConnection.LogEntry{
call: :execute,
query: statement,
params: params,
result: {:error, error},
pool_time: System.convert_time_unit(1, :millisecond, :native),
connection_time: System.convert_time_unit(5, :millisecond, :native)
}

opts[:log].(entry)
{:error, error}
end)

assert {:error, %QueryError{kind: :connection_error}} =
ClickHouseAdaptor.execute_ch_query(backend, "SELECT 1 as test")

refute_received {:telemetry_event, [:logflare, :clickhouse, :read_pool, :checkout_error], _,
_}

assert_receive {:telemetry_event, [:logflare, :clickhouse, :read_pool, :checkout],
%{connection_time: _}, _}
end

test "emits query error telemetry with the error kind", %{backend: backend} do
TestUtils.attach_forwarder([:logflare, :clickhouse, :read_pool, :query_error])

Expand Down Expand Up @@ -3118,6 +3219,21 @@ defmodule Logflare.Backends.Adaptor.ClickHouseAdaptorTest do
end)
end

defp stub_failed_checkout(%DBConnection.ConnectionError{} = error) do
expect(Ch, :query, fn _pool, statement, params, opts ->
entry = %DBConnection.LogEntry{
call: :execute,
query: statement,
params: params,
result: {:error, error},
pool_time: System.convert_time_unit(12_000, :millisecond, :native)
}

opts[:log].(entry)
{:error, error}
end)
end

defp modify_backend_with_long_token(%Backend{} = backend) do
long_token = random_string(200)

Expand Down
9 changes: 9 additions & 0 deletions test/logflare/telemetry_test.exs
Original file line number Diff line number Diff line change
Expand Up @@ -279,6 +279,15 @@ defmodule Logflare.TelemetryTest do
end
end

test "defines ClickHouse read pool checkout failures tagged by reason" do
metric = read_pool_metric([:logflare, :clickhouse, :read_pool, :checkout_error])

assert to_string(metric.__struct__) == "Elixir.Telemetry.Metrics.Sum"
assert metric.event_name == [:logflare, :clickhouse, :read_pool, :checkout_error]
assert metric.measurement == :count
assert metric.tags == [:backend_id, :read_cluster, :reason]
end

test "honors configured Broadway processor message duration sampling" do
denominator = Application.fetch_env!(:logflare, :broadway_message_sample_denominator)

Expand Down
Loading