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, &t.RetentionBlockReason, &t.RetentionCheckedAt); 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, retention_block_reason, retention_checked_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 // retention ist die Pruef-Schnittstelle gegen Archive RET-03/CMP-06 (TEN-08). // Default NoRetentionCheck{}, bis Archive angebunden ist — siehe retention.go. retention RetentionChecker } func NewLifecycle(registry *Registry, adminPool *pgxpool.Pool) *Lifecycle { return &Lifecycle{registry: registry, adminPool: adminPool, retention: NoRetentionCheck{}} } // WithRetentionChecker ersetzt den Retention-Checker (z.B. im Test durch einen // Fake, oder in Produktion durch den echten Archive-RET-03-Client). Gibt // dasselbe *Lifecycle zurueck, um Verkettung beim Aufbau zu erlauben. func (l *Lifecycle) WithRetentionChecker(checker RetentionChecker) *Lifecycle { l.retention = checker return l } // 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, slug, 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, slug, dbName string } var candidates []due for rows.Next() { var d due if err := rows.Scan(&d.id, &d.slug, &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 { // TEN-08: vor der physischen Loeschung gegen Archive RET-03/CMP-06 pruefen. // Solange eine Sperre besteht, bleibt der Tenant in pending_deletion // ("zur Loeschung vorgemerkt, aber gesperrt") — der Grund wird // festgehalten (Akzeptanzkriterium 2), die naechste Sweeper-Runde // prueft automatisch erneut (Akzeptanzkriterium 3), ohne dass ein // manueller Re-Trigger noetig waere. result, err := l.retention.CheckTenantRetention(ctx, c.id) if err != nil { return processed, fmt.Errorf("retention-pruefung fuer tenant %q: %w", c.id, err) } if result.Blocked { slog.Warn("tenant-loeschung wegen aufbewahrungspflicht/legal-hold zurueckgehalten", "tenant_slug", c.slug, "reason", result.Reason) if _, err := tx.Exec(ctx, ` UPDATE tenants SET retention_block_reason = $2, retention_checked_at = now() WHERE id = $1 `, c.id, result.Reason); err != nil { return processed, fmt.Errorf("retention-sperrgrund fuer tenant %q speichern: %w", c.id, err) } continue } 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, retention_block_reason = NULL, retention_checked_at = now() 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) } } } }