|  | @@ -1,18 +1,28 @@
 | 
	
		
			
				|  |  |  package com.genersoft.iot.vmp.gb28181.transmit.event.request.impl;
 | 
	
		
			
				|  |  |  
 | 
	
		
			
				|  |  | +import com.alibaba.fastjson.JSON;
 | 
	
		
			
				|  |  | +import com.alibaba.fastjson.JSONObject;
 | 
	
		
			
				|  |  | +import com.genersoft.iot.vmp.common.StreamInfo;
 | 
	
		
			
				|  |  |  import com.genersoft.iot.vmp.gb28181.bean.*;
 | 
	
		
			
				|  |  | +import com.genersoft.iot.vmp.gb28181.event.SipSubscribe;
 | 
	
		
			
				|  |  |  import com.genersoft.iot.vmp.gb28181.transmit.SIPProcessorObserver;
 | 
	
		
			
				|  |  | +import com.genersoft.iot.vmp.gb28181.transmit.callback.RequestMessage;
 | 
	
		
			
				|  |  | +import com.genersoft.iot.vmp.gb28181.transmit.cmd.ISIPCommander;
 | 
	
		
			
				|  |  |  import com.genersoft.iot.vmp.gb28181.transmit.cmd.impl.SIPCommander;
 | 
	
		
			
				|  |  |  import com.genersoft.iot.vmp.gb28181.transmit.cmd.impl.SIPCommanderFroPlatform;
 | 
	
		
			
				|  |  |  import com.genersoft.iot.vmp.gb28181.transmit.event.request.ISIPRequestProcessor;
 | 
	
		
			
				|  |  |  import com.genersoft.iot.vmp.gb28181.transmit.event.request.SIPRequestProcessorParent;
 | 
	
		
			
				|  |  | +import com.genersoft.iot.vmp.media.zlm.ZLMHttpHookSubscribe;
 | 
	
		
			
				|  |  |  import com.genersoft.iot.vmp.media.zlm.ZLMRTPServerFactory;
 | 
	
		
			
				|  |  |  import com.genersoft.iot.vmp.media.zlm.dto.MediaServerItem;
 | 
	
		
			
				|  |  |  import com.genersoft.iot.vmp.service.IMediaServerService;
 | 
	
		
			
				|  |  |  import com.genersoft.iot.vmp.service.IPlayService;
 | 
	
		
			
				|  |  | +import com.genersoft.iot.vmp.service.bean.SSRCInfo;
 | 
	
		
			
				|  |  |  import com.genersoft.iot.vmp.storager.IRedisCatchStorage;
 | 
	
		
			
				|  |  |  import com.genersoft.iot.vmp.storager.IVideoManagerStorager;
 | 
	
		
			
				|  |  |  import com.genersoft.iot.vmp.vmanager.gb28181.play.bean.PlayResult;
 | 
	
		
			
				|  |  | +import gov.nist.javax.sdp.TimeDescriptionImpl;
 | 
	
		
			
				|  |  | +import gov.nist.javax.sdp.fields.TimeField;
 | 
	
		
			
				|  |  |  import gov.nist.javax.sip.address.AddressImpl;
 | 
	
		
			
				|  |  |  import gov.nist.javax.sip.address.SipUri;
 | 
	
		
			
				|  |  |  import org.slf4j.Logger;
 | 
	
	
		
			
				|  | @@ -27,10 +37,13 @@ import javax.sip.RequestEvent;
 | 
	
		
			
				|  |  |  import javax.sip.ServerTransaction;
 | 
	
		
			
				|  |  |  import javax.sip.SipException;
 | 
	
		
			
				|  |  |  import javax.sip.address.SipURI;
 | 
	
		
			
				|  |  | +import javax.sip.header.CallIdHeader;
 | 
	
		
			
				|  |  |  import javax.sip.header.FromHeader;
 | 
	
		
			
				|  |  |  import javax.sip.message.Request;
 | 
	
		
			
				|  |  |  import javax.sip.message.Response;
 | 
	
		
			
				|  |  |  import java.text.ParseException;
 | 
	
		
			
				|  |  | +import java.text.SimpleDateFormat;
 | 
	
		
			
				|  |  | +import java.util.Date;
 | 
	
		
			
				|  |  |  import java.util.List;
 | 
	
		
			
				|  |  |  import java.util.Vector;
 | 
	
		
			
				|  |  |  
 | 
	
	
		
			
				|  | @@ -60,6 +73,9 @@ public class InviteRequestProcessor extends SIPRequestProcessorParent implements
 | 
	
		
			
				|  |  |  	@Autowired
 | 
	
		
			
				|  |  |  	private IPlayService playService;
 | 
	
		
			
				|  |  |  
 | 
	
		
			
				|  |  | +	@Autowired
 | 
	
		
			
				|  |  | +	private ISIPCommander commander;
 | 
	
		
			
				|  |  | +
 | 
	
		
			
				|  |  |  	@Autowired
 | 
	
		
			
				|  |  |  	private ZLMRTPServerFactory zlmrtpServerFactory;
 | 
	
		
			
				|  |  |  
 | 
	
	
		
			
				|  | @@ -69,6 +85,7 @@ public class InviteRequestProcessor extends SIPRequestProcessorParent implements
 | 
	
		
			
				|  |  |  	@Autowired
 | 
	
		
			
				|  |  |  	private SIPProcessorObserver sipProcessorObserver;
 | 
	
		
			
				|  |  |  
 | 
	
		
			
				|  |  | +
 | 
	
		
			
				|  |  |  	@Override
 | 
	
		
			
				|  |  |  	public void afterPropertiesSet() throws Exception {
 | 
	
		
			
				|  |  |  		// 添加消息处理的订阅
 | 
	
	
		
			
				|  | @@ -84,6 +101,7 @@ public class InviteRequestProcessor extends SIPRequestProcessorParent implements
 | 
	
		
			
				|  |  |  	@Override
 | 
	
		
			
				|  |  |  	public void process(RequestEvent evt) {
 | 
	
		
			
				|  |  |  		//  Invite Request消息实现,此消息一般为级联消息,上级给下级发送请求视频指令
 | 
	
		
			
				|  |  | +		Long startTimeForInvite = System.currentTimeMillis();
 | 
	
		
			
				|  |  |  		try {
 | 
	
		
			
				|  |  |  			Request request = evt.getRequest();
 | 
	
		
			
				|  |  |  			SipURI sipURI = (SipURI) request.getRequestURI();
 | 
	
	
		
			
				|  | @@ -91,6 +109,7 @@ public class InviteRequestProcessor extends SIPRequestProcessorParent implements
 | 
	
		
			
				|  |  |  			String requesterId = null;
 | 
	
		
			
				|  |  |  
 | 
	
		
			
				|  |  |  			FromHeader fromHeader = (FromHeader)request.getHeader(FromHeader.NAME);
 | 
	
		
			
				|  |  | +			CallIdHeader callIdHeader = (CallIdHeader)request.getHeader(CallIdHeader.NAME);
 | 
	
		
			
				|  |  |  			AddressImpl address = (AddressImpl) fromHeader.getAddress();
 | 
	
		
			
				|  |  |  			SipUri uri = (SipUri) address.getURI();
 | 
	
		
			
				|  |  |  			requesterId = uri.getUser();
 | 
	
	
		
			
				|  | @@ -101,7 +120,7 @@ public class InviteRequestProcessor extends SIPRequestProcessorParent implements
 | 
	
		
			
				|  |  |  				return;
 | 
	
		
			
				|  |  |  			}
 | 
	
		
			
				|  |  |  
 | 
	
		
			
				|  |  | -			// 查询请求方是否上级平台
 | 
	
		
			
				|  |  | +			// 查询请求是否来自上级平台\设备
 | 
	
		
			
				|  |  |  			ParentPlatform platform = storager.queryParentPlatByServerGBId(requesterId);
 | 
	
		
			
				|  |  |  			if (platform != null) {
 | 
	
		
			
				|  |  |  				// 查询平台下是否有该通道
 | 
	
	
		
			
				|  | @@ -158,7 +177,21 @@ public class InviteRequestProcessor extends SIPRequestProcessorParent implements
 | 
	
		
			
				|  |  |  					ssrc = ssrcDefault;
 | 
	
		
			
				|  |  |  					sdp = SdpFactory.getInstance().createSessionDescription(contentString);
 | 
	
		
			
				|  |  |  				}
 | 
	
		
			
				|  |  | -
 | 
	
		
			
				|  |  | +				String sessionName = sdp.getSessionName().getValue();
 | 
	
		
			
				|  |  | +
 | 
	
		
			
				|  |  | +				Long startTime = null;
 | 
	
		
			
				|  |  | +				Long stopTime = null;
 | 
	
		
			
				|  |  | +				Date start = null;
 | 
	
		
			
				|  |  | +				Date end = null;
 | 
	
		
			
				|  |  | +				if (sdp.getTimeDescriptions(false) != null && sdp.getTimeDescriptions(false).size() > 0) {
 | 
	
		
			
				|  |  | +					TimeDescriptionImpl timeDescription = (TimeDescriptionImpl)(sdp.getTimeDescriptions(false).get(0));
 | 
	
		
			
				|  |  | +					TimeField startTimeFiled = (TimeField)timeDescription.getTime();
 | 
	
		
			
				|  |  | +					startTime = startTimeFiled.getStartTime();
 | 
	
		
			
				|  |  | +					stopTime = startTimeFiled.getStopTime();
 | 
	
		
			
				|  |  | +
 | 
	
		
			
				|  |  | +					start = new Date(startTime*1000);
 | 
	
		
			
				|  |  | +					end = new Date(stopTime*1000);
 | 
	
		
			
				|  |  | +				}
 | 
	
		
			
				|  |  |  				//  获取支持的格式
 | 
	
		
			
				|  |  |  				Vector mediaDescriptions = sdp.getMediaDescriptions(true);
 | 
	
		
			
				|  |  |  				// 查看是否支持PS 负载96
 | 
	
	
		
			
				|  | @@ -228,23 +261,31 @@ public class InviteRequestProcessor extends SIPRequestProcessorParent implements
 | 
	
		
			
				|  |  |  						responseAck(evt, Response.BUSY_HERE);
 | 
	
		
			
				|  |  |  						return;
 | 
	
		
			
				|  |  |  					}
 | 
	
		
			
				|  |  | -
 | 
	
		
			
				|  |  | +					sendRtpItem.setCallId(callIdHeader.getCallId());
 | 
	
		
			
				|  |  | +					sendRtpItem.setPlay("Play".equals(sessionName));
 | 
	
		
			
				|  |  |  					// 写入redis, 超时时回复
 | 
	
		
			
				|  |  |  					redisCatchStorage.updateSendRTPSever(sendRtpItem);
 | 
	
		
			
				|  |  | -					// 通知下级推流,
 | 
	
		
			
				|  |  | -					PlayResult playResult = playService.play(mediaServerItem,device.getDeviceId(), channelId, (mediaServerItemInUSe, responseJSON)->{
 | 
	
		
			
				|  |  | -						// 收到推流, 回复200OK, 等待ack
 | 
	
		
			
				|  |  | +
 | 
	
		
			
				|  |  | +					Device finalDevice = device;
 | 
	
		
			
				|  |  | +					MediaServerItem finalMediaServerItem = mediaServerItem;
 | 
	
		
			
				|  |  | +					Long finalStartTime = startTime;
 | 
	
		
			
				|  |  | +					Long finalStopTime = stopTime;
 | 
	
		
			
				|  |  | +					ZLMHttpHookSubscribe.Event hookEvent = (mediaServerItemInUSe, responseJSON)->{
 | 
	
		
			
				|  |  | +						logger.info("[上级点播]收到下级开始点播订阅, {}/{}", sendRtpItem.getApp(), sendRtpItem.getStreamId());
 | 
	
		
			
				|  |  |  						// if (sendRtpItem == null) return;
 | 
	
		
			
				|  |  |  						sendRtpItem.setStatus(1);
 | 
	
		
			
				|  |  |  						redisCatchStorage.updateSendRTPSever(sendRtpItem);
 | 
	
		
			
				|  |  | -						// TODO 添加对tcp的支持
 | 
	
		
			
				|  |  |  
 | 
	
		
			
				|  |  |  						StringBuffer content = new StringBuffer(200);
 | 
	
		
			
				|  |  |  						content.append("v=0\r\n");
 | 
	
		
			
				|  |  |  						content.append("o="+ channelId +" 0 0 IN IP4 "+mediaServerItemInUSe.getSdpIp()+"\r\n");
 | 
	
		
			
				|  |  | -						content.append("s=Play\r\n");
 | 
	
		
			
				|  |  | +						content.append("s=" + sessionName+"\r\n");
 | 
	
		
			
				|  |  |  						content.append("c=IN IP4 "+mediaServerItemInUSe.getSdpIp()+"\r\n");
 | 
	
		
			
				|  |  | -						content.append("t=0 0\r\n");
 | 
	
		
			
				|  |  | +						if ("Playback".equals(sessionName)) {
 | 
	
		
			
				|  |  | +							content.append("t=" + finalStartTime + " " + finalStopTime + "\r\n");
 | 
	
		
			
				|  |  | +						}else {
 | 
	
		
			
				|  |  | +							content.append("t=0 0\r\n");
 | 
	
		
			
				|  |  | +						}
 | 
	
		
			
				|  |  |  						content.append("m=video "+ sendRtpItem.getLocalPort()+" RTP/AVP 96\r\n");
 | 
	
		
			
				|  |  |  						content.append("a=sendonly\r\n");
 | 
	
		
			
				|  |  |  						content.append("a=rtpmap:96 PS/90000\r\n");
 | 
	
	
		
			
				|  | @@ -260,7 +301,11 @@ public class InviteRequestProcessor extends SIPRequestProcessorParent implements
 | 
	
		
			
				|  |  |  						} catch (ParseException e) {
 | 
	
		
			
				|  |  |  							e.printStackTrace();
 | 
	
		
			
				|  |  |  						}
 | 
	
		
			
				|  |  | -					} ,((event) -> {
 | 
	
		
			
				|  |  | +						if ("Playback".equals(sessionName) && responseJSON != null) {
 | 
	
		
			
				|  |  | +							playService.onPublishHandlerForPlayBack(finalMediaServerItem, responseJSON, finalDevice.getDeviceId(), channelId, null);
 | 
	
		
			
				|  |  | +						}
 | 
	
		
			
				|  |  | +					};
 | 
	
		
			
				|  |  | +					SipSubscribe.Event errorEvent = ((event) -> {
 | 
	
		
			
				|  |  |  						// 未知错误。直接转发设备点播的错误
 | 
	
		
			
				|  |  |  						Response response = null;
 | 
	
		
			
				|  |  |  						try {
 | 
	
	
		
			
				|  | @@ -271,11 +316,27 @@ public class InviteRequestProcessor extends SIPRequestProcessorParent implements
 | 
	
		
			
				|  |  |  						} catch (ParseException | SipException | InvalidArgumentException e) {
 | 
	
		
			
				|  |  |  							e.printStackTrace();
 | 
	
		
			
				|  |  |  						}
 | 
	
		
			
				|  |  | -					}));
 | 
	
		
			
				|  |  | -					if (logger.isDebugEnabled()) {
 | 
	
		
			
				|  |  | -						logger.debug(playResult.getResult().toString());
 | 
	
		
			
				|  |  | +					});
 | 
	
		
			
				|  |  | +					if ("Playback".equals(sessionName)) {
 | 
	
		
			
				|  |  | +						sendRtpItem.setPlay(false);
 | 
	
		
			
				|  |  | +						SSRCInfo ssrcInfo = mediaServerService.openRTPServer(mediaServerItem, sendRtpItem.getSsrc(), true);
 | 
	
		
			
				|  |  | +						sendRtpItem.setStreamId(ssrc);
 | 
	
		
			
				|  |  | +						SimpleDateFormat format = new SimpleDateFormat("yyyy-MM-dd HH:mm:ss");
 | 
	
		
			
				|  |  | +						commander.playbackStreamCmd(mediaServerItem, ssrcInfo, device, channelId, format.format(start), format.format(end), hookEvent, errorEvent);
 | 
	
		
			
				|  |  | +					}else {
 | 
	
		
			
				|  |  | +						sendRtpItem.setPlay(true);
 | 
	
		
			
				|  |  | +						StreamInfo streamInfo = redisCatchStorage.queryPlayByDevice(device.getDeviceId(), channelId);
 | 
	
		
			
				|  |  | +						if (streamInfo == null) {
 | 
	
		
			
				|  |  | +							if (mediaServerItem.isRtpEnable()) {
 | 
	
		
			
				|  |  | +								sendRtpItem.setStreamId(String.format("%s_%s", device.getDeviceId(), channelId));
 | 
	
		
			
				|  |  | +							}
 | 
	
		
			
				|  |  | +							sendRtpItem.setPlay(false);
 | 
	
		
			
				|  |  | +							playService.play(mediaServerItem,device.getDeviceId(), channelId, hookEvent,errorEvent);
 | 
	
		
			
				|  |  | +						}else {
 | 
	
		
			
				|  |  | +							sendRtpItem.setStreamId(streamInfo.getStreamId());
 | 
	
		
			
				|  |  | +							hookEvent.response(mediaServerItem, null);
 | 
	
		
			
				|  |  | +						}
 | 
	
		
			
				|  |  |  					}
 | 
	
		
			
				|  |  | -
 | 
	
		
			
				|  |  |  				}else if (gbStream != null) {
 | 
	
		
			
				|  |  |  					SendRtpItem sendRtpItem = zlmrtpServerFactory.createSendRtpItem(mediaServerItem, addressStr, port, ssrc, requesterId,
 | 
	
		
			
				|  |  |  							gbStream.getApp(), gbStream.getStream(), channelId,
 | 
	
	
		
			
				|  | @@ -295,7 +356,6 @@ public class InviteRequestProcessor extends SIPRequestProcessorParent implements
 | 
	
		
			
				|  |  |  
 | 
	
		
			
				|  |  |  					sendRtpItem.setStatus(1);
 | 
	
		
			
				|  |  |  					redisCatchStorage.updateSendRTPSever(sendRtpItem);
 | 
	
		
			
				|  |  | -					// TODO 添加对tcp的支持
 | 
	
		
			
				|  |  |  					StringBuffer content = new StringBuffer(200);
 | 
	
		
			
				|  |  |  					content.append("v=0\r\n");
 | 
	
		
			
				|  |  |  					content.append("o="+ channelId +" 0 0 IN IP4 "+mediaServerItem.getSdpIp()+"\r\n");
 |