CFG-05: modulübergreifender-http-endpunkt-für-benachrichtigungs-ereignisse

- internal/notifyapi.Mount (POST /notify/enqueue), reiner Wrapper um
  notifyprefs.EnqueueIfAllowed (CFG-04)
- Reuse von internal/policyapi.RequireServiceToken (RBAC-06), kein
  neues Provisorium, eigener Service-Token pro Endpunkt
- Tests: 401 ohne Token, Aequivalenz HTTP vs. direkter Aufruf
  (aktiviert/deaktiviert), realer Fremd-Modul-Client mit echter
  notification_jobs-Zeile
- real deployed auf 131 (notify-api, Port 8094), Grant-Nachverfolgung
  fuer nexarch_core auf notification_preferences/notification_jobs,
  end-zu-ende per curl nachgewiesen (401, job_id+skipped:false)

Pruefungen siehe docs/CFG-05-PRUEFPROTOKOLL.md
This commit is contained in:
sysops
2026-08-30 09:04:31 +02:00
parent a2f26a6e93
commit 23a3cc1a62
5 changed files with 393 additions and 0 deletions
+176
View File
@@ -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)
}
}