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 写出 + 日志确认。 *