asman.malikov_ EN

· ~16 мин чтения

Temporal-сага платежей в проде: геораспределённые депозиты

Как 6-шаговая saga на Temporal с LIFO-компенсациями делает кросс-региональные депозиты effectively exactly-once: реальные таймауты, идемпотентность, FX-replay.

temporalsagagopaymentsidempotencydistributed-systems

Платёжный провайдер присылает success-коллбэк. Сервис зачисляет деньги на кошелёк — и падает до того, как успевает это зафиксировать. Провайдер не видит ACK и присылает коллбэк повторно. Итог: клиенту зачислили дважды, леджер разошёлся с кошельком, а кто-то из команды проводит вечер за написанием скрипта сверки. На платёжно-биллинговой платформе, над которой я работаю, этот сценарий стал заметно опаснее в тот день, когда сервисы кошельков разъехались по географическим регионам: «зачислить на кошелёк» перестало быть локальным вызовом и превратилось в gRPC-хоп через WAN. Ниже — как мы сделали кросс-региональные депозиты effectively exactly-once с помощью шестишаговой SAGA на Temporal, на Go, с реальными таймаутами и ретраями вместо дефолтов из туториалов.

TL;DR

  • Сервисы кошельков и бонусов живут в регионах (eu/asia/default), платёжный сервис остался в core. Депозит теперь — распределённая транзакция через WAN, которую оркестрирует SAGA на Temporal.
  • Шесть шагов, компенсации в порядке LIFO. Стек компенсаций выполняется на disconnected-контексте и ретраится без ограничения попыток.
  • Каждый шаг с побочным эффектом несёт детерминированный ключ идемпотентности, выведенный из саги: operation_guid = saga_id для зачисления на кошелёк, "revert-" + operation_guid для отмены, saga_id как producer key в Kafka.
  • Курс валют запрашивается один раз и при каждом ретрае воспроизводится из event history Temporal. Ретрай не может изменить сумму — это самая важная строчка во всём дизайне.
  • На входе стоит Kafka как durable-слой доставки: консьюмер коллбэков ретраит до тех пор, пока SAGA не принята, — система переживает даунтайм Temporal.
  • Сбой после зачисления на кошелёк — это forward retry, а не откат. Откатывать деньги, которые уже уехали, — верный способ устроить именно тот инцидент, от которого защищались.

Прямые вызовы перестали быть безопасными, когда кошельки уехали в регионы

Изначальный флоу был простым: платёжный сервис получает коллбэк провайдера, вызывает сервис кошельков по gRPC внутри того же кластера, обновляет строку платежа и публикует событие в Kafka для бонусного консьюмера. Три побочных эффекта, один процесс, низкая latency. Если что-то ломалось посередине, радиус поражения был маленьким, а вся retry-логика умещалась в пару if.

Потом кошельки и бонусы переехали в регионы ради data locality — и та же последовательность стала выглядеть так:

 provider ──callback──> payment-svc [core]

                             │ gRPC over WAN
                             v
                        wallet-svc [eu]        balance: +100  OK

                             X  payment-svc crashes here:
                                payment row still "pending",
                                no Kafka event produced
                             
 provider resends callback (no ACK seen)

                             v
                        wallet-svc [eu]        balance: +200  <- double credit

Три свойства нового мира сделали наивный подход неприемлемым:

  1. WAN оказался внутри транзакции. Кросс-региональный gRPC-вызов падает по причинам, которых у локального вызова никогда не было: нестабильная маршрутизация, деградация региона, целый регион в дауне. Каждый сбой оставляет вас посреди последовательности.
  2. Частичное выполнение стало нормой, а не краевым случаем. С тремя побочными эффектами в двух failure-доменах ситуация «шаг 2 выполнен, шаг 3 нет» — это состояние, под которое вы проектируете с самого начала.
  3. Консьюмер коллбэков не имеет права блокироваться. Коллбэки провайдеров приходят через Kafka. Если консьюмер блокируется на кросс-региональном вызове, один медленный регион раздувает лаг по коллбэкам для всех.

