From 356e21686f87b74cb5259f25482dddca497235b7 Mon Sep 17 00:00:00 2001 From: panxiao81 Date: Sun, 27 Sep 2026 17:37:45 +0000 Subject: [PATCH] =?UTF-8?q?refactor:=20=E6=8C=89=E8=81=8C=E8=B4=A3?= =?UTF-8?q?=E6=8B=86=E5=88=86=20controller=20=E5=90=AF=E5=8A=A8=E6=B5=81?= =?UTF-8?q?=E7=A8=8B?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- Makefile | 6 +- docs/architecture/overview.md | 6 + internal/bootstrap/database.go | 39 ++++ internal/bootstrap/manager.go | 94 ++++++++ .../bootstrap/manager_integration_test.go | 75 +++++++ internal/bootstrap/options.go | 59 ++++++ internal/bootstrap/options_test.go | 134 ++++++++++++ internal/bootstrap/run.go | 200 ++---------------- 8 files changed, 425 insertions(+), 188 deletions(-) create mode 100644 internal/bootstrap/manager.go create mode 100644 internal/bootstrap/manager_integration_test.go create mode 100644 internal/bootstrap/options.go create mode 100644 internal/bootstrap/options_test.go diff --git a/Makefile b/Makefile index f3a8a19..3841a0a 100644 --- a/Makefile +++ b/Makefile @@ -69,11 +69,13 @@ lint: golangci-lint ## Run golangci-lint linter .PHONY: test-database-integration test-database-integration: setup-envtest ## 使用临时 API server、PostgreSQL 与 OpenBao 容器验证 Database 后端。 - KUBEBUILDER_ASSETS="$(shell "$(ENVTEST)" use $(ENVTEST_K8S_VERSION) --bin-dir "$(LOCALBIN)" -p path)" go test -tags=integration -race -count=1 ./internal/database/... ./internal/infra/... + KUBEBUILDER_ASSETS="$(shell "$(ENVTEST)" use $(ENVTEST_K8S_VERSION) --bin-dir "$(LOCALBIN)" -p path)" \ + go test -tags=integration -race -count=1 ./internal/database/... ./internal/infra/... ./internal/bootstrap/... .PHONY: lint-database-integration lint-database-integration: golangci-lint ## 检查集成测试构建标签下的 Database 代码。 - "$(GOLANGCI_LINT)" run --build-tags=integration ./internal/database/... ./internal/infra/... + "$(GOLANGCI_LINT)" run --build-tags=integration \ + ./internal/database/... ./internal/infra/... ./internal/bootstrap/... .PHONY: lint-fix lint-fix: golangci-lint ## Run golangci-lint linter and perform fixes diff --git a/docs/architecture/overview.md b/docs/architecture/overview.md index 8857dc3..666a65c 100644 --- a/docs/architecture/overview.md +++ b/docs/architecture/overview.md @@ -50,6 +50,12 @@ Proxmox 作为稀缺物理基础设施可以共享,通过 pool、tag、token 不在 `cmd` 平铺组件装配文件,也不为每个领域生成独立二进制。 Makefile 与 Dockerfile 均继续构建 `cmd/main.go`。 +`Run` 只编排解析配置、创建 manager、显式装配组件、启动与退出清理。 +`options.go` 组织配置和通用 flags;`manager.go` 处理 scheme、metrics、webhook、TLS 与探针; +`database.go`、`openbao.go` 各自维护组件参数及装配细节。新增组件不向 `Run` 堆叠参数和内部 +条件分支,也不为此引入插件注册框架。组件启动失败时释放已装配资源,正常退出则先停止 +manager worker,再释放连接。 + 基础设施能力属于整个 controller-manager,不因首个消费者是 Database 就归入该领域。 `internal/infra/openbao` 管理官方 SDK client 的 TLS 配置、Kubernetes 认证及 token 生命周期, 不依赖 Database 或其他产品领域。Bao client 默认禁用自动重试,写入结果不确定时由用例处理; diff --git a/internal/bootstrap/database.go b/internal/bootstrap/database.go index 9d8f6a1..fa0ce6c 100644 --- a/internal/bootstrap/database.go +++ b/internal/bootstrap/database.go @@ -1,6 +1,11 @@ package bootstrap import ( + "context" + "flag" + "fmt" + "os" + "git.ddupan.top/panxiao81/ayatori/internal/database/adapter/kubernetes" "git.ddupan.top/panxiao81/ayatori/internal/database/adapter/postgresql" "git.ddupan.top/panxiao81/ayatori/internal/database/application" @@ -8,6 +13,40 @@ import ( ctrl "sigs.k8s.io/controller-runtime" ) +type databaseOptions struct { + secretNamespace string + rootCert string +} + +func (o *databaseOptions) bindFlags(flags *flag.FlagSet) { + flags.StringVar(&o.secretNamespace, "database-secret-namespace", os.Getenv("POD_NAMESPACE"), + "固定管理 Secret namespace;为空时不启用 Instance 观测") + flags.StringVar(&o.rootCert, "database-root-cert", "", "PostgreSQL 管理连接信任的公开 CA bundle 路径") +} + +func (o databaseOptions) configureManager(options *ctrl.Options) { + if o.secretNamespace != "" { + options.Cache = databasecontroller.InstanceCacheOptions(o.secretNamespace) + } +} + +// setupDatabase 封装 Database 的内部装配,并返回在 manager 停止后执行的清理。 +func setupDatabase(ctx context.Context, manager ctrl.Manager, options databaseOptions) (func(), error) { + cleanup := func() {} + if options.secretNamespace != "" { + service, err := setupInstanceObservation(manager, options.secretNamespace, options.rootCert) + if err != nil { + return nil, fmt.Errorf("set up Instance observation: %w", err) + } + cleanup = service.Close + } + if err := (&databasecontroller.BindingReconciler{}).SetupWithManager(ctx, manager); err != nil { + cleanup() + return nil, fmt.Errorf("set up Database binding controller: %w", err) + } + return cleanup, nil +} + func setupInstanceObservation(manager ctrl.Manager, namespace, rootCert string) (*application.InstanceService, error) { credentials, err := kubernetes.NewSecretCredentials(manager.GetAPIReader(), namespace) if err != nil { diff --git a/internal/bootstrap/manager.go b/internal/bootstrap/manager.go new file mode 100644 index 0000000..fb601eb --- /dev/null +++ b/internal/bootstrap/manager.go @@ -0,0 +1,94 @@ +package bootstrap + +import ( + "crypto/tls" + "fmt" + + "k8s.io/apimachinery/pkg/runtime" + utilruntime "k8s.io/apimachinery/pkg/util/runtime" + clientgoscheme "k8s.io/client-go/kubernetes/scheme" + + // 保留 kubeconfig 支持的官方认证插件,不自建身份加载流程。 + _ "k8s.io/client-go/plugin/pkg/client/auth" + ctrl "sigs.k8s.io/controller-runtime" + "sigs.k8s.io/controller-runtime/pkg/healthz" + "sigs.k8s.io/controller-runtime/pkg/metrics/filters" + metricsserver "sigs.k8s.io/controller-runtime/pkg/metrics/server" + "sigs.k8s.io/controller-runtime/pkg/webhook" + + databasev1alpha1 "git.ddupan.top/panxiao81/ayatori/api/database/v1alpha1" + executionv1alpha1 "git.ddupan.top/panxiao81/ayatori/api/execution/v1alpha1" +) + +func newManager(options options) (ctrl.Manager, error) { + config, err := ctrl.GetConfig() + if err != nil { + return nil, fmt.Errorf("load Kubernetes configuration: %w", err) + } + configuration := options.manager.configuration() + options.database.configureManager(&configuration) + manager, err := ctrl.NewManager(config, configuration) + if err != nil { + return nil, fmt.Errorf("create controller manager: %w", err) + } + if err := manager.AddHealthzCheck("healthz", healthz.Ping); err != nil { + return nil, fmt.Errorf("set up health check: %w", err) + } + if err := manager.AddReadyzCheck("readyz", healthz.Ping); err != nil { + return nil, fmt.Errorf("set up readiness check: %w", err) + } + return manager, nil +} + +func (o managerOptions) configuration() ctrl.Options { + scheme := runtime.NewScheme() + utilruntime.Must(clientgoscheme.AddToScheme(scheme)) + utilruntime.Must(executionv1alpha1.AddToScheme(scheme)) + utilruntime.Must(databasev1alpha1.AddToScheme(scheme)) + return ctrl.Options{ + Scheme: scheme, + Metrics: o.metricsOptions(), + WebhookServer: webhook.NewServer(o.webhookOptions()), + HealthProbeBindAddress: o.probeAddr, + LeaderElection: o.enableLeaderElection, + LeaderElectionID: "a6325ed6.ddupan.top", + // 保持默认不主动释放选主 Lease:manager 停止后还要完成组件清理。 + } +} + +func (o managerOptions) tlsOptions() []func(*tls.Config) { + if o.enableHTTP2 { + return nil + } + // 默认禁用 HTTP/2,沿用 scaffold 对 Rapid Reset 等风险的防护。 + return []func(*tls.Config){func(config *tls.Config) { + config.NextProtos = []string{"http/1.1"} + }} +} + +func (o managerOptions) metricsOptions() metricsserver.Options { + options := metricsserver.Options{ + BindAddress: o.metricsAddr, + SecureServing: o.secureMetrics, + TLSOpts: o.tlsOptions(), + } + if o.secureMetrics { + options.FilterProvider = filters.WithAuthenticationAndAuthorization + } + if o.metricsCertPath != "" { + options.CertDir = o.metricsCertPath + options.CertName = o.metricsCertName + options.KeyName = o.metricsCertKey + } + return options +} + +func (o managerOptions) webhookOptions() webhook.Options { + options := webhook.Options{Port: o.webhookPort, TLSOpts: o.tlsOptions()} + if o.webhookCertPath != "" { + options.CertDir = o.webhookCertPath + options.CertName = o.webhookCertName + options.KeyName = o.webhookCertKey + } + return options +} diff --git a/internal/bootstrap/manager_integration_test.go b/internal/bootstrap/manager_integration_test.go new file mode 100644 index 0000000..e2c1ede --- /dev/null +++ b/internal/bootstrap/manager_integration_test.go @@ -0,0 +1,75 @@ +//go:build integration + +package bootstrap + +import ( + "context" + "os" + "path/filepath" + "testing" + "time" + + "sigs.k8s.io/controller-runtime/pkg/envtest" +) + +func TestBootstrapWithRealAPIServer(t *testing.T) { + environment := &envtest.Environment{ + CRDDirectoryPaths: []string{"../../config/crd/bases"}, + ErrorIfCRDPathMissing: true, + } + config, err := environment.Start() + if err != nil { + t.Fatal(err) + } + t.Cleanup(func() { + if err := environment.Stop(); err != nil { + t.Error(err) + } + }) + user, err := environment.AddUser(envtest.User{Name: "bootstrap-fixture", Groups: []string{"system:masters"}}, config) + if err != nil { + t.Fatal(err) + } + data, err := user.KubeConfig() + if err != nil { + t.Fatal("cannot generate isolated fixture kubeconfig") + } + // 仅写临时 envtest 身份,不读取现场 kubeconfig;t.TempDir 会自动清理。 + path := filepath.Join(t.TempDir(), "kubeconfig") + if err := os.WriteFile(path, data, 0600); err != nil { + t.Fatal(err) + } + t.Setenv("KUBECONFIG", path) + t.Setenv("POD_NAMESPACE", "") + options := parseTestOptions(t, "--metrics-bind-address=0", "--health-probe-bind-address=0", "--webhook-port=-1") + manager, err := newManager(options) + if err != nil { + t.Fatal(err) + } + if err := setupOpenBaoAuthentication(manager, options.openBao); err != nil { + t.Fatal(err) + } + ctx, cancel := context.WithTimeout(t.Context(), 30*time.Second) + defer cancel() + cleanup, err := setupDatabase(ctx, manager, options.database) + if err != nil { + t.Fatal(err) + } + defer cleanup() + done := make(chan error, 1) + go func() { done <- manager.Start(ctx) }() + defer func() { + cancel() + select { + case err := <-done: + if err != nil { + t.Error(err) + } + case <-time.After(20 * time.Second): + t.Error("manager did not stop before component cleanup") + } + }() + if !manager.GetCache().WaitForCacheSync(ctx) { + t.Fatal("assembled controller cache did not synchronize") + } +} diff --git a/internal/bootstrap/options.go b/internal/bootstrap/options.go new file mode 100644 index 0000000..0ad6d9f --- /dev/null +++ b/internal/bootstrap/options.go @@ -0,0 +1,59 @@ +package bootstrap + +import ( + "flag" + + "sigs.k8s.io/controller-runtime/pkg/log/zap" +) + +type options struct { + manager managerOptions + database databaseOptions + openBao openBaoOptions + logging zap.Options +} + +func (o *options) bindFlags(flags *flag.FlagSet) { + o.manager.bindFlags(flags) + o.database.bindFlags(flags) + o.openBao.bindFlags(flags) + o.logging.Development = true + o.logging.BindFlags(flags) +} + +type managerOptions struct { + metricsAddr string + metricsCertPath string + metricsCertName string + metricsCertKey string + webhookCertPath string + webhookCertName string + webhookCertKey string + webhookPort int + enableLeaderElection bool + probeAddr string + secureMetrics bool + enableHTTP2 bool +} + +func (o *managerOptions) bindFlags(flags *flag.FlagSet) { + flags.StringVar(&o.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.") + flags.StringVar(&o.probeAddr, "health-probe-bind-address", ":8081", "The address the probe endpoint binds to.") + flags.BoolVar(&o.enableLeaderElection, "leader-elect", false, + "Enable leader election for controller manager. "+ + "Enabling this will ensure there is only one active controller manager.") + flags.BoolVar(&o.secureMetrics, "metrics-secure", true, + "If set, the metrics endpoint is served securely via HTTPS. Use --metrics-secure=false to use HTTP instead.") + flags.StringVar(&o.webhookCertPath, "webhook-cert-path", "", "The directory that contains the webhook certificate.") + flags.StringVar(&o.webhookCertName, "webhook-cert-name", "tls.crt", "The name of the webhook certificate file.") + flags.StringVar(&o.webhookCertKey, "webhook-cert-key", "tls.key", "The name of the webhook key file.") + flags.IntVar(&o.webhookPort, "webhook-port", 9443, "Port the webhook server listens on. "+ + "Defaults to 9443. Set -1 to disable the webhook server.") + flags.StringVar(&o.metricsCertPath, "metrics-cert-path", "", + "The directory that contains the metrics server certificate.") + flags.StringVar(&o.metricsCertName, "metrics-cert-name", "tls.crt", "The name of the metrics server certificate file.") + flags.StringVar(&o.metricsCertKey, "metrics-cert-key", "tls.key", "The name of the metrics server key file.") + flags.BoolVar(&o.enableHTTP2, "enable-http2", false, + "If set, HTTP/2 will be enabled for the metrics and webhook servers") +} diff --git a/internal/bootstrap/options_test.go b/internal/bootstrap/options_test.go new file mode 100644 index 0000000..60ebf4b --- /dev/null +++ b/internal/bootstrap/options_test.go @@ -0,0 +1,134 @@ +package bootstrap + +import ( + "crypto/tls" + "flag" + "slices" + "testing" + + corev1 "k8s.io/api/core/v1" + "k8s.io/apimachinery/pkg/runtime" + + databasev1alpha1 "git.ddupan.top/panxiao81/ayatori/api/database/v1alpha1" + executionv1alpha1 "git.ddupan.top/panxiao81/ayatori/api/execution/v1alpha1" +) + +func parseTestOptions(t *testing.T, args ...string) options { + t.Helper() + flags := flag.NewFlagSet("bootstrap-test", flag.ContinueOnError) + var options options + options.bindFlags(flags) + if err := flags.Parse(args); err != nil { + t.Fatal(err) + } + return options +} + +func TestDefaultConfiguration(t *testing.T) { + t.Setenv("POD_NAMESPACE", "") + options := parseTestOptions(t) + manager := options.manager.configuration() + if manager.Metrics.BindAddress != "0" || !manager.Metrics.SecureServing || manager.Metrics.FilterProvider == nil { + t.Fatal("metrics defaults or authentication changed") + } + if manager.HealthProbeBindAddress != ":8081" || manager.LeaderElection || + manager.LeaderElectionID != "a6325ed6.ddupan.top" || manager.LeaderElectionReleaseOnCancel { + t.Fatal("probe or leader election defaults changed") + } + if options.manager.webhookOptions().Port != 9443 || !options.logging.Development { + t.Fatal("webhook or logging defaults changed") + } + if options.database.secretNamespace != "" || options.openBao.address != "" { + t.Fatal("optional backends enabled by default") + } + if options.openBao.mount != "kubernetes" || options.openBao.identity.Audience != "openbao" { + t.Fatal("OpenBao defaults changed") + } + options.database.configureManager(&manager) + if len(manager.Cache.ByObject) != 0 { + t.Fatal("disabled observation unexpectedly configured Secret cache") + } + for _, object := range []runtime.Object{ + &corev1.Secret{}, &corev1.ServiceAccount{}, &executionv1alpha1.Job{}, + &databasev1alpha1.PostgreSQLInstance{}, &databasev1alpha1.PostgreSQLDatabase{}, &databasev1alpha1.PostgreSQLTenant{}, + } { + if _, _, err := manager.Scheme.ObjectKinds(object); err != nil { + t.Fatalf("missing scheme registration for %T: %v", object, err) + } + } +} + +func TestManagerFlagOverrides(t *testing.T) { + options := parseTestOptions(t, + "--metrics-bind-address=:9090", "--metrics-secure=false", "--health-probe-bind-address=:9091", + "--leader-elect", "--webhook-port=-1", "--metrics-cert-path=/fixture/metrics", + "--metrics-cert-name=server.crt", "--metrics-cert-key=server.key", "--webhook-cert-path=/fixture/webhook", + "--webhook-cert-name=hook.crt", "--webhook-cert-key=hook.key", + ) + manager := options.manager.configuration() + if manager.Metrics.BindAddress != ":9090" || manager.Metrics.SecureServing || manager.Metrics.FilterProvider != nil || + manager.HealthProbeBindAddress != ":9091" || !manager.LeaderElection { + t.Fatal("manager flags not applied") + } + if manager.Metrics.CertDir != "/fixture/metrics" || manager.Metrics.CertName != "server.crt" || manager.Metrics.KeyName != "server.key" { + t.Fatal("metrics certificate flags not applied") + } + webhook := options.manager.webhookOptions() + if webhook.Port != -1 || webhook.CertDir != "/fixture/webhook" || webhook.CertName != "hook.crt" || webhook.KeyName != "hook.key" { + t.Fatal("webhook flags not applied") + } +} + +func TestHTTP2Policy(t *testing.T) { + for _, enabled := range []bool{false, true} { + args := []string{} + if enabled { + args = append(args, "--enable-http2") + } + options := parseTestOptions(t, args...) + for _, callbacks := range [][]func(*tls.Config){options.manager.metricsOptions().TLSOpts, options.manager.webhookOptions().TLSOpts} { + config := &tls.Config{NextProtos: []string{"h2", "http/1.1"}} + for _, callback := range callbacks { + callback(config) + } + if slices.Contains(config.NextProtos, "h2") != enabled || !slices.Contains(config.NextProtos, "http/1.1") { + t.Fatal("HTTP/2 policy changed") + } + } + } +} + +func TestComponentFlagsAndNamespaceScope(t *testing.T) { + t.Setenv("POD_NAMESPACE", "from-environment") + if parseTestOptions(t).database.secretNamespace != "from-environment" { + t.Fatal("namespace environment default lost") + } + options := parseTestOptions(t, + "--database-secret-namespace=controller", "--database-root-cert=/fixture/postgres-ca.pem", + "--openbao-address=https://bao.example", "--openbao-ca-cert=/fixture/bao-ca.pem", + "--openbao-auth-mount=cluster", "--openbao-auth-role=controller", + "--openbao-service-account-namespace=identity", "--openbao-service-account-name=bao-login", + "--openbao-token-audience=bao", + ) + if options.database.secretNamespace != "controller" || options.database.rootCert != "/fixture/postgres-ca.pem" { + t.Fatal("Database flags not applied") + } + bao := options.openBao + if bao.address != "https://bao.example" || bao.caCert != "/fixture/bao-ca.pem" || bao.mount != "cluster" || + bao.role != "controller" || bao.identity.Namespace != "identity" || bao.identity.ServiceAccount != "bao-login" || bao.identity.Audience != "bao" { + t.Fatal("OpenBao flags not applied") + } + manager := options.manager.configuration() + options.database.configureManager(&manager) + if len(manager.Cache.ByObject) != 1 { + t.Fatal("Secret cache scope missing") + } + for object, config := range manager.Cache.ByObject { + if _, ok := object.(*corev1.Secret); !ok { + t.Fatal("unexpected cache object") + } + if _, ok := config.Namespaces["controller"]; !ok || len(config.Namespaces) != 1 { + t.Fatal("Secret cache escaped explicit namespace") + } + } +} diff --git a/internal/bootstrap/run.go b/internal/bootstrap/run.go index bb732b7..c7435ef 100644 --- a/internal/bootstrap/run.go +++ b/internal/bootstrap/run.go @@ -3,210 +3,38 @@ package bootstrap import ( "context" - "crypto/tls" "flag" "fmt" - "os" - // 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" - - "k8s.io/apimachinery/pkg/runtime" - utilruntime "k8s.io/apimachinery/pkg/util/runtime" - clientgoscheme "k8s.io/client-go/kubernetes/scheme" ctrl "sigs.k8s.io/controller-runtime" - "sigs.k8s.io/controller-runtime/pkg/healthz" "sigs.k8s.io/controller-runtime/pkg/log/zap" - "sigs.k8s.io/controller-runtime/pkg/metrics/filters" - metricsserver "sigs.k8s.io/controller-runtime/pkg/metrics/server" - "sigs.k8s.io/controller-runtime/pkg/webhook" - - databasev1alpha1 "git.ddupan.top/panxiao81/ayatori/api/database/v1alpha1" - executionv1alpha1 "git.ddupan.top/panxiao81/ayatori/api/execution/v1alpha1" - "git.ddupan.top/panxiao81/ayatori/internal/database/application" - databasecontroller "git.ddupan.top/panxiao81/ayatori/internal/database/controller" - // +kubebuilder:scaffold:imports ) -var ( - scheme = runtime.NewScheme() - setupLog = ctrl.Log.WithName("setup") -) +var setupLog = ctrl.Log.WithName("setup") -func init() { - utilruntime.Must(clientgoscheme.AddToScheme(scheme)) - - utilruntime.Must(executionv1alpha1.AddToScheme(scheme)) - utilruntime.Must(databasev1alpha1.AddToScheme(scheme)) - // +kubebuilder:scaffold:scheme -} - -// Run 启动唯一的 controller-manager 进程;命令行参数与信号处理只初始化一次。 +// Run 只编排启动顺序;命令行参数与信号处理在进程中初始化一次。 func Run() error { - var openBao openBaoOptions - openBao.bindFlags(flag.CommandLine) - var databaseNamespace, databaseRootCert string - flag.StringVar(&databaseNamespace, "database-secret-namespace", os.Getenv("POD_NAMESPACE"), - "固定管理 Secret namespace;为空时不启用 Instance 观测") - flag.StringVar(&databaseRootCert, "database-root-cert", "", "PostgreSQL 管理连接信任的公开 CA bundle 路径") - var metricsAddr string - var metricsCertPath, metricsCertName, metricsCertKey string - var webhookCertPath, webhookCertName, webhookCertKey string - var webhookPort int - var enableLeaderElection bool - var probeAddr string - var secureMetrics bool - var enableHTTP2 bool - 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.") - flag.BoolVar(&enableLeaderElection, "leader-elect", false, - "Enable leader election for controller manager. "+ - "Enabling this will ensure there is only one active controller manager.") - flag.BoolVar(&secureMetrics, "metrics-secure", true, - "If set, the metrics endpoint is served securely via HTTPS. Use --metrics-secure=false to use HTTP instead.") - flag.StringVar(&webhookCertPath, "webhook-cert-path", "", "The directory that contains the webhook certificate.") - flag.StringVar(&webhookCertName, "webhook-cert-name", "tls.crt", "The name of the webhook certificate file.") - flag.StringVar(&webhookCertKey, "webhook-cert-key", "tls.key", "The name of the webhook key file.") - flag.IntVar(&webhookPort, "webhook-port", 9443, "Port the webhook server listens on. "+ - "Defaults to 9443. Set -1 to disable the webhook server.") - flag.StringVar(&metricsCertPath, "metrics-cert-path", "", - "The directory that contains the metrics server certificate.") - flag.StringVar(&metricsCertName, "metrics-cert-name", "tls.crt", "The name of the metrics server certificate file.") - 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") - opts := zap.Options{ - Development: true, - } - opts.BindFlags(flag.CommandLine) + var options options + options.bindFlags(flag.CommandLine) flag.Parse() + ctrl.SetLogger(zap.New(zap.UseFlagOptions(&options.logging))) - ctrl.SetLogger(zap.New(zap.UseFlagOptions(&opts))) - - // if the enable-http2 flag is false (the default), http/2 should be disabled - // due to its vulnerabilities. More specifically, disabling http/2 will - // prevent from being vulnerable to the HTTP/2 Stream Cancellation and - // Rapid Reset CVEs. For more information see: - // - https://github.com/advisories/GHSA-qppj-fm5r-hxr3 - // - https://github.com/advisories/GHSA-4374-p667-p6c8 - disableHTTP2 := func(c *tls.Config) { - setupLog.Info("Disabling HTTP/2") - c.NextProtos = []string{"http/1.1"} - } - - if !enableHTTP2 { - tlsOpts = append(tlsOpts, disableHTTP2) - } - - // Initial webhook TLS options - webhookTLSOpts := tlsOpts - webhookServerOptions := webhook.Options{ - TLSOpts: webhookTLSOpts, - Port: webhookPort, - } - - if len(webhookCertPath) > 0 { - setupLog.Info("Initializing webhook certificate watcher using provided certificates", - "webhook-cert-path", webhookCertPath, "webhook-cert-name", webhookCertName, "webhook-cert-key", webhookCertKey) - - webhookServerOptions.CertDir = webhookCertPath - webhookServerOptions.CertName = webhookCertName - webhookServerOptions.KeyName = webhookCertKey - } - - webhookServer := webhook.NewServer(webhookServerOptions) - - // Metrics endpoint is enabled in 'config/default/kustomization.yaml'. The Metrics options configure the server. - // More info: - // - https://pkg.go.dev/sigs.k8s.io/controller-runtime@v0.25.0/pkg/metrics/server - // - https://book.kubebuilder.io/reference/metrics.html - metricsServerOptions := metricsserver.Options{ - BindAddress: metricsAddr, - SecureServing: secureMetrics, - TLSOpts: tlsOpts, - } - - if secureMetrics { - // FilterProvider is used to protect the metrics endpoint with authn/authz. - // These configurations ensure that only authorized users and service accounts - // can access the metrics endpoint. The RBAC are configured in 'config/rbac/kustomization.yaml'. More info: - // https://pkg.go.dev/sigs.k8s.io/controller-runtime@v0.25.0/pkg/metrics/filters#WithAuthenticationAndAuthorization - metricsServerOptions.FilterProvider = filters.WithAuthenticationAndAuthorization - } - - // If the certificate is not specified, controller-runtime will automatically - // generate self-signed certificates for the metrics server. While convenient for development and testing, - // this setup is not recommended for production. - // - // TODO(user): If you enable certManager, uncomment the following lines: - // - [METRICS-WITH-CERTS] at config/default/kustomization.yaml to generate and use certificates - // managed by cert-manager for the metrics server. - // - [PROMETHEUS-WITH-CERTS] at config/prometheus/kustomization.yaml for TLS certification. - if len(metricsCertPath) > 0 { - setupLog.Info("Initializing metrics certificate watcher using provided certificates", - "metrics-cert-path", metricsCertPath, "metrics-cert-name", metricsCertName, "metrics-cert-key", metricsCertKey) - - metricsServerOptions.CertDir = metricsCertPath - metricsServerOptions.CertName = metricsCertName - metricsServerOptions.KeyName = metricsCertKey - } - - managerOptions := ctrl.Options{ - Scheme: scheme, - Metrics: metricsServerOptions, - WebhookServer: webhookServer, - HealthProbeBindAddress: probeAddr, - LeaderElection: enableLeaderElection, - LeaderElectionID: "a6325ed6.ddupan.top", - // LeaderElectionReleaseOnCancel defines if the leader should step down voluntarily - // when the Manager ends. This requires the binary to immediately end when the - // Manager is stopped, otherwise, this setting is unsafe. Setting this significantly - // speeds up voluntary leader transitions as the new leader don't have to wait - // LeaseDuration time first. - // - // In the default scaffold provided, the program ends immediately after - // the manager stops, so would be fine to enable this option. However, - // if you are doing or is intended to do any operation such as perform cleanups - // after the manager stops then its usage might be unsafe. - // LeaderElectionReleaseOnCancel: true, - } - if databaseNamespace != "" { - managerOptions.Cache = databasecontroller.InstanceCacheOptions(databaseNamespace) - } - mgr, err := ctrl.NewManager(ctrl.GetConfigOrDie(), managerOptions) + manager, err := newManager(options) if err != nil { - return fmt.Errorf("create controller manager: %w", err) + return err } - - // +kubebuilder:scaffold:builder - if err := setupOpenBaoAuthentication(mgr, openBao); err != nil { + if err := setupOpenBaoAuthentication(manager, options.openBao); err != nil { return fmt.Errorf("set up OpenBao authentication: %w", err) } - var instanceService *application.InstanceService - if databaseNamespace != "" { - instanceService, err = setupInstanceObservation(mgr, databaseNamespace, databaseRootCert) - if err != nil { - return fmt.Errorf("set up Instance observation: %w", err) - } - // Run 返回前 manager 的 worker 已停止;启动中途失败也释放已装配的连接。 - defer instanceService.Close() - } - if err := (&databasecontroller.BindingReconciler{}).SetupWithManager(context.Background(), mgr); err != nil { - return fmt.Errorf("set up Database binding controller: %w", err) - } - - if err := mgr.AddHealthzCheck("healthz", healthz.Ping); err != nil { - return fmt.Errorf("set up health check: %w", err) - } - if err := mgr.AddReadyzCheck("readyz", healthz.Ping); err != nil { - return fmt.Errorf("set up readiness check: %w", err) + cleanup, err := setupDatabase(context.Background(), manager, options.database) + if err != nil { + return err } + // manager 的 worker 完全停止后才释放组件持有的资源。 + defer cleanup() setupLog.Info("Starting manager") - if err := mgr.Start(ctrl.SetupSignalHandler()); err != nil { + if err := manager.Start(ctrl.SetupSignalHandler()); err != nil { return fmt.Errorf("run controller manager: %w", err) } return nil