From e590542b0bbc2733ac397b66a17bb3a7a68e2f65 Mon Sep 17 00:00:00 2001 From: panxiao81 Date: Mon, 21 Sep 2026 09:01:05 +0000 Subject: [PATCH] =?UTF-8?q?feat:=20=E8=BF=81=E7=A7=BB=20PostgreSQL=20regis?= =?UTF-8?q?try=20=E6=89=80=E6=9C=89=E6=9D=83=E5=AD=98=E5=82=A8=E4=B8=8E?= =?UTF-8?q?=E6=81=A2=E5=A4=8D=E6=B5=8B=E8=AF=95?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- docs/database/README.md | 20 +- docs/database/development.md | 9 +- go.mod | 15 +- go.sum | 47 +- .../migrations/001_create_registry.sql | 25 + .../adapter/postgresql/registry/registry.go | 326 +++++++++++++ .../adapter/postgresql/registry/sql.go | 68 +++ .../postgresql/registry_integration_test.go | 427 ++++++++++++++++++ 8 files changed, 919 insertions(+), 18 deletions(-) create mode 100644 internal/database/adapter/postgresql/registry/migrations/001_create_registry.sql create mode 100644 internal/database/adapter/postgresql/registry/registry.go create mode 100644 internal/database/adapter/postgresql/registry/sql.go create mode 100644 internal/database/adapter/postgresql/registry_integration_test.go diff --git a/docs/database/README.md b/docs/database/README.md index ac67841..0f255b2 100644 --- a/docs/database/README.md +++ b/docs/database/README.md @@ -47,7 +47,7 @@ Secret metadata 和无关字段变化不重建连接。观测后再次读取 Sec Instance UID、endpoint 或凭据引用变化也会释放旧连接;Forget/Close 只释放本地资源。 当前只有 `ObserveVersion` 版本查询用例,不能产生完整 CapabilityObservation 或 Ready。 -controller 接入、Secret watch、finalizer、registry 与真实权限检查仍待后续切片;并发 CR 更新 +controller 接入、Secret watch、finalizer、registry 观测装配与真实权限检查仍待后续切片;并发 CR 更新 必须由调用者通过 resourceVersion 校验。应用层沿用源实现的串行处理,本阶段未引入新的调度框架。 运行 `make test-database-integration` 验证真实 API server + 一次性 PostgreSQL;fixture 不接受外部 @@ -58,6 +58,24 @@ namespace 边界、有效值轮换、metadata 无关变化、中途轮换、重 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/端到端切片验收。 + ## 设计入口 - [系统规格](specification.md):规范性行为与验收标准; diff --git a/docs/database/development.md b/docs/database/development.md index d2cc4c6..cb7bee2 100644 --- a/docs/database/development.md +++ b/docs/database/development.md @@ -3,7 +3,7 @@ > 本页迁入作为 Database 模块的测试分层与 fixture 合同。旧项目的 Make target、devcontainer > 和脚手架版本尚未适配 Ayatori;实现时应复用 Ayatori 现有工具链,并保持这里定义的测试边界。 -## Ayatori 已接入的凭据切片测试 +## Ayatori 已接入的凭据与 registry 切片测试 本节命令已在 Ayatori 接入;以下历史 Compose/Kind 操作仍属于迁入的目标合同。 @@ -28,7 +28,12 @@ runner 后端修复上线,不能从本地测试通过推断远端已经可用 fixture 启动失败会保留退出错误与 stderr,并遮蔽测试密码, 以区分缺少命令、daemon 不可达、权限和镜像拉取失败。 -这些测试尚不包含 Instance CRD/controller、Secret watch、status/finalizer 事件链、registry、 +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 diff --git a/go.mod b/go.mod index 689551f..01c9113 100644 --- a/go.mod +++ b/go.mod @@ -4,6 +4,7 @@ go 1.27.1 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/apimachinery v0.37.0 k8s.io/client-go v0.37.0 @@ -12,6 +13,10 @@ require ( require ( 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/beorn7/perks v1.0.1 // indirect github.com/blang/semver/v4 v4.0.0 // indirect @@ -44,11 +49,14 @@ require ( github.com/google/gnostic-models v0.7.0 // indirect github.com/google/uuid v1.6.0 // indirect github.com/grpc-ecosystem/grpc-gateway/v2 v2.29.0 // indirect + github.com/huandu/xstrings v1.5.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/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/reflect2 v1.0.3-0.20250322232337-35a7c28c31ee // indirect github.com/munnerz/goautoneg v0.0.0-20191010083416-a7dc8b61c822 // indirect @@ -57,6 +65,8 @@ require ( github.com/prometheus/client_model v0.6.2 // indirect github.com/prometheus/common v0.70.0 // indirect github.com/prometheus/procfs v0.21.1 // indirect + github.com/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/pflag v1.0.10 // indirect github.com/x448/float16 v0.8.4 // indirect @@ -72,14 +82,15 @@ require ( go.uber.org/multierr v1.11.0 // indirect go.uber.org/zap v1.27.1 // 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/net v0.57.0 // indirect golang.org/x/oauth2 v0.36.0 // indirect golang.org/x/sync v0.22.0 // indirect golang.org/x/sys v0.47.0 // indirect 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 gomodules.xyz/jsonpatch/v2 v2.4.0 // indirect google.golang.org/genproto/googleapis/api v0.0.0-20260526163538-3dc84a4a5aaa // indirect diff --git a/go.sum b/go.sum index 850d4aa..f8da443 100644 --- a/go.sum +++ b/go.sum @@ -1,7 +1,13 @@ cel.dev/expr v0.25.1 h1:1KrZg61W6TWSxuNZ37Xy49ps13NUovb66QLprthtwi4= cel.dev/expr v0.25.1/go.mod h1:hrXvqGP6G6gyx8UAHSHJ5RGk//1Oj5nXQ2NI02Nrsg4= -github.com/Masterminds/semver/v3 v3.4.0 h1:Zog+i5UMtVoCU8oKka5P7i9q9HgrJeGzI9SA1Xbatp0= -github.com/Masterminds/semver/v3 v3.4.0/go.mod h1:4V+yj/TJE1HU9XfppCwVMZq3I84lprf4nC11bSS5beM= +dario.cat/mergo v1.0.2 h1:85+piFYR1tMbRrLcDwR18y4UKJ3aH1Tbzi24VRW1TK8= +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/go.mod h1:GKmUxMtwp6ZgGwZSva4eWPC5mS6vUAmOABFgjdkM7Nw= 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/felixge/httpsnoop v1.0.4 h1:NFTV2Zj1bL4mc9sqWACXbQFVBBg2W3GPvqp8/ESS2Wg= 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/go.mod h1:8jBTzvmWwFyi3Pb8djgCCO5IBqzKJ/Jwo8TRcHyHii0= github.com/fxamacker/cbor/v2 v2.9.1 h1:2rWm8B193Ll4VdjsJY28jxs70IdDsHRWgQYAI80+rMQ= @@ -89,6 +97,8 @@ github.com/google/uuid v1.6.0 h1:NIvaJDMOsjHA8n1jAhLSgzrAzy1Hgr+hNrb57e+94F0= github.com/google/uuid v1.6.0/go.mod h1:TIyPZe4MgqvfeYDBFedMoGGpEw/LqOeaOT+nhxU+yHo= github.com/grpc-ecosystem/grpc-gateway/v2 v2.29.0 h1:5VipnvEpbqr2gA2VbM+nYVbkIF28c5ZQfqCBQ5g2xfk= github.com/grpc-ecosystem/grpc-gateway/v2 v2.29.0/go.mod h1:Hyl3n6Twe1hvtd9XUXDec4pTvgMSEixRuQKPTMH2bNs= +github.com/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/go.mod h1:vpF70FUmC8bwa3OWnCshd2FqLfsEA9PFc4w1p2J65bw= github.com/jackc/pgpassfile v1.0.0 h1:/6Hmqy13Ss2zCq62VdNG8tM1wchn8zjSGOBJ6icpsIM= @@ -99,6 +109,8 @@ 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/go.mod h1:e30LSqwooZae/UwlEbR2852Gd8hjQvJoHmT4TnhNGBo= github.com/klauspost/compress v1.19.0 h1:sXLILfc9jV2QYWkzFOPWStmcUVH2RHEB1JCdY2oVvCQ= @@ -109,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/kylelemons/godebug v1.1.0 h1:RPNrshWIDI6G2gRW9EHilWtl7Z6Sb1BR0xunSBf0SNc= 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-20180306012644-bacd9c7ef1dd h1:TRLaZ9cD/w8PVh93nsPXa1VrQ6jlwL5oN8l14QlcNfg= github.com/modern-go/concurrent v0.0.0-20180306012644-bacd9c7ef1dd/go.mod h1:6dJC0mAP4ikYIbvyc7fijjWJddQyLn8Ig3JB5CqoB9Q= @@ -137,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/go.mod h1:MaRKkUm5W0goXpeCfT7UZI6fk/L7L7so1lCWt35ZSgc= 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/go.mod h1:7C1pvHqHw5A4vrJfjNwvOdzYu0Gml16OCs2GRiTUUS4= github.com/spf13/pflag v1.0.9/go.mod h1:McXfInJRrz4CZXVZOBLb0bTZqETkiAhM9Iw0y3An2Bg= @@ -147,8 +167,8 @@ github.com/stretchr/objx v0.5.3 h1:jmXUvGomnU1o3W/V5h2VEradbpJDwGrzugQQvL0POH4= github.com/stretchr/objx v0.5.3/go.mod h1:rDQraq+vQZU7Fde9LOZLr8Tax6zZvy4kuNKF+QYS+U0= github.com/stretchr/testify v1.3.0/go.mod h1:M5WIy9Dh21IEIfnGCwXGc5bZfKNJtfHm1UVUgZn+9EI= github.com/stretchr/testify v1.7.0/go.mod h1:6Fq8oRcR53rry900zMqJjRRixrwX3KX962/h/Wwjteg= -github.com/stretchr/testify v1.11.1 h1:7s2iGBzp5EwR7/aIZr8ao5+dra3wiQyKjjFuvgVKu7U= -github.com/stretchr/testify v1.11.1/go.mod h1:wZwfW3scLgRK+23gO65QZefKpKQRnfz6sD981Nm4B6U= +github.com/stretchr/testify v1.12.1 h1:EuwCh5fleGS7H32xRwO3wRGT7DxrDhLAT6FF8MpWDWE= +github.com/stretchr/testify v1.12.1/go.mod h1:MDEgiDPPsNp5cuIrHPPCyornHKgEVbtFUmoNlxoYthg= github.com/x448/float16 v0.8.4 h1:qLwI1I70+NjRFUR3zs1JPUCgaCXSh3SW62uAKT1mSBM= github.com/x448/float16 v0.8.4/go.mod h1:14CWIYCyZA/cWjXOioeEpHeN/83MdbZDRQHoFcYsOfg= go.opentelemetry.io/auto/sdk v1.2.1 h1:jXsnJ4Lmnqd11kwkBV2LgLoFMZKizbCi5fNZ/ipaZ64= @@ -179,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.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/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.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/go.mod h1:J1xhfL/vlindoeF/aINzNzt2Bket5bjo9sdOYzOsU80= -golang.org/x/mod v0.37.0 h1:vF1DjpVEshcIqoEaauuHebaLk1O1forxjxBaVn884JQ= -golang.org/x/mod v0.37.0/go.mod h1:m8S8VeM9r4dzDwjrKO0a1sZP3YjeMamRRlD+fmR2Q/0= +golang.org/x/mod v0.38.0 h1:MECBjubtXD7yj4HrhIUcywNaGeNVUdfVnxmPajOk4yk= +golang.org/x/mod v0.38.0/go.mod h1:V6Xz0pq8TQ3dGqVQ1FVHuelZpAL0uNhSkk9ogYP3c40= golang.org/x/net v0.57.0 h1:K5+3DljvIuDG9/Jv9rvyMywYNFCQ9RSUY6OOTTkT+tE= golang.org/x/net v0.57.0/go.mod h1:KpXc8iv+r3XplLAG/f7Jsf9RPszJzdR0f58q9vGOuEU= golang.org/x/oauth2 v0.36.0 h1:peZ/1z27fi9hUOFCAZaHyrpWG5lwe0RJEEEeH0ThlIs= @@ -195,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/term v0.45.0 h1:NwWyBmoJCbfTHpxrWoZ9C6/VxOf7ic219I8xZZFdrf0= 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.40.0/go.mod h1:hpnzDAfGV753zIKo+wk3u1bVKCGPbrnF7+7LBF/UHVY= +golang.org/x/text v0.41.0 h1:vz/seA0lnX87Othu2f/0L24RcgrXD9/YFTSuGjj3rH8= +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/go.mod h1:Y4YMaQmXwGQZoFaVFk4YpCt4FLQMYKZe9oeV/f4MSno= -golang.org/x/tools v0.47.0 h1:7Kn5x/d1svx/PzryTsqeoZN4TZwqeH5pGWjefhLi/1Q= -golang.org/x/tools v0.47.0/go.mod h1:dFHnyTvFWY212G+h7ZY4Vsp/K3U4/7W9TyVaAul8uCA= +golang.org/x/tools v0.48.0 h1:3+hClM1aLL5mjMKm5ovokw9epgRXPuu2tILgismM6RE= +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/go.mod h1:AH3dM2RI6uoBZxn3LVrfvJ3E0/9dG4cSrbuBJT4moAY= gonum.org/v1/gonum v0.17.0 h1:VbpOemQlsSMrYmn7T2OUvQ4dqxQXU+ouZFQsZOx50z4= @@ -214,8 +237,6 @@ 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/go.mod h1:HTf+CrKn2C3g5S8VImy6tdcUvCska2kB7j23XfzDpco= 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/go.mod h1:p8EYWUEYMpynmqDbY58zCKCFZw8pRWMG4EsWvDvM72M= gopkg.in/inf.v0 v0.9.1 h1:73M5CoZyi3ZLMOyDlQh031Cx6N9NDJ2Vvfl76EDAgDc= diff --git a/internal/database/adapter/postgresql/registry/migrations/001_create_registry.sql b/internal/database/adapter/postgresql/registry/migrations/001_create_registry.sql new file mode 100644 index 0000000..24889bb --- /dev/null +++ b/internal/database/adapter/postgresql/registry/migrations/001_create_registry.sql @@ -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; diff --git a/internal/database/adapter/postgresql/registry/registry.go b/internal/database/adapter/postgresql/registry/registry.go new file mode 100644 index 0000000..12a60bb --- /dev/null +++ b/internal/database/adapter/postgresql/registry/registry.go @@ -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 +} diff --git a/internal/database/adapter/postgresql/registry/sql.go b/internal/database/adapter/postgresql/registry/sql.go new file mode 100644 index 0000000..fe75d57 --- /dev/null +++ b/internal/database/adapter/postgresql/registry/sql.go @@ -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` diff --git a/internal/database/adapter/postgresql/registry_integration_test.go b/internal/database/adapter/postgresql/registry_integration_test.go new file mode 100644 index 0000000..66c68bb --- /dev/null +++ b/internal/database/adapter/postgresql/registry_integration_test.go @@ -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 ®istryFixture{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) + } +} -- 2.54.0