refactor: 移除 Database registry 与 Instance 初始化依赖
This commit is contained in:
@@ -37,6 +37,7 @@ func TestMetadataObservesAvailableExtensionsWithoutInstalling(t *testing.T) {
|
||||
// 提供同名遮蔽对象,验证 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")
|
||||
schemasBefore := f.queryPostgres(t, "SELECT string_agg(nspname, ',' ORDER BY nspname) FROM pg_catalog.pg_namespace")
|
||||
observed, err := f.service.ObserveMetadata(f.ctx, f.target)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
@@ -56,8 +57,8 @@ func TestMetadataObservesAvailableExtensionsWithoutInstalling(t *testing.T) {
|
||||
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")
|
||||
if schemas := f.queryPostgres(t, "SELECT string_agg(nspname, ',' ORDER BY nspname) FROM pg_catalog.pg_namespace"); schemas != schemasBefore {
|
||||
t.Fatal("metadata observation changed database schemas")
|
||||
}
|
||||
|
||||
aggregate, err := instance.Reconstitute(f.target, instance.Snapshot{}, false)
|
||||
|
||||
@@ -1,25 +0,0 @@
|
||||
CREATE SCHEMA postgresql_tenant_operator;
|
||||
|
||||
CREATE TABLE postgresql_tenant_operator.tenant_ownership (
|
||||
instance_uid text NOT NULL,
|
||||
tenant_uid text NOT NULL,
|
||||
tenant_namespace text NOT NULL,
|
||||
tenant_name text NOT NULL,
|
||||
database_name text NOT NULL,
|
||||
role_name text NOT NULL,
|
||||
credential_path text NOT NULL,
|
||||
managed boolean NOT NULL DEFAULT true,
|
||||
created_at timestamptz NOT NULL DEFAULT clock_timestamp(),
|
||||
updated_at timestamptz NOT NULL DEFAULT clock_timestamp(),
|
||||
retained_at timestamptz,
|
||||
PRIMARY KEY (instance_uid, tenant_uid),
|
||||
UNIQUE (instance_uid, tenant_namespace, tenant_name),
|
||||
UNIQUE (instance_uid, database_name),
|
||||
UNIQUE (instance_uid, role_name),
|
||||
UNIQUE (credential_path),
|
||||
CHECK (managed OR retained_at IS NOT NULL)
|
||||
);
|
||||
|
||||
---- create above / drop below ----
|
||||
|
||||
DROP SCHEMA postgresql_tenant_operator CASCADE;
|
||||
@@ -1,326 +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 registry persists controller ownership evidence in the PostgreSQL
|
||||
// management database. It deliberately does not persist reconciliation phases;
|
||||
// those belong to the Kubernetes resource status.
|
||||
package registry
|
||||
|
||||
import (
|
||||
"context"
|
||||
"embed"
|
||||
"errors"
|
||||
"fmt"
|
||||
"io/fs"
|
||||
"time"
|
||||
|
||||
"github.com/jackc/pgx/v5"
|
||||
"github.com/jackc/pgx/v5/pgxpool"
|
||||
"github.com/jackc/tern/v2/migrate"
|
||||
)
|
||||
|
||||
const migrationVersionTable = "public.postgresql_tenant_operator_schema_version"
|
||||
|
||||
//go:embed migrations/*.sql
|
||||
var migrationFiles embed.FS
|
||||
|
||||
var (
|
||||
// ErrNotFound indicates that the registry has no matching ownership record.
|
||||
ErrNotFound = errors.New("registry ownership record not found")
|
||||
// ErrConflict indicates that a requested name or path belongs to another Tenant UID.
|
||||
ErrConflict = errors.New("registry ownership conflict")
|
||||
)
|
||||
|
||||
// Beginner is implemented by pgx.Conn and pgxpool.Pool.
|
||||
type Beginner interface {
|
||||
Begin(context.Context) (pgx.Tx, error)
|
||||
}
|
||||
|
||||
// Store manages ownership records in one PostgreSQLInstance management database.
|
||||
type Store struct {
|
||||
db Beginner
|
||||
}
|
||||
|
||||
// NewStore creates a registry store backed by a PostgreSQL connection or pool.
|
||||
func NewStore(db Beginner) *Store {
|
||||
return &Store{db: db}
|
||||
}
|
||||
|
||||
// Ownership identifies every external resource reserved for one Tenant UID.
|
||||
type Ownership struct {
|
||||
InstanceUID string
|
||||
TenantUID string
|
||||
TenantNamespace string
|
||||
TenantName string
|
||||
DatabaseName string
|
||||
RoleName string
|
||||
CredentialPath string
|
||||
}
|
||||
|
||||
// Record is the persisted ownership state.
|
||||
type Record struct {
|
||||
Ownership
|
||||
Managed bool
|
||||
CreatedAt time.Time
|
||||
UpdatedAt time.Time
|
||||
RetainedAt *time.Time
|
||||
}
|
||||
|
||||
// ClaimResult reports whether Claim inserted a new record or found the same claim.
|
||||
type ClaimResult string
|
||||
|
||||
const (
|
||||
ClaimCreated ClaimResult = "Created"
|
||||
ClaimOwned ClaimResult = "Owned"
|
||||
)
|
||||
|
||||
// Bootstrap applies all pending versioned registry migrations idempotently.
|
||||
func (s *Store) Bootstrap(ctx context.Context) error {
|
||||
if s == nil || s.db == nil {
|
||||
return errors.New("bootstrap registry: nil database")
|
||||
}
|
||||
|
||||
return s.withMigrationConnection(ctx, func(conn *pgx.Conn) error {
|
||||
migrations, err := fs.Sub(migrationFiles, "migrations")
|
||||
if err != nil {
|
||||
return fmt.Errorf("open embedded registry migrations: %w", err)
|
||||
}
|
||||
migrator, err := migrate.NewMigrator(ctx, conn, migrationVersionTable)
|
||||
if err != nil {
|
||||
return fmt.Errorf("initialize registry migrator: %w", err)
|
||||
}
|
||||
if err := migrator.LoadMigrations(migrations); err != nil {
|
||||
return fmt.Errorf("load registry migrations: %w", err)
|
||||
}
|
||||
if err := migrator.Migrate(ctx); err != nil {
|
||||
return fmt.Errorf("apply registry migrations: %w", err)
|
||||
}
|
||||
return nil
|
||||
})
|
||||
}
|
||||
|
||||
// Claim reserves all names for an owner. Repeating an identical claim is idempotent.
|
||||
func (s *Store) Claim(ctx context.Context, owner Ownership) (ClaimResult, error) {
|
||||
if s == nil || s.db == nil {
|
||||
return "", errors.New("claim registry ownership: nil database")
|
||||
}
|
||||
if err := owner.validate(); err != nil {
|
||||
return "", err
|
||||
}
|
||||
|
||||
tx, err := s.db.Begin(ctx)
|
||||
if err != nil {
|
||||
return "", fmt.Errorf("begin registry claim: %w", err)
|
||||
}
|
||||
defer func() { _ = tx.Rollback(ctx) }()
|
||||
|
||||
_, err = scanRecord(tx.QueryRow(ctx, claimStatement,
|
||||
owner.InstanceUID,
|
||||
owner.TenantUID,
|
||||
owner.TenantNamespace,
|
||||
owner.TenantName,
|
||||
owner.DatabaseName,
|
||||
owner.RoleName,
|
||||
owner.CredentialPath,
|
||||
))
|
||||
if err != nil && !errors.Is(err, ErrNotFound) {
|
||||
return "", fmt.Errorf("insert registry claim: %w", err)
|
||||
}
|
||||
if errors.Is(err, ErrNotFound) {
|
||||
record, err := getByTenantUID(ctx, tx, owner.InstanceUID, owner.TenantUID)
|
||||
if err != nil {
|
||||
if errors.Is(err, ErrNotFound) {
|
||||
return "", fmt.Errorf("%w: database, role, tenant identity, or credential path is already reserved", ErrConflict)
|
||||
}
|
||||
return "", err
|
||||
}
|
||||
if !record.equal(owner) || !record.Managed {
|
||||
return "", fmt.Errorf("%w: database, role, tenant identity, or credential path is already reserved", ErrConflict)
|
||||
}
|
||||
if err := tx.Commit(ctx); err != nil {
|
||||
return "", fmt.Errorf("commit registry claim: %w", err)
|
||||
}
|
||||
return ClaimOwned, nil
|
||||
}
|
||||
|
||||
if err := tx.Commit(ctx); err != nil {
|
||||
return "", fmt.Errorf("commit registry claim: %w", err)
|
||||
}
|
||||
return ClaimCreated, nil
|
||||
}
|
||||
|
||||
func (s *Store) withMigrationConnection(ctx context.Context, run func(*pgx.Conn) error) error {
|
||||
switch db := s.db.(type) {
|
||||
case *pgx.Conn:
|
||||
return run(db)
|
||||
case *pgxpool.Pool:
|
||||
conn, err := db.Acquire(ctx)
|
||||
if err != nil {
|
||||
return fmt.Errorf("acquire registry migration connection: %w", err)
|
||||
}
|
||||
defer conn.Release()
|
||||
return run(conn.Conn())
|
||||
default:
|
||||
return fmt.Errorf("bootstrap registry: database type %T cannot provide a migration connection", s.db)
|
||||
}
|
||||
}
|
||||
|
||||
// Get returns the ownership record for an Instance UID and Tenant UID.
|
||||
func (s *Store) Get(ctx context.Context, instanceUID, tenantUID string) (Record, error) {
|
||||
if s == nil || s.db == nil {
|
||||
return Record{}, errors.New("get registry record: nil database")
|
||||
}
|
||||
if instanceUID == "" || tenantUID == "" {
|
||||
return Record{}, errors.New("get registry record: instance UID and tenant UID are required")
|
||||
}
|
||||
|
||||
tx, err := s.db.Begin(ctx)
|
||||
if err != nil {
|
||||
return Record{}, fmt.Errorf("begin registry read: %w", err)
|
||||
}
|
||||
defer func() { _ = tx.Rollback(ctx) }()
|
||||
|
||||
record, err := getByTenantUID(ctx, tx, instanceUID, tenantUID)
|
||||
if err != nil {
|
||||
return Record{}, err
|
||||
}
|
||||
if err := tx.Commit(ctx); err != nil {
|
||||
return Record{}, fmt.Errorf("commit registry read: %w", err)
|
||||
}
|
||||
return record, nil
|
||||
}
|
||||
|
||||
// MarkRetained changes a matching managed ownership record into an unmanaged tombstone.
|
||||
// Repeating the operation for the same tombstone is safe.
|
||||
func (s *Store) MarkRetained(ctx context.Context, owner Ownership) error {
|
||||
return s.changeOwnership(ctx, owner, "mark registry record retained", func(ctx context.Context, tx pgx.Tx, record Record) error {
|
||||
if !record.Managed {
|
||||
return nil
|
||||
}
|
||||
tag, err := tx.Exec(ctx, markRetainedStatement, owner.InstanceUID, owner.TenantUID)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
if tag.RowsAffected() != 1 {
|
||||
return ErrConflict
|
||||
}
|
||||
return nil
|
||||
})
|
||||
}
|
||||
|
||||
// Delete removes a matching managed record after its external resources have been deleted.
|
||||
// An already absent record is treated as a successful retry; a retained record is never deleted.
|
||||
func (s *Store) Delete(ctx context.Context, owner Ownership) error {
|
||||
return s.changeOwnership(ctx, owner, "delete registry record", func(ctx context.Context, tx pgx.Tx, record Record) error {
|
||||
if !record.Managed {
|
||||
return ErrConflict
|
||||
}
|
||||
tag, err := tx.Exec(ctx, deleteStatement, owner.InstanceUID, owner.TenantUID)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
if tag.RowsAffected() != 1 {
|
||||
return ErrConflict
|
||||
}
|
||||
return nil
|
||||
})
|
||||
}
|
||||
|
||||
func (s *Store) changeOwnership(
|
||||
ctx context.Context,
|
||||
owner Ownership,
|
||||
operation string,
|
||||
change func(context.Context, pgx.Tx, Record) error,
|
||||
) error {
|
||||
if s == nil || s.db == nil {
|
||||
return fmt.Errorf("%s: nil database", operation)
|
||||
}
|
||||
if err := owner.validate(); err != nil {
|
||||
return fmt.Errorf("%s: %w", operation, err)
|
||||
}
|
||||
|
||||
tx, err := s.db.Begin(ctx)
|
||||
if err != nil {
|
||||
return fmt.Errorf("begin %s: %w", operation, err)
|
||||
}
|
||||
defer func() { _ = tx.Rollback(ctx) }()
|
||||
|
||||
record, err := getByTenantUIDForUpdate(ctx, tx, owner.InstanceUID, owner.TenantUID)
|
||||
if errors.Is(err, ErrNotFound) && operation == "delete registry record" {
|
||||
return nil
|
||||
}
|
||||
if err != nil {
|
||||
return fmt.Errorf("%s: %w", operation, err)
|
||||
}
|
||||
if !record.equal(owner) {
|
||||
return fmt.Errorf("%s: %w", operation, ErrConflict)
|
||||
}
|
||||
if err := change(ctx, tx, record); err != nil {
|
||||
return fmt.Errorf("%s: %w", operation, err)
|
||||
}
|
||||
if err := tx.Commit(ctx); err != nil {
|
||||
return fmt.Errorf("commit %s: %w", operation, err)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
func getByTenantUIDForUpdate(ctx context.Context, tx pgx.Tx, instanceUID, tenantUID string) (Record, error) {
|
||||
return scanRecord(tx.QueryRow(ctx, getByTenantUIDForUpdateStatement, instanceUID, tenantUID))
|
||||
}
|
||||
|
||||
func getByTenantUID(ctx context.Context, tx pgx.Tx, instanceUID, tenantUID string) (Record, error) {
|
||||
return scanRecord(tx.QueryRow(ctx, getByTenantUIDStatement, instanceUID, tenantUID))
|
||||
}
|
||||
|
||||
type rowScanner interface {
|
||||
Scan(...any) error
|
||||
}
|
||||
|
||||
func scanRecord(row rowScanner) (Record, error) {
|
||||
var record Record
|
||||
err := row.Scan(
|
||||
&record.InstanceUID,
|
||||
&record.TenantUID,
|
||||
&record.TenantNamespace,
|
||||
&record.TenantName,
|
||||
&record.DatabaseName,
|
||||
&record.RoleName,
|
||||
&record.CredentialPath,
|
||||
&record.Managed,
|
||||
&record.CreatedAt,
|
||||
&record.UpdatedAt,
|
||||
&record.RetainedAt,
|
||||
)
|
||||
if errors.Is(err, pgx.ErrNoRows) {
|
||||
return Record{}, ErrNotFound
|
||||
}
|
||||
if err != nil {
|
||||
return Record{}, fmt.Errorf("read registry record: %w", err)
|
||||
}
|
||||
return record, nil
|
||||
}
|
||||
|
||||
func (o Ownership) validate() error {
|
||||
if o.InstanceUID == "" || o.TenantUID == "" || o.TenantNamespace == "" || o.TenantName == "" ||
|
||||
o.DatabaseName == "" || o.RoleName == "" || o.CredentialPath == "" {
|
||||
return errors.New("claim registry ownership: all ownership fields are required")
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
func (o Ownership) equal(other Ownership) bool {
|
||||
return o == other
|
||||
}
|
||||
@@ -1,68 +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 registry
|
||||
|
||||
const claimStatement = `
|
||||
INSERT INTO postgresql_tenant_operator.tenant_ownership (
|
||||
instance_uid,
|
||||
tenant_uid,
|
||||
tenant_namespace,
|
||||
tenant_name,
|
||||
database_name,
|
||||
role_name,
|
||||
credential_path
|
||||
) VALUES ($1, $2, $3, $4, $5, $6, $7)
|
||||
ON CONFLICT DO NOTHING
|
||||
RETURNING
|
||||
instance_uid,
|
||||
tenant_uid,
|
||||
tenant_namespace,
|
||||
tenant_name,
|
||||
database_name,
|
||||
role_name,
|
||||
credential_path,
|
||||
managed,
|
||||
created_at,
|
||||
updated_at,
|
||||
retained_at`
|
||||
|
||||
const getByTenantUIDStatement = `
|
||||
SELECT
|
||||
instance_uid,
|
||||
tenant_uid,
|
||||
tenant_namespace,
|
||||
tenant_name,
|
||||
database_name,
|
||||
role_name,
|
||||
credential_path,
|
||||
managed,
|
||||
created_at,
|
||||
updated_at,
|
||||
retained_at
|
||||
FROM postgresql_tenant_operator.tenant_ownership
|
||||
WHERE instance_uid = $1 AND tenant_uid = $2`
|
||||
|
||||
const getByTenantUIDForUpdateStatement = getByTenantUIDStatement + ` FOR UPDATE`
|
||||
|
||||
const markRetainedStatement = `
|
||||
UPDATE postgresql_tenant_operator.tenant_ownership
|
||||
SET managed = false, retained_at = clock_timestamp(), updated_at = clock_timestamp()
|
||||
WHERE instance_uid = $1 AND tenant_uid = $2 AND managed = true`
|
||||
|
||||
const deleteStatement = `
|
||||
DELETE FROM postgresql_tenant_operator.tenant_ownership
|
||||
WHERE instance_uid = $1 AND tenant_uid = $2 AND managed = true`
|
||||
@@ -1,427 +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 (
|
||||
"context"
|
||||
"errors"
|
||||
"net"
|
||||
"net/url"
|
||||
"strconv"
|
||||
"sync"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"github.com/jackc/pgx/v5"
|
||||
"github.com/jackc/pgx/v5/pgxpool"
|
||||
|
||||
"git.ddupan.top/panxiao81/ayatori/internal/database/adapter/postgresql/registry"
|
||||
)
|
||||
|
||||
// registryFixture 只连接本测试创建的容器,不接受外部数据库地址。
|
||||
type registryFixture struct {
|
||||
ctx context.Context
|
||||
pool *pgxpool.Pool
|
||||
store *registry.Store
|
||||
}
|
||||
|
||||
func newRegistryFixture(t *testing.T) *registryFixture {
|
||||
t.Helper()
|
||||
ctx, cancel := context.WithTimeout(context.Background(), 90*time.Second)
|
||||
t.Cleanup(cancel)
|
||||
_, port := postgresFixture(t, ctx)
|
||||
endpoint := url.URL{
|
||||
Scheme: "postgres",
|
||||
User: url.UserPassword(fixtureUser, fixturePassword),
|
||||
Host: net.JoinHostPort("127.0.0.1", strconv.Itoa(port)),
|
||||
Path: "/postgres",
|
||||
RawQuery: "sslmode=disable",
|
||||
}
|
||||
pool, err := pgxpool.New(ctx, endpoint.String())
|
||||
if err != nil {
|
||||
t.Fatal("cannot configure registry fixture connection")
|
||||
}
|
||||
t.Cleanup(pool.Close)
|
||||
if err := pool.Ping(ctx); err != nil {
|
||||
t.Fatal("cannot connect to registry fixture")
|
||||
}
|
||||
return ®istryFixture{ctx: ctx, pool: pool, store: registry.NewStore(pool)}
|
||||
}
|
||||
|
||||
func registryOwner(suffix string) registry.Ownership {
|
||||
return registry.Ownership{
|
||||
InstanceUID: "instance-uid",
|
||||
TenantUID: "tenant-uid-" + suffix,
|
||||
TenantNamespace: "applications",
|
||||
TenantName: "tenant-" + suffix,
|
||||
DatabaseName: "database_" + suffix,
|
||||
RoleName: "role_" + suffix,
|
||||
CredentialPath: "credentials/" + suffix,
|
||||
}
|
||||
}
|
||||
|
||||
func (f *registryFixture) bootstrap(t *testing.T) {
|
||||
t.Helper()
|
||||
if err := f.store.Bootstrap(f.ctx); err != nil {
|
||||
t.Fatalf("bootstrap registry: %v", err)
|
||||
}
|
||||
}
|
||||
|
||||
func (f *registryFixture) claim(t *testing.T, owner registry.Ownership, expected registry.ClaimResult) {
|
||||
t.Helper()
|
||||
result, err := f.store.Claim(f.ctx, owner)
|
||||
if err != nil || result != expected {
|
||||
t.Fatalf("claim: result=%q, error=%v; want %q", result, err, expected)
|
||||
}
|
||||
}
|
||||
|
||||
func TestRegistryOwnershipLifecycle(t *testing.T) {
|
||||
f := newRegistryFixture(t)
|
||||
f.bootstrap(t)
|
||||
f.bootstrap(t)
|
||||
owner := registryOwner("lifecycle")
|
||||
f.claim(t, owner, registry.ClaimCreated)
|
||||
f.claim(t, owner, registry.ClaimOwned)
|
||||
record, err := f.store.Get(f.ctx, owner.InstanceUID, owner.TenantUID)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if record.Ownership != owner || !record.Managed || record.CreatedAt.IsZero() || record.RetainedAt != nil {
|
||||
t.Fatalf("unexpected ownership record: %+v", record)
|
||||
}
|
||||
|
||||
// 错误的资源归属不能删除或 Retain 原记录。
|
||||
wrongOwner := owner
|
||||
wrongOwner.RoleName = "another_role"
|
||||
if err := f.store.Delete(f.ctx, wrongOwner); !errors.Is(err, registry.ErrConflict) {
|
||||
t.Fatalf("delete mismatched owner: %v", err)
|
||||
}
|
||||
if err := f.store.MarkRetained(f.ctx, wrongOwner); !errors.Is(err, registry.ErrConflict) {
|
||||
t.Fatalf("retain mismatched owner: %v", err)
|
||||
}
|
||||
for range 2 {
|
||||
if err := f.store.Delete(f.ctx, owner); err != nil {
|
||||
t.Fatalf("delete retry: %v", err)
|
||||
}
|
||||
}
|
||||
if _, err := f.store.Get(f.ctx, owner.InstanceUID, owner.TenantUID); !errors.Is(err, registry.ErrNotFound) {
|
||||
t.Fatalf("read deleted record: %v", err)
|
||||
}
|
||||
}
|
||||
|
||||
func TestRegistryRetainedRecordCannotBeReclaimed(t *testing.T) {
|
||||
f := newRegistryFixture(t)
|
||||
f.bootstrap(t)
|
||||
owner := registryOwner("retained")
|
||||
f.claim(t, owner, registry.ClaimCreated)
|
||||
if err := f.store.MarkRetained(f.ctx, owner); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
retained, err := f.store.Get(f.ctx, owner.InstanceUID, owner.TenantUID)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if retained.Managed || retained.RetainedAt == nil {
|
||||
t.Fatalf("missing retained tombstone: %+v", retained)
|
||||
}
|
||||
if err := f.store.MarkRetained(f.ctx, owner); err != nil {
|
||||
t.Fatalf("retain retry: %v", err)
|
||||
}
|
||||
retried, err := f.store.Get(f.ctx, owner.InstanceUID, owner.TenantUID)
|
||||
if err != nil || !retried.UpdatedAt.Equal(retained.UpdatedAt) {
|
||||
t.Fatalf("retain retry changed tombstone: %+v, %v", retried, err)
|
||||
}
|
||||
if err := f.store.Delete(f.ctx, owner); !errors.Is(err, registry.ErrConflict) {
|
||||
t.Fatalf("delete retained record: %v", err)
|
||||
}
|
||||
if _, err := f.store.Claim(f.ctx, owner); !errors.Is(err, registry.ErrConflict) {
|
||||
t.Fatalf("reclaim retained record: %v", err)
|
||||
}
|
||||
owner.TenantUID = "replacement-uid"
|
||||
if _, err := f.store.Claim(f.ctx, owner); !errors.Is(err, registry.ErrConflict) {
|
||||
t.Fatalf("replacement tenant reclaimed tombstone: %v", err)
|
||||
}
|
||||
}
|
||||
|
||||
func TestRegistryUniqueReservations(t *testing.T) {
|
||||
f := newRegistryFixture(t)
|
||||
f.bootstrap(t)
|
||||
owner := registryOwner("original")
|
||||
f.claim(t, owner, registry.ClaimCreated)
|
||||
tests := []struct {
|
||||
name string
|
||||
change func(*registry.Ownership)
|
||||
}{
|
||||
{name: "tenant UID", change: func(other *registry.Ownership) { other.TenantUID = owner.TenantUID }},
|
||||
{name: "tenant name", change: func(other *registry.Ownership) { other.TenantName = owner.TenantName }},
|
||||
{name: "database", change: func(other *registry.Ownership) { other.DatabaseName = owner.DatabaseName }},
|
||||
{name: "role", change: func(other *registry.Ownership) { other.RoleName = owner.RoleName }},
|
||||
{name: "credential path", change: func(other *registry.Ownership) { other.CredentialPath = owner.CredentialPath }},
|
||||
}
|
||||
for _, tt := range tests {
|
||||
t.Run(tt.name, func(t *testing.T) {
|
||||
other := registryOwner("other")
|
||||
tt.change(&other)
|
||||
if _, err := f.store.Claim(f.ctx, other); !errors.Is(err, registry.ErrConflict) {
|
||||
t.Fatalf("conflicting reservation: %v", err)
|
||||
}
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
func TestRegistryConcurrentBootstrapAndClaims(t *testing.T) {
|
||||
f := newRegistryFixture(t)
|
||||
const workers = 8
|
||||
start := make(chan struct{})
|
||||
results := make(chan error, workers)
|
||||
var group sync.WaitGroup
|
||||
for range workers {
|
||||
group.Go(func() {
|
||||
<-start
|
||||
results <- f.store.Bootstrap(f.ctx)
|
||||
})
|
||||
}
|
||||
close(start)
|
||||
group.Wait()
|
||||
close(results)
|
||||
for err := range results {
|
||||
if err != nil {
|
||||
t.Fatalf("concurrent bootstrap: %v", err)
|
||||
}
|
||||
}
|
||||
|
||||
owner := registryOwner("concurrent")
|
||||
type claimOutcome struct {
|
||||
result registry.ClaimResult
|
||||
err error
|
||||
}
|
||||
claims := make(chan claimOutcome, workers)
|
||||
start = make(chan struct{})
|
||||
for range workers {
|
||||
group.Go(func() {
|
||||
<-start
|
||||
result, err := f.store.Claim(f.ctx, owner)
|
||||
claims <- claimOutcome{result: result, err: err}
|
||||
})
|
||||
}
|
||||
close(start)
|
||||
group.Wait()
|
||||
close(claims)
|
||||
created := 0
|
||||
for outcome := range claims {
|
||||
if outcome.err != nil {
|
||||
t.Fatal(outcome.err)
|
||||
}
|
||||
switch outcome.result {
|
||||
case registry.ClaimCreated:
|
||||
created++
|
||||
case registry.ClaimOwned:
|
||||
default:
|
||||
t.Fatalf("unexpected claim result: %q", outcome.result)
|
||||
}
|
||||
}
|
||||
if created != 1 {
|
||||
t.Fatalf("created %d records for the same owner; want 1", created)
|
||||
}
|
||||
}
|
||||
|
||||
func TestRegistryRestartAndCanceledRequestRecovery(t *testing.T) {
|
||||
f := newRegistryFixture(t)
|
||||
f.bootstrap(t)
|
||||
owner := registryOwner("restart")
|
||||
f.claim(t, owner, registry.ClaimCreated)
|
||||
config := f.pool.Config()
|
||||
f.pool.Close()
|
||||
pool, err := pgxpool.NewWithConfig(f.ctx, config)
|
||||
if err != nil {
|
||||
t.Fatal("cannot reopen fixture connection")
|
||||
}
|
||||
t.Cleanup(pool.Close)
|
||||
f.store = registry.NewStore(pool)
|
||||
f.bootstrap(t)
|
||||
// 模拟客户端丢失上次成功结果后,以新连接重试;归属证据必须来自数据库。
|
||||
f.claim(t, owner, registry.ClaimOwned)
|
||||
canceled, cancel := context.WithCancel(f.ctx)
|
||||
cancel()
|
||||
if _, err := f.store.Claim(canceled, registryOwner("canceled")); !errors.Is(err, context.Canceled) {
|
||||
t.Fatalf("canceled claim: %v", err)
|
||||
}
|
||||
f.claim(t, registryOwner("canceled"), registry.ClaimCreated)
|
||||
}
|
||||
|
||||
func TestRegistryBootstrapRejectsUnknownSchema(t *testing.T) {
|
||||
f := newRegistryFixture(t)
|
||||
// 人工创建的同名 schema 不能被初始化过程接管或覆盖。
|
||||
if _, err := f.pool.Exec(f.ctx, "CREATE SCHEMA postgresql_tenant_operator"); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if err := f.store.Bootstrap(f.ctx); err == nil {
|
||||
t.Fatal("bootstrap accepted an unmanaged schema")
|
||||
}
|
||||
var version int
|
||||
if err := f.pool.QueryRow(f.ctx, "SELECT version FROM public.postgresql_tenant_operator_schema_version").Scan(&version); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if version != 0 {
|
||||
t.Fatalf("failed migration advanced version to %d", version)
|
||||
}
|
||||
// 仅在本测试拥有的临时数据库中清除空冲突 schema,然后重试失败的迁移。
|
||||
if _, err := f.pool.Exec(f.ctx, "DROP SCHEMA postgresql_tenant_operator"); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
f.bootstrap(t)
|
||||
if _, err := f.pool.Exec(f.ctx, "UPDATE public.postgresql_tenant_operator_schema_version SET version = 100"); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if err := f.store.Bootstrap(f.ctx); err == nil {
|
||||
t.Fatal("bootstrap accepted a future migration version")
|
||||
}
|
||||
if err := f.pool.QueryRow(f.ctx, "SELECT version FROM public.postgresql_tenant_operator_schema_version").Scan(&version); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if version != 100 {
|
||||
t.Fatalf("bootstrap rewrote future version to %d", version)
|
||||
}
|
||||
}
|
||||
|
||||
// 数据库执行真实 COMMIT 后才注入客户端错误,模拟客户端无法确认提交结果。
|
||||
// 这不是网络故障测试,但能确定性覆盖已提交、调用者却收到失败的恢复分支。
|
||||
type lostCommitReply struct {
|
||||
registry.Beginner
|
||||
err error
|
||||
}
|
||||
|
||||
func (db lostCommitReply) Begin(ctx context.Context) (pgx.Tx, error) {
|
||||
tx, err := db.Beginner.Begin(ctx)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
return uncertainCommit{Tx: tx, err: db.err}, nil
|
||||
}
|
||||
|
||||
type uncertainCommit struct {
|
||||
pgx.Tx
|
||||
err error
|
||||
}
|
||||
|
||||
func (tx uncertainCommit) Commit(ctx context.Context) error {
|
||||
if err := tx.Tx.Commit(ctx); err != nil {
|
||||
return err
|
||||
}
|
||||
return tx.err
|
||||
}
|
||||
|
||||
func TestRegistryRetriesAfterLostCommitReply(t *testing.T) {
|
||||
f := newRegistryFixture(t)
|
||||
f.bootstrap(t)
|
||||
replyLost := errors.New("test-only lost commit reply")
|
||||
uncertain := registry.NewStore(lostCommitReply{Beginner: f.pool, err: replyLost})
|
||||
owner := registryOwner("uncertain")
|
||||
if _, err := uncertain.Claim(f.ctx, owner); !errors.Is(err, replyLost) {
|
||||
t.Fatalf("claim did not report lost reply: %v", err)
|
||||
}
|
||||
f.claim(t, owner, registry.ClaimOwned)
|
||||
if err := uncertain.MarkRetained(f.ctx, owner); !errors.Is(err, replyLost) {
|
||||
t.Fatalf("retain did not report lost reply: %v", err)
|
||||
}
|
||||
if err := f.store.MarkRetained(f.ctx, owner); err != nil {
|
||||
t.Fatalf("retry uncertain retain: %v", err)
|
||||
}
|
||||
if err := f.store.Delete(f.ctx, owner); !errors.Is(err, registry.ErrConflict) {
|
||||
t.Fatalf("uncertain retain lost tombstone protection: %v", err)
|
||||
}
|
||||
deletable := registryOwner("uncertain_delete")
|
||||
f.claim(t, deletable, registry.ClaimCreated)
|
||||
if err := uncertain.Delete(f.ctx, deletable); !errors.Is(err, replyLost) {
|
||||
t.Fatalf("delete did not report lost reply: %v", err)
|
||||
}
|
||||
if err := f.store.Delete(f.ctx, deletable); err != nil {
|
||||
t.Fatalf("retry uncertain delete: %v", err)
|
||||
}
|
||||
}
|
||||
|
||||
func TestRegistryConcurrentConflictingClaims(t *testing.T) {
|
||||
f := newRegistryFixture(t)
|
||||
f.bootstrap(t)
|
||||
first := registryOwner("first")
|
||||
second := registryOwner("second")
|
||||
second.DatabaseName = first.DatabaseName
|
||||
start := make(chan struct{})
|
||||
results := make(chan error, 2)
|
||||
var group sync.WaitGroup
|
||||
for _, owner := range []registry.Ownership{first, second} {
|
||||
group.Go(func() {
|
||||
<-start
|
||||
_, err := f.store.Claim(f.ctx, owner)
|
||||
results <- err
|
||||
})
|
||||
}
|
||||
close(start)
|
||||
group.Wait()
|
||||
close(results)
|
||||
created, conflicts := 0, 0
|
||||
for err := range results {
|
||||
switch {
|
||||
case err == nil:
|
||||
created++
|
||||
case errors.Is(err, registry.ErrConflict):
|
||||
conflicts++
|
||||
default:
|
||||
t.Fatalf("unexpected concurrent claim error: %v", err)
|
||||
}
|
||||
}
|
||||
if created != 1 || conflicts != 1 {
|
||||
t.Fatalf("concurrent reservation: created=%d conflicts=%d; want one of each", created, conflicts)
|
||||
}
|
||||
}
|
||||
|
||||
func TestRegistryConcurrentRetainAndDelete(t *testing.T) {
|
||||
f := newRegistryFixture(t)
|
||||
f.bootstrap(t)
|
||||
owner := registryOwner("delete_race")
|
||||
f.claim(t, owner, registry.ClaimCreated)
|
||||
start := make(chan struct{})
|
||||
retainResult := make(chan error, 1)
|
||||
deleteResult := make(chan error, 1)
|
||||
go func() {
|
||||
<-start
|
||||
retainResult <- f.store.MarkRetained(f.ctx, owner)
|
||||
}()
|
||||
go func() {
|
||||
<-start
|
||||
deleteResult <- f.store.Delete(f.ctx, owner)
|
||||
}()
|
||||
close(start)
|
||||
retainErr, deleteErr := <-retainResult, <-deleteResult
|
||||
record, readErr := f.store.Get(f.ctx, owner.InstanceUID, owner.TenantUID)
|
||||
switch {
|
||||
case retainErr == nil:
|
||||
// Retain 先取得行锁时,Delete 必须拒绝删除墓碑。
|
||||
if !errors.Is(deleteErr, registry.ErrConflict) || readErr != nil || record.Managed || record.RetainedAt == nil {
|
||||
t.Fatalf("retain won but tombstone was not protected: delete=%v read=%v record=%+v", deleteErr, readErr, record)
|
||||
}
|
||||
case errors.Is(retainErr, registry.ErrNotFound):
|
||||
// Delete 先提交时,Retain 必须报告记录已不存在,不能重建墓碑。
|
||||
if deleteErr != nil || !errors.Is(readErr, registry.ErrNotFound) {
|
||||
t.Fatalf("delete won but record remains: delete=%v read=%v", deleteErr, readErr)
|
||||
}
|
||||
default:
|
||||
t.Fatalf("unexpected retain/delete race: retain=%v delete=%v", retainErr, deleteErr)
|
||||
}
|
||||
}
|
||||
@@ -161,7 +161,7 @@ func credentialError(err error) error {
|
||||
return ErrCredentialsUnavailable
|
||||
}
|
||||
|
||||
// Forget 只释放本地连接;不删除数据库或 registry,不替代 Instance finalizer。
|
||||
// Forget 只释放本地连接;不删除数据库,不替代 Instance finalizer。
|
||||
func (s *InstanceService) Forget(name string) {
|
||||
s.mu.Lock()
|
||||
defer s.mu.Unlock()
|
||||
|
||||
@@ -18,7 +18,7 @@ package application
|
||||
|
||||
import "git.ddupan.top/panxiao81/ayatori/internal/database/domain/instance"
|
||||
|
||||
// DatabaseMetadata 是一次只读查询的事实,不包含管理权限或 registry 就绪结论。
|
||||
// DatabaseMetadata 是一次只读查询的事实,不包含管理权限或完整就绪结论。
|
||||
// AvailableExtensions 是服务器提供的可用列表,不是已安装列表或安装授权。
|
||||
type DatabaseMetadata struct {
|
||||
Version string
|
||||
|
||||
@@ -22,11 +22,10 @@ import "errors"
|
||||
type Phase string
|
||||
|
||||
const (
|
||||
PhasePending Phase = "Pending"
|
||||
PhaseValidating Phase = "Validating"
|
||||
PhaseInitializingRegistry Phase = "InitializingRegistry"
|
||||
PhaseReady Phase = "Ready"
|
||||
PhaseDeleting Phase = "Deleting"
|
||||
PhasePending Phase = "Pending"
|
||||
PhaseValidating Phase = "Validating"
|
||||
PhaseReady Phase = "Ready"
|
||||
PhaseDeleting Phase = "Deleting"
|
||||
)
|
||||
|
||||
type Readiness string
|
||||
@@ -61,7 +60,7 @@ func Reconstitute(target ObservationTarget, snapshot Snapshot, deleting bool) (*
|
||||
return nil, err
|
||||
}
|
||||
switch snapshot.Phase {
|
||||
case PhasePending, PhaseValidating, PhaseInitializingRegistry, PhaseReady, PhaseDeleting:
|
||||
case PhasePending, PhaseValidating, PhaseReady, PhaseDeleting:
|
||||
default:
|
||||
snapshot.Phase = PhasePending
|
||||
snapshot.Readiness = Unknown
|
||||
|
||||
@@ -39,7 +39,7 @@ func lifecycleInstance(t *testing.T, snapshot instance.Snapshot, deleting bool)
|
||||
// Acceptance: docs/database/domain-instance.md §3, checkpoint reconstruction and intent-only transitions.
|
||||
func TestReconstituteCheckpoints(t *testing.T) {
|
||||
for _, phase := range []instance.Phase{
|
||||
instance.PhasePending, instance.PhaseValidating, instance.PhaseInitializingRegistry,
|
||||
instance.PhasePending, instance.PhaseValidating,
|
||||
instance.PhaseReady, instance.PhaseDeleting,
|
||||
} {
|
||||
snapshot := instance.Snapshot{Phase: phase, ObservedRevision: 1, Readiness: instance.Ready, ReportedVersion: "17"}
|
||||
@@ -97,7 +97,7 @@ func TestDeletionRequiresRequestAndPreventsValidation(t *testing.T) {
|
||||
t.Fatal("rejected deletion mutated state")
|
||||
}
|
||||
for _, phase := range []instance.Phase{
|
||||
instance.PhasePending, instance.PhaseValidating, instance.PhaseInitializingRegistry,
|
||||
instance.PhasePending, instance.PhaseValidating,
|
||||
instance.PhaseReady, instance.PhaseDeleting,
|
||||
} {
|
||||
snapshot.Phase = phase
|
||||
|
||||
@@ -27,8 +27,6 @@ const (
|
||||
DependencyUnavailable
|
||||
AuthenticationFailed
|
||||
InsufficientPrivileges
|
||||
RegistryIncompatible
|
||||
RegistryNotUsable
|
||||
)
|
||||
|
||||
// CheckResult 的零值表示未观察,不能视为成功。
|
||||
@@ -70,32 +68,22 @@ func (c ManagementChecks) failure() Failure {
|
||||
return NoFailure
|
||||
}
|
||||
|
||||
type RegistryState uint8
|
||||
|
||||
const (
|
||||
RegistryUnobserved RegistryState = iota
|
||||
RegistryAbsent
|
||||
RegistryNeedsMigration
|
||||
RegistryUsable
|
||||
RegistryUnsupported
|
||||
RegistryUnavailable
|
||||
)
|
||||
|
||||
// CapabilityObservation 是值对象,不包含连接、凭据或可变集合。
|
||||
type CapabilityObservation struct {
|
||||
target ObservationTarget
|
||||
version string
|
||||
checks ManagementChecks
|
||||
registry RegistryState
|
||||
target ObservationTarget
|
||||
version string
|
||||
checks ManagementChecks
|
||||
}
|
||||
|
||||
func NewCapabilityObservation(target ObservationTarget, version string,
|
||||
checks ManagementChecks, registry RegistryState,
|
||||
func NewCapabilityObservation(
|
||||
target ObservationTarget,
|
||||
version string,
|
||||
checks ManagementChecks,
|
||||
) (CapabilityObservation, error) {
|
||||
if err := target.Validate(); err != nil {
|
||||
return CapabilityObservation{}, err
|
||||
}
|
||||
return CapabilityObservation{target: target, version: version, checks: checks, registry: registry}, nil
|
||||
return CapabilityObservation{target: target, version: version, checks: checks}, nil
|
||||
}
|
||||
|
||||
func (o CapabilityObservation) managementFailure() Failure {
|
||||
@@ -108,29 +96,6 @@ func (o CapabilityObservation) managementFailure() Failure {
|
||||
return NoFailure
|
||||
}
|
||||
|
||||
func (o CapabilityObservation) registryFailure() Failure {
|
||||
switch o.registry {
|
||||
case RegistryUsable:
|
||||
return NoFailure
|
||||
case RegistryAbsent, RegistryNeedsMigration:
|
||||
return RegistryNotUsable
|
||||
case RegistryUnsupported:
|
||||
return RegistryIncompatible
|
||||
case RegistryUnavailable:
|
||||
return DependencyUnavailable
|
||||
default:
|
||||
return ObservationIncomplete
|
||||
}
|
||||
}
|
||||
|
||||
type PreparationDecision uint8
|
||||
|
||||
const (
|
||||
PreparationDenied PreparationDecision = iota
|
||||
PreparationAllowed
|
||||
AlreadyUsable
|
||||
)
|
||||
|
||||
func (i *Instance) acceptObservation(o CapabilityObservation, phase Phase) error {
|
||||
if !i.target.Matches(o.target) {
|
||||
return errors.New("capability observation target does not match instance")
|
||||
@@ -149,75 +114,11 @@ func (i *Instance) fail(failure Failure) {
|
||||
i.snapshot.ObservedRevision = i.target.Revision().Value()
|
||||
}
|
||||
|
||||
// AssessManagement 只推进意图,不执行 registry 写入,也不完成 observedRevision。
|
||||
// AssessManagement 根据本轮完整能力观察完成验证,不执行外部写入。
|
||||
func (i *Instance) AssessManagement(o CapabilityObservation) error {
|
||||
if err := i.acceptObservation(o, PhaseValidating); err != nil {
|
||||
return err
|
||||
}
|
||||
if failure := o.managementFailure(); failure != NoFailure {
|
||||
i.fail(failure)
|
||||
return nil
|
||||
}
|
||||
if failure := o.registryFailure(); failure != NoFailure && failure != RegistryNotUsable {
|
||||
i.fail(failure)
|
||||
return nil
|
||||
}
|
||||
i.snapshot.Phase = PhaseInitializingRegistry
|
||||
i.snapshot.Readiness = Unknown
|
||||
i.snapshot.Failure = NoFailure
|
||||
i.evidence = nil
|
||||
return nil
|
||||
}
|
||||
|
||||
// PlanRegistryPreparation 不证明 checkpoint 已落盘;应用层必须先保存意图再执行写入。
|
||||
func (i *Instance) PlanRegistryPreparation(o CapabilityObservation) (PreparationDecision, error) {
|
||||
if err := i.acceptObservation(o, PhaseInitializingRegistry); err != nil {
|
||||
return PreparationDenied, err
|
||||
}
|
||||
if failure := o.managementFailure(); failure != NoFailure {
|
||||
i.fail(failure)
|
||||
return PreparationDenied, nil
|
||||
}
|
||||
switch o.registry {
|
||||
case RegistryUsable:
|
||||
return AlreadyUsable, nil
|
||||
case RegistryAbsent, RegistryNeedsMigration:
|
||||
return PreparationAllowed, nil
|
||||
default:
|
||||
i.fail(o.registryFailure())
|
||||
return PreparationDenied, nil
|
||||
}
|
||||
}
|
||||
|
||||
// RegistryPreparationResult 只能是安全失败或完整回读,不能表达裸操作成功。
|
||||
type RegistryPreparationResult struct {
|
||||
observation CapabilityObservation
|
||||
failure Failure
|
||||
}
|
||||
|
||||
func RegistryReadBack(o CapabilityObservation) RegistryPreparationResult {
|
||||
return RegistryPreparationResult{observation: o}
|
||||
}
|
||||
|
||||
func RegistryPreparationFailed(target ObservationTarget, failure Failure) (RegistryPreparationResult, error) {
|
||||
if err := target.Validate(); err != nil {
|
||||
return RegistryPreparationResult{}, err
|
||||
}
|
||||
if failure < ObservationIncomplete || failure > RegistryNotUsable {
|
||||
return RegistryPreparationResult{}, errors.New("registry preparation requires a known failure category")
|
||||
}
|
||||
return RegistryPreparationResult{observation: CapabilityObservation{target: target}, failure: failure}, nil
|
||||
}
|
||||
|
||||
func (i *Instance) AssessRegistryResult(result RegistryPreparationResult) error {
|
||||
o := result.observation
|
||||
if err := i.acceptObservation(o, PhaseInitializingRegistry); err != nil {
|
||||
return err
|
||||
}
|
||||
if result.failure != NoFailure {
|
||||
i.fail(result.failure)
|
||||
return nil
|
||||
}
|
||||
i.assessComplete(o)
|
||||
return nil
|
||||
}
|
||||
@@ -227,12 +128,12 @@ func (i *Instance) assessComplete(o CapabilityObservation) {
|
||||
i.fail(failure)
|
||||
return
|
||||
}
|
||||
if failure := o.registryFailure(); failure != NoFailure {
|
||||
i.fail(failure)
|
||||
return
|
||||
i.snapshot = Snapshot{
|
||||
Phase: PhaseReady,
|
||||
ObservedRevision: i.target.Revision().Value(),
|
||||
Readiness: Ready,
|
||||
ReportedVersion: o.version,
|
||||
}
|
||||
i.snapshot = Snapshot{Phase: PhaseReady, ObservedRevision: i.target.Revision().Value(),
|
||||
Readiness: Ready, ReportedVersion: o.version}
|
||||
i.evidence = &o
|
||||
}
|
||||
|
||||
@@ -244,10 +145,8 @@ func (i *Instance) AssessReadiness(o CapabilityObservation) error {
|
||||
if i.snapshot.ObservedRevision != i.target.Revision().Value() {
|
||||
return i.BeginValidation()
|
||||
}
|
||||
if o.managementFailure() != NoFailure || o.registry == RegistryUnavailable {
|
||||
if o.managementFailure() != NoFailure {
|
||||
i.snapshot.Phase = PhaseValidating
|
||||
} else if o.registryFailure() != NoFailure {
|
||||
i.snapshot.Phase = PhaseInitializingRegistry
|
||||
}
|
||||
i.assessComplete(o)
|
||||
return nil
|
||||
|
||||
@@ -32,11 +32,9 @@ func completeChecks() instance.ManagementChecks {
|
||||
}
|
||||
}
|
||||
|
||||
func capability(t *testing.T, value *instance.Instance, checks instance.ManagementChecks,
|
||||
registry instance.RegistryState,
|
||||
) instance.CapabilityObservation {
|
||||
func capability(t *testing.T, value *instance.Instance, checks instance.ManagementChecks) instance.CapabilityObservation {
|
||||
t.Helper()
|
||||
o, err := instance.NewCapabilityObservation(value.Target(), testServerVersion, checks, registry)
|
||||
o, err := instance.NewCapabilityObservation(value.Target(), testServerVersion, checks)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
@@ -45,8 +43,8 @@ func capability(t *testing.T, value *instance.Instance, checks instance.Manageme
|
||||
|
||||
func readyInstance(t *testing.T) *instance.Instance {
|
||||
t.Helper()
|
||||
i := lifecycleInstance(t, instance.Snapshot{Phase: instance.PhaseInitializingRegistry}, false)
|
||||
if err := i.AssessRegistryResult(instance.RegistryReadBack(capability(t, i, completeChecks(), instance.RegistryUsable))); err != nil {
|
||||
i := lifecycleInstance(t, instance.Snapshot{Phase: instance.PhaseValidating}, false)
|
||||
if err := i.AssessManagement(capability(t, i, completeChecks())); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if err := i.RequireProvisioningReady(); err != nil {
|
||||
@@ -60,38 +58,27 @@ func TestReadinessRequiresCompleteReadBack(t *testing.T) {
|
||||
if err := i.BeginValidation(); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
absent := capability(t, i, completeChecks(), instance.RegistryAbsent)
|
||||
if err := i.AssessManagement(absent); err != nil {
|
||||
if i.RequireProvisioningReady() == nil {
|
||||
t.Fatal("validation intent authorized provisioning")
|
||||
}
|
||||
if err := i.AssessManagement(capability(t, i, instance.ManagementChecks{})); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if s := i.Snapshot(); s.Phase != instance.PhaseInitializingRegistry || s.ObservedRevision != 0 || s.Readiness != instance.Unknown {
|
||||
t.Fatalf("management observation prematurely concluded readiness: %+v", s)
|
||||
if snapshot := i.Snapshot(); snapshot.Phase != instance.PhaseValidating ||
|
||||
snapshot.Readiness != instance.NotReady || snapshot.Failure != instance.ObservationIncomplete {
|
||||
t.Fatalf("incomplete observation accepted: %+v", snapshot)
|
||||
}
|
||||
for range 2 {
|
||||
decision, err := i.PlanRegistryPreparation(absent)
|
||||
if err != nil || decision != instance.PreparationAllowed {
|
||||
t.Fatalf("preparation: %v, %v", decision, err)
|
||||
}
|
||||
if i.RequireProvisioningReady() == nil {
|
||||
t.Fatal("preparation authorized provisioning")
|
||||
}
|
||||
if i.RequireProvisioningReady() == nil {
|
||||
t.Fatal("incomplete observation authorized provisioning")
|
||||
}
|
||||
if err := i.AssessRegistryResult(instance.RegistryReadBack(absent)); err != nil {
|
||||
if err := i.AssessManagement(capability(t, i, completeChecks())); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if i.Snapshot().Failure != instance.RegistryNotUsable || i.RequireProvisioningReady() == nil {
|
||||
t.Fatal("absent registry accepted as ready")
|
||||
}
|
||||
usable := capability(t, i, completeChecks(), instance.RegistryUsable)
|
||||
decision, err := i.PlanRegistryPreparation(usable)
|
||||
if err != nil || decision != instance.AlreadyUsable {
|
||||
t.Fatalf("retry after external preparation: %v, %v", decision, err)
|
||||
}
|
||||
if err := i.AssessRegistryResult(instance.RegistryReadBack(usable)); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if s := i.Snapshot(); s.Readiness != instance.Ready || s.ReportedVersion != testServerVersion || s.ObservedRevision != i.Target().Revision().Value() {
|
||||
t.Fatalf("complete observation not accepted: %+v", s)
|
||||
snapshot := i.Snapshot()
|
||||
if snapshot.Phase != instance.PhaseReady || snapshot.Readiness != instance.Ready ||
|
||||
snapshot.ReportedVersion != testServerVersion ||
|
||||
snapshot.ObservedRevision != i.Target().Revision().Value() {
|
||||
t.Fatalf("complete management observation did not establish readiness: %+v", snapshot)
|
||||
}
|
||||
if err := i.RequireProvisioningReady(); err != nil {
|
||||
t.Fatal(err)
|
||||
@@ -109,7 +96,7 @@ func TestReadinessRecoveryAndInvalidation(t *testing.T) {
|
||||
t.Fatal("persisted Ready fabricated fresh evidence")
|
||||
}
|
||||
for range 2 {
|
||||
if err := restored.AssessReadiness(capability(t, restored, completeChecks(), instance.RegistryUsable)); err != nil {
|
||||
if err := restored.AssessReadiness(capability(t, restored, completeChecks())); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if err := restored.RequireProvisioningReady(); err != nil {
|
||||
@@ -138,120 +125,77 @@ func TestReadinessRecoveryAndInvalidation(t *testing.T) {
|
||||
}
|
||||
|
||||
func TestEachManagementCheckIsRequired(t *testing.T) {
|
||||
for field := range 6 {
|
||||
for _, result := range []instance.CheckResult{instance.CheckUnobserved, instance.CheckUnavailable,
|
||||
instance.CheckAuthenticationFailed, instance.CheckInsufficientPrivileges, 255} {
|
||||
checks := completeChecks()
|
||||
fields := []*instance.CheckResult{&checks.Connection, &checks.Metadata, &checks.Roles,
|
||||
&checks.Databases, &checks.Grants, &checks.Extensions}
|
||||
*fields[field] = result
|
||||
i := readyInstance(t)
|
||||
if err := i.AssessReadiness(capability(t, i, checks, instance.RegistryUsable)); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if s := i.Snapshot(); s.Phase != instance.PhaseValidating || s.Readiness != instance.NotReady ||
|
||||
s.Failure == instance.NoFailure || i.RequireProvisioningReady() == nil {
|
||||
t.Fatalf("check %d result %d accepted: %+v", field, result, s)
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func TestRegistryDecisionsAndReadinessLoss(t *testing.T) {
|
||||
for _, tc := range []struct {
|
||||
state instance.RegistryState
|
||||
decision instance.PreparationDecision
|
||||
failure instance.Failure
|
||||
checkNames := []string{"connection", "metadata", "roles", "databases", "grants", "extensions"}
|
||||
failures := []struct {
|
||||
name string
|
||||
result instance.CheckResult
|
||||
want instance.Failure
|
||||
}{
|
||||
{instance.RegistryUsable, instance.AlreadyUsable, instance.NoFailure},
|
||||
{instance.RegistryAbsent, instance.PreparationAllowed, instance.RegistryNotUsable},
|
||||
{instance.RegistryNeedsMigration, instance.PreparationAllowed, instance.RegistryNotUsable},
|
||||
{instance.RegistryUnsupported, instance.PreparationDenied, instance.RegistryIncompatible},
|
||||
{instance.RegistryUnavailable, instance.PreparationDenied, instance.DependencyUnavailable},
|
||||
{instance.RegistryUnobserved, instance.PreparationDenied, instance.ObservationIncomplete},
|
||||
{255, instance.PreparationDenied, instance.ObservationIncomplete},
|
||||
} {
|
||||
i := lifecycleInstance(t, instance.Snapshot{Phase: instance.PhaseInitializingRegistry}, false)
|
||||
o := capability(t, i, completeChecks(), tc.state)
|
||||
decision, err := i.PlanRegistryPreparation(o)
|
||||
if err != nil || decision != tc.decision {
|
||||
t.Fatalf("registry %d: %v, %v", tc.state, decision, err)
|
||||
}
|
||||
i = readyInstance(t)
|
||||
if err := i.AssessReadiness(o); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if i.Snapshot().Failure != tc.failure {
|
||||
t.Fatalf("registry %d: %+v", tc.state, i.Snapshot())
|
||||
}
|
||||
if tc.state != instance.RegistryUsable {
|
||||
wantPhase := instance.PhaseInitializingRegistry
|
||||
if tc.state == instance.RegistryUnavailable {
|
||||
wantPhase = instance.PhaseValidating
|
||||
}
|
||||
if i.Snapshot().Phase != wantPhase || i.RequireProvisioningReady() == nil {
|
||||
t.Fatal("registry drift retained readiness")
|
||||
{"unobserved", instance.CheckUnobserved, instance.ObservationIncomplete},
|
||||
{"unavailable", instance.CheckUnavailable, instance.DependencyUnavailable},
|
||||
{"authentication", instance.CheckAuthenticationFailed, instance.AuthenticationFailed},
|
||||
{"privileges", instance.CheckInsufficientPrivileges, instance.InsufficientPrivileges},
|
||||
{"unknown", 255, instance.ObservationIncomplete},
|
||||
}
|
||||
for field, name := range checkNames {
|
||||
for _, failure := range failures {
|
||||
for _, phase := range []instance.Phase{instance.PhaseValidating, instance.PhaseReady} {
|
||||
t.Run(name+"/"+failure.name+"/"+string(phase), func(t *testing.T) {
|
||||
checks := completeChecks()
|
||||
fields := []*instance.CheckResult{
|
||||
&checks.Connection, &checks.Metadata, &checks.Roles,
|
||||
&checks.Databases, &checks.Grants, &checks.Extensions,
|
||||
}
|
||||
*fields[field] = failure.result
|
||||
value := lifecycleInstance(t, instance.Snapshot{Phase: phase}, false)
|
||||
assess := value.AssessManagement
|
||||
if phase == instance.PhaseReady {
|
||||
value = readyInstance(t)
|
||||
assess = value.AssessReadiness
|
||||
}
|
||||
if err := assess(capability(t, value, checks)); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
snapshot := value.Snapshot()
|
||||
if snapshot.Phase != instance.PhaseValidating ||
|
||||
snapshot.Readiness != instance.NotReady ||
|
||||
snapshot.Failure != failure.want ||
|
||||
snapshot.ObservedRevision != value.Target().Revision().Value() {
|
||||
t.Fatalf("incorrect failed observation: %+v", snapshot)
|
||||
}
|
||||
if value.RequireProvisioningReady() == nil {
|
||||
t.Fatal("failed check authorized provisioning")
|
||||
}
|
||||
// 依赖恢复后重新验证,不保留失败或旧就绪证据。
|
||||
if err := value.AssessManagement(capability(t, value, completeChecks())); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if err := value.RequireProvisioningReady(); err != nil {
|
||||
t.Fatal("dependency recovery did not restore readiness", err)
|
||||
}
|
||||
})
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func TestInitializationRejectsIncompleteOrFailedManagement(t *testing.T) {
|
||||
for _, tc := range []struct {
|
||||
checks instance.ManagementChecks
|
||||
registry instance.RegistryState
|
||||
failure instance.Failure
|
||||
}{
|
||||
{instance.ManagementChecks{}, instance.RegistryUsable, instance.ObservationIncomplete},
|
||||
{completeChecks(), instance.RegistryUnsupported, instance.RegistryIncompatible},
|
||||
{completeChecks(), instance.RegistryUnavailable, instance.DependencyUnavailable},
|
||||
} {
|
||||
i := lifecycleInstance(t, instance.Snapshot{Phase: instance.PhaseValidating}, false)
|
||||
if err := i.AssessManagement(capability(t, i, tc.checks, tc.registry)); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if s := i.Snapshot(); s.Phase != instance.PhaseValidating || s.Failure != tc.failure ||
|
||||
s.ObservedRevision != i.Target().Revision().Value() || s.Readiness != instance.NotReady {
|
||||
t.Fatalf("invalid management accepted: %+v", s)
|
||||
}
|
||||
}
|
||||
i := lifecycleInstance(t, instance.Snapshot{Phase: instance.PhaseInitializingRegistry}, false)
|
||||
decision, err := i.PlanRegistryPreparation(capability(t, i, instance.ManagementChecks{}, instance.RegistryAbsent))
|
||||
if err != nil || decision != instance.PreparationDenied || i.Snapshot().Failure != instance.ObservationIncomplete {
|
||||
t.Fatal("incomplete management allowed registry writes")
|
||||
}
|
||||
if err := i.AssessRegistryResult(instance.RegistryReadBack(capability(t, i, completeChecks(), instance.RegistryUsable))); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if err := i.RequireProvisioningReady(); err != nil {
|
||||
t.Fatal("dependency recovery did not restore readiness", err)
|
||||
}
|
||||
}
|
||||
|
||||
func TestReadinessMethodsRejectWrongPhaseAndDeletion(t *testing.T) {
|
||||
for _, deleting := range []bool{false, true} {
|
||||
for _, phase := range []instance.Phase{instance.PhasePending, instance.PhaseValidating,
|
||||
instance.PhaseInitializingRegistry, instance.PhaseReady, instance.PhaseDeleting} {
|
||||
instance.PhaseReady, instance.PhaseDeleting} {
|
||||
for _, operation := range []struct {
|
||||
phase instance.Phase
|
||||
apply func(*instance.Instance, instance.CapabilityObservation) error
|
||||
}{
|
||||
{instance.PhaseValidating, (*instance.Instance).AssessManagement},
|
||||
{instance.PhaseReady, (*instance.Instance).AssessReadiness},
|
||||
{instance.PhaseInitializingRegistry, func(i *instance.Instance, o instance.CapabilityObservation) error {
|
||||
_, err := i.PlanRegistryPreparation(o)
|
||||
return err
|
||||
}},
|
||||
{instance.PhaseInitializingRegistry, func(i *instance.Instance, o instance.CapabilityObservation) error {
|
||||
return i.AssessRegistryResult(instance.RegistryReadBack(o))
|
||||
}},
|
||||
} {
|
||||
if !deleting && operation.phase == phase {
|
||||
continue
|
||||
}
|
||||
i := lifecycleInstance(t, instance.Snapshot{Phase: phase}, deleting)
|
||||
before := i.Snapshot()
|
||||
if err := operation.apply(i, capability(t, i, completeChecks(), instance.RegistryUsable)); err == nil {
|
||||
if err := operation.apply(i, capability(t, i, completeChecks())); err == nil {
|
||||
t.Fatalf("phase %s deleting=%t accepted operation for %s", phase, deleting, operation.phase)
|
||||
}
|
||||
if i.Snapshot() != before {
|
||||
@@ -273,7 +217,7 @@ func TestOldGenerationObservationDoesNotReplaceEvidence(t *testing.T) {
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
o, err := instance.NewCapabilityObservation(other, testServerVersion, completeChecks(), instance.RegistryUsable)
|
||||
o, err := instance.NewCapabilityObservation(other, testServerVersion, completeChecks())
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
@@ -286,44 +230,14 @@ func TestOldGenerationObservationDoesNotReplaceEvidence(t *testing.T) {
|
||||
}
|
||||
}
|
||||
|
||||
func TestPreparationFailureCannotEstablishReadiness(t *testing.T) {
|
||||
i := lifecycleInstance(t, instance.Snapshot{Phase: instance.PhaseInitializingRegistry}, false)
|
||||
for _, failure := range []instance.Failure{instance.DependencyUnavailable, instance.AuthenticationFailed,
|
||||
instance.InsufficientPrivileges, instance.RegistryIncompatible} {
|
||||
result, err := instance.RegistryPreparationFailed(i.Target(), failure)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if err := i.AssessRegistryResult(result); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if s := i.Snapshot(); s.Failure != failure || s.Readiness != instance.NotReady ||
|
||||
s.Phase != instance.PhaseInitializingRegistry || i.RequireProvisioningReady() == nil {
|
||||
t.Fatalf("failed operation accepted: %+v", s)
|
||||
}
|
||||
}
|
||||
for _, failure := range []instance.Failure{instance.NoFailure, 255} {
|
||||
if _, err := instance.RegistryPreparationFailed(i.Target(), failure); err == nil {
|
||||
t.Fatal("invalid failure accepted")
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func TestCapabilityInputsAndLifecycleGuards(t *testing.T) {
|
||||
i := readyInstance(t)
|
||||
if _, err := instance.NewCapabilityObservation(instance.ObservationTarget{}, testServerVersion,
|
||||
completeChecks(), instance.RegistryUsable); err == nil {
|
||||
completeChecks()); err == nil {
|
||||
t.Fatal("invalid target accepted")
|
||||
}
|
||||
if _, err := instance.RegistryPreparationFailed(instance.ObservationTarget{}, instance.DependencyUnavailable); err == nil {
|
||||
t.Fatal("invalid failure target accepted")
|
||||
}
|
||||
for _, method := range []func(instance.CapabilityObservation) error{
|
||||
i.AssessManagement, i.AssessReadiness,
|
||||
func(o instance.CapabilityObservation) error { _, err := i.PlanRegistryPreparation(o); return err },
|
||||
func(o instance.CapabilityObservation) error {
|
||||
return i.AssessRegistryResult(instance.RegistryReadBack(o))
|
||||
},
|
||||
} {
|
||||
before := i.Snapshot()
|
||||
if err := method(instance.CapabilityObservation{}); err == nil || i.Snapshot() != before {
|
||||
@@ -333,13 +247,13 @@ func TestCapabilityInputsAndLifecycleGuards(t *testing.T) {
|
||||
old := i.Snapshot()
|
||||
old.ObservedRevision = 0
|
||||
changed := lifecycleInstance(t, old, false)
|
||||
if err := changed.AssessReadiness(capability(t, changed, completeChecks(), instance.RegistryUsable)); err != nil {
|
||||
if err := changed.AssessReadiness(capability(t, changed, completeChecks())); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if s := changed.Snapshot(); s.Phase != instance.PhaseValidating || s.ObservedRevision != 0 || s.Readiness != instance.Unknown {
|
||||
t.Fatalf("changed generation accepted old checkpoint: %+v", s)
|
||||
}
|
||||
o, err := instance.NewCapabilityObservation(i.Target(), "", completeChecks(), instance.RegistryUsable)
|
||||
o, err := instance.NewCapabilityObservation(i.Target(), "", completeChecks())
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user