diff --git a/.gitea/workflows/verify.yml b/.gitea/workflows/verify.yml index e7f1b3b..7a5f660 100644 --- a/.gitea/workflows/verify.yml +++ b/.gitea/workflows/verify.yml @@ -40,4 +40,32 @@ jobs: cache: true - 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 diff --git a/Makefile b/Makefile index c570f9d..750f4c7 100644 --- a/Makefile +++ b/Makefile @@ -67,6 +67,14 @@ test: manifests generate fmt vet setup-envtest ## Run tests. lint: golangci-lint ## Run golangci-lint linter "$(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 lint-fix: golangci-lint ## Run golangci-lint linter and perform fixes "$(GOLANGCI_LINT)" run --fix diff --git a/docs/database/README.md b/docs/database/README.md index 8e6c24d..ac67841 100644 --- a/docs/database/README.md +++ b/docs/database/README.md @@ -33,6 +33,31 @@ registry 准备决策与完整回读、Ready 重验及本轮 evidence 前置检 - 当前代码只检查 Instance 供应前置条件,不授予 Tenant 所有权或外部写入权限,也不表示 Database API 已经可用。 +## 管理凭据与连接切片 + +`application.InstanceService` 适配自原项目固定基线 +[`internal/instance/service.go`](https://git.ddupan.top/panxiao81/postgresql-tenant-operator/src/commit/dae546e58efa1be81e930861c87f7fb13bb12113/internal/instance/service.go), +保留 CredentialReader、Connector、Database 的装配边界及串行操作/释放规则。 +具体连接池完全由 pgxpool v5.11.0 提供,PostgreSQL adapter 不拥有凭据缓存、轮换流程或任意观测回调。 + +相对源基线有两项按已批准合同作出的必要修改:管理凭据从固定 controller namespace 的 +Kubernetes Secret 直接读取;每轮比较有效用户名、密码,检测到变化即关闭旧连接并重新装配。 +Secret metadata 和无关字段变化不重建连接。观测后再次读取 Secret,中途有效值变化则丢弃结果, +不把旧连接的成功作为新凭据有效的证据;这不构成跨 Kubernetes/PostgreSQL 的原子事务。 +Instance UID、endpoint 或凭据引用变化也会释放旧连接;Forget/Close 只释放本地资源。 + +当前只有 `ObserveVersion` 版本查询用例,不能产生完整 CapabilityObservation 或 Ready。 +controller 接入、Secret watch、finalizer、registry 与真实权限检查仍待后续切片;并发 CR 更新 +必须由调用者通过 resourceVersion 校验。应用层沿用源实现的串行处理,本阶段未引入新的调度框架。 + +运行 `make test-database-integration` 验证真实 API server + 一次性 PostgreSQL;fixture 不接受外部 +DSN,镜像固定摘要,使用随机本机回环端口并在退出时删除测试容器。覆盖缺失/错误凭据、RBAC、 +namespace 边界、有效值轮换、metadata 无关变化、中途轮换、重建/重试、并发读取、Forget/Close +与 TLS DNS/IP SAN、错误 CA/主机名和禁止明文降级。CI 使用 Pod runner 执行,由 runner 提供 +可用的 Docker,workflow 只做预检、不自行启动 daemon;不依赖 VM。普通 lint 之外还检查 +integration 标签代码。领域单测、真实 API 行为与真实 PostgreSQL 行为分别验收,不以本切片 +替代整个 Instance controller 的集成验收。 + ## 设计入口 - [系统规格](specification.md):规范性行为与验收标准; diff --git a/docs/database/development.md b/docs/database/development.md index b51cc3a..d2cc4c6 100644 --- a/docs/database/development.md +++ b/docs/database/development.md @@ -3,6 +3,34 @@ > 本页迁入作为 Database 模块的测试分层与 fixture 合同。旧项目的 Make target、devcontainer > 和脚手架版本尚未适配 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 使用 envtest 或一次性 Kind,另外两个依赖使用一次性 容器。这样既避免污染真实数据,也能把启动顺序固化为命令。 diff --git a/go.mod b/go.mod index 6e7d4b7..689551f 100644 --- a/go.mod +++ b/go.mod @@ -3,6 +3,7 @@ module git.ddupan.top/panxiao81/ayatori go 1.27.1 require ( + github.com/jackc/pgx/v5 v5.11.0 k8s.io/api v0.37.0 k8s.io/apimachinery 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/grpc-ecosystem/grpc-gateway/v2 v2.29.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/modern-go/concurrent v0.0.0-20180306012644-bacd9c7ef1dd // indirect github.com/modern-go/reflect2 v1.0.3-0.20250322232337-35a7c28c31ee // indirect diff --git a/go.sum b/go.sum index 7684ec1..850d4aa 100644 --- a/go.sum +++ b/go.sum @@ -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/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= +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/go.mod h1:e30LSqwooZae/UwlEbR2852Gd8hjQvJoHmT4TnhNGBo= 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/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/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/inf.v0 v0.9.1 h1:73M5CoZyi3ZLMOyDlQh031Cx6N9NDJ2Vvfl76EDAgDc= 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/go.mod h1:K4uyk7z7BCEPqu6E+C64Yfv1cQ7kz7rIZviUmN+EgEM= k8s.io/api v0.37.0 h1:Z//Vj9N7RA/yS2sDmxyeo7h+RR4zbUrd2vrd3Z0TbB4= diff --git a/internal/database/adapter/kubernetes/credentials.go b/internal/database/adapter/kubernetes/credentials.go new file mode 100644 index 0000000..2e01583 --- /dev/null +++ b/internal/database/adapter/kubernetes/credentials.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 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])) +} diff --git a/internal/database/adapter/postgresql/connector.go b/internal/database/adapter/postgresql/connector.go new file mode 100644 index 0000000..b0194c9 --- /dev/null +++ b/internal/database/adapter/postgresql/connector.go @@ -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 +} diff --git a/internal/database/adapter/postgresql/credentials_integration_test.go b/internal/database/adapter/postgresql/credentials_integration_test.go new file mode 100644 index 0000000..d206663 --- /dev/null +++ b/internal/database/adapter/postgresql/credentials_integration_test.go @@ -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) +} diff --git a/internal/database/adapter/postgresql/fixture_integration_test.go b/internal/database/adapter/postgresql/fixture_integration_test.go new file mode 100644 index 0000000..c791be6 --- /dev/null +++ b/internal/database/adapter/postgresql/fixture_integration_test.go @@ -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 + `) +} diff --git a/internal/database/adapter/postgresql/tls_integration_test.go b/internal/database/adapter/postgresql/tls_integration_test.go new file mode 100644 index 0000000..32bae72 --- /dev/null +++ b/internal/database/adapter/postgresql/tls_integration_test.go @@ -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") + } +} diff --git a/internal/database/application/credentials.go b/internal/database/application/credentials.go new file mode 100644 index 0000000..82dc959 --- /dev/null +++ b/internal/database/application/credentials.go @@ -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) +} diff --git a/internal/database/application/credentials_test.go b/internal/database/application/credentials_test.go new file mode 100644 index 0000000..676de63 --- /dev/null +++ b/internal/database/application/credentials_test.go @@ -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") + } + } +} diff --git a/internal/database/application/instance_service.go b/internal/database/application/instance_service.go new file mode 100644 index 0000000..489aaee --- /dev/null +++ b/internal/database/application/instance_service.go @@ -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) + } +} diff --git a/internal/database/application/instance_service_test.go b/internal/database/application/instance_service_test.go new file mode 100644 index 0000000..4a2617f --- /dev/null +++ b/internal/database/application/instance_service_test.go @@ -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") + } +}