Skip to content

Commit 9c3331f

Browse files
authored
[CONTINT] Retry event checkpoint read (#53752)
### What does this PR do? Adds retry-with-backoff when reading the `kubernetes_apiserver` event-collection checkpoint from its ConfigMap. ### Motivation A transient error reading the checkpoint was being treated the same as "checkpoint doesn't exist," which resets the watermark to 0 and causing the initial watch to potentially redeliver all events. For example: ``` 2026-07-16 14:38:39 UTC | CLUSTER | ERROR | (/home/datadog/go/src/github.com/DataDog/datadog-agent-checkpoint-baseline/pkg/util/kubernetes/apiserver/apiserver.go:509 in getOrCreateConfigMap) | Could not get the ConfigMap baseline-token: configmaps "baseline-token" not found, trying to create it. 2026-07-16 14:38:39 UTC | CLUSTER | INFO | (/home/datadog/go/src/github.com/DataDog/datadog-agent-checkpoint-baseline/pkg/util/kubernetes/apiserver/apiserver.go:519 in getOrCreateConfigMap) | Created the ConfigMap baseline-token ``` ### Describe how you validated your changes Reproduced the behavior on a kind cluster running the DCA via the Operator by denying the DCA GET requests to the checkpoint ConfigMap temporarily to hit the retry. ### Additional Notes Co-authored-by: jon.rosario <jon.rosario@datadoghq.com>
1 parent f4b370a commit 9c3331f

3 files changed

Lines changed: 29 additions & 2 deletions

File tree

pkg/collector/corechecks/cluster/kubernetesapiserver/kubernetes_apiserver.go

Lines changed: 24 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -372,15 +372,37 @@ func (k *KubeASCheck) startEventCollection() error {
372372
k.mu.Unlock()
373373

374374
// If the checkpoint is present, seed the watermark from it so the initial list only forwards newer events.
375-
if resVer, _, err := k.ac.GetTokenFromConfigmap(eventTokenKey); err != nil {
376-
log.Warnf("Could not read persisted event checkpoint, starting fresh: %s", err)
375+
resVer, err := k.readEventCheckpointWithRetry()
376+
if err != nil {
377+
log.Warnf("Could not read persisted event checkpoint after retries, starting fresh: %s", err)
377378
} else {
378379
ec.SetCheckpoint(resVer)
379380
}
380381

381382
return ec.Start(stopCh)
382383
}
383384

385+
// readEventCheckpointWithRetry reads the persisted event checkpoint, retrying with backoff
386+
func (k *KubeASCheck) readEventCheckpointWithRetry() (string, error) {
387+
const maxAttempts = 5
388+
389+
delay := time.Second
390+
var lastErr error
391+
for attempt := 1; attempt <= maxAttempts; attempt++ {
392+
resVer, _, err := k.ac.GetTokenFromConfigmap(eventTokenKey)
393+
if err == nil {
394+
return resVer, nil
395+
}
396+
lastErr = err
397+
if attempt < maxAttempts {
398+
log.Warnf("Could not read persisted event checkpoint (attempt %d/%d): %s", attempt, maxAttempts, err)
399+
time.Sleep(delay)
400+
delay *= 2
401+
}
402+
}
403+
return "", lastErr
404+
}
405+
384406
// stopEventCollection stops the running EventCollector by closing its stop
385407
// channel. It is idempotent.
386408
func (k *KubeASCheck) stopEventCollection() {

pkg/util/kubernetes/apiserver/BUILD.bazel

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -48,6 +48,7 @@ go_library(
4848
"@io_k8s_api//discovery/v1:discovery",
4949
"@io_k8s_apiextensions_apiserver//pkg/client/clientset/clientset",
5050
"@io_k8s_apiextensions_apiserver//pkg/client/informers/externalversions",
51+
"@io_k8s_apimachinery//pkg/api/errors",
5152
"@io_k8s_apimachinery//pkg/api/meta",
5253
"@io_k8s_apimachinery//pkg/apis/meta/v1:meta",
5354
"@io_k8s_apimachinery//pkg/apis/meta/v1/unstructured",

pkg/util/kubernetes/apiserver/apiserver.go

Lines changed: 4 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -21,6 +21,7 @@ import (
2121

2222
v1 "k8s.io/api/core/v1"
2323
"k8s.io/apiextensions-apiserver/pkg/client/clientset/clientset"
24+
apierrors "k8s.io/apimachinery/pkg/api/errors"
2425
"k8s.io/apimachinery/pkg/api/meta"
2526
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
2627
"k8s.io/apimachinery/pkg/fields"
@@ -506,6 +507,9 @@ func (c *APIClient) ComponentStatuses() (*v1.ComponentStatusList, error) {
506507
func (c *APIClient) getOrCreateConfigMap(name, namespace string) (cmEvent *v1.ConfigMap, err error) {
507508
cmEvent, err = c.Cl.CoreV1().ConfigMaps(namespace).Get(context.TODO(), name, metav1.GetOptions{})
508509
if err != nil {
510+
if !apierrors.IsNotFound(err) {
511+
return nil, fmt.Errorf("could not get the ConfigMap %s: %w", name, err)
512+
}
509513
log.Errorf("Could not get the ConfigMap %s: %s, trying to create it.", name, err.Error())
510514
cmEvent, err = c.Cl.CoreV1().ConfigMaps(namespace).Create(context.TODO(), &v1.ConfigMap{
511515
ObjectMeta: metav1.ObjectMeta{

0 commit comments

Comments
 (0)