feat(ledger): 거래 단계별 처리시각 기록·조회 및 pacs.002 결과통보 조회
- transfer에 pdng_at(Dior)·acsp_at(Hermes) 컬럼 추가(V5), Dior/Hermes가 각 단계 처리시각 기록 - Chanel /inquiry 응답에 접수·순번·기록·정산·완결 5단계 시각(rcvdAt/actcAt/pdngAt/acspAt/acccAt) 노출 - 경합 안전 upsert(Dior COALESCE, Hermes acsp_at 갱신·pdng_at 보존) Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
This commit is contained in:
@@ -2,6 +2,7 @@ package kr.or.bok.rtgs.chanel
|
||||
|
||||
import kr.or.bok.rtgs.chanel.jpa.AccountRepository
|
||||
import kr.or.bok.rtgs.chanel.jpa.NotificationRepository
|
||||
import kr.or.bok.rtgs.chanel.jpa.RawMessageRepository
|
||||
import kr.or.bok.rtgs.chanel.jpa.TransferRepository
|
||||
import kr.or.bok.rtgs.common.iso20022.Iso20022Codec
|
||||
import org.springframework.http.MediaType
|
||||
@@ -21,6 +22,7 @@ class ChanelController(
|
||||
private val transfers: TransferRepository,
|
||||
private val accounts: AccountRepository,
|
||||
private val notifications: NotificationRepository,
|
||||
private val rawMessages: RawMessageRepository,
|
||||
) {
|
||||
/** 자금이체 신청/접수: pacs.008 XML → pacs.002 XML. */
|
||||
@PostMapping(
|
||||
@@ -31,11 +33,13 @@ class ChanelController(
|
||||
fun payCustomer(@RequestBody rawXml: String): String =
|
||||
Iso20022Codec.writePacs002(service.accept(rawXml))
|
||||
|
||||
/** 거래 상태 조회. */
|
||||
/** 거래 상태 조회. 서비스별 단계 처리시각(타임라인) 포함. */
|
||||
@GetMapping("/inquiry/{bmi}")
|
||||
fun inquiry(@PathVariable bmi: String): Map<String, Any?> {
|
||||
val t = transfers.findById(bmi).orElse(null)
|
||||
?: return mapOf("bmi" to bmi, "status" to "IN_FLIGHT", "note" to "not yet persisted")
|
||||
val finalized = t.status.name == "ACCC" || t.status.name == "RJCT"
|
||||
val rcvdAt = rawMessages.findById(bmi).orElse(null)?.receivedAt
|
||||
return mapOf(
|
||||
"bmi" to t.bmi,
|
||||
"globalSeq" to t.globalSeq,
|
||||
@@ -46,6 +50,12 @@ class ChanelController(
|
||||
"originCenter" to t.originCenter,
|
||||
"updatedAt" to t.updatedAt,
|
||||
"reason" to t.reason,
|
||||
// 서비스별 처리시각(epoch ms). 접수=Chanel, 순번=Sequencer, 기록=Dior, 정산=Hermes, 완결=Prada.
|
||||
"rcvdAt" to rcvdAt,
|
||||
"actcAt" to t.createdAt,
|
||||
"pdngAt" to t.pdngAt,
|
||||
"acspAt" to t.acspAt,
|
||||
"acccAt" to if (finalized) t.updatedAt else null,
|
||||
)
|
||||
}
|
||||
|
||||
|
||||
@@ -47,6 +47,9 @@ class TransferRecord(
|
||||
@Column(name = "created_at") var createdAt: Long = 0,
|
||||
@Column(name = "updated_at") var updatedAt: Long = 0,
|
||||
@Column(name = "reason") var reason: String? = null,
|
||||
/** 단계별 처리시각(서비스 타임라인): 기록(Dior)·정산(Hermes). 접수=raw_message, 순번=created_at, 완결=updated_at. */
|
||||
@Column(name = "pdng_at") var pdngAt: Long? = null,
|
||||
@Column(name = "acsp_at") var acspAt: Long? = null,
|
||||
)
|
||||
|
||||
/**
|
||||
|
||||
@@ -23,17 +23,19 @@ class DiorApplication {
|
||||
|
||||
interface TransferRepository : JpaRepository<TransferRecord, String> {
|
||||
/**
|
||||
* 접수내역을 PDNG로 삽입하되, 이미 존재하면(=Hermes/타 소비자가 먼저 만든 경우) 아무것도 안 함.
|
||||
* 원자적 upsert라 Dior-Hermes 동시 삽입 경합에서 중복키 예외가 발생하지 않는다.
|
||||
* 접수내역을 PDNG로 삽입하되, 이미 존재하면(=Hermes/타 소비자가 먼저 만든 경우) 상태·금액 등은 건드리지 않고
|
||||
* 기록시각(pdng_at)만 비어있을 때 채운다(COALESCE). 원자적 upsert라 Dior-Hermes 동시 삽입 경합에서
|
||||
* 중복키 예외가 발생하지 않으며, 어느 쪽이 먼저 행을 만들어도 pdng_at은 유실 없이 기록된다.
|
||||
*/
|
||||
@Modifying
|
||||
@Query(
|
||||
value = """
|
||||
INSERT INTO transfer
|
||||
(bmi, global_seq, status, sender_code, receiver_code, amount, origin_center, created_at, updated_at)
|
||||
(bmi, global_seq, status, sender_code, receiver_code, amount, origin_center, created_at, updated_at, pdng_at)
|
||||
VALUES
|
||||
(:bmi, :seq, 'PDNG', :sender, :receiver, :amount, :origin, :createdAt, :updatedAt)
|
||||
ON CONFLICT (bmi) DO NOTHING
|
||||
(:bmi, :seq, 'PDNG', :sender, :receiver, :amount, :origin, :createdAt, :updatedAt, :pdngAt)
|
||||
ON CONFLICT (bmi) DO UPDATE SET
|
||||
pdng_at = COALESCE(transfer.pdng_at, EXCLUDED.pdng_at)
|
||||
""",
|
||||
nativeQuery = true,
|
||||
)
|
||||
@@ -46,6 +48,7 @@ interface TransferRepository : JpaRepository<TransferRecord, String> {
|
||||
@Param("origin") origin: String,
|
||||
@Param("createdAt") createdAt: Long,
|
||||
@Param("updatedAt") updatedAt: Long,
|
||||
@Param("pdngAt") pdngAt: Long,
|
||||
): Int
|
||||
}
|
||||
|
||||
|
||||
@@ -24,8 +24,9 @@ class DiorService(
|
||||
@Transactional
|
||||
fun onJournal(payload: String) {
|
||||
val e = mapper.readValue(payload, JournalEntry::class.java)
|
||||
// 원자적 upsert: 이미 있으면(Hermes/타 소비자가 먼저 만든 경우) 아무 일도 안 함 → 중복키 예외 없음.
|
||||
val inserted = transfers.insertPdngIfAbsent(
|
||||
val now = System.currentTimeMillis()
|
||||
// 원자적 upsert: 신규면 PDNG로 삽입, 이미 있으면 pdng_at(기록시각)만 채움 → 중복키 예외 없음.
|
||||
transfers.insertPdngIfAbsent(
|
||||
bmi = e.core.bmi,
|
||||
seq = e.globalSeq,
|
||||
sender = e.core.senderCode,
|
||||
@@ -33,9 +34,9 @@ class DiorService(
|
||||
amount = e.core.amount,
|
||||
origin = e.originCenter,
|
||||
createdAt = e.seqEpochMillis,
|
||||
updatedAt = System.currentTimeMillis(),
|
||||
updatedAt = now,
|
||||
pdngAt = now,
|
||||
)
|
||||
if (inserted > 0) log.info("PDNG seq=#{} bmi={}", e.globalSeq, e.core.bmi)
|
||||
else log.debug("PDNG skip(exists) seq=#{} bmi={}", e.globalSeq, e.core.bmi)
|
||||
log.info("PDNG seq=#{} bmi={}", e.globalSeq, e.core.bmi)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -31,15 +31,16 @@ interface JournalLogRepository : JpaRepository<JournalLog, Long>
|
||||
interface TransferRepository : JpaRepository<TransferRecord, String> {
|
||||
/**
|
||||
* 거래 원장 원자적 upsert. 신규면 INSERT, 이미 있으면(Dior/타 소비자가 먼저 생성) UPDATE.
|
||||
* created_at은 최초 값 유지(갱신 제외). Dior-Hermes 동시 삽입 경합에서 중복키 예외 제거.
|
||||
* created_at은 최초 값 유지(갱신 제외). acsp_at(정산시각)은 Hermes가 기록하고, pdng_at은 SET에서
|
||||
* 제외해 Dior가 기록한 값을 보존한다. Dior-Hermes 동시 삽입 경합에서 중복키 예외 제거.
|
||||
*/
|
||||
@Modifying
|
||||
@Query(
|
||||
value = """
|
||||
INSERT INTO transfer
|
||||
(bmi, global_seq, status, sender_code, receiver_code, amount, origin_center, created_at, updated_at, reason)
|
||||
(bmi, global_seq, status, sender_code, receiver_code, amount, origin_center, created_at, updated_at, reason, acsp_at)
|
||||
VALUES
|
||||
(:bmi, :seq, :status, :sender, :receiver, :amount, :origin, :createdAt, :updatedAt, :reason)
|
||||
(:bmi, :seq, :status, :sender, :receiver, :amount, :origin, :createdAt, :updatedAt, :reason, :acspAt)
|
||||
ON CONFLICT (bmi) DO UPDATE SET
|
||||
global_seq = EXCLUDED.global_seq,
|
||||
status = EXCLUDED.status,
|
||||
@@ -48,7 +49,8 @@ interface TransferRepository : JpaRepository<TransferRecord, String> {
|
||||
amount = EXCLUDED.amount,
|
||||
origin_center = EXCLUDED.origin_center,
|
||||
updated_at = EXCLUDED.updated_at,
|
||||
reason = EXCLUDED.reason
|
||||
reason = EXCLUDED.reason,
|
||||
acsp_at = EXCLUDED.acsp_at
|
||||
""",
|
||||
nativeQuery = true,
|
||||
)
|
||||
@@ -63,6 +65,7 @@ interface TransferRepository : JpaRepository<TransferRecord, String> {
|
||||
@Param("createdAt") createdAt: Long,
|
||||
@Param("updatedAt") updatedAt: Long,
|
||||
@Param("reason") reason: String?,
|
||||
@Param("acspAt") acspAt: Long,
|
||||
): Int
|
||||
}
|
||||
|
||||
|
||||
@@ -64,6 +64,7 @@ class SettlementApplier(
|
||||
bmi = e.core.bmi, seq = e.globalSeq, status = status.name,
|
||||
sender = e.core.senderCode, receiver = e.core.receiverCode, amount = e.core.amount,
|
||||
origin = e.originCenter, createdAt = e.seqEpochMillis, updatedAt = now, reason = reason,
|
||||
acspAt = now,
|
||||
)
|
||||
journalLogs.findById(e.globalSeq).ifPresent { it.applied = true; journalLogs.save(it) }
|
||||
return ResultMessage(
|
||||
|
||||
@@ -0,0 +1,6 @@
|
||||
-- 거래 단계별 처리시각(서비스별 타임라인). 기존 created_at(순번)·updated_at(완결)에 더해
|
||||
-- 기록(Dior)·정산(Hermes) 처리시각을 원장에 영속 기록한다. 추가(additive)라 기존 데이터 무해.
|
||||
ALTER TABLE transfer ADD COLUMN IF NOT EXISTS pdng_at BIGINT;
|
||||
ALTER TABLE transfer ADD COLUMN IF NOT EXISTS acsp_at BIGINT;
|
||||
COMMENT ON COLUMN transfer.pdng_at IS '기록(PDNG) 처리시각(epoch ms) — Dior';
|
||||
COMMENT ON COLUMN transfer.acsp_at IS '정산(ACSP) 처리시각(epoch ms) — Hermes';
|
||||
Reference in New Issue
Block a user