feat: 接入 Instance 版本与扩展可用性观测
Verify / test (pull_request) Successful in 6m1s
Verify / lint (pull_request) Successful in 7m19s
Verify / database-integration (pull_request) Successful in 5m11s

This commit is contained in:
2026-09-24 15:14:43 +00:00
parent 6acba0ca46
commit 9441b568da
11 changed files with 347 additions and 35 deletions
@@ -47,14 +47,6 @@ 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
@@ -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) {
t.Fatal("in-flight rotation was not detected")
}
if version != "" {
if observation.Version() != "" {
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) != "" {
t.Fatal("stale connection was retained after rotation")
}
@@ -42,6 +42,7 @@ import (
const (
fixtureHost = "fixture.invalid"
fixtureUser = "postgres"
fixtureExtension = "plpgsql"
dockerExec = "exec"
fixtureImage = "postgres@sha256:18cfe3ef5e6815560c98237d6216d1e5119702fb0f3894c8785dd58b8bbe5d73"
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 {
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()
t.Fatal("TLS metadata read failed", err)
}
@@ -33,9 +33,9 @@ var (
)
// Database 与 Connector 复用原项目 internal/instance/service.go 的能力边界。
// 版本查询只是本切片的连通性观察,不能产生领域 Ready。
// Metadata 只查询版本与可用扩展,不能产生领域 Ready。
type Database interface {
Version(context.Context) (string, error)
InspectMetadata(context.Context) (DatabaseMetadata, error)
Close()
}
@@ -74,19 +74,25 @@ func NewInstanceService(source CredentialReader, connector Connector) (*Instance
func (s *InstanceService) String() string { return "[redacted instance service]" }
func (s *InstanceService) GoString() string { return s.String() }
// ObserveVersion 返回当前目标和凭据下的版本;任何失败均返回空结果。
// 调用者仍需使用 CR resourceVersion 保存前提防止 spec 并发修改;本方法不建立跨系统事务。
// ObserveVersion 是完整 metadata 读取的便捷入口,不再维护另一条连接或查询路径。
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 {
return "", err
return MetadataObservation{}, err
}
s.mu.Lock()
defer s.mu.Unlock()
if s.closed {
return "", ErrClosed
return MetadataObservation{}, ErrClosed
}
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())
if err != nil {
s.release(name)
return "", credentialError(err)
return MetadataObservation{}, credentialError(err)
}
if credentials.username == "" || credentials.password == "" {
s.release(name)
return "", ErrCredentialsInvalid
return MetadataObservation{}, ErrCredentialsInvalid
}
// 连接身份与有效值均未变化时复用 pgxpool;generation 本身不要求换池。
@@ -111,7 +117,7 @@ func (s *InstanceService) ObserveVersion(ctx context.Context, target instance.Ob
if current == nil {
database, err := s.connector.Connect(ctx, target.Definition().Endpoint(), credentials)
if err != nil {
return "", err
return MetadataObservation{}, err
}
current = &entry{
target: target,
@@ -121,23 +127,31 @@ func (s *InstanceService) ObserveVersion(ctx context.Context, target instance.Ob
s.entries[name] = current
}
version, err := current.database.Version(ctx)
metadata, err := current.database.InspectMetadata(ctx)
if err != nil {
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())
if err != nil {
s.release(name)
return "", credentialError(err)
return MetadataObservation{}, credentialError(err)
}
if latest != credentials {
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 {
@@ -24,6 +24,8 @@ import (
"git.ddupan.top/panxiao81/ayatori/internal/database/domain/instance"
)
const serviceTestPassword = "test-only"
// 延续源项目 Service 测试,用于穷举身份与装配失败;真实行为由 adapter 集成测试验证。
type sourceStub struct {
credentials Credentials
@@ -35,11 +37,14 @@ func (s *sourceStub) Read(context.Context, instance.CredentialReference) (Creden
}
type databaseStub struct {
closes int
err error
closes int
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() {
d.closes++
}
@@ -53,7 +58,12 @@ func (c *connectorStub) Connect(context.Context, instance.Endpoint, Credentials)
if c.err != nil {
return nil, c.err
}
db := &databaseStub{}
db := &databaseStub{
metadata: DatabaseMetadata{
Version: "17",
AvailableExtensions: []string{"plpgsql"},
},
}
c.databases = append(c.databases, db)
return db, nil
}
@@ -98,7 +108,7 @@ func serviceTarget(t *testing.T, uid, host, secret string, generation int64) ins
}
func TestInstanceConnectionIdentity(t *testing.T) {
source := &sourceStub{credentials: Credentials{username: testUsername, password: "test-only"}}
source := &sourceStub{credentials: Credentials{username: testUsername, password: serviceTestPassword}}
connector := &connectorStub{}
service, err := NewInstanceService(source, connector)
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},
}
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)
}
if !observation.Target().Matches(testCase.target) {
t.Fatalf("%s: observation was bound to a previous target", testCase.name)
}
if 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) {
ctx := context.Background()
target := serviceTarget(t, "uid-1", "first", "admin", 1)
source := &sourceStub{
credentials: Credentials{username: testUsername, password: "test-only"},
credentials: Credentials{username: testUsername, password: serviceTestPassword},
err: errors.New("unsafe source error"),
}
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
}