Compare commits

..
Author SHA1 Message Date
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
77 changed files with 1007 additions and 4732 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 -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 {
-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
-43
View File
@@ -1,43 +0,0 @@
.PHONY: install run run-app run-worker lint fmt test build check
# Akzeptanzkriterium 1: ein Befehl installiert+startet App und Worker.
install: build
build:
go build ./...
run: build
@echo "Starte dms-app und dms-worker (Strg+C zum Beenden beider)"
@trap 'kill 0' EXIT; \
go run ./cmd/app & \
go run ./cmd/worker & \
wait
run-app:
go run ./cmd/app
run-worker:
go run ./cmd/worker
# Akzeptanzkriterium 2: Lint-/Format-Checks laufen lokal durch.
lint:
golangci-lint run ./...
fmt:
gofmt -l .
@test -z "$$(gofmt -l .)" || (echo "gofmt-Verstoesse gefunden, siehe oben" && exit 1)
test:
# -p 1: alle Integrationstest-Pakete teilen sich dieselbe physische
# Test-Datenbank (TEST_TENANT_DSN); parallele Paketausfuehrung wuerde
# sich gegenseitig ueberschreiben (dieselbe Konvention wie NEXARCH Core,
# siehe scripts/run-checks.sh im Core-Modul).
go test ./... -p 1 -count=1
# Setzt die geteilte Test-Datenbank zurueck, dann build/vet/test in einem
# Rutsch — analog zu NEXARCH Cores scripts/run-checks.sh.
check: build
NEXARCH_DMS_TEST_DB_PASSWORD="$${NEXARCH_DMS_TEST_DB_PASSWORD:?Setze NEXARCH_DMS_TEST_DB_PASSWORD vor dem Aufruf}" bash scripts/reset-test-env.sh
go vet ./...
golangci-lint run ./...
go test ./... -p 1 -count=1
-52
View File
@@ -1,52 +0,0 @@
# NEXARCH DMS
Dokumentenmanagement-Modul von NEXARCH. Vereint die Stärken von
paperless-ngx, Alfresco, Docspell und ecoDMS, vermeidet deren bekannte
Schwächen (siehe `known-issues-archivdms.md` im `dms-kanban/`-Ordner).
Identität, Rechte, Mandantenverwaltung, Authentifizierung, UI-Shell,
API-Grundgerüst und Benachrichtigungen kommen aus NEXARCH Core (siehe
`../` bzw. `../../core-kanban/`) — dieses Modul implementiert nur die
DMS-eigene Logik.
## Setup
Voraussetzung: Go 1.22+.
```bash
cd dms
make install # baut App und Worker
make run # startet beide (Strg+C beendet beide)
```
App läuft danach auf `:8090` (überschreibbar über
`NEXARCH_DMS_APP_LISTEN_ADDR`), `GET /healthz` liefert den Status.
## Struktur
- `cmd/app` — Anfrage-Dienst (HTTP), blockiert nie durch lange Aufgaben
- `cmd/worker` — Hintergrund-Dienst für lange laufende Aufgaben (Indexierung,
OCR, Storage-Vorgänge — folgen in FDN-02 ff.)
- `internal/shared` — von App und Worker gemeinsam genutzter Code
## Prüfungen
```bash
make fmt # gofmt-Verstöße brechen ab
make lint # golangci-lint
make test # go test ./...
```
## Branch- und Commit-Konvention
Gleiche Konvention wie NEXARCH Core:
- Branch je Ticket: `feature/<ticket-code>-<kurzbeschreibung>`, z. B.
`feature/fdn-02-datenmodell-migrationen`
- Commit-Nachricht beginnt mit dem Ticket-Code, z. B.
`FDN-02: datenmodell & migrationen`
- Ein Ticket = ein Branch. Schrittweise committen, Branch pushen, dann
anhalten (kein Merge, kein Deploy durch die bearbeitende Person selbst).
- Deutschsprachige Oberflächentexte, englischsprachige Bezeichner im Code.
- Keine Zugangsdaten/Schlüssel/Verbindungszeichenfolgen im Code —
ausschließlich über Umgebungsvariablen.
-94
View File
@@ -1,94 +0,0 @@
// app ist der Anfrage-Dienst (Request-Path) des DMS — getrennt vom Worker,
// damit lange Hintergrundaufgaben nie eine HTTP-Anfrage blockieren
// (Akzeptanzkriterium/Produkt-DNA: paperless-ngx-Trennung Dienst/Worker).
package main
import (
"context"
"log"
"net/http"
"os"
"time"
"github.com/jackc/pgx/v5/pgxpool"
"gitea.perlbach24.de/scripte/nexarch/dms/internal/jobqueue"
"gitea.perlbach24.de/scripte/nexarch/dms/internal/retentionclient"
"gitea.perlbach24.de/scripte/nexarch/dms/internal/retentiondestroy"
"gitea.perlbach24.de/scripte/nexarch/dms/internal/shared"
"gitea.perlbach24.de/scripte/nexarch/dms/internal/storage"
)
func main() {
addr := os.Getenv("NEXARCH_DMS_APP_LISTEN_ADDR")
if addr == "" {
addr = ":8090"
}
mux := http.NewServeMux()
mux.HandleFunc("/healthz", func(w http.ResponseWriter, r *http.Request) {
w.Header().Set("Content-Type", "application/json")
_, _ = w.Write([]byte(`{"status":"ok","service":"dms-app","version":"` + shared.Version + `"}`))
})
// DOC-16: RET-05-Client (Registrierung + Vernichtungs-Rückruf-Empfang).
// Optional, um cmd/app auch ohne diese Umgebungsvariablen weiter
// startbar zu halten (z. B. isolierte Tests anderer Endpunkte).
if dsn := os.Getenv("NEXARCH_DMS_APP_TENANT_DSN"); dsn != "" {
setupRetention(mux, dsn)
} else {
log.Print("dms-app: NEXARCH_DMS_APP_TENANT_DSN nicht gesetzt, RET-05-Anbindung (DOC-16) deaktiviert")
}
log.Printf("dms-app hoert auf %s (version %s)", addr, shared.Version)
if err := http.ListenAndServe(addr, mux); err != nil {
log.Fatal(err)
}
}
func setupRetention(mux *http.ServeMux, dsn string) {
ctx := context.Background()
pool, err := pgxpool.New(ctx, dsn)
if err != nil {
log.Fatalf("dms-app: datenbankverbindung: %v", err)
}
storageDir := os.Getenv("NEXARCH_DMS_APP_STORAGE_DIR")
signingSecret := os.Getenv("NEXARCH_DMS_APP_STORAGE_SIGNING_SECRET")
publicBaseURL := os.Getenv("NEXARCH_DMS_APP_STORAGE_PUBLIC_BASE_URL")
driver := storage.NewLocalDriver(storageDir, []byte(signingSecret), publicBaseURL)
usageEndpoint := os.Getenv("NEXARCH_DMS_APP_USAGE_ENDPOINT_URL")
usageClientID := os.Getenv("NEXARCH_DMS_APP_USAGE_CLIENT_ID")
usageClientSecret := os.Getenv("NEXARCH_DMS_APP_USAGE_CLIENT_SECRET")
tenantSlug := os.Getenv("NEXARCH_DMS_APP_TENANT_SLUG")
usageReporter := storage.NewHTTPUsageReporter(usageEndpoint, usageClientID, usageClientSecret, nil)
storageSvc := storage.NewService(driver, usageReporter, tenantSlug)
mux.HandleFunc("POST /retention/destroy-callback", retentiondestroy.DestroyCallbackHandler(pool, storageSvc))
retentionBaseURL := os.Getenv("NEXARCH_DMS_APP_RETENTION_BASE_URL")
retentionClass := os.Getenv("NEXARCH_DMS_APP_RETENTION_CLASS")
callbackURL := os.Getenv("NEXARCH_DMS_APP_RETENTION_CALLBACK_URL")
if retentionBaseURL == "" || retentionClass == "" || callbackURL == "" {
log.Print("dms-app: RET-05-Registrierung uebersprungen (RETENTION_BASE_URL/CLASS/CALLBACK_URL unvollstaendig)")
return
}
client := retentionclient.New(retentionBaseURL)
queue := jobqueue.NewQueue(pool, 5*time.Minute)
cfg := retentiondestroy.RegisterConfig{
ObjectType: retentiondestroy.ObjectTypeDocument,
RetentionClass: retentionClass,
CallbackURL: callbackURL,
}
regCtx, cancel := context.WithTimeout(ctx, 10*time.Second)
defer cancel()
queued, err := retentiondestroy.RegisterOrRequeue(regCtx, client, queue, cfg)
if err != nil {
log.Printf("dms-app: RET-05-registrierung fehlgeschlagen (requeued=%t): %v", queued, err)
} else {
log.Print("dms-app: bei archive RET-05 registriert")
}
}
-31
View File
@@ -1,31 +0,0 @@
package main
import (
"net/http"
"net/http/httptest"
"strings"
"testing"
)
// TestHealthz ist der Nachweis, dass der App-Dienst tatsaechlich startet und
// antwortet (Akzeptanzkriterium 1: "mit einem Befehl installieren und
// starten") — geprueft ueber den Handler direkt statt einen echten Port zu
// binden, damit der Test parallel und ohne Portkonflikte laufen kann.
func TestHealthz(t *testing.T) {
mux := http.NewServeMux()
mux.HandleFunc("/healthz", func(w http.ResponseWriter, r *http.Request) {
w.Header().Set("Content-Type", "application/json")
_, _ = w.Write([]byte(`{"status":"ok","service":"dms-app","version":"test"}`))
})
req := httptest.NewRequest(http.MethodGet, "/healthz", nil)
rec := httptest.NewRecorder()
mux.ServeHTTP(rec, req)
if rec.Code != http.StatusOK {
t.Fatalf("status = %d, want %d", rec.Code, http.StatusOK)
}
if !strings.Contains(rec.Body.String(), `"status":"ok"`) {
t.Fatalf("unerwarteter body: %s", rec.Body.String())
}
}
-86
View File
@@ -1,86 +0,0 @@
// worker ist der Hintergrund-Dienst des DMS — verarbeitet lange laufende
// Aufgaben (Indexierung, OCR, Storage-Vorgaenge in spaeteren Kacheln),
// getrennt vom App-Prozess (siehe cmd/app). DOC-16: verarbeitet zusaetzlich
// Requeue-Jobs fuer fehlgeschlagene RET-05-Registrierungen.
package main
import (
"context"
"errors"
"log"
"os"
"os/signal"
"syscall"
"time"
"github.com/jackc/pgx/v5/pgxpool"
"gitea.perlbach24.de/scripte/nexarch/dms/internal/jobqueue"
"gitea.perlbach24.de/scripte/nexarch/dms/internal/retentionclient"
"gitea.perlbach24.de/scripte/nexarch/dms/internal/retentiondestroy"
"gitea.perlbach24.de/scripte/nexarch/dms/internal/shared"
)
func main() {
log.Printf("dms-worker gestartet (version %s)", shared.Version)
ctx, stop := signal.NotifyContext(context.Background(), os.Interrupt, syscall.SIGTERM)
defer stop()
var (
queue *jobqueue.Queue
client *retentionclient.Client
)
if dsn := os.Getenv("NEXARCH_DMS_WORKER_TENANT_DSN"); dsn != "" {
pool, err := pgxpool.New(ctx, dsn)
if err != nil {
log.Fatalf("dms-worker: datenbankverbindung: %v", err)
}
queue = jobqueue.NewQueue(pool, 5*time.Minute)
client = retentionclient.New(os.Getenv("NEXARCH_DMS_WORKER_RETENTION_BASE_URL"))
} else {
log.Print("dms-worker: NEXARCH_DMS_WORKER_TENANT_DSN nicht gesetzt, DOC-16-Requeue-Verarbeitung deaktiviert")
}
ticker := time.NewTicker(30 * time.Second)
defer ticker.Stop()
for {
select {
case <-ctx.Done():
log.Println("dms-worker beendet")
return
case <-ticker.C:
if queue != nil {
processRetentionRegisterJobs(ctx, queue, client)
}
}
}
}
// processRetentionRegisterJobs holt und verarbeitet ALLE aktuell
// abholbaren doc16_retention_register-Jobs (Akzeptanzkriterium 3: Requeue
// statt Absturz/Verwerfen — Fail() haengt selbst nie ab, sondern setzt
// pending mit Backoff oder dead_letter, siehe FDN-04).
func processRetentionRegisterJobs(ctx context.Context, queue *jobqueue.Queue, client *retentionclient.Client) {
for {
job, err := queue.Dequeue(ctx, "dms-worker", []string{retentiondestroy.JobTypeRegister})
if err != nil {
if !errors.Is(err, jobqueue.ErrNoJobAvailable) {
log.Printf("dms-worker: dequeue fehlgeschlagen: %v", err)
}
return
}
if err := retentiondestroy.ProcessRegisterJob(ctx, client, job.Payload); err != nil {
log.Printf("dms-worker: retention-registrierung erneut fehlgeschlagen (job %s): %v", job.ID, err)
if failErr := queue.Fail(ctx, job.ID, err); failErr != nil {
log.Printf("dms-worker: job als fehlgeschlagen markieren: %v", failErr)
}
continue
}
if err := queue.Complete(ctx, job.ID); err != nil {
log.Printf("dms-worker: job abschliessen: %v", err)
} else {
log.Printf("dms-worker: retention-registrierung erfolgreich nachgeholt (job %s)", job.ID)
}
}
}
-78
View File
@@ -1,78 +0,0 @@
# DOC-01 Prüfprotokoll: Upload-API & Chunk-Handling
Welle 3. Voraussetzung: FDN-02, FDN-03, FDN-09 (alle Status "Fertig").
## Umsetzung
`internal/upload` verbindet FDN-02 (Datenmodell), FDN-03 (Storage) und
FDN-09 (Verschlüsselung) zur ersten echten Ingest-Strecke:
- `Validator` — MIME-Positivliste (Akzeptanzkriterium 2) und Größenlimit
(Schutz vor Speicherbomben) VOR jedem entgegengenommenen Byte geprüft.
- `SessionStore`/`upload_sessions` — Fortschritt (`bytes_received`) je
Sitzung, `GREATEST`-Update verhindert Rückschritt bei erneut zugestellten
Chunks (Akzeptanzkriterium 1).
- `Staging` — Chunks landen lokal über `os.File.WriteAt` an ihrem Offset,
beliebige Reihenfolge/Wiederholung möglich, bevor das vollständige Objekt
verschlüsselt im Storage landet.
- `Service.Complete` — SHA-256 UND Verschlüsselung lesen denselben
Byte-Strom in einem Durchlauf (`io.TeeReader`), garantiert Prüfsumme auf
exakt dem, was verschlüsselt wurde (Akzeptanzkriterium 4, Reihenfolge aus
"Bekannte Fehler vermeiden" eingehalten: Hash → verschlüsseln → ablegen).
Dokument+Revision werden in EINER Postgres-Transaktion angelegt
(Akzeptanzkriterium 3), das Storage-`Put` erfolgt innerhalb derselben
Transaktionsspanne vor dem Commit.
- Neue Spalte `file_revisions.wrapped_dek` (Migration `0004`) für den
FDN-09-Envelope-Wrapper.
## Prüfungen
| # | Prüfung | Ergebnis |
|---|---|---|
| 1 | Upload von 1 GB Datei erfolgreich | **eingeschränkt bestanden** — siehe Abschnitt "Skalierungs-Einschränkung" unten: real bis 300 MiB auf 192.168.1.131 verifiziert (Upload → Verschlüsselung → Ablage → Checksummen-Abgleich → Entschlüsselungs-Round-Trip, alles exakt), volle 1 GiB auf diesem 4-GB-RAM-Testhost mangels Arbeitsspeicher nicht möglich |
| 2 | Abbruch bei 50% und Fortsetzung ergibt identische Prüfsumme | **bestanden**`TestUploadChunk_ResumeAfterAbortProducesIdenticalChecksum`: 200 KiB Zufallsinhalt, erste Hälfte hochgeladen, `Status`-Abfrage (wie ein neu verbindender Client) bestätigt exakt die Hälfte empfangen, Fortsetzung ab genau diesem Offset, Endprüfsumme stimmt exakt mit der Prüfsumme des vollständigen Originalinhalts überein |
| 3 | Parallel-Upload von 20 Dateien ohne Datenverlust | **bestanden**`TestParallelUploads_NoDataLoss`: 20 gleichzeitige vollständige Upload-Durchläufe (eigene Session je Datei), alle 20 liefern eindeutige Dokument-IDs, jede mit korrekter, individueller Prüfsumme |
| 4 | Nach Upload ist `file_revisions.checksum_sha256` befüllt und entspricht der SHA-256 des hochgeladenen Klartexts | **bestanden**`TestCompleteUpload_CreatesDocumentAndRevisionTransactionally` UND die 300-MiB-Verifikation: gespeicherte Prüfsumme stimmt exakt mit unabhängig berechneter Prüfsumme des Originalinhalts überein |
## Skalierungs-Einschränkung (real gefunden, nicht vorab bekannt)
Ein echter 1-GiB-Durchlauf auf 192.168.1.131 (4 GB RAM, kein Swap) endete
mit `signal: killed` (OOM). Ursache: sowohl `crypto.EncryptStream`/
`DecryptStream` (FDN-09) als auch `storage.S3Driver.Put`/`HTTPUsageReporter`-
Pfad (FDN-03) puffern den gesamten Inhalt vollständig im Speicher
(`io.ReadAll`) statt echt zu streamen — bereits in den jeweiligen
Prüfprotokollen als bewusste Vereinfachung ("kleinste Lösung") dokumentiert,
hier zeigt sich der reale Preis dafür: mehrere ~1-GiB-Kopien (Klartext,
Chiffretext, ggf. weitere beim Round-Trip) gleichzeitig im Speicher
überschreiten 4 GB deutlich. Reduziert auf 300 MiB erfolgreich und
vollständig verifiziert (Upload, Verschlüsselung, Ablage, Prüfsummen-Abgleich,
Entschlüsselungs-Round-Trip — alles exakt, keine Verkürzung der eigentlichen
Prüftiefe, nur der Dateigröße).
**Nicht in dieser Kachel behoben** (wäre Umbau von FDN-03/FDN-09, „kein
Umbau angrenzender Bereiche"): echtes Streaming (segmentierte AEAD-
Verschlüsselung, `io.Copy`-basierter Storage-Pfad statt `io.ReadAll`) wäre
nötig, um sehr große Dateien auf speicherschwachen Hosts zuverlässig zu
verarbeiten. Empfehlung: eigenes Folgeticket, sobald reale Dateigrößen im
GB-Bereich Teil der Anforderungen werden — auf einem Host mit mehr RAM
wäre der volle 1-GiB-Test hingegen ohne Codeänderung durchführbar, die
Einschränkung ist eine Ressourcen-, keine Korrektheitsfrage.
## 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 -> 6/6 Pakete ok, 0 Fehlschläge (5 neue upload-Tests)
```
## Gesamtergebnis
**Bestanden, mit dokumentierter Skalierungs-Einschränkung.** Alle vier
Akzeptanzkriterien erfüllt. Von den vier Pflichtprüfungen sind drei
uneingeschränkt bestanden; Prüfung 1 (1 GB) wurde bei reduzierter, aber
vollständig verifizierter Dateigröße (300 MiB) bestanden — die Differenz
liegt nachweislich an der Testhost-Ressourcenausstattung, nicht an der
Korrektheit der Implementierung (Mechanismus bei 300 MiB exakt bewiesen,
nichts an der Logik ist größenabhängig außer dem Speicherbedarf selbst).
-79
View File
@@ -1,79 +0,0 @@
# DOC-16 Prüfprotokoll: RET-05-Client (Registrierung & Vernichtungs-Rückruf)
Voraussetzung DOC-01, FDN-04 (beide bereits Fertig), Archive RET-05 +
RET-09 (RET-05 real als Dienst gestartet — beide bereits Fertig, externes
Board).
## Umsetzung
- `dms/internal/retentionclient` schlanker HTTP-Client für Archive
RET-05/RET-09 (`POST /register`).
- `dms/internal/retentiondestroy`:
- `RegisterOrRequeue` registriert `dms_document` real gegen RET-09.
Bei Fehlschlag: FDN-04-Requeue (`JobTypeRegister`), Fehler wird
IMMER zurückgegeben (Aufrufer protokolliert, kein stilles
Verwerfen).
- `ProcessRegisterJob` vom Worker aufgerufener Retry-Handler.
- `DestroyCallbackHandler` empfängt den Vernichtungs-Rückruf,
löscht ALLE Datei-Revisionen physisch über `storage.Service.Delete`
(meldet Größenänderung an Core), markiert das Dokument als
vernichtet (`deleted_at`), antwortet erst danach mit 2xx.
- `dms/cmd/app` mountet `POST /retention/destroy-callback`, versucht
beim Start EINMAL die Registrierung (`RegisterOrRequeue`).
- `dms/cmd/worker` verarbeitet `doc16_retention_register`-Jobs aus der
FDN-04-Queue (Retry bei vorherigem Fehlschlag).
**Nur `dms_document` registriert, kein separater "Anhang"-Typ:** DOC-01/
FDN-02 modellieren keine von Dokumenten getrennte Anhang-Entität — eine
solche Unterscheidung hätte einen Umbau von FDN-02 erfordert (kein
Umbau angrenzender Bereiche). Dokumentiert als bewusste Scope-Grenze.
## Prüfungen
| # | Prüfung | Ergebnis |
|---|---|---|
| 1 | Registrierung real gegen einen Test-RET-05-Endpunkt durchgeführt und verifiziert | **bestanden** `TestRegisterOrRequeue_RealRegistrationAgainstTestEndpoint` (echter HTTP-Roundtrip gegen einen `httptest`-Server mit RET-05-Vertrag); zusätzlich real auf 131: `dms-app` gegen den laufenden `nexarch-archive-moduleadapter-api.service` (RET-09, Port 8095) gestartet, echte Zeile in `module_registrations` verifiziert (`module_name=dms, object_type=dms_document, retention_class=dms-standard, callback_url=...`) |
| 2 | Vernichtungs-Rückruf gegen einen echten Testfall ausgelöst, DMS-Löschung nachweislich erfolgt, 2xx-Antwort gesendet | **bestanden** `TestDestroyCallbackHandler_RealDeletionAnd2xx`: echtes Dokument mit echter Datei im `LocalDriver`-Speicher angelegt, Rückruf ausgelöst, Datei nachweislich nicht mehr lesbar, `deleted_at` gesetzt, Nutzungsmeldungen (+Put/-Delete) real erfasst; zusätzlich real auf 131 per `curl` reproduziert: Datei im echten Dateisystem verschwunden, `documents.deleted_at` real in Postgres gesetzt, HTTP 200 |
| 3 | Simulierter Nicht-2xx-Fehler bei der Registrierung führt nachweislich zu einem Requeue-Eintrag in der Jobqueue, kein Absturz | **bestanden** `TestRegisterOrRequeue_FailedRegistrationCreatesRequeueEntry`: fake-RET-05-Server liefert 500, echte `pending`-Zeile in `processing_jobs` nachgewiesen, Fehler wird zurückgegeben statt verschluckt; `TestProcessRegisterJob_SucceedsOnRetry` beweist zusätzlich, dass der Requeue-Job vom Worker erfolgreich nachgeholt werden kann (kein Sackgassen-Zustand) |
## Reale Betriebs-Erkenntnis (dokumentiert, nicht verschwiegen)
Beim Live-Test auf 131 zeigte sich: der konfigurierte
Nutzungsmeldungs-Endpunkt (`storage.HTTPUsageReporter`, Ziel wäre Core
API-06s `POST /internal/resync/usage`) ist auf dem aktuell laufenden
`nexarch-core.service` NICHT gemountet (`curl` liefert 404) — Core
API-06 ist als Ticket zwar fertig, aber der Dienst auf 131 läuft
offenbar auf einem älteren Stand ohne diese Route (derselbe
"fertig, aber nicht überall deployed"-Befund wie bei anderen Tickets
dieser Session, hier bei Core selbst statt bei DOC-16). Für den
Live-Beweis wurde ein simulierter Stub-Endpunkt auf Port 8097
eingesetzt (nur für die Dauer des Tests, danach entfernt) — die
Go-Integrationstests selbst verwenden ohnehin einen echten
`fakeUsageReporter`, nicht den HTTP-Reporter, und sind von dieser
Betriebslücke unberührt. DOC-16s eigener Code (`HTTPUsageReporter`) ist
korrekt implementiert und bereits vor DOC-16 fertig (FDN-03-Bestandteil);
die fehlende Route ist ein Core-seitiges Deploy-Thema, kein DOC-16-Defekt.
## Build/Test-Ergebnis (192.168.1.131)
```
go build ./... -> clean
go vet ./... -> clean
golangci-lint run ./... -> 0 issues
go test ./internal/retentiondestroy/... -> alle bestanden
```
**Hinweis:** `go test ./... -p 1` zeigt einen Fehlschlag in
`internal/migrate` (`erwartet mindestens 1 angewendete migration auf
leerer db`) — reale Altlast der geteilten Test-DB aus früheren
Testläufen dieser Session, `git diff --stat` bestätigt: DOC-16 hat
`internal/migrate` nicht berührt. Alle von DOC-16 tatsächlich berührten
Pakete (`cmd/app`, `internal/jobqueue`, `internal/retentiondestroy`,
`internal/storage`, `internal/upload`, `internal/crypto`) sind grün.
## Gesamtergebnis
**Bestanden.** Alle drei Akzeptanzkriterien und alle drei
Pflichtprüfungen real erfüllt, inklusive Live-Nachweis auf 131 gegen
den echten RET-09-Dienst (Registrierung) und einen echten, per curl
ausgelösten Vernichtungs-Rückruf mit tatsächlich gelöschter Datei.
-42
View File
@@ -1,42 +0,0 @@
# FDN-01 Prüfprotokoll: Repository & Projektgerüst
Welle 1, keine Vorbedingungen. Verzeichnis `code/dms/` im bestehenden
NEXARCH-Repository (Monorepo-Entscheidung, siehe Rückfrage im
Session-Verlauf: DMS als Unterordner statt eigenes Gitea-Repo).
## Struktur
- `cmd/app` — Anfrage-Dienst (HTTP, Port 8090 per Default)
- `cmd/worker` — Hintergrund-Dienst (getrennter Prozess)
- `internal/shared` — gemeinsam genutzter Code
- `go.mod` — eigenes Modul `gitea.perlbach24.de/scripte/nexarch/dms`,
unabhängig vom Core-Modul (kein gemeinsames `go.mod`, um Abhängigkeits-
versionen beider Module unabhängig weiterzuentwickeln)
## Prüfungen
| # | Prüfung | Zielwert | Ergebnis |
|---|---|---|---|
| 1 | Frischer Clone baut ohne manuelle Nacharbeit | `make install` läuft ohne Fehler | **bestanden**`go build ./...` clean auf 192.168.1.131 |
| 2 | Lint-Fehler brechen den Build ab | `make lint` liefert Exit-Code ≠ 0 bei echtem Verstoß | **bestanden** — absichtlich eingefügte ungenutzte Variable liefert Exit-Code 2, Fund korrekt lokalisiert (`declared and not used`) |
| 3 | README-Setupanleitung von zweiter Person nachvollzogen | — | **nicht durchgeführt** — keine zweite Person in dieser autonomen Sitzung verfügbar (gleiche Methodik-Abweichung wie QA-05 Prüfung 2/QA-09 Bildschirmleser-Durchlauf); ersatzweise die Anleitung selbst Schritt für Schritt auf einer frischen Kopie (`rsync` auf 192.168.1.131) nachvollzogen: `make install``make run``curl /healthz``{"status":"ok",...}`. |
Zusätzlich (nicht explizit gefordert, aber Teil von Akzeptanzkriterium 1
„installieren UND starten"): `make run` startet App und Worker parallel,
`GET /healthz` antwortet mit `200 {"status":"ok","service":"dms-app",...}`
innerhalb von 2 Sekunden nach Start.
## Build/Test-Ergebnis (192.168.1.131)
```
make build -> clean
make fmt -> clean (keine gofmt-Verstoesse)
make lint -> clean (golangci-lint v1.62.2: govet, staticcheck, errcheck, unused, ineffassign, gofmt, goimports)
make test -> 1/1 Pakete mit Tests ok (cmd/app), 0 Fehlschlaege
```
## Gesamtergebnis
**Bestanden**, mit einer dokumentierten Methodik-Abweichung (Prüfung 3,
Vier-Augen-Nachvollzug) mangels zweiter Person — durch Selbst-Nachvollzug auf
frischer Kopie ersetzt.
-71
View File
@@ -1,71 +0,0 @@
# FDN-02 Prüfprotokoll: Datenmodell & Migrationen
Welle 2. Voraussetzung: FDN-01 (Status "Fertig").
## Datenmodell
`migrations/tenant/0001_documents.up.sql` — läuft in der physisch isolierten
Tenant-Datenbank (Modell C, siehe Core TEN-01), keine `tenant_id`-Spalte.
| Entität | Tabelle | Beziehungen |
|---|---|---|
| Ordner | `folders` | selbstreferenzierend (`parent_folder_id`), `created_by``users(id)` |
| Dokument | `documents` | `folder_id``folders`, `current_revision_id``file_revisions`, `created_by``users(id)` |
| Datei-Revision | `file_revisions` | `document_id``documents`, `created_by``users(id)`, `UNIQUE(document_id, revision_number)` |
| Tag | `tags` | — |
| Tag-Zuordnung | `document_tags` | `document_id``documents`, `tag_id``tags` |
| Metadatenfeld | `metadata_fields` | — |
| Metadatenwert | `document_metadata_values` | `document_id``documents`, `field_id``metadata_fields` |
`created_by`/Benutzerbezug referenziert `users(id)` aus Core IAM-01 (Auth
liegt vollständig in Core, siehe "Nicht Bestandteil") — DMS legt `users`
nicht selbst an, setzt die Tabelle als bereits vorhanden voraus (dieselbe
physische Tenant-Datenbank).
**Indizes:** `folders(parent_folder_id)`, `documents(folder_id)`,
`documents(created_by)`, `file_revisions(document_id)`,
`document_tags(tag_id)`, `document_metadata_values(field_id)`.
## Migrationsmechanik
`internal/migrate` — eigenständiger, minimaler Runner (kein ORM,
`*.up.sql`/`*.down.sql`-Paare), `schema_migrations`-Tabelle als
Fortschrittsspeicher (dasselbe Prinzip wie Core, hier eigenständig
implementiert, da DMS ein eigenes Go-Modul ist und Cores `internal/`-Pakete
nicht importieren kann).
## Prüfungen
| # | Prüfung | Ergebnis |
|---|---|---|
| 1 | Migration auf leerer DB und auf bestehender DB getestet | **bestanden**`TestUp_OnEmptyAndExistingDB`: erster Lauf legt alle 7 Tabellen an, zweiter Lauf gegen dieselbe (jetzt bestehende) DB wendet 0 neue Migrationen an (über `schema_migrations` erkannt) |
| 2 | Rollback stellt Vorzustand wieder her | **bestanden**`TestDownOne_RestoresPreviousState`: nach `DownOne` existiert keine der 7 Tabellen mehr, zweiter `DownOne`-Aufruf ohne verbleibende Migration liefert korrekt leeren String statt Fehler |
| 3 | Fremdschlüssel-Constraints durch Negativtests belegt | **bestanden**`TestForeignKeyConstraints_RejectInvalidReferences`, 4 Fälle: Dokument mit unbekanntem Ordner, unbekanntem Ersteller, Datei-Revision mit unbekanntem Dokument, Tag-Zuordnung mit unbekanntem Tag — alle vier korrekt abgewiesen |
Zusätzlich (Akzeptanzkriterium 3, Seed-Datensatz): `migrations/tenant/seed/dev_seed.sql`
manuell gegen eine frische Test-DB mit einer `users`-Zeile ausgeführt (siehe
Sitzungsprotokoll) — legt Ordner, Dokument mit Revision, Tag und
Metadatenfeld+-wert an, per Abfrage bestätigt (`Beispieldokument`,
`Beispiel-Tag`, `rechnungsnummer` vorhanden). Schlägt bewusst mit
sprechender Fehlermeldung fehl, wenn noch kein Benutzer existiert (DMS legt
`users` nicht selbst an).
## Bekannte Fehler vermeiden (aus Ticket)
„Fehlende Lock-/Sum-Datei blockiert CI" — `go.sum` ist committet (siehe
`git status`/Commit-Diff), `go mod tidy` auf 192.168.1.131 ausgeführt und
Ergebnis übernommen.
## Build/Test-Ergebnis (192.168.1.131)
```
go build ./... -> clean
go vet ./... -> clean
make lint -> clean (golangci-lint)
go test ./... -v -count=1 -> 3/3 Pakete mit Tests ok (cmd/app, internal/migrate), 0 Fehlschläge
```
## Gesamtergebnis
**Bestanden.** Alle drei Akzeptanzkriterien und alle drei Pflichtprüfungen
erfüllt und belegt.
-84
View File
@@ -1,84 +0,0 @@
# FDN-03 Prüfprotokoll: Objekt-Storage-Abstraktion
Welle 2. Voraussetzung: FDN-01 (Status "Fertig"), Core LIC-05 (Status
"Fertig").
## Umsetzung
`internal/storage`:
- `Driver`-Interface (Akzeptanzkriterium 1): `Put`/`Get`/`Delete`/`SignedURL`.
- `LocalDriver` — Entwicklungs-Treiber, Dateisystem, signierte URLs über
HMAC-SHA256 (timing-safe verglichen, `crypto/subtle`, dieselbe Konvention
wie Core IAM-15).
- `S3Driver` — Produktions-Treiber, S3-kompatibel (`aws-sdk-go-v2`),
presigned URLs über `s3.PresignClient`.
- `ObjectKey(documentID, revisionID)` — Pfadschema `documents/<id>/revisions/<id>`
innerhalb des bereits mandantenspezifischen Buckets (Akzeptanzkriterium 3;
die Bucket-Trennung selbst ist Core TEN-01).
- `Service` — verbindet `Driver` mit `UsageReporter`: jeder `Put`/`Delete`
löst genau eine Nutzungsmeldung mit der tatsächlichen Objektgröße aus
(Akzeptanzkriterium 4). Repository-Code soll ausschließlich `Service`
aufrufen, nie einen `Driver` direkt.
- `HTTPUsageReporter` — meldet über Cores Service-Credential-authentifizierten
Resync-Endpunkt (`internal/resync.Handler.UsageHandler`, API-06/AUD-06-Muster),
Metrikname `storage_bytes` (gespiegelt aus Core `internal/usage.StorageBytesMetric`,
LIC-05 — DMS kann Cores `internal/`-Pakete als eigenes Go-Modul nicht
importieren).
## Wichtiger Befund: Core-Endpunkt noch nicht live verdrahtet
`internal/resync.Handler` (die Gegenstelle für `HTTPUsageReporter`) ist im
Core-Modul vollständig implementiert und getestet, aber **in keinem
`cmd/*/main.go` registriert** (per `grep` bestätigt, Stand
2026-08-29) — dieselbe Fehlerklasse wie der QA-05/AUD-06-Befund
(Bausteine existieren, sind aber nicht in einen laufenden Dienst verdrahtet).
`HTTPUsageReporter` ist daher gegen den **dokumentierten Vertrag** (exakte
Feldnamen/Header aus `internal/resync/handler.go` gelesen) getestet, nicht
gegen eine echte laufende Core-Instanz. Prüfung 4 ist damit im Rahmen dessen
erfüllt, was DMS beeinflussen kann — die Lücke auf Core-Seite ist ein
Core-Board-Thema (Empfehlung: analog AUD-06 ein Ticket "Resync-Endpunkt in
Core-Server verdrahten" anlegen), nicht Bestandteil dieser DMS-Kachel.
## Prüfungen
| # | Prüfung | Ergebnis |
|---|---|---|
| 1 | Round-Trip-Test Upload/Download je Treiber | **bestanden**`TestLocalDriver_RoundTrip` (Dateisystem) und `TestS3Driver_RoundTrip` (echtes MinIO auf 192.168.1.131, kein Mock) |
| 2 | Abgelaufene signierte URL wird abgewiesen | **bestanden**`TestLocalDriver_SignedURL_ExpiredIsRejected` (Signatur-/Ablauflogik) UND manuell gegen echtes MinIO verifiziert: presigned URL liefert `200` innerhalb der Gültigkeit, `403` nach Ablauf (2s TTL, siehe Sitzungsprotokoll) |
| 3 | Verhalten bei fehlendem Objekt liefert klaren Fehler | **bestanden**`TestLocalDriver_MissingObject`/`TestS3Driver_MissingObject`: beide Treiber liefern `ErrNotFound` für `Get` UND `Delete` eines nicht existierenden Objekts |
| 4 | Melde-Aufruf an Core LIC-05 bei Schreib-/Löschvorgang nachweislich ausgelöst, korrekte Größe | **bestanden** (mit Einschränkung s.o.) — `TestService_PutReportsPositiveDelta`/`TestService_DeleteReportsNegativeDelta` (Fake-Reporter zeichnet Aufrufe auf, prüft Tenant/Metrik/Delta) UND `TestHTTPUsageReporter_SendsCorrectContractToCore` (echter HTTP-Request gegen `httptest.Server`, der Cores Vertrag nachbildet — Header, JSON-Feldnamen) |
## Build/Test-Ergebnis (192.168.1.131)
```
go build ./... -> clean
go vet ./... -> clean
make lint -> clean (golangci-lint v2.1.6, aus Quelle mit go1.24.4 gebaut,
da v1.63.4 den Zielstand go1.24 nicht linten konnte —
.golangci.yml auf v2-Konfigurationsformat migriert)
go test ./internal/storage/... -v -count=1 -> 11/11 Tests ok (3 S3-Tests real
gegen lokal installiertes MinIO statt uebersprungen)
```
## Offener Punkt: Prüfsummen-Schreibpfad noch nicht befüllt
`file_revisions.checksum_sha256` (FDN-02) wird aktuell von **keinem**
Schreibpfad befüllt oder verifiziert — `internal/storage.Service.Put`
berechnet keine Inhalts-Prüfsumme (das einzige SHA256 im Paket ist die
HMAC-Signatur lokaler URLs, siehe oben, unabhängig vom Dateiinhalt). Die
Spalte existiert seit FDN-02 ungenutzt. Nachgetragen als
Akzeptanzkriterium/Prüfung in `DOC-01` (Upload-API), das den Hash auf
Klartext berechnen und transaktional persistieren muss — Voraussetzung für
Duplikaterkennung (`DOC-02`) und die spätere Integritätsprüfung
(Archive `BAK-08`, siehe Sitzungsprotokoll 2026-08-29 zu externem,
selbst nicht überwachtem Kunden-S3-Storage).
## Gesamtergebnis
**Bestanden**, mit einer dokumentierten Abhängigkeit auf Core-Seite
(Abschnitt "Wichtiger Befund") — Core muss `internal/resync.Handler` noch in
einen laufenden Dienst verdrahten, bevor `HTTPUsageReporter` echte
Nutzungsmeldungen an eine Produktivinstanz senden kann. Alle vier
Akzeptanzkriterien und alle vier Pflichtprüfungen im Rahmen des
DMS-seitigen Scopes erfüllt.
-77
View File
@@ -1,77 +0,0 @@
# FDN-04 Prüfprotokoll: Job-Queue & Worker-Runtime
Welle 3. Voraussetzung: FDN-02 (Status "Fertig").
## Umsetzung
`internal/jobqueue`:
- `migrations/tenant/0002_processing_jobs.up.sql``processing_jobs`-Tabelle
(Status `pending`/`processing`/`succeeded`/`failed`/`dead_letter`,
`attempts`/`max_attempts`, `available_at` für Backoff-Terminierung,
`locked_at`/`locked_by` für die Sperre, `idempotency_key` UNIQUE).
- `Queue.Enqueue` — reiht ein, mit optionalem `idempotency_key` (Dedup bei
Doppelzustellung, `ON CONFLICT DO UPDATE ... RETURNING id`).
- `Queue.Dequeue``FOR UPDATE SKIP LOCKED`, holt entweder einen fälligen
`pending`-Job oder einen `processing`-Job, dessen Sperre älter als
`staleLockAfter` ist (Absturz-Wiedervorlage). Backoff-Intervallarithmetik
über `LEAST(attempts, 10) * interval '30 seconds'` — arithmetischer
Cast, keine String-Konkatenation (siehe "Bekannte Fehler vermeiden").
- `Queue.Complete`/`Queue.Fail` — bei erschöpften Versuchen wandert der Job
in `dead_letter`.
- `Queue.RequeueDeadLetter` — manuelle Wiederholung eines DLQ-Eintrags.
- `Queue.Status` — Job-Status abfragbar.
- `Worker`/`Handler` — In-Prozess-Worker-Goroutine, pollt und ruft `Handler`
je Job auf.
## Prüfungen
| # | Prüfung | Ergebnis |
|---|---|---|
| 1 | Absturz eines Workers führt zu erneuter Zustellung | **bestanden**`TestDequeue_StaleLockIsRedelivered`: Job wird von `worker-crashed` gesperrt, NIE completed/failed (simulierter Absturz); sofortiger erneuter Dequeue-Versuch liefert `ErrNoJobAvailable` (Sperre noch frisch), nach Ablauf von `staleLockAfter` liefert `worker-2` denselben Job |
| 2 | Idempotenz bei Doppelzustellung nachgewiesen | **bestanden**`TestEnqueue_IdempotencyKeyPreventsDuplicate`: zweifache Einreihung mit gleichem `idempotency_key` erzeugt nachweislich nur 1 Zeile (per Abfrage bestätigt) |
| 3 | DLQ-Eintrag manuell wiederholbar | **bestanden**`TestRequeueDeadLetter`: Job nach erschöpften Versuchen in `dead_letter`, `RequeueDeadLetter` setzt zurück auf `pending` mit `attempts=0`; Requeue eines NICHT-DLQ-Jobs wird korrekt abgewiesen |
## Reale Fehler gefunden und behoben (kein Vorab-Wissen, beim Testen entdeckt)
1. **pgx-Typinferenz-Fehler bei ungenutztem Parameter**: `Dequeue`s SQL
übergab `workerID` als `$1`, ohne es in der Query zu referenzieren —
Postgres/pgx konnte den Typ von `$1` dadurch nicht ableiten
(`SQLSTATE 42P18`). Behoben durch Entfernen des toten Parameters
(workerID wird erst im nachfolgenden `UPDATE` gebraucht).
2. **`$2::text[]`-Cast mit untypisiertem `nil`**: `typeFilter any` (statt
`[]string`) ließ pgx den Zieltyp des Casts nicht auflösen. Behoben durch
`[]string`-Typisierung der Variable.
3. **Testinfrastruktur-Drift über Sitzungsgrenzen**: `dms_tenant_test`
sammelte über mehrere Testläufe (FDN-02/03/04) `schema_migrations`-Zustand
an, wodurch `internal/migrate`s Rollback-Test nur noch einen Teil der
Tabellen zurückrollte. Neues `scripts/reset-test-env.sh` (Datenbank
droppen+neu anlegen, analog Core `scripts/reset-test-env.sh`) sowie
`make check`-Target (Reset+vet+lint+test in einem Rutsch) behoben das
strukturell. Zusätzlich fehlte `-p 1` im `test`-Target — mehrere
Testpakete teilen sich dieselbe physische Test-DB, parallele
Paketausführung (Go-Testdefault) verursachte Querschläger zwischen
`internal/jobqueue` und `internal/migrate`.
4. **`internal/jobqueue`s Test-Fixture räumte nicht auf**: `TRUNCATE` statt
`DROP TABLE` ließ die Tabelle `processing_jobs` stehen, wodurch
`internal/migrate`s eigene, versionierte Migration mit
`relation already exists` scheiterte. Behoben durch `DROP TABLE IF EXISTS`
im Test-Cleanup.
## Build/Test-Ergebnis (192.168.1.131, `make check`)
```
go build ./... -> clean
scripts/reset-test-env.sh -> dms_tenant_test leer neu angelegt
go vet ./... -> clean
golangci-lint run ./... -> 0 issues
go test ./... -p 1 -count=1 -> 4/4 Pakete ok, 0 Fehlschläge (inkl. 8 jobqueue-Tests, 3 migrate-Tests, 6 storage-Tests real gegen MinIO)
```
## Gesamtergebnis
**Bestanden.** Alle drei Akzeptanzkriterien und alle drei Pflichtprüfungen
erfüllt. Vier reale Fehler beim Testen gefunden und behoben (zwei
Produktionscode-Bugs im SQL-Parameterhandling, zwei
Testinfrastruktur-Bugs) — bestätigt erneut den Wert, jede Prüfung
tatsächlich auf einem echten Testhost auszuführen statt nur zu behaupten.
-57
View File
@@ -1,57 +0,0 @@
# FDN-09 Prüfprotokoll: Verschlüsselung at rest & Schlüsselverwaltung
Welle 2. Voraussetzung: FDN-01 (Status "Fertig"), Core API-10 (Status
"Fertig").
## Umsetzung
`internal/crypto`:
- `GenerateDEK` — 32-Byte-Zufallsschlüssel je Objekt (Akzeptanzkriterium 1).
- `EncryptStream`/`DecryptStream` — AES-256-GCM auf dem Objektinhalt.
- `WrapDEK`/`UnwrapDEK` — Envelope-Verpackung des DEK mit dem Tenant-KEK.
- `RewrapDEK` — verpackt einen DEK-Wrapper von altem auf neuen KEK um,
ohne DEK oder Objekt-Chiffretext anzufassen (Akzeptanzkriterium 3).
- `HTTPKEKProvider` — bezieht den Tenant-KEK über Cores
`internal/kek.Handler.TenantKEKHandler` (API-10), Service-Credential-
authentifiziert (Akzeptanzkriterium 2: KEK kommt ausschließlich von Core,
wird hier nie persistiert — jeder `Seal`/`Open`-Aufruf bezieht ihn frisch).
- `Service` — verbindet `KEKProvider` mit den Envelope-Operationen
(`Seal`/`Open`), das ist die einzige öffentliche Schnittstelle, die
FDN-03/DOC-01 nutzen sollen.
## Wichtiger Befund: Core-Endpunkt noch nicht live verdrahtet
Dieselbe Fehlerklasse wie in FDN-03 (dort: `internal/resync.Handler`):
`internal/kek.Handler` (inkl. `TenantKEKHandler`) ist in Core vollständig
implementiert und eigenständig getestet, aber **in keinem `cmd/*/main.go`
registriert** (per `grep` bestätigt, Stand 2026-08-29). `HTTPKEKProvider`
ist daher gegen den **dokumentierten Vertrag** getestet (exakte
Feldnamen/Header/Query-Parameter aus `internal/kek/handler.go` gelesen),
nicht gegen eine echte laufende Core-Instanz. Empfehlung wie schon bei
FDN-03: Core-Board-Folgeticket analog `AUD-06`/dem FDN-03-Befund, das
sowohl den Resync- als auch den KEK-Endpunkt in einen laufenden Dienst
verdrahtet.
## Prüfungen
| # | Prüfung | Ergebnis |
|---|---|---|
| 1 | Round-Trip Encrypt/Decrypt liefert identischen Klartext | **bestanden**`TestEncryptDecryptStream_RoundTrip`, `TestService_SealOpen_RoundTrip` |
| 2 | Manipulierter Chiffretext wird bei Decrypt erkannt und abgelehnt (GCM-Auth-Tag) | **bestanden**`TestDecryptStream_RejectsTamperedCiphertext` (letztes Byte gekippt) UND `TestDecryptStream_WrongKeyRejected` (falscher Schlüssel, zweite mögliche Fehlerursache) |
| 3 | KEK-Rotation getestet, alte Objekte weiterhin lesbar | **bestanden**`TestRewrapDEK_RotationKeepsObjectReadable`: Objekt vor Rotation verschlüsselt, `RewrapDEK` von altem auf neuen KEK, Chiffretext bleibt UNVERÄNDERT, alter KEK kann neuen Wrapper nicht mehr entpacken (Rotation wirksam), neuer KEK + neuer Wrapper entschlüsseln das unveränderte, alte Objekt korrekt |
## 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 -> 5/5 Pakete ok, 0 Fehlschläge (10 neue crypto-Tests)
```
## Gesamtergebnis
**Bestanden**, mit derselben dokumentierten Core-seitigen Abhängigkeit wie
FDN-03 (Abschnitt "Wichtiger Befund"). Alle drei Akzeptanzkriterien und
alle drei Pflichtprüfungen im Rahmen des DMS-seitigen Scopes erfüllt.
-36
View File
@@ -1,36 +0,0 @@
module gitea.perlbach24.de/scripte/nexarch/dms
go 1.24
toolchain go1.24.4
require (
github.com/aws/aws-sdk-go-v2 v1.45.1
github.com/aws/aws-sdk-go-v2/config v1.33.1
github.com/aws/aws-sdk-go-v2/credentials v1.20.1
github.com/aws/aws-sdk-go-v2/service/s3 v1.109.1
github.com/aws/smithy-go v1.28.1
github.com/jackc/pgx/v5 v5.6.0
)
require (
github.com/aws/aws-sdk-go-v2/aws/protocol/eventstream v1.7.20 // indirect
github.com/aws/aws-sdk-go-v2/feature/ec2/imds v1.19.1 // indirect
github.com/aws/aws-sdk-go-v2/internal/configsources v1.5.1 // indirect
github.com/aws/aws-sdk-go-v2/internal/endpoints/v2 v2.8.1 // indirect
github.com/aws/aws-sdk-go-v2/internal/v4a v1.5.1 // indirect
github.com/aws/aws-sdk-go-v2/service/internal/accept-encoding v1.13.19 // indirect
github.com/aws/aws-sdk-go-v2/service/internal/checksum v1.11.1 // indirect
github.com/aws/aws-sdk-go-v2/service/internal/presigned-url v1.14.1 // indirect
github.com/aws/aws-sdk-go-v2/service/internal/s3shared v1.20.1 // indirect
github.com/aws/aws-sdk-go-v2/service/signin v1.7.1 // indirect
github.com/aws/aws-sdk-go-v2/service/sso v1.35.1 // indirect
github.com/aws/aws-sdk-go-v2/service/ssooidc v1.40.1 // indirect
github.com/aws/aws-sdk-go-v2/service/sts v1.47.1 // indirect
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
)
-64
View File
@@ -1,64 +0,0 @@
github.com/aws/aws-sdk-go-v2 v1.45.1 h1:iIoG3NaLhV6UZpPXyPXlDj2I9oS8tV/nMcMnITCC6Ks=
github.com/aws/aws-sdk-go-v2 v1.45.1/go.mod h1:bttEH6JqnUL8LepvDVfdrds/fZ5bCIxzpe3abyUrhDU=
github.com/aws/aws-sdk-go-v2/aws/protocol/eventstream v1.7.20 h1:GPRlPwz40I2B2VrBEASOA3Bi77NyeqejNLkifosX0rs=
github.com/aws/aws-sdk-go-v2/aws/protocol/eventstream v1.7.20/go.mod h1:g7PNzKcsOKWb4fkSRBA7BZVAS6Y8IcxzN+nRohhQ1Q8=
github.com/aws/aws-sdk-go-v2/config v1.33.1 h1:bq9jze1hQ5YTCLoVxNnbp0T7rglrlOE7N9YsHqjGkEw=
github.com/aws/aws-sdk-go-v2/config v1.33.1/go.mod h1:2A3HQwG4zaL5Tm80rc6RZj8LmWWv4WYT5v8raSz/L7A=
github.com/aws/aws-sdk-go-v2/credentials v1.20.1 h1:Z8GRNEx0u9sDkZOq4PUnN8mjGwbUQGRzMSXpvt3d8xQ=
github.com/aws/aws-sdk-go-v2/credentials v1.20.1/go.mod h1:uBIK00kFo95dnemqfFMTWx0X8YRqsh6ecIoCjjOkZqM=
github.com/aws/aws-sdk-go-v2/feature/ec2/imds v1.19.1 h1:YIEBqcqRnpi4Pfv0YHImtgi6czGCwKHANC7SwmUAVD0=
github.com/aws/aws-sdk-go-v2/feature/ec2/imds v1.19.1/go.mod h1:imEf0oufgAo8KAkCHhrOdqGEC0YWx1PPBQH82shSxGw=
github.com/aws/aws-sdk-go-v2/internal/configsources v1.5.1 h1:pc138gM1CW+XPc60rEwUlwwuwWFQK16CI1T7v1F9Oec=
github.com/aws/aws-sdk-go-v2/internal/configsources v1.5.1/go.mod h1:1+koxpPIbfBdfzP6vojm5/zTpTQ/micYwlxIiNB3TxI=
github.com/aws/aws-sdk-go-v2/internal/endpoints/v2 v2.8.1 h1:K0JsbZQj+1h208Ro1zHeA4l7bMp0NvRffHQ91q8Ol1s=
github.com/aws/aws-sdk-go-v2/internal/endpoints/v2 v2.8.1/go.mod h1:W3/vL6EtCIatICGy9ab29QhMuae+cOKPWcMxv02CO+Q=
github.com/aws/aws-sdk-go-v2/internal/v4a v1.5.1 h1:yhw5KD1phVyP9vijxOUzDfEtJx+bt+L63k+VfuiYFAA=
github.com/aws/aws-sdk-go-v2/internal/v4a v1.5.1/go.mod h1:ZW2e0d7DYlRxlS9hEiMXE47gTdX5KRN4byUiNbUpG+Q=
github.com/aws/aws-sdk-go-v2/service/internal/accept-encoding v1.13.19 h1:bAdDl/HkGCcGPoe25ToSHEw23VIxt6CT5fLcg111BKg=
github.com/aws/aws-sdk-go-v2/service/internal/accept-encoding v1.13.19/go.mod h1:KaUzbLxv4CeSxh6ZCl9B4m7CuFenS8kUEaDs+f/DQr4=
github.com/aws/aws-sdk-go-v2/service/internal/checksum v1.11.1 h1:s67hBfG5t9rn1NCvDuB4E3QIep3UFhHPtaIqFDjV3N8=
github.com/aws/aws-sdk-go-v2/service/internal/checksum v1.11.1/go.mod h1:FpvjBMXtSNMLPmDJsWwcY5cRnqJlpS2y1R6n4pvzs4k=
github.com/aws/aws-sdk-go-v2/service/internal/presigned-url v1.14.1 h1:RmmWQPREQdk9U+PfqeHW3MqZaBaNK7TpV9W3RY+b+7g=
github.com/aws/aws-sdk-go-v2/service/internal/presigned-url v1.14.1/go.mod h1:0A3W4F+68ZnNk5XcNL/e9HFMwnP8RlEicFfy6eOEDyw=
github.com/aws/aws-sdk-go-v2/service/internal/s3shared v1.20.1 h1:ZMbtPZZQRca+3+XYQne9PBvRiYpHZlNJJOZfE9WNfT0=
github.com/aws/aws-sdk-go-v2/service/internal/s3shared v1.20.1/go.mod h1:YAGWQdCYlVCoqrzvfv3RLxO6zKwti7gsAULOGWPLYv4=
github.com/aws/aws-sdk-go-v2/service/s3 v1.109.1 h1:kVpzaDBzOdRtOftmiSpTdQbWVqRg0kONLXijktiwXnk=
github.com/aws/aws-sdk-go-v2/service/s3 v1.109.1/go.mod h1:CUr46sCpGAg/rHaclRyhJX0LJAmH73uWSJPPSaMUrSk=
github.com/aws/aws-sdk-go-v2/service/signin v1.7.1 h1:mdMtSVKdQ3+mzBh+l0ogrFYZVQUCg6pJZOirA2ARsYE=
github.com/aws/aws-sdk-go-v2/service/signin v1.7.1/go.mod h1:9IqUlsJDbUPcg6cgx3WEzXdjrbWzLDQrak0aaSqlTcI=
github.com/aws/aws-sdk-go-v2/service/sso v1.35.1 h1:B6WFn91tobD6gG4724ONHaqrpKsoETGnv98LHe/yIGM=
github.com/aws/aws-sdk-go-v2/service/sso v1.35.1/go.mod h1:tWuiVBUtPBr8/rgRiYS8Uf85sHcAN+G7XS3D3CEoUh8=
github.com/aws/aws-sdk-go-v2/service/ssooidc v1.40.1 h1:6yeYCWFvgbI2TI3K6jr9LtBNhXgJ7g4xqD+DEiaDDmM=
github.com/aws/aws-sdk-go-v2/service/ssooidc v1.40.1/go.mod h1:naFe83jSMuYkH+QjQPX8n1MLhBkeCFM5Lsnh5m5wz3c=
github.com/aws/aws-sdk-go-v2/service/sts v1.47.1 h1:Sv2xPnRHlThSUtVujYuUBPI/Il8si6UPHXL8DMiB/F0=
github.com/aws/aws-sdk-go-v2/service/sts v1.47.1/go.mod h1:mKo/CzaCz8qytGW70NG4vIIGAx1HXTlb5lHNkC5k3lk=
github.com/aws/smithy-go v1.28.1 h1:R/nXH00c8qcfCzQVELtRw+eLQWtzv+VAIEFJ1/xxXlQ=
github.com/aws/smithy-go v1.28.1/go.mod h1:YE2RhdIuDbA5E5bTdciG9KrW3+TiEONeUWCqxX9i1Fc=
github.com/davecgh/go-spew v1.1.0/go.mod h1:J7Y8YcW2NihsgmVo/mv3lAwl/skON4iLHjSsI+c5H38=
github.com/davecgh/go-spew v1.1.1 h1:vj9j/u1bqnvCEfJOwUhtlOARqs3+rkHYY13jYWTU97c=
github.com/davecgh/go-spew v1.1.1/go.mod h1:J7Y8YcW2NihsgmVo/mv3lAwl/skON4iLHjSsI+c5H38=
github.com/jackc/pgpassfile v1.0.0 h1:/6Hmqy13Ss2zCq62VdNG8tM1wchn8zjSGOBJ6icpsIM=
github.com/jackc/pgpassfile v1.0.0/go.mod h1:CEx0iS5ambNFdcRtxPj5JhEz+xB6uRky5eyVu/W2HEg=
github.com/jackc/pgservicefile v0.0.0-20221227161230-091c0ba34f0a h1:bbPeKD0xmW/Y25WS6cokEszi5g+S0QxI/d45PkRi7Nk=
github.com/jackc/pgservicefile v0.0.0-20221227161230-091c0ba34f0a/go.mod h1:5TJZWKEWniPve33vlWYSoGYefn3gLQRzjfDlhSJ9ZKM=
github.com/jackc/pgx/v5 v5.6.0 h1:SWJzexBzPL5jb0GEsrPMLIsi/3jOo7RHlzTjcAeDrPY=
github.com/jackc/pgx/v5 v5.6.0/go.mod h1:DNZ/vlrUnhWCoFGxHAG8U2ljioxukquj7utPDgtQdTw=
github.com/jackc/puddle/v2 v2.2.1 h1:RhxXJtFG022u4ibrCSMSiu5aOq1i77R3OHKNJj77OAk=
github.com/jackc/puddle/v2 v2.2.1/go.mod h1:vriiEXHvEE654aYKXXjOvZM39qJ0q+azkZFrfEOc3H4=
github.com/pmezard/go-difflib v1.0.0 h1:4DBwDE0NGyQoBHbLQYPwSUPoCMWR5BEzIk/f1lZbAQM=
github.com/pmezard/go-difflib v1.0.0/go.mod h1:iKH77koFhYxTK1pcRnkKkqfTogsbg7gZNVY4sRDYZ/4=
github.com/stretchr/objx v0.1.0/go.mod h1:HFkY916IF+rwdDfMAkV7OtwuqBVzrE8GR6GFx+wExME=
github.com/stretchr/testify v1.3.0/go.mod h1:M5WIy9Dh21IEIfnGCwXGc5bZfKNJtfHm1UVUgZn+9EI=
github.com/stretchr/testify v1.7.0/go.mod h1:6Fq8oRcR53rry900zMqJjRRixrwX3KX962/h/Wwjteg=
github.com/stretchr/testify v1.8.1 h1:w7B6lhMri9wdJUVmEZPGGhZzrYTPvgJArz7wNPgYKsk=
github.com/stretchr/testify v1.8.1/go.mod h1:w2LPCIKwWwSfY2zedu0+kehJoqGctiVI29o6fzry7u4=
golang.org/x/crypto v0.17.0 h1:r8bRNjWL3GshPW3gkd+RpvzWrZAwPS49OmTGZ/uhM4k=
golang.org/x/crypto v0.17.0/go.mod h1:gCAAfMLgwOJRpTjQ2zCCt2OcSfYMTeZVSRtQlPC7Nq4=
golang.org/x/sync v0.1.0 h1:wsuoTGHzEhffawBOhz5CYhcrV4IdKZbEyZjBMuTp12o=
golang.org/x/sync v0.1.0/go.mod h1:RxMgew5VJxzue5/jJTE5uejpjVlOe/izrB70Jof72aM=
golang.org/x/text v0.14.0 h1:ScX5w1eTa3QqT8oi6+ziP7dTV1S2+ALU0bI+0zXKWiQ=
golang.org/x/text v0.14.0/go.mod h1:18ZOQIKpY8NJVqYksKHtTdi31H5itFRjB5/qKTNYzSU=
gopkg.in/check.v1 v0.0.0-20161208181325-20d25e280405/go.mod h1:Co6ibVJAznAaIkqp8huTwlJQCZ016jof/cbN4VW5Yz0=
gopkg.in/yaml.v3 v3.0.0-20200313102051-9f266ea9e77c/go.mod h1:K4uyk7z7BCEPqu6E+C64Yfv1cQ7kz7rIZviUmN+EgEM=
gopkg.in/yaml.v3 v3.0.1 h1:fxVm/GzAzEWqLHuvctI91KS9hhNmmWOoWu0XTYJS7CA=
gopkg.in/yaml.v3 v3.0.1/go.mod h1:K4uyk7z7BCEPqu6E+C64Yfv1cQ7kz7rIZviUmN+EgEM=
-140
View File
@@ -1,140 +0,0 @@
// Package crypto implementiert FDN-09: Envelope-Encryption fuer Objekte
// at rest. Jedes Objekt bekommt einen eigenen, zufaelligen
// Datenverschluesselungsschluessel (DEK, Akzeptanzkriterium 1), der mit
// dem Tenant-Hauptschluessel (KEK) verpackt wird — der KEK selbst kommt
// AUSSCHLIESSLICH von Core API-10 (Akzeptanzkriterium 2), wird hier nie
// persistiert, nur fluechtig fuer eine Wrap-/Unwrap-Operation gehalten.
package crypto
import (
"bytes"
"crypto/aes"
"crypto/cipher"
"crypto/rand"
"errors"
"fmt"
"io"
)
// DEKSize/KEKSize sind die geforderten Laengen fuer AES-256-GCM.
const (
DEKSize = 32
KEKSize = 32
)
// ErrDecryptFailed wird geliefert, wenn ein Chiffretext nicht entschluesselt
// werden kann — falscher Schluessel ODER manipulierte Daten (Pruefung 2:
// GCM-Auth-Tag erkennt Manipulation zuverlaessig, AEAD unterscheidet die
// beiden Ursachen bewusst nicht, um keine Seitenkanal-Information ueber
// "welcher Teil" falsch war preiszugeben).
var ErrDecryptFailed = errors.New("crypto: entschluesselung fehlgeschlagen (falscher schluessel oder manipulierte daten)")
// GenerateDEK erzeugt einen neuen, zufaelligen Datenverschluesselungs-
// schluessel — fuer JEDES Objekt neu, nie wiederverwendet
// (Akzeptanzkriterium 1).
func GenerateDEK() ([]byte, error) {
dek := make([]byte, DEKSize)
if _, err := rand.Read(dek); err != nil {
return nil, fmt.Errorf("crypto: dek erzeugen: %w", err)
}
return dek, nil
}
func seal(key, plaintext []byte) ([]byte, error) {
block, err := aes.NewCipher(key)
if err != nil {
return nil, fmt.Errorf("crypto: aes-cipher erstellen: %w", err)
}
gcm, err := cipher.NewGCM(block)
if err != nil {
return nil, fmt.Errorf("crypto: gcm erstellen: %w", err)
}
nonce := make([]byte, gcm.NonceSize())
if _, err := rand.Read(nonce); err != nil {
return nil, fmt.Errorf("crypto: nonce erzeugen: %w", err)
}
return gcm.Seal(nonce, nonce, plaintext, nil), nil
}
func open(key, sealed []byte) ([]byte, error) {
block, err := aes.NewCipher(key)
if err != nil {
return nil, fmt.Errorf("crypto: aes-cipher erstellen: %w", err)
}
gcm, err := cipher.NewGCM(block)
if err != nil {
return nil, fmt.Errorf("crypto: gcm erstellen: %w", err)
}
if len(sealed) < gcm.NonceSize() {
return nil, ErrDecryptFailed
}
nonce, ciphertext := sealed[:gcm.NonceSize()], sealed[gcm.NonceSize():]
plaintext, err := gcm.Open(nil, nonce, ciphertext, nil)
if err != nil {
return nil, ErrDecryptFailed
}
return plaintext, nil
}
// WrapDEK verpackt einen DEK mit dem Tenant-KEK (Envelope-Verfahren).
func WrapDEK(kek, dek []byte) ([]byte, error) {
wrapped, err := seal(kek, dek)
if err != nil {
return nil, fmt.Errorf("crypto: dek verpacken: %w", err)
}
return wrapped, nil
}
// UnwrapDEK entpackt einen zuvor mit WrapDEK verpackten DEK.
func UnwrapDEK(kek, wrappedDEK []byte) ([]byte, error) {
return open(kek, wrappedDEK)
}
// RewrapDEK verpackt einen bereits vorhandenen DEK-Wrapper von oldKEK auf
// newKEK um, OHNE den DEK selbst oder den Objekt-Chiffretext zu veraendern
// (Akzeptanzkriterium 3: KEK-Rotation ohne Neuverschluesselung aller
// Objekte — nur der DEK-Wrapper wird neu verpackt).
func RewrapDEK(oldKEK, newKEK, wrappedDEK []byte) ([]byte, error) {
dek, err := UnwrapDEK(oldKEK, wrappedDEK)
if err != nil {
return nil, fmt.Errorf("crypto: dek mit altem kek entpacken: %w", err)
}
rewrapped, err := WrapDEK(newKEK, dek)
if err != nil {
return nil, fmt.Errorf("crypto: dek mit neuem kek verpacken: %w", err)
}
return rewrapped, nil
}
// EncryptStream verschluesselt den gesamten Inhalt von r mit dek
// (AES-256-GCM) und liefert den Chiffretext als Stream. Liest r vollstaendig
// in den Speicher (dasselbe Muster wie internal/storage.S3Driver.Put aus
// FDN-03, das S3-PutObject ebenfalls vollstaendig puffert) — fuer sehr
// grosse Dateien waere ein segmentiertes AEAD-Verfahren noetig, das ist
// bewusst nicht Teil der "kleinsten Loesung" dieser Kachel.
func EncryptStream(dek []byte, r io.Reader) (io.Reader, error) {
plaintext, err := io.ReadAll(r)
if err != nil {
return nil, fmt.Errorf("crypto: klartext lesen: %w", err)
}
ciphertext, err := seal(dek, plaintext)
if err != nil {
return nil, fmt.Errorf("crypto: verschluesseln: %w", err)
}
return bytes.NewReader(ciphertext), nil
}
// DecryptStream entschluesselt einen zuvor mit EncryptStream erzeugten
// Chiffretext-Stream. Liefert ErrDecryptFailed bei manipuliertem Chiffretext
// (Pruefung 2, GCM-Auth-Tag).
func DecryptStream(dek []byte, r io.Reader) (io.Reader, error) {
ciphertext, err := io.ReadAll(r)
if err != nil {
return nil, fmt.Errorf("crypto: chiffretext lesen: %w", err)
}
plaintext, err := open(dek, ciphertext)
if err != nil {
return nil, err
}
return bytes.NewReader(plaintext), nil
}
-187
View File
@@ -1,187 +0,0 @@
package crypto
import (
"bytes"
"errors"
"io"
"testing"
)
// TestGenerateDEK_NeverReused ist Akzeptanzkriterium 1: DEK wird pro
// Objekt neu erzeugt, nie wiederverwendet.
func TestGenerateDEK_NeverReused(t *testing.T) {
seen := map[string]bool{}
for i := 0; i < 50; i++ {
dek, err := GenerateDEK()
if err != nil {
t.Fatalf("generatedek: %v", err)
}
if len(dek) != DEKSize {
t.Fatalf("dek-laenge = %d, want %d", len(dek), DEKSize)
}
key := string(dek)
if seen[key] {
t.Fatal("generatedek lieferte denselben schluessel zweimal")
}
seen[key] = true
}
}
// TestEncryptDecryptStream_RoundTrip ist Pruefung 1: Round-Trip liefert
// identischen Klartext.
func TestEncryptDecryptStream_RoundTrip(t *testing.T) {
dek, err := GenerateDEK()
if err != nil {
t.Fatalf("generatedek: %v", err)
}
plaintext := []byte("Rechnung 2026-0001 — vertraulicher Inhalt")
ciphertext, err := EncryptStream(dek, bytes.NewReader(plaintext))
if err != nil {
t.Fatalf("encryptstream: %v", err)
}
ciphertextBytes, err := io.ReadAll(ciphertext)
if err != nil {
t.Fatalf("chiffretext lesen: %v", err)
}
if bytes.Equal(ciphertextBytes, plaintext) {
t.Fatal("chiffretext ist identisch zum klartext - keine verschluesselung stattgefunden")
}
decrypted, err := DecryptStream(dek, bytes.NewReader(ciphertextBytes))
if err != nil {
t.Fatalf("decryptstream: %v", err)
}
got, err := io.ReadAll(decrypted)
if err != nil {
t.Fatalf("klartext lesen: %v", err)
}
if !bytes.Equal(got, plaintext) {
t.Fatalf("entschluesselter klartext = %q, want %q", got, plaintext)
}
}
// TestDecryptStream_RejectsTamperedCiphertext ist Pruefung 2: manipulierter
// Chiffretext wird bei Decrypt erkannt und abgelehnt (GCM-Auth-Tag).
func TestDecryptStream_RejectsTamperedCiphertext(t *testing.T) {
dek, err := GenerateDEK()
if err != nil {
t.Fatalf("generatedek: %v", err)
}
ciphertext, err := EncryptStream(dek, bytes.NewReader([]byte("geheimer inhalt")))
if err != nil {
t.Fatalf("encryptstream: %v", err)
}
tampered, err := io.ReadAll(ciphertext)
if err != nil {
t.Fatalf("chiffretext lesen: %v", err)
}
// Ein Byte in der Mitte des Chiffretexts kippen (nach dem Nonce-Praefix).
tampered[len(tampered)-1] ^= 0xFF
if _, err := DecryptStream(dek, bytes.NewReader(tampered)); !errors.Is(err, ErrDecryptFailed) {
t.Fatalf("manipulierter chiffretext: erwartet ErrDecryptFailed, habe %v", err)
}
}
// TestDecryptStream_WrongKeyRejected prueft den zweiten Fall, in dem
// Entschluesselung scheitern muss: falscher Schluessel statt Manipulation.
func TestDecryptStream_WrongKeyRejected(t *testing.T) {
dek1, _ := GenerateDEK()
dek2, _ := GenerateDEK()
ciphertext, err := EncryptStream(dek1, bytes.NewReader([]byte("inhalt")))
if err != nil {
t.Fatalf("encryptstream: %v", err)
}
ciphertextBytes, _ := io.ReadAll(ciphertext)
if _, err := DecryptStream(dek2, bytes.NewReader(ciphertextBytes)); !errors.Is(err, ErrDecryptFailed) {
t.Fatalf("falscher schluessel: erwartet ErrDecryptFailed, habe %v", err)
}
}
// TestWrapUnwrapDEK_RoundTrip prueft das Envelope-Verpacken des DEK selbst.
func TestWrapUnwrapDEK_RoundTrip(t *testing.T) {
kek := make([]byte, KEKSize)
for i := range kek {
kek[i] = byte(i)
}
dek, err := GenerateDEK()
if err != nil {
t.Fatalf("generatedek: %v", err)
}
wrapped, err := WrapDEK(kek, dek)
if err != nil {
t.Fatalf("wrapdek: %v", err)
}
if bytes.Equal(wrapped, dek) {
t.Fatal("wrapdek lieferte den unveraenderten dek")
}
unwrapped, err := UnwrapDEK(kek, wrapped)
if err != nil {
t.Fatalf("unwrapdek: %v", err)
}
if !bytes.Equal(unwrapped, dek) {
t.Fatal("entpackter dek stimmt nicht mit original ueberein")
}
}
// TestRewrapDEK_RotationKeepsObjectReadable ist Pruefung 3: KEK-Rotation
// getestet, alte Objekte weiterhin lesbar — OHNE dass der Chiffretext des
// Objekts selbst jemals angefasst wird (Akzeptanzkriterium 3).
func TestRewrapDEK_RotationKeepsObjectReadable(t *testing.T) {
oldKEK := bytes.Repeat([]byte{0x01}, KEKSize)
newKEK := bytes.Repeat([]byte{0x02}, KEKSize)
dek, err := GenerateDEK()
if err != nil {
t.Fatalf("generatedek: %v", err)
}
plaintext := []byte("altes objekt, vor der rotation verschluesselt")
ciphertext, err := EncryptStream(dek, bytes.NewReader(plaintext))
if err != nil {
t.Fatalf("encryptstream: %v", err)
}
ciphertextBytes, _ := io.ReadAll(ciphertext)
wrappedOld, err := WrapDEK(oldKEK, dek)
if err != nil {
t.Fatalf("wrapdek (alt): %v", err)
}
// Rotation: NUR der DEK-Wrapper wird neu verpackt, der Chiffretext
// (ciphertextBytes) wird nicht angefasst.
wrappedNew, err := RewrapDEK(oldKEK, newKEK, wrappedOld)
if err != nil {
t.Fatalf("rewrapdek: %v", err)
}
if bytes.Equal(wrappedNew, wrappedOld) {
t.Fatal("rewrapdek lieferte denselben wrapper wie vor der rotation")
}
// Der alte KEK kann den NEUEN Wrapper nicht mehr entpacken (Rotation
// war wirksam).
if _, err := UnwrapDEK(oldKEK, wrappedNew); !errors.Is(err, ErrDecryptFailed) {
t.Fatal("alter kek konnte den nach der rotation neu verpackten dek weiterhin entpacken")
}
// Mit dem NEUEN KEK und dem neuen Wrapper ist das unveraenderte, alte
// Objekt weiterhin vollstaendig lesbar.
dekAfterRotation, err := UnwrapDEK(newKEK, wrappedNew)
if err != nil {
t.Fatalf("unwrapdek (neu) nach rotation: %v", err)
}
decrypted, err := DecryptStream(dekAfterRotation, bytes.NewReader(ciphertextBytes))
if err != nil {
t.Fatalf("decryptstream nach rotation: %v", err)
}
got, err := io.ReadAll(decrypted)
if err != nil {
t.Fatalf("klartext lesen: %v", err)
}
if !bytes.Equal(got, plaintext) {
t.Fatalf("altes objekt nach rotation = %q, want %q", got, plaintext)
}
}
-80
View File
@@ -1,80 +0,0 @@
package crypto
import (
"context"
"encoding/base64"
"encoding/json"
"fmt"
"net/http"
"net/url"
)
// KEKProvider liefert den aktuellen, entschluesselten Tenant-KEK
// (Akzeptanzkriterium 2: KEK kommt ausschliesslich von Core API-10).
// Schmale Schnittstelle, damit Tests einen Fake statt eines echten
// HTTP-Aufrufs einsetzen koennen.
type KEKProvider interface {
TenantKEK(ctx context.Context, tenantSlug string) ([]byte, error)
}
// tenantKEKResponse entspricht Core internal/kek.tenantKEKResponse
// (JSON-Vertrag: tenant_kek_base64) — dieselbe Struktur, hier gespiegelt,
// da DMS Cores internal/-Pakete nicht importieren kann.
type tenantKEKResponse struct {
TenantKEKBase64 string `json:"tenant_kek_base64"`
}
// HTTPKEKProvider bezieht den Tenant-KEK ueber Cores KEK-Handler
// (internal/kek.Handler.TenantKEKHandler, API-10), authentifiziert ueber
// dasselbe Service-Credential-Verfahren wie jeder andere Modul-Core-Aufruf
// (API-02) — identisches Muster wie internal/storage.HTTPUsageReporter aus
// FDN-03.
type HTTPKEKProvider struct {
endpointURL string
clientID string
clientSecret string
httpClient *http.Client
}
func NewHTTPKEKProvider(endpointURL, clientID, clientSecret string, httpClient *http.Client) *HTTPKEKProvider {
if httpClient == nil {
httpClient = http.DefaultClient
}
return &HTTPKEKProvider{endpointURL: endpointURL, clientID: clientID, clientSecret: clientSecret, httpClient: httpClient}
}
func (p *HTTPKEKProvider) TenantKEK(ctx context.Context, tenantSlug string) ([]byte, error) {
u, err := url.Parse(p.endpointURL)
if err != nil {
return nil, fmt.Errorf("crypto: kek-endpunkt-url ungueltig: %w", err)
}
q := u.Query()
q.Set("tenant", tenantSlug)
u.RawQuery = q.Encode()
req, err := http.NewRequestWithContext(ctx, http.MethodGet, u.String(), nil)
if err != nil {
return nil, fmt.Errorf("crypto: kek-anfrage aufbauen: %w", err)
}
req.Header.Set("X-Nexarch-Client-Id", p.clientID)
req.Header.Set("X-Nexarch-Client-Secret", p.clientSecret)
resp, err := p.httpClient.Do(req)
if err != nil {
return nil, fmt.Errorf("crypto: kek-anfrage senden: %w", err)
}
defer func() { _ = resp.Body.Close() }()
if resp.StatusCode != http.StatusOK {
return nil, fmt.Errorf("crypto: kek-bezug von core abgelehnt: status %d", resp.StatusCode)
}
var body tenantKEKResponse
if err := json.NewDecoder(resp.Body).Decode(&body); err != nil {
return nil, fmt.Errorf("crypto: kek-antwort dekodieren: %w", err)
}
kek, err := base64.StdEncoding.DecodeString(body.TenantKEKBase64)
if err != nil {
return nil, fmt.Errorf("crypto: kek base64-dekodieren: %w", err)
}
return kek, nil
}
-65
View File
@@ -1,65 +0,0 @@
package crypto
import (
"context"
"encoding/base64"
"encoding/json"
"net/http"
"net/http/httptest"
"testing"
)
// TestHTTPKEKProvider_SendsCorrectContractToCore ist der Nachweis, dass
// HTTPKEKProvider exakt den Vertrag von Core internal/kek.Handler.
// TenantKEKHandler bedient (Service-Credential-Header, tenant-Query-Param,
// JSON-Feldname) — echte Vernetzung gegen einen laufenden Core-Prozess ist
// nicht Teil dieses Tests (siehe FDN-09-Pruefprotokoll: derselbe
// "nicht verdrahtet"-Befund wie bei internal/resync in FDN-03), daher hier
// gegen einen httptest-Server geprueft, der denselben Vertrag nachbildet.
func TestHTTPKEKProvider_SendsCorrectContractToCore(t *testing.T) {
wantKEK := make([]byte, KEKSize)
for i := range wantKEK {
wantKEK[i] = byte(i)
}
var gotClientID, gotClientSecret, gotTenant string
srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
gotClientID = r.Header.Get("X-Nexarch-Client-Id")
gotClientSecret = r.Header.Get("X-Nexarch-Client-Secret")
gotTenant = r.URL.Query().Get("tenant")
w.Header().Set("Content-Type", "application/json")
_ = json.NewEncoder(w).Encode(tenantKEKResponse{TenantKEKBase64: base64.StdEncoding.EncodeToString(wantKEK)})
}))
defer srv.Close()
provider := NewHTTPKEKProvider(srv.URL, "dms-service-client", "dms-service-secret", nil)
got, err := provider.TenantKEK(context.Background(), "acme")
if err != nil {
t.Fatalf("tenantkek: %v", err)
}
if gotClientID != "dms-service-client" || gotClientSecret != "dms-service-secret" {
t.Fatalf("service-credential-header falsch: id=%q secret=%q", gotClientID, gotClientSecret)
}
if gotTenant != "acme" {
t.Fatalf("tenant-query-parameter = %q, want %q", gotTenant, "acme")
}
if len(got) != KEKSize || got[0] != wantKEK[0] {
t.Fatalf("erhaltener kek stimmt nicht mit erwartetem ueberein")
}
}
// TestHTTPKEKProvider_RejectsNonOKStatus prueft, dass ein von Core
// abgelehnter Aufruf (z.B. Modul nicht fuer diesen Mandanten aktiv,
// ErrForbidden auf Core-Seite) als Fehler zurueckgegeben wird.
func TestHTTPKEKProvider_RejectsNonOKStatus(t *testing.T) {
srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
http.Error(w, "zugriff auf diesen mandanten verweigert", http.StatusForbidden)
}))
defer srv.Close()
provider := NewHTTPKEKProvider(srv.URL, "x", "y", nil)
if _, err := provider.TenantKEK(context.Background(), "fremder-tenant"); err == nil {
t.Fatal("erwartet fehler bei abgelehntem kek-bezug, habe nil")
}
}
-65
View File
@@ -1,65 +0,0 @@
package crypto
import (
"context"
"fmt"
"io"
)
// Envelope ist das Ergebnis einer Seal-Operation: der Chiffretext-Stream
// plus der mit dem Tenant-KEK verpackte DEK, der zusammen mit dem Objekt
// persistiert werden muss (z.B. in file_revisions, FDN-02 — die konkrete
// Spalte ist nicht Bestandteil dieser Kachel, siehe "Nicht Bestandteil").
type Envelope struct {
Ciphertext io.Reader
WrappedDEK []byte
}
// Service verbindet KEKProvider mit den Envelope-Operationen — die
// Storage-Abstraktion (FDN-03) ruft ausschliesslich Service auf, nie die
// Einzelfunktionen aus envelope.go direkt.
type Service struct {
kek KEKProvider
}
func NewService(kek KEKProvider) *Service {
return &Service{kek: kek}
}
// Seal erzeugt einen neuen DEK (Akzeptanzkriterium 1), verschluesselt
// plaintext damit und verpackt den DEK mit dem aktuellen Tenant-KEK
// (Akzeptanzkriterium 2 — der KEK wird bei JEDEM Aufruf frisch von Core
// bezogen, nie zwischengespeichert).
func (s *Service) Seal(ctx context.Context, tenantSlug string, plaintext io.Reader) (*Envelope, error) {
dek, err := GenerateDEK()
if err != nil {
return nil, err
}
ciphertext, err := EncryptStream(dek, plaintext)
if err != nil {
return nil, err
}
kek, err := s.kek.TenantKEK(ctx, tenantSlug)
if err != nil {
return nil, fmt.Errorf("crypto: tenant-kek beziehen: %w", err)
}
wrappedDEK, err := WrapDEK(kek, dek)
if err != nil {
return nil, err
}
return &Envelope{Ciphertext: ciphertext, WrappedDEK: wrappedDEK}, nil
}
// Open entpackt den DEK mit dem aktuellen Tenant-KEK und entschluesselt
// ciphertext damit.
func (s *Service) Open(ctx context.Context, tenantSlug string, wrappedDEK []byte, ciphertext io.Reader) (io.Reader, error) {
kek, err := s.kek.TenantKEK(ctx, tenantSlug)
if err != nil {
return nil, fmt.Errorf("crypto: tenant-kek beziehen: %w", err)
}
dek, err := UnwrapDEK(kek, wrappedDEK)
if err != nil {
return nil, err
}
return DecryptStream(dek, ciphertext)
}
-77
View File
@@ -1,77 +0,0 @@
package crypto
import (
"bytes"
"context"
"io"
"testing"
)
// fakeKEKProvider liefert einen fest hinterlegten Tenant-KEK, ohne echten
// HTTP-Aufruf — fuer Tests, die Service isoliert vom Core-Kontrakt pruefen
// wollen (der Kontrakt selbst ist in kekprovider_test.go geprueft).
type fakeKEKProvider struct {
kek []byte
calls int
}
func (f *fakeKEKProvider) TenantKEK(ctx context.Context, tenantSlug string) ([]byte, error) {
f.calls++
return f.kek, nil
}
// TestService_SealOpen_RoundTrip ist der End-to-End-Nachweis fuer
// Pruefung 1 ueber Service statt die Einzelfunktionen.
func TestService_SealOpen_RoundTrip(t *testing.T) {
kekProvider := &fakeKEKProvider{kek: bytes.Repeat([]byte{0x07}, KEKSize)}
svc := NewService(kekProvider)
ctx := context.Background()
plaintext := []byte("vertraulicher dokumentinhalt")
env, err := svc.Seal(ctx, "acme", bytes.NewReader(plaintext))
if err != nil {
t.Fatalf("seal: %v", err)
}
ciphertextBytes, err := io.ReadAll(env.Ciphertext)
if err != nil {
t.Fatalf("chiffretext lesen: %v", err)
}
if bytes.Equal(ciphertextBytes, plaintext) {
t.Fatal("chiffretext ist identisch zum klartext")
}
if len(env.WrappedDEK) == 0 {
t.Fatal("erwartet nicht-leeren wrappeddek")
}
decrypted, err := svc.Open(ctx, "acme", env.WrappedDEK, bytes.NewReader(ciphertextBytes))
if err != nil {
t.Fatalf("open: %v", err)
}
got, err := io.ReadAll(decrypted)
if err != nil {
t.Fatalf("klartext lesen: %v", err)
}
if !bytes.Equal(got, plaintext) {
t.Fatalf("entschluesselter klartext = %q, want %q", got, plaintext)
}
}
// TestService_Seal_NeverPersistsKEK ist Akzeptanzkriterium 2 (struktureller
// Nachweis): Service haelt den KEK nicht ueber einen Aufruf hinaus fest —
// jeder Seal/Open-Aufruf bezieht ihn frisch ueber KEKProvider.
func TestService_Seal_NeverPersistsKEK(t *testing.T) {
kekProvider := &fakeKEKProvider{kek: bytes.Repeat([]byte{0x09}, KEKSize)}
svc := NewService(kekProvider)
ctx := context.Background()
if _, err := svc.Seal(ctx, "acme", bytes.NewReader([]byte("a"))); err != nil {
t.Fatalf("seal 1: %v", err)
}
if _, err := svc.Seal(ctx, "acme", bytes.NewReader([]byte("b"))); err != nil {
t.Fatalf("seal 2: %v", err)
}
if kekProvider.calls != 2 {
t.Fatalf("erwartet 2 kek-abrufe (einer je Seal-Aufruf, kein Caching), habe %d", kekProvider.calls)
}
}
-238
View File
@@ -1,238 +0,0 @@
// Package jobqueue implementiert FDN-04: eine Postgres-gestuetzte
// Job-Queue mit Wiederholungslogik, Backoff und Dead-Letter-Queue — kein
// Redis/AMQP (siehe Ticket-Vorgabe). FOR UPDATE SKIP LOCKED erlaubt
// mehrere gleichzeitige Worker-Goroutinen (auch mehrinstanzfaehig, da der
// Zustand ausschliesslich in Postgres liegt, keine In-Memory-Zaehler —
// dieselbe Konvention wie Core internal/lockout, siehe "Bekannte Fehler
// vermeiden" im Ticket).
package jobqueue
import (
"context"
"encoding/json"
"errors"
"fmt"
"time"
"github.com/jackc/pgx/v5"
"github.com/jackc/pgx/v5/pgxpool"
)
// Status-Werte spiegeln den CHECK-Constraint der Migration.
const (
StatusPending = "pending"
StatusProcessing = "processing"
StatusSucceeded = "succeeded"
StatusFailed = "failed"
StatusDeadLetter = "dead_letter"
)
// ErrNotFound wird geliefert, wenn ein angefragter Job nicht existiert.
var ErrNotFound = errors.New("jobqueue: job nicht gefunden")
// ErrNoJobAvailable wird von Dequeue geliefert, wenn aktuell kein
// abholbarer Job vorhanden ist (kein Fehlerzustand, sondern der Normalfall
// bei leerer Queue).
var ErrNoJobAvailable = errors.New("jobqueue: kein job verfuegbar")
// Job ist eine einzelne Aufgabe in der Queue.
type Job struct {
ID string
JobType string
Payload json.RawMessage
Status string
Attempts int
MaxAttempts int
LastError *string
}
// DefaultMaxAttempts/DefaultStaleLockAfter sind Standardwerte, ueberschreibbar
// je Enqueue-Aufruf (MaxAttempts) bzw. am Queue selbst (StaleLockAfter).
const (
DefaultMaxAttempts = 5
)
// Queue kapselt den Zugriff auf processing_jobs.
type Queue struct {
pool *pgxpool.Pool
staleLockAfter time.Duration
}
// NewQueue erzeugt eine Queue. staleLockAfter legt fest, ab wann ein
// Job, der als "processing" markiert ist, aber dessen Worker vermutlich
// abgestuerzt ist, wieder als abholbar gilt (Pruefung 1: Absturz fuehrt zu
// erneuter Zustellung) — kein Heartbeat-Mechanismus noetig, ein grosszuegiges
// Zeitfenster genuegt fuer die "kleinste Loesung".
func NewQueue(pool *pgxpool.Pool, staleLockAfter time.Duration) *Queue {
return &Queue{pool: pool, staleLockAfter: staleLockAfter}
}
// EnqueueOptions steuert optionale Einreih-Parameter.
type EnqueueOptions struct {
// IdempotencyKey verhindert doppelte Einreihung derselben logischen
// Aufgabe (Pruefung 2: Idempotenz bei Doppelzustellung) — leer bedeutet
// kein Dedup-Anspruch.
IdempotencyKey string
MaxAttempts int
}
// Enqueue reiht einen neuen Job ein (Akzeptanzkriterium 1). Bei gesetztem
// IdempotencyKey und bereits existierendem gleichen Key wird die ID des
// BEREITS vorhandenen Jobs zurueckgegeben, kein Duplikat angelegt.
func (q *Queue) Enqueue(ctx context.Context, jobType string, payload any, opts EnqueueOptions) (string, error) {
payloadJSON, err := json.Marshal(payload)
if err != nil {
return "", fmt.Errorf("jobqueue: payload serialisieren: %w", err)
}
maxAttempts := opts.MaxAttempts
if maxAttempts <= 0 {
maxAttempts = DefaultMaxAttempts
}
var idempotencyKey any
if opts.IdempotencyKey != "" {
idempotencyKey = opts.IdempotencyKey
}
var id string
err = q.pool.QueryRow(ctx, `
INSERT INTO processing_jobs (job_type, payload, idempotency_key, max_attempts)
VALUES ($1, $2, $3, $4)
ON CONFLICT (idempotency_key) DO UPDATE SET job_type = processing_jobs.job_type
RETURNING id
`, jobType, payloadJSON, idempotencyKey, maxAttempts).Scan(&id)
if err != nil {
return "", fmt.Errorf("jobqueue: job einreihen: %w", err)
}
return id, nil
}
// Dequeue holt GENAU EINEN abholbaren Job (faellig UND nicht bereits von
// einem anderen Worker gesperrt, ODER dessen Sperre als abgestanden gilt)
// und markiert ihn atomar als "processing" (Akzeptanzkriterium 1 / Pruefung
// 1 — FOR UPDATE SKIP LOCKED erlaubt mehreren Worker-Goroutinen
// gleichzeitigen Aufruf ohne sich gegenseitig zu blockieren oder denselben
// Job doppelt zu holen).
func (q *Queue) Dequeue(ctx context.Context, workerID string, jobTypes []string) (*Job, error) {
tx, err := q.pool.Begin(ctx)
if err != nil {
return nil, fmt.Errorf("jobqueue: transaktion starten: %w", err)
}
defer func() { _ = tx.Rollback(ctx) }()
var typeFilter []string
if len(jobTypes) > 0 {
typeFilter = jobTypes
}
row := tx.QueryRow(ctx, `
SELECT id, job_type, payload, status, attempts, max_attempts, last_error
FROM processing_jobs
WHERE (
(status = 'pending' AND available_at <= now())
OR (status = 'processing' AND locked_at <= now() - ($2 * interval '1 second'))
)
AND ($1::text[] IS NULL OR job_type = ANY($1))
ORDER BY available_at
FOR UPDATE SKIP LOCKED
LIMIT 1
`, typeFilter, q.staleLockAfter.Seconds())
var j Job
if err := row.Scan(&j.ID, &j.JobType, &j.Payload, &j.Status, &j.Attempts, &j.MaxAttempts, &j.LastError); err != nil {
if errors.Is(err, pgx.ErrNoRows) {
return nil, ErrNoJobAvailable
}
return nil, fmt.Errorf("jobqueue: naechsten job lesen: %w", err)
}
if _, err := tx.Exec(ctx, `
UPDATE processing_jobs
SET status = 'processing', attempts = attempts + 1, locked_at = now(), locked_by = $2, updated_at = now()
WHERE id = $1
`, j.ID, workerID); err != nil {
return nil, fmt.Errorf("jobqueue: job sperren: %w", err)
}
if err := tx.Commit(ctx); err != nil {
return nil, fmt.Errorf("jobqueue: dequeue committen: %w", err)
}
j.Status = StatusProcessing
j.Attempts++
return &j, nil
}
// Complete markiert einen Job als erfolgreich abgeschlossen.
func (q *Queue) Complete(ctx context.Context, jobID string) error {
tag, err := q.pool.Exec(ctx, `
UPDATE processing_jobs SET status = 'succeeded', locked_at = NULL, locked_by = NULL, updated_at = now()
WHERE id = $1
`, jobID)
if err != nil {
return fmt.Errorf("jobqueue: job abschliessen: %w", err)
}
if tag.RowsAffected() == 0 {
return ErrNotFound
}
return nil
}
// Fail markiert einen Job als fehlgeschlagen (Akzeptanzkriterium 2): sind
// die maximalen Versuche erreicht, wandert der Job in die Dead-Letter-Queue
// (status='dead_letter'), sonst wird er mit exponentiellem Backoff erneut
// eingeplant. Backoff-Berechnung nutzt arithmetischen Intervall-Cast
// (attempts * interval), KEINE String-Konkatenation (siehe "Bekannte
// Fehler vermeiden" im Ticket).
func (q *Queue) Fail(ctx context.Context, jobID string, cause error) error {
errMsg := cause.Error()
tag, err := q.pool.Exec(ctx, `
UPDATE processing_jobs
SET status = CASE WHEN attempts >= max_attempts THEN 'dead_letter' ELSE 'pending' END,
available_at = now() + (LEAST(attempts, 10) * interval '30 seconds'),
locked_at = NULL, locked_by = NULL, last_error = $2, updated_at = now()
WHERE id = $1
`, jobID, errMsg)
if err != nil {
return fmt.Errorf("jobqueue: fehlschlag erfassen: %w", err)
}
if tag.RowsAffected() == 0 {
return ErrNotFound
}
return nil
}
// Status liefert den aktuellen Zustand eines Jobs (Akzeptanzkriterium 3).
func (q *Queue) Status(ctx context.Context, jobID string) (*Job, error) {
var j Job
err := q.pool.QueryRow(ctx, `
SELECT id, job_type, payload, status, attempts, max_attempts, last_error
FROM processing_jobs WHERE id = $1
`, jobID).Scan(&j.ID, &j.JobType, &j.Payload, &j.Status, &j.Attempts, &j.MaxAttempts, &j.LastError)
if err != nil {
if errors.Is(err, pgx.ErrNoRows) {
return nil, ErrNotFound
}
return nil, fmt.Errorf("jobqueue: job-status lesen: %w", err)
}
return &j, nil
}
// RequeueDeadLetter holt einen Job manuell aus der Dead-Letter-Queue zurueck
// in "pending", mit zurueckgesetztem Versuchszaehler (Pruefung 3: DLQ-Eintrag
// manuell wiederholbar). Nur fuer Jobs, die tatsaechlich in dead_letter
// stehen — verhindert versehentliches Requeue eines noch laufenden Jobs.
func (q *Queue) RequeueDeadLetter(ctx context.Context, jobID string) error {
tag, err := q.pool.Exec(ctx, `
UPDATE processing_jobs
SET status = 'pending', attempts = 0, available_at = now(), last_error = NULL, updated_at = now()
WHERE id = $1 AND status = 'dead_letter'
`, jobID)
if err != nil {
return fmt.Errorf("jobqueue: dead-letter-job erneut einreihen: %w", err)
}
if tag.RowsAffected() == 0 {
return fmt.Errorf("jobqueue: job %q steht nicht in dead_letter (oder existiert nicht): %w", jobID, ErrNotFound)
}
return nil
}
-233
View File
@@ -1,233 +0,0 @@
package jobqueue
import (
"context"
"encoding/json"
"errors"
"os"
"testing"
"time"
"github.com/jackc/pgx/v5/pgxpool"
)
func setupTest(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() })
if _, err := pool.Exec(ctx, `
CREATE EXTENSION IF NOT EXISTS pgcrypto;
CREATE TABLE IF NOT EXISTS processing_jobs (
id UUID PRIMARY KEY DEFAULT gen_random_uuid(),
job_type TEXT NOT NULL,
payload JSONB NOT NULL DEFAULT '{}'::jsonb,
idempotency_key TEXT UNIQUE,
status TEXT NOT NULL DEFAULT 'pending'
CHECK (status IN ('pending', 'processing', 'succeeded', 'failed', 'dead_letter')),
attempts INT NOT NULL DEFAULT 0,
max_attempts INT NOT NULL DEFAULT 5,
available_at TIMESTAMPTZ NOT NULL DEFAULT now(),
locked_at TIMESTAMPTZ,
locked_by TEXT,
last_error TEXT,
created_at TIMESTAMPTZ NOT NULL DEFAULT now(),
updated_at TIMESTAMPTZ NOT NULL DEFAULT now()
);
`); err != nil {
t.Fatalf("schema: %v", err)
}
// DROP statt nur TRUNCATE: internal/jobqueue und internal/migrate teilen
// sich dieselbe physische Test-Datenbank (TEST_TENANT_DSN) ueber
// Paketgrenzen hinweg. Bliebe die Tabelle stehen, wuerde internal/migrate
// spaeter mit "relation already exists" gegen die von diesem Fixture
// angelegte, aber unversionierte Tabelle scheitern.
t.Cleanup(func() {
_, _ = pool.Exec(context.Background(), `DROP TABLE IF EXISTS processing_jobs`)
})
return pool
}
// TestEnqueueDequeueComplete ist Akzeptanzkriterium 1: Jobs werden
// zuverlaessig eingereiht und verarbeitet.
func TestEnqueueDequeueComplete(t *testing.T) {
pool := setupTest(t)
q := NewQueue(pool, time.Minute)
ctx := context.Background()
id, err := q.Enqueue(ctx, "index-document", map[string]string{"document_id": "d1"}, EnqueueOptions{})
if err != nil {
t.Fatalf("enqueue: %v", err)
}
job, err := q.Dequeue(ctx, "worker-1", nil)
if err != nil {
t.Fatalf("dequeue: %v", err)
}
if job.ID != id {
t.Fatalf("dequeue lieferte job %q, want %q", job.ID, id)
}
if job.Status != StatusProcessing {
t.Fatalf("status nach dequeue = %q, want %q", job.Status, StatusProcessing)
}
var payload map[string]string
if err := json.Unmarshal(job.Payload, &payload); err != nil {
t.Fatalf("payload dekodieren: %v", err)
}
if payload["document_id"] != "d1" {
t.Fatalf("payload = %v, want document_id=d1", payload)
}
if err := q.Complete(ctx, job.ID); err != nil {
t.Fatalf("complete: %v", err)
}
status, err := q.Status(ctx, job.ID)
if err != nil {
t.Fatalf("status: %v", err)
}
if status.Status != StatusSucceeded {
t.Fatalf("endstatus = %q, want %q", status.Status, StatusSucceeded)
}
}
// TestDequeue_NoJobAvailable prueft den Leerfall.
func TestDequeue_NoJobAvailable(t *testing.T) {
pool := setupTest(t)
q := NewQueue(pool, time.Minute)
if _, err := q.Dequeue(context.Background(), "worker-1", nil); !errors.Is(err, ErrNoJobAvailable) {
t.Fatalf("erwartet ErrNoJobAvailable, habe %v", err)
}
}
// TestFail_LandsInDeadLetterAfterMaxAttempts ist Akzeptanzkriterium 2.
func TestFail_LandsInDeadLetterAfterMaxAttempts(t *testing.T) {
pool := setupTest(t)
q := NewQueue(pool, time.Minute)
ctx := context.Background()
id, err := q.Enqueue(ctx, "ocr", nil, EnqueueOptions{MaxAttempts: 2})
if err != nil {
t.Fatalf("enqueue: %v", err)
}
// Versuch 1: schlaegt fehl, geht zurueck nach "pending" (max_attempts=2 noch nicht erreicht).
job, err := q.Dequeue(ctx, "worker-1", nil)
if err != nil {
t.Fatalf("dequeue 1: %v", err)
}
if err := q.Fail(ctx, job.ID, errors.New("ocr-engine nicht erreichbar")); err != nil {
t.Fatalf("fail 1: %v", err)
}
status, err := q.Status(ctx, id)
if err != nil {
t.Fatalf("status nach fail 1: %v", err)
}
if status.Status != StatusPending {
t.Fatalf("status nach fail 1 = %q, want %q (noch nicht erschoepft)", status.Status, StatusPending)
}
// Versuch 2: Backoff manuell umgehen (available_at direkt zuruecksetzen,
// damit der Test nicht auf echten Backoff warten muss).
if _, err := pool.Exec(ctx, `UPDATE processing_jobs SET available_at = now() WHERE id = $1`, id); err != nil {
t.Fatalf("available_at zuruecksetzen: %v", err)
}
job2, err := q.Dequeue(ctx, "worker-1", nil)
if err != nil {
t.Fatalf("dequeue 2: %v", err)
}
if err := q.Fail(ctx, job2.ID, errors.New("ocr-engine weiterhin nicht erreichbar")); err != nil {
t.Fatalf("fail 2: %v", err)
}
final, err := q.Status(ctx, id)
if err != nil {
t.Fatalf("status nach fail 2: %v", err)
}
if final.Status != StatusDeadLetter {
t.Fatalf("status nach erschoepften versuchen = %q, want %q", final.Status, StatusDeadLetter)
}
if final.LastError == nil || *final.LastError == "" {
t.Fatal("erwartet gesetzten last_error im dead-letter-eintrag")
}
}
// TestRequeueDeadLetter ist Pruefung 3: DLQ-Eintrag manuell wiederholbar.
func TestRequeueDeadLetter(t *testing.T) {
pool := setupTest(t)
q := NewQueue(pool, time.Minute)
ctx := context.Background()
id, err := q.Enqueue(ctx, "convert", nil, EnqueueOptions{MaxAttempts: 1})
if err != nil {
t.Fatalf("enqueue: %v", err)
}
job, err := q.Dequeue(ctx, "worker-1", nil)
if err != nil {
t.Fatalf("dequeue: %v", err)
}
if err := q.Fail(ctx, job.ID, errors.New("konverter abgestuerzt")); err != nil {
t.Fatalf("fail: %v", err)
}
status, _ := q.Status(ctx, id)
if status.Status != StatusDeadLetter {
t.Fatalf("voraussetzung nicht erfuellt: job sollte in dead_letter stehen, ist %q", status.Status)
}
if err := q.RequeueDeadLetter(ctx, id); err != nil {
t.Fatalf("requeuedeadletter: %v", err)
}
afterRequeue, err := q.Status(ctx, id)
if err != nil {
t.Fatalf("status nach requeue: %v", err)
}
if afterRequeue.Status != StatusPending {
t.Fatalf("status nach requeue = %q, want %q", afterRequeue.Status, StatusPending)
}
if afterRequeue.Attempts != 0 {
t.Fatalf("attempts nach requeue = %d, want 0", afterRequeue.Attempts)
}
// Requeue eines NICHT in dead_letter stehenden Jobs wird abgewiesen.
if err := q.RequeueDeadLetter(ctx, id); !errors.Is(err, ErrNotFound) {
t.Fatalf("requeue eines pending-jobs: erwartet ErrNotFound, habe %v", err)
}
}
// TestEnqueue_IdempotencyKeyPreventsDuplicate ist Pruefung 2: Idempotenz
// bei Doppelzustellung nachgewiesen (auf Einreih-Ebene).
func TestEnqueue_IdempotencyKeyPreventsDuplicate(t *testing.T) {
pool := setupTest(t)
q := NewQueue(pool, time.Minute)
ctx := context.Background()
id1, err := q.Enqueue(ctx, "ocr", nil, EnqueueOptions{IdempotencyKey: "ocr:revision-42"})
if err != nil {
t.Fatalf("enqueue 1: %v", err)
}
id2, err := q.Enqueue(ctx, "ocr", nil, EnqueueOptions{IdempotencyKey: "ocr:revision-42"})
if err != nil {
t.Fatalf("enqueue 2 (doppelzustellung): %v", err)
}
if id1 != id2 {
t.Fatalf("doppelte einreihung mit gleichem idempotency-key erzeugte zwei jobs: %q != %q", id1, id2)
}
var count int
if err := pool.QueryRow(ctx, `SELECT count(*) FROM processing_jobs WHERE idempotency_key = 'ocr:revision-42'`).Scan(&count); err != nil {
t.Fatalf("zeilen zaehlen: %v", err)
}
if count != 1 {
t.Fatalf("erwartet genau 1 zeile fuer den idempotency-key, habe %d", count)
}
}
-75
View File
@@ -1,75 +0,0 @@
package jobqueue
import (
"context"
"errors"
"log"
"time"
)
// Handler verarbeitet einen einzelnen Job. Ein zurueckgegebener Fehler
// fuehrt zu Queue.Fail (Backoff/DLQ), nil zu Queue.Complete.
type Handler func(ctx context.Context, job *Job) error
// Worker pollt die Queue in einer In-Prozess-Goroutine und ruft Handler je
// abgeholtem Job auf — die "In-Prozess-Worker-Goroutinen" aus der
// Ticket-Vorgabe, kein separater Prozess/Redis noetig.
type Worker struct {
queue *Queue
workerID string
jobTypes []string
pollInterval time.Duration
handler Handler
}
func NewWorker(queue *Queue, workerID string, jobTypes []string, pollInterval time.Duration, handler Handler) *Worker {
return &Worker{queue: queue, workerID: workerID, jobTypes: jobTypes, pollInterval: pollInterval, handler: handler}
}
// Run blockiert, bis ctx beendet wird, und verarbeitet dabei fortlaufend
// Jobs. Absturzsicherheit (Pruefung 1) entsteht NICHT durch Run selbst,
// sondern dadurch, dass ein abgestuerzter Prozess (der Run gar nicht mehr
// ausfuehrt) seine "processing"-Sperren nie verlaengert — ein ANDERER
// Worker-Prozess holt den Job nach Ablauf von staleLockAfter erneut ab
// (siehe Queue.Dequeue).
func (w *Worker) Run(ctx context.Context) {
ticker := time.NewTicker(w.pollInterval)
defer ticker.Stop()
for {
select {
case <-ctx.Done():
return
case <-ticker.C:
w.processOne(ctx)
}
}
}
// processOne holt und verarbeitet EINEN Job, falls verfuegbar. Oeffentlich
// über RunOnce fuer Tests, die deterministisch (ohne Polling-Timing) einen
// einzelnen Verarbeitungsschritt auslösen wollen.
func (w *Worker) processOne(ctx context.Context) {
job, err := w.queue.Dequeue(ctx, w.workerID, w.jobTypes)
if err != nil {
if !errors.Is(err, ErrNoJobAvailable) {
log.Printf("jobqueue: dequeue fehlgeschlagen: %v", err)
}
return
}
if handlerErr := w.handler(ctx, job); handlerErr != nil {
if err := w.queue.Fail(ctx, job.ID, handlerErr); err != nil {
log.Printf("jobqueue: fehlschlag fuer job %q nicht erfassbar: %v", job.ID, err)
}
return
}
if err := w.queue.Complete(ctx, job.ID); err != nil {
log.Printf("jobqueue: abschluss fuer job %q fehlgeschlagen: %v", job.ID, err)
}
}
// RunOnce verarbeitet synchron genau einen Job (falls verfuegbar) und
// kehrt zurueck — fuer Tests, die ohne Polling-Intervall arbeiten wollen.
func (w *Worker) RunOnce(ctx context.Context) {
w.processOne(ctx)
}
-115
View File
@@ -1,115 +0,0 @@
package jobqueue
import (
"context"
"errors"
"sync/atomic"
"testing"
"time"
)
// TestWorker_RunOnce_ProcessesAndCompletesJob ist der End-to-End-Nachweis
// fuer Akzeptanzkriterium 1 ueber den Worker statt direkt ueber Queue.
func TestWorker_RunOnce_ProcessesAndCompletesJob(t *testing.T) {
pool := setupTest(t)
q := NewQueue(pool, time.Minute)
ctx := context.Background()
id, err := q.Enqueue(ctx, "index-document", nil, EnqueueOptions{})
if err != nil {
t.Fatalf("enqueue: %v", err)
}
var handled int32
w := NewWorker(q, "worker-1", nil, time.Millisecond, func(ctx context.Context, job *Job) error {
atomic.AddInt32(&handled, 1)
return nil
})
w.RunOnce(ctx)
if atomic.LoadInt32(&handled) != 1 {
t.Fatalf("handler wurde %d mal aufgerufen, want 1", handled)
}
status, err := q.Status(ctx, id)
if err != nil {
t.Fatalf("status: %v", err)
}
if status.Status != StatusSucceeded {
t.Fatalf("status = %q, want %q", status.Status, StatusSucceeded)
}
}
// TestWorker_HandlerErrorTriggersFail prueft, dass ein Handler-Fehler zu
// Queue.Fail fuehrt (Backoff/DLQ-Pfad ueber den Worker statt direkt).
func TestWorker_HandlerErrorTriggersFail(t *testing.T) {
pool := setupTest(t)
q := NewQueue(pool, time.Minute)
ctx := context.Background()
id, err := q.Enqueue(ctx, "ocr", nil, EnqueueOptions{MaxAttempts: 5})
if err != nil {
t.Fatalf("enqueue: %v", err)
}
w := NewWorker(q, "worker-1", nil, time.Millisecond, func(ctx context.Context, job *Job) error {
return errors.New("ocr fehlgeschlagen")
})
w.RunOnce(ctx)
status, err := q.Status(ctx, id)
if err != nil {
t.Fatalf("status: %v", err)
}
if status.Status != StatusPending {
t.Fatalf("status nach handler-fehler = %q, want %q (erneut eingeplant)", status.Status, StatusPending)
}
if status.LastError == nil || *status.LastError != "ocr fehlgeschlagen" {
t.Fatalf("last_error = %v, want %q", status.LastError, "ocr fehlgeschlagen")
}
}
// TestDequeue_StaleLockIsRedelivered ist Pruefung 1: Absturz eines Workers
// fuehrt zu erneuter Zustellung. Simuliert einen Absturz, indem ein Job
// dequeued (auf "processing" gesperrt), aber NIE completed/failed wird —
// nach Ablauf von staleLockAfter muss ein ANDERER Worker denselben Job
// erneut abholen koennen.
func TestDequeue_StaleLockIsRedelivered(t *testing.T) {
pool := setupTest(t)
// Sehr kurzes Stale-Fenster, damit der Test nicht lange warten muss.
q := NewQueue(pool, 50*time.Millisecond)
ctx := context.Background()
id, err := q.Enqueue(ctx, "convert", nil, EnqueueOptions{})
if err != nil {
t.Fatalf("enqueue: %v", err)
}
crashed, err := q.Dequeue(ctx, "worker-crashed", nil)
if err != nil {
t.Fatalf("dequeue (worker-crashed): %v", err)
}
if crashed.ID != id {
t.Fatalf("dequeue lieferte unerwarteten job %q", crashed.ID)
}
// worker-crashed ruft absichtlich weder Complete noch Fail auf — simuliert
// einen Prozessabsturz mitten in der Verarbeitung.
// Sofortiger erneuter Dequeue-Versuch (Sperre noch frisch) darf den Job
// NICHT liefern.
if _, err := q.Dequeue(ctx, "worker-2", nil); !errors.Is(err, ErrNoJobAvailable) {
t.Fatalf("job wurde trotz frischer sperre erneut ausgeliefert (oder anderer fehler): %v", err)
}
time.Sleep(80 * time.Millisecond) // > staleLockAfter
redelivered, err := q.Dequeue(ctx, "worker-2", nil)
if err != nil {
t.Fatalf("dequeue nach ablauf der sperre: %v", err)
}
if redelivered.ID != id {
t.Fatalf("erneut zugestellter job = %q, want %q", redelivered.ID, id)
}
if err := q.Complete(ctx, redelivered.ID); err != nil {
t.Fatalf("complete durch worker-2: %v", err)
}
}
-165
View File
@@ -1,165 +0,0 @@
// Package migrate implementiert FDN-02s Migrationsmechanik: versionierte,
// rueckrollbare SQL-Migrationsdateien (kein ORM), mit einer
// schema_migrations-Tabelle als Fortschrittsspeicher — dasselbe Prinzip wie
// NEXARCH Core (internal/migrate), hier eigenstaendig implementiert, da DMS
// ein eigenes Go-Modul ist und Cores internal/-Pakete nicht importieren
// kann.
package migrate
import (
"context"
"fmt"
"os"
"path/filepath"
"sort"
"strings"
"github.com/jackc/pgx/v5/pgxpool"
)
// Migration ist eine einzelne versionierte Migrationsdatei.
type Migration struct {
Version string // Dateiname ohne .up.sql/.down.sql, z.B. "0001_documents"
UpSQL string
DownSQL string
}
// Load liest alle *.up.sql/*.down.sql-Paare aus dir, sortiert nach
// Dateiname (Akzeptanzkriterium: Migrationen sind versioniert).
func Load(dir string) ([]Migration, error) {
entries, err := os.ReadDir(dir)
if err != nil {
return nil, fmt.Errorf("migrationsverzeichnis %q lesen: %w", dir, err)
}
var versions []string
for _, e := range entries {
if e.IsDir() || !strings.HasSuffix(e.Name(), ".up.sql") {
continue
}
versions = append(versions, strings.TrimSuffix(e.Name(), ".up.sql"))
}
sort.Strings(versions)
migrations := make([]Migration, 0, len(versions))
for _, v := range versions {
up, err := os.ReadFile(filepath.Join(dir, v+".up.sql"))
if err != nil {
return nil, fmt.Errorf("migration %q: up.sql lesen: %w", v, err)
}
down, err := os.ReadFile(filepath.Join(dir, v+".down.sql"))
if err != nil {
return nil, fmt.Errorf("migration %q: down.sql lesen (jede Migration braucht ein Rollback): %w", v, err)
}
migrations = append(migrations, Migration{Version: v, UpSQL: string(up), DownSQL: string(down)})
}
return migrations, nil
}
func ensureTrackingTable(ctx context.Context, pool *pgxpool.Pool) error {
_, err := pool.Exec(ctx, `
CREATE TABLE IF NOT EXISTS schema_migrations (
version TEXT PRIMARY KEY,
applied_at TIMESTAMPTZ NOT NULL DEFAULT now()
)
`)
if err != nil {
return fmt.Errorf("schema_migrations anlegen: %w", err)
}
return nil
}
func appliedVersions(ctx context.Context, pool *pgxpool.Pool) (map[string]bool, error) {
rows, err := pool.Query(ctx, `SELECT version FROM schema_migrations`)
if err != nil {
return nil, fmt.Errorf("angewendete migrationen lesen: %w", err)
}
defer rows.Close()
applied := map[string]bool{}
for rows.Next() {
var v string
if err := rows.Scan(&v); err != nil {
return nil, fmt.Errorf("migrationsversion lesen: %w", err)
}
applied[v] = true
}
return applied, rows.Err()
}
// Up wendet alle noch nicht angewendeten Migrationen in Reihenfolge an
// (Akzeptanzkriterium 2: vorwaerts ausfuehrbar) — bereits angewendete
// werden uebersprungen, damit Up auf einer leeren UND auf einer bestehenden
// DB funktioniert (Pruefung 1).
func Up(ctx context.Context, pool *pgxpool.Pool, migrations []Migration) (applied []string, err error) {
if err := ensureTrackingTable(ctx, pool); err != nil {
return nil, err
}
already, err := appliedVersions(ctx, pool)
if err != nil {
return nil, err
}
for _, m := range migrations {
if already[m.Version] {
continue
}
tx, err := pool.Begin(ctx)
if err != nil {
return applied, fmt.Errorf("transaktion fuer %q starten: %w", m.Version, err)
}
if _, err := tx.Exec(ctx, m.UpSQL); err != nil {
_ = tx.Rollback(ctx)
return applied, fmt.Errorf("migration %q anwenden: %w", m.Version, err)
}
if _, err := tx.Exec(ctx, `INSERT INTO schema_migrations (version) VALUES ($1)`, m.Version); err != nil {
_ = tx.Rollback(ctx)
return applied, fmt.Errorf("migration %q als angewendet markieren: %w", m.Version, err)
}
if err := tx.Commit(ctx); err != nil {
return applied, fmt.Errorf("migration %q committen: %w", m.Version, err)
}
applied = append(applied, m.Version)
}
return applied, nil
}
// DownOne macht die zuletzt angewendete Migration rueckgaengig
// (Akzeptanzkriterium 2: rueckwaerts ausfuehrbar) und liefert deren Version,
// oder "" falls keine Migration angewendet war.
func DownOne(ctx context.Context, pool *pgxpool.Pool, migrations []Migration) (version string, err error) {
if err := ensureTrackingTable(ctx, pool); err != nil {
return "", err
}
already, err := appliedVersions(ctx, pool)
if err != nil {
return "", err
}
var last *Migration
for i := len(migrations) - 1; i >= 0; i-- {
if already[migrations[i].Version] {
last = &migrations[i]
break
}
}
if last == nil {
return "", nil
}
tx, err := pool.Begin(ctx)
if err != nil {
return "", fmt.Errorf("transaktion fuer rollback von %q starten: %w", last.Version, err)
}
if _, err := tx.Exec(ctx, last.DownSQL); err != nil {
_ = tx.Rollback(ctx)
return "", fmt.Errorf("migration %q zurueckrollen: %w", last.Version, err)
}
if _, err := tx.Exec(ctx, `DELETE FROM schema_migrations WHERE version = $1`, last.Version); err != nil {
_ = tx.Rollback(ctx)
return "", fmt.Errorf("migration %q aus schema_migrations entfernen: %w", last.Version, err)
}
if err := tx.Commit(ctx); err != nil {
return "", fmt.Errorf("rollback von %q committen: %w", last.Version, err)
}
return last.Version, nil
}
-213
View File
@@ -1,213 +0,0 @@
package migrate
import (
"context"
"os"
"path/filepath"
"testing"
"github.com/jackc/pgx/v5/pgxpool"
)
// usersFixtureSQL spiegelt Core IAM-01s reales Schema
// (migrations/tenant/0001_users.up.sql im NEXARCH-Core-Modul) - DMS ist ein
// eigenes Go-Modul und kann Cores Migrationsdateien nicht importieren, daher
// hier als Testfixture kopiert, NICHT als Produktionsmigration (DMS legt
// users nicht selbst an, siehe "Nicht Bestandteil" in FDN-02).
const usersFixtureSQL = `
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,
status TEXT NOT NULL DEFAULT 'active',
created_at TIMESTAMPTZ NOT NULL DEFAULT now(),
updated_at TIMESTAMPTZ NOT NULL DEFAULT now()
);
`
func setupTest(t *testing.T) (*pgxpool.Pool, []Migration) {
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() })
if _, err := pool.Exec(ctx, usersFixtureSQL); err != nil {
t.Fatalf("users-fixture anlegen: %v", err)
}
if _, err := pool.Exec(ctx, `
INSERT INTO users (email, name) VALUES ('seed@example.test', 'Seed-Benutzer')
ON CONFLICT (email) DO NOTHING
`); err != nil {
t.Fatalf("seed-benutzer anlegen: %v", err)
}
migrations, err := Load(migrationsDir(t))
if err != nil {
t.Fatalf("migrationen laden: %v", err)
}
return pool, migrations
}
func migrationsDir(t *testing.T) string {
t.Helper()
wd, err := os.Getwd()
if err != nil {
t.Fatalf("getwd: %v", err)
}
return filepath.Join(wd, "..", "..", "migrations", "tenant")
}
func tableExists(t *testing.T, ctx context.Context, pool *pgxpool.Pool, name string) bool {
t.Helper()
var exists bool
if err := pool.QueryRow(ctx, `
SELECT EXISTS (SELECT 1 FROM information_schema.tables WHERE table_name = $1)
`, name).Scan(&exists); err != nil {
t.Fatalf("tabellenexistenz von %q pruefen: %v", name, err)
}
return exists
}
// TestUp_OnEmptyAndExistingDB ist Pruefung 1: Migration auf leerer DB und
// auf bestehender DB getestet.
func TestUp_OnEmptyAndExistingDB(t *testing.T) {
pool, migrations := setupTest(t)
ctx := context.Background()
applied, err := Up(ctx, pool, migrations)
if err != nil {
t.Fatalf("up (leere db): %v", err)
}
if len(applied) == 0 {
t.Fatal("erwartet mindestens 1 angewendete migration auf leerer db")
}
for _, table := range []string{"folders", "documents", "file_revisions", "tags", "document_tags", "metadata_fields", "document_metadata_values"} {
if !tableExists(t, ctx, pool, table) {
t.Fatalf("tabelle %q existiert nach Up nicht", table)
}
}
// Zweiter Up-Lauf gegen die JETZT BESTEHENDE db - muss ohne Fehler
// durchlaufen und darf nichts erneut anwenden (idempotent ueber
// schema_migrations).
appliedAgain, err := Up(ctx, pool, migrations)
if err != nil {
t.Fatalf("up (bestehende db, zweiter lauf): %v", err)
}
if len(appliedAgain) != 0 {
t.Fatalf("zweiter Up-Lauf haette 0 neue migrationen anwenden sollen, hat %d", len(appliedAgain))
}
}
// TestDownOne_RestoresPreviousState ist Pruefung 2: Rollback stellt den
// Vorzustand wieder her.
func TestDownOne_RestoresPreviousState(t *testing.T) {
pool, migrations := setupTest(t)
ctx := context.Background()
if _, err := Up(ctx, pool, migrations); err != nil {
t.Fatalf("up: %v", err)
}
if !tableExists(t, ctx, pool, "documents") {
t.Fatal("voraussetzung nicht erfuellt: documents sollte nach Up existieren")
}
// DownOne rollt IMMER nur die zuletzt angewendete Migration zurueck
// (dokumentiertes Verhalten) — bei mehreren Migrationen (z.B. FDN-02
// documents + FDN-04 processing_jobs) muss man entsprechend oft
// aufrufen. Anzahl der Wiederholungen richtet sich NICHT nach dem
// Rueckgabewert von Up() (der nur die in DIESEM Aufruf NEU angewendeten
// Migrationen zaehlt, siehe Kommentar an Up) — stattdessen wird
// wiederholt, bis DownOne "" liefert (schema_migrations leer).
for i := 0; i < len(migrations); i++ {
version, err := DownOne(ctx, pool, migrations)
if err != nil {
t.Fatalf("downone (lauf %d): %v", i, err)
}
if version == "" {
break
}
}
for _, table := range []string{"folders", "documents", "file_revisions", "tags", "document_tags", "metadata_fields", "document_metadata_values", "processing_jobs"} {
if tableExists(t, ctx, pool, table) {
t.Fatalf("tabelle %q existiert nach vollstaendigem Rollback noch - Vorzustand nicht wiederhergestellt", table)
}
}
// Erneutes DownOne ohne verbleibende angewendete Migration liefert "".
versionAfterAll, err := DownOne(ctx, pool, migrations)
if err != nil {
t.Fatalf("downone (nichts mehr anzuwenden): %v", err)
}
if versionAfterAll != "" {
t.Fatalf("erwartet leeren string bei leerer schema_migrations, habe %q", versionAfterAll)
}
}
// TestForeignKeyConstraints_RejectInvalidReferences ist Pruefung 3:
// Fremdschluessel-Constraints durch Negativtests belegt.
func TestForeignKeyConstraints_RejectInvalidReferences(t *testing.T) {
pool, migrations := setupTest(t)
ctx := context.Background()
if _, err := Up(ctx, pool, migrations); err != nil {
t.Fatalf("up: %v", err)
}
var userID string
if err := pool.QueryRow(ctx, `SELECT id FROM users LIMIT 1`).Scan(&userID); err != nil {
t.Fatalf("seed-benutzer lesen: %v", err)
}
t.Run("dokument mit unbekanntem ordner wird abgewiesen", func(t *testing.T) {
_, err := pool.Exec(ctx, `
INSERT INTO documents (folder_id, title, created_by) VALUES (gen_random_uuid(), 'x', $1)
`, userID)
if err == nil {
t.Fatal("insert mit unbekanntem folder_id haette scheitern muessen")
}
})
t.Run("dokument mit unbekanntem ersteller wird abgewiesen", func(t *testing.T) {
_, err := pool.Exec(ctx, `
INSERT INTO documents (title, created_by) VALUES ('x', gen_random_uuid())
`)
if err == nil {
t.Fatal("insert mit unbekanntem created_by haette scheitern muessen")
}
})
t.Run("datei-revision mit unbekanntem dokument wird abgewiesen", func(t *testing.T) {
_, err := pool.Exec(ctx, `
INSERT INTO file_revisions (document_id, revision_number, storage_key, checksum_sha256, size_bytes, mime_type, created_by)
VALUES (gen_random_uuid(), 1, 'k', repeat('0',64), 1, 'text/plain', $1)
`, userID)
if err == nil {
t.Fatal("insert mit unbekanntem document_id haette scheitern muessen")
}
})
t.Run("tag-zuordnung mit unbekanntem tag wird abgewiesen", func(t *testing.T) {
var docID string
if err := pool.QueryRow(ctx, `
INSERT INTO documents (title, created_by) VALUES ('fk-test-doc', $1) RETURNING id
`, userID).Scan(&docID); err != nil {
t.Fatalf("testdokument anlegen: %v", err)
}
_, err := pool.Exec(ctx, `
INSERT INTO document_tags (document_id, tag_id) VALUES ($1, gen_random_uuid())
`, docID)
if err == nil {
t.Fatal("insert mit unbekanntem tag_id haette scheitern muessen")
}
})
}
@@ -1,78 +0,0 @@
// Package retentionclient ist ein schlanker HTTP-Client für Archive
// RET-05 (RET-09: archive/internal/moduleadapter.RegisterHandler, real
// als Dienst gestartet). DMS ist ein physisch getrenntes Go-Modul und
// kann Archive-internen Go-Code nicht direkt importieren.
package retentionclient
import (
"bytes"
"context"
"encoding/json"
"fmt"
"net/http"
)
type Client struct {
BaseURL string
HTTPClient *http.Client
}
func New(baseURL string) *Client {
return &Client{BaseURL: baseURL, HTTPClient: http.DefaultClient}
}
type registerRequest struct {
ModuleName string `json:"module_name"`
ObjectType string `json:"object_type"`
RetentionClass string `json:"retention_class"`
CallbackURL string `json:"callback_url"`
}
// Registration spiegelt RET-05s Registration-Antwort.
type Registration struct {
ID string `json:"ID"`
ModuleName string `json:"ModuleName"`
ObjectType string `json:"ObjectType"`
RetentionClass string `json:"RetentionClass"`
CallbackURL string `json:"CallbackURL"`
}
// Register registriert einen Objekttyp bei Archive RET-05. Jeder Fehler
// (Transport, Timeout, Nicht-2xx) wird als Fehler zurückgegeben — der
// Aufrufer entscheidet ueber Requeue statt stillem Verwerfen
// (Akzeptanzkriterium 3).
func (c *Client) Register(ctx context.Context, moduleName, objectType, retentionClass, callbackURL string) (Registration, error) {
body, err := json.Marshal(registerRequest{
ModuleName: moduleName, ObjectType: objectType, RetentionClass: retentionClass, CallbackURL: callbackURL,
})
if err != nil {
return Registration{}, fmt.Errorf("retentionclient: request kodieren: %w", err)
}
req, err := http.NewRequestWithContext(ctx, http.MethodPost, c.BaseURL+"/register", bytes.NewReader(body))
if err != nil {
return Registration{}, fmt.Errorf("retentionclient: request bauen: %w", err)
}
req.Header.Set("Content-Type", "application/json")
resp, err := c.httpClient().Do(req)
if err != nil {
return Registration{}, fmt.Errorf("retentionclient: aufruf fehlgeschlagen: %w", err)
}
defer func() { _ = resp.Body.Close() }()
if resp.StatusCode < 200 || resp.StatusCode >= 300 {
return Registration{}, fmt.Errorf("retentionclient: unerwarteter status %d", resp.StatusCode)
}
var out Registration
if err := json.NewDecoder(resp.Body).Decode(&out); err != nil {
return Registration{}, fmt.Errorf("retentionclient: antwort dekodieren: %w", err)
}
return out, nil
}
func (c *Client) httpClient() *http.Client {
if c.HTTPClient != nil {
return c.HTTPClient
}
return http.DefaultClient
}
@@ -1,148 +0,0 @@
// Package retentiondestroy implementiert DOC-16: DMS als Konsument des
// Archive-RET-05-Vertrags. Registriert den Objekttyp "dms_document" bei
// Archive (RET-09, real laufender Dienst) und empfängt/verarbeitet den
// Vernichtungs-Rückruf. Kennt KEINE eigene Fristenberechnung — reiner
// Adapter zwischen DMS und dem bereits feststehenden RET-05-Vertrag.
package retentiondestroy
import (
"context"
"encoding/json"
"fmt"
"net/http"
"time"
"github.com/jackc/pgx/v5/pgxpool"
"gitea.perlbach24.de/scripte/nexarch/dms/internal/jobqueue"
"gitea.perlbach24.de/scripte/nexarch/dms/internal/retentionclient"
"gitea.perlbach24.de/scripte/nexarch/dms/internal/storage"
)
// ModuleName ist der bei Archive RET-05 registrierte Modulname.
const ModuleName = "dms"
// ObjectTypeDocument ist der einzige aktuell in DOC-01/FDN-02 modellierte
// Objekttyp — DMS führt (noch) keinen separaten "Anhang"-Entitätstyp,
// daher wird ausschließlich "dms_document" registriert (kleinste Lösung,
// kein Umbau von DOC-01/FDN-02).
const ObjectTypeDocument = "dms_document"
// JobTypeRegister ist der FDN-04-Jobtyp für Registrierungs-Retries.
const JobTypeRegister = "doc16_retention_register"
// RegisterConfig ist der Registrierungs-Auftrag, sowohl für den
// Erstversuch als auch als Requeue-Payload.
type RegisterConfig struct {
ObjectType string `json:"object_type"`
RetentionClass string `json:"retention_class"`
CallbackURL string `json:"callback_url"`
}
// RegisterOrRequeue versucht die Registrierung real gegen Archive RET-05
// (RET-09). Schlägt sie fehl (Archive nicht erreichbar/Nicht-2xx), wird
// EIN Requeue-Job über FDN-04 eingereiht, statt abzustürzen oder den
// Fehler zu verschlucken (Akzeptanzkriterium 3) — der ursprüngliche
// Fehler wird IMMER zurückgegeben, damit der Aufrufer ihn protokolliert.
func RegisterOrRequeue(ctx context.Context, client *retentionclient.Client, queue *jobqueue.Queue, cfg RegisterConfig) (queued bool, err error) {
_, regErr := client.Register(ctx, ModuleName, cfg.ObjectType, cfg.RetentionClass, cfg.CallbackURL)
if regErr == nil {
return false, nil
}
if _, qerr := queue.Enqueue(ctx, JobTypeRegister, cfg, jobqueue.EnqueueOptions{
IdempotencyKey: "doc16-register-" + cfg.ObjectType,
}); qerr != nil {
return false, fmt.Errorf("retentiondestroy: registrierung fehlgeschlagen (%v) UND requeue fehlgeschlagen: %w", regErr, qerr)
}
return true, regErr
}
// ProcessRegisterJob verarbeitet EINEN dequeuten Registrierungs-Retry —
// vom Worker (FDN-04-Konsument) aufgerufen.
func ProcessRegisterJob(ctx context.Context, client *retentionclient.Client, payload json.RawMessage) error {
var cfg RegisterConfig
if err := json.Unmarshal(payload, &cfg); err != nil {
return fmt.Errorf("retentiondestroy: job-payload lesen: %w", err)
}
_, err := client.Register(ctx, ModuleName, cfg.ObjectType, cfg.RetentionClass, cfg.CallbackURL)
return err
}
type destroyRequest struct {
ObjectType string `json:"object_type"`
ObjectReference string `json:"object_reference"`
DestroyedAt string `json:"destroyed_at"`
}
// DestroyCallbackHandler ist der RET-05-Rückruf-Empfänger
// (Akzeptanzkriterium 2): löscht physisch alle Datei-Revisionen des
// referenzierten Dokuments über storage.Service (meldet die
// Größenänderung an Core, Service.Delete-Vertrag) und markiert das
// Dokument als vernichtet (deleted_at). Antwortet erst NACH
// erfolgreicher Löschung mit 2xx.
func DestroyCallbackHandler(pool *pgxpool.Pool, storageSvc *storage.Service) http.HandlerFunc {
return func(w http.ResponseWriter, r *http.Request) {
var req destroyRequest
if err := json.NewDecoder(r.Body).Decode(&req); err != nil {
http.Error(w, "ungültiger request-body: "+err.Error(), http.StatusBadRequest)
return
}
if req.ObjectType != ObjectTypeDocument || req.ObjectReference == "" {
http.Error(w, "unbekannter object_type oder fehlende object_reference", http.StatusBadRequest)
return
}
destroyedAt, err := time.Parse(time.RFC3339, req.DestroyedAt)
if err != nil {
http.Error(w, "ungültiges destroyed_at: "+err.Error(), http.StatusBadRequest)
return
}
ctx := r.Context()
rows, err := pool.Query(ctx, `SELECT storage_key, size_bytes FROM file_revisions WHERE document_id = $1`, req.ObjectReference)
if err != nil {
http.Error(w, "revisionen lesen: "+err.Error(), http.StatusInternalServerError)
return
}
type revision struct {
key string
size int64
}
var revisions []revision
for rows.Next() {
var rv revision
if err := rows.Scan(&rv.key, &rv.size); err != nil {
rows.Close()
http.Error(w, "revisions-zeile lesen: "+err.Error(), http.StatusInternalServerError)
return
}
revisions = append(revisions, rv)
}
rowsErr := rows.Err()
rows.Close()
if rowsErr != nil {
http.Error(w, "revisionen lesen: "+rowsErr.Error(), http.StatusInternalServerError)
return
}
for _, rv := range revisions {
if err := storageSvc.Delete(ctx, rv.key, rv.size); err != nil {
http.Error(w, "objekt-storage löschen fehlgeschlagen: "+err.Error(), http.StatusInternalServerError)
return
}
}
tag, err := pool.Exec(ctx, `UPDATE documents SET deleted_at = $2 WHERE id = $1`, req.ObjectReference, destroyedAt)
if err != nil {
http.Error(w, "dokument als vernichtet markieren: "+err.Error(), http.StatusInternalServerError)
return
}
if tag.RowsAffected() == 0 {
// Dokument existiert nicht (mehr) — aus Sicht des Rückrufs
// bereits erledigt, kein Fehler (verhindert Retry-Schleifen
// bei Archive fuer laengst geloeschte Objekte).
w.WriteHeader(http.StatusOK)
return
}
w.WriteHeader(http.StatusOK)
}
}
@@ -1,228 +0,0 @@
package retentiondestroy
import (
"context"
"encoding/json"
"net/http"
"net/http/httptest"
"os"
"strings"
"testing"
"time"
"github.com/jackc/pgx/v5/pgxpool"
"gitea.perlbach24.de/scripte/nexarch/dms/internal/jobqueue"
"gitea.perlbach24.de/scripte/nexarch/dms/internal/retentionclient"
"gitea.perlbach24.de/scripte/nexarch/dms/internal/storage"
)
func setupTest(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() })
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 DEFAULT 'Test'
);
CREATE TABLE IF NOT EXISTS documents (
id UUID PRIMARY KEY DEFAULT gen_random_uuid(),
folder_id UUID,
title TEXT NOT NULL,
current_revision_id UUID,
created_by UUID NOT NULL REFERENCES users(id),
created_at TIMESTAMPTZ NOT NULL DEFAULT now(),
updated_at TIMESTAMPTZ NOT NULL DEFAULT now(),
deleted_at TIMESTAMPTZ
);
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,
revision_number INT NOT NULL,
storage_key TEXT NOT NULL,
checksum_sha256 TEXT NOT NULL,
size_bytes BIGINT NOT NULL,
mime_type TEXT NOT NULL,
created_by UUID NOT NULL REFERENCES users(id),
created_at TIMESTAMPTZ NOT NULL DEFAULT now(),
UNIQUE (document_id, revision_number)
);
CREATE TABLE IF NOT EXISTS processing_jobs (
id UUID PRIMARY KEY DEFAULT gen_random_uuid(), job_type TEXT NOT NULL,
payload JSONB NOT NULL DEFAULT '{}'::jsonb, idempotency_key TEXT UNIQUE,
status TEXT NOT NULL DEFAULT 'pending' CHECK (status IN ('pending', 'processing', 'succeeded', 'failed', 'dead_letter')),
attempts INT NOT NULL DEFAULT 0, max_attempts INT NOT NULL DEFAULT 5,
available_at TIMESTAMPTZ NOT NULL DEFAULT now(), locked_at TIMESTAMPTZ, locked_by TEXT,
last_error TEXT, created_at TIMESTAMPTZ NOT NULL DEFAULT now(), updated_at TIMESTAMPTZ NOT NULL DEFAULT now()
);
`); err != nil {
t.Fatalf("schema: %v", err)
}
t.Cleanup(func() {
_, _ = pool.Exec(context.Background(), `TRUNCATE processing_jobs, file_revisions, documents, users CASCADE`)
})
return pool
}
type fakeUsageReporter struct{ reports []int64 }
func (f *fakeUsageReporter) Report(ctx context.Context, tenantSlug, metric string, delta int64) error {
f.reports = append(f.reports, delta)
return nil
}
// fakeRET05Server simuliert Archive RET-05/RET-09 (POST /register).
func fakeRET05Server(t *testing.T, fail bool) *retentionclient.Client {
t.Helper()
mux := http.NewServeMux()
mux.HandleFunc("POST /register", func(w http.ResponseWriter, r *http.Request) {
if fail {
http.Error(w, "simulierter fehler", http.StatusInternalServerError)
return
}
w.Header().Set("Content-Type", "application/json")
_ = json.NewEncoder(w).Encode(map[string]string{
"ID": "fake-reg-id", "ModuleName": ModuleName, "ObjectType": ObjectTypeDocument,
"RetentionClass": "test-klasse", "CallbackURL": "http://test/callback",
})
})
server := httptest.NewServer(mux)
t.Cleanup(server.Close)
return retentionclient.New(server.URL)
}
// TestRegisterOrRequeue_RealRegistrationAgainstTestEndpoint ist die
// geforderte Pflichtpruefung 1: Registrierung real gegen einen
// Test-RET-05-Endpunkt durchgefuehrt und verifiziert.
func TestRegisterOrRequeue_RealRegistrationAgainstTestEndpoint(t *testing.T) {
pool := setupTest(t)
client := fakeRET05Server(t, false)
queue := jobqueue.NewQueue(pool, time.Minute)
queued, err := RegisterOrRequeue(context.Background(), client, queue, RegisterConfig{
ObjectType: ObjectTypeDocument, RetentionClass: "test-klasse", CallbackURL: "http://test/callback",
})
if err != nil {
t.Fatalf("erwartet erfolgreiche registrierung, habe fehler: %v", err)
}
if queued {
t.Fatal("erfolgreiche registrierung haette NICHT requeued werden duerfen")
}
}
// TestDestroyCallbackHandler_RealDeletionAnd2xx ist die geforderte
// Pflichtpruefung 2: Vernichtungs-Rueckruf gegen einen echten Testfall
// ausgeloest, DMS-Loeschung nachweislich erfolgt, 2xx-Antwort gesendet.
func TestDestroyCallbackHandler_RealDeletionAnd2xx(t *testing.T) {
pool := setupTest(t)
ctx := context.Background()
var userID string
if err := pool.QueryRow(ctx, `INSERT INTO users (email, name) VALUES ('destroy-test@acme.example', 'Destroy Test') RETURNING id`).Scan(&userID); err != nil {
t.Fatal(err)
}
var docID string
if err := pool.QueryRow(ctx, `INSERT INTO documents (title, created_by) VALUES ('Zu vernichtendes Dokument', $1) RETURNING id`, userID).Scan(&docID); err != nil {
t.Fatal(err)
}
storageDir := t.TempDir()
driver := storage.NewLocalDriver(storageDir, []byte("secret"), "https://files.example.test")
usage := &fakeUsageReporter{}
storageSvc := storage.NewService(driver, usage, "acme")
storageKey := "documents/" + docID + "/revisions/rev-1"
content := "vertrauliches dokument"
if _, err := storageSvc.Put(ctx, storageKey, strings.NewReader(content), int64(len(content)), "text/plain"); err != nil {
t.Fatalf("testdatei ablegen: %v", err)
}
if _, err := pool.Exec(ctx, `
INSERT INTO file_revisions (document_id, revision_number, storage_key, checksum_sha256, size_bytes, mime_type, created_by)
VALUES ($1, 1, $2, 'irrelevant', $3, 'text/plain', $4)
`, docID, storageKey, len(content), userID); err != nil {
t.Fatal(err)
}
handler := DestroyCallbackHandler(pool, storageSvc)
server := httptest.NewServer(handler)
defer server.Close()
body, _ := json.Marshal(destroyRequest{
ObjectType: ObjectTypeDocument, ObjectReference: docID, DestroyedAt: time.Now().UTC().Format(time.RFC3339),
})
resp, err := http.Post(server.URL, "application/json", strings.NewReader(string(body)))
if err != nil {
t.Fatalf("post: %v", err)
}
defer func() { _ = resp.Body.Close() }()
if resp.StatusCode < 200 || resp.StatusCode >= 300 {
t.Fatalf("status = %d, want 2xx", resp.StatusCode)
}
// Nachweis: Storage-Objekt WIRKLICH geloescht (nicht nur DB-Flag).
if _, err := driver.Get(ctx, storageKey); err == nil {
t.Fatal("erwartet geloeschtes storage-objekt, konnte es aber noch lesen")
}
var deletedAt *time.Time
if err := pool.QueryRow(ctx, `SELECT deleted_at FROM documents WHERE id = $1`, docID).Scan(&deletedAt); err != nil {
t.Fatal(err)
}
if deletedAt == nil {
t.Fatal("erwartet gesetztes deleted_at nach vernichtung")
}
if len(usage.reports) != 2 || usage.reports[0] <= 0 || usage.reports[1] >= 0 {
t.Fatalf("erwartet eine positive (put) und eine negative (delete) nutzungsmeldung, habe: %v", usage.reports)
}
}
// TestRegisterOrRequeue_FailedRegistrationCreatesRequeueEntry ist die
// geforderte Pflichtpruefung 3: simulierter Nicht-2xx-Fehler bei der
// Registrierung fuehrt nachweislich zu einem Requeue-Eintrag in der
// Jobqueue, kein Absturz.
func TestRegisterOrRequeue_FailedRegistrationCreatesRequeueEntry(t *testing.T) {
pool := setupTest(t)
client := fakeRET05Server(t, true)
queue := jobqueue.NewQueue(pool, time.Minute)
queued, err := RegisterOrRequeue(context.Background(), client, queue, RegisterConfig{
ObjectType: ObjectTypeDocument, RetentionClass: "test-klasse", CallbackURL: "http://test/callback",
})
if err == nil {
t.Fatal("erwartet fehler von der fehlgeschlagenen registrierung (fuer protokollierung), habe nil")
}
if !queued {
t.Fatal("erwartet queued=true bei fehlgeschlagener registrierung")
}
var count int
if err := pool.QueryRow(context.Background(), `SELECT count(*) FROM processing_jobs WHERE job_type = $1 AND status = 'pending'`, JobTypeRegister).Scan(&count); err != nil {
t.Fatal(err)
}
if count != 1 {
t.Fatalf("erwartet genau einen requeue-eintrag in der jobqueue, habe %d", count)
}
}
// TestProcessRegisterJob_SucceedsOnRetry beweist, dass ein Requeue-Job
// vom Worker erfolgreich nachgeholt werden kann (Requeue ist kein
// Sackgassen-Zustand).
func TestProcessRegisterJob_SucceedsOnRetry(t *testing.T) {
setupTest(t)
client := fakeRET05Server(t, false)
payload, _ := json.Marshal(RegisterConfig{ObjectType: ObjectTypeDocument, RetentionClass: "test-klasse", CallbackURL: "http://test/callback"})
if err := ProcessRegisterJob(context.Background(), client, payload); err != nil {
t.Fatalf("erwartet erfolgreichen retry, habe fehler: %v", err)
}
}
-8
View File
@@ -1,8 +0,0 @@
// Package shared enthaelt Code, der von App und Worker gemeinsam genutzt
// wird (FDN-01) — Datenmodell, Storage-Zugriff etc. kommen in spaeteren
// Kacheln (FDN-02/FDN-03) hierher, dieses Paket ist bewusst noch schlank.
package shared
// Version ist die aktuelle DMS-Version, per -ldflags ueberschreibbar
// (siehe Makefile) — Platzhalter fuer echtes Versionsmanagement.
var Version = "dev"
-45
View File
@@ -1,45 +0,0 @@
// Package storage implementiert FDN-03: eine einheitliche Objekt-Storage-
// Abstraktion mit zwei austauschbaren Treibern (lokal fuer Entwicklung,
// S3-kompatibel fuer Produktion). Verschluesselung at rest ist NICHT
// Bestandteil dieser Kachel (siehe FDN-09) — dieses Paket legt Bytes
// unveraendert ab.
package storage
import (
"context"
"errors"
"io"
"time"
)
// ErrNotFound wird geliefert, wenn ein angefragtes Objekt nicht existiert
// (Akzeptanzkriterium/Pruefung 3: klarer Fehler statt treiberspezifischer
// Fehlertypen, die der Aufrufer sonst je Treiber unterschiedlich behandeln
// muesste).
var ErrNotFound = errors.New("storage: objekt nicht gefunden")
// Driver ist die EINE Schnittstelle, gegen die der Rest von DMS arbeitet
// (Akzeptanzkriterium 1). Zwei Implementierungen: LocalDriver (Entwicklung)
// und S3Driver (Produktion, S3-kompatibel).
type Driver interface {
// Put legt die Bytes aus r unter key ab und liefert die tatsaechlich
// geschriebene Groesse in Bytes.
Put(ctx context.Context, key string, r io.Reader, size int64, contentType string) (int64, error)
// Get liefert die Bytes unter key. Existiert key nicht, liefert Get
// ErrNotFound.
Get(ctx context.Context, key string) (io.ReadCloser, error)
// Delete entfernt das Objekt unter key. Existiert key nicht, liefert
// Delete ErrNotFound.
Delete(ctx context.Context, key string) error
// SignedURL liefert eine zeitlich begrenzte, signierte URL zum Lesen des
// Objekts (Akzeptanzkriterium 2: konfigurierbare Gueltigkeit ueber ttl).
SignedURL(ctx context.Context, key string, ttl time.Duration) (string, error)
}
// ObjectKey liefert das Pfadschema fuer ein Dokument/Revision INNERHALB des
// Mandanten-Buckets (Akzeptanzkriterium 3) — die Bucket-Trennung selbst ist
// Sache von Core TEN-01, hier geht es nur um den Pfad innerhalb eines
// bereits mandantenspezifischen Buckets.
func ObjectKey(documentID, revisionID string) string {
return "documents/" + documentID + "/revisions/" + revisionID
}
-113
View File
@@ -1,113 +0,0 @@
package storage
import (
"context"
"crypto/hmac"
"crypto/sha256"
"crypto/subtle"
"encoding/base64"
"errors"
"fmt"
"io"
"os"
"path/filepath"
"strconv"
"strings"
"time"
)
// ErrURLExpired wird von VerifySignedURL geliefert, wenn eine signierte URL
// nach Ablauf ihrer Gueltigkeit verwendet wird (Pruefung 2).
var ErrURLExpired = errors.New("storage: signierte url ist abgelaufen")
// ErrInvalidSignature wird geliefert, wenn die Signatur einer URL nicht zum
// Schluessel passt (manipulierte oder falsche URL).
var ErrInvalidSignature = errors.New("storage: signatur der url ist ungueltig")
// LocalDriver legt Objekte im lokalen Dateisystem ab — der Entwicklungs-
// Treiber (Akzeptanzkriterium 1), keine externe Abhaengigkeit noetig.
type LocalDriver struct {
baseDir string
signingSecret []byte
publicBaseURL string
}
// NewLocalDriver erzeugt einen LocalDriver. signingSecret authentifiziert
// die von SignedURL ausgestellten URLs (HMAC-SHA256, konstant-zeit-
// verglichen bei der Verifikation — timing-safe wie projektweite Konvention,
// siehe Core IAM-15).
func NewLocalDriver(baseDir string, signingSecret []byte, publicBaseURL string) *LocalDriver {
return &LocalDriver{baseDir: baseDir, signingSecret: signingSecret, publicBaseURL: publicBaseURL}
}
func (d *LocalDriver) path(key string) string {
return filepath.Join(d.baseDir, filepath.FromSlash(key))
}
func (d *LocalDriver) Put(ctx context.Context, key string, r io.Reader, size int64, contentType string) (int64, error) {
full := d.path(key)
if err := os.MkdirAll(filepath.Dir(full), 0o755); err != nil {
return 0, fmt.Errorf("storage: verzeichnis anlegen: %w", err)
}
f, err := os.Create(full)
if err != nil {
return 0, fmt.Errorf("storage: datei anlegen: %w", err)
}
defer func() { _ = f.Close() }()
written, err := io.Copy(f, r)
if err != nil {
return 0, fmt.Errorf("storage: schreiben: %w", err)
}
return written, nil
}
func (d *LocalDriver) Get(ctx context.Context, key string) (io.ReadCloser, error) {
f, err := os.Open(d.path(key))
if err != nil {
if os.IsNotExist(err) {
return nil, ErrNotFound
}
return nil, fmt.Errorf("storage: lesen: %w", err)
}
return f, nil
}
func (d *LocalDriver) Delete(ctx context.Context, key string) error {
if err := os.Remove(d.path(key)); err != nil {
if os.IsNotExist(err) {
return ErrNotFound
}
return fmt.Errorf("storage: loeschen: %w", err)
}
return nil
}
func (d *LocalDriver) SignedURL(ctx context.Context, key string, ttl time.Duration) (string, error) {
expiry := time.Now().Add(ttl).Unix()
sig := d.sign(key, expiry)
return fmt.Sprintf("%s/%s?exp=%d&sig=%s", strings.TrimRight(d.publicBaseURL, "/"), key, expiry, sig), nil
}
func (d *LocalDriver) sign(key string, expiry int64) string {
mac := hmac.New(sha256.New, d.signingSecret)
mac.Write([]byte(key))
mac.Write([]byte(strconv.FormatInt(expiry, 10)))
return base64.RawURLEncoding.EncodeToString(mac.Sum(nil))
}
// VerifySignedURL prueft key/expiry/sig, wie sie z.B. aus den Query-
// Parametern einer von SignedURL ausgestellten URL stammen (Pruefung 2:
// abgelaufene URL wird abgewiesen). Die eigentliche HTTP-Auslieferung ist
// nicht Bestandteil dieser Kachel (siehe DOC-01) — hier wird nur die
// Signatur-/Ablauflogik bereitgestellt und getestet.
func (d *LocalDriver) VerifySignedURL(key string, expiry int64, sig string) error {
expected := d.sign(key, expiry)
if subtle.ConstantTimeCompare([]byte(expected), []byte(sig)) != 1 {
return ErrInvalidSignature
}
if time.Now().Unix() > expiry {
return ErrURLExpired
}
return nil
}
-108
View File
@@ -1,108 +0,0 @@
package storage
import (
"bytes"
"context"
"errors"
"io"
"net/url"
"strconv"
"testing"
"time"
)
func newTestLocalDriver(t *testing.T) *LocalDriver {
t.Helper()
return NewLocalDriver(t.TempDir(), []byte("test-signing-secret"), "https://files.example.test")
}
// TestLocalDriver_RoundTrip ist Pruefung 1 fuer den lokalen Treiber:
// Upload/Download-Roundtrip.
func TestLocalDriver_RoundTrip(t *testing.T) {
d := newTestLocalDriver(t)
ctx := context.Background()
key := "documents/doc-1/revisions/rev-1"
content := []byte("hallo welt")
written, err := d.Put(ctx, key, bytes.NewReader(content), int64(len(content)), "text/plain")
if err != nil {
t.Fatalf("put: %v", err)
}
if written != int64(len(content)) {
t.Fatalf("geschriebene groesse = %d, want %d", written, len(content))
}
rc, err := d.Get(ctx, key)
if err != nil {
t.Fatalf("get: %v", err)
}
defer func() { _ = rc.Close() }()
got, err := io.ReadAll(rc)
if err != nil {
t.Fatalf("lesen: %v", err)
}
if !bytes.Equal(got, content) {
t.Fatalf("gelesener inhalt = %q, want %q", got, content)
}
}
// TestLocalDriver_MissingObject ist Pruefung 3: klarer Fehler bei
// fehlendem Objekt, sowohl fuer Get als auch Delete.
func TestLocalDriver_MissingObject(t *testing.T) {
d := newTestLocalDriver(t)
ctx := context.Background()
if _, err := d.Get(ctx, "nie-angelegt"); !errors.Is(err, ErrNotFound) {
t.Fatalf("get eines fehlenden objekts: erwartet ErrNotFound, habe %v", err)
}
if err := d.Delete(ctx, "nie-angelegt"); !errors.Is(err, ErrNotFound) {
t.Fatalf("delete eines fehlenden objekts: erwartet ErrNotFound, habe %v", err)
}
}
// TestLocalDriver_SignedURL_ExpiredIsRejected ist Pruefung 2: eine
// abgelaufene signierte URL wird abgewiesen.
func TestLocalDriver_SignedURL_ExpiredIsRejected(t *testing.T) {
d := newTestLocalDriver(t)
ctx := context.Background()
key := "documents/doc-2/revisions/rev-1"
// Gueltige, noch nicht abgelaufene URL wird akzeptiert.
urlValid, err := d.SignedURL(ctx, key, time.Hour)
if err != nil {
t.Fatalf("signedurl (gueltig): %v", err)
}
expiry, sig := parseSignedURLQuery(t, urlValid)
if err := d.VerifySignedURL(key, expiry, sig); err != nil {
t.Fatalf("gueltige url wurde abgewiesen: %v", err)
}
// Bereits abgelaufene URL (negative TTL) wird abgewiesen.
urlExpired, err := d.SignedURL(ctx, key, -time.Hour)
if err != nil {
t.Fatalf("signedurl (abgelaufen): %v", err)
}
expiredExpiry, expiredSig := parseSignedURLQuery(t, urlExpired)
if err := d.VerifySignedURL(key, expiredExpiry, expiredSig); !errors.Is(err, ErrURLExpired) {
t.Fatalf("abgelaufene url: erwartet ErrURLExpired, habe %v", err)
}
// Manipulierte Signatur wird abgewiesen.
if err := d.VerifySignedURL(key, expiry, "manipuliert"); !errors.Is(err, ErrInvalidSignature) {
t.Fatalf("manipulierte signatur: erwartet ErrInvalidSignature, habe %v", err)
}
}
func parseSignedURLQuery(t *testing.T, rawURL string) (expiry int64, sig string) {
t.Helper()
parsed, err := url.Parse(rawURL)
if err != nil {
t.Fatalf("signierte url parsen: %v (%s)", err, rawURL)
}
q := parsed.Query()
expInt, err := strconv.ParseInt(q.Get("exp"), 10, 64)
if err != nil {
t.Fatalf("exp parsen: %v", err)
}
return expInt, q.Get("sig")
}
-135
View File
@@ -1,135 +0,0 @@
package storage
import (
"bytes"
"context"
"errors"
"fmt"
"io"
"time"
"github.com/aws/aws-sdk-go-v2/aws"
"github.com/aws/aws-sdk-go-v2/config"
"github.com/aws/aws-sdk-go-v2/credentials"
"github.com/aws/aws-sdk-go-v2/service/s3"
"github.com/aws/aws-sdk-go-v2/service/s3/types"
"github.com/aws/smithy-go"
)
// S3Driver legt Objekte in einem S3-kompatiblen Objektspeicher ab — der
// Produktions-Treiber (Akzeptanzkriterium 1). Funktioniert gegen echtes
// AWS S3 UND gegen jeden S3-kompatiblen Anbieter (MinIO etc.) ueber
// endpointURL — bewusst offenes Objektformat statt Herstellerbindung
// (Produkt-DNA: "jederzeit ohne Herstellerwerkzeug lesbar").
type S3Driver struct {
client *s3.Client
bucket string
}
// NewS3Driver verbindet zu einem S3-kompatiblen Endpunkt. endpointURL leer
// laesst den AWS-SDK-Standardendpunkt (echtes AWS S3) gelten,
// usePathStyle=true ist fuer die meisten Nicht-AWS-S3-kompatiblen Anbieter
// (MinIO, etc.) noetig.
func NewS3Driver(ctx context.Context, bucket, region, endpointURL, accessKeyID, secretAccessKey string, usePathStyle bool) (*S3Driver, error) {
cfg, err := config.LoadDefaultConfig(ctx,
config.WithRegion(region),
config.WithCredentialsProvider(credentials.NewStaticCredentialsProvider(accessKeyID, secretAccessKey, "")),
)
if err != nil {
return nil, fmt.Errorf("storage: s3-konfiguration laden: %w", err)
}
client := s3.NewFromConfig(cfg, func(o *s3.Options) {
if endpointURL != "" {
o.BaseEndpoint = aws.String(endpointURL)
}
o.UsePathStyle = usePathStyle
})
return &S3Driver{client: client, bucket: bucket}, nil
}
func (d *S3Driver) Put(ctx context.Context, key string, r io.Reader, size int64, contentType string) (int64, error) {
buf, err := io.ReadAll(r)
if err != nil {
return 0, fmt.Errorf("storage: objekt vor upload lesen: %w", err)
}
_, err = d.client.PutObject(ctx, &s3.PutObjectInput{
Bucket: aws.String(d.bucket),
Key: aws.String(key),
Body: bytesReader(buf),
ContentLength: aws.Int64(int64(len(buf))),
ContentType: aws.String(contentType),
})
if err != nil {
return 0, fmt.Errorf("storage: s3-upload: %w", err)
}
return int64(len(buf)), nil
}
func (d *S3Driver) Get(ctx context.Context, key string) (io.ReadCloser, error) {
out, err := d.client.GetObject(ctx, &s3.GetObjectInput{
Bucket: aws.String(d.bucket),
Key: aws.String(key),
})
if err != nil {
if isS3NotFound(err) {
return nil, ErrNotFound
}
return nil, fmt.Errorf("storage: s3-download: %w", err)
}
return out.Body, nil
}
func (d *S3Driver) Delete(ctx context.Context, key string) error {
// S3 liefert bei DeleteObject fuer ein nicht existierendes Objekt KEINEN
// Fehler (idempotente Semantik der S3-API) — um denselben Vertrag wie
// LocalDriver (ErrNotFound bei fehlendem Objekt) zu erfuellen, wird die
// Existenz vorher explizit geprueft (Pruefung 3).
_, err := d.client.HeadObject(ctx, &s3.HeadObjectInput{Bucket: aws.String(d.bucket), Key: aws.String(key)})
if err != nil {
if isS3NotFound(err) {
return ErrNotFound
}
return fmt.Errorf("storage: s3-existenzpruefung vor loeschen: %w", err)
}
if _, err := d.client.DeleteObject(ctx, &s3.DeleteObjectInput{
Bucket: aws.String(d.bucket),
Key: aws.String(key),
}); err != nil {
return fmt.Errorf("storage: s3-loeschen: %w", err)
}
return nil
}
func (d *S3Driver) SignedURL(ctx context.Context, key string, ttl time.Duration) (string, error) {
presignClient := s3.NewPresignClient(d.client)
req, err := presignClient.PresignGetObject(ctx, &s3.GetObjectInput{
Bucket: aws.String(d.bucket),
Key: aws.String(key),
}, s3.WithPresignExpires(ttl))
if err != nil {
return "", fmt.Errorf("storage: presigned url erzeugen: %w", err)
}
return req.URL, nil
}
func bytesReader(b []byte) *bytes.Reader {
return bytes.NewReader(b)
}
// isS3NotFound erkennt sowohl den typisierten NoSuchKey-Fehler
// (GetObject) als auch den generischen "NotFound"-API-Fehlercode
// (HeadObject liefert keinen typisierten NoSuchKey, sondern einen
// generischen smithy-API-Fehler mit Code "NotFound").
func isS3NotFound(err error) bool {
var nsk *types.NoSuchKey
if errors.As(err, &nsk) {
return true
}
var apiErr smithy.APIError
if errors.As(err, &apiErr) && apiErr.ErrorCode() == "NotFound" {
return true
}
return false
}
-98
View File
@@ -1,98 +0,0 @@
package storage
import (
"bytes"
"context"
"errors"
"io"
"os"
"testing"
"time"
)
// requireS3TestEnv liefert die S3-Testkonfiguration oder ueberspringt den
// Test — dasselbe Muster wie TEST_ADMIN_DSN im Core-Modul: kein S3-
// kompatibler Speicher in dieser Umgebung verfuegbar/geprueft (siehe
// FDN-03-Pruefprotokoll), daher hier bewusst als optional markiert statt
// den Treiber ungetestet zu lassen.
func requireS3TestEnv(t *testing.T) *S3Driver {
t.Helper()
bucket := os.Getenv("TEST_S3_BUCKET")
if bucket == "" {
t.Skip("TEST_S3_BUCKET nicht gesetzt, S3-Integrationstest uebersprungen")
}
endpoint := os.Getenv("TEST_S3_ENDPOINT")
region := os.Getenv("TEST_S3_REGION")
if region == "" {
region = "us-east-1"
}
accessKey := os.Getenv("TEST_S3_ACCESS_KEY_ID")
secretKey := os.Getenv("TEST_S3_SECRET_ACCESS_KEY")
d, err := NewS3Driver(context.Background(), bucket, region, endpoint, accessKey, secretKey, true)
if err != nil {
t.Fatalf("s3-treiber aufbauen: %v", err)
}
return d
}
// TestS3Driver_RoundTrip ist Pruefung 1 fuer den S3-Treiber.
func TestS3Driver_RoundTrip(t *testing.T) {
d := requireS3TestEnv(t)
ctx := context.Background()
key := "fdn03-test/roundtrip"
content := []byte("s3 roundtrip inhalt")
t.Cleanup(func() { _ = d.Delete(ctx, key) })
if _, err := d.Put(ctx, key, bytes.NewReader(content), int64(len(content)), "text/plain"); err != nil {
t.Fatalf("put: %v", err)
}
rc, err := d.Get(ctx, key)
if err != nil {
t.Fatalf("get: %v", err)
}
defer func() { _ = rc.Close() }()
got, err := io.ReadAll(rc)
if err != nil {
t.Fatalf("lesen: %v", err)
}
if !bytes.Equal(got, content) {
t.Fatalf("gelesener inhalt = %q, want %q", got, content)
}
}
// TestS3Driver_MissingObject ist Pruefung 3 fuer den S3-Treiber.
func TestS3Driver_MissingObject(t *testing.T) {
d := requireS3TestEnv(t)
ctx := context.Background()
if _, err := d.Get(ctx, "fdn03-test/nie-angelegt"); !errors.Is(err, ErrNotFound) {
t.Fatalf("get eines fehlenden objekts: erwartet ErrNotFound, habe %v", err)
}
if err := d.Delete(ctx, "fdn03-test/nie-angelegt"); !errors.Is(err, ErrNotFound) {
t.Fatalf("delete eines fehlenden objekts: erwartet ErrNotFound, habe %v", err)
}
}
// TestS3Driver_SignedURL ist Pruefung 2 fuer den S3-Treiber: eine
// presigned URL wird erzeugt und ist innerhalb der Gueltigkeit abrufbar.
func TestS3Driver_SignedURL(t *testing.T) {
d := requireS3TestEnv(t)
ctx := context.Background()
key := "fdn03-test/signed-url"
content := []byte("presigned")
t.Cleanup(func() { _ = d.Delete(ctx, key) })
if _, err := d.Put(ctx, key, bytes.NewReader(content), int64(len(content)), "text/plain"); err != nil {
t.Fatalf("put: %v", err)
}
url, err := d.SignedURL(ctx, key, time.Minute)
if err != nil {
t.Fatalf("signedurl: %v", err)
}
if url == "" {
t.Fatal("erwartet nicht-leere presigned url")
}
}
-59
View File
@@ -1,59 +0,0 @@
package storage
import (
"context"
"fmt"
"io"
"time"
)
// Service verbindet einen Driver mit der Nutzungsmeldung an Core
// (Akzeptanzkriterium 4) — jeder Schreib-/Loeschvorgang ueber Service loest
// GENAU EINE Meldung mit der tatsaechlich geschriebenen/geloeschten
// Objektgroesse aus. Repository-/Handler-Code (spaetere Kacheln, z.B.
// DOC-01) ruft ausschliesslich Service auf, nie einen Driver direkt — das
// verhindert einen Schreibpfad, der die Nutzungsmeldung vergisst.
type Service struct {
driver Driver
usage UsageReporter
tenantSlug string
}
func NewService(driver Driver, usage UsageReporter, tenantSlug string) *Service {
return &Service{driver: driver, usage: usage, tenantSlug: tenantSlug}
}
// Put legt das Objekt ab und meldet die geschriebene Groesse als positives
// Delta (Pruefung 4).
func (s *Service) Put(ctx context.Context, key string, r io.Reader, size int64, contentType string) (int64, error) {
written, err := s.driver.Put(ctx, key, r, size, contentType)
if err != nil {
return 0, err
}
if err := s.usage.Report(ctx, s.tenantSlug, UsageMetric, written); err != nil {
return written, fmt.Errorf("storage: objekt gespeichert, aber nutzungsmeldung fehlgeschlagen: %w", err)
}
return written, nil
}
func (s *Service) Get(ctx context.Context, key string) (io.ReadCloser, error) {
return s.driver.Get(ctx, key)
}
// Delete entfernt das Objekt und meldet dessen Groesse als negatives Delta
// (Pruefung 4) — dafuer muss der Aufrufer die Groesse kennen (z.B. aus
// file_revisions.size_bytes, FDN-02), da Delete selbst die Groesse eines
// bereits geloeschten Objekts nicht mehr ermitteln kann.
func (s *Service) Delete(ctx context.Context, key string, sizeBytes int64) error {
if err := s.driver.Delete(ctx, key); err != nil {
return err
}
if err := s.usage.Report(ctx, s.tenantSlug, UsageMetric, -sizeBytes); err != nil {
return fmt.Errorf("storage: objekt geloescht, aber nutzungsmeldung fehlgeschlagen: %w", err)
}
return nil
}
func (s *Service) SignedURL(ctx context.Context, key string, ttl time.Duration) (string, error) {
return s.driver.SignedURL(ctx, key, ttl)
}
-89
View File
@@ -1,89 +0,0 @@
package storage
import (
"bytes"
"context"
"sync"
"testing"
)
// fakeUsageReporter zeichnet jeden Report-Aufruf auf, damit Tests
// nachweisen koennen, dass Service tatsaechlich meldet (Akzeptanzkriterium
// 4 / Pruefung 4) — ohne echten HTTP-Aufruf gegen Core.
type fakeUsageReporter struct {
mu sync.Mutex
calls []reportCall
failOn int // wenn >0, schlaegt der reportCall-te Aufruf fehl
}
type reportCall struct {
tenantSlug string
metric string
delta int64
}
func (f *fakeUsageReporter) Report(ctx context.Context, tenantSlug, metric string, delta int64) error {
f.mu.Lock()
defer f.mu.Unlock()
f.calls = append(f.calls, reportCall{tenantSlug, metric, delta})
if f.failOn > 0 && len(f.calls) == f.failOn {
return context.DeadlineExceeded
}
return nil
}
// TestService_PutReportsPositiveDelta ist Pruefung 4 (Schreibvorgang):
// Melde-Aufruf an Core wird bei Put ausgeloest, mit korrekter Groesse.
func TestService_PutReportsPositiveDelta(t *testing.T) {
driver := NewLocalDriver(t.TempDir(), []byte("secret"), "https://files.example.test")
usage := &fakeUsageReporter{}
svc := NewService(driver, usage, "acme")
content := []byte("zwoelf bytes")
if _, err := svc.Put(context.Background(), "documents/d1/revisions/r1", bytes.NewReader(content), int64(len(content)), "text/plain"); err != nil {
t.Fatalf("put: %v", err)
}
if len(usage.calls) != 1 {
t.Fatalf("erwartet 1 nutzungsmeldung, habe %d", len(usage.calls))
}
call := usage.calls[0]
if call.tenantSlug != "acme" || call.metric != UsageMetric || call.delta != int64(len(content)) {
t.Fatalf("unerwarteter meldungsinhalt: %+v", call)
}
}
// TestService_DeleteReportsNegativeDelta ist Pruefung 4 (Loeschvorgang).
func TestService_DeleteReportsNegativeDelta(t *testing.T) {
driver := NewLocalDriver(t.TempDir(), []byte("secret"), "https://files.example.test")
usage := &fakeUsageReporter{}
svc := NewService(driver, usage, "acme")
ctx := context.Background()
key := "documents/d2/revisions/r1"
if _, err := svc.Put(ctx, key, bytes.NewReader([]byte("abc")), 3, "text/plain"); err != nil {
t.Fatalf("put: %v", err)
}
if err := svc.Delete(ctx, key, 3); err != nil {
t.Fatalf("delete: %v", err)
}
if len(usage.calls) != 2 {
t.Fatalf("erwartet 2 nutzungsmeldungen (put+delete), habe %d", len(usage.calls))
}
del := usage.calls[1]
if del.delta != -3 {
t.Fatalf("delete-delta = %d, want -3", del.delta)
}
}
// TestService_GetMissingObjectReturnsClearError ist Pruefung 3 auf
// Service-Ebene.
func TestService_GetMissingObjectReturnsClearError(t *testing.T) {
driver := NewLocalDriver(t.TempDir(), []byte("secret"), "https://files.example.test")
svc := NewService(driver, &fakeUsageReporter{}, "acme")
if _, err := svc.Get(context.Background(), "nie-angelegt"); err == nil {
t.Fatal("get eines fehlenden objekts haette einen fehler liefern muessen")
}
}
-75
View File
@@ -1,75 +0,0 @@
package storage
import (
"bytes"
"context"
"encoding/json"
"fmt"
"net/http"
)
// UsageMetric ist der Metrikname, unter dem Core (internal/usage, LIC-05)
// den Speicherverbrauch je Mandant fuehrt — muss exakt
// internal/usage.StorageBytesMetric aus dem NEXARCH-Core-Modul entsprechen
// (Core kann von DMS als eigenem Go-Modul nicht importiert werden, daher
// hier als Konstante gespiegelt statt importiert).
const UsageMetric = "storage_bytes"
// UsageReporter meldet Speicherverbrauchsaenderungen an Core (Akzeptanz-
// kriterium 4). Schmale Schnittstelle, damit Tests einen Fake statt eines
// echten HTTP-Aufrufs einsetzen koennen.
type UsageReporter interface {
Report(ctx context.Context, tenantSlug, metric string, delta int64) error
}
// usageDeltaDTO entspricht Core internal/resync.usageDeltaDTO
// (JSON-Vertrag: tenant_slug/metric/delta) — dieselbe Struktur, hier
// gespiegelt, da DMS Cores internal/-Pakete nicht importieren kann.
type usageDeltaDTO struct {
TenantSlug string `json:"tenant_slug"`
Metric string `json:"metric"`
Delta int64 `json:"delta"`
}
// HTTPUsageReporter meldet ueber Cores Resync-Nutzungs-Endpunkt
// (internal/resync.Handler.UsageHandler, API-06), authentifiziert ueber
// dasselbe Service-Credential-Verfahren wie jeder andere Modul-Core-Aufruf
// (API-02).
type HTTPUsageReporter struct {
endpointURL string
clientID string
clientSecret string
httpClient *http.Client
}
func NewHTTPUsageReporter(endpointURL, clientID, clientSecret string, httpClient *http.Client) *HTTPUsageReporter {
if httpClient == nil {
httpClient = http.DefaultClient
}
return &HTTPUsageReporter{endpointURL: endpointURL, clientID: clientID, clientSecret: clientSecret, httpClient: httpClient}
}
func (r *HTTPUsageReporter) Report(ctx context.Context, tenantSlug, metric string, delta int64) error {
body, err := json.Marshal([]usageDeltaDTO{{TenantSlug: tenantSlug, Metric: metric, Delta: delta}})
if err != nil {
return fmt.Errorf("storage: nutzungsmeldung serialisieren: %w", err)
}
req, err := http.NewRequestWithContext(ctx, http.MethodPost, r.endpointURL, bytes.NewReader(body))
if err != nil {
return fmt.Errorf("storage: nutzungsmeldungs-anfrage aufbauen: %w", err)
}
req.Header.Set("Content-Type", "application/json")
req.Header.Set("X-Nexarch-Client-Id", r.clientID)
req.Header.Set("X-Nexarch-Client-Secret", r.clientSecret)
resp, err := r.httpClient.Do(req)
if err != nil {
return fmt.Errorf("storage: nutzungsmeldung senden: %w", err)
}
defer func() { _ = resp.Body.Close() }()
if resp.StatusCode != http.StatusOK {
return fmt.Errorf("storage: nutzungsmeldung von core abgelehnt: status %d", resp.StatusCode)
}
return nil
}
-62
View File
@@ -1,62 +0,0 @@
package storage
import (
"context"
"encoding/json"
"net/http"
"net/http/httptest"
"testing"
)
// TestHTTPUsageReporter_SendsCorrectContractToCore ist der Nachweis, dass
// HTTPUsageReporter exakt den Vertrag von Core internal/resync.Handler.
// UsageHandler bedient (Service-Credential-Header, JSON-Feldnamen) — echte
// Vernetzung gegen einen laufenden Core-Prozess ist nicht Teil dieses
// Tests (internal/resync.Handler ist in Core aktuell in keinem cmd/*/
// main.go verdrahtet, siehe FDN-03-Pruefprotokoll), daher hier gegen einen
// httptest-Server geprueft, der denselben Vertrag nachbildet.
func TestHTTPUsageReporter_SendsCorrectContractToCore(t *testing.T) {
var gotClientID, gotClientSecret string
var gotBody []usageDeltaDTO
srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
gotClientID = r.Header.Get("X-Nexarch-Client-Id")
gotClientSecret = r.Header.Get("X-Nexarch-Client-Secret")
if err := json.NewDecoder(r.Body).Decode(&gotBody); err != nil {
t.Errorf("anfrage-koerper dekodieren: %v", err)
}
w.Header().Set("Content-Type", "application/json")
_ = json.NewEncoder(w).Encode(map[string]int{"applied": len(gotBody)})
}))
defer srv.Close()
reporter := NewHTTPUsageReporter(srv.URL, "dms-service-client", "dms-service-secret", nil)
if err := reporter.Report(context.Background(), "acme", UsageMetric, 4096); err != nil {
t.Fatalf("report: %v", err)
}
if gotClientID != "dms-service-client" || gotClientSecret != "dms-service-secret" {
t.Fatalf("service-credential-header falsch: id=%q secret=%q", gotClientID, gotClientSecret)
}
if len(gotBody) != 1 {
t.Fatalf("erwartet 1 delta im koerper, habe %d", len(gotBody))
}
if gotBody[0].TenantSlug != "acme" || gotBody[0].Metric != UsageMetric || gotBody[0].Delta != 4096 {
t.Fatalf("unerwarteter delta-inhalt: %+v", gotBody[0])
}
}
// TestHTTPUsageReporter_RejectsNonOKStatus prueft, dass ein von Core
// abgelehnter Aufruf (z.B. ungueltiges Service-Credential) als Fehler
// zurueckgegeben wird, statt stillschweigend zu verlieren.
func TestHTTPUsageReporter_RejectsNonOKStatus(t *testing.T) {
srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
http.Error(w, "ungueltiges service-credential", http.StatusUnauthorized)
}))
defer srv.Close()
reporter := NewHTTPUsageReporter(srv.URL, "x", "y", nil)
if err := reporter.Report(context.Background(), "acme", UsageMetric, 1); err == nil {
t.Fatal("erwartet fehler bei abgelehnter nutzungsmeldung, habe nil")
}
}
-157
View File
@@ -1,157 +0,0 @@
package upload
import (
"bytes"
"context"
"crypto/sha256"
"encoding/hex"
"fmt"
"io"
"github.com/jackc/pgx/v5/pgxpool"
dmscrypto "gitea.perlbach24.de/scripte/nexarch/dms/internal/crypto"
"gitea.perlbach24.de/scripte/nexarch/dms/internal/storage"
)
// Service verbindet Validierung, Staging, Verschlüsselung (FDN-09) und
// Objekt-Storage (FDN-03) zu der einen Upload-API, gegen die Handler/CLI
// später aufrufen (Akzeptanzkriterium 3: Dokument+Revision transaktional).
type Service struct {
pool *pgxpool.Pool
sessions *SessionStore
staging *Staging
storageSvc *storage.Service
cryptoSvc *dmscrypto.Service
validator *Validator
tenantSlug string
}
func NewService(pool *pgxpool.Pool, sessions *SessionStore, staging *Staging, storageSvc *storage.Service, cryptoSvc *dmscrypto.Service, validator *Validator, tenantSlug string) *Service {
return &Service{
pool: pool, sessions: sessions, staging: staging,
storageSvc: storageSvc, cryptoSvc: cryptoSvc, validator: validator, tenantSlug: tenantSlug,
}
}
// StartUpload validiert Größe/MIME-Typ (Akzeptanzkriterium 2, VOR jedem
// entgegengenommenen Byte) und legt eine neue Upload-Sitzung an.
func (s *Service) StartUpload(ctx context.Context, filename, mimeType string, totalSize int64, folderID *string, createdBy string) (string, error) {
if err := s.validator.ValidateMimeType(mimeType); err != nil {
return "", err
}
if err := s.validator.ValidateSize(totalSize); err != nil {
return "", err
}
return s.sessions.Create(ctx, filename, mimeType, totalSize, folderID, createdBy)
}
// UploadChunk nimmt einen Chunk ab offset entgegen (Akzeptanzkriterium 1:
// Chunks können nach einem Abbruch ab dem zuletzt bestätigten Offset erneut
// gesendet werden) und aktualisiert den Fortschritt.
func (s *Service) UploadChunk(ctx context.Context, sessionID string, offset int64, chunk io.Reader) (int64, error) {
newTotal, err := s.staging.WriteChunk(sessionID, offset, chunk)
if err != nil {
return 0, err
}
if err := s.sessions.RecordChunk(ctx, sessionID, newTotal); err != nil {
return 0, err
}
return newTotal, nil
}
// Status liefert den aktuellen Fortschritt — der Client fragt dies nach
// einem Verbindungsabbruch ab, um zu wissen, ab welchem Offset er
// fortsetzen muss (Akzeptanzkriterium 1).
func (s *Service) Status(ctx context.Context, sessionID string) (*Session, error) {
return s.sessions.Get(ctx, sessionID)
}
// Complete wird aufgerufen, sobald alle Bytes empfangen wurden: berechnet
// SHA-256 auf dem KLARTEXT (Akzeptanzkriterium 4, vor jeder
// Verschlüsselung — siehe "Bekannte Fehler vermeiden"), verschlüsselt
// (FDN-09), legt im Objekt-Storage ab (FDN-03) und erzeugt Dokument+
// Revision TRANSAKTIONAL (Akzeptanzkriterium 3).
func (s *Service) Complete(ctx context.Context, sessionID string) (documentID, revisionID string, err error) {
sess, err := s.sessions.Get(ctx, sessionID)
if err != nil {
return "", "", err
}
if sess.BytesReceived < sess.TotalSize {
return "", "", fmt.Errorf("upload: sitzung %q unvollstaendig (%d von %d bytes)", sessionID, sess.BytesReceived, sess.TotalSize)
}
f, err := s.staging.OpenForRead(sessionID)
if err != nil {
return "", "", err
}
defer func() { _ = f.Close() }()
// SHA-256 UND Verschluesselung lesen denselben Byte-Strom in EINEM
// Durchlauf (io.TeeReader) — die berechnete Pruefsumme bezieht sich
// garantiert exakt auf das, was tatsaechlich verschluesselt wurde.
hasher := sha256.New()
tee := io.TeeReader(f, hasher)
env, err := s.cryptoSvc.Seal(ctx, s.tenantSlug, tee)
if err != nil {
return "", "", fmt.Errorf("upload: verschluesseln: %w", err)
}
checksumHex := hex.EncodeToString(hasher.Sum(nil))
ciphertext, err := io.ReadAll(env.Ciphertext)
if err != nil {
return "", "", fmt.Errorf("upload: chiffretext lesen: %w", err)
}
tx, err := s.pool.Begin(ctx)
if err != nil {
return "", "", fmt.Errorf("upload: transaktion starten: %w", err)
}
defer func() { _ = tx.Rollback(ctx) }()
if err := tx.QueryRow(ctx, `
INSERT INTO documents (folder_id, title, created_by) VALUES ($1, $2, $3) RETURNING id
`, sess.FolderID, sess.Filename, sess.CreatedBy).Scan(&documentID); err != nil {
return "", "", fmt.Errorf("upload: dokument anlegen: %w", err)
}
objectKey := storage.ObjectKey(documentID, "pending")
if err := tx.QueryRow(ctx, `
INSERT INTO file_revisions (document_id, revision_number, storage_key, checksum_sha256, size_bytes, mime_type, created_by, wrapped_dek)
VALUES ($1, 1, $2, $3, $4, $5, $6, $7)
RETURNING id
`, documentID, objectKey, checksumHex, sess.TotalSize, sess.MimeType, sess.CreatedBy, env.WrappedDEK).Scan(&revisionID); err != nil {
return "", "", fmt.Errorf("upload: revision anlegen: %w", err)
}
// objectKey haengt vom Pfadschema ab (documents/<id>/revisions/<id>,
// FDN-03), die Revision-ID ist aber erst NACH dem INSERT bekannt — der
// vorlaeufige Key wird daher mit dem echten aktualisiert, bevor das
// Objekt tatsaechlich unter diesem Key im Storage abgelegt wird.
finalKey := storage.ObjectKey(documentID, revisionID)
if _, err := tx.Exec(ctx, `UPDATE file_revisions SET storage_key = $2 WHERE id = $1`, revisionID, finalKey); err != nil {
return "", "", fmt.Errorf("upload: storage-key aktualisieren: %w", err)
}
if _, err := tx.Exec(ctx, `UPDATE documents SET current_revision_id = $2 WHERE id = $1`, documentID, revisionID); err != nil {
return "", "", fmt.Errorf("upload: aktuelle revision setzen: %w", err)
}
if _, err := s.storageSvc.Put(ctx, finalKey, bytes.NewReader(ciphertext), int64(len(ciphertext)), "application/octet-stream"); err != nil {
return "", "", fmt.Errorf("upload: objekt ablegen: %w", err)
}
if err := tx.Commit(ctx); err != nil {
return "", "", fmt.Errorf("upload: transaktion committen: %w", err)
}
if err := s.sessions.MarkCompleted(ctx, sessionID, documentID); err != nil {
return "", "", fmt.Errorf("upload: sitzung als abgeschlossen markieren: %w", err)
}
if err := s.staging.Remove(sessionID); err != nil {
return "", "", fmt.Errorf("upload: staging-datei aufraeumen: %w", err)
}
return documentID, revisionID, nil
}
-299
View File
@@ -1,299 +0,0 @@
package upload
import (
"bytes"
"context"
"crypto/rand"
"crypto/sha256"
"encoding/hex"
"errors"
"fmt"
"os"
"sync"
"testing"
"time"
"github.com/jackc/pgx/v5/pgxpool"
dmscrypto "gitea.perlbach24.de/scripte/nexarch/dms/internal/crypto"
"gitea.perlbach24.de/scripte/nexarch/dms/internal/storage"
)
// fakeKEKProvider liefert einen fest hinterlegten Tenant-KEK — dieselbe
// Fixture wie in internal/crypto, hier lokal dupliziert, da Testhilfen
// nicht paketuebergreifend exportiert sind.
type fakeKEKProvider struct{ kek []byte }
func (f *fakeKEKProvider) TenantKEK(ctx context.Context, tenantSlug string) ([]byte, error) {
return f.kek, nil
}
func setupServiceTest(t *testing.T) (*Service, *pgxpool.Pool, string) {
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() })
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,
status TEXT NOT NULL DEFAULT 'active', created_at TIMESTAMPTZ NOT NULL DEFAULT now(), updated_at TIMESTAMPTZ NOT NULL DEFAULT now()
);
CREATE TABLE IF NOT EXISTS folders (
id UUID PRIMARY KEY DEFAULT gen_random_uuid(), parent_folder_id UUID REFERENCES folders(id) ON DELETE CASCADE,
name TEXT NOT NULL, created_by UUID NOT NULL REFERENCES users(id),
created_at TIMESTAMPTZ NOT NULL DEFAULT now(), updated_at TIMESTAMPTZ NOT NULL DEFAULT now()
);
CREATE TABLE IF NOT EXISTS documents (
id UUID PRIMARY KEY DEFAULT gen_random_uuid(), folder_id UUID REFERENCES folders(id) ON DELETE SET NULL,
title TEXT NOT NULL, current_revision_id UUID, created_by UUID NOT NULL REFERENCES users(id),
created_at TIMESTAMPTZ NOT NULL DEFAULT now(), updated_at TIMESTAMPTZ NOT NULL DEFAULT now(), deleted_at TIMESTAMPTZ
);
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,
revision_number INT NOT NULL, storage_key TEXT NOT NULL, checksum_sha256 TEXT NOT NULL,
size_bytes BIGINT NOT NULL, mime_type TEXT NOT NULL, created_by UUID NOT NULL REFERENCES users(id),
created_at TIMESTAMPTZ NOT NULL DEFAULT now(), wrapped_dek BYTEA,
UNIQUE (document_id, revision_number)
);
CREATE TABLE IF NOT EXISTS upload_sessions (
id UUID PRIMARY KEY DEFAULT gen_random_uuid(), filename TEXT NOT NULL, mime_type TEXT NOT NULL,
total_size BIGINT NOT NULL CHECK (total_size > 0), bytes_received BIGINT NOT NULL DEFAULT 0 CHECK (bytes_received >= 0),
folder_id UUID REFERENCES folders(id) ON DELETE SET NULL, created_by UUID NOT NULL REFERENCES users(id),
status TEXT NOT NULL DEFAULT 'uploading' CHECK (status IN ('uploading','completed','aborted')),
document_id UUID REFERENCES documents(id),
created_at TIMESTAMPTZ NOT NULL DEFAULT now(), updated_at TIMESTAMPTZ NOT NULL DEFAULT now()
);
`); err != nil {
t.Fatalf("schema: %v", err)
}
t.Cleanup(func() {
ctx := context.Background()
_, _ = pool.Exec(ctx, `TRUNCATE upload_sessions, file_revisions, documents, folders, users CASCADE`)
})
var userID string
if err := pool.QueryRow(ctx, `
INSERT INTO users (email, name) VALUES ($1, 'Test-Benutzer') RETURNING id
`, fmt.Sprintf("upload-test-%d@example.test", time.Now().UnixNano())).Scan(&userID); err != nil {
t.Fatalf("testbenutzer anlegen: %v", err)
}
sessions := NewSessionStore(pool)
staging := NewStaging(t.TempDir())
storageSvc := storage.NewService(storage.NewLocalDriver(t.TempDir(), []byte("secret"), "https://files.example.test"), noopUsageReporter{}, "acme")
cryptoSvc := dmscrypto.NewService(&fakeKEKProvider{kek: bytes.Repeat([]byte{0x11}, dmscrypto.KEKSize)})
validator := NewValidator(10*1024*1024, []string{"application/pdf", "text/plain"})
svc := NewService(pool, sessions, staging, storageSvc, cryptoSvc, validator, "acme")
return svc, pool, userID
}
// noopUsageReporter ersetzt den echten HTTPUsageReporter aus FDN-03 fuer
// diese Tests — Nutzungsmeldung ist bereits in FDN-03 eigenstaendig
// getestet, hier geht es nur um den Upload-Pfad selbst.
type noopUsageReporter struct{}
func (noopUsageReporter) Report(ctx context.Context, tenantSlug, metric string, delta int64) error {
return nil
}
// TestUploadChunk_DisallowedMimeTypeRejected ist Akzeptanzkriterium 2.
func TestUploadChunk_DisallowedMimeTypeRejected(t *testing.T) {
svc, _, userID := setupServiceTest(t)
ctx := context.Background()
_, err := svc.StartUpload(ctx, "schadcode.exe", "application/x-msdownload", 100, nil, userID)
if !errors.Is(err, ErrDisallowedMimeType) {
t.Fatalf("erwartet ErrDisallowedMimeType, habe %v", err)
}
}
// TestUploadChunk_OversizedRejected prueft die Groessenpruefung
// (Schutz vor Speicherbomben).
func TestUploadChunk_OversizedRejected(t *testing.T) {
svc, _, userID := setupServiceTest(t)
ctx := context.Background()
_, err := svc.StartUpload(ctx, "riesig.pdf", "application/pdf", 100*1024*1024, nil, userID)
if !errors.Is(err, ErrFileTooLarge) {
t.Fatalf("erwartet ErrFileTooLarge, habe %v", err)
}
}
// TestCompleteUpload_CreatesDocumentAndRevisionTransactionally ist
// Akzeptanzkriterium 3 UND 4 (Pruefsumme).
func TestCompleteUpload_CreatesDocumentAndRevisionTransactionally(t *testing.T) {
svc, pool, userID := setupServiceTest(t)
ctx := context.Background()
content := []byte("Rechnung 2026-0001 — Testinhalt fuer DOC-01")
sum := sha256.Sum256(content)
wantChecksum := hex.EncodeToString(sum[:])
sessionID, err := svc.StartUpload(ctx, "rechnung.pdf", "application/pdf", int64(len(content)), nil, userID)
if err != nil {
t.Fatalf("startupload: %v", err)
}
if _, err := svc.UploadChunk(ctx, sessionID, 0, bytes.NewReader(content)); err != nil {
t.Fatalf("uploadchunk: %v", err)
}
documentID, revisionID, err := svc.Complete(ctx, sessionID)
if err != nil {
t.Fatalf("complete: %v", err)
}
if documentID == "" || revisionID == "" {
t.Fatal("erwartet nicht-leere document/revision-ids")
}
var title string
var currentRevisionID *string
if err := pool.QueryRow(ctx, `SELECT title, current_revision_id FROM documents WHERE id = $1`, documentID).Scan(&title, &currentRevisionID); err != nil {
t.Fatalf("dokument lesen: %v", err)
}
if title != "rechnung.pdf" {
t.Fatalf("titel = %q, want %q", title, "rechnung.pdf")
}
if currentRevisionID == nil || *currentRevisionID != revisionID {
t.Fatalf("current_revision_id = %v, want %q", currentRevisionID, revisionID)
}
var checksum string
var wrappedDEK []byte
if err := pool.QueryRow(ctx, `SELECT checksum_sha256, wrapped_dek FROM file_revisions WHERE id = $1`, revisionID).Scan(&checksum, &wrappedDEK); err != nil {
t.Fatalf("revision lesen: %v", err)
}
if checksum != wantChecksum {
t.Fatalf("checksum_sha256 = %q, want %q (sha256 des klartexts)", checksum, wantChecksum)
}
if len(wrappedDEK) == 0 {
t.Fatal("erwartet nicht-leeren wrapped_dek (objekt wurde verschluesselt)")
}
}
// TestUploadChunk_ResumeAfterAbortProducesIdenticalChecksum ist
// Akzeptanzkriterium 1 / Pruefung 2: Abbruch bei 50% und Fortsetzung ergibt
// identische Pruefsumme.
func TestUploadChunk_ResumeAfterAbortProducesIdenticalChecksum(t *testing.T) {
svc, _, userID := setupServiceTest(t)
ctx := context.Background()
content := make([]byte, 200*1024) // 200 KiB
if _, err := rand.Read(content); err != nil {
t.Fatalf("zufallsinhalt erzeugen: %v", err)
}
sum := sha256.Sum256(content)
wantChecksum := hex.EncodeToString(sum[:])
sessionID, err := svc.StartUpload(ctx, "grosse-datei.pdf", "application/pdf", int64(len(content)), nil, userID)
if err != nil {
t.Fatalf("startupload: %v", err)
}
half := len(content) / 2
if _, err := svc.UploadChunk(ctx, sessionID, 0, bytes.NewReader(content[:half])); err != nil {
t.Fatalf("uploadchunk (erste haelfte): %v", err)
}
// Simulierter Verbindungsabbruch: Sitzung abfragen wie ein Client, der
// nach dem Abbruch neu verbindet und wissen will, wo er stand.
status, err := svc.Status(ctx, sessionID)
if err != nil {
t.Fatalf("status: %v", err)
}
if status.BytesReceived != int64(half) {
t.Fatalf("bytes_received nach abbruch = %d, want %d", status.BytesReceived, half)
}
// Fortsetzung GENAU ab dem zuletzt bestaetigten Offset.
if _, err := svc.UploadChunk(ctx, sessionID, int64(half), bytes.NewReader(content[half:])); err != nil {
t.Fatalf("uploadchunk (fortsetzung): %v", err)
}
_, revisionID, err := svc.Complete(ctx, sessionID)
if err != nil {
t.Fatalf("complete: %v", err)
}
got := checksumOf(t, svc, revisionID)
if got != wantChecksum {
t.Fatalf("checksum nach fortgesetztem upload = %q, want %q", got, wantChecksum)
}
}
func checksumOf(t *testing.T, svc *Service, revisionID string) string {
t.Helper()
var checksum string
if err := svc.pool.QueryRow(context.Background(), `SELECT checksum_sha256 FROM file_revisions WHERE id = $1`, revisionID).Scan(&checksum); err != nil {
t.Fatalf("checksum lesen: %v", err)
}
return checksum
}
// TestParallelUploads_NoDataLoss ist Pruefung 3: Parallel-Upload von 20
// Dateien ohne Datenverlust.
func TestParallelUploads_NoDataLoss(t *testing.T) {
svc, _, userID := setupServiceTest(t)
ctx := context.Background()
const n = 20
type result struct {
documentID string
checksum string
}
results := make([]result, n)
var wg sync.WaitGroup
for i := 0; i < n; i++ {
wg.Add(1)
go func(idx int) {
defer wg.Done()
content := []byte(fmt.Sprintf("paralleler inhalt nummer %d, eindeutig genug fuer eigenen hash", idx))
sum := sha256.Sum256(content)
wantChecksum := hex.EncodeToString(sum[:])
sessionID, err := svc.StartUpload(ctx, fmt.Sprintf("datei-%d.pdf", idx), "application/pdf", int64(len(content)), nil, userID)
if err != nil {
t.Errorf("startupload %d: %v", idx, err)
return
}
if _, err := svc.UploadChunk(ctx, sessionID, 0, bytes.NewReader(content)); err != nil {
t.Errorf("uploadchunk %d: %v", idx, err)
return
}
docID, revID, err := svc.Complete(ctx, sessionID)
if err != nil {
t.Errorf("complete %d: %v", idx, err)
return
}
results[idx] = result{documentID: docID, checksum: checksumOf(t, svc, revID)}
if results[idx].checksum != wantChecksum {
t.Errorf("upload %d: checksum = %q, want %q", idx, results[idx].checksum, wantChecksum)
}
}(i)
}
wg.Wait()
seen := map[string]bool{}
for i, r := range results {
if r.documentID == "" {
t.Fatalf("upload %d lieferte keine document-id (fehlgeschlagen)", i)
}
if seen[r.documentID] {
t.Fatalf("document-id %q doppelt vergeben - datenverlust/kollision", r.documentID)
}
seen[r.documentID] = true
}
if len(seen) != n {
t.Fatalf("erwartet %d eindeutige dokumente, habe %d", n, len(seen))
}
}
-105
View File
@@ -1,105 +0,0 @@
package upload
import (
"context"
"errors"
"fmt"
"github.com/jackc/pgx/v5"
"github.com/jackc/pgx/v5/pgxpool"
)
// ErrSessionNotFound wird geliefert, wenn eine angefragte Upload-Sitzung
// nicht existiert.
var ErrSessionNotFound = errors.New("upload: sitzung nicht gefunden")
// Session ist der Zustand eines laufenden oder abgeschlossenen Uploads
// (Akzeptanzkriterium 1: Grundlage fuer Fortsetzung nach Abbruch).
type Session struct {
ID string
Filename string
MimeType string
TotalSize int64
BytesReceived int64
FolderID *string
CreatedBy string
Status string
DocumentID *string
}
const (
StatusUploading = "uploading"
StatusCompleted = "completed"
StatusAborted = "aborted"
)
// SessionStore verwaltet upload_sessions.
type SessionStore struct {
pool *pgxpool.Pool
}
func NewSessionStore(pool *pgxpool.Pool) *SessionStore {
return &SessionStore{pool: pool}
}
func (s *SessionStore) Create(ctx context.Context, filename, mimeType string, totalSize int64, folderID *string, createdBy string) (string, error) {
var id string
err := s.pool.QueryRow(ctx, `
INSERT INTO upload_sessions (filename, mime_type, total_size, folder_id, created_by)
VALUES ($1, $2, $3, $4, $5)
RETURNING id
`, filename, mimeType, totalSize, folderID, createdBy).Scan(&id)
if err != nil {
return "", fmt.Errorf("upload: sitzung anlegen: %w", err)
}
return id, nil
}
func (s *SessionStore) Get(ctx context.Context, sessionID string) (*Session, error) {
var sess Session
err := s.pool.QueryRow(ctx, `
SELECT id, filename, mime_type, total_size, bytes_received, folder_id, created_by, status, document_id
FROM upload_sessions WHERE id = $1
`, sessionID).Scan(&sess.ID, &sess.Filename, &sess.MimeType, &sess.TotalSize, &sess.BytesReceived,
&sess.FolderID, &sess.CreatedBy, &sess.Status, &sess.DocumentID)
if err != nil {
if errors.Is(err, pgx.ErrNoRows) {
return nil, ErrSessionNotFound
}
return nil, fmt.Errorf("upload: sitzung lesen: %w", err)
}
return &sess, nil
}
// RecordChunk aktualisiert bytes_received auf das Maximum aus dem
// bisherigen und dem neu gemeldeten Wert (Akzeptanzkriterium 1: ein
// erneut zugestellter/uebersprungener Chunk darf den Fortschritt nie
// zurueckdrehen — GREATEST statt blindem Ueberschreiben).
func (s *SessionStore) RecordChunk(ctx context.Context, sessionID string, bytesReceivedNow int64) error {
tag, err := s.pool.Exec(ctx, `
UPDATE upload_sessions
SET bytes_received = GREATEST(bytes_received, $2), updated_at = now()
WHERE id = $1 AND status = 'uploading'
`, sessionID, bytesReceivedNow)
if err != nil {
return fmt.Errorf("upload: fortschritt aktualisieren: %w", err)
}
if tag.RowsAffected() == 0 {
return ErrSessionNotFound
}
return nil
}
func (s *SessionStore) MarkCompleted(ctx context.Context, sessionID, documentID string) error {
tag, err := s.pool.Exec(ctx, `
UPDATE upload_sessions SET status = 'completed', document_id = $2, updated_at = now()
WHERE id = $1
`, sessionID, documentID)
if err != nil {
return fmt.Errorf("upload: sitzung abschliessen: %w", err)
}
if tag.RowsAffected() == 0 {
return ErrSessionNotFound
}
return nil
}
-68
View File
@@ -1,68 +0,0 @@
package upload
import (
"fmt"
"io"
"os"
"path/filepath"
)
// Staging haelt Uploads waehrend der Uebertragung lokal auf der Platte —
// getrennt von der finalen Objekt-Storage-Ablage (FDN-03), die erst nach
// vollstaendigem, verifiziertem Empfang beschrieben wird. Chunks koennen in
// beliebiger Reihenfolge und wiederholt an einem Offset geschrieben werden
// (Akzeptanzkriterium 1: Fortsetzung nach Abbruch) — os.File.WriteAt ist
// dafuer das richtige Werkzeug, kein sequenzielles Append.
type Staging struct {
dir string
}
func NewStaging(dir string) *Staging {
return &Staging{dir: dir}
}
func (s *Staging) path(sessionID string) string {
return filepath.Join(s.dir, sessionID+".part")
}
// WriteChunk schreibt r ab byte-offset offset in die Staging-Datei der
// Sitzung. Ein bereits vorher (auch teilweise) geschriebener Bereich wird
// beim erneuten Zustellen desselben Chunks (Client-Retry) einfach identisch
// ueberschrieben — idempotent, kein Duplikat.
func (s *Staging) WriteChunk(sessionID string, offset int64, r io.Reader) (int64, error) {
if err := os.MkdirAll(s.dir, 0o755); err != nil {
return 0, fmt.Errorf("upload: staging-verzeichnis anlegen: %w", err)
}
f, err := os.OpenFile(s.path(sessionID), os.O_CREATE|os.O_WRONLY, 0o600)
if err != nil {
return 0, fmt.Errorf("upload: staging-datei oeffnen: %w", err)
}
defer func() { _ = f.Close() }()
if _, err := f.Seek(offset, io.SeekStart); err != nil {
return 0, fmt.Errorf("upload: zu offset %d springen: %w", offset, err)
}
written, err := io.Copy(f, r)
if err != nil {
return 0, fmt.Errorf("upload: chunk schreiben: %w", err)
}
return offset + written, nil
}
// OpenForRead oeffnet die vollstaendige Staging-Datei zum Lesen (nach
// Abschluss des Uploads, fuer Hash-Berechnung + Verschluesselung).
func (s *Staging) OpenForRead(sessionID string) (*os.File, error) {
f, err := os.Open(s.path(sessionID))
if err != nil {
return nil, fmt.Errorf("upload: staging-datei lesen: %w", err)
}
return f, nil
}
// Remove entfernt die Staging-Datei nach erfolgreichem Abschluss oder Abbruch.
func (s *Staging) Remove(sessionID string) error {
if err := os.Remove(s.path(sessionID)); err != nil && !os.IsNotExist(err) {
return fmt.Errorf("upload: staging-datei entfernen: %w", err)
}
return nil
}
-59
View File
@@ -1,59 +0,0 @@
// Package upload implementiert DOC-01: fortsetzbare Server-seitige
// Dateiaufnahme mit Größen-/MIME-Validierung und transaktionaler
// Dokument+Revision-Anlage. Nutzt FDN-03 (Storage) und FDN-09
// (Verschlüsselung) als bereits fertige Bausteine, dupliziert sie nicht.
package upload
import (
"errors"
"fmt"
)
// ErrFileTooLarge wird geliefert, wenn die angekündigte oder tatsächliche
// Größe das konfigurierte Limit überschreitet (Schutz vor
// Speicherbomben — siehe "Bekannte Fehler vermeiden" im Ticket).
var ErrFileTooLarge = errors.New("upload: datei ueberschreitet das erlaubte groessenlimit")
// ErrDisallowedMimeType wird geliefert, wenn der MIME-Typ nicht auf der
// Positivliste steht (Akzeptanzkriterium 2).
var ErrDisallowedMimeType = errors.New("upload: dateityp ist nicht erlaubt")
// Validator prüft Größe und MIME-Typ VOR jedem Byte, das tatsächlich
// entgegengenommen wird — die Prüfung selbst braucht keinen Datei-Inhalt,
// nur die vom Client angekündigten Metadaten.
type Validator struct {
maxSizeBytes int64
allowedMimeTypes map[string]bool
}
// NewValidator erzeugt einen Validator. maxSizeBytes<=0 bedeutet kein
// Limit (bewusst explizit statt eines "magischen" Default — siehe Ticket-
// Vorgabe "kein Start ohne Konfiguration" NUR für die Voreinstellung
// selbst, nicht für sicherheitsrelevante Limits).
func NewValidator(maxSizeBytes int64, allowedMimeTypes []string) *Validator {
allowed := make(map[string]bool, len(allowedMimeTypes))
for _, m := range allowedMimeTypes {
allowed[m] = true
}
return &Validator{maxSizeBytes: maxSizeBytes, allowedMimeTypes: allowed}
}
func (v *Validator) ValidateMimeType(mimeType string) error {
if len(v.allowedMimeTypes) == 0 {
return nil
}
if !v.allowedMimeTypes[mimeType] {
return fmt.Errorf("%w: %q", ErrDisallowedMimeType, mimeType)
}
return nil
}
func (v *Validator) ValidateSize(sizeBytes int64) error {
if v.maxSizeBytes <= 0 {
return nil
}
if sizeBytes > v.maxSizeBytes {
return fmt.Errorf("%w: %d bytes > limit %d bytes", ErrFileTooLarge, sizeBytes, v.maxSizeBytes)
}
return nil
}
@@ -1,8 +0,0 @@
DROP TABLE IF EXISTS document_metadata_values;
DROP TABLE IF EXISTS metadata_fields;
DROP TABLE IF EXISTS document_tags;
DROP TABLE IF EXISTS tags;
ALTER TABLE documents DROP CONSTRAINT IF EXISTS fk_documents_current_revision;
DROP TABLE IF EXISTS file_revisions;
DROP TABLE IF EXISTS documents;
DROP TABLE IF EXISTS folders;
@@ -1,88 +0,0 @@
-- Kern-Entitaeten des DMS (FDN-02): Dokument, Datei-Revision, Ordner, Tag,
-- Metadatenfeld. Laeuft in der DB EINES Mandanten (Modell C, siehe Core
-- TEN-01) — keine tenant_id-Spalte, die Tenant-Zugehoerigkeit ist implizit
-- durch die Datenbankverbindung gegeben. FK auf users(id) spiegelt das
-- Benutzer-Datenmodell aus Core IAM-01 (migrations/tenant/0001_users.up.sql
-- im NEXARCH-Core-Modul) — Auth/Benutzerverwaltung liegt vollstaendig in
-- Core (siehe "Nicht Bestandteil" in FDN-02), diese Migration dupliziert sie
-- NICHT, sondern setzt sie als bereits vorhanden voraus (users-Tabelle wird
-- durch Cores eigene Migration in derselben physischen Tenant-Datenbank
-- angelegt, bevor DMS-Migrationen laufen).
CREATE EXTENSION IF NOT EXISTS pgcrypto;
CREATE TABLE folders (
id UUID PRIMARY KEY DEFAULT gen_random_uuid(),
parent_folder_id UUID REFERENCES folders(id) ON DELETE CASCADE,
name TEXT NOT NULL,
created_by UUID NOT NULL REFERENCES users(id),
created_at TIMESTAMPTZ NOT NULL DEFAULT now(),
updated_at TIMESTAMPTZ NOT NULL DEFAULT now()
);
CREATE INDEX idx_folders_parent_folder_id ON folders(parent_folder_id);
CREATE TABLE documents (
id UUID PRIMARY KEY DEFAULT gen_random_uuid(),
folder_id UUID REFERENCES folders(id) ON DELETE SET NULL,
title TEXT NOT NULL,
-- current_revision_id verweist erst NACH der Anlage von file_revisions
-- auf eine Zeile (siehe ALTER TABLE unten) — beim INSERT eines Dokuments
-- existiert noch keine Revision, daher NULLable und zirkulaer per
-- nachtraeglichem FOREIGN KEY statt Inline-Referenz geloest.
current_revision_id UUID,
created_by UUID NOT NULL REFERENCES users(id),
created_at TIMESTAMPTZ NOT NULL DEFAULT now(),
updated_at TIMESTAMPTZ NOT NULL DEFAULT now(),
deleted_at TIMESTAMPTZ
);
CREATE INDEX idx_documents_folder_id ON documents(folder_id);
CREATE INDEX idx_documents_created_by ON documents(created_by);
CREATE TABLE file_revisions (
id UUID PRIMARY KEY DEFAULT gen_random_uuid(),
document_id UUID NOT NULL REFERENCES documents(id) ON DELETE CASCADE,
revision_number INT NOT NULL,
-- storage_key ist ein Platzhalter fuer die Objekt-Storage-Abstraktion
-- (FDN-03, "Nicht Bestandteil" dieser Kachel) — hier nur die Spalte, die
-- spaetere Kachel legt fest, was tatsaechlich dahinter liegt.
storage_key TEXT NOT NULL,
checksum_sha256 TEXT NOT NULL,
size_bytes BIGINT NOT NULL,
mime_type TEXT NOT NULL,
created_by UUID NOT NULL REFERENCES users(id),
created_at TIMESTAMPTZ NOT NULL DEFAULT now(),
UNIQUE (document_id, revision_number)
);
CREATE INDEX idx_file_revisions_document_id ON file_revisions(document_id);
ALTER TABLE documents
ADD CONSTRAINT fk_documents_current_revision
FOREIGN KEY (current_revision_id) REFERENCES file_revisions(id) ON DELETE SET NULL;
CREATE TABLE tags (
id UUID PRIMARY KEY DEFAULT gen_random_uuid(),
name TEXT NOT NULL UNIQUE,
created_at TIMESTAMPTZ NOT NULL DEFAULT now()
);
CREATE TABLE document_tags (
document_id UUID NOT NULL REFERENCES documents(id) ON DELETE CASCADE,
tag_id UUID NOT NULL REFERENCES tags(id) ON DELETE CASCADE,
PRIMARY KEY (document_id, tag_id)
);
CREATE INDEX idx_document_tags_tag_id ON document_tags(tag_id);
CREATE TABLE metadata_fields (
id UUID PRIMARY KEY DEFAULT gen_random_uuid(),
field_key TEXT NOT NULL UNIQUE,
label TEXT NOT NULL,
field_type TEXT NOT NULL CHECK (field_type IN ('text', 'number', 'date', 'bool', 'select')),
created_at TIMESTAMPTZ NOT NULL DEFAULT now()
);
CREATE TABLE document_metadata_values (
document_id UUID NOT NULL REFERENCES documents(id) ON DELETE CASCADE,
field_id UUID NOT NULL REFERENCES metadata_fields(id) ON DELETE CASCADE,
value TEXT NOT NULL,
PRIMARY KEY (document_id, field_id)
);
CREATE INDEX idx_document_metadata_values_field_id ON document_metadata_values(field_id);
@@ -1 +0,0 @@
DROP TABLE IF EXISTS processing_jobs;
@@ -1,30 +0,0 @@
-- FDN-04: Postgres-Jobqueue fuer asynchrone Verarbeitung (OCR, Konvertierung,
-- Indexierung, Exporte) — kein Redis/AMQP, dasselbe Muster wie das
-- projektweite Postgres-Jobqueue-Konzept (siehe SKALIERUNGSKONZEPT.md).
-- Laeuft in der DB EINES Mandanten (Modell C, siehe Core TEN-01).
CREATE TABLE processing_jobs (
id UUID PRIMARY KEY DEFAULT gen_random_uuid(),
job_type TEXT NOT NULL,
payload JSONB NOT NULL DEFAULT '{}'::jsonb,
-- idempotency_key verhindert doppelte Einreihung DERSELBEN logischen
-- Aufgabe (z.B. "ocr:<revision_id>") — NULL erlaubt mehrere Zeilen ohne
-- Dedup-Anspruch (Standard-Postgres-Verhalten: NULL ist nie gleich NULL
-- im UNIQUE-Index).
idempotency_key TEXT UNIQUE,
status TEXT NOT NULL DEFAULT 'pending'
CHECK (status IN ('pending', 'processing', 'succeeded', 'failed', 'dead_letter')),
attempts INT NOT NULL DEFAULT 0,
max_attempts INT NOT NULL DEFAULT 5,
available_at TIMESTAMPTZ NOT NULL DEFAULT now(),
locked_at TIMESTAMPTZ,
locked_by TEXT,
last_error TEXT,
created_at TIMESTAMPTZ NOT NULL DEFAULT now(),
updated_at TIMESTAMPTZ NOT NULL DEFAULT now()
);
-- Deckt genau die Zugriffsmuster von Dequeue (status+available_at) und der
-- Stale-Lock-Wiedervorlage (status+locked_at) ab.
CREATE INDEX idx_processing_jobs_pending ON processing_jobs (available_at) WHERE status = 'pending';
CREATE INDEX idx_processing_jobs_processing ON processing_jobs (locked_at) WHERE status = 'processing';
CREATE INDEX idx_processing_jobs_job_type ON processing_jobs (job_type);
@@ -1 +0,0 @@
DROP TABLE IF EXISTS upload_sessions;
@@ -1,18 +0,0 @@
-- DOC-01: Sitzungszustand fuer fortsetzbaren Upload (Akzeptanzkriterium 1).
-- Ein Sitzungseintrag je laufendem Upload, bytes_received wird bei jedem
-- angenommenen Chunk aktualisiert — nach einem Verbindungsabbruch kann der
-- Client anhand von bytes_received genau dort fortsetzen, wo er stand.
CREATE TABLE upload_sessions (
id UUID PRIMARY KEY DEFAULT gen_random_uuid(),
filename TEXT NOT NULL,
mime_type TEXT NOT NULL,
total_size BIGINT NOT NULL CHECK (total_size > 0),
bytes_received BIGINT NOT NULL DEFAULT 0 CHECK (bytes_received >= 0),
folder_id UUID REFERENCES folders(id) ON DELETE SET NULL,
created_by UUID NOT NULL REFERENCES users(id),
status TEXT NOT NULL DEFAULT 'uploading'
CHECK (status IN ('uploading', 'completed', 'aborted')),
document_id UUID REFERENCES documents(id),
created_at TIMESTAMPTZ NOT NULL DEFAULT now(),
updated_at TIMESTAMPTZ NOT NULL DEFAULT now()
);
@@ -1 +0,0 @@
ALTER TABLE file_revisions DROP COLUMN IF EXISTS wrapped_dek;
@@ -1,6 +0,0 @@
-- DOC-01: file_revisions bekommt den verpackten Datenverschluesselungs-
-- schluessel (DEK) aus FDN-09s Envelope-Encryption. Nullable, weil ein
-- Datensatz theoretisch auch unverschluesselt vorliegen kann (z.B. der
-- FDN-02-Entwicklungs-Seed) — DOC-01s regulaerer Upload-Pfad setzt ihn
-- jedoch immer.
ALTER TABLE file_revisions ADD COLUMN wrapped_dek BYTEA;
-43
View File
@@ -1,43 +0,0 @@
-- Entwicklungs-Seed (Akzeptanzkriterium 3): legt einen Beispielordner, ein
-- Beispieldokument mit einer Revision, ein Tag und ein Metadatenfeld an.
-- Setzt voraus, dass mindestens ein Benutzer existiert (Core IAM-01 legt
-- users an, DMS tut das nicht selbst - siehe "Nicht Bestandteil" in
-- FDN-02) - schlaegt sonst absichtlich mit einer sprechenden Fehlermeldung
-- fehl statt einen Platzhalter-Benutzer anzulegen, den DMS gar nicht
-- verwalten darf.
DO $$
DECLARE
seed_user_id UUID;
seed_folder_id UUID;
seed_document_id UUID;
BEGIN
SELECT id INTO seed_user_id FROM users ORDER BY created_at LIMIT 1;
IF seed_user_id IS NULL THEN
RAISE EXCEPTION 'dev_seed.sql: keine Zeile in users gefunden - zuerst Core-Seed (IAM-01) ausfuehren';
END IF;
INSERT INTO folders (name, created_by) VALUES ('Beispielordner', seed_user_id)
RETURNING id INTO seed_folder_id;
INSERT INTO documents (folder_id, title, created_by) VALUES (seed_folder_id, 'Beispieldokument', seed_user_id)
RETURNING id INTO seed_document_id;
INSERT INTO file_revisions (document_id, revision_number, storage_key, checksum_sha256, size_bytes, mime_type, created_by)
VALUES (seed_document_id, 1, 'dev-seed/beispiel.pdf', repeat('0', 64), 12345, 'application/pdf', seed_user_id);
UPDATE documents SET current_revision_id = (
SELECT id FROM file_revisions WHERE document_id = seed_document_id AND revision_number = 1
) WHERE id = seed_document_id;
INSERT INTO tags (name) VALUES ('Beispiel-Tag')
ON CONFLICT (name) DO NOTHING;
INSERT INTO document_tags (document_id, tag_id)
SELECT seed_document_id, id FROM tags WHERE name = 'Beispiel-Tag';
INSERT INTO metadata_fields (field_key, label, field_type) VALUES ('rechnungsnummer', 'Rechnungsnummer', 'text')
ON CONFLICT (field_key) DO NOTHING;
INSERT INTO document_metadata_values (document_id, field_id, value)
SELECT seed_document_id, id, 'RE-2026-0001' FROM metadata_fields WHERE field_key = 'rechnungsnummer';
END $$;
-22
View File
@@ -1,22 +0,0 @@
#!/usr/bin/env bash
# Setzt die DMS-Testumgebung zurueck: droppt die Tenant-Test-Datenbank und
# legt sie leer neu an. Noetig, weil mehrere Testpakete (internal/jobqueue,
# internal/migrate, ...) dieselbe physische Test-Datenbank ueber
# Sitzungsgrenzen hinweg teilen — ohne Reset sammelt sich Zustand
# (z.B. schema_migrations-Eintraege) an, der Migrations-/Rollback-Tests
# verfaelscht (dieselbe Fehlerklasse wie in NEXARCH Core, siehe
# [[project-nexarch-test-infra]]).
#
# Aufruf: NEXARCH_DMS_TEST_DB_PASSWORD=... TEST_TENANT_DB=dms_tenant_test ./scripts/reset-test-env.sh
set -euo pipefail
PASS="${NEXARCH_DMS_TEST_DB_PASSWORD:?Setze NEXARCH_DMS_TEST_DB_PASSWORD vor dem Aufruf}"
ROLE="${TEST_TENANT_ROLE:-nexarch_dms_test}"
DB="${TEST_TENANT_DB:-dms_tenant_test}"
export PGPASSWORD="$PASS"
psql -h localhost -U "$ROLE" -d postgres -v ON_ERROR_STOP=1 -c "DROP DATABASE IF EXISTS ${DB};"
psql -h localhost -U "$ROLE" -d postgres -v ON_ERROR_STOP=1 -c "CREATE DATABASE ${DB};"
echo "Testumgebung zurueckgesetzt: ${DB} leer neu angelegt."
+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
)
+28
View File
@@ -0,0 +1,28 @@
github.com/davecgh/go-spew v1.1.0/go.mod h1:J7Y8YcW2NihsgmVo/mv3lAwl/skON4iLHjSsI+c5H38=
github.com/davecgh/go-spew v1.1.1 h1:vj9j/u1bqnvCEfJOwUhtlOARqs3+rkHYY13jYWTU97c=
github.com/davecgh/go-spew v1.1.1/go.mod h1:J7Y8YcW2NihsgmVo/mv3lAwl/skON4iLHjSsI+c5H38=
github.com/jackc/pgpassfile v1.0.0 h1:/6Hmqy13Ss2zCq62VdNG8tM1wchn8zjSGOBJ6icpsIM=
github.com/jackc/pgpassfile v1.0.0/go.mod h1:CEx0iS5ambNFdcRtxPj5JhEz+xB6uRky5eyVu/W2HEg=
github.com/jackc/pgservicefile v0.0.0-20221227161230-091c0ba34f0a h1:bbPeKD0xmW/Y25WS6cokEszi5g+S0QxI/d45PkRi7Nk=
github.com/jackc/pgservicefile v0.0.0-20221227161230-091c0ba34f0a/go.mod h1:5TJZWKEWniPve33vlWYSoGYefn3gLQRzjfDlhSJ9ZKM=
github.com/jackc/pgx/v5 v5.6.0 h1:SWJzexBzPL5jb0GEsrPMLIsi/3jOo7RHlzTjcAeDrPY=
github.com/jackc/pgx/v5 v5.6.0/go.mod h1:DNZ/vlrUnhWCoFGxHAG8U2ljioxukquj7utPDgtQdTw=
github.com/jackc/puddle/v2 v2.2.1 h1:RhxXJtFG022u4ibrCSMSiu5aOq1i77R3OHKNJj77OAk=
github.com/jackc/puddle/v2 v2.2.1/go.mod h1:vriiEXHvEE654aYKXXjOvZM39qJ0q+azkZFrfEOc3H4=
github.com/pmezard/go-difflib v1.0.0 h1:4DBwDE0NGyQoBHbLQYPwSUPoCMWR5BEzIk/f1lZbAQM=
github.com/pmezard/go-difflib v1.0.0/go.mod h1:iKH77koFhYxTK1pcRnkKkqfTogsbg7gZNVY4sRDYZ/4=
github.com/stretchr/objx v0.1.0/go.mod h1:HFkY916IF+rwdDfMAkV7OtwuqBVzrE8GR6GFx+wExME=
github.com/stretchr/testify v1.3.0/go.mod h1:M5WIy9Dh21IEIfnGCwXGc5bZfKNJtfHm1UVUgZn+9EI=
github.com/stretchr/testify v1.7.0/go.mod h1:6Fq8oRcR53rry900zMqJjRRixrwX3KX962/h/Wwjteg=
github.com/stretchr/testify v1.8.1 h1:w7B6lhMri9wdJUVmEZPGGhZzrYTPvgJArz7wNPgYKsk=
github.com/stretchr/testify v1.8.1/go.mod h1:w2LPCIKwWwSfY2zedu0+kehJoqGctiVI29o6fzry7u4=
golang.org/x/crypto v0.17.0 h1:r8bRNjWL3GshPW3gkd+RpvzWrZAwPS49OmTGZ/uhM4k=
golang.org/x/crypto v0.17.0/go.mod h1:gCAAfMLgwOJRpTjQ2zCCt2OcSfYMTeZVSRtQlPC7Nq4=
golang.org/x/sync v0.1.0 h1:wsuoTGHzEhffawBOhz5CYhcrV4IdKZbEyZjBMuTp12o=
golang.org/x/sync v0.1.0/go.mod h1:RxMgew5VJxzue5/jJTE5uejpjVlOe/izrB70Jof72aM=
golang.org/x/text v0.14.0 h1:ScX5w1eTa3QqT8oi6+ziP7dTV1S2+ALU0bI+0zXKWiQ=
golang.org/x/text v0.14.0/go.mod h1:18ZOQIKpY8NJVqYksKHtTdi31H5itFRjB5/qKTNYzSU=
gopkg.in/check.v1 v0.0.0-20161208181325-20d25e280405/go.mod h1:Co6ibVJAznAaIkqp8huTwlJQCZ016jof/cbN4VW5Yz0=
gopkg.in/yaml.v3 v3.0.0-20200313102051-9f266ea9e77c/go.mod h1:K4uyk7z7BCEPqu6E+C64Yfv1cQ7kz7rIZviUmN+EgEM=
gopkg.in/yaml.v3 v3.0.1 h1:fxVm/GzAzEWqLHuvctI91KS9hhNmmWOoWu0XTYJS7CA=
gopkg.in/yaml.v3 v3.0.1/go.mod h1:K4uyk7z7BCEPqu6E+C64Yfv1cQ7kz7rIZviUmN+EgEM=
+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")
}
}()
}
+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()
);
+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