liusuyi
2026-06-01 a2e7e8ff9cfaa69b001d483710bddbda50d55c91
ard-modules/ard-modules-zlm/src/main/java/com/ard/zlm/service/impl/MediaServerServiceImpl.java
@@ -11,6 +11,7 @@
import com.ard.common.core.utils.file.FileMultipartFile;
import com.ard.gb28181.api.RemoteGb28181Service;
import com.ard.gb28181.api.domain.Device;
import com.ard.gb28181.api.domain.GbChannelDTO;
import com.ard.qs.api.RemoteQsDeviceService;
import com.ard.qs.api.domain.QsDevice;
import com.ard.system.api.RemoteFileService;
@@ -1811,6 +1812,101 @@
        }
    }
    /**
     * gb28181 播放(基于GbDevice,不依赖QS)
     */
    @Override
    public void startGb28181PlayByGbChannel(GbChannelDTO gbChannelDTO, Device gbDevice, ErrorCallback<StreamInfo> callback) {
        ZlmMediaServer mediaServer = getMediaServerForMinimumLoad(null);
        if (mediaServer == null) {
            callback.run(InviteErrorCode.FAIL.getCode(), "无可用的节点", null);
            return;
        }
        String streamMode = gbChannelDTO.getStreamMode() != null
                ? gbChannelDTO.getStreamMode() : gbDevice.getStreamMode();
        int tcpMode = streamMode.equals("TCP-ACTIVE") ? 2
                : (streamMode.equals("TCP-PASSIVE") ? 1 : 0);
        RTPServerParam rtpServerParam = new RTPServerParam();
        rtpServerParam.setApp("gb28181");
        rtpServerParam.setMediaServer(mediaServer);
        rtpServerParam.setType(LiveStreamType.GB28181.getCode());
        rtpServerParam.setStreamId(gbChannelDTO.getDeviceCode());
        rtpServerParam.setTcpMode(tcpMode);
        rtpServerParam.setId(gbChannelDTO.getId());
        startGb28181PlayFunByGbChannel(mediaServer, gbChannelDTO, gbDevice, rtpServerParam, null, callback);
    }
    /**
     * gb28181 停止点播(基于GbDevice,不依赖QS)
     */
    @Override
    public void stopGb28181PlayByGbChannel(InviteSessionType type, GbChannelDTO gbChannelDTO, Device device, String stream) {
        InviteInfo inviteInfo = inviteStreamService.getInviteInfo(type, gbChannelDTO.getId(), stream);
        if (inviteInfo == null) {
            if (type == InviteSessionType.PLAY) {
                GbChannelDTO update = new GbChannelDTO();
                update.setId(gbChannelDTO.getId());
                update.setStreamKey("");
                update.setMediaServerId("");
                update.setStreamStatus("0");
                R<Boolean> r = remoteGb28181Service.updateGbChannelStream(update, SecurityConstants.INNER);
                if (r.getCode() != Constants.SUCCESS) {
                    throw new RuntimeException("更新GbDevice失败");
                }
            }
            return;
        }
        inviteStreamService.removeInviteInfo(inviteInfo);
        if (InviteSessionStatus.ok == inviteInfo.getStatus()) {
            try {
                log.info("[停止点播/回放/下载] {}/{}", gbChannelDTO.getGbDeviceId(), gbChannelDTO.getGbChannelId());
                RtpServerParam rtpServer = new RtpServerParam();
                rtpServer.setApp("gb28181");
                rtpServer.setStream(gbChannelDTO.getDeviceCode());
                rtpServer.setGbDeviceId(gbChannelDTO.getGbDeviceId());
                rtpServer.setGbChannelId(gbChannelDTO.getGbChannelId());
                R<Void> r = remoteGb28181Service.streamByeCmd(rtpServer, SecurityConstants.INNER);
                if (r.getCode() != Constants.SUCCESS) {
                    log.error("[命令发送失败] 停止点播/回放/下载, deviceId:{}", gbChannelDTO.getGbDeviceId());
                    throw new RuntimeException("[命令发送失败] 停止点播/回放/下载, deviceId:" + gbChannelDTO.getGbDeviceId());
                }
            } catch (Exception e) {
                log.error("[命令发送失败] 停止点播/回放/下载, 发送BYE: {}", e.getMessage());
                throw new RuntimeException("命令发送失败: " + e.getMessage());
            }
        }
        if (inviteInfo.getType() == InviteSessionType.PLAY) {
            GbChannelDTO update = new GbChannelDTO();
            update.setId(gbChannelDTO.getId());
            update.setStreamKey("");
            update.setMediaServerId("");
            update.setStreamStatus("0");
            R<Boolean> r = remoteGb28181Service.updateGbChannelStream(update, SecurityConstants.INNER);
            if (r.getCode() != Constants.SUCCESS) {
                throw new RuntimeException("更新GbDevice失败");
            }
        }
        ZlmMediaServer mediaServer = null;
        if (inviteInfo.getStreamInfo() != null) {
            mediaServer = inviteInfo.getStreamInfo().getMediaServer();
        } else {
            mediaServer = getOne(inviteInfo.getMediaServerId());
        }
        if (mediaServer != null && inviteInfo.getSsrcInfo() != null) {
            closeRTPServer(mediaServer, inviteInfo.getSsrcInfo().getStream());
            ssrcFactory.releaseSsrc(inviteInfo.getMediaServerId(), inviteInfo.getSsrcInfo().getSsrc());
        }
    }
    /**
     * 开启国标28181播放
@@ -2029,6 +2125,197 @@
    }
    /**
     * 开启国标28181播放(基于GbDevice,不依赖QS)
     */
    private SSRCInfo startGb28181PlayFunByGbChannel(ZlmMediaServer mediaServer, GbChannelDTO gbChannelDTO,
                                                    Device gbDevice, RTPServerParam rtpServerParam,
                                                    String ssrc, ErrorCallback<StreamInfo> callback) {
        // 获取点播的状态信息
        InviteInfo inviteInfoInCatch = inviteStreamService.getInviteInfoByDeviceAndChannel(InviteSessionType.PLAY,
                gbChannelDTO.getId());
        if (inviteInfoInCatch != null) {
            if (inviteInfoInCatch.getStreamInfo() == null) {
                ssrcFactory.releaseSsrc(mediaServer.getId(), null);
                inviteStreamService.once(InviteSessionType.PLAY, gbChannelDTO.getId(), null, callback);
                log.info("[点播开始] 已经请求中,等待结果, deviceId: {}, channel: {}", gbChannelDTO.getId(), gbChannelDTO.getId());
                return inviteInfoInCatch.getSsrcInfo();
            } else {
                StreamInfo streamInfo = inviteInfoInCatch.getStreamInfo();
                String streamId = streamInfo.getStream();
                if (streamId == null) {
                    callback.run(InviteErrorCode.ERROR_FOR_CATCH_DATA.getCode(), "点播失败, redis缓存streamId等于null", null);
                    inviteStreamService.call(InviteSessionType.PLAY, gbChannelDTO.getId(), null,
                            InviteErrorCode.ERROR_FOR_CATCH_DATA.getCode(), "点播失败, redis缓存streamId等于null", null);
                    return inviteInfoInCatch.getSsrcInfo();
                }
                ZlmMediaServer mediaInfo = streamInfo.getMediaServer();
                Boolean ready = isStreamReady(mediaInfo, rtpServerParam.getApp(), streamId);
                if (ready != null && ready) {
                    if (callback != null) {
                        callback.run(InviteErrorCode.SUCCESS.getCode(), InviteErrorCode.SUCCESS.getMsg(), streamInfo);
                    }
                    inviteStreamService.call(InviteSessionType.PLAY, gbChannelDTO.getId(), null,
                            InviteErrorCode.SUCCESS.getCode(), InviteErrorCode.SUCCESS.getMsg(), streamInfo);
                    log.info("[点播已存在] 直接返回, 设备编号: {}", gbChannelDTO.getId());
                    return inviteInfoInCatch.getSsrcInfo();
                } else {
                    inviteStreamService.once(InviteSessionType.PLAY, gbChannelDTO.getId(), null, callback);
                    RTPServerParam stopRtp = new RTPServerParam();
                    stopRtp.setId(gbChannelDTO.getId());
                    stopRtp.setType(rtpServerParam.getType());
                    stopRtp.setStreamId(rtpServerParam.getStreamId());
                    stopRtpPlay(stopRtp);
                    inviteStreamService.removeInviteInfoByDeviceAndChannel(InviteSessionType.PLAY, gbChannelDTO.getId());
                }
            }
        }
        rtpServerParam.setMediaServer(mediaServer);
        if (rtpServerParam.getPresetSsrc() != null) {
            ssrc = rtpServerParam.getPresetSsrc();
        } else {
            if (rtpServerParam.isPlayback()) {
                ssrc = ssrcFactory.getPlayBackSsrc(mediaServer.getId());
            } else {
                ssrc = ssrcFactory.getPlaySsrc(mediaServer.getId());
            }
        }
        rtpServerParam.setSsrc(ssrc);
        SSRCInfo ssrcInfo = receiveRtpServerService.openRTPServer(rtpServerParam, (code, msg, result) -> {
            if (code == InviteErrorCode.SUCCESS.getCode() && result != null && result.getHookData() != null) {
                log.info("[创建RTP服务器] 成功, code: {}, msg: {}, result: {}", code, msg, result);
                StreamInfo streamInfo = getStreamInfoByAppAndStream(mediaServer, rtpServerParam.getApp(),
                        rtpServerParam.getStreamId(), result.getHookData().getMediaInfo());
                if (streamInfo == null) {
                    if (callback != null) {
                        callback.run(InviteErrorCode.ERROR_FOR_STREAM_PARSING_EXCEPTIONS.getCode(),
                                InviteErrorCode.ERROR_FOR_STREAM_PARSING_EXCEPTIONS.getMsg(), null);
                    }
                    inviteStreamService.call(InviteSessionType.PLAY, gbChannelDTO.getId(), null,
                            InviteErrorCode.ERROR_FOR_STREAM_PARSING_EXCEPTIONS.getCode(),
                            InviteErrorCode.ERROR_FOR_STREAM_PARSING_EXCEPTIONS.getMsg(), null);
                    if (result != null && result.getSsrcInfo() != null) {
                        closeRTPServer(mediaServer, result.getSsrcInfo().getStream());
                        ssrcFactory.releaseSsrc(mediaServer.getId(), result.getSsrcInfo().getSsrc());
                    }
                    return;
                }
                if (callback != null) {
                    callback.run(InviteErrorCode.SUCCESS.getCode(), InviteErrorCode.SUCCESS.getMsg(), streamInfo);
                    inviteStreamService.call(InviteSessionType.PLAY, gbChannelDTO.getId(), null,
                            InviteErrorCode.SUCCESS.getCode(), InviteErrorCode.SUCCESS.getMsg(), streamInfo);
                    InviteInfo inviteInfo = inviteStreamService.getInviteInfoByDeviceAndChannel(
                            InviteSessionType.PLAY, gbChannelDTO.getId());
                    if (inviteInfo != null) {
                        inviteInfo.setStatus(InviteSessionStatus.ok);
                        inviteInfo.setStreamInfo(streamInfo);
                        inviteStreamService.updateInviteInfo(inviteInfo);
                    }
                    String filePath = snapOnPlay(streamInfo.getMediaServer(), streamInfo.getApp(),
                            streamInfo.getStream());
                    // 更新GbDevice流状态
                    GbChannelDTO update = new GbChannelDTO();
                    update.setId(rtpServerParam.getId());
                    update.setStreamKey(rtpServerParam.getStreamId());
                    update.setMediaServerId(mediaServer.getId());
                    update.setStreamStatus("1");
                    update.setSnap(filePath);
                    R<Boolean> r = remoteGb28181Service.updateGbChannelStream(update, SecurityConstants.INNER);
                    if (r.getCode() != Constants.SUCCESS) {
                        throw new RuntimeException("更新GbDevice失败");
                    }
                }
            } else {
                log.error("[创建RTP服务器] 失败, code: {}, msg: {}, result: {}", code, msg, result);
                if (callback != null) {
                    callback.run(code, msg, null);
                }
                inviteStreamService.call(InviteSessionType.PLAY, gbChannelDTO.getId(), null, code, msg, null);
                inviteStreamService.removeInviteInfoByDeviceAndChannel(InviteSessionType.PLAY, gbChannelDTO.getId());
                if (result != null && result.getSsrcInfo() != null) {
                    closeRTPServer(mediaServer, result.getSsrcInfo().getStream());
                    ssrcFactory.releaseSsrc(mediaServer.getId(), result.getSsrcInfo().getSsrc());
                }
            }
        });
        if (ssrcInfo == null || ssrcInfo.getPort() <= 0) {
            log.info("[点播端口/SSRC]获取失败,设备编号:{}, 通道编号:{}, ssrcInfo: {}", gbChannelDTO.getId(), gbChannelDTO.getId(), ssrcInfo);
            if (rtpServerParam.getPresetSsrc() == null) {
                ssrcFactory.releaseSsrc(mediaServer.getId(), ssrc);
            }
            callback.run(InviteErrorCode.ERROR_FOR_RESOURCE_EXHAUSTION.getCode(), "获取端口或者ssrc失败", null);
            inviteStreamService.call(InviteSessionType.PLAY, gbChannelDTO.getId(), null,
                    InviteErrorCode.ERROR_FOR_RESOURCE_EXHAUSTION.getCode(),
                    InviteErrorCode.ERROR_FOR_RESOURCE_EXHAUSTION.getMsg(), null);
            return null;
        }
        int port = ssrcInfo.getPort();
        String ip = mediaServer.getIp();
        RtpServerParam rtpServer = new RtpServerParam();
        rtpServer.setPort(port);
        rtpServer.setIp(ip);
        rtpServer.setId(rtpServerParam.getId());
        rtpServer.setSsrc(rtpServerParam.getSsrc());
        rtpServer.setGbDeviceId(gbDevice.getDeviceId());
        rtpServer.setGbChannelId(gbChannelDTO.getGbChannelId());
        rtpServer.setStreamMode(gbDevice.getStreamMode());
        rtpServer.setMediaServerId(mediaServer.getId());
        rtpServer.setApp(rtpServerParam.getApp());
        rtpServer.setStream(rtpServerParam.getStreamId());
        log.info("[国标28181点播开始(基于GbDevice)] ===============================");
        log.info("[国标28181] GbDeviceId: {}, 设备国标ID: {}, 通道国标ID: {}", gbChannelDTO.getId(),
                gbDevice.getDeviceId(), gbChannelDTO.getGbChannelId());
        log.info("[国标28181] 流模式: {}, ZLM tcpMode: {}, ssrcCheck: {}", gbDevice.getStreamMode(),
                rtpServerParam.getTcpMode(), rtpServerParam.isSsrcCheck());
        log.info("[国标28181] ZLM媒体服务器IP: {}, 收流端口: {}, 流ID: {}, SSRC: {}", ip, port, ssrcInfo.getStream(),
                ssrcInfo.getSsrc());
        log.info("[国标28181] =======================================");
        InviteInfo inviteInfo = InviteInfo.getInviteInfo(gbChannelDTO.getId().toString(), gbChannelDTO.getId(),
                ssrcInfo.getStream(), ssrcInfo, mediaServer.getId(), mediaServer.getSdpIp(), ssrcInfo.getPort(),
                gbDevice.getStreamMode(), InviteSessionType.PLAY, InviteSessionStatus.ready,
                userSetting.getRecordSip());
        if ("1".equals(gbChannelDTO.getEnableMp4())) {
            inviteInfo.setRecord(true);
        }
        inviteStreamService.updateInviteInfo(inviteInfo);
        R<Void> r = remoteGb28181Service.playStreamCmd(rtpServer, SecurityConstants.INNER);
        if (r.getCode() != Constants.SUCCESS) {
            log.info("[点播失败]{}:{} deviceId: {}, channelId:{}", r.getCode(), r.getMsg(),
                    gbChannelDTO.getGbDeviceId(), gbChannelDTO.getGbChannelId());
            inviteInfo = inviteStreamService.getInviteInfo(InviteSessionType.PLAY, gbChannelDTO.getId(),
                    rtpServerParam.getStreamId());
            if (inviteInfo != null) {
                inviteStreamService.removeInviteInfo(inviteInfo);
                if (inviteInfo.getSsrcInfo() != null) {
                    ssrcFactory.releaseSsrc(mediaServer.getId(), inviteInfo.getSsrcInfo().getSsrc());
                }
            }
            closeRTPServer(mediaServer, ssrcInfo.getStream());
            ssrcFactory.releaseSsrc(mediaServer.getId(), ssrcInfo.getSsrc());
            if (callback != null) {
                callback.run(r.getCode(), r.getMsg(), null);
            }
            inviteStreamService.call(InviteSessionType.PLAY, gbChannelDTO.getId(), null,
                    r.getCode(), r.getMsg(), null);
            inviteStreamService.removeInviteInfoByDeviceAndChannel(InviteSessionType.PLAY, gbChannelDTO.getId());
            return ssrcInfo;
        }
        return ssrcInfo;
    }
    /**
     * 将 WebSocket 协议地址转换为 HTTP 协议地址
     * ws:// -> http://
     * wss:// -> https://