# 집금 배치 재설계 — 통합 구현 지침서

> **작성일**: 2026-03-22
> **DDL 버전**: v1.8
> **설계 문서**: `COLLECTION_BATCH_REDESIGN.md` (설계 확정 완료)
> **범위**: Spring Boot (common, core, open-api, admin-api) + Node.js (relayer-api, common)

---

## 목차

1. [DDL 마이그레이션](#1-ddl-마이그레이션)
2. [Spring Boot — common 모듈](#2-spring-boot--common-모듈)
3. [Spring Boot — core 모듈](#3-spring-boot--core-모듈)
4. [Spring Boot — open-api 모듈 (Webhook)](#4-spring-boot--open-api-모듈-webhook)
5. [Spring Boot — admin-api 모듈](#5-spring-boot--admin-api-모듈)
6. [Node.js — common 패키지](#6-nodejs--common-패키지)
7. [Node.js — relayer-api 패키지](#7-nodejs--relayer-api-패키지)
8. [구현 순서 체크리스트](#8-구현-순서-체크리스트)

---

## 1. DDL 마이그레이션

> DDL v1.8은 `CRYPTOMENTS_V2_DDL.sql`에 이미 반영됨. 아래는 기존 DB에 적용할 마이그레이션 SQL.

```sql
-- ══════════════════════════════════════════════════
-- DDL v1.8 마이그레이션 — 기존 DB 적용용
-- 실행 전: v1.7 collection_queue의 기존 데이터 백업 권장
-- ══════════════════════════════════════════════════

-- 1) collection_batches 신규 테이블
CREATE TABLE collection_batches (
    id BIGINT AUTO_INCREMENT PRIMARY KEY
        COMMENT 'PK',
    batch_code VARCHAR(50) NOT NULL
        COMMENT '배치 고유 코드 — cb_{YYMM}_{random8}',
    wallet_address_id BIGINT NOT NULL
        COMMENT 'wallet_addresses.id — 집금 원천 HOT 지갑',
    partner_id BIGINT NOT NULL
        COMMENT 'partners.id',
    network_id BIGINT NOT NULL
        COMMENT 'blockchain_networks.id',
    currency_id BIGINT NOT NULL
        COMMENT 'currencies.id',
    onchain_balance DECIMAL(36,18) NOT NULL
        COMMENT '집금 시점 온체인 잔액 조회 결과',
    collected_amount DECIMAL(36,18) NOT NULL
        COMMENT '실제 집금된 금액 (= 온체인 전액)',
    queue_count INT NOT NULL DEFAULT 0
        COMMENT '이 배치에 포함된 collection_queue 건수',
    estimated_gas_usd DECIMAL(18,8)
        COMMENT '예상 가스비 (USD)',
    tx_hash VARCHAR(255)
        COMMENT '집금 TX 해시',
    relayer_address_id BIGINT
        COMMENT '집금 실행 Relayer',
    status VARCHAR(20) NOT NULL DEFAULT 'PENDING'
        COMMENT 'PENDING → BROADCASTING → CONFIRMED / FAILED / STALE',
    error_message TEXT
        COMMENT '실패 시 에러',
    created_at DATETIME(6) DEFAULT CURRENT_TIMESTAMP(6),
    updated_at DATETIME(6) DEFAULT CURRENT_TIMESTAMP(6) ON UPDATE CURRENT_TIMESTAMP(6),
    UNIQUE KEY uk_batch_code (batch_code),
    KEY idx_wallet_status (wallet_address_id, status),
    KEY idx_status (status, created_at),
    KEY idx_tx_hash (tx_hash),
    KEY idx_relayer_status (relayer_address_id, status)
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COLLATE=utf8mb4_unicode_ci
  COMMENT='집금 배치 — 지갑별 물리 TX 단위 (v1.8 신규)';

-- 2) collection_queue 변경 — 컬럼 추가 먼저
ALTER TABLE collection_queue
    ADD COLUMN deposit_id BIGINT NOT NULL DEFAULT 0
        COMMENT 'deposits.id — 이 집금 건의 원본 입금'
        AFTER collection_code,
    ADD COLUMN batch_id BIGINT
        COMMENT 'collection_batches.id — 실제 집금된 배치'
        AFTER amount;

-- 3) 기존 데이터 deposit_id 역매핑 (collection_queue_id로 연결)
UPDATE collection_queue cq
  JOIN deposits d ON d.collection_queue_id = cq.id
   SET cq.deposit_id = d.id
 WHERE cq.deposit_id = 0;

-- 4) 기존 CONFIRMED 건 → COLLECTED 상태로 전환 (batch 없이)
UPDATE collection_queue SET status = 'COLLECTED'
 WHERE status IN ('CONFIRMED', 'BROADCASTING', 'COLLECTING');

-- 5) collection_queue 불필요 컬럼 제거
ALTER TABLE collection_queue
    DROP COLUMN collection_mode,
    DROP COLUMN scheduled_at,
    DROP COLUMN tx_hash,
    DROP COLUMN relayer_address_id,
    DROP INDEX idx_scheduled;

-- 6) 신규 인덱스
ALTER TABLE collection_queue
    ADD KEY idx_deposit (deposit_id),
    ADD KEY idx_batch (batch_id),
    ADD KEY idx_wallet_currency_status (wallet_address_id, currency_id, status);

-- 7) system_settings 추가
INSERT INTO system_settings (setting_key, setting_value, description, value_type) VALUES
('collection.schedule_interval_hours', '1',   '집금 배치 실행 주기 (시간)',                  'NUMBER'),
('collection.gas_threshold_ratio',     '0.1', '가스비/집금액 비율 임계값 (초과 시 SKIP)',       'NUMBER'),
('stale_tx.check_interval_minutes',    '60',  'Stale TX 점검 주기 (분)',                    'NUMBER'),
('stale_tx.threshold_minutes',         '60',  'BROADCASTING 후 미확정 임계 시간 (분)',        'NUMBER');

-- 기존 collection.min_amount_usdt → collection.min_amount_usd 변경
UPDATE system_settings
   SET setting_key = 'collection.min_amount_usd',
       description = '최소 집금 금액 (USD, 미만 시 SKIP)'
 WHERE setting_key = 'collection.min_amount_usdt';
```

---

## 2. Spring Boot — common 모듈

### 2-1. CollectionBatchStatus 신규 enum

> 파일: `common/src/main/java/com/cryptoments/common/enums/CollectionBatchStatus.java`

```java
package com.cryptoments.common.enums;

public enum CollectionBatchStatus {
    /** 배치 생성 — 아직 TX 미전송 */
    PENDING,
    /** TX 브로드캐스트 완료 — Webhook 확정 대기 */
    BROADCASTING,
    /** 블록체인 컨펌 완료 */
    CONFIRMED,
    /** TX 실패 (revert 또는 broadcast 실패) */
    FAILED,
    /** 1시간+ 미확정 — 안전망 스케줄러가 마킹 */
    STALE;
}
```

### 2-2. CollectionStatus enum 수정

> 파일: `common/src/main/java/com/cryptoments/common/enums/CollectionStatus.java`
> **기존 6개 → 2개로 단순화**

```java
package com.cryptoments.common.enums;

public enum CollectionStatus {
    /** 대기 — 배치 집금 대기 중 */
    QUEUED,
    /** 집금 완료 — batch_id 할당됨 */
    COLLECTED;
}
```

### 2-3. CollectionMode enum 삭제

> 파일: `common/src/main/java/com/cryptoments/common/enums/CollectionMode.java`
> **삭제** — 배치 스케줄이 결정하므로 IMMEDIATE/SCHEDULED 구분 불필요

### 2-4. CollectionQueue entity 수정

> 파일: `common/src/main/java/com/cryptoments/common/entity/CollectionQueue.java`

```java
package com.cryptoments.common.entity;

import com.cryptoments.common.enums.CollectionStatus;
import lombok.*;
import one.axim.framework.mybatis.annotation.XColumn;
import one.axim.framework.mybatis.annotation.XEntity;

import java.math.BigDecimal;
import java.time.LocalDateTime;

@XEntity("collection_queue")
@Getter @Setter @Builder(toBuilder = true)
@NoArgsConstructor @AllArgsConstructor
public class CollectionQueue {

    /** PK */
    @XColumn(value = "id", isPrimaryKey = true, isAutoIncrement = true)
    private Long id;

    /** 집금 고유 코드 — col_{YYMM}_{random8} */
    private String collectionCode;

    /** deposits.id — 이 집금 건의 원본 입금 */
    private Long depositId;

    /** wallet_addresses.id — 집금 대상 HOT 지갑 */
    private Long walletAddressId;

    /** partners.id 참조 */
    private Long partnerId;

    /** blockchain_networks.id 참조 */
    private Long networkId;

    /** currencies.id 참조 */
    private Long currencyId;

    /** 해당 입금 건의 집금 대상 금액 (추적용) */
    private BigDecimal amount;

    /** collection_batches.id — 실제 집금된 배치 (Webhook CONFIRMED 시 할당) */
    private Long batchId;

    /** 집금 상태 — QUEUED / COLLECTED */
    private CollectionStatus status;

    /** 참고용 에러 메시지 */
    private String errorMessage;

    /** 재시도 횟수 */
    @Builder.Default
    private Integer retryCount = 0;

    // ★ 제거됨: collectionMode, scheduledAt, txHash, relayerAddressId

    /** 생성 시각 */
    @XColumn(value = "created_at", insert = false, update = false)
    private LocalDateTime createdAt;

    /** 수정 시각 */
    @XColumn(value = "updated_at", insert = false, update = false)
    private LocalDateTime updatedAt;
}
```

### 2-5. CollectionBatch entity 신규

> 파일: `common/src/main/java/com/cryptoments/common/entity/CollectionBatch.java`

```java
package com.cryptoments.common.entity;

import com.cryptoments.common.enums.CollectionBatchStatus;
import lombok.*;
import one.axim.framework.mybatis.annotation.XColumn;
import one.axim.framework.mybatis.annotation.XEntity;

import java.math.BigDecimal;
import java.time.LocalDateTime;

@XEntity("collection_batches")
@Getter @Setter @Builder(toBuilder = true)
@NoArgsConstructor @AllArgsConstructor
public class CollectionBatch {

    /** PK */
    @XColumn(value = "id", isPrimaryKey = true, isAutoIncrement = true)
    private Long id;

    /** 배치 고유 코드 — cb_{YYMM}_{random8} */
    private String batchCode;

    /** wallet_addresses.id — 집금 원천 HOT 지갑 */
    private Long walletAddressId;

    /** partners.id */
    private Long partnerId;

    /** blockchain_networks.id */
    private Long networkId;

    /** currencies.id */
    private Long currencyId;

    /** 집금 시점 온체인 잔액 조회 결과 */
    private BigDecimal onchainBalance;

    /** 실제 집금된 금액 (= 온체인 전액) */
    private BigDecimal collectedAmount;

    /** 이 배치에 포함된 collection_queue 건수 */
    @Builder.Default
    private Integer queueCount = 0;

    /** 예상 가스비 (USD) */
    private BigDecimal estimatedGasUsd;

    /** 집금 TX 해시 */
    private String txHash;

    /** wallet_addresses.id — 집금 실행 Relayer */
    private Long relayerAddressId;

    /** 배치 상태 — PENDING / BROADCASTING / CONFIRMED / FAILED / STALE */
    private CollectionBatchStatus status;

    /** 실패 시 에러 */
    private String errorMessage;

    /** 생성 시각 */
    @XColumn(value = "created_at", insert = false, update = false)
    private LocalDateTime createdAt;

    /** 수정 시각 */
    @XColumn(value = "updated_at", insert = false, update = false)
    private LocalDateTime updatedAt;
}
```

### 2-6. CollectionQueueRepository 수정

> 파일: `common/src/main/java/com/cryptoments/common/repository/CollectionQueueRepository.java`

```java
package com.cryptoments.common.repository;

import com.cryptoments.common.entity.CollectionQueue;
import com.cryptoments.common.enums.CollectionStatus;
import one.axim.framework.mybatis.repository.IXRepository;
import one.axim.framework.mybatis.annotation.XRepository;

import java.util.List;

@XRepository
public interface CollectionQueueRepository extends IXRepository<Long, CollectionQueue> {

    CollectionQueue findByCollectionCode(String collectionCode);

    List<CollectionQueue> findByWalletAddressIdAndStatus(Long walletAddressId, CollectionStatus status);

    List<CollectionQueue> findByPartnerIdAndStatus(Long partnerId, CollectionStatus status);

    List<CollectionQueue> findByDepositId(Long depositId);

    List<CollectionQueue> findByBatchId(Long batchId);

    // ★ 제거됨: findByTxHash() — tx_hash 컬럼이 collection_batches로 이동
}
```

### 2-7. CollectionBatchRepository 신규

> 파일: `common/src/main/java/com/cryptoments/common/repository/CollectionBatchRepository.java`

```java
package com.cryptoments.common.repository;

import com.cryptoments.common.entity.CollectionBatch;
import com.cryptoments.common.enums.CollectionBatchStatus;
import one.axim.framework.mybatis.repository.IXRepository;
import one.axim.framework.mybatis.annotation.XRepository;

import java.util.List;

@XRepository
public interface CollectionBatchRepository extends IXRepository<Long, CollectionBatch> {

    CollectionBatch findByBatchCode(String batchCode);

    CollectionBatch findByTxHash(String txHash);

    List<CollectionBatch> findByWalletAddressIdAndStatus(Long walletAddressId, CollectionBatchStatus status);

    List<CollectionBatch> findByStatus(CollectionBatchStatus status);

    List<CollectionBatch> findByRelayerAddressIdAndStatus(Long relayerAddressId, CollectionBatchStatus status);
}
```

### 2-8. WithdrawalStatus enum — STALE 추가

> 파일: `common/src/main/java/com/cryptoments/common/enums/WithdrawalStatus.java`
> 기존 enum에 `STALE` 추가:

```java
// 기존 값 뒤에 추가:
    /** 1시간+ 미확정 — 안전망 스케줄러가 마킹 */
    STALE;
```

### 2-9. ErrorCodes — 집금 배치 관련 추가

> 파일: `common/src/main/java/com/cryptoments/common/exception/ErrorCodes.java`
> 기존 상수 뒤에 추가:

```java
// ── 집금 배치 ──
public static final ErrorCode COLLECTION_BATCH_NOT_FOUND =
    new ErrorCode("601", "집금 배치를 찾을 수 없습니다.");
public static final ErrorCode COLLECTION_BATCH_ALREADY_PROCESSED =
    new ErrorCode("602", "이미 처리된 집금 배치입니다.");
```

---

## 3. Spring Boot — core 모듈

### 3-1. DepositService.enqueueCollection() 수정

> 파일: `core/src/main/java/com/cryptoments/core/deposit/DepositService.java`
> 메서드 시작 위치: 약 L473

**변경 요약**:
- `collectionMode` 제거 (import `CollectionMode` 제거)
- `depositId` 필드 추가
- deposit 상태는 COLLECTING 유지 (기존과 동일)

```java
// ★ import 제거: import com.cryptoments.common.enums.CollectionMode;

/**
 * 집금 enqueue (collection_queue INSERT → Node.js BatchRunner가 DB 폴링).
 * v1.8: 배치 집금 재설계 — deposit_id 추가, collectionMode 제거
 */
public CollectionQueue enqueueCollection(Long depositId) {
    Deposit deposit = depositRepository.findOne(depositId);
    if (deposit == null) {
        throw new NotFoundException(ErrorCodes.DEPOSIT_NOT_FOUND);
    }

    if (deposit.getWalletAddressId() == null) {
        throw new ConflictException(ErrorCodes.TRANSACTION_STATUS_INVALID,
                "지갑 주소가 없어 집금 불필요");
    }

    String collectionCode = generateCode("col");

    CollectionQueue queue = CollectionQueue.builder()
            .collectionCode(collectionCode)
            .depositId(depositId)                   // ★ 신규
            .walletAddressId(deposit.getWalletAddressId())
            .partnerId(deposit.getPartnerId())
            .networkId(deposit.getNetworkId())
            .currencyId(deposit.getCurrencyId())
            .amount(deposit.getAmount())
            // ★ collectionMode 제거
            .status(CollectionStatus.QUEUED)
            .build();

    Long id = collectionQueueRepository.save(queue);
    queue.setId(id);

    Deposit updated = deposit.toBuilder()
            .collectionQueueId(id)
            .status(DepositStatus.COLLECTING)
            .build();
    depositRepository.modify(updated);

    log.info("집금 enqueue: collectionId={}, depositId={}, amount={}",
            id, depositId, deposit.getAmount());
    return queue;
}
```

### 3-2. DepositService — onCollectionConfirmed 신규 메서드

> deposit의 집금 완료 상태 전이를 위한 메서드 추가:

```java
/**
 * 집금 배치 확정 시 deposit 상태 업데이트.
 * Webhook COLLECTION_CONFIRM → batch CONFIRMED → 이 메서드 호출.
 */
public void onCollectionConfirmed(Long depositId) {
    Deposit deposit = depositRepository.findOne(depositId);
    if (deposit == null) {
        log.warn("집금 확정 대상 deposit 없음: depositId={}", depositId);
        return;
    }
    if (deposit.getStatus() != DepositStatus.COLLECTING) {
        log.info("deposit 상태 불일치 (skip): depositId={}, status={}",
                depositId, deposit.getStatus());
        return;
    }

    Deposit updated = deposit.toBuilder()
            .status(DepositStatus.COLLECTED)
            .build();
    depositRepository.modify(updated);
    log.info("입금 집금 완료: depositId={}", depositId);
}
```

> **주의**: `DepositStatus`에 `COLLECTED` 값이 없으면 추가 필요. 기존에 `COLLECTING` 다음 상태를 확인할 것.

---

## 4. Spring Boot — open-api 모듈 (Webhook)

> **핵심 변경**: collection_queue.tx_hash → collection_batches.tx_hash로 전환

### 4-1. BusinessEventClassifier 수정

> 파일: `open-api/src/main/java/com/cryptoments/webhook/service/BusinessEventClassifier.java`

**변경 요약**:
- `CollectionQueueRepository` → `CollectionBatchRepository` 교체
- `findByTxHash()` 호출 대상을 collection_batches로 변경

```java
package com.cryptoments.webhook.service;

import com.cryptoments.common.entity.WalletAddress;
import com.cryptoments.common.repository.CollectionBatchRepository;    // ★ 변경
import com.cryptoments.common.repository.WalletAddressRepository;
import com.cryptoments.common.repository.WithdrawalRepository;
import com.cryptoments.webhook.dto.TransactionWebhookDTO;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.stereotype.Service;

/**
 * Webhook TRANSFER 이벤트를 비즈니스 타입으로 분류한다.
 *
 * <p>판별 순서:
 * <ol>
 *   <li>status == "FAILED" → TX_FAILED</li>
 *   <li>양쪽 다 우리 주소? → COLLECTION_CONFIRM (내부 이체) 또는 WITHDRAWAL_CONFIRM</li>
 *   <li>to만 우리 주소? → DEPOSIT 또는 COLLECTION_CONFIRM (MASTER)</li>
 *   <li>from만 우리 주소? → WITHDRAWAL_CONFIRM 또는 COLLECTION_CONFIRM</li>
 *   <li>매칭 안 됨 → UNKNOWN</li>
 * </ol>
 *
 * <p>v1.8: collection_queue → collection_batches로 TX 매칭 변경
 */
@Service
public class BusinessEventClassifier {

    private static final Logger log = LoggerFactory.getLogger(BusinessEventClassifier.class);

    private final WalletAddressRepository walletAddressRepository;
    private final WithdrawalRepository withdrawalRepository;
    private final CollectionBatchRepository collectionBatchRepository;    // ★ 변경

    public BusinessEventClassifier(WalletAddressRepository walletAddressRepository,
                                    WithdrawalRepository withdrawalRepository,
                                    CollectionBatchRepository collectionBatchRepository) {  // ★ 변경
        this.walletAddressRepository = walletAddressRepository;
        this.withdrawalRepository = withdrawalRepository;
        this.collectionBatchRepository = collectionBatchRepository;      // ★ 변경
    }

    public String classify(TransactionWebhookDTO dto, Long networkId) {
        // 1. TX 실패
        if ("FAILED".equals(dto.getStatus())) {
            return "TX_FAILED";
        }

        WalletAddress toWallet = walletAddressRepository
                .findByAddressAndNetworkId(dto.getTo(), networkId);
        WalletAddress fromWallet = walletAddressRepository
                .findByAddressAndNetworkId(dto.getFrom(), networkId);

        boolean toIsOurs = (toWallet != null);
        boolean fromIsOurs = (fromWallet != null);

        // 2. 양쪽 다 우리 주소
        if (toIsOurs && fromIsOurs) {
            // ★ collection_batches.tx_hash로 매칭
            if (collectionBatchRepository.findByTxHash(dto.getTxHash()) != null) {
                return "COLLECTION_CONFIRM";
            }
            if (!withdrawalRepository.findByTxHash(dto.getTxHash()).isEmpty()) {
                return "WITHDRAWAL_CONFIRM";
            }
            String toWalletType = toWallet.getWalletType().name();
            if ("HOT".equals(toWalletType) || "POOL".equals(toWalletType) || "MASTER".equals(toWalletType)) {
                return "DEPOSIT";
            }
            log.info("인프라 자금 이동: txHash={}, from={}, to={}", dto.getTxHash(), dto.getFrom(), dto.getTo());
            return "FUNDING";
        }

        // 3. to만 우리 주소 — 외부 → 우리: 입금
        if (toIsOurs) {
            String walletType = toWallet.getWalletType().name();
            return switch (walletType) {
                case "HOT", "POOL" -> "DEPOSIT";
                case "MASTER" -> {
                    // ★ collection_batches.tx_hash로 매칭
                    if (collectionBatchRepository.findByTxHash(dto.getTxHash()) != null) {
                        yield "COLLECTION_CONFIRM";
                    }
                    yield "DEPOSIT";
                }
                default -> "UNKNOWN_INBOUND";
            };
        }

        // 4. from만 우리 주소 — 우리 → 외부: 출금 확인
        if (fromIsOurs) {
            if (!withdrawalRepository.findByTxHash(dto.getTxHash()).isEmpty()) {
                return "WITHDRAWAL_CONFIRM";
            }
            // ★ collection_batches.tx_hash로 매칭
            if (collectionBatchRepository.findByTxHash(dto.getTxHash()) != null) {
                return "COLLECTION_CONFIRM";
            }
            return "UNKNOWN_OUTBOUND";
        }

        return "UNKNOWN";
    }
}
```

### 4-2. WebhookProcessingService 수정

> 파일: `open-api/src/main/java/com/cryptoments/webhook/service/WebhookProcessingService.java`

**변경 요약**:
- `CollectionBatchRepository` 주입 추가
- `confirmCollection()` 전면 재작성 — batch CONFIRMED + queue 일괄 COLLECTED
- `handleFailedTx()` 집금 부분 — batch FAILED, queue는 QUEUED 유지

#### 4-2-a. import + 필드 추가

```java
// ★ 추가 import
import com.cryptoments.common.entity.CollectionBatch;
import com.cryptoments.common.enums.CollectionBatchStatus;
import com.cryptoments.common.repository.CollectionBatchRepository;

// ★ 필드 추가 (기존 collectionQueueRepository는 유지)
private final CollectionBatchRepository collectionBatchRepository;
```

#### 4-2-b. 생성자 — CollectionBatchRepository 파라미터 추가

```java
public WebhookProcessingService(BlockchainNetworkRepository blockchainNetworkRepository,
                                 BusinessEventClassifier businessEventClassifier,
                                 WebhookBalanceSyncService webhookBalanceSyncService,
                                 DepositRepository depositRepository,
                                 WithdrawalRepository withdrawalRepository,
                                 CollectionQueueRepository collectionQueueRepository,
                                 CollectionBatchRepository collectionBatchRepository,  // ★ 추가
                                 WalletAddressRepository walletAddressRepository,
                                 CurrencyRepository currencyRepository,
                                 PartnerRepository partnerRepository,
                                 DepositService depositService) {
    // ... 기존 할당 유지 ...
    this.collectionBatchRepository = collectionBatchRepository;  // ★ 추가
}
```

#### 4-2-c. confirmCollection() 전면 재작성 (L250-265 대체)

```java
/**
 * 집금 배치 확정: collection_batches CONFIRMED + collection_queue 일괄 COLLECTED.
 *
 * v1.8: tx_hash → collection_batches 매칭.
 * batch의 wallet_address_id + currency_id가 같은 QUEUED 건 전부 COLLECTED.
 */
private void confirmCollection(TransactionWebhookDTO dto) {
    // 1. tx_hash로 batch 조회
    CollectionBatch batch = collectionBatchRepository.findByTxHash(dto.getTxHash());
    if (batch == null) {
        log.warn("집금 배치 매칭 실패: txHash={}", dto.getTxHash());
        return;
    }
    if (batch.getStatus() != CollectionBatchStatus.BROADCASTING
            && batch.getStatus() != CollectionBatchStatus.STALE) {
        log.info("집금 배치 이미 처리됨 (멱등): batchId={}, status={}", batch.getId(), batch.getStatus());
        return;
    }

    // 2. batch → CONFIRMED
    CollectionBatch updated = batch.toBuilder()
            .status(CollectionBatchStatus.CONFIRMED)
            .build();
    collectionBatchRepository.modify(updated);

    // 3. 해당 지갑+통화의 QUEUED 건 전부 → COLLECTED (batch_id 할당)
    List<CollectionQueue> queuedList = collectionQueueRepository
            .findByWalletAddressIdAndStatus(batch.getWalletAddressId(), CollectionStatus.QUEUED);

    int collectedCount = 0;
    for (CollectionQueue q : queuedList) {
        // 같은 currency만 (지갑에 여러 토큰 가능)
        if (!q.getCurrencyId().equals(batch.getCurrencyId())) {
            continue;
        }
        CollectionQueue collectedQueue = q.toBuilder()
                .batchId(batch.getId())
                .status(CollectionStatus.COLLECTED)
                .build();
        collectionQueueRepository.modify(collectedQueue);

        // deposit 상태 업데이트
        depositService.onCollectionConfirmed(q.getDepositId());
        collectedCount++;
    }

    // 4. TODO: gas_cost_records 기록 (receipt에서 가스비 추출)
    // gasCostService.recordCollectionGas(batch, dto);

    log.info("집금 배치 확정: batchId={}, txHash={}, collectedCount={}",
            batch.getId(), dto.getTxHash(), collectedCount);
}
```

#### 4-2-d. handleFailedTx() 집금 부분 수정 (L282-289 대체)

```java
/**
 * TX 실패 처리: 관련 withdrawal/collection_batch 상태를 FAILED로 업데이트.
 * v1.8: collection_queue는 QUEUED 유지 — 다음 배치에서 재시도.
 */
private void handleFailedTx(TransactionWebhookDTO dto) {
    log.error("TX 실패 감지: txHash={}, from={}, to={}", dto.getTxHash(), dto.getFrom(), dto.getTo());

    // 출금 실패 처리 (기존과 동일)
    List<Withdrawal> withdrawals = withdrawalRepository.findByTxHash(dto.getTxHash());
    for (Withdrawal w : withdrawals) {
        Withdrawal updated = w.toBuilder()
                .status(WithdrawalStatus.FAILED)
                .build();
        withdrawalRepository.modify(updated);
        log.warn("출금 실패 처리: withdrawalId={}", w.getId());
    }

    // ★ 집금 실패 처리 — batch만 FAILED, queue는 QUEUED 유지
    CollectionBatch batch = collectionBatchRepository.findByTxHash(dto.getTxHash());
    if (batch != null && batch.getStatus() == CollectionBatchStatus.BROADCASTING) {
        CollectionBatch updated = batch.toBuilder()
                .status(CollectionBatchStatus.FAILED)
                .errorMessage("TX FAILED: " + dto.getTxHash())
                .build();
        collectionBatchRepository.modify(updated);
        log.warn("집금 배치 실패 처리: batchId={} (queue는 QUEUED 유지 — 다음 배치 재시도)", batch.getId());
        // ★ collection_queue는 건드리지 않음 — QUEUED 상태 유지
    }
}
```

---

## 5. Spring Boot — admin-api 모듈

### 5-1. CollectionSearchMapper 수정

> 파일: `admin-api/src/main/java/com/cryptoments/adminapi/mapper/CollectionSearchMapper.java`

**변경 요약**:
- `collection_mode` 필터 제거
- `collection_batches` LEFT JOIN 추가 (배치 정보 포함)
- 상세 조회에 batch 정보 추가

```java
@Mapper
public interface CollectionSearchMapper {

    /**
     * 집금 목록 (v1.8: collection_batches JOIN 추가).
     */
    @Select("<script>" +
            "SELECT cq.*, " +
            "  p.name AS partner_name, " +
            "  bn.chain_symbol AS network_code, " +
            "  c.symbol AS currency_code, " +
            "  cb.batch_code, " +
            "  cb.tx_hash AS batch_tx_hash, " +
            "  cb.collected_amount AS batch_collected_amount, " +
            "  cb.status AS batch_status" +
            "  FROM collection_queue cq" +
            "  LEFT JOIN partners p ON cq.partner_id = p.id" +
            "  LEFT JOIN blockchain_networks bn ON cq.network_id = bn.id" +
            "  LEFT JOIN currencies c ON cq.currency_id = c.id" +
            "  LEFT JOIN collection_batches cb ON cq.batch_id = cb.id" +
            " <where>" +
            "  <if test='search.collectionCode != null'> AND cq.collection_code = #{search.collectionCode}</if>" +
            "  <if test='search.partnerId != null'> AND cq.partner_id = #{search.partnerId}</if>" +
            "  <if test='search.networkId != null'> AND cq.network_id = #{search.networkId}</if>" +
            "  <if test='search.status != null'> AND cq.status = #{search.status}</if>" +
            "  <if test='search.batchCode != null'> AND cb.batch_code = #{search.batchCode}</if>" +
            "  <if test='search.fromDate != null'> AND cq.created_at &gt;= #{search.fromDate}</if>" +
            "  <if test='search.toDate != null'> AND cq.created_at &lt; DATE_ADD(#{search.toDate}, INTERVAL 1 DAY)</if>" +
            " </where>" +
            "</script>")
    XPage<CollectionListResponse> searchCollections(XPagination pagination,
                                                     @Param("search") CollectionSearchRequest search,
                                                     Class<?> cls);

    /**
     * 집금 단건 상세 (v1.8: batch 정보 포함).
     */
    @Select("SELECT cq.*," +
            "  p.name AS partner_name," +
            "  bn.chain_symbol AS network_code," +
            "  c.symbol AS currency_code," +
            "  wa.address AS wallet_address," +
            "  cb.batch_code," +
            "  cb.tx_hash AS batch_tx_hash," +
            "  cb.collected_amount AS batch_collected_amount," +
            "  cb.onchain_balance AS batch_onchain_balance," +
            "  cb.estimated_gas_usd AS batch_gas_usd," +
            "  cb.status AS batch_status," +
            "  cb.created_at AS batch_created_at" +
            "  FROM collection_queue cq" +
            "  LEFT JOIN partners p ON cq.partner_id = p.id" +
            "  LEFT JOIN blockchain_networks bn ON cq.network_id = bn.id" +
            "  LEFT JOIN currencies c ON cq.currency_id = c.id" +
            "  LEFT JOIN wallet_addresses wa ON cq.wallet_address_id = wa.id" +
            "  LEFT JOIN collection_batches cb ON cq.batch_id = cb.id" +
            "  WHERE cq.id = #{id}")
    CollectionDetailResponse findCollectionDetail(@Param("id") Long id);
}
```

### 5-2. CollectionSearchRequest DTO 수정

> `collectionMode` 필드 제거, `batchCode` 필드 추가:

```java
// 제거: private String collectionMode;
// 추가:
/** 배치 코드 검색 */
private String batchCode;
```

### 5-3. CollectionListResponse / CollectionDetailResponse DTO 수정

> batch 관련 필드 추가:

```java
// CollectionListResponse에 추가:
/** 배치 코드 */
private String batchCode;
/** 배치 TX 해시 */
private String batchTxHash;
/** 배치 집금액 */
private BigDecimal batchCollectedAmount;
/** 배치 상태 */
private String batchStatus;

// CollectionDetailResponse에 추가 (위 4개 + 아래):
/** 배치 온체인 잔액 */
private BigDecimal batchOnchainBalance;
/** 배치 예상 가스비 */
private BigDecimal batchGasUsd;
/** 배치 생성 시각 */
private LocalDateTime batchCreatedAt;
```

### 5-4. CollectionManagementService — retryCollection 조정

> 기존 `retryCollection`은 FAILED → QUEUED인데, v1.8에서 queue는 FAILED가 안 되므로 의미 변경 필요.
> **batch 재시도**: FAILED batch의 queue 건들은 이미 QUEUED → 별도 재시도 불필요 (다음 배치가 자동 처리).
> 다만 관리자가 STALE batch를 수동 FAILED 처리하는 기능은 필요할 수 있음.

```java
/**
 * 배치 수동 실패 처리 (STALE → FAILED).
 * STALE 배치를 관리자가 수동으로 FAILED 처리 — queue는 QUEUED 유지.
 */
public CollectionBatch failStaleBatch(Long batchId, Long adminId) {
    CollectionBatch batch = collectionBatchRepository.findOne(batchId);
    if (batch == null) {
        throw new NotFoundException(ErrorCodes.COLLECTION_BATCH_NOT_FOUND);
    }
    if (batch.getStatus() != CollectionBatchStatus.STALE) {
        throw new ConflictException(ErrorCodes.COLLECTION_BATCH_ALREADY_PROCESSED,
                "STALE 상태만 수동 실패 처리 가능 (현재: " + batch.getStatus() + ")");
    }

    CollectionBatch updated = batch.toBuilder()
            .status(CollectionBatchStatus.FAILED)
            .errorMessage("관리자 수동 FAILED 처리")
            .build();
    collectionBatchRepository.modify(updated);

    auditLogService.log(adminId, AuditAction.FAIL_STALE_BATCH, batchId,
            Map.of("previousStatus", "STALE", "newStatus", "FAILED"));

    return updated;
}
```

---

## 6. Node.js — common 패키지

### 6-1. CollectionBatchRepo 신규

> 파일: `node-service/packages/common/src/db/repositories/CollectionBatchRepo.ts`

```typescript
import { getPool } from '../pool';

export interface CollectionBatch {
  id: number;
  batch_code: string;
  wallet_address_id: number;
  partner_id: number;
  network_id: number;
  currency_id: number;
  onchain_balance: string;        // DECIMAL(36,18)
  collected_amount: string;
  queue_count: number;
  estimated_gas_usd: string | null;
  tx_hash: string | null;
  relayer_address_id: number | null;
  status: 'PENDING' | 'BROADCASTING' | 'CONFIRMED' | 'FAILED' | 'STALE';
  error_message: string | null;
  created_at: Date;
  updated_at: Date;
}

export interface WalletGroup {
  wallet_address_id: number;
  network_id: number;
  currency_id: number;
  partner_id: number;
  queue_count: number;
}

export const collectionBatchRepo = {

  async insert(data: Omit<CollectionBatch, 'id' | 'created_at' | 'updated_at'>): Promise<number> {
    const pool = getPool();
    const [result] = await pool.execute(
      `INSERT INTO collection_batches
       (batch_code, wallet_address_id, partner_id, network_id, currency_id,
        onchain_balance, collected_amount, queue_count, estimated_gas_usd,
        tx_hash, relayer_address_id, status, error_message)
       VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)`,
      [
        data.batch_code, data.wallet_address_id, data.partner_id,
        data.network_id, data.currency_id,
        data.onchain_balance, data.collected_amount, data.queue_count,
        data.estimated_gas_usd,
        data.tx_hash, data.relayer_address_id, data.status, data.error_message,
      ]
    );
    return (result as any).insertId;
  },

  async findByTxHash(txHash: string): Promise<CollectionBatch | null> {
    const pool = getPool();
    const [rows] = await pool.execute(
      'SELECT * FROM collection_batches WHERE tx_hash = ?',
      [txHash]
    );
    return (rows as CollectionBatch[])[0] || null;
  },

  async updateStatus(id: number, status: string, errorMessage?: string): Promise<void> {
    const pool = getPool();
    if (errorMessage !== undefined) {
      await pool.execute(
        'UPDATE collection_batches SET status = ?, error_message = ? WHERE id = ?',
        [status, errorMessage, id]
      );
    } else {
      await pool.execute(
        'UPDATE collection_batches SET status = ? WHERE id = ?',
        [status, id]
      );
    }
  },

  async update(id: number, data: Partial<CollectionBatch>): Promise<void> {
    const pool = getPool();
    const fields: string[] = [];
    const values: any[] = [];
    for (const [key, value] of Object.entries(data)) {
      if (key === 'id' || key === 'created_at' || key === 'updated_at') continue;
      fields.push(`${key} = ?`);
      values.push(value);
    }
    if (fields.length === 0) return;
    values.push(id);
    await pool.execute(
      `UPDATE collection_batches SET ${fields.join(', ')} WHERE id = ?`,
      values
    );
  },
};
```

### 6-2. CollectionQueueRepo 수정

> 파일: `node-service/packages/common/src/db/repositories/CollectionQueueRepo.ts`

**변경 요약**:
- `updateTxHash()` 삭제 (tx_hash 컬럼 제거)
- `findQueuedWalletGroups()` 추가 — 배치 실행을 위한 지갑 그룹핑
- `markCollected()` 추가 — batch 확정 시 일괄 COLLECTED 처리

```typescript
// ★ 삭제: updateTxHash() 메서드 전체 삭제

// ★ 추가 메서드:

/**
 * QUEUED 상태인 큐를 지갑+통화별로 그룹핑.
 * CollectionBatchRunner가 호출.
 */
async findQueuedWalletGroups(): Promise<WalletGroup[]> {
  const pool = getPool();
  const [rows] = await pool.execute(
    `SELECT wallet_address_id, network_id, currency_id, partner_id,
            COUNT(*) as queue_count
     FROM collection_queue
     WHERE status = 'QUEUED'
     GROUP BY wallet_address_id, network_id, currency_id, partner_id`
  );
  return rows as WalletGroup[];
},

/**
 * 지갑+통화의 QUEUED 건 전부 → COLLECTED (batch_id 할당).
 * Webhook COLLECTION_CONFIRM 시 호출 (Spring 쪽에서 처리하므로 Node에서는 미사용 가능).
 */
async markCollected(walletAddressId: number, currencyId: number, batchId: number): Promise<number> {
  const pool = getPool();
  const [result] = await pool.execute(
    `UPDATE collection_queue
     SET status = 'COLLECTED', batch_id = ?
     WHERE wallet_address_id = ? AND currency_id = ? AND status = 'QUEUED'`,
    [batchId, walletAddressId, currencyId]
  );
  return (result as any).affectedRows;
},
```

> `WalletGroup` 인터페이스는 CollectionBatchRepo.ts에 정의되어 있으므로 import하거나, 공통 types에 정의.

### 6-3. common/src/db/repositories/index.ts 수정

> export에 `collectionBatchRepo` 추가:

```typescript
export { collectionBatchRepo } from './CollectionBatchRepo';
```

---

## 7. Node.js — relayer-api 패키지

### 7-1. CollectionPoller.ts → CollectionBatchRunner.ts 리팩토링

> 파일: `node-service/packages/relayer-api/src/services/CollectionBatchRunner.ts`
> 기존 `CollectionPoller.ts`를 참고하되 **전면 재작성**.
> 기존 파일은 삭제하지 않고 이름만 변경: `CollectionPoller.ts.bak` (롤백 대비)

**핵심 차이점**:

| 항목 | 기존 CollectionPoller | 신규 CollectionBatchRunner |
|------|----------------------|---------------------------|
| 폴링 단위 | 건별 (collection_queue 1건씩) | 지갑별 (wallet_address_id 그룹) |
| TX 수 | 입금 건당 1 TX | 지갑당 1 TX (전액 sweep) |
| 금액 기준 | collection_queue.amount (DB) | 온체인 잔액 조회 (blockchain-api) |
| 확정 대기 | `waitForConfirmation()` | 없음 — Webhook이 처리 |
| 상태 업데이트 | collection_queue 직접 | collection_batches만 (queue는 Webhook에서) |
| 폴링 주기 | 3초 | system_settings 기반 (1시간+) |
| 스킵 조건 | 없음 | 가스비 임계값 + 최소 금액 |

**전체 코드**: `COLLECTION_BATCH_REDESIGN.md` §5-2 참조 (이미 설계 확정된 코드).

### 7-2. app.ts — CollectionPoller → CollectionBatchRunner 교체

> 파일: `node-service/packages/relayer-api/src/app.ts`

```typescript
// ★ 변경 전 (L135-148):
// const collectionPoller = new CollectionPoller({
//   pollIntervalMs: 3_000,
//   batchSize: 20,
// });
// collectionPoller.start().catch(...)

// ★ 변경 후:
import { CollectionBatchRunner } from './services/CollectionBatchRunner';

// system_settings에서 interval 로드 (또는 기본값 1시간)
const collectionRunner = new CollectionBatchRunner({
  intervalMs: 60 * 60 * 1000,  // 1시간 (system_settings.collection.schedule_interval_hours × 3600000)
});
collectionRunner.start().catch((err) => {
  logger.error('CollectionBatchRunner crashed', err);
});

// graceful shutdown에서도 교체:
// collectionPoller.stop() → collectionRunner.stop()
```

---

## 8. 구현 순서 체크리스트

> **★ = 필수 선행**, 나머지는 순서 유연

### Phase 1: DDL + Entity (선행) ★

```
□ 1-1. MySQL에 마이그레이션 SQL 실행 (§1)
□ 1-2. CollectionBatchStatus enum 신규 (§2-1)
□ 1-3. CollectionStatus enum 수정 (§2-2)
□ 1-4. CollectionMode enum 삭제 (§2-3)
□ 1-5. CollectionQueue entity 수정 (§2-4)
□ 1-6. CollectionBatch entity 신규 (§2-5)
□ 1-7. CollectionQueueRepository 수정 (§2-6)
□ 1-8. CollectionBatchRepository 신규 (§2-7)
□ 1-9. WithdrawalStatus에 STALE 추가 (§2-8)
□ 1-10. ErrorCodes 추가 (§2-9)
□ 1-11. :common:compileJava 성공 확인
```

### Phase 2: Spring Boot core + open-api ★

```
□ 2-1. DepositService.enqueueCollection() 수정 (§3-1)
□ 2-2. DepositService.onCollectionConfirmed() 추가 (§3-2)
□ 2-3. BusinessEventClassifier 수정 (§4-1)
□ 2-4. WebhookProcessingService — confirmCollection() 재작성 (§4-2-c)
□ 2-5. WebhookProcessingService — handleFailedTx() 수정 (§4-2-d)
□ 2-6. WebhookProcessingService — 생성자 수정 (§4-2-a,b)
□ 2-7. :open-api:compileJava 성공 확인
□ 2-8. :core:compileJava 성공 확인
```

### Phase 3: Spring Boot admin-api

```
□ 3-1. CollectionSearchMapper 수정 (§5-1)
□ 3-2. CollectionSearchRequest DTO 수정 (§5-2)
□ 3-3. CollectionListResponse / CollectionDetailResponse DTO 수정 (§5-3)
□ 3-4. CollectionManagementService — failStaleBatch() 추가 (§5-4)
□ 3-5. :admin-api:compileJava 성공 확인
```

### Phase 4: Node.js 리팩토링

```
□ 4-1. CollectionBatchRepo.ts 신규 (§6-1)
□ 4-2. CollectionQueueRepo.ts 수정 (§6-2)
□ 4-3. common index.ts 수정 (§6-3)
□ 4-4. CollectionBatchRunner.ts 신규 (§7-1 + REDESIGN.md §5-2)
□ 4-5. app.ts — CollectionPoller → CollectionBatchRunner 교체 (§7-2)
□ 4-6. tsc 빌드 성공 확인
```

### Phase 5: 통합 테스트

```
□ 5-1. 입금 3건 → collection_queue 3건 QUEUED 확인
□ 5-2. CollectionBatchRunner 수동 트리거 → batch 1건 BROADCASTING
□ 5-3. Webhook COLLECTION_CONFIRM → batch CONFIRMED + queue 3건 COLLECTED
□ 5-4. Deposit 상태 COLLECTING → COLLECTED 전이 확인
□ 5-5. gas_cost_records 기록 확인
□ 5-6. Admin 콘솔 집금 목록에 batch 정보 표시 확인
```

---

## 참고: 파일별 변경 요약 매트릭스

| 파일 | 변경 유형 | 핵심 내용 |
|------|-----------|-----------|
| `CollectionBatchStatus.java` | **신규** | PENDING/BROADCASTING/CONFIRMED/FAILED/STALE |
| `CollectionStatus.java` | **수정** | 6개 → 2개 (QUEUED/COLLECTED) |
| `CollectionMode.java` | **삭제** | 배치 스케줄이 결정 |
| `CollectionQueue.java` | **수정** | depositId/batchId 추가, mode/tx/relayer 제거 |
| `CollectionBatch.java` | **신규** | 지갑별 물리 TX 단위 |
| `CollectionQueueRepository.java` | **수정** | findByTxHash 제거, findByDepositId/BatchId 추가 |
| `CollectionBatchRepository.java` | **신규** | findByTxHash, findByStatus 등 |
| `WithdrawalStatus.java` | **수정** | STALE 추가 |
| `ErrorCodes.java` | **수정** | COLLECTION_BATCH_* 추가 |
| `DepositService.java` | **수정** | enqueueCollection + onCollectionConfirmed |
| `BusinessEventClassifier.java` | **수정** | collectionQueue → collectionBatch |
| `WebhookProcessingService.java` | **수정** | confirmCollection 전면 재작성 + handleFailedTx 수정 |
| `CollectionSearchMapper.java` | **수정** | batch JOIN + collectionMode 필터 제거 |
| `CollectionSearchRequest.java` | **수정** | collectionMode → batchCode |
| `CollectionListResponse.java` | **수정** | batch 필드 추가 |
| `CollectionDetailResponse.java` | **수정** | batch 필드 추가 |
| `CollectionManagementService.java` | **수정** | failStaleBatch 추가 |
| `CollectionBatchRepo.ts` | **신규** | Node.js collection_batches CRUD |
| `CollectionQueueRepo.ts` | **수정** | findQueuedWalletGroups, markCollected 추가 |
| `CollectionBatchRunner.ts` | **신규** | CollectionPoller 대체 (배치 브로드캐스트 전담) |
| `app.ts` | **수정** | CollectionPoller → CollectionBatchRunner |
