-
Notifications
You must be signed in to change notification settings - Fork 536
Expand file tree
/
Copy path_uploader.py
More file actions
286 lines (230 loc) · 10.3 KB
/
Copy path_uploader.py
File metadata and controls
286 lines (230 loc) · 10.3 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
from dataclasses import dataclass
from enum import Enum
from typing import Any
from typing import Optional
from ddtrace import config as ddconfig
from ddtrace.debugging._config import di_config
from ddtrace.debugging._encoding import LogSignalJsonEncoder
from ddtrace.debugging._encoding import SignalQueue
from ddtrace.debugging._encoding import SnapshotJsonEncoder
from ddtrace.debugging._metrics import metrics
from ddtrace.debugging._signal.collector import SignalCollector
from ddtrace.debugging._signal.model import SignalTrack
from ddtrace.internal import agent
from ddtrace.internal import logger
from ddtrace.internal.logger import get_logger
from ddtrace.internal.native import DebuggerSender
from ddtrace.internal.native import DebuggerTrackType
from ddtrace.internal.native_runtime import get_native_runtime
from ddtrace.internal.utils.formats import get_test_session_token
from ddtrace.internal.utils.retry import fibonacci_backoff_with_jitter
log = get_logger(__name__)
UNSUPPORTED_AGENT = "unsupported_agent"
logger.set_tag_rate_limit(UNSUPPORTED_AGENT, logger.HOUR)
meter = metrics.get_meter("uploader")
def build_debugger_sender() -> DebuggerSender:
"""Build a sender for the logs, snapshots and diagnostics tracks."""
timeout_ms = int(di_config.upload_timeout * 1000)
if ddconfig._agentless_enabled:
return DebuggerSender(
get_native_runtime(),
site=ddconfig._dd_site,
api_key=ddconfig._dd_api_key,
tags=di_config.tags,
timeout_ms=timeout_ms,
test_session_token=get_test_session_token(),
)
return DebuggerSender(
get_native_runtime(),
url=di_config._intake_url,
tags=di_config.tags,
timeout_ms=timeout_ms,
test_session_token=get_test_session_token(),
)
class UploaderProduct(str, Enum):
"""Uploader products."""
DEBUGGER = "dynamic_instrumentation"
EXCEPTION_REPLAY = "exception_replay"
CODE_ORIGIN_SPAN_ENTRY = "code_origin.span.entry"
@dataclass
class UploaderTrack:
track: SignalTrack
debugger_type: DebuggerTrackType
queue: SignalQueue
enabled: bool = True
class SignalUploaderError(Exception):
"""Signal uploader error."""
pass
class SignalUploader(agent.AgentCheckPeriodicService):
"""Signal uploader.
This class implements an interface with the debugger signal intake for both
the debugger and the events platform.
"""
_instance: Optional["SignalUploader"] = None
_products: set[UploaderProduct] = set()
_agent_endpoints: set[str] = set()
__queue__ = SignalQueue
__collector__ = SignalCollector
RETRY_ATTEMPTS = 3
def __init__(self, interval: Optional[float] = None) -> None:
super().__init__(interval if interval is not None else di_config.upload_interval_seconds)
self._sender = build_debugger_sender()
self._tracks = {
SignalTrack.LOGS: UploaderTrack(
track=SignalTrack.LOGS,
debugger_type=DebuggerTrackType.Logs,
queue=self.__queue__(
encoder=LogSignalJsonEncoder(di_config.service_name), on_full=self._on_buffer_full
),
),
SignalTrack.SNAPSHOT: UploaderTrack(
track=SignalTrack.SNAPSHOT,
debugger_type=DebuggerTrackType.Snapshots,
queue=self.__queue__(encoder=SnapshotJsonEncoder(di_config.service_name), on_full=self._on_buffer_full),
),
}
self._collector = self.__collector__({t: ut.queue for t, ut in self._tracks.items()})
if self._sender.agentless:
# There is no agent to negotiate endpoints with, so skip the agent
# check state and start uploading straight away.
self._state = self._online
# Make it retry-able
self._write_with_backoff = fibonacci_backoff_with_jitter(
initial_wait=0.618 * self.interval / (1.618**self.RETRY_ATTEMPTS) / 2,
attempts=self.RETRY_ATTEMPTS,
)(self._write)
log.debug("Signal uploader initialized (sender: %r, interval: %f)", self._sender, self.interval)
self._flush_full = False
def info_check(self, agent_info: Optional[dict[str, Any]]) -> bool:
if self._sender.agentless:
# Payloads go straight to the intake, on paths that are fixed at
# construction: there is nothing to negotiate with the agent.
return True
if agent_info is None:
# Agent is unreachable
return False
if "endpoints" not in agent_info:
# Agent not supported
log.debug("Unsupported Datadog agent detected. Please upgrade to 7.49.0.")
return False
# Agent /info entries may or may not carry a leading slash depending on
# the agent version, so normalize before matching (see remoteconfig).
endpoints = {endpoint.lstrip("/") for endpoint in agent_info.get("endpoints", [])}
logs_track = self._tracks[SignalTrack.LOGS]
snapshot_track = self._tracks[SignalTrack.SNAPSHOT]
logs_track.enabled = True
snapshot_track.enabled = True
if "debugger/v2/input" in endpoints:
log.debug("Detected /debugger/v2/input endpoint")
# Undo any downgrade from an earlier online cycle.
self._sender.reset_endpoints()
elif "debugger/v1/diagnostics" in endpoints:
log.debug("Detected /debugger/v1/diagnostics endpoint fallback")
self._sender.downgrade_to_diagnostics()
else:
logs_track.enabled = False
snapshot_track.enabled = False
self._throttle_agent_check()
log.warning(
UNSUPPORTED_AGENT,
extra={
"product": "debugger",
"more_info": (
"Unsupported Datadog agent detected. Logs and snapshots from Dynamic Instrumentation/"
"Exception Replay/Code Origin for Spans will not be uploaded. "
"Please upgrade to version 7.49.0 or later"
),
},
)
return True
def _write(self, payload: bytes, debugger_type: DebuggerTrackType) -> None:
try:
response = self._sender.send(payload, debugger_type)
except Exception:
# The request never completed (transport failure or timeout). Drop the
# batch: unlike a rejection, there is no endpoint to fall back to.
log.error("Failed to write payload to the %s track", debugger_type, exc_info=True)
meter.increment("error")
return
if not response.accepted:
log.error(
"Failed to upload payload to the %s track: [%d] %r",
debugger_type,
response.status,
response.body,
)
meter.increment("upload.error", tags={"status": str(response.status)})
msg = "Failed to upload payload"
raise SignalUploaderError(msg)
meter.increment("upload.success")
meter.distribution("upload.size", len(payload))
def _on_buffer_full(self, _item: Any, _encoded: bytes) -> None:
self._flush_full = True
self.upload()
def upload(self) -> None:
"""Upload request."""
self.awake()
def reset(self) -> None:
"""Reset the buffer on fork."""
super().reset()
for track in self._tracks.values():
track.queue = self.__queue__(encoder=track.queue._encoder, on_full=self._on_buffer_full)
self._collector._tracks = {t: ut.queue for t, ut in self._tracks.items()}
def _flush_track(self, track: UploaderTrack) -> None:
if (data := track.queue.flush()) is not None and track.enabled:
payload, count = data
try:
self._write_with_backoff(payload, track.debugger_type)
meter.distribution("batch.cardinality", count)
except SignalUploaderError:
if self._sender.downgrade_to_diagnostics():
log.debug("Downgraded debugger endpoints to the diagnostics endpoint")
# Retry once against the diagnostics endpoint
self._write_with_backoff(payload, track.debugger_type)
meter.distribution("batch.cardinality", count)
elif self._sender.agentless:
log.debug("Cannot upload payload to the intake", exc_info=True)
else:
raise # Propagate error to transition to agent check state
except Exception:
log.debug("Cannot upload payload", exc_info=True)
def _flush(self) -> None:
"""Upload the buffer content to the agent."""
if self._flush_full:
# We received the signal to flush a full buffer
self._flush_full = False
for uploader_track in self._tracks.values():
if uploader_track.queue.is_full():
self._flush_track(uploader_track)
for track in self._tracks.values():
if track.queue.count:
self._flush_track(track)
def online(self) -> None:
self._flush()
if not self._tracks[SignalTrack.SNAPSHOT].enabled or not self._tracks[SignalTrack.LOGS].enabled:
# If the tracks are not enabled, we raise an exception to
# transition back to the agent check state in case we detect an
# agent that can handle logs and snapshots safely.
msg = "Debugger tracks not enabled"
raise ValueError(msg)
def on_shutdown(self) -> None: # type: ignore[override]
self._flush()
@classmethod
def get_collector(cls) -> Optional[SignalCollector]:
return cls._instance._collector if cls._instance is not None else None
@classmethod
def register(cls, product: UploaderProduct) -> None:
if product in cls._products:
return
cls._products.add(product)
if cls._instance is None:
cls._instance = cls()
cls._instance.start()
@classmethod
def unregister(cls, product: UploaderProduct) -> None:
if product not in cls._products:
return
cls._products.remove(product)
if not cls._products and cls._instance is not None:
cls._instance.stop()
cls._instance = None