Zum Inhalt springen

Robuste Webhook-Delivery-Systeme aufbauen

Entwirf eine Webhook-Delivery mit Retries und exponentiellem Backoff, Signaturprüfung, Idempotenz, Delivery-Logging und Endpoint-Health-Monitoring.

5 Min. Lesezeit
Webhook-Delivery-Pipeline mit Event-Queuing, Retry-Logik mit Backoff, Signaturverifizierung und Tracking des Zustellstatus

Webhooks wirken einfach — ein HTTP POST, wenn etwas passiert. In der Praxis ist zuverlässige Webhook-Zustellung ein Problem verteilter Systeme. Endpoints fallen aus. Netzwerke werden partitioniert. Empfänger verarbeiten langsam. Dein System muss all das abfangen und gleichzeitig garantieren, dass Ereignisse ihr Ziel erreichen — ohne Verlust und ohne übermäßige Duplikate.

Der Unterschied zwischen einem Spielzeug-Webhook-System und einem produktionsreifen liegt in der Retry-Logik, der Signaturverifizierung, der Idempotenz und dem Health-Management der Endpoints.

Queuing von Webhook-Events

Events sollten sofort in eine Queue gestellt und asynchron zugestellt werden. Die Operation, die den Webhook auslöst, sollte nicht auf die Zustellung blockieren.

tstypescript
// ❌ Synchronous webhook delivery — blocks the operation
async function createOrder(order: Order): Promise<void> {
  await database.insert(order);
  // If this fails or times out, the order creation hangs
  await fetch(webhookUrl, {
    method: "POST",
    body: JSON.stringify({ event: "order.created", data: order }),
  });
}
tstypescript
// ✅ Queue-based async delivery
interface WebhookEvent {
  id: string;
  type: string;
  payload: Record<string, unknown>;
  createdAt: Date;
  subscriptionId: string;
  endpoint: string;
  attempts: number;
  maxAttempts: number;
  nextAttemptAt: Date;
  status: "pending" | "delivered" | "failed" | "exhausted";
}
 
class WebhookQueue {
  private events: WebhookEvent[] = [];
 
  enqueue(
    type: string,
    payload: Record<string, unknown>,
    subscriptions: WebhookSubscription[]
  ): string[] {
    const eventIds: string[] = [];
 
    for (const sub of subscriptions) {
      if (!sub.events.includes(type)) continue;
 
      const event: WebhookEvent = {
        id: crypto.randomUUID(),
        type,
        payload,
        createdAt: new Date(),
        subscriptionId: sub.id,
        endpoint: sub.url,
        attempts: 0,
        maxAttempts: 8,
        nextAttemptAt: new Date(),
        status: "pending",
      };
 
      this.events.push(event);
      eventIds.push(event.id);
    }
 
    return eventIds;
  }
 
  getDeliverable(limit: number): WebhookEvent[] {
    const now = new Date();
    return this.events
      .filter(
        (e) =>
          e.status === "pending" &&
          e.nextAttemptAt <= now
      )
      .slice(0, limit);
  }
}
 
interface WebhookSubscription {
  id: string;
  url: string;
  events: string[];
  secret: string;
  active: boolean;
}
 
// The order creation is now non-blocking
async function createOrder(
  order: Order,
  webhookQueue: WebhookQueue,
  subscriptions: WebhookSubscription[]
): Promise<void> {
  await database.insert(order);
  webhookQueue.enqueue("order.created", { order }, subscriptions);
  // Returns immediately — delivery happens asynchronously
}

Retry-Strategie mit exponentiellem Backoff

Wenn die Zustellung fehlschlägt, wiederhole mit steigenden Verzögerungen. Das verhindert, dass ein sich erholender Endpoint überlastet wird, und stellt gleichzeitig die letztendliche Zustellung sicher.

tstypescript
class WebhookDeliveryWorker {
  private queue: WebhookQueue;
 
  constructor(queue: WebhookQueue) {
    this.queue = queue;
  }
 
  async processNextBatch(batchSize: number = 10): Promise<void> {
    const events = this.queue.getDeliverable(batchSize);
 
    for (const event of events) {
      await this.deliver(event);
    }
  }
 
