Compare commits

..
Author SHA1 Message Date
sysopsandClaude Sonnet 5 6b5cefc20f ARC-08: verschluesselungsschluessel-rotation
Tenant-KEK-Rotation ohne Neuverschlüsselung des Archivbestands
(Envelope-Encryption bleibt aus ARC-02 unverändert, Objekt-DEKs werden
nicht angefasst).

Core (API-10, RotateTenantKEK) ersetzt den Tenant-KEK durch einen neuen
Wert und hält keine Historie vor — TenantKEKHandler liefert immer nur den
aktuellen Schlüssel. Damit Mail Altbestand nach einer Rotation weiterhin
lesen kann, versioniert Mail selbst jeden bezogenen Tenant-KEK:

- crypto/kekversions.go: KEKVersionStore, lokal verschlüsselt mit
  eigenem Wrap-Schlüssel (nur über Umgebungsvariable), erkennt Rotation
  automatisch (RecordIfNew), erlaubt gezieltes Sperren einer Version
  (Revoke).
- crypto/service.go: Service.WithVersionStore (optional, Open bleibt für
  Rückwärtskompatibilität unverändert), Seal zeichnet die verwendete
  KEK-Version auf, neue Methode OpenAtVersion liest mit historischer
  statt aktueller Version.
- encstorage.go: neuer .dek.version-Sidecar (gleiches Muster wie der
  bestehende .dek-Sidecar), GetDecrypted nutzt OpenAtVersion; fehlender
  Sidecar (Altobjekte vor ARC-08) fällt auf Version 0 zurück, identisches
  Verhalten wie vorher.

Prüfungen (alle real durchgeführt, siehe mail/docs/ARC-08-PRUEFPROTOKOLL.md):
1. TestRotation_OldArchiveStaysReadableAfterMasterKeyRotation: Altbestand
   nach realer Rotation weiterhin lesbar über OpenAtVersion, naives Open
   mit dem neuen Schlüssel schlägt für das alte Objekt real fehl.
2. TestRotation_CompromisedOldKeyCanBeRevoked: gesperrte Version blockiert
   Lesezugriff real, andere Versionen bleiben unberührt.
3. Rotationsvorgang vollständig durchgespielt (siehe Prüfprotokoll).

Kein Umbau: storage/dedup/indexworker/search unverändert, bestehende
ARC-02-Tests (encstorage_test.go) unverändert weiterhin grün.

Co-Authored-By: Claude Sonnet 5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01HhgFcLS8tYMhDJpP74C6AQ
2026-08-31 11:17:23 +02:00
sysopsandClaude Sonnet 5 86c4223855 SRC-09: suchindex-neuaufbau-reindexierung
Werkzeug für vollständigen Suchindex-Neuaufbau: neue Tabelle anlegen,
Dokumente aus der lebenden Tabelle kopieren, Trefferzahlen verifizieren,
erst dann per Manticore RENAME atomar umschalten.

- reindex.go: Reindexer.Rebuild mit Fortschritts-Callback, Cursor-
  Paginierung über id, strukturierte JSON-API (kein dynamischer
  SQL-Klauselbau). Bei Fehler vor dem Umschalten bleibt die lebende
  Tabelle unverändert, Zwischentabelle wird entfernt.
- Manticore-Verhalten entdeckt: frisch eingefügte Dokumente einer neuen
  RT-Tabelle sind für match_all-Zählungen erst nach FLUSH RAMCHUNK
  zuverlässig sichtbar — vor der Konsistenzprüfung eingebaut.
- Plattformgrenze entdeckt: kein atomares Mehrfach-RENAME in Manticore,
  Sub-Millisekunden-Fenster zwischen den zwei nötigen Einzel-RENAMEs.
  Client.Search bekam einen begrenzten Retry auf "unknown local table".
- Nebenbei echten latenten Bug in Search behoben: ohne explizites limit
  begrenzte Manticore Ergebnisse standardmäßig auf 20 Treffer, unbemerkt
  seit SRC-01 (bisherige Tests prüften nur Vorhandensein, nie Gesamtzahl).

Prüfungen (alle real durchgeführt, siehe mail/docs/SRC-09-PRUEFPROTOKOLL.md):
1. TestRebuild_SearchKeepsWorkingDuringReindex: 0 fehlgeschlagene Suchen
   während parallelem Reindex.
2. TestRebuild_AbortedReindexLeavesNoInconsistentState: abgebrochener
   Kontext hinterlässt real weder Datenverlust noch verwaiste Tabellen.
3. TestRebuild_SampleComparisonMatchesOldAndNewIndex: Stichproben vor/
   nach Reindex real identisch.

Kein Umbau: Index/Delete/Facets-Verhalten sonst unverändert,
dedup/indexworker/storage/crypto/encstorage unverändert.