У нас уже работал transactional outbox для надёжной публикации событий (см. соседнюю статью), и он никуда не делся. Но outbox даёт at-least-once доставку одного события — он не даёт многошаговую транзакцию с пошаговым откатом через регионы. Хореография — когда каждый сервис реагирует на события — размазала бы state machine депозита по трём кодовым базам, и не осталось бы единого места, где можно ответить на вопрос «где сейчас этот депозит и почему он застрял?». Для денежного флоу из шести шагов с семантикой компенсаций мы хотели оркестрацию с durable-состоянием. Это ровно та задача, под которую сделан Temporal.

Устройство системы

Два компонента, один workflow:

  • billing-saga-core (core): Temporal-воркер плюс gRPC-сервер. Владеет workflow депозита и всеми core-активностями.
  • billing-saga-region (по одному на регион): выполняет операции с кошельком локально, обращаясь к региональному сервису кошельков.
                    CORE                                REGION (eu / asia)
 ┌────────────────────────────────────────┐    ┌─────────────────────────────┐
 │  Kafka (provider callbacks)            │    │                             │
 │     │                                  │    │   billing-saga-region       │
 │     v                                  │    │   worker polls task queue   │
 │  payment-svc ──gRPC (fire-and-forget)──│    │   billing-saga-region-eu    │
 │     returns saga_id immediately        │    │          │                  │
 │     │                                  │    │          v                  │
 │     v                                  │    │   wallet-svc (local gRPC)   │
 │  billing-saga-core ◄──► Temporal ──────┼WAN─┼──────────┘                  │
 │  (workflow + core     server           │    │                             │
 │   activities)         PG persistence   │    └─────────────────────────────┘
 │                       ES visibility    │
 │     │                                  │
 │     v                                  │
 │  Kafka (deposit-events → bonus)        │
 └────────────────────────────────────────┘

Роутинг устроен в два слоя. Запросы с JWT несут region claim (default/eu/asia; отсутствие claim означает default — для обратной совместимости). Коллбэки провайдеров приходят через Kafka без JWT, поэтому выделенный Region Service резолвит GetCustomerRegion(customer_guid) -> region — Redis перед Postgres, с низкой latency, потому что он стоит на горячем пути каждого коллбэка.

Платёжный сервис остаётся тонким. Для default-региона сохранился старый прямой вызов кошелька — нет смысла платить за оркестрацию там, где вызов локальный. Для любого другого региона — fire-and-forget gRPC-вызов в saga-core, который немедленно возвращает saga_id и никогда не блокирует консьюмер коллбэков.

Гео-роутинг внутри workflow — это просто task queues. Core-активности выполняются на billing-saga-core; активности кошелька шедулятся на billing-saga-region-{name}, и эту очередь поллит только воркер, задеплоенный в соответствующем регионе. Код workflow не знает и не хочет знать, где физически находится воркер, — имя очереди и есть решение о маршрутизации.

Шесть шагов, и что каждая компенсация означает в деньгах

В депозитной SAGA шесть шагов. Компенсации выполняются в порядке LIFO, и этот порядок принципиален: компенсация каждого шага предполагает, что шаги после него уже откачены.

 forward ─────────────────────────────────────────────────────>
  1 FindOrCreatePayment   [core]         comp: mark cancelled
  2 ParseCallbackStatus   [core, local]  read-only, no comp
  3 ConvertCurrency       [core]         rate pinned in history, no comp
  4 UpdateWalletBalance   [region]       key: saga_id
                                         comp: RevertWalletBalance
                                               key: "revert-"+saga_id
  5 UpdatePaymentStatus   [core]         comp: mark failed
  6 ProduceDepositEvent   [core]         Kafka, producer key: saga_id

 <───────────────────────────────────────────── compensate (LIFO)
  fail at 6:  5' mark failed → 4' revert wallet → 1' cancel
  fail at 5:  4' revert wallet → 1' cancel
  fail at 4:  1' cancel   (money never moved)

