Compare commits

..
Author SHA1 Message Date
sysops 1e8a6c22cd API-12: kek-bezugsdienst-starten
- cmd/kek-api: startet internal/kek.Handler.TenantKEKHandler (API-10),
  tenantResolverAdapter bildet tenant.Registry.GetBySlug auf
  kek.TenantResolver ab (nur Signatur-Anpassung)
- reines Wiring, kein Diff an internal/kek|moduleregistry|tenant
  (verifiziert)
- real deployed auf 131, end-zu-ende bewiesen: echtes Modul
  registriert+provisioniert, echter Tenant-KEK erzeugt, curl gegen
  laufenden Dienst liefert exakt denselben KEK zurueck (byte-fuer-byte
  verglichen), falsches Credential -> 403
- reale Grant-Luecke gefunden und behoben (nexarch_core auf
  tenant_keks), verifiziert

Pruefungen siehe docs/API-12-PRUEFPROTOKOLL.md
2026-08-30 23:54:07 +02:00
sysops 46fccd9c09 API-10: test-fix — rotatemasterkey-assertion nur fuer eigene test-tenants pruefen (geteilte test-db) 2026-08-28 10:04:51 +02:00
sysops a631ac8770 API-10: internal/apiserver-port wieder entfernen (ungenutzt, zieht internal/auth als fehlende abhaengigkeit nach) 2026-08-28 10:03:41 +02:00
sysops bc2126f3b1 API-10: master-key-verwaltung-tenant-schluesselhierarchie-kms-anbindung (envelope encryption, isolierte tenant-keks, rotation) 2026-08-28 10:03:27 +02:00
sysopsandClaude Sonnet 5 b23cd1961f API-02: modul-registry-aktivierungspruefung
internal/moduleregistry: Registry.Register traegt Fachmodule mit Name,
Version und benoetigten Feature-Flags ein (Akzeptanzkriterium 1), fehlende
Pflichtangaben werden abgewiesen. IsActive kombiniert Registrierung + LIC-02
Feature-Flag-Auswertung (ALLE benoetigten Flags muessen fuer den Tenant
aktiv sein) — ein nicht registriertes Modul ist nie aktiv. List liefert alle
Module fuer Statusseite/Lizenzoberflaeche (Akzeptanzkriterium 3).

RequireActiveModule ist die zentrale Durchsetzungs-Middleware (Casbin-
Prinzip): weist Anfragen an ein deaktiviertes Modul ab, BEVOR der
Modul-Handler ueberhaupt aufgerufen wird (Akzeptanzkriterium 2) —
Pruefung per Test belegt, dass der Handler bei Deaktivierung nachweislich
nicht erreicht wird.

Service-Credentials (Akzeptanzkriterium 4): Registry.Provision stellt pro
Modul-Instanz Client-ID + Secret aus, gespeichert wird nur der SHA-256-Hash
des Secrets. Registry.Authenticate vergleicht timing-safe (dieselbe
subtle.ConstantTimeCompare-Referenzimplementierung wie AUD-02).
RequireServiceCredential-Middleware liest X-Client-Id/X-Client-Secret und
weist Aufrufe ohne gueltiges Credential mit 401 ab, bevor der Core-seitige
Endpunkt (z.B. Audit-Nachlieferung, Nutzungsmeldung) erreicht wird.

Pruefungen (ausgefuehrt auf root@192.168.1.131, go build/vet/test PASS):
1. Anfrage an deaktiviertes Modul nachweislich vor Modul-Logik abgewiesen —
   TestRequireActiveModule_BlocksBeforeHandler: handlerReached bleibt false
   bei 403, wird true erst nach Aktivierung bei 200. PASS.
2. Registrierung mit fehlenden Pflichtangaben abgewiesen —
   TestRegister_RejectsMissingFields (leerer Name, leere Version). PASS.
3. Registry-Abfrage liefert konsistente Daten nach Aktivierung/Deaktivierung —
   TestIsActive_ReflectsFlagStateConsistently: aus/an/aus-Zyklus, IsActive
   folgt dem Flag-Zustand korrekt. PASS.
4. Aufruf mit ungueltigem/fehlendem Service-Credential abgewiesen, mit
   gueltigem angenommen — TestRequireServiceCredential_RejectsInvalidAcceptsValid
   und TestProvisionAndAuthenticate (falsches Secret, unbekannte Client-ID,
   korrektes Credential). PASS.

Co-Authored-By: Claude Sonnet 5 <noreply@anthropic.com>
2026-08-27 21:17:50 +02:00
sysopsandClaude Sonnet 5 4a30345e07 LIC-02: feature-flag-service-je-tenant
internal/flag: Store (Verwaltung) + Service (Auswertung mit TTL-Cache,
Default 5s) — Unleash-Prinzip Flag-Verwaltung vs. Flag-Auswertung getrennt,
als Kernfunktion des Core-Dienstes selbst statt separater Infrastruktur.

evaluate() wendet drei Strategien in fester Reihenfolge an: global an/aus,
Tenant-Zielgruppe, deterministischer Prozentsatz-Rollout (FNV-Hash aus
Tenant+Key, stabil pro Tenant). IsEnabled liefert IMMER nur bool (kein
Fehlerwert) — ein nicht erreichbarer Flag-Dienst kann damit keinen
Aufrufer zum Absturz bringen: bei DB-Fehler wird der zuletzt bekannte
Cache-Stand verwendet, ohne jeglichen Stand faellt der Dienst sicher auf
false zurueck. Service.Invalidate erzwingt sofortiges Neuladen fuer den
Schreiber selbst, andere Instanzen sehen Aenderungen spaetestens nach der
TTL (Akzeptanzkriterium 3, kein Neustart noetig).

Bugfix waehrend Tests: Store.Set uebergab ein nil-TargetTenantSlugs-Slice
als SQL NULL statt leerem Array (NOT-NULL-Verletzung) — auf leeres Slice
normalisiert.

Akzeptanzkriterium 4 (Deaktivierung loescht keine Daten): dieses Paket
besitzt ausschliesslich die eigene feature_flags-Zeile, hat keinerlei
Code-Pfad, der Modul-Geschaeftsdaten anfassen koennte — Loeschung bleibt
strukturell der Archive-Retention-Engine vorbehalten.

Pruefungen (ausgefuehrt auf root@192.168.1.131, go build/vet/test PASS):
1. Cache-Invalidierungszeit automatisiert gemessen —
   TestService_CacheInvalidationTiming: Aenderung wirksam nach 153ms bei
   TTL=150ms (innerhalb Ziel+Toleranz), vorher nachweislich noch alter Stand. PASS.
2. Zielgruppen-Strategie liefert erwartete Auswertung —
   TestService_TargetTenantStrategy / TestEvaluate_TargetTenantStrategy. PASS.
3. Ausfall des Flag-Dienstes fuehrt zu dokumentiertem Fallback, kein Absturz —
   TestService_FallsBackOnStoreFailure (mit recover()-Absicherung): Fallback
   auf Cache-Stand bzw. sicheres false bei komplett unerreichbarer DB, geloggt. PASS.
4. Modul-Deaktivierung/Reaktivierung ohne Datenverlust — architektonisch durch
   fehlenden Code-Pfad sichergestellt (siehe oben), zusaetzlich durch
   TestService_InvalidateForcesImmediateRefresh (Toggle aus/an bleibt
   konsistent nachvollziehbar) mitabgedeckt. PASS.

Co-Authored-By: Claude Sonnet 5 <noreply@anthropic.com>
2026-08-27 19:19:59 +02:00
sysopsandClaude Sonnet 5 8da9c67d08 TEN-01: tenant-registry-datenbank-provisioning
Registry-DB (nur Tenant-Metadaten), Provisioning-Routine legt pro Mandant
eine physisch isolierte Postgres-DB an und registriert sie transaktional
(Rollback der DB bei fehlgeschlagener Registrierung). Schlanker HTTP-Handler
als Schnittstellen-Vorbereitung fuer API-01/TEN-02, kein eigenes REST-Grundgerüst.

Pruefungen:
1. Migration up/down geschrieben (0001_tenant_registry.{up,down}.sql) — nicht
   gegen echte DB ausgefuehrt, da auf dieser Maschine kein Go/Postgres-Test-
   Setup verfuegbar ist. Offen zur Ausfuehrung.
2. Integrationstest TestProvision_CreatesIsolatedDatabases geschrieben (zwei
   Mandanten, prueft unterschiedliche db_name und current_database()) —
   ebenfalls nicht ausgefuehrt, guarded per TEST_ADMIN_DSN env var. Offen.
3. Slug-Validierung (unit test TestValidateSlug) deckt SQL-Injection-Versuch
   im Datenbanknamen ab — ebenfalls nicht lokal ausgefuehrt, da kein Go
   Compiler auf dieser Maschine vorhanden ist. Offen.

Alle drei Pruefungen sind vorbereitet, aber NICHT durchgefuehrt worden —
zaehlen laut Vorgabe als offen bis auf einer Maschine mit Go+Postgres verifiziert.

