HardПрактика16 min

Паттерн идемпотентности

Idempotency Key, deduplication table, Redis SETNX, middleware для Symfony и Go. Идемпотентность в очередях, Saga и HTTP-методах

Что такое идемпотентность

Идемпотентность -- свойство операции, при котором многократное выполнение дает тот же результат, что и однократное. В математике: f(f(x)) = f(x). В распределенных системах это критически важная гарантия.

Зачем нужна идемпотентность

В реальных системах сбои неизбежны:

Клиент ──► Сервер: POST /api/payments
                    Сервер обработал, платёж списан
Клиент ◄── Сервер: 200 OK
                    ⚡ Сетевой разрыв! Ответ потерян
Клиент ──► Сервер: POST /api/payments (retry)
                    ❌ Без идемпотентности: двойное списание!
                    ✅ С идемпотентностью: вернул сохранённый ответ

Типичные причины дублирования запросов:

Причина Описание Частота
Network timeout Клиент не получил ответ Высокая
Client retry Браузер/мобильное приложение повторяет Высокая
Load balancer retry Nginx/HAProxy ретраит при 502/503 Средняя
Queue redelivery RabbitMQ/SQS повторно доставляет Высокая
User double-click Двойной клик на кнопке Очень высокая

HTTP-методы и идемпотентность

RFC 7231 определяет идемпотентность HTTP-методов:

Метод Идемпотентный Безопасный Пояснение
GET Да Да Чтение, никогда не меняет данные
HEAD Да Да Как GET, но без тела ответа
OPTIONS Да Да Метаинформация
PUT Да Нет Полная замена ресурса
DELETE Да Нет Удаление (повторное -- 404, но состояние то же)
POST Нет Нет Создание ресурса, побочные эффекты
PATCH Нет Нет Частичное обновление (зависит от реализации)

PUT идемпотентен по определению: PUT /users/123 {name: "Alice"} -- сколько бы раз ни вызвали, результат один. POST -- нет: каждый POST /orders создает новый заказ.

Почему PATCH не идемпотентен

PATCH /account/balance  {"operation": "add", "amount": 100}
# Первый вызов:  баланс 0   → 100
# Второй вызов:  баланс 100 → 200  ← результат отличается!

# Но если PATCH устанавливает значение напрямую:
PATCH /users/123  {"name": "Alice"}
# Первый вызов:  name = Alice
# Второй вызов:  name = Alice  ← идемпотентно!

PATCH может быть идемпотентным, если описывает конечное состояние, а не дельту.

Idempotency Key Pattern

Основная идея: клиент генерирует уникальный ключ для каждой логической операции и передает его в заголовке. Сервер сохраняет результат и при повторном запросе возвращает сохраненный ответ.

Как это работает

                    Idempotency Key Flow

Клиент                                   Сервер
  │                                         │
  │  POST /api/payments                     │
  │  Idempotency-Key: idk_abc123           │
  │ ───────────────────────────────────────►│
  │                                         │  1. Проверить ключ в БД
  │                                         │  2. Ключ не найден → обработать
  │                                         │  3. Сохранить: key → response
  │  HTTP 200 {payment_id: "pay_xyz"}      │
  │ ◄───────────────────────────────────────│
  │                                         │
  │  ⚡ Timeout! Retry...                   │
  │                                         │
  │  POST /api/payments                     │
  │  Idempotency-Key: idk_abc123           │
  │ ───────────────────────────────────────►│
  │                                         │  1. Проверить ключ в БД
  │                                         │  2. Ключ НАЙДЕН → вернуть
  │  HTTP 200 {payment_id: "pay_xyz"}      │     сохранённый ответ
  │ ◄───────────────────────────────────────│

Стандарт и реализации в индустрии

Компания Заголовок Формат ключа TTL
Stripe Idempotency-Key UUID v4 24 часа
PayPal PayPal-Request-Id UUID v4 Зависит от API
Amazon Pay x-amz-pay-idempotency-key UUID v4 24 часа
IETF Draft Idempotency-Key Произвольная строка Рекомендован

PostgreSQL schema для deduplication

