Author SHA1 Message Date
panxiao81 9441b568da feat: 接入 Instance 版本与扩展可用性观测
Verify / test (pull_request) Successful in 6m1s
Verify / lint (pull_request) Successful in 7m19s
Verify / database-integration (pull_request) Successful in 5m11s
2026-09-24 15:14:43 +00:00
11 changed files with 347 additions and 35 deletions
+27 -1
View File
@@ -46,7 +46,8 @@ Secret metadata 和无关字段变化不重建连接。观测后再次读取 Sec
不把旧连接的成功作为新凭据有效的证据;这不构成跨 Kubernetes/PostgreSQL 的原子事务。 不把旧连接的成功作为新凭据有效的证据;这不构成跨 Kubernetes/PostgreSQL 的原子事务。
Instance UID、endpoint 或凭据引用变化也会释放旧连接;Forget/Close 只释放本地资源。 Instance UID、endpoint 或凭据引用变化也会释放旧连接;Forget/Close 只释放本地资源。
当前只有 `ObserveVersion` 版本查询用例,不能产生完整 CapabilityObservation 或 Ready。 当前通过 `ObserveMetadata` 读取服务器版本和可用扩展,`ObserveVersion` 只是其版本读取便捷入口,
不能产生完整 CapabilityObservation 或 Ready。
controller 接入、Secret watch、finalizer、registry 观测装配与真实权限检查仍待后续切片;并发 CR 更新 controller 接入、Secret watch、finalizer、registry 观测装配与真实权限检查仍待后续切片;并发 CR 更新
必须由调用者通过 resourceVersion 校验。应用层沿用源实现的串行处理,本阶段未引入新的调度框架。 必须由调用者通过 resourceVersion 校验。应用层沿用源实现的串行处理,本阶段未引入新的调度框架。
@@ -58,6 +59,31 @@ namespace 边界、有效值轮换、metadata 无关变化、中途轮换、重
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`,
+6 -1
View File
@@ -3,7 +3,7 @@
> 本页迁入作为 Database 模块的测试分层与 fixture 合同。旧项目的 Make target、devcontainer > 本页迁入作为 Database 模块的测试分层与 fixture 合同。旧项目的 Make target、devcontainer
> 和脚手架版本尚未适配 Ayatori;实现时应复用 Ayatori 现有工具链,并保持这里定义的测试边界。 > 和脚手架版本尚未适配 Ayatori;实现时应复用 Ayatori 现有工具链,并保持这里定义的测试边界。
## Ayatori 已接入的凭据与 registry 切片测试 ## Ayatori 已接入的凭据、metadata 与 registry 切片测试
本节命令已在 Ayatori 接入;以下历史 Compose/Kind 操作仍属于迁入的目标合同。 本节命令已在 Ayatori 接入;以下历史 Compose/Kind 操作仍属于迁入的目标合同。
@@ -28,6 +28,11 @@ runner 后端修复上线,不能从本地测试通过推断远端已经可用
fixture 启动失败会保留退出错误与 stderr,并遮蔽测试密码, fixture 启动失败会保留退出错误与 stderr,并遮蔽测试密码,
以区分缺少命令、daemon 不可达、权限和镜像拉取失败。 以区分缺少命令、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 与
超前版本拒绝。通过在真实 COMMIT 成功后注入客户端错误,验证结果不确定时的重试;这不替代 超前版本拒绝。通过在真实 COMMIT 成功后注入客户端错误,验证结果不确定时的重试;这不替代
@@ -47,14 +47,6 @@ 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,13 +131,17 @@ func TestObservationDiscardsResultWhenCredentialsChange(t *testing.T) {
} }
} }
version, err := fixture.service.ObserveVersion(fixture.ctx, fixture.target) observation, err := fixture.service.ObserveMetadata(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 version != "" { if observation.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")
} }
@@ -42,6 +42,7 @@ 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"
@@ -0,0 +1,48 @@
/*
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
}
@@ -0,0 +1,115 @@
//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 version, err := db.Version(ctx); err != nil || version == "" { if metadata, err := db.InspectMetadata(ctx); err != nil || metadata.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 的能力边界。
// 版本查询只是本切片的连通性观察,不能产生领域 Ready。 // Metadata 只查询版本与可用扩展,不能产生领域 Ready。
type Database interface { type Database interface {
Version(context.Context) (string, error) InspectMetadata(context.Context) (DatabaseMetadata, error)
Close() Close()
} }
@@ -74,19 +74,25 @@ 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 返回当前目标和凭据下的版本;任何失败均返回空结果。 // ObserveVersion 是完整 metadata 读取的便捷入口,不再维护另一条连接或查询路径。
// 调用者仍需使用 CR resourceVersion 保存前提防止 spec 并发修改;本方法不建立跨系统事务。
func (s *InstanceService) ObserveVersion(ctx context.Context, target instance.ObservationTarget) (string, error) { 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 并发修改;本方法不建立跨系统事务。
func (s *InstanceService) ObserveMetadata(ctx context.Context, target instance.ObservationTarget) (MetadataObservation, error) {
if err := target.Validate(); err != nil { if err := target.Validate(); err != nil {
return "", err return MetadataObservation{}, err
} }
s.mu.Lock() s.mu.Lock()
defer s.mu.Unlock() defer s.mu.Unlock()
if s.closed { if s.closed {
return "", ErrClosed return MetadataObservation{}, ErrClosed
} }
if err := ctx.Err(); err != nil { if err := ctx.Err(); err != nil {
return "", err return MetadataObservation{}, err
} }
// 先读取有效凭据。读取失败时不得继续使用缓存中的旧连接。 // 先读取有效凭据。读取失败时不得继续使用缓存中的旧连接。
@@ -94,11 +100,11 @@ func (s *InstanceService) ObserveVersion(ctx context.Context, target instance.Ob
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 "", credentialError(err) return MetadataObservation{}, credentialError(err)
} }
if credentials.username == "" || credentials.password == "" { if credentials.username == "" || credentials.password == "" {
s.release(name) s.release(name)
return "", ErrCredentialsInvalid return MetadataObservation{}, ErrCredentialsInvalid
} }
// 连接身份与有效值均未变化时复用 pgxpool;generation 本身不要求换池。 // 连接身份与有效值均未变化时复用 pgxpool;generation 本身不要求换池。
@@ -111,7 +117,7 @@ func (s *InstanceService) ObserveVersion(ctx context.Context, target instance.Ob
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 "", err return MetadataObservation{}, err
} }
current = &entry{ current = &entry{
target: target, target: target,
@@ -121,23 +127,31 @@ func (s *InstanceService) ObserveVersion(ctx context.Context, target instance.Ob
s.entries[name] = current s.entries[name] = current
} }
version, err := current.database.Version(ctx) metadata, err := current.database.InspectMetadata(ctx)
if err != nil { if err != nil {
s.release(name) s.release(name)
return "", err return MetadataObservation{}, 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 "", credentialError(err) return MetadataObservation{}, credentialError(err)
} }
if latest != credentials { if latest != credentials {
s.release(name) s.release(name)
return "", ErrCredentialsChanged return MetadataObservation{}, ErrCredentialsChanged
} }
return version, nil return MetadataObservation{
target: target,
version: metadata.Version,
extensions: instance.ObserveExtensionSupport(metadata.AvailableExtensions),
}, nil
} }
func credentialError(err error) error { func credentialError(err error) error {
@@ -24,6 +24,8 @@ 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
@@ -35,11 +37,14 @@ 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) Version(context.Context) (string, error) { return "17", d.err } func (d *databaseStub) InspectMetadata(context.Context) (DatabaseMetadata, error) {
return d.metadata, d.err
}
func (d *databaseStub) Close() { func (d *databaseStub) Close() {
d.closes++ d.closes++
} }
@@ -53,7 +58,12 @@ 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
} }
@@ -98,7 +108,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: "test-only"}} source := &sourceStub{credentials: Credentials{username: testUsername, password: serviceTestPassword}}
connector := &connectorStub{} connector := &connectorStub{}
service, err := NewInstanceService(source, connector) service, err := NewInstanceService(source, connector)
if err != nil { if err != nil {
@@ -118,9 +128,13 @@ 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 {
if _, err := service.ObserveVersion(ctx, testCase.target); err != nil { observation, err := service.ObserveMetadata(ctx, testCase.target)
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)
} }
@@ -137,11 +151,64 @@ 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: "test-only"}, credentials: Credentials{username: testUsername, password: serviceTestPassword},
err: errors.New("unsafe source error"), err: errors.New("unsafe source error"),
} }
connector := &connectorStub{err: ErrConnection} connector := &connectorStub{err: ErrConnection}
+40
View File
@@ -0,0 +1,40 @@
/*
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
}