feat: 实现 PostgreSQL 所有权 registry #4

Merged
panxiao81 merged 2 commits from feature/postgresql-registry into main 2026-09-10 10:59:10 +00:00
6 changed files with 544 additions and 0 deletions
Showing only changes of commit 558df2bbcd - Show all commits
+6
View File
@@ -18,6 +18,8 @@ CONTAINER_TOOL ?= docker
# DEV_COMPOSE manages disposable PostgreSQL and OpenBao dependencies for local work.
DEV_COMPOSE ?= docker compose -f hack/dev/compose.yaml
POSTGRES_DEV_PORT ?= 15432
POSTGRES_TEST_DSN ?= postgres://postgres:[email protected]:$(POSTGRES_DEV_PORT)/postgres?sslmode=disable
# Setting SHELL to bash allows bash commands to be executed by recipes.
# Options are set to exit when a recipe line exits non-zero or a piped command fails.
@@ -59,6 +61,10 @@ dev-smoke: dev-up ## Verify PostgreSQL and OpenBao development dependencies.
dev-down: ## Remove disposable development dependencies and their data.
$(DEV_COMPOSE) down --volumes --remove-orphans
.PHONY: test-integration
test-integration: dev-up ## Run adapter integration tests against disposable dependencies.
POSTGRES_TEST_DSN="$(POSTGRES_TEST_DSN)" go test ./internal/postgresql/...
.PHONY: manifests
manifests: controller-gen ## Generate WebhookConfiguration, ClusterRole and CustomResourceDefinition objects.
"$(CONTROLLER_GEN)" rbac:roleName=manager-role crd webhook paths="./..." output:crd:artifacts:config=config/crd/bases
+4
View File
@@ -3,6 +3,7 @@ module git.ddupan.top/panxiao81/postgresql-tenant-operator
go 1.26.0
require (
github.com/jackc/pgx/v5 v5.11.0
github.com/onsi/ginkgo/v2 v2.27.4
github.com/onsi/gomega v1.39.0
k8s.io/apimachinery v0.36.0
@@ -38,6 +39,9 @@ require (
github.com/google/uuid v1.6.0 // indirect
github.com/grpc-ecosystem/grpc-gateway/v2 v2.27.7 // indirect
github.com/inconshreveable/mousetrap v1.1.0 // indirect
github.com/jackc/pgpassfile v1.0.0 // indirect
github.com/jackc/pgservicefile v0.0.0-20240606120523-5a60cdf6a761 // indirect
github.com/jackc/puddle/v2 v2.2.2 // indirect
github.com/josharian/intern v1.0.0 // indirect
github.com/json-iterator/go v1.1.12 // indirect
github.com/mailru/easyjson v0.7.7 // indirect
+9
View File
@@ -74,6 +74,14 @@ github.com/grpc-ecosystem/grpc-gateway/v2 v2.27.7 h1:X+2YciYSxvMQK0UZ7sg45ZVabVZ
github.com/grpc-ecosystem/grpc-gateway/v2 v2.27.7/go.mod h1:lW34nIZuQ8UDPdkon5fmfp2l3+ZkQ2me/+oecHYLOII=
github.com/inconshreveable/mousetrap v1.1.0 h1:wN+x4NVGpMsO7ErUn/mUI3vEoE6Jt13X2s0bqwp9tc8=
github.com/inconshreveable/mousetrap v1.1.0/go.mod h1:vpF70FUmC8bwa3OWnCshd2FqLfsEA9PFc4w1p2J65bw=
github.com/jackc/pgpassfile v1.0.0 h1:/6Hmqy13Ss2zCq62VdNG8tM1wchn8zjSGOBJ6icpsIM=
github.com/jackc/pgpassfile v1.0.0/go.mod h1:CEx0iS5ambNFdcRtxPj5JhEz+xB6uRky5eyVu/W2HEg=
github.com/jackc/pgservicefile v0.0.0-20240606120523-5a60cdf6a761 h1:iCEnooe7UlwOQYpKFhBabPMi4aNAfoODPEFNiAnClxo=
github.com/jackc/pgservicefile v0.0.0-20240606120523-5a60cdf6a761/go.mod h1:5TJZWKEWniPve33vlWYSoGYefn3gLQRzjfDlhSJ9ZKM=
github.com/jackc/pgx/v5 v5.11.0 h1:IzBBtyK9AHqf98cctWFifYSci2hgQR/cd56wB4p+ogg=
github.com/jackc/pgx/v5 v5.11.0/go.mod h1:mal1tBGAFfLHvZzaYh77YS/eC6IX9OWbRV1QIIM0Jn4=
github.com/jackc/puddle/v2 v2.2.2 h1:PR8nw+E/1w0GLuRFSmiioY6UooMp6KJv0/61nB7icHo=
github.com/jackc/puddle/v2 v2.2.2/go.mod h1:vriiEXHvEE654aYKXXjOvZM39qJ0q+azkZFrfEOc3H4=
github.com/josharian/intern v1.0.0 h1:vlS4z54oSdjm0bgjRigI+G1HpF+tI+9rE5LLzOg8HmY=
github.com/josharian/intern v1.0.0/go.mod h1:5DoeVV0s6jJacbCEi61lwdGj/aVlrQvzHFFd8Hwg//Y=
github.com/joshdk/go-junit v1.0.0 h1:S86cUKIdwBHWwA6xCmFlf3RTLfVXYQfvanM5Uh+K6GE=
@@ -137,6 +145,7 @@ github.com/stretchr/objx v0.5.0/go.mod h1:Yh+to48EsGEfYuaHDzXPcE3xhTkx73EhmCGUpE
github.com/stretchr/objx v0.5.2 h1:xuMeJ0Sdp5ZMRXx/aWO6RZxdr3beISkG5/G/aIRr3pY=
github.com/stretchr/objx v0.5.2/go.mod h1:FRsXN1f5AsAjCGJKqEizvkpNtU+EGNCLh3NxZ/8L+MA=
github.com/stretchr/testify v1.3.0/go.mod h1:M5WIy9Dh21IEIfnGCwXGc5bZfKNJtfHm1UVUgZn+9EI=
github.com/stretchr/testify v1.7.0/go.mod h1:6Fq8oRcR53rry900zMqJjRRixrwX3KX962/h/Wwjteg=
github.com/stretchr/testify v1.7.1/go.mod h1:6Fq8oRcR53rry900zMqJjRRixrwX3KX962/h/Wwjteg=
github.com/stretchr/testify v1.8.0/go.mod h1:yNjHg4UonilssWZ8iaSj1OCr/vHnekPRkoO+kdMU+MU=
github.com/stretchr/testify v1.8.1/go.mod h1:w2LPCIKwWwSfY2zedu0+kehJoqGctiVI29o6fzry7u4=
+307
View File
@@ -0,0 +1,307 @@
/*
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"
"errors"
"fmt"
"time"
"github.com/jackc/pgx/v5"
)
const schemaVersion = 1
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 creates the private schema and the current registry table idempotently.
panxiao81 marked this conversation as resolved
Review

要么用bootstrap脚本在init容器里跑,要么用migrate,不要用土办法

要么用bootstrap脚本在init容器里跑,要么用migrate,不要用土办法
func (s *Store) Bootstrap(ctx context.Context) error {
if s == nil || s.db == nil {
return errors.New("bootstrap registry: nil database")
}
tx, err := s.db.Begin(ctx)
if err != nil {
return fmt.Errorf("begin registry bootstrap: %w", err)
}
defer func() { _ = tx.Rollback(ctx) }()
for _, statement := range bootstrapStatements {
if _, err := tx.Exec(ctx, statement); err != nil {
return fmt.Errorf("bootstrap registry schema version %d: %w", schemaVersion, err)
}
}
var version int
if err := tx.QueryRow(ctx, `SELECT version FROM postgresql_tenant_operator.schema_version WHERE singleton`).Scan(&version); err != nil {
return fmt.Errorf("read registry schema version: %w", err)
}
if version != schemaVersion {
return fmt.Errorf("registry schema version %d is unsupported; expected %d", version, schemaVersion)
}
if err := tx.Commit(ctx); err != nil {
return fmt.Errorf("commit registry bootstrap: %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) }()
commandTag, err := tx.Exec(ctx, claimStatement,
panxiao81 marked this conversation as resolved
Review

pg sql有RETURN的

pg sql有RETURN的
owner.InstanceUID,
owner.TenantUID,
owner.TenantNamespace,
owner.TenantName,
owner.DatabaseName,
owner.RoleName,
owner.CredentialPath,
)
if err != nil {
return "", fmt.Errorf("insert registry claim: %w", err)
}
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: existing record does not match the requested managed owner", ErrConflict)
}
if err := tx.Commit(ctx); err != nil {
return "", fmt.Errorf("commit registry claim: %w", err)
}
if commandTag.RowsAffected() == 1 {
return ClaimCreated, nil
}
return ClaimOwned, nil
}
// 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
}
@@ -0,0 +1,132 @@
/*
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
import (
"context"
"errors"
"os"
"testing"
"github.com/jackc/pgx/v5/pgconn"
"github.com/jackc/pgx/v5/pgxpool"
)
func TestRegistryLifecycle(t *testing.T) {
dsn := os.Getenv("POSTGRES_TEST_DSN")
if dsn == "" {
t.Skip("POSTGRES_TEST_DSN is not set")
}
ctx := context.Background()
pool, err := pgxpool.New(ctx, dsn)
if err != nil {
t.Fatalf("create PostgreSQL pool: %v", err)
}
t.Cleanup(pool.Close)
if err := pool.Ping(ctx); err != nil {
t.Fatalf("ping PostgreSQL: %v", err)
}
dropRegistrySchema(t, ctx, pool)
t.Cleanup(func() { dropRegistrySchema(t, ctx, pool) })
store := NewStore(pool)
if err := store.Bootstrap(ctx); err != nil {
t.Fatalf("Bootstrap() error = %v", err)
}
if err := store.Bootstrap(ctx); err != nil {
t.Fatalf("second Bootstrap() error = %v", err)
}
owner := testOwnership("tenant-uid-1", "netbox", "netbox", "netbox", "postgresql-tenants/netbox/netbox")
result, err := store.Claim(ctx, owner)
if err != nil || result != ClaimCreated {
t.Fatalf("first Claim() = %q, %v; want %q, nil", result, err, ClaimCreated)
}
result, err = store.Claim(ctx, owner)
if err != nil || result != ClaimOwned {
t.Fatalf("second Claim() = %q, %v; want %q, nil", result, err, ClaimOwned)
}
record, err := store.Get(ctx, owner.InstanceUID, owner.TenantUID)
if err != nil {
t.Fatalf("Get() error = %v", err)
}
if record.Ownership != owner || !record.Managed || record.RetainedAt != nil {
t.Fatalf("Get() = %#v; want matching managed owner", record)
}
conflict := testOwnership("tenant-uid-2", "other", owner.DatabaseName, "other", "postgresql-tenants/other/other")
if _, err := store.Claim(ctx, conflict); !errors.Is(err, ErrConflict) {
t.Fatalf("conflicting Claim() error = %v, want ErrConflict", err)
}
if err := store.MarkRetained(ctx, owner); err != nil {
t.Fatalf("MarkRetained() error = %v", err)
}
if err := store.MarkRetained(ctx, owner); err != nil {
t.Fatalf("second MarkRetained() error = %v", err)
}
record, err = store.Get(ctx, owner.InstanceUID, owner.TenantUID)
if err != nil || record.Managed || record.RetainedAt == nil {
t.Fatalf("retained Get() = %#v, %v; want unmanaged tombstone", record, err)
}
if err := store.Delete(ctx, owner); !errors.Is(err, ErrConflict) {
t.Fatalf("Delete(retained) error = %v, want ErrConflict", err)
}
if _, err := store.Claim(ctx, owner); !errors.Is(err, ErrConflict) {
t.Fatalf("Claim(retained) error = %v, want ErrConflict", err)
}
deletable := testOwnership("tenant-uid-3", "gitea", "gitea", "gitea", "postgresql-tenants/gitea/gitea")
if _, err := store.Claim(ctx, deletable); err != nil {
t.Fatalf("Claim(deletable) error = %v", err)
}
if err := store.Delete(ctx, deletable); err != nil {
t.Fatalf("Delete() error = %v", err)
}
if err := store.Delete(ctx, deletable); err != nil {
t.Fatalf("second Delete() error = %v", err)
}
if _, err := store.Get(ctx, deletable.InstanceUID, deletable.TenantUID); !errors.Is(err, ErrNotFound) {
t.Fatalf("Get(deleted) error = %v, want ErrNotFound", err)
}
}
func testOwnership(tenantUID, tenantName, databaseName, roleName, credentialPath string) Ownership {
return Ownership{
InstanceUID: "instance-uid-1",
TenantUID: tenantUID,
TenantNamespace: tenantName,
TenantName: tenantName,
DatabaseName: databaseName,
RoleName: roleName,
CredentialPath: credentialPath,
}
}
type schemaDropper interface {
Exec(context.Context, string, ...any) (pgconn.CommandTag, error)
}
func dropRegistrySchema(t *testing.T, ctx context.Context, db schemaDropper) {
t.Helper()
if _, err := db.Exec(ctx, `DROP SCHEMA IF EXISTS postgresql_tenant_operator CASCADE`); err != nil {
t.Fatalf("drop registry schema: %v", err)
}
}
+86
View File
@@ -0,0 +1,86 @@
/*
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
var bootstrapStatements = []string{
panxiao81 marked this conversation as resolved
Review

索引怎么说,两个ID复合够吗(应该够了吧

索引怎么说,两个ID复合够吗(应该够了吧
`CREATE SCHEMA IF NOT EXISTS postgresql_tenant_operator`,
`CREATE TABLE IF NOT EXISTS postgresql_tenant_operator.schema_version (
singleton boolean PRIMARY KEY DEFAULT true CHECK (singleton),
version integer NOT NULL
)`,
`INSERT INTO postgresql_tenant_operator.schema_version (singleton, version)
VALUES (true, 1)
ON CONFLICT (singleton) DO NOTHING`,
`CREATE TABLE IF NOT EXISTS 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)
)`,
}
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`
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`