# CQRS+ES知識

## CQRS+ES採用判断

CQRS+ES は、状態変更をドメイン上の出来事として保存し、そこから現在状態やRead Modelを導出する設計である。バックエンド全体やワークフローが CQRS+ES を扱う場合でも、すべての新機能をイベントソーシングで実装する必要はない。

CQRS+ES の採用は要件から導く。既存システムが CQRS+ES を含むことは、依存方向や境界をそろえる理由にはなるが、単純な設定テーブルまでイベントソーシング化する理由にはならない。

### 要件変換時の扱い

元要件やユーザー要求が CRUD 相当の業務要件だけを述べている場合、仕様に「コマンド・イベント・プロジェクション」を新しい要件として追加しない。CQRS+ES が必要か不明な場合は、採用理由を明示するか、未確認事項として残す。

| 元要求 | 意味・選択肢 |
|--------|--------------------|
| 「施設ごとに許可IPを管理したい」 | CRUDの管理設定として扱う。ドメイン語彙が「追加・削除」だけで業務ルールがない |
| 「注文の承認・取消・返品を管理し、状態に応じて請求や在庫が連動する」 | CQRS+ESの候補。複雑な状態遷移と業務不変条件があり、複数集約が連動する |
| 「保険の契約変更で、変更種別ごとに審査ルールが異なり、過去の査定履歴が将来の判断に影響する」 | CQRS+ESの候補。ビジネスルールが複雑で変化し、履歴そのものが業務判断の入力になる |
| 「誰がいつ変更したかを画面に出したい」 | CRUD + 監査ログで足りるか確認する。変更履歴の表示だけなら監査列で十分 |
| 「通知設定のON/OFFを切り替えたい」 | CRUDの管理設定として扱う。現在値の参照・更新のみ |

CQRS+ES は複雑なビジネスドメイン（金融、保険、医療など、ビジネスルールが複雑で変化するドメイン）でその真価を発揮する。単純な監査要件や技術的な非同期処理は、それだけでは CQRS+ES の十分条件にならない。判断の軸はビジネスロジックの複雑さにある。

## Aggregate設計

Aggregateは判断に必要なフィールドのみ保持する。

Command Model（Aggregate）の役割は「コマンドを受けて判断し、イベントを発行する」こと。クエリ用データはRead Model（Projection）が担当する。

「判断に必要」とは:
- `if`/`require`の条件分岐に使う
- インスタンスメソッドでイベント発行時にフィールド値を参照する

`if`/`require` に使われていることだけでは、Aggregate state に保持する根拠にならない。先に、その分岐や検証が Aggregate 全体の本質的な不変条件かを確認する。

### 由来メタデータと不変条件

入力元、チャネル、生成元、連携元などの由来メタデータは、表示・検索・監査・連携追跡で必要になることがある。ただし、それだけでは Aggregate の状態として復元する理由にならない。

```kotlin
// 避ける例: 由来メタデータで既存Aggregateの通常ライフサイクルを狭めている
data class Note(
    val noteId: String,
    val sourceType: SourceType?,
    val targetIds: List<String>,
) {
    fun update(text: String, targetIds: List<String>): NoteUpdatedEvent {
        if (sourceType == SourceType.EXTERNAL_IMPORT) {
            require(targetIds.isNotEmpty())
        }
        return NoteUpdatedEvent(noteId, text, targetIds)
    }
}

// 例: 由来はイベント・Read Modelで追跡し、Aggregateの不変条件は通常ライフサイクルに合わせる
data class Note(
    val noteId: String,
    val confirmed: Boolean,
) {
    fun update(text: String, targetIds: List<String>): NoteUpdatedEvent {
        check(!confirmed)
        return NoteUpdatedEvent(noteId, text, targetIds)
    }
}

data class NoteCreatedEvent(
    val noteId: String,
    val text: String,
    val targetIds: List<String>,
    val sourceType: SourceType?, // Projectionや監査で使う由来事実
)
```

### 既存ライフサイクルの優先

既存 Aggregate に新しい入力フローを統合する場合、既存の通常ライフサイクルを優先する。入力元が違うだけで専用 command、専用 wrapper、専用 service、専用削除処理を増やさない。

良いAggregate:
```kotlin
// 判断に必要なフィールドのみ
data class Order(
    val orderId: String,      // イベント発行時に使用
    val status: OrderStatus   // 状態チェックに使用
) {
    fun confirm(confirmedBy: String): OrderConfirmedEvent {
        require(status == OrderStatus.PENDING) { "確定できる状態ではありません" }
        return OrderConfirmedEvent(
            orderId = orderId,
            confirmedBy = confirmedBy,
            confirmedAt = LocalDateTime.now()
        )
    }
}

// 判断に使わないフィールドを保持（NG）
data class Order(
    val orderId: String,
    val customerId: String,     // 判断に未使用
    val shippingAddress: Address, // 判断に未使用
    val status: OrderStatus
)
```

