From 5679462ea332f8b33fa0fee679498f8f1af56803 Mon Sep 17 00:00:00 2001 From: Peter Solnica Date: Fri, 24 Jul 2026 10:59:21 +0000 Subject: [PATCH] feat: support for log_byte and trace_metric_byte --- lib/sentry/client_report/sender.ex | 50 +++- lib/sentry/envelope.ex | 40 +++ lib/sentry/logger_handler/logs_backend.ex | 4 +- lib/sentry/metrics.ex | 4 +- lib/sentry/telemetry/buffer.ex | 26 +- lib/sentry/telemetry/category.ex | 71 ++++++ lib/sentry/telemetry/scheduler.ex | 12 +- lib/sentry/telemetry_processor.ex | 10 +- lib/sentry/transport.ex | 7 +- test/envelope_test.exs | 41 +++ test/sentry/client_report/sender_test.exs | 88 ++++++- .../telemetry_processor_integration_test.exs | 235 +++++++++++++++++- test/sentry/transport_test.exs | 52 ++++ 13 files changed, 610 insertions(+), 30 deletions(-) diff --git a/lib/sentry/client_report/sender.ex b/lib/sentry/client_report/sender.ex index c688ada2..c3e1de48 100644 --- a/lib/sentry/client_report/sender.ex +++ b/lib/sentry/client_report/sender.ex @@ -6,7 +6,19 @@ defmodule Sentry.ClientReport.Sender do use GenServer - alias Sentry.{Client, ClientReport, Config, Envelope, Transaction} + alias Sentry.{ + Client, + ClientReport, + Config, + Envelope, + LogBatch, + LogEvent, + Metric, + MetricBatch, + Transaction + } + + alias Sentry.Telemetry.Category @send_interval 30_000 @@ -39,6 +51,10 @@ defmodule Sentry.ClientReport.Sender do | Sentry.CheckIn.t() | ClientReport.t() | Sentry.Event.t() + | LogBatch.t() + | LogEvent.t() + | Metric.t() + | MetricBatch.t() | Sentry.Transaction.t() def record_discarded_events(reason, event_items, genserver) when is_list(event_items) do @@ -65,6 +81,38 @@ defmodule Sentry.ClientReport.Sender do [{Envelope.get_data_category(transaction), 1}, {"span", span_count}] end + defp data_categories(%LogEvent{} = log_event) do + [ + {Category.data_category(:log), 1}, + {Category.byte_data_category(:log), Envelope.log_event_byte_size(log_event)} + ] + end + + defp data_categories(%Metric{} = metric) do + [ + {Category.data_category(:metric), 1}, + {Category.byte_data_category(:metric), Envelope.metric_byte_size(metric)} + ] + end + + defp data_categories(%LogBatch{log_events: log_events}) do + bytes = Enum.reduce(log_events, 0, &(Envelope.log_event_byte_size(&1) + &2)) + + [ + {Category.data_category(:log), length(log_events)}, + {Category.byte_data_category(:log), bytes} + ] + end + + defp data_categories(%MetricBatch{metrics: metrics}) do + bytes = Enum.reduce(metrics, 0, &(Envelope.metric_byte_size(&1) + &2)) + + [ + {Category.data_category(:metric), length(metrics)}, + {Category.byte_data_category(:metric), bytes} + ] + end + defp data_categories(item) do [{Envelope.get_data_category(item), 1}] end diff --git a/lib/sentry/envelope.ex b/lib/sentry/envelope.ex index 0e83597c..120cfe11 100644 --- a/lib/sentry/envelope.ex +++ b/lib/sentry/envelope.ex @@ -138,6 +138,46 @@ defmodule Sentry.Envelope do def get_data_category(%LogBatch{}), do: "log_item" def get_data_category(%MetricBatch{}), do: "trace_metric" + @doc """ + Approximates the serialized byte size of a single log event. + + The size is computed by encoding the log event with the same serialization used + when the event is placed in an envelope item (see `Sentry.LogEvent.to_map/1`), + then measuring the byte size of the resulting JSON. Per the client report spec, + a serialized-size approximation like this is acceptable for `log_byte` outcomes. + + Returns `0` if the log event cannot be encoded. + """ + @doc since: "13.4.0" + @spec log_event_byte_size(LogEvent.t()) :: non_neg_integer() + def log_event_byte_size(%LogEvent{} = log_event) do + log_event |> LogEvent.to_map() |> item_byte_size() + end + + @doc """ + Approximates the serialized byte size of a single metric. + + The size is computed by encoding the metric with the same serialization used + when the metric is placed in an envelope item (see `Sentry.Metric.to_map/1`), + then measuring the byte size of the resulting JSON. Per the client report spec, + a serialized-size approximation like this is acceptable for `trace_metric_byte` + outcomes. + + Returns `0` if the metric cannot be encoded. + """ + @doc since: "13.4.0" + @spec metric_byte_size(Metric.t()) :: non_neg_integer() + def metric_byte_size(%Metric{} = metric) do + metric |> Metric.to_map() |> item_byte_size() + end + + defp item_byte_size(map) do + case Sentry.JSON.encode(map, Config.json_library()) do + {:ok, encoded} -> byte_size(encoded) + {:error, _reason} -> 0 + end + end + @doc """ Returns the total number of payload items in the envelope. diff --git a/lib/sentry/logger_handler/logs_backend.ex b/lib/sentry/logger_handler/logs_backend.ex index ad50ee13..84f33dee 100644 --- a/lib/sentry/logger_handler/logs_backend.ex +++ b/lib/sentry/logger_handler/logs_backend.ex @@ -37,8 +37,8 @@ defmodule Sentry.LoggerHandler.LogsBackend do log_event_struct = LogEvent.from_logger_event(log_event, attributes, parameters) case TelemetryProcessor.add(log_event_struct) do - {:ok, {:rate_limited, data_category}} -> - Sentry.ClientReport.Sender.record_discarded_events(:ratelimit_backoff, data_category) + {:ok, {:rate_limited, _data_category}} -> + Sentry.ClientReport.Sender.record_discarded_events(:ratelimit_backoff, [log_event_struct]) :ok -> :ok diff --git a/lib/sentry/metrics.ex b/lib/sentry/metrics.ex index 55b29d69..65e7bc5e 100644 --- a/lib/sentry/metrics.ex +++ b/lib/sentry/metrics.ex @@ -135,8 +135,8 @@ defmodule Sentry.Metrics do metric = Metric.attach_default_attributes(metric) case TelemetryProcessor.add(metric) do - {:ok, {:rate_limited, data_category}} -> - ClientReport.Sender.record_discarded_events(:ratelimit_backoff, data_category) + {:ok, {:rate_limited, _data_category}} -> + ClientReport.Sender.record_discarded_events(:ratelimit_backoff, [metric]) :ok -> :ok diff --git a/lib/sentry/telemetry/buffer.ex b/lib/sentry/telemetry/buffer.ex index b6244e14..7008fd69 100644 --- a/lib/sentry/telemetry/buffer.ex +++ b/lib/sentry/telemetry/buffer.ex @@ -24,6 +24,7 @@ defmodule Sentry.Telemetry.Buffer do alias Sentry.ClientReport alias Sentry.Telemetry.Category + alias Sentry.{LogEvent, Metric} @enforce_keys [:category, :capacity, :batch_size] defstruct [ @@ -178,12 +179,9 @@ defmodule Sentry.Telemetry.Buffer do defp offer(%Buffer{size: size, capacity: capacity} = state, item) when size >= capacity do - {{:value, _dropped}, items} = :queue.out(state.items) + {{:value, dropped}, items} = :queue.out(state.items) - ClientReport.Sender.record_discarded_events( - :cache_overflow, - Category.data_category(state.category) - ) + record_overflow_discard(state.category, dropped) %{state | items: :queue.in(item, items)} end @@ -192,6 +190,24 @@ defmodule Sentry.Telemetry.Buffer do %{state | items: :queue.in(item, state.items), size: state.size + 1} end + # Log and metric structs go through the item-based recorder for their byte + # outcome; other categories fall back to the category string because the + # buffer is generic and the item-based recorder raises on unknown items. + defp record_overflow_discard(:log, %LogEvent{} = dropped) do + ClientReport.Sender.record_discarded_events(:cache_overflow, [dropped]) + end + + defp record_overflow_discard(:metric, %Metric{} = dropped) do + ClientReport.Sender.record_discarded_events(:cache_overflow, [dropped]) + end + + defp record_overflow_discard(category, _dropped) do + ClientReport.Sender.record_discarded_events( + :cache_overflow, + Category.data_category(category) + ) + end + defp poll_batch(state, count), do: poll_batch(state, count, []) defp poll_batch(state, 0, acc), do: {Enum.reverse(acc), state} defp poll_batch(%{size: 0} = state, _count, acc), do: {Enum.reverse(acc), state} diff --git a/lib/sentry/telemetry/category.ex b/lib/sentry/telemetry/category.ex index 51acee5f..82b8bff8 100644 --- a/lib/sentry/telemetry/category.ex +++ b/lib/sentry/telemetry/category.ex @@ -186,4 +186,75 @@ defmodule Sentry.Telemetry.Category do def data_category(:transaction), do: "transaction" def data_category(:log), do: "log_item" def data_category(:metric), do: "trace_metric" + + @doc """ + Returns the byte-based Sentry data category string for a given telemetry category. + + Some data categories have a companion "byte" category used to report the total + serialized size of dropped items in client reports and to honor byte-based rate + limits. Only `:log` and `:metric` currently have such companion categories. + + These strings are used in client reports and rate limiting alongside the + count-based category returned by `data_category/1`. + + ## Examples + + iex> Sentry.Telemetry.Category.byte_data_category(:log) + "log_byte" + + iex> Sentry.Telemetry.Category.byte_data_category(:metric) + "trace_metric_byte" + + """ + @spec byte_data_category(:log | :metric) :: String.t() + def byte_data_category(:log), do: "log_byte" + def byte_data_category(:metric), do: "trace_metric_byte" + + @doc """ + Returns every rate-limit data category that gates the given count-based data + category, including any companion byte category. + + Logs and metrics have a companion byte category (`log_byte` / + `trace_metric_byte`) that Sentry can rate-limit independently, so an active + limit on either the count or the byte category must suppress sending. All + other categories gate on themselves only. + + These strings are matched against the limits stored from the + `X-Sentry-Rate-Limits` response header. + + ## Examples + + iex> Sentry.Telemetry.Category.rate_limit_categories("log_item") + ["log_item", "log_byte"] + + iex> Sentry.Telemetry.Category.rate_limit_categories("trace_metric") + ["trace_metric", "trace_metric_byte"] + + iex> Sentry.Telemetry.Category.rate_limit_categories("error") + ["error"] + + """ + @spec rate_limit_categories(String.t()) :: [String.t(), ...] + def rate_limit_categories("log_item"), do: ["log_item", "log_byte"] + def rate_limit_categories("trace_metric"), do: ["trace_metric", "trace_metric_byte"] + def rate_limit_categories(category) when is_binary(category), do: [category] + + @doc """ + Returns all Sentry data category strings recognized by the SDK. + + This includes the count-based categories returned by `data_category/1` as well + as the byte-based categories returned by `byte_data_category/1`. Any category + outside this set is unknown to the SDK. + + ## Examples + + iex> categories = Sentry.Telemetry.Category.data_categories() + iex> "log_byte" in categories and "trace_metric_byte" in categories + true + + """ + @spec data_categories() :: [String.t(), ...] + def data_categories do + Enum.map(@categories, &data_category/1) ++ ["log_byte", "trace_metric_byte"] + end end diff --git a/lib/sentry/telemetry/scheduler.ex b/lib/sentry/telemetry/scheduler.ex index a64c6871..6cac8c8b 100644 --- a/lib/sentry/telemetry/scheduler.ex +++ b/lib/sentry/telemetry/scheduler.ex @@ -531,9 +531,9 @@ defmodule Sentry.Telemetry.Scheduler do items == [] -> :ok - # Transactions carry spans, so pass the actual structs through the - # list-based recorder to also record the discarded "span" outcomes. - category == :transaction -> + # These categories expand to paired outcomes (span / log_byte / + # trace_metric_byte) that only the struct-based recorder emits. + category in [:transaction, :log, :metric] -> ClientReport.Sender.record_discarded_events(:ratelimit_backoff, items) true -> @@ -551,8 +551,10 @@ defmodule Sentry.Telemetry.Scheduler do defp category_rate_limited?(%{on_envelope: cb}, _category) when is_function(cb, 1), do: false defp category_rate_limited?(_state, category) do - data_category = Category.data_category(category) - RateLimiter.rate_limited?(data_category) + category + |> Category.data_category() + |> Category.rate_limit_categories() + |> Enum.any?(&RateLimiter.rate_limited?/1) end defp default_weights do diff --git a/lib/sentry/telemetry_processor.ex b/lib/sentry/telemetry_processor.ex index de78419a..b99b17e0 100644 --- a/lib/sentry/telemetry_processor.ex +++ b/lib/sentry/telemetry_processor.ex @@ -253,7 +253,7 @@ defmodule Sentry.TelemetryProcessor do defp add_to_buffer(processor, category, item) when is_atom(processor) do data_category = Category.data_category(category) - if RateLimiter.rate_limited?(data_category) do + if data_category_rate_limited?(data_category) do {:ok, {:rate_limited, data_category}} else Buffer.add(buffer_name(processor, category), item) @@ -265,7 +265,7 @@ defmodule Sentry.TelemetryProcessor do defp add_to_buffer(processor, category, item) do data_category = Category.data_category(category) - if RateLimiter.rate_limited?(data_category) do + if data_category_rate_limited?(data_category) do {:ok, {:rate_limited, data_category}} else buffer = get_buffer(processor, category) @@ -276,6 +276,12 @@ defmodule Sentry.TelemetryProcessor do end end + defp data_category_rate_limited?(data_category) do + data_category + |> Category.rate_limit_categories() + |> Enum.any?(&RateLimiter.rate_limited?/1) + end + defp safe_get_buffer(processor, category) when is_atom(processor) do try do {:ok, diff --git a/lib/sentry/transport.ex b/lib/sentry/transport.ex index 36ae4fb0..835fa672 100644 --- a/lib/sentry/transport.ex +++ b/lib/sentry/transport.ex @@ -4,6 +4,7 @@ defmodule Sentry.Transport do # This module is exclusively responsible for encoding and POSTing envelopes to Sentry. alias Sentry.{ClientError, ClientReport, Config, Envelope, LoggerUtils} + alias Sentry.Telemetry.Category alias Sentry.Transport.RateLimiter @default_retries [1000, 2000, 4000, 8000] @@ -86,8 +87,10 @@ defmodule Sentry.Transport do defp check_rate_limited(envelope_items) do rate_limited? = Enum.any?(envelope_items, fn item -> - category = Envelope.get_data_category(item) - RateLimiter.rate_limited?(category) + item + |> Envelope.get_data_category() + |> Category.rate_limit_categories() + |> Enum.any?(&RateLimiter.rate_limited?/1) end) if rate_limited?, do: {:error, :rate_limited}, else: :ok diff --git a/test/envelope_test.exs b/test/envelope_test.exs index feea8995..18f1fc1a 100644 --- a/test/envelope_test.exs +++ b/test/envelope_test.exs @@ -336,4 +336,45 @@ defmodule Sentry.EnvelopeTest do assert Envelope.get_data_category(metric_batch) == "trace_metric" end end + + describe "log_event_byte_size/1" do + test "approximates the serialized size of a single log event" do + log_event = %LogEvent{ + timestamp: 1_588_601_261.535_386, + level: :info, + body: "something happened" + } + + assert Envelope.log_event_byte_size(log_event) > 0 + end + + test "grows with the size of the log body" do + base = %LogEvent{timestamp: 1_588_601_261.535_386, level: :info, body: "hi"} + large = %LogEvent{base | body: String.duplicate("x", 1_000)} + + assert Envelope.log_event_byte_size(large) > + Envelope.log_event_byte_size(base) + 900 + end + end + + describe "metric_byte_size/1" do + test "approximates the serialized size of a single metric" do + metric = %Metric{ + type: :counter, + name: "test.counter", + value: 1, + timestamp: 1_588_601_261.535_386 + } + + assert Envelope.metric_byte_size(metric) > 0 + end + + test "grows with the size of the metric name" do + base = %Metric{type: :counter, name: "m", value: 1, timestamp: 1_588_601_261.535_386} + large = %Metric{base | name: String.duplicate("n", 1_000)} + + assert Envelope.metric_byte_size(large) > + Envelope.metric_byte_size(base) + 900 + end + end end diff --git a/test/sentry/client_report/sender_test.exs b/test/sentry/client_report/sender_test.exs index 940a72cf..b9641d49 100644 --- a/test/sentry/client_report/sender_test.exs +++ b/test/sentry/client_report/sender_test.exs @@ -4,7 +4,7 @@ defmodule Sentry.ClientReportTest do import Sentry.TestHelpers alias Sentry.ClientReport.Sender - alias Sentry.Event + alias Sentry.{Envelope, Event, LogBatch, LogEvent, Metric, MetricBatch} setup do setup_bypass() @@ -139,5 +139,91 @@ defmodule Sentry.ClientReportTest do {:before_send, "span"} => 1 } end + + test "records both log_item and log_byte outcomes when a log event is discarded" do + start_supervised!({Sender, name: :test_log_report}) + + log_event = build_log_event("hello world") + expected_bytes = Envelope.log_event_byte_size(log_event) + assert expected_bytes > 0 + + assert :ok = + Sender.record_discarded_events(:ratelimit_backoff, [log_event], :test_log_report) + + assert :sys.get_state(:test_log_report) == %{ + {:ratelimit_backoff, "log_item"} => 1, + {:ratelimit_backoff, "log_byte"} => expected_bytes + } + end + + test "records both trace_metric and trace_metric_byte outcomes when a metric is discarded" do + start_supervised!({Sender, name: :test_metric_report}) + + metric = build_metric("requests", 1) + expected_bytes = Envelope.metric_byte_size(metric) + assert expected_bytes > 0 + + assert :ok = + Sender.record_discarded_events(:ratelimit_backoff, [metric], :test_metric_report) + + assert :sys.get_state(:test_metric_report) == %{ + {:ratelimit_backoff, "trace_metric"} => 1, + {:ratelimit_backoff, "trace_metric_byte"} => expected_bytes + } + end + + test "records aggregate log_item and log_byte outcomes when a log batch is discarded" do + start_supervised!({Sender, name: :test_log_batch_report}) + + log_events = [build_log_event("first"), build_log_event("second")] + expected_bytes = Enum.reduce(log_events, 0, &(Envelope.log_event_byte_size(&1) + &2)) + log_batch = %LogBatch{log_events: log_events} + + assert :ok = + Sender.record_discarded_events(:send_error, [log_batch], :test_log_batch_report) + + assert :sys.get_state(:test_log_batch_report) == %{ + {:send_error, "log_item"} => 2, + {:send_error, "log_byte"} => expected_bytes + } + end + + test "records aggregate trace_metric and trace_metric_byte outcomes when a metric batch is discarded" do + start_supervised!({Sender, name: :test_metric_batch_report}) + + metrics = [build_metric("a", 1), build_metric("b", 2)] + expected_bytes = Enum.reduce(metrics, 0, &(Envelope.metric_byte_size(&1) + &2)) + metric_batch = %MetricBatch{metrics: metrics} + + assert :ok = + Sender.record_discarded_events( + :send_error, + [metric_batch], + :test_metric_batch_report + ) + + assert :sys.get_state(:test_metric_batch_report) == %{ + {:send_error, "trace_metric"} => 2, + {:send_error, "trace_metric_byte"} => expected_bytes + } + end + end + + defp build_log_event(body) do + %LogEvent{ + timestamp: System.system_time(:nanosecond) / 1_000_000_000, + level: :info, + body: body + } + end + + defp build_metric(name, value) do + %Metric{ + type: :counter, + name: name, + value: value, + timestamp: System.system_time(:nanosecond) / 1_000_000_000, + attributes: %{} + } end end diff --git a/test/sentry/telemetry_processor_integration_test.exs b/test/sentry/telemetry_processor_integration_test.exs index 61be3b53..10b218d8 100644 --- a/test/sentry/telemetry_processor_integration_test.exs +++ b/test/sentry/telemetry_processor_integration_test.exs @@ -3,6 +3,8 @@ defmodule Sentry.TelemetryProcessorIntegrationTest do import Sentry.TestHelpers + require Logger + alias Sentry.TelemetryProcessor alias Sentry.Telemetry.Buffer alias Sentry.{LogEvent, Metric, Transaction} @@ -245,7 +247,12 @@ defmodule Sentry.TelemetryProcessorIntegrationTest do describe "buffer overflow client reports" do setup ctx do - Sentry.Test.setup_telemetry_processor(buffer_configs: %{log: %{capacity: 2, batch_size: 1}}) + Sentry.Test.setup_telemetry_processor( + buffer_configs: %{ + log: %{capacity: 2, batch_size: 1}, + metric: %{capacity: 2, batch_size: 1} + } + ) Sentry.ClientReport.Sender.flush() flush_ref_messages(ctx.ref) @@ -257,12 +264,12 @@ defmodule Sentry.TelemetryProcessorIntegrationTest do scheduler = TelemetryProcessor.get_scheduler(ctx.processor) :sys.suspend(scheduler) - TelemetryProcessor.add(ctx.processor, make_log_event("log-1")) + dropped_log = make_log_event("log-1") + TelemetryProcessor.add(ctx.processor, dropped_log) TelemetryProcessor.add(ctx.processor, make_log_event("log-2")) TelemetryProcessor.add(ctx.processor, make_log_event("log-3")) - log_buffer = TelemetryProcessor.get_buffer(ctx.processor, :log) - _ = Buffer.size(log_buffer) + _ = Buffer.size(TelemetryProcessor.get_buffer(ctx.processor, :log)) Sentry.ClientReport.Sender.flush() @@ -272,11 +279,53 @@ defmodule Sentry.TelemetryProcessorIntegrationTest do items = decode_envelope!(body) assert [{%{"type" => "client_report"}, client_report}] = items - cache_overflow = - Enum.find(client_report["discarded_events"], &(&1["reason"] == "cache_overflow")) + discarded = client_report["discarded_events"] + + log_item = + Enum.find(discarded, &(&1["reason"] == "cache_overflow" and &1["category"] == "log_item")) + + log_byte = + Enum.find(discarded, &(&1["reason"] == "cache_overflow" and &1["category"] == "log_byte")) + + assert log_item["quantity"] == 1 + assert log_byte["quantity"] == Sentry.Envelope.log_event_byte_size(dropped_log) + + :sys.resume(scheduler) + end + + test "sends trace_metric and trace_metric_byte reports when metric buffer overflows", ctx do + scheduler = TelemetryProcessor.get_scheduler(ctx.processor) + :sys.suspend(scheduler) + + dropped_metric = make_metric("metric-1", 1) + TelemetryProcessor.add(ctx.processor, dropped_metric) + TelemetryProcessor.add(ctx.processor, make_metric("metric-2", 2)) + TelemetryProcessor.add(ctx.processor, make_metric("metric-3", 3)) + + _ = Buffer.size(TelemetryProcessor.get_buffer(ctx.processor, :metric)) - assert cache_overflow["category"] == "log_item" - assert cache_overflow["quantity"] == 1 + Sentry.ClientReport.Sender.flush() + + ref = ctx.ref + assert_receive {:bypass_envelope, ^ref, body}, 2000 + + assert [{%{"type" => "client_report"}, client_report}] = decode_envelope!(body) + discarded = client_report["discarded_events"] + + trace_metric = + Enum.find( + discarded, + &(&1["reason"] == "cache_overflow" and &1["category"] == "trace_metric") + ) + + trace_metric_byte = + Enum.find( + discarded, + &(&1["reason"] == "cache_overflow" and &1["category"] == "trace_metric_byte") + ) + + assert trace_metric["quantity"] == 1 + assert trace_metric_byte["quantity"] == Sentry.Envelope.metric_byte_size(dropped_metric) :sys.resume(scheduler) end @@ -330,7 +379,7 @@ defmodule Sentry.TelemetryProcessorIntegrationTest do Sentry.capture_message("first-error", result: :none) - assert_receive {^ref, body}, 2000 + assert_receive {^ref, body}, 5000 assert [{%{"type" => "event"}, event}] = decode_envelope!(body) assert event["message"]["formatted"] == "first-error" @@ -347,7 +396,7 @@ defmodule Sentry.TelemetryProcessorIntegrationTest do Sentry.ClientReport.Sender.flush() - assert_receive {^ref, body}, 2000 + assert_receive {^ref, body}, 5000 items = decode_envelope!(body) assert [{%{"type" => "client_report"}, client_report}] = items @@ -360,6 +409,65 @@ defmodule Sentry.TelemetryProcessorIntegrationTest do end end + describe "byte-based rate limiting through public APIs" do + setup ctx do + Sentry.Test.setup_telemetry_processor( + buffer_configs: %{log: %{batch_size: 1}, metric: %{batch_size: 1}} + ) + + put_test_config(client: Sentry.FinchClient) + Sentry.ClientReport.Sender.flush() + flush_ref_messages(ctx.ref) + + on_exit(fn -> + for category <- ~w(log_item log_byte trace_metric trace_metric_byte) do + try do + :ets.delete(Sentry.Transport.RateLimiter, category) + catch + :error, :badarg -> :ok + end + end + end) + + :ok + end + + test "an incoming trace_metric_byte rate limit drops metrics emitted via Sentry.Metrics", + ctx do + put_test_config(enable_metrics: true) + + ref = install_rate_limit_response(ctx.bypass, "trace_metric_byte") + + Sentry.Metrics.count("first.metric", 1) + assert_receive {^ref, _body}, 5000 + + wait_for_scheduler_idle(ctx.processor) + + Sentry.Metrics.count("rate.limited.metric", 1) + refute_receive {^ref, _body}, 200 + + assert_paired_rate_limit_outcomes(ref, "trace_metric", "trace_metric_byte") + end + + @tag :capture_log + test "an incoming log_byte rate limit drops logs emitted via Logger", ctx do + put_test_config(enable_logs: true, logs: [level: :info]) + attach_sentry_logs_handler() + + ref = install_rate_limit_response(ctx.bypass, "log_byte") + + Logger.info("first log") + assert_receive {^ref, _body}, 5000 + + wait_for_scheduler_idle(ctx.processor) + + Logger.info("rate limited log") + refute_receive {^ref, _body}, 200 + + assert_paired_rate_limit_outcomes(ref, "log_item", "log_byte") + end + end + describe "pre-buffer rate limit checks" do setup ctx do Sentry.ClientReport.Sender.flush() @@ -498,6 +606,10 @@ defmodule Sentry.TelemetryProcessorIntegrationTest do on_exit(fn -> try do :ets.delete(Sentry.Transport.RateLimiter, "transaction") + :ets.delete(Sentry.Transport.RateLimiter, "log_item") + :ets.delete(Sentry.Transport.RateLimiter, "log_byte") + :ets.delete(Sentry.Transport.RateLimiter, "trace_metric") + :ets.delete(Sentry.Transport.RateLimiter, "trace_metric_byte") catch :error, :badarg -> :ok end @@ -545,6 +657,109 @@ defmodule Sentry.TelemetryProcessorIntegrationTest do assert outcomes["transaction"] == 1 assert outcomes["span"] == 3 end + + test "records log_item and log_byte outcomes when a buffered log is dropped by rate limiting", + ctx do + scheduler = TelemetryProcessor.get_scheduler(ctx.processor) + log_buffer = TelemetryProcessor.get_buffer(ctx.processor, :log) + + :sys.suspend(scheduler) + + dropped_log = make_log_event("rate-limited-log") + TelemetryProcessor.add(ctx.processor, dropped_log) + assert Buffer.size(log_buffer) == 1 + + :ets.insert(Sentry.Transport.RateLimiter, {"log_item", System.system_time(:second) + 60}) + + :sys.resume(scheduler) + GenServer.cast(scheduler, :signal) + + poll_until(fn -> Buffer.size(log_buffer) == 0 end) + + Sentry.ClientReport.Sender.flush() + + ref = ctx.ref + assert_receive {:bypass_envelope, ^ref, body}, 2000 + assert [{%{"type" => "client_report"}, client_report}] = decode_envelope!(body) + + outcomes = + for event <- client_report["discarded_events"], + event["reason"] == "ratelimit_backoff", + into: %{}, + do: {event["category"], event["quantity"]} + + assert outcomes["log_item"] == 1 + assert outcomes["log_byte"] == Sentry.Envelope.log_event_byte_size(dropped_log) + end + end + + defp install_rate_limit_response(bypass, category) do + test_pid = self() + ref = make_ref() + request_count = :counters.new(1, []) + + Bypass.expect(bypass, "POST", "/api/1/envelope/", fn conn -> + count = :counters.get(request_count, 1) + :counters.add(request_count, 1, 1) + {:ok, body, conn} = Plug.Conn.read_body(conn) + send(test_pid, {ref, body}) + + if count == 0 do + conn + |> Plug.Conn.put_resp_header("X-Sentry-Rate-Limits", "60:#{category}:organization") + |> Plug.Conn.resp(200, ~s<{"id": "340"}>) + else + Plug.Conn.resp(conn, 200, ~s<{"id": "340"}>) + end + end) + + ref + end + + defp wait_for_scheduler_idle(processor) do + scheduler = TelemetryProcessor.get_scheduler(processor) + + poll_until(fn -> + %{active_ref: active_ref} = :sys.get_state(scheduler) + is_nil(active_ref) + end) + end + + defp assert_paired_rate_limit_outcomes(ref, count_category, byte_category) do + Sentry.ClientReport.Sender.flush() + + assert_receive {^ref, body}, 5000 + assert [{%{"type" => "client_report"}, client_report}] = decode_envelope!(body) + + outcomes = + for event <- client_report["discarded_events"], + event["reason"] == "ratelimit_backoff", + into: %{}, + do: {event["category"], event["quantity"]} + + assert outcomes[count_category] == 1 + assert outcomes[byte_category] > 0 + end + + defp attach_sentry_logs_handler do + logs = Sentry.Config.logs() + + config = [ + enable_logs: true, + capture_log_messages: Keyword.fetch!(logs, :capture_log_messages), + capture_level: Keyword.fetch!(logs, :capture_level), + capture_metadata: Keyword.fetch!(logs, :capture_metadata), + capture_excluded_domains: Keyword.fetch!(logs, :capture_excluded_domains), + logs_level: Keyword.fetch!(logs, :level), + logs_metadata: Keyword.fetch!(logs, :metadata), + logs_excluded_domains: Keyword.fetch!(logs, :excluded_domains) + ] + + handler_name = :"sentry_logs_handler_#{System.unique_integer([:positive])}" + :ok = :logger.add_handler(handler_name, Sentry.LoggerHandler, %{config: config}) + on_exit(fn -> _ = :logger.remove_handler(handler_name) end) + + handler_name end defp make_transaction do diff --git a/test/sentry/transport_test.exs b/test/sentry/transport_test.exs index 9b222449..3e5d2934 100644 --- a/test/sentry/transport_test.exs +++ b/test/sentry/transport_test.exs @@ -393,6 +393,58 @@ defmodule Sentry.TransportTest do # Other categories should not be rate-limited refute Transport.RateLimiter.rate_limited?("session") end + + test "drops log envelopes when a log_byte rate limit is active", %{bypass: bypass} do + Bypass.expect_once(bypass, "POST", "/api/1/envelope/", fn conn -> + conn + |> Plug.Conn.put_resp_header("X-Sentry-Rate-Limits", "60:log_byte:key") + |> Plug.Conn.resp(200, ~s<{"id":"first"}>) + end) + + first = Envelope.from_event(Event.create_event(message: "First")) + assert {:ok, "first"} = Transport.encode_and_post_envelope(first, HackneyClient) + assert Transport.RateLimiter.rate_limited?("log_byte") + + log_envelope = Envelope.from_log_events([make_log_event("dropped")]) + + assert {:error, %ClientError{reason: :rate_limited}} = + Transport.encode_and_post_envelope(log_envelope, HackneyClient, _retries = []) + end + + test "drops metric envelopes when a trace_metric_byte rate limit is active", %{bypass: bypass} do + Bypass.expect_once(bypass, "POST", "/api/1/envelope/", fn conn -> + conn + |> Plug.Conn.put_resp_header("X-Sentry-Rate-Limits", "60:trace_metric_byte:key") + |> Plug.Conn.resp(200, ~s<{"id":"first"}>) + end) + + first = Envelope.from_event(Event.create_event(message: "First")) + assert {:ok, "first"} = Transport.encode_and_post_envelope(first, HackneyClient) + assert Transport.RateLimiter.rate_limited?("trace_metric_byte") + + metric_envelope = Envelope.from_metric_events([make_metric("dropped", 1)]) + + assert {:error, %ClientError{reason: :rate_limited}} = + Transport.encode_and_post_envelope(metric_envelope, HackneyClient, _retries = []) + end + end + + defp make_log_event(body) do + %Sentry.LogEvent{ + timestamp: System.system_time(:nanosecond) / 1_000_000_000, + level: :info, + body: body + } + end + + defp make_metric(name, value) do + %Sentry.Metric{ + type: :counter, + name: name, + value: value, + timestamp: System.system_time(:nanosecond) / 1_000_000_000, + attributes: %{} + } end defp error(fun) do