Skip to content
Open
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
292 changes: 287 additions & 5 deletions datadog/_client.py
Original file line number Diff line number Diff line change
Expand Up @@ -4,23 +4,57 @@
import shutil
import subprocess
import time
from typing import Any, Dict, List, Literal, Optional, Tuple, Type, Union, cast
from typing import (
TYPE_CHECKING,
Any,
Awaitable,
Callable,
Dict,
List,
Literal,
Optional,
Sequence,
Tuple,
Type,
Union,
cast,
)


import ddtrace
from ddtrace.internal.http import HTTPConnection
from ddtrace.internal.utils.formats import asbool
from ddtrace.profiling import Profiler
from ddtrace.runtime import RuntimeMetrics
from ddtrace.trace import Span
from ddtrace._logger import DD_LOG_FORMAT

from ._metrics import MetricsClient
from ._logging import V2LogWriter


if TYPE_CHECKING:
from ddtrace.llmobs._experiment import Dataset
from ddtrace.llmobs._experiment import DatasetRecordNew
from ddtrace.llmobs._experiment import Experiment
from ddtrace.llmobs._experiment import ExperimentResult
from ddtrace.llmobs._experiment import SyncExperiment


logger = logging.getLogger(__name__)

TraceSampleRule = Tuple[str, str, float]
# The value of an evaluation metric, and the metric_type it implies.
EvalValue = Union[str, int, float, bool, Dict[str, Any]]
MetricType = Literal["categorical", "score", "boolean", "json"]
# A span to join an evaluation to: a live span, or the dict produced by
# ``export_span()`` (which is what crosses a process boundary).
EvalSpan = Union[Span, Dict[str, str]]
# Experiment building blocks. Kept loose here; ``datadog.llmobs`` re-exports
# the real ddtrace types for annotating user code.
Task = Callable[..., Any]
Evaluator = Any
DEFAULT_PROJECT_NAME = "default"
# recursive types aren't supported (yet): https://github.com/python/mypy/issues/731
# _JSON = Union[str, float, int, List["_JSON"], Dict[str, "_JSON"], None]

Expand Down Expand Up @@ -50,6 +84,7 @@ def __bool__(self):
llmobs_enabled=False,
llmobs_agentless_enabled=False,
llmobs_integrations_enabled=True,
llmobs_project_name=DEFAULT_PROJECT_NAME,
) # type: Dict[str, Any]


Expand Down Expand Up @@ -81,6 +116,8 @@ def __init__(
llmobs_ml_app=_sentinel, # type: Union[_Sentinel, str]
llmobs_agentless_enabled=_sentinel, # type: Union[_Sentinel, bool]
llmobs_integrations_enabled=_sentinel, # type: Union[_Sentinel, bool]
llmobs_app_key=_sentinel, # type: Union[_Sentinel, str]
llmobs_project_name=_sentinel, # type: Union[_Sentinel, str]
default_config=_DEFAULT_CONFIG, # type: Dict[str, Any]
):
# type: (...) -> None
Expand Down Expand Up @@ -241,6 +278,36 @@ def __init__(
)
self.llmobs_integrations_enabled = llmobs_integrations_enabled

# The app key is only needed by the eval APIs (datasets, experiments,
# publishing evaluators) which talk to the Datadog API directly, so an
# empty one is not an error until one of those is used.
if isinstance(llmobs_app_key, _Sentinel):
llmobs_app_key = os.getenv("DD_APP_KEY", "")
self.llmobs_app_key = cast(str, llmobs_app_key)

if isinstance(llmobs_project_name, _Sentinel):
llmobs_project_name = os.getenv(
"DD_LLMOBS_PROJECT_NAME", default_config["llmobs_project_name"]
)
self.llmobs_project_name = cast(str, llmobs_project_name)


def _infer_metric_type(label, value):
# type: (str, EvalValue) -> MetricType
# bool first: it is a subclass of int and would otherwise read as a score.
if isinstance(value, bool):
return "boolean"
if isinstance(value, (int, float)):
return "score"
if isinstance(value, str):
return "categorical"
if isinstance(value, dict):
return "json"
raise ValueError(
"Cannot infer the metric type of the %r evaluation from a %s value. "
"Pass metric_type= explicitly." % (label, type(value).__name__)
)


