Author SHA1 Message Date
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 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
15 changed files with 1468 additions and 1 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
+25
View File
@@ -33,6 +33,31 @@ registry 准备决策与完整回读、Ready 重验及本轮 evidence 前置检
- 当前代码只检查 Instance 供应前置条件,不授予 Tenant 所有权或外部写入权限,也不表示 - 当前代码只检查 Instance 供应前置条件,不授予 Tenant 所有权或外部写入权限,也不表示
Database API 已经可用。 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 的集成验收。
## 设计入口 ## 设计入口
- [系统规格](specification.md):规范性行为与验收标准; - [系统规格](specification.md):规范性行为与验收标准;
+28
View File
@@ -3,6 +3,34 @@
> 本页迁入作为 Database 模块的测试分层与 fixture 合同。旧项目的 Make target、devcontainer > 本页迁入作为 Database 模块的测试分层与 fixture 合同。旧项目的 Make target、devcontainer
> 和脚手架版本尚未适配 Ayatori;实现时应复用 Ayatori 现有工具链,并保持这里定义的测试边界。 > 和脚手架版本尚未适配 Ayatori;实现时应复用 Ayatori 现有工具链,并保持这里定义的测试边界。
## Ayatori 已接入的凭据切片测试
本节命令已在 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 不可达、权限和镜像拉取失败。
这些测试尚不包含 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,另外两个依赖使用一次性
容器。这样既避免污染真实数据,也能把启动顺序固化为命令。 容器。这样既避免污染真实数据,也能把启动顺序固化为命令。
+4
View File
@@ -3,6 +3,7 @@ module git.ddupan.top/panxiao81/ayatori
go 1.27.1 go 1.27.1
require ( require (
github.com/jackc/pgx/v5 v5.11.0
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
@@ -44,6 +45,9 @@ require (
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/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/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
+10
View File
@@ -91,6 +91,14 @@ github.com/grpc-ecosystem/grpc-gateway/v2 v2.29.0 h1:5VipnvEpbqr2gA2VbM+nYVbkIF2
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/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/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=
@@ -138,6 +146,7 @@ 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.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 h1:7s2iGBzp5EwR7/aIZr8ao5+dra3wiQyKjjFuvgVKu7U=
github.com/stretchr/testify v1.11.1/go.mod h1:wZwfW3scLgRK+23gO65QZefKpKQRnfz6sD981Nm4B6U= github.com/stretchr/testify v1.11.1/go.mod h1:wZwfW3scLgRK+23gO65QZefKpKQRnfz6sD981Nm4B6U=
github.com/x448/float16 v0.8.4 h1:qLwI1I70+NjRFUR3zs1JPUCgaCXSh3SW62uAKT1mSBM= github.com/x448/float16 v0.8.4 h1:qLwI1I70+NjRFUR3zs1JPUCgaCXSh3SW62uAKT1mSBM=
@@ -211,6 +220,7 @@ gopkg.in/evanphx/json-patch.v4 v4.13.0 h1:czT3CmqEaQ1aanPc5SdlgQrrEIb8w/wwCvWWnf
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,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")
}
}