追加操作がないAggregateはIDのみ:
```kotlin
// 作成のみで追加操作がない場合
data class Notification(val notificationId: String) {
    companion object {
        fun create(customerId: String, message: String): NotificationCreatedEvent {
            return NotificationCreatedEvent(
                notificationId = UUID.randomUUID().toString(),
                customerId = customerId,
                message = message
            )
        }
    }
}
```

### Adapterパターン（ドメインとフレームワークの分離）

ドメインモデルにフレームワークのアノテーション（`@Aggregate`, `@CommandHandler`等）を直接付けない。Adapterクラスがフレームワーク統合を担当し、ドメインモデルはビジネスロジックに専念する。

```kotlin
// ドメインモデル: フレームワーク非依存。ビジネスロジックのみ
data class Order(
    val orderId: String,
    val status: OrderStatus = OrderStatus.PENDING
) {
    companion object {
        fun place(orderId: String, customerId: String): OrderPlacedEvent {
            require(customerId.isNotBlank()) { "Customer ID cannot be blank" }
            return OrderPlacedEvent(orderId, customerId)
        }

        fun from(event: OrderPlacedEvent): Order {
            return Order(orderId = event.orderId, status = OrderStatus.PENDING)
        }
    }

    fun confirm(confirmedBy: String): OrderConfirmedEvent {
        require(status == OrderStatus.PENDING) { "確定できる状態ではありません" }
        return OrderConfirmedEvent(orderId, confirmedBy, LocalDateTime.now())
    }

    fun apply(event: OrderEvent): Order = when (event) {
        is OrderPlacedEvent -> from(event)
        is OrderConfirmedEvent -> copy(status = OrderStatus.CONFIRMED)
        is OrderCancelledEvent -> copy(status = OrderStatus.CANCELLED)
    }
}

// Adapter: フレームワーク統合。ドメイン呼び出し → イベント発行の中継
@Aggregate
class OrderAggregateAdapter() {
    private var order: Order? = null

    @AggregateIdentifier
    fun orderId(): String? = order?.orderId

    @CommandHandler
    constructor(command: PlaceOrderCommand) : this() {
        val event = Order.place(command.orderId, command.customerId)
        AggregateLifecycle.apply(event)
    }

    @CommandHandler
    fun handle(command: ConfirmOrderCommand) {
        val event = order!!.confirm(command.confirmedBy)
        AggregateLifecycle.apply(event)
    }

    @EventSourcingHandler
    fun on(event: OrderEvent) {
        this.order = when (event) {
            is OrderPlacedEvent -> Order.from(event)
            else -> order?.apply(event)
        }
    }
}
```

分離の利点:
- ドメインモデル単体でユニットテスト可能（フレームワーク不要）
- フレームワーク移行時にドメインモデルは変更不要
- Adapterはコマンド受信 → ドメイン呼び出し → イベント発行の定型コード

### apply/from パターン（イベント再生）

ドメインモデルが自身の状態をイベントから再構築するパターン。

- `from(event)`: 生成イベントから初期状態を構築するファクトリ
- `apply(event)`: イベントを受けて新しい状態を返す（`copy()` でイミュータブルに更新）
- `when` 式 + sealed interface で全イベント型の網羅性をコンパイラが保証

```kotlin
fun apply(event: OrderEvent): Order = when (event) {
    is OrderPlacedEvent -> from(event)
    is OrderConfirmedEvent -> copy(status = OrderStatus.CONFIRMED)
    is OrderShippedEvent -> copy(status = OrderStatus.SHIPPED)
    // sealed interface なので、イベント型の追加漏れはコンパイルエラーになる
}
```

## イベント設計

良いイベント:
```kotlin
// Good: ドメインの意図が明確
OrderPlaced, PaymentReceived, ItemShipped

// 避ける例: CRUDスタイル
OrderUpdated, OrderDeleted
```

### 事実イベントと要求イベント

イベントは発生した事実を表し、名前はその業務上の意味から決める。`〜Requested` という接尾辞や現在の消費者数だけでは、事実イベントか偽装コマンドかを判定できない。要求を受理・開始したこと自体が業務上の出来事なら事実になり得るが、既知の宛先へ処理を命令する以外の意味を持たない通知はコマンドとして表す。

| 比較軸 | 事実イベント | コマンドを検討 |
|--------|-------------|----------------|
| 業務上の意味 | 受理・開始・却下など、監査や状態復元にも意味がある出来事 | 特定処理を実行させることだけが目的 |
| 発行元のライフサイクル | 待機、重複拒否、期限切れ等の後続判断に使う | 発行元は結果や状態を追跡しない |
| 結果 | 完了・失敗等の別の事実へ続く、不確実な処理 | 同一境界内で直ちに実行でき、結果もその場で扱える |
| 消費者 | 増減してもイベントの意味が変わらない | 宛先と処理内容がイベントの意味そのもの |

