diff --git a/cmd/database.go b/cmd/database.go new file mode 100644 index 0000000..9603c38 --- /dev/null +++ b/cmd/database.go @@ -0,0 +1,26 @@ +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 +} diff --git a/cmd/main.go b/cmd/main.go index 8a676d6..4144bd6 100644 --- a/cmd/main.go +++ b/cmd/main.go @@ -22,6 +22,7 @@ import ( 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 ) @@ -41,6 +42,10 @@ func init() { // nolint:gocyclo func main() { + var databaseNamespace, databaseRootCert string + flag.StringVar(&databaseNamespace, "database-secret-namespace", os.Getenv("POD_NAMESPACE"), + "固定管理 Secret namespace;为空时不启用 Instance 观测") + flag.StringVar(&databaseRootCert, "database-root-cert", "", "PostgreSQL 管理连接信任的公开 CA bundle 路径") var metricsAddr string var metricsCertPath, metricsCertName, metricsCertKey string var webhookCertPath, webhookCertName, webhookCertKey string @@ -145,7 +150,7 @@ func main() { metricsServerOptions.KeyName = metricsCertKey } - mgr, err := ctrl.NewManager(ctrl.GetConfigOrDie(), ctrl.Options{ + managerOptions := ctrl.Options{ Scheme: scheme, Metrics: metricsServerOptions, WebhookServer: webhookServer, @@ -163,13 +168,25 @@ func main() { // 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) @@ -185,7 +202,12 @@ func main() { } setupLog.Info("Starting manager") - if err := mgr.Start(ctrl.SetupSignalHandler()); err != nil { + 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) } diff --git a/config/manager/manager.yaml b/config/manager/manager.yaml index e72e1a5..2520c4b 100644 --- a/config/manager/manager.yaml +++ b/config/manager/manager.yaml @@ -65,6 +65,11 @@ spec: - --health-probe-bind-address=:8081 image: controller:latest name: manager + env: + - name: POD_NAMESPACE + valueFrom: + fieldRef: + fieldPath: metadata.namespace ports: - containerPort: 8081 name: health diff --git a/config/rbac/database_credentials_role.yaml b/config/rbac/database_credentials_role.yaml new file mode 100644 index 0000000..327eeec --- /dev/null +++ b/config/rbac/database_credentials_role.yaml @@ -0,0 +1,23 @@ +apiVersion: rbac.authorization.k8s.io/v1 +kind: Role +metadata: + name: database-management-credentials + namespace: system +rules: +- apiGroups: [""] + resources: [secrets] + verbs: [get, list, watch] +--- +apiVersion: rbac.authorization.k8s.io/v1 +kind: RoleBinding +metadata: + name: database-management-credentials + namespace: system +roleRef: + apiGroup: rbac.authorization.k8s.io + kind: Role + name: database-management-credentials +subjects: +- kind: ServiceAccount + name: controller-manager + namespace: system diff --git a/config/rbac/kustomization.yaml b/config/rbac/kustomization.yaml index 71fc603..3da2091 100644 --- a/config/rbac/kustomization.yaml +++ b/config/rbac/kustomization.yaml @@ -7,6 +7,7 @@ resources: - service_account.yaml - role.yaml - role_binding.yaml +- database_credentials_role.yaml - leader_election_role.yaml - leader_election_role_binding.yaml # The following RBAC configurations are used to protect diff --git a/config/rbac/role.yaml b/config/rbac/role.yaml index 00be889..e81e13b 100644 --- a/config/rbac/role.yaml +++ b/config/rbac/role.yaml @@ -19,6 +19,7 @@ rules: - database.ayatori.ddupan.top resources: - postgresqldatabases/finalizers + - postgresqlinstances/finalizers - postgresqltenants/finalizers verbs: - update @@ -26,6 +27,7 @@ rules: - database.ayatori.ddupan.top resources: - postgresqldatabases/status + - postgresqlinstances/status - postgresqltenants/status verbs: - get @@ -35,13 +37,6 @@ rules: - database.ayatori.ddupan.top resources: - postgresqlinstances - verbs: - - get - - list - - watch -- apiGroups: - - database.ayatori.ddupan.top - resources: - postgresqltenants verbs: - get diff --git a/docs/database/README.md b/docs/database/README.md index c812f5f..0113e2a 100644 --- a/docs/database/README.md +++ b/docs/database/README.md @@ -1,7 +1,7 @@ # Database 模块 -Database 是 Ayatori 首批实际产品领域之一。当前已包含 Instance 领域基础、管理凭据连接与 -metadata 观察切片,尚未完成 Database API/controller 和 Tenant 供应链路。 +Database 是 Ayatori 首批实际产品领域之一。当前已包含三资源 API、分层绑定与 Instance 原生 +管理能力观测;尚未完成 Database 供应/导入、Tenant 凭据交付与资源回收链路。 ## 当前设计(2026-09-24) @@ -12,7 +12,7 @@ Retain 后人工重新绑定与资源侧 Delete。撤销 PostgreSQL ownership re 依据 [ADR-0009](../decisions/0009-database-resource-and-claim.md),当前合同见 [系统规格](specification.md)。下面的迁移来源与已存在代码不反向约束新设计。 registry adapter、专属迁移/测试及 Instance 的 registry 判定现已撤除;Instance 根据完整管理 -能力观察直接判定 Ready。新增 Database 资源、绑定、导入、角色/凭据管理边界与回收链路尚未实现。wiki 同步位置见 +能力观察直接判定 Ready。Database 资源与绑定已接入,导入、角色/凭据供应及回收仍未完成。wiki 同步位置见 `homelab-wiki/services/postgresql-tenant-operator.md`,跨仓库发布状态由 wiki 的同步记录维护。 ## 来源基线 @@ -59,10 +59,10 @@ Secret metadata 和无关字段变化不重建连接。观测后再次读取 Sec 不把旧连接的成功作为新凭据有效的证据;这不构成跨 Kubernetes/PostgreSQL 的原子事务。 Instance UID、endpoint 或凭据引用变化也会释放旧连接;Forget/Close 只释放本地资源。 -当前通过 `ObserveMetadata` 读取服务器版本和可用扩展,`ObserveVersion` 只是其版本读取便捷入口, -不能产生完整 CapabilityObservation 或 Ready。 -controller 接入、Secret watch、finalizer 与真实权限检查仍待后续切片;并发 CR 更新 -必须由调用者通过 resourceVersion 校验。应用层沿用源实现的串行处理,本阶段未引入新的调度框架。 +`ObserveMetadata` 读取服务器版本和可用扩展,`ObserveVersion` 是版本读取便捷入口;两者 +不会填充管理检查,不能产生 Ready。`ObserveManagement` 使用同一凭据/连接边界读取完整 +原生管理检查,返回绑定当前 target 的 `InstanceObservation`。controller 的资源呈现适配器 +使用 resourceVersion 拒绝过期写入。应用层沿用串行处理,不引入新的调度框架。 运行 `make test-database-integration` 验证真实 API server + 一次性 PostgreSQL;fixture 不接受外部 DSN,镜像固定摘要,使用随机本机回环端口并在退出时删除测试容器。覆盖缺失/错误凭据、RBAC、 @@ -83,14 +83,13 @@ SQL adapter 通过一条只读语句读取 `pg_catalog.current_setting('server_v `InstanceService` 保留原有 CredentialReader → Connector → Database 边界,复用同一个 凭据读取、连接刷新和串行释放流程,不新增连接池封装或任意查询回调。每次重新查询 metadata, -并在 Secret 有效值回读一致后生成不可变的 `MetadataObservation`,绑定本次 target(含当前 +并在 Secret 有效值回读一致后生成不可变的 `InstanceObservation`,绑定本次 target(含当前 generation),不绑定建池时的旧 target。结果不包含凭据,扩展集合不与 driver 的可变 slice 共享。 凭据中途变化、读取失败或查询失败时,返回零值观察并释放连接,不复用旧的扩展列表。 调用方可将 `Target()` 与 `Extensions()` 交给 Instance 的 `ObserveExtensions`;应用调用链 仍负责同轮次使用,不能持久化或跨轮缓存这份证据。metadata 读取不安装扩展、不初始化 registry、 -不设置 Ready,也不授予 Tenant 写权限。管理权限矩阵及 controller 的 -checkpoint/status/finalizer 链路仍是后续切片。 +不设置 Ready,也不授予 Tenant 写权限。管理权限检查使用下面的独立入口。 真实 API server + PostgreSQL 测试验证未安装扩展可被观察、名称保持大小写、search_path 遮蔽 不改变查询来源、低权限账号读取、权限撤回失败与恢复、Secret 中途变化丢弃扩展结果。 @@ -106,7 +105,51 @@ Instance 不再具有 InitializingRegistry 阶段、RegistryState、准备决策 保留凭据读取与连接刷新、TLS、metadata/扩展观察及其真实后端测试。领域测试覆盖每项能力 在初次验证和 Ready 重验时失败、依赖恢复、重启后重新取证、错误目标/阶段及删除保护。 -这不等于管理权限探测矩阵或三资源 controller 已实现。 +这不等于三资源供应、交付和回收已实现。 + +## Instance 原生管理观测 + +2026-09-25 维护者确认先使用原生非 superuser 方案,不引入 SECURITY DEFINER 接口。 +`InspectManagement` 通过同一条只读语句读取当前执行角色的 `CREATEROLE`、`CREATEDB`、 +superuser 属性、服务器可写状态、版本和扩展列表。不创建探针数据库/角色,不初始化 schema。 +只有非 superuser、具备两项原生属性且当前服务器/会话可写时,基础管理能力才通过。 +角色属性不可从继承成员关系推导;具体已有资源仍须检查 owner、membership 与授权范围。 + +权限依据和真实测试对应: + +- role:当前角色具有 CREATEROLE,可创建普通登录角色; +- database:当前角色具有 CREATEDB,且会话非只读、服务器不在 recovery; +- grant:使用自己新建角色的管理权限,显式建立 SET membership,再以 owner 管理数据库 ACL; +- extension:新建数据库 owner 可安装 trusted 扩展;可用列表不是安装授权,非 trusted 或 + 其他前提不满足的扩展仍可能失败,必须逐请求执行和回读。导入不继承此动态供应授权。 + +参考 PostgreSQL 官方 [CREATE ROLE](https://www.postgresql.org/docs/18/sql-createrole.html)、 +[CREATE DATABASE](https://www.postgresql.org/docs/18/sql-createdatabase.html) 与 +[CREATE EXTENSION](https://www.postgresql.org/docs/18/sql-createextension.html)。这些检查是基础 +能力观察,不是未来操作必然成功的保证;权限、容量、连接数等仍可能在执行时变化。 + +`InstanceReconciliation` 协调 finalizer、观察、领域判定和删除引用检查,controller 只连接 +事件、用例、状态呈现与重试。每轮重建无证据的领域对象;旧 Ready 不授权新一轮操作。 +失败清除当前版本结果并撤销 Ready;CR 在观察期间被修改则拒绝旧结果,下一轮重新读取。 +连接/权限变化由 Secret watch 和 30 秒重查驱动,单轮 IO 最长 15 秒;不持续写入相同状态。 + +manager 通过 `--database-secret-namespace`(默认 `POD_NAMESPACE`)启用 Instance 观测; +为空时不启用。本地运行需显式提供该参数。Deployment 使用 downward API 获取自身 namespace; +Secret 的 get/list/watch 权限由该 namespace 的 Role 单独授予,不放入 ClusterRole。 +watch 使用 controller-runtime 的 metadata-only cache,读取有效凭据仍直连 API server。 +TLS 使用既有 endpoint 合同,公开 CA bundle 可由 `--database-root-cert` 指定;不会自动挂载 +生产证书或创建管理 Secret。manager worker 停止后统一关闭 pgxpool。 + +Instance 删除首先释放本地连接并撤销 Ready。任何引用它的 Database(含 Released、删除中) +或动态 Tenant 申请都会阻止 finalizer 解除;列表查询失败也等待。仅在引用全部解除后移除 +`database.ayatori.ddupan.top/instance-protection`,不删除 PostgreSQL、账号或凭据。 +引用查询与删除不是跨对象事务;后续供应仍必须拒绝已删除/删除中的 Instance。 + +验收使用真实 PostgreSQL + API server:原生管理账号实际建库、owner 授权、trusted 扩展 +安装/回读,拒绝非 trusted 扩展;权限撤回/恢复、只读会话、superuser 拒绝和中途轮换。 +实际 manager 在生成的资源 RBAC 和 namespaced Secret Role 下验证缺失 Secret 后出现、 +轮换、删除、跨 namespace 拒绝和 watch。API 测试另覆盖写入版本冲突、幂等、新 reconciler +恢复与引用删除保护。完整 DBaaS 仍需供应/导入、OpenBao/ESO、Retain/Delete 集成验收。 ## 设计入口 diff --git a/docs/database/api-reference.md b/docs/database/api-reference.md index 5d719cd..be022e7 100644 --- a/docs/database/api-reference.md +++ b/docs/database/api-reference.md @@ -53,7 +53,9 @@ Tenant 进入 `status.phase=Binding` 后由 CEL 固定申请目标;Database 绑定顺序由 application service 协调,纯资格规则在领域层;Kubernetes adapter 负责快照 映射、finalizer 和状态呈现。呈现前若资源版本已变化,返回冲突供下一轮重读,不覆盖其他修改。 双向记录完成后 Tenant 为 Bound,Ready=False/BindingComplete,明确尚未供应或交付。 -生成的 manager RBAC 仅授予绑定所需资源读写,不包含 Secret 读取或后端凭据权限。 +生成的 manager ClusterRole 授予资源读写,不包含 Secret 读取;Instance 观测的管理 Secret +权限由固定 namespace 的独立 Role 授予。Instance controller 已接入原生管理观察与引用删除 +保护,启用方式及 Ready 边界见 [模块说明](README.md#instance-原生管理观测)。 当前有 Tenant/Database finalizer 保护,但**删除清理尚未实现**:Tenant 删除报告 Ready=False/DeletionPending 并保留绑定与 finalizer,Database 的保护也不会被自动移除。 diff --git a/docs/database/domain-instance.md b/docs/database/domain-instance.md index e163927..cf6bb79 100644 --- a/docs/database/domain-instance.md +++ b/docs/database/domain-instance.md @@ -1,8 +1,8 @@ # Instance 领域对象规格 -日期:2026-09-24。资源模型修订依据 +日期:2026-09-25。资源模型修订依据 [ADR-0009](../decisions/0009-database-resource-and-claim.md),行为以 -[系统规格](specification.md) 为准。本页替代原 registry 准备与恢复合同;领域依赖已撤除,完整应用/controller 链路尚未接入。 +[系统规格](specification.md) 为准。本页替代原 registry 准备与恢复合同;领域依赖已撤除,Instance 应用/controller 观测链路已接入。 ## 职责 @@ -85,4 +85,6 @@ finalizer 不阻止并发申请 CR 创建;新请求见 Instance 删除中/不 Instance 领域代码、adapter 与测试的 registry 依赖已撤除。AssessManagement 根据完整观察 直接完成验证;AssessReadiness 失败进入 Validating,依赖恢复后重新验证。领域测试覆盖各检查项 -在这两个入口的失败与恢复,但完整权限探测、Instance controller 和三资源生命周期尚未完成。 +在这两个入口的失败与恢复。原生权限检查、Instance controller、metadata-only Secret watch +和引用删除保护已接入,验证矩阵见 [模块说明](README.md#instance-原生管理观测); +Database/Tenant 的供应、交付和回收仍未完成。 diff --git a/docs/database/operations.md b/docs/database/operations.md index 5f3260e..51a8507 100644 --- a/docs/database/operations.md +++ b/docs/database/operations.md @@ -14,6 +14,13 @@ Database 的 `database.ayatori.ddupan.top/database-protection` 也尚无清理 ## 日常检查 +Instance 观察已实现:先确认 manager 配置了 `--database-secret-namespace` 或 `POD_NAMESPACE`, +再检查 Ready Reason、observedGeneration 与管理 Secret 名称/字段映射,切勿导出其 data。 +`InsufficientPrivileges` 表示当前原生方案要求的非 superuser、CREATEDB/CREATEROLE 不满足; +`CredentialsChanged` 会丢弃中途轮换的结果并重验;`InstanceInUse` 消息定位阻塞删除的资源。 +Secret 事件立即入队,30 秒重查覆盖 PostgreSQL 权限等没有 Kubernetes 事件的外部变化。 +Instance 删除不要求 PostgreSQL 可达,但必须可读取所有 Database/Tenant 引用。 + 先看 Instance、Database、Tenant 的 Ready Condition、绑定 UID、阶段与 observedGeneration, 再核对 PostgreSQL catalog、OpenBao metadata、ExternalSecret 与 Secret 投射状态。 具体 kubectl 资源名、finalizer 名称与人工确认字段在 API 实现后补齐,不提供猜测的 patch 命令。 diff --git a/docs/database/security.md b/docs/database/security.md index 24e641e..05ebbe0 100644 --- a/docs/database/security.md +++ b/docs/database/security.md @@ -3,7 +3,7 @@ | 项目 | 内容 | | --- | --- | | 状态 | Review | -| 最后更新 | 2026-09-24 | +| 最后更新 | 2026-09-25 | ## 保护目标 @@ -47,9 +47,11 @@ ESO 身份只读管理路径,租户 ESO 身份只读 tenant base path,二者 使用管理凭据 Store。controller 对管理 Secret 的读取限于自身 namespace,Instance 不能指定其他 namespace;controller 不创建或修改管理 Secret/ExternalSecret。 -PostgreSQL 管理 role 不应是 superuser。若平台选择 SECURITY DEFINER 函数承载创建或 -删除操作,函数必须固定 `search_path`、严格校验 identifier、拒绝任意 SQL,并仅向 -controller role 授予 EXECUTE。controller 不调用 shell 或 `psql` 拼接用户输入。 +2026-09-25 维护者确认第一版使用原生非 superuser 管理 role,具有 CREATEDB/CREATEROLE, +不引入 SECURITY DEFINER 接口。Instance 检查拒绝 superuser;具体已有资源的 owner 和 +membership 仍需逐资源验证,不能把基础能力用于接管他人资源。扩展按实际权限安装, +不因可用列表包含某个扩展就默认能安装它。controller 不调用 shell 或 `psql` 拼接用户输入。 +当前检查与真实权限矩阵见 [Instance 原生管理观测](README.md#instance-原生管理观测)。 Kubernetes RBAC 应把 Instance 管理、Database 导入、Released 重新开放和回收限制给平台管理员。 有权创建 Tenant 的申请者可显式申请未绑定且可用的 Database,不增加资源侧允许绑定名单 diff --git a/docs/database/specification.md b/docs/database/specification.md index 356ac62..5d8c3fc 100644 --- a/docs/database/specification.md +++ b/docs/database/specification.md @@ -171,6 +171,9 @@ status 缺失不假定发生于正常重启。Instance 可重新探测能力;D ## 9. 权限、凭据与扩展 +2026-09-25 确认第一版管理账号使用原生非 superuser + CREATEDB/CREATEROLE 方案, +不引入 SECURITY DEFINER 接口;权限检查与限制见 [安全合同](security.md#最小权限)。 + 动态供应继续使用一个兼任 database owner 的 LOGIN role;应用角色不得具备 superuser、 CREATEDB、CREATEROLE 或 replication 权限。撤销 PUBLIC CONNECT,再授予目标角色; 不修改无关数据库和角色。identifier 匹配 `^[a-z][a-z0-9_]{0,62}$`,SQL 安全引用。 diff --git a/internal/database/adapter/kubernetes/instance_mapping.go b/internal/database/adapter/kubernetes/instance_mapping.go new file mode 100644 index 0000000..7205968 --- /dev/null +++ b/internal/database/adapter/kubernetes/instance_mapping.go @@ -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 +} diff --git a/internal/database/adapter/kubernetes/instance_resources.go b/internal/database/adapter/kubernetes/instance_resources.go new file mode 100644 index 0000000..2c4799a --- /dev/null +++ b/internal/database/adapter/kubernetes/instance_resources.go @@ -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 +} diff --git a/internal/database/adapter/postgresql/fixture_integration_test.go b/internal/database/adapter/postgresql/fixture_integration_test.go index a3feda1..593c487 100644 --- a/internal/database/adapter/postgresql/fixture_integration_test.go +++ b/internal/database/adapter/postgresql/fixture_integration_test.go @@ -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 { diff --git a/internal/database/adapter/postgresql/instance_controller_integration_test.go b/internal/database/adapter/postgresql/instance_controller_integration_test.go new file mode 100644 index 0000000..78e8f5c --- /dev/null +++ b/internal/database/adapter/postgresql/instance_controller_integration_test.go @@ -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 +} diff --git a/internal/database/adapter/postgresql/management.go b/internal/database/adapter/postgresql/management.go new file mode 100644 index 0000000..8ef62eb --- /dev/null +++ b/internal/database/adapter/postgresql/management.go @@ -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 +} diff --git a/internal/database/adapter/postgresql/management_integration_test.go b/internal/database/adapter/postgresql/management_integration_test.go new file mode 100644 index 0000000..c59ecd7 --- /dev/null +++ b/internal/database/adapter/postgresql/management_integration_test.go @@ -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") +} diff --git a/internal/database/application/instance_reconciliation.go b/internal/database/application/instance_reconciliation.go new file mode 100644 index 0000000..ec4d953 --- /dev/null +++ b/internal/database/application/instance_reconciliation.go @@ -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" + } +} diff --git a/internal/database/application/instance_reconciliation_test.go b/internal/database/application/instance_reconciliation_test.go new file mode 100644 index 0000000..73403b7 --- /dev/null +++ b/internal/database/application/instance_reconciliation_test.go @@ -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 入口不应携带完整管理检查") + } +} diff --git a/internal/database/application/instance_service.go b/internal/database/application/instance_service.go index 781d2f8..b1498e0 100644 --- a/internal/database/application/instance_service.go +++ b/internal/database/application/instance_service.go @@ -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 } diff --git a/internal/database/application/instance_service_test.go b/internal/database/application/instance_service_test.go index e6a8e42..d9c2a68 100644 --- a/internal/database/application/instance_service_test.go +++ b/internal/database/application/instance_service_test.go @@ -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 } diff --git a/internal/database/application/metadata.go b/internal/database/application/metadata.go index 2f6f4e4..fd54ec2 100644 --- a/internal/database/application/metadata.go +++ b/internal/database/application/metadata.go @@ -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 } diff --git a/internal/database/controller/instance_controller.go b/internal/database/controller/instance_controller.go new file mode 100644 index 0000000..ff20446 --- /dev/null +++ b/internal/database/controller/instance_controller.go @@ -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 +} diff --git a/internal/database/controller/instance_controller_test.go b/internal/database/controller/instance_controller_test.go new file mode 100644 index 0000000..1c4479e --- /dev/null +++ b/internal/database/controller/instance_controller_test.go @@ -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) +} diff --git a/internal/database/controller/instance_setup.go b/internal/database/controller/instance_setup.go new file mode 100644 index 0000000..fe9097f --- /dev/null +++ b/internal/database/controller/instance_setup.go @@ -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} +}