Теперь то же самое в терминах денег:

  1. FindOrCreatePayment — найти платёж по external_id; создать, если это webhook-only флоу, где коллбэк — первое, что мы вообще слышим об этом депозите. Компенсация: пометить отменённым. Денег здесь ещё нет, но осиротевшая строка в статусе «pending» — это тикет в саппорт.
  2. ParseCallbackStatus — распарсить статус провайдера. Read-only, выполняется как local activity (это чистое вычисление; полноценный круг активности через сервер был бы расточительством). Компенсации нет, потому что откатывать нечего.
  3. ConvertCurrency — получить курс из валютного сервиса и посчитать сумму зачисления. Об этом подробнее ниже — это центральный элемент всей конструкции.
  4. UpdateWalletBalance — gRPC-депозит в региональный кошелёк. Здесь двигаются деньги. Компенсация — RevertWalletBalance, отдельная операция со своим ключом. Не «delete»: кошелёк хранит в леджере обе записи.
  5. UpdatePaymentStatus — пометить платёж успешным, приложить логи провайдера и конвертации, номер транзакции. Компенсация: пометить failed.
  6. ProduceDepositEvent — опубликовать в топик deposit-events для бонусного консьюмера. Собственной компенсации нет; если шаг исчерпал ретраи, LIFO-стек раскручивает всё, что выше.

Почему LIFO, а не «откатываем в любом порядке»? Потому что компенсация шага 5 (пометить failed) честна только в том случае, если событие из шага 6 так и не стало видимым для консьюмеров, а revert из шага 4 должен произойти, пока строка платежа ещё отражает реальность. Раскрутка в обратном порядке зависимостей означает, что каждая компенсация выполняется против того состояния, которое создал её forward-шаг.

Код workflow — с реальными числами

Урезанный, но честный Go. Числа — продовые, а не заглушки.

