diff --git a/contrib/envoyproxy/go-control-plane/cmd/serviceextensions/README.md b/contrib/envoyproxy/go-control-plane/cmd/serviceextensions/README.md index 10e67c85220..020e2902b04 100644 --- a/contrib/envoyproxy/go-control-plane/cmd/serviceextensions/README.md +++ b/contrib/envoyproxy/go-control-plane/cmd/serviceextensions/README.md @@ -22,6 +22,18 @@ The ASM Service Extension expose some configuration. The configuration can be tw >**GCP requires that the default configuration for the Service Extension should not change.** +### Forward the source IP attribute + +The Service Extension must be configured to forward the GCP client address attribute: + +```yaml +forwardAttributes: [source.ip] +``` + +The forwarded `source.ip` is authoritative for requests sent directly from the client to the Google Cloud Load Balancer. If `source.ip` is omitted, client IP detection retains its legacy header and connection-address fallback behavior. + +This configuration does not recover the original client address when a third-party CDN is deployed in front of the load balancer. Supporting that topology is outside the scope of this integration; configure and trust the CDN-specific forwarding mechanism separately. + | Environment variable | Default value | Description | |-------------------------------------------|-----------------|---------------------------------------------------------------------------------------------------------------| | `DD_SERVICE_EXTENSION_HOST` | `0.0.0.0` | Host on where the gRPC and HTTP server should listen to. | diff --git a/contrib/envoyproxy/go-control-plane/envoy_messages.go b/contrib/envoyproxy/go-control-plane/envoy_messages.go index c35a24f0084..ae0a6a3b9fc 100644 --- a/contrib/envoyproxy/go-control-plane/envoy_messages.go +++ b/contrib/envoyproxy/go-control-plane/envoy_messages.go @@ -8,6 +8,7 @@ package gocontrolplane import ( "context" "fmt" + "net/netip" "os" "strconv" "sync" @@ -15,6 +16,8 @@ import ( extproc "github.com/envoyproxy/go-control-plane/envoy/service/ext_proc/v3" "google.golang.org/grpc/metadata" + "google.golang.org/protobuf/types/known/structpb" + "github.com/DataDog/dd-trace-go/v2/ddtrace/ext" "github.com/DataDog/dd-trace-go/v2/ddtrace/tracer" "github.com/DataDog/dd-trace-go/v2/instrumentation/appsec/proxy" @@ -55,6 +58,40 @@ func (m messageRequestHeaders) ExtractRequest(ctx context.Context) (proxy.Pseudo }, nil } +const ( + gcpServiceExtensionAttributesNamespace = "envoy.filters.http.ext_proc" + gcpServiceExtensionSourceIPAttribute = "source.ip" +) + +// ClientIPOverride returns the authoritative GCP Service Extension source IP when present. +func (m messageRequestHeaders) ClientIPOverride(ctx context.Context) (netip.Addr, bool) { + if m.component(ctx) != componentNameGCPServiceExtension { + return netip.Addr{}, false + } + + namespace, ok := m.ProcessingRequest.GetAttributes()[gcpServiceExtensionAttributesNamespace] + if !ok || namespace == nil { + return netip.Addr{}, false + } + value, ok := namespace.GetFields()[gcpServiceExtensionSourceIPAttribute] + if !ok { + return netip.Addr{}, false + } + + if value == nil { + return netip.Addr{}, true + } + stringValue, ok := value.GetKind().(*structpb.Value_StringValue) + if !ok { + return netip.Addr{}, true + } + ip, err := netip.ParseAddr(stringValue.StringValue) + if err != nil || ip.Zone() != "" { + return netip.Addr{}, true + } + return ip.Unmap(), true +} + func (m messageRequestHeaders) MessageType() proxy.MessageType { return proxy.MessageTypeRequestHeaders } diff --git a/contrib/envoyproxy/go-control-plane/go.mod b/contrib/envoyproxy/go-control-plane/go.mod index e708047e596..dfa42cc8daf 100644 --- a/contrib/envoyproxy/go-control-plane/go.mod +++ b/contrib/envoyproxy/go-control-plane/go.mod @@ -8,6 +8,7 @@ require ( github.com/stretchr/testify v1.11.1 golang.org/x/sync v0.20.0 google.golang.org/grpc v1.80.0 + google.golang.org/protobuf v1.36.11 ) require ( @@ -82,7 +83,6 @@ require ( golang.org/x/time v0.15.0 // indirect golang.org/x/xerrors v0.0.0-20240903120638-7835f813f4da // indirect google.golang.org/genproto/googleapis/rpc v0.0.0-20260406210006-6f92a3bedf2d // indirect - google.golang.org/protobuf v1.36.11 // indirect gopkg.in/yaml.v3 v3.0.1 // indirect ) diff --git a/contrib/envoyproxy/go-control-plane/source_ip_test.go b/contrib/envoyproxy/go-control-plane/source_ip_test.go new file mode 100644 index 00000000000..6fbdd8f1f3c --- /dev/null +++ b/contrib/envoyproxy/go-control-plane/source_ip_test.go @@ -0,0 +1,332 @@ +// Unless explicitly stated otherwise all files in this repository are licensed +// under the Apache License Version 2.0. +// This product includes software developed at Datadog (https://www.datadoghq.com/). +// Copyright 2026 Datadog, Inc. + +package gocontrolplane + +import ( + "context" + "io" + "maps" + "testing" + + envoyextproc "github.com/envoyproxy/go-control-plane/envoy/service/ext_proc/v3" + "github.com/stretchr/testify/require" + "google.golang.org/grpc/metadata" + "google.golang.org/protobuf/types/known/structpb" + + "github.com/DataDog/dd-trace-go/v2/ddtrace/ext" + "github.com/DataDog/dd-trace-go/v2/ddtrace/mocktracer" + "github.com/DataDog/dd-trace-go/v2/instrumentation/httptrace" + "github.com/DataDog/dd-trace-go/v2/instrumentation/testutils" +) + +const ( + gcpServiceExtensionAttributeNamespace = "envoy.filters.http.ext_proc" + forgedXForwardedFor = "198.51.100.42" +) + +func TestGCPServiceExtensionSourceIPAuthoritative(t *testing.T) { + t.Setenv("DD_APPSEC_RULES", "../../../internal/appsec/testdata/user_rules.json") + t.Setenv("DD_APPSEC_WAF_TIMEOUT", "10ms") + t.Cleanup(httptrace.ResetCfg) + testutils.StartAppSec(t) + httptrace.ResetCfg() + + t.Run("valid", func(t *testing.T) { + tests := []struct { + name string + sourceIP string + canonical string + }{ + {name: "public IPv4", sourceIP: "203.0.113.10", canonical: "203.0.113.10"}, + {name: "private IPv4", sourceIP: "10.20.30.40", canonical: "10.20.30.40"}, + {name: "IPv6", sourceIP: "2001:0db8:0000:0000:0000:0000:0000:0001", canonical: "2001:db8::1"}, + {name: "IPv4-mapped IPv6", sourceIP: "::ffff:203.0.113.11", canonical: "203.0.113.11"}, + } + + for _, tc := range tests { + t.Run(tc.name, func(t *testing.T) { + _, tags := runSourceIPRequest(t, GCPServiceExtensionIntegration, nil, sourceIPAttributes(structpb.NewStringValue(tc.sourceIP))) + + require.Equal(t, tc.canonical, tags[ext.NetworkClientIP]) + require.Equal(t, tc.canonical, tags[ext.HTTPClientIP]) + require.Equal(t, forgedXForwardedFor, tags["http.request.headers.x-forwarded-for"]) + }) + } + }) + + t.Run("present invalid suppresses XFF", func(t *testing.T) { + tests := []struct { + name string + value *structpb.Value + }{ + {name: "nil", value: nil}, + {name: "empty", value: structpb.NewStringValue("")}, + {name: "malformed", value: structpb.NewStringValue("not-an-ip")}, + {name: "scoped IPv6", value: structpb.NewStringValue("fe80::1%eth0")}, + {name: "number", value: structpb.NewNumberValue(42)}, + {name: "boolean", value: structpb.NewBoolValue(true)}, + {name: "null", value: structpb.NewNullValue()}, + {name: "struct", value: structpb.NewStructValue(&structpb.Struct{})}, + {name: "list", value: structpb.NewListValue(&structpb.ListValue{})}, + } + + for _, tc := range tests { + t.Run(tc.name, func(t *testing.T) { + _, tags := runSourceIPRequest(t, GCPServiceExtensionIntegration, nil, sourceIPAttributes(tc.value)) + + require.NotContains(t, tags, ext.NetworkClientIP) + require.NotContains(t, tags, ext.HTTPClientIP) + require.Equal(t, forgedXForwardedFor, tags["http.request.headers.x-forwarded-for"]) + }) + } + }) + + t.Run("absent preserves legacy resolution", func(t *testing.T) { + tests := []struct { + name string + attributes map[string]*structpb.Struct + }{ + {name: "missing namespace"}, + {name: "missing field", attributes: map[string]*structpb.Struct{ + gcpServiceExtensionAttributeNamespace: {}, + }}, + } + + for _, tc := range tests { + t.Run(tc.name, func(t *testing.T) { + _, tags := runSourceIPRequest(t, GCPServiceExtensionIntegration, nil, tc.attributes) + + require.NotContains(t, tags, ext.NetworkClientIP) + require.Equal(t, forgedXForwardedFor, tags[ext.HTTPClientIP]) + require.Equal(t, forgedXForwardedFor, tags["http.request.headers.x-forwarded-for"]) + }) + } + }) + + t.Run("ignored outside effective GCP", func(t *testing.T) { + tests := []struct { + name string + integration Integration + metadata metadata.MD + }{ + {name: "Envoy integration", integration: EnvoyIntegration}, + {name: "Envoy Gateway integration", integration: EnvoyGatewayIntegration}, + {name: "Istio integration", integration: IstioIntegration}, + {name: "Envoy metadata override", integration: GCPServiceExtensionIntegration, metadata: metadata.Pairs(datadogEnvoyIntegrationHeader, "1")}, + {name: "Istio metadata override", integration: GCPServiceExtensionIntegration, metadata: metadata.Pairs(datadogIntegrationHeader, "1")}, + } + + for _, tc := range tests { + t.Run(tc.name, func(t *testing.T) { + _, tags := runSourceIPRequest(t, tc.integration, tc.metadata, sourceIPAttributes(structpb.NewStringValue("203.0.113.12"))) + + require.NotContains(t, tags, ext.NetworkClientIP) + require.Equal(t, forgedXForwardedFor, tags[ext.HTTPClientIP]) + require.Equal(t, forgedXForwardedFor, tags["http.request.headers.x-forwarded-for"]) + }) + } + }) +} + +func TestMessageRequestHeadersIgnoresSourceIPForNonGCP(t *testing.T) { + tests := []struct { + name string + integration Integration + }{ + {name: "Envoy", integration: EnvoyIntegration}, + {name: "Envoy Gateway", integration: EnvoyGatewayIntegration}, + {name: "Istio", integration: IstioIntegration}, + } + + for _, tc := range tests { + t.Run(tc.name, func(t *testing.T) { + processingRequest := &envoyextproc.ProcessingRequest{ + Attributes: sourceIPAttributes(structpb.NewStringValue("203.0.113.14")), + } + headers := &envoyextproc.HttpHeaders{ + Headers: makeRequestHeaders(t, nil, "GET", "/"), + } + + message := messageRequestHeaders{ + ProcessingRequest: processingRequest, + HttpHeaders: headers, + integration: tc.integration, + } + ip, set := message.ClientIPOverride(context.Background()) + + require.False(t, set) + require.False(t, ip.IsValid()) + }) + } +} + +func TestGCPServiceExtensionSourceIPIsWAFIdentity(t *testing.T) { + t.Setenv("DD_APPSEC_RULES", "../../../internal/appsec/testdata/user_rules.json") + t.Setenv("DD_APPSEC_WAF_TIMEOUT", "10ms") + t.Cleanup(httptrace.ResetCfg) + testutils.StartAppSec(t) + httptrace.ResetCfg() + + response, tags := runSourceIPRequest(t, GCPServiceExtensionIntegration, nil, sourceIPAttributes(structpb.NewStringValue("111.222.111.222"))) + + require.IsType(t, &envoyextproc.ProcessingResponse_ImmediateResponse{}, response.GetResponse()) + require.Equal(t, "111.222.111.222", tags[ext.NetworkClientIP]) + require.Equal(t, "111.222.111.222", tags[ext.HTTPClientIP]) + require.Equal(t, forgedXForwardedFor, tags["http.request.headers.x-forwarded-for"]) + require.Equal(t, "true", tags["appsec.blocked"]) +} + +func TestGCPServiceExtensionSourceIPSuppressesBlockedXFF(t *testing.T) { + t.Setenv("DD_APPSEC_RULES", "../../../internal/appsec/testdata/user_rules.json") + t.Setenv("DD_APPSEC_WAF_TIMEOUT", "10ms") + t.Cleanup(httptrace.ResetCfg) + testutils.StartAppSec(t) + httptrace.ResetCfg() + + const blockedXForwardedFor = "111.222.111.222" + tests := []struct { + name string + value *structpb.Value + expectedIP string + }{ + {name: "valid safe source", value: structpb.NewStringValue("203.0.113.15"), expectedIP: "203.0.113.15"}, + {name: "present invalid source", value: structpb.NewStringValue("not-an-ip")}, + } + + for _, tc := range tests { + t.Run(tc.name, func(t *testing.T) { + response, tags := runSourceIPRequestWithXFF(t, GCPServiceExtensionIntegration, nil, sourceIPAttributes(tc.value), blockedXForwardedFor) + + require.IsType(t, &envoyextproc.ProcessingResponse_RequestHeaders{}, response.GetResponse()) + require.Nil(t, response.GetImmediateResponse()) + require.NotContains(t, tags, "appsec.blocked") + require.Equal(t, blockedXForwardedFor, tags["http.request.headers.x-forwarded-for"]) + if tc.expectedIP == "" { + require.NotContains(t, tags, ext.NetworkClientIP) + require.NotContains(t, tags, ext.HTTPClientIP) + return + } + require.Equal(t, tc.expectedIP, tags[ext.NetworkClientIP]) + require.Equal(t, tc.expectedIP, tags[ext.HTTPClientIP]) + }) + } +} + +func TestGCPServiceExtensionSourceIPCollectionActivation(t *testing.T) { + t.Cleanup(httptrace.ResetCfg) + + tests := []struct { + name string + clientIPEnabled string + expectedSourceTag bool + }{ + {name: "client IP collection enabled", clientIPEnabled: "true", expectedSourceTag: true}, + {name: "AppSec and client IP collection disabled", clientIPEnabled: "false", expectedSourceTag: false}, + } + + for _, tc := range tests { + t.Run(tc.name, func(t *testing.T) { + t.Setenv("DD_APPSEC_ENABLED", "false") + t.Setenv("DD_TRACE_CLIENT_IP_ENABLED", tc.clientIPEnabled) + httptrace.ResetCfg() + + _, tags := runSourceIPRequest(t, GCPServiceExtensionIntegration, nil, sourceIPAttributes(structpb.NewStringValue("203.0.113.13"))) + + if tc.expectedSourceTag { + require.Equal(t, "203.0.113.13", tags[ext.NetworkClientIP]) + require.Equal(t, "203.0.113.13", tags[ext.HTTPClientIP]) + } else { + require.NotContains(t, tags, ext.NetworkClientIP) + require.NotContains(t, tags, ext.HTTPClientIP) + } + require.NotContains(t, tags, "http.request.headers.x-forwarded-for") + }) + } +} + +func TestGCPServiceExtensionSourceIPAbsentUsesMetadataConnectionAddress(t *testing.T) { + t.Setenv("DD_APPSEC_ENABLED", "false") + t.Setenv("DD_TRACE_CLIENT_IP_ENABLED", "true") + t.Cleanup(httptrace.ResetCfg) + httptrace.ResetCfg() + + const metadataAddress = "192.0.2.20" + _, tags := runSourceIPRequestWithXFF( + t, + GCPServiceExtensionIntegration, + metadata.Pairs("x-forwarded-for", metadataAddress), + nil, + "", + ) + + require.Equal(t, metadataAddress, tags[ext.NetworkClientIP]) + require.Equal(t, metadataAddress, tags[ext.HTTPClientIP]) + require.NotContains(t, tags, "http.request.headers.x-forwarded-for") +} + +func sourceIPAttributes(value *structpb.Value) map[string]*structpb.Struct { + return map[string]*structpb.Struct{ + gcpServiceExtensionAttributeNamespace: { + Fields: map[string]*structpb.Value{"source.ip": value}, + }, + } +} + +func runSourceIPRequest(t *testing.T, integration Integration, md metadata.MD, attributes map[string]*structpb.Struct) (*envoyextproc.ProcessingResponse, map[string]any) { + t.Helper() + return runSourceIPRequestWithXFF(t, integration, md, attributes, forgedXForwardedFor) +} + +func runSourceIPRequestWithXFF(t *testing.T, integration Integration, md metadata.MD, attributes map[string]*structpb.Struct, xForwardedFor string) (*envoyextproc.ProcessingResponse, map[string]any) { + t.Helper() + var headers map[string]string + if xForwardedFor != "" { + headers = map[string]string{"X-Forwarded-For": xForwardedFor} + } + + rig, err := newEnvoyAppsecRig(t, integration, false, nil) + require.NoError(t, err) + defer rig.Close() + + mt := mocktracer.Start() + defer mt.Stop() + + ctx := context.Background() + if md != nil { + ctx = metadata.NewOutgoingContext(ctx, md) + } + stream, err := rig.client.Process(ctx) + require.NoError(t, err) + + err = stream.Send(&envoyextproc.ProcessingRequest{ + Attributes: attributes, + Request: &envoyextproc.ProcessingRequest_RequestHeaders{ + RequestHeaders: &envoyextproc.HttpHeaders{ + Headers: makeRequestHeaders(t, headers, "GET", "/"), + EndOfStream: true, + }, + }, + }) + require.NoError(t, err) + + response, err := stream.Recv() + require.NoError(t, err) + if response.GetImmediateResponse() == nil { + require.NotNil(t, response.GetRequestHeaders()) + sendProcessingResponseHeaders(t, stream, nil, "200", false) + _, err = stream.Recv() + require.ErrorIs(t, err, io.EOF) + } + + require.NoError(t, stream.CloseSend()) + _, _ = stream.Recv() + + finished := mt.FinishedSpans() + require.Len(t, finished, 1) + tags := make(map[string]any, len(finished[0].Tags())) + maps.Copy(tags, finished[0].Tags()) + return response, tags +} diff --git a/instrumentation/appsec/emitter/httpsec/http.go b/instrumentation/appsec/emitter/httpsec/http.go index f4ecaaa526b..33246f2f0e1 100644 --- a/instrumentation/appsec/emitter/httpsec/http.go +++ b/instrumentation/appsec/emitter/httpsec/http.go @@ -12,6 +12,7 @@ package httpsec import ( "context" + "net/netip" "sync" // Blank import needed to use embed for the default blocked response payloads @@ -48,6 +49,9 @@ type ( method string // route is the HTTP route for the current handler operation (or the URL if no route is available). route string + // clientIPOverride is an optional authoritative client IP resolution result. + clientIPOverride netip.Addr + clientIPOverrideSet bool // downstreamRequestBodyAnalysis is the number of times a call to a downstream request body monitoring function was made. downstreamRequestBodyAnalysis atomic.Int32 @@ -88,22 +92,50 @@ type ( EarlyBlock struct{} ) +type clientIPOverrideContextKey struct{} + +type clientIPOverrideContextValue struct { + ip netip.Addr + set bool +} + +// ContextWithClientIPOverride returns a context carrying an authoritative client IP resolution result. +// The set value distinguishes an absent override from a present but invalid IP. +func ContextWithClientIPOverride(ctx context.Context, ip netip.Addr, set bool) context.Context { + if !set { + return ctx + } + return context.WithValue(ctx, clientIPOverrideContextKey{}, clientIPOverrideContextValue{ip: ip, set: set}) +} + +// ClientIPOverrideFromContext returns the authoritative client IP resolution result carried by ctx. +func ClientIPOverrideFromContext(ctx context.Context) (netip.Addr, bool) { + value, ok := ctx.Value(clientIPOverrideContextKey{}).(clientIPOverrideContextValue) + if !ok { + return netip.Addr{}, false + } + return value.ip, value.set +} + func (HandlerOperationArgs) IsArgOf(*HandlerOperation) {} func (HandlerOperationRes) IsResultOf(*HandlerOperation) {} func StartOperation(ctx context.Context, args HandlerOperationArgs, span trace.TagSetter) (*HandlerOperation, *atomic.Pointer[actions.BlockHTTP], context.Context) { + clientIPOverride, clientIPOverrideSet := ClientIPOverrideFromContext(ctx) wafOp, found := dyngo.FindOperation[waf.ContextOperation](ctx) if !found { wafOp, ctx = waf.StartContextOperation(ctx, span) } op := &HandlerOperation{ - Operation: dyngo.NewOperation(wafOp), - ContextOperation: wafOp, - wafContextOwner: !found, // If we started the parent operation, we finish it, otherwise we don't - framework: args.Framework, - method: args.Method, - route: args.RequestRoute, + Operation: dyngo.NewOperation(wafOp), + ContextOperation: wafOp, + wafContextOwner: !found, // If we started the parent operation, we finish it, otherwise we don't + framework: args.Framework, + method: args.Method, + route: args.RequestRoute, + clientIPOverride: clientIPOverride, + clientIPOverrideSet: clientIPOverrideSet, } // We need to use an atomic pointer to store the action because the action may be created asynchronously in the future @@ -140,6 +172,11 @@ func (op *HandlerOperation) Route() string { return op.route } +// ClientIPOverride returns the authoritative client IP resolution result for this operation. +func (op *HandlerOperation) ClientIPOverride() (netip.Addr, bool) { + return op.clientIPOverride, op.clientIPOverrideSet +} + // DownstreamRequestBodyAnalysis returns the number of times a call to a downstream request body monitoring function was made. func (op *HandlerOperation) DownstreamRequestBodyAnalysis() int { return int(op.downstreamRequestBodyAnalysis.Load()) diff --git a/instrumentation/appsec/proxy/message_processor.go b/instrumentation/appsec/proxy/message_processor.go index 4535edf0c95..8818aaa7044 100644 --- a/instrumentation/appsec/proxy/message_processor.go +++ b/instrumentation/appsec/proxy/message_processor.go @@ -9,18 +9,24 @@ import ( "context" "fmt" "io" + "net/netip" "sync" "sync/atomic" "github.com/DataDog/dd-trace-go/v2/appsec" "github.com/DataDog/dd-trace-go/v2/instrumentation" "github.com/DataDog/dd-trace-go/v2/instrumentation/appsec/dyngo" + httpsecemitter "github.com/DataDog/dd-trace-go/v2/instrumentation/appsec/emitter/httpsec" "github.com/DataDog/dd-trace-go/v2/instrumentation/appsec/emitter/waf/actions" "github.com/DataDog/dd-trace-go/v2/instrumentation/httptrace" "github.com/DataDog/dd-trace-go/v2/internal/appsec/body" "github.com/DataDog/dd-trace-go/v2/internal/appsec/body/json" ) +type clientIPOverrideProvider interface { + ClientIPOverride(context.Context) (netip.Addr, bool) +} + // Processor is a state machine that handles incoming HTTP request and response in a streaming manner, // made for proxy external-processing protocols like Envoy's External Processing or HAProxy's SPOP. // @@ -84,6 +90,11 @@ func (mp *Processor) OnRequestHeaders(ctx context.Context, req RequestHeaders) ( return reqState, fmt.Errorf("error extracting request header from input message: %w", err) } + if provider, ok := req.(clientIPOverrideProvider); ok { + ip, set := provider.ClientIPOverride(ctx) + ctx = httpsecemitter.ContextWithClientIPOverride(ctx, ip, set) + } + httpRequest, err := pseudoRequest.toNetHTTP(ctx) if err != nil { return reqState, fmt.Errorf("error converting to net/http request: %w", err) diff --git a/instrumentation/httptrace/httptrace.go b/instrumentation/httptrace/httptrace.go index 1fdaf821216..41613ad90e0 100644 --- a/instrumentation/httptrace/httptrace.go +++ b/instrumentation/httptrace/httptrace.go @@ -19,6 +19,7 @@ import ( "github.com/DataDog/dd-trace-go/v2/ddtrace/ext" "github.com/DataDog/dd-trace-go/v2/ddtrace/tracer" "github.com/DataDog/dd-trace-go/v2/instrumentation" + emitterhttpsec "github.com/DataDog/dd-trace-go/v2/instrumentation/appsec/emitter/httpsec" appsechttpsec "github.com/DataDog/dd-trace-go/v2/instrumentation/appsec/httpsec" listenerhttpsec "github.com/DataDog/dd-trace-go/v2/internal/appsec/listener/httpsec" "github.com/DataDog/dd-trace-go/v2/internal/log" @@ -68,7 +69,8 @@ func StartRequestSpan(r *http.Request, opts ...tracer.StartSpanOption) (*tracer. var ipTags map[string]string if cfg.traceClientIP { - ipTags, _ = listenerhttpsec.ClientIPTags(r.Header, true, r.RemoteAddr) + ip, set := emitterhttpsec.ClientIPOverrideFromContext(r.Context()) + ipTags, _ = listenerhttpsec.ClientIPTagsWithOverride(r.Header, true, r.RemoteAddr, ip, set) } var inferredProxySpan *tracer.Span diff --git a/internal/appsec/listener/httpsec/http.go b/internal/appsec/listener/httpsec/http.go index c45fa0ebb94..9c425188195 100644 --- a/internal/appsec/listener/httpsec/http.go +++ b/internal/appsec/listener/httpsec/http.go @@ -140,7 +140,8 @@ func (*HeaderExtractionFeature) OnResponse(op *httpsec.HandlerOperation, resp ht } func extractRequestHeaders(op *httpsec.HandlerOperation, args httpsec.HandlerOperationArgs) (map[string][]string, netip.Addr) { - tags, ip := ClientIPTags(args.Headers, true, args.RemoteAddr) + override, overrideSet := op.ClientIPOverride() + tags, ip := ClientIPTagsWithOverride(args.Headers, true, args.RemoteAddr, override, overrideSet) op.SetStringTags(tags) headers := headersRemoveCookies(args.Headers) diff --git a/internal/appsec/listener/httpsec/request.go b/internal/appsec/listener/httpsec/request.go index f4d3d78e45f..a865b2b039a 100644 --- a/internal/appsec/listener/httpsec/request.go +++ b/internal/appsec/listener/httpsec/request.go @@ -103,6 +103,15 @@ func ClientIPTags(headers map[string][]string, hasCanonicalHeaders bool, remoteA return ClientIPTagsFor(remoteIP, clientIP), clientIP } +// ClientIPTagsWithOverride resolves client IP tags using an authoritative override when set. +// A set but invalid override intentionally suppresses legacy header and remote-address resolution. +func ClientIPTagsWithOverride(headers map[string][]string, hasCanonicalHeaders bool, remoteAddr string, override netip.Addr, overrideSet bool) (tags map[string]string, clientIP netip.Addr) { + if overrideSet { + return ClientIPTagsFor(override, override), override + } + return ClientIPTags(headers, hasCanonicalHeaders, remoteAddr) +} + func ClientIPTagsFor(remoteIP netip.Addr, clientIP netip.Addr) map[string]string { remoteIPValid := remoteIP.IsValid() clientIPValid := clientIP.IsValid()