class DDAgent:
def __init__(self, version: str, config: DDConfig):
Expand Down Expand Up @@ -362,13 +429,15 @@ def __init__(
agentless_enabled=config.llmobs_agentless_enabled,
site=config.site,
api_key=config.api_key,
app_key=config.llmobs_app_key,
project_name=config.llmobs_project_name,
service=config.service,
env=config.env,
)
self._llmobs = LLMObs

def trace(self, *args, **kwargs):
# type: (...) -> ddtrace.Span
# type: (...) -> Span
return self._tracer.trace(*args, **kwargs)

def traced(self, *args, **kwargs):
Expand All @@ -386,6 +455,17 @@ def _require_llmobs(self):
)
return self._llmobs

def _require_llmobs_app_key(self):
# Datasets, experiments and remote evaluators are read/written through
# the Datadog API, which needs an app key on top of the api key.
llmobs = self._require_llmobs()
if not self._config.llmobs_app_key:
raise RuntimeError(
"An app key is required for the datasets/experiments API. Pass "
"llmobs_app_key= to DDConfig (or set DD_APP_KEY)."
)
return llmobs

# -- LLM Observability decorators ---------------------------------
# Use as ``@ddclient.workflow(name="…")`` / ``@ddclient.llm(...)`` etc.
# Each proxies to ``ddtrace.llmobs.decorators.<kind>``, which produces
Expand Down Expand Up @@ -437,12 +517,214 @@ def annotate(self, *args, **kwargs):
def annotation_context(self, *args, **kwargs):
return self._require_llmobs().annotation_context(*args, **kwargs)

def submit_evaluation(self, *args, **kwargs):
return self._require_llmobs().submit_evaluation(*args, **kwargs)

def export_span(self, *args, **kwargs):
return self._require_llmobs().export_span(*args, **kwargs)

# -- LLM Observability evaluations --------------------------------
# Evaluations come in two flavours: metrics submitted against spans that
# already happened (``submit_evaluation``), and experiments, which run a
# task over a dataset and score each row with evaluators.

def submit_evaluation(
self,
label: str,
value: EvalValue,
metric_type: Optional[MetricType] = None,
span: Optional[EvalSpan] = None,
span_with_tag_value: Optional[Dict[str, str]] = None,
tags: Optional[Dict[str, str]] = None,
metadata: Optional[Dict[str, object]] = None,
assessment: Optional[Literal["pass", "fail"]] = None,
reasoning: Optional[str] = None,
timestamp_ms: Optional[int] = None,
eval_scope: Literal["span", "trace"] = "span",
) -> None:
"""Attach an evaluation metric to a span.

``metric_type`` is inferred from ``value`` unless given, and ``span``
defaults to the span currently being traced, so the common case is
just ``ddclient.submit_evaluation("relevance", 0.9)``. A span from a
previous process can be joined by passing the dict that
``export_span()`` produced, or by tag with ``span_with_tag_value``.
"""
llmobs = self._require_llmobs()
if metric_type is None:
metric_type = _infer_metric_type(label, value)
if span is None and span_with_tag_value is None:
# export_span() raises rather than returning None when there is no
# LLM Obs span in scope (a plain traced span doesn't count).
from ddtrace.llmobs._llmobs import LLMObsExportSpanError

try:
span = llmobs.export_span()
except LLMObsExportSpanError:
span = None
if span is None:
raise ValueError(
"No LLM Obs span is currently active to attach the %r "
"evaluation to. Pass span= (a span, or the dict from "
"export_span()) or span_with_tag_value=." % label
)
elif isinstance(span, Span):
span = llmobs.export_span(span)
llmobs.submit_evaluation(
label=label,
metric_type=metric_type,
value=value,
span=span,
span_with_tag_value=span_with_tag_value,
tags=tags,
metadata=metadata,
assessment=assessment,
reasoning=reasoning,
timestamp_ms=timestamp_ms,
eval_scope=eval_scope,
)