  private async deliver(event: WebhookEvent): Promise<void> {
    const signature = this.sign(event);
 
    try {
      const controller = new AbortController();
      const timeout = setTimeout(
        () => controller.abort(),
        10_000
      );
 
      const response = await fetch(event.endpoint, {
        method: "POST",
        headers: {
          "Content-Type": "application/json",
          "X-Webhook-ID": event.id,
          "X-Webhook-Signature": signature,
          "X-Webhook-Timestamp": event.createdAt.toISOString(),
        },
        body: JSON.stringify({
          id: event.id,
          type: event.type,
          data: event.payload,
          created_at: event.createdAt.toISOString(),
        }),
        signal: controller.signal,
      });
 
      clearTimeout(timeout);
 
      if (response.ok) {
        event.status = "delivered";
        this.logDelivery(event, "success", response.status);
      } else if (response.status >= 500) {
        this.scheduleRetry(event);
      } else if (response.status >= 400) {
        // Client error — don't retry
        event.status = "failed";
        this.logDelivery(event, "client_error", response.status);
      }
    } catch (error) {
      this.scheduleRetry(event);
    }
  }
 
  private scheduleRetry(event: WebhookEvent): void {
    event.attempts++;
 
    if (event.attempts >= event.maxAttempts) {
      event.status = "exhausted";
      this.logDelivery(event, "exhausted", 0);
      return;
    }
 
    // Exponential backoff: 1m, 2m, 4m, 8m, 16m, 32m, 64m, 128m
    const delayMs = Math.min(
      60_000 * Math.pow(2, event.attempts),
      128 * 60_000
    );
 
    // Add jitter to prevent thundering herd
    const jitter = Math.random() * delayMs * 0.1;
 
    event.nextAttemptAt = new Date(
      Date.now() + delayMs + jitter
    );
    event.status = "pending";
  }
 
  private sign(event: WebhookEvent): string {
    // HMAC-SHA256 signature for verification
    const payload = JSON.stringify({
      id: event.id,
      type: event.type,
      data: event.payload,
      created_at: event.createdAt.toISOString(),
    });
 
    return `sha256=${computeHmac(payload, event.subscriptionId)}`;
  }
 
  private logDelivery(
    event: WebhookEvent,
    result: string,
    statusCode: number
  ): void {
    console.log(
      JSON.stringify({
        eventId: event.id,
        type: event.type,
        endpoint: event.endpoint,
        attempt: event.attempts,
        result,
        statusCode,
        timestamp: new Date().toISOString(),
      })
    );
  }
}
 
function computeHmac(payload: string, secret: string): string {
  // Placeholder for HMAC-SHA256
  return `hmac_${payload.length}_${secret.slice(0, 4)}`;
}

Signaturverifizierung auf Empfängerseite

Empfänger müssen verifizieren, dass Webhooks tatsächlich von deinem System stammen und unterwegs nicht manipuliert wurden.

tstypescript
import { createHmac, timingSafeEqual } from "crypto";
 
function verifyWebhookSignature(
  payload: string,
  signature: string,
  secret: string,
  toleranceSeconds: number = 300
): { valid: boolean; reason?: string } {
  // Check timestamp to prevent replay attacks
  const timestampHeader = extractTimestamp(signature);
  if (timestampHeader) {
    const eventTime = new Date(timestampHeader).getTime();
    const now = Date.now();
    const age = Math.abs(now - eventTime);
 
    if (age > toleranceSeconds * 1000) {
      return {
        valid: false,
        reason: `Event too old: ${Math.round(age / 1000)}s`,
      };
    }
  }
 
  // Compute expected signature
  const expectedSignature = createHmac("sha256", secret)
    .update(payload)
    .digest("hex");
 
  const expected = `sha256=${expectedSignature}`;
 
  // Timing-safe comparison prevents timing attacks
  const signatureBuffer = Buffer.from(signature);
  const expectedBuffer = Buffer.from(expected);
 
  if (signatureBuffer.length !== expectedBuffer.length) {
    return { valid: false, reason: "Signature length mismatch" };
  }
 
  const isValid = timingSafeEqual(signatureBuffer, expectedBuffer);
 
  return {
    valid: isValid,
    reason: isValid ? undefined : "Signature mismatch",
  };
}
 
function extractTimestamp(headers: string): string | null {
  // Extract from X-Webhook-Timestamp header
  return null; // Simplified
}
 