Co-Authored-By: Claude Sonnet 5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01HhgFcLS8tYMhDJpP74C6AQ
2026-08-31 11:01:45 +02:00
10 changed files with 1202 additions and 22 deletions
+87
View File
@@ -0,0 +1,87 @@
# ARC-08 Prüfprotokoll: Verschlüsselungsschlüssel-Rotation
Voraussetzung ARC-02 (Fertig).
## Architektur-Ausgangslage (real geprüft)
Core (API-10, `internal/kek.Store.RotateTenantKEK`, bereits Fertig)
ersetzt den Tenant-KEK bei Rotation durch einen komplett NEUEN Wert und
hält KEINE Historie vor — der laufende `kek-api`-Dienst (192.168.1.131,
Port 8102) exponiert ausschließlich `TenantKEKHandler`, der immer nur den
AKTUELLEN KEK liefert (real im Quelltext von
`/root/nexarch-code/internal/kek/handler.go` und `cmd/kek-api/main.go`
auf 131 verifiziert). Damit Mail nach einer Core-seitigen Rotation
Altbestand weiterhin lesen kann, MUSS Mail selbst jeden bezogenen
Tenant-KEK versioniert zwischenspeichern — das ist der Kern dieser
Kachel.
## Umsetzung
- `mail/internal/crypto/kekversions.go``KEKVersionStore`: persistiert
jede vom Core bezogene Tenant-KEK-Version lokal, verschlüsselt mit
einem eigenen, ausschließlich über Umgebungsvariable bezogenen
Wrap-Schlüssel (kein Klartext-KEK in der Datenbank). `RecordIfNew`
erkennt Rotation (neuer KEK-Wert ≠ letzter bekannter) und legt nur dann
eine neue Version an (Akzeptanzkriterium 1). `Revoke` sperrt gezielt
eine einzelne Version (Pflichtprüfung 2).
- `mail/internal/crypto/service.go``Service.WithVersionStore`
(optional, Rückwärtskompatibilität: ohne Aufruf verhält sich `Service`
exakt wie vor ARC-08). `Seal` zeichnet bei aktivierter Versionierung
die verwendete KEK-Version im `Envelope` auf. Neue Methode
`OpenAtVersion` entpackt mit der historischen statt der aktuellen
Tenant-KEK-Version (Akzeptanzkriterium 3) — `Open` bleibt unverändert
für Rückwärtskompatibilität.
- `mail/internal/encstorage/encstorage.go` — neuer Sidecar
`<key>.dek.version` (gleiches Muster wie der bestehende `.dek`-Sidecar
aus ARC-02) speichert die KEK-Version je Objekt. `GetDecrypted` nutzt
jetzt `OpenAtVersion` statt `Open`; fehlt der Sidecar (vor ARC-08
geschriebene Objekte), wird Version 0 angenommen (identisches
Verhalten wie vorher).
- Kein Umbau: `mail/internal/storage`/`mail/internal/dedup`/
`mail/internal/indexworker`/`mail/internal/search` unverändert;
bestehende ARC-02-Tests (`encstorage_test.go`) unverändert lauffähig
ohne Codeänderung an ihnen.
## Prüfungen
| # | Prüfung | Ergebnis |
|---|---|---|
| 1 | Test: Rotation des Hauptschlüssels lässt Altbestand weiterhin lesbar | **bestanden** `TestRotation_OldArchiveStaysReadableAfterMasterKeyRotation`: Objekt vor Rotation versiegelt (Version 1), Tenant-Hauptschlüssel real rotiert (Provider liefert ab dann einen anderen Wert, exakt wie `RotateTenantKEK` es bei Core bewirkt), neues Objekt nach Rotation versiegelt (Version 2), Altbestand über `OpenAtVersion` real weiterhin korrekt entschlüsselt — zusätzlich real bestätigt, dass der naive `Open()` mit dem neuen aktuellen KEK für das alte Objekt fehlschlägt (beweist, dass `OpenAtVersion` tatsächlich etwas leistet) |
| 2 | Test: kompromittierter alter Schlüssel kann gezielt gesperrt werden | **bestanden** `TestRotation_CompromisedOldKeyCanBeRevoked`: Version gesperrt, `OpenAtVersion` liefert danach real `ErrKEKVersionRevoked`; `TestKEKVersionStore_RevokeBlocksOnlyThatVersion` bestätigt zusätzlich, dass eine ANDERE Version davon unberührt bleibt |
| 3 | Dokumentierter Rotationsvorgang wurde einmal vollständig durchgespielt | **bestanden** siehe Abschnitt "Rotationsvorgang" unten, real durchlaufen als `TestRotation_OldArchiveStaysReadableAfterMasterKeyRotation` |
### Rotationsvorgang (Pflichtprüfung 3, vollständig durchgespielt)
1. Objekt A wird mit Tenant-KEK-Version 1 versiegelt (`Seal`, Envelope
trägt `KEKVersion=1`, `KEKVersionStore` legt Version 1 real an).
2. Core rotiert den Tenant-Hauptschlüssel (in diesem Test durch den
`KEKProvider` simuliert, exakt am selben Punkt, an dem `Service` mit
dem echten `HTTPKEKProvider`/Core API-12 interagieren würde).
3. Objekt B wird versiegelt — automatisch mit der NEUEN Version 2, ohne
dass Objekt A angefasst wird (Akzeptanzkriterium 2: kein
Neuverschlüsseln des Bestands).
4. Objekt A wird über `OpenAtVersion(..., kekVersion=1, ...)` gelesen —
real erfolgreich, Klartext identisch zum Original.
5. Ein naiver Lesezugriff über `Open()` (aktueller KEK) auf Objekt A
schlägt real fehl — zeigt, dass ohne Versionsverfolgung der
Altbestand nach Rotation unlesbar geworden wäre.
## Build/Test-Ergebnis (192.168.1.131)
```
go build ./... -> clean
go vet ./... -> clean
golangci-lint run ./... -> 0 issues
TEST_TENANT_DSN=postgresql://nexarch_test:***@localhost:5432/tenant_acme?sslmode=disable \
TEST_MANTICORE_URL=http://127.0.0.1:9308 \
go test ./... -v -p 1 -> alle Pakete bestanden, inkl. internal/crypto (4 Tests, neu)
und internal/encstorage (4 Tests, unverändert weiterhin grün — Rückwärtskompatibilität
real bestätigt)
```
## Gesamtergebnis
**Bestanden.** Alle drei Akzeptanzkriterien und alle drei Pflichtprüfungen
real erfüllt. Trägt (gemeinsam mit SRC-02, SRC-04, SRC-05, SRC-09) zu
QA-03 bei — QA-03 bleibt weiterhin blockiert, bis auch SRC-08 und SRC-10
fertig sind.
+72
View File
@@ -0,0 +1,72 @@
# SRC-09 Prüfprotokoll: Suchindex-Neuaufbau/Reindexierung
Voraussetzung SRC-01 (Fertig).
## Umsetzung
- `mail/internal/search/reindex.go``Reindexer.Rebuild(ctx, onProgress)`:
1. legt eine neue physische Manticore-Tabelle an (Name aus striktem
Muster `mail_documents_reindex_<Ziffern>`, per Regex validiert —
Verteidigung in der Tiefe, obwohl der Wert ausschließlich
paketintern erzeugt wird),
2. kopiert alle Dokumente aus der lebenden Tabelle seitenweise
(Cursor-Paginierung über `id`, strukturierte JSON-API, kein
dynamischer SQL-Klauselbau) — die lebende Tabelle wird dabei nur
gelesen, nie verändert (Akzeptanzkriterium 1),
3. meldet Fortschritt über einen `onProgress`-Callback
(Akzeptanzkriterium 2),
4. vergleicht Trefferzahlen alt/neu — bei Abweichung kein Umschalten,
5. schaltet erst danach per Manticore `ALTER TABLE ... RENAME`
(reine Metadaten-Operation) atomar um. Schlägt ein Schritt vor dem
Umschalten fehl, wird die Zwischentabelle entfernt, die lebende
Tabelle bleibt unverändert (Akzeptanzkriterium 3 / Pflichtprüfung 2).
- Echtes Manticore-Verhalten entdeckt und behandelt: frisch eingefügte
Dokumente einer neu angelegten RT-Tabelle sind für `match_all`-Zählungen
erst nach explizitem `FLUSH RAMCHUNK` zuverlässig sichtbar (SQL-`SELECT`
sah sie sofort, `/search`-Zählung zeigte 0) — vor der
Konsistenzprüfung eingebaut.
- Echte Plattformgrenze gefunden und abgefangen: Manticore unterstützt kein
atomares Mehrfach-`RENAME` in einer Anweisung — zwischen den zwei
nötigen Einzel-`RENAME`s existiert ein Sub-Millisekunden-Fenster ohne
`mail_documents`-Tabelle. `Client.Search` bekam dafür einen begrenzten
Retry (bis zu 2 Wiederholungen, 20ms Pause) speziell auf den
Manticore-Fehler `"unknown local table"` — real durch eine parallele
Suchlast während des Umschaltens nachgewiesen (Pflichtprüfung 1).
- Nebenbei einen echten, latenten Fehler in `Search` gefunden und behoben:
ohne explizites `limit` begrenzte Manticore Ergebnisse standardmäßig auf
20 Treffer — unbemerkt, weil bisherige Tests (SRC-01/03/05) nur auf das
Vorhandensein einzelner Treffer prüften, nie auf die Gesamtzahl. Jetzt
`searchResultLimit = 1000`.
- Kein Umbau: `Index`/`Delete`/`Facets`-Verhalten sonst unverändert,
`mail/internal/dedup`/`indexworker`/`storage`/`crypto`/`encstorage`
unverändert.
## Prüfungen
| # | Prüfung | Ergebnis |
|---|---|---|
| 1 | Test: Reindex während laufender Suchanfragen unterbricht die Suche nicht | **bestanden** `TestRebuild_SearchKeepsWorkingDuringReindex`: 30 reale Dokumente indexiert, parallele Sucher-Goroutine (alle 2ms) läuft während `Rebuild` mit — 0 fehlgeschlagene Suchen über den gesamten Umschaltvorgang, danach weiterhin real alle 30 Treffer auffindbar |
| 2 | Test: abgebrochener Reindex hinterlässt keinen inkonsistenten Zustand | **bestanden** `TestRebuild_AbortedReindexLeavesNoInconsistentState`: Kontext vor `Rebuild` abgebrochen, Fehler kommt real zurück, lebende Tabelle bleibt danach unverändert (weiterhin 1 Treffer real auffindbar), keine verwaisten Zwischentabellen über `SHOW TABLES` real bestätigt |
| 3 | Stichprobenvergleich Alt-/Neuindex bestätigt gleiche Trefferzahlen | **bestanden** `TestRebuild_SampleComparisonMatchesOldAndNewIndex`: 3 unterschiedliche Suchbegriffe vor und nach Reindex real verglichen, identische Trefferzahlen je Stichprobe |
Zusätzlich (Akzeptanzkriterium 2): `TestRebuild_ReportsProgress` bestätigt
reale Fortschrittsmeldungen bis zum vollständigen Abschluss.
## Build/Test-Ergebnis (192.168.1.131)
```
go build ./... -> clean
go vet ./... -> clean
golangci-lint run ./... -> 0 issues
TEST_TENANT_DSN=postgresql://nexarch_test:***@localhost:5432/tenant_acme?sslmode=disable \
TEST_MANTICORE_URL=http://127.0.0.1:9308 \
go test ./... -v -p 1 -> alle Pakete bestanden, inkl. internal/search (14 Tests,
keine Regression in dedup/indexworker/storage/encstorage/example/mimeparse/pflichttestgate)
```
## Gesamtergebnis
**Bestanden.** Alle drei Akzeptanzkriterien und alle drei Pflichtprüfungen
real erfüllt. Trägt (gemeinsam mit ARC-08, SRC-02, SRC-04, SRC-05,
SRC-08, SRC-10) zu QA-03 bei — QA-03 bleibt weiterhin blockiert, bis auch
ARC-08, SRC-08 und SRC-10 fertig sind.
+144
View File
@@ -0,0 +1,144 @@
// ARC-08: Tenant-KEK-Rotation. Core (API-10, `internal/kek.Store.
// RotateTenantKEK`) ersetzt den Tenant-KEK durch einen komplett neuen
// Wert — Core selbst hält KEINE Historie vor, `TenantKEKHandler` liefert
// immer nur den AKTUELLEN Schlüssel (siehe kekprovider.go). Damit ARC-08s
// Akzeptanzkriterium 3 ("alte Schlüsselversionen bleiben für
// Lesezugriff kontrolliert verfügbar") erfüllbar ist, muss Mail selbst
// jeden von Core bezogenen Tenant-KEK versioniert zwischenspeichern —
// KEKVersionStore übernimmt genau das, lokal mit einem eigenen,
// ausschließlich über Umgebungsvariable bezogenen Wrap-Schlüssel
// verschlüsselt (kein Klartext-KEK in der Datenbank).
package crypto
import (
"bytes"
"context"
"errors"
"fmt"
"github.com/jackc/pgx/v5"
"github.com/jackc/pgx/v5/pgxpool"
)
// ErrKEKVersionNotFound wird geliefert, wenn die angefragte Version für
// den Mandanten nicht existiert.
var ErrKEKVersionNotFound = errors.New("crypto: kek-version nicht gefunden")
// ErrKEKVersionRevoked wird geliefert, wenn die angefragte Version gezielt
// gesperrt wurde (Pflichtprüfung 2: kompromittierter alter Schlüssel kann
// gezielt gesperrt werden) — der Lesezugriff auf mit dieser Version
// verschlüsselte Altobjekte ist dann bewusst blockiert.
var ErrKEKVersionRevoked = errors.New("crypto: kek-version wurde gesperrt")
// KEKVersionStore verwaltet die Versionshistorie der Tenant-KEKs, die
// dieses Mail-Modul im Lauf der Zeit von Core bezogen hat.
type KEKVersionStore struct {
pool *pgxpool.Pool
localWrapKey []byte
}
// NewKEKVersionStore erzeugt einen Store. localWrapKey verschlüsselt die
// zwischengespeicherten Tenant-KEKs lokal at rest (KEKSize Bytes,
// ausschließlich über Umgebungsvariable bezogen — nie im Code).
func NewKEKVersionStore(pool *pgxpool.Pool, localWrapKey []byte) *KEKVersionStore {
return &KEKVersionStore{pool: pool, localWrapKey: localWrapKey}
}
// EnsureSchema legt die Tabelle an, falls sie noch nicht existiert —
// gleiches Muster wie mail/internal/dedup/indexworker (kein zentraler
// Migrationsläufer für Mandanten-Datenbanken im Mail-Modul vorhanden).
func (s *KEKVersionStore) EnsureSchema(ctx context.Context) error {
if _, err := s.pool.Exec(ctx, `
CREATE TABLE IF NOT EXISTS mail_kek_versions (
tenant_slug TEXT NOT NULL,
version INT NOT NULL,
wrapped_kek BYTEA NOT NULL,
revoked BOOLEAN NOT NULL DEFAULT false,
created_at TIMESTAMPTZ NOT NULL DEFAULT now(),
PRIMARY KEY (tenant_slug, version)
)
`); err != nil {
return fmt.Errorf("crypto: kek-versionsschema anlegen: %w", err)
}
return nil
}
// RecordIfNew merkt sich plainKEK als neue Version für tenantSlug, FALLS
// er sich vom zuletzt gespeicherten Wert unterscheidet (Rotation
// erkannt) — bei unverändertem KEK wird keine neue Version angelegt,
// sondern die bestehende Versionsnummer zurückgegeben (Akzeptanzkriterium
// 1: rotierbar verwaltet, nicht bei jedem Aufruf eine neue Version).
func (s *KEKVersionStore) RecordIfNew(ctx context.Context, tenantSlug string, plainKEK []byte) (version int, err error) {
var latestVersion int
var latestWrapped []byte
err = s.pool.QueryRow(ctx, `
SELECT version, wrapped_kek FROM mail_kek_versions
WHERE tenant_slug = $1 ORDER BY version DESC LIMIT 1
`, tenantSlug).Scan(&latestVersion, &latestWrapped)
switch {
case errors.Is(err, pgx.ErrNoRows):
return s.insertVersion(ctx, tenantSlug, 1, plainKEK)
case err != nil:
return 0, fmt.Errorf("crypto: letzte kek-version lesen: %w", err)
}
latestPlain, err := open(s.localWrapKey, latestWrapped)
if err != nil {
return 0, fmt.Errorf("crypto: zwischengespeicherten kek entpacken: %w", err)
}
if bytes.Equal(latestPlain, plainKEK) {
return latestVersion, nil
}
return s.insertVersion(ctx, tenantSlug, latestVersion+1, plainKEK)
}
func (s *KEKVersionStore) insertVersion(ctx context.Context, tenantSlug string, version int, plainKEK []byte) (int, error) {
wrapped, err := seal(s.localWrapKey, plainKEK)
if err != nil {
return 0, fmt.Errorf("crypto: kek für zwischenspeicherung verpacken: %w", err)
}
if _, err := s.pool.Exec(ctx, `
INSERT INTO mail_kek_versions (tenant_slug, version, wrapped_kek) VALUES ($1, $2, $3)
`, tenantSlug, version, wrapped); err != nil {
return 0, fmt.Errorf("crypto: kek-version speichern: %w", err)
}
return version, nil
}
// Get liefert den entschlüsselten historischen Tenant-KEK einer
// bestimmten Version. Liefert ErrKEKVersionRevoked, wenn die Version
// gezielt gesperrt wurde (Pflichtprüfung 2).
func (s *KEKVersionStore) Get(ctx context.Context, tenantSlug string, version int) ([]byte, error) {
var wrapped []byte
var revoked bool
err := s.pool.QueryRow(ctx, `
SELECT wrapped_kek, revoked FROM mail_kek_versions
WHERE tenant_slug = $1 AND version = $2
`, tenantSlug, version).Scan(&wrapped, &revoked)
if err != nil {
if errors.Is(err, pgx.ErrNoRows) {
return nil, ErrKEKVersionNotFound
}
return nil, fmt.Errorf("crypto: kek-version lesen: %w", err)
}
if revoked {
return nil, ErrKEKVersionRevoked
}
return open(s.localWrapKey, wrapped)
}
// Revoke sperrt eine Tenant-KEK-Version gezielt (Pflichtprüfung 2):
// nachfolgende Get-Aufrufe für genau diese Version schlagen mit
// ErrKEKVersionRevoked fehl, andere Versionen bleiben unberührt.
func (s *KEKVersionStore) Revoke(ctx context.Context, tenantSlug string, version int) error {
tag, err := s.pool.Exec(ctx, `
UPDATE mail_kek_versions SET revoked = true WHERE tenant_slug = $1 AND version = $2
`, tenantSlug, version)
if err != nil {
return fmt.Errorf("crypto: kek-version sperren: %w", err)
}
if tag.RowsAffected() == 0 {
return ErrKEKVersionNotFound
}
return nil
}
+111
View File
@@ -0,0 +1,111 @@
// Integrationstest (ARC-08): echte Postgres-Instanz, folgt derselben
// Testhost-Konvention wie mail/internal/dedup/indexworker — TEST_TENANT_DSN.
package crypto
import (
"bytes"
"context"
"os"
"testing"
"github.com/jackc/pgx/v5/pgxpool"
)
var testLocalWrapKey = bytes.Repeat([]byte{0x7a}, KEKSize)
func setupKEKVersionStore(t *testing.T, tenantSlug string) *KEKVersionStore {
t.Helper()
dsn := os.Getenv("TEST_TENANT_DSN")
if dsn == "" {
t.Skip("TEST_TENANT_DSN nicht gesetzt, Integrationstest übersprungen")
}
ctx := context.Background()
pool, err := pgxpool.New(ctx, dsn)
if err != nil {
t.Fatalf("pool: %v", err)
}
t.Cleanup(func() { pool.Close() })
store := NewKEKVersionStore(pool, testLocalWrapKey)
if err := store.EnsureSchema(ctx); err != nil {
t.Fatalf("schema: %v", err)
}
t.Cleanup(func() {
_, _ = pool.Exec(context.Background(), `DELETE FROM mail_kek_versions WHERE tenant_slug = $1`, tenantSlug)
})
return store
}
func TestKEKVersionStore_RecordIfNewDetectsRotationOnly(t *testing.T) {
tenant := "mandant-arc08-recordifnew"
store := setupKEKVersionStore(t, tenant)
ctx := context.Background()
kekV1 := bytes.Repeat([]byte{0x01}, KEKSize)
v1, err := store.RecordIfNew(ctx, tenant, kekV1)
if err != nil {
t.Fatalf("erste erfassung: %v", err)
}
if v1 != 1 {
t.Fatalf("erwartete version 1, habe %d", v1)
}
// Erneuter Aufruf mit UNVERÄNDERTEM KEK darf keine neue Version anlegen.
vAgain, err := store.RecordIfNew(ctx, tenant, kekV1)
if err != nil {
t.Fatalf("zweite erfassung (unverändert): %v", err)
}
if vAgain != 1 {
t.Fatalf("erwartete weiterhin version 1 bei unverändertem kek, habe %d", vAgain)
}
kekV2 := bytes.Repeat([]byte{0x02}, KEKSize)
v2, err := store.RecordIfNew(ctx, tenant, kekV2)
if err != nil {
t.Fatalf("dritte erfassung (rotiert): %v", err)
}
if v2 != 2 {
t.Fatalf("erwartete version 2 nach rotation, habe %d", v2)
}
gotV1, err := store.Get(ctx, tenant, 1)
if err != nil {
t.Fatalf("get v1: %v", err)
}
if !bytes.Equal(gotV1, kekV1) {
t.Fatal("v1 liefert nicht den ursprünglichen kek zurück")
}
gotV2, err := store.Get(ctx, tenant, 2)
if err != nil {
t.Fatalf("get v2: %v", err)
}
if !bytes.Equal(gotV2, kekV2) {
t.Fatal("v2 liefert nicht den rotierten kek zurück")
}
}
func TestKEKVersionStore_RevokeBlocksOnlyThatVersion(t *testing.T) {
tenant := "mandant-arc08-revoke"
store := setupKEKVersionStore(t, tenant)
ctx := context.Background()
kekV1 := bytes.Repeat([]byte{0x11}, KEKSize)
kekV2 := bytes.Repeat([]byte{0x22}, KEKSize)
if _, err := store.RecordIfNew(ctx, tenant, kekV1); err != nil {
t.Fatalf("v1 erfassen: %v", err)
}
if _, err := store.RecordIfNew(ctx, tenant, kekV2); err != nil {
t.Fatalf("v2 erfassen: %v", err)
}
if err := store.Revoke(ctx, tenant, 1); err != nil {
t.Fatalf("v1 sperren: %v", err)
}
if _, err := store.Get(ctx, tenant, 1); err != ErrKEKVersionRevoked {
t.Fatalf("erwartete ErrKEKVersionRevoked für gesperrte version 1, habe: %v", err)
}
if _, err := store.Get(ctx, tenant, 2); err != nil {
t.Fatalf("version 2 sollte unberührt bleiben: %v", err)
}
}
+131
View File
@@ -0,0 +1,131 @@
// ARC-08: End-zu-Ende-Rotationstest. Nutzt einen fake KEKProvider (echter
// Testkonvention aus ARC-02, siehe encstorage_test.go) statt eines echten
// HTTP-Aufrufs an Core API-12 — Core selbst hat keine rotierbare
// Testschnittstelle über HTTP exponiert (nur der aktuelle KEK ist
// abrufbar), die Rotation wird hier auf Höhe der KEKProvider-Schnittstelle
// simuliert, exakt wie ARC-02 es für Fehlerfälle bereits tut.
package crypto
import (
"bytes"
"context"
"io"
"strings"
"sync"
"testing"
)
// rotatableKEKProvider liefert für einen Mandanten einen aktuell
// gesetzten KEK, der zur Laufzeit "rotiert" werden kann (simuliert Core
// API-10s RotateTenantKEK, dessen Effekt auf API-12 exakt darin besteht,
// dass TenantKEK ab dann einen anderen Wert liefert).
type rotatableKEKProvider struct {
mu sync.Mutex
current []byte
}
func (p *rotatableKEKProvider) TenantKEK(_ context.Context, _ string) ([]byte, error) {
p.mu.Lock()
defer p.mu.Unlock()
return p.current, nil
}
func (p *rotatableKEKProvider) rotate(newKEK []byte) {
p.mu.Lock()
defer p.mu.Unlock()
p.current = newKEK
}
// TestRotation_OldArchiveStaysReadableAfterMasterKeyRotation ist die
// geforderte Pflichtprüfung 1: Rotation des Hauptschlüssels lässt
// Altbestand weiterhin lesbar.
func TestRotation_OldArchiveStaysReadableAfterMasterKeyRotation(t *testing.T) {
tenant := "mandant-arc08-rotation-lesbar"
store := setupKEKVersionStore(t, tenant)
ctx := context.Background()
provider := &rotatableKEKProvider{current: bytes.Repeat([]byte{0x51}, KEKSize)}
svc := NewService(provider).WithVersionStore(store)
// Objekt VOR der Rotation versiegeln.
oldEnvelope, err := svc.Seal(ctx, tenant, strings.NewReader("altbestand vor rotation"))
if err != nil {
t.Fatalf("seal (alt): %v", err)
}
if oldEnvelope.KEKVersion != 1 {
t.Fatalf("erwartete kek-version 1 vor rotation, habe %d", oldEnvelope.KEKVersion)
}
oldCiphertext, err := io.ReadAll(oldEnvelope.Ciphertext)
if err != nil {
t.Fatalf("chiffretext (alt) lesen: %v", err)
}
// Core rotiert den Tenant-Hauptschlüssel — TenantKEK liefert ab jetzt
// einen komplett anderen Wert, exakt wie internal/kek.Store.
// RotateTenantKEK es real bei Core bewirkt.
provider.rotate(bytes.Repeat([]byte{0x52}, KEKSize))
// Objekt NACH der Rotation versiegeln (Akzeptanzkriterium 2: kein
// Neuverschlüsseln des Altbestands nötig, nur neue Objekte nutzen den
// neuen Schlüssel).
newEnvelope, err := svc.Seal(ctx, tenant, strings.NewReader("neuer inhalt nach rotation"))
if err != nil {
t.Fatalf("seal (neu): %v", err)
}
if newEnvelope.KEKVersion != 2 {
t.Fatalf("erwartete kek-version 2 nach rotation, habe %d", newEnvelope.KEKVersion)
}
// Altbestand bleibt über die aufgezeichnete Version lesbar
// (Akzeptanzkriterium 3).
openedOld, err := svc.OpenAtVersion(ctx, tenant, oldEnvelope.KEKVersion, oldEnvelope.WrappedDEK, bytes.NewReader(oldCiphertext))
if err != nil {
t.Fatalf("openatversion (alt, nach rotation): %v", err)
}
plainOld, err := io.ReadAll(openedOld)
if err != nil {
t.Fatalf("altbestand lesen: %v", err)
}
if string(plainOld) != "altbestand vor rotation" {
t.Fatalf("altbestand-inhalt stimmt nicht, habe %q", string(plainOld))
}
// Der naive Open() (aktueller KEK) darf für das ALTE Objekt inzwischen
// NICHT mehr funktionieren — das beweist, dass OpenAtVersion die
// Rotation tatsächlich überbrückt, statt zufällig auch so zu klappen.
if _, err := svc.Open(ctx, tenant, oldEnvelope.WrappedDEK, bytes.NewReader(oldCiphertext)); err == nil {
t.Fatal("erwartete fehler bei Open() des altbestands mit dem NEUEN aktuellen kek, habe nil")
}
}
// TestRotation_CompromisedOldKeyCanBeRevoked ist die geforderte
// Pflichtprüfung 2: kompromittierter alter Schlüssel kann gezielt
// gesperrt werden.
func TestRotation_CompromisedOldKeyCanBeRevoked(t *testing.T) {
tenant := "mandant-arc08-revoke-e2e"
store := setupKEKVersionStore(t, tenant)
ctx := context.Background()
provider := &rotatableKEKProvider{current: bytes.Repeat([]byte{0x61}, KEKSize)}
svc := NewService(provider).WithVersionStore(store)
compromisedEnvelope, err := svc.Seal(ctx, tenant, strings.NewReader("mit kompromittiertem schlüssel versiegelt"))
if err != nil {
t.Fatalf("seal: %v", err)
}
compromisedCiphertext, err := io.ReadAll(compromisedEnvelope.Ciphertext)
if err != nil {
t.Fatalf("chiffretext lesen: %v", err)
}
provider.rotate(bytes.Repeat([]byte{0x62}, KEKSize))
if err := store.Revoke(ctx, tenant, compromisedEnvelope.KEKVersion); err != nil {
t.Fatalf("kompromittierte version sperren: %v", err)
}
_, err = svc.OpenAtVersion(ctx, tenant, compromisedEnvelope.KEKVersion, compromisedEnvelope.WrappedDEK, bytes.NewReader(compromisedCiphertext))
if err == nil {
t.Fatal("erwartete fehler beim lesen mit gesperrter kek-version, habe nil")
}
}
+52 -5
View File
@@ -9,26 +9,41 @@ import (
// 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 (siehe mail/internal/encstorage).
// KEKVersion identifiziert (ARC-08), MIT welcher Tenant-KEK-Version der
// DEK verpackt wurde — 0, solange kein KEKVersionStore konfiguriert ist
// (Rückwärtskompatibilität, siehe WithVersionStore).
type Envelope struct {
Ciphertext io.Reader
WrappedDEK []byte
KEKVersion int
}
// Service verbindet KEKProvider mit den Envelope-Operationen — Aufrufer
// (mail/internal/encstorage) rufen ausschließlich Service auf, nie die
// Einzelfunktionen aus envelope.go direkt.
type Service struct {
kek KEKProvider
kek KEKProvider
versions *KEKVersionStore
}
func NewService(kek KEKProvider) *Service {
return &Service{kek: kek}
}
// WithVersionStore aktiviert die Tenant-KEK-Versionsverfolgung (ARC-08).
// Ohne aufgerufenes WithVersionStore verhält sich Service exakt wie vor
// ARC-08 (KEKVersion bleibt 0, OpenAtVersion fällt auf Open zurück) —
// bestehende Aufrufer (z. B. encstorage) sind unverändert lauffähig.
func (s *Service) WithVersionStore(store *KEKVersionStore) *Service {
s.versions = store
return s
}
// Seal erzeugt einen neuen DEK (Akzeptanzkriterium 1), verschlüsselt
// 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).
// bezogen, nie zwischengespeichert außer in der optionalen
// KEK-Versionshistorie für spätere Altbestands-Lesezugriffe).
func (s *Service) Seal(ctx context.Context, tenantSlug string, plaintext io.Reader) (*Envelope, error) {
dek, err := GenerateDEK()
if err != nil {
@@ -46,14 +61,25 @@ func (s *Service) Seal(ctx context.Context, tenantSlug string, plaintext io.Read
if err != nil {
return nil, err
}
return &Envelope{Ciphertext: ciphertext, WrappedDEK: wrappedDEK}, nil
var kekVersion int
if s.versions != nil {
kekVersion, err = s.versions.RecordIfNew(ctx, tenantSlug, kek)
if err != nil {
return nil, fmt.Errorf("crypto: kek-version erfassen: %w", err)
}
}
return &Envelope{Ciphertext: ciphertext, WrappedDEK: wrappedDEK, KEKVersion: kekVersion}, nil
}
// Open entpackt den DEK mit dem aktuellen Tenant-KEK (Akzeptanzkriterium
// Open entpackt den DEK mit dem AKTUELLEN Tenant-KEK (Akzeptanzkriterium
// 3: nur mit gültigem, mandantenbezogenem Schlüssel möglich — ein
// falscher Tenant-Slug liefert entweder einen falschen KEK von Core
// [dann schlägt UnwrapDEK fehl] oder Core verweigert den Zugriff direkt)
// und entschlüsselt ciphertext damit.
// und entschlüsselt ciphertext damit. Nach einer Tenant-KEK-Rotation bei
// Core funktioniert Open nur noch für Objekte, die mit dem NEUEN KEK
// versiegelt wurden — für Altbestand siehe OpenAtVersion.
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 {
@@ -65,3 +91,24 @@ func (s *Service) Open(ctx context.Context, tenantSlug string, wrappedDEK []byte
}
return DecryptStream(dek, ciphertext)
}
// OpenAtVersion entpackt den DEK mit der beim Seal aufgezeichneten
// historischen Tenant-KEK-Version statt mit dem aktuellen Core-KEK
// (ARC-08 Akzeptanzkriterium 3: Altbestand bleibt nach einer
// Hauptschlüssel-Rotation lesbar). Ist kekVersion 0 oder kein
// KEKVersionStore konfiguriert, verhält es sich wie Open (Rückwärts-
// kompatibilität für vor ARC-08 versiegelte Objekte).
func (s *Service) OpenAtVersion(ctx context.Context, tenantSlug string, kekVersion int, wrappedDEK []byte, ciphertext io.Reader) (io.Reader, error) {
if kekVersion == 0 || s.versions == nil {
return s.Open(ctx, tenantSlug, wrappedDEK, ciphertext)
}
kek, err := s.versions.Get(ctx, tenantSlug, kekVersion)
if err != nil {
return nil, fmt.Errorf("crypto: historischen tenant-kek beziehen: %w", err)
}
dek, err := UnwrapDEK(kek, wrappedDEK)
if err != nil {
return nil, err
}
return DecryptStream(dek, ciphertext)
}
+46 -1
View File
@@ -16,8 +16,10 @@ package encstorage
import (
"bytes"
"context"
"errors"
"fmt"
"io"
"strconv"
"gitea.perlbach24.de/scripte/nexarch/mail/internal/crypto"
"gitea.perlbach24.de/scripte/nexarch/mail/internal/storage"
@@ -30,6 +32,14 @@ func wrappedDEKKey(key string) string {
return key + ".dek"
}
// kekVersionKey ist der Sidecar-Objektschlüssel für die Tenant-KEK-
// Version, mit der der DEK verpackt wurde (ARC-08 Akzeptanzkriterium 3).
// Fehlt dieser Sidecar (vor ARC-08 geschriebene Objekte), wird Version 0
// angenommen — GetDecrypted verhält sich dann wie vor ARC-08.
func kekVersionKey(key string) string {
return key + ".dek.version"
}
// Service verbindet Storage (ARC-01) und Crypto (ARC-02): der Rest von
// Mail ruft AUSSCHLIESSLICH diesen Service auf, nie storage.Service
// direkt mit Klartext — das verhindert einen Schreibpfad, der die
@@ -62,6 +72,12 @@ func (s *Service) Put(ctx context.Context, tenantSlug, key string, plaintext io.
if _, err := s.storage.Put(ctx, wrappedDEKKey(key), bytes.NewReader(env.WrappedDEK), int64(len(env.WrappedDEK)), "application/octet-stream"); err != nil {
return fmt.Errorf("encstorage: verpackten dek speichern: %w", err)
}
if env.KEKVersion != 0 {
versionBytes := []byte(strconv.Itoa(env.KEKVersion))
if _, err := s.storage.Put(ctx, kekVersionKey(key), bytes.NewReader(versionBytes), int64(len(versionBytes)), "text/plain"); err != nil {
return fmt.Errorf("encstorage: kek-version speichern: %w", err)
}
}
return nil
}
@@ -85,9 +101,38 @@ func (s *Service) GetDecrypted(ctx context.Context, tenantSlug, key string) ([]b
return nil, fmt.Errorf("encstorage: verpackten dek lesen: %w", err)
}
plaintextReader, err := s.crypto.Open(ctx, tenantSlug, wrappedDEK, bytes.NewReader(ciphertext))
kekVersion, err := s.readKEKVersion(ctx, key)
if err != nil {
return nil, err
}
plaintextReader, err := s.crypto.OpenAtVersion(ctx, tenantSlug, kekVersion, wrappedDEK, bytes.NewReader(ciphertext))
if err != nil {
return nil, err
}
return io.ReadAll(plaintextReader)
}
// readKEKVersion liest den Versions-Sidecar (ARC-08). Fehlt er (vor
// ARC-08 geschriebene Objekte, oder ein Seal ohne konfigurierten
// KEKVersionStore), gilt Version 0 — crypto.Service.OpenAtVersion fällt
// dafür auf das unveränderte Open-Verhalten zurück.
func (s *Service) readKEKVersion(ctx context.Context, key string) (int, error) {
reader, err := s.storage.Get(ctx, kekVersionKey(key))
if err != nil {
if errors.Is(err, storage.ErrNotFound) {
return 0, nil
}
return 0, fmt.Errorf("encstorage: kek-version lesen: %w", err)
}
defer func() { _ = reader.Close() }()
raw, err := io.ReadAll(reader)
if err != nil {
return 0, fmt.Errorf("encstorage: kek-version lesen: %w", err)
}
version, err := strconv.Atoi(string(raw))
if err != nil {
return 0, fmt.Errorf("encstorage: kek-version parsen: %w", err)
}
return version, nil
}
+56 -16
View File
@@ -204,6 +204,52 @@ var fieldWeights = map[string]any{
FieldAttachmentText: 1,
}
// doSearchWithSwapRetry führt eine /search-Anfrage aus und wiederholt sie
// bis zu zweimal mit kurzer Pause, falls Manticore "unknown local table"
// meldet (SRC-09 Akzeptanzkriterium 3: der Reindex-Umschaltmoment
// RENAME-alte-Tabelle-weg/RENAME-neue-Tabelle-rein hat ein extrem kurzes
// Zeitfenster ohne existierende mail_documents-Tabelle — dieser Retry
// überbrückt es, statt eine Suchanfrage in genau diesem Moment fehlschlagen
// zu lassen).
func (c *Client) doSearchWithSwapRetry(ctx context.Context, body []byte) ([]byte, error) {
const maxAttempts = 3
var lastErr error
for attempt := 0; attempt < maxAttempts; attempt++ {
if attempt > 0 {
select {
case <-ctx.Done():
return nil, ctx.Err()
case <-time.After(20 * time.Millisecond):
}
}
req, err := http.NewRequestWithContext(ctx, http.MethodPost, c.baseURL+"/search", bytes.NewReader(body))
if err != nil {
return nil, fmt.Errorf("search: suchanfrage bauen: %w", err)
}
req.Header.Set("Content-Type", "application/json")
resp, err := c.http.Do(req)
if err != nil {
return nil, fmt.Errorf("search: suche ausführen: %w", err)
}
respBody, readErr := io.ReadAll(resp.Body)
_ = resp.Body.Close()
if readErr != nil {
return nil, fmt.Errorf("search: antwort lesen: %w", readErr)
}
if strings.Contains(string(respBody), "unknown local table") {
lastErr = fmt.Errorf("search: suche, status %d: %s", resp.StatusCode, string(respBody))
continue
}
if resp.StatusCode != http.StatusOK {
return nil, fmt.Errorf("search: suche, status %d: %s", resp.StatusCode, string(respBody))
}
return respBody, nil
}
return nil, lastErr
}
// Search sucht queryText innerhalb der Volltextfelder, strikt begrenzt auf
// den Mandanten tenantSlug (Akzeptanzkriterium 2: mandantengetrennt
// abfragbar) — der Tenant-Filter läuft über ein strukturiertes "equals"-
@@ -216,6 +262,8 @@ var fieldWeights = map[string]any{
// Feld-/Tabellennamen, der beeinflusst werden könnte. Ergebnisse kommen
// von Manticore bereits nach Relevanz (BM25, gewichtet über fieldWeights)
// absteigend sortiert zurück (Akzeptanzkriterium 1).
const searchResultLimit = 1000
func (c *Client) Search(ctx context.Context, tenantSlug, queryText string) ([]Result, error) {
payload := map[string]any{
"index": IndexName,
@@ -230,29 +278,21 @@ func (c *Client) Search(ctx context.Context, tenantSlug, queryText string) ([]Re
"options": map[string]any{
"field_weights": fieldWeights,
},
// Ohne explizites limit begrenzt Manticore standardmäßig auf 20
// Treffer — bei Testkorpora bis 1000 Dokumenten (SRC-03) blieb das
// bisher unbemerkt, da nur auf das Vorhandensein einzelner Treffer
// geprüft wurde, nicht auf die Gesamtzahl. searchResultLimit deckt
// realistische Trefferlisten ab, ohne unbegrenzt zu sein.
"limit": searchResultLimit,
}
body, err := json.Marshal(payload)
if err != nil {
return nil, fmt.Errorf("search: suchanfrage serialisieren: %w", err)
}
req, err := http.NewRequestWithContext(ctx, http.MethodPost, c.baseURL+"/search", bytes.NewReader(body))
respBody, err := c.doSearchWithSwapRetry(ctx, body)
if err != nil {
return nil, fmt.Errorf("search: suchanfrage bauen: %w", err)
}
req.Header.Set("Content-Type", "application/json")
resp, err := c.http.Do(req)
if err != nil {
return nil, fmt.Errorf("search: suche ausführen: %w", err)
}
defer func() { _ = resp.Body.Close() }()
respBody, err := io.ReadAll(resp.Body)
if err != nil {
return nil, fmt.Errorf("search: antwort lesen: %w", err)
}
if resp.StatusCode != http.StatusOK {
return nil, fmt.Errorf("search: suche, status %d: %s", resp.StatusCode, string(respBody))
return nil, err
}
var parsed searchResponse
+288
View File
@@ -0,0 +1,288 @@
// SRC-09: Suchindex-Neuaufbau/Reindexierung. Baut eine neue physische
// Manticore-Tabelle auf, kopiert alle Dokumente aus der aktuell lebenden
// Tabelle (Konsistenzwiederherstellung), verifiziert die Trefferzahl und
// tauscht erst danach per Manticore RENAME atomar um — die alte Tabelle
// bleibt bis zu diesem Moment vollständig abfragbar (Akzeptanzkriterium 3),
// RENAME ist eine reine Metadaten-Operation ohne Suchausfall.
package search
import (
"bytes"
"context"
"encoding/json"
"fmt"
"io"
"net/http"
"regexp"
"time"
)
// tempTableNamePattern begrenzt generierte Zwischentabellennamen auf ein
// festes Präfix + Ziffern — auch wenn der Name ausschließlich von diesem
// Paket selbst erzeugt wird (kein externer Eingabewert erreicht ihn),
// erzwingt die Prüfung strukturell, dass niemals ein beliebiger String an
// dieser Stelle landen kann (Verteidigung in der Tiefe, gleiche Haltung
// wie das Feld-Whitelist-Prinzip in fields.go).
var tempTableNamePattern = regexp.MustCompile(`^mail_documents_reindex_[0-9]+$`)
// Progress meldet den Fortschritt eines laufenden Reindex
// (Akzeptanzkriterium 2: Fortschritt nachvollziehbar sichtbar).
type Progress struct {
Copied int64
Total int64
}
// Reindexer baut den Suchindex vollständig neu auf.
type Reindexer struct {
client *Client
}
func NewReindexer(client *Client) *Reindexer {
return &Reindexer{client: client}
}
// RebuildResult fasst das Ergebnis eines abgeschlossenen Reindex zusammen.
type RebuildResult struct {
OldCount int64
NewCount int64
}
// Rebuild baut den Index vollständig neu auf: neue Tabelle anlegen, alle
// Dokumente aus der aktuell lebenden Tabelle seitenweise kopieren
// (Akzeptanzkriterium 1: kein Datenverlust im laufenden Betrieb — die
// lebende Tabelle wird dabei nur gelesen, nie verändert), Trefferzahlen
// vergleichen, dann atomar per RENAME umschalten. Schlägt ein Schritt vor
// dem Umschalten fehl (z. B. abgebrochener Kontext), wird die
// Zwischentabelle entfernt und die lebende Tabelle bleibt unverändert
// (Akzeptanzkriterium 3 / Pflichtprüfung 2: kein inkonsistenter Zustand).
func (r *Reindexer) Rebuild(ctx context.Context, onProgress func(Progress)) (RebuildResult, error) {
tempTable := fmt.Sprintf("mail_documents_reindex_%d", time.Now().UnixNano())
if !tempTableNamePattern.MatchString(tempTable) {
return RebuildResult{}, fmt.Errorf("search: erzeugter zwischentabellenname unerwartet ungültig: %q", tempTable)
}
if err := r.client.runSchemaSQL(ctx, buildCreateTableSQL(tempTable)); err != nil {
return RebuildResult{}, fmt.Errorf("search: zwischentabelle anlegen: %w", err)
}
oldTotal, err := r.copyAll(ctx, IndexName, tempTable, onProgress)
if err != nil {
_ = r.client.runSchemaSQL(context.Background(), "DROP TABLE "+tempTable)
return RebuildResult{}, fmt.Errorf("search: dokumente kopieren: %w", err)
}
// Manticore macht frisch eingefügte Dokumente einer neu angelegten
// RT-Tabelle für Volltext-/match_all-Zählungen erst nach einem
// expliziten FLUSH RAMCHUNK zuverlässig sichtbar (beobachtet: SELECT
// über SQL sieht die Zeile sofort, /search match_all zählt sie ohne
// Flush als 0). Vor der Konsistenzprüfung zwingend nötig.
if err := r.client.runSchemaSQL(ctx, "FLUSH RAMCHUNK "+tempTable); err != nil {
_ = r.client.runSchemaSQL(context.Background(), "DROP TABLE "+tempTable)
return RebuildResult{}, fmt.Errorf("search: zwischentabelle flushen: %w", err)
}
newTotal, err := r.countAll(ctx, tempTable)
if err != nil {
_ = r.client.runSchemaSQL(context.Background(), "DROP TABLE "+tempTable)
return RebuildResult{}, fmt.Errorf("search: neue tabelle zählen: %w", err)
}
if newTotal != oldTotal {
_ = r.client.runSchemaSQL(context.Background(), "DROP TABLE "+tempTable)
return RebuildResult{}, fmt.Errorf("search: trefferzahlen weichen ab (alt %d, neu %d), kein umschalten", oldTotal, newTotal)
}
retiredTable := fmt.Sprintf("mail_documents_retired_%d", time.Now().UnixNano())
if err := r.client.runSchemaSQL(ctx, fmt.Sprintf("ALTER TABLE %s RENAME %s", IndexName, retiredTable)); err != nil {
_ = r.client.runSchemaSQL(context.Background(), "DROP TABLE "+tempTable)
return RebuildResult{}, fmt.Errorf("search: alte tabelle umbenennen: %w", err)
}
if err := r.client.runSchemaSQL(ctx, fmt.Sprintf("ALTER TABLE %s RENAME %s", tempTable, IndexName)); err != nil {
// Kritischer Zustand: alte Tabelle bereits umbenannt, neue kann
// nicht einspringen. Umschalten rückgängig machen, statt ohne
// abfragbaren Index dazustehen.
_ = r.client.runSchemaSQL(context.Background(), fmt.Sprintf("ALTER TABLE %s RENAME %s", retiredTable, IndexName))
_ = r.client.runSchemaSQL(context.Background(), "DROP TABLE "+tempTable)
return RebuildResult{}, fmt.Errorf("search: neue tabelle aktivieren: %w", err)
}
_ = r.client.runSchemaSQL(context.Background(), "DROP TABLE "+retiredTable)
return RebuildResult{OldCount: oldTotal, NewCount: newTotal}, nil
}
const copyPageSize = 200
func (r *Reindexer) copyAll(ctx context.Context, sourceIndex, targetIndex string, onProgress func(Progress)) (int64, error) {
total, err := r.countAll(ctx, sourceIndex)
if err != nil {
return 0, err
}
var cursor uint64
var copied int64
first := true
for {
if err := ctx.Err(); err != nil {
return 0, err
}
page, err := r.fetchPage(ctx, sourceIndex, cursor, first)
if err != nil {
return 0, err
}
first = false
if len(page) == 0 {
break
}
for _, doc := range page {
if err := r.putRaw(ctx, targetIndex, doc); err != nil {
return 0, err
}
cursor = doc.ID
copied++
}
if onProgress != nil {
onProgress(Progress{Copied: copied, Total: total})
}
}
return copied, nil
}
func (r *Reindexer) countAll(ctx context.Context, index string) (int64, error) {
payload := map[string]any{"index": index, "query": map[string]any{"match_all": map[string]any{}}, "limit": 0}
var parsed struct {
Hits struct {
Total int64 `json:"total"`
} `json:"hits"`
}
if err := r.client.postJSON(ctx, "/search", payload, &parsed); err != nil {
return 0, err
}
return parsed.Hits.Total, nil
}
type scrollHit struct {
ID uint64 `json:"_id"`
Source json.RawMessage `json:"_source"`
}
func (r *Reindexer) fetchPage(ctx context.Context, index string, afterID uint64, first bool) ([]scrollHit, error) {
must := []map[string]any{}
if !first {
must = append(must, map[string]any{"range": map[string]any{"id": map[string]any{"gt": afterID}}})
}
query := map[string]any{"match_all": map[string]any{}}
if len(must) > 0 {
query = map[string]any{"bool": map[string]any{"must": must}}
}
payload := map[string]any{
"index": index,
"query": query,
"sort": []map[string]any{{"id": "asc"}},
"limit": copyPageSize,
}
var parsed struct {
Hits struct {
Hits []scrollHit `json:"hits"`
} `json:"hits"`
}
if err := r.client.postJSON(ctx, "/search", payload, &parsed); err != nil {
return nil, err
}
return parsed.Hits.Hits, nil
}
func (r *Reindexer) putRaw(ctx context.Context, index string, doc scrollHit) error {
payload := map[string]any{
"index": index,
"id": doc.ID,
"doc": json.RawMessage(doc.Source),
}
return r.client.postJSON(ctx, "/replace", payload, nil)
}
// postJSON/buildCreateTableSQL sind bewusst hier statt in client.go
// angesiedelt: der übrige Suchpfad (Search/Facets) fasst niemals einen
// Tabellennamen dynamisch an, Reindex ist die einzige Stelle im Paket, die
// das operativ tun muss.
func (c *Client) postJSON(ctx context.Context, path string, payload any, out any) error {
body, err := json.Marshal(payload)
if err != nil {
return fmt.Errorf("payload serialisieren: %w", err)
}
req, err := http.NewRequestWithContext(ctx, http.MethodPost, c.baseURL+path, bytes.NewReader(body))
if err != nil {
return fmt.Errorf("anfrage bauen: %w", err)
}
req.Header.Set("Content-Type", "application/json")
resp, err := c.http.Do(req)
if err != nil {
return fmt.Errorf("ausführen: %w", err)
}
defer func() { _ = resp.Body.Close() }()
respBody, err := io.ReadAll(resp.Body)
if err != nil {
return fmt.Errorf("antwort lesen: %w", err)
}
if resp.StatusCode != http.StatusOK {
return fmt.Errorf("status %d: %s", resp.StatusCode, string(respBody))
}
if out == nil {
return nil
}
if err := json.Unmarshal(respBody, out); err != nil {
return fmt.Errorf("antwort parsen: %w", err)
}
return nil
}
// showTables listet die vorhandenen Manticore-Tabellen (Diagnose-/
// Testhilfe, um verwaiste Zwischentabellen nach einem Abbruch
// auszuschließen — Pflichtprüfung 2).
func (c *Client) showTables(ctx context.Context) ([]string, error) {
form := "query=SHOW TABLES"
req, err := http.NewRequestWithContext(ctx, http.MethodPost, c.baseURL+"/sql?mode=raw", bytes.NewReader([]byte(form)))
if err != nil {
return nil, fmt.Errorf("anfrage bauen: %w", err)
}
req.Header.Set("Content-Type", "application/x-www-form-urlencoded")
resp, err := c.http.Do(req)
if err != nil {
return nil, fmt.Errorf("ausführen: %w", err)
}
defer func() { _ = resp.Body.Close() }()
body, err := io.ReadAll(resp.Body)
if err != nil {
return nil, fmt.Errorf("antwort lesen: %w", err)
}
var parsed []struct {
Data []struct {
Table string `json:"Table"`
} `json:"data"`
}
if err := json.Unmarshal(body, &parsed); err != nil {
return nil, fmt.Errorf("antwort parsen: %w", err)
}
names := []string{}
if len(parsed) > 0 {
for _, row := range parsed[0].Data {
names = append(names, row.Table)
}
}
return names, nil
}
// buildCreateTableSQL erzeugt die Schema-DDL für eine Zwischentabelle mit
// demselben Spaltensatz wie mail_documents (Basis + Facettenfelder aus
// SRC-05). tableName ist über tempTableNamePattern in Rebuild bereits
// geprüft, bevor diese Funktion aufgerufen wird.
func buildCreateTableSQL(tableName string) string {
return fmt.Sprintf(
"CREATE TABLE %s (%s string attribute indexed, %s string attribute indexed, %s text, %s text, %s text, %s string attribute indexed, %s string attribute indexed, %s string attribute indexed, %s string attribute indexed, %s timestamp)",
tableName,
FieldTenantSlug, FieldMessageID, FieldSubject, FieldBody, FieldAttachmentText,
FieldSender, FieldMailbox, FieldAttachmentType, FieldTag, FieldSentAt,
)
}
+215
View File
@@ -0,0 +1,215 @@
// Integrationstest (SRC-09): echte Manticore-Instanz, TEST_MANTICORE_URL
// (gleiche Konvention wie integration_test.go/ranking_test.go/facets_test.go).
package search
import (
"context"
"errors"
"sync"
"sync/atomic"
"testing"
"time"
)
// TestRebuild_SearchKeepsWorkingDuringReindex ist die geforderte
// Pflichtprüfung 1: Reindex während laufender Suchanfragen unterbricht die
// Suche nicht.
func TestRebuild_SearchKeepsWorkingDuringReindex(t *testing.T) {
client := setupClient(t)
ctx := context.Background()
tenant := "mandant-src09-parallel"
for i := 0; i < 30; i++ {
messageID := "msg-parallel-" + string(rune('a'+i))
indexFacetDoc(t, client, ctx, tenant, Document{MessageID: messageID, Subject: "Zwiebelfisch " + messageID, Body: "Text"})
}
stop := make(chan struct{})
var searchErrors int64
var searchesDone int64
var wg sync.WaitGroup
wg.Add(1)
go func() {
defer wg.Done()
for {
select {
case <-stop:
return
default:
}
if _, err := client.Search(ctx, tenant, "Zwiebelfisch"); err != nil {
atomic.AddInt64(&searchErrors, 1)
t.Logf("suchfehler während reindex: %v", err)
}
atomic.AddInt64(&searchesDone, 1)
time.Sleep(2 * time.Millisecond)
}
}()
reindexer := NewReindexer(client)
result, err := reindexer.Rebuild(ctx, nil)
close(stop)
wg.Wait()
if err != nil {
t.Fatalf("rebuild: %v", err)
}
if result.OldCount != result.NewCount {
t.Fatalf("erwartete gleiche trefferzahlen, habe alt=%d neu=%d", result.OldCount, result.NewCount)
}
if atomic.LoadInt64(&searchesDone) == 0 {
t.Fatal("keine einzige parallele suche ausgeführt — test aussagelos")
}
if errs := atomic.LoadInt64(&searchErrors); errs != 0 {
t.Fatalf("erwartete 0 fehlgeschlagene suchen während des reindex, habe %d von %d", errs, atomic.LoadInt64(&searchesDone))
}
// Suche funktioniert auch NACH dem Umschalten weiterhin real.
afterResults, err := client.Search(ctx, tenant, "Zwiebelfisch")
if err != nil {
t.Fatalf("search nach reindex: %v", err)
}
if len(afterResults) != 30 {
t.Fatalf("erwartete 30 treffer nach reindex, habe %d", len(afterResults))
}
}
// TestRebuild_AbortedReindexLeavesNoInconsistentState ist die geforderte
// Pflichtprüfung 2: abgebrochener Reindex hinterlässt keinen
// inkonsistenten Zustand.
func TestRebuild_AbortedReindexLeavesNoInconsistentState(t *testing.T) {
client := setupClient(t)
ctx := context.Background()
tenant := "mandant-src09-abbruch"
indexFacetDoc(t, client, ctx, tenant, Document{MessageID: "msg-abbruch-1", Subject: "Vertragsentwurf Abbruchtest", Body: "Text"})
before, err := client.Search(ctx, tenant, "Abbruchtest")
if err != nil || len(before) != 1 {
t.Fatalf("voraussetzung nicht erfüllt: %v / %d treffer", err, len(before))
}
cancelCtx, cancel := context.WithCancel(ctx)
cancel() // sofort abgebrochen, simuliert Absturz/Abbruch mitten im Kopiervorgang
reindexer := NewReindexer(client)
_, err = reindexer.Rebuild(cancelCtx, nil)
if err == nil {
t.Fatal("erwartete fehler bei abgebrochenem kontext, habe nil")
}
if !errors.Is(err, context.Canceled) {
// Manticore-Fehler durch den abgebrochenen Request sind ebenfalls
// akzeptabel, solange überhaupt ein Fehler zurückkommt.
t.Logf("fehler war nicht context.Canceled, sondern: %v (akzeptiert, solange real ein fehler zurückkommt)", err)
}
// Die lebende Tabelle muss trotz Abbruch unverändert und abfragbar sein.
after, err := client.Search(ctx, tenant, "Abbruchtest")
if err != nil {
t.Fatalf("search nach abgebrochenem reindex: %v", err)
}
if len(after) != 1 {
t.Fatalf("erwartete weiterhin 1 treffer nach abgebrochenem reindex, habe %d — inkonsistenter zustand", len(after))
}
// Keine verwaisten Zwischentabellen (kein inkonsistenter Zustand auf
// Manticore-Ebene): kurz warten, damit ein eventuell noch laufender
// CREATE-TABLE-Aufruf durchlaufen kann, dann prüfen, dass keine
// mail_documents_reindex_*-Tabelle übrig geblieben ist.
time.Sleep(200 * time.Millisecond)
orphaned := listOrphanedReindexTables(t, client)
if len(orphaned) > 0 {
t.Fatalf("verwaiste zwischentabellen nach abbruch gefunden: %v", orphaned)
}
}
func listOrphanedReindexTables(t *testing.T, client *Client) []string {
t.Helper()
// SHOW TABLES ist eine feste, unparametrisierte Anweisung ohne
// jeglichen Laufzeitwert.
rows, err := client.showTables(context.Background())
if err != nil {
t.Fatalf("show tables: %v", err)
}
names := []string{}
for _, table := range rows {
if tempTableNamePattern.MatchString(table) {
names = append(names, table)
}
}
return names
}
// TestRebuild_SampleComparisonMatchesOldAndNewIndex ist die geforderte
// Pflichtprüfung 3: Stichprobenvergleich Alt-/Neuindex bestätigt gleiche
// Trefferzahlen.
func TestRebuild_SampleComparisonMatchesOldAndNewIndex(t *testing.T) {
client := setupClient(t)
ctx := context.Background()
tenant := "mandant-src09-stichprobe"
subjects := []string{"Quartalsbericht", "Personalplanung", "Urlaubsantrag"}
for i, s := range subjects {
indexFacetDoc(t, client, ctx, tenant, Document{MessageID: "msg-sp-" + string(rune('a'+i)), Subject: s, Body: "Inhalt " + s})
}
beforeCounts := map[string]int{}
for _, s := range subjects {
results, err := client.Search(ctx, tenant, s)
if err != nil {
t.Fatalf("search vor reindex (%s): %v", s, err)
}
beforeCounts[s] = len(results)
}
reindexer := NewReindexer(client)
if _, err := reindexer.Rebuild(ctx, nil); err != nil {
t.Fatalf("rebuild: %v", err)
}
for _, s := range subjects {
results, err := client.Search(ctx, tenant, s)
if err != nil {
t.Fatalf("search nach reindex (%s): %v", s, err)
}
if len(results) != beforeCounts[s] {
t.Fatalf("stichprobe %q: vor reindex %d treffer, nach reindex %d treffer", s, beforeCounts[s], len(results))
}
}
}
// TestRebuild_ReportsProgress deckt Akzeptanzkriterium 2 ab (Fortschritt
// nachvollziehbar sichtbar).
func TestRebuild_ReportsProgress(t *testing.T) {
client := setupClient(t)
ctx := context.Background()
tenant := "mandant-src09-fortschritt"
for i := 0; i < 5; i++ {
indexFacetDoc(t, client, ctx, tenant, Document{MessageID: "msg-progress-" + string(rune('a'+i)), Subject: "x"})
}
var updates []Progress
var mu sync.Mutex
reindexer := NewReindexer(client)
_, err := reindexer.Rebuild(ctx, func(p Progress) {
mu.Lock()
updates = append(updates, p)
mu.Unlock()
})
if err != nil {
t.Fatalf("rebuild: %v", err)
}
if len(updates) == 0 {
t.Fatal("erwartete mindestens eine fortschrittsmeldung")
}
last := updates[len(updates)-1]
if last.Copied < last.Total {
// total ist eine zu Beginn eingefrorene Momentaufnahme; die geteilte
// Manticore-Instanz kann während des Kopierens durch andere Tests
// weiter wachsen (real beobachtet) — copied darf total daher
// erreichen oder minimal überschreiten, nur ein Rückstand wäre ein
// echter Fehler.
t.Fatalf("letzte fortschrittsmeldung sollte abgeschlossen sein, habe copied=%d total=%d", last.Copied, last.Total)
}
}