-- Deduplication table for idempotency keys
CREATE TABLE idempotency_keys (
    id          UUID PRIMARY KEY DEFAULT gen_random_uuid(),
    key         VARCHAR(255) NOT NULL,
    user_id     UUID NOT NULL,

    -- Request fingerprint
    method      VARCHAR(10) NOT NULL,
    path        VARCHAR(512) NOT NULL,
    request_hash VARCHAR(64),           -- SHA-256 of request body

    -- Stored response
    status_code SMALLINT NOT NULL,
    response_headers JSONB DEFAULT '{}',
    response_body JSONB,

    -- Lifecycle
    locked_at   TIMESTAMPTZ,            -- Processing lock
    completed_at TIMESTAMPTZ,           -- When processing finished
    created_at  TIMESTAMPTZ NOT NULL DEFAULT NOW(),
    expires_at  TIMESTAMPTZ NOT NULL DEFAULT NOW() + INTERVAL '24 hours',

    -- Unique per user to prevent cross-user collision
    CONSTRAINT uq_idempotency_key_user UNIQUE (key, user_id)
);

-- Index for cleanup job
CREATE INDEX idx_idempotency_keys_expires ON idempotency_keys (expires_at)
    WHERE completed_at IS NOT NULL;

-- Index for lock detection (find stale locks)
CREATE INDEX idx_idempotency_keys_locked ON idempotency_keys (locked_at)
    WHERE locked_at IS NOT NULL AND completed_at IS NULL;

-- TTL cleanup: run periodically via cron or pg_cron
-- DELETE FROM idempotency_keys WHERE expires_at < NOW();

PHP: Symfony Event Subscriber для Idempotency-Key

<?php

declare(strict_types=1);

namespace App\EventSubscriber;

use Symfony\Component\EventDispatcher\EventSubscriberInterface;
use Symfony\Component\HttpFoundation\JsonResponse;
use Symfony\Component\HttpFoundation\Request;
use Symfony\Component\HttpFoundation\Response;
use Symfony\Component\HttpKernel\Event\RequestEvent;
use Symfony\Component\HttpKernel\Event\ResponseEvent;
use Symfony\Component\HttpKernel\KernelEvents;

/**
 * Handles Idempotency-Key header for POST/PATCH requests.
 *
 * Flow:
 * 1. onKernelRequest: check if key exists in DB, return cached response
 * 2. Controller processes request normally
 * 3. onKernelResponse: store response in DB for future deduplication
 */
final class IdempotencySubscriber implements EventSubscriberInterface
{
    private const HEADER = 'Idempotency-Key';
    private const MAX_KEY_LENGTH = 255;
    private const LOCK_TIMEOUT_SECONDS = 30;

    public function __construct(
        private readonly \PDO $db,
        private readonly \Psr\Log\LoggerInterface $logger,
    ) {}

    public static function getSubscribedEvents(): array
    {
        return [
            KernelEvents::REQUEST => ['onKernelRequest', 100],
            KernelEvents::RESPONSE => ['onKernelResponse', -100],
        ];
    }

    public function onKernelRequest(RequestEvent $event): void
    {
        $request = $event->getRequest();

        if (!$event->isMainRequest()) {
            return;
        }

        // Only apply to non-idempotent methods
        if (!in_array($request->getMethod(), ['POST', 'PATCH'], true)) {
            return;
        }

        $idempotencyKey = $request->headers->get(self::HEADER);
        if ($idempotencyKey === null) {
            return; // No key provided: process normally
        }

        // Validate key format
        if (strlen($idempotencyKey) > self::MAX_KEY_LENGTH) {
            $event->setResponse(new JsonResponse(
                ['error' => 'Idempotency key too long'],
                Response::HTTP_BAD_REQUEST,
            ));
            return;
        }

        $userId = $this->resolveUserId($request);
        $existing = $this->findExistingKey($idempotencyKey, $userId);

        if ($existing === null) {
            // New key: lock it for processing
            $this->lockKey($idempotencyKey, $userId, $request);
            $request->attributes->set('_idempotency_key', $idempotencyKey);
            return;
        }

        // Key exists but still processing (concurrent request)
        if ($existing['completed_at'] === null) {
            if ($this->isLockStale($existing['locked_at'])) {
                // Stale lock: allow retry by re-locking
                $this->relockKey($idempotencyKey, $userId);
                $request->attributes->set('_idempotency_key', $idempotencyKey);
                return;
            }

            $event->setResponse(new JsonResponse(
                ['error' => 'Request is already being processed'],
                Response::HTTP_CONFLICT,
            ));
            return;
        }

        // Verify request body matches (prevent key reuse for different payload)
        $currentHash = $this->hashRequestBody($request);
        if ($existing['request_hash'] !== null && $existing['request_hash'] !== $currentHash) {
            $event->setResponse(new JsonResponse(
                ['error' => 'Idempotency key reused with different request body'],
                Response::HTTP_UNPROCESSABLE_ENTITY,
            ));
            return;
        }

        // Return cached response
        $this->logger->info('Returning cached idempotent response', [
            'key' => $idempotencyKey,
        ]);

        $event->setResponse(new JsonResponse(
            json_decode($existing['response_body'], true),
            (int) $existing['status_code'],
            json_decode($existing['response_headers'], true),
        ));
    }

