|  |  |  | 
|---|
|  |  |  | 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.gb28181.utils.SipUtils; | 
|---|
|  |  |  | import com.genersoft.iot.vmp.media.event.hook.Hook; | 
|---|
|  |  |  | import com.genersoft.iot.vmp.media.event.hook.HookType; | 
|---|
|  |  |  | import com.genersoft.iot.vmp.media.service.IMediaServerService; | 
|---|
|  |  |  | import com.genersoft.iot.vmp.media.zlm.ZLMMediaListManager; | 
|---|
|  |  |  | import com.genersoft.iot.vmp.media.zlm.ZLMServerFactory; | 
|---|
|  |  |  | import com.genersoft.iot.vmp.media.zlm.ZlmHttpHookSubscribe; | 
|---|
|  |  |  | import com.genersoft.iot.vmp.media.event.hook.HookSubscribe; | 
|---|
|  |  |  | 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.utils.DateUtil; | 
|---|
|  |  |  | import gov.nist.javax.sdp.TimeDescriptionImpl; | 
|---|
|  |  |  | import gov.nist.javax.sdp.fields.TimeField; | 
|---|
|  |  |  | import gov.nist.javax.sdp.fields.URIField; | 
|---|
|  |  |  | import gov.nist.javax.sip.message.SIPRequest; | 
|---|
|  |  |  | import gov.nist.javax.sip.message.SIPResponse; | 
|---|
|  |  |  | import org.apache.commons.lang3.StringUtils; | 
|---|
|  |  |  | import org.slf4j.Logger; | 
|---|
|  |  |  | import org.slf4j.LoggerFactory; | 
|---|
|  |  |  | import org.springframework.beans.factory.InitializingBean; | 
|---|
|  |  |  | 
|---|
|  |  |  | private IMediaServerService mediaServerService; | 
|---|
|  |  |  |  | 
|---|
|  |  |  | @Autowired | 
|---|
|  |  |  | private ZlmHttpHookSubscribe zlmHttpHookSubscribe; | 
|---|
|  |  |  | private HookSubscribe hookSubscribe; | 
|---|
|  |  |  |  | 
|---|
|  |  |  | @Autowired | 
|---|
|  |  |  | private SIPProcessorObserver sipProcessorObserver; | 
|---|
|  |  |  | 
|---|
|  |  |  | public void process(RequestEvent evt) { | 
|---|
|  |  |  | //  Invite Request消息实现,此消息一般为级联消息,上级给下级发送请求视频指令 | 
|---|
|  |  |  | try { | 
|---|
|  |  |  | SIPRequest request = (SIPRequest) evt.getRequest(); | 
|---|
|  |  |  | String channelId = SipUtils.getChannelIdFromRequest(request); | 
|---|
|  |  |  | SIPRequest request = (SIPRequest)evt.getRequest(); | 
|---|
|  |  |  | String channelIdFromSub = SipUtils.getChannelIdFromRequest(request); | 
|---|
|  |  |  |  | 
|---|
|  |  |  | // 解析sdp消息, 使用jainsip 自带的sdp解析方式 | 
|---|
|  |  |  | String contentString = new String(request.getRawContent()); | 
|---|
|  |  |  | Gb28181Sdp gb28181Sdp = SipUtils.parseSDP(contentString); | 
|---|
|  |  |  | SessionDescription sdp = gb28181Sdp.getBaseSdb(); | 
|---|
|  |  |  | String sessionName = sdp.getSessionName().getValue(); | 
|---|
|  |  |  | String channelIdFromSdp = null; | 
|---|
|  |  |  | if(StringUtils.equalsIgnoreCase("Playback", sessionName)){ | 
|---|
|  |  |  | URIField uriField = (URIField)sdp.getURI(); | 
|---|
|  |  |  | channelIdFromSdp = uriField.getURI().split(":")[0]; | 
|---|
|  |  |  | } | 
|---|
|  |  |  | final String channelId = StringUtils.isNotBlank(channelIdFromSdp) ? channelIdFromSdp : channelIdFromSub; | 
|---|
|  |  |  |  | 
|---|
|  |  |  | String requesterId = SipUtils.getUserIdFromFromHeader(request); | 
|---|
|  |  |  | CallIdHeader callIdHeader = (CallIdHeader) request.getHeader(CallIdHeader.NAME); | 
|---|
|  |  |  | if (requesterId == null || channelId == null) { | 
|---|
|  |  |  | 
|---|
|  |  |  | GbStream gbStream = storager.queryStreamInParentPlatform(requesterId, channelId); | 
|---|
|  |  |  | PlatformCatalog catalog = storager.getCatalog(requesterId, channelId); | 
|---|
|  |  |  |  | 
|---|
|  |  |  | MediaServerItem mediaServerItem = null; | 
|---|
|  |  |  | MediaServer mediaServerItem = null; | 
|---|
|  |  |  | StreamPushItem streamPushItem = null; | 
|---|
|  |  |  | StreamProxyItem proxyByAppAndStream = null; | 
|---|
|  |  |  | // 不是通道可能是直播流 | 
|---|
|  |  |  | 
|---|
|  |  |  | } | 
|---|
|  |  |  | return; | 
|---|
|  |  |  | } | 
|---|
|  |  |  | // 解析sdp消息, 使用jainsip 自带的sdp解析方式 | 
|---|
|  |  |  | String contentString = new String(request.getRawContent()); | 
|---|
|  |  |  |  | 
|---|
|  |  |  | Gb28181Sdp gb28181Sdp = SipUtils.parseSDP(contentString); | 
|---|
|  |  |  | SessionDescription sdp = gb28181Sdp.getBaseSdb(); | 
|---|
|  |  |  | String sessionName = sdp.getSessionName().getValue(); | 
|---|
|  |  |  |  | 
|---|
|  |  |  | Long startTime = null; | 
|---|
|  |  |  | Long stopTime = null; | 
|---|
|  |  |  | 
|---|
|  |  |  | Long finalStopTime = stopTime; | 
|---|
|  |  |  | ErrorCallback<Object> hookEvent = (code, msg, data) -> { | 
|---|
|  |  |  | StreamInfo streamInfo = (StreamInfo)data; | 
|---|
|  |  |  | MediaServerItem mediaServerItemInUSe = mediaServerService.getOne(streamInfo.getMediaServerId()); | 
|---|
|  |  |  | MediaServer mediaServerItemInUSe = mediaServerService.getOne(streamInfo.getMediaServerId()); | 
|---|
|  |  |  | logger.info("[上级Invite]下级已经开始推流。 回复200OK(SDP), {}/{}", streamInfo.getApp(), streamInfo.getStream()); | 
|---|
|  |  |  | //     * 0 等待设备推流上来 | 
|---|
|  |  |  | //     * 1 下级已经推流,等待上级平台回复ack | 
|---|
|  |  |  | 
|---|
|  |  |  | responseSdpAck(request, content.toString(), platform); | 
|---|
|  |  |  | // tcp主动模式,回复sdp后开启监听 | 
|---|
|  |  |  | if (sendRtpItem.isTcpActive()) { | 
|---|
|  |  |  | MediaServerItem mediaInfo = mediaServerService.getOne(sendRtpItem.getMediaServerId()); | 
|---|
|  |  |  | MediaServer mediaInfo = mediaServerService.getOne(sendRtpItem.getMediaServerId()); | 
|---|
|  |  |  | Map<String, Object> param = new HashMap<>(12); | 
|---|
|  |  |  | param.put("vhost","__defaultVhost__"); | 
|---|
|  |  |  | param.put("app",sendRtpItem.getApp()); | 
|---|
|  |  |  | 
|---|
|  |  |  | String endTimeStr = DateUtil.urlFormatter.format(end); | 
|---|
|  |  |  | String stream = device.getDeviceId() + "_" + channelId + "_" + startTimeStr + "_" + endTimeStr; | 
|---|
|  |  |  | SSRCInfo ssrcInfo = mediaServerService.openRTPServer(mediaServerItem, stream, null, device.isSsrcCheck(), true, 0,false, false, device.getStreamModeForParam()); | 
|---|
|  |  |  | sendRtpItem.setStream(stream); | 
|---|
|  |  |  | // 写入redis, 超时时回复 | 
|---|
|  |  |  | redisCatchStorage.updateSendRTPSever(sendRtpItem); | 
|---|
|  |  |  | playService.playBack(mediaServerItem, ssrcInfo, device.getDeviceId(), channelId, DateUtil.formatter.format(start), | 
|---|
|  |  |  | 
|---|
|  |  |  | } else { | 
|---|
|  |  |  | sendRtpItem.setPlayType(InviteStreamType.PLAY); | 
|---|
|  |  |  | String streamId = String.format("%s_%s", device.getDeviceId(), channelId); | 
|---|
|  |  |  | sendRtpItem.setStreamId(streamId); | 
|---|
|  |  |  | sendRtpItem.setStream(streamId); | 
|---|
|  |  |  | redisCatchStorage.updateSendRTPSever(sendRtpItem); | 
|---|
|  |  |  | SSRCInfo ssrcInfo = playService.play(mediaServerItem, device.getDeviceId(), channelId, ssrc, ((code, msg, data) -> { | 
|---|
|  |  |  | if (code == InviteErrorCode.SUCCESS.getCode()) { | 
|---|
|  |  |  | 
|---|
|  |  |  | * 安排推流 | 
|---|
|  |  |  | */ | 
|---|
|  |  |  | private void pushProxyStream(RequestEvent evt, SIPRequest request, GbStream gbStream, ParentPlatform platform, | 
|---|
|  |  |  | CallIdHeader callIdHeader, MediaServerItem mediaServerItem, | 
|---|
|  |  |  | CallIdHeader callIdHeader, MediaServer mediaServerItem, | 
|---|
|  |  |  | int port, Boolean tcpActive, boolean mediaTransmissionTCP, | 
|---|
|  |  |  | String channelId, String addressStr, String ssrc, String requesterId) { | 
|---|
|  |  |  | Boolean streamReady = zlmServerFactory.isStreamReady(mediaServerItem, gbStream.getApp(), gbStream.getStream()); | 
|---|
|  |  |  | 
|---|
|  |  |  | } | 
|---|
|  |  |  |  | 
|---|
|  |  |  | private void pushStream(RequestEvent evt, SIPRequest request, GbStream gbStream, StreamPushItem streamPushItem, ParentPlatform platform, | 
|---|
|  |  |  | CallIdHeader callIdHeader, MediaServerItem mediaServerItem, | 
|---|
|  |  |  | CallIdHeader callIdHeader, MediaServer mediaServerItem, | 
|---|
|  |  |  | int port, Boolean tcpActive, boolean mediaTransmissionTCP, | 
|---|
|  |  |  | String channelId, String addressStr, String ssrc, String requesterId) { | 
|---|
|  |  |  | // 推流 | 
|---|
|  |  |  | 
|---|
|  |  |  | * 通知流上线 | 
|---|
|  |  |  | */ | 
|---|
|  |  |  | private void notifyStreamOnline(RequestEvent evt, SIPRequest request, GbStream gbStream, StreamPushItem streamPushItem, ParentPlatform platform, | 
|---|
|  |  |  | CallIdHeader callIdHeader, MediaServerItem mediaServerItem, | 
|---|
|  |  |  | 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()); | 
|---|
|  |  |  | // 监听流上线 | 
|---|
|  |  |  | HookSubscribeForStreamChange hookSubscribe = HookSubscribeFactory.on_stream_changed(gbStream.getApp(), gbStream.getStream(), true, "rtsp", mediaServerItem.getId()); | 
|---|
|  |  |  | zlmHttpHookSubscribe.addSubscribe(hookSubscribe, (mediaServerItemInUSe, hookParam) -> { | 
|---|
|  |  |  | OnStreamChangedHookParam streamChangedHookParam = (OnStreamChangedHookParam)hookParam; | 
|---|
|  |  |  | logger.info("[上级点播]拉流代理已经就绪, {}/{}", streamChangedHookParam.getApp(), streamChangedHookParam.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()); | 
|---|
|  |  |  | zlmHttpHookSubscribe.removeSubscribe(hookSubscribe); | 
|---|
|  |  |  | this.hookSubscribe.removeSubscribe(hook); | 
|---|
|  |  |  | }, userSetting.getPlatformPlayTimeout()); | 
|---|
|  |  |  | boolean start = streamProxyService.start(gbStream.getApp(), gbStream.getStream()); | 
|---|
|  |  |  | if (!start) { | 
|---|
|  |  |  | 
|---|
|  |  |  | } catch (SipException | InvalidArgumentException | ParseException e) { | 
|---|
|  |  |  | logger.error("[命令发送失败] invite 通道未推流: {}", e.getMessage()); | 
|---|
|  |  |  | } | 
|---|
|  |  |  | zlmHttpHookSubscribe.removeSubscribe(hookSubscribe); | 
|---|
|  |  |  | this.hookSubscribe.removeSubscribe(hook); | 
|---|
|  |  |  | dynamicTask.stop(callIdHeader.getCallId()); | 
|---|
|  |  |  | } | 
|---|
|  |  |  | } else if ("push".equals(gbStream.getStreamType())) { | 
|---|
|  |  |  | 
|---|
|  |  |  | * 来自其他wvp的推流 | 
|---|
|  |  |  | */ | 
|---|
|  |  |  | private void otherWvpPushStream(RequestEvent evt, SIPRequest request, GbStream gbStream, StreamPushItem streamPushItem, ParentPlatform platform, | 
|---|
|  |  |  | CallIdHeader callIdHeader, MediaServerItem mediaServerItem, | 
|---|
|  |  |  | CallIdHeader callIdHeader, MediaServer mediaServerItem, | 
|---|
|  |  |  | int port, Boolean tcpActive, boolean mediaTransmissionTCP, | 
|---|
|  |  |  | String channelId, String addressStr, String ssrc, String requesterId) { | 
|---|
|  |  |  | logger.info("[级联点播]直播流来自其他平台,发送redis消息"); | 
|---|
|  |  |  | 
|---|
|  |  |  | }); | 
|---|
|  |  |  | } | 
|---|
|  |  |  |  | 
|---|
|  |  |  | public SIPResponse sendStreamAck(MediaServerItem mediaServerItem, SIPRequest request, SendRtpItem sendRtpItem, ParentPlatform platform, RequestEvent evt) { | 
|---|
|  |  |  | public SIPResponse sendStreamAck(MediaServer mediaServerItem, SIPRequest request, SendRtpItem sendRtpItem, ParentPlatform platform, RequestEvent evt) { | 
|---|
|  |  |  |  | 
|---|
|  |  |  | StringBuffer content = new StringBuffer(200); | 
|---|
|  |  |  | content.append("v=0\r\n"); | 
|---|
|  |  |  | 
|---|
|  |  |  | } | 
|---|
|  |  |  | if (device != null) { | 
|---|
|  |  |  | logger.info("收到设备" + requesterId + "的语音广播Invite请求"); | 
|---|
|  |  |  | String key = VideoManagerConstants.BROADCAST_WAITE_INVITE + device.getDeviceId() + broadcastCatch.getChannelId(); | 
|---|
|  |  |  | String key = VideoManagerConstants.BROADCAST_WAITE_INVITE + device.getDeviceId(); | 
|---|
|  |  |  | if (!SipUtils.isFrontEnd(device.getDeviceId())) { | 
|---|
|  |  |  | key += broadcastCatch.getChannelId(); | 
|---|
|  |  |  | } | 
|---|
|  |  |  | dynamicTask.stop(key); | 
|---|
|  |  |  | try { | 
|---|
|  |  |  | responseAck(request, Response.TRYING); | 
|---|
|  |  |  | 
|---|
|  |  |  | return; | 
|---|
|  |  |  | } | 
|---|
|  |  |  | String contentString = new String(request.getRawContent()); | 
|---|
|  |  |  | // jainSip不支持y=字段, 移除移除以解析。 | 
|---|
|  |  |  | String ssrc = "0000000404"; | 
|---|
|  |  |  |  | 
|---|
|  |  |  | try { | 
|---|
|  |  |  | Gb28181Sdp gb28181Sdp = SipUtils.parseSDP(contentString); | 
|---|
|  |  |  | 
|---|
|  |  |  | Media media = mediaDescription.getMedia(); | 
|---|
|  |  |  |  | 
|---|
|  |  |  | Vector mediaFormats = media.getMediaFormats(false); | 
|---|
|  |  |  | if (mediaFormats.contains("8")) { | 
|---|
|  |  |  | //                    if (mediaFormats.contains("8")) { | 
|---|
|  |  |  | port = media.getMediaPort(); | 
|---|
|  |  |  | String protocol = media.getProtocol(); | 
|---|
|  |  |  | // 区分TCP发流还是udp, 当前默认udp | 
|---|
|  |  |  | 
|---|
|  |  |  | } | 
|---|
|  |  |  | } | 
|---|
|  |  |  | break; | 
|---|
|  |  |  | } | 
|---|
|  |  |  | //                    } | 
|---|
|  |  |  | } | 
|---|
|  |  |  | if (port == -1) { | 
|---|
|  |  |  | logger.info("不支持的媒体格式,返回415"); | 
|---|
|  |  |  | 
|---|
|  |  |  | return; | 
|---|
|  |  |  | } | 
|---|
|  |  |  | String addressStr = sdp.getOrigin().getAddress(); | 
|---|
|  |  |  | logger.info("设备{}请求语音流,地址:{}:{},ssrc:{}, {}", requesterId, addressStr, port, ssrc, | 
|---|
|  |  |  | logger.info("设备{}请求语音流,地址:{}:{},ssrc:{}, {}", requesterId, addressStr, port, gb28181Sdp.getSsrc(), | 
|---|
|  |  |  | mediaTransmissionTCP ? (tcpActive ? "TCP主动" : "TCP被动") : "UDP"); | 
|---|
|  |  |  |  | 
|---|
|  |  |  | MediaServerItem mediaServerItem = broadcastCatch.getMediaServerItem(); | 
|---|
|  |  |  | MediaServer mediaServerItem = broadcastCatch.getMediaServerItem(); | 
|---|
|  |  |  | if (mediaServerItem == null) { | 
|---|
|  |  |  | logger.warn("未找到语音喊话使用的zlm"); | 
|---|
|  |  |  | try { | 
|---|
|  |  |  | 
|---|
|  |  |  | } | 
|---|
|  |  |  | return; | 
|---|
|  |  |  | } | 
|---|
|  |  |  | logger.info("设备{}请求语音流, 收流地址:{}:{},ssrc:{}, {}, 对讲方式:{}", requesterId, addressStr, port, ssrc, | 
|---|
|  |  |  | logger.info("设备{}请求语音流, 收流地址:{}:{},ssrc:{}, {}, 对讲方式:{}", requesterId, addressStr, port, gb28181Sdp.getSsrc(), | 
|---|
|  |  |  | mediaTransmissionTCP ? (tcpActive ? "TCP主动" : "TCP被动") : "UDP", sdp.getSessionName().getValue()); | 
|---|
|  |  |  | CallIdHeader callIdHeader = (CallIdHeader) request.getHeader(CallIdHeader.NAME); | 
|---|
|  |  |  |  | 
|---|
|  |  |  | SendRtpItem sendRtpItem = zlmServerFactory.createSendRtpItem(mediaServerItem, addressStr, port, ssrc, requesterId, | 
|---|
|  |  |  | SendRtpItem sendRtpItem = zlmServerFactory.createSendRtpItem(mediaServerItem, addressStr, port, gb28181Sdp.getSsrc(), requesterId, | 
|---|
|  |  |  | device.getDeviceId(), broadcastCatch.getChannelId(), | 
|---|
|  |  |  | mediaTransmissionTCP, false); | 
|---|
|  |  |  |  | 
|---|
|  |  |  | 
|---|
|  |  |  |  | 
|---|
|  |  |  | Boolean streamReady = zlmServerFactory.isStreamReady(mediaServerItem, broadcastCatch.getApp(), broadcastCatch.getStream()); | 
|---|
|  |  |  | if (streamReady) { | 
|---|
|  |  |  | sendOk(device, sendRtpItem, sdp, request, mediaServerItem, mediaTransmissionTCP, ssrc); | 
|---|
|  |  |  | sendOk(device, sendRtpItem, sdp, request, mediaServerItem, mediaTransmissionTCP, gb28181Sdp.getSsrc()); | 
|---|
|  |  |  | } else { | 
|---|
|  |  |  | logger.warn("[语音通话], 未发现待推送的流,app={},stream={}", broadcastCatch.getApp(), broadcastCatch.getStream()); | 
|---|
|  |  |  | try { | 
|---|
|  |  |  | 
|---|
|  |  |  | } | 
|---|
|  |  |  | } | 
|---|
|  |  |  |  | 
|---|
|  |  |  | SIPResponse sendOk(Device device, SendRtpItem sendRtpItem, SessionDescription sdp, SIPRequest request, MediaServerItem mediaServerItem, boolean mediaTransmissionTCP, String ssrc) { | 
|---|
|  |  |  | SIPResponse sendOk(Device device, SendRtpItem sendRtpItem, SessionDescription sdp, SIPRequest request, MediaServer mediaServerItem, boolean mediaTransmissionTCP, String ssrc) { | 
|---|
|  |  |  | SIPResponse sipResponse = null; | 
|---|
|  |  |  | try { | 
|---|
|  |  |  | sendRtpItem.setStatus(2); | 
|---|