diff --git a/.changelog/5676.fixed b/.changelog/5676.fixed new file mode 100644 index 00000000000..149f94505f6 --- /dev/null +++ b/.changelog/5676.fixed @@ -0,0 +1 @@ +`opentelemetry-sdk`: count logs dropped by `SimpleLogRecordProcessor`'s recursion guard on `otel.sdk.processor.log.processed` with `error.type=recursion` diff --git a/opentelemetry-sdk/src/opentelemetry/sdk/_logs/_internal/export/__init__.py b/opentelemetry-sdk/src/opentelemetry/sdk/_logs/_internal/export/__init__.py index d6ff7717eff..55921e043db 100644 --- a/opentelemetry-sdk/src/opentelemetry/sdk/_logs/_internal/export/__init__.py +++ b/opentelemetry-sdk/src/opentelemetry/sdk/_logs/_internal/export/__init__.py @@ -208,6 +208,7 @@ def on_emit(self, log_record: ReadWriteLogRecord): _propagate_false_logger.warning( "SimpleLogRecordProcessor.on_emit has entered a recursive loop. Dropping log and exiting the loop." ) + self._metrics.drop_items(1, "recursion") return token = attach( set_value( diff --git a/opentelemetry-sdk/src/opentelemetry/sdk/_shared_internal/_processor_metrics.py b/opentelemetry-sdk/src/opentelemetry/sdk/_shared_internal/_processor_metrics.py index 699f7e1821c..aae2b14d967 100644 --- a/opentelemetry-sdk/src/opentelemetry/sdk/_shared_internal/_processor_metrics.py +++ b/opentelemetry-sdk/src/opentelemetry/sdk/_shared_internal/_processor_metrics.py @@ -76,6 +76,11 @@ def __init__( ERROR_TYPE: "already_shutdown", } + self._recursion_attrs = { + **self._standard_attrs, + ERROR_TYPE: "recursion", + } + if signal == "traces": create_processed = create_otel_sdk_processor_span_processed create_queue_capacity = create_otel_sdk_processor_span_queue_capacity @@ -114,6 +119,8 @@ def record_queue_size( def drop_items(self, count: int, error_type: str = "queue_full") -> None: if error_type == "already_shutdown": self._processed.add(count, self._already_shutdown_attrs) + elif error_type == "recursion": + self._processed.add(count, self._recursion_attrs) else: self._processed.add(count, self._dropped_attrs) diff --git a/opentelemetry-sdk/tests/logs/test_export.py b/opentelemetry-sdk/tests/logs/test_export.py index cc0b22d29c7..a3265772eff 100644 --- a/opentelemetry-sdk/tests/logs/test_export.py +++ b/opentelemetry-sdk/tests/logs/test_export.py @@ -95,6 +95,60 @@ def export(self, batch: Sequence[ReadableLogRecord]): finally: root_logger.removeHandler(handler) + @patch.dict("os.environ", {OTEL_PYTHON_SDK_INTERNAL_METRICS_ENABLED: "true"}) + @mark.skipif( + (3, 13, 0) <= sys.version_info <= (3, 13, 5), + reason="This will fail on 3.13.5 due to https://github.com/python/cpython/pull/131812 which prevents the recursion from being detected.", + ) + def test_metrics_recursive_loop(self): + metric_reader = InMemoryMetricReader() + meter_provider = MeterProvider(metric_readers=[metric_reader]) + + class Exporter(LogRecordExporter): + def shutdown(self): + pass + + def force_flush(self, timeout_millis: int = 10_000) -> bool: + return True + + def export(self, batch: Sequence[ReadableLogRecord]): + logger = logging.getLogger("any logger..") + logger.warning("Something happened.") + + exporter = Exporter() + logger_provider = LoggerProvider() + logger_provider.add_log_record_processor(SimpleLogRecordProcessor(exporter, meter_provider=meter_provider)) + root_logger = logging.getLogger() + handler = LoggingHandler(level=logging.DEBUG, logger_provider=logger_provider) + root_logger.addHandler(handler) + propagate_false_logger = logging.getLogger("opentelemetry.sdk._logs._internal.export.propagate.false") + try: + with self.assertLogs(propagate_false_logger) as cm: + root_logger.warning("hello!") + assert "SimpleLogRecordProcessor.on_emit has entered a recursive loop" in cm.output[0] + finally: + root_logger.removeHandler(handler) + + metrics_data = metric_reader.get_metrics_data() + scope_metrics = metrics_data.resource_metrics[0].scope_metrics[0] + metrics = scope_metrics.metrics + self.assertEqual(len(metrics), 1) + self.assertEqual(metrics[0].name, "otel.sdk.processor.log.processed") + data_points = sorted( + metrics[0].data.data_points, + key=lambda dp: dp.attributes.get("error.type", ""), + ) + # The record the recursion guard discards is counted as processed with + # error.type=recursion, so the total stays reconcilable with the number + # of records the processor accepted. + recursion_points = [dp for dp in data_points if dp.attributes.get("error.type") == "recursion"] + self.assertEqual( + len(recursion_points), + 1, + "the log dropped by the recursion guard was not counted", + ) + self.assertEqual(recursion_points[0].value, 1) + def test_simple_log_record_processor_default_level(self): exporter = InMemoryLogRecordExporter() logger_provider = LoggerProvider()