Compare commits

..
Author SHA1 Message Date
sysopsandClaude Sonnet 5 b08f6a49fc TEN-07: migrations-orchestrierung-tenant-datenbanken
Orchestrator.RolloutAll wendet eine geordnete Liste von Migrationen auf JEDE
registrierte Tenant-Datenbank an (LoadMigrations liest *.up.sql aus einem
Verzeichnis nach der bestehenden 000N_name-Namenskonvention). Jeder Tenant
laeuft unabhaengig in eigener Verbindung — ein Fehlschlag bei einem Mandanten
bricht nur dessen eigenen Rollout ab (spaetere Migrationen bauen typischerweise
auf frueheren auf) und blockiert die uebrigen Tenants nicht.

schema_migrations-Tabelle pro Tenant-Datenbank (version PK, applied_at,
success, error) haelt den Stand pro Version einzeln nachvollziehbar fest.
Bereits erfolgreiche Versionen werden bei einem erneuten Rollout uebersprungen
(isAlreadySuccessful-Check vor jeder Anwendung), fehlgeschlagene werden beim
naechsten Versuch automatisch erneut probiert (kein manuelles Zuruecksetzen
noetig) — ON CONFLICT DO UPDATE haelt jeweils nur den letzten Versuch fest.

Bewusst ohne Abhaengigkeit von TEN-06 (Router): Migrations-Rollouts sind
seltene Batch-Vorgaenge, ein kurzlebiger Pool pro Tenant und Lauf reicht.

Pruefungen (ausgefuehrt auf root@192.168.1.131, go build/vet/test PASS,
go test ./... mit -p 1 noetig da mehrere Pakete die geteilte Registry-Tabelle
auf derselben Postgres-Instanz nutzen — siehe scripts/reset-test-env.sh):
1. Rollout gegen 3 Test-Tenants, einer absichtlich inkompatibel (Tabellen-
   Konflikt bei Migration 2) — TestRolloutAll_IsolatesFailurePerTenant: die
   anderen beiden erhalten beide Migrationen, der inkompatible bekommt
   Migration 1 trotzdem, scheitert nur an Migration 2, faellt nicht die
   anderen um. PASS.
2. Migrationsstand-Abfrage liefert korrekten Stand pro Tenant —
   TestStatus_ReflectsPerTenantState: Version 1 success=true, Version 2
   success=false mit Fehlertext. PASS.
3. Wiederholter Rollout fuer fehlgeschlagene Migration moeglich, ohne bereits
   erfolgreiche erneut anzuwenden — TestRolloutAll_RetryDoesNotReapplySuccessful:
   Migration 1 nutzt bewusst kein IF NOT EXISTS, ein Reapply haette den
   zweiten Lauf scheitern lassen; zweiter Lauf ist fehlerfrei und wendet nur
   die zuvor fehlgeschlagene Version an. PASS.