def get_spans(self, *args, **kwargs) -> List[Dict[str, Any]]:
"""Query already-submitted LLM Obs spans, e.g. to evaluate them offline."""
return self._require_llmobs().get_spans(*args, **kwargs)

def publish_evaluator(
self,
evaluator: Evaluator,
eval_name: Optional[str] = None,
variable_mapping: Optional[Dict[str, str]] = None,
) -> Dict[str, str]:
"""Publish an evaluator so Datadog runs it against live spans."""
return self._require_llmobs_app_key().publish_evaluator(
evaluator=evaluator,
ml_app=self._config.llmobs_ml_app,
eval_name=eval_name,
variable_mapping=variable_mapping,
)

# -- LLM Observability datasets -----------------------------------

def create_dataset(
self,
name: str,
records: Optional[List["DatasetRecordNew"]] = None,
description: str = "",
project_name: Optional[str] = None,
**kwargs,
) -> "Dataset":
return self._require_llmobs_app_key().create_dataset(
dataset_name=name,
project_name=self._project_name(project_name),
description=description,
records=records,
**kwargs,
)

def create_dataset_from_csv(
self,
csv_path: str,
name: str,
input_data_columns: List[str],
expected_output_columns: Optional[List[str]] = None,
description: str = "",
project_name: Optional[str] = None,
**kwargs,
) -> "Dataset":
return self._require_llmobs_app_key().create_dataset_from_csv(
csv_path=csv_path,
dataset_name=name,
input_data_columns=input_data_columns,
expected_output_columns=expected_output_columns,
description=description,
project_name=self._project_name(project_name),
**kwargs,
)

def pull_dataset(
self,
name: str,
project_name: Optional[str] = None,
version: Optional[int] = None,
tags: Optional[List[str]] = None,
) -> "Dataset":
return self._require_llmobs_app_key().pull_dataset(
dataset_name=name,
project_name=self._project_name(project_name),
version=version,
tags=tags,
)

# -- LLM Observability experiments --------------------------------

def experiment(
self,
name: str,
task: Task,
dataset: "Dataset",
evaluators: Sequence[Evaluator],
description: str = "",
project_name: Optional[str] = None,
**kwargs,
) -> "SyncExperiment":
"""Build an experiment. Call ``.run()`` on it, or use ``run_experiment``."""
return self._require_llmobs_app_key().experiment(
name=name,
task=task,
dataset=dataset,
evaluators=evaluators,
description=description,
project_name=self._project_name(project_name),
**kwargs,
)

def async_experiment(
self,
name: str,
task: Callable[..., Awaitable[Any]],
dataset: "Dataset",
evaluators: Sequence[Evaluator],
description: str = "",
project_name: Optional[str] = None,
**kwargs,
) -> "Experiment":
"""``experiment`` for a coroutine task."""
return self._require_llmobs_app_key().async_experiment(
name=name,
task=task,
dataset=dataset,
evaluators=evaluators,
description=description,
project_name=self._project_name(project_name),
**kwargs,
)

def run_experiment(
self,
name: str,
task: Task,
dataset: "Dataset",
evaluators: Sequence[Evaluator],
jobs: int = 1,
raise_errors: bool = False,
sample_size: Optional[int] = None,
**kwargs,
) -> "ExperimentResult":
"""Build and run an experiment in one call, returning its results."""
experiment = self.experiment(
name=name,
task=task,
dataset=dataset,
evaluators=evaluators,
**kwargs,
)
return experiment.run(
jobs=jobs, raise_errors=raise_errors, sample_size=sample_size
)

def pull_experiment(self, experiment_id: str) -> "SyncExperiment":
return self._require_llmobs_app_key().pull_experiment(experiment_id)

def _project_name(self, project_name=None):
# type: (Optional[str]) -> str
return project_name or self._config.llmobs_project_name

def _dd_log(self, log_level, msg, tags=_sentinel):
# TODO: timestamp
log = {
Expand Down
Loading