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
4 changes: 3 additions & 1 deletion comp/core/configstreamconsumer/impl/consumer.go
Original file line number Diff line number Diff line change
Expand Up @@ -589,8 +589,10 @@ func (c *consumer) handleConfigEvent(event *pb.ConfigEvent) error {
// After applySnapshot's retraction loop, so a remapped value is not retracted out from under itself.
c.applyOverrides()
if snapshotApplied {
// Signalled last: waitForReady must not release before the remap has folded in the override.
// After applyOverrides: waitForReady must not release before the remap has folded in the override.
c.markReady()
// Diagnostic only, so it trails readiness; after applyOverrides, or a remap reads as a loss.
configstreambootstrap.ReportDroppedEnvOverrides(c.params.ClientName)
}
return nil
}
Expand Down
53 changes: 53 additions & 0 deletions comp/core/configstreamconsumer/impl/integration_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -21,6 +21,7 @@ import (
"net"
"os"
"path/filepath"
"slices"
"sync"
"testing"
"time"
Expand Down Expand Up @@ -789,3 +790,55 @@ func TestLogLevelOverrideIsPerClient(t *testing.T) {
t.Fatal("OneShot did not complete")
}
}

func TestRemappedKeyIsNotReportedAsADroppedEnvVar(t *testing.T) {
t.Setenv("DD_LOG_LEVEL", "debug")
t.Setenv("DD_SITE", "datadoghq.eu")
configstreambootstrap.ResetGlobalConfig(t)
configstreambootstrap.UseDynamicSchema(t)

dir := t.TempDir()
addr, mock, cleanup := setupFakeCoreAgent(t, dir)
defer cleanup()

datadogPath := overrideTestConfig(t, dir, addr)

opts := fx.Options(
fx.Provide(func() log.Component { return logmock.New(t) }),
telemetryfx.Module(),
fx.Supply(configstreamconsumer.NewParams("security-agent", datadogPath, configstreamconsumer.WithReadyTimeout(10*time.Second))),
configstreamconsumerfx.Module(),
)

testRun := func(_ configstreamconsumer.Component) error {
// The report trails readiness, so it may land just after OneShot hands control back.
require.Eventually(t, func() bool {
return slices.Contains(configstreambootstrap.LastEnvOverrideReport(), "site (DD_SITE)")
}, 10*time.Second, 20*time.Millisecond, "a setting the core Agent never streamed must be reported")
require.NotContains(t, configstreambootstrap.LastEnvOverrideReport(), "log_level (DD_LOG_LEVEL)",
"the per-agent remap reproduced the local value, so nothing was lost")
return nil
}

done := make(chan error, 1)
go func() { done <- fxutil.OneShot(testRun, opts) }()

mock.events <- &pb.ConfigEvent{
Event: &pb.ConfigEvent_Snapshot{
Snapshot: &pb.ConfigSnapshot{
SequenceId: 1,
Settings: []*pb.ConfigSetting{
{Key: "log_level", Value: mustNewValue(t, "info"), Source: string(model.SourceFile)},
{Key: "security_agent.log_level", Value: mustNewValue(t, "debug"), Source: string(model.SourceFile)},
},
},
},
}

select {
case err := <-done:
require.NoError(t, err)
case <-time.After(30 * time.Second):
t.Fatal("OneShot did not complete")
}
}
11 changes: 11 additions & 0 deletions pkg/config/nodetreemodel/getter.go
Original file line number Diff line number Diff line change
Expand Up @@ -40,6 +40,17 @@ func (c *ntmConfig) GetEnvVars() []string {
return vars
}

// ConfigEnvVars returns a copy of the env vars bound to each config key, highest priority first.
func (c *ntmConfig) ConfigEnvVars() map[string][]string {
c.RLock()
defer c.RUnlock()
out := make(map[string][]string, len(c.configEnvVars))
for key, envVars := range c.configEnvVars {
out[key] = slices.Clone(envVars)
}
return out
}

