|
|
@@ -31,11 +31,17 @@ import com.genersoft.iot.vmp.service.IInviteStreamService;
|
|
|
import com.genersoft.iot.vmp.service.IPlayService;
|
|
|
import com.genersoft.iot.vmp.service.IStreamProxyService;
|
|
|
import com.genersoft.iot.vmp.service.IStreamPushService;
|
|
|
+import com.genersoft.iot.vmp.media.zlm.SendRtpPortManager;
|
|
|
+import com.genersoft.iot.vmp.service.redisMsg.IRedisRpcService;
|
|
|
+import com.genersoft.iot.vmp.media.zlm.ZLMServerFactory;
|
|
|
+import com.genersoft.iot.vmp.media.zlm.ZlmHttpHookSubscribe;
|
|
|
+import com.genersoft.iot.vmp.media.zlm.dto.*;
|
|
|
+import com.genersoft.iot.vmp.media.zlm.dto.hook.OnStreamChangedHookParam;
|
|
|
+import com.genersoft.iot.vmp.service.*;
|
|
|
import com.genersoft.iot.vmp.service.bean.ErrorCallback;
|
|
|
import com.genersoft.iot.vmp.service.bean.InviteErrorCode;
|
|
|
import com.genersoft.iot.vmp.service.bean.MessageForPushChannel;
|
|
|
import com.genersoft.iot.vmp.service.bean.SSRCInfo;
|
|
|
-import com.genersoft.iot.vmp.service.redisMsg.RedisGbPlayMsgListener;
|
|
|
import com.genersoft.iot.vmp.service.redisMsg.RedisPushStreamResponseListener;
|
|
|
import com.genersoft.iot.vmp.storager.IRedisCatchStorage;
|
|
|
import com.genersoft.iot.vmp.storager.IVideoManagerStorage;
|
|
|
@@ -51,6 +57,7 @@ import org.slf4j.Logger;
|
|
|
import org.slf4j.LoggerFactory;
|
|
|
import org.springframework.beans.factory.InitializingBean;
|
|
|
import org.springframework.beans.factory.annotation.Autowired;
|
|
|
+import org.springframework.data.redis.core.RedisTemplate;
|
|
|
import org.springframework.stereotype.Component;
|
|
|
|
|
|
import javax.sdp.*;
|
|
|
@@ -64,6 +71,7 @@ import java.time.Instant;
|
|
|
import java.util.Map;
|
|
|
import java.util.Random;
|
|
|
import java.util.Vector;
|
|
|
+import java.util.*;
|
|
|
|
|
|
/**
|
|
|
* SIP命令类型: INVITE请求
|
|
|
@@ -92,16 +100,16 @@ public class InviteRequestProcessor extends SIPRequestProcessorParent implements
|
|
|
private IRedisCatchStorage redisCatchStorage;
|
|
|
|
|
|
@Autowired
|
|
|
- private IInviteStreamService inviteStreamService;
|
|
|
+ private IRedisRpcService redisRpcService;
|
|
|
|
|
|
@Autowired
|
|
|
- private SSRCFactory ssrcFactory;
|
|
|
+ private RedisTemplate<Object, Object> redisTemplate;
|
|
|
|
|
|
@Autowired
|
|
|
- private DynamicTask dynamicTask;
|
|
|
+ private SSRCFactory ssrcFactory;
|
|
|
|
|
|
@Autowired
|
|
|
- private RedisPushStreamResponseListener redisPushStreamResponseListener;
|
|
|
+ private DynamicTask dynamicTask;
|
|
|
|
|
|
@Autowired
|
|
|
private IPlayService playService;
|
|
|
@@ -124,18 +132,17 @@ public class InviteRequestProcessor extends SIPRequestProcessorParent implements
|
|
|
@Autowired
|
|
|
private UserSetting userSetting;
|
|
|
|
|
|
- @Autowired
|
|
|
- private ZLMMediaListManager mediaListManager;
|
|
|
-
|
|
|
@Autowired
|
|
|
private SipConfig config;
|
|
|
|
|
|
+ @Autowired
|
|
|
+ private VideoStreamSessionManager streamSession;
|
|
|
|
|
|
@Autowired
|
|
|
- private RedisGbPlayMsgListener redisGbPlayMsgListener;
|
|
|
+ private SendRtpPortManager sendRtpPortManager;
|
|
|
|
|
|
@Autowired
|
|
|
- private VideoStreamSessionManager streamSession;
|
|
|
+ private RedisPushStreamResponseListener redisPushStreamResponseListener;
|
|
|
|
|
|
|
|
|
@Override
|
|
|
@@ -552,43 +559,79 @@ public class InviteRequestProcessor extends SIPRequestProcessorParent implements
|
|
|
|
|
|
}
|
|
|
} else if (gbStream != null) {
|
|
|
-
|
|
|
- String ssrc;
|
|
|
- if (userSetting.getUseCustomSsrcForParentInvite() || gb28181Sdp.getSsrc() == null) {
|
|
|
- // 上级平台点播时不使用上级平台指定的ssrc,使用自定义的ssrc,参考国标文档-点播外域设备媒体流SSRC处理方式
|
|
|
- ssrc = "Play".equalsIgnoreCase(sessionName) ? ssrcFactory.getPlaySsrc(mediaServerItem.getId()) : ssrcFactory.getPlayBackSsrc(mediaServerItem.getId());
|
|
|
- }else {
|
|
|
- ssrc = gb28181Sdp.getSsrc();
|
|
|
+ SendRtpItem sendRtpItem = new SendRtpItem();
|
|
|
+ if (!userSetting.getUseCustomSsrcForParentInvite() && gb28181Sdp.getSsrc() != null) {
|
|
|
+ sendRtpItem.setSsrc(gb28181Sdp.getSsrc());
|
|
|
}
|
|
|
|
|
|
+ if (tcpActive != null) {
|
|
|
+ sendRtpItem.setTcpActive(tcpActive);
|
|
|
+ }
|
|
|
+ sendRtpItem.setTcp(mediaTransmissionTCP);
|
|
|
+ sendRtpItem.setRtcp(platform.isRtcp());
|
|
|
+ sendRtpItem.setPlatformName(platform.getName());
|
|
|
+ sendRtpItem.setPlatformId(platform.getServerGBId());
|
|
|
+ sendRtpItem.setMediaServerId(mediaServerItem.getId());
|
|
|
+ sendRtpItem.setChannelId(channelId);
|
|
|
+ sendRtpItem.setIp(addressStr);
|
|
|
+ sendRtpItem.setPort(port);
|
|
|
+ sendRtpItem.setUsePs(true);
|
|
|
+ sendRtpItem.setApp(gbStream.getApp());
|
|
|
+ sendRtpItem.setStream(gbStream.getStream());
|
|
|
+ sendRtpItem.setCallId(callIdHeader.getCallId());
|
|
|
+ sendRtpItem.setFromTag(request.getFromTag());
|
|
|
+ sendRtpItem.setOnlyAudio(false);
|
|
|
+ sendRtpItem.setStatus(0);
|
|
|
+ sendRtpItem.setSessionName(sessionName);
|
|
|
+ // 清理可能存在的缓存避免用到旧的数据
|
|
|
+ List<SendRtpItem> sendRtpItemList = redisCatchStorage.querySendRTPServer(platform.getServerGBId(), channelId, gbStream.getStream());
|
|
|
+ if (!sendRtpItemList.isEmpty()) {
|
|
|
+ for (SendRtpItem rtpItem : sendRtpItemList) {
|
|
|
+ redisCatchStorage.deleteSendRTPServer(rtpItem);
|
|
|
+ }
|
|
|
+ }
|
|
|
if ("push".equals(gbStream.getStreamType())) {
|
|
|
+ sendRtpItem.setPlayType(InviteStreamType.PUSH);
|
|
|
if (streamPushItem != null) {
|
|
|
// 从redis查询是否正在接收这个推流
|
|
|
StreamPushItem pushListItem = redisCatchStorage.getPushListItem(gbStream.getApp(), gbStream.getStream());
|
|
|
if (pushListItem != null) {
|
|
|
- pushListItem.setSelf(userSetting.getServerId().equals(pushListItem.getServerId()));
|
|
|
- // 推流状态
|
|
|
- pushStream(evt, request, gbStream, pushListItem, platform, callIdHeader, mediaServerItem, port, tcpActive,
|
|
|
- mediaTransmissionTCP, channelId, addressStr, ssrc, requesterId);
|
|
|
+ sendRtpItem.setServerId(pushListItem.getSeverId());
|
|
|
+ sendRtpItem.setMediaServerId(pushListItem.getMediaServerId());
|
|
|
+
|
|
|
+ StreamPushItem transform = streamPushService.transform(pushListItem);
|
|
|
+ transform.setSelf(userSetting.getServerId().equals(pushListItem.getSeverId()));
|
|
|
+ redisCatchStorage.updateSendRTPSever(sendRtpItem);
|
|
|
+ // 开始推流
|
|
|
+ sendPushStream(sendRtpItem, mediaServerItem, platform, request);
|
|
|
}else {
|
|
|
- // 未推流 拉起
|
|
|
- notifyStreamOnline(evt, request, gbStream, streamPushItem, platform, callIdHeader, mediaServerItem, port, tcpActive,
|
|
|
- mediaTransmissionTCP, channelId, addressStr, ssrc, requesterId);
|
|
|
+ if (!platform.isStartOfflinePush()) {
|
|
|
+ // 平台设置中关闭了拉起离线的推流则直接回复
|
|
|
+ try {
|
|
|
+ logger.info("[上级点播] 失败,推流设备未推流,channel: {}, app: {}, stream: {}", sendRtpItem.getChannelId(), sendRtpItem.getApp(), sendRtpItem.getStream());
|
|
|
+ responseAck(request, Response.TEMPORARILY_UNAVAILABLE, "channel stream not pushing");
|
|
|
+ } catch (SipException | InvalidArgumentException | ParseException e) {
|
|
|
+ logger.error("[命令发送失败] invite 通道未推流: {}", e.getMessage());
|
|
|
+ }
|
|
|
+ return;
|
|
|
+ }
|
|
|
+ notifyPushStreamOnline(sendRtpItem, mediaServerItem, platform, request);
|
|
|
}
|
|
|
}
|
|
|
} else if ("proxy".equals(gbStream.getStreamType())) {
|
|
|
if (null != proxyByAppAndStream) {
|
|
|
+ if (sendRtpItem.getSsrc() == null) {
|
|
|
+ // 上级平台点播时不使用上级平台指定的ssrc,使用自定义的ssrc,参考国标文档-点播外域设备媒体流SSRC处理方式
|
|
|
+ String ssrc = "Play".equalsIgnoreCase(sessionName) ? ssrcFactory.getPlaySsrc(mediaServerItem.getId()) : ssrcFactory.getPlayBackSsrc(mediaServerItem.getId());
|
|
|
+ sendRtpItem.setSsrc(ssrc);
|
|
|
+ }
|
|
|
if (proxyByAppAndStream.isStatus()) {
|
|
|
- pushProxyStream(evt, request, gbStream, platform, callIdHeader, mediaServerItem, port, tcpActive,
|
|
|
- mediaTransmissionTCP, channelId, addressStr, ssrc, requesterId);
|
|
|
+ sendProxyStream(sendRtpItem, mediaServerItem, platform, request);
|
|
|
} else {
|
|
|
//开启代理拉流
|
|
|
- notifyStreamOnline(evt, request, gbStream, null, platform, callIdHeader, mediaServerItem, port, tcpActive,
|
|
|
- mediaTransmissionTCP, channelId, addressStr, ssrc, requesterId);
|
|
|
+ notifyProxyStreamOnline(sendRtpItem, mediaServerItem, platform, request);
|
|
|
}
|
|
|
}
|
|
|
-
|
|
|
-
|
|
|
}
|
|
|
}
|
|
|
}
|
|
|
@@ -614,58 +657,42 @@ public class InviteRequestProcessor extends SIPRequestProcessorParent implements
|
|
|
/**
|
|
|
* 安排推流
|
|
|
*/
|
|
|
- private void pushProxyStream(RequestEvent evt, SIPRequest request, GbStream gbStream, ParentPlatform platform,
|
|
|
- CallIdHeader callIdHeader, MediaServer mediaServer,
|
|
|
- int port, Boolean tcpActive, boolean mediaTransmissionTCP,
|
|
|
- String channelId, String addressStr, String ssrc, String requesterId) {
|
|
|
- Boolean streamReady = mediaServerService.isStreamReady(mediaServer, gbStream.getApp(), gbStream.getStream());
|
|
|
+ private void sendProxyStream(SendRtpItem sendRtpItem, MediaServerItem mediaServerItem, ParentPlatform platform, SIPRequest request) {
|
|
|
+ Boolean streamReady = zlmServerFactory.isStreamReady(mediaServerItem, sendRtpItem.getApp(), sendRtpItem.getStream());
|
|
|
if (streamReady != null && streamReady) {
|
|
|
|
|
|
// 自平台内容
|
|
|
- SendRtpItem sendRtpItem = mediaServerService.createSendRtpItem(mediaServer, addressStr, port, ssrc, requesterId,
|
|
|
- gbStream.getApp(), gbStream.getStream(), channelId, mediaTransmissionTCP, platform.isRtcp());
|
|
|
-
|
|
|
- if (sendRtpItem == null) {
|
|
|
- logger.warn("服务器端口资源不足");
|
|
|
- try {
|
|
|
- responseAck(request, Response.BUSY_HERE);
|
|
|
- } catch (SipException | InvalidArgumentException | ParseException e) {
|
|
|
- logger.error("[命令发送失败] invite 服务器端口资源不足: {}", e.getMessage());
|
|
|
+ int localPort = sendRtpPortManager.getNextPort(mediaServerItem);
|
|
|
+ if (localPort == 0) {
|
|
|
+ logger.warn("服务器端口资源不足");
|
|
|
+ try {
|
|
|
+ responseAck(request, Response.BUSY_HERE);
|
|
|
+ } catch (SipException | InvalidArgumentException | ParseException e) {
|
|
|
+ logger.error("[命令发送失败] invite 服务器端口资源不足: {}", e.getMessage());
|
|
|
+ }
|
|
|
+ return;
|
|
|
}
|
|
|
- return;
|
|
|
- }
|
|
|
- if (tcpActive != null) {
|
|
|
- sendRtpItem.setTcpActive(tcpActive);
|
|
|
- }
|
|
|
- sendRtpItem.setPlayType(InviteStreamType.PUSH);
|
|
|
+ sendRtpItem.setPlayType(InviteStreamType.PROXY);
|
|
|
// 写入redis, 超时时回复
|
|
|
sendRtpItem.setStatus(1);
|
|
|
- sendRtpItem.setCallId(callIdHeader.getCallId());
|
|
|
- sendRtpItem.setFromTag(request.getFromTag());
|
|
|
+ sendRtpItem.setLocalIp(mediaServerItem.getSdpIp());
|
|
|
|
|
|
SIPResponse response = sendStreamAck(mediaServer, request, sendRtpItem, platform, evt);
|
|
|
if (response != null) {
|
|
|
sendRtpItem.setToTag(response.getToTag());
|
|
|
}
|
|
|
redisCatchStorage.updateSendRTPSever(sendRtpItem);
|
|
|
-
|
|
|
}
|
|
|
-
|
|
|
}
|
|
|
|
|
|
- private void pushStream(RequestEvent evt, SIPRequest request, GbStream gbStream, StreamPushItem streamPushItem, ParentPlatform platform,
|
|
|
- CallIdHeader callIdHeader, MediaServer mediaServerItem,
|
|
|
- int port, Boolean tcpActive, boolean mediaTransmissionTCP,
|
|
|
- String channelId, String addressStr, String ssrc, String requesterId) {
|
|
|
+ private void sendPushStream(SendRtpItem sendRtpItem, MediaServerItem mediaServerItem, ParentPlatform platform, SIPRequest request) {
|
|
|
// 推流
|
|
|
- if (streamPushItem.isSelf()) {
|
|
|
- Boolean streamReady = mediaServerService.isStreamReady(mediaServerItem, gbStream.getApp(), gbStream.getStream());
|
|
|
+ if (sendRtpItem.getServerId().equals(userSetting.getServerId())) {
|
|
|
+ Boolean streamReady = zlmServerFactory.isStreamReady(mediaServerItem, sendRtpItem.getApp(), sendRtpItem.getStream());
|
|
|
if (streamReady != null && streamReady) {
|
|
|
// 自平台内容
|
|
|
- SendRtpItem sendRtpItem = mediaServerService.createSendRtpItem(mediaServerItem, addressStr, port, ssrc, requesterId,
|
|
|
- gbStream.getApp(), gbStream.getStream(), channelId, mediaTransmissionTCP, platform.isRtcp());
|
|
|
-
|
|
|
- if (sendRtpItem == null) {
|
|
|
+ int localPort = sendRtpPortManager.getNextPort(mediaServerItem);
|
|
|
+ if (localPort == 0) {
|
|
|
logger.warn("服务器端口资源不足");
|
|
|
try {
|
|
|
responseAck(request, Response.BUSY_HERE);
|
|
|
@@ -674,226 +701,168 @@ public class InviteRequestProcessor extends SIPRequestProcessorParent implements
|
|
|
}
|
|
|
return;
|
|
|
}
|
|
|
- if (tcpActive != null) {
|
|
|
- sendRtpItem.setTcpActive(tcpActive);
|
|
|
- }
|
|
|
- sendRtpItem.setPlayType(InviteStreamType.PUSH);
|
|
|
// 写入redis, 超时时回复
|
|
|
sendRtpItem.setStatus(1);
|
|
|
- sendRtpItem.setCallId(callIdHeader.getCallId());
|
|
|
-
|
|
|
- sendRtpItem.setFromTag(request.getFromTag());
|
|
|
- SIPResponse response = sendStreamAck(mediaServerItem, request, sendRtpItem, platform, evt);
|
|
|
+ SIPResponse response = sendStreamAck(request, sendRtpItem, platform);
|
|
|
if (response != null) {
|
|
|
sendRtpItem.setToTag(response.getToTag());
|
|
|
}
|
|
|
+ if (sendRtpItem.getSsrc() == null) {
|
|
|
+ // 上级平台点播时不使用上级平台指定的ssrc,使用自定义的ssrc,参考国标文档-点播外域设备媒体流SSRC处理方式
|
|
|
+ String ssrc = "Play".equalsIgnoreCase(sendRtpItem.getSessionName()) ? ssrcFactory.getPlaySsrc(mediaServerItem.getId()) : ssrcFactory.getPlayBackSsrc(mediaServerItem.getId());
|
|
|
+ sendRtpItem.setSsrc(ssrc);
|
|
|
+ }
|
|
|
redisCatchStorage.updateSendRTPSever(sendRtpItem);
|
|
|
-
|
|
|
} else {
|
|
|
// 不在线 拉起
|
|
|
- notifyStreamOnline(evt, request, gbStream, streamPushItem, platform, callIdHeader, mediaServerItem, port, tcpActive,
|
|
|
- mediaTransmissionTCP, channelId, addressStr, ssrc, requesterId);
|
|
|
+ notifyPushStreamOnline(sendRtpItem, mediaServerItem, platform, request);
|
|
|
}
|
|
|
-
|
|
|
} else {
|
|
|
// 其他平台内容
|
|
|
- otherWvpPushStream(evt, request, gbStream, streamPushItem, platform, callIdHeader, mediaServerItem, port, tcpActive,
|
|
|
- mediaTransmissionTCP, channelId, addressStr, ssrc, requesterId);
|
|
|
+ otherWvpPushStream(sendRtpItem, request, platform);
|
|
|
}
|
|
|
}
|
|
|
|
|
|
/**
|
|
|
* 通知流上线
|
|
|
*/
|
|
|
- private void notifyStreamOnline(RequestEvent evt, SIPRequest request, GbStream gbStream, StreamPushItem streamPushItem, ParentPlatform platform,
|
|
|
- CallIdHeader callIdHeader, MediaServer mediaServerItem,
|
|
|
- int port, Boolean tcpActive, boolean mediaTransmissionTCP,
|
|
|
- String channelId, String addressStr, String ssrc, String requesterId) {
|
|
|
- if ("proxy".equals(gbStream.getStreamType())) {
|
|
|
- // TODO 控制启用以使设备上线
|
|
|
- logger.info("[ app={}, stream={} ]通道未推流,启用流后开始推流", gbStream.getApp(), gbStream.getStream());
|
|
|
- // 监听流上线
|
|
|
- Hook hook = Hook.getInstance(HookType.on_media_arrival, gbStream.getApp(), gbStream.getStream(), mediaServerItem.getId());
|
|
|
- this.hookSubscribe.addSubscribe(hook, (hookData) -> {
|
|
|
- logger.info("[上级点播]拉流代理已经就绪, {}/{}", hookData.getApp(), hookData.getStream());
|
|
|
- dynamicTask.stop(callIdHeader.getCallId());
|
|
|
- pushProxyStream(evt, request, gbStream, platform, callIdHeader, mediaServerItem, port, tcpActive,
|
|
|
- mediaTransmissionTCP, channelId, addressStr, ssrc, requesterId);
|
|
|
- });
|
|
|
- dynamicTask.startDelay(callIdHeader.getCallId(), () -> {
|
|
|
- logger.info("[ app={}, stream={} ] 等待拉流代理流超时", gbStream.getApp(), gbStream.getStream());
|
|
|
- this.hookSubscribe.removeSubscribe(hook);
|
|
|
- }, userSetting.getPlatformPlayTimeout());
|
|
|
- boolean start = streamProxyService.start(gbStream.getApp(), gbStream.getStream());
|
|
|
- if (!start) {
|
|
|
- try {
|
|
|
- responseAck(request, Response.BUSY_HERE, "channel [" + gbStream.getGbId() + "] offline");
|
|
|
- } catch (SipException | InvalidArgumentException | ParseException e) {
|
|
|
- logger.error("[命令发送失败] invite 通道未推流: {}", e.getMessage());
|
|
|
- }
|
|
|
- this.hookSubscribe.removeSubscribe(hook);
|
|
|
- dynamicTask.stop(callIdHeader.getCallId());
|
|
|
+ private void notifyProxyStreamOnline(SendRtpItem sendRtpItem, MediaServerItem mediaServerItem, ParentPlatform platform, SIPRequest request) {
|
|
|
+ // TODO 控制启用以使设备上线
|
|
|
+ logger.info("[ app={}, stream={} ]通道未推流,启用流后开始推流", sendRtpItem.getApp(), sendRtpItem.getStream());
|
|
|
+ // 监听流上线
|
|
|
+ HookSubscribeForStreamChange hookSubscribe = HookSubscribeFactory.on_stream_changed(sendRtpItem.getApp(), sendRtpItem.getStream(), true, "rtsp", mediaServerItem.getId());
|
|
|
+ zlmHttpHookSubscribe.addSubscribe(hookSubscribe, (mediaServerItemInUSe, hookParam) -> {
|
|
|
+ OnStreamChangedHookParam streamChangedHookParam = (OnStreamChangedHookParam)hookParam;
|
|
|
+ logger.info("[上级点播]拉流代理已经就绪, {}/{}", streamChangedHookParam.getApp(), streamChangedHookParam.getStream());
|
|
|
+ dynamicTask.stop(sendRtpItem.getCallId());
|
|
|
+ sendProxyStream(sendRtpItem, mediaServerItem, platform, request);
|
|
|
+ });
|
|
|
+ dynamicTask.startDelay(sendRtpItem.getCallId(), () -> {
|
|
|
+ logger.info("[ app={}, stream={} ] 等待拉流代理流超时", sendRtpItem.getApp(), sendRtpItem.getStream());
|
|
|
+ zlmHttpHookSubscribe.removeSubscribe(hookSubscribe);
|
|
|
+ }, userSetting.getPlatformPlayTimeout());
|
|
|
+ boolean start = streamProxyService.start(sendRtpItem.getApp(), sendRtpItem.getStream());
|
|
|
+ if (!start) {
|
|
|
+ try {
|
|
|
+ responseAck(request, Response.BUSY_HERE, "channel [" + sendRtpItem.getChannelId() + "] offline");
|
|
|
+ } catch (SipException | InvalidArgumentException | ParseException e) {
|
|
|
+ logger.error("[命令发送失败] invite 通道未推流: {}", e.getMessage());
|
|
|
+ }
|
|
|
+ zlmHttpHookSubscribe.removeSubscribe(hookSubscribe);
|
|
|
+ dynamicTask.stop(sendRtpItem.getCallId());
|
|
|
+ }
|
|
|
+ }
|
|
|
+
|
|
|
+ /**
|
|
|
+ * 通知流上线
|
|
|
+ */
|
|
|
+ private void notifyPushStreamOnline(SendRtpItem sendRtpItem, MediaServerItem mediaServerItem, ParentPlatform platform, SIPRequest request) {
|
|
|
+ // 发送redis消息以使设备上线,流上线后被
|
|
|
+ logger.info("[ app={}, stream={} ]通道未推流,发送redis信息控制设备开始推流", sendRtpItem.getApp(), sendRtpItem.getStream());
|
|
|
+ MessageForPushChannel messageForPushChannel = MessageForPushChannel.getInstance(1,
|
|
|
+ sendRtpItem.getApp(), sendRtpItem.getStream(), sendRtpItem.getChannelId(), sendRtpItem.getPlatformId(),
|
|
|
+ platform.getName(), userSetting.getServerId(), sendRtpItem.getMediaServerId());
|
|
|
+ redisCatchStorage.sendStreamPushRequestedMsg(messageForPushChannel);
|
|
|
+ // 设置超时
|
|
|
+ dynamicTask.startDelay(sendRtpItem.getCallId(), () -> {
|
|
|
+ redisRpcService.stopWaitePushStreamOnline(sendRtpItem);
|
|
|
+ logger.info("[ app={}, stream={} ] 等待设备开始推流超时", sendRtpItem.getApp(), sendRtpItem.getStream());
|
|
|
+ try {
|
|
|
+ responseAck(request, Response.REQUEST_TIMEOUT); // 超时
|
|
|
+ } catch (SipException | InvalidArgumentException | ParseException e) {
|
|
|
+ logger.error("未处理的异常 ", e);
|
|
|
}
|
|
|
- } else if ("push".equals(gbStream.getStreamType())) {
|
|
|
- if (!platform.isStartOfflinePush()) {
|
|
|
- // 平台设置中关闭了拉起离线的推流则直接回复
|
|
|
+ }, userSetting.getPlatformPlayTimeout());
|
|
|
+ //
|
|
|
+ long key = redisRpcService.waitePushStreamOnline(sendRtpItem, (sendRtpItemKey) -> {
|
|
|
+ dynamicTask.stop(sendRtpItem.getCallId());
|
|
|
+ if (sendRtpItemKey == null) {
|
|
|
+ logger.warn("[级联点播] 等待推流得到结果未空: {}/{}", sendRtpItem.getApp(), sendRtpItem.getStream());
|
|
|
try {
|
|
|
- logger.info("[上级点播] 失败,推流设备未推流,channel: {}, app: {}, stream: {}", gbStream.getGbId(), gbStream.getApp(), gbStream.getStream());
|
|
|
- responseAck(request, Response.TEMPORARILY_UNAVAILABLE, "channel stream not pushing");
|
|
|
+ responseAck(request, Response.BUSY_HERE);
|
|
|
} catch (SipException | InvalidArgumentException | ParseException e) {
|
|
|
- logger.error("[命令发送失败] invite 通道未推流: {}", e.getMessage());
|
|
|
+ logger.error("未处理的异常 ", e);
|
|
|
}
|
|
|
return;
|
|
|
}
|
|
|
- // 发送redis消息以使设备上线
|
|
|
- logger.info("[ app={}, stream={} ]通道未推流,发送redis信息控制设备开始推流", gbStream.getApp(), gbStream.getStream());
|
|
|
-
|
|
|
- MessageForPushChannel messageForPushChannel = MessageForPushChannel.getInstance(1,
|
|
|
- gbStream.getApp(), gbStream.getStream(), gbStream.getGbId(), gbStream.getPlatformId(),
|
|
|
- platform.getName(), null, gbStream.getMediaServerId());
|
|
|
- redisCatchStorage.sendStreamPushRequestedMsg(messageForPushChannel);
|
|
|
- // 设置超时
|
|
|
- dynamicTask.startDelay(callIdHeader.getCallId(), () -> {
|
|
|
- logger.info("[ app={}, stream={} ] 等待设备开始推流超时", gbStream.getApp(), gbStream.getStream());
|
|
|
+ SendRtpItem sendRtpItemFromRedis = (SendRtpItem)redisTemplate.opsForValue().get(sendRtpItemKey);
|
|
|
+ if (sendRtpItemFromRedis == null) {
|
|
|
+ logger.warn("[级联点播] 等待推流, 未找到redis中缓存的发流信息: {}/{}", sendRtpItem.getApp(), sendRtpItem.getStream());
|
|
|
try {
|
|
|
- redisPushStreamResponseListener.removeEvent(gbStream.getApp(), gbStream.getStream());
|
|
|
- mediaListManager.removedChannelOnlineEventLister(gbStream.getApp(), gbStream.getStream());
|
|
|
- responseAck(request, Response.REQUEST_TIMEOUT); // 超时
|
|
|
+ responseAck(request, Response.BUSY_HERE);
|
|
|
} catch (SipException | InvalidArgumentException | ParseException e) {
|
|
|
logger.error("未处理的异常 ", e);
|
|
|
}
|
|
|
- }, userSetting.getPlatformPlayTimeout());
|
|
|
- // 添加监听
|
|
|
- int finalPort = port;
|
|
|
- Boolean finalTcpActive = tcpActive;
|
|
|
-
|
|
|
- // 添加在本机上线的通知
|
|
|
- mediaListManager.addChannelOnlineEventLister(gbStream.getApp(), gbStream.getStream(), (app, stream, serverId) -> {
|
|
|
- dynamicTask.stop(callIdHeader.getCallId());
|
|
|
- redisPushStreamResponseListener.removeEvent(gbStream.getApp(), gbStream.getStream());
|
|
|
- if (serverId.equals(userSetting.getServerId())) {
|
|
|
- SendRtpItem sendRtpItem = mediaServerService.createSendRtpItem(mediaServerItem, addressStr, finalPort, ssrc, requesterId,
|
|
|
- app, stream, channelId, mediaTransmissionTCP, platform.isRtcp());
|
|
|
-
|
|
|
- if (sendRtpItem == null) {
|
|
|
- logger.warn("上级点时创建sendRTPItem失败,可能是服务器端口资源不足");
|
|
|
- try {
|
|
|
- responseAck(request, Response.BUSY_HERE);
|
|
|
- } catch (SipException e) {
|
|
|
- logger.error("未处理的异常 ", e);
|
|
|
- } catch (InvalidArgumentException e) {
|
|
|
- logger.error("未处理的异常 ", e);
|
|
|
- } catch (ParseException e) {
|
|
|
- logger.error("未处理的异常 ", e);
|
|
|
- }
|
|
|
- return;
|
|
|
- }
|
|
|
- if (finalTcpActive != null) {
|
|
|
- sendRtpItem.setTcpActive(finalTcpActive);
|
|
|
- }
|
|
|
- sendRtpItem.setPlayType(InviteStreamType.PUSH);
|
|
|
- // 写入redis, 超时时回复
|
|
|
- sendRtpItem.setStatus(1);
|
|
|
- sendRtpItem.setCallId(callIdHeader.getCallId());
|
|
|
-
|
|
|
- sendRtpItem.setFromTag(request.getFromTag());
|
|
|
- SIPResponse response = sendStreamAck(mediaServerItem, request, sendRtpItem, platform, evt);
|
|
|
- if (response != null) {
|
|
|
- sendRtpItem.setToTag(response.getToTag());
|
|
|
- }
|
|
|
- redisCatchStorage.updateSendRTPSever(sendRtpItem);
|
|
|
- } else {
|
|
|
- // 其他平台内容
|
|
|
- otherWvpPushStream(evt, request, gbStream, streamPushItem, platform, callIdHeader, mediaServerItem, port, tcpActive,
|
|
|
- mediaTransmissionTCP, channelId, addressStr, ssrc, requesterId);
|
|
|
- }
|
|
|
- });
|
|
|
-
|
|
|
- // 添加回复的拒绝或者错误的通知
|
|
|
- redisPushStreamResponseListener.addEvent(gbStream.getApp(), gbStream.getStream(), response -> {
|
|
|
- if (response.getCode() != 0) {
|
|
|
- dynamicTask.stop(callIdHeader.getCallId());
|
|
|
- mediaListManager.removedChannelOnlineEventLister(gbStream.getApp(), gbStream.getStream());
|
|
|
+ return;
|
|
|
+ }
|
|
|
+ if (sendRtpItemFromRedis.getServerId().equals(userSetting.getServerId())) {
|
|
|
+ logger.info("[级联点播] 等待的推流在本平台上线 {}/{}", sendRtpItem.getApp(), sendRtpItem.getStream());
|
|
|
+ int localPort = sendRtpPortManager.getNextPort(mediaServerItem);
|
|
|
+ if (localPort == 0) {
|
|
|
+ logger.warn("上级点时创建sendRTPItem失败,可能是服务器端口资源不足");
|
|
|
try {
|
|
|
- responseAck(request, Response.TEMPORARILY_UNAVAILABLE, response.getMsg());
|
|
|
+ responseAck(request, Response.BUSY_HERE);
|
|
|
} catch (SipException | InvalidArgumentException | ParseException e) {
|
|
|
- logger.error("[命令发送失败] 国标级联 点播回复: {}", e.getMessage());
|
|
|
+ logger.error("未处理的异常 ", e);
|
|
|
}
|
|
|
+ return;
|
|
|
}
|
|
|
- });
|
|
|
- }
|
|
|
+ sendRtpItem.setLocalPort(localPort);
|
|
|
+ if (!ObjectUtils.isEmpty(platform.getSendStreamIp())) {
|
|
|
+ sendRtpItem.setLocalIp(platform.getSendStreamIp());
|
|
|
+ }
|
|
|
+
|
|
|
+ // 写入redis, 超时时回复
|
|
|
+ sendRtpItem.setStatus(1);
|
|
|
+ SIPResponse response = sendStreamAck(request, sendRtpItem, platform);
|
|
|
+ if (response != null) {
|
|
|
+ sendRtpItem.setToTag(response.getToTag());
|
|
|
+ }
|
|
|
+ redisCatchStorage.updateSendRTPSever(sendRtpItem);
|
|
|
+ } else {
|
|
|
+ // 其他平台内容
|
|
|
+ otherWvpPushStream(sendRtpItemFromRedis, request, platform);
|
|
|
+ }
|
|
|
+ });
|
|
|
+ // 添加回复的拒绝或者错误的通知
|
|
|
+ // redis消息例如: PUBLISH VM_MSG_STREAM_PUSH_RESPONSE '{"code":1,"msg":"失败","app":"1","stream":"2"}'
|
|
|
+ redisPushStreamResponseListener.addEvent(sendRtpItem.getApp(), sendRtpItem.getStream(), response -> {
|
|
|
+ if (response.getCode() != 0) {
|
|
|
+ dynamicTask.stop(sendRtpItem.getCallId());
|
|
|
+ redisRpcService.stopWaitePushStreamOnline(sendRtpItem);
|
|
|
+ redisRpcService.removeCallback(key);
|
|
|
+ try {
|
|
|
+ responseAck(request, Response.TEMPORARILY_UNAVAILABLE, response.getMsg());
|
|
|
+ } catch (SipException | InvalidArgumentException | ParseException e) {
|
|
|
+ logger.error("[命令发送失败] 国标级联 点播回复: {}", e.getMessage());
|
|
|
+ }
|
|
|
+ }
|
|
|
+ });
|
|
|
}
|
|
|
|
|
|
+
|
|
|
+
|
|
|
/**
|
|
|
* 来自其他wvp的推流
|
|
|
*/
|
|
|
- private void otherWvpPushStream(RequestEvent evt, SIPRequest request, GbStream gbStream, StreamPushItem streamPushItem, ParentPlatform platform,
|
|
|
- CallIdHeader callIdHeader, MediaServer mediaServerItem,
|
|
|
- int port, Boolean tcpActive, boolean mediaTransmissionTCP,
|
|
|
- String channelId, String addressStr, String ssrc, String requesterId) {
|
|
|
- logger.info("[级联点播]直播流来自其他平台,发送redis消息");
|
|
|
- // 发送redis消息
|
|
|
- redisGbPlayMsgListener.sendMsg(streamPushItem.getServerId(), streamPushItem.getMediaServerId(),
|
|
|
- streamPushItem.getApp(), streamPushItem.getStream(), addressStr, port, ssrc, requesterId,
|
|
|
- channelId, mediaTransmissionTCP, platform.isRtcp(),platform.getName(), responseSendItemMsg -> {
|
|
|
- SendRtpItem sendRtpItem = responseSendItemMsg.getSendRtpItem();
|
|
|
- if (sendRtpItem == null || responseSendItemMsg.getMediaServerItem() == null) {
|
|
|
- logger.warn("服务器端口资源不足");
|
|
|
- try {
|
|
|
- responseAck(request, Response.BUSY_HERE);
|
|
|
- } catch (SipException e) {
|
|
|
- logger.error("未处理的异常 ", e);
|
|
|
- } catch (InvalidArgumentException e) {
|
|
|
- logger.error("未处理的异常 ", e);
|
|
|
- } catch (ParseException e) {
|
|
|
- logger.error("未处理的异常 ", e);
|
|
|
- }
|
|
|
- return;
|
|
|
- }
|
|
|
- // 收到sendItem
|
|
|
- if (tcpActive != null) {
|
|
|
- sendRtpItem.setTcpActive(tcpActive);
|
|
|
- }
|
|
|
- sendRtpItem.setPlayType(InviteStreamType.PUSH);
|
|
|
- // 写入redis, 超时时回复
|
|
|
- sendRtpItem.setStatus(1);
|
|
|
- sendRtpItem.setCallId(callIdHeader.getCallId());
|
|
|
-
|
|
|
- sendRtpItem.setFromTag(request.getFromTag());
|
|
|
- SIPResponse response = sendStreamAck(responseSendItemMsg.getMediaServerItem(), request, sendRtpItem, platform, evt);
|
|
|
- if (response != null) {
|
|
|
- sendRtpItem.setToTag(response.getToTag());
|
|
|
- }
|
|
|
- redisCatchStorage.updateSendRTPSever(sendRtpItem);
|
|
|
- }, (wvpResult) -> {
|
|
|
-
|
|
|
- // 错误
|
|
|
- if (wvpResult.getCode() == RedisGbPlayMsgListener.ERROR_CODE_OFFLINE) {
|
|
|
- // 离线
|
|
|
- // 查询是否在本机上线了
|
|
|
- StreamPushItem currentStreamPushItem = streamPushService.getPush(streamPushItem.getApp(), streamPushItem.getStream());
|
|
|
- if (currentStreamPushItem.isPushIng()) {
|
|
|
- // 在线状态
|
|
|
- pushStream(evt, request, gbStream, streamPushItem, platform, callIdHeader, mediaServerItem, port, tcpActive,
|
|
|
- mediaTransmissionTCP, channelId, addressStr, ssrc, requesterId);
|
|
|
-
|
|
|
- } else {
|
|
|
- // 不在线 拉起
|
|
|
- notifyStreamOnline(evt, request, gbStream, streamPushItem, platform, callIdHeader, mediaServerItem, port, tcpActive,
|
|
|
- mediaTransmissionTCP, channelId, addressStr, ssrc, requesterId);
|
|
|
- }
|
|
|
- }
|
|
|
- try {
|
|
|
- responseAck(request, Response.BUSY_HERE);
|
|
|
- } catch (InvalidArgumentException | ParseException | SipException e) {
|
|
|
- logger.error("[命令发送失败] 国标级联 点播回复 BUSY_HERE: {}", e.getMessage());
|
|
|
- }
|
|
|
- });
|
|
|
+ private void otherWvpPushStream(SendRtpItem sendRtpItem, SIPRequest request, ParentPlatform platform) {
|
|
|
+ logger.info("[级联点播] 来自其他wvp的推流 {}/{}", sendRtpItem.getApp(), sendRtpItem.getStream());
|
|
|
+ sendRtpItem = redisRpcService.getSendRtpItem(sendRtpItem.getRedisKey());
|
|
|
+ if (sendRtpItem == null) {
|
|
|
+ return;
|
|
|
+ }
|
|
|
+ // 写入redis, 超时时回复
|
|
|
+ sendRtpItem.setStatus(1);
|
|
|
+ SIPResponse response = sendStreamAck(request, sendRtpItem, platform);
|
|
|
+ if (response != null) {
|
|
|
+ sendRtpItem.setToTag(response.getToTag());
|
|
|
+ }
|
|
|
+ redisCatchStorage.updateSendRTPSever(sendRtpItem);
|
|
|
}
|
|
|
|
|
|
- public SIPResponse sendStreamAck(MediaServer mediaServerItem, SIPRequest request, SendRtpItem sendRtpItem, ParentPlatform platform, RequestEvent evt) {
|
|
|
+ public SIPResponse sendStreamAck(MediaServerItem mediaServerItem, SIPRequest request, SendRtpItem sendRtpItem, ParentPlatform platform, RequestEvent evt) {
|
|
|
|
|
|
- String sdpIp = mediaServerItem.getSdpIp();
|
|
|
+ String sdpIp = sendRtpItem.getLocalIp();
|
|
|
if (!ObjectUtils.isEmpty(platform.getSendStreamIp())) {
|
|
|
sdpIp = platform.getSendStreamIp();
|
|
|
}
|