Skip to content
Draft
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: 2 additions & 0 deletions include/datadog/sampling_decision.h
Original file line number Diff line number Diff line change
Expand Up @@ -39,6 +39,8 @@ struct SamplingDecision {
// The per-second maximum allowed number of "keeps" configured for the limiter
// consulted in this decision, if any.
Optional<double> limiter_max_per_second;
// The outcome of the probability comparison before rate limiting, if any.
Optional<bool> probability_sampled;
// The provenance of this decision.
Origin origin;
};
Expand Down
2 changes: 2 additions & 0 deletions include/datadog/trace_segment.h
Original file line number Diff line number Diff line change
Expand Up @@ -77,6 +77,7 @@ class TraceSegment {
std::vector<std::unique_ptr<SpanData>> spans_;
std::size_t num_finished_spans_;
Optional<SamplingDecision> sampling_decision_;
const Optional<std::string> otel_w3c_tracestate_;
const Optional<std::string> additional_w3c_tracestate_;
const Optional<std::string> additional_datadog_w3c_tracestate_;

Expand All @@ -99,6 +100,7 @@ class TraceSegment {
Optional<std::string> origin, std::size_t tags_header_max_size,
std::vector<std::pair<std::string, std::string>> trace_tags,
Optional<SamplingDecision> sampling_decision,
Optional<std::string> otel_w3c_tracestate,
Optional<std::string> additional_w3c_tracestate,
Optional<std::string> additional_datadog_w3c_tracestate,
std::unique_ptr<SpanData> local_root,
Expand Down
2 changes: 2 additions & 0 deletions src/datadog/extracted_data.h
Original file line number Diff line number Diff line change
Expand Up @@ -31,6 +31,8 @@ struct ExtractedData {
// then `additional_w3c_tracestate` is null.
// `additional_w3c_tracestate` is used for the `W3C` injection style.
Optional<std::string> additional_w3c_tracestate;
// The raw value of the OpenTelemetry `ot` tracestate member, if present.
Optional<std::string> otel_w3c_tracestate;
// If this `ExtractedData` was created on account of `PropagationStyle::W3C`,
// and if the "tracestate" header contained a "dd" (Datadog) entry, then
// `additional_datadog_w3c_tracestate` contains fields from within the "dd"
Expand Down
6 changes: 4 additions & 2 deletions src/datadog/trace_sampler.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -51,7 +51,8 @@ SamplingDecision TraceSampler::decide(const SpanData& span) {
decision.limiter_max_per_second = limiter_max_per_second_;
decision.configured_rate = rule.rate;
const std::uint64_t threshold = max_id_from_rate(rule.rate);
if (knuth_hash(span.trace_id.low) <= threshold) {
decision.probability_sampled = knuth_hash(span.trace_id.low) <= threshold;
if (*decision.probability_sampled) {
if (rule.bypass_limiter) {
decision.priority = int(SamplingPriority::USER_KEEP);
return decision;
Expand Down Expand Up @@ -91,7 +92,8 @@ SamplingDecision TraceSampler::decide(const SpanData& span) {
}

const std::uint64_t threshold = max_id_from_rate(*decision.configured_rate);
if (knuth_hash(span.trace_id.low) <= threshold) {
decision.probability_sampled = knuth_hash(span.trace_id.low) <= threshold;
if (*decision.probability_sampled) {
decision.priority = int(SamplingPriority::AUTO_KEEP);
} else {
decision.priority = int(SamplingPriority::AUTO_DROP);
Expand Down
98 changes: 79 additions & 19 deletions src/datadog/trace_segment.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -6,6 +6,7 @@
#include <datadog/injection_options.h>
#include <datadog/logger.h>
#include <datadog/optional.h>
#include <datadog/sampling_mechanism.h>
#include <datadog/span_defaults.h>
#include <datadog/telemetry/metrics.h>
#include <datadog/telemetry/telemetry.h>
Expand All @@ -14,6 +15,8 @@
#include <array>
#include <cassert>
#include <charconv>
#include <cmath>
#include <cstdint>
#include <string>
#include <unordered_map>
#include <utility>
Expand All @@ -23,6 +26,7 @@
#include "endpoint_inferral.h"
#include "hex.h"
#include "platform_util.h"
#include "sampling_util.h"
#include "span_data.h"
#include "span_sampler.h"
#include "tag_propagation.h"
Expand Down Expand Up @@ -151,6 +155,54 @@ Optional<std::string> format_rate(double rate, Logger& logger) {
return std::string(begin, end);
}

bool is_probability_mechanism(int mechanism) {
switch (static_cast<SamplingMechanism>(mechanism)) {
case SamplingMechanism::DEFAULT:
case SamplingMechanism::AGENT_RATE:
case SamplingMechanism::REMOTE_RATE_AUTO:
case SamplingMechanism::RULE:
case SamplingMechanism::REMOTE_RATE_USER_DEFINED:
case SamplingMechanism::REMOTE_RATE_EMERGENCY:
case SamplingMechanism::REMOTE_RULE:
case SamplingMechanism::REMOTE_ADAPTIVE_RULE:
return true;
default:
return false;
}
}

Optional<std::string> resolve_otel_tracestate(
TraceID trace_id, const SamplingDecision& decision,
const Optional<std::string>& inherited) {
const StringView raw = inherited ? StringView(*inherited) : StringView{};
if (decision.origin != SamplingDecision::Origin::LOCAL) {
return inherited ? sanitize_otel_tracestate(raw) : nullopt;
}

if (!decision.mechanism || !decision.configured_rate ||
!decision.probability_sampled ||
!is_probability_mechanism(*decision.mechanism) ||
(*decision.probability_sampled && decision.priority <= 0)) {
return rewrite_otel_tracestate(raw, extract_otel_random_value(raw),
nullopt);
}

constexpr std::uint64_t max_value = UINT64_C(1) << 56;
auto threshold = static_cast<std::uint64_t>(
std::round((1.0 - decision.configured_rate->value()) *
static_cast<double>(max_value)));
threshold = std::min(threshold, max_value - 1);

std::uint64_t random_value = (~knuth_hash(trace_id.low)) >> 8;
if (*decision.probability_sampled && random_value < threshold) {
random_value = threshold;
} else if (!*decision.probability_sampled && random_value >= threshold) {
random_value = threshold == 0 ? 0 : threshold - 1;
}

return rewrite_otel_tracestate(raw, random_value, threshold);
}

} // anonymous namespace

TraceSegment::TraceSegment(
Expand All @@ -166,6 +218,7 @@ TraceSegment::TraceSegment(
std::size_t tags_header_max_size,
std::vector<std::pair<std::string, std::string>> trace_tags,
Optional<SamplingDecision> sampling_decision,
Optional<std::string> otel_w3c_tracestate,
Optional<std::string> additional_w3c_tracestate,
Optional<std::string> additional_datadog_w3c_tracestate,
std::unique_ptr<SpanData> local_root,
Expand All @@ -184,6 +237,7 @@ TraceSegment::TraceSegment(
trace_tags_(std::move(trace_tags)),
num_finished_spans_(0),
sampling_decision_(std::move(sampling_decision)),
otel_w3c_tracestate_(std::move(otel_w3c_tracestate)),
additional_w3c_tracestate_(std::move(additional_w3c_tracestate)),
additional_datadog_w3c_tracestate_(
std::move(additional_datadog_w3c_tracestate)),
Expand Down Expand Up @@ -216,14 +270,14 @@ Optional<SamplingDecision> TraceSegment::sampling_decision() const {

Optional<std::pair<std::string, std::uint32_t>> TraceSegment::w3c_link_context(
const SpanData& span) const {
int sampling_priority;
SamplingDecision sampling_decision;
std::vector<std::pair<std::string, std::string>> trace_tags;
{
std::lock_guard<std::mutex> lock(mutex_);
if (!sampling_decision_) {
return nullopt;
}
sampling_priority = sampling_decision_->priority;
sampling_decision = *sampling_decision_;
trace_tags = trace_tags_;

const Optional<std::string> trace_source_tag =
Expand All @@ -234,10 +288,13 @@ Optional<std::pair<std::string, std::uint32_t>> TraceSegment::w3c_link_context(
}

return std::make_pair(
encode_tracestate(span.span_id, sampling_priority, origin_, trace_tags,
additional_datadog_w3c_tracestate_,
additional_w3c_tracestate_),
sampling_priority > 0 ? 1u : 0u);
encode_tracestate(
span.span_id, sampling_decision.priority, origin_, trace_tags,
additional_datadog_w3c_tracestate_,
resolve_otel_tracestate(span.trace_id, sampling_decision,
otel_w3c_tracestate_),
additional_w3c_tracestate_),
sampling_decision.priority > 0 ? 1u : 0u);
}

Logger& TraceSegment::logger() const { return *logger_; }
Expand Down Expand Up @@ -443,13 +500,13 @@ bool TraceSegment::inject(DictWriter& writer, const SpanData& span,
// and trace tags might change when that happens ("_dd.p.dm").
// So, we lock here, make a sampling decision if necessary, and then copy the
// decision and trace tags before unlocking.
int sampling_priority;
SamplingDecision sampling_decision;
std::vector<std::pair<std::string, std::string>> trace_tags;
{
std::lock_guard<std::mutex> lock(mutex_);
make_sampling_decision_if_null();
assert(sampling_decision_);
sampling_priority = sampling_decision_->priority;
sampling_decision = *sampling_decision_;
trace_tags = trace_tags_;
}

Expand All @@ -464,7 +521,7 @@ bool TraceSegment::inject(DictWriter& writer, const SpanData& span,
// - the local root span is NOT created by another product (no `_dd.p.ts`)
// - sampling priority is DROP
if (!tracing_enabled_) {
if (!trace_source_tag && sampling_priority <= 0) {
if (!trace_source_tag && sampling_decision.priority <= 0) {
writer.erase("x-datadog-trace-id");
writer.erase("x-datadog-parent-id");
writer.erase("x-datadog-sampling-priority");
Expand Down Expand Up @@ -492,7 +549,7 @@ bool TraceSegment::inject(DictWriter& writer, const SpanData& span,
writer.set("x-datadog-trace-id", std::to_string(span.trace_id.low));
writer.set("x-datadog-parent-id", std::to_string(span.span_id));
writer.set("x-datadog-sampling-priority",
std::to_string(sampling_priority));
std::to_string(sampling_decision.priority));
if (origin_) {
writer.set("x-datadog-origin", *origin_);
}
Expand All @@ -509,7 +566,8 @@ bool TraceSegment::inject(DictWriter& writer, const SpanData& span,
writer.set("x-b3-traceid", hex_padded(span.trace_id.low));
}
writer.set("x-b3-spanid", hex_padded(span.span_id));
writer.set("x-b3-sampled", std::to_string(int(sampling_priority > 0)));
writer.set("x-b3-sampled",
std::to_string(int(sampling_decision.priority > 0)));
if (origin_) {
writer.set("x-datadog-origin", *origin_);
}
Expand All @@ -519,14 +577,16 @@ bool TraceSegment::inject(DictWriter& writer, const SpanData& span,
{"header_style:b3multi"});
break;
case PropagationStyle::W3C:
writer.set(
"traceparent",
encode_traceparent(span.trace_id, span.span_id, sampling_priority));
writer.set(
"tracestate",
encode_tracestate(span.span_id, sampling_priority, origin_,
trace_tags, additional_datadog_w3c_tracestate_,
additional_w3c_tracestate_));
writer.set("traceparent",
encode_traceparent(span.trace_id, span.span_id,
sampling_decision.priority));
writer.set("tracestate",
encode_tracestate(
span.span_id, sampling_decision.priority, origin_,
trace_tags, additional_datadog_w3c_tracestate_,
resolve_otel_tracestate(span.trace_id, sampling_decision,
otel_w3c_tracestate_),
additional_w3c_tracestate_));
telemetry::counter::increment(metrics::tracer::trace_context::injected,
{"header_style:tracecontext"});
break;
Expand Down
4 changes: 3 additions & 1 deletion src/datadog/tracer.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -247,7 +247,8 @@ Span Tracer::create_span(const SpanConfig& config) {
logger_, collector_, config_manager_->trace_sampler(), span_sampler_,
defaults, config_manager_, runtime_id_, injection_styles_, hostname_,
nullopt /* origin */, tags_header_max_size_, std::move(trace_tags),
nullopt /* sampling_decision */, nullopt /* additional_w3c_tracestate */,
nullopt /* sampling_decision */, nullopt /* otel_w3c_tracestate */,
nullopt /* additional_w3c_tracestate */,
nullopt /* additional_datadog_w3c_tracestate*/, std::move(span_data),
resource_renaming_mode_, tracing_enabled_);
Span span{span_data_ptr, segment,
Expand Down Expand Up @@ -483,6 +484,7 @@ Expected<Span> Tracer::extract_span(const DictReader& reader,
injection_styles_, hostname_, std::move(merged_context.origin),
tags_header_max_size_, std::move(merged_context.trace_tags),
std::move(sampling_decision),
std::move(merged_context.otel_w3c_tracestate),
std::move(merged_context.additional_w3c_tracestate),
std::move(merged_context.additional_datadog_w3c_tracestate),
std::move(span_data), resource_renaming_mode_, tracing_enabled_);
Expand Down
Loading
Loading