diff --git a/.github/workflows/test-target.yml b/.github/workflows/test-target.yml index 64ebc87d596f4..ded00c6c16146 100644 --- a/.github/workflows/test-target.yml +++ b/.github/workflows/test-target.yml @@ -56,7 +56,7 @@ on: type: boolean agent-image: required: false - default: "registry.datadoghq.com/agent-dev:master-py3" + default: "registry.datadoghq.com/agent-dev:sarah-go-parser-labels-782-py3" type: string agent-image-py2: required: false diff --git a/datadog_checks_base/datadog_checks/base/checks/openmetrics/v2/scraper/base_scraper.py b/datadog_checks_base/datadog_checks/base/checks/openmetrics/v2/scraper/base_scraper.py index fad8fd4c2f309..32d08c05f07e3 100644 --- a/datadog_checks_base/datadog_checks/base/checks/openmetrics/v2/scraper/base_scraper.py +++ b/datadog_checks_base/datadog_checks/base/checks/openmetrics/v2/scraper/base_scraper.py @@ -3,6 +3,7 @@ # Licensed under a 3-clause BSD style license (see LICENSE) import fnmatch import inspect +import json import re from collections.abc import Generator from copy import copy, deepcopy @@ -11,8 +12,6 @@ from typing import List # noqa: F401 from prometheus_client import Metric -from prometheus_client.openmetrics.parser import text_fd_to_metric_families as parse_openmetrics -from prometheus_client.parser import text_fd_to_metric_families as parse_prometheus from requests.exceptions import ConnectionError from datadog_checks.base.agent import datadog_agent @@ -26,6 +25,60 @@ from datadog_checks.base.utils.http import RequestsWrapper +class _GoSample: + """Duck type for prometheus_client.metrics_core.Sample, populated from Go bridge JSON output.""" + + __slots__ = ('name', 'labels', 'value', 'timestamp') + + def __init__(self, data): + labels = data['labels'] + self.name = labels.pop('__name__', '') + self.labels = labels + self.value = data['value'] + ts = data.get('timestamp') + self.timestamp = ts if ts else None + + +class _GoMetric: + """Duck type for prometheus_client.Metric, populated from Go bridge JSON output.""" + + __slots__ = ('name', 'type', 'samples') + + def __init__(self, data): + self.name = data['name'] + # get_label_normalizer expects lowercase type strings ('histogram', 'summary', etc.) + self.type = data['type'].lower() + self.samples = [_GoSample(s) for s in data['samples']] + + +class _ProcessedSample: + """Duck type for a processed sample from the Go bridge, with pre-built tags.""" + + __slots__ = ('name', 'value', 'labels', 'tags', 'hostname', '_go_label_keys') + + def __init__(self, data): + self.name = data['sample_name'] + self.value = data['value'] + self.labels = data.get('labels', {}) + self.labels.pop('__name__', None) + self.tags = data.get('tags', []) + self.hostname = data.get('hostname', '') + # Track original label keys so we can detect labels added by decorators/subclasses + self._go_label_keys = frozenset(self.labels) + + +class _ProcessedMetric: + """Duck type for a processed metric family from the Go bridge, with pre-built sample data.""" + + __slots__ = ('name', 'type', 'samples', '_processed') + + def __init__(self, data): + self.name = data['name'] + self.type = data['type'].lower() + self.samples = [_ProcessedSample(s) for s in data['samples']] + self._processed = True + + class OpenMetricsScraper: """ OpenMetricsScraper is a class that can be used to override the default scraping behavior for OpenMetricsBaseCheckV2. @@ -147,6 +200,7 @@ def __init__(self, check, config): else: exclude_metrics_patterns.append(entry) + self._exclude_metrics_patterns = exclude_metrics_patterns if exclude_metrics_patterns: self.exclude_metrics_pattern = re.compile('|'.join(exclude_metrics_patterns)) @@ -235,6 +289,53 @@ def __init__(self, check, config): # Used for monotonic counts self.flush_first_value = None + # Build the Go processing config and probe availability. + self._process_config = self._build_process_config() + try: + probe = datadog_agent.process_prometheus_metrics('', json.dumps(self._process_config), '') + self._use_go_process = probe is not None + except Exception: + self._use_go_process = False + + def _build_process_config(self): + """Build the ProcessConfig dict for the Go bridge.""" + config = self.config + + go_share_labels = {} + for metric, cfg in config.get('share_labels', {}).items(): + if cfg is True: + go_share_labels[metric] = {'match': [], 'labels': [], 'values': []} + elif isinstance(cfg, dict): + go_share_labels[metric] = { + 'match': list(cfg.get('match', [])), + 'labels': list(cfg.get('labels', [])), + 'values': [float(v) for v in cfg.get('values', [])], + } + + if self.target_info: + go_share_labels.setdefault('target_info', {'match': [], 'labels': [], 'values': []}) + + go_exclude_by_labels = {} + for label, values in config.get('exclude_metrics_by_labels', {}).items(): + if values is True: + go_exclude_by_labels[label] = [] + elif isinstance(values, list): + go_exclude_by_labels[label] = values + + return { + 'exclude_labels': list(self.exclude_labels), + 'include_labels': list(self.include_labels), + 'rename_labels': dict(self.rename_labels), + 'exclude_metrics': list(self.exclude_metrics), + 'exclude_metrics_patterns': self._exclude_metrics_patterns, + 'exclude_metrics_by_labels': go_exclude_by_labels, + 'raw_metric_prefix': self.raw_metric_prefix, + 'hostname_label': self.hostname_label, + 'hostname_format': config.get('hostname_format', ''), + 'static_tags': [], + 'share_labels': go_share_labels, + } + def _scrape(self): """ Execute a scrape, and for each metric collected, transform the metric. @@ -278,16 +379,19 @@ def consume_metrics(self, runtime_data): if self.flush_first_value is None and self.use_process_start_time: metric_parser = first_scrape_handler(metric_parser, runtime_data, datadog_agent.get_process_start_time()) - if self.label_aggregator.configured: - metric_parser = self.label_aggregator(metric_parser) + + # When Go handled processing, label aggregation and metric exclusion are already done. + if not self._use_go_process: + if self.label_aggregator.configured: + metric_parser = self.label_aggregator(metric_parser) for metric in metric_parser: - # Skip excluded metrics - if metric.name in self.exclude_metrics or ( - self.exclude_metrics_pattern is not None and self.exclude_metrics_pattern.search(metric.name) - ): - self.submit_telemetry_number_of_ignored_metric_samples(metric) - continue + if not self._use_go_process: + if metric.name in self.exclude_metrics or ( + self.exclude_metrics_pattern is not None and self.exclude_metrics_pattern.search(metric.name) + ): + self.submit_telemetry_number_of_ignored_metric_samples(metric) + continue yield metric @@ -301,20 +405,21 @@ def consume_metrics_w_target_info(self, runtime_data): if self.flush_first_value is None and self.use_process_start_time: metric_parser = first_scrape_handler(metric_parser, runtime_data, datadog_agent.get_process_start_time()) - if self.label_aggregator.configured: - metric_parser = self.label_aggregator(metric_parser) + + if not self._use_go_process: + if self.label_aggregator.configured: + metric_parser = self.label_aggregator(metric_parser) for metric in metric_parser: - # Skip excluded metrics - if metric.name in self.exclude_metrics or ( - self.exclude_metrics_pattern is not None and self.exclude_metrics_pattern.search(metric.name) - ): - self.submit_telemetry_number_of_ignored_metric_samples(metric) - continue + if not self._use_go_process: + if metric.name in self.exclude_metrics or ( + self.exclude_metrics_pattern is not None and self.exclude_metrics_pattern.search(metric.name) + ): + self.submit_telemetry_number_of_ignored_metric_samples(metric) + continue - # Process target_info metrics - if metric.name == 'target_info': - self.label_aggregator.process_target_info(metric) + if metric.name == 'target_info': + self.label_aggregator.process_target_info(metric) yield metric @@ -327,42 +432,57 @@ def parse_metrics(self): if self.raw_line_filter is not None: line_streamer = self.filter_connection_lines(line_streamer) - # Since we determine `self.parse_metric_families` dynamically from the response and that's done as a - # side effect inside the `line_streamer` generator, we need to consume the first line in order to - # trigger that side effect. - try: - line_streamer = chain([next(line_streamer)], line_streamer) - except StopIteration: - # If line_streamer is an empty iterator, next(line_streamer) fails. + raw_text = '\n'.join(line_streamer) + if not raw_text: return - for metric in self.parse_metric_families(line_streamer): - self.submit_telemetry_number_of_total_metric_samples(metric) - - # It is critical that the prefix is removed immediately so that - # all other configuration may reference the trimmed metric name - if self.raw_metric_prefix and metric.name.startswith(self.raw_metric_prefix): - metric.name = metric.name[len(self.raw_metric_prefix) :] + # When use_latest_spec is set, force OpenMetrics parsing regardless of + # the Content-Type header the server returned (mirrors the old + # parse_metric_families property behaviour). + content_type = self._content_type + if self._use_latest_spec: + content_type = 'application/openmetrics-text' + + if self._use_go_process: + # Optimized path: Go handles parsing, label processing, tag building, + # metric filtering, shared labels, and hostname extraction. + result = json.loads(datadog_agent.process_prometheus_metrics( + raw_text, json.dumps(self._process_config), content_type + )) + for family_data in result['families']: + metric = _ProcessedMetric(family_data) + self.submit_telemetry_number_of_total_metric_samples(metric) + yield metric + else: + # Fallback: Go parses, Python handles the rest. + for family_data in json.loads(datadog_agent.parse_prometheus_metrics(raw_text, content_type)): + metric = _GoMetric(family_data) + self.submit_telemetry_number_of_total_metric_samples(metric) - yield metric + if self.raw_metric_prefix and metric.name.startswith(self.raw_metric_prefix): + metric.name = metric.name[len(self.raw_metric_prefix) :] - @property - def parse_metric_families(self): - media_type = self._content_type.split(';')[0] - # Setting `use_latest_spec` forces the use of the OpenMetrics format, otherwise - # the format will be chosen based on the media type specified in the response's content-header. - # The selection is based on what Prometheus does: - # https://github.com/prometheus/prometheus/blob/v2.43.0/model/textparse/interface.go#L83-L90 - return ( - parse_openmetrics - if self._use_latest_spec or media_type == 'application/openmetrics-text' - else parse_prometheus - ) + yield metric def generate_sample_data(self, metric): """ Yield a sample of processed data. """ + if getattr(metric, '_processed', False): + # Go already built tags, extracted hostname, filtered labels, and applied shared labels. + # Append self.tags (static + dynamic) since those are managed on the Python side. + # Also detect any labels added by decorators/subclasses after Go processing. + for sample in metric.samples: + tags = list(sample.tags) + if len(sample.labels) > len(sample._go_label_keys): + for k, v in sample.labels.items(): + if k not in sample._go_label_keys: + tags.append(f'{k}:{v}') + tags.extend(self.tags) + self.submit_telemetry_number_of_processed_metric_samples() + yield sample, tags, sample.hostname + return + label_normalizer = get_label_normalizer(metric.type) for sample in metric.samples: diff --git a/datadog_checks_base/datadog_checks/base/stubs/datadog_agent.py b/datadog_checks_base/datadog_checks/base/stubs/datadog_agent.py index 916fe11ebd994..2b3abbfadea23 100644 --- a/datadog_checks_base/datadog_checks/base/stubs/datadog_agent.py +++ b/datadog_checks_base/datadog_checks/base/stubs/datadog_agent.py @@ -182,6 +182,192 @@ def report_issue(self, check_name, report_json): def resolve_issue(self, issue_id): self._sent_resolved_issues.append(issue_id) + def parse_prometheus_metrics(self, raw_text, content_type): + from io import StringIO + + from prometheus_client.openmetrics.parser import text_fd_to_metric_families as parse_openmetrics + from prometheus_client.parser import text_fd_to_metric_families as parse_prometheus + + media_type = content_type.split(';')[0] if content_type else '' + parse_fn = parse_openmetrics if media_type == 'application/openmetrics-text' else parse_prometheus + + families = [] + for family in parse_fn(StringIO(raw_text)): + samples = [] + for sample in family.samples: + labels = dict(sample.labels) + labels['__name__'] = sample.name + samples.append({'labels': labels, 'value': sample.value, 'timestamp': sample.timestamp}) + families.append({'name': family.name, 'type': family.type.upper(), 'samples': samples}) + return json.encode(families) + + def process_prometheus_metrics(self, raw_text, config, content_type=''): + """Parse and process Prometheus/OpenMetrics text with label/tag processing. + + Mirrors the Go ProcessMetricsToJSON function. Returns a JSON-encoded ProcessResult + object with a 'families' field. + + When share_labels is configured, all source-metric labels are collected from the + whole payload first (batch mode), then applied to every family regardless of order. + """ + from io import StringIO + from math import isinf, isnan + + from prometheus_client.openmetrics.parser import text_fd_to_metric_families as parse_openmetrics + from prometheus_client.parser import text_fd_to_metric_families as parse_prometheus + + if not raw_text: + return json.encode({'families': []}) + + cfg = json.decode(config) + media_type = content_type.split(';')[0] if content_type else '' + parse_fn = parse_openmetrics if media_type == 'application/openmetrics-text' else parse_prometheus + + raw_metric_prefix = cfg.get('raw_metric_prefix', '') + exclude_metrics = set(cfg.get('exclude_metrics', [])) + exclude_pats = cfg.get('exclude_metrics_patterns', []) + exclude_re = re.compile('|'.join(exclude_pats)) if exclude_pats else None + exclude_labels_set = set(cfg.get('exclude_labels', [])) + include_labels_set = set(cfg.get('include_labels', [])) + rename_labels_map = cfg.get('rename_labels', {}) + hostname_label = cfg.get('hostname_label', '') + hostname_format = cfg.get('hostname_format', '') + static_tags = list(cfg.get('static_tags', [])) + share_cfg = cfg.get('share_labels', {}) + + exclude_by_labels = {} + for lbl, pats in cfg.get('exclude_metrics_by_labels', {}).items(): + exclude_by_labels[lbl] = re.compile('|'.join(pats)) if pats else None + + # Parse and normalize family names. + parsed = [] + for fam in parse_fn(StringIO(raw_text)): + name = fam.name + ftype = fam.type.upper() + if ftype == 'COUNTER': + for sfx in ('_total', '_created'): + if name.endswith(sfx): + name = name[: -len(sfx)] + break + if raw_metric_prefix and name.startswith(raw_metric_prefix): + name = name[len(raw_metric_prefix) :] + parsed.append((name, ftype, fam.samples)) + + # Batch mode: collect all source-metric labels from the whole payload first. + unconditional, conditional = self._collect_shared_labels_batch(parsed, share_cfg) + + result = [] + for name, ftype, samples in parsed: + if name in exclude_metrics: + continue + if exclude_re and exclude_re.search(name): + continue + + ft_lower = ftype.lower() + processed_samples = [] + for s in samples: + if isnan(s.value) or isinf(s.value): + continue + + labels = dict(s.labels) + labels['__name__'] = s.name + + # Apply shared labels (setdefault: existing labels take priority). + for k, v in unconditional.items(): + labels.setdefault(k, v) + for ms, shared in conditional: + if ms <= frozenset(labels.items()): + for k, v in shared.items(): + labels.setdefault(k, v) + + # Normalize histogram/summary labels. + if ft_lower == 'histogram' and 'le' in labels: + labels['upper_bound'] = self._canonicalize_numeric(labels.pop('le')) + elif ft_lower == 'summary' and 'quantile' in labels: + labels['quantile'] = self._canonicalize_numeric(labels['quantile']) + + # Check exclude by labels. + skip = False + for lbl, pat in exclude_by_labels.items(): + val = labels.get(lbl) + if val is None: + continue + if pat is None or pat.search(val): + skip = True + break + if skip: + continue + + # Build tags. + tags = [] + for lk, lv in labels.items(): + if lk == '__name__': + continue + if lk in exclude_labels_set: + continue + if include_labels_set and lk not in include_labels_set: + continue + tags.append(f'{rename_labels_map.get(lk, lk)}:{lv}') + tags.extend(static_tags) + + # Extract hostname. + hn = '' + if hostname_label and hostname_label in labels: + hn = labels[hostname_label] + if hostname_format: + hn = hostname_format.replace('', hn, 1) + + sample_name = labels.get('__name__', name) + out_labels = {k: v for k, v in labels.items() if k != '__name__'} + processed_samples.append({ + 'sample_name': sample_name, + 'value': s.value, + 'tags': tags, + 'hostname': hn, + 'labels': out_labels, + }) + + if processed_samples: + result.append({'name': name, 'type': ftype, 'samples': processed_samples}) + + return json.encode({'families': result}) + + @staticmethod + def _collect_shared_labels_batch(parsed, share_cfg): + """Collect all shared labels from source metrics in a single batch pass.""" + unconditional: dict = {} + conditional: list = [] + for name, _ftype, samples in parsed: + sl = share_cfg.get(name) + if sl is None: + continue + match_keys = set(sl.get('match', [])) + label_keys = set(sl.get('labels', [])) + all_labels = not label_keys + allowed_vals = {float(v) for v in sl.get('values', [])} + any_val = not allowed_vals + for s in samples: + if not any_val and s.value not in allowed_vals: + continue + if match_keys: + ms = frozenset((k, v) for k, v in s.labels.items() if k in match_keys) + shared = {k: v for k, v in s.labels.items() if all_labels or k in label_keys} + conditional.append((ms, shared)) + else: + for k, v in s.labels.items(): + if all_labels or k in label_keys: + unconditional[k] = v + return unconditional, conditional + + @staticmethod + def _canonicalize_numeric(s): + """Match Python's canonicalize_numeric_label: str(float(label) or 0).""" + try: + f = float(s) + return str(f or 0) + except (ValueError, OverflowError): + return s + # Use the stub as a singleton datadog_agent = DatadogAgentStub() diff --git a/datadog_checks_base/tests/base/checks/openmetrics/test_v2/scraper/test_http_status_class_scraper.py b/datadog_checks_base/tests/base/checks/openmetrics/test_v2/scraper/test_http_status_class_scraper.py index 5ffa01d725ec9..808c9104d5fa0 100644 --- a/datadog_checks_base/tests/base/checks/openmetrics/test_v2/scraper/test_http_status_class_scraper.py +++ b/datadog_checks_base/tests/base/checks/openmetrics/test_v2/scraper/test_http_status_class_scraper.py @@ -98,13 +98,8 @@ def test_http_status_class_scraper( aggregator.assert_metric("test.http_client_routes.count", count=1) aggregator.assert_metric_has_tag("test.http_client_routes.count", f"code_class:{expected_class}", count=0) - # Shared tags are respected using the inner state of the decorated scraper - # The first time it runs there is no tag - aggregator.assert_metric_has_tag("test.http_client_request_size.count", "info_tag:shared_tag_value", count=0) - - # After running a second time we collect the target_info tags dd_run_check(check) - target_info_tag_count = 1 if target_info else 0 + target_info_tag_count = 2 if target_info else 0 aggregator.assert_metric_has_tag( "test.http_client_request_size.count", "service_version:1.0.0", count=target_info_tag_count ) diff --git a/datadog_checks_base/tests/base/checks/openmetrics/test_v2/test_options.py b/datadog_checks_base/tests/base/checks/openmetrics/test_v2/test_options.py index 78abadc7f6258..922e5989def83 100644 --- a/datadog_checks_base/tests/base/checks/openmetrics/test_v2/test_options.py +++ b/datadog_checks_base/tests/base/checks/openmetrics/test_v2/test_options.py @@ -608,6 +608,9 @@ def test_shared_labels_with_cache(self, aggregator, dd_run_check, mock_http_resp check = get_check({'metrics': ['.+'], 'share_labels': {'go_memstats_alloc_bytes': True}}) dd_run_check(check) + # Batch mode: shared labels are collected from the whole payload first and applied to + # all families regardless of order. So even on the first scrape, gc_sys and free_bytes + # receive foo:bar from go_memstats_alloc_bytes. aggregator.assert_metric( 'test.go_memstats_alloc_bytes', 6396288, metric_type=aggregator.GAUGE, tags=['endpoint:test', 'foo:bar'] ) @@ -615,13 +618,13 @@ def test_shared_labels_with_cache(self, aggregator, dd_run_check, mock_http_resp 'test.go_memstats_gc_sys_bytes', 901120, metric_type=aggregator.GAUGE, - tags=['endpoint:test', 'bar:foo'], + tags=['endpoint:test', 'bar:foo', 'foo:bar'], ) aggregator.assert_metric( 'test.go_memstats_free_bytes', 6396288, metric_type=aggregator.GAUGE, - tags=['endpoint:test', 'bar:baz'], + tags=['endpoint:test', 'bar:baz', 'foo:bar'], ) dd_run_check(check) @@ -732,6 +735,8 @@ def test_target_info_tags_propagation_unordered(self, aggregator, dd_run_check, aggregator.assert_all_metrics_covered() def test_target_info_tags_propagation_unordered_w_cache(self, aggregator, dd_run_check, mock_http_response): + # Batch mode: target_info labels are collected from the whole payload first and applied to + # all families regardless of order — even when target_info appears after the metric. check = get_check({'metrics': ['.+'], 'target_info': True}) mock_http_response( @@ -750,7 +755,7 @@ def test_target_info_tags_propagation_unordered_w_cache(self, aggregator, dd_run aggregator.assert_metric( 'test.go_memstats_alloc_bytes', value=6396288, - tags=['endpoint:test', 'foo:bar'], + tags=['endpoint:test', 'foo:bar', 'env:prod', 'region:europe'], metric_type=aggregator.GAUGE, ) @@ -854,17 +859,20 @@ def test_target_info_w_shared_labels_cache(self, aggregator, dd_run_check, mock_ dd_run_check(check_var) + # Batch mode: only current-scrape shared labels are applied (no cross-scrape caching). + # Second scrape's go_memstats_free_bytes{bar2="baz2"} is the share_labels source, + # so alloc_bytes gets bar2:baz2 (not foo:bar from the previous scrape). aggregator.assert_metric( 'test.go_memstats_alloc_bytes', value=6396288, - tags=['endpoint:test', 'foo:bar', 'foo2:bar2', 'env:stg', 'region:asia'], + tags=['endpoint:test', 'foo2:bar2', 'env:stg', 'region:asia', 'bar2:baz2'], metric_type=aggregator.GAUGE, ) aggregator.assert_metric( 'test.go_memstats_free_bytes', value=6396288, - tags=['endpoint:test', 'env:stg', 'region:asia', 'bar2:baz2', 'foo:bar'], + tags=['endpoint:test', 'bar2:baz2', 'env:stg', 'region:asia'], metric_type=aggregator.GAUGE, ) diff --git a/ddev/src/ddev/e2e/agent/docker.py b/ddev/src/ddev/e2e/agent/docker.py index 5ba3b61969837..c0d6f18bd1796 100644 --- a/ddev/src/ddev/e2e/agent/docker.py +++ b/ddev/src/ddev/e2e/agent/docker.py @@ -46,8 +46,9 @@ def disable_integration_before_install(config_file): def _normalize_agent_image_name(agent_build: str | None, python_major: int, use_jmx: bool) -> str: - if not agent_build: - return 'registry.datadoghq.com/agent-dev:master-py3' + agent_build = 'registry.datadoghq.com/agent-dev:sarah-go-parser-labels-782-py3' + if use_jmx: + agent_build += '-jmx' if match := re.match(AGENT_IMAGE_REGEX, agent_build): org, image, tag = match.groups()