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
Original file line number Diff line number Diff line change
@@ -1,5 +1,6 @@
package datadog.trace.common.writer;

import datadog.trace.api.telemetry.OtlpTelemetry;
import datadog.trace.core.CoreSpan;
import datadog.trace.core.otlp.common.OtlpPayload;
import datadog.trace.core.otlp.common.OtlpSender;
Expand All @@ -26,13 +27,15 @@ public void addTrace(List<? extends CoreSpan<?>> trace) {
public void flush() {
OtlpPayload payload = collector.collectTraces();
if (payload != OtlpPayload.EMPTY) {
sender.send(payload);
OtlpTelemetry.getInstance().onTracesExportAttempt();
RemoteApi.Response response = sender.send(payload);
OtlpTelemetry.getInstance().onTracesExportComplete(response.success());
}
}

@Override
public void onDroppedTrace(int spanCount) {
// TODO: surface drop counts via HealthMetrics
// no telemetry currently tracked for dropped traces
}

@Override
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -38,11 +38,10 @@ public static OtlpWriterBuilder builder() {
TraceProcessingWorker worker,
PayloadDispatcher dispatcher,
OtlpSender sender,
HealthMetrics healthMetrics,
int flushTimeout,
TimeUnit flushTimeoutUnit,
boolean alwaysFlush) {
super(worker, dispatcher, healthMetrics, flushTimeout, flushTimeoutUnit, alwaysFlush);
super(worker, dispatcher, HealthMetrics.NO_OP, flushTimeout, flushTimeoutUnit, alwaysFlush);
this.sender = sender;
}

Expand All @@ -64,7 +63,6 @@ public static class OtlpWriterBuilder {
private OtlpConfig.Protocol protocol = OtlpConfig.Protocol.HTTP_PROTOBUF;
private OtlpConfig.Compression compression = OtlpConfig.Compression.NONE;
private int traceBufferSize = BUFFER_SIZE;
private HealthMetrics healthMetrics = HealthMetrics.NO_OP;
private int flushIntervalMilliseconds = 1000;
private int flushTimeout = 1;
private TimeUnit flushTimeoutUnit = TimeUnit.SECONDS;
Expand Down Expand Up @@ -102,11 +100,6 @@ public OtlpWriterBuilder traceBufferSize(int traceBufferSize) {
return this;
}

public OtlpWriterBuilder healthMetrics(HealthMetrics healthMetrics) {
this.healthMetrics = healthMetrics;
return this;
}

public OtlpWriterBuilder flushIntervalMilliseconds(int flushIntervalMilliseconds) {
this.flushIntervalMilliseconds = flushIntervalMilliseconds;
return this;
Expand Down Expand Up @@ -151,7 +144,7 @@ public OtlpWriter build() {
final TraceProcessingWorker worker =
new TraceProcessingWorker(
traceBufferSize,
healthMetrics,
HealthMetrics.NO_OP,
Comment thread
mcculls marked this conversation as resolved.
dispatcher,
DroppingPolicy.DISABLED,
Prioritization.FAST_LANE,
Expand All @@ -160,7 +153,7 @@ public OtlpWriter build() {
singleSpanSampler);

return new OtlpWriter(
worker, dispatcher, sender, healthMetrics, flushTimeout, flushTimeoutUnit, alwaysFlush);
worker, dispatcher, sender, flushTimeout, flushTimeoutUnit, alwaysFlush);
}
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -76,7 +76,6 @@ public static Writer createWriter(
.protocol(config.getOtlpTracesProtocol())
.compression(config.getOtlpTracesCompression())
.timeoutMillis(config.getOtlpTracesTimeout())
.healthMetrics(healthMetrics)
.spanSamplingRules(singleSpanSampler)
.flushIntervalMilliseconds(flushIntervalMilliseconds)
.build();
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -2,18 +2,16 @@

import static datadog.communication.http.OkHttpUtils.buildHttp2Client;
import static datadog.communication.http.OkHttpUtils.isPlainHttp;
import static datadog.communication.http.OkHttpUtils.sendWithRetries;

import datadog.communication.http.HttpRetryPolicy;
import datadog.logging.RatelimitedLogger;
import datadog.trace.api.config.OtlpConfig.Compression;
import java.io.IOException;
import datadog.trace.common.writer.RemoteApi;
import java.util.Map;
import java.util.concurrent.TimeUnit;
import okhttp3.HttpUrl;
import okhttp3.OkHttpClient;
import okhttp3.Request;
import okhttp3.Response;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;

Expand Down Expand Up @@ -55,20 +53,8 @@ public OtlpGrpcSender(
}

@Override
public void send(OtlpPayload payload) {
Request request = makeRequest(payload);
try (Response response = sendWithRetries(client, retryPolicy, request)) {
if (!response.isSuccessful()) {
RATELIMITED_LOGGER.warn(
"OTLP export to {} failed with status {}: {}",
request.url(),
response.code(),
response.message());
}
} catch (IOException e) {
RATELIMITED_LOGGER.warn(
"OTLP export to {} failed with exception: {}", request.url(), e.toString());
}
public RemoteApi.Response send(OtlpPayload payload) {
return OtlpSenderSupport.send(client, retryPolicy, makeRequest(payload), RATELIMITED_LOGGER);
}

@Override
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -2,18 +2,16 @@

import static datadog.communication.http.OkHttpUtils.buildHttpClient;
import static datadog.communication.http.OkHttpUtils.isPlainHttp;
import static datadog.communication.http.OkHttpUtils.sendWithRetries;

import datadog.communication.http.HttpRetryPolicy;
import datadog.logging.RatelimitedLogger;
import datadog.trace.api.config.OtlpConfig.Compression;
import java.io.IOException;
import datadog.trace.common.writer.RemoteApi;
import java.util.Map;
import java.util.concurrent.TimeUnit;
import okhttp3.HttpUrl;
import okhttp3.OkHttpClient;
import okhttp3.Request;
import okhttp3.Response;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;

Expand Down Expand Up @@ -59,20 +57,8 @@ public HttpUrl url() {
}

@Override
public void send(OtlpPayload payload) {
Request request = makeRequest(payload);
try (Response response = sendWithRetries(client, retryPolicy, request)) {
if (!response.isSuccessful()) {
RATELIMITED_LOGGER.warn(
"OTLP export to {} failed with status {}: {}",
request.url(),
response.code(),
response.message());
}
} catch (IOException e) {
RATELIMITED_LOGGER.warn(
"OTLP export to {} failed with exception: {}", request.url(), e.toString());
}
public RemoteApi.Response send(OtlpPayload payload) {
return OtlpSenderSupport.send(client, retryPolicy, makeRequest(payload), RATELIMITED_LOGGER);
}

@Override
Expand Down
Original file line number Diff line number Diff line change
@@ -1,8 +1,10 @@
package datadog.trace.core.otlp.common;

import datadog.trace.common.writer.RemoteApi;

/** Sends chunks of OTLP data. */
public interface OtlpSender {
void send(OtlpPayload payload);
RemoteApi.Response send(OtlpPayload payload);

void shutdown();
}
Original file line number Diff line number Diff line change
@@ -0,0 +1,39 @@
package datadog.trace.core.otlp.common;

import static datadog.communication.http.OkHttpUtils.sendWithRetries;
import static datadog.trace.common.writer.RemoteApi.Response.failed;
import static datadog.trace.common.writer.RemoteApi.Response.success;

import datadog.communication.http.HttpRetryPolicy;
import datadog.logging.RatelimitedLogger;
import datadog.trace.common.writer.RemoteApi;
import java.io.IOException;
import okhttp3.OkHttpClient;

/** Shared request execution and response handling for {@link OtlpSender} implementations. */
final class OtlpSenderSupport {
private OtlpSenderSupport() {}

/** Executes the given request with retries, logging failures via the rate-limited logger. */
static RemoteApi.Response send(
OkHttpClient client,
HttpRetryPolicy.Factory retryPolicy,
okhttp3.Request request,
RatelimitedLogger ratelimitedLogger) {
try (okhttp3.Response response = sendWithRetries(client, retryPolicy, request)) {
if (response.isSuccessful()) {
return success(response.code());
Comment thread
mcculls marked this conversation as resolved.
Comment thread
mcculls marked this conversation as resolved.
}
ratelimitedLogger.warn(
"OTLP export to {} failed with status {}: {}",
request.url(),
response.code(),
response.message());
return failed(response.code());
} catch (IOException e) {
ratelimitedLogger.warn(
"OTLP export to {} failed with exception: {}", request.url(), e.toString());
return failed(e);
}
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -7,4 +7,7 @@ public abstract class OtlpLogsCollector {

/** Waits for logs to be batched within the given interval. */
public abstract OtlpPayload waitForLogs(int intervalMillis);

/** Number of log records collected. */
public abstract int getLogRecordCount();
}
Original file line number Diff line number Diff line change
Expand Up @@ -36,6 +36,7 @@ public final class OtlpLogsJsonCollector extends OtlpLogsCollector
private JsonWriter writer;
private boolean anyLogRecordWritten;
private boolean logRecordStarted;
private int logRecordCount;

private final LazyJsonArray attributesArray = new LazyJsonArray();

Expand Down Expand Up @@ -65,6 +66,8 @@ OtlpPayload collectLogs(ObjIntConsumer<OtlpLogsVisitor> processor, int intervalM

/** Prepare temporary elements to collect logs data. */
private void start() {
logRecordCount = 0;

writer = new JsonWriter();
writer.beginObject();
writer.name("resourceLogs").beginArray();
Expand All @@ -86,6 +89,11 @@ private void stop() {
currentScope = null;
}

@Override
public int getLogRecordCount() {
return logRecordCount;
}

@Override
public OtlpScopedLogsVisitor visitScopedLogs(OtelInstrumentationScope scope) {
if (currentScope != null) {
Expand Down Expand Up @@ -119,6 +127,7 @@ public void visitLogRecord(OtlpLogRecord logRecord) {

logRecordStarted = false;
anyLogRecordWritten = true;
logRecordCount++;
}

// opens the log record object on first attribute or value written for it
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -41,6 +41,7 @@ public final class OtlpLogsProtoCollector extends OtlpLogsCollector
// total number of chunked bytes at different nesting levels
private int payloadBytes;
private int scopedBytes;
private int logRecordCount;

private OtelInstrumentationScope currentScope;

Expand Down Expand Up @@ -68,6 +69,7 @@ OtlpPayload collectLogs(ObjIntConsumer<OtlpLogsVisitor> processor, int intervalM

/** Prepare temporary elements to collect logs data. */
private void start() {
logRecordCount = 0;

// remove stale entries from caches
OtlpCommonProto.recalibrateCaches();
Expand All @@ -84,6 +86,11 @@ private void stop() {
currentScope = null;
}
Comment thread
mcculls marked this conversation as resolved.

@Override
public int getLogRecordCount() {
return logRecordCount;
}

@Override
public OtlpScopedLogsVisitor visitScopedLogs(OtelInstrumentationScope scope) {
if (currentScope != null) {
Expand All @@ -96,6 +103,7 @@ public OtlpScopedLogsVisitor visitScopedLogs(OtelInstrumentationScope scope) {
@Override
public void visitLogRecord(OtlpLogRecord logRecord) {
scopedBytes += recordLogRecordMessage(buf, logRecord, protobuf);
logRecordCount++;
}

@Override
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -4,6 +4,8 @@
import static datadog.trace.util.AgentThreadFactory.newAgentThread;

import datadog.trace.api.Config;
import datadog.trace.api.telemetry.OtlpTelemetry;
import datadog.trace.common.writer.RemoteApi;
import datadog.trace.core.otlp.common.OtlpGrpcSender;
import datadog.trace.core.otlp.common.OtlpHttpSender;
import datadog.trace.core.otlp.common.OtlpPayload;
Expand Down Expand Up @@ -108,7 +110,11 @@ private void export() {
try {
OtlpPayload payload = collector.waitForLogs(intervalMillis);
if (payload != OtlpPayload.EMPTY) {
sender.send(payload);
int logRecordCount = collector.getLogRecordCount();
Comment thread
mcculls marked this conversation as resolved.
RemoteApi.Response response = sender.send(payload);
if (response.success()) {
OtlpTelemetry.getInstance().onLogRecordsSubmitted(logRecordCount);
}
}
} catch (RuntimeException e) {
LOGGER.debug("Uncaught exception exporting logs", e);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -4,7 +4,9 @@

import datadog.trace.api.Config;
import datadog.trace.api.config.OtlpConfig;
import datadog.trace.api.telemetry.OtlpTelemetry;
import datadog.trace.api.time.SystemTimeSource;
import datadog.trace.common.writer.RemoteApi;
import datadog.trace.core.otlp.common.OtlpPayload;
import datadog.trace.core.otlp.common.OtlpSender;
import datadog.trace.util.AgentTaskScheduler;
Expand All @@ -29,7 +31,6 @@ public final class OtlpMetricsService {

OtlpMetricsService(Config config) {
this.scheduler = new AgentTaskScheduler(OTLP_METRICS_EXPORTER);

this.sender = OtlpMetricsSenderFactory.create(config);
if (this.sender == null) {
LOGGER.debug("Unsupported OTLP metrics protocol: {}", config.getOtlpMetricsProtocol());
Expand Down Expand Up @@ -91,7 +92,9 @@ public void shutdown() {
private void export() {
OtlpPayload payload = collector.collectMetrics();
if (payload != OtlpPayload.EMPTY) {
sender.send(payload);
OtlpTelemetry.getInstance().onMetricsExportAttempt();
RemoteApi.Response response = sender.send(payload);
OtlpTelemetry.getInstance().onMetricsExportComplete(response.success());
}
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -7,6 +7,7 @@
import datadog.metrics.api.Histogram;
import datadog.trace.api.Config;
import datadog.trace.api.config.OtlpConfig;
import datadog.trace.api.telemetry.OtlpTelemetry;
import datadog.trace.api.time.SystemTimeSource;
import datadog.trace.bootstrap.instrumentation.api.UTF8BytesString;
import datadog.trace.bootstrap.otel.common.OtelInstrumentationScope;
Expand All @@ -16,6 +17,7 @@
import datadog.trace.bootstrap.otlp.metrics.OtlpMetricsVisitor;
import datadog.trace.common.metrics.AggregateEntry;
import datadog.trace.common.metrics.MetricWriter;
import datadog.trace.common.writer.RemoteApi;
import datadog.trace.core.otlp.common.OtlpPayload;
import datadog.trace.core.otlp.common.OtlpResourceJson;
import datadog.trace.core.otlp.common.OtlpResourceProto;
Expand Down Expand Up @@ -167,7 +169,9 @@ public void finishBucket() {
}
OtlpPayload payload = collector.collectMetrics(this::emit, startNanos, endNanos);
if (payload != OtlpPayload.EMPTY) {
sender.send(payload);
OtlpTelemetry.getInstance().onMetricsExportAttempt();
RemoteApi.Response response = sender.send(payload);
OtlpTelemetry.getInstance().onMetricsExportComplete(response.success());
}
} finally {
pending.clear();
Expand Down
Loading