イベントは技術的な消費先ではなく、独立した業務上の事実を単位に分ける。同じ出来事を状態復元用と処理起動用に重複発行せず、発行元集約が所有する事実を EventHandler や Projection がそれぞれ購読する。複数のイベントが必要なのは、別々に命名・監査・再生する意味を持つ事実が同時に成立した場合である。別集約の内部状態や初期化詳細はイベントに詰め込まず、安定したIDや参照を境界側で解決する。

```kotlin
// 避ける例: 1つの業務上の事実を、状態用と技術的な処理起動用に重複分割
fun addItem(itemId: String, productId: String, quantity: Int): List<OrderEvent> = listOf(
    OrderItemLinkedEvent(orderId, itemId),                 // 状態用
    OrderItemCreationRequestedEvent(orderId, itemId, productId, quantity), // トリガー用（実質コマンド）
)

// 例: 発行元集約の事実として必要な内容を持つ単一イベント
fun addItem(itemId: String, productId: String, quantity: Int): OrderItemAddedEvent =
    OrderItemAddedEvent(orderId, itemId, productId, quantity)
```

### sealed interface によるイベント型階層

集約のイベントは sealed interface で型階層化する。集約ルートIDを共通フィールドとして強制し、`when` 式の網羅性チェックを有効にする。

```kotlin
sealed interface OrderEvent {
    val orderId: String  // 全イベントに必須
}

data class OrderPlacedEvent(
    override val orderId: String,
    val customerId: String
) : OrderEvent

data class OrderConfirmedEvent(
    override val orderId: String,
    val approvalInfo: ApprovalInfo
) : OrderEvent

data class OrderCancelledEvent(
    override val orderId: String,
    val cancellationInfo: CancellationInfo
) : OrderEvent
```

利点:
- `when (event)` で全イベント型を列挙しないとコンパイルエラー（`apply` メソッドで特に重要）
- 集約ルートIDの存在をコンパイラが保証
- 型ベースのイベントハンドラ分岐が安全

イベント粒度:
- 細かすぎ: `OrderFieldChanged` → ドメインの意図が不明
- 適切: `ShippingAddressChanged` → 意図が明確
- 粗すぎ: `OrderModified` → 何が変わったか不明

## Event Evolution

イベント進化では、現行イベント契約、履歴payloadの変換、イベント再生による状態復元を別の責務として扱う。現行イベント型とドメインロジックは現在の意味だけを表す。履歴payloadの変換を行う場合は、イベントストアから復元する境界で replay 前に変換する。

イベント進化で分ける責務:

| 責務 | 置き場所 |
|------|----------|
| 現行イベントの意味とフィールド | イベント型 |
| 設計対象となる場合の履歴payloadの読み替え | event-store 復元境界の upcaster |
| イベント再生による状態復元 | Aggregate の `apply` |
| 履歴payload読み替えの振る舞い証跡 | upcaster テスト |

```kotlin
// 現行イベント型
data class OrderAssignedEvent(
    override val orderId: String,
    val assigneeIds: List<String>
) : OrderEvent
```

```kotlin
// 例 - 履歴 payload を復元境界の upcaster で現行 payload へ変換する
when (eventType) {
    OrderAssignedEvent::class.java.typeName -> {
        event.moveTextFieldToArray("assigneeId", "assigneeIds")
    }
}
```

履歴変換が設計対象となる場合、旧イベント型そのものをアプリケーションコードに残すかは、利用フレームワークと移行方式で決まる。旧 serialized type と payload は、現行ドメインイベントに含めずに upcaster の入力契約として扱える。

### migration の責務境界

CQRS+ES では、DB schema migration、data migration、event upcaster、Read Model rebuild、API互換対応がそれぞれ異なる契約と実行境界を持つ。

| migration 種別 | 責務境界 |
|----------------|----------|
| DB schema migration | relational schema の変更 |
| data migration / backfill | relational data の変換 |
| event upcaster | event-store 復元時の履歴payload変換 |
| Read Model rebuild | イベントから導出可能な projection の再生成 |
| API compatibility | 外部利用側との契約境界 |

## コマンドハンドラ

### コマンドとイベントの契約寿命

イベントは履歴として永続化される長寿命の契約である。履歴payloadの変換を行う場合は、現行イベントの型識別子・payloadと、履歴payloadを replay 可能な形へ変換する境界を分け、変換方式はイベントストアとシリアライズ方式から選ぶ。

コマンドは通常、application 境界で生成・処理される短寿命のメッセージだが、予約実行、outbox、再試行、dead-letter、監査等で永続化される構成もある。永続参照の有無は、移動・改名時に調べる影響境界である。ドメインモデルは配送方式やフレームワークのコマンド型に依存せず、application / adapter 境界でドメインの引数・値オブジェクトへ変換する。

良いコマンドハンドラ:
```
1. コマンドを受け取る
2. Aggregateをイベントストアから復元
3. Aggregateにコマンドを適用
4. 発行されたイベントを保存
```

