From 4d8dd54e591f92a51e2dda9c2083ee9764974b14 Mon Sep 17 00:00:00 2001 From: kkfluous Date: Mon, 20 Apr 2026 15:41:59 +0800 Subject: [PATCH] refactor(gb32960): extract channelRead0 command branches into private methods MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 把 Gb32960ChannelHandler.channelRead0 的 140 行 if/else-if 级联拆成按 CommandType 的 switch 分派到 7 个独立私有方法: - handleVehicleLogin (0x01) - handleVehicleLogout (0x04) - handlePlatformLogin (0x05,返回 boolean 表示是否继续 dispatch) - handlePlatformLogout (0x06) - handleReportOrHeartbeat (0x02/0x03/0x07) - handleTimeCalibration (0x08) - logOtherFrame (其它命令的调试日志) 同时抽出 eventOrNow(msg) 静态 helper,消除 6 处 "msg.header().eventTime() != null ? msg.header().eventTime() : Instant.now()" 的重复。 纯重构,语义零变化:ACK 写入顺序、PlatformLogin 鉴权失败短路、日志 tag、VIN 白名单 预检都保持原行为。全模块 41 tests + 全仓 E2E 全绿,无回归。 channelRead0 从 140 行收敛到约 45 行,每个分支处理器 20~40 行并带 §6.3.2 规范注释, 未来新增命令或调整 ack 时序只动对应私有方法。 Co-Authored-By: Claude Opus 4.7 (1M context) --- .../inbound/Gb32960ChannelHandler.java | 245 +++++++++++------- 1 file changed, 147 insertions(+), 98 deletions(-) diff --git a/protocol-gb32960/src/main/java/com/lingniu/ingest/protocol/gb32960/inbound/Gb32960ChannelHandler.java b/protocol-gb32960/src/main/java/com/lingniu/ingest/protocol/gb32960/inbound/Gb32960ChannelHandler.java index 97c07101..86fbbe08 100644 --- a/protocol-gb32960/src/main/java/com/lingniu/ingest/protocol/gb32960/inbound/Gb32960ChannelHandler.java +++ b/protocol-gb32960/src/main/java/com/lingniu/ingest/protocol/gb32960/inbound/Gb32960ChannelHandler.java @@ -125,105 +125,18 @@ public class Gb32960ChannelHandler extends SimpleChannelInboundHandler { return; } - // 鉴权通过:对 0x01 登入主动下发成功应答(某些终端不收到应答会一直重发) - if (cmd == CommandType.VEHICLE_LOGIN) { - byte[] ack = Gb32960FrameEncoder.buildResponse( - msg.header().protocolVersion(), - CommandType.VEHICLE_LOGIN, - ResponseFlag.SUCCESS, - rawVin, - msg.header().eventTime() != null ? msg.header().eventTime() : Instant.now(), - null); - writeAck(ctx, ack, "vehicle-login-ack"); - log.info("[gb32960] vehicle login peer={} vin={} protocolVersion={}", - addr(ctx), vin, msg.header().protocolVersion()); - } else if (cmd == CommandType.VEHICLE_LOGOUT) { - log.info("[gb32960] vehicle logout peer={} vin={}", addr(ctx), vin); - writeAck(ctx, Gb32960FrameEncoder.buildResponse( - msg.header().protocolVersion(), - CommandType.VEHICLE_LOGOUT, - ResponseFlag.SUCCESS, - rawVin, - msg.header().eventTime() != null ? msg.header().eventTime() : Instant.now(), - null), "vehicle-logout-ack"); - } else if (cmd == CommandType.PLATFORM_LOGIN) { - // 平台登入 0x05(表 29):含 12B username + 20B password + 加密规则 - if (!(msg.commandBody() instanceof CommandBody.PlatformLogin pl)) { - log.warn("[gb32960] platform login peer={} body unparsed, closing", addr(ctx)); - ctx.close(); - return; + // 鉴权通过:按命令类型回 ACK,写成功后统一 dispatch 到通用管线。 + // 平台登入鉴权失败会短路 return,不进 dispatch。 + switch (cmd) { + case VEHICLE_LOGIN -> handleVehicleLogin(ctx, msg, rawVin); + case VEHICLE_LOGOUT -> handleVehicleLogout(ctx, msg, rawVin); + case PLATFORM_LOGIN -> { + if (!handlePlatformLogin(ctx, msg, rawVin)) return; } - String peerIp = peerIp(ctx); - Gb32960PlatformAuthorizer.Result result = platformAuthorizer.authenticate( - pl.username(), pl.password(), peerIp); - if (!result.accepted()) { - log.warn("[gb32960] platform login REJECTED peer={} username={} reason={}", - addr(ctx), pl.username(), result); - byte[] nack = Gb32960FrameEncoder.buildResponse( - msg.header().protocolVersion(), - CommandType.PLATFORM_LOGIN, - ResponseFlag.OTHER_ERROR, - rawVin, - msg.header().eventTime() != null ? msg.header().eventTime() : Instant.now(), - null); - ctx.writeAndFlush(Unpooled.wrappedBuffer(nack)).addListener(f -> { - log.info("[gb32960] platform login NACK flushed success={} peer={}", - f.isSuccess(), addr(ctx)); - ctx.close(); - }); - return; - } - log.info("[gb32960] platform login peer={} username={} encryptRule={} serial={} time={} policy={}", - addr(ctx), pl.username(), pl.encryptRule(), pl.serialNo(), pl.loginTime(), result); - // 把 username 钉到 channel attribute,供后续 0x02/0x03 帧的 vendor profile 路由 - ctx.channel().attr(PLATFORM_ACCOUNT_ATTR).set(pl.username()); - byte[] ack = Gb32960FrameEncoder.buildResponse( - msg.header().protocolVersion(), - CommandType.PLATFORM_LOGIN, - ResponseFlag.SUCCESS, - rawVin, - msg.header().eventTime() != null ? msg.header().eventTime() : Instant.now(), - null); - writeAck(ctx, ack, "platform-login-ack"); - } else if (cmd == CommandType.PLATFORM_LOGOUT) { - log.info("[gb32960] platform logout peer={}", addr(ctx)); - writeAck(ctx, Gb32960FrameEncoder.buildResponse( - msg.header().protocolVersion(), - CommandType.PLATFORM_LOGOUT, - ResponseFlag.SUCCESS, - rawVin, - msg.header().eventTime() != null ? msg.header().eventTime() : Instant.now(), - null), "platform-logout-ack"); - } else if (cmd == CommandType.REALTIME_REPORT - || cmd == CommandType.RESEND_REPORT - || cmd == CommandType.HEARTBEAT) { - // 实时 / 补发 / 心跳:按 §6.3.2 回应答帧,保留原帧采集时间,data 段仅 6B 时间。 - // 部分车端实现严格按规范,没收到应答会停止后续上报或重连。 - writeAck(ctx, Gb32960FrameEncoder.buildResponse( - msg.header().protocolVersion(), - cmd, - ResponseFlag.SUCCESS, - rawVin, - msg.header().eventTime() != null ? msg.header().eventTime() : Instant.now(), - null), "report-ack-0x" + Integer.toHexString(cmd.code())); - if (log.isDebugEnabled()) { - log.debug("[gb32960] frame peer={} vin={} cmd=0x{} infoBlocks={} json={}", - addr(ctx), vin, Integer.toHexString(cmd.code()), - msg.infoBlocks().size(), toJson(msg)); - } - } else if (cmd == CommandType.TIME_CALIBRATION) { - // 0x08 终端校时:平台应答 data 段为平台当前时间 6B(用 Instant.now())。 - writeAck(ctx, Gb32960FrameEncoder.buildResponse( - msg.header().protocolVersion(), - CommandType.TIME_CALIBRATION, - ResponseFlag.SUCCESS, - rawVin, - Instant.now(), - null), "time-calibration-ack"); - } else if (log.isDebugEnabled()) { - log.debug("[gb32960] frame peer={} vin={} cmd=0x{} infoBlocks={} json={}", - addr(ctx), vin, Integer.toHexString(cmd.code()), - msg.infoBlocks().size(), toJson(msg)); + case PLATFORM_LOGOUT -> handlePlatformLogout(ctx, msg, rawVin); + case REALTIME_REPORT, RESEND_REPORT, HEARTBEAT -> handleReportOrHeartbeat(ctx, msg, rawVin); + case TIME_CALIBRATION -> handleTimeCalibration(ctx, msg, rawVin); + default -> logOtherFrame(ctx, msg); } Map sourceMeta = new HashMap<>(4); @@ -240,6 +153,142 @@ public class Gb32960ChannelHandler extends SimpleChannelInboundHandler { dispatcher.dispatch(rf); } + /** 0x01 车辆登入:按 §6.3.2 带原采集时间回成功 ACK + INFO 日志。 */ + private void handleVehicleLogin(ChannelHandlerContext ctx, Gb32960Message msg, byte[] rawVin) { + byte[] ack = Gb32960FrameEncoder.buildResponse( + msg.header().protocolVersion(), + CommandType.VEHICLE_LOGIN, + ResponseFlag.SUCCESS, + rawVin, + eventOrNow(msg), + null); + writeAck(ctx, ack, "vehicle-login-ack"); + log.info("[gb32960] vehicle login peer={} vin={} protocolVersion={}", + addr(ctx), msg.header().vin(), msg.header().protocolVersion()); + } + + /** 0x04 车辆登出:INFO 日志 + 带原采集时间回成功 ACK。 */ + private void handleVehicleLogout(ChannelHandlerContext ctx, Gb32960Message msg, byte[] rawVin) { + log.info("[gb32960] vehicle logout peer={} vin={}", addr(ctx), msg.header().vin()); + writeAck(ctx, Gb32960FrameEncoder.buildResponse( + msg.header().protocolVersion(), + CommandType.VEHICLE_LOGOUT, + ResponseFlag.SUCCESS, + rawVin, + eventOrNow(msg), + null), "vehicle-logout-ack"); + } + + /** + * 0x05 平台登入(表 29):含 12B username + 20B password + 加密规则。 + * + *
    + *
  • 消息体解析缺失:关闭连接,返回 false 短路 + *
  • 鉴权失败:写 {@code 0x02 OTHER_ERROR} NACK + 关闭连接,返回 false 短路 + *
  • 鉴权通过:把 username 钉到 channel attribute(供后续 vendor profile 路由), + * 写成功 ACK + INFO 日志,返回 true 继续 dispatch + *
