From f0a6aacc4cd971ed724862f28e608359d2d4e6d0 Mon Sep 17 00:00:00 2001 From: Andrew Kent Date: Mon, 21 Sep 2026 02:13:26 -0600 Subject: [PATCH] feat: add span customizers (SDK-316) --- src/Braintrust.Sdk/Braintrust.Sdk.csproj | 10 + src/Braintrust.Sdk/Config/BraintrustConfig.cs | 37 +- .../Trace/BraintrustSpanCustomizerHandler.cs | 159 ++++ src/Braintrust.Sdk/Trace/BraintrustTracing.cs | 24 +- src/Braintrust.Sdk/Trace/ISpanCustomizer.cs | 21 + src/Braintrust.Sdk/Trace/OtlpJsonPayload.cs | 167 ++++ .../collector/trace/v1/trace_service.proto | 79 ++ .../proto/common/v1/common.proto | 81 ++ .../proto/resource/v1/resource.proto | 37 + .../opentelemetry/proto/trace/v1/trace.proto | 357 ++++++++ .../Trace/SpanCustomizerTest.cs | 764 ++++++++++++++++++ 11 files changed, 1727 insertions(+), 9 deletions(-) create mode 100644 src/Braintrust.Sdk/Trace/BraintrustSpanCustomizerHandler.cs create mode 100644 src/Braintrust.Sdk/Trace/ISpanCustomizer.cs create mode 100644 src/Braintrust.Sdk/Trace/OtlpJsonPayload.cs create mode 100644 src/Braintrust.Sdk/Trace/Protos/opentelemetry/proto/collector/trace/v1/trace_service.proto create mode 100644 src/Braintrust.Sdk/Trace/Protos/opentelemetry/proto/common/v1/common.proto create mode 100644 src/Braintrust.Sdk/Trace/Protos/opentelemetry/proto/resource/v1/resource.proto create mode 100644 src/Braintrust.Sdk/Trace/Protos/opentelemetry/proto/trace/v1/trace.proto create mode 100644 tests/Braintrust.Sdk.Tests/Trace/SpanCustomizerTest.cs diff --git a/src/Braintrust.Sdk/Braintrust.Sdk.csproj b/src/Braintrust.Sdk/Braintrust.Sdk.csproj index e5d238f..3bb7411 100644 --- a/src/Braintrust.Sdk/Braintrust.Sdk.csproj +++ b/src/Braintrust.Sdk/Braintrust.Sdk.csproj @@ -55,6 +55,8 @@ + + @@ -65,4 +67,12 @@ + + + + + diff --git a/src/Braintrust.Sdk/Config/BraintrustConfig.cs b/src/Braintrust.Sdk/Config/BraintrustConfig.cs index 8322cba..1d82e83 100644 --- a/src/Braintrust.Sdk/Config/BraintrustConfig.cs +++ b/src/Braintrust.Sdk/Config/BraintrustConfig.cs @@ -1,3 +1,5 @@ +using Braintrust.Sdk.Trace; + namespace Braintrust.Sdk.Config; public sealed record SpanOriginEnvironment(string Type, string? Name = null); @@ -29,12 +31,33 @@ public sealed class BraintrustConfig : BaseConfig public TimeSpan RequestTimeout { get; } public SpanOriginEnvironment? Environment { get; } + /// + /// Ordered registration snapshot. Hooks run synchronously on each detached outgoing span. + /// + public IReadOnlyList SpanCustomizers { get; } + public static BraintrustConfig FromEnvironment() { return Of(); } + public static BraintrustConfig FromEnvironment(IEnumerable spanCustomizers) + { + return Of(spanCustomizers); + } + public static BraintrustConfig Of(params (string Key, string? Value)[] envOverrides) + { + return Of(Array.Empty(), envOverrides); + } + + /// + /// Creates a configuration with an immutable copy of the ordered customizer registrations. + /// Customizer instances themselves are not cloned and must be safe for export threads. + /// + public static BraintrustConfig Of( + IEnumerable spanCustomizers, + params (string Key, string? Value)[] envOverrides) { var overridesMap = new Dictionary(); @@ -43,11 +66,21 @@ public static BraintrustConfig Of(params (string Key, string? Value)[] envOverri overridesMap[key] = value; } - return new BraintrustConfig(overridesMap); + return new BraintrustConfig(overridesMap, spanCustomizers); } - private BraintrustConfig(IDictionary envOverrides) : base(envOverrides) + private BraintrustConfig( + IDictionary envOverrides, + IEnumerable spanCustomizers) : base(envOverrides) { + ArgumentNullException.ThrowIfNull(spanCustomizers); + var customizers = spanCustomizers.ToArray(); + if (customizers.Any(customizer => customizer is null)) + { + throw new ArgumentException("Span customizers cannot contain null.", nameof(spanCustomizers)); + } + SpanCustomizers = Array.AsReadOnly(customizers); + try { _braintrustEnvSearchRoot = Directory.GetCurrentDirectory(); diff --git a/src/Braintrust.Sdk/Trace/BraintrustSpanCustomizerHandler.cs b/src/Braintrust.Sdk/Trace/BraintrustSpanCustomizerHandler.cs new file mode 100644 index 0000000..dde6b43 --- /dev/null +++ b/src/Braintrust.Sdk/Trace/BraintrustSpanCustomizerHandler.cs @@ -0,0 +1,159 @@ +using Braintrust.Sdk.Trace.Protos.Collector.Trace.V1; +using Braintrust.Sdk.Trace.Protos.Common.V1; +using Braintrust.Sdk.Trace.Protos.Resource.V1; +using Braintrust.Sdk.Trace.Protos.Trace.V1; +using OtlpSpan = Braintrust.Sdk.Trace.Protos.Trace.V1.Span; +using Google.Protobuf; + +namespace Braintrust.Sdk.Trace; + +/// +/// Adapts the upstream Activity-only exporter to detached per-span OTLP JSON data. +/// Hooks run before final serialization and before auth or network transmission. +/// +internal sealed class BraintrustSpanCustomizerHandler : DelegatingHandler +{ + private readonly IReadOnlyList _customizers; + + internal BraintrustSpanCustomizerHandler(IReadOnlyList customizers) + { + _customizers = customizers; + } + + protected override HttpResponseMessage Send(HttpRequestMessage request, CancellationToken cancellationToken) + { + CustomizeAsync(request, cancellationToken).GetAwaiter().GetResult(); + return base.Send(request, cancellationToken); + } + + protected override async Task SendAsync( + HttpRequestMessage request, CancellationToken cancellationToken) + { + await CustomizeAsync(request, cancellationToken).ConfigureAwait(false); + return await base.SendAsync(request, cancellationToken).ConfigureAwait(false); + } + + private async Task CustomizeAsync(HttpRequestMessage request, CancellationToken cancellationToken) + { + try + { + var original = request.Content ?? throw new InvalidOperationException("Missing OTLP request content."); + if (original.Headers.ContentType?.MediaType != "application/x-protobuf" || + original.Headers.ContentEncoding.Count != 0) + { + throw new InvalidOperationException("Span customization requires uncompressed OTLP protobuf content."); + } + + var bytes = await original.ReadAsByteArrayAsync(cancellationToken).ConfigureAwait(false); + var batch = ExportTraceServiceRequest.Parser.ParseFrom(bytes); + batch = CustomizeBatch(batch); + + var replacement = new ByteArrayContent(batch.ToByteArray()); + try + { + foreach (var header in original.Headers) + { + if (!header.Key.Equals("Content-Length", StringComparison.OrdinalIgnoreCase)) + { + replacement.Headers.TryAddWithoutValidation(header.Key, header.Value); + } + } + request.Content = replacement; + } + catch + { + replacement.Dispose(); + throw; + } + original.Dispose(); + } + catch (Exception ex) + { + // Never log the snapshot or exception message, which may contain secrets. + Console.Error.WriteLine($"[Braintrust] Span customization failed; no spans sent ({ex.GetType().Name})."); + // The OTel exporter catches this and reports ExportResult.Failure. Do not + // use HttpRequestException: customization failures are not network retries. + throw new InvalidOperationException("Braintrust span customization failed; the batch was not sent.", ex); + } + } + + private ExportTraceServiceRequest CustomizeBatch(ExportTraceServiceRequest batch) + { + OtlpJsonPayload.ValidateKnownFields(batch); + var result = new ExportTraceServiceRequest(); + var groups = new Dictionary<(Resource? Resource, string SchemaUrl), + (ResourceSpans Resource, Dictionary<(InstrumentationScope? Scope, string SchemaUrl), ScopeSpans> Scopes)>(); + + foreach (var resource in batch.ResourceSpans) + { + if (resource.ScopeSpans.Count == 0) + { + Add(resource); + } + foreach (var scope in resource.ScopeSpans) + { + if (scope.Spans.Count == 0) + { + Add(new ResourceSpans + { + Resource = resource.Resource, + SchemaUrl = resource.SchemaUrl, + ScopeSpans = { scope } + }); + } + foreach (var span in scope.Spans) + { + var data = OtlpJsonPayload.ToSpanJson(resource, scope, span); + for (var i = 0; i < _customizers.Count; i++) + { + data = _customizers[i].OnSpanExport(data) + ?? throw new InvalidOperationException("A span customizer returned null."); + var customized = OtlpJsonPayload.FromSpanJson(data); + ValidateIdentity(customized.ScopeSpans[0].Spans[0], span); + if (i == _customizers.Count - 1) + { + // Snapshot the final result before invoking hooks on another span. + // Retained JsonObjects cannot mutate a previously customized span. + Add(customized); + } + } + } + } + } + return result; + + void Add(ResourceSpans data) + { + // Group only after customization. Metadata objects are private protobuf + // snapshots, never mutated after becoming dictionary keys. + var key = (data.Resource, data.SchemaUrl); + if (!groups.TryGetValue(key, out var group)) + { + group = (new ResourceSpans { Resource = data.Resource, SchemaUrl = data.SchemaUrl }, new()); + groups.Add(key, group); + result.ResourceSpans.Add(group.Resource); + } + foreach (var incomingScope in data.ScopeSpans) + { + var scopeKey = (incomingScope.Scope, incomingScope.SchemaUrl); + if (!group.Scopes.TryGetValue(scopeKey, out var scope)) + { + scope = new ScopeSpans { Scope = incomingScope.Scope, SchemaUrl = incomingScope.SchemaUrl }; + group.Scopes.Add(scopeKey, scope); + group.Resource.ScopeSpans.Add(scope); + } + scope.Spans.Add(incomingScope.Spans); + } + } + } + + private static void ValidateIdentity(OtlpSpan span, OtlpSpan original) + { + if (!span.TraceId.Equals(original.TraceId) || + !span.SpanId.Equals(original.SpanId) || + !span.ParentSpanId.Equals(original.ParentSpanId)) + { + throw new InvalidOperationException("A span customizer changed the trace, span or parent span ID."); + } + } +} diff --git a/src/Braintrust.Sdk/Trace/BraintrustTracing.cs b/src/Braintrust.Sdk/Trace/BraintrustTracing.cs index 739679e..4f384a6 100644 --- a/src/Braintrust.Sdk/Trace/BraintrustTracing.cs +++ b/src/Braintrust.Sdk/Trace/BraintrustTracing.cs @@ -83,17 +83,27 @@ public static void Enable(BraintrustConfig config, TracerProviderBuilder tracerP otlpOptions.Protocol = OtlpExportProtocol.HttpProtobuf; otlpOptions.Endpoint = new Uri($"{config.ApiUrl}{config.TracesPath}"); otlpOptions.TimeoutMilliseconds = (int)config.RequestTimeout.TotalMilliseconds; - otlpOptions.HttpClientFactory = () => new HttpClient(new BraintrustOtlpAuthHandler(config) - { - InnerHandler = new HttpClientHandler() - }) - { - Timeout = config.RequestTimeout - }; + otlpOptions.HttpClientFactory = () => CreateHttpClient(config); }) .SetSampler(new AlwaysOnSampler()); } + internal static HttpClient CreateHttpClient(BraintrustConfig config, HttpMessageHandler? transport = null) + { + HttpMessageHandler handler = new BraintrustOtlpAuthHandler(config) + { + InnerHandler = transport ?? new HttpClientHandler() + }; + if (config.SpanCustomizers.Count != 0) + { + handler = new BraintrustSpanCustomizerHandler(config.SpanCustomizers) + { + InnerHandler = handler + }; + } + return new HttpClient(handler) { Timeout = config.RequestTimeout }; + } + /// /// Flush all pending spans to Braintrust. /// Returns true if the flush completed within the timeout, false otherwise. diff --git a/src/Braintrust.Sdk/Trace/ISpanCustomizer.cs b/src/Braintrust.Sdk/Trace/ISpanCustomizer.cs new file mode 100644 index 0000000..2dca38d --- /dev/null +++ b/src/Braintrust.Sdk/Trace/ISpanCustomizer.cs @@ -0,0 +1,21 @@ +using System.Text.Json.Nodes; + +namespace Braintrust.Sdk.Trace; + +/// +/// Transforms detached outgoing span data, never the application's Activity. +/// Hooks run synchronously in registration order for each completed span. +/// +public interface ISpanCustomizer +{ + /// + /// Returns the supplied span or a replacement, for example from . + /// The object contains OTLP span fields, plus optional resource, scope, resourceSchemaUrl + /// and scopeSchemaUrl metadata. Metadata changes apply only to this span. + /// traceId, spanId and parentSpanId must remain unchanged. IDs are hexadecimal strings, + /// enums are integers and 64-bit integers are decimal strings. Null, exceptions, invalid + /// data or changed IDs fail the entire batch before transmission. Do not retain and + /// mutate the object after returning. + /// + JsonObject OnSpanExport(JsonObject span) => span; +} diff --git a/src/Braintrust.Sdk/Trace/OtlpJsonPayload.cs b/src/Braintrust.Sdk/Trace/OtlpJsonPayload.cs new file mode 100644 index 0000000..c0d7ce7 --- /dev/null +++ b/src/Braintrust.Sdk/Trace/OtlpJsonPayload.cs @@ -0,0 +1,167 @@ +using System.Text.Json; +using System.Text.Json.Nodes; +using Braintrust.Sdk.Trace.Protos.Collector.Trace.V1; +using Braintrust.Sdk.Trace.Protos.Trace.V1; +using Google.Protobuf; + +namespace Braintrust.Sdk.Trace; + +/// +/// Bridges protobuf JSON and OTLP JSON, whose span and link IDs use hex rather +/// than base64. Other bytes remain base64 and 64-bit integers remain strings. +/// +internal static class OtlpJsonPayload +{ + private static readonly JsonFormatter Formatter = new( + JsonFormatter.Settings.Default.WithFormatEnumsAsIntegers(true)); + private static readonly JsonParser Parser = new( + JsonParser.Settings.Default.WithIgnoreUnknownFields(false)); + + internal static void ValidateKnownFields(ExportTraceServiceRequest batch) + { + var json = Formatter.Format(batch); + if (!batch.Equals(Parser.Parse(json))) + { + throw new InvalidOperationException( + "The OTLP request contains fields that cannot be represented losslessly as JSON."); + } + } + + internal static JsonObject ToSpanJson(ResourceSpans resource, ScopeSpans scope, Span span) + { + var payload = JsonNode.Parse(Formatter.Format(span))!.AsObject(); + ConvertIds(payload, toHex: true); + + if (resource.Resource is not null) + { + payload["resource"] = JsonNode.Parse(Formatter.Format(resource.Resource)); + } + if (scope.Scope is not null) + { + payload["scope"] = JsonNode.Parse(Formatter.Format(scope.Scope)); + } + if (resource.SchemaUrl.Length != 0) + { + payload["resourceSchemaUrl"] = resource.SchemaUrl; + } + if (scope.SchemaUrl.Length != 0) + { + payload["scopeSchemaUrl"] = scope.SchemaUrl; + } + + return payload; + } + + internal static ResourceSpans FromSpanJson(JsonObject data) + { + // The hook owns its tree: neither ID conversion nor regrouping may alter it. + var span = (JsonObject)data.DeepClone(); + var resource = span["resource"]; + var scope = span["scope"]; + var resourceSchemaUrl = span["resourceSchemaUrl"]; + var scopeSchemaUrl = span["scopeSchemaUrl"]; + span.Remove("resource"); + span.Remove("scope"); + span.Remove("resourceSchemaUrl"); + span.Remove("scopeSchemaUrl"); + ConvertIds(span, toHex: false); + + // Parse metadata and span fields together so the protobuf parser enforces + // their exact schemas, including unknown fields and explicit null values. + var payload = new JsonObject + { + ["resource"] = resource, + ["schemaUrl"] = resourceSchemaUrl, + ["scopeSpans"] = new JsonArray(new JsonObject + { + ["scope"] = scope, + ["schemaUrl"] = scopeSchemaUrl, + ["spans"] = new JsonArray(span), + }), + }; + return Parser.Parse(payload.ToJsonString()); + } + + private static void ConvertIds(JsonObject span, bool toHex) + { + ConvertId(span, "traceId", "trace_id", 16, toHex); + ConvertId(span, "spanId", "span_id", 8, toHex); + ConvertId(span, "parentSpanId", "parent_span_id", 8, toHex, optional: true); + foreach (var link in Objects(span, "links")) + { + ConvertId(link, "traceId", "trace_id", 16, toHex); + ConvertId(link, "spanId", "span_id", 8, toHex); + } + } + + private static IEnumerable Objects(JsonObject parent, string name) + { + var node = parent[name]; + if (node is null) + { + yield break; + } + + if (node is not JsonArray array) + { + throw new JsonException($"OTLP field '{name}' must be an array."); + } + + foreach (var item in array) + { + yield return item as JsonObject + ?? throw new JsonException($"OTLP field '{name}' must contain objects."); + } + } + + private static void ConvertId( + JsonObject parent, string name, string protoName, int byteLength, bool toHex, bool optional = false) + { + name = FieldName(parent, name, protoName); + var node = parent[name]; + if (optional && node is null) + { + return; + } + + if (node is not JsonValue value || !value.TryGetValue(out var text)) + { + throw new JsonException($"OTLP field '{name}' must be a string."); + } + + if (optional && text.Length == 0) + { + return; + } + + if (!toHex && text.Length != byteLength * 2) + { + throw new JsonException($"OTLP field '{name}' must contain {byteLength * 2} hexadecimal characters."); + } + + var bytes = toHex ? Convert.FromBase64String(text) : Convert.FromHexString(text); + if (bytes.Length != byteLength) + { + throw new JsonException($"OTLP field '{name}' must contain exactly {byteLength} bytes."); + } + + parent[name] = toHex ? Convert.ToHexString(bytes).ToLowerInvariant() : Convert.ToBase64String(bytes); + } + + private static string FieldName(JsonObject parent, string name, string protoName) + { + // JsonParser accepts both names. Recognize aliases so they cannot bypass + // ID conversion, but reject duplicate spellings of the same field. + if (!parent.ContainsKey(protoName)) + { + return name; + } + + if (parent.ContainsKey(name)) + { + throw new JsonException($"OTLP field '{name}' is specified more than once."); + } + + return protoName; + } +} diff --git a/src/Braintrust.Sdk/Trace/Protos/opentelemetry/proto/collector/trace/v1/trace_service.proto b/src/Braintrust.Sdk/Trace/Protos/opentelemetry/proto/collector/trace/v1/trace_service.proto new file mode 100644 index 0000000..f2cf8bb --- /dev/null +++ b/src/Braintrust.Sdk/Trace/Protos/opentelemetry/proto/collector/trace/v1/trace_service.proto @@ -0,0 +1,79 @@ +// Copyright 2019, OpenTelemetry Authors +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +syntax = "proto3"; + +package opentelemetry.proto.collector.trace.v1; + +import "opentelemetry/proto/trace/v1/trace.proto"; + +option csharp_namespace = "Braintrust.Sdk.Trace.Protos.Collector.Trace.V1"; +option java_multiple_files = true; +option java_package = "io.opentelemetry.proto.collector.trace.v1"; +option java_outer_classname = "TraceServiceProto"; +option go_package = "go.opentelemetry.io/proto/otlp/collector/trace/v1"; + +// Service that can be used to push spans between one Application instrumented with +// OpenTelemetry and a collector, or between a collector and a central collector (in this +// case spans are sent/received to/from multiple Applications). +service TraceService { + // For performance reasons, it is recommended to keep this RPC + // alive for the entire life of the application. + rpc Export(ExportTraceServiceRequest) returns (ExportTraceServiceResponse) {} +} + +message ExportTraceServiceRequest { + // An array of ResourceSpans. + // For data coming from a single resource this array will typically contain one + // element. Intermediary nodes (such as OpenTelemetry Collector) that receive + // data from multiple origins typically batch the data before forwarding further and + // in that case this array will contain multiple elements. + repeated opentelemetry.proto.trace.v1.ResourceSpans resource_spans = 1; +} + +message ExportTraceServiceResponse { + // The details of a partially successful export request. + // + // If the request is only partially accepted + // (i.e. when the server accepts only parts of the data and rejects the rest) + // the server MUST initialize the `partial_success` field and MUST + // set the `rejected_` with the number of items it rejected. + // + // Servers MAY also make use of the `partial_success` field to convey + // warnings/suggestions to senders even when the request was fully accepted. + // In such cases, the `rejected_` MUST have a value of `0` and + // the `error_message` MUST be non-empty. + // + // A `partial_success` message with an empty value (rejected_ = 0 and + // `error_message` = "") is equivalent to it not being set/present. Senders + // SHOULD interpret it the same way as in the full success case. + ExportTracePartialSuccess partial_success = 1; +} + +message ExportTracePartialSuccess { + // The number of rejected spans. + // + // A `rejected_` field holding a `0` value indicates that the + // request was fully accepted. + int64 rejected_spans = 1; + + // A developer-facing human-readable message in English. It should be used + // either to explain why the server rejected parts of the data during a partial + // success or to convey warnings/suggestions during a full success. The message + // should offer guidance on how users can address such issues. + // + // error_message is an optional field. An error_message with an empty value + // is equivalent to it not being set. + string error_message = 2; +} diff --git a/src/Braintrust.Sdk/Trace/Protos/opentelemetry/proto/common/v1/common.proto b/src/Braintrust.Sdk/Trace/Protos/opentelemetry/proto/common/v1/common.proto new file mode 100644 index 0000000..f9f7307 --- /dev/null +++ b/src/Braintrust.Sdk/Trace/Protos/opentelemetry/proto/common/v1/common.proto @@ -0,0 +1,81 @@ +// Copyright 2019, OpenTelemetry Authors +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +syntax = "proto3"; + +package opentelemetry.proto.common.v1; + +option csharp_namespace = "Braintrust.Sdk.Trace.Protos.Common.V1"; +option java_multiple_files = true; +option java_package = "io.opentelemetry.proto.common.v1"; +option java_outer_classname = "CommonProto"; +option go_package = "go.opentelemetry.io/proto/otlp/common/v1"; + +// AnyValue is used to represent any type of attribute value. AnyValue may contain a +// primitive value such as a string or integer or it may contain an arbitrary nested +// object containing arrays, key-value lists and primitives. +message AnyValue { + // The value is one of the listed fields. It is valid for all values to be unspecified + // in which case this AnyValue is considered to be "empty". + oneof value { + string string_value = 1; + bool bool_value = 2; + int64 int_value = 3; + double double_value = 4; + ArrayValue array_value = 5; + KeyValueList kvlist_value = 6; + bytes bytes_value = 7; + } +} + +// ArrayValue is a list of AnyValue messages. We need ArrayValue as a message +// since oneof in AnyValue does not allow repeated fields. +message ArrayValue { + // Array of values. The array may be empty (contain 0 elements). + repeated AnyValue values = 1; +} + +// KeyValueList is a list of KeyValue messages. We need KeyValueList as a message +// since `oneof` in AnyValue does not allow repeated fields. Everywhere else where we need +// a list of KeyValue messages (e.g. in Span) we use `repeated KeyValue` directly to +// avoid unnecessary extra wrapping (which slows down the protocol). The 2 approaches +// are semantically equivalent. +message KeyValueList { + // A collection of key/value pairs of key-value pairs. The list may be empty (may + // contain 0 elements). + // The keys MUST be unique (it is not allowed to have more than one + // value with the same key). + repeated KeyValue values = 1; +} + +// KeyValue is a key-value pair that is used to store Span attributes, Link +// attributes, etc. +message KeyValue { + string key = 1; + AnyValue value = 2; +} + +// InstrumentationScope is a message representing the instrumentation scope information +// such as the fully qualified name and version. +message InstrumentationScope { + // An empty instrumentation scope name means the name is unknown. + string name = 1; + string version = 2; + + // Additional attributes that describe the scope. [Optional]. + // Attribute keys MUST be unique (it is not allowed to have more than one + // attribute with the same key). + repeated KeyValue attributes = 3; + uint32 dropped_attributes_count = 4; +} diff --git a/src/Braintrust.Sdk/Trace/Protos/opentelemetry/proto/resource/v1/resource.proto b/src/Braintrust.Sdk/Trace/Protos/opentelemetry/proto/resource/v1/resource.proto new file mode 100644 index 0000000..83f744b --- /dev/null +++ b/src/Braintrust.Sdk/Trace/Protos/opentelemetry/proto/resource/v1/resource.proto @@ -0,0 +1,37 @@ +// Copyright 2019, OpenTelemetry Authors +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +syntax = "proto3"; + +package opentelemetry.proto.resource.v1; + +import "opentelemetry/proto/common/v1/common.proto"; + +option csharp_namespace = "Braintrust.Sdk.Trace.Protos.Resource.V1"; +option java_multiple_files = true; +option java_package = "io.opentelemetry.proto.resource.v1"; +option java_outer_classname = "ResourceProto"; +option go_package = "go.opentelemetry.io/proto/otlp/resource/v1"; + +// Resource information. +message Resource { + // Set of attributes that describe the resource. + // Attribute keys MUST be unique (it is not allowed to have more than one + // attribute with the same key). + repeated opentelemetry.proto.common.v1.KeyValue attributes = 1; + + // dropped_attributes_count is the number of dropped attributes. If the value is 0, then + // no attributes were dropped. + uint32 dropped_attributes_count = 2; +} diff --git a/src/Braintrust.Sdk/Trace/Protos/opentelemetry/proto/trace/v1/trace.proto b/src/Braintrust.Sdk/Trace/Protos/opentelemetry/proto/trace/v1/trace.proto new file mode 100644 index 0000000..6a9fdca --- /dev/null +++ b/src/Braintrust.Sdk/Trace/Protos/opentelemetry/proto/trace/v1/trace.proto @@ -0,0 +1,357 @@ +// Copyright 2019, OpenTelemetry Authors +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +syntax = "proto3"; + +package opentelemetry.proto.trace.v1; + +import "opentelemetry/proto/common/v1/common.proto"; +import "opentelemetry/proto/resource/v1/resource.proto"; + +option csharp_namespace = "Braintrust.Sdk.Trace.Protos.Trace.V1"; +option java_multiple_files = true; +option java_package = "io.opentelemetry.proto.trace.v1"; +option java_outer_classname = "TraceProto"; +option go_package = "go.opentelemetry.io/proto/otlp/trace/v1"; + +// TracesData represents the traces data that can be stored in a persistent storage, +// OR can be embedded by other protocols that transfer OTLP traces data but do +// not implement the OTLP protocol. +// +// The main difference between this message and collector protocol is that +// in this message there will not be any "control" or "metadata" specific to +// OTLP protocol. +// +// When new fields are added into this message, the OTLP request MUST be updated +// as well. +message TracesData { + // An array of ResourceSpans. + // For data coming from a single resource this array will typically contain + // one element. Intermediary nodes that receive data from multiple origins + // typically batch the data before forwarding further and in that case this + // array will contain multiple elements. + repeated ResourceSpans resource_spans = 1; +} + +// A collection of ScopeSpans from a Resource. +message ResourceSpans { + reserved 1000; + + // The resource for the spans in this message. + // If this field is not set then no resource info is known. + opentelemetry.proto.resource.v1.Resource resource = 1; + + // A list of ScopeSpans that originate from a resource. + repeated ScopeSpans scope_spans = 2; + + // The Schema URL, if known. This is the identifier of the Schema that the resource data + // is recorded in. Notably, the last part of the URL path is the version number of the + // schema: http[s]://server[:port]/path/. To learn more about Schema URL see + // https://opentelemetry.io/docs/specs/otel/schemas/#schema-url + // This schema_url applies to the data in the "resource" field. It does not apply + // to the data in the "scope_spans" field which have their own schema_url field. + string schema_url = 3; +} + +// A collection of Spans produced by an InstrumentationScope. +message ScopeSpans { + // The instrumentation scope information for the spans in this message. + // Semantically when InstrumentationScope isn't set, it is equivalent with + // an empty instrumentation scope name (unknown). + opentelemetry.proto.common.v1.InstrumentationScope scope = 1; + + // A list of Spans that originate from an instrumentation scope. + repeated Span spans = 2; + + // The Schema URL, if known. This is the identifier of the Schema that the span data + // is recorded in. Notably, the last part of the URL path is the version number of the + // schema: http[s]://server[:port]/path/. To learn more about Schema URL see + // https://opentelemetry.io/docs/specs/otel/schemas/#schema-url + // This schema_url applies to all spans and span events in the "spans" field. + string schema_url = 3; +} + +// A Span represents a single operation performed by a single component of the system. +// +// The next available field id is 17. +message Span { + // A unique identifier for a trace. All spans from the same trace share + // the same `trace_id`. The ID is a 16-byte array. An ID with all zeroes OR + // of length other than 16 bytes is considered invalid (empty string in OTLP/JSON + // is zero-length and thus is also invalid). + // + // This field is required. + bytes trace_id = 1; + + // A unique identifier for a span within a trace, assigned when the span + // is created. The ID is an 8-byte array. An ID with all zeroes OR of length + // other than 8 bytes is considered invalid (empty string in OTLP/JSON + // is zero-length and thus is also invalid). + // + // This field is required. + bytes span_id = 2; + + // trace_state conveys information about request position in multiple distributed tracing graphs. + // It is a trace_state in w3c-trace-context format: https://www.w3.org/TR/trace-context/#tracestate-header + // See also https://github.com/w3c/distributed-tracing for more details about this field. + string trace_state = 3; + + // The `span_id` of this span's parent span. If this is a root span, then this + // field must be empty. The ID is an 8-byte array. + bytes parent_span_id = 4; + + // Flags, a bit field. + // + // Bits 0-7 (8 least significant bits) are the trace flags as defined in W3C Trace + // Context specification. To read the 8-bit W3C trace flag, use + // `flags & SPAN_FLAGS_TRACE_FLAGS_MASK`. + // + // See https://www.w3.org/TR/trace-context-2/#trace-flags for the flag definitions. + // + // Bits 8 and 9 represent the 3 states of whether a span's parent + // is remote. The states are (unknown, is not remote, is remote). + // To read whether the value is known, use `(flags & SPAN_FLAGS_CONTEXT_HAS_IS_REMOTE_MASK) != 0`. + // To read whether the span is remote, use `(flags & SPAN_FLAGS_CONTEXT_IS_REMOTE_MASK) != 0`. + // + // When creating span messages, if the message is logically forwarded from another source + // with an equivalent flags fields (i.e., usually another OTLP span message), the field SHOULD + // be copied as-is. If creating from a source that does not have an equivalent flags field + // (such as a runtime representation of an OpenTelemetry span), the high 22 bits MUST + // be set to zero. + // Readers MUST NOT assume that bits 10-31 (22 most significant bits) will be zero. + // + // [Optional]. + fixed32 flags = 16; + + // A description of the span's operation. + // + // For example, the name can be a qualified method name or a file name + // and a line number where the operation is called. A best practice is to use + // the same display name at the same call point in an application. + // This makes it easier to correlate spans in different traces. + // + // This field is semantically required to be set to non-empty string. + // Empty value is equivalent to an unknown span name. + // + // This field is required. + string name = 5; + + // SpanKind is the type of span. Can be used to specify additional relationships between spans + // in addition to a parent/child relationship. + enum SpanKind { + // Unspecified. Do NOT use as default. + // Implementations MAY assume SpanKind to be INTERNAL when receiving UNSPECIFIED. + SPAN_KIND_UNSPECIFIED = 0; + + // Indicates that the span represents an internal operation within an application, + // as opposed to an operation happening at the boundaries. Default value. + SPAN_KIND_INTERNAL = 1; + + // Indicates that the span covers server-side handling of an RPC or other + // remote network request. + SPAN_KIND_SERVER = 2; + + // Indicates that the span describes a request to some remote service. + SPAN_KIND_CLIENT = 3; + + // Indicates that the span describes a producer sending a message to a broker. + // Unlike CLIENT and SERVER, there is often no direct critical path latency relationship + // between producer and consumer spans. A PRODUCER span ends when the message was accepted + // by the broker while the logical processing of the message might span a much longer time. + SPAN_KIND_PRODUCER = 4; + + // Indicates that the span describes consumer receiving a message from a broker. + // Like the PRODUCER kind, there is often no direct critical path latency relationship + // between producer and consumer spans. + SPAN_KIND_CONSUMER = 5; + } + + // Distinguishes between spans generated in a particular context. For example, + // two spans with the same name may be distinguished using `CLIENT` (caller) + // and `SERVER` (callee) to identify queueing latency associated with the span. + SpanKind kind = 6; + + // start_time_unix_nano is the start time of the span. On the client side, this is the time + // kept by the local machine where the span execution starts. On the server side, this + // is the time when the server's application handler starts running. + // Value is UNIX Epoch time in nanoseconds since 00:00:00 UTC on 1 January 1970. + // + // This field is semantically required and it is expected that end_time >= start_time. + fixed64 start_time_unix_nano = 7; + + // end_time_unix_nano is the end time of the span. On the client side, this is the time + // kept by the local machine where the span execution ends. On the server side, this + // is the time when the server application handler stops running. + // Value is UNIX Epoch time in nanoseconds since 00:00:00 UTC on 1 January 1970. + // + // This field is semantically required and it is expected that end_time >= start_time. + fixed64 end_time_unix_nano = 8; + + // attributes is a collection of key/value pairs. Note, global attributes + // like server name can be set using the resource API. Examples of attributes: + // + // "/http/user_agent": "Mozilla/5.0 (Macintosh; Intel Mac OS X 10_14_2) AppleWebKit/537.36 (KHTML, like Gecko) Chrome/71.0.3578.98 Safari/537.36" + // "/http/server_latency": 300 + // "example.com/myattribute": true + // "example.com/score": 10.239 + // + // The OpenTelemetry API specification further restricts the allowed value types: + // https://github.com/open-telemetry/opentelemetry-specification/blob/main/specification/common/README.md#attribute + // Attribute keys MUST be unique (it is not allowed to have more than one + // attribute with the same key). + repeated opentelemetry.proto.common.v1.KeyValue attributes = 9; + + // dropped_attributes_count is the number of attributes that were discarded. Attributes + // can be discarded because their keys are too long or because there are too many + // attributes. If this value is 0, then no attributes were dropped. + uint32 dropped_attributes_count = 10; + + // Event is a time-stamped annotation of the span, consisting of user-supplied + // text description and key-value pairs. + message Event { + // time_unix_nano is the time the event occurred. + fixed64 time_unix_nano = 1; + + // name of the event. + // This field is semantically required to be set to non-empty string. + string name = 2; + + // attributes is a collection of attribute key/value pairs on the event. + // Attribute keys MUST be unique (it is not allowed to have more than one + // attribute with the same key). + repeated opentelemetry.proto.common.v1.KeyValue attributes = 3; + + // dropped_attributes_count is the number of dropped attributes. If the value is 0, + // then no attributes were dropped. + uint32 dropped_attributes_count = 4; + } + + // events is a collection of Event items. + repeated Event events = 11; + + // dropped_events_count is the number of dropped events. If the value is 0, then no + // events were dropped. + uint32 dropped_events_count = 12; + + // A pointer from the current span to another span in the same trace or in a + // different trace. For example, this can be used in batching operations, + // where a single batch handler processes multiple requests from different + // traces or when the handler receives a request from a different project. + message Link { + // A unique identifier of a trace that this linked span is part of. The ID is a + // 16-byte array. + bytes trace_id = 1; + + // A unique identifier for the linked span. The ID is an 8-byte array. + bytes span_id = 2; + + // The trace_state associated with the link. + string trace_state = 3; + + // attributes is a collection of attribute key/value pairs on the link. + // Attribute keys MUST be unique (it is not allowed to have more than one + // attribute with the same key). + repeated opentelemetry.proto.common.v1.KeyValue attributes = 4; + + // dropped_attributes_count is the number of dropped attributes. If the value is 0, + // then no attributes were dropped. + uint32 dropped_attributes_count = 5; + + // Flags, a bit field. + // + // Bits 0-7 (8 least significant bits) are the trace flags as defined in W3C Trace + // Context specification. To read the 8-bit W3C trace flag, use + // `flags & SPAN_FLAGS_TRACE_FLAGS_MASK`. + // + // See https://www.w3.org/TR/trace-context-2/#trace-flags for the flag definitions. + // + // Bits 8 and 9 represent the 3 states of whether the link is remote. + // The states are (unknown, is not remote, is remote). + // To read whether the value is known, use `(flags & SPAN_FLAGS_CONTEXT_HAS_IS_REMOTE_MASK) != 0`. + // To read whether the link is remote, use `(flags & SPAN_FLAGS_CONTEXT_IS_REMOTE_MASK) != 0`. + // + // Readers MUST NOT assume that bits 10-31 (22 most significant bits) will be zero. + // When creating new spans, bits 10-31 (most-significant 22-bits) MUST be zero. + // + // [Optional]. + fixed32 flags = 6; + } + + // links is a collection of Links, which are references from this span to a span + // in the same or different trace. + repeated Link links = 13; + + // dropped_links_count is the number of dropped links after the maximum size was + // enforced. If this value is 0, then no links were dropped. + uint32 dropped_links_count = 14; + + // An optional final status for this span. Semantically when Status isn't set, it means + // span's status code is unset, i.e. assume STATUS_CODE_UNSET (code = 0). + Status status = 15; +} + +// The Status type defines a logical error model that is suitable for different +// programming environments, including REST APIs and RPC APIs. +message Status { + reserved 1; + + // A developer-facing human readable error message. + string message = 2; + + // For the semantics of status codes see + // https://github.com/open-telemetry/opentelemetry-specification/blob/main/specification/trace/api.md#set-status + enum StatusCode { + // The default status. + STATUS_CODE_UNSET = 0; + // The Span has been validated by an Application developer or Operator to + // have completed successfully. + STATUS_CODE_OK = 1; + // The Span contains an error. + STATUS_CODE_ERROR = 2; + }; + + // The status code. + StatusCode code = 3; +} + +// SpanFlags represents constants used to interpret the +// Span.flags field, which is protobuf 'fixed32' type and is to +// be used as bit-fields. Each non-zero value defined in this enum is +// a bit-mask. To extract the bit-field, for example, use an +// expression like: +// +// (span.flags & SPAN_FLAGS_TRACE_FLAGS_MASK) +// +// See https://www.w3.org/TR/trace-context-2/#trace-flags for the flag definitions. +// +// Note that Span flags were introduced in version 1.1 of the +// OpenTelemetry protocol. Older Span producers do not set this +// field, consequently consumers should not rely on the absence of a +// particular flag bit to indicate the presence of a particular feature. +enum SpanFlags { + // The zero value for the enum. Should not be used for comparisons. + // Instead use bitwise "and" with the appropriate mask as shown above. + SPAN_FLAGS_DO_NOT_USE = 0; + + // Bits 0-7 are used for trace flags. + SPAN_FLAGS_TRACE_FLAGS_MASK = 0x000000FF; + + // Bits 8 and 9 are used to indicate that the parent span or link span is remote. + // Bit 8 (`HAS_IS_REMOTE`) indicates whether the value is known. + // Bit 9 (`IS_REMOTE`) indicates whether the span or link is remote. + SPAN_FLAGS_CONTEXT_HAS_IS_REMOTE_MASK = 0x00000100; + SPAN_FLAGS_CONTEXT_IS_REMOTE_MASK = 0x00000200; + + // Bits 10-31 are reserved for future use. +} diff --git a/tests/Braintrust.Sdk.Tests/Trace/SpanCustomizerTest.cs b/tests/Braintrust.Sdk.Tests/Trace/SpanCustomizerTest.cs new file mode 100644 index 0000000..e785c38 --- /dev/null +++ b/tests/Braintrust.Sdk.Tests/Trace/SpanCustomizerTest.cs @@ -0,0 +1,764 @@ +using System.Diagnostics; +using System.Net; +using System.Net.Http.Headers; +using System.Text.Json.Nodes; +using Braintrust.Sdk.Config; +using Braintrust.Sdk.Trace; +using Braintrust.Sdk.Trace.Protos.Collector.Trace.V1; +using Braintrust.Sdk.Trace.Protos.Common.V1; +using Braintrust.Sdk.Trace.Protos.Trace.V1; +using Google.Protobuf; +using OpenTelemetry; +using OpenTelemetry.Exporter; +using OpenTelemetry.Resources; +using OpenTelemetry.Trace; +using OtlpSpan = Braintrust.Sdk.Trace.Protos.Trace.V1.Span; +using Resource = Braintrust.Sdk.Trace.Protos.Resource.V1.Resource; +using Status = Braintrust.Sdk.Trace.Protos.Trace.V1.Status; + +namespace Braintrust.Sdk.Tests.Trace; + +public class SpanCustomizerTest : IDisposable +{ + private readonly ActivitySource _source = new($"secret-source-{Guid.NewGuid()}", "secret-version"); + private readonly ActivityListener _listener; + + public SpanCustomizerTest() + { + _listener = new ActivityListener + { + ShouldListenTo = source => source == _source, + Sample = (ref ActivityCreationOptions _) => ActivitySamplingResult.AllDataAndRecorded + }; + ActivitySource.AddActivityListener(_listener); + } + + [Fact] + public void Export_ScrubsEachSpanInRegistrationOrderWithoutMutatingActivityOrOtherExports() + { + var calls = new List(); + var originals = new List<(JsonObject Payload, string Json)>(); + var results = new List<(JsonObject Payload, string Json)>(); + JsonObject? replacement = null; + var registrations = new List + { + new OmittedHook(), + new Customizer(payload => + { + calls.Add($"first({payload["name"]!.GetValue()})"); + Assert.False(payload.ContainsKey("resourceSpans")); + originals.Add((payload, payload.ToJsonString())); + replacement = payload.DeepClone().AsObject(); + Scrub(replacement); + return replacement; + }), + new Customizer(payload => + { + calls.Add($"second({payload["name"]!.GetValue()})"); + Assert.Same(replacement, payload); + Assert.StartsWith("safe-", payload["name"]!.GetValue()); + payload["name"] = payload["name"]!.GetValue() + "-final"; + results.Add((payload, payload.ToJsonString())); + return payload; + }) + }; + var config = Config(registrations); + registrations.Clear(); + registrations.Add(new Customizer(_ => throw new InvalidOperationException("must not be registered"))); + using var first = CompletedSpan("secret-first"); + using var second = CompletedSpan("secret-second"); + var baseline = new CaptureHandler(); + Assert.Equal(ExportResult.Success, Export(Config(), baseline, first, second)); + var transport = new CaptureHandler(); + + Assert.Equal(ExportResult.Success, Export(config, transport, first, second)); + + Assert.Equal(new[] { "first(secret-first)", "second(safe-first)", "first(secret-second)", "second(safe-second)" }, calls); + Assert.All(originals, original => Assert.Equal(original.Json, original.Payload.ToJsonString())); + Assert.All(results, result => Assert.Equal(result.Json, result.Payload.ToJsonString())); + var originalBatch = Assert.Single(baseline.Batches); + var batch = Assert.Single(transport.Batches); + var resource = Assert.Single(batch.ResourceSpans); + Assert.Equal("safe-service", resource.Resource.Attributes.Single(a => a.Key == "service.name").Value.StringValue); + Assert.Equal("safe-resource-value", resource.Resource.Attributes.Single(a => a.Key == "safe-resource-key").Value.StringValue); + var scope = Assert.Single(resource.ScopeSpans); + Assert.Equal(_source.Name.Replace("secret", "safe"), scope.Scope.Name); + Assert.Equal("safe-version", scope.Scope.Version); + var originalSpans = originalBatch.ResourceSpans[0].ScopeSpans[0].Spans; + for (var i = 0; i < scope.Spans.Count; i++) + { + var exported = scope.Spans[i]; + var original = originalSpans[i]; + Assert.Equal(original.Name.Replace("secret", "safe") + "-final", exported.Name); + Assert.Equal("safe-input", exported.Attributes.Single(a => a.Key == "safe-tag-key").Value.StringValue); + Assert.Equal("safe-event", exported.Events[0].Name); + Assert.Equal("safe-detail", exported.Events[0].Attributes.Single(a => a.Key == "safe-event-key").Value.StringValue); + Assert.Equal("safe-status", exported.Status.Message); + Assert.Equal("safe-link-value", exported.Links[0].Attributes.Single(a => a.Key == "safe-link-key").Value.StringValue); + Assert.Equal(original.TraceId, exported.TraceId); + Assert.Equal(original.SpanId, exported.SpanId); + Assert.Equal(original.ParentSpanId, exported.ParentSpanId); + Assert.Equal(original.StartTimeUnixNano, exported.StartTimeUnixNano); + Assert.Equal(original.EndTimeUnixNano, exported.EndTimeUnixNano); + Assert.Equal(original.Kind, exported.Kind); + Assert.Equal(original.Flags, exported.Flags); + Assert.Equal(original.TraceState, exported.TraceState); + Assert.Equal(original.Status.Code, exported.Status.Code); + Assert.Equal(original.Links[0].TraceId, exported.Links[0].TraceId); + Assert.Equal(original.Links[0].SpanId, exported.Links[0].SpanId); + } + Assert.Equal("Bearer test-key", transport.Authorization); + Assert.Equal("project_id:default-project", transport.Parent); + Assert.Equal("secret-first", first.DisplayName); + Assert.Equal("secret-input", first.GetTagItem("secret-tag-key")); + Assert.Equal("secret-event", first.Events.Single().Name); + Assert.Equal("secret-detail", first.Events.Single().Tags.Single().Value); + Assert.Equal("secret-status", first.StatusDescription); + Assert.Equal("secret-link-value", first.Links.Single().Tags!.Single().Value); + var laterBaseline = new CaptureHandler(); + Assert.Equal(ExportResult.Success, Export(Config(), laterBaseline, first, second)); + Assert.Equal(originalBatch, Assert.Single(laterBaseline.Batches)); + } + + [Fact] + public async Task Handler_PreservesHexIdsAndIntegerPrecisionWhileScrubbingNestedValues() + { + var batch = WireBatch(); + var original = batch.Clone(); + var span = batch.ResourceSpans[0].ScopeSpans[0].Spans[0]; + var calls = 0; + var config = Config(new[] + { + new Customizer(payload => + { + calls++; + var jsonSpan = payload; + Assert.Equal(Convert.ToHexString(span.TraceId.ToByteArray()).ToLowerInvariant(), jsonSpan["traceId"]!.GetValue()); + Assert.Equal(Convert.ToHexString(span.SpanId.ToByteArray()).ToLowerInvariant(), jsonSpan["spanId"]!.GetValue()); + Assert.Equal(Convert.ToHexString(span.ParentSpanId.ToByteArray()).ToLowerInvariant(), jsonSpan["parentSpanId"]!.GetValue()); + Assert.Equal("18446744073709551615", jsonSpan["startTimeUnixNano"]!.GetValue()); + Assert.Equal("9007199254740993", jsonSpan["endTimeUnixNano"]!.GetValue()); + Assert.Equal((int)span.Kind, jsonSpan["kind"]!.GetValue()); + Assert.Equal((int)span.Status.Code, jsonSpan["status"]!["code"]!.GetValue()); + var values = jsonSpan["attributes"]!.AsArray(); + Assert.Equal("9223372036854775807", values[0]!["value"]!["intValue"]!.GetValue()); + Assert.Equal("-9223372036854775808", values[1]!["value"]!["intValue"]!.GetValue()); + Assert.Equal("AAH+/w==", values[2]!["value"]!["bytesValue"]!.GetValue()); + Assert.Equal(Convert.ToHexString(span.Links[0].TraceId.ToByteArray()).ToLowerInvariant(), jsonSpan["links"]![0]!["traceId"]!.GetValue()); + Assert.Equal(Convert.ToHexString(span.Links[0].SpanId.ToByteArray()).ToLowerInvariant(), jsonSpan["links"]![0]!["spanId"]!.GetValue()); + Scrub(payload); + return payload; + }) + }); + var transport = new CaptureHandler(); + + await SendAsync(config, transport, batch.ToByteArray()); + + Assert.Equal(1, calls); + Assert.Equal(original, batch); + var result = Assert.Single(transport.Batches); + var scope = result.ResourceSpans[0].ScopeSpans[0]; + Assert.Equal("safe-scope-value", scope.Scope.Attributes[0].Value.StringValue); + Assert.Equal("safe-scope-key", scope.Scope.Attributes[0].Key); + var nested = scope.Spans[0].Attributes[3].Value.ArrayValue.Values[0].KvlistValue.Values[0]; + Assert.Equal("safe-nested-key", nested.Key); + Assert.Equal("safe-nested-value", nested.Value.StringValue); + Assert.Equal(span.Attributes.Take(3), scope.Spans[0].Attributes.Take(3)); + Assert.Equal(span.StartTimeUnixNano, scope.Spans[0].StartTimeUnixNano); + Assert.Equal(span.EndTimeUnixNano, scope.Spans[0].EndTimeUnixNano); + Assert.Equal(span.Links, scope.Spans[0].Links); + } + + [Theory] + [InlineData("resource")] + [InlineData("scope")] + [InlineData("resourceSchemaUrl")] + [InlineData("scopeSchemaUrl")] + public async Task Handler_IsolatesSiblingMetadataAndSplitsOnlyChangedGroups(string field) + { + var batch = WireBatch(); + var resource = batch.ResourceSpans[0]; + var scope = resource.ScopeSpans[0]; + var first = scope.Spans[0]; + first.Name = "first"; + var second = first.Clone(); + second.Name = "second"; + second.SpanId = ByteString.CopyFrom(Convert.FromHexString("1123456789abcdef")); + var third = first.Clone(); + third.Name = "third"; + third.SpanId = ByteString.CopyFrom(Convert.FromHexString("2123456789abcdef")); + scope.Spans.Add(second); + scope.Spans.Add(third); + var original = batch.Clone(); + JsonObject? changedPayload = null; + JsonNode? originalMetadata = null; + var config = Config(new[] + { + new Customizer(payload => + { + if (payload["name"]!.GetValue() == "first") + { + changedPayload = payload; + originalMetadata = payload[field]?.DeepClone(); + switch (field) + { + case "resource": + payload["resource"]!["attributes"]![0]!["value"]!["stringValue"] = "changed-resource"; + break; + case "scope": payload["scope"]!["name"] = "changed-scope"; break; + default: payload[field] = "https://changed.example/schema"; break; + } + } + else + { + Assert.True(JsonNode.DeepEquals(originalMetadata, payload[field])); + Assert.NotSame(changedPayload!["resource"], payload["resource"]); + Assert.NotSame(changedPayload["scope"], payload["scope"]); + } + return payload; + }) + }); + var transport = new CaptureHandler(); + + await SendAsync(config, transport, batch.ToByteArray()); + + Assert.Equal(original, batch); + var exported = Assert.Single(transport.Batches); + var changedResource = exported.ResourceSpans.Single(r => r.ScopeSpans.Any(s => s.Spans.Any(p => p.Name == "first"))); + var unchangedResource = exported.ResourceSpans.Single(r => r.ScopeSpans.Any(s => s.Spans.Any(p => p.Name == "second"))); + var changedScope = changedResource.ScopeSpans.Single(s => s.Spans.Any(p => p.Name == "first")); + var unchangedScope = unchangedResource.ScopeSpans.Single(s => s.Spans.Any(p => p.Name == "second")); + Assert.Equal(first, Assert.Single(changedScope.Spans)); + Assert.Equal(new[] { second, third }, unchangedScope.Spans); + Assert.Equal(resource.Resource, unchangedResource.Resource); + Assert.Equal(resource.SchemaUrl, unchangedResource.SchemaUrl); + Assert.Equal(scope.Scope, unchangedScope.Scope); + Assert.Equal(scope.SchemaUrl, unchangedScope.SchemaUrl); + if (field is "resource" or "resourceSchemaUrl") + { + Assert.Equal(2, exported.ResourceSpans.Count); + Assert.Single(changedResource.ScopeSpans); + Assert.Single(unchangedResource.ScopeSpans); + } + else + { + Assert.Single(exported.ResourceSpans); + Assert.Equal(2, changedResource.ScopeSpans.Count); + } + var expectedResource = resource.Resource.Clone(); + var expectedScope = scope.Scope.Clone(); + if (field == "resource") expectedResource.Attributes[0].Value.StringValue = "changed-resource"; + if (field == "scope") expectedScope.Name = "changed-scope"; + Assert.Equal(expectedResource, changedResource.Resource); + Assert.Equal(expectedScope, changedScope.Scope); + Assert.Equal(field == "resourceSchemaUrl" ? "https://changed.example/schema" : resource.SchemaUrl, changedResource.SchemaUrl); + Assert.Equal(field == "scopeSchemaUrl" ? "https://changed.example/schema" : scope.SchemaUrl, changedScope.SchemaUrl); + } + + [Theory] + [InlineData("traceId")] + [InlineData("spanId")] + [InlineData("parentSpanId")] + [InlineData("root-parent")] + public void Export_RejectsLaterHookIdentityChangesBeforeTransmitting(string field) + { + var calls = 0; + var config = Config(new ISpanCustomizer[] + { + new OmittedHook(), + new Customizer(payload => + { + calls++; + payload[field == "root-parent" ? "parentSpanId" : field] = new string('1', field == "traceId" ? 32 : 16); + return payload; + }) + }); + using var activity = CompletedSpan("unchanged", root: field == "root-parent"); + var transport = new CaptureHandler(); + + Assert.Equal(ExportResult.Failure, Export(config, transport, activity)); + + Assert.Equal(1, calls); + Assert.Equal(0, transport.Requests); + Assert.Empty(transport.Batches); + } + + [Fact] + public void Export_RejectsSwappingSiblingIdentitiesEvenThoughBatchIdentitySetWouldBeUnchanged() + { + using var first = CompletedSpan("first"); + using var second = CompletedSpan("second"); + var calls = 0; + var config = Config(new[] + { + new Customizer(payload => + { + calls++; + var other = payload["name"]!.GetValue() == "first" ? second : first; + payload["traceId"] = other.TraceId.ToHexString(); + payload["spanId"] = other.SpanId.ToHexString(); + payload["parentSpanId"] = other.ParentSpanId.ToHexString(); + return payload; + }) + }); + var transport = new CaptureHandler(); + + Assert.Equal(ExportResult.Failure, Export(config, transport, first, second)); + + Assert.Equal(1, calls); + Assert.Equal(0, transport.Requests); + Assert.Empty(transport.Batches); + } + + [Fact] + public void Export_ValidatesEachHookBeforeLaterHookCanRestoreIdentity() + { + string? originalId = null; + var laterCalled = false; + var config = Config(new[] + { + new Customizer(payload => + { + originalId = payload["spanId"]!.GetValue(); + payload["spanId"] = new string('1', 16); + return payload; + }), + new Customizer(payload => + { + laterCalled = true; + payload["spanId"] = originalId; + return payload; + }) + }); + using var activity = CompletedSpan("unchanged"); + var transport = new CaptureHandler(); + + Assert.Equal(ExportResult.Failure, Export(config, transport, activity)); + + Assert.False(laterCalled); + Assert.Equal(0, transport.Requests); + } + + [Fact] + public void Export_OmittedHookPreservesPayloadAndExplicitResubmissionRunsHooksAgain() + { + using var activity = CompletedSpan("unchanged", root: true); + var baseline = new CaptureHandler(); + Assert.Equal(ExportResult.Success, Export(Config(), baseline, activity)); + var calls = 0; + var config = Config(new ISpanCustomizer[] { new OmittedHook(), new Customizer(payload => { calls++; return payload; }) }); + var customized = new CaptureHandler(); + Assert.Equal(ExportResult.Success, Export(config, customized, activity)); + Assert.Equal(Assert.Single(baseline.Batches), Assert.Single(customized.Batches)); + Assert.Equal(ExportResult.Success, Export(config, new CaptureHandler(), activity)); + Assert.Equal(2, calls); + } + + [Fact] + public void Export_LaterSpanCannotMutateAnEarlierAcceptedResultThroughRetainedJson() + { + using var first = CompletedSpan("first"); + using var second = CompletedSpan("second"); + JsonObject? retained = null; + var config = Config(new[] + { + new Customizer(payload => + { + if (payload["name"]!.GetValue() == "first") + { + payload["name"] = "accepted-first"; + retained = payload; + } + else + { + retained!["name"] = "corrupted"; + retained["spanId"] = new string('1', 16); + retained["scope"]!["name"] = "corrupted-scope"; + retained["resource"] = null; + } + return payload; + }) + }); + var transport = new CaptureHandler(); + + Assert.Equal(ExportResult.Success, Export(config, transport, first, second)); + + var resource = Assert.Single(Assert.Single(transport.Batches).ResourceSpans); + Assert.Equal("secret-service", resource.Resource.Attributes.Single(a => a.Key == "service.name").Value.StringValue); + var scope = Assert.Single(resource.ScopeSpans); + Assert.Equal(_source.Name, scope.Scope.Name); + Assert.Equal(new[] { "accepted-first", "second" }, scope.Spans.Select(span => span.Name)); + Assert.Equal(first.SpanId.ToHexString(), Convert.ToHexString(scope.Spans[0].SpanId.ToByteArray()).ToLowerInvariant()); + Assert.Equal("corrupted", retained!["name"]!.GetValue()); + } + + [Theory] + [InlineData("throw")] + [InlineData("null")] + [InlineData("unknown-field")] + [InlineData("unknown-nested-field")] + [InlineData("wrong-type")] + [InlineData("invalid-hex")] + [InlineData("invalid-shape")] + [InlineData("unknown-resource-field")] + [InlineData("unknown-scope-field")] + public async Task Handler_LaterSpanHookFailurePreventsSendingEarlierCustomizedSpans(string failure) + { + var batch = WireBatch(); + var spans = batch.ResourceSpans[0].ScopeSpans[0].Spans; + spans[0].Name = "first"; + var second = spans[0].Clone(); + second.Name = "second"; + second.SpanId = ByteString.CopyFrom(Convert.FromHexString("1123456789abcdef")); + spans.Add(second); + var calls = new List(); + var config = Config(new[] + { + new Customizer(payload => + { + calls.Add($"first({payload["name"]!.GetValue()})"); + payload["name"] = payload["name"]!.GetValue() + "-customized"; + return payload; + }), + new Customizer(payload => + { + calls.Add($"second({payload["name"]!.GetValue()})"); + if (payload["name"]!.GetValue() != "second-customized") return payload; + switch (failure) + { + case "throw": throw new InvalidOperationException("scrubbing failed"); + case "null": return null!; + case "unknown-field": payload["unexpected"] = true; break; + case "unknown-nested-field": payload["status"]!["unexpected"] = true; break; + case "wrong-type": payload["name"] = 17; break; + case "invalid-hex": payload["traceId"] = new string('z', 32); break; + case "invalid-shape": payload["scope"] = "not an object"; break; + case "unknown-resource-field": payload["resource"]!["unexpected"] = true; break; + case "unknown-scope-field": payload["scope"]!["unexpected"] = true; break; + } + return payload; + }) + }); + var transport = new CaptureHandler(); + + await Assert.ThrowsAsync(() => SendAsync(config, transport, batch.ToByteArray())); + + Assert.Equal(new[] { "first(first)", "second(first-customized)", "first(second)", "second(second-customized)" }, calls); + Assert.Equal(0, transport.Requests); + Assert.Empty(transport.Batches); + } + + [Fact] + public async Task Handler_RejectsInvalidIntermediateResultBeforeLaterHookCanRepairIt() + { + var laterCalled = false; + var config = Config(new[] + { + new Customizer(payload => + { + payload["unexpected"] = true; + return payload; + }), + new Customizer(payload => + { + laterCalled = true; + payload.Remove("unexpected"); + return payload; + }) + }); + var transport = new CaptureHandler(); + + await Assert.ThrowsAsync(() => SendAsync(config, transport, WireBatch().ToByteArray())); + + Assert.False(laterCalled); + Assert.Equal(0, transport.Requests); + Assert.Empty(transport.Batches); + } + + [Theory] + [InlineData(false)] + [InlineData(true)] + public async Task Handler_RejectsMalformedOrCompressedPayloadBeforeTransport(bool compressed) + { + var transport = new CaptureHandler(); + var bytes = compressed ? new ExportTraceServiceRequest().ToByteArray() : new byte[] { 0xff }; + + await Assert.ThrowsAsync(() => SendAsync(Config(new[] { new OmittedHook() }), transport, bytes, compressed)); + + Assert.Equal(0, transport.Requests); + } + + [Theory] + [InlineData("batch")] + [InlineData("span")] + [InlineData("attribute-value")] + public async Task Handler_RejectsUnknownProtobufFieldsBeforeCallingHooks(string location) + { + var batch = WireBatch(); + var spans = batch.ResourceSpans[0].ScopeSpans[0].Spans; + switch (location) + { + case "batch": batch = ExportTraceServiceRequest.Parser.ParseFrom(WithUnknownField(batch)); break; + case "span": + var later = OtlpSpan.Parser.ParseFrom(WithUnknownField(spans[0])); + later.SpanId = ByteString.CopyFrom(Convert.FromHexString("1123456789abcdef")); + spans.Add(later); + break; + case "attribute-value": + spans[0].Attributes[0].Value = AnyValue.Parser.ParseFrom(WithUnknownField(spans[0].Attributes[0].Value)); + break; + } + var called = false; + var config = Config(new[] { new Customizer(payload => { called = true; return payload; }) }); + var transport = new CaptureHandler(); + + await Assert.ThrowsAsync(() => SendAsync(config, transport, batch.ToByteArray())); + + Assert.False(called); + Assert.Equal(0, transport.Requests); + } + + [Theory] + [InlineData(false)] + [InlineData(true)] + public async Task Handler_RemovedOrNullMetadataDoesNotRestoreOriginalContent(bool nullMetadata) + { + var batch = WireBatch(); + var config = Config(new[] + { + new Customizer(payload => + { + foreach (var field in new[] { "resource", "scope", "resourceSchemaUrl", "scopeSchemaUrl" }) + { + if (nullMetadata) payload[field] = null; + else payload.Remove(field); + } + foreach (var field in new[] { "attributes", "events", "links", "status" }) + payload.Remove(field); + payload["kind"] = (int)OtlpSpan.Types.SpanKind.Server; + payload["traceState"] = "sanitized=state"; + payload["startTimeUnixNano"] = "9007199254740993"; + payload["endTimeUnixNano"] = "9007199254740994"; + return payload; + }) + }); + var transport = new CaptureHandler(); + + await SendAsync(config, transport, batch.ToByteArray()); + + var resource = Assert.Single(Assert.Single(transport.Batches).ResourceSpans); + Assert.Null(resource.Resource); + Assert.Equal("", resource.SchemaUrl); + var scope = Assert.Single(resource.ScopeSpans); + Assert.Null(scope.Scope); + Assert.Equal("", scope.SchemaUrl); + var exported = Assert.Single(scope.Spans); + Assert.Empty(exported.Attributes); + Assert.Empty(exported.Events); + Assert.Empty(exported.Links); + Assert.Null(exported.Status); + Assert.Equal(OtlpSpan.Types.SpanKind.Server, exported.Kind); + Assert.Equal("sanitized=state", exported.TraceState); + Assert.Equal(9007199254740993UL, exported.StartTimeUnixNano); + Assert.Equal(9007199254740994UL, exported.EndTimeUnixNano); + } + + private static void Scrub(JsonNode node) + { + switch (node) + { + case JsonObject obj: + foreach (var (key, value) in obj.ToArray()) + { + if (key is "traceId" or "spanId" or "parentSpanId" || value == null) continue; + if (value is JsonValue scalar && scalar.TryGetValue(out var text)) + obj[key] = text.Replace("secret", "safe", StringComparison.Ordinal); + else + Scrub(value); + } + break; + case JsonArray array: + for (var i = 0; i < array.Count; i++) + { + if (array[i] is JsonValue scalar && scalar.TryGetValue(out var text)) + array[i] = text.Replace("secret", "safe", StringComparison.Ordinal); + else if (array[i] is { } child) + Scrub(child); + } + break; + } + } + + private Activity CompletedSpan(string name, bool root = false) + { + var parent = root ? default : new ActivityContext(ActivityTraceId.CreateRandom(), ActivitySpanId.CreateRandom(), ActivityTraceFlags.Recorded, "vendor=value", true); + var linkContext = new ActivityContext(ActivityTraceId.CreateRandom(), ActivitySpanId.CreateRandom(), ActivityTraceFlags.Recorded); + var activity = _source.StartActivity(name, ActivityKind.Client, parent, + links: new[] { new ActivityLink(linkContext, new ActivityTagsCollection { { "secret-link-key", "secret-link-value" } }) })!; + activity.SetTag("secret-tag-key", "secret-input"); + activity.SetTag("braintrust.parent", "project_id:original"); + activity.AddEvent(new ActivityEvent("secret-event", tags: new ActivityTagsCollection { { "secret-event-key", "secret-detail" } })); + activity.SetStatus(ActivityStatusCode.Error, "secret-status"); + activity.Stop(); + return activity; + } + + private static ExportTraceServiceRequest WireBatch() + { + var span = new OtlpSpan + { + TraceId = ByteString.CopyFrom(Convert.FromHexString("00112233445566778899aabbccddeeff")), + SpanId = ByteString.CopyFrom(Convert.FromHexString("0123456789abcdef")), + ParentSpanId = ByteString.CopyFrom(Convert.FromHexString("fedcba9876543210")), + Name = "secret-span", + Kind = OtlpSpan.Types.SpanKind.Client, + StartTimeUnixNano = ulong.MaxValue, + EndTimeUnixNano = 9007199254740993UL, + Status = new Status { Code = Status.Types.StatusCode.Error, Message = "secret-status" }, + Attributes = + { + new KeyValue { Key = "maximum", Value = new AnyValue { IntValue = long.MaxValue } }, + new KeyValue { Key = "minimum", Value = new AnyValue { IntValue = long.MinValue } }, + new KeyValue { Key = "bytes", Value = new AnyValue { BytesValue = ByteString.CopyFrom(new byte[] { 0, 1, 254, 255 }) } }, + new KeyValue + { + Key = "nested", + Value = new AnyValue + { + ArrayValue = new ArrayValue + { + Values = { new AnyValue { KvlistValue = new KeyValueList { Values = { new KeyValue { Key = "secret-nested-key", Value = new AnyValue { StringValue = "secret-nested-value" } } } } } } + } + } + } + }, + Events = + { + new OtlpSpan.Types.Event + { + Name = "secret-event", + TimeUnixNano = 9007199254740993UL, + Attributes = { new KeyValue { Key = "secret-event-key", Value = new AnyValue { StringValue = "secret-event-value" } } } + } + }, + Links = + { + new OtlpSpan.Types.Link + { + TraceId = ByteString.CopyFrom(Convert.FromHexString("ffeeddccbbaa99887766554433221100")), + SpanId = ByteString.CopyFrom(Convert.FromHexString("8877665544332211")) + } + } + }; + return new ExportTraceServiceRequest + { + ResourceSpans = + { + new ResourceSpans + { + Resource = new Resource + { + Attributes = { new KeyValue { Key = "secret-resource-key", Value = new AnyValue { StringValue = "secret-resource-value" } } } + }, + SchemaUrl = "https://secret.example/resource", + ScopeSpans = + { + new ScopeSpans + { + SchemaUrl = "https://secret.example/scope", + Scope = new InstrumentationScope + { + Name = "secret-scope", + Attributes = { new KeyValue { Key = "secret-scope-key", Value = new AnyValue { StringValue = "secret-scope-value" } } } + }, + Spans = { span } + } + } + } + } + }; + } + + private static byte[] WithUnknownField(IMessage message) + { + using var stream = new MemoryStream(); + stream.Write(message.ToByteArray()); + using (var output = new CodedOutputStream(stream, leaveOpen: true)) + { + output.WriteTag(127, WireFormat.WireType.Varint); + output.WriteUInt32(1); + output.Flush(); + } + return stream.ToArray(); + } + + private static BraintrustConfig Config(IEnumerable? customizers = null) => BraintrustConfig.Of( + customizers ?? Array.Empty(), + ("BRAINTRUST_API_KEY", "test-key"), + ("BRAINTRUST_DEFAULT_PROJECT_ID", "default-project")); + + private static ExportResult Export(BraintrustConfig config, CaptureHandler transport, params Activity[] activities) + { + using var client = BraintrustTracing.CreateHttpClient(config, transport); + var exporter = new OtlpTraceExporter(new OtlpExporterOptions + { + Endpoint = new Uri("https://example.invalid/otel/v1/traces"), + Protocol = OtlpExportProtocol.HttpProtobuf, + HttpClientFactory = () => client + }); + using var provider = OpenTelemetry.Sdk.CreateTracerProviderBuilder() + .SetResourceBuilder(ResourceBuilder.CreateEmpty() + .AddService("secret-service", serviceInstanceId: "test-instance") + .AddAttributes(new[] { new KeyValuePair("secret-resource-key", "secret-resource-value") })) + .AddProcessor(new SimpleActivityExportProcessor(exporter)) + .Build(); + using var batch = new Batch(activities, activities.Length); + return exporter.Export(in batch); + } + + private static async Task SendAsync(BraintrustConfig config, CaptureHandler transport, byte[] bytes, bool compressed = false) + { + using var client = BraintrustTracing.CreateHttpClient(config, transport); + using var request = new HttpRequestMessage(HttpMethod.Post, "https://example.invalid/otel/v1/traces") + { + Content = new ByteArrayContent(bytes) + }; + request.Content.Headers.ContentType = new MediaTypeHeaderValue("application/x-protobuf"); + if (compressed) request.Content.Headers.ContentEncoding.Add("gzip"); + using var response = await client.SendAsync(request); + Assert.Equal(HttpStatusCode.OK, response.StatusCode); + } + + private sealed class OmittedHook : ISpanCustomizer { } + + private sealed class Customizer(Func hook) : ISpanCustomizer + { + public JsonObject OnSpanExport(JsonObject payload) => hook(payload); + } + + private sealed class CaptureHandler : HttpMessageHandler + { + internal List Batches { get; } = new(); + internal int Requests { get; private set; } + internal string? Authorization { get; private set; } + internal string? Parent { get; private set; } + + protected override HttpResponseMessage Send(HttpRequestMessage request, CancellationToken cancellationToken) + { + Requests++; + Authorization = request.Headers.Authorization?.ToString(); + Parent = request.Headers.GetValues("x-bt-parent").Single(); + Batches.Add(ExportTraceServiceRequest.Parser.ParseFrom(request.Content!.ReadAsByteArrayAsync(cancellationToken).GetAwaiter().GetResult())); + return new HttpResponseMessage(HttpStatusCode.OK); + } + + protected override Task SendAsync(HttpRequestMessage request, CancellationToken cancellationToken) + => Task.FromResult(Send(request, cancellationToken)); + } + + public void Dispose() + { + _listener.Dispose(); + _source.Dispose(); + } +}