### 多層バリデーション

バリデーションは層ごとに役割が異なる。すべてを1箇所に集めない。

| 層 | 責務 | 手段 | 例 |
|----|------|------|-----|
| API層 | 構造的バリデーション | `@NotBlank`, `init` ブロック | 必須項目、型、フォーマット |
| UseCase層 | ビジネスルール検証 | Read Modelへの問い合わせ | 重複チェック、前提条件の存在確認 |
| ドメイン層 | 状態遷移の不変条件 | `require` | 「PENDINGでないと承認できない」 |

### Aggregateの判断境界

Aggregate は、自身のイベント履歴から復元できる状態と、コマンドとして明示された事実だけで判断する。境界由来の入力を解釈・正規化・所有権確認する場所ではない。

Aggregate に入れてよい検証は「イベント再生だけで再現できる状態」に基づくものに限る。それ以外の検証は、コマンド送信前に境界側で解決し、Aggregate には解決済みの事実を渡す。

| 対象 | 置き場所 |
|---------|---------|
| 現在状態でその操作が可能か | Aggregate |
| コマンド実行者がAggregate ownerと一致するか | Aggregate |
| HTTP/API入力の形式が正しいか | API層 |
| object key、URL、path などの外部識別子の形式解釈 | UseCase層または境界側Policy/Verifier |
| 外部識別子が現在user/tenantに属するか | UseCase層または境界側Policy/Verifier |
| 他AggregateのRead Modelや外部事実の確認 | UseCase層 |
| 同じAggregateの現在状態に基づく状態遷移判断 | Aggregate |
| 外部サービス上に実体があるか | Application層の外部サービス連携 |

例: アップロード完了コマンドでは、Aggregate は「このセッションのownerと実行者が一致するか」「現在状態で完了可能か」を判断する。保存先object keyの文字列形式や、そのkeyが現在user/tenantの領域かどうかは、コマンド送信前にUseCase層で検証する。

### Command 意図と事前 Query

Command は「現在の状態を見て何の command を送るか」ではなく、「利用者や外部処理が何をしたいか」を表す。現在状態に基づく Add / Update / Delete / Noop の判断は、同じ Aggregate の Read Model ではなく、復元済み Aggregate に寄せる。

```kotlin
// 避ける例: Query 結果で command 種別を投げ分けている
if (readService.exists(orderId)) {
    commandGateway.send(UpdateOrderCommand(orderId, value))
} else {
    commandGateway.send(AddOrderCommand(orderId, value))
}

// 例: 意図 command を送り、Aggregate が復元済み状態で判断する
commandGateway.send(SetOrderValueCommand(orderId, value))
```

```kotlin
// API層: 構造的バリデーション
data class OrderPostRequest(
    @field:NotBlank val customerId: String,
    @field:NotNull val items: List<OrderItemRequest>
) {
    init {
        require(items.isNotEmpty()) { "注文には1つ以上の商品が必要です" }
    }
}

// UseCase層: ビジネスルール検証（Read Model参照）
@Service
class PlaceOrderUseCase(
    private val commandGateway: CommandGateway,
    private val customerRepository: CustomerRepository,
    private val inventoryRepository: InventoryRepository
) {
    fun execute(input: PlaceOrderInput): Mono<PlaceOrderOutput> {
        return Mono.fromCallable {
            // 顧客の存在確認
            customerRepository.findById(input.customerId)
                ?: throw CustomerNotFoundException("顧客が存在しません")
            // 在庫の事前確認
            validateInventory(input.items)
            // コマンド送信
            val orderId = UUID.randomUUID().toString()
            commandGateway.send<Any>(PlaceOrderCommand(orderId, input.customerId, input.items))
            PlaceOrderOutput(orderId)
        }
    }
}

// ドメイン層: 状態遷移の不変条件
fun confirm(confirmedBy: String): OrderConfirmedEvent {
    require(status == OrderStatus.PENDING) { "確定できる状態ではありません" }
    return OrderConfirmedEvent(orderId, confirmedBy, LocalDateTime.now())
}
```

## UseCase層（オーケストレーション）

Controller と CommandGateway の間にUseCase層を置く。UseCase層は境界で解決すべき事実を集め、原則として1つの意図 command を送る。後続の状態変更は、確定済みイベントを起点に EventHandler が進める。

```
Controller → UseCase → CommandGateway → Aggregate
                ↓
          QueryGateway / Repository（Read Model参照）
```

UseCaseが必要なケース:
- コマンド発行前に他集約のRead Modelや外部事実を確認する
- 複数のバリデーションを直列に実行する
- 同期 API 契約のため、コマンド送信後の結果整合性を待機する

UseCaseが不要なケース:
- Controllerからコマンドを1つ送るだけで完結する単純な操作
- ControllerからQuery側へ問い合わせてレスポンスへ変換するだけの単純な参照
- 既存リソースの存在確認・スコープ確認後にコマンドを1つ送るだけの操作

