Skip to content

Commit 7bccdfc

Browse files
authored
Count post-shutdown-dropped records as already_shutdown on processor.processed (#5509)
* Count post-shutdown-dropped records as processed with error.type=already_shutdown Batch (span+log) and simple log processors previously dropped records silently after shutdown. Count them on otel.sdk.processor.{span,log}.processed with error.type=already_shutdown, matching the semantic conventions and the .NET SDK. SimpleSpanProcessor is unchanged (it has no shutdown gate). Assisted-by: Claude Opus 4.8 * Add changelog fragment for already_shutdown processed counting Assisted-by: Claude Opus 4.8 * Fix lint: ruff-format wrap and too-many-lines disable Assisted-by: Claude Opus 4.8 * Re-trigger CI Signed-off-by: cijothomas <cijo.thomas@gmail.com> --------- Signed-off-by: cijothomas <cijo.thomas@gmail.com>
1 parent 91e49a5 commit 7bccdfc

6 files changed

Lines changed: 118 additions & 12 deletions

File tree

.changelog/5509.added

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1 @@
1+
`opentelemetry-sdk`: count records dropped after shutdown on `otel.sdk.processor.{span,log}.processed` with `error.type=already_shutdown` (batch span/log and simple log processors), which the semantic conventions define as a valid value for this metric.

opentelemetry-sdk/src/opentelemetry/sdk/_logs/_internal/export/__init__.py

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -233,6 +233,7 @@ def on_emit(self, log_record: ReadWriteLogRecord):
233233
try:
234234
if self._shutdown:
235235
_logger.warning("Processor is already shutdown, ignoring call")
236+
self._metrics.drop_items(1, "already_shutdown")
236237
return
237238
# Convert ReadWriteLogRecord to ReadableLogRecord before exporting
238239
# Note: resource should not be None at this point as it's set during Logger.emit()

opentelemetry-sdk/src/opentelemetry/sdk/_shared_internal/__init__.py

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -191,6 +191,7 @@ def _export(self, batch_strategy: BatchExportStrategy) -> None:
191191
def emit(self, data: Telemetry) -> None:
192192
if self._shutdown:
193193
_logger.info("Shutdown called, ignoring %s.", self._exporting)
194+
self._metrics.drop_items(1, "already_shutdown")
194195
return
195196
if self._pid != os.getpid():
196197
self._bsp_reset_once.do_once(self._at_fork_reinit)

opentelemetry-sdk/src/opentelemetry/sdk/_shared_internal/_processor_metrics.py

Lines changed: 14 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -31,7 +31,9 @@ def register_queue_size(
3131
self, get_queue_size: Callable[[], int]
3232
) -> None: ...
3333

34-
def drop_items(self, count: int) -> None: ...
34+
def drop_items(
35+
self, count: int, error_type: str = "queue_full"
36+
) -> None: ...
3537

3638
def finish_items(self, count: int) -> None: ...
3739

@@ -40,7 +42,7 @@ class NoOpProcessorMetrics:
4042
def register_queue_size(self, get_queue_size: Callable[[], int]) -> None:
4143
pass
4244

43-
def drop_items(self, count: int) -> None:
45+
def drop_items(self, count: int, error_type: str = "queue_full") -> None:
4446
pass
4547

4648
def finish_items(self, count: int) -> None:
@@ -73,6 +75,11 @@ def __init__(
7375
ERROR_TYPE: "queue_full",
7476
}
7577

78+
self._already_shutdown_attrs = {
79+
**self._standard_attrs,
80+
ERROR_TYPE: "already_shutdown",
81+
}
82+
7683
if signal == "traces":
7784
create_processed = create_otel_sdk_processor_span_processed
7885
create_queue_capacity = (
@@ -112,8 +119,11 @@ def record_queue_size(
112119
unit=queue_size_unit,
113120
)
114121

115-
def drop_items(self, count: int) -> None:
116-
self._processed.add(count, self._dropped_attrs)
122+
def drop_items(self, count: int, error_type: str = "queue_full") -> None:
123+
if error_type == "already_shutdown":
124+
self._processed.add(count, self._already_shutdown_attrs)
125+
else:
126+
self._processed.add(count, self._dropped_attrs)
117127

118128
def finish_items(self, count: int) -> None:
119129
self._processed.add(count, self._standard_attrs)

opentelemetry-sdk/tests/logs/test_export.py

Lines changed: 56 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -1,7 +1,7 @@
11
# Copyright The OpenTelemetry Authors
22
# SPDX-License-Identifier: Apache-2.0
33

4-
# pylint: disable=protected-access
4+
# pylint: disable=protected-access,too-many-lines
55
import logging
66
import os
77
import sys
@@ -453,7 +453,7 @@ def export_logs(_logs):
453453
@patch.dict(
454454
"os.environ", {OTEL_PYTHON_SDK_INTERNAL_METRICS_ENABLED: "true"}
455455
)
456-
def test_metrics_not_counted_after_shutdown(self):
456+
def test_metrics_already_shutdown(self):
457457
metric_reader = InMemoryMetricReader()
458458
meter_provider = MeterProvider(metric_readers=[metric_reader])
459459

@@ -466,7 +466,8 @@ def test_metrics_not_counted_after_shutdown(self):
466466
processor.on_emit(EMPTY_LOG)
467467

468468
# Shut only the processor down; the record emitted afterwards hits the
469-
# already-shutdown early return and must not be counted as processed.
469+
# already-shutdown early return and is counted as processed with
470+
# error.type=already_shutdown (never handed to the exporter).
470471
processor.shutdown()
471472
processor.on_emit(EMPTY_LOG)
472473

@@ -475,11 +476,18 @@ def test_metrics_not_counted_after_shutdown(self):
475476
metrics = scope_metrics.metrics
476477
self.assertEqual(len(metrics), 1)
477478
self.assertEqual(metrics[0].name, "otel.sdk.processor.log.processed")
478-
processed_data_points = metrics[0].data.data_points
479-
self.assertEqual(len(processed_data_points), 1)
480-
self.assertEqual(processed_data_points[0].value, 1)
481-
self.assertIsNone(
482-
processed_data_points[0].attributes.get("error.type")
479+
data_points = sorted(
480+
metrics[0].data.data_points,
481+
key=lambda dp: dp.attributes.get("error.type", ""),
482+
)
483+
self.assertEqual(len(data_points), 2)
484+
# Successful submission (before shutdown), no error.type.
485+
self.assertEqual(data_points[0].value, 1)
486+
self.assertIsNone(data_points[0].attributes.get("error.type"))
487+
# Dropped after shutdown.
488+
self.assertEqual(data_points[1].value, 1)
489+
self.assertEqual(
490+
data_points[1].attributes.get("error.type"), "already_shutdown"
483491
)
484492
self.assertEqual(exporter.export.call_count, 1)
485493

@@ -709,6 +717,46 @@ def test_validation_negative_max_queue_size(self):
709717
max_export_batch_size=101,
710718
)
711719

720+
@patch.dict(
721+
"os.environ", {OTEL_PYTHON_SDK_INTERNAL_METRICS_ENABLED: "true"}
722+
)
723+
def test_metrics_already_shutdown(self):
724+
metric_reader = InMemoryMetricReader()
725+
meter_provider = MeterProvider(metric_readers=[metric_reader])
726+
727+
exporter = mock.MagicMock()
728+
exporter.export.return_value = LogRecordExportResult.SUCCESS
729+
processor = BatchLogRecordProcessor(
730+
exporter, meter_provider=meter_provider
731+
)
732+
733+
# Emitted before shutdown: drained on shutdown and counted as a
734+
# successful submit to the exporter.
735+
processor.on_emit(EMPTY_LOG)
736+
processor.shutdown()
737+
738+
# Emitted after shutdown: dropped and counted as already_shutdown.
739+
processor.on_emit(EMPTY_LOG)
740+
741+
metrics_data = metric_reader.get_metrics_data()
742+
scope_metrics = metrics_data.resource_metrics[0].scope_metrics[0]
743+
processed = next(
744+
m
745+
for m in scope_metrics.metrics
746+
if m.name == "otel.sdk.processor.log.processed"
747+
)
748+
data_points = sorted(
749+
processed.data.data_points,
750+
key=lambda dp: dp.attributes.get("error.type", ""),
751+
)
752+
self.assertEqual(len(data_points), 2)
753+
self.assertEqual(data_points[0].value, 1)
754+
self.assertIsNone(data_points[0].attributes.get("error.type"))
755+
self.assertEqual(data_points[1].value, 1)
756+
self.assertEqual(
757+
data_points[1].attributes.get("error.type"), "already_shutdown"
758+
)
759+
712760
@patch.dict(
713761
"os.environ", {OTEL_PYTHON_SDK_INTERNAL_METRICS_ENABLED: "true"}
714762
)

opentelemetry-sdk/tests/trace/export/test_export.py

Lines changed: 45 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -375,6 +375,51 @@ def test_batch_span_processor_parameters(self):
375375
max_export_batch_size=512,
376376
)
377377

378+
@mock.patch.dict(
379+
"os.environ", {OTEL_PYTHON_SDK_INTERNAL_METRICS_ENABLED: "true"}
380+
)
381+
def test_metrics_already_shutdown(self):
382+
metric_reader = InMemoryMetricReader()
383+
meter_provider = MeterProvider(metric_readers=[metric_reader])
384+
385+
exporter = mock.MagicMock()
386+
exporter.export.return_value = export.SpanExportResult.SUCCESS
387+
span_processor = export.BatchSpanProcessor(
388+
exporter, meter_provider=meter_provider
389+
)
390+
provider = trace.TracerProvider()
391+
tracer = provider.get_tracer(__name__)
392+
provider.add_span_processor(span_processor)
393+
394+
# Ended before shutdown: drained on shutdown and counted as a
395+
# successful submit to the exporter.
396+
with tracer.start_as_current_span("foo"):
397+
pass
398+
span_processor.shutdown()
399+
400+
# Ended after shutdown: dropped and counted as already_shutdown.
401+
with tracer.start_as_current_span("bar"):
402+
pass
403+
404+
metrics_data = metric_reader.get_metrics_data()
405+
scope_metrics = metrics_data.resource_metrics[0].scope_metrics[0]
406+
processed = next(
407+
m
408+
for m in scope_metrics.metrics
409+
if m.name == "otel.sdk.processor.span.processed"
410+
)
411+
data_points = sorted(
412+
processed.data.data_points,
413+
key=lambda dp: dp.attributes.get("error.type", ""),
414+
)
415+
self.assertEqual(len(data_points), 2)
416+
self.assertEqual(data_points[0].value, 1)
417+
self.assertIsNone(data_points[0].attributes.get("error.type"))
418+
self.assertEqual(data_points[1].value, 1)
419+
self.assertEqual(
420+
data_points[1].attributes.get("error.type"), "already_shutdown"
421+
)
422+
378423
@mock.patch.dict(
379424
"os.environ", {OTEL_PYTHON_SDK_INTERNAL_METRICS_ENABLED: "true"}
380425
)

0 commit comments

Comments
 (0)