| | |
| | | package com.genersoft.iot.vmp.media.zlm; |
| | | |
| | | import com.alibaba.fastjson.JSONArray; |
| | | import com.alibaba.fastjson.JSONObject; |
| | | import com.genersoft.iot.vmp.gb28181.bean.SendRtpItem; |
| | | import com.genersoft.iot.vmp.gb28181.session.SsrcUtil; |
| | | import com.genersoft.iot.vmp.media.zlm.dto.MediaServerItem; |
| | | import org.slf4j.Logger; |
| | | import org.slf4j.LoggerFactory; |
| | | import org.springframework.beans.factory.annotation.Autowired; |
| | | import org.springframework.beans.factory.annotation.Value; |
| | | import org.springframework.stereotype.Component; |
| | | import org.springframework.util.StringUtils; |
| | | |
| | | import java.util.HashMap; |
| | | import java.util.Map; |
| | |
| | | |
| | | private Logger logger = LoggerFactory.getLogger("ZLMRTPServerFactory"); |
| | | |
| | | @Value("${media.rtp.udpPortRange}") |
| | | private String udpPortRange; |
| | | |
| | | @Autowired |
| | | private ZLMRESTfulUtils zlmresTfulUtils; |
| | | |
| | | private int[] udpPortRangeArray = new int[2]; |
| | | private int[] portRangeArray = new int[2]; |
| | | |
| | | private int currentPort = 0; |
| | | public int createRTPServer(MediaServerItem mediaServerItem, String streamId) { |
| | | Map<String, Integer> currentStreams = new HashMap<>(); |
| | | JSONObject listRtpServerJsonResult = zlmresTfulUtils.listRtpServer(mediaServerItem); |
| | | if (listRtpServerJsonResult != null) { |
| | | JSONArray data = listRtpServerJsonResult.getJSONArray("data"); |
| | | if (data != null) { |
| | | for (int i = 0; i < data.size(); i++) { |
| | | JSONObject dataItem = data.getJSONObject(i); |
| | | currentStreams.put(dataItem.getString("stream_id"), dataItem.getInteger("port")); |
| | | } |
| | | } |
| | | } |
| | | // 已经在推流 |
| | | if (currentStreams.get(streamId) != null) { |
| | | Map<String, Object> closeRtpServerParam = new HashMap<>(); |
| | | closeRtpServerParam.put("stream_id", streamId); |
| | | zlmresTfulUtils.closeRtpServer(mediaServerItem, closeRtpServerParam); |
| | | currentStreams.remove(streamId); |
| | | } |
| | | |
| | | public int createRTPServer(String streamId) { |
| | | Map<String, Object> param = new HashMap<>(); |
| | | int result = -1; |
| | | int newPort = getPortFromUdpPortRange(); |
| | | param.put("port", newPort); |
| | | /** |
| | | * 不设置推流端口端则使用随机端口 |
| | | */ |
| | | if (StringUtils.isEmpty(mediaServerItem.getSendRtpPortRange())){ |
| | | param.put("port", 0); |
| | | }else { |
| | | int newPort = getPortFromportRange(mediaServerItem); |
| | | param.put("port", newPort); |
| | | } |
| | | param.put("enable_tcp", 1); |
| | | param.put("stream_id", streamId); |
| | | JSONObject jsonObject = zlmresTfulUtils.openRtpServer(param); |
| | | System.out.println(jsonObject); |
| | | JSONObject openRtpServerResultJson = zlmresTfulUtils.openRtpServer(mediaServerItem, param); |
| | | |
| | | if (jsonObject != null) { |
| | | switch (jsonObject.getInteger("code")){ |
| | | if (openRtpServerResultJson != null) { |
| | | switch (openRtpServerResultJson.getInteger("code")){ |
| | | case 0: |
| | | result= newPort; |
| | | result= openRtpServerResultJson.getInteger("port"); |
| | | break; |
| | | case -300: // id已经存在 |
| | | result = newPort; |
| | | case -300: // id已经存在, 可能已经在其他端口推流 |
| | | Map<String, Object> closeRtpServerParam = new HashMap<>(); |
| | | closeRtpServerParam.put("stream_id", streamId); |
| | | zlmresTfulUtils.closeRtpServer(mediaServerItem, closeRtpServerParam); |
| | | result = createRTPServer(mediaServerItem, streamId);; |
| | | break; |
| | | case -400: // 端口占用 |
| | | result= createRTPServer(streamId); |
| | | result= createRTPServer(mediaServerItem, streamId); |
| | | break; |
| | | default: |
| | | logger.error("创建RTP Server 失败: " + jsonObject.getString("msg")); |
| | | logger.error("创建RTP Server 失败 {}: " + openRtpServerResultJson.getString("msg"), param.get("port")); |
| | | break; |
| | | } |
| | | }else { |
| | | // 检查ZLM状态 |
| | | logger.error("创建RTP Server 失败: 请检查ZLM服务"); |
| | | logger.error("创建RTP Server 失败 {}: 请检查ZLM服务", param.get("port")); |
| | | } |
| | | return result; |
| | | } |
| | | |
| | | public boolean closeRTPServer(String streamId) { |
| | | public boolean closeRTPServer(MediaServerItem serverItem, String streamId) { |
| | | boolean result = false; |
| | | Map<String, Object> param = new HashMap<>(); |
| | | param.put("stream_id", streamId); |
| | | JSONObject jsonObject = zlmresTfulUtils.closeRtpServer(param); |
| | | if (jsonObject != null ) { |
| | | if (jsonObject.getInteger("code") == 0) { |
| | | result = jsonObject.getInteger("hit") == 1; |
| | | if (serverItem !=null){ |
| | | Map<String, Object> param = new HashMap<>(); |
| | | param.put("stream_id", streamId); |
| | | JSONObject jsonObject = zlmresTfulUtils.closeRtpServer(serverItem, param); |
| | | if (jsonObject != null ) { |
| | | if (jsonObject.getInteger("code") == 0) { |
| | | result = jsonObject.getInteger("hit") == 1; |
| | | }else { |
| | | logger.error("关闭RTP Server 失败: " + jsonObject.getString("msg")); |
| | | } |
| | | }else { |
| | | logger.error("关闭RTP Server 失败: " + jsonObject.getString("msg")); |
| | | // 检查ZLM状态 |
| | | logger.error("关闭RTP Server 失败: 请检查ZLM服务"); |
| | | } |
| | | }else { |
| | | // 检查ZLM状态 |
| | | logger.error("关闭RTP Server 失败: 请检查ZLM服务"); |
| | | } |
| | | return result; |
| | | } |
| | | |
| | | private int getPortFromUdpPortRange() { |
| | | private int getPortFromportRange(MediaServerItem mediaServerItem) { |
| | | int currentPort = mediaServerItem.getCurrentPort(); |
| | | if (currentPort == 0) { |
| | | String[] udpPortRangeStrArray = udpPortRange.split(","); |
| | | udpPortRangeArray[0] = Integer.parseInt(udpPortRangeStrArray[0]); |
| | | udpPortRangeArray[1] = Integer.parseInt(udpPortRangeStrArray[1]); |
| | | String[] portRangeStrArray = mediaServerItem.getSendRtpPortRange().split(","); |
| | | if (portRangeStrArray.length != 2) { |
| | | portRangeArray[0] = 30000; |
| | | portRangeArray[1] = 30500; |
| | | }else { |
| | | portRangeArray[0] = Integer.parseInt(portRangeStrArray[0]); |
| | | portRangeArray[1] = Integer.parseInt(portRangeStrArray[1]); |
| | | } |
| | | } |
| | | |
| | | if (currentPort == 0 || currentPort++ > udpPortRangeArray[1]) { |
| | | currentPort = udpPortRangeArray[0]; |
| | | return udpPortRangeArray[0]; |
| | | if (currentPort == 0 || currentPort++ > portRangeArray[1]) { |
| | | currentPort = portRangeArray[0]; |
| | | mediaServerItem.setCurrentPort(currentPort); |
| | | return portRangeArray[0]; |
| | | } else { |
| | | if (currentPort % 2 == 1) { |
| | | currentPort++; |
| | | } |
| | | return currentPort++; |
| | | currentPort++; |
| | | mediaServerItem.setCurrentPort(currentPort); |
| | | return currentPort; |
| | | } |
| | | } |
| | | |
| | | /** |
| | | * 创建一个推流 |
| | | * 创建一个国标推流 |
| | | * @param ip 推流ip |
| | | * @param port 推流端口 |
| | | * @param ssrc 推流唯一标识 |
| | |
| | | * @param tcp 是否为tcp |
| | | * @return SendRtpItem |
| | | */ |
| | | public SendRtpItem createSendRtpItem(String ip, int port, String ssrc, String platformId, String deviceId, String channelId, boolean tcp){ |
| | | String playSsrc = SsrcUtil.getPlaySsrc(); |
| | | int localPort = createRTPServer(SsrcUtil.getPlaySsrc()); |
| | | public SendRtpItem createSendRtpItem(MediaServerItem serverItem, String ip, int port, String ssrc, String platformId, String deviceId, String channelId, boolean tcp){ |
| | | |
| | | // 使用RTPServer 功能找一个可用的端口 |
| | | String playSsrc = serverItem.getSsrcConfig().getPlaySsrc(); |
| | | int localPort = createRTPServer(serverItem, playSsrc); |
| | | if (localPort != -1) { |
| | | closeRTPServer(playSsrc); |
| | | // TODO 高并发时可能因为未放入缓存而ssrc冲突 |
| | | serverItem.getSsrcConfig().releaseSsrc(playSsrc); |
| | | closeRTPServer(serverItem, playSsrc); |
| | | }else { |
| | | logger.error("没有可用的端口"); |
| | | return null; |
| | |
| | | sendRtpItem.setDeviceId(deviceId); |
| | | sendRtpItem.setChannelId(channelId); |
| | | sendRtpItem.setTcp(tcp); |
| | | sendRtpItem.setApp("rtp"); |
| | | sendRtpItem.setLocalPort(localPort); |
| | | sendRtpItem.setMediaServerId(serverItem.getId()); |
| | | return sendRtpItem; |
| | | } |
| | | |
| | | /** |
| | | * |
| | | * 创建一个直播推流 |
| | | * @param ip 推流ip |
| | | * @param port 推流端口 |
| | | * @param ssrc 推流唯一标识 |
| | | * @param platformId 平台id |
| | | * @param channelId 通道id |
| | | * @param tcp 是否为tcp |
| | | * @return SendRtpItem |
| | | */ |
| | | public Boolean startSendRtpStream(Map<String, Object>param) { |
| | | Boolean result = false; |
| | | JSONObject jsonObject = zlmresTfulUtils.startSendRtp(param); |
| | | System.out.println(jsonObject); |
| | | if (jsonObject != null) { |
| | | switch (jsonObject.getInteger("code")){ |
| | | case 0: |
| | | result= true; |
| | | logger.error("RTP推流请求成功,本地推流端口:" + jsonObject.getString("local_port")); |
| | | break; |
| | | // case -300: // id已经存在 |
| | | // result = false; |
| | | // break; |
| | | // case -400: // 端口占用 |
| | | // result= false; |
| | | // break; |
| | | default: |
| | | logger.error("RTP推流失败: " + jsonObject.getString("msg")); |
| | | break; |
| | | } |
| | | public SendRtpItem createSendRtpItem(MediaServerItem serverItem, String ip, int port, String ssrc, String platformId, String app, String stream, String channelId, boolean tcp){ |
| | | String playSsrc = serverItem.getSsrcConfig().getPlaySsrc(); |
| | | int localPort = createRTPServer(serverItem, playSsrc); |
| | | if (localPort != -1) { |
| | | // TODO 高并发时可能因为未放入缓存而ssrc冲突 |
| | | serverItem.getSsrcConfig().releaseSsrc(ssrc); |
| | | closeRTPServer(serverItem, playSsrc); |
| | | }else { |
| | | // 检查ZLM状态 |
| | | logger.error("没有可用的端口"); |
| | | return null; |
| | | } |
| | | SendRtpItem sendRtpItem = new SendRtpItem(); |
| | | sendRtpItem.setIp(ip); |
| | | sendRtpItem.setPort(port); |
| | | sendRtpItem.setSsrc(ssrc); |
| | | sendRtpItem.setApp(app); |
| | | sendRtpItem.setStreamId(stream); |
| | | sendRtpItem.setPlatformId(platformId); |
| | | sendRtpItem.setChannelId(channelId); |
| | | sendRtpItem.setTcp(tcp); |
| | | sendRtpItem.setLocalPort(localPort); |
| | | sendRtpItem.setMediaServerId(serverItem.getId()); |
| | | return sendRtpItem; |
| | | } |
| | | |
| | | /** |
| | | * 调用zlm RESTful API —— startSendRtp |
| | | */ |
| | | public Boolean startSendRtpStream(MediaServerItem mediaServerItem, Map<String, Object>param) { |
| | | Boolean result = false; |
| | | JSONObject jsonObject = zlmresTfulUtils.startSendRtp(mediaServerItem, param); |
| | | if (jsonObject == null) { |
| | | logger.error("RTP推流失败: 请检查ZLM服务"); |
| | | } else if (jsonObject.getInteger("code") == 0) { |
| | | result= true; |
| | | logger.info("RTP推流[ {}/{} ]请求成功,本地推流端口:{}" ,param.get("app"), param.get("stream"), jsonObject.getString("local_port")); |
| | | } else { |
| | | logger.error("RTP推流失败: " + jsonObject.getString("msg")); |
| | | } |
| | | return result; |
| | | } |
| | | |
| | | /** |
| | | * |
| | | * 查询待转推的流是否就绪 |
| | | */ |
| | | public Boolean isRtpReady(String streamId) { |
| | | JSONObject mediaInfo = zlmresTfulUtils.getMediaInfo("rtp", "rtmp", streamId); |
| | | if (mediaInfo.getInteger("code") == 0 && mediaInfo.getBoolean("online")) { |
| | | logger.info("设备RTP推流成功"); |
| | | return true; |
| | | } else { |
| | | logger.info("设备RTP推流未完成"); |
| | | return false; |
| | | public Boolean isRtpReady(MediaServerItem mediaServerItem, String streamId) { |
| | | JSONObject mediaInfo = zlmresTfulUtils.getMediaInfo(mediaServerItem,"rtp", "rtmp", streamId); |
| | | return (mediaInfo.getInteger("code") == 0 && mediaInfo.getBoolean("online")); |
| | | } |
| | | |
| | | /** |
| | | * 查询待转推的流是否就绪 |
| | | */ |
| | | public Boolean isStreamReady(MediaServerItem mediaServerItem, String app, String streamId) { |
| | | JSONObject mediaInfo = zlmresTfulUtils.getMediaInfo(mediaServerItem, app, "rtmp", streamId); |
| | | return (mediaInfo.getInteger("code") == 0 && mediaInfo.getBoolean("online")); |
| | | } |
| | | |
| | | /** |
| | | * 查询转推的流是否有其它观看者 |
| | | * @param streamId |
| | | * @return |
| | | */ |
| | | public int totalReaderCount(MediaServerItem mediaServerItem, String app, String streamId) { |
| | | JSONObject mediaInfo = zlmresTfulUtils.getMediaInfo(mediaServerItem, app, "rtmp", streamId); |
| | | if (mediaInfo == null) { |
| | | return 0; |
| | | } |
| | | return mediaInfo.getInteger("totalReaderCount"); |
| | | } |
| | | |
| | | /** |
| | | * 调用zlm RESTful API —— stopSendRtp |
| | | */ |
| | | public Boolean stopSendRtpStream(MediaServerItem mediaServerItem, Map<String, Object>param) { |
| | | Boolean result = false; |
| | | JSONObject jsonObject = zlmresTfulUtils.stopSendRtp(mediaServerItem, param); |
| | | if (jsonObject == null) { |
| | | logger.error("停止RTP推流失败: 请检查ZLM服务"); |
| | | } else if (jsonObject.getInteger("code") == 0) { |
| | | result= true; |
| | | logger.info("停止RTP推流成功"); |
| | | } else { |
| | | logger.error("停止RTP推流失败: " + jsonObject.getString("msg")); |
| | | } |
| | | return result; |
| | | } |
| | | |
| | | public void closeAllSendRtpStream() { |
| | | |
| | | } |
| | | } |