## イベントドリブン連鎖

CQRS+ES では、状態変更の連鎖は確定済みイベントを起点に進める。Application Service / UseCase / Controller が同じ状態遷移のために command を直列に投げて、複数 Aggregate の変更順序を同期制御しない。

基本形:

```text
UseCase → Command → Aggregate → Event
                              ↓
                         EventHandler → Command → 別Aggregate
                              ↓
                         Projection → Read Model
```

## プロジェクション設計

良いプロジェクション:
- 特定の読み取りユースケースに最適化
- イベントから冪等に再構築可能
- Writeモデルから完全に独立

### Projection と EventHandler（サイドエフェクト）の区別

どちらも `@EventHandler` を使うが、責務が異なる。混同しない。

| 種類 | 責務 | やること | やらないこと |
|------|------|---------|-------------|
| Projection | Read Model 更新 | Entity の保存・更新 | コマンド送信、外部API呼び出し |
| EventHandler | サイドエフェクト | 他集約へのコマンド送信 | Read Model 更新 |

```kotlin
// Projection: Read Model 更新のみ
@Component
class OrderProjection(private val orderRepository: OrderRepository) {
    @EventHandler
    fun on(event: OrderPlacedEvent) {
        val entity = OrderEntity(
            orderId = event.orderId,
            customerId = event.customerId,
            status = OrderStatus.PENDING
        )
        orderRepository.save(entity)
    }

    @EventHandler
    fun on(event: OrderConfirmedEvent) {
        orderRepository.findById(event.orderId).ifPresent { entity ->
            entity.status = OrderStatus.CONFIRMED
            orderRepository.save(entity)
        }
    }
}

// EventHandler: サイドエフェクト（他集約へのコマンド送信）
@Component
class InventoryReleaseHandler(private val commandGateway: CommandGateway) {
    @EventHandler
    fun on(event: OrderCancelledEvent) {
        val command = ReleaseInventoryCommand(
            productId = event.productId,
            quantity = event.quantity
        )
        commandGateway.send<Any>(command)
    }
}
```

### 外部処理の起動

外部ワーカーや非同期処理の起動は、Aggregate が確定したドメインイベントを起点にする。Application Service や Coordinator が、コマンド送信と外部副作用を同じ制御フローで束ねない。

## Query側の設計

Query側はイベント駆動のPubSubモデルで動作する。Projection が EventHandler でRead Modelを更新し、Query側はRead Modelを参照する。

イベント配信はPubSub（メッセージブローカー経由）で全インスタンスに配信する。同一インスタンスへの配信を前提とする仕組みは、配送保証が確認できない限り使わない。

- **Subscription Query**（たとえばAxonの `subscriptionQuery()`）: クエリ結果の変更通知を購読元へ返す仕組み。既存基盤として採用され、購読元への通知配送が保証される構成でのみ使う。tracking processor や tracker を前提にした構成では、機能実装のためだけに subscription query を新設しない。
- **Subscribing イベントプロセッサ**（たとえばAxonの `SubscribingEventProcessor`）: ローカルのイベントバスからの直接購読に依存し、イベントを発行したインスタンスのみがイベントを受け取る。分散環境では他インスタンスの Projection が更新されない。PubSubで全インスタンスにイベントが配信される構成にする。

### QueryHandler と ApplicationService の命名

CQRSではクエリを受けるコンポーネントを QueryHandler と呼び、クエリを送る入口は QueryGateway / QueryBus として扱う。Controller から読み取りユースケースを呼ぶ facade は、QueryHandler と混同しないよう ApplicationService または ReadService と名付ける。

レイヤー間の型:
- `application/query/` - Query結果の型（例: `OrderDetail`）
- `adapter/protocol/` - RESTレスポンスの型（例: `OrderDetailResponse`）
- QueryHandler は application層の型を返し、Controller が adapter層の型に変換

```kotlin
// application/query/OrderDetail.kt
data class OrderDetail(
    val orderId: String,
    val customerName: String,
    val totalAmount: Money
)

// adapter/protocol/OrderDetailResponse.kt
data class OrderDetailResponse(...) {
    companion object {
        fun from(detail: OrderDetail) = OrderDetailResponse(...)
    }
}

// QueryHandler - application層の型を返す
@QueryHandler
fun handle(query: GetOrderDetailQuery): OrderDetail? {
    val entity = repository.findById(query.id) ?: return null
    return OrderDetail(...)
}

// Controller - 単純な参照は同期返却で十分
@GetMapping("/{id}")
fun getById(@PathVariable id: String): ResponseEntity<OrderDetailResponse> {
    val detail = queryGateway.query(
        GetOrderDetailQuery(id),
        OrderDetail::class.java
    ).join() ?: throw NotFoundException("...")

    return ResponseEntity.ok(OrderDetailResponse.from(detail))
}
```