    public function onKernelResponse(ResponseEvent $event): void
    {
        $request = $event->getRequest();

        if (!$event->isMainRequest()) {
            return;
        }

        $idempotencyKey = $request->attributes->get('_idempotency_key');
        if ($idempotencyKey === null) {
            return;
        }

        $response = $event->getResponse();
        $userId = $this->resolveUserId($request);

        // Store the response for future deduplication
        $this->completeKey(
            $idempotencyKey,
            $userId,
            $response->getStatusCode(),
            $this->extractHeaders($response),
            $response->getContent(),
            $this->hashRequestBody($request),
        );
    }

    private function findExistingKey(string $key, string $userId): ?array
    {
        $stmt = $this->db->prepare(
            'SELECT status_code, response_headers, response_body,
                    request_hash, locked_at, completed_at
             FROM idempotency_keys
             WHERE key = :key AND user_id = :userId AND expires_at > NOW()'
        );
        $stmt->execute(['key' => $key, 'userId' => $userId]);

        return $stmt->fetch(\PDO::FETCH_ASSOC) ?: null;
    }

    private function lockKey(string $key, string $userId, Request $request): void
    {
        $this->db->prepare(
            'INSERT INTO idempotency_keys (key, user_id, method, path, locked_at)
             VALUES (:key, :userId, :method, :path, NOW())
             ON CONFLICT (key, user_id) DO NOTHING'
        )->execute([
            'key' => $key,
            'userId' => $userId,
            'method' => $request->getMethod(),
            'path' => $request->getPathInfo(),
        ]);
    }

    private function completeKey(
        string $key,
        string $userId,
        int $statusCode,
        string $headers,
        string $body,
        string $requestHash,
    ): void {
        $this->db->prepare(
            'UPDATE idempotency_keys
             SET status_code = :statusCode,
                 response_headers = :headers,
                 response_body = :body,
                 request_hash = :requestHash,
                 completed_at = NOW()
             WHERE key = :key AND user_id = :userId'
        )->execute([
            'statusCode' => $statusCode,
            'headers' => $headers,
            'body' => $body,
            'requestHash' => $requestHash,
            'key' => $key,
            'userId' => $userId,
        ]);
    }

    private function relockKey(string $key, string $userId): void
    {
        $this->db->prepare(
            'UPDATE idempotency_keys
             SET locked_at = NOW(), completed_at = NULL
             WHERE key = :key AND user_id = :userId'
        )->execute(['key' => $key, 'userId' => $userId]);
    }

    private function isLockStale(?string $lockedAt): bool
    {
        if ($lockedAt === null) {
            return true;
        }

        $locked = new \DateTimeImmutable($lockedAt);
        $threshold = new \DateTimeImmutable("-" . self::LOCK_TIMEOUT_SECONDS . " seconds");

        return $locked < $threshold;
    }

    private function resolveUserId(Request $request): string
    {
        // Extract user ID from JWT token or session
        return $request->attributes->get('_user_id', 'anonymous');
    }

    private function hashRequestBody(Request $request): string
    {
        return hash('sha256', $request->getContent());
    }

    private function extractHeaders(Response $response): string
    {
        $headers = [];
        foreach (['Content-Type', 'X-Request-Id'] as $name) {
            $value = $response->headers->get($name);
            if ($value !== null) {
                $headers[$name] = $value;
            }
        }

        return json_encode($headers);
    }
}

Go: HTTP Middleware для Idempotency-Key

package middleware

import (
	"context"
	"crypto/sha256"
	"database/sql"
	"encoding/hex"
	"encoding/json"
	"io"
	"log/slog"
	"net/http"
	"strings"
	"time"
)

// IdempotencyRecord stores the cached response for a given key.
type IdempotencyRecord struct {
	StatusCode      int               `json:"status_code"`
	ResponseHeaders map[string]string `json:"response_headers"`
	ResponseBody    json.RawMessage   `json:"response_body"`
	RequestHash     string            `json:"request_hash"`
	CompletedAt     *time.Time        `json:"completed_at"`
	LockedAt        *time.Time        `json:"locked_at"`
}

// IdempotencyStore abstracts storage for idempotency keys.
type IdempotencyStore interface {
	Find(ctx context.Context, key, userID string) (*IdempotencyRecord, error)
	Lock(ctx context.Context, key, userID, method, path string) error
	Complete(ctx context.Context, key, userID string, rec IdempotencyRecord) error
}

