Compare commits
1
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
f6d552f944
|
@@ -45,8 +45,7 @@ jobs:
|
|||||||
make lint-database-integration
|
make lint-database-integration
|
||||||
|
|
||||||
database-integration:
|
database-integration:
|
||||||
runs-on: [self-hosted, pod]
|
runs-on: [self-hosted, vm]
|
||||||
timeout-minutes: 30
|
|
||||||
steps:
|
steps:
|
||||||
- name: Checkout
|
- name: Checkout
|
||||||
uses: actions/checkout@de0fac2e4500dabe0009e67214ff5f5447ce83dd
|
uses: actions/checkout@de0fac2e4500dabe0009e67214ff5f5447ce83dd
|
||||||
@@ -59,13 +58,5 @@ jobs:
|
|||||||
go-version-file: go.mod
|
go-version-file: go.mod
|
||||||
cache: true
|
cache: true
|
||||||
|
|
||||||
# Runner 提供本 job 可用的 Docker;workflow 只验证,不重复启动 daemon。
|
- name: Test Database credentials with real backends
|
||||||
- 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
|
run: make test-database-integration
|
||||||
|
|||||||
+2
-29
@@ -46,44 +46,17 @@ Secret metadata 和无关字段变化不重建连接。观测后再次读取 Sec
|
|||||||
不把旧连接的成功作为新凭据有效的证据;这不构成跨 Kubernetes/PostgreSQL 的原子事务。
|
不把旧连接的成功作为新凭据有效的证据;这不构成跨 Kubernetes/PostgreSQL 的原子事务。
|
||||||
Instance UID、endpoint 或凭据引用变化也会释放旧连接;Forget/Close 只释放本地资源。
|
Instance UID、endpoint 或凭据引用变化也会释放旧连接;Forget/Close 只释放本地资源。
|
||||||
|
|
||||||
当前通过 `ObserveMetadata` 读取服务器版本和可用扩展,`ObserveVersion` 只是其版本读取便捷入口,
|
当前只有 `ObserveVersion` 版本查询用例,不能产生完整 CapabilityObservation 或 Ready。
|
||||||
不能产生完整 CapabilityObservation 或 Ready。
|
|
||||||
controller 接入、Secret watch、finalizer、registry 观测装配与真实权限检查仍待后续切片;并发 CR 更新
|
controller 接入、Secret watch、finalizer、registry 观测装配与真实权限检查仍待后续切片;并发 CR 更新
|
||||||
必须由调用者通过 resourceVersion 校验。应用层沿用源实现的串行处理,本阶段未引入新的调度框架。
|
必须由调用者通过 resourceVersion 校验。应用层沿用源实现的串行处理,本阶段未引入新的调度框架。
|
||||||
|
|
||||||
运行 `make test-database-integration` 验证真实 API server + 一次性 PostgreSQL;fixture 不接受外部
|
运行 `make test-database-integration` 验证真实 API server + 一次性 PostgreSQL;fixture 不接受外部
|
||||||
DSN,镜像固定摘要,使用随机本机回环端口并在退出时删除测试容器。覆盖缺失/错误凭据、RBAC、
|
DSN,镜像固定摘要,使用随机本机回环端口并在退出时删除测试容器。覆盖缺失/错误凭据、RBAC、
|
||||||
namespace 边界、有效值轮换、metadata 无关变化、中途轮换、重建/重试、并发读取、Forget/Close
|
namespace 边界、有效值轮换、metadata 无关变化、中途轮换、重建/重试、并发读取、Forget/Close
|
||||||
与 TLS DNS/IP SAN、错误 CA/主机名和禁止明文降级。CI 使用 Pod runner 执行,由 runner 提供
|
与 TLS DNS/IP SAN、错误 CA/主机名和禁止明文降级。CI 使用 VM runner 执行,普通 lint 之外还检查
|
||||||
可用的 Docker,workflow 只做预检、不自行启动 daemon;不依赖 VM。普通 lint 之外还检查
|
|
||||||
integration 标签代码。领域单测、真实 API 行为与真实 PostgreSQL 行为分别验收,不以本切片
|
integration 标签代码。领域单测、真实 API 行为与真实 PostgreSQL 行为分别验收,不以本切片
|
||||||
替代整个 Instance controller 的集成验收。
|
替代整个 Instance controller 的集成验收。
|
||||||
|
|
||||||
## 服务器 metadata 与扩展观测
|
|
||||||
|
|
||||||
SQL adapter 通过一条只读语句读取 `pg_catalog.current_setting('server_version')` 和
|
|
||||||
`pg_catalog.pg_available_extensions`,避免从已安装列表推断可用列表,且不依赖可修改的
|
|
||||||
`search_path`。查询失败丢弃整份结果;成功返回空列表与尚未观察严格区分。
|
|
||||||
参考 PostgreSQL 的 [pg_available_extensions](https://www.postgresql.org/docs/18/view-pg-available-extensions.html)
|
|
||||||
与 [CREATE EXTENSION](https://www.postgresql.org/docs/18/sql-createextension.html) 合同:可用列表
|
|
||||||
表示服务器提供的扩展,不证明管理账号有安装权限,也不保证依赖和其他安装前提满足。
|
|
||||||
|
|
||||||
`InstanceService` 保留原有 CredentialReader → Connector → Database 边界,复用同一个
|
|
||||||
凭据读取、连接刷新和串行释放流程,不新增连接池封装或任意查询回调。每次重新查询 metadata,
|
|
||||||
并在 Secret 有效值回读一致后生成不可变的 `MetadataObservation`,绑定本次 target(含当前
|
|
||||||
generation),不绑定建池时的旧 target。结果不包含凭据,扩展集合不与 driver 的可变 slice 共享。
|
|
||||||
凭据中途变化、读取失败或查询失败时,返回零值观察并释放连接,不复用旧的扩展列表。
|
|
||||||
|
|
||||||
调用方可将 `Target()` 与 `Extensions()` 交给 Instance 的 `ObserveExtensions`;应用调用链
|
|
||||||
仍负责同轮次使用,不能持久化或跨轮缓存这份证据。metadata 读取不安装扩展、不初始化 registry、
|
|
||||||
不设置 Ready,也不授予 Tenant 写权限。管理权限矩阵、registry 完整回读及 controller 的
|
|
||||||
checkpoint/status/finalizer 链路仍是后续切片。
|
|
||||||
|
|
||||||
真实 API server + PostgreSQL 测试验证未安装扩展可被观察、名称保持大小写、search_path 遮蔽
|
|
||||||
不改变查询来源、低权限账号读取、权限撤回失败与恢复、Secret 中途变化丢弃扩展结果。
|
|
||||||
单元测试补充成功空列表、查询附带部分数据时丢弃、结果与可变 slice 隔离、每轮重新读取和
|
|
||||||
generation 变化时的目标绑定;原凭据/TLS/超时/并发测试沿同一 metadata 路径继续运行。
|
|
||||||
|
|
||||||
## Registry 所有权存储切片
|
## Registry 所有权存储切片
|
||||||
|
|
||||||
`adapter/postgresql/registry` 直接迁入上述固定基线的 `internal/postgresql/registry`,
|
`adapter/postgresql/registry` 直接迁入上述固定基线的 `internal/postgresql/registry`,
|
||||||
|
|||||||
@@ -3,7 +3,7 @@
|
|||||||
> 本页迁入作为 Database 模块的测试分层与 fixture 合同。旧项目的 Make target、devcontainer
|
> 本页迁入作为 Database 模块的测试分层与 fixture 合同。旧项目的 Make target、devcontainer
|
||||||
> 和脚手架版本尚未适配 Ayatori;实现时应复用 Ayatori 现有工具链,并保持这里定义的测试边界。
|
> 和脚手架版本尚未适配 Ayatori;实现时应复用 Ayatori 现有工具链,并保持这里定义的测试边界。
|
||||||
|
|
||||||
## Ayatori 已接入的凭据、metadata 与 registry 切片测试
|
## Ayatori 已接入的凭据与 registry 切片测试
|
||||||
|
|
||||||
本节命令已在 Ayatori 接入;以下历史 Compose/Kind 操作仍属于迁入的目标合同。
|
本节命令已在 Ayatori 接入;以下历史 Compose/Kind 操作仍属于迁入的目标合同。
|
||||||
|
|
||||||
@@ -21,17 +21,7 @@ TLS 测试在临时目录生成一次性证书与私钥,不使用生产 CA。
|
|||||||
|
|
||||||
当前覆盖固定 namespace 的 Secret 读取与 RBAC、缺失/无效凭据恢复、有效凭据变化后的重连、
|
当前覆盖固定 namespace 的 Secret 读取与 RBAC、缺失/无效凭据恢复、有效凭据变化后的重连、
|
||||||
无关字段更新不重连、中途轮换时丢弃观察、会话重建、并发读取、本地连接释放和 TLS 验证。
|
无关字段更新不重连、中途轮换时丢弃观察、会话重建、并发读取、本地连接释放和 TLS 验证。
|
||||||
快速测试、lint 和 Database 集成测试均使用 Pod runner。按维护者于 2026-09-21 更新的接口
|
默认 Pod runner 跑快速测试与 lint,VM runner 跑带 Docker 的 Database 集成测试。
|
||||||
约定,runner 提供默认可用的 Docker;workflow 通过 `docker version` 和 `docker info` 预检,
|
|
||||||
不自行启动 daemon、不强制 storage driver 或覆盖 Docker endpoint。该约定的 CI 验收依赖
|
|
||||||
runner 后端修复上线,不能从本地测试通过推断远端已经可用。
|
|
||||||
fixture 启动失败会保留退出错误与 stderr,并遮蔽测试密码,
|
|
||||||
以区分缺少命令、daemon 不可达、权限和镜像拉取失败。
|
|
||||||
|
|
||||||
metadata 测试验证版本与可用扩展的只读查询,包括未安装扩展、大小写保持、search_path 遮蔽、
|
|
||||||
低权限账号读取、catalog 访问被撤回后的失败与恢复,以及凭据中途变化时同时丢弃版本和扩展。
|
|
||||||
权限撤回只修改每个场景自建 PostgreSQL 容器的 ACL;不连接现有服务。
|
|
||||||
可用列表不等于安装权限,这些检查不替代后续的完整管理权限矩阵或 Instance Ready 验收。
|
|
||||||
|
|
||||||
registry 测试独立使用一次性 PostgreSQL,不启动 Kubernetes API server。覆盖重复和并发迁移、
|
registry 测试独立使用一次性 PostgreSQL,不启动 Kubernetes API server。覆盖重复和并发迁移、
|
||||||
所有权唯一约束、并发占用、Retain 墓碑与 Delete 幂等、连接重建、取消恢复、未知 schema 与
|
所有权唯一约束、并发占用、Retain 墓碑与 Delete 幂等、连接重建、取消恢复、未知 schema 与
|
||||||
|
|||||||
@@ -47,6 +47,14 @@ func (d *database) Close() {
|
|||||||
d.pool.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) {
|
func (c Connector) Connect(ctx context.Context, endpoint instance.Endpoint, credentials application.Credentials) (application.Database, error) {
|
||||||
if err := endpoint.Validate(); err != nil {
|
if err := endpoint.Validate(); err != nil {
|
||||||
return nil, err
|
return nil, err
|
||||||
|
|||||||
@@ -131,17 +131,13 @@ func TestObservationDiscardsResultWhenCredentialsChange(t *testing.T) {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
observation, err := fixture.service.ObserveMetadata(fixture.ctx, fixture.target)
|
version, err := fixture.service.ObserveVersion(fixture.ctx, fixture.target)
|
||||||
if !errors.Is(err, application.ErrCredentialsChanged) {
|
if !errors.Is(err, application.ErrCredentialsChanged) {
|
||||||
t.Fatal("in-flight rotation was not detected")
|
t.Fatal("in-flight rotation was not detected")
|
||||||
}
|
}
|
||||||
if observation.Version() != "" {
|
if version != "" {
|
||||||
t.Fatal("observation returned data obtained with stale credentials")
|
t.Fatal("observation returned data obtained with stale credentials")
|
||||||
}
|
}
|
||||||
requested := instance.NewExtensionSet([]string{fixtureExtension})
|
|
||||||
if observation.Extensions().Check(requested).Decision != instance.ExtensionSupportUnobserved {
|
|
||||||
t.Fatal("observation returned extension support obtained with stale credentials")
|
|
||||||
}
|
|
||||||
if fixture.backendIDs(t) != "" {
|
if fixture.backendIDs(t) != "" {
|
||||||
t.Fatal("stale connection was retained after rotation")
|
t.Fatal("stale connection was retained after rotation")
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -20,7 +20,6 @@ package postgresql_test
|
|||||||
|
|
||||||
import (
|
import (
|
||||||
"context"
|
"context"
|
||||||
"errors"
|
|
||||||
"os/exec"
|
"os/exec"
|
||||||
"regexp"
|
"regexp"
|
||||||
"strconv"
|
"strconv"
|
||||||
@@ -42,7 +41,6 @@ import (
|
|||||||
const (
|
const (
|
||||||
fixtureHost = "fixture.invalid"
|
fixtureHost = "fixture.invalid"
|
||||||
fixtureUser = "postgres"
|
fixtureUser = "postgres"
|
||||||
fixtureExtension = "plpgsql"
|
|
||||||
dockerExec = "exec"
|
dockerExec = "exec"
|
||||||
fixtureImage = "postgres@sha256:18cfe3ef5e6815560c98237d6216d1e5119702fb0f3894c8785dd58b8bbe5d73"
|
fixtureImage = "postgres@sha256:18cfe3ef5e6815560c98237d6216d1e5119702fb0f3894c8785dd58b8bbe5d73"
|
||||||
fixturePassword = "AYATORI-TEST-ONLY-initial-password"
|
fixturePassword = "AYATORI-TEST-ONLY-initial-password"
|
||||||
@@ -57,7 +55,7 @@ func postgresFixture(t *testing.T, ctx context.Context) (string, int) {
|
|||||||
output, err := exec.CommandContext(ctx, "docker", "run", "--rm", "-d", "-p", "127.0.0.1::5432",
|
output, err := exec.CommandContext(ctx, "docker", "run", "--rm", "-d", "-p", "127.0.0.1::5432",
|
||||||
"-e", "POSTGRES_PASSWORD="+fixturePassword, fixtureImage).Output()
|
"-e", "POSTGRES_PASSWORD="+fixturePassword, fixtureImage).Output()
|
||||||
if err != nil {
|
if err != nil {
|
||||||
t.Fatalf("cannot start isolated PostgreSQL fixture: %s", fixtureCommandError(err))
|
t.Fatal("cannot start isolated PostgreSQL fixture")
|
||||||
}
|
}
|
||||||
id := strings.TrimSpace(string(output))
|
id := strings.TrimSpace(string(output))
|
||||||
if !regexp.MustCompile(`^[a-f0-9]{64}$`).MatchString(id) {
|
if !regexp.MustCompile(`^[a-f0-9]{64}$`).MatchString(id) {
|
||||||
@@ -72,7 +70,7 @@ func postgresFixture(t *testing.T, ctx context.Context) (string, int) {
|
|||||||
})
|
})
|
||||||
output, err = exec.CommandContext(ctx, "docker", "inspect", "--format", `{{(index (index .NetworkSettings.Ports "5432/tcp") 0).HostPort}}`, id).Output()
|
output, err = exec.CommandContext(ctx, "docker", "inspect", "--format", `{{(index (index .NetworkSettings.Ports "5432/tcp") 0).HostPort}}`, id).Output()
|
||||||
if err != nil {
|
if err != nil {
|
||||||
t.Fatalf("cannot inspect fixture port: %s", fixtureCommandError(err))
|
t.Fatal("cannot inspect fixture port")
|
||||||
}
|
}
|
||||||
port, err := strconv.Atoi(strings.TrimSpace(string(output)))
|
port, err := strconv.Atoi(strings.TrimSpace(string(output)))
|
||||||
if err != nil {
|
if err != nil {
|
||||||
@@ -89,59 +87,6 @@ func postgresFixture(t *testing.T, ctx context.Context) (string, int) {
|
|||||||
return id, port
|
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 {
|
func target(t *testing.T, port int, mode instance.TLSMode) instance.ObservationTarget {
|
||||||
t.Helper()
|
t.Helper()
|
||||||
id, err := instance.NewIdentity("fixture-uid", "fixture")
|
id, err := instance.NewIdentity("fixture-uid", "fixture")
|
||||||
|
|||||||
@@ -1,48 +0,0 @@
|
|||||||
/*
|
|
||||||
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
|
|
||||||
|
|
||||||
import (
|
|
||||||
"context"
|
|
||||||
|
|
||||||
"git.ddupan.top/panxiao81/ayatori/internal/database/application"
|
|
||||||
)
|
|
||||||
|
|
||||||
// 使用 pg_catalog 限定名称,避免管理账号的 search_path 改变查询来源。
|
|
||||||
// 一个语句读取版本和可用列表;ARRAY 子查询在无行时返回空数组,而非 NULL。
|
|
||||||
// 不查询 pg_extension:已安装集合不能代表服务器提供的全部扩展。
|
|
||||||
const inspectMetadataStatement = `
|
|
||||||
SELECT
|
|
||||||
pg_catalog.current_setting('server_version'),
|
|
||||||
ARRAY(
|
|
||||||
SELECT name::text
|
|
||||||
FROM pg_catalog.pg_available_extensions
|
|
||||||
ORDER BY name
|
|
||||||
)`
|
|
||||||
|
|
||||||
func (d *database) InspectMetadata(ctx context.Context) (application.DatabaseMetadata, error) {
|
|
||||||
var metadata application.DatabaseMetadata
|
|
||||||
err := d.pool.QueryRow(ctx, inspectMetadataStatement).Scan(
|
|
||||||
&metadata.Version,
|
|
||||||
&metadata.AvailableExtensions,
|
|
||||||
)
|
|
||||||
if err != nil {
|
|
||||||
// 不返回部分结果,也不把查询失败转换为“成功观察到空列表”。
|
|
||||||
return application.DatabaseMetadata{}, safeError(err, application.ErrObservation)
|
|
||||||
}
|
|
||||||
return metadata, nil
|
|
||||||
}
|
|
||||||
@@ -1,115 +0,0 @@
|
|||||||
//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 (
|
|
||||||
"errors"
|
|
||||||
"testing"
|
|
||||||
|
|
||||||
corev1 "k8s.io/api/core/v1"
|
|
||||||
|
|
||||||
"git.ddupan.top/panxiao81/ayatori/internal/database/application"
|
|
||||||
"git.ddupan.top/panxiao81/ayatori/internal/database/domain/instance"
|
|
||||||
)
|
|
||||||
|
|
||||||
func TestMetadataObservesAvailableExtensionsWithoutInstalling(t *testing.T) {
|
|
||||||
f := newCredentialFixture(t)
|
|
||||||
f.createSecret(t, controllerNamespace)
|
|
||||||
if installed := f.queryPostgres(t, "SELECT count(*) FROM pg_catalog.pg_extension WHERE extname = 'hstore'"); installed != "0" {
|
|
||||||
t.Fatal("fixture unexpectedly has hstore installed")
|
|
||||||
}
|
|
||||||
// 提供同名遮蔽对象,验证 adapter 不依赖管理账号可修改的 search_path。
|
|
||||||
f.queryPostgres(t, "CREATE VIEW public.pg_available_extensions AS SELECT 'fake_extension'::name AS name")
|
|
||||||
f.queryPostgres(t, "ALTER ROLE postgres SET search_path = public, pg_catalog")
|
|
||||||
observed, err := f.service.ObserveMetadata(f.ctx, f.target)
|
|
||||||
if err != nil {
|
|
||||||
t.Fatal(err)
|
|
||||||
}
|
|
||||||
if !observed.Target().Matches(f.target) || observed.Version() == "" {
|
|
||||||
t.Fatal("metadata was not bound to the current target")
|
|
||||||
}
|
|
||||||
requested := instance.NewExtensionSet([]string{"hstore", fixtureExtension})
|
|
||||||
if observed.Extensions().Check(requested).Decision != instance.ExtensionsAccepted {
|
|
||||||
t.Fatal("available but uninstalled extension was omitted")
|
|
||||||
}
|
|
||||||
unsupported := instance.NewExtensionSet([]string{"fake_extension", "HSTORE"})
|
|
||||||
check := observed.Extensions().Check(unsupported)
|
|
||||||
if check.Decision != instance.ExtensionsUnsupported || len(check.Unsupported) != 2 {
|
|
||||||
t.Fatal("metadata accepted shadowed or case-normalized extension names")
|
|
||||||
}
|
|
||||||
if installed := f.queryPostgres(t, "SELECT count(*) FROM pg_catalog.pg_extension WHERE extname = 'hstore'"); installed != "0" {
|
|
||||||
t.Fatal("metadata observation installed an extension")
|
|
||||||
}
|
|
||||||
if schemas := f.queryPostgres(t, "SELECT count(*) FROM pg_catalog.pg_namespace WHERE nspname = 'postgresql_tenant_operator'"); schemas != "0" {
|
|
||||||
t.Fatal("metadata observation initialized the registry")
|
|
||||||
}
|
|
||||||
|
|
||||||
aggregate, err := instance.Reconstitute(f.target, instance.Snapshot{}, false)
|
|
||||||
if err != nil {
|
|
||||||
t.Fatal(err)
|
|
||||||
}
|
|
||||||
if err := aggregate.ObserveExtensions(observed.Target(), observed.Extensions()); err != nil {
|
|
||||||
t.Fatal(err)
|
|
||||||
}
|
|
||||||
if aggregate.CheckExtensions(requested).Decision != instance.ExtensionsAccepted {
|
|
||||||
t.Fatal("domain rejected observed extension availability")
|
|
||||||
}
|
|
||||||
if err := aggregate.RequireProvisioningReady(); err == nil {
|
|
||||||
t.Fatal("extension availability incorrectly authorized provisioning")
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
func TestMetadataPermissionFailureAndRecovery(t *testing.T) {
|
|
||||||
f := newCredentialFixture(t)
|
|
||||||
f.createSecret(t, controllerNamespace)
|
|
||||||
// 低权限账号也能读取可用列表;这不能证明具备 role/database/extension 管理权限。
|
|
||||||
f.queryPostgres(t, "CREATE ROLE metadata_reader LOGIN PASSWORD '"+fixturePassword+"'")
|
|
||||||
f.updateSecret(t, func(secret *corev1.Secret) {
|
|
||||||
secret.Data["login"] = []byte("metadata_reader")
|
|
||||||
})
|
|
||||||
if flags := f.queryPostgres(t, "SELECT rolsuper, rolcreaterole, rolcreatedb FROM pg_catalog.pg_roles WHERE rolname = 'metadata_reader'"); flags != "f|f|f" {
|
|
||||||
t.Fatal("metadata reader unexpectedly has management privileges")
|
|
||||||
}
|
|
||||||
requested := instance.NewExtensionSet([]string{fixtureExtension})
|
|
||||||
observed, err := f.service.ObserveMetadata(f.ctx, f.target)
|
|
||||||
if err != nil || observed.Extensions().Check(requested).Decision != instance.ExtensionsAccepted {
|
|
||||||
t.Fatalf("read-only account could not observe metadata: %v", err)
|
|
||||||
}
|
|
||||||
|
|
||||||
// 仅操作本测试独占容器的 catalog ACL;失败不能转换成“不支持任何扩展”。
|
|
||||||
f.queryPostgres(t, "REVOKE SELECT ON pg_catalog.pg_available_extensions FROM PUBLIC")
|
|
||||||
failed, err := f.service.ObserveMetadata(f.ctx, f.target)
|
|
||||||
if !errors.Is(err, application.ErrObservation) {
|
|
||||||
t.Fatalf("metadata permission failure was not reported safely: %v", err)
|
|
||||||
}
|
|
||||||
if failed.Version() != "" || failed.Target().Validate() == nil {
|
|
||||||
t.Fatal("permission failure returned partial metadata")
|
|
||||||
}
|
|
||||||
if failed.Extensions().Check(requested).Decision != instance.ExtensionSupportUnobserved {
|
|
||||||
t.Fatal("permission failure returned an observed empty set")
|
|
||||||
}
|
|
||||||
if f.backendIDs(t) != "" {
|
|
||||||
t.Fatal("failed metadata connection was retained")
|
|
||||||
}
|
|
||||||
f.queryPostgres(t, "GRANT SELECT ON pg_catalog.pg_available_extensions TO PUBLIC")
|
|
||||||
recovered, err := f.service.ObserveMetadata(f.ctx, f.target)
|
|
||||||
if err != nil || recovered.Extensions().Check(requested).Decision != instance.ExtensionsAccepted {
|
|
||||||
t.Fatalf("metadata observation did not recover: %v", err)
|
|
||||||
}
|
|
||||||
}
|
|
||||||
@@ -107,7 +107,7 @@ func TestPostgreSQLTLSHostIdentity(t *testing.T) {
|
|||||||
if err != nil {
|
if err != nil {
|
||||||
t.Fatal("trusted DNS SAN connection failed", err)
|
t.Fatal("trusted DNS SAN connection failed", err)
|
||||||
}
|
}
|
||||||
if metadata, err := db.InspectMetadata(ctx); err != nil || metadata.Version == "" {
|
if version, err := db.Version(ctx); err != nil || version == "" {
|
||||||
db.Close()
|
db.Close()
|
||||||
t.Fatal("TLS metadata read failed", err)
|
t.Fatal("TLS metadata read failed", err)
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -33,9 +33,9 @@ var (
|
|||||||
)
|
)
|
||||||
|
|
||||||
// Database 与 Connector 复用原项目 internal/instance/service.go 的能力边界。
|
// Database 与 Connector 复用原项目 internal/instance/service.go 的能力边界。
|
||||||
// Metadata 只查询版本与可用扩展,不能产生领域 Ready。
|
// 版本查询只是本切片的连通性观察,不能产生领域 Ready。
|
||||||
type Database interface {
|
type Database interface {
|
||||||
InspectMetadata(context.Context) (DatabaseMetadata, error)
|
Version(context.Context) (string, error)
|
||||||
Close()
|
Close()
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -74,25 +74,19 @@ func NewInstanceService(source CredentialReader, connector Connector) (*Instance
|
|||||||
func (s *InstanceService) String() string { return "[redacted instance service]" }
|
func (s *InstanceService) String() string { return "[redacted instance service]" }
|
||||||
func (s *InstanceService) GoString() string { return s.String() }
|
func (s *InstanceService) GoString() string { return s.String() }
|
||||||
|
|
||||||
// ObserveVersion 是完整 metadata 读取的便捷入口,不再维护另一条连接或查询路径。
|
// ObserveVersion 返回当前目标和凭据下的版本;任何失败均返回空结果。
|
||||||
func (s *InstanceService) ObserveVersion(ctx context.Context, target instance.ObservationTarget) (string, error) {
|
|
||||||
observation, err := s.ObserveMetadata(ctx, target)
|
|
||||||
return observation.Version(), err
|
|
||||||
}
|
|
||||||
|
|
||||||
// ObserveMetadata 返回当前目标和凭据下的版本与扩展;任何失败均丢弃全部结果。
|
|
||||||
// 调用者仍需使用 CR resourceVersion 保存前提防止 spec 并发修改;本方法不建立跨系统事务。
|
// 调用者仍需使用 CR resourceVersion 保存前提防止 spec 并发修改;本方法不建立跨系统事务。
|
||||||
func (s *InstanceService) ObserveMetadata(ctx context.Context, target instance.ObservationTarget) (MetadataObservation, error) {
|
func (s *InstanceService) ObserveVersion(ctx context.Context, target instance.ObservationTarget) (string, error) {
|
||||||
if err := target.Validate(); err != nil {
|
if err := target.Validate(); err != nil {
|
||||||
return MetadataObservation{}, err
|
return "", err
|
||||||
}
|
}
|
||||||
s.mu.Lock()
|
s.mu.Lock()
|
||||||
defer s.mu.Unlock()
|
defer s.mu.Unlock()
|
||||||
if s.closed {
|
if s.closed {
|
||||||
return MetadataObservation{}, ErrClosed
|
return "", ErrClosed
|
||||||
}
|
}
|
||||||
if err := ctx.Err(); err != nil {
|
if err := ctx.Err(); err != nil {
|
||||||
return MetadataObservation{}, err
|
return "", err
|
||||||
}
|
}
|
||||||
|
|
||||||
// 先读取有效凭据。读取失败时不得继续使用缓存中的旧连接。
|
// 先读取有效凭据。读取失败时不得继续使用缓存中的旧连接。
|
||||||
@@ -100,11 +94,11 @@ func (s *InstanceService) ObserveMetadata(ctx context.Context, target instance.O
|
|||||||
credentials, err := s.source.Read(ctx, target.Definition().AdminCredential())
|
credentials, err := s.source.Read(ctx, target.Definition().AdminCredential())
|
||||||
if err != nil {
|
if err != nil {
|
||||||
s.release(name)
|
s.release(name)
|
||||||
return MetadataObservation{}, credentialError(err)
|
return "", credentialError(err)
|
||||||
}
|
}
|
||||||
if credentials.username == "" || credentials.password == "" {
|
if credentials.username == "" || credentials.password == "" {
|
||||||
s.release(name)
|
s.release(name)
|
||||||
return MetadataObservation{}, ErrCredentialsInvalid
|
return "", ErrCredentialsInvalid
|
||||||
}
|
}
|
||||||
|
|
||||||
// 连接身份与有效值均未变化时复用 pgxpool;generation 本身不要求换池。
|
// 连接身份与有效值均未变化时复用 pgxpool;generation 本身不要求换池。
|
||||||
@@ -117,7 +111,7 @@ func (s *InstanceService) ObserveMetadata(ctx context.Context, target instance.O
|
|||||||
if current == nil {
|
if current == nil {
|
||||||
database, err := s.connector.Connect(ctx, target.Definition().Endpoint(), credentials)
|
database, err := s.connector.Connect(ctx, target.Definition().Endpoint(), credentials)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return MetadataObservation{}, err
|
return "", err
|
||||||
}
|
}
|
||||||
current = &entry{
|
current = &entry{
|
||||||
target: target,
|
target: target,
|
||||||
@@ -127,31 +121,23 @@ func (s *InstanceService) ObserveMetadata(ctx context.Context, target instance.O
|
|||||||
s.entries[name] = current
|
s.entries[name] = current
|
||||||
}
|
}
|
||||||
|
|
||||||
metadata, err := current.database.InspectMetadata(ctx)
|
version, err := current.database.Version(ctx)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
s.release(name)
|
s.release(name)
|
||||||
return MetadataObservation{}, err
|
return "", err
|
||||||
}
|
|
||||||
if metadata.Version == "" {
|
|
||||||
s.release(name)
|
|
||||||
return MetadataObservation{}, ErrObservation
|
|
||||||
}
|
}
|
||||||
|
|
||||||
// 回读后再检查凭据,避免把轮换前取得的结果交给新凭据的调用链。
|
// 回读后再检查凭据,避免把轮换前取得的结果交给新凭据的调用链。
|
||||||
latest, err := s.source.Read(ctx, target.Definition().AdminCredential())
|
latest, err := s.source.Read(ctx, target.Definition().AdminCredential())
|
||||||
if err != nil {
|
if err != nil {
|
||||||
s.release(name)
|
s.release(name)
|
||||||
return MetadataObservation{}, credentialError(err)
|
return "", credentialError(err)
|
||||||
}
|
}
|
||||||
if latest != credentials {
|
if latest != credentials {
|
||||||
s.release(name)
|
s.release(name)
|
||||||
return MetadataObservation{}, ErrCredentialsChanged
|
return "", ErrCredentialsChanged
|
||||||
}
|
}
|
||||||
return MetadataObservation{
|
return version, nil
|
||||||
target: target,
|
|
||||||
version: metadata.Version,
|
|
||||||
extensions: instance.ObserveExtensionSupport(metadata.AvailableExtensions),
|
|
||||||
}, nil
|
|
||||||
}
|
}
|
||||||
|
|
||||||
func credentialError(err error) error {
|
func credentialError(err error) error {
|
||||||
|
|||||||
@@ -24,8 +24,6 @@ import (
|
|||||||
"git.ddupan.top/panxiao81/ayatori/internal/database/domain/instance"
|
"git.ddupan.top/panxiao81/ayatori/internal/database/domain/instance"
|
||||||
)
|
)
|
||||||
|
|
||||||
const serviceTestPassword = "test-only"
|
|
||||||
|
|
||||||
// 延续源项目 Service 测试,用于穷举身份与装配失败;真实行为由 adapter 集成测试验证。
|
// 延续源项目 Service 测试,用于穷举身份与装配失败;真实行为由 adapter 集成测试验证。
|
||||||
type sourceStub struct {
|
type sourceStub struct {
|
||||||
credentials Credentials
|
credentials Credentials
|
||||||
@@ -37,14 +35,11 @@ func (s *sourceStub) Read(context.Context, instance.CredentialReference) (Creden
|
|||||||
}
|
}
|
||||||
|
|
||||||
type databaseStub struct {
|
type databaseStub struct {
|
||||||
closes int
|
closes int
|
||||||
err error
|
err error
|
||||||
metadata DatabaseMetadata
|
|
||||||
}
|
}
|
||||||
|
|
||||||
func (d *databaseStub) InspectMetadata(context.Context) (DatabaseMetadata, error) {
|
func (d *databaseStub) Version(context.Context) (string, error) { return "17", d.err }
|
||||||
return d.metadata, d.err
|
|
||||||
}
|
|
||||||
func (d *databaseStub) Close() {
|
func (d *databaseStub) Close() {
|
||||||
d.closes++
|
d.closes++
|
||||||
}
|
}
|
||||||
@@ -58,12 +53,7 @@ func (c *connectorStub) Connect(context.Context, instance.Endpoint, Credentials)
|
|||||||
if c.err != nil {
|
if c.err != nil {
|
||||||
return nil, c.err
|
return nil, c.err
|
||||||
}
|
}
|
||||||
db := &databaseStub{
|
db := &databaseStub{}
|
||||||
metadata: DatabaseMetadata{
|
|
||||||
Version: "17",
|
|
||||||
AvailableExtensions: []string{"plpgsql"},
|
|
||||||
},
|
|
||||||
}
|
|
||||||
c.databases = append(c.databases, db)
|
c.databases = append(c.databases, db)
|
||||||
return db, nil
|
return db, nil
|
||||||
}
|
}
|
||||||
@@ -108,7 +98,7 @@ func serviceTarget(t *testing.T, uid, host, secret string, generation int64) ins
|
|||||||
}
|
}
|
||||||
|
|
||||||
func TestInstanceConnectionIdentity(t *testing.T) {
|
func TestInstanceConnectionIdentity(t *testing.T) {
|
||||||
source := &sourceStub{credentials: Credentials{username: testUsername, password: serviceTestPassword}}
|
source := &sourceStub{credentials: Credentials{username: testUsername, password: "test-only"}}
|
||||||
connector := &connectorStub{}
|
connector := &connectorStub{}
|
||||||
service, err := NewInstanceService(source, connector)
|
service, err := NewInstanceService(source, connector)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
@@ -128,13 +118,9 @@ func TestInstanceConnectionIdentity(t *testing.T) {
|
|||||||
{"same name with new UID", serviceTarget(t, "uid-2", "second", "replacement", 1), 4},
|
{"same name with new UID", serviceTarget(t, "uid-2", "second", "replacement", 1), 4},
|
||||||
}
|
}
|
||||||
for _, testCase := range cases {
|
for _, testCase := range cases {
|
||||||
observation, err := service.ObserveMetadata(ctx, testCase.target)
|
if _, err := service.ObserveVersion(ctx, testCase.target); err != nil {
|
||||||
if err != nil {
|
|
||||||
t.Fatal(err)
|
t.Fatal(err)
|
||||||
}
|
}
|
||||||
if !observation.Target().Matches(testCase.target) {
|
|
||||||
t.Fatalf("%s: observation was bound to a previous target", testCase.name)
|
|
||||||
}
|
|
||||||
if len(connector.databases) != testCase.wantConnections {
|
if len(connector.databases) != testCase.wantConnections {
|
||||||
t.Fatalf("%s: got %d connections, want %d", testCase.name, len(connector.databases), testCase.wantConnections)
|
t.Fatalf("%s: got %d connections, want %d", testCase.name, len(connector.databases), testCase.wantConnections)
|
||||||
}
|
}
|
||||||
@@ -151,64 +137,11 @@ func TestInstanceConnectionIdentity(t *testing.T) {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
func TestMetadataObservationFreshnessAndFailure(t *testing.T) {
|
|
||||||
source := &sourceStub{credentials: Credentials{username: testUsername, password: serviceTestPassword}}
|
|
||||||
connector := &connectorStub{}
|
|
||||||
service, err := NewInstanceService(source, connector)
|
|
||||||
if err != nil {
|
|
||||||
t.Fatal(err)
|
|
||||||
}
|
|
||||||
defer service.Close()
|
|
||||||
ctx := context.Background()
|
|
||||||
target := serviceTarget(t, "metadata-uid", "first", "admin", 1)
|
|
||||||
observed, err := service.ObserveMetadata(ctx, target)
|
|
||||||
if err != nil || observed.Version() != "17" || !observed.Target().Matches(target) {
|
|
||||||
t.Fatalf("metadata observation: %v", err)
|
|
||||||
}
|
|
||||||
requested := instance.NewExtensionSet([]string{"plpgsql"})
|
|
||||||
if observed.Extensions().Check(requested).Decision != instance.ExtensionsAccepted {
|
|
||||||
t.Fatal("extension list was not observed")
|
|
||||||
}
|
|
||||||
|
|
||||||
// 连接可以复用,但每轮必须重新查询;旧观察还必须与 adapter 的可变 slice 脱离。
|
|
||||||
database := connector.databases[0]
|
|
||||||
database.metadata.AvailableExtensions[0] = "replacement"
|
|
||||||
if observed.Extensions().Check(requested).Decision != instance.ExtensionsAccepted {
|
|
||||||
t.Fatal("adapter mutation changed a completed observation")
|
|
||||||
}
|
|
||||||
refreshed, err := service.ObserveMetadata(ctx, target)
|
|
||||||
if err != nil || refreshed.Extensions().Check(requested).Decision != instance.ExtensionsUnsupported {
|
|
||||||
t.Fatalf("extension list was cached across observations: %v", err)
|
|
||||||
}
|
|
||||||
database.metadata.AvailableExtensions = nil
|
|
||||||
empty, err := service.ObserveMetadata(ctx, target)
|
|
||||||
if err != nil || empty.Extensions().Check(requested).Decision != instance.ExtensionsUnsupported {
|
|
||||||
t.Fatalf("successful empty list was treated as unobserved: %v", err)
|
|
||||||
}
|
|
||||||
|
|
||||||
// 即使 adapter 附带部分数据,错误仍使整个观察失效。
|
|
||||||
database.err = ErrObservation
|
|
||||||
failed, err := service.ObserveMetadata(ctx, target)
|
|
||||||
if !errors.Is(err, ErrObservation) || failed.Version() != "" || failed.Target().Validate() == nil {
|
|
||||||
t.Fatal("failed query returned a bound observation")
|
|
||||||
}
|
|
||||||
if failed.Extensions().Check(requested).Decision != instance.ExtensionSupportUnobserved {
|
|
||||||
t.Fatal("failed query was interpreted as an empty extension list")
|
|
||||||
}
|
|
||||||
if _, err := service.ObserveMetadata(ctx, target); err != nil {
|
|
||||||
t.Fatalf("retry after query failure: %v", err)
|
|
||||||
}
|
|
||||||
connector.databases[1].metadata.Version = ""
|
|
||||||
if _, err := service.ObserveMetadata(ctx, target); !errors.Is(err, ErrObservation) {
|
|
||||||
t.Fatal("missing server version was accepted as complete metadata")
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
func TestInstanceAssemblyFailureRecovery(t *testing.T) {
|
func TestInstanceAssemblyFailureRecovery(t *testing.T) {
|
||||||
ctx := context.Background()
|
ctx := context.Background()
|
||||||
target := serviceTarget(t, "uid-1", "first", "admin", 1)
|
target := serviceTarget(t, "uid-1", "first", "admin", 1)
|
||||||
source := &sourceStub{
|
source := &sourceStub{
|
||||||
credentials: Credentials{username: testUsername, password: serviceTestPassword},
|
credentials: Credentials{username: testUsername, password: "test-only"},
|
||||||
err: errors.New("unsafe source error"),
|
err: errors.New("unsafe source error"),
|
||||||
}
|
}
|
||||||
connector := &connectorStub{err: ErrConnection}
|
connector := &connectorStub{err: ErrConnection}
|
||||||
|
|||||||
@@ -1,40 +0,0 @@
|
|||||||
/*
|
|
||||||
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 "git.ddupan.top/panxiao81/ayatori/internal/database/domain/instance"
|
|
||||||
|
|
||||||
// DatabaseMetadata 是一次只读查询的事实,不包含管理权限或 registry 就绪结论。
|
|
||||||
// AvailableExtensions 是服务器提供的可用列表,不是已安装列表或安装授权。
|
|
||||||
type DatabaseMetadata struct {
|
|
||||||
Version string
|
|
||||||
AvailableExtensions []string
|
|
||||||
}
|
|
||||||
|
|
||||||
// MetadataObservation 只在查询成功且有效凭据再次核对一致后产生。
|
|
||||||
// target 绑定本次调用,而非连接最初创建时的 generation;零值表示没有观察。
|
|
||||||
type MetadataObservation struct {
|
|
||||||
target instance.ObservationTarget
|
|
||||||
version string
|
|
||||||
extensions instance.ExtensionSupport
|
|
||||||
}
|
|
||||||
|
|
||||||
func (o MetadataObservation) Target() instance.ObservationTarget { return o.target }
|
|
||||||
func (o MetadataObservation) Version() string { return o.version }
|
|
||||||
func (o MetadataObservation) Extensions() instance.ExtensionSupport {
|
|
||||||
return o.extensions
|
|
||||||
}
|
|
||||||
Reference in New Issue
Block a user