优化适配zlm的hook保活

This commit is contained in:
648540858 2021-12-08 16:45:50 +08:00
parent 9b1af8ef13
commit ab81136765
20 changed files with 441 additions and 79 deletions

View File

@ -1,9 +1,7 @@
package com.genersoft.iot.vmp.gb28181; package com.genersoft.iot.vmp.gb28181;
import com.genersoft.iot.vmp.conf.SipConfig; import com.genersoft.iot.vmp.conf.SipConfig;
import com.genersoft.iot.vmp.gb28181.event.SipSubscribe;
import com.genersoft.iot.vmp.gb28181.transmit.ISIPProcessorObserver; import com.genersoft.iot.vmp.gb28181.transmit.ISIPProcessorObserver;
import com.genersoft.iot.vmp.gb28181.transmit.SIPProcessorObserver;
import gov.nist.javax.sip.SipProviderImpl; import gov.nist.javax.sip.SipProviderImpl;
import gov.nist.javax.sip.SipStackImpl; import gov.nist.javax.sip.SipStackImpl;
import org.slf4j.Logger; import org.slf4j.Logger;

View File

@ -5,6 +5,7 @@ import com.genersoft.iot.vmp.gb28181.event.offline.OfflineEvent;
import com.genersoft.iot.vmp.gb28181.event.platformKeepaliveExpire.PlatformKeepaliveExpireEvent; import com.genersoft.iot.vmp.gb28181.event.platformKeepaliveExpire.PlatformKeepaliveExpireEvent;
import com.genersoft.iot.vmp.gb28181.event.platformNotRegister.PlatformNotRegisterEvent; import com.genersoft.iot.vmp.gb28181.event.platformNotRegister.PlatformNotRegisterEvent;
import com.genersoft.iot.vmp.media.zlm.event.ZLMOfflineEvent; import com.genersoft.iot.vmp.media.zlm.event.ZLMOfflineEvent;
import com.genersoft.iot.vmp.media.zlm.event.ZLMOnlineEvent;
import org.springframework.beans.factory.annotation.Autowired; import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.context.ApplicationEventPublisher; import org.springframework.context.ApplicationEventPublisher;
import org.springframework.stereotype.Component; import org.springframework.stereotype.Component;
@ -74,4 +75,9 @@ public class EventPublisher {
applicationEventPublisher.publishEvent(outEvent); applicationEventPublisher.publishEvent(outEvent);
} }
public void zlmOnlineEventPublish(String mediaServerId) {
ZLMOnlineEvent outEvent = new ZLMOnlineEvent(this);
outEvent.setMediaServerId(mediaServerId);
applicationEventPublisher.publishEvent(outEvent);
}
} }

View File

@ -179,29 +179,33 @@ public class ZLMHttpHookListener {
public ResponseEntity<String> onPublish(@RequestBody JSONObject json) { public ResponseEntity<String> onPublish(@RequestBody JSONObject json) {
logger.debug("[ ZLM HOOK ]on_publish API调用参数" + json.toString()); logger.debug("[ ZLM HOOK ]on_publish API调用参数" + json.toString());
JSONObject ret = new JSONObject();
ret.put("code", 0);
ret.put("msg", "success");
ret.put("enableHls", true);
ret.put("enableMP4", userSetup.isRecordPushLive());
String mediaServerId = json.getString("mediaServerId"); String mediaServerId = json.getString("mediaServerId");
ZLMHttpHookSubscribe.Event subscribe = this.subscribe.getSubscribe(ZLMHttpHookSubscribe.HookType.on_publish, json); ZLMHttpHookSubscribe.Event subscribe = this.subscribe.getSubscribe(ZLMHttpHookSubscribe.HookType.on_publish, json);
if (subscribe != null) { if (subscribe != null) {
MediaServerItem mediaInfo = mediaServerService.getOne(mediaServerId); MediaServerItem mediaInfo = mediaServerService.getOne(mediaServerId);
if (mediaInfo != null) { if (mediaInfo != null) {
subscribe.response(mediaInfo, json); subscribe.response(mediaInfo, json);
}else {
ret.put("code", 1);
ret.put("msg", "zlm not register");
} }
} }
String app = json.getString("app"); String app = json.getString("app");
String stream = json.getString("stream"); String stream = json.getString("stream");
StreamInfo streamInfo = redisCatchStorage.queryPlaybackByStreamId(stream); StreamInfo streamInfo = redisCatchStorage.queryPlaybackByStreamId(stream);
JSONObject ret = new JSONObject();
// 录像回放时不进行录像下载 // 录像回放时不进行录像下载
if (streamInfo != null) { if (streamInfo != null) {
ret.put("enableMP4", false); ret.put("enableMP4", false);
}else { }else {
ret.put("enableMP4", userSetup.isRecordPushLive()); ret.put("enableMP4", userSetup.isRecordPushLive());
} }
ret.put("code", 0);
ret.put("msg", "success");
ret.put("enableHls", true);
ret.put("enableMP4", userSetup.isRecordPushLive());
return new ResponseEntity<String>(ret.toString(), HttpStatus.OK); return new ResponseEntity<String>(ret.toString(), HttpStatus.OK);
} }
@ -340,6 +344,7 @@ public class ZLMHttpHookListener {
if (!"rtp".equals(app)){ if (!"rtp".equals(app)){
String type = OriginType.values()[item.getOriginType()].getType(); String type = OriginType.values()[item.getOriginType()].getType();
MediaServerItem mediaServerItem = mediaServerService.getOne(mediaServerId); MediaServerItem mediaServerItem = mediaServerService.getOne(mediaServerId);
if (mediaServerItem != null){
if (regist) { if (regist) {
StreamInfo streamInfo = mediaService.getStreamInfoByAppAndStream(mediaServerItem, app, streamId, tracks); StreamInfo streamInfo = mediaService.getStreamInfoByAppAndStream(mediaServerItem, app, streamId, tracks);
redisCatchStorage.addStream(mediaServerItem, type, app, streamId, streamInfo); redisCatchStorage.addStream(mediaServerItem, type, app, streamId, streamInfo);
@ -360,9 +365,8 @@ public class ZLMHttpHookListener {
} }
} }
zlmMediaListManager.removeMedia(app, streamId); zlmMediaListManager.removeMedia(app, streamId);
redisCatchStorage.removeStream(mediaServerItem, type, app, streamId); redisCatchStorage.removeStream(mediaServerItem.getId(), type, app, streamId);
} }
// 发送流变化redis消息 // 发送流变化redis消息
JSONObject jsonObject = new JSONObject(); JSONObject jsonObject = new JSONObject();
jsonObject.put("serverId", userSetup.getServerId()); jsonObject.put("serverId", userSetup.getServerId());
@ -374,6 +378,7 @@ public class ZLMHttpHookListener {
} }
} }
} }
}
JSONObject ret = new JSONObject(); JSONObject ret = new JSONObject();
ret.put("code", 0); ret.put("code", 0);

