Part 09 · Distributed Systems

Code snippets và bài thực hành Order–Inventory–Payment

Các đoạn code là skeleton để tự build, không phải project copy-paste hoàn chỉnh. Mỗi lab yêu cầu bổ sung schema, wiring, tests và fault injection để chứng minh invariant.

1. Domain state

enum OrderStatus {
  PENDING, INVENTORY_PENDING, INVENTORY_RESERVED,
  PAYMENT_PENDING, PAYMENT_AUTHORIZED,
  CONFIRMING, CONFIRMED, FULFILLMENT_PENDING, COMPLETED,
  CANCELLING, CANCELLED, COMPENSATION_PENDING, MANUAL_REVIEW
}

enum ReservationStatus { ACTIVE, CONFIRMED, RELEASED, EXPIRED }
enum PaymentStatus { CREATED, AUTHORIZING, AUTHORIZED, CAPTURED,
                     DECLINED, VOIDED, REFUNDED, UNKNOWN }

record Money(BigDecimal amount, Currency currency) {
  Money {
    if (amount.signum() <= 0) throw new IllegalArgumentException("amount");
    if (amount.scale() > currency.getDefaultFractionDigits())
      throw new IllegalArgumentException("scale");
  }
}

Không dùng một boolean success; state phải biểu diễn pending, terminal failure và compensation. Transition method trong aggregate từ chối state quay lùi hoặc event đến muộn.

2. Create Order + Outbox trong một transaction

@Service
class CreateOrderHandler {
  private final OrderRepository orders;
  private final OutboxRepository outbox;
  private final IdempotencyRepository idempotency;

  @Transactional
  public OrderResult handle(CreateOrder cmd) {
    return idempotency.find(cmd.customerId(), cmd.key())
        .map(record -> record.replayOrConflict(cmd.fingerprint()))
        .orElseGet(() -> {
          Order order = Order.pending(cmd.customerId(), cmd.lines());
          orders.save(order);
          outbox.save(OutboxMessage.command(
              order.id(), "ReserveInventory",
              new ReserveInventory(order.id(), order.reservationId(), cmd.lines())));
          idempotency.save(IdempotencyRecord.accepted(cmd, order.id()));
          return OrderResult.accepted(order.id());
        });
  }
}

Nếu insert Outbox fail, Order cũng rollback. Nếu HTTP response mất sau commit, client retry cùng key và nhận lại cùng orderId.

3. Inventory reservation entity

@Entity
@Table(name = "inventory_reservations",
  uniqueConstraints = @UniqueConstraint(
      name = "uk_reservation_operation", columnNames = "reservation_id"))
class InventoryReservation {
  @Id UUID id;
  UUID reservationId;
  UUID orderId;
  String sku;
  int quantity;
  Instant expiresAt;
  @Enumerated(EnumType.STRING) ReservationStatus status;
  @Version long version;

  void confirm(Instant now) {
    if (status == ReservationStatus.CONFIRMED) return;
    if (status != ReservationStatus.ACTIVE || !now.isBefore(expiresAt))
      throw new InvalidReservationTransition();
    status = ReservationStatus.CONFIRMED;
  }

  void release() {
    if (status == ReservationStatus.RELEASED) return;
    if (status == ReservationStatus.CONFIRMED)
      throw new InvalidReservationTransition();
    status = ReservationStatus.RELEASED;
  }
}

4. Atomic reserve để tránh oversell

@Modifying
@Query("""
  update Stock s
     set s.reserved = s.reserved + :qty,
         s.version = s.version + 1
   where s.sku = :sku
     and s.onHand - s.reserved >= :qty
  """)
int tryReserve(String sku, int qty);

@Transactional
public ReservationResult reserve(ReserveInventory cmd) {
  return reservations.findByReservationId(cmd.reservationId())
      .map(ReservationResult::replay)
      .orElseGet(() -> {
        if (stock.tryReserve(cmd.sku(), cmd.quantity()) != 1)
          return ReservationResult.outOfStock();
        InventoryReservation r = reservations.save(
            Reservation.active(cmd, clock.instant().plus(holdTtl)));
        outbox.save(OutboxMessage.event(r.orderId(), "InventoryReserved", r));
        return ReservationResult.reserved(r.id());
      });
}

