feat: 接入 Instance 原生管理观测与删除保护
Verify / test (pull_request) Successful in 12m41s
Verify / lint (pull_request) Successful in 14m14s
Verify / database-integration (pull_request) Successful in 16m9s

This commit is contained in:
2026-09-25 11:35:04 +00:00
parent 55b269ce2e
commit bc227bfdb4
26 changed files with 1348 additions and 64 deletions
@@ -0,0 +1,45 @@
package kubernetes
import (
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/instance"
)
func instanceRecord(object *databasev1alpha1.PostgreSQLInstance) (*application.InstanceRecord, error) {
identity, err := instance.NewIdentity(string(object.UID), object.Name)
if err != nil {
return nil, err
}
revision, err := instance.NewRevision(object.Generation)
if err != nil {
return nil, err
}
spec := object.Spec
endpoint, err := instance.NewEndpoint(instance.EndpointValues{
Host: spec.Endpoint.Host, HostAddr: spec.Endpoint.HostAddr,
Port: int(spec.Endpoint.Port), ManagementDatabase: string(spec.Endpoint.Database),
TLSMode: instance.TLSMode(spec.Endpoint.SSLMode),
})
if err != nil {
return nil, err
}
credential, err := instance.NewCredentialReference(instance.CredentialReferenceValues{
Name: string(spec.AdminCredentialRef.Name),
UsernameKey: spec.AdminCredentialRef.UsernameKey, PasswordKey: spec.AdminCredentialRef.PasswordKey,
})
if err != nil {
return nil, err
}
definition, err := instance.NewDefinition(endpoint, credential)
if err != nil {
return nil, err
}
target, err := instance.NewObservationTarget(identity, revision, definition)
if err != nil {
return nil, err
}
return &application.InstanceRecord{
Target: target, Revision: object.ResourceVersion, Deleting: !object.DeletionTimestamp.IsZero(),
}, nil
}
@@ -0,0 +1,112 @@
package kubernetes
import (
"context"
"errors"
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/instance"
"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"
"sigs.k8s.io/controller-runtime/pkg/client"
"sigs.k8s.io/controller-runtime/pkg/controller/controllerutil"
)
const InstanceFinalizer = "database.ayatori.ddupan.top/instance-protection"
type InstanceResources struct {
Client client.Client
Reader client.Reader
}
func (r *InstanceResources) LoadInstance(ctx context.Context, name string) (*application.InstanceRecord, error) {
object := &databasev1alpha1.PostgreSQLInstance{}
if err := r.Reader.Get(ctx, client.ObjectKey{Name: name}, object); err != nil {
return nil, client.IgnoreNotFound(err)
}
return instanceRecord(object)
}
func (r *InstanceResources) ProtectInstance(ctx context.Context, record *application.InstanceRecord) (*application.InstanceRecord, error) {
object, err := r.instanceAtVersion(ctx, record)
if err != nil {
return nil, err
}
if controllerutil.AddFinalizer(object, InstanceFinalizer) {
if err := r.Client.Update(ctx, object); err != nil {
return nil, err
}
}
return instanceRecord(object)
}
func (r *InstanceResources) InstanceReferences(ctx context.Context, name string) (string, error) {
// 删除判断必须直读 API;包含 Released、删除中的 Database 和尚未绑定的申请。
// 不按旧 Instance UID 忽略引用,也不依赖 informer 索引的及时性。
databases := &databasev1alpha1.PostgreSQLDatabaseList{}
if err := r.Reader.List(ctx, databases); err != nil {
return "", err
}
for _, database := range databases.Items {
if string(database.Spec.InstanceRef.Name) == name {
return "Database/" + database.Name, nil
}
}
tenants := &databasev1alpha1.PostgreSQLTenantList{}
if err := r.Reader.List(ctx, tenants); err != nil {
return "", err
}
for _, tenant := range tenants.Items {
if tenant.Spec.Provision != nil && string(tenant.Spec.Provision.InstanceRef.Name) == name {
return "Tenant/" + tenant.Namespace + "/" + tenant.Name, nil
}
}
return "", nil
}
func (r *InstanceResources) PresentInstance(ctx context.Context, result application.InstanceResult) error {
if result.Record == nil {
return nil
}
object, err := r.instanceAtVersion(ctx, result.Record)
if err != nil {
return client.IgnoreNotFound(err)
}
previous := object.Status.DeepCopy()
object.Status.Phase = string(result.Snapshot.Phase)
object.Status.ObservedGeneration = object.Generation
object.Status.PostgreSQLVersion = result.Snapshot.ReportedVersion
ready := metav1.ConditionFalse
if result.Snapshot.Readiness == instance.Ready {
ready = metav1.ConditionTrue
}
meta.SetStatusCondition(&object.Status.Conditions, metav1.Condition{
Type: "Ready", Status: ready, ObservedGeneration: object.Generation,
Reason: result.Reason, Message: result.Message,
})
if !equality.Semantic.DeepEqual(*previous, object.Status) {
if err := r.Client.Status().Update(ctx, object); err != nil {
return err
}
}
if result.RemoveProtection && controllerutil.RemoveFinalizer(object, InstanceFinalizer) {
return r.Client.Update(ctx, object)
}
return nil
}
func (r *InstanceResources) instanceAtVersion(ctx context.Context, record *application.InstanceRecord) (*databasev1alpha1.PostgreSQLInstance, error) {
object := &databasev1alpha1.PostgreSQLInstance{}
name := record.Target.Identity().Name()
if err := r.Reader.Get(ctx, client.ObjectKey{Name: name}, object); err != nil {
return nil, err
}
if string(object.UID) != record.Target.Identity().UID() || object.ResourceVersion != record.Revision {
return nil, apierrors.NewConflict(databasev1alpha1.GroupVersion.WithResource("postgresqlinstances").GroupResource(),
name, errors.New("Instance 快照已过期,请重新观察"))
}
return object, nil
}
@@ -31,6 +31,7 @@ import (
corev1 "k8s.io/api/core/v1"
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
"k8s.io/client-go/kubernetes"
"k8s.io/client-go/rest"
"sigs.k8s.io/controller-runtime/pkg/envtest"
secretadapter "git.ddupan.top/panxiao81/ayatori/internal/database/adapter/kubernetes"
@@ -40,15 +41,19 @@ import (
)
const (
fixtureHost = "fixture.invalid"
fixtureUser = "postgres"
fixtureExtension = "plpgsql"
dockerExec = "exec"
fixtureImage = "postgres@sha256:18cfe3ef5e6815560c98237d6216d1e5119702fb0f3894c8785dd58b8bbe5d73"
fixturePassword = "AYATORI-TEST-ONLY-initial-password"
rotatedPassword = "AYATORI-TEST-ONLY-rotated-password"
controllerNamespace = "database-controller"
secretName = "management"
fixtureAddress = "127.0.0.1"
managementUsernameKey = "login"
managementPasswordKey = "credential"
unrelatedNamespace = "unrelated"
fixtureHost = "fixture.invalid"
fixtureUser = "postgres"
fixtureExtension = "plpgsql"
dockerExec = "exec"
fixtureImage = "postgres@sha256:18cfe3ef5e6815560c98237d6216d1e5119702fb0f3894c8785dd58b8bbe5d73"
fixturePassword = "AYATORI-TEST-ONLY-initial-password"
rotatedPassword = "AYATORI-TEST-ONLY-rotated-password"
controllerNamespace = "database-controller"
secretName = "management"
)
// fixture 不接受外部 DSN,只创建自己的临时容器并按确切 ID 清理。
@@ -79,7 +84,7 @@ func postgresFixture(t *testing.T, ctx context.Context) (string, int) {
t.Fatal("invalid fixture port")
}
// 初次 init 的临时服务器只监听 Unix socket,必须等最终 TCP listener。
for exec.CommandContext(ctx, "docker", dockerExec, id, "pg_isready", "-h", "127.0.0.1", "-U", fixtureUser).Run() != nil {
for exec.CommandContext(ctx, "docker", dockerExec, id, "pg_isready", "-h", fixtureAddress, "-U", fixtureUser).Run() != nil {
select {
case <-ctx.Done():
t.Fatal("fixture startup timed out")
@@ -154,7 +159,7 @@ func target(t *testing.T, port int, mode instance.TLSMode) instance.ObservationT
}
endpoint, err := instance.NewEndpoint(instance.EndpointValues{
Host: fixtureHost,
HostAddr: "127.0.0.1",
HostAddr: fixtureAddress,
Port: port,
ManagementDatabase: fixtureUser,
TLSMode: mode,
@@ -164,8 +169,8 @@ func target(t *testing.T, port int, mode instance.TLSMode) instance.ObservationT
}
ref, err := instance.NewCredentialReference(instance.CredentialReferenceValues{
Name: secretName,
UsernameKey: "login",
PasswordKey: "credential",
UsernameKey: managementUsernameKey,
PasswordKey: managementPasswordKey,
})
if err != nil {
t.Fatal(err)
@@ -196,6 +201,7 @@ func (r *gatedReader) Read(ctx context.Context, ref instance.CredentialReference
// credentialFixture 为每个场景创建独立 API server、PostgreSQL 和应用服务。
type credentialFixture struct {
config *rest.Config
ctx context.Context
client *kubernetes.Clientset
reader *secretadapter.SecretCredentials
@@ -227,7 +233,7 @@ func newCredentialFixture(t *testing.T) *credentialFixture {
if err != nil {
t.Fatal("cannot create test client")
}
for _, namespace := range []string{controllerNamespace, "unrelated"} {
for _, namespace := range []string{controllerNamespace, unrelatedNamespace} {
_, err := client.CoreV1().Namespaces().Create(
ctx,
&corev1.Namespace{Name: namespace},
@@ -260,6 +266,7 @@ func newCredentialFixture(t *testing.T) *credentialFixture {
t.Cleanup(service.Close)
return &credentialFixture{
config: config,
ctx: ctx,
client: client,
reader: reader,
@@ -277,8 +284,8 @@ func (f *credentialFixture) createSecret(t *testing.T, namespace string) {
secret := &corev1.Secret{
Name: secretName,
Data: map[string][]byte{
"login": []byte(fixtureUser),
"credential": []byte(fixturePassword),
managementUsernameKey: []byte(fixtureUser),
managementPasswordKey: []byte(fixturePassword),
},
}
if _, err := f.client.CoreV1().Secrets(namespace).Create(f.ctx, secret, metav1.CreateOptions{}); err != nil {
@@ -0,0 +1,210 @@
//go:build integration
package postgresql_test
import (
"bytes"
"context"
"os"
"testing"
"time"
databasev1alpha1 "git.ddupan.top/panxiao81/ayatori/api/database/v1alpha1"
secretadapter "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"
corev1 "k8s.io/api/core/v1"
rbacv1 "k8s.io/api/rbac/v1"
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/runtime"
yamlutil "k8s.io/apimachinery/pkg/util/yaml"
"k8s.io/client-go/rest"
ctrl "sigs.k8s.io/controller-runtime"
"sigs.k8s.io/controller-runtime/pkg/client"
controllerconfig "sigs.k8s.io/controller-runtime/pkg/config"
"sigs.k8s.io/controller-runtime/pkg/envtest"
metricsserver "sigs.k8s.io/controller-runtime/pkg/metrics/server"
"sigs.k8s.io/yaml"
)
const watchRevisionAnnotation = "test.ayatori/observation"
func TestInstanceControllerWithRealPostgreSQL(t *testing.T) {
f := newCredentialFixture(t)
if _, err := envtest.InstallCRDs(f.config, envtest.CRDInstallOptions{
Paths: []string{"../../../../config/crd/bases"}, ErrorIfPathMissing: true,
}); err != nil {
t.Fatal(err)
}
scheme := runtime.NewScheme()
for _, install := range []func(*runtime.Scheme) error{databasev1alpha1.AddToScheme, corev1.AddToScheme, rbacv1.AddToScheme} {
if err := install(scheme); err != nil {
t.Fatal(err)
}
}
apiClient, err := client.New(f.config, client.Options{Scheme: scheme})
if err != nil {
t.Fatal(err)
}
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 名称唯一。
skipRepeatedName := true
manager, err := ctrl.NewManager(restricted, ctrl.Options{
Scheme: scheme, Cache: databasecontroller.InstanceCacheOptions(controllerNamespace),
Metrics: metricsserver.Options{BindAddress: "0"}, HealthProbeBindAddress: "0",
Controller: controllerconfig.Controller{SkipNameValidation: &skipRepeatedName},
})
if err != nil {
t.Fatal(err)
}
reconciler := &databasecontroller.InstanceReconciler{Observer: service, SecretNamespace: controllerNamespace}
if err := reconciler.SetupWithManager(manager); err != nil {
t.Fatal(err)
}
managerContext, stop := context.WithCancel(f.ctx)
done := make(chan error, 1)
go func() { done <- manager.Start(managerContext) }()
t.Cleanup(func() {
stop()
select {
case err := <-done:
if err != nil {
t.Error(err)
}
case <-time.After(20 * time.Second):
t.Error("Instance manager 未停止")
}
service.Close()
})
object := &databasev1alpha1.PostgreSQLInstance{}
object.Name = "native-instance"
object.Spec.Endpoint = databasev1alpha1.PostgreSQLEndpoint{
Host: fixtureHost, HostAddr: fixtureAddress, Port: int32(f.port), SSLMode: "disable",
}
object.Spec.AdminCredentialRef = databasev1alpha1.AdminCredentialReference{
Name: secretName, UsernameKey: managementUsernameKey, PasswordKey: managementPasswordKey,
}
if err := apiClient.Create(f.ctx, object); err != nil {
t.Fatal(err)
}
awaitInstanceReason(t, f, apiClient, object, "DependencyUnavailable")
// 30 秒轮询前必须收到 Secret 创建事件;实际 controller 使用 namespace Role + metadata watch。
useNativeManager(t, f)
awaitInstanceReason(t, f, apiClient, object, "ManagementReady")
before := f.backendIDs(t)
f.updateSecret(t, func(secret *corev1.Secret) {
secret.Annotations = map[string]string{watchRevisionAnnotation: "changed"}
})
// 用实际 API 事件触发重验,metadata 改动不应换池。
time.Sleep(200 * time.Millisecond)
if f.backendIDs(t) != before {
t.Fatal("无关 Secret metadata 修改重建了连接")
}
f.queryPostgres(t, "ALTER ROLE native_manager PASSWORD '"+rotatedPassword+"'")
f.updateSecret(t, func(secret *corev1.Secret) { secret.Data[managementPasswordKey] = []byte("invalid-test-password") })
awaitInstanceReason(t, f, apiClient, object, "AuthenticationFailed")
f.updateSecret(t, func(secret *corev1.Secret) { secret.Data[managementPasswordKey] = []byte(rotatedPassword) })
awaitInstanceReason(t, f, apiClient, object, "ManagementReady")
if f.backendIDs(t) == before {
t.Fatal("凭据轮换没有替换旧连接")
}
f.queryPostgres(t, "ALTER ROLE native_manager NOCREATEROLE")
f.updateSecret(t, func(secret *corev1.Secret) { secret.Annotations[watchRevisionAnnotation] = "recheck" })
awaitInstanceReason(t, f, apiClient, object, "InsufficientPrivileges")
f.queryPostgres(t, "ALTER ROLE native_manager CREATEROLE")
f.updateSecret(t, func(secret *corev1.Secret) { secret.Annotations[watchRevisionAnnotation] = "recovered" })
awaitInstanceReason(t, f, apiClient, object, "ManagementReady")
if err := f.client.CoreV1().Secrets(controllerNamespace).Delete(f.ctx, secretName, metav1.DeleteOptions{}); err != nil {
t.Fatal(err)
}
awaitInstanceReason(t, f, apiClient, object, "DependencyUnavailable")
if f.backendIDs(t) != "" {
t.Fatal("Secret 删除后旧连接未释放")
}
}
func awaitInstanceReason(t *testing.T, f *credentialFixture, apiClient client.Client,
object *databasev1alpha1.PostgreSQLInstance, reason string) {
t.Helper()
deadline := time.Now().Add(10 * time.Second)
for time.Now().Before(deadline) {
if err := apiClient.Get(f.ctx, client.ObjectKeyFromObject(object), object); err != nil {
t.Fatal(err)
}
condition := meta.FindStatusCondition(object.Status.Conditions, "Ready")
if condition != nil && condition.Reason == reason && condition.ObservedGeneration == object.Generation {
if (condition.Status == metav1.ConditionTrue) != (reason == "ManagementReady") {
t.Fatal("Ready 与检查结果不一致")
}
return
}
time.Sleep(50 * time.Millisecond)
}
t.Fatalf("watch 未及时推进到 %s", reason)
}
func instanceControllerRBAC(t *testing.T, f *credentialFixture, apiClient client.Client) *rest.Config {
t.Helper()
roleBytes, err := os.ReadFile("../../../../config/rbac/role.yaml")
if err != nil {
t.Fatal(err)
}
role := &rbacv1.ClusterRole{}
if err := yaml.Unmarshal(roleBytes, role); err != nil {
t.Fatal(err)
}
if err := apiClient.Create(f.ctx, role); err != nil {
t.Fatal(err)
}
user := "instance-controller-test"
binding := &rbacv1.ClusterRoleBinding{}
binding.Name = user
binding.RoleRef = rbacv1.RoleRef{APIGroup: rbacv1.GroupName, Kind: "ClusterRole", Name: role.Name}
binding.Subjects = []rbacv1.Subject{{Kind: "User", APIGroup: rbacv1.GroupName, Name: user}}
if err := apiClient.Create(f.ctx, binding); err != nil {
t.Fatal(err)
}
credentialBytes, err := os.ReadFile("../../../../config/rbac/database_credentials_role.yaml")
if err != nil {
t.Fatal(err)
}
namespaceRole := &rbacv1.Role{}
decoder := yamlutil.NewYAMLOrJSONDecoder(bytes.NewReader(credentialBytes), 4096)
if err := decoder.Decode(namespaceRole); err != nil {
t.Fatal(err)
}
namespaceRole.Namespace = controllerNamespace
if err := apiClient.Create(f.ctx, namespaceRole); err != nil {
t.Fatal(err)
}
namespaceBinding := &rbacv1.RoleBinding{}
namespaceBinding.Name, namespaceBinding.Namespace = user, controllerNamespace
namespaceBinding.RoleRef = rbacv1.RoleRef{APIGroup: rbacv1.GroupName, Kind: "Role", Name: namespaceRole.Name}
namespaceBinding.Subjects = binding.Subjects
if err := apiClient.Create(f.ctx, namespaceBinding); err != nil {
t.Fatal(err)
}
config := rest.CopyConfig(f.config)
config.Impersonate.UserName = user
restrictedClient, err := client.New(config, client.Options{Scheme: apiClient.Scheme()})
if err != nil {
t.Fatal(err)
}
secret := &corev1.Secret{}
err = restrictedClient.Get(f.ctx, client.ObjectKey{Namespace: unrelatedNamespace, Name: secretName}, secret)
if !apierrors.IsForbidden(err) {
t.Fatal("Instance controller 可以跨 namespace 读取 Secret")
}
return config
}
@@ -0,0 +1,58 @@
package postgresql
import (
"context"
"git.ddupan.top/panxiao81/ayatori/internal/database/application"
"git.ddupan.top/panxiao81/ayatori/internal/database/domain/instance"
)
// 只读取当前执行角色的属性,不能从可继承的角色成员关系推导 CREATEDB/CREATEROLE。
// 所有事实来自同一条语句;不创建探针数据库,不修改管理账号或持久 schema。
const inspectManagementStatement = `
SELECT
pg_catalog.current_setting('server_version'),
ARRAY(SELECT name::text FROM pg_catalog.pg_available_extensions ORDER BY name),
role.rolsuper,
role.rolcreaterole,
role.rolcreatedb,
pg_catalog.pg_is_in_recovery() OR
pg_catalog.current_setting('transaction_read_only')::boolean
FROM pg_catalog.pg_roles AS role
WHERE role.rolname = current_user`
func (d *database) InspectManagement(ctx context.Context) (application.DatabaseMetadata, error) {
var metadata application.DatabaseMetadata
var superuser, createRole, createDatabase, readOnly bool
err := d.pool.QueryRow(ctx, inspectManagementStatement).Scan(
&metadata.Version, &metadata.AvailableExtensions,
&superuser, &createRole, &createDatabase, &readOnly,
)
if err != nil {
return application.DatabaseMetadata{}, safeError(err, application.ErrObservation)
}
checks := instance.ManagementChecks{
Connection: instance.CheckPassed,
Metadata: instance.CheckPassed,
Roles: nativePrivilege(createRole && !superuser),
Databases: nativePrivilege(createDatabase && !superuser),
// CREATEROLE 可管理自己新建角色的 membership;供应时必须显式取得 SET 权限,
// 再以 owner 操作数据库 ACL。这里不授权操作任意导入角色或他人数据库。
Grants: nativePrivilege(createRole && createDatabase && !superuser),
// 新建数据库 owner 可安装 trusted 扩展。具体扩展仍需逐请求执行和回读,
// 非 trusted 扩展不能因出现在 available 列表就视为可安装。
Extensions: nativePrivilege(createRole && createDatabase && !superuser),
}
if readOnly {
checks.Databases = instance.CheckUnavailable
}
metadata.Management = checks
return metadata, nil
}
func nativePrivilege(allowed bool) instance.CheckResult {
if allowed {
return instance.CheckPassed
}
return instance.CheckInsufficientPrivileges
}
@@ -0,0 +1,152 @@
//go:build integration
package postgresql_test
import (
"context"
"errors"
"testing"
"github.com/jackc/pgx/v5"
corev1 "k8s.io/api/core/v1"
"git.ddupan.top/panxiao81/ayatori/internal/database/application"
"git.ddupan.top/panxiao81/ayatori/internal/database/domain/instance"
)
const nativeManager = "native_manager"
func useNativeManager(t *testing.T, f *credentialFixture) {
t.Helper()
f.queryPostgres(t, "CREATE ROLE native_manager LOGIN CREATEDB CREATEROLE PASSWORD '"+fixturePassword+"'")
f.createSecret(t, controllerNamespace)
f.updateSecret(t, func(secret *corev1.Secret) { secret.Data[managementUsernameKey] = []byte(nativeManager) })
}
func assessManagement(t *testing.T, f *credentialFixture) instance.Snapshot {
t.Helper()
observation, err := f.service.ObserveManagement(f.ctx, f.target)
if err != nil {
t.Fatal(err)
}
aggregate, err := instance.Reconstitute(f.target, instance.Snapshot{}, false)
if err != nil {
t.Fatal(err)
}
if err := aggregate.BeginValidation(); err != nil {
t.Fatal(err)
}
capabilities, err := observation.Capabilities()
if err != nil {
t.Fatal(err)
}
if err := aggregate.AssessManagement(capabilities); err != nil {
t.Fatal(err)
}
return aggregate.Snapshot()
}
func TestNativeManagementPrivileges(t *testing.T) {
f := newCredentialFixture(t)
useNativeManager(t, f)
if snapshot := assessManagement(t, f); snapshot.Readiness != instance.Ready {
t.Fatal("原生非 superuser 管理账号未通过检查")
}
before := f.backendIDs(t)
for _, attribute := range []string{"NOCREATEROLE", "NOCREATEDB"} {
f.queryPostgres(t, "ALTER ROLE native_manager "+attribute)
if snapshot := assessManagement(t, f); snapshot.Failure != instance.InsufficientPrivileges {
t.Fatal("已有连接忽略了管理权限撤回")
}
f.queryPostgres(t, "ALTER ROLE native_manager CREATEROLE CREATEDB")
if snapshot := assessManagement(t, f); snapshot.Readiness != instance.Ready {
t.Fatal("管理权限恢复后无法重新就绪")
}
}
if f.backendIDs(t) != before {
t.Fatal("权限检查不应要求重建连接才生效")
}
f.queryPostgres(t, "ALTER ROLE native_manager SET default_transaction_read_only = on")
f.service.Forget(f.target.Identity().Name())
if snapshot := assessManagement(t, f); snapshot.Failure != instance.DependencyUnavailable {
t.Fatal("只读会话不应标记可供应")
}
f.queryPostgres(t, "ALTER ROLE native_manager RESET default_transaction_read_only")
f.service.Forget(f.target.Identity().Name())
if snapshot := assessManagement(t, f); snapshot.Readiness != instance.Ready {
t.Fatal("恢复可写会话后没有就绪")
}
f.updateSecret(t, func(secret *corev1.Secret) { secret.Data[managementUsernameKey] = []byte(fixtureUser) })
if snapshot := assessManagement(t, f); snapshot.Failure != instance.InsufficientPrivileges {
t.Fatal("不应以 superuser 绕过非特权账号合同")
}
reads := 0
f.gate.beforeRead = func() {
reads++
if reads == 2 {
f.updateSecret(t, func(secret *corev1.Secret) { secret.Data[managementPasswordKey] = []byte(rotatedPassword) })
}
}
observation, err := f.service.ObserveManagement(f.ctx, f.target)
if !errors.Is(err, application.ErrCredentialsChanged) || observation.Target().Validate() == nil {
t.Fatal("管理观察期间凭据轮换应丢弃全部能力结果")
}
if f.backendIDs(t) != "" {
t.Fatal("中途轮换后不应保留旧管理连接")
}
}
// 以实际非 superuser 会话验证能力矩阵的依据,不用超级用户执行 SQL 模拟管理账号。
// 这些固定名称只存在于本测试独占容器,生产观察本身不会创建探针对象。
func TestNativeManagementSupplyContract(t *testing.T) {
f := newCredentialFixture(t)
useNativeManager(t, f)
config, err := pgx.ParseConfig("")
if err != nil {
t.Fatal("无法装配隔离测试连接")
}
config.Host, config.Port = fixtureAddress, uint16(f.port)
config.Database, config.User, config.Password = fixtureUser, nativeManager, fixturePassword
config.TLSConfig, config.Fallbacks = nil, nil
connection, err := pgx.ConnectConfig(f.ctx, config)
if err != nil {
t.Fatal("非 superuser 测试连接失败")
}
t.Cleanup(func() { _ = connection.Close(context.Background()) })
execute := func(statement string) {
t.Helper()
if _, err := connection.Exec(f.ctx, statement); err != nil {
t.Fatalf("原生管理能力合同未满足,步骤 %q", statement)
}
}
execute("CREATE ROLE managed_owner LOGIN NOSUPERUSER NOCREATEDB NOCREATEROLE NOREPLICATION")
execute("GRANT managed_owner TO native_manager WITH SET TRUE")
execute("CREATE DATABASE managed_database OWNER managed_owner")
execute("SET ROLE managed_owner")
execute("REVOKE CONNECT ON DATABASE managed_database FROM PUBLIC")
execute("GRANT CONNECT ON DATABASE managed_database TO managed_owner")
execute("RESET ROLE")
config.Database = "managed_database"
tenantConnection, err := pgx.ConnectConfig(f.ctx, config)
if err != nil {
t.Fatal("管理账号无法访问其受管数据库")
}
defer func() { _ = tenantConnection.Close(context.Background()) }()
if _, err := tenantConnection.Exec(f.ctx, "SET ROLE managed_owner; CREATE EXTENSION hstore"); err != nil {
t.Fatal("owner 无法安装 trusted 扩展")
}
var installed bool
if err := tenantConnection.QueryRow(f.ctx, "SELECT EXISTS (SELECT FROM pg_catalog.pg_extension WHERE extname = 'hstore')").Scan(&installed); err != nil || !installed {
t.Fatal("扩展安装后实际回读失败")
}
if _, err := tenantConnection.Exec(f.ctx, "CREATE EXTENSION file_fdw"); err == nil {
t.Fatal("非 trusted 扩展不应被 Ready 隐式授权")
}
if err := tenantConnection.Close(f.ctx); err != nil {
t.Fatal("关闭目标数据库连接失败")
}
execute("SET ROLE managed_owner")
execute("DROP DATABASE managed_database")
execute("RESET ROLE")
execute("DROP ROLE managed_owner")
}
@@ -0,0 +1,149 @@
package application
import (
"context"
"errors"
"git.ddupan.top/panxiao81/ayatori/internal/database/domain/instance"
)
const (
instanceDependencyUnavailable = "DependencyUnavailable"
instanceAuthenticationFailed = "AuthenticationFailed"
)
// InstanceRecord 是 API 快照;Revision 仅用于持久化并发保护,不是领域版本。
type InstanceRecord struct {
Target instance.ObservationTarget
Revision string
Deleting bool
}
type InstanceResources interface {
LoadInstance(context.Context, string) (*InstanceRecord, error)
ProtectInstance(context.Context, *InstanceRecord) (*InstanceRecord, error)
// InstanceReferences 返回一个可定位的阻塞引用;空字符串表示没有引用。
InstanceReferences(context.Context, string) (string, error)
}
type InstanceObserver interface {
ObserveManagement(context.Context, instance.ObservationTarget) (InstanceObservation, error)
Forget(string)
}
type InstanceResult struct {
Record *InstanceRecord
Snapshot instance.Snapshot
Reason string
Message string
RemoveProtection bool
}
// InstanceReconciliation 协调 API 保护、实时观察和领域判断,不拼装 Kubernetes status。
type InstanceReconciliation struct {
Resources InstanceResources
Observer InstanceObserver
}
func (s *InstanceReconciliation) Reconcile(ctx context.Context, name string) (InstanceResult, error) {
record, err := s.Resources.LoadInstance(ctx, name)
if err != nil {
return InstanceResult{}, err
}
if record == nil {
s.Observer.Forget(name)
return InstanceResult{}, nil
}
if record.Deleting {
return s.deleting(ctx, record)
}
record, err = s.Resources.ProtectInstance(ctx, record)
if err != nil {
return InstanceResult{}, err
}
// 每轮从无证据的领域对象开始;持久化 Ready 和连接存活不能替代本轮检查。
aggregate, err := instance.Reconstitute(record.Target, instance.Snapshot{}, false)
if err != nil {
return InstanceResult{}, err
}
if err := aggregate.BeginValidation(); err != nil {
return InstanceResult{}, err
}
observation, observationErr := s.Observer.ObserveManagement(ctx, record.Target)
result := InstanceResult{Record: record}
if observationErr != nil {
result.Snapshot = aggregate.Snapshot()
result.Snapshot.Readiness = instance.NotReady
result.Snapshot.ObservedRevision = record.Target.Revision().Value()
result.Reason, result.Message = observationFailure(observationErr)
return result, nil
}
capabilities, err := observation.Capabilities()
if err != nil {
return InstanceResult{}, err
}
if err := aggregate.AssessManagement(capabilities); err != nil {
return InstanceResult{}, err
}
result.Snapshot = aggregate.Snapshot()
result.Snapshot.ReportedVersion = observation.Version()
result.Reason, result.Message = managementResult(result.Snapshot.Failure)
return result, nil
}
func (s *InstanceReconciliation) deleting(ctx context.Context, record *InstanceRecord) (InstanceResult, error) {
name := record.Target.Identity().Name()
s.Observer.Forget(name)
aggregate, err := instance.Reconstitute(record.Target, instance.Snapshot{}, true)
if err != nil {
return InstanceResult{}, err
}
if err := aggregate.BeginDeletion(); err != nil {
return InstanceResult{}, err
}
result := InstanceResult{Record: record, Snapshot: aggregate.Snapshot(), Reason: "Deleting"}
reference, err := s.Resources.InstanceReferences(ctx, name)
if err != nil {
result.Reason = instanceDependencyUnavailable
result.Message = "无法确认 Database/Tenant 引用已解除;保留 Instance 删除保护并重试"
return result, nil
}
if reference != "" {
result.Reason = "InstanceInUse"
result.Message = "仍被 " + reference + " 引用;先处理该资源,不会级联删除外部数据库"
return result, nil
}
result.Message = "引用已解除,仅移除登记保护;不删除 PostgreSQL 或凭据"
result.RemoveProtection = true
return result, nil
}
func observationFailure(err error) (string, string) {
switch {
case errors.Is(err, ErrAuthentication):
return instanceAuthenticationFailed, "管理连接认证或 TLS 校验失败;检查管理 Secret 和 CA/证书配置"
case errors.Is(err, ErrCredentialsInvalid):
return "InvalidCredentials", "管理 Secret 的用户名或密码字段缺失;检查引用字段映射"
case errors.Is(err, ErrCredentialsChanged):
return "CredentialsChanged", "观察期间管理凭据变化,已丢弃结果并关闭旧连接;等待重新验证"
case errors.Is(err, ErrCredentialsUnavailable):
return instanceDependencyUnavailable, "无法读取管理 Secret;检查其是否存在及 controller namespace 内的读取权限"
default:
return instanceDependencyUnavailable, "管理连接或能力查询失败;检查 PostgreSQL 可达性、catalog 读取权限和超时"
}
}
func managementResult(failure instance.Failure) (string, string) {
switch failure {
case instance.NoFailure:
return "ManagementReady", "当前管理能力检查通过;具体资源授权和扩展安装仍需执行时验证"
case instance.InsufficientPrivileges:
return "InsufficientPrivileges", "原生管理要求非 superuser 且具备 CREATEDB/CREATEROLE;不会自动修改账号权限"
case instance.DependencyUnavailable:
return instanceDependencyUnavailable, "当前 PostgreSQL 不可写或所需管理能力暂不可用"
case instance.AuthenticationFailed:
return instanceAuthenticationFailed, "当前管理能力检查未通过认证"
default:
return "ObservationIncomplete", "管理能力检查尚有缺项,不能仅凭 metadata 查询成功标记 Ready"
}
}
@@ -0,0 +1,79 @@
package application
import (
"context"
"errors"
"testing"
"git.ddupan.top/panxiao81/ayatori/internal/database/domain/instance"
)
func TestInstanceFailurePresentation(t *testing.T) {
cases := []struct {
err error
reason string
}{
{ErrAuthentication, "AuthenticationFailed"},
{ErrCredentialsInvalid, "InvalidCredentials"},
{ErrCredentialsChanged, "CredentialsChanged"},
{ErrCredentialsUnavailable, instanceDependencyUnavailable},
{context.DeadlineExceeded, instanceDependencyUnavailable},
{errors.New("private backend detail"), instanceDependencyUnavailable},
}
for _, test := range cases {
reason, message := observationFailure(test.err)
if reason != test.reason || message == "" || message == test.err.Error() {
t.Fatal("观察失败没有安全且可诊断的状态")
}
}
for _, failure := range []instance.Failure{
instance.NoFailure, instance.ObservationIncomplete, instance.DependencyUnavailable,
instance.AuthenticationFailed, instance.InsufficientPrivileges,
} {
reason, message := managementResult(failure)
if reason == "" || message == "" || (reason == "ManagementReady") != (failure == instance.NoFailure) {
t.Fatal("领域能力判定与状态不一致")
}
}
}
func TestMetadataCannotEstablishManagementReadiness(t *testing.T) {
source := &sourceStub{}
source.credentials, _ = NewCredentials("test", serviceTestPassword)
connector := &connectorStub{}
service, err := NewInstanceService(source, connector)
if err != nil {
t.Fatal(err)
}
defer service.Close()
target := serviceTarget(t, "uid", "postgres.test", "management", 1)
if _, err := service.ObserveMetadata(t.Context(), target); err != nil {
t.Fatal(err)
}
connector.databases[0].metadata.Management = instance.ManagementChecks{
Connection: instance.CheckPassed, Metadata: instance.CheckPassed,
Roles: instance.CheckPassed, Databases: instance.CheckPassed,
Grants: instance.CheckPassed, Extensions: instance.CheckPassed,
}
observation, err := service.ObserveMetadata(t.Context(), target)
if err != nil {
t.Fatal(err)
}
capabilities, err := observation.Capabilities()
if err != nil {
t.Fatal(err)
}
aggregate, err := instance.Reconstitute(target, instance.Snapshot{}, false)
if err != nil {
t.Fatal(err)
}
if err := aggregate.BeginValidation(); err != nil {
t.Fatal(err)
}
if err := aggregate.AssessManagement(capabilities); err != nil {
t.Fatal(err)
}
if aggregate.Snapshot().Failure != instance.ObservationIncomplete {
t.Fatal("metadata 入口不应携带完整管理检查")
}
}
@@ -36,6 +36,7 @@ var (
// Metadata 只查询版本与可用扩展,不能产生领域 Ready。
type Database interface {
InspectMetadata(context.Context) (DatabaseMetadata, error)
InspectManagement(context.Context) (DatabaseMetadata, error)
Close()
}
@@ -82,17 +83,26 @@ func (s *InstanceService) ObserveVersion(ctx context.Context, target instance.Ob
// ObserveMetadata 返回当前目标和凭据下的版本与扩展;任何失败均丢弃全部结果。
// 调用者仍需使用 CR resourceVersion 保存前提防止 spec 并发修改;本方法不建立跨系统事务。
func (s *InstanceService) ObserveMetadata(ctx context.Context, target instance.ObservationTarget) (MetadataObservation, error) {
func (s *InstanceService) ObserveMetadata(ctx context.Context, target instance.ObservationTarget) (InstanceObservation, error) {
return s.observe(ctx, target, false)
}
// ObserveManagement 复用同一凭据刷新与回读边界,但每轮重新检查原生管理能力。
func (s *InstanceService) ObserveManagement(ctx context.Context, target instance.ObservationTarget) (InstanceObservation, error) {
return s.observe(ctx, target, true)
}
func (s *InstanceService) observe(ctx context.Context, target instance.ObservationTarget, management bool) (InstanceObservation, error) {
if err := target.Validate(); err != nil {
return MetadataObservation{}, err
return InstanceObservation{}, err
}
s.mu.Lock()
defer s.mu.Unlock()
if s.closed {
return MetadataObservation{}, ErrClosed
return InstanceObservation{}, ErrClosed
}
if err := ctx.Err(); err != nil {
return MetadataObservation{}, err
return InstanceObservation{}, err
}
// 先读取有效凭据。读取失败时不得继续使用缓存中的旧连接。
@@ -100,11 +110,11 @@ func (s *InstanceService) ObserveMetadata(ctx context.Context, target instance.O
credentials, err := s.source.Read(ctx, target.Definition().AdminCredential())
if err != nil {
s.release(name)
return MetadataObservation{}, credentialError(err)
return InstanceObservation{}, credentialError(err)
}
if credentials.username == "" || credentials.password == "" {
s.release(name)
return MetadataObservation{}, ErrCredentialsInvalid
return InstanceObservation{}, ErrCredentialsInvalid
}
// 连接身份与有效值均未变化时复用 pgxpool;generation 本身不要求换池。
@@ -117,7 +127,7 @@ func (s *InstanceService) ObserveMetadata(ctx context.Context, target instance.O
if current == nil {
database, err := s.connector.Connect(ctx, target.Definition().Endpoint(), credentials)
if err != nil {
return MetadataObservation{}, err
return InstanceObservation{}, err
}
current = &entry{
target: target,
@@ -127,30 +137,38 @@ func (s *InstanceService) ObserveMetadata(ctx context.Context, target instance.O
s.entries[name] = current
}
metadata, err := current.database.InspectMetadata(ctx)
var metadata DatabaseMetadata
if management {
metadata, err = current.database.InspectManagement(ctx)
} else {
metadata, err = current.database.InspectMetadata(ctx)
// 即使 adapter 误填权限,也不能把只读 metadata 入口升级为 Ready。
metadata.Management = instance.ManagementChecks{}
}
if err != nil {
s.release(name)
return MetadataObservation{}, err
return InstanceObservation{}, err
}
if metadata.Version == "" {
s.release(name)
return MetadataObservation{}, ErrObservation
return InstanceObservation{}, ErrObservation
}
// 回读后再检查凭据,避免把轮换前取得的结果交给新凭据的调用链。
latest, err := s.source.Read(ctx, target.Definition().AdminCredential())
if err != nil {
s.release(name)
return MetadataObservation{}, credentialError(err)
return InstanceObservation{}, credentialError(err)
}
if latest != credentials {
s.release(name)
return MetadataObservation{}, ErrCredentialsChanged
return InstanceObservation{}, ErrCredentialsChanged
}
return MetadataObservation{
return InstanceObservation{
target: target,
version: metadata.Version,
extensions: instance.ObserveExtensionSupport(metadata.AvailableExtensions),
management: metadata.Management,
}, nil
}
@@ -42,6 +42,10 @@ type databaseStub struct {
metadata DatabaseMetadata
}
func (d *databaseStub) InspectManagement(ctx context.Context) (DatabaseMetadata, error) {
return d.InspectMetadata(ctx)
}
func (d *databaseStub) InspectMetadata(context.Context) (DatabaseMetadata, error) {
return d.metadata, d.err
}
+14 -6
View File
@@ -18,23 +18,31 @@ package application
import "git.ddupan.top/panxiao81/ayatori/internal/database/domain/instance"
// DatabaseMetadata 是一次只读查询的事实,不包含管理权限或完整就绪结论。
// DatabaseMetadata 是一次只读查询的事实,可附带原生管理检查,但不包含就绪结论。
// AvailableExtensions 是服务器提供的可用列表,不是已安装列表或安装授权。
type DatabaseMetadata struct {
Version string
AvailableExtensions []string
// Management 仅由 InspectManagement 填充;metadata 查询必须保持未观察。
Management instance.ManagementChecks
}
// MetadataObservation 只在查询成功且有效凭据再次核对一致后产生。
// InstanceObservation 只在查询成功且有效凭据再次核对一致后产生。
// target 绑定本次调用,而非连接最初创建时的 generation;零值表示没有观察。
type MetadataObservation struct {
type InstanceObservation struct {
target instance.ObservationTarget
version string
extensions instance.ExtensionSupport
management instance.ManagementChecks
}
func (o MetadataObservation) Target() instance.ObservationTarget { return o.target }
func (o MetadataObservation) Version() string { return o.version }
func (o MetadataObservation) Extensions() instance.ExtensionSupport {
// Capabilities 保留缺项为未观察;不能从 metadata 的成功补齐管理检查。
func (o InstanceObservation) Capabilities() (instance.CapabilityObservation, error) {
return instance.NewCapabilityObservation(o.target, o.version, o.management)
}
func (o InstanceObservation) Target() instance.ObservationTarget { return o.target }
func (o InstanceObservation) Version() string { return o.version }
func (o InstanceObservation) Extensions() instance.ExtensionSupport {
return o.extensions
}
@@ -0,0 +1,42 @@
package controller
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"
)
type InstanceReconciler struct {
Client client.Client
Reader client.Reader
Observer application.InstanceObserver
SecretNamespace string
}
// +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)
if err != nil {
return ctrl.Result{}, err
}
// 查询超时后仍用 worker context 保存安全失败结果;manager 停止时不强行写入。
if err := resources.PresentInstance(ctx, result); err != nil {
return ctrl.Result{}, err
}
if result.Record == nil || result.RemoveProtection {
return ctrl.Result{}, nil
}
return ctrl.Result{RequeueAfter: dependencyRetry}, nil
}
@@ -0,0 +1,192 @@
package controller
import (
"context"
"errors"
"testing"
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"
"git.ddupan.top/panxiao81/ayatori/internal/database/domain/instance"
apierrors "k8s.io/apimachinery/pkg/api/errors"
"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"
"sigs.k8s.io/controller-runtime/pkg/controller/controllerutil"
)
type instanceBackend struct {
checks instance.ManagementChecks
err error
inspect func()
closed int
}
func (b *instanceBackend) Read(context.Context, instance.CredentialReference) (application.Credentials, error) {
return application.NewCredentials("fixture", "test-only-instance-password")
}
func (b *instanceBackend) Connect(context.Context, instance.Endpoint, application.Credentials) (application.Database, error) {
return b, nil
}
func (b *instanceBackend) InspectMetadata(context.Context) (application.DatabaseMetadata, error) {
return application.DatabaseMetadata{Version: "18"}, nil
}
func (b *instanceBackend) InspectManagement(context.Context) (application.DatabaseMetadata, error) {
if b.inspect != nil {
b.inspect()
}
return application.DatabaseMetadata{Version: "18", Management: b.checks}, b.err
}
func (b *instanceBackend) Close() { b.closed++ }
func newInstanceReconciler(t *testing.T, apiClient client.Client, backend *instanceBackend) *InstanceReconciler {
t.Helper()
service, err := application.NewInstanceService(backend, backend)
if err != nil {
t.Fatal(err)
}
t.Cleanup(service.Close)
return &InstanceReconciler{Client: apiClient, Reader: apiClient, Observer: service}
}
func reconcileInstance(t *testing.T, reconciler *InstanceReconciler, object *databasev1alpha1.PostgreSQLInstance) {
t.Helper()
if _, err := reconciler.Reconcile(t.Context(), ctrl.Request{NamespacedName: client.ObjectKeyFromObject(object)}); err != nil {
t.Fatal(err)
}
}
func assertInstanceReason(t *testing.T, object *databasev1alpha1.PostgreSQLInstance, reason string) {
t.Helper()
condition := meta.FindStatusCondition(object.Status.Conditions, "Ready")
if condition == nil || condition.Reason != reason || condition.ObservedGeneration != object.Generation {
t.Fatalf("Instance 状态不是当前 generation 的 %s", reason)
}
if reason != "ManagementReady" && condition.Status != metav1.ConditionFalse {
t.Fatal("失败状态仍为 Ready")
}
}
func TestInstanceObservationAPI(t *testing.T) {
apiClient, _, _ := bindingEnvironment(t)
backend := &instanceBackend{checks: instance.ManagementChecks{
Connection: instance.CheckPassed, Metadata: instance.CheckPassed,
Roles: instance.CheckPassed, Databases: instance.CheckPassed,
Grants: instance.CheckPassed, Extensions: instance.CheckPassed,
}}
reconciler := newInstanceReconciler(t, apiClient, backend)
object := readyInstance(t, apiClient, "observed-instance")
backend.inspect = func() {
current := &databasev1alpha1.PostgreSQLInstance{}
current.Name = object.Name
reload(t, apiClient, current)
if !controllerutil.ContainsFinalizer(current, kubernetes.InstanceFinalizer) {
t.Fatal("观察早于 finalizer 持久化")
}
}
reconcileInstance(t, reconciler, object)
reload(t, apiClient, object)
assertInstanceReason(t, object, "ManagementReady")
if object.Status.Phase != string(instance.PhaseReady) || object.Status.PostgreSQLVersion != "18" {
t.Fatal("当前成功观察未呈现")
}
before := object.ResourceVersion
reconcileInstance(t, reconciler, object)
reload(t, apiClient, object)
if object.ResourceVersion != before {
t.Fatal("相同观察不应反复写入 status")
}
backend.err = application.ErrAuthentication
reconcileInstance(t, reconciler, object)
reload(t, apiClient, object)
assertInstanceReason(t, object, "AuthenticationFailed")
if backend.closed != 1 || object.Status.PostgreSQLVersion != "" {
t.Fatal("观察失败应释放连接并清除旧版本结果")
}
backend.err = nil
backend.checks.Grants = instance.CheckUnobserved
reconcileInstance(t, reconciler, object)
reload(t, apiClient, object)
assertInstanceReason(t, object, "ObservationIncomplete")
backend.checks.Grants = instance.CheckPassed
// 用新 service/reconciler 恢复;不依赖上轮领域对象或 Ready。
reconciler = newInstanceReconciler(t, apiClient, backend)
reconcileInstance(t, reconciler, object)
reload(t, apiClient, object)
assertInstanceReason(t, object, "ManagementReady")
backend.inspect = func() {
reload(t, apiClient, object)
object.Annotations = map[string]string{"concurrent": "kept-by-instance-test"}
if err := apiClient.Update(t.Context(), object); err != nil {
t.Fatal(err)
}
}
_, err := reconciler.Reconcile(t.Context(), ctrl.Request{NamespacedName: client.ObjectKeyFromObject(object)})
if !apierrors.IsConflict(err) {
t.Fatal("旧观察不应覆盖在途 API 修改")
}
backend.inspect = nil
reconcileInstance(t, reconciler, object)
reload(t, apiClient, object)
if object.Annotations["concurrent"] != "kept-by-instance-test" {
t.Fatal("重试覆盖了其他字段")
}
}
type failedReferenceReader struct{ client.Reader }
func (*failedReferenceReader) List(context.Context, client.ObjectList, ...client.ListOption) error {
return errors.New("injected reference list failure")
}
func TestInstanceDeletionProtection(t *testing.T) {
apiClient, _, _ := bindingEnvironment(t)
backend := &instanceBackend{}
reconciler := newInstanceReconciler(t, apiClient, backend)
object := readyInstance(t, apiClient, "protected-instance")
reconcileInstance(t, reconciler, object)
reload(t, apiClient, object)
database := availableDatabase(t, apiClient, "retained-database", object)
database.Status.Phase = "Released"
if err := apiClient.Status().Update(t.Context(), database); err != nil {
t.Fatal(err)
}
tenant := provisionTenant("pending-request", object.Name)
requireCreate(t, apiClient, tenant)
if err := apiClient.Delete(t.Context(), object); err != nil {
t.Fatal(err)
}
backend.inspect = func() { t.Fatal("删除中不应连接 PostgreSQL") }
reconciler.Reader = &failedReferenceReader{Reader: apiClient}
reconcileInstance(t, reconciler, object)
reload(t, apiClient, object)
assertInstanceReason(t, object, reasonDependency)
reconciler.Reader = apiClient
reconcileInstance(t, reconciler, object)
reload(t, apiClient, object)
assertInstanceReason(t, object, "InstanceInUse")
if object.Status.Phase != string(instance.PhaseDeleting) || backend.closed != 1 {
t.Fatal("删除没有停止本地观察")
}
if err := apiClient.Delete(t.Context(), database); err != nil {
t.Fatal(err)
}
reconcileInstance(t, reconciler, object)
reload(t, apiClient, object)
assertInstanceReason(t, object, "InstanceInUse")
if err := apiClient.Delete(t.Context(), tenant); err != nil {
t.Fatal(err)
}
reconcileInstance(t, reconciler, object)
if err := apiClient.Get(t.Context(), client.ObjectKeyFromObject(object), object); !apierrors.IsNotFound(err) {
t.Fatal("最后一个引用解除后 Instance 应可删除")
}
reconcileInstance(t, reconciler, object)
}
@@ -0,0 +1,77 @@
package controller
import (
"context"
"errors"
databasev1alpha1 "git.ddupan.top/panxiao81/ayatori/api/database/v1alpha1"
corev1 "k8s.io/api/core/v1"
"k8s.io/apimachinery/pkg/util/validation"
ctrl "sigs.k8s.io/controller-runtime"
"sigs.k8s.io/controller-runtime/pkg/cache"
"sigs.k8s.io/controller-runtime/pkg/client"
"sigs.k8s.io/controller-runtime/pkg/handler"
)
// InstanceCacheOptions 必须在创建 manager 时使用;只 watch 固定 namespace 的 Secret metadata。
// SecretCredentials 始终直读 API,不会令共享 cache 保存密码。
func InstanceCacheOptions(namespace string) cache.Options {
return cache.Options{ByObject: map[client.Object]cache.ByObject{
&corev1.Secret{}: {Namespaces: map[string]cache.Config{namespace: {}}},
}}
}
func (r *InstanceReconciler) SetupWithManager(manager ctrl.Manager) error {
if r.Observer == 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{}).
WatchesMetadata(&corev1.Secret{}, handler.EnqueueRequestsFromMapFunc(r.instancesForSecret)).
Watches(&databasev1alpha1.PostgreSQLDatabase{}, handler.EnqueueRequestsFromMapFunc(r.instanceForReference)).
Watches(&databasev1alpha1.PostgreSQLTenant{}, handler.EnqueueRequestsFromMapFunc(r.instanceForReference)).
Complete(r)
}
func (r *InstanceReconciler) instancesForSecret(ctx context.Context, object client.Object) []ctrl.Request {
if object.GetNamespace() != r.SecretNamespace {
return nil
}
instances := &databasev1alpha1.PostgreSQLInstanceList{}
if err := r.Client.List(ctx, instances); err != nil {
ctrl.LoggerFrom(ctx).Error(err, "无法映射管理 Secret 事件;等待低频重试")
return nil
}
var requests []ctrl.Request
for _, item := range instances.Items {
if string(item.Spec.AdminCredentialRef.Name) == object.GetName() {
request := ctrl.Request{Name: item.Name}
requests = append(requests, request)
}
}
return requests
}
func (r *InstanceReconciler) instanceForReference(_ context.Context, object client.Object) []ctrl.Request {
var name string
switch item := object.(type) {
case *databasev1alpha1.PostgreSQLDatabase:
name = string(item.Spec.InstanceRef.Name)
case *databasev1alpha1.PostgreSQLTenant:
if item.Spec.Provision != nil {
name = string(item.Spec.Provision.InstanceRef.Name)
}
}
if name == "" {
return nil
}
request := ctrl.Request{Name: name}
return []ctrl.Request{request}
}