const (
	idempotencyHeader  = "Idempotency-Key"
	maxKeyLength       = 255
	lockTimeoutSeconds = 30
)

// Idempotency returns middleware that handles Idempotency-Key header.
func Idempotency(store IdempotencyStore, logger *slog.Logger) func(http.Handler) http.Handler {
	return func(next http.Handler) http.Handler {
		return http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
			// Only apply to non-idempotent methods
			if r.Method != http.MethodPost && r.Method != http.MethodPatch {
				next.ServeHTTP(w, r)
				return
			}

			key := r.Header.Get(idempotencyHeader)
			if key == "" {
				next.ServeHTTP(w, r)
				return
			}

			if len(key) > maxKeyLength {
				writeJSON(w, http.StatusBadRequest, map[string]string{
					"error": "Idempotency key too long",
				})
				return
			}

			userID := resolveUserID(r)

			// Read body for hashing (buffer it for re-read)
			body, err := io.ReadAll(r.Body)
			if err != nil {
				writeJSON(w, http.StatusBadRequest, map[string]string{
					"error": "Failed to read request body",
				})
				return
			}
			r.Body = io.NopCloser(strings.NewReader(string(body)))
			bodyHash := hashBody(body)

			existing, err := store.Find(r.Context(), key, userID)
			if err != nil {
				logger.Error("idempotency store lookup failed", "error", err)
				next.ServeHTTP(w, r)
				return
			}

			if existing != nil {
				// Still processing by another request
				if existing.CompletedAt == nil {
					if !isLockStale(existing.LockedAt) {
						writeJSON(w, http.StatusConflict, map[string]string{
							"error": "Request is already being processed",
						})
						return
					}
				} else {
					// Verify request body matches
					if existing.RequestHash != "" && existing.RequestHash != bodyHash {
						writeJSON(w, http.StatusUnprocessableEntity, map[string]string{
							"error": "Idempotency key reused with different request body",
						})
						return
					}

					// Return cached response
					logger.Info("returning cached idempotent response", "key", key)
					for k, v := range existing.ResponseHeaders {
						w.Header().Set(k, v)
					}
					w.WriteHeader(existing.StatusCode)
					w.Write(existing.ResponseBody) //nolint:errcheck
					return
				}
			}

			// Lock the key for processing
			if err := store.Lock(r.Context(), key, userID, r.Method, r.URL.Path); err != nil {
				logger.Error("failed to lock idempotency key", "error", err)
				next.ServeHTTP(w, r)
				return
			}

			// Wrap ResponseWriter to capture the response
			recorder := &responseRecorder{
				ResponseWriter: w,
				statusCode:     http.StatusOK,
			}

			next.ServeHTTP(recorder, r)

			// Store the response
			rec := IdempotencyRecord{
				StatusCode:      recorder.statusCode,
				ResponseHeaders: extractHeaders(recorder.Header()),
				ResponseBody:    recorder.body,
				RequestHash:     bodyHash,
			}

			if err := store.Complete(r.Context(), key, userID, rec); err != nil {
				logger.Error("failed to store idempotent response", "error", err)
			}
		})
	}
}

// responseRecorder captures the response for storage.
type responseRecorder struct {
	http.ResponseWriter
	statusCode int
	body       []byte
}

func (r *responseRecorder) WriteHeader(code int) {
	r.statusCode = code
	r.ResponseWriter.WriteHeader(code)
}

func (r *responseRecorder) Write(b []byte) (int, error) {
	r.body = append(r.body, b...)
	return r.ResponseWriter.Write(b)
}

func hashBody(body []byte) string {
	h := sha256.Sum256(body)
	return hex.EncodeToString(h[:])
}

func isLockStale(lockedAt *time.Time) bool {
	if lockedAt == nil {
		return true
	}
	return time.Since(*lockedAt) > lockTimeoutSeconds*time.Second
}

func resolveUserID(r *http.Request) string {
	if uid := r.Context().Value("user_id"); uid != nil {
		return uid.(string)
	}
	return "anonymous"
}

func extractHeaders(h http.Header) map[string]string {
	result := make(map[string]string)
	for _, name := range []string{"Content-Type", "X-Request-Id"} {
		if v := h.Get(name); v != "" {
			result[name] = v
		}
	}
	return result
}

func writeJSON(w http.ResponseWriter, status int, data any) {
	w.Header().Set("Content-Type", "application/json")
	w.WriteHeader(status)
	json.NewEncoder(w).Encode(data) //nolint:errcheck
}

