Compare commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
e2016d3727
|
||
|
|
d71285a3f4 | ||
|
|
7834cab97f
|
||
|
|
72ce3eda40 |
@@ -38,6 +38,16 @@ type PostgreSQLDatabaseStatus struct {
|
|||||||
// +kubebuilder:validation:Minimum=1
|
// +kubebuilder:validation:Minimum=1
|
||||||
// +optional
|
// +optional
|
||||||
CredentialVersion int64 `json:"credentialVersion,omitempty"`
|
CredentialVersion int64 `json:"credentialVersion,omitempty"`
|
||||||
|
// RoleOID 只在角色创建提交并回读成功后保存,后续不按同名认领。
|
||||||
|
// +kubebuilder:validation:Minimum=1
|
||||||
|
// +kubebuilder:validation:Maximum=4294967295
|
||||||
|
// +optional
|
||||||
|
RoleOID int64 `json:"roleOID,omitempty"`
|
||||||
|
// DatabaseOID 只在数据库创建并回读成功后保存;缺失不代表允许认领已有对象。
|
||||||
|
// +kubebuilder:validation:Minimum=1
|
||||||
|
// +kubebuilder:validation:Maximum=4294967295
|
||||||
|
// +optional
|
||||||
|
DatabaseOID int64 `json:"databaseOID,omitempty"`
|
||||||
// Phase 暂不冻结供应子阶段枚举;它不是操作授权或绑定的替代记录。
|
// Phase 暂不冻结供应子阶段枚举;它不是操作授权或绑定的替代记录。
|
||||||
// +optional
|
// +optional
|
||||||
Phase string `json:"phase,omitempty"`
|
Phase string `json:"phase,omitempty"`
|
||||||
@@ -50,6 +60,10 @@ type PostgreSQLDatabaseStatus struct {
|
|||||||
// +kubebuilder:object:root=true
|
// +kubebuilder:object:root=true
|
||||||
// +kubebuilder:subresource:status
|
// +kubebuilder:subresource:status
|
||||||
// +kubebuilder:resource:scope=Cluster
|
// +kubebuilder:resource:scope=Cluster
|
||||||
|
// +kubebuilder:validation:XValidation:rule="!has(self.status) || !has(self.status.roleOID) || (has(self.status.credentialVersion) && has(self.status.instanceUID))",message="roleOID requires confirmed credentials and instanceUID"
|
||||||
|
// +kubebuilder:validation:XValidation:rule="!has(self.status) || !has(self.status.databaseOID) || has(self.status.roleOID)",message="databaseOID requires roleOID"
|
||||||
|
// +kubebuilder:validation:XValidation:rule="!(has(oldSelf.status) && has(oldSelf.status.roleOID)) || (has(self.status) && has(self.status.roleOID) && self.status.roleOID == oldSelf.status.roleOID)",message="confirmed roleOID cannot change or be removed"
|
||||||
|
// +kubebuilder:validation:XValidation:rule="!(has(oldSelf.status) && has(oldSelf.status.databaseOID)) || (has(self.status) && has(self.status.databaseOID) && self.status.databaseOID == oldSelf.status.databaseOID)",message="confirmed databaseOID cannot change or be removed"
|
||||||
// +kubebuilder:validation:XValidation:rule="!has(self.status) || !has(self.status.credentialVersion) || has(self.status.credentialRef)",message="credentialVersion requires credentialRef"
|
// +kubebuilder:validation:XValidation:rule="!has(self.status) || !has(self.status.credentialVersion) || has(self.status.credentialRef)",message="credentialVersion requires credentialRef"
|
||||||
// +kubebuilder:validation:XValidation:rule="!(has(oldSelf.status) && has(oldSelf.status.credentialRef)) || (has(self.status) && has(self.status.credentialRef) && self.status.credentialRef == oldSelf.status.credentialRef)",message="recorded credentialRef cannot change or be removed"
|
// +kubebuilder:validation:XValidation:rule="!(has(oldSelf.status) && has(oldSelf.status.credentialRef)) || (has(self.status) && has(self.status.credentialRef) && self.status.credentialRef == oldSelf.status.credentialRef)",message="recorded credentialRef cannot change or be removed"
|
||||||
// +kubebuilder:validation:XValidation:rule="!(has(oldSelf.status) && has(oldSelf.status.credentialVersion)) || (has(self.status) && has(self.status.credentialVersion) && self.status.credentialVersion == oldSelf.status.credentialVersion)",message="confirmed credentialVersion cannot change or be removed"
|
// +kubebuilder:validation:XValidation:rule="!(has(oldSelf.status) && has(oldSelf.status.credentialVersion)) || (has(self.status) && has(self.status.credentialVersion) && self.status.credentialVersion == oldSelf.status.credentialVersion)",message="confirmed credentialVersion cannot change or be removed"
|
||||||
|
|||||||
@@ -0,0 +1,40 @@
|
|||||||
|
package v1alpha1_test
|
||||||
|
|
||||||
|
import (
|
||||||
|
"testing"
|
||||||
|
|
||||||
|
databasev1alpha1 "git.ddupan.top/panxiao81/ayatori/api/database/v1alpha1"
|
||||||
|
apierrors "k8s.io/apimachinery/pkg/api/errors"
|
||||||
|
"sigs.k8s.io/controller-runtime/pkg/client"
|
||||||
|
)
|
||||||
|
|
||||||
|
func testResourceStatus(t *testing.T, api client.Client) {
|
||||||
|
object := validDatabase("resource-status")
|
||||||
|
if err := api.Create(t.Context(), object); err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
invalid := object.DeepCopy()
|
||||||
|
invalid.Status.RoleOID = 11
|
||||||
|
if err := api.Status().Update(t.Context(), invalid); !apierrors.IsInvalid(err) {
|
||||||
|
t.Fatal("无凭据不能确认角色", err)
|
||||||
|
}
|
||||||
|
object.Status.InstanceUID = "instance-uid"
|
||||||
|
object.Status.CredentialRef = &databasev1alpha1.CredentialReference{Mount: "application-secrets", Path: "applications/test"}
|
||||||
|
object.Status.CredentialVersion = 1
|
||||||
|
object.Status.RoleOID, object.Status.DatabaseOID = 11, 22
|
||||||
|
if err := api.Status().Update(t.Context(), object); err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
for _, change := range []func(*databasev1alpha1.PostgreSQLDatabase){
|
||||||
|
func(d *databasev1alpha1.PostgreSQLDatabase) { d.Status.RoleOID = 0 },
|
||||||
|
func(d *databasev1alpha1.PostgreSQLDatabase) { d.Status.RoleOID = 12 },
|
||||||
|
func(d *databasev1alpha1.PostgreSQLDatabase) { d.Status.DatabaseOID = 0 },
|
||||||
|
func(d *databasev1alpha1.PostgreSQLDatabase) { d.Status.DatabaseOID = 23 },
|
||||||
|
} {
|
||||||
|
invalid = object.DeepCopy()
|
||||||
|
change(invalid)
|
||||||
|
if err := api.Status().Update(t.Context(), invalid); !apierrors.IsInvalid(err) {
|
||||||
|
t.Fatal("已确认 OID 不得覆盖或移除", err)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -62,6 +62,7 @@ func TestDatabaseAPI(t *testing.T) {
|
|||||||
t.Run("拒绝非法声明", func(t *testing.T) { testInvalidDeclarations(t, client) })
|
t.Run("拒绝非法声明", func(t *testing.T) { testInvalidDeclarations(t, client) })
|
||||||
t.Run("status隔离和绑定并发", func(t *testing.T) { testBindingWrites(t, client) })
|
t.Run("status隔离和绑定并发", func(t *testing.T) { testBindingWrites(t, client) })
|
||||||
t.Run("凭据位置与确认版本", func(t *testing.T) { testCredentialStatus(t, client) })
|
t.Run("凭据位置与确认版本", func(t *testing.T) { testCredentialStatus(t, client) })
|
||||||
|
t.Run("资源创建确认", func(t *testing.T) { testResourceStatus(t, client) })
|
||||||
t.Run("仓库示例", func(t *testing.T) { testSamples(t, client, scheme) })
|
t.Run("仓库示例", func(t *testing.T) { testSamples(t, client, scheme) })
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
@@ -221,6 +221,12 @@ spec:
|
|||||||
format: int64
|
format: int64
|
||||||
minimum: 1
|
minimum: 1
|
||||||
type: integer
|
type: integer
|
||||||
|
databaseOID:
|
||||||
|
description: DatabaseOID 只在数据库创建并回读成功后保存;缺失不代表允许认领已有对象。
|
||||||
|
format: int64
|
||||||
|
maximum: 4294967295
|
||||||
|
minimum: 1
|
||||||
|
type: integer
|
||||||
instanceUID:
|
instanceUID:
|
||||||
description: InstanceUID 记录观察时的实例身份,不把同名新实例视为原目标。
|
description: InstanceUID 记录观察时的实例身份,不把同名新实例视为原目标。
|
||||||
type: string
|
type: string
|
||||||
@@ -230,11 +236,28 @@ spec:
|
|||||||
phase:
|
phase:
|
||||||
description: Phase 暂不冻结供应子阶段枚举;它不是操作授权或绑定的替代记录。
|
description: Phase 暂不冻结供应子阶段枚举;它不是操作授权或绑定的替代记录。
|
||||||
type: string
|
type: string
|
||||||
|
roleOID:
|
||||||
|
description: RoleOID 只在角色创建提交并回读成功后保存,后续不按同名认领。
|
||||||
|
format: int64
|
||||||
|
maximum: 4294967295
|
||||||
|
minimum: 1
|
||||||
|
type: integer
|
||||||
type: object
|
type: object
|
||||||
required:
|
required:
|
||||||
- spec
|
- spec
|
||||||
type: object
|
type: object
|
||||||
x-kubernetes-validations:
|
x-kubernetes-validations:
|
||||||
|
- message: roleOID requires confirmed credentials and instanceUID
|
||||||
|
rule: '!has(self.status) || !has(self.status.roleOID) || (has(self.status.credentialVersion)
|
||||||
|
&& has(self.status.instanceUID))'
|
||||||
|
- message: databaseOID requires roleOID
|
||||||
|
rule: '!has(self.status) || !has(self.status.databaseOID) || has(self.status.roleOID)'
|
||||||
|
- message: confirmed roleOID cannot change or be removed
|
||||||
|
rule: '!(has(oldSelf.status) && has(oldSelf.status.roleOID)) || (has(self.status)
|
||||||
|
&& has(self.status.roleOID) && self.status.roleOID == oldSelf.status.roleOID)'
|
||||||
|
- message: confirmed databaseOID cannot change or be removed
|
||||||
|
rule: '!(has(oldSelf.status) && has(oldSelf.status.databaseOID)) || (has(self.status)
|
||||||
|
&& has(self.status.databaseOID) && self.status.databaseOID == oldSelf.status.databaseOID)'
|
||||||
- message: credentialVersion requires credentialRef
|
- message: credentialVersion requires credentialRef
|
||||||
rule: '!has(self.status) || !has(self.status.credentialVersion) || has(self.status.credentialRef)'
|
rule: '!has(self.status) || !has(self.status.credentialVersion) || has(self.status.credentialRef)'
|
||||||
- message: recorded credentialRef cannot change or be removed
|
- message: recorded credentialRef cannot change or be removed
|
||||||
|
|||||||
@@ -56,6 +56,17 @@ Makefile 与 Dockerfile 均继续构建 `cmd/main.go`。
|
|||||||
条件分支,也不为此引入插件注册框架。组件启动失败时释放已装配资源,正常退出则先停止
|
条件分支,也不为此引入插件注册框架。组件启动失败时释放已装配资源,正常退出则先停止
|
||||||
manager worker,再释放连接。
|
manager worker,再释放连接。
|
||||||
|
|
||||||
|
Database 使用 Wire 式的显式构造函数注入,目前手写装配,不引入 Wire 生成器、Dig/Fx
|
||||||
|
容器或运行时服务查找。`database_wiring.go` 集中展示 repository → 用例 → controller 的
|
||||||
|
对象构造;`registerDatabaseControllers` 负责注册,返回的 `closeDatabaseConnections`
|
||||||
|
只在 worker 停止后关闭 Instance 连接。controller 在启动时接收完整依赖,不在 Reconcile
|
||||||
|
或 SetupWithManager 中补建 adapter/service;缺失依赖在注册时失败。
|
||||||
|
|
||||||
|
domain 持有凭据值对象、准备资格与创建恢复规则;application 只组织读写、调用领域判断
|
||||||
|
和保存结果。用例需要的 repository 接口与并发快照仍由消费方定义,不把所有类型都强塞
|
||||||
|
进领域。Kubernetes adapter 将领域阶段映射为已有 Conditions,领域不依赖其字符串协议。
|
||||||
|
单测检查 domain/controller 的依赖边界,并验证共享 writer、直连 reader 和服务的注入。
|
||||||
|
|
||||||
基础设施能力属于整个 controller-manager,不因首个消费者是 Database 就归入该领域。
|
基础设施能力属于整个 controller-manager,不因首个消费者是 Database 就归入该领域。
|
||||||
`internal/infra/openbao` 管理官方 SDK client 的 TLS 配置、Kubernetes 认证及 token 生命周期,
|
`internal/infra/openbao` 管理官方 SDK client 的 TLS 配置、Kubernetes 认证及 token 生命周期,
|
||||||
不依赖 Database 或其他产品领域。Bao client 默认禁用自动重试,写入结果不确定时由用例处理;
|
不依赖 Database 或其他产品领域。Bao client 默认禁用自动重试,写入结果不确定时由用例处理;
|
||||||
|
|||||||
+51
-4
@@ -1,7 +1,7 @@
|
|||||||
# Database 模块
|
# Database 模块
|
||||||
|
|
||||||
Database 是 Ayatori 首批实际产品领域之一。当前已包含三资源 API、分层绑定与 Instance 原生
|
Database 是 Ayatori 首批实际产品领域之一。当前已包含三资源 API、分层绑定与 Instance 原生
|
||||||
管理能力观测;尚未完成 Database 供应/导入、Tenant 凭据交付与资源回收链路。
|
管理能力观测、凭据准备及显式启用的角色/数据库创建;扩展、导入、Tenant 交付与回收尚未完成。
|
||||||
|
|
||||||
## 当前设计(2026-09-24)
|
## 当前设计(2026-09-24)
|
||||||
|
|
||||||
@@ -12,7 +12,7 @@ Retain 后人工重新绑定与资源侧 Delete。撤销 PostgreSQL ownership re
|
|||||||
依据 [ADR-0009](../decisions/0009-database-resource-and-claim.md),当前合同见
|
依据 [ADR-0009](../decisions/0009-database-resource-and-claim.md),当前合同见
|
||||||
[系统规格](specification.md)。下面的迁移来源与已存在代码不反向约束新设计。
|
[系统规格](specification.md)。下面的迁移来源与已存在代码不反向约束新设计。
|
||||||
registry adapter、专属迁移/测试及 Instance 的 registry 判定现已撤除;Instance 根据完整管理
|
registry adapter、专属迁移/测试及 Instance 的 registry 判定现已撤除;Instance 根据完整管理
|
||||||
能力观察直接判定 Ready。Database 资源与绑定已接入,导入、角色/凭据供应及回收仍未完成。wiki 同步位置见
|
能力观察直接判定 Ready。Database 资源与绑定、角色/凭据供应已接入,导入及回收仍未完成。wiki 同步位置见
|
||||||
`homelab-wiki/services/postgresql-tenant-operator.md`,跨仓库发布状态由 wiki 的同步记录维护。
|
`homelab-wiki/services/postgresql-tenant-operator.md`,跨仓库发布状态由 wiki 的同步记录维护。
|
||||||
|
|
||||||
## 来源基线
|
## 来源基线
|
||||||
@@ -183,6 +183,10 @@ Database 已有 `status.credentialRef` 和 `status.credentialVersion` 的字段
|
|||||||
顺序为固定位置 → 确认后端尚无凭据 → 保存 CreationStarted → 创建并回读 → 保存确认版本。
|
顺序为固定位置 → 确认后端尚无凭据 → 保存 CreationStarted → 创建并回读 → 保存确认版本。
|
||||||
`CredentialReconciler` 只负责 Database/Tenant/Instance watch 和 30 秒依赖重查;
|
`CredentialReconciler` 只负责 Database/Tenant/Instance watch 和 30 秒依赖重查;
|
||||||
Kubernetes repository 负责直接读取及有 resourceVersion 保护的状态更新。
|
Kubernetes repository 负责直接读取及有 resourceVersion 保护的状态更新。
|
||||||
|
三类 controller 的依赖均在 bootstrap 显式组装,不在 Reconcile 中构造服务。
|
||||||
|
`domain/credential` 保存应用凭据值对象、准备资格、固定位置与未确认创建的恢复规则;
|
||||||
|
application 保留 I/O 顺序、消费方接口及并发快照,不再复用绑定用例的快照类型。
|
||||||
|
领域阶段与 `CredentialsReady`/Reason 的转换由 Kubernetes adapter 负责,已有 API 保持兼容。
|
||||||
|
|
||||||
外部操作前后回查目标:Database 必须仍是同一 UID/resourceVersion,Tenant/Instance 必须
|
外部操作前后回查目标:Database 必须仍是同一 UID/resourceVersion,Tenant/Instance 必须
|
||||||
保持绑定、spec generation、删除状态、finalizer 保护与有效 Ready;无关 Conditions 刷新
|
保持绑定、spec generation、删除状态、finalizer 保护与有效 Ready;无关 Conditions 刷新
|
||||||
@@ -197,14 +201,57 @@ Kubernetes repository 负责直接读取及有 resourceVersion 保护的状态
|
|||||||
同样可能保守地要求人工处理,不承诺无损接续;多副本部署应启用既有 leader election。
|
同样可能保守地要求人工处理,不承诺无损接续;多副本部署应启用既有 leader election。
|
||||||
|
|
||||||
成功只设置 `CredentialsReady=True`,Database Ready 仍为 False/ProvisioningIncomplete,
|
成功只设置 `CredentialsReady=True`,Database Ready 仍为 False/ProvisioningIncomplete,
|
||||||
Tenant 仍未完成交付。本切片没有 PostgreSQL role/database 创建、扩展安装、ESO 投射、
|
Tenant 仍未完成交付。凭据准备本身不创建 PostgreSQL 资源;后续创建需另行启用下述用例。
|
||||||
Retain 释放或 Delete 清理,也不会解除 finalizer。不要作为完整 DBaaS 部署。
|
扩展安装、ESO 投射、Retain 释放或 Delete 清理尚未接入,也不会解除 finalizer。不要作为完整 DBaaS 部署。
|
||||||
|
|
||||||
单元测试穷举前置条件;真实 API server + 隔离 Bao 验证创建、状态确认、幂等、重启、并发、
|
单元测试穷举前置条件;真实 API server + 隔离 Bao 验证创建、状态确认、幂等、重启、并发、
|
||||||
依赖恢复、固定位置、不确定结果、确认保存失败和删除边界;实际 manager 验证 watch 驱动
|
依赖恢复、固定位置、不确定结果、确认保存失败和删除边界;实际 manager 验证 watch 驱动
|
||||||
及重启。该 fixture 只声明 Instance 前置 Ready,真实 PostgreSQL 管理能力由既有 Instance
|
及重启。该 fixture 只声明 Instance 前置 Ready,真实 PostgreSQL 管理能力由既有 Instance
|
||||||
集成测试覆盖,不把凭据准备验收当成实际建库或应用登录验收。
|
集成测试覆盖,不把凭据准备验收当成实际建库或应用登录验收。
|
||||||
|
|
||||||
|
## PostgreSQL 资源创建闭环
|
||||||
|
|
||||||
|
在管理 Secret、OpenBao 认证和凭据准备配置之上,显式设置 `--database-provision-resources`
|
||||||
|
才启用外部角色/数据库写入;默认关闭。升级时先安装新 CRD,供应字段见
|
||||||
|
[API 合同](api-reference.md#当前-api-切片)。只处理 Provision 来源,导入资源不进入创建流程。
|
||||||
|
|
||||||
|
`DatabaseProvisioning` 在每轮验证双向绑定、UID、删除/Released 状态、保护和当前 Instance
|
||||||
|
Ready,读取固定位置的已确认凭据。所有 SQL 操作复用 InstanceService 已有的管理连接及
|
||||||
|
Secret 刷新机制,不再开一个管理连接池。执行前重新查询非 superuser 管理能力;外部操作
|
||||||
|
前后直接回读 API 快照,发生并发变更则停止确认,不盲目覆盖 status。
|
||||||
|
启用资源创建时,`DatabaseReconciliation` 在同一条 reconcile 中先运行凭据准备/检查,
|
||||||
|
再运行资源供应用例,不同时注册独立凭据 controller,避免两个 worker 争写同一资源状态。
|
||||||
|
|
||||||
|
创建分为三个可观察步骤:
|
||||||
|
|
||||||
|
1. 保存 CreatingRole,事务内创建无管理特权的 LOGIN owner 并授予管理账号 SET 权限;
|
||||||
|
提交并回读成功后保存 `status.roleOID`。
|
||||||
|
2. 保存 CreatingDatabase,创建以已确认角色为 owner 且 `ALLOW_CONNECTIONS false` 的数据库;
|
||||||
|
回读成功后保存 `status.databaseOID`。
|
||||||
|
3. 在已确认对象上以 owner 收紧 ACL:撤销 PUBLIC CONNECT、授予 owner CONNECT,再开放
|
||||||
|
连接入口。ACL 收敛可幂等重试,不重置密码或更换 owner。
|
||||||
|
|
||||||
|
OID 只保存在 CR,表示本次成功创建的回读身份,不是 registry 或认领机制。未知同名对象、
|
||||||
|
已确认对象消失/被重建、owner/角色特权漂移均报 Conflict。CreatingRole/CreatingDatabase
|
||||||
|
未留下对应 OID 时,下轮保守停止;即使请求可能尚未发出,也不尝试推断或自动补记。
|
||||||
|
网络结果不确定、创建成功后确认写入失败同样交给人工处理。诊断包含 Database、Instance、
|
||||||
|
实际 database/role、失败步骤;确认 OID 独立保留。依赖或明确权限拒绝可等待恢复。
|
||||||
|
OID 不是跨集群/备份恢复的稳定身份,恢复后须人工核对,不能靠匹配 OID 推导管理权。
|
||||||
|
|
||||||
|
`ResourcesReady=True` 仅表示角色、数据库及连接 ACL 已确认,不代表扩展或 Tenant 交付完成;
|
||||||
|
Database/Tenant 总体 Ready 仍为 False。当前不轮换密码、不做运行期应用登录健康检查,也不
|
||||||
|
撤销其他数据库的 PUBLIC 权限。管理员须保证共享实例中其他数据库的接入策略满足隔离要求。
|
||||||
|
|
||||||
|
真实 API server + PostgreSQL + OpenBao 测试覆盖实际密码登录、owner 建表、应用管理权限拒绝、
|
||||||
|
无关账号连接拒绝、重试/重启、权限恢复、同名冲突、角色重建、并发授权、创建结果丢失、
|
||||||
|
确认持久化失败、删除停止与实际 manager 的观察链路。保留外部操作/API 写入间的非原子边界,
|
||||||
|
不承诺控制面与数据库之间的事务或对任意管理员并发 DDL 的无损恢复。
|
||||||
|
|
||||||
|
遵循 PostgreSQL 官方的
|
||||||
|
[CREATE DATABASE 非事务与 owner 权限合同](https://www.postgresql.org/docs/17/sql-createdatabase.html)
|
||||||
|
和 [CREATEROLE 的成员授权](https://www.postgresql.org/docs/17/role-attributes.html)。
|
||||||
|
不新增 controller 通用生命周期框架,继续采用既有直接观察、resourceVersion 保护和条件报告。
|
||||||
|
|
||||||
## OpenBao Kubernetes 认证会话
|
## OpenBao Kubernetes 认证会话
|
||||||
|
|
||||||
公共 `internal/infra/openbao.KubernetesSession` 复用官方 Kubernetes auth helper 和 `LifetimeWatcher`
|
公共 `internal/infra/openbao.KubernetesSession` 复用官方 Kubernetes auth helper 和 `LifetimeWatcher`
|
||||||
|
|||||||
@@ -2,9 +2,9 @@
|
|||||||
|
|
||||||
| 项目 | 内容 |
|
| 项目 | 内容 |
|
||||||
| --- | --- |
|
| --- | --- |
|
||||||
| 状态 | API schema 与绑定 controller 已实现;供应、交付与删除清理未接入 |
|
| 状态 | API、绑定、凭据准备及角色/数据库创建已实现;扩展、导入、交付与删除清理未接入 |
|
||||||
| API group/version | `database.ayatori.ddupan.top/v1alpha1` |
|
| API group/version | `database.ayatori.ddupan.top/v1alpha1` |
|
||||||
| 最后更新 | 2026-09-27 |
|
| 最后更新 | 2026-09-29 |
|
||||||
|
|
||||||
以 [系统规格](specification.md) 与
|
以 [系统规格](specification.md) 与
|
||||||
[ADR-0009](../decisions/0009-database-resource-and-claim.md) 为准。类型与生成的 CRD 已纳入源码,
|
[ADR-0009](../decisions/0009-database-resource-and-claim.md) 为准。类型与生成的 CRD 已纳入源码,
|
||||||
@@ -26,6 +26,7 @@ Go 类型位于 `api/database/v1alpha1`,CRD 随 `config/crd` 发布;manager
|
|||||||
| Database | `status.instanceUID` | 观察时的 Instance 身份 |
|
| Database | `status.instanceUID` | 观察时的 Instance 身份 |
|
||||||
| Database | `status.credentialRef.mount/path` | 首次写入前固定的 KV v2 位置,不随部署配置迁移 |
|
| Database | `status.credentialRef.mount/path` | 首次写入前固定的 KV v2 位置,不随部署配置迁移 |
|
||||||
| Database | `status.credentialVersion` | 创建并回读成功后确认的正整数版本;省略表示未确认 |
|
| Database | `status.credentialVersion` | 创建并回读成功后确认的正整数版本;省略表示未确认 |
|
||||||
|
| Database | `status.roleOID`、`status.databaseOID` | 分别创建并回读成功后确认的 PostgreSQL OID;不能凭同名补记 |
|
||||||
| Tenant | `spec.provision.instanceRef.name` | 动态申请来源,与 `spec.databaseRef` 互斥且必须二选一 |
|
| Tenant | `spec.provision.instanceRef.name` | 动态申请来源,与 `spec.databaseRef` 互斥且必须二选一 |
|
||||||
| Tenant | `spec.provision.database/loginRole` | 可省略,语义默认值由 controller 解析,不由 CRD 推导 |
|
| Tenant | `spec.provision.database/loginRole` | 可省略,语义默认值由 controller 解析,不由 CRD 推导 |
|
||||||
| Tenant | `spec.databaseRef.name` | 显式申请已有 Database,不额外指定 Instance |
|
| Tenant | `spec.databaseRef.name` | 显式申请已有 Database,不额外指定 Instance |
|
||||||
@@ -47,6 +48,13 @@ KV 版本,删除或版本漂移均需人工处理,不回退旧版本或生
|
|||||||
`CreationStarted` 且无确认版本表示创建未完成确认,重入时停在 Conflict;不尝试推断
|
`CreationStarted` 且无确认版本表示创建未完成确认,重入时停在 Conflict;不尝试推断
|
||||||
进程中断前请求是否发出。完整执行和测试边界见[凭据准备闭环](README.md#凭据准备闭环)。
|
进程中断前请求是否发出。完整执行和测试边界见[凭据准备闭环](README.md#凭据准备闭环)。
|
||||||
|
|
||||||
|
角色/数据库 OID 为 1–4294967295,写入后不可修改或移除。角色确认要求已有凭据确认与
|
||||||
|
Instance UID;数据库确认要求先有角色确认。`ResourcesReady` 的 CreatingRole/CreatingDatabase
|
||||||
|
没有相应 OID 时,下轮转 Conflict;依赖故障不清空确认。旧 CRD 裁剪 OID 会阻止继续下一步,
|
||||||
|
但不能回滚已发出的外部创建,因此必须先升级 CRD 再启用供应。
|
||||||
|
`ResourcesReady=True/Available` 只表示角色、数据库和 ACL 已确认,总体 Ready 仍不放行。
|
||||||
|
详见[资源创建闭环](README.md#postgresql-资源创建闭环)。
|
||||||
|
|
||||||
示例:[Instance](../../config/samples/database_v1alpha1_postgresqlinstance.yaml)、
|
示例:[Instance](../../config/samples/database_v1alpha1_postgresqlinstance.yaml)、
|
||||||
[导入 Database](../../config/samples/database_v1alpha1_postgresqldatabase.yaml)、
|
[导入 Database](../../config/samples/database_v1alpha1_postgresqldatabase.yaml)、
|
||||||
[动态/已有资源申请](../../config/samples/database_v1alpha1_postgresqltenant.yaml)。
|
[动态/已有资源申请](../../config/samples/database_v1alpha1_postgresqltenant.yaml)。
|
||||||
@@ -131,7 +139,7 @@ Database 是平台管理的集群级资源,不归属于应用 namespace,也
|
|||||||
普通申请者通过 Tenant 申请使用,不能自行修改 Database 回收策略或将 Released 资源重新开放;
|
普通申请者通过 Tenant 申请使用,不能自行修改 Database 回收策略或将 Released 资源重新开放;
|
||||||
这些资源管理操作由平台管理员授权。controller 的绑定协调权限与用户申请权限分别配置。
|
这些资源管理操作由平台管理员授权。controller 的绑定协调权限与用户申请权限分别配置。
|
||||||
|
|
||||||
以下是行为合同,具体 schema 见当前 API 切片与生成的 CRD;后端行为尚未实现:
|
以下是行为合同,具体 schema 见当前 API 切片与生成的 CRD;导入、交付与回收尚未实现:
|
||||||
|
|
||||||
| 内容 | 合同 |
|
| 内容 | 合同 |
|
||||||
| --- | --- |
|
| --- | --- |
|
||||||
|
|||||||
@@ -1,6 +1,6 @@
|
|||||||
# 部署与配置
|
# 部署与配置
|
||||||
|
|
||||||
> 本页区分已实现的 Instance 观测、认证和凭据准备配置,与尚未接入的 PostgreSQL 供应/交付合同。
|
> 本页区分已实现的 Instance 观测、认证、凭据准备和资源创建配置,与尚未接入的完整供应/交付合同。
|
||||||
> 完整 Database 服务仍不可部署使用;当前可执行入口见 [模块说明](README.md)。
|
> 完整 Database 服务仍不可部署使用;当前可执行入口见 [模块说明](README.md)。
|
||||||
|
|
||||||
| 项目 | 内容 |
|
| 项目 | 内容 |
|
||||||
@@ -46,6 +46,7 @@ Instance 观测)与 `--database-root-cert`(公开 PostgreSQL CA PEM 路径
|
|||||||
| `--openbao-auth-role` | 已实现,启用时必填 | OpenBao 登录 role |
|
| `--openbao-auth-role` | 已实现,启用时必填 | OpenBao 登录 role |
|
||||||
| `--openbao-ca-cert` | 已实现,默认系统信任根 | OpenBao 公开 CA PEM 路径 |
|
| `--openbao-ca-cert` | 已实现,默认系统信任根 | OpenBao 公开 CA PEM 路径 |
|
||||||
| `--database-credential-mount` | 已实现,默认空 | 显式设置后启用应用凭据准备,要求已配置 OpenBao 认证 |
|
| `--database-credential-mount` | 已实现,默认空 | 显式设置后启用应用凭据准备,要求已配置 OpenBao 认证 |
|
||||||
|
| `--database-provision-resources` | 已实现,默认 false | 启用角色/数据库及 ACL 创建;要求管理 Secret namespace、凭据 mount 和认证配置,先升级 CRD |
|
||||||
| `--openbao-service-account-namespace` | 已实现,启用时必填 | TokenRequest 目标 SA 的固定 namespace |
|
| `--openbao-service-account-namespace` | 已实现,启用时必填 | TokenRequest 目标 SA 的固定 namespace |
|
||||||
| `--openbao-service-account-name` | 已实现,启用时必填 | TokenRequest 目标 SA 名称 |
|
| `--openbao-service-account-name` | 已实现,启用时必填 | TokenRequest 目标 SA 名称 |
|
||||||
| `--openbao-token-audience` | 已实现,`openbao` | SA JWT audience,必须匹配 OpenBao role |
|
| `--openbao-token-audience` | 已实现,`openbao` | SA JWT audience,必须匹配 OpenBao role |
|
||||||
|
|||||||
@@ -21,6 +21,7 @@ type databaseOptions struct {
|
|||||||
rootCert string
|
rootCert string
|
||||||
credentialMount string
|
credentialMount string
|
||||||
credentialPrefix string
|
credentialPrefix string
|
||||||
|
provisionResources bool
|
||||||
}
|
}
|
||||||
|
|
||||||
func (o *databaseOptions) bindFlags(flags *flag.FlagSet) {
|
func (o *databaseOptions) bindFlags(flags *flag.FlagSet) {
|
||||||
@@ -29,6 +30,7 @@ func (o *databaseOptions) bindFlags(flags *flag.FlagSet) {
|
|||||||
flags.StringVar(&o.rootCert, "database-root-cert", "", "PostgreSQL 管理连接信任的公开 CA bundle 路径")
|
flags.StringVar(&o.rootCert, "database-root-cert", "", "PostgreSQL 管理连接信任的公开 CA bundle 路径")
|
||||||
flags.StringVar(&o.credentialMount, "database-credential-mount", "", "应用凭据 KV v2 mount;为空时不启用凭据准备")
|
flags.StringVar(&o.credentialMount, "database-credential-mount", "", "应用凭据 KV v2 mount;为空时不启用凭据准备")
|
||||||
flags.StringVar(&o.credentialPrefix, "database-credential-prefix", "applications", "应用凭据路径前缀;已有固定位置不随配置变化迁移")
|
flags.StringVar(&o.credentialPrefix, "database-credential-prefix", "applications", "应用凭据路径前缀;已有固定位置不随配置变化迁移")
|
||||||
|
flags.BoolVar(&o.provisionResources, "database-provision-resources", false, "启用 PostgreSQL 角色和数据库创建;需要 Instance 管理凭据及 OpenBao 凭据准备")
|
||||||
}
|
}
|
||||||
|
|
||||||
func (o databaseOptions) configureManager(options *ctrl.Options) {
|
func (o databaseOptions) configureManager(options *ctrl.Options) {
|
||||||
@@ -37,28 +39,33 @@ func (o databaseOptions) configureManager(options *ctrl.Options) {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
// setupDatabase 封装 Database 的内部装配,并返回在 manager 停止后执行的清理。
|
// registerDatabaseControllers 封装 Database 的内部装配,并返回在 manager 停止后执行的清理。
|
||||||
func setupDatabase(ctx context.Context, manager ctrl.Manager, options databaseOptions, baoClient *bao.Client) (func(), error) {
|
func registerDatabaseControllers(ctx context.Context, manager ctrl.Manager, options databaseOptions, baoClient *bao.Client) (func(), error) {
|
||||||
cleanup := func() {}
|
closeDatabaseConnections := func() {}
|
||||||
|
var backend *application.InstanceService
|
||||||
if options.secretNamespace != "" {
|
if options.secretNamespace != "" {
|
||||||
service, err := setupInstanceObservation(manager, options.secretNamespace, options.rootCert)
|
service, err := setupInstanceObservation(manager, options.secretNamespace, options.rootCert)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return nil, fmt.Errorf("set up Instance observation: %w", err)
|
return nil, fmt.Errorf("set up Instance observation: %w", err)
|
||||||
}
|
}
|
||||||
cleanup = service.Close
|
closeDatabaseConnections = service.Close
|
||||||
|
backend = service
|
||||||
}
|
}
|
||||||
if err := (&databasecontroller.BindingReconciler{}).SetupWithManager(ctx, manager); err != nil {
|
if err := wireBindingController(manager.GetClient(), manager.GetAPIReader()).SetupWithManager(ctx, manager); err != nil {
|
||||||
cleanup()
|
closeDatabaseConnections()
|
||||||
return nil, fmt.Errorf("set up Database binding controller: %w", err)
|
return nil, fmt.Errorf("set up Database binding controller: %w", err)
|
||||||
}
|
}
|
||||||
if err := setupCredentialPreparation(manager, options, baoClient); err != nil {
|
if err := registerDatabaseSupplyController(manager, options, baoClient, backend); err != nil {
|
||||||
cleanup()
|
closeDatabaseConnections()
|
||||||
return nil, fmt.Errorf("set up Database credential preparation: %w", err)
|
return nil, fmt.Errorf("register Database supply controller: %w", err)
|
||||||
}
|
}
|
||||||
return cleanup, nil
|
return closeDatabaseConnections, nil
|
||||||
}
|
}
|
||||||
|
|
||||||
func setupCredentialPreparation(manager ctrl.Manager, options databaseOptions, baoClient *bao.Client) error {
|
func registerDatabaseSupplyController(manager ctrl.Manager, options databaseOptions, baoClient *bao.Client, backend *application.InstanceService) error {
|
||||||
|
if options.provisionResources && (backend == nil || options.credentialMount == "") {
|
||||||
|
return fmt.Errorf("database resource provisioning requires management Secret namespace and credential mount")
|
||||||
|
}
|
||||||
if options.credentialMount == "" {
|
if options.credentialMount == "" {
|
||||||
return nil
|
return nil
|
||||||
}
|
}
|
||||||
@@ -69,7 +76,10 @@ func setupCredentialPreparation(manager ctrl.Manager, options databaseOptions, b
|
|||||||
if err != nil {
|
if err != nil {
|
||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
return (&databasecontroller.CredentialReconciler{Store: store}).SetupWithManager(manager)
|
if options.provisionResources {
|
||||||
|
return wireProvisioningController(manager.GetClient(), manager.GetAPIReader(), store, backend).SetupWithManager(manager)
|
||||||
|
}
|
||||||
|
return wireCredentialController(manager.GetClient(), manager.GetAPIReader(), store).SetupWithManager(manager)
|
||||||
}
|
}
|
||||||
|
|
||||||
func setupInstanceObservation(manager ctrl.Manager, namespace, rootCert string) (*application.InstanceService, error) {
|
func setupInstanceObservation(manager ctrl.Manager, namespace, rootCert string) (*application.InstanceService, error) {
|
||||||
@@ -81,7 +91,7 @@ func setupInstanceObservation(manager ctrl.Manager, namespace, rootCert string)
|
|||||||
if err != nil {
|
if err != nil {
|
||||||
return nil, err
|
return nil, err
|
||||||
}
|
}
|
||||||
reconciler := &databasecontroller.InstanceReconciler{Observer: service, SecretNamespace: namespace}
|
reconciler := wireInstanceController(manager.GetClient(), manager.GetAPIReader(), service, namespace)
|
||||||
if err := reconciler.SetupWithManager(manager); err != nil {
|
if err := reconciler.SetupWithManager(manager); err != nil {
|
||||||
service.Close()
|
service.Close()
|
||||||
return nil, err
|
return nil, err
|
||||||
|
|||||||
@@ -0,0 +1,38 @@
|
|||||||
|
package bootstrap
|
||||||
|
|
||||||
|
import (
|
||||||
|
"git.ddupan.top/panxiao81/ayatori/internal/database/adapter/kubernetes"
|
||||||
|
"git.ddupan.top/panxiao81/ayatori/internal/database/application"
|
||||||
|
databasecontroller "git.ddupan.top/panxiao81/ayatori/internal/database/controller"
|
||||||
|
"sigs.k8s.io/controller-runtime/pkg/client"
|
||||||
|
)
|
||||||
|
|
||||||
|
// 本文件是 Database 的显式依赖图:adapter → 用例 → controller。
|
||||||
|
// 构造只在启动时执行;用例和 controller 都不持有容器或动态查找依赖。
|
||||||
|
func wireBindingController(writer client.Client, reader client.Reader) *databasecontroller.BindingReconciler {
|
||||||
|
resources := &kubernetes.BindingResources{Client: writer, Reader: reader}
|
||||||
|
service := &application.BindingService{Resources: resources}
|
||||||
|
return databasecontroller.NewBindingReconciler(writer, service, resources)
|
||||||
|
}
|
||||||
|
|
||||||
|
func wireCredentialController(writer client.Client, reader client.Reader, store application.CredentialStore) *databasecontroller.CredentialReconciler {
|
||||||
|
resources := &kubernetes.CredentialResources{Client: writer, Reader: reader}
|
||||||
|
service := &application.CredentialPreparation{Resources: resources, Store: store}
|
||||||
|
return databasecontroller.NewCredentialReconciler(writer, service)
|
||||||
|
}
|
||||||
|
|
||||||
|
func wireInstanceController(writer client.Client, reader client.Reader, observer application.InstanceObserver, namespace string) *databasecontroller.InstanceReconciler {
|
||||||
|
resources := &kubernetes.InstanceResources{Client: writer, Reader: reader}
|
||||||
|
service := &application.InstanceReconciliation{Resources: resources, Observer: observer}
|
||||||
|
return databasecontroller.NewInstanceReconciler(writer, service, resources, namespace)
|
||||||
|
}
|
||||||
|
|
||||||
|
func wireProvisioningController(writer client.Client, reader client.Reader, store application.CredentialStore, backend application.ProvisioningBackend) *databasecontroller.ProvisioningReconciler {
|
||||||
|
resources := &kubernetes.ProvisioningResources{Client: writer, Reader: reader}
|
||||||
|
credentialResources := &kubernetes.CredentialResources{Client: writer, Reader: reader}
|
||||||
|
service := &application.DatabaseReconciliation{
|
||||||
|
Credentials: &application.CredentialPreparation{Resources: credentialResources, Store: store},
|
||||||
|
Provisioning: &application.DatabaseProvisioning{Resources: resources, Credentials: store, Backend: backend},
|
||||||
|
}
|
||||||
|
return databasecontroller.NewProvisioningReconciler(writer, service)
|
||||||
|
}
|
||||||
@@ -0,0 +1,79 @@
|
|||||||
|
package bootstrap
|
||||||
|
|
||||||
|
import (
|
||||||
|
"go/parser"
|
||||||
|
"go/token"
|
||||||
|
"io/fs"
|
||||||
|
"path/filepath"
|
||||||
|
"strconv"
|
||||||
|
"strings"
|
||||||
|
"testing"
|
||||||
|
|
||||||
|
"git.ddupan.top/panxiao81/ayatori/internal/database/adapter/kubernetes"
|
||||||
|
"git.ddupan.top/panxiao81/ayatori/internal/database/adapter/postgresql"
|
||||||
|
"git.ddupan.top/panxiao81/ayatori/internal/database/application"
|
||||||
|
"sigs.k8s.io/controller-runtime/pkg/client/fake"
|
||||||
|
)
|
||||||
|
|
||||||
|
const controllerLayer = "controller"
|
||||||
|
|
||||||
|
func TestDatabaseExplicitWiring(t *testing.T) {
|
||||||
|
writer := fake.NewClientBuilder().Build()
|
||||||
|
directReader := fake.NewClientBuilder().Build()
|
||||||
|
binder := wireBindingController(writer, directReader)
|
||||||
|
bindingResources, ok := binder.Service.Resources.(*kubernetes.BindingResources)
|
||||||
|
if !ok || binder.Client != writer || bindingResources.Reader != directReader || bindingResources.Client != writer || binder.Presenter != bindingResources {
|
||||||
|
t.Fatal("绑定用例没有共享显式注入的 writer、直连 reader 与 presenter")
|
||||||
|
}
|
||||||
|
credentials, err := kubernetes.NewSecretCredentials(directReader, "wiring-tests")
|
||||||
|
if err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
observer, err := application.NewInstanceService(credentials, postgresql.Connector{})
|
||||||
|
if err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
t.Cleanup(observer.Close)
|
||||||
|
reconciler := wireInstanceController(writer, directReader, observer, "wiring-tests")
|
||||||
|
resources, ok := reconciler.Service.Resources.(*kubernetes.InstanceResources)
|
||||||
|
if !ok || resources.Reader != directReader || resources.Client != writer || reconciler.Service.Observer != observer || reconciler.Presenter != resources {
|
||||||
|
t.Fatal("Instance 的服务或读取边界未按依赖图注入")
|
||||||
|
}
|
||||||
|
provisioner := wireProvisioningController(writer, directReader, nil, observer)
|
||||||
|
provisioningResources, ok := provisioner.Service.Provisioning.Resources.(*kubernetes.ProvisioningResources)
|
||||||
|
if !ok || provisioningResources.Client != writer || provisioningResources.Reader != directReader || provisioner.Service.Provisioning.Backend != observer {
|
||||||
|
t.Fatal("资源供应必须复用已有 Instance 管理服务和 API 客户端")
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// 防止领域重新依赖用例/存储,也防止 controller 再次私自构造具体 adapter。
|
||||||
|
func TestDatabaseLayerImports(t *testing.T) {
|
||||||
|
for _, layer := range []string{"domain", controllerLayer} {
|
||||||
|
err := filepath.WalkDir(filepath.Join("../database", layer), func(path string, entry fs.DirEntry, err error) error {
|
||||||
|
if err != nil {
|
||||||
|
return err
|
||||||
|
}
|
||||||
|
if entry.IsDir() || !strings.HasSuffix(path, ".go") || strings.HasSuffix(path, "_test.go") {
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
file, err := parser.ParseFile(token.NewFileSet(), path, nil, parser.ImportsOnly)
|
||||||
|
if err != nil {
|
||||||
|
return err
|
||||||
|
}
|
||||||
|
for _, dependency := range file.Imports {
|
||||||
|
name, err := strconv.Unquote(dependency.Path.Value)
|
||||||
|
if err != nil {
|
||||||
|
return err
|
||||||
|
}
|
||||||
|
if strings.Contains(name, "/database/adapter/") || (layer == "domain" &&
|
||||||
|
(strings.Contains(name, "/database/application") || strings.Contains(name, "k8s.io/") || strings.Contains(name, "/internal/infra/"))) {
|
||||||
|
t.Errorf("%s 不得导入 %s", path, name)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
return nil
|
||||||
|
})
|
||||||
|
if err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -54,11 +54,11 @@ func TestBootstrapWithRealAPIServer(t *testing.T) {
|
|||||||
}
|
}
|
||||||
ctx, cancel := context.WithTimeout(t.Context(), 30*time.Second)
|
ctx, cancel := context.WithTimeout(t.Context(), 30*time.Second)
|
||||||
defer cancel()
|
defer cancel()
|
||||||
cleanup, err := setupDatabase(ctx, manager, options.database, baoClient)
|
closeDatabaseConnections, err := registerDatabaseControllers(ctx, manager, options.database, baoClient)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
t.Fatal(err)
|
t.Fatal(err)
|
||||||
}
|
}
|
||||||
defer cleanup()
|
defer closeDatabaseConnections()
|
||||||
// 空 API 中没有供应目标;此处验证启用路径确实注册 controller,不访问外部 Bao。
|
// 空 API 中没有供应目标;此处验证启用路径确实注册 controller,不访问外部 Bao。
|
||||||
fixtureConfig := bao.NewConfig()
|
fixtureConfig := bao.NewConfig()
|
||||||
fixtureConfig.Address = "http://127.0.0.1:1"
|
fixtureConfig.Address = "http://127.0.0.1:1"
|
||||||
@@ -67,7 +67,7 @@ func TestBootstrapWithRealAPIServer(t *testing.T) {
|
|||||||
t.Fatal(err)
|
t.Fatal(err)
|
||||||
}
|
}
|
||||||
options.database.credentialMount = "secret"
|
options.database.credentialMount = "secret"
|
||||||
if err := setupCredentialPreparation(manager, options.database, fixtureClient); err != nil {
|
if err := registerDatabaseSupplyController(manager, options.database, fixtureClient, nil); err != nil {
|
||||||
t.Fatal(err)
|
t.Fatal(err)
|
||||||
}
|
}
|
||||||
done := make(chan error, 1)
|
done := make(chan error, 1)
|
||||||
@@ -80,7 +80,7 @@ func TestBootstrapWithRealAPIServer(t *testing.T) {
|
|||||||
t.Error(err)
|
t.Error(err)
|
||||||
}
|
}
|
||||||
case <-time.After(20 * time.Second):
|
case <-time.After(20 * time.Second):
|
||||||
t.Error("manager did not stop before component cleanup")
|
t.Error("manager did not stop before component closeDatabaseConnections")
|
||||||
}
|
}
|
||||||
}()
|
}()
|
||||||
if !manager.GetCache().WaitForCacheSync(ctx) {
|
if !manager.GetCache().WaitForCacheSync(ctx) {
|
||||||
|
|||||||
@@ -63,12 +63,15 @@ func TestCredentialPreparationOptions(t *testing.T) {
|
|||||||
if options.database.credentialMount != "applications-kv" || options.database.credentialPrefix != "database" {
|
if options.database.credentialMount != "applications-kv" || options.database.credentialPrefix != "database" {
|
||||||
t.Fatal("凭据准备参数未传入领域装配")
|
t.Fatal("凭据准备参数未传入领域装配")
|
||||||
}
|
}
|
||||||
if err := setupCredentialPreparation(nil, options.database, nil); err == nil {
|
if err := registerDatabaseSupplyController(nil, options.database, nil, nil); err == nil {
|
||||||
t.Fatal("启用凭据准备必须有显式配置的认证 client")
|
t.Fatal("启用凭据准备必须有显式配置的认证 client")
|
||||||
}
|
}
|
||||||
if err := setupCredentialPreparation(nil, databaseOptions{}, nil); err != nil {
|
if err := registerDatabaseSupplyController(nil, databaseOptions{}, nil, nil); err != nil {
|
||||||
t.Fatal("默认停用凭据准备不应要求后端")
|
t.Fatal("默认停用凭据准备不应要求后端")
|
||||||
}
|
}
|
||||||
|
if err := registerDatabaseSupplyController(nil, databaseOptions{provisionResources: true}, nil, nil); err == nil {
|
||||||
|
t.Fatal("启用资源供应必须有管理连接及凭据准备")
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
func TestManagerFlagOverrides(t *testing.T) {
|
func TestManagerFlagOverrides(t *testing.T) {
|
||||||
|
|||||||
@@ -27,12 +27,12 @@ func Run() error {
|
|||||||
if err != nil {
|
if err != nil {
|
||||||
return fmt.Errorf("set up OpenBao authentication: %w", err)
|
return fmt.Errorf("set up OpenBao authentication: %w", err)
|
||||||
}
|
}
|
||||||
cleanup, err := setupDatabase(context.Background(), manager, options.database, baoClient)
|
closeDatabaseConnections, err := registerDatabaseControllers(context.Background(), manager, options.database, baoClient)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
// manager 的 worker 完全停止后才释放组件持有的资源。
|
// manager 的 worker 完全停止后才释放组件持有的资源。
|
||||||
defer cleanup()
|
defer closeDatabaseConnections()
|
||||||
|
|
||||||
setupLog.Info("Starting manager")
|
setupLog.Info("Starting manager")
|
||||||
if err := manager.Start(ctrl.SetupSignalHandler()); err != nil {
|
if err := manager.Start(ctrl.SetupSignalHandler()); err != nil {
|
||||||
|
|||||||
@@ -51,11 +51,3 @@ func tenantReference(tenant binding.TenantIdentity) *databasev1alpha1.TenantRefe
|
|||||||
Namespace: tenant.Namespace, Name: databasev1alpha1.ObjectName(tenant.Name), UID: types.UID(tenant.UID),
|
Namespace: tenant.Namespace, Name: databasev1alpha1.ObjectName(tenant.Name), UID: types.UID(tenant.UID),
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
// BindingTargetName 供 informer 索引使用;不把无效请求丢出事件映射。
|
|
||||||
func BindingTargetName(tenant *databasev1alpha1.PostgreSQLTenant) string {
|
|
||||||
if tenant.Spec.DatabaseRef != nil {
|
|
||||||
return string(tenant.Spec.DatabaseRef.Name)
|
|
||||||
}
|
|
||||||
return binding.DynamicDatabaseName(string(tenant.UID))
|
|
||||||
}
|
|
||||||
|
|||||||
@@ -0,0 +1,28 @@
|
|||||||
|
package kubernetes
|
||||||
|
|
||||||
|
import credentialdomain "git.ddupan.top/panxiao81/ayatori/internal/database/domain/credential"
|
||||||
|
|
||||||
|
const credentialsReadyCondition = "CredentialsReady"
|
||||||
|
|
||||||
|
// 持久化 API 的 Reason 保持兼容,领域内部只使用准备阶段。
|
||||||
|
var credentialReasons = map[credentialdomain.Phase]string{
|
||||||
|
credentialdomain.Pending: "PreparationPending",
|
||||||
|
credentialdomain.Pinned: "LocationPinned",
|
||||||
|
credentialdomain.Creating: "CreationStarted",
|
||||||
|
credentialdomain.Prepared: "CredentialPrepared",
|
||||||
|
credentialdomain.Conflict: "Conflict",
|
||||||
|
credentialdomain.Unavailable: "DependencyUnavailable",
|
||||||
|
credentialdomain.Stopped: "PreparationStopped",
|
||||||
|
credentialdomain.InvalidTarget: "InvalidTarget",
|
||||||
|
}
|
||||||
|
|
||||||
|
func credentialReason(phase credentialdomain.Phase) string { return credentialReasons[phase] }
|
||||||
|
|
||||||
|
func credentialPhase(reason string) credentialdomain.Phase {
|
||||||
|
for phase, value := range credentialReasons {
|
||||||
|
if value == reason {
|
||||||
|
return phase
|
||||||
|
}
|
||||||
|
}
|
||||||
|
return credentialdomain.Pending
|
||||||
|
}
|
||||||
@@ -4,6 +4,8 @@ import (
|
|||||||
"context"
|
"context"
|
||||||
"fmt"
|
"fmt"
|
||||||
|
|
||||||
|
credentialdomain "git.ddupan.top/panxiao81/ayatori/internal/database/domain/credential"
|
||||||
|
|
||||||
databasev1alpha1 "git.ddupan.top/panxiao81/ayatori/api/database/v1alpha1"
|
databasev1alpha1 "git.ddupan.top/panxiao81/ayatori/api/database/v1alpha1"
|
||||||
"git.ddupan.top/panxiao81/ayatori/internal/database/application"
|
"git.ddupan.top/panxiao81/ayatori/internal/database/application"
|
||||||
"git.ddupan.top/panxiao81/ayatori/internal/database/domain/binding"
|
"git.ddupan.top/panxiao81/ayatori/internal/database/domain/binding"
|
||||||
@@ -30,7 +32,8 @@ func (r *CredentialResources) Load(ctx context.Context, name string) (*applicati
|
|||||||
return nil, client.IgnoreNotFound(err)
|
return nil, client.IgnoreNotFound(err)
|
||||||
}
|
}
|
||||||
record := &application.CredentialRecord{
|
record := &application.CredentialRecord{
|
||||||
Database: *bindingDatabase(database), DatabaseProtected: controllerutil.ContainsFinalizer(database, DatabaseFinalizer),
|
Database: bindingDatabase(database).Database, Revision: database.ResourceVersion,
|
||||||
|
DatabaseProtected: controllerutil.ContainsFinalizer(database, DatabaseFinalizer),
|
||||||
Status: credentialStatus(database),
|
Status: credentialStatus(database),
|
||||||
}
|
}
|
||||||
if database.Spec.Source != "Provision" {
|
if database.Spec.Source != "Provision" {
|
||||||
@@ -43,7 +46,8 @@ func (r *CredentialResources) Load(ctx context.Context, name string) (*applicati
|
|||||||
return nil, err
|
return nil, err
|
||||||
}
|
}
|
||||||
if err == nil {
|
if err == nil {
|
||||||
record.Tenant = bindingTenant(tenant)
|
record.Tenant = &bindingTenant(tenant).Tenant
|
||||||
|
record.TenantGeneration = tenant.Generation
|
||||||
record.TenantProtected = controllerutil.ContainsFinalizer(tenant, TenantFinalizer)
|
record.TenantProtected = controllerutil.ContainsFinalizer(tenant, TenantFinalizer)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
@@ -56,29 +60,29 @@ func (r *CredentialResources) Load(ctx context.Context, name string) (*applicati
|
|||||||
if err != nil {
|
if err != nil {
|
||||||
return nil, err
|
return nil, err
|
||||||
}
|
}
|
||||||
record.Instance = &application.CredentialInstance{
|
record.Instance = &credentialdomain.Instance{
|
||||||
Identity: binding.Identity{Name: instance.Name, UID: string(instance.UID)},
|
Identity: binding.Identity{Name: instance.Name, UID: string(instance.UID)},
|
||||||
Deleting: !instance.DeletionTimestamp.IsZero(), Ready: currentReady(instance.Generation, instance.Status.Conditions),
|
Deleting: !instance.DeletionTimestamp.IsZero(), Ready: currentReady(instance.Generation, instance.Status.Conditions),
|
||||||
Generation: instance.Generation, Endpoint: observed.Target.Definition().Endpoint(),
|
Endpoint: observed.Target.Definition().Endpoint(),
|
||||||
}
|
}
|
||||||
|
record.InstanceGeneration = instance.Generation
|
||||||
return record, nil
|
return record, nil
|
||||||
}
|
}
|
||||||
|
|
||||||
func credentialStatus(database *databasev1alpha1.PostgreSQLDatabase) application.CredentialStatus {
|
func credentialStatus(database *databasev1alpha1.PostgreSQLDatabase) credentialdomain.State {
|
||||||
status := application.CredentialStatus{Version: database.Status.CredentialVersion}
|
status := credentialdomain.State{Version: database.Status.CredentialVersion}
|
||||||
if ref := database.Status.CredentialRef; ref != nil {
|
if ref := database.Status.CredentialRef; ref != nil {
|
||||||
status.Location = &application.CredentialLocation{Mount: ref.Mount, Path: ref.Path}
|
status.Location = &credentialdomain.Location{Mount: ref.Mount, Path: ref.Path}
|
||||||
}
|
}
|
||||||
if condition := meta.FindStatusCondition(database.Status.Conditions, application.CredentialsReady); condition != nil {
|
if condition := meta.FindStatusCondition(database.Status.Conditions, credentialsReadyCondition); condition != nil {
|
||||||
status.Ready = condition.Status == metav1.ConditionTrue && condition.ObservedGeneration == database.Generation
|
status.Phase, status.Message = credentialPhase(condition.Reason), condition.Message
|
||||||
status.Reason, status.Message = condition.Reason, condition.Message
|
|
||||||
}
|
}
|
||||||
return status
|
return status
|
||||||
}
|
}
|
||||||
|
|
||||||
func (r *CredentialResources) Save(ctx context.Context, record *application.CredentialRecord, status application.CredentialStatus) (*application.CredentialRecord, error) {
|
func (r *CredentialResources) Save(ctx context.Context, record *application.CredentialRecord, status credentialdomain.State) (*application.CredentialRecord, error) {
|
||||||
bindingResources := &BindingResources{Client: r.Client, Reader: r.Reader}
|
bindingResources := &BindingResources{Client: r.Client, Reader: r.Reader}
|
||||||
database, err := bindingResources.databaseAtVersion(ctx, &record.Database)
|
database, err := bindingResources.databaseAtVersion(ctx, &application.BindingDatabase{Database: record.Database, Revision: record.Revision})
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return nil, err
|
return nil, err
|
||||||
}
|
}
|
||||||
@@ -88,11 +92,11 @@ func (r *CredentialResources) Save(ctx context.Context, record *application.Cred
|
|||||||
}
|
}
|
||||||
database.Status.CredentialVersion = status.Version
|
database.Status.CredentialVersion = status.Version
|
||||||
conditionStatus := metav1.ConditionFalse
|
conditionStatus := metav1.ConditionFalse
|
||||||
if status.Ready {
|
if status.Phase == credentialdomain.Prepared {
|
||||||
conditionStatus = metav1.ConditionTrue
|
conditionStatus = metav1.ConditionTrue
|
||||||
}
|
}
|
||||||
meta.SetStatusCondition(&database.Status.Conditions, metav1.Condition{
|
meta.SetStatusCondition(&database.Status.Conditions, metav1.Condition{
|
||||||
Type: application.CredentialsReady, Status: conditionStatus, Reason: status.Reason, Message: status.Message,
|
Type: credentialsReadyCondition, Status: conditionStatus, Reason: credentialReason(status.Phase), Message: status.Message,
|
||||||
ObservedGeneration: database.Generation,
|
ObservedGeneration: database.Generation,
|
||||||
})
|
})
|
||||||
meta.SetStatusCondition(&database.Status.Conditions, metav1.Condition{
|
meta.SetStatusCondition(&database.Status.Conditions, metav1.Condition{
|
||||||
@@ -106,7 +110,8 @@ func (r *CredentialResources) Save(ctx context.Context, record *application.Cred
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
updated := *record
|
updated := *record
|
||||||
updated.Database = *bindingDatabase(database)
|
updated.Database = bindingDatabase(database).Database
|
||||||
|
updated.Revision = database.ResourceVersion
|
||||||
updated.Status = credentialStatus(database)
|
updated.Status = credentialStatus(database)
|
||||||
// 旧 CRD 会裁剪未知 status 字段。不能把 HTTP 成功当作位置/版本已保存后继续写后端。
|
// 旧 CRD 会裁剪未知 status 字段。不能把 HTTP 成功当作位置/版本已保存后继续写后端。
|
||||||
if !equality.Semantic.DeepEqual(updated.Status.Location, status.Location) || updated.Status.Version != status.Version {
|
if !equality.Semantic.DeepEqual(updated.Status.Location, status.Location) || updated.Status.Version != status.Version {
|
||||||
@@ -120,7 +125,7 @@ func (r *CredentialResources) CheckCurrent(ctx context.Context, record *applicat
|
|||||||
if err != nil {
|
if err != nil {
|
||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
if current == nil || current.Database.Identity != record.Database.Identity || current.Database.Revision != record.Database.Revision ||
|
if current == nil || current.Database.Identity != record.Database.Identity || current.Revision != record.Revision ||
|
||||||
!sameCredentialDependencies(current, record) {
|
!sameCredentialDependencies(current, record) {
|
||||||
return apierrors.NewConflict(databasev1alpha1.GroupVersion.WithResource("postgresqldatabases").GroupResource(), record.Database.Identity.Name,
|
return apierrors.NewConflict(databasev1alpha1.GroupVersion.WithResource("postgresqldatabases").GroupResource(), record.Database.Identity.Name,
|
||||||
fmt.Errorf("凭据准备的资源快照已变化;停止本轮操作并重新观察"))
|
fmt.Errorf("凭据准备的资源快照已变化;停止本轮操作并重新观察"))
|
||||||
@@ -134,9 +139,9 @@ func sameCredentialDependencies(current, previous *application.CredentialRecord)
|
|||||||
}
|
}
|
||||||
// Tenant Ready 的诊断变化、Instance 对同一 generation 的观测刷新不改变写入目标。
|
// Tenant Ready 的诊断变化、Instance 对同一 generation 的观测刷新不改变写入目标。
|
||||||
// 仍检查申请 spec generation、完整绑定身份、删除状态、保护和当前 Instance Ready。
|
// 仍检查申请 spec generation、完整绑定身份、删除状态、保护和当前 Instance Ready。
|
||||||
return current.Tenant.Generation == previous.Tenant.Generation &&
|
return current.TenantGeneration == previous.TenantGeneration &&
|
||||||
equality.Semantic.DeepEqual(current.Tenant.Tenant, previous.Tenant.Tenant) &&
|
equality.Semantic.DeepEqual(current.Tenant, previous.Tenant) &&
|
||||||
current.TenantProtected == previous.TenantProtected && current.DatabaseProtected == previous.DatabaseProtected &&
|
current.TenantProtected == previous.TenantProtected && current.DatabaseProtected == previous.DatabaseProtected &&
|
||||||
current.Instance.Instance == previous.Instance.Instance && current.Instance.Generation == previous.Instance.Generation &&
|
current.Instance.Instance == previous.Instance.Instance && current.InstanceGeneration == previous.InstanceGeneration &&
|
||||||
current.Instance.Endpoint == previous.Instance.Endpoint
|
current.Instance.Endpoint == previous.Instance.Endpoint
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -0,0 +1,111 @@
|
|||||||
|
package kubernetes
|
||||||
|
|
||||||
|
import (
|
||||||
|
"context"
|
||||||
|
"fmt"
|
||||||
|
"reflect"
|
||||||
|
|
||||||
|
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/provisioning"
|
||||||
|
"k8s.io/apimachinery/pkg/api/equality"
|
||||||
|
apierrors "k8s.io/apimachinery/pkg/api/errors"
|
||||||
|
"k8s.io/apimachinery/pkg/api/meta"
|
||||||
|
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
|
||||||
|
"k8s.io/apimachinery/pkg/types"
|
||||||
|
"sigs.k8s.io/controller-runtime/pkg/client"
|
||||||
|
)
|
||||||
|
|
||||||
|
const resourcesReadyCondition = "ResourcesReady"
|
||||||
|
|
||||||
|
type ProvisioningResources struct {
|
||||||
|
Client client.Client
|
||||||
|
Reader client.Reader
|
||||||
|
}
|
||||||
|
|
||||||
|
func (r *ProvisioningResources) Load(ctx context.Context, name string) (*application.ProvisioningRecord, error) {
|
||||||
|
// 同一组 API 事实复用映射;供应用例不持有凭据准备用例的快照或 repository。
|
||||||
|
facts, err := (&CredentialResources{Client: r.Client, Reader: r.Reader}).Load(ctx, name)
|
||||||
|
if err != nil || facts == nil {
|
||||||
|
return nil, err
|
||||||
|
}
|
||||||
|
object := &databasev1alpha1.PostgreSQLDatabase{}
|
||||||
|
if err := r.Reader.Get(ctx, types.NamespacedName{Name: name}, object); err != nil {
|
||||||
|
return nil, err
|
||||||
|
}
|
||||||
|
if object.ResourceVersion != facts.Revision {
|
||||||
|
return nil, provisioningSnapshotConflict(name)
|
||||||
|
}
|
||||||
|
record := &application.ProvisioningRecord{Target: facts.Target, Revision: facts.Revision, TenantGeneration: facts.TenantGeneration,
|
||||||
|
InstanceGeneration: facts.InstanceGeneration, Credentials: facts.Status, State: provisioningStatus(object)}
|
||||||
|
if facts.Instance == nil {
|
||||||
|
return record, nil
|
||||||
|
}
|
||||||
|
instance := &databasev1alpha1.PostgreSQLInstance{}
|
||||||
|
if err := r.Reader.Get(ctx, types.NamespacedName{Name: facts.Database.Instance}, instance); err != nil {
|
||||||
|
return nil, err
|
||||||
|
}
|
||||||
|
if string(instance.UID) != facts.Instance.Identity.UID || instance.Generation != facts.InstanceGeneration ||
|
||||||
|
(!instance.DeletionTimestamp.IsZero()) != facts.Instance.Deleting || currentReady(instance.Generation, instance.Status.Conditions) != facts.Instance.Ready {
|
||||||
|
return nil, provisioningSnapshotConflict(name)
|
||||||
|
}
|
||||||
|
mapped, err := instanceRecord(instance)
|
||||||
|
if err != nil {
|
||||||
|
return nil, err
|
||||||
|
}
|
||||||
|
record.InstanceTarget = mapped.Target
|
||||||
|
return record, nil
|
||||||
|
}
|
||||||
|
|
||||||
|
func provisioningStatus(object *databasev1alpha1.PostgreSQLDatabase) provisioning.State {
|
||||||
|
state := provisioning.State{RoleOID: uint32(object.Status.RoleOID), DatabaseOID: uint32(object.Status.DatabaseOID), Phase: provisioning.Pending}
|
||||||
|
if condition := meta.FindStatusCondition(object.Status.Conditions, resourcesReadyCondition); condition != nil {
|
||||||
|
state.Phase, state.Message = provisioning.Phase(condition.Reason), condition.Message
|
||||||
|
}
|
||||||
|
return state
|
||||||
|
}
|
||||||
|
|
||||||
|
func (r *ProvisioningResources) Save(ctx context.Context, record *application.ProvisioningRecord, state provisioning.State) (*application.ProvisioningRecord, error) {
|
||||||
|
object := &databasev1alpha1.PostgreSQLDatabase{}
|
||||||
|
if err := r.Reader.Get(ctx, types.NamespacedName{Name: record.Database.Identity.Name}, object); err != nil {
|
||||||
|
return nil, err
|
||||||
|
}
|
||||||
|
if object.ResourceVersion != record.Revision || string(object.UID) != record.Database.Identity.UID {
|
||||||
|
return nil, provisioningSnapshotConflict(object.Name)
|
||||||
|
}
|
||||||
|
previous := object.Status.DeepCopy()
|
||||||
|
object.Status.RoleOID, object.Status.DatabaseOID = int64(state.RoleOID), int64(state.DatabaseOID)
|
||||||
|
status := metav1.ConditionFalse
|
||||||
|
if state.Phase == provisioning.Available {
|
||||||
|
status = metav1.ConditionTrue
|
||||||
|
}
|
||||||
|
meta.SetStatusCondition(&object.Status.Conditions, metav1.Condition{Type: resourcesReadyCondition, Status: status,
|
||||||
|
Reason: string(state.Phase), Message: state.Message, ObservedGeneration: object.Generation})
|
||||||
|
if !equality.Semantic.DeepEqual(*previous, object.Status) {
|
||||||
|
if err := r.Client.Status().Update(ctx, object); err != nil {
|
||||||
|
return nil, err
|
||||||
|
}
|
||||||
|
}
|
||||||
|
if object.Status.RoleOID != int64(state.RoleOID) || object.Status.DatabaseOID != int64(state.DatabaseOID) {
|
||||||
|
return nil, fmt.Errorf("资源 OID 未被 API 保留;请先升级 Database CRD,未继续供应")
|
||||||
|
}
|
||||||
|
updated := *record
|
||||||
|
updated.Revision, updated.State = object.ResourceVersion, provisioningStatus(object)
|
||||||
|
return &updated, nil
|
||||||
|
}
|
||||||
|
|
||||||
|
func (r *ProvisioningResources) CheckCurrent(ctx context.Context, record *application.ProvisioningRecord) error {
|
||||||
|
current, err := r.Load(ctx, record.Database.Identity.Name)
|
||||||
|
if err != nil {
|
||||||
|
return err
|
||||||
|
}
|
||||||
|
if current == nil || !reflect.DeepEqual(current, record) {
|
||||||
|
return provisioningSnapshotConflict(record.Database.Identity.Name)
|
||||||
|
}
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
|
||||||
|
func provisioningSnapshotConflict(name string) error {
|
||||||
|
return apierrors.NewConflict(databasev1alpha1.GroupVersion.WithResource("postgresqldatabases").GroupResource(), name,
|
||||||
|
fmt.Errorf("供应快照发生变化;停止本轮并重新观察"))
|
||||||
|
}
|
||||||
@@ -3,22 +3,24 @@ package openbao
|
|||||||
import (
|
import (
|
||||||
"context"
|
"context"
|
||||||
|
|
||||||
|
credentialdomain "git.ddupan.top/panxiao81/ayatori/internal/database/domain/credential"
|
||||||
|
|
||||||
"git.ddupan.top/panxiao81/ayatori/internal/database/application"
|
"git.ddupan.top/panxiao81/ayatori/internal/database/application"
|
||||||
)
|
)
|
||||||
|
|
||||||
var _ application.CredentialStore = (*Credentials)(nil)
|
var _ application.CredentialStore = (*Credentials)(nil)
|
||||||
|
|
||||||
func (c *Credentials) ProvisionLocation(uid string) (application.CredentialLocation, error) {
|
func (c *Credentials) ProvisionLocation(uid string) (credentialdomain.Location, error) {
|
||||||
path, err := c.ProvisionPath(uid)
|
path, err := c.ProvisionPath(uid)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return application.CredentialLocation{}, err
|
return credentialdomain.Location{}, err
|
||||||
}
|
}
|
||||||
return application.CredentialLocation{Mount: c.mount, Path: path}, nil
|
return credentialdomain.Location{Mount: c.mount, Path: path}, nil
|
||||||
}
|
}
|
||||||
|
|
||||||
func (c *Credentials) ReadCredential(ctx context.Context, location application.CredentialLocation, version int64) (application.ApplicationCredential, error) {
|
func (c *Credentials) ReadCredential(ctx context.Context, location credentialdomain.Location, version int64) (credentialdomain.ApplicationCredential, error) {
|
||||||
if location.Mount != c.mount {
|
if location.Mount != c.mount {
|
||||||
return application.ApplicationCredential{}, ErrInvalidLocation
|
return credentialdomain.ApplicationCredential{}, ErrInvalidLocation
|
||||||
}
|
}
|
||||||
if version == 0 {
|
if version == 0 {
|
||||||
return c.Read(ctx, location.Path)
|
return c.Read(ctx, location.Path)
|
||||||
@@ -26,7 +28,7 @@ func (c *Credentials) ReadCredential(ctx context.Context, location application.C
|
|||||||
return c.ReadConfirmed(ctx, location.Path, version)
|
return c.ReadConfirmed(ctx, location.Path, version)
|
||||||
}
|
}
|
||||||
|
|
||||||
func (c *Credentials) CreateCredential(ctx context.Context, location application.CredentialLocation, credential application.ApplicationCredential) (int64, error) {
|
func (c *Credentials) CreateCredential(ctx context.Context, location credentialdomain.Location, credential credentialdomain.ApplicationCredential) (int64, error) {
|
||||||
if location.Mount != c.mount {
|
if location.Mount != c.mount {
|
||||||
return 0, ErrInvalidLocation
|
return 0, ErrInvalidLocation
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -26,6 +26,8 @@ import (
|
|||||||
"slices"
|
"slices"
|
||||||
"strings"
|
"strings"
|
||||||
|
|
||||||
|
credentialdomain "git.ddupan.top/panxiao81/ayatori/internal/database/domain/credential"
|
||||||
|
|
||||||
bao "github.com/openbao/openbao/api/v2"
|
bao "github.com/openbao/openbao/api/v2"
|
||||||
|
|
||||||
"git.ddupan.top/panxiao81/ayatori/internal/database/application"
|
"git.ddupan.top/panxiao81/ayatori/internal/database/application"
|
||||||
@@ -80,33 +82,33 @@ func (c *Credentials) accepts(path string) bool {
|
|||||||
}
|
}
|
||||||
|
|
||||||
// Read 只读取调用方已确认关联的路径;成功读取不构成对既有凭据的自动认领。
|
// Read 只读取调用方已确认关联的路径;成功读取不构成对既有凭据的自动认领。
|
||||||
func (c *Credentials) Read(ctx context.Context, path string) (application.ApplicationCredential, error) {
|
func (c *Credentials) Read(ctx context.Context, path string) (credentialdomain.ApplicationCredential, error) {
|
||||||
secret, err := c.read(ctx, path)
|
secret, err := c.read(ctx, path)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return application.ApplicationCredential{}, err
|
return credentialdomain.ApplicationCredential{}, err
|
||||||
}
|
}
|
||||||
return application.ParseApplicationCredential(secret.Data)
|
return credentialdomain.ParseApplicationCredential(secret.Data)
|
||||||
}
|
}
|
||||||
|
|
||||||
// ReadConfirmed 读取最新值并核对已持久化的确认版本,不回退读取历史版本。
|
// ReadConfirmed 读取最新值并核对已持久化的确认版本,不回退读取历史版本。
|
||||||
// 确认后的删除或改写需要人工处理,不能因此重新生成密码。
|
// 确认后的删除或改写需要人工处理,不能因此重新生成密码。
|
||||||
func (c *Credentials) ReadConfirmed(ctx context.Context, path string, version int64) (application.ApplicationCredential, error) {
|
func (c *Credentials) ReadConfirmed(ctx context.Context, path string, version int64) (credentialdomain.ApplicationCredential, error) {
|
||||||
if version < 1 {
|
if version < 1 {
|
||||||
return application.ApplicationCredential{}, ErrConflict
|
return credentialdomain.ApplicationCredential{}, ErrConflict
|
||||||
}
|
}
|
||||||
secret, err := c.read(ctx, path)
|
secret, err := c.read(ctx, path)
|
||||||
if errors.Is(err, ErrNotFound) {
|
if errors.Is(err, ErrNotFound) {
|
||||||
return application.ApplicationCredential{}, ErrConflict
|
return credentialdomain.ApplicationCredential{}, ErrConflict
|
||||||
}
|
}
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return application.ApplicationCredential{}, err
|
return credentialdomain.ApplicationCredential{}, err
|
||||||
}
|
}
|
||||||
if secret.VersionMetadata == nil || int64(secret.VersionMetadata.Version) != version {
|
if secret.VersionMetadata == nil || int64(secret.VersionMetadata.Version) != version {
|
||||||
return application.ApplicationCredential{}, ErrConflict
|
return credentialdomain.ApplicationCredential{}, ErrConflict
|
||||||
}
|
}
|
||||||
credential, err := application.ParseApplicationCredential(secret.Data)
|
credential, err := credentialdomain.ParseApplicationCredential(secret.Data)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return application.ApplicationCredential{}, ErrConflict
|
return credentialdomain.ApplicationCredential{}, ErrConflict
|
||||||
}
|
}
|
||||||
return credential, nil
|
return credential, nil
|
||||||
}
|
}
|
||||||
@@ -130,7 +132,7 @@ func (c *Credentials) read(ctx context.Context, path string) (*bao.KVSecret, err
|
|||||||
|
|
||||||
// Create 只创建从未存在过的路径,并验证回读七键与提交值完全一致。
|
// Create 只创建从未存在过的路径,并验证回读七键与提交值完全一致。
|
||||||
// 任何不确定写入都不返回凭据;上层必须停止供应并持久化冲突,不能重新生成密码。
|
// 任何不确定写入都不返回凭据;上层必须停止供应并持久化冲突,不能重新生成密码。
|
||||||
func (c *Credentials) Create(ctx context.Context, path string, credential application.ApplicationCredential) error {
|
func (c *Credentials) Create(ctx context.Context, path string, credential credentialdomain.ApplicationCredential) error {
|
||||||
if !c.accepts(path) {
|
if !c.accepts(path) {
|
||||||
return ErrInvalidLocation
|
return ErrInvalidLocation
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -25,10 +25,11 @@ import (
|
|||||||
"sync/atomic"
|
"sync/atomic"
|
||||||
"testing"
|
"testing"
|
||||||
|
|
||||||
|
credentialdomain "git.ddupan.top/panxiao81/ayatori/internal/database/domain/credential"
|
||||||
|
|
||||||
bao "github.com/openbao/openbao/api/v2"
|
bao "github.com/openbao/openbao/api/v2"
|
||||||
|
|
||||||
"git.ddupan.top/panxiao81/ayatori/internal/database/adapter/openbao"
|
"git.ddupan.top/panxiao81/ayatori/internal/database/adapter/openbao"
|
||||||
"git.ddupan.top/panxiao81/ayatori/internal/database/application"
|
|
||||||
)
|
)
|
||||||
|
|
||||||
const (
|
const (
|
||||||
@@ -39,9 +40,9 @@ const (
|
|||||||
kvVersionKey = "version"
|
kvVersionKey = "version"
|
||||||
)
|
)
|
||||||
|
|
||||||
func fixtureCredential(t *testing.T) application.ApplicationCredential {
|
func fixtureCredential(t *testing.T) credentialdomain.ApplicationCredential {
|
||||||
t.Helper()
|
t.Helper()
|
||||||
credential, err := application.ParseApplicationCredential(map[string]any{
|
credential, err := credentialdomain.ParseApplicationCredential(map[string]any{
|
||||||
"username": "app_owner", "password": fixturePassword, "database": "app",
|
"username": "app_owner", "password": fixturePassword, "database": "app",
|
||||||
"host": "postgres.example", "hostaddr": "192.0.2.1", "port": "5432", "sslmode": "verify-full",
|
"host": "postgres.example", "hostaddr": "192.0.2.1", "port": "5432", "sslmode": "verify-full",
|
||||||
})
|
})
|
||||||
|
|||||||
@@ -54,7 +54,7 @@ func testPreparationConcurrency(t *testing.T, f *preparationFixture) {
|
|||||||
if succeeded != 1 || conflicted != 1 {
|
if succeeded != 1 || conflicted != 1 {
|
||||||
t.Fatal("同一快照只能有一个用例成功固定位置并继续创建")
|
t.Fatal("同一快照只能有一个用例成功固定位置并继续创建")
|
||||||
}
|
}
|
||||||
f.status(t, database, 1, application.CredentialPrepared)
|
f.status(t, database, 1, "CredentialPrepared")
|
||||||
stored, err := f.bao.KVv2("secret").Get(t.Context(), database.Status.CredentialRef.Path)
|
stored, err := f.bao.KVv2("secret").Get(t.Context(), database.Status.CredentialRef.Path)
|
||||||
if err != nil || stored.VersionMetadata.Version != 1 {
|
if err != nil || stored.VersionMetadata.Version != 1 {
|
||||||
t.Fatal("并发准备用例只能产生一个凭据版本")
|
t.Fatal("并发准备用例只能产生一个凭据版本")
|
||||||
|
|||||||
@@ -80,7 +80,8 @@ func (f *preparationFixture) bound(t *testing.T, name string) (*databasev1alpha1
|
|||||||
if err := f.api.Create(t.Context(), tenant); err != nil {
|
if err := f.api.Create(t.Context(), tenant); err != nil {
|
||||||
t.Fatal(err)
|
t.Fatal(err)
|
||||||
}
|
}
|
||||||
binder := &databasecontroller.BindingReconciler{Client: f.api, Reader: f.api}
|
resources := &kubernetes.BindingResources{Client: f.api, Reader: f.api}
|
||||||
|
binder := databasecontroller.NewBindingReconciler(f.api, &application.BindingService{Resources: resources}, resources)
|
||||||
if _, err := binder.Reconcile(t.Context(), ctrl.Request{NamespacedName: client.ObjectKeyFromObject(tenant)}); err != nil {
|
if _, err := binder.Reconcile(t.Context(), ctrl.Request{NamespacedName: client.ObjectKeyFromObject(tenant)}); err != nil {
|
||||||
t.Fatal(err)
|
t.Fatal(err)
|
||||||
}
|
}
|
||||||
@@ -106,7 +107,7 @@ func (f *preparationFixture) status(t *testing.T, database *databasev1alpha1.Pos
|
|||||||
if err := f.api.Get(t.Context(), client.ObjectKeyFromObject(database), database); err != nil {
|
if err := f.api.Get(t.Context(), client.ObjectKeyFromObject(database), database); err != nil {
|
||||||
t.Fatal(err)
|
t.Fatal(err)
|
||||||
}
|
}
|
||||||
condition := meta.FindStatusCondition(database.Status.Conditions, application.CredentialsReady)
|
condition := meta.FindStatusCondition(database.Status.Conditions, "CredentialsReady")
|
||||||
if database.Status.CredentialVersion != version || condition == nil || condition.Reason != reason {
|
if database.Status.CredentialVersion != version || condition == nil || condition.Reason != reason {
|
||||||
t.Fatalf("凭据版本或条件不符:version=%d,期望 reason=%s", database.Status.CredentialVersion, reason)
|
t.Fatalf("凭据版本或条件不符:version=%d,期望 reason=%s", database.Status.CredentialVersion, reason)
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -8,6 +8,8 @@ import (
|
|||||||
"maps"
|
"maps"
|
||||||
"testing"
|
"testing"
|
||||||
|
|
||||||
|
credentialdomain "git.ddupan.top/panxiao81/ayatori/internal/database/domain/credential"
|
||||||
|
|
||||||
bao "github.com/openbao/openbao/api/v2"
|
bao "github.com/openbao/openbao/api/v2"
|
||||||
"sigs.k8s.io/controller-runtime/pkg/client"
|
"sigs.k8s.io/controller-runtime/pkg/client"
|
||||||
|
|
||||||
@@ -36,7 +38,7 @@ func testPreparationRestart(t *testing.T, f *preparationFixture) {
|
|||||||
if err := f.service(t).Reconcile(t.Context(), database.Name); err != nil {
|
if err := f.service(t).Reconcile(t.Context(), database.Name); err != nil {
|
||||||
t.Fatal(err)
|
t.Fatal(err)
|
||||||
}
|
}
|
||||||
f.status(t, database, 1, application.CredentialPrepared)
|
f.status(t, database, 1, "CredentialPrepared")
|
||||||
path := database.Status.CredentialRef.Path
|
path := database.Status.CredentialRef.Path
|
||||||
before, err := f.bao.KVv2("secret").Get(t.Context(), path)
|
before, err := f.bao.KVv2("secret").Get(t.Context(), path)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
@@ -51,12 +53,12 @@ func testPreparationRestart(t *testing.T, f *preparationFixture) {
|
|||||||
t.Fatal(err)
|
t.Fatal(err)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
f.status(t, database, 1, application.CredentialPrepared)
|
f.status(t, database, 1, "CredentialPrepared")
|
||||||
revision := database.ResourceVersion
|
revision := database.ResourceVersion
|
||||||
if err := f.service(t).Reconcile(t.Context(), database.Name); err != nil {
|
if err := f.service(t).Reconcile(t.Context(), database.Name); err != nil {
|
||||||
t.Fatal(err)
|
t.Fatal(err)
|
||||||
}
|
}
|
||||||
f.status(t, database, 1, application.CredentialPrepared)
|
f.status(t, database, 1, "CredentialPrepared")
|
||||||
if database.ResourceVersion != revision {
|
if database.ResourceVersion != revision {
|
||||||
t.Fatal("幂等重试不应改写 status")
|
t.Fatal("幂等重试不应改写 status")
|
||||||
}
|
}
|
||||||
@@ -135,7 +137,7 @@ func testPreparationLostConfirmation(t *testing.T, f *preparationFixture) {
|
|||||||
if err := service.Reconcile(t.Context(), database.Name); err == nil {
|
if err := service.Reconcile(t.Context(), database.Name); err == nil {
|
||||||
t.Fatal("确认写入失败应返回 API 错误")
|
t.Fatal("确认写入失败应返回 API 错误")
|
||||||
}
|
}
|
||||||
f.status(t, database, 0, application.CredentialCreationStarted)
|
f.status(t, database, 0, "CreationStarted")
|
||||||
if err := f.service(t).Reconcile(t.Context(), database.Name); err != nil {
|
if err := f.service(t).Reconcile(t.Context(), database.Name); err != nil {
|
||||||
t.Fatal(err)
|
t.Fatal(err)
|
||||||
}
|
}
|
||||||
@@ -151,7 +153,7 @@ type afterCreateStore struct {
|
|||||||
after func() error
|
after func() error
|
||||||
}
|
}
|
||||||
|
|
||||||
func (s afterCreateStore) CreateCredential(ctx context.Context, location application.CredentialLocation, credential application.ApplicationCredential) (int64, error) {
|
func (s afterCreateStore) CreateCredential(ctx context.Context, location credentialdomain.Location, credential credentialdomain.ApplicationCredential) (int64, error) {
|
||||||
version, err := s.CredentialStore.CreateCredential(ctx, location, credential)
|
version, err := s.CredentialStore.CreateCredential(ctx, location, credential)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return 0, err
|
return 0, err
|
||||||
@@ -185,7 +187,7 @@ func testPreparationChangedBinding(t *testing.T, f *preparationFixture) {
|
|||||||
if err := service.Reconcile(t.Context(), database.Name); err == nil {
|
if err := service.Reconcile(t.Context(), database.Name); err == nil {
|
||||||
t.Fatal("中途删除 Tenant 后不得确认凭据")
|
t.Fatal("中途删除 Tenant 后不得确认凭据")
|
||||||
}
|
}
|
||||||
f.status(t, database, 0, application.CredentialCreationStarted)
|
f.status(t, database, 0, "CreationStarted")
|
||||||
if err := f.service(t).Reconcile(t.Context(), database.Name); err != nil {
|
if err := f.service(t).Reconcile(t.Context(), database.Name); err != nil {
|
||||||
t.Fatal(err)
|
t.Fatal(err)
|
||||||
}
|
}
|
||||||
@@ -205,7 +207,7 @@ func testPreparationDependencies(t *testing.T, f *preparationFixture) {
|
|||||||
if err := f.service(t).Reconcile(t.Context(), database.Name); err != nil {
|
if err := f.service(t).Reconcile(t.Context(), database.Name); err != nil {
|
||||||
t.Fatal(err)
|
t.Fatal(err)
|
||||||
}
|
}
|
||||||
f.status(t, database, 1, application.CredentialPrepared)
|
f.status(t, database, 1, "CredentialPrepared")
|
||||||
if err := service.Reconcile(t.Context(), database.Name); err != nil {
|
if err := service.Reconcile(t.Context(), database.Name); err != nil {
|
||||||
t.Fatal(err)
|
t.Fatal(err)
|
||||||
}
|
}
|
||||||
@@ -213,7 +215,7 @@ func testPreparationDependencies(t *testing.T, f *preparationFixture) {
|
|||||||
if err := f.service(t).Reconcile(t.Context(), database.Name); err != nil {
|
if err := f.service(t).Reconcile(t.Context(), database.Name); err != nil {
|
||||||
t.Fatal(err)
|
t.Fatal(err)
|
||||||
}
|
}
|
||||||
f.status(t, database, 1, application.CredentialPrepared)
|
f.status(t, database, 1, "CredentialPrepared")
|
||||||
moved, err := openbao.NewCredentials(f.bao, "other", "elsewhere")
|
moved, err := openbao.NewCredentials(f.bao, "other", "elsewhere")
|
||||||
if err != nil {
|
if err != nil {
|
||||||
t.Fatal(err)
|
t.Fatal(err)
|
||||||
|
|||||||
@@ -12,7 +12,6 @@ import (
|
|||||||
|
|
||||||
databasev1alpha1 "git.ddupan.top/panxiao81/ayatori/api/database/v1alpha1"
|
databasev1alpha1 "git.ddupan.top/panxiao81/ayatori/api/database/v1alpha1"
|
||||||
"git.ddupan.top/panxiao81/ayatori/internal/database/adapter/kubernetes"
|
"git.ddupan.top/panxiao81/ayatori/internal/database/adapter/kubernetes"
|
||||||
"git.ddupan.top/panxiao81/ayatori/internal/database/application"
|
|
||||||
)
|
)
|
||||||
|
|
||||||
// 在 API 边界模拟旧 schema 裁剪位置字段;剩余写入仍由真实 API server 处理。
|
// 在 API 边界模拟旧 schema 裁剪位置字段;剩余写入仍由真实 API server 处理。
|
||||||
@@ -48,5 +47,5 @@ func testPreparationPruning(t *testing.T, f *preparationFixture) {
|
|||||||
if err := f.service(t).Reconcile(t.Context(), database.Name); err != nil {
|
if err := f.service(t).Reconcile(t.Context(), database.Name); err != nil {
|
||||||
t.Fatal(err)
|
t.Fatal(err)
|
||||||
}
|
}
|
||||||
f.status(t, database, 1, application.CredentialPrepared)
|
f.status(t, database, 1, "CredentialPrepared")
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -7,6 +7,9 @@ import (
|
|||||||
"testing"
|
"testing"
|
||||||
"time"
|
"time"
|
||||||
|
|
||||||
|
"git.ddupan.top/panxiao81/ayatori/internal/database/adapter/kubernetes"
|
||||||
|
"git.ddupan.top/panxiao81/ayatori/internal/database/application"
|
||||||
|
|
||||||
"k8s.io/apimachinery/pkg/api/meta"
|
"k8s.io/apimachinery/pkg/api/meta"
|
||||||
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
|
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
|
||||||
ctrl "sigs.k8s.io/controller-runtime"
|
ctrl "sigs.k8s.io/controller-runtime"
|
||||||
@@ -15,7 +18,6 @@ import (
|
|||||||
metricsserver "sigs.k8s.io/controller-runtime/pkg/metrics/server"
|
metricsserver "sigs.k8s.io/controller-runtime/pkg/metrics/server"
|
||||||
|
|
||||||
databasev1alpha1 "git.ddupan.top/panxiao81/ayatori/api/database/v1alpha1"
|
databasev1alpha1 "git.ddupan.top/panxiao81/ayatori/api/database/v1alpha1"
|
||||||
"git.ddupan.top/panxiao81/ayatori/internal/database/application"
|
|
||||||
databasecontroller "git.ddupan.top/panxiao81/ayatori/internal/database/controller"
|
databasecontroller "git.ddupan.top/panxiao81/ayatori/internal/database/controller"
|
||||||
)
|
)
|
||||||
|
|
||||||
@@ -26,7 +28,7 @@ func testPreparationWatch(t *testing.T, f *preparationFixture) {
|
|||||||
return
|
return
|
||||||
}
|
}
|
||||||
if err := f.api.Get(context.Background(), client.ObjectKeyFromObject(database), database); err == nil {
|
if err := f.api.Get(context.Background(), client.ObjectKeyFromObject(database), database); err == nil {
|
||||||
if condition := meta.FindStatusCondition(database.Status.Conditions, application.CredentialsReady); condition != nil {
|
if condition := meta.FindStatusCondition(database.Status.Conditions, "CredentialsReady"); condition != nil {
|
||||||
t.Logf("失败时凭据条件: %s: %s", condition.Reason, condition.Message)
|
t.Logf("失败时凭据条件: %s: %s", condition.Reason, condition.Message)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
@@ -50,7 +52,7 @@ func testPreparationWatch(t *testing.T, f *preparationFixture) {
|
|||||||
if err := f.api.Get(t.Context(), client.ObjectKeyFromObject(database), database); err != nil {
|
if err := f.api.Get(t.Context(), client.ObjectKeyFromObject(database), database); err != nil {
|
||||||
return false
|
return false
|
||||||
}
|
}
|
||||||
condition := meta.FindStatusCondition(database.Status.Conditions, application.CredentialsReady)
|
condition := meta.FindStatusCondition(database.Status.Conditions, "CredentialsReady")
|
||||||
return condition != nil && condition.Reason == "DependencyUnavailable"
|
return condition != nil && condition.Reason == "DependencyUnavailable"
|
||||||
})
|
})
|
||||||
if database.Status.CredentialRef != nil {
|
if database.Status.CredentialRef != nil {
|
||||||
@@ -69,7 +71,7 @@ func testPreparationWatch(t *testing.T, f *preparationFixture) {
|
|||||||
if err := f.api.Get(t.Context(), client.ObjectKeyFromObject(database), database); err != nil {
|
if err := f.api.Get(t.Context(), client.ObjectKeyFromObject(database), database); err != nil {
|
||||||
return false
|
return false
|
||||||
}
|
}
|
||||||
condition := meta.FindStatusCondition(database.Status.Conditions, application.CredentialsReady)
|
condition := meta.FindStatusCondition(database.Status.Conditions, "CredentialsReady")
|
||||||
return condition != nil && condition.Status == metav1.ConditionFalse && condition.Reason == "DependencyUnavailable"
|
return condition != nil && condition.Status == metav1.ConditionFalse && condition.Reason == "DependencyUnavailable"
|
||||||
})
|
})
|
||||||
setReady(metav1.ConditionTrue)
|
setReady(metav1.ConditionTrue)
|
||||||
@@ -77,7 +79,7 @@ func testPreparationWatch(t *testing.T, f *preparationFixture) {
|
|||||||
if err := f.api.Get(t.Context(), client.ObjectKeyFromObject(database), database); err != nil {
|
if err := f.api.Get(t.Context(), client.ObjectKeyFromObject(database), database); err != nil {
|
||||||
return false
|
return false
|
||||||
}
|
}
|
||||||
condition := meta.FindStatusCondition(database.Status.Conditions, application.CredentialsReady)
|
condition := meta.FindStatusCondition(database.Status.Conditions, "CredentialsReady")
|
||||||
return condition != nil && condition.Status == metav1.ConditionTrue && database.Status.CredentialVersion == 1
|
return condition != nil && condition.Status == metav1.ConditionTrue && database.Status.CredentialVersion == 1
|
||||||
})
|
})
|
||||||
stored, err := f.bao.KVv2("secret").Get(t.Context(), database.Status.CredentialRef.Path)
|
stored, err := f.bao.KVv2("secret").Get(t.Context(), database.Status.CredentialRef.Path)
|
||||||
@@ -96,10 +98,14 @@ func startPreparationManager(t *testing.T, f *preparationFixture) func() {
|
|||||||
if err != nil {
|
if err != nil {
|
||||||
t.Fatal(err)
|
t.Fatal(err)
|
||||||
}
|
}
|
||||||
if err := (&databasecontroller.BindingReconciler{}).SetupWithManager(t.Context(), manager); err != nil {
|
bindingResources := &kubernetes.BindingResources{Client: manager.GetClient(), Reader: manager.GetAPIReader()}
|
||||||
|
binder := databasecontroller.NewBindingReconciler(manager.GetClient(), &application.BindingService{Resources: bindingResources}, bindingResources)
|
||||||
|
if err := binder.SetupWithManager(t.Context(), manager); err != nil {
|
||||||
t.Fatal(err)
|
t.Fatal(err)
|
||||||
}
|
}
|
||||||
if err := (&databasecontroller.CredentialReconciler{Store: fixtureStore(t, f.bao)}).SetupWithManager(manager); err != nil {
|
credentialResources := &kubernetes.CredentialResources{Client: manager.GetClient(), Reader: manager.GetAPIReader()}
|
||||||
|
preparation := &application.CredentialPreparation{Resources: credentialResources, Store: fixtureStore(t, f.bao)}
|
||||||
|
if err := databasecontroller.NewCredentialReconciler(manager.GetClient(), preparation).SetupWithManager(manager); err != nil {
|
||||||
t.Fatal(err)
|
t.Fatal(err)
|
||||||
}
|
}
|
||||||
ctx, cancel := context.WithCancel(t.Context())
|
ctx, cancel := context.WithCancel(t.Context())
|
||||||
|
|||||||
@@ -68,7 +68,9 @@ func TestInstanceControllerWithRealPostgreSQL(t *testing.T) {
|
|||||||
if err != nil {
|
if err != nil {
|
||||||
t.Fatal(err)
|
t.Fatal(err)
|
||||||
}
|
}
|
||||||
reconciler := &databasecontroller.InstanceReconciler{Observer: service, SecretNamespace: controllerNamespace}
|
resources := &secretadapter.InstanceResources{Client: manager.GetClient(), Reader: manager.GetAPIReader()}
|
||||||
|
usecase := &application.InstanceReconciliation{Resources: resources, Observer: service}
|
||||||
|
reconciler := databasecontroller.NewInstanceReconciler(manager.GetClient(), usecase, resources, controllerNamespace)
|
||||||
if err := reconciler.SetupWithManager(manager); err != nil {
|
if err := reconciler.SetupWithManager(manager); err != nil {
|
||||||
t.Fatal(err)
|
t.Fatal(err)
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -0,0 +1,184 @@
|
|||||||
|
//go:build integration
|
||||||
|
|
||||||
|
package postgresql_test
|
||||||
|
|
||||||
|
import (
|
||||||
|
"context"
|
||||||
|
"os/exec"
|
||||||
|
"regexp"
|
||||||
|
"strings"
|
||||||
|
"testing"
|
||||||
|
"time"
|
||||||
|
|
||||||
|
databasev1alpha1 "git.ddupan.top/panxiao81/ayatori/api/database/v1alpha1"
|
||||||
|
"git.ddupan.top/panxiao81/ayatori/internal/database/adapter/kubernetes"
|
||||||
|
"git.ddupan.top/panxiao81/ayatori/internal/database/adapter/openbao"
|
||||||
|
"git.ddupan.top/panxiao81/ayatori/internal/database/application"
|
||||||
|
databasecontroller "git.ddupan.top/panxiao81/ayatori/internal/database/controller"
|
||||||
|
"git.ddupan.top/panxiao81/ayatori/internal/database/domain/provisioning"
|
||||||
|
bao "github.com/openbao/openbao/api/v2"
|
||||||
|
"k8s.io/apimachinery/pkg/api/meta"
|
||||||
|
"k8s.io/apimachinery/pkg/runtime"
|
||||||
|
ctrl "sigs.k8s.io/controller-runtime"
|
||||||
|
"sigs.k8s.io/controller-runtime/pkg/client"
|
||||||
|
"sigs.k8s.io/controller-runtime/pkg/envtest"
|
||||||
|
)
|
||||||
|
|
||||||
|
type provisioningFixture struct {
|
||||||
|
*credentialFixture
|
||||||
|
api client.Client
|
||||||
|
scheme *runtime.Scheme
|
||||||
|
store *openbao.Credentials
|
||||||
|
resources *kubernetes.ProvisioningResources
|
||||||
|
usecase application.DatabaseProvisioning
|
||||||
|
}
|
||||||
|
|
||||||
|
func newProvisioningFixture(t *testing.T) *provisioningFixture {
|
||||||
|
t.Helper()
|
||||||
|
f := newCredentialFixture(t)
|
||||||
|
useNativeManager(t, f)
|
||||||
|
if _, err := envtest.InstallCRDs(f.config, envtest.CRDInstallOptions{Paths: []string{"../../../../config/crd/bases"}, ErrorIfPathMissing: true}); err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
scheme := runtime.NewScheme()
|
||||||
|
if err := databasev1alpha1.AddToScheme(scheme); err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
api, err := client.New(f.config, client.Options{Scheme: scheme})
|
||||||
|
if err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
object := &databasev1alpha1.PostgreSQLInstance{Name: "supply-instance", Spec: databasev1alpha1.PostgreSQLInstanceSpec{
|
||||||
|
Endpoint: databasev1alpha1.PostgreSQLEndpoint{Host: fixtureHost, HostAddr: fixtureAddress, Port: int32(f.port), SSLMode: "disable"},
|
||||||
|
AdminCredentialRef: databasev1alpha1.AdminCredentialReference{Name: secretName, UsernameKey: managementUsernameKey, PasswordKey: managementPasswordKey},
|
||||||
|
}}
|
||||||
|
if err := api.Create(f.ctx, object); err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
instanceResources := &kubernetes.InstanceResources{Client: api, Reader: api}
|
||||||
|
instanceService := &application.InstanceReconciliation{Resources: instanceResources, Observer: f.service}
|
||||||
|
result, err := instanceService.Reconcile(f.ctx, object.Name)
|
||||||
|
if err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
if err := instanceResources.PresentInstance(f.ctx, result); err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
if err := api.Get(f.ctx, client.ObjectKeyFromObject(object), object); err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
if !meta.IsStatusConditionTrue(object.Status.Conditions, "Ready") {
|
||||||
|
t.Fatal("实际管理账号未就绪")
|
||||||
|
}
|
||||||
|
store := provisioningBao(t)
|
||||||
|
resources := &kubernetes.ProvisioningResources{Client: api, Reader: api}
|
||||||
|
return &provisioningFixture{credentialFixture: f, api: api, scheme: scheme, store: store, resources: resources,
|
||||||
|
usecase: application.DatabaseProvisioning{Resources: resources, Credentials: store, Backend: f.service}}
|
||||||
|
}
|
||||||
|
|
||||||
|
// 三后端组合验收只启动自己的 dev Bao;不读取环境 token 或外部地址。
|
||||||
|
func provisioningBao(t *testing.T) *openbao.Credentials {
|
||||||
|
t.Helper()
|
||||||
|
const image = "openbao/openbao@sha256:5b2486ab0fb90bbc788cc345b0a08616dfb375873ee8be5df3a2fd4d378a67e0"
|
||||||
|
const token = "AYATORI-TEST-ONLY-supply-token"
|
||||||
|
prepare, cancel := context.WithTimeout(t.Context(), 5*time.Minute)
|
||||||
|
defer cancel()
|
||||||
|
if exec.CommandContext(prepare, "docker", "image", "inspect", image).Run() != nil {
|
||||||
|
if output, err := exec.CommandContext(prepare, "docker", "pull", image).CombinedOutput(); err != nil {
|
||||||
|
t.Fatalf("隔离 Bao 镜像准备失败:%s", output)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
ctx, stop := context.WithTimeout(t.Context(), time.Minute)
|
||||||
|
defer stop()
|
||||||
|
output, err := exec.CommandContext(ctx, "docker", "run", "--pull=never", "--rm", "-d", "-p", "127.0.0.1::8200", image,
|
||||||
|
"server", "-dev", "-dev-root-token-id="+token, "-dev-listen-address=0.0.0.0:8200").Output()
|
||||||
|
if err != nil {
|
||||||
|
t.Fatal("隔离 Bao 启动失败")
|
||||||
|
}
|
||||||
|
id := strings.TrimSpace(string(output))
|
||||||
|
if !regexp.MustCompile(`^[a-f0-9]{64}$`).MatchString(id) {
|
||||||
|
t.Fatal("无效容器 ID")
|
||||||
|
}
|
||||||
|
t.Cleanup(func() {
|
||||||
|
cleanup, cancel := context.WithTimeout(context.Background(), 30*time.Second)
|
||||||
|
defer cancel()
|
||||||
|
if exec.CommandContext(cleanup, "docker", "rm", "-f", id).Run() != nil {
|
||||||
|
t.Error("隔离 Bao 清理失败")
|
||||||
|
}
|
||||||
|
})
|
||||||
|
output, err = exec.CommandContext(ctx, "docker", "inspect", "--format", `{{(index (index .NetworkSettings.Ports "8200/tcp") 0).HostPort}}`, id).Output()
|
||||||
|
if err != nil {
|
||||||
|
t.Fatal("隔离 Bao 端口不可读")
|
||||||
|
}
|
||||||
|
backend, err := bao.NewClient(&bao.Config{Address: "http://127.0.0.1:" + strings.TrimSpace(string(output)), Timeout: 5 * time.Second})
|
||||||
|
if err != nil {
|
||||||
|
t.Fatal("隔离 Bao client 构造失败")
|
||||||
|
}
|
||||||
|
backend.SetToken(token)
|
||||||
|
for {
|
||||||
|
if _, err := backend.Sys().HealthWithContext(ctx); err == nil {
|
||||||
|
break
|
||||||
|
}
|
||||||
|
select {
|
||||||
|
case <-ctx.Done():
|
||||||
|
t.Fatal("隔离 Bao 健康检查超时")
|
||||||
|
case <-time.After(100 * time.Millisecond):
|
||||||
|
}
|
||||||
|
}
|
||||||
|
store, err := openbao.NewCredentials(backend, "secret", "applications")
|
||||||
|
if err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
return store
|
||||||
|
}
|
||||||
|
|
||||||
|
func (f *provisioningFixture) bound(t *testing.T, name string) *databasev1alpha1.PostgreSQLDatabase {
|
||||||
|
t.Helper()
|
||||||
|
tenant := &databasev1alpha1.PostgreSQLTenant{Name: strings.ReplaceAll(name, "_", "-"), Namespace: controllerNamespace,
|
||||||
|
Spec: databasev1alpha1.PostgreSQLTenantSpec{Provision: &databasev1alpha1.DatabaseProvisionRequest{
|
||||||
|
InstanceRef: databasev1alpha1.InstanceReference{Name: "supply-instance"},
|
||||||
|
Database: databasev1alpha1.PostgreSQLIdentifier(name), LoginRole: databasev1alpha1.PostgreSQLIdentifier(name),
|
||||||
|
}}}
|
||||||
|
if err := f.api.Create(f.ctx, tenant); err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
resources := &kubernetes.BindingResources{Client: f.api, Reader: f.api}
|
||||||
|
binder := databasecontroller.NewBindingReconciler(f.api, &application.BindingService{Resources: resources}, resources)
|
||||||
|
if _, err := binder.Reconcile(f.ctx, ctrl.Request{NamespacedName: client.ObjectKeyFromObject(tenant)}); err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
if err := f.api.Get(f.ctx, client.ObjectKeyFromObject(tenant), tenant); err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
database := &databasev1alpha1.PostgreSQLDatabase{}
|
||||||
|
if err := f.api.Get(f.ctx, client.ObjectKey{Name: string(tenant.Status.DatabaseRef.Name)}, database); err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
prepare := application.CredentialPreparation{Resources: &kubernetes.CredentialResources{Client: f.api, Reader: f.api}, Store: f.store}
|
||||||
|
if err := prepare.Reconcile(f.ctx, database.Name); err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
return database
|
||||||
|
}
|
||||||
|
|
||||||
|
func (f *provisioningFixture) reconcile(t *testing.T, database *databasev1alpha1.PostgreSQLDatabase, phase provisioning.Phase) {
|
||||||
|
t.Helper()
|
||||||
|
if err := f.usecase.Reconcile(f.ctx, database.Name); err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
f.status(t, database, phase)
|
||||||
|
}
|
||||||
|
|
||||||
|
func (f *provisioningFixture) status(t *testing.T, database *databasev1alpha1.PostgreSQLDatabase, phase provisioning.Phase) {
|
||||||
|
t.Helper()
|
||||||
|
if err := f.api.Get(f.ctx, client.ObjectKeyFromObject(database), database); err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
condition := meta.FindStatusCondition(database.Status.Conditions, "ResourcesReady")
|
||||||
|
if condition == nil || condition.Reason != string(phase) {
|
||||||
|
t.Fatalf("期望资源阶段 %s,实际 %v", phase, condition)
|
||||||
|
}
|
||||||
|
if meta.IsStatusConditionTrue(database.Status.Conditions, "Ready") {
|
||||||
|
t.Fatal("角色建库不代表扩展与交付完成")
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -0,0 +1,243 @@
|
|||||||
|
//go:build integration
|
||||||
|
|
||||||
|
package postgresql_test
|
||||||
|
|
||||||
|
import (
|
||||||
|
"context"
|
||||||
|
"errors"
|
||||||
|
"sync"
|
||||||
|
"sync/atomic"
|
||||||
|
"testing"
|
||||||
|
|
||||||
|
"git.ddupan.top/panxiao81/ayatori/internal/database/application"
|
||||||
|
"git.ddupan.top/panxiao81/ayatori/internal/database/domain/credential"
|
||||||
|
"git.ddupan.top/panxiao81/ayatori/internal/database/domain/instance"
|
||||||
|
"git.ddupan.top/panxiao81/ayatori/internal/database/domain/provisioning"
|
||||||
|
"github.com/jackc/pgx/v5"
|
||||||
|
apierrors "k8s.io/apimachinery/pkg/api/errors"
|
||||||
|
)
|
||||||
|
|
||||||
|
func testProvisioningLogin(t *testing.T, f *provisioningFixture) {
|
||||||
|
t.Helper()
|
||||||
|
database := f.bound(t, "supplied")
|
||||||
|
f.reconcile(t, database, provisioning.Pending)
|
||||||
|
if database.Status.RoleOID == 0 || database.Status.DatabaseOID != 0 {
|
||||||
|
t.Fatal("角色确认未独立保存")
|
||||||
|
}
|
||||||
|
f.service.Forget("supply-instance")
|
||||||
|
f.reconcile(t, database, provisioning.Pending)
|
||||||
|
if database.Status.DatabaseOID == 0 {
|
||||||
|
t.Fatal("数据库未独立确认")
|
||||||
|
}
|
||||||
|
if got := f.queryPostgres(t, "SELECT datallowconn FROM pg_database WHERE datname='supplied'"); got != "f" {
|
||||||
|
t.Fatal("ACL 收紧前连接入口应关闭")
|
||||||
|
}
|
||||||
|
f.reconcile(t, database, provisioning.Available)
|
||||||
|
beforeRole, beforeDB := database.Status.RoleOID, database.Status.DatabaseOID
|
||||||
|
for range 2 {
|
||||||
|
f.reconcile(t, database, provisioning.Available)
|
||||||
|
}
|
||||||
|
if beforeRole != database.Status.RoleOID || beforeDB != database.Status.DatabaseOID {
|
||||||
|
t.Fatal("幂等协调替换了对象")
|
||||||
|
}
|
||||||
|
value, err := f.store.ReadCredential(f.ctx, credential.Location{Mount: database.Status.CredentialRef.Mount, Path: database.Status.CredentialRef.Path}, database.Status.CredentialVersion)
|
||||||
|
if err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
config, err := pgx.ParseConfig("")
|
||||||
|
if err != nil {
|
||||||
|
t.Fatal("fixture 配置失败")
|
||||||
|
}
|
||||||
|
config.Host, config.Port = fixtureAddress, uint16(f.port)
|
||||||
|
config.User, config.Database, config.Password = "supplied", "supplied", value.SecretData()["password"].(string)
|
||||||
|
config.TLSConfig, config.Fallbacks = nil, nil
|
||||||
|
connection, err := pgx.ConnectConfig(f.ctx, config)
|
||||||
|
if err != nil {
|
||||||
|
t.Fatal("实际应用密码不能登录已供应数据库")
|
||||||
|
}
|
||||||
|
defer func() { _ = connection.Close(context.Background()) }()
|
||||||
|
if _, err := connection.Exec(f.ctx, "CREATE TABLE app_data (id integer)"); err != nil {
|
||||||
|
t.Fatal("应用 owner 不能创建表")
|
||||||
|
}
|
||||||
|
if _, err := connection.Exec(f.ctx, "CREATE ROLE should_be_denied"); err == nil {
|
||||||
|
t.Fatal("应用具有 CREATEROLE")
|
||||||
|
}
|
||||||
|
if _, err := connection.Exec(f.ctx, "CREATE DATABASE should_be_denied"); err == nil {
|
||||||
|
t.Fatal("应用具有 CREATEDB")
|
||||||
|
}
|
||||||
|
f.queryPostgres(t, "CREATE ROLE unrelated_login LOGIN PASSWORD '"+fixturePassword+"'")
|
||||||
|
config.User, config.Password = "unrelated_login", fixturePassword
|
||||||
|
if outsider, err := pgx.ConnectConfig(f.ctx, config); err == nil {
|
||||||
|
_ = outsider.Close(f.ctx)
|
||||||
|
t.Fatal("其他账号可连接受管数据库")
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
func TestDatabaseProvisioningWithRealBackends(t *testing.T) {
|
||||||
|
f := newProvisioningFixture(t)
|
||||||
|
t.Run("创建重启幂等与实际登录", func(t *testing.T) { testProvisioningLogin(t, f) })
|
||||||
|
t.Run("未知同名不认领", func(t *testing.T) {
|
||||||
|
database := f.bound(t, "occupied")
|
||||||
|
f.queryPostgres(t, "CREATE ROLE occupied LOGIN")
|
||||||
|
before := f.queryPostgres(t, "SELECT oid FROM pg_roles WHERE rolname='occupied'")
|
||||||
|
f.reconcile(t, database, provisioning.Conflict)
|
||||||
|
f.reconcile(t, database, provisioning.Conflict)
|
||||||
|
if database.Status.RoleOID != 0 || f.queryPostgres(t, "SELECT oid FROM pg_roles WHERE rolname='occupied'") != before {
|
||||||
|
t.Fatal("未知角色被认领或修改")
|
||||||
|
}
|
||||||
|
})
|
||||||
|
t.Run("未知同名数据库不认领", func(t *testing.T) {
|
||||||
|
database := f.bound(t, "occupied_database")
|
||||||
|
f.queryPostgres(t, "CREATE DATABASE occupied_database")
|
||||||
|
before := f.queryPostgres(t, "SELECT oid FROM pg_database WHERE datname='occupied_database'")
|
||||||
|
f.reconcile(t, database, provisioning.Conflict)
|
||||||
|
if database.Status.RoleOID != 0 || database.Status.DatabaseOID != 0 ||
|
||||||
|
f.queryPostgres(t, "SELECT oid FROM pg_database WHERE datname='occupied_database'") != before ||
|
||||||
|
f.queryPostgres(t, "SELECT count(*) FROM pg_roles WHERE rolname='occupied_database'") != "0" {
|
||||||
|
t.Fatal("未知数据库被认领、修改或继续创建了角色")
|
||||||
|
}
|
||||||
|
})
|
||||||
|
t.Run("已确认角色后权限恢复", func(t *testing.T) {
|
||||||
|
database := f.bound(t, "permission_restore")
|
||||||
|
f.reconcile(t, database, provisioning.Pending)
|
||||||
|
f.queryPostgres(t, "ALTER ROLE native_manager NOCREATEDB")
|
||||||
|
f.reconcile(t, database, provisioning.Unavailable)
|
||||||
|
f.queryPostgres(t, "ALTER ROLE native_manager CREATEDB")
|
||||||
|
f.reconcile(t, database, provisioning.Pending)
|
||||||
|
f.reconcile(t, database, provisioning.Available)
|
||||||
|
})
|
||||||
|
t.Run("删除前置阻止写入", func(t *testing.T) {
|
||||||
|
database := f.bound(t, "deleting_supply")
|
||||||
|
if err := f.api.Delete(f.ctx, database); err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
f.reconcile(t, database, provisioning.Stopped)
|
||||||
|
if len(database.Finalizers) == 0 || database.Status.RoleOID != 0 {
|
||||||
|
t.Fatal("删除边界被供应绕过")
|
||||||
|
}
|
||||||
|
})
|
||||||
|
t.Run("确认保存失败转人工冲突", func(t *testing.T) {
|
||||||
|
database := f.bound(t, "lost_confirmation")
|
||||||
|
failed := f.usecase
|
||||||
|
failed.Resources = &failRoleConfirmation{ProvisioningResources: f.resources}
|
||||||
|
if err := failed.Reconcile(f.ctx, database.Name); err == nil {
|
||||||
|
t.Fatal("未注入确认写入故障")
|
||||||
|
}
|
||||||
|
f.reconcile(t, database, provisioning.Conflict)
|
||||||
|
if database.Status.RoleOID != 0 || f.queryPostgres(t, "SELECT count(*) FROM pg_roles WHERE rolname='lost_confirmation'") != "1" {
|
||||||
|
t.Fatal("失败恢复认领或清理了残留角色")
|
||||||
|
}
|
||||||
|
})
|
||||||
|
t.Run("数据库确认失败保留关闭入口", func(t *testing.T) {
|
||||||
|
database := f.bound(t, "lost_database_confirmation")
|
||||||
|
f.reconcile(t, database, provisioning.Pending)
|
||||||
|
failed := f.usecase
|
||||||
|
failed.Resources = &failDatabaseConfirmation{ProvisioningResources: f.resources}
|
||||||
|
if err := failed.Reconcile(f.ctx, database.Name); err == nil {
|
||||||
|
t.Fatal("未注入数据库确认写入故障")
|
||||||
|
}
|
||||||
|
f.reconcile(t, database, provisioning.Conflict)
|
||||||
|
if database.Status.RoleOID == 0 || database.Status.DatabaseOID != 0 || f.queryPostgres(t, "SELECT datallowconn FROM pg_database WHERE datname='lost_database_confirmation'") != "f" {
|
||||||
|
t.Fatal("未确认数据库被认领或开放连接")
|
||||||
|
}
|
||||||
|
})
|
||||||
|
t.Run("创建成功但调用方丢失结果", func(t *testing.T) {
|
||||||
|
database := f.bound(t, "lost_role_response")
|
||||||
|
failed := f.usecase
|
||||||
|
failed.Backend = &lostRoleResponse{ProvisioningBackend: f.service}
|
||||||
|
if err := failed.Reconcile(f.ctx, database.Name); err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
f.reconcile(t, database, provisioning.Conflict)
|
||||||
|
if database.Status.RoleOID != 0 || f.queryPostgres(t, "SELECT count(*) FROM pg_roles WHERE rolname='lost_role_response'") != "1" {
|
||||||
|
t.Fatal("丢失结果后应保留未认领角色")
|
||||||
|
}
|
||||||
|
})
|
||||||
|
t.Run("同名角色被重建", func(t *testing.T) {
|
||||||
|
database := f.bound(t, "replaced_role")
|
||||||
|
f.reconcile(t, database, provisioning.Pending)
|
||||||
|
f.queryPostgres(t, "DROP ROLE replaced_role; CREATE ROLE replaced_role LOGIN")
|
||||||
|
f.reconcile(t, database, provisioning.Conflict)
|
||||||
|
if database.Status.DatabaseOID != 0 {
|
||||||
|
t.Fatal("重建角色被当成原 owner 继续供应")
|
||||||
|
}
|
||||||
|
})
|
||||||
|
t.Run("创建步骤并发只有一个获准", func(t *testing.T) {
|
||||||
|
database := f.bound(t, "concurrent_supply")
|
||||||
|
service := f.usecase
|
||||||
|
resources := &concurrentProvisioningResources{ProvisioningResources: f.resources}
|
||||||
|
resources.loaded.Add(2)
|
||||||
|
service.Resources = resources
|
||||||
|
results := make(chan error, 2)
|
||||||
|
for range 2 {
|
||||||
|
go func() { results <- service.Reconcile(f.ctx, database.Name) }()
|
||||||
|
}
|
||||||
|
success, conflict := 0, 0
|
||||||
|
for range 2 {
|
||||||
|
err := <-results
|
||||||
|
if err == nil {
|
||||||
|
success++
|
||||||
|
} else if apierrors.IsConflict(err) {
|
||||||
|
conflict++
|
||||||
|
} else {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
if success != 1 || conflict != 1 {
|
||||||
|
t.Fatal("同一快照不得重复授权创建")
|
||||||
|
}
|
||||||
|
f.status(t, database, provisioning.Pending)
|
||||||
|
if database.Status.RoleOID == 0 {
|
||||||
|
t.Fatal("胜方未确认角色")
|
||||||
|
}
|
||||||
|
})
|
||||||
|
t.Run("实际manager观察与重启", func(t *testing.T) { testProvisioningManager(t, f) })
|
||||||
|
}
|
||||||
|
|
||||||
|
type failRoleConfirmation struct {
|
||||||
|
application.ProvisioningResources
|
||||||
|
}
|
||||||
|
|
||||||
|
func (r *failRoleConfirmation) Save(ctx context.Context, record *application.ProvisioningRecord, state provisioning.State) (*application.ProvisioningRecord, error) {
|
||||||
|
if state.RoleOID != 0 {
|
||||||
|
return nil, errors.New("injected confirmation persistence failure")
|
||||||
|
}
|
||||||
|
return r.ProvisioningResources.Save(ctx, record, state)
|
||||||
|
}
|
||||||
|
|
||||||
|
type failDatabaseConfirmation struct {
|
||||||
|
application.ProvisioningResources
|
||||||
|
}
|
||||||
|
|
||||||
|
func (r *failDatabaseConfirmation) Save(ctx context.Context, record *application.ProvisioningRecord, state provisioning.State) (*application.ProvisioningRecord, error) {
|
||||||
|
if state.DatabaseOID != 0 {
|
||||||
|
return nil, errors.New("injected database confirmation persistence failure")
|
||||||
|
}
|
||||||
|
return r.ProvisioningResources.Save(ctx, record, state)
|
||||||
|
}
|
||||||
|
|
||||||
|
type lostRoleResponse struct {
|
||||||
|
application.ProvisioningBackend
|
||||||
|
}
|
||||||
|
|
||||||
|
func (b *lostRoleResponse) CreateLoginRole(ctx context.Context, target instance.ObservationTarget, value credential.ApplicationCredential) (uint32, error) {
|
||||||
|
if _, err := b.ProvisioningBackend.CreateLoginRole(ctx, target, value); err != nil {
|
||||||
|
return 0, err
|
||||||
|
}
|
||||||
|
return 0, application.ErrResourceUncertain
|
||||||
|
}
|
||||||
|
|
||||||
|
type concurrentProvisioningResources struct {
|
||||||
|
application.ProvisioningResources
|
||||||
|
readers atomic.Int32
|
||||||
|
loaded sync.WaitGroup
|
||||||
|
}
|
||||||
|
|
||||||
|
func (r *concurrentProvisioningResources) Load(ctx context.Context, name string) (*application.ProvisioningRecord, error) {
|
||||||
|
record, err := r.ProvisioningResources.Load(ctx, name)
|
||||||
|
if r.readers.Add(1) <= 2 {
|
||||||
|
r.loaded.Done()
|
||||||
|
r.loaded.Wait()
|
||||||
|
}
|
||||||
|
return record, err
|
||||||
|
}
|
||||||
@@ -0,0 +1,97 @@
|
|||||||
|
//go:build integration
|
||||||
|
|
||||||
|
package postgresql_test
|
||||||
|
|
||||||
|
import (
|
||||||
|
"context"
|
||||||
|
"sync/atomic"
|
||||||
|
"testing"
|
||||||
|
"time"
|
||||||
|
|
||||||
|
"git.ddupan.top/panxiao81/ayatori/internal/database/adapter/kubernetes"
|
||||||
|
"git.ddupan.top/panxiao81/ayatori/internal/database/application"
|
||||||
|
databasecontroller "git.ddupan.top/panxiao81/ayatori/internal/database/controller"
|
||||||
|
"git.ddupan.top/panxiao81/ayatori/internal/database/domain/instance"
|
||||||
|
"git.ddupan.top/panxiao81/ayatori/internal/database/domain/provisioning"
|
||||||
|
"k8s.io/apimachinery/pkg/api/meta"
|
||||||
|
ctrl "sigs.k8s.io/controller-runtime"
|
||||||
|
"sigs.k8s.io/controller-runtime/pkg/client"
|
||||||
|
controllerconfig "sigs.k8s.io/controller-runtime/pkg/config"
|
||||||
|
metricsserver "sigs.k8s.io/controller-runtime/pkg/metrics/server"
|
||||||
|
)
|
||||||
|
|
||||||
|
func testProvisioningManager(t *testing.T, f *provisioningFixture) {
|
||||||
|
t.Helper()
|
||||||
|
database := f.bound(t, "watch_supply")
|
||||||
|
// 首次 list/watch 自动供应;第二个新 manager 只观察确认记录,不重新建库。
|
||||||
|
var roleOID, databaseOID int64
|
||||||
|
for range 2 {
|
||||||
|
skipRepeatedName := true
|
||||||
|
manager, err := ctrl.NewManager(f.config, ctrl.Options{Scheme: f.scheme, Metrics: metricsserver.Options{BindAddress: "0"}, HealthProbeBindAddress: "0",
|
||||||
|
Controller: controllerconfig.Controller{SkipNameValidation: &skipRepeatedName}})
|
||||||
|
if err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
resources := &kubernetes.ProvisioningResources{Client: manager.GetClient(), Reader: manager.GetAPIReader()}
|
||||||
|
backend := &countProvisioningObservations{ProvisioningBackend: f.service}
|
||||||
|
service := &application.DatabaseReconciliation{
|
||||||
|
Credentials: &application.CredentialPreparation{Resources: &kubernetes.CredentialResources{Client: manager.GetClient(), Reader: manager.GetAPIReader()}, Store: f.store},
|
||||||
|
Provisioning: &application.DatabaseProvisioning{Resources: resources, Credentials: f.store, Backend: backend},
|
||||||
|
}
|
||||||
|
if err := databasecontroller.NewProvisioningReconciler(manager.GetClient(), service).SetupWithManager(manager); err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
ctx, stop := context.WithCancel(f.ctx)
|
||||||
|
done := make(chan error, 1)
|
||||||
|
go func() { done <- manager.Start(ctx) }()
|
||||||
|
// 使用 scope 确保任何失败都先停止 worker,再由 fixture 关闭共享连接。
|
||||||
|
func() {
|
||||||
|
defer func() {
|
||||||
|
stop()
|
||||||
|
select {
|
||||||
|
case err := <-done:
|
||||||
|
if err != nil {
|
||||||
|
t.Error(err)
|
||||||
|
}
|
||||||
|
case <-time.After(10 * time.Second):
|
||||||
|
t.Error("供应 manager 未停止")
|
||||||
|
}
|
||||||
|
}()
|
||||||
|
if !manager.GetCache().WaitForCacheSync(ctx) {
|
||||||
|
t.Fatal("供应 cache 未同步")
|
||||||
|
}
|
||||||
|
deadline := time.After(15 * time.Second)
|
||||||
|
for {
|
||||||
|
if err := f.api.Get(f.ctx, client.ObjectKeyFromObject(database), database); err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
if meta.IsStatusConditionTrue(database.Status.Conditions, "ResourcesReady") && backend.reads.Load() > 0 {
|
||||||
|
break
|
||||||
|
}
|
||||||
|
select {
|
||||||
|
case <-deadline:
|
||||||
|
t.Fatal("watch 未推动资源创建")
|
||||||
|
case <-time.After(100 * time.Millisecond):
|
||||||
|
}
|
||||||
|
}
|
||||||
|
if roleOID != 0 && (roleOID != database.Status.RoleOID || databaseOID != database.Status.DatabaseOID) {
|
||||||
|
t.Fatal("manager 重启替换已确认资源")
|
||||||
|
}
|
||||||
|
roleOID, databaseOID = database.Status.RoleOID, database.Status.DatabaseOID
|
||||||
|
}()
|
||||||
|
}
|
||||||
|
f.status(t, database, provisioning.Available)
|
||||||
|
}
|
||||||
|
|
||||||
|
type countProvisioningObservations struct {
|
||||||
|
application.ProvisioningBackend
|
||||||
|
reads atomic.Int32
|
||||||
|
}
|
||||||
|
|
||||||
|
func (b *countProvisioningObservations) InspectResources(ctx context.Context, target instance.ObservationTarget, name, role string) (provisioning.Observation, error) {
|
||||||
|
observation, err := b.ProvisioningBackend.InspectResources(ctx, target, name, role)
|
||||||
|
if err == nil && name == "watch_supply" {
|
||||||
|
b.reads.Add(1)
|
||||||
|
}
|
||||||
|
return observation, err
|
||||||
|
}
|
||||||
@@ -0,0 +1,159 @@
|
|||||||
|
package postgresql
|
||||||
|
|
||||||
|
import (
|
||||||
|
"context"
|
||||||
|
"errors"
|
||||||
|
"regexp"
|
||||||
|
|
||||||
|
"git.ddupan.top/panxiao81/ayatori/internal/database/application"
|
||||||
|
"git.ddupan.top/panxiao81/ayatori/internal/database/domain/credential"
|
||||||
|
"git.ddupan.top/panxiao81/ayatori/internal/database/domain/provisioning"
|
||||||
|
"github.com/jackc/pgx/v5"
|
||||||
|
"github.com/jackc/pgx/v5/pgconn"
|
||||||
|
)
|
||||||
|
|
||||||
|
var resourceIdentifier = regexp.MustCompile(`^[a-z][a-z0-9_]{0,62}$`)
|
||||||
|
|
||||||
|
const inspectResourcesSQL = `
|
||||||
|
SELECT COALESCE(r.oid, 0), COALESCE(d.oid, 0),
|
||||||
|
COALESCE(r.rolcanlogin AND NOT (r.rolsuper OR r.rolcreatedb OR r.rolcreaterole OR r.rolreplication OR r.rolbypassrls)
|
||||||
|
AND NOT EXISTS (SELECT FROM pg_catalog.pg_auth_members m WHERE m.member = r.oid), false),
|
||||||
|
COALESCE(d.datdba, 0), COALESCE(d.datallowconn, false),
|
||||||
|
EXISTS (SELECT FROM pg_catalog.aclexplode(COALESCE(d.datacl, pg_catalog.acldefault('d', d.datdba))) a
|
||||||
|
WHERE a.grantee = 0 AND a.privilege_type = 'CONNECT')
|
||||||
|
FROM (SELECT 1) seed
|
||||||
|
LEFT JOIN pg_catalog.pg_roles r ON r.rolname = $1
|
||||||
|
LEFT JOIN pg_catalog.pg_database d ON d.datname = $2`
|
||||||
|
|
||||||
|
type resourceReader interface {
|
||||||
|
QueryRow(context.Context, string, ...any) pgx.Row
|
||||||
|
}
|
||||||
|
|
||||||
|
func inspectResources(ctx context.Context, reader resourceReader, name, role string) (provisioning.Observation, error) {
|
||||||
|
var result provisioning.Observation
|
||||||
|
if !resourceIdentifier.MatchString(name) || !resourceIdentifier.MatchString(role) {
|
||||||
|
return result, application.ErrResourceConflict
|
||||||
|
}
|
||||||
|
err := reader.QueryRow(ctx, inspectResourcesSQL, role, name).Scan(
|
||||||
|
&result.RoleOID, &result.DatabaseOID, &result.RoleSafe, &result.OwnerOID, &result.AllowConnections, &result.PublicConnect,
|
||||||
|
)
|
||||||
|
if err != nil {
|
||||||
|
return provisioning.Observation{}, application.ErrResourceUnavailable
|
||||||
|
}
|
||||||
|
return result, nil
|
||||||
|
}
|
||||||
|
|
||||||
|
func (d *database) InspectResources(ctx context.Context, name, role string) (provisioning.Observation, error) {
|
||||||
|
return inspectResources(ctx, d.pool, name, role)
|
||||||
|
}
|
||||||
|
|
||||||
|
// CREATE ROLE 与 membership 在一个原生事务提交;不修改任何已有角色。
|
||||||
|
// simple protocol 使用 pgx 的参数转义,避免自行拼接密码字面量;错误不带 SQL 或驱动响应。
|
||||||
|
func (d *database) CreateLoginRole(ctx context.Context, value credential.ApplicationCredential) (uint32, error) {
|
||||||
|
if value.Validate() != nil {
|
||||||
|
return 0, application.ErrResourceConflict
|
||||||
|
}
|
||||||
|
data := value.SecretData()
|
||||||
|
role := data["username"].(string)
|
||||||
|
transaction, err := d.pool.Begin(ctx)
|
||||||
|
if err != nil {
|
||||||
|
return 0, application.ErrResourceUnavailable
|
||||||
|
}
|
||||||
|
defer func() { _ = transaction.Rollback(ctx) }()
|
||||||
|
statement := "CREATE ROLE " + pgx.Identifier{role}.Sanitize() + " LOGIN NOSUPERUSER NOCREATEDB NOCREATEROLE NOREPLICATION NOBYPASSRLS NOINHERIT PASSWORD $1"
|
||||||
|
if _, err := transaction.Exec(ctx, statement, pgx.QueryExecModeSimpleProtocol, data["password"]); err != nil {
|
||||||
|
return 0, creationError(err)
|
||||||
|
}
|
||||||
|
if _, err := transaction.Exec(ctx, "GRANT "+pgx.Identifier{role}.Sanitize()+" TO CURRENT_USER WITH SET TRUE, INHERIT FALSE"); err != nil {
|
||||||
|
return 0, application.ErrResourceUncertain
|
||||||
|
}
|
||||||
|
var oid uint32
|
||||||
|
if err := transaction.QueryRow(ctx, "SELECT oid FROM pg_catalog.pg_roles WHERE rolname=$1", role).Scan(&oid); err != nil {
|
||||||
|
return 0, application.ErrResourceUncertain
|
||||||
|
}
|
||||||
|
if err := transaction.Commit(ctx); err != nil {
|
||||||
|
return 0, application.ErrResourceUncertain
|
||||||
|
}
|
||||||
|
observed, err := d.InspectResources(ctx, data["database"].(string), role)
|
||||||
|
if err != nil || observed.RoleOID != oid || !observed.RoleSafe {
|
||||||
|
return 0, application.ErrResourceUncertain
|
||||||
|
}
|
||||||
|
return oid, nil
|
||||||
|
}
|
||||||
|
|
||||||
|
func (d *database) CreateOwnedDatabase(ctx context.Context, name, role string, roleOID uint32) (uint32, error) {
|
||||||
|
observed, err := d.InspectResources(ctx, name, role)
|
||||||
|
if err != nil {
|
||||||
|
return 0, err
|
||||||
|
}
|
||||||
|
if roleOID == 0 || observed.RoleOID != roleOID || !observed.RoleSafe || observed.DatabaseOID != 0 {
|
||||||
|
return 0, application.ErrResourceConflict
|
||||||
|
}
|
||||||
|
// 不能放入事务。先关闭连接入口,避免默认 PUBLIC CONNECT 暴露未收紧的数据库。
|
||||||
|
statement := "CREATE DATABASE " + pgx.Identifier{name}.Sanitize() + " OWNER " + pgx.Identifier{role}.Sanitize() + " ALLOW_CONNECTIONS false"
|
||||||
|
if _, err := d.pool.Exec(ctx, statement); err != nil {
|
||||||
|
return 0, creationError(err)
|
||||||
|
}
|
||||||
|
observed, err = d.InspectResources(ctx, name, role)
|
||||||
|
if err != nil || observed.DatabaseOID == 0 || observed.OwnerOID != roleOID || observed.RoleOID != roleOID {
|
||||||
|
return 0, application.ErrResourceUncertain
|
||||||
|
}
|
||||||
|
return observed.DatabaseOID, nil
|
||||||
|
}
|
||||||
|
|
||||||
|
func (d *database) ConfigureAccess(ctx context.Context, name, role string, state provisioning.State) error {
|
||||||
|
transaction, err := d.pool.Begin(ctx)
|
||||||
|
if err != nil {
|
||||||
|
return application.ErrResourceUnavailable
|
||||||
|
}
|
||||||
|
defer func() { _ = transaction.Rollback(ctx) }()
|
||||||
|
observed, err := inspectResources(ctx, transaction, name, role)
|
||||||
|
if err != nil {
|
||||||
|
return err
|
||||||
|
}
|
||||||
|
if state.RoleOID == 0 || state.DatabaseOID == 0 || state.Check(observed) != nil {
|
||||||
|
return application.ErrResourceConflict
|
||||||
|
}
|
||||||
|
if _, err := transaction.Exec(ctx, "SET LOCAL ROLE "+pgx.Identifier{role}.Sanitize()); err != nil {
|
||||||
|
return application.ErrResourceUnavailable
|
||||||
|
}
|
||||||
|
for _, statement := range []string{
|
||||||
|
"REVOKE CONNECT ON DATABASE " + pgx.Identifier{name}.Sanitize() + " FROM PUBLIC",
|
||||||
|
"GRANT CONNECT ON DATABASE " + pgx.Identifier{name}.Sanitize() + " TO " + pgx.Identifier{role}.Sanitize(),
|
||||||
|
"ALTER DATABASE " + pgx.Identifier{name}.Sanitize() + " ALLOW_CONNECTIONS true",
|
||||||
|
} {
|
||||||
|
if _, err := transaction.Exec(ctx, statement); err != nil {
|
||||||
|
return application.ErrResourceUnavailable
|
||||||
|
}
|
||||||
|
}
|
||||||
|
// SET LOCAL 不污染池内连接;ACL 失败可在已确认对象上幂等重试,不改密码或 owner。
|
||||||
|
if err := transaction.Commit(ctx); err != nil {
|
||||||
|
return application.ErrResourceUnavailable
|
||||||
|
}
|
||||||
|
observed, err = d.InspectResources(ctx, name, role)
|
||||||
|
if err != nil {
|
||||||
|
return err
|
||||||
|
}
|
||||||
|
if state.Check(observed) != nil {
|
||||||
|
return application.ErrResourceConflict
|
||||||
|
}
|
||||||
|
if observed.PublicConnect || !observed.AllowConnections {
|
||||||
|
return application.ErrResourceUnavailable
|
||||||
|
}
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
|
||||||
|
func creationError(err error) error {
|
||||||
|
if serverError, ok := errors.AsType[*pgconn.PgError](err); ok {
|
||||||
|
switch serverError.Code {
|
||||||
|
case "42710", "42P04", "23505":
|
||||||
|
return application.ErrResourceConflict
|
||||||
|
case "42501", "25006", "28P01", "28000":
|
||||||
|
return application.ErrResourceUnavailable
|
||||||
|
}
|
||||||
|
}
|
||||||
|
if pgconn.SafeToRetry(err) {
|
||||||
|
return application.ErrResourceUnavailable
|
||||||
|
}
|
||||||
|
return application.ErrResourceUncertain
|
||||||
|
}
|
||||||
@@ -0,0 +1,35 @@
|
|||||||
|
package postgresql
|
||||||
|
|
||||||
|
import (
|
||||||
|
"errors"
|
||||||
|
"testing"
|
||||||
|
|
||||||
|
"git.ddupan.top/panxiao81/ayatori/internal/database/application"
|
||||||
|
"github.com/jackc/pgx/v5/pgconn"
|
||||||
|
)
|
||||||
|
|
||||||
|
func TestCreationErrorsDoNotExposeServerDetails(t *testing.T) {
|
||||||
|
for _, test := range []struct {
|
||||||
|
code string
|
||||||
|
want error
|
||||||
|
}{
|
||||||
|
{"42710", application.ErrResourceConflict},
|
||||||
|
{"42P04", application.ErrResourceConflict},
|
||||||
|
{"23505", application.ErrResourceConflict},
|
||||||
|
{"42501", application.ErrResourceUnavailable},
|
||||||
|
{"25006", application.ErrResourceUnavailable},
|
||||||
|
{"57014", application.ErrResourceUncertain},
|
||||||
|
{"XX000", application.ErrResourceUncertain},
|
||||||
|
} {
|
||||||
|
actual := creationError(&pgconn.PgError{Code: test.code, Message: "unsafe SQL and credential detail"})
|
||||||
|
if !errors.Is(actual, test.want) {
|
||||||
|
t.Fatalf("SQLSTATE %s 分类错误", test.code)
|
||||||
|
}
|
||||||
|
if actual.Error() != test.want.Error() {
|
||||||
|
t.Fatal("后端错误携带原始响应")
|
||||||
|
}
|
||||||
|
}
|
||||||
|
if !errors.Is(creationError(errors.New("connection lost after send")), application.ErrResourceUncertain) {
|
||||||
|
t.Fatal("未知网络结果不允许自动重试创建")
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -5,8 +5,7 @@ import (
|
|||||||
"errors"
|
"errors"
|
||||||
"fmt"
|
"fmt"
|
||||||
|
|
||||||
"git.ddupan.top/panxiao81/ayatori/internal/database/domain/binding"
|
credentialdomain "git.ddupan.top/panxiao81/ayatori/internal/database/domain/credential"
|
||||||
"git.ddupan.top/panxiao81/ayatori/internal/database/domain/instance"
|
|
||||||
)
|
)
|
||||||
|
|
||||||
var (
|
var (
|
||||||
@@ -17,54 +16,28 @@ var (
|
|||||||
ErrCredentialUncertain = errors.New("credential creation outcome is uncertain; manual resolution required")
|
ErrCredentialUncertain = errors.New("credential creation outcome is uncertain; manual resolution required")
|
||||||
)
|
)
|
||||||
|
|
||||||
const (
|
|
||||||
CredentialsReady = "CredentialsReady"
|
|
||||||
CredentialCreationStarted = "CreationStarted"
|
|
||||||
CredentialPrepared = "CredentialPrepared"
|
|
||||||
)
|
|
||||||
|
|
||||||
type CredentialLocation struct {
|
|
||||||
Mount string
|
|
||||||
Path string
|
|
||||||
}
|
|
||||||
|
|
||||||
// CredentialStore 只表达本用例需要的凭据操作,不提供覆盖或删除。
|
// CredentialStore 只表达本用例需要的凭据操作,不提供覆盖或删除。
|
||||||
// version=0 的读取只用于观察是否已有值,成功不能作为认领依据。
|
// version=0 的读取只用于观察是否已有值,成功不能作为认领依据。
|
||||||
type CredentialStore interface {
|
type CredentialStore interface {
|
||||||
ProvisionLocation(string) (CredentialLocation, error)
|
ProvisionLocation(string) (credentialdomain.Location, error)
|
||||||
ReadCredential(context.Context, CredentialLocation, int64) (ApplicationCredential, error)
|
ReadCredential(context.Context, credentialdomain.Location, int64) (credentialdomain.ApplicationCredential, error)
|
||||||
CreateCredential(context.Context, CredentialLocation, ApplicationCredential) (int64, error)
|
CreateCredential(context.Context, credentialdomain.Location, credentialdomain.ApplicationCredential) (int64, error)
|
||||||
}
|
|
||||||
|
|
||||||
type CredentialInstance struct {
|
|
||||||
binding.Instance
|
|
||||||
Generation int64
|
|
||||||
Endpoint instance.Endpoint
|
|
||||||
}
|
}
|
||||||
|
|
||||||
// CredentialRecord 是同一轮观察的事实,状态中永远不保存密码。
|
// CredentialRecord 是同一轮观察的事实,状态中永远不保存密码。
|
||||||
type CredentialRecord struct {
|
type CredentialRecord struct {
|
||||||
Database BindingDatabase
|
credentialdomain.Target
|
||||||
Tenant *BindingTenant
|
Revision string
|
||||||
Instance *CredentialInstance
|
TenantGeneration int64
|
||||||
DatabaseProtected bool
|
InstanceGeneration int64
|
||||||
TenantProtected bool
|
Status credentialdomain.State
|
||||||
Status CredentialStatus
|
|
||||||
}
|
|
||||||
|
|
||||||
type CredentialStatus struct {
|
|
||||||
Location *CredentialLocation
|
|
||||||
Version int64
|
|
||||||
Ready bool
|
|
||||||
Reason string
|
|
||||||
Message string
|
|
||||||
}
|
}
|
||||||
|
|
||||||
// CredentialResources 的写入必须检查 Database UID/resourceVersion,保留其他状态。
|
// CredentialResources 的写入必须检查 Database UID/resourceVersion,保留其他状态。
|
||||||
// CheckCurrent 在外部操作前后回读本轮三个资源,拒绝陈旧快照;它不是跨系统事务。
|
// CheckCurrent 在外部操作前后回读本轮三个资源,拒绝陈旧快照;它不是跨系统事务。
|
||||||
type CredentialResources interface {
|
type CredentialResources interface {
|
||||||
Load(context.Context, string) (*CredentialRecord, error)
|
Load(context.Context, string) (*CredentialRecord, error)
|
||||||
Save(context.Context, *CredentialRecord, CredentialStatus) (*CredentialRecord, error)
|
Save(context.Context, *CredentialRecord, credentialdomain.State) (*CredentialRecord, error)
|
||||||
CheckCurrent(context.Context, *CredentialRecord) error
|
CheckCurrent(context.Context, *CredentialRecord) error
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -75,37 +48,35 @@ type CredentialPreparation struct {
|
|||||||
|
|
||||||
func (s CredentialPreparation) Reconcile(ctx context.Context, name string) error {
|
func (s CredentialPreparation) Reconcile(ctx context.Context, name string) error {
|
||||||
record, err := s.Resources.Load(ctx, name)
|
record, err := s.Resources.Load(ctx, name)
|
||||||
if err != nil || record == nil || record.Database.Source != "Provision" {
|
if err != nil || record == nil || !record.RequiresPreparation() {
|
||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
// 未完成创建的重入不猜测后端结果。即使进程在实际发请求前退出,也需要人工核实。
|
if state, canContinue := record.Status.Resume(); !canContinue {
|
||||||
if record.Status.Version == 0 && record.Status.Reason == binding.Conflict {
|
if state == record.Status {
|
||||||
return nil // 保留首次冲突的具体原因,不因后端恢复而重入创建。
|
return nil
|
||||||
}
|
}
|
||||||
if record.Status.Version == 0 && record.Status.Reason == CredentialCreationStarted {
|
return s.report(ctx, record, state.Phase, state.Message)
|
||||||
return s.report(ctx, record, binding.Conflict,
|
|
||||||
"凭据创建未留下成功确认;请核对固定位置与后端历史并人工处理,未重新生成密码")
|
|
||||||
}
|
}
|
||||||
if issue := record.check(); issue != nil {
|
if issue := record.Check(); issue != nil {
|
||||||
return s.report(ctx, record, issue.Reason, issue.Message)
|
return s.report(ctx, record, issue.Phase, issue.Message)
|
||||||
}
|
}
|
||||||
location, err := s.Store.ProvisionLocation(record.Database.Identity.UID)
|
location, err := s.Store.ProvisionLocation(record.Database.Identity.UID)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return s.report(ctx, record, binding.DependencyUnavailable, "凭据存储位置配置无效,未执行外部写入")
|
return s.report(ctx, record, credentialdomain.Unavailable, "凭据存储位置配置无效,未执行外部写入")
|
||||||
|
}
|
||||||
|
if issue := record.Status.CheckLocation(location); issue != nil {
|
||||||
|
return s.report(ctx, record, issue.Phase, issue.Message)
|
||||||
}
|
}
|
||||||
if record.Status.Location == nil {
|
if record.Status.Location == nil {
|
||||||
status := record.Status
|
status := record.Status
|
||||||
status.Location = &location
|
status.Location = &location
|
||||||
status.Ready, status.Reason, status.Message = false, "LocationPinned", "凭据位置已固定,等待创建"
|
status = status.WithPhase(credentialdomain.Pinned, "凭据位置已固定,等待创建")
|
||||||
record, err = s.Resources.Save(ctx, record, status)
|
record, err = s.Resources.Save(ctx, record, status)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
} else if *record.Status.Location != location {
|
|
||||||
return s.report(ctx, record, binding.DependencyUnavailable,
|
|
||||||
"部署配置与固定凭据位置不一致;请恢复原 mount/path 配置,未迁移或改密")
|
|
||||||
}
|
}
|
||||||
if record.Status.Version > 0 {
|
if record.Status.Confirmed() {
|
||||||
return s.observe(ctx, record)
|
return s.observe(ctx, record)
|
||||||
}
|
}
|
||||||
return s.create(ctx, record)
|
return s.create(ctx, record)
|
||||||
@@ -113,19 +84,18 @@ func (s CredentialPreparation) Reconcile(ctx context.Context, name string) error
|
|||||||
|
|
||||||
func (s CredentialPreparation) create(ctx context.Context, record *CredentialRecord) error {
|
func (s CredentialPreparation) create(ctx context.Context, record *CredentialRecord) error {
|
||||||
_, err := s.Store.ReadCredential(ctx, *record.Status.Location, 0)
|
_, err := s.Store.ReadCredential(ctx, *record.Status.Location, 0)
|
||||||
if err == nil || errors.Is(err, ErrApplicationCredentialInvalid) {
|
if err == nil || errors.Is(err, credentialdomain.ErrApplicationCredentialInvalid) {
|
||||||
return s.report(ctx, record, binding.Conflict, "固定位置已有未确认的凭据;请人工核实,未认领或覆盖")
|
return s.report(ctx, record, credentialdomain.Conflict, "固定位置已有未确认的凭据;请人工核实,未认领或覆盖")
|
||||||
}
|
}
|
||||||
if !errors.Is(err, ErrCredentialNotFound) {
|
if !errors.Is(err, ErrCredentialNotFound) {
|
||||||
return s.report(ctx, record, binding.DependencyUnavailable, "创建前无法确认凭据位置是否为空,等待依赖恢复")
|
return s.report(ctx, record, credentialdomain.Unavailable, "创建前无法确认凭据位置是否为空,等待依赖恢复")
|
||||||
}
|
}
|
||||||
credential, err := GenerateApplicationCredential(record.Database.LoginRole, record.Database.Name, record.Instance.Endpoint)
|
credential, err := credentialdomain.GenerateApplicationCredential(record.Database.LoginRole, record.Database.Name, record.Instance.Endpoint)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return s.report(ctx, record, "InvalidTarget", "应用凭据目标无效,未执行外部写入")
|
return s.report(ctx, record, credentialdomain.InvalidTarget, "应用凭据目标无效,未执行外部写入")
|
||||||
}
|
}
|
||||||
status := record.Status
|
status := record.Status
|
||||||
status.Ready, status.Reason = false, CredentialCreationStarted
|
status = status.WithPhase(credentialdomain.Creating, "凭据创建已开始;尚无成功确认时不得重入创建")
|
||||||
status.Message = "凭据创建已开始;尚无成功确认时不得重入创建"
|
|
||||||
record, err = s.Resources.Save(ctx, record, status)
|
record, err = s.Resources.Save(ctx, record, status)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return err
|
return err
|
||||||
@@ -136,46 +106,47 @@ func (s CredentialPreparation) create(ctx context.Context, record *CredentialRec
|
|||||||
version, err := s.Store.CreateCredential(ctx, *record.Status.Location, credential)
|
version, err := s.Store.CreateCredential(ctx, *record.Status.Location, credential)
|
||||||
if errors.Is(err, ErrCredentialUnavailable) {
|
if errors.Is(err, ErrCredentialUnavailable) {
|
||||||
// 适配器只在明确未执行写入(认证拒绝或请求前取消)时返回此错误。
|
// 适配器只在明确未执行写入(认证拒绝或请求前取消)时返回此错误。
|
||||||
return s.report(ctx, record, binding.DependencyUnavailable, "凭据创建在执行前被拒绝,等待认证或权限恢复")
|
return s.report(ctx, record, credentialdomain.Unavailable, "凭据创建在执行前被拒绝,等待认证或权限恢复")
|
||||||
}
|
}
|
||||||
if err != nil || version != 1 {
|
if err != nil {
|
||||||
return s.report(ctx, record, binding.Conflict,
|
return s.report(ctx, record, credentialdomain.Conflict,
|
||||||
"凭据创建冲突或结果不确定;请核对固定位置的版本历史,未认领、覆盖或重新生成密码")
|
"凭据创建冲突或结果不确定;请核对固定位置的版本历史,未认领、覆盖或重新生成密码")
|
||||||
}
|
}
|
||||||
|
confirmed, issue := record.Status.Created(version)
|
||||||
|
if issue != nil {
|
||||||
|
return s.report(ctx, record, issue.Phase, issue.Message)
|
||||||
|
}
|
||||||
if err := s.Resources.CheckCurrent(ctx, record); err != nil {
|
if err := s.Resources.CheckCurrent(ctx, record); err != nil {
|
||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
status = record.Status
|
_, err = s.Resources.Save(ctx, record, confirmed)
|
||||||
status.Version, status.Ready, status.Reason = version, true, CredentialPrepared
|
|
||||||
status.Message = "凭据已创建并回读确认;尚未创建 PostgreSQL 资源或交付给 Tenant"
|
|
||||||
_, err = s.Resources.Save(ctx, record, status)
|
|
||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
|
|
||||||
func (s CredentialPreparation) observe(ctx context.Context, record *CredentialRecord) error {
|
func (s CredentialPreparation) observe(ctx context.Context, record *CredentialRecord) error {
|
||||||
credential, err := s.Store.ReadCredential(ctx, *record.Status.Location, record.Status.Version)
|
credential, err := s.Store.ReadCredential(ctx, *record.Status.Location, record.Status.Version)
|
||||||
if errors.Is(err, ErrCredentialConflict) || errors.Is(err, ErrCredentialNotFound) {
|
if errors.Is(err, ErrCredentialConflict) || errors.Is(err, ErrCredentialNotFound) {
|
||||||
return s.report(ctx, record, binding.Conflict, "已确认凭据消失、版本变化或内容无效;请人工核实,未生成替代密码")
|
return s.report(ctx, record, credentialdomain.Conflict, "已确认凭据消失、版本变化或内容无效;请人工核实,未生成替代密码")
|
||||||
}
|
}
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return s.report(ctx, record, binding.DependencyUnavailable, "已确认凭据暂时无法读取;保留确认版本,等待依赖恢复")
|
return s.report(ctx, record, credentialdomain.Unavailable, "已确认凭据暂时无法读取;保留确认版本,等待依赖恢复")
|
||||||
}
|
}
|
||||||
if !credential.MatchesTarget(record.Database.LoginRole, record.Database.Name, record.Instance.Endpoint) {
|
if issue := record.CheckCredential(credential); issue != nil {
|
||||||
return s.report(ctx, record, binding.Conflict, "已确认凭据与当前 Instance/database/loginRole 不一致;请人工核实,未修改凭据")
|
return s.report(ctx, record, issue.Phase, issue.Message)
|
||||||
}
|
}
|
||||||
if err := s.Resources.CheckCurrent(ctx, record); err != nil {
|
if err := s.Resources.CheckCurrent(ctx, record); err != nil {
|
||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
status := record.Status
|
status := record.Status
|
||||||
status.Ready, status.Reason = true, CredentialPrepared
|
status.Phase = credentialdomain.Prepared
|
||||||
status.Message = "已确认凭据可读取;尚未验证 PostgreSQL 资源或完成 Tenant 交付"
|
status.Message = "已确认凭据可读取;尚未验证 PostgreSQL 资源或完成 Tenant 交付"
|
||||||
_, err = s.Resources.Save(ctx, record, status)
|
_, err = s.Resources.Save(ctx, record, status)
|
||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
|
|
||||||
func (s CredentialPreparation) report(ctx context.Context, record *CredentialRecord, reason, message string) error {
|
func (s CredentialPreparation) report(ctx context.Context, record *CredentialRecord, phase credentialdomain.Phase, message string) error {
|
||||||
status := record.Status
|
status := record.Status
|
||||||
status.Ready, status.Reason = false, reason
|
status.Phase = phase
|
||||||
status.Message = fmt.Sprintf("Database %s:%s", record.Database.Identity.Name, message)
|
status.Message = fmt.Sprintf("Database %s:%s", record.Database.Identity.Name, message)
|
||||||
_, err := s.Resources.Save(ctx, record, status)
|
_, err := s.Resources.Save(ctx, record, status)
|
||||||
return err
|
return err
|
||||||
|
|||||||
@@ -5,6 +5,8 @@ import (
|
|||||||
"errors"
|
"errors"
|
||||||
"testing"
|
"testing"
|
||||||
|
|
||||||
|
credentialdomain "git.ddupan.top/panxiao81/ayatori/internal/database/domain/credential"
|
||||||
|
|
||||||
"git.ddupan.top/panxiao81/ayatori/internal/database/domain/binding"
|
"git.ddupan.top/panxiao81/ayatori/internal/database/domain/binding"
|
||||||
"git.ddupan.top/panxiao81/ayatori/internal/database/domain/instance"
|
"git.ddupan.top/panxiao81/ayatori/internal/database/domain/instance"
|
||||||
)
|
)
|
||||||
@@ -25,14 +27,14 @@ func preparationRecord(t *testing.T) *CredentialRecord {
|
|||||||
t.Fatal(err)
|
t.Fatal(err)
|
||||||
}
|
}
|
||||||
return &CredentialRecord{
|
return &CredentialRecord{
|
||||||
Database: BindingDatabase{Database: binding.Database{
|
Database: binding.Database{
|
||||||
Identity: database, Instance: preparationInstanceName, InstanceUID: preparationInstanceUID, Name: bindingTestName, LoginRole: bindingTestName, Source: "Provision", Tenant: &tenant,
|
Identity: database, Instance: preparationInstanceName, InstanceUID: preparationInstanceUID, Name: bindingTestName, LoginRole: bindingTestName, Source: "Provision", Tenant: &tenant,
|
||||||
}},
|
},
|
||||||
Tenant: &BindingTenant{
|
Tenant: &binding.Tenant{
|
||||||
Identity: tenant, Phase: binding.Bound, Database: &database,
|
Identity: tenant, Phase: binding.Bound, Database: &database,
|
||||||
Request: binding.Request{Provision: &binding.ProvisionRequest{Instance: preparationInstanceName}},
|
Request: binding.Request{Provision: &binding.ProvisionRequest{Instance: preparationInstanceName}},
|
||||||
},
|
},
|
||||||
Instance: &CredentialInstance{Identity: binding.Identity{Name: preparationInstanceName, UID: preparationInstanceUID}, Ready: true, Endpoint: endpoint},
|
Instance: &credentialdomain.Instance{Identity: binding.Identity{Name: preparationInstanceName, UID: preparationInstanceUID}, Ready: true, Endpoint: endpoint},
|
||||||
DatabaseProtected: true, TenantProtected: true,
|
DatabaseProtected: true, TenantProtected: true,
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
@@ -48,7 +50,7 @@ func (r *memoryCredentialResources) Load(context.Context, string) (*CredentialRe
|
|||||||
return ©, nil
|
return ©, nil
|
||||||
}
|
}
|
||||||
|
|
||||||
func (r *memoryCredentialResources) Save(_ context.Context, record *CredentialRecord, status CredentialStatus) (*CredentialRecord, error) {
|
func (r *memoryCredentialResources) Save(_ context.Context, record *CredentialRecord, status credentialdomain.State) (*CredentialRecord, error) {
|
||||||
if r.saveError != nil {
|
if r.saveError != nil {
|
||||||
return nil, r.saveError
|
return nil, r.saveError
|
||||||
}
|
}
|
||||||
@@ -69,16 +71,16 @@ type preparationStore struct {
|
|||||||
createError error
|
createError error
|
||||||
}
|
}
|
||||||
|
|
||||||
func (*preparationStore) ProvisionLocation(uid string) (CredentialLocation, error) {
|
func (*preparationStore) ProvisionLocation(uid string) (credentialdomain.Location, error) {
|
||||||
return CredentialLocation{Mount: "applications", Path: "database/" + uid}, nil
|
return credentialdomain.Location{Mount: "applications", Path: "database/" + uid}, nil
|
||||||
}
|
}
|
||||||
|
|
||||||
func (s *preparationStore) ReadCredential(context.Context, CredentialLocation, int64) (ApplicationCredential, error) {
|
func (s *preparationStore) ReadCredential(context.Context, credentialdomain.Location, int64) (credentialdomain.ApplicationCredential, error) {
|
||||||
s.reads++
|
s.reads++
|
||||||
return ApplicationCredential{}, s.readError
|
return credentialdomain.ApplicationCredential{}, s.readError
|
||||||
}
|
}
|
||||||
|
|
||||||
func (s *preparationStore) CreateCredential(context.Context, CredentialLocation, ApplicationCredential) (int64, error) {
|
func (s *preparationStore) CreateCredential(context.Context, credentialdomain.Location, credentialdomain.ApplicationCredential) (int64, error) {
|
||||||
s.creates++
|
s.creates++
|
||||||
if s.createError != nil {
|
if s.createError != nil {
|
||||||
return 0, s.createError
|
return 0, s.createError
|
||||||
@@ -159,7 +161,7 @@ func TestCredentialPreparationWriteBoundary(t *testing.T) {
|
|||||||
if err := service.Reconcile(t.Context(), resources.record.Database.Identity.Name); err != nil {
|
if err := service.Reconcile(t.Context(), resources.record.Database.Identity.Name); err != nil {
|
||||||
t.Fatal(err)
|
t.Fatal(err)
|
||||||
}
|
}
|
||||||
if store.creates != test.wantCreates || resources.record.Status.Reason != binding.Conflict {
|
if store.creates != test.wantCreates || resources.record.Status.Phase != credentialdomain.Conflict {
|
||||||
t.Fatal("未确认创建重入时不得生成替代密码")
|
t.Fatal("未确认创建重入时不得生成替代密码")
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -1,27 +0,0 @@
|
|||||||
package application
|
|
||||||
|
|
||||||
import "git.ddupan.top/panxiao81/ayatori/internal/database/domain/binding"
|
|
||||||
|
|
||||||
func (r *CredentialRecord) check() *binding.Issue {
|
|
||||||
database := r.Database
|
|
||||||
if database.Deleting || database.Phase == binding.Deleting || database.Phase == "Released" {
|
|
||||||
return &binding.Issue{Reason: "PreparationStopped", Message: "Database 正在删除或已释放;保留凭据与 finalizer,不执行供应或清理"}
|
|
||||||
}
|
|
||||||
if database.Tenant == nil || r.Tenant == nil || r.Tenant.Database == nil {
|
|
||||||
return &binding.Issue{Reason: binding.DependencyUnavailable, Message: "等待 Database 与 Tenant 双向绑定完成"}
|
|
||||||
}
|
|
||||||
if *database.Tenant != r.Tenant.Identity || *r.Tenant.Database != database.Identity {
|
|
||||||
return &binding.Issue{Reason: binding.Conflict, Message: "双向绑定的名称或 UID 不匹配,未创建凭据"}
|
|
||||||
}
|
|
||||||
if r.Tenant.Deleting || r.Tenant.Phase != binding.Bound || !r.DatabaseProtected || !r.TenantProtected {
|
|
||||||
return &binding.Issue{Reason: "PreparationStopped", Message: "Tenant 未完成绑定、正在删除或缺少 finalizer 保护,未创建凭据"}
|
|
||||||
}
|
|
||||||
target, err := r.Tenant.Request.Resolve(r.Tenant.Identity)
|
|
||||||
if err != nil || (target.Provision != nil && !database.MatchesProvision(target, r.Tenant.Identity)) || target.Name != database.Identity.Name {
|
|
||||||
return &binding.Issue{Reason: binding.Conflict, Message: "Tenant 申请与 Database 目标不一致,未创建凭据"}
|
|
||||||
}
|
|
||||||
if r.Instance == nil || database.InstanceUID == "" {
|
|
||||||
return &binding.Issue{Reason: binding.DependencyUnavailable, Message: "等待 Instance 与已记录的实例身份"}
|
|
||||||
}
|
|
||||||
return r.Instance.Check(&database.Database)
|
|
||||||
}
|
|
||||||
@@ -0,0 +1,152 @@
|
|||||||
|
package application
|
||||||
|
|
||||||
|
import (
|
||||||
|
"context"
|
||||||
|
"errors"
|
||||||
|
"fmt"
|
||||||
|
|
||||||
|
"git.ddupan.top/panxiao81/ayatori/internal/database/domain/credential"
|
||||||
|
"git.ddupan.top/panxiao81/ayatori/internal/database/domain/instance"
|
||||||
|
"git.ddupan.top/panxiao81/ayatori/internal/database/domain/provisioning"
|
||||||
|
)
|
||||||
|
|
||||||
|
type ProvisioningRecord struct {
|
||||||
|
credential.Target
|
||||||
|
Revision string
|
||||||
|
TenantGeneration int64
|
||||||
|
InstanceGeneration int64
|
||||||
|
InstanceTarget instance.ObservationTarget
|
||||||
|
Credentials credential.State
|
||||||
|
State provisioning.State
|
||||||
|
}
|
||||||
|
|
||||||
|
type ProvisioningResources interface {
|
||||||
|
Load(context.Context, string) (*ProvisioningRecord, error)
|
||||||
|
Save(context.Context, *ProvisioningRecord, provisioning.State) (*ProvisioningRecord, error)
|
||||||
|
CheckCurrent(context.Context, *ProvisioningRecord) error
|
||||||
|
}
|
||||||
|
|
||||||
|
type DatabaseProvisioning struct {
|
||||||
|
Resources ProvisioningResources
|
||||||
|
Credentials CredentialStore
|
||||||
|
Backend ProvisioningBackend
|
||||||
|
}
|
||||||
|
|
||||||
|
func (s DatabaseProvisioning) Reconcile(ctx context.Context, name string) error {
|
||||||
|
record, err := s.Resources.Load(ctx, name)
|
||||||
|
if err != nil || record == nil || !record.RequiresPreparation() {
|
||||||
|
return err
|
||||||
|
}
|
||||||
|
if state, resume := record.State.Resume(); !resume {
|
||||||
|
if state == record.State {
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
_, err := s.Resources.Save(ctx, record, state)
|
||||||
|
return err
|
||||||
|
}
|
||||||
|
if issue := record.Check(); issue != nil {
|
||||||
|
phase := provisioning.Unavailable
|
||||||
|
if issue.Phase == credential.Conflict {
|
||||||
|
phase = provisioning.Conflict
|
||||||
|
}
|
||||||
|
if issue.Phase == credential.Stopped {
|
||||||
|
phase = provisioning.Stopped
|
||||||
|
}
|
||||||
|
return s.report(ctx, record, phase, issue.Message)
|
||||||
|
}
|
||||||
|
if !record.Credentials.Confirmed() || record.Credentials.Location == nil || record.Credentials.Phase != credential.Prepared {
|
||||||
|
return s.report(ctx, record, provisioning.Unavailable, "等待已确认且当前可用的应用凭据")
|
||||||
|
}
|
||||||
|
value, err := s.Credentials.ReadCredential(ctx, *record.Credentials.Location, record.Credentials.Version)
|
||||||
|
if err != nil {
|
||||||
|
return s.report(ctx, record, provisioning.Unavailable, "无法读取已确认凭据;未执行 PostgreSQL 写入")
|
||||||
|
}
|
||||||
|
if issue := record.CheckCredential(value); issue != nil {
|
||||||
|
return s.report(ctx, record, provisioning.Conflict, issue.Message)
|
||||||
|
}
|
||||||
|
observed, err := s.Backend.InspectResources(ctx, record.InstanceTarget, record.Database.Name, record.Database.LoginRole)
|
||||||
|
if err != nil {
|
||||||
|
return s.report(ctx, record, provisioning.Unavailable, "无法验证当前 PostgreSQL 资源和管理权限")
|
||||||
|
}
|
||||||
|
if err := record.State.Check(observed); err != nil {
|
||||||
|
return s.report(ctx, record, provisioning.Conflict, err.Error()+";请人工核对,未认领或覆盖")
|
||||||
|
}
|
||||||
|
if err := s.Resources.CheckCurrent(ctx, record); err != nil {
|
||||||
|
return err
|
||||||
|
}
|
||||||
|
if record.State.RoleOID == 0 {
|
||||||
|
return s.createRole(ctx, record, value)
|
||||||
|
}
|
||||||
|
if record.State.DatabaseOID == 0 {
|
||||||
|
return s.createDatabase(ctx, record)
|
||||||
|
}
|
||||||
|
if err := s.Backend.ConfigureAccess(ctx, record.InstanceTarget, record.Database.Name, record.Database.LoginRole, record.State); err != nil {
|
||||||
|
return s.backendFailure(ctx, record, err, "数据库访问权限收敛")
|
||||||
|
}
|
||||||
|
if err := s.Resources.CheckCurrent(ctx, record); err != nil {
|
||||||
|
return err
|
||||||
|
}
|
||||||
|
return s.report(ctx, record, provisioning.Available, "角色和数据库已确认,PUBLIC CONNECT 已撤销;扩展与 Tenant 交付尚未完成")
|
||||||
|
}
|
||||||
|
|
||||||
|
func (s DatabaseProvisioning) createRole(ctx context.Context, record *ProvisioningRecord, value credential.ApplicationCredential) error {
|
||||||
|
record, err := s.Resources.Save(ctx, record, record.State.WithPhase(provisioning.CreatingRole, "开始创建登录角色;尚未持久确认"))
|
||||||
|
if err != nil {
|
||||||
|
return err
|
||||||
|
}
|
||||||
|
if err := s.Resources.CheckCurrent(ctx, record); err != nil {
|
||||||
|
return err
|
||||||
|
}
|
||||||
|
oid, err := s.Backend.CreateLoginRole(ctx, record.InstanceTarget, value)
|
||||||
|
if err != nil {
|
||||||
|
return s.backendFailure(ctx, record, err, "登录角色创建")
|
||||||
|
}
|
||||||
|
if oid == 0 {
|
||||||
|
return s.backendFailure(ctx, record, ErrResourceUncertain, "登录角色创建")
|
||||||
|
}
|
||||||
|
if err := s.Resources.CheckCurrent(ctx, record); err != nil {
|
||||||
|
return err
|
||||||
|
}
|
||||||
|
state := record.State.WithPhase(provisioning.Pending, "登录角色已创建并回读,等待创建数据库")
|
||||||
|
state.RoleOID = oid
|
||||||
|
_, err = s.Resources.Save(ctx, record, state)
|
||||||
|
return err
|
||||||
|
}
|
||||||
|
|
||||||
|
func (s DatabaseProvisioning) createDatabase(ctx context.Context, record *ProvisioningRecord) error {
|
||||||
|
record, err := s.Resources.Save(ctx, record, record.State.WithPhase(provisioning.CreatingDatabase, "开始创建数据库;尚未持久确认"))
|
||||||
|
if err != nil {
|
||||||
|
return err
|
||||||
|
}
|
||||||
|
if err := s.Resources.CheckCurrent(ctx, record); err != nil {
|
||||||
|
return err
|
||||||
|
}
|
||||||
|
oid, err := s.Backend.CreateOwnedDatabase(ctx, record.InstanceTarget, record.Database.Name, record.Database.LoginRole, record.State.RoleOID)
|
||||||
|
if err != nil {
|
||||||
|
return s.backendFailure(ctx, record, err, "数据库创建")
|
||||||
|
}
|
||||||
|
if oid == 0 {
|
||||||
|
return s.backendFailure(ctx, record, ErrResourceUncertain, "数据库创建")
|
||||||
|
}
|
||||||
|
if err := s.Resources.CheckCurrent(ctx, record); err != nil {
|
||||||
|
return err
|
||||||
|
}
|
||||||
|
state := record.State.WithPhase(provisioning.Pending, "数据库已创建并回读,等待收紧访问权限;连接入口仍关闭")
|
||||||
|
state.DatabaseOID = oid
|
||||||
|
_, err = s.Resources.Save(ctx, record, state)
|
||||||
|
return err
|
||||||
|
}
|
||||||
|
|
||||||
|
func (s DatabaseProvisioning) backendFailure(ctx context.Context, record *ProvisioningRecord, err error, step string) error {
|
||||||
|
if errors.Is(err, ErrResourceUnavailable) {
|
||||||
|
return s.report(ctx, record, provisioning.Unavailable, step+"被拒绝或暂不可用,保留已确认步骤并等待依赖恢复")
|
||||||
|
}
|
||||||
|
return s.report(ctx, record, provisioning.Conflict, step+"冲突或结果不确定;请核对目标与已确认 OID,未自动认领、改密或清理")
|
||||||
|
}
|
||||||
|
|
||||||
|
func (s DatabaseProvisioning) report(ctx context.Context, record *ProvisioningRecord, phase provisioning.Phase, message string) error {
|
||||||
|
message = fmt.Sprintf("Database %s / Instance %s / database %s / role %s:%s", record.Database.Identity.Name,
|
||||||
|
record.Database.Instance, record.Database.Name, record.Database.LoginRole, message)
|
||||||
|
_, err := s.Resources.Save(ctx, record, record.State.WithPhase(phase, message))
|
||||||
|
return err
|
||||||
|
}
|
||||||
@@ -0,0 +1,17 @@
|
|||||||
|
package application
|
||||||
|
|
||||||
|
import "context"
|
||||||
|
|
||||||
|
// DatabaseReconciliation 顺序编排两个已有用例,不复制它们的领域判断。
|
||||||
|
// 启用资源创建时由同一个 controller 驱动,避免两个 worker 争写同一 Database status。
|
||||||
|
type DatabaseReconciliation struct {
|
||||||
|
Credentials *CredentialPreparation
|
||||||
|
Provisioning *DatabaseProvisioning
|
||||||
|
}
|
||||||
|
|
||||||
|
func (s DatabaseReconciliation) Reconcile(ctx context.Context, name string) error {
|
||||||
|
if err := s.Credentials.Reconcile(ctx, name); err != nil {
|
||||||
|
return err
|
||||||
|
}
|
||||||
|
return s.Provisioning.Reconcile(ctx, name)
|
||||||
|
}
|
||||||
@@ -35,6 +35,7 @@ var (
|
|||||||
// Database 与 Connector 复用原项目 internal/instance/service.go 的能力边界。
|
// Database 与 Connector 复用原项目 internal/instance/service.go 的能力边界。
|
||||||
// Metadata 只查询版本与可用扩展,不能产生领域 Ready。
|
// Metadata 只查询版本与可用扩展,不能产生领域 Ready。
|
||||||
type Database interface {
|
type Database interface {
|
||||||
|
ResourceDatabase
|
||||||
InspectMetadata(context.Context) (DatabaseMetadata, error)
|
InspectMetadata(context.Context) (DatabaseMetadata, error)
|
||||||
InspectManagement(context.Context) (DatabaseMetadata, error)
|
InspectManagement(context.Context) (DatabaseMetadata, error)
|
||||||
Close()
|
Close()
|
||||||
@@ -93,16 +94,45 @@ func (s *InstanceService) ObserveManagement(ctx context.Context, target instance
|
|||||||
}
|
}
|
||||||
|
|
||||||
func (s *InstanceService) observe(ctx context.Context, target instance.ObservationTarget, management bool) (InstanceObservation, error) {
|
func (s *InstanceService) observe(ctx context.Context, target instance.ObservationTarget, management bool) (InstanceObservation, error) {
|
||||||
if err := target.Validate(); err != nil {
|
|
||||||
return InstanceObservation{}, err
|
|
||||||
}
|
|
||||||
s.mu.Lock()
|
s.mu.Lock()
|
||||||
defer s.mu.Unlock()
|
defer s.mu.Unlock()
|
||||||
|
current, err := s.connection(ctx, target)
|
||||||
|
if err != nil {
|
||||||
|
return InstanceObservation{}, err
|
||||||
|
}
|
||||||
|
name := target.Identity().Name()
|
||||||
|
var metadata DatabaseMetadata
|
||||||
|
if management {
|
||||||
|
metadata, err = current.database.InspectManagement(ctx)
|
||||||
|
} else {
|
||||||
|
metadata, err = current.database.InspectMetadata(ctx)
|
||||||
|
metadata.Management = instance.ManagementChecks{}
|
||||||
|
}
|
||||||
|
if err != nil {
|
||||||
|
s.release(name)
|
||||||
|
return InstanceObservation{}, err
|
||||||
|
}
|
||||||
|
if metadata.Version == "" {
|
||||||
|
s.release(name)
|
||||||
|
return InstanceObservation{}, ErrObservation
|
||||||
|
}
|
||||||
|
if err := s.checkCredentials(ctx, target, current.credentials); err != nil {
|
||||||
|
return InstanceObservation{}, err
|
||||||
|
}
|
||||||
|
return InstanceObservation{target: target, version: metadata.Version,
|
||||||
|
extensions: instance.ObserveExtensionSupport(metadata.AvailableExtensions), management: metadata.Management}, nil
|
||||||
|
}
|
||||||
|
|
||||||
|
// connection 由持有 mu 的调用者使用;不缓存观察结果或管理权限。
|
||||||
|
func (s *InstanceService) connection(ctx context.Context, target instance.ObservationTarget) (*entry, error) {
|
||||||
|
if err := target.Validate(); err != nil {
|
||||||
|
return nil, err
|
||||||
|
}
|
||||||
if s.closed {
|
if s.closed {
|
||||||
return InstanceObservation{}, ErrClosed
|
return nil, ErrClosed
|
||||||
}
|
}
|
||||||
if err := ctx.Err(); err != nil {
|
if err := ctx.Err(); err != nil {
|
||||||
return InstanceObservation{}, err
|
return nil, err
|
||||||
}
|
}
|
||||||
|
|
||||||
// 先读取有效凭据。读取失败时不得继续使用缓存中的旧连接。
|
// 先读取有效凭据。读取失败时不得继续使用缓存中的旧连接。
|
||||||
@@ -110,11 +140,11 @@ func (s *InstanceService) observe(ctx context.Context, target instance.Observati
|
|||||||
credentials, err := s.source.Read(ctx, target.Definition().AdminCredential())
|
credentials, err := s.source.Read(ctx, target.Definition().AdminCredential())
|
||||||
if err != nil {
|
if err != nil {
|
||||||
s.release(name)
|
s.release(name)
|
||||||
return InstanceObservation{}, credentialError(err)
|
return nil, credentialError(err)
|
||||||
}
|
}
|
||||||
if credentials.username == "" || credentials.password == "" {
|
if credentials.username == "" || credentials.password == "" {
|
||||||
s.release(name)
|
s.release(name)
|
||||||
return InstanceObservation{}, ErrCredentialsInvalid
|
return nil, ErrCredentialsInvalid
|
||||||
}
|
}
|
||||||
|
|
||||||
// 连接身份与有效值均未变化时复用 pgxpool;generation 本身不要求换池。
|
// 连接身份与有效值均未变化时复用 pgxpool;generation 本身不要求换池。
|
||||||
@@ -127,7 +157,7 @@ func (s *InstanceService) observe(ctx context.Context, target instance.Observati
|
|||||||
if current == nil {
|
if current == nil {
|
||||||
database, err := s.connector.Connect(ctx, target.Definition().Endpoint(), credentials)
|
database, err := s.connector.Connect(ctx, target.Definition().Endpoint(), credentials)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return InstanceObservation{}, err
|
return nil, err
|
||||||
}
|
}
|
||||||
current = &entry{
|
current = &entry{
|
||||||
target: target,
|
target: target,
|
||||||
@@ -137,39 +167,22 @@ func (s *InstanceService) observe(ctx context.Context, target instance.Observati
|
|||||||
s.entries[name] = current
|
s.entries[name] = current
|
||||||
}
|
}
|
||||||
|
|
||||||
var metadata DatabaseMetadata
|
return current, nil
|
||||||
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 InstanceObservation{}, err
|
|
||||||
}
|
|
||||||
if metadata.Version == "" {
|
|
||||||
s.release(name)
|
|
||||||
return InstanceObservation{}, ErrObservation
|
|
||||||
}
|
}
|
||||||
|
|
||||||
|
func (s *InstanceService) checkCredentials(ctx context.Context, target instance.ObservationTarget, credentials Credentials) error {
|
||||||
|
name := target.Identity().Name()
|
||||||
// 回读后再检查凭据,避免把轮换前取得的结果交给新凭据的调用链。
|
// 回读后再检查凭据,避免把轮换前取得的结果交给新凭据的调用链。
|
||||||
latest, err := s.source.Read(ctx, target.Definition().AdminCredential())
|
latest, err := s.source.Read(ctx, target.Definition().AdminCredential())
|
||||||
if err != nil {
|
if err != nil {
|
||||||
s.release(name)
|
s.release(name)
|
||||||
return InstanceObservation{}, credentialError(err)
|
return credentialError(err)
|
||||||
}
|
}
|
||||||
if latest != credentials {
|
if latest != credentials {
|
||||||
s.release(name)
|
s.release(name)
|
||||||
return InstanceObservation{}, ErrCredentialsChanged
|
return ErrCredentialsChanged
|
||||||
}
|
}
|
||||||
return InstanceObservation{
|
return nil
|
||||||
target: target,
|
|
||||||
version: metadata.Version,
|
|
||||||
extensions: instance.ObserveExtensionSupport(metadata.AvailableExtensions),
|
|
||||||
management: metadata.Management,
|
|
||||||
}, nil
|
|
||||||
}
|
}
|
||||||
|
|
||||||
func credentialError(err error) error {
|
func credentialError(err error) error {
|
||||||
|
|||||||
@@ -37,6 +37,7 @@ func (s *sourceStub) Read(context.Context, instance.CredentialReference) (Creden
|
|||||||
}
|
}
|
||||||
|
|
||||||
type databaseStub struct {
|
type databaseStub struct {
|
||||||
|
ResourceDatabase
|
||||||
closes int
|
closes int
|
||||||
err error
|
err error
|
||||||
metadata DatabaseMetadata
|
metadata DatabaseMetadata
|
||||||
|
|||||||
@@ -0,0 +1,102 @@
|
|||||||
|
package application
|
||||||
|
|
||||||
|
import (
|
||||||
|
"context"
|
||||||
|
"errors"
|
||||||
|
|
||||||
|
"git.ddupan.top/panxiao81/ayatori/internal/database/domain/credential"
|
||||||
|
"git.ddupan.top/panxiao81/ayatori/internal/database/domain/instance"
|
||||||
|
"git.ddupan.top/panxiao81/ayatori/internal/database/domain/provisioning"
|
||||||
|
)
|
||||||
|
|
||||||
|
var (
|
||||||
|
ErrResourceConflict = errors.New("PostgreSQL resource conflict; manual resolution required")
|
||||||
|
ErrResourceUncertain = errors.New("PostgreSQL creation outcome uncertain; manual resolution required")
|
||||||
|
ErrResourceUnavailable = errors.New("PostgreSQL operation unavailable before creation")
|
||||||
|
)
|
||||||
|
|
||||||
|
// ResourceDatabase 是供应实际需要的后端能力,不暴露 SQL 或任意回调给用例。
|
||||||
|
type ResourceDatabase interface {
|
||||||
|
InspectResources(context.Context, string, string) (provisioning.Observation, error)
|
||||||
|
CreateLoginRole(context.Context, credential.ApplicationCredential) (uint32, error)
|
||||||
|
CreateOwnedDatabase(context.Context, string, string, uint32) (uint32, error)
|
||||||
|
ConfigureAccess(context.Context, string, string, provisioning.State) error
|
||||||
|
}
|
||||||
|
|
||||||
|
type ProvisioningBackend interface {
|
||||||
|
InspectResources(context.Context, instance.ObservationTarget, string, string) (provisioning.Observation, error)
|
||||||
|
CreateLoginRole(context.Context, instance.ObservationTarget, credential.ApplicationCredential) (uint32, error)
|
||||||
|
CreateOwnedDatabase(context.Context, instance.ObservationTarget, string, string, uint32) (uint32, error)
|
||||||
|
ConfigureAccess(context.Context, instance.ObservationTarget, string, string, provisioning.State) error
|
||||||
|
}
|
||||||
|
|
||||||
|
func (s *InstanceService) InspectResources(ctx context.Context, target instance.ObservationTarget, name, role string) (provisioning.Observation, error) {
|
||||||
|
var result provisioning.Observation
|
||||||
|
err := s.withManagementConnection(ctx, target, false, func(database Database) (err error) {
|
||||||
|
result, err = database.InspectResources(ctx, name, role)
|
||||||
|
return err
|
||||||
|
})
|
||||||
|
if err != nil {
|
||||||
|
return provisioning.Observation{}, err
|
||||||
|
}
|
||||||
|
return result, nil
|
||||||
|
}
|
||||||
|
|
||||||
|
func (s *InstanceService) CreateLoginRole(ctx context.Context, target instance.ObservationTarget, value credential.ApplicationCredential) (uint32, error) {
|
||||||
|
var oid uint32
|
||||||
|
err := s.withManagementConnection(ctx, target, true, func(database Database) (err error) {
|
||||||
|
oid, err = database.CreateLoginRole(ctx, value)
|
||||||
|
return err
|
||||||
|
})
|
||||||
|
if err != nil {
|
||||||
|
return 0, err
|
||||||
|
}
|
||||||
|
return oid, nil
|
||||||
|
}
|
||||||
|
|
||||||
|
func (s *InstanceService) CreateOwnedDatabase(ctx context.Context, target instance.ObservationTarget, name, role string, roleOID uint32) (uint32, error) {
|
||||||
|
var oid uint32
|
||||||
|
err := s.withManagementConnection(ctx, target, true, func(database Database) (err error) {
|
||||||
|
oid, err = database.CreateOwnedDatabase(ctx, name, role, roleOID)
|
||||||
|
return err
|
||||||
|
})
|
||||||
|
if err != nil {
|
||||||
|
return 0, err
|
||||||
|
}
|
||||||
|
return oid, nil
|
||||||
|
}
|
||||||
|
|
||||||
|
func (s *InstanceService) ConfigureAccess(ctx context.Context, target instance.ObservationTarget, name, role string, state provisioning.State) error {
|
||||||
|
return s.withManagementConnection(ctx, target, false, func(database Database) error {
|
||||||
|
return database.ConfigureAccess(ctx, name, role, state)
|
||||||
|
})
|
||||||
|
}
|
||||||
|
|
||||||
|
// 所有供应操作与观察共用同一个连接登记和 Secret 刷新边界。
|
||||||
|
// 创建后凭据回读失败不能冒充明确未执行;用例必须保留不确定诊断。
|
||||||
|
func (s *InstanceService) withManagementConnection(ctx context.Context, target instance.ObservationTarget, creating bool, operation func(Database) error) error {
|
||||||
|
s.mu.Lock()
|
||||||
|
defer s.mu.Unlock()
|
||||||
|
current, err := s.connection(ctx, target)
|
||||||
|
if err != nil {
|
||||||
|
return ErrResourceUnavailable
|
||||||
|
}
|
||||||
|
metadata, err := current.database.InspectManagement(ctx)
|
||||||
|
checks := metadata.Management
|
||||||
|
if err != nil || checks.Connection != instance.CheckPassed || checks.Roles != instance.CheckPassed || checks.Databases != instance.CheckPassed || checks.Grants != instance.CheckPassed {
|
||||||
|
return ErrResourceUnavailable
|
||||||
|
}
|
||||||
|
if err := s.checkCredentials(ctx, target, current.credentials); err != nil {
|
||||||
|
return ErrResourceUnavailable
|
||||||
|
}
|
||||||
|
if err := operation(current.database); err != nil {
|
||||||
|
return err
|
||||||
|
}
|
||||||
|
if err := s.checkCredentials(ctx, target, current.credentials); err != nil {
|
||||||
|
if creating {
|
||||||
|
return ErrResourceUncertain
|
||||||
|
}
|
||||||
|
return ErrResourceUnavailable
|
||||||
|
}
|
||||||
|
return nil
|
||||||
|
}
|
||||||
@@ -5,7 +5,6 @@ import (
|
|||||||
"context"
|
"context"
|
||||||
"time"
|
"time"
|
||||||
|
|
||||||
"git.ddupan.top/panxiao81/ayatori/internal/database/adapter/kubernetes"
|
|
||||||
"git.ddupan.top/panxiao81/ayatori/internal/database/application"
|
"git.ddupan.top/panxiao81/ayatori/internal/database/application"
|
||||||
ctrl "sigs.k8s.io/controller-runtime"
|
ctrl "sigs.k8s.io/controller-runtime"
|
||||||
"sigs.k8s.io/controller-runtime/pkg/client"
|
"sigs.k8s.io/controller-runtime/pkg/client"
|
||||||
@@ -15,7 +14,16 @@ const dependencyRetry = 30 * time.Second
|
|||||||
|
|
||||||
type BindingReconciler struct {
|
type BindingReconciler struct {
|
||||||
Client client.Client
|
Client client.Client
|
||||||
Reader client.Reader
|
Service *application.BindingService
|
||||||
|
Presenter BindingPresenter
|
||||||
|
}
|
||||||
|
|
||||||
|
type BindingPresenter interface {
|
||||||
|
Present(context.Context, application.BindingResult) error
|
||||||
|
}
|
||||||
|
|
||||||
|
func NewBindingReconciler(cache client.Client, service *application.BindingService, presenter BindingPresenter) *BindingReconciler {
|
||||||
|
return &BindingReconciler{Client: cache, Service: service, Presenter: presenter}
|
||||||
}
|
}
|
||||||
|
|
||||||
// +kubebuilder:rbac:groups=database.ayatori.ddupan.top,resources=postgresqltenants,verbs=get;list;watch;update;patch
|
// +kubebuilder:rbac:groups=database.ayatori.ddupan.top,resources=postgresqltenants,verbs=get;list;watch;update;patch
|
||||||
@@ -27,13 +35,11 @@ type BindingReconciler struct {
|
|||||||
// +kubebuilder:rbac:groups=database.ayatori.ddupan.top,resources=postgresqlinstances,verbs=get;list;watch
|
// +kubebuilder:rbac:groups=database.ayatori.ddupan.top,resources=postgresqlinstances,verbs=get;list;watch
|
||||||
|
|
||||||
func (r *BindingReconciler) Reconcile(ctx context.Context, request ctrl.Request) (ctrl.Result, error) {
|
func (r *BindingReconciler) Reconcile(ctx context.Context, request ctrl.Request) (ctrl.Result, error) {
|
||||||
resources := &kubernetes.BindingResources{Client: r.Client, Reader: r.Reader}
|
result, err := r.Service.Reconcile(ctx, request.Namespace, request.Name)
|
||||||
service := application.BindingService{Resources: resources}
|
|
||||||
result, err := service.Reconcile(ctx, request.Namespace, request.Name)
|
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return ctrl.Result{}, err
|
return ctrl.Result{}, err
|
||||||
}
|
}
|
||||||
if err := resources.Present(ctx, result); err != nil {
|
if err := r.Presenter.Present(ctx, result); err != nil {
|
||||||
return ctrl.Result{}, err
|
return ctrl.Result{}, err
|
||||||
}
|
}
|
||||||
if result.RetrySoon {
|
if result.RetrySoon {
|
||||||
|
|||||||
@@ -42,6 +42,11 @@ func targetDatabaseName(tenant *databasev1alpha1.PostgreSQLTenant) string {
|
|||||||
return "tenant-" + string(tenant.UID)
|
return "tenant-" + string(tenant.UID)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
func bindingTestReconciler(writer client.Client, reader client.Reader) *BindingReconciler {
|
||||||
|
resources := &kubernetes.BindingResources{Client: writer, Reader: reader}
|
||||||
|
return NewBindingReconciler(writer, &application.BindingService{Resources: resources}, resources)
|
||||||
|
}
|
||||||
|
|
||||||
func tenantReference(tenant *databasev1alpha1.PostgreSQLTenant) *databasev1alpha1.TenantReference {
|
func tenantReference(tenant *databasev1alpha1.PostgreSQLTenant) *databasev1alpha1.TenantReference {
|
||||||
return &databasev1alpha1.TenantReference{
|
return &databasev1alpha1.TenantReference{
|
||||||
Namespace: tenant.Namespace, Name: databasev1alpha1.ObjectName(tenant.Name), UID: tenant.UID,
|
Namespace: tenant.Namespace, Name: databasev1alpha1.ObjectName(tenant.Name), UID: tenant.UID,
|
||||||
@@ -103,7 +108,7 @@ func testDynamicBinding(t *testing.T, apiClient client.Client) {
|
|||||||
instance := readyInstance(t, apiClient, "dynamic-instance")
|
instance := readyInstance(t, apiClient, "dynamic-instance")
|
||||||
tenant := provisionTenant("dynamic", instance.Name)
|
tenant := provisionTenant("dynamic", instance.Name)
|
||||||
requireCreate(t, apiClient, tenant)
|
requireCreate(t, apiClient, tenant)
|
||||||
reconciler := &BindingReconciler{Client: apiClient, Reader: apiClient}
|
reconciler := bindingTestReconciler(apiClient, apiClient)
|
||||||
reconcileOK(t, reconciler, tenant)
|
reconcileOK(t, reconciler, tenant)
|
||||||
reload(t, apiClient, tenant)
|
reload(t, apiClient, tenant)
|
||||||
if tenant.Status.DatabaseRef == nil || tenant.Status.Phase != phaseBound {
|
if tenant.Status.DatabaseRef == nil || tenant.Status.Phase != phaseBound {
|
||||||
@@ -154,7 +159,7 @@ func testBindingRestart(t *testing.T, apiClient client.Client) {
|
|||||||
instance := readyInstance(t, apiClient, "restart-instance")
|
instance := readyInstance(t, apiClient, "restart-instance")
|
||||||
tenant := provisionTenant("restart", instance.Name)
|
tenant := provisionTenant("restart", instance.Name)
|
||||||
requireCreate(t, apiClient, tenant)
|
requireCreate(t, apiClient, tenant)
|
||||||
first := &BindingReconciler{Client: &failedTenantStatusClient{Client: apiClient}, Reader: apiClient}
|
first := bindingTestReconciler(&failedTenantStatusClient{Client: apiClient}, apiClient)
|
||||||
if _, err := first.Reconcile(t.Context(), ctrl.Request{NamespacedName: client.ObjectKeyFromObject(tenant)}); err == nil {
|
if _, err := first.Reconcile(t.Context(), ctrl.Request{NamespacedName: client.ObjectKeyFromObject(tenant)}); err == nil {
|
||||||
t.Fatal("预期第二次绑定写入失败")
|
t.Fatal("预期第二次绑定写入失败")
|
||||||
}
|
}
|
||||||
@@ -169,7 +174,7 @@ func testBindingRestart(t *testing.T, apiClient client.Client) {
|
|||||||
t.Fatal("失败后资源侧绑定不应回滚")
|
t.Fatal("失败后资源侧绑定不应回滚")
|
||||||
}
|
}
|
||||||
// 新建 reconciler,无旧内存,只从 API 中读取进度。
|
// 新建 reconciler,无旧内存,只从 API 中读取进度。
|
||||||
restarted := &BindingReconciler{Client: apiClient, Reader: apiClient}
|
restarted := bindingTestReconciler(apiClient, apiClient)
|
||||||
reconcileOK(t, restarted, tenant)
|
reconcileOK(t, restarted, tenant)
|
||||||
reload(t, apiClient, tenant)
|
reload(t, apiClient, tenant)
|
||||||
if tenant.Status.DatabaseRef == nil || tenant.Status.DatabaseRef.UID != database.UID {
|
if tenant.Status.DatabaseRef == nil || tenant.Status.DatabaseRef.UID != database.UID {
|
||||||
@@ -190,7 +195,7 @@ func testConcurrentBinding(t *testing.T, apiClient client.Client) {
|
|||||||
results := make(chan error, len(tenants))
|
results := make(chan error, len(tenants))
|
||||||
for _, tenant := range tenants {
|
for _, tenant := range tenants {
|
||||||
workers.Go(func() {
|
workers.Go(func() {
|
||||||
reconciler := &BindingReconciler{Client: apiClient, Reader: apiClient}
|
reconciler := bindingTestReconciler(apiClient, apiClient)
|
||||||
_, err := reconciler.Reconcile(t.Context(), ctrl.Request{NamespacedName: client.ObjectKeyFromObject(tenant)})
|
_, err := reconciler.Reconcile(t.Context(), ctrl.Request{NamespacedName: client.ObjectKeyFromObject(tenant)})
|
||||||
results <- err
|
results <- err
|
||||||
})
|
})
|
||||||
@@ -202,7 +207,7 @@ func testConcurrentBinding(t *testing.T, apiClient client.Client) {
|
|||||||
t.Fatalf("并发协调出现非版本冲突错误: %v", err)
|
t.Fatalf("并发协调出现非版本冲突错误: %v", err)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
reconciler := &BindingReconciler{Client: apiClient, Reader: apiClient}
|
reconciler := bindingTestReconciler(apiClient, apiClient)
|
||||||
bound := 0
|
bound := 0
|
||||||
for _, tenant := range tenants {
|
for _, tenant := range tenants {
|
||||||
reconcileOK(t, reconciler, tenant)
|
reconcileOK(t, reconciler, tenant)
|
||||||
@@ -233,7 +238,7 @@ func testBindingIdentity(t *testing.T, apiClient client.Client) {
|
|||||||
}
|
}
|
||||||
tenant := existingTenant("identity", database.Name)
|
tenant := existingTenant("identity", database.Name)
|
||||||
requireCreate(t, apiClient, tenant)
|
requireCreate(t, apiClient, tenant)
|
||||||
reconciler := &BindingReconciler{Client: apiClient, Reader: apiClient}
|
reconciler := bindingTestReconciler(apiClient, apiClient)
|
||||||
reconcileOK(t, reconciler, tenant)
|
reconcileOK(t, reconciler, tenant)
|
||||||
reload(t, apiClient, tenant)
|
reload(t, apiClient, tenant)
|
||||||
assertNotReady(t, tenant, reasonConflict)
|
assertNotReady(t, tenant, reasonConflict)
|
||||||
@@ -256,7 +261,7 @@ func testBindingIdentity(t *testing.T, apiClient client.Client) {
|
|||||||
func testBindingProtection(t *testing.T, apiClient client.Client) {
|
func testBindingProtection(t *testing.T, apiClient client.Client) {
|
||||||
tenant := provisionTenant("protection", "missing-instance")
|
tenant := provisionTenant("protection", "missing-instance")
|
||||||
requireCreate(t, apiClient, tenant)
|
requireCreate(t, apiClient, tenant)
|
||||||
reconciler := &BindingReconciler{Client: apiClient, Reader: apiClient}
|
reconciler := bindingTestReconciler(apiClient, apiClient)
|
||||||
reconcileOK(t, reconciler, tenant)
|
reconcileOK(t, reconciler, tenant)
|
||||||
reload(t, apiClient, tenant)
|
reload(t, apiClient, tenant)
|
||||||
assertNotReady(t, tenant, reasonDependency)
|
assertNotReady(t, tenant, reasonDependency)
|
||||||
@@ -292,7 +297,7 @@ func testStaleObservation(t *testing.T, apiClient client.Client) {
|
|||||||
}
|
}
|
||||||
tenant := existingTenant("stale", database.Name)
|
tenant := existingTenant("stale", database.Name)
|
||||||
requireCreate(t, apiClient, tenant)
|
requireCreate(t, apiClient, tenant)
|
||||||
reconciler := &BindingReconciler{Client: apiClient, Reader: apiClient}
|
reconciler := bindingTestReconciler(apiClient, apiClient)
|
||||||
reconcileOK(t, reconciler, tenant)
|
reconcileOK(t, reconciler, tenant)
|
||||||
reload(t, apiClient, tenant)
|
reload(t, apiClient, tenant)
|
||||||
assertNotReady(t, tenant, reasonDependency)
|
assertNotReady(t, tenant, reasonDependency)
|
||||||
@@ -318,7 +323,7 @@ func testBindingWatch(t *testing.T, apiClient client.Client, config *rest.Config
|
|||||||
if err != nil {
|
if err != nil {
|
||||||
t.Fatal(err)
|
t.Fatal(err)
|
||||||
}
|
}
|
||||||
reconciler := &BindingReconciler{}
|
reconciler := bindingTestReconciler(manager.GetClient(), manager.GetAPIReader())
|
||||||
if err := reconciler.SetupWithManager(t.Context(), manager); err != nil {
|
if err := reconciler.SetupWithManager(t.Context(), manager); err != nil {
|
||||||
t.Fatal(err)
|
t.Fatal(err)
|
||||||
}
|
}
|
||||||
@@ -373,7 +378,7 @@ func testPresentationVersion(t *testing.T, apiClient client.Client) {
|
|||||||
if err := resources.Present(t.Context(), result); !apierrors.IsConflict(err) {
|
if err := resources.Present(t.Context(), result); !apierrors.IsConflict(err) {
|
||||||
t.Fatalf("过期结果呈现 = %v, want Conflict", err)
|
t.Fatalf("过期结果呈现 = %v, want Conflict", err)
|
||||||
}
|
}
|
||||||
reconcileOK(t, &BindingReconciler{Client: apiClient, Reader: apiClient}, tenant)
|
reconcileOK(t, bindingTestReconciler(apiClient, apiClient), tenant)
|
||||||
reload(t, apiClient, tenant)
|
reload(t, apiClient, tenant)
|
||||||
if tenant.Status.Phase != phaseBound || tenant.Spec.SecretName != "updated-delivery" ||
|
if tenant.Status.Phase != phaseBound || tenant.Spec.SecretName != "updated-delivery" ||
|
||||||
tenant.Annotations["example.test/keep"] != "preserved" {
|
tenant.Annotations["example.test/keep"] != "preserved" {
|
||||||
|
|||||||
@@ -2,9 +2,10 @@ package controller
|
|||||||
|
|
||||||
import (
|
import (
|
||||||
"context"
|
"context"
|
||||||
|
"errors"
|
||||||
|
|
||||||
databasev1alpha1 "git.ddupan.top/panxiao81/ayatori/api/database/v1alpha1"
|
databasev1alpha1 "git.ddupan.top/panxiao81/ayatori/api/database/v1alpha1"
|
||||||
"git.ddupan.top/panxiao81/ayatori/internal/database/adapter/kubernetes"
|
"git.ddupan.top/panxiao81/ayatori/internal/database/domain/binding"
|
||||||
"k8s.io/apimachinery/pkg/types"
|
"k8s.io/apimachinery/pkg/types"
|
||||||
ctrl "sigs.k8s.io/controller-runtime"
|
ctrl "sigs.k8s.io/controller-runtime"
|
||||||
"sigs.k8s.io/controller-runtime/pkg/client"
|
"sigs.k8s.io/controller-runtime/pkg/client"
|
||||||
@@ -14,16 +15,16 @@ import (
|
|||||||
const targetDatabaseIndex = "database.bindingTarget"
|
const targetDatabaseIndex = "database.bindingTarget"
|
||||||
|
|
||||||
func (r *BindingReconciler) SetupWithManager(ctx context.Context, manager ctrl.Manager) error {
|
func (r *BindingReconciler) SetupWithManager(ctx context.Context, manager ctrl.Manager) error {
|
||||||
if r.Client == nil {
|
if r.Client == nil || r.Service == nil || r.Presenter == nil {
|
||||||
r.Client = manager.GetClient()
|
return errors.New("binding controller requires injected client, use case and presenter")
|
||||||
}
|
|
||||||
if r.Reader == nil {
|
|
||||||
r.Reader = manager.GetAPIReader()
|
|
||||||
}
|
}
|
||||||
if err := manager.GetFieldIndexer().IndexField(ctx, &databasev1alpha1.PostgreSQLTenant{},
|
if err := manager.GetFieldIndexer().IndexField(ctx, &databasev1alpha1.PostgreSQLTenant{},
|
||||||
targetDatabaseIndex, func(object client.Object) []string {
|
targetDatabaseIndex, func(object client.Object) []string {
|
||||||
tenant := object.(*databasev1alpha1.PostgreSQLTenant)
|
tenant := object.(*databasev1alpha1.PostgreSQLTenant)
|
||||||
return []string{kubernetes.BindingTargetName(tenant)}
|
if tenant.Spec.DatabaseRef != nil {
|
||||||
|
return []string{string(tenant.Spec.DatabaseRef.Name)}
|
||||||
|
}
|
||||||
|
return []string{binding.DynamicDatabaseName(string(tenant.UID))}
|
||||||
}); err != nil {
|
}); err != nil {
|
||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -2,9 +2,9 @@ package controller
|
|||||||
|
|
||||||
import (
|
import (
|
||||||
"context"
|
"context"
|
||||||
|
"errors"
|
||||||
|
|
||||||
databasev1alpha1 "git.ddupan.top/panxiao81/ayatori/api/database/v1alpha1"
|
databasev1alpha1 "git.ddupan.top/panxiao81/ayatori/api/database/v1alpha1"
|
||||||
"git.ddupan.top/panxiao81/ayatori/internal/database/adapter/kubernetes"
|
|
||||||
"git.ddupan.top/panxiao81/ayatori/internal/database/application"
|
"git.ddupan.top/panxiao81/ayatori/internal/database/application"
|
||||||
ctrl "sigs.k8s.io/controller-runtime"
|
ctrl "sigs.k8s.io/controller-runtime"
|
||||||
"sigs.k8s.io/controller-runtime/pkg/client"
|
"sigs.k8s.io/controller-runtime/pkg/client"
|
||||||
@@ -14,15 +14,15 @@ import (
|
|||||||
// CredentialReconciler 只连接事件、用例和重试,不在控制器中编排凭据写入。
|
// CredentialReconciler 只连接事件、用例和重试,不在控制器中编排凭据写入。
|
||||||
type CredentialReconciler struct {
|
type CredentialReconciler struct {
|
||||||
Client client.Client
|
Client client.Client
|
||||||
Reader client.Reader
|
Service *application.CredentialPreparation
|
||||||
Store application.CredentialStore
|
}
|
||||||
|
|
||||||
|
func NewCredentialReconciler(cache client.Client, service *application.CredentialPreparation) *CredentialReconciler {
|
||||||
|
return &CredentialReconciler{Client: cache, Service: service}
|
||||||
}
|
}
|
||||||
|
|
||||||
func (r *CredentialReconciler) Reconcile(ctx context.Context, request ctrl.Request) (ctrl.Result, error) {
|
func (r *CredentialReconciler) Reconcile(ctx context.Context, request ctrl.Request) (ctrl.Result, error) {
|
||||||
service := application.CredentialPreparation{
|
if err := r.Service.Reconcile(ctx, request.Name); err != nil {
|
||||||
Resources: &kubernetes.CredentialResources{Client: r.Client, Reader: r.Reader}, Store: r.Store,
|
|
||||||
}
|
|
||||||
if err := service.Reconcile(ctx, request.Name); err != nil {
|
|
||||||
return ctrl.Result{}, err
|
return ctrl.Result{}, err
|
||||||
}
|
}
|
||||||
// Bao 的可用性和版本变化没有 Kubernetes watch;与已有依赖重查保持一致。
|
// Bao 的可用性和版本变化没有 Kubernetes watch;与已有依赖重查保持一致。
|
||||||
@@ -30,20 +30,19 @@ func (r *CredentialReconciler) Reconcile(ctx context.Context, request ctrl.Reque
|
|||||||
}
|
}
|
||||||
|
|
||||||
func (r *CredentialReconciler) SetupWithManager(manager ctrl.Manager) error {
|
func (r *CredentialReconciler) SetupWithManager(manager ctrl.Manager) error {
|
||||||
if r.Client == nil {
|
if r.Client == nil || r.Service == nil {
|
||||||
r.Client = manager.GetClient()
|
return errors.New("credential controller requires injected client and use case")
|
||||||
}
|
|
||||||
if r.Reader == nil {
|
|
||||||
r.Reader = manager.GetAPIReader()
|
|
||||||
}
|
}
|
||||||
return ctrl.NewControllerManagedBy(manager).
|
return ctrl.NewControllerManagedBy(manager).
|
||||||
Named("database-credentials").For(&databasev1alpha1.PostgreSQLDatabase{}).
|
Named("database-credentials").For(&databasev1alpha1.PostgreSQLDatabase{}).
|
||||||
Watches(&databasev1alpha1.PostgreSQLTenant{}, handler.EnqueueRequestsFromMapFunc(r.requestsForTenant)).
|
Watches(&databasev1alpha1.PostgreSQLTenant{}, handler.EnqueueRequestsFromMapFunc(databaseRequestsForTenant)).
|
||||||
Watches(&databasev1alpha1.PostgreSQLInstance{}, handler.EnqueueRequestsFromMapFunc(r.requestsForInstance)).
|
Watches(&databasev1alpha1.PostgreSQLInstance{}, handler.EnqueueRequestsFromMapFunc(func(ctx context.Context, object client.Object) []ctrl.Request {
|
||||||
|
return databaseRequestsForInstance(ctx, r.Client, object)
|
||||||
|
})).
|
||||||
Complete(r)
|
Complete(r)
|
||||||
}
|
}
|
||||||
|
|
||||||
func (r *CredentialReconciler) requestsForTenant(_ context.Context, object client.Object) []ctrl.Request {
|
func databaseRequestsForTenant(_ context.Context, object client.Object) []ctrl.Request {
|
||||||
tenant := object.(*databasev1alpha1.PostgreSQLTenant)
|
tenant := object.(*databasev1alpha1.PostgreSQLTenant)
|
||||||
if tenant.Status.DatabaseRef == nil {
|
if tenant.Status.DatabaseRef == nil {
|
||||||
return nil
|
return nil
|
||||||
@@ -51,9 +50,9 @@ func (r *CredentialReconciler) requestsForTenant(_ context.Context, object clien
|
|||||||
return []ctrl.Request{{Name: string(tenant.Status.DatabaseRef.Name)}}
|
return []ctrl.Request{{Name: string(tenant.Status.DatabaseRef.Name)}}
|
||||||
}
|
}
|
||||||
|
|
||||||
func (r *CredentialReconciler) requestsForInstance(ctx context.Context, object client.Object) []ctrl.Request {
|
func databaseRequestsForInstance(ctx context.Context, cache client.Client, object client.Object) []ctrl.Request {
|
||||||
databases := &databasev1alpha1.PostgreSQLDatabaseList{}
|
databases := &databasev1alpha1.PostgreSQLDatabaseList{}
|
||||||
if err := r.Client.List(ctx, databases); err != nil {
|
if err := cache.List(ctx, databases); err != nil {
|
||||||
ctrl.LoggerFrom(ctx).Error(err, "无法映射 Instance 凭据准备事件;等待低频重试")
|
ctrl.LoggerFrom(ctx).Error(err, "无法映射 Instance 凭据准备事件;等待低频重试")
|
||||||
return nil
|
return nil
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -4,7 +4,6 @@ import (
|
|||||||
"context"
|
"context"
|
||||||
"time"
|
"time"
|
||||||
|
|
||||||
"git.ddupan.top/panxiao81/ayatori/internal/database/adapter/kubernetes"
|
|
||||||
"git.ddupan.top/panxiao81/ayatori/internal/database/application"
|
"git.ddupan.top/panxiao81/ayatori/internal/database/application"
|
||||||
ctrl "sigs.k8s.io/controller-runtime"
|
ctrl "sigs.k8s.io/controller-runtime"
|
||||||
"sigs.k8s.io/controller-runtime/pkg/client"
|
"sigs.k8s.io/controller-runtime/pkg/client"
|
||||||
@@ -12,27 +11,33 @@ import (
|
|||||||
|
|
||||||
type InstanceReconciler struct {
|
type InstanceReconciler struct {
|
||||||
Client client.Client
|
Client client.Client
|
||||||
Reader client.Reader
|
Service *application.InstanceReconciliation
|
||||||
Observer application.InstanceObserver
|
Presenter InstancePresenter
|
||||||
SecretNamespace string
|
SecretNamespace string
|
||||||
}
|
}
|
||||||
|
|
||||||
|
type InstancePresenter interface {
|
||||||
|
PresentInstance(context.Context, application.InstanceResult) error
|
||||||
|
}
|
||||||
|
|
||||||
|
func NewInstanceReconciler(cache client.Client, service *application.InstanceReconciliation, presenter InstancePresenter, namespace string) *InstanceReconciler {
|
||||||
|
return &InstanceReconciler{Client: cache, Service: service, Presenter: presenter, SecretNamespace: namespace}
|
||||||
|
}
|
||||||
|
|
||||||
// +kubebuilder:rbac:groups=database.ayatori.ddupan.top,resources=postgresqlinstances,verbs=get;list;watch;update;patch
|
// +kubebuilder:rbac:groups=database.ayatori.ddupan.top,resources=postgresqlinstances,verbs=get;list;watch;update;patch
|
||||||
// +kubebuilder:rbac:groups=database.ayatori.ddupan.top,resources=postgresqlinstances/status,verbs=get;update;patch
|
// +kubebuilder:rbac:groups=database.ayatori.ddupan.top,resources=postgresqlinstances/status,verbs=get;update;patch
|
||||||
// +kubebuilder:rbac:groups=database.ayatori.ddupan.top,resources=postgresqlinstances/finalizers,verbs=update
|
// +kubebuilder:rbac:groups=database.ayatori.ddupan.top,resources=postgresqlinstances/finalizers,verbs=update
|
||||||
// Secret 权限单独声明为 namespace Role,不放入生成的 ClusterRole。
|
// Secret 权限单独声明为 namespace Role,不放入生成的 ClusterRole。
|
||||||
|
|
||||||
func (r *InstanceReconciler) Reconcile(ctx context.Context, request ctrl.Request) (ctrl.Result, error) {
|
func (r *InstanceReconciler) Reconcile(ctx context.Context, request ctrl.Request) (ctrl.Result, error) {
|
||||||
resources := &kubernetes.InstanceResources{Client: r.Client, Reader: r.Reader}
|
|
||||||
service := application.InstanceReconciliation{Resources: resources, Observer: r.Observer}
|
|
||||||
observationContext, cancel := context.WithTimeout(ctx, 15*time.Second)
|
observationContext, cancel := context.WithTimeout(ctx, 15*time.Second)
|
||||||
defer cancel()
|
defer cancel()
|
||||||
result, err := service.Reconcile(observationContext, request.Name)
|
result, err := r.Service.Reconcile(observationContext, request.Name)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return ctrl.Result{}, err
|
return ctrl.Result{}, err
|
||||||
}
|
}
|
||||||
// 查询超时后仍用 worker context 保存安全失败结果;manager 停止时不强行写入。
|
// 查询超时后仍用 worker context 保存安全失败结果;manager 停止时不强行写入。
|
||||||
if err := resources.PresentInstance(ctx, result); err != nil {
|
if err := r.Presenter.PresentInstance(ctx, result); err != nil {
|
||||||
return ctrl.Result{}, err
|
return ctrl.Result{}, err
|
||||||
}
|
}
|
||||||
if result.Record == nil || result.RemoveProtection {
|
if result.Record == nil || result.RemoveProtection {
|
||||||
|
|||||||
@@ -18,6 +18,7 @@ import (
|
|||||||
)
|
)
|
||||||
|
|
||||||
type instanceBackend struct {
|
type instanceBackend struct {
|
||||||
|
application.ResourceDatabase
|
||||||
checks instance.ManagementChecks
|
checks instance.ManagementChecks
|
||||||
err error
|
err error
|
||||||
inspect func()
|
inspect func()
|
||||||
@@ -52,7 +53,8 @@ func newInstanceReconciler(t *testing.T, apiClient client.Client, backend *insta
|
|||||||
t.Fatal(err)
|
t.Fatal(err)
|
||||||
}
|
}
|
||||||
t.Cleanup(service.Close)
|
t.Cleanup(service.Close)
|
||||||
return &InstanceReconciler{Client: apiClient, Reader: apiClient, Observer: service}
|
resources := &kubernetes.InstanceResources{Client: apiClient, Reader: apiClient}
|
||||||
|
return NewInstanceReconciler(apiClient, &application.InstanceReconciliation{Resources: resources, Observer: service}, resources, "")
|
||||||
}
|
}
|
||||||
|
|
||||||
func reconcileInstance(t *testing.T, reconciler *InstanceReconciler, object *databasev1alpha1.PostgreSQLInstance) {
|
func reconcileInstance(t *testing.T, reconciler *InstanceReconciler, object *databasev1alpha1.PostgreSQLInstance) {
|
||||||
@@ -164,11 +166,11 @@ func TestInstanceDeletionProtection(t *testing.T) {
|
|||||||
t.Fatal(err)
|
t.Fatal(err)
|
||||||
}
|
}
|
||||||
backend.inspect = func() { t.Fatal("删除中不应连接 PostgreSQL") }
|
backend.inspect = func() { t.Fatal("删除中不应连接 PostgreSQL") }
|
||||||
reconciler.Reader = &failedReferenceReader{Reader: apiClient}
|
reconciler.Service.Resources.(*kubernetes.InstanceResources).Reader = &failedReferenceReader{Reader: apiClient}
|
||||||
reconcileInstance(t, reconciler, object)
|
reconcileInstance(t, reconciler, object)
|
||||||
reload(t, apiClient, object)
|
reload(t, apiClient, object)
|
||||||
assertInstanceReason(t, object, reasonDependency)
|
assertInstanceReason(t, object, reasonDependency)
|
||||||
reconciler.Reader = apiClient
|
reconciler.Service.Resources.(*kubernetes.InstanceResources).Reader = apiClient
|
||||||
reconcileInstance(t, reconciler, object)
|
reconcileInstance(t, reconciler, object)
|
||||||
reload(t, apiClient, object)
|
reload(t, apiClient, object)
|
||||||
assertInstanceReason(t, object, "InstanceInUse")
|
assertInstanceReason(t, object, "InstanceInUse")
|
||||||
|
|||||||
@@ -22,15 +22,9 @@ func InstanceCacheOptions(namespace string) cache.Options {
|
|||||||
}
|
}
|
||||||
|
|
||||||
func (r *InstanceReconciler) SetupWithManager(manager ctrl.Manager) error {
|
func (r *InstanceReconciler) SetupWithManager(manager ctrl.Manager) error {
|
||||||
if r.Observer == nil || len(validation.IsDNS1123Label(r.SecretNamespace)) != 0 {
|
if r.Client == nil || r.Service == nil || r.Presenter == nil || len(validation.IsDNS1123Label(r.SecretNamespace)) != 0 {
|
||||||
return errors.New("instance observer and valid management Secret namespace required")
|
return errors.New("instance observer and valid management Secret namespace required")
|
||||||
}
|
}
|
||||||
if r.Client == nil {
|
|
||||||
r.Client = manager.GetClient()
|
|
||||||
}
|
|
||||||
if r.Reader == nil {
|
|
||||||
r.Reader = manager.GetAPIReader()
|
|
||||||
}
|
|
||||||
return ctrl.NewControllerManagedBy(manager).
|
return ctrl.NewControllerManagedBy(manager).
|
||||||
Named("database-instance").
|
Named("database-instance").
|
||||||
For(&databasev1alpha1.PostgreSQLInstance{}).
|
For(&databasev1alpha1.PostgreSQLInstance{}).
|
||||||
|
|||||||
@@ -0,0 +1,43 @@
|
|||||||
|
package controller
|
||||||
|
|
||||||
|
import (
|
||||||
|
"context"
|
||||||
|
"errors"
|
||||||
|
"time"
|
||||||
|
|
||||||
|
databasev1alpha1 "git.ddupan.top/panxiao81/ayatori/api/database/v1alpha1"
|
||||||
|
"git.ddupan.top/panxiao81/ayatori/internal/database/application"
|
||||||
|
ctrl "sigs.k8s.io/controller-runtime"
|
||||||
|
"sigs.k8s.io/controller-runtime/pkg/client"
|
||||||
|
"sigs.k8s.io/controller-runtime/pkg/handler"
|
||||||
|
)
|
||||||
|
|
||||||
|
type ProvisioningReconciler struct {
|
||||||
|
Client client.Client
|
||||||
|
Service *application.DatabaseReconciliation
|
||||||
|
}
|
||||||
|
|
||||||
|
func NewProvisioningReconciler(cache client.Client, service *application.DatabaseReconciliation) *ProvisioningReconciler {
|
||||||
|
return &ProvisioningReconciler{Client: cache, Service: service}
|
||||||
|
}
|
||||||
|
|
||||||
|
func (r *ProvisioningReconciler) Reconcile(ctx context.Context, request ctrl.Request) (ctrl.Result, error) {
|
||||||
|
operationContext, cancel := context.WithTimeout(ctx, 30*time.Second)
|
||||||
|
defer cancel()
|
||||||
|
if err := r.Service.Reconcile(operationContext, request.Name); err != nil {
|
||||||
|
return ctrl.Result{}, err
|
||||||
|
}
|
||||||
|
return ctrl.Result{RequeueAfter: dependencyRetry}, nil
|
||||||
|
}
|
||||||
|
|
||||||
|
func (r *ProvisioningReconciler) SetupWithManager(manager ctrl.Manager) error {
|
||||||
|
if r.Client == nil || r.Service == nil {
|
||||||
|
return errors.New("provisioning controller requires injected client and use case")
|
||||||
|
}
|
||||||
|
return ctrl.NewControllerManagedBy(manager).Named("database-provisioning").
|
||||||
|
For(&databasev1alpha1.PostgreSQLDatabase{}).
|
||||||
|
Watches(&databasev1alpha1.PostgreSQLTenant{}, handler.EnqueueRequestsFromMapFunc(databaseRequestsForTenant)).
|
||||||
|
Watches(&databasev1alpha1.PostgreSQLInstance{}, handler.EnqueueRequestsFromMapFunc(func(ctx context.Context, object client.Object) []ctrl.Request {
|
||||||
|
return databaseRequestsForInstance(ctx, r.Client, object)
|
||||||
|
})).Complete(r)
|
||||||
|
}
|
||||||
+1
-1
@@ -14,7 +14,7 @@ See the License for the specific language governing permissions and
|
|||||||
limitations under the License.
|
limitations under the License.
|
||||||
*/
|
*/
|
||||||
|
|
||||||
package application
|
package credential
|
||||||
|
|
||||||
import (
|
import (
|
||||||
"crypto/rand"
|
"crypto/rand"
|
||||||
+9
-8
@@ -14,7 +14,7 @@ See the License for the specific language governing permissions and
|
|||||||
limitations under the License.
|
limitations under the License.
|
||||||
*/
|
*/
|
||||||
|
|
||||||
package application_test
|
package credential_test
|
||||||
|
|
||||||
import (
|
import (
|
||||||
"encoding/json"
|
"encoding/json"
|
||||||
@@ -23,7 +23,8 @@ import (
|
|||||||
"strings"
|
"strings"
|
||||||
"testing"
|
"testing"
|
||||||
|
|
||||||
"git.ddupan.top/panxiao81/ayatori/internal/database/application"
|
credentialdomain "git.ddupan.top/panxiao81/ayatori/internal/database/domain/credential"
|
||||||
|
|
||||||
"git.ddupan.top/panxiao81/ayatori/internal/database/domain/instance"
|
"git.ddupan.top/panxiao81/ayatori/internal/database/domain/instance"
|
||||||
)
|
)
|
||||||
|
|
||||||
@@ -35,11 +36,11 @@ func TestApplicationCredential(t *testing.T) {
|
|||||||
if err != nil {
|
if err != nil {
|
||||||
t.Fatal(err)
|
t.Fatal(err)
|
||||||
}
|
}
|
||||||
first, err := application.GenerateApplicationCredential("owner", "app", endpoint)
|
first, err := credentialdomain.GenerateApplicationCredential("owner", "app", endpoint)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
t.Fatal(err)
|
t.Fatal(err)
|
||||||
}
|
}
|
||||||
second, err := application.GenerateApplicationCredential("owner", "app", endpoint)
|
second, err := credentialdomain.GenerateApplicationCredential("owner", "app", endpoint)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
t.Fatal(err)
|
t.Fatal(err)
|
||||||
}
|
}
|
||||||
@@ -47,7 +48,7 @@ func TestApplicationCredential(t *testing.T) {
|
|||||||
if len(data) != 7 || data["password"] == second.SecretData()["password"] || len(data["password"].(string)) != 43 {
|
if len(data) != 7 || data["password"] == second.SecretData()["password"] || len(data["password"].(string)) != 43 {
|
||||||
t.Fatal("expected seven keys and independent 256-bit passwords")
|
t.Fatal("expected seven keys and independent 256-bit passwords")
|
||||||
}
|
}
|
||||||
parsed, err := application.ParseApplicationCredential(data)
|
parsed, err := credentialdomain.ParseApplicationCredential(data)
|
||||||
if err != nil || !maps.Equal(parsed.SecretData(), data) {
|
if err != nil || !maps.Equal(parsed.SecretData(), data) {
|
||||||
t.Fatal("credential did not round trip")
|
t.Fatal("credential did not round trip")
|
||||||
}
|
}
|
||||||
@@ -67,15 +68,15 @@ func TestApplicationCredential(t *testing.T) {
|
|||||||
for key := range data {
|
for key := range data {
|
||||||
invalid := maps.Clone(data)
|
invalid := maps.Clone(data)
|
||||||
delete(invalid, key)
|
delete(invalid, key)
|
||||||
if _, err := application.ParseApplicationCredential(invalid); err == nil {
|
if _, err := credentialdomain.ParseApplicationCredential(invalid); err == nil {
|
||||||
t.Fatalf("accepted missing %s", key)
|
t.Fatalf("accepted missing %s", key)
|
||||||
}
|
}
|
||||||
invalid[key] = 42
|
invalid[key] = 42
|
||||||
if _, err := application.ParseApplicationCredential(invalid); err == nil {
|
if _, err := credentialdomain.ParseApplicationCredential(invalid); err == nil {
|
||||||
t.Fatalf("accepted non-string %s", key)
|
t.Fatalf("accepted non-string %s", key)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
if (application.ApplicationCredential{}).Validate() == nil {
|
if (credentialdomain.ApplicationCredential{}).Validate() == nil {
|
||||||
t.Fatal("accepted zero credential")
|
t.Fatal("accepted zero credential")
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
@@ -0,0 +1,129 @@
|
|||||||
|
package credential
|
||||||
|
|
||||||
|
import (
|
||||||
|
"git.ddupan.top/panxiao81/ayatori/internal/database/domain/binding"
|
||||||
|
"git.ddupan.top/panxiao81/ayatori/internal/database/domain/instance"
|
||||||
|
)
|
||||||
|
|
||||||
|
type Location struct {
|
||||||
|
Mount string
|
||||||
|
Path string
|
||||||
|
}
|
||||||
|
|
||||||
|
// Phase 表达凭据准备进度,不依赖 Kubernetes Condition 的类型或 Reason。
|
||||||
|
type Phase uint8
|
||||||
|
|
||||||
|
const (
|
||||||
|
Pending Phase = iota
|
||||||
|
Pinned
|
||||||
|
Creating
|
||||||
|
Prepared
|
||||||
|
Conflict
|
||||||
|
Unavailable
|
||||||
|
Stopped
|
||||||
|
InvalidTarget
|
||||||
|
)
|
||||||
|
|
||||||
|
type State struct {
|
||||||
|
Location *Location
|
||||||
|
Version int64
|
||||||
|
Phase Phase
|
||||||
|
Message string
|
||||||
|
}
|
||||||
|
|
||||||
|
func (s State) WithPhase(phase Phase, message string) State {
|
||||||
|
s.Phase, s.Message = phase, message
|
||||||
|
return s
|
||||||
|
}
|
||||||
|
|
||||||
|
func (s State) Confirmed() bool { return s.Version > 0 }
|
||||||
|
|
||||||
|
// Created 只接受首次创建并回读得到的版本,不能把后续写入认作首次供应。
|
||||||
|
func (s State) Created(version int64) (State, *Issue) {
|
||||||
|
if version != 1 {
|
||||||
|
return s, &Issue{Conflict, "凭据创建冲突或结果不确定;请核对固定位置的版本历史,未认领、覆盖或重新生成密码"}
|
||||||
|
}
|
||||||
|
s.Version = version
|
||||||
|
return s.WithPhase(Prepared, "凭据已创建并回读确认;尚未创建 PostgreSQL 资源或交付给 Tenant"), nil
|
||||||
|
}
|
||||||
|
|
||||||
|
// Resume 决定新一轮是否可以继续。未确认的创建不能靠读取成功认领。
|
||||||
|
func (s State) Resume() (State, bool) {
|
||||||
|
if s.Version != 0 {
|
||||||
|
return s, true
|
||||||
|
}
|
||||||
|
switch s.Phase {
|
||||||
|
case Conflict:
|
||||||
|
return s, false
|
||||||
|
case Creating:
|
||||||
|
return s.WithPhase(Conflict, "凭据创建未留下成功确认;请核对固定位置与后端历史并人工处理,未重新生成密码"), false
|
||||||
|
default:
|
||||||
|
return s, true
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
func (s State) CheckLocation(configured Location) *Issue {
|
||||||
|
if s.Location != nil && *s.Location != configured {
|
||||||
|
return &Issue{Unavailable, "部署配置与固定凭据位置不一致;请恢复原 mount/path 配置,未迁移或改密"}
|
||||||
|
}
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
|
||||||
|
type Instance struct {
|
||||||
|
binding.Instance
|
||||||
|
Endpoint instance.Endpoint
|
||||||
|
}
|
||||||
|
|
||||||
|
// Target 只包含供应资格所需事实,不含 resourceVersion、Conditions 或 repository 对象。
|
||||||
|
type Target struct {
|
||||||
|
Database binding.Database
|
||||||
|
Tenant *binding.Tenant
|
||||||
|
Instance *Instance
|
||||||
|
DatabaseProtected bool
|
||||||
|
TenantProtected bool
|
||||||
|
}
|
||||||
|
|
||||||
|
type Issue struct {
|
||||||
|
Phase Phase
|
||||||
|
Message string
|
||||||
|
}
|
||||||
|
|
||||||
|
func (t Target) RequiresPreparation() bool { return t.Database.Source == "Provision" }
|
||||||
|
|
||||||
|
func (t Target) CheckCredential(value ApplicationCredential) *Issue {
|
||||||
|
if t.Instance == nil || !value.MatchesTarget(t.Database.LoginRole, t.Database.Name, t.Instance.Endpoint) {
|
||||||
|
return &Issue{Conflict, "已确认凭据与当前 Instance/database/loginRole 不一致;请人工核实,未修改凭据"}
|
||||||
|
}
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
|
||||||
|
func (t Target) Check() *Issue {
|
||||||
|
database := t.Database
|
||||||
|
if database.Deleting || database.Phase == binding.Deleting || database.Phase == "Released" {
|
||||||
|
return &Issue{Stopped, "Database 正在删除或已释放;保留凭据与 finalizer,不执行供应或清理"}
|
||||||
|
}
|
||||||
|
if database.Tenant == nil || t.Tenant == nil || t.Tenant.Database == nil {
|
||||||
|
return &Issue{Unavailable, "等待 Database 与 Tenant 双向绑定完成"}
|
||||||
|
}
|
||||||
|
if *database.Tenant != t.Tenant.Identity || *t.Tenant.Database != database.Identity {
|
||||||
|
return &Issue{Conflict, "双向绑定的名称或 UID 不匹配,未继续供应"}
|
||||||
|
}
|
||||||
|
if t.Tenant.Deleting || t.Tenant.Phase != binding.Bound || !t.DatabaseProtected || !t.TenantProtected {
|
||||||
|
return &Issue{Stopped, "Tenant 未完成绑定、正在删除或缺少 finalizer 保护,未继续供应"}
|
||||||
|
}
|
||||||
|
request, err := t.Tenant.Request.Resolve(t.Tenant.Identity)
|
||||||
|
if err != nil || (request.Provision != nil && !database.MatchesProvision(request, t.Tenant.Identity)) || request.Name != database.Identity.Name {
|
||||||
|
return &Issue{Conflict, "Tenant 申请与 Database 目标不一致,未继续供应"}
|
||||||
|
}
|
||||||
|
if t.Instance == nil || database.InstanceUID == "" {
|
||||||
|
return &Issue{Unavailable, "等待 Instance 与已记录的实例身份"}
|
||||||
|
}
|
||||||
|
if issue := t.Instance.Check(&database); issue != nil {
|
||||||
|
phase := Unavailable
|
||||||
|
if issue.Reason == binding.Conflict {
|
||||||
|
phase = Conflict
|
||||||
|
}
|
||||||
|
return &Issue{phase, issue.Message}
|
||||||
|
}
|
||||||
|
return nil
|
||||||
|
}
|
||||||
@@ -0,0 +1,59 @@
|
|||||||
|
package credential_test
|
||||||
|
|
||||||
|
import (
|
||||||
|
"testing"
|
||||||
|
|
||||||
|
credential "git.ddupan.top/panxiao81/ayatori/internal/database/domain/credential"
|
||||||
|
)
|
||||||
|
|
||||||
|
func TestPreparationResume(t *testing.T) {
|
||||||
|
for _, test := range []struct {
|
||||||
|
name string
|
||||||
|
phase credential.Phase
|
||||||
|
version int64
|
||||||
|
continueAllowed bool
|
||||||
|
result credential.Phase
|
||||||
|
}{
|
||||||
|
{"尚未创建", credential.Pinned, 0, true, credential.Pinned},
|
||||||
|
{"依赖恢复", credential.Unavailable, 0, true, credential.Unavailable},
|
||||||
|
{"中断创建", credential.Creating, 0, false, credential.Conflict},
|
||||||
|
{"未确认冲突", credential.Conflict, 0, false, credential.Conflict},
|
||||||
|
{"已确认后读取失败", credential.Unavailable, 1, true, credential.Unavailable},
|
||||||
|
{"已确认后冲突重验", credential.Conflict, 1, true, credential.Conflict},
|
||||||
|
} {
|
||||||
|
t.Run(test.name, func(t *testing.T) {
|
||||||
|
original := credential.State{Phase: test.phase, Version: test.version, Message: "保留原诊断"}
|
||||||
|
state, allowed := original.Resume()
|
||||||
|
if allowed != test.continueAllowed || state.Phase != test.result || state.Version != original.Version {
|
||||||
|
t.Fatal("恢复判定或确认版本发生变化")
|
||||||
|
}
|
||||||
|
if test.phase == credential.Conflict && state.Message != original.Message {
|
||||||
|
t.Fatal("冲突重入应保留原诊断")
|
||||||
|
}
|
||||||
|
})
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
func TestPreparationLocationAndConfirmation(t *testing.T) {
|
||||||
|
location := credential.Location{Mount: "applications", Path: "database/uid"}
|
||||||
|
state := credential.State{Location: &location, Phase: credential.Creating}
|
||||||
|
if issue := state.CheckLocation(location); issue != nil {
|
||||||
|
t.Fatal("固定位置不应被拒绝")
|
||||||
|
}
|
||||||
|
if issue := state.CheckLocation(credential.Location{Mount: "other", Path: location.Path}); issue == nil || issue.Phase != credential.Unavailable {
|
||||||
|
t.Fatal("配置变化必须停止,不迁移已固定位置")
|
||||||
|
}
|
||||||
|
for _, version := range []int64{0, -1, 2} {
|
||||||
|
result, issue := state.Created(version)
|
||||||
|
if issue == nil || issue.Phase != credential.Conflict || result.Confirmed() {
|
||||||
|
t.Fatal("错误版本不得确认创建")
|
||||||
|
}
|
||||||
|
}
|
||||||
|
result, issue := state.Created(1)
|
||||||
|
if issue != nil || !result.Confirmed() || result.Phase != credential.Prepared || result.Location != state.Location {
|
||||||
|
t.Fatal("首次写入回读应确认并保留位置")
|
||||||
|
}
|
||||||
|
if state.Version != 0 {
|
||||||
|
t.Fatal("领域判定不得修改调用方的旧状态")
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -0,0 +1,52 @@
|
|||||||
|
package credential_test
|
||||||
|
|
||||||
|
import (
|
||||||
|
"testing"
|
||||||
|
|
||||||
|
"git.ddupan.top/panxiao81/ayatori/internal/database/domain/binding"
|
||||||
|
credential "git.ddupan.top/panxiao81/ayatori/internal/database/domain/credential"
|
||||||
|
)
|
||||||
|
|
||||||
|
const (
|
||||||
|
targetTestApplication = "sampleapp"
|
||||||
|
targetTestInstance = "test-instance"
|
||||||
|
)
|
||||||
|
|
||||||
|
func TestPreparationTarget(t *testing.T) {
|
||||||
|
for _, test := range []struct {
|
||||||
|
name string
|
||||||
|
change func(*credential.Target)
|
||||||
|
want credential.Phase
|
||||||
|
}{
|
||||||
|
{"完整绑定", func(*credential.Target) {}, credential.Pending},
|
||||||
|
{"单向绑定", func(target *credential.Target) { target.Tenant.Database = nil }, credential.Unavailable},
|
||||||
|
{"旧租户身份", func(target *credential.Target) { target.Tenant.Identity.UID = "new" }, credential.Conflict},
|
||||||
|
{"旧实例身份", func(target *credential.Target) { target.Instance.Identity.UID = "new" }, credential.Conflict},
|
||||||
|
{"资源删除", func(target *credential.Target) { target.Database.Deleting = true }, credential.Stopped},
|
||||||
|
{"申请删除", func(target *credential.Target) { target.Tenant.Deleting = true }, credential.Stopped},
|
||||||
|
{"Released", func(target *credential.Target) { target.Database.Phase = "Released" }, credential.Stopped},
|
||||||
|
{"缺少保护", func(target *credential.Target) { target.DatabaseProtected = false }, credential.Stopped},
|
||||||
|
{"实例未就绪", func(target *credential.Target) { target.Instance.Ready = false }, credential.Unavailable},
|
||||||
|
{"目标变化", func(target *credential.Target) { target.Database.LoginRole = "other" }, credential.Conflict},
|
||||||
|
} {
|
||||||
|
t.Run(test.name, func(t *testing.T) {
|
||||||
|
tenantID := binding.TenantIdentity{Namespace: "apps", Name: targetTestApplication, UID: "tenant"}
|
||||||
|
databaseID := binding.Identity{Name: binding.DynamicDatabaseName(tenantID.UID), UID: "database"}
|
||||||
|
target := credential.Target{
|
||||||
|
Database: binding.Database{Identity: databaseID, Tenant: &tenantID, Instance: targetTestInstance, InstanceUID: "instance-id", Name: targetTestApplication, LoginRole: targetTestApplication, Source: "Provision"},
|
||||||
|
Tenant: &binding.Tenant{Identity: tenantID, Database: &databaseID, Phase: binding.Bound, Request: binding.Request{Provision: &binding.ProvisionRequest{Instance: targetTestInstance}}},
|
||||||
|
Instance: &credential.Instance{Identity: binding.Identity{Name: targetTestInstance, UID: "instance-id"}, Ready: true},
|
||||||
|
DatabaseProtected: true, TenantProtected: true,
|
||||||
|
}
|
||||||
|
test.change(&target)
|
||||||
|
issue := target.Check()
|
||||||
|
if test.want == credential.Pending {
|
||||||
|
if issue != nil {
|
||||||
|
t.Fatalf("有效绑定被拒绝: %s", issue.Message)
|
||||||
|
}
|
||||||
|
} else if issue == nil || issue.Phase != test.want {
|
||||||
|
t.Fatal("领域资格判定不符")
|
||||||
|
}
|
||||||
|
})
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -0,0 +1,66 @@
|
|||||||
|
// Package provisioning 保存 PostgreSQL 资源供应的确认与恢复规则,不依赖后端或 API。
|
||||||
|
package provisioning
|
||||||
|
|
||||||
|
import "fmt"
|
||||||
|
|
||||||
|
type Phase string
|
||||||
|
|
||||||
|
const (
|
||||||
|
Pending Phase = "Pending"
|
||||||
|
CreatingRole Phase = "CreatingRole"
|
||||||
|
CreatingDatabase Phase = "CreatingDatabase"
|
||||||
|
Available Phase = "Available"
|
||||||
|
Conflict Phase = "Conflict"
|
||||||
|
Unavailable Phase = "Unavailable"
|
||||||
|
Stopped Phase = "Stopped"
|
||||||
|
)
|
||||||
|
|
||||||
|
// OID 是已成功创建并回读的对象身份,不是从名称推导出的管理授权。
|
||||||
|
// 仅保存在 CR;不建立 PostgreSQL registry,也不承诺备份还原后的自动认领。
|
||||||
|
type State struct {
|
||||||
|
RoleOID uint32
|
||||||
|
DatabaseOID uint32
|
||||||
|
Phase Phase
|
||||||
|
Message string
|
||||||
|
}
|
||||||
|
|
||||||
|
type Observation struct {
|
||||||
|
RoleOID uint32
|
||||||
|
DatabaseOID uint32
|
||||||
|
RoleSafe bool
|
||||||
|
OwnerOID uint32
|
||||||
|
PublicConnect bool
|
||||||
|
AllowConnections bool
|
||||||
|
}
|
||||||
|
|
||||||
|
func (s State) WithPhase(phase Phase, message string) State {
|
||||||
|
s.Phase, s.Message = phase, message
|
||||||
|
return s
|
||||||
|
}
|
||||||
|
|
||||||
|
func (s State) Resume() (State, bool) {
|
||||||
|
if s.Phase == Conflict {
|
||||||
|
return s, false
|
||||||
|
}
|
||||||
|
if (s.Phase == CreatingRole && s.RoleOID == 0) || (s.Phase == CreatingDatabase && s.DatabaseOID == 0) {
|
||||||
|
return s.WithPhase(Conflict, "外部创建未留下成功确认;请核对目标角色和数据库,不自动认领或重复创建"), false
|
||||||
|
}
|
||||||
|
return s, true
|
||||||
|
}
|
||||||
|
|
||||||
|
// Check 既阻止未知同名对象,也拒绝已确认对象消失、被重建或权限漂移。
|
||||||
|
func (s State) Check(o Observation) error {
|
||||||
|
if s.RoleOID != o.RoleOID {
|
||||||
|
return fmt.Errorf("角色身份不匹配:记录 OID=%d,观察 OID=%d", s.RoleOID, o.RoleOID)
|
||||||
|
}
|
||||||
|
if s.DatabaseOID != o.DatabaseOID {
|
||||||
|
return fmt.Errorf("数据库身份不匹配:记录 OID=%d,观察 OID=%d", s.DatabaseOID, o.DatabaseOID)
|
||||||
|
}
|
||||||
|
if o.RoleOID != 0 && !o.RoleSafe {
|
||||||
|
return fmt.Errorf("已确认角色的登录属性、特权或成员关系发生变化")
|
||||||
|
}
|
||||||
|
if o.DatabaseOID != 0 && o.OwnerOID != s.RoleOID {
|
||||||
|
return fmt.Errorf("数据库 owner 与已确认角色不匹配")
|
||||||
|
}
|
||||||
|
return nil
|
||||||
|
}
|
||||||
@@ -0,0 +1,52 @@
|
|||||||
|
package provisioning
|
||||||
|
|
||||||
|
import "testing"
|
||||||
|
|
||||||
|
func TestResume(t *testing.T) {
|
||||||
|
for _, test := range []struct {
|
||||||
|
name string
|
||||||
|
state State
|
||||||
|
resume bool
|
||||||
|
}{
|
||||||
|
{"初次", State{}, true},
|
||||||
|
{"角色创建不确定", State{Phase: CreatingRole}, false},
|
||||||
|
{"库创建不确定", State{RoleOID: 11, Phase: CreatingDatabase}, false},
|
||||||
|
{"已确认角色", State{RoleOID: 11, Phase: Pending}, true},
|
||||||
|
{"已确认资源暂不可用", State{RoleOID: 11, DatabaseOID: 22, Phase: Unavailable}, true},
|
||||||
|
{"冲突保持", State{RoleOID: 11, Phase: Conflict, Message: "首次诊断"}, false},
|
||||||
|
} {
|
||||||
|
t.Run(test.name, func(t *testing.T) {
|
||||||
|
state, resume := test.state.Resume()
|
||||||
|
if resume != test.resume || state.RoleOID != test.state.RoleOID || state.DatabaseOID != test.state.DatabaseOID {
|
||||||
|
t.Fatal("恢复决策丢失确认或错误放行")
|
||||||
|
}
|
||||||
|
if !resume && state.Phase != Conflict {
|
||||||
|
t.Fatal("不确定创建未转冲突")
|
||||||
|
}
|
||||||
|
if test.state.Phase == Conflict && state != test.state {
|
||||||
|
t.Fatal("冲突诊断被覆盖")
|
||||||
|
}
|
||||||
|
})
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
func TestObservedIdentity(t *testing.T) {
|
||||||
|
state := State{RoleOID: 11, DatabaseOID: 22}
|
||||||
|
good := Observation{RoleOID: 11, DatabaseOID: 22, RoleSafe: true, OwnerOID: 11}
|
||||||
|
if err := state.Check(good); err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
for _, observed := range []Observation{
|
||||||
|
{}, {RoleOID: 12, DatabaseOID: 22, RoleSafe: true, OwnerOID: 12},
|
||||||
|
{RoleOID: 11, DatabaseOID: 23, RoleSafe: true, OwnerOID: 11},
|
||||||
|
{RoleOID: 11, DatabaseOID: 22, RoleSafe: false, OwnerOID: 11},
|
||||||
|
{RoleOID: 11, DatabaseOID: 22, RoleSafe: true, OwnerOID: 99},
|
||||||
|
} {
|
||||||
|
if state.Check(observed) == nil {
|
||||||
|
t.Fatal("对象身份或权限漂移被接受")
|
||||||
|
}
|
||||||
|
}
|
||||||
|
if (State{}).Check(good) == nil {
|
||||||
|
t.Fatal("未确认同名对象被认领")
|
||||||
|
}
|
||||||
|
}
|
||||||
Reference in New Issue
Block a user