Compare commits
5
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
2a4f622f44
|
||
|
|
f6bb9e4599
|
||
|
|
f4deb98a7f | ||
|
|
bc227bfdb4
|
||
|
|
55b269ce2e |
@@ -68,7 +68,7 @@ lint: golangci-lint ## Run golangci-lint linter
|
||||
"$(GOLANGCI_LINT)" run
|
||||
|
||||
.PHONY: test-database-integration
|
||||
test-database-integration: setup-envtest ## 使用临时 API server 与独立 PostgreSQL 容器验证凭据读取和连接更新。
|
||||
test-database-integration: setup-envtest ## 使用临时 API server、PostgreSQL 与 OpenBao 容器验证 Database 后端。
|
||||
KUBEBUILDER_ASSETS="$(shell "$(ENVTEST)" use $(ENVTEST_K8S_VERSION) --bin-dir "$(LOCALBIN)" -p path)" go test -tags=integration -race -count=1 ./internal/database/...
|
||||
|
||||
.PHONY: lint-database-integration
|
||||
|
||||
@@ -0,0 +1,26 @@
|
||||
package main
|
||||
|
||||
import (
|
||||
"git.ddupan.top/panxiao81/ayatori/internal/database/adapter/kubernetes"
|
||||
"git.ddupan.top/panxiao81/ayatori/internal/database/adapter/postgresql"
|
||||
"git.ddupan.top/panxiao81/ayatori/internal/database/application"
|
||||
databasecontroller "git.ddupan.top/panxiao81/ayatori/internal/database/controller"
|
||||
ctrl "sigs.k8s.io/controller-runtime"
|
||||
)
|
||||
|
||||
func setupInstanceObservation(manager ctrl.Manager, namespace, rootCert string) (*application.InstanceService, error) {
|
||||
credentials, err := kubernetes.NewSecretCredentials(manager.GetConfig(), namespace)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
service, err := application.NewInstanceService(credentials, postgresql.Connector{RootCert: rootCert})
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
reconciler := &databasecontroller.InstanceReconciler{Observer: service, SecretNamespace: namespace}
|
||||
if err := reconciler.SetupWithManager(manager); err != nil {
|
||||
service.Close()
|
||||
return nil, err
|
||||
}
|
||||
return service, nil
|
||||
}
|
||||
+25
-3
@@ -22,6 +22,7 @@ import (
|
||||
|
||||
databasev1alpha1 "git.ddupan.top/panxiao81/ayatori/api/database/v1alpha1"
|
||||
executionv1alpha1 "git.ddupan.top/panxiao81/ayatori/api/execution/v1alpha1"
|
||||
"git.ddupan.top/panxiao81/ayatori/internal/database/application"
|
||||
databasecontroller "git.ddupan.top/panxiao81/ayatori/internal/database/controller"
|
||||
// +kubebuilder:scaffold:imports
|
||||
)
|
||||
@@ -41,6 +42,10 @@ func init() {
|
||||
|
||||
// nolint:gocyclo
|
||||
func main() {
|
||||
var databaseNamespace, databaseRootCert string
|
||||
flag.StringVar(&databaseNamespace, "database-secret-namespace", os.Getenv("POD_NAMESPACE"),
|
||||
"固定管理 Secret namespace;为空时不启用 Instance 观测")
|
||||
flag.StringVar(&databaseRootCert, "database-root-cert", "", "PostgreSQL 管理连接信任的公开 CA bundle 路径")
|
||||
var metricsAddr string
|
||||
var metricsCertPath, metricsCertName, metricsCertKey string
|
||||
var webhookCertPath, webhookCertName, webhookCertKey string
|
||||
@@ -145,7 +150,7 @@ func main() {
|
||||
metricsServerOptions.KeyName = metricsCertKey
|
||||
}
|
||||
|
||||
mgr, err := ctrl.NewManager(ctrl.GetConfigOrDie(), ctrl.Options{
|
||||
managerOptions := ctrl.Options{
|
||||
Scheme: scheme,
|
||||
Metrics: metricsServerOptions,
|
||||
WebhookServer: webhookServer,
|
||||
@@ -163,13 +168,25 @@ func main() {
|
||||
// if you are doing or is intended to do any operation such as perform cleanups
|
||||
// after the manager stops then its usage might be unsafe.
|
||||
// LeaderElectionReleaseOnCancel: true,
|
||||
})
|
||||
}
|
||||
if databaseNamespace != "" {
|
||||
managerOptions.Cache = databasecontroller.InstanceCacheOptions(databaseNamespace)
|
||||
}
|
||||
mgr, err := ctrl.NewManager(ctrl.GetConfigOrDie(), managerOptions)
|
||||
if err != nil {
|
||||
setupLog.Error(err, "Failed to start manager")
|
||||
os.Exit(1)
|
||||
}
|
||||
|
||||
// +kubebuilder:scaffold:builder
|
||||
var instanceService *application.InstanceService
|
||||
if databaseNamespace != "" {
|
||||
instanceService, err = setupInstanceObservation(mgr, databaseNamespace, databaseRootCert)
|
||||
if err != nil {
|
||||
setupLog.Error(err, "Failed to set up Instance observation")
|
||||
os.Exit(1)
|
||||
}
|
||||
}
|
||||
if err := (&databasecontroller.BindingReconciler{}).SetupWithManager(context.Background(), mgr); err != nil {
|
||||
setupLog.Error(err, "Failed to set up Database binding controller")
|
||||
os.Exit(1)
|
||||
@@ -185,7 +202,12 @@ func main() {
|
||||
}
|
||||
|
||||
setupLog.Info("Starting manager")
|
||||
if err := mgr.Start(ctrl.SetupSignalHandler()); err != nil {
|
||||
err = mgr.Start(ctrl.SetupSignalHandler())
|
||||
// worker 完全停止后才释放 pgxpool,避免与在途观察竞争。
|
||||
if instanceService != nil {
|
||||
instanceService.Close()
|
||||
}
|
||||
if err != nil {
|
||||
setupLog.Error(err, "Failed to run manager")
|
||||
os.Exit(1)
|
||||
}
|
||||
|
||||
@@ -65,6 +65,11 @@ spec:
|
||||
- --health-probe-bind-address=:8081
|
||||
image: controller:latest
|
||||
name: manager
|
||||
env:
|
||||
- name: POD_NAMESPACE
|
||||
valueFrom:
|
||||
fieldRef:
|
||||
fieldPath: metadata.namespace
|
||||
ports:
|
||||
- containerPort: 8081
|
||||
name: health
|
||||
|
||||
@@ -0,0 +1,23 @@
|
||||
apiVersion: rbac.authorization.k8s.io/v1
|
||||
kind: Role
|
||||
metadata:
|
||||
name: database-management-credentials
|
||||
namespace: system
|
||||
rules:
|
||||
- apiGroups: [""]
|
||||
resources: [secrets]
|
||||
verbs: [get, list, watch]
|
||||
---
|
||||
apiVersion: rbac.authorization.k8s.io/v1
|
||||
kind: RoleBinding
|
||||
metadata:
|
||||
name: database-management-credentials
|
||||
namespace: system
|
||||
roleRef:
|
||||
apiGroup: rbac.authorization.k8s.io
|
||||
kind: Role
|
||||
name: database-management-credentials
|
||||
subjects:
|
||||
- kind: ServiceAccount
|
||||
name: controller-manager
|
||||
namespace: system
|
||||
@@ -7,6 +7,7 @@ resources:
|
||||
- service_account.yaml
|
||||
- role.yaml
|
||||
- role_binding.yaml
|
||||
- database_credentials_role.yaml
|
||||
- leader_election_role.yaml
|
||||
- leader_election_role_binding.yaml
|
||||
# The following RBAC configurations are used to protect
|
||||
|
||||
@@ -19,6 +19,7 @@ rules:
|
||||
- database.ayatori.ddupan.top
|
||||
resources:
|
||||
- postgresqldatabases/finalizers
|
||||
- postgresqlinstances/finalizers
|
||||
- postgresqltenants/finalizers
|
||||
verbs:
|
||||
- update
|
||||
@@ -26,6 +27,7 @@ rules:
|
||||
- database.ayatori.ddupan.top
|
||||
resources:
|
||||
- postgresqldatabases/status
|
||||
- postgresqlinstances/status
|
||||
- postgresqltenants/status
|
||||
verbs:
|
||||
- get
|
||||
@@ -35,13 +37,6 @@ rules:
|
||||
- database.ayatori.ddupan.top
|
||||
resources:
|
||||
- postgresqlinstances
|
||||
verbs:
|
||||
- get
|
||||
- list
|
||||
- watch
|
||||
- apiGroups:
|
||||
- database.ayatori.ddupan.top
|
||||
resources:
|
||||
- postgresqltenants
|
||||
verbs:
|
||||
- get
|
||||
|
||||
+75
-11
@@ -1,7 +1,7 @@
|
||||
# Database 模块
|
||||
|
||||
Database 是 Ayatori 首批实际产品领域之一。当前已包含 Instance 领域基础、管理凭据连接与
|
||||
metadata 观察切片,尚未完成 Database API/controller 和 Tenant 供应链路。
|
||||
Database 是 Ayatori 首批实际产品领域之一。当前已包含三资源 API、分层绑定与 Instance 原生
|
||||
管理能力观测;尚未完成 Database 供应/导入、Tenant 凭据交付与资源回收链路。
|
||||
|
||||
## 当前设计(2026-09-24)
|
||||
|
||||
@@ -12,7 +12,7 @@ Retain 后人工重新绑定与资源侧 Delete。撤销 PostgreSQL ownership re
|
||||
依据 [ADR-0009](../decisions/0009-database-resource-and-claim.md),当前合同见
|
||||
[系统规格](specification.md)。下面的迁移来源与已存在代码不反向约束新设计。
|
||||
registry adapter、专属迁移/测试及 Instance 的 registry 判定现已撤除;Instance 根据完整管理
|
||||
能力观察直接判定 Ready。新增 Database 资源、绑定、导入、角色/凭据管理边界与回收链路尚未实现。wiki 同步位置见
|
||||
能力观察直接判定 Ready。Database 资源与绑定已接入,导入、角色/凭据供应及回收仍未完成。wiki 同步位置见
|
||||
`homelab-wiki/services/postgresql-tenant-operator.md`,跨仓库发布状态由 wiki 的同步记录维护。
|
||||
|
||||
## 来源基线
|
||||
@@ -59,10 +59,10 @@ Secret metadata 和无关字段变化不重建连接。观测后再次读取 Sec
|
||||
不把旧连接的成功作为新凭据有效的证据;这不构成跨 Kubernetes/PostgreSQL 的原子事务。
|
||||
Instance UID、endpoint 或凭据引用变化也会释放旧连接;Forget/Close 只释放本地资源。
|
||||
|
||||
当前通过 `ObserveMetadata` 读取服务器版本和可用扩展,`ObserveVersion` 只是其版本读取便捷入口,
|
||||
不能产生完整 CapabilityObservation 或 Ready。
|
||||
controller 接入、Secret watch、finalizer 与真实权限检查仍待后续切片;并发 CR 更新
|
||||
必须由调用者通过 resourceVersion 校验。应用层沿用源实现的串行处理,本阶段未引入新的调度框架。
|
||||
`ObserveMetadata` 读取服务器版本和可用扩展,`ObserveVersion` 是版本读取便捷入口;两者
|
||||
不会填充管理检查,不能产生 Ready。`ObserveManagement` 使用同一凭据/连接边界读取完整
|
||||
原生管理检查,返回绑定当前 target 的 `InstanceObservation`。controller 的资源呈现适配器
|
||||
使用 resourceVersion 拒绝过期写入。应用层沿用串行处理,不引入新的调度框架。
|
||||
|
||||
运行 `make test-database-integration` 验证真实 API server + 一次性 PostgreSQL;fixture 不接受外部
|
||||
DSN,镜像固定摘要,使用随机本机回环端口并在退出时删除测试容器。覆盖缺失/错误凭据、RBAC、
|
||||
@@ -83,14 +83,13 @@ SQL adapter 通过一条只读语句读取 `pg_catalog.current_setting('server_v
|
||||
|
||||
`InstanceService` 保留原有 CredentialReader → Connector → Database 边界,复用同一个
|
||||
凭据读取、连接刷新和串行释放流程,不新增连接池封装或任意查询回调。每次重新查询 metadata,
|
||||
并在 Secret 有效值回读一致后生成不可变的 `MetadataObservation`,绑定本次 target(含当前
|
||||
并在 Secret 有效值回读一致后生成不可变的 `InstanceObservation`,绑定本次 target(含当前
|
||||
generation),不绑定建池时的旧 target。结果不包含凭据,扩展集合不与 driver 的可变 slice 共享。
|
||||
凭据中途变化、读取失败或查询失败时,返回零值观察并释放连接,不复用旧的扩展列表。
|
||||
|
||||
调用方可将 `Target()` 与 `Extensions()` 交给 Instance 的 `ObserveExtensions`;应用调用链
|
||||
仍负责同轮次使用,不能持久化或跨轮缓存这份证据。metadata 读取不安装扩展、不初始化 registry、
|
||||
不设置 Ready,也不授予 Tenant 写权限。管理权限矩阵及 controller 的
|
||||
checkpoint/status/finalizer 链路仍是后续切片。
|
||||
不设置 Ready,也不授予 Tenant 写权限。管理权限检查使用下面的独立入口。
|
||||
|
||||
真实 API server + PostgreSQL 测试验证未安装扩展可被观察、名称保持大小写、search_path 遮蔽
|
||||
不改变查询来源、低权限账号读取、权限撤回失败与恢复、Secret 中途变化丢弃扩展结果。
|
||||
@@ -106,7 +105,72 @@ Instance 不再具有 InitializingRegistry 阶段、RegistryState、准备决策
|
||||
|
||||
保留凭据读取与连接刷新、TLS、metadata/扩展观察及其真实后端测试。领域测试覆盖每项能力
|
||||
在初次验证和 Ready 重验时失败、依赖恢复、重启后重新取证、错误目标/阶段及删除保护。
|
||||
这不等于管理权限探测矩阵或三资源 controller 已实现。
|
||||
这不等于三资源供应、交付和回收已实现。
|
||||
|
||||
## Instance 原生管理观测
|
||||
|
||||
2026-09-25 维护者确认先使用原生非 superuser 方案,不引入 SECURITY DEFINER 接口。
|
||||
`InspectManagement` 通过同一条只读语句读取当前执行角色的 `CREATEROLE`、`CREATEDB`、
|
||||
superuser 属性、服务器可写状态、版本和扩展列表。不创建探针数据库/角色,不初始化 schema。
|
||||
只有非 superuser、具备两项原生属性且当前服务器/会话可写时,基础管理能力才通过。
|
||||
角色属性不可从继承成员关系推导;具体已有资源仍须检查 owner、membership 与授权范围。
|
||||
|
||||
权限依据和真实测试对应:
|
||||
|
||||
- role:当前角色具有 CREATEROLE,可创建普通登录角色;
|
||||
- database:当前角色具有 CREATEDB,且会话非只读、服务器不在 recovery;
|
||||
- grant:使用自己新建角色的管理权限,显式建立 SET membership,再以 owner 管理数据库 ACL;
|
||||
- extension:新建数据库 owner 可安装 trusted 扩展;可用列表不是安装授权,非 trusted 或
|
||||
其他前提不满足的扩展仍可能失败,必须逐请求执行和回读。导入不继承此动态供应授权。
|
||||
|
||||
参考 PostgreSQL 官方 [CREATE ROLE](https://www.postgresql.org/docs/18/sql-createrole.html)、
|
||||
[CREATE DATABASE](https://www.postgresql.org/docs/18/sql-createdatabase.html) 与
|
||||
[CREATE EXTENSION](https://www.postgresql.org/docs/18/sql-createextension.html)。这些检查是基础
|
||||
能力观察,不是未来操作必然成功的保证;权限、容量、连接数等仍可能在执行时变化。
|
||||
|
||||
`InstanceReconciliation` 协调 finalizer、观察、领域判定和删除引用检查,controller 只连接
|
||||
事件、用例、状态呈现与重试。每轮重建无证据的领域对象;旧 Ready 不授权新一轮操作。
|
||||
失败清除当前版本结果并撤销 Ready;CR 在观察期间被修改则拒绝旧结果,下一轮重新读取。
|
||||
连接/权限变化由 Secret watch 和 30 秒重查驱动,单轮 IO 最长 15 秒;不持续写入相同状态。
|
||||
|
||||
manager 通过 `--database-secret-namespace`(默认 `POD_NAMESPACE`)启用 Instance 观测;
|
||||
为空时不启用。本地运行需显式提供该参数。Deployment 使用 downward API 获取自身 namespace;
|
||||
Secret 的 get/list/watch 权限由该 namespace 的 Role 单独授予,不放入 ClusterRole。
|
||||
watch 使用 controller-runtime 的 metadata-only cache,读取有效凭据仍直连 API server。
|
||||
TLS 使用既有 endpoint 合同,公开 CA bundle 可由 `--database-root-cert` 指定;不会自动挂载
|
||||
生产证书或创建管理 Secret。manager worker 停止后统一关闭 pgxpool。
|
||||
|
||||
Instance 删除首先释放本地连接并撤销 Ready。任何引用它的 Database(含 Released、删除中)
|
||||
或动态 Tenant 申请都会阻止 finalizer 解除;列表查询失败也等待。仅在引用全部解除后移除
|
||||
`database.ayatori.ddupan.top/instance-protection`,不删除 PostgreSQL、账号或凭据。
|
||||
引用查询与删除不是跨对象事务;后续供应仍必须拒绝已删除/删除中的 Instance。
|
||||
|
||||
验收使用真实 PostgreSQL + API server:原生管理账号实际建库、owner 授权、trusted 扩展
|
||||
安装/回读,拒绝非 trusted 扩展;权限撤回/恢复、只读会话、superuser 拒绝和中途轮换。
|
||||
实际 manager 在生成的资源 RBAC 和 namespaced Secret Role 下验证缺失 Secret 后出现、
|
||||
轮换、删除、跨 namespace 拒绝和 watch。API 测试另覆盖写入版本冲突、幂等、新 reconciler
|
||||
恢复与引用删除保护。完整 DBaaS 仍需供应/导入、OpenBao/ESO、Retain/Delete 集成验收。
|
||||
|
||||
## 应用凭据存储切片
|
||||
|
||||
`adapter/openbao` 使用官方 Go SDK `api/v2 v2.7.0` 的 KV v2 API,只有创建和读取,
|
||||
不维护 registry、不覆盖已有密码。动态路径由固定前缀与 Database UID 组成;所有访问都校验
|
||||
配置前缀,已有导入位置也不能绕过 controller 的凭据权限范围。
|
||||
|
||||
创建使用 CAS=0,随后回读七键和版本 1;已有值或软删除历史报冲突。关闭 SDK 自动重试,
|
||||
写入响应丢失、回读失败或内容变化均返回不确定结果,上层不得生成第二份密码或自动认领。
|
||||
`Read` 只适用于调用方已确认关联的路径,读取成功本身不是管理权证据。错误不传播 SDK
|
||||
响应体;内存凭据的普通格式化及 JSON 输出均脱敏,明确的 `SecretData` 才返回明文七键。
|
||||
|
||||
依据官方 [KV v2 CAS 合同](https://github.com/openbao/openbao/blob/main/internal/builtin/logical/kv/path_data.go)
|
||||
与 [Go SDK](https://github.com/openbao/openbao/tree/main/api)。`make test-database-integration`
|
||||
现包含独立 OpenBao dev 容器,固定摘要、随机回环端口、无持久卷,不接受外部地址。
|
||||
真实后端覆盖创建/回读、并发唯一创建、重建适配器读取、软删除冲突、固定前缀 token
|
||||
拒绝管理路径,以及成功写入后丢失响应;HTTP 故障测试补充不重试和错误脱敏。
|
||||
|
||||
这一切片尚未接入 manager:Kubernetes auth/token 生命周期、Database 状态中的稳定位置和
|
||||
已确认步骤、供应 service/controller、PostgreSQL 创建以及 ESO 交付仍未完成。
|
||||
测试 token 只用于临时 fixture,不是生产静态 token 配置接口。现有绑定不会触发外部写入。
|
||||
|
||||
## 设计入口
|
||||
|
||||
|
||||
@@ -53,7 +53,9 @@ Tenant 进入 `status.phase=Binding` 后由 CEL 固定申请目标;Database
|
||||
绑定顺序由 application service 协调,纯资格规则在领域层;Kubernetes adapter 负责快照
|
||||
映射、finalizer 和状态呈现。呈现前若资源版本已变化,返回冲突供下一轮重读,不覆盖其他修改。
|
||||
双向记录完成后 Tenant 为 Bound,Ready=False/BindingComplete,明确尚未供应或交付。
|
||||
生成的 manager RBAC 仅授予绑定所需资源读写,不包含 Secret 读取或后端凭据权限。
|
||||
生成的 manager ClusterRole 授予资源读写,不包含 Secret 读取;Instance 观测的管理 Secret
|
||||
权限由固定 namespace 的独立 Role 授予。Instance controller 已接入原生管理观察与引用删除
|
||||
保护,启用方式及 Ready 边界见 [模块说明](README.md#instance-原生管理观测)。
|
||||
|
||||
当前有 Tenant/Database finalizer 保护,但**删除清理尚未实现**:Tenant 删除报告
|
||||
Ready=False/DeletionPending 并保留绑定与 finalizer,Database 的保护也不会被自动移除。
|
||||
|
||||
+18
-11
@@ -1,16 +1,16 @@
|
||||
# 部署与配置
|
||||
|
||||
> 本页迁入作为 Database 模块的目标部署合同。Ayatori manager flags、manifests 与发布装配尚未
|
||||
> 实现;当前行为以修订后的系统规格为准,本页不能直接用于部署。
|
||||
> 本页区分已实现的 Instance 观测配置与尚未接入的供应/交付目标合同。
|
||||
> 完整 Database 服务仍不可部署使用;当前可执行入口见 [模块说明](README.md)。
|
||||
|
||||
| 项目 | 内容 |
|
||||
| --- | --- |
|
||||
| 状态 | Review |
|
||||
| 环境 | homelab Kubernetes + 外部 PostgreSQL/OpenBao |
|
||||
| 最后更新 | 2026-09-24 |
|
||||
| 最后更新 | 2026-09-25 |
|
||||
|
||||
本文定义 v1alpha1 的运行依赖、启动顺序和部署级配置。当前 manifests 尚未实现这些
|
||||
配置,示例是后续实现合同,不可直接用于现有脚手架。
|
||||
本文定义 v1alpha1 的运行依赖、启动顺序和部署级配置。Instance 观测已接入 manager;
|
||||
OpenBao、ESO 与完整供应装配仍是后续实现合同。
|
||||
|
||||
## 依赖与顺序
|
||||
|
||||
@@ -30,7 +30,11 @@
|
||||
|
||||
## Controller 配置合同
|
||||
|
||||
以下是尚待实现的部署配置合同,凭据定位随三资源 API 继续细化。controller 使用这些 CLI flags。
|
||||
当前 manager 支持 `--database-secret-namespace`(默认 `POD_NAMESPACE`,为空则停用
|
||||
Instance 观测)与 `--database-root-cert`(公开 PostgreSQL CA PEM 路径)。Deployment
|
||||
通过 downward API 获取 namespace,Secret 权限由该 namespace 的 Role 授予。
|
||||
|
||||
以下是尚待实现的供应/交付配置合同,不表示当前 manager 接受这些 CLI flags。
|
||||
必填项缺失、路径无效或 duration 不为正数时,进程必须在启动 manager 前失败;
|
||||
不得等到 reconcile 时才逐个资源报告配置错误。
|
||||
|
||||
@@ -44,7 +48,7 @@
|
||||
| `--openbao-service-account-token-path` | `/var/run/secrets/kubernetes.io/serviceaccount/token` | Kubernetes auth 使用的投射 token 文件 |
|
||||
| `--openbao-tenant-base-path` | 默认 `postgresql-tenants` | controller 专属 mount-relative 前缀 |
|
||||
| `--external-secret-store-name` | 必填 | controller 创建的 ExternalSecret 固定引用 |
|
||||
| `--postgresql-ca-bundle-path` | PostgreSQL TLS 模式必填 | 只读 PEM trust bundle,不含私钥 |
|
||||
| `--database-root-cert` | 已实现 | 只读 PEM trust bundle,不含私钥;沿用 Instance 连接配置 |
|
||||
| `--reconcile-timeout` | `30s` | 单轮 reconcile 中外部操作的总期限,必须大于零 |
|
||||
|
||||
address 必须是绝对 `http` 或 `https` URL,不允许 userinfo、query 或 fragment,末尾 `/`
|
||||
@@ -54,7 +58,9 @@ API 层。生产环境的 `--openbao-address` 必须使用 HTTPS;HTTP 只用
|
||||
|
||||
Tenant 不能选择任意凭据路径。凭据必须能随 Database 保留并安全交付给被授权的新 Tenant;
|
||||
原 `<base-path>/<namespace>/<metadata.name>` 定位规则不再直接作为新 API 合同。
|
||||
稳定位置与导入关联方式待 API 评审;consumer URL 仍使用无认证信息的 KV v2 API URL。
|
||||
动态供应位置使用 `<base-path>/<Database UID>`;导入使用 Database 的显式 credentialRef,
|
||||
不要求搬迁已有凭据。供应流程须先记录原 mount/path,不能在配置变化后重新推导位置。
|
||||
consumer URL 仍使用无认证信息的 KV v2 API URL。
|
||||
|
||||
base path 必须是合法 mount-relative path,不以 `/` 开头且不包含空段、`.`、`..`、
|
||||
`data`/`metadata` API 层。ExternalSecret 固定命名为
|
||||
@@ -77,9 +83,10 @@ base path 必须是合法 mount-relative path,不以 `/` 开头且不包含空
|
||||
- `Delete` 时禁止连接、终止目标 database session、删除已验证归属的 database/role。
|
||||
|
||||
部分 PostgreSQL 操作天然要求较高权限,尤其终止其他 session 和安装某些 extension。
|
||||
应优先使用 PostgreSQL 预定义角色、受控 SECURITY DEFINER 管理函数或限定数据库的
|
||||
授权;任何不得不使用 superuser 的 extension 都必须按实例单独记录,不得扩大默认
|
||||
controller 权限。最终可执行 SQL grant 将随 PostgreSQL adapter 集成测试固化。
|
||||
第一版使用原生非 superuser 的 CREATEDB/CREATEROLE 方案,不引入 SECURITY DEFINER
|
||||
管理接口。对自行创建的 owner 显式建立 SET membership,再以 owner 管理 ACL 与扩展;
|
||||
已有对象仍须逐资源核实授权,不能凭基础属性接管。需要 superuser 的扩展不能扩大 controller
|
||||
权限。真实权限矩阵见 [Instance 原生管理观测](README.md#instance-原生管理观测)。
|
||||
|
||||
## OpenBao 与 ESO
|
||||
|
||||
|
||||
@@ -1,8 +1,8 @@
|
||||
# Instance 领域对象规格
|
||||
|
||||
日期:2026-09-24。资源模型修订依据
|
||||
日期:2026-09-25。资源模型修订依据
|
||||
[ADR-0009](../decisions/0009-database-resource-and-claim.md),行为以
|
||||
[系统规格](specification.md) 为准。本页替代原 registry 准备与恢复合同;领域依赖已撤除,完整应用/controller 链路尚未接入。
|
||||
[系统规格](specification.md) 为准。本页替代原 registry 准备与恢复合同;领域依赖已撤除,Instance 应用/controller 观测链路已接入。
|
||||
|
||||
## 职责
|
||||
|
||||
@@ -85,4 +85,6 @@ finalizer 不阻止并发申请 CR 创建;新请求见 Instance 删除中/不
|
||||
|
||||
Instance 领域代码、adapter 与测试的 registry 依赖已撤除。AssessManagement 根据完整观察
|
||||
直接完成验证;AssessReadiness 失败进入 Validating,依赖恢复后重新验证。领域测试覆盖各检查项
|
||||
在这两个入口的失败与恢复,但完整权限探测、Instance controller 和三资源生命周期尚未完成。
|
||||
在这两个入口的失败与恢复。原生权限检查、Instance controller、metadata-only Secret watch
|
||||
和引用删除保护已接入,验证矩阵见 [模块说明](README.md#instance-原生管理观测);
|
||||
Database/Tenant 的供应、交付和回收仍未完成。
|
||||
|
||||
@@ -14,6 +14,13 @@ Database 的 `database.ayatori.ddupan.top/database-protection` 也尚无清理
|
||||
|
||||
## 日常检查
|
||||
|
||||
Instance 观察已实现:先确认 manager 配置了 `--database-secret-namespace` 或 `POD_NAMESPACE`,
|
||||
再检查 Ready Reason、observedGeneration 与管理 Secret 名称/字段映射,切勿导出其 data。
|
||||
`InsufficientPrivileges` 表示当前原生方案要求的非 superuser、CREATEDB/CREATEROLE 不满足;
|
||||
`CredentialsChanged` 会丢弃中途轮换的结果并重验;`InstanceInUse` 消息定位阻塞删除的资源。
|
||||
Secret 事件立即入队,30 秒重查覆盖 PostgreSQL 权限等没有 Kubernetes 事件的外部变化。
|
||||
Instance 删除不要求 PostgreSQL 可达,但必须可读取所有 Database/Tenant 引用。
|
||||
|
||||
先看 Instance、Database、Tenant 的 Ready Condition、绑定 UID、阶段与 observedGeneration,
|
||||
再核对 PostgreSQL catalog、OpenBao metadata、ExternalSecret 与 Secret 投射状态。
|
||||
具体 kubectl 资源名、finalizer 名称与人工确认字段在 API 实现后补齐,不提供猜测的 patch 命令。
|
||||
|
||||
@@ -3,7 +3,7 @@
|
||||
| 项目 | 内容 |
|
||||
| --- | --- |
|
||||
| 状态 | Review |
|
||||
| 最后更新 | 2026-09-24 |
|
||||
| 最后更新 | 2026-09-25 |
|
||||
|
||||
## 保护目标
|
||||
|
||||
@@ -47,9 +47,11 @@ ESO 身份只读管理路径,租户 ESO 身份只读 tenant base path,二者
|
||||
使用管理凭据 Store。controller 对管理 Secret 的读取限于自身 namespace,Instance
|
||||
不能指定其他 namespace;controller 不创建或修改管理 Secret/ExternalSecret。
|
||||
|
||||
PostgreSQL 管理 role 不应是 superuser。若平台选择 SECURITY DEFINER 函数承载创建或
|
||||
删除操作,函数必须固定 `search_path`、严格校验 identifier、拒绝任意 SQL,并仅向
|
||||
controller role 授予 EXECUTE。controller 不调用 shell 或 `psql` 拼接用户输入。
|
||||
2026-09-25 维护者确认第一版使用原生非 superuser 管理 role,具有 CREATEDB/CREATEROLE,
|
||||
不引入 SECURITY DEFINER 接口。Instance 检查拒绝 superuser;具体已有资源的 owner 和
|
||||
membership 仍需逐资源验证,不能把基础能力用于接管他人资源。扩展按实际权限安装,
|
||||
不因可用列表包含某个扩展就默认能安装它。controller 不调用 shell 或 `psql` 拼接用户输入。
|
||||
当前检查与真实权限矩阵见 [Instance 原生管理观测](README.md#instance-原生管理观测)。
|
||||
|
||||
Kubernetes RBAC 应把 Instance 管理、Database 导入、Released 重新开放和回收限制给平台管理员。
|
||||
有权创建 Tenant 的申请者可显式申请未绑定且可用的 Database,不增加资源侧允许绑定名单
|
||||
|
||||
@@ -171,6 +171,9 @@ status 缺失不假定发生于正常重启。Instance 可重新探测能力;D
|
||||
|
||||
## 9. 权限、凭据与扩展
|
||||
|
||||
2026-09-25 确认第一版管理账号使用原生非 superuser + CREATEDB/CREATEROLE 方案,
|
||||
不引入 SECURITY DEFINER 接口;权限检查与限制见 [安全合同](security.md#最小权限)。
|
||||
|
||||
动态供应继续使用一个兼任 database owner 的 LOGIN role;应用角色不得具备 superuser、
|
||||
CREATEDB、CREATEROLE 或 replication 权限。撤销 PUBLIC CONNECT,再授予目标角色;
|
||||
不修改无关数据库和角色。identifier 匹配 `^[a-z][a-z0-9_]{0,62}$`,SQL 安全引用。
|
||||
|
||||
@@ -4,10 +4,12 @@ go 1.27.1
|
||||
|
||||
require (
|
||||
github.com/jackc/pgx/v5 v5.11.0
|
||||
github.com/openbao/openbao/api/v2 v2.7.0
|
||||
k8s.io/api v0.37.0
|
||||
k8s.io/apimachinery v0.37.0
|
||||
k8s.io/client-go v0.37.0
|
||||
sigs.k8s.io/controller-runtime v0.25.0
|
||||
sigs.k8s.io/yaml v1.6.0
|
||||
)
|
||||
|
||||
require (
|
||||
@@ -24,6 +26,7 @@ require (
|
||||
github.com/felixge/httpsnoop v1.0.4 // indirect
|
||||
github.com/fsnotify/fsnotify v1.9.0 // indirect
|
||||
github.com/fxamacker/cbor/v2 v2.9.1 // indirect
|
||||
github.com/go-jose/go-jose/v4 v4.1.4 // indirect
|
||||
github.com/go-logr/logr v1.4.3 // indirect
|
||||
github.com/go-logr/stdr v1.2.2 // indirect
|
||||
github.com/go-logr/zapr v1.3.0 // indirect
|
||||
@@ -41,15 +44,25 @@ require (
|
||||
github.com/go-openapi/swag/stringutils v0.27.1 // indirect
|
||||
github.com/go-openapi/swag/typeutils v0.27.1 // indirect
|
||||
github.com/go-openapi/swag/yamlutils v0.27.1 // indirect
|
||||
github.com/go-viper/mapstructure/v2 v2.5.0 // indirect
|
||||
github.com/google/cel-go v0.29.2 // indirect
|
||||
github.com/google/gnostic-models v0.7.0 // indirect
|
||||
github.com/google/uuid v1.6.0 // indirect
|
||||
github.com/grpc-ecosystem/grpc-gateway/v2 v2.29.0 // indirect
|
||||
github.com/hashicorp/errwrap v1.1.0 // indirect
|
||||
github.com/hashicorp/go-cleanhttp v0.5.2 // indirect
|
||||
github.com/hashicorp/go-multierror v1.1.1 // indirect
|
||||
github.com/hashicorp/go-retryablehttp v0.7.8 // indirect
|
||||
github.com/hashicorp/go-secure-stdlib/parseutil v0.2.0 // indirect
|
||||
github.com/hashicorp/go-secure-stdlib/strutil v0.1.2 // indirect
|
||||
github.com/hashicorp/go-sockaddr v1.0.7 // indirect
|
||||
github.com/hashicorp/hcl v1.0.1-vault-7 // indirect
|
||||
github.com/inconshreveable/mousetrap v1.1.0 // indirect
|
||||
github.com/jackc/pgpassfile v1.0.0 // indirect
|
||||
github.com/jackc/pgservicefile v0.0.0-20240606120523-5a60cdf6a761 // indirect
|
||||
github.com/jackc/puddle/v2 v2.2.2 // indirect
|
||||
github.com/json-iterator/go v1.1.12 // indirect
|
||||
github.com/mitchellh/mapstructure v1.5.0 // indirect
|
||||
github.com/modern-go/concurrent v0.0.0-20180306012644-bacd9c7ef1dd // indirect
|
||||
github.com/modern-go/reflect2 v1.0.3-0.20250322232337-35a7c28c31ee // indirect
|
||||
github.com/munnerz/goautoneg v0.0.0-20191010083416-a7dc8b61c822 // indirect
|
||||
@@ -58,6 +71,7 @@ require (
|
||||
github.com/prometheus/client_model v0.6.2 // indirect
|
||||
github.com/prometheus/common v0.70.0 // indirect
|
||||
github.com/prometheus/procfs v0.21.1 // indirect
|
||||
github.com/ryanuber/go-glob v1.0.0 // indirect
|
||||
github.com/spf13/cobra v1.10.2 // indirect
|
||||
github.com/spf13/pflag v1.0.10 // indirect
|
||||
github.com/x448/float16 v0.8.4 // indirect
|
||||
@@ -75,7 +89,7 @@ require (
|
||||
go.yaml.in/yaml/v2 v2.4.4 // indirect
|
||||
go.yaml.in/yaml/v3 v3.0.5 // indirect
|
||||
golang.org/x/exp v0.0.0-20260410095643-746e56fc9e2f // indirect
|
||||
golang.org/x/net v0.57.0 // indirect
|
||||
golang.org/x/net v0.58.0 // indirect
|
||||
golang.org/x/oauth2 v0.36.0 // indirect
|
||||
golang.org/x/sync v0.22.0 // indirect
|
||||
golang.org/x/sys v0.47.0 // indirect
|
||||
@@ -100,5 +114,4 @@ require (
|
||||
sigs.k8s.io/json v0.0.0-20250730193827-2d320260d730 // indirect
|
||||
sigs.k8s.io/randfill v1.0.0 // indirect
|
||||
sigs.k8s.io/structured-merge-diff/v6 v6.4.2 // indirect
|
||||
sigs.k8s.io/yaml v1.6.0 // indirect
|
||||
)
|
||||
|
||||
@@ -23,12 +23,16 @@ github.com/evanphx/json-patch v0.5.2 h1:xVCHIVMUu1wtM/VkR9jVZ45N3FhZfYMMYGorLCR8
|
||||
github.com/evanphx/json-patch v0.5.2/go.mod h1:ZWS5hhDbVDyob71nXKNL0+PWn6ToqBHMikGIFbs31qQ=
|
||||
github.com/evanphx/json-patch/v5 v5.9.11 h1:/8HVnzMq13/3x9TPvjG08wUGqBTmZBsCWzjTM0wiaDU=
|
||||
github.com/evanphx/json-patch/v5 v5.9.11/go.mod h1:3j+LviiESTElxA4p3EMKAB9HXj3/XEtnUf6OZxqIQTM=
|
||||
github.com/fatih/color v1.19.0 h1:Zp3PiM21/9Ld6FzSKyL5c/BULoe/ONr9KlbYVOfG8+w=
|
||||
github.com/fatih/color v1.19.0/go.mod h1:zNk67I0ZUT1bEGsSGyCZYZNrHuTkJJB+r6Q9VuMi0LE=
|
||||
github.com/felixge/httpsnoop v1.0.4 h1:NFTV2Zj1bL4mc9sqWACXbQFVBBg2W3GPvqp8/ESS2Wg=
|
||||
github.com/felixge/httpsnoop v1.0.4/go.mod h1:m8KPJKqk1gH5J9DgRY2ASl2lWCfGKXixSwevea8zH2U=
|
||||
github.com/fsnotify/fsnotify v1.9.0 h1:2Ml+OJNzbYCTzsxtv8vKSFD9PbJjmhYF14k/jKC7S9k=
|
||||
github.com/fsnotify/fsnotify v1.9.0/go.mod h1:8jBTzvmWwFyi3Pb8djgCCO5IBqzKJ/Jwo8TRcHyHii0=
|
||||
github.com/fxamacker/cbor/v2 v2.9.1 h1:2rWm8B193Ll4VdjsJY28jxs70IdDsHRWgQYAI80+rMQ=
|
||||
github.com/fxamacker/cbor/v2 v2.9.1/go.mod h1:vM4b+DJCtHn+zz7h3FFp/hDAI9WNWCsZj23V5ytsSxQ=
|
||||
github.com/go-jose/go-jose/v4 v4.1.4 h1:moDMcTHmvE6Groj34emNPLs/qtYXRVcd6S7NHbHz3kA=
|
||||
github.com/go-jose/go-jose/v4 v4.1.4/go.mod h1:x4oUasVrzR7071A4TnHLGSPpNOm2a21K9Kf04k1rs08=
|
||||
github.com/go-logr/logr v1.2.2/go.mod h1:jdQByPbusPIv2/zmleS9BjJVeZ6kBagPoEUsqbVz/1A=
|
||||
github.com/go-logr/logr v1.4.3 h1:CjnDlHq8ikf6E492q6eKboGOC0T8CDaOvkHCIg8idEI=
|
||||
github.com/go-logr/logr v1.4.3/go.mod h1:9T104GzyrTigFIr8wt5mBrctHMim0Nb2HLGrmQ40KvY=
|
||||
@@ -72,6 +76,10 @@ github.com/go-openapi/testify/v2 v2.6.0 h1:5PKH2HE7YJ/LuRPQGvSxBRlFXNQhSetBLlGAg
|
||||
github.com/go-openapi/testify/v2 v2.6.0/go.mod h1:SgsVHtfooshd0tublTtJ50FPKhujf47YRqauXXOUxfw=
|
||||
github.com/go-task/slim-sprig/v3 v3.0.0 h1:sUs3vkvUymDpBKi3qH1YSqBQk9+9D/8M2mN1vB6EwHI=
|
||||
github.com/go-task/slim-sprig/v3 v3.0.0/go.mod h1:W848ghGpv3Qj3dhTPRyJypKRiqCdHZiAzKg9hl15HA8=
|
||||
github.com/go-test/deep v1.1.1 h1:0r/53hagsehfO4bzD2Pgr/+RgHqhmf+k1Bpse2cTu1U=
|
||||
github.com/go-test/deep v1.1.1/go.mod h1:5C2ZWiW0ErCdrYzpqxLbTX7MG14M9iiw8DgHncVwcsE=
|
||||
github.com/go-viper/mapstructure/v2 v2.5.0 h1:vM5IJoUAy3d7zRSVtIwQgBj7BiWtMPfmPEgAXnvj1Ro=
|
||||
github.com/go-viper/mapstructure/v2 v2.5.0/go.mod h1:oJDH3BJKyqBA2TXFhDsKDGDTlndYOZ6rGS0BRZIxGhM=
|
||||
github.com/golang/protobuf v1.5.4 h1:i7eJL8qZTpSEXOPTxNKhASYpMn+8e5Q6AdndVa1dWek=
|
||||
github.com/golang/protobuf v1.5.4/go.mod h1:lnTiLA8Wa4RWRcIUkrtSVa5nRhsEGBg48fD6rSs7xps=
|
||||
github.com/google/cel-go v0.29.2 h1:ZtDxkeiMmz0mxbKDYiNkE5Lk7V5edMRcaaDf2jX002k=
|
||||
@@ -89,6 +97,25 @@ github.com/google/uuid v1.6.0 h1:NIvaJDMOsjHA8n1jAhLSgzrAzy1Hgr+hNrb57e+94F0=
|
||||
github.com/google/uuid v1.6.0/go.mod h1:TIyPZe4MgqvfeYDBFedMoGGpEw/LqOeaOT+nhxU+yHo=
|
||||
github.com/grpc-ecosystem/grpc-gateway/v2 v2.29.0 h1:5VipnvEpbqr2gA2VbM+nYVbkIF28c5ZQfqCBQ5g2xfk=
|
||||
github.com/grpc-ecosystem/grpc-gateway/v2 v2.29.0/go.mod h1:Hyl3n6Twe1hvtd9XUXDec4pTvgMSEixRuQKPTMH2bNs=
|
||||
github.com/hashicorp/errwrap v1.0.0/go.mod h1:YH+1FKiLXxHSkmPseP+kNlulaMuP3n2brvKWEqk/Jc4=
|
||||
github.com/hashicorp/errwrap v1.1.0 h1:OxrOeh75EUXMY8TBjag2fzXGZ40LB6IKw45YeGUDY2I=
|
||||
github.com/hashicorp/errwrap v1.1.0/go.mod h1:YH+1FKiLXxHSkmPseP+kNlulaMuP3n2brvKWEqk/Jc4=
|
||||
github.com/hashicorp/go-cleanhttp v0.5.2 h1:035FKYIWjmULyFRBKPs8TBQoi0x6d9G4xc9neXJWAZQ=
|
||||
github.com/hashicorp/go-cleanhttp v0.5.2/go.mod h1:kO/YDlP8L1346E6Sodw+PrpBSV4/SoxCXGY6BqNFT48=
|
||||
github.com/hashicorp/go-hclog v1.6.3 h1:Qr2kF+eVWjTiYmU7Y31tYlP1h0q/X3Nl3tPGdaB11/k=
|
||||
github.com/hashicorp/go-hclog v1.6.3/go.mod h1:W4Qnvbt70Wk/zYJryRzDRU/4r0kIg0PVHBcfoyhpF5M=
|
||||
github.com/hashicorp/go-multierror v1.1.1 h1:H5DkEtf6CXdFp0N0Em5UCwQpXMWke8IA0+lD48awMYo=
|
||||
github.com/hashicorp/go-multierror v1.1.1/go.mod h1:iw975J/qwKPdAO1clOe2L8331t/9/fmwbPZ6JB6eMoM=
|
||||
github.com/hashicorp/go-retryablehttp v0.7.8 h1:ylXZWnqa7Lhqpk0L1P1LzDtGcCR0rPVUrx/c8Unxc48=
|
||||
github.com/hashicorp/go-retryablehttp v0.7.8/go.mod h1:rjiScheydd+CxvumBsIrFKlx3iS0jrZ7LvzFGFmuKbw=
|
||||
github.com/hashicorp/go-secure-stdlib/parseutil v0.2.0 h1:U+kC2dOhMFQctRfhK0gRctKAPTloZdMU5ZJxaesJ/VM=
|
||||
github.com/hashicorp/go-secure-stdlib/parseutil v0.2.0/go.mod h1:Ll013mhdmsVDuoIXVfBtvgGJsXDYkTw1kooNcoCXuE0=
|
||||
github.com/hashicorp/go-secure-stdlib/strutil v0.1.2 h1:kes8mmyCpxJsI7FTwtzRqEy9CdjCtrXrXGuOpxEA7Ts=
|
||||
github.com/hashicorp/go-secure-stdlib/strutil v0.1.2/go.mod h1:Gou2R9+il93BqX25LAKCLuM+y9U2T4hlwvT1yprcna4=
|
||||
github.com/hashicorp/go-sockaddr v1.0.7 h1:G+pTkSO01HpR5qCxg7lxfsFEZaG+C0VssTy/9dbT+Fw=
|
||||
github.com/hashicorp/go-sockaddr v1.0.7/go.mod h1:FZQbEYa1pxkQ7WLpyXJ6cbjpT8q0YgQaK/JakXqGyWw=
|
||||
github.com/hashicorp/hcl v1.0.1-vault-7 h1:ag5OxFVy3QYTFTJODRzTKVZ6xvdfLLCA1cy/Y6xGI0I=
|
||||
github.com/hashicorp/hcl v1.0.1-vault-7/go.mod h1:XYhtn6ijBSAj6n4YqAaf7RBPS4I06AItNorpy+MoQNM=
|
||||
github.com/inconshreveable/mousetrap v1.1.0 h1:wN+x4NVGpMsO7ErUn/mUI3vEoE6Jt13X2s0bqwp9tc8=
|
||||
github.com/inconshreveable/mousetrap v1.1.0/go.mod h1:vpF70FUmC8bwa3OWnCshd2FqLfsEA9PFc4w1p2J65bw=
|
||||
github.com/jackc/pgpassfile v1.0.0 h1:/6Hmqy13Ss2zCq62VdNG8tM1wchn8zjSGOBJ6icpsIM=
|
||||
@@ -105,6 +132,12 @@ github.com/klauspost/compress v1.19.0 h1:sXLILfc9jV2QYWkzFOPWStmcUVH2RHEB1JCdY2o
|
||||
github.com/klauspost/compress v1.19.0/go.mod h1:cwPg85FWrGar70rWktvGQj8/hthj3wpl0PGDogxkrSQ=
|
||||
github.com/kylelemons/godebug v1.1.0 h1:RPNrshWIDI6G2gRW9EHilWtl7Z6Sb1BR0xunSBf0SNc=
|
||||
github.com/kylelemons/godebug v1.1.0/go.mod h1:9/0rRGxNHcop5bhtWyNeEfOS8JIWk580+fNqagV/RAw=
|
||||
github.com/mattn/go-colorable v0.1.15 h1:+u9SLTRGnXv73cEsnsmoZBom+dMU88B2M0aDcWy0/jY=
|
||||
github.com/mattn/go-colorable v0.1.15/go.mod h1:6LmQG8QLFO4G5z1gPvYEzlUgJ2wF+stgPZH1UqBm1s8=
|
||||
github.com/mattn/go-isatty v0.0.24 h1:tGZZoVgT/KiqK1c8ocVLeDS8BSWMRd47J3Lbz7vsReI=
|
||||
github.com/mattn/go-isatty v0.0.24/go.mod h1:nMCL3Zebbrt45jsMDgnfIwz6ydEQApk5oEI3HqDio6A=
|
||||
github.com/mitchellh/mapstructure v1.5.0 h1:jeMsZIYE/09sWLaz43PL7Gy6RuMjD2eJVyuac5Z2hdY=
|
||||
github.com/mitchellh/mapstructure v1.5.0/go.mod h1:bFUtVrKA4DC2yAKiSyO/QUcy7e+RRV2QTWOzhPopBRo=
|
||||
github.com/modern-go/concurrent v0.0.0-20180228061459-e0a39a4cb421/go.mod h1:6dJC0mAP4ikYIbvyc7fijjWJddQyLn8Ig3JB5CqoB9Q=
|
||||
github.com/modern-go/concurrent v0.0.0-20180306012644-bacd9c7ef1dd h1:TRLaZ9cD/w8PVh93nsPXa1VrQ6jlwL5oN8l14QlcNfg=
|
||||
github.com/modern-go/concurrent v0.0.0-20180306012644-bacd9c7ef1dd/go.mod h1:6dJC0mAP4ikYIbvyc7fijjWJddQyLn8Ig3JB5CqoB9Q=
|
||||
@@ -117,6 +150,8 @@ github.com/onsi/ginkgo/v2 v2.27.4 h1:fcEcQW/A++6aZAZQNUmNjvA9PSOzefMJBerHJ4t8v8Y
|
||||
github.com/onsi/ginkgo/v2 v2.27.4/go.mod h1:ArE1D/XhNXBXCBkKOLkbsb2c81dQHCRcF5zwn/ykDRo=
|
||||
github.com/onsi/gomega v1.39.0 h1:y2ROC3hKFmQZJNFeGAMeHZKkjBL65mIZcvrLQBF9k6Q=
|
||||
github.com/onsi/gomega v1.39.0/go.mod h1:ZCU1pkQcXDO5Sl9/VVEGlDyp+zm0m1cmeG5TOzLgdh4=
|
||||
github.com/openbao/openbao/api/v2 v2.7.0 h1:3CD1l3tr39nQraCgFGAWA5vYvPFzZoZrt3NL7DMQKAc=
|
||||
github.com/openbao/openbao/api/v2 v2.7.0/go.mod h1:uXbMoyH2pjSvNyTepinUvLde8pOJB82EuhUCfOKnKbo=
|
||||
github.com/pkg/errors v0.9.1 h1:FEBLx1zS214owpjy7qsBeixbURkuhQAwrK5UwLGTwt4=
|
||||
github.com/pkg/errors v0.9.1/go.mod h1:bwawxfHBFNV+L2hUp1rHADufV3IMtnDRdf1r5NINEl0=
|
||||
github.com/pmezard/go-difflib v1.0.0/go.mod h1:iKH77koFhYxTK1pcRnkKkqfTogsbg7gZNVY4sRDYZ/4=
|
||||
@@ -131,6 +166,8 @@ github.com/prometheus/common v0.70.0/go.mod h1:S/SFasQmgGiYH6C81LKCtYa8QACgthGg5
|
||||
github.com/prometheus/procfs v0.21.1 h1:GljZCt+zSTS+NZq88cyQ1LjZ+RCHp3uVuabBWA5+OJI=
|
||||
github.com/prometheus/procfs v0.21.1/go.mod h1:aB55Cww9pdSJVHk0hUf0inxWyyjPogFIjmHKYgMKmtY=
|
||||
github.com/russross/blackfriday/v2 v2.1.0/go.mod h1:+Rmxgy9KzJVeS9/2gXHxylqXiyQDYRxCVz55jmeOWTM=
|
||||
github.com/ryanuber/go-glob v1.0.0 h1:iQh3xXAumdQ+4Ufa5b25cRpC5TYKlno6hsv6Cb3pkBk=
|
||||
github.com/ryanuber/go-glob v1.0.0/go.mod h1:807d1WSdnB0XRJzKNil9Om6lcp/3a0v4qIHxIXzX/Yc=
|
||||
github.com/spf13/cobra v1.10.2 h1:DMTTonx5m65Ic0GOoRY2c16WCbHxOOw6xxezuLaBpcU=
|
||||
github.com/spf13/cobra v1.10.2/go.mod h1:7C1pvHqHw5A4vrJfjNwvOdzYu0Gml16OCs2GRiTUUS4=
|
||||
github.com/spf13/pflag v1.0.9/go.mod h1:McXfInJRrz4CZXVZOBLb0bTZqETkiAhM9Iw0y3An2Bg=
|
||||
@@ -141,8 +178,8 @@ github.com/stretchr/objx v0.5.3 h1:jmXUvGomnU1o3W/V5h2VEradbpJDwGrzugQQvL0POH4=
|
||||
github.com/stretchr/objx v0.5.3/go.mod h1:rDQraq+vQZU7Fde9LOZLr8Tax6zZvy4kuNKF+QYS+U0=
|
||||
github.com/stretchr/testify v1.3.0/go.mod h1:M5WIy9Dh21IEIfnGCwXGc5bZfKNJtfHm1UVUgZn+9EI=
|
||||
github.com/stretchr/testify v1.7.0/go.mod h1:6Fq8oRcR53rry900zMqJjRRixrwX3KX962/h/Wwjteg=
|
||||
github.com/stretchr/testify v1.11.1 h1:7s2iGBzp5EwR7/aIZr8ao5+dra3wiQyKjjFuvgVKu7U=
|
||||
github.com/stretchr/testify v1.11.1/go.mod h1:wZwfW3scLgRK+23gO65QZefKpKQRnfz6sD981Nm4B6U=
|
||||
github.com/stretchr/testify v1.12.1 h1:EuwCh5fleGS7H32xRwO3wRGT7DxrDhLAT6FF8MpWDWE=
|
||||
github.com/stretchr/testify v1.12.1/go.mod h1:MDEgiDPPsNp5cuIrHPPCyornHKgEVbtFUmoNlxoYthg=
|
||||
github.com/x448/float16 v0.8.4 h1:qLwI1I70+NjRFUR3zs1JPUCgaCXSh3SW62uAKT1mSBM=
|
||||
github.com/x448/float16 v0.8.4/go.mod h1:14CWIYCyZA/cWjXOioeEpHeN/83MdbZDRQHoFcYsOfg=
|
||||
go.opentelemetry.io/auto/sdk v1.2.1 h1:jXsnJ4Lmnqd11kwkBV2LgLoFMZKizbCi5fNZ/ipaZ64=
|
||||
@@ -180,8 +217,8 @@ golang.org/x/exp v0.0.0-20260410095643-746e56fc9e2f h1:W3F4c+6OLc6H2lb//N1q4WpJk
|
||||
golang.org/x/exp v0.0.0-20260410095643-746e56fc9e2f/go.mod h1:J1xhfL/vlindoeF/aINzNzt2Bket5bjo9sdOYzOsU80=
|
||||
golang.org/x/mod v0.38.0 h1:MECBjubtXD7yj4HrhIUcywNaGeNVUdfVnxmPajOk4yk=
|
||||
golang.org/x/mod v0.38.0/go.mod h1:V6Xz0pq8TQ3dGqVQ1FVHuelZpAL0uNhSkk9ogYP3c40=
|
||||
golang.org/x/net v0.57.0 h1:K5+3DljvIuDG9/Jv9rvyMywYNFCQ9RSUY6OOTTkT+tE=
|
||||
golang.org/x/net v0.57.0/go.mod h1:KpXc8iv+r3XplLAG/f7Jsf9RPszJzdR0f58q9vGOuEU=
|
||||
golang.org/x/net v0.58.0 h1:ynWG7rqYi4ccpTEuPZ2QGWHktVEM9DMCj9yzDE0Q7To=
|
||||
golang.org/x/net v0.58.0/go.mod h1:YwCddHnFlT7eLQqVprV19OnhLGtc5xOKgE0RyqgfWAU=
|
||||
golang.org/x/oauth2 v0.36.0 h1:peZ/1z27fi9hUOFCAZaHyrpWG5lwe0RJEEEeH0ThlIs=
|
||||
golang.org/x/oauth2 v0.36.0/go.mod h1:YDBUJMTkDnJS+A4BP4eZBjCqtokkg1hODuPjwiGPO7Q=
|
||||
golang.org/x/sync v0.22.0 h1:SZjpbeLmrCk4xhRSZFNZW5gFUeCeFgjekvI/+gfScek=
|
||||
|
||||
@@ -0,0 +1,45 @@
|
||||
package kubernetes
|
||||
|
||||
import (
|
||||
databasev1alpha1 "git.ddupan.top/panxiao81/ayatori/api/database/v1alpha1"
|
||||
"git.ddupan.top/panxiao81/ayatori/internal/database/application"
|
||||
"git.ddupan.top/panxiao81/ayatori/internal/database/domain/instance"
|
||||
)
|
||||
|
||||
func instanceRecord(object *databasev1alpha1.PostgreSQLInstance) (*application.InstanceRecord, error) {
|
||||
identity, err := instance.NewIdentity(string(object.UID), object.Name)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
revision, err := instance.NewRevision(object.Generation)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
spec := object.Spec
|
||||
endpoint, err := instance.NewEndpoint(instance.EndpointValues{
|
||||
Host: spec.Endpoint.Host, HostAddr: spec.Endpoint.HostAddr,
|
||||
Port: int(spec.Endpoint.Port), ManagementDatabase: string(spec.Endpoint.Database),
|
||||
TLSMode: instance.TLSMode(spec.Endpoint.SSLMode),
|
||||
})
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
credential, err := instance.NewCredentialReference(instance.CredentialReferenceValues{
|
||||
Name: string(spec.AdminCredentialRef.Name),
|
||||
UsernameKey: spec.AdminCredentialRef.UsernameKey, PasswordKey: spec.AdminCredentialRef.PasswordKey,
|
||||
})
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
definition, err := instance.NewDefinition(endpoint, credential)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
target, err := instance.NewObservationTarget(identity, revision, definition)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
return &application.InstanceRecord{
|
||||
Target: target, Revision: object.ResourceVersion, Deleting: !object.DeletionTimestamp.IsZero(),
|
||||
}, nil
|
||||
}
|
||||
@@ -0,0 +1,112 @@
|
||||
package kubernetes
|
||||
|
||||
import (
|
||||
"context"
|
||||
"errors"
|
||||
|
||||
databasev1alpha1 "git.ddupan.top/panxiao81/ayatori/api/database/v1alpha1"
|
||||
"git.ddupan.top/panxiao81/ayatori/internal/database/application"
|
||||
"git.ddupan.top/panxiao81/ayatori/internal/database/domain/instance"
|
||||
"k8s.io/apimachinery/pkg/api/equality"
|
||||
apierrors "k8s.io/apimachinery/pkg/api/errors"
|
||||
"k8s.io/apimachinery/pkg/api/meta"
|
||||
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
|
||||
"sigs.k8s.io/controller-runtime/pkg/client"
|
||||
"sigs.k8s.io/controller-runtime/pkg/controller/controllerutil"
|
||||
)
|
||||
|
||||
const InstanceFinalizer = "database.ayatori.ddupan.top/instance-protection"
|
||||
|
||||
type InstanceResources struct {
|
||||
Client client.Client
|
||||
Reader client.Reader
|
||||
}
|
||||
|
||||
func (r *InstanceResources) LoadInstance(ctx context.Context, name string) (*application.InstanceRecord, error) {
|
||||
object := &databasev1alpha1.PostgreSQLInstance{}
|
||||
if err := r.Reader.Get(ctx, client.ObjectKey{Name: name}, object); err != nil {
|
||||
return nil, client.IgnoreNotFound(err)
|
||||
}
|
||||
return instanceRecord(object)
|
||||
}
|
||||
|
||||
func (r *InstanceResources) ProtectInstance(ctx context.Context, record *application.InstanceRecord) (*application.InstanceRecord, error) {
|
||||
object, err := r.instanceAtVersion(ctx, record)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
if controllerutil.AddFinalizer(object, InstanceFinalizer) {
|
||||
if err := r.Client.Update(ctx, object); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
}
|
||||
return instanceRecord(object)
|
||||
}
|
||||
|
||||
func (r *InstanceResources) InstanceReferences(ctx context.Context, name string) (string, error) {
|
||||
// 删除判断必须直读 API;包含 Released、删除中的 Database 和尚未绑定的申请。
|
||||
// 不按旧 Instance UID 忽略引用,也不依赖 informer 索引的及时性。
|
||||
databases := &databasev1alpha1.PostgreSQLDatabaseList{}
|
||||
if err := r.Reader.List(ctx, databases); err != nil {
|
||||
return "", err
|
||||
}
|
||||
for _, database := range databases.Items {
|
||||
if string(database.Spec.InstanceRef.Name) == name {
|
||||
return "Database/" + database.Name, nil
|
||||
}
|
||||
}
|
||||
tenants := &databasev1alpha1.PostgreSQLTenantList{}
|
||||
if err := r.Reader.List(ctx, tenants); err != nil {
|
||||
return "", err
|
||||
}
|
||||
for _, tenant := range tenants.Items {
|
||||
if tenant.Spec.Provision != nil && string(tenant.Spec.Provision.InstanceRef.Name) == name {
|
||||
return "Tenant/" + tenant.Namespace + "/" + tenant.Name, nil
|
||||
}
|
||||
}
|
||||
return "", nil
|
||||
}
|
||||
|
||||
func (r *InstanceResources) PresentInstance(ctx context.Context, result application.InstanceResult) error {
|
||||
if result.Record == nil {
|
||||
return nil
|
||||
}
|
||||
object, err := r.instanceAtVersion(ctx, result.Record)
|
||||
if err != nil {
|
||||
return client.IgnoreNotFound(err)
|
||||
}
|
||||
previous := object.Status.DeepCopy()
|
||||
object.Status.Phase = string(result.Snapshot.Phase)
|
||||
object.Status.ObservedGeneration = object.Generation
|
||||
object.Status.PostgreSQLVersion = result.Snapshot.ReportedVersion
|
||||
ready := metav1.ConditionFalse
|
||||
if result.Snapshot.Readiness == instance.Ready {
|
||||
ready = metav1.ConditionTrue
|
||||
}
|
||||
meta.SetStatusCondition(&object.Status.Conditions, metav1.Condition{
|
||||
Type: "Ready", Status: ready, ObservedGeneration: object.Generation,
|
||||
Reason: result.Reason, Message: result.Message,
|
||||
})
|
||||
if !equality.Semantic.DeepEqual(*previous, object.Status) {
|
||||
if err := r.Client.Status().Update(ctx, object); err != nil {
|
||||
return err
|
||||
}
|
||||
}
|
||||
if result.RemoveProtection && controllerutil.RemoveFinalizer(object, InstanceFinalizer) {
|
||||
return r.Client.Update(ctx, object)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
func (r *InstanceResources) instanceAtVersion(ctx context.Context, record *application.InstanceRecord) (*databasev1alpha1.PostgreSQLInstance, error) {
|
||||
object := &databasev1alpha1.PostgreSQLInstance{}
|
||||
name := record.Target.Identity().Name()
|
||||
if err := r.Reader.Get(ctx, client.ObjectKey{Name: name}, object); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
if string(object.UID) != record.Target.Identity().UID() || object.ResourceVersion != record.Revision {
|
||||
return nil, apierrors.NewConflict(databasev1alpha1.GroupVersion.WithResource("postgresqlinstances").GroupResource(),
|
||||
name, errors.New("Instance 快照已过期,请重新观察"))
|
||||
}
|
||||
return object, nil
|
||||
}
|
||||
@@ -0,0 +1,140 @@
|
||||
/*
|
||||
Copyright 2026.
|
||||
|
||||
Licensed under the Apache License, Version 2.0 (the "License");
|
||||
you may not use this file except in compliance with the License.
|
||||
You may obtain a copy of the License at
|
||||
|
||||
http://www.apache.org/licenses/LICENSE-2.0
|
||||
|
||||
Unless required by applicable law or agreed to in writing, software
|
||||
distributed under the License is distributed on an "AS IS" BASIS,
|
||||
WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
|
||||
See the License for the specific language governing permissions and
|
||||
limitations under the License.
|
||||
*/
|
||||
|
||||
// Package openbao 通过官方 SDK 适配应用凭据,不保存资源归属或重建供应状态。
|
||||
package openbao
|
||||
|
||||
import (
|
||||
"context"
|
||||
"errors"
|
||||
"maps"
|
||||
"net/http"
|
||||
"regexp"
|
||||
"slices"
|
||||
"strings"
|
||||
|
||||
bao "github.com/openbao/openbao/api/v2"
|
||||
|
||||
"git.ddupan.top/panxiao81/ayatori/internal/database/application"
|
||||
)
|
||||
|
||||
var (
|
||||
ErrInvalidLocation = errors.New("credential location is outside the configured scope")
|
||||
ErrUnavailable = errors.New("credential backend unavailable")
|
||||
ErrNotFound = errors.New("application credential not found")
|
||||
ErrConflict = errors.New("credential creation requires manual conflict resolution")
|
||||
ErrUncertain = errors.New("credential creation outcome is uncertain; manual resolution required")
|
||||
)
|
||||
|
||||
var pathSegment = regexp.MustCompile(`^[A-Za-z0-9_-]+$`)
|
||||
|
||||
// Credentials 使用独立的 SDK client;认证与短期 token 生命周期由部署装配负责。
|
||||
// 本适配器既不自动认领已有值,也不提供覆盖、轮换或删除操作。
|
||||
type Credentials struct {
|
||||
kv *bao.KVv2
|
||||
basePath string
|
||||
}
|
||||
|
||||
// NewCredentials 不登录、不读取环境 token。调用方必须提供专用的已认证 client。
|
||||
// 禁用 SDK 写入重试,防止第一次结果丢失后被 CAS 错误掩盖。
|
||||
func NewCredentials(client *bao.Client, mount, basePath string) (*Credentials, error) {
|
||||
if client == nil || !validPath(mount) || !validPath(basePath) {
|
||||
return nil, ErrInvalidLocation
|
||||
}
|
||||
client.SetMaxRetries(0)
|
||||
return &Credentials{kv: client.KVv2(mount), basePath: basePath}, nil
|
||||
}
|
||||
|
||||
func validPath(value string) bool {
|
||||
for segment := range strings.SplitSeq(value, "/") {
|
||||
if !pathSegment.MatchString(segment) || segment == "data" || segment == "metadata" {
|
||||
return false
|
||||
}
|
||||
}
|
||||
return true
|
||||
}
|
||||
|
||||
// ProvisionPath 只按 Database UID 定位;调用方须先持久化位置,再执行外部写入。
|
||||
func (c *Credentials) ProvisionPath(databaseUID string) (string, error) {
|
||||
if !pathSegment.MatchString(databaseUID) {
|
||||
return "", ErrInvalidLocation
|
||||
}
|
||||
return c.basePath + "/" + databaseUID, nil
|
||||
}
|
||||
|
||||
func (c *Credentials) accepts(path string) bool {
|
||||
return validPath(path) && strings.HasPrefix(path, c.basePath+"/")
|
||||
}
|
||||
|
||||
// Read 只读取调用方已确认关联的路径;成功读取不构成对既有凭据的自动认领。
|
||||
func (c *Credentials) Read(ctx context.Context, path string) (application.ApplicationCredential, error) {
|
||||
if !c.accepts(path) {
|
||||
return application.ApplicationCredential{}, ErrInvalidLocation
|
||||
}
|
||||
secret, err := c.kv.Get(ctx, path)
|
||||
if errors.Is(err, bao.ErrSecretNotFound) {
|
||||
return application.ApplicationCredential{}, ErrNotFound
|
||||
}
|
||||
if err != nil {
|
||||
return application.ApplicationCredential{}, ErrUnavailable
|
||||
}
|
||||
if secret == nil || secret.Data == nil {
|
||||
return application.ApplicationCredential{}, ErrNotFound
|
||||
}
|
||||
return application.ParseApplicationCredential(secret.Data)
|
||||
}
|
||||
|
||||
// Create 只创建从未存在过的路径,并验证回读七键与提交值完全一致。
|
||||
// 任何不确定写入都不返回凭据;上层必须停止供应并持久化冲突,不能重新生成密码。
|
||||
func (c *Credentials) Create(ctx context.Context, path string, credential application.ApplicationCredential) error {
|
||||
if !c.accepts(path) {
|
||||
return ErrInvalidLocation
|
||||
}
|
||||
if err := credential.Validate(); err != nil {
|
||||
return err
|
||||
}
|
||||
if ctx.Err() != nil {
|
||||
return ErrUnavailable
|
||||
}
|
||||
data := credential.SecretData()
|
||||
created, err := c.kv.Put(ctx, path, data, bao.WithCheckAndSet(0))
|
||||
if err != nil {
|
||||
// 明确的认证/权限拒绝没有发生写入,可以等待依赖恢复。
|
||||
// SDK 的原始错误可能携带路径及响应体,不向外传播。
|
||||
if response, ok := errors.AsType[*bao.ResponseError](err); ok {
|
||||
switch response.StatusCode {
|
||||
case http.StatusUnauthorized, http.StatusForbidden:
|
||||
return ErrUnavailable
|
||||
case http.StatusBadRequest:
|
||||
if slices.Contains(response.Errors, "check-and-set parameter did not match the current version") {
|
||||
return ErrConflict
|
||||
}
|
||||
}
|
||||
}
|
||||
return ErrUncertain
|
||||
}
|
||||
if created == nil || created.VersionMetadata == nil || created.VersionMetadata.Version != 1 {
|
||||
return ErrUncertain
|
||||
}
|
||||
observed, err := c.kv.Get(ctx, path)
|
||||
if err != nil || observed == nil || observed.VersionMetadata == nil || observed.VersionMetadata.Version != 1 {
|
||||
return ErrUncertain
|
||||
}
|
||||
if !maps.Equal(data, observed.Data) {
|
||||
return ErrUncertain
|
||||
}
|
||||
return nil
|
||||
}
|
||||
@@ -0,0 +1,227 @@
|
||||
//go:build integration
|
||||
|
||||
/*
|
||||
Copyright 2026.
|
||||
|
||||
Licensed under the Apache License, Version 2.0 (the "License");
|
||||
you may not use this file except in compliance with the License.
|
||||
You may obtain a copy of the License at
|
||||
|
||||
http://www.apache.org/licenses/LICENSE-2.0
|
||||
|
||||
Unless required by applicable law or agreed to in writing, software
|
||||
distributed under the License is distributed on an "AS IS" BASIS,
|
||||
WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
|
||||
See the License for the specific language governing permissions and
|
||||
limitations under the License.
|
||||
*/
|
||||
|
||||
package openbao_test
|
||||
|
||||
import (
|
||||
"context"
|
||||
"errors"
|
||||
"maps"
|
||||
"net/http"
|
||||
"net/http/httptest"
|
||||
"net/http/httputil"
|
||||
"net/url"
|
||||
"os/exec"
|
||||
"regexp"
|
||||
"strings"
|
||||
"sync"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
bao "github.com/openbao/openbao/api/v2"
|
||||
|
||||
"git.ddupan.top/panxiao81/ayatori/internal/database/adapter/openbao"
|
||||
)
|
||||
|
||||
// 只连接本测试创建的无持久卷 dev server,不接受生产地址或环境 token。
|
||||
func baoFixture(t *testing.T) *bao.Client {
|
||||
t.Helper()
|
||||
const image = "openbao/openbao@sha256:5b2486ab0fb90bbc788cc345b0a08616dfb375873ee8be5df3a2fd4d378a67e0"
|
||||
prepareBaoImage(t, image)
|
||||
// 冷缓存拉取不占用容器启动和健康检查的一分钟预算。
|
||||
ctx, cancel := context.WithTimeout(t.Context(), time.Minute)
|
||||
defer cancel()
|
||||
output, err := exec.CommandContext(ctx, "docker", "run", "--pull=never", "--rm", "-d", "-p", "127.0.0.1::8200",
|
||||
image, "server", "-dev", "-dev-root-token-id="+fixtureToken, "-dev-listen-address=0.0.0.0:8200").Output()
|
||||
if err != nil {
|
||||
t.Fatalf("cannot start isolated OpenBao fixture: %s", baoCommandError(ctx, err))
|
||||
}
|
||||
id := strings.TrimSpace(string(output))
|
||||
if !regexp.MustCompile(`^[a-f0-9]{64}$`).MatchString(id) {
|
||||
t.Fatal("unexpected fixture container ID")
|
||||
}
|
||||
t.Cleanup(func() {
|
||||
cleanup, stop := context.WithTimeout(context.Background(), 30*time.Second)
|
||||
defer stop()
|
||||
if exec.CommandContext(cleanup, "docker", "rm", "-f", id).Run() != nil {
|
||||
t.Error("OpenBao fixture cleanup failed")
|
||||
}
|
||||
})
|
||||
output, err = exec.CommandContext(ctx, "docker", "inspect", "--format",
|
||||
`{{(index (index .NetworkSettings.Ports "8200/tcp") 0).HostPort}}`, id).Output()
|
||||
if err != nil {
|
||||
t.Fatalf("cannot inspect fixture port: %s", baoCommandError(ctx, err))
|
||||
}
|
||||
client := fixtureClient(t, "http://127.0.0.1:"+strings.TrimSpace(string(output)))
|
||||
client.SetMaxRetries(0)
|
||||
for {
|
||||
if _, err := client.Sys().HealthWithContext(ctx); err == nil {
|
||||
return client
|
||||
}
|
||||
select {
|
||||
case <-ctx.Done():
|
||||
t.Fatal("OpenBao fixture startup timed out")
|
||||
case <-time.After(100 * time.Millisecond):
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func prepareBaoImage(t *testing.T, image string) {
|
||||
t.Helper()
|
||||
ctx, cancel := context.WithTimeout(t.Context(), 5*time.Minute)
|
||||
defer cancel()
|
||||
if exec.CommandContext(ctx, "docker", "image", "inspect", image).Run() == nil {
|
||||
return
|
||||
}
|
||||
t.Log("pulling isolated OpenBao fixture image (timeout: 5m)")
|
||||
if _, err := exec.CommandContext(ctx, "docker", "pull", image).Output(); err != nil {
|
||||
t.Fatalf("cannot pull OpenBao fixture image: %s", baoCommandError(ctx, err))
|
||||
}
|
||||
}
|
||||
|
||||
// 保留 Docker stderr 与超时原因,但不泄露测试 token/password 或完整命令参数。
|
||||
func baoCommandError(ctx context.Context, err error) string {
|
||||
detail := err.Error()
|
||||
if exitErr, ok := errors.AsType[*exec.ExitError](err); ok {
|
||||
detail += ": " + strings.TrimSpace(string(exitErr.Stderr))
|
||||
}
|
||||
if ctx.Err() != nil {
|
||||
detail += ": " + ctx.Err().Error()
|
||||
}
|
||||
return strings.NewReplacer(fixtureToken, "[REDACTED]", fixturePassword, "[REDACTED]").Replace(detail)
|
||||
}
|
||||
|
||||
func TestBaoCommandError(t *testing.T) {
|
||||
err := &exec.ExitError{Stderr: []byte("registry unavailable " + fixtureToken + " " + fixturePassword)}
|
||||
ctx, cancel := context.WithCancel(t.Context())
|
||||
cancel()
|
||||
detail := baoCommandError(ctx, err)
|
||||
if !strings.Contains(detail, "registry unavailable") || !strings.Contains(detail, "context canceled") {
|
||||
t.Fatal("Docker diagnostic or context failure was lost")
|
||||
}
|
||||
if strings.Contains(detail, fixtureToken) || strings.Contains(detail, fixturePassword) {
|
||||
t.Fatal("Docker diagnostic exposed fixture credentials")
|
||||
}
|
||||
}
|
||||
|
||||
func TestCredentialConcurrentCreateWithRealOpenBao(t *testing.T) {
|
||||
root := baoFixture(t)
|
||||
store := fixtureStore(t, root)
|
||||
credential := fixtureCredential(t)
|
||||
results := make(chan error, 2)
|
||||
var workers sync.WaitGroup
|
||||
for range 2 {
|
||||
workers.Go(func() { results <- store.Create(t.Context(), credentialPath, credential) })
|
||||
}
|
||||
workers.Wait()
|
||||
close(results)
|
||||
succeeded, conflicted := 0, 0
|
||||
for err := range results {
|
||||
switch err {
|
||||
case nil:
|
||||
succeeded++
|
||||
case openbao.ErrConflict:
|
||||
conflicted++
|
||||
default:
|
||||
t.Fatal("unexpected concurrent create result")
|
||||
}
|
||||
}
|
||||
if succeeded != 1 || conflicted != 1 {
|
||||
t.Fatal("CAS must allow exactly one creator")
|
||||
}
|
||||
}
|
||||
|
||||
func TestCredentialLostWriteResponseWithRealOpenBao(t *testing.T) {
|
||||
root := baoFixture(t)
|
||||
address, err := url.Parse(root.Address())
|
||||
if err != nil {
|
||||
t.Fatal("invalid fixture address")
|
||||
}
|
||||
proxy := httputil.NewSingleHostReverseProxy(address)
|
||||
proxy.ModifyResponse = func(response *http.Response) error {
|
||||
if response.Request.Method == http.MethodPut && response.StatusCode == http.StatusOK {
|
||||
return errors.New("fixture drops successful write response")
|
||||
}
|
||||
return nil
|
||||
}
|
||||
proxy.ErrorHandler = func(w http.ResponseWriter, _ *http.Request, _ error) {
|
||||
w.WriteHeader(http.StatusBadGateway)
|
||||
}
|
||||
server := httptest.NewServer(proxy)
|
||||
defer server.Close()
|
||||
store := fixtureStore(t, fixtureClient(t, server.URL))
|
||||
credential := fixtureCredential(t)
|
||||
if err := store.Create(t.Context(), credentialPath, credential); err != openbao.ErrUncertain {
|
||||
t.Fatal("lost response must stop provisioning")
|
||||
}
|
||||
confirmed, err := root.KVv2("secret").Get(t.Context(), credentialPath)
|
||||
if err != nil || !maps.Equal(confirmed.Data, credential.SecretData()) || confirmed.VersionMetadata.Version != 1 {
|
||||
t.Fatal("fault injection did not preserve the original write")
|
||||
}
|
||||
if err := fixtureStore(t, root).Create(t.Context(), credentialPath, credential); err != openbao.ErrConflict {
|
||||
t.Fatal("restart must not adopt an unconfirmed write")
|
||||
}
|
||||
}
|
||||
|
||||
func TestCredentialsWithRealOpenBao(t *testing.T) {
|
||||
root := baoFixture(t)
|
||||
ctx := t.Context()
|
||||
// root 仅用于 fixture 装配;实际读写使用固定前缀的短期 token。
|
||||
policy := `path "secret/data/applications/*" { capabilities = ["create", "update", "read"] }`
|
||||
if err := root.Sys().PutPolicyWithContext(ctx, "application-fixture", policy); err != nil {
|
||||
t.Fatal("cannot configure fixture policy")
|
||||
}
|
||||
secret, err := root.Auth().Token().CreateWithContext(ctx, &bao.TokenCreateRequest{
|
||||
Policies: []string{"application-fixture"}, NoDefaultPolicy: true, TTL: "5m",
|
||||
})
|
||||
if err != nil {
|
||||
t.Fatal("cannot create scoped fixture token")
|
||||
}
|
||||
client := fixtureClient(t, root.Address())
|
||||
client.SetToken(secret.Auth.ClientToken)
|
||||
store := fixtureStore(t, client)
|
||||
credential := fixtureCredential(t)
|
||||
if err := store.Create(ctx, credentialPath, credential); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
// 重建适配器读取已确认路径;重复 Create 仍报冲突,不把读取当作认领。
|
||||
restarted := fixtureStore(t, client)
|
||||
observed, err := restarted.Read(ctx, credentialPath)
|
||||
if err != nil || !maps.Equal(observed.SecretData(), credential.SecretData()) {
|
||||
t.Fatal("confirmed credential was not preserved across adapter restart")
|
||||
}
|
||||
if err := restarted.Create(ctx, credentialPath, credential); !errors.Is(err, openbao.ErrConflict) {
|
||||
t.Fatal("existing credential must conflict even if contents match")
|
||||
}
|
||||
metadata, err := root.KVv2("secret").GetMetadata(ctx, credentialPath)
|
||||
if err != nil || metadata.CurrentVersion != 1 {
|
||||
t.Fatal("duplicate create changed credential version")
|
||||
}
|
||||
if _, err := client.KVv2("secret").Get(ctx, "management/instance"); err == nil {
|
||||
t.Fatal("scoped token accessed management credentials")
|
||||
}
|
||||
if err := root.KVv2("secret").Delete(ctx, credentialPath); err != nil {
|
||||
t.Fatal("cannot soft-delete fixture credential")
|
||||
}
|
||||
if _, err := store.Read(ctx, credentialPath); err != openbao.ErrNotFound {
|
||||
t.Fatal("soft-deleted credential must not be usable")
|
||||
}
|
||||
if err := store.Create(ctx, credentialPath, credential); err != openbao.ErrConflict {
|
||||
t.Fatal("soft-deleted credential must not be recreated")
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,190 @@
|
||||
/*
|
||||
Copyright 2026.
|
||||
|
||||
Licensed under the Apache License, Version 2.0 (the "License");
|
||||
you may not use this file except in compliance with the License.
|
||||
You may obtain a copy of the License at
|
||||
|
||||
http://www.apache.org/licenses/LICENSE-2.0
|
||||
|
||||
Unless required by applicable law or agreed to in writing, software
|
||||
distributed under the License is distributed on an "AS IS" BASIS,
|
||||
WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
|
||||
See the License for the specific language governing permissions and
|
||||
limitations under the License.
|
||||
*/
|
||||
|
||||
package openbao_test
|
||||
|
||||
import (
|
||||
"context"
|
||||
"encoding/json"
|
||||
"errors"
|
||||
"net/http"
|
||||
"net/http/httptest"
|
||||
"sync/atomic"
|
||||
"testing"
|
||||
|
||||
bao "github.com/openbao/openbao/api/v2"
|
||||
|
||||
"git.ddupan.top/panxiao81/ayatori/internal/database/adapter/openbao"
|
||||
"git.ddupan.top/panxiao81/ayatori/internal/database/application"
|
||||
)
|
||||
|
||||
const (
|
||||
credentialPath = "applications/database-uid"
|
||||
fixturePassword = "AYATORI-TEST-ONLY-application-password"
|
||||
fixtureToken = "AYATORI-TEST-ONLY-bao-token"
|
||||
kvDataKey = "data"
|
||||
)
|
||||
|
||||
func fixtureCredential(t *testing.T) application.ApplicationCredential {
|
||||
t.Helper()
|
||||
credential, err := application.ParseApplicationCredential(map[string]any{
|
||||
"username": "app_owner", "password": fixturePassword, "database": "app",
|
||||
"host": "postgres.example", "hostaddr": "192.0.2.1", "port": "5432", "sslmode": "verify-full",
|
||||
})
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
return credential
|
||||
}
|
||||
|
||||
func TestCredentialReadbackMustConfirmTheWrite(t *testing.T) {
|
||||
for _, scenario := range []string{"read failure", "changed version", "changed password", "missing metadata"} {
|
||||
t.Run(scenario, func(t *testing.T) {
|
||||
credential := fixtureCredential(t)
|
||||
var writes atomic.Int32
|
||||
server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
|
||||
if r.Method == http.MethodPut {
|
||||
writes.Add(1)
|
||||
var request struct {
|
||||
Options struct {
|
||||
CAS *int `json:"cas"`
|
||||
} `json:"options"`
|
||||
}
|
||||
if json.NewDecoder(r.Body).Decode(&request) != nil || request.Options.CAS == nil || *request.Options.CAS != 0 {
|
||||
t.Error("create request must explicitly require CAS=0")
|
||||
}
|
||||
if err := json.NewEncoder(w).Encode(map[string]any{kvDataKey: map[string]any{"version": 1}}); err != nil {
|
||||
t.Error("cannot encode fixture write response")
|
||||
}
|
||||
return
|
||||
}
|
||||
if scenario == "read failure" {
|
||||
w.WriteHeader(http.StatusServiceUnavailable)
|
||||
return
|
||||
}
|
||||
data := credential.SecretData()
|
||||
version := 1
|
||||
if scenario == "changed version" {
|
||||
version = 2
|
||||
}
|
||||
if scenario == "changed password" {
|
||||
data["password"] = "modified"
|
||||
}
|
||||
response := map[string]any{kvDataKey: data}
|
||||
if scenario != "missing metadata" {
|
||||
response["metadata"] = map[string]any{"version": version}
|
||||
}
|
||||
if err := json.NewEncoder(w).Encode(map[string]any{kvDataKey: response}); err != nil {
|
||||
t.Error("cannot encode fixture read response")
|
||||
}
|
||||
}))
|
||||
defer server.Close()
|
||||
store := fixtureStore(t, fixtureClient(t, server.URL))
|
||||
if err := store.Create(t.Context(), credentialPath, credential); err != openbao.ErrUncertain || writes.Load() != 1 {
|
||||
t.Fatal("unconfirmed readback must stop after one write")
|
||||
}
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
func fixtureClient(t *testing.T, address string) *bao.Client {
|
||||
t.Helper()
|
||||
config := bao.DefaultConfig()
|
||||
config.Address = address
|
||||
client, err := bao.NewClient(config)
|
||||
if err != nil {
|
||||
t.Fatal("cannot construct fixture client")
|
||||
}
|
||||
client.SetToken(fixtureToken)
|
||||
return client
|
||||
}
|
||||
|
||||
func fixtureStore(t *testing.T, client *bao.Client) *openbao.Credentials {
|
||||
t.Helper()
|
||||
store, err := openbao.NewCredentials(client, "secret", "applications")
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
return store
|
||||
}
|
||||
|
||||
func TestCredentialLocationScope(t *testing.T) {
|
||||
client := fixtureClient(t, "http://127.0.0.1:1")
|
||||
store := fixtureStore(t, client)
|
||||
path, err := store.ProvisionPath("database-uid")
|
||||
if err != nil || path != credentialPath {
|
||||
t.Fatal("unexpected stable location")
|
||||
}
|
||||
for _, path := range []string{"", "/absolute", "applications", "applications-other/key", "applications/../management", "applications/%2e%2e/key", "applications//key", "applications/data/key"} {
|
||||
if _, err := store.Read(t.Context(), path); !errors.Is(err, openbao.ErrInvalidLocation) {
|
||||
t.Fatal("accepted invalid location")
|
||||
}
|
||||
if err := store.Create(t.Context(), path, fixtureCredential(t)); !errors.Is(err, openbao.ErrInvalidLocation) {
|
||||
t.Fatal("accepted invalid create location")
|
||||
}
|
||||
}
|
||||
for _, uid := range []string{"", "../key", "a/b", "a?b"} {
|
||||
if _, err := store.ProvisionPath(uid); err == nil {
|
||||
t.Fatal("accepted invalid UID")
|
||||
}
|
||||
}
|
||||
for _, invalid := range []string{"", "data", "metadata", "../secret", "secret/", "secret?query"} {
|
||||
if _, err := openbao.NewCredentials(client, invalid, "applications"); err == nil {
|
||||
t.Fatal("accepted invalid mount")
|
||||
}
|
||||
if _, err := openbao.NewCredentials(client, "secret", invalid); err == nil {
|
||||
t.Fatal("accepted invalid base path")
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func TestCredentialWriteFailureDoesNotRetryOrLeak(t *testing.T) {
|
||||
var requests atomic.Int32
|
||||
server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, _ *http.Request) {
|
||||
requests.Add(1)
|
||||
http.Error(w, fixturePassword+fixtureToken, http.StatusInternalServerError)
|
||||
}))
|
||||
defer server.Close()
|
||||
store := fixtureStore(t, fixtureClient(t, server.URL))
|
||||
if err := store.Create(t.Context(), credentialPath, fixtureCredential(t)); err != openbao.ErrUncertain {
|
||||
t.Fatal("write error must be a redacted uncertain outcome")
|
||||
}
|
||||
if requests.Load() != 1 {
|
||||
t.Fatal("SDK retried an uncertain write")
|
||||
}
|
||||
if _, err := store.Read(t.Context(), credentialPath); err != openbao.ErrUnavailable {
|
||||
t.Fatal("read error must be redacted")
|
||||
}
|
||||
ctx, cancel := context.WithCancel(t.Context())
|
||||
cancel()
|
||||
if err := store.Create(ctx, credentialPath, fixtureCredential(t)); err != openbao.ErrUnavailable || requests.Load() != 2 {
|
||||
t.Fatal("canceled operation must not write")
|
||||
}
|
||||
}
|
||||
|
||||
func TestCredentialWriteDeniedBeforeExecution(t *testing.T) {
|
||||
for _, status := range []int{http.StatusUnauthorized, http.StatusForbidden} {
|
||||
server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, _ *http.Request) {
|
||||
http.Error(w, fixtureToken, status)
|
||||
}))
|
||||
store := fixtureStore(t, fixtureClient(t, server.URL))
|
||||
err := store.Create(t.Context(), credentialPath, fixtureCredential(t))
|
||||
server.Close()
|
||||
if err != openbao.ErrUnavailable {
|
||||
t.Fatalf("status %d: definite rejection should wait for dependency recovery, got %v", status, err)
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -31,6 +31,7 @@ import (
|
||||
corev1 "k8s.io/api/core/v1"
|
||||
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
|
||||
"k8s.io/client-go/kubernetes"
|
||||
"k8s.io/client-go/rest"
|
||||
"sigs.k8s.io/controller-runtime/pkg/envtest"
|
||||
|
||||
secretadapter "git.ddupan.top/panxiao81/ayatori/internal/database/adapter/kubernetes"
|
||||
@@ -40,15 +41,19 @@ import (
|
||||
)
|
||||
|
||||
const (
|
||||
fixtureHost = "fixture.invalid"
|
||||
fixtureUser = "postgres"
|
||||
fixtureExtension = "plpgsql"
|
||||
dockerExec = "exec"
|
||||
fixtureImage = "postgres@sha256:18cfe3ef5e6815560c98237d6216d1e5119702fb0f3894c8785dd58b8bbe5d73"
|
||||
fixturePassword = "AYATORI-TEST-ONLY-initial-password"
|
||||
rotatedPassword = "AYATORI-TEST-ONLY-rotated-password"
|
||||
controllerNamespace = "database-controller"
|
||||
secretName = "management"
|
||||
fixtureAddress = "127.0.0.1"
|
||||
managementUsernameKey = "login"
|
||||
managementPasswordKey = "credential"
|
||||
unrelatedNamespace = "unrelated"
|
||||
fixtureHost = "fixture.invalid"
|
||||
fixtureUser = "postgres"
|
||||
fixtureExtension = "plpgsql"
|
||||
dockerExec = "exec"
|
||||
fixtureImage = "postgres@sha256:18cfe3ef5e6815560c98237d6216d1e5119702fb0f3894c8785dd58b8bbe5d73"
|
||||
fixturePassword = "AYATORI-TEST-ONLY-initial-password"
|
||||
rotatedPassword = "AYATORI-TEST-ONLY-rotated-password"
|
||||
controllerNamespace = "database-controller"
|
||||
secretName = "management"
|
||||
)
|
||||
|
||||
// fixture 不接受外部 DSN,只创建自己的临时容器并按确切 ID 清理。
|
||||
@@ -79,7 +84,7 @@ func postgresFixture(t *testing.T, ctx context.Context) (string, int) {
|
||||
t.Fatal("invalid fixture port")
|
||||
}
|
||||
// 初次 init 的临时服务器只监听 Unix socket,必须等最终 TCP listener。
|
||||
for exec.CommandContext(ctx, "docker", dockerExec, id, "pg_isready", "-h", "127.0.0.1", "-U", fixtureUser).Run() != nil {
|
||||
for exec.CommandContext(ctx, "docker", dockerExec, id, "pg_isready", "-h", fixtureAddress, "-U", fixtureUser).Run() != nil {
|
||||
select {
|
||||
case <-ctx.Done():
|
||||
t.Fatal("fixture startup timed out")
|
||||
@@ -154,7 +159,7 @@ func target(t *testing.T, port int, mode instance.TLSMode) instance.ObservationT
|
||||
}
|
||||
endpoint, err := instance.NewEndpoint(instance.EndpointValues{
|
||||
Host: fixtureHost,
|
||||
HostAddr: "127.0.0.1",
|
||||
HostAddr: fixtureAddress,
|
||||
Port: port,
|
||||
ManagementDatabase: fixtureUser,
|
||||
TLSMode: mode,
|
||||
@@ -164,8 +169,8 @@ func target(t *testing.T, port int, mode instance.TLSMode) instance.ObservationT
|
||||
}
|
||||
ref, err := instance.NewCredentialReference(instance.CredentialReferenceValues{
|
||||
Name: secretName,
|
||||
UsernameKey: "login",
|
||||
PasswordKey: "credential",
|
||||
UsernameKey: managementUsernameKey,
|
||||
PasswordKey: managementPasswordKey,
|
||||
})
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
@@ -196,6 +201,7 @@ func (r *gatedReader) Read(ctx context.Context, ref instance.CredentialReference
|
||||
|
||||
// credentialFixture 为每个场景创建独立 API server、PostgreSQL 和应用服务。
|
||||
type credentialFixture struct {
|
||||
config *rest.Config
|
||||
ctx context.Context
|
||||
client *kubernetes.Clientset
|
||||
reader *secretadapter.SecretCredentials
|
||||
@@ -227,7 +233,7 @@ func newCredentialFixture(t *testing.T) *credentialFixture {
|
||||
if err != nil {
|
||||
t.Fatal("cannot create test client")
|
||||
}
|
||||
for _, namespace := range []string{controllerNamespace, "unrelated"} {
|
||||
for _, namespace := range []string{controllerNamespace, unrelatedNamespace} {
|
||||
_, err := client.CoreV1().Namespaces().Create(
|
||||
ctx,
|
||||
&corev1.Namespace{Name: namespace},
|
||||
@@ -260,6 +266,7 @@ func newCredentialFixture(t *testing.T) *credentialFixture {
|
||||
t.Cleanup(service.Close)
|
||||
|
||||
return &credentialFixture{
|
||||
config: config,
|
||||
ctx: ctx,
|
||||
client: client,
|
||||
reader: reader,
|
||||
@@ -277,8 +284,8 @@ func (f *credentialFixture) createSecret(t *testing.T, namespace string) {
|
||||
secret := &corev1.Secret{
|
||||
Name: secretName,
|
||||
Data: map[string][]byte{
|
||||
"login": []byte(fixtureUser),
|
||||
"credential": []byte(fixturePassword),
|
||||
managementUsernameKey: []byte(fixtureUser),
|
||||
managementPasswordKey: []byte(fixturePassword),
|
||||
},
|
||||
}
|
||||
if _, err := f.client.CoreV1().Secrets(namespace).Create(f.ctx, secret, metav1.CreateOptions{}); err != nil {
|
||||
|
||||
@@ -0,0 +1,210 @@
|
||||
//go:build integration
|
||||
|
||||
package postgresql_test
|
||||
|
||||
import (
|
||||
"bytes"
|
||||
"context"
|
||||
"os"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
databasev1alpha1 "git.ddupan.top/panxiao81/ayatori/api/database/v1alpha1"
|
||||
secretadapter "git.ddupan.top/panxiao81/ayatori/internal/database/adapter/kubernetes"
|
||||
"git.ddupan.top/panxiao81/ayatori/internal/database/adapter/postgresql"
|
||||
"git.ddupan.top/panxiao81/ayatori/internal/database/application"
|
||||
databasecontroller "git.ddupan.top/panxiao81/ayatori/internal/database/controller"
|
||||
corev1 "k8s.io/api/core/v1"
|
||||
rbacv1 "k8s.io/api/rbac/v1"
|
||||
apierrors "k8s.io/apimachinery/pkg/api/errors"
|
||||
"k8s.io/apimachinery/pkg/api/meta"
|
||||
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
|
||||
"k8s.io/apimachinery/pkg/runtime"
|
||||
yamlutil "k8s.io/apimachinery/pkg/util/yaml"
|
||||
"k8s.io/client-go/rest"
|
||||
ctrl "sigs.k8s.io/controller-runtime"
|
||||
"sigs.k8s.io/controller-runtime/pkg/client"
|
||||
controllerconfig "sigs.k8s.io/controller-runtime/pkg/config"
|
||||
"sigs.k8s.io/controller-runtime/pkg/envtest"
|
||||
metricsserver "sigs.k8s.io/controller-runtime/pkg/metrics/server"
|
||||
"sigs.k8s.io/yaml"
|
||||
)
|
||||
|
||||
const watchRevisionAnnotation = "test.ayatori/observation"
|
||||
|
||||
func TestInstanceControllerWithRealPostgreSQL(t *testing.T) {
|
||||
f := newCredentialFixture(t)
|
||||
if _, err := envtest.InstallCRDs(f.config, envtest.CRDInstallOptions{
|
||||
Paths: []string{"../../../../config/crd/bases"}, ErrorIfPathMissing: true,
|
||||
}); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
scheme := runtime.NewScheme()
|
||||
for _, install := range []func(*runtime.Scheme) error{databasev1alpha1.AddToScheme, corev1.AddToScheme, rbacv1.AddToScheme} {
|
||||
if err := install(scheme); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
}
|
||||
apiClient, err := client.New(f.config, client.Options{Scheme: scheme})
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
restricted := instanceControllerRBAC(t, f, apiClient)
|
||||
credentials, err := secretadapter.NewSecretCredentials(restricted, controllerNamespace)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
service, err := application.NewInstanceService(credentials, postgresql.Connector{})
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
// 同进程 -count 重复启动测试 manager;生产继续校验 controller 名称唯一。
|
||||
skipRepeatedName := true
|
||||
manager, err := ctrl.NewManager(restricted, ctrl.Options{
|
||||
Scheme: scheme, Cache: databasecontroller.InstanceCacheOptions(controllerNamespace),
|
||||
Metrics: metricsserver.Options{BindAddress: "0"}, HealthProbeBindAddress: "0",
|
||||
Controller: controllerconfig.Controller{SkipNameValidation: &skipRepeatedName},
|
||||
})
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
reconciler := &databasecontroller.InstanceReconciler{Observer: service, SecretNamespace: controllerNamespace}
|
||||
if err := reconciler.SetupWithManager(manager); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
managerContext, stop := context.WithCancel(f.ctx)
|
||||
done := make(chan error, 1)
|
||||
go func() { done <- manager.Start(managerContext) }()
|
||||
t.Cleanup(func() {
|
||||
stop()
|
||||
select {
|
||||
case err := <-done:
|
||||
if err != nil {
|
||||
t.Error(err)
|
||||
}
|
||||
case <-time.After(20 * time.Second):
|
||||
t.Error("Instance manager 未停止")
|
||||
}
|
||||
service.Close()
|
||||
})
|
||||
object := &databasev1alpha1.PostgreSQLInstance{}
|
||||
object.Name = "native-instance"
|
||||
object.Spec.Endpoint = databasev1alpha1.PostgreSQLEndpoint{
|
||||
Host: fixtureHost, HostAddr: fixtureAddress, Port: int32(f.port), SSLMode: "disable",
|
||||
}
|
||||
object.Spec.AdminCredentialRef = databasev1alpha1.AdminCredentialReference{
|
||||
Name: secretName, UsernameKey: managementUsernameKey, PasswordKey: managementPasswordKey,
|
||||
}
|
||||
if err := apiClient.Create(f.ctx, object); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
awaitInstanceReason(t, f, apiClient, object, "DependencyUnavailable")
|
||||
// 30 秒轮询前必须收到 Secret 创建事件;实际 controller 使用 namespace Role + metadata watch。
|
||||
useNativeManager(t, f)
|
||||
awaitInstanceReason(t, f, apiClient, object, "ManagementReady")
|
||||
before := f.backendIDs(t)
|
||||
f.updateSecret(t, func(secret *corev1.Secret) {
|
||||
secret.Annotations = map[string]string{watchRevisionAnnotation: "changed"}
|
||||
})
|
||||
// 用实际 API 事件触发重验,metadata 改动不应换池。
|
||||
time.Sleep(200 * time.Millisecond)
|
||||
if f.backendIDs(t) != before {
|
||||
t.Fatal("无关 Secret metadata 修改重建了连接")
|
||||
}
|
||||
f.queryPostgres(t, "ALTER ROLE native_manager PASSWORD '"+rotatedPassword+"'")
|
||||
f.updateSecret(t, func(secret *corev1.Secret) { secret.Data[managementPasswordKey] = []byte("invalid-test-password") })
|
||||
awaitInstanceReason(t, f, apiClient, object, "AuthenticationFailed")
|
||||
f.updateSecret(t, func(secret *corev1.Secret) { secret.Data[managementPasswordKey] = []byte(rotatedPassword) })
|
||||
awaitInstanceReason(t, f, apiClient, object, "ManagementReady")
|
||||
if f.backendIDs(t) == before {
|
||||
t.Fatal("凭据轮换没有替换旧连接")
|
||||
}
|
||||
f.queryPostgres(t, "ALTER ROLE native_manager NOCREATEROLE")
|
||||
f.updateSecret(t, func(secret *corev1.Secret) { secret.Annotations[watchRevisionAnnotation] = "recheck" })
|
||||
awaitInstanceReason(t, f, apiClient, object, "InsufficientPrivileges")
|
||||
f.queryPostgres(t, "ALTER ROLE native_manager CREATEROLE")
|
||||
f.updateSecret(t, func(secret *corev1.Secret) { secret.Annotations[watchRevisionAnnotation] = "recovered" })
|
||||
awaitInstanceReason(t, f, apiClient, object, "ManagementReady")
|
||||
if err := f.client.CoreV1().Secrets(controllerNamespace).Delete(f.ctx, secretName, metav1.DeleteOptions{}); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
awaitInstanceReason(t, f, apiClient, object, "DependencyUnavailable")
|
||||
if f.backendIDs(t) != "" {
|
||||
t.Fatal("Secret 删除后旧连接未释放")
|
||||
}
|
||||
}
|
||||
|
||||
func awaitInstanceReason(t *testing.T, f *credentialFixture, apiClient client.Client,
|
||||
object *databasev1alpha1.PostgreSQLInstance, reason string) {
|
||||
t.Helper()
|
||||
deadline := time.Now().Add(10 * time.Second)
|
||||
for time.Now().Before(deadline) {
|
||||
if err := apiClient.Get(f.ctx, client.ObjectKeyFromObject(object), object); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
condition := meta.FindStatusCondition(object.Status.Conditions, "Ready")
|
||||
if condition != nil && condition.Reason == reason && condition.ObservedGeneration == object.Generation {
|
||||
if (condition.Status == metav1.ConditionTrue) != (reason == "ManagementReady") {
|
||||
t.Fatal("Ready 与检查结果不一致")
|
||||
}
|
||||
return
|
||||
}
|
||||
time.Sleep(50 * time.Millisecond)
|
||||
}
|
||||
t.Fatalf("watch 未及时推进到 %s", reason)
|
||||
}
|
||||
|
||||
func instanceControllerRBAC(t *testing.T, f *credentialFixture, apiClient client.Client) *rest.Config {
|
||||
t.Helper()
|
||||
roleBytes, err := os.ReadFile("../../../../config/rbac/role.yaml")
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
role := &rbacv1.ClusterRole{}
|
||||
if err := yaml.Unmarshal(roleBytes, role); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if err := apiClient.Create(f.ctx, role); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
user := "instance-controller-test"
|
||||
binding := &rbacv1.ClusterRoleBinding{}
|
||||
binding.Name = user
|
||||
binding.RoleRef = rbacv1.RoleRef{APIGroup: rbacv1.GroupName, Kind: "ClusterRole", Name: role.Name}
|
||||
binding.Subjects = []rbacv1.Subject{{Kind: "User", APIGroup: rbacv1.GroupName, Name: user}}
|
||||
if err := apiClient.Create(f.ctx, binding); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
credentialBytes, err := os.ReadFile("../../../../config/rbac/database_credentials_role.yaml")
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
namespaceRole := &rbacv1.Role{}
|
||||
decoder := yamlutil.NewYAMLOrJSONDecoder(bytes.NewReader(credentialBytes), 4096)
|
||||
if err := decoder.Decode(namespaceRole); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
namespaceRole.Namespace = controllerNamespace
|
||||
if err := apiClient.Create(f.ctx, namespaceRole); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
namespaceBinding := &rbacv1.RoleBinding{}
|
||||
namespaceBinding.Name, namespaceBinding.Namespace = user, controllerNamespace
|
||||
namespaceBinding.RoleRef = rbacv1.RoleRef{APIGroup: rbacv1.GroupName, Kind: "Role", Name: namespaceRole.Name}
|
||||
namespaceBinding.Subjects = binding.Subjects
|
||||
if err := apiClient.Create(f.ctx, namespaceBinding); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
config := rest.CopyConfig(f.config)
|
||||
config.Impersonate.UserName = user
|
||||
restrictedClient, err := client.New(config, client.Options{Scheme: apiClient.Scheme()})
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
secret := &corev1.Secret{}
|
||||
err = restrictedClient.Get(f.ctx, client.ObjectKey{Namespace: unrelatedNamespace, Name: secretName}, secret)
|
||||
if !apierrors.IsForbidden(err) {
|
||||
t.Fatal("Instance controller 可以跨 namespace 读取 Secret")
|
||||
}
|
||||
return config
|
||||
}
|
||||
@@ -0,0 +1,58 @@
|
||||
package postgresql
|
||||
|
||||
import (
|
||||
"context"
|
||||
|
||||
"git.ddupan.top/panxiao81/ayatori/internal/database/application"
|
||||
"git.ddupan.top/panxiao81/ayatori/internal/database/domain/instance"
|
||||
)
|
||||
|
||||
// 只读取当前执行角色的属性,不能从可继承的角色成员关系推导 CREATEDB/CREATEROLE。
|
||||
// 所有事实来自同一条语句;不创建探针数据库,不修改管理账号或持久 schema。
|
||||
const inspectManagementStatement = `
|
||||
SELECT
|
||||
pg_catalog.current_setting('server_version'),
|
||||
ARRAY(SELECT name::text FROM pg_catalog.pg_available_extensions ORDER BY name),
|
||||
role.rolsuper,
|
||||
role.rolcreaterole,
|
||||
role.rolcreatedb,
|
||||
pg_catalog.pg_is_in_recovery() OR
|
||||
pg_catalog.current_setting('transaction_read_only')::boolean
|
||||
FROM pg_catalog.pg_roles AS role
|
||||
WHERE role.rolname = current_user`
|
||||
|
||||
func (d *database) InspectManagement(ctx context.Context) (application.DatabaseMetadata, error) {
|
||||
var metadata application.DatabaseMetadata
|
||||
var superuser, createRole, createDatabase, readOnly bool
|
||||
err := d.pool.QueryRow(ctx, inspectManagementStatement).Scan(
|
||||
&metadata.Version, &metadata.AvailableExtensions,
|
||||
&superuser, &createRole, &createDatabase, &readOnly,
|
||||
)
|
||||
if err != nil {
|
||||
return application.DatabaseMetadata{}, safeError(err, application.ErrObservation)
|
||||
}
|
||||
checks := instance.ManagementChecks{
|
||||
Connection: instance.CheckPassed,
|
||||
Metadata: instance.CheckPassed,
|
||||
Roles: nativePrivilege(createRole && !superuser),
|
||||
Databases: nativePrivilege(createDatabase && !superuser),
|
||||
// CREATEROLE 可管理自己新建角色的 membership;供应时必须显式取得 SET 权限,
|
||||
// 再以 owner 操作数据库 ACL。这里不授权操作任意导入角色或他人数据库。
|
||||
Grants: nativePrivilege(createRole && createDatabase && !superuser),
|
||||
// 新建数据库 owner 可安装 trusted 扩展。具体扩展仍需逐请求执行和回读,
|
||||
// 非 trusted 扩展不能因出现在 available 列表就视为可安装。
|
||||
Extensions: nativePrivilege(createRole && createDatabase && !superuser),
|
||||
}
|
||||
if readOnly {
|
||||
checks.Databases = instance.CheckUnavailable
|
||||
}
|
||||
metadata.Management = checks
|
||||
return metadata, nil
|
||||
}
|
||||
|
||||
func nativePrivilege(allowed bool) instance.CheckResult {
|
||||
if allowed {
|
||||
return instance.CheckPassed
|
||||
}
|
||||
return instance.CheckInsufficientPrivileges
|
||||
}
|
||||
@@ -0,0 +1,152 @@
|
||||
//go:build integration
|
||||
|
||||
package postgresql_test
|
||||
|
||||
import (
|
||||
"context"
|
||||
"errors"
|
||||
"testing"
|
||||
|
||||
"github.com/jackc/pgx/v5"
|
||||
corev1 "k8s.io/api/core/v1"
|
||||
|
||||
"git.ddupan.top/panxiao81/ayatori/internal/database/application"
|
||||
"git.ddupan.top/panxiao81/ayatori/internal/database/domain/instance"
|
||||
)
|
||||
|
||||
const nativeManager = "native_manager"
|
||||
|
||||
func useNativeManager(t *testing.T, f *credentialFixture) {
|
||||
t.Helper()
|
||||
f.queryPostgres(t, "CREATE ROLE native_manager LOGIN CREATEDB CREATEROLE PASSWORD '"+fixturePassword+"'")
|
||||
f.createSecret(t, controllerNamespace)
|
||||
f.updateSecret(t, func(secret *corev1.Secret) { secret.Data[managementUsernameKey] = []byte(nativeManager) })
|
||||
}
|
||||
|
||||
func assessManagement(t *testing.T, f *credentialFixture) instance.Snapshot {
|
||||
t.Helper()
|
||||
observation, err := f.service.ObserveManagement(f.ctx, f.target)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
aggregate, err := instance.Reconstitute(f.target, instance.Snapshot{}, false)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if err := aggregate.BeginValidation(); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
capabilities, err := observation.Capabilities()
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if err := aggregate.AssessManagement(capabilities); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
return aggregate.Snapshot()
|
||||
}
|
||||
|
||||
func TestNativeManagementPrivileges(t *testing.T) {
|
||||
f := newCredentialFixture(t)
|
||||
useNativeManager(t, f)
|
||||
if snapshot := assessManagement(t, f); snapshot.Readiness != instance.Ready {
|
||||
t.Fatal("原生非 superuser 管理账号未通过检查")
|
||||
}
|
||||
before := f.backendIDs(t)
|
||||
for _, attribute := range []string{"NOCREATEROLE", "NOCREATEDB"} {
|
||||
f.queryPostgres(t, "ALTER ROLE native_manager "+attribute)
|
||||
if snapshot := assessManagement(t, f); snapshot.Failure != instance.InsufficientPrivileges {
|
||||
t.Fatal("已有连接忽略了管理权限撤回")
|
||||
}
|
||||
f.queryPostgres(t, "ALTER ROLE native_manager CREATEROLE CREATEDB")
|
||||
if snapshot := assessManagement(t, f); snapshot.Readiness != instance.Ready {
|
||||
t.Fatal("管理权限恢复后无法重新就绪")
|
||||
}
|
||||
}
|
||||
if f.backendIDs(t) != before {
|
||||
t.Fatal("权限检查不应要求重建连接才生效")
|
||||
}
|
||||
f.queryPostgres(t, "ALTER ROLE native_manager SET default_transaction_read_only = on")
|
||||
f.service.Forget(f.target.Identity().Name())
|
||||
if snapshot := assessManagement(t, f); snapshot.Failure != instance.DependencyUnavailable {
|
||||
t.Fatal("只读会话不应标记可供应")
|
||||
}
|
||||
f.queryPostgres(t, "ALTER ROLE native_manager RESET default_transaction_read_only")
|
||||
f.service.Forget(f.target.Identity().Name())
|
||||
if snapshot := assessManagement(t, f); snapshot.Readiness != instance.Ready {
|
||||
t.Fatal("恢复可写会话后没有就绪")
|
||||
}
|
||||
f.updateSecret(t, func(secret *corev1.Secret) { secret.Data[managementUsernameKey] = []byte(fixtureUser) })
|
||||
if snapshot := assessManagement(t, f); snapshot.Failure != instance.InsufficientPrivileges {
|
||||
t.Fatal("不应以 superuser 绕过非特权账号合同")
|
||||
}
|
||||
reads := 0
|
||||
f.gate.beforeRead = func() {
|
||||
reads++
|
||||
if reads == 2 {
|
||||
f.updateSecret(t, func(secret *corev1.Secret) { secret.Data[managementPasswordKey] = []byte(rotatedPassword) })
|
||||
}
|
||||
}
|
||||
observation, err := f.service.ObserveManagement(f.ctx, f.target)
|
||||
if !errors.Is(err, application.ErrCredentialsChanged) || observation.Target().Validate() == nil {
|
||||
t.Fatal("管理观察期间凭据轮换应丢弃全部能力结果")
|
||||
}
|
||||
if f.backendIDs(t) != "" {
|
||||
t.Fatal("中途轮换后不应保留旧管理连接")
|
||||
}
|
||||
}
|
||||
|
||||
// 以实际非 superuser 会话验证能力矩阵的依据,不用超级用户执行 SQL 模拟管理账号。
|
||||
// 这些固定名称只存在于本测试独占容器,生产观察本身不会创建探针对象。
|
||||
func TestNativeManagementSupplyContract(t *testing.T) {
|
||||
f := newCredentialFixture(t)
|
||||
useNativeManager(t, f)
|
||||
config, err := pgx.ParseConfig("")
|
||||
if err != nil {
|
||||
t.Fatal("无法装配隔离测试连接")
|
||||
}
|
||||
config.Host, config.Port = fixtureAddress, uint16(f.port)
|
||||
config.Database, config.User, config.Password = fixtureUser, nativeManager, fixturePassword
|
||||
config.TLSConfig, config.Fallbacks = nil, nil
|
||||
connection, err := pgx.ConnectConfig(f.ctx, config)
|
||||
if err != nil {
|
||||
t.Fatal("非 superuser 测试连接失败")
|
||||
}
|
||||
t.Cleanup(func() { _ = connection.Close(context.Background()) })
|
||||
execute := func(statement string) {
|
||||
t.Helper()
|
||||
if _, err := connection.Exec(f.ctx, statement); err != nil {
|
||||
t.Fatalf("原生管理能力合同未满足,步骤 %q", statement)
|
||||
}
|
||||
}
|
||||
execute("CREATE ROLE managed_owner LOGIN NOSUPERUSER NOCREATEDB NOCREATEROLE NOREPLICATION")
|
||||
execute("GRANT managed_owner TO native_manager WITH SET TRUE")
|
||||
execute("CREATE DATABASE managed_database OWNER managed_owner")
|
||||
execute("SET ROLE managed_owner")
|
||||
execute("REVOKE CONNECT ON DATABASE managed_database FROM PUBLIC")
|
||||
execute("GRANT CONNECT ON DATABASE managed_database TO managed_owner")
|
||||
execute("RESET ROLE")
|
||||
config.Database = "managed_database"
|
||||
tenantConnection, err := pgx.ConnectConfig(f.ctx, config)
|
||||
if err != nil {
|
||||
t.Fatal("管理账号无法访问其受管数据库")
|
||||
}
|
||||
defer func() { _ = tenantConnection.Close(context.Background()) }()
|
||||
if _, err := tenantConnection.Exec(f.ctx, "SET ROLE managed_owner; CREATE EXTENSION hstore"); err != nil {
|
||||
t.Fatal("owner 无法安装 trusted 扩展")
|
||||
}
|
||||
var installed bool
|
||||
if err := tenantConnection.QueryRow(f.ctx, "SELECT EXISTS (SELECT FROM pg_catalog.pg_extension WHERE extname = 'hstore')").Scan(&installed); err != nil || !installed {
|
||||
t.Fatal("扩展安装后实际回读失败")
|
||||
}
|
||||
if _, err := tenantConnection.Exec(f.ctx, "CREATE EXTENSION file_fdw"); err == nil {
|
||||
t.Fatal("非 trusted 扩展不应被 Ready 隐式授权")
|
||||
}
|
||||
if err := tenantConnection.Close(f.ctx); err != nil {
|
||||
t.Fatal("关闭目标数据库连接失败")
|
||||
}
|
||||
execute("SET ROLE managed_owner")
|
||||
execute("DROP DATABASE managed_database")
|
||||
execute("RESET ROLE")
|
||||
execute("DROP ROLE managed_owner")
|
||||
}
|
||||
@@ -0,0 +1,116 @@
|
||||
/*
|
||||
Copyright 2026.
|
||||
|
||||
Licensed under the Apache License, Version 2.0 (the "License");
|
||||
you may not use this file except in compliance with the License.
|
||||
You may obtain a copy of the License at
|
||||
|
||||
http://www.apache.org/licenses/LICENSE-2.0
|
||||
|
||||
Unless required by applicable law or agreed to in writing, software
|
||||
distributed under the License is distributed on an "AS IS" BASIS,
|
||||
WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
|
||||
See the License for the specific language governing permissions and
|
||||
limitations under the License.
|
||||
*/
|
||||
|
||||
package application
|
||||
|
||||
import (
|
||||
"crypto/rand"
|
||||
"encoding/base64"
|
||||
"errors"
|
||||
"regexp"
|
||||
"strconv"
|
||||
|
||||
"git.ddupan.top/panxiao81/ayatori/internal/database/domain/instance"
|
||||
)
|
||||
|
||||
var ErrApplicationCredentialInvalid = errors.New("application credential is invalid")
|
||||
|
||||
var applicationIdentifier = regexp.MustCompile(`^[a-z][a-z0-9_]{0,62}$`)
|
||||
|
||||
// ApplicationCredential 是内存中的应用连接凭据,不得放入 CR 或普通日志。
|
||||
// 它与 Instance 管理凭据分开,固定输出交付合同中的七键,不生成带密码的 URI。
|
||||
type ApplicationCredential struct {
|
||||
username string
|
||||
password string
|
||||
database string
|
||||
endpoint instance.Endpoint
|
||||
}
|
||||
|
||||
func NewApplicationCredential(username, password, database string, endpoint instance.Endpoint) (ApplicationCredential, error) {
|
||||
if !applicationIdentifier.MatchString(username) || !applicationIdentifier.MatchString(database) || password == "" {
|
||||
return ApplicationCredential{}, ErrApplicationCredentialInvalid
|
||||
}
|
||||
if endpoint.Validate() != nil {
|
||||
return ApplicationCredential{}, ErrApplicationCredentialInvalid
|
||||
}
|
||||
return ApplicationCredential{
|
||||
username: username,
|
||||
password: password,
|
||||
database: database,
|
||||
endpoint: endpoint,
|
||||
}, nil
|
||||
}
|
||||
|
||||
// GenerateApplicationCredential 仅供已获准首次创建凭据的供应步骤调用。
|
||||
// 不能在读取失败、写入结果不确定或重启后无条件重新调用。
|
||||
func GenerateApplicationCredential(username, database string, endpoint instance.Endpoint) (ApplicationCredential, error) {
|
||||
password := make([]byte, 32)
|
||||
rand.Read(password)
|
||||
return NewApplicationCredential(username, base64.RawURLEncoding.EncodeToString(password), database, endpoint)
|
||||
}
|
||||
|
||||
func (c ApplicationCredential) String() string { return "[redacted application credential]" }
|
||||
func (c ApplicationCredential) GoString() string { return c.String() }
|
||||
func (c ApplicationCredential) MarshalJSON() ([]byte, error) {
|
||||
return []byte(`"[redacted application credential]"`), nil
|
||||
}
|
||||
|
||||
// SecretData 只在凭据后端或数据库连接边界使用;返回值包含明文密码,禁止记录日志。
|
||||
// 每次返回独立 map,调用方不能修改已经构造的凭据。
|
||||
func (c ApplicationCredential) SecretData() map[string]any {
|
||||
endpoint := c.endpoint.Values()
|
||||
return map[string]any{
|
||||
"username": c.username,
|
||||
"password": c.password,
|
||||
"database": c.database,
|
||||
"host": endpoint.Host,
|
||||
"hostaddr": endpoint.HostAddr,
|
||||
"port": strconv.Itoa(endpoint.Port),
|
||||
"sslmode": string(endpoint.TLSMode),
|
||||
}
|
||||
}
|
||||
|
||||
func (c ApplicationCredential) Validate() error {
|
||||
_, err := NewApplicationCredential(c.username, c.password, c.database, c.endpoint)
|
||||
return err
|
||||
}
|
||||
|
||||
// ParseApplicationCredential 拒绝缺键、非字符串或非法连接参数,不回显后端内容。
|
||||
func ParseApplicationCredential(data map[string]any) (ApplicationCredential, error) {
|
||||
values := make(map[string]string, 7)
|
||||
for _, key := range []string{"username", "password", "database", "host", "hostaddr", "port", "sslmode"} {
|
||||
value, ok := data[key].(string)
|
||||
if !ok || value == "" {
|
||||
return ApplicationCredential{}, ErrApplicationCredentialInvalid
|
||||
}
|
||||
values[key] = value
|
||||
}
|
||||
port, err := strconv.Atoi(values["port"])
|
||||
if err != nil {
|
||||
return ApplicationCredential{}, ErrApplicationCredentialInvalid
|
||||
}
|
||||
endpoint, err := instance.NewEndpoint(instance.EndpointValues{
|
||||
Host: values["host"],
|
||||
HostAddr: values["hostaddr"],
|
||||
Port: port,
|
||||
ManagementDatabase: values["database"],
|
||||
TLSMode: instance.TLSMode(values["sslmode"]),
|
||||
})
|
||||
if err != nil {
|
||||
return ApplicationCredential{}, ErrApplicationCredentialInvalid
|
||||
}
|
||||
return NewApplicationCredential(values["username"], values["password"], values["database"], endpoint)
|
||||
}
|
||||
@@ -0,0 +1,81 @@
|
||||
/*
|
||||
Copyright 2026.
|
||||
|
||||
Licensed under the Apache License, Version 2.0 (the "License");
|
||||
you may not use this file except in compliance with the License.
|
||||
You may obtain a copy of the License at
|
||||
|
||||
http://www.apache.org/licenses/LICENSE-2.0
|
||||
|
||||
Unless required by applicable law or agreed to in writing, software
|
||||
distributed under the License is distributed on an "AS IS" BASIS,
|
||||
WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
|
||||
See the License for the specific language governing permissions and
|
||||
limitations under the License.
|
||||
*/
|
||||
|
||||
package application_test
|
||||
|
||||
import (
|
||||
"encoding/json"
|
||||
"fmt"
|
||||
"maps"
|
||||
"strings"
|
||||
"testing"
|
||||
|
||||
"git.ddupan.top/panxiao81/ayatori/internal/database/application"
|
||||
"git.ddupan.top/panxiao81/ayatori/internal/database/domain/instance"
|
||||
)
|
||||
|
||||
func TestApplicationCredential(t *testing.T) {
|
||||
endpoint, err := instance.NewEndpoint(instance.EndpointValues{
|
||||
Host: "postgres.example", HostAddr: "192.0.2.1", Port: 5432,
|
||||
ManagementDatabase: "postgres", TLSMode: instance.TLSVerifyFull,
|
||||
})
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
first, err := application.GenerateApplicationCredential("owner", "app", endpoint)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
second, err := application.GenerateApplicationCredential("owner", "app", endpoint)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
data := first.SecretData()
|
||||
if len(data) != 7 || data["password"] == second.SecretData()["password"] || len(data["password"].(string)) != 43 {
|
||||
t.Fatal("expected seven keys and independent 256-bit passwords")
|
||||
}
|
||||
parsed, err := application.ParseApplicationCredential(data)
|
||||
if err != nil || !maps.Equal(parsed.SecretData(), data) {
|
||||
t.Fatal("credential did not round trip")
|
||||
}
|
||||
encoded, err := json.Marshal(first)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
for _, output := range []string{fmt.Sprint(first), fmt.Sprintf("%+v", first), fmt.Sprintf("%#v", first), string(encoded)} {
|
||||
if strings.Contains(output, data["password"].(string)) {
|
||||
t.Fatal("credential formatting leaked the password")
|
||||
}
|
||||
}
|
||||
data["password"] = "changed"
|
||||
if first.SecretData()["password"] == "changed" {
|
||||
t.Fatal("caller mutated credential")
|
||||
}
|
||||
for key := range data {
|
||||
invalid := maps.Clone(data)
|
||||
delete(invalid, key)
|
||||
if _, err := application.ParseApplicationCredential(invalid); err == nil {
|
||||
t.Fatalf("accepted missing %s", key)
|
||||
}
|
||||
invalid[key] = 42
|
||||
if _, err := application.ParseApplicationCredential(invalid); err == nil {
|
||||
t.Fatalf("accepted non-string %s", key)
|
||||
}
|
||||
}
|
||||
if (application.ApplicationCredential{}).Validate() == nil {
|
||||
t.Fatal("accepted zero credential")
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,149 @@
|
||||
package application
|
||||
|
||||
import (
|
||||
"context"
|
||||
"errors"
|
||||
|
||||
"git.ddupan.top/panxiao81/ayatori/internal/database/domain/instance"
|
||||
)
|
||||
|
||||
const (
|
||||
instanceDependencyUnavailable = "DependencyUnavailable"
|
||||
instanceAuthenticationFailed = "AuthenticationFailed"
|
||||
)
|
||||
|
||||
// InstanceRecord 是 API 快照;Revision 仅用于持久化并发保护,不是领域版本。
|
||||
type InstanceRecord struct {
|
||||
Target instance.ObservationTarget
|
||||
Revision string
|
||||
Deleting bool
|
||||
}
|
||||
|
||||
type InstanceResources interface {
|
||||
LoadInstance(context.Context, string) (*InstanceRecord, error)
|
||||
ProtectInstance(context.Context, *InstanceRecord) (*InstanceRecord, error)
|
||||
// InstanceReferences 返回一个可定位的阻塞引用;空字符串表示没有引用。
|
||||
InstanceReferences(context.Context, string) (string, error)
|
||||
}
|
||||
|
||||
type InstanceObserver interface {
|
||||
ObserveManagement(context.Context, instance.ObservationTarget) (InstanceObservation, error)
|
||||
Forget(string)
|
||||
}
|
||||
|
||||
type InstanceResult struct {
|
||||
Record *InstanceRecord
|
||||
Snapshot instance.Snapshot
|
||||
Reason string
|
||||
Message string
|
||||
RemoveProtection bool
|
||||
}
|
||||
|
||||
// InstanceReconciliation 协调 API 保护、实时观察和领域判断,不拼装 Kubernetes status。
|
||||
type InstanceReconciliation struct {
|
||||
Resources InstanceResources
|
||||
Observer InstanceObserver
|
||||
}
|
||||
|
||||
func (s *InstanceReconciliation) Reconcile(ctx context.Context, name string) (InstanceResult, error) {
|
||||
record, err := s.Resources.LoadInstance(ctx, name)
|
||||
if err != nil {
|
||||
return InstanceResult{}, err
|
||||
}
|
||||
if record == nil {
|
||||
s.Observer.Forget(name)
|
||||
return InstanceResult{}, nil
|
||||
}
|
||||
if record.Deleting {
|
||||
return s.deleting(ctx, record)
|
||||
}
|
||||
record, err = s.Resources.ProtectInstance(ctx, record)
|
||||
if err != nil {
|
||||
return InstanceResult{}, err
|
||||
}
|
||||
// 每轮从无证据的领域对象开始;持久化 Ready 和连接存活不能替代本轮检查。
|
||||
aggregate, err := instance.Reconstitute(record.Target, instance.Snapshot{}, false)
|
||||
if err != nil {
|
||||
return InstanceResult{}, err
|
||||
}
|
||||
if err := aggregate.BeginValidation(); err != nil {
|
||||
return InstanceResult{}, err
|
||||
}
|
||||
observation, observationErr := s.Observer.ObserveManagement(ctx, record.Target)
|
||||
result := InstanceResult{Record: record}
|
||||
if observationErr != nil {
|
||||
result.Snapshot = aggregate.Snapshot()
|
||||
result.Snapshot.Readiness = instance.NotReady
|
||||
result.Snapshot.ObservedRevision = record.Target.Revision().Value()
|
||||
result.Reason, result.Message = observationFailure(observationErr)
|
||||
return result, nil
|
||||
}
|
||||
capabilities, err := observation.Capabilities()
|
||||
if err != nil {
|
||||
return InstanceResult{}, err
|
||||
}
|
||||
if err := aggregate.AssessManagement(capabilities); err != nil {
|
||||
return InstanceResult{}, err
|
||||
}
|
||||
result.Snapshot = aggregate.Snapshot()
|
||||
result.Snapshot.ReportedVersion = observation.Version()
|
||||
result.Reason, result.Message = managementResult(result.Snapshot.Failure)
|
||||
return result, nil
|
||||
}
|
||||
|
||||
func (s *InstanceReconciliation) deleting(ctx context.Context, record *InstanceRecord) (InstanceResult, error) {
|
||||
name := record.Target.Identity().Name()
|
||||
s.Observer.Forget(name)
|
||||
aggregate, err := instance.Reconstitute(record.Target, instance.Snapshot{}, true)
|
||||
if err != nil {
|
||||
return InstanceResult{}, err
|
||||
}
|
||||
if err := aggregate.BeginDeletion(); err != nil {
|
||||
return InstanceResult{}, err
|
||||
}
|
||||
result := InstanceResult{Record: record, Snapshot: aggregate.Snapshot(), Reason: "Deleting"}
|
||||
reference, err := s.Resources.InstanceReferences(ctx, name)
|
||||
if err != nil {
|
||||
result.Reason = instanceDependencyUnavailable
|
||||
result.Message = "无法确认 Database/Tenant 引用已解除;保留 Instance 删除保护并重试"
|
||||
return result, nil
|
||||
}
|
||||
if reference != "" {
|
||||
result.Reason = "InstanceInUse"
|
||||
result.Message = "仍被 " + reference + " 引用;先处理该资源,不会级联删除外部数据库"
|
||||
return result, nil
|
||||
}
|
||||
result.Message = "引用已解除,仅移除登记保护;不删除 PostgreSQL 或凭据"
|
||||
result.RemoveProtection = true
|
||||
return result, nil
|
||||
}
|
||||
|
||||
func observationFailure(err error) (string, string) {
|
||||
switch {
|
||||
case errors.Is(err, ErrAuthentication):
|
||||
return instanceAuthenticationFailed, "管理连接认证或 TLS 校验失败;检查管理 Secret 和 CA/证书配置"
|
||||
case errors.Is(err, ErrCredentialsInvalid):
|
||||
return "InvalidCredentials", "管理 Secret 的用户名或密码字段缺失;检查引用字段映射"
|
||||
case errors.Is(err, ErrCredentialsChanged):
|
||||
return "CredentialsChanged", "观察期间管理凭据变化,已丢弃结果并关闭旧连接;等待重新验证"
|
||||
case errors.Is(err, ErrCredentialsUnavailable):
|
||||
return instanceDependencyUnavailable, "无法读取管理 Secret;检查其是否存在及 controller namespace 内的读取权限"
|
||||
default:
|
||||
return instanceDependencyUnavailable, "管理连接或能力查询失败;检查 PostgreSQL 可达性、catalog 读取权限和超时"
|
||||
}
|
||||
}
|
||||
|
||||
func managementResult(failure instance.Failure) (string, string) {
|
||||
switch failure {
|
||||
case instance.NoFailure:
|
||||
return "ManagementReady", "当前管理能力检查通过;具体资源授权和扩展安装仍需执行时验证"
|
||||
case instance.InsufficientPrivileges:
|
||||
return "InsufficientPrivileges", "原生管理要求非 superuser 且具备 CREATEDB/CREATEROLE;不会自动修改账号权限"
|
||||
case instance.DependencyUnavailable:
|
||||
return instanceDependencyUnavailable, "当前 PostgreSQL 不可写或所需管理能力暂不可用"
|
||||
case instance.AuthenticationFailed:
|
||||
return instanceAuthenticationFailed, "当前管理能力检查未通过认证"
|
||||
default:
|
||||
return "ObservationIncomplete", "管理能力检查尚有缺项,不能仅凭 metadata 查询成功标记 Ready"
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,79 @@
|
||||
package application
|
||||
|
||||
import (
|
||||
"context"
|
||||
"errors"
|
||||
"testing"
|
||||
|
||||
"git.ddupan.top/panxiao81/ayatori/internal/database/domain/instance"
|
||||
)
|
||||
|
||||
func TestInstanceFailurePresentation(t *testing.T) {
|
||||
cases := []struct {
|
||||
err error
|
||||
reason string
|
||||
}{
|
||||
{ErrAuthentication, "AuthenticationFailed"},
|
||||
{ErrCredentialsInvalid, "InvalidCredentials"},
|
||||
{ErrCredentialsChanged, "CredentialsChanged"},
|
||||
{ErrCredentialsUnavailable, instanceDependencyUnavailable},
|
||||
{context.DeadlineExceeded, instanceDependencyUnavailable},
|
||||
{errors.New("private backend detail"), instanceDependencyUnavailable},
|
||||
}
|
||||
for _, test := range cases {
|
||||
reason, message := observationFailure(test.err)
|
||||
if reason != test.reason || message == "" || message == test.err.Error() {
|
||||
t.Fatal("观察失败没有安全且可诊断的状态")
|
||||
}
|
||||
}
|
||||
for _, failure := range []instance.Failure{
|
||||
instance.NoFailure, instance.ObservationIncomplete, instance.DependencyUnavailable,
|
||||
instance.AuthenticationFailed, instance.InsufficientPrivileges,
|
||||
} {
|
||||
reason, message := managementResult(failure)
|
||||
if reason == "" || message == "" || (reason == "ManagementReady") != (failure == instance.NoFailure) {
|
||||
t.Fatal("领域能力判定与状态不一致")
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func TestMetadataCannotEstablishManagementReadiness(t *testing.T) {
|
||||
source := &sourceStub{}
|
||||
source.credentials, _ = NewCredentials("test", serviceTestPassword)
|
||||
connector := &connectorStub{}
|
||||
service, err := NewInstanceService(source, connector)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
defer service.Close()
|
||||
target := serviceTarget(t, "uid", "postgres.test", "management", 1)
|
||||
if _, err := service.ObserveMetadata(t.Context(), target); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
connector.databases[0].metadata.Management = instance.ManagementChecks{
|
||||
Connection: instance.CheckPassed, Metadata: instance.CheckPassed,
|
||||
Roles: instance.CheckPassed, Databases: instance.CheckPassed,
|
||||
Grants: instance.CheckPassed, Extensions: instance.CheckPassed,
|
||||
}
|
||||
observation, err := service.ObserveMetadata(t.Context(), target)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
capabilities, err := observation.Capabilities()
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
aggregate, err := instance.Reconstitute(target, instance.Snapshot{}, false)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if err := aggregate.BeginValidation(); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if err := aggregate.AssessManagement(capabilities); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if aggregate.Snapshot().Failure != instance.ObservationIncomplete {
|
||||
t.Fatal("metadata 入口不应携带完整管理检查")
|
||||
}
|
||||
}
|
||||
@@ -36,6 +36,7 @@ var (
|
||||
// Metadata 只查询版本与可用扩展,不能产生领域 Ready。
|
||||
type Database interface {
|
||||
InspectMetadata(context.Context) (DatabaseMetadata, error)
|
||||
InspectManagement(context.Context) (DatabaseMetadata, error)
|
||||
Close()
|
||||
}
|
||||
|
||||
@@ -82,17 +83,26 @@ func (s *InstanceService) ObserveVersion(ctx context.Context, target instance.Ob
|
||||
|
||||
// ObserveMetadata 返回当前目标和凭据下的版本与扩展;任何失败均丢弃全部结果。
|
||||
// 调用者仍需使用 CR resourceVersion 保存前提防止 spec 并发修改;本方法不建立跨系统事务。
|
||||
func (s *InstanceService) ObserveMetadata(ctx context.Context, target instance.ObservationTarget) (MetadataObservation, error) {
|
||||
func (s *InstanceService) ObserveMetadata(ctx context.Context, target instance.ObservationTarget) (InstanceObservation, error) {
|
||||
return s.observe(ctx, target, false)
|
||||
}
|
||||
|
||||
// ObserveManagement 复用同一凭据刷新与回读边界,但每轮重新检查原生管理能力。
|
||||
func (s *InstanceService) ObserveManagement(ctx context.Context, target instance.ObservationTarget) (InstanceObservation, error) {
|
||||
return s.observe(ctx, target, true)
|
||||
}
|
||||
|
||||
func (s *InstanceService) observe(ctx context.Context, target instance.ObservationTarget, management bool) (InstanceObservation, error) {
|
||||
if err := target.Validate(); err != nil {
|
||||
return MetadataObservation{}, err
|
||||
return InstanceObservation{}, err
|
||||
}
|
||||
s.mu.Lock()
|
||||
defer s.mu.Unlock()
|
||||
if s.closed {
|
||||
return MetadataObservation{}, ErrClosed
|
||||
return InstanceObservation{}, ErrClosed
|
||||
}
|
||||
if err := ctx.Err(); err != nil {
|
||||
return MetadataObservation{}, err
|
||||
return InstanceObservation{}, err
|
||||
}
|
||||
|
||||
// 先读取有效凭据。读取失败时不得继续使用缓存中的旧连接。
|
||||
@@ -100,11 +110,11 @@ func (s *InstanceService) ObserveMetadata(ctx context.Context, target instance.O
|
||||
credentials, err := s.source.Read(ctx, target.Definition().AdminCredential())
|
||||
if err != nil {
|
||||
s.release(name)
|
||||
return MetadataObservation{}, credentialError(err)
|
||||
return InstanceObservation{}, credentialError(err)
|
||||
}
|
||||
if credentials.username == "" || credentials.password == "" {
|
||||
s.release(name)
|
||||
return MetadataObservation{}, ErrCredentialsInvalid
|
||||
return InstanceObservation{}, ErrCredentialsInvalid
|
||||
}
|
||||
|
||||
// 连接身份与有效值均未变化时复用 pgxpool;generation 本身不要求换池。
|
||||
@@ -117,7 +127,7 @@ func (s *InstanceService) ObserveMetadata(ctx context.Context, target instance.O
|
||||
if current == nil {
|
||||
database, err := s.connector.Connect(ctx, target.Definition().Endpoint(), credentials)
|
||||
if err != nil {
|
||||
return MetadataObservation{}, err
|
||||
return InstanceObservation{}, err
|
||||
}
|
||||
current = &entry{
|
||||
target: target,
|
||||
@@ -127,30 +137,38 @@ func (s *InstanceService) ObserveMetadata(ctx context.Context, target instance.O
|
||||
s.entries[name] = current
|
||||
}
|
||||
|
||||
metadata, err := current.database.InspectMetadata(ctx)
|
||||
var metadata DatabaseMetadata
|
||||
if management {
|
||||
metadata, err = current.database.InspectManagement(ctx)
|
||||
} else {
|
||||
metadata, err = current.database.InspectMetadata(ctx)
|
||||
// 即使 adapter 误填权限,也不能把只读 metadata 入口升级为 Ready。
|
||||
metadata.Management = instance.ManagementChecks{}
|
||||
}
|
||||
if err != nil {
|
||||
s.release(name)
|
||||
return MetadataObservation{}, err
|
||||
return InstanceObservation{}, err
|
||||
}
|
||||
if metadata.Version == "" {
|
||||
s.release(name)
|
||||
return MetadataObservation{}, ErrObservation
|
||||
return InstanceObservation{}, ErrObservation
|
||||
}
|
||||
|
||||
// 回读后再检查凭据,避免把轮换前取得的结果交给新凭据的调用链。
|
||||
latest, err := s.source.Read(ctx, target.Definition().AdminCredential())
|
||||
if err != nil {
|
||||
s.release(name)
|
||||
return MetadataObservation{}, credentialError(err)
|
||||
return InstanceObservation{}, credentialError(err)
|
||||
}
|
||||
if latest != credentials {
|
||||
s.release(name)
|
||||
return MetadataObservation{}, ErrCredentialsChanged
|
||||
return InstanceObservation{}, ErrCredentialsChanged
|
||||
}
|
||||
return MetadataObservation{
|
||||
return InstanceObservation{
|
||||
target: target,
|
||||
version: metadata.Version,
|
||||
extensions: instance.ObserveExtensionSupport(metadata.AvailableExtensions),
|
||||
management: metadata.Management,
|
||||
}, nil
|
||||
}
|
||||
|
||||
|
||||
@@ -42,6 +42,10 @@ type databaseStub struct {
|
||||
metadata DatabaseMetadata
|
||||
}
|
||||
|
||||
func (d *databaseStub) InspectManagement(ctx context.Context) (DatabaseMetadata, error) {
|
||||
return d.InspectMetadata(ctx)
|
||||
}
|
||||
|
||||
func (d *databaseStub) InspectMetadata(context.Context) (DatabaseMetadata, error) {
|
||||
return d.metadata, d.err
|
||||
}
|
||||
|
||||
@@ -18,23 +18,31 @@ package application
|
||||
|
||||
import "git.ddupan.top/panxiao81/ayatori/internal/database/domain/instance"
|
||||
|
||||
// DatabaseMetadata 是一次只读查询的事实,不包含管理权限或完整就绪结论。
|
||||
// DatabaseMetadata 是一次只读查询的事实,可附带原生管理检查,但不包含就绪结论。
|
||||
// AvailableExtensions 是服务器提供的可用列表,不是已安装列表或安装授权。
|
||||
type DatabaseMetadata struct {
|
||||
Version string
|
||||
AvailableExtensions []string
|
||||
// Management 仅由 InspectManagement 填充;metadata 查询必须保持未观察。
|
||||
Management instance.ManagementChecks
|
||||
}
|
||||
|
||||
// MetadataObservation 只在查询成功且有效凭据再次核对一致后产生。
|
||||
// InstanceObservation 只在查询成功且有效凭据再次核对一致后产生。
|
||||
// target 绑定本次调用,而非连接最初创建时的 generation;零值表示没有观察。
|
||||
type MetadataObservation struct {
|
||||
type InstanceObservation struct {
|
||||
target instance.ObservationTarget
|
||||
version string
|
||||
extensions instance.ExtensionSupport
|
||||
management instance.ManagementChecks
|
||||
}
|
||||
|
||||
func (o MetadataObservation) Target() instance.ObservationTarget { return o.target }
|
||||
func (o MetadataObservation) Version() string { return o.version }
|
||||
func (o MetadataObservation) Extensions() instance.ExtensionSupport {
|
||||
// Capabilities 保留缺项为未观察;不能从 metadata 的成功补齐管理检查。
|
||||
func (o InstanceObservation) Capabilities() (instance.CapabilityObservation, error) {
|
||||
return instance.NewCapabilityObservation(o.target, o.version, o.management)
|
||||
}
|
||||
|
||||
func (o InstanceObservation) Target() instance.ObservationTarget { return o.target }
|
||||
func (o InstanceObservation) Version() string { return o.version }
|
||||
func (o InstanceObservation) Extensions() instance.ExtensionSupport {
|
||||
return o.extensions
|
||||
}
|
||||
|
||||
@@ -0,0 +1,42 @@
|
||||
package controller
|
||||
|
||||
import (
|
||||
"context"
|
||||
"time"
|
||||
|
||||
"git.ddupan.top/panxiao81/ayatori/internal/database/adapter/kubernetes"
|
||||
"git.ddupan.top/panxiao81/ayatori/internal/database/application"
|
||||
ctrl "sigs.k8s.io/controller-runtime"
|
||||
"sigs.k8s.io/controller-runtime/pkg/client"
|
||||
)
|
||||
|
||||
type InstanceReconciler struct {
|
||||
Client client.Client
|
||||
Reader client.Reader
|
||||
Observer application.InstanceObserver
|
||||
SecretNamespace string
|
||||
}
|
||||
|
||||
// +kubebuilder:rbac:groups=database.ayatori.ddupan.top,resources=postgresqlinstances,verbs=get;list;watch;update;patch
|
||||
// +kubebuilder:rbac:groups=database.ayatori.ddupan.top,resources=postgresqlinstances/status,verbs=get;update;patch
|
||||
// +kubebuilder:rbac:groups=database.ayatori.ddupan.top,resources=postgresqlinstances/finalizers,verbs=update
|
||||
// Secret 权限单独声明为 namespace Role,不放入生成的 ClusterRole。
|
||||
|
||||
func (r *InstanceReconciler) Reconcile(ctx context.Context, request ctrl.Request) (ctrl.Result, error) {
|
||||
resources := &kubernetes.InstanceResources{Client: r.Client, Reader: r.Reader}
|
||||
service := application.InstanceReconciliation{Resources: resources, Observer: r.Observer}
|
||||
observationContext, cancel := context.WithTimeout(ctx, 15*time.Second)
|
||||
defer cancel()
|
||||
result, err := service.Reconcile(observationContext, request.Name)
|
||||
if err != nil {
|
||||
return ctrl.Result{}, err
|
||||
}
|
||||
// 查询超时后仍用 worker context 保存安全失败结果;manager 停止时不强行写入。
|
||||
if err := resources.PresentInstance(ctx, result); err != nil {
|
||||
return ctrl.Result{}, err
|
||||
}
|
||||
if result.Record == nil || result.RemoveProtection {
|
||||
return ctrl.Result{}, nil
|
||||
}
|
||||
return ctrl.Result{RequeueAfter: dependencyRetry}, nil
|
||||
}
|
||||
@@ -0,0 +1,192 @@
|
||||
package controller
|
||||
|
||||
import (
|
||||
"context"
|
||||
"errors"
|
||||
"testing"
|
||||
|
||||
databasev1alpha1 "git.ddupan.top/panxiao81/ayatori/api/database/v1alpha1"
|
||||
"git.ddupan.top/panxiao81/ayatori/internal/database/adapter/kubernetes"
|
||||
"git.ddupan.top/panxiao81/ayatori/internal/database/application"
|
||||
"git.ddupan.top/panxiao81/ayatori/internal/database/domain/instance"
|
||||
apierrors "k8s.io/apimachinery/pkg/api/errors"
|
||||
"k8s.io/apimachinery/pkg/api/meta"
|
||||
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
|
||||
ctrl "sigs.k8s.io/controller-runtime"
|
||||
"sigs.k8s.io/controller-runtime/pkg/client"
|
||||
"sigs.k8s.io/controller-runtime/pkg/controller/controllerutil"
|
||||
)
|
||||
|
||||
type instanceBackend struct {
|
||||
checks instance.ManagementChecks
|
||||
err error
|
||||
inspect func()
|
||||
closed int
|
||||
}
|
||||
|
||||
func (b *instanceBackend) Read(context.Context, instance.CredentialReference) (application.Credentials, error) {
|
||||
return application.NewCredentials("fixture", "test-only-instance-password")
|
||||
}
|
||||
|
||||
func (b *instanceBackend) Connect(context.Context, instance.Endpoint, application.Credentials) (application.Database, error) {
|
||||
return b, nil
|
||||
}
|
||||
|
||||
func (b *instanceBackend) InspectMetadata(context.Context) (application.DatabaseMetadata, error) {
|
||||
return application.DatabaseMetadata{Version: "18"}, nil
|
||||
}
|
||||
|
||||
func (b *instanceBackend) InspectManagement(context.Context) (application.DatabaseMetadata, error) {
|
||||
if b.inspect != nil {
|
||||
b.inspect()
|
||||
}
|
||||
return application.DatabaseMetadata{Version: "18", Management: b.checks}, b.err
|
||||
}
|
||||
|
||||
func (b *instanceBackend) Close() { b.closed++ }
|
||||
|
||||
func newInstanceReconciler(t *testing.T, apiClient client.Client, backend *instanceBackend) *InstanceReconciler {
|
||||
t.Helper()
|
||||
service, err := application.NewInstanceService(backend, backend)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
t.Cleanup(service.Close)
|
||||
return &InstanceReconciler{Client: apiClient, Reader: apiClient, Observer: service}
|
||||
}
|
||||
|
||||
func reconcileInstance(t *testing.T, reconciler *InstanceReconciler, object *databasev1alpha1.PostgreSQLInstance) {
|
||||
t.Helper()
|
||||
if _, err := reconciler.Reconcile(t.Context(), ctrl.Request{NamespacedName: client.ObjectKeyFromObject(object)}); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
}
|
||||
|
||||
func assertInstanceReason(t *testing.T, object *databasev1alpha1.PostgreSQLInstance, reason string) {
|
||||
t.Helper()
|
||||
condition := meta.FindStatusCondition(object.Status.Conditions, "Ready")
|
||||
if condition == nil || condition.Reason != reason || condition.ObservedGeneration != object.Generation {
|
||||
t.Fatalf("Instance 状态不是当前 generation 的 %s", reason)
|
||||
}
|
||||
if reason != "ManagementReady" && condition.Status != metav1.ConditionFalse {
|
||||
t.Fatal("失败状态仍为 Ready")
|
||||
}
|
||||
}
|
||||
|
||||
func TestInstanceObservationAPI(t *testing.T) {
|
||||
apiClient, _, _ := bindingEnvironment(t)
|
||||
backend := &instanceBackend{checks: instance.ManagementChecks{
|
||||
Connection: instance.CheckPassed, Metadata: instance.CheckPassed,
|
||||
Roles: instance.CheckPassed, Databases: instance.CheckPassed,
|
||||
Grants: instance.CheckPassed, Extensions: instance.CheckPassed,
|
||||
}}
|
||||
reconciler := newInstanceReconciler(t, apiClient, backend)
|
||||
object := readyInstance(t, apiClient, "observed-instance")
|
||||
backend.inspect = func() {
|
||||
current := &databasev1alpha1.PostgreSQLInstance{}
|
||||
current.Name = object.Name
|
||||
reload(t, apiClient, current)
|
||||
if !controllerutil.ContainsFinalizer(current, kubernetes.InstanceFinalizer) {
|
||||
t.Fatal("观察早于 finalizer 持久化")
|
||||
}
|
||||
}
|
||||
reconcileInstance(t, reconciler, object)
|
||||
reload(t, apiClient, object)
|
||||
assertInstanceReason(t, object, "ManagementReady")
|
||||
if object.Status.Phase != string(instance.PhaseReady) || object.Status.PostgreSQLVersion != "18" {
|
||||
t.Fatal("当前成功观察未呈现")
|
||||
}
|
||||
before := object.ResourceVersion
|
||||
reconcileInstance(t, reconciler, object)
|
||||
reload(t, apiClient, object)
|
||||
if object.ResourceVersion != before {
|
||||
t.Fatal("相同观察不应反复写入 status")
|
||||
}
|
||||
backend.err = application.ErrAuthentication
|
||||
reconcileInstance(t, reconciler, object)
|
||||
reload(t, apiClient, object)
|
||||
assertInstanceReason(t, object, "AuthenticationFailed")
|
||||
if backend.closed != 1 || object.Status.PostgreSQLVersion != "" {
|
||||
t.Fatal("观察失败应释放连接并清除旧版本结果")
|
||||
}
|
||||
backend.err = nil
|
||||
backend.checks.Grants = instance.CheckUnobserved
|
||||
reconcileInstance(t, reconciler, object)
|
||||
reload(t, apiClient, object)
|
||||
assertInstanceReason(t, object, "ObservationIncomplete")
|
||||
backend.checks.Grants = instance.CheckPassed
|
||||
// 用新 service/reconciler 恢复;不依赖上轮领域对象或 Ready。
|
||||
reconciler = newInstanceReconciler(t, apiClient, backend)
|
||||
reconcileInstance(t, reconciler, object)
|
||||
reload(t, apiClient, object)
|
||||
assertInstanceReason(t, object, "ManagementReady")
|
||||
|
||||
backend.inspect = func() {
|
||||
reload(t, apiClient, object)
|
||||
object.Annotations = map[string]string{"concurrent": "kept-by-instance-test"}
|
||||
if err := apiClient.Update(t.Context(), object); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
}
|
||||
_, err := reconciler.Reconcile(t.Context(), ctrl.Request{NamespacedName: client.ObjectKeyFromObject(object)})
|
||||
if !apierrors.IsConflict(err) {
|
||||
t.Fatal("旧观察不应覆盖在途 API 修改")
|
||||
}
|
||||
backend.inspect = nil
|
||||
reconcileInstance(t, reconciler, object)
|
||||
reload(t, apiClient, object)
|
||||
if object.Annotations["concurrent"] != "kept-by-instance-test" {
|
||||
t.Fatal("重试覆盖了其他字段")
|
||||
}
|
||||
}
|
||||
|
||||
type failedReferenceReader struct{ client.Reader }
|
||||
|
||||
func (*failedReferenceReader) List(context.Context, client.ObjectList, ...client.ListOption) error {
|
||||
return errors.New("injected reference list failure")
|
||||
}
|
||||
|
||||
func TestInstanceDeletionProtection(t *testing.T) {
|
||||
apiClient, _, _ := bindingEnvironment(t)
|
||||
backend := &instanceBackend{}
|
||||
reconciler := newInstanceReconciler(t, apiClient, backend)
|
||||
object := readyInstance(t, apiClient, "protected-instance")
|
||||
reconcileInstance(t, reconciler, object)
|
||||
reload(t, apiClient, object)
|
||||
database := availableDatabase(t, apiClient, "retained-database", object)
|
||||
database.Status.Phase = "Released"
|
||||
if err := apiClient.Status().Update(t.Context(), database); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
tenant := provisionTenant("pending-request", object.Name)
|
||||
requireCreate(t, apiClient, tenant)
|
||||
if err := apiClient.Delete(t.Context(), object); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
backend.inspect = func() { t.Fatal("删除中不应连接 PostgreSQL") }
|
||||
reconciler.Reader = &failedReferenceReader{Reader: apiClient}
|
||||
reconcileInstance(t, reconciler, object)
|
||||
reload(t, apiClient, object)
|
||||
assertInstanceReason(t, object, reasonDependency)
|
||||
reconciler.Reader = apiClient
|
||||
reconcileInstance(t, reconciler, object)
|
||||
reload(t, apiClient, object)
|
||||
assertInstanceReason(t, object, "InstanceInUse")
|
||||
if object.Status.Phase != string(instance.PhaseDeleting) || backend.closed != 1 {
|
||||
t.Fatal("删除没有停止本地观察")
|
||||
}
|
||||
if err := apiClient.Delete(t.Context(), database); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
reconcileInstance(t, reconciler, object)
|
||||
reload(t, apiClient, object)
|
||||
assertInstanceReason(t, object, "InstanceInUse")
|
||||
if err := apiClient.Delete(t.Context(), tenant); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
reconcileInstance(t, reconciler, object)
|
||||
if err := apiClient.Get(t.Context(), client.ObjectKeyFromObject(object), object); !apierrors.IsNotFound(err) {
|
||||
t.Fatal("最后一个引用解除后 Instance 应可删除")
|
||||
}
|
||||
reconcileInstance(t, reconciler, object)
|
||||
}
|
||||
@@ -0,0 +1,77 @@
|
||||
package controller
|
||||
|
||||
import (
|
||||
"context"
|
||||
"errors"
|
||||
|
||||
databasev1alpha1 "git.ddupan.top/panxiao81/ayatori/api/database/v1alpha1"
|
||||
corev1 "k8s.io/api/core/v1"
|
||||
"k8s.io/apimachinery/pkg/util/validation"
|
||||
ctrl "sigs.k8s.io/controller-runtime"
|
||||
"sigs.k8s.io/controller-runtime/pkg/cache"
|
||||
"sigs.k8s.io/controller-runtime/pkg/client"
|
||||
"sigs.k8s.io/controller-runtime/pkg/handler"
|
||||
)
|
||||
|
||||
// InstanceCacheOptions 必须在创建 manager 时使用;只 watch 固定 namespace 的 Secret metadata。
|
||||
// SecretCredentials 始终直读 API,不会令共享 cache 保存密码。
|
||||
func InstanceCacheOptions(namespace string) cache.Options {
|
||||
return cache.Options{ByObject: map[client.Object]cache.ByObject{
|
||||
&corev1.Secret{}: {Namespaces: map[string]cache.Config{namespace: {}}},
|
||||
}}
|
||||
}
|
||||
|
||||
func (r *InstanceReconciler) SetupWithManager(manager ctrl.Manager) error {
|
||||
if r.Observer == nil || len(validation.IsDNS1123Label(r.SecretNamespace)) != 0 {
|
||||
return errors.New("instance observer and valid management Secret namespace required")
|
||||
}
|
||||
if r.Client == nil {
|
||||
r.Client = manager.GetClient()
|
||||
}
|
||||
if r.Reader == nil {
|
||||
r.Reader = manager.GetAPIReader()
|
||||
}
|
||||
return ctrl.NewControllerManagedBy(manager).
|
||||
Named("database-instance").
|
||||
For(&databasev1alpha1.PostgreSQLInstance{}).
|
||||
WatchesMetadata(&corev1.Secret{}, handler.EnqueueRequestsFromMapFunc(r.instancesForSecret)).
|
||||
Watches(&databasev1alpha1.PostgreSQLDatabase{}, handler.EnqueueRequestsFromMapFunc(r.instanceForReference)).
|
||||
Watches(&databasev1alpha1.PostgreSQLTenant{}, handler.EnqueueRequestsFromMapFunc(r.instanceForReference)).
|
||||
Complete(r)
|
||||
}
|
||||
|
||||
func (r *InstanceReconciler) instancesForSecret(ctx context.Context, object client.Object) []ctrl.Request {
|
||||
if object.GetNamespace() != r.SecretNamespace {
|
||||
return nil
|
||||
}
|
||||
instances := &databasev1alpha1.PostgreSQLInstanceList{}
|
||||
if err := r.Client.List(ctx, instances); err != nil {
|
||||
ctrl.LoggerFrom(ctx).Error(err, "无法映射管理 Secret 事件;等待低频重试")
|
||||
return nil
|
||||
}
|
||||
var requests []ctrl.Request
|
||||
for _, item := range instances.Items {
|
||||
if string(item.Spec.AdminCredentialRef.Name) == object.GetName() {
|
||||
request := ctrl.Request{Name: item.Name}
|
||||
requests = append(requests, request)
|
||||
}
|
||||
}
|
||||
return requests
|
||||
}
|
||||
|
||||
func (r *InstanceReconciler) instanceForReference(_ context.Context, object client.Object) []ctrl.Request {
|
||||
var name string
|
||||
switch item := object.(type) {
|
||||
case *databasev1alpha1.PostgreSQLDatabase:
|
||||
name = string(item.Spec.InstanceRef.Name)
|
||||
case *databasev1alpha1.PostgreSQLTenant:
|
||||
if item.Spec.Provision != nil {
|
||||
name = string(item.Spec.Provision.InstanceRef.Name)
|
||||
}
|
||||
}
|
||||
if name == "" {
|
||||
return nil
|
||||
}
|
||||
request := ctrl.Request{Name: name}
|
||||
return []ctrl.Request{request}
|
||||
}
|
||||
Reference in New Issue
Block a user