View File

@ -141,7 +141,6 @@ public class ZLMMediaListManager {
}else { }else {
gbStreamMapper.add(transform); gbStreamMapper.add(transform);
} }
} }
} }

View File

@ -4,6 +4,7 @@ import com.alibaba.fastjson.JSON;
import com.alibaba.fastjson.JSONArray; import com.alibaba.fastjson.JSONArray;
import com.alibaba.fastjson.JSONObject; import com.alibaba.fastjson.JSONObject;
import com.genersoft.iot.vmp.conf.MediaConfig; import com.genersoft.iot.vmp.conf.MediaConfig;
import com.genersoft.iot.vmp.gb28181.event.EventPublisher;
import com.genersoft.iot.vmp.media.zlm.dto.MediaServerItem; import com.genersoft.iot.vmp.media.zlm.dto.MediaServerItem;
import com.genersoft.iot.vmp.service.IMediaServerService; import com.genersoft.iot.vmp.service.IMediaServerService;
import com.genersoft.iot.vmp.service.IStreamProxyService; import com.genersoft.iot.vmp.service.IStreamProxyService;
@ -17,6 +18,7 @@ import org.springframework.core.annotation.Order;
import org.springframework.scheduling.annotation.Async; import org.springframework.scheduling.annotation.Async;
import org.springframework.scheduling.concurrent.ThreadPoolTaskExecutor; import org.springframework.scheduling.concurrent.ThreadPoolTaskExecutor;
import org.springframework.stereotype.Component; import org.springframework.stereotype.Component;
import org.springframework.util.StringUtils;
import java.util.*; import java.util.*;
@ -37,6 +39,9 @@ public class ZLMRunner implements CommandLineRunner {
@Autowired @Autowired
private IStreamProxyService streamProxyService; private IStreamProxyService streamProxyService;
@Autowired
private EventPublisher publisher;
@Autowired @Autowired
private IMediaServerService mediaServerService; private IMediaServerService mediaServerService;
@ -117,7 +122,7 @@ public class ZLMRunner implements CommandLineRunner {
@Async @Async
public void connectZlmServer(MediaServerItem mediaServerItem){ public void connectZlmServer(MediaServerItem mediaServerItem){
ZLMServerConfig zlmServerConfig = getMediaServerConfig(mediaServerItem); ZLMServerConfig zlmServerConfig = getMediaServerConfig(mediaServerItem, 1);
if (zlmServerConfig != null) { if (zlmServerConfig != null) {
zlmServerConfig.setIp(mediaServerItem.getIp()); zlmServerConfig.setIp(mediaServerItem.getIp());
zlmServerConfig.setHttpPort(mediaServerItem.getHttpPort()); zlmServerConfig.setHttpPort(mediaServerItem.getHttpPort());
@ -126,7 +131,7 @@ public class ZLMRunner implements CommandLineRunner {
} }
} }
public ZLMServerConfig getMediaServerConfig(MediaServerItem mediaServerItem) { public ZLMServerConfig getMediaServerConfig(MediaServerItem mediaServerItem, int index) {
if (startGetMedia == null) { return null;} if (startGetMedia == null) { return null;}
if (!mediaServerItem.isDefaultServer() && mediaServerService.getOne(mediaServerItem.getId()) == null) { if (!mediaServerItem.isDefaultServer() && mediaServerService.getOne(mediaServerItem.getId()) == null) {
return null; return null;
@ -143,14 +148,19 @@ public class ZLMRunner implements CommandLineRunner {
ZLMServerConfig.setIp(mediaServerItem.getIp()); ZLMServerConfig.setIp(mediaServerItem.getIp());
} }
} else { } else {
logger.error("[ {} ]-[ {}:{} ]主动连接失败失败, 2s后重试", logger.error("[ {} ]-[ {}:{} ]第{}次主动连接失败, 2s后重试",
mediaServerItem.getId(), mediaServerItem.getIp(), mediaServerItem.getHttpPort()); mediaServerItem.getId(), mediaServerItem.getIp(), mediaServerItem.getHttpPort(), index);
if (index == 1 && !StringUtils.isEmpty(mediaServerItem.getId())) {
logger.info("[ {} ]-[ {}:{} ]第{}次主动连接失败, 开始清理相关资源",
mediaServerItem.getId(), mediaServerItem.getIp(), mediaServerItem.getHttpPort(), index);
publisher.zlmOfflineEventPublish(mediaServerItem.getId());
}
try { try {
Thread.sleep(2000); Thread.sleep(2000);
} catch (InterruptedException e) { } catch (InterruptedException e) {
e.printStackTrace(); e.printStackTrace();
} }
ZLMServerConfig = getMediaServerConfig(mediaServerItem); ZLMServerConfig = getMediaServerConfig(mediaServerItem, index += 1);
} }
return ZLMServerConfig; return ZLMServerConfig;

View File

@ -0,0 +1,52 @@
package com.genersoft.iot.vmp.media.zlm.event;
import com.genersoft.iot.vmp.common.VideoManagerConstants;
import com.genersoft.iot.vmp.conf.UserSetup;
import com.genersoft.iot.vmp.gb28181.event.EventPublisher;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.data.redis.connection.Message;
import org.springframework.data.redis.listener.KeyExpirationEventMessageListener;
import org.springframework.data.redis.listener.RedisMessageListenerContainer;
import org.springframework.stereotype.Component;
/**
* @description:设备心跳超时监听,借助redis过期特性进行监听监听到说明设备心跳超时发送离线事件
* @author: swwheihei
* @date: 2020年5月6日 上午11:35:46
*/
@Component
public class ZLMKeepliveTimeoutListener extends KeyExpirationEventMessageListener {
private Logger logger = LoggerFactory.getLogger(ZLMKeepliveTimeoutListener.class);
@Autowired
private EventPublisher publisher;
@Autowired
private UserSetup userSetup;
public ZLMKeepliveTimeoutListener(RedisMessageListenerContainer listenerContainer) {
super(listenerContainer);
}
/**
* 监听失效的keykey格式为keeplive_deviceId
* @param message
* @param pattern
*/
@Override
public void onMessage(Message message, byte[] pattern) {
// 获取失效的key
String expiredKey = message.toString();
String KEEPLIVEKEY_PREFIX = VideoManagerConstants.MEDIA_SERVER_KEEPALIVE_PREFIX + userSetup.getServerId() + "_";
if(!expiredKey.startsWith(KEEPLIVEKEY_PREFIX)){
return;
}
String mediaServerId = expiredKey.substring(KEEPLIVEKEY_PREFIX.length(),expiredKey.length());
publisher.zlmOfflineEventPublish(mediaServerId);
}
}

View File

@ -0,0 +1,11 @@
package com.genersoft.iot.vmp.media.zlm.event;
/**
* zlm离线事件类
*/
public class ZLMOfflineEvent extends ZLMEventAbstract {
public ZLMOfflineEvent(Object source) {
super(source);
}
}

View File

@ -0,0 +1,44 @@
package com.genersoft.iot.vmp.media.zlm.event;
import com.genersoft.iot.vmp.conf.UserSetup;
import com.genersoft.iot.vmp.service.IMediaServerService;
import com.genersoft.iot.vmp.service.IStreamProxyService;
import com.genersoft.iot.vmp.service.IStreamPushService;
import com.genersoft.iot.vmp.storager.IVideoManagerStorager;
import com.genersoft.iot.vmp.utils.redis.RedisUtil;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.context.ApplicationListener;
import org.springframework.stereotype.Component;
/**
*
*/
@Component
public class ZLMOfflineEventListener implements ApplicationListener<ZLMOfflineEvent> {
private final static Logger logger = LoggerFactory.getLogger(ZLMOfflineEventListener.class);
@Autowired
private IMediaServerService mediaServerService;
@Autowired
private IStreamPushService streamPushService;
@Autowired
private IStreamProxyService streamProxyService;
@Override
public void onApplicationEvent(ZLMOfflineEvent event) {
if (logger.isDebugEnabled()) {
logger.debug("ZLM离线事件触发ID" + event.getMediaServerId());
}
// 处理ZLM离线
mediaServerService.zlmServerOffline(event.getMediaServerId());
streamProxyService.zlmServerOffline(event.getMediaServerId());
streamPushService.zlmServerOffline(event.getMediaServerId());
// TODO 处理对国标的影响
}
}

View File

@ -0,0 +1,11 @@
package com.genersoft.iot.vmp.media.zlm.event;
/**
* zlm在线事件
*/
public class ZLMOnlineEvent extends ZLMEventAbstract {
public ZLMOnlineEvent(Object source) {
super(source);
}
}

View File

@ -0,0 +1,65 @@
package com.genersoft.iot.vmp.media.zlm.event;
import com.genersoft.iot.vmp.conf.SipConfig;
import com.genersoft.iot.vmp.conf.UserSetup;
import com.genersoft.iot.vmp.service.IMediaServerService;
import com.genersoft.iot.vmp.service.IStreamProxyService;
import com.genersoft.iot.vmp.service.IStreamPushService;
import com.genersoft.iot.vmp.storager.IVideoManagerStorager;
import com.genersoft.iot.vmp.utils.redis.RedisUtil;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.context.ApplicationListener;
import org.springframework.stereotype.Component;
import java.text.SimpleDateFormat;
/**
* @description: 在线事件监听器监听到离线后修改设备离在线状态 设备在线有两个来源
* 1设备主动注销发送注销指令
* 2设备未知原因离线心跳超时
* @author: swwheihei
* @date: 2020年5月6日 下午1:51:23
*/
@Component
public class ZLMOnlineEventListener implements ApplicationListener<ZLMOnlineEvent> {
private final static Logger logger = LoggerFactory.getLogger(ZLMOnlineEventListener.class);
@Autowired
private IVideoManagerStorager storager;
@Autowired
private RedisUtil redis;
@Autowired
private SipConfig sipConfig;
@Autowired
private UserSetup userSetup;
@Autowired
private IMediaServerService mediaServerService;
@Autowired
private IStreamPushService streamPushService;
@Autowired
private IStreamProxyService streamProxyService;
private SimpleDateFormat format = new SimpleDateFormat("yyyy-MM-dd HH:mm:ss");
@Override
public void onApplicationEvent(ZLMOnlineEvent event) {
if (logger.isDebugEnabled()) {
logger.debug("ZLM上线事件触发ID" + event.getMediaServerId());
}
streamPushService.zlmServerOnline(event.getMediaServerId());
streamProxyService.zlmServerOnline(event.getMediaServerId());
}
}

View File

@ -78,10 +78,10 @@ public interface IStreamProxyService {
/** /**
* 新的节点加入 * 新的节点加入
* @param zlmServerConfig * @param mediaServerId
* @return * @return
*/ */
void zlmServerOnline(ZLMServerConfig zlmServerConfig); void zlmServerOnline(String mediaServerId);
/** /**
* 节点离线 * 节点离线
@ -89,4 +89,6 @@ public interface IStreamProxyService {
* @return * @return
*/ */
void zlmServerOffline(String mediaServerId); void zlmServerOffline(String mediaServerId);
void clean();
} }

View File

@ -34,6 +34,7 @@ public interface IStreamPushService {
* @return * @return
*/ */
PageInfo<StreamPushItem> getPushList(Integer page, Integer count); PageInfo<StreamPushItem> getPushList(Integer page, Integer count);
List<StreamPushItem> getPushList(String mediaSererId);
StreamPushItem transform(MediaItem item); StreamPushItem transform(MediaItem item);
@ -49,10 +50,10 @@ public interface IStreamPushService {
/** /**
* 新的节点加入 * 新的节点加入
* @param zlmServerConfig * @param mediaServerId
* @return * @return
*/ */
void zlmServerOnline(ZLMServerConfig zlmServerConfig); void zlmServerOnline(String mediaServerId);
/** /**
* 节点离线 * 节点离线
@ -61,4 +62,5 @@ public interface IStreamPushService {
*/ */
void zlmServerOffline(String mediaServerId); void zlmServerOffline(String mediaServerId);
void clean();
} }

View File

@ -4,10 +4,10 @@ import com.alibaba.fastjson.JSON;
import com.alibaba.fastjson.JSONArray; import com.alibaba.fastjson.JSONArray;
import com.alibaba.fastjson.JSONObject; import com.alibaba.fastjson.JSONObject;
import com.genersoft.iot.vmp.common.VideoManagerConstants; import com.genersoft.iot.vmp.common.VideoManagerConstants;
import com.genersoft.iot.vmp.conf.MediaConfig;
import com.genersoft.iot.vmp.conf.SipConfig; import com.genersoft.iot.vmp.conf.SipConfig;
import com.genersoft.iot.vmp.conf.UserSetup; import com.genersoft.iot.vmp.conf.UserSetup;
import com.genersoft.iot.vmp.gb28181.bean.Device; import com.genersoft.iot.vmp.gb28181.bean.Device;
import com.genersoft.iot.vmp.gb28181.event.EventPublisher;
import com.genersoft.iot.vmp.gb28181.session.SsrcConfig; import com.genersoft.iot.vmp.gb28181.session.SsrcConfig;
import com.genersoft.iot.vmp.gb28181.session.VideoStreamSessionManager; import com.genersoft.iot.vmp.gb28181.session.VideoStreamSessionManager;
import com.genersoft.iot.vmp.media.zlm.ZLMRESTfulUtils; import com.genersoft.iot.vmp.media.zlm.ZLMRESTfulUtils;
@ -70,6 +70,9 @@ public class MediaServerServiceImpl implements IMediaServerService, CommandLineR
@Autowired @Autowired
private RedisUtil redisUtil; private RedisUtil redisUtil;
@Autowired
private EventPublisher publisher;
@Autowired @Autowired
JedisUtil jedisUtil; JedisUtil jedisUtil;
@ -312,8 +315,6 @@ public class MediaServerServiceImpl implements IMediaServerService, CommandLineR
return mediaServerMapper.update(mediaSerItem); return mediaServerMapper.update(mediaSerItem);
} }
/** /**
* 处理zlm上线 * 处理zlm上线
* @param zlmServerConfig zlm上线携带的参数 * @param zlmServerConfig zlm上线携带的参数
@ -353,28 +354,31 @@ public class MediaServerServiceImpl implements IMediaServerService, CommandLineR
if (serverItem.getRtpProxyPort() == 0) { if (serverItem.getRtpProxyPort() == 0) {
serverItem.setRtpProxyPort(zlmServerConfig.getRtpProxyPort()); serverItem.setRtpProxyPort(zlmServerConfig.getRtpProxyPort());
} }
if (StringUtils.isEmpty(serverItem.getId())) {
serverItem.setId(zlmServerConfig.getGeneralMediaServerId());
}
serverItem.setStatus(true); serverItem.setStatus(true);
if (StringUtils.isEmpty(serverItem.getId())) { if (StringUtils.isEmpty(serverItem.getId())) {
serverItem.setId(zlmServerConfig.getGeneralMediaServerId()); serverItem.setId(zlmServerConfig.getGeneralMediaServerId());
mediaServerMapper.updateByHostAndPort(serverItem); mediaServerMapper.updateByHostAndPort(serverItem);
}else { }else {
mediaServerMapper.update(serverItem); mediaServerMapper.update(serverItem);
} }
String key = VideoManagerConstants.MEDIA_SERVER_PREFIX + userSetup.getServerId() + "_" + serverItem.getId(); String key = VideoManagerConstants.MEDIA_SERVER_PREFIX + userSetup.getServerId() + "_" + zlmServerConfig.getGeneralMediaServerId();
if (redisUtil.get(key) == null) { if (redisUtil.get(key) == null) {
SsrcConfig ssrcConfig = new SsrcConfig(serverItem.getId(), null, sipConfig.getDomain()); SsrcConfig ssrcConfig = new SsrcConfig(zlmServerConfig.getGeneralMediaServerId(), null, sipConfig.getDomain());
serverItem.setSsrcConfig(ssrcConfig); serverItem.setSsrcConfig(ssrcConfig);
redisUtil.set(key, serverItem); }else {
MediaServerItem mediaServerItemInRedis = (MediaServerItem)redisUtil.get(key);
serverItem.setSsrcConfig(mediaServerItemInRedis.getSsrcConfig());
} }
redisUtil.set(key, serverItem);
resetOnlineServerItem(serverItem); resetOnlineServerItem(serverItem);
updateMediaServerKeepalive(serverItem.getId(), null); updateMediaServerKeepalive(serverItem.getId(), null);
setZLMConfig(serverItem); setZLMConfig(serverItem);
publisher.zlmOnlineEventPublish(serverItem.getId());
} }
@Override @Override
public void zlmServerOffline(String mediaServerId) { public void zlmServerOffline(String mediaServerId) {
delete(mediaServerId); delete(mediaServerId);
@ -567,6 +571,10 @@ public class MediaServerServiceImpl implements IMediaServerService, CommandLineR
@Override @Override
public void updateMediaServerKeepalive(String mediaServerId, JSONObject data) { public void updateMediaServerKeepalive(String mediaServerId, JSONObject data) {
MediaServerItem mediaServerItem = getOne(mediaServerId); MediaServerItem mediaServerItem = getOne(mediaServerId);
if (mediaServerItem == null) {
logger.warn("[更新ZLM 保活信息]失败,未找到流媒体信息");
return;
}
String key = VideoManagerConstants.MEDIA_SERVER_KEEPALIVE_PREFIX + userSetup.getServerId() + "_" + mediaServerId; String key = VideoManagerConstants.MEDIA_SERVER_KEEPALIVE_PREFIX + userSetup.getServerId() + "_" + mediaServerId;
int hookAliveInterval = mediaServerItem.getHookAliveInterval() + 2; int hookAliveInterval = mediaServerItem.getHookAliveInterval() + 2;
redisUtil.set(key, data, hookAliveInterval); redisUtil.set(key, data, hookAliveInterval);

View File

@ -2,6 +2,7 @@ package com.genersoft.iot.vmp.service.impl;
import com.alibaba.fastjson.JSONObject; import com.alibaba.fastjson.JSONObject;
import com.genersoft.iot.vmp.common.StreamInfo; import com.genersoft.iot.vmp.common.StreamInfo;
import com.genersoft.iot.vmp.conf.UserSetup;
import com.genersoft.iot.vmp.gb28181.bean.GbStream; import com.genersoft.iot.vmp.gb28181.bean.GbStream;
import com.genersoft.iot.vmp.gb28181.bean.ParentPlatform; import com.genersoft.iot.vmp.gb28181.bean.ParentPlatform;
import com.genersoft.iot.vmp.media.zlm.ZLMRESTfulUtils; import com.genersoft.iot.vmp.media.zlm.ZLMRESTfulUtils;
@ -28,8 +29,7 @@ import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.stereotype.Service; import org.springframework.stereotype.Service;
import org.springframework.util.StringUtils; import org.springframework.util.StringUtils;
import java.util.ArrayList; import java.util.*;
import java.util.List;
/** /**
* 视频代理业务 * 视频代理业务
@ -54,6 +54,9 @@ public class StreamProxyServiceImpl implements IStreamProxyService {
@Autowired @Autowired
private IRedisCatchStorage redisCatchStorage; private IRedisCatchStorage redisCatchStorage;
@Autowired
private UserSetup userSetup;
@Autowired @Autowired
private GbStreamMapper gbStreamMapper; private GbStreamMapper gbStreamMapper;
@ -160,6 +163,9 @@ public class StreamProxyServiceImpl implements IStreamProxyService {
}else { }else {
mediaServerItem = mediaServerService.getOne(param.getMediaServerId()); mediaServerItem = mediaServerService.getOne(param.getMediaServerId());
} }
if (mediaServerItem == null) {
return null;
}
if ("default".equals(param.getType())){ if ("default".equals(param.getType())){
result = zlmresTfulUtils.addStreamProxy(mediaServerItem, param.getApp(), param.getStream(), param.getUrl(), result = zlmresTfulUtils.addStreamProxy(mediaServerItem, param.getApp(), param.getStream(), param.getUrl(),
param.isEnable_hls(), param.isEnable_mp4(), param.getRtp_type()); param.isEnable_hls(), param.isEnable_mp4(), param.getRtp_type());
@ -244,7 +250,6 @@ public class StreamProxyServiceImpl implements IStreamProxyService {
} }
} }
} }
return result; return result;
} }
@ -255,18 +260,41 @@ public class StreamProxyServiceImpl implements IStreamProxyService {
} }
@Override @Override
public void zlmServerOnline(ZLMServerConfig zlmServerConfig) { public void zlmServerOnline(String mediaServerId) {
zlmServerOffline(mediaServerId);
} }
@Override @Override
public void zlmServerOffline(String mediaServerId) { public void zlmServerOffline(String mediaServerId) {
// 移除开启了无人观看自动移除的流 // 移除开启了无人观看自动移除的流
List<StreamProxyItem> streamProxyItemList = streamProxyMapper.selecAutoRemoveItemByMediaServerId(mediaServerId);
if (streamProxyItemList.size() > 0) {
gbStreamMapper.batchDel(streamProxyItemList);
}
streamProxyMapper.deleteAutoRemoveItemByMediaServerId(mediaServerId); streamProxyMapper.deleteAutoRemoveItemByMediaServerId(mediaServerId);
// 其他的流设置未启用 // 其他的流设置未启用
streamProxyMapper.updateStatus(false, mediaServerId); streamProxyMapper.updateStatus(false, mediaServerId);
String type = "PULL";
// 发送redis消息
List<StreamInfo> streamInfoList = redisCatchStorage.getStreams(mediaServerId, type);
if (streamInfoList.size() > 0) {
for (StreamInfo streamInfo : streamInfoList) {
JSONObject jsonObject = new JSONObject();
jsonObject.put("serverId", userSetup.getServerId());
jsonObject.put("app", streamInfo.getApp());
jsonObject.put("stream", streamInfo.getStreamId());
jsonObject.put("register", false);
jsonObject.put("mediaServerId", mediaServerId);
redisCatchStorage.sendStreamChangeMsg(type, jsonObject);
// 移除redis内流的信息 // 移除redis内流的信息
redisCatchStorage.removeStream(mediaServerId, "PULL"); redisCatchStorage.removeStream(mediaServerId, type, streamInfo.getApp(), streamInfo.getStreamId());
}
}
}
@Override
public void clean() {
} }
} }

View File

@ -3,11 +3,15 @@ package com.genersoft.iot.vmp.service.impl;
import com.alibaba.fastjson.JSON; import com.alibaba.fastjson.JSON;
import com.alibaba.fastjson.JSONObject; import com.alibaba.fastjson.JSONObject;
import com.alibaba.fastjson.TypeReference; import com.alibaba.fastjson.TypeReference;
import com.genersoft.iot.vmp.common.StreamInfo;
import com.genersoft.iot.vmp.conf.UserSetup;
import com.genersoft.iot.vmp.gb28181.bean.GbStream; import com.genersoft.iot.vmp.gb28181.bean.GbStream;
import com.genersoft.iot.vmp.media.zlm.ZLMHttpHookSubscribe;
import com.genersoft.iot.vmp.media.zlm.ZLMRESTfulUtils; import com.genersoft.iot.vmp.media.zlm.ZLMRESTfulUtils;
import com.genersoft.iot.vmp.media.zlm.ZLMServerConfig; import com.genersoft.iot.vmp.media.zlm.ZLMServerConfig;
import com.genersoft.iot.vmp.media.zlm.dto.MediaItem; import com.genersoft.iot.vmp.media.zlm.dto.MediaItem;
import com.genersoft.iot.vmp.media.zlm.dto.MediaServerItem; import com.genersoft.iot.vmp.media.zlm.dto.MediaServerItem;
import com.genersoft.iot.vmp.media.zlm.dto.OriginType;
import com.genersoft.iot.vmp.media.zlm.dto.StreamPushItem; import com.genersoft.iot.vmp.media.zlm.dto.StreamPushItem;
import com.genersoft.iot.vmp.service.IMediaServerService; import com.genersoft.iot.vmp.service.IMediaServerService;
import com.genersoft.iot.vmp.service.IStreamPushService; import com.genersoft.iot.vmp.service.IStreamPushService;
@ -20,10 +24,7 @@ import com.github.pagehelper.PageInfo;
import org.springframework.beans.factory.annotation.Autowired; import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.stereotype.Service; import org.springframework.stereotype.Service;
import java.util.ArrayList; import java.util.*;
import java.util.HashMap;
import java.util.List;
import java.util.Map;
@Service @Service
public class StreamPushServiceImpl implements IStreamPushService { public class StreamPushServiceImpl implements IStreamPushService {
@ -43,6 +44,9 @@ public class StreamPushServiceImpl implements IStreamPushService {
@Autowired @Autowired
private IRedisCatchStorage redisCatchStorage; private IRedisCatchStorage redisCatchStorage;
@Autowired
private UserSetup userSetup;
@Autowired @Autowired
private IMediaServerService mediaServerService; private IMediaServerService mediaServerService;
@ -56,7 +60,9 @@ public class StreamPushServiceImpl implements IStreamPushService {
for (MediaItem item : mediaItems) { for (MediaItem item : mediaItems) {
// 不保存国标推理以及拉流代理的流 // 不保存国标推理以及拉流代理的流
if (item.getOriginType() == 1 || item.getOriginType() == 2 || item.getOriginType() == 8) { if (item.getOriginType() == OriginType.RTSP_PUSH.ordinal()
|| item.getOriginType() == OriginType.RTMP_PUSH.ordinal()
|| item.getOriginType() == OriginType.RTC_PUSH.ordinal() ) {
String key = item.getApp() + "_" + item.getStream(); String key = item.getApp() + "_" + item.getStream();
StreamPushItem streamPushItem = result.get(key); StreamPushItem streamPushItem = result.get(key);
if (streamPushItem == null) { if (streamPushItem == null) {
@ -97,6 +103,11 @@ public class StreamPushServiceImpl implements IStreamPushService {
return new PageInfo<>(all); return new PageInfo<>(all);
} }
@Override
public List<StreamPushItem> getPushList(String mediaServerId) {
return streamPushMapper.selectAllByMediaServerId(mediaServerId);
}
@Override @Override
public boolean saveToGB(GbStream stream) { public boolean saveToGB(GbStream stream) {
stream.setStreamType("push"); stream.setStreamType("push");
@ -137,17 +148,84 @@ public class StreamPushServiceImpl implements IStreamPushService {
} }
@Override @Override
public void zlmServerOnline(ZLMServerConfig zlmServerConfig) { public void zlmServerOnline(String mediaServerId) {
// 似乎没啥需要做的 // 同步zlm推流信息
MediaServerItem mediaServerItem = mediaServerService.getOne(mediaServerId);
if (mediaServerItem == null) {
return;
}
List<StreamPushItem> pushList = getPushList(mediaServerId);
if (pushList.size() > 0) {
Map<String, StreamPushItem> pushItemMap = new HashMap<>();
for (StreamPushItem streamPushItem : pushList) {
pushItemMap.put(streamPushItem.getApp() + streamPushItem.getStream(), streamPushItem);
}
zlmresTfulUtils.getMediaList(mediaServerItem, (mediaList ->{
if (mediaList == null) return;
String dataStr = mediaList.getString("data");
Integer code = mediaList.getInteger("code");
List<StreamPushItem> streamPushItems = null;
if (code == 0 ) {
if (dataStr != null) {
streamPushItems = handleJSON(dataStr, mediaServerItem);
}
}
if (streamPushItems != null) {
for (StreamPushItem streamPushItem : streamPushItems) {
pushItemMap.remove(streamPushItem.getApp() + streamPushItem.getStream());
}
}
Collection<StreamPushItem> offlinePushItems = pushItemMap.values();
if (offlinePushItems.size() > 0) {
String type = "PUSH";
streamPushMapper.delAll(new ArrayList<>(offlinePushItems));
for (StreamPushItem offlinePushItem : offlinePushItems) {
JSONObject jsonObject = new JSONObject();
jsonObject.put("serverId", userSetup.getServerId());
jsonObject.put("app", offlinePushItem.getApp());
jsonObject.put("stream", offlinePushItem.getStream());
jsonObject.put("register", false);
jsonObject.put("mediaServerId", mediaServerId);
redisCatchStorage.sendStreamChangeMsg(type, jsonObject);
// 移除redis内流的信息
redisCatchStorage.removeStream(mediaServerItem.getId(), "PUSH", offlinePushItem.getApp(), offlinePushItem.getStream());
}
}
}));
}
} }
@Override @Override
public void zlmServerOffline(String mediaServerId) { public void zlmServerOffline(String mediaServerId) {
// 移除没有serverId的推流 List<StreamPushItem> streamPushItems = streamPushMapper.selectAllByMediaServerId(mediaServerId);
// 移除没有GBId的推流
streamPushMapper.deleteWithoutGBId(mediaServerId); streamPushMapper.deleteWithoutGBId(mediaServerId);
gbStreamMapper.deleteWithoutGBId("push", mediaServerId);
// 其他的流设置未启用 // 其他的流设置未启用
gbStreamMapper.updateStatusByMediaServerId(mediaServerId, false); gbStreamMapper.updateStatusByMediaServerId(mediaServerId, false);
// 发送流停止消息
String type = "PUSH";
// 发送redis消息
List<StreamInfo> streamInfoList = redisCatchStorage.getStreams(mediaServerId, type);
if (streamInfoList.size() > 0) {
for (StreamInfo streamInfo : streamInfoList) {
JSONObject jsonObject = new JSONObject();
jsonObject.put("serverId", userSetup.getServerId());
jsonObject.put("app", streamInfo.getApp());
jsonObject.put("stream", streamInfo.getStreamId());
jsonObject.put("register", false);
jsonObject.put("mediaServerId", mediaServerId);
redisCatchStorage.sendStreamChangeMsg(type, jsonObject);
// 移除redis内流的信息 // 移除redis内流的信息
redisCatchStorage.removeStream(mediaServerId, "PUSH"); redisCatchStorage.removeStream(mediaServerId, type, streamInfo.getApp(), streamInfo.getStreamId());
}
}
}
@Override
public void clean() {
} }
} }

View File

@ -140,11 +140,11 @@ public interface IRedisCatchStorage {
/** /**
* 移除流信息从redis * 移除流信息从redis
* @param mediaServerItem * @param mediaServerId
* @param app * @param app
* @param streamId * @param streamId
*/ */
void removeStream(MediaServerItem mediaServerItem, String type, String app, String streamId); void removeStream(String mediaServerId, String type, String app, String streamId);
/** /**
@ -167,4 +167,6 @@ public interface IRedisCatchStorage {
* @return * @return
*/ */
ThirdPartyGB queryMemberNoGBId(String queryKey); ThirdPartyGB queryMemberNoGBId(String queryKey);
List<StreamInfo> getStreams(String mediaServerId, String pull);
} }

