-
Notifications
You must be signed in to change notification settings - Fork 527
Expand file tree
/
Copy path_llmobs.py
More file actions
3507 lines (3165 loc) · 165 KB
/
Copy path_llmobs.py
File metadata and controls
3507 lines (3165 loc) · 165 KB
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
404
405
406
407
408
409
410
411
412
413
414
415
416
417
418
419
420
421
422
423
424
425
426
427
428
429
430
431
432
433
434
435
436
437
438
439
440
441
442
443
444
445
446
447
448
449
450
451
452
453
454
455
456
457
458
459
460
461
462
463
464
465
466
467
468
469
470
471
472
473
474
475
476
477
478
479
480
481
482
483
484
485
486
487
488
489
490
491
492
493
494
495
496
497
498
499
500
501
502
503
504
505
506
507
508
509
510
511
512
513
514
515
516
517
518
519
520
521
522
523
524
525
526
527
528
529
530
531
532
533
534
535
536
537
538
539
540
541
542
543
544
545
546
547
548
549
550
551
552
553
554
555
556
557
558
559
560
561
562
563
564
565
566
567
568
569
570
571
572
573
574
575
576
577
578
579
580
581
582
583
584
585
586
587
588
589
590
591
592
593
594
595
596
597
598
599
600
601
602
603
604
605
606
607
608
609
610
611
612
613
614
615
616
617
618
619
620
621
622
623
624
625
626
627
628
629
630
631
632
633
634
635
636
637
638
639
640
641
642
643
644
645
646
647
648
649
650
651
652
653
654
655
656
657
658
659
660
661
662
663
664
665
666
667
668
669
670
671
672
673
674
675
676
677
678
679
680
681
682
683
684
685
686
687
688
689
690
691
692
693
694
695
696
697
698
699
700
701
702
703
704
705
706
707
708
709
710
711
712
713
714
715
716
717
718
719
720
721
722
723
724
725
726
727
728
729
730
731
732
733
734
735
736
737
738
739
740
741
742
743
744
745
746
747
748
749
750
751
752
753
754
755
756
757
758
759
760
761
762
763
764
765
766
767
768
769
770
771
772
773
774
775
776
777
778
779
780
781
782
783
784
785
786
787
788
789
790
791
792
793
794
795
796
797
798
799
800
801
802
803
804
805
806
807
808
809
810
811
812
813
814
815
816
817
818
819
820
821
822
823
824
825
826
827
828
829
830
831
832
833
834
835
836
837
838
839
840
841
842
843
844
845
846
847
848
849
850
851
852
853
854
855
856
857
858
859
860
861
862
863
864
865
866
867
868
869
870
871
872
873
874
875
876
877
878
879
880
881
882
883
884
885
886
887
888
889
890
891
892
893
894
895
896
897
898
899
900
901
902
903
904
905
906
907
908
909
910
911
912
913
914
915
916
917
918
919
920
921
922
923
924
925
926
927
928
929
930
931
932
933
934
935
936
937
938
939
940
941
942
943
944
945
946
947
948
949
950
951
952
953
954
955
956
957
958
959
960
961
962
963
964
965
966
967
968
969
970
971
972
973
974
975
976
977
978
979
980
981
982
983
984
985
986
987
988
989
990
991
992
993
994
995
996
997
998
999
1000
import csv
from dataclasses import dataclass
from dataclasses import field
import inspect
import json
import math
import sys
import time
from typing import Any
from typing import Callable
from typing import Literal
from typing import Optional
from typing import Sequence
from typing import Union
from typing import cast
import urllib.parse
import ddtrace
from ddtrace import config
from ddtrace import patch
from ddtrace._trace.context import Context
from ddtrace._trace.processor import _NoopTraceProcessor
from ddtrace._trace.sampler import RateSampler
from ddtrace._trace.span import Span
from ddtrace._trace.tracer import Tracer
from ddtrace.constants import ERROR_MSG
from ddtrace.constants import ERROR_STACK
from ddtrace.constants import ERROR_TYPE
from ddtrace.ext import SpanTypes
from ddtrace.ext import git
from ddtrace.internal import atexit
from ddtrace.internal import core
from ddtrace.internal import forksafe
from ddtrace.internal.compat import ensure_text
from ddtrace.internal.logger import get_logger
from ddtrace.internal.native import generate_128bit_trace_id
from ddtrace.internal.native import rand64bits
from ddtrace.internal.remoteconfig.worker import remoteconfig_poller
from ddtrace.internal.sampling import format_rate
from ddtrace.internal.service import Service
from ddtrace.internal.service import ServiceStatusError
from ddtrace.internal.settings import env as _env
from ddtrace.internal.settings.integration import _integration_env_var_id
from ddtrace.internal.telemetry import get_config as _get_config
from ddtrace.internal.telemetry import telemetry_writer
from ddtrace.internal.telemetry.constants import TELEMETRY_APM_PRODUCT
from ddtrace.internal.threads import RLock
from ddtrace.internal.utils.deprecations import DDTraceDeprecationWarning
from ddtrace.internal.utils.formats import asbool
from ddtrace.internal.utils.formats import format_trace_id
from ddtrace.internal.utils.formats import parse_tags_str
from ddtrace.llmobs import _telemetry as telemetry
from ddtrace.llmobs._constants import ANNOTATIONS_CONTEXT_ID
from ddtrace.llmobs._constants import CACHED_LLMOBS_EVENT_CTX_KEY
from ddtrace.llmobs._constants import CACHED_LLMOBS_EXPORT_MODE_CTX_KEY
from ddtrace.llmobs._constants import CLAUDE_AGENT_SDK_APM_SPAN_NAME
from ddtrace.llmobs._constants import CREWAI_APM_SPAN_NAME
from ddtrace.llmobs._constants import DEFAULT_PROJECT_NAME
from ddtrace.llmobs._constants import DEFAULT_PROMPTS_CACHE_TTL
from ddtrace.llmobs._constants import DEFAULT_PROMPTS_TIMEOUT
from ddtrace.llmobs._constants import DISPATCH_ON_GUARDRAIL_SPAN_START
from ddtrace.llmobs._constants import DISPATCH_ON_LLM_SPAN_FINISH
from ddtrace.llmobs._constants import DISPATCH_ON_LLM_TOOL_CHOICE
from ddtrace.llmobs._constants import DISPATCH_ON_OPENAI_AGENT_SPAN_FINISH
from ddtrace.llmobs._constants import DISPATCH_ON_TOOL_CALL
from ddtrace.llmobs._constants import DISPATCH_ON_TOOL_CALL_OUTPUT_USED
from ddtrace.llmobs._constants import EXPERIMENT_CSV_FIELD_MAX_SIZE
from ddtrace.llmobs._constants import EXPERIMENT_DATASET_ID_KEY
from ddtrace.llmobs._constants import EXPERIMENT_DATASET_NAME_KEY
from ddtrace.llmobs._constants import EXPERIMENT_ID_KEY
from ddtrace.llmobs._constants import EXPERIMENT_NAME_KEY
from ddtrace.llmobs._constants import EXPERIMENT_PROJECT_ID_KEY
from ddtrace.llmobs._constants import EXPERIMENT_PROJECT_NAME_KEY
from ddtrace.llmobs._constants import EXPERIMENT_RUN_ID_KEY
from ddtrace.llmobs._constants import EXPERIMENT_RUN_ITERATION_KEY
from ddtrace.llmobs._constants import GEMINI_APM_SPAN_NAME
from ddtrace.llmobs._constants import INSTRUMENTATION_METHOD_ANNOTATED
from ddtrace.llmobs._constants import LANGCHAIN_APM_SPAN_NAME
from ddtrace.llmobs._constants import LITELLM_APM_SPAN_NAME
from ddtrace.llmobs._constants import LLMOBS_STRUCT
from ddtrace.llmobs._constants import ML_APP
from ddtrace.llmobs._constants import PROMPT_TRACKING_INSTRUMENTATION_METHOD
from ddtrace.llmobs._constants import PROPAGATED_LLMOBS_TRACE_ID_KEY
from ddtrace.llmobs._constants import PROPAGATED_ML_APP_KEY
from ddtrace.llmobs._constants import PROPAGATED_PARENT_ID_KEY
from ddtrace.llmobs._constants import PROPAGATED_SAMPLE_RATE
from ddtrace.llmobs._constants import PROPAGATED_SAMPLING_DECISION
from ddtrace.llmobs._constants import PROPAGATED_SESSION_ID_KEY
from ddtrace.llmobs._constants import ROOT_PARENT_ID
from ddtrace.llmobs._constants import SESSION_ID
from ddtrace.llmobs._constants import SPAN_START_WHILE_DISABLED_WARNING
from ddtrace.llmobs._constants import SUPPORTED_LLMOBS_INTEGRATIONS
from ddtrace.llmobs._constants import UNKNOWN_MODEL_NAME
from ddtrace.llmobs._constants import UNKNOWN_MODEL_PROVIDER
from ddtrace.llmobs._constants import VERTEXAI_APM_SPAN_NAME
from ddtrace.llmobs._constants import LLMObsExportMode
from ddtrace.llmobs._constants import LLMObsSamplingDecision
from ddtrace.llmobs._context import LLMObsContextProvider
from ddtrace.llmobs._evaluators.runner import EvaluatorRunner
from ddtrace.llmobs._experiment import AsyncEvaluatorType
from ddtrace.llmobs._experiment import AsyncSummaryEvaluatorType
from ddtrace.llmobs._experiment import AsyncTaskType
from ddtrace.llmobs._experiment import BaseAsyncEvaluator
from ddtrace.llmobs._experiment import BaseAsyncSummaryEvaluator
from ddtrace.llmobs._experiment import BaseEvaluator
from ddtrace.llmobs._experiment import BaseSummaryEvaluator
from ddtrace.llmobs._experiment import ConfigType
from ddtrace.llmobs._experiment import Dataset
from ddtrace.llmobs._experiment import DatasetRecord
from ddtrace.llmobs._experiment import DatasetRecordNew
from ddtrace.llmobs._experiment import EvaluatorType
from ddtrace.llmobs._experiment import Experiment
from ddtrace.llmobs._experiment import ExperimentResult
from ddtrace.llmobs._experiment import JSONType
from ddtrace.llmobs._experiment import Project
from ddtrace.llmobs._experiment import SummaryEvaluatorType
from ddtrace.llmobs._experiment import SyncExperiment
from ddtrace.llmobs._experiment import TaskType
from ddtrace.llmobs._experiment import _deep_eval_async_evaluator_wrapper
from ddtrace.llmobs._experiment import _deep_eval_evaluator_wrapper
from ddtrace.llmobs._experiment import _get_base_url
from ddtrace.llmobs._experiment import _is_deep_eval_evaluator
from ddtrace.llmobs._experiment import _is_pydantic_evaluator
from ddtrace.llmobs._experiment import _is_pydantic_report_evaluator
from ddtrace.llmobs._experiment import _parse_experiment_result
from ddtrace.llmobs._experiment import _pydantic_async_evaluator_wrapper
from ddtrace.llmobs._experiment import _pydantic_async_report_evaluator_wrapper
from ddtrace.llmobs._experiment import _pydantic_evaluator_wrapper
from ddtrace.llmobs._experiment import _pydantic_report_evaluator_wrapper
from ddtrace.llmobs._processor import LLMObsProcessor
from ddtrace.llmobs._prompt_optimization import PromptOptimization
from ddtrace.llmobs._prompt_optimization import validate_dataset
from ddtrace.llmobs._prompt_optimization import validate_dataset_split
from ddtrace.llmobs._prompt_optimization import validate_evaluators
from ddtrace.llmobs._prompt_optimization import validate_optimization_task
from ddtrace.llmobs._prompt_optimization import validate_task
from ddtrace.llmobs._prompt_optimization import validate_test_dataset
from ddtrace.llmobs._prompts import ManagedPrompt
from ddtrace.llmobs._prompts.cache import WarmCache
from ddtrace.llmobs._prompts.manager import PromptManager
from ddtrace.llmobs._utils import AnnotationContext
from ddtrace.llmobs._utils import LinkTracker
from ddtrace.llmobs._utils import _annotate_llmobs_span_data
from ddtrace.llmobs._utils import _batched
from ddtrace.llmobs._utils import _get_llmobs_data_metastruct
from ddtrace.llmobs._utils import _get_nearest_llmobs_ancestor
from ddtrace.llmobs._utils import _get_parent_prompt
from ddtrace.llmobs._utils import _normalize_wire_trace_id_to_hex
from ddtrace.llmobs._utils import _sanitize_span_event_depth
from ddtrace.llmobs._utils import _trace_id_to_wire
from ddtrace.llmobs._utils import _validate_prompt
from ddtrace.llmobs._utils import add_span_link
from ddtrace.llmobs._utils import enforce_message_role
from ddtrace.llmobs._utils import get_asyncio
from ddtrace.llmobs._utils import get_llmobs_ml_app
from ddtrace.llmobs._utils import get_llmobs_sample_rate
from ddtrace.llmobs._utils import get_llmobs_sampling_decision
from ddtrace.llmobs._utils import get_llmobs_session_id
from ddtrace.llmobs._utils import get_llmobs_span_kind
from ddtrace.llmobs._utils import get_llmobs_span_links
from ddtrace.llmobs._utils import get_llmobs_span_name
from ddtrace.llmobs._utils import get_llmobs_tags
from ddtrace.llmobs._utils import get_llmobs_trace_id
from ddtrace.llmobs._utils import get_tool_version_from_llm_span
from ddtrace.llmobs._utils import resolve_llmobs_git_metadata
from ddtrace.llmobs._utils import resolve_ml_app
from ddtrace.llmobs._utils import safe_json
from ddtrace.llmobs._writer import LLMObsAPIClient
from ddtrace.llmobs._writer import LLMObsEvalMetricWriter
from ddtrace.llmobs._writer import LLMObsEvaluationMetricEvent
from ddtrace.llmobs._writer import LLMObsExperimentsClient
from ddtrace.llmobs._writer import LLMObsSpanEvent
from ddtrace.llmobs._writer import LLMObsSpanWriter
from ddtrace.llmobs._writer import should_use_agentless
from ddtrace.llmobs.types import ChatMessage
from ddtrace.llmobs.types import DeletedPromptResponse
from ddtrace.llmobs.types import ExportedLLMObsSpan
from ddtrace.llmobs.types import Message
from ddtrace.llmobs.types import Prompt
from ddtrace.llmobs.types import PromptAuthError
from ddtrace.llmobs.types import PromptFallback
from ddtrace.llmobs.types import PromptResponse
from ddtrace.llmobs.types import PromptVersionResponse
from ddtrace.llmobs.types import _ErrorField
from ddtrace.llmobs.types import _Meta
from ddtrace.llmobs.types import _MetaIO
from ddtrace.llmobs.types import _SpanField
from ddtrace.llmobs.types import _ToolField
from ddtrace.llmobs.utils import Documents
from ddtrace.llmobs.utils import Messages
from ddtrace.llmobs.utils import extract_tool_definitions
from ddtrace.propagation.http import HTTPPropagator
from ddtrace.vendor.debtcollector import deprecate
from ddtrace.version import __version__
log = get_logger(__name__)
_STANDARD_INTEGRATION_SPAN_NAMES = (
CLAUDE_AGENT_SDK_APM_SPAN_NAME,
CREWAI_APM_SPAN_NAME,
GEMINI_APM_SPAN_NAME,
LANGCHAIN_APM_SPAN_NAME,
LITELLM_APM_SPAN_NAME,
VERTEXAI_APM_SPAN_NAME,
)
# requests/concurrent frameworks for distributed injection/extraction
_INTEGRATIONS_W_PROPAGATION_SUPPORT: dict[str, str] = {
"requests": "requests",
"httpx": "httpx",
"urllib3": "urllib3",
"grpc": "grpc",
"flask": "flask",
"starlette": "starlette",
"fastapi": "fastapi",
"aiohttp": "aiohttp",
"asyncio": "asyncio",
"futures": "futures",
}
# Constants for validation
_TASK_REQUIRED_PARAMS = {"input_data", "config"}
_EVALUATOR_REQUIRED_PARAMS = ("input_data", "output_data", "expected_output")
_SUMMARY_EVALUATOR_REQUIRED_PARAMS = (
"inputs",
"outputs",
"expected_outputs",
"evaluators_results",
)
_ml_app_deprecation_warned = False
def _resolve_agent_service(agent_service: Optional[str], ml_app: Optional[str]) -> Optional[str]:
global _ml_app_deprecation_warned
if ml_app and not _ml_app_deprecation_warned:
_ml_app_deprecation_warned = True
deprecate(
"The `ml_app` argument is deprecated",
message="Use `agent_service` instead. `ml_app` will be removed in a future major version.",
removal_version="5.0.0",
category=DDTraceDeprecationWarning,
)
return agent_service or ml_app
def _validate_task_signature(task: Callable, is_async: bool) -> None:
if not callable(task):
raise TypeError("task must be a callable function.")
if is_async and not get_asyncio().iscoroutinefunction(task):
raise TypeError("task must be an async function (coroutine function).")
sig = inspect.signature(task)
params = sig.parameters
if not all(param in params for param in _TASK_REQUIRED_PARAMS):
raise TypeError("Task function must have 'input_data' and 'config' parameters.")
def _validate_evaluator_signature(evaluator: Any, is_async: bool) -> None:
valid_base_classes: tuple[type, ...] = (BaseEvaluator,)
if is_async:
# async experiment allows both sync and async evaluators
valid_base_classes = (BaseEvaluator, BaseAsyncEvaluator)
if isinstance(evaluator, valid_base_classes):
return
if _is_deep_eval_evaluator(evaluator):
return
if _is_pydantic_evaluator(evaluator):
return
if not callable(evaluator):
if is_async:
raise TypeError(
f"Evaluator {evaluator} must be callable or an instance of BaseEvaluator/BaseAsyncEvaluator."
)
else:
raise TypeError(f"Evaluator {evaluator} must be callable or an instance of BaseEvaluator.")
sig = inspect.signature(evaluator)
params = sig.parameters
if not all(param in params for param in _EVALUATOR_REQUIRED_PARAMS):
raise TypeError("Evaluator function must have parameters {}.".format(tuple(_EVALUATOR_REQUIRED_PARAMS)))
def _validate_summary_evaluator_signature(evaluator: Any, is_async: bool) -> None:
valid_base_classes: tuple[type, ...] = (BaseSummaryEvaluator,)
if is_async:
# async experiment allows both sync and async summary evaluators
valid_base_classes = (BaseSummaryEvaluator, BaseAsyncSummaryEvaluator)
if isinstance(evaluator, valid_base_classes):
return
if not callable(evaluator):
if is_async:
raise TypeError(
f"Summary evaluator {evaluator} must be callable "
"or an instance of BaseSummaryEvaluator/BaseAsyncSummaryEvaluator."
)
else:
raise TypeError(f"Summary evaluator {evaluator} must be callable or an instance of BaseSummaryEvaluator.")
sig = inspect.signature(evaluator)
params = sig.parameters
if not all(param in params for param in _SUMMARY_EVALUATOR_REQUIRED_PARAMS):
raise TypeError(
"Summary evaluator function must have parameters {}.".format(tuple(_SUMMARY_EVALUATOR_REQUIRED_PARAMS))
)
class LLMObsExportSpanError(Exception):
"""Error raised when exporting a span."""
pass
class LLMObsAnnotateSpanError(Exception):
"""Error raised when annotating a span."""
pass
class LLMObsSubmitEvaluationError(Exception):
"""Error raised when submitting an evaluation."""
pass
class LLMObsInjectDistributedHeadersError(Exception):
"""Error raised when injecting distributed headers."""
pass
class LLMObsActivateDistributedHeadersError(Exception):
"""Error raised when activating distributed headers."""
pass
def _deprecate_prompt_label(method: str) -> None:
deprecate(
prefix="The 'label' parameter of LLMObs.{}() is deprecated".format(method),
message="Set DD_ENV instead; the prompt version is resolved for that environment.",
category=DDTraceDeprecationWarning,
)
@dataclass
class LLMObsSpan:
"""LLMObs span object.
Passed to the `span_processor` function in the `enable` or `register_processor` methods.
Example::
def span_processor(span: LLMObsSpan) -> Optional[LLMObsSpan]:
# Modify input/output
if span.get_tag("omit_span") == "1":
return None
if span.get_tag("no_input") == "1":
span.input = []
# Redact or drop sensitive top-level metadata
span.metadata.pop("tool_config", None)
return span
``metadata`` exposes the span's top-level metadata for redaction. Internal Datadog
fields (such as cost and agent manifest data) are never included and cannot be
modified through this object.
"""
input: list[Message] = field(default_factory=list)
output: list[Message] = field(default_factory=list)
_tags: dict[str, str] = field(default_factory=dict)
metadata: dict[str, Any] = field(default_factory=dict)
def get_tag(self, key: str) -> Optional[str]:
"""Get a tag from the span.
:param str key: The key of the tag to get.
:return: The value of the tag or None if the tag does not exist.
:rtype: Optional[str]
"""
return self._tags.get(key)
def _build_llmobs_span(
span_kind: str,
llmobs_input: _MetaIO,
llmobs_output: _MetaIO,
) -> tuple[LLMObsSpan, Literal["value", "messages", "documents", ""], Literal["value", "messages", "documents", ""]]:
"""Build an LLMObsSpan populated for the user span processor.
Routes input/output to messages or value depending on span kind.
Returns (llmobs_span, input_type, output_type).
"""
# experiment spans can store arbitrary (non-dict) values in output
if not isinstance(llmobs_input, dict):
llmobs_input = _MetaIO()
if not isinstance(llmobs_output, dict):
llmobs_output = _MetaIO()
llmobs_span = LLMObsSpan()
input_type: Literal["value", "messages", "documents", ""] = ""
output_type: Literal["value", "messages", "documents", ""] = ""
input_value = llmobs_input.get(LLMOBS_STRUCT.VALUE)
if input_value is not None:
input_type = "value"
llmobs_span.input = [Message(content=safe_json(input_value, ensure_ascii=False) or "", role="")]
input_messages = llmobs_input.get(LLMOBS_STRUCT.MESSAGES)
if span_kind == "llm" and input_messages is not None:
input_type = "messages"
llmobs_span.input = enforce_message_role(input_messages)
input_documents = llmobs_input.get(LLMOBS_STRUCT.DOCUMENTS)
if input_documents is not None:
input_type = "documents"
llmobs_span.input = [Message(content=doc.get("text", ""), role="") for doc in input_documents]
output_value = llmobs_output.get(LLMOBS_STRUCT.VALUE)
if output_value is not None:
output_type = "value"
llmobs_span.output = [Message(content=safe_json(output_value, ensure_ascii=False) or "", role="")]
output_messages = llmobs_output.get(LLMOBS_STRUCT.MESSAGES)
if span_kind == "llm" and output_messages is not None:
output_type = "messages"
llmobs_span.output = enforce_message_role(output_messages)
output_documents = llmobs_output.get(LLMOBS_STRUCT.DOCUMENTS)
if output_documents is not None:
output_type = "documents"
llmobs_span.output = [Message(content=doc.get("text", ""), role="") for doc in output_documents]
return llmobs_span, input_type, output_type
def _reconstruct_documents(meta_io: _MetaIO, messages: list[Message]) -> None:
"""Merge processor-modified message content back into documents, preserving metadata."""
original_docs = meta_io.get(LLMOBS_STRUCT.DOCUMENTS) or []
if messages:
meta_io[LLMOBS_STRUCT.DOCUMENTS] = [
{**original_docs[i], "text": msg.get("content", "")}
for i, msg in enumerate(messages)
if i < len(original_docs)
]
else:
meta_io.pop(LLMOBS_STRUCT.DOCUMENTS, None)
def _normalize_llmobs_meta(
span: Span,
llmobs_span: LLMObsSpan,
llmobs_meta: _Meta,
span_kind: str,
input_type: Literal["value", "messages", "documents", ""],
output_type: Literal["value", "messages", "documents", ""],
export_to_llmobs: bool,
) -> None:
"""Normalize the llmobs meta dict in place so `_llmobs_span_event()` can read it directly.
Writes post-user-processor I/O back, inherits parent prompts for LLM spans, drops
invalid prompts, populates the error field, normalizes model_provider, and removes
empty optional fields.
"""
llmobs_meta[LLMOBS_STRUCT.SPAN] = _SpanField(kind=span_kind)
current_metadata = llmobs_meta.get(LLMOBS_STRUCT.METADATA)
if not isinstance(current_metadata, dict):
current_metadata = {}
preserved_dd = current_metadata.get(LLMOBS_STRUCT.METADATA_DD)
user_metadata = llmobs_span.metadata
if not isinstance(user_metadata, dict):
log.warning("LLMObs span processor set non-dict metadata (%r); ignoring.", type(user_metadata))
user_metadata = current_metadata
# Drop any user-supplied `_dd` so internal metadata cannot be spoofed, then restore the real one.
new_metadata = {k: v for k, v in user_metadata.items() if k != LLMOBS_STRUCT.METADATA_DD}
if preserved_dd is not None:
new_metadata[LLMOBS_STRUCT.METADATA_DD] = preserved_dd
llmobs_meta[LLMOBS_STRUCT.METADATA] = new_metadata
model_name = llmobs_meta.pop(LLMOBS_STRUCT.MODEL_NAME, None)
model_provider = llmobs_meta.pop(LLMOBS_STRUCT.MODEL_PROVIDER, None)
if span_kind in ("llm", "embedding"):
llmobs_meta[LLMOBS_STRUCT.MODEL_NAME] = model_name or UNKNOWN_MODEL_NAME
llmobs_meta[LLMOBS_STRUCT.MODEL_PROVIDER] = (model_provider or UNKNOWN_MODEL_PROVIDER).lower()
if span_kind != "llm":
llmobs_meta.pop(LLMOBS_STRUCT.TOOL_DEFINITIONS, None)
if span_kind != "tool":
llmobs_meta.pop(LLMOBS_STRUCT.TOOL, None)
elif LLMOBS_STRUCT.TOOL not in llmobs_meta:
tool_name = get_llmobs_span_name(span)
if tool_name:
ancestor = _get_nearest_llmobs_ancestor(span)
if ancestor is not None and get_llmobs_span_kind(ancestor) == "llm":
version = get_tool_version_from_llm_span(ancestor, tool_name)
if version is not None:
llmobs_meta[LLMOBS_STRUCT.TOOL] = _ToolField(version=version)
intent = llmobs_meta.pop(LLMOBS_STRUCT.INTENT, None)
if intent:
llmobs_meta[LLMOBS_STRUCT.INTENT] = str(intent)
if span.error and export_to_llmobs:
llmobs_meta[LLMOBS_STRUCT.ERROR] = _ErrorField(
message=span.get_tag(ERROR_MSG) or "",
stack=span.get_tag(ERROR_STACK) or "",
type=span.get_tag(ERROR_TYPE) or "",
)
if span.context.get_baggage_item(EXPERIMENT_ID_KEY) and span_kind == "experiment":
# experiment i/o is stored as raw values — already in place from annotation
return
meta_input: _MetaIO = cast(_MetaIO, dict(llmobs_meta.get(LLMOBS_STRUCT.INPUT) or {}))
meta_output: _MetaIO = cast(_MetaIO, dict(llmobs_meta.get(LLMOBS_STRUCT.OUTPUT) or {}))
input_prompt = meta_input.get(LLMOBS_STRUCT.PROMPT)
if input_prompt is not None and span_kind != "llm":
log.warning("Dropping prompt on non-LLM span kind, annotating prompts is only supported for LLM span kinds.")
meta_input.pop(LLMOBS_STRUCT.PROMPT, None)
elif input_prompt is None and span_kind == "llm":
parent_prompt = _get_parent_prompt(span)
if parent_prompt is not None:
meta_input[LLMOBS_STRUCT.PROMPT] = parent_prompt
if input_type == "messages":
meta_input[LLMOBS_STRUCT.MESSAGES] = llmobs_span.input
elif input_type == "value" and llmobs_span.input:
meta_input[LLMOBS_STRUCT.VALUE] = llmobs_span.input[0].get("content", "")
elif input_type == "documents":
_reconstruct_documents(meta_input, llmobs_span.input)
if meta_input:
llmobs_meta[LLMOBS_STRUCT.INPUT] = meta_input
else:
llmobs_meta.pop(LLMOBS_STRUCT.INPUT, None)
if output_type == "messages":
meta_output[LLMOBS_STRUCT.MESSAGES] = llmobs_span.output
elif output_type == "value" and llmobs_span.output:
meta_output[LLMOBS_STRUCT.VALUE] = llmobs_span.output[0].get("content", "")
elif output_type == "documents":
_reconstruct_documents(meta_output, llmobs_span.output)
if meta_output:
llmobs_meta[LLMOBS_STRUCT.OUTPUT] = meta_output
else:
llmobs_meta.pop(LLMOBS_STRUCT.OUTPUT, None)
class LLMObs(Service):
_instance = None # type: LLMObs
enabled = False
_app_key: str = _env.get("DD_APP_KEY", "")
_project_name: str = _env.get("DD_LLMOBS_PROJECT_NAME", DEFAULT_PROJECT_NAME)
_git_repository_url: str = ""
_git_commit_sha: str = ""
def __init__(
self,
tracer: Optional[Tracer] = None,
span_processor: Optional[Callable[[LLMObsSpan], Optional[LLMObsSpan]]] = None,
) -> None:
super(LLMObs, self).__init__()
self.tracer = tracer or ddtrace.tracer
self._llmobs_context_provider = LLMObsContextProvider()
self._user_span_processor = span_processor
agentless_enabled = should_use_agentless(user_defined_agentless_enabled=config._llmobs_agentless_enabled)
if not asbool(_env.get("DD_APM_TRACING_ENABLED", "true")):
# APMTracingEnabledFilter drops every trace.
self._export_mode = (
LLMObsExportMode.LLMOBS_AGENTLESS if agentless_enabled else LLMObsExportMode.LLMOBS_AGENT_PROXY
)
elif agentless_enabled:
self._export_mode = LLMObsExportMode.APM_AGENTLESS
else:
self._export_mode = LLMObsExportMode.APM_AGENT
self._llmobs_span_writer = LLMObsSpanWriter(
interval=float(_env.get("_DD_LLMOBS_WRITER_INTERVAL", 1.0)),
timeout=float(_env.get("_DD_LLMOBS_WRITER_TIMEOUT", 5.0)),
is_agentless=agentless_enabled,
)
self._llmobs_eval_metric_writer = LLMObsEvalMetricWriter(
interval=float(_env.get("_DD_LLMOBS_WRITER_INTERVAL", 1.0)),
timeout=float(_env.get("_DD_LLMOBS_WRITER_TIMEOUT", 5.0)),
is_agentless=agentless_enabled,
)
self._evaluator_runner = EvaluatorRunner(
interval=float(_env.get("_DD_LLMOBS_EVALUATOR_INTERVAL", 1.0)),
llmobs_service=self,
)
self._dne_client = LLMObsExperimentsClient(
interval=float(_env.get("_DD_LLMOBS_WRITER_INTERVAL", 1.0)),
timeout=float(_env.get("_DD_LLMOBS_WRITER_TIMEOUT", 5.0)),
_app_key=self._app_key,
_default_project=Project(name=self._project_name, _id=""),
is_agentless=True, # agent proxy doesn't seem to work for experiments
)
self._api_client = LLMObsAPIClient(app_key=self._app_key)
forksafe.register(self._child_after_fork)
self._link_tracker = LinkTracker()
self._annotations: list[tuple[str, str, dict[str, Any]]] = []
self._annotation_context_lock = RLock()
# True if enable() switched the APM writer to agentless; disable() reverts it.
self._apm_writer_switched_to_agentless = False
self._sampler = RateSampler(sample_rate=config._llmobs_sample_rate)
def _on_span_start(self, span: Span) -> None:
if self.enabled and span.span_type == SpanTypes.LLM:
self._activate_llmobs_span(span)
telemetry.record_span_started()
self._do_annotations(span)
def _on_span_finish(self, span: Span) -> None:
if not self.enabled or span.span_type != SpanTypes.LLM:
return
span_kind = get_llmobs_span_kind(span)
if span_kind == "llm":
core.dispatch(DISPATCH_ON_LLM_SPAN_FINISH, (span,))
span_event = None
try:
if self._prepare_llmobs_span_data(span, span_kind):
span_event = self._llmobs_span_event(span)
except (KeyError, TypeError, ValueError):
log.error(
"Error generating LLMObs span event for span %s, likely due to malformed span",
span,
exc_info=True,
)
if not span_event:
# clear meta_struct if no event to export (dropped by user processor / error during preparation/assembly)
span._remove_struct_tag(LLMOBS_STRUCT.KEY)
return
if self._evaluator_runner and span_kind == "llm":
self._evaluator_runner.enqueue(span_event, span)
span._set_ctx_item(CACHED_LLMOBS_EXPORT_MODE_CTX_KEY, self._export_mode)
span._set_ctx_item(CACHED_LLMOBS_EVENT_CTX_KEY, span_event)
def _apply_user_span_processor(self, span: Span, llmobs_span: LLMObsSpan) -> Optional[LLMObsSpan]:
"""Run the user span processor.
Returns the possibly mutated span, or None if the span should be dropped.
On error, logs and returns the original span unchanged.
"""
if self._user_span_processor is None:
return llmobs_span
error = False
try:
llmobs_span._tags = get_llmobs_tags(span) or {}
llmobs_span._tags["span.kind"] = get_llmobs_span_kind(span) or ""
result = self._user_span_processor(llmobs_span)
if result is None:
return None
if not isinstance(result, LLMObsSpan):
raise TypeError("User span processor must return an LLMObsSpan or None, got %r" % type(result))
return result
except Exception as e:
log.error("Error in LLMObs span processor (%r): %r", self._user_span_processor, e)
error = True
return llmobs_span
finally:
telemetry.record_llmobs_user_processor_called(error)
def _prepare_llmobs_span_data(self, span: Span, span_kind: Optional[str]) -> bool:
"""Commit final I/O and meta values to meta_struct before event assembly.
Runs the user span processor and folds its mutations — plus parent-prompt
inheritance, error fields, and other meta-level rules — directly into
the LLMObsSpanData stored on the span meta_struct.
Returns True if the span is ready to be serialized into an event; False if
the span has no LLMObs data, the user processor dropped it, or finalization
failed.
"""
llmobs_data = _get_llmobs_data_metastruct(span)
if not llmobs_data:
log.error(
"Error preparing LLMObs span event for span %s, missing LLMObs data in span context.",
span,
)
return False
if not span_kind:
log.error(
"Error preparing LLMObs span event for span %s, missing span kind in span context.",
span,
)
return False
llmobs_meta = llmobs_data.setdefault(LLMOBS_STRUCT.META, _Meta())
llmobs_input = llmobs_meta.get(LLMOBS_STRUCT.INPUT) or _MetaIO()
llmobs_output = llmobs_meta.get(LLMOBS_STRUCT.OUTPUT) or _MetaIO()
llmobs_span, input_type, output_type = _build_llmobs_span(span_kind, llmobs_input, llmobs_output)
original_metadata = llmobs_meta.get(LLMOBS_STRUCT.METADATA)
if not isinstance(original_metadata, dict):
original_metadata = {}
llmobs_span.metadata = {k: v for k, v in original_metadata.items() if k != LLMOBS_STRUCT.METADATA_DD}
user_processed_span = self._apply_user_span_processor(span, llmobs_span)
if user_processed_span is None:
log.debug("LLMObs span %s dropped by user processor", span)
return False
_normalize_llmobs_meta(
span,
user_processed_span,
llmobs_meta,
span_kind,
input_type,
output_type,
export_to_llmobs=self._export_mode != LLMObsExportMode.APM_AGENTLESS,
)
llmobs_data[LLMOBS_STRUCT.META] = _sanitize_span_event_depth(llmobs_meta)
if self._export_mode == LLMObsExportMode.APM_AGENTLESS:
# APM agentless ingestion treats dots in tag keys as nested-path separators;
# replace them with underscores before encoding.
tags = {k.replace(".", "_"): v for k, v in llmobs_data.get(LLMOBS_STRUCT.TAGS, {}).items()}
llmobs_data[LLMOBS_STRUCT.TAGS] = tags
span._set_struct_tag(LLMOBS_STRUCT.KEY, cast(dict[str, Any], llmobs_data))
return True
def _llmobs_span_event(self, span: Span) -> Optional[LLMObsSpanEvent]:
"""Assemble the LLMObs span event from the finalized meta_struct contents.
This function is a pure reader: all I/O mutations, user-processor hooks, and
meta-level finalization must already have been committed to
``meta_struct`` by ``_prepare_llmobs_span_data``.
"""
llmobs_data = _get_llmobs_data_metastruct(span)
if not llmobs_data:
return None
parent_id = llmobs_data.get(LLMOBS_STRUCT.PARENT_ID) or ROOT_PARENT_ID
llmobs_trace_id = llmobs_data.get(LLMOBS_STRUCT.TRACE_ID)
if llmobs_trace_id is None:
raise ValueError("Failed to extract LLMObs trace ID from span context.")
meta = llmobs_data.get(LLMOBS_STRUCT.META) or _Meta()
metrics = llmobs_data.get(LLMOBS_STRUCT.METRICS) or {}
tags = self._llmobs_tags(span)
_dd_attrs = {
**(llmobs_data.get("_dd") or {}),
"span_id": str(span.span_id),
"trace_id": format_trace_id(span.trace_id),
"apm_trace_id": format_trace_id(span.trace_id),
}
llmobs_span_event: LLMObsSpanEvent = {
"trace_id": llmobs_trace_id,
"span_id": str(span.span_id),
"parent_id": parent_id,
"name": get_llmobs_span_name(span) or span.name,
"start_ns": span.start_ns,
"duration": cast(int, span.duration_ns),
"status": "error" if span.error else "ok",
"meta": meta,
"metrics": metrics,
"tags": tags,
"_dd": _dd_attrs,
}
experiment_config = llmobs_data.get(LLMOBS_STRUCT.CONFIG)
if experiment_config:
llmobs_span_event["config"] = experiment_config
session_id = get_llmobs_session_id(span)
if session_id:
llmobs_span_event["session_id"] = session_id
span_links = get_llmobs_span_links(span) or []
if span_links:
llmobs_span_event["span_links"] = span_links
return llmobs_span_event
def _llmobs_tags(self, span: Span) -> list[str]:
tags = dict(get_llmobs_tags(span) or {})
if self._export_mode != LLMObsExportMode.APM_AGENTLESS:
tags["error"] = str(span.error)
err_type = span.get_tag(ERROR_TYPE)
if err_type:
tags["error_type"] = err_type
return sorted("{}:{}".format(k, v) for k, v in tags.items())
def _do_annotations(self, span: Span) -> None:
# get the current span context
# only do the annotations if it matches the context
if span.span_type != SpanTypes.LLM: # do this check to avoid the warning log in `annotate`
return
current_context = self._instance.tracer.current_trace_context()
if current_context is None:
return
current_context_id = current_context.get_baggage_item(ANNOTATIONS_CONTEXT_ID)
with self._annotation_context_lock:
for _, context_id, annotation_kwargs in self._instance._annotations:
if current_context_id == context_id:
self.annotate(span, **annotation_kwargs, _suppress_span_kind_error=True)
def _child_after_fork(self) -> None:
self._llmobs_span_writer = self._llmobs_span_writer.recreate()
self._llmobs_eval_metric_writer = self._llmobs_eval_metric_writer.recreate()
self._evaluator_runner = self._evaluator_runner.recreate()
LLMObs._prompt_manager = None
# The tracer's fork handler keeps whatever writer was active; we didn't swap it here,
# so clear the flag to stop disable() reverting it in the child.
self._apm_writer_switched_to_agentless = False
if self.enabled:
# Rebind: the processor holds the pre-fork writer whose worker thread is dead
# after fork(), so leaving it would silently buffer rescued events in the child.
self.tracer._span_aggregator.llmobs_processor = LLMObsProcessor(self._llmobs_span_writer, self.tracer)
self._start_service()
def _start_service(self) -> None:
try:
self._llmobs_span_writer.start()
self._llmobs_eval_metric_writer.start()
except ServiceStatusError:
log.debug("Error starting LLMObs writers")
try:
self._evaluator_runner.start()
except ServiceStatusError:
log.debug("Error starting evaluator runner")
def _stop_service(self) -> None:
try:
self._evaluator_runner.stop()
# flush remaining evaluation spans & evaluations
self._instance._llmobs_span_writer.periodic()
self._instance._llmobs_eval_metric_writer.periodic()
except ServiceStatusError:
log.debug("Error stopping evaluator runner")
try:
self._llmobs_span_writer.stop()
self._llmobs_eval_metric_writer.stop()
except ServiceStatusError:
log.debug("Error stopping LLMObs writers")
# Remove listener hooks for span events
core.reset_listeners("trace.span_start", self._on_span_start)
core.reset_listeners("trace.span_finish", self._on_span_finish)
core.reset_listeners("http.span_inject", self._inject_llmobs_context)
core.reset_listeners(
"http.activate_distributed_headers",
self._activate_llmobs_distributed_context_soft_fail,
)
core.reset_listeners("threading.submit", self._current_trace_context)
core.reset_listeners("threading.execution", self._llmobs_context_provider.activate)
core.reset_listeners("asyncio.create_task", self._on_asyncio_create_task)
core.reset_listeners("asyncio.execute_task", self._on_asyncio_execute_task)
core.reset_listeners(DISPATCH_ON_LLM_TOOL_CHOICE, self._link_tracker.on_llm_tool_choice)
core.reset_listeners(DISPATCH_ON_TOOL_CALL, self._link_tracker.on_tool_call)
core.reset_listeners(
DISPATCH_ON_TOOL_CALL_OUTPUT_USED,
self._link_tracker.on_tool_call_output_used,
)
core.reset_listeners(DISPATCH_ON_GUARDRAIL_SPAN_START, self._link_tracker.on_guardrail_span_start)
core.reset_listeners(DISPATCH_ON_LLM_SPAN_FINISH, self._link_tracker.on_llm_span_finish)
core.reset_listeners(
DISPATCH_ON_OPENAI_AGENT_SPAN_FINISH,
self._link_tracker.on_openai_agent_span_finish,
)
forksafe.unregister(self._child_after_fork)
@classmethod
def enable(
cls,
ml_app: Optional[str] = None,
integrations_enabled: bool = True,
agentless_enabled: Optional[bool] = None,
instrumented_proxy_urls: Optional[set[str]] = None,
site: Optional[str] = None,
api_key: Optional[str] = None,
app_key: Optional[str] = None,
project_name: Optional[str] = None,
env: Optional[str] = None,
service: Optional[str] = None,
span_processor: Optional[Callable[[LLMObsSpan], Optional[LLMObsSpan]]] = None,
sample_rate: Optional[float] = None,
agent_service: Optional[str] = None,
_tracer: Optional[Tracer] = None,
_auto: bool = False,
) -> None:
"""
Enable LLM Observability tracing.
:param str ml_app: Deprecated. Use ``agent_service`` instead.
:param str agent_service: The name of your agent service. Takes precedence over ``ml_app`` and
``service``, in that order.
:param bool integrations_enabled: set to `true` to enable LLM integrations.
:param bool agentless_enabled: set to `true` to disable sending data that requires a Datadog Agent.
:param set[str] instrumented_proxy_urls: A set of instrumented proxy URLs to help detect when to emit LLM spans.
:param str site: Your datadog site.
:param str api_key: Your datadog api key.
:param str app_key: Your datadog application key.
:param str project_name: Your project name used for experiments.
:param str env: Your environment name.
:param str service: Your service name.
:param Callable[[LLMObsSpan], Optional[LLMObsSpan]] span_processor: A function that takes an LLMObsSpan and
returns an LLMObsSpan or None. If None is returned, the span will be omitted and not sent to LLMObs.
:param float sample_rate: The proportion of LLMObs traces to sample, between 0.0 and 1.0 (inclusive).
Takes precedence over the DD_LLMOBS_SAMPLE_RATE environment variable. Defaults to that env var, or 1.0
(sample everything) if neither is set or within the valid range.
"""
if cls.enabled:
log.debug("%s already enabled", cls.__name__)
return
cls._warn_if_litellm_was_imported()
if _env.get("DD_LLMOBS_ENABLED") and not asbool(_env.get("DD_LLMOBS_ENABLED")):
log.debug("LLMObs.enable() called when DD_LLMOBS_ENABLED is set to false or 0, not starting LLMObs service")
return
# grab required values for LLMObs
config._dd_site = site or config._dd_site
config._dd_api_key = api_key or config._dd_api_key
cls._app_key = app_key or cls._app_key
if app_key:
# Invalidate any prompt manager cached by a read path (e.g. get_prompt)
# before the app key was configured, so it rebuilds with the new key.
with cls._prompt_manager_lock:
cls._prompt_manager = None
cls._project_name = project_name or cls._project_name or DEFAULT_PROJECT_NAME
cls._git_repository_url, cls._git_commit_sha = resolve_llmobs_git_metadata()
config.env = env or config.env
config.service = service or config.service
config._llmobs_ml_app = _resolve_agent_service(agent_service, ml_app) or config._llmobs_ml_app
config._llmobs_instrumented_proxy_urls = instrumented_proxy_urls or config._llmobs_instrumented_proxy_urls
# Validate and fallback sample rates (inline arg --> env var --> 1.0)
if sample_rate is not None and 0.0 <= sample_rate <= 1.0:
config._llmobs_sample_rate = sample_rate
else:
if sample_rate is not None:
log.warning(
"Invalid LLMObs sample rate argument (%r outside valid range [0.0, 1.0]); ignoring it.",
sample_rate,
)
if not 0.0 <= config._llmobs_sample_rate <= 1.0:
log.warning(
"Invalid LLMObs sample rate (%r outside valid range [0.0, 1.0]). Falling back to 1.0.",
config._llmobs_sample_rate,
)
config._llmobs_sample_rate = 1.0
error = None
start_ns = time.time_ns()
try:
config._llmobs_agentless_enabled = should_use_agentless(
user_defined_agentless_enabled=(
agentless_enabled if agentless_enabled is not None else config._llmobs_agentless_enabled
)
)
if config._llmobs_agentless_enabled:
# validate required values for agentless LLMObs
if not config._dd_api_key:
error = "missing_api_key"
raise ValueError(
"DD_API_KEY is required for sending LLMObs data when agentless mode is enabled. "
"Ensure this configuration is set before running your application."
)
if not config._dd_site:
error = "missing_site"
raise ValueError(
"DD_SITE is required for sending LLMObs data when agentless mode is enabled. "
"Ensure this configuration is set before running your application."
)
if not _env.get("DD_REMOTE_CONFIGURATION_ENABLED"):
config._remote_config_enabled = False
log.debug("Remote configuration disabled because DD_LLMOBS_AGENTLESS_ENABLED is set to true.")
remoteconfig_poller.disable()
# Since the API key can be set programmatically and TelemetryWriter is already initialized by now,
# we need to force telemetry to use agentless configuration
telemetry_writer.enable_agentless_client(True)
if integrations_enabled:
cls._patch_integrations()
# override the default _instance with a new tracer
cls._instance = cls(tracer=_tracer, span_processor=span_processor)
cls.enabled = True
# Align config._llmobs_enabled with effective state for user-initiated calls.
# When _auto=True, the caller (RC handler, env-var auto-start) has already
# written the appropriate source; stamping "code" would mask it on precedence.