func DepositWorkflow(ctx workflow.Context, in DepositInput) error {
	sagaID := workflow.GetInfo(ctx).WorkflowExecution.ID

	// Forward activities: fail fast, bounded retries.
	ao := workflow.ActivityOptions{
		ScheduleToCloseTimeout: 30 * time.Second,
		StartToCloseTimeout:    10 * time.Second,
		// Our activities call activity.RecordHeartbeat; drop this if yours don't,
		// or a legal 6-10s call will die on heartbeat timeout instead.
		HeartbeatTimeout: 5 * time.Second,
		RetryPolicy: &temporal.RetryPolicy{
			InitialInterval:    time.Second,
			BackoffCoefficient: 5.0, // waits 1s, then 5s (3 attempts total)
			MaximumAttempts:    3,
		},
	}
	coreCtx := workflow.WithActivityOptions(ctx, ao)

	regionAO := ao
	regionAO.TaskQueue = "billing-saga-region-" + in.Region
	regionCtx := workflow.WithActivityOptions(ctx, regionAO)

	// LIFO compensation stack. Runs on a disconnected context so it
	// still executes after cancellation; compensations retry forever.
	var comps []func(workflow.Context) error
	unwind := func() {
		dc, _ := workflow.NewDisconnectedContext(ctx)
		dc = workflow.WithActivityOptions(dc, workflow.ActivityOptions{
			StartToCloseTimeout: 10 * time.Second,
			RetryPolicy:         &temporal.RetryPolicy{MaximumAttempts: 0}, // unlimited
		})
		for i := len(comps) - 1; i >= 0; i-- {
			_ = comps[i](dc)
		}
	}

	// 1. Find or create the payment.
	var pay Payment
	if err := workflow.ExecuteActivity(coreCtx,
		a.FindOrCreatePayment, in.ExternalID).Get(ctx, &pay); err != nil {
		return err
	}
	comps = append(comps, func(c workflow.Context) error {
		return workflow.ExecuteActivity(c,
			a.MarkPaymentCancelled, pay.GUID).Get(c, nil)
	})

	// 2. Parse provider status (local activity: pure computation).
	lctx := workflow.WithLocalActivityOptions(ctx, workflow.LocalActivityOptions{
		StartToCloseTimeout: 10 * time.Second,
	})
	var st CallbackStatus
	if err := workflow.ExecuteLocalActivity(lctx,
		a.ParseCallbackStatus, in.RawCallback).Get(ctx, &st); err != nil {
		unwind()
		return err
	}

	// wait-confirm flow: park until the psp-callback signal arrives.
	// WorkflowExecutionTimeout (1h) bounds the wait.
	if st.WaitConfirm {
		var cb ProviderCallback
		ch := workflow.GetSignalChannel(ctx, "psp-callback")
		ch.Receive(ctx, &cb)
		st = cb.Status
	}

	// 3. FX conversion. Executes once; the result lives in event
	// history and is replayed on every retry. Never re-fetched.
	var amt ConvertedAmount
	if err := workflow.ExecuteActivity(coreCtx,
		a.ConvertCurrency, pay.Currency, in.WalletCurrency,
		pay.Amount).Get(ctx, &amt); err != nil {
		unwind()
		return err
	}

	// 4. Credit the regional wallet. operation_guid = saga_id.
	if err := workflow.ExecuteActivity(regionCtx,
		a.UpdateWalletBalance, WalletOp{
			OperationGUID: sagaID,
			CustomerGUID:  in.CustomerGUID,
			Amount:        amt.Value,
		}).Get(ctx, nil); err != nil {
		unwind()
		return err
	}
	comps = append(comps, func(c workflow.Context) error {
		// Same regional task queue, but MaximumAttempts: 0 — the revert
		// must eventually land, so it inherits the unlimited-retry policy.
		rc := workflow.WithActivityOptions(c, regionAOForComp(in.Region))
		return workflow.ExecuteActivity(rc,
			a.RevertWalletBalance, WalletOp{
				OperationGUID: "revert-" + sagaID,
				CustomerGUID:  in.CustomerGUID,
				Amount:        amt.Value,
			}).Get(rc, nil)
	})

	// 5. Mark the payment successful (idempotent UPDATE).
	if err := workflow.ExecuteActivity(coreCtx,
		a.UpdatePaymentStatus, pay.GUID, st, amt).Get(ctx, nil); err != nil {
		unwind()
		return err
	}
	comps = append(comps, func(c workflow.Context) error {
		return workflow.ExecuteActivity(c,
			a.MarkPaymentFailed, pay.GUID).Get(c, nil)
	})

	// 6. Publish for the bonus consumer (idempotent producer, key=saga_id).
	if err := workflow.ExecuteActivity(coreCtx,
		a.ProduceDepositEvent, sagaID, pay, amt).Get(ctx, nil); err != nil {
		unwind()
		return err
	}
	return nil
}

Комментарии к числам — потому что именно на дефолтах туториалы и врут:

  • StartToClose 10s, ScheduleToClose 30s. Зачисление на кошелёк, которое идёт дольше 10 секунд, — это не «медленно», это «сломано». ScheduleToClose ограничивает весь бюджет ретраев целиком: мёртвый регион валит шаг за полминуты, а не висит бесконечно.
  • 3 попытки, backoff 1s, затем 5s. Достаточно, чтобы пережить кратковременный сбой gRPC; недостаточно, чтобы замаскировать настоящий outage. Платёжная сага, которая молча ретраится час, хуже той, что громко падает внутри своего 30-секундного окна.
  • Компенсации ретраятся без ограничений (MaximumAttempts: 0). Продвижение вперёд — опционально; откат зачисления на кошелёк — нет. Если revert не удаётся применить, мы хотим видеть красную активность на дашборде, пока регион не восстановится или не вмешается человек, — а не проглоченную ошибку и разъехавшийся баланс.
  • WorkflowExecutionTimeout 1h ограничивает всё, включая ожидание сигнала в wait-confirm флоу. Ни одна сага не может «застрять» дольше часа — по построению.