Подключение Go middleware

package main

import (
	"net/http"

	"github.com/go-chi/chi/v5"
)

func main() {
	r := chi.NewRouter()

	store := NewPostgresIdempotencyStore(db)
	logger := slog.Default()

	// Apply idempotency middleware to API routes
	r.Route("/api/v1", func(r chi.Router) {
		r.Use(middleware.Idempotency(store, logger))

		r.Post("/payments", handleCreatePayment)
		r.Post("/orders", handleCreateOrder)
	})

	http.ListenAndServe(":8080", r)
}

Redis-based Deduplication (SETNX + TTL)

Для высоконагруженных систем, где обращение к PostgreSQL при каждом запросе слишком дорого, используют Redis.

<?php

declare(strict_types=1);

/**
 * Redis-based idempotency store using SETNX + TTL.
 *
 * Faster than PostgreSQL but less durable:
 * - Redis may lose data on restart (unless AOF is configured)
 * - Suitable for non-critical deduplication
 * - Use PostgreSQL for financial transactions
 */
final class RedisIdempotencyStore
{
    private const PREFIX = 'idem:';
    private const LOCK_SUFFIX = ':lock';
    private const DEFAULT_TTL = 86400; // 24 hours

    public function __construct(
        private readonly \Redis $redis,
        private readonly int $ttl = self::DEFAULT_TTL,
        private readonly int $lockTtl = 30,
    ) {}

    /**
     * Try to acquire processing lock for the key.
     * Returns true if this is a new request.
     */
    public function tryLock(string $key, string $userId): bool
    {
        $lockKey = self::PREFIX . $userId . ':' . $key . self::LOCK_SUFFIX;

        // SETNX: only sets if key doesn't exist (atomic)
        $acquired = $this->redis->set($lockKey, time(), [
            'NX' => true,          // Only set if Not eXists
            'EX' => $this->lockTtl, // Auto-expire lock
        ]);

        return (bool) $acquired;
    }

    /**
     * Store the response for future deduplication.
     */
    public function store(string $key, string $userId, array $response): void
    {
        $dataKey = self::PREFIX . $userId . ':' . $key;
        $lockKey = $dataKey . self::LOCK_SUFFIX;

        // Atomic: store response and remove lock in pipeline
        $this->redis->multi();
        $this->redis->setex($dataKey, $this->ttl, json_encode($response));
        $this->redis->del($lockKey);
        $this->redis->exec();
    }

    /**
     * Get cached response if exists.
     */
    public function find(string $key, string $userId): ?array
    {
        $dataKey = self::PREFIX . $userId . ':' . $key;
        $data = $this->redis->get($dataKey);

        if ($data === false) {
            return null;
        }

        return json_decode($data, true);
    }

    /**
     * Check if key is currently being processed.
     */
    public function isLocked(string $key, string $userId): bool
    {
        $lockKey = self::PREFIX . $userId . ':' . $key . self::LOCK_SUFFIX;

        return (bool) $this->redis->exists($lockKey);
    }
}

Сравнение: PostgreSQL vs Redis для дедупликации

Критерий PostgreSQL Redis
Надежность Высокая (ACID) Средняя (может потерять при рестарте)
Скорость ~2-5ms ~0.1-0.5ms
Подходит для Платежи, финансы Некритичные операции
Масштабирование Вертикальное Горизонтальное (Cluster)
TTL cleanup Нужен cron Встроенный (EXPIRE)
Concurrent locks FOR UPDATE / Advisory SETNX (атомарный)

Go: Redis Idempotency Store

package idempotency

import (
	"context"
	"encoding/json"
	"fmt"
	"time"

	"github.com/redis/go-redis/v9"
)

// RedisStore implements IdempotencyStore using Redis SETNX.
type RedisStore struct {
	client  *redis.Client
	ttl     time.Duration
	lockTTL time.Duration
}

func NewRedisStore(client *redis.Client) *RedisStore {
	return &RedisStore{
		client:  client,
		ttl:     24 * time.Hour,
		lockTTL: 30 * time.Second,
	}
}

