Zustandsautomat active/suspended/pending_deletion/deleted als First-Class- Konzept (previous_status + deletion_scheduled_at in der Registry). Alle Uebergaenge in Registry.transition als atomarer Check-and-Set (UPDATE ... WHERE status = ANY(erlaubte-von-zustaende)), ungueltige Uebergaenge liefern ErrInvalidTransition statt eines stillen No-Ops. ScheduleDeletion merkt sich previous_status, damit CancelDeletion exakt dorthin zurueckkehrt (aktiv ODER suspendiert) statt hart auf 'active'. Lifecycle.ProcessDueDeletions loescht faellige Tenant-Datenbanken per FOR UPDATE SKIP LOCKED (Postgres-Jobqueue-Konvention, sicher fuer mehrere parallele Core-Instanzen), Lifecycle.RunSweeper triggert das periodisch per In-Prozess-Goroutine. Lifecycle.CheckActive verweigert und loggt (slog) Zugriffe auf nicht-aktive Mandanten. Bugfix nebenbei: Registry.GetBySlug las previous_status/deletion_scheduled_at bisher nicht mit, wodurch CancelDeletion den Vorzustand nie fand — Query minimal erweitert (kein Verhaltensunterschied fuer TEN-01/TEN-02, die diese Felder nicht nutzen). Neu: scripts/reset-test-env.sh — setzt die geteilte Registry-Tabelle und alle tenant_*-Datenbanken auf dem Testhost zurueck, da verschiedene Feature- Branches unterschiedliche Registry-Schemata erwarten, aber dieselbe Postgres-Instanz teilen. Pruefungen (ausgefuehrt auf root@192.168.1.131, go build/vet/test PASS): 1. Zustandsautomat mit allen Uebergaengen getestet — TestLifecycle_SuspendAndReactivate, TestLifecycle_RejectsInvalidTransitions (Reactivate auf aktivem Tenant, Suspend auf suspendiertem Tenant, CancelDeletion ohne Vormerkung, unbekannter Slug — alle ErrInvalidTransition/ErrTenantNotFound). PASS. 2. Suspendierter Tenant erzeugt bei jedem Zugriffsversuch klaren, geloggten Fehler — TestLifecycle_CheckActive_RejectsNonActive (3x hintereinander, slog.Warn nachweislich pro Aufruf). PASS. 3. Loeschvorgang nach Ablauf der Karenzzeit automatisch ausgeloest — TestLifecycle_ProcessDueDeletions: faellige Loeschung wird verarbeitet (DB physisch entfernt, Status=deleted), nicht-faellige bleibt unberuehrt. PASS. Co-Authored-By: Claude Sonnet 5 <noreply@anthropic.com>
207 lines
7.4 KiB
Go
207 lines
7.4 KiB
Go
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)
|
|
}
|
|
}
|
|
}
|
|
}
|