View File

@ -65,4 +65,18 @@ public interface GbStreamMapper {
"SET status=${status} " + "SET status=${status} " +
"WHERE mediaServerId=#{mediaServerId} ") "WHERE mediaServerId=#{mediaServerId} ")
void updateStatusByMediaServerId(String mediaServerId, boolean status); void updateStatusByMediaServerId(String mediaServerId, boolean status);
@Select("SELECT * FROM gb_stream WHERE mediaServerId=#{mediaServerId}")
void delByMediaServerId(String mediaServerId);
@Delete("DELETE FROM gb_stream WHERE streamType=#{type} AND gbId=NULL AND mediaServerId=#{mediaServerId}")
void deleteWithoutGBId(String type, String mediaServerId);
@Delete("<script> "+
"DELETE FROM gb_stream where " +
"<foreach collection='streamProxyItemList' item='item' separator='or'>" +
"(app=#{item.app} and stream=#{item.stream}) " +
"</foreach>" +
"</script>")
void batchDel(List<StreamProxyItem> streamProxyItemList);
} }

View File

@ -62,6 +62,9 @@ public interface StreamProxyMapper {
"WHERE mediaServerId=#{mediaServerId}") "WHERE mediaServerId=#{mediaServerId}")
void updateStatus(boolean status, String mediaServerId); void updateStatus(boolean status, String mediaServerId);
@Delete("DELETE FROM stream_proxy WHERE mediaServerId=#{mediaServerId}") @Delete("DELETE FROM stream_proxy WHERE enable_remove_none_reader=true AND mediaServerId=#{mediaServerId}")
void deleteAutoRemoveItemByMediaServerId(String mediaServerId); void deleteAutoRemoveItemByMediaServerId(String mediaServerId);
@Select("SELECT st.*, pgs.gbId, pgs.name, pgs.longitude, pgs.latitude FROM stream_proxy st LEFT JOIN gb_stream pgs on st.app = pgs.app AND st.stream = pgs.stream WHERE st.enable_remove_none_reader=true AND st.mediaServerId=#{mediaServerId} order by st.createTime desc")
List<StreamProxyItem> selecAutoRemoveItemByMediaServerId(String mediaServerId);
} }

