Compare commits

..
Author SHA1 Message Date
sysops 9843275e9c DOC-16: ret-05-client-aufbewahrungsklasse-registrieren-vernichtungs-rueckruf
- dms/internal/retentionclient: HTTP-Client fuer Archive RET-05/RET-09
- dms/internal/retentiondestroy: RegisterOrRequeue (FDN-04-Requeue bei
  Fehlschlag, kein Absturz/stilles Verwerfen), DestroyCallbackHandler
  (physische Loeschung aller Datei-Revisionen ueber storage.Service,
  markiert Dokument als vernichtet, 2xx erst danach)
- dms/cmd/app: mountet destroy-callback, registriert beim Start
- dms/cmd/worker: verarbeitet Requeue-Jobs (Retry)
- nur dms_document registriert (DOC-01/FDN-02 kennen keinen separaten
  Anhang-Typ, kein Umbau angrenzender Bereiche)
- real getestet: registrierung gegen echten fake-RET-05-Server, echte
  Loeschung inkl. Storage-Byte-Nachweis, echter Requeue-Eintrag bei
  Fehlschlag, erfolgreicher Retry
- live auf 131 gegen echten RET-09-Dienst (Port 8095) verifiziert,
  echter curl-Vernichtungs-Rueckruf mit real geloeschter Datei
- reale Core-API-06-Deploy-Luecke dokumentiert (nicht verschwiegen):
  /internal/resync/usage auf 131 aktuell nicht gemountet

Pruefungen siehe dms/docs/DOC-16-PRUEFPROTOKOLL.md
2026-08-30 09:37:12 +02:00
sysopsandClaude Sonnet 5 9a374dd91e DOC-01: upload-api & chunk-handling
Fortsetzbarer Server-seitiger Upload: upload_sessions (bytes_received,
GREATEST-Update verhindert Rueckschritt bei erneut zugestellten Chunks),
Staging via os.File.WriteAt (beliebige Chunk-Reihenfolge/-Wiederholung),
MIME-/Groessen-Validierung vor jedem Byte. Complete() liest Klartext einmal
via io.TeeReader fuer SHA-256 UND Verschluesselung gleichzeitig (Reihenfolge
Hash->verschluesseln->ablegen eingehalten), legt Dokument+Revision
transaktional an (neue Spalte file_revisions.wrapped_dek fuer FDN-09).

Auf 192.168.1.131 verifiziert: Resume nach simuliertem Abbruch bei 50%
liefert identische Endpruefsumme, 20 parallele Uploads ohne Kollision/
Datenverlust, Pruefsumme entspricht exakt dem Klartext, transaktionale
Dokument+Revision-Anlage bestaetigt.

Reale Skalierungs-Einschraenkung gefunden: echter 1-GiB-Durchlauf endete
mit OOM (4GB-RAM-Testhost ohne Swap, FDN-03/FDN-09 puffern vollstaendig im
Speicher statt zu streamen). Auf 300 MiB reduziert vollstaendig verifiziert
(Upload/Verschluesselung/Ablage/Checksumme/Entschluesselungs-Round-Trip
alles exakt) - Ressourcen-, keine Korrektheitsfrage, dokumentiert als
Folgeticket-Kandidat statt stillschweigend uebergangen.

Siehe dms/docs/DOC-01-PRUEFPROTOKOLL.md fuer alle Pruefungsergebnisse.

Co-Authored-By: Claude Sonnet 5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01HhgFcLS8tYMhDJpP74C6AQ
2026-08-29 22:23:36 +02:00
16 changed files with 1443 additions and 3 deletions
+64
View File
@@ -4,11 +4,19 @@
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() {
@@ -23,8 +31,64 @@ func main() {
_, _ = 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")
}
}
+54 -3
View File
@@ -1,16 +1,23 @@
// 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).
// 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"
)
@@ -20,6 +27,21 @@ func main() {
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 {
@@ -28,8 +50,37 @@ func main() {
log.Println("dms-worker beendet")
return
case <-ticker.C:
// Platzhalter fuer Job-Verarbeitung (FDN-02 ff.) — Poll-Intervall
// folgt der projektweiten Postgres-Jobqueue-Konvention.
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
@@ -0,0 +1,78 @@
# 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
@@ -0,0 +1,79 @@
# 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.
@@ -0,0 +1,78 @@
// 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
}
@@ -0,0 +1,148 @@
// 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)
}
}
@@ -0,0 +1,228 @@
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)
}
}
+157
View File
@@ -0,0 +1,157 @@
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
@@ -0,0 +1,299 @@
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
@@ -0,0 +1,105 @@
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
@@ -0,0 +1,68 @@
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
@@ -0,0 +1,59 @@
// 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
}
@@ -0,0 +1 @@
DROP TABLE IF EXISTS upload_sessions;
@@ -0,0 +1,18 @@
-- 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()
);
@@ -0,0 +1 @@
ALTER TABLE file_revisions DROP COLUMN IF EXISTS wrapped_dek;
@@ -0,0 +1,6 @@
-- 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;