diff --git a/cmd/objectbucket-notifications-adapter/main.go b/cmd/objectbucket-notifications-adapter/main.go index b547458..2ce9763 100644 --- a/cmd/objectbucket-notifications-adapter/main.go +++ b/cmd/objectbucket-notifications-adapter/main.go @@ -20,21 +20,16 @@ import ( "context" "crypto/tls" "flag" - "fmt" "os" "path/filepath" - "regexp" - "strconv" "strings" // Import all Kubernetes client auth plugins (e.g. Azure, GCP, OIDC, etc.) // to ensure that exec-entrypoint and run can make use of them. _ "k8s.io/client-go/plugin/pkg/client/auth" - metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" "k8s.io/apimachinery/pkg/runtime" utilruntime "k8s.io/apimachinery/pkg/util/runtime" - "k8s.io/client-go/kubernetes" clientgoscheme "k8s.io/client-go/kubernetes/scheme" ctrl "sigs.k8s.io/controller-runtime" "sigs.k8s.io/controller-runtime/pkg/certwatcher" @@ -44,11 +39,9 @@ import ( metricsserver "sigs.k8s.io/controller-runtime/pkg/metrics/server" "sigs.k8s.io/controller-runtime/pkg/webhook" - "github.com/IBM/sarama" - sourcesv1alpha1 "github.com/functions-dev/func-operator/api/sources/v1alpha1" + "github.com/functions-dev/func-operator/internal/objectbucketsource/config" "github.com/functions-dev/func-operator/internal/objectbucketsource/controller" - kafkaconfig "github.com/functions-dev/func-operator/internal/objectbucketsource/kafka" "github.com/functions-dev/func-operator/internal/objectbucketsource/notificationserver" // +kubebuilder:scaffold:imports ) @@ -74,7 +67,15 @@ func main() { var probeAddr string var secureMetrics bool var enableHTTP2 bool + var configMapName string + var createConfig bool + var adapterPort int + var notificationsMode string + var kafkaBrokers string + var kafkaNotificationsTopics string + var kafkaNotificationsGroupID string var tlsOpts []func(*tls.Config) + flag.StringVar(&metricsAddr, "metrics-bind-address", "0", "The address the metrics endpoint binds to. "+ "Use :8443 for HTTPS or :8080 for HTTP, or leave as 0 to disable the metrics service.") flag.StringVar(&probeAddr, "health-probe-bind-address", ":8081", "The address the probe endpoint binds to.") @@ -92,6 +93,25 @@ func main() { flag.StringVar(&metricsCertKey, "metrics-cert-key", "tls.key", "The name of the metrics server key file.") flag.BoolVar(&enableHTTP2, "enable-http2", false, "If set, HTTP/2 will be enabled for the metrics and webhook servers") + flag.StringVar(&configMapName, "config", "objectbucket-notifications-adapter-config", + "Name of the ConfigMap containing adapter configuration") + flag.BoolVar(&createConfig, "create-config", false, + "If set, create the default configuration ConfigMap at startup if it does not already exist.") + flag.IntVar(&adapterPort, "adapter-port", 8888, + "Port the notification HTTP server listens on (HTTP mode only)") + flag.StringVar(¬ificationsMode, "notifications-mode", "http", + "http or kafka - selects how the adapter receives notifications. "+ + "Default for the NOTIFICATIONS_MODE ConfigMap key, which can override it at runtime.") + flag.StringVar(&kafkaBrokers, "kafka-brokers", "", + "Comma-separated list of Kafka broker addresses (required for Kafka mode). "+ + "Default for the KAFKA_BROKERS ConfigMap key, which can override it at runtime.") + flag.StringVar(&kafkaNotificationsTopics, "kafka-notifications-topics", "", + "Comma-separated list of Kafka topics to consume notifications from (required for Kafka mode). "+ + "Default for the KAFKA_NOTIFICATIONS_TOPICS ConfigMap key, which can override it at runtime.") + flag.StringVar(&kafkaNotificationsGroupID, "kafka-notifications-group-id", "", + "Consumer group ID for consuming notifications (required for Kafka mode). "+ + "Default for the KAFKA_NOTIFICATIONS_GROUP_ID ConfigMap key, which can override it at runtime.") + opts := zap.Options{ Development: true, } @@ -178,136 +198,64 @@ func main() { os.Exit(1) } - noobaaAdapterID := envOrDefault("NOOBAA_ADAPTER_ID", "mcg-adapter") - noobaaAdapterTopic := envOrDefault("NOOBAA_ADAPTER_TOPIC_ARN", "mcg-adapter-connection/connect.json") - noobaaStorageClassPattern := envOrDefault("NOOBAA_ADAPTER_STORAGECLASS_PATTERN", `.*noobaa\.io$`) - - radosgwAdapterID := envOrDefault("RADOSGW_ADAPTER_ID", "rgw-adapter") - radosgwAdapterTopic := envOrDefault("RADOSGW_ADAPTER_TOPIC_ARN", - "arn:aws:sns:ocs-storagecluster-cephobjectstore::rgw-adapter-notifications") - radosgwStorageClassPattern := envOrDefault("RADOSGW_ADAPTER_STORAGECLASS_PATTERN", `.*ceph-rgw$`) - - adapterConfigs := make([]controller.AdapterConfig, 0, 2) - for _, cfg := range []struct { - id, topic, pattern string - }{ - {noobaaAdapterID, noobaaAdapterTopic, noobaaStorageClassPattern}, - {radosgwAdapterID, radosgwAdapterTopic, radosgwStorageClassPattern}, - } { - re, err := regexp.Compile(cfg.pattern) - if err != nil { - setupLog.Error(err, "invalid storageclass pattern", "pattern", cfg.pattern) - os.Exit(1) - } - adapterConfigs = append(adapterConfigs, controller.AdapterConfig{ - ID: cfg.id, - Topic: cfg.topic, - StorageClassPattern: re, - }) + // The command-line flags provide the defaults for the notification settings. + // The actual values are resolved from the ConfigMap (falling back to these + // defaults) and can be changed at runtime. Validation happens in the config + // provider when the ConfigMap is loaded. + notificationDefaults := config.NotificationSettings{ + Mode: notificationsMode, + KafkaBrokers: splitAndTrim(kafkaBrokers), + KafkaNotificationsTopics: splitAndTrim(kafkaNotificationsTopics), + KafkaNotificationsGroupID: kafkaNotificationsGroupID, } - adapterPort := 8888 - if portStr := os.Getenv("ADAPTER_PORT"); portStr != "" { - var err error - adapterPort, err = strconv.Atoi(portStr) + // Determine namespace for ConfigMap + ns := os.Getenv("POD_NAMESPACE") + if ns == "" { + nsBytes, err := os.ReadFile("/var/run/secrets/kubernetes.io/serviceaccount/namespace") if err != nil { - setupLog.Error(err, "invalid ADAPTER_PORT") + setupLog.Error(err, "cannot determine pod namespace") os.Exit(1) } - } - notificationsMode := os.Getenv("NOTIFICATIONS_MODE") - if notificationsMode == "" { - notificationsMode = "http" - } - if notificationsMode != "http" && notificationsMode != "kafka" { - setupLog.Error(fmt.Errorf("invalid NOTIFICATIONS_MODE %q", notificationsMode), "must be \"http\" or \"kafka\"") - os.Exit(1) + ns = strings.TrimSpace(string(nsBytes)) } - var kafkaNotificationsTopics []string - if topicsStr := os.Getenv("KAFKA_NOTIFICATIONS_TOPIC"); topicsStr != "" { - for _, t := range strings.Split(topicsStr, ",") { - if trimmed := strings.TrimSpace(t); trimmed != "" { - kafkaNotificationsTopics = append(kafkaNotificationsTopics, trimmed) - } + // Optionally create the default configuration ConfigMap before the provider + // tries to read it. + if createConfig { + if err := config.EnsureDefaultConfigMap(context.Background(), ns, configMapName, notificationDefaults); err != nil { + setupLog.Error(err, "failed to ensure default configuration ConfigMap") + os.Exit(1) } } - kafkaNotificationsGroupID := os.Getenv("KAFKA_NOTIFICATIONS_GROUP_ID") - var kafkaBrokers []string - if brokersStr := os.Getenv("KAFKA_BROKERS"); brokersStr != "" { - kafkaBrokers = strings.Split(brokersStr, ",") + // Create configuration provider + configProvider, err := config.NewProvider(context.Background(), ns, configMapName, notificationDefaults) + if err != nil { + setupLog.Error(err, "failed to create configuration provider") + os.Exit(1) } - var kafkaCfg *sarama.Config - if kafkaSecretName := os.Getenv("KAFKA_SECRET"); kafkaSecretName != "" { - ns := os.Getenv("POD_NAMESPACE") - if ns == "" { - nsBytes, err := os.ReadFile("/var/run/secrets/kubernetes.io/serviceaccount/namespace") - if err != nil { - setupLog.Error(err, "cannot determine pod namespace for KAFKA_SECRET") - os.Exit(1) - } - ns = strings.TrimSpace(string(nsBytes)) - } - clientset, err := kubernetes.NewForConfig(ctrl.GetConfigOrDie()) - if err != nil { - setupLog.Error(err, "creating kubernetes clientset for KAFKA_SECRET") - os.Exit(1) - } - secret, err := clientset.CoreV1().Secrets(ns).Get(context.Background(), kafkaSecretName, metav1.GetOptions{}) - if err != nil { - setupLog.Error(err, "reading KAFKA_SECRET", "name", kafkaSecretName, "namespace", ns) - os.Exit(1) - } - kafkaCfg, err = kafkaconfig.NewConfig(secret.Data) - if err != nil { - setupLog.Error(err, "configuring kafka from secret", "name", kafkaSecretName) - os.Exit(1) - } - setupLog.Info("kafka configured from secret", "name", kafkaSecretName, "namespace", ns) - } else { - var err error - kafkaCfg, err = kafkaconfig.NewConfig(nil) - if err != nil { - setupLog.Error(err, "creating default kafka config") - os.Exit(1) - } + // Add config provider to manager + if err := mgr.Add(configProvider); err != nil { + setupLog.Error(err, "unable to add config provider to manager") + os.Exit(1) } if err := (&controller.ObjectBucketSourceReconciler{ Client: mgr.GetClient(), Scheme: mgr.GetScheme(), - AdapterConfigs: adapterConfigs, + ConfigProvider: configProvider, }).SetupWithManager(mgr); err != nil { setupLog.Error(err, "unable to create controller", "controller", "ObjectBucketSource") os.Exit(1) } // +kubebuilder:scaffold:builder - if notificationsMode == "kafka" { - if len(kafkaNotificationsTopics) == 0 { - setupLog.Error(fmt.Errorf("KAFKA_NOTIFICATIONS_TOPIC is required when NOTIFICATIONS_MODE=kafka"), "missing env") - os.Exit(1) - } - if kafkaNotificationsGroupID == "" { - setupLog.Error(fmt.Errorf("KAFKA_NOTIFICATIONS_GROUP_ID is required when NOTIFICATIONS_MODE=kafka"), "missing env") - os.Exit(1) - } - if len(kafkaBrokers) == 0 { - setupLog.Error(fmt.Errorf("KAFKA_BROKERS is required when NOTIFICATIONS_MODE=kafka"), "missing env") - os.Exit(1) - } - } - notifServer := ¬ificationserver.NotificationServer{ - Client: mgr.GetClient(), - Port: adapterPort, - KafkaBrokers: kafkaBrokers, - KafkaConfig: kafkaCfg, - NotificationsMode: notificationsMode, - KafkaNotificationsTopics: kafkaNotificationsTopics, - KafkaNotificationsGroupID: kafkaNotificationsGroupID, + Client: mgr.GetClient(), + Port: adapterPort, + ConfigProvider: configProvider, } if err := mgr.Add(notifServer); err != nil { setupLog.Error(err, "unable to add notification server") @@ -346,9 +294,14 @@ func main() { } } -func envOrDefault(key, defaultValue string) string { - if v := os.Getenv(key); v != "" { - return v +// splitAndTrim splits a comma-separated flag value, trimming whitespace and +// dropping empty entries. +func splitAndTrim(s string) []string { + var out []string + for _, part := range strings.Split(s, ",") { + if trimmed := strings.TrimSpace(part); trimmed != "" { + out = append(out, trimmed) + } } - return defaultValue + return out } diff --git a/config/combined/source-objectbucket/manager.yaml b/config/combined/source-objectbucket/manager.yaml index b47b6f1..eb6811d 100644 --- a/config/combined/source-objectbucket/manager.yaml +++ b/config/combined/source-objectbucket/manager.yaml @@ -31,6 +31,7 @@ spec: args: - --leader-elect - --health-probe-bind-address=:8081 + - --create-config image: source-objectbucket-adapter:latest name: manager env: diff --git a/config/rbac/role.yaml b/config/rbac/role.yaml index 5f28be4..e8cae5f 100644 --- a/config/rbac/role.yaml +++ b/config/rbac/role.yaml @@ -9,6 +9,7 @@ rules: resources: - configmaps verbs: + - create - get - list - watch diff --git a/config/samples/objectbucket-notifications-adapter-config.yaml b/config/samples/objectbucket-notifications-adapter-config.yaml new file mode 100644 index 0000000..044dcd9 --- /dev/null +++ b/config/samples/objectbucket-notifications-adapter-config.yaml @@ -0,0 +1,27 @@ +apiVersion: v1 +kind: ConfigMap +metadata: + name: objectbucket-notifications-adapter-config + namespace: system +data: + # NooBaa adapter configuration + NOOBAA_ADAPTER_ID: "mcg-adapter" + NOOBAA_ADAPTER_TOPIC_ARN: "mcg-adapter-connection/connect.json" + NOOBAA_ADAPTER_STORAGECLASS_PATTERN: ".*noobaa\\.io$" + + # RadosGW adapter configuration + RADOSGW_ADAPTER_ID: "rgw-adapter" + RADOSGW_ADAPTER_TOPIC_ARN: "arn:aws:sns:ocs-storagecluster-cephobjectstore::rgw-adapter-notifications" + RADOSGW_ADAPTER_STORAGECLASS_PATTERN: ".*ceph-rgw$" + + # Notification transport configuration (dynamic). + # These override the corresponding command-line flag defaults and are applied + # at runtime. Changing any Kafka setting gracefully restarts the consumer. + # NOTIFICATIONS_MODE: "http" # http or kafka + # KAFKA_BROKERS: "broker1:9092,broker2:9092" # required for kafka mode + # KAFKA_NOTIFICATIONS_TOPICS: "mcg-notifications,rgw-notifications" # required for kafka mode + # KAFKA_NOTIFICATIONS_GROUP_ID: "adapter-consumer-group" # required for kafka mode + + # Kafka secret for authentication (optional) + # Uncomment to use a Kafka secret for authentication + # KAFKA_SECRET: "kafka-credentials" diff --git a/config/sources/objectbucket/manager/manager.yaml b/config/sources/objectbucket/manager/manager.yaml index dba4d90..2c1d48b 100644 --- a/config/sources/objectbucket/manager/manager.yaml +++ b/config/sources/objectbucket/manager/manager.yaml @@ -40,6 +40,7 @@ spec: args: - --leader-elect - --health-probe-bind-address=:8081 + - --create-config image: source-objectbucket-adapter:latest name: manager env: diff --git a/config/sources/objectbucket/rbac/role.yaml b/config/sources/objectbucket/rbac/role.yaml index c57dba4..f584acf 100644 --- a/config/sources/objectbucket/rbac/role.yaml +++ b/config/sources/objectbucket/rbac/role.yaml @@ -8,6 +8,14 @@ rules: - "" resources: - configmaps + verbs: + - create + - get + - list + - watch +- apiGroups: + - "" + resources: - secrets verbs: - get diff --git a/docs/objectbucket-notifications-adapter-configuration.md b/docs/objectbucket-notifications-adapter-configuration.md new file mode 100644 index 0000000..1bfac3d --- /dev/null +++ b/docs/objectbucket-notifications-adapter-configuration.md @@ -0,0 +1,203 @@ +# ObjectBucket Notifications Adapter Configuration + +The `objectbucket-notifications-adapter` supports runtime configuration through a Kubernetes ConfigMap. This allows you to modify adapter settings without restarting the pod. + +## Configuration Overview + +The adapter uses two types of configuration: + +### Static Configuration (Command-line flags) + +These settings are provided at pod startup and require a pod restart to change: + +| Flag | Default | Description | +|------|---------|-------------| +| `--config` | `objectbucket-notifications-adapter-config` | Name of the ConfigMap containing adapter configuration | +| `--create-config` | `false` | If set, create the default configuration ConfigMap at startup when it does not already exist (requires `create` permission on ConfigMaps) | +| `--adapter-port` | `8888` | Port the notification HTTP server listens on (HTTP mode only) | +| `--notifications-mode` | `http` | Default for the `NOTIFICATIONS_MODE` ConfigMap key (see below) | +| `--kafka-brokers` | _(none)_ | Default for the `KAFKA_BROKERS` ConfigMap key (see below) | +| `--kafka-notifications-topics` | _(none)_ | Default for the `KAFKA_NOTIFICATIONS_TOPICS` ConfigMap key (see below) | +| `--kafka-notifications-group-id` | _(none)_ | Default for the `KAFKA_NOTIFICATIONS_GROUP_ID` ConfigMap key (see below) | + +> **Note:** The notification transport settings (`--notifications-mode`, `--kafka-brokers`, +> `--kafka-notifications-topics`, `--kafka-notifications-group-id`) are now **dynamic**. The +> command-line flags only supply the *defaults*; the effective values come from the ConfigMap +> and can be changed at runtime. When any Kafka setting changes, the adapter gracefully +> restarts its Kafka consumer with the new settings. + +### Dynamic Configuration (ConfigMap) + +These settings can be changed at runtime by modifying the ConfigMap. The adapter watches for changes and reloads automatically: + +| ConfigMap Key | Default | Description | +|---------------|---------|-------------| +| `NOOBAA_ADAPTER_ID` | `mcg-adapter` | Identifier used in the S3 bucket notification configuration for NooBaa-managed OBCs | +| `NOOBAA_ADAPTER_TOPIC_ARN` | `mcg-adapter-connection/connect.json` | NooBaa connection secret reference used as TopicArn in put-bucket-notification calls | +| `NOOBAA_ADAPTER_STORAGECLASS_PATTERN` | `.*noobaa\.io$` | Regex matched against OBC `spec.storageClassName` to classify as NooBaa-managed | +| `RADOSGW_ADAPTER_ID` | `rgw-adapter` | Identifier used in the S3 bucket notification configuration for RadosGW-managed OBCs | +| `RADOSGW_ADAPTER_TOPIC_ARN` | `arn:aws:sns:ocs-storagecluster-cephobjectstore::rgw-adapter-notifications` | RadosGW SNS TopicArn used in put-bucket-notification calls | +| `RADOSGW_ADAPTER_STORAGECLASS_PATTERN` | `.*ceph-rgw$` | Regex matched against OBC `spec.storageClassName` to classify as RadosGW-managed | +| `NOTIFICATIONS_MODE` | value of `--notifications-mode` (`http`) | `http` or `kafka` — selects how the adapter receives NooBaa/RadosGW notifications. Switching modes restarts the notification runner. | +| `KAFKA_BROKERS` | value of `--kafka-brokers` | Comma-separated list of Kafka broker addresses (required for Kafka mode). Changing it gracefully restarts the Kafka consumer. | +| `KAFKA_NOTIFICATIONS_TOPICS` | value of `--kafka-notifications-topics` | Comma-separated list of Kafka topics to consume notifications from (required for Kafka mode). Changing it gracefully restarts the Kafka consumer. | +| `KAFKA_NOTIFICATIONS_GROUP_ID` | value of `--kafka-notifications-group-id` | Consumer group ID for consuming notifications (required for Kafka mode). Changing it gracefully restarts the Kafka consumer. | +| `KAFKA_SECRET` | _(none)_ | Name of a Kubernetes Secret (in the adapter's namespace) containing Kafka credentials. See `kafka-secret-format.md`. | + +## Example ConfigMap + +```yaml +apiVersion: v1 +kind: ConfigMap +metadata: + name: objectbucket-notifications-adapter-config + namespace: your-namespace +data: + # NooBaa adapter configuration + NOOBAA_ADAPTER_ID: "mcg-adapter" + NOOBAA_ADAPTER_TOPIC_ARN: "mcg-adapter-connection/connect.json" + NOOBAA_ADAPTER_STORAGECLASS_PATTERN: ".*noobaa\\.io$" + + # RadosGW adapter configuration + RADOSGW_ADAPTER_ID: "rgw-adapter" + RADOSGW_ADAPTER_TOPIC_ARN: "arn:aws:sns:ocs-storagecluster-cephobjectstore::rgw-adapter-notifications" + RADOSGW_ADAPTER_STORAGECLASS_PATTERN: ".*ceph-rgw$" + + # Optional: Kafka secret for authentication + KAFKA_SECRET: "kafka-credentials" +``` + +## Deployment + +### 1. Create the ConfigMap + +Create the ConfigMap in the same namespace as the adapter: + +```bash +kubectl apply -f config/samples/objectbucket-notifications-adapter-config.yaml -n your-namespace +``` + +### 2. Deploy the Adapter + +When deploying the adapter, specify the ConfigMap name and static configuration: + +```yaml +apiVersion: apps/v1 +kind: Deployment +metadata: + name: objectbucket-notifications-adapter +spec: + template: + spec: + containers: + - name: adapter + image: your-registry/objectbucket-notifications-adapter:latest + args: + - --config=objectbucket-notifications-adapter-config + - --adapter-port=8888 + - --notifications-mode=http + # For Kafka mode, add: + # - --notifications-mode=kafka + # - --kafka-brokers=broker1:9092,broker2:9092 + # - --kafka-notifications-topics=mcg-notifications,rgw-notifications + # - --kafka-notifications-group-id=adapter-consumer-group + env: + - name: POD_NAMESPACE + valueFrom: + fieldRef: + fieldPath: metadata.namespace +``` + +## Runtime Configuration Updates + +To update the adapter configuration at runtime: + +1. Edit the ConfigMap: + ```bash + kubectl edit configmap objectbucket-notifications-adapter-config -n your-namespace + ``` + +2. The adapter will detect the change and reload the configuration automatically. You'll see log messages like: + ``` + ConfigMap changed, reloading configuration + configuration reloaded successfully + ``` + +3. The new configuration applies immediately to new reconciliation loops. Existing ObjectBucketSource resources will use the updated configuration on their next reconciliation. + +## Dynamic Notification Transport + +The notification transport settings (`NOTIFICATIONS_MODE`, `KAFKA_BROKERS`, +`KAFKA_NOTIFICATIONS_TOPICS`, `KAFKA_NOTIFICATIONS_GROUP_ID`) can be changed at runtime +via the ConfigMap. The same applies to the Kafka credentials referenced by `KAFKA_SECRET`. +When any of them change, the adapter: + +1. Stops the current notification runner (the HTTP server or the Kafka consumer), waiting + for it to shut down gracefully. +2. Starts a new runner using the updated settings/credentials. + +This means you can, for example, switch the adapter from `http` to `kafka` mode, point the +consumer at different brokers, subscribe to different topics, change the consumer group ID, +or rotate the Kafka credentials — all without restarting the pod. + +The runner is restarted only when a change actually affects it: changes to unrelated +ConfigMap keys (e.g. adapter IDs) do not restart it, and a Kafka credential change only +restarts the runner when Kafka is actually in use (`kafka` mode, or `http` mode with +`KAFKA_BROKERS` set for `kafka:` sinks). + +If a ConfigMap change produces an invalid notification configuration (for example +`NOTIFICATIONS_MODE=kafka` without `KAFKA_BROKERS`), the reload is rejected: the adapter +logs an error and keeps running with the previous valid configuration. + +## Configuration Validation + +The adapter validates the configuration when loading: + +- Storage class patterns must be valid regular expressions +- All required fields have sensible defaults +- If `KAFKA_SECRET` is specified, the secret must exist in the adapter's namespace + +If validation fails, the adapter logs an error and continues using the previous valid configuration. + +## Kafka Credential Rotation + +To rotate Kafka credentials without restarting the adapter pod: + +1. Update the Kafka Secret with new credentials in place, **or** update the ConfigMap to + reference a new secret via `KAFKA_SECRET`. +2. The adapter watches both the ConfigMap and the referenced Kafka Secret, so it detects the + change automatically, rebuilds the Kafka configuration, and — if Kafka is in use — + gracefully restarts its Kafka producer/consumer so the new credentials take effect + immediately. + +Note: The adapter watches the Secret named by the current `KAFKA_SECRET`. If you point +`KAFKA_SECRET` at a different Secret, the watcher automatically follows the new reference. + +## Migration from Environment Variables + +If you're migrating from the old environment variable configuration: + +1. Create a ConfigMap with values from your current environment variables +2. Update your deployment to use command-line flags instead of environment variables for static settings +3. Remove the environment variable definitions from your deployment +4. Dynamic settings (adapter IDs, topics, patterns) can now be updated via the ConfigMap + +Example migration: + +**Old (env vars):** +```yaml +env: +- name: NOOBAA_ADAPTER_ID + value: "mcg-adapter" +- name: ADAPTER_PORT + value: "8888" +``` + +**New (ConfigMap + flags):** +```yaml +args: +- --adapter-port=8888 +- --config=objectbucket-notifications-adapter-config +``` + +And create a ConfigMap with `NOOBAA_ADAPTER_ID: "mcg-adapter"`. diff --git a/go.mod b/go.mod index 2b1561d..b1f96a5 100644 --- a/go.mod +++ b/go.mod @@ -5,9 +5,9 @@ go 1.26.0 require ( code.gitea.io/sdk/gitea v0.25.1 github.com/IBM/sarama v1.60.1 - github.com/aws/aws-sdk-go-v2 v1.43.4 - github.com/aws/aws-sdk-go-v2/credentials v1.19.34 - github.com/aws/aws-sdk-go-v2/service/s3 v1.107.0 + github.com/aws/aws-sdk-go-v2 v1.43.5 + github.com/aws/aws-sdk-go-v2/credentials v1.19.35 + github.com/aws/aws-sdk-go-v2/service/s3 v1.107.1 github.com/cloudevents/sdk-go/v2 v2.16.2 github.com/go-git/go-git/v6 v6.0.0-alpha.5 github.com/go-logr/logr v1.4.4 @@ -17,7 +17,7 @@ require ( github.com/prometheus/client_golang v1.24.1 github.com/stretchr/testify v1.11.1 github.com/xdg-go/scram v1.2.0 - golang.org/x/crypto v0.54.0 + golang.org/x/crypto v0.55.0 gopkg.in/yaml.v3 v3.0.1 k8s.io/api v0.36.3 k8s.io/apimachinery v0.36.3 @@ -36,15 +36,15 @@ require ( github.com/Microsoft/go-winio v0.6.2 // indirect github.com/ProtonMail/go-crypto v1.4.1 // indirect github.com/antlr4-go/antlr/v4 v4.13.1 // indirect - github.com/aws/aws-sdk-go-v2/aws/protocol/eventstream v1.7.16 // indirect - github.com/aws/aws-sdk-go-v2/internal/configsources v1.4.35 // indirect - github.com/aws/aws-sdk-go-v2/internal/endpoints/v2 v2.7.35 // indirect - github.com/aws/aws-sdk-go-v2/internal/v4a v1.4.36 // indirect - github.com/aws/aws-sdk-go-v2/service/internal/accept-encoding v1.13.15 // indirect - github.com/aws/aws-sdk-go-v2/service/internal/checksum v1.9.28 // indirect - github.com/aws/aws-sdk-go-v2/service/internal/presigned-url v1.13.35 // indirect - github.com/aws/aws-sdk-go-v2/service/internal/s3shared v1.19.36 // indirect - github.com/aws/smithy-go v1.27.6 // indirect + github.com/aws/aws-sdk-go-v2/aws/protocol/eventstream v1.7.17 // indirect + github.com/aws/aws-sdk-go-v2/internal/configsources v1.4.36 // indirect + github.com/aws/aws-sdk-go-v2/internal/endpoints/v2 v2.7.36 // indirect + github.com/aws/aws-sdk-go-v2/internal/v4a v1.4.37 // indirect + github.com/aws/aws-sdk-go-v2/service/internal/accept-encoding v1.13.16 // indirect + github.com/aws/aws-sdk-go-v2/service/internal/checksum v1.9.29 // indirect + github.com/aws/aws-sdk-go-v2/service/internal/presigned-url v1.13.36 // indirect + github.com/aws/aws-sdk-go-v2/service/internal/s3shared v1.19.37 // indirect + github.com/aws/smithy-go v1.27.7 // indirect github.com/beorn7/perks v1.0.1 // indirect github.com/blang/semver/v4 v4.0.0 // indirect github.com/cenkalti/backoff/v5 v5.0.3 // indirect @@ -148,15 +148,15 @@ require ( go.yaml.in/yaml/v2 v2.4.4 // indirect go.yaml.in/yaml/v3 v3.0.4 // indirect golang.org/x/exp v0.0.0-20260410095643-746e56fc9e2f // indirect - golang.org/x/mod v0.37.0 // indirect + golang.org/x/mod v0.38.0 // indirect golang.org/x/net v0.57.0 // indirect golang.org/x/oauth2 v0.36.0 // indirect golang.org/x/sync v0.22.0 // indirect golang.org/x/sys v0.47.0 // indirect golang.org/x/term v0.45.0 // indirect - golang.org/x/text v0.40.0 // indirect + golang.org/x/text v0.41.0 // indirect golang.org/x/time v0.15.0 // indirect - golang.org/x/tools v0.47.0 // indirect + golang.org/x/tools v0.48.0 // indirect gomodules.xyz/jsonpatch/v2 v2.5.0 // indirect google.golang.org/genproto/googleapis/api v0.0.0-20260414002931-afd174a4e478 // indirect google.golang.org/genproto/googleapis/rpc v0.0.0-20260414002931-afd174a4e478 // indirect diff --git a/go.sum b/go.sum index 23d3be6..ed51f7f 100644 --- a/go.sum +++ b/go.sum @@ -29,30 +29,30 @@ github.com/antlr4-go/antlr/v4 v4.13.1 h1:SqQKkuVZ+zWkMMNkjy5FZe5mr5WURWnlpmOuzYW github.com/antlr4-go/antlr/v4 v4.13.1/go.mod h1:GKmUxMtwp6ZgGwZSva4eWPC5mS6vUAmOABFgjdkM7Nw= github.com/armon/go-socks5 v0.0.0-20160902184237-e75332964ef5 h1:0CwZNZbxp69SHPdPJAN/hZIm0C4OItdklCFmMRWYpio= github.com/armon/go-socks5 v0.0.0-20160902184237-e75332964ef5/go.mod h1:wHh0iHkYZB8zMSxRWpUBQtwG5a7fFgvEO+odwuTv2gs= -github.com/aws/aws-sdk-go-v2 v1.43.4 h1:b9FTvbRwy+JCsfp2Wp6wV/KbOx3Aj7nkoFb2cRX0IhE= -github.com/aws/aws-sdk-go-v2 v1.43.4/go.mod h1:70vwSy16txshwG+g55WkpgPKDIByzHI8ccBsOteo3bQ= -github.com/aws/aws-sdk-go-v2/aws/protocol/eventstream v1.7.16 h1:aiuaKlDweRC5qExJondpWjOgyzMHpofpwspGXUtwn4c= -github.com/aws/aws-sdk-go-v2/aws/protocol/eventstream v1.7.16/go.mod h1:nG/LOlmox9BDe9HvQnXWzgcK8uKbgBMZ/Hp5pVt/21I= -github.com/aws/aws-sdk-go-v2/credentials v1.19.34 h1:y6GkSmcv5myd1ngrYbGmiLlwQqB6TQhOuN/tbSSuWDY= -github.com/aws/aws-sdk-go-v2/credentials v1.19.34/go.mod h1:w3dTcnDVoQIewjo7JG45hduAToikiIFLC4FIO7fndvw= -github.com/aws/aws-sdk-go-v2/internal/configsources v1.4.35 h1:kzVuGlatQtYinwBJEEyLAbggepCoavosiaHHX9+fD+c= -github.com/aws/aws-sdk-go-v2/internal/configsources v1.4.35/go.mod h1:0yLx0yEI+SfqeJMPvOtIEFoZbiQYXMGszBueiutQyaI= -github.com/aws/aws-sdk-go-v2/internal/endpoints/v2 v2.7.35 h1:WK6CjihTuLisCjSKKbildJ79sGZZgbBz3iNa7VsKIhU= -github.com/aws/aws-sdk-go-v2/internal/endpoints/v2 v2.7.35/go.mod h1:KYleN57luLoe97R7vTnx8PMcVrr9gAcRECtOjl91DNg= -github.com/aws/aws-sdk-go-v2/internal/v4a v1.4.36 h1:jbGY4CXLzZElOXgGsexlC3Hi+3YM0rSmk4opFXKqg/k= -github.com/aws/aws-sdk-go-v2/internal/v4a v1.4.36/go.mod h1:uBu/9aKsS/UQGc72RAt3y54kjgYQxmhut8ZD2dXCDNE= -github.com/aws/aws-sdk-go-v2/service/internal/accept-encoding v1.13.15 h1:JJLBQxwY+AFwuPAi5ivGc1ChnTdUt4cXMv7e76m2c/Y= -github.com/aws/aws-sdk-go-v2/service/internal/accept-encoding v1.13.15/go.mod h1:lQknBIe78MVL0cQOQDlag8KGflMbMEVFx9mB6O8ENvk= -github.com/aws/aws-sdk-go-v2/service/internal/checksum v1.9.28 h1:Q1TF1J9jVD+vFo0LzNnmNdQ9EAt52TS+MQlq9Ir+Yxo= -github.com/aws/aws-sdk-go-v2/service/internal/checksum v1.9.28/go.mod h1:4KqXXC/p1hrotmouDFbrRoWaLy962b9PMUReCG6+uWo= -github.com/aws/aws-sdk-go-v2/service/internal/presigned-url v1.13.35 h1:BBEElKh4a+rKshvjrfpajTe9CbpZvrbb4Jkg2PB7RzA= -github.com/aws/aws-sdk-go-v2/service/internal/presigned-url v1.13.35/go.mod h1:zaZk983w//8beSruBVec/mr4CmDwgZitW/qzGhAAX0g= -github.com/aws/aws-sdk-go-v2/service/internal/s3shared v1.19.36 h1:EUIwBoN+q7UmhAejxgD27APiRjh1vwCFo53gSqdT0BM= -github.com/aws/aws-sdk-go-v2/service/internal/s3shared v1.19.36/go.mod h1:6u00gmlTGR6W0b2k9NBrld7MnOEmf1Spqx0VVt6AqyE= -github.com/aws/aws-sdk-go-v2/service/s3 v1.107.0 h1:OkYV+1171za+ab9otU1tGxMXhx6uZvwVEtVddjLuYTg= -github.com/aws/aws-sdk-go-v2/service/s3 v1.107.0/go.mod h1:5FTZoQxhmLEiCAtYVk6V+t0iS/B5yGZVLZ3Wq5FDJZI= -github.com/aws/smithy-go v1.27.6 h1:0zjT8jgK3jbrTT7JJ3EE6JsMhX8JTrZ+f1sEndYDXrA= -github.com/aws/smithy-go v1.27.6/go.mod h1:YE2RhdIuDbA5E5bTdciG9KrW3+TiEONeUWCqxX9i1Fc= +github.com/aws/aws-sdk-go-v2 v1.43.5 h1:yKT5GYnFWhuDo+DqKvE5ZPwVn3RjC4MAeBtZGlh6AVM= +github.com/aws/aws-sdk-go-v2 v1.43.5/go.mod h1:wZjAJppCntyOGgVSmgVTfDyRJK5PHOasO6Wsy8U7Axk= +github.com/aws/aws-sdk-go-v2/aws/protocol/eventstream v1.7.17 h1:mn+Vxb9zgz/FE/yDTcFim3DZ1qpcrxR+qBQkBrl6bzA= +github.com/aws/aws-sdk-go-v2/aws/protocol/eventstream v1.7.17/go.mod h1:eDfmEFxu+BSVsUGLbzJhWjpOurv1mqczClS97yI8wdk= +github.com/aws/aws-sdk-go-v2/credentials v1.19.35 h1:Cxua2RVdRwL0sfjHM/SnQoOnQ7xKng9m5EQBO8BnZlg= +github.com/aws/aws-sdk-go-v2/credentials v1.19.35/go.mod h1:9XQ+RSIGPkycr+oCJYnB1uTv5kMVVR+rd2vYK0Hxj2w= +github.com/aws/aws-sdk-go-v2/internal/configsources v1.4.36 h1:5CrzwxDqf4w3x1Vs3/NiZ0nsC34Hbm3pIDMWbsLebOE= +github.com/aws/aws-sdk-go-v2/internal/configsources v1.4.36/go.mod h1:A3gHdKZIvG/QXERzZwcxNS3RNDFcRCuhhTFBYp+V/nw= +github.com/aws/aws-sdk-go-v2/internal/endpoints/v2 v2.7.36 h1:A4N2f4YPcST0v+dWtX+xrpPPCL9VTBhoIFFUWYqbacE= +github.com/aws/aws-sdk-go-v2/internal/endpoints/v2 v2.7.36/go.mod h1:B/Qr859uxWUEfZeGotK5KAEoof4Q9YWgNtPSwV6jcyk= +github.com/aws/aws-sdk-go-v2/internal/v4a v1.4.37 h1:oyd3ke4V9AhKcRR7rRgxk1VyI+DjK2CBQtbxh3OkdaA= +github.com/aws/aws-sdk-go-v2/internal/v4a v1.4.37/go.mod h1:aA9D7SqfG9IC1b7FLD7Iyc8Q4JN0a8gHhNjN4zPlIaI= +github.com/aws/aws-sdk-go-v2/service/internal/accept-encoding v1.13.16 h1:iE4NGbvqUZnHDqddQAauZzCILYtFjOHwRM5MOOKLB5A= +github.com/aws/aws-sdk-go-v2/service/internal/accept-encoding v1.13.16/go.mod h1:VsjEgrP+ibcou8TlWA4tYaB+0OojuhirsmCe+U60hTA= +github.com/aws/aws-sdk-go-v2/service/internal/checksum v1.9.29 h1:E65Hj648dOV6FuUfI0mYXXhQRHbsi7n+B9h6fZPJO/E= +github.com/aws/aws-sdk-go-v2/service/internal/checksum v1.9.29/go.mod h1:xLrF9yNTCs92VZSpdEd68EJbgcdw3SMR74RO6QDzWHE= +github.com/aws/aws-sdk-go-v2/service/internal/presigned-url v1.13.36 h1:fx2ujmozWn+C/GtfXfz5k6Ckzza40ElOpIW7d92fLWQ= +github.com/aws/aws-sdk-go-v2/service/internal/presigned-url v1.13.36/go.mod h1:QT2ufGVJ+xTRxtXPHTQ1kHkAdWIKPCmD+BqYAXWv8/4= +github.com/aws/aws-sdk-go-v2/service/internal/s3shared v1.19.37 h1:KGHa9iZCrgtkOsFfXb0S4ywsjostA/hau7WE9aSb43E= +github.com/aws/aws-sdk-go-v2/service/internal/s3shared v1.19.37/go.mod h1:FV79f0DSnZIEGsQjWenENGtUycrasyAaJZO+zRanLHA= +github.com/aws/aws-sdk-go-v2/service/s3 v1.107.1 h1:VUTtUJMuRNMkb/7NIKmd8NQaeQLPGCMoTJxkYKre4qM= +github.com/aws/aws-sdk-go-v2/service/s3 v1.107.1/go.mod h1:WvUaO0lP5GNMs1R6cs6qvB3mqo16GLta8yfOuf55Rpc= +github.com/aws/smithy-go v1.27.7 h1:Zgj5z4LfcDYoQIVk+n/yGdTkP/2y6ZT5vYxe0fp7bqE= +github.com/aws/smithy-go v1.27.7/go.mod h1:YE2RhdIuDbA5E5bTdciG9KrW3+TiEONeUWCqxX9i1Fc= github.com/beorn7/perks v1.0.1 h1:VlbKKnNfV8bJzeqoa4cOKqO6bYr3WgKZxO8Z16+hsOM= github.com/beorn7/perks v1.0.1/go.mod h1:G2ZrVWU2WbWT9wwq4/hrbKbnv/1ERSJQ0ibhJ6rlkpw= github.com/blang/semver/v4 v4.0.0 h1:1PFHFE6yCCTv8C1TeyNNarDzntLi7wMI5i/pzqYIsAM= @@ -469,13 +469,13 @@ golang.org/x/crypto v0.0.0-20210513164829-c07d793c2f9a/go.mod h1:P+XmwS30IXTQdn5 golang.org/x/crypto v0.0.0-20210921155107-089bfa567519/go.mod h1:GvvjBRRGRdwPK5ydBHafDWAxML/pGHZbMvKqRZ5+Abc= golang.org/x/crypto v0.0.0-20220622213112-05595931fe9d/go.mod h1:IxCIyHEi3zRg3s0A5j5BB6A9Jmi73HwBIUl50j+osU4= golang.org/x/crypto v0.6.0/go.mod h1:OFC/31mSvZgRz0V1QTNCzfAI1aIRzbiufJtkMIlEp58= -golang.org/x/crypto v0.54.0 h1:YLIA59K4fiNzHzjnZt2tUJQjQtUWfWbeHBqKtk3eScw= -golang.org/x/crypto v0.54.0/go.mod h1:KWL8ny2AZdGR2cWmzeHrp2azQPGogOv+HeQaVEXC2dk= +golang.org/x/crypto v0.55.0 h1:+KWHjbgOaAQ66dh/YlkZKHlz9ZUlq61AFirAR9ntP8M= +golang.org/x/crypto v0.55.0/go.mod h1:uq0V9dE/fzQuJtbnL+2EhWOE63vo164FY8xqEnV9xis= golang.org/x/exp v0.0.0-20260410095643-746e56fc9e2f h1:W3F4c+6OLc6H2lb//N1q4WpJkhzJCK5J6kUi1NTVXfM= golang.org/x/exp v0.0.0-20260410095643-746e56fc9e2f/go.mod h1:J1xhfL/vlindoeF/aINzNzt2Bket5bjo9sdOYzOsU80= golang.org/x/mod v0.6.0-dev.0.20220419223038-86c51ed26bb4/go.mod h1:jJ57K6gSWd91VN4djpZkiMVwK6gcyfeH4XE8wZrZaV4= -golang.org/x/mod v0.37.0 h1:vF1DjpVEshcIqoEaauuHebaLk1O1forxjxBaVn884JQ= -golang.org/x/mod v0.37.0/go.mod h1:m8S8VeM9r4dzDwjrKO0a1sZP3YjeMamRRlD+fmR2Q/0= +golang.org/x/mod v0.38.0 h1:MECBjubtXD7yj4HrhIUcywNaGeNVUdfVnxmPajOk4yk= +golang.org/x/mod v0.38.0/go.mod h1:V6Xz0pq8TQ3dGqVQ1FVHuelZpAL0uNhSkk9ogYP3c40= golang.org/x/net v0.0.0-20190404232315-eb5bcb51f2a3/go.mod h1:t9HGtf8HONx5eT2rtn7q6eTqICYqUVnKs3thJo3Qplg= golang.org/x/net v0.0.0-20190620200207-3b0461eec859/go.mod h1:z5CRVTTTmAJ677TzLLGU+0bjPO0LkuOLi4/5GtJWs/s= golang.org/x/net v0.0.0-20200114155413-6afb5195e5aa/go.mod h1:z5CRVTTTmAJ677TzLLGU+0bjPO0LkuOLi4/5GtJWs/s= @@ -516,15 +516,15 @@ golang.org/x/text v0.3.6/go.mod h1:5Zoc/QRtKVWzQhOtBMvqHzDpF6irO9z98xDceosuGiQ= golang.org/x/text v0.3.7/go.mod h1:u+2+/6zg+i71rQMx5EYifcz6MCKuco9NR6JIITiCfzQ= golang.org/x/text v0.3.8/go.mod h1:E6s5w1FMmriuDzIBO73fBruAKo1PCIq6d2Q6DHfQ8WQ= golang.org/x/text v0.7.0/go.mod h1:mrYo+phRRbMaCq/xk9113O4dZlRixOauAjOtrjsXDZ8= -golang.org/x/text v0.40.0 h1:Ub2Z6/xjgF1WrYQz2nuITOEegKFtiIy+rieRJ5lHZKs= -golang.org/x/text v0.40.0/go.mod h1:hpnzDAfGV753zIKo+wk3u1bVKCGPbrnF7+7LBF/UHVY= +golang.org/x/text v0.41.0 h1:vz/seA0lnX87Othu2f/0L24RcgrXD9/YFTSuGjj3rH8= +golang.org/x/text v0.41.0/go.mod h1:jvf1O8ajNzZqhSrQBPbutR/EB83Cc0CFrezNQIwbb5M= golang.org/x/time v0.15.0 h1:bbrp8t3bGUeFOx08pvsMYRTCVSMk89u4tKbNOZbp88U= golang.org/x/time v0.15.0/go.mod h1:Y4YMaQmXwGQZoFaVFk4YpCt4FLQMYKZe9oeV/f4MSno= golang.org/x/tools v0.0.0-20180917221912-90fa682c2a6e/go.mod h1:n7NCudcB/nEzxVGmLbDWY5pfWTLqBcC2KZ6jyYvM4mQ= golang.org/x/tools v0.0.0-20191119224855-298f0cb1881e/go.mod h1:b+2E5dAYhXwXZwtnZ6UAqBI28+e2cm9otk0dWdXHAEo= golang.org/x/tools v0.1.12/go.mod h1:hNGJHUnrk76NpqgfD5Aqm5Crs+Hm0VOH/i9J2+nxYbc= -golang.org/x/tools v0.47.0 h1:7Kn5x/d1svx/PzryTsqeoZN4TZwqeH5pGWjefhLi/1Q= -golang.org/x/tools v0.47.0/go.mod h1:dFHnyTvFWY212G+h7ZY4Vsp/K3U4/7W9TyVaAul8uCA= +golang.org/x/tools v0.48.0 h1:3+hClM1aLL5mjMKm5ovokw9epgRXPuu2tILgismM6RE= +golang.org/x/tools v0.48.0/go.mod h1:08xX0orndb/F7jJxGDicx061tyd5pcMto75YMAXr6lk= golang.org/x/xerrors v0.0.0-20190717185122-a985d3407aa7/go.mod h1:I/5z698sn9Ka8TeJc9MKroUUfqBBauWjQqLJ2OPfmY0= gomodules.xyz/jsonpatch/v2 v2.5.0 h1:JELs8RLM12qJGXU4u/TO3V25KW8GreMKl9pdkk14RM0= gomodules.xyz/jsonpatch/v2 v2.5.0/go.mod h1:AH3dM2RI6uoBZxn3LVrfvJ3E0/9dG4cSrbuBJT4moAY= diff --git a/internal/objectbucketsource/config/config.go b/internal/objectbucketsource/config/config.go new file mode 100644 index 0000000..7d75af1 --- /dev/null +++ b/internal/objectbucketsource/config/config.go @@ -0,0 +1,585 @@ +package config + +import ( + "context" + "crypto/sha256" + "encoding/hex" + "fmt" + "regexp" + "sort" + "strings" + "sync" + "time" + + corev1 "k8s.io/api/core/v1" + apierrors "k8s.io/apimachinery/pkg/api/errors" + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" + "k8s.io/apimachinery/pkg/watch" + "k8s.io/client-go/kubernetes" + ctrl "sigs.k8s.io/controller-runtime" + logf "sigs.k8s.io/controller-runtime/pkg/log" + + "github.com/IBM/sarama" + kafkaconfig "github.com/functions-dev/func-operator/internal/objectbucketsource/kafka" +) + +var log = logf.Log.WithName("adapter-config") + +// Default values for the adapter configuration. They are used both when resolving +// missing ConfigMap keys and when creating the default ConfigMap (--create-config). +const ( + defaultNoobaaAdapterID = "mcg-adapter" + defaultNoobaaTopicARN = "mcg-adapter-connection/connect.json" + defaultNoobaaStorageClassRE = `.*noobaa\.io$` + defaultRadosgwAdapterID = "rgw-adapter" + defaultRadosgwTopicARN = "arn:aws:sns:ocs-storagecluster-cephobjectstore::rgw-adapter-notifications" + defaultRadosgwStorageClassRE = `.*ceph-rgw$` +) + +// AdapterBackendConfig holds configuration for a single storage backend adapter +type AdapterBackendConfig struct { + ID string + TopicARN string + StorageClassPattern *regexp.Regexp +} + +// NotificationSettings holds the transport configuration that controls how the +// adapter receives NooBaa/RadosGW notifications. These settings can be changed at +// runtime via the ConfigMap; the notification server restarts its Kafka consumer +// when any of them change. +type NotificationSettings struct { + // Mode is "http" or "kafka". + Mode string + KafkaBrokers []string + KafkaNotificationsTopics []string + KafkaNotificationsGroupID string +} + +// Config holds runtime-configurable settings for the objectbucket-notifications-adapter +type Config struct { + NoobaaAdapter AdapterBackendConfig + RadosgwAdapter AdapterBackendConfig + Notifications NotificationSettings +} + +// Provider provides access to the current configuration and watches for changes +type Provider struct { + mu sync.RWMutex + config Config + + namespace string + configMapName string + clientset *kubernetes.Clientset + cancelWatch context.CancelFunc + + // defaults holds the notification settings supplied via command-line flags. + // They are used whenever the corresponding ConfigMap keys are absent. + defaults NotificationSettings + + kafkaConfigMu sync.RWMutex + kafkaConfig *sarama.Config + kafkaSecret string + // kafkaFingerprint is a content hash of the Kafka secret name and its data. + // It changes whenever the Kafka credentials/connection settings change, which + // lets consumers detect credential rotation and restart their connections. + kafkaFingerprint string + + subscribersMu sync.Mutex + subscribers []chan struct{} + + // reloadSignal is fired internally after every successful reload so the + // Secret watcher can re-evaluate which Secret it should be watching. + reloadSignal chan struct{} +} + +// NewProvider creates a new configuration provider that watches a ConfigMap. +// The defaults are used for any notification settings not present in the ConfigMap. +func NewProvider(ctx context.Context, namespace, configMapName string, defaults NotificationSettings) (*Provider, error) { + clientset, err := kubernetes.NewForConfig(ctrl.GetConfigOrDie()) + if err != nil { + return nil, fmt.Errorf("creating kubernetes clientset: %w", err) + } + + p := &Provider{ + namespace: namespace, + configMapName: configMapName, + clientset: clientset, + defaults: defaults, + reloadSignal: make(chan struct{}, 1), + } + + if err := p.loadConfig(ctx); err != nil { + return nil, fmt.Errorf("loading initial config: %w", err) + } + + watchCtx, cancel := context.WithCancel(context.Background()) + p.cancelWatch = cancel + go p.watchConfigMap(watchCtx) + go p.watchSecret(watchCtx) + + return p, nil +} + +// EnsureDefaultConfigMap creates the adapter configuration ConfigMap with default +// values if it does not already exist. It is a no-op when the ConfigMap is already +// present. The notification-related defaults (from the command-line flags) are +// written into the ConfigMap so it reflects the effective startup configuration. +func EnsureDefaultConfigMap(ctx context.Context, namespace, configMapName string, defaults NotificationSettings) error { + clientset, err := kubernetes.NewForConfig(ctrl.GetConfigOrDie()) + if err != nil { + return fmt.Errorf("creating kubernetes clientset: %w", err) + } + + _, err = clientset.CoreV1().ConfigMaps(namespace).Get(ctx, configMapName, metav1.GetOptions{}) + if err == nil { + log.Info("configuration ConfigMap already exists, not creating", "name", configMapName, "namespace", namespace) + return nil + } + if !apierrors.IsNotFound(err) { + return fmt.Errorf("checking for ConfigMap %s/%s: %w", namespace, configMapName, err) + } + + cm := &corev1.ConfigMap{ + ObjectMeta: metav1.ObjectMeta{ + Name: configMapName, + Namespace: namespace, + }, + Data: defaultConfigMapData(defaults), + } + + if _, err := clientset.CoreV1().ConfigMaps(namespace).Create(ctx, cm, metav1.CreateOptions{}); err != nil { + if apierrors.IsAlreadyExists(err) { + // Another replica created it concurrently; treat as success. + log.Info("configuration ConfigMap already created concurrently", "name", configMapName, "namespace", namespace) + return nil + } + return fmt.Errorf("creating default ConfigMap %s/%s: %w", namespace, configMapName, err) + } + + log.Info("created default configuration ConfigMap", "name", configMapName, "namespace", namespace) + return nil +} + +// defaultConfigMapData builds the data map for the default configuration ConfigMap. +func defaultConfigMapData(defaults NotificationSettings) map[string]string { + data := map[string]string{ + "NOOBAA_ADAPTER_ID": defaultNoobaaAdapterID, + "NOOBAA_ADAPTER_TOPIC_ARN": defaultNoobaaTopicARN, + "NOOBAA_ADAPTER_STORAGECLASS_PATTERN": defaultNoobaaStorageClassRE, + "RADOSGW_ADAPTER_ID": defaultRadosgwAdapterID, + "RADOSGW_ADAPTER_TOPIC_ARN": defaultRadosgwTopicARN, + "RADOSGW_ADAPTER_STORAGECLASS_PATTERN": defaultRadosgwStorageClassRE, + } + + mode := defaults.Mode + if mode == "" { + mode = "http" + } + data["NOTIFICATIONS_MODE"] = mode + + if len(defaults.KafkaBrokers) > 0 { + data["KAFKA_BROKERS"] = strings.Join(defaults.KafkaBrokers, ",") + } + if len(defaults.KafkaNotificationsTopics) > 0 { + data["KAFKA_NOTIFICATIONS_TOPICS"] = strings.Join(defaults.KafkaNotificationsTopics, ",") + } + if defaults.KafkaNotificationsGroupID != "" { + data["KAFKA_NOTIFICATIONS_GROUP_ID"] = defaults.KafkaNotificationsGroupID + } + + return data +} + +// GetConfig returns a copy of the current configuration +func (p *Provider) GetConfig() Config { + p.mu.RLock() + defer p.mu.RUnlock() + return p.config +} + +// GetKafkaConfig returns the current Kafka configuration +func (p *Provider) GetKafkaConfig() *sarama.Config { + p.kafkaConfigMu.RLock() + defer p.kafkaConfigMu.RUnlock() + return p.kafkaConfig +} + +// GetKafkaFingerprint returns a content hash of the current Kafka secret. It +// changes whenever the Kafka credentials/connection settings change, allowing +// callers to detect credential rotation. +func (p *Provider) GetKafkaFingerprint() string { + p.kafkaConfigMu.RLock() + defer p.kafkaConfigMu.RUnlock() + return p.kafkaFingerprint +} + +// Subscribe returns a channel that receives a signal whenever the configuration +// is successfully reloaded. The channel is buffered (size 1) and signals are +// coalesced, so a slow subscriber never blocks the config watcher. +func (p *Provider) Subscribe() <-chan struct{} { + ch := make(chan struct{}, 1) + p.subscribersMu.Lock() + p.subscribers = append(p.subscribers, ch) + p.subscribersMu.Unlock() + return ch +} + +func (p *Provider) notifySubscribers() { + p.subscribersMu.Lock() + defer p.subscribersMu.Unlock() + for _, ch := range p.subscribers { + select { + case ch <- struct{}{}: + default: + } + } +} + +// Stop stops watching the ConfigMap +func (p *Provider) Stop() { + if p.cancelWatch != nil { + p.cancelWatch() + } +} + +// NeedLeaderElection implements the manager.Runnable interface +func (p *Provider) NeedLeaderElection() bool { + return false +} + +// Start implements the manager.Runnable interface +func (p *Provider) Start(ctx context.Context) error { + <-ctx.Done() + p.Stop() + return nil +} + +func (p *Provider) loadConfig(ctx context.Context) error { + cm, err := p.clientset.CoreV1().ConfigMaps(p.namespace).Get(ctx, p.configMapName, metav1.GetOptions{}) + if err != nil { + return fmt.Errorf("getting ConfigMap %s/%s: %w", p.namespace, p.configMapName, err) + } + + config, err := parseConfig(cm, p.defaults) + if err != nil { + return fmt.Errorf("parsing ConfigMap: %w", err) + } + + kafkaSecret := cm.Data["KAFKA_SECRET"] + var kafkaCfg *sarama.Config + var secretData map[string][]byte + if kafkaSecret != "" { + secret, err := p.clientset.CoreV1().Secrets(p.namespace).Get(ctx, kafkaSecret, metav1.GetOptions{}) + if err != nil { + return fmt.Errorf("reading KAFKA_SECRET %s/%s: %w", p.namespace, kafkaSecret, err) + } + secretData = secret.Data + kafkaCfg, err = kafkaconfig.NewConfig(secretData) + if err != nil { + return fmt.Errorf("configuring kafka from secret %s: %w", kafkaSecret, err) + } + log.Info("kafka configured from secret", "name", kafkaSecret, "namespace", p.namespace) + } else { + kafkaCfg, err = kafkaconfig.NewConfig(nil) + if err != nil { + return fmt.Errorf("creating default kafka config: %w", err) + } + } + fingerprint := kafkaFingerprint(kafkaSecret, secretData) + + p.mu.Lock() + p.config = config + p.mu.Unlock() + + p.kafkaConfigMu.Lock() + p.kafkaConfig = kafkaCfg + p.kafkaSecret = kafkaSecret + p.kafkaFingerprint = fingerprint + p.kafkaConfigMu.Unlock() + + log.Info("configuration loaded", + "noobaa-adapter-id", config.NoobaaAdapter.ID, + "radosgw-adapter-id", config.RadosgwAdapter.ID, + "notifications-mode", config.Notifications.Mode, + "kafka-brokers", config.Notifications.KafkaBrokers, + "kafka-notifications-topics", config.Notifications.KafkaNotificationsTopics, + "kafka-notifications-group-id", config.Notifications.KafkaNotificationsGroupID) + + p.notifySubscribers() + p.signalReload() + + return nil +} + +// signalReload notifies the Secret watcher (non-blocking) that the configuration +// was reloaded so it can re-evaluate which Secret to watch. +func (p *Provider) signalReload() { + if p.reloadSignal == nil { + return + } + select { + case p.reloadSignal <- struct{}{}: + default: + } +} + +func (p *Provider) currentKafkaSecret() string { + p.kafkaConfigMu.RLock() + defer p.kafkaConfigMu.RUnlock() + return p.kafkaSecret +} + +func (p *Provider) watchConfigMap(ctx context.Context) { + for { + select { + case <-ctx.Done(): + return + default: + } + + watcher, err := p.clientset.CoreV1().ConfigMaps(p.namespace).Watch(ctx, metav1.ListOptions{ + FieldSelector: fmt.Sprintf("metadata.name=%s", p.configMapName), + }) + if err != nil { + log.Error(err, "failed to create ConfigMap watcher, retrying in 5s") + select { + case <-ctx.Done(): + return + case <-time.After(5 * time.Second): + } + continue + } + + log.Info("watching ConfigMap for changes", "name", p.configMapName, "namespace", p.namespace) + + func() { + defer watcher.Stop() + + for { + select { + case <-ctx.Done(): + return + case event, ok := <-watcher.ResultChan(): + if !ok { + log.Info("ConfigMap watch channel closed, restarting watcher") + return + } + + if event.Type == watch.Modified || event.Type == watch.Added { + cm, ok := event.Object.(*corev1.ConfigMap) + if !ok { + log.Error(fmt.Errorf("unexpected object type"), "failed to cast to ConfigMap") + continue + } + + log.Info("ConfigMap changed, reloading configuration", "name", cm.Name) + if err := p.loadConfig(ctx); err != nil { + log.Error(err, "failed to reload configuration") + } else { + log.Info("configuration reloaded successfully") + } + } else if event.Type == watch.Deleted { + log.Error(fmt.Errorf("ConfigMap deleted"), "adapter configuration unavailable", "name", p.configMapName) + } + } + } + }() + } +} + +// watchSecret watches the Kafka Secret currently referenced by KAFKA_SECRET and +// reloads the configuration when its contents change, so an in-place credential +// rotation is picked up without requiring a ConfigMap edit. The watched Secret +// name is dynamic: when a ConfigMap reload changes (or clears) the reference, the +// watcher re-establishes itself against the new Secret. +func (p *Provider) watchSecret(ctx context.Context) { + for { + select { + case <-ctx.Done(): + return + default: + } + + secretName := p.currentKafkaSecret() + if secretName == "" { + // No Secret referenced; wait until a reload might add one. + select { + case <-ctx.Done(): + return + case <-p.reloadSignal: + continue + } + } + + watcher, err := p.clientset.CoreV1().Secrets(p.namespace).Watch(ctx, metav1.ListOptions{ + FieldSelector: fmt.Sprintf("metadata.name=%s", secretName), + }) + if err != nil { + log.Error(err, "failed to create Secret watcher, retrying in 5s", "secret", secretName) + select { + case <-ctx.Done(): + return + case <-time.After(5 * time.Second): + } + continue + } + + log.Info("watching Kafka Secret for changes", "name", secretName, "namespace", p.namespace) + p.consumeSecretEvents(ctx, watcher, secretName) + } +} + +func (p *Provider) consumeSecretEvents(ctx context.Context, watcher watch.Interface, secretName string) { + defer watcher.Stop() + + for { + select { + case <-ctx.Done(): + return + case <-p.reloadSignal: + // A reload happened; the referenced Secret may have changed. If so, + // restart the watcher against the new Secret. + if p.currentKafkaSecret() != secretName { + log.Info("Kafka Secret reference changed, restarting Secret watcher", + "old", secretName, "new", p.currentKafkaSecret()) + return + } + case event, ok := <-watcher.ResultChan(): + if !ok { + log.Info("Secret watch channel closed, restarting watcher", "secret", secretName) + return + } + + switch event.Type { + case watch.Modified, watch.Added: + log.Info("Kafka Secret changed, reloading configuration", "name", secretName) + if err := p.loadConfig(ctx); err != nil { + log.Error(err, "failed to reload configuration after Secret change") + } else { + log.Info("configuration reloaded successfully after Secret change") + } + case watch.Deleted: + log.Error(fmt.Errorf("kafka secret deleted"), "kafka credentials unavailable", "name", secretName) + } + } + } +} + +func parseConfig(cm *corev1.ConfigMap, defaults NotificationSettings) (Config, error) { + noobaaPattern := getOrDefault(cm.Data, "NOOBAA_ADAPTER_STORAGECLASS_PATTERN", defaultNoobaaStorageClassRE) + radosgwPattern := getOrDefault(cm.Data, "RADOSGW_ADAPTER_STORAGECLASS_PATTERN", defaultRadosgwStorageClassRE) + + noobaaRe, err := regexp.Compile(noobaaPattern) + if err != nil { + return Config{}, fmt.Errorf("invalid NOOBAA_ADAPTER_STORAGECLASS_PATTERN: %w", err) + } + + radosgwRe, err := regexp.Compile(radosgwPattern) + if err != nil { + return Config{}, fmt.Errorf("invalid RADOSGW_ADAPTER_STORAGECLASS_PATTERN: %w", err) + } + + notifications, err := parseNotificationSettings(cm.Data, defaults) + if err != nil { + return Config{}, err + } + + config := Config{ + NoobaaAdapter: AdapterBackendConfig{ + ID: getOrDefault(cm.Data, "NOOBAA_ADAPTER_ID", defaultNoobaaAdapterID), + TopicARN: getOrDefault(cm.Data, "NOOBAA_ADAPTER_TOPIC_ARN", defaultNoobaaTopicARN), + StorageClassPattern: noobaaRe, + }, + RadosgwAdapter: AdapterBackendConfig{ + ID: getOrDefault(cm.Data, "RADOSGW_ADAPTER_ID", defaultRadosgwAdapterID), + TopicARN: getOrDefault(cm.Data, "RADOSGW_ADAPTER_TOPIC_ARN", defaultRadosgwTopicARN), + StorageClassPattern: radosgwRe, + }, + Notifications: notifications, + } + + return config, nil +} + +// parseNotificationSettings resolves the notification transport settings from the +// ConfigMap, falling back to the provided defaults (from command-line flags) when +// a key is absent. It validates the resulting settings so that invalid ConfigMap +// changes are rejected and the previous valid configuration is retained. +func parseNotificationSettings(data map[string]string, defaults NotificationSettings) (NotificationSettings, error) { + mode := getOrDefault(data, "NOTIFICATIONS_MODE", defaults.Mode) + if mode == "" { + mode = "http" + } + if mode != "http" && mode != "kafka" { + return NotificationSettings{}, fmt.Errorf("invalid NOTIFICATIONS_MODE %q: must be \"http\" or \"kafka\"", mode) + } + + settings := NotificationSettings{ + Mode: mode, + KafkaBrokers: defaults.KafkaBrokers, + KafkaNotificationsTopics: defaults.KafkaNotificationsTopics, + KafkaNotificationsGroupID: getOrDefault(data, "KAFKA_NOTIFICATIONS_GROUP_ID", defaults.KafkaNotificationsGroupID), + } + if v, ok := data["KAFKA_BROKERS"]; ok && strings.TrimSpace(v) != "" { + settings.KafkaBrokers = splitAndTrim(v) + } + if v, ok := data["KAFKA_NOTIFICATIONS_TOPICS"]; ok && strings.TrimSpace(v) != "" { + settings.KafkaNotificationsTopics = splitAndTrim(v) + } + + if mode == "kafka" { + if len(settings.KafkaBrokers) == 0 { + return NotificationSettings{}, fmt.Errorf("KAFKA_BROKERS is required when NOTIFICATIONS_MODE=kafka") + } + if len(settings.KafkaNotificationsTopics) == 0 { + return NotificationSettings{}, fmt.Errorf("KAFKA_NOTIFICATIONS_TOPICS is required when NOTIFICATIONS_MODE=kafka") + } + if settings.KafkaNotificationsGroupID == "" { + return NotificationSettings{}, fmt.Errorf("KAFKA_NOTIFICATIONS_GROUP_ID is required when NOTIFICATIONS_MODE=kafka") + } + } + + return settings, nil +} + +func getOrDefault(data map[string]string, key, defaultValue string) string { + if v, ok := data[key]; ok && v != "" { + return v + } + return defaultValue +} + +// kafkaFingerprint returns a stable content hash of the Kafka secret name and its +// data. Two calls return the same value iff the secret reference and its contents +// are identical, so a change indicates the Kafka credentials/connection settings +// were rotated. +func kafkaFingerprint(secretName string, data map[string][]byte) string { + h := sha256.New() + h.Write([]byte(secretName)) + h.Write([]byte{0}) + + keys := make([]string, 0, len(data)) + for k := range data { + keys = append(keys, k) + } + sort.Strings(keys) + for _, k := range keys { + h.Write([]byte(k)) + h.Write([]byte{0}) + h.Write(data[k]) + h.Write([]byte{0}) + } + return hex.EncodeToString(h.Sum(nil)) +} + +// splitAndTrim splits a comma-separated string, trimming whitespace and dropping +// empty entries. +func splitAndTrim(s string) []string { + var out []string + for _, part := range strings.Split(s, ",") { + if trimmed := strings.TrimSpace(part); trimmed != "" { + out = append(out, trimmed) + } + } + return out +} diff --git a/internal/objectbucketsource/config/config_test.go b/internal/objectbucketsource/config/config_test.go new file mode 100644 index 0000000..46991d7 --- /dev/null +++ b/internal/objectbucketsource/config/config_test.go @@ -0,0 +1,215 @@ +package config + +import ( + "reflect" + "testing" + + corev1 "k8s.io/api/core/v1" +) + +func TestParseNotificationSettings_DefaultsWhenAbsent(t *testing.T) { + defaults := NotificationSettings{ + Mode: "http", + KafkaBrokers: []string{"b1:9092"}, + KafkaNotificationsTopics: []string{"t1"}, + KafkaNotificationsGroupID: "g1", + } + + got, err := parseNotificationSettings(map[string]string{}, defaults) + if err != nil { + t.Fatalf("unexpected error: %v", err) + } + if !reflect.DeepEqual(got, defaults) { + t.Fatalf("expected defaults %+v, got %+v", defaults, got) + } +} + +func TestParseNotificationSettings_ConfigMapOverrides(t *testing.T) { + defaults := NotificationSettings{Mode: "http"} + data := map[string]string{ + "NOTIFICATIONS_MODE": "kafka", + "KAFKA_BROKERS": "b1:9092, b2:9092 ", + "KAFKA_NOTIFICATIONS_TOPICS": "t1,t2", + "KAFKA_NOTIFICATIONS_GROUP_ID": "grp", + } + + got, err := parseNotificationSettings(data, defaults) + if err != nil { + t.Fatalf("unexpected error: %v", err) + } + + want := NotificationSettings{ + Mode: "kafka", + KafkaBrokers: []string{"b1:9092", "b2:9092"}, + KafkaNotificationsTopics: []string{"t1", "t2"}, + KafkaNotificationsGroupID: "grp", + } + if !reflect.DeepEqual(got, want) { + t.Fatalf("expected %+v, got %+v", want, got) + } +} + +func TestParseNotificationSettings_InvalidMode(t *testing.T) { + _, err := parseNotificationSettings(map[string]string{"NOTIFICATIONS_MODE": "bogus"}, NotificationSettings{}) + if err == nil { + t.Fatal("expected error for invalid mode, got nil") + } +} + +func TestParseNotificationSettings_KafkaRequirements(t *testing.T) { + tests := []struct { + name string + data map[string]string + }{ + { + name: "missing brokers", + data: map[string]string{ + "NOTIFICATIONS_MODE": "kafka", + "KAFKA_NOTIFICATIONS_TOPICS": "t1", + "KAFKA_NOTIFICATIONS_GROUP_ID": "g1", + }, + }, + { + name: "missing topics", + data: map[string]string{ + "NOTIFICATIONS_MODE": "kafka", + "KAFKA_BROKERS": "b1:9092", + "KAFKA_NOTIFICATIONS_GROUP_ID": "g1", + }, + }, + { + name: "missing group id", + data: map[string]string{ + "NOTIFICATIONS_MODE": "kafka", + "KAFKA_BROKERS": "b1:9092", + "KAFKA_NOTIFICATIONS_TOPICS": "t1", + }, + }, + } + + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + if _, err := parseNotificationSettings(tt.data, NotificationSettings{}); err == nil { + t.Fatalf("expected error for %s, got nil", tt.name) + } + }) + } +} + +func TestParseNotificationSettings_EmptyModeDefaultsToHTTP(t *testing.T) { + got, err := parseNotificationSettings(map[string]string{}, NotificationSettings{}) + if err != nil { + t.Fatalf("unexpected error: %v", err) + } + if got.Mode != "http" { + t.Fatalf("expected mode http, got %q", got.Mode) + } +} + +func TestKafkaFingerprint_Stability(t *testing.T) { + data := map[string][]byte{ + "protocol": []byte("SASL_SSL"), + "user": []byte("alice"), + "password": []byte("s3cret"), + } + + // Equal inputs from a separately-constructed map produce the same fingerprint + // (independent of map iteration order). + same := map[string][]byte{ + "password": []byte("s3cret"), + "protocol": []byte("SASL_SSL"), + "user": []byte("alice"), + } + if kafkaFingerprint("kafka-creds", data) != kafkaFingerprint("kafka-creds", same) { + t.Fatal("fingerprint is not stable for equal inputs") + } + + // A changed value produces a different fingerprint (credential rotation). + rotated := map[string][]byte{ + "protocol": []byte("SASL_SSL"), + "user": []byte("alice"), + "password": []byte("rotated"), + } + if kafkaFingerprint("kafka-creds", data) == kafkaFingerprint("kafka-creds", rotated) { + t.Fatal("fingerprint did not change after rotating a value") + } + + // A changed secret name produces a different fingerprint. + if kafkaFingerprint("kafka-creds", data) == kafkaFingerprint("other-creds", data) { + t.Fatal("fingerprint did not change after changing the secret name") + } + + // No secret (nil data, empty name) is stable and non-panicking, and differs + // from any non-empty secret. + if kafkaFingerprint("", nil) == kafkaFingerprint("kafka-creds", data) { + t.Fatal("empty fingerprint collided with a populated secret") + } +} + +func TestDefaultConfigMapData(t *testing.T) { + // With no kafka defaults, only NOTIFICATIONS_MODE is added (defaulting to http), + // and no empty kafka keys are written. + data := defaultConfigMapData(NotificationSettings{}) + if data["NOOBAA_ADAPTER_ID"] != defaultNoobaaAdapterID { + t.Fatalf("expected default noobaa adapter id, got %q", data["NOOBAA_ADAPTER_ID"]) + } + if data["NOTIFICATIONS_MODE"] != "http" { + t.Fatalf("expected NOTIFICATIONS_MODE=http, got %q", data["NOTIFICATIONS_MODE"]) + } + for _, k := range []string{"KAFKA_BROKERS", "KAFKA_NOTIFICATIONS_TOPICS", "KAFKA_NOTIFICATIONS_GROUP_ID"} { + if _, ok := data[k]; ok { + t.Fatalf("did not expect key %q when kafka defaults are empty", k) + } + } + + // With kafka defaults, they are serialized into the ConfigMap data. + data = defaultConfigMapData(NotificationSettings{ + Mode: "kafka", + KafkaBrokers: []string{"b1:9092", "b2:9092"}, + KafkaNotificationsTopics: []string{"t1", "t2"}, + KafkaNotificationsGroupID: "grp", + }) + if data["NOTIFICATIONS_MODE"] != "kafka" { + t.Fatalf("expected NOTIFICATIONS_MODE=kafka, got %q", data["NOTIFICATIONS_MODE"]) + } + if data["KAFKA_BROKERS"] != "b1:9092,b2:9092" { + t.Fatalf("unexpected KAFKA_BROKERS: %q", data["KAFKA_BROKERS"]) + } + if data["KAFKA_NOTIFICATIONS_TOPICS"] != "t1,t2" { + t.Fatalf("unexpected KAFKA_NOTIFICATIONS_TOPICS: %q", data["KAFKA_NOTIFICATIONS_TOPICS"]) + } + if data["KAFKA_NOTIFICATIONS_GROUP_ID"] != "grp" { + t.Fatalf("unexpected KAFKA_NOTIFICATIONS_GROUP_ID: %q", data["KAFKA_NOTIFICATIONS_GROUP_ID"]) + } + + // The produced data must parse back cleanly, yielding the same effective settings. + cfg, err := parseConfig(&corev1.ConfigMap{Data: data}, NotificationSettings{}) + if err != nil { + t.Fatalf("default ConfigMap data failed to parse: %v", err) + } + if cfg.Notifications.Mode != "kafka" { + t.Fatalf("round-trip mode mismatch: %q", cfg.Notifications.Mode) + } +} + +func TestParseConfig_IncludesNotifications(t *testing.T) { + cm := &corev1.ConfigMap{ + Data: map[string]string{ + "NOTIFICATIONS_MODE": "kafka", + "KAFKA_BROKERS": "b1:9092", + "KAFKA_NOTIFICATIONS_TOPICS": "t1", + "KAFKA_NOTIFICATIONS_GROUP_ID": "g1", + }, + } + + cfg, err := parseConfig(cm, NotificationSettings{Mode: "http"}) + if err != nil { + t.Fatalf("unexpected error: %v", err) + } + if cfg.Notifications.Mode != "kafka" { + t.Fatalf("expected kafka mode, got %q", cfg.Notifications.Mode) + } + if cfg.NoobaaAdapter.ID != "mcg-adapter" { + t.Fatalf("expected default noobaa adapter id, got %q", cfg.NoobaaAdapter.ID) + } +} diff --git a/internal/objectbucketsource/config/interface.go b/internal/objectbucketsource/config/interface.go new file mode 100644 index 0000000..7d3238f --- /dev/null +++ b/internal/objectbucketsource/config/interface.go @@ -0,0 +1,9 @@ +package config + +import "github.com/IBM/sarama" + +// ConfigProvider provides access to adapter configuration +type ConfigProvider interface { + GetConfig() Config + GetKafkaConfig() *sarama.Config +} diff --git a/internal/objectbucketsource/config/mock.go b/internal/objectbucketsource/config/mock.go new file mode 100644 index 0000000..fd2dd53 --- /dev/null +++ b/internal/objectbucketsource/config/mock.go @@ -0,0 +1,105 @@ +package config + +import ( + "regexp" + "sync" + + "github.com/IBM/sarama" +) + +// MockProvider is a simple mock implementation of the config provider for testing +type MockProvider struct { + mu sync.RWMutex + config Config + kafkaConfig *sarama.Config + kafkaFingerprint string + subscribers []chan struct{} +} + +// NewMockProvider creates a mock provider with default test configuration +func NewMockProvider() *MockProvider { + noobaaRe := regexp.MustCompile(`.*noobaa\.io$`) + radosgwRe := regexp.MustCompile(`.*ceph-rgw$`) + + return &MockProvider{ + config: Config{ + NoobaaAdapter: AdapterBackendConfig{ + ID: "mcg-adapter", + TopicARN: "mcg-adapter-connection/connect.json", + StorageClassPattern: noobaaRe, + }, + RadosgwAdapter: AdapterBackendConfig{ + ID: "rgw-adapter", + TopicARN: "arn:aws:sns:ocs-storagecluster-cephobjectstore::rgw-adapter-notifications", + StorageClassPattern: radosgwRe, + }, + Notifications: NotificationSettings{ + Mode: "http", + }, + }, + kafkaConfig: sarama.NewConfig(), + } +} + +// GetConfig returns the mock configuration +func (m *MockProvider) GetConfig() Config { + m.mu.RLock() + defer m.mu.RUnlock() + return m.config +} + +// GetKafkaConfig returns the mock Kafka configuration +func (m *MockProvider) GetKafkaConfig() *sarama.Config { + m.mu.RLock() + defer m.mu.RUnlock() + return m.kafkaConfig +} + +// GetKafkaFingerprint returns the mock Kafka fingerprint +func (m *MockProvider) GetKafkaFingerprint() string { + m.mu.RLock() + defer m.mu.RUnlock() + return m.kafkaFingerprint +} + +// Subscribe returns a channel that is signaled whenever the mock configuration changes. +func (m *MockProvider) Subscribe() <-chan struct{} { + ch := make(chan struct{}, 1) + m.mu.Lock() + m.subscribers = append(m.subscribers, ch) + m.mu.Unlock() + return ch +} + +func (m *MockProvider) notify() { + for _, ch := range m.subscribers { + select { + case ch <- struct{}{}: + default: + } + } +} + +// SetConfig updates the mock configuration (for testing) +func (m *MockProvider) SetConfig(cfg Config) { + m.mu.Lock() + defer m.mu.Unlock() + m.config = cfg + m.notify() +} + +// SetKafkaConfig updates the mock Kafka configuration (for testing) +func (m *MockProvider) SetKafkaConfig(cfg *sarama.Config) { + m.mu.Lock() + defer m.mu.Unlock() + m.kafkaConfig = cfg + m.notify() +} + +// SetKafkaFingerprint updates the mock Kafka fingerprint (for testing) +func (m *MockProvider) SetKafkaFingerprint(fingerprint string) { + m.mu.Lock() + defer m.mu.Unlock() + m.kafkaFingerprint = fingerprint + m.notify() +} diff --git a/internal/objectbucketsource/controller/objectbucketsource_controller.go b/internal/objectbucketsource/controller/objectbucketsource_controller.go index 1a1bde6..c02aa94 100644 --- a/internal/objectbucketsource/controller/objectbucketsource_controller.go +++ b/internal/objectbucketsource/controller/objectbucketsource_controller.go @@ -36,6 +36,7 @@ import ( logf "sigs.k8s.io/controller-runtime/pkg/log" sourcesv1alpha1 "github.com/functions-dev/func-operator/api/sources/v1alpha1" + "github.com/functions-dev/func-operator/internal/objectbucketsource/config" "github.com/functions-dev/func-operator/internal/objectbucketsource/s3client" ) @@ -56,13 +57,13 @@ var obcGVR = schema.GroupVersionResource{ type ObjectBucketSourceReconciler struct { client.Client Scheme *runtime.Scheme - AdapterConfigs []AdapterConfig + ConfigProvider config.ConfigProvider } // +kubebuilder:rbac:groups=sources.functions.dev,resources=objectbucketsources,verbs=get;list;watch;create;update;patch;delete // +kubebuilder:rbac:groups=sources.functions.dev,resources=objectbucketsources/status,verbs=get;update;patch // +kubebuilder:rbac:groups=sources.functions.dev,resources=objectbucketsources/finalizers,verbs=update -// +kubebuilder:rbac:groups="",resources=configmaps,verbs=get;list;watch +// +kubebuilder:rbac:groups="",resources=configmaps,verbs=get;list;watch;create // +kubebuilder:rbac:groups="",resources=secrets,verbs=get;list;watch // +kubebuilder:rbac:groups=objectbucket.io,resources=objectbucketclaims,verbs=get;list;watch @@ -271,18 +272,32 @@ func (r *ObjectBucketSourceReconciler) readOBCStorageClassName(ctx context.Conte } func (r *ObjectBucketSourceReconciler) resolveAdapterConfig(ctx context.Context, namespace, obcName string) (*AdapterConfig, error) { - if len(r.AdapterConfigs) == 0 { - return nil, fmt.Errorf("no adapter configs defined") + cfg := r.ConfigProvider.GetConfig() + + adapterConfigs := []AdapterConfig{ + { + ID: cfg.NoobaaAdapter.ID, + Topic: cfg.NoobaaAdapter.TopicARN, + StorageClassPattern: cfg.NoobaaAdapter.StorageClassPattern, + }, + { + ID: cfg.RadosgwAdapter.ID, + Topic: cfg.RadosgwAdapter.TopicARN, + StorageClassPattern: cfg.RadosgwAdapter.StorageClassPattern, + }, } - if len(r.AdapterConfigs) == 1 { - return &r.AdapterConfigs[0], nil + + if len(adapterConfigs) == 1 { + return &adapterConfigs[0], nil } + storageClass, err := r.readOBCStorageClassName(ctx, namespace, obcName) if err != nil { return nil, err } - for i := range r.AdapterConfigs { - cfg := &r.AdapterConfigs[i] + + for i := range adapterConfigs { + cfg := &adapterConfigs[i] if cfg.StorageClassPattern != nil && cfg.StorageClassPattern.MatchString(storageClass) { return cfg, nil } diff --git a/internal/objectbucketsource/controller/objectbucketsource_controller_test.go b/internal/objectbucketsource/controller/objectbucketsource_controller_test.go index d0bf536..49c5d72 100644 --- a/internal/objectbucketsource/controller/objectbucketsource_controller_test.go +++ b/internal/objectbucketsource/controller/objectbucketsource_controller_test.go @@ -28,6 +28,7 @@ import ( metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" sourcesv1alpha1 "github.com/functions-dev/func-operator/api/sources/v1alpha1" + "github.com/functions-dev/func-operator/internal/objectbucketsource/config" ) var _ = Describe("ObjectBucketSource Controller", func() { @@ -76,8 +77,9 @@ var _ = Describe("ObjectBucketSource Controller", func() { It("should successfully reconcile the resource", func() { By("Reconciling the created resource") controllerReconciler := &ObjectBucketSourceReconciler{ - Client: k8sClient, - Scheme: k8sClient.Scheme(), + Client: k8sClient, + Scheme: k8sClient.Scheme(), + ConfigProvider: config.NewMockProvider(), } _, err := controllerReconciler.Reconcile(ctx, reconcile.Request{ diff --git a/internal/objectbucketsource/notificationserver/server.go b/internal/objectbucketsource/notificationserver/server.go index 9ffbf51..5880ba7 100644 --- a/internal/objectbucketsource/notificationserver/server.go +++ b/internal/objectbucketsource/notificationserver/server.go @@ -11,36 +11,133 @@ import ( logf "sigs.k8s.io/controller-runtime/pkg/log" ceDispatch "github.com/functions-dev/func-operator/internal/objectbucketsource/cloudevents" + "github.com/functions-dev/func-operator/internal/objectbucketsource/config" ) var log = logf.Log.WithName("notification-server") +// ConfigProvider provides access to configuration +type ConfigProvider interface { + GetConfig() config.Config + GetKafkaConfig() *sarama.Config + // GetKafkaFingerprint returns a content hash of the Kafka secret; it changes + // whenever the Kafka credentials/connection settings are rotated. + GetKafkaFingerprint() string + // Subscribe returns a channel that is signaled whenever the configuration is reloaded. + Subscribe() <-chan struct{} +} + type NotificationServer struct { - Client client.Client - Port int - KafkaBrokers []string - KafkaConfig *sarama.Config - NotificationsMode string - KafkaNotificationsTopics []string - KafkaNotificationsGroupID string + Client client.Client + Port int + ConfigProvider ConfigProvider +} + +// runSnapshot captures everything that, when changed, requires the notification +// runner to be restarted. +type runSnapshot struct { + settings config.NotificationSettings + kafkaFingerprint string +} + +func (s *NotificationServer) snapshot() runSnapshot { + return runSnapshot{ + settings: s.ConfigProvider.GetConfig().Notifications, + kafkaFingerprint: s.ConfigProvider.GetKafkaFingerprint(), + } } +// Start runs the notification runner (HTTP server or Kafka consumer) according to +// the current configuration and supervises it: when the notification-related +// settings or the Kafka credentials change in the ConfigMap/Secret, it gracefully +// stops the current runner and starts a new one with the updated settings. func (s *NotificationServer) Start(ctx context.Context) error { + changes := s.ConfigProvider.Subscribe() + + for { + if ctx.Err() != nil { + return nil + } + + snap := s.snapshot() + if shutdown := s.superviseRun(ctx, changes, snap); shutdown { + return nil + } + } +} + +// superviseRun starts the notification runner for the given snapshot and blocks +// until either the parent context is cancelled (returns true, indicating +// shutdown) or a restart is required because the settings/credentials changed or +// the runner exited (returns false). It always stops the runner before returning. +func (s *NotificationServer) superviseRun(ctx context.Context, changes <-chan struct{}, snap runSnapshot) (shutdown bool) { + runCtx, cancel := context.WithCancel(ctx) + defer cancel() + + errCh := make(chan error, 1) + go func() { + errCh <- s.run(runCtx, snap.settings) + }() + + for { + select { + case <-ctx.Done(): + <-errCh + return true + case err := <-errCh: + // The runner exited on its own (fatal setup error). Restart it, + // unless we're shutting down. + if ctx.Err() != nil { + return true + } + if err != nil { + log.Error(err, "notification runner stopped unexpectedly, restarting in 5s") + select { + case <-ctx.Done(): + return true + case <-time.After(5 * time.Second): + } + } + return false + case <-changes: + newSnap := s.snapshot() + if !needsRestart(snap, newSnap) { + // Unrelated configuration change (e.g. adapter IDs); keep running. + continue + } + kafkaCredsChanged := snap.kafkaFingerprint != newSnap.kafkaFingerprint + log.Info("notification configuration changed, restarting notification runner", + "old-mode", snap.settings.Mode, "new-mode", newSnap.settings.Mode, + "old-brokers", snap.settings.KafkaBrokers, "new-brokers", newSnap.settings.KafkaBrokers, + "old-topics", snap.settings.KafkaNotificationsTopics, "new-topics", newSnap.settings.KafkaNotificationsTopics, + "old-group-id", snap.settings.KafkaNotificationsGroupID, "new-group-id", newSnap.settings.KafkaNotificationsGroupID, + "kafka-credentials-changed", kafkaCredsChanged) + cancel() + <-errCh + return false + } + } +} + +// run starts the notification transport for the given settings and blocks until +// ctx is cancelled or a fatal error occurs. +func (s *NotificationServer) run(ctx context.Context, settings config.NotificationSettings) error { var kafkaProducer sarama.SyncProducer - if len(s.KafkaBrokers) > 0 { + if len(settings.KafkaBrokers) > 0 { var err error - kafkaProducer, err = ceDispatch.NewKafkaProducer(s.KafkaBrokers, s.KafkaConfig) + kafkaCfg := s.ConfigProvider.GetKafkaConfig() + kafkaProducer, err = ceDispatch.NewKafkaProducer(settings.KafkaBrokers, kafkaCfg) if err != nil { return fmt.Errorf("creating kafka producer: %w", err) } defer func() { _ = kafkaProducer.Close() }() - log.Info("kafka producer initialized", "brokers", s.KafkaBrokers) + log.Info("kafka producer initialized", "brokers", settings.KafkaBrokers) } handler := ¬ificationHandler{client: s.Client, kafkaProducer: kafkaProducer} - if s.NotificationsMode == "kafka" { - return s.startKafkaConsumer(ctx, handler) + if settings.Mode == "kafka" { + return s.startKafkaConsumer(ctx, handler, settings) } return s.startHTTPServer(ctx, handler) } @@ -71,21 +168,22 @@ func (s *NotificationServer) startHTTPServer(ctx context.Context, handler *notif return nil } -func (s *NotificationServer) startKafkaConsumer(ctx context.Context, handler *notificationHandler) error { - consumerConfig := *s.KafkaConfig +func (s *NotificationServer) startKafkaConsumer(ctx context.Context, handler *notificationHandler, settings config.NotificationSettings) error { + kafkaCfg := s.ConfigProvider.GetKafkaConfig() + consumerConfig := *kafkaCfg consumerConfig.Consumer.Return.Errors = true consumerConfig.Consumer.Offsets.Initial = sarama.OffsetNewest - consumerGroup, err := sarama.NewConsumerGroup(s.KafkaBrokers, s.KafkaNotificationsGroupID, &consumerConfig) + consumerGroup, err := sarama.NewConsumerGroup(settings.KafkaBrokers, settings.KafkaNotificationsGroupID, &consumerConfig) if err != nil { return fmt.Errorf("creating kafka consumer group: %w", err) } defer func() { _ = consumerGroup.Close() }() log.Info("starting kafka notification consumer", - "topics", s.KafkaNotificationsTopics, - "group", s.KafkaNotificationsGroupID, - "brokers", s.KafkaBrokers) + "topics", settings.KafkaNotificationsTopics, + "group", settings.KafkaNotificationsGroupID, + "brokers", settings.KafkaBrokers) go func() { for err := range consumerGroup.Errors() { @@ -96,7 +194,7 @@ func (s *NotificationServer) startKafkaConsumer(ctx context.Context, handler *no cgHandler := &consumerGroupHandler{handler: handler} for { - if err := consumerGroup.Consume(ctx, s.KafkaNotificationsTopics, cgHandler); err != nil { + if err := consumerGroup.Consume(ctx, settings.KafkaNotificationsTopics, cgHandler); err != nil { if ctx.Err() != nil { return nil } @@ -111,3 +209,44 @@ func (s *NotificationServer) startKafkaConsumer(ctx context.Context, handler *no func (s *NotificationServer) NeedLeaderElection() bool { return false } + +// needsRestart reports whether the change from old to new requires restarting the +// notification runner. The runner is restarted when the notification transport +// settings change, or when the Kafka credentials change while Kafka is in use +// (Kafka mode, or HTTP mode with brokers configured for kafka: sinks). +func needsRestart(old, new runSnapshot) bool { + if !notificationSettingsEqual(old.settings, new.settings) { + return true + } + if usesKafka(new.settings) && old.kafkaFingerprint != new.kafkaFingerprint { + return true + } + return false +} + +// usesKafka reports whether the given settings establish any Kafka connection +// (a consumer in kafka mode, or a producer for kafka: sinks when brokers are set). +func usesKafka(s config.NotificationSettings) bool { + return s.Mode == "kafka" || len(s.KafkaBrokers) > 0 +} + +// notificationSettingsEqual reports whether the two notification settings are +// equivalent for the purposes of deciding whether the runner must be restarted. +func notificationSettingsEqual(a, b config.NotificationSettings) bool { + return a.Mode == b.Mode && + a.KafkaNotificationsGroupID == b.KafkaNotificationsGroupID && + stringSlicesEqual(a.KafkaBrokers, b.KafkaBrokers) && + stringSlicesEqual(a.KafkaNotificationsTopics, b.KafkaNotificationsTopics) +} + +func stringSlicesEqual(a, b []string) bool { + if len(a) != len(b) { + return false + } + for i := range a { + if a[i] != b[i] { + return false + } + } + return true +}