refactor: 集中显式注入并归位凭据领域规则
This commit is contained in:
@@ -5,7 +5,6 @@ import (
|
||||
"context"
|
||||
"time"
|
||||
|
||||
"git.ddupan.top/panxiao81/ayatori/internal/database/adapter/kubernetes"
|
||||
"git.ddupan.top/panxiao81/ayatori/internal/database/application"
|
||||
ctrl "sigs.k8s.io/controller-runtime"
|
||||
"sigs.k8s.io/controller-runtime/pkg/client"
|
||||
@@ -14,8 +13,17 @@ import (
|
||||
const dependencyRetry = 30 * time.Second
|
||||
|
||||
type BindingReconciler struct {
|
||||
Client client.Client
|
||||
Reader client.Reader
|
||||
Client client.Client
|
||||
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
|
||||
@@ -27,13 +35,11 @@ type BindingReconciler struct {
|
||||
// +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) {
|
||||
resources := &kubernetes.BindingResources{Client: r.Client, Reader: r.Reader}
|
||||
service := application.BindingService{Resources: resources}
|
||||
result, err := service.Reconcile(ctx, request.Namespace, request.Name)
|
||||
result, err := r.Service.Reconcile(ctx, request.Namespace, request.Name)
|
||||
if err != nil {
|
||||
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
|
||||
}
|
||||
if result.RetrySoon {
|
||||
|
||||
@@ -42,6 +42,11 @@ func targetDatabaseName(tenant *databasev1alpha1.PostgreSQLTenant) string {
|
||||
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 {
|
||||
return &databasev1alpha1.TenantReference{
|
||||
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")
|
||||
tenant := provisionTenant("dynamic", instance.Name)
|
||||
requireCreate(t, apiClient, tenant)
|
||||
reconciler := &BindingReconciler{Client: apiClient, Reader: apiClient}
|
||||
reconciler := bindingTestReconciler(apiClient, apiClient)
|
||||
reconcileOK(t, reconciler, tenant)
|
||||
reload(t, apiClient, tenant)
|
||||
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")
|
||||
tenant := provisionTenant("restart", instance.Name)
|
||||
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 {
|
||||
t.Fatal("预期第二次绑定写入失败")
|
||||
}
|
||||
@@ -169,7 +174,7 @@ func testBindingRestart(t *testing.T, apiClient client.Client) {
|
||||
t.Fatal("失败后资源侧绑定不应回滚")
|
||||
}
|
||||
// 新建 reconciler,无旧内存,只从 API 中读取进度。
|
||||
restarted := &BindingReconciler{Client: apiClient, Reader: apiClient}
|
||||
restarted := bindingTestReconciler(apiClient, apiClient)
|
||||
reconcileOK(t, restarted, tenant)
|
||||
reload(t, apiClient, tenant)
|
||||
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))
|
||||
for _, tenant := range tenants {
|
||||
workers.Go(func() {
|
||||
reconciler := &BindingReconciler{Client: apiClient, Reader: apiClient}
|
||||
reconciler := bindingTestReconciler(apiClient, apiClient)
|
||||
_, err := reconciler.Reconcile(t.Context(), ctrl.Request{NamespacedName: client.ObjectKeyFromObject(tenant)})
|
||||
results <- err
|
||||
})
|
||||
@@ -202,7 +207,7 @@ func testConcurrentBinding(t *testing.T, apiClient client.Client) {
|
||||
t.Fatalf("并发协调出现非版本冲突错误: %v", err)
|
||||
}
|
||||
}
|
||||
reconciler := &BindingReconciler{Client: apiClient, Reader: apiClient}
|
||||
reconciler := bindingTestReconciler(apiClient, apiClient)
|
||||
bound := 0
|
||||
for _, tenant := range tenants {
|
||||
reconcileOK(t, reconciler, tenant)
|
||||
@@ -233,7 +238,7 @@ func testBindingIdentity(t *testing.T, apiClient client.Client) {
|
||||
}
|
||||
tenant := existingTenant("identity", database.Name)
|
||||
requireCreate(t, apiClient, tenant)
|
||||
reconciler := &BindingReconciler{Client: apiClient, Reader: apiClient}
|
||||
reconciler := bindingTestReconciler(apiClient, apiClient)
|
||||
reconcileOK(t, reconciler, tenant)
|
||||
reload(t, apiClient, tenant)
|
||||
assertNotReady(t, tenant, reasonConflict)
|
||||
@@ -256,7 +261,7 @@ func testBindingIdentity(t *testing.T, apiClient client.Client) {
|
||||
func testBindingProtection(t *testing.T, apiClient client.Client) {
|
||||
tenant := provisionTenant("protection", "missing-instance")
|
||||
requireCreate(t, apiClient, tenant)
|
||||
reconciler := &BindingReconciler{Client: apiClient, Reader: apiClient}
|
||||
reconciler := bindingTestReconciler(apiClient, apiClient)
|
||||
reconcileOK(t, reconciler, tenant)
|
||||
reload(t, apiClient, tenant)
|
||||
assertNotReady(t, tenant, reasonDependency)
|
||||
@@ -292,7 +297,7 @@ func testStaleObservation(t *testing.T, apiClient client.Client) {
|
||||
}
|
||||
tenant := existingTenant("stale", database.Name)
|
||||
requireCreate(t, apiClient, tenant)
|
||||
reconciler := &BindingReconciler{Client: apiClient, Reader: apiClient}
|
||||
reconciler := bindingTestReconciler(apiClient, apiClient)
|
||||
reconcileOK(t, reconciler, tenant)
|
||||
reload(t, apiClient, tenant)
|
||||
assertNotReady(t, tenant, reasonDependency)
|
||||
@@ -318,7 +323,7 @@ func testBindingWatch(t *testing.T, apiClient client.Client, config *rest.Config
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
reconciler := &BindingReconciler{}
|
||||
reconciler := bindingTestReconciler(manager.GetClient(), manager.GetAPIReader())
|
||||
if err := reconciler.SetupWithManager(t.Context(), manager); err != nil {
|
||||
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) {
|
||||
t.Fatalf("过期结果呈现 = %v, want Conflict", err)
|
||||
}
|
||||
reconcileOK(t, &BindingReconciler{Client: apiClient, Reader: apiClient}, tenant)
|
||||
reconcileOK(t, bindingTestReconciler(apiClient, apiClient), tenant)
|
||||
reload(t, apiClient, tenant)
|
||||
if tenant.Status.Phase != phaseBound || tenant.Spec.SecretName != "updated-delivery" ||
|
||||
tenant.Annotations["example.test/keep"] != "preserved" {
|
||||
|
||||
@@ -2,9 +2,10 @@ package controller
|
||||
|
||||
import (
|
||||
"context"
|
||||
"errors"
|
||||
|
||||
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"
|
||||
ctrl "sigs.k8s.io/controller-runtime"
|
||||
"sigs.k8s.io/controller-runtime/pkg/client"
|
||||
@@ -14,16 +15,16 @@ import (
|
||||
const targetDatabaseIndex = "database.bindingTarget"
|
||||
|
||||
func (r *BindingReconciler) SetupWithManager(ctx context.Context, manager ctrl.Manager) error {
|
||||
if r.Client == nil {
|
||||
r.Client = manager.GetClient()
|
||||
}
|
||||
if r.Reader == nil {
|
||||
r.Reader = manager.GetAPIReader()
|
||||
if r.Client == nil || r.Service == nil || r.Presenter == nil {
|
||||
return errors.New("binding controller requires injected client, use case and presenter")
|
||||
}
|
||||
if err := manager.GetFieldIndexer().IndexField(ctx, &databasev1alpha1.PostgreSQLTenant{},
|
||||
targetDatabaseIndex, func(object client.Object) []string {
|
||||
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 {
|
||||
return err
|
||||
}
|
||||
|
||||
@@ -2,9 +2,9 @@ package controller
|
||||
|
||||
import (
|
||||
"context"
|
||||
"errors"
|
||||
|
||||
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"
|
||||
ctrl "sigs.k8s.io/controller-runtime"
|
||||
"sigs.k8s.io/controller-runtime/pkg/client"
|
||||
@@ -13,16 +13,16 @@ import (
|
||||
|
||||
// CredentialReconciler 只连接事件、用例和重试,不在控制器中编排凭据写入。
|
||||
type CredentialReconciler struct {
|
||||
Client client.Client
|
||||
Reader client.Reader
|
||||
Store application.CredentialStore
|
||||
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) {
|
||||
service := application.CredentialPreparation{
|
||||
Resources: &kubernetes.CredentialResources{Client: r.Client, Reader: r.Reader}, Store: r.Store,
|
||||
}
|
||||
if err := service.Reconcile(ctx, request.Name); err != nil {
|
||||
if err := r.Service.Reconcile(ctx, request.Name); err != nil {
|
||||
return ctrl.Result{}, err
|
||||
}
|
||||
// Bao 的可用性和版本变化没有 Kubernetes watch;与已有依赖重查保持一致。
|
||||
@@ -30,11 +30,8 @@ func (r *CredentialReconciler) Reconcile(ctx context.Context, request ctrl.Reque
|
||||
}
|
||||
|
||||
func (r *CredentialReconciler) SetupWithManager(manager ctrl.Manager) error {
|
||||
if r.Client == nil {
|
||||
r.Client = manager.GetClient()
|
||||
}
|
||||
if r.Reader == nil {
|
||||
r.Reader = manager.GetAPIReader()
|
||||
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{}).
|
||||
|
||||
@@ -4,7 +4,6 @@ import (
|
||||
"context"
|
||||
"time"
|
||||
|
||||
"git.ddupan.top/panxiao81/ayatori/internal/database/adapter/kubernetes"
|
||||
"git.ddupan.top/panxiao81/ayatori/internal/database/application"
|
||||
ctrl "sigs.k8s.io/controller-runtime"
|
||||
"sigs.k8s.io/controller-runtime/pkg/client"
|
||||
@@ -12,27 +11,33 @@ import (
|
||||
|
||||
type InstanceReconciler struct {
|
||||
Client client.Client
|
||||
Reader client.Reader
|
||||
Observer application.InstanceObserver
|
||||
Service *application.InstanceReconciliation
|
||||
Presenter InstancePresenter
|
||||
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/status,verbs=get;update;patch
|
||||
// +kubebuilder:rbac:groups=database.ayatori.ddupan.top,resources=postgresqlinstances/finalizers,verbs=update
|
||||
// Secret 权限单独声明为 namespace Role,不放入生成的 ClusterRole。
|
||||
|
||||
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)
|
||||
defer cancel()
|
||||
result, err := service.Reconcile(observationContext, request.Name)
|
||||
result, err := r.Service.Reconcile(observationContext, request.Name)
|
||||
if err != nil {
|
||||
return ctrl.Result{}, err
|
||||
}
|
||||
// 查询超时后仍用 worker context 保存安全失败结果;manager 停止时不强行写入。
|
||||
if err := resources.PresentInstance(ctx, result); err != nil {
|
||||
if err := r.Presenter.PresentInstance(ctx, result); err != nil {
|
||||
return ctrl.Result{}, err
|
||||
}
|
||||
if result.Record == nil || result.RemoveProtection {
|
||||
|
||||
@@ -52,7 +52,8 @@ func newInstanceReconciler(t *testing.T, apiClient client.Client, backend *insta
|
||||
t.Fatal(err)
|
||||
}
|
||||
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) {
|
||||
@@ -164,11 +165,11 @@ func TestInstanceDeletionProtection(t *testing.T) {
|
||||
t.Fatal(err)
|
||||
}
|
||||
backend.inspect = func() { t.Fatal("删除中不应连接 PostgreSQL") }
|
||||
reconciler.Reader = &failedReferenceReader{Reader: apiClient}
|
||||
reconciler.Service.Resources.(*kubernetes.InstanceResources).Reader = &failedReferenceReader{Reader: apiClient}
|
||||
reconcileInstance(t, reconciler, object)
|
||||
reload(t, apiClient, object)
|
||||
assertInstanceReason(t, object, reasonDependency)
|
||||
reconciler.Reader = apiClient
|
||||
reconciler.Service.Resources.(*kubernetes.InstanceResources).Reader = apiClient
|
||||
reconcileInstance(t, reconciler, object)
|
||||
reload(t, apiClient, object)
|
||||
assertInstanceReason(t, object, "InstanceInUse")
|
||||
|
||||
@@ -22,15 +22,9 @@ func InstanceCacheOptions(namespace string) cache.Options {
|
||||
}
|
||||
|
||||
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")
|
||||
}
|
||||
if r.Client == nil {
|
||||
r.Client = manager.GetClient()
|
||||
}
|
||||
if r.Reader == nil {
|
||||
r.Reader = manager.GetAPIReader()
|
||||
}
|
||||
return ctrl.NewControllerManagedBy(manager).
|
||||
Named("database-instance").
|
||||
For(&databasev1alpha1.PostgreSQLInstance{}).
|
||||
|
||||
Reference in New Issue
Block a user