+ * + * @return true 表示鉴权通过应继续 dispatch;false 表示调用方应 short-circuit return + */ + private boolean handlePlatformLogin(ChannelHandlerContext ctx, Gb32960Message msg, byte[] rawVin) { + if (!(msg.commandBody() instanceof CommandBody.PlatformLogin pl)) { + log.warn("[gb32960] platform login peer={} body unparsed, closing", addr(ctx)); + ctx.close(); + return false; + } + String peerIp = peerIp(ctx); + Gb32960PlatformAuthorizer.Result result = platformAuthorizer.authenticate( + pl.username(), pl.password(), peerIp); + if (!result.accepted()) { + log.warn("[gb32960] platform login REJECTED peer={} username={} reason={}", + addr(ctx), pl.username(), result); + byte[] nack = Gb32960FrameEncoder.buildResponse( + msg.header().protocolVersion(), + CommandType.PLATFORM_LOGIN, + ResponseFlag.OTHER_ERROR, + rawVin, + eventOrNow(msg), + null); + ctx.writeAndFlush(Unpooled.wrappedBuffer(nack)).addListener(f -> { + log.info("[gb32960] platform login NACK flushed success={} peer={}", + f.isSuccess(), addr(ctx)); + ctx.close(); + }); + return false; + } + log.info("[gb32960] platform login peer={} username={} encryptRule={} serial={} time={} policy={}", + addr(ctx), pl.username(), pl.encryptRule(), pl.serialNo(), pl.loginTime(), result); + // 把 username 钉到 channel attribute,供后续 0x02/0x03 帧的 vendor profile 路由 + ctx.channel().attr(PLATFORM_ACCOUNT_ATTR).set(pl.username()); + byte[] ack = Gb32960FrameEncoder.buildResponse( + msg.header().protocolVersion(), + CommandType.PLATFORM_LOGIN, + ResponseFlag.SUCCESS, + rawVin, + eventOrNow(msg), + null); + writeAck(ctx, ack, "platform-login-ack"); + return true; + } + + /** 0x06 平台登出:INFO 日志 + 带原采集时间回成功 ACK。 */ + private void handlePlatformLogout(ChannelHandlerContext ctx, Gb32960Message msg, byte[] rawVin) { + log.info("[gb32960] platform logout peer={}", addr(ctx)); + writeAck(ctx, Gb32960FrameEncoder.buildResponse( + msg.header().protocolVersion(), + CommandType.PLATFORM_LOGOUT, + ResponseFlag.SUCCESS, + rawVin, + eventOrNow(msg), + null), "platform-logout-ack"); + } + + /** + * 0x02/0x03/0x07 实时 / 补发 / 心跳:按 §6.3.2 回应答帧,保留原帧采集时间。 + * 部分车端实现严格按规范,没收到应答会停止后续上报或重连。ack-tag 按 cmd.code 区分便于日志筛选。 + */ + private void handleReportOrHeartbeat(ChannelHandlerContext ctx, Gb32960Message msg, byte[] rawVin) { + CommandType cmd = msg.header().command(); + writeAck(ctx, Gb32960FrameEncoder.buildResponse( + msg.header().protocolVersion(), + cmd, + ResponseFlag.SUCCESS, + rawVin, + eventOrNow(msg), + null), "report-ack-0x" + Integer.toHexString(cmd.code())); + if (log.isDebugEnabled()) { + log.debug("[gb32960] frame peer={} vin={} cmd=0x{} infoBlocks={} json={}", + addr(ctx), msg.header().vin(), Integer.toHexString(cmd.code()), + msg.infoBlocks().size(), toJson(msg)); + } + } + + /** 0x08 终端校时:平台应答 data 段为平台当前时间 6B(用 Instant.now())。 */ + private void handleTimeCalibration(ChannelHandlerContext ctx, Gb32960Message msg, byte[] rawVin) { + writeAck(ctx, Gb32960FrameEncoder.buildResponse( + msg.header().protocolVersion(), + CommandType.TIME_CALIBRATION, + ResponseFlag.SUCCESS, + rawVin, + Instant.now(), + null), "time-calibration-ack"); + } + + /** 其它未单独处理的命令(如激活 0x09、密钥交换 0x0A):仅 DEBUG 日志,仍会进 dispatcher。 */ + private void logOtherFrame(ChannelHandlerContext ctx, Gb32960Message msg) { + if (log.isDebugEnabled()) { + log.debug("[gb32960] frame peer={} vin={} cmd=0x{} infoBlocks={} json={}", + addr(ctx), msg.header().vin(), Integer.toHexString(msg.header().command().code()), + msg.infoBlocks().size(), toJson(msg)); + } + } + + /** 优先使用帧内采集时间;为 null 时回落到 {@link Instant#now()}。 */ + private static Instant eventOrNow(Gb32960Message msg) { + return msg.header().eventTime() != null ? msg.header().eventTime() : Instant.now(); + } + /** * 统一的 ack 写出 + 日志确认。 *