Coverage for anaconda_opentelemetry/metrics.py: 88%

112 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# metrics.py 

6""" 

7Anaconda Telemetry - Metrics signal class. 

8""" 

9 

10import re 

11from typing import Dict, Any 

12 

13from opentelemetry import metrics 

14from opentelemetry.sdk.metrics import MeterProvider, Counter, UpDownCounter, Histogram, ObservableCounter, ObservableUpDownCounter 

15from opentelemetry.sdk.metrics.export import PeriodicExportingMetricReader, ConsoleMetricExporter, AggregationTemporality 

16 

17from ._compat import GaugeInstrument 

18 

19from .common import _AnacondaCommon, MetricsNotInitialized 

20from .config import Configuration as Config 

21from .attributes import ResourceAttributes as Attributes 

22from .exporter_shim import OTLPMetricExporterShim 

23from .formatting import AttrDict 

24 

25 

26class _AnacondaMetrics(_AnacondaCommon): 

27 # Singleton instance (internal only); provide a single instance of the metrics class 

28 _instance = None 

29 

30 _default_temporality: dict[type,AggregationTemporality] = { 

31 Counter: AggregationTemporality.DELTA, 

32 ObservableCounter: AggregationTemporality.DELTA, 

33 Histogram: AggregationTemporality.CUMULATIVE, 

34 UpDownCounter: AggregationTemporality.DELTA, 

35 ObservableUpDownCounter: AggregationTemporality.CUMULATIVE, 

36 GaugeInstrument: AggregationTemporality.CUMULATIVE, # A gauge is a last-value metric; DELTA is meaningless for it. 

37 } 

38 

39 _cumulative_temporality: dict[type,AggregationTemporality] = { 

40 Counter: AggregationTemporality.CUMULATIVE, 

41 ObservableCounter: AggregationTemporality.CUMULATIVE, 

42 Histogram: AggregationTemporality.CUMULATIVE, 

43 UpDownCounter: AggregationTemporality.CUMULATIVE, 

44 ObservableUpDownCounter: AggregationTemporality.CUMULATIVE, 

45 GaugeInstrument: AggregationTemporality.CUMULATIVE, 

46 } 

47 

48 _temporalityValue: dict[bool,str] = { 

49 False: "DELTA", 

50 True: "CUMULATIVE" 

51 } 

52 

53 def __init__(self, config: Config, attributes: Attributes): 

54 super().__init__(config, attributes) 

55 

56 self.metrics_endpoint = config._get_metrics_endpoint() 

57 self.telemetry_export_interval_millis = config._get_metrics_export_interval_ms() 

58 self.counter_objects: Dict[str, Any] = {} 

59 self.up_down_counter_objects: Dict[str, Any] = {} 

60 self.histogram_objects: Dict[str, Any] = {} 

61 self.gauge_objects: Dict[str, Any] = {} 

62 

63 self.meter = self._setup_metrics(config) 

64 self.create_dispatcher = { 

65 'simple_counter': self.meter.create_counter, 

66 'simple_up_down_counter': self.meter.create_up_down_counter, 

67 'histogram': self.meter.create_histogram, 

68 'gauge': self.meter.create_gauge 

69 } 

70 self.type_list = { 

71 'simple_counter': self.counter_objects, 

72 'simple_up_down_counter': self.up_down_counter_objects, 

73 'histogram': self.histogram_objects, 

74 'gauge': self.gauge_objects 

75 } 

76 

77 def _setup_metrics(self, config: Config) -> metrics.Meter: 

78 if self.use_console_exporters: 

79 exporter = ConsoleMetricExporter(preferred_temporality=self._get_temporality()) 

80 else: 

81 auth_token = config._get_auth_token_metrics() 

82 headers: Dict[str, str] = {} 

83 if auth_token is not None: 

84 headers['authorization'] = f'Bearer {auth_token}' 

85 if config._get_request_protocol_metrics() in ['grpc', 'grpcs']: # gRPC 

86 from opentelemetry.exporter.otlp.proto.grpc.metric_exporter import OTLPMetricExporter as OTLPMetricExportergRPC 

87 insecure = not config._get_TLS_metrics() 

88 exporter = OTLPMetricExporterShim( 

89 OTLPMetricExportergRPC, 

90 endpoint=self.metrics_endpoint, 

91 insecure=insecure, 

92 credentials=config._get_ca_cert_metrics() if not insecure else None, 

93 headers=headers, 

94 preferred_temporality=self._get_temporality() 

95 ) 

96 else: # HTTP 

97 from opentelemetry.exporter.otlp.proto.http.metric_exporter import OTLPMetricExporter as OTLPMetricExporterHTTP 

98 http_kwargs = self._build_http_exporter_kwargs( 

99 'metrics', self.metrics_endpoint, headers, 

100 preferred_temporality=self._get_temporality() 

101 ) 

102 exporter = OTLPMetricExporterShim( 

103 OTLPMetricExporterHTTP, 

104 **http_kwargs 

105 ) 

