getOrderHistory(@PathVariable UUID orderId) {
+ return orderEventService.getHistory(orderId);
+ }
+
+ /**
+ * SSE endpoint — streams real-time order status updates to the client.
+ *
+ * Usage (JavaScript):
+ *
+ * const es = new EventSource('/api/orders/{orderId}/status-stream');
+ * es.addEventListener('status-update', e => console.log(JSON.parse(e.data)));
+ * es.addEventListener('complete', () => es.close());
+ *
+ *
+ * Events emitted:
+ * connected — immediately on subscription
+ * status-update — on every saga step (INVENTORY_RESERVED, PAYMENT_COMPLETED, SHIPPED, etc.)
+ * complete — when the order reaches a terminal state (stream then closes)
+ *
+ * The connection is held open for up to 5 minutes. If no terminal state is
+ * reached by then, the client should reconnect and poll {@link #getOrderById}.
+ */
+ @GetMapping(value = "/{orderId}/status-stream", produces = MediaType.TEXT_EVENT_STREAM_VALUE)
+ @Operation(
+ summary = "Stream real-time order status via SSE",
+ description = "Opens a Server-Sent Events stream that pushes status updates as the saga progresses"
+ )
+ public SseEmitter streamOrderStatus(
+ @Parameter(description = "Order ID to subscribe to") @PathVariable UUID orderId) {
+ return orderStatusEmitter.subscribe(orderId);
+ }
+}
diff --git a/order-service/src/main/java/com/hacisimsek/order/dto/OrderResponse.java b/order-service/src/main/java/com/hacisimsek/order/dto/OrderResponse.java
index 44208cb..264fbb2 100644
--- a/order-service/src/main/java/com/hacisimsek/order/dto/OrderResponse.java
+++ b/order-service/src/main/java/com/hacisimsek/order/dto/OrderResponse.java
@@ -5,7 +5,6 @@
import lombok.Builder;
import lombok.Data;
import lombok.NoArgsConstructor;
-
import java.math.BigDecimal;
import java.time.Instant;
import java.util.List;
diff --git a/order-service/src/main/java/com/hacisimsek/order/eventsourcing/OrderEvent.java b/order-service/src/main/java/com/hacisimsek/order/eventsourcing/OrderEvent.java
new file mode 100644
index 0000000..80da1bc
--- /dev/null
+++ b/order-service/src/main/java/com/hacisimsek/order/eventsourcing/OrderEvent.java
@@ -0,0 +1,91 @@
+package com.hacisimsek.order.eventsourcing;
+
+import com.hacisimsek.order.model.Order;
+import jakarta.persistence.*;
+import lombok.*;
+
+import java.time.Instant;
+import java.util.UUID;
+
+/**
+ * Immutable event record that captures every state transition of an Order.
+ *
+ * This is the Event Sourcing append-only log. Every time an order's status
+ * changes, a new OrderEvent row is inserted — never updated, never deleted.
+ *
+ * The full history of an order is the ordered sequence of its OrderEvents.
+ * The current state can always be rebuilt by replaying them in sequence.
+ *
+ * Indexing:
+ * - orderId + occurredAt for chronological history queries
+ * - correlationId for saga-level tracing across services
+ */
+@Entity
+@Table(name = "order_events",
+ indexes = {
+ @Index(name = "idx_order_events_order_id", columnList = "orderId, occurredAt"),
+ @Index(name = "idx_order_events_correlation", columnList = "correlationId")
+ })
+@Data
+@Builder
+@NoArgsConstructor
+@AllArgsConstructor
+public class OrderEvent {
+
+ @Id
+ @GeneratedValue(strategy = GenerationType.UUID)
+ private UUID id;
+
+ /** The order this event belongs to */
+ @Column(nullable = false)
+ private UUID orderId;
+
+ /** Correlation ID that ties this event to the saga and to gateway logs */
+ private UUID correlationId;
+
+ /** The type of event — maps to OrderStatus transitions */
+ @Enumerated(EnumType.STRING)
+ @Column(nullable = false)
+ private EventType eventType;
+
+ /** The new status after this event */
+ @Enumerated(EnumType.STRING)
+ @Column(nullable = false)
+ private Order.OrderStatus newStatus;
+
+ /** The previous status before this event (null for ORDER_CREATED) */
+ @Enumerated(EnumType.STRING)
+ private Order.OrderStatus previousStatus;
+
+ /** Who or what triggered this event (e.g. "inventory-service", "payment-service") */
+ @Column(length = 100)
+ private String triggeredBy;
+
+ /** Optional reason/detail (e.g. failure reason, tracking number) */
+ @Column(length = 500)
+ private String details;
+
+ /** When this event occurred — immutable, set at insert time */
+ @Column(nullable = false, updatable = false)
+ private Instant occurredAt;
+
+ @PrePersist
+ protected void onCreate() {
+ this.occurredAt = Instant.now();
+ }
+
+ public enum EventType {
+ ORDER_CREATED,
+ INVENTORY_CHECKING,
+ INVENTORY_RESERVED,
+ INVENTORY_RESERVATION_FAILED,
+ PAYMENT_PROCESSING,
+ PAYMENT_COMPLETED,
+ PAYMENT_FAILED,
+ SHIPPING_PROCESSING,
+ ORDER_SHIPPED,
+ ORDER_COMPLETED,
+ ORDER_CANCELLED,
+ ORDER_FAILED
+ }
+}
diff --git a/order-service/src/main/java/com/hacisimsek/order/eventsourcing/OrderEventRepository.java b/order-service/src/main/java/com/hacisimsek/order/eventsourcing/OrderEventRepository.java
new file mode 100644
index 0000000..3400279
--- /dev/null
+++ b/order-service/src/main/java/com/hacisimsek/order/eventsourcing/OrderEventRepository.java
@@ -0,0 +1,21 @@
+package com.hacisimsek.order.eventsourcing;
+
+import org.springframework.data.jpa.repository.JpaRepository;
+
+import java.util.List;
+import java.util.UUID;
+
+/**
+ * Repository for the append-only order_events table.
+ *
+ * Events are always inserted, never updated or deleted.
+ * Queries return events in chronological order.
+ */
+public interface OrderEventRepository extends JpaRepository {
+
+ /** Full audit trail for one order — ordered oldest first */
+ List findByOrderIdOrderByOccurredAtAsc(UUID orderId);
+
+ /** All events tied to a saga correlation ID — cross-service audit */
+ List findByCorrelationIdOrderByOccurredAtAsc(UUID correlationId);
+}
diff --git a/order-service/src/main/java/com/hacisimsek/order/eventsourcing/OrderEventService.java b/order-service/src/main/java/com/hacisimsek/order/eventsourcing/OrderEventService.java
new file mode 100644
index 0000000..4817726
--- /dev/null
+++ b/order-service/src/main/java/com/hacisimsek/order/eventsourcing/OrderEventService.java
@@ -0,0 +1,78 @@
+package com.hacisimsek.order.eventsourcing;
+
+import com.hacisimsek.order.model.Order;
+import lombok.RequiredArgsConstructor;
+import lombok.extern.slf4j.Slf4j;
+import org.springframework.stereotype.Service;
+import org.springframework.transaction.annotation.Transactional;
+
+import java.util.List;
+import java.util.UUID;
+
+/**
+ * Service for appending and querying order events (Event Sourcing log).
+ *
+ * Every call to {@link #append} inserts one immutable row into order_events.
+ * The event log is the source of truth for what happened to an order and when.
+ */
+@Service
+@RequiredArgsConstructor
+@Slf4j
+public class OrderEventService {
+
+ private final OrderEventRepository orderEventRepository;
+
+ /**
+ * Append a new event to the order's event log.
+ * Called inside the same transaction as the Order status update.
+ */
+ @Transactional
+ public OrderEvent append(UUID orderId,
+ UUID correlationId,
+ OrderEvent.EventType eventType,
+ Order.OrderStatus previousStatus,
+ Order.OrderStatus newStatus,
+ String triggeredBy,
+ String details) {
+ OrderEvent event = OrderEvent.builder()
+ .orderId(orderId)
+ .correlationId(correlationId)
+ .eventType(eventType)
+ .previousStatus(previousStatus)
+ .newStatus(newStatus)
+ .triggeredBy(triggeredBy)
+ .details(details)
+ .build();
+
+ OrderEvent saved = orderEventRepository.save(event);
+ log.debug("[EventStore] Appended {} for order {} ({} → {})",
+ eventType, orderId, previousStatus, newStatus);
+ return saved;
+ }
+
+ /** Retrieve the full immutable event log for an order */
+ @Transactional(readOnly = true)
+ public List getHistory(UUID orderId) {
+ return orderEventRepository.findByOrderIdOrderByOccurredAtAsc(orderId);
+ }
+
+ /** Retrieve all events tied to a saga correlation ID */
+ @Transactional(readOnly = true)
+ public List getByCorrelationId(UUID correlationId) {
+ return orderEventRepository.findByCorrelationIdOrderByOccurredAtAsc(correlationId);
+ }
+
+ /**
+ * Rebuild current order status by replaying the event log.
+ * Useful for auditing or reconciling against the Order table.
+ */
+ @Transactional(readOnly = true)
+ public Order.OrderStatus rebuildCurrentStatus(UUID orderId) {
+ List events = getHistory(orderId);
+ if (events.isEmpty()) {
+ throw new RuntimeException("No events found for order: " + orderId);
+ }
+ // The last event's newStatus is the current state
+ return events.get(events.size() - 1).getNewStatus();
+ }
+}
diff --git a/order-service/src/main/java/com/hacisimsek/order/outbox/OutboxEvent.java b/order-service/src/main/java/com/hacisimsek/order/outbox/OutboxEvent.java
new file mode 100644
index 0000000..f00c55b
--- /dev/null
+++ b/order-service/src/main/java/com/hacisimsek/order/outbox/OutboxEvent.java
@@ -0,0 +1,79 @@
+package com.hacisimsek.order.outbox;
+
+import jakarta.persistence.*;
+import lombok.*;
+
+import java.time.Instant;
+import java.util.UUID;
+
+/**
+ * Outbox table entry — written in the same DB transaction as the Order.
+ *
+ * The OutboxPublisher reads unpublished rows on a fixed schedule and
+ * publishes them to Kafka. On success the row is marked PUBLISHED.
+ * On failure it stays PENDING and is retried on the next schedule tick.
+ *
+ * This guarantees at-least-once Kafka delivery even if the service
+ * crashes between saving the order and sending to Kafka.
+ */
+@Entity
+@Table(name = "outbox_events",
+ indexes = {
+ @Index(name = "idx_outbox_status_created", columnList = "status, createdAt"),
+ @Index(name = "idx_outbox_aggregate", columnList = "aggregateId")
+ })
+@Data
+@Builder
+@NoArgsConstructor
+@AllArgsConstructor
+public class OutboxEvent {
+
+ @Id
+ @GeneratedValue(strategy = GenerationType.UUID)
+ private UUID id;
+
+ /** The Kafka topic this event should be published to */
+ @Column(nullable = false)
+ private String topic;
+
+ /** The entity this event belongs to — used as the Kafka message key */
+ @Column(nullable = false)
+ private UUID aggregateId;
+
+ /** Fully-qualified Java class name of the payload (e.g. OrderCreatedEvent) */
+ @Column(nullable = false)
+ private String eventType;
+
+ /** JSON-serialized event payload */
+ @Column(nullable = false, columnDefinition = "TEXT")
+ private String payload;
+
+ @Enumerated(EnumType.STRING)
+ @Column(nullable = false)
+ @Builder.Default
+ private Status status = Status.PENDING;
+
+ @Column(nullable = false, updatable = false)
+ private Instant createdAt;
+
+ private Instant publishedAt;
+
+ /** Number of failed publish attempts — for observability */
+ @Builder.Default
+ private int retryCount = 0;
+
+ /** Last error message from a failed publish attempt */
+ @Column(length = 1000)
+ private String lastError;
+
+ @PrePersist
+ protected void onCreate() {
+ this.createdAt = Instant.now();
+ }
+
+ public enum Status {
+ PENDING,
+ PUBLISHED,
+ FAILED // after max retries exceeded (currently 5)
+ }
+}
diff --git a/order-service/src/main/java/com/hacisimsek/order/outbox/OutboxEventRepository.java b/order-service/src/main/java/com/hacisimsek/order/outbox/OutboxEventRepository.java
new file mode 100644
index 0000000..21e6a6d
--- /dev/null
+++ b/order-service/src/main/java/com/hacisimsek/order/outbox/OutboxEventRepository.java
@@ -0,0 +1,21 @@
+package com.hacisimsek.order.outbox;
+
+import org.springframework.data.jpa.repository.JpaRepository;
+import org.springframework.data.jpa.repository.Modifying;
+import org.springframework.data.jpa.repository.Query;
+import org.springframework.data.repository.query.Param;
+
+import java.time.Instant;
+import java.util.List;
+import java.util.UUID;
+
+public interface OutboxEventRepository extends JpaRepository {
+
+ /** Fetch all pending events ordered oldest-first (for FIFO delivery) */
+ List findByStatusOrderByCreatedAtAsc(OutboxEvent.Status status);
+
+ /** Clean up published events older than the given cutoff to keep the table small */
+ @Modifying
+ @Query("DELETE FROM OutboxEvent e WHERE e.status = 'PUBLISHED' AND e.publishedAt < :cutoff")
+ void deletePublishedBefore(@Param("cutoff") Instant cutoff);
+}
diff --git a/order-service/src/main/java/com/hacisimsek/order/outbox/OutboxPublisher.java b/order-service/src/main/java/com/hacisimsek/order/outbox/OutboxPublisher.java
new file mode 100644
index 0000000..e4ee0b7
--- /dev/null
+++ b/order-service/src/main/java/com/hacisimsek/order/outbox/OutboxPublisher.java
@@ -0,0 +1,103 @@
+package com.hacisimsek.order.outbox;
+
+import com.fasterxml.jackson.databind.ObjectMapper;
+import lombok.RequiredArgsConstructor;
+import lombok.extern.slf4j.Slf4j;
+import org.springframework.kafka.core.KafkaTemplate;
+import org.springframework.kafka.support.SendResult;
+import org.springframework.scheduling.annotation.Scheduled;
+import org.springframework.stereotype.Component;
+import org.springframework.transaction.annotation.Transactional;
+
+import java.time.Instant;
+import java.time.temporal.ChronoUnit;
+import java.util.List;
+
+/**
+ * Polls the outbox table every 5 seconds and publishes pending events to Kafka.
+ *
+ * Flow:
+ * 1. Read all PENDING rows (oldest first)
+ * 2. Deserialize payload back to the original event object
+ * 3. Send to the target Kafka topic synchronously (get() with 10s timeout)
+ * 4. On success → mark PUBLISHED
+ * 5. On failure → increment retryCount; after 5 failures mark FAILED
+ *
+ * A nightly cleanup job removes PUBLISHED rows older than 7 days.
+ *
+ * Why synchronous send? Because we must know whether Kafka accepted the
+ * message before marking it published. An async callback arriving after
+ * a crash would leave the row in PENDING (safe — it will be retried).
+ */
+@Component
+@RequiredArgsConstructor
+@Slf4j
+public class OutboxPublisher {
+
+ private static final int MAX_RETRIES = 5;
+
+ private final OutboxEventRepository outboxEventRepository;
+ private final KafkaTemplate kafkaTemplate;
+ private final ObjectMapper objectMapper;
+
+ @Scheduled(fixedDelay = 5000) // runs 5s after the previous run completes
+ @Transactional
+ public void publishPendingEvents() {
+ List pending =
+ outboxEventRepository.findByStatusOrderByCreatedAtAsc(OutboxEvent.Status.PENDING);
+
+ if (pending.isEmpty()) return;
+
+ log.debug("[Outbox] Processing {} pending event(s)", pending.size());
+
+ for (OutboxEvent event : pending) {
+ try {
+ // Deserialize the stored JSON payload back to the original event class
+ Class> eventClass = Class.forName(event.getEventType());
+ Object eventPayload = objectMapper.readValue(event.getPayload(), eventClass);
+
+ // Synchronous send — waits for broker ACK (or throws on timeout/error)
+ SendResult result = kafkaTemplate
+ .send(event.getTopic(), event.getAggregateId().toString(), eventPayload)
+ .get();
+
+ // Mark published
+ event.setStatus(OutboxEvent.Status.PUBLISHED);
+ event.setPublishedAt(Instant.now());
+ outboxEventRepository.save(event);
+
+ log.info("[Outbox] Published {} → topic={} partition={} offset={}",
+ event.getEventType(),
+ event.getTopic(),
+ result.getRecordMetadata().partition(),
+ result.getRecordMetadata().offset());
+
+ } catch (Exception ex) {
+ int retries = event.getRetryCount() + 1;
+ event.setRetryCount(retries);
+ event.setLastError(ex.getMessage() != null
+ ? ex.getMessage().substring(0, Math.min(ex.getMessage().length(), 1000))
+ : "unknown");
+
+ if (retries >= MAX_RETRIES) {
+ event.setStatus(OutboxEvent.Status.FAILED);
+ log.error("[Outbox] Event {} FAILED after {} retries. Manual intervention required. Error: {}",
+ event.getId(), retries, ex.getMessage());
+ } else {
+ log.warn("[Outbox] Publish attempt {}/{} failed for event {} ({}): {}",
+ retries, MAX_RETRIES, event.getId(), event.getEventType(), ex.getMessage());
+ }
+ outboxEventRepository.save(event);
+ }
+ }
+ }
+
+ /** Runs nightly to clean up old published events and keep the table small */
+ @Scheduled(cron = "0 0 2 * * *") // 02:00 every day
+ @Transactional
+ public void purgePublishedEvents() {
+ Instant cutoff = Instant.now().minus(7, ChronoUnit.DAYS);
+ outboxEventRepository.deletePublishedBefore(cutoff);
+ log.info("[Outbox] Purged published events older than 7 days");
+ }
+}
diff --git a/order-service/src/main/java/com/hacisimsek/order/saga/OrderSagaHandler.java b/order-service/src/main/java/com/hacisimsek/order/saga/OrderSagaHandler.java
index e83aedd..64e0405 100644
--- a/order-service/src/main/java/com/hacisimsek/order/saga/OrderSagaHandler.java
+++ b/order-service/src/main/java/com/hacisimsek/order/saga/OrderSagaHandler.java
@@ -1,113 +1,78 @@
package com.hacisimsek.order.saga;
-import java.util.Map;
-
-import org.springframework.kafka.annotation.KafkaListener;
-import org.springframework.stereotype.Component;
-
import com.hacisimsek.common.event.inventory.InventoryReservationFailedEvent;
import com.hacisimsek.common.event.inventory.InventoryReservedEvent;
import com.hacisimsek.common.event.payment.PaymentFailedEvent;
import com.hacisimsek.common.event.payment.PaymentProcessedEvent;
import com.hacisimsek.common.event.shipping.ShipmentFailedEvent;
import com.hacisimsek.common.event.shipping.ShipmentProcessedEvent;
-import com.hacisimsek.common.logging.LogPublisher;
-import com.hacisimsek.order.model.Order;
-import com.hacisimsek.order.service.OrderService;
-
+import com.hacisimsek.order.saga.orchestrator.OrderSagaOrchestrator;
import lombok.RequiredArgsConstructor;
import lombok.extern.slf4j.Slf4j;
+import org.springframework.kafka.annotation.KafkaListener;
+import org.springframework.stereotype.Component;
+/**
+ * Kafka listeners for the Order Saga.
+ *
+ * This class is now a thin adapter layer — it receives Kafka events and
+ * immediately delegates to the {@link OrderSagaOrchestrator} which owns
+ * all the state machine logic and compensation decisions.
+ *
+ * All business logic that was previously inline here has moved to the orchestrator.
+ */
@Component
@RequiredArgsConstructor
@Slf4j
public class OrderSagaHandler {
- private static final String SERVICE_NAME = "order-service";
+ private final OrderSagaOrchestrator orchestrator;
- private final OrderService orderService;
- private final LogPublisher logPublisher;
+ // ── Inventory events ──────────────────────────────────────────────────────
@KafkaListener(topics = "inventory-events", groupId = "order-service-group",
containerFactory = "kafkaListenerContainerFactory")
public void handleInventoryEvents(Object event) {
- log.info("Received inventory event: {}", event.getClass().getSimpleName());
+ log.debug("Received inventory event: {}", event.getClass().getSimpleName());
- if (event instanceof InventoryReservedEvent reservedEvent) {
- orderService.updateOrderStatus(reservedEvent.getOrderId(), Order.OrderStatus.INVENTORY_RESERVED);
- log.info("Inventory reserved for order: {}", reservedEvent.getOrderId());
- logPublisher.info(SERVICE_NAME,
- reservedEvent.getCorrelationId() != null ? reservedEvent.getCorrelationId().toString() : null,
- "Inventory reserved for order: " + reservedEvent.getOrderId(),
- Map.of("orderId", reservedEvent.getOrderId().toString(), "status", "INVENTORY_RESERVED"));
+ if (event instanceof InventoryReservedEvent e) {
+ orchestrator.onInventoryReserved(e.getOrderId(), e.getCorrelationId());
- } else if (event instanceof InventoryReservationFailedEvent failedEvent) {
- orderService.updateOrderStatus(failedEvent.getOrderId(), Order.OrderStatus.CANCELLED);
- log.error("Inventory reservation failed for order: {}, reason: {}",
- failedEvent.getOrderId(), failedEvent.getReason());
- logPublisher.error(SERVICE_NAME,
- failedEvent.getCorrelationId() != null ? failedEvent.getCorrelationId().toString() : null,
- "Inventory reservation failed — order cancelled: " + failedEvent.getOrderId(),
- Map.of("orderId", failedEvent.getOrderId().toString(),
- "reason", failedEvent.getReason() != null ? failedEvent.getReason() : "unknown",
- "status", "CANCELLED"));
+ } else if (event instanceof InventoryReservationFailedEvent e) {
+ orchestrator.onInventoryFailed(e.getOrderId(), e.getCorrelationId(),
+ e.getReason() != null ? e.getReason() : "unknown");
}
}
+ // ── Payment events ────────────────────────────────────────────────────────
+
@KafkaListener(topics = "payment-events", groupId = "order-service-group",
containerFactory = "kafkaListenerContainerFactory")
public void handlePaymentEvents(Object event) {
- log.info("Received payment event: {}", event.getClass().getSimpleName());
+ log.debug("Received payment event: {}", event.getClass().getSimpleName());
- if (event instanceof PaymentProcessedEvent processedEvent) {
- orderService.updateOrderStatus(processedEvent.getOrderId(), Order.OrderStatus.PAYMENT_COMPLETED);
- log.info("Payment processed for order: {}", processedEvent.getOrderId());
- logPublisher.info(SERVICE_NAME,
- processedEvent.getCorrelationId() != null ? processedEvent.getCorrelationId().toString() : null,
- "Payment completed for order: " + processedEvent.getOrderId(),
- Map.of("orderId", processedEvent.getOrderId().toString(),
- "paymentId", processedEvent.getPaymentId() != null ? processedEvent.getPaymentId().toString() : "unknown",
- "status", "PAYMENT_COMPLETED"));
+ if (event instanceof PaymentProcessedEvent e) {
+ orchestrator.onPaymentCompleted(e.getOrderId(), e.getCorrelationId(), e.getPaymentId());
- } else if (event instanceof PaymentFailedEvent failedEvent) {
- orderService.updateOrderStatus(failedEvent.getOrderId(), Order.OrderStatus.FAILED);
- log.error("Payment failed for order: {}, reason: {}",
- failedEvent.getOrderId(), failedEvent.getReason());
- logPublisher.error(SERVICE_NAME,
- failedEvent.getCorrelationId() != null ? failedEvent.getCorrelationId().toString() : null,
- "Payment failed — order marked FAILED: " + failedEvent.getOrderId(),
- Map.of("orderId", failedEvent.getOrderId().toString(),
- "reason", failedEvent.getReason() != null ? failedEvent.getReason() : "unknown",
- "status", "FAILED"));
+ } else if (event instanceof PaymentFailedEvent e) {
+ orchestrator.onPaymentFailed(e.getOrderId(), e.getCorrelationId(),
+ e.getReason() != null ? e.getReason() : "unknown");
}
}
+ // ── Shipping events ───────────────────────────────────────────────────────
+
@KafkaListener(topics = "shipping-events", groupId = "order-service-group",
containerFactory = "kafkaListenerContainerFactory")
public void handleShippingEvents(Object event) {
- log.info("Received shipping event: {}", event.getClass().getSimpleName());
+ log.debug("Received shipping event: {}", event.getClass().getSimpleName());
- if (event instanceof ShipmentProcessedEvent processedEvent) {
- orderService.updateOrderStatus(processedEvent.getOrderId(), Order.OrderStatus.SHIPPED);
- log.info("Order shipped: {}, tracking number: {}",
- processedEvent.getOrderId(), processedEvent.getTrackingNumber());
- logPublisher.info(SERVICE_NAME,
- processedEvent.getCorrelationId() != null ? processedEvent.getCorrelationId().toString() : null,
- "Order shipped: " + processedEvent.getOrderId() + " | tracking: " + processedEvent.getTrackingNumber(),
- Map.of("orderId", processedEvent.getOrderId().toString(),
- "trackingNumber", processedEvent.getTrackingNumber() != null ? processedEvent.getTrackingNumber() : "unknown",
- "status", "SHIPPED"));
+ if (event instanceof ShipmentProcessedEvent e) {
+ orchestrator.onShipmentCreated(e.getOrderId(), e.getCorrelationId(), e.getTrackingNumber());
- } else if (event instanceof ShipmentFailedEvent failedEvent) {
- orderService.updateOrderStatus(failedEvent.getOrderId(), Order.OrderStatus.FAILED);
- log.error("Shipping failed for order: {}, reason: {}",
- failedEvent.getOrderId(), failedEvent.getReason());
- logPublisher.error(SERVICE_NAME,
- failedEvent.getCorrelationId() != null ? failedEvent.getCorrelationId().toString() : null,
- "Shipment failed — order marked FAILED: " + failedEvent.getOrderId(),
- Map.of("orderId", failedEvent.getOrderId().toString(),
- "reason", failedEvent.getReason() != null ? failedEvent.getReason() : "unknown",
- "status", "FAILED"));
+ } else if (event instanceof ShipmentFailedEvent e) {
+ orchestrator.onShipmentFailed(e.getOrderId(), e.getCorrelationId(),
+ e.getReason() != null ? e.getReason() : "unknown");
}
}
-}
\ No newline at end of file
+}
diff --git a/order-service/src/main/java/com/hacisimsek/order/saga/orchestrator/OrderSagaOrchestrator.java b/order-service/src/main/java/com/hacisimsek/order/saga/orchestrator/OrderSagaOrchestrator.java
new file mode 100644
index 0000000..52d5c5d
--- /dev/null
+++ b/order-service/src/main/java/com/hacisimsek/order/saga/orchestrator/OrderSagaOrchestrator.java
@@ -0,0 +1,160 @@
+package com.hacisimsek.order.saga.orchestrator;
+
+import com.hacisimsek.common.logging.LogPublisher;
+import com.hacisimsek.order.model.Order;
+import com.hacisimsek.order.service.OrderService;
+import lombok.RequiredArgsConstructor;
+import lombok.extern.slf4j.Slf4j;
+import org.springframework.messaging.Message;
+import org.springframework.messaging.support.MessageBuilder;
+import org.springframework.statemachine.StateMachine;
+import org.springframework.statemachine.config.StateMachineFactory;
+import org.springframework.statemachine.support.DefaultStateMachineContext;
+import org.springframework.stereotype.Service;
+import reactor.core.publisher.Mono;
+
+import java.util.Map;
+import java.util.UUID;
+
+/**
+ * Central Saga Orchestrator — drives the Order state machine.
+ *
+ * Instead of each service reacting to events independently (choreography),
+ * this orchestrator:
+ * 1. Maintains the authoritative state machine per order
+ * 2. Receives all saga outcomes (inventory/payment/shipment results)
+ * 3. Decides the next step and updates order status
+ * 4. Handles compensation centrally (e.g. release inventory on payment failure)
+ *
+ * The Kafka listeners in OrderSagaHandler delegate to this orchestrator.
+ * This keeps all saga logic in one place instead of scattered across handlers.
+ *
+ * State machine instances are stateless in this implementation — the current
+ * state is always loaded from the database (Order.status) before processing
+ * each event, which makes it crash-safe and idempotent.
+ */
+@Service
+@RequiredArgsConstructor
+@Slf4j
+public class OrderSagaOrchestrator {
+
+ private static final String SERVICE_NAME = "order-service";
+ private static final String ORDER_ID_KEY = "orderId";
+ private static final String CORRELATION_ID_KEY = "correlationId";
+
+ private final StateMachineFactory stateMachineFactory;
+ private final OrderService orderService;
+ private final LogPublisher logPublisher;
+
+ // ── Public API called by Kafka listeners ──────────────────────────────────
+
+ public void onInventoryReserved(UUID orderId, UUID correlationId) {
+ processEvent(orderId, correlationId, SagaEvent.INVENTORY_RESERVED,
+ Order.OrderStatus.INVENTORY_RESERVED,
+ "inventory-service", "Inventory reserved — initiating payment");
+ }
+
+ public void onInventoryFailed(UUID orderId, UUID correlationId, String reason) {
+ processEvent(orderId, correlationId, SagaEvent.INVENTORY_FAILED,
+ Order.OrderStatus.CANCELLED,
+ "inventory-service", "Inventory reservation failed: " + reason);
+ log.warn("[Orchestrator] Order {} CANCELLED — inventory failed: {}", orderId, reason);
+ }
+
+ public void onPaymentCompleted(UUID orderId, UUID correlationId, UUID paymentId) {
+ processEvent(orderId, correlationId, SagaEvent.PAYMENT_COMPLETED,
+ Order.OrderStatus.PAYMENT_COMPLETED,
+ "payment-service", "Payment completed, paymentId=" + paymentId);
+ }
+
+ public void onPaymentFailed(UUID orderId, UUID correlationId, String reason) {
+ processEvent(orderId, correlationId, SagaEvent.PAYMENT_FAILED,
+ Order.OrderStatus.FAILED,
+ "payment-service", "Payment failed: " + reason);
+ // Compensation is handled by inventory-service which listens to payment-events
+ // (already implemented in InventorySagaHandler)
+ log.warn("[Orchestrator] Order {} FAILED — payment failed: {}", orderId, reason);
+ }
+
+ public void onShipmentCreated(UUID orderId, UUID correlationId, String trackingNumber) {
+ processEvent(orderId, correlationId, SagaEvent.SHIPMENT_CREATED,
+ Order.OrderStatus.SHIPPED,
+ "shipping-service", "Shipment created, tracking=" + trackingNumber);
+ }
+
+ public void onShipmentFailed(UUID orderId, UUID correlationId, String reason) {
+ processEvent(orderId, correlationId, SagaEvent.SHIPMENT_FAILED,
+ Order.OrderStatus.FAILED,
+ "shipping-service", "Shipment failed: " + reason);
+ log.warn("[Orchestrator] Order {} FAILED — shipment failed: {}", orderId, reason);
+ }
+
+ // ── Core state machine processing ─────────────────────────────────────────
+
+ private void processEvent(UUID orderId,
+ UUID correlationId,
+ SagaEvent event,
+ Order.OrderStatus targetStatus,
+ String triggeredBy,
+ String details) {
+ try {
+ // Load current state from DB (crash-safe: state machine is rebuilt each time)
+ Order.OrderStatus currentStatus = orderService.getOrderById(orderId).getStatus();
+
+ // Build a state machine pre-loaded at the current state
+ StateMachine sm = buildStateMachine(orderId, currentStatus);
+
+ // Send the event
+ Message message = MessageBuilder.withPayload(event)
+ .setHeader(ORDER_ID_KEY, orderId.toString())
+ .setHeader(CORRELATION_ID_KEY, correlationId != null ? correlationId.toString() : "")
+ .build();
+
+ sm.sendEvent(Mono.just(message)).subscribe();
+
+ Order.OrderStatus newState = sm.getState().getId();
+
+ // Persist the new state
+ orderService.updateOrderStatus(orderId, newState);
+
+ log.info("[Orchestrator] Order {} | event={} | {} → {}",
+ orderId, event, currentStatus, newState);
+
+ logPublisher.info(SERVICE_NAME,
+ correlationId != null ? correlationId.toString() : null,
+ "[Orchestrator] " + details,
+ Map.of("orderId", orderId.toString(),
+ "event", event.name(),
+ "from", currentStatus.name(),
+ "to", newState.name(),
+ "triggeredBy", triggeredBy));
+
+ } catch (Exception e) {
+ log.error("[Orchestrator] Failed to process event {} for order {}: {}",
+ event, orderId, e.getMessage());
+ }
+ }
+
+ /**
+ * Build a StateMachine instance pre-restored to the given state.
+ * Using the factory (not a singleton) ensures each order gets isolated state.
+ */
+ private StateMachine buildStateMachine(
+ UUID orderId, Order.OrderStatus currentState) throws Exception {
+
+ StateMachine sm =
+ stateMachineFactory.getStateMachine(orderId.toString());
+
+ sm.stopReactively().block();
+
+ sm.getStateMachineAccessor()
+ .doWithAllRegions(accessor ->
+ accessor.resetStateMachineReactively(
+ new DefaultStateMachineContext<>(currentState, null, null, null)
+ ).block()
+ );
+
+ sm.startReactively().block();
+ return sm;
+ }
+}
diff --git a/order-service/src/main/java/com/hacisimsek/order/saga/orchestrator/OrderSagaStateMachineConfig.java b/order-service/src/main/java/com/hacisimsek/order/saga/orchestrator/OrderSagaStateMachineConfig.java
new file mode 100644
index 0000000..f6bcc54
--- /dev/null
+++ b/order-service/src/main/java/com/hacisimsek/order/saga/orchestrator/OrderSagaStateMachineConfig.java
@@ -0,0 +1,124 @@
+package com.hacisimsek.order.saga.orchestrator;
+
+import com.hacisimsek.order.model.Order;
+import lombok.extern.slf4j.Slf4j;
+import org.springframework.context.annotation.Configuration;
+import org.springframework.statemachine.config.EnableStateMachineFactory;
+import org.springframework.statemachine.config.StateMachineConfigurerAdapter;
+import org.springframework.statemachine.config.builders.StateMachineStateConfigurer;
+import org.springframework.statemachine.config.builders.StateMachineTransitionConfigurer;
+
+import java.util.EnumSet;
+
+/**
+ * Spring State Machine configuration for the Order Saga Orchestrator.
+ *
+ * This replaces the scattered choreography listeners with a single, explicit
+ * state machine that drives the entire order lifecycle from one place.
+ *
+ * States map 1:1 with Order.OrderStatus.
+ * Events (SagaEvent) represent outcomes from downstream services.
+ *
+ * Transitions:
+ *
+ * PENDING ──[ORDER_PLACED]──► INVENTORY_CHECKING
+ * INVENTORY_CHECKING ──[INVENTORY_RESERVED]──► PAYMENT_PROCESSING
+ * INVENTORY_CHECKING ──[INVENTORY_FAILED]──► CANCELLED
+ * PAYMENT_PROCESSING ──[PAYMENT_COMPLETED]──► SHIPPING_PROCESSING
+ * PAYMENT_PROCESSING ──[PAYMENT_FAILED]──► FAILED (+ compensate inventory)
+ * SHIPPING_PROCESSING ──[SHIPMENT_CREATED]──► SHIPPED
+ * SHIPPING_PROCESSING ──[SHIPMENT_FAILED]──► FAILED (+ compensate inventory)
+ * SHIPPED ──[DELIVERY_CONFIRMED]──► COMPLETED
+ *
+ * The factory produces one StateMachine per order (keyed by orderId).
+ */
+@Configuration
+@EnableStateMachineFactory
+@Slf4j
+public class OrderSagaStateMachineConfig
+ extends StateMachineConfigurerAdapter {
+
+ @Override
+ public void configure(StateMachineStateConfigurer states)
+ throws Exception {
+ states
+ .withStates()
+ .initial(Order.OrderStatus.PENDING)
+ .states(EnumSet.allOf(Order.OrderStatus.class))
+ .end(Order.OrderStatus.COMPLETED)
+ .end(Order.OrderStatus.CANCELLED)
+ .end(Order.OrderStatus.FAILED);
+ }
+
+ @Override
+ public void configure(StateMachineTransitionConfigurer transitions)
+ throws Exception {
+ transitions
+ // Order placed → start inventory check
+ .withExternal()
+ .source(Order.OrderStatus.PENDING)
+ .target(Order.OrderStatus.INVENTORY_CHECKING)
+ .event(SagaEvent.ORDER_PLACED)
+ .and()
+
+ // Inventory reserved → initiate payment
+ .withExternal()
+ .source(Order.OrderStatus.INVENTORY_CHECKING)
+ .target(Order.OrderStatus.INVENTORY_RESERVED)
+ .event(SagaEvent.INVENTORY_RESERVED)
+ .and()
+
+ .withExternal()
+ .source(Order.OrderStatus.INVENTORY_RESERVED)
+ .target(Order.OrderStatus.PAYMENT_PROCESSING)
+ .event(SagaEvent.INVENTORY_RESERVED)
+ .and()
+
+ // Inventory failed → cancel order (no compensation needed)
+ .withExternal()
+ .source(Order.OrderStatus.INVENTORY_CHECKING)
+ .target(Order.OrderStatus.CANCELLED)
+ .event(SagaEvent.INVENTORY_FAILED)
+ .and()
+
+ // Payment completed → initiate shipping
+ .withExternal()
+ .source(Order.OrderStatus.PAYMENT_PROCESSING)
+ .target(Order.OrderStatus.PAYMENT_COMPLETED)
+ .event(SagaEvent.PAYMENT_COMPLETED)
+ .and()
+
+ .withExternal()
+ .source(Order.OrderStatus.PAYMENT_COMPLETED)
+ .target(Order.OrderStatus.SHIPPING_PROCESSING)
+ .event(SagaEvent.PAYMENT_COMPLETED)
+ .and()
+
+ // Payment failed → fail order + compensate inventory
+ .withExternal()
+ .source(Order.OrderStatus.PAYMENT_PROCESSING)
+ .target(Order.OrderStatus.FAILED)
+ .event(SagaEvent.PAYMENT_FAILED)
+ .and()
+
+ // Shipment created → order shipped
+ .withExternal()
+ .source(Order.OrderStatus.SHIPPING_PROCESSING)
+ .target(Order.OrderStatus.SHIPPED)
+ .event(SagaEvent.SHIPMENT_CREATED)
+ .and()
+
+ // Shipment failed → fail order + compensate inventory
+ .withExternal()
+ .source(Order.OrderStatus.SHIPPING_PROCESSING)
+ .target(Order.OrderStatus.FAILED)
+ .event(SagaEvent.SHIPMENT_FAILED)
+ .and()
+
+ // Delivery confirmed → complete
+ .withExternal()
+ .source(Order.OrderStatus.SHIPPED)
+ .target(Order.OrderStatus.COMPLETED)
+ .event(SagaEvent.DELIVERY_CONFIRMED);
+ }
+}
diff --git a/order-service/src/main/java/com/hacisimsek/order/saga/orchestrator/SagaEvent.java b/order-service/src/main/java/com/hacisimsek/order/saga/orchestrator/SagaEvent.java
new file mode 100644
index 0000000..8d526c6
--- /dev/null
+++ b/order-service/src/main/java/com/hacisimsek/order/saga/orchestrator/SagaEvent.java
@@ -0,0 +1,16 @@
+package com.hacisimsek.order.saga.orchestrator;
+
+/**
+ * Events that drive the Order Saga state machine.
+ * Published by downstream services and received via Kafka listeners.
+ */
+public enum SagaEvent {
+ ORDER_PLACED,
+ INVENTORY_RESERVED,
+ INVENTORY_FAILED,
+ PAYMENT_COMPLETED,
+ PAYMENT_FAILED,
+ SHIPMENT_CREATED,
+ SHIPMENT_FAILED,
+ DELIVERY_CONFIRMED
+}
diff --git a/order-service/src/main/java/com/hacisimsek/order/service/impl/OrderServiceImpl.java b/order-service/src/main/java/com/hacisimsek/order/service/impl/OrderServiceImpl.java
index a57fab8..6493bb1 100644
--- a/order-service/src/main/java/com/hacisimsek/order/service/impl/OrderServiceImpl.java
+++ b/order-service/src/main/java/com/hacisimsek/order/service/impl/OrderServiceImpl.java
@@ -1,5 +1,7 @@
package com.hacisimsek.order.service.impl;
+import com.fasterxml.jackson.core.JsonProcessingException;
+import com.fasterxml.jackson.databind.ObjectMapper;
import com.hacisimsek.common.dto.OrderItemDto;
import com.hacisimsek.common.event.order.OrderCreatedEvent;
import com.hacisimsek.order.dto.OrderItemResponse;
@@ -7,12 +9,13 @@
import com.hacisimsek.order.dto.OrderResponse;
import com.hacisimsek.order.model.Order;
import com.hacisimsek.order.model.OrderItem;
+import com.hacisimsek.order.outbox.OutboxEvent;
+import com.hacisimsek.order.outbox.OutboxEventRepository;
import com.hacisimsek.order.repository.OrderRepository;
import com.hacisimsek.order.service.OrderService;
import io.micrometer.core.instrument.Counter;
import io.micrometer.core.instrument.MeterRegistry;
import lombok.extern.slf4j.Slf4j;
-import org.springframework.kafka.core.KafkaTemplate;
import org.springframework.stereotype.Service;
import org.springframework.transaction.annotation.Transactional;
@@ -21,19 +24,32 @@
import java.util.UUID;
import java.util.stream.Collectors;
+import com.hacisimsek.order.eventsourcing.OrderEvent;
+import com.hacisimsek.order.eventsourcing.OrderEventService;
+import com.hacisimsek.order.sse.OrderStatusEmitter;
+
@Service
@Slf4j
public class OrderServiceImpl implements OrderService {
private final OrderRepository orderRepository;
- private final KafkaTemplate kafkaTemplate;
+ private final OutboxEventRepository outboxEventRepository;
+ private final ObjectMapper objectMapper;
+ private final OrderStatusEmitter orderStatusEmitter;
+ private final OrderEventService orderEventService;
private final Counter ordersCreatedCounter;
public OrderServiceImpl(OrderRepository orderRepository,
- KafkaTemplate kafkaTemplate,
+ OutboxEventRepository outboxEventRepository,
+ ObjectMapper objectMapper,
+ OrderStatusEmitter orderStatusEmitter,
+ OrderEventService orderEventService,
MeterRegistry meterRegistry) {
this.orderRepository = orderRepository;
- this.kafkaTemplate = kafkaTemplate;
+ this.outboxEventRepository = outboxEventRepository;
+ this.objectMapper = objectMapper;
+ this.orderStatusEmitter = orderStatusEmitter;
+ this.orderEventService = orderEventService;
this.ordersCreatedCounter = Counter.builder("zexxity.orders.created")
.description("Total number of orders successfully created")
.register(meterRegistry);
@@ -42,7 +58,7 @@ public OrderServiceImpl(OrderRepository orderRepository,
@Override
@Transactional
public OrderResponse createOrder(OrderRequest orderRequest) {
- // Convert order items from request
+ // ── 1. Build and save the Order ──────────────────────────────────────
List orderItems = orderRequest.getItems().stream()
.map(item -> OrderItem.builder()
.productId(item.getProductId())
@@ -52,12 +68,10 @@ public OrderResponse createOrder(OrderRequest orderRequest) {
.build())
.collect(Collectors.toList());
- // Calculate total amount
BigDecimal totalAmount = orderItems.stream()
.map(item -> item.getPrice().multiply(BigDecimal.valueOf(item.getQuantity())))
.reduce(BigDecimal.ZERO, BigDecimal::add);
- // Create and save order
Order order = Order.builder()
.customerId(orderRequest.getCustomerId())
.customerEmail(orderRequest.getCustomerEmail())
@@ -68,8 +82,9 @@ public OrderResponse createOrder(OrderRequest orderRequest) {
Order savedOrder = orderRepository.save(order);
- // Start the saga by sending OrderCreatedEvent
+ // ── 2. Build the Kafka event ──────────────────────────────────────────
UUID correlationId = UUID.randomUUID();
+
List itemDtos = savedOrder.getItems().stream()
.map(item -> new OrderItemDto(
item.getProductId(),
@@ -87,22 +102,51 @@ public OrderResponse createOrder(OrderRequest orderRequest) {
savedOrder.getTotalAmount()
);
- log.info("Sending OrderCreatedEvent for order {}", savedOrder.getId());
-
- // Increment Prometheus counter
+ // ── 3. Write to the Outbox in the SAME transaction ───────────────────
+ //
+ // By writing the OutboxEvent inside the same @Transactional method,
+ // both the Order row and the OutboxEvent row are committed atomically.
+ // If Kafka is unavailable, the OutboxPublisher scheduler will pick up
+ // and publish the pending row on the next tick (every 5 seconds).
+ // This eliminates the "dual-write" race condition in the original code.
+ try {
+ OutboxEvent outboxEntry = OutboxEvent.builder()
+ .topic("order-events")
+ .aggregateId(savedOrder.getId())
+ .eventType(OrderCreatedEvent.class.getName())
+ .payload(objectMapper.writeValueAsString(event))
+ .status(OutboxEvent.Status.PENDING)
+ .build();
+
+ outboxEventRepository.save(outboxEntry);
+ log.info("Order {} saved with outbox entry (correlationId={})",
+ savedOrder.getId(), correlationId);
+ } catch (JsonProcessingException ex) {
+ // This would be a programming error (unparseable event) — rethrow
+ throw new IllegalStateException("Failed to serialize OrderCreatedEvent for outbox", ex);
+ }
+
+ // ── 4. Update status and metrics ─────────────────────────────────────
ordersCreatedCounter.increment();
-
- // Update order status to indicate saga started
savedOrder.setStatus(Order.OrderStatus.INVENTORY_CHECKING);
orderRepository.save(savedOrder);
- // Publish event to Kafka
- log.info("BEFORE Kafka Send");
+ // Append ORDER_CREATED event to the immutable event log
+ orderEventService.append(
+ savedOrder.getId(), correlationId,
+ OrderEvent.EventType.ORDER_CREATED,
+ null, Order.OrderStatus.PENDING,
+ "order-service", "Order created with " + itemDtos.size() + " item(s)");
-// Publish event to Kafka
- kafkaTemplate.send("order-events", event);
+ // Append INVENTORY_CHECKING event
+ orderEventService.append(
+ savedOrder.getId(), correlationId,
+ OrderEvent.EventType.INVENTORY_CHECKING,
+ Order.OrderStatus.PENDING, Order.OrderStatus.INVENTORY_CHECKING,
+ "order-service", "Saga started — checking inventory");
- log.info("AFTER Kafka Send");
+ // Push initial status to any SSE subscriber
+ orderStatusEmitter.push(savedOrder.getId(), Order.OrderStatus.INVENTORY_CHECKING.name(), false);
return mapToOrderResponse(savedOrder);
}
@@ -133,9 +177,49 @@ public List getOrdersByCustomerId(UUID customerId) {
public void updateOrderStatus(UUID orderId, Order.OrderStatus status) {
Order order = orderRepository.findById(orderId)
.orElseThrow(() -> new RuntimeException("Order not found with id: " + orderId));
+
+ Order.OrderStatus previousStatus = order.getStatus();
order.setStatus(status);
orderRepository.save(order);
log.info("Updated order {} status to {}", orderId, status);
+
+ // Append transition event to the immutable event log
+ OrderEvent.EventType eventType = resolveEventType(status);
+ orderEventService.append(
+ orderId, null,
+ eventType,
+ previousStatus, status,
+ "saga", null);
+
+ // Push real-time status update via SSE
+ boolean terminal = isTerminalStatus(status);
+ orderStatusEmitter.push(orderId, status.name(), terminal);
+ }
+
+ private OrderEvent.EventType resolveEventType(Order.OrderStatus status) {
+ return switch (status) {
+ case INVENTORY_CHECKING -> OrderEvent.EventType.INVENTORY_CHECKING;
+ case INVENTORY_RESERVED -> OrderEvent.EventType.INVENTORY_RESERVED;
+ case PAYMENT_PROCESSING -> OrderEvent.EventType.PAYMENT_PROCESSING;
+ case PAYMENT_COMPLETED -> OrderEvent.EventType.PAYMENT_COMPLETED;
+ case SHIPPING_PROCESSING -> OrderEvent.EventType.SHIPPING_PROCESSING;
+ case SHIPPED -> OrderEvent.EventType.ORDER_SHIPPED;
+ case COMPLETED -> OrderEvent.EventType.ORDER_COMPLETED;
+ case CANCELLED -> OrderEvent.EventType.ORDER_CANCELLED;
+ case FAILED -> OrderEvent.EventType.ORDER_FAILED;
+ default -> OrderEvent.EventType.ORDER_CREATED;
+ };
+ }
+
+ /**
+ * Terminal statuses — the saga has reached a final state.
+ * After these, no further status changes will occur.
+ */
+ private boolean isTerminalStatus(Order.OrderStatus status) {
+ return status == Order.OrderStatus.SHIPPED
+ || status == Order.OrderStatus.COMPLETED
+ || status == Order.OrderStatus.FAILED
+ || status == Order.OrderStatus.CANCELLED;
}
private OrderResponse mapToOrderResponse(Order order) {
@@ -160,4 +244,4 @@ private OrderResponse mapToOrderResponse(Order order) {
.lastModifiedAt(order.getLastModifiedAt())
.build();
}
-}
\ No newline at end of file
+}
diff --git a/order-service/src/main/java/com/hacisimsek/order/sse/OrderStatusEmitter.java b/order-service/src/main/java/com/hacisimsek/order/sse/OrderStatusEmitter.java
new file mode 100644
index 0000000..f3f466e
--- /dev/null
+++ b/order-service/src/main/java/com/hacisimsek/order/sse/OrderStatusEmitter.java
@@ -0,0 +1,121 @@
+package com.hacisimsek.order.sse;
+
+import lombok.extern.slf4j.Slf4j;
+import org.springframework.stereotype.Component;
+import org.springframework.web.servlet.mvc.method.annotation.SseEmitter;
+
+import java.io.IOException;
+import java.time.Instant;
+import java.util.Map;
+import java.util.UUID;
+import java.util.concurrent.ConcurrentHashMap;
+
+/**
+ * Manages active SSE (Server-Sent Events) connections for order status updates.
+ *
+ * Each client subscribes to a specific orderId. When the order status changes
+ * (driven by saga events), the saga handler calls {@link #push} and the update
+ * is streamed instantly to any connected browser/client — no polling needed.
+ *
+ * Connection lifecycle:
+ * - Client opens GET /api/orders/{orderId}/status-stream
+ * - Server holds the connection open (SseEmitter with 5-min timeout)
+ * - On status change → push event to that orderId's emitter
+ * - On COMPLETED/CANCELLED/FAILED → push final event and complete the stream
+ * - On timeout or client disconnect → emitter is cleaned up automatically
+ *
+ * Thread safety: ConcurrentHashMap handles concurrent subscribe/push/cleanup.
+ */
+@Component
+@Slf4j
+public class OrderStatusEmitter {
+
+ // orderId → active SseEmitter for that order
+ private final Map emitters = new ConcurrentHashMap<>();
+
+ /** SSE connection timeout — 5 minutes. Client should reconnect if needed. */
+ private static final long TIMEOUT_MS = 5 * 60 * 1000L;
+
+ /**
+ * Register a new SSE connection for the given orderId.
+ * Returns the emitter to be written directly to the HTTP response.
+ */
+ public SseEmitter subscribe(UUID orderId) {
+ SseEmitter emitter = new SseEmitter(TIMEOUT_MS);
+
+ // Clean up on completion, timeout, or error
+ emitter.onCompletion(() -> {
+ emitters.remove(orderId);
+ log.debug("[SSE] Connection completed for order {}", orderId);
+ });
+ emitter.onTimeout(() -> {
+ emitters.remove(orderId);
+ log.debug("[SSE] Connection timed out for order {}", orderId);
+ });
+ emitter.onError(ex -> {
+ emitters.remove(orderId);
+ log.debug("[SSE] Connection error for order {}: {}", orderId, ex.getMessage());
+ });
+
+ emitters.put(orderId, emitter);
+ log.info("[SSE] Client subscribed to order {} status stream (active connections: {})",
+ orderId, emitters.size());
+
+ // Send an initial "connected" event so the client knows the stream is live
+ try {
+ emitter.send(SseEmitter.event()
+ .name("connected")
+ .data("{\"orderId\":\"" + orderId + "\",\"message\":\"Subscribed to order status stream\"}"));
+ } catch (IOException e) {
+ emitters.remove(orderId);
+ }
+
+ return emitter;
+ }
+
+ /**
+ * Push a status update to the client subscribed to this orderId.
+ * Called by {@link com.hacisimsek.order.service.impl.OrderServiceImpl}
+ * whenever the order status changes.
+ *
+ * @param orderId the order that changed
+ * @param newStatus the new status string (e.g. "PAYMENT_COMPLETED")
+ * @param terminal true if this is the final state (SHIPPED, COMPLETED, FAILED, CANCELLED)
+ */
+ public void push(UUID orderId, String newStatus, boolean terminal) {
+ SseEmitter emitter = emitters.get(orderId);
+ if (emitter == null) {
+ // No connected client — that's fine, most users poll instead of streaming
+ return;
+ }
+
+ String payload = String.format(
+ "{\"orderId\":\"%s\",\"status\":\"%s\",\"timestamp\":\"%s\"}",
+ orderId, newStatus, Instant.now());
+
+ try {
+ emitter.send(SseEmitter.event()
+ .name("status-update")
+ .data(payload));
+
+ log.info("[SSE] Pushed status {} to order {} subscriber", newStatus, orderId);
+
+ // Complete the stream on terminal states — no more updates coming
+ if (terminal) {
+ emitter.send(SseEmitter.event()
+ .name("complete")
+ .data("{\"message\":\"Order reached terminal state: " + newStatus + "\"}"));
+ emitter.complete();
+ emitters.remove(orderId);
+ }
+ } catch (IOException e) {
+ log.warn("[SSE] Failed to push to order {} subscriber — removing: {}", orderId, e.getMessage());
+ emitters.remove(orderId);
+ }
+ }
+
+ /** Returns the number of active SSE connections — useful for monitoring */
+ public int activeConnections() {
+ return emitters.size();
+ }
+}
diff --git a/order-service/src/main/resources/application.yml b/order-service/src/main/resources/application.yml
new file mode 100644
index 0000000..7e90175
--- /dev/null
+++ b/order-service/src/main/resources/application.yml
@@ -0,0 +1,75 @@
+server:
+ port: 8081
+
+spring:
+ application:
+ name: order-service
+
+ datasource:
+ url: jdbc:postgresql://${DB_HOST:localhost}:${DB_PORT:5432}/order_db
+ username: ${DB_USER:postgres}
+ password: ${DB_PASSWORD:postgres}
+ driver-class-name: org.postgresql.Driver
+ hikari:
+ maximum-pool-size: 10
+ minimum-idle: 2
+ connection-timeout: 30000
+ idle-timeout: 600000
+ max-lifetime: 1800000
+
+ jpa:
+ hibernate:
+ ddl-auto: update
+ show-sql: false
+ properties:
+ hibernate:
+ format_sql: true
+ dialect: org.hibernate.dialect.PostgreSQLDialect
+
+ kafka:
+ bootstrap-servers: ${KAFKA_BOOTSTRAP_SERVERS:localhost:9095}
+ producer:
+ key-serializer: org.apache.kafka.common.serialization.StringSerializer
+ value-serializer: org.springframework.kafka.support.serializer.JsonSerializer
+ properties:
+ spring.json.add.type.headers: true
+ consumer:
+ group-id: order-service-group
+ auto-offset-reset: earliest
+ enable-auto-commit: false
+ properties:
+ session.timeout.ms: 30000
+ heartbeat.interval.ms: 10000
+ max.poll.interval.ms: 300000
+ listener:
+ ack-mode: record
+
+eureka:
+ client:
+ service-url:
+ defaultZone: http://${EUREKA_HOST:localhost}:8761/eureka/
+ instance:
+ prefer-ip-address: true
+
+management:
+ endpoints:
+ web:
+ exposure:
+ include: health,info,prometheus,metrics
+ endpoint:
+ health:
+ show-details: always
+ prometheus:
+ enabled: true
+ metrics:
+ tags:
+ application: ${spring.application.name}
+ distribution:
+ percentiles-histogram:
+ http.server.requests: true
+ percentiles:
+ http.server.requests: 0.5, 0.95, 0.99
+
+logging:
+ level:
+ com.hacisimsek.order: INFO
diff --git a/payment-service/Dockerfile b/payment-service/Dockerfile
new file mode 100644
index 0000000..0ead949
--- /dev/null
+++ b/payment-service/Dockerfile
@@ -0,0 +1,29 @@
+# -- Stage 1: Build ----------------------------------------------------------
+FROM eclipse-temurin:21-jdk-alpine AS builder
+WORKDIR /build
+
+COPY pom.xml .
+COPY common-library/pom.xml common-library/
+COPY common-library/src common-library/src
+COPY payment-service/pom.xml payment-service/
+COPY payment-service/src payment-service/src
+
+RUN --mount=type=cache,target=/root/.m2 `
+ ./mvnw -pl common-library,payment-service -am clean package -DskipTests --no-transfer-progress
+
+# -- Stage 2: Runtime --------------------------------------------------------
+FROM eclipse-temurin:21-jre-alpine
+WORKDIR /app
+
+RUN addgroup -S appgroup && adduser -S appuser -G appgroup
+USER appuser
+
+COPY --from=builder /build/payment-service/target/*.jar app.jar
+
+EXPOSE 8083
+
+ENTRYPOINT ["java", `
+ "-XX:+UseContainerSupport", `
+ "-XX:MaxRAMPercentage=75.0", `
+ "-Djava.security.egd=file:/dev/./urandom", `
+ "-jar", "app.jar"]
diff --git a/payment-service/pom.xml b/payment-service/pom.xml
index d2e3be4..0f51b1f 100644
--- a/payment-service/pom.xml
+++ b/payment-service/pom.xml
@@ -66,6 +66,11 @@
io.micrometer
micrometer-registry-prometheus
+
+
+ org.springdoc
+ springdoc-openapi-starter-webmvc-ui
+
diff --git a/payment-service/src/main/java/com/hacisimsek/payment/controller/PaymentController.java b/payment-service/src/main/java/com/hacisimsek/payment/controller/PaymentController.java
index 93b398e..5ac866c 100644
--- a/payment-service/src/main/java/com/hacisimsek/payment/controller/PaymentController.java
+++ b/payment-service/src/main/java/com/hacisimsek/payment/controller/PaymentController.java
@@ -1,4 +1,4 @@
-package com.hacisimsek.payment.controller;
+package com.hacisimsek.payment.controller;
import com.hacisimsek.payment.dto.GatewayOrderResponse;
import com.hacisimsek.payment.dto.InitiatePaymentRequest;
@@ -24,7 +24,7 @@
import java.util.UUID;
@RestController
-@RequestMapping("/api/payments")
+@RequestMapping("/api/v1/payments")
@RequiredArgsConstructor
@Slf4j
public class PaymentController {
@@ -32,7 +32,7 @@ public class PaymentController {
private final PaymentService paymentService;
private final RestTemplate restTemplate;
- /** Initiate a payment — returns gateway order/session for frontend checkout UI. */
+ /** Initiate a payment — returns gateway order/session for frontend checkout UI. */
@PostMapping("/initiate")
public ResponseEntity initiatePayment(
@Valid @RequestBody InitiatePaymentRequest request) {
diff --git a/payment-service/src/main/java/com/hacisimsek/payment/controller/WebhookController.java b/payment-service/src/main/java/com/hacisimsek/payment/controller/WebhookController.java
index 0aeb8a7..2917383 100644
--- a/payment-service/src/main/java/com/hacisimsek/payment/controller/WebhookController.java
+++ b/payment-service/src/main/java/com/hacisimsek/payment/controller/WebhookController.java
@@ -1,4 +1,4 @@
-package com.hacisimsek.payment.controller;
+package com.hacisimsek.payment.controller;
import com.fasterxml.jackson.databind.JsonNode;
import com.fasterxml.jackson.databind.ObjectMapper;
@@ -23,13 +23,13 @@
/**
* Receives webhook events pushed by Razorpay to your server.
*
- * Razorpay Dashboard → Settings → Webhooks → Add new webhook:
+ * Razorpay Dashboard → Settings → Webhooks → Add new webhook:
* URL: https:///api/payments/webhook/razorpay
* Events: payment.captured, payment.failed, refund.created
* Secret: value of RAZORPAY_WEBHOOK_SECRET env var
*
* IMPORTANT: This endpoint is intentionally excluded from JWT auth in the
- * API Gateway / Security config — Razorpay calls it directly, not the user.
+ * API Gateway / Security config — Razorpay calls it directly, not the user.
* Security is provided solely by HMAC-SHA256 signature verification.
*
* Spring must receive the raw bytes (not a parsed object) so the signature
@@ -37,7 +37,7 @@
* and convert to String only after verification passes.
*/
@RestController
-@RequestMapping("/api/payments/webhook")
+@RequestMapping("/api/v1/payments/webhook")
@RequiredArgsConstructor
@Slf4j
public class WebhookController {
@@ -45,7 +45,7 @@ public class WebhookController {
private final PaymentService paymentService;
private final ObjectMapper objectMapper;
- /** All gateway adapters — used to look up the Razorpay adapter by type. */
+ /** All gateway adapters — used to look up the Razorpay adapter by type. */
private final List gatewayAdapters;
/**
@@ -56,9 +56,9 @@ public class WebhookController {
* X-Razorpay-Signature:
*
* Response contract:
- * 200 OK → event acknowledged (Razorpay will not retry)
- * 400 → signature invalid (logged, no retry by Razorpay for bad sig)
- * 500 → processing error (Razorpay WILL retry — safe to throw on transient errors)
+ * 200 OK → event acknowledged (Razorpay will not retry)
+ * 400 → signature invalid (logged, no retry by Razorpay for bad sig)
+ * 500 → processing error (Razorpay WILL retry — safe to throw on transient errors)
*/
@PostMapping(
value = "/razorpay",
@@ -73,24 +73,24 @@ public ResponseEntity