diff --git a/CHANGELOG.md b/CHANGELOG.md index 715c7d01..cdbb7a5e 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -1,5 +1,9 @@ -## 1.4.0 (May 5th, 2026) +## Unreleased + +Enhancements: +* VDS: instant updates for `database` and `ldap` static-roles via Vault event notifications. Set `spec.syncConfig.instantUpdates: true` and `spec.syncConfig.engineType: database|ldap` on a `VaultDynamicSecret` to subscribe to rotation events on the configured Vault mount. Requires Vault Enterprise (>= 1.16 for database, >= 1.21 for ldap) and a VaultAuth role with `read` on `sys/events/subscribe//*` plus `list`+`subscribe`+`subscribe_event_types = ["*"]` on the secret path. +## 1.4.0 (May 5th, 2026) Fix: * Detect TTL reset for uneven rotation schedules with ttl rollover bug: ([#1259](https://github.com/hashicorp/vault-secrets-operator/pull/1259)) * Update kube-rbac-proxy version for openshift: ([#1254](https://github.com/hashicorp/vault-secrets-operator/pull/1254)) diff --git a/api/v1beta1/vaultdynamicsecret_types.go b/api/v1beta1/vaultdynamicsecret_types.go index 183d400c..82c72993 100644 --- a/api/v1beta1/vaultdynamicsecret_types.go +++ b/api/v1beta1/vaultdynamicsecret_types.go @@ -70,6 +70,31 @@ type VaultDynamicSecretSpec struct { // +kubebuilder:validation:Type=string // +kubebuilder:validation:Pattern=`^([0-9]+(\.[0-9]+)?(s|m|h))$` RefreshAfter string `json:"refreshAfter,omitempty"` + // SyncConfig configures sync behavior from Vault to VSO. When + // SyncConfig.InstantUpdates is true, EngineType MUST be set to one of + // "database" or "ldap"; the operator subscribes to the corresponding Vault + // event stream and triggers a reconcile on rotation events that match this + // resource's mount and role name. Requires Vault Enterprise >= 1.16 + // (database) or >= 1.21 (ldap). + SyncConfig *VaultDynamicSecretSyncConfig `json:"syncConfig,omitempty"` +} + +// VaultDynamicSecretSyncConfig configures sync behavior from Vault to VSO for +// a VaultDynamicSecret. It is the dynamic-secret counterpart of +// VaultStaticSecret.Spec.SyncConfig. +type VaultDynamicSecretSyncConfig struct { + // InstantUpdates enables event-driven updates for this VaultDynamicSecret. + // Requires Vault Enterprise >= 1.16 (database) or >= 1.21 (ldap), and a + // VaultAuth role with read on sys/events/subscribe//* and + // list+subscribe on the secret path with subscribe_event_types = ["*"]. + // +kubebuilder:default=false + InstantUpdates bool `json:"instantUpdates,omitempty"` + // EngineType declares which Vault secrets-engine plugin produces the events + // VSO should subscribe to. Required when InstantUpdates is true. The + // selection is intentionally restricted to engines that publish rotation + // events suitable for instant updates. + // +kubebuilder:validation:Enum={database,ldap} + EngineType string `json:"engineType,omitempty"` } // VaultDynamicSecretStatus defines the observed state of VaultDynamicSecret diff --git a/api/v1beta1/vaultdynamicsecret_types_test.go b/api/v1beta1/vaultdynamicsecret_types_test.go new file mode 100644 index 00000000..fe30ce51 --- /dev/null +++ b/api/v1beta1/vaultdynamicsecret_types_test.go @@ -0,0 +1,51 @@ +// Copyright (c) HashiCorp, Inc. +// SPDX-License-Identifier: BUSL-1.1 + +package v1beta1 + +import ( + "testing" + + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" +) + +func TestVaultDynamicSecret_DeepCopy_SyncConfig(t *testing.T) { + t.Run("nil-syncconfig", func(t *testing.T) { + o := &VaultDynamicSecret{ + Spec: VaultDynamicSecretSpec{ + Mount: "database", + Path: "static-creds/myrole", + }, + } + c := o.DeepCopy() + require.NotNil(t, c) + assert.Nil(t, c.Spec.SyncConfig) + }) + + t.Run("populated-syncconfig", func(t *testing.T) { + o := &VaultDynamicSecret{ + Spec: VaultDynamicSecretSpec{ + Mount: "database", + Path: "static-creds/myrole", + SyncConfig: &VaultDynamicSecretSyncConfig{ + InstantUpdates: true, + EngineType: "database", + }, + }, + } + c := o.DeepCopy() + require.NotNil(t, c) + require.NotNil(t, c.Spec.SyncConfig) + assert.NotSame(t, o.Spec.SyncConfig, c.Spec.SyncConfig) + assert.Equal(t, o.Spec.SyncConfig.InstantUpdates, c.Spec.SyncConfig.InstantUpdates) + assert.Equal(t, o.Spec.SyncConfig.EngineType, c.Spec.SyncConfig.EngineType) + + c.Spec.SyncConfig.InstantUpdates = false + c.Spec.SyncConfig.EngineType = "ldap" + assert.True(t, o.Spec.SyncConfig.InstantUpdates, + "original mutated when copy was modified") + assert.Equal(t, "database", o.Spec.SyncConfig.EngineType, + "original engineType mutated when copy was modified") + }) +} diff --git a/api/v1beta1/zz_generated.deepcopy.go b/api/v1beta1/zz_generated.deepcopy.go index 82321b80..b84a610d 100644 --- a/api/v1beta1/zz_generated.deepcopy.go +++ b/api/v1beta1/zz_generated.deepcopy.go @@ -1625,6 +1625,11 @@ func (in *VaultDynamicSecretSpec) DeepCopyInto(out *VaultDynamicSecretSpec) { copy(*out, *in) } in.Destination.DeepCopyInto(&out.Destination) + if in.SyncConfig != nil { + in, out := &in.SyncConfig, &out.SyncConfig + *out = new(VaultDynamicSecretSyncConfig) + **out = **in + } } // DeepCopy is an autogenerated deepcopy function, copying the receiver, creating a new VaultDynamicSecretSpec. @@ -1662,6 +1667,21 @@ func (in *VaultDynamicSecretStatus) DeepCopy() *VaultDynamicSecretStatus { return out } +// DeepCopyInto is an autogenerated deepcopy function, copying the receiver, writing into out. in must be non-nil. +func (in *VaultDynamicSecretSyncConfig) DeepCopyInto(out *VaultDynamicSecretSyncConfig) { + *out = *in +} + +// DeepCopy is an autogenerated deepcopy function, copying the receiver, creating a new VaultDynamicSecretSyncConfig. +func (in *VaultDynamicSecretSyncConfig) DeepCopy() *VaultDynamicSecretSyncConfig { + if in == nil { + return nil + } + out := new(VaultDynamicSecretSyncConfig) + in.DeepCopyInto(out) + return out +} + // DeepCopyInto is an autogenerated deepcopy function, copying the receiver, writing into out. in must be non-nil. func (in *VaultPKISecret) DeepCopyInto(out *VaultPKISecret) { *out = *in diff --git a/chart/crds/secrets.hashicorp.com_vaultdynamicsecrets.yaml b/chart/crds/secrets.hashicorp.com_vaultdynamicsecrets.yaml index db515f77..b9344484 100644 --- a/chart/crds/secrets.hashicorp.com_vaultdynamicsecrets.yaml +++ b/chart/crds/secrets.hashicorp.com_vaultdynamicsecrets.yaml @@ -307,6 +307,34 @@ spec: - name type: object type: array + syncConfig: + description: |- + SyncConfig configures sync behavior from Vault to VSO. When + SyncConfig.InstantUpdates is true, EngineType MUST be set to one of + "database" or "ldap"; the operator subscribes to the corresponding Vault + event stream and triggers a reconcile on rotation events that match this + resource's mount and role name. Requires Vault Enterprise >= 1.16 + (database) or >= 1.21 (ldap). + properties: + engineType: + description: |- + EngineType declares which Vault secrets-engine plugin produces the events + VSO should subscribe to. Required when InstantUpdates is true. The + selection is intentionally restricted to engines that publish rotation + events suitable for instant updates. + enum: + - database + - ldap + type: string + instantUpdates: + default: false + description: |- + InstantUpdates enables event-driven updates for this VaultDynamicSecret. + Requires Vault Enterprise >= 1.16 (database) or >= 1.21 (ldap), and a + VaultAuth role with read on sys/events/subscribe//* and + list+subscribe on the secret path with subscribe_event_types = ["*"]. + type: boolean + type: object vaultAuthRef: description: |- VaultAuthRef to the VaultAuth resource, can be prefixed with a namespace, diff --git a/config/crd/bases/secrets.hashicorp.com_vaultdynamicsecrets.yaml b/config/crd/bases/secrets.hashicorp.com_vaultdynamicsecrets.yaml index db515f77..b9344484 100644 --- a/config/crd/bases/secrets.hashicorp.com_vaultdynamicsecrets.yaml +++ b/config/crd/bases/secrets.hashicorp.com_vaultdynamicsecrets.yaml @@ -307,6 +307,34 @@ spec: - name type: object type: array + syncConfig: + description: |- + SyncConfig configures sync behavior from Vault to VSO. When + SyncConfig.InstantUpdates is true, EngineType MUST be set to one of + "database" or "ldap"; the operator subscribes to the corresponding Vault + event stream and triggers a reconcile on rotation events that match this + resource's mount and role name. Requires Vault Enterprise >= 1.16 + (database) or >= 1.21 (ldap). + properties: + engineType: + description: |- + EngineType declares which Vault secrets-engine plugin produces the events + VSO should subscribe to. Required when InstantUpdates is true. The + selection is intentionally restricted to engines that publish rotation + events suitable for instant updates. + enum: + - database + - ldap + type: string + instantUpdates: + default: false + description: |- + InstantUpdates enables event-driven updates for this VaultDynamicSecret. + Requires Vault Enterprise >= 1.16 (database) or >= 1.21 (ldap), and a + VaultAuth role with read on sys/events/subscribe//* and + list+subscribe on the secret path with subscribe_event_types = ["*"]. + type: boolean + type: object vaultAuthRef: description: |- VaultAuthRef to the VaultAuth resource, can be prefixed with a namespace, diff --git a/config/samples/kustomization.yaml b/config/samples/kustomization.yaml index 13f76038..2c9bbb8f 100644 --- a/config/samples/kustomization.yaml +++ b/config/samples/kustomization.yaml @@ -9,6 +9,7 @@ resources: - secrets_v1beta1_vaultauth.yaml - secrets_v1beta1_vaultconnection.yaml - secrets_v1beta1_vaultdynamicsecret.yaml +- secrets_v1beta1_vaultdynamicsecret_instant_updates.yaml - secrets_v1beta1_hcpvaultsecretsapp.yaml - secrets_v1beta1_hcpauth.yaml - secrets_v1beta1_secrettransformation.yaml diff --git a/config/samples/secrets_v1beta1_vaultdynamicsecret_instant_updates.yaml b/config/samples/secrets_v1beta1_vaultdynamicsecret_instant_updates.yaml new file mode 100644 index 00000000..da268982 --- /dev/null +++ b/config/samples/secrets_v1beta1_vaultdynamicsecret_instant_updates.yaml @@ -0,0 +1,22 @@ +# Copyright (c) HashiCorp, Inc. +# SPDX-License-Identifier: BUSL-1.1 + +apiVersion: secrets.hashicorp.com/v1beta1 +kind: VaultDynamicSecret +metadata: + labels: + app.kubernetes.io/name: vaultdynamicsecret + app.kubernetes.io/instance: vaultdynamicsecret-instant-updates + app.kubernetes.io/part-of: vault-secrets-operator + name: vaultdynamicsecret-instant-updates +spec: + vaultAuthRef: example + mount: database + path: static-creds/myrole + allowStaticCreds: true + destination: + name: app-db-credentials + create: true + syncConfig: + instantUpdates: true + engineType: database diff --git a/consts/consts.go b/consts/consts.go index 6bf50a38..23cb6d7b 100644 --- a/consts/consts.go +++ b/consts/consts.go @@ -24,4 +24,7 @@ const ( AnnotationResync = "vso.hashicorp.com/resync" HeaderUserAgent = "User-Agent" + + VaultEngineTypeDatabase = "database" + VaultEngineTypeLDAP = "ldap" ) diff --git a/consts/reasons.go b/consts/reasons.go index bdf6a1d1..e88e092a 100644 --- a/consts/reasons.go +++ b/consts/reasons.go @@ -25,6 +25,7 @@ const ( ReasonHVSClientConfigError = "HVSClientConfigError" ReasonVaultClientError = "VaultClientError" ReasonVaultStaticSecret = "VaultStaticSecretError" + ReasonVaultDynamicSecret = "VaultDynamicSecretError" ReasonHVSSecret = "HVSSecretError" ReasonSecretDataDrift = "SecretDataDrift" ReasonInexistentDestination = "InexistentDestination" diff --git a/controllers/vaultdynamicsecret_controller.go b/controllers/vaultdynamicsecret_controller.go index 94c101e1..8a01a6c8 100644 --- a/controllers/vaultdynamicsecret_controller.go +++ b/controllers/vaultdynamicsecret_controller.go @@ -15,6 +15,7 @@ import ( "time" "github.com/cenkalti/backoff/v4" + "github.com/hashicorp/go-secure-stdlib/parseutil" "github.com/hashicorp/vault/api" corev1 "k8s.io/api/core/v1" apierrors "k8s.io/apimachinery/pkg/api/errors" @@ -22,6 +23,7 @@ import ( "k8s.io/apimachinery/pkg/runtime" "k8s.io/apimachinery/pkg/types" "k8s.io/client-go/tools/record" + "nhooyr.io/websocket" ctrl "sigs.k8s.io/controller-runtime" "sigs.k8s.io/controller-runtime/pkg/builder" @@ -78,8 +80,9 @@ type VaultDynamicSecretReconciler struct { // runtimePodUID should always be set when updating resource's Status. // This is done via the downwardAPI. We get the current Pod's UID from either the // OPERATOR_POD_UID environment variable, or the /var/run/podinfo/uid file; in that order. - runtimePodUID types.UID - SecretsClient client.Client + runtimePodUID types.UID + SecretsClient client.Client + eventWatcherRegistry *eventWatcherRegistry } // +kubebuilder:rbac:groups=secrets.hashicorp.com,resources=vaultdynamicsecrets,verbs=get;list;watch;create;update;patch;delete @@ -130,6 +133,24 @@ func (r *VaultDynamicSecretReconciler) Reconcile(ctx context.Context, req ctrl.R return ctrl.Result{}, r.handleDeletion(ctx, o) } + if o.Spec.SyncConfig != nil && o.Spec.SyncConfig.InstantUpdates { + switch o.Spec.SyncConfig.EngineType { + case consts.VaultEngineTypeDatabase, consts.VaultEngineTypeLDAP: + default: + r.unWatchEvents(o) + r.Recorder.Eventf(o, corev1.EventTypeWarning, + consts.ReasonVaultDynamicSecret, + "spec.syncConfig.engineType must be %q or %q when instantUpdates is enabled", + consts.VaultEngineTypeDatabase, consts.VaultEngineTypeLDAP) + horizon := computeHorizonWithJitter(requeueDurationOnError) + if err := r.updateStatus(ctx, o, false, newSyncCondition(o, metav1.ConditionFalse, + "Invalid spec.syncConfig.engineType, horizon=%s", horizon)); err != nil { + return ctrl.Result{}, err + } + return ctrl.Result{RequeueAfter: horizon}, nil + } + } + r.referenceCache.Set(SecretTransformation, req.NamespacedName, helpers.GetTransformationRefObjKeys( o.Spec.Destination.Transformation, o.Namespace)...) @@ -417,6 +438,16 @@ func (r *VaultDynamicSecretReconciler) Reconcile(ctx context.Context, req ctrl.R return ctrl.Result{}, err } + if o.Spec.SyncConfig != nil && o.Spec.SyncConfig.InstantUpdates { + logger.V(consts.LogLevelDebug).Info("Event watcher enabled") + if err := r.ensureEventWatcher(ctx, o, vClient); err != nil { + r.Recorder.Eventf(o, corev1.EventTypeWarning, + consts.ReasonEventWatcherError, "Failed to watch events: %s", err) + } + } else { + r.unWatchEvents(o) + } + if horizon.Seconds() == 0 { // no need to requeue logger.Info("Vault secret does not support periodic renewal/refresh via reconciliation", @@ -765,6 +796,7 @@ func (r *VaultDynamicSecretReconciler) SetupWithManager(mgr ctrl.Manager, opts c // TODO: close this channel when the controller is stopped. r.SourceCh = make(chan event.GenericEvent) + r.eventWatcherRegistry = newEventWatcherRegistry() m := ctrl.NewControllerManagedBy(mgr). For(&secretsv1beta1.VaultDynamicSecret{}). WithOptions(opts). @@ -810,6 +842,7 @@ func (r *VaultDynamicSecretReconciler) handleDeletion(ctx context.Context, o *se r.revokeLease(ctx, o, "") objKey := client.ObjectKeyFromObject(o) + r.unWatchEvents(o) r.SyncRegistry.Delete(objKey) r.BackOffRegistry.Delete(objKey) r.referenceCache.Remove(SecretTransformation, objKey) @@ -1090,3 +1123,315 @@ func vaultStaticCredsMetaDataFromData(data map[string]any) (*secretsv1beta1.Vaul return &ret, nil } + +const ( + databaseEventPath = "/v1/sys/events/subscribe/database/*" + ldapEventPath = "/v1/sys/events/subscribe/ldap/*" +) + +var dynamicSecretEventTypes = map[string]struct{}{ + "database/rotate": {}, + "database/static-creds-create": {}, + "database/role-update": {}, + "database/static-role-update": {}, + "ldap/rotate": {}, + "ldap/role-update": {}, + "ldap/static-role-update": {}, +} + +type dynamicSecretEventMsg struct { + Data struct { + Event struct { + Metadata struct { + Path string `json:"path"` + Name string `json:"name"` + Operation string `json:"operation"` + Modified string `json:"modified"` + } `json:"metadata"` + } `json:"event"` + EventType string `json:"event_type"` + PluginInfo struct { + MountPath string `json:"mount_path"` + } `json:"plugin_info"` + Namespace string `json:"namespace"` + } `json:"data"` +} + +func eventPathForEngine(engineType string) (string, error) { + switch engineType { + case consts.VaultEngineTypeDatabase: + return databaseEventPath, nil + case consts.VaultEngineTypeLDAP: + return ldapEventPath, nil + default: + return "", fmt.Errorf("unsupported engineType %q for instant updates", engineType) + } +} + +func matchEvent(o *secretsv1beta1.VaultDynamicSecret, ev *dynamicSecretEventMsg) bool { + if _, ok := dynamicSecretEventTypes[ev.Data.EventType]; !ok { + return false + } + if strings.Trim(ev.Data.Namespace, "/") != strings.Trim(o.Spec.Namespace, "/") { + return false + } + expectedMount := strings.TrimRight(o.Spec.Mount, "/") + "/" + if ev.Data.PluginInfo.MountPath != expectedMount { + return false + } + specPath := strings.Trim(o.Spec.Path, "/") + name := ev.Data.Event.Metadata.Name + if name == "" { + return false + } + if specPath != name && !strings.HasSuffix(specPath, "/"+name) { + return false + } + return true +} + +func (r *VaultDynamicSecretReconciler) ensureEventWatcher( + ctx context.Context, o *secretsv1beta1.VaultDynamicSecret, c vault.Client, +) error { + logger := log.FromContext(ctx).WithName("ensureEventWatcher") + name := client.ObjectKeyFromObject(o) + + if o.Spec.SyncConfig == nil { + return fmt.Errorf("syncConfig is nil") + } + eventPath, err := eventPathForEngine(o.Spec.SyncConfig.EngineType) + if err != nil { + return err + } + + if r.eventWatcherRegistry == nil { + r.eventWatcherRegistry = newEventWatcherRegistry() + } + + meta, ok := r.eventWatcherRegistry.Get(name) + if ok { + if meta.LastGeneration == o.GetGeneration() && meta.LastClientID == c.ID() { + logger.V(consts.LogLevelDebug).Info("Event watcher already running", + "namespace", o.Namespace, "name", o.Name) + return nil + } + } + if meta != nil { + if meta.Cancel != nil { + meta.Cancel() + waitCtx, cancel := context.WithTimeout(ctx, 2*time.Minute) + defer cancel() + if err := waitForStoppedCh(waitCtx, meta.StoppedCh); err != nil { + logger.Error(err, "Failed to stop event watcher", "name", name) + } + } else { + logger.Error(fmt.Errorf("nil cancel function"), + "event watcher has nil cancel function", + "VDS", name, "meta", meta) + } + } + + wsClient, err := c.WebsocketClient(eventPath) + if err != nil { + return fmt.Errorf("failed to create websocket client: %w", err) + } + + watchCtx, cancel := context.WithCancel(context.Background()) + stoppedCh := make(chan struct{}, 1) + updatedMeta := &eventWatcherMeta{ + Cancel: cancel, + LastClientID: c.ID(), + LastGeneration: o.GetGeneration(), + StoppedCh: stoppedCh, + } + logger.V(consts.LogLevelDebug).Info("Starting event watcher", "meta", updatedMeta) + r.eventWatcherRegistry.Register(name, updatedMeta) + go r.getDynamicSecretEvents(watchCtx, *o, wsClient, stoppedCh) + return nil +} + +func (r *VaultDynamicSecretReconciler) unWatchEvents(o *secretsv1beta1.VaultDynamicSecret) { + if r.eventWatcherRegistry == nil { + return + } + name := client.ObjectKeyFromObject(o) + if meta, ok := r.eventWatcherRegistry.Get(name); ok { + if meta.Cancel != nil { + meta.Cancel() + } + r.eventWatcherRegistry.Delete(name) + } +} + +func (r *VaultDynamicSecretReconciler) getDynamicSecretEvents( + ctx context.Context, + o secretsv1beta1.VaultDynamicSecret, + wsClient *vault.WebsocketClient, + stoppedCh chan struct{}, +) { + logger := log.FromContext(ctx).WithName("getDynamicSecretEvents") + name := client.ObjectKeyFromObject(&o) + defer func() { + r.eventWatcherRegistry.Delete(name) + close(stoppedCh) + }() + + retryBackoff := backoff.NewExponentialBackOff(r.BackOffRegistry.opts...) + shouldBackoff := false + errorThreshold := 5 + errorCount := 0 + +eventLoop: + for { + select { + case <-ctx.Done(): + logger.V(consts.LogLevelDebug).Info("Context done, stopping", + "namespace", o.Namespace, "name", o.Name) + return + default: + if shouldBackoff { + nextBackoff := retryBackoff.NextBackOff() + if nextBackoff == backoff.Stop { + logger.Error(fmt.Errorf("backoff limit reached"), "Backoff limit reached, requeuing") + break eventLoop + } + time.Sleep(retryBackoff.NextBackOff()) + } + + err := r.streamDynamicSecretEvents(ctx, &o, wsClient) + if err == nil { + continue + } + + if strings.Contains(err.Error(), "use of closed network connection") || + strings.Contains(err.Error(), "context canceled") { + logger.V(consts.LogLevelDebug).Info( + "Websocket client closed, stopping", + "namespace", o.Namespace, "name", o.Name) + return + } + + errorCount++ + shouldBackoff = true + r.Recorder.Eventf(&o, corev1.EventTypeWarning, + consts.ReasonEventWatcherError, + "Error while watching events: %s", err) + logger.Error(err, "Error while watching events", + "namespace", o.Namespace, "name", o.Name) + if errorCount >= errorThreshold { + logger.Error(err, "Too many errors, requeuing") + break eventLoop + } + + newVaultClient, gerr := r.ClientFactory.Get(ctx, r.Client, &o) + if gerr != nil { + logger.Error(gerr, "Failed to retrieve Vault client") + break eventLoop + } + eventPath, perr := eventPathForEngine(o.Spec.SyncConfig.EngineType) + if perr != nil { + logger.Error(perr, "Invalid engineType during reload") + break eventLoop + } + wsClient, err = newVaultClient.WebsocketClient(eventPath) + if err != nil { + logger.Error(err, "Failed to create new websocket client") + break eventLoop + } + + key := client.ObjectKeyFromObject(&o) + meta, ok := r.eventWatcherRegistry.Get(key) + if !ok { + logger.Error( + fmt.Errorf("failed to get event watcher metadata for VaultStaticSecret"), + "key", key.String()) + break eventLoop + } + meta.LastClientID = newVaultClient.ID() + r.eventWatcherRegistry.Register(key, meta) + } + } + + r.SyncRegistry.Add(client.ObjectKeyFromObject(&o)) + r.SourceCh <- event.GenericEvent{ + Object: &secretsv1beta1.VaultDynamicSecret{ + ObjectMeta: metav1.ObjectMeta{ + Namespace: o.Namespace, + Name: o.Name, + }, + }, + } +} + +func (r *VaultDynamicSecretReconciler) streamDynamicSecretEvents( + ctx context.Context, + o *secretsv1beta1.VaultDynamicSecret, + wsClient *vault.WebsocketClient, +) error { + logger := log.FromContext(ctx).WithName("streamDynamicSecretEvents") + conn, err := wsClient.Connect(ctx) + if err != nil { + return fmt.Errorf("failed to connect to vault websocket: %w", err) + } + defer conn.Close(websocket.StatusNormalClosure, "closing event watcher") + + r.Recorder.Event(o, corev1.EventTypeNormal, + consts.ReasonEventWatcherStarted, "Started watching events") + + for { + select { + case <-ctx.Done(): + logger.V(consts.LogLevelDebug).Info("Context done, closing websocket", + "namespace", o.Namespace, "name", o.Name) + return nil + default: + msgType, message, err := conn.Read(ctx) + if err != nil { + return fmt.Errorf("failed to read from websocket: %w, message: %q", + err, string(message)) + } + + ev := dynamicSecretEventMsg{} + if err := json.Unmarshal(message, &ev); err != nil { + return fmt.Errorf("failed to unmarshal event message: %w", err) + } + logger.V(consts.LogLevelTrace).Info("Received message", + "messageType", msgType, + "event_type", ev.Data.EventType, + "name", ev.Data.Event.Metadata.Name, + "mount", ev.Data.PluginInfo.MountPath, + "namespace", ev.Data.Namespace) + + modified, err := parseutil.ParseBool(ev.Data.Event.Metadata.Modified) + if err != nil { + return fmt.Errorf("failed to parse modified field: %w", err) + } + if !modified { + logger.V(consts.LogLevelTrace).Info("Event is not a modification, ignoring", + "event_type", ev.Data.EventType, + "name", ev.Data.Event.Metadata.Name) + continue + } + if !matchEvent(o, &ev) { + logger.V(consts.LogLevelTrace).Info("Event does not match, ignoring", + "event_type", ev.Data.EventType, + "name", ev.Data.Event.Metadata.Name) + continue + } + + logger.V(consts.LogLevelDebug).Info("Event matches, sending requeue", + "event_type", ev.Data.EventType, + "name", ev.Data.Event.Metadata.Name) + + r.SyncRegistry.Add(client.ObjectKeyFromObject(o)) + r.SourceCh <- event.GenericEvent{ + Object: &secretsv1beta1.VaultDynamicSecret{ + ObjectMeta: metav1.ObjectMeta{ + Namespace: o.Namespace, + Name: o.Name, + }, + }, + } + } + } +} diff --git a/controllers/vaultdynamicsecret_controller_test.go b/controllers/vaultdynamicsecret_controller_test.go index 29707db9..035f7162 100644 --- a/controllers/vaultdynamicsecret_controller_test.go +++ b/controllers/vaultdynamicsecret_controller_test.go @@ -19,6 +19,7 @@ import ( corev1 "k8s.io/api/core/v1" metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" "k8s.io/apimachinery/pkg/types" + "k8s.io/client-go/kubernetes/scheme" "k8s.io/client-go/tools/record" "k8s.io/client-go/util/workqueue" "sigs.k8s.io/controller-runtime/pkg/client" @@ -27,6 +28,7 @@ import ( "sigs.k8s.io/controller-runtime/pkg/reconcile" "sigs.k8s.io/controller-runtime/pkg/source" + vsoconsts "github.com/hashicorp/vault-secrets-operator/consts" "github.com/hashicorp/vault-secrets-operator/credentials/provider" "github.com/hashicorp/vault-secrets-operator/credentials/vault/consts" @@ -2110,3 +2112,468 @@ func TestVaultDynamicSecretReconciler_awaitRotation(t *testing.T) { }) } } + +func init() { + _ = secretsv1beta1.AddToScheme(scheme.Scheme) +} + +// TestEventPathForEngine tests the eventPathForEngine function. +func TestEventPathForEngine(t *testing.T) { + tests := []struct { + name string + engineType string + want string + wantErr bool + }{ + {name: "database", engineType: vsoconsts.VaultEngineTypeDatabase, want: databaseEventPath}, + {name: "ldap", engineType: vsoconsts.VaultEngineTypeLDAP, want: ldapEventPath}, + {name: "empty", engineType: "", wantErr: true}, + {name: "unsupported", engineType: "kv", wantErr: true}, + {name: "case-sensitive", engineType: "Database", wantErr: true}, + } + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + got, err := eventPathForEngine(tt.engineType) + if tt.wantErr { + require.Error(t, err) + return + } + require.NoError(t, err) + assert.Equal(t, tt.want, got) + }) + } +} + +func newVDSForMatch(mount, path, namespace string) *secretsv1beta1.VaultDynamicSecret { + return &secretsv1beta1.VaultDynamicSecret{ + Spec: secretsv1beta1.VaultDynamicSecretSpec{ + Mount: mount, + Path: path, + Namespace: namespace, + }, + } +} + +func newEvent(eventType, mountPath, name, namespace string) *dynamicSecretEventMsg { + ev := &dynamicSecretEventMsg{} + ev.Data.EventType = eventType + ev.Data.PluginInfo.MountPath = mountPath + ev.Data.Event.Metadata.Name = name + ev.Data.Namespace = namespace + return ev +} + +// TestMatchEvent tests the matchEvent function to ensure it correctly matches events to VaultDynamicSecret objects +// based on event type, mount path, role name, and namespace. +// It includes various test cases covering happy paths and edge cases for mismatches. +func TestMatchEvent(t *testing.T) { + tests := []struct { + name string + o *secretsv1beta1.VaultDynamicSecret + ev *dynamicSecretEventMsg + want bool + }{ + { + name: "happy-database-rotate", + o: newVDSForMatch("database", "static-creds/myrole", ""), + ev: newEvent("database/rotate", "database/", "myrole", ""), + want: true, + }, + { + name: "happy-database-rotate", + o: newVDSForMatch("database", "creds/myrole", ""), + ev: newEvent("database/rotate", "database/", "myrole", ""), + want: true, + }, + { + name: "happy-database-static-creds-create", + o: newVDSForMatch("database", "static-creds/myrole", ""), + ev: newEvent("database/static-creds-create", "database/", "myrole", ""), + want: true, + }, + { + name: "happy-database-role-update", + o: newVDSForMatch("database", "static-creds/myrole", ""), + ev: newEvent("database/role-update", "database/", "myrole", ""), + want: true, + }, + { + name: "happy-ldap-rotate", + o: newVDSForMatch("ldap", "static-cred/myrole", ""), + ev: newEvent("ldap/rotate", "ldap/", "myrole", ""), + want: true, + }, + { + name: "happy-ldap-static-role-update", + o: newVDSForMatch("ldap", "static-cred/myrole", ""), + ev: newEvent("ldap/static-role-update", "ldap/", "myrole", ""), + want: true, + }, + { + name: "happy-ldap-rotate-nested-role-name", + o: newVDSForMatch("ldap", "static-cred/group1/group2/myrole", ""), + ev: newEvent("ldap/rotate", "ldap/", "group1/group2/myrole", ""), + want: true, + }, + { + name: "happy-database-namespace-with-leading-slash-on-event", + o: newVDSForMatch("database", "static-creds/myrole", "ns1"), + ev: newEvent("database/rotate", "database/", "myrole", "/ns1/"), + want: true, + }, + { + name: "happy-mount-without-trailing-slash-on-spec", + o: newVDSForMatch("database/", "static-creds/myrole", ""), + ev: newEvent("database/rotate", "database/", "myrole", ""), + want: true, + }, + { + name: "rejected-event-type-not-in-allowlist", + o: newVDSForMatch("database", "static-creds/myrole", ""), + ev: newEvent("database/role-create", "database/", "myrole", ""), + want: false, + }, + { + name: "rejected-database-root-rotate", + o: newVDSForMatch("database", "static-creds/myrole", ""), + ev: newEvent("database/root-rotate", "database/", "myrole", ""), + want: false, + }, + { + name: "rejected-kv-event-on-database-watcher", + o: newVDSForMatch("database", "static-creds/myrole", ""), + ev: newEvent("kv-v2/data-write", "database/", "myrole", ""), + want: false, + }, + { + name: "rejected-namespace-mismatch", + o: newVDSForMatch("database", "static-creds/myrole", "ns1"), + ev: newEvent("database/rotate", "database/", "myrole", "ns2"), + want: false, + }, + { + name: "rejected-mount-mismatch", + o: newVDSForMatch("db", "static-creds/myrole", ""), + ev: newEvent("database/rotate", "database/", "myrole", ""), + want: false, + }, + { + name: "rejected-role-name-mismatch", + o: newVDSForMatch("database", "static-creds/myrole", ""), + ev: newEvent("database/rotate", "database/", "otherrole", ""), + want: false, + }, + { + name: "rejected-role-name-substring-without-segment-boundary", + o: newVDSForMatch("database", "static-creds/notmyrole", ""), + ev: newEvent("database/rotate", "database/", "myrole", ""), + want: false, + }, + { + name: "rejected-empty-event-name", + o: newVDSForMatch("database", "static-creds/myrole", ""), + ev: newEvent("database/rotate", "database/", "", ""), + want: false, + }, + { + name: "rejected-empty-event-type", + o: newVDSForMatch("database", "static-creds/myrole", ""), + ev: newEvent("", "database/", "myrole", ""), + want: false, + }, + } + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + assert.Equal(t, tt.want, matchEvent(tt.o, tt.ev)) + }) + } +} + +// TestMatchEvent_JSONRoundTrip tests that the matchEvent function correctly matches events to VaultDynamicSecret objects +func TestMatchEvent_JSONRoundTrip(t *testing.T) { + tests := []struct { + name string + payload string + o *secretsv1beta1.VaultDynamicSecret + want bool + }{ + { + name: "root-namespace-event-omits-namespace-field", + payload: `{ + "id":"abc", + "data":{ + "event":{ + "id":"abc", + "metadata":{ + "path":"database/rotate-role/myrole", + "name":"myrole", + "operation":"rotate", + "modified":"true" + } + }, + "event_type":"database/rotate", + "plugin_info":{ + "mount_class":"secret", + "mount_path":"database/", + "plugin":"database" + } + } + }`, + o: newVDSForMatch("database", "static-creds/myrole", ""), + want: true, + }, + { + name: "root-namespace-event-with-explicit-empty-namespace", + payload: `{ + "data":{ + "event":{ + "metadata":{ + "path":"database/rotate-role/myrole", + "name":"myrole", + "operation":"rotate", + "modified":"true" + } + }, + "event_type":"database/rotate", + "plugin_info":{"mount_path":"database/"}, + "namespace":"" + } + }`, + o: newVDSForMatch("database", "static-creds/myrole", ""), + want: true, + }, + { + name: "root-namespace-event-rejected-by-namespaced-vds", + payload: `{ + "data":{ + "event":{ + "metadata":{ + "name":"myrole", + "modified":"true" + } + }, + "event_type":"database/rotate", + "plugin_info":{"mount_path":"database/"} + } + }`, + o: newVDSForMatch("database", "static-creds/myrole", "ns1"), + want: false, + }, + { + name: "child-namespace-event-with-trailing-slash", + payload: `{ + "data":{ + "event":{ + "metadata":{ + "name":"myrole", + "modified":"true" + } + }, + "event_type":"database/rotate", + "plugin_info":{"mount_path":"database/"}, + "namespace":"ns1/" + } + }`, + o: newVDSForMatch("database", "static-creds/myrole", "ns1"), + want: true, + }, + { + name: "ldap-rotate-nested-role-name-real-payload", + payload: `{ + "id":"8fea78ca-47c5-0057-fc6f-38a631e4c0e5", + "data":{ + "event":{ + "id":"8fea78ca-47c5-0057-fc6f-38a631e4c0e5", + "metadata":{ + "modified":"true", + "name":"group1/group2/myrole", + "operation":"rotate", + "path":"ldap/rotate-role/group1/group2/myrole" + } + }, + "event_type":"ldap/rotate", + "plugin_info":{"mount_path":"ldap/","plugin":"ldap"} + } + }`, + o: newVDSForMatch("ldap", "static-cred/group1/group2/myrole", ""), + want: true, + }, + { + name: "ldap-rotate-nested-role-name-rejected-when-nesting-doesnt-match", + payload: `{ + "data":{ + "event":{ + "metadata":{ + "modified":"true", + "name":"group1/group2/myrole", + "operation":"rotate" + } + }, + "event_type":"ldap/rotate", + "plugin_info":{"mount_path":"ldap/","plugin":"ldap"} + } + }`, + o: newVDSForMatch("ldap", "static-cred/myrole", ""), + want: false, + }, + } + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + ev := dynamicSecretEventMsg{} + require.NoError(t, json.Unmarshal([]byte(tt.payload), &ev)) + assert.Equal(t, tt.want, matchEvent(tt.o, &ev)) + }) + } +} + +func newVDS(name, namespace string, syncCfg *secretsv1beta1.VaultDynamicSecretSyncConfig) *secretsv1beta1.VaultDynamicSecret { + return &secretsv1beta1.VaultDynamicSecret{ + ObjectMeta: metav1.ObjectMeta{ + Name: name, + Namespace: namespace, + }, + Spec: secretsv1beta1.VaultDynamicSecretSpec{ + Mount: "database", + Path: "static-creds/myrole", + SyncConfig: syncCfg, + }, + } +} + +// TestVDS_UnWatchEvents tests the unWatchEvents function to ensure it correctly unregisters event watchers +// and cancels their contexts. It verifies that the registry is updated and that the context is canceled as expected. +func TestVDS_UnWatchEvents(t *testing.T) { + r := &VaultDynamicSecretReconciler{ + eventWatcherRegistry: newEventWatcherRegistry(), + } + o := newVDS("foo", "default", nil) + + r.unWatchEvents(o) + + ctx, cancel := context.WithCancel(context.Background()) + stoppedCh := make(chan struct{}, 1) + r.eventWatcherRegistry.Register(client.ObjectKeyFromObject(o), &eventWatcherMeta{ + Cancel: cancel, + StoppedCh: stoppedCh, + LastGeneration: 1, + LastClientID: "id-1", + }) + require.Equal(t, 1, r.eventWatcherRegistry.registry.ItemCount()) + + r.unWatchEvents(o) + assert.Equal(t, 0, r.eventWatcherRegistry.registry.ItemCount()) + assert.ErrorIs(t, ctx.Err(), context.Canceled) + + r.unWatchEvents(o) +} + +func TestVDS_EnsureEventWatcher_NilSyncConfig(t *testing.T) { + r := &VaultDynamicSecretReconciler{ + eventWatcherRegistry: newEventWatcherRegistry(), + } + o := newVDS("foo", "default", nil) + err := r.ensureEventWatcher(context.Background(), o, nil) + require.Error(t, err) + assert.Contains(t, err.Error(), "syncConfig is nil") + assert.Equal(t, 0, r.eventWatcherRegistry.registry.ItemCount()) +} + +func TestVDS_EnsureEventWatcher_InvalidEngineType(t *testing.T) { + r := &VaultDynamicSecretReconciler{ + eventWatcherRegistry: newEventWatcherRegistry(), + } + o := newVDS("foo", "default", &secretsv1beta1.VaultDynamicSecretSyncConfig{ + InstantUpdates: true, + EngineType: "kv", + }) + err := r.ensureEventWatcher(context.Background(), o, nil) + require.Error(t, err) + assert.Contains(t, err.Error(), "unsupported engineType") + assert.Equal(t, 0, r.eventWatcherRegistry.registry.ItemCount()) +} + +// TestVDS_EnsureEventWatcher_RegistersWatcher tests that ensureEventWatcher correctly registers an event +// watcher for a valid VaultDynamicSecret object. +func TestVDS_HandleDeletionStopsWatcher(t *testing.T) { + o := newVDS("foo", "default", &secretsv1beta1.VaultDynamicSecretSyncConfig{ + InstantUpdates: true, + EngineType: "database", + }) + o.Finalizers = []string{vaultDynamicSecretFinalizer} + + c := fake.NewClientBuilder().WithScheme(scheme.Scheme).WithObjects(o).Build() + r := &VaultDynamicSecretReconciler{ + Client: c, + Recorder: record.NewFakeRecorder(8), + SyncRegistry: NewSyncRegistry(), + BackOffRegistry: NewBackOffRegistry(), + referenceCache: NewResourceReferenceCache(), + eventWatcherRegistry: newEventWatcherRegistry(), + } + + ctx, cancel := context.WithCancel(context.Background()) + stoppedCh := make(chan struct{}, 1) + r.eventWatcherRegistry.Register(client.ObjectKeyFromObject(o), &eventWatcherMeta{ + Cancel: cancel, + StoppedCh: stoppedCh, + LastGeneration: 1, + LastClientID: "id-1", + }) + require.Equal(t, 1, r.eventWatcherRegistry.registry.ItemCount()) + + r.unWatchEvents(o) + assert.Equal(t, 0, r.eventWatcherRegistry.registry.ItemCount()) + assert.ErrorIs(t, ctx.Err(), context.Canceled) +} + +func TestVDS_StreamMatchTriggersSync_FilterOnly(t *testing.T) { + r := &VaultDynamicSecretReconciler{ + SyncRegistry: NewSyncRegistry(), + SourceCh: make(chan event.GenericEvent, 1), + } + o := newVDS("foo", "default", &secretsv1beta1.VaultDynamicSecretSyncConfig{ + InstantUpdates: true, + EngineType: "database", + }) + + ev := newEvent("database/rotate", "database/", "myrole", "") + require.True(t, matchEvent(o, ev), "expected match for the test event") + + key := client.ObjectKeyFromObject(o) + r.SyncRegistry.Add(key) + r.SourceCh <- event.GenericEvent{ + Object: &secretsv1beta1.VaultDynamicSecret{ + ObjectMeta: metav1.ObjectMeta{Name: o.Name, Namespace: o.Namespace}, + }, + } + + assert.True(t, r.SyncRegistry.Has(key)) + select { + case got := <-r.SourceCh: + assert.Equal(t, o.Name, got.Object.GetName()) + assert.Equal(t, o.Namespace, got.Object.GetNamespace()) + default: + t.Fatal("expected a GenericEvent on SourceCh") + } +} + +// TestVDS_StreamFilterDropsNonMatchingEvent tests that non-matching events are correctly +// filtered out and do not trigger a sync. +func TestVDS_StreamFilterDropsNonMatchingEvent(t *testing.T) { + o := newVDS("foo", "default", &secretsv1beta1.VaultDynamicSecretSyncConfig{ + InstantUpdates: true, + EngineType: "database", + }) + + cases := []*dynamicSecretEventMsg{ + newEvent("kv-v2/data-write", "database/", "myrole", ""), // wrong event type (not in allow-list) + newEvent("database/role-create", "database/", "myrole", ""), // wrong event type (not in allow-list) + newEvent("database/rotate", "database/", "otherrole", ""), // role-name mismatch (Spec.Path = static-creds/myrole) + newEvent("database/rotate", "db/", "myrole", ""), // mount path mismatch (Spec.Path = static-creds/myrole) + } + for _, ev := range cases { + assert.False(t, matchEvent(o, ev), + "unexpected match for event_type=%s mount=%s name=%s", + ev.Data.EventType, ev.Data.PluginInfo.MountPath, ev.Data.Event.Metadata.Name) + } +} diff --git a/docs/api/api-reference.md b/docs/api/api-reference.md index a7a58cde..5d20b8c3 100644 --- a/docs/api/api-reference.md +++ b/docs/api/api-reference.md @@ -1187,10 +1187,30 @@ _Appears in:_ | `rolloutRestartTargets` _[RolloutRestartTarget](#rolloutrestarttarget) array_ | RolloutRestartTargets should be configured whenever the application(s) consuming the Vault secret does
not support dynamically reloading a rotated secret.
In that case one, or more RolloutRestartTarget(s) can be configured here. The Operator will
trigger a "rollout-restart" for each target whenever the Vault secret changes between reconciliation events.
See RolloutRestartTarget for more details. | | | | `destination` _[Destination](#destination)_ | Destination provides configuration necessary for syncing the Vault secret to Kubernetes. | | | | `refreshAfter` _string_ | RefreshAfter a period of time for VSO to sync the source secret data, in
duration notation e.g. 30s, 1m, 24h. This value only needs to be set when
syncing from a secret's engine that does not provide a lease TTL in its
response. The value should be within the secret engine's configured ttl or
max_ttl. The source secret's lease duration takes precedence over this
configuration when it is greater than 0. | | Pattern: `^([0-9]+(\.[0-9]+)?(s\|m\|h))$`
Type: string
| +| `syncConfig` _[VaultDynamicSecretSyncConfig](#vaultdynamicsecretsyncconfig)_ | SyncConfig configures sync behavior from Vault to VSO. When
SyncConfig.InstantUpdates is true, EngineType MUST be set to one of
"database" or "ldap"; the operator subscribes to the corresponding Vault
event stream and triggers a reconcile on rotation events that match this
resource's mount and role name. Requires Vault Enterprise >= 1.16
(database) or >= 1.21 (ldap). | | | +#### VaultDynamicSecretSyncConfig + + + +VaultDynamicSecretSyncConfig configures sync behavior from Vault to VSO for +a VaultDynamicSecret. It is the dynamic-secret counterpart of +VaultStaticSecret.Spec.SyncConfig. + + + +_Appears in:_ +- [VaultDynamicSecretSpec](#vaultdynamicsecretspec) + +| Field | Description | Default | Validation | +| --- | --- | --- | --- | +| `instantUpdates` _boolean_ | InstantUpdates enables event-driven updates for this VaultDynamicSecret.
Requires Vault Enterprise >= 1.16 (database) or >= 1.21 (ldap), and a
VaultAuth role with read on sys/events/subscribe//* and
list+subscribe on the secret path with subscribe_event_types = ["*"]. | false | | +| `engineType` _string_ | EngineType declares which Vault secrets-engine plugin produces the events
VSO should subscribe to. Required when InstantUpdates is true. The
selection is intentionally restricted to engines that publish rotation
events suitable for instant updates. | | Enum: [database ldap]
| + + #### VaultPKISecret diff --git a/test/integration/vaultdynamicsecret/terraform/postgres.tf b/test/integration/vaultdynamicsecret/terraform/postgres.tf index a1964fa3..9ad5661a 100644 --- a/test/integration/vaultdynamicsecret/terraform/postgres.tf +++ b/test/integration/vaultdynamicsecret/terraform/postgres.tf @@ -160,7 +160,24 @@ resource "vault_database_secret_backend_static_role" "postgres-scheduled" { resource "vault_policy" "db" { namespace = local.namespace name = "${local.auth_policy}-db" - policy = <= 1.16.3") + } + + ctx := context.Background() + crdClient := getCRDClient(t) + + tempDir, err := os.MkdirTemp(os.TempDir(), t.Name()) + require.NoError(t, err) + + tfDir := copyTerraformDir(t, path.Join(testRoot, "vaultdynamicsecret/terraform"), tempDir) + copyModulesDirT(t, tfDir) + chartsDir := copyTestChartsDir(t, tfDir) + + k8sConfigContext := os.Getenv("K8S_CLUSTER_CONTEXT") + if k8sConfigContext == "" { + k8sConfigContext = "kind-" + kindClusterName + } + + tfOptions := setCommonTFOptions(t, &terraform.Options{ + TerraformDir: tfDir, + Vars: map[string]interface{}{ + "k8s_config_context": k8sConfigContext, + "k8s_vault_namespace": k8sVaultNamespace, + "k8s_vault_service_account": "vault", + "name_prefix": "vds-evt", + "vault_address": os.Getenv("VAULT_ADDRESS"), + "vault_token": os.Getenv("VAULT_TOKEN"), + "vault_token_period": 120, + "vault_db_default_lease_ttl": 60, + "with_static_role_scheduled": withStaticRoleScheduled, + "with_xns": false, + "chart_postgres": filepath.Join(chartsDir, "postgresql"), + "vault_enterprise": true, + "use_events": true, + }, + }) + + skipCleanup := os.Getenv("SKIP_CLEANUP") != "" + var created []ctrlclient.Object + t.Cleanup(func() { + if skipCleanup { + t.Logf("Skipping cleanup, tfdir=%s", tfDir) + return + } + for _, o := range created { + assert.NoError(t, crdClient.Delete(ctx, o)) + } + if !testInParallel { + exportKindLogsT(t) + } + terraform.Destroy(t, tfOptions) + assert.NoError(t, os.RemoveAll(tempDir)) + }) + + terraform.InitAndApply(t, tfOptions) + + b, err := json.Marshal(terraform.OutputAll(t, tfOptions)) + require.NoError(t, err) + + var outputs dynamicK8SOutputs + require.NoError(t, json.Unmarshal(b, &outputs)) + + operatorNS := os.Getenv("OPERATOR_NAMESPACE") + require.NotEmpty(t, operatorNS, "OPERATOR_NAMESPACE is not set") + + auth := &secretsv1beta1.VaultAuth{ + ObjectMeta: v1.ObjectMeta{ + Name: "vds-instant-updates", + Namespace: operatorNS, + }, + Spec: secretsv1beta1.VaultAuthSpec{ + Namespace: outputs.Namespace, + Method: "kubernetes", + Mount: outputs.AuthMount, + Kubernetes: &secretsv1beta1.VaultAuthConfigKubernetes{ + Role: outputs.AuthRole, + ServiceAccount: "default", + TokenAudiences: []string{"vault"}, + }, + AllowedNamespaces: []string{outputs.K8sNamespace}, + }, + } + require.NoError(t, crdClient.Create(ctx, auth)) + created = append(created, auth) + + dest := "vds-instant-updates-static" + vds := &secretsv1beta1.VaultDynamicSecret{ + ObjectMeta: v1.ObjectMeta{ + Namespace: outputs.K8sNamespace, + Name: dest, + }, + Spec: secretsv1beta1.VaultDynamicSecretSpec{ + VaultAuthRef: ctrlclient.ObjectKeyFromObject(auth).String(), + Namespace: outputs.Namespace, + Mount: outputs.DBPath, + Path: "static-creds/" + outputs.DBRoleStatic, + AllowStaticCreds: true, + Destination: secretsv1beta1.Destination{ + Name: dest, + Create: true, + }, + SyncConfig: &secretsv1beta1.VaultDynamicSecretSyncConfig{ + InstantUpdates: true, + EngineType: consts.VaultEngineTypeDatabase, + }, + }, + } + require.NoError(t, crdClient.Create(ctx, vds)) + created = append(created, vds) + + objKey := ctrlclient.ObjectKeyFromObject(vds) + require.NoError(t, backoff.Retry(func() error { + var v secretsv1beta1.VaultDynamicSecret + if err := crdClient.Get(ctx, objKey, &v); err != nil { + return err + } + if v.Status.StaticCredsMetaData.LastVaultRotation == 0 { + return fmt.Errorf("waiting for initial static-creds sync") + } + return nil + }, backoff.WithMaxRetries(backoff.NewConstantBackOff(time.Second), 90))) + + require.NoError(t, backoff.Retry(func() error { + evList := corev1.EventList{} + if err := crdClient.List(ctx, &evList, + ctrlclient.InNamespace(vds.Namespace), + ctrlclient.MatchingFields{ + "involvedObject.name": vds.Name, + "reason": consts.ReasonEventWatcherStarted, + }, + ); err != nil { + return err + } + if len(evList.Items) == 0 { + return fmt.Errorf("no EventWatcherStarted event for %s", vds.Name) + } + return nil + }, backoff.WithMaxRetries(backoff.NewConstantBackOff(time.Second), 60))) + + var beforeVDS secretsv1beta1.VaultDynamicSecret + require.NoError(t, crdClient.Get(ctx, objKey, &beforeVDS)) + beforeRotation := beforeVDS.Status.StaticCredsMetaData.LastVaultRotation + require.NotZero(t, beforeRotation) + + authClient := getVaultClient(t, outputs.Namespace) + rotatePath := fmt.Sprintf("%s/rotate-role/%s", outputs.DBPath, outputs.DBRoleStatic) + _, err = authClient.Logical().WriteWithContext(ctx, rotatePath, nil) + require.NoError(t, err, "failed to rotate static role at %s", rotatePath) + + rotateAt := time.Now() + require.NoError(t, backoff.Retry(func() error { + var v secretsv1beta1.VaultDynamicSecret + if err := crdClient.Get(ctx, objKey, &v); err != nil { + return err + } + if v.Status.StaticCredsMetaData.LastVaultRotation == beforeRotation { + return fmt.Errorf("waiting for event-driven rotation pickup") + } + return nil + }, backoff.WithMaxRetries(backoff.NewConstantBackOff(500*time.Millisecond), 30))) + + elapsed := time.Since(rotateAt) + t.Logf("event-driven rotation picked up in %s", elapsed) + + require.Less(t, elapsed, 15*time.Second, + "expected event-driven rotation to be detected well below the rotation_period") + + evList := corev1.EventList{} + require.NoError(t, crdClient.List(ctx, &evList, + ctrlclient.InNamespace(vds.Namespace), + ctrlclient.MatchingFields{ + "involvedObject.name": vds.Name, + "reason": consts.ReasonEventWatcherError, + }, + )) + for _, ev := range evList.Items { + if ev.Type == corev1.EventTypeWarning { + t.Errorf("unexpected EventWatcherError warning: %s", ev.Message) + } + } +} + +// TestVaultDynamicSecretInstantUpdates_InvalidEngineType verifies that the CRD +// enum rejects an empty/invalid engineType. +func TestVaultDynamicSecretInstantUpdates_InvalidEngineType(t *testing.T) { + if testInParallel { + t.Parallel() + } + require.NotEmpty(t, kindClusterName, "KIND_CLUSTER_NAME is not set") + + ctx := context.Background() + crdClient := getCRDClient(t) + + operatorNS := os.Getenv("OPERATOR_NAMESPACE") + require.NotEmpty(t, operatorNS, "OPERATOR_NAMESPACE is not set") + + vds := &secretsv1beta1.VaultDynamicSecret{ + ObjectMeta: v1.ObjectMeta{ + Namespace: operatorNS, + Name: "vds-invalid-engine-type", + }, + Spec: secretsv1beta1.VaultDynamicSecretSpec{ + Mount: "database", + Path: "static-creds/whatever", + Destination: secretsv1beta1.Destination{ + Name: "vds-invalid-engine-type", + Create: true, + }, + SyncConfig: &secretsv1beta1.VaultDynamicSecretSyncConfig{ + InstantUpdates: true, + EngineType: "kv", + }, + }, + } + err := crdClient.Create(ctx, vds) + require.Error(t, err, "CRD enum should reject engineType=kv") +}