構成:
```
Controller (adapter) → QueryGateway → QueryHandler (application) → Repository
     ↓                                      ↓
Response.from(detail)                  OrderDetail

イベント流（PubSub）:
Aggregate → Event Bus → Projection(@EventHandler) → Repository(Read Model)
                                                          ↑
                                          QueryHandler がここを参照
```

### 非同期コールバックと並行制御

非同期処理の完了通知は重複・遅延・順序逆転を前提に設計する。Controller や単一プロセス内のロックではなく、Aggregate の状態遷移とコマンドの冪等性で守る。

## 結果整合性

コマンド発行後の Projection 待機は、同一 API レスポンスで更新後 Read Model を返す明示的な同期契約がある場合に限る。画面側が入力値やIDを保持できる場合、サーバーは待機せず、Read Model の収束を通常の参照 API で扱う。

### リアクティブポーリング

コマンド発行 → Projection更新完了を非ブロッキングなポーリングで待機するパターン。リアクティブポーリングはリクエストスレッドを占有しない待機であり、`while` ループと `Thread.sleep` で同期的に待つ実装ではない。

ポーリングの判定はイベント通知ではなく、Read Model を再取得して期待する状態になったかを predicate で確認する。条件を満たすまで一定間隔で再取得し、timeout または maxAttempts に達したら待機を打ち切る。

```kotlin
// UseCase: コマンド送信 → ポーリングで完了待機
fun execute(input: PlaceOrderInput): Mono<PlaceOrderOutput> {
    val orderId = UUID.randomUUID().toString()
    return Mono.fromCallable { validatePreConditions(input) }
        .subscribeOn(Schedulers.boundedElastic())
        .flatMap {
            Mono.fromFuture(commandGateway.send<Any>(
                PlaceOrderCommand(orderId, input.customerId, input.items)
            ))
        }
        .then(pollForCompletion(orderId))
        .thenReturn(PlaceOrderOutput(orderId))
}

// ポーリング: Projection の更新を待機
private fun pollForCompletion(orderId: String): Mono<Void> {
    return ReactivePolling.waitFor(
        supplier = { orderRepository.findById(orderId).orElse(null) },
        condition = { it.sagaCompleted || it.status == OrderStatus.CONFIRMED },
        timeout = Duration.ofSeconds(60),
        maxAttempts = 300
    )
}
```

ブロッキング待機は避ける:

```kotlin
// 避ける例: リクエストスレッドを占有し、負荷時にスレッド枯渇を起こす
while (Instant.now().isBefore(deadline)) {
    val order = orderRepository.findById(orderId).orElse(null)
    if (order?.status == OrderStatus.CONFIRMED) return PlaceOrderOutput(orderId)
    Thread.sleep(100)
}

// 例: 同一レスポンスで待つならリアクティブな待機へ載せる
return pollForCompletion(orderId).thenReturn(PlaceOrderOutput(orderId))
```

ポーリングが適切なケース:
- Saga が完了するまでレスポンスを返したくない場合
- コマンド発行後に作成されたリソースのIDを返す場合

ポーリングが不要なケース:
- コマンド発行だけで完了する単純な操作（結果を待たない）
- UIがリアルタイム更新を必要としない場合

サーバー側で待たない場合は、コマンド受付後に `202 Accepted` と追跡IDを返し、フロントエンドが読み取りAPIをロングポーリングまたは通常ポーリングする。ユーザー体験上の即時性が必要なら SSE や WebSocket も選択肢に含める。

## Saga vs EventHandler

Sagaは「競合が発生する複数アグリゲート間の操作」にのみ使用する。

Sagaが必要なケース:
```
複数のアクターが同じリソースを取り合う場合
例: 在庫確保（10人が同時に同じ商品を注文）

OrderPlacedEvent
  ↓ InventoryReservationSaga
ReserveInventoryCommand → Inventory集約（同時実行を直列化）
  ↓
InventoryReservedEvent → ConfirmOrderCommand
InventoryReservationFailedEvent → CancelOrderCommand
```

Sagaが不要なケース:
```
競合が発生しない操作
例: 注文キャンセル時の在庫解放

OrderCancelledEvent
  ↓ InventoryReleaseHandler（単純なEventHandler）
ReleaseInventoryCommand
  ↓
InventoryReleasedEvent
```

判断基準:

| 状況 | Saga | EventHandler |
|------|------|--------------|
| リソースの取り合いがある | 使う | - |
| 補償トランザクションが必要 | 使う | - |
| 競合しない単純な連携 | - | 使う |
| 失敗時は再試行で十分 | - | 使う |

アンチパターン:
```kotlin
// 避ける例: ライフサイクル管理のためにSagaを使う
@Saga
class OrderLifecycleSaga {
    // 注文の全状態遷移をSagaで追跡
    // PLACED → CONFIRMED → SHIPPED → DELIVERED
}

// 例: 結果整合性が必要な操作だけをSagaで処理
@Saga
class InventoryReservationSaga {
    // 在庫確保の同時実行制御のみ
}
```