Дубликаты коллбэков отсекаются ещё до старта workflow — через Workflow ID:

run, err := c.ExecuteWorkflow(ctx, client.StartWorkflowOptions{
	ID:        fmt.Sprintf("deposit-%s-%s", cb.ExternalID, cb.EventType),
	TaskQueue: "billing-saga-core",
	WorkflowExecutionTimeout: time.Hour,
	WorkflowIDReusePolicy: enumspb.WORKFLOW_ID_REUSE_POLICY_REJECT_DUPLICATE,
	TypedSearchAttributes: temporal.NewSearchAttributes(
		saPaymentGUID.ValueSet(cb.PaymentGUID),
		saCustomerGUID.ValueSet(cb.CustomerGUID),
		saRegion.ValueSet(region),
	),
}, DepositWorkflow, input)

Провайдер прислал коллбэк повторно? Тот же external_id, тот же event_type, тот же Workflow ID — Temporal отклоняет дублирующий старт. Дедупликация стоит ноль строк прикладного кода. А search attributes, выставленные на старте (payment_guid, customer_guid, region), — плюс amount и saga_status, которые обновляются по ходу выполнения, — позволяют саппорту найти любой депозит в UI Temporal, не притрагиваясь к базе.

Курс валют закреплён в event history — и в этом весь смысл

Депозит приходит в EUR, кошелёк в другой валюте, и ConvertCurrency запрашивает курс у валютного сервиса. Курсы двигаются постоянно. Теперь представьте: зачисление на кошелёк в шаге 4 отваливается по таймауту, ретраится и проходит через 26 секунд. Если бы что-нибудь на этом retry-пути перезапрашивало курс, зачисленная сумма могла бы разойтись с той, что мы вот-вот запишем в платёж и опубликуем в Kafka. Клиенту зачислили X, в книгах записано Y — и кто-то снова садится писать тот самый скрипт сверки.

Модель исполнения Temporal закрывает эту дыру бесплатно — если её уважать. Результат активности, однажды записанный в event history, воспроизводится при каждом последующем replay workflow и никогда не выполняется заново. ConvertCurrency выполняется один раз, его результат становится фактом в истории, и каждый ретрай каждого последующего шага — и каждое падение воркера с полным replay всего workflow — видит идентичную сумму. Курс фиксируется в момент первого выполнения и больше не меняется.

Следствие — классическое правило детерминизма, только с деньгами на кону: никогда не запрашивайте курс, не читайте часы и не генерируйте ID прямо в коде workflow. Делайте это в активности — или выводите детерминированно из состояния workflow. В момент, когда кто-то «упрощает» конвертацию до обычного вызова функции внутри workflow, replay перестаёт быть детерминированным, а ретраи начинают менять суммы. У нас это review-blocking правило, и это первое, что я проверяю в любом диффе с workflow.

Ключи идемпотентности — вот и весь фокус

Temporal даёт effectively-once семантику для workflow, но активности — это at-least-once: воркер может выполнить зачисление на кошелёк и умереть до того, как отчитается о завершении, — и Temporal зашедулит активность снова. Exactly-once эффекты поэтому обязаны приходить со стороны целевых систем, а это значит — детерминированные ключи:

  • UpdateWalletBalance: operation_guid = saga_id. ID саги генерируется один раз, до любого побочного эффекта, и стабилен через все ретраи и replay. Сервис кошельков дедуплицирует по нему, так что at-least-once исполнение схлопывается ровно в одно зачисление. Выводить ключ из чего-либо, сгенерированного во время исполнения (таймстемп, случайный UUID внутри активности), — значит тихо вернуть двойные зачисления обратно.
  • RevertWalletBalance: "revert-" + operation_guid. Revert — самостоятельная идемпотентная операция со своим ключом. Бесконечно ретраить компенсацию безопасно только потому, что её повтор безвреден.
  • UpdatePaymentStatus: идемпотентный UPDATE по payment_guid + статусу:
