Compare commits
1
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
9843275e9c |
@@ -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
@@ -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)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user