From e5449a5cb1218ae8fec732ae145ec9b89f992b91 Mon Sep 17 00:00:00 2001 From: windyboy Date: Mon, 21 Sep 2026 16:01:06 +0800 Subject: [PATCH] =?UTF-8?q?feat(refdata):=20Q22=20basicdata=20=E6=98=A0?= =?UTF-8?q?=E5=B0=84=E4=B8=8E=20ReferenceDataProcessor=EF=BC=88ACM2-93?= =?UTF-8?q?=EF=BC=89?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 闭合 Q22;V6 建 basicdata 表组;RefData 走真处理器(DNLD/RESP/ADD/UPD/DEL/RSTA)。 Co-authored-by: Cursor --- docs/contracts/interface-contract.md | 6 +- docs/implementation.md | 23 + docs/reference.md | 1 - docs/specification.md | 4 +- .../omms/msgexchange/codec/JacksonXmlCodec.kt | 3 +- .../omms/msgexchange/codec/RefDataWire.kt | 205 +++++++++ .../omms/msgexchange/codec/SisMessageBody.kt | 14 + .../msgexchange/domain/ref/RefDataBody.kt | 36 ++ .../domain/ref/RefDataFieldMapping.kt | 37 ++ .../domain/ref/RefDataValidator.kt | 50 +++ .../infra/persistence/RefDataRepository.kt | 15 + .../persistence/jdbc/JdbcRefDataRepository.kt | 419 ++++++++++++++++++ .../infra/stub/StubRefDataRepository.kt | 177 ++++++++ .../processing/OutboundRequestService.kt | 3 + .../gzzn/omms/msgexchange/processing/Pump.kt | 8 +- .../msgexchange/processing/RefDataCommit.kt | 28 ++ .../processing/ReferenceDataProcessor.kt | 52 +++ .../db/migration/V6__basicdata_ref_data.sql | 178 ++++++++ .../omms/msgexchange/PipelineSmokeTest.kt | 17 +- .../processing/FlightCommitTest.kt | 1 + .../processing/IgnoreBranchTest.kt | 60 ++- .../processing/ReferenceDataProcessorTest.kt | 112 +++++ .../msgexchange/processing/RespGuardTest.kt | 1 + .../processing/TestFlightCommit.kt | 26 ++ 24 files changed, 1442 insertions(+), 34 deletions(-) create mode 100644 src/main/kotlin/com/gzzn/omms/msgexchange/codec/RefDataWire.kt create mode 100644 src/main/kotlin/com/gzzn/omms/msgexchange/domain/ref/RefDataBody.kt create mode 100644 src/main/kotlin/com/gzzn/omms/msgexchange/domain/ref/RefDataFieldMapping.kt create mode 100644 src/main/kotlin/com/gzzn/omms/msgexchange/domain/ref/RefDataValidator.kt create mode 100644 src/main/kotlin/com/gzzn/omms/msgexchange/infra/persistence/RefDataRepository.kt create mode 100644 src/main/kotlin/com/gzzn/omms/msgexchange/infra/persistence/jdbc/JdbcRefDataRepository.kt create mode 100644 src/main/kotlin/com/gzzn/omms/msgexchange/infra/stub/StubRefDataRepository.kt create mode 100644 src/main/kotlin/com/gzzn/omms/msgexchange/processing/RefDataCommit.kt create mode 100644 src/main/kotlin/com/gzzn/omms/msgexchange/processing/ReferenceDataProcessor.kt create mode 100644 src/main/resources/db/migration/V6__basicdata_ref_data.sql create mode 100644 src/test/kotlin/com/gzzn/omms/msgexchange/processing/ReferenceDataProcessorTest.kt diff --git a/docs/contracts/interface-contract.md b/docs/contracts/interface-contract.md index 1b44ddf..45cdda8 100644 --- a/docs/contracts/interface-contract.md +++ b/docs/contracts/interface-contract.md @@ -27,7 +27,7 @@ AODB 经 CIIMS adapter 把 XML 报文写入 `CMINMSGS`,格式以架构指定的 [SIS 接口规范](../legacy/SIS_AODB_RMS-V0.1.md) 与 [XSD](../legacy/unisysaodbsis.xsd) 为依据。 -入站报文按类型处理:`ADFT` 建立计划外航班,`FLOP` 改航班动态,`FDEL` 删航班,`SCHD-DNLD` 与 `SCHD-RESP` 同步日计划快照,参考数据写入自有 PG 的静态参考数据表组(逻辑视图 `REF_MASTER`,物理形态由内部迁移确定;报文到表的映射待 `Q22` 定稿;联调栈 `basicdata` schema 的基础数据表对齐 admin-api 实体注解,不预设 `REF_MASTER` 落表方式;`US-04`~`US-07`、`US-13`);报文不合法进死信,合法但本系统不支持的类型跳过留档并按已处理写回信箱(`US-03`)。 +入站报文按类型处理:`ADFT` 建立计划外航班,`FLOP` 改航班动态,`FDEL` 删航班,`SCHD-DNLD` 与 `SCHD-RESP` 同步日计划快照,参考数据写入自有 PG 的 schema `basicdata` 表组(逻辑视图 `REF_MASTER`,物理多表映射见 `implementation.md`「静态参考数据」、`Q22` 已定;联调栈同 schema 对齐 admin-api 实体注解;`US-04`~`US-07`、`US-13`);报文不合法进死信,合法但本系统不支持的类型跳过留档并按已处理写回信箱(`US-03`)。 | 环节 | 已确定的边界 | 尚需确定 | |---|---|---| @@ -87,7 +87,7 @@ AODB 经 CIIMS adapter 把 XML 报文写入 `CMINMSGS`,格式以架构指定 | 表或表组 | 边界 | 尚需确定的字段级契约 | |---|---|---| -| 静态参考数据表组 | 静态参考数据与资源状态保存在独立数据表组(逻辑视图 `REF_MASTER`,物理形态由内部迁移确定;报文到表的映射待 `Q22` 定稿;联调栈 `basicdata` schema 的基础数据表对齐 admin-api 实体注解,不预设 `REF_MASTER` 落表方式);表与列直接取 admin-api 实体注解,不改名、不合并,admin-api 直接只读;新消息覆盖旧记录,全量消息整体替换,增删改消息逐条处理(`US-13`);一类校验不通过只停这一类、其他类照常,校验失败类别的已有记录不变;字段为空表示「当前没有值」,不是删除(`US-13` AC4)。 | 13 类报文与资源状态到表组的映射(`Q22`);admin-api 需要哪些字段。类别码与消息中的识别标签见下表。 | +| 静态参考数据表组 | 静态参考数据与资源状态保存在 schema `basicdata`(逻辑视图 `REF_MASTER`;`RTYPE`→表映射见 `implementation.md`「静态参考数据」、`Q22` 已定);表与列取 admin-api 实体注解;全量 `DNLD`/`RESP` 整类替换,增量 `ADD`/`UPD`/`DEL` 逐条;单类校验失败不写入该类;空标签表示无值非删除(`US-13`)。 | admin-api 生产只读接入与字段裁剪需求。类别码与识别标签见下表。 | | 航班当前态表 | 航班当前态的唯一权威;`FLID` 唯一(`INV-6`);Redis 和 Kafka 从处理结果派生,不反向覆盖这些表。 | 主键、字段与类型、外键/索引。 | | 内部处理表 | 管道处理、请求跟踪、留痕与互斥由本系统维护;不对外提供直接读写接口。 | 字段与约束由内部实现设计确定。 | @@ -97,7 +97,7 @@ admin-api 还从本系统数据库只读季度计划;供数方与报文形态 ### 静态参考数据类别与编号来源 -类别码与识别标签来自架构引用的 [SIS 接口规范](../legacy/SIS_AODB_RMS-V0.1.md);识别同一条参考记录时,用类别码加识别标签值。标签值的格式、标签在哪个范围内唯一,以及报文到表的映射,仍需 `Q22` 定稿。 +类别码与识别标签来自 [SIS 接口规范](../legacy/SIS_AODB_RMS-V0.1.md);识别同一条参考记录时用类别码加识别标签值。报文到 `basicdata` 列的映射见 `implementation.md`「静态参考数据」。 | 类别码 | 类别 | 消息中的识别标签 | SIS 依据 | |---|---|---|---| diff --git a/docs/implementation.md b/docs/implementation.md index fc2f848..5a2454e 100644 --- a/docs/implementation.md +++ b/docs/implementation.md @@ -482,3 +482,26 @@ SIS 声明的上游忽略与截断口径见 `SIS:3.1`、`SIS:3.2`、`SIS:3.4`、 ### 13.5 admin-api 下游读取边界 SIS → 本网关 → PG → admin-api(`C-10`)。本网关不调用 admin-api。Oracle 待 `Q14` 验证。MySQL 只是信箱。 + +### 13.6 物理表映射(`Q22`) + +逻辑视图仍为 `REF_MASTER`(键 `(RTYPE, RKEY)`);物理落 schema **`basicdata`** 的多张表,列名取自 admin-api secondary 实体 `@Table`/`@Column`,**不**使用单表 `REF_MASTER`。Flyway:`V6__basicdata_ref_data.sql`;联调栈同名表见 `deploy/integration/postgres-init/02-basicdata-schema.sql`。 + +| `RTYPE` | `RKEY` 标签 | 物理表 | 主键 / 业务键列 | SIS → 列(常用) | +|---|---|---|---|---| +| `COUL` | `COUC` | `basicdata.sys_country` | `COUNTRY_ID` / `COUNTRY_CODE` | `COUC`→`COUNTRY_CODE`;`COUN`→`COUNTRY_NAME_ENG`;`CNMC`→`COUNTRY_NAME_CHN` | +| `ARPT` | `ITCD` | `basicdata.sys_airport` | `AIRPORT_ID` / `AIRPORT_CODE_IATA` | `ITCD`→`AIRPORT_CODE_IATA`;`ICCD`→`AIRPORT_CODE_ICAO`;`ANAM`/`ANMC`→英/中名;`HAUL`→`HAUL`;`ATYP`→`INTERNATIONAL_FLAG` | +| `AIRL` | `ITOP` | `basicdata.sys_airline` | `AIRLINE_ID` / `AIRLINE_CODE_IATA` | `ITOP`/`ICOP`→IATA/ICAO;`ONAM`/`ONMC`→名称;`DORI`→`INTERNATIONAL_FLAG` | +| `AIRC` | `ITAT` | `basicdata.sys_aircrafttype` | `AIRCRAFT_TYPE_CODE` | `ITAT`/`ICAT`→IATA/ICAO;`DESC`→`AIRCRAFT_TYPE_NAME`;`MAXP`/`MFWT`/`MTWT`/`MABR`→限额字段 | +| `REGN` | `RNUM` | `basicdata.sys_aircraft` | `AIRCRAFT_ID` / `AIRCRAFT_NUMBER` | `RNUM`→`AIRCRAFT_NUMBER`;`ITAT`→`AIRCRAFT_TYPE_CODE`;`ACAL`→`AIRLINE_CODE` | +| `ORGN` | `OGID` | `basicdata.sys_flightagent` | `FLIGHTAGENT_ID` | `OGID`→`FLIGHTAGENT_ID`/`OGID`;`ONAM`→`FLIGHTAGENT_NAME` | +| `FLTL` | `FTYP` | `basicdata.sys_flight_type` | `FLIGHT_TYPE_CODE` | `FTYP`→码;`FDES`/`FDSC`→中英文名;`CTYP`→`FLIGHT_TYPE_CAA_CODE` | +| `TLST` | `TCOD` | `basicdata.orms_terminal` | `TERMINAL_CODE` | `TCOD`/`TNAM`→码与名称 | +| `GLST` | `GCOD` | `basicdata.orms_gate` | `GATE_CODE` | `GCOD`/`GTNM`/`GCAT`/`GTML`→门码、名称、类别、航站楼 | +| `SLST` | `SCOD` | `basicdata.orms_stand` | `STAND_CODE` | `SCOD`/`STNM`/`STML`/`SATC` | +| `CLST` | `CCOD` | `basicdata.orms_checkindesk` | `CHECKINDESK_CODE` | `CCOD`/`CTNM`/`CCAT`(勿与 `CHLT` 混类) | +| `BLST` | `BCOD` | `basicdata.orms_carousel` | `CAROUSEL_CODE` | `BCOD`/`BTNM`/`BCAT`/`BTML` | +| `CHLT` | `CCOD` | `basicdata.orms_chut` | `CHUT_CODE` | `CCOD`/`CHNM`/`CCAT`/`CTML` | +| `RSTA` | `RTYP`+`RSID` | 上表资源行 | 资源码列 | `STAT` `E`/`D`→`GATE`/`STND` 用数值态,转盘/柜台/滑槽用 `E`/`D` 字符串列;不以包内缺席删其他资源态 | + +数值主键表在首次写入时用 `(RTYPE,RKEY)` 稳定派生 `*_ID`(同键可重复 upsert);varchar 主键表直接以 SIS 业务键为 PK。 diff --git a/docs/reference.md b/docs/reference.md index f89f27f..b2b746e 100644 --- a/docs/reference.md +++ b/docs/reference.md @@ -138,7 +138,6 @@ | `PROTOCOL` | 整份报文未通过业务校验,或运营日不一致 | 整份不写数据库,记 `DEAD` | 否 | | `CODEC_ERROR` | 解码程序尚不能识别报文结构 | 记 `FAILED`,等待后重试 | 是 | | `UNSUPPORTED`(合法未知类型) | 报文种类可识别,但没有对应处理程序(`MsgKind.Unsupported`) | 记 `SKIPPED`,写回信箱为已处理(`US-03` AC2) | 否(终态,不进重放白名单) | -| `UNSUPPORTED`(参考数据) | 参考数据类别已识别,但处理器尚未落地(`MsgKind.RefData`) | 记 `FAILED(UNSUPPORTED)` 并重试 | 是 | | `INFRA` | 数据库、网络或程序执行出错 | 记 `FAILED`,等待后重试 | 是 | | `EXHAUSTED` | 自动重试次数已用尽 | 记 `DEAD`;原错误类别被覆盖,原因留在 `LAST_ERROR` | 是 | diff --git a/docs/specification.md b/docs/specification.md index 8912fca..bc4e255 100644 --- a/docs/specification.md +++ b/docs/specification.md @@ -134,7 +134,7 @@ | Q19 | — | (未分配) | — | | Q20 | — | (未分配) | — | | Q21 | 本系统 | `GET /all/flights` HTTP 约定 | 见 `C-11` | -| Q22 | 本系统 | 静态参考数据对应哪些表 | 见 `C-10`、接口契约 | +| Q22 | 已定 | 静态参考数据对应哪些表 | `C-10`;物理为 schema `basicdata` 多表(非 `REF_MASTER` 单表),映射见 `implementation.md`「静态参考数据」 | | Q23 | 本系统 | Elasticsearch 历史怎么写、保留多久 | 见 `US-14`、`G-FLIGHT-HIST-RETENTION` | | Q24 | 本系统 | `REQ_TRACK` 已结案记录保留多久 | 见 `G-REQ-TRACK-RETENTION` | | Q25 | 本系统 | 季度计划从哪来、什么格式 | SIS 只有日计划;旧系统读 Oracle 季度表 | @@ -154,7 +154,7 @@ | `G-MAFL` | 主航班共享列表未做 | `US-06` AC2 | | `G-PROC-CLEANUP` | 处理记录到期清理未做 | `US-11` | | `G-REDIS-PROJECTION` | Redis 快照写入与删除未做:投影端口与三步提交时序已就位,真实 Redis 客户端未接通,日计划覆盖范围的缺席清扫也未做 | `US-05`、`US-06`、`US-07`、`US-12` | -| `G-REF-DATA` | 静态参考数据处理与 admin-api 直读未做 | `US-13` | +| `G-REF-DATA` | admin-api 生产侧只读接入与联调验收未闭合 | `US-13`;本网关落库与 `ReferenceDataProcessor` 已做(ACM2-93) | | `G-REQ-OPEN-UNIQUE` | 同类型未处理完时不允许再发,未做 | `US-09` | | `G-REQ-TRACK` | 出站请求跟踪未做 | `US-09` | | `G-REQ-TRACK-RETENTION` | `REQ_TRACK` 已结案保留期未定(`Q24`) | `US-09` | diff --git a/src/main/kotlin/com/gzzn/omms/msgexchange/codec/JacksonXmlCodec.kt b/src/main/kotlin/com/gzzn/omms/msgexchange/codec/JacksonXmlCodec.kt index b18f411..716ea30 100644 --- a/src/main/kotlin/com/gzzn/omms/msgexchange/codec/JacksonXmlCodec.kt +++ b/src/main/kotlin/com/gzzn/omms/msgexchange/codec/JacksonXmlCodec.kt @@ -67,8 +67,7 @@ class JacksonXmlCodec : XmlCodec { is MsgKind.Schd -> msg.schd?.let(SisWireMapper::scheduleBody) is MsgKind.Flop -> msg.flop?.let(SisWireMapper::flopPayload) MsgKind.Fdel -> msg.flop?.let(SisWireMapper::flopPayload) - // 参考数据类别报文的正文结构随 13 类各异(ACM2-93 实装 ReferenceDataProcessor 时绑定) - is MsgKind.RefData -> null + is MsgKind.RefData -> RefDataWireMapper.body(type, styp, msg) MsgKind.Eror -> msg.eror?.let { val seqs = it.seqs ?: return DecodeResult.Err(DecodeFailure(ErrorClass.MALFORMED, "missing-eror-seqs")) val typs = it.typs?.trim()?.takeIf { t -> t.isNotEmpty() } diff --git a/src/main/kotlin/com/gzzn/omms/msgexchange/codec/RefDataWire.kt b/src/main/kotlin/com/gzzn/omms/msgexchange/codec/RefDataWire.kt new file mode 100644 index 0000000..312a110 --- /dev/null +++ b/src/main/kotlin/com/gzzn/omms/msgexchange/codec/RefDataWire.kt @@ -0,0 +1,205 @@ +package com.gzzn.omms.msgexchange.codec + +import com.fasterxml.jackson.annotation.JsonIgnoreProperties +import com.fasterxml.jackson.dataformat.xml.annotation.JacksonXmlElementWrapper +import com.fasterxml.jackson.dataformat.xml.annotation.JacksonXmlProperty +import com.gzzn.omms.msgexchange.domain.ref.RefDataBody +import com.gzzn.omms.msgexchange.domain.ref.RefDataRecord +import com.gzzn.omms.msgexchange.domain.ref.ResourceStatusRecord + +@JsonIgnoreProperties(ignoreUnknown = true) +data class CoulRecordXml( + @param:JacksonXmlProperty(localName = "COUC") val couc: String? = null, + @param:JacksonXmlProperty(localName = "COUN") val coun: String? = null, + @param:JacksonXmlProperty(localName = "CNMC") val cnmc: String? = null, +) + +@JsonIgnoreProperties(ignoreUnknown = true) +data class ArptRecordXml( + @param:JacksonXmlProperty(localName = "ITCD") val itcd: String? = null, + @param:JacksonXmlProperty(localName = "ICCD") val iccd: String? = null, + @param:JacksonXmlProperty(localName = "ANAM") val anam: String? = null, + @param:JacksonXmlProperty(localName = "ANMC") val anmc: String? = null, + @param:JacksonXmlProperty(localName = "CTRY") val ctry: String? = null, + @param:JacksonXmlProperty(localName = "ACTY") val acty: String? = null, + @param:JacksonXmlProperty(localName = "ATYP") val atyp: String? = null, + @param:JacksonXmlProperty(localName = "HAUL") val haul: String? = null, +) + +@JsonIgnoreProperties(ignoreUnknown = true) +data class AirlRecordXml( + @param:JacksonXmlProperty(localName = "ITOP") val itop: String? = null, + @param:JacksonXmlProperty(localName = "ICOP") val icop: String? = null, + @param:JacksonXmlProperty(localName = "ONAM") val onam: String? = null, + @param:JacksonXmlProperty(localName = "ONMC") val onmc: String? = null, + @param:JacksonXmlProperty(localName = "CTRY") val ctry: String? = null, + @param:JacksonXmlProperty(localName = "DORI") val dori: String? = null, +) + +@JsonIgnoreProperties(ignoreUnknown = true) +data class AircRecordXml( + @param:JacksonXmlProperty(localName = "ITAT") val itat: String? = null, + @param:JacksonXmlProperty(localName = "ICAT") val icat: String? = null, + @param:JacksonXmlProperty(localName = "DESC") val desc: String? = null, + @param:JacksonXmlProperty(localName = "CDSC") val cdsc: String? = null, + @param:JacksonXmlProperty(localName = "CHAP") val chap: String? = null, + @param:JacksonXmlProperty(localName = "MAXP") val maxp: String? = null, + @param:JacksonXmlProperty(localName = "MFWT") val mfwt: String? = null, + @param:JacksonXmlProperty(localName = "MTWT") val mtwt: String? = null, + @param:JacksonXmlProperty(localName = "MABR") val mabr: String? = null, +) + +@JsonIgnoreProperties(ignoreUnknown = true) +data class RegnRecordXml( + @param:JacksonXmlProperty(localName = "RNUM") val rnum: String? = null, + @param:JacksonXmlProperty(localName = "ITAT") val itat: String? = null, + @param:JacksonXmlProperty(localName = "OWID") val owid: String? = null, + @param:JacksonXmlProperty(localName = "ACAL") val acal: String? = null, + @param:JacksonXmlProperty(localName = "MAXP") val maxp: String? = null, + @param:JacksonXmlProperty(localName = "MFWT") val mfwt: String? = null, + @param:JacksonXmlProperty(localName = "MTWT") val mtwt: String? = null, +) + +@JsonIgnoreProperties(ignoreUnknown = true) +data class OrgnRecordXml( + @param:JacksonXmlProperty(localName = "OGID") val ogid: String? = null, + @param:JacksonXmlProperty(localName = "ONAM") val onam: String? = null, + @param:JacksonXmlProperty(localName = "ONMC") val onmc: String? = null, +) + +@JsonIgnoreProperties(ignoreUnknown = true) +data class FltlRecordXml( + @param:JacksonXmlProperty(localName = "FTYP") val ftyp: String? = null, + @param:JacksonXmlProperty(localName = "CTYP") val ctyp: String? = null, + @param:JacksonXmlProperty(localName = "FDES") val fdes: String? = null, + @param:JacksonXmlProperty(localName = "FDSC") val fdsc: String? = null, + @param:JacksonXmlProperty(localName = "FCML") val fcml: String? = null, +) + +@JsonIgnoreProperties(ignoreUnknown = true) +data class TlstRecordXml( + @param:JacksonXmlProperty(localName = "TCOD") val tcod: String? = null, + @param:JacksonXmlProperty(localName = "TNAM") val tnam: String? = null, + @param:JacksonXmlProperty(localName = "TNMC") val tnmc: String? = null, +) + +@JsonIgnoreProperties(ignoreUnknown = true) +data class GlstRecordXml( + @param:JacksonXmlProperty(localName = "GCOD") val gcod: String? = null, + @param:JacksonXmlProperty(localName = "GTNM") val gtnm: String? = null, + @param:JacksonXmlProperty(localName = "GNMC") val gnmc: String? = null, + @param:JacksonXmlProperty(localName = "GCAT") val gcat: String? = null, + @param:JacksonXmlProperty(localName = "GTML") val gtml: String? = null, +) + +@JsonIgnoreProperties(ignoreUnknown = true) +data class SlstRecordXml( + @param:JacksonXmlProperty(localName = "SCOD") val scod: String? = null, + @param:JacksonXmlProperty(localName = "STNM") val stnm: String? = null, + @param:JacksonXmlProperty(localName = "SNMC") val snmc: String? = null, + @param:JacksonXmlProperty(localName = "STML") val stml: String? = null, + @param:JacksonXmlProperty(localName = "SATC") val satc: String? = null, +) + +@JsonIgnoreProperties(ignoreUnknown = true) +data class ClstRecordXml( + @param:JacksonXmlProperty(localName = "CCOD") val ccod: String? = null, + @param:JacksonXmlProperty(localName = "CTNM") val ctnm: String? = null, + @param:JacksonXmlProperty(localName = "CNMC") val cnmc: String? = null, + @param:JacksonXmlProperty(localName = "CCAT") val ccat: String? = null, + @param:JacksonXmlProperty(localName = "CTML") val ctml: String? = null, + @param:JacksonXmlProperty(localName = "CTRA") val ctra: String? = null, +) + +@JsonIgnoreProperties(ignoreUnknown = true) +data class BlstRecordXml( + @param:JacksonXmlProperty(localName = "BCOD") val bcod: String? = null, + @param:JacksonXmlProperty(localName = "BTNM") val btnm: String? = null, + @param:JacksonXmlProperty(localName = "BNMC") val bnmc: String? = null, + @param:JacksonXmlProperty(localName = "BCAT") val bcat: String? = null, + @param:JacksonXmlProperty(localName = "BTML") val btml: String? = null, +) + +@JsonIgnoreProperties(ignoreUnknown = true) +data class ChltRecordXml( + @param:JacksonXmlProperty(localName = "CCOD") val ccod: String? = null, + @param:JacksonXmlProperty(localName = "CHNM") val chnm: String? = null, + @param:JacksonXmlProperty(localName = "CNMC") val cnmc: String? = null, + @param:JacksonXmlProperty(localName = "CTML") val ctml: String? = null, + @param:JacksonXmlProperty(localName = "CCAT") val ccat: String? = null, +) + +@JsonIgnoreProperties(ignoreUnknown = true) +data class RstaRecordXml( + @param:JacksonXmlProperty(localName = "RTYP") val rtyp: String? = null, + @param:JacksonXmlProperty(localName = "RSID") val rsid: String? = null, + @param:JacksonXmlProperty(localName = "STAT") val stat: String? = null, + @param:JacksonXmlProperty(localName = "RDST") val rdst: String? = null, + @param:JacksonXmlProperty(localName = "RDET") val rdet: String? = null, +) + +object RefDataWireMapper { + fun body(category: String, styp: String, msg: SisMessageXml): RefDataBody { + val records = when (category) { + "COUL" -> msg.coul.map { recordOf("COUC" to it.couc, "COUN" to it.coun, "CNMC" to it.cnmc) } + "ARPT" -> msg.arpt.map { + recordOf( + "ITCD" to it.itcd, "ICCD" to it.iccd, "ANAM" to it.anam, "ANMC" to it.anmc, + "CTRY" to it.ctry, "ACTY" to it.acty, "ATYP" to it.atyp, "HAUL" to it.haul, + ) + } + "AIRL" -> msg.airl.map { + recordOf("ITOP" to it.itop, "ICOP" to it.icop, "ONAM" to it.onam, "ONMC" to it.onmc, "CTRY" to it.ctry, "DORI" to it.dori) + } + "AIRC" -> msg.airc.map { + recordOf( + "ITAT" to it.itat, "ICAT" to it.icat, "DESC" to it.desc, "CDSC" to it.cdsc, + "CHAP" to it.chap, "MAXP" to it.maxp, "MFWT" to it.mfwt, "MTWT" to it.mtwt, "MABR" to it.mabr, + ) + } + "REGN" -> msg.regn.map { + recordOf( + "RNUM" to it.rnum, "ITAT" to it.itat, "OWID" to it.owid, "ACAL" to it.acal, + "MAXP" to it.maxp, "MFWT" to it.mfwt, "MTWT" to it.mtwt, + ) + } + "ORGN" -> msg.orgn.map { recordOf("OGID" to it.ogid, "ONAM" to it.onam, "ONMC" to it.onmc) } + "FLTL" -> msg.fltl.map { + recordOf("FTYP" to it.ftyp, "CTYP" to it.ctyp, "FDES" to it.fdes, "FDSC" to it.fdsc, "FCML" to it.fcml) + } + "TLST" -> msg.tlst.map { recordOf("TCOD" to it.tcod, "TNAM" to it.tnam, "TNMC" to it.tnmc) } + "GLST" -> msg.glst.map { + recordOf("GCOD" to it.gcod, "GTNM" to it.gtnm, "GNMC" to it.gnmc, "GCAT" to it.gcat, "GTML" to it.gtml) + } + "SLST" -> msg.slst.map { + recordOf("SCOD" to it.scod, "STNM" to it.stnm, "SNMC" to it.snmc, "STML" to it.stml, "SATC" to it.satc) + } + "CLST" -> msg.clst.map { + recordOf( + "CCOD" to it.ccod, "CTNM" to it.ctnm, "CNMC" to it.cnmc, + "CCAT" to it.ccat, "CTML" to it.ctml, "CTRA" to it.ctra, + ) + } + "BLST" -> msg.blst.map { + recordOf("BCOD" to it.bcod, "BTNM" to it.btnm, "BNMC" to it.bnmc, "BCAT" to it.bcat, "BTML" to it.btml) + } + "CHLT" -> msg.chlt.map { + recordOf("CCOD" to it.ccod, "CHNM" to it.chnm, "CNMC" to it.cnmc, "CTML" to it.ctml, "CCAT" to it.ccat) + } + "RSTA" -> emptyList() + else -> emptyList() + } + val rsta = msg.rsta.mapNotNull { + val rtyp = it.rtyp?.trim()?.uppercase()?.takeIf { r -> r.isNotEmpty() } ?: return@mapNotNull null + val rsid = it.rsid?.trim()?.takeIf { r -> r.isNotEmpty() } ?: return@mapNotNull null + val stat = it.stat?.trim()?.uppercase()?.takeIf { r -> r.isNotEmpty() } ?: return@mapNotNull null + ResourceStatusRecord(rtyp, rsid, stat, blankToNull(it.rdst), blankToNull(it.rdet)) + } + return RefDataBody(category, styp, records, rsta) + } + + private fun recordOf(vararg pairs: Pair): RefDataRecord = + RefDataRecord(pairs.associate { (k, v) -> k to (v ?: "") }) + + private fun blankToNull(v: String?) = v?.trim()?.takeIf { it.isNotEmpty() } +} diff --git a/src/main/kotlin/com/gzzn/omms/msgexchange/codec/SisMessageBody.kt b/src/main/kotlin/com/gzzn/omms/msgexchange/codec/SisMessageBody.kt index 48dec9f..6921508 100644 --- a/src/main/kotlin/com/gzzn/omms/msgexchange/codec/SisMessageBody.kt +++ b/src/main/kotlin/com/gzzn/omms/msgexchange/codec/SisMessageBody.kt @@ -59,6 +59,20 @@ data class SisMessageXml( @param:JacksonXmlProperty(localName = "EROR") val eror: ErorXml? = null, @param:JacksonXmlProperty(localName = "RQFD") val rqfd: RqfdXml? = null, @param:JacksonXmlProperty(localName = "RQRD") val rqrd: RqrdXml? = null, + @param:JacksonXmlElementWrapper(useWrapping = false) @param:JacksonXmlProperty(localName = "COUL") val coul: List = emptyList(), + @param:JacksonXmlElementWrapper(useWrapping = false) @param:JacksonXmlProperty(localName = "ARPT") val arpt: List = emptyList(), + @param:JacksonXmlElementWrapper(useWrapping = false) @param:JacksonXmlProperty(localName = "AIRL") val airl: List = emptyList(), + @param:JacksonXmlElementWrapper(useWrapping = false) @param:JacksonXmlProperty(localName = "AIRC") val airc: List = emptyList(), + @param:JacksonXmlElementWrapper(useWrapping = false) @param:JacksonXmlProperty(localName = "REGN") val regn: List = emptyList(), + @param:JacksonXmlElementWrapper(useWrapping = false) @param:JacksonXmlProperty(localName = "ORGN") val orgn: List = emptyList(), + @param:JacksonXmlElementWrapper(useWrapping = false) @param:JacksonXmlProperty(localName = "FLTL") val fltl: List = emptyList(), + @param:JacksonXmlElementWrapper(useWrapping = false) @param:JacksonXmlProperty(localName = "TLST") val tlst: List = emptyList(), + @param:JacksonXmlElementWrapper(useWrapping = false) @param:JacksonXmlProperty(localName = "GLST") val glst: List = emptyList(), + @param:JacksonXmlElementWrapper(useWrapping = false) @param:JacksonXmlProperty(localName = "SLST") val slst: List = emptyList(), + @param:JacksonXmlElementWrapper(useWrapping = false) @param:JacksonXmlProperty(localName = "CLST") val clst: List = emptyList(), + @param:JacksonXmlElementWrapper(useWrapping = false) @param:JacksonXmlProperty(localName = "BLST") val blst: List = emptyList(), + @param:JacksonXmlElementWrapper(useWrapping = false) @param:JacksonXmlProperty(localName = "CHLT") val chlt: List = emptyList(), + @param:JacksonXmlElementWrapper(useWrapping = false) @param:JacksonXmlProperty(localName = "RSTA") val rsta: List = emptyList(), ) @JsonIgnoreProperties(ignoreUnknown = true) diff --git a/src/main/kotlin/com/gzzn/omms/msgexchange/domain/ref/RefDataBody.kt b/src/main/kotlin/com/gzzn/omms/msgexchange/domain/ref/RefDataBody.kt new file mode 100644 index 0000000..43a60c0 --- /dev/null +++ b/src/main/kotlin/com/gzzn/omms/msgexchange/domain/ref/RefDataBody.kt @@ -0,0 +1,36 @@ +package com.gzzn.omms.msgexchange.domain.ref + +/** 解码后的静态参考数据载荷(`US-13`);一条 MSG 对应一个类别码 [category]。 */ +data class RefDataBody( + val category: String, + val styp: String, + val records: List, + val resourceStatuses: List = emptyList(), +) + +/** 一条参考记录:键为 SIS 叶子标签,空串表示「无值」语义(`US-13` AC4)。 */ +data class RefDataRecord(val fields: Map) + +/** `TYPE=RSTA` 资源状态行(`SIS:3.14`)。 */ +data class ResourceStatusRecord( + val rtyp: String, + val rsid: String, + val stat: String, + val rdst: String? = null, + val rdet: String? = null, +) + +enum class RefDataMode { + FULL_REPLACE, + ADD, + UPDATE, + DELETE, +} + +fun refDataModeOf(styp: String): RefDataMode? = when (styp.uppercase()) { + "DNLD", "RESP" -> RefDataMode.FULL_REPLACE + "ADD" -> RefDataMode.ADD + "UPD" -> RefDataMode.UPDATE + "DEL" -> RefDataMode.DELETE + else -> null +} diff --git a/src/main/kotlin/com/gzzn/omms/msgexchange/domain/ref/RefDataFieldMapping.kt b/src/main/kotlin/com/gzzn/omms/msgexchange/domain/ref/RefDataFieldMapping.kt new file mode 100644 index 0000000..ced9cb0 --- /dev/null +++ b/src/main/kotlin/com/gzzn/omms/msgexchange/domain/ref/RefDataFieldMapping.kt @@ -0,0 +1,37 @@ +package com.gzzn.omms.msgexchange.domain.ref + +/** SIS 标签 → `basicdata` 列(`Q22` 映射表;列名取自 admin-api 实体注解)。 */ +internal object RefDataFieldMapping { + fun businessKey(category: String, rec: RefDataRecord): String = when (category) { + "COUL" -> tag(rec, "COUC") + "ARPT" -> tag(rec, "ITCD") + "AIRL" -> tag(rec, "ITOP") + "AIRC" -> tag(rec, "ITAT") + "REGN" -> tag(rec, "RNUM") + "ORGN" -> tag(rec, "OGID") + "FLTL" -> tag(rec, "FTYP") + "TLST" -> tag(rec, "TCOD") + "GLST" -> tag(rec, "GCOD") + "SLST" -> tag(rec, "SCOD") + "CLST" -> tag(rec, "CCOD") + "BLST" -> tag(rec, "BCOD") + "CHLT" -> tag(rec, "CCOD") + else -> error("no key for $category") + } + + fun stableNumericId(category: String, key: String): Long { + val h = "$category:$key".hashCode().toLong() + return (if (h == Long.MIN_VALUE) 1L else kotlin.math.abs(h)).coerceAtLeast(1L) + } + + private fun tag(rec: RefDataRecord, name: String) = rec.fields[name]?.trim().orEmpty() + + fun blankToNull(raw: String?): String? = raw?.trim()?.takeIf { it.isNotEmpty() } + + fun domIntFlag(code: String?): java.math.BigDecimal? = when (blankToNull(code)?.uppercase()) { + "I" -> java.math.BigDecimal.ONE + "D" -> java.math.BigDecimal.ZERO + "R" -> java.math.BigDecimal.valueOf(2) + else -> null + } +} diff --git a/src/main/kotlin/com/gzzn/omms/msgexchange/domain/ref/RefDataValidator.kt b/src/main/kotlin/com/gzzn/omms/msgexchange/domain/ref/RefDataValidator.kt new file mode 100644 index 0000000..41c94f2 --- /dev/null +++ b/src/main/kotlin/com/gzzn/omms/msgexchange/domain/ref/RefDataValidator.kt @@ -0,0 +1,50 @@ +package com.gzzn.omms.msgexchange.domain.ref + +/** 类别校验:失败时整类不写入(`US-13` AC2)。 */ +object RefDataValidator { + private val KEY_TAG = mapOf( + "COUL" to "COUC", + "ARPT" to "ITCD", + "AIRL" to "ITOP", + "AIRC" to "ITAT", + "REGN" to "RNUM", + "ORGN" to "OGID", + "FLTL" to "FTYP", + "TLST" to "TCOD", + "GLST" to "GCOD", + "SLST" to "SCOD", + "CLST" to "CCOD", + "BLST" to "BCOD", + "CHLT" to "CCOD", + ) + + private val ALLOWED_RSTA_RTYP = setOf("BELT", "CNTR", "GATE", "STND", "CHUT") + private val ALLOWED_STAT = setOf("E", "D") + + fun validate(category: String, body: RefDataBody, mode: RefDataMode): String? { + if (category == "RSTA") { + if (body.resourceStatuses.isEmpty()) return "rsta-empty" + body.resourceStatuses.forEach { r -> + if (r.rtyp !in ALLOWED_RSTA_RTYP) return "rsta-rtyp:${r.rtyp}" + if (r.stat !in ALLOWED_STAT) return "rsta-stat:${r.stat}" + if (r.rsid.isBlank()) return "rsta-rsid-missing" + } + return null + } + val keyTag = KEY_TAG[category] ?: return "unknown-category:$category" + if (mode == RefDataMode.FULL_REPLACE) { + body.records.forEach { rec -> + validateRecord(keyTag, rec)?.let { return it } + } + return null + } + if (body.records.size != 1) return "incremental-needs-single-record" + return validateRecord(keyTag, body.records.single()) + } + + private fun validateRecord(keyTag: String, rec: RefDataRecord): String? { + val key = rec.fields[keyTag]?.trim().orEmpty() + if (key.isEmpty()) return "missing-key:$keyTag" + return null + } +} diff --git a/src/main/kotlin/com/gzzn/omms/msgexchange/infra/persistence/RefDataRepository.kt b/src/main/kotlin/com/gzzn/omms/msgexchange/infra/persistence/RefDataRepository.kt new file mode 100644 index 0000000..5b268af --- /dev/null +++ b/src/main/kotlin/com/gzzn/omms/msgexchange/infra/persistence/RefDataRepository.kt @@ -0,0 +1,15 @@ +package com.gzzn.omms.msgexchange.infra.persistence + +import com.gzzn.omms.msgexchange.domain.ref.RefDataRecord +import com.gzzn.omms.msgexchange.domain.ref.ResourceStatusRecord + +/** 静态参考数据表组(schema `basicdata`)写入端口(`US-13`、`C-10`)。 */ +interface RefDataRepository { + fun replaceCategory(category: String, records: List) + + fun upsertCategory(category: String, record: RefDataRecord) + + fun deleteCategory(category: String, businessKey: String) + + fun applyResourceStatuses(updates: List) +} diff --git a/src/main/kotlin/com/gzzn/omms/msgexchange/infra/persistence/jdbc/JdbcRefDataRepository.kt b/src/main/kotlin/com/gzzn/omms/msgexchange/infra/persistence/jdbc/JdbcRefDataRepository.kt new file mode 100644 index 0000000..d9300d1 --- /dev/null +++ b/src/main/kotlin/com/gzzn/omms/msgexchange/infra/persistence/jdbc/JdbcRefDataRepository.kt @@ -0,0 +1,419 @@ +package com.gzzn.omms.msgexchange.infra.persistence.jdbc + +import com.gzzn.omms.msgexchange.domain.ref.RefDataFieldMapping +import com.gzzn.omms.msgexchange.domain.ref.RefDataRecord +import com.gzzn.omms.msgexchange.domain.ref.ResourceStatusRecord +import com.gzzn.omms.msgexchange.infra.persistence.RefDataRepository +import io.micronaut.context.annotation.Requires +import jakarta.inject.Singleton +import java.math.BigDecimal +import javax.sql.DataSource + +@Singleton +@Requires(property = "datasources.default.enabled", value = "true") +@Requires(missingProperty = "msgx.stubs") +class JdbcRefDataRepository( + private val ds: DataSource, +) : RefDataRepository { + + override fun replaceCategory(category: String, records: List) { + ds.withTransaction { + deleteAll(category) + records.forEach { upsertCategory(category, it) } + } + } + + override fun upsertCategory(category: String, record: RefDataRecord) { + when (category) { + "COUL" -> upsertCountry(record) + "ARPT" -> upsertAirport(record) + "AIRL" -> upsertAirline(record) + "AIRC" -> upsertAircraftType(record) + "REGN" -> upsertAircraft(record) + "ORGN" -> upsertFlightAgent(record) + "FLTL" -> upsertFlightType(record) + "TLST" -> upsertTerminal(record) + "GLST" -> upsertGate(record) + "SLST" -> upsertStand(record) + "CLST" -> upsertCheckinDesk(record) + "BLST" -> upsertCarousel(record) + "CHLT" -> upsertChut(record) + else -> error("unsupported category $category") + } + } + + override fun deleteCategory(category: String, businessKey: String) { + val sql = when (category) { + "COUL" -> "DELETE FROM basicdata.sys_country WHERE COUNTRY_CODE = ?" + "ARPT" -> "DELETE FROM basicdata.sys_airport WHERE AIRPORT_CODE_IATA = ?" + "AIRL" -> "DELETE FROM basicdata.sys_airline WHERE AIRLINE_CODE_IATA = ?" + "AIRC" -> "DELETE FROM basicdata.sys_aircrafttype WHERE AIRCRAFT_TYPE_CODE = ?" + "REGN" -> "DELETE FROM basicdata.sys_aircraft WHERE AIRCRAFT_NUMBER = ?" + "ORGN" -> "DELETE FROM basicdata.sys_flightagent WHERE FLIGHTAGENT_ID = ?" + "FLTL" -> "DELETE FROM basicdata.sys_flight_type WHERE FLIGHT_TYPE_CODE = ?" + "TLST" -> "DELETE FROM basicdata.orms_terminal WHERE TERMINAL_CODE = ?" + "GLST" -> "DELETE FROM basicdata.orms_gate WHERE GATE_CODE = ?" + "SLST" -> "DELETE FROM basicdata.orms_stand WHERE STAND_CODE = ?" + "CLST" -> "DELETE FROM basicdata.orms_checkindesk WHERE CHECKINDESK_CODE = ?" + "BLST" -> "DELETE FROM basicdata.orms_carousel WHERE CAROUSEL_CODE = ?" + "CHLT" -> "DELETE FROM basicdata.orms_chut WHERE CHUT_CODE = ?" + else -> error("unsupported category $category") + } + ds.update(sql) { ps -> ps.setString(1, businessKey) } + } + + override fun applyResourceStatuses(updates: List) { + updates.forEach { u -> + when (u.rtyp) { + "GATE" -> ds.update( + "UPDATE basicdata.orms_gate SET GATE_STATUS = ? WHERE GATE_CODE = ?", + ) { ps -> + ps.setBigDecimal(1, numericStatus(u.stat)) + ps.setString(2, u.rsid) + } + "STND" -> ds.update( + "UPDATE basicdata.orms_stand SET STAND_STATUS = ? WHERE STAND_CODE = ?", + ) { ps -> + ps.setBigDecimal(1, numericStatus(u.stat)) + ps.setString(2, u.rsid) + } + "BELT" -> ds.update( + "UPDATE basicdata.orms_carousel SET CAROUSEL_STATUS = ? WHERE CAROUSEL_CODE = ?", + ) { ps -> + ps.setString(1, u.stat) + ps.setString(2, u.rsid) + } + "CNTR" -> ds.update( + "UPDATE basicdata.orms_checkindesk SET CHECKINDESK_STATUS = ? WHERE CHECKINDESK_CODE = ?", + ) { ps -> + ps.setString(1, u.stat) + ps.setString(2, u.rsid) + } + "CHUT" -> ds.update( + "UPDATE basicdata.orms_chut SET CHUT_STATUS = ? WHERE CHUT_CODE = ?", + ) { ps -> + ps.setString(1, u.stat) + ps.setString(2, u.rsid) + } + } + } + } + + private fun deleteAll(category: String) { + val sql = when (category) { + "COUL" -> "DELETE FROM basicdata.sys_country" + "ARPT" -> "DELETE FROM basicdata.sys_airport" + "AIRL" -> "DELETE FROM basicdata.sys_airline" + "AIRC" -> "DELETE FROM basicdata.sys_aircrafttype" + "REGN" -> "DELETE FROM basicdata.sys_aircraft" + "ORGN" -> "DELETE FROM basicdata.sys_flightagent" + "FLTL" -> "DELETE FROM basicdata.sys_flight_type" + "TLST" -> "DELETE FROM basicdata.orms_terminal" + "GLST" -> "DELETE FROM basicdata.orms_gate" + "SLST" -> "DELETE FROM basicdata.orms_stand" + "CLST" -> "DELETE FROM basicdata.orms_checkindesk" + "BLST" -> "DELETE FROM basicdata.orms_carousel" + "CHLT" -> "DELETE FROM basicdata.orms_chut" + else -> error("unsupported category $category") + } + ds.update(sql) {} + } + + private fun upsertCountry(rec: RefDataRecord) { + val f = rec.fields + val code = tag(f, "COUC") ?: return + val id = RefDataFieldMapping.stableNumericId("COUL", code) + ds.update( + """ + INSERT INTO basicdata.sys_country (COUNTRY_ID, COUNTRY_CODE, COUNTRY_NAME_ENG, COUNTRY_NAME_CHN) + VALUES (?, ?, ?, ?) + ON CONFLICT (COUNTRY_ID) DO UPDATE SET + COUNTRY_CODE = EXCLUDED.COUNTRY_CODE, + COUNTRY_NAME_ENG = EXCLUDED.COUNTRY_NAME_ENG, + COUNTRY_NAME_CHN = EXCLUDED.COUNTRY_NAME_CHN + """.trimIndent(), + ) { ps -> + ps.setLong(1, id) + ps.setString(2, code) + ps.setString(3, tag(f, "COUN")) + ps.setString(4, tag(f, "CNMC")) + } + } + + private fun upsertAirport(rec: RefDataRecord) { + val f = rec.fields + val iata = tag(f, "ITCD") ?: return + val id = RefDataFieldMapping.stableNumericId("ARPT", iata) + ds.update( + """ + INSERT INTO basicdata.sys_airport ( + AIRPORT_ID, AIRPORT_CODE_IATA, AIRPORT_CODE_ICAO, AIRPORT_NAME_ENG, AIRPORT_NAME_CHN, HAUL, INTERNATIONAL_FLAG + ) VALUES (?, ?, ?, ?, ?, ?, ?) + ON CONFLICT (AIRPORT_ID) DO UPDATE SET + AIRPORT_CODE_IATA = EXCLUDED.AIRPORT_CODE_IATA, + AIRPORT_CODE_ICAO = EXCLUDED.AIRPORT_CODE_ICAO, + AIRPORT_NAME_ENG = EXCLUDED.AIRPORT_NAME_ENG, + AIRPORT_NAME_CHN = EXCLUDED.AIRPORT_NAME_CHN, + HAUL = EXCLUDED.HAUL, + INTERNATIONAL_FLAG = EXCLUDED.INTERNATIONAL_FLAG + """.trimIndent(), + ) { ps -> + ps.setLong(1, id) + ps.setString(2, iata) + ps.setString(3, tag(f, "ICCD")) + ps.setString(4, tag(f, "ANAM")) + ps.setString(5, tag(f, "ANMC")) + ps.setString(6, tag(f, "HAUL")) + ps.setObject(7, RefDataFieldMapping.domIntFlag(f["ATYP"])) + } + } + + private fun upsertAirline(rec: RefDataRecord) { + val f = rec.fields + val iata = tag(f, "ITOP") ?: return + val id = RefDataFieldMapping.stableNumericId("AIRL", iata) + ds.update( + """ + INSERT INTO basicdata.sys_airline ( + AIRLINE_ID, AIRLINE_CODE_IATA, AIRLINE_CODE_ICAO, AIRLINE_NAME_ENG, BRIEF_NAME_CHN, INTERNATIONAL_FLAG + ) VALUES (?, ?, ?, ?, ?, ?) + ON CONFLICT (AIRLINE_ID) DO UPDATE SET + AIRLINE_CODE_IATA = EXCLUDED.AIRLINE_CODE_IATA, + AIRLINE_CODE_ICAO = EXCLUDED.AIRLINE_CODE_ICAO, + AIRLINE_NAME_ENG = EXCLUDED.AIRLINE_NAME_ENG, + BRIEF_NAME_CHN = EXCLUDED.BRIEF_NAME_CHN, + INTERNATIONAL_FLAG = EXCLUDED.INTERNATIONAL_FLAG + """.trimIndent(), + ) { ps -> + ps.setLong(1, id) + ps.setString(2, iata) + ps.setString(3, tag(f, "ICOP")) + ps.setString(4, tag(f, "ONAM")) + ps.setString(5, tag(f, "ONMC")) + ps.setObject(6, RefDataFieldMapping.domIntFlag(f["DORI"])) + } + } + + private fun upsertAircraftType(rec: RefDataRecord) { + val f = rec.fields + val code = tag(f, "ITAT") ?: return + ds.update( + """ + INSERT INTO basicdata.sys_aircrafttype ( + AIRCRAFT_TYPE_CODE, AIRCRAFT_TYPE_CODE_IATA, AIRCRAFT_TYPE_CODE_ICAO, AIRCRAFT_TYPE_NAME, + CHAP, MAX_PASSENGERS, MAX_FREIGHT_WEIGHT, MAX_TAKEOFF_WEIGHT, MAX_AIRBRIDGES + ) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?) + ON CONFLICT (AIRCRAFT_TYPE_CODE) DO UPDATE SET + AIRCRAFT_TYPE_CODE_IATA = EXCLUDED.AIRCRAFT_TYPE_CODE_IATA, + AIRCRAFT_TYPE_CODE_ICAO = EXCLUDED.AIRCRAFT_TYPE_CODE_ICAO, + AIRCRAFT_TYPE_NAME = EXCLUDED.AIRCRAFT_TYPE_NAME, + CHAP = EXCLUDED.CHAP, + MAX_PASSENGERS = EXCLUDED.MAX_PASSENGERS, + MAX_FREIGHT_WEIGHT = EXCLUDED.MAX_FREIGHT_WEIGHT, + MAX_TAKEOFF_WEIGHT = EXCLUDED.MAX_TAKEOFF_WEIGHT, + MAX_AIRBRIDGES = EXCLUDED.MAX_AIRBRIDGES + """.trimIndent(), + ) { ps -> + ps.setString(1, code) + ps.setString(2, code) + ps.setString(3, tag(f, "ICAT")) + ps.setString(4, tag(f, "DESC")) + ps.setString(5, tag(f, "CHAP")) + ps.setObject(6, tag(f, "MAXP")?.toBigDecimalOrNull()) + ps.setObject(7, tag(f, "MFWT")?.toBigDecimalOrNull()) + ps.setObject(8, tag(f, "MTWT")?.toBigDecimalOrNull()) + ps.setObject(9, tag(f, "MABR")?.toBigDecimalOrNull()) + } + } + + private fun upsertAircraft(rec: RefDataRecord) { + val f = rec.fields + val tail = tag(f, "RNUM") ?: return + val id = RefDataFieldMapping.stableNumericId("REGN", tail) + ds.update( + """ + INSERT INTO basicdata.sys_aircraft ( + AIRCRAFT_ID, AIRCRAFT_NUMBER, AIRCRAFT_TYPE_CODE, AIRLINE_CODE, OWID, + MAX_PASSENGERS, MAX_FREIGHT_WEIGHT, MAX_TAKEOFF_WEIGHT + ) VALUES (?, ?, ?, ?, ?, ?, ?, ?) + ON CONFLICT (AIRCRAFT_ID) DO UPDATE SET + AIRCRAFT_NUMBER = EXCLUDED.AIRCRAFT_NUMBER, + AIRCRAFT_TYPE_CODE = EXCLUDED.AIRCRAFT_TYPE_CODE, + AIRLINE_CODE = EXCLUDED.AIRLINE_CODE, + OWID = EXCLUDED.OWID, + MAX_PASSENGERS = EXCLUDED.MAX_PASSENGERS, + MAX_FREIGHT_WEIGHT = EXCLUDED.MAX_FREIGHT_WEIGHT, + MAX_TAKEOFF_WEIGHT = EXCLUDED.MAX_TAKEOFF_WEIGHT + """.trimIndent(), + ) { ps -> + ps.setLong(1, id) + ps.setString(2, tail) + ps.setString(3, tag(f, "ITAT")) + ps.setString(4, tag(f, "ACAL")) + ps.setObject(5, tag(f, "OWID")?.toBigDecimalOrNull()) + ps.setObject(6, tag(f, "MAXP")?.toBigDecimalOrNull()) + ps.setObject(7, tag(f, "MFWT")?.toBigDecimalOrNull()) + ps.setObject(8, tag(f, "MTWT")?.toBigDecimalOrNull()) + } + } + + private fun upsertFlightAgent(rec: RefDataRecord) { + val f = rec.fields + val ogid = tag(f, "OGID") ?: return + ds.update( + """ + INSERT INTO basicdata.sys_flightagent (FLIGHTAGENT_ID, OGID, FLIGHTAGENT_NAME) + VALUES (?, ?, ?) + ON CONFLICT (FLIGHTAGENT_ID) DO UPDATE SET + OGID = EXCLUDED.OGID, + FLIGHTAGENT_NAME = EXCLUDED.FLIGHTAGENT_NAME + """.trimIndent(), + ) { ps -> + ps.setString(1, ogid) + ps.setString(2, ogid) + ps.setString(3, tag(f, "ONAM")) + } + } + + private fun upsertFlightType(rec: RefDataRecord) { + val f = rec.fields + val code = tag(f, "FTYP") ?: return + ds.update( + """ + INSERT INTO basicdata.sys_flight_type ( + FLIGHT_TYPE_CODE, FLIGHT_TYPE_CAA_CODE, FLIGHT_TYPE_NAME, FLIGHT_TYPE_NAME_CN, C_TAG + ) VALUES (?, ?, ?, ?, ?) + ON CONFLICT (FLIGHT_TYPE_CODE) DO UPDATE SET + FLIGHT_TYPE_CAA_CODE = EXCLUDED.FLIGHT_TYPE_CAA_CODE, + FLIGHT_TYPE_NAME = EXCLUDED.FLIGHT_TYPE_NAME, + FLIGHT_TYPE_NAME_CN = EXCLUDED.FLIGHT_TYPE_NAME_CN, + C_TAG = EXCLUDED.C_TAG + """.trimIndent(), + ) { ps -> + ps.setString(1, code) + ps.setString(2, tag(f, "CTYP")) + ps.setString(3, tag(f, "FDES")) + ps.setString(4, tag(f, "FDSC")) + ps.setString(5, tag(f, "FCML")) + } + } + + private fun upsertTerminal(rec: RefDataRecord) { + val f = rec.fields + val code = tag(f, "TCOD") ?: return + ds.update( + """ + INSERT INTO basicdata.orms_terminal (TERMINAL_CODE, TERMINAL_NAME) + VALUES (?, ?) + ON CONFLICT (TERMINAL_CODE) DO UPDATE SET TERMINAL_NAME = EXCLUDED.TERMINAL_NAME + """.trimIndent(), + ) { ps -> + ps.setString(1, code) + ps.setString(2, tag(f, "TNAM")) + } + } + + private fun upsertGate(rec: RefDataRecord) { + val f = rec.fields + val code = tag(f, "GCOD") ?: return + ds.update( + """ + INSERT INTO basicdata.orms_gate (GATE_CODE, GATE_NAME, GATE_CATEGORY, GTML, GATE_STATUS) + VALUES (?, ?, ?, ?, 0) + ON CONFLICT (GATE_CODE) DO UPDATE SET + GATE_NAME = EXCLUDED.GATE_NAME, + GATE_CATEGORY = EXCLUDED.GATE_CATEGORY, + GTML = EXCLUDED.GTML + """.trimIndent(), + ) { ps -> + ps.setString(1, code) + ps.setString(2, tag(f, "GTNM")) + ps.setString(3, tag(f, "GCAT")) + ps.setString(4, tag(f, "GTML")) + } + } + + private fun upsertStand(rec: RefDataRecord) { + val f = rec.fields + val code = tag(f, "SCOD") ?: return + ds.update( + """ + INSERT INTO basicdata.orms_stand (STAND_CODE, STAND_NAME, STML, SATC, STAND_STATUS) + VALUES (?, ?, ?, ?, 0) + ON CONFLICT (STAND_CODE) DO UPDATE SET + STAND_NAME = EXCLUDED.STAND_NAME, + STML = EXCLUDED.STML, + SATC = EXCLUDED.SATC + """.trimIndent(), + ) { ps -> + ps.setString(1, code) + ps.setString(2, tag(f, "STNM")) + ps.setString(3, tag(f, "STML")) + ps.setString(4, tag(f, "SATC")) + } + } + + private fun upsertCheckinDesk(rec: RefDataRecord) { + val f = rec.fields + val code = tag(f, "CCOD") ?: return + ds.update( + """ + INSERT INTO basicdata.orms_checkindesk (CHECKINDESK_CODE, CHECKINDESK_NAME, CCAT, CHECKINDESK_STATUS) + VALUES (?, ?, ?, 'E') + ON CONFLICT (CHECKINDESK_CODE) DO UPDATE SET + CHECKINDESK_NAME = EXCLUDED.CHECKINDESK_NAME, + CCAT = EXCLUDED.CCAT + """.trimIndent(), + ) { ps -> + ps.setString(1, code) + ps.setString(2, tag(f, "CTNM")) + ps.setString(3, tag(f, "CCAT")) + } + } + + private fun upsertCarousel(rec: RefDataRecord) { + val f = rec.fields + val code = tag(f, "BCOD") ?: return + ds.update( + """ + INSERT INTO basicdata.orms_carousel (CAROUSEL_CODE, CAROUSEL_NAME, BCAT, BTML, CAROUSEL_STATUS) + VALUES (?, ?, ?, ?, 'E') + ON CONFLICT (CAROUSEL_CODE) DO UPDATE SET + CAROUSEL_NAME = EXCLUDED.CAROUSEL_NAME, + BCAT = EXCLUDED.BCAT, + BTML = EXCLUDED.BTML + """.trimIndent(), + ) { ps -> + ps.setString(1, code) + ps.setString(2, tag(f, "BTNM")) + ps.setString(3, tag(f, "BCAT")) + ps.setString(4, tag(f, "BTML")) + } + } + + private fun upsertChut(rec: RefDataRecord) { + val f = rec.fields + val code = tag(f, "CCOD") ?: return + ds.update( + """ + INSERT INTO basicdata.orms_chut (CHUT_CODE, CHUT_NAME, BCAT, BTML, CHUT_STATUS) + VALUES (?, ?, ?, ?, 'E') + ON CONFLICT (CHUT_CODE) DO UPDATE SET + CHUT_NAME = EXCLUDED.CHUT_NAME, + BCAT = EXCLUDED.BCAT, + BTML = EXCLUDED.BTML + """.trimIndent(), + ) { ps -> + ps.setString(1, code) + ps.setString(2, tag(f, "CHNM")) + ps.setString(3, tag(f, "CCAT")) + ps.setString(4, tag(f, "CTML")) + } + } + + private fun tag(fields: Map, name: String): String? = + RefDataFieldMapping.blankToNull(fields[name]) + + private fun numericStatus(stat: String): BigDecimal = + if (stat == "D") BigDecimal.ONE else BigDecimal.ZERO +} diff --git a/src/main/kotlin/com/gzzn/omms/msgexchange/infra/stub/StubRefDataRepository.kt b/src/main/kotlin/com/gzzn/omms/msgexchange/infra/stub/StubRefDataRepository.kt new file mode 100644 index 0000000..7b8609d --- /dev/null +++ b/src/main/kotlin/com/gzzn/omms/msgexchange/infra/stub/StubRefDataRepository.kt @@ -0,0 +1,177 @@ +package com.gzzn.omms.msgexchange.infra.stub + +import com.gzzn.omms.msgexchange.domain.ref.RefDataFieldMapping +import com.gzzn.omms.msgexchange.domain.ref.RefDataRecord +import com.gzzn.omms.msgexchange.domain.ref.ResourceStatusRecord +import com.gzzn.omms.msgexchange.infra.persistence.RefDataRepository +import io.micronaut.context.annotation.Requires +import jakarta.inject.Singleton +import java.math.BigDecimal +import java.util.concurrent.ConcurrentHashMap + +@Singleton +@Requires(property = "msgx.stubs", value = "true") +class StubRefDataRepository : RefDataRepository { + val countries = ConcurrentHashMap>() + val airports = ConcurrentHashMap>() + val airlines = ConcurrentHashMap>() + val aircraftTypes = ConcurrentHashMap>() + val aircraft = ConcurrentHashMap>() + val flightAgents = ConcurrentHashMap>() + val flightTypes = ConcurrentHashMap>() + val terminals = ConcurrentHashMap>() + val gates = ConcurrentHashMap>() + val stands = ConcurrentHashMap>() + val checkinDesks = ConcurrentHashMap>() + val carousels = ConcurrentHashMap>() + val chuts = ConcurrentHashMap>() + + override fun replaceCategory(category: String, records: List) { + tableOf(category).clear() + records.forEach { upsertCategory(category, it) } + } + + override fun upsertCategory(category: String, record: RefDataRecord) { + val key = RefDataFieldMapping.businessKey(category, record) + val row = mapRow(category, record).toMutableMap() + tableOf(category)[key] = row + } + + override fun deleteCategory(category: String, businessKey: String) { + tableOf(category).remove(businessKey) + } + + override fun applyResourceStatuses(updates: List) { + updates.forEach { u -> + when (u.rtyp) { + "GATE" -> gates[u.rsid]?.let { it["GATE_STATUS"] = gateStatus(u.stat) } + "STND" -> stands[u.rsid]?.let { it["STAND_STATUS"] = gateStatus(u.stat) } + "BELT" -> carousels[u.rsid]?.let { it["CAROUSEL_STATUS"] = u.stat } + "CNTR" -> checkinDesks[u.rsid]?.let { it["CHECKINDESK_STATUS"] = u.stat } + "CHUT" -> chuts[u.rsid]?.let { it["CHUT_STATUS"] = u.stat } + } + } + } + + private fun tableOf(category: String): ConcurrentHashMap> = when (category) { + "COUL" -> countries + "ARPT" -> airports + "AIRL" -> airlines + "AIRC" -> aircraftTypes + "REGN" -> aircraft + "ORGN" -> flightAgents + "FLTL" -> flightTypes + "TLST" -> terminals + "GLST" -> gates + "SLST" -> stands + "CLST" -> checkinDesks + "BLST" -> carousels + "CHLT" -> chuts + else -> error("unsupported category $category") + } + + private fun mapRow(category: String, rec: RefDataRecord): Map { + val f = rec.fields + fun s(tag: String) = RefDataFieldMapping.blankToNull(f[tag]) + return when (category) { + "COUL" -> mapOf( + "COUNTRY_ID" to RefDataFieldMapping.stableNumericId(category, s("COUC")!!), + "COUNTRY_CODE" to s("COUC"), + "COUNTRY_NAME_ENG" to s("COUN"), + "COUNTRY_NAME_CHN" to s("CNMC"), + ) + "ARPT" -> mapOf( + "AIRPORT_ID" to RefDataFieldMapping.stableNumericId(category, s("ITCD")!!), + "AIRPORT_CODE_IATA" to s("ITCD"), + "AIRPORT_CODE_ICAO" to s("ICCD"), + "AIRPORT_NAME_ENG" to s("ANAM"), + "AIRPORT_NAME_CHN" to s("ANMC"), + "HAUL" to s("HAUL"), + "INTERNATIONAL_FLAG" to RefDataFieldMapping.domIntFlag(s("ATYP")), + ) + "AIRL" -> mapOf( + "AIRLINE_ID" to RefDataFieldMapping.stableNumericId(category, s("ITOP")!!), + "AIRLINE_CODE_IATA" to s("ITOP"), + "AIRLINE_CODE_ICAO" to s("ICOP"), + "AIRLINE_NAME_ENG" to s("ONAM"), + "BRIEF_NAME_CHN" to s("ONMC"), + "INTERNATIONAL_FLAG" to RefDataFieldMapping.domIntFlag(s("DORI")), + ) + "AIRC" -> mapOf( + "AIRCRAFT_TYPE_CODE" to s("ITAT"), + "AIRCRAFT_TYPE_CODE_IATA" to s("ITAT"), + "AIRCRAFT_TYPE_CODE_ICAO" to s("ICAT"), + "AIRCRAFT_TYPE_NAME" to s("DESC"), + "CHAP" to s("CHAP"), + "MAX_PASSENGERS" to s("MAXP")?.toBigDecimalOrNull(), + "MAX_FREIGHT_WEIGHT" to s("MFWT")?.toBigDecimalOrNull(), + "MAX_TAKEOFF_WEIGHT" to s("MTWT")?.toBigDecimalOrNull(), + "MAX_AIRBRIDGES" to s("MABR")?.toBigDecimalOrNull(), + ) + "REGN" -> mapOf( + "AIRCRAFT_ID" to RefDataFieldMapping.stableNumericId(category, s("RNUM")!!), + "AIRCRAFT_NUMBER" to s("RNUM"), + "AIRCRAFT_TYPE_CODE" to s("ITAT"), + "AIRLINE_CODE" to s("ACAL"), + "OWID" to s("OWID")?.toBigDecimalOrNull(), + "MAX_PASSENGERS" to s("MAXP")?.toBigDecimalOrNull(), + "MAX_FREIGHT_WEIGHT" to s("MFWT")?.toBigDecimalOrNull(), + "MAX_TAKEOFF_WEIGHT" to s("MTWT")?.toBigDecimalOrNull(), + ) + "ORGN" -> mapOf( + "FLIGHTAGENT_ID" to s("OGID"), + "OGID" to s("OGID"), + "FLIGHTAGENT_NAME" to s("ONAM"), + ) + "FLTL" -> mapOf( + "FLIGHT_TYPE_CODE" to s("FTYP"), + "FLIGHT_TYPE_CAA_CODE" to s("CTYP"), + "FLIGHT_TYPE_NAME" to s("FDES"), + "FLIGHT_TYPE_NAME_CN" to s("FDSC"), + "C_TAG" to s("FCML"), + ) + "TLST" -> mapOf( + "TERMINAL_CODE" to s("TCOD"), + "TERMINAL_NAME" to s("TNAM"), + ) + "GLST" -> mapOf( + "GATE_CODE" to s("GCOD"), + "GATE_NAME" to s("GTNM"), + "GATE_CATEGORY" to s("GCAT"), + "GTML" to s("GTML"), + "GATE_STATUS" to BigDecimal.ZERO, + ) + "SLST" -> mapOf( + "STAND_CODE" to s("SCOD"), + "STAND_NAME" to s("STNM"), + "STML" to s("STML"), + "SATC" to s("SATC"), + "STAND_STATUS" to BigDecimal.ZERO, + ) + "CLST" -> mapOf( + "CHECKINDESK_CODE" to s("CCOD"), + "CHECKINDESK_NAME" to s("CTNM"), + "CCAT" to s("CCAT"), + "CHECKINDESK_STATUS" to "E", + ) + "BLST" -> mapOf( + "CAROUSEL_CODE" to s("BCOD"), + "CAROUSEL_NAME" to s("BTNM"), + "BCAT" to s("BCAT"), + "BTML" to s("BTML"), + "CAROUSEL_STATUS" to "E", + ) + "CHLT" -> mapOf( + "CHUT_CODE" to s("CCOD"), + "CHUT_NAME" to s("CHNM"), + "BCAT" to s("CCAT"), + "BTML" to s("CTML"), + "CHUT_STATUS" to "E", + ) + else -> error("unsupported category $category") + } + } + + private fun gateStatus(stat: String): BigDecimal = + if (stat == "D") BigDecimal.ONE else BigDecimal.ZERO +} diff --git a/src/main/kotlin/com/gzzn/omms/msgexchange/processing/OutboundRequestService.kt b/src/main/kotlin/com/gzzn/omms/msgexchange/processing/OutboundRequestService.kt index 3a246ec..73e430b 100644 --- a/src/main/kotlin/com/gzzn/omms/msgexchange/processing/OutboundRequestService.kt +++ b/src/main/kotlin/com/gzzn/omms/msgexchange/processing/OutboundRequestService.kt @@ -69,6 +69,9 @@ class OutboundRequestService( fun completeRqfdResponse(operationDay: LocalDate = currentOperationDay()): Boolean = reqTrack.completeLatest(OutboundRequestKeys.RQFD_REQ_TYPE, operationDay, OutboundRequestKeys.SENDER) + fun completeRqrdResponse(operationDay: LocalDate = currentOperationDay()): Boolean = + reqTrack.completeLatest(OutboundRequestKeys.RQRD_REQ_TYPE, operationDay, OutboundRequestKeys.SENDER) + fun dispatchPending(limit: Int = 10): Int { var sent = 0 reqTrack.listPendingDispatch(limit).forEach { req -> diff --git a/src/main/kotlin/com/gzzn/omms/msgexchange/processing/Pump.kt b/src/main/kotlin/com/gzzn/omms/msgexchange/processing/Pump.kt index 24f3114..7a3b157 100644 --- a/src/main/kotlin/com/gzzn/omms/msgexchange/processing/Pump.kt +++ b/src/main/kotlin/com/gzzn/omms/msgexchange/processing/Pump.kt @@ -124,6 +124,7 @@ class MessageProcessor( private val flopProcessor: FlopProcessor, private val fdelProcessor: FdelProcessor, private val adftProcessor: AdftProcessor, + private val referenceDataProcessor: ReferenceDataProcessor, private val outbound: OutboundRequestService, private val procFailure: ProcFailure, private val props: PipelineProps, @@ -237,10 +238,9 @@ class MessageProcessor( flopProcessor.apply(head, decoded, payload) } is MsgKind.RefData -> { - // 参考数据必须走 US-13(implementation.md「分派」),不得按"不支持"跳过; - // ReferenceDataProcessor 尚未实装(ACM2-93,待 Q22),先按可重试失败留队并告警。 - log.warn("refdata without handler -> FAILED(UNSUPPORTED) msgId={} type={}", head.msgId, kind.type) - return procFailure.fail(head, ErrorClass.UNSUPPORTED, "refdata-pending:${kind.type}") + val body = decoded.body as? com.gzzn.omms.msgexchange.domain.ref.RefDataBody // validated below + if (body == null) return deadMalformed(head, "missing-refdata-body") + referenceDataProcessor.apply(head, decoded, kind.type) } MsgKind.Eror -> { val payload = decoded.body as? ErorPayload diff --git a/src/main/kotlin/com/gzzn/omms/msgexchange/processing/RefDataCommit.kt b/src/main/kotlin/com/gzzn/omms/msgexchange/processing/RefDataCommit.kt new file mode 100644 index 0000000..1b4e603 --- /dev/null +++ b/src/main/kotlin/com/gzzn/omms/msgexchange/processing/RefDataCommit.kt @@ -0,0 +1,28 @@ +package com.gzzn.omms.msgexchange.processing + +import com.gzzn.omms.msgexchange.domain.ProcState +import com.gzzn.omms.msgexchange.domain.ProcStatus +import com.gzzn.omms.msgexchange.infra.persistence.PipelineLockRepository +import com.gzzn.omms.msgexchange.infra.persistence.PipelineTransactionManager +import com.gzzn.omms.msgexchange.infra.persistence.ProcStateRepository +import jakarta.inject.Singleton +import java.time.Clock + +/** 静态参考数据单事务提交:落库 + 终态 + 回填意图(无 Redis / MSG_EVENT,`US-13`)。 */ +@Singleton +class RefDataCommit( + private val txManager: PipelineTransactionManager, + private val lock: PipelineLockRepository, + private val procState: ProcStateRepository, + private val clock: Clock, +) { + fun commit(head: ProcState, work: () -> Unit) { + lock.holdAcrossTransactions { + txManager.inTransaction { + lock.lock() + work() + procState.markTerminal(head.msgId, ProcStatus.SUCCEEDED, now = clock.instant()) + } + } + } +} diff --git a/src/main/kotlin/com/gzzn/omms/msgexchange/processing/ReferenceDataProcessor.kt b/src/main/kotlin/com/gzzn/omms/msgexchange/processing/ReferenceDataProcessor.kt new file mode 100644 index 0000000..24c0826 --- /dev/null +++ b/src/main/kotlin/com/gzzn/omms/msgexchange/processing/ReferenceDataProcessor.kt @@ -0,0 +1,52 @@ +package com.gzzn.omms.msgexchange.processing + +import com.gzzn.omms.msgexchange.domain.DecodedMessage +import com.gzzn.omms.msgexchange.domain.ProcState +import com.gzzn.omms.msgexchange.domain.ref.RefDataBody +import com.gzzn.omms.msgexchange.domain.ref.RefDataFieldMapping +import com.gzzn.omms.msgexchange.domain.ref.RefDataMode +import com.gzzn.omms.msgexchange.domain.ref.RefDataValidator +import com.gzzn.omms.msgexchange.domain.ref.refDataModeOf +import com.gzzn.omms.msgexchange.infra.persistence.RefDataRepository +import jakarta.inject.Singleton +import org.slf4j.LoggerFactory + +@Singleton +class ReferenceDataProcessor( + private val commit: RefDataCommit, + private val refData: RefDataRepository, + private val outbound: OutboundRequestService, +) { + private val log = LoggerFactory.getLogger(ReferenceDataProcessor::class.java) + + fun apply(head: ProcState, msg: DecodedMessage, category: String): ApplyResult { + val body = msg.body as? RefDataBody ?: return ApplyResult.DeadProtocol("missing-refdata-body") + val mode = refDataModeOf(body.styp) ?: return ApplyResult.DeadProtocol("refdata-styp:${body.styp}") + val reason = RefDataValidator.validate(category, body, mode) + if (reason != null) { + log.warn("refdata validation failed msgId={} category={} reason={}", head.msgId, category, reason) + return ApplyResult.DeadProtocol(reason) + } + + commit.commit(head) { + if (category == "RSTA") { + refData.applyResourceStatuses(body.resourceStatuses) + } else { + when (mode) { + RefDataMode.FULL_REPLACE -> refData.replaceCategory(category, body.records) + RefDataMode.ADD, RefDataMode.UPDATE -> + body.records.forEach { refData.upsertCategory(category, it) } + RefDataMode.DELETE -> + body.records.forEach { + refData.deleteCategory(category, RefDataFieldMapping.businessKey(category, it)) + } + } + } + } + + if (body.styp.equals("RESP", ignoreCase = true)) { + outbound.completeRqrdResponse() + } + return ApplyResult.Succeeded + } +} diff --git a/src/main/resources/db/migration/V6__basicdata_ref_data.sql b/src/main/resources/db/migration/V6__basicdata_ref_data.sql new file mode 100644 index 0000000..e630004 --- /dev/null +++ b/src/main/resources/db/migration/V6__basicdata_ref_data.sql @@ -0,0 +1,178 @@ +-- 静态参考数据物理表组(`Q22`):schema `basicdata`,表列对齐 admin-api secondary 实体注解, +-- 与 deploy/integration/postgres-init/02-basicdata-schema.sql 同源;本系统 Flyway 管辖。 + +CREATE SCHEMA IF NOT EXISTS basicdata; + +COMMENT ON SCHEMA basicdata IS + '静态参考数据与 admin-api 只读基础数据(表列取自 admin-api 实体注解)。'; + +SET search_path = basicdata; + +CREATE TABLE sys_country ( + COUNTRY_ID BIGINT NOT NULL PRIMARY KEY, + COUNTRY_CODE VARCHAR(16), + COUNTRY_NAME_CHN VARCHAR(64), + COUNTRY_NAME_ENG VARCHAR(64), + OPERATE SMALLINT +); + +CREATE TABLE sys_city ( + CITY_ID BIGINT NOT NULL PRIMARY KEY, + CITY_CODE VARCHAR(16), + CITY_ICCDE VARCHAR(16), + CITY_NAME_CHN VARCHAR(64), + CITY_NAME_ENG VARCHAR(64), + COUNTRY_ID BIGINT, + MAIN_CITY_TAG NUMERIC(18,0), + OPERATE NUMERIC(18,0), + TEL_CITY_CODE VARCHAR(16) +); + +CREATE TABLE sys_airport ( + AIRPORT_ID BIGINT NOT NULL PRIMARY KEY, + AIRPORT_CODE VARCHAR(16), + AIRPORT_CODE_CAA VARCHAR(16), + AIRPORT_CODE_IATA VARCHAR(16), + AIRPORT_CODE_ICAO VARCHAR(16), + AIRPORT_GROUP_ID BIGINT, + AIRPORT_NAME_CHN VARCHAR(64), + AIRPORT_NAME_ENG VARCHAR(64), + BRIEF_NAME_CHN VARCHAR(64), + BRIEF_NAME_CHN_COPY VARCHAR(64), + BRIEF_NAME_ENG VARCHAR(64), + CITY_ID BIGINT, + HAUL VARCHAR(8), + HOST_FLAG NUMERIC(18,0), + INTERNATIONAL_FLAG NUMERIC(18,0), + OPERATE NUMERIC(18,0) +); + +CREATE UNIQUE INDEX uq_sys_airport_iata ON sys_airport (AIRPORT_CODE_IATA) WHERE AIRPORT_CODE_IATA IS NOT NULL; + +CREATE TABLE sys_airline ( + AIRLINE_ID BIGINT NOT NULL PRIMARY KEY, + AIRLINE_CODE VARCHAR(16), + AIRLINE_CODE_CAA VARCHAR(16), + AIRLINE_CODE_IATA VARCHAR(16), + AIRLINE_CODE_ICAO VARCHAR(16), + AIRLINE_GROUP_ID BIGINT, + AIRLINE_NAME VARCHAR(64), + AIRLINE_NAME_ENG VARCHAR(64), + BRIEF_NAME_CHN VARCHAR(64), + BRIEF_NAME_ENG VARCHAR(64), + CODE VARCHAR(16), + COUNTRY_ID BIGINT, + HOST_FLAG NUMERIC(18,0), + INTERNATIONAL_FLAG NUMERIC(18,0), + NAMC VARCHAR(64), + NAME VARCHAR(64), + OPERATE NUMERIC(18,0), + PARENT_AIRLINE_ID BIGINT +); + +CREATE UNIQUE INDEX uq_sys_airline_iata ON sys_airline (AIRLINE_CODE_IATA) WHERE AIRLINE_CODE_IATA IS NOT NULL; + +CREATE TABLE sys_aircrafttype ( + AIRCRAFT_TYPE_CODE VARCHAR(16) NOT NULL PRIMARY KEY, + AIRCRAFT_HEIGHT NUMERIC(18,4), + AIRCRAFT_LENGTH NUMERIC(18,4), + AIRCRAFT_TYPE_CODE_CAA VARCHAR(16), + AIRCRAFT_TYPE_CODE_IATA VARCHAR(16), + AIRCRAFT_TYPE_CODE_ICAO VARCHAR(16), + AIRCRAFT_TYPE_NAME VARCHAR(64), + AIRCRAFT_WIDTH NUMERIC(18,4), + AIRCRAFTTYPEGROUP_ID BIGINT, + CHAP VARCHAR(16), + MAX_AIRBRIDGES NUMERIC(18,0), + MAX_FREIGHT_WEIGHT NUMERIC(18,4), + MAX_PASSENGERS NUMERIC(18,0), + MAX_TAKEOFF_WEIGHT NUMERIC(18,4), + OPERATE NUMERIC(18,0), + TRANSITSTOPTIME NUMERIC(18,0) +); + +CREATE TABLE sys_aircraft ( + AIRCRAFT_ID BIGINT NOT NULL PRIMARY KEY, + AIRCRAFT_NUMBER VARCHAR(32), + AIRCRAFT_OWNER_CODE VARCHAR(16), + AIRCRAFT_TYPE_CODE VARCHAR(16), + AIRLINE_CODE VARCHAR(16), + MAX_AIRBRIDGES NUMERIC(18,0), + MAX_FREIGHT_WEIGHT NUMERIC(18,4), + MAX_PASSENGERS NUMERIC(18,0), + MAX_TAKEOFF_WEIGHT NUMERIC(18,4), + OPERATE NUMERIC(18,0), + OWID NUMERIC(18,0) +); + +CREATE UNIQUE INDEX uq_sys_aircraft_number ON sys_aircraft (AIRCRAFT_NUMBER) WHERE AIRCRAFT_NUMBER IS NOT NULL; + +CREATE TABLE sys_flight_type ( + FLIGHT_TYPE_CODE VARCHAR(8) NOT NULL PRIMARY KEY, + FLIGHT_TYPE_CAA_CODE VARCHAR(8), + FLIGHT_TYPE_NAME VARCHAR(64), + C_TAG VARCHAR(8), + VIP_TAG VARCHAR(8), + FLIGHT_TYPE_NAME_CN VARCHAR(64), + OPERATE SMALLINT +); + +CREATE TABLE sys_flightagent ( + FLIGHTAGENT_ID VARCHAR(32) NOT NULL PRIMARY KEY, + FLIGHTAGENT_NAME VARCHAR(64), + OGID VARCHAR(32) +); + +CREATE TABLE orms_terminal ( + TERMINAL_CODE VARCHAR(16) NOT NULL PRIMARY KEY, + OPERATE NUMERIC(18,0), + TERMINAL_NAME VARCHAR(64) +); + +CREATE TABLE orms_stand ( + STAND_CODE VARCHAR(16) NOT NULL PRIMARY KEY, + STAND_NAME VARCHAR(64), + STAND_TYPE VARCHAR(16), + STAND_STATUS NUMERIC(18,0), + STML VARCHAR(16), + SATC VARCHAR(16), + OPERATE NUMERIC(18,0) +); + +CREATE TABLE orms_gate ( + GATE_CODE VARCHAR(16) NOT NULL PRIMARY KEY, + GATE_NAME VARCHAR(64), + GATE_CATEGORY VARCHAR(8), + GATE_STATUS NUMERIC(18,0), + GTML VARCHAR(16), + OPERATE NUMERIC(18,0) +); + +CREATE TABLE orms_checkindesk ( + CHECKINDESK_CODE VARCHAR(16) NOT NULL PRIMARY KEY, + CHECKINDESK_NAME VARCHAR(64), + CHECKINDESK_STATUS VARCHAR(8), + CCAT VARCHAR(16), + CTRA VARCHAR(8), + TERMINAL_AREA_CODE VARCHAR(16), + OPERATE NUMERIC(18,0) +); + +CREATE TABLE orms_carousel ( + CAROUSEL_CODE VARCHAR(16) NOT NULL PRIMARY KEY, + CAROUSEL_NAME VARCHAR(64), + CAROUSEL_STATUS VARCHAR(32), + BCAT VARCHAR(16), + BTML VARCHAR(16) +); + +CREATE TABLE orms_chut ( + CHUT_CODE VARCHAR(16) NOT NULL PRIMARY KEY, + CHUT_NAME VARCHAR(64), + CHUT_STATUS VARCHAR(8), + BCAT VARCHAR(16), + BTML VARCHAR(16), + OPERATE NUMERIC(18,0) +); + +RESET search_path; diff --git a/src/test/kotlin/com/gzzn/omms/msgexchange/PipelineSmokeTest.kt b/src/test/kotlin/com/gzzn/omms/msgexchange/PipelineSmokeTest.kt index c78e964..0a981d9 100644 --- a/src/test/kotlin/com/gzzn/omms/msgexchange/PipelineSmokeTest.kt +++ b/src/test/kotlin/com/gzzn/omms/msgexchange/PipelineSmokeTest.kt @@ -84,7 +84,6 @@ class PipelineSmokeTest { """.trimIndent() - /** 参考数据类别(SIS:3.2):ReferenceDataProcessor 未实装(ACM2-93),先按可重试失败留队 */ val REFDATA_XML = """ AODB220260908120000ARPTDNLD @@ -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),由回填扫描把标记写回信箱。 * 少了这一步,这些行永远占着每批的名额,攒够一批就再也发现不了新消息了。 diff --git a/src/test/kotlin/com/gzzn/omms/msgexchange/processing/FlightCommitTest.kt b/src/test/kotlin/com/gzzn/omms/msgexchange/processing/FlightCommitTest.kt index c0c2734..3af989a 100644 --- a/src/test/kotlin/com/gzzn/omms/msgexchange/processing/FlightCommitTest.kt +++ b/src/test/kotlin/com/gzzn/omms/msgexchange/processing/FlightCommitTest.kt @@ -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(), diff --git a/src/test/kotlin/com/gzzn/omms/msgexchange/processing/IgnoreBranchTest.kt b/src/test/kotlin/com/gzzn/omms/msgexchange/processing/IgnoreBranchTest.kt index 0db742d..3a5c9c8 100644 --- a/src/test/kotlin/com/gzzn/omms/msgexchange/processing/IgnoreBranchTest.kt +++ b/src/test/kotlin/com/gzzn/omms/msgexchange/processing/IgnoreBranchTest.kt @@ -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 inTransaction(block: () -> T) = block() }, + object : PipelineLockRepository { + override fun lock() = Unit + override fun 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 = "", + 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 = "", + 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] = "" - val msg = DecodedMessage( - meta = MetaFields("AODB", type, "X", 300L, 1L), - kind = MsgKind.RefData(type), - rawXml = "", - ) 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] = "" val msg = DecodedMessage( - meta = MetaFields("AODB", "ARPT", "DNLD", 400L, 1L), + meta = MetaFields("AODB", "ARPT", "ADD", 400L, 1L), kind = MsgKind.RefData("ARPT"), rawXml = "", + 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 diff --git a/src/test/kotlin/com/gzzn/omms/msgexchange/processing/ReferenceDataProcessorTest.kt b/src/test/kotlin/com/gzzn/omms/msgexchange/processing/ReferenceDataProcessorTest.kt new file mode 100644 index 0000000..b309685 --- /dev/null +++ b/src/test/kotlin/com/gzzn/omms/msgexchange/processing/ReferenceDataProcessorTest.kt @@ -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 = """ + + AODB120260908120000COULDNLD + CNChina中国11 + + """.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 = """ + + AODB220260908120000ARPTUPD + CTUChengdu成都CNCTU + + """.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 = """ + + AODB320260908120000RSTADNLD + GATEG1D20260908120000 + + """.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 = """ + + AODB420260908120000GLSTRESP + G9Gate 99号门DT2 + + """.trimIndent() + processor(req).apply(head(4), msg("GLST", "RESP", xml), "GLST") + assertEquals(ReqTrackRepository.ReqState.DONE, req.rows[reqId]!!.state) + } +} diff --git a/src/test/kotlin/com/gzzn/omms/msgexchange/processing/RespGuardTest.kt b/src/test/kotlin/com/gzzn/omms/msgexchange/processing/RespGuardTest.kt index a5c6f0c..f32690b 100644 --- a/src/test/kotlin/com/gzzn/omms/msgexchange/processing/RespGuardTest.kt +++ b/src/test/kotlin/com/gzzn/omms/msgexchange/processing/RespGuardTest.kt @@ -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, diff --git a/src/test/kotlin/com/gzzn/omms/msgexchange/processing/TestFlightCommit.kt b/src/test/kotlin/com/gzzn/omms/msgexchange/processing/TestFlightCommit.kt index 7f641f4..4fdb419 100644 --- a/src/test/kotlin/com/gzzn/omms/msgexchange/processing/TestFlightCommit.kt +++ b/src/test/kotlin/com/gzzn/omms/msgexchange/processing/TestFlightCommit.kt @@ -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) +}