Sagaはライフサイクル管理ツールではない。結果整合性が必要な「操作」単位で作成する。

## 例外 vs イベント（失敗時の選択）

監査不要な失敗は例外、監査が必要な失敗はイベント。

例外アプローチ（推奨: ほとんどのケース）:
```kotlin
// ドメインモデル: バリデーション失敗時に例外をスロー
fun reserveInventory(orderId: String, quantity: Int): InventoryReservedEvent {
    if (availableQuantity < quantity) {
        throw InsufficientInventoryException("在庫が不足しています")
    }
    return InventoryReservedEvent(productId, orderId, quantity)
}

// Saga: exceptionally でキャッチして補償アクション
commandGateway.send<Any>(command)
    .exceptionally { ex ->
        commandGateway.send<Any>(CancelOrderCommand(
            orderId = orderId,
            reason = ex.cause?.message ?: "在庫確保に失敗しました"
        ))
        null
    }
```

イベントアプローチ（稀なケース）:
```kotlin
// 監査が必要な場合のみ
data class PaymentFailedEvent(
    val paymentId: String,
    val reason: String,
    val attemptedAmount: Money
) : PaymentEvent
```

判断基準:

| 質問 | 例外 | イベント |
|------|------|----------|
| この失敗を後で確認する必要があるか? | No | Yes |
| 規制やコンプライアンスで記録が必要か? | No | Yes |
| Sagaだけが失敗を気にするか? | Yes | No |
| Event Storeに残すと価値があるか? | No | Yes |

デフォルトは例外アプローチ。監査要件がある場合のみイベントを検討する。

## 抽象化レベルの評価

**条件分岐と抽象化**

分岐数だけで Strategy、State、ポリモーフィズムを選ばない。同じドメイン上の意味・契約・変更理由を持つ処理が2つ確認できたら、Aggregate、EventHandler、Projection など本来の所有者へ集約するか判断する。イベント種別や状態が異なる理由で変わる処理は分離を保つ。

**抽象度の不一致検出**

| パターン | 問題 | 修正案 |
|---------|------|--------|
| CommandHandlerにDB操作詳細 | 責務違反 | Repository層に分離 |
| EventHandlerにビジネスロジック | 責務違反 | ドメインサービスに抽出 |
| Aggregateに永続化処理 | レイヤー違反 | EventStore経由に変更 |
| Projectionに計算ロジック | 保守困難 | 専用サービスに抽出 |

良い抽象化の例:

```kotlin
// イベント種別による分岐の増殖（NG）
@EventHandler
fun on(event: DomainEvent) {
    when (event) {
        is OrderPlacedEvent -> handleOrderPlaced(event)
        is OrderConfirmedEvent -> handleOrderConfirmed(event)
        is OrderShippedEvent -> handleOrderShipped(event)
        // ...どんどん増える
    }
}

// イベントごとにハンドラを分離（OK）
@EventHandler
fun on(event: OrderPlacedEvent) { ... }

@EventHandler
fun on(event: OrderConfirmedEvent) { ... }

@EventHandler
fun on(event: OrderShippedEvent) { ... }
```

```kotlin
// 状態による分岐が複雑（NG）
fun process(command: ProcessCommand) {
    when (status) {
        PENDING -> if (command.type == "approve") { ... } else if (command.type == "reject") { ... }
        APPROVED -> if (command.type == "ship") { ... }
        // ...複雑化
    }
}

// State Patternで抽象化（OK）
sealed class OrderState {
    abstract fun handle(command: ProcessCommand): List<DomainEvent>
}
class PendingState : OrderState() {
    override fun handle(command: ProcessCommand) = when (command) {
        is ApproveCommand -> listOf(OrderApprovedEvent(...))
        is RejectCommand -> listOf(OrderRejectedEvent(...))
        else -> throw InvalidCommandException()
    }
}
```

## アンチパターンの観察

CQRS+ES では、CRUD の形だけをなぞる実装、意味のないイベントの乱発、イベント順序への暗黙依存、重要な事実の欠落、Aggregate への責務集中を設計上の確認対象とする。

## テスト戦略

レイヤーごとにテスト方針を分ける。

テストピラミッド:
```
        ┌─────────────┐
        │   E2E Test  │  ← 少数: 全体フロー確認
        ├─────────────┤
        │ Integration │  ← Command→Event→Projection→Query の連携確認
        ├─────────────┤
        │  Unit Test  │  ← 多数: 各レイヤー独立テスト
        └─────────────┘
```

Command側（Aggregate）:
```kotlin
// AggregateTestFixture使用
@Test
fun `確定コマンドでイベントが発行される`() {
    fixture
        .given(OrderPlacedEvent(...))
        .`when`(ConfirmOrderCommand(orderId, confirmedBy))
        .expectSuccessfulHandlerExecution()
        .expectEvents(OrderConfirmedEvent(...))
}
```

