diff --git a/README.md b/README.md index e53b81c..b29b862 100644 --- a/README.md +++ b/README.md @@ -73,7 +73,7 @@ 仓库根目录提供兼容 Podman Compose 与 Docker Compose 的开发中间件栈 `compose.yaml`,包含: - **共享信箱 MySQL**(`mysql:8.4` LTS,端口 3306,库 `cdairport`):容器启动时自动执行 `deploy/dev/mysql-init/01-mailbox.sql` 创建本地联调所需的 `CMINMSGS`、`CMINMSGS_HST`、`COUTMSGS` 模拟表。 - **自有 PostgreSQL**(`postgres:17-alpine`,端口 5432,库 `msgx`):容器提供干净数据库,应用启动时由 Flyway(`src/main/resources/db/migration/V1__flight_state_baseline.sql`)自动建自有表。 -- **Valkey**(`valkey/valkey:8-alpine`,端口 6379):本地兼容服务;应用仍不连接它(航班状态权威在自有 PG)。处理侧的投影写入时序已就位,但 `msgx.redis.enabled` 默认关闭、写入是空操作,真实客户端与 `GET /all/flights` 仍未接通(`G-REDIS-PROJECTION`)。 +- **Valkey**(`valkey/valkey:8-alpine`,端口 6379):本地兼容服务;航班状态权威仍在自有 PG。设置 `msgx.redis.enabled=true` 后,应用通过 `PARAM:redis.uri` 读写航班查询投影;默认关闭时写入是空操作。 - **Kafka**(`apache/kafka:3.8.0` KRaft 单节点,端口 9092):listener `PLAINTEXT://localhost:9092`,`default.replication.factor=1`,已预配幂等生产者与 acks=all 所需的单节点参数。 ### 快速启动 diff --git a/build.gradle.kts b/build.gradle.kts index 861c159..c28b32e 100644 --- a/build.gradle.kts +++ b/build.gradle.kts @@ -48,6 +48,7 @@ dependencies { implementation(libs.micronaut.discovery.eureka) // 指标:把管道积压/回填/水位等暴露为 MeterRegistry 指标(版本由 platform BOM 管) implementation(libs.micronaut.micrometer.core) + implementation(libs.micronaut.redis.lettuce) implementation(libs.jackson.dataformat.xml) implementation(libs.jackson.module.kotlin) implementation(libs.logstash.logback.encoder) diff --git a/docs/implementation.md b/docs/implementation.md index 5a2454e..ccc5f55 100644 --- a/docs/implementation.md +++ b/docs/implementation.md @@ -198,7 +198,7 @@ LIMIT PARAM:msgx.pipeline.backfill-batch 3. 分批写:跨运营日整包失败;**覆盖范围内**快照缺席航班标删、发删除事件、删 Redis,范围外的不受影响(`INV-7`)。每批同事务写变更 + 事件(`INV-3`)。 4. 整包成功 → `SUCCEEDED` + 回填意图;留痕在事务外。 -字段语义见「航班域」。`RESP` 须匹配开放请求,否则不更新(`G-RESP-GUARD`)。 +字段语义见「航班域」。`RESP` 须匹配开放请求,否则不更新。 ### 7.2 上游请求与静态数据 @@ -283,7 +283,7 @@ PENDING → SENT → DONE | 航班历史存储 | 保留期 | — | — | 我们 | `G-FLIGHT-HIST-RETENTION` | | `SCHD_SNAP_LOG` | 保留期 | 无(本地可重建) | 无 | 我们 | — | | `MSG_EVENT` 已发送行 | `SENT` | 无 | 无 | 我们 | — | -| `PROC_STATE` 终态行 | 见下 | 无(到期直接删除) | 回填了结 | 我们 | `G-PROC-CLEANUP` | +| `PROC_STATE` 终态行 | 见下 | 无(到期直接删除) | 回填了结 | 我们 | — | | `REQ_TRACK` 关闭态行 | 保留期 | 无 | 无 | 我们 | `G-REQ-TRACK-RETENTION` | **`PROC_STATE` 清理**(`US-11`):终态 + `BACKFILL_AT` 非空 + 超保留期(从 `UPDATED_AT` 起)。删除按 `STATE` 条件执行,0 行则跳过(与人工重放互斥)。批量删除不持 `PIPELINE_LOCK`。 @@ -357,7 +357,7 @@ PG 是航班数据源(`INV-5`);Redis 从 PG 同步,处理完成前写入 | `ORDINAL` | 从 1 起;单元素多属性(如 `FDIV/@DDES` 与 `FDIV/@DDIR`)各自 `ORDINAL=1`,因 `PATH` 不同而不冲突 | | `RAW_VALUE` | 原值 UTF-8 字符串;显式空标签记空串,不省略行 | -写入与航班变更同一 PG 事务:对该 `MSG_ID` 先删后插,重处理结果一致即幂等。清理保留期不含本表(`G-PROC-CLEANUP` 只清 `PROC_STATE`)。 +写入与航班变更同一 PG 事务:对该 `MSG_ID` 先删后插,重处理结果一致即幂等。清理保留期不含本表;`US-11` 只清 `PROC_STATE`。 - `ORDINAL` 是持久化顺序,从 1 开始;`SOURCE_SEQ` 是上游序号,允许为空或重复。 - 相同资源号不代表同一条分配,禁止按资源号去重。 diff --git a/docs/reference.md b/docs/reference.md index e9ef679..4c36509 100644 --- a/docs/reference.md +++ b/docs/reference.md @@ -55,8 +55,12 @@ | `msgx.identity.include-day-boundary` | `false` | 去重时是否把日期算进消息身份;受 `C-3` 约束,不能随意改变 | 默认关闭 | | `msgx.health.backlog-cache-ttl-ms` | `30000` ms | `/health` 与 `/metrics` 共用的未处理消息统计最多缓存多久;`0` 表示每次重算 | 暂定 | | `msgx.history.history-store-enabled` | `false` | 是否启用历史航班写入;关闭时不删除实时航班 | 默认关闭 | -| `msgx.redis.enabled` | `false` | 是否把航班快照写进 Redis 查询投影;关闭时投影写入是空操作(`G-REDIS-PROJECTION`) | 默认关闭 | +| `msgx.redis.enabled` | `false` | 是否把航班快照写进 Redis 查询投影;关闭时投影写入是空操作 | 默认关闭 | | `msgx.redis.flight-key` | `flightInfo` | 航班快照所在的 Redis 哈希键,字段名是 `FLID` | 沿用旧系统 | +| `redis.uri` | `redis://127.0.0.1:6379` | Redis 连接地址;生产通过 `MSGX_REDIS_URI` 覆盖 | 当前实现 | +| `msgx.outbound.response-timeout` | `30m` | 出站请求落信后等待应答的最长时间 | 暂定 | +| `msgx.outbound.routing-rqfd` | `OSH5RQFD` | `RQFD` 写入 `COUTMSGS` 使用的路由标识 | 接口契约 | +| `msgx.outbound.routing-rqrd` | `OMMSRQRD` | `RQRD` 写入 `COUTMSGS` 使用的路由标识 | 当前实现 | | `msgx.history.planned-age-days` | `3` 天 | `US-14` AC2 条件 1:计划时间早于当前超过该天数视为已结束 | `US-14` AC2 | | `msgx.history.cancelled-hours` | `1` 小时 | `US-14` AC2 条件 2:取消时间早于当前超过该窗口视为已结束 | `US-14` AC2 | | `msgx.history.diverted-hours` | `1` 小时 | `US-14` AC2 条件 3:备降且计划时间早于当前超过该窗口;备降依据 `FDIV`(`DDES`/`DDIR`)落地前不参与判定(`Q3`、ACM2-100) | `US-14` AC2 | diff --git a/docs/specification.md b/docs/specification.md index bc4e255..b3ea6dd 100644 --- a/docs/specification.md +++ b/docs/specification.md @@ -152,14 +152,10 @@ | `G-FLOP-SEMANTICS` | 动态报文处理与 `US-05` 要求不一致 | `US-05` | | `G-FLOP-UNMAPPED` | [XSD](legacy/unisysaodbsis.xsd) FLOP 字段映射不全 | `US-05` | | `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`;本网关落库与 `ReferenceDataProcessor` 已做(ACM2-93) | -| `G-REQ-OPEN-UNIQUE` | 同类型未处理完时不允许再发,未做 | `US-09` | -| `G-REQ-TRACK` | 出站请求跟踪未做 | `US-09` | +| `G-REQ-OPEN-UNIQUE` | 新请求等待在途请求结案的登记模型未闭合;当前开放态唯一约束会拒绝新登记 | `US-09` | +| `G-REQ-TRACK` | `RQFD` 跟踪已落地;14 类 `RQRD` 人工登记、子类型作废与等待投递未做 | `US-09` | | `G-REQ-TRACK-RETENTION` | `REQ_TRACK` 已结案保留期未定(`Q24`) | `US-09` | -| `G-RESP-GUARD` | `SCHD-RESP` 过期判断未做 | `US-07` | -| `G-SCHD-SNAPSHOT` | 日计划快照删除、清空、分批未做 | `INV-7`、`INV-9` | | `G-SRVT-VIPF` | `SRVT`、`VIPF` 缺席是否清除待 `Q2`;段出现时已落 `FLIGHT_SRVT`/`FLIGHT_VIPF` | `US-05` | ## 8. 验证映射 @@ -181,13 +177,13 @@ | INV-6 | implementation.md「数据模型」 | `FLID` 唯一 | | INV-7 | `US-07` AC2/AC3 | 覆盖范围内快照缺席航班标删并从 Redis 删,范围外不动;未带字段清空 | | INV-8 | `US-06` AC1 | 标删后从 Redis 删;失败下轮重做 | -| INV-9 | `US-07` AC4/AC5 | 分批失败全部重来(`G-SCHD-SNAPSHOT` 做完前不可验);PG 整份写完后再刷 Redis | +| INV-9 | `US-07` AC4/AC5 | 分批失败全部重来;PG 整份写完后再刷 Redis | | INV-11 | `US-12` AC1/AC2 | 返回全部非共享航班且与 Redis 一致;Redis 故障返回错误 | | CLM-1(重复不重复生效) | `US-03` AC2/AC3/AC5 | 跳过/无法处理留档且无业务改动;失败可查;PG 已提交不因写回或 Kafka 失败撤销;待 `Q3`、`G-FLOP-IDEMPOTENT`、`G-FLOP-SEMANTICS` | | CLM-2(升序) | `US-01` AC3、`US-03` AC1 | 按编号从小到大;最前面未完成时后面的不处理 | | CLM-7(启动拒绝) | `OPS-1` | 配置错误时启动失败 | | US-04(ADFT 更新) | `US-04` AC2 | ADFT 未带字段保持原值 | | US-14、D1 | `US-14` AC3/AC4 | 历史写入成功才删实时数据;删事件先于实时删除登记 | -| G-PROC-CLEANUP(未完成不删) | `US-11` AC1/AC2 | 未处理完的记录不清理 | +| US-11(未完成不删) | `US-11` AC1/AC2 | 未处理完的记录不清理 | | OPS-4(切换与回退) | `OPS-4` | 写回处理时间后再切回;在途消息见 `Q18` | | PRE-1(测试隔离) | `OPS-3` | 测试环境独立 DB、Redis、Kafka,不连生产信箱 | diff --git a/gradle/libs.versions.toml b/gradle/libs.versions.toml index c55d045..60c944e 100644 --- a/gradle/libs.versions.toml +++ b/gradle/libs.versions.toml @@ -34,6 +34,7 @@ micronaut-discovery-eureka = { module = "io.micronaut.discovery:micronaut-discov micronaut-management = { module = "io.micronaut:micronaut-management" } # 指标接入(Micrometer):版本由 micronaut-platform BOM 统一管理,勿在此钉版本。 micronaut-micrometer-core = { module = "io.micronaut.micrometer:micronaut-micrometer-core" } +micronaut-redis-lettuce = { module = "io.micronaut.redis:micronaut-redis-lettuce" } # Kotlin 注解处理(KSP)处理器——与 kotlin-ksp 插件配套(U01) micronaut-inject-kotlin = { module = "io.micronaut:micronaut-inject-kotlin", version.ref = "micronaut" } jackson-dataformat-xml = { module = "com.fasterxml.jackson.dataformat:jackson-dataformat-xml", version = "2.18.2" } diff --git a/src/main/kotlin/com/gzzn/omms/msgexchange/config/PipelineProps.kt b/src/main/kotlin/com/gzzn/omms/msgexchange/config/PipelineProps.kt index ab7a5a8..cf4591c 100644 --- a/src/main/kotlin/com/gzzn/omms/msgexchange/config/PipelineProps.kt +++ b/src/main/kotlin/com/gzzn/omms/msgexchange/config/PipelineProps.kt @@ -121,7 +121,7 @@ class PipelineProps { /** 出站请求落信后等待应答的最长时限(`US-09`;具体取值待 Q)。 */ var responseTimeout: Duration = Duration.ofMinutes(30) - var routingRqfd: String = "OMMSRQFD" + var routingRqfd: String = "OSH5RQFD" var routingRqrd: String = "OMMSRQRD" } diff --git a/src/main/kotlin/com/gzzn/omms/msgexchange/config/RedisProps.kt b/src/main/kotlin/com/gzzn/omms/msgexchange/config/RedisProps.kt index c808a79..a06b685 100644 --- a/src/main/kotlin/com/gzzn/omms/msgexchange/config/RedisProps.kt +++ b/src/main/kotlin/com/gzzn/omms/msgexchange/config/RedisProps.kt @@ -5,7 +5,7 @@ import io.micronaut.context.annotation.ConfigurationProperties /** 航班查询投影所在的 Redis,对应 `msgx.redis.*`;键位沿用旧系统(`C-11`、`Q6`)。 */ @ConfigurationProperties("msgx.redis") class RedisProps { - /** 有没有 Redis 可用。关着的时候投影写是空操作(`G-REDIS-PROJECTION`)。 */ + /** 有没有 Redis 可用。关着的时候投影写是空操作。 */ var enabled: Boolean = false /** 航班快照所在的哈希键,字段名是 `FLID`。 */ diff --git a/src/main/kotlin/com/gzzn/omms/msgexchange/infra/persistence/jdbc/JdbcPgRepositories.kt b/src/main/kotlin/com/gzzn/omms/msgexchange/infra/persistence/jdbc/JdbcPgRepositories.kt index 54989ec..9aef4e6 100644 --- a/src/main/kotlin/com/gzzn/omms/msgexchange/infra/persistence/jdbc/JdbcPgRepositories.kt +++ b/src/main/kotlin/com/gzzn/omms/msgexchange/infra/persistence/jdbc/JdbcPgRepositories.kt @@ -965,11 +965,18 @@ class JdbcFlightStateRepository( } else { "SELECT * FROM ${spec.table} WHERE flid = ? ORDER BY ordinal ASC" } - return ds.query(sql, { ps -> ps.setString(1, flid); spec.routeKind?.let { ps.setString(2, it) } }, ::detailItem) + return ds.query( + sql, + { ps -> ps.setString(1, flid); spec.routeKind?.let { ps.setString(2, it) } }, + { rs -> detailItem(rs, spec) }, + ) } - private fun detailItem(rs: ResultSet): Map { + private fun detailItem(rs: ResultSet, spec: DetailSpec): Map { val item = linkedMapOf() + spec.seqAttr?.let { seqAttr -> + rs.getString("source_seq")?.takeIf { it.isNotEmpty() }?.let { item[seqAttr] = it } + } val md = rs.metaData for (i in 1..md.columnCount) { val col = md.getColumnName(i).uppercase() @@ -997,7 +1004,7 @@ class JdbcFlightStateRepository( /** 明细读取时排除的追踪列(业务键以报文线格式大写键回传)。 */ private val IGNORED_DETAIL_COLUMNS = - setOf("CREATED_AT", "UPDATED_AT", "FLID", "ORDINAL", "ROUTE_KIND") + setOf("CREATED_AT", "UPDATED_AT", "FLID", "ORDINAL", "ROUTE_KIND", "SOURCE_SEQ") private const val SELECT_MAIN = "SELECT flid, operation_day, state, state_version, last_msg_id, updated_at FROM flight_schd" diff --git a/src/main/kotlin/com/gzzn/omms/msgexchange/infra/projection/FlightProjectionAdapters.kt b/src/main/kotlin/com/gzzn/omms/msgexchange/infra/projection/FlightProjectionAdapters.kt index dc82a40..0a61ddc 100644 --- a/src/main/kotlin/com/gzzn/omms/msgexchange/infra/projection/FlightProjectionAdapters.kt +++ b/src/main/kotlin/com/gzzn/omms/msgexchange/infra/projection/FlightProjectionAdapters.kt @@ -4,42 +4,52 @@ import com.gzzn.omms.msgexchange.config.RedisProps import com.gzzn.omms.msgexchange.processing.FlightProjectionPort import com.gzzn.omms.msgexchange.processing.FlightProjectionWrite import io.micronaut.context.annotation.Requires +import io.lettuce.core.api.StatefulRedisConnection import jakarta.inject.Singleton /** - * Redis 投影适配器骨架:键位与写法在这里定死,**但客户端还没接进来**(`G-REDIS-PROJECTION`)。 + * Redis 投影适配器:一个哈希键保存全部航班,批次用 Redis 事务原子提交。 * * 写入形态沿用旧系统:一个哈希键(`PARAM:msgx.redis.flight-key`)装全部航班, * field 是 `FLID`,value 是整态 JSON,不设过期;删除航班就删掉这个 field(`INV-8`)。 * - * 开了 `msgx.redis.enabled` 却没有客户端时,写入直接抛异常而不是假装成功: - * 按 `INV-10`,这条消息就停在未完成、下轮重试,不会被当成已处理写回信箱。 + * Redis 调用失败会直接抛异常;按 `INV-10`,消息保持未完成并在下轮重试。 */ @Requires(property = "msgx.stubs", notEquals = "true") @Requires(property = "msgx.redis.enabled", value = "true") @Singleton class RedisFlightProjectionPort( private val props: RedisProps, + private val connection: StatefulRedisConnection, ) : FlightProjectionPort { override fun write(writes: List) { - throw UnsupportedOperationException( - "redis client not wired yet (G-REDIS-PROJECTION); pending writes=${writes.size} key=${props.flightKey}", - ) + if (writes.isEmpty()) return + val commands = connection.sync() + commands.multi() + try { + writes.forEach { write -> + when (write) { + is FlightProjectionWrite.Upsert -> commands.hset(props.flightKey, write.flid, write.payloadJson) + is FlightProjectionWrite.Delete -> commands.hdel(props.flightKey, write.flid) + } + } + commands.exec() + } catch (failure: RuntimeException) { + runCatching { commands.discard() } + throw failure + } } - override fun readAll(): List = throw UnsupportedOperationException( - "redis client not wired yet (G-REDIS-PROJECTION); key=${props.flightKey}", - ) + override fun readAll(): List = connection.sync().hvals(props.flightKey) - override fun ping(): Boolean = false + override fun ping(): Boolean = connection.sync().ping().equals("PONG", ignoreCase = true) } /** * 没接 Redis 时的投影出口:什么都不做。 * - * 这不是"投影写成功"的承诺——`INV-10` 在真实客户端接通之前只能空转(`G-REDIS-PROJECTION`)。 - * 摆这个 bean 是为了让三步提交的时序在没有 Redis 的环境里也照常跑,而不是让每条航班报文都失败。 + * 这不是"投影写成功"的承诺。这个 bean 只让明确关闭投影的环境继续运行。 */ @Requires(property = "msgx.stubs", notEquals = "true") @Requires(property = "msgx.redis.enabled", notEquals = "true") @@ -48,7 +58,7 @@ class NoopFlightProjectionPort : FlightProjectionPort { private val log = org.slf4j.LoggerFactory.getLogger(NoopFlightProjectionPort::class.java) override fun write(writes: List) { - log.debug("redis projection disabled, skipping {} write(s) [G-REDIS-PROJECTION]", writes.size) + log.debug("redis projection disabled, skipping {} write(s)", writes.size) } /** 没有投影可读时报错而不是回空列表:空列表会被查询方当成"现在没有航班"(`INV-11`)。 */ 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 73e430b..c0961f3 100644 --- a/src/main/kotlin/com/gzzn/omms/msgexchange/processing/OutboundRequestService.kt +++ b/src/main/kotlin/com/gzzn/omms/msgexchange/processing/OutboundRequestService.kt @@ -14,7 +14,7 @@ import java.time.ZoneId import java.time.format.DateTimeFormatter /** - * 出站 REQ_TRACK 协调:登记、COUTMSGS 落信、超时与 EROR 失败(`G-REQ-TRACK`、`C-4`)。 + * 出站 REQ_TRACK 协调:登记、COUTMSGS 落信、超时与 EROR 失败(`C-4`)。 */ @Singleton class OutboundRequestService( 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 7a3b157..48d1ca0 100644 --- a/src/main/kotlin/com/gzzn/omms/msgexchange/processing/Pump.kt +++ b/src/main/kotlin/com/gzzn/omms/msgexchange/processing/Pump.kt @@ -206,7 +206,7 @@ class MessageProcessor( scheduleProcessor.applyScheduleRecords(head, decoded) MsgKind.SchdSubtype.RESP -> { if (!outbound.hasOpenSentRqfd()) { - log.info("SCHD-RESP without open RQFD -> SKIPPED msgId={} [G-RESP-GUARD]", head.msgId) + log.info("SCHD-RESP without open RQFD -> SKIPPED msgId={}", head.msgId) procState.markTerminal( head.msgId, ProcStatus.SKIPPED, lastError = "resp-guard:no-open-req", diff --git a/src/main/resources/application.yml b/src/main/resources/application.yml index 42c034d..3790a3c 100644 --- a/src/main/resources/application.yml +++ b/src/main/resources/application.yml @@ -36,10 +36,15 @@ msgx: arrived-hours: 1 # 条件 5:实际到港早于当前超过 1 小时 snap-log-retention-days: 90 history-store-enabled: false # 未接通历史存储时 HISTORY_SWEEP 删 0 条 - redis: # 航班查询投影(C-11);真实客户端未接通(G-REDIS-PROJECTION) - enabled: false # 关闭时投影写是空操作;开启但无客户端时写入直接失败,不伪装成功 + redis: # 航班查询投影(C-11) + enabled: false # 关闭时投影写是空操作 flight-key: flightInfo # 航班快照哈希键,field = FLID(沿用旧系统) +# Lettuce 连接由 msgx.redis.enabled 同步启停;地址必须由环境覆盖后再启用。 +redis: + enabled: ${msgx.redis.enabled} + uri: ${MSGX_REDIS_URI:redis://127.0.0.1:6379} + micronaut: application: name: msgexchange-nextgen diff --git a/src/test/kotlin/com/gzzn/omms/msgexchange/infra/persistence/jdbc/FlywayMigrationTest.kt b/src/test/kotlin/com/gzzn/omms/msgexchange/infra/persistence/jdbc/FlywayMigrationTest.kt index 60d4bd0..d1b8c6c 100644 --- a/src/test/kotlin/com/gzzn/omms/msgexchange/infra/persistence/jdbc/FlywayMigrationTest.kt +++ b/src/test/kotlin/com/gzzn/omms/msgexchange/infra/persistence/jdbc/FlywayMigrationTest.kt @@ -11,7 +11,7 @@ import java.sql.DriverManager /** * 在真实 PostgreSQL 上跑一遍迁移链,确认结果符合预期:`V1__flight_state_baseline.sql` - * 基线加 `V2__flight_chute_class_type_rename.sql` 更名、`V3__msg_event_hold.sql` 事件归属列 + * 基线加 V2~V6 增量迁移 * 依序执行成功,该建的表和 PIPELINE_LOCK 单行种子都在,回填事实落在基线里,`INBOX_CURSOR`、 * `BACKFILL_TODO`、`idx_evt_flid`、`PROC_STATE` 的处理开始时间列都不复存在;FLIGHT_CHUTE * 的类字段列已由 V2 更名为 CCLS/CTYP(SIS 口径)。 @@ -51,13 +51,17 @@ class FlywayMigrationTest { while (rs.next()) { records.add(Triple(rs.getString("version"), rs.getString("script"), rs.getBoolean("success"))) } - assertEquals(3, records.size, "迁移链应为 V1 基线 + V2 更名 + V3 事件归属三条") - assertEquals("1", records[0].first) - assertEquals("V1__flight_state_baseline.sql", records[0].second) - assertEquals("2", records[1].first) - assertEquals("V2__flight_chute_class_type_rename.sql", records[1].second) - assertEquals("3", records[2].first) - assertEquals("V3__msg_event_hold.sql", records[2].second) + assertEquals( + listOf( + "1" to "V1__flight_state_baseline.sql", + "2" to "V2__flight_chute_class_type_rename.sql", + "3" to "V3__msg_event_hold.sql", + "4" to "V4__req_track_outbound_seqn.sql", + "5" to "V5__unmapped_field_srvt_vipf.sql", + "6" to "V6__basicdata_ref_data.sql", + ), + records.map { it.first to it.second }, + ) assertTrue(records.all { it.third }) } diff --git a/src/test/kotlin/com/gzzn/omms/msgexchange/infra/persistence/jdbc/JdbcFlightStateRoutePgTest.kt b/src/test/kotlin/com/gzzn/omms/msgexchange/infra/persistence/jdbc/JdbcFlightStateRoutePgTest.kt index 00d63d9..ed875c7 100644 --- a/src/test/kotlin/com/gzzn/omms/msgexchange/infra/persistence/jdbc/JdbcFlightStateRoutePgTest.kt +++ b/src/test/kotlin/com/gzzn/omms/msgexchange/infra/persistence/jdbc/JdbcFlightStateRoutePgTest.kt @@ -79,6 +79,26 @@ class JdbcFlightStateRoutePgTest { assertEquals(4, loaded.collections["ERUT"]!!.size) } + @Test + fun `CHDT persists and loads SIS class fields`() { + assumeTrue(PgTestSupport.canConnect(), PgTestSupport.skipMessage()) + val repo = JdbcFlightStateRepository(dataSource(), Clock.fixed(t0, ZoneOffset.UTC)) + val flid = "CH-" + UUID.randomUUID().toString().take(8) + val chdt = mapOf("CHNO" to "1", "CHUT" to "C01", "CCLS" to "A", "CTYP" to "IN") + + repo.persistFullState( + FlightSnapshot( + flid, LocalDate.of(2026, 9, 12), FlightState.ACTIVE, 1, + mapOf("SODT" to "12Sep261200"), + mapOf("CHDT" to listOf(chdt)), + ), + msgId = 2, + now = t0, + ) + + assertEquals(listOf(chdt), repo.loadFullSnapshot(flid)!!.collections["CHDT"]) + } + private fun dataSource(): javax.sql.DataSource { Flyway.configure() .dataSource(PgTestSupport.jdbcUrl, PgTestSupport.user, PgTestSupport.password) diff --git a/src/test/kotlin/com/gzzn/omms/msgexchange/infra/projection/RedisFlightProjectionPortTest.kt b/src/test/kotlin/com/gzzn/omms/msgexchange/infra/projection/RedisFlightProjectionPortTest.kt new file mode 100644 index 0000000..4c8a221 --- /dev/null +++ b/src/test/kotlin/com/gzzn/omms/msgexchange/infra/projection/RedisFlightProjectionPortTest.kt @@ -0,0 +1,40 @@ +package com.gzzn.omms.msgexchange.infra.projection + +import com.gzzn.omms.msgexchange.config.RedisProps +import com.gzzn.omms.msgexchange.processing.FlightProjectionWrite +import com.gzzn.omms.msgexchange.support.RedisTestSupport +import io.lettuce.core.RedisClient +import org.junit.jupiter.api.Assertions.assertEquals +import org.junit.jupiter.api.Assertions.assertTrue +import org.junit.jupiter.api.Assumptions.assumeTrue +import org.junit.jupiter.api.Test +import java.util.UUID + +class RedisFlightProjectionPortTest { + + @Test + fun `writes reads and deletes flight hash in real Valkey`() { + assumeTrue(RedisTestSupport.canConnect(), RedisTestSupport.skipMessage()) + val client = RedisClient.create(RedisTestSupport.uri) + client.connect().use { connection -> + val key = "flightInfo:test:${UUID.randomUUID()}" + val port = RedisFlightProjectionPort(RedisProps().apply { flightKey = key }, connection) + try { + port.write( + listOf( + FlightProjectionWrite.Upsert("F1", 1, """{"flid":"F1","stateVersion":1}"""), + FlightProjectionWrite.Upsert("F2", 1, """{"flid":"F2","stateVersion":1}"""), + ), + ) + assertEquals(setOf("F1", "F2"), port.readAll().map { Regex("F\\d").find(it)!!.value }.toSet()) + + port.write(listOf(FlightProjectionWrite.Delete("F1", 2))) + assertEquals(listOf("""{"flid":"F2","stateVersion":1}"""), port.readAll()) + assertTrue(port.ping()) + } finally { + connection.sync().del(key) + } + } + client.shutdown() + } +} diff --git a/src/test/kotlin/com/gzzn/omms/msgexchange/support/RedisTestSupport.kt b/src/test/kotlin/com/gzzn/omms/msgexchange/support/RedisTestSupport.kt new file mode 100644 index 0000000..a89ae19 --- /dev/null +++ b/src/test/kotlin/com/gzzn/omms/msgexchange/support/RedisTestSupport.kt @@ -0,0 +1,31 @@ +package com.gzzn.omms.msgexchange.support + +import org.testcontainers.containers.GenericContainer + +private class ValkeyContainer(image: String) : GenericContainer(image) + +/** Real Valkey endpoint for projection adapter tests; skips when neither env nor Docker is available. */ +object RedisTestSupport { + private val envUri = System.getenv("MSGX_REDIS_URI") + + private val container: ValkeyContainer? by lazy { + if (envUri != null) null else runCatching { + ValkeyContainer("valkey/valkey:8-alpine").withExposedPorts(6379).also { it.start() } + }.getOrNull() + } + + val uri: String + get() = envUri ?: container?.let { "redis://${it.host}:${it.getMappedPort(6379)}" } + ?: "redis://127.0.0.1:6379" + + fun canConnect(): Boolean = runCatching { + val client = io.lettuce.core.RedisClient.create(uri) + try { + client.connect().use { it.sync().ping() } + } finally { + client.shutdown() + } + }.isSuccess + + fun skipMessage(): String = "Valkey not accessible (set MSGX_REDIS_URI or enable Docker for Testcontainers)" +}