UPDATE payments
SET status = 'successful',
    transaction_number = $2,
    provider_log = $3,
    conversion_log = $4,
    updated_at = now()
WHERE payment_guid = $1
  AND status <> 'successful';
  • ProduceDepositEvent: идемпотентный Kafka-продюсер с ключом saga_id, чтобы бонусный консьюмер мог дедуплицировать, даже если produce-активность ретраится после успешной, но неподтверждённой отправки.

Один ID саги, прошитый через все побочные эффекты. Когда что-то идёт не так, этот единственный ID соединяет историю Temporal, леджер кошелька, строку платежа и событие в Kafka.

Kafka на входе — значит, Temporal имеет право падать

Оркестратор — это зависимость, а зависимости отказывают. Мы сознательно не поставили Temporal на путь приёма коллбэков. Коллбэки провайдеров попадают в Kafka; задача консьюмера — только добиться, чтобы сага была принята, и он ретраит до успеха. Если Temporal лежит десять минут, коллбэки копятся в топике с durability-гарантиями Kafka и разгребаются, когда он поднимется. Ни потерянных депозитов, ни шторма провайдерских ретраев в мёртвый эндпоинт. Temporal оркестрирует движение денег; Kafka гарантирует, что запрос на это движение вообще будет доставлен. Система переживает даунтайм оркестратора потому, что оркестратор никогда не был входной дверью.

Отказы, под которые проектировали

Военных историй не будет — у этой системы не было публичного инцидента, который можно пересказать, а выдумывать я не собираюсь. Зато могу показать матрицу отказов, против которой мы проектировали, — реальные уроки живут именно здесь:

  • Регион лежит. Активности кошелька на billing-saga-region-eu отваливаются по таймауту (ScheduleToClose 30s), сага падает, компенсации отрабатывают, срабатывает алёрт. Ограниченно, громко, чисто. Альтернатива — бесконечный ретрай в мёртвый регион — прячет outage внутри «pending»-саг.
  • gRPC кошелька лежит, регион жив. 3 попытки внутри 30-секундного окна ScheduleToClose, затем компенсация. Кратковременные сбои поглощаются, настоящие отказы всплывают быстро.
  • Сбой после зачисления на кошелёк. Вот случай, который отделяет хорошие дизайны от плохих. Деньги уже уехали. Шаги 5 и 6 — идемпотентная бухгалтерия. Правильный ход — forward retry, а не откат: ретраить пометку платежа и публикацию события, пока они не выполнятся. Откатить зачисление, которое клиент, возможно, уже видит, — из-за того, что икнул produce в Kafka, — значит превратить ретраебельный глюк в видимое клиенту падение баланса. Rollback — для «депозит не может завершиться», а не для «шаг 6 что-то тормозит».
  • Дубликаты коллбэков. Отсекаются на уровне Workflow ID, до любого бизнес-кода.
  • Застрявшая сага. Структурно невозможна дольше часа: WorkflowExecutionTimeout принудительно закрывает выполнение со статусом timed out, и на это терминальное состояние летит алёрт. Важно, чего при этом не происходит: execution timeout убивает workflow, не запуская компенсации, — именно поэтому каждый шаг идемпотентен и прошит ключом саги, а сага с таймаутом попадает в очередь на разбор человеком, а не делает вид, что прибрала за собой.

Транспорт через WAN: Temporal long-poll против NATS request/reply

Для транспорта core → регион на доске было два дизайна.

Вариант A — региональный Temporal-воркер. В регионе крутится обычный Temporal-воркер, который long-poll’ит core-сервер Temporal по gRPC через WAN, на своей очереди billing-saga-region-{name}. Плюсы: одна система, одна модель исполнения; ретраи, таймауты и история достаются бесплатно; соединения из региона исходящие, так что файрволы и NAT перестают быть проблемой; мониторинг — те же метрики Temporal везде. Цена: трафик воркер—сервер идёт через WAN, и нестабильность канала проявляется как schedule-to-start latency, за которой нужно следить.