View File

@ -4,6 +4,7 @@ import com.genersoft.iot.vmp.media.zlm.dto.StreamPushItem;
import org.apache.ibatis.annotations.*; import org.apache.ibatis.annotations.*;
import org.springframework.stereotype.Repository; import org.springframework.stereotype.Repository;
import java.util.Collection;
import java.util.List; import java.util.List;
@Mapper @Mapper
@ -31,6 +32,14 @@ public interface StreamPushMapper {
@Delete("DELETE FROM stream_push WHERE app=#{app} AND stream=#{stream}") @Delete("DELETE FROM stream_push WHERE app=#{app} AND stream=#{stream}")
int del(String app, String stream); int del(String app, String stream);
@Delete("<script> "+
"DELETE FROM stream_push where " +
"<foreach collection='streamPushItems' item='item' separator='or'>" +
"(app=#{item.app} and stream=#{item.stream}) " +
"</foreach>" +
"</script>")
int delAll(List<StreamPushItem> streamPushItems);
@Select("SELECT st.*, pgs.gbId, pgs.status, pgs.name, pgs.longitude, pgs.latitude FROM stream_push st LEFT JOIN gb_stream pgs on st.app = pgs.app AND st.stream = pgs.stream") @Select("SELECT st.*, pgs.gbId, pgs.status, pgs.name, pgs.longitude, pgs.latitude FROM stream_push st LEFT JOIN gb_stream pgs on st.app = pgs.app AND st.stream = pgs.stream")
List<StreamPushItem> selectAll(); List<StreamPushItem> selectAll();
@ -56,4 +65,7 @@ public interface StreamPushMapper {
@Delete("DELETE FROM stream_push WHERE mediaServerId=#{mediaServerId}") @Delete("DELETE FROM stream_push WHERE mediaServerId=#{mediaServerId}")
void deleteWithoutGBId(String mediaServerId); void deleteWithoutGBId(String mediaServerId);
@Select("SELECT * FROM stream_push WHERE mediaServerId=#{mediaServerId}")
List<StreamPushItem> selectAllByMediaServerId(String mediaServerId);
} }

