Skip to content
Closed
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
40 changes: 38 additions & 2 deletions packages/traceloop-sdk/traceloop/sdk/logging/logging.py
Original file line number Diff line number Diff line change
Expand Up @@ -14,6 +14,9 @@
from opentelemetry.instrumentation.logging import LoggingInstrumentor


_FORTIFYROOT_LOGGING_HANDLER_MARKER = "_fortifyroot_logging_handler"


class LoggerWrapper(object):
resource_attributes: Dict[Any, Any] = {}
endpoint: Optional[str] = None
Expand All @@ -35,7 +38,9 @@ def __new__(cls, exporter: Optional[LogExporter] = None) -> "LoggerWrapper":
obj.__logging_provider = init_logging_provider(
obj.__logging_exporter, LoggerWrapper.resource_attributes
)
LoggingInstrumentor().instrument(set_logging_format=True)
instrumentor = LoggingInstrumentor()
if not instrumentor.is_instrumented_by_opentelemetry:
instrumentor.instrument(set_logging_format=True)

return cls.instance

Expand All @@ -49,6 +54,12 @@ def set_static_params(
LoggerWrapper.endpoint = endpoint
LoggerWrapper.headers = headers

@classmethod
def get_logging_provider(cls) -> Optional[LoggerProvider]:
if not hasattr(cls, "instance"):
return None
return cls.instance.__logging_provider


def init_logging_exporter(endpoint: str, headers: Dict[str, str]) -> LogExporter:
if "http" in endpoint.lower() or "https" in endpoint.lower():
Expand All @@ -70,6 +81,31 @@ def init_logging_provider(
logger_provider.add_log_record_processor(BatchLogRecordProcessor(exporter))

logging_handler = LoggingHandler(level=logging.NOTSET, logger_provider=logger_provider)
logging.basicConfig(level=logging.INFO, handlers=[logging_handler])
_attach_root_logging_handler(logging_handler)

return logger_provider


def is_fortifyroot_logging_handler(handler: logging.Handler) -> bool:
return bool(getattr(handler, _FORTIFYROOT_LOGGING_HANDLER_MARKER, False))


def _attach_root_logging_handler(logging_handler: LoggingHandler) -> None:
root_logger = logging.getLogger()
had_non_fortifyroot_handlers = any(
not is_fortifyroot_logging_handler(handler) for handler in root_logger.handlers
)

for handler in list(root_logger.handlers):
if is_fortifyroot_logging_handler(handler):
root_logger.removeHandler(handler)

setattr(logging_handler, _FORTIFYROOT_LOGGING_HANDLER_MARKER, True)
root_logger.addHandler(logging_handler)

# Preserve existing app logging levels/formatters when they are already
# configured. If the app has not configured logging at all, keep the prior
# default of exporting INFO and above.
if not had_non_fortifyroot_handlers and root_logger.level > logging.INFO:
root_logger.setLevel(logging.INFO)

43 changes: 31 additions & 12 deletions packages/traceloop-sdk/traceloop/sdk/tracing/tracing.py
Original file line number Diff line number Diff line change
Expand Up @@ -99,6 +99,11 @@ def __new__(
obj.__tracer_provider = init_tracer_provider(
resource=obj.__resource, sampler=sampler
)
callback_processor = (
_SpanPostprocessCallbackProcessor(span_postprocess_callback)
if span_postprocess_callback
else None
)

# Handle multiple processors case
if processor is not None and isinstance(processor, list):
Expand All @@ -118,6 +123,9 @@ def chained_on_start(

obj.__tracer_provider.add_span_processor(proc)

if callback_processor is not None:
obj.__tracer_provider.add_span_processor(callback_processor)

# Handle single processor case (backward compatibility)
elif processor is not None:
obj.__spans_processor: SpanProcessor = processor
Expand All @@ -132,27 +140,21 @@ def chained_on_start(span, parent_context=None, orig=original_on_start):

obj.__tracer_provider.add_span_processor(obj.__spans_processor)

if callback_processor is not None:
obj.__tracer_provider.add_span_processor(callback_processor)

# Handle default processor case
else:
obj.__spans_processor = get_default_span_processor(
disable_batch=disable_batch, exporter=exporter
)

if span_postprocess_callback:
# Create a wrapper that calls both the custom and original methods
original_on_end = obj.__spans_processor.on_end

def wrapped_on_end(span):
# Call the custom on_end first
span_postprocess_callback(span)
# Then call the original to ensure normal processing
original_on_end(span)

obj.__spans_processor.on_end = wrapped_on_end

obj.__spans_processor.on_start = obj._span_processor_on_start
obj.__tracer_provider.add_span_processor(obj.__spans_processor)

if callback_processor is not None:
obj.__tracer_provider.add_span_processor(callback_processor)

if propagator:
set_global_textmap(propagator)

Expand Down Expand Up @@ -234,6 +236,23 @@ def get_tracer(self):
return self.__tracer_provider.get_tracer(TRACER_NAME)


class _SpanPostprocessCallbackProcessor(SpanProcessor):
def __init__(self, callback: Callable[[ReadableSpan], None]) -> None:
self._callback = callback

def on_start(self, span: Span, parent_context: Optional[Context] = None) -> None:
return None

def on_end(self, span: ReadableSpan) -> None:
self._callback(span)

def shutdown(self) -> None:
return None

def force_flush(self, timeout_millis: int = 30000) -> bool:
return True


def set_association_properties(properties: dict) -> None:
attach(set_value("association_properties", properties))

Expand Down
Loading