Coverage for anaconda_opentelemetry/signals.py: 92%
225 statements
« prev ^ index » next coverage.py v7.16.2, created at 2026-10-01 20:06 +0000
« prev ^ index » next coverage.py v7.16.2, created at 2026-10-01 20:06 +0000
1# -*- coding: utf-8 -*-
2# SPDX-FileCopyrightText: 2025 Anaconda, Inc
3# SPDX-License-Identifier: Apache-2.0
5# signals.py
6"""
7Anaconda Telemetry - Metrics Module
9This module provides functionality for logging, metrics, and tracing (together called
10signals) using OpenTelemetry. It includes classes for handling logging, metrics, and
11tracing, as well as functions for initializing the telemetry system and recording metrics.
12"""
14import logging, socket, threading
15from typing import Dict, Iterator, List, Optional
16from contextlib import contextmanager
18from opentelemetry import trace, metrics
19from opentelemetry.sdk.trace import TracerProvider
20from opentelemetry.sdk.metrics import MeterProvider
22from ._compat import OTEL_PRIVATE_LOGS_AVAILABLE, LoggingHandler
24from .config import Configuration as Config
25from .attributes import ResourceAttributes as Attributes
26from .formatting import AttrDict
28from .common import MetricsNotInitialized
29from .logging import _AnacondaLogger
30from .metrics import _AnacondaMetrics
31from .tracing import _AnacondaTrace, _ASpan
34_SUPPRESSED_LOGGER_ROOTS = ('opentelemetry',)
37def _suppress_otel_export_errors():
38 for root in _SUPPRESSED_LOGGER_ROOTS:
39 logging.getLogger(root).setLevel(logging.CRITICAL)
40 prefix = root + '.'
41 for name, logger in list(logging.Logger.manager.loggerDict.items()):
42 if name.startswith(prefix) and isinstance(logger, logging.Logger):
43 logger.setLevel(logging.NOTSET)
46# Internet and endpoint access check method
47def __check_internet_status(config: Config, timeout: float = 5.0) -> tuple[bool,bool]: # seconds max to pause....
48 # Relies on Configuration to validate the endpoint...
49 internet = True
50 access = True
51 if config._get_skip_internet_check():
52 return True, True
53 endpoint = config._get_default_endpoint()
54 try:
55 # Access to a highly available DNS site...
56 socket.create_connection(('8.8.8.8', 53), timeout=timeout / 2).close()
57 except OSError:
58 logging.getLogger(__package__).warning("Anaconda OpenTelemetry: No Internet was detected!")
59 internet = False # No internet, but internet is not an absolute requirement for on-prem solutions.
60 try:
61 socket.create_connection((config._endpoints['default_endpoint'].host, config._endpoints['default_endpoint']._internet_check_port), timeout=timeout / 2).close()
62 except OSError:
63 logging.getLogger(__package__).fatal(f"Anaconda OpenTelemetry: No access to the endpoint '{endpoint}'!")
64 access = False # This could be fatal, not endpoint for telemetry.
65 if access == True:
66 logging.getLogger(__package__).info(f"Anaconda OpenTelemetry: Successful access to the endpoint '{endpoint}'!")
67 return internet, access
69__ANACONDA_TELEMETRY_INITIALIZED = False
70__SIGNALS = None
71__CONFIG = None
74################################################################################
75# Exposed APIs
76def initialize_telemetry(config: Config,
77 attributes: Attributes = None,
78 signal_types: List[str] = ['metrics']):
79 """
80 Initializes the telemetry system.
82 Args:
83 service_name (str): The name of the service.
84 service_version (str): The version of the service.
85 config (Configuration): The configuration for the telemetry. At a minimum, the Configuration must have a default endpoint
86 for connection to the collector.
87 attributes (ResourceAttributes, optional): A class containing common attributes. If provided,
88 it will override any values shared with configuration file.
89 signal_types (list, optional): List of metric types to initialize. Defaults to ['logging','metrics','tracing'].
90 Supported values are 'logging', 'metrics', and 'tracing'. If an empty list is provided, no metrics will be initialized.
92 Raises:
93 ValueError: If the config passed is None or the attributes passed are None.
94 """
95 global __ANACONDA_TELEMETRY_INITIALIZED
96 global __SIGNALS
97 global __CONFIG
99 if __ANACONDA_TELEMETRY_INITIALIZED is True:
100 return # Already initialized
101 if config is None:
102 raise ValueError(f"The config argument is required but was None")
103 if attributes is None:
104 raise ValueError(f"The attributes argument is required but was None")
106 __CONFIG = config
107 __SIGNALS = signal_types
109 if type(attributes.parameters) != dict:
110 raise ValueError(f"The parameters attribute in ResourceAttributes must be a dictionary")
112 if not config._get_verbose_export_errors():
113 _suppress_otel_export_errors()
115 # Right now, no action is taken but it possible to disable telemetry with no access to the endpoint...
116 _, _ = __check_internet_status(config, timeout=4) # Max wait 4 seconds...
118 # all params are the same currently so only write them once
119 init_params = (config, attributes)
121 # Initialize logging here...
122 signal_type_count = 0
123 if 'logging' in signal_types:
124 # The OTel logs SDK lives behind a private namespace; _compat degrades to
125 # no-ops if it ever disappears. Skip the signal rather than wire up
126 # inert objects the caller would have to debug.
127 if OTEL_PRIVATE_LOGS_AVAILABLE:
128 _AnacondaLogger._instance = _AnacondaLogger(*init_params)
129 signal_type_count += 1
130 else:
131 logging.getLogger(__package__).warning(
132 "Anaconda OpenTelemetry: the 'logging' signal was requested but the "
133 "installed opentelemetry-sdk does not provide the required logs API. "
134 "Skipping log telemetry."
135 )
137 # Initialize the telemetry system here
138 if 'metrics' in signal_types:
139 _AnacondaMetrics._instance = _AnacondaMetrics(*init_params)
140 signal_type_count += 1
142 # Initialize tracing here...
143 if 'tracing' in signal_types:
144 _AnacondaTrace._instance = _AnacondaTrace(*init_params)
145 signal_type_count += 1
147 if signal_type_count == 0:
148 logging.getLogger(__package__).warning(
149 "No signal types were initialized. Was this intended? If not please check the " +
150 "'metrics' section in the configuration file and/or the list of " +
151 "metric types in the parameter 'signal_types'."
152 )
153 __ANACONDA_TELEMETRY_INITIALIZED = True
155_SHUTDOWN_DONE = False
156_shutdown_lock = threading.Lock()
159def flush_telemetry() -> bool:
160 """Force-flush the providers this package emits through.
162 Spans and metrics are recorded via the OTel globals, so the global getters are correct
163 for them. Log records are not: send_event and get_telemetry_logger_handler write to the
164 LoggerProvider owned by _AnacondaLogger, which is not the global one when another
165 library set that first, so the owned provider is flushed instead.
167 Returns True if every provider we own flushed successfully.
168 """
169 if not __ANACONDA_TELEMETRY_INITIALIZED:
170 return False
171 success = True
172 try:
173 tp = trace.get_tracer_provider()
174 if isinstance(tp, TracerProvider):
175 try:
176 tp.force_flush()
177 except Exception:
178 logging.getLogger(__package__).debug("Tracer flush failed", exc_info=True)
179 success = False
181 mp = metrics.get_meter_provider()
182 if isinstance(mp, MeterProvider):
183 try:
184 mp.force_flush()
185 except Exception:
186 logging.getLogger(__package__).debug("Meter flush failed", exc_info=True)
187 success = False
189 if _AnacondaLogger._instance is not None:
190 lp = _AnacondaLogger._instance._provider
191 try:
192 lp.force_flush()
193 except Exception:
194 logging.getLogger(__package__).debug("Logger flush failed", exc_info=True)
195 success = False
196 except Exception:
197 logging.getLogger(__package__).debug("flush_telemetry failed", exc_info=True)
198 success = False
199 return success
202def shutdown_telemetry(*, timeout_seconds: Optional[float] = None) -> bool:
203 """Flush all telemetry providers at process shutdown, optionally time-bounded.
205 Performs a bounded *force-flush* (via :func:`flush_telemetry`). It intentionally does
206 not call ``provider.shutdown()``: at process exit that only adds worker-thread joins
207 (more blocking) with no benefit. Pair with
208 ``config.set_shutdown_on_exit(False)`` to control flush timing from a
209 signal handler or atexit path.
211 With ``timeout_seconds=None`` the flush runs synchronously (unbounded). When set, the
212 flush runs on a daemon thread joined for at most ``timeout_seconds``; only the
213 caller's wait is bounded (a still-running flush thread is reaped at interpreter exit).
215 Idempotent and thread-safe: once a flush completes it is not repeated; concurrent or
216 re-entrant (signal-handler) calls never double-flush or block. A call that times out
217 does not mark completion, so a later call may retry.
219 The ``join`` blocks the calling thread, so do not call this with a ``timeout_seconds``
220 from inside an async event loop; use it from a signal handler, an atexit path, or a
221 dedicated thread.
223 Returns True if the flush completed (now or previously); False if telemetry was never
224 initialized, the flush timed out, or another call is already in progress.
225 """
226 global _SHUTDOWN_DONE
227 if not __ANACONDA_TELEMETRY_INITIALIZED:
228 return False
229 if _SHUTDOWN_DONE:
230 return True
231 # Non-blocking so a concurrent or re-entrant caller returns immediately instead of
232 # double-flushing or deadlocking (this may run inside a signal handler).
233 if not _shutdown_lock.acquire(blocking=False):
234 return _SHUTDOWN_DONE
235 try:
236 if _SHUTDOWN_DONE:
237 return True
238 if timeout_seconds is None:
239 completed = flush_telemetry()
240 else:
241 flush_thread = threading.Thread(target=flush_telemetry, daemon=True)
242 flush_thread.start()
243 flush_thread.join(timeout=timeout_seconds)
244 completed = not flush_thread.is_alive()
245 if completed:
246 _SHUTDOWN_DONE = True
247 return completed
248 finally:
249 _shutdown_lock.release()
252def change_signal_endpoint(signal_type: str,
253 new_endpoint: str,
254 auth_token: str = None):
255 """
256 Updates the endpoint for the passed signal
258 Args:
259 signal_type (str): signal type to update endpoint for. Supported values are 'logging', 'metrics', and 'tracing'
261 Returns:
262 boolean: value indicating whether the update was successful or not
263 """
264 if signal_type.lower() == 'metrics':
265 _AnacondaTelInstance = _AnacondaMetrics
266 batch_access = _AnacondaTelInstance._instance.metric_reader
267 elif signal_type.lower() == 'tracing':
268 _AnacondaTelInstance = _AnacondaTrace
269 batch_access = _AnacondaTelInstance._instance._processor
270 elif signal_type.lower() == 'logging':
271 _AnacondaTelInstance = _AnacondaLogger
272 batch_access = _AnacondaTelInstance._instance._processor
273 else:
274 logging.getLogger(__package__).warning(f"{signal_type} not a valid signal type.")
275 return False
277 # execute OpenTelemetry changes
278 updated_endpoint = _AnacondaTelInstance._instance.exporter.change_signal_endpoint(
279 batch_access,
280 _AnacondaTelInstance._instance._config,
281 new_endpoint,
282 auth_token=auth_token
283 )
285 if not updated_endpoint:
286 logging.getLogger(__package__).warning(f"Endpoint for {signal_type} failed to update.")
287 return False
288 else:
289 logging.getLogger(__package__).info(f"Endpoint for {signal_type} was successfully updated.")
290 return True
292def record_histogram(metric_name, value, attributes: AttrDict={}) -> bool:
293 """
294 Records a increasing only metric with the given name and value. The value will
295 always appear in the attributes section in the raw OTLP output and the timestamp
296 will be the histogram value.
298 Will catch any exceptions generated by metric usage.
300 Args:
301 metric_name (str): The name of the metric.
302 value (float): The value of the metric. Can be any float since the timestamp is the ever increasing value of the histogram.
303 attributes (dict, optional): Additional attributes for the metric. Defaults to {}.
305 Returns:
306 bool: True if the metric was recorded successfully, False otherwise (logging the error).
307 """
308 if __ANACONDA_TELEMETRY_INITIALIZED is False:
309 logging.getLogger(__package__).error("Anaconda telemetry system not initialized.") # Since init didn't happen this is not exported in OTel!!!
310 return False
311 try:
312 return _AnacondaMetrics._instance.record_histogram(metric_name, value, _AnacondaMetrics._instance._process_attributes(attributes))
313 except MetricsNotInitialized as me:
314 logging.getLogger(__package__).warning(f"An attempt was made to record a histogram metric when metrics were not configured.")
315 return False
316 except Exception as e:
317 logging.getLogger(__package__).error(f"UNCAUGHT EXCEPTION:\n{e}")
318 return False
320def set_gauge(metric_name, value, attributes: AttrDict={}) -> bool:
321 """
322 Sets a gauge metric with the given name to the given value. A gauge records the last
323 value set rather than a sum, so it is the right choice for values that go up and down
324 and are sampled rather than accumulated (queue depth, memory in use, temperature).
326 Will catch any exceptions generated by metric usage.
328 Args:
329 metric_name (str): The name of the metric.
330 value (int | float): The current value of the metric. Must be an int or a float; negative
331 values are allowed. A value of any other type (including None and bool) is rejected.
332 attributes (dict, optional): Additional attributes for the metric. Defaults to {}.
334 Returns:
335 bool: True if the metric was recorded successfully, False otherwise (logging the error).
336 """
337 if __ANACONDA_TELEMETRY_INITIALIZED is False:
338 logging.getLogger(__package__).error("Anaconda telemetry system not initialized.") # Since init didn't happen this is not exported in OTel!!!
339 return False
340 try:
341 return _AnacondaMetrics._instance.set_gauge(metric_name, value, _AnacondaMetrics._instance._process_attributes(attributes))
342 except MetricsNotInitialized:
343 logging.getLogger(__package__).warning("An attempt was made to set a gauge metric when metrics were not configured.")
344 return False
345 except Exception as e:
346 logging.getLogger(__package__).error(f"UNCAUGHT EXCEPTION:\n{e}")
347 return False
349def increment_counter(counter_name, by=1, attributes: AttrDict={}) -> bool:
350 """
351 Increments a counter or up down counter by the given parameter 'by'.
353 Will catch any exceptions generated by metric usage.
355 Args:
356 counter_name (str): The name of the counter.
357 by (int, optional): The value to increment by. Defaults to 1. The abs(by) is used to protect from negative numbers.
358 attributes (dict, optional): Additional attributes for the counter. Defaults to {}.
360 Returns:
361 bool: True if the counter was incremented successfully, False otherwise (logging the error).
362 """
363 if __ANACONDA_TELEMETRY_INITIALIZED is False:
364 logging.getLogger(__package__).error("Anaconda telemetry system not initialized.") # Since init didn't happen this is not exported in OTel!!!
365 return False
366 try:
367 return _AnacondaMetrics._instance.increment_counter(counter_name, by, _AnacondaMetrics._instance._process_attributes(attributes))
368 except MetricsNotInitialized:
369 logging.getLogger(__package__).warning(f"An attempt was made to change/create a counter metric when metrics were not configured.")
370 return False
371 except Exception as e:
372 logging.getLogger(__package__).error(f"UNCAUGHT EXCEPTION:\n{e}")
373 return False
375def decrement_counter(counter_name, by=1, attributes: AttrDict={}) -> bool:
376 """
377 Decrements a up down counter with the given name and value. If applied to a regular counter it will log a warning and silently fail.
379 Will catch any exceptions generated by metric usage.
381 Args:
382 counter_name (str): The name of the counter.
383 by (int, optional): The value to decrement by. Defaults to 1. abs(by) is used to protect from negative numbers.
384 attributes (dict, optional): Additional attributes for the counter. Defaults to {}.
386 Returns:
387 bool: True if the counter was decremented successfully, False otherwise (logging the error).
388 """
389 if __ANACONDA_TELEMETRY_INITIALIZED is False:
390 logging.getLogger(__package__).error("Anaconda telemetry system not initialized.") # Since init didn't happen this is not exported in OTel!!!
391 return False
392 try:
393 return _AnacondaMetrics._instance.decrement_counter(counter_name, by, _AnacondaMetrics._instance._process_attributes(attributes))
394 except MetricsNotInitialized:
395 logging.getLogger(__package__).warning(f"An attempt was made to change/create a counter metric when metrics were not configured.")
396 return False
397 except Exception as e:
398 logging.getLogger(__package__).error(f"UNCAUGHT EXCEPTION:\n{e}")
399 return False
401@contextmanager
402def get_trace(name: str, attributes: AttrDict = {}, carrier: Dict[str,str] = None) -> Iterator[_ASpan]:
403 """
404 Create or continue a named trace (based on the 'carrier' parameter).
406 Use the function like a Python I/O object (keyword 'with') to ensure the span is closed properly.
408 Will catch any exceptions generated by tracing usage.
410 Args:
411 name (str): The name of the trace.
412 attributes (dict, optional): Additional attributes for the trace. Defaults to {}.
413 carrier (dict, optional): The carrier used to continue a trace context in the output data. Defaults to None.
415 Example:
416 with get_trace("my_trace_name", {"key": "value"}) as span:
417 # Do some work here
418 pass
419 # The span will be closed automatically when exiting the 'with' block.
421 Returns:
422 Iterator[Tracer]: An iterator for the tracer.
423 """
424 if __ANACONDA_TELEMETRY_INITIALIZED is False:
425 logging.getLogger(__package__).error("Anaconda telemetry system not initialized.") # Since init didn't happen this is not exported in OTel!!!
426 return None
428 try:
429 aspan = _AnacondaTrace._instance.get_span(name, _AnacondaTrace._instance._process_attributes(attributes), carrier)
430 except: # Trace is different than the other signals, there is no easy way to log and continue.
431 logging.getLogger(__package__).warning(f"Attempt to trace a with-block when tracing was not configured.")
432 aspan = _ASpan("UNKNOWN", span=None, noop=True)
433 try:
434 yield aspan
435 except Exception as e:
436 aspan.add_exception(e)
437 aspan.set_error_status()
438 _AnacondaTrace._instance.logger.error(f"Error in trace span {name}: {e}")
439 finally:
440 aspan._close()
442def get_telemetry_logger_handler() -> LoggingHandler:
443 """
444 Returns the telemetry logger handler. This lets the package user control how the application uses the telemetry logger.
445 Insert this handler into your named logger.
447 log = logging.getLogger("my_logger")
448 log.addHandler(get_telemetry_logger_handler())
450 Previously, this was injected into the root logger, but this turned out to be problematic for some applications that
451 wanted to control the logging configuration more precisely. This injection behavior is now disabled. If you wish to
452 inject the handler into the root logger, you can do so manually. See the Python logging documentation for more information.
454 Returns:
455 logging.Logger: The telemetry logger handler if logging was enabled via signal_types in initialize_telemetry,
456 otherwise this function returns None.
457 Raises:
458 RuntimeError: if `initialize_telemetry` has not been called
459 """
460 global __ANACONDA_TELEMETRY_INITIALIZED
461 if __ANACONDA_TELEMETRY_INITIALIZED is False:
462 logging.getLogger(__package__).error("Anaconda telemetry system not initialized.") # Since init didn't happen this is not exported in OTel!!!
463 raise RuntimeError("Anaconda telemetry system not initialized.")
464 if _AnacondaLogger._instance is not None:
465 return _AnacondaLogger._instance._get_log_handler()
466 return None # No logger handler available, logging not initialized or not configured.
468def send_event(body: str, event_name: str, attributes: AttrDict={}) -> bool:
469 """
470 Sends a log event directly to the OpenTelemetry pipeline without using Python's logging module.
471 This is useful when you want to export log telemetry but don't want the output mixing with
472 your application's output or developer logs.
474 Params:
475 body (str): the log message body
476 event_name (str): mandatory event name added to attributes
477 attributes (AttrDict): optional attributes dict
478 Returns:
479 bool: True if the event was sent, False if logging was not initialized
480 Raises:
481 RuntimeError: if `initialize_telemetry` has not been called
482 """
483 global __ANACONDA_TELEMETRY_INITIALIZED
484 if __ANACONDA_TELEMETRY_INITIALIZED is False:
485 logging.getLogger(__package__).error("Anaconda telemetry system not initialized.")
486 raise RuntimeError("Anaconda telemetry system not initialized.")
487 if _AnacondaLogger._instance is not None:
488 event_logger = _AnacondaLogger._instance._get_event_logger()
489 event_logger._send_event(body, event_name, _AnacondaLogger._instance._process_attributes(attributes))
490 return True
491 return False