Co-Authored-By: Claude Sonnet 5 <noreply@anthropic.com>
2026-08-27 19:05:47 +02:00
9 changed files with 451 additions and 505 deletions
+211
View File
@@ -0,0 +1,211 @@
// Package migrate implementiert Core TEN-07: das automatisierte, pro Tenant
// fehlerisolierte Ausrollen von SQL-Migrationen ueber alle registrierten
// Mandanten-Datenbanken (Modell C — jede Migration muss N-mal statt einmal
// laufen, siehe TEN-01).
package migrate
import (
"context"
"errors"
"fmt"
"os"
"path/filepath"
"sort"
"strings"
"time"
"github.com/jackc/pgx/v5"
"github.com/jackc/pgx/v5/pgxpool"
"gitea.perlbach24.de/scripte/nexarch/internal/tenant"
)
const upSuffix = ".up.sql"
type Migration struct {
Version string
SQL string
}
// LoadMigrations liest alle *.up.sql-Dateien aus dir und sortiert sie nach
// Dateiname (die bestehende Namenskonvention 0001_..., 0002_... aus TEN-01
// sorgt fuer eine stabile Reihenfolge).
func LoadMigrations(dir string) ([]Migration, error) {
entries, err := os.ReadDir(dir)
if err != nil {
return nil, fmt.Errorf("migrationsverzeichnis lesen: %w", err)
}
var names []string
for _, e := range entries {
if !e.IsDir() && strings.HasSuffix(e.Name(), upSuffix) {
names = append(names, e.Name())
}
}
sort.Strings(names)
migrations := make([]Migration, 0, len(names))
for _, name := range names {
content, err := os.ReadFile(filepath.Join(dir, name))
if err != nil {
return nil, fmt.Errorf("migration %q lesen: %w", name, err)
}
version := strings.TrimSuffix(name, upSuffix)
migrations = append(migrations, Migration{Version: version, SQL: string(content)})
}
return migrations, nil
}
// MigrationStatus ist der pro Tenant und Version nachvollziehbare Stand
// (Akzeptanzkriterium 3): welche Version, wann zuletzt versucht, Erfolg oder
// Fehler.
type MigrationStatus struct {
Version string
AppliedAt time.Time
Success bool
Error string
}
// TenantResult fasst das Ergebnis eines Rollout-Versuchs fuer EINEN Tenant
// zusammen — wird von Orchestrator.RolloutAll pro Tenant gesammelt, damit ein
// Fehlschlag bei einem Mandanten die anderen nicht blockiert (Akzeptanzkriterium 2).
type TenantResult struct {
TenantSlug string
Applied []string // erfolgreich in diesem Lauf angewendete Versionen
FailedAt string // Version, bei der abgebrochen wurde; leer wenn kein Fehlschlag
Err error
}
// Orchestrator rollt Migrationen ueber alle in der Registry gefuehrten
// Mandanten aus. Baut bewusst NICHT auf TEN-06 (Router) auf — Migrations-
// Rollouts sind seltene Batch-Vorgaenge, kein Hot-Path, ein kurzlebiger Pool
// pro Tenant und Lauf ist hier einfacher und unabhaengig von TEN-06 testbar.
type Orchestrator struct {
registry *tenant.Registry
}
func NewOrchestrator(registry *tenant.Registry) *Orchestrator {
return &Orchestrator{registry: registry}
}
const ensureTableSQL = `
CREATE TABLE IF NOT EXISTS schema_migrations (
version TEXT PRIMARY KEY,
applied_at TIMESTAMPTZ NOT NULL DEFAULT now(),
success BOOLEAN NOT NULL,
error TEXT
);`
// RolloutAll wendet migrations auf JEDE registrierte Tenant-Datenbank an.
// Jeder Tenant laeuft unabhaengig — ein Fehlschlag bei einem Mandanten wird
// im jeweiligen TenantResult festgehalten und blockiert die uebrigen nicht
// (Akzeptanzkriterium 1 + 2).
func (o *Orchestrator) RolloutAll(ctx context.Context, migrations []Migration) ([]TenantResult, error) {
tenants, err := o.registry.List(ctx)
if err != nil {
return nil, fmt.Errorf("tenants fuer rollout auflisten: %w", err)
}
results := make([]TenantResult, 0, len(tenants))
for _, t := range tenants {
results = append(results, o.rolloutForTenant(ctx, t, migrations))
}
return results, nil
}
func (o *Orchestrator) rolloutForTenant(ctx context.Context, t tenant.Tenant, migrations []Migration) TenantResult {
result := TenantResult{TenantSlug: t.Slug}
pool, err := pgxpool.New(ctx, t.DBDSN)
if err != nil {
result.Err = fmt.Errorf("verbindung zu tenant %q: %w", t.Slug, err)
return result
}
defer pool.Close()
if _, err := pool.Exec(ctx, ensureTableSQL); err != nil {
result.Err = fmt.Errorf("schema_migrations anlegen fuer tenant %q: %w", t.Slug, err)
return result
}
for _, m := range migrations {
alreadyApplied, err := isAlreadySuccessful(ctx, pool, m.Version)
if err != nil {
result.Err = fmt.Errorf("migrationsstand lesen fuer tenant %q, version %q: %w", t.Slug, m.Version, err)
return result
}
if alreadyApplied {
continue // Akzeptanzkriterium 3 / Pruefung 3: keine erneute Anwendung.
}
_, execErr := pool.Exec(ctx, m.SQL)
if execErr != nil {
recordAttempt(ctx, pool, m.Version, false, execErr.Error())
result.FailedAt = m.Version
result.Err = fmt.Errorf("migration %q fuer tenant %q fehlgeschlagen: %w", m.Version, t.Slug, execErr)
return result // spaetere Migrationen bauen typischerweise auf dieser auf — Abbruch NUR fuer diesen Tenant.
}
recordAttempt(ctx, pool, m.Version, true, "")
result.Applied = append(result.Applied, m.Version)
}
return result
}
func isAlreadySuccessful(ctx context.Context, pool *pgxpool.Pool, version string) (bool, error) {
var success bool
err := pool.QueryRow(ctx, `SELECT success FROM schema_migrations WHERE version = $1`, version).Scan(&success)
if err != nil {
if errors.Is(err, pgx.ErrNoRows) {
return false, nil
}
return false, err
}
return success, nil
}
func recordAttempt(ctx context.Context, pool *pgxpool.Pool, version string, success bool, errMsg string) {
var errVal *string
if errMsg != "" {
errVal = &errMsg
}
_, _ = pool.Exec(ctx, `
INSERT INTO schema_migrations (version, applied_at, success, error)
VALUES ($1, now(), $2, $3)
ON CONFLICT (version) DO UPDATE SET applied_at = now(), success = $2, error = $3
`, version, success, errVal)
}
// Status liefert den Migrationsstand eines einzelnen Tenants (Akzeptanzkriterium 3).
func (o *Orchestrator) Status(ctx context.Context, tenantSlug string) ([]MigrationStatus, error) {
t, err := o.registry.GetBySlug(ctx, tenantSlug)
if err != nil {
return nil, fmt.Errorf("tenant %q nicht gefunden: %w", tenantSlug, err)
}
pool, err := pgxpool.New(ctx, t.DBDSN)
if err != nil {
return nil, fmt.Errorf("verbindung zu tenant %q: %w", tenantSlug, err)
}
defer pool.Close()
rows, err := pool.Query(ctx, `
SELECT version, applied_at, success, COALESCE(error, '')
FROM schema_migrations ORDER BY version
`)
if err != nil {
return nil, fmt.Errorf("migrationsstand abfragen: %w", err)
}
defer rows.Close()
var out []MigrationStatus
for rows.Next() {
var s MigrationStatus
if err := rows.Scan(&s.Version, &s.AppliedAt, &s.Success, &s.Error); err != nil {
return nil, fmt.Errorf("migrationsstand lesen: %w", err)
}
out = append(out, s)
}
return out, rows.Err()
}
+238
View File
@@ -0,0 +1,238 @@
package migrate
import (
"context"
"fmt"
"os"
"strings"
"testing"
"github.com/jackc/pgx/v5/pgxpool"
"gitea.perlbach24.de/scripte/nexarch/internal/tenant"
)
func TestLoadMigrations_SortsByFilename(t *testing.T) {
dir := t.TempDir()
writeFile(t, dir, "0002_second.up.sql", "SELECT 2;")
writeFile(t, dir, "0001_first.up.sql", "SELECT 1;")
writeFile(t, dir, "0001_first.down.sql", "SELECT 'ignored';") // muss ignoriert werden
migrations, err := LoadMigrations(dir)
if err != nil {
t.Fatalf("load: %v", err)
}
if len(migrations) != 2 {
t.Fatalf("erwartet 2 migrationen, habe %d", len(migrations))
}
if migrations[0].Version != "0001_first" || migrations[1].Version != "0002_second" {
t.Fatalf("unerwartete reihenfolge: %+v", migrations)
}
}
func writeFile(t *testing.T, dir, name, content string) {
t.Helper()
if err := os.WriteFile(dir+"/"+name, []byte(content), 0o644); err != nil {
t.Fatalf("write %s: %v", name, err)
}
}
func setupOrchestratorTest(t *testing.T, tenantCount int) (*Orchestrator, []tenant.Tenant, *pgxpool.Pool, func()) {
t.Helper()
adminDSN := os.Getenv("TEST_ADMIN_DSN")
if adminDSN == "" {
t.Skip("TEST_ADMIN_DSN nicht gesetzt, Integrationstest uebersprungen")
}
ctx := context.Background()
adminPool, err := pgxpool.New(ctx, adminDSN)
if err != nil {
t.Fatalf("admin pool: %v", err)
}
registryPool, err := pgxpool.New(ctx, adminDSN)
if err != nil {
t.Fatalf("registry pool: %v", err)
}
if _, err := registryPool.Exec(ctx, `
CREATE TABLE IF NOT EXISTS tenants (
id UUID PRIMARY KEY DEFAULT gen_random_uuid(),
slug TEXT NOT NULL UNIQUE,
name TEXT NOT NULL,
db_name TEXT NOT NULL UNIQUE,
db_dsn TEXT NOT NULL,
status TEXT NOT NULL DEFAULT 'active',
created_at TIMESTAMPTZ NOT NULL DEFAULT now()
)`); err != nil {
t.Fatalf("registry-schema: %v", err)
}
registry := tenant.NewRegistry(registryPool)
dsnTemplate := strings.Replace(adminDSN, "/postgres?", "/%s?", 1)
provisioner := tenant.NewProvisioner(adminPool, registry, dsnTemplate)
var tenants []tenant.Tenant
var slugs []string
for i := 0; i < tenantCount; i++ {
slug := fmt.Sprintf("mig_t%d", i)
slugs = append(slugs, slug)
tn, err := provisioner.Provision(ctx, slug, slug)
if err != nil {
t.Fatalf("provision %s: %v", slug, err)
}
tenants = append(tenants, tn)
}
orchestrator := NewOrchestrator(registry)
cleanup := func() {
for _, slug := range slugs {
_, _ = adminPool.Exec(ctx, fmt.Sprintf(`DROP DATABASE IF EXISTS %q`, "tenant_"+slug))
}
_, _ = registryPool.Exec(ctx, `DELETE FROM tenants WHERE slug = ANY($1)`, slugs)
registryPool.Close()
adminPool.Close()
}
return orchestrator, tenants, adminPool, cleanup
}
// Akzeptanzkriterium 1 + 2 + Pruefung 1: Rollout gegen mehrere Tenants, einer
// davon absichtlich inkompatibel — die anderen laufen trotzdem durch.
func TestRolloutAll_IsolatesFailurePerTenant(t *testing.T) {
orchestrator, tenants, _, cleanup := setupOrchestratorTest(t, 3)
defer cleanup()
ctx := context.Background()
migrations := []Migration{
{Version: "0001_demo_a", SQL: "CREATE TABLE demo_a (id INT);"},
{Version: "0002_demo_b", SQL: "CREATE TABLE demo_b (id INT);"},
}
// tenants[1] absichtlich inkompatibel machen: demo_b existiert schon,
// migration 0002 schlaegt dort mit "already exists" fehl.
badTenant := tenants[1]
badPool, err := pgxpool.New(ctx, badTenant.DBDSN)
if err != nil {
t.Fatalf("connect bad tenant: %v", err)
}
if _, err := badPool.Exec(ctx, "CREATE TABLE demo_b (id INT);"); err != nil {
t.Fatalf("inkompatiblen zustand vorbereiten: %v", err)
}
badPool.Close()
results, err := orchestrator.RolloutAll(ctx, migrations)
if err != nil {
t.Fatalf("rollout: %v", err)
}
if len(results) != 3 {
t.Fatalf("erwartet 3 ergebnisse, habe %d", len(results))
}
byslug := map[string]TenantResult{}
for _, r := range results {
byslug[r.TenantSlug] = r
}
good0 := byslug[tenants[0].Slug]
if good0.Err != nil || len(good0.Applied) != 2 {
t.Fatalf("tenant[0] sollte beide migrationen erhalten, habe %+v", good0)
}
good2 := byslug[tenants[2].Slug]
if good2.Err != nil || len(good2.Applied) != 2 {
t.Fatalf("tenant[2] sollte beide migrationen erhalten, habe %+v", good2)
}
bad := byslug[badTenant.Slug]
if bad.Err == nil {
t.Fatal("erwartet fehler fuer den inkompatiblen tenant")
}
if bad.FailedAt != "0002_demo_b" {
t.Fatalf("failedAt = %q, want 0002_demo_b", bad.FailedAt)
}
if len(bad.Applied) != 1 || bad.Applied[0] != "0001_demo_a" {
t.Fatalf("erwartet dass 0001_demo_a trotzdem erfolgreich war, habe %+v", bad.Applied)
}
}
// Akzeptanzkriterium 3 + Pruefung 2: Migrationsstand-Abfrage liefert fuer
// jeden Tenant den korrekten, unabhaengigen Stand.
func TestStatus_ReflectsPerTenantState(t *testing.T) {
orchestrator, tenants, _, cleanup := setupOrchestratorTest(t, 1)
defer cleanup()
ctx := context.Background()
migrations := []Migration{
{Version: "0001_ok", SQL: "CREATE TABLE ok_table (id INT);"},
{Version: "0002_fail", SQL: "SELECT this_column_does_not_exist FROM ok_table;"},
}
if _, err := orchestrator.RolloutAll(ctx, migrations); err != nil {
t.Fatalf("rollout: %v", err)
}
status, err := orchestrator.Status(ctx, tenants[0].Slug)
if err != nil {
t.Fatalf("status: %v", err)
}
if len(status) != 2 {
t.Fatalf("erwartet 2 status-eintraege, habe %d", len(status))
}
if !status[0].Success || status[0].Version != "0001_ok" {
t.Fatalf("status[0] unerwartet: %+v", status[0])
}
if status[1].Success || status[1].Version != "0002_fail" || status[1].Error == "" {
t.Fatalf("status[1] sollte fehlgeschlagen sein mit fehlertext: %+v", status[1])
}
}
// Pruefung 3: wiederholter Rollout-Versuch wendet bereits erfolgreiche
// Migrationen NICHT erneut an und kann die zuvor fehlgeschlagene nachholen,
// sobald die Ursache behoben ist.
func TestRolloutAll_RetryDoesNotReapplySuccessful(t *testing.T) {
orchestrator, tenants, _, cleanup := setupOrchestratorTest(t, 1)
defer cleanup()
ctx := context.Background()
// 0002 schlaegt beim ersten Versuch fehl, weil demo_conflict schon
// existiert (wir legen sie vorher an, um den Fehlschlag zu erzwingen).
pool, err := pgxpool.New(ctx, tenants[0].DBDSN)
if err != nil {
t.Fatalf("connect: %v", err)
}
if _, err := pool.Exec(ctx, "CREATE TABLE demo_conflict (id INT);"); err != nil {
t.Fatalf("vorbedingung: %v", err)
}
migrations := []Migration{
{Version: "0001_ok", SQL: "CREATE TABLE demo_first (id INT);"}, // OHNE IF NOT EXISTS,
// damit ein erneutes Anwenden nachweislich fehlschlagen wuerde.
{Version: "0002_conflict", SQL: "CREATE TABLE demo_conflict (id INT);"},
}
firstRun, err := orchestrator.RolloutAll(ctx, migrations)
if err != nil {
t.Fatalf("rollout 1: %v", err)
}
if firstRun[0].FailedAt != "0002_conflict" {
t.Fatalf("erwartet fehlschlag bei 0002_conflict im ersten lauf, habe %+v", firstRun[0])
}
// Ursache beheben.
if _, err := pool.Exec(ctx, "DROP TABLE demo_conflict;"); err != nil {
t.Fatalf("ursache beheben: %v", err)
}
pool.Close()
secondRun, err := orchestrator.RolloutAll(ctx, migrations)
if err != nil {
t.Fatalf("rollout 2: %v", err)
}
// Waere 0001_ok erneut angewendet worden ("CREATE TABLE demo_first" ohne
// IF NOT EXISTS), haette das einen Fehler erzeugt statt eines sauberen
// Applied-Eintrags fuer 0002_conflict.
if secondRun[0].Err != nil {
t.Fatalf("zweiter lauf sollte fehlerfrei sein, habe %+v", secondRun[0])
}
if len(secondRun[0].Applied) != 1 || secondRun[0].Applied[0] != "0002_conflict" {
t.Fatalf("erwartet nur 0002_conflict im zweiten lauf angewendet, habe %+v", secondRun[0].Applied)
}
}
-206
View File
@@ -1,206 +0,0 @@
package tenant
import (
"context"
"errors"
"fmt"
"log/slog"
"time"
"github.com/jackc/pgx/v5"
"github.com/jackc/pgx/v5/pgxpool"
)
var (
ErrTenantNotFound = errors.New("tenant: nicht gefunden")
ErrInvalidTransition = errors.New("tenant: ungueltiger zustandsuebergang")
// ErrTenantNotActive wird von Lifecycle.CheckActive verwendet — bewusst
// EIN Fehler fuer suspendiert/zur-Loeschung-vorgemerkt/geloescht, da der
// Aufrufer (z.B. Login) nur wissen muss "kein Zugriff", nicht welcher der
// Nicht-aktiv-Zustaende genau vorliegt.
ErrTenantNotActive = errors.New("tenant: nicht aktiv")
)
func scanTenantWithLifecycle(row pgx.Row) (Tenant, error) {
var t Tenant
if err := row.Scan(&t.ID, &t.Slug, &t.Name, &t.DBName, &t.DBDSN, &t.Status,
&t.CreatedAt, &t.PreviousStatus, &t.DeletionScheduledAt); err != nil {
return Tenant{}, err
}
return t, nil
}
// transition fuehrt einen bewachten Zustandsuebergang aus: das UPDATE greift
// nur, wenn der aktuelle Status einer von allowedFrom ist (atomarer
// Check-and-Set, kein Race zwischen Lesen und Schreiben). Greift es nicht,
// wird zwischen "Tenant existiert nicht" und "Uebergang nicht erlaubt"
// unterschieden, damit AC1 ("ungueltige Uebergaenge werden abgewiesen") einen
// sprechenden Fehler liefert statt eines stillen No-Ops.
func (r *Registry) transition(ctx context.Context, slug string, allowedFrom []Status, to Status, previousStatus *string, deletionAt *time.Time) (Tenant, error) {
from := make([]string, len(allowedFrom))
for i, s := range allowedFrom {
from[i] = string(s)
}
row := r.pool.QueryRow(ctx, `
UPDATE tenants
SET status = $2, previous_status = $3, deletion_scheduled_at = $4
WHERE slug = $1 AND status = ANY($5)
RETURNING id, slug, name, db_name, db_dsn, status, created_at, previous_status, deletion_scheduled_at
`, slug, string(to), previousStatus, deletionAt, from)
t, err := scanTenantWithLifecycle(row)
if err == nil {
return t, nil
}
if !errors.Is(err, pgx.ErrNoRows) {
return Tenant{}, fmt.Errorf("zustandsuebergang: %w", err)
}
existing, getErr := r.GetBySlug(ctx, slug)
if getErr != nil {
return Tenant{}, ErrTenantNotFound
}
return Tenant{}, fmt.Errorf("%w: von %q nach %q (aktuell: %q)", ErrInvalidTransition, allowedFrom, to, existing.Status)
}
// Suspend haelt die Daten des Mandanten unveraendert, sperrt aber den Zugriff
// (Akzeptanzkriterium 1) — es findet keine Loeschung/Migration statt.
func (r *Registry) Suspend(ctx context.Context, slug string) (Tenant, error) {
return r.transition(ctx, slug, []Status{StatusActive}, StatusSuspended, nil, nil)
}
// Reactivate stellt den Zustand vor der Suspendierung vollstaendig wieder her
// (Akzeptanzkriterium 2) — da Suspend keine weiteren Daten veraendert, genuegt
// die Rueckkehr nach StatusActive.
func (r *Registry) Reactivate(ctx context.Context, slug string) (Tenant, error) {
return r.transition(ctx, slug, []Status{StatusSuspended}, StatusActive, nil, nil)
}
// ScheduleDeletion merkt den Mandanten zur Loeschung vor und startet die
// Karenzzeit (Akzeptanzkriterium 3). previous_status wird festgehalten, damit
// CancelDeletion exakt dorthin zurueckkehren kann (aktiv ODER suspendiert).
func (r *Registry) ScheduleDeletion(ctx context.Context, slug string, grace time.Duration) (Tenant, error) {
existing, err := r.GetBySlug(ctx, slug)
if err != nil {
return Tenant{}, ErrTenantNotFound
}
prev := string(existing.Status)
deletionAt := time.Now().Add(grace)
return r.transition(ctx, slug, []Status{StatusActive, StatusSuspended}, StatusPendingDeletion, &prev, &deletionAt)
}
// CancelDeletion widerruft eine Loeschvormerkung innerhalb der Karenzzeit und
// stellt exakt den zuvor gesicherten Zustand wieder her.
func (r *Registry) CancelDeletion(ctx context.Context, slug string) (Tenant, error) {
existing, err := r.GetBySlug(ctx, slug)
if err != nil {
return Tenant{}, ErrTenantNotFound
}
if existing.Status != StatusPendingDeletion || existing.PreviousStatus == nil {
return Tenant{}, fmt.Errorf("%w: von %q nach aktiv/suspendiert (aktuell: %q)", ErrInvalidTransition, StatusPendingDeletion, existing.Status)
}
restoreTo := Status(*existing.PreviousStatus)
return r.transition(ctx, slug, []Status{StatusPendingDeletion}, restoreTo, nil, nil)
}
// Lifecycle fuehrt die tatsaechliche, physische Loeschung nach Ablauf der
// Karenzzeit aus (Datenbank-Drop) und stellt die Zugriffsschutz-Pruefung
// bereit. Getrennt von Registry, weil hierfuer zusaetzlich der adminPool
// (fuer DROP DATABASE) noetig ist, siehe internal/tenant.Provisioner.
type Lifecycle struct {
registry *Registry
adminPool *pgxpool.Pool
}
func NewLifecycle(registry *Registry, adminPool *pgxpool.Pool) *Lifecycle {
return &Lifecycle{registry: registry, adminPool: adminPool}
}
// CheckActive verweigert Zugriff fuer jeden Nicht-aktiv-Zustand und loggt den
// Vorgang strukturiert (Akzeptanzkriterium 1 / Pruefung 2).
func (l *Lifecycle) CheckActive(ctx context.Context, slug string) error {
t, err := l.registry.GetBySlug(ctx, slug)
if err != nil {
return ErrTenantNotFound
}
if t.Status != StatusActive {
slog.Warn("zugriff auf nicht-aktiven mandanten verweigert",
"tenant_slug", slug, "tenant_status", t.Status)
return ErrTenantNotActive
}
return nil
}
// ProcessDueDeletions loescht alle Mandanten-Datenbanken, deren Karenzzeit
// abgelaufen ist (Akzeptanzkriterium 3 / Pruefung 3). FOR UPDATE SKIP LOCKED
// folgt der projektweiten Postgres-Jobqueue-Konvention (siehe
// SKALIERUNGSKONZEPT.md) und macht die Funktion sicher fuer mehrere parallel
// laufende Core-Instanzen.
func (l *Lifecycle) ProcessDueDeletions(ctx context.Context) (int, error) {
tx, err := l.registry.pool.Begin(ctx)
if err != nil {
return 0, fmt.Errorf("sweep-transaktion starten: %w", err)
}
defer func() { _ = tx.Rollback(ctx) }()
rows, err := tx.Query(ctx, `
SELECT id, db_name FROM tenants
WHERE status = $1 AND deletion_scheduled_at <= now()
FOR UPDATE SKIP LOCKED
`, string(StatusPendingDeletion))
if err != nil {
return 0, fmt.Errorf("faellige loeschungen abfragen: %w", err)
}
type due struct{ id, dbName string }
var candidates []due
for rows.Next() {
var d due
if err := rows.Scan(&d.id, &d.dbName); err != nil {
rows.Close()
return 0, fmt.Errorf("faellige loeschung lesen: %w", err)
}
candidates = append(candidates, d)
}
rows.Close()
if err := rows.Err(); err != nil {
return 0, err
}
processed := 0
for _, c := range candidates {
if _, err := l.adminPool.Exec(ctx, fmt.Sprintf(`DROP DATABASE IF EXISTS %q`, c.dbName)); err != nil {
return processed, fmt.Errorf("tenant-datenbank %q loeschen: %w", c.dbName, err)
}
if _, err := tx.Exec(ctx, `
UPDATE tenants SET status = $2, previous_status = NULL, deletion_scheduled_at = NULL
WHERE id = $1
`, c.id, string(StatusDeleted)); err != nil {
return processed, fmt.Errorf("tenant %q als geloescht markieren: %w", c.id, err)
}
processed++
}
if err := tx.Commit(ctx); err != nil {
return 0, fmt.Errorf("sweep-transaktion committen: %w", err)
}
return processed, nil
}
// RunSweeper triggert ProcessDueDeletions periodisch, bis ctx beendet wird —
// die "In-Prozess-Worker-Goroutine" aus der projektweiten Jobqueue-Konvention.
func (l *Lifecycle) RunSweeper(ctx context.Context, interval time.Duration) {
ticker := time.NewTicker(interval)
defer ticker.Stop()
for {
select {
case <-ctx.Done():
return
case <-ticker.C:
if _, err := l.ProcessDueDeletions(ctx); err != nil {
slog.Error("tenant-loeschung-sweep fehlgeschlagen", "error", err)
}
}
}
}
-253
View File
@@ -1,253 +0,0 @@
package tenant
import (
"context"
"errors"
"fmt"
"os"
"strings"
"testing"
"time"
"github.com/jackc/pgx/v5/pgxpool"
)
func newLifecycleTestSetup(t *testing.T) (*Registry, *Lifecycle, *pgxpool.Pool, func()) {
t.Helper()
adminDSN := os.Getenv("TEST_ADMIN_DSN")
if adminDSN == "" {
t.Skip("TEST_ADMIN_DSN nicht gesetzt, Integrationstest uebersprungen")
}
ctx := context.Background()
adminPool, err := pgxpool.New(ctx, adminDSN)
if err != nil {
t.Fatalf("admin pool: %v", err)
}
registryPool, err := pgxpool.New(ctx, adminDSN)
if err != nil {
t.Fatalf("registry pool: %v", err)
}
if _, err := registryPool.Exec(ctx, `
CREATE TABLE IF NOT EXISTS tenants (
id UUID PRIMARY KEY DEFAULT gen_random_uuid(),
slug TEXT NOT NULL UNIQUE,
name TEXT NOT NULL,
db_name TEXT NOT NULL UNIQUE,
db_dsn TEXT NOT NULL,
status TEXT NOT NULL DEFAULT 'active',
created_at TIMESTAMPTZ NOT NULL DEFAULT now(),
previous_status TEXT,
deletion_scheduled_at TIMESTAMPTZ
)`); err != nil {
t.Fatalf("registry-schema: %v", err)
}
registry := NewRegistry(registryPool)
dsnTemplate := strings.Replace(adminDSN, "/postgres?", "/%s?", 1)
provisioner := NewProvisioner(adminPool, registry, dsnTemplate)
lifecycle := NewLifecycle(registry, adminPool)
cleanup := func() {
registryPool.Close()
adminPool.Close()
}
_ = provisioner
return registry, lifecycle, adminPool, cleanup
}
func provisionTestTenant(t *testing.T, registry *Registry, adminPool *pgxpool.Pool, slug string) {
t.Helper()
dsnTemplate := strings.Replace(os.Getenv("TEST_ADMIN_DSN"), "/postgres?", "/%s?", 1)
provisioner := NewProvisioner(adminPool, registry, dsnTemplate)
if _, err := provisioner.Provision(context.Background(), slug, slug); err != nil {
t.Fatalf("provision %s: %v", slug, err)
}
t.Cleanup(func() {
ctx := context.Background()
_, _ = adminPool.Exec(ctx, fmt.Sprintf(`DROP DATABASE IF EXISTS %q`, dbNameForSlug(slug)))
_, _ = registry.pool.Exec(ctx, `DELETE FROM tenants WHERE slug = $1`, slug)
})
}
// Akzeptanzkriterium 1 (Suspend) + 2 (Reactivate) + Pruefung 1 (Uebergaenge).
func TestLifecycle_SuspendAndReactivate(t *testing.T) {
registry, _, adminPool, cleanup := newLifecycleTestSetup(t)
defer cleanup()
provisionTestTenant(t, registry, adminPool, "lc_suspend")
ctx := context.Background()
suspended, err := registry.Suspend(ctx, "lc_suspend")
if err != nil {
t.Fatalf("suspend: %v", err)
}
if suspended.Status != StatusSuspended {
t.Fatalf("status = %q, want suspended", suspended.Status)
}
reactivated, err := registry.Reactivate(ctx, "lc_suspend")
if err != nil {
t.Fatalf("reactivate: %v", err)
}
if reactivated.Status != StatusActive {
t.Fatalf("status = %q, want active", reactivated.Status)
}
}
// Pruefung 1: ungueltige Uebergaenge werden abgewiesen.
func TestLifecycle_RejectsInvalidTransitions(t *testing.T) {
registry, _, adminPool, cleanup := newLifecycleTestSetup(t)
defer cleanup()
provisionTestTenant(t, registry, adminPool, "lc_invalid")
ctx := context.Background()
// Reactivate auf einem bereits aktiven Tenant ist kein gueltiger Uebergang.
if _, err := registry.Reactivate(ctx, "lc_invalid"); !errors.Is(err, ErrInvalidTransition) {
t.Fatalf("erwartet ErrInvalidTransition, habe %v", err)
}
if _, err := registry.Suspend(ctx, "lc_invalid"); err != nil {
t.Fatalf("suspend: %v", err)
}
// Suspend auf einem bereits suspendierten Tenant ist ebenfalls ungueltig.
if _, err := registry.Suspend(ctx, "lc_invalid"); !errors.Is(err, ErrInvalidTransition) {
t.Fatalf("erwartet ErrInvalidTransition, habe %v", err)
}
// CancelDeletion ohne vorherige Loeschvormerkung ist ungueltig.
if _, err := registry.CancelDeletion(ctx, "lc_invalid"); !errors.Is(err, ErrInvalidTransition) {
t.Fatalf("erwartet ErrInvalidTransition, habe %v", err)
}
if _, err := registry.Suspend(ctx, "unbekannter-slug-xyz"); !errors.Is(err, ErrTenantNotFound) {
t.Fatalf("erwartet ErrTenantNotFound, habe %v", err)
}
}
// Akzeptanzkriterium 3: Loeschung zweistufig mit Karenzzeit, innerhalb der
// Frist widerrufbar — sowohl aus 'active' als auch aus 'suspended' heraus,
// mit exakter Wiederherstellung des jeweiligen Vorzustands.
func TestLifecycle_ScheduleAndCancelDeletion_RestoresExactPreviousState(t *testing.T) {
registry, _, adminPool, cleanup := newLifecycleTestSetup(t)
defer cleanup()
provisionTestTenant(t, registry, adminPool, "lc_cancel_active")
provisionTestTenant(t, registry, adminPool, "lc_cancel_suspended")
ctx := context.Background()
// Fall 1: aus 'active' heraus vorgemerkt und widerrufen.
scheduled, err := registry.ScheduleDeletion(ctx, "lc_cancel_active", time.Hour)
if err != nil {
t.Fatalf("schedule deletion: %v", err)
}
if scheduled.Status != StatusPendingDeletion {
t.Fatalf("status = %q, want pending_deletion", scheduled.Status)
}
if scheduled.DeletionScheduledAt == nil {
t.Fatal("erwartet gesetzte deletion_scheduled_at")
}
restored, err := registry.CancelDeletion(ctx, "lc_cancel_active")
if err != nil {
t.Fatalf("cancel deletion: %v", err)
}
if restored.Status != StatusActive {
t.Fatalf("status = %q, want active (vorheriger zustand)", restored.Status)
}
// Fall 2: aus 'suspended' heraus vorgemerkt und widerrufen — muss zu
// 'suspended' zurueckkehren, NICHT zu 'active'.
if _, err := registry.Suspend(ctx, "lc_cancel_suspended"); err != nil {
t.Fatalf("suspend: %v", err)
}
if _, err := registry.ScheduleDeletion(ctx, "lc_cancel_suspended", time.Hour); err != nil {
t.Fatalf("schedule deletion: %v", err)
}
restoredSuspended, err := registry.CancelDeletion(ctx, "lc_cancel_suspended")
if err != nil {
t.Fatalf("cancel deletion: %v", err)
}
if restoredSuspended.Status != StatusSuspended {
t.Fatalf("status = %q, want suspended (vorheriger zustand)", restoredSuspended.Status)
}
}
// Akzeptanzkriterium 1 + Pruefung 2: suspendierter Tenant erzeugt bei jedem
// Zugriffsversuch einen klaren Fehler.
func TestLifecycle_CheckActive_RejectsNonActive(t *testing.T) {
registry, lifecycle, adminPool, cleanup := newLifecycleTestSetup(t)
defer cleanup()
provisionTestTenant(t, registry, adminPool, "lc_checkactive")
ctx := context.Background()
if err := lifecycle.CheckActive(ctx, "lc_checkactive"); err != nil {
t.Fatalf("aktiver tenant sollte durchgehen, habe %v", err)
}
if _, err := registry.Suspend(ctx, "lc_checkactive"); err != nil {
t.Fatalf("suspend: %v", err)
}
for i := 0; i < 3; i++ {
if err := lifecycle.CheckActive(ctx, "lc_checkactive"); !errors.Is(err, ErrTenantNotActive) {
t.Fatalf("versuch %d: erwartet ErrTenantNotActive, habe %v", i, err)
}
}
if err := lifecycle.CheckActive(ctx, "nie-registriert"); !errors.Is(err, ErrTenantNotFound) {
t.Fatalf("erwartet ErrTenantNotFound, habe %v", err)
}
}
// Akzeptanzkriterium 3 + Pruefung 3: Loeschvorgang nach Ablauf der Karenzzeit
// automatisch ausgeloest (hier durch direkten Aufruf von ProcessDueDeletions,
// das RunSweeper periodisch aufruft).
func TestLifecycle_ProcessDueDeletions(t *testing.T) {
registry, lifecycle, adminPool, cleanup := newLifecycleTestSetup(t)
defer cleanup()
provisionTestTenant(t, registry, adminPool, "lc_due")
provisionTestTenant(t, registry, adminPool, "lc_not_due")
ctx := context.Background()
// lc_due: Karenzzeit liegt bereits in der Vergangenheit -> faellig.
if _, err := registry.ScheduleDeletion(ctx, "lc_due", -time.Minute); err != nil {
t.Fatalf("schedule deletion (due): %v", err)
}
// lc_not_due: Karenzzeit liegt weit in der Zukunft -> nicht faellig.
if _, err := registry.ScheduleDeletion(ctx, "lc_not_due", time.Hour); err != nil {
t.Fatalf("schedule deletion (not due): %v", err)
}
processed, err := lifecycle.ProcessDueDeletions(ctx)
if err != nil {
t.Fatalf("process due deletions: %v", err)
}
if processed != 1 {
t.Fatalf("erwartet genau 1 verarbeitete loeschung, habe %d", processed)
}
due, err := registry.GetBySlug(ctx, "lc_due")
if err != nil {
t.Fatalf("get lc_due: %v", err)
}
if due.Status != StatusDeleted {
t.Fatalf("lc_due status = %q, want deleted", due.Status)
}
notDue, err := registry.GetBySlug(ctx, "lc_not_due")
if err != nil {
t.Fatalf("get lc_not_due: %v", err)
}
if notDue.Status != StatusPendingDeletion {
t.Fatalf("lc_not_due status = %q, want pending_deletion (noch nicht faellig)", notDue.Status)
}
// Datenbank von lc_due wurde tatsaechlich physisch entfernt.
var exists bool
if err := adminPool.QueryRow(ctx, `SELECT EXISTS(SELECT 1 FROM pg_database WHERE datname = $1)`,
dbNameForSlug("lc_due")).Scan(&exists); err != nil {
t.Fatalf("pg_database pruefen: %v", err)
}
if exists {
t.Fatal("erwartet, dass die tenant-datenbank von lc_due geloescht wurde")
}
}
+2 -6
View File
@@ -35,17 +35,13 @@ func (r *Registry) insertTx(ctx context.Context, tx pgx.Tx, t Tenant) (Tenant, e
}
func (r *Registry) GetBySlug(ctx context.Context, slug string) (Tenant, error) {
// previous_status/deletion_scheduled_at werden mitgelesen, damit TEN-04
// (internal/tenant/lifecycle.go) den vollstaendigen Lebenszyklus-Zustand
// ueber GetBySlug ansehen kann, statt eine eigene Abfrage zu duplizieren.
var t Tenant
row := r.pool.QueryRow(ctx, `
SELECT id, slug, name, db_name, db_dsn, status, created_at, previous_status, deletion_scheduled_at
SELECT id, slug, name, db_name, db_dsn, status, created_at
FROM tenants WHERE slug = $1
`, slug)
if err := row.Scan(&t.ID, &t.Slug, &t.Name, &t.DBName, &t.DBDSN, &t.Status, &t.CreatedAt,
&t.PreviousStatus, &t.DeletionScheduledAt); err != nil {
if err := row.Scan(&t.ID, &t.Slug, &t.Name, &t.DBName, &t.DBDSN, &t.Status, &t.CreatedAt); err != nil {
return Tenant{}, fmt.Errorf("tenant laden: %w", err)
}
return t, nil
-9
View File
@@ -12,10 +12,6 @@ type Status string
const (
StatusActive Status = "active"
// Lebenszyklus-Zustaende aus TEN-04 (siehe internal/tenant/lifecycle.go).
StatusSuspended Status = "suspended"
StatusPendingDeletion Status = "pending_deletion"
StatusDeleted Status = "deleted"
)
type Tenant struct {
@@ -26,11 +22,6 @@ type Tenant struct {
DBDSN string
Status Status
CreatedAt time.Time
// PreviousStatus und DeletionScheduledAt sind nur waehrend
// StatusPendingDeletion gesetzt (TEN-04) — sie halten fest, in welchen
// Zustand CancelDeletion zurueckkehrt und wann die Karenzzeit ablaeuft.
PreviousStatus *string
DeletionScheduledAt *time.Time
}
// slugPattern erzwingt sichere, als SQL-Identifier verwendbare Slugs, damit
@@ -1,2 +0,0 @@
ALTER TABLE tenants DROP COLUMN previous_status;
ALTER TABLE tenants DROP COLUMN deletion_scheduled_at;
-6
View File
@@ -1,6 +0,0 @@
-- Lebenszyklus-Zustaende fuer Mandanten (TEN-04, siehe core-kanban/tickets/TEN-04.md).
-- previous_status haelt den Zustand VOR einer Loeschvormerkung, damit
-- CancelDeletion "den vorherigen Zustand vollstaendig wiederherstellt"
-- (aktiv ODER suspendiert), statt hart auf 'active' zurueckzusetzen.
ALTER TABLE tenants ADD COLUMN previous_status TEXT;
ALTER TABLE tenants ADD COLUMN deletion_scheduled_at TIMESTAMPTZ;
-23
View File
@@ -1,23 +0,0 @@
#!/usr/bin/env bash
# Setzt die nexarch-Testumgebung zurueck: loescht die geteilte
# Registry-Tabelle "tenants" in der postgres-Wartungsdatenbank sowie alle
# tenant_*-Datenbanken. Noetig, weil verschiedene Feature-Branches
# unterschiedliche Registry-Schemata erwarten, aber dieselbe physische
# Postgres-Instanz auf dem Testhost teilen (siehe [[project-nexarch-test-infra]]).
#
# Aufruf: NEXARCH_TEST_DB_PASSWORD=... ./scripts/reset-test-env.sh
set -euo pipefail
PASS="${NEXARCH_TEST_DB_PASSWORD:?Setze NEXARCH_TEST_DB_PASSWORD vor dem Aufruf}"
ROLE="nexarch_test"
export PGPASSWORD="$PASS"
psql -h localhost -U "$ROLE" -d postgres -v ON_ERROR_STOP=1 -c "DROP TABLE IF EXISTS tenants;"
dbs=$(psql -h localhost -U "$ROLE" -d postgres -tAc "SELECT datname FROM pg_database WHERE datname LIKE 'tenant\_%' ESCAPE '\'")
for db in $dbs; do
psql -h localhost -U "$ROLE" -d postgres -v ON_ERROR_STOP=1 -c "DROP DATABASE IF EXISTS \"${db}\";"
done
echo "Testumgebung zurueckgesetzt: registry-tabelle + $(echo "$dbs" | grep -c . || true) tenant-datenbank(en) entfernt."