Вариант B — NATS request/reply. Сабжекты billing.{region}.wallet.{update|revert|get}, stateless-подписчик в каждом регионе, региональная активность workflow превращается в core-активность, делающую NATS request. Плюсы: лёгкий транспорт, заточенный ровно под этот хоп; региональный footprint — тривиальный подписчик без зависимости от Temporal. Цена: вторая messaging-система, которую нужно деплоить, защищать и мониторить в каждом регионе; семантика ретраев и таймаутов теперь живёт в двух местах; вы заново изобрели кусок того, что task queue уже давала.

Честное резюме: вариант A не наращивает операционную сложность и переносит риск в качество WAN-канала; вариант B покупает независимость транспорта ценой дополнительной инфраструктуры. Если NATS у вас и так везде — B можно защитить. Если Temporal — единственная новая система в картине, добавлять вторую только ради того, чтобы избежать long-poll через WAN, — это сложность без выигрыша.

Что мы отвергли — и когда стоит отвергнуть Temporal

  • Самописная сага на Kafka. Таблицы state machine, sweeper’ы таймаутов, диспетчеры компенсаций — в итоге вы пишете худший Temporal без гарантий replay и со всей поддержкой на себе.
  • Только outbox + хореография. Отлично для «это событие обязано быть опубликовано»; неправильная форма для «эти шесть шагов обязаны завершиться или быть откачены по порядку». Состояние депозита становится неявным, размазанным по консьюмерам.
  • Двухфазный коммит. Локи, удерживаемые через WAN, координатор как единая точка отказа, и половина участников (Kafka, почти-внешние API кошельков) всё равно не говорят на XA.
  • Легковесные Temporal-подобные библиотеки. Меньше эксплуатационной нагрузки, но персистентность, visibility и тулинг для replay — ровно та часть, которая нужна, когда речь о деньгах. И ровно её легковесные варианты и вырезают.

И честная обратная сторона — когда Temporal избыточен: однорегиональные флоу, где одна транзакция БД накрывает все побочные эффекты; флоу, где единственное требование — надёжная публикация события (outbox сам по себе проще и дешевле); команды без ресурса оперировать ещё одной stateful-системой — кластер Temporal с персистентностью на PostgreSQL и visibility на Elasticsearch — это настоящая инфраструктура со своими настоящими отказами. Default-регион мы оставили на прямом вызове кошелька ровно по этой причине: оркестрация там, где её требует WAN, и скучные локальные вызовы там, где не требует.

Прежде чем выкатывать такое

  • У каждой активности с побочным эффектом есть детерминированный ключ идемпотентности, выведенный из ID саги, и целевая система реально по нему дедуплицирует. Тестируйте дедупликацию, а не только happy path.
  • Всё недетерминированное — курсы валют, часы, генерируемые ID — живёт в активностях и никогда не инлайнится в код workflow. Проверяется replay-тестами против экспортированных историй в CI.
  • Forward-ретраи ограничены (попытки + ScheduleToClose); ретраи компенсаций — безлимитные и под алёртами.
  • Сбой после того, как деньги уехали, — это явно forward retry, и команда может объяснить почему.
  • WorkflowExecutionTimeout выставлен — и команда знает, что он завершает выполнение без компенсаций, и кому прилетит пейдж, когда он сработает.
  • Workflow ID кодируют естественные ключи дедупликации (external_id + тип события) с reject-duplicate reuse policy.
  • Путь приёма переживает даунтайм оркестратора: durable-очередь спереди и консьюмер, который ретраит, пока workflow не принят.
  • Search attributes покрывают поля, по которым саппорт реально будет искать: ID платежа, ID клиента, регион, сумма, статус саги.
  • У вас записано, что происходит, когда компенсация не может выполниться N часов, — и в этом цикле есть человек.

Смежное

← Назад в блог