func (s *RedisStore) Find(ctx context.Context, key, userID string) (*IdempotencyRecord, error) {
	dataKey := s.dataKey(key, userID)

	data, err := s.client.Get(ctx, dataKey).Bytes()
	if err == redis.Nil {
		// Check if locked (being processed)
		lockKey := s.lockKey(key, userID)
		exists, _ := s.client.Exists(ctx, lockKey).Result()
		if exists > 0 {
			now := time.Now()
			return &IdempotencyRecord{LockedAt: &now}, nil
		}
		return nil, nil
	}
	if err != nil {
		return nil, fmt.Errorf("redis get: %w", err)
	}

	var rec IdempotencyRecord
	if err := json.Unmarshal(data, &rec); err != nil {
		return nil, fmt.Errorf("unmarshal record: %w", err)
	}
	now := time.Now()
	rec.CompletedAt = &now

	return &rec, nil
}

func (s *RedisStore) Lock(ctx context.Context, key, userID, method, path string) error {
	lockKey := s.lockKey(key, userID)

	ok, err := s.client.SetNX(ctx, lockKey, time.Now().Unix(), s.lockTTL).Result()
	if err != nil {
		return fmt.Errorf("redis setnx: %w", err)
	}
	if !ok {
		return fmt.Errorf("key already locked: %s", key)
	}

	return nil
}

func (s *RedisStore) Complete(ctx context.Context, key, userID string, rec IdempotencyRecord) error {
	dataKey := s.dataKey(key, userID)
	lockKey := s.lockKey(key, userID)

	data, err := json.Marshal(rec)
	if err != nil {
		return fmt.Errorf("marshal record: %w", err)
	}

	// Pipeline: store data + remove lock atomically
	pipe := s.client.Pipeline()
	pipe.Set(ctx, dataKey, data, s.ttl)
	pipe.Del(ctx, lockKey)
	_, err = pipe.Exec(ctx)

	return err
}

func (s *RedisStore) dataKey(key, userID string) string {
	return fmt.Sprintf("idem:%s:%s", userID, key)
}

func (s *RedisStore) lockKey(key, userID string) string {
	return fmt.Sprintf("idem:%s:%s:lock", userID, key)
}

Retry Safety: клиентская сторона

Идемпотентность на сервере бесполезна без правильных ретраев на клиенте.

Exponential Backoff с Retry-After

package httpclient

import (
	"context"
	"fmt"
	"math"
	"math/rand/v2"
	"net/http"
	"strconv"
	"time"
)

// RetryConfig configures retry behavior.
type RetryConfig struct {
	MaxRetries     int
	BaseDelay      time.Duration
	MaxDelay       time.Duration
	RetryableCodes map[int]bool
}

func DefaultRetryConfig() RetryConfig {
	return RetryConfig{
		MaxRetries: 3,
		BaseDelay:  200 * time.Millisecond,
		MaxDelay:   30 * time.Second,
		RetryableCodes: map[int]bool{
			http.StatusRequestTimeout:      true, // 408
			http.StatusTooManyRequests:      true, // 429
			http.StatusInternalServerError:  true, // 500
			http.StatusBadGateway:           true, // 502
			http.StatusServiceUnavailable:   true, // 503
			http.StatusGatewayTimeout:       true, // 504
		},
	}
}

// DoWithRetry executes an HTTP request with retries and exponential backoff.
// Respects Retry-After header from the server.
func DoWithRetry(
	ctx context.Context,
	client *http.Client,
	req *http.Request,
	cfg RetryConfig,
) (*http.Response, error) {
	var lastErr error

	for attempt := 0; attempt <= cfg.MaxRetries; attempt++ {
		if attempt > 0 {
			delay := calculateBackoff(attempt-1, cfg.BaseDelay, cfg.MaxDelay)

			select {
			case <-ctx.Done():
				return nil, ctx.Err()
			case <-time.After(delay):
			}
		}

		resp, err := client.Do(req)
		if err != nil {
			lastErr = err
			continue // Network error: retry
		}

		// Success or non-retryable status
		if !cfg.RetryableCodes[resp.StatusCode] {
			return resp, nil
		}

		// Check Retry-After header
		if retryAfter := resp.Header.Get("Retry-After"); retryAfter != "" {
			if seconds, err := strconv.Atoi(retryAfter); err == nil {
				delay := time.Duration(seconds) * time.Second
				select {
				case <-ctx.Done():
					return nil, ctx.Err()
				case <-time.After(delay):
				}
			}
		}

		lastErr = fmt.Errorf("HTTP %d from %s", resp.StatusCode, req.URL)
		resp.Body.Close()
	}

	return nil, fmt.Errorf("max retries exceeded: %w", lastErr)
}

// calculateBackoff returns delay with exponential backoff + full jitter.
func calculateBackoff(attempt int, base, max time.Duration) time.Duration {
	exp := math.Pow(2, float64(attempt))
	delay := time.Duration(float64(base) * exp)

	if delay > max {
		delay = max
	}

	// Full jitter: random between 0 and calculated delay
	jitter := time.Duration(rand.Int64N(int64(delay)))
	return jitter
}

