Skip to content
Draft
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
Original file line number Diff line number Diff line change
Expand Up @@ -34,7 +34,7 @@ Homepage = "https://github.com/open-telemetry/opentelemetry-python/tree/main/exp
Repository = "https://github.com/open-telemetry/opentelemetry-python"

[tool.hatch.version]
path = "src/opentelemetry/exporter/otlp/proto/common/version/__init__.py"
path = "src/opentelemetry/exporter/otlp/_proto/common/version/__init__.py"

[tool.hatch.build.targets.sdist]
include = [
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,153 @@
# Copyright The OpenTelemetry Authors
# SPDX-License-Identifier: Apache-2.0

from __future__ import annotations

from collections import Counter
from collections.abc import Iterator
from contextlib import AbstractContextManager, contextmanager
from dataclasses import dataclass
from time import perf_counter
from typing import TYPE_CHECKING, Protocol

from opentelemetry.metrics import MeterProvider, get_meter_provider
from opentelemetry.semconv._incubating.attributes.otel_attributes import (
OTEL_COMPONENT_NAME,
OTEL_COMPONENT_TYPE,
OtelComponentTypeValues,
)
from opentelemetry.semconv._incubating.metrics.otel_metrics import (
create_otel_sdk_exporter_log_exported,
create_otel_sdk_exporter_log_inflight,
create_otel_sdk_exporter_metric_data_point_exported,
create_otel_sdk_exporter_metric_data_point_inflight,
create_otel_sdk_exporter_operation_duration,
create_otel_sdk_exporter_span_exported,
create_otel_sdk_exporter_span_inflight,
)
from opentelemetry.semconv.attributes.error_attributes import ERROR_TYPE
from opentelemetry.semconv.attributes.server_attributes import (
SERVER_ADDRESS,
SERVER_PORT,
)

if TYPE_CHECKING:
from typing import Literal
from urllib.parse import ParseResult as UrlParseResult

from opentelemetry.util.types import Attributes, AttributeValue

_component_counter = Counter()


@dataclass
class ExportResult:
error: Exception | None = None
error_attrs: Attributes = None


class ExporterMetricsT(Protocol):
def export_operation(
self, num_items: int
) -> AbstractContextManager[ExportResult]: ...


class NoOpExporterMetrics:
@contextmanager
def export_operation(self, num_items: int) -> Iterator[ExportResult]:
yield ExportResult()


class ExporterMetrics:
def __init__(
self,
component_type: OtelComponentTypeValues | None,
signal: Literal["traces", "metrics", "logs"],
endpoint: UrlParseResult,
meter_provider: MeterProvider | None,
) -> None:
if signal == "traces":
create_exported = create_otel_sdk_exporter_span_exported
create_inflight = create_otel_sdk_exporter_span_inflight
elif signal == "logs":
create_exported = create_otel_sdk_exporter_log_exported
create_inflight = create_otel_sdk_exporter_log_inflight
else:
create_exported = (
create_otel_sdk_exporter_metric_data_point_exported
)
create_inflight = (
create_otel_sdk_exporter_metric_data_point_inflight
)

port = endpoint.port
if port is None:
if endpoint.scheme == "https":
port = 443
elif endpoint.scheme == "http":
port = 80

component_type_value = (
component_type.value if component_type else "unknown_otlp_exporter"
)
count = _component_counter[component_type_value]
_component_counter[component_type_value] = count + 1
self._standard_attrs: dict[str, AttributeValue] = {
OTEL_COMPONENT_TYPE: component_type_value,
OTEL_COMPONENT_NAME: f"{component_type_value}/{count}",
}
if endpoint.hostname:
self._standard_attrs[SERVER_ADDRESS] = endpoint.hostname
if port is not None:
self._standard_attrs[SERVER_PORT] = port

meter_provider = meter_provider or get_meter_provider()
meter = meter_provider.get_meter("opentelemetry-sdk")
self._inflight = create_inflight(meter)
self._exported = create_exported(meter)
self._duration = create_otel_sdk_exporter_operation_duration(meter)

@contextmanager
def export_operation(self, num_items: int) -> Iterator[ExportResult]:
start_time = perf_counter()
self._inflight.add(num_items, self._standard_attrs)

result = ExportResult()
try:
yield result
finally:
error = result.error
error_attrs = result.error_attrs

end_time = perf_counter()
self._inflight.add(-num_items, self._standard_attrs)
exported_attrs = (
{**self._standard_attrs, ERROR_TYPE: type(error).__qualname__}
if error
else self._standard_attrs
)
self._exported.add(num_items, exported_attrs)
duration_attrs = (
{**exported_attrs, **error_attrs}
if error_attrs
else exported_attrs
)
self._duration.record(end_time - start_time, duration_attrs)


def create_exporter_metrics(
component_type: OtelComponentTypeValues | None,
signal: Literal["traces", "metrics", "logs"],
endpoint: UrlParseResult,
meter_provider: MeterProvider | None,
enabled: bool,
) -> ExporterMetricsT:
if not enabled:
return NoOpExporterMetrics()

return ExporterMetrics(
component_type,
signal,
endpoint,
meter_provider,
)
Original file line number Diff line number Diff line change
@@ -0,0 +1,114 @@
# Copyright The OpenTelemetry Authors
# SPDX-License-Identifier: Apache-2.0

from __future__ import annotations

from logging import getLogger
from collections.abc import Callable, Mapping, Sequence
from typing import Any, TypeVar

from opentelemetry._proto.common.v1.common_pb2 import (
AnyValue,
ArrayValue,
InstrumentationScope as PB2InstrumentationScope,
KeyValue,
KeyValueList,
)
from opentelemetry._proto.resource.v1.resource_pb2 import Resource as PB2Resource
from opentelemetry.sdk.trace import Resource
from opentelemetry.sdk.util.instrumentation import InstrumentationScope
from opentelemetry.util.types import _ExtendedAttributes

_logger = getLogger(__name__)

_TypingResourceT = TypeVar("_TypingResourceT")
_ResourceDataT = TypeVar("_ResourceDataT")


def _encode_instrumentation_scope(
instrumentation_scope: InstrumentationScope,
) -> PB2InstrumentationScope:
if instrumentation_scope is None:
return PB2InstrumentationScope()
return PB2InstrumentationScope(
name=instrumentation_scope.name,
version=instrumentation_scope.version,
attributes=_encode_attributes(instrumentation_scope.attributes),
)


def _encode_resource(resource: Resource) -> PB2Resource:
return PB2Resource(attributes=_encode_attributes(resource.attributes))


def _encode_value(value: Any) -> AnyValue:
if value is None:
return AnyValue()
if isinstance(value, bool):
return AnyValue(bool_value=value)
if isinstance(value, str):
return AnyValue(string_value=value)
if isinstance(value, int):
return AnyValue(int_value=value)
if isinstance(value, float):
return AnyValue(double_value=value)
if isinstance(value, bytes):
return AnyValue(bytes_value=value)
if isinstance(value, Sequence):
return AnyValue(
array_value=ArrayValue(values=[_encode_value(v) for v in value])
)
if isinstance(value, Mapping):
return AnyValue(
kvlist_value=KeyValueList(
values=[_encode_key_value(str(k), v) for k, v in value.items()]
)
)
raise Exception(f"Invalid type {type(value)} of value {value}")


def _encode_key_value(key: str, value: Any) -> KeyValue:
return KeyValue(key=key, value=_encode_value(value))


def _encode_span_id(span_id: int) -> bytes:
return span_id.to_bytes(length=8, byteorder="big", signed=False)


def _encode_trace_id(trace_id: int) -> bytes:
return trace_id.to_bytes(length=16, byteorder="big", signed=False)


def _encode_attributes(
attributes: _ExtendedAttributes | None,
) -> list[KeyValue]:
if not attributes:
return []
pb2_attributes = []
for key, value in attributes.items():
try:
pb2_attributes.append(_encode_key_value(key, value))
except Exception as error:
_logger.exception("Failed to encode key %s: %s", key, error)
return pb2_attributes


def _get_resource_data(
sdk_resource_scope_data: dict[Resource, _ResourceDataT],
resource_class: Callable[..., _TypingResourceT],
name: str,
) -> list[_TypingResourceT]:
resource_data = []
for sdk_resource, scope_data in sdk_resource_scope_data.items():
collector_resource = PB2Resource(
attributes=_encode_attributes(sdk_resource.attributes)
)
resource_data.append(
resource_class(
**{
"resource": collector_resource,
f"scope_{name}": scope_data.values(),
}
)
)
return resource_data
Original file line number Diff line number Diff line change
@@ -0,0 +1,87 @@
# Copyright The OpenTelemetry Authors
# SPDX-License-Identifier: Apache-2.0
from collections import defaultdict
from collections.abc import Sequence

from opentelemetry.exporter.otlp._proto.common._internal import (
_encode_attributes,
_encode_instrumentation_scope,
_encode_resource,
_encode_span_id,
_encode_trace_id,
_encode_value,
)
from opentelemetry._proto.collector.logs.v1.logs_service_pb2 import (
ExportLogsServiceRequest,
)
from opentelemetry._proto.logs.v1.logs_pb2 import (
LogRecord,
ResourceLogs,
ScopeLogs,
)
from opentelemetry.sdk._logs import ReadableLogRecord


def encode_logs(batch: Sequence[ReadableLogRecord]) -> ExportLogsServiceRequest:
return ExportLogsServiceRequest(resource_logs=_encode_resource_logs(batch))


def _encode_log(readable_log_record: ReadableLogRecord) -> LogRecord:
log = readable_log_record.log_record
span_id = (
b""
if log.span_id == 0
else _encode_span_id(log.span_id)
)
trace_id = (
b""
if log.trace_id == 0
else _encode_trace_id(log.trace_id)
)
return LogRecord(
time_unix_nano=log.timestamp,
observed_time_unix_nano=log.observed_timestamp,
span_id=span_id,
trace_id=trace_id,
flags=int(log.trace_flags),
body=_encode_value(log.body),
severity_text=log.severity_text or "",
attributes=_encode_attributes(log.attributes),
dropped_attributes_count=readable_log_record.dropped_attributes,
severity_number=getattr(log.severity_number, "value", 0) or 0,
event_name=log.event_name or "",
)


def _encode_resource_logs(
batch: Sequence[ReadableLogRecord],
) -> list[ResourceLogs]:
sdk_resource_logs: dict = defaultdict(lambda: defaultdict(list))

for readable_log in batch:
sdk_resource = readable_log.resource
sdk_instrumentation = readable_log.instrumentation_scope or None
pb2_log = _encode_log(readable_log)
sdk_resource_logs[sdk_resource][sdk_instrumentation].append(pb2_log)

pb2_resource_logs = []
for sdk_resource, sdk_instrumentations in sdk_resource_logs.items():
scope_logs = []
for sdk_instrumentation, pb2_logs in sdk_instrumentations.items():
scope_logs.append(
ScopeLogs(
scope=_encode_instrumentation_scope(sdk_instrumentation),
log_records=pb2_logs,
schema_url=sdk_instrumentation.schema_url
if sdk_instrumentation
else "",
)
)
pb2_resource_logs.append(
ResourceLogs(
resource=_encode_resource(sdk_resource),
scope_logs=scope_logs,
schema_url=sdk_resource.schema_url,
)
)
return pb2_resource_logs
Loading
Loading