106 

107 self.exporter = exporter 

108 self.metric_reader = PeriodicExportingMetricReader(self.exporter, export_interval_millis=self.telemetry_export_interval_millis) 

109 # Create and set meter provider 

110 meter_provider = MeterProvider( 

111 resource=self.resource, 

112 metric_readers=[self.metric_reader], 

113 shutdown_on_exit=self._shutdown_on_exit 

114 ) 

115 self._provider = meter_provider 

116 try: 

117 metrics.set_meter_provider(meter_provider) 

118 except Exception as e: 

119 self.logger.warning(f"The metrics provider was previously set and will take precidence over this call.") 

120 # Get meter for this service 

121 return metrics.get_meter(self.service_name, self.service_version) 

122 

123 def _get_temporality(self) -> dict[type,AggregationTemporality]: 

124 if self._config._get_use_cumulative_metrics() == True: 

125 return _AnacondaMetrics._cumulative_temporality 

126 return _AnacondaMetrics._default_temporality 

127 

128 def _check_for_metric(self, metric_name: str, metric_type: str) -> bool: 

129 bucket_list = self.type_list.get(metric_type, None) 

130 if bucket_list is None: 

131 return False 

132 return bucket_list.get(metric_name, None) is not None 

133 

134 def _get_or_create_metric(self, metric_name: str, metric_type: str = 'simple_up_down_counter', units: str = '#', description='No description.') -> Any: 

135 bucket_list = self.type_list.get(metric_type, None) 

136 if bucket_list is None: 

137 raise MetricsNotInitialized(f"Metric type '{metric_type}' is unknown!") 

138 metric = bucket_list.get(metric_name, None) 

139 if metric is None: 

140 if not re.fullmatch(r"^[A-Za-z][A-Za-z_0-9.]*[A-Za-z_0-9]$", metric_name): 

141 self.logger.warning(f"Metric {metric_name} does not match valid regex: r\"^[A-Za-z][A-Za-z_0-9.]*[A-Za-z_0-9]$\"") 

142 return None 

143 create = self.create_dispatcher.get(metric_type, None) 

144 if create is None: 

145 self.logger.warning(f"Metric '{metric_name}' has an invalid type '{metric_type}'; cannot create metric.") 

146 return None 

147 metric = create( 

148 metric_name, 

149 unit=units, 

150 description=description 

151 ) 

152 if metric is None: 

153 self.logger.error(f"Failed to create metric '{metric_name}'!") 

154 bucket_list[metric_name] = metric 

155 return metric 

156 

157 def record_histogram(self, metric_name, value, attributes: AttrDict={}) -> bool: 

158 # Record a histogram metric with the given name and value. 

159 metric = self._get_or_create_metric(metric_name, metric_type='histogram', units='#', description='Dynamically create histogram metric.') 

160 if metric is None: 

161 self.logger.error(f"Metric '{metric_name}' failed to be created.") 

162 return False 

163 metric.record(value, attributes) 

164 return True 

165 

166 def set_gauge(self, metric_name, value, attributes: AttrDict={}) -> bool: 

167 # Set a gauge metric with the given name to the given value; the last value set wins. 

168 # Unlike the other instruments, OTel's Gauge.set() does no arithmetic on the value, so a 

169 # non-numeric value is accepted here and only fails later during serialization -- which 

170 # discards the whole export batch, not just this metric. Reject it up front instead. 

171 if isinstance(value, bool) or not isinstance(value, (int, float)): 

172 self.logger.error(f"Metric '{metric_name}' gauge value must be an int or float, not {type(value).__name__}.") 

173 return False 

174 metric = self._get_or_create_metric(metric_name, metric_type='gauge', units='#', description='Dynamically create gauge metric.') 

175 if metric is None: 

176 self.logger.error(f"Metric '{metric_name}' failed to be created.") 

177 return False 

178 metric.set(value, attributes) 

179 return True 

180 

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

182 # Increment a counter with the given name by the 'by' parameter. abs(by) is used. 

183 metric = None 

184 if self._check_for_metric(metric_name=counter_name, metric_type='simple_counter'): 

185 metric = self._get_or_create_metric(counter_name, metric_type='simple_counter') 

186 if metric is None: 

187 metric = self._get_or_create_metric(counter_name, metric_type='simple_up_down_counter') 

188 if metric is None: 

189 self.logger.error(f"Metric '{counter_name}' failed to be created.") 

190 return False 

191 metric.add(abs(by), attributes) 

192 return True 

193 

194 def decrement_counter(self, counter_name, by=1, attributes:AttrDict={}) -> bool: 

195 # Decrement a up down counter with the given name by the 'by' parameter. abs(by) is used. 

196 metric = self._get_or_create_metric(counter_name) 

197 if metric is None: 

198 self.logger.error(f"Metric '{counter_name}' failed to be created.") 

199 return False 

200 metric.add(-abs(by), attributes) 

201 return True