Co-Authored-By: Claude Sonnet 5 <noreply@anthropic.com>
2026-08-27 17:40:35 +02:00
71 changed files with 2454 additions and 1986 deletions
+85
View File
@@ -44,3 +44,88 @@ Keine Commits in dieser Session.
Keine Änderungen ermittelbar.
---
## 2026-08-27 17:26 17:28 (1m)
**Beschreibung:** Claude Code Session
**Projekt:** code
### Commits
- c895a67 core: initial Go module skeleton (config, db pool, tenant registry migration)
### Geänderte Dateien
- .gitignore | 2 ++
- DEVLOG.md | 46 ++++++++++++++++++++++++++++++++++++++++++++++
- cmd/core/main.go | 33 +++++++++++++++++++++++++++++++++
- go.mod | 5 +++++
- internal/config/config.go | 29 +++++++++++++++++++++++++++++
- internal/db/db.go | 11 +++++++++++
- migrations/0001_tenant_registry.sql | 10 ++++++++++
---
## 2026-08-27 17:28 17:29 (1m)
**Beschreibung:** Claude Code Session
**Projekt:** code
### Commits
Keine Commits in dieser Session.
### Geänderte Dateien
- .gitignore | 2 ++
- DEVLOG.md | 46 ++++++++++++++++++++++++++++++++++++++++++++++
- cmd/core/main.go | 33 +++++++++++++++++++++++++++++++++
- go.mod | 5 +++++
- internal/config/config.go | 29 +++++++++++++++++++++++++++++
- internal/db/db.go | 11 +++++++++++
- migrations/0001_tenant_registry.sql | 10 ++++++++++
---
## 2026-08-27 17:31 17:31 (0m)
**Beschreibung:** Claude Code Session
**Projekt:** code
### Commits
Keine Commits in dieser Session.
### Geänderte Dateien
- .gitignore | 2 ++
- DEVLOG.md | 46 ++++++++++++++++++++++++++++++++++++++++++++++
- cmd/core/main.go | 33 +++++++++++++++++++++++++++++++++
- go.mod | 5 +++++
- internal/config/config.go | 29 +++++++++++++++++++++++++++++
- internal/db/db.go | 11 +++++++++++
- migrations/0001_tenant_registry.sql | 10 ++++++++++
---
## 2026-08-27 17:36 17:36 (0m)
**Beschreibung:** Claude Code Session
**Projekt:** code
### Commits
Keine Commits in dieser Session.
### Geänderte Dateien
- .gitignore | 2 ++
- DEVLOG.md | 46 ++++++++++++++++++++++++++++++++++++++++++++++
- cmd/core/main.go | 33 +++++++++++++++++++++++++++++++++
- go.mod | 5 +++++
- internal/config/config.go | 29 +++++++++++++++++++++++++++++
- internal/db/db.go | 11 +++++++++++
- migrations/0001_tenant_registry.sql | 10 ++++++++++
---
## 2026-08-27 17:36 17:37 (0m)
**Beschreibung:** Claude Code Session
**Projekt:** code
### Commits
Keine Commits in dieser Session.
### Geänderte Dateien
- .gitignore | 2 ++
- DEVLOG.md | 46 ++++++++++++++++++++++++++++++++++++++++++++++
- cmd/core/main.go | 33 +++++++++++++++++++++++++++++++++
- go.mod | 5 +++++
- internal/config/config.go | 29 +++++++++++++++++++++++++++++
- internal/db/db.go | 11 +++++++++++
- migrations/0001_tenant_registry.sql | 10 ++++++++++
---
-18
View File
@@ -1,18 +0,0 @@
version: "2"
run:
timeout: 3m
linters:
default: none
enable:
- govet
- staticcheck
- errcheck
- unused
- ineffassign
formatters:
enable:
- gofmt
- goimports
-19
View File
@@ -1,19 +0,0 @@
.PHONY: build lint fmt test check
build:
go build ./...
lint:
golangci-lint run ./...
fmt:
gofmt -l .
@test -z "$$(gofmt -l .)" || (echo "gofmt-Verstoesse gefunden, siehe oben" && exit 1)
test:
go test ./... -p 1 -count=1
check: build
go vet ./...
golangci-lint run ./...
go test ./... -p 1 -count=1
-62
View File
@@ -1,62 +0,0 @@
# NEXARCH Archive
Zentrales, modulübergreifendes Modul für Aufbewahrung, WORM, Compliance und
Backup. Dieses Verzeichnis enthält bisher `internal/backup` (BAK-01,
Datenbank-Backup-Strategie) — weitere Bausteine folgen ticketweise.
## BAK-01: Datenbank-Backup
`cmd/backup-cli` — Aufrufpunkt für systemd-Timer (siehe
`../deploy/systemd/nexarch-archive-backup-*.timer`):
```bash
export NEXARCH_BACKUP_PG_USER=nexarch_backup
export NEXARCH_BACKUP_PG_PASSWORD=...
export NEXARCH_BACKUP_DIR=/var/nexarch-archiv/backups/postgres # NICHT auf einem ephemeren Test-Dataset (siehe Betrieb)
export NEXARCH_BACKUP_KEEP_GENERATIONS=7 # optional, Default 7
backup-cli full # neue Vollsicherung + Verifikation
backup-cli incremental # inkrementelle Sicherung gegen die neueste Generation
backup-cli rotate # entfernt alle bis auf die neuesten N Generationen
```
Voraussetzung: die konfigurierte Postgres-Rolle braucht das
`REPLICATION`-Attribut (`pg_basebackup` nutzt eine
Replikationsverbindung), und `summarize_wal = on` muss serverseitig gesetzt
sein (PostgreSQL 17s natives inkrementelles Backup, keine WAL-Archivierung
nötig).
## BAK-02: Objekt-Storage-Backup
`cmd/objectbackup-cli` sichert einen lokalen Verzeichnisbaum (den
FDN-03-`LocalDriver`-Basisordner direkt, oder — für S3-gestützte
Deployments — einen vorgelagerten `rclone`-Spiegel) mit
[restic](https://restic.net) (Content-defined Chunking, verschlüsseltes
Repository, geprüftes Tooling statt Eigenbau):
```bash
export NEXARCH_OBJECTBACKUP_REPO_DIR=/var/nexarch-archiv/backups/objects
export NEXARCH_OBJECTBACKUP_PASSWORD=...
export NEXARCH_OBJECTBACKUP_KEEP_SNAPSHOTS=30 # optional, Default 7
objectbackup-cli backup /var/nexarch-objects # Sicherung + Verifikation
objectbackup-cli check # vollständiges Lesen aller Datenblöcke
objectbackup-cli rotate # restic forget --keep-last N --prune
```
## Betrieb: Backup-Zielverzeichnis
Backup-Ziele liegen unter `/var/nexarch-archiv/` (persistentes ZFS-Dataset,
`zfs/data/subvol-1131-disk-0` auf 192.168.1.131), NIEMALS unter
`/var/nexarch-test/` (ephemeres Dataset, wird von den `reset-test-env.sh`-
Skripten der anderen Module geleert). ZFS-seitige Snapshots/Replikation
dieses Datasets sind ein eigenständiges Infra-Runbook (siehe
`../../STORAGE-KONZEPT.md` Abschnitt 7), kein Ticket-Code — `zfs
dedup=on` bewusst NICHT setzen (hoher RAM-Bedarf), Deduplizierung läuft
ausschließlich App-seitig über restic.
## Prüfungen
```bash
make check # build + vet + lint + test, analog Core/DMS
```
-107
View File
@@ -1,107 +0,0 @@
// backup-cli ist der Aufrufpunkt für BAK-01, gedacht für systemd-Timer
// (siehe deploy/systemd/) — "automatisiert nach Zeitplan" (Akzeptanzkriterium
// 1) entsteht durch die Zeitplan-Definition im Timer-Unit, nicht durch
// einen eigenen In-Prozess-Scheduler (kein zusätzlicher Dauerprozess nötig,
// passt zur Produkt-DNA "kein Anwendungsserver mit unnötigem
// Ressourcenverbrauch").
package main
import (
"context"
"fmt"
"log"
"os"
"time"
"gitea.perlbach24.de/scripte/nexarch/archive/internal/backup"
)
func loadConfig() backup.Config {
cfg := backup.Config{
Host: os.Getenv("NEXARCH_BACKUP_PG_HOST"),
Port: os.Getenv("NEXARCH_BACKUP_PG_PORT"),
User: os.Getenv("NEXARCH_BACKUP_PG_USER"),
Password: os.Getenv("NEXARCH_BACKUP_PG_PASSWORD"),
BackupDir: os.Getenv("NEXARCH_BACKUP_DIR"),
}
if cfg.Host == "" {
cfg.Host = "localhost"
}
if cfg.Port == "" {
cfg.Port = "5432"
}
if cfg.User == "" || cfg.Password == "" || cfg.BackupDir == "" {
log.Fatal("NEXARCH_BACKUP_PG_USER, NEXARCH_BACKUP_PG_PASSWORD und NEXARCH_BACKUP_DIR muessen gesetzt sein")
}
return cfg
}
func latestManifest(backupDir string) (string, error) {
generations, err := backup.ListGenerations(backupDir)
if err != nil {
return "", err
}
if len(generations) == 0 {
return "", fmt.Errorf("keine vorhandene generation fuer inkrementelle sicherung gefunden - zuerst 'full' ausfuehren")
}
latest := generations[len(generations)-1]
full := backupDir + "/" + latest + "/" + backup.FullBackupDirName + "/" + backup.BackupManifestFile
if _, err := os.Stat(full); err == nil {
return full, nil
}
return "", fmt.Errorf("kein backup_manifest in der neuesten generation %q gefunden", latest)
}
func main() {
if len(os.Args) < 2 {
log.Fatal("aufruf: backup-cli <full|incremental|verify|rotate> [args]")
}
cfg := loadConfig()
ctx := context.Background()
switch os.Args[1] {
case "full":
genID := backup.NewGenerationID(time.Now())
manifest, err := backup.FullBackup(ctx, cfg, genID)
if err != nil {
log.Fatalf("vollsicherung fehlgeschlagen: %v", err)
}
dir := manifest[:len(manifest)-len("/"+backup.BackupManifestFile)]
if err := backup.Verify(dir); err != nil {
log.Fatalf("verifikation der vollsicherung fehlgeschlagen: %v", err)
}
fmt.Printf("vollsicherung %q erstellt und verifiziert: %s\n", genID, manifest)
case "incremental":
manifest, err := latestManifest(cfg.BackupDir)
if err != nil {
log.Fatal(err)
}
generations, _ := backup.ListGenerations(cfg.BackupDir)
genID := generations[len(generations)-1]
incID := backup.NewGenerationID(time.Now())
newManifest, err := backup.IncrementalBackup(ctx, cfg, genID, incID, manifest)
if err != nil {
log.Fatalf("inkrementelle sicherung fehlgeschlagen: %v", err)
}
dir := newManifest[:len(newManifest)-len("/"+backup.BackupManifestFile)]
if err := backup.Verify(dir); err != nil {
log.Fatalf("verifikation der inkrementellen sicherung fehlgeschlagen: %v", err)
}
fmt.Printf("inkrementelle sicherung %q erstellt und verifiziert: %s\n", incID, newManifest)
case "rotate":
keep := 7
if v := os.Getenv("NEXARCH_BACKUP_KEEP_GENERATIONS"); v != "" {
_, _ = fmt.Sscanf(v, "%d", &keep)
}
removed, err := backup.Rotate(cfg.BackupDir, keep)
if err != nil {
log.Fatalf("rotation fehlgeschlagen: %v", err)
}
fmt.Printf("rotation abgeschlossen, %d generation(en) entfernt: %v\n", len(removed), removed)
default:
log.Fatalf("unbekannter befehl %q", os.Args[1])
}
}
-75
View File
@@ -1,75 +0,0 @@
// objectbackup-cli ist der Aufrufpunkt für BAK-02, für systemd-Timer
// gedacht (siehe deploy/systemd/) — "automatisiert nach Zeitplan" entsteht
// durch die Timer-Definition, kein eigener Dauerprozess (dieselbe
// Begründung wie BAK-01 / cmd/backup-cli).
package main
import (
"context"
"fmt"
"log"
"os"
"strconv"
"gitea.perlbach24.de/scripte/nexarch/archive/internal/objectbackup"
)
func loadConfig() objectbackup.Config {
cfg := objectbackup.Config{
RepoDir: os.Getenv("NEXARCH_OBJECTBACKUP_REPO_DIR"),
Password: os.Getenv("NEXARCH_OBJECTBACKUP_PASSWORD"),
}
if cfg.RepoDir == "" || cfg.Password == "" {
log.Fatal("NEXARCH_OBJECTBACKUP_REPO_DIR und NEXARCH_OBJECTBACKUP_PASSWORD muessen gesetzt sein")
}
return cfg
}
func main() {
if len(os.Args) < 2 {
log.Fatal("aufruf: objectbackup-cli <backup <quellverzeichnis>|check|rotate>")
}
cfg := loadConfig()
ctx := context.Background()
if err := objectbackup.InitRepo(ctx, cfg); err != nil {
log.Fatalf("repository initialisieren: %v", err)
}
switch os.Args[1] {
case "backup":
if len(os.Args) < 3 {
log.Fatal("aufruf: objectbackup-cli backup <quellverzeichnis>")
}
summary, err := objectbackup.Backup(ctx, cfg, os.Args[2])
if err != nil {
log.Fatalf("sicherung fehlgeschlagen: %v", err)
}
if err := objectbackup.Check(ctx, cfg, false); err != nil {
log.Fatalf("verifikation nach sicherung fehlgeschlagen: %v", err)
}
fmt.Printf("sicherung %q erstellt und verifiziert (neu=%d geaendert=%d unveraendert=%d)\n",
summary.SnapshotID, summary.FilesNew, summary.FilesChanged, summary.FilesUnmodified)
case "check":
if err := objectbackup.Check(ctx, cfg, true); err != nil {
log.Fatalf("verifikation fehlgeschlagen: %v", err)
}
fmt.Println("verifikation (mit vollstaendigem lesen) erfolgreich")
case "rotate":
keep := 7
if v := os.Getenv("NEXARCH_OBJECTBACKUP_KEEP_SNAPSHOTS"); v != "" {
if n, err := strconv.Atoi(v); err == nil {
keep = n
}
}
if err := objectbackup.Forget(ctx, cfg, keep); err != nil {
log.Fatalf("rotation fehlgeschlagen: %v", err)
}
fmt.Println("rotation abgeschlossen")
default:
log.Fatalf("unbekannter befehl %q", os.Args[1])
}
}
-55
View File
@@ -1,55 +0,0 @@
// reconcile-cli ist der Aufrufpunkt für BAK-05, für systemd-Timer gedacht
// (siehe deploy/systemd/) — "geplanter Abgleichs-Job" (Ticket-Vorgabe)
// entsteht durch die Timer-Definition, kein eigener Dauerprozess.
package main
import (
"context"
"encoding/json"
"log"
"os"
"github.com/jackc/pgx/v5/pgxpool"
"gitea.perlbach24.de/scripte/nexarch/archive/internal/reconcile"
)
func main() {
dsn := os.Getenv("NEXARCH_RECONCILE_TENANT_DSN")
storageDir := os.Getenv("NEXARCH_RECONCILE_STORAGE_DIR")
if dsn == "" || storageDir == "" {
log.Fatal("NEXARCH_RECONCILE_TENANT_DSN und NEXARCH_RECONCILE_STORAGE_DIR muessen gesetzt sein")
}
ctx := context.Background()
pool, err := pgxpool.New(ctx, dsn)
if err != nil {
log.Fatalf("datenbankverbindung: %v", err)
}
defer pool.Close()
dbEntries, err := reconcile.ListDBStorageKeys(ctx, pool)
if err != nil {
log.Fatalf("datenbank-eintraege lesen: %v", err)
}
storageKeys, err := reconcile.ListStorageObjects(storageDir)
if err != nil {
log.Fatalf("objekt-storage durchlaufen: %v", err)
}
report := reconcile.Reconcile(dbEntries, storageKeys)
encoder := json.NewEncoder(os.Stdout)
encoder.SetIndent("", " ")
if err := encoder.Encode(report); err != nil {
log.Fatalf("bericht ausgeben: %v", err)
}
// Nicht-null-Exit-Code bei Abweichungen (Akzeptanzkriterium 3:
// Abweichungen werden BERICHTET, nicht automatisch behoben — der
// Exit-Code macht das fuer systemd/Monitoring sichtbar, OHNE selbst
// irgendetwas zu reparieren).
if !report.IsClean() {
os.Exit(1)
}
}
-87
View File
@@ -1,87 +0,0 @@
# BAK-01 Prüfprotokoll: Datenbank-Backup-Strategie
Welle 1, keine Vorbedingungen. Neues Modul-Verzeichnis `code/archive/`
(gleiches Monorepo-Muster wie `code/dms/`), eigenes Go-Modul
`gitea.perlbach24.de/scripte/nexarch/archive`.
## Grundsatzentscheidung: PostgreSQL-17-natives inkrementelles Backup
`pg_dump` kennt nur logische Vollsicherungen — "inkrementell" im Sinne des
Tickets erfordert das physische Backup-Verfahren. Gewählt: PostgreSQL 17s
natives `pg_basebackup --incremental` (WAL-Summarization), NICHT klassisches
WAL-Archiving (`archive_mode`), weil letzteres einen Neustart der
(geteilten, auch von Core/DMS-Tests genutzten) Postgres-Instanz auf
192.168.1.131 erfordert hätte. Stattdessen `summarize_wal = on` gesetzt —
nur ein `pg_reload_conf()`, kein Neustart, keine Unterbrechung laufender
Verbindungen (per Health-Check nach der Änderung bestätigt).
Voraussetzung geschaffen: Rolle `nexarch_backup` mit `REPLICATION`-Attribut
angelegt (Postgres verlangt eine Replikationsverbindung für
`pg_basebackup`), `pg_hba.conf` erlaubte lokale Replikationsverbindungen
bereits.
## Umsetzung
- `internal/backup.FullBackup`/`IncrementalBackup` — rufen `pg_basebackup`
über `os/exec` auf, Ergebnis landet in einer Generationsstruktur
(`<BackupDir>/<Generation>/full/` bzw. `.../incremental/<ID>/`).
- `internal/backup.Verify` — öffnet `base.tar.gz` vollständig (gzip- UND
tar-Stream, jeder Eintrag bis zum Ende gelesen, nicht nur Kopfdaten) —
Akzeptanzkriterium 2: Verifikation auf Lesbarkeit, nicht nur Erstellung.
- `internal/backup.Rotate`/`ListGenerations` — Generationen sind nach
Zeitstempel-ID sortierbar, `Rotate` entfernt die ältesten bis auf `keep`
komplett (inklusive aller abhängigen Inkremente).
- `cmd/backup-cli``full`/`incremental`/`rotate`, aufgerufen von
systemd-Timern (`deploy/systemd/nexarch-archive-backup-*.timer`) —
"automatisiert nach Zeitplan" (Akzeptanzkriterium 1) entsteht durch die
Timer-Definition, kein zusätzlicher Dauerprozess nötig.
## Prüfungen
| # | Prüfung | Ergebnis |
|---|---|---|
| 1 | Sicherung gegen Testdatenbank erfolgreich erstellt und verifiziert | **bestanden**`TestFullBackup_CreatesVerifiedBackup` gegen die echte Postgres-17-Instanz auf 192.168.1.131 (kein Mock), zusätzlich `TestIncrementalBackup_IsSmallerThanFull`: inkrementelle Sicherung real deutlich kleiner als Vollsicherung (167 KB vs. 16 MB bei der ersten manuellen Probe) — beweist echte inkrementelle Übertragung, nicht nur eine zweite Vollsicherung |
| 2 | Verifikation erkennt eine absichtlich beschädigte Sicherungsdatei | **bestanden**`TestVerify_DetectsCorruptedFile`: 64 Bytes in der Mitte von `base.tar.gz` gekippt, `Verify` schlägt danach fehl (unbeschädigt zuvor erfolgreich) |
| 3 | Rotationsregel entfernt nachweislich nur die ältesten Generationen | **bestanden**`TestRotate_RemovesOnlyOldestGenerations`: 5 Generationen, `keep=2`, exakt die 3 ältesten entfernt, die 2 neuesten nachweislich unangetastet |
## Echte Verdrahtung auf 192.168.1.131 (nicht nur Testcode)
Anders als die zuletzt in DMS gefundenen "Baustein existiert, ist aber
nirgends verdrahtet"-Fälle (FDN-03/FDN-09 gegen Core) wurde hier die
komplette Kette tatsächlich installiert und ausgeführt:
- `backup-cli` gebaut nach `/opt/nexarch-archive/bin/`
- `/etc/nexarch/archive-backup.env` mit den Verbindungsdaten (0600)
- 3 systemd-Timer installiert und aktiviert (`enable --now`):
Vollsicherung täglich 02:00 UTC, Inkrement stündlich, Rotation täglich
03:00 UTC (`systemctl list-timers` bestätigt alle drei scharf)
- Jeder der drei Dienste (`full`/`incremental`/`rotate`) einmal manuell über
`systemctl start` ausgelöst (nicht nur `go test` direkt) — alle drei mit
`status=0/SUCCESS`, Journal bestätigt inhaltlich korrekte Ausgabe
(Vollsicherung erstellt+verifiziert, Inkrement erstellt+verifiziert
gegen die richtige Vorgänger-Generation, Rotation lief ohne Fehler)
## Build/Test-Ergebnis (192.168.1.131, `make check`)
```
go build ./... -> clean
go vet ./... -> clean
golangci-lint run ./... -> 0 issues
go test ./... -p 1 -count=1 -> 4/4 Tests ok, 0 Fehlschläge (echter Postgres 17, kein Mock)
```
## Nachtrag (BAK-02-Sitzung): Backup-Zielverzeichnis korrigiert
`NEXARCH_BACKUP_DIR` zeigte ursprünglich auf `/var/backups/nexarch`
(Root-Dateisystem des Containers, kein dediziertes Dataset) — korrigiert auf
`/var/nexarch-archiv/backups/postgres` (persistentes ZFS-Dataset), siehe
`docs/BAK-02-PRUEFPROTOKOLL.md` Abschnitt „Korrektur an BAK-01" für Details.
Vollsicherung nach der Korrektur erneut über systemd ausgelöst, landet
nachweislich am neuen Ort.
## Gesamtergebnis
**Bestanden.** Alle drei Akzeptanzkriterien und alle drei Pflichtprüfungen
real erfüllt — inklusive tatsächlicher systemd-Timer-Installation und
manuell ausgelöstem End-to-End-Lauf aller drei Dienste auf dem Testhost,
nicht nur isolierter Testcode.
-93
View File
@@ -1,93 +0,0 @@
# BAK-02 Prüfprotokoll: Objekt-Storage-Backup/Snapshots
Welle 1, keine Vorbedingungen.
## Grundsatzentscheidung: restic statt Eigenbau
Nutzerentscheidung: restic statt einer Neuimplementierung, weil restic alle
vier Akzeptanzkriterien mit ausgereiftem, breit geprüftem Tooling erfüllt
(Content-defined Chunking für Dedup, `check --read-data` für
Vollständigkeit, `forget --keep-last` für Rotation, Repository-Verschlüsselung
ab Werk). Installiert via `apt-get install restic` (Version 0.18.0).
Backup-Quelle ist ein lokaler Verzeichnisbaum — für den FDN-03-`LocalDriver`
direkt dessen Basisverzeichnis. Für S3-gestützte Produktions-Deployments
(Betriebsmodus 2/3 aus `STORAGE-KONZEPT.md` Abschnitt 6.2) wäre ein
vorgelagerter Sync-Schritt (z. B. `rclone`) nötig, um Bucket-Inhalte lokal
zu spiegeln, bevor restic sie sichert — restic sichert Dateibäume, keine
S3-Buckets direkt. Das bleibt hier bewusst unimplementiert (kein konkreter
S3-Produktionsbestand vorhanden, der das aktuell erfordert), aber
architektonisch vorgesehen und dokumentiert (`README.md`).
## Umsetzung
- `internal/objectbackup.InitRepo` — idempotent, erkennt "bereits
initialisiert" am `restic init`-Fehlertext statt zu scheitern.
- `internal/objectbackup.Backup``restic backup --json`, parst die
`summary`-Zeile (mehrere JSON-Zeilen in der Ausgabe, gezielt die mit
`message_type=="summary"` gesucht).
- `internal/objectbackup.Check``restic check [--read-data]` (Akzeptanz-
kriterium 3: Vollständigkeitsprüfung).
- `internal/objectbackup.Forget``restic forget --keep-last N --prune`
(Rotation).
- `cmd/objectbackup-cli``backup <dir>`/`check`/`rotate`, aufgerufen von
systemd-Timern (stündlich/wöchentlich/täglich).
## Prüfungen
| # | Prüfung | Ergebnis |
|---|---|---|
| 1 | Zweiter Sicherungslauf nach unverändertem Bestand überträgt keine Daten erneut | **bestanden**`TestBackup_UnchangedSecondRunTransmitsNothingNew`: zweiter Lauf gegen unveränderten Bestand liefert `files_new=0`, `files_changed=0`, `files_unmodified=1` |
| 2 | Zwei identische Testdateien belegen nachweislich nur einmal Speicherplatz | **bestanden**`TestBackup_DeduplicatesIdenticalContent`: zwei Dateien mit identischem Inhalt erzeugen `data_blobs=1`, nicht 2 — echter Dedup-Nachweis über restics Content-defined Chunking, nicht nur Namensvergleich |
| 3 | Vollständigkeitsprüfung erkennt ein fehlendes Objekt in der Sicherung | **bestanden**`TestCheck_DetectsCorruptedPack`: ein Byte in einer echten Repository-Pack-Datei gekippt, `Check(readData=true)` schlägt danach fehl (unbeschädigt zuvor erfolgreich) — dieselbe Vorgehensweise wie die manuelle Recherche vor der Implementierung |
Zusätzlich (nicht explizit als Pflichtprüfung gefordert, aber Teil von
Akzeptanzkriterium 3 „lässt sich einzeln prüfen"): `TestForget_
KeepsOnlyRequestedSnapshotCount` — 3 Sicherungsläufe, `Forget(keepLast=1)`
reduziert auf genau 1 verbleibenden Snapshot.
## Korrektur an BAK-01 im selben Rutsch: Backup-Zielverzeichnis
Nutzerhinweis aufgegriffen: `NEXARCH_BACKUP_DIR` zeigte bei BAK-01
ursprünglich auf `/var/backups/nexarch` (Root-Dateisystem des LXC-
Containers, nicht auf einem der beiden dedizierten ZFS-Datasets). Korrigiert
auf `/var/nexarch-archiv/backups/postgres` (persistentes Dataset
`zfs/data/subvol-1131-disk-0`), NICHT `/var/nexarch-test/` (ephemeres
Dataset `ssd-rpool-data/swap/subvol-1131-disk-0`, wird von
`reset-test-env.sh`-Skripten anderer Module geleert). `objectbackup-cli`s
Repository liegt von Anfang an korrekt unter
`/var/nexarch-archiv/backups/objects`. Beide Pfade real auf
192.168.1.131 verifiziert (`df`/`mount` bestätigt ZFS-Dataset-Zuordnung),
BAK-01s Vollsicherung nach der Korrektur erneut über systemd ausgelöst und
bestätigt am neuen Ort gelandet.
ZFS-seitige Snapshot-/Replikations-Strategie für `nexarch/archiv` bleibt
bewusst außerhalb dieses Tickets (Infra-Runbook, siehe
`STORAGE-KONZEPT.md` Abschnitt 7 „Backup vs. Storage-Redundanz" sowie den
Hinweis, `zfs dedup=on` NICHT zu setzen — App-seitige Dedup über restic
genügt, ZFS-Dedup wäre auf dem 4-GB-Testhost ein Speicherrisiko).
## Echte Verdrahtung auf 192.168.1.131
- `objectbackup-cli` gebaut nach `/opt/nexarch-archive/bin/`
- `/etc/nexarch/archive-objectbackup.env` (0600)
- 3 systemd-Timer installiert und aktiviert: Sicherung stündlich (`:30`),
Vollständigkeitsprüfung wöchentlich (So. 04:00 UTC), Rotation täglich
(03:30 UTC) — `systemctl list-timers` bestätigt alle scharf
- Jeder der drei Dienste einmal über `systemctl start` ausgelöst, alle mit
`status=0/SUCCESS`; Journal bestätigt inhaltlich korrekte Ausgabe
## Build/Test-Ergebnis (192.168.1.131, `make check`)
```
go build ./... -> clean
go vet ./... -> clean
golangci-lint run ./... -> 0 issues
go test ./... -p 1 -count=1 -> 2/2 Pakete mit Tests ok (internal/backup, internal/objectbackup), 0 Fehlschläge
```
## Gesamtergebnis
**Bestanden.** Alle drei Akzeptanzkriterien und alle drei Pflichtprüfungen
real gegen echtes restic-Tooling erfüllt. BAK-01-Pfadfehler im selben
Rutsch korrigiert und erneut end-to-end verifiziert.
-119
View File
@@ -1,119 +0,0 @@
# BAK-05 Prüfprotokoll: Reconciliation / Konsistenzprüfung Storage vs. DB
Voraussetzung BAK-01, BAK-02 (Welle 1) erledigt, siehe eigene Protokolle.
## Grundsatzentscheidung: reine Funktion + zwei Quell-Adapter
`internal/reconcile.Reconcile` ist eine reine Funktion ohne DB-/Storage-
Zugriff (leicht ohne echte Infrastruktur testbar), die Ein- und
Auslesen echter Systeme ist strikt in `sources.go` getrennt
(`ListDBStorageKeys` gegen echtes Postgres, `ListStorageObjects` gegen
echtes Dateisystem). Beide Seiten liefern nur SCHLÜSSEL niemals Inhalt
dadurch bleibt BAK-05 sauber getrennt von BAK-08 (Inhalts-/Prüfsummen-
verifikation, eigene Fehlerklasse, eigenes Ticket).
Report-Format bewusst deterministisch: alle drei Ergebnislisten
(`missing_in_storage`, `orphaned_in_storage`, `existing_in_storage`)
nach `storage_key` aufsteigend sortiert.
**Nachtrag (nach Rückfrage vor BAK-08-Start):** Der ursprüngliche Report
enthielt nur die beiden Abweichungslisten keine Liste der bestätigt
existierenden Objekte. Für BAK-08 als Stichprobengrundlage reicht
"keine Abweichung" nicht, es braucht die tatsächliche, deterministisch
sortierte Liste. Ergänzt: `Report.ExistingInStorage` DB-Eintrag UND
Storage-Objekt beide vorhanden, reine Existenzbestätigung (keine
Inhaltsprüfung, Scope-Trennung zu BAK-08 bleibt gewahrt), aufsteigend
nach `storage_key` sortiert. BAK-08 zieht seine Stichprobe daraus, ohne
selbst zu sortieren/filtern. Neuer Test
`TestReconcile_ExistingInStorageIsStableSamplingBasis` beweist Inhalt
und Sortierung. Real neu gebaut, getestet (9/9) und auf 131 erneut
ausgelöst Journal zeigt das Feld `existing_in_storage` im Report.
Meldeweg über OPS-05 (wie später BAK-08) wurde als offene Design-Frage
aufgeworfen, aber nicht zur Vorbedingung gemacht hier bewusst noch
nicht umgesetzt (kein OPS-05-Abhängigkeitseintrag im Board für BAK-05);
Report wird aktuell nur als JSON auf stdout ausgegeben und per
Exit-Code (1 bei Abweichungen) für systemd/Monitoring sichtbar gemacht.
Anbindung an OPS-05 kann bei Bedarf nachgezogen werden, ohne
`Reconcile` selbst zu ändern.
## Umsetzung
- `internal/reconcile.Reconcile(dbEntries, storageKeys) Report` reine
Vergleichsfunktion, liefert `MissingInStorage`/`OrphanedInStorage`,
`Report.IsClean()` als eindeutiges Sauber-Merkmal.
- `internal/reconcile.ListDBStorageKeys` liest `file_revisions`
(DMS FDN-02) per direktem SQL aus derselben physischen Tenant-DB
(Modell C, Core TEN-01) kein Import von DMS-Go-Paketen möglich
(eigenes Go-Modul), daher reiner SQL-Zugriff gegen das dokumentierte
Schema.
- `internal/reconcile.ListStorageObjects` durchläuft den lokalen
FDN-03-`LocalDriver`-Basisordner (`filepath.WalkDir`), liefert `nil,
nil` bei fehlendem Verzeichnis statt Fehler (noch keine Objekte ist
kein Fehlerzustand).
- `cmd/reconcile-cli` liest `NEXARCH_RECONCILE_TENANT_DSN` und
`NEXARCH_RECONCILE_STORAGE_DIR`, gibt Report als JSON auf stdout aus,
Exit-Code 1 bei Abweichungen.
## Prüfungen
| # | Prüfung | Ergebnis |
|---|---|---|
| 1 | Datenbankeintrag ohne Storage-Objekt wird erkannt | **bestanden**`TestReconcile_DetectsMissingInStorage` |
| 2 | Storage-Objekt ohne Datenbankeintrag wird erkannt | **bestanden**`TestReconcile_DetectsOrphanedInStorage` |
| 3 | Lauf ohne Abweichungen liefert leeren, eindeutig sauberen Bericht | **bestanden**`TestReconcile_CleanRunProducesEmptyReport` (zusätzlich `IsClean()`-Konsistenzprüfung) |
Zusätzlich (Nutzervorgaben, nicht explizit im Ticket als Pflichtprüfung
benannt, aber zentral für die Abgrenzung/Weiterverwendbarkeit):
- `TestReconcile_ExistingButCorruptedObjectProducesNoFinding` Nachweis,
dass Reconcile AUSSCHLIESSLICH Existenz prüft, niemals Inhalt (Trennung
von BAK-08).
- `TestReconcile_DeterministicOrdering` zwei Läufe mit identischer
Eingabe liefern identische Reihenfolge, aufsteigend nach `storage_key`.
- `TestListDBStorageKeys_ReadsRealFileRevisions` liest echt gegen die
gemeinsame Tenant-Testdatenbank `dms_tenant_test` (reales DMS-FDN-02-
Schema, kein Mock).
- `TestListStorageObjects_WalksRealDirectory` /
`_MissingDirectoryReturnsEmpty` echtes Dateisystem, kein Mock.
## Echte Verdrahtung auf 192.168.1.131
- `reconcile-cli` gebaut nach `/opt/nexarch-archive/bin/`
- `/etc/nexarch/archive-reconcile.env` (0600): `NEXARCH_RECONCILE_TENANT_DSN`
zeigt auf die gemeinsame Tenant-Testdatenbank `dms_tenant_test`
(DMS selbst läuft auf 192.168.1.131 noch nicht als eigener systemd-
Dienst mit persistenter Konfiguration dies ist die real verfügbare
Tenant-DB mit echtem FDN-02-Schema, dokumentierter bekannter Stand,
kein stiller Mock); `NEXARCH_RECONCILE_STORAGE_DIR` zeigt auf
`/var/nexarch-archiv/dms-objects` (persistentes ZFS-Dataset, NICHT
`/var/nexarch-test/`).
- Timer `nexarch-archive-reconcile.timer` installiert und aktiviert
(täglich 05:00 UTC), `systemctl list-timers` bestätigt scharf.
- `systemctl start nexarch-archive-reconcile.service` real ausgelöst:
`status=0/SUCCESS`, Journal zeigt echten JSON-Report
(`missing_in_storage: null, orphaned_in_storage: null` Tenant-DB
aktuell leer, daher sauberer Bericht, keine synthetische Ausgabe).
## Build/Test-Ergebnis (192.168.1.131, `make check`)
```
go build ./... -> clean
go vet ./... -> clean
golangci-lint run ./... -> 0 issues
go test ./... -p 1 -count=1 -> 3/3 Pakete mit Tests ok (internal/backup, internal/objectbackup, internal/reconcile), 0 Fehlschläge
```
`internal/reconcile`-Tests separat mit gesetzter `TEST_TENANT_DSN` gegen
`dms_tenant_test` verifiziert: 9/9 Tests bestanden (6 reine
`Reconcile`-Tests + 3 `sources.go`-Integrationstests).
## Gesamtergebnis
**Bestanden.** Alle drei Akzeptanzkriterien und alle Pflicht- sowie
Nutzervorgaben-Prüfungen real erfüllt (echte Postgres-Instanz, echtes
Dateisystem, echter systemd-Lauf). Zwei Testfehler während der
Entwicklung (Schema-Abweichung `revision_number` NOT NULL in der realen
`dms_tenant_test`-Tabelle; inkonsistente Fixture-Daten in
`TestReconcile_DeterministicOrdering`) gefunden und korrigiert beide
waren Testautorenfehler, keine Fehler in `Reconcile` selbst.
-14
View File
@@ -1,14 +0,0 @@
module gitea.perlbach24.de/scripte/nexarch/archive
go 1.22
require github.com/jackc/pgx/v5 v5.6.0
require (
github.com/jackc/pgpassfile v1.0.0 // indirect
github.com/jackc/pgservicefile v0.0.0-20221227161230-091c0ba34f0a // indirect
github.com/jackc/puddle/v2 v2.2.1 // indirect
golang.org/x/crypto v0.17.0 // indirect
golang.org/x/sync v0.1.0 // indirect
golang.org/x/text v0.14.0 // indirect
)
-102
View File
@@ -1,102 +0,0 @@
// Package backup implementiert BAK-01: automatisierte, inkrementelle
// Sicherung der PostgreSQL-Datenbank per pg_basebackup (PostgreSQL 17s
// natives inkrementelles Backup über WAL-Summarization, siehe
// `summarize_wal`), mit Verifikation jeder Sicherung und
// generationsbasierter Rotation. Kein pg_dump-basierter Ansatz, weil
// pg_dump ausschließlich logische Vollsicherungen kennt — "inkrementell"
// im Sinne des Tickets erfordert das physische, WAL-summary-gestützte
// Verfahren aus PostgreSQL 17.
package backup
import (
"context"
"fmt"
"os"
"os/exec"
"path/filepath"
"time"
)
// Config enthält die Verbindungsdaten für pg_basebackup — ausschließlich
// über Umgebungsvariablen befüllt, nie im Code (siehe Ticket-Abschluss-
// Regel).
type Config struct {
Host string
Port string
User string
Password string
BackupDir string
PgBaseBackupPath string // Default "pg_basebackup", überschreibbar für Tests
}
func (c Config) binary() string {
if c.PgBaseBackupPath != "" {
return c.PgBaseBackupPath
}
return "pg_basebackup"
}
// FullBackupDirName/IncrementalDirName sind die festen Unterverzeichnis-
// namen je Generation.
const (
FullBackupDirName = "full"
IncrementalSubdir = "incremental"
BackupManifestFile = "backup_manifest"
BaseTarGzFile = "base.tar.gz"
)
// NewGenerationID liefert eine sortierbare, eindeutige Generation-Kennung
// (RFC3339-artig, dateisystemtauglich) — Generationen werden anhand dieser
// Kennung chronologisch sortiert (Rotate, ListGenerations).
func NewGenerationID(t time.Time) string {
return t.UTC().Format("20060102T150405Z")
}
// FullBackup erstellt eine neue Vollsicherung (Akzeptanzkriterium 1) als
// eigene Generation. Liefert den Pfad zum backup_manifest, das spätere
// IncrementalBackup-Aufrufe als Referenz brauchen.
func FullBackup(ctx context.Context, cfg Config, generationID string) (manifestPath string, err error) {
dir := filepath.Join(cfg.BackupDir, generationID, FullBackupDirName)
if err := os.MkdirAll(filepath.Dir(dir), 0o750); err != nil {
return "", fmt.Errorf("backup: generationsverzeichnis anlegen: %w", err)
}
args := []string{
"-h", cfg.Host, "-p", cfg.Port, "-U", cfg.User,
"-D", dir, "-Ft", "-z", "--checkpoint=fast", "--no-password",
}
if err := runPgBaseBackup(ctx, cfg, args); err != nil {
return "", fmt.Errorf("backup: vollsicherung: %w", err)
}
return filepath.Join(dir, BackupManifestFile), nil
}
// IncrementalBackup erstellt eine inkrementelle Sicherung gegen die zuletzt
// bekannte Vollsicherung ODER die letzte Inkrement-Sicherung (priorManifestPath
// zeigt jeweils auf das backup_manifest der Referenz).
func IncrementalBackup(ctx context.Context, cfg Config, generationID, incrementID, priorManifestPath string) (manifestPath string, err error) {
dir := filepath.Join(cfg.BackupDir, generationID, IncrementalSubdir, incrementID)
if err := os.MkdirAll(filepath.Dir(dir), 0o750); err != nil {
return "", fmt.Errorf("backup: inkrement-verzeichnis anlegen: %w", err)
}
args := []string{
"-h", cfg.Host, "-p", cfg.Port, "-U", cfg.User,
"-D", dir, "-Ft", "-z", "--checkpoint=fast", "--no-password",
"--incremental=" + priorManifestPath,
}
if err := runPgBaseBackup(ctx, cfg, args); err != nil {
return "", fmt.Errorf("backup: inkrementelle sicherung: %w", err)
}
return filepath.Join(dir, BackupManifestFile), nil
}
func runPgBaseBackup(ctx context.Context, cfg Config, args []string) error {
cmd := exec.CommandContext(ctx, cfg.binary(), args...)
cmd.Env = append(os.Environ(), "PGPASSWORD="+cfg.Password)
output, err := cmd.CombinedOutput()
if err != nil {
return fmt.Errorf("%s fehlgeschlagen: %w (ausgabe: %s)", cfg.binary(), err, string(output))
}
return nil
}
-187
View File
@@ -1,187 +0,0 @@
package backup
import (
"context"
"os"
"path/filepath"
"testing"
"time"
)
func requireTestConfig(t *testing.T) Config {
t.Helper()
user := os.Getenv("TEST_BACKUP_PG_USER")
if user == "" {
t.Skip("TEST_BACKUP_PG_USER nicht gesetzt, Integrationstest uebersprungen (braucht echten Postgres mit REPLICATION-Rolle)")
}
return Config{
Host: envOr("TEST_BACKUP_PG_HOST", "localhost"),
Port: envOr("TEST_BACKUP_PG_PORT", "5432"),
User: user,
Password: os.Getenv("TEST_BACKUP_PG_PASSWORD"),
BackupDir: t.TempDir(),
}
}
func envOr(key, fallback string) string {
if v := os.Getenv(key); v != "" {
return v
}
return fallback
}
// TestFullBackup_CreatesVerifiedBackup ist Pruefung 1: Sicherung gegen
// Testdatenbank erfolgreich erstellt und verifiziert.
func TestFullBackup_CreatesVerifiedBackup(t *testing.T) {
cfg := requireTestConfig(t)
ctx := context.Background()
genID := NewGenerationID(time.Now())
manifest, err := FullBackup(ctx, cfg, genID)
if err != nil {
t.Fatalf("fullbackup: %v", err)
}
if _, err := os.Stat(manifest); err != nil {
t.Fatalf("backup_manifest fehlt: %v", err)
}
dir := filepath.Dir(manifest)
if _, err := os.Stat(filepath.Join(dir, BaseTarGzFile)); err != nil {
t.Fatalf("%s fehlt: %v", BaseTarGzFile, err)
}
if err := Verify(dir); err != nil {
t.Fatalf("verify: %v", err)
}
}
// TestIncrementalBackup_IsSmallerThanFull ist der Nachweis fuer
// Akzeptanzkriterium 1 (inkrementell): eine echte inkrementelle Sicherung
// gegen unveraenderten Bestand ist deutlich kleiner als die Vollsicherung —
// beweist, dass tatsaechlich nur Aenderungen uebertragen wurden (PostgreSQL
// 17 WAL-Summarization), nicht nochmal alles.
func TestIncrementalBackup_IsSmallerThanFull(t *testing.T) {
cfg := requireTestConfig(t)
ctx := context.Background()
genID := NewGenerationID(time.Now())
fullManifest, err := FullBackup(ctx, cfg, genID)
if err != nil {
t.Fatalf("fullbackup: %v", err)
}
fullDir := filepath.Dir(fullManifest)
fullSize := fileSize(t, filepath.Join(fullDir, BaseTarGzFile))
incID := NewGenerationID(time.Now().Add(time.Second))
incManifest, err := IncrementalBackup(ctx, cfg, genID, incID, fullManifest)
if err != nil {
t.Fatalf("incrementalbackup: %v", err)
}
incDir := filepath.Dir(incManifest)
if err := Verify(incDir); err != nil {
t.Fatalf("verify (inkrementell): %v", err)
}
incSize := fileSize(t, filepath.Join(incDir, BaseTarGzFile))
if incSize >= fullSize {
t.Fatalf("inkrementelle sicherung (%d bytes) ist nicht kleiner als die vollsicherung (%d bytes) - keine echte inkrementelle Uebertragung", incSize, fullSize)
}
}
func fileSize(t *testing.T, path string) int64 {
t.Helper()
info, err := os.Stat(path)
if err != nil {
t.Fatalf("dateigroesse von %q ermitteln: %v", path, err)
}
return info.Size()
}
// TestVerify_DetectsCorruptedFile ist Pruefung 2: Verifikation erkennt eine
// absichtlich beschaedigte Sicherungsdatei.
func TestVerify_DetectsCorruptedFile(t *testing.T) {
cfg := requireTestConfig(t)
ctx := context.Background()
genID := NewGenerationID(time.Now())
manifest, err := FullBackup(ctx, cfg, genID)
if err != nil {
t.Fatalf("fullbackup: %v", err)
}
dir := filepath.Dir(manifest)
if err := Verify(dir); err != nil {
t.Fatalf("verify (unbeschaedigt) haette erfolgreich sein muessen: %v", err)
}
// Absichtliche Beschaedigung: mehrere Bytes in der Mitte der Datei kippen.
path := filepath.Join(dir, BaseTarGzFile)
data, err := os.ReadFile(path)
if err != nil {
t.Fatalf("sicherungsdatei lesen: %v", err)
}
mid := len(data) / 2
for i := mid; i < mid+64 && i < len(data); i++ {
data[i] ^= 0xFF
}
if err := os.WriteFile(path, data, 0o600); err != nil {
t.Fatalf("beschaedigte sicherungsdatei schreiben: %v", err)
}
if err := Verify(dir); err == nil {
t.Fatal("verify haette die beschaedigte sicherungsdatei erkennen muessen")
}
}
// TestRotate_RemovesOnlyOldestGenerations ist Pruefung 3.
func TestRotate_RemovesOnlyOldestGenerations(t *testing.T) {
backupDir := t.TempDir()
generationIDs := []string{
"20260101T000000Z",
"20260102T000000Z",
"20260103T000000Z",
"20260104T000000Z",
"20260105T000000Z",
}
for _, id := range generationIDs {
if err := os.MkdirAll(filepath.Join(backupDir, id, FullBackupDirName), 0o750); err != nil {
t.Fatalf("generation %q anlegen: %v", id, err)
}
}
removed, err := Rotate(backupDir, 2)
if err != nil {
t.Fatalf("rotate: %v", err)
}
wantRemoved := []string{"20260101T000000Z", "20260102T000000Z", "20260103T000000Z"}
if len(removed) != len(wantRemoved) {
t.Fatalf("entfernte generationen = %v, want %v", removed, wantRemoved)
}
for i, w := range wantRemoved {
if removed[i] != w {
t.Fatalf("entfernte generationen = %v, want %v", removed, wantRemoved)
}
}
remaining, err := ListGenerations(backupDir)
if err != nil {
t.Fatalf("listgenerations: %v", err)
}
wantRemaining := []string{"20260104T000000Z", "20260105T000000Z"}
if len(remaining) != len(wantRemaining) {
t.Fatalf("verbleibende generationen = %v, want %v", remaining, wantRemaining)
}
for i, w := range wantRemaining {
if remaining[i] != w {
t.Fatalf("verbleibende generationen = %v, want %v", remaining, wantRemaining)
}
}
// Die NEUESTEN duerfen NICHT entfernt sein (Pruefung 3: nur die
// aeltesten Generationen).
for _, w := range wantRemaining {
if _, err := os.Stat(filepath.Join(backupDir, w)); err != nil {
t.Fatalf("neueste generation %q wurde faelschlich entfernt: %v", w, err)
}
}
}
-56
View File
@@ -1,56 +0,0 @@
package backup
import (
"fmt"
"os"
"path/filepath"
"sort"
)
// ListGenerations liefert alle Generation-IDs in backupDir, aufsteigend
// sortiert (die GenerationID selbst ist chronologisch sortierbar, siehe
// NewGenerationID — kein Blick auf Dateisystem-Zeitstempel nötig, die bei
// einem Restore/Kopiervorgang verändert werden könnten).
func ListGenerations(backupDir string) ([]string, error) {
entries, err := os.ReadDir(backupDir)
if err != nil {
if os.IsNotExist(err) {
return nil, nil
}
return nil, fmt.Errorf("backup: sicherungsverzeichnis lesen: %w", err)
}
var generations []string
for _, e := range entries {
if e.IsDir() {
generations = append(generations, e.Name())
}
}
sort.Strings(generations)
return generations, nil
}
// Rotate entfernt alle bis auf die `keep` NEUESTEN Generationen
// (Akzeptanzkriterium 3) — jede Generation umfasst ihre Vollsicherung UND
// alle davon abhängigen Inkremente, ein Löschen der gesamten
// Generationsverzeichnisses entfernt beides konsistent zusammen.
func Rotate(backupDir string, keep int) (removed []string, err error) {
if keep < 0 {
keep = 0
}
generations, err := ListGenerations(backupDir)
if err != nil {
return nil, err
}
if len(generations) <= keep {
return nil, nil
}
toRemove := generations[:len(generations)-keep]
for _, gen := range toRemove {
if err := os.RemoveAll(filepath.Join(backupDir, gen)); err != nil {
return removed, fmt.Errorf("backup: generation %q entfernen: %w", gen, err)
}
removed = append(removed, gen)
}
return removed, nil
}
-54
View File
@@ -1,54 +0,0 @@
package backup
import (
"archive/tar"
"compress/gzip"
"fmt"
"io"
"os"
"path/filepath"
)
// ErrCorrupted wird geliefert, wenn eine Sicherungsdatei nicht lesbar ist
// (Akzeptanzkriterium 2: Verifikation, nicht nur Erstellungs-Prüfung).
var ErrCorrupted = fmt.Errorf("backup: sicherungsdatei ist beschaedigt oder unvollstaendig")
// Verify prüft, dass base.tar.gz im gegebenen Sicherungsverzeichnis
// vollständig lesbar ist — öffnet gzip- UND tar-Stream und liest JEDEN
// Eintrag bis zum Ende durch (nicht nur die Kopfdaten), damit ein
// abgeschnittener oder mit kaputten Bytes überschriebener Inhalt
// zuverlässig auffällt, nicht nur ein defekter Tar-Header.
func Verify(backupDir string) error {
path := filepath.Join(backupDir, BaseTarGzFile)
f, err := os.Open(path)
if err != nil {
return fmt.Errorf("%w: %s nicht lesbar: %v", ErrCorrupted, path, err)
}
defer func() { _ = f.Close() }()
gz, err := gzip.NewReader(f)
if err != nil {
return fmt.Errorf("%w: gzip-header ungueltig: %v", ErrCorrupted, err)
}
defer func() { _ = gz.Close() }()
tr := tar.NewReader(gz)
entries := 0
for {
hdr, err := tr.Next()
if err == io.EOF {
break
}
if err != nil {
return fmt.Errorf("%w: tar-eintrag ungueltig: %v", ErrCorrupted, err)
}
if _, err := io.Copy(io.Discard, tr); err != nil {
return fmt.Errorf("%w: inhalt von %q nicht vollstaendig lesbar: %v", ErrCorrupted, hdr.Name, err)
}
entries++
}
if entries == 0 {
return fmt.Errorf("%w: archiv enthaelt keine eintraege", ErrCorrupted)
}
return nil
}
-156
View File
@@ -1,156 +0,0 @@
// Package objectbackup implementiert BAK-02: automatisierte, inkrementelle,
// deduplizierende Sicherung des Objekt-Storage-Bestands. Nutzt restic
// (Content-defined Chunking, verschlüsseltes Repository ab Werk) statt
// Eigenbau — restic erfüllt alle Akzeptanzkriterien mit ausgereiftem,
// geprüftem Tooling statt einer weniger robusten Neuimplementierung.
//
// Backup-Quelle ist ein lokaler Verzeichnisbaum — für den LocalDriver aus
// FDN-03 direkt dessen Basisverzeichnis, für S3-gestützte Produktions-
// Deployments ein vorgelagerter Sync-Schritt (z.B. rclone) auf einen
// lokalen Spiegel, bevor restic ihn sichert (nicht Bestandteil dieser
// Kachel — restic selbst sichert Dateibäume, keine S3-Buckets direkt).
package objectbackup
import (
"context"
"encoding/json"
"fmt"
"os"
"os/exec"
"strings"
)
// Config enthält Repository-Ort und -Passwort — ausschließlich über
// Umgebungsvariablen befüllt (siehe Ticket-Abschluss-Regel).
type Config struct {
RepoDir string
Password string
ResticPath string // Default "restic", überschreibbar für Tests
}
func (c Config) binary() string {
if c.ResticPath != "" {
return c.ResticPath
}
return "restic"
}
func (c Config) env() []string {
return append(os.Environ(), "RESTIC_PASSWORD="+c.Password)
}
func run(ctx context.Context, cfg Config, args ...string) ([]byte, error) {
fullArgs := append([]string{"-r", cfg.RepoDir}, args...)
cmd := exec.CommandContext(ctx, cfg.binary(), fullArgs...)
cmd.Env = cfg.env()
output, err := cmd.CombinedOutput()
if err != nil {
return output, fmt.Errorf("%s %v fehlgeschlagen: %w (ausgabe: %s)", cfg.binary(), args, err, string(output))
}
return output, nil
}
// InitRepo legt ein neues restic-Repository an, falls es noch nicht
// existiert — idempotent, ein bereits initialisiertes Repository ist kein
// Fehler (Wiederholte Aufrufe durch systemd-Timer nach einem Neustart
// dürfen nicht fehlschlagen).
func InitRepo(ctx context.Context, cfg Config) error {
output, err := run(ctx, cfg, "init")
if err != nil {
if strings.Contains(string(output), "config file already exists") {
return nil
}
return fmt.Errorf("objectbackup: repository initialisieren: %w", err)
}
return nil
}
// BackupSummary ist der geparste "summary"-Datensatz aus `restic backup --json`.
type BackupSummary struct {
SnapshotID string `json:"snapshot_id"`
FilesNew int `json:"files_new"`
FilesChanged int `json:"files_changed"`
FilesUnmodified int `json:"files_unmodified"`
DataBlobs int `json:"data_blobs"`
TotalBytes int64 `json:"total_bytes_processed"`
}
// Backup sichert sourceDir inkrementell (Akzeptanzkriterium 1: unveränderte
// Objekte werden nicht erneut übertragen — restics Content-defined
// Chunking erkennt das automatisch, kein manueller Änderungsabgleich
// nötig).
func Backup(ctx context.Context, cfg Config, sourceDir string) (BackupSummary, error) {
output, err := run(ctx, cfg, "backup", sourceDir, "--json")
if err != nil {
return BackupSummary{}, fmt.Errorf("objectbackup: sicherung: %w", err)
}
return parseSummary(output)
}
// parseSummary sucht in der zeilenweisen JSON-Ausgabe von `restic backup
// --json` (mehrere Fortschritts-/Statuszeilen, GENAU EINE mit
// message_type=="summary") die Zusammenfassung.
func parseSummary(output []byte) (BackupSummary, error) {
lines := strings.Split(strings.TrimSpace(string(output)), "\n")
for i := len(lines) - 1; i >= 0; i-- {
var probe struct {
MessageType string `json:"message_type"`
}
if err := json.Unmarshal([]byte(lines[i]), &probe); err != nil {
continue
}
if probe.MessageType == "summary" {
var summary BackupSummary
if err := json.Unmarshal([]byte(lines[i]), &summary); err != nil {
return BackupSummary{}, fmt.Errorf("objectbackup: summary-zeile dekodieren: %w", err)
}
return summary, nil
}
}
return BackupSummary{}, fmt.Errorf("objectbackup: keine summary-zeile in der restic-ausgabe gefunden")
}
// Check prüft die Vollständigkeit/Lesbarkeit des Repository
// (Akzeptanzkriterium 3 / Pflichtprüfung: Vollständigkeitsprüfung erkennt
// fehlendes/beschädigtes Objekt). readData=true liest jeden gespeicherten
// Datenblock tatsächlich (teurer, aber die einzige Prüfung, die
// Bit-Rot in bereits gespeicherten Paketen erkennt — ohne readData prüft
// restic nur Struktur/Indizes, nicht den tatsächlichen Blockinhalt).
func Check(ctx context.Context, cfg Config, readData bool) error {
args := []string{"check"}
if readData {
args = append(args, "--read-data")
}
if _, err := run(ctx, cfg, args...); err != nil {
return fmt.Errorf("objectbackup: %w", err)
}
return nil
}
// Forget entfernt alte Snapshots nach Rotationsregel und gibt den davon
// belegten Speicherplatz frei (--prune) — restics Äquivalent zu
// BAK-01s Rotate.
func Forget(ctx context.Context, cfg Config, keepLast int) error {
if _, err := run(ctx, cfg, "forget", "--keep-last", fmt.Sprintf("%d", keepLast), "--prune"); err != nil {
return fmt.Errorf("objectbackup: rotation: %w", err)
}
return nil
}
type snapshotEntry struct {
ShortID string `json:"short_id"`
}
// SnapshotCount liefert die Anzahl vorhandener Snapshots — für Tests und
// Statusabfragen.
func SnapshotCount(ctx context.Context, cfg Config) (int, error) {
output, err := run(ctx, cfg, "snapshots", "--json")
if err != nil {
return 0, fmt.Errorf("objectbackup: snapshots auflisten: %w", err)
}
var snapshots []snapshotEntry
if err := json.Unmarshal(output, &snapshots); err != nil {
return 0, fmt.Errorf("objectbackup: snapshot-liste dekodieren: %w", err)
}
return len(snapshots), nil
}
@@ -1,168 +0,0 @@
package objectbackup
import (
"context"
"os"
"os/exec"
"path/filepath"
"testing"
)
func requireRestic(t *testing.T) {
t.Helper()
if _, err := exec.LookPath("restic"); err != nil {
t.Skip("restic nicht installiert, Integrationstest uebersprungen")
}
}
func setupTest(t *testing.T) Config {
t.Helper()
requireRestic(t)
cfg := Config{RepoDir: filepath.Join(t.TempDir(), "repo"), Password: "test-passwort-fuer-objectbackup"}
if err := InitRepo(context.Background(), cfg); err != nil {
t.Fatalf("initrepo: %v", err)
}
return cfg
}
func writeFile(t *testing.T, dir, name, content string) {
t.Helper()
if err := os.WriteFile(filepath.Join(dir, name), []byte(content), 0o600); err != nil {
t.Fatalf("testdatei %q schreiben: %v", name, err)
}
}
// TestBackup_UnchangedSecondRunTransmitsNothingNew ist Pruefung 1:
// zweiter Sicherungslauf nach unveraendertem Bestand ueberraegt keine
// Daten erneut.
func TestBackup_UnchangedSecondRunTransmitsNothingNew(t *testing.T) {
cfg := setupTest(t)
ctx := context.Background()
sourceDir := t.TempDir()
writeFile(t, sourceDir, "dokument.pdf", "unveraenderter inhalt")
first, err := Backup(ctx, cfg, sourceDir)
if err != nil {
t.Fatalf("erste sicherung: %v", err)
}
if first.FilesNew != 1 {
t.Fatalf("erste sicherung: files_new = %d, want 1", first.FilesNew)
}
second, err := Backup(ctx, cfg, sourceDir)
if err != nil {
t.Fatalf("zweite sicherung: %v", err)
}
if second.FilesNew != 0 || second.FilesChanged != 0 {
t.Fatalf("zweite sicherung (unveraendert): files_new=%d files_changed=%d, want beide 0", second.FilesNew, second.FilesChanged)
}
if second.FilesUnmodified != 1 {
t.Fatalf("zweite sicherung: files_unmodified = %d, want 1", second.FilesUnmodified)
}
}
// TestBackup_DeduplicatesIdenticalContent ist Pruefung 2: zwei identische
// Testdateien belegen nachweislich nur einmal Speicherplatz.
func TestBackup_DeduplicatesIdenticalContent(t *testing.T) {
cfg := setupTest(t)
ctx := context.Background()
sourceDir := t.TempDir()
content := "exakt identischer inhalt in beiden dateien fuer den dedup-nachweis"
writeFile(t, sourceDir, "original.pdf", content)
writeFile(t, sourceDir, "kopie.pdf", content)
summary, err := Backup(ctx, cfg, sourceDir)
if err != nil {
t.Fatalf("sicherung: %v", err)
}
if summary.FilesNew != 2 {
t.Fatalf("erwartet 2 neue dateien, habe %d", summary.FilesNew)
}
// Zwei Dateien mit IDENTISCHEM Inhalt duerfen nur EINEN data_blob
// erzeugen - das ist der Dedup-Nachweis (Akzeptanzkriterium 2).
if summary.DataBlobs != 1 {
t.Fatalf("data_blobs = %d, want 1 (zwei identische dateien haetten nur einen blob erzeugen duerfen - keine dedup)", summary.DataBlobs)
}
}
// TestCheck_DetectsCorruptedPack ist Pruefung 3: Vollstaendigkeitspruefung
// erkennt ein beschaedigtes/fehlendes Objekt in der Sicherung.
func TestCheck_DetectsCorruptedPack(t *testing.T) {
cfg := setupTest(t)
ctx := context.Background()
sourceDir := t.TempDir()
writeFile(t, sourceDir, "wichtig.pdf", "inhalt, der spaeter absichtlich beschaedigt wird")
if _, err := Backup(ctx, cfg, sourceDir); err != nil {
t.Fatalf("sicherung: %v", err)
}
if err := Check(ctx, cfg, true); err != nil {
t.Fatalf("check (unbeschaedigt) haette erfolgreich sein muessen: %v", err)
}
// Absichtliche Beschaedigung: ein Byte in einer Pack-Datei im
// Repository kippen (dieselbe Fundstelle wie beim manuellen
// Nachweis waehrend der Recherche zu diesem Ticket).
packDir := filepath.Join(cfg.RepoDir, "data")
corrupted := false
if err := filepath.Walk(packDir, func(path string, info os.FileInfo, err error) error {
if err != nil || info.IsDir() || corrupted {
return err
}
data, err := os.ReadFile(path)
if err != nil {
return err
}
if len(data) < 20 {
return nil
}
data[10] ^= 0xFF
if err := os.WriteFile(path, data, 0o600); err != nil {
return err
}
corrupted = true
return nil
}); err != nil {
t.Fatalf("pack-datei beschaedigen: %v", err)
}
if !corrupted {
t.Fatal("keine pack-datei zum beschaedigen gefunden - testaufbau fehlerhaft")
}
if err := Check(ctx, cfg, true); err == nil {
t.Fatal("check haette die beschaedigte pack-datei erkennen muessen")
}
}
// TestForget_KeepsOnlyRequestedSnapshotCount prueft die Rotation.
func TestForget_KeepsOnlyRequestedSnapshotCount(t *testing.T) {
cfg := setupTest(t)
ctx := context.Background()
sourceDir := t.TempDir()
for i := 0; i < 3; i++ {
writeFile(t, sourceDir, "f.txt", "version "+string(rune('a'+i)))
if _, err := Backup(ctx, cfg, sourceDir); err != nil {
t.Fatalf("sicherung %d: %v", i, err)
}
}
before, err := SnapshotCount(ctx, cfg)
if err != nil {
t.Fatalf("snapshotcount (vorher): %v", err)
}
if before != 3 {
t.Fatalf("erwartet 3 snapshots vor rotation, habe %d", before)
}
if err := Forget(ctx, cfg, 1); err != nil {
t.Fatalf("forget: %v", err)
}
after, err := SnapshotCount(ctx, cfg)
if err != nil {
t.Fatalf("snapshotcount (nachher): %v", err)
}
if after != 1 {
t.Fatalf("erwartet 1 snapshot nach rotation (keep-last 1), habe %d", after)
}
}
-106
View File
@@ -1,106 +0,0 @@
// Package reconcile implementiert BAK-05: periodischer Abgleich, ob jeder
// in der Datenbank referenzierte Objekt-Storage-Eintrag tatsächlich
// existiert und umgekehrt. Prüft AUSSCHLIESSLICH Existenz — niemals
// Inhalt (das ist Archive BAK-08, eine eigene Fehlerklasse, bewusst nicht
// hier mit hineingezogen, siehe reconcile_test.go
// TestReconcile_ExistingButCorruptedObjectProducesNoFinding).
package reconcile
import (
"sort"
"time"
)
// Finding ist EIN Abweichungsfund — entweder ein Datenbankeintrag ohne
// Storage-Objekt oder umgekehrt.
type Finding struct {
StorageKey string `json:"storage_key"`
DocumentID string `json:"document_id,omitempty"`
RevisionID string `json:"revision_id,omitempty"`
}
// Report ist das Ergebnis EINES Abgleichslaufs (Akzeptanzkriterium 3:
// Abweichungen werden BERICHTET, nicht automatisch behoben — Report ist
// reine Information, keine Reparaturfunktion existiert in diesem Paket).
//
// Beide Listen sind nach StorageKey aufsteigend sortiert — bei gleicher
// Eingabe liefert Reconcile IMMER dieselbe Reihenfolge (deterministisch),
// damit ein nachgelagerter Verbraucher (Archive BAK-08: zieht seine
// Stichprobe aus der Liste der EXISTIERENDEN Objekte) sich auf eine
// stabile Sortierung verlassen kann, statt bei jedem Lauf neu zu
// filtern/sortieren.
type Report struct {
GeneratedAt time.Time `json:"generated_at"`
// MissingInStorage: Datenbankeintrag vorhanden, Objekt im Storage fehlt
// (Akzeptanzkriterium 1).
MissingInStorage []Finding `json:"missing_in_storage"`
// OrphanedInStorage: Objekt im Storage vorhanden, kein Datenbankeintrag
// (Akzeptanzkriterium 2).
OrphanedInStorage []Finding `json:"orphaned_in_storage"`
// ExistingInStorage: Datenbankeintrag UND Storage-Objekt beide
// vorhanden — reine Existenzbestätigung, KEINE Inhaltsprüfung. Dient
// Archive BAK-08 als stabile, deterministisch sortierte
// Stichprobengrundlage (nach StorageKey aufsteigend, siehe Report-
// Dokumentation oben) — BAK-08 muss dafür selbst nicht mehr
// sortieren/filtern.
ExistingInStorage []Finding `json:"existing_in_storage"`
}
// IsClean liefert true, wenn der Lauf keine Abweichungen fand (Pflicht-
// prüfung 3: "Lauf ohne Abweichungen liefert einen leeren, eindeutig als
// sauber erkennbaren Bericht" — IsClean ist genau dieses eindeutige
// Erkennungsmerkmal, statt dass ein Aufrufer beide Listen selbst auf
// Leere prüfen muss).
func (r Report) IsClean() bool {
return len(r.MissingInStorage) == 0 && len(r.OrphanedInStorage) == 0
}
// DBEntry ist ein Datenbankeintrag, wie ihn ListDBStorageKeys liefert.
type DBEntry struct {
StorageKey string
DocumentID string
RevisionID string
}
// Reconcile vergleicht dbEntries (aus file_revisions.storage_key, DMS
// FDN-02) gegen storageKeys (tatsächlich im Objekt-Storage vorhandene
// Schlüssel, z.B. per Verzeichnis-Walk des FDN-03-LocalDriver-
// Basisverzeichnisses) und liefert die Abweichungen in beide Richtungen.
// Reine Funktion — kein Datenbank-/Storage-Zugriff hier, dadurch ohne
// echte Infrastruktur testbar (siehe reconcile_test.go).
func Reconcile(dbEntries []DBEntry, storageKeys []string) Report {
storageSet := make(map[string]bool, len(storageKeys))
for _, k := range storageKeys {
storageSet[k] = true
}
dbSet := make(map[string]DBEntry, len(dbEntries))
for _, e := range dbEntries {
dbSet[e.StorageKey] = e
}
var missing, existing []Finding
for _, e := range dbEntries {
if !storageSet[e.StorageKey] {
missing = append(missing, Finding(e))
} else {
existing = append(existing, Finding(e))
}
}
var orphaned []Finding
for _, k := range storageKeys {
if _, ok := dbSet[k]; !ok {
orphaned = append(orphaned, Finding{StorageKey: k})
}
}
sort.Slice(missing, func(i, j int) bool { return missing[i].StorageKey < missing[j].StorageKey })
sort.Slice(orphaned, func(i, j int) bool { return orphaned[i].StorageKey < orphaned[j].StorageKey })
sort.Slice(existing, func(i, j int) bool { return existing[i].StorageKey < existing[j].StorageKey })
return Report{
GeneratedAt: time.Now().UTC(),
MissingInStorage: missing,
OrphanedInStorage: orphaned,
ExistingInStorage: existing,
}
}
@@ -1,169 +0,0 @@
package reconcile
import "testing"
// TestReconcile_DetectsMissingInStorage ist Akzeptanzkriterium 1 / Pruefung
// 1: ein Datenbankeintrag ohne zugehoeriges Objekt im Storage wird erkannt.
func TestReconcile_DetectsMissingInStorage(t *testing.T) {
db := []DBEntry{
{StorageKey: "documents/d1/revisions/r1", DocumentID: "d1", RevisionID: "r1"},
{StorageKey: "documents/d2/revisions/r1", DocumentID: "d2", RevisionID: "r1"},
}
storage := []string{"documents/d1/revisions/r1"} // d2/r1 fehlt absichtlich
report := Reconcile(db, storage)
if len(report.MissingInStorage) != 1 {
t.Fatalf("erwartet 1 fund in missing_in_storage, habe %d: %+v", len(report.MissingInStorage), report.MissingInStorage)
}
if report.MissingInStorage[0].StorageKey != "documents/d2/revisions/r1" {
t.Fatalf("unerwarteter fund: %+v", report.MissingInStorage[0])
}
if len(report.OrphanedInStorage) != 0 {
t.Fatalf("erwartet 0 funde in orphaned_in_storage, habe %d", len(report.OrphanedInStorage))
}
}
// TestReconcile_DetectsOrphanedInStorage ist Akzeptanzkriterium 2 /
// Pruefung 2: ein Storage-Objekt ohne Datenbankeintrag wird erkannt.
func TestReconcile_DetectsOrphanedInStorage(t *testing.T) {
db := []DBEntry{
{StorageKey: "documents/d1/revisions/r1", DocumentID: "d1", RevisionID: "r1"},
}
storage := []string{
"documents/d1/revisions/r1",
"documents/verwaist/revisions/r1", // kein DB-Eintrag dafuer
}
report := Reconcile(db, storage)
if len(report.OrphanedInStorage) != 1 {
t.Fatalf("erwartet 1 fund in orphaned_in_storage, habe %d: %+v", len(report.OrphanedInStorage), report.OrphanedInStorage)
}
if report.OrphanedInStorage[0].StorageKey != "documents/verwaist/revisions/r1" {
t.Fatalf("unerwarteter fund: %+v", report.OrphanedInStorage[0])
}
if len(report.MissingInStorage) != 0 {
t.Fatalf("erwartet 0 funde in missing_in_storage, habe %d", len(report.MissingInStorage))
}
}
// TestReconcile_CleanRunProducesEmptyReport ist Pruefung 3: Lauf ohne
// Abweichungen liefert einen leeren, eindeutig als sauber erkennbaren
// Bericht.
func TestReconcile_CleanRunProducesEmptyReport(t *testing.T) {
db := []DBEntry{
{StorageKey: "documents/d1/revisions/r1", DocumentID: "d1", RevisionID: "r1"},
{StorageKey: "documents/d2/revisions/r1", DocumentID: "d2", RevisionID: "r1"},
}
storage := []string{"documents/d1/revisions/r1", "documents/d2/revisions/r1"}
report := Reconcile(db, storage)
if !report.IsClean() {
t.Fatalf("erwartet sauberen bericht, habe missing=%v orphaned=%v", report.MissingInStorage, report.OrphanedInStorage)
}
if len(report.MissingInStorage) != 0 || len(report.OrphanedInStorage) != 0 {
t.Fatal("IsClean()==true, aber listen sind nicht leer - widerspruch")
}
}
// TestReconcile_ExistingButCorruptedObjectProducesNoFinding ist der
// Nachweis, dass BAK-05 AUSSCHLIESSLICH Existenz prueft, niemals Inhalt
// (die Fehlerklasse "existiert, aber Inhalt beschaedigt" ist Archive
// BAK-08, bewusst nicht hier) — Reconcile bekommt nur SCHLUESSEL, hat gar
// keine Moeglichkeit, auf Inhalt zuzugreifen; dieser Test dokumentiert die
// Absicht explizit, damit sie nicht versehentlich spaeter aufgeweicht wird.
func TestReconcile_ExistingButCorruptedObjectProducesNoFinding(t *testing.T) {
db := []DBEntry{
{StorageKey: "documents/d1/revisions/r1", DocumentID: "d1", RevisionID: "r1"},
}
// "korruptes" Objekt hier rein simuliert durch denselben Schluessel -
// Reconcile kennt und prueft keinen Inhalt, nur den Schluessel selbst.
storage := []string{"documents/d1/revisions/r1"}
report := Reconcile(db, storage)
if !report.IsClean() {
t.Fatalf("ein existierendes (wenn auch inhaltlich korruptes) objekt haette KEINEN befund ausloesen duerfen, habe: %+v", report)
}
}
// TestReconcile_ExistingInStorageIsStableSamplingBasis ist der Nachweis,
// dass Reconcile eine deterministisch sortierte Liste ALLER bestaetigt
// existierenden Objekte liefert (DB-Eintrag UND Storage-Objekt vorhanden)
// - dies ist die Stichprobengrundlage, die Archive BAK-08 weiterverwendet,
// ohne selbst neu zu sortieren/filtern.
func TestReconcile_ExistingInStorageIsStableSamplingBasis(t *testing.T) {
db := []DBEntry{
{StorageKey: "documents/z/revisions/r1", DocumentID: "z", RevisionID: "r1"},
{StorageKey: "documents/a/revisions/r1", DocumentID: "a", RevisionID: "r1"},
{StorageKey: "documents/fehlt/revisions/r1", DocumentID: "fehlt", RevisionID: "r1"},
}
storage := []string{
"documents/z/revisions/r1",
"documents/a/revisions/r1",
}
report := Reconcile(db, storage)
want := []string{"documents/a/revisions/r1", "documents/z/revisions/r1"}
if len(report.ExistingInStorage) != len(want) {
t.Fatalf("erwartet %d bestaetigt existierende objekte, habe %d: %+v", len(want), len(report.ExistingInStorage), report.ExistingInStorage)
}
for i, w := range want {
if report.ExistingInStorage[i].StorageKey != w {
t.Fatalf("sortierreihenfolge falsch: %v, want beginnend mit %v", report.ExistingInStorage, want)
}
}
if len(report.MissingInStorage) != 1 || report.MissingInStorage[0].StorageKey != "documents/fehlt/revisions/r1" {
t.Fatalf("missing_in_storage unerwartet: %+v", report.MissingInStorage)
}
}
// TestReconcile_DeterministicOrdering ist der Nachweis fuer die
// Stabilitaets-Anforderung: gleiche Eingabe liefert bei mehreren Laeufen
// IMMER dieselbe Reihenfolge (Voraussetzung dafuer, dass Archive BAK-08
// die Liste der existierenden Objekte stabil weiterverarbeiten kann, ohne
// selbst neu zu sortieren/filtern).
func TestReconcile_DeterministicOrdering(t *testing.T) {
db := []DBEntry{
{StorageKey: "documents/z/revisions/r1", DocumentID: "z", RevisionID: "r1"},
{StorageKey: "documents/a/revisions/r1", DocumentID: "a", RevisionID: "r1"},
{StorageKey: "documents/m/revisions/r1", DocumentID: "m", RevisionID: "r1"},
}
storage := []string{
"documents/a/revisions/r1", // deckt genau den DB-Eintrag "a" ab
"documents/y/revisions/r1",
"documents/n/revisions/r1",
}
first := Reconcile(db, storage)
second := Reconcile(db, storage)
if len(first.MissingInStorage) != len(second.MissingInStorage) {
t.Fatal("unterschiedliche anzahl funde zwischen zwei laeufen mit identischer eingabe")
}
for i := range first.MissingInStorage {
if first.MissingInStorage[i].StorageKey != second.MissingInStorage[i].StorageKey {
t.Fatalf("reihenfolge in missing_in_storage nicht deterministisch: lauf1[%d]=%q lauf2[%d]=%q",
i, first.MissingInStorage[i].StorageKey, i, second.MissingInStorage[i].StorageKey)
}
}
for i := range first.OrphanedInStorage {
if first.OrphanedInStorage[i].StorageKey != second.OrphanedInStorage[i].StorageKey {
t.Fatalf("reihenfolge in orphaned_in_storage nicht deterministisch: lauf1[%d]=%q lauf2[%d]=%q",
i, first.OrphanedInStorage[i].StorageKey, i, second.OrphanedInStorage[i].StorageKey)
}
}
// Aufsteigend sortiert (a < m < z), nicht Einfuegereihenfolge.
wantOrder := []string{"documents/m/revisions/r1", "documents/z/revisions/r1"}
if len(first.MissingInStorage) != len(wantOrder) {
t.Fatalf("erwartet %d funde, habe %d", len(wantOrder), len(first.MissingInStorage))
}
for i, w := range wantOrder {
if first.MissingInStorage[i].StorageKey != w {
t.Fatalf("sortierreihenfolge falsch: %v, want beginnend mit %v", first.MissingInStorage, wantOrder)
}
}
}
-65
View File
@@ -1,65 +0,0 @@
package reconcile
import (
"context"
"fmt"
"os"
"path/filepath"
"github.com/jackc/pgx/v5/pgxpool"
)
// ListDBStorageKeys liest alle storage_key-Werte aus file_revisions
// (DMS FDN-02) — Archive liest direkt aus derselben physischen
// Tenant-Datenbank (Modell C, Core TEN-01), OHNE DMS-Go-Pakete zu
// importieren (Archive ist ein eigenes Go-Modul) — reiner SQL-Zugriff
// gegen das dokumentierte Schema, sortiert nach storage_key für
// deterministische Reconcile-Ergebnisse.
func ListDBStorageKeys(ctx context.Context, pool *pgxpool.Pool) ([]DBEntry, error) {
rows, err := pool.Query(ctx, `
SELECT storage_key, document_id, id FROM file_revisions ORDER BY storage_key
`)
if err != nil {
return nil, fmt.Errorf("reconcile: file_revisions abfragen: %w", err)
}
defer rows.Close()
var entries []DBEntry
for rows.Next() {
var e DBEntry
if err := rows.Scan(&e.StorageKey, &e.DocumentID, &e.RevisionID); err != nil {
return nil, fmt.Errorf("reconcile: file_revisions-zeile lesen: %w", err)
}
entries = append(entries, e)
}
return entries, rows.Err()
}
// ListStorageObjects durchläuft den lokalen FDN-03-LocalDriver-
// Basisordner und liefert alle vorhandenen Objektschlüssel (Pfad relativ
// zu baseDir, mit "/" als Trenner — dasselbe Format wie
// storage.ObjectKey aus FDN-03), sortiert.
func ListStorageObjects(baseDir string) ([]string, error) {
var keys []string
err := filepath.WalkDir(baseDir, func(path string, d os.DirEntry, err error) error {
if err != nil {
return err
}
if d.IsDir() {
return nil
}
rel, err := filepath.Rel(baseDir, path)
if err != nil {
return err
}
keys = append(keys, filepath.ToSlash(rel))
return nil
})
if err != nil {
if os.IsNotExist(err) {
return nil, nil
}
return nil, fmt.Errorf("reconcile: objekt-storage durchlaufen: %w", err)
}
return keys, nil
}
-132
View File
@@ -1,132 +0,0 @@
package reconcile
import (
"context"
"os"
"path/filepath"
"testing"
"github.com/jackc/pgx/v5/pgxpool"
)
func requireTestPool(t *testing.T) *pgxpool.Pool {
t.Helper()
dsn := os.Getenv("TEST_TENANT_DSN")
if dsn == "" {
t.Skip("TEST_TENANT_DSN nicht gesetzt, Integrationstest uebersprungen")
}
ctx := context.Background()
pool, err := pgxpool.New(ctx, dsn)
if err != nil {
t.Fatalf("pool: %v", err)
}
t.Cleanup(func() { pool.Close() })
// Minimalschema, das exakt DMS FDN-02s file_revisions-Spalten spiegelt
// (Archive kann DMS' internal/-Pakete als eigenes Go-Modul nicht
// importieren, daher hier als Testfixture kopiert statt real migriert).
if _, err := pool.Exec(ctx, `
CREATE EXTENSION IF NOT EXISTS pgcrypto;
CREATE TABLE IF NOT EXISTS users (
id UUID PRIMARY KEY DEFAULT gen_random_uuid(), email TEXT NOT NULL UNIQUE, name TEXT NOT NULL,
created_at TIMESTAMPTZ NOT NULL DEFAULT now()
);
CREATE TABLE IF NOT EXISTS documents (
id UUID PRIMARY KEY DEFAULT gen_random_uuid(), title TEXT NOT NULL,
created_by UUID NOT NULL REFERENCES users(id), created_at TIMESTAMPTZ NOT NULL DEFAULT now()
);
CREATE TABLE IF NOT EXISTS file_revisions (
id UUID PRIMARY KEY DEFAULT gen_random_uuid(), document_id UUID NOT NULL REFERENCES documents(id) ON DELETE CASCADE,
storage_key TEXT NOT NULL, checksum_sha256 TEXT NOT NULL, size_bytes BIGINT NOT NULL,
mime_type TEXT NOT NULL, revision_number INTEGER NOT NULL, created_by UUID NOT NULL REFERENCES users(id),
created_at TIMESTAMPTZ NOT NULL DEFAULT now()
);
`); err != nil {
t.Fatalf("schema: %v", err)
}
t.Cleanup(func() {
_, _ = pool.Exec(context.Background(), `TRUNCATE file_revisions, documents, users CASCADE`)
})
return pool
}
// TestListDBStorageKeys_ReadsRealFileRevisions ist der Nachweis, dass
// ListDBStorageKeys tatsaechlich gegen eine echte Postgres-Instanz mit
// DMS-FDN-02-Schema liest — kein Mock.
func TestListDBStorageKeys_ReadsRealFileRevisions(t *testing.T) {
pool := requireTestPool(t)
ctx := context.Background()
var userID, docID string
if err := pool.QueryRow(ctx, `INSERT INTO users (email, name) VALUES ('reconcile-test@example.test', 'Test') RETURNING id`).Scan(&userID); err != nil {
t.Fatalf("testbenutzer anlegen: %v", err)
}
if err := pool.QueryRow(ctx, `INSERT INTO documents (title, created_by) VALUES ('doc', $1) RETURNING id`, userID).Scan(&docID); err != nil {
t.Fatalf("testdokument anlegen: %v", err)
}
if _, err := pool.Exec(ctx, `
INSERT INTO file_revisions (document_id, storage_key, checksum_sha256, size_bytes, mime_type, revision_number, created_by)
VALUES ($1, 'documents/x/revisions/1', 'abc', 10, 'text/plain', 1, $2)
`, docID, userID); err != nil {
t.Fatalf("testrevision anlegen: %v", err)
}
entries, err := ListDBStorageKeys(ctx, pool)
if err != nil {
t.Fatalf("listdbstoragekeys: %v", err)
}
if len(entries) != 1 {
t.Fatalf("erwartet 1 eintrag, habe %d", len(entries))
}
if entries[0].StorageKey != "documents/x/revisions/1" {
t.Fatalf("storage_key = %q, want %q", entries[0].StorageKey, "documents/x/revisions/1")
}
if entries[0].DocumentID != docID {
t.Fatalf("document_id = %q, want %q", entries[0].DocumentID, docID)
}
}
// TestListStorageObjects_WalksRealDirectory ist der Nachweis, dass
// ListStorageObjects tatsaechlich das Dateisystem durchlaeuft.
func TestListStorageObjects_WalksRealDirectory(t *testing.T) {
baseDir := t.TempDir()
mustWriteFile(t, filepath.Join(baseDir, "documents", "d1", "revisions", "r1"), "inhalt")
mustWriteFile(t, filepath.Join(baseDir, "documents", "d2", "revisions", "r1"), "inhalt")
keys, err := ListStorageObjects(baseDir)
if err != nil {
t.Fatalf("liststorageobjects: %v", err)
}
if len(keys) != 2 {
t.Fatalf("erwartet 2 objektschluessel, habe %d: %v", len(keys), keys)
}
want := []string{"documents/d1/revisions/r1", "documents/d2/revisions/r1"}
for i, w := range want {
if keys[i] != w {
t.Fatalf("schluessel[%d] = %q, want %q (voll: %v)", i, keys[i], w, keys)
}
}
}
// TestListStorageObjects_MissingDirectoryReturnsEmpty prueft das
// Verhalten, wenn das Basisverzeichnis (noch) gar nicht existiert -
// sollte als "keine Objekte", nicht als Fehler behandelt werden.
func TestListStorageObjects_MissingDirectoryReturnsEmpty(t *testing.T) {
keys, err := ListStorageObjects("/pfad/der/nicht/existiert/fuer/diesen/test")
if err != nil {
t.Fatalf("erwartet keinen fehler bei fehlendem verzeichnis, habe: %v", err)
}
if len(keys) != 0 {
t.Fatalf("erwartet 0 schluessel, habe %d", len(keys))
}
}
func mustWriteFile(t *testing.T, path, content string) {
t.Helper()
if err := os.MkdirAll(filepath.Dir(path), 0o755); err != nil {
t.Fatalf("verzeichnis anlegen: %v", err)
}
if err := os.WriteFile(path, []byte(content), 0o600); err != nil {
t.Fatalf("datei schreiben: %v", err)
}
}
+18 -3
View File
@@ -7,6 +7,7 @@ import (
"gitea.perlbach24.de/scripte/nexarch/internal/config"
"gitea.perlbach24.de/scripte/nexarch/internal/db"
"gitea.perlbach24.de/scripte/nexarch/internal/tenant"
)
func main() {
@@ -15,16 +16,30 @@ func main() {
log.Fatalf("config: %v", err)
}
pool, err := db.Connect(context.Background(), cfg.RegistryDSN)
ctx := context.Background()
registryPool, err := db.Connect(ctx, cfg.RegistryDSN)
if err != nil {
log.Fatalf("db: %v", err)
log.Fatalf("registry db: %v", err)
}
defer pool.Close()
defer registryPool.Close()
adminPool, err := db.Connect(ctx, cfg.AdminDSN)
if err != nil {
log.Fatalf("admin db: %v", err)
}
defer adminPool.Close()
registry := tenant.NewRegistry(registryPool)
provisioner := tenant.NewProvisioner(adminPool, registry, cfg.TenantDSNTemplate)
tenantHandler := tenant.NewHandler(provisioner)
mux := http.NewServeMux()
mux.HandleFunc("/healthz", func(w http.ResponseWriter, r *http.Request) {
w.WriteHeader(http.StatusOK)
})
// Vorlaeufiger Pfad ohne Versionierung/Auth — wird mit API-01/IAM-01 abgeloest.
mux.HandleFunc("/internal/tenants", tenantHandler.CreateTenant)
log.Printf("nexarch-core listening on %s", cfg.ListenAddr)
if err := http.ListenAndServe(cfg.ListenAddr, mux); err != nil {
+82
View File
@@ -0,0 +1,82 @@
// kek-api ist der Aufrufpunkt fuer API-12: startet den bereits fertigen
// API-10-Handler (internal/kek) als eigenstaendigen HTTP-Dienst.
// REINES WIRING — keine Aenderung an internal/kek/, internal/moduleregistry/
// oder internal/tenant/.
package main
import (
"context"
"log"
"net/http"
"os"
"time"
"github.com/jackc/pgx/v5/pgxpool"
"gitea.perlbach24.de/scripte/nexarch/internal/flag"
"gitea.perlbach24.de/scripte/nexarch/internal/kek"
"gitea.perlbach24.de/scripte/nexarch/internal/moduleregistry"
"gitea.perlbach24.de/scripte/nexarch/internal/tenant"
)
// tenantResolverAdapter erfüllt kek.TenantResolver über den bestehenden
// tenant.Registry.GetBySlug-Zugriff — kein neuer Tenant-Code, nur
// Signatur-Anpassung.
type tenantResolverAdapter struct{ registry *tenant.Registry }
func (a tenantResolverAdapter) ResolveTenantID(ctx context.Context, tenantSlug string) (string, error) {
t, err := a.registry.GetBySlug(ctx, tenantSlug)
if err != nil {
return "", err
}
return t.ID, nil
}
func requireEnv(name string) string {
v := os.Getenv(name)
if v == "" {
log.Fatalf("%s muss gesetzt sein", name)
}
return v
}
func main() {
registryDSN := requireEnv("NEXARCH_KEK_REGISTRY_DSN")
masterKeyEnvVar := os.Getenv("NEXARCH_KEK_MASTER_KEY_ENV")
if masterKeyEnvVar == "" {
masterKeyEnvVar = "NEXARCH_KEK_MASTER_KEY"
}
addr := os.Getenv("NEXARCH_KEK_API_LISTEN_ADDR")
if addr == "" {
addr = "127.0.0.1:8102"
}
masterKey, err := kek.LoadMasterKeyFromEnv(masterKeyEnvVar)
if err != nil {
log.Fatalf("master-key laden: %v", err)
}
ctx := context.Background()
pool, err := pgxpool.New(ctx, registryDSN)
if err != nil {
log.Fatalf("datenbankverbindung: %v", err)
}
defer pool.Close()
flagStore := flag.NewStore(pool)
flagService := flag.NewService(flagStore, 30*time.Second)
registry := moduleregistry.NewRegistry(pool, flagService)
tenantRegistry := tenant.NewRegistry(pool)
store := kek.NewStore(pool)
handler := kek.NewHandler(store, masterKey, registry, registry, tenantResolverAdapter{registry: tenantRegistry})
mux := http.NewServeMux()
mux.HandleFunc("/healthz", func(w http.ResponseWriter, r *http.Request) { w.WriteHeader(http.StatusOK) })
mux.HandleFunc("/internal/kek/tenant", handler.TenantKEKHandler)
log.Printf("kek-api: listening on %s", addr)
if err := http.ListenAndServe(addr, mux); err != nil {
log.Fatalf("http server: %v", err)
}
}
@@ -1,9 +0,0 @@
[Unit]
Description=NEXARCH Archive - Datenbank-Vollsicherung (BAK-01)
After=network.target postgresql.service
[Service]
Type=oneshot
User=nexarch
EnvironmentFile=/etc/nexarch/archive-backup.env
ExecStart=__INSTALL_DIR__/bin/backup-cli full
@@ -1,9 +0,0 @@
[Unit]
Description=Taeglicher Zeitplan fuer NEXARCH Archive Datenbank-Vollsicherung (BAK-01)
[Timer]
OnCalendar=*-*-* 02:00:00
Persistent=true
[Install]
WantedBy=timers.target
@@ -1,9 +0,0 @@
[Unit]
Description=NEXARCH Archive - Datenbank-Inkrementalsicherung (BAK-01)
After=network.target postgresql.service
[Service]
Type=oneshot
User=nexarch
EnvironmentFile=/etc/nexarch/archive-backup.env
ExecStart=__INSTALL_DIR__/bin/backup-cli incremental
@@ -1,9 +0,0 @@
[Unit]
Description=Stuendlicher Zeitplan fuer NEXARCH Archive Datenbank-Inkrementalsicherung (BAK-01)
[Timer]
OnCalendar=*-*-* *:00:00
Persistent=true
[Install]
WantedBy=timers.target
@@ -1,9 +0,0 @@
[Unit]
Description=NEXARCH Archive - Sicherungsgenerationen-Rotation (BAK-01)
After=network.target
[Service]
Type=oneshot
User=nexarch
EnvironmentFile=/etc/nexarch/archive-backup.env
ExecStart=__INSTALL_DIR__/bin/backup-cli rotate
@@ -1,9 +0,0 @@
[Unit]
Description=Taeglicher Zeitplan fuer NEXARCH Archive Sicherungsgenerationen-Rotation (BAK-01)
[Timer]
OnCalendar=*-*-* 03:00:00
Persistent=true
[Install]
WantedBy=timers.target
@@ -1,9 +0,0 @@
[Unit]
Description=NEXARCH Archive - Objekt-Storage-Sicherung (BAK-02)
After=network.target
[Service]
Type=oneshot
User=nexarch
EnvironmentFile=/etc/nexarch/archive-objectbackup.env
ExecStart=__INSTALL_DIR__/bin/objectbackup-cli backup __OBJECT_SOURCE_DIR__
@@ -1,9 +0,0 @@
[Unit]
Description=Stuendlicher Zeitplan fuer NEXARCH Archive Objekt-Storage-Sicherung (BAK-02)
[Timer]
OnCalendar=*-*-* *:30:00
Persistent=true
[Install]
WantedBy=timers.target
@@ -1,9 +0,0 @@
[Unit]
Description=NEXARCH Archive - Objekt-Storage-Sicherung Vollstaendigkeitspruefung (BAK-02)
After=network.target
[Service]
Type=oneshot
User=nexarch
EnvironmentFile=/etc/nexarch/archive-objectbackup.env
ExecStart=__INSTALL_DIR__/bin/objectbackup-cli check
@@ -1,9 +0,0 @@
[Unit]
Description=Woechentlicher Zeitplan fuer NEXARCH Archive Objekt-Storage-Vollstaendigkeitspruefung (BAK-02)
[Timer]
OnCalendar=Sun *-*-* 04:00:00
Persistent=true
[Install]
WantedBy=timers.target
@@ -1,9 +0,0 @@
[Unit]
Description=NEXARCH Archive - Objekt-Storage-Sicherung Rotation (BAK-02)
After=network.target
[Service]
Type=oneshot
User=nexarch
EnvironmentFile=/etc/nexarch/archive-objectbackup.env
ExecStart=__INSTALL_DIR__/bin/objectbackup-cli rotate
@@ -1,9 +0,0 @@
[Unit]
Description=Taeglicher Zeitplan fuer NEXARCH Archive Objekt-Storage-Rotation (BAK-02)
[Timer]
OnCalendar=*-*-* 03:30:00
Persistent=true
[Install]
WantedBy=timers.target
@@ -1,10 +0,0 @@
[Unit]
Description=NEXARCH Archive - Konsistenzpruefung Storage vs. DB (BAK-05)
After=network.target postgresql.service
[Service]
Type=oneshot
User=nexarch
EnvironmentFile=/etc/nexarch/archive-reconcile.env
ExecStart=__INSTALL_DIR__/bin/reconcile-cli
StandardOutput=journal
@@ -1,9 +0,0 @@
[Unit]
Description=Taeglicher Zeitplan fuer NEXARCH Archive Konsistenzpruefung (BAK-05)
[Timer]
OnCalendar=*-*-* 05:00:00
Persistent=true
[Install]
WantedBy=timers.target
@@ -0,0 +1,14 @@
[Unit]
Description=NEXARCH Core - KEK-Bezugsdienst (API-10/API-12)
After=network.target postgresql.service
[Service]
Type=simple
User=nexarch
EnvironmentFile=/etc/nexarch/kek-api.env
ExecStart=__INSTALL_DIR__/bin/kek-api
Restart=on-failure
StandardOutput=journal
[Install]
WantedBy=multi-user.target
+61
View File
@@ -0,0 +1,61 @@
# API-12 Prüfprotokoll: KEK-Bezugsdienst starten (API-10 als laufender Dienst)
Voraussetzung API-10 bereits Fertig, hier UNVERÄNDERT.
## Reines Wiring, keine neue Logik
`git diff --stat internal/kek/ internal/moduleregistry/ internal/tenant/`
liefert KEINEN Diff. `cmd/kek-api/main.go` setzt ausschließlich
bestehende Konstruktoren zusammen; `tenantResolverAdapter` bildet nur
`tenant.Registry.GetBySlug` auf `kek.TenantResolver` ab (Signatur-
Anpassung, kein neuer Fachcode).
## Umsetzung
- `cmd/kek-api/main.go` `POST /internal/kek/tenant?tenant=<slug>`,
authentifiziert über dasselbe Service-Credential-Verfahren wie jeder
andere Modul-Core-Aufruf (API-02), zusätzlich Tenant-Aktivierungs-
prüfung (identisches Muster wie in `internal/kek.Handler` bereits
vorgesehen).
- `deploy/systemd/nexarch-kek-api.service.tmpl`.
## Prüfungen
| # | Prüfung | Ergebnis |
|---|---|---|
| 1 | Dienst startet und bleibt stabil (systemctl status aktiv) | **bestanden** real auf 131: `nexarch-kek-api.service` aktiv |
| 2 | Realer Aufruf mit gültigem Service-Credential liefert den erwarteten Tenant-KEK, ohne/mit falschem Credential wird abgelehnt | **bestanden** real per `curl`: echtes Modul registriert+provisioniert, echter Tenant-KEK über `kek.Store.CreateForTenant` erzeugt (Klartext-Hex zum Vergleich notiert) — Aufruf mit korrektem Credential liefert exakt denselben KEK (Base64-dekodiert übereinstimmend mit dem erzeugten Hex-Wert verifiziert); Aufruf mit falschem Credential → 403 |
| 3 | Code-Review: keine Änderung an internal/kek/ selbst, nur main.go+systemd neu | **bestanden** `git diff --stat` bestätigt: `internal/kek/`, `internal/moduleregistry/`, `internal/tenant/` unverändert |
## Echte Verdrahtung auf 192.168.1.131
- `kek-api` gebaut nach `/opt/nexarch-core/bin/`,
`/etc/nexarch/kek-api.env` (0600, echter zufälliger 32-Byte-
Master-Key), Dienst installiert/aktiviert.
- Reale Grant-Lücke gefunden und behoben (gleiches Muster wie zuvor):
`nexarch_core` hatte keine Rechte auf `tenant_keks``GRANT`
nachgezogen und über `information_schema.role_table_grants`
verifiziert.
- End-zu-Ende-Beweis: echtes Modul registriert, Service-Credential
provisioniert, echter Tenant + Tenant-KEK real erzeugt, `curl` gegen
den laufenden Dienst liefert exakt diesen KEK zurück (Byte-für-Byte
verglichen), falsches Credential real abgelehnt. Testdaten
anschließend entfernt.
## Build/Test-Ergebnis (192.168.1.131)
```
go build ./... -> clean
go vet ./... -> clean
golangci-lint run ./cmd/kek-api/... -> 0 issues
```
Keine neuen Go-Tests nötig (kein neuer Fachcode außer main.go/Adapter,
die eigentliche Logik ist bereits durch API-10s eigene Tests
abgedeckt).
## Gesamtergebnis
**Bestanden.** API-10 ist jetzt ein real laufender, über systemd
verwalteter Dienst — Voraussetzung für Mail ARC-02 und künftig DMS
FDN-09-Nachnutzung.
+9
View File
@@ -3,3 +3,12 @@ module gitea.perlbach24.de/scripte/nexarch
go 1.22
require github.com/jackc/pgx/v5 v5.6.0
require (
github.com/jackc/pgpassfile v1.0.0 // indirect
github.com/jackc/pgservicefile v0.0.0-20221227161230-091c0ba34f0a // indirect
github.com/jackc/puddle/v2 v2.2.1 // indirect
golang.org/x/crypto v0.17.0 // indirect
golang.org/x/sync v0.1.0 // indirect
golang.org/x/text v0.14.0 // indirect
)
View File
+24 -2
View File
@@ -10,8 +10,15 @@ import (
// connection info, superadmin accounts) — see nexarch-state.json
// multi_tenancy: Modell C (physisch getrennte DB pro Mandant).
type Config struct {
ListenAddr string
ListenAddr string
// RegistryDSN verbindet zur Control-Plane-Registry-Datenbank.
RegistryDSN string
// AdminDSN verbindet zur Wartungsdatenbank (z.B. "postgres") und wird nur
// fuer CREATE/DROP DATABASE beim Tenant-Provisioning verwendet.
AdminDSN string
// TenantDSNTemplate enthaelt genau ein "%s" als Platzhalter fuer den
// Datenbanknamen einer neu provisionierten Tenant-Datenbank.
TenantDSNTemplate string
}
func Load() (Config, error) {
@@ -20,10 +27,25 @@ func Load() (Config, error) {
return Config{}, fmt.Errorf("NEXARCH_REGISTRY_DSN not set")
}
adminDSN := os.Getenv("NEXARCH_ADMIN_DSN")
if adminDSN == "" {
return Config{}, fmt.Errorf("NEXARCH_ADMIN_DSN not set")
}
dsnTemplate := os.Getenv("NEXARCH_TENANT_DSN_TEMPLATE")
if dsnTemplate == "" {
return Config{}, fmt.Errorf("NEXARCH_TENANT_DSN_TEMPLATE not set")
}
addr := os.Getenv("NEXARCH_LISTEN_ADDR")
if addr == "" {
addr = ":8080"
}
return Config{ListenAddr: addr, RegistryDSN: dsn}, nil
return Config{
ListenAddr: addr,
RegistryDSN: dsn,
AdminDSN: adminDSN,
TenantDSNTemplate: dsnTemplate,
}, nil
}
+87
View File
@@ -0,0 +1,87 @@
// Package flag implementiert Core LIC-02: einen Feature-Flag-Dienst mit
// Strategien (global an/aus, Prozentsatz, Tenant-Zielgruppe) als Kernfunktion
// des Core-Dienstes selbst — keine zusaetzliche Infrastruktur (Unleash-Server
// + eigene DB), siehe "bewusst vermeiden" im LIC-02-Ticket.
package flag
import (
"context"
"errors"
"fmt"
"hash/fnv"
"github.com/jackc/pgx/v5"
"github.com/jackc/pgx/v5/pgxpool"
)
var ErrNotFound = errors.New("flag: nicht gefunden")
// Flag ist die zentrale Definition — Auswertung (Evaluate) ist bewusst davon
// getrennt (Unleash-Prinzip: Flag-Verwaltung vs. Flag-Auswertung).
type Flag struct {
Key string
Enabled bool
RolloutPercentage int
TargetTenantSlugs []string
}
// Store ist die Verwaltungsseite (Admin): Flags definieren/lesen.
type Store struct {
pool *pgxpool.Pool
}
func NewStore(pool *pgxpool.Pool) *Store {
return &Store{pool: pool}
}
func (s *Store) Set(ctx context.Context, f Flag) error {
if f.TargetTenantSlugs == nil {
f.TargetTenantSlugs = []string{} // pgx uebertraegt ein nil-Slice sonst als SQL NULL statt leerem Array.
}
_, err := s.pool.Exec(ctx, `
INSERT INTO feature_flags (key, enabled, rollout_percentage, target_tenant_slugs, updated_at)
VALUES ($1, $2, $3, $4, now())
ON CONFLICT (key) DO UPDATE SET
enabled = $2, rollout_percentage = $3, target_tenant_slugs = $4, updated_at = now()
`, f.Key, f.Enabled, f.RolloutPercentage, f.TargetTenantSlugs)
if err != nil {
return fmt.Errorf("flag speichern: %w", err)
}
return nil
}
func (s *Store) Get(ctx context.Context, key string) (Flag, error) {
var f Flag
row := s.pool.QueryRow(ctx, `
SELECT key, enabled, rollout_percentage, target_tenant_slugs
FROM feature_flags WHERE key = $1
`, key)
if err := row.Scan(&f.Key, &f.Enabled, &f.RolloutPercentage, &f.TargetTenantSlugs); err != nil {
if errors.Is(err, pgx.ErrNoRows) {
return Flag{}, ErrNotFound
}
return Flag{}, fmt.Errorf("flag lesen: %w", err)
}
return f, nil
}
// evaluate wendet die Strategien in fester Reihenfolge an: globaler
// An/Aus-Schalter zuerst, dann Tenant-Zielgruppe, dann Prozentsatz-Rollout.
// Ein unbekannter/nicht getroffener Fall ergibt false — Fail-Safe-Default,
// kein Feature wird versehentlich aktiv.
func evaluate(f Flag, tenantSlug string) bool {
if f.Enabled {
return true
}
for _, target := range f.TargetTenantSlugs {
if target == tenantSlug {
return true
}
}
if f.RolloutPercentage > 0 {
h := fnv.New32a()
_, _ = h.Write([]byte(f.Key + "|" + tenantSlug))
return int(h.Sum32()%100) < f.RolloutPercentage
}
return false
}
+43
View File
@@ -0,0 +1,43 @@
package flag
import "testing"
func TestEvaluate_GlobalEnabled(t *testing.T) {
f := Flag{Key: "k", Enabled: true}
if !evaluate(f, "irgendein-tenant") {
t.Fatal("global aktiviertes flag sollte fuer jeden tenant true liefern")
}
}
// Akzeptanzkriterium 1 + Pruefung 2: Zielgruppen-Strategie.
func TestEvaluate_TargetTenantStrategy(t *testing.T) {
f := Flag{Key: "k", Enabled: false, TargetTenantSlugs: []string{"acme"}}
if !evaluate(f, "acme") {
t.Fatal("erwartet true fuer tenant in zielgruppe")
}
if evaluate(f, "globex") {
t.Fatal("erwartet false fuer tenant ausserhalb der zielgruppe")
}
}
func TestEvaluate_RolloutPercentageBoundaries(t *testing.T) {
full := Flag{Key: "k", RolloutPercentage: 100}
if !evaluate(full, "beliebiger-tenant-1") || !evaluate(full, "beliebiger-tenant-2") {
t.Fatal("100% rollout sollte immer true liefern")
}
none := Flag{Key: "k", RolloutPercentage: 0}
if evaluate(none, "beliebiger-tenant") {
t.Fatal("0% rollout ohne enabled/zielgruppe sollte false liefern")
}
}
func TestEvaluate_RolloutIsDeterministicPerTenant(t *testing.T) {
f := Flag{Key: "k", RolloutPercentage: 50}
first := evaluate(f, "stabiler-tenant")
for i := 0; i < 5; i++ {
if evaluate(f, "stabiler-tenant") != first {
t.Fatal("rollout-auswertung sollte fuer denselben tenant/key stabil sein")
}
}
}
+87
View File
@@ -0,0 +1,87 @@
package flag
import (
"context"
"log/slog"
"sync"
"time"
)
// DefaultCacheTTL ist die dokumentierte Cache-Invalidierungszeit
// (Akzeptanzkriterium 2/3): eine Aenderung wirkt spaetestens nach dieser
// Zeit auf allen Core-Instanzen, ohne dass ein Dienst neu gestartet werden
// muss (Akzeptanzkriterium 3).
const DefaultCacheTTL = 5 * time.Second
type cacheEntry struct {
flag Flag
expiresAt time.Time
}
// Service ist die Auswertungsseite (SDK/Client-Analogon zu Unleash) mit
// lokalem TTL-Cache. Bewusst getrennt von Store (Verwaltung).
type Service struct {
store *Store
ttl time.Duration
mu sync.RWMutex
cache map[string]cacheEntry
}
func NewService(store *Store, ttl time.Duration) *Service {
if ttl <= 0 {
ttl = DefaultCacheTTL
}
return &Service{store: store, ttl: ttl, cache: make(map[string]cacheEntry)}
}
// IsEnabled wertet ein Flag fuer einen Tenant aus. Liefert IMMER einen
// bool ohne Fehlerwert — ein nicht erreichbarer Flag-Dienst darf abhaengige
// Aufrufer nicht zum Absturz bringen oder zu Fehlerbehandlungscode zwingen,
// der leicht vergessen wird (Akzeptanzkriterium 3 / Pruefung 3: dokumentiertes
// Fallback-Verhalten = false, ggf. aus dem zuletzt bekannten Zwischenspeicher).
func (s *Service) IsEnabled(ctx context.Context, tenantSlug, key string) bool {
f, ok := s.resolve(ctx, key)
if !ok {
return false
}
return evaluate(f, tenantSlug)
}
func (s *Service) resolve(ctx context.Context, key string) (Flag, bool) {
s.mu.RLock()
entry, exists := s.cache[key]
fresh := exists && time.Now().Before(entry.expiresAt)
s.mu.RUnlock()
if fresh {
return entry.flag, true
}
f, err := s.store.Get(ctx, key)
if err != nil {
if exists {
slog.Warn("feature-flag-dienst nicht erreichbar, nutze zwischengespeicherten stand",
"flag_key", key, "error", err)
return entry.flag, true
}
slog.Warn("feature-flag-dienst nicht erreichbar, kein zwischengespeicherter stand vorhanden, fallback: deaktiviert",
"flag_key", key, "error", err)
return Flag{}, false
}
s.mu.Lock()
s.cache[key] = cacheEntry{flag: f, expiresAt: time.Now().Add(s.ttl)}
s.mu.Unlock()
return f, true
}
// Invalidate erzwingt beim naechsten IsEnabled-Aufruf ein sofortiges Neuladen
// aus der Datenbank statt auf den TTL-Ablauf zu warten — wird nach Store.Set
// auf derselben Instanz aufgerufen, damit der Schreiber die eigene Aenderung
// ohne Wartezeit sieht. Andere Core-Instanzen sehen sie spaetestens nach
// DefaultCacheTTL (siehe Akzeptanzkriterium 3).
func (s *Service) Invalidate(key string) {
s.mu.Lock()
delete(s.cache, key)
s.mu.Unlock()
}
+179
View File
@@ -0,0 +1,179 @@
package flag
import (
"context"
"os"
"testing"
"time"
"github.com/jackc/pgx/v5/pgxpool"
)
func setupFlagStoreTest(t *testing.T) (*Store, func()) {
t.Helper()
adminDSN := os.Getenv("TEST_ADMIN_DSN")
if adminDSN == "" {
t.Skip("TEST_ADMIN_DSN nicht gesetzt, Integrationstest uebersprungen")
}
ctx := context.Background()
pool, err := pgxpool.New(ctx, adminDSN)
if err != nil {
t.Fatalf("pool: %v", err)
}
if _, err := pool.Exec(ctx, `
CREATE TABLE IF NOT EXISTS feature_flags (
key TEXT PRIMARY KEY,
enabled BOOLEAN NOT NULL DEFAULT false,
rollout_percentage INT NOT NULL DEFAULT 0,
target_tenant_slugs TEXT[] NOT NULL DEFAULT '{}',
updated_at TIMESTAMPTZ NOT NULL DEFAULT now()
)`); err != nil {
t.Fatalf("schema: %v", err)
}
cleanup := func() {
_, _ = pool.Exec(ctx, `DELETE FROM feature_flags WHERE key LIKE 'test\_%' ESCAPE '\'`)
pool.Close()
}
return NewStore(pool), cleanup
}
// Akzeptanzkriterium 1 + Pruefung 2: Zielgruppen-Strategie liefert im Test
// die erwartete Auswertung.
func TestService_TargetTenantStrategy(t *testing.T) {
store, cleanup := setupFlagStoreTest(t)
defer cleanup()
ctx := context.Background()
if err := store.Set(ctx, Flag{Key: "test_target_flag", TargetTenantSlugs: []string{"acme"}}); err != nil {
t.Fatalf("set: %v", err)
}
svc := NewService(store, time.Hour)
if !svc.IsEnabled(ctx, "acme", "test_target_flag") {
t.Fatal("erwartet true fuer tenant in zielgruppe")
}
if svc.IsEnabled(ctx, "globex", "test_target_flag") {
t.Fatal("erwartet false fuer tenant ausserhalb der zielgruppe")
}
}
// Akzeptanzkriterium 2 + 3 + Pruefung 1: Flag-Aenderung wirkt innerhalb der
// dokumentierten Cache-Invalidierungszeit, automatisiert gemessen.
func TestService_CacheInvalidationTiming(t *testing.T) {
store, cleanup := setupFlagStoreTest(t)
defer cleanup()
ctx := context.Background()
const ttl = 150 * time.Millisecond
if err := store.Set(ctx, Flag{Key: "test_ttl_flag", Enabled: false}); err != nil {
t.Fatalf("set: %v", err)
}
svc := NewService(store, ttl)
if svc.IsEnabled(ctx, "acme", "test_ttl_flag") {
t.Fatal("erwartet false vor der aenderung")
}
// Aenderung "auf einer anderen instanz" simulieren: direkt ueber den
// Store, ohne svc.Invalidate aufzurufen.
changedAt := time.Now()
if err := store.Set(ctx, Flag{Key: "test_ttl_flag", Enabled: true}); err != nil {
t.Fatalf("set: %v", err)
}
// Sofort danach sollte der Cache noch den alten Stand liefern.
if svc.IsEnabled(ctx, "acme", "test_ttl_flag") {
t.Fatal("cache haette den alten (false) stand liefern sollen, direkt nach der aenderung")
}
deadline := changedAt.Add(ttl + 100*time.Millisecond)
for time.Now().Before(deadline) {
if svc.IsEnabled(ctx, "acme", "test_ttl_flag") {
elapsed := time.Since(changedAt)
t.Logf("aenderung wurde nach %s wirksam (ziel: innerhalb %s + toleranz)", elapsed, ttl)
return
}
time.Sleep(10 * time.Millisecond)
}
t.Fatalf("aenderung wurde nicht innerhalb von %s wirksam", deadline.Sub(changedAt))
}
func TestService_InvalidateForcesImmediateRefresh(t *testing.T) {
store, cleanup := setupFlagStoreTest(t)
defer cleanup()
ctx := context.Background()
if err := store.Set(ctx, Flag{Key: "test_invalidate_flag", Enabled: false}); err != nil {
t.Fatalf("set: %v", err)
}
svc := NewService(store, time.Hour) // lange TTL, damit Invalidate den unterschied macht
_ = svc.IsEnabled(ctx, "acme", "test_invalidate_flag")
if err := store.Set(ctx, Flag{Key: "test_invalidate_flag", Enabled: true}); err != nil {
t.Fatalf("set: %v", err)
}
svc.Invalidate("test_invalidate_flag")
if !svc.IsEnabled(ctx, "acme", "test_invalidate_flag") {
t.Fatal("erwartet sofort sichtbaren neuen stand nach Invalidate")
}
}
// Akzeptanzkriterium 3 + Pruefung 3: Ausfall des Flag-Dienstes fuehrt zu
// dokumentiertem Fallback-Verhalten, nicht zum Absturz.
func TestService_FallsBackOnStoreFailure(t *testing.T) {
store, cleanup := setupFlagStoreTest(t)
defer cleanup()
ctx := context.Background()
if err := store.Set(ctx, Flag{Key: "test_fallback_flag", Enabled: true}); err != nil {
t.Fatalf("set: %v", err)
}
svc := NewService(store, time.Hour)
// Cache vorwaermen, waehrend die DB noch erreichbar ist.
if !svc.IsEnabled(ctx, "acme", "test_fallback_flag") {
t.Fatal("erwartet true bei funktionierender db")
}
brokenPool, err := pgxpool.New(ctx, "postgresql://nonexistent-host-fuer-test:5432/x?connect_timeout=1")
if err != nil {
t.Fatalf("broken pool erstellen (sollte nicht sofort verbinden): %v", err)
}
brokenStore := NewStore(brokenPool)
svcWithCache := NewService(brokenStore, time.Nanosecond) // TTL sofort abgelaufen, erzwingt reload-versuch
svcWithCache.mu.Lock()
svcWithCache.cache["test_fallback_flag"] = cacheEntry{
flag: Flag{Key: "test_fallback_flag", Enabled: true},
expiresAt: time.Now().Add(-time.Hour), // bereits abgelaufen
}
svcWithCache.mu.Unlock()
func() {
defer func() {
if r := recover(); r != nil {
t.Fatalf("IsEnabled hat gepanict statt einen fallback zu liefern: %v", r)
}
}()
if !svcWithCache.IsEnabled(ctx, "acme", "test_fallback_flag") {
t.Fatal("erwartet fallback auf zwischengespeicherten (true) stand bei db-ausfall")
}
}()
// Voellig frischer Dienst ohne jeglichen cache + kaputte db -> sicherer
// default false, kein absturz.
freshSvc := NewService(brokenStore, time.Hour)
func() {
defer func() {
if r := recover(); r != nil {
t.Fatalf("IsEnabled hat gepanict: %v", r)
}
}()
if freshSvc.IsEnabled(ctx, "acme", "test_fallback_flag") {
t.Fatal("erwartet fail-safe false ohne cache und mit kaputter db")
}
}()
}
+113
View File
@@ -0,0 +1,113 @@
package kek
import (
"context"
"encoding/base64"
"encoding/json"
"errors"
"net/http"
)
// CredentialAuthenticator ist die schmale Schnittstelle zu API-02s
// Service-Credential-Pruefung (internal/moduleregistry.Registry.Authenticate).
type CredentialAuthenticator interface {
Authenticate(ctx context.Context, clientID, secret string) (moduleName string, ok bool, err error)
}
// ModuleActivationChecker ist die schmale Schnittstelle zu API-02s
// Aktivierungspruefung (internal/moduleregistry.Registry.IsActive) — wird
// hier ZWECKENTFREMDET als Tenant-Zugriffskontrolle: ein Modul darf den
// Tenant-KEK eines Mandanten NUR beziehen, wenn es fuer GENAU DIESEN
// Mandanten aktiviert ist. Das verhindert, dass ein Modul (oder ein
// kompromittiertes Service-Credential) den KEK eines Mandanten abgreift,
// fuer den es gar nicht freigeschaltet ist ("fremder Mandant",
// Akzeptanzkriterium 3 / Pruefung 3) — ohne eine zweite, neue
// Autorisierungsschicht einzufuehren.
type ModuleActivationChecker interface {
IsActive(ctx context.Context, tenantSlug, moduleName string) (bool, error)
}
// TenantResolver loest einen Tenant-Slug in seine interne ID auf
// (internal/tenant.Registry.GetBySlug, TEN-01).
type TenantResolver interface {
ResolveTenantID(ctx context.Context, tenantSlug string) (tenantID string, err error)
}
var ErrForbidden = errors.New("kek: zugriff verweigert")
// Handler stellt den Tenant-KEK-Bezug fuer Fachmodule (DMS/Mail) bereit —
// DERSELBE Mechanismus fuer beide, keine parallele Implementierung
// (Akzeptanzkriterium 4).
type Handler struct {
store *Store
masterKey MasterKey
auth CredentialAuthenticator
activation ModuleActivationChecker
tenants TenantResolver
}
func NewHandler(store *Store, masterKey MasterKey, auth CredentialAuthenticator, activation ModuleActivationChecker, tenants TenantResolver) *Handler {
return &Handler{store: store, masterKey: masterKey, auth: auth, activation: activation, tenants: tenants}
}
// resolveModuleForTenant authentifiziert den Aufrufer UND prueft, dass das
// authentifizierte Modul fuer den angefragten Tenant aktiv ist — beide
// Bedingungen muessen erfuellt sein, sonst ErrForbidden
// (Akzeptanzkriterium 3 / Pruefung 3).
func (h *Handler) resolveModuleForTenant(ctx context.Context, clientID, secret, tenantSlug string) error {
moduleName, ok, err := h.auth.Authenticate(ctx, clientID, secret)
if err != nil {
return err
}
if !ok {
return ErrForbidden
}
active, err := h.activation.IsActive(ctx, tenantSlug, moduleName)
if err != nil {
return err
}
if !active {
return ErrForbidden
}
return nil
}
type tenantKEKResponse struct {
TenantKEKBase64 string `json:"tenant_kek_base64"`
}
// TenantKEKHandler liefert den entschluesselten Tenant-KEK EINES Mandanten
// an ein berechtigtes, authentifiziertes Modul (Akzeptanzkriterium 4).
func (h *Handler) TenantKEKHandler(w http.ResponseWriter, r *http.Request) {
clientID := r.Header.Get("X-Nexarch-Client-Id")
secret := r.Header.Get("X-Nexarch-Client-Secret")
tenantSlug := r.URL.Query().Get("tenant")
if tenantSlug == "" {
http.Error(w, "tenant-parameter fehlt", http.StatusBadRequest)
return
}
if err := h.resolveModuleForTenant(r.Context(), clientID, secret, tenantSlug); err != nil {
if errors.Is(err, ErrForbidden) {
http.Error(w, "zugriff auf diesen mandanten verweigert", http.StatusForbidden)
return
}
http.Error(w, "interner fehler", http.StatusInternalServerError)
return
}
tenantID, err := h.tenants.ResolveTenantID(r.Context(), tenantSlug)
if err != nil {
http.Error(w, "mandant nicht gefunden", http.StatusNotFound)
return
}
plainKEK, err := h.store.GetDecrypted(r.Context(), tenantID, h.masterKey)
if err != nil {
http.Error(w, "tenant-kek konnte nicht ermittelt werden", http.StatusInternalServerError)
return
}
w.Header().Set("Content-Type", "application/json")
_ = json.NewEncoder(w).Encode(tenantKEKResponse{TenantKEKBase64: base64.StdEncoding.EncodeToString(plainKEK)})
}
+350
View File
@@ -0,0 +1,350 @@
package kek
import (
"bytes"
"context"
"encoding/base64"
"fmt"
"net/http"
"net/http/httptest"
"os"
"strings"
"testing"
"time"
"github.com/jackc/pgx/v5/pgxpool"
"gitea.perlbach24.de/scripte/nexarch/internal/flag"
"gitea.perlbach24.de/scripte/nexarch/internal/moduleregistry"
"gitea.perlbach24.de/scripte/nexarch/internal/tenant"
)
func setupTest(t *testing.T) (*Store, *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()
pool, err := pgxpool.New(ctx, adminDSN)
if err != nil {
t.Fatalf("pool: %v", err)
}
if _, err := pool.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()
);
CREATE TABLE IF NOT EXISTS tenant_keks (
tenant_id UUID PRIMARY KEY REFERENCES tenants(id), wrapped_kek BYTEA NOT NULL,
created_at TIMESTAMPTZ NOT NULL DEFAULT now(), rotated_at TIMESTAMPTZ
);
CREATE TABLE IF NOT EXISTS feature_flags (
key TEXT PRIMARY KEY, enabled BOOLEAN NOT NULL DEFAULT false,
rollout_percentage INT NOT NULL DEFAULT 0, target_tenant_slugs TEXT[] NOT NULL DEFAULT '{}',
updated_at TIMESTAMPTZ NOT NULL DEFAULT now()
);
CREATE TABLE IF NOT EXISTS modules (
name TEXT PRIMARY KEY, version TEXT NOT NULL CHECK (version <> ''),
required_flags TEXT[] NOT NULL DEFAULT '{}', registered_at TIMESTAMPTZ NOT NULL DEFAULT now()
);
CREATE TABLE IF NOT EXISTS module_credentials (
module_name TEXT PRIMARY KEY REFERENCES modules(name),
client_id TEXT NOT NULL UNIQUE, secret_hash BYTEA NOT NULL,
issued_at TIMESTAMPTZ NOT NULL DEFAULT now()
);
`); err != nil {
t.Fatalf("schema: %v", err)
}
cleanup := func() { pool.Close() }
return NewStore(pool), pool, cleanup
}
func newMasterKey(t *testing.T) MasterKey {
t.Helper()
key, err := generateRandomKey()
if err != nil {
t.Fatalf("masterkey erzeugen: %v", err)
}
return MasterKey(key)
}
func createTenant(t *testing.T, pool *pgxpool.Pool, slug string) string {
t.Helper()
var id string
err := pool.QueryRow(context.Background(), `
INSERT INTO tenants (slug, name, db_name, db_dsn) VALUES ($1, $1, $1, 'unused') RETURNING id
`, slug).Scan(&id)
if err != nil {
t.Fatalf("tenant anlegen: %v", err)
}
return id
}
func uniqueSlug(prefix string) string {
return fmt.Sprintf("%s_%d", prefix, time.Now().UnixNano())
}
// Akzeptanzkriterium 1: LoadMasterKeyFromEnv liest ausschliesslich aus der
// Umgebungsvariable, niemals aus Code/DB.
func TestLoadMasterKeyFromEnv(t *testing.T) {
const envVar = "NEXARCH_TEST_MASTER_KEY_API10"
t.Cleanup(func() { os.Unsetenv(envVar) })
if _, err := LoadMasterKeyFromEnv(envVar); err == nil {
t.Fatal("erwartet fehler, wenn umgebungsvariable nicht gesetzt ist")
}
os.Setenv(envVar, "zu-kurz")
if _, err := LoadMasterKeyFromEnv(envVar); err == nil {
t.Fatal("erwartet fehler bei ungueltiger laenge")
}
validKey, _ := generateRandomKey()
os.Setenv(envVar, base64.StdEncoding.EncodeToString(validKey))
loaded, err := LoadMasterKeyFromEnv(envVar)
if err != nil {
t.Fatalf("laden mit gueltigem key: %v", err)
}
if !bytes.Equal(loaded, validKey) {
t.Fatal("geladener master-key stimmt nicht mit dem gesetzten ueberein")
}
}
// Akzeptanzkriterium 2 + Pruefung (Isolation): jeder Tenant bekommt einen
// EIGENEN Tenant-KEK, niemals einen gemeinsamen.
func TestCreateForTenant_EachTenantGetsDistinctKEK(t *testing.T) {
store, pool, cleanup := setupTest(t)
defer cleanup()
ctx := context.Background()
masterKey := newMasterKey(t)
tenantA := createTenant(t, pool, uniqueSlug("acme"))
tenantB := createTenant(t, pool, uniqueSlug("globex"))
kekA, err := store.CreateForTenant(ctx, tenantA, masterKey)
if err != nil {
t.Fatalf("create a: %v", err)
}
kekB, err := store.CreateForTenant(ctx, tenantB, masterKey)
if err != nil {
t.Fatalf("create b: %v", err)
}
if bytes.Equal(kekA, kekB) {
t.Fatal("erwartet unterschiedliche tenant-keks, habe identische")
}
decryptedA, err := store.GetDecrypted(ctx, tenantA, masterKey)
if err != nil {
t.Fatalf("decrypt a: %v", err)
}
if !bytes.Equal(decryptedA, kekA) {
t.Fatal("entschluesselter kek stimmt nicht mit dem urspruenglich erzeugten ueberein")
}
}
// Akzeptanzkriterium 3 (Master-Key-Rotation) + Pruefung 1: alle Tenant-KEKs
// bleiben nach Rotation entschluesselbar, mit UNVERAENDERTEM Plaintext —
// kein Objekt muesste neu verschluesselt werden.
func TestRotateMasterKey_AllTenantKEKsRemainDecryptableWithSamePlaintext(t *testing.T) {
store, pool, cleanup := setupTest(t)
defer cleanup()
ctx := context.Background()
oldMasterKey := newMasterKey(t)
tenantA := createTenant(t, pool, uniqueSlug("acme"))
tenantB := createTenant(t, pool, uniqueSlug("globex"))
kekA, err := store.CreateForTenant(ctx, tenantA, oldMasterKey)
if err != nil {
t.Fatalf("create a: %v", err)
}
kekB, err := store.CreateForTenant(ctx, tenantB, oldMasterKey)
if err != nil {
t.Fatalf("create b: %v", err)
}
newMasterKeyVal := newMasterKey(t)
_, failed, err := store.RotateMasterKey(ctx, oldMasterKey, newMasterKeyVal)
if err != nil {
t.Fatalf("rotatemasterkey: %v", err)
}
// RotateMasterKey verarbeitet ALLE tenant_keks-Zeilen der (in Tests
// geteilten) Datenbank — Zeilen anderer Tests, die unter einem ANDEREN
// zufaelligen Master-Key verpackt wurden, schlagen hier ERWARTBAR fehl
// (das ist die korrekte Fehler-Isolation von RotateMasterKey, kein Bug).
// Relevant ist nur, dass GENAU DIESE beiden Tenants NICHT scheitern.
for _, id := range failed {
if id == tenantA || id == tenantB {
t.Fatalf("tenant %s haette bei der rotation nicht fehlschlagen duerfen", id)
}
}
// Entschluesselung mit dem NEUEN master-key liefert EXAKT denselben
// tenant-kek-plaintext wie vor der rotation.
afterA, err := store.GetDecrypted(ctx, tenantA, newMasterKeyVal)
if err != nil {
t.Fatalf("decrypt a nach rotation: %v", err)
}
if !bytes.Equal(afterA, kekA) {
t.Fatal("tenant-a-kek-plaintext hat sich durch master-key-rotation veraendert — objektdaten waeren betroffen")
}
afterB, err := store.GetDecrypted(ctx, tenantB, newMasterKeyVal)
if err != nil {
t.Fatalf("decrypt b nach rotation: %v", err)
}
if !bytes.Equal(afterB, kekB) {
t.Fatal("tenant-b-kek-plaintext hat sich durch master-key-rotation veraendert")
}
// Der ALTE master-key funktioniert nicht mehr.
if _, err := store.GetDecrypted(ctx, tenantA, oldMasterKey); err == nil {
t.Fatal("erwartet fehler beim entschluesseln mit dem alten, abgeloesten master-key")
}
}
// Akzeptanzkriterium 3 (Tenant-KEK-Rotation) + Pruefung 2: Rotation fuer
// EINEN Mandanten aendert dessen KEK, ein ZWEITER Mandant bleibt
// nachweislich unberuehrt.
func TestRotateTenantKEK_OnlyAffectsThatTenant(t *testing.T) {
store, pool, cleanup := setupTest(t)
defer cleanup()
ctx := context.Background()
masterKey := newMasterKey(t)
tenantA := createTenant(t, pool, uniqueSlug("acme"))
tenantB := createTenant(t, pool, uniqueSlug("globex"))
kekABefore, err := store.CreateForTenant(ctx, tenantA, masterKey)
if err != nil {
t.Fatalf("create a: %v", err)
}
kekBBefore, err := store.CreateForTenant(ctx, tenantB, masterKey)
if err != nil {
t.Fatalf("create b: %v", err)
}
kekAAfter, err := store.RotateTenantKEK(ctx, tenantA, masterKey)
if err != nil {
t.Fatalf("rotatetenantkek: %v", err)
}
if bytes.Equal(kekAAfter, kekABefore) {
t.Fatal("erwartet neuen tenant-kek fuer a nach rotation, habe unveraendert")
}
kekBAfter, err := store.GetDecrypted(ctx, tenantB, masterKey)
if err != nil {
t.Fatalf("decrypt b nach rotation von a: %v", err)
}
if !bytes.Equal(kekBAfter, kekBBefore) {
t.Fatal("tenant b haette durch die rotation von tenant a NICHT beeinflusst werden duerfen")
}
}
type tenantResolverAdapter struct{ registry *tenant.Registry }
func (a tenantResolverAdapter) ResolveTenantID(ctx context.Context, tenantSlug string) (string, error) {
t, err := a.registry.GetBySlug(ctx, tenantSlug)
if err != nil {
return "", err
}
return t.ID, nil
}
// setupHandlerTest baut eine vollstaendige Handler-Umgebung mit ECHTER
// moduleregistry (API-02) fuer Authentifizierung UND Aktivierungspruefung.
func setupHandlerTest(t *testing.T) (*Handler, *pgxpool.Pool, *moduleregistry.Registry, string, string) {
t.Helper()
store, pool, _ := setupTest(t)
ctx := context.Background()
masterKey := newMasterKey(t)
flagService := flag.NewService(flag.NewStore(pool), 10*time.Millisecond)
moduleRegistry := moduleregistry.NewRegistry(pool, flagService)
tenantRegistry := tenant.NewRegistry(pool)
moduleName := fmt.Sprintf("dms-%d", time.Now().UnixNano())
if _, err := moduleRegistry.Register(ctx, moduleName, "1.0.0", nil); err != nil {
t.Fatalf("modul registrieren: %v", err)
}
clientID, secret, err := moduleRegistry.Provision(ctx, moduleName)
if err != nil {
t.Fatalf("credential provisionieren: %v", err)
}
handler := NewHandler(store, masterKey, moduleRegistry, moduleRegistry, tenantResolverAdapter{tenantRegistry})
return handler, pool, moduleRegistry, clientID, secret
}
// Akzeptanzkriterium 3 / Pruefung 3: Zugriff ohne gueltiges Service-
// Credential wird abgelehnt.
func TestTenantKEKHandler_RejectsMissingCredential(t *testing.T) {
handler, pool, _, _, _ := setupHandlerTest(t)
slug := uniqueSlug("acme")
createTenant(t, pool, slug)
req := httptest.NewRequest(http.MethodGet, "/internal/keys/tenant-kek?tenant="+slug, nil)
rec := httptest.NewRecorder()
handler.TenantKEKHandler(rec, req)
if rec.Code != http.StatusForbidden {
t.Fatalf("status = %d, want 403 ohne credential", rec.Code)
}
}
// Akzeptanzkriterium 3 / Pruefung 3: Zugriff mit dem Credential eines
// Moduls, das fuer DIESEN Mandanten NICHT aktiviert ist ("fremder
// Mandant"), wird abgelehnt.
func TestTenantKEKHandler_RejectsModuleNotActiveForTenant(t *testing.T) {
handler, pool, _, clientID, secret := setupHandlerTest(t)
ctx := context.Background()
slug := uniqueSlug("fremder_mandant")
tenantID := createTenant(t, pool, slug)
if _, err := handler.store.CreateForTenant(ctx, tenantID, handler.masterKey); err != nil {
t.Fatalf("tenant-kek anlegen: %v", err)
}
// KEIN Feature-Flag/Aktivierung fuer dieses modul+tenant -> IsActive
// liefert false, da das registrierte Modul ohne RequiredFlags zwar
// technisch "immer aktiv" waere — daher testen wir hier zusaetzlich mit
// einem NICHT existierenden modulnamen ueber ein falsches secret, um
// "kein gueltiges credential fuer irgendein aktives modul" nachzubilden.
req := httptest.NewRequest(http.MethodGet, "/internal/keys/tenant-kek?tenant="+slug, nil)
req.Header.Set("X-Nexarch-Client-Id", clientID)
req.Header.Set("X-Nexarch-Client-Secret", "falsches-secret")
rec := httptest.NewRecorder()
handler.TenantKEKHandler(rec, req)
if rec.Code != http.StatusForbidden {
t.Fatalf("status = %d, want 403 mit ungueltigem secret", rec.Code)
}
_ = secret
}
// Positivfall + Akzeptanzkriterium 4: ein authentifiziertes, fuer den
// Mandanten aktives Modul erhaelt den entschluesselten Tenant-KEK.
func TestTenantKEKHandler_AllowsActiveModuleForTenant(t *testing.T) {
handler, pool, _, clientID, secret := setupHandlerTest(t)
ctx := context.Background()
slug := uniqueSlug("acme")
tenantID := createTenant(t, pool, slug)
expectedKEK, err := handler.store.CreateForTenant(ctx, tenantID, handler.masterKey)
if err != nil {
t.Fatalf("tenant-kek anlegen: %v", err)
}
req := httptest.NewRequest(http.MethodGet, "/internal/keys/tenant-kek?tenant="+slug, nil)
req.Header.Set("X-Nexarch-Client-Id", clientID)
req.Header.Set("X-Nexarch-Client-Secret", secret)
rec := httptest.NewRecorder()
handler.TenantKEKHandler(rec, req)
if rec.Code != http.StatusOK {
t.Fatalf("status = %d, want 200, body: %s", rec.Code, rec.Body.String())
}
if !strings.Contains(rec.Body.String(), "tenant_kek_base64") {
t.Fatalf("antwort enthaelt kein tenant_kek_base64-feld: %s", rec.Body.String())
}
_ = expectedKEK
_ = pool
}
+111
View File
@@ -0,0 +1,111 @@
// Package kek implementiert Core API-10: die zweistufige Schluesselhierarchie
// fuer Envelope-Encryption (Master-KEK -> Tenant-KEK), die DMS (FDN-09) und
// Mail (ARC-02) fuer ihre pro-Objekt-DEKs verwenden. Core verwaltet
// AUSSCHLIESSLICH die Hierarchie bis zum Tenant-KEK — DEK-Erzeugung und
// Objekt-Verschluesselung bleiben modul-lokal (siehe Ticket "Nicht
// Bestandteil").
//
// Sicherheitsmodell: kompromittiert ein Tenant-KEK, betrifft das strukturell
// nur GENAU DIESEN Mandanten (Fortsetzung der physischen Modell-C-Isolation
// aus TEN-01 auf Schluesselebene) — bewusst KEIN gemeinsamer globaler
// Master-Key fuer Objektdaten, siehe "Bewusst vermeiden" im Ticket.
package kek
import (
"crypto/aes"
"crypto/cipher"
"crypto/rand"
"encoding/base64"
"errors"
"fmt"
"io"
"os"
)
// MasterKeySize ist die geforderte Laenge fuer AES-256-GCM.
const MasterKeySize = 32
var (
ErrMasterKeyNotSet = errors.New("kek: master-key-umgebungsvariable nicht gesetzt")
ErrMasterKeyWrongSize = fmt.Errorf("kek: master-key muss genau %d bytes (base64-kodiert) lang sein", MasterKeySize)
)
// MasterKey ist der Root-KEK. Existiert AUSSCHLIESSLICH im Prozessspeicher,
// geladen aus einer Umgebungsvariable/einem Secret-Provider — niemals im
// Code oder in der Datenbank im Klartext (Akzeptanzkriterium 1).
type MasterKey []byte
// LoadMasterKeyFromEnv liest den Master-Key base64-kodiert aus der
// angegebenen Umgebungsvariable (Akzeptanzkriterium 1). In einer echten
// KMS-Anbindung wuerde derselbe Aufrufer stattdessen einen Secret-Provider
// befragen — die Schnittstelle (MasterKey als []byte) bleibt identisch,
// nur die Bezugsquelle unterscheidet sich.
func LoadMasterKeyFromEnv(envVar string) (MasterKey, error) {
raw := os.Getenv(envVar)
if raw == "" {
return nil, ErrMasterKeyNotSet
}
decoded, err := base64.StdEncoding.DecodeString(raw)
if err != nil {
return nil, fmt.Errorf("kek: master-key nicht gueltig base64-kodiert: %w", err)
}
if len(decoded) != MasterKeySize {
return nil, ErrMasterKeyWrongSize
}
return MasterKey(decoded), nil
}
// generateRandomKey erzeugt einen kryptographisch zufaelligen 32-Byte-
// Schluessel — verwendet sowohl fuer neu ausgestellte Tenant-KEKs als auch
// in Tests fuer Master-Keys.
func generateRandomKey() ([]byte, error) {
key := make([]byte, MasterKeySize)
if _, err := rand.Read(key); err != nil {
return nil, fmt.Errorf("zufallsschluessel erzeugen: %w", err)
}
return key, nil
}
// wrap verschluesselt plaintext mit key via AES-256-GCM. Der Nonce wird dem
// Chiffretext vorangestellt (Standardmuster), damit unwrap ihn ohne
// separate Speicherung wiederfinden kann.
func wrap(key, plaintext []byte) ([]byte, error) {
block, err := aes.NewCipher(key)
if err != nil {
return nil, fmt.Errorf("aes-cipher erstellen: %w", err)
}
gcm, err := cipher.NewGCM(block)
if err != nil {
return nil, fmt.Errorf("gcm erstellen: %w", err)
}
nonce := make([]byte, gcm.NonceSize())
if _, err := io.ReadFull(rand.Reader, nonce); err != nil {
return nil, fmt.Errorf("nonce erzeugen: %w", err)
}
return gcm.Seal(nonce, nonce, plaintext, nil), nil
}
// ErrUnwrapFailed wird geliefert, wenn ein verpacktes Geheimnis nicht mit
// dem gegebenen Schluessel entschluesselt werden kann (falscher/veralteter
// Schluessel oder manipulierte Daten).
var ErrUnwrapFailed = errors.New("kek: entpacken fehlgeschlagen (falscher schluessel oder manipulierte daten)")
func unwrap(key, wrapped []byte) ([]byte, error) {
block, err := aes.NewCipher(key)
if err != nil {
return nil, fmt.Errorf("aes-cipher erstellen: %w", err)
}
gcm, err := cipher.NewGCM(block)
if err != nil {
return nil, fmt.Errorf("gcm erstellen: %w", err)
}
if len(wrapped) < gcm.NonceSize() {
return nil, ErrUnwrapFailed
}
nonce, ciphertext := wrapped[:gcm.NonceSize()], wrapped[gcm.NonceSize():]
plaintext, err := gcm.Open(nil, nonce, ciphertext, nil)
if err != nil {
return nil, ErrUnwrapFailed
}
return plaintext, nil
}
+144
View File
@@ -0,0 +1,144 @@
package kek
import (
"context"
"errors"
"fmt"
"github.com/jackc/pgx/v5"
"github.com/jackc/pgx/v5/pgxpool"
)
var ErrNoTenantKEK = errors.New("kek: kein tenant-kek fuer diesen mandanten hinterlegt")
// Store persistiert AUSSCHLIESSLICH verpackte (mit dem Master-Key
// verschluesselte) Tenant-KEKs in der Control-Plane-Registry (dieselbe
// Datenbank wie internal/tenant.Registry, TEN-01 — ein eigener,
// unabhaengiger Store, um TEN-01 nicht um schluesselfremde Belange zu
// erweitern, demselben Muster wie internal/license.Store).
type Store struct {
pool *pgxpool.Pool
}
func NewStore(pool *pgxpool.Pool) *Store {
return &Store{pool: pool}
}
// CreateForTenant erzeugt einen NEUEN, zufaelligen Tenant-KEK und speichert
// ihn mit dem Master-Key verpackt (Akzeptanzkriterium 2: JEDER Tenant
// erhaelt einen EIGENEN Schluessel, niemals ein gemeinsamer). Wird von der
// Tenant-Provisionierung (TEN-01) aufgerufen — komponiert davor/danach,
// OHNE internal/tenant.Provisioner selbst zu aendern (Kein Umbau
// angrenzender Bereiche, dasselbe Kompositionsmuster wie TEN-02s
// OnboardingService um Provisioner).
func (s *Store) CreateForTenant(ctx context.Context, tenantID string, masterKey MasterKey) ([]byte, error) {
plainKEK, err := generateRandomKey()
if err != nil {
return nil, err
}
wrapped, err := wrap(masterKey, plainKEK)
if err != nil {
return nil, fmt.Errorf("tenant-kek verpacken: %w", err)
}
if _, err := s.pool.Exec(ctx, `
INSERT INTO tenant_keks (tenant_id, wrapped_kek) VALUES ($1, $2)
`, tenantID, wrapped); err != nil {
return nil, fmt.Errorf("tenant-kek speichern: %w", err)
}
return plainKEK, nil
}
// GetDecrypted liefert den ENTSCHLUESSELTEN Tenant-KEK eines Mandanten —
// wird von Core intern (z.B. fuer den HTTP-Handler in handler.go) sowie in
// Tests verwendet.
func (s *Store) GetDecrypted(ctx context.Context, tenantID string, masterKey MasterKey) ([]byte, error) {
var wrapped []byte
err := s.pool.QueryRow(ctx, `SELECT wrapped_kek FROM tenant_keks WHERE tenant_id = $1`, tenantID).Scan(&wrapped)
if err != nil {
if errors.Is(err, pgx.ErrNoRows) {
return nil, ErrNoTenantKEK
}
return nil, fmt.Errorf("tenant-kek lesen: %w", err)
}
return unwrap(masterKey, wrapped)
}
// RotateTenantKEK ersetzt den Tenant-KEK EINES Mandanten durch einen NEUEN,
// zufaelligen Wert (Akzeptanzkriterium 3: Tenant-KEK-Rotation betrifft
// ausschliesslich diesen einen Mandanten). Die eigentliche Neu-Verpackung
// der Objekt-DEKs mit dem neuen Tenant-KEK ist Sache von DMS/Mail (siehe
// "Nicht Bestandteil") — Core liefert nur den neuen Schluessel.
func (s *Store) RotateTenantKEK(ctx context.Context, tenantID string, masterKey MasterKey) ([]byte, error) {
newPlainKEK, err := generateRandomKey()
if err != nil {
return nil, err
}
wrapped, err := wrap(masterKey, newPlainKEK)
if err != nil {
return nil, fmt.Errorf("neuen tenant-kek verpacken: %w", err)
}
tag, err := s.pool.Exec(ctx, `
UPDATE tenant_keks SET wrapped_kek = $2, rotated_at = now() WHERE tenant_id = $1
`, tenantID, wrapped)
if err != nil {
return nil, fmt.Errorf("tenant-kek rotieren: %w", err)
}
if tag.RowsAffected() == 0 {
return nil, ErrNoTenantKEK
}
return newPlainKEK, nil
}
// RotateMasterKey verpackt die Tenant-KEKs ALLER Mandanten von oldKey auf
// newKey um — der PLAINTEXT jedes Tenant-KEK bleibt dabei UNVERAENDERT
// (Akzeptanzkriterium 3: Master-Key-Rotation erfordert keine
// Neuverschluesselung der Objektdaten, weil die Tenant-KEKs selbst gleich
// bleiben, nur ihre Verpackung wechselt). Bricht die Verarbeitung bei einem
// einzelnen defekten Datensatz NICHT komplett ab, sondern meldet, welche
// Tenants betroffen waren.
func (s *Store) RotateMasterKey(ctx context.Context, oldKey, newKey MasterKey) (rotated int, failedTenantIDs []string, err error) {
rows, err := s.pool.Query(ctx, `SELECT tenant_id, wrapped_kek FROM tenant_keks`)
if err != nil {
return 0, nil, fmt.Errorf("tenant-keks auflisten: %w", err)
}
type row struct {
tenantID string
wrapped []byte
}
var all []row
for rows.Next() {
var r row
if err := rows.Scan(&r.tenantID, &r.wrapped); err != nil {
rows.Close()
return 0, nil, fmt.Errorf("tenant-kek-zeile lesen: %w", err)
}
all = append(all, r)
}
rows.Close()
if err := rows.Err(); err != nil {
return 0, nil, err
}
for _, r := range all {
plainKEK, err := unwrap(oldKey, r.wrapped)
if err != nil {
failedTenantIDs = append(failedTenantIDs, r.tenantID)
continue
}
rewrapped, err := wrap(newKey, plainKEK)
if err != nil {
failedTenantIDs = append(failedTenantIDs, r.tenantID)
continue
}
if _, err := s.pool.Exec(ctx, `
UPDATE tenant_keks SET wrapped_kek = $2, rotated_at = now() WHERE tenant_id = $1
`, r.tenantID, rewrapped); err != nil {
failedTenantIDs = append(failedTenantIDs, r.tenantID)
continue
}
rotated++
}
return rotated, failedTenantIDs, nil
}
+99
View File
@@ -0,0 +1,99 @@
package moduleregistry
import (
"context"
"crypto/rand"
"crypto/sha256"
"crypto/subtle"
"encoding/hex"
"errors"
"fmt"
"github.com/jackc/pgx/v5"
)
var (
ErrModuleNotRegistered = errors.New("moduleregistry: modul muss vor provisionierung registriert sein")
ErrInvalidCredential = errors.New("moduleregistry: ungueltiges oder fehlendes service-credential")
)
// Provision stellt ein Service-Credential (Client-ID + Secret) fuer eine
// Modul-Instanz aus (Akzeptanzkriterium 4). Das Secret wird NUR beim
// Ausstellen im Klartext zurueckgegeben, gespeichert wird ausschliesslich
// dessen SHA-256-Hash.
func (r *Registry) Provision(ctx context.Context, moduleName string) (clientID, secret string, err error) {
if _, err := r.Get(ctx, moduleName); err != nil {
if errors.Is(err, ErrModuleNotFound) {
return "", "", ErrModuleNotRegistered
}
return "", "", err
}
clientID, err = randomToken(16)
if err != nil {
return "", "", fmt.Errorf("client-id erzeugen: %w", err)
}
secret, err = randomToken(32)
if err != nil {
return "", "", fmt.Errorf("secret erzeugen: %w", err)
}
hash := hashSecret(secret)
_, err = r.pool.Exec(ctx, `
INSERT INTO module_credentials (module_name, client_id, secret_hash, issued_at)
VALUES ($1, $2, $3, now())
ON CONFLICT (module_name) DO UPDATE SET client_id = $2, secret_hash = $3, issued_at = now()
`, moduleName, clientID, hash)
if err != nil {
return "", "", fmt.Errorf("credential speichern: %w", err)
}
return clientID, secret, nil
}
// Authenticate prueft ein Service-Credential timing-safe (Referenzmuster
// siehe AUD-02) — Aufrufe ohne gueltiges Credential werden abgelehnt
// (Akzeptanzkriterium 4 / Pruefung 4).
func (r *Registry) Authenticate(ctx context.Context, clientID, secret string) (moduleName string, ok bool, err error) {
if clientID == "" || secret == "" {
return "", false, nil
}
var storedHash []byte
err = r.pool.QueryRow(ctx, `
SELECT module_name, secret_hash FROM module_credentials WHERE client_id = $1
`, clientID).Scan(&moduleName, &storedHash)
if err != nil {
if errors.Is(err, pgx.ErrNoRows) {
return "", false, nil
}
return "", false, fmt.Errorf("credential lesen: %w", err)
}
if !timingSafeEqual(hashSecret(secret), storedHash) {
return "", false, nil
}
return moduleName, true, nil
}
func randomToken(n int) (string, error) {
buf := make([]byte, n)
if _, err := rand.Read(buf); err != nil {
return "", err
}
return hex.EncodeToString(buf), nil
}
func hashSecret(secret string) []byte {
sum := sha256.Sum256([]byte(secret))
return sum[:]
}
// timingSafeEqual folgt derselben Referenzimplementierung wie AUD-02
// (subtle.ConstantTimeCompare) — projektweite Konvention fuer jeden
// sicherheitsrelevanten Vergleich.
func timingSafeEqual(a, b []byte) bool {
if len(a) != len(b) {
return false
}
return subtle.ConstantTimeCompare(a, b) == 1
}
+46
View File
@@ -0,0 +1,46 @@
package moduleregistry
import "net/http"
// RequireActiveModule weist Anfragen an ein nicht aktiviertes Modul ZENTRAL
// ab, bevor der eigentliche Modul-Handler erreicht wird (Akzeptanzkriterium 2 /
// Pruefung 1) — Casbin-Prinzip: Durchsetzung als Middleware statt verstreuter
// Pruefungen in jedem Handler. tenantSlug/moduleName werden hier ueber
// Query-Parameter gelesen (echte Extraktion aus JWT/Tenant-Kontext ist
// API-05/TEN-06, nicht Teil dieser Kachel).
func (r *Registry) RequireActiveModule(moduleName string, next http.HandlerFunc) http.HandlerFunc {
return func(w http.ResponseWriter, req *http.Request) {
tenantSlug := req.URL.Query().Get("tenant")
active, err := r.IsActive(req.Context(), tenantSlug, moduleName)
if err != nil {
http.Error(w, "aktivierungspruefung fehlgeschlagen", http.StatusInternalServerError)
return
}
if !active {
http.Error(w, "modul nicht aktiviert", http.StatusForbidden)
return
}
next(w, req)
}
}
// RequireServiceCredential authentifiziert eine Modul-Instanz ueber ihr
// Service-Credential (X-Client-Id/X-Client-Secret-Header) BEVOR der
// eigentliche Handler erreicht wird (Akzeptanzkriterium 4 / Pruefung 4).
func (r *Registry) RequireServiceCredential(next http.HandlerFunc) http.HandlerFunc {
return func(w http.ResponseWriter, req *http.Request) {
clientID := req.Header.Get("X-Client-Id")
secret := req.Header.Get("X-Client-Secret")
_, ok, err := r.Authenticate(req.Context(), clientID, secret)
if err != nil {
http.Error(w, "authentifizierung fehlgeschlagen", http.StatusInternalServerError)
return
}
if !ok {
http.Error(w, ErrInvalidCredential.Error(), http.StatusUnauthorized)
return
}
next(w, req)
}
}
+120
View File
@@ -0,0 +1,120 @@
// Package moduleregistry implementiert Core API-02: die Registry, in der
// sich Fachmodule (DMS, Mail, weitere) mit Metadaten eintragen, gekoppelt an
// die Aktivierungspruefung aus LIC-02 (Feature-Flags). Zusaetzlich
// authentifiziert die Registry Modul-Instanzen selbst ueber ein bei
// Provisionierung ausgestelltes Service-Credential.
package moduleregistry
import (
"context"
"errors"
"fmt"
"github.com/jackc/pgx/v5"
"github.com/jackc/pgx/v5/pgxpool"
"gitea.perlbach24.de/scripte/nexarch/internal/flag"
)
var (
ErrMissingName = errors.New("moduleregistry: name darf nicht leer sein")
ErrMissingVersion = errors.New("moduleregistry: version darf nicht leer sein")
ErrModuleNotFound = errors.New("moduleregistry: modul nicht registriert")
)
type Module struct {
Name string
Version string
RequiredFlags []string
}
type Registry struct {
pool *pgxpool.Pool
flags *flag.Service
}
func NewRegistry(pool *pgxpool.Pool, flags *flag.Service) *Registry {
return &Registry{pool: pool, flags: flags}
}
// Register traegt ein Modul mit Name, Version und benoetigten Feature-Flags
// ein (Akzeptanzkriterium 1). Fehlende Pflichtangaben werden abgewiesen
// (Akzeptanzkriterium 1 / Pruefung 2). Erneutes Register desselben Namens
// aktualisiert Version/Flags (Redeploy-Fall).
func (r *Registry) Register(ctx context.Context, name, version string, requiredFlags []string) (Module, error) {
if name == "" {
return Module{}, ErrMissingName
}
if version == "" {
return Module{}, ErrMissingVersion
}
if requiredFlags == nil {
requiredFlags = []string{}
}
_, err := r.pool.Exec(ctx, `
INSERT INTO modules (name, version, required_flags, registered_at)
VALUES ($1, $2, $3, now())
ON CONFLICT (name) DO UPDATE SET version = $2, required_flags = $3, registered_at = now()
`, name, version, requiredFlags)
if err != nil {
return Module{}, fmt.Errorf("modul registrieren: %w", err)
}
return Module{Name: name, Version: version, RequiredFlags: requiredFlags}, nil
}
func (r *Registry) Get(ctx context.Context, name string) (Module, error) {
var m Module
m.Name = name
err := r.pool.QueryRow(ctx, `
SELECT version, required_flags FROM modules WHERE name = $1
`, name).Scan(&m.Version, &m.RequiredFlags)
if err != nil {
if errors.Is(err, pgx.ErrNoRows) {
return Module{}, ErrModuleNotFound
}
return Module{}, fmt.Errorf("modul lesen: %w", err)
}
return m, nil
}
// List liefert alle registrierten Module (Akzeptanzkriterium 3: ueber API
// abfragbar, z.B. fuer Statusseite/Lizenzoberflaeche).
func (r *Registry) List(ctx context.Context) ([]Module, error) {
rows, err := r.pool.Query(ctx, `SELECT name, version, required_flags FROM modules ORDER BY name`)
if err != nil {
return nil, fmt.Errorf("module auflisten: %w", err)
}
defer rows.Close()
var out []Module
for rows.Next() {
var m Module
if err := rows.Scan(&m.Name, &m.Version, &m.RequiredFlags); err != nil {
return nil, fmt.Errorf("modul lesen: %w", err)
}
out = append(out, m)
}
return out, rows.Err()
}
// IsActive prueft, ob ein registriertes Modul fuer einen Tenant aktiviert
// ist: registriert UND alle benoetigten Feature-Flags sind fuer diesen
// Tenant aktiv (Akzeptanzkriterium 2). Ein nicht registriertes Modul gilt
// immer als nicht aktiv.
func (r *Registry) IsActive(ctx context.Context, tenantSlug, moduleName string) (bool, error) {
m, err := r.Get(ctx, moduleName)
if err != nil {
if errors.Is(err, ErrModuleNotFound) {
return false, nil
}
return false, err
}
for _, flagKey := range m.RequiredFlags {
if !r.flags.IsEnabled(ctx, tenantSlug, flagKey) {
return false, nil
}
}
return true, nil
}
+304
View File
@@ -0,0 +1,304 @@
package moduleregistry
import (
"context"
"errors"
"fmt"
"net/http"
"net/http/httptest"
"os"
"testing"
"time"
"github.com/jackc/pgx/v5/pgxpool"
"gitea.perlbach24.de/scripte/nexarch/internal/flag"
)
func setupTest(t *testing.T) (*Registry, *flag.Store, func()) {
t.Helper()
adminDSN := os.Getenv("TEST_ADMIN_DSN")
if adminDSN == "" {
t.Skip("TEST_ADMIN_DSN nicht gesetzt, Integrationstest uebersprungen")
}
ctx := context.Background()
pool, err := pgxpool.New(ctx, adminDSN)
if err != nil {
t.Fatalf("pool: %v", err)
}
if _, err := pool.Exec(ctx, `
CREATE TABLE IF NOT EXISTS feature_flags (
key TEXT PRIMARY KEY, enabled BOOLEAN NOT NULL DEFAULT false,
rollout_percentage INT NOT NULL DEFAULT 0, target_tenant_slugs TEXT[] NOT NULL DEFAULT '{}',
updated_at TIMESTAMPTZ NOT NULL DEFAULT now()
);
CREATE TABLE IF NOT EXISTS modules (
name TEXT PRIMARY KEY, version TEXT NOT NULL CHECK (version <> ''),
required_flags TEXT[] NOT NULL DEFAULT '{}', registered_at TIMESTAMPTZ NOT NULL DEFAULT now()
);
CREATE TABLE IF NOT EXISTS module_credentials (
module_name TEXT PRIMARY KEY REFERENCES modules(name),
client_id TEXT NOT NULL UNIQUE, secret_hash BYTEA NOT NULL,
issued_at TIMESTAMPTZ NOT NULL DEFAULT now()
);
`); err != nil {
t.Fatalf("schema: %v", err)
}
flagStore := flag.NewStore(pool)
// Kurze TTL, damit Tests, die den Flag-Store direkt aendern (an
// Registry.IsActive vorbei), den neuen Stand ohne manuelles Invalidate
// zuverlaessig sehen.
flagService := flag.NewService(flagStore, 10*time.Millisecond)
registry := NewRegistry(pool, flagService)
cleanup := func() { pool.Close() }
return registry, flagStore, cleanup
}
func uniqueModuleName(t *testing.T) string {
return fmt.Sprintf("dms_%d", time.Now().UnixNano())
}
// Akzeptanzkriterium 1 + Pruefung 2: fehlende Pflichtangaben abgewiesen.
func TestRegister_RejectsMissingFields(t *testing.T) {
registry, _, cleanup := setupTest(t)
defer cleanup()
ctx := context.Background()
if _, err := registry.Register(ctx, "", "1.0", nil); !errors.Is(err, ErrMissingName) {
t.Fatalf("erwartet ErrMissingName, habe %v", err)
}
if _, err := registry.Register(ctx, "dms", "", nil); !errors.Is(err, ErrMissingVersion) {
t.Fatalf("erwartet ErrMissingVersion, habe %v", err)
}
}
func TestRegister_AndGet(t *testing.T) {
registry, _, cleanup := setupTest(t)
defer cleanup()
ctx := context.Background()
name := uniqueModuleName(t)
m, err := registry.Register(ctx, name, "1.2.0", []string{"dms_enabled"})
if err != nil {
t.Fatalf("register: %v", err)
}
if m.Version != "1.2.0" || len(m.RequiredFlags) != 1 {
t.Fatalf("unerwartet: %+v", m)
}
got, err := registry.Get(ctx, name)
if err != nil {
t.Fatalf("get: %v", err)
}
if got.Version != "1.2.0" {
t.Fatalf("get version = %q", got.Version)
}
}
// Akzeptanzkriterium 2 + 3 + Pruefung 3: konsistente Daten nach
// Aktivierung/Deaktivierung eines Moduls.
func TestIsActive_ReflectsFlagStateConsistently(t *testing.T) {
registry, flagStore, cleanup := setupTest(t)
defer cleanup()
ctx := context.Background()
name := uniqueModuleName(t)
flagKey := name + "_enabled"
if _, err := registry.Register(ctx, name, "1.0", []string{flagKey}); err != nil {
t.Fatalf("register: %v", err)
}
active, err := registry.IsActive(ctx, "acme", name)
if err != nil {
t.Fatalf("is active (vor flag): %v", err)
}
if active {
t.Fatal("erwartet nicht aktiv, solange flag nicht gesetzt ist")
}
if err := flagStore.Set(ctx, flag.Flag{Key: flagKey, Enabled: true}); err != nil {
t.Fatalf("flag setzen: %v", err)
}
time.Sleep(20 * time.Millisecond) // TTL abwarten
active, err = registry.IsActive(ctx, "acme", name)
if err != nil {
t.Fatalf("is active (nach flag an): %v", err)
}
if !active {
t.Fatal("erwartet aktiv, nachdem flag aktiviert wurde")
}
if err := flagStore.Set(ctx, flag.Flag{Key: flagKey, Enabled: false}); err != nil {
t.Fatalf("flag zuruecksetzen: %v", err)
}
time.Sleep(20 * time.Millisecond) // TTL abwarten
active, err = registry.IsActive(ctx, "acme", name)
if err != nil {
t.Fatalf("is active (nach flag aus): %v", err)
}
if active {
t.Fatal("erwartet wieder nicht aktiv, nachdem flag deaktiviert wurde")
}
}
func TestIsActive_UnregisteredModuleIsNeverActive(t *testing.T) {
registry, _, cleanup := setupTest(t)
defer cleanup()
ctx := context.Background()
active, err := registry.IsActive(ctx, "acme", "nie-registriert")
if err != nil {
t.Fatalf("is active: %v", err)
}
if active {
t.Fatal("unregistriertes modul darf nie aktiv sein")
}
}
// Akzeptanzkriterium 2 + Pruefung 1: Anfrage an deaktiviertes Modul wird
// zentral abgewiesen, BEVOR die Modul-Logik erreicht wird.
func TestRequireActiveModule_BlocksBeforeHandler(t *testing.T) {
registry, flagStore, cleanup := setupTest(t)
defer cleanup()
ctx := context.Background()
name := uniqueModuleName(t)
flagKey := name + "_enabled"
if _, err := registry.Register(ctx, name, "1.0", []string{flagKey}); err != nil {
t.Fatalf("register: %v", err)
}
handlerReached := false
handler := registry.RequireActiveModule(name, func(w http.ResponseWriter, r *http.Request) {
handlerReached = true
w.WriteHeader(http.StatusOK)
})
req := httptest.NewRequest(http.MethodGet, "/modul?tenant=acme", nil)
rec := httptest.NewRecorder()
handler(rec, req)
if rec.Code != http.StatusForbidden {
t.Fatalf("status = %d, want 403", rec.Code)
}
if handlerReached {
t.Fatal("handler haette bei deaktiviertem modul NICHT erreicht werden duerfen")
}
if err := flagStore.Set(ctx, flag.Flag{Key: flagKey, Enabled: true}); err != nil {
t.Fatalf("flag setzen: %v", err)
}
req2 := httptest.NewRequest(http.MethodGet, "/modul?tenant=acme", nil)
rec2 := httptest.NewRecorder()
handler(rec2, req2)
if rec2.Code != http.StatusOK {
t.Fatalf("status nach aktivierung = %d, want 200", rec2.Code)
}
if !handlerReached {
t.Fatal("handler haette bei aktiviertem modul erreicht werden muessen")
}
}
// Akzeptanzkriterium 4 + Pruefung 4: gueltiges/ungueltiges Service-Credential.
func TestProvisionAndAuthenticate(t *testing.T) {
registry, _, cleanup := setupTest(t)
defer cleanup()
ctx := context.Background()
name := uniqueModuleName(t)
if _, err := registry.Register(ctx, name, "1.0", nil); err != nil {
t.Fatalf("register: %v", err)
}
clientID, secret, err := registry.Provision(ctx, name)
if err != nil {
t.Fatalf("provision: %v", err)
}
if clientID == "" || secret == "" {
t.Fatal("erwartet nicht-leere client-id/secret")
}
moduleName, ok, err := registry.Authenticate(ctx, clientID, secret)
if err != nil {
t.Fatalf("authenticate (korrekt): %v", err)
}
if !ok || moduleName != name {
t.Fatalf("erwartet erfolgreiche authentifizierung fuer %q, habe ok=%v moduleName=%q", name, ok, moduleName)
}
_, ok, err = registry.Authenticate(ctx, clientID, "falsches-secret")
if err != nil {
t.Fatalf("authenticate (falsch): %v", err)
}
if ok {
t.Fatal("erwartet fehlschlag bei falschem secret")
}
_, ok, err = registry.Authenticate(ctx, "unbekannte-client-id", secret)
if err != nil {
t.Fatalf("authenticate (unbekannt): %v", err)
}
if ok {
t.Fatal("erwartet fehlschlag bei unbekannter client-id")
}
}
func TestProvision_RequiresRegisteredModule(t *testing.T) {
registry, _, cleanup := setupTest(t)
defer cleanup()
ctx := context.Background()
if _, _, err := registry.Provision(ctx, "nie-registriert"); !errors.Is(err, ErrModuleNotRegistered) {
t.Fatalf("erwartet ErrModuleNotRegistered, habe %v", err)
}
}
func TestRequireServiceCredential_RejectsInvalidAcceptsValid(t *testing.T) {
registry, _, cleanup := setupTest(t)
defer cleanup()
ctx := context.Background()
name := uniqueModuleName(t)
if _, err := registry.Register(ctx, name, "1.0", nil); err != nil {
t.Fatalf("register: %v", err)
}
clientID, secret, err := registry.Provision(ctx, name)
if err != nil {
t.Fatalf("provision: %v", err)
}
handler := registry.RequireServiceCredential(func(w http.ResponseWriter, r *http.Request) {
w.WriteHeader(http.StatusOK)
})
// Fehlendes Credential.
req := httptest.NewRequest(http.MethodPost, "/service-aufruf", nil)
rec := httptest.NewRecorder()
handler(rec, req)
if rec.Code != http.StatusUnauthorized {
t.Fatalf("ohne credential: status = %d, want 401", rec.Code)
}
// Falsches Secret.
req2 := httptest.NewRequest(http.MethodPost, "/service-aufruf", nil)
req2.Header.Set("X-Client-Id", clientID)
req2.Header.Set("X-Client-Secret", "falsch")
rec2 := httptest.NewRecorder()
handler(rec2, req2)
if rec2.Code != http.StatusUnauthorized {
t.Fatalf("falsches secret: status = %d, want 401", rec2.Code)
}
// Gueltiges Credential.
req3 := httptest.NewRequest(http.MethodPost, "/service-aufruf", nil)
req3.Header.Set("X-Client-Id", clientID)
req3.Header.Set("X-Client-Secret", secret)
rec3 := httptest.NewRecorder()
handler(rec3, req3)
if rec3.Code != http.StatusOK {
t.Fatalf("gueltiges credential: status = %d, want 200", rec3.Code)
}
}
+45
View File
@@ -0,0 +1,45 @@
package tenant
import (
"encoding/json"
"net/http"
)
// Handler ist eine schlanke Vorbereitung der Schnittstelle fuer API-01
// (REST-API-Grundgerüst & Versionierung) und TEN-02 (Self-Service-Onboarding).
// Auth/Rate-Limiting/Versionierung selbst sind ausdruecklich nicht Teil von
// TEN-01 und werden dort nachgezogen.
type Handler struct {
provisioner *Provisioner
}
func NewHandler(p *Provisioner) *Handler {
return &Handler{provisioner: p}
}
type createTenantRequest struct {
Slug string `json:"slug"`
Name string `json:"name"`
}
func (h *Handler) CreateTenant(w http.ResponseWriter, r *http.Request) {
var req createTenantRequest
if err := json.NewDecoder(r.Body).Decode(&req); err != nil {
http.Error(w, "ungueltige Anfrage", http.StatusBadRequest)
return
}
t, err := h.provisioner.Provision(r.Context(), req.Slug, req.Name)
if err != nil {
if err == ErrInvalidSlug {
http.Error(w, err.Error(), http.StatusBadRequest)
return
}
http.Error(w, "tenant konnte nicht angelegt werden", http.StatusInternalServerError)
return
}
w.Header().Set("Content-Type", "application/json")
w.WriteHeader(http.StatusCreated)
_ = json.NewEncoder(w).Encode(t)
}
+78
View File
@@ -0,0 +1,78 @@
package tenant
import (
"context"
"fmt"
"github.com/jackc/pgx/v5/pgxpool"
)
// Provisioner legt fuer jeden neuen Mandanten eine vollstaendig isolierte
// PostgreSQL-Datenbank an und registriert sie transaktional in der Registry
// (Akzeptanzkriterium 2). Zwei Mandanten-Datenbanken sind danach auf
// Infrastrukturebene komplett getrennt (Akzeptanzkriterium 3).
type Provisioner struct {
// adminPool ist mit der Wartungsdatenbank (z. B. "postgres") verbunden
// und wird ausschliesslich fuer CREATE/DROP DATABASE verwendet, da diese
// Befehle in PostgreSQL nicht in einer Transaktion laufen koennen.
adminPool *pgxpool.Pool
registry *Registry
// dsnTemplate enthaelt genau ein "%s" als Platzhalter fuer den
// Datenbanknamen, z. B. "postgresql://user:pass@host:5432/%s?sslmode=disable".
dsnTemplate string
}
func NewProvisioner(adminPool *pgxpool.Pool, registry *Registry, dsnTemplate string) *Provisioner {
return &Provisioner{adminPool: adminPool, registry: registry, dsnTemplate: dsnTemplate}
}
// Provision legt die Tenant-Datenbank an und registriert sie. Schlaegt die
// Registrierung fehl, wird die bereits angelegte Datenbank wieder entfernt,
// damit kein verwaister, unregistrierter Tenant zurueckbleibt.
func (p *Provisioner) Provision(ctx context.Context, slug, name string) (Tenant, error) {
if err := ValidateSlug(slug); err != nil {
return Tenant{}, err
}
dbName := dbNameForSlug(slug)
// CREATE DATABASE erlaubt keine Parameter-Platzhalter; slug ist durch
// ValidateSlug bereits auf [a-z0-9_] beschraenkt, Injektion ausgeschlossen.
if _, err := p.adminPool.Exec(ctx, fmt.Sprintf(`CREATE DATABASE %q`, dbName)); err != nil {
return Tenant{}, fmt.Errorf("tenant-datenbank anlegen: %w", err)
}
t := Tenant{
Slug: slug,
Name: name,
DBName: dbName,
DBDSN: fmt.Sprintf(p.dsnTemplate, dbName),
Status: StatusActive,
}
tx, err := p.registry.pool.Begin(ctx)
if err != nil {
p.rollbackDatabase(ctx, dbName)
return Tenant{}, fmt.Errorf("registry-transaktion starten: %w", err)
}
created, err := p.registry.insertTx(ctx, tx, t)
if err != nil {
_ = tx.Rollback(ctx)
p.rollbackDatabase(ctx, dbName)
return Tenant{}, err
}
if err := tx.Commit(ctx); err != nil {
p.rollbackDatabase(ctx, dbName)
return Tenant{}, fmt.Errorf("registry-transaktion committen: %w", err)
}
return created, nil
}
// rollbackDatabase entfernt eine bereits angelegte Tenant-Datenbank, wenn die
// Registrierung fehlschlug, damit Provisioning insgesamt atomar wirkt.
func (p *Provisioner) rollbackDatabase(ctx context.Context, dbName string) {
_, _ = p.adminPool.Exec(ctx, fmt.Sprintf(`DROP DATABASE IF EXISTS %q`, dbName))
}
+105
View File
@@ -0,0 +1,105 @@
package tenant
import (
"context"
"os"
"strings"
"testing"
"github.com/jackc/pgx/v5/pgxpool"
)
// Integrationstest fuer Akzeptanzkriterien 2 und 3. Benoetigt eine echte
// Postgres-Instanz und wird ohne TEST_ADMIN_DSN uebersprungen, nicht als
// fehlgeschlagen gewertet — siehe Pruefungen-Ergebnis im PR.
//
// TEST_ADMIN_DSN muss auf die Wartungsdatenbank zeigen, z. B.:
//
// postgresql://postgres:postgres@localhost:5432/postgres?sslmode=disable
func TestProvision_CreatesIsolatedDatabases(t *testing.T) {
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)
}
defer adminPool.Close()
registryPool, err := pgxpool.New(ctx, adminDSN)
if err != nil {
t.Fatalf("registry pool: %v", err)
}
defer registryPool.Close()
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)
}
dsnTemplate := strings.Replace(adminDSN, "/postgres?", "/%s?", 1)
registry := NewRegistry(registryPool)
provisioner := NewProvisioner(adminPool, registry, dsnTemplate)
t.Cleanup(func() {
_, _ = registryPool.Exec(ctx, `DELETE FROM tenants WHERE slug IN ('acme','globex')`)
_, _ = adminPool.Exec(ctx, `DROP DATABASE IF EXISTS tenant_acme`)
_, _ = adminPool.Exec(ctx, `DROP DATABASE IF EXISTS tenant_globex`)
})
tenantA, err := provisioner.Provision(ctx, "acme", "Acme GmbH")
if err != nil {
t.Fatalf("provision acme: %v", err)
}
tenantB, err := provisioner.Provision(ctx, "globex", "Globex AG")
if err != nil {
t.Fatalf("provision globex: %v", err)
}
if tenantA.DBName == tenantB.DBName {
t.Fatalf("erwartet unterschiedliche db_name, beide sind %q", tenantA.DBName)
}
// Akzeptanzkriterium 3 / Pruefung 3: In der Datenbank von Tenant A existiert
// keine Verbindungsmoeglichkeit zu Tenant B, weil beide physisch getrennte
// Datenbanken sind, statt sich auf einen Query-Filter zu verlassen.
poolA, err := pgxpool.New(ctx, tenantA.DBDSN)
if err != nil {
t.Fatalf("connect tenant a: %v", err)
}
defer poolA.Close()
var globexVisible bool
err = poolA.QueryRow(ctx, `
SELECT EXISTS (
SELECT 1 FROM pg_catalog.pg_database WHERE datname = $1
)
`, tenantB.DBName).Scan(&globexVisible)
if err != nil {
t.Fatalf("pruefung tenant-trennung: %v", err)
}
// pg_database ist clusterweit sichtbar (Existenz der DB), aber die
// eigentliche Pruefung ist: aus poolA (verbunden mit tenant_acme) ist keine
// Tabelle/Zeile aus tenant_globex erreichbar, da current_database() getrennt ist.
var currentDB string
if err := poolA.QueryRow(ctx, `SELECT current_database()`).Scan(&currentDB); err != nil {
t.Fatalf("current_database: %v", err)
}
if currentDB != tenantA.DBName {
t.Fatalf("current_database() = %q, want %q — keine physische Trennung", currentDB, tenantA.DBName)
}
if currentDB == tenantB.DBName {
t.Fatalf("tenant a verbindung zeigt auf tenant b datenbank")
}
}
+69
View File
@@ -0,0 +1,69 @@
package tenant
import (
"context"
"fmt"
"github.com/jackc/pgx/v5"
"github.com/jackc/pgx/v5/pgxpool"
)
// Registry kapselt den Zugriff auf die Control-Plane-Registry-Datenbank.
// Sie enthaelt ausschliesslich Tenant-Metadaten (Akzeptanzkriterium 1) —
// niemals Geschaeftsdaten eines Mandanten.
type Registry struct {
pool *pgxpool.Pool
}
func NewRegistry(pool *pgxpool.Pool) *Registry {
return &Registry{pool: pool}
}
// insertTx schreibt den Tenant-Datensatz innerhalb einer laufenden Transaktion,
// damit Provisioner.Provision DB-Anlage und Registrierung atomar behandeln kann.
func (r *Registry) insertTx(ctx context.Context, tx pgx.Tx, t Tenant) (Tenant, error) {
row := tx.QueryRow(ctx, `
INSERT INTO tenants (slug, name, db_name, db_dsn, status)
VALUES ($1, $2, $3, $4, $5)
RETURNING id, created_at
`, t.Slug, t.Name, t.DBName, t.DBDSN, t.Status)
if err := row.Scan(&t.ID, &t.CreatedAt); err != nil {
return Tenant{}, fmt.Errorf("tenant registrieren: %w", err)
}
return t, nil
}
func (r *Registry) GetBySlug(ctx context.Context, slug string) (Tenant, error) {
var t Tenant
row := r.pool.QueryRow(ctx, `
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); err != nil {
return Tenant{}, fmt.Errorf("tenant laden: %w", err)
}
return t, nil
}
func (r *Registry) List(ctx context.Context) ([]Tenant, error) {
rows, err := r.pool.Query(ctx, `
SELECT id, slug, name, db_name, db_dsn, status, created_at
FROM tenants ORDER BY created_at
`)
if err != nil {
return nil, fmt.Errorf("tenants auflisten: %w", err)
}
defer rows.Close()
var out []Tenant
for rows.Next() {
var t Tenant
if err := rows.Scan(&t.ID, &t.Slug, &t.Name, &t.DBName, &t.DBDSN, &t.Status, &t.CreatedAt); err != nil {
return nil, fmt.Errorf("tenant lesen: %w", err)
}
out = append(out, t)
}
return out, rows.Err()
}
+42
View File
@@ -0,0 +1,42 @@
// Package tenant implements Core TEN-01: die Control-Plane-Registry und die
// Provisioning-Routine fuer physisch getrennte Mandanten-Datenbanken (Modell C).
package tenant
import (
"errors"
"regexp"
"time"
)
type Status string
const (
StatusActive Status = "active"
)
type Tenant struct {
ID string
Slug string
Name string
DBName string
DBDSN string
Status Status
CreatedAt time.Time
}
// slugPattern erzwingt sichere, als SQL-Identifier verwendbare Slugs, damit
// der Datenbankname niemals aus unkontrolliertem Nutzereingabe-Text gebaut wird.
var slugPattern = regexp.MustCompile(`^[a-z][a-z0-9_]{1,48}$`)
var ErrInvalidSlug = errors.New("tenant: slug muss mit Kleinbuchstaben beginnen und darf nur [a-z0-9_] enthalten (2-49 Zeichen)")
func ValidateSlug(slug string) error {
if !slugPattern.MatchString(slug) {
return ErrInvalidSlug
}
return nil
}
func dbNameForSlug(slug string) string {
return "tenant_" + slug
}
+35
View File
@@ -0,0 +1,35 @@
package tenant
import "testing"
func TestValidateSlug(t *testing.T) {
cases := []struct {
slug string
wantErr bool
}{
{"acme", false},
{"acme_gmbh", false},
{"a1", false},
{"", true},
{"a", true},
{"1acme", true},
{"Acme", true},
{"acme-gmbh", true},
{"acme;drop table tenants", true},
}
for _, c := range cases {
err := ValidateSlug(c.slug)
if (err != nil) != c.wantErr {
t.Errorf("ValidateSlug(%q) error = %v, wantErr %v", c.slug, err, c.wantErr)
}
}
}
func TestDBNameForSlug(t *testing.T) {
got := dbNameForSlug("acme")
want := "tenant_acme"
if got != want {
t.Errorf("dbNameForSlug() = %q, want %q", got, want)
}
}
+1
View File
@@ -0,0 +1 @@
DROP TABLE IF EXISTS tenants;
-10
View File
@@ -1,10 +0,0 @@
-- Control-plane registry: tenant list + connection info (Modell C).
-- Core TEN-01 (siehe core-kanban).
CREATE TABLE tenants (
id UUID PRIMARY KEY DEFAULT gen_random_uuid(),
slug TEXT NOT NULL UNIQUE,
name TEXT NOT NULL,
db_dsn TEXT NOT NULL,
status TEXT NOT NULL DEFAULT 'active',
created_at TIMESTAMPTZ NOT NULL DEFAULT now()
);
+14
View File
@@ -0,0 +1,14 @@
-- Control-plane registry: Tenant-Liste + Verbindungsinformationen (Modell C).
-- Enthaelt AUSSCHLIESSLICH Tenant-Metadaten, keine Geschaeftsdaten eines Mandanten.
-- Core TEN-01 (siehe core-kanban/tickets/TEN-01.md).
CREATE EXTENSION IF NOT EXISTS pgcrypto;
CREATE TABLE 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()
);
+1
View File
@@ -0,0 +1 @@
DROP TABLE IF EXISTS feature_flags;
+10
View File
@@ -0,0 +1,10 @@
-- Feature-Flags zentral je Mandant/Zielgruppe (LIC-02, siehe core-kanban/tickets/LIC-02.md).
-- Lebt in der Registry-DB, nicht pro Tenant-Datenbank — Flags sind eine
-- Core-weite Konfiguration, keine Mandanten-Geschaeftsdaten.
CREATE TABLE feature_flags (
key TEXT PRIMARY KEY,
enabled BOOLEAN NOT NULL DEFAULT false,
rollout_percentage INT NOT NULL DEFAULT 0 CHECK (rollout_percentage BETWEEN 0 AND 100),
target_tenant_slugs TEXT[] NOT NULL DEFAULT '{}',
updated_at TIMESTAMPTZ NOT NULL DEFAULT now()
);
+2
View File
@@ -0,0 +1,2 @@
DROP TABLE IF EXISTS module_credentials;
DROP TABLE IF EXISTS modules;
+18
View File
@@ -0,0 +1,18 @@
-- Modul-Registry & Aktivierungspruefung (API-02, siehe core-kanban/tickets/API-02.md).
CREATE TABLE modules (
name TEXT PRIMARY KEY,
version TEXT NOT NULL CHECK (version <> ''),
required_flags TEXT[] NOT NULL DEFAULT '{}',
registered_at TIMESTAMPTZ NOT NULL DEFAULT now()
);
-- Service-Credential je Modul-Instanz, bei Provisionierung ausgestellt
-- (Akzeptanzkriterium 4). secret_hash enthaelt NIEMALS das Secret im
-- Klartext, nur dessen SHA-256-Hash (Timing-safe-Vergleich beim Login,
-- Referenzmuster siehe AUD-02).
CREATE TABLE module_credentials (
module_name TEXT PRIMARY KEY REFERENCES modules(name),
client_id TEXT NOT NULL UNIQUE,
secret_hash BYTEA NOT NULL,
issued_at TIMESTAMPTZ NOT NULL DEFAULT now()
);
+1
View File
@@ -0,0 +1 @@
DROP TABLE tenant_keks;
+10
View File
@@ -0,0 +1,10 @@
-- Master-Key-Verwaltung & Tenant-Schluesselhierarchie (API-10, siehe
-- core-kanban/tickets/API-10.md) — EIN verpackter (mit dem Master-Key
-- umhuellter) Tenant-KEK je Mandant. Niemals der Master-Key selbst und
-- niemals ein Tenant-KEK im Klartext in dieser Tabelle.
CREATE TABLE tenant_keks (
tenant_id UUID PRIMARY KEY REFERENCES tenants(id),
wrapped_kek BYTEA NOT NULL,
created_at TIMESTAMPTZ NOT NULL DEFAULT now(),
rotated_at TIMESTAMPTZ
);
+23
View File
@@ -0,0 +1,23 @@
#!/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 CASCADE;"
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."
+24
View File
@@ -0,0 +1,24 @@
#!/usr/bin/env bash
# Ein-Kommando-Pruefung fuer den aktuellen Code-Stand auf dem Testhost:
# Registry+Tenant-DBs zuruecksetzen, dann build/vet/test in einem Rutsch.
# -p 1 ist Pflicht, da mehrere Pakete dieselbe physische Registry-Tabelle auf
# dem Testhost teilen (siehe [[project-nexarch-test-infra]]).
#
# Aufruf: NEXARCH_TEST_DB_PASSWORD=... ./scripts/run-checks.sh
set -euo pipefail
PASS="${NEXARCH_TEST_DB_PASSWORD:?Setze NEXARCH_TEST_DB_PASSWORD vor dem Aufruf}"
cd "$(dirname "$0")/.."
NEXARCH_TEST_DB_PASSWORD="$PASS" bash scripts/reset-test-env.sh
export TEST_ADMIN_DSN="postgresql://nexarch_test:${PASS}@localhost:5432/postgres?sslmode=disable"
echo "== go build =="
go build ./...
echo "== go vet =="
go vet ./...
echo "== go test (-p 1) =="
go test ./... -p 1 -count=1