From d7c0a3dd9199fe9f5d1104be98ea6d82bdf26f9e Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=E7=8E=8B=E5=AF=85?= <97163845@qq.com> Date: Wed, 9 Sep 2026 18:04:05 +0800 Subject: [PATCH] =?UTF-8?q?=E4=BF=AE=E6=94=B9=E7=BA=A2=E5=A4=96ptzcontroll?= =?UTF-8?q?er?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../com/inspect/nvr/config/HikLibConfig.java | 9 +- .../nvr/controller/IvsCameraController.java | 7 + .../HikNvrPlaybackExceptionHandler.java | 1 + .../hik/domain/HikNvrLiveStreamResponse.java | 3 + .../domain/HikNvrPlaybackErrorResponse.java | 3 + .../exception/HikNvrPlaybackException.java | 70 +- .../nvr/hik/service/HikIvsVendorService.java | 12 +- .../hik/service/HikNvrLiveStreamService.java | 1218 ++++++++++++++--- .../nvr/hik/vo/HikNvrLiveStreamConfigVo.java | 21 +- .../inspect/nvr/service/HikLoginService.java | 331 +++-- 10 files changed, 1370 insertions(+), 305 deletions(-) diff --git a/src/main/java/com/inspect/nvr/config/HikLibConfig.java b/src/main/java/com/inspect/nvr/config/HikLibConfig.java index 2c978ccc..c2835817 100644 --- a/src/main/java/com/inspect/nvr/config/HikLibConfig.java +++ b/src/main/java/com/inspect/nvr/config/HikLibConfig.java @@ -9,6 +9,7 @@ import com.inspect.nvr.service.impl.HikLoginResultCallBack; import com.sun.jna.Native; import com.sun.jna.Pointer; import com.sun.jna.ptr.IntByReference; +import org.springframework.beans.factory.annotation.Value; import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; import org.springframework.stereotype.Component; @@ -18,8 +19,9 @@ import org.springframework.stereotype.Component; @Component public class HikLibConfig { - /** ClientDemo 使用的连接超时时间,避免失联设备长时间占用控制请求。 */ - private static final int SDK_CONNECT_TIMEOUT_MILLIS = 5000; + /** 海康 SDK 建连超时,失联 NVR 必须在实时预览总超时内快速返回。 */ + @Value("${hik.sdk.connect-timeout-millis:1000}") + private int sdkConnectTimeoutMillis; /** 海康 Linux64 SDK 6.1.11.5 随包提供的 OpenSSL 加密库文件名。 */ private static final String LINUX_CRYPTO_LIBRARY_NAME = "libcrypto.so.3"; @@ -101,7 +103,8 @@ public class HikLibConfig { //启动SDK写日志 hcNetSDK.NET_DVR_SetLogToFile(3, "./sdkLog", false); //设置连接时间与重连时间 - hcNetSDK.NET_DVR_SetConnectTime(SDK_CONNECT_TIMEOUT_MILLIS, 1); + // 将 SDK 单次建连限制在可配置的短超时内,避免断网时占满 Tomcat 请求线程。 + hcNetSDK.NET_DVR_SetConnectTime(Math.max(500, sdkConnectTimeoutMillis), 1); hcNetSDK.NET_DVR_SetReconnect(100000,true); //连接所有nvr diff --git a/src/main/java/com/inspect/nvr/controller/IvsCameraController.java b/src/main/java/com/inspect/nvr/controller/IvsCameraController.java index 62b79f52..c5f45e95 100644 --- a/src/main/java/com/inspect/nvr/controller/IvsCameraController.java +++ b/src/main/java/com/inspect/nvr/controller/IvsCameraController.java @@ -63,6 +63,13 @@ public class IvsCameraController { return ResponseEntity.ok().body(ivsOpenApiService.ptzControl(param)); } + @PostMapping({"/v1/device/ptzcontrol"}) + public ResponseEntity ptzControl(@RequestBody PtzControlParam param) { + log.info("Ptz control request param: {}", param); + return ResponseEntity + .ok() + .body(ivsCameraService.ptzControl(param)); + } /** * 接收前端视频框选坐标,并按框选方向执行 3D 定位放大或缩小。 * diff --git a/src/main/java/com/inspect/nvr/hik/controller/HikNvrPlaybackExceptionHandler.java b/src/main/java/com/inspect/nvr/hik/controller/HikNvrPlaybackExceptionHandler.java index 95d91e81..d03a44d4 100644 --- a/src/main/java/com/inspect/nvr/hik/controller/HikNvrPlaybackExceptionHandler.java +++ b/src/main/java/com/inspect/nvr/hik/controller/HikNvrPlaybackExceptionHandler.java @@ -26,6 +26,7 @@ public class HikNvrPlaybackExceptionHandler { // 组装错误码和中文错误消息。 HikNvrPlaybackErrorResponse response = HikNvrPlaybackErrorResponse.builder() .code(exception.getCode()) + .errorCode(exception.getErrorCode()) .message(exception.getMessage()) .build(); // 使用异常携带的 HTTP 状态码返回错误响应。 diff --git a/src/main/java/com/inspect/nvr/hik/domain/HikNvrLiveStreamResponse.java b/src/main/java/com/inspect/nvr/hik/domain/HikNvrLiveStreamResponse.java index fe696599..6a0d20d2 100644 --- a/src/main/java/com/inspect/nvr/hik/domain/HikNvrLiveStreamResponse.java +++ b/src/main/java/com/inspect/nvr/hik/domain/HikNvrLiveStreamResponse.java @@ -29,6 +29,9 @@ public class HikNvrLiveStreamResponse { /** 会话状态。 */ private String status; + /** 启动失败时返回的稳定字符串错误码。 */ + private String errorCode; + /** HLS 播放地址。 */ private String hlsUrl; diff --git a/src/main/java/com/inspect/nvr/hik/domain/HikNvrPlaybackErrorResponse.java b/src/main/java/com/inspect/nvr/hik/domain/HikNvrPlaybackErrorResponse.java index fcf0af13..7a7bbc4a 100644 --- a/src/main/java/com/inspect/nvr/hik/domain/HikNvrPlaybackErrorResponse.java +++ b/src/main/java/com/inspect/nvr/hik/domain/HikNvrPlaybackErrorResponse.java @@ -14,6 +14,9 @@ public class HikNvrPlaybackErrorResponse { /** 参数错误或海康 SDK 错误码。 */ private int code; + /** 可供前端和运维直接判断故障类型的稳定字符串错误码。 */ + private String errorCode; + /** 中文错误说明。 */ private String message; } diff --git a/src/main/java/com/inspect/nvr/hik/exception/HikNvrPlaybackException.java b/src/main/java/com/inspect/nvr/hik/exception/HikNvrPlaybackException.java index 44039c25..c65caa46 100644 --- a/src/main/java/com/inspect/nvr/hik/exception/HikNvrPlaybackException.java +++ b/src/main/java/com/inspect/nvr/hik/exception/HikNvrPlaybackException.java @@ -6,29 +6,34 @@ public class HikNvrPlaybackException extends RuntimeException { /** 对外返回的业务或 SDK 错误码。 */ private final int code; + /** 对外返回的稳定字符串错误码;旧接口异常未分类时为空。 */ + private final String errorCode; /** 对应的 HTTP 响应状态。 */ private final HttpStatus status; + /** 保存统一异常的数值码、字符串码、HTTP 状态和原始原因。 */ private HikNvrPlaybackException( int code, + String errorCode, HttpStatus status, String message, Throwable cause) { super(message, cause); this.code = code; + this.errorCode = errorCode; this.status = status; } /** 创建参数错误异常。 */ public static HikNvrPlaybackException badRequest(String message) { // 创建 HTTP 400 参数错误。 - return new HikNvrPlaybackException(-1, HttpStatus.BAD_REQUEST, message, null); + return new HikNvrPlaybackException(-1, null, HttpStatus.BAD_REQUEST, message, null); } /** 创建不带原始原因的 SDK/设备错误异常。 */ public static HikNvrPlaybackException sdkFailure(int code, String message) { // 创建 HTTP 502 厂商 SDK/设备调用错误。 - return new HikNvrPlaybackException(code, HttpStatus.BAD_GATEWAY, message, null); + return new HikNvrPlaybackException(code, null, HttpStatus.BAD_GATEWAY, message, null); } /** 创建带原始原因的 SDK/设备错误异常。 */ @@ -37,37 +42,80 @@ public class HikNvrPlaybackException extends RuntimeException { String message, Throwable cause) { // 创建带原始异常原因的 HTTP 502 设备错误。 - return new HikNvrPlaybackException(code, HttpStatus.BAD_GATEWAY, message, cause); + return new HikNvrPlaybackException(code, null, HttpStatus.BAD_GATEWAY, message, cause); + } + + /** 创建带稳定字符串错误码的 HTTP 502 NVR/通道异常。 */ + public static HikNvrPlaybackException sdkFailure( + String errorCode, + String message, + Throwable cause) { + // 保留实时预览历史 502 语义,同时增加可机读故障分类。 + return new HikNvrPlaybackException( + -1, errorCode, HttpStatus.BAD_GATEWAY, message, cause); + } + + /** 创建带稳定字符串错误码且不含原始原因的 HTTP 502 NVR/通道异常。 */ + public static HikNvrPlaybackException sdkFailure( + String errorCode, + String message) { + // 没有底层异常时仍保留可机读故障分类。 + return sdkFailure(errorCode, message, null); } /** 创建带原始原因的媒体链路错误异常。 */ public static HikNvrPlaybackException mediaFailure(String message, Throwable cause) { // 创建 HTTP 503 媒体链路错误。 - return new HikNvrPlaybackException(-3, HttpStatus.SERVICE_UNAVAILABLE, message, cause); + return new HikNvrPlaybackException(-3, null, HttpStatus.SERVICE_UNAVAILABLE, message, cause); } /** 创建不带原始原因的媒体链路错误异常。 */ public static HikNvrPlaybackException mediaFailure(String message) { // 创建不带原始异常的 HTTP 503 媒体错误。 - return new HikNvrPlaybackException(-3, HttpStatus.SERVICE_UNAVAILABLE, message, null); + return new HikNvrPlaybackException(-3, null, HttpStatus.SERVICE_UNAVAILABLE, message, null); + } + + /** 创建带稳定字符串错误码的媒体链路异常。 */ + public static HikNvrPlaybackException mediaFailure( + String errorCode, + String message, + Throwable cause) { + // 实时流调用方依赖字符串码区分 NVR、通道与 ZLMediaKit 故障。 + return new HikNvrPlaybackException( + -3, errorCode, HttpStatus.SERVICE_UNAVAILABLE, message, cause); + } + + /** 创建带稳定字符串错误码且不含原始原因的媒体链路异常。 */ + public static HikNvrPlaybackException mediaFailure( + String errorCode, + String message) { + // 没有底层异常时仍保留可机读错误类型。 + return mediaFailure(errorCode, message, null); } /** 创建资源不存在异常。 */ public static HikNvrPlaybackException notFound(String message) { // 创建 HTTP 404 资源不存在错误。 - return new HikNvrPlaybackException(-4, HttpStatus.NOT_FOUND, message, null); + return new HikNvrPlaybackException(-4, null, HttpStatus.NOT_FOUND, message, null); } /** 创建状态冲突异常。 */ public static HikNvrPlaybackException conflict(String message) { // 创建 HTTP 409 状态冲突错误。 - return new HikNvrPlaybackException(-5, HttpStatus.CONFLICT, message, null); + return new HikNvrPlaybackException(-5, null, HttpStatus.CONFLICT, message, null); + } + + /** 创建带稳定字符串错误码的状态冲突异常。 */ + public static HikNvrPlaybackException conflict(String errorCode, String message) { + // 起流队列满时用 409 快速释放 Tomcat 请求线程。 + return new HikNvrPlaybackException( + -5, errorCode, HttpStatus.CONFLICT, message, null); } /** 创建云台控制锁被占用异常,提示持锁角色和剩余秒数。 */ public static HikNvrPlaybackException operationLocked(String holderRoleName, int remainSeconds) { // 创建 HTTP 423 控制权被占错误。 - return new HikNvrPlaybackException(-423, HttpStatus.LOCKED, + return new HikNvrPlaybackException(-423, null, HttpStatus.LOCKED, holderRoleName + "正在进行操作,操作锁定剩余时间" + Math.max(remainSeconds, 0) + "秒,请稍后重试", null); } @@ -77,6 +125,12 @@ public class HikNvrPlaybackException extends RuntimeException { return code; } + /** 返回可供前端稳定判断故障类型的字符串错误码。 */ + public String getErrorCode() { + // 未分类的历史异常保持为空,兼容原有响应。 + return errorCode; + } + /** 返回 HTTP 状态。 */ public HttpStatus getStatus() { // 返回 HTTP 状态码。 diff --git a/src/main/java/com/inspect/nvr/hik/service/HikIvsVendorService.java b/src/main/java/com/inspect/nvr/hik/service/HikIvsVendorService.java index 8488e118..e629f752 100644 --- a/src/main/java/com/inspect/nvr/hik/service/HikIvsVendorService.java +++ b/src/main/java/com/inspect/nvr/hik/service/HikIvsVendorService.java @@ -686,8 +686,8 @@ public class HikIvsVendorService implements IvsVendorAdapter { if (request == null) { throw HikNvrPlaybackException.badRequest("预置位请求不能为空"); } - // 校验完整摄像机编码,确保新增动作落到正确的 NVR 通道。 - validateFullCameraCode(request.getCameraCode()); + // 校验摄像机编码,兼容完整“设备编码#域编码”和仅设备编码两种调用方式。 + validatePtzCameraCode(request.getCameraCode()); // 读取名称并去掉首尾空白,统一设备端保存的文本。 String name = request.getPresetName() == null ? "" : request.getPresetName().trim(); @@ -752,8 +752,8 @@ public class HikIvsVendorService implements IvsVendorAdapter { if (request == null || !hasText(request.getCameraCode())) { throw HikNvrPlaybackException.badRequest("cameraCode不能为空"); } - // 校验完整 cameraCode 的格式,确保能解析出设备域和通道信息。 - validateFullCameraCode(request.getCameraCode()); + // 校验 cameraCode 格式,兼容完整“设备编码#域编码”和仅设备编码两种调用方式。 + validatePtzCameraCode(request.getCameraCode()); // 读取预置位对象,后续从其中获取编号、名称和协议保留字段。 IvsPresetModifyRequest.PtzPresetInfo info = request.getPtzPresetInfo(); // 预置位对象或编号缺失时无法确定要修改哪一个预置位。 @@ -915,9 +915,9 @@ public class HikIvsVendorService implements IvsVendorAdapter { } /** - * 校验云台控制允许使用的摄像机编码格式。 + * 校验云台和预置位接口允许使用的摄像机编码格式。 * - *

云台控制兼容完整的“设备编码#域编码”和仅含设备编码的调用方式; + *

云台控制与预置位增改兼容完整的“设备编码#域编码”和仅含设备编码的调用方式; * 不带域编码时,后续根据设备编码前缀匹配唯一的 NVR。

* * @param cameraCode 完整摄像机编码或不带域后缀的设备编码 diff --git a/src/main/java/com/inspect/nvr/hik/service/HikNvrLiveStreamService.java b/src/main/java/com/inspect/nvr/hik/service/HikNvrLiveStreamService.java index 1d140fc7..55d9583c 100644 --- a/src/main/java/com/inspect/nvr/hik/service/HikNvrLiveStreamService.java +++ b/src/main/java/com/inspect/nvr/hik/service/HikNvrLiveStreamService.java @@ -1,304 +1,904 @@ package com.inspect.nvr.hik.service; +import com.alibaba.fastjson.JSON; +import com.alibaba.fastjson.JSONArray; +import com.alibaba.fastjson.JSONObject; +import com.inspect.nvr.domain.Infrared.NvrInfo; import com.inspect.nvr.domain.ivs.IvsRtspUrlResponse; import com.inspect.nvr.domain.ivs.IvsVideoRtspVo; +import com.inspect.nvr.hik.domain.HikNvrCameraTarget; import com.inspect.nvr.hik.domain.HikNvrLiveStreamResponse; import com.inspect.nvr.hik.exception.HikNvrPlaybackException; import com.inspect.nvr.hik.vo.HikNvrLiveStreamConfigVo; import com.inspect.nvr.ivs.service.IvsOpenApiService; +import com.inspect.nvr.service.HikLoginService; import lombok.extern.slf4j.Slf4j; import org.springframework.stereotype.Service; -import javax.annotation.PreDestroy; import javax.annotation.PostConstruct; +import javax.annotation.PreDestroy; import javax.annotation.Resource; +import java.io.ByteArrayOutputStream; import java.io.IOException; import java.io.InputStream; import java.net.HttpURLConnection; import java.net.URL; +import java.net.URLEncoder; import java.nio.charset.StandardCharsets; import java.nio.file.Files; import java.nio.file.Path; import java.nio.file.Paths; import java.util.ArrayList; import java.util.List; +import java.util.Locale; import java.util.Map; +import java.util.concurrent.ArrayBlockingQueue; +import java.util.concurrent.CompletableFuture; +import java.util.concurrent.ExecutionException; +import java.util.concurrent.RejectedExecutionException; +import java.util.concurrent.ThreadPoolExecutor; +import java.util.concurrent.TimeUnit; +import java.util.concurrent.TimeoutException; +import java.util.concurrent.atomic.AtomicLong; +import java.util.regex.Matcher; +import java.util.regex.Pattern; import java.util.concurrent.ConcurrentHashMap; /** - * 通过现有 IVS 门面按 cameraCode 获取 RTSP,再由 FFmpeg 推送到本地 ZLMediaKit - * 的 RTMP/HLS 浏览器实时预览服务。 + * 通过现有 IVS 门面取得 RTSP,由 FFmpeg 推送到 ZLMediaKit,并按 NVR 串行起流。 */ @Slf4j @Service public class HikNvrLiveStreamService { + /** NVR 拒绝并发 RTSP SETUP 时返回的稳定错误码。 */ + private static final String NVR_RTSP_BUSY = "NVR_RTSP_BUSY"; + /** 通道离线或在时限内没有任何视频数据时返回的稳定错误码。 */ + private static final String CHANNEL_OFFLINE = "CHANNEL_OFFLINE"; + /** NVR 网络不可达、拒绝连接或连接超时时返回的稳定错误码。 */ + private static final String NVR_UNREACHABLE = "NVR_UNREACHABLE"; + /** FFmpeg 无法向 ZLMediaKit 发布或 ZLM API 不可达时返回的稳定错误码。 */ + private static final String ZLM_PUSH_FAILED = "ZLM_PUSH_FAILED"; + /** SDK 登录或 RTSP 鉴权失败时返回的稳定错误码。 */ + private static final String NVR_LOGIN_FAILED = "NVR_LOGIN_FAILED"; + /** 单台 NVR 起流队列已满时返回的稳定错误码。 */ + private static final String NVR_START_QUEUE_FULL = "NVR_START_QUEUE_FULL"; + /** 从 SDK 异常链中提取海康原始错误码。 */ + private static final Pattern SDK_ERROR_PATTERN = + Pattern.compile("错误码[::]\\s*(\\d+)"); + /** 复用现有 IVS 媒体接口按 cameraCode 获取实时 RTSP 地址。 */ @Resource private IvsOpenApiService ivsOpenApiService; - /** 实时拉流配置 VO。 */ + /** 按 cameraCode 解析所属 NVR、通道和服务端数据库凭据。 */ + @Resource + private HikNvrCameraTargetService cameraTargetService; + /** 管理海康 SDK 登录缓存并在起流失败时执行一次重登。 */ + @Resource + private HikLoginService hikLoginService; + /** 实时拉流配置。 */ @Resource private HikNvrLiveStreamConfigVo config; + /** FFmpeg 可执行文件路径。 */ private String ffmpegPath; /** ZLMediaKit RTMP 基础地址。 */ private String rtmpBaseUrl; - /** ZLMediaKit HTTP/HLS 基础地址。 */ + /** 返回给浏览器的 ZLMediaKit HTTP 基础地址。 */ private String httpBaseUrl; - /** 等待实时流就绪的最大时间。 */ + /** 服务端访问 ZLMediaKit API 的 HTTP 基础地址。 */ + private String zlmApiBaseUrl; + /** 服务端调用 ZLMediaKit API 使用的密钥。 */ + private String zlmApiSecret; + /** 单个实时流从启动到明确成功或失败的总时限。 */ private long startupTimeoutMillis; - /** 按 IVS 完整 cameraCode 保存实时拉流会话。 */ + /** 同一 NVR 相邻起流握手的最小间隔。 */ + private long startSpacingMillis; + /** 每台 NVR 等待队列的最大任务数。 */ + private int maxStartQueueSize; + /** ZLMediaKit API 单次连接和读取超时。 */ + private int zlmApiTimeoutMillis; + + /** 按 IVS 完整 cameraCode 保存本服务拥有的 FFmpeg 拉流会话。 */ private final Map sessions = new ConcurrentHashMap<>(); + /** 按 cameraCode 保存排队或正在执行的起流任务。 */ + private final Map pendingStarts = new ConcurrentHashMap<>(); + /** 保存异步排队任务最近一次明确失败,供状态接口返回具体错误。 */ + private final Map failedStarts = + new ConcurrentHashMap<>(); + /** 每台 NVR 独立维护一个有界单线程起流通道。 */ + private final Map startLanes = new ConcurrentHashMap<>(); - /** 将实时拉流配置 VO 转换为服务运行时使用的规范化配置。 */ + /** 将实时拉流配置转换为服务运行时使用的规范化参数。 */ @PostConstruct public void initialize() { this.ffmpegPath = config.getFfmpegPath(); this.rtmpBaseUrl = trimTrailingSlash(config.getRtmpBaseUrl()); this.httpBaseUrl = trimTrailingSlash(config.getHttpBaseUrl()); - this.startupTimeoutMillis = Math.max(3000L, config.getStartupTimeoutMillis()); + this.zlmApiBaseUrl = trimTrailingSlash(config.getApiBaseUrl()); + this.zlmApiSecret = config.getApiSecret() == null + ? "" : config.getApiSecret().trim(); + // 最低一秒避免正常 LAN 握手被误杀;现场实测流注册需 2~4 秒、闪断后更慢,默认 10 秒。 + this.startupTimeoutMillis = Math.max(1000L, config.getStartupTimeoutMillis()); + // 相邻握手间隔限制在不小于零,默认使用现场验证稳定的 600ms。 + this.startSpacingMillis = Math.max(0L, config.getStartSpacingMillis()); + // 至少保留一个等待位;默认八个等待位加一个执行位可接收九路同时起流。 + this.maxStartQueueSize = Math.max(1, config.getMaxStartQueueSize()); + // ZLM 探测本身必须显著短于总启动时限,避免探测请求反向拖死队列。 + this.zlmApiTimeoutMillis = Math.max(100, config.getZlmApiTimeoutMillis()); } /** - * 按页面选择的 cameraCode 通过现有 IVS 接口取得 RTSP,并启动 FFmpeg 转为浏览器流。 + * 按页面选择的 cameraCode 提交实时流;首个任务短暂等待,其余任务立即返回排队状态。 * * @param cameraCode 页面设备树提供的 IVS 完整摄像机编码 - * @param streamType 实时码流类型 - * @return 实时流会话 + * @param streamType 主码流 1 或子码流 2 + * @return 已播放、正在排队或正在启动的实时流状态 */ - public synchronized HikNvrLiveStreamResponse start( - String cameraCode, Integer streamType) { - // cameraCode 是实时预览唯一设备标识,NVR 与通道由现有 IVS 能力内部解析。 + public HikNvrLiveStreamResponse start(String cameraCode, Integer streamType) { + // cameraCode 是实时预览唯一设备标识,NVR 与通道由服务端数据库解析。 String targetCameraCode = requireCameraCode(cameraCode); - // 启动前校验 FFmpeg 与媒体网关基础配置。 + // 启动前校验 FFmpeg、ZLM API 和媒体网关配置,缺项时立即返回明确错误。 ensureConfigured(); Integer type = streamType; - // 实时预览仅支持主码流 1 和子码流 2。 + // 实时预览只接受海康主码流 1 和子码流 2。 if (type == null || (type != 1 && type != 2)) { - throw HikNvrPlaybackException.badRequest("码流类型只能是1(主码流)或2(子码流)"); + throw HikNvrPlaybackException.badRequest( + "码流类型只能是1(主码流)或2(子码流)"); + } + + // 解析真实 NVR IP 和通道,队列必须以设备而不是全局或摄像机为粒度。 + HikNvrCameraTarget target = cameraTargetService.resolve(targetCameraCode); + NvrInfo nvrInfo = target.getNvrInfo(); + // 数据库目标缺少 IP 时无法建立每台 NVR 的队列,也不能安全起流。 + if (nvrInfo == null || !hasText(nvrInfo.getNvrIp())) { + throw HikNvrPlaybackException.sdkFailure( + NVR_LOGIN_FAILED, + NVR_LOGIN_FAILED + ":cameraCode未解析到有效NVR"); } - // 完整 cameraCode 自带设备、通道与域信息,可直接作为会话唯一键。 String sessionKey = sessionKey(targetCameraCode); + String streamId = buildStreamId(targetCameraCode); + PendingStart pending = new PendingStart( + target, type, streamId, hlsUrl(streamId), flvUrl(streamId)); + + // 相同摄像机已有任务时复用;码流切换则取消旧任务并用新任务替换。 + while (true) { + PendingStart previous = pendingStarts.putIfAbsent(sessionKey, pending); + // 没有旧任务时已经取得该摄像机的排队所有权。 + if (previous == null) { + break; + } + // 相同码流正在排队或启动时直接返回现有任务,避免重复拉流。 + if (previous.streamType == type) { + return toPendingResponse(previous, "实时预览正在排队或启动"); + } + // 码流类型变化时取消旧任务,旧任务会在下一探测点释放已启动的 FFmpeg。 + previous.cancelled = true; + // 原子替换成功后由当前请求提交新码流任务;竞争失败则重新检查最新任务。 + if (pendingStarts.replace(sessionKey, previous, pending)) { + break; + } + } + + // 新起流覆盖该摄像机上一次失败,状态接口从此返回本次任务状态。 + failedStarts.remove(sessionKey); + String nvrKey = nvrInfo.getNvrIp().trim().toLowerCase(Locale.ROOT); + // 为目标 NVR 创建或复用有界单线程队列,其他 NVR 不受该设备阻塞影响。 + NvrStartLane lane = startLanes.computeIfAbsent( + nvrKey, key -> new NvrStartLane(key, maxStartQueueSize)); + boolean firstInLane; + try { + // 同步提交计数,准确判断本请求是否是当前 NVR 的首个任务。 + synchronized (lane) { + firstInLane = lane.outstanding == 0; + lane.outstanding++; + // 真正的 SDK 登录、FFmpeg 和 ZLM 探测都在 NVR 专用线程执行。 + lane.executor.execute(() -> executePending(lane, sessionKey, pending)); + } + } catch (RejectedExecutionException rejection) { + // 队列满时撤销占位并快速返回,不占用 Tomcat 线程等待设备。 + synchronized (lane) { + lane.outstanding = Math.max(0, lane.outstanding - 1); + } + pendingStarts.remove(sessionKey, pending); + throw HikNvrPlaybackException.conflict( + NVR_START_QUEUE_FULL, + NVR_START_QUEUE_FULL + ":该NVR起流队列已满,请稍后重试"); + } + + // 非首个任务已安全进入队列,立即返回 STARTING 让前端轮询状态。 + if (!firstInLane) { + return toPendingResponse(pending, "实时预览排队中"); + } + try { + // 首个任务最多等待总启动时限加少量调度余量,避免 Tomcat 请求长时间挂死。 + return pending.future.get( + startupTimeoutMillis + 500L, TimeUnit.MILLISECONDS); + } catch (TimeoutException timeout) { + // 原生 SDK 若超出预期仍在工作,HTTP 先返回 STARTING,后台结果由状态接口提供。 + return toPendingResponse(pending, "实时预览启动中"); + } catch (InterruptedException interrupted) { + // 请求线程被中断时保留后台任务,并恢复线程中断标志。 + Thread.currentThread().interrupt(); + return toPendingResponse(pending, "实时预览启动中"); + } catch (ExecutionException execution) { + Throwable cause = execution.getCause(); + // 起流任务已返回统一业务异常时保持原始字符串错误码和 HTTP 状态。 + if (cause instanceof HikNvrPlaybackException) { + throw (HikNvrPlaybackException) cause; + } + // 未预期异常统一按媒体链路错误返回,避免泄漏线程实现细节。 + throw HikNvrPlaybackException.mediaFailure( + ZLM_PUSH_FAILED, + ZLM_PUSH_FAILED + ":实时预览后台任务异常", + cause); + } + } + + /** + * 在目标 NVR 的专用线程中执行起流,并保存异步失败供状态接口查询。 + * + * @param lane 目标 NVR 的串行起流通道 + * @param sessionKey 摄像机会话键 + * @param pending 当前排队任务 + */ + private void executePending( + NvrStartLane lane, + String sessionKey, + PendingStart pending) { + try { + // 等待期间已被停止或码流切换取消时不再连接 NVR。 + if (pending.cancelled) { + pending.future.complete(toStoppedResponse(pending, "实时预览启动已取消")); + return; + } + // 执行受总时限约束的两次起流尝试,第二次会先重建海康登录会话。 + HikNvrLiveStreamResponse response = startInternal(lane, sessionKey, pending); + pending.future.complete(response); + } catch (HikNvrPlaybackException failure) { + // 排队请求无法直接接收异常,因此把稳定错误码和消息保留给状态轮询接口。 + failedStarts.put(sessionKey, toFailureResponse(pending, failure)); + pending.future.completeExceptionally(failure); + log.warn("[live] start failed cameraCode={} nvrIp={} errorCode={} log={}", + pending.target.getCameraCode(), + pending.target.getNvrInfo().getNvrIp(), + failure.getErrorCode(), + logFile(pending.streamId)); + } catch (RuntimeException failure) { + // 未分类运行时错误也转换成 ZLM_PUSH_FAILED,保证前端得到明确类型。 + HikNvrPlaybackException classified = HikNvrPlaybackException.mediaFailure( + ZLM_PUSH_FAILED, + ZLM_PUSH_FAILED + ":实时预览内部处理失败", + failure); + failedStarts.put(sessionKey, toFailureResponse(pending, classified)); + pending.future.completeExceptionally(classified); + } finally { + // 只移除当前任务,码流切换时不能误删后来替换的新任务。 + pendingStarts.remove(sessionKey, pending); + synchronized (lane) { + lane.outstanding = Math.max(0, lane.outstanding - 1); + } + } + } + + /** + * 串行执行一次起流任务;失败后释放进程、重建 SDK 登录并重试一次。 + * + * @param lane 目标 NVR 的串行起流通道 + * @param sessionKey 摄像机会话键 + * @param pending 当前起流任务 + * @return 已经由 ZLMediaKit API 确认注册的播放响应 + */ + private HikNvrLiveStreamResponse startInternal( + NvrStartLane lane, + String sessionKey, + PendingStart pending) { LiveSession existing = sessions.get(sessionKey); - // 已有相同码流类型的存活会话时只等待其地址就绪,不重复启动 FFmpeg。 - if (existing != null && existing.isAlive() && existing.streamType == type) { - // 轮询已有 HLS 地址并返回当前播放状态。 - boolean ready = waitUntilReady(existing.hlsUrl, 2500L); - return toResponse(existing, ready ? "PLAYING" : "STARTING", - ready ? "实时预览已在播放" : "实时预览启动中"); - } - // 旧会话已退出或码流类型发生变化时先清理进程,确保新请求真正切换主辅码流。 + // 本地相同码流仍存活且 ZLM 已注册时直接复用,不重复发起 RTSP 握手。 + if (existing != null + && existing.isAlive() + && existing.streamType == pending.streamType + && queryZlmStream(existing.streamId) == ZlmProbeResult.REGISTERED) { + return toResponse(existing, "PLAYING", "实时预览已在播放"); + } + // 旧进程已退出或码流类型变化时先释放旧流,再启动本次请求。 if (existing != null) { stopInternal(sessionKey); + // 短暂等待 ZLM 注销旧发布,避免主辅码流切换撞上同名 publisher。 + waitUntilUnregistered(existing.streamId, 500L); + } + // 即使本地会话丢失,也先查 ZLM;同名流已存在时直接返回现有地址。 + if (queryZlmStream(pending.streamId) == ZlmProbeResult.REGISTERED) { + return toExternalResponse(pending, "PLAYING", "ZLMediaKit中已有同名实时流"); } - // 流标识由完整 cameraCode 生成,避免不同设备的同号通道发布冲突。 - String streamId = buildStreamId(targetCameraCode); - // 通过现有 IVS 实时媒体能力生成 RTSP,海康实时服务不再自行解析 NVR 和通道。 - String rtspUrl = resolveLiveRtspUrl(targetCameraCode, type); - String rtmpUrl = rtmpBaseUrl + "/live/" + streamId; - String hlsUrl = httpBaseUrl + "/live/" + streamId + "/hls.m3u8"; - String flvUrl = httpBaseUrl + "/live/" + streamId + ".live.flv"; - - Process process; - Path logFile = Paths.get("target", "live-" + streamId + ".log"); + Path logFile = logFile(pending.streamId); try { - // 确保 FFmpeg 日志目录存在。 + // 每次新的高层起流任务清空旧日志;内部两次尝试使用追加模式保留完整 stderr。 Files.createDirectories(logFile.getParent()); - // 根据厂商 RTSP 路径构造推流命令;大华 HEVC 转 H.264 后才能交给浏览器播放。 - List command = buildFfmpegCommand(rtspUrl, rtmpUrl); - - // 创建 FFmpeg 进程并把标准错误重定向到日志文件。 - ProcessBuilder builder = new ProcessBuilder(command); - builder.redirectErrorStream(true); - builder.redirectOutput(logFile.toFile()); - // 启动 FFmpeg 拉流推送进程。 - process = builder.start(); - log.info("[live] started cameraCode={} streamType={} streamId={} log={}", - targetCameraCode, type, streamId, logFile); - } catch (IOException ex) { - throw HikNvrPlaybackException.sdkFailure(-1, "启动实时预览失败:" + ex.getMessage()); - } - - LiveSession session = new LiveSession( - targetCameraCode, type, streamId, process, - hlsUrl, flvUrl, System.currentTimeMillis()); - // 保存新会话,供停止和状态查询复用。 - sessions.put(sessionKey, session); - - // 等待 HLS 地址生成,避免前端立即播放时拿到空资源。 - boolean ready = waitUntilReady(hlsUrl, startupTimeoutMillis); - // FFmpeg 已退出且地址未就绪时判定启动失败并清理进程。 - if (!ready && !session.isAlive()) { + Files.deleteIfExists(logFile); + } catch (IOException failure) { + throw HikNvrPlaybackException.mediaFailure( + ZLM_PUSH_FAILED, + ZLM_PUSH_FAILED + ":无法创建FFmpeg日志文件", + failure); + } + + long deadline = System.currentTimeMillis() + startupTimeoutMillis; + HikNvrPlaybackException lastFailure = null; + // 最多尝试两次;第二次之前会主动失效海康登录缓存并重登。 + for (int attempt = 1; attempt <= 2; attempt++) { + // 停止或码流切换已取消本任务时立即结束,不再占用 NVR 会话。 + if (pending.cancelled) { + return toStoppedResponse(pending, "实时预览启动已取消"); + } + // 第二次尝试前主动重建海康 SDK 登录会话,清除缓存中的僵死 userID。 + if (attempt == 2 && usesHikLogin(pending.target.getNvrInfo())) { + try { + // 先登出旧句柄再登录,后续 IVS RTSP 解析会复用这次新会话。 + hikLoginService.relogin(pending.target.getNvrInfo()); + } catch (RuntimeException loginFailure) { + throw classifyFailure("", ZlmProbeResult.ABSENT, loginFailure); + } + } + // 相同 NVR 的相邻握手保持配置间隔,避免设备瞬时并发 SETUP 500。 + if (!lane.awaitSpacing(deadline, startSpacingMillis)) { + break; + } + + String rtspUrl; + try { + // 通过 IVS 厂商路由生成当前码流的 RTSP 地址,海康路径会复用 SDK 登录布局。 + rtspUrl = resolveLiveRtspUrl( + pending.target.getCameraCode(), pending.streamType); + } catch (RuntimeException resolveFailure) { + lastFailure = classifyFailure( + readLog(logFile), ZlmProbeResult.ABSENT, resolveFailure); + // 第一次解析失败仍执行一次登录重建;第二次失败直接返回分类结果。 + if (attempt == 1) { + continue; + } + throw lastFailure; + } + + String rtmpUrl = rtmpBaseUrl + "/live/" + pending.streamId; + Process process; + try { + // 启动 FFmpeg 并把两次尝试的原始 stderr 追加到同一个每路日志文件。 + process = startFfmpeg(buildFfmpegCommand(rtspUrl, rtmpUrl), logFile); + } catch (IOException processFailure) { + lastFailure = classifyFailure( + readLog(logFile), ZlmProbeResult.UNAVAILABLE, processFailure); + // 首次本地启动失败仍按统一自愈流程重建登录并重试一次。 + if (attempt == 1) { + continue; + } + throw lastFailure; + } + + LiveSession session = new LiveSession( + pending.target.getCameraCode(), + pending.target.getNvrInfo().getNvrIp(), + pending.target.getChannel(), + pending.streamType, + pending.streamId, + process, + pending.hlsUrl, + pending.flvUrl); + // 注册本地会话后,停止和状态请求即可立即定位并释放正在启动的进程。 + sessions.put(sessionKey, session); + log.info("[live] started cameraCode={} nvrIp={} streamType={} streamId={} attempt={} log={}", + session.cameraCode, session.nvrIp, session.streamType, + session.streamId, attempt, logFile); + + long remaining = Math.max(0L, deadline - System.currentTimeMillis()); + // 第一次最多使用剩余时间的一半,为失效登录并重试保留预算。 + long attemptBudget = attempt == 1 + ? Math.max(400L, remaining / 2L) : remaining; + long attemptDeadline = Math.min( + deadline, System.currentTimeMillis() + attemptBudget); + // 轮询 ZLM API;只有媒体流已经注册才能向前端返回 PLAYING。 + ZlmProbeResult probe = waitUntilRegistered( + session, pending, attemptDeadline); + // ZLM 已确认同名 RTMP 流注册时才认为起流成功。 + if (probe == ZlmProbeResult.REGISTERED) { + failedStarts.remove(sessionKey); + return toResponse(session, "PLAYING", "实时预览已启动"); + } + + // 未注册或进程退出时立即销毁 FFmpeg,快速断开 NVR RTSP 会话。 stopInternal(sessionKey); - throw HikNvrPlaybackException.sdkFailure(-1, "实时预览拉流失败,请检查NVR RTSP与ZLMediaKit"); + lastFailure = classifyFailure(readLog(logFile), probe, null); + // SETUP 500 是已确认的设备带宽预算拒绝,重登不会释放带宽且会继续冲击 NVR。 + if (NVR_RTSP_BUSY.equals(lastFailure.getErrorCode())) { + throw lastFailure; + } + // 第一次失败进入一次登录失效和重试;第二次返回最终明确分类。 + if (attempt == 2) { + throw lastFailure; + } + } + // 总时限不足以继续时返回最后一次分类;走到这里说明尚未发起任何握手。 + if (lastFailure != null) { + throw lastFailure; } - return toResponse(session, ready ? "PLAYING" : "STARTING", - ready ? "实时预览已启动" : "实时预览启动中,请稍候"); + throw HikNvrPlaybackException.sdkFailure( + CHANNEL_OFFLINE, + CHANNEL_OFFLINE + ":启动时限" + + (startupTimeoutMillis / 1000L) + + "秒内未等到本NVR起流时隙,请稍后重试"); } /** - * 停止指定 cameraCode 的 FFmpeg 拉流进程并清理会话。 + * 停止指定 cameraCode 的排队任务或 FFmpeg 拉流进程。 * * @param cameraCode 页面设备树提供的 IVS 完整摄像机编码 * @return 幂等停止结果 */ - public synchronized HikNvrLiveStreamResponse stop(String cameraCode) { + public HikNvrLiveStreamResponse stop(String cameraCode) { // 规范化 cameraCode,确保启动、停止和状态查询命中同一会话键。 String targetCameraCode = requireCameraCode(cameraCode); - // 使用完整 cameraCode 定位唯一实时会话。 String sessionKey = sessionKey(targetCameraCode); - // 查找该摄像机的活动会话。 + PendingStart pending = pendingStarts.remove(sessionKey); + // 排队或正在启动的任务打上取消标记,工作线程会在最近探测点释放资源。 + if (pending != null) { + pending.cancelled = true; + } + failedStarts.remove(sessionKey); LiveSession session = sessions.get(sessionKey); - // 没有活动会话时返回幂等停止结果。 - if (session == null) { + // 找到本地 FFmpeg 会话时立即销毁,避免 RTSP 连接继续占用 NVR 会话池。 + if (session != null) { + stopInternal(sessionKey); return HikNvrLiveStreamResponse.builder() .cameraCode(targetCameraCode) + .nvrIp(session.nvrIp) + .channel(session.channel) + .streamId(session.streamId) + .streamType(session.streamType) .status("STOPPED") - .message("摄像机没有活跃实时预览") + .message("实时预览已停止") .build(); } - // 销毁活动 FFmpeg 进程并从会话表移除。 - stopInternal(sessionKey); + // 只有排队任务时返回其流标识,便于前端核对取消结果。 + if (pending != null) { + return toStoppedResponse(pending, "实时预览启动已取消"); + } return HikNvrLiveStreamResponse.builder() .cameraCode(targetCameraCode) - .streamId(session.streamId) - .streamType(session.streamType) .status("STOPPED") - .message("实时预览已停止") + .message("摄像机没有活跃实时预览") .build(); } /** - * 查询指定 cameraCode 当前实时拉流会话状态。 + * 查询指定 cameraCode 当前排队、播放或失败状态。 * * @param cameraCode 页面设备树提供的 IVS 完整摄像机编码 * @return 实时流状态 */ public HikNvrLiveStreamResponse status(String cameraCode) { - // 规范化 cameraCode,确保状态查询只读取目标摄像机的会话。 + // 规范化 cameraCode,确保状态查询读取唯一摄像机会话。 String targetCameraCode = requireCameraCode(cameraCode); - // 使用完整 cameraCode 定位唯一实时会话。 String sessionKey = sessionKey(targetCameraCode); LiveSession session = sessions.get(sessionKey); - // 会话不存在或进程已退出时清理旧映射并返回 STOPPED。 - if (session == null || !session.isAlive()) { - // 只有存在旧对象时才从并发会话表移除。 - if (session != null) { + // 本地进程存在时仍以 ZLM API 注册结果作为 PLAYING 的唯一依据。 + if (session != null) { + ZlmProbeResult probe = queryZlmStream(session.streamId); + // ZLM 已注册说明前端现在可以安全拉取 FLV/HLS。 + if (probe == ZlmProbeResult.REGISTERED) { + return toResponse(session, "PLAYING", "实时预览播放中"); + } + // 进程仍在时统一返回 STARTING;即使 ZLM API 短暂不可用,也不能误报为 STOPPED。 + if (session.isAlive()) { + String message = probe == ZlmProbeResult.UNAVAILABLE + ? "ZLMediaKit状态暂时无法确认" + : "实时预览启动中"; + return toResponse(session, "STARTING", message); + } + // 进程已经退出时移除旧映射,后面返回异步失败或停止状态。 + if (!session.isAlive()) { sessions.remove(sessionKey, session); } + } + + PendingStart pending = pendingStarts.get(sessionKey); + // 排队任务尚未产生进程时向前端返回地址和 STARTING,供统一轮询。 + if (pending != null) { + return toPendingResponse(pending, "实时预览排队中"); + } + HikNvrLiveStreamResponse failure = failedStarts.get(sessionKey); + // 异步任务已经失败时返回稳定错误码和具体消息,避免前端只显示通用超时。 + if (failure != null) { + return failure; + } + + String streamId = buildStreamId(targetCameraCode); + // 本地状态丢失但 ZLM 仍有同名流时覆盖容器重启或其他入口创建的流复用场景。 + if (queryZlmStream(streamId) == ZlmProbeResult.REGISTERED) { return HikNvrLiveStreamResponse.builder() .cameraCode(targetCameraCode) - .status("STOPPED") - .message("未在播放") + .streamId(streamId) + .status("PLAYING") + .hlsUrl(hlsUrl(streamId)) + .flvUrl(flvUrl(streamId)) + .message("ZLMediaKit中已有同名实时流") .build(); } - // 探测 HLS 内容是否已经可播放。 - boolean ready = isHlsReady(session.hlsUrl); - return toResponse(session, ready ? "PLAYING" : "STARTING", - ready ? "实时预览播放中" : "实时预览启动中"); + return HikNvrLiveStreamResponse.builder() + .cameraCode(targetCameraCode) + .streamId(streamId) + .status("STOPPED") + .message("未在播放") + .build(); } - /** 应用关闭时停止所有实时拉流进程。 */ + /** 应用关闭时停止队列和所有实时拉流进程。 */ @PreDestroy - public synchronized void destroy() { - // 应用关闭时遍历并停止所有实时拉流会话。 + public void destroy() { + // 先中断每台 NVR 的排队任务,防止应用关闭期间继续创建 FFmpeg。 + for (NvrStartLane lane : new ArrayList<>(startLanes.values())) { + lane.executor.shutdownNow(); + } + startLanes.clear(); + // 再遍历并停止已启动进程,确保容器退出后没有残留 RTSP 会话。 for (String sessionKey : new ArrayList<>(sessions.keySet())) { stopInternal(sessionKey); } + pendingStarts.clear(); } /** - * 内部停止指定 cameraCode 并销毁对应会话对象。 + * 移除并销毁指定本地会话。 * * @param sessionKey 规范化 cameraCode 会话键 */ private void stopInternal(String sessionKey) { - // 先移除会话,避免停止过程中被其他请求再次复用。 + // 先从会话表移除,避免销毁过程中被其他请求复用。 LiveSession session = sessions.remove(sessionKey); - // 没有会话时保持幂等返回。 + // 会话已经被并发停止时保持幂等。 if (session == null) { return; } - // 关闭对应 FFmpeg 进程。 + // 尽快结束 FFmpeg,关闭其持有的 NVR RTSP socket。 destroyProcess(session.process); log.info("[live] stopped cameraCode={} streamId={}", session.cameraCode, session.streamId); } - /** 停止并等待 FFmpeg 进程退出,超时后强制结束。 */ + /** + * 停止 FFmpeg;短暂等待正常退出,随后强制结束以满足快速释放要求。 + * + * @param process FFmpeg 进程 + */ private void destroyProcess(Process process) { - // 进程为空时无需执行销毁操作。 - if (process == null) { + // 空进程或已退出进程没有可释放的系统资源。 + if (process == null || !process.isAlive()) { return; } try { - // 先请求进程正常退出,最多等待两秒。 + // 先发送正常终止信号,让 FFmpeg 有机会关闭 RTSP 和日志句柄。 process.destroy(); - // 超时仍未退出时强制结束 FFmpeg。 - if (!process.waitFor(2, java.util.concurrent.TimeUnit.SECONDS)) { + // 300ms 内未退出就强制结束,不能继续占用 NVR 会话数秒。 + if (!process.waitFor(300L, TimeUnit.MILLISECONDS)) { process.destroyForcibly(); } - } catch (Exception ex) { + } catch (InterruptedException interrupted) { + // 中断场景也必须强制释放 FFmpeg,并恢复当前线程中断状态。 + process.destroyForcibly(); + Thread.currentThread().interrupt(); + } catch (RuntimeException failure) { + // 进程状态竞争或平台异常时强制结束作为最后释放手段。 process.destroyForcibly(); } } - /** 在限定时间内轮询 HLS 地址是否已经生成。 */ - private boolean waitUntilReady(String hlsUrl, long timeoutMillis) { - long deadline = System.currentTimeMillis() + timeoutMillis; + /** + * 在本次尝试时限内轮询 ZLM API,并同时监控 FFmpeg 是否退出或任务是否取消。 + * + * @param session 正在启动的本地会话 + * @param pending 当前排队任务 + * @param deadline 本次尝试截止时间戳 + * @return 已注册、未注册或 API 不可用 + */ + private ZlmProbeResult waitUntilRegistered( + LiveSession session, + PendingStart pending, + long deadline) { + ZlmProbeResult lastProbe = ZlmProbeResult.ABSENT; + // 在总时限内快速轮询,流注册或进程退出都会提前结束。 while (System.currentTimeMillis() < deadline) { - // 每轮探测 HLS 是否已经返回媒体清单。 - if (isHlsReady(hlsUrl)) { - return true; + // 用户停止或码流切换取消后立即结束探测,由调用方销毁进程。 + if (pending.cancelled) { + return ZlmProbeResult.ABSENT; + } + // 调用 ZLM getMediaList 精确确认 live/streamId 是否已经注册。 + lastProbe = queryZlmStream(session.streamId); + // 已注册时前端可安全播放,立即返回成功。 + if (lastProbe == ZlmProbeResult.REGISTERED) { + return lastProbe; + } + // FFmpeg 提前退出说明 stderr 已有确定错误,无需等满启动超时。 + if (!session.isAlive()) { + return lastProbe; } try { - Thread.sleep(400L); - } catch (InterruptedException ex) { + // 100ms 轮询兼顾快速错误反馈和 ZLM API 压力。 + Thread.sleep(100L); + } catch (InterruptedException interrupted) { + // 中断时恢复标志,让应用关闭或任务取消能够快速生效。 Thread.currentThread().interrupt(); - return false; + return lastProbe; } } - // 超时后再做一次最终探测,减少边界时刻误判。 - return isHlsReady(hlsUrl); + // 截止点再探测一次,减少恰好在边界注册造成的误判。 + return queryZlmStream(session.streamId); } - /** 探测 HLS 地址是否返回可播放内容。 */ - private boolean isHlsReady(String hlsUrl) { + /** + * 短暂等待旧流从 ZLM 注销,供同名流切换主辅码流使用。 + * + * @param streamId ZLMediaKit 流标识 + * @param timeoutMillis 最大等待时间 + */ + private void waitUntilUnregistered(String streamId, long timeoutMillis) { + long deadline = System.currentTimeMillis() + Math.max(0L, timeoutMillis); + // 旧流仍注册时短暂轮询,超时后继续由新 FFmpeg 返回真实发布结果。 + while (System.currentTimeMillis() < deadline) { + // 已注销或 API 暂时不可用时不继续阻塞同 NVR 队列。 + if (queryZlmStream(streamId) != ZlmProbeResult.REGISTERED) { + return; + } + try { + Thread.sleep(50L); + } catch (InterruptedException interrupted) { + // 应用关闭或任务取消时立即退出等待。 + Thread.currentThread().interrupt(); + return; + } + } + } + + /** + * 调用 ZLMediaKit getMediaList 确认同名 RTMP 流是否注册。 + * + * @param streamId ZLMediaKit 流标识 + * @return 已注册、未注册或 API 不可用 + */ + private ZlmProbeResult queryZlmStream(String streamId) { HttpURLConnection connection = null; try { - connection = (HttpURLConnection) new URL(hlsUrl).openConnection(); - connection.setConnectTimeout(1500); - connection.setReadTimeout(1500); + String query = "secret=" + encode(zlmApiSecret) + + "&schema=rtmp&vhost=__defaultVhost__&app=live&stream=" + + encode(streamId); + URL url = new URL(zlmApiBaseUrl + "/index/api/getMediaList?" + query); + connection = (HttpURLConnection) url.openConnection(); + connection.setConnectTimeout(zlmApiTimeoutMillis); + connection.setReadTimeout(zlmApiTimeoutMillis); connection.setRequestMethod("GET"); - // 获取 HLS 清单 HTTP 状态码。 - int code = connection.getResponseCode(); - // 非 200 表示流地址尚未准备好。 - if (code != 200) { - return false; + connection.setUseCaches(false); + // 非 200 表示 ZLM API 本身不可用,不能把流误报为 PLAYING。 + if (connection.getResponseCode() != HttpURLConnection.HTTP_OK) { + return ZlmProbeResult.UNAVAILABLE; } - try (InputStream in = connection.getInputStream()) { - byte[] buf = new byte[256]; - int n = in.read(buf); - // 没有内容时不能认为 HLS 已就绪。 - if (n <= 0) { - return false; + String body; + // 读取 ZLM JSON 响应,原始 API 密钥不会写入日志。 + try (InputStream input = connection.getInputStream()) { + body = readText(input); + } + JSONObject response = JSON.parseObject(body); + // ZLM code 非 0 表示鉴权或 API 调用失败,不等同于流不存在。 + if (response == null || response.getIntValue("code") != 0) { + return ZlmProbeResult.UNAVAILABLE; + } + JSONArray data = response.getJSONArray("data"); + // data 为空表示精确查询的同名流尚未注册。 + if (data == null || data.isEmpty()) { + return ZlmProbeResult.ABSENT; + } + // 即使 ZLM 忽略部分过滤参数,也只接受 live/streamId 的精确匹配。 + for (int index = 0; index < data.size(); index++) { + JSONObject media = data.getJSONObject(index); + // app、stream 和 schema 同时匹配才说明 FFmpeg 的 RTMP 发布已经就绪。 + if (media != null + && "live".equals(media.getString("app")) + && streamId.equals(media.getString("stream")) + && "rtmp".equalsIgnoreCase(media.getString("schema"))) { + return ZlmProbeResult.REGISTERED; } - String body = new String(buf, 0, n, StandardCharsets.UTF_8); - // HLS 播放清单必须包含 #EXTM3U 标记。 - return body.contains("#EXTM3U"); } - } catch (Exception ex) { - return false; + return ZlmProbeResult.ABSENT; + } catch (Exception failure) { + // 连接、超时或 JSON 解析失败都表示 ZLM API 不可确认,禁止返回 PLAYING。 + return ZlmProbeResult.UNAVAILABLE; } finally { - // 无论探测成功与否都关闭 HTTP 连接。 + // 每次短轮询都主动关闭 HTTP 连接,避免探测自身耗尽连接池。 if (connection != null) { connection.disconnect(); } } } + /** + * 读取 HTTP 文本响应。 + * + * @param input HTTP 输入流 + * @return UTF-8 文本 + * @throws IOException 读取失败时抛出 + */ + private String readText(InputStream input) throws IOException { + ByteArrayOutputStream output = new ByteArrayOutputStream(); + byte[] buffer = new byte[1024]; + int count; + // 读取短小 ZLM JSON 响应直到 EOF。 + while ((count = input.read(buffer)) >= 0) { + output.write(buffer, 0, count); + } + return new String(output.toByteArray(), StandardCharsets.UTF_8); + } + + /** + * URL 编码 ZLM API 查询参数。 + * + * @param value 查询参数原值 + * @return UTF-8 编码值 + * @throws IOException 当前 JVM 不支持 UTF-8 时抛出 + */ + private String encode(String value) throws IOException { + // JDK 8 使用字符集名称重载,保证 secret 和流名不会破坏查询字符串。 + return URLEncoder.encode(value == null ? "" : value, StandardCharsets.UTF_8.name()); + } + + /** + * 根据 FFmpeg stderr、ZLM 探测和 SDK 异常链分类实时流失败。 + * + * @param ffmpegLog 当前流的原始 FFmpeg stderr + * @param probe 最后一次 ZLM 探测结果 + * @param cause SDK、进程或内部异常 + * @return 携带稳定字符串错误码的统一异常 + */ + private HikNvrPlaybackException classifyFailure( + String ffmpegLog, + ZlmProbeResult probe, + Throwable cause) { + String logText = ffmpegLog == null + ? "" : ffmpegLog.toLowerCase(Locale.ROOT); + String causeText = throwableText(cause).toLowerCase(Locale.ROOT); + String combined = logText + "\n" + causeText; + + // RTSP SETUP 500、带宽不足或连接数过多都属于 NVR 瞬时会话拒绝。 + if (combined.contains("setup failed: 500") + || combined.contains("setup: 500") + || combined.contains("rtsp/1.0 500") + || combined.contains("453 not enough bandwidth") + || combined.contains("too many connections") + || combined.contains("maximum number of clients")) { + return HikNvrPlaybackException.sdkFailure( + NVR_RTSP_BUSY, + NVR_RTSP_BUSY + ":NVR拒绝RTSP SETUP或实时会话已满", + cause); + } + // RTSP 401/Unauthorized 与 SDK 密码、权限或登录句柄失败统一归为登录失败。 + if (combined.contains("401 unauthorized") + || combined.contains("server returned 401") + || combined.contains("method describe failed: 401") + || isNonNetworkLoginFailure(causeText)) { + return HikNvrPlaybackException.sdkFailure( + NVR_LOGIN_FAILED, + NVR_LOGIN_FAILED + ":NVR登录或RTSP鉴权失败", + cause); + } + // 输出端、RTMP 或 publisher 错误明确指向 ZLMediaKit 发布链路。 + if (combined.contains("already publishing") + || combined.contains("error opening output") + || combined.contains("failed to update header") + || combined.contains("broken pipe") + || (combined.contains("rtmp://") + && (combined.contains("connection refused") + || combined.contains("connection reset") + || combined.contains("timed out")))) { + return HikNvrPlaybackException.mediaFailure( + ZLM_PUSH_FAILED, + ZLM_PUSH_FAILED + ":FFmpeg无法向ZLMediaKit发布实时流", + cause); + } + // 网络类 SDK 错误或 RTSP 输入端连接异常说明 NVR 当前不可达。 + if (isNetworkSdkFailure(causeText) + || combined.contains("no route to host") + || combined.contains("network is unreachable") + || combined.contains("connection timed out") + || combined.contains("operation timed out") + || combined.contains("connection refused") + || combined.contains("could not connect to server")) { + return HikNvrPlaybackException.sdkFailure( + NVR_UNREACHABLE, + NVR_UNREACHABLE + ":无法连接NVR,请检查网络和端口", + cause); + } + // DESCRIBE 404、无视频流或输入数据无效说明目标通道当前没有可用媒体。 + if (combined.contains("404 not found") + || combined.contains("channel offline") + || combined.contains("no video") + || combined.contains("could not find codec parameters") + || combined.contains("invalid data found when processing input")) { + return HikNvrPlaybackException.sdkFailure( + CHANNEL_OFFLINE, + CHANNEL_OFFLINE + ":目标通道离线或没有视频数据", + cause); + } + // ZLM API 不可访问或鉴权失败时无法完成强制就绪确认,应明确返回推流端故障。 + if (probe == ZlmProbeResult.UNAVAILABLE) { + return HikNvrPlaybackException.mediaFailure( + ZLM_PUSH_FAILED, + ZLM_PUSH_FAILED + ":无法通过ZLMediaKit API确认流已注册", + cause); + } + // 没有明确 stderr 且 ZLM 始终无流,按启动时限内没有媒体数据处理; + // 消息带时限秒数并说明已重试,避免误导向“通道离线硬件故障”方向排查。 + return HikNvrPlaybackException.sdkFailure( + CHANNEL_OFFLINE, + CHANNEL_OFFLINE + ":启动时限" + + (startupTimeoutMillis / 1000L) + + "秒内未收到媒体数据(已自动重登并重试一次;" + + "网络闪断后流注册偏慢或通道暂无视频都会出现此提示,请稍后重试)", + cause); + } + + /** 判断异常链中的海康 SDK 错误是否属于网络不可达。 */ + private boolean isNetworkSdkFailure(String causeText) { + Matcher matcher = SDK_ERROR_PATTERN.matcher(causeText); + // 遍历异常链文本中的所有 SDK 错误码,任一网络码都按 NVR 不可达处理。 + while (matcher.find()) { + int code = Integer.parseInt(matcher.group(1)); + // 7-10 分别表示连接、发送、接收失败和接收超时。 + if (code >= 7 && code <= 10) { + return true; + } + } + return false; + } + + /** 判断异常链是否包含非网络类 SDK 登录失败。 */ + private boolean isNonNetworkLoginFailure(String causeText) { + // 没有登录语义时不应把普通媒体错误误分为鉴权失败。 + if (!causeText.contains("登录nvr失败") + && !causeText.contains("登录失败")) { + return false; + } + // 网络错误由 NVR_UNREACHABLE 分支处理,其余登录错误归为登录失败。 + return !isNetworkSdkFailure(causeText); + } + + /** 把异常及 cause 链消息连接成分类文本。 */ + private String throwableText(Throwable cause) { + StringBuilder text = new StringBuilder(); + Throwable current = cause; + // 最多读取八层 cause,避免异常链意外成环或过深。 + for (int depth = 0; current != null && depth < 8; depth++) { + // 只收集消息,不输出堆栈和敏感 RTSP 凭据。 + if (current.getMessage() != null) { + text.append(current.getMessage()).append('\n'); + } + current = current.getCause(); + } + return text.toString(); + } + + /** 读取当前流的 FFmpeg 原始日志,文件本身不做改写。 */ + private String readLog(Path logFile) { + // 日志尚未生成时返回空文本,由 ZLM 探测结果补充分类依据。 + if (logFile == null || !Files.isRegularFile(logFile)) { + return ""; + } + try { + // 启动窗口内 stderr 很小,直接按 UTF-8 读取用于错误分类。 + return new String(Files.readAllBytes(logFile), StandardCharsets.UTF_8); + } catch (IOException failure) { + // 分类辅助日志读取失败不能覆盖真正的媒体故障。 + return ""; + } + } + /** * 调用现有 IVS 实时媒体能力取得 cameraCode 对应的 RTSP 地址。 * @@ -307,13 +907,13 @@ public class HikNvrLiveStreamService { * @return 可交给 FFmpeg 拉取的 RTSP 地址 */ private String resolveLiveRtspUrl(String cameraCode, int streamType) { - // 使用现有 IVS 请求对象保持 /video/rtspurl/v1.0 的媒体参数语义一致。 + // 使用现有 IVS 请求对象保持北向实时媒体参数语义一致。 IvsVideoRtspVo request = IvsVideoRtspVo.realStream(cameraCode); - // 页面选择的码流类型写入 IVS 媒体参数,由 IVS 厂商路由解析实际 NVR 和通道。 + // 页面选择的码流类型交给厂商路由生成 Channels/N01 或 Channels/N02。 request.getMediaURLParam().setStreamType(streamType); - // 调用现有 IVS 门面,不在实时转码服务中复制厂商解析和鉴权逻辑。 + // 调用 IVS 门面,不在实时服务中复制厂商 URL 和鉴权规则。 IvsRtspUrlResponse response = ivsOpenApiService.rtspUrl(request); - // IVS 未返回成功结果或有效 RTSP 时不能启动空 FFmpeg 进程。 + // IVS 没有返回成功 RTSP 时不能启动空 FFmpeg 进程。 if (response == null || response.getResultCode() == null || response.getResultCode() != 0 @@ -334,42 +934,42 @@ public class HikNvrLiveStreamService { private List buildFfmpegCommand(String rtspUrl, String rtmpUrl) { // 参数缺失时无法建立确定的拉流和发布目标,按内部调用错误立即失败。 if (!hasText(rtspUrl) || !hasText(rtmpUrl)) { - throw HikNvrPlaybackException.badRequest("实时预览RTSP或RTMP地址不能为空"); + throw HikNvrPlaybackException.badRequest( + "实时预览RTSP或RTMP地址不能为空"); } List command = new ArrayList<>(); command.add(ffmpegPath); command.add("-hide_banner"); command.add("-loglevel"); - // 只保留真正的 FFmpeg 错误,避免实时流逐包调试信息持续写入日志文件。 + // 保留完整错误但关闭周期性进度,日志可用于稳定分类且不会持续刷盘。 command.add("error"); - // 关闭周期性进度行,限制长时间实时预览产生的无诊断价值日志。 command.add("-nostats"); + // 海康 RTSP 不在 SDP 中给出帧率,ffmpeg 会自愿读满默认 5 秒探测窗口来估帧率; + // 显式压到 3 秒让推流提前约 2 秒开始,否则首次尝试 5 秒预算会在探测完成前杀掉进程。 + command.add("-analyzeduration"); + command.add("3000000"); command.add("-rtsp_transport"); command.add("tcp"); command.add("-i"); command.add(rtspUrl); command.add("-c:v"); - // 大华实时地址使用 /cam/realmonitor,现场主辅码流均为 HEVC,必须转成 FLV/HLS 可播放的 H.264。 - if (rtspUrl.toLowerCase().contains("/cam/realmonitor")) { + // 大华 realmonitor 现场流为 HEVC,需要转 H.264 后才能交给浏览器。 + if (rtspUrl.toLowerCase(Locale.ROOT).contains("/cam/realmonitor")) { command.add("libx264"); command.add("-vf"); - // 大华主码流可能达到 2K/4K;限制到最高 1080p,避免旧浏览器因 H.264 Level 5 黑屏。 command.add("scale=w='min(1920,iw)':h=-2"); command.add("-profile:v"); - // Baseline 不依赖 B 帧,兼容旧版 Chromium、低功耗终端和常见安防 WebView。 command.add("baseline"); command.add("-level:v"); - // 1080p 25fps 使用 Level 4.1,在浏览器软硬件解码器中的覆盖面更稳定。 command.add("4.1"); command.add("-preset"); - // ultrafast 优先降低多画面实时转码 CPU 延迟,画质仍由 NVR 原始码率决定。 command.add("ultrafast"); command.add("-tune"); command.add("zerolatency"); command.add("-pix_fmt"); command.add("yuv420p"); } else { - // 海康现有 H.264 路径继续原码流转发,避免改变已经稳定的多画面资源占用。 + // 海康 H.264 路径继续原码流转发,避免九画面增加转码 CPU。 command.add("copy"); } command.add("-an"); @@ -379,11 +979,32 @@ public class HikNvrLiveStreamService { return command; } - /** 将内部拉流会话转换为接口响应对象。 */ - private HikNvrLiveStreamResponse toResponse(LiveSession session, String status, String message) { - // 响应回传完整 cameraCode,让前端按设备树唯一标识核对实时会话。 + /** + * 创建 FFmpeg 进程并把原始 stderr 追加到每路日志。 + * + * @param command ProcessBuilder 参数列表 + * @param logFile 当前流日志路径 + * @return 已启动进程 + * @throws IOException 进程无法创建时抛出 + */ + Process startFfmpeg(List command, Path logFile) throws IOException { + ProcessBuilder builder = new ProcessBuilder(command); + // FFmpeg 的 stdout/stderr 合并后完整追加,重试不会覆盖第一次失败证据。 + builder.redirectErrorStream(true); + builder.redirectOutput(ProcessBuilder.Redirect.appendTo(logFile.toFile())); + return builder.start(); + } + + /** 将本地会话转换为接口响应。 */ + private HikNvrLiveStreamResponse toResponse( + LiveSession session, + String status, + String message) { + // 返回服务端解析出的 NVR 和通道,但绝不回传账号密码或 RTSP 地址。 return HikNvrLiveStreamResponse.builder() .cameraCode(session.cameraCode) + .nvrIp(session.nvrIp) + .channel(session.channel) .streamType(session.streamType) .streamId(session.streamId) .status(status) @@ -393,111 +1014,205 @@ public class HikNvrLiveStreamService { .build(); } - /** - * 校验并规范化页面提交的完整 cameraCode。 - * - * @param cameraCode 页面设备树提供的摄像机编码 - * @return 去除首尾空白后的完整编码 - */ + /** 将排队任务转换为 STARTING 响应。 */ + private HikNvrLiveStreamResponse toPendingResponse( + PendingStart pending, + String message) { + return HikNvrLiveStreamResponse.builder() + .cameraCode(pending.target.getCameraCode()) + .nvrIp(pending.target.getNvrInfo().getNvrIp()) + .channel(pending.target.getChannel()) + .streamType(pending.streamType) + .streamId(pending.streamId) + .status("STARTING") + .hlsUrl(pending.hlsUrl) + .flvUrl(pending.flvUrl) + .message(message) + .build(); + } + + /** 将已取消的排队任务转换为 STOPPED 响应。 */ + private HikNvrLiveStreamResponse toStoppedResponse( + PendingStart pending, + String message) { + // 取消结果保留目标信息,前端可准确清理对应槽位。 + return HikNvrLiveStreamResponse.builder() + .cameraCode(pending.target.getCameraCode()) + .nvrIp(pending.target.getNvrInfo().getNvrIp()) + .channel(pending.target.getChannel()) + .streamType(pending.streamType) + .streamId(pending.streamId) + .status("STOPPED") + .message(message) + .build(); + } + + /** 将异步起流异常转换为状态接口可返回的 FAILED 响应。 */ + private HikNvrLiveStreamResponse toFailureResponse( + PendingStart pending, + HikNvrPlaybackException failure) { + return HikNvrLiveStreamResponse.builder() + .cameraCode(pending.target.getCameraCode()) + .nvrIp(pending.target.getNvrInfo().getNvrIp()) + .channel(pending.target.getChannel()) + .streamType(pending.streamType) + .streamId(pending.streamId) + .status("FAILED") + .errorCode(failure.getErrorCode()) + .message(failure.getMessage()) + .build(); + } + + /** 将 ZLM 中已有但非本服务创建的同名流转换为复用响应。 */ + private HikNvrLiveStreamResponse toExternalResponse( + PendingStart pending, + String status, + String message) { + return HikNvrLiveStreamResponse.builder() + .cameraCode(pending.target.getCameraCode()) + .nvrIp(pending.target.getNvrInfo().getNvrIp()) + .channel(pending.target.getChannel()) + .streamType(pending.streamType) + .streamId(pending.streamId) + .status(status) + .hlsUrl(pending.hlsUrl) + .flvUrl(pending.flvUrl) + .message(message) + .build(); + } + + /** 校验并规范化页面提交的完整 cameraCode。 */ private String requireCameraCode(String cameraCode) { - // 空 cameraCode 无法调用 IVS 实时接口或定位已有拉流会话。 + // 空 cameraCode 无法解析 NVR、通道或会话键。 if (!hasText(cameraCode)) { throw HikNvrPlaybackException.badRequest("cameraCode不能为空"); } return cameraCode.trim(); } - /** - * 生成不同摄像机间不会冲突的实时会话键。 - * - * @param cameraCode IVS 完整摄像机编码 - * @return 忽略字母大小写的内部会话键 - */ + /** 生成忽略字母大小写的摄像机会话键。 */ private String sessionKey(String cameraCode) { - // cameraCode 去除首尾空格并忽略域编码字母大小写,保证三类操作命中同一键。 - return cameraCode.trim().toLowerCase(); + // 域编码大小写差异不应创建两路同名设备流。 + return cameraCode.trim().toLowerCase(Locale.ROOT); } - /** - * 生成可用于 ZLMediaKit 路径和日志文件名的唯一流标识。 - * - * @param cameraCode IVS 完整摄像机编码 - * @return 仅包含安全字符的摄像机流标识 - */ + /** 生成可用于 ZLMediaKit 路径和日志文件名的唯一流标识。 */ private String buildStreamId(String cameraCode) { - // 将 cameraCode 中的 # 等路径分隔字符统一替换为下划线。 - String safeCameraCode = cameraCode.trim() - .replaceAll("[^A-Za-z0-9_-]", "_"); - // 流名直接包含完整摄像机编码,确保跨 NVR 仍然唯一。 - return "camera-" + safeCameraCode; + // 将 # 等路径分隔字符替换为下划线,保留完整设备与域编码避免跨 NVR 冲突。 + return "camera-" + cameraCode.trim().replaceAll("[^A-Za-z0-9_-]", "_"); + } + + /** 返回当前流的 HLS 浏览器地址。 */ + private String hlsUrl(String streamId) { + return httpBaseUrl + "/live/" + streamId + "/hls.m3u8"; + } + + /** 返回当前流的 HTTP-FLV 浏览器地址。 */ + private String flvUrl(String streamId) { + return httpBaseUrl + "/live/" + streamId + ".live.flv"; } - /** 校验实时转码所需的本地 FFmpeg 配置。 */ + /** 返回当前流保存原始 FFmpeg stderr 的日志文件。 */ + private Path logFile(String streamId) { + return Paths.get("target", "live-" + streamId + ".log"); + } + + /** 校验实时预览所需的本地与 ZLMediaKit 配置。 */ private void ensureConfigured() { - Path ffmpeg = Paths.get(ffmpegPath); - // FFmpeg 文件必须真实存在且为普通文件。 + Path ffmpeg = Paths.get(ffmpegPath == null ? "" : ffmpegPath); + // FFmpeg 文件必须真实存在,避免把本地启动错误误报为 NVR 故障。 if (!Files.isRegularFile(ffmpeg)) { - throw HikNvrPlaybackException.badRequest("未找到FFmpeg:" + ffmpegPath); + throw HikNvrPlaybackException.mediaFailure( + ZLM_PUSH_FAILED, + ZLM_PUSH_FAILED + ":未找到FFmpeg:" + ffmpegPath); + } + // RTMP、浏览器地址、API 地址和 secret 任一缺失都无法完成可靠就绪确认。 + if (!hasText(rtmpBaseUrl) + || !hasText(httpBaseUrl) + || !hasText(zlmApiBaseUrl) + || !hasText(zlmApiSecret)) { + throw HikNvrPlaybackException.mediaFailure( + ZLM_PUSH_FAILED, + ZLM_PUSH_FAILED + ":ZLMediaKit地址或API密钥未配置"); } } - /** - * 判断页面或 IVS 返回文本是否包含有效内容。 - * - * @param value 待检查文本 - * @return 非空且不全为空白时返回 true - */ + /** 判断目标是否使用海康登录会话;大华实时流不调用 HCNetSDK。 */ + private boolean usesHikLogin(NvrInfo nvrInfo) { + String vendor = nvrInfo == null ? null : nvrInfo.getNvrCompany(); + // 厂商为空沿用本服务历史海康语义,显式大华目标才跳过海康重登。 + return vendor == null + || (!vendor.toUpperCase(Locale.ROOT).contains("DAHUA") + && !vendor.contains("大华")); + } + + /** 判断文本是否包含非空白内容。 */ private boolean hasText(String value) { - // 同时处理 null、空串和全空白输入。 + // 同时处理 null、空串和纯空白文本。 return value != null && !value.trim().isEmpty(); } /** 移除 URL 末尾多余的斜杠。 */ private String trimTrailingSlash(String value) { - // 基础地址为空时返回空字符串,避免调用方拼接 null。 + // 基础地址为空时返回空字符串,由配置校验统一给出明确错误。 if (value == null) { return ""; } String trimmed = value.trim(); - // 循环移除末尾斜杠,统一基础 URL 格式。 + // 循环移除所有末尾斜杠,避免拼接 API 和流路径时出现双斜杠。 while (trimmed.endsWith("/")) { trimmed = trimmed.substring(0, trimmed.length() - 1); } return trimmed; } + /** ZLMediaKit 精确流查询的三种结果。 */ + private enum ZlmProbeResult { + /** live/streamId 已经注册,允许返回 PLAYING。 */ + REGISTERED, + /** API 正常但同名流尚未注册。 */ + ABSENT, + /** API 连接、鉴权或响应解析失败。 */ + UNAVAILABLE + } + + /** 本服务拥有的一路 FFmpeg 实时拉流会话。 */ private static final class LiveSession { - /** 当前实时流绑定的 IVS 完整摄像机编码。 */ + /** IVS 完整摄像机编码。 */ private final String cameraCode; - /** 码流类型。 */ + /** 摄像机所属 NVR IP。 */ + private final String nvrIp; + /** 页面逻辑通道号。 */ + private final Integer channel; + /** 主码流 1 或子码流 2。 */ private final int streamType; - /** 推流标识。 */ + /** ZLMediaKit 流标识。 */ private final String streamId; - /** FFmpeg 进程。 */ + /** 持有 RTSP 会话的 FFmpeg 进程。 */ private final Process process; - /** HLS 播放地址。 */ + /** HLS 浏览器地址。 */ private final String hlsUrl; - /** HTTP-FLV 播放地址。 */ + /** HTTP-FLV 浏览器地址。 */ private final String flvUrl; - /** 会话创建时间。 */ - private final long startedAt; - /** 创建并保存一个实时拉流会话。 */ + /** 保存一路本地 FFmpeg 实时流的完整生命周期信息。 */ private LiveSession( String cameraCode, + String nvrIp, + Integer channel, int streamType, String streamId, Process process, String hlsUrl, - String flvUrl, - long startedAt) { + String flvUrl) { this.cameraCode = cameraCode; + this.nvrIp = nvrIp; + this.channel = channel; this.streamType = streamType; this.streamId = streamId; this.process = process; this.hlsUrl = hlsUrl; this.flvUrl = flvUrl; - this.startedAt = startedAt; } /** 判断 FFmpeg 拉流进程是否仍在运行。 */ @@ -505,4 +1220,95 @@ public class HikNvrLiveStreamService { return process != null && process.isAlive(); } } + + /** 一次排队或正在执行的起流任务。 */ + private static final class PendingStart { + /** 已解析的摄像机、NVR 和通道目标。 */ + private final HikNvrCameraTarget target; + /** 主码流 1 或子码流 2。 */ + private final int streamType; + /** ZLMediaKit 流标识。 */ + private final String streamId; + /** HLS 浏览器地址。 */ + private final String hlsUrl; + /** HTTP-FLV 浏览器地址。 */ + private final String flvUrl; + /** 首个请求等待结果,排队请求通过状态接口读取相同结果。 */ + private final CompletableFuture future = + new CompletableFuture<>(); + /** 停止或码流切换设置的取消标志。 */ + private volatile boolean cancelled; + + /** 保存起流目标和预先可返回给前端的流地址。 */ + private PendingStart( + HikNvrCameraTarget target, + int streamType, + String streamId, + String hlsUrl, + String flvUrl) { + this.target = target; + this.streamType = streamType; + this.streamId = streamId; + this.hlsUrl = hlsUrl; + this.flvUrl = flvUrl; + } + } + + /** 每台 NVR 独立的有界单线程起流队列和握手节流状态。 */ + private static final class NvrStartLane { + /** 串行执行该 NVR 起流任务的有界线程池。 */ + private final ThreadPoolExecutor executor; + /** 最近一次 RTSP 起流尝试的开始时间。 */ + private final AtomicLong lastStartAt = new AtomicLong(); + /** 已提交但未完成的任务数,用于判断 HTTP 是否需要等待首个结果。 */ + private int outstanding; + + /** 为一台 NVR 创建单线程、有界等待队列。 */ + private NvrStartLane(String nvrKey, int queueSize) { + this.executor = new ThreadPoolExecutor( + 1, + 1, + 0L, + TimeUnit.MILLISECONDS, + new ArrayBlockingQueue<>(queueSize), + runnable -> { + Thread thread = new Thread( + runnable, + "hik-live-start-" + nvrKey.replaceAll("[^A-Za-z0-9_-]", "_")); + // 守护线程不能阻止 Spring 容器退出,PreDestroy 仍会主动清理。 + thread.setDaemon(true); + return thread; + }, + new ThreadPoolExecutor.AbortPolicy()); + } + + /** + * 等待满足同 NVR 相邻握手间隔,并记录本次尝试开始时间。 + * + * @param deadline 当前起流总截止时间 + * @param spacingMillis 配置的相邻握手间隔 + * @return 截止时间前取得起流时隙时返回 true + */ + private boolean awaitSpacing(long deadline, long spacingMillis) { + long now = System.currentTimeMillis(); + long waitMillis = Math.max( + 0L, lastStartAt.get() + spacingMillis - now); + // 等待时间已经超过总启动时限时不再发起新的 NVR 握手。 + if (now + waitMillis >= deadline) { + return false; + } + // 需要节流时只阻塞该 NVR 的专用线程,不占用 Tomcat 请求线程。 + if (waitMillis > 0L) { + try { + Thread.sleep(waitMillis); + } catch (InterruptedException interrupted) { + // 应用关闭时恢复中断标志并取消本次握手。 + Thread.currentThread().interrupt(); + return false; + } + } + lastStartAt.set(System.currentTimeMillis()); + return true; + } + } } diff --git a/src/main/java/com/inspect/nvr/hik/vo/HikNvrLiveStreamConfigVo.java b/src/main/java/com/inspect/nvr/hik/vo/HikNvrLiveStreamConfigVo.java index 42cc3b3b..9e876d0a 100644 --- a/src/main/java/com/inspect/nvr/hik/vo/HikNvrLiveStreamConfigVo.java +++ b/src/main/java/com/inspect/nvr/hik/vo/HikNvrLiveStreamConfigVo.java @@ -18,7 +18,22 @@ public class HikNvrLiveStreamConfigVo { /** ZLMediaKit HTTP/HLS 基础地址。 */ @Value("${hik.playback.hls.http-base-url:http://127.0.0.1:18081}") private String httpBaseUrl = "http://127.0.0.1:18081"; - /** 实时流启动探测超时时间。 */ - @Value("${hik.playback.live.startup-timeout-millis:20000}") - private long startupTimeoutMillis = 20000L; + /** 服务端访问 ZLMediaKit HTTP API 的基础地址。 */ + @Value("${hik.playback.hls.api-base-url:http://127.0.0.1:18081}") + private String apiBaseUrl = "http://127.0.0.1:18081"; + /** ZLMediaKit API 密钥,仅用于服务端就绪探测。 */ + @Value("${hik.playback.hls.api-secret:}") + private String apiSecret = ""; + /** 实时流启动探测超时时间。现场实测流注册需 2~4 秒,网络闪断后更慢,默认 10 秒。 */ + @Value("${hik.playback.live.startup-timeout-millis:10000}") + private long startupTimeoutMillis = 10000L; + /** 同一 NVR 相邻 RTSP 握手的最小间隔。 */ + @Value("${hik.playback.live.start-spacing-millis:600}") + private long startSpacingMillis = 600L; + /** 每台 NVR 等待起流的最大任务数,不含正在执行的任务。 */ + @Value("${hik.playback.live.max-start-queue-size:8}") + private int maxStartQueueSize = 8; + /** 单次 ZLMediaKit API 连接和读取超时时间。 */ + @Value("${hik.playback.live.zlm-api-timeout-millis:300}") + private int zlmApiTimeoutMillis = 300; } diff --git a/src/main/java/com/inspect/nvr/service/HikLoginService.java b/src/main/java/com/inspect/nvr/service/HikLoginService.java index 5f452711..72bf4c16 100644 --- a/src/main/java/com/inspect/nvr/service/HikLoginService.java +++ b/src/main/java/com/inspect/nvr/service/HikLoginService.java @@ -10,147 +10,320 @@ import com.inspect.nvr.hikVision.utils.jna.HikVisionUtils; import com.inspect.nvr.utils.redis.RedisService; import lombok.extern.slf4j.Slf4j; import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.scheduling.annotation.Scheduled; import org.springframework.stereotype.Service; import javax.annotation.Resource; +import java.util.ArrayList; +import java.util.Map; +import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.TimeUnit; /** - * 海康登录服务统一管理 + * 统一管理海康登录会话,负责按 NVR 复用、失效、健康检查和自动重登。 */ @Slf4j @Service public class HikLoginService { + + /** Redis 中记录 SDK 登出失败事件的有序集合键。 */ private static final String ERROR_LOGOUT_KEY = "hik:error"; + + /** 登出失败时记录原始 SDK 错误,Redis 不可用不能阻断会话释放。 */ @Resource private RedisService redisService; + /** 当前进程共享的海康 SDK 实例。 */ @Autowired private HCNetSDK hcNetSDK; - // 使用 Caffeine 缓存,10分钟未访问自动移除并登出 + /** 保存已登录 NVR 的服务端连接参数,供健康检查失败后自动重登。 */ + private final Map loginTargets = new ConcurrentHashMap<>(); + + /** + * 按 NVR IP 缓存登录句柄;同一个 key 的加载由 Caffeine 原子串行, + * 不同 NVR 的首次登录可以并行。 + */ private final Cache sessionCache = Caffeine.newBuilder() - // 10分钟未被get()就过期 - .expireAfterAccess(10, TimeUnit.MINUTES) - .removalListener((String ip, HikLoginSession session, RemovalCause cause) -> { - if (session != null - && (cause == RemovalCause.EXPIRED || cause == RemovalCause.SIZE)) { - log.info("[海康]会话超时自动登出,ip: {},userID: {}", - ip, session.getUserId()); - doLogout(ip, session.getUserId()); - } - }).build(); + .expireAfterAccess(10, TimeUnit.MINUTES) + .removalListener((String ip, HikLoginSession session, RemovalCause cause) -> { + // 只有缓存自然过期或容量淘汰时由监听器负责 SDK 登出。 + if (session != null + && (cause == RemovalCause.EXPIRED + || cause == RemovalCause.SIZE)) { + loginTargets.remove(ip); + log.info("[海康]会话超时自动登出,ip: {},userID: {}", + ip, session.getUserId()); + // 释放 NVR 端登录连接,避免缓存淘汰后遗留设备会话。 + doLogout(ip, session.getUserId()); + } + }).build(); /** - * 登录 + * 登录目标 NVR 并返回 SDK 用户句柄。 + * + * @param nvrInfo 数据库读取的 NVR 连接信息 + * @return 海康 SDK 用户句柄 */ public int login(NvrInfo nvrInfo) { + // 统一经过会话对象入口,确保通道布局也被缓存。 return loginSession(nvrInfo).getUserId(); } /** - * 登录并返回包含设备通道布局的会话。 + * 登录并返回包含设备通道布局的会话;同一 NVR 的并发首次登录只执行一次。 + * + * @param nvrInfo 数据库读取的 NVR 连接信息 + * @return 可复用的登录会话 */ public HikLoginSession loginSession(NvrInfo nvrInfo) { - String ip = nvrInfo.getNvrIp(); + // 登录参数缺失时无法创建稳定缓存键,也不能安全调用原生 SDK。 + if (nvrInfo == null || !hasText(nvrInfo.getNvrIp())) { + throw new IllegalArgumentException("NVR IP不能为空"); + } + String ip = nvrInfo.getNvrIp().trim(); + // 保存最近一次数据库连接参数,健康检查发现僵死句柄后用它自动重登。 + loginTargets.put(ip, nvrInfo); HikLoginSession existingSession = sessionCache.getIfPresent(ip); + // 缓存命中时直接返回;周期健康检查和起流失败重试负责剔除僵死句柄。 if (existingSession != null) { log.info("[海康]登录命中缓存,ip: {},userID: {}", ip, existingSession.getUserId()); return existingSession; } - // 缓存命中不再经过整个 synchronized 方法;只有真正建立新会话时串行化, - // 避免多个并发 PTZ 请求在缓存命中时互相排队。 - synchronized (this) { - existingSession = sessionCache.getIfPresent(ip); - if (existingSession != null) { - log.info("[海康]登录命中缓存,ip: {},userID: {}", - ip, existingSession.getUserId()); - return existingSession; - } + // Caffeine 对同一 IP 的映射计算只执行一次,避免瞬时并发重复登录设备。 + return sessionCache.get(ip, ignored -> createSession(nvrInfo, ip)); + } - // 执行登录 - HCNetSDK.NET_DVR_USER_LOGIN_INFO m_strLoginInfo = HikVisionUtils.login_V40(nvrInfo.getNvrIp(), nvrInfo.getServerPort().shortValue(), nvrInfo.getAccount(), nvrInfo.getPassword()); - HCNetSDK.NET_DVR_DEVICEINFO_V40 m_strDeviceInfo = new HCNetSDK.NET_DVR_DEVICEINFO_V40(); - int userID = hcNetSDK.NET_DVR_Login_V40(m_strLoginInfo, m_strDeviceInfo); - if (userID < 0) { - int errorCode = hcNetSDK.NET_DVR_GetLastError(); - throw new RuntimeException("登录失败,错误码:" + errorCode); + /** + * 主动注销旧句柄并重新登录一次,供起流失败自愈使用。 + * + * @param nvrInfo 数据库读取的 NVR 连接信息 + * @return 新建立的登录会话 + */ + public HikLoginSession relogin(NvrInfo nvrInfo) { + // 先失效缓存并立即登出设备端旧句柄,避免重新登录仍命中僵死会话。 + logout(nvrInfo == null ? null : nvrInfo.getNvrIp()); + // 使用同一份数据库凭据重新建立会话,成功后重新进入缓存。 + return loginSession(nvrInfo); + } + + /** 定时用轻量工作状态接口检查缓存句柄,失败时自动注销并重登。 */ + @Scheduled(fixedDelayString = "${hik.login.health-check-interval-millis:60000}") + public void healthCheckSessions() { + // 遍历已知登录目标而非仅遍历会话缓存,断网重登失败后仍会在下一轮继续恢复。 + for (Map.Entry entry + : new ArrayList<>(loginTargets.entrySet())) { + String ip = entry.getKey(); + NvrInfo target = entry.getValue(); + HikLoginSession session = sessionCache.getIfPresent(ip); + // 当前存在且健康的句柄继续复用,不制造额外登录连接。 + if (session != null && isSessionHealthy(ip, session)) { + continue; + } + log.warn("[海康]登录会话健康检查失败,准备自动重登,ip: {},userID: {}", + ip, session == null ? null : session.getUserId()); + // 先释放缓存和设备端旧句柄,确保后续登录不会再次命中失效对象。 + logout(ip); + // 没有可用数据库连接参数时只能完成失效,等待下一次业务请求重新提供参数。 + if (target == null) { + continue; } - m_strDeviceInfo.read(); - HCNetSDK.NET_DVR_DEVICEINFO_V30 deviceInfo = m_strDeviceInfo.struDeviceV30; - HikLoginSession session = new HikLoginSession( - userID, - Byte.toUnsignedInt(deviceInfo.byChanNum), - Byte.toUnsignedInt(deviceInfo.byStartChan), - Byte.toUnsignedInt(deviceInfo.byIPChanNum) - + Byte.toUnsignedInt(deviceInfo.byHighDChanNum) * 256, - Byte.toUnsignedInt(deviceInfo.byStartDChan)); - // 放入缓存,自动开始计时10分钟 - sessionCache.put(nvrInfo.getNvrIp(), session); - log.info("[海康]登录成功,ip:{},userID:{}", ip, userID); - return session; + try { + // 自动重登恢复会话,设备网络恢复后无需重启容器。 + loginSession(target); + log.info("[海康]登录会话自动重登成功,ip: {}", ip); + } catch (RuntimeException failure) { + // 重登失败保留告警并等待下一轮健康检查或业务请求重试。 + log.warn("[海康]登录会话自动重登失败,ip: {},原因: {}", + ip, failure.getMessage()); + } + } + } + + /** + * 对指定缓存会话调用海康轻量工作状态接口。 + * + * @param ip NVR IP,仅用于诊断日志 + * @param session 待检查的缓存会话 + * @return SDK 确认句柄有效时返回 true + */ + private boolean isSessionHealthy(String ip, HikLoginSession session) { + // 空会话或负数句柄必定无效,不再调用原生接口。 + if (session == null || session.getUserId() < 0) { + return false; + } + HCNetSDK.NET_DVR_WORKSTATE_V30 workState = + new HCNetSDK.NET_DVR_WORKSTATE_V30(); + // 把结构体内存同步给 JNA,供 SDK 写入设备工作状态。 + workState.write(); + boolean healthy; + try { + // 轻量查询验证登录句柄和 NVR 连接是否仍然有效。 + healthy = hcNetSDK.NET_DVR_GetDVRWorkState_V30( + session.getUserId(), workState); + } catch (RuntimeException sdkFailure) { + // 原生 SDK 抛异常同样视为句柄失效,让定时任务继续执行注销和重登。 + log.warn("[海康]会话健康检查调用异常,ip: {},userID: {},原因: {}", + ip, session.getUserId(), sdkFailure.getMessage()); + return false; } + // 健康检查失败时记录原始 SDK 错误码,便于区分网络与句柄失效。 + if (!healthy) { + log.warn("[海康]会话失效,ip: {},userID: {},错误码: {}", + ip, session.getUserId(), hcNetSDK.NET_DVR_GetLastError()); + return false; + } + return true; } /** - * 记录注销失败信息到 Redis + * 调用海康 SDK 建立新会话并读取设备通道布局。 + * + * @param nvrInfo 数据库读取的 NVR 连接信息 + * @param ip 已规范化的缓存键 + * @return 新登录会话 */ - private void recordLogoutError(String ip, Integer userID, int errorCode) { - JSONObject json = new JSONObject(); - json.put("ip", ip); - json.put("userID", userID); - json.put("errorCode", errorCode); - json.put("time", System.currentTimeMillis()); - redisService.redisTemplate.opsForZSet().add(ERROR_LOGOUT_KEY, json.toJSONString(), System.currentTimeMillis()); + private HikLoginSession createSession(NvrInfo nvrInfo, String ip) { + // 组装海康 V40 登录结构体,账号密码只在服务端内存和 SDK 调用中使用。 + HCNetSDK.NET_DVR_USER_LOGIN_INFO loginInfo = HikVisionUtils.login_V40( + ip, + nvrInfo.getServerPort().shortValue(), + nvrInfo.getAccount(), + nvrInfo.getPassword()); + HCNetSDK.NET_DVR_DEVICEINFO_V40 deviceInfoV40 = + new HCNetSDK.NET_DVR_DEVICEINFO_V40(); + // 发起 SDK 登录并取得后续设备调用使用的用户句柄。 + int userId = hcNetSDK.NET_DVR_Login_V40(loginInfo, deviceInfoV40); + // 负数表示 SDK 登录失败,原始错误码保留在异常消息中供上层分类。 + if (userId < 0) { + int errorCode = hcNetSDK.NET_DVR_GetLastError(); + throw new RuntimeException("登录失败,错误码:" + errorCode); + } + // 从原生内存读取 NVR 返回的模拟与数字通道布局。 + deviceInfoV40.read(); + HCNetSDK.NET_DVR_DEVICEINFO_V30 deviceInfo = deviceInfoV40.struDeviceV30; + HikLoginSession session = new HikLoginSession( + userId, + Byte.toUnsignedInt(deviceInfo.byChanNum), + Byte.toUnsignedInt(deviceInfo.byStartChan), + Byte.toUnsignedInt(deviceInfo.byIPChanNum) + + Byte.toUnsignedInt(deviceInfo.byHighDChanNum) * 256, + Byte.toUnsignedInt(deviceInfo.byStartDChan)); + log.info("[海康]登录成功,ip: {},userID: {}", ip, userId); + return session; } /** - * 登出具体实现 + * 记录注销失败信息到 Redis;Redis 异常不能影响会话释放主流程。 + * + * @param ip NVR IP + * @param userId SDK 用户句柄 + * @param errorCode SDK 原始错误码 */ - public void doLogout(String ip, Integer userID) { - if (userID != null) { - // 调用登出SDK - boolean isLogout = hcNetSDK.NET_DVR_Logout(userID); - if (isLogout) { - log.info("[海康]登出成功,ip: {},userID: {}", ip, userID); - } else { - int errorCode = hcNetSDK.NET_DVR_GetLastError(); - log.error("[海康]登出失败,ip: {},userID: {},错误码: {}", ip, userID, errorCode); - // 登出失败日志记录到Redis中 - recordLogoutError(ip, userID, errorCode); - } + private void recordLogoutError(String ip, Integer userId, int errorCode) { + // 单元测试或降级运行没有 Redis 时只保留应用日志。 + if (redisService == null || redisService.redisTemplate == null) { + return; + } + try { + JSONObject json = new JSONObject(); + json.put("ip", ip); + json.put("userID", userId); + json.put("errorCode", errorCode); + json.put("time", System.currentTimeMillis()); + // 按发生时间写入有序集合,供运维追查设备端未释放句柄。 + redisService.redisTemplate.opsForZSet().add( + ERROR_LOGOUT_KEY, json.toJSONString(), System.currentTimeMillis()); + } catch (RuntimeException failure) { + // Redis 暂时不可用时不反向阻塞 SDK 会话清理。 + log.warn("[海康]记录登出失败事件到Redis失败,ip: {},原因: {}", + ip, failure.getMessage()); } } /** - * 登出 + * 立即调用 SDK 登出指定用户句柄。 + * + * @param ip NVR IP + * @param userId SDK 用户句柄 */ - public synchronized void logout(String ip) { - HikLoginSession session = sessionCache.getIfPresent(ip); - if (session != null) { - sessionCache.invalidate(ip); + public void doLogout(String ip, Integer userId) { + // 空句柄没有可释放的设备会话。 + if (userId == null) { + return; + } + boolean logoutSucceeded; + try { + // 调用海康 SDK 断开设备端登录连接。 + logoutSucceeded = hcNetSDK.NET_DVR_Logout(userId); + } catch (RuntimeException sdkFailure) { + // 原生注销异常不能阻断缓存失效或后续重新登录。 + log.error("[海康]登出调用异常,ip: {},userID: {},原因: {}", + ip, userId, sdkFailure.getMessage()); + return; + } + // SDK 确认成功时记录完整会话生命周期。 + if (logoutSucceeded) { + log.info("[海康]登出成功,ip: {},userID: {}", ip, userId); + return; } + int errorCode = hcNetSDK.NET_DVR_GetLastError(); + log.error("[海康]登出失败,ip: {},userID: {},错误码: {}", + ip, userId, errorCode); + // 保留设备端登出失败证据,但 Redis 故障不影响当前调用返回。 + recordLogoutError(ip, userId, errorCode); } /** - * 登出所有用户 + * 失效并立即登出指定 NVR 的缓存会话。 + * + * @param ip NVR IP */ + public void logout(String ip) { + // 空 IP 不能定位缓存会话,按幂等操作直接返回。 + if (!hasText(ip)) { + return; + } + String normalizedIp = ip.trim(); + // 原子移除会话并取得旧句柄,避免只删缓存却遗留设备端连接。 + HikLoginSession session = sessionCache.asMap().remove(normalizedIp); + loginTargets.remove(normalizedIp); + // 缓存没有该 NVR 时保持幂等,不调用无效 SDK 句柄。 + if (session == null) { + return; + } + // 显式注销必须立即释放设备端登录会话。 + doLogout(normalizedIp, session.getUserId()); + } + + /** 登出当前进程缓存的所有海康用户。 */ public void logoutAll() { - // 获取所有缓存的IP和userID - sessionCache.asMap().forEach((ip, session) -> { - doLogout(ip, session.getUserId()); - }); - // 清空整个缓存 - sessionCache.invalidateAll(); + // 使用 IP 快照逐个走显式登出,确保每个设备句柄都真正释放。 + for (String ip : new ArrayList<>(sessionCache.asMap().keySet())) { + logout(ip); + } log.info("[海康]所有用户已登出"); } /** - * 检查是否已登录(同时刷新过期时间) + * 检查指定 NVR 是否存在缓存会话,同时刷新访问过期时间。 + * + * @param ip NVR IP + * @return 缓存中存在会话时返回 true */ public boolean isLoggedIn(String ip) { - return sessionCache.getIfPresent(ip) != null; + // 空 IP 不可能命中登录缓存。 + if (!hasText(ip)) { + return false; + } + return sessionCache.getIfPresent(ip.trim()) != null; + } + + /** 判断文本是否包含非空白内容。 */ + private boolean hasText(String value) { + // 同时覆盖 null、空串和纯空白输入。 + return value != null && !value.trim().isEmpty(); } }