package webhook import ( "context" "encoding/json" "io" "net/http" "net/http/httptest" "os" "sync/atomic" "testing" "time" "github.com/jackc/pgx/v5/pgxpool" ) func setupTest(t *testing.T) (*Store, *pgxpool.Pool, func()) { t.Helper() adminDSN := os.Getenv("TEST_ADMIN_DSN") if adminDSN == "" { t.Skip("TEST_ADMIN_DSN nicht gesetzt, Integrationstest uebersprungen") } ctx := context.Background() pool, err := pgxpool.New(ctx, adminDSN) if err != nil { t.Fatalf("pool: %v", err) } if _, err := pool.Exec(ctx, ` CREATE TABLE IF NOT EXISTS webhook_subscriptions ( id UUID PRIMARY KEY DEFAULT gen_random_uuid(), event_type TEXT NOT NULL, target_url TEXT NOT NULL, secret TEXT NOT NULL, created_at TIMESTAMPTZ NOT NULL DEFAULT now() ); CREATE TABLE IF NOT EXISTS webhook_deliveries ( id UUID PRIMARY KEY DEFAULT gen_random_uuid(), subscription_id UUID NOT NULL REFERENCES webhook_subscriptions(id), event_type TEXT NOT NULL, payload JSONB NOT NULL, status TEXT NOT NULL DEFAULT 'pending', attempt INT NOT NULL DEFAULT 0, next_attempt_at TIMESTAMPTZ NOT NULL DEFAULT now(), last_error TEXT, created_at TIMESTAMPTZ NOT NULL DEFAULT now(), delivered_at TIMESTAMPTZ ); `); err != nil { t.Fatalf("schema: %v", err) } cleanup := func() { pool.Close() } return NewStore(pool), pool, cleanup } func uniqueEventType(prefix string) string { return prefix + "-" + time.Now().Format("150405.000000000") } // Akzeptanzkriterium 1 + Pruefung 1: ein Modul reicht ein Ereignis ein // (Enqueue) ohne eigene Zustellungslogik, der zentrale Dispatcher liefert // zuverlaessig aus. func TestEnqueueAndProcessDue_DeliversSuccessfully(t *testing.T) { store, pool, cleanup := setupTest(t) defer cleanup() ctx := context.Background() var receivedBody []byte var receivedSignature string target := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { body, _ := io.ReadAll(r.Body) receivedBody = body receivedSignature = r.Header.Get(SignatureHeader) w.WriteHeader(http.StatusOK) })) defer target.Close() eventType := uniqueEventType("dms.file.created") sub, err := store.Subscribe(ctx, eventType, target.URL, "geheimes-secret") if err != nil { t.Fatalf("subscribe: %v", err) } n, err := store.Enqueue(ctx, eventType, map[string]string{"file_id": "42"}) if err != nil { t.Fatalf("enqueue: %v", err) } if n != 1 { t.Fatalf("erwartet 1 eingereihte zustellung, habe %d", n) } dispatcher := NewDispatcher(pool, target.Client()) processed, err := dispatcher.ProcessDue(ctx) if err != nil { t.Fatalf("processdue: %v", err) } if processed != 1 { t.Fatalf("erwartet 1 verarbeitete zustellung, habe %d", processed) } var status string if err := pool.QueryRow(ctx, `SELECT status FROM webhook_deliveries WHERE subscription_id = $1`, sub.ID).Scan(&status); err != nil { t.Fatalf("status lesen: %v", err) } if status != StatusDelivered { t.Fatalf("status = %q, want %q", status, StatusDelivered) } // Postgres' JSONB-Spalte kann die Byte-Repraesentation des Payloads // gegenueber dem urspruenglichen json.Marshal kanonisieren (z.B. // Leerzeichen) — das ist unschaedlich, denn der Dispatcher signiert // IMMER exakt die Bytes, die er auch sendet. Die Pruefung vergleicht // deshalb Signatur gegen tatsaechlich empfangene Bytes (Selbstkonsistenz), // nicht gegen eine unabhaengig neu marshalte Referenz. if !VerifySignature("geheimes-secret", receivedBody, receivedSignature) { t.Fatalf("empfangene signatur %q passt nicht zum empfangenen payload %q", receivedSignature, receivedBody) } var decoded map[string]string if err := json.Unmarshal(receivedBody, &decoded); err != nil { t.Fatalf("empfangener payload nicht als json lesbar: %v", err) } if decoded["file_id"] != "42" { t.Fatalf("empfangener payload = %v, want file_id=42", decoded) } } // Akzeptanzkriterium 2 + Pruefung 2: fehlschlagendes Ziel loest Retry mit // wachsendem Backoff aus und endet nach der konfigurierten Obergrenze in // "failed". func TestProcessDue_RetriesWithBackoffThenMarksFailed(t *testing.T) { store, pool, cleanup := setupTest(t) defer cleanup() ctx := context.Background() var callCount int32 target := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { atomic.AddInt32(&callCount, 1) w.WriteHeader(http.StatusInternalServerError) })) defer target.Close() eventType := uniqueEventType("mail.send.failed") sub, err := store.Subscribe(ctx, eventType, target.URL, "secret") if err != nil { t.Fatalf("subscribe: %v", err) } if _, err := store.Enqueue(ctx, eventType, map[string]string{"x": "y"}); err != nil { t.Fatalf("enqueue: %v", err) } dispatcher := NewDispatcher(pool, target.Client()) dispatcher.MaxAttempts = 2 dispatcher.BaseBackoff = 30 * time.Millisecond // 1. Versuch: schlaegt fehl, ist aber noch nicht die Obergrenze. if _, err := dispatcher.ProcessDue(ctx); err != nil { t.Fatalf("processdue 1: %v", err) } var status string var attempt int if err := pool.QueryRow(ctx, `SELECT status, attempt FROM webhook_deliveries WHERE subscription_id = $1`, sub.ID).Scan(&status, &attempt); err != nil { t.Fatalf("status lesen 1: %v", err) } if status != StatusPending || attempt != 1 { t.Fatalf("nach 1. fehlschlag: status=%q attempt=%d, want pending/1", status, attempt) } // Sofort erneut verarbeiten: Backoff ist noch nicht abgelaufen -> nichts faellig. processedTooEarly, err := dispatcher.ProcessDue(ctx) if err != nil { t.Fatalf("processdue (zu frueh): %v", err) } if processedTooEarly != 0 { t.Fatal("erwartet 0 verarbeitete zustellungen, solange backoff nicht abgelaufen ist") } time.Sleep(dispatcher.backoffFor(1) + 20*time.Millisecond) // 2. Versuch: erreicht MaxAttempts=2 -> endgueltig fehlgeschlagen. if _, err := dispatcher.ProcessDue(ctx); err != nil { t.Fatalf("processdue 2: %v", err) } if err := pool.QueryRow(ctx, `SELECT status, attempt FROM webhook_deliveries WHERE subscription_id = $1`, sub.ID).Scan(&status, &attempt); err != nil { t.Fatalf("status lesen 2: %v", err) } if status != StatusFailed || attempt != 2 { t.Fatalf("nach 2. fehlschlag: status=%q attempt=%d, want failed/2", status, attempt) } if atomic.LoadInt32(&callCount) != 2 { t.Fatalf("erwartet genau 2 tatsaechliche zustellversuche, habe %d", callCount) } } // Akzeptanzkriterium 3 + Pruefung 3: Signaturpruefung erkennt eine // manipulierte Payload zuverlaessig. func TestVerifySignature_DetectsTamperedPayload(t *testing.T) { secret := "geteiltes-geheimnis" payload := []byte(`{"file_id":"42"}`) signature := Sign(secret, payload) if !VerifySignature(secret, payload, signature) { t.Fatal("erwartet gueltige signatur fuer unveraenderte payload") } tampered := []byte(`{"file_id":"99"}`) if VerifySignature(secret, tampered, signature) { t.Fatal("erwartet ungueltige signatur fuer manipulierte payload") } if VerifySignature("falsches-secret", payload, signature) { t.Fatal("erwartet ungueltige signatur bei falschem secret") } }