// GetProxies returns the proxy settings from the configuration
func (c *ntmConfig) GetProxies() *model.Proxy {
c.Lock()
Expand Down
1 change: 1 addition & 0 deletions pkg/configstreambootstrap/BUILD.bazel
Original file line number Diff line number Diff line change
Expand Up @@ -23,6 +23,7 @@ dd_agent_go_test(
srcs = ["bootstrap_test.go"],
embed = [":configstreambootstrap"],
deps = [
"//pkg/config/model",
"//pkg/config/setup",
"@com_github_stretchr_testify//require",
],
Expand Down
100 changes: 97 additions & 3 deletions pkg/configstreambootstrap/bootstrap.go
Original file line number Diff line number Diff line change
Expand Up @@ -9,6 +9,13 @@
package configstreambootstrap

import (
"fmt"
"os"
"reflect"
"slices"
"strings"
"sync"

pkgtoken "github.com/DataDog/datadog-agent/pkg/api/security"
"github.com/DataDog/datadog-agent/pkg/api/security/cert"
pkgconfigmodel "github.com/DataDog/datadog-agent/pkg/config/model"
Expand Down Expand Up @@ -66,16 +73,103 @@ func SeedGlobalBuilder(s Settings, configFile string) {
cert.PersistCertFilepath(b)
}

// envOverride is a setting the local env layer was deciding the value of before the wipe.
type envOverride struct {
key string
envVar string
value any
}

var (
// DisableLocalEnvLayer runs on the main goroutine at startup and the report runs on the
// consumer's stream goroutine; the ordering is not enforced by this package's API, so lock.
envOverridesMu sync.Mutex
capturedEnvOverrides []envOverride
lastEnvOverrideReport []string
)

// DisableLocalEnvLayer drops the env layer (nodetreemodel only) so local DD_* vars
// can't override streamed values. Viper-backed configs cannot clear env vars.
func DisableLocalEnvLayer(clientName string) {
b := pkgconfigsetup.Datadog()
type envVarClearer interface{ ClearEnvVars() }
if clearer, ok := b.(envVarClearer); ok {
clearer.ClearEnvVars()
pkglog.Infof("configstreamconsumer[%s]: local env-var layer disabled", clientName)
clearer, ok := b.(envVarClearer)
if !ok {
return
}
type configEnvVarLister interface{ ConfigEnvVars() map[string][]string }
if lister, ok := b.(configEnvVarLister); ok {
envOverridesMu.Lock()
capturedEnvOverrides = captureEnvOverrides(b, lister.ConfigEnvVars())
envOverridesMu.Unlock()
}
clearer.ClearEnvVars()
pkglog.Infof("configstreamconsumer[%s]: local env-var layer disabled", clientName)
}

// captureEnvOverrides records the settings the env layer is actually deciding. A key whose env var
// is set but loses to a higher-precedence source was not being overridden by the env, so it is skipped.
func captureEnvOverrides(cfg pkgconfigmodel.Reader, configEnvVars map[string][]string) []envOverride {
captured := make([]envOverride, 0, len(configEnvVars))
for key, envVars := range configEnvVars {
if cfg.GetSource(key) != pkgconfigmodel.SourceEnvVar {
continue
}
if name := winningEnvVar(envVars); name != "" {
captured = append(captured, envOverride{key: key, envVar: name, value: cfg.Get(key)})
}
}
slices.SortFunc(captured, func(a, b envOverride) int { return strings.Compare(a.key, b.key) })
return captured
}

// winningEnvVar mirrors nodetreemodel.buildEnvVars: the first var that is set and non-empty wins.
func winningEnvVar(envVars []string) string {
for _, name := range envVars {
if value, isSet := os.LookupEnv(name); isSet && value != "" {
return name
}
}
return ""
}

// ReportDroppedEnvOverrides warns about settings the stream did not reproduce. It must run after
// the first snapshot and after any post-snapshot remapping, and consumes the captured state so
// that later value changes, which are ordinary operation, never warn.
func ReportDroppedEnvOverrides(clientName string) {
envOverridesMu.Lock()
captured := capturedEnvOverrides
capturedEnvOverrides = nil
envOverridesMu.Unlock()
if len(captured) == 0 {
return
}

dropped := diffEnvOverrides(pkgconfigsetup.Datadog(), captured)

envOverridesMu.Lock()
lastEnvOverrideReport = dropped
envOverridesMu.Unlock()

if len(dropped) == 0 {
return
}
pkglog.Warnf("configstreamconsumer[%s]: these settings were set by DD_* env vars on this process and the core Agent streamed a different value, so the local value is lost; set them on the core Agent instead: %s",
clientName, strings.Join(dropped, ", "))
}

// diffEnvOverrides names the captured settings whose current value differs from the pre-wipe one.
// Names only, never values: several of these settings are credentials, and "differs" is the signal.
func diffEnvOverrides(cfg pkgconfigmodel.Reader, captured []envOverride) []string {
var dropped []string
for _, o := range captured {
// Values can be maps or slices, so == would panic.
if reflect.DeepEqual(cfg.Get(o.key), o.value) {
continue
}
dropped = append(dropped, fmt.Sprintf("%s (%s)", o.key, o.envVar))
}
return dropped
}

// AuthTokenFilepath resolves the auth-token path via pkg/api/security's fallback rules.
Expand Down
72 changes: 72 additions & 0 deletions pkg/configstreambootstrap/bootstrap_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -11,6 +11,7 @@ import (

"github.com/stretchr/testify/require"

pkgconfigmodel "github.com/DataDog/datadog-agent/pkg/config/model"
pkgconfigsetup "github.com/DataDog/datadog-agent/pkg/config/setup"
)

Expand All @@ -21,3 +22,74 @@ func TestSeedGlobalBuilderResolvesIPCArtifactsNextToDatadogYaml(t *testing.T) {
SeedGlobalBuilder(Settings{CmdHost: "localhost", CmdPort: 5001}, yamlPath)
require.Equal(t, filepath.Join(dir, "auth_token"), AuthTokenFilepath())
}

// configEnvVars reaches the accessor the same way bootstrap.go does.
func configEnvVars(t *testing.T, cfg pkgconfigmodel.Reader) map[string][]string {
t.Helper()
lister, ok := cfg.(interface{ ConfigEnvVars() map[string][]string })
require.True(t, ok, "the global config must expose ConfigEnvVars")
return lister.ConfigEnvVars()
}

func capturedByKey(captured []envOverride) map[string]envOverride {
byKey := make(map[string]envOverride, len(captured))
for _, o := range captured {
byKey[o.key] = o
}
return byKey
}

func TestCaptureEnvOverridesOnlyTakesKeysTheEnvLayerDecides(t *testing.T) {
t.Setenv("DD_SITE", "datadoghq.eu")
t.Setenv("DD_LOG_LEVEL", "debug")
t.Setenv("DD_API_KEY", "some-secret-value")
pkgconfigsetup.InitConfigObjects()
cfg := pkgconfigsetup.Datadog()

// SourceRC (10) outranks SourceEnvVar (4), so the env var was never deciding this value.
cfg.Set("log_level", "info", pkgconfigmodel.SourceRC)
cfg.Set("api_key", "from-cli", pkgconfigmodel.SourceCLI)

captured := capturedByKey(captureEnvOverrides(cfg, configEnvVars(t, cfg)))
require.Equal(t, envOverride{key: "site", envVar: "DD_SITE", value: "datadoghq.eu"}, captured["site"])
require.NotContains(t, captured, "log_level", "a higher-precedence source was already overriding the env var")
require.NotContains(t, captured, "api_key", "a higher-precedence source was already overriding the env var")
}

func TestDiffEnvOverridesNamesOnlyDifferingSettings(t *testing.T) {
pkgconfigsetup.InitConfigObjects()
cfg := pkgconfigsetup.Datadog()
cfg.Set("site", "datadoghq.eu", pkgconfigmodel.SourceFile)

captured := []envOverride{
{key: "site", envVar: "DD_SITE", value: "datadoghq.eu"},
// Never streamed, so it reads back as the default: the incident-60263 shape.
{key: "runtime_security_config.enabled", envVar: "DD_RUNTIME_SECURITY_CONFIG_ENABLED", value: true},
}

require.Equal(t,
[]string{"runtime_security_config.enabled (DD_RUNTIME_SECURITY_CONFIG_ENABLED)"},
diffEnvOverrides(cfg, captured),
"a byte-identical streamed value must not warn")
}

func TestDiffEnvOverridesComparesMapsAndSlices(t *testing.T) {
pkgconfigsetup.InitConfigObjects()
cfg := pkgconfigsetup.Datadog()
cfg.Set("tags", []string{"env:prod", "team:agent"}, pkgconfigmodel.SourceFile)
cfg.Set("docker_labels_as_tags", map[string]string{"app": "kube_app"}, pkgconfigmodel.SourceFile)

same := []envOverride{
{key: "tags", envVar: "DD_TAGS", value: cfg.Get("tags")},
{key: "docker_labels_as_tags", envVar: "DD_DOCKER_LABELS_AS_TAGS", value: cfg.Get("docker_labels_as_tags")},
}
require.Empty(t, diffEnvOverrides(cfg, same))

differing := []envOverride{
{key: "tags", envVar: "DD_TAGS", value: []string{"env:staging"}},
{key: "docker_labels_as_tags", envVar: "DD_DOCKER_LABELS_AS_TAGS", value: map[string]string{"app": "other"}},
}
require.Equal(t,
[]string{"tags (DD_TAGS)", "docker_labels_as_tags (DD_DOCKER_LABELS_AS_TAGS)"},
diffEnvOverrides(cfg, differing))
}
15 changes: 15 additions & 0 deletions pkg/configstreambootstrap/testhelpers.go
Original file line number Diff line number Diff line change
Expand Up @@ -13,9 +13,24 @@ import (
pkgconfigsetup "github.com/DataDog/datadog-agent/pkg/config/setup"
)

// ResetGlobalConfig rebuilds the global config so its env layer reflects the current environment.
// Lives here because the pkgconfigusage depguard blocks pkg/config/setup imports from comp/.
func ResetGlobalConfig(t testing.TB) {
t.Helper()
pkgconfigsetup.InitConfigObjects()
t.Cleanup(pkgconfigsetup.InitConfigObjects)
}

// UseDynamicSchema makes the global config auto-rebuild the env layer when any DD_ var changes.
func UseDynamicSchema(t testing.TB) {
t.Helper()
pkgconfigsetup.Datadog().SetTestOnlyDynamicSchema(true)
t.Cleanup(func() { pkgconfigsetup.Datadog().SetTestOnlyDynamicSchema(false) })
}

// LastEnvOverrideReport returns the settings named by the most recent ReportDroppedEnvOverrides call.
func LastEnvOverrideReport() []string {
envOverridesMu.Lock()
defer envOverridesMu.Unlock()
return lastEnvOverrideReport
}
Original file line number Diff line number Diff line change
@@ -0,0 +1,11 @@
---
deprecations:
- |
When config streaming is enabled, ``DD_*`` environment variables set directly on a
remote Agent process (trace-agent, process-agent, security-agent, system-probe) are
ignored: the core Agent is the source of truth for configuration and its streamed
snapshot replaces the remote Agent's local environment layer. Such variables must be
set on the core Agent instead. After the first snapshot, remote Agents log a single
warning naming only the settings whose streamed value differs from the local
environment variable that was setting them; settings the core Agent streams back
unchanged are not reported. The warning names settings and variables, never values.
Loading