Идемпотентность в очередях сообщений

Очереди гарантируют at-least-once delivery -- сообщение может прийти дважды. Consumer обязан быть идемпотентным.

RabbitMQ: Message Deduplication

<?php

declare(strict_types=1);

/**
 * Idempotent message consumer for RabbitMQ.
 *
 * RabbitMQ guarantees at-least-once delivery:
 * - Consumer crashes after processing but before ack → redelivery
 * - Network partition → redelivery
 * - Consumer must handle duplicates safely
 */
final class IdempotentConsumer
{
    public function __construct(
        private readonly \PDO $db,
        private readonly \Psr\Log\LoggerInterface $logger,
    ) {}

    /**
     * Process a message idempotently using message_id as dedup key.
     *
     * @param callable(array): void $handler
     */
    public function handle(string $messageId, array $payload, callable $handler): void
    {
        // Check if already processed
        if ($this->isProcessed($messageId)) {
            $this->logger->info('Skipping duplicate message', [
                'message_id' => $messageId,
            ]);
            return; // Already processed: skip
        }

        // Process within transaction
        $this->db->beginTransaction();
        try {
            // Mark as processing (INSERT with unique constraint)
            $this->markAsProcessing($messageId);

            // Execute business logic
            $handler($payload);

            // Mark as completed
            $this->markAsCompleted($messageId);

            $this->db->commit();
        } catch (\Throwable $e) {
            $this->db->rollBack();

            // Remove the processing mark so message can be retried
            $this->removeProcessingMark($messageId);

            throw $e;
        }
    }

    private function isProcessed(string $messageId): bool
    {
        $stmt = $this->db->prepare(
            'SELECT 1 FROM processed_messages
             WHERE message_id = :id AND status = :status'
        );
        $stmt->execute(['id' => $messageId, 'status' => 'completed']);

        return $stmt->fetchColumn() !== false;
    }

    private function markAsProcessing(string $messageId): void
    {
        $this->db->prepare(
            'INSERT INTO processed_messages (message_id, status, created_at)
             VALUES (:id, :status, NOW())
             ON CONFLICT (message_id) DO NOTHING'
        )->execute(['id' => $messageId, 'status' => 'processing']);
    }

    private function markAsCompleted(string $messageId): void
    {
        $this->db->prepare(
            'UPDATE processed_messages
             SET status = :status, completed_at = NOW()
             WHERE message_id = :id'
        )->execute(['id' => $messageId, 'status' => 'completed']);
    }

    private function removeProcessingMark(string $messageId): void
    {
        $this->db->prepare(
            'DELETE FROM processed_messages
             WHERE message_id = :id AND status = :status'
        )->execute(['id' => $messageId, 'status' => 'processing']);
    }
}

// Schema for processed messages:
// CREATE TABLE processed_messages (
//     message_id VARCHAR(255) PRIMARY KEY,
//     status     VARCHAR(20) NOT NULL DEFAULT 'processing',
//     created_at TIMESTAMPTZ NOT NULL DEFAULT NOW(),
//     completed_at TIMESTAMPTZ
// );

SQS: встроенная дедупликация

Amazon SQS FIFO очереди имеют встроенную дедупликацию с 5-минутным окном:

Producer → SQS FIFO (MessageDeduplicationId: "order-123-payment")
                     ↓ Окно 5 минут
Повторное сообщение с тем же DeduplicationId → автоматически отклонено

Идемпотентность в Saga Pattern

В Saga каждый шаг должен быть идемпотентным, потому что compensating transaction может выполняться повторно.

<?php

declare(strict_types=1);

/**
 * Saga step with idempotent execution and compensation.
 *
 * Each step tracks its execution state.
 * Re-executing a completed step returns cached result.
 * Re-compensating a compensated step is a no-op.
 */
final class IdempotentSagaStep
{
    public function __construct(
        private readonly string $sagaId,
        private readonly string $stepName,
        private readonly \PDO $db,
    ) {}

    /**
     * Execute step idempotently.
     * Returns cached result if step was already completed.
     */
    public function execute(callable $action): mixed
    {
        $state = $this->getState();

        // Already completed: return cached result
        if ($state === 'completed') {
            return $this->getCachedResult();
        }

        // Already compensated: cannot re-execute
        if ($state === 'compensated') {
            throw new \LogicException("Cannot execute compensated step: {$this->stepName}");
        }

        $this->db->beginTransaction();
        try {
            $this->setState('executing');
            $result = $action();
            $this->setState('completed', $result);
            $this->db->commit();

            return $result;
        } catch (\Throwable $e) {
            $this->db->rollBack();
            $this->setState('failed');
            throw $e;
        }
    }

