133 lines
4.3 KiB
Go
133 lines
4.3 KiB
Go
// Package webhook implementiert Core API-07: eine zentrale Webhook-Registry
|
|
// und Zustellungs-Engine fuer alle Fachmodule (DMS/Mail/Archive/Workflow/AI).
|
|
// Module reichen Ereignisse EINMAL zur Zustellung ein und implementieren
|
|
// selbst KEINE eigene Retry-/Signatur-Logik — das ist der zentrale Zweck
|
|
// dieser Kachel ("Bewusst vermeiden: jedes Modul baut seine eigene
|
|
// Webhook-Zustellungs-Engine"). Zustellung laeuft ueber dieselbe
|
|
// Postgres-Jobqueue-Konvention (SELECT ... FOR UPDATE SKIP LOCKED) wie
|
|
// internal/tenant.Lifecycle.ProcessDueDeletions (TEN-04) und
|
|
// internal/notify.Dispatcher (CFG-02) — kein Redis/AMQP.
|
|
package webhook
|
|
|
|
import (
|
|
"context"
|
|
"crypto/hmac"
|
|
"crypto/sha256"
|
|
"encoding/hex"
|
|
"encoding/json"
|
|
"fmt"
|
|
"time"
|
|
|
|
"github.com/jackc/pgx/v5"
|
|
"github.com/jackc/pgx/v5/pgxpool"
|
|
)
|
|
|
|
const (
|
|
StatusPending = "pending"
|
|
StatusDelivered = "delivered"
|
|
StatusFailed = "failed"
|
|
)
|
|
|
|
// Subscription ist EIN externer Abonnent fuer einen Ereignistyp
|
|
// (Akzeptanzkriterium 1: Module registrieren Ereignistypen, externe
|
|
// Abonnenten registrieren Ziel-URLs — dieses Paket modelliert die
|
|
// Abonnenten-Seite; welche Ereignistypen ein Modul anbietet, ist bewusst
|
|
// NICHT Teil dieser Kachel).
|
|
type Subscription struct {
|
|
ID string
|
|
EventType string
|
|
TargetURL string
|
|
Secret string
|
|
CreatedAt time.Time
|
|
}
|
|
|
|
// Store persistiert Abonnements und Zustellversuche.
|
|
type Store struct {
|
|
pool *pgxpool.Pool
|
|
}
|
|
|
|
func NewStore(pool *pgxpool.Pool) *Store {
|
|
return &Store{pool: pool}
|
|
}
|
|
|
|
// Subscribe registriert einen Abonnenten fuer einen Ereignistyp. secret wird
|
|
// spaeter zur HMAC-Signierung jeder Zustellung an diesen Abonnenten
|
|
// verwendet (Akzeptanzkriterium 3).
|
|
func (s *Store) Subscribe(ctx context.Context, eventType, targetURL, secret string) (Subscription, error) {
|
|
var sub Subscription
|
|
sub.EventType, sub.TargetURL, sub.Secret = eventType, targetURL, secret
|
|
row := s.pool.QueryRow(ctx, `
|
|
INSERT INTO webhook_subscriptions (event_type, target_url, secret)
|
|
VALUES ($1, $2, $3)
|
|
RETURNING id, created_at
|
|
`, eventType, targetURL, secret)
|
|
if err := row.Scan(&sub.ID, &sub.CreatedAt); err != nil {
|
|
return Subscription{}, fmt.Errorf("abonnement anlegen: %w", err)
|
|
}
|
|
return sub, nil
|
|
}
|
|
|
|
// Enqueue reicht EIN Ereignis zur Zustellung an ALLE Abonnenten des
|
|
// angegebenen Ereignistyps ein — dies ist die EINZIGE Schnittstelle, die
|
|
// ein Fachmodul braucht (Akzeptanzkriterium 1). Jeder Abonnent erhaelt
|
|
// einen eigenen, unabhaengigen Zustellversuch-Datensatz.
|
|
func (s *Store) Enqueue(ctx context.Context, eventType string, payload any) (int, error) {
|
|
payloadJSON, err := json.Marshal(payload)
|
|
if err != nil {
|
|
return 0, fmt.Errorf("payload serialisieren: %w", err)
|
|
}
|
|
|
|
rows, err := s.pool.Query(ctx, `
|
|
SELECT id FROM webhook_subscriptions WHERE event_type = $1
|
|
`, eventType)
|
|
if err != nil {
|
|
return 0, fmt.Errorf("abonnenten ermitteln: %w", err)
|
|
}
|
|
var subscriptionIDs []string
|
|
for rows.Next() {
|
|
var id string
|
|
if err := rows.Scan(&id); err != nil {
|
|
rows.Close()
|
|
return 0, fmt.Errorf("abonnent lesen: %w", err)
|
|
}
|
|
subscriptionIDs = append(subscriptionIDs, id)
|
|
}
|
|
rows.Close()
|
|
if err := rows.Err(); err != nil {
|
|
return 0, err
|
|
}
|
|
|
|
for _, subID := range subscriptionIDs {
|
|
if _, err := s.pool.Exec(ctx, `
|
|
INSERT INTO webhook_deliveries (subscription_id, event_type, payload, status, next_attempt_at)
|
|
VALUES ($1, $2, $3, $4, now())
|
|
`, subID, eventType, payloadJSON, StatusPending); err != nil {
|
|
return 0, fmt.Errorf("zustellung einreihen: %w", err)
|
|
}
|
|
}
|
|
return len(subscriptionIDs), nil
|
|
}
|
|
|
|
// delivery ist ein interner Datensatz fuer EINEN Zustellversuch, inklusive
|
|
// der zugehoerigen Abonnentendaten (per JOIN geladen).
|
|
type delivery struct {
|
|
ID string
|
|
TargetURL string
|
|
Secret string
|
|
Payload []byte
|
|
Attempt int
|
|
}
|
|
|
|
// Sign berechnet die HMAC-SHA256-Signatur des Payloads (Akzeptanzkriterium
|
|
// 3) — hex-kodiert, damit sie problemlos als HTTP-Header uebertragen werden
|
|
// kann.
|
|
func Sign(secret string, payload []byte) string {
|
|
mac := hmac.New(sha256.New, []byte(secret))
|
|
mac.Write(payload)
|
|
return hex.EncodeToString(mac.Sum(nil))
|
|
}
|
|
|
|
// SignatureHeader ist der HTTP-Header, unter dem die Signatur uebertragen
|
|
// wird — dokumentierter Vertrag fuer Empfaenger (Akzeptanzkriterium 3).
|
|
const SignatureHeader = "X-Nexarch-Signature-256"
|