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
2 changes: 0 additions & 2 deletions tracer/missing-nullability-files.csv
Original file line number Diff line number Diff line change
Expand Up @@ -34,7 +34,6 @@ src/Datadog.Trace/TracerManagerFactory.cs
src/Datadog.Trace/Agent/Api.cs
src/Datadog.Trace/Agent/ApiOtlp.cs
src/Datadog.Trace/Agent/IApi.cs
src/Datadog.Trace/Agent/IApiRequest.cs
src/Datadog.Trace/Agent/IApiRequestFactory.cs
src/Datadog.Trace/Agent/IApiResponseTelemetryExtensions.cs
src/Datadog.Trace/Agent/IKeepRateCalculator.cs
Expand Down Expand Up @@ -212,7 +211,6 @@ src/Datadog.Trace/Agent/TraceSamplers/ErrorSampler.cs
src/Datadog.Trace/Agent/TraceSamplers/ITraceChunkSampler.cs
src/Datadog.Trace/Agent/TraceSamplers/PrioritySampler.cs
src/Datadog.Trace/Agent/TraceSamplers/RareSampler.cs
src/Datadog.Trace/Agent/Transports/ApiWebRequest.cs
src/Datadog.Trace/Agent/Transports/ApiWebRequestFactory.cs
src/Datadog.Trace/Agent/Transports/HttpClientRequest.cs
src/Datadog.Trace/Agent/Transports/HttpClientRequestFactory.cs
Expand Down
4 changes: 2 additions & 2 deletions tracer/src/Datadog.Trace/Agent/Api.cs
Original file line number Diff line number Diff line change
Expand Up @@ -238,7 +238,7 @@ private async Task<SendResult> SendStatsAsyncImpl(IApiRequest request, bool isFi
}
catch (Exception ex)
{
var tag = ex is TimeoutException ? MetricTags.ApiError.Timeout : MetricTags.ApiError.NetworkError;
var tag = ex is TimeoutException or OperationCanceledException ? MetricTags.ApiError.Timeout : MetricTags.ApiError.NetworkError;
TelemetryFactory.Metrics.RecordCountStatsApiErrors(tag);
throw;
}
Expand Down Expand Up @@ -321,7 +321,7 @@ private async Task<SendResult> SendTracesAsyncImpl(IApiRequest request, bool fin
{
// count only network/infrastructure errors, not valid responses with error status codes
// (which are handled below)
var tag = ex is TimeoutException ? MetricTags.ApiError.Timeout : MetricTags.ApiError.NetworkError;
var tag = ex is TimeoutException or OperationCanceledException ? MetricTags.ApiError.Timeout : MetricTags.ApiError.NetworkError;
TelemetryFactory.Metrics.RecordCountTraceApiErrors(tag);
healthStats?.Increment(TracerMetricNames.Api.Errors);
throw;
Expand Down
6 changes: 4 additions & 2 deletions tracer/src/Datadog.Trace/Agent/IApiRequest.cs
Original file line number Diff line number Diff line change
Expand Up @@ -3,6 +3,8 @@
// This product includes software developed at Datadog (https://www.datadoghq.com/). Copyright 2017 Datadog, Inc.
// </copyright>

#nullable enable

using System;
using System.IO;
using System.Threading.Tasks;
Expand All @@ -19,13 +21,13 @@ internal interface IApiRequest

Task<IApiResponse> PostAsync(ArraySegment<byte> bytes, string contentType);

Task<IApiResponse> PostAsync(ArraySegment<byte> bytes, string contentType, string contentEncoding);
Task<IApiResponse> PostAsync(ArraySegment<byte> bytes, string contentType, string? contentEncoding);

Task<IApiResponse> PostAsJsonAsync<T>(T payload, MultipartCompression compression);

Task<IApiResponse> PostAsJsonAsync<T>(T payload, MultipartCompression compression, JsonSerializerSettings settings);

Task<IApiResponse> PostAsync(Func<Stream, Task> writeToRequestStream, string contentType, string contentEncoding, string multipartBoundary);
Task<IApiResponse> PostAsync(Func<Stream, Task> writeToRequestStream, string contentType, string? contentEncoding, string multipartBoundary);

Task<IApiResponse> PostAsync(MultipartFormItem[] items, MultipartCompression multipartCompression = MultipartCompression.None);
}
Expand Down
215 changes: 113 additions & 102 deletions tracer/src/Datadog.Trace/Agent/Transports/ApiWebRequest.cs
Original file line number Diff line number Diff line change
Expand Up @@ -3,13 +3,17 @@
// This product includes software developed at Datadog (https://www.datadoghq.com/). Copyright 2017 Datadog, Inc.
// </copyright>

#nullable enable

using System;
using System.IO;
using System.IO.Compression;
using System.Net;
using System.Text;
using System.Threading;
using System.Threading.Tasks;
using Datadog.Trace.Logging;
using Datadog.Trace.SourceGenerators;
using Datadog.Trace.Util;
using Datadog.Trace.Vendors.Newtonsoft.Json;
using Datadog.Trace.Vendors.Serilog.Events;
Expand All @@ -19,15 +23,9 @@ namespace Datadog.Trace.Agent.Transports
{
internal sealed class ApiWebRequest : IApiRequest
{
private const string BoundarySeparator = $"{CrLf}--{Boundary}{CrLf}";
private const string BoundaryTrailer = $"{CrLf}--{Boundary}--{CrLf}";

private static readonly IDatadogLogger Log = DatadogLogging.GetLoggerFor<ApiWebRequest>();
private readonly HttpWebRequest _request;

private byte[] _boundarySeparatorInBytes;
private byte[] _boundaryTrailerInBytes;

public ApiWebRequest(HttpWebRequest request)
{
_request = request;
Expand All @@ -39,59 +37,45 @@ public void AddHeader(string name, string value)
}

public Task<IApiResponse> GetAsync()
{
ResetRequest(method: "GET", contentType: null, contentEncoding: null);

return FinishAndGetResponse();
}
=> SendAsync(method: "GET", contentType: null, contentEncoding: null, state: this, writeBody: null);

public Task<IApiResponse> PostAsync(ArraySegment<byte> bytes, string contentType)
=> PostAsync(bytes, contentType, null);

public async Task<IApiResponse> PostAsync(ArraySegment<byte> bytes, string contentType, string contentEncoding)
{
ResetRequest(method: "POST", contentType, contentEncoding);

using (var requestStream = await _request.GetRequestStreamAsync().ConfigureAwait(false))
{
await requestStream.WriteAsync(bytes.Array, bytes.Offset, bytes.Count).ConfigureAwait(false);
}

return await FinishAndGetResponse().ConfigureAwait(false);
}
public Task<IApiResponse> PostAsync(ArraySegment<byte> bytes, string contentType, string? contentEncoding)
=> SendAsync(
method: "POST",
contentType,
contentEncoding,
state: bytes,
writeBody: static (requestStream, body) => requestStream.WriteAsync(body.Array!, body.Offset, body.Count));

public Task<IApiResponse> PostAsJsonAsync<T>(T payload, MultipartCompression compression)
=> PostAsJsonAsync(payload, compression, SerializationHelpers.DefaultJsonSettings);

public async Task<IApiResponse> PostAsJsonAsync<T>(T payload, MultipartCompression compression, JsonSerializerSettings settings)
public Task<IApiResponse> PostAsJsonAsync<T>(T payload, MultipartCompression compression, JsonSerializerSettings settings)
{
var contentEncoding = compression == MultipartCompression.GZip ? "gzip" : null;
if (Log.IsEnabled(LogEventLevel.Debug))
{
Log.Debug("Sending {Type} data as JSON with compression '{Compression}'", typeof(T).FullName, contentEncoding ?? "none");
}

ResetRequest(method: "POST", contentType: MimeTypes.Json, contentEncoding: contentEncoding);

using (var reqStream = await _request.GetRequestStreamAsync().ConfigureAwait(false))
{
await SerializationHelpers.WriteAsJson(reqStream, payload, settings, compression).ConfigureAwait(false);
}

return await FinishAndGetResponse().ConfigureAwait(false);
return SendAsync(
method: "POST",
contentType: MimeTypes.Json,
contentEncoding,
state: new JsonState<T>(payload, settings, compression),
writeBody: static (reqStream, state) => SerializationHelpers.WriteAsJson(reqStream, state.Payload, state.Settings, state.Compression));
}

public async Task<IApiResponse> PostAsync(Func<Stream, Task> writeToRequestStream, string contentType, string contentEncoding, string multipartBoundary)
{
ResetRequest(method: "POST", ContentTypeHelper.GetContentType(contentType, multipartBoundary), contentEncoding);

using (var requestStream = await _request.GetRequestStreamAsync().ConfigureAwait(false))
{
await writeToRequestStream(requestStream).ConfigureAwait(false);
}

return await FinishAndGetResponse().ConfigureAwait(false);
}
public Task<IApiResponse> PostAsync(Func<Stream, Task> writeToRequestStream, string contentType, string? contentEncoding, string multipartBoundary)
=> SendAsync(
method: "POST",
ContentTypeHelper.GetContentType(contentType, multipartBoundary),
contentEncoding,
state: writeToRequestStream,
writeBody: static (stream, wb) => wb(stream));

/// <summary>
/// Send a Post request using multipart form data.
Expand All @@ -100,7 +84,7 @@ public async Task<IApiResponse> PostAsync(Func<Stream, Task> writeToRequestStrea
/// <param name="items">Multipart form data items</param>
/// <param name="multipartCompression">Multipart compression</param>
/// <returns>Task with the response</returns>
public async Task<IApiResponse> PostAsync(MultipartFormItem[] items, MultipartCompression multipartCompression = MultipartCompression.None)
public Task<IApiResponse> PostAsync(MultipartFormItem[] items, MultipartCompression multipartCompression = MultipartCompression.None)
{
if (items is null)
{
Expand All @@ -109,80 +93,84 @@ public async Task<IApiResponse> PostAsync(MultipartFormItem[] items, MultipartCo

Log.Debug<int>("Sending multipart form request with {Count} items.", items.Length);

ResetRequest(method: "POST", contentType: "multipart/form-data; boundary=" + Boundary, contentEncoding: multipartCompression == MultipartCompression.GZip ? "gzip" : null);
using (var reqStream = await _request.GetRequestStreamAsync().ConfigureAwait(false))
return SendAsync(
method: "POST",
contentType: "multipart/form-data; boundary=" + Boundary,
contentEncoding: multipartCompression == MultipartCompression.GZip ? "gzip" : null,
state: new MultipartState(items, multipartCompression),
writeBody: static (requestStream, state) => WriteMultipartAsync(state.Items, requestStream, state.Compression));
}

private static async Task WriteMultipartAsync(MultipartFormItem[] items, Stream reqStream, MultipartCompression multipartCompression)
{
if (multipartCompression == MultipartCompression.GZip)
{
if (multipartCompression == MultipartCompression.GZip)
{
Log.Debug("Using MultipartCompression.GZip");
using var gzip = new GZipStream(reqStream, CompressionMode.Compress, leaveOpen: true);
await WriteToStreamAsync(items, gzip).ConfigureAwait(false);
await gzip.FlushAsync().ConfigureAwait(false);
Log.Debug("Compressing multipart payload...");
}
else
{
await WriteToStreamAsync(items, reqStream).ConfigureAwait(false);
}
Log.Debug("Using MultipartCompression.GZip");
using var gzip = new GZipStream(reqStream, CompressionMode.Compress, leaveOpen: true);
await WriteToStreamAsync(items, gzip).ConfigureAwait(false);
await gzip.FlushAsync().ConfigureAwait(false);
Log.Debug("Compressing multipart payload...");
}
else
{
await WriteToStreamAsync(items, reqStream).ConfigureAwait(false);
}
}

return await FinishAndGetResponse().ConfigureAwait(false);
private static async Task WriteToStreamAsync(MultipartFormItem[] multipartItems, Stream requestStream)
{
// Write form request using the boundary
var boundaryBytes = MultipartBytes.BoundarySeparator;
var trailerBytes = MultipartBytes.BoundaryTrailer;

async Task WriteToStreamAsync(MultipartFormItem[] multipartItems, Stream requestStream)
// Write each MultipartFormItem
var itemsWritten = 0;
foreach (var item in multipartItems)
{
// Write form request using the boundary
var boundaryBytes = _boundarySeparatorInBytes ??= Encoding.ASCII.GetBytes(BoundarySeparator);
var trailerBytes = _boundaryTrailerInBytes ??= Encoding.ASCII.GetBytes(BoundaryTrailer);

// Write each MultipartFormItem
var itemsWritten = 0;
foreach (var item in multipartItems)
if (!item.IsValid(Log))
{
if (!item.IsValid(Log))
{
continue;
}

var headerBytes = Encoding.ASCII.GetBytes(
item.FileName is not null
? $"Content-Type: {item.ContentType}\r\nContent-Disposition: form-data; name=\"{item.Name}\"; filename=\"{item.FileName}\"\r\n\r\n"
: $"Content-Type: {item.ContentType}\r\nContent-Disposition: form-data; name=\"{item.Name}\"\r\n\r\n");

if (itemsWritten == 0)
{
// If we are writing the first item, we skip the initial `\r\n` in the array
await requestStream.WriteAsync(boundaryBytes, 2, boundaryBytes.Length - 2).ConfigureAwait(false);
}
else
{
await requestStream.WriteAsync(boundaryBytes, 0, boundaryBytes.Length).ConfigureAwait(false);
}

await requestStream.WriteAsync(headerBytes, 0, headerBytes.Length).ConfigureAwait(false);
if (item.ContentInBytes is { } arraySegment)
{
Log.Debug("Adding to Multipart Byte Array | Name: {Name} | FileName: {FileName} | ContentType: {ContentType}", item.Name, item.FileName, item.ContentType);
await requestStream.WriteAsync(arraySegment.Array, arraySegment.Offset, arraySegment.Count).ConfigureAwait(false);
}
else if (item.ContentInStream is { } stream)
{
Log.Debug("Adding to Multipart Stream | Name: {Name} | FileName: {FileName} | ContentType: {ContentType}", item.Name, item.FileName, item.ContentType);
await stream.CopyToAsync(requestStream).ConfigureAwait(false);
}

itemsWritten++;
continue;
}

var headerBytes = Encoding.ASCII.GetBytes(
item.FileName is not null
? $"Content-Type: {item.ContentType}\r\nContent-Disposition: form-data; name=\"{item.Name}\"; filename=\"{item.FileName}\"\r\n\r\n"
: $"Content-Type: {item.ContentType}\r\nContent-Disposition: form-data; name=\"{item.Name}\"\r\n\r\n");

if (itemsWritten == 0)
{
// If we are writing the first item, we skip the initial `\r\n` in the array
await requestStream.WriteAsync(boundaryBytes, 2, boundaryBytes.Length - 2).ConfigureAwait(false);
}
else
{
await requestStream.WriteAsync(boundaryBytes, 0, boundaryBytes.Length).ConfigureAwait(false);
}

await requestStream.WriteAsync(trailerBytes, 0, trailerBytes.Length).ConfigureAwait(false);
await requestStream.WriteAsync(headerBytes, 0, headerBytes.Length).ConfigureAwait(false);
if (item.ContentInBytes is { } arraySegment)
{
Log.Debug("Adding to Multipart Byte Array | Name: {Name} | FileName: {FileName} | ContentType: {ContentType}", item.Name, item.FileName, item.ContentType);
await requestStream.WriteAsync(arraySegment.Array!, arraySegment.Offset, arraySegment.Count).ConfigureAwait(false);
}
else if (item.ContentInStream is { } stream)
{
Log.Debug("Adding to Multipart Stream | Name: {Name} | FileName: {FileName} | ContentType: {ContentType}", item.Name, item.FileName, item.ContentType);
await stream.CopyToAsync(requestStream).ConfigureAwait(false);
}

itemsWritten++;
}

if (itemsWritten == 0)
{
await requestStream.WriteAsync(boundaryBytes, 2, boundaryBytes.Length - 2).ConfigureAwait(false);
}

await requestStream.WriteAsync(trailerBytes, 0, trailerBytes.Length).ConfigureAwait(false);
}

private void ResetRequest(string method, string contentType, string contentEncoding)
private void ResetRequest(string method, string? contentType, string? contentEncoding)
{
_request.Method = method;
_request.ContentType = string.IsNullOrEmpty(contentType) ? null : contentType;
Expand All @@ -196,10 +184,20 @@ private void ResetRequest(string method, string contentType, string contentEncod
}
}

private async Task<IApiResponse> FinishAndGetResponse()
private async Task<IApiResponse> SendAsync<TState>(string method, string? contentType, string? contentEncoding, TState state, Func<Stream, TState, Task>? writeBody)
{
try
{
ResetRequest(method, contentType, contentEncoding);

if (writeBody is not null)
{
using (var requestStream = await _request.GetRequestStreamAsync().ConfigureAwait(false))
{
await writeBody(requestStream, state).ConfigureAwait(false);
}
}

var httpWebResponse = (HttpWebResponse)await _request.GetResponseAsync().ConfigureAwait(false);
return new ApiWebResponse(httpWebResponse);
}
Expand All @@ -210,5 +208,18 @@ private async Task<IApiResponse> FinishAndGetResponse()
return new ApiWebResponse((HttpWebResponse)exception.Response);
}
}

private readonly struct JsonState<T>(T payload, JsonSerializerSettings settings, MultipartCompression compression)
{
public readonly T Payload = payload;
public readonly JsonSerializerSettings Settings = settings;
public readonly MultipartCompression Compression = compression;
}

private readonly struct MultipartState(MultipartFormItem[] items, MultipartCompression compression)
{
public readonly MultipartFormItem[] Items = items;
public readonly MultipartCompression Compression = compression;
}
}
}
4 changes: 2 additions & 2 deletions tracer/src/Datadog.Trace/Ci/Agent/CIWriterHttpSender.cs
Original file line number Diff line number Diff line change
Expand Up @@ -111,9 +111,9 @@ private async Task SendPayloadAsync<T>(EventPlatformPayload payload, Func<IApiRe
{
request.AddHeader(EvpSubdomainHeader, payload.EventPlatformSubdomain);
}
else
else if (TestOptimization.Instance.Settings.ApiKey is { } apiKey)
{
request.AddHeader(ApiKeyHeader, TestOptimization.Instance.Settings.ApiKey);
request.AddHeader(ApiKeyHeader, apiKey);
}
}
catch (Exception ex)
Expand Down
4 changes: 2 additions & 2 deletions tracer/src/Datadog.Trace/Ci/Net/TestOptimizationClient.cs
Original file line number Diff line number Diff line change
Expand Up @@ -414,9 +414,9 @@ private void SetRequestHeader(IApiRequest request)
{
request.AddHeader(EvpSubdomainHeader, "api");
}
else
else if (_testOptimization.Settings.ApiKey is { } apiKey)
{
request.AddHeader(ApiKeyHeader, _testOptimization.Settings.ApiKey);
request.AddHeader(ApiKeyHeader, apiKey);
}
}

Expand Down
Loading
Loading