Author SHA1 Message Date
panxiao81 6acba0ca46 Merge pull request 'feat: 迁移 PostgreSQL registry 所有权存储' (#7) from feat/database-registry into main
Verify / test (push) Successful in 8m36s
Verify / lint (push) Successful in 10m32s
Verify / database-integration (push) Successful in 6m9s
2026-09-24 14:46:52 +00:00
panxiao81 77848c7dc8 Merge pull request 'feat: 接入管理 Secret 凭据与连接刷新' (#6) from feat/database-admin-credentials into main
Verify / test (push) Successful in 5m45s
Verify / database-integration (push) Successful in 11m41s
Verify / lint (push) Successful in 12m38s
2026-09-24 14:45:58 +00:00
panxiao81 e590542b0b feat: 迁移 PostgreSQL registry 所有权存储与恢复测试
Verify / test (pull_request) Successful in 9m18s
Verify / lint (pull_request) Successful in 10m19s
Verify / database-integration (pull_request) Successful in 11m32s
2026-09-21 15:59:14 +00:00
panxiao81 3d91abe380 fix: 按 runner 默认提供 Docker 的约定仅执行预检
Verify / test (pull_request) Successful in 5m45s
Verify / lint (pull_request) Successful in 7m44s
Verify / database-integration (pull_request) Successful in 11m42s
2026-09-21 15:59:12 +00:00
panxiao81 fb2c34ce46 fix: 不将 Docker 数据目录独立挂载作为启动前置条件
Verify / database-integration (pull_request) Failing after 11m16s
Verify / lint (pull_request) Successful in 15m38s
Verify / test (pull_request) Successful in 16m2s
2026-09-21 15:51:10 +00:00
panxiao81 505aeb8e50 fix: 在 Pod CI 初始化 Docker 并保留 fixture 诊断
Verify / database-integration (pull_request) Failing after 58s
Verify / test (pull_request) Successful in 11m30s
Verify / lint (pull_request) Successful in 13m59s
2026-09-21 15:48:06 +00:00
panxiao81 c110aceb0f feat: 接入管理 Secret 凭据与连接刷新
Verify / test (pull_request) Successful in 4m45s
Verify / lint (pull_request) Successful in 7m5s
Verify / database-integration (pull_request) Failing after 11m1s
2026-09-21 08:35:24 +00:00
panxiao81 341c095a97 Merge pull request 'feat: 实现 Instance Ready 领域判定与恢复规则' (#5) from feat/database-instance-readiness into main
Verify / test (push) Successful in 11m36s
Verify / lint (push) Successful in 12m3s
Reviewed-on: #5
2026-09-21 07:58:31 +00:00
panxiao81 cd0d3a70ae feat: 实现 Instance Ready 领域判定与恢复规则
Verify / test (pull_request) Successful in 7m20s
Verify / lint (pull_request) Successful in 7m51s
2026-09-21 07:49:41 +00:00
panxiao81 984c0aee73 Merge pull request 'feat: 迁移 Database Instance 领域基线与扩展观测' (#4) from feat/database-instance-domain into main
Verify / test (push) Successful in 6m0s
Verify / lint (push) Successful in 6m23s
Reviewed-on: #4
2026-09-21 07:42:08 +00:00
23 changed files with 3034 additions and 22 deletions
+29 -1
View File
@@ -40,4 +40,32 @@ jobs:
cache: true cache: true
- name: Lint - name: Lint
run: make lint run: |
make lint
make lint-database-integration
database-integration:
runs-on: [self-hosted, pod]
timeout-minutes: 30
steps:
- name: Checkout
uses: actions/checkout@de0fac2e4500dabe0009e67214ff5f5447ce83dd
with:
persist-credentials: false
- name: Set up Go
uses: actions/setup-go@4b73464bb391d4059bd26b0524d20df3927bd417
with:
go-version-file: go.mod
cache: true
# Runner 提供本 job 可用的 Docker;workflow 只验证,不重复启动 daemon。
- name: Verify Docker availability
shell: bash
run: |
set -euo pipefail
docker version
docker info --format 'Server={{.ServerVersion}} StorageDriver={{.Driver}}'
- name: Test Database integration with real backends
run: make test-database-integration
+8
View File
@@ -67,6 +67,14 @@ test: manifests generate fmt vet setup-envtest ## Run tests.
lint: golangci-lint ## Run golangci-lint linter lint: golangci-lint ## Run golangci-lint linter
"$(GOLANGCI_LINT)" run "$(GOLANGCI_LINT)" run
.PHONY: test-database-integration
test-database-integration: setup-envtest ## 使用临时 API server 与独立 PostgreSQL 容器验证凭据读取和连接更新。
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
lint-database-integration: golangci-lint ## 检查集成测试构建标签下的 Database 代码。
"$(GOLANGCI_LINT)" run --build-tags=integration ./internal/database/...
.PHONY: lint-fix .PHONY: lint-fix
lint-fix: golangci-lint ## Run golangci-lint linter and perform fixes lint-fix: golangci-lint ## Run golangci-lint linter and perform fixes
"$(GOLANGCI_LINT)" run --fix "$(GOLANGCI_LINT)" run --fix
+48 -3
View File
@@ -21,15 +21,60 @@ Database 是 Ayatori 首批实际产品领域之一。第一个迁移切片只
代码被移动到 Ayatori 的 `internal/database/domain/instance`,测试 import 和文档链接相应更新; 代码被移动到 Ayatori 的 `internal/database/domain/instance`,测试 import 和文档链接相应更新;
首个后续切片按已批准合同增加 Instance extension observation:观测与当前 target 绑定,进入重新 首个后续切片按已批准合同增加 Instance extension observation:观测与当前 target 绑定,进入重新
验证或删除时失效,且支持判定不授权 Tenant provisioning。其余 Ready/observation 行为仍应先更新 验证或删除时失效,且支持判定不授权 Tenant provisioning。后续 Ready 切片实现管理能力判定、
合同与测试再实现,不能把旧运行链路接回该模型。 registry 准备决策与完整回读、Ready 重验及本轮 evidence 前置检查;沿用已批准合同,不能把旧运行
链路接回该模型。各层验证边界见 [Instance 领域规格](domain-instance.md)。
## 边界 ## 边界
- 领域层不依赖 Kubernetes types、数据库 driver 或凭据 provider。 - 领域层不依赖 Kubernetes types、数据库 driver 或凭据 provider。
- CredentialReference 只携带管理 Secret 的名称与字段映射,不包含 Secret 内容或 OpenBao path。 - CredentialReference 只携带管理 Secret 的名称与字段映射,不包含 Secret 内容或 OpenBao path。
- Instance checkpoint 不是外部事实;实际能力必须由 application/adapter 观察后交给领域对象判断。 - Instance checkpoint 不是外部事实;实际能力必须由 application/adapter 观察后交给领域对象判断。
- 当前代码不授权 Tenant provisioning,也不表示 Database API 已经可用。 - 当前代码只检查 Instance 供应前置条件,不授予 Tenant 所有权或外部写入权限,也不表示
Database API 已经可用。
## 管理凭据与连接切片
`application.InstanceService` 适配自原项目固定基线
[`internal/instance/service.go`](https://git.ddupan.top/panxiao81/postgresql-tenant-operator/src/commit/dae546e58efa1be81e930861c87f7fb13bb12113/internal/instance/service.go),
保留 CredentialReader、Connector、Database 的装配边界及串行操作/释放规则。
具体连接池完全由 pgxpool v5.11.0 提供,PostgreSQL adapter 不拥有凭据缓存、轮换流程或任意观测回调。
相对源基线有两项按已批准合同作出的必要修改:管理凭据从固定 controller namespace 的
Kubernetes Secret 直接读取;每轮比较有效用户名、密码,检测到变化即关闭旧连接并重新装配。
Secret metadata 和无关字段变化不重建连接。观测后再次读取 Secret,中途有效值变化则丢弃结果,
不把旧连接的成功作为新凭据有效的证据;这不构成跨 Kubernetes/PostgreSQL 的原子事务。
Instance UID、endpoint 或凭据引用变化也会释放旧连接;Forget/Close 只释放本地资源。
当前只有 `ObserveVersion` 版本查询用例,不能产生完整 CapabilityObservation 或 Ready。
controller 接入、Secret watch、finalizer、registry 观测装配与真实权限检查仍待后续切片;并发 CR 更新
必须由调用者通过 resourceVersion 校验。应用层沿用源实现的串行处理,本阶段未引入新的调度框架。
运行 `make test-database-integration` 验证真实 API server + 一次性 PostgreSQL;fixture 不接受外部
DSN,镜像固定摘要,使用随机本机回环端口并在退出时删除测试容器。覆盖缺失/错误凭据、RBAC、
namespace 边界、有效值轮换、metadata 无关变化、中途轮换、重建/重试、并发读取、Forget/Close
与 TLS DNS/IP SAN、错误 CA/主机名和禁止明文降级。CI 使用 Pod runner 执行,由 runner 提供
可用的 Docker,workflow 只做预检、不自行启动 daemon;不依赖 VM。普通 lint 之外还检查
integration 标签代码。领域单测、真实 API 行为与真实 PostgreSQL 行为分别验收,不以本切片
替代整个 Instance controller 的集成验收。
## Registry 所有权存储切片
`adapter/postgresql/registry` 直接迁入上述固定基线的 `internal/postgresql/registry`,
保留原 schema、版本表和唯一约束,不接回旧 controller。Store 使用调用方提供的 pgx 连接或
连接池,不创建连接池、不读取凭据、不管理重连。迁移复用 tern v2.4.3 的事务和 advisory lock,
所有权变更复用 PostgreSQL 唯一约束与行锁,没有新增迁移或锁框架。
Claim 支持相同归属的幂等重试;Retain 保留不可重新占用的墓碑;Delete 仅删除匹配归属的
managed 记录,且只应在调用方确认外部资源已删除后调用。Store 本身不删除 database、role
或凭据,也不授予 Tenant provisioning 权限。驱动错误仅供内部调用方分类,接入 controller 时
仍须按安全合同转换为脱敏的 Condition/Event,不能直接发布原始错误。
真实 PostgreSQL 测试覆盖重复/并发初始化、唯一约束冲突、并发 Claim、归属不匹配、Retain
与 Delete 重试、连接重建、取消恢复、未知 schema/超前版本拒绝,以及真实事务提交后注入
客户端失败的恢复分支。后者确定性模拟“提交成功但调用方不知道”,不声称覆盖真实网络分区。
registry 尚未装配到 InstanceService;初始化成功不能当作完整 registry 回读或 Instance Ready。
controller 重启、Secret watch、finalizer 与跨后端删除仍由后续 envtest/端到端切片验收。
## 设计入口 ## 设计入口
+33
View File
@@ -3,6 +3,39 @@
> 本页迁入作为 Database 模块的测试分层与 fixture 合同。旧项目的 Make target、devcontainer > 本页迁入作为 Database 模块的测试分层与 fixture 合同。旧项目的 Make target、devcontainer
> 和脚手架版本尚未适配 Ayatori;实现时应复用 Ayatori 现有工具链,并保持这里定义的测试边界。 > 和脚手架版本尚未适配 Ayatori;实现时应复用 Ayatori 现有工具链,并保持这里定义的测试边界。
## Ayatori 已接入的凭据与 registry 切片测试
本节命令已在 Ayatori 接入;以下历史 Compose/Kind 操作仍属于迁入的目标合同。
```sh
make test
make lint
make lint-database-integration
make test-database-integration
```
最后一项要求本机 Docker 可用。它启动 envtest 的真实 API server/etcd 和固定镜像摘要的临时
PostgreSQL 容器,随机绑定回环端口,不读取 kubeconfig,也不接受指向现有数据库的 DSN。
每个凭据场景使用独立环境;测试清理仅关闭自己的进程和按确切 ID 删除自己的容器。
TLS 测试在临时目录生成一次性证书与私钥,不使用生产 CA。
当前覆盖固定 namespace 的 Secret 读取与 RBAC、缺失/无效凭据恢复、有效凭据变化后的重连、
无关字段更新不重连、中途轮换时丢弃观察、会话重建、并发读取、本地连接释放和 TLS 验证。
快速测试、lint 和 Database 集成测试均使用 Pod runner。按维护者于 2026-09-21 更新的接口
约定,runner 提供默认可用的 Docker;workflow 通过 `docker version` 和 `docker info` 预检,
不自行启动 daemon、不强制 storage driver 或覆盖 Docker endpoint。该约定的 CI 验收依赖
runner 后端修复上线,不能从本地测试通过推断远端已经可用。
fixture 启动失败会保留退出错误与 stderr,并遮蔽测试密码,
以区分缺少命令、daemon 不可达、权限和镜像拉取失败。
registry 测试独立使用一次性 PostgreSQL,不启动 Kubernetes API server。覆盖重复和并发迁移、
所有权唯一约束、并发占用、Retain 墓碑与 Delete 幂等、连接重建、取消恢复、未知 schema 与
超前版本拒绝。通过在真实 COMMIT 成功后注入客户端错误,验证结果不确定时的重试;这不替代
真实网络故障测试,也不覆盖跨后端的删除步骤。
这些测试尚不包含 Instance CRD/controller、Secret watch、status/finalizer 事件链、registry 观测装配、
权限探测矩阵、ESO 或 Tenant 供应。版本查询成功不意味着 Instance Ready。
本项目同时依赖 Kubernetes API、PostgreSQL、OpenBao 和 ESO。日常开发不连接 homelab 本项目同时依赖 Kubernetes API、PostgreSQL、OpenBao 和 ESO。日常开发不连接 homelab
中的真实服务:Kubernetes 使用 envtest 或一次性 Kind,另外两个依赖使用一次性 中的真实服务:Kubernetes 使用 envtest 或一次性 Kind,另外两个依赖使用一次性
容器。这样既避免污染真实数据,也能把启动顺序固化为命令。 容器。这样既避免污染真实数据,也能把启动顺序固化为命令。
+24 -1
View File
@@ -2,7 +2,7 @@
状态:Draft,含已确认决策。日期:2026-09-13。 状态:Draft,含已确认决策。日期:2026-09-13。
上层合并边界见 [ADR-0008](../decisions/0008-merge-Ayatori Database controller.md)。本文只展开 Instance,不包含 Tenant 的供应 上层合并边界见 [ADR-0008](../decisions/0008-merge-postgresql-tenant-operator.md)。本文只展开 Instance,不包含 Tenant 的供应
实现,也不新增 CRD 字段。设计签名用于评审职责与行为,不是待复制的 Go 接口代码。 实现,也不新增 CRD 字段。设计签名用于评审职责与行为,不是待复制的 Go 接口代码。
## 1. 对象职责与生命周期 ## 1. 对象职责与生命周期
@@ -228,3 +228,26 @@ Ready --registry 需修复/保存--> InitializingRegistry
均已确认。其他决策及未决项见总体草案,不增加后台清扫器或状态字段。 均已确认。其他决策及未决项见总体草案,不增加后台清扫器或状态字段。
批准本对象结构不等于批准这些未决行为,也不意味着立刻实现完整供应链路。 批准本对象结构不等于批准这些未决行为,也不意味着立刻实现完整供应链路。
## 8. 领域实现与验证边界(2026-09-21)
`internal/database/domain/instance` 按上述方法合同实现 Ready 纯判定。输入分别表达连接、
metadata、role、database、grant、extension 管理能力,以及 registry 的未观察、缺失、需迁移、
可用、不兼容和不可访问状态。检查零值或未知值按证据不足处理;操作失败只使用封闭的安全
类别,不接收驱动错误。`Snapshot.Failure` 是 Condition 映射的领域输入,不新增 CRD/status 字段。
管理能力分别指目标连接可用、服务器 metadata 可读,以及执行规格 §7 所要求的角色、数据库、
授权和扩展管理操作的能力;不是仅凭版本查询或扩展可用列表判定权限。具体 SQL 权限探测矩阵、
最小权限角色和扩展权限例外仍须在 PostgreSQL adapter 切片定义并用真实后端验证。
registry 不兼容独立保留为领域失败类别,不将其误报为权限不足;公开 Condition Reason 的映射
留待 API/application 切片按原合同评审。
领域测试验证完整回读、缺少检查项、状态重建、重复判定、目标不匹配、配置变化、依赖失败、
registry 丢失/不兼容、操作结果不确定和删除限制。只有本轮完整能力判定通过后,Instance
前置条件检查才通过;这不授予 Tenant 所有权,也不替代实际写入前的并发校验。
本切片不新增 controller、adapter、Secret 读取或外部生命周期操作。checkpoint 保存失败、
resourceVersion 冲突、watch 与 finalizer 事件链需由后续 application/envtest 验证;SQL 探测、
registry 初始化/迁移、超时后的真实状态回读和并发幂等由 PostgreSQL 集成测试验证;Secret
变化后的连接刷新由 Kubernetes API 加真实 PostgreSQL 的集成测试验证。纯领域测试不能证明
Database 已可运行或这些集成合同已完成。
+17 -2
View File
@@ -3,6 +3,8 @@ module git.ddupan.top/panxiao81/ayatori
go 1.27.1 go 1.27.1
require ( require (
github.com/jackc/pgx/v5 v5.11.0
github.com/jackc/tern/v2 v2.4.3
k8s.io/api v0.37.0 k8s.io/api v0.37.0
k8s.io/apimachinery v0.37.0 k8s.io/apimachinery v0.37.0
k8s.io/client-go v0.37.0 k8s.io/client-go v0.37.0
@@ -11,6 +13,10 @@ require (
require ( require (
cel.dev/expr v0.25.1 // indirect cel.dev/expr v0.25.1 // indirect
dario.cat/mergo v1.0.2 // indirect
github.com/Masterminds/goutils v1.1.1 // indirect
github.com/Masterminds/semver/v3 v3.5.0 // indirect
github.com/Masterminds/sprig/v3 v3.3.0 // indirect
github.com/antlr4-go/antlr/v4 v4.13.1 // indirect github.com/antlr4-go/antlr/v4 v4.13.1 // indirect
github.com/beorn7/perks v1.0.1 // indirect github.com/beorn7/perks v1.0.1 // indirect
github.com/blang/semver/v4 v4.0.0 // indirect github.com/blang/semver/v4 v4.0.0 // indirect
@@ -43,8 +49,14 @@ require (
github.com/google/gnostic-models v0.7.0 // indirect github.com/google/gnostic-models v0.7.0 // indirect
github.com/google/uuid v1.6.0 // indirect github.com/google/uuid v1.6.0 // indirect
github.com/grpc-ecosystem/grpc-gateway/v2 v2.29.0 // indirect github.com/grpc-ecosystem/grpc-gateway/v2 v2.29.0 // indirect
github.com/huandu/xstrings v1.5.0 // indirect
github.com/inconshreveable/mousetrap v1.1.0 // 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/json-iterator/go v1.1.12 // indirect
github.com/mitchellh/copystructure v1.2.0 // indirect
github.com/mitchellh/reflectwalk v1.0.2 // indirect
github.com/modern-go/concurrent v0.0.0-20180306012644-bacd9c7ef1dd // 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/modern-go/reflect2 v1.0.3-0.20250322232337-35a7c28c31ee // indirect
github.com/munnerz/goautoneg v0.0.0-20191010083416-a7dc8b61c822 // indirect github.com/munnerz/goautoneg v0.0.0-20191010083416-a7dc8b61c822 // indirect
@@ -53,6 +65,8 @@ require (
github.com/prometheus/client_model v0.6.2 // indirect github.com/prometheus/client_model v0.6.2 // indirect
github.com/prometheus/common v0.70.0 // indirect github.com/prometheus/common v0.70.0 // indirect
github.com/prometheus/procfs v0.21.1 // indirect github.com/prometheus/procfs v0.21.1 // indirect
github.com/shopspring/decimal v1.4.0 // indirect
github.com/spf13/cast v1.10.0 // indirect
github.com/spf13/cobra v1.10.2 // indirect github.com/spf13/cobra v1.10.2 // indirect
github.com/spf13/pflag v1.0.10 // indirect github.com/spf13/pflag v1.0.10 // indirect
github.com/x448/float16 v0.8.4 // indirect github.com/x448/float16 v0.8.4 // indirect
@@ -68,14 +82,15 @@ require (
go.uber.org/multierr v1.11.0 // indirect go.uber.org/multierr v1.11.0 // indirect
go.uber.org/zap v1.27.1 // indirect go.uber.org/zap v1.27.1 // indirect
go.yaml.in/yaml/v2 v2.4.4 // indirect go.yaml.in/yaml/v2 v2.4.4 // indirect
go.yaml.in/yaml/v3 v3.0.4 // indirect go.yaml.in/yaml/v3 v3.0.5 // indirect
golang.org/x/crypto v0.55.0 // indirect
golang.org/x/exp v0.0.0-20260410095643-746e56fc9e2f // indirect golang.org/x/exp v0.0.0-20260410095643-746e56fc9e2f // indirect
golang.org/x/net v0.57.0 // indirect golang.org/x/net v0.57.0 // indirect
golang.org/x/oauth2 v0.36.0 // indirect golang.org/x/oauth2 v0.36.0 // indirect
golang.org/x/sync v0.22.0 // indirect golang.org/x/sync v0.22.0 // indirect
golang.org/x/sys v0.47.0 // indirect golang.org/x/sys v0.47.0 // indirect
golang.org/x/term v0.45.0 // indirect golang.org/x/term v0.45.0 // indirect
golang.org/x/text v0.40.0 // indirect golang.org/x/text v0.41.0 // indirect
golang.org/x/time v0.15.0 // indirect golang.org/x/time v0.15.0 // indirect
gomodules.xyz/jsonpatch/v2 v2.4.0 // indirect gomodules.xyz/jsonpatch/v2 v2.4.0 // indirect
google.golang.org/genproto/googleapis/api v0.0.0-20260526163538-3dc84a4a5aaa // indirect google.golang.org/genproto/googleapis/api v0.0.0-20260526163538-3dc84a4a5aaa // indirect
+44 -13
View File
@@ -1,7 +1,13 @@
cel.dev/expr v0.25.1 h1:1KrZg61W6TWSxuNZ37Xy49ps13NUovb66QLprthtwi4= cel.dev/expr v0.25.1 h1:1KrZg61W6TWSxuNZ37Xy49ps13NUovb66QLprthtwi4=
cel.dev/expr v0.25.1/go.mod h1:hrXvqGP6G6gyx8UAHSHJ5RGk//1Oj5nXQ2NI02Nrsg4= cel.dev/expr v0.25.1/go.mod h1:hrXvqGP6G6gyx8UAHSHJ5RGk//1Oj5nXQ2NI02Nrsg4=
github.com/Masterminds/semver/v3 v3.4.0 h1:Zog+i5UMtVoCU8oKka5P7i9q9HgrJeGzI9SA1Xbatp0= dario.cat/mergo v1.0.2 h1:85+piFYR1tMbRrLcDwR18y4UKJ3aH1Tbzi24VRW1TK8=
github.com/Masterminds/semver/v3 v3.4.0/go.mod h1:4V+yj/TJE1HU9XfppCwVMZq3I84lprf4nC11bSS5beM= dario.cat/mergo v1.0.2/go.mod h1:E/hbnu0NxMFBjpMIE34DRGLWqDy0g5FuKDhCb31ngxA=
github.com/Masterminds/goutils v1.1.1 h1:5nUrii3FMTL5diU80unEVvNevw1nH4+ZV4DSLVJLSYI=
github.com/Masterminds/goutils v1.1.1/go.mod h1:8cTjp+g8YejhMuvIA5y2vz3BpJxksy863GQaJW2MFNU=
github.com/Masterminds/semver/v3 v3.5.0 h1:kQceYJfbupGfZOKZQg0kou0DgAKhzDg2NZPAwZ/2OOE=
github.com/Masterminds/semver/v3 v3.5.0/go.mod h1:4V+yj/TJE1HU9XfppCwVMZq3I84lprf4nC11bSS5beM=
github.com/Masterminds/sprig/v3 v3.3.0 h1:mQh0Yrg1XPo6vjYXgtf5OtijNAKJRNcTdOOGZe3tPhs=
github.com/Masterminds/sprig/v3 v3.3.0/go.mod h1:Zy1iXRYNqNLUolqCpL4uhk6SHUMAOSCzdgBfDb35Lz0=
github.com/antlr4-go/antlr/v4 v4.13.1 h1:SqQKkuVZ+zWkMMNkjy5FZe5mr5WURWnlpmOuzYWrPrQ= github.com/antlr4-go/antlr/v4 v4.13.1 h1:SqQKkuVZ+zWkMMNkjy5FZe5mr5WURWnlpmOuzYWrPrQ=
github.com/antlr4-go/antlr/v4 v4.13.1/go.mod h1:GKmUxMtwp6ZgGwZSva4eWPC5mS6vUAmOABFgjdkM7Nw= github.com/antlr4-go/antlr/v4 v4.13.1/go.mod h1:GKmUxMtwp6ZgGwZSva4eWPC5mS6vUAmOABFgjdkM7Nw=
github.com/beorn7/perks v1.0.1 h1:VlbKKnNfV8bJzeqoa4cOKqO6bYr3WgKZxO8Z16+hsOM= github.com/beorn7/perks v1.0.1 h1:VlbKKnNfV8bJzeqoa4cOKqO6bYr3WgKZxO8Z16+hsOM=
@@ -25,6 +31,8 @@ github.com/evanphx/json-patch/v5 v5.9.11 h1:/8HVnzMq13/3x9TPvjG08wUGqBTmZBsCWzjT
github.com/evanphx/json-patch/v5 v5.9.11/go.mod h1:3j+LviiESTElxA4p3EMKAB9HXj3/XEtnUf6OZxqIQTM= github.com/evanphx/json-patch/v5 v5.9.11/go.mod h1:3j+LviiESTElxA4p3EMKAB9HXj3/XEtnUf6OZxqIQTM=
github.com/felixge/httpsnoop v1.0.4 h1:NFTV2Zj1bL4mc9sqWACXbQFVBBg2W3GPvqp8/ESS2Wg= github.com/felixge/httpsnoop v1.0.4 h1:NFTV2Zj1bL4mc9sqWACXbQFVBBg2W3GPvqp8/ESS2Wg=
github.com/felixge/httpsnoop v1.0.4/go.mod h1:m8KPJKqk1gH5J9DgRY2ASl2lWCfGKXixSwevea8zH2U= github.com/felixge/httpsnoop v1.0.4/go.mod h1:m8KPJKqk1gH5J9DgRY2ASl2lWCfGKXixSwevea8zH2U=
github.com/frankban/quicktest v1.14.6 h1:7Xjx+VpznH+oBnejlPUj8oUpdxnVs4f8XU8WnHkI4W8=
github.com/frankban/quicktest v1.14.6/go.mod h1:4ptaffx2x8+WTWXmUCuVU6aPUX1/Mz7zb5vbUoiM6w0=
github.com/fsnotify/fsnotify v1.9.0 h1:2Ml+OJNzbYCTzsxtv8vKSFD9PbJjmhYF14k/jKC7S9k= 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/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 h1:2rWm8B193Ll4VdjsJY28jxs70IdDsHRWgQYAI80+rMQ=
@@ -89,8 +97,20 @@ 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/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 h1:5VipnvEpbqr2gA2VbM+nYVbkIF28c5ZQfqCBQ5g2xfk=
github.com/grpc-ecosystem/grpc-gateway/v2 v2.29.0/go.mod h1:Hyl3n6Twe1hvtd9XUXDec4pTvgMSEixRuQKPTMH2bNs= github.com/grpc-ecosystem/grpc-gateway/v2 v2.29.0/go.mod h1:Hyl3n6Twe1hvtd9XUXDec4pTvgMSEixRuQKPTMH2bNs=
github.com/huandu/xstrings v1.5.0 h1:2ag3IFq9ZDANvthTwTiqSSZLjDc+BedvHPAp5tJy2TI=
github.com/huandu/xstrings v1.5.0/go.mod h1:y5/lhBue+AyNmUVz9RLU9xbLR0o4KIIExikq4ovT0aE=
github.com/inconshreveable/mousetrap v1.1.0 h1:wN+x4NVGpMsO7ErUn/mUI3vEoE6Jt13X2s0bqwp9tc8= github.com/inconshreveable/mousetrap v1.1.0 h1:wN+x4NVGpMsO7ErUn/mUI3vEoE6Jt13X2s0bqwp9tc8=
github.com/inconshreveable/mousetrap v1.1.0/go.mod h1:vpF70FUmC8bwa3OWnCshd2FqLfsEA9PFc4w1p2J65bw= github.com/inconshreveable/mousetrap v1.1.0/go.mod h1:vpF70FUmC8bwa3OWnCshd2FqLfsEA9PFc4w1p2J65bw=
github.com/jackc/pgpassfile v1.0.0 h1:/6Hmqy13Ss2zCq62VdNG8tM1wchn8zjSGOBJ6icpsIM=
github.com/jackc/pgpassfile v1.0.0/go.mod h1:CEx0iS5ambNFdcRtxPj5JhEz+xB6uRky5eyVu/W2HEg=
github.com/jackc/pgservicefile v0.0.0-20240606120523-5a60cdf6a761 h1:iCEnooe7UlwOQYpKFhBabPMi4aNAfoODPEFNiAnClxo=
github.com/jackc/pgservicefile v0.0.0-20240606120523-5a60cdf6a761/go.mod h1:5TJZWKEWniPve33vlWYSoGYefn3gLQRzjfDlhSJ9ZKM=
github.com/jackc/pgx/v5 v5.11.0 h1:IzBBtyK9AHqf98cctWFifYSci2hgQR/cd56wB4p+ogg=
github.com/jackc/pgx/v5 v5.11.0/go.mod h1:mal1tBGAFfLHvZzaYh77YS/eC6IX9OWbRV1QIIM0Jn4=
github.com/jackc/puddle/v2 v2.2.2 h1:PR8nw+E/1w0GLuRFSmiioY6UooMp6KJv0/61nB7icHo=
github.com/jackc/puddle/v2 v2.2.2/go.mod h1:vriiEXHvEE654aYKXXjOvZM39qJ0q+azkZFrfEOc3H4=
github.com/jackc/tern/v2 v2.4.3 h1:g293d3OZgW7OFEhsYXgEv0C21jea2boNr0VR4K5I7OY=
github.com/jackc/tern/v2 v2.4.3/go.mod h1:rMpMuRYcff5wWLptoTSO1qcDxJ4OodysvK17i2SVBys=
github.com/json-iterator/go v1.1.12 h1:PV8peI4a0ysnczrg+LtxykD8LfKY9ML6u2jnxaEnrnM= github.com/json-iterator/go v1.1.12 h1:PV8peI4a0ysnczrg+LtxykD8LfKY9ML6u2jnxaEnrnM=
github.com/json-iterator/go v1.1.12/go.mod h1:e30LSqwooZae/UwlEbR2852Gd8hjQvJoHmT4TnhNGBo= github.com/json-iterator/go v1.1.12/go.mod h1:e30LSqwooZae/UwlEbR2852Gd8hjQvJoHmT4TnhNGBo=
github.com/klauspost/compress v1.19.0 h1:sXLILfc9jV2QYWkzFOPWStmcUVH2RHEB1JCdY2oVvCQ= github.com/klauspost/compress v1.19.0 h1:sXLILfc9jV2QYWkzFOPWStmcUVH2RHEB1JCdY2oVvCQ=
@@ -101,6 +121,10 @@ github.com/kr/text v0.2.0 h1:5Nx0Ya0ZqY2ygV366QzturHI13Jq95ApcVaJBhpS+AY=
github.com/kr/text v0.2.0/go.mod h1:eLer722TekiGuMkidMxC/pM04lWEeraHUUmBw8l2grE= github.com/kr/text v0.2.0/go.mod h1:eLer722TekiGuMkidMxC/pM04lWEeraHUUmBw8l2grE=
github.com/kylelemons/godebug v1.1.0 h1:RPNrshWIDI6G2gRW9EHilWtl7Z6Sb1BR0xunSBf0SNc= 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/kylelemons/godebug v1.1.0/go.mod h1:9/0rRGxNHcop5bhtWyNeEfOS8JIWk580+fNqagV/RAw=
github.com/mitchellh/copystructure v1.2.0 h1:vpKXTN4ewci03Vljg/q9QvCGUDttBOGBIa15WveJJGw=
github.com/mitchellh/copystructure v1.2.0/go.mod h1:qLl+cE2AmVv+CoeAwDPye/v+N2HKCj9FbZEVFJRxO9s=
github.com/mitchellh/reflectwalk v1.0.2 h1:G2LzWKi524PWgd3mLHV8Y5k7s6XUvT0Gef6zxSIeXaQ=
github.com/mitchellh/reflectwalk v1.0.2/go.mod h1:mSTlrgnPZtwu0c4WaC2kGObEpuNDbx0jmZXqmk4esnw=
github.com/modern-go/concurrent v0.0.0-20180228061459-e0a39a4cb421/go.mod h1:6dJC0mAP4ikYIbvyc7fijjWJddQyLn8Ig3JB5CqoB9Q= 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 h1:TRLaZ9cD/w8PVh93nsPXa1VrQ6jlwL5oN8l14QlcNfg=
github.com/modern-go/concurrent v0.0.0-20180306012644-bacd9c7ef1dd/go.mod h1:6dJC0mAP4ikYIbvyc7fijjWJddQyLn8Ig3JB5CqoB9Q= github.com/modern-go/concurrent v0.0.0-20180306012644-bacd9c7ef1dd/go.mod h1:6dJC0mAP4ikYIbvyc7fijjWJddQyLn8Ig3JB5CqoB9Q=
@@ -129,6 +153,10 @@ github.com/prometheus/procfs v0.21.1/go.mod h1:aB55Cww9pdSJVHk0hUf0inxWyyjPogFIj
github.com/rogpeppe/go-internal v1.14.1 h1:UQB4HGPB6osV0SQTLymcB4TgvyWu6ZyliaW0tI/otEQ= github.com/rogpeppe/go-internal v1.14.1 h1:UQB4HGPB6osV0SQTLymcB4TgvyWu6ZyliaW0tI/otEQ=
github.com/rogpeppe/go-internal v1.14.1/go.mod h1:MaRKkUm5W0goXpeCfT7UZI6fk/L7L7so1lCWt35ZSgc= github.com/rogpeppe/go-internal v1.14.1/go.mod h1:MaRKkUm5W0goXpeCfT7UZI6fk/L7L7so1lCWt35ZSgc=
github.com/russross/blackfriday/v2 v2.1.0/go.mod h1:+Rmxgy9KzJVeS9/2gXHxylqXiyQDYRxCVz55jmeOWTM= github.com/russross/blackfriday/v2 v2.1.0/go.mod h1:+Rmxgy9KzJVeS9/2gXHxylqXiyQDYRxCVz55jmeOWTM=
github.com/shopspring/decimal v1.4.0 h1:bxl37RwXBklmTi0C79JfXCEBD1cqqHt0bbgBAGFp81k=
github.com/shopspring/decimal v1.4.0/go.mod h1:gawqmDU56v4yIKSwfBSFip1HdCCXN8/+DMd9qYNcwME=
github.com/spf13/cast v1.10.0 h1:h2x0u2shc1QuLHfxi+cTJvs30+ZAHOGRic8uyGTDWxY=
github.com/spf13/cast v1.10.0/go.mod h1:jNfB8QC9IA6ZuY2ZjDp0KtFO2LZZlg4S/7bzP6qqeHo=
github.com/spf13/cobra v1.10.2 h1:DMTTonx5m65Ic0GOoRY2c16WCbHxOOw6xxezuLaBpcU= github.com/spf13/cobra v1.10.2 h1:DMTTonx5m65Ic0GOoRY2c16WCbHxOOw6xxezuLaBpcU=
github.com/spf13/cobra v1.10.2/go.mod h1:7C1pvHqHw5A4vrJfjNwvOdzYu0Gml16OCs2GRiTUUS4= github.com/spf13/cobra v1.10.2/go.mod h1:7C1pvHqHw5A4vrJfjNwvOdzYu0Gml16OCs2GRiTUUS4=
github.com/spf13/pflag v1.0.9/go.mod h1:McXfInJRrz4CZXVZOBLb0bTZqETkiAhM9Iw0y3An2Bg= github.com/spf13/pflag v1.0.9/go.mod h1:McXfInJRrz4CZXVZOBLb0bTZqETkiAhM9Iw0y3An2Bg=
@@ -138,8 +166,9 @@ github.com/stretchr/objx v0.1.0/go.mod h1:HFkY916IF+rwdDfMAkV7OtwuqBVzrE8GR6GFx+
github.com/stretchr/objx v0.5.3 h1:jmXUvGomnU1o3W/V5h2VEradbpJDwGrzugQQvL0POH4= 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/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.3.0/go.mod h1:M5WIy9Dh21IEIfnGCwXGc5bZfKNJtfHm1UVUgZn+9EI=
github.com/stretchr/testify v1.11.1 h1:7s2iGBzp5EwR7/aIZr8ao5+dra3wiQyKjjFuvgVKu7U= github.com/stretchr/testify v1.7.0/go.mod h1:6Fq8oRcR53rry900zMqJjRRixrwX3KX962/h/Wwjteg=
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 h1:qLwI1I70+NjRFUR3zs1JPUCgaCXSh3SW62uAKT1mSBM=
github.com/x448/float16 v0.8.4/go.mod h1:14CWIYCyZA/cWjXOioeEpHeN/83MdbZDRQHoFcYsOfg= github.com/x448/float16 v0.8.4/go.mod h1:14CWIYCyZA/cWjXOioeEpHeN/83MdbZDRQHoFcYsOfg=
go.opentelemetry.io/auto/sdk v1.2.1 h1:jXsnJ4Lmnqd11kwkBV2LgLoFMZKizbCi5fNZ/ipaZ64= go.opentelemetry.io/auto/sdk v1.2.1 h1:jXsnJ4Lmnqd11kwkBV2LgLoFMZKizbCi5fNZ/ipaZ64=
@@ -170,12 +199,15 @@ go.uber.org/zap v1.27.1 h1:08RqriUEv8+ArZRYSTXy1LeBScaMpVSTBhCeaZYfMYc=
go.uber.org/zap v1.27.1/go.mod h1:GB2qFLM7cTU87MWRP2mPIjqfIDnGu+VIO4V/SdhGo2E= go.uber.org/zap v1.27.1/go.mod h1:GB2qFLM7cTU87MWRP2mPIjqfIDnGu+VIO4V/SdhGo2E=
go.yaml.in/yaml/v2 v2.4.4 h1:tuyd0P+2Ont/d6e2rl3be67goVK4R6deVxCUX5vyPaQ= go.yaml.in/yaml/v2 v2.4.4 h1:tuyd0P+2Ont/d6e2rl3be67goVK4R6deVxCUX5vyPaQ=
go.yaml.in/yaml/v2 v2.4.4/go.mod h1:gMZqIpDtDqOfM0uNfy0SkpRhvUryYH0Z6wdMYcacYXQ= go.yaml.in/yaml/v2 v2.4.4/go.mod h1:gMZqIpDtDqOfM0uNfy0SkpRhvUryYH0Z6wdMYcacYXQ=
go.yaml.in/yaml/v3 v3.0.4 h1:tfq32ie2Jv2UxXFdLJdh3jXuOzWiL1fo0bu/FbuKpbc=
go.yaml.in/yaml/v3 v3.0.4/go.mod h1:DhzuOOF2ATzADvBadXxruRBLzYTpT36CKvDb3+aBEFg= go.yaml.in/yaml/v3 v3.0.4/go.mod h1:DhzuOOF2ATzADvBadXxruRBLzYTpT36CKvDb3+aBEFg=
go.yaml.in/yaml/v3 v3.0.5 h1:N6y/pJk8buWs9NY5ERU2HSMfm+IuD/OtfdAnq6kESPw=
go.yaml.in/yaml/v3 v3.0.5/go.mod h1:HVTZu1O7/Vkt2N+BFy8Zza+lnLsABggaTM2ZpNIGuKg=
golang.org/x/crypto v0.55.0 h1:+KWHjbgOaAQ66dh/YlkZKHlz9ZUlq61AFirAR9ntP8M=
golang.org/x/crypto v0.55.0/go.mod h1:uq0V9dE/fzQuJtbnL+2EhWOE63vo164FY8xqEnV9xis=
golang.org/x/exp v0.0.0-20260410095643-746e56fc9e2f h1:W3F4c+6OLc6H2lb//N1q4WpJkhzJCK5J6kUi1NTVXfM= golang.org/x/exp v0.0.0-20260410095643-746e56fc9e2f h1:W3F4c+6OLc6H2lb//N1q4WpJkhzJCK5J6kUi1NTVXfM=
golang.org/x/exp v0.0.0-20260410095643-746e56fc9e2f/go.mod h1:J1xhfL/vlindoeF/aINzNzt2Bket5bjo9sdOYzOsU80= golang.org/x/exp v0.0.0-20260410095643-746e56fc9e2f/go.mod h1:J1xhfL/vlindoeF/aINzNzt2Bket5bjo9sdOYzOsU80=
golang.org/x/mod v0.37.0 h1:vF1DjpVEshcIqoEaauuHebaLk1O1forxjxBaVn884JQ= golang.org/x/mod v0.38.0 h1:MECBjubtXD7yj4HrhIUcywNaGeNVUdfVnxmPajOk4yk=
golang.org/x/mod v0.37.0/go.mod h1:m8S8VeM9r4dzDwjrKO0a1sZP3YjeMamRRlD+fmR2Q/0= 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 h1:K5+3DljvIuDG9/Jv9rvyMywYNFCQ9RSUY6OOTTkT+tE=
golang.org/x/net v0.57.0/go.mod h1:KpXc8iv+r3XplLAG/f7Jsf9RPszJzdR0f58q9vGOuEU= golang.org/x/net v0.57.0/go.mod h1:KpXc8iv+r3XplLAG/f7Jsf9RPszJzdR0f58q9vGOuEU=
golang.org/x/oauth2 v0.36.0 h1:peZ/1z27fi9hUOFCAZaHyrpWG5lwe0RJEEEeH0ThlIs= golang.org/x/oauth2 v0.36.0 h1:peZ/1z27fi9hUOFCAZaHyrpWG5lwe0RJEEEeH0ThlIs=
@@ -186,12 +218,12 @@ golang.org/x/sys v0.47.0 h1:o7XGOvZQCADBQQ4Y7VNq2dRWQR7JmOUW8Kxx4ZsNgWs=
golang.org/x/sys v0.47.0/go.mod h1:4GL1E5IUh+htKOUEOaiffhrAeqysfVGipDYzABqnCmw= golang.org/x/sys v0.47.0/go.mod h1:4GL1E5IUh+htKOUEOaiffhrAeqysfVGipDYzABqnCmw=
golang.org/x/term v0.45.0 h1:NwWyBmoJCbfTHpxrWoZ9C6/VxOf7ic219I8xZZFdrf0= golang.org/x/term v0.45.0 h1:NwWyBmoJCbfTHpxrWoZ9C6/VxOf7ic219I8xZZFdrf0=
golang.org/x/term v0.45.0/go.mod h1:9aqxs0blBcrm/n0L9QW0aRVD+ktan8ssZromtqJC43w= golang.org/x/term v0.45.0/go.mod h1:9aqxs0blBcrm/n0L9QW0aRVD+ktan8ssZromtqJC43w=
golang.org/x/text v0.40.0 h1:Ub2Z6/xjgF1WrYQz2nuITOEegKFtiIy+rieRJ5lHZKs= golang.org/x/text v0.41.0 h1:vz/seA0lnX87Othu2f/0L24RcgrXD9/YFTSuGjj3rH8=
golang.org/x/text v0.40.0/go.mod h1:hpnzDAfGV753zIKo+wk3u1bVKCGPbrnF7+7LBF/UHVY= golang.org/x/text v0.41.0/go.mod h1:jvf1O8ajNzZqhSrQBPbutR/EB83Cc0CFrezNQIwbb5M=
golang.org/x/time v0.15.0 h1:bbrp8t3bGUeFOx08pvsMYRTCVSMk89u4tKbNOZbp88U= golang.org/x/time v0.15.0 h1:bbrp8t3bGUeFOx08pvsMYRTCVSMk89u4tKbNOZbp88U=
golang.org/x/time v0.15.0/go.mod h1:Y4YMaQmXwGQZoFaVFk4YpCt4FLQMYKZe9oeV/f4MSno= golang.org/x/time v0.15.0/go.mod h1:Y4YMaQmXwGQZoFaVFk4YpCt4FLQMYKZe9oeV/f4MSno=
golang.org/x/tools v0.47.0 h1:7Kn5x/d1svx/PzryTsqeoZN4TZwqeH5pGWjefhLi/1Q= golang.org/x/tools v0.48.0 h1:3+hClM1aLL5mjMKm5ovokw9epgRXPuu2tILgismM6RE=
golang.org/x/tools v0.47.0/go.mod h1:dFHnyTvFWY212G+h7ZY4Vsp/K3U4/7W9TyVaAul8uCA= golang.org/x/tools v0.48.0/go.mod h1:08xX0orndb/F7jJxGDicx061tyd5pcMto75YMAXr6lk=
gomodules.xyz/jsonpatch/v2 v2.4.0 h1:Ci3iUJyx9UeRx7CeFN8ARgGbkESwJK+KB9lLcWxY/Zw= gomodules.xyz/jsonpatch/v2 v2.4.0 h1:Ci3iUJyx9UeRx7CeFN8ARgGbkESwJK+KB9lLcWxY/Zw=
gomodules.xyz/jsonpatch/v2 v2.4.0/go.mod h1:AH3dM2RI6uoBZxn3LVrfvJ3E0/9dG4cSrbuBJT4moAY= gomodules.xyz/jsonpatch/v2 v2.4.0/go.mod h1:AH3dM2RI6uoBZxn3LVrfvJ3E0/9dG4cSrbuBJT4moAY=
gonum.org/v1/gonum v0.17.0 h1:VbpOemQlsSMrYmn7T2OUvQ4dqxQXU+ouZFQsZOx50z4= gonum.org/v1/gonum v0.17.0 h1:VbpOemQlsSMrYmn7T2OUvQ4dqxQXU+ouZFQsZOx50z4=
@@ -205,12 +237,11 @@ google.golang.org/grpc v1.82.1/go.mod h1:yzTZ1TB1Z3SG+LIYaI+WiE8D5+PZ3ArnrSp8zF3
google.golang.org/protobuf v1.36.12-0.20260120151049-f2248ac996af h1:+5/Sw3GsDNlEmu7TfklWKPdQ0Ykja5VEmq2i817+jbI= google.golang.org/protobuf v1.36.12-0.20260120151049-f2248ac996af h1:+5/Sw3GsDNlEmu7TfklWKPdQ0Ykja5VEmq2i817+jbI=
google.golang.org/protobuf v1.36.12-0.20260120151049-f2248ac996af/go.mod h1:HTf+CrKn2C3g5S8VImy6tdcUvCska2kB7j23XfzDpco= google.golang.org/protobuf v1.36.12-0.20260120151049-f2248ac996af/go.mod h1:HTf+CrKn2C3g5S8VImy6tdcUvCska2kB7j23XfzDpco=
gopkg.in/check.v1 v0.0.0-20161208181325-20d25e280405/go.mod h1:Co6ibVJAznAaIkqp8huTwlJQCZ016jof/cbN4VW5Yz0= gopkg.in/check.v1 v0.0.0-20161208181325-20d25e280405/go.mod h1:Co6ibVJAznAaIkqp8huTwlJQCZ016jof/cbN4VW5Yz0=
gopkg.in/check.v1 v1.0.0-20201130134442-10cb98267c6c h1:Hei/4ADfdWqJk1ZMxUNpqntNwaWcugrBjAiHlqqRiVk=
gopkg.in/check.v1 v1.0.0-20201130134442-10cb98267c6c/go.mod h1:JHkPIbrfpd72SG/EVd6muEfDQjcINNoR0C8j2r3qZ4Q=
gopkg.in/evanphx/json-patch.v4 v4.13.0 h1:czT3CmqEaQ1aanPc5SdlgQrrEIb8w/wwCvWWnfEbYzo= gopkg.in/evanphx/json-patch.v4 v4.13.0 h1:czT3CmqEaQ1aanPc5SdlgQrrEIb8w/wwCvWWnfEbYzo=
gopkg.in/evanphx/json-patch.v4 v4.13.0/go.mod h1:p8EYWUEYMpynmqDbY58zCKCFZw8pRWMG4EsWvDvM72M= gopkg.in/evanphx/json-patch.v4 v4.13.0/go.mod h1:p8EYWUEYMpynmqDbY58zCKCFZw8pRWMG4EsWvDvM72M=
gopkg.in/inf.v0 v0.9.1 h1:73M5CoZyi3ZLMOyDlQh031Cx6N9NDJ2Vvfl76EDAgDc= gopkg.in/inf.v0 v0.9.1 h1:73M5CoZyi3ZLMOyDlQh031Cx6N9NDJ2Vvfl76EDAgDc=
gopkg.in/inf.v0 v0.9.1/go.mod h1:cWUDdTG/fYaXco+Dcufb5Vnc6Gp2YChqWtbxRZE0mXw= gopkg.in/inf.v0 v0.9.1/go.mod h1:cWUDdTG/fYaXco+Dcufb5Vnc6Gp2YChqWtbxRZE0mXw=
gopkg.in/yaml.v3 v3.0.0-20200313102051-9f266ea9e77c/go.mod h1:K4uyk7z7BCEPqu6E+C64Yfv1cQ7kz7rIZviUmN+EgEM=
gopkg.in/yaml.v3 v3.0.1 h1:fxVm/GzAzEWqLHuvctI91KS9hhNmmWOoWu0XTYJS7CA= gopkg.in/yaml.v3 v3.0.1 h1:fxVm/GzAzEWqLHuvctI91KS9hhNmmWOoWu0XTYJS7CA=
gopkg.in/yaml.v3 v3.0.1/go.mod h1:K4uyk7z7BCEPqu6E+C64Yfv1cQ7kz7rIZviUmN+EgEM= gopkg.in/yaml.v3 v3.0.1/go.mod h1:K4uyk7z7BCEPqu6E+C64Yfv1cQ7kz7rIZviUmN+EgEM=
k8s.io/api v0.37.0 h1:Z//Vj9N7RA/yS2sDmxyeo7h+RR4zbUrd2vrd3Z0TbB4= k8s.io/api v0.37.0 h1:Z//Vj9N7RA/yS2sDmxyeo7h+RR4zbUrd2vrd3Z0TbB4=
@@ -0,0 +1,68 @@
/*
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 kubernetes 提供 Database 所需的 Kubernetes API 薄适配。
package kubernetes
import (
"context"
"errors"
corev1 "k8s.io/api/core/v1"
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
"k8s.io/apimachinery/pkg/util/validation"
typedcore "k8s.io/client-go/kubernetes/typed/core/v1"
"k8s.io/client-go/rest"
"git.ddupan.top/panxiao81/ayatori/internal/database/application"
"git.ddupan.top/panxiao81/ayatori/internal/database/domain/instance"
)
// SecretCredentials 直接读取 API server,不将 Secret 数据纳入共享 informer cache。
// namespace 在装配时固定,Instance 不能选择跨 namespace 读取。
type SecretCredentials struct {
secrets typedcore.SecretInterface
}
func NewSecretCredentials(config *rest.Config, namespace string) (*SecretCredentials, error) {
if config == nil || len(validation.IsDNS1123Label(namespace)) != 0 {
return nil, errors.New("valid controller namespace and API configuration required")
}
client, err := typedcore.NewForConfig(config)
if err != nil {
return nil, application.ErrCredentialsUnavailable
}
return &SecretCredentials{secrets: client.Secrets(namespace)}, nil
}
func (r *SecretCredentials) Read(ctx context.Context, ref instance.CredentialReference) (application.Credentials, error) {
if err := ref.Validate(); err != nil {
return application.Credentials{}, application.ErrCredentialsInvalid
}
keys := ref.Values()
secret, err := r.secrets.Get(ctx, keys.Name, metav1.GetOptions{})
if err != nil {
return application.Credentials{}, application.ErrCredentialsUnavailable
}
return decode(secret, keys)
}
func decode(secret *corev1.Secret, keys instance.CredentialReferenceValues) (application.Credentials, error) {
if secret.DeletionTimestamp != nil {
return application.Credentials{}, application.ErrCredentialsUnavailable
}
return application.NewCredentials(string(secret.Data[keys.UsernameKey]), string(secret.Data[keys.PasswordKey]))
}
@@ -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 postgresql 使用 pgxpool 提供 PostgreSQL 能力的薄适配。
package postgresql
import (
"context"
"crypto/tls"
"errors"
"net"
"net/url"
"strconv"
"github.com/jackc/pgx/v5/pgconn"
"github.com/jackc/pgx/v5/pgxpool"
"git.ddupan.top/panxiao81/ayatori/internal/database/application"
"git.ddupan.top/panxiao81/ayatori/internal/database/domain/instance"
)
// Connector 不读取 Secret、不决定连接何时替换;池本身由 pgxpool 实现。
type Connector struct {
RootCert string
}
type database struct {
pool *pgxpool.Pool
}
func (*database) String() string { return "[redacted PostgreSQL database]" }
func (d *database) GoString() string { return d.String() }
func (d *database) Close() {
d.pool.Close()
}
func (d *database) Version(ctx context.Context) (string, error) {
var version string
if err := d.pool.QueryRow(ctx, "SHOW server_version").Scan(&version); err != nil {
return "", safeError(err, application.ErrObservation)
}
return version, nil
}
func (c Connector) Connect(ctx context.Context, endpoint instance.Endpoint, credentials application.Credentials) (application.Database, error) {
if err := endpoint.Validate(); err != nil {
return nil, err
}
if credentials.Username() == "" || credentials.Password() == "" {
return nil, application.ErrCredentialsInvalid
}
endpointValues := endpoint.Values()
query := url.Values{
"sslmode": {string(endpointValues.TLSMode)},
"connect_timeout": {"5"},
"application_name": {"ayatori-database-management"},
}
if c.RootCert != "" {
query.Set("sslrootcert", c.RootCert)
}
connectionURL := url.URL{
Scheme: "postgresql",
Host: net.JoinHostPort(endpointValues.Host, strconv.Itoa(endpointValues.Port)),
Path: "/" + endpointValues.ManagementDatabase,
User: url.UserPassword(credentials.Username(), credentials.Password()),
RawQuery: query.Encode(),
}
config, err := pgxpool.ParseConfig(connectionURL.String())
if err != nil {
return nil, application.ErrConnection
}
// pgx 不实现 libpq hostaddr;复用其 LookupFunc 扩展点,TLS 验证身份仍采用 host。
config.ConnConfig.LookupFunc = func(context.Context, string) ([]string, error) {
return []string{endpointValues.HostAddr}, nil
}
config.ConnConfig.Fallbacks = nil
pool, err := pgxpool.NewWithConfig(ctx, config)
if err != nil {
return nil, safeError(err, application.ErrConnection)
}
if err := pool.Ping(ctx); err != nil {
pool.Close()
return nil, safeError(err, application.ErrConnection)
}
return &database{pool: pool}, nil
}
func safeError(err, fallback error) error {
if errors.Is(err, context.Canceled) {
return context.Canceled
}
if errors.Is(err, context.DeadlineExceeded) {
return context.DeadlineExceeded
}
var pgerr *pgconn.PgError
if errors.As(err, &pgerr) && (pgerr.Code == "28P01" || pgerr.Code == "28000") {
return application.ErrAuthentication
}
if _, ok := errors.AsType[*tls.CertificateVerificationError](err); ok {
return application.ErrAuthentication
}
return fallback
}
@@ -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 postgresql_test
import (
"context"
"errors"
"os/exec"
"sync"
"testing"
"time"
corev1 "k8s.io/api/core/v1"
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
"git.ddupan.top/panxiao81/ayatori/internal/database/adapter/postgresql"
"git.ddupan.top/panxiao81/ayatori/internal/database/application"
"git.ddupan.top/panxiao81/ayatori/internal/database/domain/instance"
)
func TestManagementSecretScopeAndMissingDependencyRecovery(t *testing.T) {
fixture := newCredentialFixture(t)
reference := fixture.target.Definition().AdminCredential()
fixture.createSecret(t, "unrelated")
if _, err := fixture.reader.Read(fixture.ctx, reference); !errors.Is(err, application.ErrCredentialsUnavailable) {
t.Fatal("a Secret in another namespace satisfied the reference")
}
if _, err := fixture.service.ObserveVersion(fixture.ctx, fixture.target); !errors.Is(err, application.ErrCredentialsUnavailable) {
t.Fatal("missing Secret did not fail closed")
}
fixture.createSecret(t, controllerNamespace)
if _, err := fixture.deniedReader.Read(fixture.ctx, reference); !errors.Is(err, application.ErrCredentialsUnavailable) {
t.Fatal("API server did not enforce Secret RBAC")
}
fixture.observeVersion(t)
fixture.updateSecret(t, func(secret *corev1.Secret) {
delete(secret.Data, "credential")
})
if _, err := fixture.service.ObserveVersion(fixture.ctx, fixture.target); !errors.Is(err, application.ErrCredentialsInvalid) {
t.Fatal("missing credential field reused a cached connection")
}
fixture.updateSecret(t, func(secret *corev1.Secret) {
secret.Data["credential"] = []byte(fixturePassword)
})
fixture.observeVersion(t)
err := fixture.client.CoreV1().Secrets(controllerNamespace).Delete(fixture.ctx, secretName, metav1.DeleteOptions{})
if err != nil {
t.Fatal("cannot delete fixture Secret")
}
if _, err := fixture.service.ObserveVersion(fixture.ctx, fixture.target); !errors.Is(err, application.ErrCredentialsUnavailable) {
t.Fatal("deleted Secret retained access")
}
}
func TestEffectiveCredentialChangesReplaceConnection(t *testing.T) {
fixture := newCredentialFixture(t)
fixture.createSecret(t, controllerNamespace)
fixture.observeVersion(t)
originalBackend := fixture.backendIDs(t)
if originalBackend == "" {
t.Fatal("management connection not visible in PostgreSQL")
}
fixture.updateSecret(t, func(secret *corev1.Secret) {
secret.Labels = map[string]string{"changed": "true"}
secret.Data["unrelated"] = []byte("ignored")
})
fixture.observeVersion(t)
if fixture.backendIDs(t) != originalBackend {
t.Fatal("metadata or unrelated fields rebuilt the connection")
}
// 先改变 Secret、暂不改变服务器密码:旧连接必须失效,新认证必须失败。
fixture.updateSecret(t, func(secret *corev1.Secret) {
secret.Data["credential"] = []byte(rotatedPassword)
})
version, err := fixture.service.ObserveVersion(fixture.ctx, fixture.target)
if !errors.Is(err, application.ErrAuthentication) || version != "" {
t.Fatal("old connection bypassed changed credentials")
}
fixture.queryPostgres(t, "ALTER ROLE postgres PASSWORD '"+rotatedPassword+"'")
fixture.observeVersion(t)
if fixture.backendIDs(t) == originalBackend {
t.Fatal("password rotation reused the old backend")
}
fixture.updateSecret(t, func(secret *corev1.Secret) {
secret.Data["login"] = []byte("nonexistent")
})
if _, err := fixture.service.ObserveVersion(fixture.ctx, fixture.target); !errors.Is(err, application.ErrAuthentication) {
t.Fatal("username change did not require a new authentication")
}
fixture.updateSecret(t, func(secret *corev1.Secret) {
secret.Data["login"] = []byte(fixtureUser)
})
fixture.observeVersion(t)
}
func TestObservationDiscardsResultWhenCredentialsChange(t *testing.T) {
fixture := newCredentialFixture(t)
fixture.createSecret(t, controllerNamespace)
reads := 0
fixture.gate.beforeRead = func() {
reads++
if reads == 2 {
fixture.updateSecret(t, func(secret *corev1.Secret) {
secret.Data["credential"] = []byte(rotatedPassword)
})
}
}
version, err := fixture.service.ObserveVersion(fixture.ctx, fixture.target)
if !errors.Is(err, application.ErrCredentialsChanged) {
t.Fatal("in-flight rotation was not detected")
}
if version != "" {
t.Fatal("observation returned data obtained with stale credentials")
}
if fixture.backendIDs(t) != "" {
t.Fatal("stale connection was retained after rotation")
}
}
func TestConnectionReleaseAndServiceRestart(t *testing.T) {
fixture := newCredentialFixture(t)
fixture.createSecret(t, controllerNamespace)
fixture.observeVersion(t)
fixture.service.Forget(fixture.target.Identity().Name())
if fixture.backendIDs(t) != "" {
t.Fatal("Forget retained a connection")
}
fixture.observeVersion(t)
fixture.service.Close()
fixture.service.Close()
if _, err := fixture.service.ObserveVersion(fixture.ctx, fixture.target); !errors.Is(err, application.ErrClosed) {
t.Fatal("closed service accepted work")
}
restarted, err := application.NewInstanceService(fixture.reader, postgresql.Connector{})
if err != nil {
t.Fatal(err)
}
t.Cleanup(restarted.Close)
if _, err := restarted.ObserveVersion(fixture.ctx, fixture.target); err != nil {
t.Fatal("new service could not recover from stored Secret", err)
}
// 此 fixture 未启用 TLS;各加密模式均不得偷偷回退到明文连接。
for _, mode := range []instance.TLSMode{instance.TLSRequire, instance.TLSVerifyCA, instance.TLSVerifyFull} {
securedTarget := target(t, fixture.port, mode)
if _, err := restarted.ObserveVersion(fixture.ctx, securedTarget); err == nil {
t.Fatal("TLS policy silently downgraded to plaintext")
}
}
}
func TestConcurrentVersionObservations(t *testing.T) {
fixture := newCredentialFixture(t)
fixture.createSecret(t, controllerNamespace)
var workers sync.WaitGroup
for range 4 {
workers.Go(func() {
version, err := fixture.service.ObserveVersion(fixture.ctx, fixture.target)
if err != nil || version == "" {
t.Error("concurrent observation failed", err)
}
})
}
workers.Go(func() {
fixture.service.Forget(fixture.target.Identity().Name())
})
workers.Wait()
fixture.observeVersion(t)
}
func TestManagementConnectionRecoversAfterTimeout(t *testing.T) {
fixture := newCredentialFixture(t)
fixture.createSecret(t, controllerNamespace)
fixture.observeVersion(t)
if err := exec.CommandContext(fixture.ctx, "docker", "pause", fixture.containerID).Run(); err != nil {
t.Fatal("cannot pause isolated PostgreSQL fixture")
}
// 即使断言失败,也先恢复容器,再由 fixture 按原 ID 清理。
t.Cleanup(func() {
cleanupContext, cancel := context.WithTimeout(context.Background(), 10*time.Second)
defer cancel()
_ = exec.CommandContext(cleanupContext, "docker", "unpause", fixture.containerID).Run()
})
queryContext, cancel := context.WithTimeout(fixture.ctx, 500*time.Millisecond)
version, err := fixture.service.ObserveVersion(queryContext, fixture.target)
cancel()
if err == nil || version != "" {
t.Fatal("timed out PostgreSQL observation returned a successful result")
}
if err := exec.CommandContext(fixture.ctx, "docker", "unpause", fixture.containerID).Run(); err != nil {
t.Fatal("cannot resume isolated PostgreSQL fixture")
}
fixture.observeVersion(t)
}
@@ -0,0 +1,332 @@
//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 postgresql_test
import (
"context"
"errors"
"os/exec"
"regexp"
"strconv"
"strings"
"testing"
"time"
corev1 "k8s.io/api/core/v1"
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
"k8s.io/client-go/kubernetes"
"sigs.k8s.io/controller-runtime/pkg/envtest"
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"
"git.ddupan.top/panxiao81/ayatori/internal/database/domain/instance"
)
const (
fixtureHost = "fixture.invalid"
fixtureUser = "postgres"
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 清理。
func postgresFixture(t *testing.T, ctx context.Context) (string, int) {
t.Helper()
output, err := exec.CommandContext(ctx, "docker", "run", "--rm", "-d", "-p", "127.0.0.1::5432",
"-e", "POSTGRES_PASSWORD="+fixturePassword, fixtureImage).Output()
if err != nil {
t.Fatalf("cannot start isolated PostgreSQL fixture: %s", fixtureCommandError(err))
}
id := strings.TrimSpace(string(output))
if !regexp.MustCompile(`^[a-f0-9]{64}$`).MatchString(id) {
t.Fatal("unexpected container identifier")
}
t.Cleanup(func() {
cleanup, cancel := context.WithTimeout(context.Background(), 30*time.Second)
defer cancel()
if exec.CommandContext(cleanup, "docker", "rm", "-f", id).Run() != nil {
t.Error("fixture cleanup failed")
}
})
output, err = exec.CommandContext(ctx, "docker", "inspect", "--format", `{{(index (index .NetworkSettings.Ports "5432/tcp") 0).HostPort}}`, id).Output()
if err != nil {
t.Fatalf("cannot inspect fixture port: %s", fixtureCommandError(err))
}
port, err := strconv.Atoi(strings.TrimSpace(string(output)))
if err != nil {
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 {
select {
case <-ctx.Done():
t.Fatal("fixture startup timed out")
case <-time.After(200 * time.Millisecond):
}
}
return id, port
}
// Output 将 stderr 保存在 ExitError 中;保留诊断,但不打印命令参数和测试密码。
func fixtureCommandError(err error) string {
detail := err.Error()
if exitErr, ok := errors.AsType[*exec.ExitError](err); ok {
detail += ": " + strings.TrimSpace(string(exitErr.Stderr))
}
redactor := strings.NewReplacer(
fixturePassword, "[REDACTED]",
rotatedPassword, "[REDACTED]",
)
return redactor.Replace(detail)
}
func TestFixtureCommandErrorPreservesDiagnosticsAndRedactsPasswords(t *testing.T) {
tests := []struct {
name string
err error
want string
}{
{
name: "missing docker executable",
err: &exec.Error{Name: "docker", Err: exec.ErrNotFound},
want: "executable file not found",
},
{
name: "daemon failure from stderr",
err: &exec.ExitError{Stderr: []byte("Cannot connect to the Docker daemon")},
want: "Cannot connect to the Docker daemon",
},
{
name: "passwords in stderr",
err: &exec.ExitError{Stderr: []byte("failure: " + fixturePassword + " " + rotatedPassword)},
want: "failure: [REDACTED] [REDACTED]",
},
{
name: "password in error text",
err: errors.New("failure: " + fixturePassword),
want: "failure: [REDACTED]",
},
}
for _, tt := range tests {
t.Run(tt.name, func(t *testing.T) {
detail := fixtureCommandError(tt.err)
if !strings.Contains(detail, tt.want) {
t.Fatalf("diagnostic lost expected information: %q", tt.want)
}
if strings.Contains(detail, fixturePassword) || strings.Contains(detail, rotatedPassword) {
t.Fatal("diagnostic exposed a fixture password")
}
})
}
}
func target(t *testing.T, port int, mode instance.TLSMode) instance.ObservationTarget {
t.Helper()
id, err := instance.NewIdentity("fixture-uid", "fixture")
if err != nil {
t.Fatal(err)
}
revision, err := instance.NewRevision(1)
if err != nil {
t.Fatal(err)
}
endpoint, err := instance.NewEndpoint(instance.EndpointValues{
Host: fixtureHost,
HostAddr: "127.0.0.1",
Port: port,
ManagementDatabase: fixtureUser,
TLSMode: mode,
})
if err != nil {
t.Fatal(err)
}
ref, err := instance.NewCredentialReference(instance.CredentialReferenceValues{
Name: secretName,
UsernameKey: "login",
PasswordKey: "credential",
})
if err != nil {
t.Fatal(err)
}
definition, err := instance.NewDefinition(endpoint, ref)
if err != nil {
t.Fatal(err)
}
value, err := instance.NewObservationTarget(id, revision, definition)
if err != nil {
t.Fatal(err)
}
return value
}
// 在真实读取前设置屏障,确定性验证观测期间 Secret 变化;实际数据仍来自 API server。
type gatedReader struct {
application.CredentialReader
beforeRead func()
}
func (r *gatedReader) Read(ctx context.Context, ref instance.CredentialReference) (application.Credentials, error) {
if r.beforeRead != nil {
r.beforeRead()
}
return r.CredentialReader.Read(ctx, ref)
}
// credentialFixture 为每个场景创建独立 API server、PostgreSQL 和应用服务。
type credentialFixture struct {
ctx context.Context
client *kubernetes.Clientset
reader *secretadapter.SecretCredentials
deniedReader *secretadapter.SecretCredentials
gate *gatedReader
service *application.InstanceService
target instance.ObservationTarget
containerID string
port int
}
func newCredentialFixture(t *testing.T) *credentialFixture {
t.Helper()
ctx, cancel := context.WithTimeout(context.Background(), 2*time.Minute)
t.Cleanup(cancel)
environment := &envtest.Environment{}
config, err := environment.Start()
if err != nil {
t.Fatal("envtest startup failed", err)
}
t.Cleanup(func() {
if err := environment.Stop(); err != nil {
t.Error("envtest cleanup failed", err)
}
})
client, err := kubernetes.NewForConfig(config)
if err != nil {
t.Fatal("cannot create test client")
}
for _, namespace := range []string{controllerNamespace, "unrelated"} {
_, err := client.CoreV1().Namespaces().Create(
ctx,
&corev1.Namespace{Name: namespace},
metav1.CreateOptions{},
)
if err != nil {
t.Fatal("cannot create fixture namespace")
}
}
reader, err := secretadapter.NewSecretCredentials(config, controllerNamespace)
if err != nil {
t.Fatal(err)
}
user, err := environment.AddUser(envtest.User{Name: "without-secret-access"}, config)
if err != nil {
t.Fatal(err)
}
deniedReader, err := secretadapter.NewSecretCredentials(user.Config(), controllerNamespace)
if err != nil {
t.Fatal(err)
}
containerID, port := postgresFixture(t, ctx)
gate := &gatedReader{CredentialReader: reader}
service, err := application.NewInstanceService(gate, postgresql.Connector{})
if err != nil {
t.Fatal(err)
}
t.Cleanup(service.Close)
return &credentialFixture{
ctx: ctx,
client: client,
reader: reader,
deniedReader: deniedReader,
gate: gate,
service: service,
target: target(t, port, instance.TLSDisable),
containerID: containerID,
port: port,
}
}
func (f *credentialFixture) createSecret(t *testing.T, namespace string) {
t.Helper()
secret := &corev1.Secret{
Name: secretName,
Data: map[string][]byte{
"login": []byte(fixtureUser),
"credential": []byte(fixturePassword),
},
}
if _, err := f.client.CoreV1().Secrets(namespace).Create(f.ctx, secret, metav1.CreateOptions{}); err != nil {
t.Fatal("cannot create fixture Secret")
}
}
func (f *credentialFixture) updateSecret(t *testing.T, change func(*corev1.Secret)) {
t.Helper()
secrets := f.client.CoreV1().Secrets(controllerNamespace)
secret, err := secrets.Get(f.ctx, secretName, metav1.GetOptions{})
if err != nil {
t.Fatal("cannot read fixture Secret")
}
change(secret)
if _, err := secrets.Update(f.ctx, secret, metav1.UpdateOptions{}); err != nil {
t.Fatal("cannot update fixture Secret")
}
}
func (f *credentialFixture) observeVersion(t *testing.T) {
t.Helper()
version, err := f.service.ObserveVersion(f.ctx, f.target)
if err != nil {
t.Fatal("version observation failed", err)
}
if version == "" {
t.Fatal("successful observation returned an empty version")
}
}
func (f *credentialFixture) queryPostgres(t *testing.T, sql string) string {
t.Helper()
output, err := exec.CommandContext(
f.ctx, "docker", dockerExec, f.containerID,
"psql", "-U", fixtureUser, "-tAc", sql,
).Output()
if err != nil {
t.Fatal("fixture SQL failed")
}
return strings.TrimSpace(string(output))
}
func (f *credentialFixture) backendIDs(t *testing.T) string {
t.Helper()
return f.queryPostgres(t, `
SELECT pid
FROM pg_stat_activity
WHERE application_name = 'ayatori-database-management'
ORDER BY pid
`)
}
@@ -0,0 +1,25 @@
CREATE SCHEMA postgresql_tenant_operator;
CREATE TABLE postgresql_tenant_operator.tenant_ownership (
instance_uid text NOT NULL,
tenant_uid text NOT NULL,
tenant_namespace text NOT NULL,
tenant_name text NOT NULL,
database_name text NOT NULL,
role_name text NOT NULL,
credential_path text NOT NULL,
managed boolean NOT NULL DEFAULT true,
created_at timestamptz NOT NULL DEFAULT clock_timestamp(),
updated_at timestamptz NOT NULL DEFAULT clock_timestamp(),
retained_at timestamptz,
PRIMARY KEY (instance_uid, tenant_uid),
UNIQUE (instance_uid, tenant_namespace, tenant_name),
UNIQUE (instance_uid, database_name),
UNIQUE (instance_uid, role_name),
UNIQUE (credential_path),
CHECK (managed OR retained_at IS NOT NULL)
);
---- create above / drop below ----
DROP SCHEMA postgresql_tenant_operator CASCADE;
@@ -0,0 +1,326 @@
/*
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 registry persists controller ownership evidence in the PostgreSQL
// management database. It deliberately does not persist reconciliation phases;
// those belong to the Kubernetes resource status.
package registry
import (
"context"
"embed"
"errors"
"fmt"
"io/fs"
"time"
"github.com/jackc/pgx/v5"
"github.com/jackc/pgx/v5/pgxpool"
"github.com/jackc/tern/v2/migrate"
)
const migrationVersionTable = "public.postgresql_tenant_operator_schema_version"
//go:embed migrations/*.sql
var migrationFiles embed.FS
var (
// ErrNotFound indicates that the registry has no matching ownership record.
ErrNotFound = errors.New("registry ownership record not found")
// ErrConflict indicates that a requested name or path belongs to another Tenant UID.
ErrConflict = errors.New("registry ownership conflict")
)
// Beginner is implemented by pgx.Conn and pgxpool.Pool.
type Beginner interface {
Begin(context.Context) (pgx.Tx, error)
}
// Store manages ownership records in one PostgreSQLInstance management database.
type Store struct {
db Beginner
}
// NewStore creates a registry store backed by a PostgreSQL connection or pool.
func NewStore(db Beginner) *Store {
return &Store{db: db}
}
// Ownership identifies every external resource reserved for one Tenant UID.
type Ownership struct {
InstanceUID string
TenantUID string
TenantNamespace string
TenantName string
DatabaseName string
RoleName string
CredentialPath string
}
// Record is the persisted ownership state.
type Record struct {
Ownership
Managed bool
CreatedAt time.Time
UpdatedAt time.Time
RetainedAt *time.Time
}
// ClaimResult reports whether Claim inserted a new record or found the same claim.
type ClaimResult string
const (
ClaimCreated ClaimResult = "Created"
ClaimOwned ClaimResult = "Owned"
)
// Bootstrap applies all pending versioned registry migrations idempotently.
func (s *Store) Bootstrap(ctx context.Context) error {
if s == nil || s.db == nil {
return errors.New("bootstrap registry: nil database")
}
return s.withMigrationConnection(ctx, func(conn *pgx.Conn) error {
migrations, err := fs.Sub(migrationFiles, "migrations")
if err != nil {
return fmt.Errorf("open embedded registry migrations: %w", err)
}
migrator, err := migrate.NewMigrator(ctx, conn, migrationVersionTable)
if err != nil {
return fmt.Errorf("initialize registry migrator: %w", err)
}
if err := migrator.LoadMigrations(migrations); err != nil {
return fmt.Errorf("load registry migrations: %w", err)
}
if err := migrator.Migrate(ctx); err != nil {
return fmt.Errorf("apply registry migrations: %w", err)
}
return nil
})
}
// Claim reserves all names for an owner. Repeating an identical claim is idempotent.
func (s *Store) Claim(ctx context.Context, owner Ownership) (ClaimResult, error) {
if s == nil || s.db == nil {
return "", errors.New("claim registry ownership: nil database")
}
if err := owner.validate(); err != nil {
return "", err
}
tx, err := s.db.Begin(ctx)
if err != nil {
return "", fmt.Errorf("begin registry claim: %w", err)
}
defer func() { _ = tx.Rollback(ctx) }()
_, err = scanRecord(tx.QueryRow(ctx, claimStatement,
owner.InstanceUID,
owner.TenantUID,
owner.TenantNamespace,
owner.TenantName,
owner.DatabaseName,
owner.RoleName,
owner.CredentialPath,
))
if err != nil && !errors.Is(err, ErrNotFound) {
return "", fmt.Errorf("insert registry claim: %w", err)
}
if errors.Is(err, ErrNotFound) {
record, err := getByTenantUID(ctx, tx, owner.InstanceUID, owner.TenantUID)
if err != nil {
if errors.Is(err, ErrNotFound) {
return "", fmt.Errorf("%w: database, role, tenant identity, or credential path is already reserved", ErrConflict)
}
return "", err
}
if !record.equal(owner) || !record.Managed {
return "", fmt.Errorf("%w: database, role, tenant identity, or credential path is already reserved", ErrConflict)
}
if err := tx.Commit(ctx); err != nil {
return "", fmt.Errorf("commit registry claim: %w", err)
}
return ClaimOwned, nil
}
if err := tx.Commit(ctx); err != nil {
return "", fmt.Errorf("commit registry claim: %w", err)
}
return ClaimCreated, nil
}
func (s *Store) withMigrationConnection(ctx context.Context, run func(*pgx.Conn) error) error {
switch db := s.db.(type) {
case *pgx.Conn:
return run(db)
case *pgxpool.Pool:
conn, err := db.Acquire(ctx)
if err != nil {
return fmt.Errorf("acquire registry migration connection: %w", err)
}
defer conn.Release()
return run(conn.Conn())
default:
return fmt.Errorf("bootstrap registry: database type %T cannot provide a migration connection", s.db)
}
}
// Get returns the ownership record for an Instance UID and Tenant UID.
func (s *Store) Get(ctx context.Context, instanceUID, tenantUID string) (Record, error) {
if s == nil || s.db == nil {
return Record{}, errors.New("get registry record: nil database")
}
if instanceUID == "" || tenantUID == "" {
return Record{}, errors.New("get registry record: instance UID and tenant UID are required")
}
tx, err := s.db.Begin(ctx)
if err != nil {
return Record{}, fmt.Errorf("begin registry read: %w", err)
}
defer func() { _ = tx.Rollback(ctx) }()
record, err := getByTenantUID(ctx, tx, instanceUID, tenantUID)
if err != nil {
return Record{}, err
}
if err := tx.Commit(ctx); err != nil {
return Record{}, fmt.Errorf("commit registry read: %w", err)
}
return record, nil
}
// MarkRetained changes a matching managed ownership record into an unmanaged tombstone.
// Repeating the operation for the same tombstone is safe.
func (s *Store) MarkRetained(ctx context.Context, owner Ownership) error {
return s.changeOwnership(ctx, owner, "mark registry record retained", func(ctx context.Context, tx pgx.Tx, record Record) error {
if !record.Managed {
return nil
}
tag, err := tx.Exec(ctx, markRetainedStatement, owner.InstanceUID, owner.TenantUID)
if err != nil {
return err
}
if tag.RowsAffected() != 1 {
return ErrConflict
}
return nil
})
}
// Delete removes a matching managed record after its external resources have been deleted.
// An already absent record is treated as a successful retry; a retained record is never deleted.
func (s *Store) Delete(ctx context.Context, owner Ownership) error {
return s.changeOwnership(ctx, owner, "delete registry record", func(ctx context.Context, tx pgx.Tx, record Record) error {
if !record.Managed {
return ErrConflict
}
tag, err := tx.Exec(ctx, deleteStatement, owner.InstanceUID, owner.TenantUID)
if err != nil {
return err
}
if tag.RowsAffected() != 1 {
return ErrConflict
}
return nil
})
}
func (s *Store) changeOwnership(
ctx context.Context,
owner Ownership,
operation string,
change func(context.Context, pgx.Tx, Record) error,
) error {
if s == nil || s.db == nil {
return fmt.Errorf("%s: nil database", operation)
}
if err := owner.validate(); err != nil {
return fmt.Errorf("%s: %w", operation, err)
}
tx, err := s.db.Begin(ctx)
if err != nil {
return fmt.Errorf("begin %s: %w", operation, err)
}
defer func() { _ = tx.Rollback(ctx) }()
record, err := getByTenantUIDForUpdate(ctx, tx, owner.InstanceUID, owner.TenantUID)
if errors.Is(err, ErrNotFound) && operation == "delete registry record" {
return nil
}
if err != nil {
return fmt.Errorf("%s: %w", operation, err)
}
if !record.equal(owner) {
return fmt.Errorf("%s: %w", operation, ErrConflict)
}
if err := change(ctx, tx, record); err != nil {
return fmt.Errorf("%s: %w", operation, err)
}
if err := tx.Commit(ctx); err != nil {
return fmt.Errorf("commit %s: %w", operation, err)
}
return nil
}
func getByTenantUIDForUpdate(ctx context.Context, tx pgx.Tx, instanceUID, tenantUID string) (Record, error) {
return scanRecord(tx.QueryRow(ctx, getByTenantUIDForUpdateStatement, instanceUID, tenantUID))
}
func getByTenantUID(ctx context.Context, tx pgx.Tx, instanceUID, tenantUID string) (Record, error) {
return scanRecord(tx.QueryRow(ctx, getByTenantUIDStatement, instanceUID, tenantUID))
}
type rowScanner interface {
Scan(...any) error
}
func scanRecord(row rowScanner) (Record, error) {
var record Record
err := row.Scan(
&record.InstanceUID,
&record.TenantUID,
&record.TenantNamespace,
&record.TenantName,
&record.DatabaseName,
&record.RoleName,
&record.CredentialPath,
&record.Managed,
&record.CreatedAt,
&record.UpdatedAt,
&record.RetainedAt,
)
if errors.Is(err, pgx.ErrNoRows) {
return Record{}, ErrNotFound
}
if err != nil {
return Record{}, fmt.Errorf("read registry record: %w", err)
}
return record, nil
}
func (o Ownership) validate() error {
if o.InstanceUID == "" || o.TenantUID == "" || o.TenantNamespace == "" || o.TenantName == "" ||
o.DatabaseName == "" || o.RoleName == "" || o.CredentialPath == "" {
return errors.New("claim registry ownership: all ownership fields are required")
}
return nil
}
func (o Ownership) equal(other Ownership) bool {
return o == other
}
@@ -0,0 +1,68 @@
/*
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 registry
const claimStatement = `
INSERT INTO postgresql_tenant_operator.tenant_ownership (
instance_uid,
tenant_uid,
tenant_namespace,
tenant_name,
database_name,
role_name,
credential_path
) VALUES ($1, $2, $3, $4, $5, $6, $7)
ON CONFLICT DO NOTHING
RETURNING
instance_uid,
tenant_uid,
tenant_namespace,
tenant_name,
database_name,
role_name,
credential_path,
managed,
created_at,
updated_at,
retained_at`
const getByTenantUIDStatement = `
SELECT
instance_uid,
tenant_uid,
tenant_namespace,
tenant_name,
database_name,
role_name,
credential_path,
managed,
created_at,
updated_at,
retained_at
FROM postgresql_tenant_operator.tenant_ownership
WHERE instance_uid = $1 AND tenant_uid = $2`
const getByTenantUIDForUpdateStatement = getByTenantUIDStatement + ` FOR UPDATE`
const markRetainedStatement = `
UPDATE postgresql_tenant_operator.tenant_ownership
SET managed = false, retained_at = clock_timestamp(), updated_at = clock_timestamp()
WHERE instance_uid = $1 AND tenant_uid = $2 AND managed = true`
const deleteStatement = `
DELETE FROM postgresql_tenant_operator.tenant_ownership
WHERE instance_uid = $1 AND tenant_uid = $2 AND managed = true`
@@ -0,0 +1,427 @@
//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 postgresql_test
import (
"context"
"errors"
"net"
"net/url"
"strconv"
"sync"
"testing"
"time"
"github.com/jackc/pgx/v5"
"github.com/jackc/pgx/v5/pgxpool"
"git.ddupan.top/panxiao81/ayatori/internal/database/adapter/postgresql/registry"
)
// registryFixture 只连接本测试创建的容器,不接受外部数据库地址。
type registryFixture struct {
ctx context.Context
pool *pgxpool.Pool
store *registry.Store
}
func newRegistryFixture(t *testing.T) *registryFixture {
t.Helper()
ctx, cancel := context.WithTimeout(context.Background(), 90*time.Second)
t.Cleanup(cancel)
_, port := postgresFixture(t, ctx)
endpoint := url.URL{
Scheme: "postgres",
User: url.UserPassword(fixtureUser, fixturePassword),
Host: net.JoinHostPort("127.0.0.1", strconv.Itoa(port)),
Path: "/postgres",
RawQuery: "sslmode=disable",
}
pool, err := pgxpool.New(ctx, endpoint.String())
if err != nil {
t.Fatal("cannot configure registry fixture connection")
}
t.Cleanup(pool.Close)
if err := pool.Ping(ctx); err != nil {
t.Fatal("cannot connect to registry fixture")
}
return &registryFixture{ctx: ctx, pool: pool, store: registry.NewStore(pool)}
}
func registryOwner(suffix string) registry.Ownership {
return registry.Ownership{
InstanceUID: "instance-uid",
TenantUID: "tenant-uid-" + suffix,
TenantNamespace: "applications",
TenantName: "tenant-" + suffix,
DatabaseName: "database_" + suffix,
RoleName: "role_" + suffix,
CredentialPath: "credentials/" + suffix,
}
}
func (f *registryFixture) bootstrap(t *testing.T) {
t.Helper()
if err := f.store.Bootstrap(f.ctx); err != nil {
t.Fatalf("bootstrap registry: %v", err)
}
}
func (f *registryFixture) claim(t *testing.T, owner registry.Ownership, expected registry.ClaimResult) {
t.Helper()
result, err := f.store.Claim(f.ctx, owner)
if err != nil || result != expected {
t.Fatalf("claim: result=%q, error=%v; want %q", result, err, expected)
}
}
func TestRegistryOwnershipLifecycle(t *testing.T) {
f := newRegistryFixture(t)
f.bootstrap(t)
f.bootstrap(t)
owner := registryOwner("lifecycle")
f.claim(t, owner, registry.ClaimCreated)
f.claim(t, owner, registry.ClaimOwned)
record, err := f.store.Get(f.ctx, owner.InstanceUID, owner.TenantUID)
if err != nil {
t.Fatal(err)
}
if record.Ownership != owner || !record.Managed || record.CreatedAt.IsZero() || record.RetainedAt != nil {
t.Fatalf("unexpected ownership record: %+v", record)
}
// 错误的资源归属不能删除或 Retain 原记录。
wrongOwner := owner
wrongOwner.RoleName = "another_role"
if err := f.store.Delete(f.ctx, wrongOwner); !errors.Is(err, registry.ErrConflict) {
t.Fatalf("delete mismatched owner: %v", err)
}
if err := f.store.MarkRetained(f.ctx, wrongOwner); !errors.Is(err, registry.ErrConflict) {
t.Fatalf("retain mismatched owner: %v", err)
}
for range 2 {
if err := f.store.Delete(f.ctx, owner); err != nil {
t.Fatalf("delete retry: %v", err)
}
}
if _, err := f.store.Get(f.ctx, owner.InstanceUID, owner.TenantUID); !errors.Is(err, registry.ErrNotFound) {
t.Fatalf("read deleted record: %v", err)
}
}
func TestRegistryRetainedRecordCannotBeReclaimed(t *testing.T) {
f := newRegistryFixture(t)
f.bootstrap(t)
owner := registryOwner("retained")
f.claim(t, owner, registry.ClaimCreated)
if err := f.store.MarkRetained(f.ctx, owner); err != nil {
t.Fatal(err)
}
retained, err := f.store.Get(f.ctx, owner.InstanceUID, owner.TenantUID)
if err != nil {
t.Fatal(err)
}
if retained.Managed || retained.RetainedAt == nil {
t.Fatalf("missing retained tombstone: %+v", retained)
}
if err := f.store.MarkRetained(f.ctx, owner); err != nil {
t.Fatalf("retain retry: %v", err)
}
retried, err := f.store.Get(f.ctx, owner.InstanceUID, owner.TenantUID)
if err != nil || !retried.UpdatedAt.Equal(retained.UpdatedAt) {
t.Fatalf("retain retry changed tombstone: %+v, %v", retried, err)
}
if err := f.store.Delete(f.ctx, owner); !errors.Is(err, registry.ErrConflict) {
t.Fatalf("delete retained record: %v", err)
}
if _, err := f.store.Claim(f.ctx, owner); !errors.Is(err, registry.ErrConflict) {
t.Fatalf("reclaim retained record: %v", err)
}
owner.TenantUID = "replacement-uid"
if _, err := f.store.Claim(f.ctx, owner); !errors.Is(err, registry.ErrConflict) {
t.Fatalf("replacement tenant reclaimed tombstone: %v", err)
}
}
func TestRegistryUniqueReservations(t *testing.T) {
f := newRegistryFixture(t)
f.bootstrap(t)
owner := registryOwner("original")
f.claim(t, owner, registry.ClaimCreated)
tests := []struct {
name string
change func(*registry.Ownership)
}{
{name: "tenant UID", change: func(other *registry.Ownership) { other.TenantUID = owner.TenantUID }},
{name: "tenant name", change: func(other *registry.Ownership) { other.TenantName = owner.TenantName }},
{name: "database", change: func(other *registry.Ownership) { other.DatabaseName = owner.DatabaseName }},
{name: "role", change: func(other *registry.Ownership) { other.RoleName = owner.RoleName }},
{name: "credential path", change: func(other *registry.Ownership) { other.CredentialPath = owner.CredentialPath }},
}
for _, tt := range tests {
t.Run(tt.name, func(t *testing.T) {
other := registryOwner("other")
tt.change(&other)
if _, err := f.store.Claim(f.ctx, other); !errors.Is(err, registry.ErrConflict) {
t.Fatalf("conflicting reservation: %v", err)
}
})
}
}
func TestRegistryConcurrentBootstrapAndClaims(t *testing.T) {
f := newRegistryFixture(t)
const workers = 8
start := make(chan struct{})
results := make(chan error, workers)
var group sync.WaitGroup
for range workers {
group.Go(func() {
<-start
results <- f.store.Bootstrap(f.ctx)
})
}
close(start)
group.Wait()
close(results)
for err := range results {
if err != nil {
t.Fatalf("concurrent bootstrap: %v", err)
}
}
owner := registryOwner("concurrent")
type claimOutcome struct {
result registry.ClaimResult
err error
}
claims := make(chan claimOutcome, workers)
start = make(chan struct{})
for range workers {
group.Go(func() {
<-start
result, err := f.store.Claim(f.ctx, owner)
claims <- claimOutcome{result: result, err: err}
})
}
close(start)
group.Wait()
close(claims)
created := 0
for outcome := range claims {
if outcome.err != nil {
t.Fatal(outcome.err)
}
switch outcome.result {
case registry.ClaimCreated:
created++
case registry.ClaimOwned:
default:
t.Fatalf("unexpected claim result: %q", outcome.result)
}
}
if created != 1 {
t.Fatalf("created %d records for the same owner; want 1", created)
}
}
func TestRegistryRestartAndCanceledRequestRecovery(t *testing.T) {
f := newRegistryFixture(t)
f.bootstrap(t)
owner := registryOwner("restart")
f.claim(t, owner, registry.ClaimCreated)
config := f.pool.Config()
f.pool.Close()
pool, err := pgxpool.NewWithConfig(f.ctx, config)
if err != nil {
t.Fatal("cannot reopen fixture connection")
}
t.Cleanup(pool.Close)
f.store = registry.NewStore(pool)
f.bootstrap(t)
// 模拟客户端丢失上次成功结果后,以新连接重试;归属证据必须来自数据库。
f.claim(t, owner, registry.ClaimOwned)
canceled, cancel := context.WithCancel(f.ctx)
cancel()
if _, err := f.store.Claim(canceled, registryOwner("canceled")); !errors.Is(err, context.Canceled) {
t.Fatalf("canceled claim: %v", err)
}
f.claim(t, registryOwner("canceled"), registry.ClaimCreated)
}
func TestRegistryBootstrapRejectsUnknownSchema(t *testing.T) {
f := newRegistryFixture(t)
// 人工创建的同名 schema 不能被初始化过程接管或覆盖。
if _, err := f.pool.Exec(f.ctx, "CREATE SCHEMA postgresql_tenant_operator"); err != nil {
t.Fatal(err)
}
if err := f.store.Bootstrap(f.ctx); err == nil {
t.Fatal("bootstrap accepted an unmanaged schema")
}
var version int
if err := f.pool.QueryRow(f.ctx, "SELECT version FROM public.postgresql_tenant_operator_schema_version").Scan(&version); err != nil {
t.Fatal(err)
}
if version != 0 {
t.Fatalf("failed migration advanced version to %d", version)
}
// 仅在本测试拥有的临时数据库中清除空冲突 schema,然后重试失败的迁移。
if _, err := f.pool.Exec(f.ctx, "DROP SCHEMA postgresql_tenant_operator"); err != nil {
t.Fatal(err)
}
f.bootstrap(t)
if _, err := f.pool.Exec(f.ctx, "UPDATE public.postgresql_tenant_operator_schema_version SET version = 100"); err != nil {
t.Fatal(err)
}
if err := f.store.Bootstrap(f.ctx); err == nil {
t.Fatal("bootstrap accepted a future migration version")
}
if err := f.pool.QueryRow(f.ctx, "SELECT version FROM public.postgresql_tenant_operator_schema_version").Scan(&version); err != nil {
t.Fatal(err)
}
if version != 100 {
t.Fatalf("bootstrap rewrote future version to %d", version)
}
}
// 数据库执行真实 COMMIT 后才注入客户端错误,模拟客户端无法确认提交结果。
// 这不是网络故障测试,但能确定性覆盖已提交、调用者却收到失败的恢复分支。
type lostCommitReply struct {
registry.Beginner
err error
}
func (db lostCommitReply) Begin(ctx context.Context) (pgx.Tx, error) {
tx, err := db.Beginner.Begin(ctx)
if err != nil {
return nil, err
}
return uncertainCommit{Tx: tx, err: db.err}, nil
}
type uncertainCommit struct {
pgx.Tx
err error
}
func (tx uncertainCommit) Commit(ctx context.Context) error {
if err := tx.Tx.Commit(ctx); err != nil {
return err
}
return tx.err
}
func TestRegistryRetriesAfterLostCommitReply(t *testing.T) {
f := newRegistryFixture(t)
f.bootstrap(t)
replyLost := errors.New("test-only lost commit reply")
uncertain := registry.NewStore(lostCommitReply{Beginner: f.pool, err: replyLost})
owner := registryOwner("uncertain")
if _, err := uncertain.Claim(f.ctx, owner); !errors.Is(err, replyLost) {
t.Fatalf("claim did not report lost reply: %v", err)
}
f.claim(t, owner, registry.ClaimOwned)
if err := uncertain.MarkRetained(f.ctx, owner); !errors.Is(err, replyLost) {
t.Fatalf("retain did not report lost reply: %v", err)
}
if err := f.store.MarkRetained(f.ctx, owner); err != nil {
t.Fatalf("retry uncertain retain: %v", err)
}
if err := f.store.Delete(f.ctx, owner); !errors.Is(err, registry.ErrConflict) {
t.Fatalf("uncertain retain lost tombstone protection: %v", err)
}
deletable := registryOwner("uncertain_delete")
f.claim(t, deletable, registry.ClaimCreated)
if err := uncertain.Delete(f.ctx, deletable); !errors.Is(err, replyLost) {
t.Fatalf("delete did not report lost reply: %v", err)
}
if err := f.store.Delete(f.ctx, deletable); err != nil {
t.Fatalf("retry uncertain delete: %v", err)
}
}
func TestRegistryConcurrentConflictingClaims(t *testing.T) {
f := newRegistryFixture(t)
f.bootstrap(t)
first := registryOwner("first")
second := registryOwner("second")
second.DatabaseName = first.DatabaseName
start := make(chan struct{})
results := make(chan error, 2)
var group sync.WaitGroup
for _, owner := range []registry.Ownership{first, second} {
group.Go(func() {
<-start
_, err := f.store.Claim(f.ctx, owner)
results <- err
})
}
close(start)
group.Wait()
close(results)
created, conflicts := 0, 0
for err := range results {
switch {
case err == nil:
created++
case errors.Is(err, registry.ErrConflict):
conflicts++
default:
t.Fatalf("unexpected concurrent claim error: %v", err)
}
}
if created != 1 || conflicts != 1 {
t.Fatalf("concurrent reservation: created=%d conflicts=%d; want one of each", created, conflicts)
}
}
func TestRegistryConcurrentRetainAndDelete(t *testing.T) {
f := newRegistryFixture(t)
f.bootstrap(t)
owner := registryOwner("delete_race")
f.claim(t, owner, registry.ClaimCreated)
start := make(chan struct{})
retainResult := make(chan error, 1)
deleteResult := make(chan error, 1)
go func() {
<-start
retainResult <- f.store.MarkRetained(f.ctx, owner)
}()
go func() {
<-start
deleteResult <- f.store.Delete(f.ctx, owner)
}()
close(start)
retainErr, deleteErr := <-retainResult, <-deleteResult
record, readErr := f.store.Get(f.ctx, owner.InstanceUID, owner.TenantUID)
switch {
case retainErr == nil:
// Retain 先取得行锁时,Delete 必须拒绝删除墓碑。
if !errors.Is(deleteErr, registry.ErrConflict) || readErr != nil || record.Managed || record.RetainedAt == nil {
t.Fatalf("retain won but tombstone was not protected: delete=%v read=%v record=%+v", deleteErr, readErr, record)
}
case errors.Is(retainErr, registry.ErrNotFound):
// Delete 先提交时,Retain 必须报告记录已不存在,不能重建墓碑。
if deleteErr != nil || !errors.Is(readErr, registry.ErrNotFound) {
t.Fatalf("delete won but record remains: delete=%v read=%v", deleteErr, readErr)
}
default:
t.Fatalf("unexpected retain/delete race: retain=%v delete=%v", retainErr, deleteErr)
}
}
@@ -0,0 +1,140 @@
//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 postgresql_test
import (
"context"
"crypto/ecdsa"
"crypto/elliptic"
"crypto/rand"
"crypto/x509"
"crypto/x509/pkix"
"encoding/pem"
"math/big"
"net"
"os"
"os/exec"
"path/filepath"
"testing"
"time"
"git.ddupan.top/panxiao81/ayatori/internal/database/adapter/postgresql"
"git.ddupan.top/panxiao81/ayatori/internal/database/application"
"git.ddupan.top/panxiao81/ayatori/internal/database/domain/instance"
)
func fixtureCertificate(t *testing.T) (string, string) {
t.Helper()
key, err := ecdsa.GenerateKey(elliptic.P256(), rand.Reader)
if err != nil {
t.Fatal(err)
}
cert := &x509.Certificate{
SerialNumber: big.NewInt(1),
Subject: pkix.Name{CommonName: fixtureHost},
NotBefore: time.Now().Add(-time.Hour),
NotAfter: time.Now().Add(time.Hour),
DNSNames: []string{fixtureHost},
IPAddresses: []net.IP{net.ParseIP("127.0.0.1")},
IsCA: true,
BasicConstraintsValid: true,
KeyUsage: x509.KeyUsageCertSign | x509.KeyUsageDigitalSignature,
ExtKeyUsage: []x509.ExtKeyUsage{x509.ExtKeyUsageServerAuth},
}
der, err := x509.CreateCertificate(rand.Reader, cert, cert, &key.PublicKey, key)
if err != nil {
t.Fatal(err)
}
encodedKey, err := x509.MarshalECPrivateKey(key)
if err != nil {
t.Fatal(err)
}
dir := t.TempDir()
certPath := filepath.Join(dir, "server.crt")
keyPath := filepath.Join(dir, "server.key")
if err := os.WriteFile(certPath, pem.EncodeToMemory(&pem.Block{Type: "CERTIFICATE", Bytes: der}), 0600); err != nil {
t.Fatal(err)
}
if err := os.WriteFile(keyPath, pem.EncodeToMemory(&pem.Block{Type: "EC PRIVATE KEY", Bytes: encodedKey}), 0600); err != nil {
t.Fatal(err)
}
return certPath, keyPath
}
func TestPostgreSQLTLSHostIdentity(t *testing.T) {
const psql = "psql"
ctx, cancel := context.WithTimeout(context.Background(), time.Minute)
defer cancel()
id, port := postgresFixture(t, ctx)
certPath, keyPath := fixtureCertificate(t)
commands := [][]string{
{"cp", certPath, id + ":/tmp/server.crt"},
{"cp", keyPath, id + ":/tmp/server.key"},
{dockerExec, "-u", "0", id, "chown", "postgres:postgres", "/tmp/server.crt", "/tmp/server.key"},
{dockerExec, id, psql, "-U", fixtureUser, "-c", "ALTER SYSTEM SET ssl_cert_file='/tmp/server.crt'"},
{dockerExec, id, psql, "-U", fixtureUser, "-c", "ALTER SYSTEM SET ssl_key_file='/tmp/server.key'"},
{dockerExec, id, psql, "-U", fixtureUser, "-c", "ALTER SYSTEM SET ssl=on"},
{dockerExec, id, psql, "-U", fixtureUser, "-c", "SELECT pg_reload_conf()"},
}
for _, args := range commands {
if exec.CommandContext(ctx, "docker", args...).Run() != nil {
t.Fatal("TLS fixture setup failed")
}
}
credentials, err := application.NewCredentials(fixtureUser, fixturePassword)
if err != nil {
t.Fatal(err)
}
connector := postgresql.Connector{RootCert: certPath}
endpoint := target(t, port, instance.TLSVerifyFull).Definition().Endpoint()
db, err := connector.Connect(ctx, endpoint, credentials)
if err != nil {
t.Fatal("trusted DNS SAN connection failed", err)
}
if version, err := db.Version(ctx); err != nil || version == "" {
db.Close()
t.Fatal("TLS metadata read failed", err)
}
db.Close()
values := endpoint.Values()
values.Host = "127.0.0.1"
ipEndpoint, err := instance.NewEndpoint(values)
if err != nil {
t.Fatal(err)
}
db, err = connector.Connect(ctx, ipEndpoint, credentials)
if err != nil {
t.Fatal("trusted IP SAN connection failed", err)
}
db.Close()
values.Host = "wrong.invalid"
wrongEndpoint, err := instance.NewEndpoint(values)
if err != nil {
t.Fatal(err)
}
if db, err := connector.Connect(ctx, wrongEndpoint, credentials); err == nil {
db.Close()
t.Fatal("wrong TLS hostname accepted")
}
otherCA, _ := fixtureCertificate(t)
if db, err := (postgresql.Connector{RootCert: otherCA}).Connect(ctx, endpoint, credentials); err == nil {
db.Close()
t.Fatal("wrong CA accepted")
}
}
@@ -0,0 +1,58 @@
/*
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 定义 Database 用例与适配器之间的边界。
package application
import (
"context"
"errors"
"git.ddupan.top/panxiao81/ayatori/internal/database/domain/instance"
)
var (
ErrCredentialsUnavailable = errors.New("management credentials unavailable")
ErrCredentialsInvalid = errors.New("management credentials invalid")
)
// Credentials 只存在于应用与连接适配器内存,不进入领域对象或持久化状态。
type Credentials struct {
username string
password string
}
func NewCredentials(username, password string) (Credentials, error) {
if username == "" || password == "" {
return Credentials{}, ErrCredentialsInvalid
}
return Credentials{username: username, password: password}, nil
}
func (c Credentials) Username() string { return c.username }
func (c Credentials) Password() string { return c.password }
func (c Credentials) String() string { return "[redacted management credentials]" }
func (c Credentials) GoString() string { return c.String() }
// MarshalJSON 显式隐藏内容,避免未来字段调整意外改变日志或序列化行为。
func (c Credentials) MarshalJSON() ([]byte, error) {
return []byte(`"[redacted management credentials]"`), nil
}
// CredentialReader 返回本次读取的有效值;metadata 不参与凭据相等比较。
type CredentialReader interface {
Read(context.Context, instance.CredentialReference) (Credentials, error)
}
@@ -0,0 +1,70 @@
/*
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 (
"encoding/json"
"fmt"
"strings"
"testing"
)
const testUsername = "test-user"
func TestCredentialsRejectEmptyValues(t *testing.T) {
for _, values := range [][2]string{
{"", "test-password"},
{testUsername, ""},
{"", ""},
} {
if _, err := NewCredentials(values[0], values[1]); err != ErrCredentialsInvalid {
t.Fatal("empty credential was accepted")
}
}
}
func TestCredentialAndServiceFormattingIsRedacted(t *testing.T) {
const canary = "SECRET-CANARY-never-log-this"
credentials, err := NewCredentials(canary, canary)
if err != nil {
t.Fatal(err)
}
if credentials.Username() != canary || credentials.Password() != canary {
t.Fatal("explicit credential access changed values")
}
service, err := NewInstanceService(&sourceStub{credentials: credentials}, &connectorStub{})
if err != nil {
t.Fatal(err)
}
t.Cleanup(service.Close)
encoded, err := json.Marshal(credentials)
if err != nil {
t.Fatal(err)
}
outputs := []string{
string(encoded),
fmt.Sprintf("%v %+v %#v", credentials, credentials, credentials),
fmt.Sprintf("%v %+v %#v", service, service, service),
}
for _, output := range outputs {
if strings.Contains(output, canary) {
t.Fatal("formatting leaked credential data")
}
}
}
@@ -0,0 +1,171 @@
/*
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 (
"context"
"errors"
"sync"
"git.ddupan.top/panxiao81/ayatori/internal/database/domain/instance"
)
var (
ErrConnection = errors.New("management connection unavailable")
ErrAuthentication = errors.New("management authentication failed")
ErrObservation = errors.New("management observation failed")
ErrCredentialsChanged = errors.New("management credentials changed during observation")
ErrClosed = errors.New("instance service closed")
)
// Database 与 Connector 复用原项目 internal/instance/service.go 的能力边界。
// 版本查询只是本切片的连通性观察,不能产生领域 Ready。
type Database interface {
Version(context.Context) (string, error)
Close()
}
type Connector interface {
Connect(context.Context, instance.Endpoint, Credentials) (Database, error)
}
type entry struct {
target instance.ObservationTarget
credentials Credentials
database Database
}
// InstanceService 由原 Service 迁移:连接复用与释放属于应用装配,不属于 SQL adapter。
// 保留原实现串行操作的约束,防止 Close 与查询并发;controller 停止 worker 后调用 Close。
// 不缓存能力观察,不把连接存活等同于 Ready。凭据每轮重新读取,而非只在引用变化时读取。
type InstanceService struct {
mu sync.Mutex
source CredentialReader
connector Connector
entries map[string]*entry
closed bool
}
func NewInstanceService(source CredentialReader, connector Connector) (*InstanceService, error) {
if source == nil || connector == nil {
return nil, errors.New("credential source and connector required")
}
return &InstanceService{
source: source,
connector: connector,
entries: make(map[string]*entry),
}, nil
}
func (s *InstanceService) String() string { return "[redacted instance service]" }
func (s *InstanceService) GoString() string { return s.String() }
// ObserveVersion 返回当前目标和凭据下的版本;任何失败均返回空结果。
// 调用者仍需使用 CR resourceVersion 保存前提防止 spec 并发修改;本方法不建立跨系统事务。
func (s *InstanceService) ObserveVersion(ctx context.Context, target instance.ObservationTarget) (string, error) {
if err := target.Validate(); err != nil {
return "", err
}
s.mu.Lock()
defer s.mu.Unlock()
if s.closed {
return "", ErrClosed
}
if err := ctx.Err(); err != nil {
return "", err
}
// 先读取有效凭据。读取失败时不得继续使用缓存中的旧连接。
name := target.Identity().Name()
credentials, err := s.source.Read(ctx, target.Definition().AdminCredential())
if err != nil {
s.release(name)
return "", credentialError(err)
}
if credentials.username == "" || credentials.password == "" {
s.release(name)
return "", ErrCredentialsInvalid
}
// 连接身份与有效值均未变化时复用 pgxpool;generation 本身不要求换池。
current := s.entries[name]
if current != nil && (current.target.Identity() != target.Identity() ||
current.target.Definition() != target.Definition() || current.credentials != credentials) {
s.release(name)
current = nil
}
if current == nil {
database, err := s.connector.Connect(ctx, target.Definition().Endpoint(), credentials)
if err != nil {
return "", err
}
current = &entry{
target: target,
credentials: credentials,
database: database,
}
s.entries[name] = current
}
version, err := current.database.Version(ctx)
if err != nil {
s.release(name)
return "", err
}
// 回读后再检查凭据,避免把轮换前取得的结果交给新凭据的调用链。
latest, err := s.source.Read(ctx, target.Definition().AdminCredential())
if err != nil {
s.release(name)
return "", credentialError(err)
}
if latest != credentials {
s.release(name)
return "", ErrCredentialsChanged
}
return version, nil
}
func credentialError(err error) error {
if errors.Is(err, ErrCredentialsInvalid) {
return ErrCredentialsInvalid
}
return ErrCredentialsUnavailable
}
// Forget 只释放本地连接;不删除数据库或 registry,不替代 Instance finalizer。
func (s *InstanceService) Forget(name string) {
s.mu.Lock()
defer s.mu.Unlock()
s.release(name)
}
func (s *InstanceService) release(name string) {
if current := s.entries[name]; current != nil {
current.database.Close()
}
delete(s.entries, name)
}
func (s *InstanceService) Close() {
s.mu.Lock()
defer s.mu.Unlock()
s.closed = true
for name := range s.entries {
s.release(name)
}
}
@@ -0,0 +1,182 @@
/*
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 (
"context"
"errors"
"testing"
"git.ddupan.top/panxiao81/ayatori/internal/database/domain/instance"
)
// 延续源项目 Service 测试,用于穷举身份与装配失败;真实行为由 adapter 集成测试验证。
type sourceStub struct {
credentials Credentials
err error
}
func (s *sourceStub) Read(context.Context, instance.CredentialReference) (Credentials, error) {
return s.credentials, s.err
}
type databaseStub struct {
closes int
err error
}
func (d *databaseStub) Version(context.Context) (string, error) { return "17", d.err }
func (d *databaseStub) Close() {
d.closes++
}
type connectorStub struct {
databases []*databaseStub
err error
}
func (c *connectorStub) Connect(context.Context, instance.Endpoint, Credentials) (Database, error) {
if c.err != nil {
return nil, c.err
}
db := &databaseStub{}
c.databases = append(c.databases, db)
return db, nil
}
func serviceTarget(t *testing.T, uid, host, secret string, generation int64) instance.ObservationTarget {
t.Helper()
id, err := instance.NewIdentity(uid, "shared")
if err != nil {
t.Fatal(err)
}
revision, err := instance.NewRevision(generation)
if err != nil {
t.Fatal(err)
}
endpoint, err := instance.NewEndpoint(instance.EndpointValues{
Host: host,
HostAddr: "127.0.0.1",
Port: 5432,
ManagementDatabase: "postgres",
TLSMode: instance.TLSDisable,
})
if err != nil {
t.Fatal(err)
}
ref, err := instance.NewCredentialReference(instance.CredentialReferenceValues{
Name: secret,
UsernameKey: "user",
PasswordKey: "pass",
})
if err != nil {
t.Fatal(err)
}
definition, err := instance.NewDefinition(endpoint, ref)
if err != nil {
t.Fatal(err)
}
target, err := instance.NewObservationTarget(id, revision, definition)
if err != nil {
t.Fatal(err)
}
return target
}
func TestInstanceConnectionIdentity(t *testing.T) {
source := &sourceStub{credentials: Credentials{username: testUsername, password: "test-only"}}
connector := &connectorStub{}
service, err := NewInstanceService(source, connector)
if err != nil {
t.Fatal(err)
}
defer service.Close()
ctx := context.Background()
cases := []struct {
name string
target instance.ObservationTarget
wantConnections int
}{
{"initial connection", serviceTarget(t, "uid-1", "first", "admin", 1), 1},
{"generation alone", serviceTarget(t, "uid-1", "first", "admin", 2), 1},
{"endpoint changed", serviceTarget(t, "uid-1", "second", "admin", 3), 2},
{"reference changed", serviceTarget(t, "uid-1", "second", "replacement", 4), 3},
{"same name with new UID", serviceTarget(t, "uid-2", "second", "replacement", 1), 4},
}
for _, testCase := range cases {
if _, err := service.ObserveVersion(ctx, testCase.target); err != nil {
t.Fatal(err)
}
if len(connector.databases) != testCase.wantConnections {
t.Fatalf("%s: got %d connections, want %d", testCase.name, len(connector.databases), testCase.wantConnections)
}
}
for _, db := range connector.databases[:3] {
if db.closes != 1 {
t.Fatal("replaced connection not closed exactly once")
}
}
service.Forget("shared")
service.Forget("shared")
if connector.databases[3].closes != 1 {
t.Fatal("forget did not close exactly once")
}
}
func TestInstanceAssemblyFailureRecovery(t *testing.T) {
ctx := context.Background()
target := serviceTarget(t, "uid-1", "first", "admin", 1)
source := &sourceStub{
credentials: Credentials{username: testUsername, password: "test-only"},
err: errors.New("unsafe source error"),
}
connector := &connectorStub{err: ErrConnection}
service, err := NewInstanceService(source, connector)
if err != nil {
t.Fatal(err)
}
defer service.Close()
if version, err := service.ObserveVersion(ctx, target); version != "" || !errors.Is(err, ErrCredentialsUnavailable) {
t.Fatal("unsafe source error escaped")
}
source.err = nil
if version, err := service.ObserveVersion(ctx, target); version != "" || !errors.Is(err, ErrConnection) {
t.Fatal("connection failure returned evidence")
}
connector.err = nil
if _, err := service.ObserveVersion(ctx, target); err != nil {
t.Fatal(err)
}
connector.databases[0].err = ErrObservation
if version, err := service.ObserveVersion(ctx, target); version != "" || !errors.Is(err, ErrObservation) {
t.Fatal("failed query returned evidence")
}
if connector.databases[0].closes != 1 {
t.Fatal("failed connection retained")
}
if _, err := service.ObserveVersion(ctx, target); err != nil {
t.Fatal("retry failed", err)
}
service.Close()
service.Close()
if connector.databases[1].closes != 1 {
t.Fatal("shutdown did not close once")
}
if _, err := service.ObserveVersion(ctx, target); !errors.Is(err, ErrClosed) {
t.Fatal("closed service accepted work")
}
}
@@ -38,22 +38,22 @@ const (
) )
// Snapshot contains persisted observations only, without credentials or live evidence. // Snapshot contains persisted observations only, without credentials or live evidence.
// Failure detail mapping will be added with capability assessment, not intent transitions.
type Snapshot struct { type Snapshot struct {
Phase Phase Phase Phase
ObservedRevision int64 ObservedRevision int64
Readiness Readiness Readiness Readiness
ReportedVersion string ReportedVersion string
Failure Failure
} }
// Instance protects registration state and pure lifecycle transitions. // Instance protects registration state and pure lifecycle transitions.
// Reconstitution does not establish live capability evidence, even for a Ready snapshot. // Reconstitution does not establish live capability evidence, even for a Ready snapshot.
// This initial slice deliberately exposes no operation that authorizes provisioning.
type Instance struct { type Instance struct {
target ObservationTarget target ObservationTarget
snapshot Snapshot snapshot Snapshot
deleting bool deleting bool
extensions ExtensionSupport extensions ExtensionSupport
evidence *CapabilityObservation
} }
func Reconstitute(target ObservationTarget, snapshot Snapshot, deleting bool) (*Instance, error) { func Reconstitute(target ObservationTarget, snapshot Snapshot, deleting bool) (*Instance, error) {
@@ -85,6 +85,8 @@ func (i *Instance) BeginValidation() error {
i.snapshot.Phase = PhaseValidating i.snapshot.Phase = PhaseValidating
i.snapshot.Readiness = Unknown i.snapshot.Readiness = Unknown
i.extensions = ExtensionSupport{} i.extensions = ExtensionSupport{}
i.evidence = nil
i.snapshot.Failure = NoFailure
return nil return nil
} }
@@ -100,6 +102,8 @@ func (i *Instance) BeginDeletion() error {
i.snapshot.Phase = PhaseDeleting i.snapshot.Phase = PhaseDeleting
i.snapshot.Readiness = Unknown i.snapshot.Readiness = Unknown
i.extensions = ExtensionSupport{} i.extensions = ExtensionSupport{}
i.evidence = nil
i.snapshot.Failure = NoFailure
return nil return nil
} }
@@ -0,0 +1,263 @@
/*
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 instance
import "errors"
// Failure 只表示安全类别;驱动错误、凭据和 Condition 文案留在应用边界。
type Failure uint8
const (
NoFailure Failure = iota
ObservationIncomplete
DependencyUnavailable
AuthenticationFailed
InsufficientPrivileges
RegistryIncompatible
RegistryNotUsable
)
// CheckResult 的零值表示未观察,不能视为成功。
type CheckResult uint8
const (
CheckUnobserved CheckResult = iota
CheckPassed
CheckUnavailable
CheckAuthenticationFailed
CheckInsufficientPrivileges
)
// ManagementChecks 分别记录所需能力;SQL 探测和同轮次关联由 adapter/application 保证。
// Extensions 不代表任意扩展均可安装;具体请求仍需支持检查、执行及回读。
type ManagementChecks struct {
Connection CheckResult
Metadata CheckResult
Roles CheckResult
Databases CheckResult
Grants CheckResult
Extensions CheckResult
}
func (c ManagementChecks) failure() Failure {
for _, check := range []CheckResult{c.Connection, c.Metadata, c.Roles, c.Databases, c.Grants, c.Extensions} {
switch check {
case CheckPassed:
case CheckUnavailable:
return DependencyUnavailable
case CheckAuthenticationFailed:
return AuthenticationFailed
case CheckInsufficientPrivileges:
return InsufficientPrivileges
default:
return ObservationIncomplete
}
}
return NoFailure
}
type RegistryState uint8
const (
RegistryUnobserved RegistryState = iota
RegistryAbsent
RegistryNeedsMigration
RegistryUsable
RegistryUnsupported
RegistryUnavailable
)
// CapabilityObservation 是值对象,不包含连接、凭据或可变集合。
type CapabilityObservation struct {
target ObservationTarget
version string
checks ManagementChecks
registry RegistryState
}
func NewCapabilityObservation(target ObservationTarget, version string,
checks ManagementChecks, registry RegistryState,
) (CapabilityObservation, error) {
if err := target.Validate(); err != nil {
return CapabilityObservation{}, err
}
return CapabilityObservation{target: target, version: version, checks: checks, registry: registry}, nil
}
func (o CapabilityObservation) managementFailure() Failure {
if failure := o.checks.failure(); failure != NoFailure {
return failure
}
if o.version == "" {
return ObservationIncomplete
}
return NoFailure
}
func (o CapabilityObservation) registryFailure() Failure {
switch o.registry {
case RegistryUsable:
return NoFailure
case RegistryAbsent, RegistryNeedsMigration:
return RegistryNotUsable
case RegistryUnsupported:
return RegistryIncompatible
case RegistryUnavailable:
return DependencyUnavailable
default:
return ObservationIncomplete
}
}
type PreparationDecision uint8
const (
PreparationDenied PreparationDecision = iota
PreparationAllowed
AlreadyUsable
)
func (i *Instance) acceptObservation(o CapabilityObservation, phase Phase) error {
if !i.target.Matches(o.target) {
return errors.New("capability observation target does not match instance")
}
if i.deleting || i.snapshot.Phase != phase {
return errors.New("capability observation is not allowed in current lifecycle")
}
return nil
}
func (i *Instance) fail(failure Failure) {
i.evidence = nil
i.extensions = ExtensionSupport{}
i.snapshot.Readiness = NotReady
i.snapshot.Failure = failure
i.snapshot.ObservedRevision = i.target.Revision().Value()
}
// AssessManagement 只推进意图,不执行 registry 写入,也不完成 observedRevision。
func (i *Instance) AssessManagement(o CapabilityObservation) error {
if err := i.acceptObservation(o, PhaseValidating); err != nil {
return err
}
if failure := o.managementFailure(); failure != NoFailure {
i.fail(failure)
return nil
}
if failure := o.registryFailure(); failure != NoFailure && failure != RegistryNotUsable {
i.fail(failure)
return nil
}
i.snapshot.Phase = PhaseInitializingRegistry
i.snapshot.Readiness = Unknown
i.snapshot.Failure = NoFailure
i.evidence = nil
return nil
}
// PlanRegistryPreparation 不证明 checkpoint 已落盘;应用层必须先保存意图再执行写入。
func (i *Instance) PlanRegistryPreparation(o CapabilityObservation) (PreparationDecision, error) {
if err := i.acceptObservation(o, PhaseInitializingRegistry); err != nil {
return PreparationDenied, err
}
if failure := o.managementFailure(); failure != NoFailure {
i.fail(failure)
return PreparationDenied, nil
}
switch o.registry {
case RegistryUsable:
return AlreadyUsable, nil
case RegistryAbsent, RegistryNeedsMigration:
return PreparationAllowed, nil
default:
i.fail(o.registryFailure())
return PreparationDenied, nil
}
}
// RegistryPreparationResult 只能是安全失败或完整回读,不能表达裸操作成功。
type RegistryPreparationResult struct {
observation CapabilityObservation
failure Failure
}
func RegistryReadBack(o CapabilityObservation) RegistryPreparationResult {
return RegistryPreparationResult{observation: o}
}
func RegistryPreparationFailed(target ObservationTarget, failure Failure) (RegistryPreparationResult, error) {
if err := target.Validate(); err != nil {
return RegistryPreparationResult{}, err
}
if failure < ObservationIncomplete || failure > RegistryNotUsable {
return RegistryPreparationResult{}, errors.New("registry preparation requires a known failure category")
}
return RegistryPreparationResult{observation: CapabilityObservation{target: target}, failure: failure}, nil
}
func (i *Instance) AssessRegistryResult(result RegistryPreparationResult) error {
o := result.observation
if err := i.acceptObservation(o, PhaseInitializingRegistry); err != nil {
return err
}
if result.failure != NoFailure {
i.fail(result.failure)
return nil
}
i.assessComplete(o)
return nil
}
func (i *Instance) assessComplete(o CapabilityObservation) {
if failure := o.managementFailure(); failure != NoFailure {
i.fail(failure)
return
}
if failure := o.registryFailure(); failure != NoFailure {
i.fail(failure)
return
}
i.snapshot = Snapshot{Phase: PhaseReady, ObservedRevision: i.target.Revision().Value(),
Readiness: Ready, ReportedVersion: o.version}
i.evidence = &o
}
// AssessReadiness 每轮接收完整事实,失败立即撤销本轮供应能力。
func (i *Instance) AssessReadiness(o CapabilityObservation) error {
if err := i.acceptObservation(o, PhaseReady); err != nil {
return err
}
if i.snapshot.ObservedRevision != i.target.Revision().Value() {
return i.BeginValidation()
}
if o.managementFailure() != NoFailure || o.registry == RegistryUnavailable {
i.snapshot.Phase = PhaseValidating
} else if o.registryFailure() != NoFailure {
i.snapshot.Phase = PhaseInitializingRegistry
}
i.assessComplete(o)
return nil
}
// RequireProvisioningReady 仅检查 Instance 前置条件,不授予 Tenant 所有权或外部写入许可。
func (i *Instance) RequireProvisioningReady() error {
if i.deleting || i.evidence == nil || i.snapshot.Phase != PhaseReady || i.snapshot.Readiness != Ready ||
i.snapshot.ObservedRevision != i.target.Revision().Value() {
return errors.New("instance is not ready for provisioning")
}
return nil
}
@@ -0,0 +1,352 @@
/*
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 instance_test
import (
"testing"
"git.ddupan.top/panxiao81/ayatori/internal/database/domain/instance"
)
const testServerVersion = "17.6"
func completeChecks() instance.ManagementChecks {
return instance.ManagementChecks{
Connection: instance.CheckPassed, Metadata: instance.CheckPassed,
Roles: instance.CheckPassed, Databases: instance.CheckPassed,
Grants: instance.CheckPassed, Extensions: instance.CheckPassed,
}
}
func capability(t *testing.T, value *instance.Instance, checks instance.ManagementChecks,
registry instance.RegistryState,
) instance.CapabilityObservation {
t.Helper()
o, err := instance.NewCapabilityObservation(value.Target(), testServerVersion, checks, registry)
if err != nil {
t.Fatal(err)
}
return o
}
func readyInstance(t *testing.T) *instance.Instance {
t.Helper()
i := lifecycleInstance(t, instance.Snapshot{Phase: instance.PhaseInitializingRegistry}, false)
if err := i.AssessRegistryResult(instance.RegistryReadBack(capability(t, i, completeChecks(), instance.RegistryUsable))); err != nil {
t.Fatal(err)
}
if err := i.RequireProvisioningReady(); err != nil {
t.Fatal(err)
}
return i
}
func TestReadinessRequiresCompleteReadBack(t *testing.T) {
i := lifecycleInstance(t, instance.Snapshot{}, false)
if err := i.BeginValidation(); err != nil {
t.Fatal(err)
}
absent := capability(t, i, completeChecks(), instance.RegistryAbsent)
if err := i.AssessManagement(absent); err != nil {
t.Fatal(err)
}
if s := i.Snapshot(); s.Phase != instance.PhaseInitializingRegistry || s.ObservedRevision != 0 || s.Readiness != instance.Unknown {
t.Fatalf("management observation prematurely concluded readiness: %+v", s)
}
for range 2 {
decision, err := i.PlanRegistryPreparation(absent)
if err != nil || decision != instance.PreparationAllowed {
t.Fatalf("preparation: %v, %v", decision, err)
}
if i.RequireProvisioningReady() == nil {
t.Fatal("preparation authorized provisioning")
}
}
if err := i.AssessRegistryResult(instance.RegistryReadBack(absent)); err != nil {
t.Fatal(err)
}
if i.Snapshot().Failure != instance.RegistryNotUsable || i.RequireProvisioningReady() == nil {
t.Fatal("absent registry accepted as ready")
}
usable := capability(t, i, completeChecks(), instance.RegistryUsable)
decision, err := i.PlanRegistryPreparation(usable)
if err != nil || decision != instance.AlreadyUsable {
t.Fatalf("retry after external preparation: %v, %v", decision, err)
}
if err := i.AssessRegistryResult(instance.RegistryReadBack(usable)); err != nil {
t.Fatal(err)
}
if s := i.Snapshot(); s.Readiness != instance.Ready || s.ReportedVersion != testServerVersion || s.ObservedRevision != i.Target().Revision().Value() {
t.Fatalf("complete observation not accepted: %+v", s)
}
if err := i.RequireProvisioningReady(); err != nil {
t.Fatal(err)
}
}
// 重启只恢复 checkpoint;依赖稍后恢复时必须重新取得完整事实。
func TestReadinessRecoveryAndInvalidation(t *testing.T) {
i := readyInstance(t)
restored, err := instance.Reconstitute(i.Target(), i.Snapshot(), false)
if err != nil {
t.Fatal(err)
}
if restored.RequireProvisioningReady() == nil {
t.Fatal("persisted Ready fabricated fresh evidence")
}
for range 2 {
if err := restored.AssessReadiness(capability(t, restored, completeChecks(), instance.RegistryUsable)); err != nil {
t.Fatal(err)
}
if err := restored.RequireProvisioningReady(); err != nil {
t.Fatal(err)
}
}
if err := restored.BeginValidation(); err != nil {
t.Fatal(err)
}
if restored.RequireProvisioningReady() == nil {
t.Fatal("validation retained evidence")
}
deleted, err := instance.Reconstitute(i.Target(), i.Snapshot(), true)
if err != nil {
t.Fatal(err)
}
if deleted.RequireProvisioningReady() == nil {
t.Fatal("deletion allowed provisioning")
}
if err := deleted.BeginDeletion(); err != nil {
t.Fatal(err)
}
if deleted.RequireProvisioningReady() == nil {
t.Fatal("deleting checkpoint allowed provisioning")
}
}
func TestEachManagementCheckIsRequired(t *testing.T) {
for field := range 6 {
for _, result := range []instance.CheckResult{instance.CheckUnobserved, instance.CheckUnavailable,
instance.CheckAuthenticationFailed, instance.CheckInsufficientPrivileges, 255} {
checks := completeChecks()
fields := []*instance.CheckResult{&checks.Connection, &checks.Metadata, &checks.Roles,
&checks.Databases, &checks.Grants, &checks.Extensions}
*fields[field] = result
i := readyInstance(t)
if err := i.AssessReadiness(capability(t, i, checks, instance.RegistryUsable)); err != nil {
t.Fatal(err)
}
if s := i.Snapshot(); s.Phase != instance.PhaseValidating || s.Readiness != instance.NotReady ||
s.Failure == instance.NoFailure || i.RequireProvisioningReady() == nil {
t.Fatalf("check %d result %d accepted: %+v", field, result, s)
}
}
}
}
func TestRegistryDecisionsAndReadinessLoss(t *testing.T) {
for _, tc := range []struct {
state instance.RegistryState
decision instance.PreparationDecision
failure instance.Failure
}{
{instance.RegistryUsable, instance.AlreadyUsable, instance.NoFailure},
{instance.RegistryAbsent, instance.PreparationAllowed, instance.RegistryNotUsable},
{instance.RegistryNeedsMigration, instance.PreparationAllowed, instance.RegistryNotUsable},
{instance.RegistryUnsupported, instance.PreparationDenied, instance.RegistryIncompatible},
{instance.RegistryUnavailable, instance.PreparationDenied, instance.DependencyUnavailable},
{instance.RegistryUnobserved, instance.PreparationDenied, instance.ObservationIncomplete},
{255, instance.PreparationDenied, instance.ObservationIncomplete},
} {
i := lifecycleInstance(t, instance.Snapshot{Phase: instance.PhaseInitializingRegistry}, false)
o := capability(t, i, completeChecks(), tc.state)
decision, err := i.PlanRegistryPreparation(o)
if err != nil || decision != tc.decision {
t.Fatalf("registry %d: %v, %v", tc.state, decision, err)
}
i = readyInstance(t)
if err := i.AssessReadiness(o); err != nil {
t.Fatal(err)
}
if i.Snapshot().Failure != tc.failure {
t.Fatalf("registry %d: %+v", tc.state, i.Snapshot())
}
if tc.state != instance.RegistryUsable {
wantPhase := instance.PhaseInitializingRegistry
if tc.state == instance.RegistryUnavailable {
wantPhase = instance.PhaseValidating
}
if i.Snapshot().Phase != wantPhase || i.RequireProvisioningReady() == nil {
t.Fatal("registry drift retained readiness")
}
}
}
}
func TestInitializationRejectsIncompleteOrFailedManagement(t *testing.T) {
for _, tc := range []struct {
checks instance.ManagementChecks
registry instance.RegistryState
failure instance.Failure
}{
{instance.ManagementChecks{}, instance.RegistryUsable, instance.ObservationIncomplete},
{completeChecks(), instance.RegistryUnsupported, instance.RegistryIncompatible},
{completeChecks(), instance.RegistryUnavailable, instance.DependencyUnavailable},
} {
i := lifecycleInstance(t, instance.Snapshot{Phase: instance.PhaseValidating}, false)
if err := i.AssessManagement(capability(t, i, tc.checks, tc.registry)); err != nil {
t.Fatal(err)
}
if s := i.Snapshot(); s.Phase != instance.PhaseValidating || s.Failure != tc.failure ||
s.ObservedRevision != i.Target().Revision().Value() || s.Readiness != instance.NotReady {
t.Fatalf("invalid management accepted: %+v", s)
}
}
i := lifecycleInstance(t, instance.Snapshot{Phase: instance.PhaseInitializingRegistry}, false)
decision, err := i.PlanRegistryPreparation(capability(t, i, instance.ManagementChecks{}, instance.RegistryAbsent))
if err != nil || decision != instance.PreparationDenied || i.Snapshot().Failure != instance.ObservationIncomplete {
t.Fatal("incomplete management allowed registry writes")
}
if err := i.AssessRegistryResult(instance.RegistryReadBack(capability(t, i, completeChecks(), instance.RegistryUsable))); err != nil {
t.Fatal(err)
}
if err := i.RequireProvisioningReady(); err != nil {
t.Fatal("dependency recovery did not restore readiness", err)
}
}
func TestReadinessMethodsRejectWrongPhaseAndDeletion(t *testing.T) {
for _, deleting := range []bool{false, true} {
for _, phase := range []instance.Phase{instance.PhasePending, instance.PhaseValidating,
instance.PhaseInitializingRegistry, instance.PhaseReady, instance.PhaseDeleting} {
for _, operation := range []struct {
phase instance.Phase
apply func(*instance.Instance, instance.CapabilityObservation) error
}{
{instance.PhaseValidating, (*instance.Instance).AssessManagement},
{instance.PhaseReady, (*instance.Instance).AssessReadiness},
{instance.PhaseInitializingRegistry, func(i *instance.Instance, o instance.CapabilityObservation) error {
_, err := i.PlanRegistryPreparation(o)
return err
}},
{instance.PhaseInitializingRegistry, func(i *instance.Instance, o instance.CapabilityObservation) error {
return i.AssessRegistryResult(instance.RegistryReadBack(o))
}},
} {
if !deleting && operation.phase == phase {
continue
}
i := lifecycleInstance(t, instance.Snapshot{Phase: phase}, deleting)
before := i.Snapshot()
if err := operation.apply(i, capability(t, i, completeChecks(), instance.RegistryUsable)); err == nil {
t.Fatalf("phase %s deleting=%t accepted operation for %s", phase, deleting, operation.phase)
}
if i.Snapshot() != before {
t.Fatal("rejected operation mutated snapshot")
}
}
}
}
}
func TestOldGenerationObservationDoesNotReplaceEvidence(t *testing.T) {
i := readyInstance(t)
target := i.Target()
revision, err := instance.NewRevision(target.Revision().Value() + 1)
if err != nil {
t.Fatal(err)
}
other, err := instance.NewObservationTarget(target.Identity(), revision, target.Definition())
if err != nil {
t.Fatal(err)
}
o, err := instance.NewCapabilityObservation(other, testServerVersion, completeChecks(), instance.RegistryUsable)
if err != nil {
t.Fatal(err)
}
before := i.Snapshot()
if err := i.AssessReadiness(o); err == nil || i.Snapshot() != before {
t.Fatal("mismatched generation observation was accepted")
}
if err := i.RequireProvisioningReady(); err != nil {
t.Fatal("rejected unrelated input changed previously accepted evidence", err)
}
}
func TestPreparationFailureCannotEstablishReadiness(t *testing.T) {
i := lifecycleInstance(t, instance.Snapshot{Phase: instance.PhaseInitializingRegistry}, false)
for _, failure := range []instance.Failure{instance.DependencyUnavailable, instance.AuthenticationFailed,
instance.InsufficientPrivileges, instance.RegistryIncompatible} {
result, err := instance.RegistryPreparationFailed(i.Target(), failure)
if err != nil {
t.Fatal(err)
}
if err := i.AssessRegistryResult(result); err != nil {
t.Fatal(err)
}
if s := i.Snapshot(); s.Failure != failure || s.Readiness != instance.NotReady ||
s.Phase != instance.PhaseInitializingRegistry || i.RequireProvisioningReady() == nil {
t.Fatalf("failed operation accepted: %+v", s)
}
}
for _, failure := range []instance.Failure{instance.NoFailure, 255} {
if _, err := instance.RegistryPreparationFailed(i.Target(), failure); err == nil {
t.Fatal("invalid failure accepted")
}
}
}
func TestCapabilityInputsAndLifecycleGuards(t *testing.T) {
i := readyInstance(t)
if _, err := instance.NewCapabilityObservation(instance.ObservationTarget{}, testServerVersion,
completeChecks(), instance.RegistryUsable); err == nil {
t.Fatal("invalid target accepted")
}
if _, err := instance.RegistryPreparationFailed(instance.ObservationTarget{}, instance.DependencyUnavailable); err == nil {
t.Fatal("invalid failure target accepted")
}
for _, method := range []func(instance.CapabilityObservation) error{
i.AssessManagement, i.AssessReadiness,
func(o instance.CapabilityObservation) error { _, err := i.PlanRegistryPreparation(o); return err },
func(o instance.CapabilityObservation) error {
return i.AssessRegistryResult(instance.RegistryReadBack(o))
},
} {
before := i.Snapshot()
if err := method(instance.CapabilityObservation{}); err == nil || i.Snapshot() != before {
t.Fatal("mismatched observation accepted or mutated state")
}
}
old := i.Snapshot()
old.ObservedRevision = 0
changed := lifecycleInstance(t, old, false)
if err := changed.AssessReadiness(capability(t, changed, completeChecks(), instance.RegistryUsable)); err != nil {
t.Fatal(err)
}
if s := changed.Snapshot(); s.Phase != instance.PhaseValidating || s.ObservedRevision != 0 || s.Readiness != instance.Unknown {
t.Fatalf("changed generation accepted old checkpoint: %+v", s)
}
o, err := instance.NewCapabilityObservation(i.Target(), "", completeChecks(), instance.RegistryUsable)
if err != nil {
t.Fatal(err)
}
if err := i.AssessReadiness(o); err != nil {
t.Fatal(err)
}
if i.Snapshot().Failure != instance.ObservationIncomplete {
t.Fatal("missing version accepted")
}
}