Compare commits

..
Author SHA1 Message Date
panxiao81 7834cab97f refactor: 集中显式注入并归位凭据领域规则
Verify / test (pull_request) Successful in 13m29s
Verify / lint (pull_request) Successful in 14m13s
Verify / database-integration (pull_request) Successful in 15m48s
2026-09-27 19:16:48 +00:00
panxiao81 72ce3eda40 Merge pull request 'feat: 接通 Database 凭据准备闭环' (#14) from feat/database-credential-checkpoint into main 2026-09-27 18:57:35 +00:00
panxiao81 35ada6d7eb feat: 接通 Database 凭据准备闭环
Verify / test (pull_request) Successful in 11m23s
Verify / lint (pull_request) Successful in 12m18s
Verify / database-integration (pull_request) Successful in 13m49s
2026-09-27 18:35:13 +00:00
panxiao81 371248659b feat: 固定 Database 凭据位置与确认版本
Verify / test (pull_request) Successful in 13m42s
Verify / lint (pull_request) Successful in 14m26s
Verify / database-integration (pull_request) Successful in 16m16s
2026-09-27 18:08:09 +00:00
panxiao81 6b2808ce91 Merge pull request 'feat: OpenBao Kubernetes 认证与短期会话续期' (#13) from feat/database-openbao-auth into main 2026-09-27 17:48:22 +00:00
panxiao81 356e21686f refactor: 按职责拆分 controller 启动流程
Verify / test (pull_request) Successful in 9m48s
Verify / lint (pull_request) Successful in 10m42s
Verify / database-integration (pull_request) Successful in 12m1s
2026-09-27 17:37:45 +00:00
panxiao81 c7b80890fb refactor: 将 controller 装配收拢到 bootstrap 包
Verify / lint (pull_request) Successful in 9m21s
Verify / test (pull_request) Successful in 10m8s
Verify / database-integration (pull_request) Successful in 10m3s
2026-09-27 17:06:35 +00:00
panxiao81 65c60cca45 refactor: 分离公共基础设施与 Database 适配器
Verify / lint (pull_request) Successful in 13m32s
Verify / test (pull_request) Successful in 14m5s
Verify / database-integration (pull_request) Successful in 15m1s
2026-09-25 21:55:24 +00:00
panxiao81 c20f8930f0 refactor: OpenBao 认证复用 manager Kubernetes client
Verify / test (pull_request) Successful in 4m36s
Verify / lint (pull_request) Successful in 12m25s
Verify / database-integration (pull_request) Successful in 15m49s
2026-09-25 21:38:56 +00:00
panxiao81 f0aa86f676 fix: 复用 manager kubeconfig 通过 TokenRequest 登录 OpenBao
Verify / lint (pull_request) Successful in 11m47s
Verify / database-integration (pull_request) Successful in 15m24s
Verify / test (pull_request) Successful in 16m41s
2026-09-25 21:02:06 +00:00
panxiao81 dcf9ab50df feat: 接入 OpenBao Kubernetes 认证与短期会话续期
Verify / test (pull_request) Successful in 15m55s
Verify / lint (pull_request) Successful in 17m56s
Verify / database-integration (pull_request) Successful in 19m54s
2026-09-25 19:39:56 +00:00
panxiao81 22ab72ec60 Merge pull request 'feat: OpenBao 应用凭据安全存储切片' (#12) from feat/database-provisioning into main 2026-09-25 19:26:15 +00:00
panxiao81 2a4f622f44 test: 分离 OpenBao 镜像拉取与启动超时并保留诊断
Verify / test (pull_request) Successful in 5m24s
Verify / lint (pull_request) Successful in 11m53s
Verify / database-integration (pull_request) Successful in 13m32s
2026-09-25 16:30:09 +00:00
panxiao81 f6bb9e4599 feat: 增加 OpenBao 应用凭据安全存储切片
Verify / test (pull_request) Successful in 5m43s
Verify / lint (pull_request) Successful in 14m8s
Verify / database-integration (pull_request) Failing after 13m35s
2026-09-25 15:41:53 +00:00
panxiao81 f4deb98a7f Merge pull request 'feat: Instance 原生管理观测与删除保护' (#11) from feat/database-instance-observation into main 2026-09-25 14:32:25 +00:00
67 changed files with 4565 additions and 347 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:
+5 -3
View File
@@ -68,12 +68,14 @@ lint: golangci-lint ## Run golangci-lint linter
"$(GOLANGCI_LINT)" run "$(GOLANGCI_LINT)" run
.PHONY: test-database-integration .PHONY: test-database-integration
test-database-integration: setup-envtest ## 使用临时 API server 与独立 PostgreSQL 容器验证凭据读取和连接更新。 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
+2 -2
View File
@@ -54,8 +54,8 @@ const (
ReclaimDelete ReclaimPolicy = "Delete" ReclaimDelete ReclaimPolicy = "Delete"
) )
// CredentialReference 定位已有 OpenBao KV v2 凭据,不包含任何秘密值。 // CredentialReference 定位 OpenBao KV v2 凭据,不包含任何秘密值。
// 只由资源管理员在导入时填写;controller 必须检查部署允许的 mount/path 范围。 // 管理员在导入声明中指定,controller 在 status 中固定位置;两者均须检查部署允许的范围。
type CredentialReference struct { type CredentialReference struct {
// +kubebuilder:validation:MinLength=1 // +kubebuilder:validation:MinLength=1
// +kubebuilder:validation:MaxLength=253 // +kubebuilder:validation:MaxLength=253
@@ -0,0 +1,87 @@
package v1alpha1_test
import (
"testing"
databasev1alpha1 "git.ddupan.top/panxiao81/ayatori/api/database/v1alpha1"
apierrors "k8s.io/apimachinery/pkg/api/errors"
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
ctrlclient "sigs.k8s.io/controller-runtime/pkg/client"
)
func testCredentialStatus(t *testing.T, client ctrlclient.Client) {
database := validDatabase("credential-status")
if err := client.Create(t.Context(), database); err != nil {
t.Fatal(err)
}
withoutLocation := database.DeepCopy()
withoutLocation.Status.CredentialVersion = 1
if err := client.Status().Update(t.Context(), withoutLocation); !apierrors.IsInvalid(err) {
t.Fatalf("没有位置不能确认版本: %v", err)
}
database.Status.CredentialRef = &databasev1alpha1.CredentialReference{
Mount: "application-secrets", Path: "applications/" + string(database.UID),
}
if err := client.Status().Update(t.Context(), database); err != nil {
t.Fatal(err)
}
testCredentialStatusRemoval(t, client, database)
changedTarget := database.DeepCopy()
changedTarget.Spec.Database = "another_database"
if err := client.Update(t.Context(), changedTarget); !apierrors.IsInvalid(err) {
t.Fatalf("凭据位置固定后不得更换实际目标: %v", err)
}
stale := database.DeepCopy()
database.Status.CredentialVersion = 1
if err := client.Status().Update(t.Context(), database); err != nil {
t.Fatal(err)
}
stale.Status.Phase = "Binding"
if err := client.Status().Update(t.Context(), stale); !apierrors.IsConflict(err) {
t.Fatalf("旧 resourceVersion 不得覆盖确认结果: %v", err)
}
testCredentialStatusRemoval(t, client, database)
for _, version := range []int64{0, -1, 2} {
changed := database.DeepCopy()
changed.Status.CredentialVersion = version
if err := client.Status().Update(t.Context(), changed); !apierrors.IsInvalid(err) {
t.Fatalf("不可移除或更改确认版本 %d: %v", version, err)
}
}
database.Status.Conditions = []metav1.Condition{{
Type: "CredentialsReady", Status: metav1.ConditionFalse,
Reason: "DependencyUnavailable", Message: "凭据暂时无法读取",
ObservedGeneration: database.Generation, LastTransitionTime: metav1.Now(),
}}
if err := client.Status().Update(t.Context(), database); err != nil {
t.Fatalf("当前不可用不应阻止保留确认记录: %v", err)
}
reloaded := &databasev1alpha1.PostgreSQLDatabase{}
if err := client.Get(t.Context(), ctrlclient.ObjectKeyFromObject(database), reloaded); err != nil {
t.Fatal(err)
}
if reloaded.Status.CredentialVersion != 1 || reloaded.Status.CredentialRef == nil {
t.Fatal("重读丢失凭据确认记录")
}
}
func testCredentialStatusRemoval(t *testing.T, client ctrlclient.Client, database *databasev1alpha1.PostgreSQLDatabase) {
t.Helper()
for _, change := range []struct {
name string
edit func(*databasev1alpha1.PostgreSQLDatabase)
}{
{"移除整个 status", func(d *databasev1alpha1.PostgreSQLDatabase) { d.Status = databasev1alpha1.PostgreSQLDatabaseStatus{} }},
{"移除位置", func(d *databasev1alpha1.PostgreSQLDatabase) { d.Status.CredentialRef = nil }},
{"修改 mount", func(d *databasev1alpha1.PostgreSQLDatabase) { d.Status.CredentialRef.Mount = "other" }},
{"修改 path", func(d *databasev1alpha1.PostgreSQLDatabase) { d.Status.CredentialRef.Path = "applications/other" }},
} {
t.Run(change.name, func(t *testing.T) {
changed := database.DeepCopy()
change.edit(changed)
if err := client.Status().Update(t.Context(), changed); !apierrors.IsInvalid(err) {
t.Fatalf("不应接受已固定凭据位置的更改: %v", err)
}
})
}
}
@@ -30,6 +30,14 @@ type PostgreSQLDatabaseStatus struct {
// InstanceUID 记录观察时的实例身份,不把同名新实例视为原目标。 // InstanceUID 记录观察时的实例身份,不把同名新实例视为原目标。
// +optional // +optional
InstanceUID types.UID `json:"instanceUID,omitempty"` InstanceUID types.UID `json:"instanceUID,omitempty"`
// CredentialRef 在首次外部写入前固定凭据位置;部署配置变化不迁移此位置。
// +optional
CredentialRef *CredentialReference `json:"credentialRef,omitempty"`
// CredentialVersion 只在创建并回读成功后记录,不表示凭据当前仍可用。
// 省略表示尚未确认;已有凭据不能仅凭读取成功补记确认。
// +kubebuilder:validation:Minimum=1
// +optional
CredentialVersion int64 `json:"credentialVersion,omitempty"`
// Phase 暂不冻结供应子阶段枚举;它不是操作授权或绑定的替代记录。 // Phase 暂不冻结供应子阶段枚举;它不是操作授权或绑定的替代记录。
// +optional // +optional
Phase string `json:"phase,omitempty"` Phase string `json:"phase,omitempty"`
@@ -42,7 +50,10 @@ type PostgreSQLDatabaseStatus struct {
// +kubebuilder:object:root=true // +kubebuilder:object:root=true
// +kubebuilder:subresource:status // +kubebuilder:subresource:status
// +kubebuilder:resource:scope=Cluster // +kubebuilder:resource:scope=Cluster
// +kubebuilder:validation:XValidation:rule="!(has(oldSelf.spec.tenantRef) || (has(oldSelf.status) && has(oldSelf.status.instanceUID))) || (self.spec.instanceRef == oldSelf.spec.instanceRef && self.spec.database == oldSelf.spec.database && self.spec.loginRole == oldSelf.spec.loginRole && self.spec.source == oldSelf.spec.source && has(self.spec.credentialRef) == has(oldSelf.spec.credentialRef) && (!has(oldSelf.spec.credentialRef) || self.spec.credentialRef == oldSelf.spec.credentialRef))",message="managed database target cannot change after observation or binding starts" // +kubebuilder:validation:XValidation:rule="!has(self.status) || !has(self.status.credentialVersion) || has(self.status.credentialRef)",message="credentialVersion requires credentialRef"
// +kubebuilder:validation:XValidation:rule="!(has(oldSelf.status) && has(oldSelf.status.credentialRef)) || (has(self.status) && has(self.status.credentialRef) && self.status.credentialRef == oldSelf.status.credentialRef)",message="recorded credentialRef cannot change or be removed"
// +kubebuilder:validation:XValidation:rule="!(has(oldSelf.status) && has(oldSelf.status.credentialVersion)) || (has(self.status) && has(self.status.credentialVersion) && self.status.credentialVersion == oldSelf.status.credentialVersion)",message="confirmed credentialVersion cannot change or be removed"
// +kubebuilder:validation:XValidation:rule="!(has(oldSelf.spec.tenantRef) || (has(oldSelf.status) && (has(oldSelf.status.instanceUID) || has(oldSelf.status.credentialRef)))) || (self.spec.instanceRef == oldSelf.spec.instanceRef && self.spec.database == oldSelf.spec.database && self.spec.loginRole == oldSelf.spec.loginRole && self.spec.source == oldSelf.spec.source && has(self.spec.credentialRef) == has(oldSelf.spec.credentialRef) && (!has(oldSelf.spec.credentialRef) || self.spec.credentialRef == oldSelf.spec.credentialRef))",message="managed database target cannot change after observation or binding starts"
// +kubebuilder:printcolumn:name="Instance",type=string,JSONPath=`.spec.instanceRef.name` // +kubebuilder:printcolumn:name="Instance",type=string,JSONPath=`.spec.instanceRef.name`
// +kubebuilder:printcolumn:name="Database",type=string,JSONPath=`.spec.database` // +kubebuilder:printcolumn:name="Database",type=string,JSONPath=`.spec.database`
// +kubebuilder:printcolumn:name="Ready",type=string,JSONPath=`.status.conditions[?(@.type=='Ready')].status` // +kubebuilder:printcolumn:name="Ready",type=string,JSONPath=`.status.conditions[?(@.type=='Ready')].status`
+1
View File
@@ -61,6 +61,7 @@ func TestDatabaseAPI(t *testing.T) {
t.Run("作用域和默认值", func(t *testing.T) { testDefaults(t, client) }) t.Run("作用域和默认值", func(t *testing.T) { testDefaults(t, client) })
t.Run("拒绝非法声明", func(t *testing.T) { testInvalidDeclarations(t, client) }) t.Run("拒绝非法声明", func(t *testing.T) { testInvalidDeclarations(t, client) })
t.Run("status隔离和绑定并发", func(t *testing.T) { testBindingWrites(t, client) }) t.Run("status隔离和绑定并发", func(t *testing.T) { testBindingWrites(t, client) })
t.Run("凭据位置与确认版本", func(t *testing.T) { testCredentialStatus(t, client) })
t.Run("仓库示例", func(t *testing.T) { testSamples(t, client, scheme) }) t.Run("仓库示例", func(t *testing.T) { testSamples(t, client, scheme) })
} }
@@ -188,6 +188,11 @@ func (in *PostgreSQLDatabaseSpec) DeepCopy() *PostgreSQLDatabaseSpec {
// DeepCopyInto is an autogenerated deepcopy function, copying the receiver, writing into out. in must be non-nil. // DeepCopyInto is an autogenerated deepcopy function, copying the receiver, writing into out. in must be non-nil.
func (in *PostgreSQLDatabaseStatus) DeepCopyInto(out *PostgreSQLDatabaseStatus) { func (in *PostgreSQLDatabaseStatus) DeepCopyInto(out *PostgreSQLDatabaseStatus) {
*out = *in *out = *in
if in.CredentialRef != nil {
in, out := &in.CredentialRef, &out.CredentialRef
*out = new(CredentialReference)
**out = **in
}
if in.Conditions != nil { if in.Conditions != nil {
in, out := &in.Conditions, &out.Conditions in, out := &in.Conditions, &out.Conditions
*out = make([]v1.Condition, len(*in)) *out = make([]v1.Condition, len(*in))
-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)
} }
} }
@@ -50,8 +50,8 @@ spec:
properties: properties:
credentialRef: credentialRef:
description: |- description: |-
CredentialReference 定位已有 OpenBao KV v2 凭据,不包含任何秘密值。 CredentialReference 定位 OpenBao KV v2 凭据,不包含任何秘密值。
只由资源管理员在导入时填写;controller 必须检查部署允许的 mount/path 范围。 管理员在导入声明中指定,controller 在 status 中固定位置;两者均须检查部署允许的范围。
properties: properties:
mount: mount:
maxLength: 253 maxLength: 253
@@ -198,6 +198,29 @@ spec:
x-kubernetes-list-map-keys: x-kubernetes-list-map-keys:
- type - type
x-kubernetes-list-type: map x-kubernetes-list-type: map
credentialRef:
description: CredentialRef 在首次外部写入前固定凭据位置;部署配置变化不迁移此位置。
properties:
mount:
maxLength: 253
minLength: 1
type: string
path:
description: Path 是 mount 内的逻辑路径,不含 KV v2 的 data/ API 前缀。
maxLength: 1024
minLength: 1
type: string
required:
- mount
- path
type: object
credentialVersion:
description: |-
CredentialVersion 只在创建并回读成功后记录,不表示凭据当前仍可用。
省略表示尚未确认;已有凭据不能仅凭读取成功补记确认。
format: int64
minimum: 1
type: integer
instanceUID: instanceUID:
description: InstanceUID 记录观察时的实例身份,不把同名新实例视为原目标。 description: InstanceUID 记录观察时的实例身份,不把同名新实例视为原目标。
type: string type: string
@@ -212,13 +235,22 @@ spec:
- spec - spec
type: object type: object
x-kubernetes-validations: x-kubernetes-validations:
- message: credentialVersion requires credentialRef
rule: '!has(self.status) || !has(self.status.credentialVersion) || has(self.status.credentialRef)'
- message: recorded credentialRef cannot change or be removed
rule: '!(has(oldSelf.status) && has(oldSelf.status.credentialRef)) || (has(self.status)
&& has(self.status.credentialRef) && self.status.credentialRef == oldSelf.status.credentialRef)'
- message: confirmed credentialVersion cannot change or be removed
rule: '!(has(oldSelf.status) && has(oldSelf.status.credentialVersion)) ||
(has(self.status) && has(self.status.credentialVersion) && self.status.credentialVersion
== oldSelf.status.credentialVersion)'
- message: managed database target cannot change after observation or binding - message: managed database target cannot change after observation or binding
starts starts
rule: '!(has(oldSelf.spec.tenantRef) || (has(oldSelf.status) && has(oldSelf.status.instanceUID))) rule: '!(has(oldSelf.spec.tenantRef) || (has(oldSelf.status) && (has(oldSelf.status.instanceUID)
|| (self.spec.instanceRef == oldSelf.spec.instanceRef && self.spec.database || has(oldSelf.status.credentialRef)))) || (self.spec.instanceRef == oldSelf.spec.instanceRef
== oldSelf.spec.database && self.spec.loginRole == oldSelf.spec.loginRole && self.spec.database == oldSelf.spec.database && self.spec.loginRole
&& self.spec.source == oldSelf.spec.source && has(self.spec.credentialRef) == oldSelf.spec.loginRole && self.spec.source == oldSelf.spec.source &&
== has(oldSelf.spec.credentialRef) && (!has(oldSelf.spec.credentialRef) has(self.spec.credentialRef) == has(oldSelf.spec.credentialRef) && (!has(oldSelf.spec.credentialRef)
|| self.spec.credentialRef == oldSelf.spec.credentialRef))' || self.spec.credentialRef == oldSelf.spec.credentialRef))'
served: true served: true
storage: true storage: true
@@ -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
+39
View File
@@ -43,6 +43,45 @@ 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,再释放连接。
Database 使用 Wire 式的显式构造函数注入,目前手写装配,不引入 Wire 生成器、Dig/Fx
容器或运行时服务查找。`database_wiring.go` 集中展示 repository → 用例 → controller 的
对象构造;`registerDatabaseControllers` 负责注册,返回的 `closeDatabaseConnections`
只在 worker 停止后关闭 Instance 连接。controller 在启动时接收完整依赖,不在 Reconcile
或 SetupWithManager 中补建 adapter/service;缺失依赖在注册时失败。
domain 持有凭据值对象、准备资格与创建恢复规则;application 只组织读写、调用领域判断
和保存结果。用例需要的 repository 接口与并发快照仍由消费方定义,不把所有类型都强塞
进领域。Kubernetes adapter 将领域阶段映射为已有 Conditions,领域不依赖其字符串协议。
单测检查 domain/controller 的依赖边界,并验证共享 writer、直连 reader 和服务的注入。
基础设施能力属于整个 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、
+107 -2
View File
@@ -43,8 +43,8 @@ registry 准备决策;这一依赖现已从代码移除,不能把旧运行
- 领域层不依赖 Kubernetes types、数据库 driver 或凭据 provider。 - 领域层不依赖 Kubernetes types、数据库 driver 或凭据 provider。
- CredentialReference 只携带管理 Secret 的名称与字段映射,不包含 Secret 内容或 OpenBao path。 - CredentialReference 只携带管理 Secret 的名称与字段映射,不包含 Secret 内容或 OpenBao path。
- Instance checkpoint 不是外部事实;实际能力必须由 application/adapter 观察后交给领域对象判断。 - Instance checkpoint 不是外部事实;实际能力必须由 application/adapter 观察后交给领域对象判断。
- 当前代码只检查 Instance 供应前置条件,不授予 Tenant 所有权或外部写入权限,也不表示 - Instance 观察只检查供应前置条件,不单独授予 Tenant 所有权或外部写入权限;凭据准备
Database API 已经可用。 另行校验双向绑定和保护,尚不表示完整 Database API 已经可用。
## 管理凭据与连接切片 ## 管理凭据与连接切片
@@ -151,6 +151,111 @@ Instance 删除首先释放本地连接并撤销 Ready。任何引用它的 Data
轮换、删除、跨 namespace 拒绝和 watch。API 测试另覆盖写入版本冲突、幂等、新 reconciler 轮换、删除、跨 namespace 拒绝和 watch。API 测试另覆盖写入版本冲突、幂等、新 reconciler
恢复与引用删除保护。完整 DBaaS 仍需供应/导入、OpenBao/ESO、Retain/Delete 集成验收。 恢复与引用删除保护。完整 DBaaS 仍需供应/导入、OpenBao/ESO、Retain/Delete 集成验收。
## 应用凭据存储切片
`adapter/openbao` 使用官方 Go SDK `api/v2 v2.7.0` 的 KV v2 API,只有创建和读取,
不维护 registry、不覆盖已有密码。动态路径由固定前缀与 Database UID 组成;所有访问都校验
配置前缀,已有导入位置也不能绕过 controller 的凭据权限范围。
创建使用 CAS=0,随后回读七键和版本 1;已有值或软删除历史报冲突。关闭 SDK 自动重试,
写入响应丢失、回读失败或内容变化均返回不确定结果,上层不得生成第二份密码或自动认领。
`Read` 只适用于调用方已确认关联的路径,读取成功本身不是管理权证据。错误不传播 SDK
响应体;内存凭据的普通格式化及 JSON 输出均脱敏,明确的 `SecretData` 才返回明文七键。
依据官方 [KV v2 CAS 合同](https://github.com/openbao/openbao/blob/main/internal/builtin/logical/kv/path_data.go)
与 [Go SDK](https://github.com/openbao/openbao/tree/main/api)。`make test-database-integration`
现包含独立 OpenBao dev 容器,固定摘要、随机回环端口、无持久卷,不接受外部地址。
真实后端覆盖创建/回读、并发唯一创建、重建适配器读取、软删除冲突、固定前缀 token
拒绝管理路径,以及成功写入后丢失响应;HTTP 故障测试补充不重试和错误脱敏。
认证会话和凭据准备用例已接入 manager,但默认不启用外部写入。
Database 已有 `status.credentialRef` 和 `status.credentialVersion` 的字段与 CEL 校验:
固定位置、只在创建并回读成功后确认版本,二者写入后不可清空或修改。
`ReadConfirmed` 按已确认版本检查最新 KV 值,删除或版本漂移报 Conflict,不读取旧版本掩盖变化。
测试 token 只用于临时 fixture,不是生产静态 token 配置接口。
### 凭据准备闭环
配置认证后,显式设置 `--database-credential-mount` 启用准备;路径默认
`applications/<Database UID>`,前缀由 `--database-credential-prefix` 配置。
`CredentialPreparation` 用例检查 Provision 来源、双向绑定 UID、申请目标、finalizer、
删除/Released 状态及当前 Instance Ready;导入资源不执行本流程。
顺序为固定位置 → 确认后端尚无凭据 → 保存 CreationStarted → 创建并回读 → 保存确认版本。
`CredentialReconciler` 只负责 Database/Tenant/Instance watch 和 30 秒依赖重查;
Kubernetes repository 负责直接读取及有 resourceVersion 保护的状态更新。
三类 controller 的依赖均在 bootstrap 显式组装,不在 Reconcile 中构造服务。
`domain/credential` 保存应用凭据值对象、准备资格、固定位置与未确认创建的恢复规则;
application 保留 I/O 顺序、消费方接口及并发快照,不再复用绑定用例的快照类型。
领域阶段与 `CredentialsReady`/Reason 的转换由 Kubernetes adapter 负责,已有 API 保持兼容。
外部操作前后回查目标:Database 必须仍是同一 UID/resourceVersion,Tenant/Instance 必须
保持绑定、spec generation、删除状态、finalizer 保护与有效 Ready;无关 Conditions 刷新
不构成目标变化。不存在跨 Kubernetes/OpenBao 原子事务:中途发生相关变化时停止确认,
留下可观察状态,交给后续协调或人工核实,不盲目重试过期的确认写入。
`CredentialsReady=False/CreationStarted` 是正在进行而未确认的诊断,不是成功证明。
任何新一轮协调遇到该状态且没有确认版本,都转为 Conflict;即使进程在发出请求前退出
也采用这一保守边界。明确的认证/权限拒绝可等待恢复;响应丢失、回读失败、确认保存失败
及未确认的已有值均不会自动认领或重新生成密码。Conflict 保留首次原因,普通依赖错误
保留已确认版本;配置变化停止在原位置,不搬迁或覆盖。并行实例中途撞上已开始的创建
同样可能保守地要求人工处理,不承诺无损接续;多副本部署应启用既有 leader election。
成功只设置 `CredentialsReady=True`,Database Ready 仍为 False/ProvisioningIncomplete,
Tenant 仍未完成交付。本切片没有 PostgreSQL role/database 创建、扩展安装、ESO 投射、
Retain 释放或 Delete 清理,也不会解除 finalizer。不要作为完整 DBaaS 部署。
单元测试穷举前置条件;真实 API server + 隔离 Bao 验证创建、状态确认、幂等、重启、并发、
依赖恢复、固定位置、不确定结果、确认保存失败和删除边界;实际 manager 验证 watch 驱动
及重启。该 fixture 只声明 Instance 前置 Ready,真实 PostgreSQL 管理能力由既有 Instance
集成测试覆盖,不把凭据准备验收当成实际建库或应用登录验收。
## 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):规范性行为与验收标准;
+13 -2
View File
@@ -4,7 +4,7 @@
| --- | --- | | --- | --- |
| 状态 | API schema 与绑定 controller 已实现;供应、交付与删除清理未接入 | | 状态 | API schema 与绑定 controller 已实现;供应、交付与删除清理未接入 |
| API group/version | `database.ayatori.ddupan.top/v1alpha1` | | API group/version | `database.ayatori.ddupan.top/v1alpha1` |
| 最后更新 | 2026-09-25 | | 最后更新 | 2026-09-27 |
以 [系统规格](specification.md) 与 以 [系统规格](specification.md) 与
[ADR-0009](../decisions/0009-database-resource-and-claim.md) 为准。类型与生成的 CRD 已纳入源码, [ADR-0009](../decisions/0009-database-resource-and-claim.md) 为准。类型与生成的 CRD 已纳入源码,
@@ -24,6 +24,8 @@ Go 类型位于 `api/database/v1alpha1`,CRD 随 `config/crd` 发布;manager
| Database | `spec.reclaimPolicy` | Retain 默认或 Delete | | Database | `spec.reclaimPolicy` | Retain 默认或 Delete |
| Database | `spec.tenantRef.namespace/name/uid` | controller 写入的完整绑定身份,不是允许名单 | | Database | `spec.tenantRef.namespace/name/uid` | controller 写入的完整绑定身份,不是允许名单 |
| Database | `status.instanceUID` | 观察时的 Instance 身份 | | Database | `status.instanceUID` | 观察时的 Instance 身份 |
| Database | `status.credentialRef.mount/path` | 首次写入前固定的 KV v2 位置,不随部署配置迁移 |
| Database | `status.credentialVersion` | 创建并回读成功后确认的正整数版本;省略表示未确认 |
| Tenant | `spec.provision.instanceRef.name` | 动态申请来源,与 `spec.databaseRef` 互斥且必须二选一 | | Tenant | `spec.provision.instanceRef.name` | 动态申请来源,与 `spec.databaseRef` 互斥且必须二选一 |
| Tenant | `spec.provision.database/loginRole` | 可省略,语义默认值由 controller 解析,不由 CRD 推导 | | Tenant | `spec.provision.database/loginRole` | 可省略,语义默认值由 controller 解析,不由 CRD 推导 |
| Tenant | `spec.databaseRef.name` | 显式申请已有 Database,不额外指定 Instance | | Tenant | `spec.databaseRef.name` | 显式申请已有 Database,不额外指定 Instance |
@@ -36,6 +38,15 @@ Instance phase 沿用已批准枚举;Database/Tenant phase 暂不冻结供应
`credentialRef.path` 是 mount 内逻辑路径,不包含 KV v2 的 `data/` 前缀。 `credentialRef.path` 是 mount 内逻辑路径,不包含 KV v2 的 `data/` 前缀。
其部署允许范围、实际凭据读取和 URL 安全构造仍由后续 adapter/controller 验证。 其部署允许范围、实际凭据读取和 URL 安全构造仍由后续 adapter/controller 验证。
凭据位置与确认版本一旦写入便不可更改或移除;确认版本必须有对应位置。Conditions 描述
当前可用性,不替代确认记录,也不能因读取暂时失败而清空记录。上述 schema 已有真实 API
server 校验;独立的凭据准备用例在显式启用后填写这些字段,绑定 controller 不负责外部写入。
未确认的已有值必须报 Conflict,不能用读取成功补记版本;已确认版本的恢复读取检查最新
KV 版本,删除或版本漂移均需人工处理,不回退旧版本或生成替代密码。第一版不提供轮换入口。
`CredentialsReady` 条件只描述凭据准备结果;它不授权实际数据库交付。
`CreationStarted` 且无确认版本表示创建未完成确认,重入时停在 Conflict;不尝试推断
进程中断前请求是否发出。完整执行和测试边界见[凭据准备闭环](README.md#凭据准备闭环)。
示例:[Instance](../../config/samples/database_v1alpha1_postgresqlinstance.yaml)、 示例:[Instance](../../config/samples/database_v1alpha1_postgresqlinstance.yaml)、
[导入 Database](../../config/samples/database_v1alpha1_postgresqldatabase.yaml)、 [导入 Database](../../config/samples/database_v1alpha1_postgresqldatabase.yaml)、
[动态/已有资源申请](../../config/samples/database_v1alpha1_postgresqltenant.yaml)。 [动态/已有资源申请](../../config/samples/database_v1alpha1_postgresqltenant.yaml)。
@@ -47,7 +58,7 @@ API 接受两个 Tenant 引用同一 Database 不表示允许双重绑定;排
已有资源必须有当前版本 Ready 观察、匹配的 Instance UID,并处于未绑定的 Available 状态。 已有资源必须有当前版本 Ready 观察、匹配的 Instance UID,并处于未绑定的 Available 状态。
同一 Tenant 的资源侧记录已写入时,允许回读后补齐申请侧,不重新争抢资源。 同一 Tenant 的资源侧记录已写入时,允许回读后补齐申请侧,不重新争抢资源。
Tenant 进入 `status.phase=Binding` 后由 CEL 固定申请目标;Database 有实例身份观察或 Tenant 进入 `status.phase=Binding` 后由 CEL 固定申请目标;Database 有实例身份观察、凭据位置或
绑定后固定实际 database、loginRole、来源和凭据引用,回收策略仍可修改。 绑定后固定实际 database、loginRole、来源和凭据引用,回收策略仍可修改。
读取绑定判断使用 APIReader,写入依靠 resourceVersion;watch/cache 负责触发协调。 读取绑定判断使用 APIReader,写入依靠 resourceVersion;watch/cache 负责触发协调。
绑定顺序由 application service 协调,纯资格规则在领域层;Kubernetes adapter 负责快照 绑定顺序由 application service 协调,纯资格规则在领域层;Kubernetes adapter 负责快照
+70 -19
View File
@@ -1,16 +1,16 @@
# 部署与配置 # 部署与配置
> 本页迁入作为 Database 模块的目标部署合同。Ayatori manager flags、manifests 与发布装配尚未 > 本页区分已实现的 Instance 观测、认证和凭据准备配置,与尚未接入的 PostgreSQL 供应/交付合同。
> 实现;当前行为以修订后的系统规格为准,本页不能直接用于部署。 > 完整 Database 服务仍不可部署使用;当前可执行入口见 [模块说明](README.md)。
| 项目 | 内容 | | 项目 | 内容 |
| --- | --- | | --- | --- |
| 状态 | Review | | 状态 | Review |
| 环境 | homelab Kubernetes + 外部 PostgreSQL/OpenBao | | 环境 | homelab Kubernetes + 外部 PostgreSQL/OpenBao |
| 最后更新 | 2026-09-24 | | 最后更新 | 2026-09-25 |
本文定义 v1alpha1 的运行依赖、启动顺序和部署级配置。当前 manifests 尚未实现这些 本文定义 v1alpha1 的运行依赖、启动顺序和部署级配置。Instance 观测已接入 manager;
配置,示例是后续实现合同,不可直接用于现有脚手架。 OpenBao 认证可显式启用;ESO 与完整供应装配仍是后续实现合同。
## 依赖与顺序 ## 依赖与顺序
@@ -30,21 +30,28 @@
## Controller 配置合同 ## Controller 配置合同
以下是尚待实现的部署配置合同,凭据定位随三资源 API 继续细化。controller 使用这些 CLI flags。 当前 manager 支持 `--database-secret-namespace`(默认 `POD_NAMESPACE`,为空则停用
Instance 观测)与 `--database-root-cert`(公开 PostgreSQL CA PEM 路径)。Deployment
通过 downward API 获取 namespace,Secret 权限由该 namespace 的 Role 授予。
以下表格区分已实现的认证参数与尚待实现的供应/交付参数。
必填项缺失、路径无效或 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-kv-mount` | `kv` | KV v2 mount;开发可显式用 `secret` | | `--openbao-ca-cert` | 已实现,默认系统信任根 | OpenBao 公开 CA PEM 路径 |
| `--openbao-service-account-token-path` | `/var/run/secrets/kubernetes.io/serviceaccount/token` | Kubernetes auth 使用的投射 token 文件 | | `--database-credential-mount` | 已实现,默认空 | 显式设置后启用应用凭据准备,要求已配置 OpenBao 认证 |
| `--openbao-tenant-base-path` | 默认 `postgresql-tenants` | controller 专属 mount-relative 前缀 | | `--openbao-service-account-namespace` | 已实现,启用时必填 | TokenRequest 目标 SA 的固定 namespace |
| `--openbao-service-account-name` | 已实现,启用时必填 | TokenRequest 目标 SA 名称 |
| `--openbao-token-audience` | 已实现,`openbao` | SA JWT audience,必须匹配 OpenBao role |
| `--database-credential-prefix` | 已实现,默认 `applications` | Database UID 路径的 mount-relative 前缀;不迁移已有位置 |
| `--external-secret-store-name` | 必填 | controller 创建的 ExternalSecret 固定引用 | | `--external-secret-store-name` | 必填 | controller 创建的 ExternalSecret 固定引用 |
| `--postgresql-ca-bundle-path` | PostgreSQL TLS 模式必填 | 只读 PEM trust bundle,不含私钥 | | `--database-root-cert` | 已实现 | 只读 PEM trust bundle,不含私钥;沿用 Instance 连接配置 |
| `--reconcile-timeout` | `30s` | 单轮 reconcile 中外部操作的总期限,必须大于零 | | `--reconcile-timeout` | `30s` | 单轮 reconcile 中外部操作的总期限,必须大于零 |
address 必须是绝对 `http` 或 `https` URL,不允许 userinfo、query 或 fragment,末尾 `/` address 必须是绝对 `http` 或 `https` URL,不允许 userinfo、query 或 fragment,末尾 `/`
@@ -52,9 +59,40 @@ 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 合同。
稳定位置与导入关联方式待 API 评审;consumer URL 仍使用无认证信息的 KV v2 API URL。 动态供应位置使用 `<base-path>/<Database UID>`;导入使用 Database 的显式 credentialRef,
不要求搬迁已有凭据。供应流程须先记录原 mount/path,不能在配置变化后重新推导位置。
consumer URL 仍使用无认证信息的 KV v2 API URL。
base path 必须是合法 mount-relative path,不以 `/` 开头且不包含空段、`.`、`..`、 base path 必须是合法 mount-relative path,不以 `/` 开头且不包含空段、`.`、`..`、
`data`/`metadata` API 层。ExternalSecret 固定命名为 `data`/`metadata` API 层。ExternalSecret 固定命名为
@@ -77,14 +115,27 @@ base path 必须是合法 mount-relative path,不以 `/` 开头且不包含空
- `Delete` 时禁止连接、终止目标 database session、删除已验证归属的 database/role。 - `Delete` 时禁止连接、终止目标 database session、删除已验证归属的 database/role。
部分 PostgreSQL 操作天然要求较高权限,尤其终止其他 session 和安装某些 extension。 部分 PostgreSQL 操作天然要求较高权限,尤其终止其他 session 和安装某些 extension。
应优先使用 PostgreSQL 预定义角色、受控 SECURITY DEFINER 管理函数或限定数据库的 第一版使用原生非 superuser 的 CREATEDB/CREATEROLE 方案,不引入 SECURITY DEFINER
授权;任何不得不使用 superuser 的 extension 都必须按实例单独记录,不得扩大默认 管理接口。对自行创建的 owner 显式建立 SET membership,再以 owner 管理 ACL 与扩展;
controller 权限。最终可执行 SQL grant 将随 PostgreSQL adapter 集成测试固化。 已有对象仍须逐资源核实授权,不能凭基础属性接管。需要 superuser 的扩展不能扩大 controller
权限。真实权限矩阵见 [Instance 原生管理观测](README.md#instance-原生管理观测)。
## OpenBao 与 ESO ## OpenBao 与 ESO
controller policy 仅允许在固定 tenant base path 下 create/read/update/delete KV v2 当前凭据准备只需固定前缀下 KV v2 data 的 create/read/update 权限,不需要 metadata 或
data 和 metadata,Delete 必须能永久删除全部版本及 metadata;不读取管理凭据路径。 delete 权限,也不读取管理凭据路径。示例(`secret` 为测试 mount,需替换为实际配置):
```hcl
path "secret/data/applications/*" {
capabilities = ["create", "read", "update"]
}
```
在上面的认证启动参数基础上增加 `--database-credential-mount=secret` 即启用凭据准备;
需要不同前缀时同时配置 `--database-credential-prefix` 和对应 policy。仅设置
`--openbao-address` 不会启用凭据创建。缺少认证配置或非法 mount/prefix 在启动时失败。
新字段 CRD 必须先升级再启用 controller,否则 API pruning 会使确认记录无法保存。
未来 Delete 清理才需要永久删除全部版本及 metadata 的权限,本切片不提前授予。
管理凭据由管理员维护的 ExternalSecret 同步到 controller namespace;其 ESO 身份 管理凭据由管理员维护的 ExternalSecret 同步到 controller namespace;其 ESO 身份
只读对应管理路径,不能供 Tenant 使用。租户 ESO 身份只读 tenant base path,不得 只读对应管理路径,不能供 Tenant 使用。租户 ESO 身份只读 tenant base path,不得
+10 -1
View File
@@ -5,13 +5,22 @@
## 当前绑定切片的限制 ## 当前绑定切片的限制
源码已接入绑定 controller,未接入 PostgreSQL 供应、OpenBao/ESO 交付或删除清理。 源码已接入绑定、Instance 观测与可选的 Bao 凭据准备,未接入 PostgreSQL 供应、ESO 交付或删除清理。
Bound/BindingComplete 只表示 Kubernetes 双向记录一致,Ready 仍为 False。 Bound/BindingComplete 只表示 Kubernetes 双向记录一致,Ready 仍为 False。
Tenant 删除会保留 `database.ayatori.ddupan.top/tenant-protection` 并报告 DeletionPending; Tenant 删除会保留 `database.ayatori.ddupan.top/tenant-protection` 并报告 DeletionPending;
Database 的 `database.ayatori.ddupan.top/database-protection` 也尚无清理后移除路径。 Database 的 `database.ayatori.ddupan.top/database-protection` 也尚无清理后移除路径。
这是未完成能力的明确边界,不是已经实现的 Retain/Delete 恢复逻辑。不要将此切片部署为 这是未完成能力的明确边界,不是已经实现的 Retain/Delete 恢复逻辑。不要将此切片部署为
业务 DBaaS,也不要为了消除等待状态直接移除 finalizer;后续必须补齐清理与验收。 业务 DBaaS,也不要为了消除等待状态直接移除 finalizer;后续必须补齐清理与验收。
凭据准备现在可单独启用,见[部署参数](deployment.md)。观察 Database 的
`CredentialsReady`、固定 `status.credentialRef` 和 `status.credentialVersion`,不要导出
凭据内容。Prepared 只代表 Bao 凭据可用,不代表已建库或 Tenant Ready。
`CreationStarted` 在创建结果确认前持久化;重启或失败留下该状态时报告 Conflict。
这包括“状态已写但请求还没发出”的保守停止;不要因为当前位置暂时为空就重试生成密码。
无确认版本的 Conflict 不会自行消失,后端恢复也不自动重入;应暂停该 controller、等待
在途请求结束后核实版本历史、残留与绑定,再决定清理或显式导入。不要清空整个 status
或更改固定引用来绕过保护。导入和完整人工恢复入口仍待对应切片实现。
## 日常检查 ## 日常检查
Instance 观察已实现:先确认 manager 配置了 `--database-secret-namespace` 或 `POD_NAMESPACE`, Instance 观察已实现:先确认 manager 配置了 `--database-secret-namespace` 或 `POD_NAMESPACE`,
+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 现值,结果不确定时停止并报冲突。
+6
View File
@@ -192,6 +192,12 @@ mount/base path 属部署配置,Tenant 不得自选任意路径;原按 Tenan
位置,不要求搬迁已有凭据。Released 不自动改密,管理员处理旧访问后才重新开放资源。 位置,不要求搬迁已有凭据。Released 不自动改密,管理员处理旧访问后才重新开放资源。
不得因换 Tenant、改部署参数或重新绑定就隐式搬迁凭据或改密。 不得因换 Tenant、改部署参数或重新绑定就隐式搬迁凭据或改密。
2026-09-27 确认最小凭据记录:Database `status.credentialRef` 在首次外部写入前固定
mount/path;`status.credentialVersion` 仅在成功创建并回读后保存 KV 版本。位置和确认版本
不得自动更改或清空,Conditions 只描述当前可用性。已有值但无确认版本时报告 Conflict,
不能靠读取成功认领;已确认凭据消失或最新版本不一致同样停止,等待人工核实。
部署参数变化不得搬迁旧位置;确认记录写入的 resourceVersion 冲突不能通过盲目重试覆盖。
TLS、OpenBao Kubernetes auth、controller/ESO 身份隔离、Secret 读取范围和防泄漏要求 TLS、OpenBao Kubernetes auth、controller/ESO 身份隔离、Secret 读取范围和防泄漏要求
见 [安全模型](security.md)。这些安全约束继续适用。 见 [安全模型](security.md)。这些安全约束继续适用。
+16 -2
View File
@@ -4,10 +4,13 @@ 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
k8s.io/api v0.37.0 k8s.io/api v0.37.0
k8s.io/apimachinery v0.37.0 k8s.io/apimachinery v0.37.0
k8s.io/client-go v0.37.0 k8s.io/client-go v0.37.0
sigs.k8s.io/controller-runtime v0.25.0 sigs.k8s.io/controller-runtime v0.25.0
sigs.k8s.io/yaml v1.6.0
) )
require ( require (
@@ -24,6 +27,7 @@ require (
github.com/felixge/httpsnoop v1.0.4 // indirect github.com/felixge/httpsnoop v1.0.4 // indirect
github.com/fsnotify/fsnotify v1.9.0 // indirect github.com/fsnotify/fsnotify v1.9.0 // indirect
github.com/fxamacker/cbor/v2 v2.9.1 // indirect github.com/fxamacker/cbor/v2 v2.9.1 // indirect
github.com/go-jose/go-jose/v4 v4.1.4 // indirect
github.com/go-logr/logr v1.4.3 // indirect github.com/go-logr/logr v1.4.3 // indirect
github.com/go-logr/stdr v1.2.2 // indirect github.com/go-logr/stdr v1.2.2 // indirect
github.com/go-logr/zapr v1.3.0 // indirect github.com/go-logr/zapr v1.3.0 // indirect
@@ -41,15 +45,25 @@ require (
github.com/go-openapi/swag/stringutils v0.27.1 // indirect github.com/go-openapi/swag/stringutils v0.27.1 // indirect
github.com/go-openapi/swag/typeutils v0.27.1 // indirect github.com/go-openapi/swag/typeutils v0.27.1 // indirect
github.com/go-openapi/swag/yamlutils v0.27.1 // indirect github.com/go-openapi/swag/yamlutils v0.27.1 // indirect
github.com/go-viper/mapstructure/v2 v2.5.0 // indirect
github.com/google/cel-go v0.29.2 // indirect github.com/google/cel-go v0.29.2 // indirect
github.com/google/gnostic-models v0.7.0 // indirect github.com/google/gnostic-models v0.7.0 // indirect
github.com/google/uuid v1.6.0 // indirect github.com/google/uuid v1.6.0 // indirect
github.com/grpc-ecosystem/grpc-gateway/v2 v2.29.0 // indirect github.com/grpc-ecosystem/grpc-gateway/v2 v2.29.0 // indirect
github.com/hashicorp/errwrap v1.1.0 // indirect
github.com/hashicorp/go-cleanhttp v0.5.2 // indirect
github.com/hashicorp/go-multierror v1.1.1 // indirect
github.com/hashicorp/go-retryablehttp v0.7.8 // indirect
github.com/hashicorp/go-secure-stdlib/parseutil v0.2.0 // indirect
github.com/hashicorp/go-secure-stdlib/strutil v0.1.2 // indirect
github.com/hashicorp/go-sockaddr v1.0.7 // indirect
github.com/hashicorp/hcl v1.0.1-vault-7 // indirect
github.com/inconshreveable/mousetrap v1.1.0 // indirect github.com/inconshreveable/mousetrap v1.1.0 // indirect
github.com/jackc/pgpassfile v1.0.0 // indirect github.com/jackc/pgpassfile v1.0.0 // indirect
github.com/jackc/pgservicefile v0.0.0-20240606120523-5a60cdf6a761 // indirect github.com/jackc/pgservicefile v0.0.0-20240606120523-5a60cdf6a761 // indirect
github.com/jackc/puddle/v2 v2.2.2 // indirect github.com/jackc/puddle/v2 v2.2.2 // indirect
github.com/json-iterator/go v1.1.12 // indirect github.com/json-iterator/go v1.1.12 // indirect
github.com/mitchellh/mapstructure v1.5.0 // indirect
github.com/modern-go/concurrent v0.0.0-20180306012644-bacd9c7ef1dd // indirect github.com/modern-go/concurrent v0.0.0-20180306012644-bacd9c7ef1dd // indirect
github.com/modern-go/reflect2 v1.0.3-0.20250322232337-35a7c28c31ee // indirect github.com/modern-go/reflect2 v1.0.3-0.20250322232337-35a7c28c31ee // indirect
github.com/munnerz/goautoneg v0.0.0-20191010083416-a7dc8b61c822 // indirect github.com/munnerz/goautoneg v0.0.0-20191010083416-a7dc8b61c822 // indirect
@@ -58,6 +72,7 @@ require (
github.com/prometheus/client_model v0.6.2 // indirect github.com/prometheus/client_model v0.6.2 // indirect
github.com/prometheus/common v0.70.0 // indirect github.com/prometheus/common v0.70.0 // indirect
github.com/prometheus/procfs v0.21.1 // indirect github.com/prometheus/procfs v0.21.1 // indirect
github.com/ryanuber/go-glob v1.0.0 // indirect
github.com/spf13/cobra v1.10.2 // indirect github.com/spf13/cobra v1.10.2 // indirect
github.com/spf13/pflag v1.0.10 // indirect github.com/spf13/pflag v1.0.10 // indirect
github.com/x448/float16 v0.8.4 // indirect github.com/x448/float16 v0.8.4 // indirect
@@ -75,7 +90,7 @@ require (
go.yaml.in/yaml/v2 v2.4.4 // indirect go.yaml.in/yaml/v2 v2.4.4 // indirect
go.yaml.in/yaml/v3 v3.0.5 // indirect go.yaml.in/yaml/v3 v3.0.5 // indirect
golang.org/x/exp v0.0.0-20260410095643-746e56fc9e2f // indirect golang.org/x/exp v0.0.0-20260410095643-746e56fc9e2f // indirect
golang.org/x/net v0.57.0 // indirect golang.org/x/net v0.58.0 // indirect
golang.org/x/oauth2 v0.36.0 // indirect golang.org/x/oauth2 v0.36.0 // indirect
golang.org/x/sync v0.22.0 // indirect golang.org/x/sync v0.22.0 // indirect
golang.org/x/sys v0.47.0 // indirect golang.org/x/sys v0.47.0 // indirect
@@ -100,5 +115,4 @@ require (
sigs.k8s.io/json v0.0.0-20250730193827-2d320260d730 // indirect sigs.k8s.io/json v0.0.0-20250730193827-2d320260d730 // indirect
sigs.k8s.io/randfill v1.0.0 // indirect sigs.k8s.io/randfill v1.0.0 // indirect
sigs.k8s.io/structured-merge-diff/v6 v6.4.2 // indirect sigs.k8s.io/structured-merge-diff/v6 v6.4.2 // indirect
sigs.k8s.io/yaml v1.6.0 // indirect
) )
+43 -4
View File
@@ -23,12 +23,16 @@ github.com/evanphx/json-patch v0.5.2 h1:xVCHIVMUu1wtM/VkR9jVZ45N3FhZfYMMYGorLCR8
github.com/evanphx/json-patch v0.5.2/go.mod h1:ZWS5hhDbVDyob71nXKNL0+PWn6ToqBHMikGIFbs31qQ= github.com/evanphx/json-patch v0.5.2/go.mod h1:ZWS5hhDbVDyob71nXKNL0+PWn6ToqBHMikGIFbs31qQ=
github.com/evanphx/json-patch/v5 v5.9.11 h1:/8HVnzMq13/3x9TPvjG08wUGqBTmZBsCWzjTM0wiaDU= github.com/evanphx/json-patch/v5 v5.9.11 h1:/8HVnzMq13/3x9TPvjG08wUGqBTmZBsCWzjTM0wiaDU=
github.com/evanphx/json-patch/v5 v5.9.11/go.mod h1:3j+LviiESTElxA4p3EMKAB9HXj3/XEtnUf6OZxqIQTM= github.com/evanphx/json-patch/v5 v5.9.11/go.mod h1:3j+LviiESTElxA4p3EMKAB9HXj3/XEtnUf6OZxqIQTM=
github.com/fatih/color v1.19.0 h1:Zp3PiM21/9Ld6FzSKyL5c/BULoe/ONr9KlbYVOfG8+w=
github.com/fatih/color v1.19.0/go.mod h1:zNk67I0ZUT1bEGsSGyCZYZNrHuTkJJB+r6Q9VuMi0LE=
github.com/felixge/httpsnoop v1.0.4 h1:NFTV2Zj1bL4mc9sqWACXbQFVBBg2W3GPvqp8/ESS2Wg= github.com/felixge/httpsnoop v1.0.4 h1:NFTV2Zj1bL4mc9sqWACXbQFVBBg2W3GPvqp8/ESS2Wg=
github.com/felixge/httpsnoop v1.0.4/go.mod h1:m8KPJKqk1gH5J9DgRY2ASl2lWCfGKXixSwevea8zH2U= github.com/felixge/httpsnoop v1.0.4/go.mod h1:m8KPJKqk1gH5J9DgRY2ASl2lWCfGKXixSwevea8zH2U=
github.com/fsnotify/fsnotify v1.9.0 h1:2Ml+OJNzbYCTzsxtv8vKSFD9PbJjmhYF14k/jKC7S9k= github.com/fsnotify/fsnotify v1.9.0 h1:2Ml+OJNzbYCTzsxtv8vKSFD9PbJjmhYF14k/jKC7S9k=
github.com/fsnotify/fsnotify v1.9.0/go.mod h1:8jBTzvmWwFyi3Pb8djgCCO5IBqzKJ/Jwo8TRcHyHii0= github.com/fsnotify/fsnotify v1.9.0/go.mod h1:8jBTzvmWwFyi3Pb8djgCCO5IBqzKJ/Jwo8TRcHyHii0=
github.com/fxamacker/cbor/v2 v2.9.1 h1:2rWm8B193Ll4VdjsJY28jxs70IdDsHRWgQYAI80+rMQ= github.com/fxamacker/cbor/v2 v2.9.1 h1:2rWm8B193Ll4VdjsJY28jxs70IdDsHRWgQYAI80+rMQ=
github.com/fxamacker/cbor/v2 v2.9.1/go.mod h1:vM4b+DJCtHn+zz7h3FFp/hDAI9WNWCsZj23V5ytsSxQ= github.com/fxamacker/cbor/v2 v2.9.1/go.mod h1:vM4b+DJCtHn+zz7h3FFp/hDAI9WNWCsZj23V5ytsSxQ=
github.com/go-jose/go-jose/v4 v4.1.4 h1:moDMcTHmvE6Groj34emNPLs/qtYXRVcd6S7NHbHz3kA=
github.com/go-jose/go-jose/v4 v4.1.4/go.mod h1:x4oUasVrzR7071A4TnHLGSPpNOm2a21K9Kf04k1rs08=
github.com/go-logr/logr v1.2.2/go.mod h1:jdQByPbusPIv2/zmleS9BjJVeZ6kBagPoEUsqbVz/1A= github.com/go-logr/logr v1.2.2/go.mod h1:jdQByPbusPIv2/zmleS9BjJVeZ6kBagPoEUsqbVz/1A=
github.com/go-logr/logr v1.4.3 h1:CjnDlHq8ikf6E492q6eKboGOC0T8CDaOvkHCIg8idEI= github.com/go-logr/logr v1.4.3 h1:CjnDlHq8ikf6E492q6eKboGOC0T8CDaOvkHCIg8idEI=
github.com/go-logr/logr v1.4.3/go.mod h1:9T104GzyrTigFIr8wt5mBrctHMim0Nb2HLGrmQ40KvY= github.com/go-logr/logr v1.4.3/go.mod h1:9T104GzyrTigFIr8wt5mBrctHMim0Nb2HLGrmQ40KvY=
@@ -72,6 +76,10 @@ github.com/go-openapi/testify/v2 v2.6.0 h1:5PKH2HE7YJ/LuRPQGvSxBRlFXNQhSetBLlGAg
github.com/go-openapi/testify/v2 v2.6.0/go.mod h1:SgsVHtfooshd0tublTtJ50FPKhujf47YRqauXXOUxfw= github.com/go-openapi/testify/v2 v2.6.0/go.mod h1:SgsVHtfooshd0tublTtJ50FPKhujf47YRqauXXOUxfw=
github.com/go-task/slim-sprig/v3 v3.0.0 h1:sUs3vkvUymDpBKi3qH1YSqBQk9+9D/8M2mN1vB6EwHI= github.com/go-task/slim-sprig/v3 v3.0.0 h1:sUs3vkvUymDpBKi3qH1YSqBQk9+9D/8M2mN1vB6EwHI=
github.com/go-task/slim-sprig/v3 v3.0.0/go.mod h1:W848ghGpv3Qj3dhTPRyJypKRiqCdHZiAzKg9hl15HA8= github.com/go-task/slim-sprig/v3 v3.0.0/go.mod h1:W848ghGpv3Qj3dhTPRyJypKRiqCdHZiAzKg9hl15HA8=
github.com/go-test/deep v1.1.1 h1:0r/53hagsehfO4bzD2Pgr/+RgHqhmf+k1Bpse2cTu1U=
github.com/go-test/deep v1.1.1/go.mod h1:5C2ZWiW0ErCdrYzpqxLbTX7MG14M9iiw8DgHncVwcsE=
github.com/go-viper/mapstructure/v2 v2.5.0 h1:vM5IJoUAy3d7zRSVtIwQgBj7BiWtMPfmPEgAXnvj1Ro=
github.com/go-viper/mapstructure/v2 v2.5.0/go.mod h1:oJDH3BJKyqBA2TXFhDsKDGDTlndYOZ6rGS0BRZIxGhM=
github.com/golang/protobuf v1.5.4 h1:i7eJL8qZTpSEXOPTxNKhASYpMn+8e5Q6AdndVa1dWek= github.com/golang/protobuf v1.5.4 h1:i7eJL8qZTpSEXOPTxNKhASYpMn+8e5Q6AdndVa1dWek=
github.com/golang/protobuf v1.5.4/go.mod h1:lnTiLA8Wa4RWRcIUkrtSVa5nRhsEGBg48fD6rSs7xps= github.com/golang/protobuf v1.5.4/go.mod h1:lnTiLA8Wa4RWRcIUkrtSVa5nRhsEGBg48fD6rSs7xps=
github.com/google/cel-go v0.29.2 h1:ZtDxkeiMmz0mxbKDYiNkE5Lk7V5edMRcaaDf2jX002k= github.com/google/cel-go v0.29.2 h1:ZtDxkeiMmz0mxbKDYiNkE5Lk7V5edMRcaaDf2jX002k=
@@ -89,6 +97,25 @@ github.com/google/uuid v1.6.0 h1:NIvaJDMOsjHA8n1jAhLSgzrAzy1Hgr+hNrb57e+94F0=
github.com/google/uuid v1.6.0/go.mod h1:TIyPZe4MgqvfeYDBFedMoGGpEw/LqOeaOT+nhxU+yHo= github.com/google/uuid v1.6.0/go.mod h1:TIyPZe4MgqvfeYDBFedMoGGpEw/LqOeaOT+nhxU+yHo=
github.com/grpc-ecosystem/grpc-gateway/v2 v2.29.0 h1:5VipnvEpbqr2gA2VbM+nYVbkIF28c5ZQfqCBQ5g2xfk= github.com/grpc-ecosystem/grpc-gateway/v2 v2.29.0 h1:5VipnvEpbqr2gA2VbM+nYVbkIF28c5ZQfqCBQ5g2xfk=
github.com/grpc-ecosystem/grpc-gateway/v2 v2.29.0/go.mod h1:Hyl3n6Twe1hvtd9XUXDec4pTvgMSEixRuQKPTMH2bNs= github.com/grpc-ecosystem/grpc-gateway/v2 v2.29.0/go.mod h1:Hyl3n6Twe1hvtd9XUXDec4pTvgMSEixRuQKPTMH2bNs=
github.com/hashicorp/errwrap v1.0.0/go.mod h1:YH+1FKiLXxHSkmPseP+kNlulaMuP3n2brvKWEqk/Jc4=
github.com/hashicorp/errwrap v1.1.0 h1:OxrOeh75EUXMY8TBjag2fzXGZ40LB6IKw45YeGUDY2I=
github.com/hashicorp/errwrap v1.1.0/go.mod h1:YH+1FKiLXxHSkmPseP+kNlulaMuP3n2brvKWEqk/Jc4=
github.com/hashicorp/go-cleanhttp v0.5.2 h1:035FKYIWjmULyFRBKPs8TBQoi0x6d9G4xc9neXJWAZQ=
github.com/hashicorp/go-cleanhttp v0.5.2/go.mod h1:kO/YDlP8L1346E6Sodw+PrpBSV4/SoxCXGY6BqNFT48=
github.com/hashicorp/go-hclog v1.6.3 h1:Qr2kF+eVWjTiYmU7Y31tYlP1h0q/X3Nl3tPGdaB11/k=
github.com/hashicorp/go-hclog v1.6.3/go.mod h1:W4Qnvbt70Wk/zYJryRzDRU/4r0kIg0PVHBcfoyhpF5M=
github.com/hashicorp/go-multierror v1.1.1 h1:H5DkEtf6CXdFp0N0Em5UCwQpXMWke8IA0+lD48awMYo=
github.com/hashicorp/go-multierror v1.1.1/go.mod h1:iw975J/qwKPdAO1clOe2L8331t/9/fmwbPZ6JB6eMoM=
github.com/hashicorp/go-retryablehttp v0.7.8 h1:ylXZWnqa7Lhqpk0L1P1LzDtGcCR0rPVUrx/c8Unxc48=
github.com/hashicorp/go-retryablehttp v0.7.8/go.mod h1:rjiScheydd+CxvumBsIrFKlx3iS0jrZ7LvzFGFmuKbw=
github.com/hashicorp/go-secure-stdlib/parseutil v0.2.0 h1:U+kC2dOhMFQctRfhK0gRctKAPTloZdMU5ZJxaesJ/VM=
github.com/hashicorp/go-secure-stdlib/parseutil v0.2.0/go.mod h1:Ll013mhdmsVDuoIXVfBtvgGJsXDYkTw1kooNcoCXuE0=
github.com/hashicorp/go-secure-stdlib/strutil v0.1.2 h1:kes8mmyCpxJsI7FTwtzRqEy9CdjCtrXrXGuOpxEA7Ts=
github.com/hashicorp/go-secure-stdlib/strutil v0.1.2/go.mod h1:Gou2R9+il93BqX25LAKCLuM+y9U2T4hlwvT1yprcna4=
github.com/hashicorp/go-sockaddr v1.0.7 h1:G+pTkSO01HpR5qCxg7lxfsFEZaG+C0VssTy/9dbT+Fw=
github.com/hashicorp/go-sockaddr v1.0.7/go.mod h1:FZQbEYa1pxkQ7WLpyXJ6cbjpT8q0YgQaK/JakXqGyWw=
github.com/hashicorp/hcl v1.0.1-vault-7 h1:ag5OxFVy3QYTFTJODRzTKVZ6xvdfLLCA1cy/Y6xGI0I=
github.com/hashicorp/hcl v1.0.1-vault-7/go.mod h1:XYhtn6ijBSAj6n4YqAaf7RBPS4I06AItNorpy+MoQNM=
github.com/inconshreveable/mousetrap v1.1.0 h1:wN+x4NVGpMsO7ErUn/mUI3vEoE6Jt13X2s0bqwp9tc8= github.com/inconshreveable/mousetrap v1.1.0 h1:wN+x4NVGpMsO7ErUn/mUI3vEoE6Jt13X2s0bqwp9tc8=
github.com/inconshreveable/mousetrap v1.1.0/go.mod h1:vpF70FUmC8bwa3OWnCshd2FqLfsEA9PFc4w1p2J65bw= github.com/inconshreveable/mousetrap v1.1.0/go.mod h1:vpF70FUmC8bwa3OWnCshd2FqLfsEA9PFc4w1p2J65bw=
github.com/jackc/pgpassfile v1.0.0 h1:/6Hmqy13Ss2zCq62VdNG8tM1wchn8zjSGOBJ6icpsIM= github.com/jackc/pgpassfile v1.0.0 h1:/6Hmqy13Ss2zCq62VdNG8tM1wchn8zjSGOBJ6icpsIM=
@@ -105,6 +132,12 @@ github.com/klauspost/compress v1.19.0 h1:sXLILfc9jV2QYWkzFOPWStmcUVH2RHEB1JCdY2o
github.com/klauspost/compress v1.19.0/go.mod h1:cwPg85FWrGar70rWktvGQj8/hthj3wpl0PGDogxkrSQ= github.com/klauspost/compress v1.19.0/go.mod h1:cwPg85FWrGar70rWktvGQj8/hthj3wpl0PGDogxkrSQ=
github.com/kylelemons/godebug v1.1.0 h1:RPNrshWIDI6G2gRW9EHilWtl7Z6Sb1BR0xunSBf0SNc= github.com/kylelemons/godebug v1.1.0 h1:RPNrshWIDI6G2gRW9EHilWtl7Z6Sb1BR0xunSBf0SNc=
github.com/kylelemons/godebug v1.1.0/go.mod h1:9/0rRGxNHcop5bhtWyNeEfOS8JIWk580+fNqagV/RAw= github.com/kylelemons/godebug v1.1.0/go.mod h1:9/0rRGxNHcop5bhtWyNeEfOS8JIWk580+fNqagV/RAw=
github.com/mattn/go-colorable v0.1.15 h1:+u9SLTRGnXv73cEsnsmoZBom+dMU88B2M0aDcWy0/jY=
github.com/mattn/go-colorable v0.1.15/go.mod h1:6LmQG8QLFO4G5z1gPvYEzlUgJ2wF+stgPZH1UqBm1s8=
github.com/mattn/go-isatty v0.0.24 h1:tGZZoVgT/KiqK1c8ocVLeDS8BSWMRd47J3Lbz7vsReI=
github.com/mattn/go-isatty v0.0.24/go.mod h1:nMCL3Zebbrt45jsMDgnfIwz6ydEQApk5oEI3HqDio6A=
github.com/mitchellh/mapstructure v1.5.0 h1:jeMsZIYE/09sWLaz43PL7Gy6RuMjD2eJVyuac5Z2hdY=
github.com/mitchellh/mapstructure v1.5.0/go.mod h1:bFUtVrKA4DC2yAKiSyO/QUcy7e+RRV2QTWOzhPopBRo=
github.com/modern-go/concurrent v0.0.0-20180228061459-e0a39a4cb421/go.mod h1:6dJC0mAP4ikYIbvyc7fijjWJddQyLn8Ig3JB5CqoB9Q= github.com/modern-go/concurrent v0.0.0-20180228061459-e0a39a4cb421/go.mod h1:6dJC0mAP4ikYIbvyc7fijjWJddQyLn8Ig3JB5CqoB9Q=
github.com/modern-go/concurrent v0.0.0-20180306012644-bacd9c7ef1dd h1:TRLaZ9cD/w8PVh93nsPXa1VrQ6jlwL5oN8l14QlcNfg= github.com/modern-go/concurrent v0.0.0-20180306012644-bacd9c7ef1dd h1:TRLaZ9cD/w8PVh93nsPXa1VrQ6jlwL5oN8l14QlcNfg=
github.com/modern-go/concurrent v0.0.0-20180306012644-bacd9c7ef1dd/go.mod h1:6dJC0mAP4ikYIbvyc7fijjWJddQyLn8Ig3JB5CqoB9Q= github.com/modern-go/concurrent v0.0.0-20180306012644-bacd9c7ef1dd/go.mod h1:6dJC0mAP4ikYIbvyc7fijjWJddQyLn8Ig3JB5CqoB9Q=
@@ -117,6 +150,10 @@ 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/go.mod h1:uXbMoyH2pjSvNyTepinUvLde8pOJB82EuhUCfOKnKbo=
github.com/pkg/errors v0.9.1 h1:FEBLx1zS214owpjy7qsBeixbURkuhQAwrK5UwLGTwt4= github.com/pkg/errors v0.9.1 h1:FEBLx1zS214owpjy7qsBeixbURkuhQAwrK5UwLGTwt4=
github.com/pkg/errors v0.9.1/go.mod h1:bwawxfHBFNV+L2hUp1rHADufV3IMtnDRdf1r5NINEl0= github.com/pkg/errors v0.9.1/go.mod h1:bwawxfHBFNV+L2hUp1rHADufV3IMtnDRdf1r5NINEl0=
github.com/pmezard/go-difflib v1.0.0/go.mod h1:iKH77koFhYxTK1pcRnkKkqfTogsbg7gZNVY4sRDYZ/4= github.com/pmezard/go-difflib v1.0.0/go.mod h1:iKH77koFhYxTK1pcRnkKkqfTogsbg7gZNVY4sRDYZ/4=
@@ -131,6 +168,8 @@ github.com/prometheus/common v0.70.0/go.mod h1:S/SFasQmgGiYH6C81LKCtYa8QACgthGg5
github.com/prometheus/procfs v0.21.1 h1:GljZCt+zSTS+NZq88cyQ1LjZ+RCHp3uVuabBWA5+OJI= github.com/prometheus/procfs v0.21.1 h1:GljZCt+zSTS+NZq88cyQ1LjZ+RCHp3uVuabBWA5+OJI=
github.com/prometheus/procfs v0.21.1/go.mod h1:aB55Cww9pdSJVHk0hUf0inxWyyjPogFIjmHKYgMKmtY= github.com/prometheus/procfs v0.21.1/go.mod h1:aB55Cww9pdSJVHk0hUf0inxWyyjPogFIjmHKYgMKmtY=
github.com/russross/blackfriday/v2 v2.1.0/go.mod h1:+Rmxgy9KzJVeS9/2gXHxylqXiyQDYRxCVz55jmeOWTM= github.com/russross/blackfriday/v2 v2.1.0/go.mod h1:+Rmxgy9KzJVeS9/2gXHxylqXiyQDYRxCVz55jmeOWTM=
github.com/ryanuber/go-glob v1.0.0 h1:iQh3xXAumdQ+4Ufa5b25cRpC5TYKlno6hsv6Cb3pkBk=
github.com/ryanuber/go-glob v1.0.0/go.mod h1:807d1WSdnB0XRJzKNil9Om6lcp/3a0v4qIHxIXzX/Yc=
github.com/spf13/cobra v1.10.2 h1:DMTTonx5m65Ic0GOoRY2c16WCbHxOOw6xxezuLaBpcU= github.com/spf13/cobra v1.10.2 h1:DMTTonx5m65Ic0GOoRY2c16WCbHxOOw6xxezuLaBpcU=
github.com/spf13/cobra v1.10.2/go.mod h1:7C1pvHqHw5A4vrJfjNwvOdzYu0Gml16OCs2GRiTUUS4= github.com/spf13/cobra v1.10.2/go.mod h1:7C1pvHqHw5A4vrJfjNwvOdzYu0Gml16OCs2GRiTUUS4=
github.com/spf13/pflag v1.0.9/go.mod h1:McXfInJRrz4CZXVZOBLb0bTZqETkiAhM9Iw0y3An2Bg= github.com/spf13/pflag v1.0.9/go.mod h1:McXfInJRrz4CZXVZOBLb0bTZqETkiAhM9Iw0y3An2Bg=
@@ -141,8 +180,8 @@ github.com/stretchr/objx v0.5.3 h1:jmXUvGomnU1o3W/V5h2VEradbpJDwGrzugQQvL0POH4=
github.com/stretchr/objx v0.5.3/go.mod h1:rDQraq+vQZU7Fde9LOZLr8Tax6zZvy4kuNKF+QYS+U0= github.com/stretchr/objx v0.5.3/go.mod h1:rDQraq+vQZU7Fde9LOZLr8Tax6zZvy4kuNKF+QYS+U0=
github.com/stretchr/testify v1.3.0/go.mod h1:M5WIy9Dh21IEIfnGCwXGc5bZfKNJtfHm1UVUgZn+9EI= github.com/stretchr/testify v1.3.0/go.mod h1:M5WIy9Dh21IEIfnGCwXGc5bZfKNJtfHm1UVUgZn+9EI=
github.com/stretchr/testify v1.7.0/go.mod h1:6Fq8oRcR53rry900zMqJjRRixrwX3KX962/h/Wwjteg= github.com/stretchr/testify v1.7.0/go.mod h1:6Fq8oRcR53rry900zMqJjRRixrwX3KX962/h/Wwjteg=
github.com/stretchr/testify v1.11.1 h1:7s2iGBzp5EwR7/aIZr8ao5+dra3wiQyKjjFuvgVKu7U= github.com/stretchr/testify v1.12.1 h1:EuwCh5fleGS7H32xRwO3wRGT7DxrDhLAT6FF8MpWDWE=
github.com/stretchr/testify v1.11.1/go.mod h1:wZwfW3scLgRK+23gO65QZefKpKQRnfz6sD981Nm4B6U= github.com/stretchr/testify v1.12.1/go.mod h1:MDEgiDPPsNp5cuIrHPPCyornHKgEVbtFUmoNlxoYthg=
github.com/x448/float16 v0.8.4 h1:qLwI1I70+NjRFUR3zs1JPUCgaCXSh3SW62uAKT1mSBM= github.com/x448/float16 v0.8.4 h1:qLwI1I70+NjRFUR3zs1JPUCgaCXSh3SW62uAKT1mSBM=
github.com/x448/float16 v0.8.4/go.mod h1:14CWIYCyZA/cWjXOioeEpHeN/83MdbZDRQHoFcYsOfg= github.com/x448/float16 v0.8.4/go.mod h1:14CWIYCyZA/cWjXOioeEpHeN/83MdbZDRQHoFcYsOfg=
go.opentelemetry.io/auto/sdk v1.2.1 h1:jXsnJ4Lmnqd11kwkBV2LgLoFMZKizbCi5fNZ/ipaZ64= go.opentelemetry.io/auto/sdk v1.2.1 h1:jXsnJ4Lmnqd11kwkBV2LgLoFMZKizbCi5fNZ/ipaZ64=
@@ -180,8 +219,8 @@ golang.org/x/exp v0.0.0-20260410095643-746e56fc9e2f h1:W3F4c+6OLc6H2lb//N1q4WpJk
golang.org/x/exp v0.0.0-20260410095643-746e56fc9e2f/go.mod h1:J1xhfL/vlindoeF/aINzNzt2Bket5bjo9sdOYzOsU80= golang.org/x/exp v0.0.0-20260410095643-746e56fc9e2f/go.mod h1:J1xhfL/vlindoeF/aINzNzt2Bket5bjo9sdOYzOsU80=
golang.org/x/mod v0.38.0 h1:MECBjubtXD7yj4HrhIUcywNaGeNVUdfVnxmPajOk4yk= golang.org/x/mod v0.38.0 h1:MECBjubtXD7yj4HrhIUcywNaGeNVUdfVnxmPajOk4yk=
golang.org/x/mod v0.38.0/go.mod h1:V6Xz0pq8TQ3dGqVQ1FVHuelZpAL0uNhSkk9ogYP3c40= golang.org/x/mod v0.38.0/go.mod h1:V6Xz0pq8TQ3dGqVQ1FVHuelZpAL0uNhSkk9ogYP3c40=
golang.org/x/net v0.57.0 h1:K5+3DljvIuDG9/Jv9rvyMywYNFCQ9RSUY6OOTTkT+tE= golang.org/x/net v0.58.0 h1:ynWG7rqYi4ccpTEuPZ2QGWHktVEM9DMCj9yzDE0Q7To=
golang.org/x/net v0.57.0/go.mod h1:KpXc8iv+r3XplLAG/f7Jsf9RPszJzdR0f58q9vGOuEU= golang.org/x/net v0.58.0/go.mod h1:YwCddHnFlT7eLQqVprV19OnhLGtc5xOKgE0RyqgfWAU=
golang.org/x/oauth2 v0.36.0 h1:peZ/1z27fi9hUOFCAZaHyrpWG5lwe0RJEEEeH0ThlIs= golang.org/x/oauth2 v0.36.0 h1:peZ/1z27fi9hUOFCAZaHyrpWG5lwe0RJEEEeH0ThlIs=
golang.org/x/oauth2 v0.36.0/go.mod h1:YDBUJMTkDnJS+A4BP4eZBjCqtokkg1hODuPjwiGPO7Q= golang.org/x/oauth2 v0.36.0/go.mod h1:YDBUJMTkDnJS+A4BP4eZBjCqtokkg1hODuPjwiGPO7Q=
golang.org/x/sync v0.22.0 h1:SZjpbeLmrCk4xhRSZFNZW5gFUeCeFgjekvI/+gfScek= golang.org/x/sync v0.22.0 h1:SZjpbeLmrCk4xhRSZFNZW5gFUeCeFgjekvI/+gfScek=
+90
View File
@@ -0,0 +1,90 @@
package bootstrap
import (
"context"
"flag"
"fmt"
"os"
bao "github.com/openbao/openbao/api/v2"
"git.ddupan.top/panxiao81/ayatori/internal/database/adapter/kubernetes"
databasebao "git.ddupan.top/panxiao81/ayatori/internal/database/adapter/openbao"
"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
credentialMount string
credentialPrefix 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 路径")
flags.StringVar(&o.credentialMount, "database-credential-mount", "", "应用凭据 KV v2 mount;为空时不启用凭据准备")
flags.StringVar(&o.credentialPrefix, "database-credential-prefix", "applications", "应用凭据路径前缀;已有固定位置不随配置变化迁移")
}
func (o databaseOptions) configureManager(options *ctrl.Options) {
if o.secretNamespace != "" {
options.Cache = databasecontroller.InstanceCacheOptions(o.secretNamespace)
}
}
// registerDatabaseControllers 封装 Database 的内部装配,并返回在 manager 停止后执行的清理。
func registerDatabaseControllers(ctx context.Context, manager ctrl.Manager, options databaseOptions, baoClient *bao.Client) (func(), error) {
closeDatabaseConnections := 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)
}
closeDatabaseConnections = service.Close
}
if err := wireBindingController(manager.GetClient(), manager.GetAPIReader()).SetupWithManager(ctx, manager); err != nil {
closeDatabaseConnections()
return nil, fmt.Errorf("set up Database binding controller: %w", err)
}
if err := setupCredentialPreparation(manager, options, baoClient); err != nil {
closeDatabaseConnections()
return nil, fmt.Errorf("set up Database credential preparation: %w", err)
}
return closeDatabaseConnections, nil
}
func setupCredentialPreparation(manager ctrl.Manager, options databaseOptions, baoClient *bao.Client) error {
if options.credentialMount == "" {
return nil
}
if baoClient == nil {
return fmt.Errorf("database credential preparation requires OpenBao authentication configuration")
}
store, err := databasebao.NewCredentials(baoClient, options.credentialMount, options.credentialPrefix)
if err != nil {
return err
}
return wireCredentialController(manager.GetClient(), manager.GetAPIReader(), store).SetupWithManager(manager)
}
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 := wireInstanceController(manager.GetClient(), manager.GetAPIReader(), service, namespace)
if err := reconciler.SetupWithManager(manager); err != nil {
service.Close()
return nil, err
}
return service, nil
}
+28
View File
@@ -0,0 +1,28 @@
package bootstrap
import (
"git.ddupan.top/panxiao81/ayatori/internal/database/adapter/kubernetes"
"git.ddupan.top/panxiao81/ayatori/internal/database/application"
databasecontroller "git.ddupan.top/panxiao81/ayatori/internal/database/controller"
"sigs.k8s.io/controller-runtime/pkg/client"
)
// 本文件是 Database 的显式依赖图:adapter → 用例 → controller。
// 构造只在启动时执行;用例和 controller 都不持有容器或动态查找依赖。
func wireBindingController(writer client.Client, reader client.Reader) *databasecontroller.BindingReconciler {
resources := &kubernetes.BindingResources{Client: writer, Reader: reader}
service := &application.BindingService{Resources: resources}
return databasecontroller.NewBindingReconciler(writer, service, resources)
}
func wireCredentialController(writer client.Client, reader client.Reader, store application.CredentialStore) *databasecontroller.CredentialReconciler {
resources := &kubernetes.CredentialResources{Client: writer, Reader: reader}
service := &application.CredentialPreparation{Resources: resources, Store: store}
return databasecontroller.NewCredentialReconciler(writer, service)
}
func wireInstanceController(writer client.Client, reader client.Reader, observer application.InstanceObserver, namespace string) *databasecontroller.InstanceReconciler {
resources := &kubernetes.InstanceResources{Client: writer, Reader: reader}
service := &application.InstanceReconciliation{Resources: resources, Observer: observer}
return databasecontroller.NewInstanceReconciler(writer, service, resources, namespace)
}
@@ -0,0 +1,74 @@
package bootstrap
import (
"go/parser"
"go/token"
"io/fs"
"path/filepath"
"strconv"
"strings"
"testing"
"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"
"sigs.k8s.io/controller-runtime/pkg/client/fake"
)
const controllerLayer = "controller"
func TestDatabaseExplicitWiring(t *testing.T) {
writer := fake.NewClientBuilder().Build()
directReader := fake.NewClientBuilder().Build()
binder := wireBindingController(writer, directReader)
bindingResources, ok := binder.Service.Resources.(*kubernetes.BindingResources)
if !ok || binder.Client != writer || bindingResources.Reader != directReader || bindingResources.Client != writer || binder.Presenter != bindingResources {
t.Fatal("绑定用例没有共享显式注入的 writer、直连 reader 与 presenter")
}
credentials, err := kubernetes.NewSecretCredentials(directReader, "wiring-tests")
if err != nil {
t.Fatal(err)
}
observer, err := application.NewInstanceService(credentials, postgresql.Connector{})
if err != nil {
t.Fatal(err)
}
t.Cleanup(observer.Close)
reconciler := wireInstanceController(writer, directReader, observer, "wiring-tests")
resources, ok := reconciler.Service.Resources.(*kubernetes.InstanceResources)
if !ok || resources.Reader != directReader || resources.Client != writer || reconciler.Service.Observer != observer || reconciler.Presenter != resources {
t.Fatal("Instance 的服务或读取边界未按依赖图注入")
}
}
// 防止领域重新依赖用例/存储,也防止 controller 再次私自构造具体 adapter。
func TestDatabaseLayerImports(t *testing.T) {
for _, layer := range []string{"domain", controllerLayer} {
err := filepath.WalkDir(filepath.Join("../database", layer), func(path string, entry fs.DirEntry, err error) error {
if err != nil {
return err
}
if entry.IsDir() || !strings.HasSuffix(path, ".go") || strings.HasSuffix(path, "_test.go") {
return nil
}
file, err := parser.ParseFile(token.NewFileSet(), path, nil, parser.ImportsOnly)
if err != nil {
return err
}
for _, dependency := range file.Imports {
name, err := strconv.Unquote(dependency.Path.Value)
if err != nil {
return err
}
if strings.Contains(name, "/database/adapter/") || (layer == "domain" &&
(strings.Contains(name, "/database/application") || strings.Contains(name, "k8s.io/") || strings.Contains(name, "/internal/infra/"))) {
t.Errorf("%s 不得导入 %s", path, name)
}
}
return nil
})
if err != nil {
t.Fatal(err)
}
}
}
+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,89 @@
//go:build integration
package bootstrap
import (
"context"
"os"
"path/filepath"
"testing"
"time"
bao "github.com/openbao/openbao/api/v2"
"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)
}
baoClient, err := setupOpenBaoAuthentication(manager, options.openBao)
if err != nil {
t.Fatal(err)
}
ctx, cancel := context.WithTimeout(t.Context(), 30*time.Second)
defer cancel()
closeDatabaseConnections, err := registerDatabaseControllers(ctx, manager, options.database, baoClient)
if err != nil {
t.Fatal(err)
}
defer closeDatabaseConnections()
// 空 API 中没有供应目标;此处验证启用路径确实注册 controller,不访问外部 Bao。
fixtureConfig := bao.NewConfig()
fixtureConfig.Address = "http://127.0.0.1:1"
fixtureClient, err := bao.NewClient(fixtureConfig)
if err != nil {
t.Fatal(err)
}
options.database.credentialMount = "secret"
if err := setupCredentialPreparation(manager, options.database, fixtureClient); err != nil {
t.Fatal(err)
}
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 closeDatabaseConnections")
}
}()
if !manager.GetCache().WaitForCacheSync(ctx) {
t.Fatal("assembled controller cache did not synchronize")
}
}
+78
View File
@@ -0,0 +1,78 @@
/*
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"
bao "github.com/openbao/openbao/api/v2"
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) (*bao.Client, error) {
if options.address == "" {
return nil, nil
}
client, err := openbao.NewClient(options.address, options.caCert)
if err != nil {
return nil, err
}
// 复用 manager 已装配的 Kubernetes client,不重复加载配置或创建客户端。
session, err := openbao.NewKubernetesSession(
client, manager.GetClient(), options.mount, options.role, options.identity,
)
if err != nil {
return nil, err
}
if err := manager.Add(session); err != nil {
return nil, err
}
err = manager.AddReadyzCheck("openbao-auth", func(_ *http.Request) error {
if !session.Ready() {
return errors.New("OpenBao Kubernetes authentication unavailable")
}
return nil
})
if err != nil {
return nil, err
}
return client, 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")
}
+147
View File
@@ -0,0 +1,147 @@
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.database.credentialMount != "" || 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 TestCredentialPreparationOptions(t *testing.T) {
options := parseTestOptions(t, "--database-credential-mount=applications-kv", "--database-credential-prefix=database")
if options.database.credentialMount != "applications-kv" || options.database.credentialPrefix != "database" {
t.Fatal("凭据准备参数未传入领域装配")
}
if err := setupCredentialPreparation(nil, options.database, nil); err == nil {
t.Fatal("启用凭据准备必须有显式配置的认证 client")
}
if err := setupCredentialPreparation(nil, databaseOptions{}, nil); err != nil {
t.Fatal("默认停用凭据准备不应要求后端")
}
}
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")
}
}
}
+42
View File
@@ -0,0 +1,42 @@
// 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
}
baoClient, err := setupOpenBaoAuthentication(manager, options.openBao)
if err != nil {
return fmt.Errorf("set up OpenBao authentication: %w", err)
}
closeDatabaseConnections, err := registerDatabaseControllers(context.Background(), manager, options.database, baoClient)
if err != nil {
return err
}
// manager 的 worker 完全停止后才释放组件持有的资源。
defer closeDatabaseConnections()
setupLog.Info("Starting manager")
if err := manager.Start(ctrl.SetupSignalHandler()); err != nil {
return fmt.Errorf("run controller manager: %w", err)
}
return nil
}
@@ -51,11 +51,3 @@ func tenantReference(tenant binding.TenantIdentity) *databasev1alpha1.TenantRefe
Namespace: tenant.Namespace, Name: databasev1alpha1.ObjectName(tenant.Name), UID: types.UID(tenant.UID), Namespace: tenant.Namespace, Name: databasev1alpha1.ObjectName(tenant.Name), UID: types.UID(tenant.UID),
} }
} }
// BindingTargetName 供 informer 索引使用;不把无效请求丢出事件映射。
func BindingTargetName(tenant *databasev1alpha1.PostgreSQLTenant) string {
if tenant.Spec.DatabaseRef != nil {
return string(tenant.Spec.DatabaseRef.Name)
}
return binding.DynamicDatabaseName(string(tenant.UID))
}
@@ -19,6 +19,7 @@ import (
const ( const (
TenantFinalizer = "database.ayatori.ddupan.top/tenant-protection" TenantFinalizer = "database.ayatori.ddupan.top/tenant-protection"
DatabaseFinalizer = "database.ayatori.ddupan.top/database-protection" DatabaseFinalizer = "database.ayatori.ddupan.top/database-protection"
readyCondition = "Ready"
) )
// BindingResources 读取领域所需事实,并把用例结果呈现为 CR、finalizer 与 Conditions。 // BindingResources 读取领域所需事实,并把用例结果呈现为 CR、finalizer 与 Conditions。
@@ -148,7 +149,7 @@ func (r *BindingResources) presentStatus(ctx context.Context, object *databasev1
} }
} }
meta.SetStatusCondition(&object.Status.Conditions, metav1.Condition{ meta.SetStatusCondition(&object.Status.Conditions, metav1.Condition{
Type: "Ready", Status: metav1.ConditionFalse, Reason: status.Reason, Message: status.Message, Type: readyCondition, Status: metav1.ConditionFalse, Reason: status.Reason, Message: status.Message,
ObservedGeneration: object.Generation, ObservedGeneration: object.Generation,
}) })
if equality.Semantic.DeepEqual(*previous, object.Status) { if equality.Semantic.DeepEqual(*previous, object.Status) {
@@ -186,6 +187,6 @@ func bindingVersionConflict(resource, name string) error {
} }
func currentReady(generation int64, conditions []metav1.Condition) bool { func currentReady(generation int64, conditions []metav1.Condition) bool {
condition := meta.FindStatusCondition(conditions, "Ready") condition := meta.FindStatusCondition(conditions, readyCondition)
return condition != nil && condition.Status == metav1.ConditionTrue && condition.ObservedGeneration == generation return condition != nil && condition.Status == metav1.ConditionTrue && condition.ObservedGeneration == generation
} }
@@ -0,0 +1,28 @@
package kubernetes
import credentialdomain "git.ddupan.top/panxiao81/ayatori/internal/database/domain/credential"
const credentialsReadyCondition = "CredentialsReady"
// 持久化 API 的 Reason 保持兼容,领域内部只使用准备阶段。
var credentialReasons = map[credentialdomain.Phase]string{
credentialdomain.Pending: "PreparationPending",
credentialdomain.Pinned: "LocationPinned",
credentialdomain.Creating: "CreationStarted",
credentialdomain.Prepared: "CredentialPrepared",
credentialdomain.Conflict: "Conflict",
credentialdomain.Unavailable: "DependencyUnavailable",
credentialdomain.Stopped: "PreparationStopped",
credentialdomain.InvalidTarget: "InvalidTarget",
}
func credentialReason(phase credentialdomain.Phase) string { return credentialReasons[phase] }
func credentialPhase(reason string) credentialdomain.Phase {
for phase, value := range credentialReasons {
if value == reason {
return phase
}
}
return credentialdomain.Pending
}
@@ -0,0 +1,147 @@
package kubernetes
import (
"context"
"fmt"
credentialdomain "git.ddupan.top/panxiao81/ayatori/internal/database/domain/credential"
databasev1alpha1 "git.ddupan.top/panxiao81/ayatori/api/database/v1alpha1"
"git.ddupan.top/panxiao81/ayatori/internal/database/application"
"git.ddupan.top/panxiao81/ayatori/internal/database/domain/binding"
"k8s.io/apimachinery/pkg/api/equality"
apierrors "k8s.io/apimachinery/pkg/api/errors"
"k8s.io/apimachinery/pkg/api/meta"
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
"k8s.io/apimachinery/pkg/types"
"sigs.k8s.io/controller-runtime/pkg/client"
"sigs.k8s.io/controller-runtime/pkg/controller/controllerutil"
)
// CredentialResources 只映射资源事实与状态写入,供应资格和执行顺序由用例决定。
type CredentialResources struct {
Client client.Client
Reader client.Reader
}
var _ application.CredentialResources = (*CredentialResources)(nil)
func (r *CredentialResources) Load(ctx context.Context, name string) (*application.CredentialRecord, error) {
database := &databasev1alpha1.PostgreSQLDatabase{}
if err := r.Reader.Get(ctx, types.NamespacedName{Name: name}, database); err != nil {
return nil, client.IgnoreNotFound(err)
}
record := &application.CredentialRecord{
Database: bindingDatabase(database).Database, Revision: database.ResourceVersion,
DatabaseProtected: controllerutil.ContainsFinalizer(database, DatabaseFinalizer),
Status: credentialStatus(database),
}
if database.Spec.Source != "Provision" {
return record, nil
}
if ref := database.Spec.TenantRef; ref != nil {
tenant := &databasev1alpha1.PostgreSQLTenant{}
err := r.Reader.Get(ctx, types.NamespacedName{Namespace: ref.Namespace, Name: string(ref.Name)}, tenant)
if err != nil && !apierrors.IsNotFound(err) {
return nil, err
}
if err == nil {
record.Tenant = &bindingTenant(tenant).Tenant
record.TenantGeneration = tenant.Generation
record.TenantProtected = controllerutil.ContainsFinalizer(tenant, TenantFinalizer)
}
}
instance := &databasev1alpha1.PostgreSQLInstance{}
err := r.Reader.Get(ctx, types.NamespacedName{Name: string(database.Spec.InstanceRef.Name)}, instance)
if err != nil {
return record, client.IgnoreNotFound(err)
}
observed, err := instanceRecord(instance)
if err != nil {
return nil, err
}
record.Instance = &credentialdomain.Instance{
Identity: binding.Identity{Name: instance.Name, UID: string(instance.UID)},
Deleting: !instance.DeletionTimestamp.IsZero(), Ready: currentReady(instance.Generation, instance.Status.Conditions),
Endpoint: observed.Target.Definition().Endpoint(),
}
record.InstanceGeneration = instance.Generation
return record, nil
}
func credentialStatus(database *databasev1alpha1.PostgreSQLDatabase) credentialdomain.State {
status := credentialdomain.State{Version: database.Status.CredentialVersion}
if ref := database.Status.CredentialRef; ref != nil {
status.Location = &credentialdomain.Location{Mount: ref.Mount, Path: ref.Path}
}
if condition := meta.FindStatusCondition(database.Status.Conditions, credentialsReadyCondition); condition != nil {
status.Phase, status.Message = credentialPhase(condition.Reason), condition.Message
}
return status
}
func (r *CredentialResources) Save(ctx context.Context, record *application.CredentialRecord, status credentialdomain.State) (*application.CredentialRecord, error) {
bindingResources := &BindingResources{Client: r.Client, Reader: r.Reader}
database, err := bindingResources.databaseAtVersion(ctx, &application.BindingDatabase{Database: record.Database, Revision: record.Revision})
if err != nil {
return nil, err
}
previous := database.Status.DeepCopy()
if status.Location != nil {
database.Status.CredentialRef = &databasev1alpha1.CredentialReference{Mount: status.Location.Mount, Path: status.Location.Path}
}
database.Status.CredentialVersion = status.Version
conditionStatus := metav1.ConditionFalse
if status.Phase == credentialdomain.Prepared {
conditionStatus = metav1.ConditionTrue
}
meta.SetStatusCondition(&database.Status.Conditions, metav1.Condition{
Type: credentialsReadyCondition, Status: conditionStatus, Reason: credentialReason(status.Phase), Message: status.Message,
ObservedGeneration: database.Generation,
})
meta.SetStatusCondition(&database.Status.Conditions, metav1.Condition{
Type: readyCondition, Status: metav1.ConditionFalse, Reason: "ProvisioningIncomplete",
Message: "凭据准备不代表 PostgreSQL 供应、应用登录或 Tenant 交付已经完成", ObservedGeneration: database.Generation,
})
database.Status.ObservedGeneration = database.Generation
if !equality.Semantic.DeepEqual(*previous, database.Status) {
if err := r.Client.Status().Update(ctx, database); err != nil {
return nil, err
}
}
updated := *record
updated.Database = bindingDatabase(database).Database
updated.Revision = database.ResourceVersion
updated.Status = credentialStatus(database)
// 旧 CRD 会裁剪未知 status 字段。不能把 HTTP 成功当作位置/版本已保存后继续写后端。
if !equality.Semantic.DeepEqual(updated.Status.Location, status.Location) || updated.Status.Version != status.Version {
return nil, fmt.Errorf("凭据位置或版本未被 API 保留;请先升级 Database CRD,未继续外部操作")
}
return &updated, nil
}
func (r *CredentialResources) CheckCurrent(ctx context.Context, record *application.CredentialRecord) error {
current, err := r.Load(ctx, record.Database.Identity.Name)
if err != nil {
return err
}
if current == nil || current.Database.Identity != record.Database.Identity || current.Revision != record.Revision ||
!sameCredentialDependencies(current, record) {
return apierrors.NewConflict(databasev1alpha1.GroupVersion.WithResource("postgresqldatabases").GroupResource(), record.Database.Identity.Name,
fmt.Errorf("凭据准备的资源快照已变化;停止本轮操作并重新观察"))
}
return nil
}
func sameCredentialDependencies(current, previous *application.CredentialRecord) bool {
if current.Tenant == nil || previous.Tenant == nil || current.Instance == nil || previous.Instance == nil {
return false
}
// Tenant Ready 的诊断变化、Instance 对同一 generation 的观测刷新不改变写入目标。
// 仍检查申请 spec generation、完整绑定身份、删除状态、保护和当前 Instance Ready。
return current.TenantGeneration == previous.TenantGeneration &&
equality.Semantic.DeepEqual(current.Tenant, previous.Tenant) &&
current.TenantProtected == previous.TenantProtected && current.DatabaseProtected == previous.DatabaseProtected &&
current.Instance.Instance == previous.Instance.Instance && current.InstanceGeneration == previous.InstanceGeneration &&
current.Instance.Endpoint == previous.Instance.Endpoint
}
@@ -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
} }
@@ -84,7 +84,7 @@ func (r *InstanceResources) PresentInstance(ctx context.Context, result applicat
ready = metav1.ConditionTrue ready = metav1.ConditionTrue
} }
meta.SetStatusCondition(&object.Status.Conditions, metav1.Condition{ meta.SetStatusCondition(&object.Status.Conditions, metav1.Condition{
Type: "Ready", Status: ready, ObservedGeneration: object.Generation, Type: readyCondition, Status: ready, ObservedGeneration: object.Generation,
Reason: result.Reason, Message: result.Message, Reason: result.Reason, Message: result.Message,
}) })
if !equality.Semantic.DeepEqual(*previous, object.Status) { if !equality.Semantic.DeepEqual(*previous, object.Status) {
@@ -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() == "" })
}
@@ -0,0 +1,47 @@
//go:build integration
package openbao_test
import (
"maps"
"testing"
bao "github.com/openbao/openbao/api/v2"
"git.ddupan.top/panxiao81/ayatori/internal/database/adapter/openbao"
)
func TestConfirmedCredentialRecoveryWithRealOpenBao(t *testing.T) {
root := baoFixture(t)
store := fixtureStore(t, root)
credential := fixtureCredential(t)
if err := store.Create(t.Context(), credentialPath, credential); err != nil {
t.Fatal(err)
}
// 新 adapter 模拟进程重启:只有已持久化的版本才能恢复读取。
restarted := fixtureStore(t, root)
if _, err := restarted.ReadConfirmed(t.Context(), credentialPath, 0); err != openbao.ErrConflict {
t.Fatal("已有值不能替代缺失的确认记录")
}
observed, err := restarted.ReadConfirmed(t.Context(), credentialPath, 1)
if err != nil || !maps.Equal(observed.SecretData(), credential.SecretData()) {
t.Fatal("确认版本的凭据不能恢复读取")
}
// 即使内容相同,新增版本也不属于已确认写入;历史版本仍在不能掩盖改写。
if _, err := root.KVv2("secret").Put(t.Context(), credentialPath, credential.SecretData(), bao.WithCheckAndSet(1)); err != nil {
t.Fatal("无法准备测试中的版本变化")
}
if _, err := restarted.ReadConfirmed(t.Context(), credentialPath, 1); err != openbao.ErrConflict {
t.Fatal("最新版本漂移必须报冲突")
}
if err := root.KVv2("secret").Delete(t.Context(), credentialPath); err != nil {
t.Fatal("无法准备测试中的软删除")
}
if _, err := restarted.ReadConfirmed(t.Context(), credentialPath, 1); err != openbao.ErrConflict {
t.Fatal("已确认凭据被删除必须报冲突")
}
metadata, err := root.KVv2("secret").GetMetadata(t.Context(), credentialPath)
if err != nil || metadata.CurrentVersion != 2 {
t.Fatal("恢复检查不得生成或覆盖凭据")
}
}
@@ -0,0 +1,69 @@
package openbao_test
import (
"encoding/json"
"net/http"
"net/http/httptest"
"testing"
"git.ddupan.top/panxiao81/ayatori/internal/database/adapter/openbao"
)
func credentialVersionMetadata(version int) map[string]any {
return map[string]any{kvVersionKey: version}
}
func TestConfirmedCredentialRead(t *testing.T) {
for _, test := range []struct {
name string
version int64
status int
meta map[string]any
invalid bool
want error
}{
{name: "已确认", version: 1, status: 200, meta: credentialVersionMetadata(1)},
{name: "未确认禁止读取", version: 0, want: openbao.ErrConflict},
{name: "版本已变化", version: 1, status: 200, meta: credentialVersionMetadata(2), want: openbao.ErrConflict},
{name: "版本为零", version: 1, status: 200, meta: credentialVersionMetadata(0), want: openbao.ErrConflict},
{name: "响应无法解析", version: 1, status: 200, want: openbao.ErrUnavailable},
{name: "内容损坏", version: 1, status: 200, meta: credentialVersionMetadata(1), invalid: true, want: openbao.ErrConflict},
{name: "凭据丢失", version: 1, status: 404, want: openbao.ErrConflict},
{name: "读取被拒绝", version: 1, status: 403, want: openbao.ErrUnavailable},
{name: "后端不可用", version: 1, status: 503, want: openbao.ErrUnavailable},
} {
t.Run(test.name, func(t *testing.T) {
credential := fixtureCredential(t)
server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
if test.version == 0 {
t.Error("没有持久化确认记录不能读取已有值")
}
if r.Method != http.MethodGet || r.URL.RawQuery != "" {
t.Error("只允许读取最新值,不能写入或回退历史版本")
}
w.WriteHeader(test.status)
if test.status != http.StatusOK {
return
}
data := credential.SecretData()
if test.invalid {
delete(data, "password")
}
if err := json.NewEncoder(w).Encode(map[string]any{
kvDataKey: map[string]any{kvDataKey: data, "metadata": test.meta},
}); err != nil {
t.Error("测试响应编码失败")
}
}))
defer server.Close()
store := fixtureStore(t, fixtureClient(t, server.URL))
observed, err := store.ReadConfirmed(t.Context(), credentialPath, test.version)
if err != test.want {
t.Fatalf("期望 %v,得到 %v", test.want, err)
}
if err != nil && observed.Validate() == nil {
t.Fatal("失败时不能返回可用凭据")
}
})
}
}
@@ -0,0 +1,40 @@
package openbao
import (
"context"
credentialdomain "git.ddupan.top/panxiao81/ayatori/internal/database/domain/credential"
"git.ddupan.top/panxiao81/ayatori/internal/database/application"
)
var _ application.CredentialStore = (*Credentials)(nil)
func (c *Credentials) ProvisionLocation(uid string) (credentialdomain.Location, error) {
path, err := c.ProvisionPath(uid)
if err != nil {
return credentialdomain.Location{}, err
}
return credentialdomain.Location{Mount: c.mount, Path: path}, nil
}
func (c *Credentials) ReadCredential(ctx context.Context, location credentialdomain.Location, version int64) (credentialdomain.ApplicationCredential, error) {
if location.Mount != c.mount {
return credentialdomain.ApplicationCredential{}, ErrInvalidLocation
}
if version == 0 {
return c.Read(ctx, location.Path)
}
return c.ReadConfirmed(ctx, location.Path, version)
}
func (c *Credentials) CreateCredential(ctx context.Context, location credentialdomain.Location, credential credentialdomain.ApplicationCredential) (int64, error) {
if location.Mount != c.mount {
return 0, ErrInvalidLocation
}
if err := c.Create(ctx, location.Path, credential); err != nil {
return 0, err
}
// Create 的合同是 CAS=0 且回读版本 1 和七键完全一致。
return 1, nil
}
@@ -0,0 +1,173 @@
/*
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 通过官方 SDK 适配应用凭据,不保存资源归属或重建供应状态。
package openbao
import (
"context"
"errors"
"maps"
"net/http"
"regexp"
"slices"
"strings"
credentialdomain "git.ddupan.top/panxiao81/ayatori/internal/database/domain/credential"
bao "github.com/openbao/openbao/api/v2"
"git.ddupan.top/panxiao81/ayatori/internal/database/application"
)
var (
ErrInvalidLocation = application.ErrCredentialLocation
ErrUnavailable = application.ErrCredentialUnavailable
ErrNotFound = application.ErrCredentialNotFound
ErrConflict = application.ErrCredentialConflict
ErrUncertain = application.ErrCredentialUncertain
)
var pathSegment = regexp.MustCompile(`^[A-Za-z0-9_-]+$`)
// Credentials 使用独立的 SDK client;认证与短期 token 生命周期由部署装配负责。
// 本适配器既不自动认领已有值,也不提供覆盖、轮换或删除操作。
type Credentials struct {
kv *bao.KVv2
mount string
basePath string
}
// NewCredentials 不登录、不读取环境 token。调用方必须提供专用的已认证 client。
// client 由公共 infra 禁用自动重试,防止第一次结果丢失后被 CAS 错误掩盖。
func NewCredentials(client *bao.Client, mount, basePath string) (*Credentials, error) {
if client == nil || !validPath(mount) || !validPath(basePath) {
return nil, ErrInvalidLocation
}
return &Credentials{kv: client.KVv2(mount), mount: mount, basePath: basePath}, nil
}
func validPath(value string) bool {
for segment := range strings.SplitSeq(value, "/") {
if !pathSegment.MatchString(segment) || segment == "data" || segment == "metadata" {
return false
}
}
return true
}
// ProvisionPath 只按 Database UID 定位;调用方须先持久化位置,再执行外部写入。
func (c *Credentials) ProvisionPath(databaseUID string) (string, error) {
if !pathSegment.MatchString(databaseUID) {
return "", ErrInvalidLocation
}
return c.basePath + "/" + databaseUID, nil
}
func (c *Credentials) accepts(path string) bool {
return validPath(path) && strings.HasPrefix(path, c.basePath+"/")
}
// Read 只读取调用方已确认关联的路径;成功读取不构成对既有凭据的自动认领。
func (c *Credentials) Read(ctx context.Context, path string) (credentialdomain.ApplicationCredential, error) {
secret, err := c.read(ctx, path)
if err != nil {
return credentialdomain.ApplicationCredential{}, err
}
return credentialdomain.ParseApplicationCredential(secret.Data)
}
// ReadConfirmed 读取最新值并核对已持久化的确认版本,不回退读取历史版本。
// 确认后的删除或改写需要人工处理,不能因此重新生成密码。
func (c *Credentials) ReadConfirmed(ctx context.Context, path string, version int64) (credentialdomain.ApplicationCredential, error) {
if version < 1 {
return credentialdomain.ApplicationCredential{}, ErrConflict
}
secret, err := c.read(ctx, path)
if errors.Is(err, ErrNotFound) {
return credentialdomain.ApplicationCredential{}, ErrConflict
}
if err != nil {
return credentialdomain.ApplicationCredential{}, err
}
if secret.VersionMetadata == nil || int64(secret.VersionMetadata.Version) != version {
return credentialdomain.ApplicationCredential{}, ErrConflict
}
credential, err := credentialdomain.ParseApplicationCredential(secret.Data)
if err != nil {
return credentialdomain.ApplicationCredential{}, ErrConflict
}
return credential, nil
}
func (c *Credentials) read(ctx context.Context, path string) (*bao.KVSecret, error) {
if !c.accepts(path) {
return nil, ErrInvalidLocation
}
secret, err := c.kv.Get(ctx, path)
if errors.Is(err, bao.ErrSecretNotFound) {
return nil, ErrNotFound
}
if err != nil {
return nil, ErrUnavailable
}
if secret == nil || secret.Data == nil {
return nil, ErrNotFound
}
return secret, nil
}
// Create 只创建从未存在过的路径,并验证回读七键与提交值完全一致。
// 任何不确定写入都不返回凭据;上层必须停止供应并持久化冲突,不能重新生成密码。
func (c *Credentials) Create(ctx context.Context, path string, credential credentialdomain.ApplicationCredential) error {
if !c.accepts(path) {
return ErrInvalidLocation
}
if err := credential.Validate(); err != nil {
return err
}
if ctx.Err() != nil {
return ErrUnavailable
}
data := credential.SecretData()
created, err := c.kv.Put(ctx, path, data, bao.WithCheckAndSet(0))
if err != nil {
// 明确的认证/权限拒绝没有发生写入,可以等待依赖恢复。
// SDK 的原始错误可能携带路径及响应体,不向外传播。
if response, ok := errors.AsType[*bao.ResponseError](err); ok {
switch response.StatusCode {
case http.StatusUnauthorized, http.StatusForbidden:
return ErrUnavailable
case http.StatusBadRequest:
if slices.Contains(response.Errors, "check-and-set parameter did not match the current version") {
return ErrConflict
}
}
}
return ErrUncertain
}
if created == nil || created.VersionMetadata == nil || created.VersionMetadata.Version != 1 {
return ErrUncertain
}
observed, err := c.kv.Get(ctx, path)
if err != nil || observed == nil || observed.VersionMetadata == nil || observed.VersionMetadata.Version != 1 {
return ErrUncertain
}
if !maps.Equal(data, observed.Data) {
return ErrUncertain
}
return nil
}
@@ -0,0 +1,227 @@
//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"
"errors"
"maps"
"net/http"
"net/http/httptest"
"net/http/httputil"
"net/url"
"os/exec"
"regexp"
"strings"
"sync"
"testing"
"time"
bao "github.com/openbao/openbao/api/v2"
"git.ddupan.top/panxiao81/ayatori/internal/database/adapter/openbao"
)
// 只连接本测试创建的无持久卷 dev server,不接受生产地址或环境 token。
func baoFixture(t *testing.T) *bao.Client {
t.Helper()
const image = "openbao/openbao@sha256:5b2486ab0fb90bbc788cc345b0a08616dfb375873ee8be5df3a2fd4d378a67e0"
prepareBaoImage(t, image)
// 冷缓存拉取不占用容器启动和健康检查的一分钟预算。
ctx, cancel := context.WithTimeout(t.Context(), time.Minute)
defer cancel()
output, err := exec.CommandContext(ctx, "docker", "run", "--pull=never", "--rm", "-d", "-p", "127.0.0.1::8200",
image, "server", "-dev", "-dev-root-token-id="+fixtureToken, "-dev-listen-address=0.0.0.0:8200").Output()
if err != nil {
t.Fatalf("cannot start isolated OpenBao fixture: %s", baoCommandError(ctx, err))
}
id := strings.TrimSpace(string(output))
if !regexp.MustCompile(`^[a-f0-9]{64}$`).MatchString(id) {
t.Fatal("unexpected fixture container ID")
}
t.Cleanup(func() {
cleanup, stop := context.WithTimeout(context.Background(), 30*time.Second)
defer stop()
if exec.CommandContext(cleanup, "docker", "rm", "-f", id).Run() != nil {
t.Error("OpenBao fixture cleanup failed")
}
})
output, err = exec.CommandContext(ctx, "docker", "inspect", "--format",
`{{(index (index .NetworkSettings.Ports "8200/tcp") 0).HostPort}}`, id).Output()
if err != nil {
t.Fatalf("cannot inspect fixture port: %s", baoCommandError(ctx, err))
}
client := fixtureClient(t, "http://127.0.0.1:"+strings.TrimSpace(string(output)))
client.SetMaxRetries(0)
for {
if _, err := client.Sys().HealthWithContext(ctx); err == nil {
return client
}
select {
case <-ctx.Done():
t.Fatal("OpenBao fixture startup timed out")
case <-time.After(100 * time.Millisecond):
}
}
}
func prepareBaoImage(t *testing.T, image string) {
t.Helper()
ctx, cancel := context.WithTimeout(t.Context(), 5*time.Minute)
defer cancel()
if exec.CommandContext(ctx, "docker", "image", "inspect", image).Run() == nil {
return
}
t.Log("pulling isolated OpenBao fixture image (timeout: 5m)")
if _, err := exec.CommandContext(ctx, "docker", "pull", image).Output(); err != nil {
t.Fatalf("cannot pull OpenBao fixture image: %s", baoCommandError(ctx, err))
}
}
// 保留 Docker stderr 与超时原因,但不泄露测试 token/password 或完整命令参数。
func baoCommandError(ctx context.Context, err error) string {
detail := err.Error()
if exitErr, ok := errors.AsType[*exec.ExitError](err); ok {
detail += ": " + strings.TrimSpace(string(exitErr.Stderr))
}
if ctx.Err() != nil {
detail += ": " + ctx.Err().Error()
}
return strings.NewReplacer(fixtureToken, "[REDACTED]", fixturePassword, "[REDACTED]").Replace(detail)
}
func TestBaoCommandError(t *testing.T) {
err := &exec.ExitError{Stderr: []byte("registry unavailable " + fixtureToken + " " + fixturePassword)}
ctx, cancel := context.WithCancel(t.Context())
cancel()
detail := baoCommandError(ctx, err)
if !strings.Contains(detail, "registry unavailable") || !strings.Contains(detail, "context canceled") {
t.Fatal("Docker diagnostic or context failure was lost")
}
if strings.Contains(detail, fixtureToken) || strings.Contains(detail, fixturePassword) {
t.Fatal("Docker diagnostic exposed fixture credentials")
}
}
func TestCredentialConcurrentCreateWithRealOpenBao(t *testing.T) {
root := baoFixture(t)
store := fixtureStore(t, root)
credential := fixtureCredential(t)
results := make(chan error, 2)
var workers sync.WaitGroup
for range 2 {
workers.Go(func() { results <- store.Create(t.Context(), credentialPath, credential) })
}
workers.Wait()
close(results)
succeeded, conflicted := 0, 0
for err := range results {
switch err {
case nil:
succeeded++
case openbao.ErrConflict:
conflicted++
default:
t.Fatal("unexpected concurrent create result")
}
}
if succeeded != 1 || conflicted != 1 {
t.Fatal("CAS must allow exactly one creator")
}
}
func TestCredentialLostWriteResponseWithRealOpenBao(t *testing.T) {
root := baoFixture(t)
address, err := url.Parse(root.Address())
if err != nil {
t.Fatal("invalid fixture address")
}
proxy := httputil.NewSingleHostReverseProxy(address)
proxy.ModifyResponse = func(response *http.Response) error {
if response.Request.Method == http.MethodPut && response.StatusCode == http.StatusOK {
return errors.New("fixture drops successful write response")
}
return nil
}
proxy.ErrorHandler = func(w http.ResponseWriter, _ *http.Request, _ error) {
w.WriteHeader(http.StatusBadGateway)
}
server := httptest.NewServer(proxy)
defer server.Close()
store := fixtureStore(t, fixtureClient(t, server.URL))
credential := fixtureCredential(t)
if err := store.Create(t.Context(), credentialPath, credential); err != openbao.ErrUncertain {
t.Fatal("lost response must stop provisioning")
}
confirmed, err := root.KVv2("secret").Get(t.Context(), credentialPath)
if err != nil || !maps.Equal(confirmed.Data, credential.SecretData()) || confirmed.VersionMetadata.Version != 1 {
t.Fatal("fault injection did not preserve the original write")
}
if err := fixtureStore(t, root).Create(t.Context(), credentialPath, credential); err != openbao.ErrConflict {
t.Fatal("restart must not adopt an unconfirmed write")
}
}
func TestCredentialsWithRealOpenBao(t *testing.T) {
root := baoFixture(t)
ctx := t.Context()
// root 仅用于 fixture 装配;实际读写使用固定前缀的短期 token。
policy := `path "secret/data/applications/*" { capabilities = ["create", "update", "read"] }`
if err := root.Sys().PutPolicyWithContext(ctx, "application-fixture", policy); err != nil {
t.Fatal("cannot configure fixture policy")
}
secret, err := root.Auth().Token().CreateWithContext(ctx, &bao.TokenCreateRequest{
Policies: []string{"application-fixture"}, NoDefaultPolicy: true, TTL: "5m",
})
if err != nil {
t.Fatal("cannot create scoped fixture token")
}
client := fixtureClient(t, root.Address())
client.SetToken(secret.Auth.ClientToken)
store := fixtureStore(t, client)
credential := fixtureCredential(t)
if err := store.Create(ctx, credentialPath, credential); err != nil {
t.Fatal(err)
}
// 重建适配器读取已确认路径;重复 Create 仍报冲突,不把读取当作认领。
restarted := fixtureStore(t, client)
observed, err := restarted.Read(ctx, credentialPath)
if err != nil || !maps.Equal(observed.SecretData(), credential.SecretData()) {
t.Fatal("confirmed credential was not preserved across adapter restart")
}
if err := restarted.Create(ctx, credentialPath, credential); !errors.Is(err, openbao.ErrConflict) {
t.Fatal("existing credential must conflict even if contents match")
}
metadata, err := root.KVv2("secret").GetMetadata(ctx, credentialPath)
if err != nil || metadata.CurrentVersion != 1 {
t.Fatal("duplicate create changed credential version")
}
if _, err := client.KVv2("secret").Get(ctx, "management/instance"); err == nil {
t.Fatal("scoped token accessed management credentials")
}
if err := root.KVv2("secret").Delete(ctx, credentialPath); err != nil {
t.Fatal("cannot soft-delete fixture credential")
}
if _, err := store.Read(ctx, credentialPath); err != openbao.ErrNotFound {
t.Fatal("soft-deleted credential must not be usable")
}
if err := store.Create(ctx, credentialPath, credential); err != openbao.ErrConflict {
t.Fatal("soft-deleted credential must not be recreated")
}
}
@@ -0,0 +1,193 @@
/*
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"
"errors"
"net/http"
"net/http/httptest"
"sync/atomic"
"testing"
credentialdomain "git.ddupan.top/panxiao81/ayatori/internal/database/domain/credential"
bao "github.com/openbao/openbao/api/v2"
"git.ddupan.top/panxiao81/ayatori/internal/database/adapter/openbao"
)
const (
credentialPath = "applications/database-uid"
fixturePassword = "AYATORI-TEST-ONLY-application-password"
fixtureToken = "AYATORI-TEST-ONLY-bao-token"
kvDataKey = "data"
kvVersionKey = "version"
)
func fixtureCredential(t *testing.T) credentialdomain.ApplicationCredential {
t.Helper()
credential, err := credentialdomain.ParseApplicationCredential(map[string]any{
"username": "app_owner", "password": fixturePassword, "database": "app",
"host": "postgres.example", "hostaddr": "192.0.2.1", "port": "5432", "sslmode": "verify-full",
})
if err != nil {
t.Fatal(err)
}
return credential
}
func TestCredentialReadbackMustConfirmTheWrite(t *testing.T) {
for _, scenario := range []string{"read failure", "changed version", "changed password", "missing metadata"} {
t.Run(scenario, func(t *testing.T) {
credential := fixtureCredential(t)
var writes atomic.Int32
server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
if r.Method == http.MethodPut {
writes.Add(1)
var request struct {
Options struct {
CAS *int `json:"cas"`
} `json:"options"`
}
if json.NewDecoder(r.Body).Decode(&request) != nil || request.Options.CAS == nil || *request.Options.CAS != 0 {
t.Error("create request must explicitly require CAS=0")
}
if err := json.NewEncoder(w).Encode(map[string]any{kvDataKey: map[string]any{kvVersionKey: 1}}); err != nil {
t.Error("cannot encode fixture write response")
}
return
}
if scenario == "read failure" {
w.WriteHeader(http.StatusServiceUnavailable)
return
}
data := credential.SecretData()
version := 1
if scenario == "changed version" {
version = 2
}
if scenario == "changed password" {
data["password"] = "modified"
}
response := map[string]any{kvDataKey: data}
if scenario != "missing metadata" {
response["metadata"] = map[string]any{kvVersionKey: version}
}
if err := json.NewEncoder(w).Encode(map[string]any{kvDataKey: response}); err != nil {
t.Error("cannot encode fixture read response")
}
}))
defer server.Close()
store := fixtureStore(t, fixtureClient(t, server.URL))
if err := store.Create(t.Context(), credentialPath, credential); err != openbao.ErrUncertain || writes.Load() != 1 {
t.Fatal("unconfirmed readback must stop after one write")
}
})
}
}
func fixtureClient(t *testing.T, address string) *bao.Client {
t.Helper()
config := bao.DefaultConfig()
config.Address = address
config.MaxRetries = 0
client, err := bao.NewClient(config)
if err != nil {
t.Fatal("cannot construct fixture client")
}
client.SetToken(fixtureToken)
return client
}
func fixtureStore(t *testing.T, client *bao.Client) *openbao.Credentials {
t.Helper()
store, err := openbao.NewCredentials(client, "secret", "applications")
if err != nil {
t.Fatal(err)
}
return store
}
func TestCredentialLocationScope(t *testing.T) {
client := fixtureClient(t, "http://127.0.0.1:1")
store := fixtureStore(t, client)
path, err := store.ProvisionPath("database-uid")
if err != nil || path != credentialPath {
t.Fatal("unexpected stable location")
}
for _, path := range []string{"", "/absolute", "applications", "applications-other/key", "applications/../management", "applications/%2e%2e/key", "applications//key", "applications/data/key"} {
if _, err := store.Read(t.Context(), path); !errors.Is(err, openbao.ErrInvalidLocation) {
t.Fatal("accepted invalid location")
}
if err := store.Create(t.Context(), path, fixtureCredential(t)); !errors.Is(err, openbao.ErrInvalidLocation) {
t.Fatal("accepted invalid create location")
}
}
for _, uid := range []string{"", "../key", "a/b", "a?b"} {
if _, err := store.ProvisionPath(uid); err == nil {
t.Fatal("accepted invalid UID")
}
}
for _, invalid := range []string{"", "data", "metadata", "../secret", "secret/", "secret?query"} {
if _, err := openbao.NewCredentials(client, invalid, "applications"); err == nil {
t.Fatal("accepted invalid mount")
}
if _, err := openbao.NewCredentials(client, "secret", invalid); err == nil {
t.Fatal("accepted invalid base path")
}
}
}
func TestCredentialWriteFailureDoesNotRetryOrLeak(t *testing.T) {
var requests atomic.Int32
server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, _ *http.Request) {
requests.Add(1)
http.Error(w, fixturePassword+fixtureToken, http.StatusInternalServerError)
}))
defer server.Close()
store := fixtureStore(t, fixtureClient(t, server.URL))
if err := store.Create(t.Context(), credentialPath, fixtureCredential(t)); err != openbao.ErrUncertain {
t.Fatal("write error must be a redacted uncertain outcome")
}
if requests.Load() != 1 {
t.Fatal("SDK retried an uncertain write")
}
if _, err := store.Read(t.Context(), credentialPath); err != openbao.ErrUnavailable {
t.Fatal("read error must be redacted")
}
ctx, cancel := context.WithCancel(t.Context())
cancel()
if err := store.Create(ctx, credentialPath, fixtureCredential(t)); err != openbao.ErrUnavailable || requests.Load() != 2 {
t.Fatal("canceled operation must not write")
}
}
func TestCredentialWriteDeniedBeforeExecution(t *testing.T) {
for _, status := range []int{http.StatusUnauthorized, http.StatusForbidden} {
server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, _ *http.Request) {
http.Error(w, fixtureToken, status)
}))
store := fixtureStore(t, fixtureClient(t, server.URL))
err := store.Create(t.Context(), credentialPath, fixtureCredential(t))
server.Close()
if err != openbao.ErrUnavailable {
t.Fatalf("status %d: definite rejection should wait for dependency recovery, got %v", status, err)
}
}
}
@@ -0,0 +1,62 @@
//go:build integration
package openbao_test
import (
"context"
"sync"
"sync/atomic"
"testing"
apierrors "k8s.io/apimachinery/pkg/api/errors"
"git.ddupan.top/panxiao81/ayatori/internal/database/application"
)
// 两个用例都先拿到同一真实 resourceVersion,再竞争固定位置;后续写入仍经过 API server。
type concurrentCredentialResources struct {
application.CredentialResources
readers atomic.Int32
loaded sync.WaitGroup
}
func (r *concurrentCredentialResources) Load(ctx context.Context, name string) (*application.CredentialRecord, error) {
record, err := r.CredentialResources.Load(ctx, name)
if r.readers.Add(1) <= 2 {
r.loaded.Done()
r.loaded.Wait()
}
return record, err
}
func testPreparationConcurrency(t *testing.T, f *preparationFixture) {
database, _ := f.bound(t, "concurrent")
service := f.service(t)
resources := &concurrentCredentialResources{CredentialResources: service.Resources}
resources.loaded.Add(2)
service.Resources = resources
results := make(chan error, 2)
for range 2 {
go func() { results <- service.Reconcile(t.Context(), database.Name) }()
}
succeeded, conflicted := 0, 0
for range 2 {
err := <-results
switch {
case err == nil:
succeeded++
case apierrors.IsConflict(err):
conflicted++
default:
t.Fatalf("并发用例返回意外错误: %v", err)
}
}
if succeeded != 1 || conflicted != 1 {
t.Fatal("同一快照只能有一个用例成功固定位置并继续创建")
}
f.status(t, database, 1, "CredentialPrepared")
stored, err := f.bao.KVv2("secret").Get(t.Context(), database.Status.CredentialRef.Path)
if err != nil || stored.VersionMetadata.Version != 1 {
t.Fatal("并发准备用例只能产生一个凭据版本")
}
}
@@ -0,0 +1,117 @@
//go:build integration
package openbao_test
import (
"testing"
bao "github.com/openbao/openbao/api/v2"
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"
"k8s.io/client-go/rest"
ctrl "sigs.k8s.io/controller-runtime"
"sigs.k8s.io/controller-runtime/pkg/client"
"sigs.k8s.io/controller-runtime/pkg/envtest"
databasev1alpha1 "git.ddupan.top/panxiao81/ayatori/api/database/v1alpha1"
"git.ddupan.top/panxiao81/ayatori/internal/database/adapter/kubernetes"
"git.ddupan.top/panxiao81/ayatori/internal/database/application"
databasecontroller "git.ddupan.top/panxiao81/ayatori/internal/database/controller"
)
const preparationNamespace = "credential-preparation"
type preparationFixture struct {
api client.Client
config *rest.Config
scheme *runtime.Scheme
bao *bao.Client
}
func newPreparationFixture(t *testing.T) *preparationFixture {
t.Helper()
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)
}
})
scheme := runtime.NewScheme()
if err := databasev1alpha1.AddToScheme(scheme); err != nil {
t.Fatal(err)
}
if err := corev1.AddToScheme(scheme); err != nil {
t.Fatal(err)
}
api, err := client.New(config, client.Options{Scheme: scheme})
if err != nil {
t.Fatal(err)
}
if err := api.Create(t.Context(), &corev1.Namespace{Name: preparationNamespace}); err != nil {
t.Fatal(err)
}
return &preparationFixture{api: api, config: config, scheme: scheme, bao: baoFixture(t)}
}
func (f *preparationFixture) bound(t *testing.T, name string) (*databasev1alpha1.PostgreSQLDatabase, *databasev1alpha1.PostgreSQLTenant) {
t.Helper()
instance := &databasev1alpha1.PostgreSQLInstance{Name: name, Spec: databasev1alpha1.PostgreSQLInstanceSpec{
Endpoint: databasev1alpha1.PostgreSQLEndpoint{Host: "postgres.example", HostAddr: "192.0.2.1"},
AdminCredentialRef: databasev1alpha1.AdminCredentialReference{Name: "management"},
}}
if err := f.api.Create(t.Context(), instance); err != nil {
t.Fatal(err)
}
// 本切片不连接 PG;这里只声明供应前置观察,真实管理能力由既有 PG 集成测试覆盖。
meta.SetStatusCondition(&instance.Status.Conditions, metav1.Condition{
Type: "Ready", Status: metav1.ConditionTrue, Reason: "FixtureReady", Message: "隔离测试前置观察", ObservedGeneration: instance.Generation,
})
if err := f.api.Status().Update(t.Context(), instance); err != nil {
t.Fatal(err)
}
tenant := &databasev1alpha1.PostgreSQLTenant{Name: name, Namespace: preparationNamespace}
tenant.Spec.Provision = &databasev1alpha1.DatabaseProvisionRequest{InstanceRef: databasev1alpha1.InstanceReference{Name: databasev1alpha1.ObjectName(name)}}
if err := f.api.Create(t.Context(), tenant); err != nil {
t.Fatal(err)
}
resources := &kubernetes.BindingResources{Client: f.api, Reader: f.api}
binder := databasecontroller.NewBindingReconciler(f.api, &application.BindingService{Resources: resources}, resources)
if _, err := binder.Reconcile(t.Context(), ctrl.Request{NamespacedName: client.ObjectKeyFromObject(tenant)}); err != nil {
t.Fatal(err)
}
if err := f.api.Get(t.Context(), client.ObjectKeyFromObject(tenant), tenant); err != nil {
t.Fatal(err)
}
database := &databasev1alpha1.PostgreSQLDatabase{}
if err := f.api.Get(t.Context(), client.ObjectKey{Name: string(tenant.Status.DatabaseRef.Name)}, database); err != nil {
t.Fatal(err)
}
return database, tenant
}
func (f *preparationFixture) service(t *testing.T) application.CredentialPreparation {
t.Helper()
return application.CredentialPreparation{
Resources: &kubernetes.CredentialResources{Client: f.api, Reader: f.api}, Store: fixtureStore(t, f.bao),
}
}
func (f *preparationFixture) status(t *testing.T, database *databasev1alpha1.PostgreSQLDatabase, version int64, reason string) {
t.Helper()
if err := f.api.Get(t.Context(), client.ObjectKeyFromObject(database), database); err != nil {
t.Fatal(err)
}
condition := meta.FindStatusCondition(database.Status.Conditions, "CredentialsReady")
if database.Status.CredentialVersion != version || condition == nil || condition.Reason != reason {
t.Fatalf("凭据版本或条件不符:version=%d,期望 reason=%s", database.Status.CredentialVersion, reason)
}
if meta.IsStatusConditionTrue(database.Status.Conditions, "Ready") {
t.Fatal("凭据准备不得宣告 Database Ready")
}
}
@@ -0,0 +1,238 @@
//go:build integration
package openbao_test
import (
"context"
"errors"
"maps"
"testing"
credentialdomain "git.ddupan.top/panxiao81/ayatori/internal/database/domain/credential"
bao "github.com/openbao/openbao/api/v2"
"sigs.k8s.io/controller-runtime/pkg/client"
databasev1alpha1 "git.ddupan.top/panxiao81/ayatori/api/database/v1alpha1"
"git.ddupan.top/panxiao81/ayatori/internal/database/adapter/kubernetes"
"git.ddupan.top/panxiao81/ayatori/internal/database/adapter/openbao"
"git.ddupan.top/panxiao81/ayatori/internal/database/application"
"git.ddupan.top/panxiao81/ayatori/internal/database/domain/binding"
)
func TestCredentialPreparationWithRealBackends(t *testing.T) {
fixture := newPreparationFixture(t)
t.Run("创建重启和幂等", func(t *testing.T) { testPreparationRestart(t, fixture) })
t.Run("已有值不认领", func(t *testing.T) { testPreparationExisting(t, fixture) })
t.Run("确认写入失败", func(t *testing.T) { testPreparationLostConfirmation(t, fixture) })
t.Run("后端结果不确定", func(t *testing.T) { testPreparationUncertain(t, fixture) })
t.Run("创建中绑定变化", func(t *testing.T) { testPreparationChangedBinding(t, fixture) })
t.Run("依赖恢复与固定位置", func(t *testing.T) { testPreparationDependencies(t, fixture) })
t.Run("并发用例只有一个写入", func(t *testing.T) { testPreparationConcurrency(t, fixture) })
t.Run("位置被API裁剪时拒绝外部写入", func(t *testing.T) { testPreparationPruning(t, fixture) })
t.Run("实际manager的watch与重启", func(t *testing.T) { testPreparationWatch(t, fixture) })
}
func testPreparationRestart(t *testing.T, f *preparationFixture) {
database, tenant := f.bound(t, "restart")
if err := f.service(t).Reconcile(t.Context(), database.Name); err != nil {
t.Fatal(err)
}
f.status(t, database, 1, "CredentialPrepared")
path := database.Status.CredentialRef.Path
before, err := f.bao.KVv2("secret").Get(t.Context(), path)
if err != nil {
t.Fatal("无法读取隔离测试凭据")
}
if before.Data["username"] != "restart" || before.Data["database"] != "restart" || len(before.Data) != 7 {
t.Fatal("用例未使用绑定目标生成七键凭据")
}
// 两次新用例实例模拟重启,第二轮必须没有无意义状态更新。
for range 2 {
if err := f.service(t).Reconcile(t.Context(), database.Name); err != nil {
t.Fatal(err)
}
}
f.status(t, database, 1, "CredentialPrepared")
revision := database.ResourceVersion
if err := f.service(t).Reconcile(t.Context(), database.Name); err != nil {
t.Fatal(err)
}
f.status(t, database, 1, "CredentialPrepared")
if database.ResourceVersion != revision {
t.Fatal("幂等重试不应改写 status")
}
after, err := f.bao.KVv2("secret").Get(t.Context(), path)
if err != nil || after.VersionMetadata.Version != 1 || !maps.Equal(before.Data, after.Data) {
t.Fatal("重启后不应生成或覆盖凭据")
}
if err := f.api.Delete(t.Context(), tenant); err != nil {
t.Fatal(err)
}
if err := f.service(t).Reconcile(t.Context(), database.Name); err != nil {
t.Fatal(err)
}
f.status(t, database, 1, "PreparationStopped")
if err := f.api.Delete(t.Context(), database); err != nil {
t.Fatal(err)
}
if err := f.service(t).Reconcile(t.Context(), database.Name); err != nil {
t.Fatal(err)
}
f.status(t, database, 1, "PreparationStopped")
if len(database.Finalizers) == 0 {
t.Fatal("本切片不得解除删除保护")
}
if _, err := f.bao.KVv2("secret").Get(t.Context(), path); err != nil {
t.Fatal("删除流程不得在本切片清理凭据")
}
}
func testPreparationExisting(t *testing.T, f *preparationFixture) {
database, _ := f.bound(t, "existing")
store := fixtureStore(t, f.bao)
path, err := store.ProvisionPath(string(database.UID))
if err != nil {
t.Fatal(err)
}
if err := store.Create(t.Context(), path, fixtureCredential(t)); err != nil {
t.Fatal(err)
}
if err := f.service(t).Reconcile(t.Context(), database.Name); err != nil {
t.Fatal(err)
}
f.status(t, database, 0, binding.Conflict)
// 即使管理员删除了后端值,未解除的不确定状态也不能自动创建。
if err := f.bao.KVv2("secret").DeleteMetadata(t.Context(), path); err != nil {
t.Fatal("无法清理隔离测试 key")
}
if err := f.service(t).Reconcile(t.Context(), database.Name); err != nil {
t.Fatal(err)
}
if _, err := f.bao.KVv2("secret").Get(t.Context(), path); !errors.Is(err, bao.ErrSecretNotFound) {
t.Fatal("Conflict 重入不应生成替代凭据")
}
}
// 故障只注入确认保存边界;此前写入使用真实 API server 和 Bao。
type failedConfirmationClient struct{ client.Client }
func (c failedConfirmationClient) Status() client.SubResourceWriter {
return failedConfirmationWriter{SubResourceWriter: c.Client.Status()}
}
type failedConfirmationWriter struct{ client.SubResourceWriter }
func (w failedConfirmationWriter) Update(ctx context.Context, object client.Object, options ...client.SubResourceUpdateOption) error {
if database, ok := object.(*databasev1alpha1.PostgreSQLDatabase); ok && database.Status.CredentialVersion > 0 {
return errors.New("fixture refuses confirmation update")
}
return w.SubResourceWriter.Update(ctx, object, options...)
}
func testPreparationLostConfirmation(t *testing.T, f *preparationFixture) {
database, _ := f.bound(t, "confirmation")
service := f.service(t)
service.Resources = &kubernetes.CredentialResources{Client: failedConfirmationClient{f.api}, Reader: f.api}
if err := service.Reconcile(t.Context(), database.Name); err == nil {
t.Fatal("确认写入失败应返回 API 错误")
}
f.status(t, database, 0, "CreationStarted")
if err := f.service(t).Reconcile(t.Context(), database.Name); err != nil {
t.Fatal(err)
}
f.status(t, database, 0, binding.Conflict)
stored, err := f.bao.KVv2("secret").Get(t.Context(), database.Status.CredentialRef.Path)
if err != nil || stored.VersionMetadata.Version != 1 {
t.Fatal("写入结果必须保留且不能被重启认领")
}
}
type afterCreateStore struct {
application.CredentialStore
after func() error
}
func (s afterCreateStore) CreateCredential(ctx context.Context, location credentialdomain.Location, credential credentialdomain.ApplicationCredential) (int64, error) {
version, err := s.CredentialStore.CreateCredential(ctx, location, credential)
if err != nil {
return 0, err
}
if err := s.after(); err != nil {
return 0, err
}
return version, nil
}
func testPreparationUncertain(t *testing.T, f *preparationFixture) {
database, _ := f.bound(t, "uncertain")
service := f.service(t)
service.Store = afterCreateStore{CredentialStore: service.Store, after: func() error { return application.ErrCredentialUncertain }}
if err := service.Reconcile(t.Context(), database.Name); err != nil {
t.Fatal(err)
}
f.status(t, database, 0, binding.Conflict)
if err := f.service(t).Reconcile(t.Context(), database.Name); err != nil {
t.Fatal(err)
}
f.status(t, database, 0, binding.Conflict)
}
func testPreparationChangedBinding(t *testing.T, f *preparationFixture) {
database, tenant := f.bound(t, "changed")
service := f.service(t)
service.Store = afterCreateStore{CredentialStore: service.Store, after: func() error {
return f.api.Delete(t.Context(), tenant)
}}
if err := service.Reconcile(t.Context(), database.Name); err == nil {
t.Fatal("中途删除 Tenant 后不得确认凭据")
}
f.status(t, database, 0, "CreationStarted")
if err := f.service(t).Reconcile(t.Context(), database.Name); err != nil {
t.Fatal(err)
}
f.status(t, database, 0, binding.Conflict)
}
func testPreparationDependencies(t *testing.T, f *preparationFixture) {
database, _ := f.bound(t, "dependency")
denied := fixtureClient(t, f.bao.Address())
denied.ClearToken()
service := f.service(t)
service.Store = fixtureStore(t, denied)
if err := service.Reconcile(t.Context(), database.Name); err != nil {
t.Fatal(err)
}
f.status(t, database, 0, binding.DependencyUnavailable)
if err := f.service(t).Reconcile(t.Context(), database.Name); err != nil {
t.Fatal(err)
}
f.status(t, database, 1, "CredentialPrepared")
if err := service.Reconcile(t.Context(), database.Name); err != nil {
t.Fatal(err)
}
f.status(t, database, 1, binding.DependencyUnavailable)
if err := f.service(t).Reconcile(t.Context(), database.Name); err != nil {
t.Fatal(err)
}
f.status(t, database, 1, "CredentialPrepared")
moved, err := openbao.NewCredentials(f.bao, "other", "elsewhere")
if err != nil {
t.Fatal(err)
}
service.Store = moved
if err := service.Reconcile(t.Context(), database.Name); err != nil {
t.Fatal(err)
}
f.status(t, database, 1, binding.DependencyUnavailable)
if database.Status.CredentialRef.Mount != "secret" {
t.Fatal("配置变化不得迁移固定位置")
}
if err := f.bao.KVv2("secret").Delete(t.Context(), database.Status.CredentialRef.Path); err != nil {
t.Fatal("无法准备测试软删除")
}
if err := f.service(t).Reconcile(t.Context(), database.Name); err != nil {
t.Fatal(err)
}
f.status(t, database, 1, binding.Conflict)
}
@@ -0,0 +1,51 @@
//go:build integration
package openbao_test
import (
"context"
"errors"
"testing"
bao "github.com/openbao/openbao/api/v2"
"sigs.k8s.io/controller-runtime/pkg/client"
databasev1alpha1 "git.ddupan.top/panxiao81/ayatori/api/database/v1alpha1"
"git.ddupan.top/panxiao81/ayatori/internal/database/adapter/kubernetes"
)
// 在 API 边界模拟旧 schema 裁剪位置字段;剩余写入仍由真实 API server 处理。
type prunedLocationClient struct{ client.Client }
func (c prunedLocationClient) Status() client.SubResourceWriter {
return prunedLocationWriter{SubResourceWriter: c.Client.Status()}
}
type prunedLocationWriter struct{ client.SubResourceWriter }
func (w prunedLocationWriter) Update(ctx context.Context, object client.Object, options ...client.SubResourceUpdateOption) error {
if database, ok := object.(*databasev1alpha1.PostgreSQLDatabase); ok {
database.Status.CredentialRef = nil
}
return w.SubResourceWriter.Update(ctx, object, options...)
}
func testPreparationPruning(t *testing.T, f *preparationFixture) {
database, _ := f.bound(t, "pruning")
service := f.service(t)
service.Resources = &kubernetes.CredentialResources{Client: prunedLocationClient{f.api}, Reader: f.api}
if err := service.Reconcile(t.Context(), database.Name); err == nil {
t.Fatal("API 未保留位置时不得继续供应")
}
location, err := service.Store.ProvisionLocation(string(database.UID))
if err != nil {
t.Fatal(err)
}
if _, err := f.bao.KVv2(location.Mount).Get(t.Context(), location.Path); !errors.Is(err, bao.ErrSecretNotFound) {
t.Fatal("没有持久化位置时产生了外部凭据")
}
if err := f.service(t).Reconcile(t.Context(), database.Name); err != nil {
t.Fatal(err)
}
f.status(t, database, 1, "CredentialPrepared")
}
@@ -0,0 +1,144 @@
//go:build integration
package openbao_test
import (
"context"
"testing"
"time"
"git.ddupan.top/panxiao81/ayatori/internal/database/adapter/kubernetes"
"git.ddupan.top/panxiao81/ayatori/internal/database/application"
"k8s.io/apimachinery/pkg/api/meta"
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
ctrl "sigs.k8s.io/controller-runtime"
"sigs.k8s.io/controller-runtime/pkg/client"
controllerconfig "sigs.k8s.io/controller-runtime/pkg/config"
metricsserver "sigs.k8s.io/controller-runtime/pkg/metrics/server"
databasev1alpha1 "git.ddupan.top/panxiao81/ayatori/api/database/v1alpha1"
databasecontroller "git.ddupan.top/panxiao81/ayatori/internal/database/controller"
)
func testPreparationWatch(t *testing.T, f *preparationFixture) {
database, _ := f.bound(t, "watch")
t.Cleanup(func() {
if !t.Failed() {
return
}
if err := f.api.Get(context.Background(), client.ObjectKeyFromObject(database), database); err == nil {
if condition := meta.FindStatusCondition(database.Status.Conditions, "CredentialsReady"); condition != nil {
t.Logf("失败时凭据条件: %s: %s", condition.Reason, condition.Message)
}
}
})
instance := &databasev1alpha1.PostgreSQLInstance{}
if err := f.api.Get(t.Context(), client.ObjectKey{Name: "watch"}, instance); err != nil {
t.Fatal(err)
}
setReady := func(ready metav1.ConditionStatus) {
t.Helper()
meta.SetStatusCondition(&instance.Status.Conditions, metav1.Condition{
Type: "Ready", Status: ready, Reason: "FixtureObservation", Message: "隔离测试前置观察", ObservedGeneration: instance.Generation,
})
if err := f.api.Status().Update(t.Context(), instance); err != nil {
t.Fatal(err)
}
}
setReady(metav1.ConditionFalse)
stop := startPreparationManager(t, f)
waitForPreparation(t, func() bool {
if err := f.api.Get(t.Context(), client.ObjectKeyFromObject(database), database); err != nil {
return false
}
condition := meta.FindStatusCondition(database.Status.Conditions, "CredentialsReady")
return condition != nil && condition.Reason == "DependencyUnavailable"
})
if database.Status.CredentialRef != nil {
t.Fatal("Instance 未 Ready 时不应固定或写入凭据")
}
setReady(metav1.ConditionTrue)
waitForPreparation(t, func() bool {
return f.api.Get(t.Context(), client.ObjectKeyFromObject(database), database) == nil && database.Status.CredentialVersion == 1
})
stop()
// 重启真实 manager 后继续观察;不得只依赖旧进程中的记忆。
setReady(metav1.ConditionFalse)
stop = startPreparationManager(t, f)
defer stop()
waitForPreparation(t, func() bool {
if err := f.api.Get(t.Context(), client.ObjectKeyFromObject(database), database); err != nil {
return false
}
condition := meta.FindStatusCondition(database.Status.Conditions, "CredentialsReady")
return condition != nil && condition.Status == metav1.ConditionFalse && condition.Reason == "DependencyUnavailable"
})
setReady(metav1.ConditionTrue)
waitForPreparation(t, func() bool {
if err := f.api.Get(t.Context(), client.ObjectKeyFromObject(database), database); err != nil {
return false
}
condition := meta.FindStatusCondition(database.Status.Conditions, "CredentialsReady")
return condition != nil && condition.Status == metav1.ConditionTrue && database.Status.CredentialVersion == 1
})
stored, err := f.bao.KVv2("secret").Get(t.Context(), database.Status.CredentialRef.Path)
if err != nil || stored.VersionMetadata.Version != 1 {
t.Fatal("manager 重启不得重建凭据")
}
}
func startPreparationManager(t *testing.T, f *preparationFixture) func() {
t.Helper()
manager, err := ctrl.NewManager(f.config, ctrl.Options{
Scheme: f.scheme, Metrics: metricsserver.Options{BindAddress: "0"}, HealthProbeBindAddress: "0",
// 测试在同一进程顺序重启 manager;旧 worker 已停止,但名称注册表仍是进程级。
Controller: controllerconfig.Controller{SkipNameValidation: new(true)},
})
if err != nil {
t.Fatal(err)
}
bindingResources := &kubernetes.BindingResources{Client: manager.GetClient(), Reader: manager.GetAPIReader()}
binder := databasecontroller.NewBindingReconciler(manager.GetClient(), &application.BindingService{Resources: bindingResources}, bindingResources)
if err := binder.SetupWithManager(t.Context(), manager); err != nil {
t.Fatal(err)
}
credentialResources := &kubernetes.CredentialResources{Client: manager.GetClient(), Reader: manager.GetAPIReader()}
preparation := &application.CredentialPreparation{Resources: credentialResources, Store: fixtureStore(t, f.bao)}
if err := databasecontroller.NewCredentialReconciler(manager.GetClient(), preparation).SetupWithManager(manager); err != nil {
t.Fatal(err)
}
ctx, cancel := context.WithCancel(t.Context())
done := make(chan error, 1)
go func() { done <- manager.Start(ctx) }()
stopped := false
stop := func() {
if stopped {
return
}
stopped = true
cancel()
select {
case err := <-done:
if err != nil {
t.Error(err)
}
case <-time.After(20 * time.Second):
t.Error("凭据 manager 未停止")
}
}
t.Cleanup(stop)
return stop
}
func waitForPreparation(t *testing.T, check func() bool) {
t.Helper()
deadline := time.Now().Add(10 * time.Second)
for time.Now().Before(deadline) {
if check() {
return
}
time.Sleep(20 * time.Millisecond)
}
t.Fatal("凭据 watch 没有在低频重查之前驱动协调")
}
@@ -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,7 +60,17 @@ func TestInstanceControllerWithRealPostgreSQL(t *testing.T) {
if err != nil { if err != nil {
t.Fatal(err) t.Fatal(err)
} }
reconciler := &databasecontroller.InstanceReconciler{Observer: service, SecretNamespace: controllerNamespace} 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)
}
resources := &secretadapter.InstanceResources{Client: manager.GetClient(), Reader: manager.GetAPIReader()}
usecase := &application.InstanceReconciliation{Resources: resources, Observer: service}
reconciler := databasecontroller.NewInstanceReconciler(manager.GetClient(), usecase, resources, controllerNamespace)
if err := reconciler.SetupWithManager(manager); err != nil { if err := reconciler.SetupWithManager(manager); err != nil {
t.Fatal(err) t.Fatal(err)
} }
@@ -0,0 +1,153 @@
package application
import (
"context"
"errors"
"fmt"
credentialdomain "git.ddupan.top/panxiao81/ayatori/internal/database/domain/credential"
)
var (
ErrCredentialLocation = errors.New("credential location is outside the configured scope")
ErrCredentialUnavailable = errors.New("credential backend unavailable")
ErrCredentialNotFound = errors.New("application credential not found")
ErrCredentialConflict = errors.New("credential creation requires manual conflict resolution")
ErrCredentialUncertain = errors.New("credential creation outcome is uncertain; manual resolution required")
)
// CredentialStore 只表达本用例需要的凭据操作,不提供覆盖或删除。
// version=0 的读取只用于观察是否已有值,成功不能作为认领依据。
type CredentialStore interface {
ProvisionLocation(string) (credentialdomain.Location, error)
ReadCredential(context.Context, credentialdomain.Location, int64) (credentialdomain.ApplicationCredential, error)
CreateCredential(context.Context, credentialdomain.Location, credentialdomain.ApplicationCredential) (int64, error)
}
// CredentialRecord 是同一轮观察的事实,状态中永远不保存密码。
type CredentialRecord struct {
credentialdomain.Target
Revision string
TenantGeneration int64
InstanceGeneration int64
Status credentialdomain.State
}
// CredentialResources 的写入必须检查 Database UID/resourceVersion,保留其他状态。
// CheckCurrent 在外部操作前后回读本轮三个资源,拒绝陈旧快照;它不是跨系统事务。
type CredentialResources interface {
Load(context.Context, string) (*CredentialRecord, error)
Save(context.Context, *CredentialRecord, credentialdomain.State) (*CredentialRecord, error)
CheckCurrent(context.Context, *CredentialRecord) error
}
type CredentialPreparation struct {
Resources CredentialResources
Store CredentialStore
}
func (s CredentialPreparation) Reconcile(ctx context.Context, name string) error {
record, err := s.Resources.Load(ctx, name)
if err != nil || record == nil || !record.RequiresPreparation() {
return err
}
if state, canContinue := record.Status.Resume(); !canContinue {
if state == record.Status {
return nil
}
return s.report(ctx, record, state.Phase, state.Message)
}
if issue := record.Check(); issue != nil {
return s.report(ctx, record, issue.Phase, issue.Message)
}
location, err := s.Store.ProvisionLocation(record.Database.Identity.UID)
if err != nil {
return s.report(ctx, record, credentialdomain.Unavailable, "凭据存储位置配置无效,未执行外部写入")
}
if issue := record.Status.CheckLocation(location); issue != nil {
return s.report(ctx, record, issue.Phase, issue.Message)
}
if record.Status.Location == nil {
status := record.Status
status.Location = &location
status = status.WithPhase(credentialdomain.Pinned, "凭据位置已固定,等待创建")
record, err = s.Resources.Save(ctx, record, status)
if err != nil {
return err
}
}
if record.Status.Confirmed() {
return s.observe(ctx, record)
}
return s.create(ctx, record)
}
func (s CredentialPreparation) create(ctx context.Context, record *CredentialRecord) error {
_, err := s.Store.ReadCredential(ctx, *record.Status.Location, 0)
if err == nil || errors.Is(err, credentialdomain.ErrApplicationCredentialInvalid) {
return s.report(ctx, record, credentialdomain.Conflict, "固定位置已有未确认的凭据;请人工核实,未认领或覆盖")
}
if !errors.Is(err, ErrCredentialNotFound) {
return s.report(ctx, record, credentialdomain.Unavailable, "创建前无法确认凭据位置是否为空,等待依赖恢复")
}
credential, err := credentialdomain.GenerateApplicationCredential(record.Database.LoginRole, record.Database.Name, record.Instance.Endpoint)
if err != nil {
return s.report(ctx, record, credentialdomain.InvalidTarget, "应用凭据目标无效,未执行外部写入")
}
status := record.Status
status = status.WithPhase(credentialdomain.Creating, "凭据创建已开始;尚无成功确认时不得重入创建")
record, err = s.Resources.Save(ctx, record, status)
if err != nil {
return err
}
if err := s.Resources.CheckCurrent(ctx, record); err != nil {
return err
}
version, err := s.Store.CreateCredential(ctx, *record.Status.Location, credential)
if errors.Is(err, ErrCredentialUnavailable) {
// 适配器只在明确未执行写入(认证拒绝或请求前取消)时返回此错误。
return s.report(ctx, record, credentialdomain.Unavailable, "凭据创建在执行前被拒绝,等待认证或权限恢复")
}
if err != nil {
return s.report(ctx, record, credentialdomain.Conflict,
"凭据创建冲突或结果不确定;请核对固定位置的版本历史,未认领、覆盖或重新生成密码")
}
confirmed, issue := record.Status.Created(version)
if issue != nil {
return s.report(ctx, record, issue.Phase, issue.Message)
}
if err := s.Resources.CheckCurrent(ctx, record); err != nil {
return err
}
_, err = s.Resources.Save(ctx, record, confirmed)
return err
}
func (s CredentialPreparation) observe(ctx context.Context, record *CredentialRecord) error {
credential, err := s.Store.ReadCredential(ctx, *record.Status.Location, record.Status.Version)
if errors.Is(err, ErrCredentialConflict) || errors.Is(err, ErrCredentialNotFound) {
return s.report(ctx, record, credentialdomain.Conflict, "已确认凭据消失、版本变化或内容无效;请人工核实,未生成替代密码")
}
if err != nil {
return s.report(ctx, record, credentialdomain.Unavailable, "已确认凭据暂时无法读取;保留确认版本,等待依赖恢复")
}
if issue := record.CheckCredential(credential); issue != nil {
return s.report(ctx, record, issue.Phase, issue.Message)
}
if err := s.Resources.CheckCurrent(ctx, record); err != nil {
return err
}
status := record.Status
status.Phase = credentialdomain.Prepared
status.Message = "已确认凭据可读取;尚未验证 PostgreSQL 资源或完成 Tenant 交付"
_, err = s.Resources.Save(ctx, record, status)
return err
}
func (s CredentialPreparation) report(ctx context.Context, record *CredentialRecord, phase credentialdomain.Phase, message string) error {
status := record.Status
status.Phase = phase
status.Message = fmt.Sprintf("Database %s:%s", record.Database.Identity.Name, message)
_, err := s.Resources.Save(ctx, record, status)
return err
}
@@ -0,0 +1,170 @@
package application
import (
"context"
"errors"
"testing"
credentialdomain "git.ddupan.top/panxiao81/ayatori/internal/database/domain/credential"
"git.ddupan.top/panxiao81/ayatori/internal/database/domain/binding"
"git.ddupan.top/panxiao81/ayatori/internal/database/domain/instance"
)
const (
preparationInstanceName = "credential-instance"
preparationInstanceUID = "credential-instance-uid"
)
func preparationRecord(t *testing.T) *CredentialRecord {
t.Helper()
tenant := binding.TenantIdentity{Namespace: bindingTestNamespace, Name: bindingTestName, UID: "credential-tenant-uid"}
database := binding.Identity{Name: binding.DynamicDatabaseName(tenant.UID), UID: "credential-database-uid"}
endpoint, err := instance.NewEndpoint(instance.EndpointValues{
Host: "postgres.example", HostAddr: "192.0.2.1", Port: 5432, ManagementDatabase: "management", TLSMode: instance.TLSVerifyFull,
})
if err != nil {
t.Fatal(err)
}
return &CredentialRecord{
Database: binding.Database{
Identity: database, Instance: preparationInstanceName, InstanceUID: preparationInstanceUID, Name: bindingTestName, LoginRole: bindingTestName, Source: "Provision", Tenant: &tenant,
},
Tenant: &binding.Tenant{
Identity: tenant, Phase: binding.Bound, Database: &database,
Request: binding.Request{Provision: &binding.ProvisionRequest{Instance: preparationInstanceName}},
},
Instance: &credentialdomain.Instance{Identity: binding.Identity{Name: preparationInstanceName, UID: preparationInstanceUID}, Ready: true, Endpoint: endpoint},
DatabaseProtected: true, TenantProtected: true,
}
}
type memoryCredentialResources struct {
record *CredentialRecord
saveError error
checkError error
}
func (r *memoryCredentialResources) Load(context.Context, string) (*CredentialRecord, error) {
copy := *r.record
return &copy, nil
}
func (r *memoryCredentialResources) Save(_ context.Context, record *CredentialRecord, status credentialdomain.State) (*CredentialRecord, error) {
if r.saveError != nil {
return nil, r.saveError
}
copy := *record
copy.Status = status
r.record = &copy
return &copy, nil
}
func (r *memoryCredentialResources) CheckCurrent(context.Context, *CredentialRecord) error {
return r.checkError
}
type preparationStore struct {
reads int
creates int
readError error
createError error
}
func (*preparationStore) ProvisionLocation(uid string) (credentialdomain.Location, error) {
return credentialdomain.Location{Mount: "applications", Path: "database/" + uid}, nil
}
func (s *preparationStore) ReadCredential(context.Context, credentialdomain.Location, int64) (credentialdomain.ApplicationCredential, error) {
s.reads++
return credentialdomain.ApplicationCredential{}, s.readError
}
func (s *preparationStore) CreateCredential(context.Context, credentialdomain.Location, credentialdomain.ApplicationCredential) (int64, error) {
s.creates++
if s.createError != nil {
return 0, s.createError
}
return 1, nil
}
func TestCredentialPreparationGates(t *testing.T) {
for _, test := range []struct {
name string
change func(*CredentialRecord)
}{
{"导入不供应", func(r *CredentialRecord) { r.Database.Source = "Import" }},
{"删除中的资源", func(r *CredentialRecord) { r.Database.Deleting = true }},
{"已释放资源", func(r *CredentialRecord) { r.Database.Phase = "Released" }},
{"缺少资源保护", func(r *CredentialRecord) { r.DatabaseProtected = false }},
{"缺少申请保护", func(r *CredentialRecord) { r.TenantProtected = false }},
{"单向绑定", func(r *CredentialRecord) { r.Tenant.Database = nil }},
{"删除中的申请", func(r *CredentialRecord) { r.Tenant.Deleting = true }},
{"申请尚未Bound", func(r *CredentialRecord) { r.Tenant.Phase = binding.Binding }},
{"旧申请身份", func(r *CredentialRecord) { r.Tenant.Identity.UID = "new-tenant" }},
{"旧资源身份", func(r *CredentialRecord) { r.Tenant.Database.UID = "new-database" }},
{"Instance未出现", func(r *CredentialRecord) { r.Instance = nil }},
{"Instance正在删除", func(r *CredentialRecord) { r.Instance.Deleting = true }},
{"Instance未Ready", func(r *CredentialRecord) { r.Instance.Ready = false }},
{"Instance身份未记录", func(r *CredentialRecord) { r.Database.InstanceUID = "" }},
{"Instance同名重建", func(r *CredentialRecord) { r.Instance.Identity.UID = "new-instance" }},
{"目标不一致", func(r *CredentialRecord) { r.Database.LoginRole = "another_owner" }},
} {
t.Run(test.name, func(t *testing.T) {
record := preparationRecord(t)
test.change(record)
resources := &memoryCredentialResources{record: record}
store := &preparationStore{readError: ErrCredentialNotFound}
if err := (CredentialPreparation{Resources: resources, Store: store}).Reconcile(t.Context(), record.Database.Identity.Name); err != nil {
t.Fatal(err)
}
if store.reads != 0 || store.creates != 0 || resources.record.Status.Version != 0 {
t.Fatal("前置条件不满足时不得读取、创建或确认凭据")
}
})
}
}
func TestCredentialPreparationWriteBoundary(t *testing.T) {
for _, test := range []struct {
name string
saveError bool
stale bool
createError error
wantError bool
wantCreates int
wantVersion int64
}{
{name: "位置无法保存", saveError: true, wantError: true},
{name: "外部操作前快照变化", stale: true, wantError: true},
{name: "明确权限拒绝", createError: ErrCredentialUnavailable, wantCreates: 1},
{name: "创建结果不确定", createError: ErrCredentialUncertain, wantCreates: 1},
{name: "并发创建冲突", createError: ErrCredentialConflict, wantCreates: 1},
{name: "创建与回读成功", wantCreates: 1, wantVersion: 1},
} {
t.Run(test.name, func(t *testing.T) {
resources := &memoryCredentialResources{record: preparationRecord(t)}
if test.saveError {
resources.saveError = errors.New("fixture write failure")
}
if test.stale {
resources.checkError = errors.New("fixture stale observation")
}
store := &preparationStore{readError: ErrCredentialNotFound, createError: test.createError}
service := CredentialPreparation{Resources: resources, Store: store}
err := service.Reconcile(t.Context(), resources.record.Database.Identity.Name)
if (err != nil) != test.wantError || store.creates != test.wantCreates || resources.record.Status.Version != test.wantVersion {
t.Fatal("外部写入边界或确认时机不符合预期")
}
if test.createError == ErrCredentialUncertain || test.stale {
resources.checkError = nil
if err := service.Reconcile(t.Context(), resources.record.Database.Identity.Name); err != nil {
t.Fatal(err)
}
if store.creates != test.wantCreates || resources.record.Status.Phase != credentialdomain.Conflict {
t.Fatal("未确认创建重入时不得生成替代密码")
}
}
})
}
}
@@ -5,7 +5,6 @@ import (
"context" "context"
"time" "time"
"git.ddupan.top/panxiao81/ayatori/internal/database/adapter/kubernetes"
"git.ddupan.top/panxiao81/ayatori/internal/database/application" "git.ddupan.top/panxiao81/ayatori/internal/database/application"
ctrl "sigs.k8s.io/controller-runtime" ctrl "sigs.k8s.io/controller-runtime"
"sigs.k8s.io/controller-runtime/pkg/client" "sigs.k8s.io/controller-runtime/pkg/client"
@@ -15,7 +14,16 @@ const dependencyRetry = 30 * time.Second
type BindingReconciler struct { type BindingReconciler struct {
Client client.Client Client client.Client
Reader client.Reader Service *application.BindingService
Presenter BindingPresenter
}
type BindingPresenter interface {
Present(context.Context, application.BindingResult) error
}
func NewBindingReconciler(cache client.Client, service *application.BindingService, presenter BindingPresenter) *BindingReconciler {
return &BindingReconciler{Client: cache, Service: service, Presenter: presenter}
} }
// +kubebuilder:rbac:groups=database.ayatori.ddupan.top,resources=postgresqltenants,verbs=get;list;watch;update;patch // +kubebuilder:rbac:groups=database.ayatori.ddupan.top,resources=postgresqltenants,verbs=get;list;watch;update;patch
@@ -27,13 +35,11 @@ type BindingReconciler struct {
// +kubebuilder:rbac:groups=database.ayatori.ddupan.top,resources=postgresqlinstances,verbs=get;list;watch // +kubebuilder:rbac:groups=database.ayatori.ddupan.top,resources=postgresqlinstances,verbs=get;list;watch
func (r *BindingReconciler) Reconcile(ctx context.Context, request ctrl.Request) (ctrl.Result, error) { func (r *BindingReconciler) Reconcile(ctx context.Context, request ctrl.Request) (ctrl.Result, error) {
resources := &kubernetes.BindingResources{Client: r.Client, Reader: r.Reader} result, err := r.Service.Reconcile(ctx, request.Namespace, request.Name)
service := application.BindingService{Resources: resources}
result, err := service.Reconcile(ctx, request.Namespace, request.Name)
if err != nil { if err != nil {
return ctrl.Result{}, err return ctrl.Result{}, err
} }
if err := resources.Present(ctx, result); err != nil { if err := r.Presenter.Present(ctx, result); err != nil {
return ctrl.Result{}, err return ctrl.Result{}, err
} }
if result.RetrySoon { if result.RetrySoon {
@@ -42,6 +42,11 @@ func targetDatabaseName(tenant *databasev1alpha1.PostgreSQLTenant) string {
return "tenant-" + string(tenant.UID) return "tenant-" + string(tenant.UID)
} }
func bindingTestReconciler(writer client.Client, reader client.Reader) *BindingReconciler {
resources := &kubernetes.BindingResources{Client: writer, Reader: reader}
return NewBindingReconciler(writer, &application.BindingService{Resources: resources}, resources)
}
func tenantReference(tenant *databasev1alpha1.PostgreSQLTenant) *databasev1alpha1.TenantReference { func tenantReference(tenant *databasev1alpha1.PostgreSQLTenant) *databasev1alpha1.TenantReference {
return &databasev1alpha1.TenantReference{ return &databasev1alpha1.TenantReference{
Namespace: tenant.Namespace, Name: databasev1alpha1.ObjectName(tenant.Name), UID: tenant.UID, Namespace: tenant.Namespace, Name: databasev1alpha1.ObjectName(tenant.Name), UID: tenant.UID,
@@ -103,7 +108,7 @@ func testDynamicBinding(t *testing.T, apiClient client.Client) {
instance := readyInstance(t, apiClient, "dynamic-instance") instance := readyInstance(t, apiClient, "dynamic-instance")
tenant := provisionTenant("dynamic", instance.Name) tenant := provisionTenant("dynamic", instance.Name)
requireCreate(t, apiClient, tenant) requireCreate(t, apiClient, tenant)
reconciler := &BindingReconciler{Client: apiClient, Reader: apiClient} reconciler := bindingTestReconciler(apiClient, apiClient)
reconcileOK(t, reconciler, tenant) reconcileOK(t, reconciler, tenant)
reload(t, apiClient, tenant) reload(t, apiClient, tenant)
if tenant.Status.DatabaseRef == nil || tenant.Status.Phase != phaseBound { if tenant.Status.DatabaseRef == nil || tenant.Status.Phase != phaseBound {
@@ -154,7 +159,7 @@ func testBindingRestart(t *testing.T, apiClient client.Client) {
instance := readyInstance(t, apiClient, "restart-instance") instance := readyInstance(t, apiClient, "restart-instance")
tenant := provisionTenant("restart", instance.Name) tenant := provisionTenant("restart", instance.Name)
requireCreate(t, apiClient, tenant) requireCreate(t, apiClient, tenant)
first := &BindingReconciler{Client: &failedTenantStatusClient{Client: apiClient}, Reader: apiClient} first := bindingTestReconciler(&failedTenantStatusClient{Client: apiClient}, apiClient)
if _, err := first.Reconcile(t.Context(), ctrl.Request{NamespacedName: client.ObjectKeyFromObject(tenant)}); err == nil { if _, err := first.Reconcile(t.Context(), ctrl.Request{NamespacedName: client.ObjectKeyFromObject(tenant)}); err == nil {
t.Fatal("预期第二次绑定写入失败") t.Fatal("预期第二次绑定写入失败")
} }
@@ -169,7 +174,7 @@ func testBindingRestart(t *testing.T, apiClient client.Client) {
t.Fatal("失败后资源侧绑定不应回滚") t.Fatal("失败后资源侧绑定不应回滚")
} }
// 新建 reconciler,无旧内存,只从 API 中读取进度。 // 新建 reconciler,无旧内存,只从 API 中读取进度。
restarted := &BindingReconciler{Client: apiClient, Reader: apiClient} restarted := bindingTestReconciler(apiClient, apiClient)
reconcileOK(t, restarted, tenant) reconcileOK(t, restarted, tenant)
reload(t, apiClient, tenant) reload(t, apiClient, tenant)
if tenant.Status.DatabaseRef == nil || tenant.Status.DatabaseRef.UID != database.UID { if tenant.Status.DatabaseRef == nil || tenant.Status.DatabaseRef.UID != database.UID {
@@ -190,7 +195,7 @@ func testConcurrentBinding(t *testing.T, apiClient client.Client) {
results := make(chan error, len(tenants)) results := make(chan error, len(tenants))
for _, tenant := range tenants { for _, tenant := range tenants {
workers.Go(func() { workers.Go(func() {
reconciler := &BindingReconciler{Client: apiClient, Reader: apiClient} reconciler := bindingTestReconciler(apiClient, apiClient)
_, err := reconciler.Reconcile(t.Context(), ctrl.Request{NamespacedName: client.ObjectKeyFromObject(tenant)}) _, err := reconciler.Reconcile(t.Context(), ctrl.Request{NamespacedName: client.ObjectKeyFromObject(tenant)})
results <- err results <- err
}) })
@@ -202,7 +207,7 @@ func testConcurrentBinding(t *testing.T, apiClient client.Client) {
t.Fatalf("并发协调出现非版本冲突错误: %v", err) t.Fatalf("并发协调出现非版本冲突错误: %v", err)
} }
} }
reconciler := &BindingReconciler{Client: apiClient, Reader: apiClient} reconciler := bindingTestReconciler(apiClient, apiClient)
bound := 0 bound := 0
for _, tenant := range tenants { for _, tenant := range tenants {
reconcileOK(t, reconciler, tenant) reconcileOK(t, reconciler, tenant)
@@ -233,7 +238,7 @@ func testBindingIdentity(t *testing.T, apiClient client.Client) {
} }
tenant := existingTenant("identity", database.Name) tenant := existingTenant("identity", database.Name)
requireCreate(t, apiClient, tenant) requireCreate(t, apiClient, tenant)
reconciler := &BindingReconciler{Client: apiClient, Reader: apiClient} reconciler := bindingTestReconciler(apiClient, apiClient)
reconcileOK(t, reconciler, tenant) reconcileOK(t, reconciler, tenant)
reload(t, apiClient, tenant) reload(t, apiClient, tenant)
assertNotReady(t, tenant, reasonConflict) assertNotReady(t, tenant, reasonConflict)
@@ -256,7 +261,7 @@ func testBindingIdentity(t *testing.T, apiClient client.Client) {
func testBindingProtection(t *testing.T, apiClient client.Client) { func testBindingProtection(t *testing.T, apiClient client.Client) {
tenant := provisionTenant("protection", "missing-instance") tenant := provisionTenant("protection", "missing-instance")
requireCreate(t, apiClient, tenant) requireCreate(t, apiClient, tenant)
reconciler := &BindingReconciler{Client: apiClient, Reader: apiClient} reconciler := bindingTestReconciler(apiClient, apiClient)
reconcileOK(t, reconciler, tenant) reconcileOK(t, reconciler, tenant)
reload(t, apiClient, tenant) reload(t, apiClient, tenant)
assertNotReady(t, tenant, reasonDependency) assertNotReady(t, tenant, reasonDependency)
@@ -292,7 +297,7 @@ func testStaleObservation(t *testing.T, apiClient client.Client) {
} }
tenant := existingTenant("stale", database.Name) tenant := existingTenant("stale", database.Name)
requireCreate(t, apiClient, tenant) requireCreate(t, apiClient, tenant)
reconciler := &BindingReconciler{Client: apiClient, Reader: apiClient} reconciler := bindingTestReconciler(apiClient, apiClient)
reconcileOK(t, reconciler, tenant) reconcileOK(t, reconciler, tenant)
reload(t, apiClient, tenant) reload(t, apiClient, tenant)
assertNotReady(t, tenant, reasonDependency) assertNotReady(t, tenant, reasonDependency)
@@ -318,7 +323,7 @@ func testBindingWatch(t *testing.T, apiClient client.Client, config *rest.Config
if err != nil { if err != nil {
t.Fatal(err) t.Fatal(err)
} }
reconciler := &BindingReconciler{} reconciler := bindingTestReconciler(manager.GetClient(), manager.GetAPIReader())
if err := reconciler.SetupWithManager(t.Context(), manager); err != nil { if err := reconciler.SetupWithManager(t.Context(), manager); err != nil {
t.Fatal(err) t.Fatal(err)
} }
@@ -373,7 +378,7 @@ func testPresentationVersion(t *testing.T, apiClient client.Client) {
if err := resources.Present(t.Context(), result); !apierrors.IsConflict(err) { if err := resources.Present(t.Context(), result); !apierrors.IsConflict(err) {
t.Fatalf("过期结果呈现 = %v, want Conflict", err) t.Fatalf("过期结果呈现 = %v, want Conflict", err)
} }
reconcileOK(t, &BindingReconciler{Client: apiClient, Reader: apiClient}, tenant) reconcileOK(t, bindingTestReconciler(apiClient, apiClient), tenant)
reload(t, apiClient, tenant) reload(t, apiClient, tenant)
if tenant.Status.Phase != phaseBound || tenant.Spec.SecretName != "updated-delivery" || if tenant.Status.Phase != phaseBound || tenant.Spec.SecretName != "updated-delivery" ||
tenant.Annotations["example.test/keep"] != "preserved" { tenant.Annotations["example.test/keep"] != "preserved" {
@@ -2,9 +2,10 @@ package controller
import ( import (
"context" "context"
"errors"
databasev1alpha1 "git.ddupan.top/panxiao81/ayatori/api/database/v1alpha1" databasev1alpha1 "git.ddupan.top/panxiao81/ayatori/api/database/v1alpha1"
"git.ddupan.top/panxiao81/ayatori/internal/database/adapter/kubernetes" "git.ddupan.top/panxiao81/ayatori/internal/database/domain/binding"
"k8s.io/apimachinery/pkg/types" "k8s.io/apimachinery/pkg/types"
ctrl "sigs.k8s.io/controller-runtime" ctrl "sigs.k8s.io/controller-runtime"
"sigs.k8s.io/controller-runtime/pkg/client" "sigs.k8s.io/controller-runtime/pkg/client"
@@ -14,16 +15,16 @@ import (
const targetDatabaseIndex = "database.bindingTarget" const targetDatabaseIndex = "database.bindingTarget"
func (r *BindingReconciler) SetupWithManager(ctx context.Context, manager ctrl.Manager) error { func (r *BindingReconciler) SetupWithManager(ctx context.Context, manager ctrl.Manager) error {
if r.Client == nil { if r.Client == nil || r.Service == nil || r.Presenter == nil {
r.Client = manager.GetClient() return errors.New("binding controller requires injected client, use case and presenter")
}
if r.Reader == nil {
r.Reader = manager.GetAPIReader()
} }
if err := manager.GetFieldIndexer().IndexField(ctx, &databasev1alpha1.PostgreSQLTenant{}, if err := manager.GetFieldIndexer().IndexField(ctx, &databasev1alpha1.PostgreSQLTenant{},
targetDatabaseIndex, func(object client.Object) []string { targetDatabaseIndex, func(object client.Object) []string {
tenant := object.(*databasev1alpha1.PostgreSQLTenant) tenant := object.(*databasev1alpha1.PostgreSQLTenant)
return []string{kubernetes.BindingTargetName(tenant)} if tenant.Spec.DatabaseRef != nil {
return []string{string(tenant.Spec.DatabaseRef.Name)}
}
return []string{binding.DynamicDatabaseName(string(tenant.UID))}
}); err != nil { }); err != nil {
return err return err
} }
@@ -0,0 +1,64 @@
package controller
import (
"context"
"errors"
databasev1alpha1 "git.ddupan.top/panxiao81/ayatori/api/database/v1alpha1"
"git.ddupan.top/panxiao81/ayatori/internal/database/application"
ctrl "sigs.k8s.io/controller-runtime"
"sigs.k8s.io/controller-runtime/pkg/client"
"sigs.k8s.io/controller-runtime/pkg/handler"
)
// CredentialReconciler 只连接事件、用例和重试,不在控制器中编排凭据写入。
type CredentialReconciler struct {
Client client.Client
Service *application.CredentialPreparation
}
func NewCredentialReconciler(cache client.Client, service *application.CredentialPreparation) *CredentialReconciler {
return &CredentialReconciler{Client: cache, Service: service}
}
func (r *CredentialReconciler) Reconcile(ctx context.Context, request ctrl.Request) (ctrl.Result, error) {
if err := r.Service.Reconcile(ctx, request.Name); err != nil {
return ctrl.Result{}, err
}
// Bao 的可用性和版本变化没有 Kubernetes watch;与已有依赖重查保持一致。
return ctrl.Result{RequeueAfter: dependencyRetry}, nil
}
func (r *CredentialReconciler) SetupWithManager(manager ctrl.Manager) error {
if r.Client == nil || r.Service == nil {
return errors.New("credential controller requires injected client and use case")
}
return ctrl.NewControllerManagedBy(manager).
Named("database-credentials").For(&databasev1alpha1.PostgreSQLDatabase{}).
Watches(&databasev1alpha1.PostgreSQLTenant{}, handler.EnqueueRequestsFromMapFunc(r.requestsForTenant)).
Watches(&databasev1alpha1.PostgreSQLInstance{}, handler.EnqueueRequestsFromMapFunc(r.requestsForInstance)).
Complete(r)
}
func (r *CredentialReconciler) requestsForTenant(_ context.Context, object client.Object) []ctrl.Request {
tenant := object.(*databasev1alpha1.PostgreSQLTenant)
if tenant.Status.DatabaseRef == nil {
return nil
}
return []ctrl.Request{{Name: string(tenant.Status.DatabaseRef.Name)}}
}
func (r *CredentialReconciler) requestsForInstance(ctx context.Context, object client.Object) []ctrl.Request {
databases := &databasev1alpha1.PostgreSQLDatabaseList{}
if err := r.Client.List(ctx, databases); err != nil {
ctrl.LoggerFrom(ctx).Error(err, "无法映射 Instance 凭据准备事件;等待低频重试")
return nil
}
var requests []ctrl.Request
for _, database := range databases.Items {
if string(database.Spec.InstanceRef.Name) == object.GetName() {
requests = append(requests, ctrl.Request{NamespacedName: client.ObjectKeyFromObject(&database)})
}
}
return requests
}
@@ -4,7 +4,6 @@ import (
"context" "context"
"time" "time"
"git.ddupan.top/panxiao81/ayatori/internal/database/adapter/kubernetes"
"git.ddupan.top/panxiao81/ayatori/internal/database/application" "git.ddupan.top/panxiao81/ayatori/internal/database/application"
ctrl "sigs.k8s.io/controller-runtime" ctrl "sigs.k8s.io/controller-runtime"
"sigs.k8s.io/controller-runtime/pkg/client" "sigs.k8s.io/controller-runtime/pkg/client"
@@ -12,27 +11,33 @@ import (
type InstanceReconciler struct { type InstanceReconciler struct {
Client client.Client Client client.Client
Reader client.Reader Service *application.InstanceReconciliation
Observer application.InstanceObserver Presenter InstancePresenter
SecretNamespace string SecretNamespace string
} }
type InstancePresenter interface {
PresentInstance(context.Context, application.InstanceResult) error
}
func NewInstanceReconciler(cache client.Client, service *application.InstanceReconciliation, presenter InstancePresenter, namespace string) *InstanceReconciler {
return &InstanceReconciler{Client: cache, Service: service, Presenter: presenter, SecretNamespace: namespace}
}
// +kubebuilder:rbac:groups=database.ayatori.ddupan.top,resources=postgresqlinstances,verbs=get;list;watch;update;patch // +kubebuilder:rbac:groups=database.ayatori.ddupan.top,resources=postgresqlinstances,verbs=get;list;watch;update;patch
// +kubebuilder:rbac:groups=database.ayatori.ddupan.top,resources=postgresqlinstances/status,verbs=get;update;patch // +kubebuilder:rbac:groups=database.ayatori.ddupan.top,resources=postgresqlinstances/status,verbs=get;update;patch
// +kubebuilder:rbac:groups=database.ayatori.ddupan.top,resources=postgresqlinstances/finalizers,verbs=update // +kubebuilder:rbac:groups=database.ayatori.ddupan.top,resources=postgresqlinstances/finalizers,verbs=update
// Secret 权限单独声明为 namespace Role,不放入生成的 ClusterRole。 // Secret 权限单独声明为 namespace Role,不放入生成的 ClusterRole。
func (r *InstanceReconciler) Reconcile(ctx context.Context, request ctrl.Request) (ctrl.Result, error) { func (r *InstanceReconciler) Reconcile(ctx context.Context, request ctrl.Request) (ctrl.Result, error) {
resources := &kubernetes.InstanceResources{Client: r.Client, Reader: r.Reader}
service := application.InstanceReconciliation{Resources: resources, Observer: r.Observer}
observationContext, cancel := context.WithTimeout(ctx, 15*time.Second) observationContext, cancel := context.WithTimeout(ctx, 15*time.Second)
defer cancel() defer cancel()
result, err := service.Reconcile(observationContext, request.Name) result, err := r.Service.Reconcile(observationContext, request.Name)
if err != nil { if err != nil {
return ctrl.Result{}, err return ctrl.Result{}, err
} }
// 查询超时后仍用 worker context 保存安全失败结果;manager 停止时不强行写入。 // 查询超时后仍用 worker context 保存安全失败结果;manager 停止时不强行写入。
if err := resources.PresentInstance(ctx, result); err != nil { if err := r.Presenter.PresentInstance(ctx, result); err != nil {
return ctrl.Result{}, err return ctrl.Result{}, err
} }
if result.Record == nil || result.RemoveProtection { if result.Record == nil || result.RemoveProtection {
@@ -52,7 +52,8 @@ func newInstanceReconciler(t *testing.T, apiClient client.Client, backend *insta
t.Fatal(err) t.Fatal(err)
} }
t.Cleanup(service.Close) t.Cleanup(service.Close)
return &InstanceReconciler{Client: apiClient, Reader: apiClient, Observer: service} resources := &kubernetes.InstanceResources{Client: apiClient, Reader: apiClient}
return NewInstanceReconciler(apiClient, &application.InstanceReconciliation{Resources: resources, Observer: service}, resources, "")
} }
func reconcileInstance(t *testing.T, reconciler *InstanceReconciler, object *databasev1alpha1.PostgreSQLInstance) { func reconcileInstance(t *testing.T, reconciler *InstanceReconciler, object *databasev1alpha1.PostgreSQLInstance) {
@@ -164,11 +165,11 @@ func TestInstanceDeletionProtection(t *testing.T) {
t.Fatal(err) t.Fatal(err)
} }
backend.inspect = func() { t.Fatal("删除中不应连接 PostgreSQL") } backend.inspect = func() { t.Fatal("删除中不应连接 PostgreSQL") }
reconciler.Reader = &failedReferenceReader{Reader: apiClient} reconciler.Service.Resources.(*kubernetes.InstanceResources).Reader = &failedReferenceReader{Reader: apiClient}
reconcileInstance(t, reconciler, object) reconcileInstance(t, reconciler, object)
reload(t, apiClient, object) reload(t, apiClient, object)
assertInstanceReason(t, object, reasonDependency) assertInstanceReason(t, object, reasonDependency)
reconciler.Reader = apiClient reconciler.Service.Resources.(*kubernetes.InstanceResources).Reader = apiClient
reconcileInstance(t, reconciler, object) reconcileInstance(t, reconciler, object)
reload(t, apiClient, object) reload(t, apiClient, object)
assertInstanceReason(t, object, "InstanceInUse") assertInstanceReason(t, object, "InstanceInUse")
@@ -22,15 +22,9 @@ func InstanceCacheOptions(namespace string) cache.Options {
} }
func (r *InstanceReconciler) SetupWithManager(manager ctrl.Manager) error { func (r *InstanceReconciler) SetupWithManager(manager ctrl.Manager) error {
if r.Observer == nil || len(validation.IsDNS1123Label(r.SecretNamespace)) != 0 { if r.Client == nil || r.Service == nil || r.Presenter == nil || len(validation.IsDNS1123Label(r.SecretNamespace)) != 0 {
return errors.New("instance observer and valid management Secret namespace required") return errors.New("instance observer and valid management Secret namespace required")
} }
if r.Client == nil {
r.Client = manager.GetClient()
}
if r.Reader == nil {
r.Reader = manager.GetAPIReader()
}
return ctrl.NewControllerManagedBy(manager). return ctrl.NewControllerManagedBy(manager).
Named("database-instance"). Named("database-instance").
For(&databasev1alpha1.PostgreSQLInstance{}). For(&databasev1alpha1.PostgreSQLInstance{}).
@@ -0,0 +1,123 @@
/*
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 credential
import (
"crypto/rand"
"encoding/base64"
"errors"
"regexp"
"strconv"
"git.ddupan.top/panxiao81/ayatori/internal/database/domain/instance"
)
var ErrApplicationCredentialInvalid = errors.New("application credential is invalid")
var applicationIdentifier = regexp.MustCompile(`^[a-z][a-z0-9_]{0,62}$`)
// ApplicationCredential 是内存中的应用连接凭据,不得放入 CR 或普通日志。
// 它与 Instance 管理凭据分开,固定输出交付合同中的七键,不生成带密码的 URI。
type ApplicationCredential struct {
username string
password string
database string
endpoint instance.Endpoint
}
func NewApplicationCredential(username, password, database string, endpoint instance.Endpoint) (ApplicationCredential, error) {
if !applicationIdentifier.MatchString(username) || !applicationIdentifier.MatchString(database) || password == "" {
return ApplicationCredential{}, ErrApplicationCredentialInvalid
}
if endpoint.Validate() != nil {
return ApplicationCredential{}, ErrApplicationCredentialInvalid
}
return ApplicationCredential{
username: username,
password: password,
database: database,
endpoint: endpoint,
}, nil
}
// GenerateApplicationCredential 仅供已获准首次创建凭据的供应步骤调用。
// 不能在读取失败、写入结果不确定或重启后无条件重新调用。
func GenerateApplicationCredential(username, database string, endpoint instance.Endpoint) (ApplicationCredential, error) {
password := make([]byte, 32)
rand.Read(password)
return NewApplicationCredential(username, base64.RawURLEncoding.EncodeToString(password), database, endpoint)
}
func (c ApplicationCredential) String() string { return "[redacted application credential]" }
func (c ApplicationCredential) GoString() string { return c.String() }
func (c ApplicationCredential) MarshalJSON() ([]byte, error) {
return []byte(`"[redacted application credential]"`), nil
}
// SecretData 只在凭据后端或数据库连接边界使用;返回值包含明文密码,禁止记录日志。
// 每次返回独立 map,调用方不能修改已经构造的凭据。
func (c ApplicationCredential) SecretData() map[string]any {
endpoint := c.endpoint.Values()
return map[string]any{
"username": c.username,
"password": c.password,
"database": c.database,
"host": endpoint.Host,
"hostaddr": endpoint.HostAddr,
"port": strconv.Itoa(endpoint.Port),
"sslmode": string(endpoint.TLSMode),
}
}
func (c ApplicationCredential) Validate() error {
_, err := NewApplicationCredential(c.username, c.password, c.database, c.endpoint)
return err
}
// MatchesTarget 只比较连接目标,不向用例暴露密码;管理库名不是应用连接目标的一部分。
func (c ApplicationCredential) MatchesTarget(username, database string, endpoint instance.Endpoint) bool {
actual, wanted := c.endpoint.Values(), endpoint.Values()
return c.username == username && c.database == database && actual.Host == wanted.Host &&
actual.HostAddr == wanted.HostAddr && actual.Port == wanted.Port && actual.TLSMode == wanted.TLSMode
}
// ParseApplicationCredential 拒绝缺键、非字符串或非法连接参数,不回显后端内容。
func ParseApplicationCredential(data map[string]any) (ApplicationCredential, error) {
values := make(map[string]string, 7)
for _, key := range []string{"username", "password", "database", "host", "hostaddr", "port", "sslmode"} {
value, ok := data[key].(string)
if !ok || value == "" {
return ApplicationCredential{}, ErrApplicationCredentialInvalid
}
values[key] = value
}
port, err := strconv.Atoi(values["port"])
if err != nil {
return ApplicationCredential{}, ErrApplicationCredentialInvalid
}
endpoint, err := instance.NewEndpoint(instance.EndpointValues{
Host: values["host"],
HostAddr: values["hostaddr"],
Port: port,
ManagementDatabase: values["database"],
TLSMode: instance.TLSMode(values["sslmode"]),
})
if err != nil {
return ApplicationCredential{}, ErrApplicationCredentialInvalid
}
return NewApplicationCredential(values["username"], values["password"], values["database"], endpoint)
}
@@ -0,0 +1,82 @@
/*
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 credential_test
import (
"encoding/json"
"fmt"
"maps"
"strings"
"testing"
credentialdomain "git.ddupan.top/panxiao81/ayatori/internal/database/domain/credential"
"git.ddupan.top/panxiao81/ayatori/internal/database/domain/instance"
)
func TestApplicationCredential(t *testing.T) {
endpoint, err := instance.NewEndpoint(instance.EndpointValues{
Host: "postgres.example", HostAddr: "192.0.2.1", Port: 5432,
ManagementDatabase: "postgres", TLSMode: instance.TLSVerifyFull,
})
if err != nil {
t.Fatal(err)
}
first, err := credentialdomain.GenerateApplicationCredential("owner", "app", endpoint)
if err != nil {
t.Fatal(err)
}
second, err := credentialdomain.GenerateApplicationCredential("owner", "app", endpoint)
if err != nil {
t.Fatal(err)
}
data := first.SecretData()
if len(data) != 7 || data["password"] == second.SecretData()["password"] || len(data["password"].(string)) != 43 {
t.Fatal("expected seven keys and independent 256-bit passwords")
}
parsed, err := credentialdomain.ParseApplicationCredential(data)
if err != nil || !maps.Equal(parsed.SecretData(), data) {
t.Fatal("credential did not round trip")
}
encoded, err := json.Marshal(first)
if err != nil {
t.Fatal(err)
}
for _, output := range []string{fmt.Sprint(first), fmt.Sprintf("%+v", first), fmt.Sprintf("%#v", first), string(encoded)} {
if strings.Contains(output, data["password"].(string)) {
t.Fatal("credential formatting leaked the password")
}
}
data["password"] = "changed"
if first.SecretData()["password"] == "changed" {
t.Fatal("caller mutated credential")
}
for key := range data {
invalid := maps.Clone(data)
delete(invalid, key)
if _, err := credentialdomain.ParseApplicationCredential(invalid); err == nil {
t.Fatalf("accepted missing %s", key)
}
invalid[key] = 42
if _, err := credentialdomain.ParseApplicationCredential(invalid); err == nil {
t.Fatalf("accepted non-string %s", key)
}
}
if (credentialdomain.ApplicationCredential{}).Validate() == nil {
t.Fatal("accepted zero credential")
}
}
@@ -0,0 +1,129 @@
package credential
import (
"git.ddupan.top/panxiao81/ayatori/internal/database/domain/binding"
"git.ddupan.top/panxiao81/ayatori/internal/database/domain/instance"
)
type Location struct {
Mount string
Path string
}
// Phase 表达凭据准备进度,不依赖 Kubernetes Condition 的类型或 Reason。
type Phase uint8
const (
Pending Phase = iota
Pinned
Creating
Prepared
Conflict
Unavailable
Stopped
InvalidTarget
)
type State struct {
Location *Location
Version int64
Phase Phase
Message string
}
func (s State) WithPhase(phase Phase, message string) State {
s.Phase, s.Message = phase, message
return s
}
func (s State) Confirmed() bool { return s.Version > 0 }
// Created 只接受首次创建并回读得到的版本,不能把后续写入认作首次供应。
func (s State) Created(version int64) (State, *Issue) {
if version != 1 {
return s, &Issue{Conflict, "凭据创建冲突或结果不确定;请核对固定位置的版本历史,未认领、覆盖或重新生成密码"}
}
s.Version = version
return s.WithPhase(Prepared, "凭据已创建并回读确认;尚未创建 PostgreSQL 资源或交付给 Tenant"), nil
}
// Resume 决定新一轮是否可以继续。未确认的创建不能靠读取成功认领。
func (s State) Resume() (State, bool) {
if s.Version != 0 {
return s, true
}
switch s.Phase {
case Conflict:
return s, false
case Creating:
return s.WithPhase(Conflict, "凭据创建未留下成功确认;请核对固定位置与后端历史并人工处理,未重新生成密码"), false
default:
return s, true
}
}
func (s State) CheckLocation(configured Location) *Issue {
if s.Location != nil && *s.Location != configured {
return &Issue{Unavailable, "部署配置与固定凭据位置不一致;请恢复原 mount/path 配置,未迁移或改密"}
}
return nil
}
type Instance struct {
binding.Instance
Endpoint instance.Endpoint
}
// Target 只包含供应资格所需事实,不含 resourceVersion、Conditions 或 repository 对象。
type Target struct {
Database binding.Database
Tenant *binding.Tenant
Instance *Instance
DatabaseProtected bool
TenantProtected bool
}
type Issue struct {
Phase Phase
Message string
}
func (t Target) RequiresPreparation() bool { return t.Database.Source == "Provision" }
func (t Target) CheckCredential(value ApplicationCredential) *Issue {
if t.Instance == nil || !value.MatchesTarget(t.Database.LoginRole, t.Database.Name, t.Instance.Endpoint) {
return &Issue{Conflict, "已确认凭据与当前 Instance/database/loginRole 不一致;请人工核实,未修改凭据"}
}
return nil
}
func (t Target) Check() *Issue {
database := t.Database
if database.Deleting || database.Phase == binding.Deleting || database.Phase == "Released" {
return &Issue{Stopped, "Database 正在删除或已释放;保留凭据与 finalizer,不执行供应或清理"}
}
if database.Tenant == nil || t.Tenant == nil || t.Tenant.Database == nil {
return &Issue{Unavailable, "等待 Database 与 Tenant 双向绑定完成"}
}
if *database.Tenant != t.Tenant.Identity || *t.Tenant.Database != database.Identity {
return &Issue{Conflict, "双向绑定的名称或 UID 不匹配,未创建凭据"}
}
if t.Tenant.Deleting || t.Tenant.Phase != binding.Bound || !t.DatabaseProtected || !t.TenantProtected {
return &Issue{Stopped, "Tenant 未完成绑定、正在删除或缺少 finalizer 保护,未创建凭据"}
}
request, err := t.Tenant.Request.Resolve(t.Tenant.Identity)
if err != nil || (request.Provision != nil && !database.MatchesProvision(request, t.Tenant.Identity)) || request.Name != database.Identity.Name {
return &Issue{Conflict, "Tenant 申请与 Database 目标不一致,未创建凭据"}
}
if t.Instance == nil || database.InstanceUID == "" {
return &Issue{Unavailable, "等待 Instance 与已记录的实例身份"}
}
if issue := t.Instance.Check(&database); issue != nil {
phase := Unavailable
if issue.Reason == binding.Conflict {
phase = Conflict
}
return &Issue{phase, issue.Message}
}
return nil
}
@@ -0,0 +1,59 @@
package credential_test
import (
"testing"
credential "git.ddupan.top/panxiao81/ayatori/internal/database/domain/credential"
)
func TestPreparationResume(t *testing.T) {
for _, test := range []struct {
name string
phase credential.Phase
version int64
continueAllowed bool
result credential.Phase
}{
{"尚未创建", credential.Pinned, 0, true, credential.Pinned},
{"依赖恢复", credential.Unavailable, 0, true, credential.Unavailable},
{"中断创建", credential.Creating, 0, false, credential.Conflict},
{"未确认冲突", credential.Conflict, 0, false, credential.Conflict},
{"已确认后读取失败", credential.Unavailable, 1, true, credential.Unavailable},
{"已确认后冲突重验", credential.Conflict, 1, true, credential.Conflict},
} {
t.Run(test.name, func(t *testing.T) {
original := credential.State{Phase: test.phase, Version: test.version, Message: "保留原诊断"}
state, allowed := original.Resume()
if allowed != test.continueAllowed || state.Phase != test.result || state.Version != original.Version {
t.Fatal("恢复判定或确认版本发生变化")
}
if test.phase == credential.Conflict && state.Message != original.Message {
t.Fatal("冲突重入应保留原诊断")
}
})
}
}
func TestPreparationLocationAndConfirmation(t *testing.T) {
location := credential.Location{Mount: "applications", Path: "database/uid"}
state := credential.State{Location: &location, Phase: credential.Creating}
if issue := state.CheckLocation(location); issue != nil {
t.Fatal("固定位置不应被拒绝")
}
if issue := state.CheckLocation(credential.Location{Mount: "other", Path: location.Path}); issue == nil || issue.Phase != credential.Unavailable {
t.Fatal("配置变化必须停止,不迁移已固定位置")
}
for _, version := range []int64{0, -1, 2} {
result, issue := state.Created(version)
if issue == nil || issue.Phase != credential.Conflict || result.Confirmed() {
t.Fatal("错误版本不得确认创建")
}
}
result, issue := state.Created(1)
if issue != nil || !result.Confirmed() || result.Phase != credential.Prepared || result.Location != state.Location {
t.Fatal("首次写入回读应确认并保留位置")
}
if state.Version != 0 {
t.Fatal("领域判定不得修改调用方的旧状态")
}
}
@@ -0,0 +1,52 @@
package credential_test
import (
"testing"
"git.ddupan.top/panxiao81/ayatori/internal/database/domain/binding"
credential "git.ddupan.top/panxiao81/ayatori/internal/database/domain/credential"
)
const (
targetTestApplication = "sampleapp"
targetTestInstance = "test-instance"
)
func TestPreparationTarget(t *testing.T) {
for _, test := range []struct {
name string
change func(*credential.Target)
want credential.Phase
}{
{"完整绑定", func(*credential.Target) {}, credential.Pending},
{"单向绑定", func(target *credential.Target) { target.Tenant.Database = nil }, credential.Unavailable},
{"旧租户身份", func(target *credential.Target) { target.Tenant.Identity.UID = "new" }, credential.Conflict},
{"旧实例身份", func(target *credential.Target) { target.Instance.Identity.UID = "new" }, credential.Conflict},
{"资源删除", func(target *credential.Target) { target.Database.Deleting = true }, credential.Stopped},
{"申请删除", func(target *credential.Target) { target.Tenant.Deleting = true }, credential.Stopped},
{"Released", func(target *credential.Target) { target.Database.Phase = "Released" }, credential.Stopped},
{"缺少保护", func(target *credential.Target) { target.DatabaseProtected = false }, credential.Stopped},
{"实例未就绪", func(target *credential.Target) { target.Instance.Ready = false }, credential.Unavailable},
{"目标变化", func(target *credential.Target) { target.Database.LoginRole = "other" }, credential.Conflict},
} {
t.Run(test.name, func(t *testing.T) {
tenantID := binding.TenantIdentity{Namespace: "apps", Name: targetTestApplication, UID: "tenant"}
databaseID := binding.Identity{Name: binding.DynamicDatabaseName(tenantID.UID), UID: "database"}
target := credential.Target{
Database: binding.Database{Identity: databaseID, Tenant: &tenantID, Instance: targetTestInstance, InstanceUID: "instance-id", Name: targetTestApplication, LoginRole: targetTestApplication, Source: "Provision"},
Tenant: &binding.Tenant{Identity: tenantID, Database: &databaseID, Phase: binding.Bound, Request: binding.Request{Provision: &binding.ProvisionRequest{Instance: targetTestInstance}}},
Instance: &credential.Instance{Identity: binding.Identity{Name: targetTestInstance, UID: "instance-id"}, Ready: true},
DatabaseProtected: true, TenantProtected: true,
}
test.change(&target)
issue := target.Check()
if test.want == credential.Pending {
if issue != nil {
t.Fatalf("有效绑定被拒绝: %s", issue.Message)
}
} else if issue == nil || issue.Phase != test.want {
t.Fatal("领域资格判定不符")
}
})
}
}
+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")
}
}