Query側:
```kotlin
// Read Model直接セットアップ + QueryGateway
@Test
fun `注文詳細が取得できる`() {
    // Given: Read Modelを直接セットアップ
    orderRepository.save(OrderEntity(...))

    // When: QueryGateway経由でクエリ実行
    val detail = queryGateway.query(GetOrderDetailQuery(orderId), ...).join()

    // Then
    assertEquals(expectedDetail, detail)
}
```

## 値オブジェクト設計

Aggregate とイベントの構成要素として値オブジェクトを使う。プリミティブ型（String, Int）で済ませない。

```kotlin
// 避ける例: プリミティブ型のまま
data class OrderPlacedEvent(
    val orderId: String,
    val categoryId: String,      // ただの文字列
    val from: LocalDateTime,     // 意味が不明確
    val to: LocalDateTime
)

// 例: 値オブジェクトで意味と制約を表現
data class OrderPlacedEvent(
    val orderId: String,
    val categoryId: CategoryId,
    val period: OrderPeriod
)
```

値オブジェクトの設計ルール:
- `data class` で equals/hashCode を自動生成（同値性で比較）
- `init` ブロックで不変条件を保証（生成時に必ず検証）
- ドメインロジック（計算）は含まない（純粋なデータホルダー）
- `@JsonValue` でシリアライゼーションを制御

```kotlin
// ID系: 単一値ラッパー
data class CategoryId(@get:JsonValue val value: String) {
    init {
        require(value.isNotBlank()) { "Category ID cannot be blank" }
    }
    override fun toString(): String = value
}

// 範囲系: 複数値の不変条件を保証
data class OrderPeriod(
    val from: LocalDateTime,
    val to: LocalDateTime
) {
    init {
        require(!to.isBefore(from)) { "終了日は開始日以降でなければなりません" }
    }
}

// メタ情報系: イベントペイロード内の付随情報
data class ApprovalInfo(
    val approvedBy: String,
    val approvalTime: LocalDateTime
)
```

## マスタデータ・設定値と CRUD の使い分け

CQRS+ES システム内でも、すべてをイベントソーシングで実装する必要はない。マスタデータ（参照データ）、管理設定、許可リストのように性質が単純なものは、通常の CRUD で実装した方がシンプルで保守しやすい。

ただし、「マスタデータだから CRUD」と機械的に判断しない。以下の基準で該当するものが多いほど CRUD が適している。逆に、CQRS+ES 採用判断の基準に該当する明示要件があれば、採用を検討する。

**CRUD で十分と判断する基準:**

| 観点 | CRUD 寄り | CQRS+ES 寄り |
|------|----------|-------------|
| ビジネス要件 | 「〜を管理したい」程度で特別な言及がない | 固有のビジネスルールや制約がある |
| ロジックの発展 | 単純な参照・更新で完結し、発展が見込めない | 状態遷移やライフサイクルが複雑化しうる |
| 変更履歴・監査 | 「いつ誰が変えたか」の追跡が不要 | 変更履歴の参照や監査証跡が必要 |
| ドメインイベント | この変更が他の集約やプロセスに影響しない | 変更が下流プロセスをトリガーする |
| 整合性の範囲 | 単体で完結し、他集約との整合性が不要 | 他の集約と整合性を保つ必要がある |
| 時点参照 | 「過去のある時点の状態」を問われない | 時点指定のクエリが必要 |

**典型的な CRUD 対象の例:**
- 都道府県・国コードなどのコードマスタ
- カテゴリ・タグなどの分類マスタ
- 設定値・定数テーブル
- IP許可リスト、機能フラグ、通知設定などの現在値ベースの管理設定

**CQRS+ES が必要と判断できる例:**
- 商品マスタだが、価格変更履歴の追跡が必要
- 組織マスタだが、変更時に権限の再計算をトリガーする
- 取引先マスタだが、与信審査の状態遷移がある

```kotlin
// CRUD で十分: 単純なカテゴリマスタ
@Entity
data class Category(
    @Id val categoryId: String,
    val name: String,
    val displayOrder: Int
)

// CQRS+ES が適切: 価格変更履歴の追跡が必要な商品
data class Product(
    val productId: String,
    val currentPrice: Money
) {
    fun changePrice(newPrice: Money, reason: String): PriceChangedEvent {
        require(newPrice.amount > BigDecimal.ZERO) { "価格は正の値でなければなりません" }
        return PriceChangedEvent(productId, currentPrice, newPrice, reason)
    }
}
```

CRUD で実装する場合も、CQRS+ES システム内の他集約からは ID 参照で利用する。CRUD エンティティが集約の内部状態を直接参照しない点は同じ。

## インフラ層

確認事項:
- イベントストアの選択は適切か
- メッセージング基盤は要件を満たすか
- スナップショット戦略は定義されているか
- イベントのシリアライズ形式は適切か
