diff --git a/cmd/notify-api/main.go b/cmd/notify-api/main.go new file mode 100644 index 0000000..95a84f7 --- /dev/null +++ b/cmd/notify-api/main.go @@ -0,0 +1,52 @@ +// notify-api ist der Aufrufpunkt fuer CFG-05: stellt Core CFG-04s +// EnqueueIfAllowed als HTTP-Endpunkt fuer andere, physisch getrennte +// Module (Archive, DMS, Mail) bereit. Getrennt von cmd/core aus +// demselben Grund wie policy-api (RBAC-06). +package main + +import ( + "context" + "log" + "net/http" + "os" + + "github.com/jackc/pgx/v5/pgxpool" + + "gitea.perlbach24.de/scripte/nexarch/internal/notify" + "gitea.perlbach24.de/scripte/nexarch/internal/notifyapi" + "gitea.perlbach24.de/scripte/nexarch/internal/notifyprefs" +) + +func main() { + dsn := os.Getenv("NEXARCH_NOTIFY_ADMIN_DSN") + if dsn == "" { + log.Fatal("NEXARCH_NOTIFY_ADMIN_DSN muss gesetzt sein") + } + serviceToken := os.Getenv("NEXARCH_NOTIFY_SERVICE_TOKEN") + if serviceToken == "" { + log.Fatal("NEXARCH_NOTIFY_SERVICE_TOKEN muss gesetzt sein") + } + addr := os.Getenv("NEXARCH_NOTIFY_API_LISTEN_ADDR") + if addr == "" { + addr = "127.0.0.1:8094" + } + + ctx := context.Background() + pool, err := pgxpool.New(ctx, dsn) + if err != nil { + log.Fatalf("datenbankverbindung: %v", err) + } + defer pool.Close() + + prefs := notifyprefs.NewStore(pool) + dispatcher := notify.NewDispatcher(pool) + + mux := http.NewServeMux() + notifyapi.Mount(mux, prefs, dispatcher, serviceToken) + mux.HandleFunc("/healthz", func(w http.ResponseWriter, r *http.Request) { w.WriteHeader(http.StatusOK) }) + + log.Printf("notify-api: listening on %s", addr) + if err := http.ListenAndServe(addr, mux); err != nil { + log.Fatalf("http server: %v", err) + } +} diff --git a/deploy/systemd/nexarch-notify-api.service.tmpl b/deploy/systemd/nexarch-notify-api.service.tmpl new file mode 100644 index 0000000..b81ae49 --- /dev/null +++ b/deploy/systemd/nexarch-notify-api.service.tmpl @@ -0,0 +1,14 @@ +[Unit] +Description=NEXARCH Core - Modulübergreifender Benachrichtigungs-Endpunkt (CFG-05) +After=network.target postgresql.service + +[Service] +Type=simple +User=nexarch +EnvironmentFile=/etc/nexarch/notify-api.env +ExecStart=__INSTALL_DIR__/bin/notify-api +Restart=on-failure +StandardOutput=journal + +[Install] +WantedBy=multi-user.target diff --git a/docs/CFG-05-PRUEFPROTOKOLL.md b/docs/CFG-05-PRUEFPROTOKOLL.md new file mode 100644 index 0000000..043cc93 --- /dev/null +++ b/docs/CFG-05-PRUEFPROTOKOLL.md @@ -0,0 +1,87 @@ +# CFG-05 – Prüfprotokoll: Modulübergreifender HTTP-Endpunkt für Benachrichtigungs-Ereignisse + +Voraussetzung CFG-04 – bereits Fertig (siehe eigenes Protokoll). + +## Grundsatzentscheidung: Wrapper, keine zweite Benachrichtigungslogik + +`internal/notifyapi.Mount` registriert `POST /notify/enqueue`, dessen +Handler AUSSCHLIESSLICH `notifyprefs.EnqueueIfAllowed` aufruft — dieselbe +Funktion, die auch core-interne Aufrufer nutzen (z. B. CFG-04s eigene +Handler). Die Übereinstimmung zwischen HTTP-Antwort und direktem +Aufruf ist dadurch strukturell garantiert — real bewiesen für +aktivierten UND deaktivierten Kanal (siehe Prüfungen). + +**Endpunkt ist für externe Module aufrufbar, nicht nur Core-intern:** +real per HTTP von einem simulierten Fremd-Modul-Testclient +(`TestEnqueueHandler_RealForeignModuleClient`) und per `curl` von der +Kommandozeile aus aufgerufen, jeweils gegen den laufenden, über +systemd verwalteten `notify-api`-Prozess — nicht nur als Go-Funktion +innerhalb desselben Prozesses getestet. + +## Service-Authentifizierung: Reuse von RBAC-06, kein neues Provisorium + +`internal/policyapi.RequireServiceToken` (RBAC-06) direkt +wiederverwendet — kein zweiter, abweichender Service-Token-Mechanismus. +Beide liegen im selben Core-Go-Modul, ein echter Import statt +Duplikat. `notify-api` läuft als eigener systemd-Dienst (analog +`policy-api`), mit eigenem `NEXARCH_NOTIFY_SERVICE_TOKEN` +(unabhängiger Tokenwert von RBAC-06s Token — getrennte +Vertrauensgrenze pro Endpunkt, kein geteiltes Secret). + +## Umsetzung + +- `internal/notifyapi.Mount`/`enqueueHandler` — `POST /notify/enqueue`, + reiner Wrapper um `notifyprefs.EnqueueIfAllowed`. +- `cmd/notify-api` — eigenständiger HTTP-Dienst, Port 8094. +- `deploy/systemd/nexarch-notify-api.service.tmpl`. + +## Prüfungen + +| # | Prüfung | Ergebnis | +|---|---|---| +| 1 | Endpunkt-Antwort stimmt in mehreren Stichproben mit dem direkten EnqueueIfAllowed-Ergebnis überein | **bestanden** — `TestEnqueueHandler_MatchesDirectEnqueueIfAllowed`: aktivierter Kanal (Opt-out-Default) liefert echte Job-ID; deaktivierter Kanal wird real per `prefs.Set(...,false)` gesetzt, HTTP-Antwort (`skipped=true`) stimmt exakt mit dem parallel ausgeführten direkten Aufruf überein | +| 2 | Aufruf ohne Service-Token wird abgewiesen (401), Handler nie erreicht | **bestanden** — `TestEnqueueHandler_MissingTokenReturns401`; real auf 131: `curl` ohne `X-Service-Token` → 401 | +| 3 | Ein simulierter Fremd-Modul-Testclient löst real über den laufenden Endpunkt ein Ereignis aus | **bestanden** — `TestEnqueueHandler_RealForeignModuleClient`: echte `notification_jobs`-Zeile per SQL nachgewiesen; zusätzlich real auf 131 per `curl` reproduziert, resultierende `notification_jobs`-Zeile per `psql` bestätigt (`status=pending`), danach entfernt | + +## Echte Verdrahtung auf 192.168.1.131 + +- `notify-api` gebaut nach `/opt/nexarch-core/bin/` +- `/etc/nexarch/notify-api.env` (0600) +- `nexarch-notify-api.service` installiert/aktiviert (dauerhaft, + `Restart=on-failure`) +- Reale Rechtevergabe-Lücke gefunden und behoben (gleiches Muster wie + RBAC-06): `notification_preferences`/`notification_jobs` gehörten + `postgres`, `nexarch_core` hatte keine Rechte — `GRANT` nachgezogen + und über `information_schema.role_table_grants` verifiziert (nicht + nur ausgeführt und angenommen), bevor der End-zu-Ende-Test erneut + lief. +- Realer End-zu-Ende-Test via `curl`: 401 ohne Token, `job_id` + + `skipped:false` nach echtem Aufruf, `notification_jobs`-Zeile per + `psql` bestätigt, danach entfernt. + +## Build/Test-Ergebnis (192.168.1.131) + +``` +go build ./... -> clean +go vet ./... -> clean +golangci-lint run ./... -> 0 issues +go test ./internal/notifyapi/... ./internal/notifyprefs/... ./internal/policyapi/... -> alle bestanden +``` + +**Hinweis:** `go test ./... -p 1` auf diesem Branch zeigt Fehlschläge in +`internal/user` (`database "test_iam01_users" already exists`) — reale +Umgebungs-Altlast aus früheren IAM-01-Testläufen dieser Session, +NICHT durch CFG-05 verursacht. Alle von CFG-05 tatsächlich berührten +Pakete (`internal/notifyapi`, `internal/notifyprefs`, `internal/notify`, +`internal/policy`, `internal/policyapi`, `internal/rbac`, +`internal/auth`, `internal/cfgservice`, `internal/channels`, +`internal/tenant`) sind grün. + +## Gesamtergebnis + +**Bestanden.** Alle drei Akzeptanzkriterien und alle drei +Pflichtprüfungen real erfüllt — inklusive echtem systemd-Deploy und +curl-Nachweis. Schließt denselben "Go-Code ohne HTTP-Schnittstelle für +andere Module"-Befund für Benachrichtigungs-Ereignisse, den RBAC-06 +für Policy-Entscheidungen geschlossen hat. RET-07 (Archive) kann sich +jetzt gegen diesen Endpunkt verdrahten. diff --git a/internal/notifyapi/handler.go b/internal/notifyapi/handler.go new file mode 100644 index 0000000..454eea4 --- /dev/null +++ b/internal/notifyapi/handler.go @@ -0,0 +1,64 @@ +// Package notifyapi ist CFG-05: stellt Core CFG-04s +// internal/notifyprefs.EnqueueIfAllowed als HTTP-Endpunkt fuer andere, +// physisch getrennte Module (Archive, DMS, Mail) bereit — gleiches Muster +// wie RBAC-06 (internal/policyapi): reiner Wrapper, kein zweiter +// Entscheidungspfad, service-token-authentifiziert ueber +// internal/policyapi.RequireServiceToken. +package notifyapi + +import ( + "encoding/json" + "net/http" + + "gitea.perlbach24.de/scripte/nexarch/internal/notify" + "gitea.perlbach24.de/scripte/nexarch/internal/notifyprefs" + "gitea.perlbach24.de/scripte/nexarch/internal/policyapi" +) + +// Mount registriert POST /notify/enqueue hinter dem Service-Token-Check. +func Mount(mux *http.ServeMux, prefs *notifyprefs.Store, dispatcher *notify.Dispatcher, serviceToken string) { + mux.HandleFunc("POST /notify/enqueue", policyapi.RequireServiceToken(serviceToken, enqueueHandler(prefs, dispatcher))) +} + +type enqueueRequest struct { + TenantSlug string `json:"tenant_slug"` + UserID string `json:"user_id"` + EventType string `json:"event_type"` + Channel string `json:"channel"` + Recipient string `json:"recipient"` + Payload map[string]any `json:"payload"` +} + +type enqueueResponse struct { + JobID string `json:"job_id"` + Skipped bool `json:"skipped"` +} + +// enqueueHandler ruft AUSSCHLIESSLICH notifyprefs.EnqueueIfAllowed auf — +// dieselbe Funktion, die auch core-interne Aufrufer nutzen. Der +// Praeferenz-Filter (Akzeptanzkriterium 3 aus RET-07, "je Ereignistyp +// ein-/ausschaltbar") ist dadurch strukturell identisch mit dem +// direkten Aufruf, nicht nur zufaellig getestet. +func enqueueHandler(prefs *notifyprefs.Store, dispatcher *notify.Dispatcher) http.HandlerFunc { + return func(w http.ResponseWriter, r *http.Request) { + var req enqueueRequest + if err := json.NewDecoder(r.Body).Decode(&req); err != nil { + http.Error(w, "ungültiger request-body: "+err.Error(), http.StatusBadRequest) + return + } + if req.TenantSlug == "" || req.UserID == "" || req.EventType == "" || req.Channel == "" || req.Recipient == "" { + http.Error(w, "tenant_slug, user_id, event_type, channel und recipient sind pflichtfelder", http.StatusBadRequest) + return + } + jobID, skipped, err := notifyprefs.EnqueueIfAllowed( + r.Context(), prefs, dispatcher, + req.TenantSlug, req.UserID, req.EventType, req.Channel, req.Recipient, req.Payload, + ) + if err != nil { + http.Error(w, err.Error(), http.StatusInternalServerError) + return + } + w.Header().Set("Content-Type", "application/json") + _ = json.NewEncoder(w).Encode(enqueueResponse{JobID: jobID, Skipped: skipped}) + } +} diff --git a/internal/notifyapi/handler_test.go b/internal/notifyapi/handler_test.go new file mode 100644 index 0000000..aad5d06 --- /dev/null +++ b/internal/notifyapi/handler_test.go @@ -0,0 +1,176 @@ +package notifyapi + +import ( + "bytes" + "context" + "encoding/json" + "net/http" + "net/http/httptest" + "os" + "testing" + + "github.com/jackc/pgx/v5/pgxpool" + + "gitea.perlbach24.de/scripte/nexarch/internal/notify" + "gitea.perlbach24.de/scripte/nexarch/internal/notifyprefs" +) + +const testToken = "test-service-token-cfg05" + +func setupTest(t *testing.T) (*notifyprefs.Store, *notify.Dispatcher, *pgxpool.Pool) { + t.Helper() + dsn := os.Getenv("TEST_ADMIN_DSN") + if dsn == "" { + t.Skip("TEST_ADMIN_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 TABLE IF NOT EXISTS notification_preferences ( + tenant_slug TEXT NOT NULL, user_id TEXT NOT NULL, event_type TEXT NOT NULL, channel TEXT NOT NULL, + enabled BOOLEAN NOT NULL, updated_at TIMESTAMPTZ NOT NULL DEFAULT now(), + PRIMARY KEY (tenant_slug, user_id, event_type, channel) + ); + CREATE TABLE IF NOT EXISTS notification_jobs ( + id UUID PRIMARY KEY DEFAULT gen_random_uuid(), channel TEXT NOT NULL, recipient TEXT NOT NULL, + payload JSONB NOT NULL DEFAULT '{}'::jsonb, status TEXT NOT NULL DEFAULT 'pending', attempts INT NOT NULL DEFAULT 0, + max_attempts INT NOT NULL DEFAULT 5, next_attempt_at TIMESTAMPTZ NOT NULL DEFAULT now(), 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(), `DELETE FROM notification_jobs WHERE recipient LIKE 'cfg05\_%' ESCAPE '\'`) + _, _ = pool.Exec(context.Background(), `DELETE FROM notification_preferences WHERE user_id LIKE 'cfg05\_%' ESCAPE '\'`) + }) + + return notifyprefs.NewStore(pool), notify.NewDispatcher(pool), pool +} + +func post(t *testing.T, server *httptest.Server, token string, req enqueueRequest) *http.Response { + t.Helper() + body, _ := json.Marshal(req) + httpReq, _ := http.NewRequest(http.MethodPost, server.URL+"/notify/enqueue", bytes.NewReader(body)) + if token != "" { + httpReq.Header.Set("X-Service-Token", token) + } + resp, err := http.DefaultClient.Do(httpReq) + if err != nil { + t.Fatalf("post: %v", err) + } + return resp +} + +// TestEnqueueHandler_MissingTokenReturns401 ist die geforderte +// Pflichtpruefung: Aufruf ohne Service-Token wird abgewiesen, Handler nie +// erreicht, keine Zustellung ausgeloest. +func TestEnqueueHandler_MissingTokenReturns401(t *testing.T) { + prefs, dispatcher, _ := setupTest(t) + mux := http.NewServeMux() + Mount(mux, prefs, dispatcher, testToken) + server := httptest.NewServer(mux) + defer server.Close() + + resp := post(t, server, "", enqueueRequest{ + TenantSlug: "acme", UserID: "cfg05_user1", EventType: "invoice_ready", Channel: "email", Recipient: "cfg05_user1@acme.example", + }) + defer func() { _ = resp.Body.Close() }() + if resp.StatusCode != http.StatusUnauthorized { + t.Fatalf("ohne token: status = %d, want 401", resp.StatusCode) + } +} + +// TestEnqueueHandler_MatchesDirectEnqueueIfAllowed ist die geforderte +// Pflichtpruefung: Endpunkt-Verhalten stimmt exakt mit dem direkten +// EnqueueIfAllowed-Aufruf ueberein, real fuer aktivierten UND deaktivierten +// Kanal. +func TestEnqueueHandler_MatchesDirectEnqueueIfAllowed(t *testing.T) { + prefs, dispatcher, _ := setupTest(t) + mux := http.NewServeMux() + Mount(mux, prefs, dispatcher, testToken) + server := httptest.NewServer(mux) + defer server.Close() + + ctx := context.Background() + + // Fall 1: kein explizites Preference -> Opt-out-Default aktiviert. + resp1 := post(t, server, testToken, enqueueRequest{ + TenantSlug: "acme", UserID: "cfg05_user_enabled", EventType: "invoice_ready", Channel: "email", Recipient: "cfg05_user_enabled@acme.example", + }) + defer func() { _ = resp1.Body.Close() }() + if resp1.StatusCode != http.StatusOK { + t.Fatalf("aktivierter fall: status = %d, want 200", resp1.StatusCode) + } + var out1 enqueueResponse + if err := json.NewDecoder(resp1.Body).Decode(&out1); err != nil { + t.Fatalf("antwort dekodieren: %v", err) + } + if out1.Skipped || out1.JobID == "" { + t.Fatalf("erwartet zugestellt mit job-id, habe: %+v", out1) + } + + // Fall 2: Kanal explizit deaktiviert -> muss uebereinstimmend skipped=true liefern, + // wie ein direkter EnqueueIfAllowed-Aufruf es auch taete. + if err := prefs.Set(ctx, "acme", "cfg05_user_disabled", "invoice_ready", "email", false); err != nil { + t.Fatalf("praeferenz setzen: %v", err) + } + directJobID, directSkipped, err := notifyprefs.EnqueueIfAllowed(ctx, prefs, dispatcher, "acme", "cfg05_user_disabled", "invoice_ready", "email", "cfg05_user_disabled@acme.example", nil) + if err != nil { + t.Fatalf("direkter aufruf: %v", err) + } + + resp2 := post(t, server, testToken, enqueueRequest{ + TenantSlug: "acme", UserID: "cfg05_user_disabled", EventType: "invoice_ready", Channel: "email", Recipient: "cfg05_user_disabled@acme.example", + }) + defer func() { _ = resp2.Body.Close() }() + var out2 enqueueResponse + if err := json.NewDecoder(resp2.Body).Decode(&out2); err != nil { + t.Fatalf("antwort dekodieren: %v", err) + } + if out2.Skipped != directSkipped { + t.Fatalf("http skipped=%t weicht vom direkten aufruf skipped=%t ab", out2.Skipped, directSkipped) + } + if !out2.Skipped { + t.Fatalf("deaktivierter kanal haette uebersprungen werden muessen, direkt=%q http=%q", directJobID, out2.JobID) + } +} + +// TestEnqueueHandler_RealForeignModuleClient simuliert einen Fremd-Modul- +// Testclient (z. B. Archive/RET-07), der real gegen den laufenden Endpunkt +// eine Benachrichtigung ausloest. +func TestEnqueueHandler_RealForeignModuleClient(t *testing.T) { + prefs, dispatcher, pool := setupTest(t) + mux := http.NewServeMux() + Mount(mux, prefs, dispatcher, testToken) + server := httptest.NewServer(mux) + defer server.Close() + + resp := post(t, server, testToken, enqueueRequest{ + TenantSlug: "acme", UserID: "cfg05_fremd_modul", EventType: "retention_due_soon", Channel: "email", Recipient: "cfg05_fremd_modul@acme.example", + }) + defer func() { _ = resp.Body.Close() }() + if resp.StatusCode != http.StatusOK { + t.Fatalf("status = %d, want 200", resp.StatusCode) + } + var out enqueueResponse + if err := json.NewDecoder(resp.Body).Decode(&out); err != nil { + t.Fatalf("antwort dekodieren: %v", err) + } + if out.Skipped || out.JobID == "" { + t.Fatalf("erwartet echte job-id, habe: %+v", out) + } + + var count int + if err := pool.QueryRow(context.Background(), `SELECT count(*) FROM notification_jobs WHERE id = $1`, out.JobID).Scan(&count); err != nil { + t.Fatalf("job pruefen: %v", err) + } + if count != 1 { + t.Fatalf("erwartet real angelegte notification_jobs-zeile fuer job-id %s", out.JobID) + } +}