diff --git a/apps/ingest/Cargo.lock b/apps/ingest/Cargo.lock index 8c01ff6ff3..da56ce3ab9 100644 --- a/apps/ingest/Cargo.lock +++ b/apps/ingest/Cargo.lock @@ -1811,6 +1811,7 @@ version = "0.32.0" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "9966929966d17620d7c316c643ba62631826e10021409357772d5eea84f62c35" dependencies = [ + "flate2", "http", "opentelemetry", "opentelemetry-http", diff --git a/apps/ingest/Cargo.toml b/apps/ingest/Cargo.toml index 7e75b968ce..a04dc0151e 100644 --- a/apps/ingest/Cargo.toml +++ b/apps/ingest/Cargo.toml @@ -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", diff --git a/apps/ingest/alchemy.run.ts b/apps/ingest/alchemy.run.ts index cbbf22d427..056dcc79f8 100644 --- a/apps/ingest/alchemy.run.ts +++ b/apps/ingest/alchemy.run.ts @@ -23,6 +23,7 @@ import { resolveIngestDesiredCount, resolveIngestNamespaceName, resolveIngestScaling, + resolveIngestSelfTraceSampleRatio, resolveIngestTaskSize, stageDeploysCollector, stageEnablesReplayBlobs, @@ -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) @@ -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 diff --git a/apps/ingest/src/main.rs b/apps/ingest/src/main.rs index b5423ad1d6..d6bb773272 100644 --- a/apps/ingest/src/main.rs +++ b/apps/ingest/src/main.rs @@ -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; @@ -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}; @@ -148,6 +153,8 @@ struct AppConfig { otlp_grpc_port: Option, 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, @@ -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() @@ -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, @@ -1849,6 +1862,7 @@ fn init_tracing( bind_port: u16, service_instance_id: &str, internal_org_id: &str, + trace_sample_ratio: f64, ) -> Option { // Two filters on purpose. The registry-wide filter is pinned at `info` // because every gateway span (`ingest`, `ingest.authenticate`, the Postgres @@ -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, @@ -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, @@ -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) @@ -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 @@ -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, +} + +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 @@ -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, @@ -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() { @@ -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, @@ -6509,6 +6570,25 @@ fn parse_u64(name: &str, raw: Option, default: u64) -> Result, default: f64) -> Result { + let Some(raw) = raw else { + return Ok(default); + }; + + let value = raw.trim(); + if value.is_empty() { + return Ok(default); + } + + value + .parse::() + .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, default: u32) -> Result { let Some(raw) = raw else { return Ok(default); @@ -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")); @@ -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, diff --git a/docs/otel-spec/resource-and-config.md b/docs/otel-spec/resource-and-config.md index 852c1d5419..9e521081ea 100644 --- a/docs/otel-spec/resource-and-config.md +++ b/docs/otel-spec/resource-and-config.md @@ -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) diff --git a/packages/infra/src/aws/stage.test.ts b/packages/infra/src/aws/stage.test.ts index 2a7404d831..f60b321142 100644 --- a/packages/infra/src/aws/stage.test.ts +++ b/packages/infra/src/aws/stage.test.ts @@ -12,6 +12,7 @@ import { resolveIngestDesiredCount, resolveIngestNamespaceName, resolveIngestScaling, + resolveIngestSelfTraceSampleRatio, stageDeploysCollector, stageDeploysIngest, stageEnablesReplayBlobs, @@ -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 }) diff --git a/packages/infra/src/aws/stage.ts b/packages/infra/src/aws/stage.ts index c716f1db64..164dbf14b8 100644 --- a/packages/infra/src/aws/stage.ts +++ b/packages/infra/src/aws/stage.ts @@ -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,