    /**
     * Compensate step idempotently.
     * No-op if already compensated.
     */
    public function compensate(callable $compensation): void
    {
        $state = $this->getState();

        // Not executed or already compensated: nothing to do
        if ($state === null || $state === 'compensated') {
            return;
        }

        $this->db->beginTransaction();
        try {
            $compensation();
            $this->setState('compensated');
            $this->db->commit();
        } catch (\Throwable $e) {
            $this->db->rollBack();
            throw $e;
        }
    }

    private function getState(): ?string
    {
        $stmt = $this->db->prepare(
            'SELECT state FROM saga_steps
             WHERE saga_id = :sagaId AND step_name = :stepName'
        );
        $stmt->execute([
            'sagaId' => $this->sagaId,
            'stepName' => $this->stepName,
        ]);

        return $stmt->fetchColumn() ?: null;
    }

    private function setState(string $state, mixed $result = null): void
    {
        $this->db->prepare(
            'INSERT INTO saga_steps (saga_id, step_name, state, result, updated_at)
             VALUES (:sagaId, :stepName, :state, :result, NOW())
             ON CONFLICT (saga_id, step_name)
             DO UPDATE SET state = :state2, result = :result2, updated_at = NOW()'
        )->execute([
            'sagaId' => $this->sagaId,
            'stepName' => $this->stepName,
            'state' => $state,
            'result' => $result !== null ? json_encode($result) : null,
            'state2' => $state,
            'result2' => $result !== null ? json_encode($result) : null,
        ]);
    }

    private function getCachedResult(): mixed
    {
        $stmt = $this->db->prepare(
            'SELECT result FROM saga_steps
             WHERE saga_id = :sagaId AND step_name = :stepName'
        );
        $stmt->execute([
            'sagaId' => $this->sagaId,
            'stepName' => $this->stepName,
        ]);

        $result = $stmt->fetchColumn();
        return $result !== false ? json_decode($result, true) : null;
    }
}

Anti-patterns: когда идемпотентность НЕ нужна

Не все операции требуют идемпотентности. Добавление сложности без необходимости -- это тоже ошибка.

Сценарий Нужна идемпотентность Почему
Платежи, переводы Да Дублирование = финансовые потери
Создание заказов Да Двойной заказ = проблемы
Отправка email Да Двойное письмо = спам
Чтение данных (GET) Нет Уже идемпотентно по определению
Логирование/аналитика Нет Дублирование безвредно (обычно)
Обновление счетчика лайков Зависит Повторный лайк уже ограничен бизнес-логикой
Генерация отчетов Нет Повторная генерация = тот же результат
DELETE с фиксированным ID Нет HTTP DELETE уже идемпотентен

Ошибки при реализации

Anti-pattern 1: Idempotency key без TTL
─────────────────────────────────────────
❌ Ключи накапливаются бесконечно → БД растёт
✅ Всегда устанавливайте TTL (обычно 24 часа)

Anti-pattern 2: Один ключ на все запросы
─────────────────────────────────────────
❌ Клиент переиспользует ключ для разных операций
✅ Генерировать уникальный ключ для каждой операции

Anti-pattern 3: Idempotency без проверки тела запроса
─────────────────────────────────────────
❌ Тот же ключ, другое тело → неправильный ответ
✅ Хеш тела запроса как дополнительная проверка

Anti-pattern 4: Блокировка без timeout
─────────────────────────────────────────
❌ Процесс упал → ключ залочен навсегда
✅ TTL на lock + проверка stale locks

Выводы

Идемпотентность -- фундаментальное свойство надежных распределенных систем. Для POST/PATCH запросов используйте Idempotency-Key паттерн с deduplication table. Для очередей -- отслеживайте message_id обработанных сообщений. Для Saga -- каждый шаг должен быть безопасен при повторном выполнении. Начинайте с PostgreSQL для критичных данных и переходите на Redis для высоконагруженных сценариев.

Проверь себя

5 из 8

Какой TTL рекомендуется для Idempotency-Key в большинстве платежных систем (Stripe, Amazon Pay)?

Какую команду Redis использует Redis-based дедупликация для атомарного захвата ключа?

Что такое stale lock в контексте идемпотентности?

Почему RabbitMQ consumer должен быть идемпотентным?

Для чего используется поле request_hash в таблице idempotency_keys?