View File

@ -338,8 +338,8 @@ public class RedisCatchStorageImpl implements IRedisCatchStorage {
} }
@Override @Override
public void removeStream(MediaServerItem mediaServerItem, String type, String app, String streamId) { public void removeStream(String mediaServerId, String type, String app, String streamId) {
String key = VideoManagerConstants.WVP_SERVER_STREAM_PREFIX + userSetup.getServerId() + "_" + type + "_" + app + "_" + streamId + "_" + mediaServerItem.getId(); String key = VideoManagerConstants.WVP_SERVER_STREAM_PREFIX + userSetup.getServerId() + "_" + type + "_" + app + "_" + streamId + "_" + mediaServerId;
redis.del(key); redis.del(key);
} }
@ -365,4 +365,16 @@ public class RedisCatchStorageImpl implements IRedisCatchStorage {
redis.del((String) stream); redis.del((String) stream);
} }
} }
@Override
public List<StreamInfo> getStreams(String mediaServerId, String type) {
List<StreamInfo> result = new ArrayList<>();
String key = VideoManagerConstants.WVP_SERVER_STREAM_PREFIX + userSetup.getServerId() + "_" + type + "_*_*_" + mediaServerId;
List<Object> streams = redis.scan(key);
for (Object stream : streams) {
StreamInfo streamInfo = (StreamInfo)redis.get((String) stream);
result.add(streamInfo);
}
return result;
}
} }