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

1# -*- coding: utf-8 -*- 

2# SPDX-FileCopyrightText: 2025 Anaconda, Inc 

3# SPDX-License-Identifier: Apache-2.0 

4 

5# signals.py 

6""" 

7Anaconda Telemetry - Metrics Module 

8 

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""" 

13 

14import logging, socket, threading 

15from typing import Dict, Iterator, List, Optional 

16from contextlib import contextmanager 

17 

18from opentelemetry import trace, metrics 

19from opentelemetry.sdk.trace import TracerProvider 

20from opentelemetry.sdk.metrics import MeterProvider 

21 

22from ._compat import OTEL_PRIVATE_LOGS_AVAILABLE, LoggingHandler 

23 

24from .config import Configuration as Config 

25from .attributes import ResourceAttributes as Attributes 

26from .formatting import AttrDict 

27 

28from .common import MetricsNotInitialized 

29from .logging import _AnacondaLogger 

30from .metrics import _AnacondaMetrics 

31from .tracing import _AnacondaTrace, _ASpan 

32 

33 

34_SUPPRESSED_LOGGER_ROOTS = ('opentelemetry',) 

35 

36 

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) 

44 

45 

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 

68 

69__ANACONDA_TELEMETRY_INITIALIZED = False 

70__SIGNALS = None 

71__CONFIG = None 

72 

73 

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. 

81 

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. 

91 

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 

98 

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") 

105 

106 __CONFIG = config 

107 __SIGNALS = signal_types 

108 

109 if type(attributes.parameters) != dict: 

110 raise ValueError(f"The parameters attribute in ResourceAttributes must be a dictionary") 

111 

112 if not config._get_verbose_export_errors(): 

113 _suppress_otel_export_errors() 

114 

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... 

117 

118 # all params are the same currently so only write them once 

119 init_params = (config, attributes) 

120 

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 ) 

136 

137 # Initialize the telemetry system here 

138 if 'metrics' in signal_types: 

139 _AnacondaMetrics._instance = _AnacondaMetrics(*init_params) 

140 signal_type_count += 1 

141 

142 # Initialize tracing here... 

143 if 'tracing' in signal_types: 

144 _AnacondaTrace._instance = _AnacondaTrace(*init_params) 

145 signal_type_count += 1 

146 

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 

154 

155_SHUTDOWN_DONE = False 

156_shutdown_lock = threading.Lock() 

157 

158 

159def flush_telemetry() -> bool: 

160 """Force-flush the providers this package emits through. 

161 

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. 

166 

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 

180 

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 

188 

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 

200 

201 

202def shutdown_telemetry(*, timeout_seconds: Optional[float] = None) -> bool: 

203 """Flush all telemetry providers at process shutdown, optionally time-bounded. 

204 

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. 

210 

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). 

214 

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. 

218 

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. 

222 

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() 

250 

251 

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 

257 

258 Args: 

259 signal_type (str): signal type to update endpoint for. Supported values are 'logging', 'metrics', and 'tracing' 

260 

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 

276 

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 ) 

284 

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 

291 

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. 

297 

298 Will catch any exceptions generated by metric usage. 

299 

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 {}. 

304 

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 

319 

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). 

325 

326 Will catch any exceptions generated by metric usage. 

327 

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 {}. 

333 

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 

348 

349def increment_counter(counter_name, by=1, attributes: AttrDict={}) -> bool: 

350 """ 

351 Increments a counter or up down counter by the given parameter 'by'. 

352 

353 Will catch any exceptions generated by metric usage. 

354 

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 {}. 

359 

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 

374 

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. 

378 

379 Will catch any exceptions generated by metric usage. 

380 

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 {}. 

385 

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 

400 

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). 

405 

406 Use the function like a Python I/O object (keyword 'with') to ensure the span is closed properly. 

407 

408 Will catch any exceptions generated by tracing usage. 

409 

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. 

414 

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. 

420 

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 

427 

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() 

441 

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. 

446 

447 log = logging.getLogger("my_logger") 

448 log.addHandler(get_telemetry_logger_handler()) 

449 

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. 

453 

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. 

467 

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. 

473 

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