feat(refdata): Q22 basicdata 映射与 ReferenceDataProcessor(ACM2-93)
闭合 Q22;V6 建 basicdata 表组;RefData 走真处理器(DNLD/RESP/ADD/UPD/DEL/RSTA)。 Co-authored-by: Cursor <cursoragent@cursor.com>
This commit is contained in:
@@ -84,7 +84,6 @@ class PipelineSmokeTest {
|
||||
</MSG>
|
||||
""".trimIndent()
|
||||
|
||||
/** 参考数据类别(SIS:3.2):ReferenceDataProcessor 未实装(ACM2-93),先按可重试失败留队 */
|
||||
val REFDATA_XML = """
|
||||
<MSG>
|
||||
<META><SNDR>AODB</SNDR><SEQN>2</SEQN><DTTM>20260908120000</DTTM><TYPE>ARPT</TYPE><STYP>DNLD</STYP></META>
|
||||
@@ -110,11 +109,10 @@ class PipelineSmokeTest {
|
||||
|
||||
@Test
|
||||
fun `replay reopens failed UNSUPPORTED row to PENDING`() {
|
||||
val receipt = controller.send(REFDATA_XML)
|
||||
val id = receipt.body()!!.toLong()
|
||||
pump.tick()
|
||||
val stub = ctx.getBean(StubProcState::class.java)
|
||||
assertEquals(ProcStatus.FAILED, stub.snapshotOf(id)!!.state)
|
||||
val id = 88001L
|
||||
stub.insertIfAbsent(id, java.time.Instant.now())
|
||||
stub.update(id, ProcStatus.FAILED, errorClass = ErrorClass.UNSUPPORTED, lastError = "test-unsupported")
|
||||
|
||||
val n = ctx.getBean(com.gzzn.omms.msgexchange.infra.retry.ReplayService::class.java)
|
||||
.replay(listOf(ErrorClass.UNSUPPORTED))
|
||||
@@ -125,6 +123,15 @@ class PipelineSmokeTest {
|
||||
assertEquals(0, reopened.attempts)
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `reference data empty DNLD succeeds`() {
|
||||
val receipt = controller.send(REFDATA_XML)
|
||||
val id = receipt.body()!!.toLong()
|
||||
pump.tick()
|
||||
val stub = ctx.getBean(StubProcState::class.java)
|
||||
assertEquals(ProcStatus.SUCCEEDED, stub.snapshotOf(id)!!.state)
|
||||
}
|
||||
|
||||
/**
|
||||
* 回归用例:死信在进入终态时同时记下回填意图(同一条 UPDATE),由回填扫描把标记写回信箱。
|
||||
* 少了这一步,这些行永远占着每批的名额,攒够一批就再也发现不了新消息了。
|
||||
|
||||
@@ -206,6 +206,7 @@ class FlightCommitTest {
|
||||
flopProcessor = FlopProcessor(commit, f.flights, f.events, mapper, clock),
|
||||
fdelProcessor = FdelProcessor(commit, f.flights, f.events, mapper, clock),
|
||||
adftProcessor = AdftProcessor(commit, f.flights, f.events, opDay, mapper, clock),
|
||||
referenceDataProcessor = testRefDataProcessor(f.procState, clock = clock),
|
||||
outbound = OutboundRequestService(
|
||||
com.gzzn.omms.msgexchange.infra.stub.StubReqTrack(clock),
|
||||
com.gzzn.omms.msgexchange.infra.stub.StubCoutmsgOutbox(),
|
||||
|
||||
@@ -24,7 +24,11 @@ import com.gzzn.omms.msgexchange.infra.stub.StubReqTrack
|
||||
import com.gzzn.omms.msgexchange.processing.OutboundRequestService
|
||||
import com.gzzn.omms.msgexchange.infra.stub.StubMsgEvents
|
||||
import com.gzzn.omms.msgexchange.infra.stub.StubProcState
|
||||
import com.gzzn.omms.msgexchange.infra.stub.StubRefDataRepository
|
||||
import com.gzzn.omms.msgexchange.infra.stub.StubSnapshotLog
|
||||
import com.gzzn.omms.msgexchange.domain.ref.RefDataBody
|
||||
import com.gzzn.omms.msgexchange.domain.ref.RefDataRecord
|
||||
import com.gzzn.omms.msgexchange.domain.ref.ResourceStatusRecord
|
||||
import com.fasterxml.jackson.databind.ObjectMapper
|
||||
import org.junit.jupiter.api.Assertions.assertEquals
|
||||
import org.junit.jupiter.api.Assertions.assertNotNull
|
||||
@@ -36,7 +40,7 @@ import java.time.Instant
|
||||
/**
|
||||
* 忽略类报文(US-04):命中忽略清单的报文先绑定身份再写 SKIPPED,
|
||||
* 不产生航班或 outbox 副作用;未命中且无处理器的合法类型同样落 SKIPPED(unsupported)
|
||||
* 终态(US-03 AC2),参考数据类别(US-13)在处理器落地前保持可重试失败。
|
||||
* 终态(US-03 AC2);参考数据(US-13)走 ReferenceDataProcessor,校验失败为 DEAD(PROTOCOL)。
|
||||
*/
|
||||
class IgnoreBranchTest {
|
||||
|
||||
@@ -70,6 +74,19 @@ class IgnoreBranchTest {
|
||||
val log = StubSnapshotLog()
|
||||
val procFailure = ProcFailure(proc, FailureScheduler(props, clock))
|
||||
val commit = testCommit(procState = proc, msgEvents = events, clock = clock)
|
||||
val refCommit = RefDataCommit(
|
||||
object : PipelineTransactionManager { override fun <T> inTransaction(block: () -> T) = block() },
|
||||
object : PipelineLockRepository {
|
||||
override fun lock() = Unit
|
||||
override fun <T> holdAcrossTransactions(block: () -> T) = block()
|
||||
},
|
||||
proc,
|
||||
clock,
|
||||
)
|
||||
val outbound = OutboundRequestService(
|
||||
StubReqTrack(clock), StubCoutmsgOutbox(), JacksonXmlCodec(), props, opDay, clock,
|
||||
)
|
||||
val refData = ReferenceDataProcessor(refCommit, StubRefDataRepository(), outbound)
|
||||
return MessageProcessor(
|
||||
inbox = inbox,
|
||||
procState = proc,
|
||||
@@ -78,9 +95,8 @@ class IgnoreBranchTest {
|
||||
flopProcessor = FlopProcessor(commit, flights, events, ObjectMapper(), clock),
|
||||
fdelProcessor = FdelProcessor(commit, flights, events, ObjectMapper(), clock),
|
||||
adftProcessor = AdftProcessor(commit, flights, events, opDay, ObjectMapper(), clock),
|
||||
outbound = OutboundRequestService(
|
||||
StubReqTrack(clock), StubCoutmsgOutbox(), JacksonXmlCodec(), props, opDay, clock,
|
||||
),
|
||||
referenceDataProcessor = refData,
|
||||
outbound = outbound,
|
||||
procFailure = procFailure,
|
||||
props = props,
|
||||
clock = clock,
|
||||
@@ -127,25 +143,33 @@ class IgnoreBranchTest {
|
||||
|
||||
@Test
|
||||
fun `REGN and RSTA route to refdata path not ignore list`() {
|
||||
for (type in listOf("REGN", "RSTA")) {
|
||||
val regnMsg = DecodedMessage(
|
||||
meta = MetaFields("AODB", "REGN", "ADD", 300L, 1L),
|
||||
kind = MsgKind.RefData("REGN"),
|
||||
rawXml = "<MSG/>",
|
||||
body = RefDataBody("REGN", "ADD", listOf(RefDataRecord(mapOf("RNUM" to "B1234", "ITAT" to "738")))),
|
||||
)
|
||||
val rstaMsg = DecodedMessage(
|
||||
meta = MetaFields("AODB", "RSTA", "DNLD", 301L, 1L),
|
||||
kind = MsgKind.RefData("RSTA"),
|
||||
rawXml = "<MSG/>",
|
||||
body = RefDataBody(
|
||||
"RSTA", "DNLD", emptyList(),
|
||||
listOf(ResourceStatusRecord("GATE", "C05", "D")),
|
||||
),
|
||||
)
|
||||
for ((type, msg) in listOf("REGN" to regnMsg, "RSTA" to rstaMsg)) {
|
||||
val proc = StubProcState()
|
||||
val id = msgId + type.hashCode().toLong()
|
||||
proc.insertIfAbsent(id, null)
|
||||
val inbox = StubInbox()
|
||||
inbox.raws[id] = "<RAW/>"
|
||||
val msg = DecodedMessage(
|
||||
meta = MetaFields("AODB", type, "X", 300L, 1L),
|
||||
kind = MsgKind.RefData(type),
|
||||
rawXml = "<MSG/>",
|
||||
)
|
||||
val p = processor(proc = proc, inbox = inbox, codec = codecReturning(msg))
|
||||
|
||||
p.processOne(ProcState(id, ProcStatus.PENDING, updatedAt = Instant.EPOCH))
|
||||
|
||||
val row = proc.find(id)!!
|
||||
assertEquals(ProcStatus.FAILED, row.state)
|
||||
assertEquals(ErrorClass.UNSUPPORTED, row.errorClass)
|
||||
assertEquals("refdata-pending:$type", row.lastError)
|
||||
assertEquals(ProcStatus.SUCCEEDED, row.state, "type=$type")
|
||||
}
|
||||
}
|
||||
|
||||
@@ -202,23 +226,25 @@ class IgnoreBranchTest {
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `reference data type without handler stays retryable instead of being swallowed`() {
|
||||
fun `reference data validation failure is DEAD PROTOCOL not SKIPPED`() {
|
||||
val proc = StubProcState()
|
||||
proc.insertIfAbsent(msgId, null)
|
||||
val inbox = StubInbox()
|
||||
inbox.raws[msgId] = "<RAW/>"
|
||||
val msg = DecodedMessage(
|
||||
meta = MetaFields("AODB", "ARPT", "DNLD", 400L, 1L),
|
||||
meta = MetaFields("AODB", "ARPT", "ADD", 400L, 1L),
|
||||
kind = MsgKind.RefData("ARPT"),
|
||||
rawXml = "<RAW/>",
|
||||
body = RefDataBody("ARPT", "ADD", listOf(RefDataRecord(mapOf("ITCD" to "")))),
|
||||
)
|
||||
val p = processor(proc = proc, inbox = inbox, codec = codecReturning(msg))
|
||||
|
||||
p.processOne(head())
|
||||
|
||||
val row = proc.find(msgId)!!
|
||||
assertEquals(ProcStatus.FAILED, row.state)
|
||||
assertEquals("refdata-pending:ARPT", row.lastError)
|
||||
assertEquals(ProcStatus.DEAD, row.state)
|
||||
assertEquals(ErrorClass.PROTOCOL, row.errorClass)
|
||||
assertTrue(row.lastError!!.contains("missing-key"))
|
||||
}
|
||||
|
||||
@Test
|
||||
|
||||
@@ -0,0 +1,112 @@
|
||||
package com.gzzn.omms.msgexchange.processing
|
||||
|
||||
import com.gzzn.omms.msgexchange.MutableClock
|
||||
import com.gzzn.omms.msgexchange.codec.JacksonXmlCodec
|
||||
import com.gzzn.omms.msgexchange.config.OperationDayProps
|
||||
import com.gzzn.omms.msgexchange.config.PipelineProps
|
||||
import com.gzzn.omms.msgexchange.domain.DecodedMessage
|
||||
import com.gzzn.omms.msgexchange.domain.MetaFields
|
||||
import com.gzzn.omms.msgexchange.domain.MsgKind
|
||||
import com.gzzn.omms.msgexchange.domain.ProcState
|
||||
import com.gzzn.omms.msgexchange.domain.ProcStatus
|
||||
import com.gzzn.omms.msgexchange.domain.ref.RefDataBody
|
||||
import com.gzzn.omms.msgexchange.infra.persistence.ReqTrackRepository
|
||||
import com.gzzn.omms.msgexchange.infra.stub.StubCoutmsgOutbox
|
||||
import com.gzzn.omms.msgexchange.infra.stub.StubProcState
|
||||
import com.gzzn.omms.msgexchange.infra.stub.StubPipelineLock
|
||||
import com.gzzn.omms.msgexchange.infra.stub.StubPipelineTx
|
||||
import com.gzzn.omms.msgexchange.infra.stub.StubRefDataRepository
|
||||
import com.gzzn.omms.msgexchange.infra.stub.StubReqTrack
|
||||
import org.junit.jupiter.api.Assertions.assertEquals
|
||||
import org.junit.jupiter.api.Assertions.assertNull
|
||||
import org.junit.jupiter.api.Test
|
||||
import java.math.BigDecimal
|
||||
|
||||
class ReferenceDataProcessorTest {
|
||||
private val clock = MutableClock(MutableClock.BASE)
|
||||
private val proc = StubProcState()
|
||||
private val refRepo = StubRefDataRepository()
|
||||
|
||||
private fun processor(req: StubReqTrack = StubReqTrack(clock)): ReferenceDataProcessor {
|
||||
val refCommit = RefDataCommit(StubPipelineTx(), StubPipelineLock(), proc, clock) // stub tx/lock
|
||||
val props = PipelineProps()
|
||||
val opDay = OperationDayProps().apply { zone = "Asia/Shanghai"; cutoffHour = 0 }
|
||||
val outbound = OutboundRequestService(req, StubCoutmsgOutbox(), JacksonXmlCodec(), props, opDay, clock)
|
||||
return ReferenceDataProcessor(refCommit, refRepo, outbound)
|
||||
}
|
||||
|
||||
private fun head(id: Long = 1L) = ProcState(id, ProcStatus.PENDING, updatedAt = clock.instant())
|
||||
|
||||
private fun msg(type: String, styp: String, xml: String): DecodedMessage {
|
||||
val decoded = JacksonXmlCodec().decode(xml)
|
||||
require(decoded is com.gzzn.omms.msgexchange.codec.DecodeResult.Ok)
|
||||
return decoded.message
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `COUL DNLD full replace swaps table contents`() {
|
||||
refRepo.countries["XX"] = mutableMapOf("COUNTRY_CODE" to "XX")
|
||||
val xml = """
|
||||
<MSG>
|
||||
<META><SNDR>AODB</SNDR><SEQN>1</SEQN><DTTM>20260908120000</DTTM><TYPE>COUL</TYPE><STYP>DNLD</STYP></META>
|
||||
<COUL><COUC>CN</COUC><COUN>China</COUN><CNMC>中国</CNMC><REGC>11</REGC></COUL>
|
||||
</MSG>
|
||||
""".trimIndent()
|
||||
processor().apply(head(), msg("COUL", "DNLD", xml), "COUL")
|
||||
assertEquals(1, refRepo.countries.size)
|
||||
assertEquals("CN", refRepo.countries["CN"]!!["COUNTRY_CODE"])
|
||||
assertNull(refRepo.countries["XX"])
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `ARPT UPD empty tag clears nullable column`() {
|
||||
refRepo.airports["CTU"] = mutableMapOf(
|
||||
"AIRPORT_CODE_IATA" to "CTU",
|
||||
"AIRPORT_NAME_ENG" to "Old",
|
||||
"HAUL" to "S",
|
||||
)
|
||||
val xml = """
|
||||
<MSG>
|
||||
<META><SNDR>AODB</SNDR><SEQN>2</SEQN><DTTM>20260908120000</DTTM><TYPE>ARPT</TYPE><STYP>UPD</STYP></META>
|
||||
<ARPT><ITCD>CTU</ITCD><ICCD></ICCD><ANAM>Chengdu</ANAM><ANMC>成都</ANMC><CTRY>CN</CTRY><ACTY>CTU</ACTY><ATYP></ATYP><HAUL></HAUL></ARPT>
|
||||
</MSG>
|
||||
""".trimIndent()
|
||||
processor().apply(head(2), msg("ARPT", "UPD", xml), "ARPT")
|
||||
val row = refRepo.airports["CTU"]!!
|
||||
assertEquals("Chengdu", row["AIRPORT_NAME_ENG"])
|
||||
assertNull(row["HAUL"])
|
||||
assertNull(row["INTERNATIONAL_FLAG"])
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `RSTA DNLD updates gate status by RTYP and RSID`() {
|
||||
refRepo.gates["G1"] = mutableMapOf("GATE_CODE" to "G1", "GATE_STATUS" to BigDecimal.ZERO)
|
||||
val xml = """
|
||||
<MSG>
|
||||
<META><SNDR>AODB</SNDR><SEQN>3</SEQN><DTTM>20260908120000</DTTM><TYPE>RSTA</TYPE><STYP>DNLD</STYP></META>
|
||||
<RSTA><RTYP>GATE</RTYP><RSID>G1</RSID><STAT>D</STAT><RDST>20260908120000</RDST><RDET></RDET></RSTA>
|
||||
</MSG>
|
||||
""".trimIndent()
|
||||
processor().apply(head(3), msg("RSTA", "DNLD", xml), "RSTA")
|
||||
assertEquals(BigDecimal.ONE, refRepo.gates["G1"]!!["GATE_STATUS"])
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `RESP completes open RQRD track`() {
|
||||
val req = StubReqTrack(clock)
|
||||
val props = PipelineProps()
|
||||
val opDay = OperationDayProps().apply { zone = "Asia/Shanghai"; cutoffHour = 0 }
|
||||
val outbound = OutboundRequestService(req, StubCoutmsgOutbox(), JacksonXmlCodec(), props, opDay, clock)
|
||||
val day = outbound.currentOperationDay()
|
||||
val reqId = req.insert(com.gzzn.omms.msgexchange.domain.OutboundRequestKeys.RQRD_REQ_TYPE, day, "OMMS")
|
||||
req.markSent(reqId, clock.instant())
|
||||
val xml = """
|
||||
<MSG>
|
||||
<META><SNDR>AODB</SNDR><SEQN>4</SEQN><DTTM>20260908120000</DTTM><TYPE>GLST</TYPE><STYP>RESP</STYP></META>
|
||||
<GLST><GCOD>G9</GCOD><GTNM>Gate 9</GTNM><GNMC>9号门</GNMC><GCAT>D</GCAT><GTML>T2</GTML></GLST>
|
||||
</MSG>
|
||||
""".trimIndent()
|
||||
processor(req).apply(head(4), msg("GLST", "RESP", xml), "GLST")
|
||||
assertEquals(ReqTrackRepository.ReqState.DONE, req.rows[reqId]!!.state)
|
||||
}
|
||||
}
|
||||
@@ -70,6 +70,7 @@ class RespGuardTest {
|
||||
flopProcessor = FlopProcessor(commit, flights, events, ObjectMapper(), clock),
|
||||
fdelProcessor = FdelProcessor(commit, flights, events, ObjectMapper(), clock),
|
||||
adftProcessor = AdftProcessor(commit, flights, events, opDay, ObjectMapper(), clock),
|
||||
referenceDataProcessor = testRefDataProcessor(proc, clock = clock),
|
||||
outbound = OutboundRequestService(req, StubCoutmsgOutbox(), JacksonXmlCodec(), props, opDay, clock),
|
||||
procFailure = ProcFailure(proc, FailureScheduler(props, clock)),
|
||||
props = props,
|
||||
|
||||
@@ -8,7 +8,14 @@ import com.gzzn.omms.msgexchange.infra.stub.StubFlightProjectionPort
|
||||
import com.gzzn.omms.msgexchange.infra.stub.StubMsgEvents
|
||||
import com.gzzn.omms.msgexchange.infra.stub.StubPipelineLock
|
||||
import com.gzzn.omms.msgexchange.infra.stub.StubPipelineTx
|
||||
import com.gzzn.omms.msgexchange.codec.JacksonXmlCodec
|
||||
import com.gzzn.omms.msgexchange.config.OperationDayProps
|
||||
import com.gzzn.omms.msgexchange.config.PipelineProps
|
||||
import com.gzzn.omms.msgexchange.infra.persistence.RefDataRepository
|
||||
import com.gzzn.omms.msgexchange.infra.stub.StubCoutmsgOutbox
|
||||
import com.gzzn.omms.msgexchange.infra.stub.StubProcState
|
||||
import com.gzzn.omms.msgexchange.infra.stub.StubRefDataRepository
|
||||
import com.gzzn.omms.msgexchange.infra.stub.StubReqTrack
|
||||
import java.time.Clock
|
||||
|
||||
/**
|
||||
@@ -23,3 +30,22 @@ internal fun testCommit(
|
||||
txManager: PipelineTransactionManager = StubPipelineTx(),
|
||||
clock: Clock = Clock.systemUTC(),
|
||||
) = FlightCommit(txManager, lock, procState, msgEvents, projection, clock)
|
||||
|
||||
internal fun testRefDataProcessor(
|
||||
procState: ProcStateRepository = StubProcState(),
|
||||
refData: RefDataRepository = StubRefDataRepository(),
|
||||
clock: Clock = Clock.systemUTC(),
|
||||
): ReferenceDataProcessor {
|
||||
val refCommit = RefDataCommit(
|
||||
StubPipelineTx(),
|
||||
StubPipelineLock(),
|
||||
procState,
|
||||
clock,
|
||||
)
|
||||
val props = PipelineProps()
|
||||
val opDay = OperationDayProps().apply { zone = "Asia/Shanghai"; cutoffHour = 0 }
|
||||
val outbound = OutboundRequestService(
|
||||
StubReqTrack(clock), StubCoutmsgOutbox(), JacksonXmlCodec(), props, opDay, clock,
|
||||
)
|
||||
return ReferenceDataProcessor(refCommit, refData, outbound)
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user