Skip to content
Merged
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
1 change: 1 addition & 0 deletions apps/ingest/Cargo.lock

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

1 change: 1 addition & 0 deletions apps/ingest/Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -28,6 +28,7 @@ opentelemetry_sdk = { version = "0.32", features = [
] }
opentelemetry-otlp = { version = "0.32", default-features = false, features = [
"http-proto",
"gzip-http",
"reqwest-client",
"reqwest-rustls",
"trace",
Expand Down
7 changes: 7 additions & 0 deletions apps/ingest/alchemy.run.ts
Original file line number Diff line number Diff line change
Expand Up @@ -23,6 +23,7 @@ import {
resolveIngestDesiredCount,
resolveIngestNamespaceName,
resolveIngestScaling,
resolveIngestSelfTraceSampleRatio,
resolveIngestTaskSize,
stageDeploysCollector,
stageEnablesReplayBlobs,
Expand Down Expand Up @@ -225,6 +226,7 @@ export const createMapleIngest = ({ stage, domains, region, dbRole }: CreateMapl
Effect.gen(function* () {
const replayBlobs = yield* replayBlobWriterCredentials(stage, region)
const taskSize = resolveIngestTaskSize(stage)
const selfTraceSampleRatio = resolveIngestSelfTraceSampleRatio(stage)
const scaling = resolveIngestScaling(stage)
const name = (base: string) => resolveAwsResourceName(base, stage, region)
const fleets = parseIngestFleets((yield* optionalPlain("MAPLE_INGEST_FLEETS")).MAPLE_INGEST_FLEETS)
Expand Down Expand Up @@ -848,6 +850,11 @@ export const createMapleIngest = ({ stage, domains, region, dbRole }: CreateMapl
...(collectorEndpoint
? { INGEST_FORWARD_OTLP_ENDPOINT: collectorEndpoint }
: yield* optionalPlain("INGEST_FORWARD_OTLP_ENDPOINT")),
// The gateway's own traces are most of that collector's volume; prd
// head-samples them (logs and metrics stay complete).
...(selfTraceSampleRatio
? { INGEST_SELF_TRACE_SAMPLE_RATIO: selfTraceSampleRatio }
: yield* optionalPlain("INGEST_SELF_TRACE_SAMPLE_RATIO")),
// Every optional entry below is `yield*`-ed. `optionalPlain` returns a
// `Config`, not a record, so a bare `...optionalPlain("X")` spreads the
// Config's own fields (`_tag`, `original`, `mapOrFail`) into the task
Expand Down
159 changes: 152 additions & 7 deletions apps/ingest/src/main.rs
Original file line number Diff line number Diff line change
Expand Up @@ -68,7 +68,10 @@ use maple_ingest::wal_store::WalSegmentStore;
use moka::future::Cache;
use opentelemetry::trace::TracerProvider as _;
use opentelemetry_appender_tracing::layer::OpenTelemetryTracingBridge;
use opentelemetry_otlp::{LogExporter, MetricExporter, Protocol, SpanExporter, WithExportConfig};
use opentelemetry_otlp::{
Compression as OtlpCompression, LogExporter, MetricExporter, Protocol, SpanExporter,
WithExportConfig, WithHttpConfig,
};
use opentelemetry_proto::tonic::collector::logs::v1::ExportLogsServiceRequest;
use opentelemetry_proto::tonic::collector::metrics::v1::ExportMetricsServiceRequest;
use opentelemetry_proto::tonic::collector::trace::v1::ExportTraceServiceRequest;
Expand All @@ -81,7 +84,9 @@ use opentelemetry_sdk::metrics::periodic_reader_with_async_runtime::PeriodicRead
use opentelemetry_sdk::metrics::{SdkMeterProvider, Temporality};
use opentelemetry_sdk::runtime::Tokio as OtelTokio;
use opentelemetry_sdk::trace::span_processor_with_async_runtime::BatchSpanProcessor;
use opentelemetry_sdk::trace::{BatchConfigBuilder, SdkTracerProvider};
use opentelemetry_sdk::trace::{
BatchConfigBuilder, Sampler, SamplingDecision, SamplingResult, SdkTracerProvider, ShouldSample,
};
use prost::Message;
use reqwest::Client;
use serde::{Deserialize, Serialize};
Expand Down Expand Up @@ -148,6 +153,8 @@ struct AppConfig {
otlp_grpc_port: Option<u16>,
forward_endpoint: String,
forward_timeout: Duration,
/// Head-sampling ratio for the gateway's own traces, in (0, 1].
self_trace_sample_ratio: f64,
write_mode: WriteMode,
tinybird: TinybirdConfig,
max_request_body_bytes: usize,
Expand Down Expand Up @@ -278,6 +285,11 @@ impl AppConfig {
if forward_endpoint.is_empty() {
return Err("INGEST_FORWARD_OTLP_ENDPOINT is required".to_owned());
}
let self_trace_sample_ratio = parse_sample_ratio(
"INGEST_SELF_TRACE_SAMPLE_RATIO",
std::env::var("INGEST_SELF_TRACE_SAMPLE_RATIO").ok(),
1.0,
)?;

let internal_org_id = std::env::var("MAPLE_INTERNAL_ORG_ID")
.unwrap_or_default()
Expand Down Expand Up @@ -651,6 +663,7 @@ impl AppConfig {
otlp_grpc_port,
forward_endpoint,
forward_timeout: Duration::from_millis(forward_timeout_ms),
self_trace_sample_ratio,
write_mode,
tinybird,
max_request_body_bytes,
Expand Down Expand Up @@ -1849,6 +1862,7 @@ fn init_tracing(
bind_port: u16,
service_instance_id: &str,
internal_org_id: &str,
trace_sample_ratio: f64,
) -> Option<TelemetryProviders> {
// Two filters on purpose. The registry-wide filter is pinned at `info`
// because every gateway span (`ingest`, `ingest.authenticate`, the Postgres
Expand Down Expand Up @@ -1898,6 +1912,7 @@ fn init_tracing(
.with_http()
.with_endpoint(format!("{forward_endpoint}/v1/traces"))
.with_protocol(Protocol::HttpBinary)
.with_compression(OtlpCompression::Gzip)
.build()
{
Ok(exporter) => exporter,
Expand All @@ -1916,6 +1931,7 @@ fn init_tracing(
.with_http()
.with_endpoint(format!("{forward_endpoint}/v1/logs"))
.with_protocol(Protocol::HttpBinary)
.with_compression(OtlpCompression::Gzip)
.build()
{
Ok(exporter) => exporter,
Expand All @@ -1932,11 +1948,11 @@ fn init_tracing(
};

// Sized for the per-request child spans (ingest.authenticate / .decode /
// .parse / .accept / .encode_rows / .wal_commit): a request emits ~7 spans,
// not 1. At 2048 the queue held only a few seconds of that volume, so a
// brief export stall silently dropped spans — parents and children alike.
// Head volume is controlled by OTEL_TRACES_SAMPLER / _ARG, which the SDK
// reads via TracerProviderBuilder's derived Default.
// .parse / .accept / .encode_rows / .wal_commit): a sampled request emits ~7
// spans, not 1. At 2048 the queue held only a few seconds of that volume, so
// a brief export stall silently dropped spans, parents and children alike.
// Head volume is INGEST_SELF_TRACE_SAMPLE_RATIO (see `SelfTraceSampler`);
// the queue is sized for the unsampled default.
let batch_config = BatchConfigBuilder::default()
.with_max_queue_size(8192)
.with_max_export_batch_size(1024)
Expand All @@ -1949,6 +1965,7 @@ fn init_tracing(

let provider = SdkTracerProvider::builder()
.with_resource(resource.clone())
.with_sampler(SelfTraceSampler::new(trace_sample_ratio))
.with_span_processor(processor)
.build();
// The runtime argument is not optional here: the runtime-less
Expand Down Expand Up @@ -1989,6 +2006,47 @@ fn init_tracing(
})
}

/// Head sampler for the gateway's own traces: whole traces by trace id, children
/// following their parent. Sampled spans carry `SampleRate` = 1/ratio, which
/// Maple's throughput math weights by, so request rates still read true.
#[derive(Clone, Debug)]
struct SelfTraceSampler {
inner: Sampler,
sample_rate: Option<opentelemetry::KeyValue>,
}

impl SelfTraceSampler {
fn new(ratio: f64) -> Self {
Self {
inner: Sampler::ParentBased(Box::new(Sampler::TraceIdRatioBased(ratio))),
sample_rate: (ratio < 1.0)
.then(|| opentelemetry::KeyValue::new("SampleRate", 1.0 / ratio)),
}
}
}

impl ShouldSample for SelfTraceSampler {
fn should_sample(
&self,
parent_context: Option<&opentelemetry::Context>,
trace_id: opentelemetry::trace::TraceId,
name: &str,
span_kind: &opentelemetry::trace::SpanKind,
attributes: &[opentelemetry::KeyValue],
links: &[opentelemetry::trace::Link],
) -> SamplingResult {
let mut result =
self.inner
.should_sample(parent_context, trace_id, name, span_kind, attributes, links);
if let (SamplingDecision::RecordAndSample, Some(rate)) =
(&result.decision, &self.sample_rate)
{
result.attributes.push(rate.clone());
}
result
}
}

/// Wire up OTLP metric export, mirroring `init_tracing`. The gateway's own
/// operational metrics are pushed to `{forward_endpoint}/v1/metrics` on a
/// periodic interval — the same downstream collector → Tinybird pipeline that
Expand Down Expand Up @@ -2020,6 +2078,7 @@ fn init_metrics(
.with_http()
.with_endpoint(format!("{forward_endpoint}/v1/metrics"))
.with_protocol(Protocol::HttpBinary)
.with_compression(OtlpCompression::Gzip)
.build()
{
Ok(exporter) => exporter,
Expand Down Expand Up @@ -2082,6 +2141,7 @@ fn init_usage_metrics(
.with_http()
.with_endpoint(format!("{forward_endpoint}/v1/metrics"))
.with_protocol(Protocol::HttpBinary)
.with_compression(OtlpCompression::Gzip)
.with_temporality(Temporality::Delta)
.build()
{
Expand Down Expand Up @@ -2145,6 +2205,7 @@ async fn main() {
config.port,
&service_instance_id,
&config.internal_org_id,
config.self_trace_sample_ratio,
);
let meter_provider = init_metrics(
&config.forward_endpoint,
Expand Down Expand Up @@ -6509,6 +6570,25 @@ fn parse_u64(name: &str, raw: Option<String>, default: u64) -> Result<u64, Strin
.map_err(|_| format!("{name} must be a positive integer"))
}

/// A sampling ratio in (0, 1]. Zero is refused rather than read as "off": it
/// would silence the gateway's traces while every other signal keeps flowing.
fn parse_sample_ratio(name: &str, raw: Option<String>, default: f64) -> Result<f64, String> {
let Some(raw) = raw else {
return Ok(default);
};

let value = raw.trim();
if value.is_empty() {
return Ok(default);
}

value
.parse::<f64>()
.ok()
.filter(|ratio| *ratio > 0.0 && *ratio <= 1.0)
.ok_or_else(|| format!("{name} must be a number in (0, 1]"))
}

fn parse_u32(name: &str, raw: Option<String>, default: u32) -> Result<u32, String> {
let Some(raw) = raw else {
return Ok(default);
Expand Down Expand Up @@ -6547,6 +6627,70 @@ mod tests {
use opentelemetry_proto::tonic::metrics::v1::metric;
use std::sync::atomic::AtomicBool;

#[test]
fn self_trace_sample_ratio_must_be_in_the_unit_interval() {
let parse = |raw: &str| parse_sample_ratio("RATIO", Some(raw.to_owned()), 1.0);
assert_eq!(parse_sample_ratio("RATIO", None, 1.0), Ok(1.0));
assert_eq!(parse(" "), Ok(1.0));
assert_eq!(parse("0.05"), Ok(0.05));
assert_eq!(parse("1"), Ok(1.0));
for invalid in ["0", "-0.1", "1.5", "NaN", "inf", "five percent"] {
assert!(parse(invalid).is_err(), "{invalid} must refuse to boot");
}
}

#[test]
fn sampled_self_spans_carry_their_sample_rate() {
use opentelemetry::trace::{
SpanContext, SpanId, SpanKind, TraceContextExt, TraceFlags, TraceId, TraceState,
};

let sample = |sampler: &SelfTraceSampler, parent: Option<&opentelemetry::Context>, id| {
sampler.should_sample(
parent,
TraceId::from(id),
"ingest",
&SpanKind::Server,
&[],
&[],
)
};
let sample_rate = |result: &SamplingResult| {
result
.attributes
.iter()
.find(|kv| kv.key.as_str() == "SampleRate")
.map(|kv| kv.value.clone())
};

let sampled = SelfTraceSampler::new(0.05);
// The ratio sampler keeps a trace when its id's low 64 bits fall under
// ratio * 2^64, so 1 is always kept and u128::MAX never is.
let root = sample(&sampled, None, 1);
assert_eq!(root.decision, SamplingDecision::RecordAndSample);
assert_eq!(sample_rate(&root), Some(opentelemetry::Value::F64(20.0)));
let dropped = sample(&sampled, None, u128::MAX);
assert_eq!(dropped.decision, SamplingDecision::Drop);
assert_eq!(sample_rate(&dropped), None);

// Children follow the root's decision and are weighted the same.
let parent = opentelemetry::Context::new().with_remote_span_context(SpanContext::new(
TraceId::from(u128::MAX),
SpanId::from(7),
TraceFlags::SAMPLED,
false,
TraceState::default(),
));
let child = sample(&sampled, Some(&parent), u128::MAX);
assert_eq!(child.decision, SamplingDecision::RecordAndSample);
assert_eq!(sample_rate(&child), Some(opentelemetry::Value::F64(20.0)));

// Unsampled, nothing is added: a SampleRate of 1 would only cost bytes.
let unsampled = sample(&SelfTraceSampler::new(1.0), None, u128::MAX);
assert_eq!(unsampled.decision, SamplingDecision::RecordAndSample);
assert_eq!(sample_rate(&unsampled), None);
}

#[test]
fn invalid_key_message_points_at_the_other_region() {
assert_eq!(invalid_ingest_key_message(None), Ok("Invalid ingest key"));
Expand Down Expand Up @@ -7558,6 +7702,7 @@ mod tests {
internal_org_id: "org_test_internal".to_owned(),
forward_endpoint,
forward_timeout: Duration::from_secs(5),
self_trace_sample_ratio: 1.0,
write_mode: WriteMode::Forward,
tinybird,
max_request_body_bytes: 1024 * 1024,
Expand Down
2 changes: 1 addition & 1 deletion docs/otel-spec/resource-and-config.md
Original file line number Diff line number Diff line change
Expand Up @@ -126,7 +126,7 @@ Source: https://opentelemetry.io/docs/specs/otel/configuration/sdk-environment-v
| Variable | Default | Notes |
| ------------------------- | ----------------------- | ---------------------------------------------------------------------------------------------------------------------------------------------------------------------- |
| `OTEL_TRACES_SAMPLER` | `parentbased_always_on` | e.g. `parentbased_traceidratio`, `always_on`, `always_off`, `traceidratio`. |
| `OTEL_TRACES_SAMPLER_ARG` | _(unset)_ | Argument shape depends on the sampler (e.g. a ratio `0.1` for `traceidratio`/`parentbased_traceidratio`). `apps/ingest` sets its head sampling through these two vars. |
| `OTEL_TRACES_SAMPLER_ARG` | _(unset)_ | Argument shape depends on the sampler (e.g. a ratio `0.1` for `traceidratio`/`parentbased_traceidratio`). `apps/ingest` ignores both: its head sampling is `INGEST_SELF_TRACE_SAMPLE_RATIO`, which also stamps `SampleRate`. |

### Batch Span Processor (BSP)

Expand Down
9 changes: 9 additions & 0 deletions packages/infra/src/aws/stage.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -12,6 +12,7 @@ import {
resolveIngestDesiredCount,
resolveIngestNamespaceName,
resolveIngestScaling,
resolveIngestSelfTraceSampleRatio,
stageDeploysCollector,
stageDeploysIngest,
stageEnablesReplayBlobs,
Expand Down Expand Up @@ -129,6 +130,14 @@ describe("resolveIngestScaling", () => {
})
})

describe("resolveIngestSelfTraceSampleRatio", () => {
it("samples the gateway's own traces in production only", () => {
expect(resolveIngestSelfTraceSampleRatio(parseMapleStage("prd"))).toBe("0.05")
expect(resolveIngestSelfTraceSampleRatio(parseMapleStage("pr-12"))).toBeUndefined()
expect(resolveIngestSelfTraceSampleRatio(parseMapleStage("dev-alice"))).toBeUndefined()
})
})

describe("parseIngestFleets", () => {
it("is EC2 only when unset", () => {
expect(parseIngestFleets(undefined)).toEqual({ fargate: false, ec2: true })
Expand Down
9 changes: 9 additions & 0 deletions packages/infra/src/aws/stage.ts
Original file line number Diff line number Diff line change
Expand Up @@ -131,6 +131,15 @@ export function resolveIngestTaskSize(stage: MapleStage): IngestTaskSize {
return stage.kind === "prd" ? { cpu: 1024, memory: 2048 } : { cpu: 512, memory: 1024 }
}

/**
* Head-sampling ratio for the gateway's own traces (`INGEST_SELF_TRACE_SAMPLE_RATIO`),
* or `undefined` for the gateway default of 1.0. prd emits ~7 spans per request at
* ~650k req/h; sampled spans carry `SampleRate`, so throughput still reads true.
*/
export function resolveIngestSelfTraceSampleRatio(stage: MapleStage): string | undefined {
return stage.kind === "prd" ? "0.05" : undefined
}

/**
* Which fleets run the gateway. Both can run at once, each behind its own ALB,
* which is how a fleet cutover works: bring the new one up beside the old,
Expand Down
Loading