Merge pull request 'feat: OpenBao Kubernetes 认证与短期会话续期' (#13) from feat/database-openbao-auth into main

This commit was merged in pull request #13.
This commit is contained in:
2026-09-27 17:48:22 +00:00
28 changed files with 1645 additions and 263 deletions
+5
View File
@@ -25,6 +25,11 @@ jobs:
make test make test
git diff --exit-code git diff --exit-code
- name: Build controller entrypoint
run: |
make build
./bin/manager --help
lint: lint:
runs-on: [self-hosted, pod] runs-on: [self-hosted, pod]
steps: steps:
+4 -2
View File
@@ -69,11 +69,13 @@ lint: golangci-lint ## Run golangci-lint linter
.PHONY: test-database-integration .PHONY: test-database-integration
test-database-integration: setup-envtest ## 使用临时 API server、PostgreSQL 与 OpenBao 容器验证 Database 后端。 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/... 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 .PHONY: lint-database-integration
lint-database-integration: golangci-lint ## 检查集成测试构建标签下的 Database 代码。 lint-database-integration: golangci-lint ## 检查集成测试构建标签下的 Database 代码。
"$(GOLANGCI_LINT)" run --build-tags=integration ./internal/database/... "$(GOLANGCI_LINT)" run --build-tags=integration \
./internal/database/... ./internal/infra/... ./internal/bootstrap/...
.PHONY: lint-fix .PHONY: lint-fix
lint-fix: golangci-lint ## Run golangci-lint linter and perform fixes lint-fix: golangci-lint ## Run golangci-lint linter and perform fixes
-26
View File
@@ -1,26 +0,0 @@
package main
import (
"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"
databasecontroller "git.ddupan.top/panxiao81/ayatori/internal/database/controller"
ctrl "sigs.k8s.io/controller-runtime"
)
func setupInstanceObservation(manager ctrl.Manager, namespace, rootCert string) (*application.InstanceService, error) {
credentials, err := kubernetes.NewSecretCredentials(manager.GetConfig(), namespace)
if err != nil {
return nil, err
}
service, err := application.NewInstanceService(credentials, postgresql.Connector{RootCert: rootCert})
if err != nil {
return nil, err
}
reconciler := &databasecontroller.InstanceReconciler{Observer: service, SecretNamespace: namespace}
if err := reconciler.SetupWithManager(manager); err != nil {
service.Close()
return nil, err
}
return service, nil
}
+3 -202
View File
@@ -1,214 +1,15 @@
package main package main
import ( import (
"context"
"crypto/tls"
"flag"
"os" "os"
// Import all Kubernetes client auth plugins (e.g. Azure, GCP, OIDC, etc.) "git.ddupan.top/panxiao81/ayatori/internal/bootstrap"
// 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" 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")
)
func init() {
utilruntime.Must(clientgoscheme.AddToScheme(scheme))
utilruntime.Must(executionv1alpha1.AddToScheme(scheme))
utilruntime.Must(databasev1alpha1.AddToScheme(scheme))
// +kubebuilder:scaffold:scheme
}
// nolint:gocyclo
func main() { func main() {
var databaseNamespace, databaseRootCert string if err := bootstrap.Run(); err != nil {
flag.StringVar(&databaseNamespace, "database-secret-namespace", os.Getenv("POD_NAMESPACE"), ctrl.Log.WithName("setup").Error(err, "Controller manager exited")
"固定管理 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)
flag.Parse()
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/[email protected]/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/[email protected]/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)
if err != nil {
setupLog.Error(err, "Failed to start manager")
os.Exit(1)
}
// +kubebuilder:scaffold:builder
var instanceService *application.InstanceService
if databaseNamespace != "" {
instanceService, err = setupInstanceObservation(mgr, databaseNamespace, databaseRootCert)
if err != nil {
setupLog.Error(err, "Failed to set up Instance observation")
os.Exit(1)
}
}
if err := (&databasecontroller.BindingReconciler{}).SetupWithManager(context.Background(), mgr); err != nil {
setupLog.Error(err, "Failed to set up Database binding controller")
os.Exit(1)
}
if err := mgr.AddHealthzCheck("healthz", healthz.Ping); err != nil {
setupLog.Error(err, "Failed to set up health check")
os.Exit(1)
}
if err := mgr.AddReadyzCheck("readyz", healthz.Ping); err != nil {
setupLog.Error(err, "Failed to set up ready check")
os.Exit(1)
}
setupLog.Info("Starting manager")
err = mgr.Start(ctrl.SetupSignalHandler())
// worker 完全停止后才释放 pgxpool,避免与在途观察竞争。
if instanceService != nil {
instanceService.Close()
}
if err != nil {
setupLog.Error(err, "Failed to run manager")
os.Exit(1) os.Exit(1)
} }
} }
@@ -0,0 +1,37 @@
# 示例:只授权 controller 为固定登录 SA 创建短期 JWT;由管理员替换 namespace/subject 后应用。
# 不自动纳入 config/default,不包含 kubeconfig、长期 token 或 OpenBao 管理权限。
apiVersion: v1
kind: ServiceAccount
metadata:
name: database-openbao-login
namespace: ayatori-system
---
apiVersion: rbac.authorization.k8s.io/v1
kind: Role
metadata:
name: database-openbao-token
namespace: ayatori-system
rules:
- apiGroups: [""]
resources: [serviceaccounts/token]
resourceNames: [database-openbao-login]
verbs: [create]
---
apiVersion: rbac.authorization.k8s.io/v1
kind: RoleBinding
metadata:
name: database-openbao-token
namespace: ayatori-system
roleRef:
apiGroup: rbac.authorization.k8s.io
kind: Role
name: database-openbao-token
subjects:
# 集群外示例:对应管理员签发 kubeconfig 的实际用户名,不是登录目标 SA 的名字。
- kind: User
apiGroup: rbac.authorization.k8s.io
name: ayatori-controller
# 集群内可以改为 manager 自己的 ServiceAccount;不要求两个 subject 同时授权。
# - kind: ServiceAccount
# name: controller-manager
# namespace: ayatori-system
+28
View File
@@ -43,6 +43,34 @@ Dev 与 Prod 使用独立的 Kubernetes API、数据库、身份和 controller
Proxmox 作为稀缺物理基础设施可以共享,通过 pool、tag、token 和明确的资源范围区分 Proxmox 作为稀缺物理基础设施可以共享,通过 pool、tag、token 和明确的资源范围区分
环境。其他后端尽量使用独立数据库、角色、地址池、DNS 空间与凭据。 环境。其他后端尽量使用独立数据库、角色、地址池、DNS 空间与凭据。
## 进程内依赖边界
`cmd/main.go` 是唯一程序入口,只调用 `internal/bootstrap.Run` 并处理退出状态。
命令行参数、manager 创建、infra 与领域装配集中在 `internal/bootstrap`,
不在 `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 默认禁用自动重试,写入结果不确定时由用例处理;
领域适配器不修改共享 client 的全局配置。Kubernetes 客户端、cache 和直连 reader 由 manager 管理;
启动入口负责装配与注入,不在领域适配器内重复创建客户端。
读写能力优先直接使用官方 `client.Reader`、`client.Client`、OpenBao KV API 等接口,
不为统一命名再包一层通用 reader/writer,也不引入全局注册中心。共享连接不表示扩大授权;
不同身份或权限边界仍由启动装配显式隔离。
领域按用例需要维护 repository 契约,其 adapter 负责 CR/领域对象映射及业务结果转换。
例如 Database 的七键凭据格式、UID 路径、禁止覆盖和不确定结果处理仍由 Database 维护;
它们不是公共 KV 存储的业务规则。Secret 管理凭据读取注入 `manager.GetAPIReader()`,
保持直连 API server、不缓存 Secret 内容的安全边界;资源写入复用 `manager.GetClient()`。
## 数据面 ## 数据面
Ayatori 不承载或重新实现数据面。控制面故障只应阻止创建与变更,不应停止已有 VM、 Ayatori 不承载或重新实现数据面。控制面故障只应阻止创建与变更,不应停止已有 VM、
+49 -2
View File
@@ -168,10 +168,57 @@ Instance 删除首先释放本地连接并撤销 Ready。任何引用它的 Data
真实后端覆盖创建/回读、并发唯一创建、重建适配器读取、软删除冲突、固定前缀 token 真实后端覆盖创建/回读、并发唯一创建、重建适配器读取、软删除冲突、固定前缀 token
拒绝管理路径,以及成功写入后丢失响应;HTTP 故障测试补充不重试和错误脱敏。 拒绝管理路径,以及成功写入后丢失响应;HTTP 故障测试补充不重试和错误脱敏。
这一切片尚未接入 manager:Kubernetes auth/token 生命周期、Database 状态中的稳定位置和 认证会话已按下面的显式参数接入 manager;凭据存储尚未接入供应用例。
已确认步骤、供应 service/controller、PostgreSQL 创建以及 ESO 交付仍未完成。 Database 状态中的稳定位置和已确认步骤、供应 service/controller、PostgreSQL 创建以及 ESO 交付仍未完成。
测试 token 只用于临时 fixture,不是生产静态 token 配置接口。现有绑定不会触发外部写入。 测试 token 只用于临时 fixture,不是生产静态 token 配置接口。现有绑定不会触发外部写入。
## OpenBao Kubernetes 认证会话
公共 `internal/infra/openbao.KubernetesSession` 复用官方 Kubernetes auth helper 和 `LifetimeWatcher`
(均为 v2.7.0)。认证直接注入 `manager.GetClient()`,与 reconcile 共用已装配的 Kubernetes
client,不从配置另建客户端。标准 `--kubeconfig` /
`KUBECONFIG` 支持 systemd 或其他集群外运行方式,集群内使用 in-cluster 配置,不要求存在 Pod。
Kubernetes 身份的签发和更新由部署管理及 client-go 的认证机制负责,不另建 kubeconfig 读取器。
每次登录前,通过该 client 的 `SubResource("token").Create` 调用固定 namespace/name 的
ServiceAccount TokenRequest;写入直连 API server,不读取 cache 或要求额外的 SA get 权限。申请
audience 匹配 OpenBao role、期望有效期 600 秒的短期 JWT;检查返回值非空且未过期,再交给
官方 Kubernetes auth helper。JWT 不缓存,不读取投射文件,也不回退静态 OpenBao token;
实际 JWT 有效期由 API server 决定。RBAC 拒绝或 TokenRequest 失败时不会继续 Bao 登录。
OpenBao 登录结果必须包含有效 token 和有限 TTL。
续期、到期阈值与等待时间由 SDK 管理;可续期 token 的续期失败就撤下本地 token,不可续期
token 由 SDK 监测剩余寿命。会话结束后最多每 5 秒重新登录一次,并重新申请 Kubernetes JWT。
`Ready()` 仅表示当前 lease 正受 SDK 管理,不授权任何 Database 写入,也不能保证下一次请求
必然成功。认证/续期响应不写日志、不返回给调用方;后端操作仍独立检查并返回脱敏错误。
`Start(ctx)` 退出时清空 client token 并等待续期 goroutine 结束。SDK Stop 不取消已经发出的
续期 HTTP 请求,因此专用 client 的请求期限固定为 15 秒;不增加新连接池或自己的续期算法。
同一会话拒绝并发 Start。manager 使用 Runnable 管理生命周期,并增加 `openbao-auth` readiness
检查;认证故障不影响 liveness。`--openbao-address` 为空时不启用,不自动修改生产 auth/RBAC。
HTTPS 和显式 CA/系统信任根不可通过 BAO 环境变量降级,参数见 [部署合同](deployment.md)。
依据官方 [Kubernetes auth](https://openbao.org/docs/auth/kubernetes/) 与
[token 生命周期](https://openbao.org/docs/concepts/auth/)。真实测试使用 envtest 签发 SA token,
OpenBao 通过专用 reviewer 调用真实 TokenReview;集群外受限 kubeconfig 启动实际 manager,
验证共享 client 申请 JWT、短 TTL 续期、RBAC 撤回/恢复与重新登录、跨 namespace/其他 SA 拒绝、
错误 OpenBao audience 拒绝以及凭据访问恢复。临时 TokenReview 入口仅允许对应 POST,
两段连接均验证 TLS;其 Docker bridge 入口仅为隔离测试,不修改生产 OpenBao 或 Kubernetes。
单元测试补充 TokenRequest 失败/空 token 无回退、重新申请 JWT、无期限 lease 拒绝、
并发生命周期、退出清理及 manager 显式参数不受 BAO 环境身份覆盖。
## 公共基础设施与领域适配
Bao client/TLS 与认证生命周期已移至 `internal/infra/openbao`,与任何产品领域无关。
Kubernetes 读写客户端由 manager 管理,Secret 凭据适配器只接收直连的 `client.Reader`。
`adapter/openbao.Credentials` 仍属于 Database:它直接使用官方 KV v2 API,实现七键凭据、
UID 路径、CAS=0 与回读确认的领域合同,不把这些规则推广为公共存储语义。
后续供应用例需要的 repository 接口由领域侧按实际操作定义,不提前增加通用仓储抽象。
分层约定见[总体架构](../architecture/overview.md#进程内依赖边界)。
认证单元测试归公共 infra;真实认证与 Database 凭据读写的组合测试仍在 Database adapter。
Database 集成测试和 lint 入口同时覆盖 `internal/infra/...`,避免拆包导致 CI 漏测。
## 设计入口 ## 设计入口
- [系统规格](specification.md):规范性行为与验收标准; - [系统规格](specification.md):规范性行为与验收标准;
+38 -6
View File
@@ -10,7 +10,7 @@
| 最后更新 | 2026-09-25 | | 最后更新 | 2026-09-25 |
本文定义 v1alpha1 的运行依赖、启动顺序和部署级配置。Instance 观测已接入 manager; 本文定义 v1alpha1 的运行依赖、启动顺序和部署级配置。Instance 观测已接入 manager;
OpenBao、ESO 与完整供应装配仍是后续实现合同。 OpenBao 认证可显式启用;ESO 与完整供应装配仍是后续实现合同。
## 依赖与顺序 ## 依赖与顺序
@@ -34,18 +34,21 @@ OpenBao、ESO 与完整供应装配仍是后续实现合同。
Instance 观测)与 `--database-root-cert`(公开 PostgreSQL CA PEM 路径)。Deployment Instance 观测)与 `--database-root-cert`(公开 PostgreSQL CA PEM 路径)。Deployment
通过 downward API 获取 namespace,Secret 权限由该 namespace 的 Role 授予。 通过 downward API 获取 namespace,Secret 权限由该 namespace 的 Role 授予。
以下是尚待实现的供应/交付配置合同,不表示当前 manager 接受这些 CLI flags。 以下表格区分已实现的认证参数与尚待实现的供应/交付参数。
必填项缺失、路径无效或 duration 不为正数时,进程必须在启动 manager 前失败; 必填项缺失、路径无效或 duration 不为正数时,进程必须在启动 manager 前失败;
不得等到 reconcile 时才逐个资源报告配置错误。 不得等到 reconcile 时才逐个资源报告配置错误。
| CLI flag | 必填/默认 | 说明 | | CLI flag | 必填/默认 | 说明 |
| --- | --- | --- | | --- | --- | --- |
| `--openbao-address` | 必填 | controller 可访问的 OpenBao API address | | `--openbao-address` | 已实现,默认空 | HTTPS API 地址;为空时关闭认证会话 |
| `--openbao-consumer-address` | 默认同 `--openbao-address` | 写入 Tenant status,必须能被预期外部消费者解析 | | `--openbao-consumer-address` | 默认同 `--openbao-address` | 写入 Tenant status,必须能被预期外部消费者解析 |
| `--openbao-auth-mount` | `kubernetes` | Kubernetes auth mount 名称 | | `--openbao-auth-mount` | 已实现,`kubernetes` | Kubernetes auth mount 名称 |
| `--openbao-auth-role` | 必填 | controller ServiceAccount 对应 role | | `--openbao-auth-role` | 已实现,启用时必填 | OpenBao 登录 role |
| `--openbao-ca-cert` | 已实现,默认系统信任根 | OpenBao 公开 CA PEM 路径 |
| `--openbao-kv-mount` | `kv` | KV v2 mount;开发可显式用 `secret` | | `--openbao-kv-mount` | `kv` | KV v2 mount;开发可显式用 `secret` |
| `--openbao-service-account-token-path` | `/var/run/secrets/kubernetes.io/serviceaccount/token` | Kubernetes auth 使用的投射 token 文件 | | `--openbao-service-account-namespace` | 已实现,启用时必填 | TokenRequest 目标 SA 的固定 namespace |
| `--openbao-service-account-name` | 已实现,启用时必填 | TokenRequest 目标 SA 名称 |
| `--openbao-token-audience` | 已实现,`openbao` | SA JWT audience,必须匹配 OpenBao role |
| `--openbao-tenant-base-path` | 默认 `postgresql-tenants` | controller 专属 mount-relative 前缀 | | `--openbao-tenant-base-path` | 默认 `postgresql-tenants` | controller 专属 mount-relative 前缀 |
| `--external-secret-store-name` | 必填 | controller 创建的 ExternalSecret 固定引用 | | `--external-secret-store-name` | 必填 | controller 创建的 ExternalSecret 固定引用 |
| `--database-root-cert` | 已实现 | 只读 PEM trust bundle,不含私钥;沿用 Instance 连接配置 | | `--database-root-cert` | 已实现 | 只读 PEM trust bundle,不含私钥;沿用 Instance 连接配置 |
@@ -56,6 +59,35 @@ address 必须是绝对 `http` 或 `https` URL,不允许 userinfo、query 或
`/` 开头,不含空段、`.` 或 `..`;base path 还不得编码 KV v2 的 `data`/`metadata` `/` 开头,不含空段、`.` 或 `..`;base path 还不得编码 KV v2 的 `data`/`metadata`
API 层。生产环境的 `--openbao-address` 必须使用 HTTPS;HTTP 只用于明确的开发 fixture。 API 层。生产环境的 `--openbao-address` 必须使用 HTTPS;HTTP 只用于明确的开发 fixture。
### 集群内与 systemd 共用 Kubernetes 认证
controller 是 API 客户端,不要求部署为 Pod。OpenBao 认证直接复用 manager 已加载的
Kubernetes 配置:集群外可用标准 `--kubeconfig`(或 `KUBECONFIG`),集群内可使用
in-cluster 配置。禁止另要求 `/var/run/secrets/.../token` 文件或解析 kubeconfig 中的 bearer token。
每次 Bao 登录前通过 TokenRequest 申请新的短期 SA JWT,Kubernetes JWT 与 Bao token 的
生命周期分别由 API 签发和官方 SDK 续期管理;不把 kubeconfig 本身当成永久有效凭据。
管理员为 controller 的实际 Kubernetes 身份授予目标 namespace 内
`create serviceaccounts/token`,用 `resourceNames` 限定目标 SA;示例见
[最小 RBAC](../../config/samples/database_openbao_auth_rbac.yaml)。集群外 RoleBinding subject
对应 kubeconfig 的用户/组,集群内可绑定 manager SA;登录目标 SA 可以独立于调用者身份。
controller 不创建 SA、Role/RoleBinding,也不向自己授予权限。不要求 `get secrets` 来获取 JWT。
OpenBao role 还需限制 SA 名称、namespace 与 audience,TokenReview reviewer 身份由管理员配置。
示例启动参数(仅示意,不包含真实 kubeconfig 或凭据):
```sh
manager --kubeconfig=/etc/ayatori/controller.kubeconfig \
--database-secret-namespace=ayatori-system \
--openbao-address=https://bao.example:8200 \
--openbao-auth-role=ayatori-database \
--openbao-service-account-namespace=ayatori-system \
--openbao-service-account-name=database-openbao-login
```
这里的认证成功只开放 manager readiness,不代表 Database 已具备供应或交付能力。
实例管理 Secret 的 namespace 同样由参数指定,systemd 模式不依赖 `POD_NAMESPACE` 环境变量。
Tenant 不能选择任意凭据路径。凭据必须能随 Database 保留并安全交付给被授权的新 Tenant; Tenant 不能选择任意凭据路径。凭据必须能随 Database 保留并安全交付给被授权的新 Tenant;
原 `<base-path>/<namespace>/<metadata.name>` 定位规则不再直接作为新 API 合同。 原 `<base-path>/<namespace>/<metadata.name>` 定位规则不再直接作为新 API 合同。
动态供应位置使用 `<base-path>/<Database UID>`;导入使用 Database 的显式 credentialRef, 动态供应位置使用 `<base-path>/<Database UID>`;导入使用 Database 的显式 credentialRef,
+3
View File
@@ -24,6 +24,9 @@ Kubernetes 管理员、OpenBao 管理员和 PostgreSQL 管理员是平台信任
## 凭据处理 ## 凭据处理
- controller 使用 Kubernetes auth 获取短期 OpenBao token,不配置长期静态 token。 - controller 使用 Kubernetes auth 获取短期 OpenBao token,不配置长期静态 token。
- Kubernetes auth 不等于部署在 Kubernetes 内:复用 manager 的 kubeconfig/in-cluster 身份,
通过最小 RBAC 的指定 ServiceAccount TokenRequest 获取 JWT,不依赖 Pod 投射文件。
kubeconfig 的签发、更新与撤销由部署管理负责;申请失败不回退其他机器身份。
- 管理凭据只从 Instance 引用的 controller namespace Secret 读取,不复制到 - 管理凭据只从 Instance 引用的 controller namespace Secret 读取,不复制到
CR/status/Event/metric/trace;管理员维护 ExternalSecret,由 ESO 同步该 Secret。 CR/status/Event/metric/trace;管理员维护 ExternalSecret,由 ESO 同步该 Secret。
- 动态供应密码使用密码学安全随机源;已有可靠关联时复用 OpenBao 现值,结果不确定时停止并报冲突。 - 动态供应密码使用密码学安全随机源;已有可靠关联时复用 OpenBao 现值,结果不确定时停止并报冲突。
+1
View File
@@ -4,6 +4,7 @@ go 1.27.1
require ( require (
github.com/jackc/pgx/v5 v5.11.0 github.com/jackc/pgx/v5 v5.11.0
github.com/openbao/openbao/api/auth/kubernetes/v2 v2.7.0
github.com/openbao/openbao/api/v2 v2.7.0 github.com/openbao/openbao/api/v2 v2.7.0
k8s.io/api v0.37.0 k8s.io/api v0.37.0
k8s.io/apimachinery v0.37.0 k8s.io/apimachinery v0.37.0
+2
View File
@@ -150,6 +150,8 @@ github.com/onsi/ginkgo/v2 v2.27.4 h1:fcEcQW/A++6aZAZQNUmNjvA9PSOzefMJBerHJ4t8v8Y
github.com/onsi/ginkgo/v2 v2.27.4/go.mod h1:ArE1D/XhNXBXCBkKOLkbsb2c81dQHCRcF5zwn/ykDRo= github.com/onsi/ginkgo/v2 v2.27.4/go.mod h1:ArE1D/XhNXBXCBkKOLkbsb2c81dQHCRcF5zwn/ykDRo=
github.com/onsi/gomega v1.39.0 h1:y2ROC3hKFmQZJNFeGAMeHZKkjBL65mIZcvrLQBF9k6Q= github.com/onsi/gomega v1.39.0 h1:y2ROC3hKFmQZJNFeGAMeHZKkjBL65mIZcvrLQBF9k6Q=
github.com/onsi/gomega v1.39.0/go.mod h1:ZCU1pkQcXDO5Sl9/VVEGlDyp+zm0m1cmeG5TOzLgdh4= github.com/onsi/gomega v1.39.0/go.mod h1:ZCU1pkQcXDO5Sl9/VVEGlDyp+zm0m1cmeG5TOzLgdh4=
github.com/openbao/openbao/api/auth/kubernetes/v2 v2.7.0 h1:Fw/pJRMpMTH83pMByCyikRHhxuBDYcnyiNSiK8OqJW0=
github.com/openbao/openbao/api/auth/kubernetes/v2 v2.7.0/go.mod h1:LkXPq4+8aLyQ+qoNBHcJF7nZFx0PYt2FOu+m7sdpAXU=
github.com/openbao/openbao/api/v2 v2.7.0 h1:3CD1l3tr39nQraCgFGAWA5vYvPFzZoZrt3NL7DMQKAc= github.com/openbao/openbao/api/v2 v2.7.0 h1:3CD1l3tr39nQraCgFGAWA5vYvPFzZoZrt3NL7DMQKAc=
github.com/openbao/openbao/api/v2 v2.7.0/go.mod h1:uXbMoyH2pjSvNyTepinUvLde8pOJB82EuhUCfOKnKbo= github.com/openbao/openbao/api/v2 v2.7.0/go.mod h1:uXbMoyH2pjSvNyTepinUvLde8pOJB82EuhUCfOKnKbo=
github.com/pkg/errors v0.9.1 h1:FEBLx1zS214owpjy7qsBeixbURkuhQAwrK5UwLGTwt4= github.com/pkg/errors v0.9.1 h1:FEBLx1zS214owpjy7qsBeixbURkuhQAwrK5UwLGTwt4=
+65
View File
@@ -0,0 +1,65 @@
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"
databasecontroller "git.ddupan.top/panxiao81/ayatori/internal/database/controller"
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 {
return nil, err
}
service, err := application.NewInstanceService(credentials, postgresql.Connector{RootCert: rootCert})
if err != nil {
return nil, err
}
reconciler := &databasecontroller.InstanceReconciler{Observer: service, SecretNamespace: namespace}
if err := reconciler.SetupWithManager(manager); err != nil {
service.Close()
return nil, err
}
return service, nil
}
+94
View File
@@ -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
}
@@ -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")
}
}
+72
View File
@@ -0,0 +1,72 @@
/*
Copyright 2026.
Licensed under the Apache License, Version 2.0 (the "License");
you may not use this file except in compliance with the License.
You may obtain a copy of the License at
http://www.apache.org/licenses/LICENSE-2.0
Unless required by applicable law or agreed to in writing, software
distributed under the License is distributed on an "AS IS" BASIS,
WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
See the License for the specific language governing permissions and
limitations under the License.
*/
package bootstrap
import (
"errors"
"flag"
"net/http"
ctrl "sigs.k8s.io/controller-runtime"
"git.ddupan.top/panxiao81/ayatori/internal/infra/openbao"
)
type openBaoOptions struct {
address string
caCert string
mount string
role string
identity openbao.KubernetesIdentity
}
func (o *openBaoOptions) bindFlags(flags *flag.FlagSet) {
flags.StringVar(&o.address, "openbao-address", "", "OpenBao HTTPS 地址;为空时不启用认证会话")
flags.StringVar(&o.caCert, "openbao-ca-cert", "", "OpenBao 公开 CA PEM 路径;默认使用系统信任根")
flags.StringVar(&o.mount, "openbao-auth-mount", "kubernetes", "OpenBao Kubernetes auth mount")
flags.StringVar(&o.role, "openbao-auth-role", "", "OpenBao 登录 role")
flags.StringVar(&o.identity.Namespace, "openbao-service-account-namespace", "",
"TokenRequest 的固定 ServiceAccount namespace")
flags.StringVar(&o.identity.ServiceAccount, "openbao-service-account-name", "", "TokenRequest 的固定 ServiceAccount 名称")
flags.StringVar(&o.identity.Audience, "openbao-token-audience", "openbao", "SA JWT audience,须匹配 OpenBao role")
}
func setupOpenBaoAuthentication(manager ctrl.Manager, options openBaoOptions) error {
if options.address == "" {
return nil
}
client, err := openbao.NewClient(options.address, options.caCert)
if err != nil {
return err
}
// 复用 manager 已装配的 Kubernetes client,不重复加载配置或创建客户端。
session, err := openbao.NewKubernetesSession(
client, manager.GetClient(), options.mount, options.role, options.identity,
)
if err != nil {
return err
}
if err := manager.Add(session); err != nil {
return err
}
return manager.AddReadyzCheck("openbao-auth", func(_ *http.Request) error {
if !session.Ready() {
return errors.New("OpenBao Kubernetes authentication unavailable")
}
return nil
})
}
+59
View File
@@ -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")
}
+134
View File
@@ -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")
}
}
}
+41
View File
@@ -0,0 +1,41 @@
// Package bootstrap 负责 controller-manager 的参数解析、依赖装配与启动。
package bootstrap
import (
"context"
"flag"
"fmt"
ctrl "sigs.k8s.io/controller-runtime"
"sigs.k8s.io/controller-runtime/pkg/log/zap"
)
var setupLog = ctrl.Log.WithName("setup")
// Run 只编排启动顺序;命令行参数与信号处理在进程中初始化一次。
func Run() error {
var options options
options.bindFlags(flag.CommandLine)
flag.Parse()
ctrl.SetLogger(zap.New(zap.UseFlagOptions(&options.logging)))
manager, err := newManager(options)
if err != nil {
return err
}
if err := setupOpenBaoAuthentication(manager, options.openBao); err != nil {
return fmt.Errorf("set up OpenBao authentication: %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 := manager.Start(ctrl.SetupSignalHandler()); err != nil {
return fmt.Errorf("run controller manager: %w", err)
}
return nil
}
@@ -22,10 +22,8 @@ import (
"errors" "errors"
corev1 "k8s.io/api/core/v1" corev1 "k8s.io/api/core/v1"
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
"k8s.io/apimachinery/pkg/util/validation" "k8s.io/apimachinery/pkg/util/validation"
typedcore "k8s.io/client-go/kubernetes/typed/core/v1" "sigs.k8s.io/controller-runtime/pkg/client"
"k8s.io/client-go/rest"
"git.ddupan.top/panxiao81/ayatori/internal/database/application" "git.ddupan.top/panxiao81/ayatori/internal/database/application"
"git.ddupan.top/panxiao81/ayatori/internal/database/domain/instance" "git.ddupan.top/panxiao81/ayatori/internal/database/domain/instance"
@@ -34,18 +32,16 @@ import (
// SecretCredentials 直接读取 API server,不将 Secret 数据纳入共享 informer cache。 // SecretCredentials 直接读取 API server,不将 Secret 数据纳入共享 informer cache。
// namespace 在装配时固定,Instance 不能选择跨 namespace 读取。 // namespace 在装配时固定,Instance 不能选择跨 namespace 读取。
type SecretCredentials struct { type SecretCredentials struct {
secrets typedcore.SecretInterface reader client.Reader
namespace string
} }
func NewSecretCredentials(config *rest.Config, namespace string) (*SecretCredentials, error) { // NewSecretCredentials 要求注入 manager.GetAPIReader() 或等价直连 reader,不可使用缓存 reader。
if config == nil || len(validation.IsDNS1123Label(namespace)) != 0 { func NewSecretCredentials(reader client.Reader, namespace string) (*SecretCredentials, error) {
return nil, errors.New("valid controller namespace and API configuration required") if reader == nil || len(validation.IsDNS1123Label(namespace)) != 0 {
return nil, errors.New("valid controller namespace and API reader required")
} }
client, err := typedcore.NewForConfig(config) return &SecretCredentials{reader: reader, namespace: namespace}, nil
if err != nil {
return nil, application.ErrCredentialsUnavailable
}
return &SecretCredentials{secrets: client.Secrets(namespace)}, nil
} }
func (r *SecretCredentials) Read(ctx context.Context, ref instance.CredentialReference) (application.Credentials, error) { func (r *SecretCredentials) Read(ctx context.Context, ref instance.CredentialReference) (application.Credentials, error) {
@@ -53,7 +49,8 @@ func (r *SecretCredentials) Read(ctx context.Context, ref instance.CredentialRef
return application.Credentials{}, application.ErrCredentialsInvalid return application.Credentials{}, application.ErrCredentialsInvalid
} }
keys := ref.Values() keys := ref.Values()
secret, err := r.secrets.Get(ctx, keys.Name, metav1.GetOptions{}) secret := &corev1.Secret{}
err := r.reader.Get(ctx, client.ObjectKey{Namespace: r.namespace, Name: keys.Name}, secret)
if err != nil { if err != nil {
return application.Credentials{}, application.ErrCredentialsUnavailable return application.Credentials{}, application.ErrCredentialsUnavailable
} }
@@ -0,0 +1,361 @@
//go:build integration
/*
Copyright 2026.
Licensed under the Apache License, Version 2.0 (the "License");
you may not use this file except in compliance with the License.
You may obtain a copy of the License at
http://www.apache.org/licenses/LICENSE-2.0
Unless required by applicable law or agreed to in writing, software
distributed under the License is distributed on an "AS IS" BASIS,
WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
See the License for the specific language governing permissions and
limitations under the License.
*/
package openbao_test
import (
"context"
"crypto/ecdsa"
"crypto/elliptic"
"crypto/rand"
"crypto/tls"
"crypto/x509"
"encoding/pem"
"math/big"
"net"
"net/http"
"net/http/httptest"
"net/http/httputil"
"net/url"
"os/exec"
"strconv"
"strings"
"testing"
"time"
bao "github.com/openbao/openbao/api/v2"
authenticationv1 "k8s.io/api/authentication/v1"
corev1 "k8s.io/api/core/v1"
rbacv1 "k8s.io/api/rbac/v1"
apierrors "k8s.io/apimachinery/pkg/api/errors"
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
"k8s.io/client-go/kubernetes"
"k8s.io/client-go/rest"
"k8s.io/client-go/tools/clientcmd"
ctrl "sigs.k8s.io/controller-runtime"
"sigs.k8s.io/controller-runtime/pkg/envtest"
metricsserver "sigs.k8s.io/controller-runtime/pkg/metrics/server"
"git.ddupan.top/panxiao81/ayatori/internal/infra/openbao"
)
const (
unrelatedAuthNamespace = "bao-unrelated"
tokenRequestRole = "request-openbao-token"
authNamespace = "bao-controller"
authAudience = "openbao"
)
// Bao 容器通过 Docker bridge 访问这个仅转发 TokenReview 的临时入口。
// 上游仍是带 CA 验证的真实 envtest API;不模拟 JWT 签名、audience 或 RBAC 判定。
func tokenReviewEndpoint(t *testing.T, config *rest.Config) (string, string) {
t.Helper()
upstream, err := url.Parse(config.Host)
if err != nil {
t.Fatal("invalid envtest address")
}
transport, err := rest.TransportFor(rest.AnonymousClientConfig(config))
if err != nil {
t.Fatal("cannot construct TokenReview transport")
}
proxy := httputil.NewSingleHostReverseProxy(upstream)
proxy.Transport = transport
proxy.ErrorHandler = func(w http.ResponseWriter, _ *http.Request, _ error) { w.WriteHeader(http.StatusBadGateway) }
server := httptest.NewUnstartedServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
if r.Method != http.MethodPost || r.URL.Path != "/apis/authentication.k8s.io/v1/tokenreviews" {
w.WriteHeader(http.StatusForbidden)
return
}
proxy.ServeHTTP(w, r)
}))
if err := server.Listener.Close(); err != nil {
t.Fatal("cannot replace fixture listener")
}
server.Listener, err = net.Listen("tcp", "0.0.0.0:0")
if err != nil {
t.Fatal("cannot expose TokenReview fixture")
}
ctx, cancel := context.WithTimeout(t.Context(), 10*time.Second)
defer cancel()
output, err := exec.CommandContext(ctx, "docker", "network", "inspect", "bridge", "--format",
`{{(index .IPAM.Config 0).Gateway}}`).Output()
if err != nil {
t.Fatal("cannot locate fixture Docker bridge")
}
gateway := strings.TrimSpace(string(output))
if net.ParseIP(gateway) == nil {
t.Fatal("invalid fixture bridge gateway")
}
certificate, caPEM := tokenReviewCertificate(t, net.ParseIP(gateway))
server.TLS = &tls.Config{Certificates: []tls.Certificate{certificate}, MinVersion: tls.VersionTLS12}
server.StartTLS()
t.Cleanup(server.Close)
port := server.Listener.Addr().(*net.TCPAddr).Port
return "https://" + net.JoinHostPort(gateway, strconv.Itoa(port)), caPEM
}
func tokenReviewCertificate(t *testing.T, address net.IP) (tls.Certificate, string) {
t.Helper()
key, err := ecdsa.GenerateKey(elliptic.P256(), rand.Reader)
if err != nil {
t.Fatal("cannot create fixture TLS key")
}
template := &x509.Certificate{
SerialNumber: big.NewInt(1),
NotBefore: time.Now().Add(-time.Minute),
NotAfter: time.Now().Add(time.Hour),
IPAddresses: []net.IP{address},
KeyUsage: x509.KeyUsageDigitalSignature | x509.KeyUsageCertSign,
ExtKeyUsage: []x509.ExtKeyUsage{x509.ExtKeyUsageServerAuth},
IsCA: true, BasicConstraintsValid: true,
}
der, err := x509.CreateCertificate(rand.Reader, template, template, &key.PublicKey, key)
if err != nil {
t.Fatal("cannot create fixture TLS certificate")
}
certificate := tls.Certificate{Certificate: [][]byte{der}, PrivateKey: key}
return certificate, string(pem.EncodeToMemory(&pem.Block{Type: "CERTIFICATE", Bytes: der}))
}
type kubernetesAuthFixture struct {
config *rest.Config
controllerConfig *rest.Config
admin *kubernetes.Clientset
controller *kubernetes.Clientset
grant func()
requestToken func(string, []string) string
}
func newKubernetesAuthFixture(t *testing.T) *kubernetesAuthFixture {
t.Helper()
environment := &envtest.Environment{}
config, err := environment.Start()
if err != nil {
t.Fatal("cannot start authentication API fixture", err)
}
t.Cleanup(func() {
if err := environment.Stop(); err != nil {
t.Error("cannot stop authentication API fixture")
}
})
clientset, err := kubernetes.NewForConfig(config)
if err != nil {
t.Fatal("cannot construct fixture API client")
}
ctx := t.Context()
for _, namespace := range []string{authNamespace, unrelatedAuthNamespace} {
if _, err := clientset.CoreV1().Namespaces().Create(ctx, &corev1.Namespace{Name: namespace}, metav1.CreateOptions{}); err != nil {
t.Fatal("cannot create fixture namespace")
}
if _, err := clientset.CoreV1().ServiceAccounts(namespace).Create(ctx, &corev1.ServiceAccount{Name: authRole}, metav1.CreateOptions{}); err != nil {
t.Fatal("cannot create fixture service account")
}
}
// 集群外身份只有指定 SA 的 TokenRequest 权限,不使用 envtest 管理员凭据运行会话。
user, err := environment.AddUser(envtest.User{Name: "systemd-controller"}, config)
if err != nil {
t.Fatal("cannot create external controller identity")
}
kubeconfig, err := user.KubeConfig()
if err != nil {
t.Fatal("cannot build external kubeconfig")
}
externalConfig, err := clientcmd.RESTConfigFromKubeConfig(kubeconfig)
if err != nil {
t.Fatal("cannot load external kubeconfig")
}
if _, err := clientset.RbacV1().Roles(authNamespace).Create(ctx, &rbacv1.Role{
Name: tokenRequestRole,
Rules: []rbacv1.PolicyRule{{APIGroups: []string{""}, Resources: []string{"serviceaccounts/token"},
ResourceNames: []string{authRole}, Verbs: []string{"create"}}},
}, metav1.CreateOptions{}); err != nil {
t.Fatal("cannot create TokenRequest Role")
}
permission := &rbacv1.RoleBinding{
Name: tokenRequestRole,
RoleRef: rbacv1.RoleRef{APIGroup: rbacv1.GroupName, Kind: "Role", Name: tokenRequestRole},
Subjects: []rbacv1.Subject{{Kind: "User", APIGroup: rbacv1.GroupName, Name: "systemd-controller"}},
}
grant := func() {
t.Helper()
if _, err := clientset.RbacV1().RoleBindings(authNamespace).Create(ctx, permission.DeepCopy(), metav1.CreateOptions{}); err != nil {
t.Fatal("cannot grant TokenRequest permission")
}
}
grant()
externalClient, err := kubernetes.NewForConfig(externalConfig)
if err != nil {
t.Fatal("cannot construct restricted controller client")
}
for _, target := range []struct{ namespace, name string }{{unrelatedAuthNamespace, authRole}, {authNamespace, "another-account"}} {
_, err := externalClient.CoreV1().ServiceAccounts(target.namespace).CreateToken(ctx, target.name,
&authenticationv1.TokenRequest{Spec: authenticationv1.TokenRequestSpec{Audiences: []string{authAudience}}}, metav1.CreateOptions{})
if !apierrors.IsForbidden(err) {
t.Fatal("TokenRequest escaped Role namespace/resourceNames restriction")
}
}
if _, err := clientset.RbacV1().ClusterRoleBindings().Create(ctx, &rbacv1.ClusterRoleBinding{
Name: "bao-fixture-reviewer",
RoleRef: rbacv1.RoleRef{APIGroup: rbacv1.GroupName, Kind: "ClusterRole", Name: "system:auth-delegator"},
Subjects: []rbacv1.Subject{{Kind: "ServiceAccount", Namespace: authNamespace, Name: authRole}},
}, metav1.CreateOptions{}); err != nil {
t.Fatal("cannot authorize fixture TokenReview")
}
requestToken := func(namespace string, audiences []string) string {
t.Helper()
response, err := clientset.CoreV1().ServiceAccounts(namespace).CreateToken(ctx, authRole,
&authenticationv1.TokenRequest{Spec: authenticationv1.TokenRequestSpec{Audiences: audiences}}, metav1.CreateOptions{})
if err != nil {
t.Fatal("cannot issue fixture service account token")
}
return response.Status.Token
}
return &kubernetesAuthFixture{
config: config, controllerConfig: externalConfig, admin: clientset, controller: externalClient,
grant: grant, requestToken: requestToken,
}
}
const authRole = "controller"
var testIdentity = openbao.KubernetesIdentity{Namespace: authNamespace, ServiceAccount: authRole, Audience: authAudience}
func waitForAuthentication(t *testing.T, check func() bool) {
t.Helper()
deadline := time.NewTimer(20 * time.Second)
defer deadline.Stop()
for !check() {
select {
case <-deadline.C:
t.Fatal("authentication condition timed out")
case <-time.After(20 * time.Millisecond):
}
}
}
func startSession(t *testing.T, session interface{ Start(context.Context) error }) context.CancelFunc {
t.Helper()
ctx, cancel := context.WithCancel(t.Context())
done := make(chan error, 1)
go func() { done <- session.Start(ctx) }()
t.Cleanup(func() {
cancel()
select {
case <-done:
case <-time.After(20 * time.Second):
t.Error("authentication did not stop")
}
})
return cancel
}
// 此处验证公共认证会话与 Database 凭据适配器的跨层集成。
func TestKubernetesSessionWithRealTokenReview(t *testing.T) {
api := newKubernetesAuthFixture(t)
ctx := t.Context()
root := baoFixture(t)
if err := root.Sys().EnableAuthWithOptionsWithContext(ctx, "kubernetes", &bao.EnableAuthOptions{Type: "kubernetes"}); err != nil {
t.Fatal("cannot enable fixture Kubernetes auth")
}
reviewerToken := api.requestToken(authNamespace, nil)
reviewURL, reviewCA := tokenReviewEndpoint(t, api.config)
if _, err := root.Logical().WriteWithContext(ctx, "auth/kubernetes/config", map[string]any{
"kubernetes_host": reviewURL,
"kubernetes_ca_cert": reviewCA,
"token_reviewer_jwt": reviewerToken,
"disable_local_ca_jwt": true,
}); err != nil {
t.Fatal("cannot configure fixture TokenReview:", strings.NewReplacer(reviewerToken, "[REDACTED]", fixtureToken, "[REDACTED]").Replace(err.Error()))
}
if err := root.Sys().PutPolicyWithContext(ctx, authRole, `path "secret/data/applications/*" { capabilities = ["create", "update", "read"] }`); err != nil {
t.Fatal("cannot configure fixture credential policy")
}
if _, err := root.Logical().WriteWithContext(ctx, "auth/kubernetes/role/controller", map[string]any{
"bound_service_account_names": []string{authRole},
"bound_service_account_namespaces": []string{authNamespace},
"audience": authAudience,
"token_policies": []string{authRole},
"token_ttl": "3s", "token_max_ttl": "60s",
}); err != nil {
t.Fatal("cannot configure fixture auth role")
}
client := fixtureClient(t, root.Address())
manager, err := ctrl.NewManager(api.controllerConfig, ctrl.Options{Metrics: metricsserver.Options{BindAddress: "0"}, HealthProbeBindAddress: "0"})
if err != nil {
t.Fatal("cannot construct external controller manager")
}
session, err := openbao.NewKubernetesSession(client, manager.GetClient(), "kubernetes", authRole, testIdentity)
if err != nil {
t.Fatal(err)
}
if err := manager.Add(session); err != nil {
t.Fatal("cannot register authentication lifecycle")
}
cancel := startSession(t, manager)
waitForAuthentication(t, session.Ready)
store := fixtureStore(t, client)
if err := store.Create(ctx, credentialPath, fixtureCredential(t)); err != nil {
t.Fatal("Kubernetes identity cannot create scoped credential", err)
}
initialToken := client.Token()
// 超过初始 TTL 后同一 token 仍可用,证明发生真实 renew-self,而非只登录一次。
start := time.Now()
waitForAuthentication(t, func() bool { return time.Since(start) > 4*time.Second })
if client.Token() != initialToken {
t.Fatal("token was replaced before renewal could be verified")
}
if _, err := client.Auth().Token().LookupSelfWithContext(ctx); err != nil {
t.Fatal("short-lived token was not renewed")
}
for _, invalid := range []struct {
namespace string
audience string
}{
{unrelatedAuthNamespace, authAudience}, {authNamespace, "wrong-audience"},
} {
// 用相同真实登录入口直接确认 namespace/audience 拒绝,不依赖定时轮询推断。
if _, err := root.Logical().WriteWithContext(ctx, "auth/kubernetes/login", map[string]any{
"role": authRole, "jwt": api.requestToken(invalid.namespace, []string{invalid.audience}),
}); err == nil {
t.Fatal("invalid Kubernetes identity was accepted")
}
}
if err := api.admin.RbacV1().RoleBindings(authNamespace).Delete(ctx, tokenRequestRole, metav1.DeleteOptions{}); err != nil {
t.Fatal("cannot revoke TokenRequest permission")
}
if err := root.Auth().Token().RevokeOrphanWithContext(ctx, client.Token()); err != nil {
t.Fatal("cannot revoke fixture OpenBao token")
}
waitForAuthentication(t, func() bool { return !session.Ready() && client.Token() == "" })
denied, err := api.controller.CoreV1().ServiceAccounts(authNamespace).CreateToken(ctx, authRole,
&authenticationv1.TokenRequest{Spec: authenticationv1.TokenRequestSpec{Audiences: []string{authAudience}}}, metav1.CreateOptions{})
if !apierrors.IsForbidden(err) || (denied != nil && denied.Status.Token != "") {
t.Fatal("revoked caller still obtained a token")
}
api.grant()
waitForAuthentication(t, session.Ready)
if client.Token() == initialToken {
t.Fatal("reauthentication reused revoked OpenBao token")
}
if _, err := store.Read(ctx, credentialPath); err != nil {
t.Fatal("credential access did not recover", err)
}
cancel()
waitForAuthentication(t, func() bool { return !session.Ready() && client.Token() == "" })
}
@@ -49,12 +49,11 @@ type Credentials struct {
} }
// NewCredentials 不登录、不读取环境 token。调用方必须提供专用的已认证 client。 // NewCredentials 不登录、不读取环境 token。调用方必须提供专用的已认证 client。
// 禁用 SDK 写入重试,防止第一次结果丢失后被 CAS 错误掩盖。 // client 由公共 infra 禁用自动重试,防止第一次结果丢失后被 CAS 错误掩盖。
func NewCredentials(client *bao.Client, mount, basePath string) (*Credentials, error) { func NewCredentials(client *bao.Client, mount, basePath string) (*Credentials, error) {
if client == nil || !validPath(mount) || !validPath(basePath) { if client == nil || !validPath(mount) || !validPath(basePath) {
return nil, ErrInvalidLocation return nil, ErrInvalidLocation
} }
client.SetMaxRetries(0)
return &Credentials{kv: client.KVv2(mount), basePath: basePath}, nil return &Credentials{kv: client.KVv2(mount), basePath: basePath}, nil
} }
@@ -104,6 +104,7 @@ func fixtureClient(t *testing.T, address string) *bao.Client {
t.Helper() t.Helper()
config := bao.DefaultConfig() config := bao.DefaultConfig()
config.Address = address config.Address = address
config.MaxRetries = 0
client, err := bao.NewClient(config) client, err := bao.NewClient(config)
if err != nil { if err != nil {
t.Fatal("cannot construct fixture client") t.Fatal("cannot construct fixture client")
@@ -32,6 +32,7 @@ import (
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
"k8s.io/client-go/kubernetes" "k8s.io/client-go/kubernetes"
"k8s.io/client-go/rest" "k8s.io/client-go/rest"
kubeclient "sigs.k8s.io/controller-runtime/pkg/client"
"sigs.k8s.io/controller-runtime/pkg/envtest" "sigs.k8s.io/controller-runtime/pkg/envtest"
secretadapter "git.ddupan.top/panxiao81/ayatori/internal/database/adapter/kubernetes" secretadapter "git.ddupan.top/panxiao81/ayatori/internal/database/adapter/kubernetes"
@@ -244,7 +245,11 @@ func newCredentialFixture(t *testing.T) *credentialFixture {
} }
} }
reader, err := secretadapter.NewSecretCredentials(config, controllerNamespace) apiReader, err := kubeclient.New(config, kubeclient.Options{})
if err != nil {
t.Fatal(err)
}
reader, err := secretadapter.NewSecretCredentials(apiReader, controllerNamespace)
if err != nil { if err != nil {
t.Fatal(err) t.Fatal(err)
} }
@@ -252,7 +257,11 @@ func newCredentialFixture(t *testing.T) *credentialFixture {
if err != nil { if err != nil {
t.Fatal(err) t.Fatal(err)
} }
deniedReader, err := secretadapter.NewSecretCredentials(user.Config(), controllerNamespace) deniedAPIReader, err := kubeclient.New(user.Config(), kubeclient.Options{})
if err != nil {
t.Fatal(err)
}
deniedReader, err := secretadapter.NewSecretCredentials(deniedAPIReader, controllerNamespace)
if err != nil { if err != nil {
t.Fatal(err) t.Fatal(err)
} }
@@ -50,14 +50,6 @@ func TestInstanceControllerWithRealPostgreSQL(t *testing.T) {
t.Fatal(err) t.Fatal(err)
} }
restricted := instanceControllerRBAC(t, f, apiClient) restricted := instanceControllerRBAC(t, f, apiClient)
credentials, err := secretadapter.NewSecretCredentials(restricted, controllerNamespace)
if err != nil {
t.Fatal(err)
}
service, err := application.NewInstanceService(credentials, postgresql.Connector{})
if err != nil {
t.Fatal(err)
}
// 同进程 -count 重复启动测试 manager;生产继续校验 controller 名称唯一。 // 同进程 -count 重复启动测试 manager;生产继续校验 controller 名称唯一。
skipRepeatedName := true skipRepeatedName := true
manager, err := ctrl.NewManager(restricted, ctrl.Options{ manager, err := ctrl.NewManager(restricted, ctrl.Options{
@@ -68,6 +60,14 @@ func TestInstanceControllerWithRealPostgreSQL(t *testing.T) {
if err != nil { if err != nil {
t.Fatal(err) t.Fatal(err)
} }
credentials, err := secretadapter.NewSecretCredentials(manager.GetAPIReader(), controllerNamespace)
if err != nil {
t.Fatal(err)
}
service, err := application.NewInstanceService(credentials, postgresql.Connector{})
if err != nil {
t.Fatal(err)
}
reconciler := &databasecontroller.InstanceReconciler{Observer: service, SecretNamespace: controllerNamespace} reconciler := &databasecontroller.InstanceReconciler{Observer: service, SecretNamespace: controllerNamespace}
if err := reconciler.SetupWithManager(manager); err != nil { if err := reconciler.SetupWithManager(manager); err != nil {
t.Fatal(err) t.Fatal(err)
+192
View File
@@ -0,0 +1,192 @@
/*
Copyright 2026.
Licensed under the Apache License, Version 2.0 (the "License");
you may not use this file except in compliance with the License.
You may obtain a copy of the License at
http://www.apache.org/licenses/LICENSE-2.0
Unless required by applicable law or agreed to in writing, software
distributed under the License is distributed on an "AS IS" BASIS,
WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
See the License for the specific language governing permissions and
limitations under the License.
*/
package openbao
import (
"context"
"errors"
"regexp"
"strings"
"sync"
"sync/atomic"
"time"
kubernetesauth "github.com/openbao/openbao/api/auth/kubernetes/v2"
bao "github.com/openbao/openbao/api/v2"
authenticationv1 "k8s.io/api/authentication/v1"
corev1 "k8s.io/api/core/v1"
"k8s.io/apimachinery/pkg/util/validation"
kubeclient "sigs.k8s.io/controller-runtime/pkg/client"
)
var ErrAuthenticationConfiguration = errors.New("invalid OpenBao Kubernetes authentication configuration")
var authPathSegment = regexp.MustCompile(`^[A-Za-z0-9_-]+$`)
func validAuthMount(mount string) bool {
for segment := range strings.SplitSeq(mount, "/") {
if !authPathSegment.MatchString(segment) {
return false
}
}
return true
}
// KubernetesSession 为专用 SDK client 维护短期登录,不持久化或对外返回 token。
// 每次登录通过 manager 的 Kubernetes 身份申请新的 SA JWT,不依赖 controller 的部署位置。
// 续期调度由官方 LifetimeWatcher 负责,不实现自己的 lease 算法。
type KubernetesSession struct {
client *bao.Client
mount string
role string
kubernetes kubeclient.Client
identity KubernetesIdentity
running sync.Mutex
ready atomic.Bool
}
// KubernetesIdentity 是部署固定的登录目标,不由业务请求选择。
type KubernetesIdentity struct {
Namespace string
ServiceAccount string
Audience string
}
func NewKubernetesSession(
client *bao.Client,
kubernetes kubeclient.Client,
mount, role string,
identity KubernetesIdentity,
) (*KubernetesSession, error) {
if client == nil || kubernetes == nil || !validAuthMount(mount) || !authPathSegment.MatchString(role) ||
len(validation.IsDNS1123Label(identity.Namespace)) != 0 ||
len(validation.IsDNS1123Subdomain(identity.ServiceAccount)) != 0 || strings.TrimSpace(identity.Audience) == "" {
return nil, ErrAuthenticationConfiguration
}
client.ClearToken()
client.SetMaxRetries(0)
client.SetClientTimeout(15 * time.Second)
return &KubernetesSession{
client: client,
mount: mount,
role: role,
kubernetes: kubernetes,
identity: identity,
}, nil
}
// Ready 仅代表当前登录 lease 正由 SDK 管理,不保证下一次后端请求一定成功。
func (s *KubernetesSession) Ready() bool { return s.ready.Load() }
// Start 可交给 manager 管理;关闭时清空本地 token,不撤销共享后端数据。
// SDK 的 Stop 不取消已发出的续期 HTTP 请求,因此等待该请求结束后才退出,最长受 client timeout 限制。
func (s *KubernetesSession) Start(ctx context.Context) error {
if !s.running.TryLock() {
return errors.New("OpenBao authentication is already running")
}
defer s.running.Unlock()
defer s.clear()
for ctx.Err() == nil {
s.clear()
secret := s.login(ctx)
if secret != nil {
s.watch(ctx, secret)
}
s.clear()
// 登录失败及不能继续续期均有限速,防止依赖故障时形成请求忙循环。
select {
case <-ctx.Done():
return nil
case <-time.After(5 * time.Second):
}
}
return nil
}
func (s *KubernetesSession) clear() {
s.ready.Store(false)
s.client.ClearToken()
}
func (s *KubernetesSession) login(ctx context.Context) *bao.Secret {
ctx, cancel := context.WithTimeout(ctx, 15*time.Second)
defer cancel()
// JWT 只用于本次登录,不缓存或自行解析 kubeconfig 中的凭据。
// client-go 负责 kubeconfig/in-cluster 身份与凭据更新;API server 按 RBAC 签发。
expirationSeconds := int64(600)
account := &corev1.ServiceAccount{
Namespace: s.identity.Namespace,
Name: s.identity.ServiceAccount,
}
token := &authenticationv1.TokenRequest{
Spec: authenticationv1.TokenRequestSpec{
Audiences: []string{s.identity.Audience},
ExpirationSeconds: &expirationSeconds,
},
}
// 子资源写入直接请求 API server,不读 cache,也不需要额外的 ServiceAccount get 权限。
err := s.kubernetes.SubResource("token").Create(ctx, account, token)
if err != nil || strings.TrimSpace(token.Status.Token) == "" ||
!token.Status.ExpirationTimestamp.After(time.Now()) {
return nil
}
// helper 会缓存 token,不能跨登录轮次复用。
method, err := kubernetesauth.NewKubernetesAuth(s.role,
kubernetesauth.WithMountPath(s.mount),
kubernetesauth.WithServiceAccountToken(token.Status.Token),
)
if err != nil {
return nil
}
// 先检查短期 lease,再发布到共享 client,避免暴露不合规的登录结果。
secret, err := method.Login(ctx, s.client)
if err != nil || secret == nil || secret.Auth == nil || secret.Auth.ClientToken == "" || secret.Auth.LeaseDuration <= 0 {
return nil
}
if ctx.Err() != nil {
return nil
}
s.client.SetToken(secret.Auth.ClientToken)
return secret
}
func (s *KubernetesSession) watch(ctx context.Context, secret *bao.Secret) {
behavior := bao.RenewBehaviorErrorOnErrors
if !secret.Auth.Renewable {
behavior = bao.RenewBehaviorRenewDisabled
}
watcher, err := s.client.NewLifetimeWatcher(&bao.LifetimeWatcherInput{Secret: secret, RenewBehavior: behavior})
if err != nil {
return
}
s.ready.Store(true)
go watcher.Start()
for {
select {
case <-ctx.Done():
s.clear()
watcher.Stop()
<-watcher.DoneCh()
return
case <-watcher.DoneCh():
watcher.Stop()
return
case <-watcher.RenewCh():
// 不记录 SDK Secret 或 token;无需复制 SDK 已处理的续期数据。
}
}
}
@@ -0,0 +1,233 @@
/*
Copyright 2026.
Licensed under the Apache License, Version 2.0 (the "License");
you may not use this file except in compliance with the License.
You may obtain a copy of the License at
http://www.apache.org/licenses/LICENSE-2.0
Unless required by applicable law or agreed to in writing, software
distributed under the License is distributed on an "AS IS" BASIS,
WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
See the License for the specific language governing permissions and
limitations under the License.
*/
package openbao_test
import (
"context"
"encoding/json"
"net/http"
"net/http/httptest"
"sync/atomic"
"testing"
"time"
bao "github.com/openbao/openbao/api/v2"
authenticationv1 "k8s.io/api/authentication/v1"
corev1 "k8s.io/api/core/v1"
"k8s.io/apimachinery/pkg/api/meta"
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
"k8s.io/apimachinery/pkg/runtime/schema"
"k8s.io/client-go/rest"
kubeclient "sigs.k8s.io/controller-runtime/pkg/client"
"git.ddupan.top/panxiao81/ayatori/internal/infra/openbao"
)
const (
authRole = "controller"
fixtureToken = "AYATORI-TEST-ONLY-bao-token"
)
func fixtureClient(t *testing.T, address string) *bao.Client {
t.Helper()
config := bao.NewConfig()
config.Address = address
client, err := bao.NewClient(config)
if err != nil {
t.Fatal("cannot construct fixture client")
}
client.SetToken(fixtureToken)
return client
}
var testIdentity = openbao.KubernetesIdentity{Namespace: "bao-controller", ServiceAccount: authRole, Audience: "openbao"}
func authenticationClient(t *testing.T, address string) kubeclient.Client {
t.Helper()
// HTTP 单元 fixture 只提供 TokenRequest;静态映射避免额外模拟 discovery API。
mapper := meta.NewDefaultRESTMapper([]schema.GroupVersion{corev1.SchemeGroupVersion})
mapper.Add(corev1.SchemeGroupVersion.WithKind("ServiceAccount"), meta.RESTScopeNamespace)
client, err := kubeclient.New(&rest.Config{Host: address, ContentType: "application/json"}, kubeclient.Options{
Mapper: mapper,
})
if err != nil {
t.Fatal("cannot construct fixture Kubernetes client")
}
return client
}
func authenticationServer(t *testing.T, token func() (string, bool), login http.HandlerFunc) *httptest.Server {
t.Helper()
return httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
if r.URL.Path != "/api/v1/namespaces/bao-controller/serviceaccounts/controller/token" {
login(w, r)
return
}
var request authenticationv1.TokenRequest
if json.NewDecoder(r.Body).Decode(&request) != nil || r.Method != http.MethodPost ||
len(request.Spec.Audiences) != 1 || request.Spec.Audiences[0] != testIdentity.Audience ||
request.Spec.ExpirationSeconds == nil || *request.Spec.ExpirationSeconds != 600 {
t.Error("unexpected TokenRequest target or lifetime")
w.WriteHeader(http.StatusBadRequest)
return
}
jwt, allowed := token()
if !allowed {
w.WriteHeader(http.StatusForbidden)
return
}
w.Header().Set("Content-Type", "application/json")
if err := json.NewEncoder(w).Encode(authenticationv1.TokenRequest{
APIVersion: "authentication.k8s.io/v1", Kind: "TokenRequest",
Status: authenticationv1.TokenRequestStatus{Token: jwt, ExpirationTimestamp: metav1.NewTime(time.Now().Add(10 * time.Minute))},
}); err != nil {
t.Error("cannot encode fixture TokenRequest response")
}
}))
}
func TestKubernetesSessionDoesNotFallbackFromMissingToken(t *testing.T) {
var requests atomic.Int32
for _, missing := range []bool{true, false} {
server := authenticationServer(t, func() (string, bool) { return "", !missing }, func(w http.ResponseWriter, _ *http.Request) {
requests.Add(1)
w.WriteHeader(http.StatusForbidden)
})
defer server.Close()
client := fixtureClient(t, server.URL)
session, err := openbao.NewKubernetesSession(client, authenticationClient(t, server.URL), "kubernetes", authRole, testIdentity)
if err != nil {
t.Fatal(err)
}
ctx, cancel := context.WithTimeout(t.Context(), 100*time.Millisecond)
err = session.Start(ctx)
cancel()
if err != nil || session.Ready() || client.Token() != "" || requests.Load() != 0 {
t.Fatal("denied or empty TokenRequest must not fall back to another identity")
}
}
}
func waitForAuthentication(t *testing.T, check func() bool) {
t.Helper()
deadline := time.NewTimer(20 * time.Second)
defer deadline.Stop()
for !check() {
select {
case <-deadline.C:
t.Fatal("authentication condition timed out")
case <-time.After(20 * time.Millisecond):
}
}
}
func startSession(t *testing.T, session interface{ Start(context.Context) error }) context.CancelFunc {
t.Helper()
ctx, cancel := context.WithCancel(t.Context())
done := make(chan error, 1)
go func() { done <- session.Start(ctx) }()
t.Cleanup(func() {
cancel()
select {
case <-done:
case <-time.After(20 * time.Second):
t.Error("authentication did not stop")
}
})
return cancel
}
func TestKubernetesSessionRequestsNewTokenAfterFailure(t *testing.T) {
var attempts atomic.Int32
var accepted atomic.Bool
var rotated atomic.Bool
server := authenticationServer(t, func() (string, bool) {
if rotated.Load() {
return "rotated-test-jwt", true
}
return "expired-test-jwt", true
}, func(w http.ResponseWriter, r *http.Request) {
if r.URL.Path != "/v1/auth/kubernetes/login" {
w.WriteHeader(http.StatusNotFound)
return
}
var request struct{ JWT, Role string }
if json.NewDecoder(r.Body).Decode(&request) != nil || request.Role != authRole {
t.Error("unexpected authentication request")
w.WriteHeader(http.StatusBadRequest)
return
}
attempts.Add(1)
if request.JWT != "rotated-test-jwt" {
w.WriteHeader(http.StatusForbidden)
return
}
accepted.Store(true)
if err := json.NewEncoder(w).Encode(map[string]any{"auth": map[string]any{
"client_token": fixtureToken, "lease_duration": 60, "renewable": false,
}}); err != nil {
t.Error("cannot encode authentication fixture response")
}
})
defer server.Close()
client := fixtureClient(t, server.URL)
session, err := openbao.NewKubernetesSession(client, authenticationClient(t, server.URL), "kubernetes", authRole, testIdentity)
if err != nil {
t.Fatal(err)
}
if client.Token() != "" {
t.Fatal("constructor retained preexisting static token")
}
cancel := startSession(t, session)
waitForAuthentication(t, func() bool { return attempts.Load() > 0 })
if session.Ready() || client.Token() != "" {
t.Fatal("failed login retained credentials")
}
rotated.Store(true)
waitForAuthentication(t, func() bool { return accepted.Load() && session.Ready() })
if client.Token() != fixtureToken {
t.Fatal("successful login did not configure the client")
}
if session.Start(t.Context()) == nil {
t.Fatal("allowed concurrent lifecycle owners")
}
cancel()
waitForAuthentication(t, func() bool { return !session.Ready() && client.Token() == "" })
}
func TestKubernetesSessionRejectsUnboundedLease(t *testing.T) {
var attempts atomic.Int32
server := authenticationServer(t, func() (string, bool) { return "test-jwt", true }, func(w http.ResponseWriter, _ *http.Request) {
attempts.Add(1)
if err := json.NewEncoder(w).Encode(map[string]any{"auth": map[string]any{
"client_token": fixtureToken, "lease_duration": 0,
}}); err != nil {
t.Error("cannot encode authentication fixture response")
}
})
defer server.Close()
client := fixtureClient(t, server.URL)
session, err := openbao.NewKubernetesSession(client, authenticationClient(t, server.URL), "kubernetes", authRole, testIdentity)
if err != nil {
t.Fatal(err)
}
startSession(t, session)
waitForAuthentication(t, func() bool { return attempts.Load() > 0 })
if session.Ready() || client.Token() != "" {
t.Fatal("accepted a token without a finite lease")
}
}
+51
View File
@@ -0,0 +1,51 @@
/*
Copyright 2026.
Licensed under the Apache License, Version 2.0 (the "License");
you may not use this file except in compliance with the License.
You may obtain a copy of the License at
http://www.apache.org/licenses/LICENSE-2.0
Unless required by applicable law or agreed to in writing, software
distributed under the License is distributed on an "AS IS" BASIS,
WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
See the License for the specific language governing permissions and
limitations under the License.
*/
// Package openbao 管理控制面共享的 OpenBao 连接与认证,不依赖任何产品领域。
package openbao
import (
"errors"
"net/url"
"strings"
"time"
bao "github.com/openbao/openbao/api/v2"
)
// NewClient 创建供读写适配器共享的官方 SDK client;身份由 KubernetesSession 管理。
func NewClient(addressURL, caCert string) (*bao.Client, error) {
address, err := url.Parse(addressURL)
if err != nil || address.Scheme != "https" || address.Host == "" || address.User != nil ||
address.RawQuery != "" || address.ForceQuery || address.Fragment != "" ||
(address.Path != "" && address.Path != "/") {
return nil, errors.New("OpenBao address must be an absolute HTTPS URL without credentials, query, fragment or path")
}
// NewConfig 不读取 BAO_TOKEN/BAO_SKIP_VERIFY 等环境配置,不允许旁路 Kubernetes 身份或 TLS。
config := bao.NewConfig()
config.Address = strings.TrimSuffix(addressURL, "/")
// 公共读写 client 不自动重试:写入结果不确定时交由具体用例决定恢复行为。
config.MaxRetries = 0
config.Timeout = 15 * time.Second
if config.Error != nil || config.ConfigureTLS(&bao.TLSConfig{CACert: caCert}) != nil {
return nil, errors.New("cannot configure OpenBao TLS trust")
}
client, err := bao.NewClient(config)
if err != nil {
return nil, errors.New("cannot construct OpenBao client")
}
return client, nil
}
+67
View File
@@ -0,0 +1,67 @@
/*
Copyright 2026.
Licensed under the Apache License, Version 2.0 (the "License");
you may not use this file except in compliance with the License.
You may obtain a copy of the License at
http://www.apache.org/licenses/LICENSE-2.0
Unless required by applicable law or agreed to in writing, software
distributed under the License is distributed on an "AS IS" BASIS,
WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
See the License for the specific language governing permissions and
limitations under the License.
*/
package openbao_test
import (
"net/http"
"net/http/httptest"
"testing"
"git.ddupan.top/panxiao81/ayatori/internal/infra/openbao"
)
func TestOpenBaoClientDoesNotUseEnvironmentIdentityOrAddress(t *testing.T) {
t.Setenv("BAO_TOKEN", "TEST-ONLY-unwanted-static-token")
t.Setenv("BAO_ADDR", "http://unwanted.invalid")
t.Setenv("BAO_SKIP_VERIFY", "true")
client, err := openbao.NewClient("https://bao.example/", "")
if err != nil {
t.Fatal(err)
}
if client.Address() != "https://bao.example" || client.Token() != "" {
t.Fatal("ambient environment replaced the explicit connection or identity")
}
if client.MaxRetries() != 0 {
t.Fatal("shared client must not automatically retry uncertain writes")
}
server := httptest.NewTLSServer(http.HandlerFunc(func(w http.ResponseWriter, _ *http.Request) {
w.WriteHeader(http.StatusOK)
}))
defer server.Close()
untrusted, err := openbao.NewClient(server.URL, "")
if err != nil {
t.Fatal(err)
}
untrusted.SetMaxRetries(0)
if _, err := untrusted.Sys().HealthWithContext(t.Context()); err == nil {
t.Fatal("BAO_SKIP_VERIFY bypassed TLS validation")
}
}
func TestOpenBaoClientRejectsUnsafeConfiguration(t *testing.T) {
for _, address := range []string{
"", "http://bao.example", "https://user:[email protected]",
"https://bao.example/?token=secret", "https://bao.example/#secret", "https://bao.example/path",
} {
if _, err := openbao.NewClient(address, ""); err == nil {
t.Fatal("accepted unsafe OpenBao address")
}
}
if _, err := openbao.NewClient("https://bao.example", "/nonexistent/fixture-ca"); err == nil {
t.Fatal("accepted missing explicit CA")
}
}