Skip to content
Merged
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
1 change: 1 addition & 0 deletions datadog_checks_base/changelog.d/23963.added
Original file line number Diff line number Diff line change
@@ -0,0 +1 @@
Add configuration discovery runtime (Service/Port types, candidate_ports, probing harness, and discovery entry points).
51 changes: 42 additions & 9 deletions datadog_checks_base/datadog_checks/base/checks/base.py
Original file line number Diff line number Diff line change
Expand Up @@ -10,6 +10,7 @@
import os
import re
from collections import deque
from collections.abc import Iterable
from os.path import basename
from pathlib import Path
from typing import (
Expand Down Expand Up @@ -65,6 +66,7 @@
import unicodedata as _module_unicodedata

from datadog_checks.base.utils.diagnose import Diagnosis
from datadog_checks.base.utils.discovery import Service
from datadog_checks.base.utils.http import RequestsWrapper
from datadog_checks.base.utils.metadata import MetadataManager

Expand Down Expand Up @@ -179,6 +181,33 @@ def __init_subclass__(cls, *args, **kwargs):
except Exception:
return cls

@classmethod
def generate_configs(cls, service: Service) -> Iterable[dict[str, Any]]:
"""
Yield candidate full configurations for service discovery.

Integrations can opt into config discovery by declaring a discovery
stanza in their spec and generating config_models.discovery.
"""
from datadog_checks.base.utils.discovery.probe import generated_discovery_candidates

return generated_discovery_candidates(cls, service)

@classmethod
def discover_config(cls, service_json: str) -> str:
"""
Return discovered configurations for an AD service payload.

The Agent calls this classmethod through rtloader. Candidate configs
are generated by ``generate_configs`` and accepted only when the real
check can run against their metric instances successfully. Returns the
first accepted candidate only (first-match-wins); remaining candidates
are not evaluated.
"""
from datadog_checks.base.utils.discovery.probe import run_discovery

return run_discovery(cls, service_json)

def __init__(self, *args, **kwargs):
# type: (*Any, **Any) -> None
"""
Expand Down Expand Up @@ -752,7 +781,7 @@ def submit_histogram_bucket(
if hostname is None:
hostname = ''

aggregator.submit_histogram_bucket(
self._aggregator().submit_histogram_bucket(
self,
self.check_id,
self._format_namespace(name, raw),
Expand All @@ -770,28 +799,28 @@ def database_monitoring_query_sample(self, raw_event):
if raw_event is None:
return

aggregator.submit_event_platform_event(self, self.check_id, to_native_string(raw_event), "dbm-samples")
self._aggregator().submit_event_platform_event(self, self.check_id, to_native_string(raw_event), "dbm-samples")

def database_monitoring_query_metrics(self, raw_event):
# type: (str) -> None
if raw_event is None:
return

aggregator.submit_event_platform_event(self, self.check_id, to_native_string(raw_event), "dbm-metrics")
self._aggregator().submit_event_platform_event(self, self.check_id, to_native_string(raw_event), "dbm-metrics")

def database_monitoring_query_activity(self, raw_event):
# type: (str) -> None
if raw_event is None:
return

aggregator.submit_event_platform_event(self, self.check_id, to_native_string(raw_event), "dbm-activity")
self._aggregator().submit_event_platform_event(self, self.check_id, to_native_string(raw_event), "dbm-activity")

def database_monitoring_metadata(self, raw_event):
# type: (str) -> None
if raw_event is None:
return

aggregator.submit_event_platform_event(self, self.check_id, to_native_string(raw_event), "dbm-metadata")
self._aggregator().submit_event_platform_event(self, self.check_id, to_native_string(raw_event), "dbm-metadata")

def event_platform_event(self, raw_event, event_track_type):
# type: (str | bytes, str) -> None
Expand All @@ -810,7 +839,7 @@ def event_platform_event(self, raw_event, event_track_type):
raw_event = bytes(raw_event)
elif not isinstance(raw_event, bytes):
raw_event = to_native_string(raw_event)
aggregator.submit_event_platform_event(self, self.check_id, raw_event, event_track_type)
self._aggregator().submit_event_platform_event(self, self.check_id, raw_event, event_track_type)

def submit_generic_resource(self, *, type, key, fields, include, seen_at=None, expire_at=None):
# type: (str, str, dict | None, dict, int | None, int | None) -> None
Expand Down Expand Up @@ -974,6 +1003,10 @@ def _metric_excluded(self, metric_name):

return self.exclude_metrics_pattern.search(metric_name) is not None

def _aggregator(self) -> Any:
"""Return the active aggregator: proxy during a discovery probe, module singleton otherwise."""
return getattr(self, '_discovery_aggregator', None) or aggregator

def _submit_metric(
self, mtype, name, value, tags=None, hostname=None, device_name=None, raw=False, flush_first_value=False
):
Expand Down Expand Up @@ -1012,7 +1045,7 @@ def _submit_metric(
self.warning(err_msg)
return

aggregator.submit_metric(self, self.check_id, mtype, name, value, tags, hostname, flush_first_value)
self._aggregator().submit_metric(self, self.check_id, mtype, name, value, tags, hostname, flush_first_value)

def gauge(self, name, value, tags=None, hostname=None, device_name=None, raw=False):
# type: (str, float, Sequence[str], str, str, bool) -> None
Expand Down Expand Up @@ -1233,7 +1266,7 @@ def service_check(self, name, status, tags=None, hostname=None, message=None, ra

message = self.sanitize(message)

aggregator.submit_service_check(
self._aggregator().submit_service_check(
self, self.check_id, self._format_namespace(name, raw), status, tags, hostname, message
)

Expand Down Expand Up @@ -1649,7 +1682,7 @@ def event(self, event):
if self.__NAMESPACE__:
event.setdefault('source_type_name', self.__NAMESPACE__)

aggregator.submit_event(self, self.check_id, event)
self._aggregator().submit_event(self, self.check_id, event)

def _normalize_tags_type(self, tags, device_name=None, metric_name=None):
# type: (Sequence[Union[None, str, bytes]], str, str) -> List[str]
Expand Down
Original file line number Diff line number Diff line change
@@ -1,6 +1,7 @@
# (C) Datadog, Inc. 2025-present
# All rights reserved
# Licensed under a 3-clause BSD style license (see LICENSE)
from .discovery import Discovery
from .discovery import Discovery, Port, Service, candidate_ports
from .strategies import discovery_strategy

__all__ = ['Discovery']
__all__ = ['Discovery', 'Port', 'Service', 'candidate_ports', 'discovery_strategy']
Original file line number Diff line number Diff line change
@@ -1,10 +1,33 @@
# (C) Datadog, Inc. 2023-present
# All rights reserved
# Licensed under a 3-clause BSD style license (see LICENSE)
from collections.abc import Iterable, Iterator

from pydantic import BaseModel, ConfigDict

from .cache import Cache
from .filter import Filter


class Port(BaseModel):
"""An Autodiscovery-exposed port on a service."""

model_config = ConfigDict(frozen=True)

number: int
name: str = ""


class Service(BaseModel):
"""An Autodiscovery-discovered service instance."""

model_config = ConfigDict(frozen=True)

id: str
host: str
ports: tuple[Port, ...] = ()


class Discovery:
def __init__(
self,
Expand All @@ -21,3 +44,19 @@ def __init__(
def get_items(self):
items = self._cache.get_items()
return self._filter.get_items(items)


def candidate_ports(service: Service, hints: Iterable[int]) -> Iterator[Port]:
"""Yield hinted ports first, then remaining service ports."""
by_number = {port.number: port for port in service.ports}
seen: set[int] = set()

for hint in hints:
if hint in by_number and hint not in seen:
seen.add(hint)
yield by_number[hint]

for port in service.ports:
if port.number not in seen:
seen.add(port.number)
yield port
Loading
Loading