Skip to content

Commit b306862

Browse files
authored
Merge pull request DataDog#6276 from DataDog/jpbempel/diagnostic-endpoint
Send probe statuses to new debugger track
2 parents 4908b2c + aa39e08 commit b306862

15 files changed

Lines changed: 220 additions & 54 deletions

File tree

‎communication/src/main/java/datadog/communication/ddagent/DDAgentFeaturesDiscovery.java‎

Lines changed: 21 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -46,6 +46,7 @@ public class DDAgentFeaturesDiscovery implements DroppingPolicy {
4646
public static final String DATADOG_AGENT_STATE = "Datadog-Agent-State";
4747

4848
public static final String DEBUGGER_ENDPOINT = "debugger/v1/input";
49+
public static final String DEBUGGER_DIAGNOSTICS_ENDPOINT = "debugger/v1/diagnostics";
4950

5051
public static final String TELEMETRY_PROXY_ENDPOINT = "telemetry/proxy/";
5152

@@ -70,6 +71,7 @@ public class DDAgentFeaturesDiscovery implements DroppingPolicy {
7071
private volatile String state;
7172
private volatile String configEndpoint;
7273
private volatile String debuggerEndpoint;
74+
private volatile String debuggerDiagnosticsEndpoint;
7375
private volatile String evpProxyEndpoint;
7476
private volatile String version;
7577
private volatile String telemetryProxyEndpoint;
@@ -100,6 +102,7 @@ private void reset() {
100102
state = null;
101103
configEndpoint = null;
102104
debuggerEndpoint = null;
105+
debuggerDiagnosticsEndpoint = null;
103106
dataStreamsEndpoint = null;
104107
evpProxyEndpoint = null;
105108
version = null;
@@ -211,7 +214,7 @@ private boolean processInfoResponse(String response) {
211214
String foundMetricsEndpoint = null;
212215
if (metricsEnabled) {
213216
for (String endpoint : metricsEndpoints) {
214-
if (endpoints.contains(endpoint) || endpoints.contains("/" + endpoint)) {
217+
if (containsEndpoint(endpoints, endpoint)) {
215218
foundMetricsEndpoint = endpoint;
216219
break;
217220
}
@@ -222,39 +225,42 @@ private boolean processInfoResponse(String response) {
222225
metricsEndpoint = foundMetricsEndpoint;
223226

224227
for (String endpoint : traceEndpoints) {
225-
if (endpoints.contains(endpoint) || endpoints.contains("/" + endpoint)) {
228+
if (containsEndpoint(endpoints, endpoint)) {
226229
traceEndpoint = endpoint;
227230
break;
228231
}
229232
}
230233

231234
for (String endpoint : configEndpoints) {
232-
if (endpoints.contains(endpoint) || endpoints.contains("/" + endpoint)) {
235+
if (containsEndpoint(endpoints, endpoint)) {
233236
configEndpoint = endpoint;
234237
break;
235238
}
236239
}
237240

238-
if (endpoints.contains(DEBUGGER_ENDPOINT) || endpoints.contains("/" + DEBUGGER_ENDPOINT)) {
241+
if (containsEndpoint(endpoints, DEBUGGER_ENDPOINT)) {
239242
debuggerEndpoint = DEBUGGER_ENDPOINT;
240243
}
244+
if (containsEndpoint(endpoints, DEBUGGER_DIAGNOSTICS_ENDPOINT)) {
245+
debuggerDiagnosticsEndpoint = DEBUGGER_DIAGNOSTICS_ENDPOINT;
246+
}
241247

242248
for (String endpoint : dataStreamsEndpoints) {
243-
if (endpoints.contains(endpoint) || endpoints.contains("/" + endpoint)) {
249+
if (containsEndpoint(endpoints, endpoint)) {
244250
dataStreamsEndpoint = endpoint;
245251
break;
246252
}
247253
}
248254

249255
for (String endpoint : evpProxyEndpoints) {
250-
if (endpoints.contains(endpoint) || endpoints.contains("/" + endpoint)) {
256+
if (containsEndpoint(endpoints, endpoint)) {
251257
evpProxyEndpoint = endpoint;
252258
break;
253259
}
254260
}
255261

256262
for (String endpoint : telemetryProxyEndpoints) {
257-
if (endpoints.contains(endpoint) || endpoints.contains("/" + endpoint)) {
263+
if (containsEndpoint(endpoints, endpoint)) {
258264
telemetryProxyEndpoint = endpoint;
259265
break;
260266
}
@@ -281,6 +287,10 @@ private boolean processInfoResponse(String response) {
281287
return false;
282288
}
283289

290+
private static boolean containsEndpoint(Set<String> endpoints, String endpoint) {
291+
return endpoints.contains(endpoint) || endpoints.contains("/" + endpoint);
292+
}
293+
284294
@SuppressWarnings("unchecked")
285295
private static void discoverStatsDPort(final Map<String, Object> info) {
286296
try {
@@ -309,6 +319,10 @@ public boolean supportsDebugger() {
309319
return debuggerEndpoint != null;
310320
}
311321

322+
public boolean supportsDebuggerDiagnostics() {
323+
return debuggerDiagnosticsEndpoint != null;
324+
}
325+
312326
boolean supportsDropping() {
313327
return supportsDropping;
314328
}

‎communication/src/test/groovy/datadog/communication/ddagent/DDAgentFeaturesDiscoveryTest.groovy‎

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -59,6 +59,7 @@ class DDAgentFeaturesDiscoveryTest extends DDSpecification {
5959
features.state() == INFO_STATE
6060
features.getConfigEndpoint() == V7_CONFIG_ENDPOINT
6161
features.supportsDebugger()
62+
features.supportsDebuggerDiagnostics()
6263
features.supportsEvpProxy()
6364
features.getVersion() == "0.99.0"
6465
!features.supportsLongRunning()
@@ -87,6 +88,7 @@ class DDAgentFeaturesDiscoveryTest extends DDSpecification {
8788
features.state() == INFO_STATE
8889
features.getConfigEndpoint() == V7_CONFIG_ENDPOINT
8990
features.supportsDebugger()
91+
features.supportsDebuggerDiagnostics()
9092
features.supportsEvpProxy()
9193
features.getVersion() == "0.99.0"
9294
!features.supportsLongRunning()

‎communication/src/test/resources/agent-features/agent-info.json‎

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -14,6 +14,7 @@
1414
"/evp_proxy/v1/",
1515
"/evp_proxy/v2/",
1616
"/debugger/v1/input",
17+
"/debugger/v1/diagnostics",
1718
"/v0.7/config"
1819
],
1920
"feature_flags": [

‎dd-java-agent/agent-debugger/src/main/java/com/datadog/debugger/agent/ConfigurationUpdater.java‎

Lines changed: 8 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -59,12 +59,18 @@ DebuggerTransformer supply(
5959

6060
private Configuration currentConfiguration;
6161

62-
public ConfigurationUpdater(
62+
// used only for tests
63+
ConfigurationUpdater(
6364
Instrumentation instrumentation,
6465
TransformerSupplier transformerSupplier,
6566
Config config,
6667
ClassesToRetransformFinder finder) {
67-
this(instrumentation, transformerSupplier, config, new DebuggerSink(config), finder);
68+
this(
69+
instrumentation,
70+
transformerSupplier,
71+
config,
72+
new DebuggerSink(config, config.getFinalDebuggerSnapshotUrl(), false),
73+
finder);
6874
}
6975

7076
public ConfigurationUpdater(

‎dd-java-agent/agent-debugger/src/main/java/com/datadog/debugger/agent/DebuggerAgent.java‎

Lines changed: 12 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -50,7 +50,10 @@ public static synchronized void run(
5050
DDAgentFeaturesDiscovery ddAgentFeaturesDiscovery = sco.featuresDiscovery(config);
5151
ddAgentFeaturesDiscovery.discoverIfOutdated();
5252
agentVersion = ddAgentFeaturesDiscovery.getVersion();
53-
DebuggerSink debuggerSink = new DebuggerSink(config);
53+
String diagnosticEndpoint = getDiagnosticEndpoint(config, ddAgentFeaturesDiscovery);
54+
DebuggerSink debuggerSink =
55+
new DebuggerSink(
56+
config, diagnosticEndpoint, ddAgentFeaturesDiscovery.supportsDebuggerDiagnostics());
5457
debuggerSink.start();
5558
ConfigurationUpdater configurationUpdater =
5659
new ConfigurationUpdater(
@@ -106,6 +109,14 @@ public static synchronized void run(
106109
}
107110
}
108111

112+
private static String getDiagnosticEndpoint(
113+
Config config, DDAgentFeaturesDiscovery ddAgentFeaturesDiscovery) {
114+
if (ddAgentFeaturesDiscovery.supportsDebuggerDiagnostics()) {
115+
return config.getAgentUrl() + "/" + DDAgentFeaturesDiscovery.DEBUGGER_DIAGNOSTICS_ENDPOINT;
116+
}
117+
return config.getFinalDebuggerSnapshotUrl();
118+
}
119+
109120
private static void setupSourceFileTracking(
110121
Instrumentation instrumentation, ClassesToRetransformFinder finder) {
111122
instrumentation.addTransformer(new SourceFileTrackingTransformer(finder));

‎dd-java-agent/agent-debugger/src/main/java/com/datadog/debugger/agent/DebuggerTransformer.java‎

Lines changed: 5 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -103,7 +103,11 @@ public DebuggerTransformer(
103103

104104
// Used only for tests
105105
DebuggerTransformer(Config config, Configuration configuration) {
106-
this(config, configuration, null, new DebuggerSink(config));
106+
this(
107+
config,
108+
configuration,
109+
null,
110+
new DebuggerSink(config, config.getFinalDebuggerSnapshotUrl(), false));
107111
}
108112

109113
private void readExcludeFiles(String commaSeparatedFileNames) {

‎dd-java-agent/agent-debugger/src/main/java/com/datadog/debugger/sink/DebuggerSink.java‎

Lines changed: 16 additions & 18 deletions
Original file line numberDiff line numberDiff line change
@@ -32,30 +32,31 @@ public class DebuggerSink {
3232
private final SnapshotSink snapshotSink;
3333
private final SymbolSink symbolSink;
3434
private final DebuggerMetrics debuggerMetrics;
35-
private final BatchUploader batchUploader;
35+
private final BatchUploader snapshotUploader;
3636
private final String tags;
3737
private final int uploadFlushInterval;
3838

3939
private volatile AgentTaskScheduler.Scheduled<DebuggerSink> scheduled;
4040
private volatile AgentTaskScheduler.Scheduled<DebuggerSink> flushIntervalScheduled;
4141
private volatile long currentFlushInterval = INITIAL_FLUSH_INTERVAL;
4242

43-
public DebuggerSink(Config config) {
43+
public DebuggerSink(Config config, String diagnosticsEndpoint, boolean useMultiPart) {
4444
this(
4545
config,
4646
new BatchUploader(config, config.getFinalDebuggerSnapshotUrl()),
4747
DebuggerMetrics.getInstance(config),
48-
new ProbeStatusSink(config),
48+
new ProbeStatusSink(config, diagnosticsEndpoint, useMultiPart),
4949
new SnapshotSink(config),
5050
new SymbolSink(config));
5151
}
5252

53-
DebuggerSink(Config config, BatchUploader batchUploader) {
53+
// Used only for tests
54+
DebuggerSink(Config config, BatchUploader snapshotUploader) {
5455
this(
5556
config,
56-
batchUploader,
57+
snapshotUploader,
5758
DebuggerMetrics.getInstance(config),
58-
new ProbeStatusSink(config),
59+
new ProbeStatusSink(config, config.getFinalDebuggerSnapshotUrl(), false),
5960
new SnapshotSink(config),
6061
new SymbolSink(config));
6162
}
@@ -70,24 +71,24 @@ public DebuggerSink(Config config, ProbeStatusSink probeStatusSink) {
7071
new SymbolSink(config));
7172
}
7273

73-
DebuggerSink(Config config, BatchUploader batchUploader, DebuggerMetrics debuggerMetrics) {
74+
DebuggerSink(Config config, BatchUploader snapshotUploader, DebuggerMetrics debuggerMetrics) {
7475
this(
7576
config,
76-
batchUploader,
77+
snapshotUploader,
7778
debuggerMetrics,
78-
new ProbeStatusSink(config),
79+
new ProbeStatusSink(config, config.getFinalDebuggerSnapshotUrl(), false),
7980
new SnapshotSink(config),
8081
new SymbolSink(config));
8182
}
8283

8384
public DebuggerSink(
8485
Config config,
85-
BatchUploader batchUploader,
86+
BatchUploader snapshotUploader,
8687
DebuggerMetrics debuggerMetrics,
8788
ProbeStatusSink probeStatusSink,
8889
SnapshotSink snapshotSink,
8990
SymbolSink symbolSink) {
90-
this.batchUploader = batchUploader;
91+
this.snapshotUploader = snapshotUploader;
9192
tags = getDefaultTagsMergedWithGlobalTags(config);
9293
this.debuggerMetrics = debuggerMetrics;
9394
this.probeStatusSink = probeStatusSink;
@@ -140,7 +141,7 @@ public void stop() {
140141
}
141142

142143
public BatchUploader getSnapshotUploader() {
143-
return batchUploader;
144+
return snapshotUploader;
144145
}
145146

146147
public SymbolSink getSymbolSink() {
@@ -171,23 +172,20 @@ private void reschedule() {
171172
// visible for testing
172173
void flush(DebuggerSink ignored) {
173174
symbolSink.flush();
174-
List<String> diagnostics = probeStatusSink.getSerializedDiagnostics();
175+
probeStatusSink.flush(tags);
175176
List<String> snapshots = snapshotSink.getSerializedSnapshots();
176-
if (snapshots.size() + diagnostics.size() == 0) {
177+
if (snapshots.isEmpty()) {
177178
return;
178179
}
179180
if (snapshots.size() > 0) {
180181
uploadPayloads(snapshots);
181182
}
182-
if (diagnostics.size() > 0) {
183-
uploadPayloads(diagnostics);
184-
}
185183
}
186184

187185
private void uploadPayloads(List<String> payloads) {
188186
List<byte[]> batches = IntakeBatchHelper.createBatches(payloads);
189187
for (byte[] batch : batches) {
190-
batchUploader.upload(batch, tags);
188+
snapshotUploader.upload(batch, tags);
191189
}
192190
}
193191

‎dd-java-agent/agent-debugger/src/main/java/com/datadog/debugger/sink/ProbeStatusSink.java‎

Lines changed: 24 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -3,6 +3,7 @@
33
import com.datadog.debugger.agent.ProbeStatus;
44
import com.datadog.debugger.agent.ProbeStatus.Builder;
55
import com.datadog.debugger.agent.ProbeStatus.Status;
6+
import com.datadog.debugger.uploader.BatchUploader;
67
import com.datadog.debugger.util.ExceptionHelper;
78
import com.datadog.debugger.util.MoshiHelper;
89
import com.squareup.moshi.JsonAdapter;
@@ -26,14 +27,22 @@ public class ProbeStatusSink {
2627
private static final JsonAdapter<ProbeStatus> PROBE_STATUS_ADAPTER =
2728
MoshiHelper.createMoshiProbeStatus().adapter(ProbeStatus.class);
2829

30+
private final BatchUploader diagnosticUploader;
2931
private final Builder messageBuilder;
3032
private final Map<String, TimedMessage> probeStatuses = new ConcurrentHashMap<>();
3133
private final ArrayBlockingQueue<ProbeStatus> queue;
3234
private final Duration interval;
3335
private final int batchSize;
3436
private final boolean isInstrumentTheWorld;
37+
private final boolean useMultiPart;
3538

36-
ProbeStatusSink(Config config) {
39+
ProbeStatusSink(Config config, String diagnosticsEndpoint, boolean useMultiPart) {
40+
this(config, new BatchUploader(config, diagnosticsEndpoint), useMultiPart);
41+
}
42+
43+
ProbeStatusSink(Config config, BatchUploader diagnosticUploader, boolean useMultiPart) {
44+
this.diagnosticUploader = diagnosticUploader;
45+
this.useMultiPart = useMultiPart;
3746
this.messageBuilder = new Builder(config);
3847
this.interval = Duration.ofSeconds(config.getDebuggerDiagnosticsInterval());
3948
this.batchSize = config.getDebuggerUploadBatchSize();
@@ -61,7 +70,20 @@ public void addError(ProbeId probeId, String message) {
6170
addDiagnostics(messageBuilder.errorMessage(probeId, message));
6271
}
6372

64-
public List<String> getSerializedDiagnostics() {
73+
public void flush(String tags) {
74+
List<String> serializedDiagnostics = getSerializedDiagnostics();
75+
List<byte[]> batches = IntakeBatchHelper.createBatches(serializedDiagnostics);
76+
for (byte[] batch : batches) {
77+
if (useMultiPart) {
78+
diagnosticUploader.uploadAsMultipart(
79+
tags, new BatchUploader.MultiPartContent(batch, "event", "event.json"));
80+
} else {
81+
diagnosticUploader.upload(batch, tags);
82+
}
83+
}
84+
}
85+
86+
private List<String> getSerializedDiagnostics() {
6587
List<ProbeStatus> diagnostics = getDiagnostics();
6688
List<String> serializedDiagnostics = new ArrayList<>();
6789
for (ProbeStatus message : diagnostics) {

‎dd-java-agent/agent-debugger/src/main/java/com/datadog/debugger/sink/SymbolSink.java‎

Lines changed: 3 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -34,8 +34,8 @@ public class SymbolSink {
3434
private final String env;
3535
private final String version;
3636
private final BatchUploader symbolUploader;
37-
private final BlockingQueue<ServiceVersion> scopes = new ArrayBlockingQueue<>(CAPACITY);
3837
private final BatchUploader.MultiPartContent event;
38+
private final BlockingQueue<ServiceVersion> scopes = new ArrayBlockingQueue<>(CAPACITY);
3939

4040
public SymbolSink(Config config) {
4141
this(config, new BatchUploader(config, config.getFinalDebuggerSymDBUrl()));
@@ -47,7 +47,8 @@ public SymbolSink(Config config) {
4747
this.version = config.getVersion();
4848
this.symbolUploader = symbolUploader;
4949
byte[] eventContent =
50-
String.format(EVENT_FORMAT, serviceName, config.getRuntimeId())
50+
String.format(
51+
EVENT_FORMAT, TagsHelper.sanitize(config.getServiceName()), config.getRuntimeId())
5152
.getBytes(StandardCharsets.UTF_8);
5253
this.event = new BatchUploader.MultiPartContent(eventContent, "event", "event.json");
5354
}

0 commit comments

Comments
 (0)