Skip to content

Commit 1e49b37

Browse files
authored
CAP-3858 fix missing AgentVersion on manifest payload (#53915)
<!--Please give us some feedback on your experience writing this PR ! https://app.datadoghq.com/forms/43db4c02-6837-400c-8083-692e141b1b88 !--> ### What does this PR do? Fix missing `AgentVersion` on `CollectorManifest` payloads for all non-pod Kubernetes resources. For buffered manifest payloads (Deployments, Services, Nodes, and every other non-pod resource collected by the Cluster Agent orchestrator check), we have `AgentVersion = nil`. Root cause: `ManifestBuffer.flushManifest` built a synthetic processor context from a snapshot of a handful of config fields and that snapshot silently omitted `AgentVersion`. Pods weren't affected because they're collected by the process-agent on each node, not through this buffer. ### What changed - **Fix the drop.** `ManifestBuffer` now holds a reference to the owning `OrchestratorCheck` and builds the flush context via a new `newFlushContext()` helper that reads shared per-check fields live from the check. - **Prevent the class of bug.** Removed the identity fields from `ManifestBufferConfig` so it only holds buffer tuning knobs. Any future per-check field will flow through `newFlushContext` automatically instead of needing to be manually copied into a snapshot. - **Also fixed on the way through:** - Terminated-resource run config was dropping `AgentVersion` too — propagated it from the main run config. - `newOrchestratorCheck` wasn't copying the `Meta` field from `version.Version` into `model.AgentVersion` — added it. ### Describe how you validated your changes Added a regression assertion in `TestFlushManifest` that fails if `sentManifest.AgentVersion` is nil or doesn't match the check's version. All 647 package tests pass. ### Additional Notes Co-authored-by: kangyi.li <kangyi.li@datadoghq.com>
1 parent 7a98cab commit 1e49b37

4 files changed

Lines changed: 28 additions & 31 deletions

File tree

pkg/collector/corechecks/cluster/orchestrator/collector_bundle.go

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -117,6 +117,7 @@ func NewCollectorBundle(chk *OrchestratorCheck) *CollectorBundle {
117117
ClusterID: runCfg.ClusterID,
118118
Config: runCfg.Config,
119119
MsgGroupRef: runCfg.MsgGroupRef,
120+
AgentVersion: runCfg.AgentVersion,
120121
TerminatedResources: true,
121122
}
122123

pkg/collector/corechecks/cluster/orchestrator/manifest_buffer.go

Lines changed: 21 additions & 25 deletions
Original file line numberDiff line numberDiff line change
@@ -20,7 +20,6 @@ import (
2020
"github.com/DataDog/datadog-agent/pkg/collector/corechecks/cluster/orchestrator/processors"
2121
"github.com/DataDog/datadog-agent/pkg/collector/corechecks/cluster/orchestrator/processors/common"
2222
"github.com/DataDog/datadog-agent/pkg/orchestrator"
23-
"github.com/DataDog/datadog-agent/pkg/orchestrator/config"
2423
pkgorchestratormodel "github.com/DataDog/datadog-agent/pkg/orchestrator/model"
2524
"github.com/DataDog/datadog-agent/pkg/util/log"
2625
)
@@ -40,17 +39,14 @@ func init() {
4039
bufferExpVars.Set("BufferFlushed", bufferFlushedTotal)
4140
}
4241

43-
// ManifestBufferConfig contains information about buffering manifests.
42+
// ManifestBufferConfig holds buffer tuning knobs. Per-check identity
43+
// (cluster, agent version, config, ...) lives on ManifestBuffer.chk and is
44+
// read live at flush time via newFlushContext.
4445
type ManifestBufferConfig struct {
45-
KubeClusterName string
46-
ClusterID string
47-
MaxPerMessage int
48-
MaxWeightPerMessageBytes int
4946
MsgGroupRef *atomic.Int32
5047
BufferedManifestEnabled bool
5148
MaxBufferedManifests int
5249
ManifestBufferFlushInterval time.Duration
53-
ExtraTags []string
5450
}
5551

5652
// ManifestBuffer is a buffer of manifest sent from all collectors
@@ -59,6 +55,7 @@ type ManifestBufferConfig struct {
5955
// and gets stopped after the check is done.
6056
type ManifestBuffer struct {
6157
Cfg *ManifestBufferConfig
58+
chk *OrchestratorCheck
6259
ManifestChan chan interface{}
6360
bufferedManifests []interface{}
6461
stopCh chan struct{}
@@ -68,16 +65,12 @@ type ManifestBuffer struct {
6865
// NewManifestBuffer returns a new ManifestBuffer
6966
func NewManifestBuffer(chk *OrchestratorCheck) *ManifestBuffer {
7067
manifestBuffer := &ManifestBuffer{
68+
chk: chk,
7169
Cfg: &ManifestBufferConfig{
72-
ClusterID: chk.clusterID,
73-
KubeClusterName: chk.orchestratorConfig.KubeClusterName,
7470
MsgGroupRef: chk.groupID,
75-
MaxPerMessage: chk.orchestratorConfig.MaxPerMessage,
76-
MaxWeightPerMessageBytes: chk.orchestratorConfig.MaxWeightPerMessageBytes,
7771
BufferedManifestEnabled: chk.orchestratorConfig.BufferedManifestEnabled,
7872
MaxBufferedManifests: chk.orchestratorConfig.MaxPerMessage,
7973
ManifestBufferFlushInterval: chk.orchestratorConfig.ManifestBufferFlushInterval,
80-
ExtraTags: chk.orchestratorConfig.ExtraTags,
8174
},
8275
ManifestChan: make(chan interface{}),
8376
stopCh: make(chan struct{}),
@@ -87,23 +80,26 @@ func NewManifestBuffer(chk *OrchestratorCheck) *ManifestBuffer {
8780
return manifestBuffer
8881
}
8982

90-
// flushManifest flushes manifests by chunking them first then sending them to the sender
91-
func (cb *ManifestBuffer) flushManifest(sender sender.Sender) {
92-
manifests := cb.bufferedManifests
93-
ctx := &processors.K8sProcessorContext{
83+
// newFlushContext is the single source of truth for shared per-check fields
84+
// on the buffered manifest flush path.
85+
func (cb *ManifestBuffer) newFlushContext() *processors.K8sProcessorContext {
86+
return &processors.K8sProcessorContext{
9487
BaseProcessorContext: processors.BaseProcessorContext{
95-
MsgGroupID: cb.Cfg.MsgGroupRef.Inc(),
96-
Cfg: &config.OrchestratorConfig{
97-
KubeClusterName: cb.Cfg.KubeClusterName,
98-
MaxPerMessage: cb.Cfg.MaxPerMessage,
99-
MaxWeightPerMessageBytes: cb.Cfg.MaxWeightPerMessageBytes,
100-
ExtraTags: cb.Cfg.ExtraTags,
101-
},
102-
ClusterID: cb.Cfg.ClusterID,
88+
Cfg: cb.chk.orchestratorConfig,
89+
MsgGroupID: cb.Cfg.MsgGroupRef.Inc(),
90+
ClusterID: cb.chk.clusterID,
91+
AgentVersion: cb.chk.agentVersion,
10392
},
93+
APIClient: cb.chk.apiClient,
10494
}
95+
}
96+
97+
// flushManifest flushes manifests by chunking them first then sending them to the sender
98+
func (cb *ManifestBuffer) flushManifest(sender sender.Sender) {
99+
manifests := cb.bufferedManifests
100+
ctx := cb.newFlushContext()
105101
manifestMessages := processors.ChunkManifest(ctx, common.BaseHandlers{}.BuildManifestMessageBody, manifests)
106-
sender.OrchestratorManifest(manifestMessages, cb.Cfg.ClusterID)
102+
sender.OrchestratorManifest(manifestMessages, cb.chk.clusterID)
107103
setManifestStats(manifests)
108104
cb.bufferedManifests = cb.bufferedManifests[:0]
109105
}

pkg/collector/corechecks/cluster/orchestrator/manifest_buffer_test.go

Lines changed: 5 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -123,11 +123,8 @@ func TestNewManifestBuffer(t *testing.T) {
123123
assert.Equal(t, 0, len(mb.bufferedManifests))
124124
assert.Equal(t, cap(mb.bufferedManifests), mb.Cfg.MaxBufferedManifests)
125125

126-
// Verify configuration was copied correctly
127-
assert.Equal(t, orchCheck.clusterID, mb.Cfg.ClusterID)
128-
assert.Equal(t, orchCheck.orchestratorConfig.KubeClusterName, mb.Cfg.KubeClusterName)
129-
assert.Equal(t, orchCheck.orchestratorConfig.MaxPerMessage, mb.Cfg.MaxPerMessage)
130-
assert.Equal(t, orchCheck.orchestratorConfig.MaxWeightPerMessageBytes, mb.Cfg.MaxWeightPerMessageBytes)
126+
assert.Equal(t, orchCheck.orchestratorConfig.MaxPerMessage, mb.Cfg.MaxBufferedManifests)
127+
assert.Same(t, orchCheck, mb.chk)
131128
}
132129

133130
func TestFlushManifest(t *testing.T) {
@@ -171,10 +168,12 @@ func TestFlushManifest(t *testing.T) {
171168

172169
// Verify the manifest contains expected data
173170
assert.Equal(t, "buffer-cluster", sentManifest.ClusterName)
174-
assert.Equal(t, mb.Cfg.ClusterID, sentManifest.ClusterId)
171+
assert.Equal(t, mb.chk.clusterID, sentManifest.ClusterId)
175172
assert.Equal(t, []string{"tag:low"}, sentManifest.Tags) // From ExtraTags in test setup
176173
assert.Equal(t, int32(1), sentManifest.GroupId) // MsgGroupRef.Inc() should return 1 for first call
177174
assert.Equal(t, int32(1), sentManifest.GroupSize) // Only one chunk
175+
require.NotNil(t, sentManifest.AgentVersion)
176+
assert.Equal(t, mb.chk.agentVersion, sentManifest.AgentVersion)
178177

179178
// Verify manifests are correctly included
180179
assert.Len(t, sentManifest.Manifests, 2)

pkg/collector/corechecks/cluster/orchestrator/orchestrator.go

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -113,6 +113,7 @@ func newOrchestratorCheck(base core.CheckBase, instance *OrchestratorInstance, c
113113
Minor: agentVersion.Minor,
114114
Patch: agentVersion.Patch,
115115
Pre: agentVersion.Pre,
116+
Meta: agentVersion.Meta,
116117
Commit: agentVersion.Commit,
117118
},
118119
}

0 commit comments

Comments
 (0)