Trong project thật, xử lý nhiều SKU cần lock order/batch strategy và rollback toàn bộ nếu một line thiếu hàng. Unique reservation ID chống reserve lặp; affected rows bằng 0 là business outcome, không phải technical retry.

5. Inbox consumer

@Transactional
public void onReserveInventory(Envelope<ReserveInventory> envelope) {
  if (!inbox.tryInsert("inventory", envelope.eventId())) return;
  ReservationResult result = inventory.reserve(envelope.payload());
  outbox.save(OutboxMessage.event(
      envelope.aggregateId(), result.eventType(), result.payload()));
} // Inbox + inventory effect + result Outbox commit cùng nhau

Không ghi Inbox ở transaction riêng trước effect: crash ở giữa sẽ làm redelivery bị bỏ qua dù effect chưa chạy.

6. Saga orchestrator transition

@Transactional
public void onInventoryReserved(Envelope<InventoryReserved> event) {
  if (!inbox.tryInsert("order-saga", event.eventId())) return;
  OrderSaga saga = sagas.lockByOrderId(event.aggregateId());
  if (!saga.canApply(event.aggregateVersion())) return; // stale/late event

  saga.inventoryReserved(event.payload().reservationId());
  outbox.save(OutboxMessage.command(
      saga.orderId(), "AuthorizePayment",
      new AuthorizePayment(saga.paymentOperationId(), saga.orderId(), saga.money())));
}

Orchestrator chỉ quản workflow state/commands, không sửa database Inventory hoặc Payment. Mỗi handler persist state + next command atomically.

7. Payment unknown outcome

public PaymentResult authorize(AuthorizePayment cmd) {
  try {
    return gateway.authorize(cmd.paymentOperationId(), cmd.money());
  } catch (HttpTimeoutException | ConnectionResetException ex) {
    return PaymentResult.unknown(cmd.paymentOperationId());
  }
}

@Scheduled(fixedDelayString = "${payment.reconcile-delay}")
void reconcileUnknownPayments() {
  payments.findUnknownBatch().forEach(payment -> {
    ProviderStatus status = gateway.query(payment.operationId());
    transactionTemplate.executeWithoutResult(tx ->
        paymentReconciler.apply(payment.id(), status));
  });
}

Không gọi lại authorize bằng operation ID mới. Query/retry cùng key và persist UNKNOWN; reconciliation có batch limit, backoff và metric oldest-unknown age.

8. Compensation command

@Transactional
public void cancelAfterPaymentDeclined(OrderSaga saga) {
  saga.startCancelling("PAYMENT_DECLINED");
  outbox.save(OutboxMessage.command(
      saga.orderId(), "ReleaseInventory",
      new ReleaseInventory(saga.reservationId(), saga.orderId())));
}

@Transactional
public void onInventoryReleased(InventoryReleased event) {
  OrderSaga saga = sagas.lockByOrderId(event.orderId());
  saga.cancelled();
  orders.markCancelled(event.orderId(), "PAYMENT_DECLINED");
}

Nếu release fail, Saga ở COMPENSATION_PENDING; retry cùng reservation ID. Không đánh dấu CANCELLED trước khi compensation bắt buộc đã hoàn tất, trừ khi external API contract chủ động tách “customer cancelled” khỏi “cleanup pending”.

9. Reservation expiry worker

@Modifying(clearAutomatically = true)
@Query("""
  update InventoryReservation r
     set r.status = com.acme.ReservationStatus.EXPIRED
   where r.status = com.acme.ReservationStatus.ACTIVE
     and r.expiresAt < :now
  """)
int expireBatch(Instant now);

// Sau bulk update cần phát events/recalculate stock theo design đã chọn;
// bulk DML bỏ qua entity callbacks và entities đang managed có thể stale.

Thực tế nên claim rows theo batch và tạo Outbox event cho từng expiry trong cùng transaction, hoặc dùng entity loop bounded. Không dựa vào entity callback bị bypass bởi bulk update.

10. Outbox relay skeleton

