178 lines
5.7 KiB
Go
178 lines
5.7 KiB
Go
package webhook
|
|
|
|
import (
|
|
"bytes"
|
|
"context"
|
|
"crypto/subtle"
|
|
"encoding/hex"
|
|
"fmt"
|
|
"net/http"
|
|
"time"
|
|
|
|
"github.com/jackc/pgx/v5"
|
|
)
|
|
|
|
// DefaultMaxAttempts ist die konfigurierbare Obergrenze, ab der eine
|
|
// Zustellung endgueltig als fehlgeschlagen gilt (Akzeptanzkriterium 2).
|
|
const DefaultMaxAttempts = 5
|
|
|
|
// DefaultBaseBackoff ist die Basisdauer fuer exponentielles Backoff:
|
|
// naechster Versuch nach BaseBackoff * 2^attempt (Akzeptanzkriterium 2).
|
|
const DefaultBaseBackoff = 2 * time.Second
|
|
|
|
// Dispatcher liefert faellige Zustellungen aus. Konfigurierbar in Tests
|
|
// (kleine BaseBackoff, kleine MaxAttempts), damit Retry/Backoff/Obergrenze
|
|
// ohne minutenlange Wartezeit real durchlaufen werden koennen.
|
|
type Dispatcher struct {
|
|
pool pgxIface
|
|
client *http.Client
|
|
MaxAttempts int
|
|
BaseBackoff time.Duration
|
|
}
|
|
|
|
// pgxIface ist die schmale Teilmenge von *pgxpool.Pool, die der Dispatcher
|
|
// braucht — als Interface, damit Tests keine echte Verbindung fuer reine
|
|
// Signatur-/Backoff-Logik brauchen (wird hier aber durchgehend mit echten
|
|
// Integrationstests gegen Postgres verwendet, siehe dispatcher_test.go).
|
|
type pgxIface interface {
|
|
Begin(ctx context.Context) (pgx.Tx, error)
|
|
}
|
|
|
|
func NewDispatcher(pool pgxIface, client *http.Client) *Dispatcher {
|
|
if client == nil {
|
|
client = &http.Client{Timeout: 5 * time.Second}
|
|
}
|
|
return &Dispatcher{pool: pool, client: client, MaxAttempts: DefaultMaxAttempts, BaseBackoff: DefaultBaseBackoff}
|
|
}
|
|
|
|
// backoffFor berechnet die Wartezeit vor dem naechsten Versuch: exponentiell
|
|
// wachsend mit der Anzahl bereits unternommener Versuche.
|
|
func (d *Dispatcher) backoffFor(attempt int) time.Duration {
|
|
return d.BaseBackoff * time.Duration(1<<uint(attempt))
|
|
}
|
|
|
|
// ProcessDue liefert ALLE derzeit faelligen Zustellungen aus — dasselbe
|
|
// SELECT ... FOR UPDATE SKIP LOCKED-Muster wie
|
|
// internal/tenant.Lifecycle.ProcessDueDeletions, damit mehrere Dispatcher-
|
|
// Instanzen dieselbe Zustellung nie doppelt bearbeiten.
|
|
func (d *Dispatcher) ProcessDue(ctx context.Context) (int, error) {
|
|
tx, err := d.pool.Begin(ctx)
|
|
if err != nil {
|
|
return 0, fmt.Errorf("transaktion starten: %w", err)
|
|
}
|
|
defer func() { _ = tx.Rollback(ctx) }()
|
|
|
|
rows, err := tx.Query(ctx, `
|
|
SELECT wd.id, wd.payload, wd.attempt, ws.target_url, ws.secret
|
|
FROM webhook_deliveries wd
|
|
JOIN webhook_subscriptions ws ON ws.id = wd.subscription_id
|
|
WHERE wd.status = $1 AND wd.next_attempt_at <= now()
|
|
FOR UPDATE OF wd SKIP LOCKED
|
|
`, StatusPending)
|
|
if err != nil {
|
|
return 0, fmt.Errorf("faellige zustellungen abfragen: %w", err)
|
|
}
|
|
|
|
var deliveries []delivery
|
|
for rows.Next() {
|
|
var del delivery
|
|
if err := rows.Scan(&del.ID, &del.Payload, &del.Attempt, &del.TargetURL, &del.Secret); err != nil {
|
|
rows.Close()
|
|
return 0, fmt.Errorf("zustellung lesen: %w", err)
|
|
}
|
|
deliveries = append(deliveries, del)
|
|
}
|
|
rows.Close()
|
|
if err := rows.Err(); err != nil {
|
|
return 0, err
|
|
}
|
|
|
|
for _, del := range deliveries {
|
|
d.attemptOne(ctx, tx, del)
|
|
}
|
|
|
|
if err := tx.Commit(ctx); err != nil {
|
|
return 0, fmt.Errorf("transaktion committen: %w", err)
|
|
}
|
|
return len(deliveries), nil
|
|
}
|
|
|
|
// attemptOne fuehrt GENAU EINEN Zustellversuch aus und aktualisiert den
|
|
// Zustellungsdatensatz entsprechend — Erfolg (Akzeptanzkriterium 1/Pruefung
|
|
// 1), erneuter Fehlversuch mit Backoff, oder endgueltiges Scheitern nach
|
|
// DefaultMaxAttempts (Akzeptanzkriterium 2/Pruefung 2). Ein Fehler bei
|
|
// GENAU EINER Zustellung darf die anderen in diesem Batch nicht verhindern.
|
|
func (d *Dispatcher) attemptOne(ctx context.Context, tx pgx.Tx, del delivery) {
|
|
signature := Sign(del.Secret, del.Payload)
|
|
|
|
req, err := http.NewRequestWithContext(ctx, http.MethodPost, del.TargetURL, bytes.NewReader(del.Payload))
|
|
deliveryErr := err
|
|
var statusCode int
|
|
if err == nil {
|
|
req.Header.Set("Content-Type", "application/json")
|
|
req.Header.Set(SignatureHeader, signature)
|
|
resp, err := d.client.Do(req)
|
|
if err != nil {
|
|
deliveryErr = err
|
|
} else {
|
|
statusCode = resp.StatusCode
|
|
resp.Body.Close()
|
|
if statusCode < 200 || statusCode >= 300 {
|
|
deliveryErr = fmt.Errorf("unerwarteter statuscode %d", statusCode)
|
|
}
|
|
}
|
|
}
|
|
|
|
if deliveryErr == nil {
|
|
_, _ = tx.Exec(ctx, `
|
|
UPDATE webhook_deliveries SET status = $2, delivered_at = now(), attempt = attempt + 1
|
|
WHERE id = $1
|
|
`, del.ID, StatusDelivered)
|
|
return
|
|
}
|
|
|
|
nextAttempt := del.Attempt + 1
|
|
if nextAttempt >= d.MaxAttempts {
|
|
_, _ = tx.Exec(ctx, `
|
|
UPDATE webhook_deliveries SET status = $2, attempt = $3, last_error = $4
|
|
WHERE id = $1
|
|
`, del.ID, StatusFailed, nextAttempt, deliveryErr.Error())
|
|
return
|
|
}
|
|
|
|
nextAttemptAt := time.Now().Add(d.backoffFor(nextAttempt))
|
|
_, _ = tx.Exec(ctx, `
|
|
UPDATE webhook_deliveries SET attempt = $2, next_attempt_at = $3, last_error = $4
|
|
WHERE id = $1
|
|
`, del.ID, nextAttempt, nextAttemptAt, deliveryErr.Error())
|
|
}
|
|
|
|
// Run ruft ProcessDue in festen Abstaenden auf, bis ctx beendet wird —
|
|
// dieselbe Konvention wie internal/tenant.Lifecycle.RunSweeper.
|
|
func (d *Dispatcher) Run(ctx context.Context, interval time.Duration) {
|
|
ticker := time.NewTicker(interval)
|
|
defer ticker.Stop()
|
|
for {
|
|
select {
|
|
case <-ctx.Done():
|
|
return
|
|
case <-ticker.C:
|
|
_, _ = d.ProcessDue(ctx)
|
|
}
|
|
}
|
|
}
|
|
|
|
// VerifySignature prueft empfaengerseitig, ob signature zu payload und
|
|
// secret passt — timing-safe (dasselbe Muster wie internal/audit.timingsafe),
|
|
// damit ein Empfaenger die Authentizitaet einer Zustellung pruefen kann
|
|
// (Akzeptanzkriterium 3).
|
|
func VerifySignature(secret string, payload []byte, signature string) bool {
|
|
expected := Sign(secret, payload)
|
|
expectedBytes, err1 := hex.DecodeString(expected)
|
|
gotBytes, err2 := hex.DecodeString(signature)
|
|
if err1 != nil || err2 != nil {
|
|
return false
|
|
}
|
|
return subtle.ConstantTimeCompare(expectedBytes, gotBytes) == 1
|
|
}
|