// Express middleware for webhook verification
function webhookVerificationMiddleware(secret: string) {
  return (req: Request, res: Response, next: NextFunction) => {
    const signature = req.headers["x-webhook-signature"] as string;
 
    if (!signature) {
      return res.status(401).json({ error: "Missing signature" });
    }
 
    const result = verifyWebhookSignature(
      JSON.stringify(req.body),
      signature,
      secret
    );
 
    if (!result.valid) {
      return res.status(403).json({
        error: "Invalid signature",
        reason: result.reason,
      });
    }
 
    next();
  };
}

Health-Monitoring der Endpoints

Verfolge den Gesundheitszustand der Endpoints, um keine Ressourcen für Zustellungen an dauerhaft fehlschlagende Endpoints zu verschwenden.

tstypescript
interface EndpointHealth {
  url: string;
  consecutiveFailures: number;
  lastSuccessAt: Date | null;
  lastFailureAt: Date | null;
  totalDeliveries: number;
  totalFailures: number;
  status: "healthy" | "degraded" | "disabled";
}
 
class EndpointMonitor {
  private health: Map<string, EndpointHealth> = new Map();
 
  recordSuccess(url: string): void {
    const h = this.getHealth(url);
    h.consecutiveFailures = 0;
    h.lastSuccessAt = new Date();
    h.totalDeliveries++;
    h.status = "healthy";
  }
 
  recordFailure(url: string): void {
    const h = this.getHealth(url);
    h.consecutiveFailures++;
    h.lastFailureAt = new Date();
    h.totalDeliveries++;
    h.totalFailures++;
 
    if (h.consecutiveFailures >= 10) {
      h.status = "disabled";
    } else if (h.consecutiveFailures >= 3) {
      h.status = "degraded";
    }
  }
 
  shouldDeliver(url: string): boolean {
    const h = this.health.get(url);
    if (!h) return true;
 
    if (h.status === "disabled") {
      // Check if enough time has passed to retry
      const cooldownMs = 30 * 60 * 1000; // 30 minutes
      if (
        h.lastFailureAt &&
        Date.now() - h.lastFailureAt.getTime() > cooldownMs
      ) {
        // Allow one probe delivery
        return true;
      }
      return false;
    }
 
    return true;
  }
 
  getHealthReport(): EndpointHealth[] {
    return [...this.health.values()];
  }
 
  private getHealth(url: string): EndpointHealth {
    let h = this.health.get(url);
    if (!h) {
      h = {
        url,
        consecutiveFailures: 0,
        lastSuccessAt: null,
        lastFailureAt: null,
        totalDeliveries: 0,
        totalFailures: 0,
        status: "healthy",
      };
      this.health.set(url, h);
    }
    return h;
  }
}

Wichtigste Erkenntnisse

Webhook-Zustellung ist ein asynchrones Problem verteilter Systeme — stelle Events sofort in eine Queue und liefere sie außerhalb des Request-Pfads zu, damit die auslösende Operation niemals auf die Verfügbarkeit eines Endpoints blockiert. Wiederhole mit exponentiellem Backoff und Jitter, um transiente Fehler abzufangen, ohne sich erholende Endpoints zu überlasten, mit einer sinnvollen Obergrenze für maximale Verzögerung und Versuchszahl. Signiere jedes Webhook-Payload mit HMAC-SHA256, damit Empfänger die Authentizität verifizieren können, und füge Zeitstempel hinzu, um Replay-Angriffe zu verhindern. Empfänger müssen Signaturen mit zeitkonsistentem Vergleich prüfen, um Timing-Seitenkanalangriffe zu verhindern. Überwache den Gesundheitszustand der Endpoints durch Zählen aufeinanderfolgender Fehler, deaktiviere die Zustellung an dauerhaft fehlschlagende Endpoints automatisch und prüfe periodisch auf Wiederherstellung. Protokolliere jeden Zustellversuch mit Event-ID, Versuchsnummer und Ergebnis, um eine vollständige Audit-Spur für das Debugging von Zustellproblemen zu liefern. Das Ziel ist ein System, in dem Events nie verloren gehen, Endpoints nie überlastet werden und sowohl Sender als auch Empfänger die Integrität jeder Zustellung verifizieren können.

Wilfredo Rujel

Wilfredo Rujel

Full-Stack-Softwareentwickler

Diesen Beitrag teilenX