void publishBatch() {
  List<OutboxMessage> batch = transactionTemplate.execute(tx ->
      outbox.claimUnpublished(workerId, batchSize));

  for (OutboxMessage message : batch) {
    broker.publish(message); // có thể thành công rồi process crash
    transactionTemplate.executeWithoutResult(tx ->
        outbox.markPublished(message.id()));
  }
}

Crash sau publish trước mark tạo duplicate; consumer phải idempotent. Claim có lease/timeout để worker khác lấy lại. Production có thể polling hoặc CDC/Debezium.

11. Reconciliation query

select o.id, o.status, r.status as reservation_status, p.status as payment_status
from orders o
left join inventory_reservations r on r.order_id = o.id
left join payments p on p.order_id = o.id
where o.updated_at < now() - interval '5 minutes'
  and o.status in ('INVENTORY_PENDING', 'PAYMENT_PENDING',
                   'CANCELLING', 'COMPENSATION_PENDING');

Job phân loại discrepancy thành safe auto-repair, retry command hoặc manual review. Không join cross-service production databases trực tiếp; snippet biểu diễn logic đối chiếu—thực tế dùng read model/API/export có ownership rõ.

12. REST status resource

@PostMapping("/orders")
ResponseEntity<OrderView> create(
    @RequestHeader("Idempotency-Key") String key,
    @Valid @RequestBody CreateOrderRequest request) {
  OrderResult result = handler.handle(map(key, request));
  return ResponseEntity.accepted()
      .location(URI.create("/orders/" + result.orderId()))
      .body(result.view());
}

@GetMapping("/orders/{id}")
OrderView status(@PathVariable UUID id) { return query.find(id); }

13. Fault injection interface

interface FailureInjector {
  void beforeCommit(String checkpoint);
  void afterCommitBeforeResponse(String checkpoint);
}

// Test profiles inject: throw-before-commit, delay-after-commit,
// drop-response, duplicate-event, reorder-event, crash-before-ack.

Trong production code không để random failure flag. Lab có deterministic checkpoint/latch để tái hiện chính xác boundary.

14. Integration test: response lost after commit

@Test
void retryAfterLostResponseReturnsSameOrderWithoutDoubleReserve() {
  String key = UUID.randomUUID().toString();
  client.createOrder(key, request, DROP_RESPONSE_AFTER_COMMIT);

  OrderView retried = client.createOrder(key, request, NORMAL);

  assertThat(orderRowsByKey(key)).isEqualTo(1);
  assertThat(reservationsByOrder(retried.id())).hasSize(1);
  assertThat(stockReserved("SKU-1")).isEqualTo(1);
}

15. Integration test: compensation failure

@Test
void declinedPaymentAndTemporaryReleaseFailureEventuallyCancel() {
  payment.stubDeclined();
  inventory.failNextReleaseAttempts(2);

  UUID orderId = createOrder();
  await().untilAsserted(() ->
      assertThat(orderStatus(orderId)).isEqualTo(COMPENSATION_PENDING));
  await().untilAsserted(() -> {
      assertThat(orderStatus(orderId)).isEqualTo(CANCELLED);
      assertThat(activeReservations(orderId)).isZero();
  });
}

16. Bài thực hành theo cấp độ

LabYêu cầuGate
1 · ReservationReserve/confirm/release/expire + concurrent requestsKhông oversell; duplicate không tăng reserved.
2 · SagaOrder→Inventory→Payment state machine + Outbox/InboxKill process tại từng commit boundary vẫn hội tụ.
3 · Unknown outcomeDrop response sau Inventory/Payment commitRetry cùng ID, không duplicate effect.
4 · CompensationPayment declined, release fail hai lầnQuan sát COMPENSATION_PENDING rồi CANCELLED.
5 · OrderingDuplicate/reorder PaymentAuthorized và OrderCancelledState không quay lùi; stale event bị bỏ/audit.
6 · ReconciliationTạo stuck/missing/duplicate fixturesReport đúng, repair idempotent, ambiguous vào review.
7 · OverloadInventory chậm + retries + bounded poolsKhông retry storm; critical throughput được bảo vệ.
8 · ObservabilityMetrics/traces/runbook cho một incidentTìm được stuck stage và recovery bằng evidence.

17. Expected deliverables