Skip to content

Commit 94bc64e

Browse files
authored
[autodiscovery] Wait for tag completeness before creating AD services (#48732)
### What does this PR do? This PR is part of the effort to address the "missing tags" issue in the Agent. It introduces a new configuration option (`ad_tag_completeness_max_wait`) that makes autodiscovery wait until all tags are complete before processing an entity. The goal is to avoid scheduling checks with incomplete tags. When tags are not complete before the specified max wait, the check is scheduled as it is today with whatever tags are available at that moment. This PR also adds a new telemetry metric (`autodiscovery.tag_completeness_delay`) that captures the delay in processing discovered entities. The new option only applies to Kubernetes environments, because it's the only environment where tags can be detected as incomplete by workloadmeta. The other environment where this could happen is ECS EC2, but that has not been implemented yet. The new option is disabled by default for now and not exposed. The reason is that, with the default config, tags can be delayed for several seconds. Mainly because the kubemetadata workloadmeta collector is pull-based and runs every 60s. It was recently converted to stream-based, but the option to use it is not enabled by default yet (needs more testing). ### Describe how you validated your changes New unit tests plus tests in a local kind cluster: - Deployed with the default config and verified that checks are scheduled immediately. The telemetry doesn't show any delays. - Deployed with `ad_tag_completeness_max_wait=120`. Checks are a bit delayed as expected and the telemetry confirms it. - Deployed with `ad_tag_completeness_max_wait=120` and kubemetadata streaming enabled (`kubernetes_metadata_streaming=true`). Checks are scheduled quickly and the telemetry confirms it. Co-authored-by: david.ortiz <david.ortiz@datadoghq.com>
1 parent 7349bff commit 94bc64e

6 files changed

Lines changed: 393 additions & 1 deletion

File tree

comp/core/autodiscovery/listeners/kubelet.go

Lines changed: 34 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -55,7 +55,8 @@ func NewKubeletListener(options ServiceListernerDeps) (ServiceListener, error) {
5555
return nil, errors.New("workloadmeta store is not initialized")
5656
}
5757
var err error
58-
l.workloadmetaListener, err = newWorkloadmetaListener(name, wmetaFilter, l.processPod, wmetaInstance, options.Telemetry)
58+
maxWait := time.Duration(pkgconfigsetup.Datadog().GetInt("ad_tag_completeness_max_wait")) * time.Second
59+
l.workloadmetaListener, err = newWorkloadmetaListenerWithTagWait(name, wmetaFilter, l.processPod, wmetaInstance, options.Telemetry, l.areTagsComplete, maxWait)
5960
if err != nil {
6061
return nil, err
6162
}
@@ -215,3 +216,35 @@ func (l *KubeletListener) createContainerService(
215216
podSvcID := buildSvcID(pod.GetID())
216217
l.AddService(svcID, svc, podSvcID)
217218
}
219+
220+
func (l *KubeletListener) areTagsComplete(entity workloadmeta.Entity) bool {
221+
pod, ok := entity.(*workloadmeta.KubernetesPod)
222+
if !ok {
223+
log.Errorf("expected KubernetesPod entity, got %T", entity)
224+
return true
225+
}
226+
227+
podTaggerID := common.BuildTaggerEntityID(pod.GetID())
228+
_, podComplete, err := l.tagger.TagWithCompleteness(podTaggerID, types.ChecksConfigCardinality)
229+
if err != nil {
230+
log.Debugf("error checking tag completeness for pod %s: %s", pod.ID, err)
231+
return false
232+
}
233+
if !podComplete {
234+
return false
235+
}
236+
237+
for _, podContainer := range pod.GetAllContainers() {
238+
containerTaggerID := types.NewEntityID(types.ContainerID, podContainer.ID)
239+
_, containerComplete, err := l.tagger.TagWithCompleteness(containerTaggerID, types.ChecksConfigCardinality)
240+
if err != nil {
241+
log.Debugf("error checking tag completeness for container %s: %s", podContainer.ID, err)
242+
return false
243+
}
244+
if !containerComplete {
245+
return false
246+
}
247+
}
248+
249+
return true
250+
}

comp/core/autodiscovery/listeners/kubelet_test.go

Lines changed: 128 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -12,8 +12,11 @@ import (
1212
"testing"
1313
"time"
1414

15+
"github.com/stretchr/testify/assert"
16+
1517
tagger "github.com/DataDog/datadog-agent/comp/core/tagger/def"
1618
taggerfxmock "github.com/DataDog/datadog-agent/comp/core/tagger/fx-mock"
19+
"github.com/DataDog/datadog-agent/comp/core/tagger/types"
1720
workloadfilter "github.com/DataDog/datadog-agent/comp/core/workloadfilter/def"
1821
workloadfilterfxmock "github.com/DataDog/datadog-agent/comp/core/workloadfilter/fx-mock"
1922
workloadmeta "github.com/DataDog/datadog-agent/comp/core/workloadmeta/def"
@@ -663,6 +666,131 @@ func TestProcessPodWithEphemeralContainer(t *testing.T) {
663666
wlm.assertServices(expectedServices)
664667
}
665668

669+
func TestAreTagsComplete(t *testing.T) {
670+
podEntityID := types.NewEntityID(types.KubernetesPodUID, podID)
671+
containerEntityID := types.NewEntityID(types.ContainerID, containerID)
672+
673+
tests := []struct {
674+
name string
675+
pod *workloadmeta.KubernetesPod
676+
tagInfos []*types.TagInfo
677+
expected bool
678+
}{
679+
{
680+
name: "pod and container complete",
681+
pod: &workloadmeta.KubernetesPod{
682+
EntityID: workloadmeta.EntityID{
683+
Kind: workloadmeta.KindKubernetesPod,
684+
ID: podID,
685+
},
686+
Containers: []workloadmeta.OrchestratorContainer{
687+
{
688+
ID: containerID,
689+
Name: containerName,
690+
},
691+
},
692+
},
693+
tagInfos: []*types.TagInfo{
694+
{
695+
Source: "source",
696+
EntityID: podEntityID,
697+
LowCardTags: []string{"kube_namespace:default"},
698+
IsComplete: true,
699+
},
700+
{
701+
Source: "source",
702+
EntityID: containerEntityID,
703+
LowCardTags: []string{"container_name:agent"},
704+
IsComplete: true,
705+
},
706+
},
707+
expected: true,
708+
},
709+
{
710+
name: "pod incomplete",
711+
pod: &workloadmeta.KubernetesPod{
712+
EntityID: workloadmeta.EntityID{
713+
Kind: workloadmeta.KindKubernetesPod,
714+
ID: podID,
715+
},
716+
},
717+
tagInfos: []*types.TagInfo{
718+
{
719+
Source: "source",
720+
EntityID: podEntityID,
721+
LowCardTags: []string{"kube_namespace:default"},
722+
IsComplete: false,
723+
},
724+
},
725+
expected: false,
726+
},
727+
{
728+
name: "container incomplete",
729+
pod: &workloadmeta.KubernetesPod{
730+
EntityID: workloadmeta.EntityID{
731+
Kind: workloadmeta.KindKubernetesPod,
732+
ID: podID,
733+
},
734+
Containers: []workloadmeta.OrchestratorContainer{
735+
{
736+
ID: containerID,
737+
Name: containerName,
738+
},
739+
},
740+
},
741+
tagInfos: []*types.TagInfo{
742+
{
743+
Source: "source",
744+
EntityID: podEntityID,
745+
LowCardTags: []string{"kube_namespace:default"},
746+
IsComplete: true,
747+
},
748+
{
749+
Source: "source",
750+
EntityID: containerEntityID,
751+
LowCardTags: []string{"container_name:agent"},
752+
IsComplete: false,
753+
},
754+
},
755+
expected: false,
756+
},
757+
{
758+
name: "container not in tagger",
759+
pod: &workloadmeta.KubernetesPod{
760+
EntityID: workloadmeta.EntityID{
761+
Kind: workloadmeta.KindKubernetesPod,
762+
ID: podID,
763+
},
764+
Containers: []workloadmeta.OrchestratorContainer{
765+
{
766+
ID: "unknown-container",
767+
Name: "unknown",
768+
},
769+
},
770+
},
771+
tagInfos: []*types.TagInfo{
772+
{
773+
Source: "source",
774+
EntityID: podEntityID,
775+
LowCardTags: []string{"kube_namespace:default"},
776+
IsComplete: true,
777+
},
778+
},
779+
expected: false,
780+
},
781+
}
782+
783+
for _, test := range tests {
784+
t.Run(test.name, func(t *testing.T) {
785+
taggerMock := taggerfxmock.SetupFakeTagger(t)
786+
taggerMock.GetTagStore().ProcessTagInfo(test.tagInfos)
787+
listener, _ := newKubeletListener(t, taggerMock)
788+
789+
assert.Equal(t, test.expected, listener.areTagsComplete(test.pod))
790+
})
791+
}
792+
}
793+
666794
func newKubeletListener(t *testing.T, tagger tagger.Component) (*KubeletListener, *testWorkloadmetaListener) {
667795
wlm := newTestWorkloadmetaListener(t)
668796
filterStore := workloadfilterfxmock.SetupMockFilter(t)

comp/core/autodiscovery/listeners/workloadmeta.go

Lines changed: 124 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -10,13 +10,16 @@ package listeners
1010
import (
1111
"fmt"
1212
"strings"
13+
"time"
1314

1415
"github.com/DataDog/datadog-agent/comp/core/autodiscovery/telemetry"
1516
workloadmeta "github.com/DataDog/datadog-agent/comp/core/workloadmeta/def"
1617
"github.com/DataDog/datadog-agent/pkg/status/health"
1718
"github.com/DataDog/datadog-agent/pkg/util/log"
1819
)
1920

21+
const tagCompletenessRetryInterval = 1 * time.Second
22+
2023
// workloadmetaListener is a generic subscriber to workloadmeta events that
2124
// generates AD services.
2225
type workloadmetaListener interface {
@@ -33,6 +36,11 @@ type workloadmetaListener interface {
3336
AddService(svcID string, svc Service, parentSvcID string)
3437
}
3538

39+
type pendingEntityInfo struct {
40+
entity workloadmeta.Entity
41+
firstSeen time.Time
42+
}
43+
3644
// workloadmetaListenerImpl implements workloadmetaListener.
3745
type workloadmetaListenerImpl struct {
3846
name string
@@ -50,6 +58,10 @@ type workloadmetaListenerImpl struct {
5058
delService chan<- Service
5159

5260
telemetryStore *telemetry.Store
61+
62+
isReadyFn func(workloadmeta.Entity) bool // when nil, consider ready
63+
pendingEntities map[string]pendingEntityInfo // svcID → pending info
64+
maxWaitDuration time.Duration
5365
}
5466

5567
var _ workloadmetaListener = &workloadmetaListenerImpl{}
@@ -82,6 +94,36 @@ func newWorkloadmetaListener(
8294
}, nil
8395
}
8496

97+
// newWorkloadmetaListenerWithTagWait is like newWorkloadmetaListener but defers
98+
// processing until isReadyFn returns true. Disabled when maxWait is 0.
99+
func newWorkloadmetaListenerWithTagWait(
100+
name string,
101+
workloadFilters *workloadmeta.Filter,
102+
processFn func(workloadmeta.Entity),
103+
wmeta workloadmeta.Component,
104+
telemetryStore *telemetry.Store,
105+
isReadyFn func(workloadmeta.Entity) bool,
106+
maxWait time.Duration,
107+
) (workloadmetaListener, error) {
108+
base, err := newWorkloadmetaListener(name, workloadFilters, processFn, wmeta, telemetryStore)
109+
if err != nil {
110+
return nil, err
111+
}
112+
113+
if maxWait > 0 {
114+
listener, ok := base.(*workloadmetaListenerImpl)
115+
if !ok {
116+
return nil, fmt.Errorf("unexpected listener type %T", base)
117+
}
118+
119+
listener.isReadyFn = isReadyFn
120+
listener.pendingEntities = make(map[string]pendingEntityInfo)
121+
listener.maxWaitDuration = maxWait
122+
}
123+
124+
return base, nil
125+
}
126+
85127
func (l *workloadmetaListenerImpl) Store() workloadmeta.Component {
86128
return l.store
87129
}
@@ -126,7 +168,17 @@ func (l *workloadmetaListenerImpl) Listen(newSvc chan<- Service, delSvc chan<- S
126168
log.Infof("%s initialized successfully", l.name)
127169

128170
go func() {
171+
var retryTicker *time.Ticker
172+
var retryChan <-chan time.Time
173+
if l.isReadyFn != nil {
174+
retryTicker = time.NewTicker(tagCompletenessRetryInterval)
175+
retryChan = retryTicker.C
176+
}
177+
129178
defer func() {
179+
if retryTicker != nil {
180+
retryTicker.Stop()
181+
}
130182
err := health.Deregister()
131183
if err != nil {
132184
log.Warnf("error de-registering health check: %s", err)
@@ -142,6 +194,9 @@ func (l *workloadmetaListenerImpl) Listen(newSvc chan<- Service, delSvc chan<- S
142194

143195
l.processEvents(evBundle)
144196

197+
case <-retryChan:
198+
l.retryPendingEntities()
199+
145200
case <-health.C:
146201

147202
case <-l.stop:
@@ -181,6 +236,10 @@ func (l *workloadmetaListenerImpl) processEvents(evBundle workloadmeta.EventBund
181236
func (l *workloadmetaListenerImpl) processSetEntity(entity workloadmeta.Entity) {
182237
svcID := buildSvcID(entity.GetID())
183238

239+
if l.isReadyFn != nil && l.waitIfNotReady(svcID, entity) {
240+
return
241+
}
242+
184243
// keep track of children of this entity from previous iterations ...
185244
unseen := make(map[string]struct{})
186245
for childSvcID := range l.children[svcID] {
@@ -204,10 +263,75 @@ func (l *workloadmetaListenerImpl) processSetEntity(entity workloadmeta.Entity)
204263
}
205264
}
206265

266+
// waitIfNotReady returns true if the entity is not ready and should be retried
267+
// later.
268+
func (l *workloadmetaListenerImpl) waitIfNotReady(svcID string, entity workloadmeta.Entity) bool {
269+
if l.isReadyFn(entity) {
270+
_, wasPending := l.pendingEntities[svcID]
271+
l.resolvePending(svcID)
272+
if !wasPending && l.telemetryStore != nil {
273+
l.telemetryStore.TagCompletenessDelay.Observe(0, l.name)
274+
}
275+
return false
276+
}
277+
278+
now := time.Now()
279+
280+
pending, exists := l.pendingEntities[svcID]
281+
if !exists {
282+
pending = pendingEntityInfo{
283+
entity: entity,
284+
firstSeen: now,
285+
}
286+
} else {
287+
pending.entity = entity
288+
}
289+
290+
if now.Sub(pending.firstSeen) < l.maxWaitDuration {
291+
l.pendingEntities[svcID] = pending
292+
log.Debugf("%s not adding entity %s: tags not complete yet", l.name, svcID)
293+
return true
294+
}
295+
296+
// Timeout exceeded
297+
log.Warnf("%s adding entity %s with potentially incomplete tags", l.name, svcID)
298+
l.resolvePending(svcID)
299+
300+
return false
301+
}
302+
303+
func (l *workloadmetaListenerImpl) resolvePending(svcID string) {
304+
pending, wasPending := l.pendingEntities[svcID]
305+
if !wasPending {
306+
return
307+
}
308+
309+
delay := time.Since(pending.firstSeen).Seconds()
310+
if l.telemetryStore != nil {
311+
l.telemetryStore.TagCompletenessDelay.Observe(delay, l.name)
312+
}
313+
delete(l.pendingEntities, svcID)
314+
}
315+
316+
func (l *workloadmetaListenerImpl) retryPendingEntities() {
317+
pending := make([]workloadmeta.Entity, 0, len(l.pendingEntities))
318+
for _, pendingEntity := range l.pendingEntities {
319+
pending = append(pending, pendingEntity.entity)
320+
}
321+
322+
for _, entity := range pending {
323+
l.processSetEntity(entity)
324+
}
325+
}
326+
207327
func (l *workloadmetaListenerImpl) processUnsetEntity(entity workloadmeta.Entity) {
208328
entityID := entity.GetID()
209329
parentSvcID := buildSvcID(entityID)
210330

331+
if l.pendingEntities != nil {
332+
delete(l.pendingEntities, parentSvcID)
333+
}
334+
211335
l.removeService(parentSvcID)
212336

213337
childrenSvcIDs := l.children[parentSvcID]

0 commit comments

Comments
 (0)