From 8ee7211ba3f5dd53901e1be73db7c7c2199fe12e Mon Sep 17 00:00:00 2001
From: 648540858 <648540858@qq.com>
Date: 星期五, 22 七月 2022 16:05:23 +0800
Subject: [PATCH] 更新mysql.sql
---
src/main/java/com/genersoft/iot/vmp/service/impl/StreamProxyServiceImpl.java | 198 ++++++++++++++++++++++++++++++++++++++-----------
1 files changed, 154 insertions(+), 44 deletions(-)
diff --git a/src/main/java/com/genersoft/iot/vmp/service/impl/StreamProxyServiceImpl.java b/src/main/java/com/genersoft/iot/vmp/service/impl/StreamProxyServiceImpl.java
index 564deb5..40c37c2 100644
--- a/src/main/java/com/genersoft/iot/vmp/service/impl/StreamProxyServiceImpl.java
+++ b/src/main/java/com/genersoft/iot/vmp/service/impl/StreamProxyServiceImpl.java
@@ -1,38 +1,41 @@
package com.genersoft.iot.vmp.service.impl;
+import com.alibaba.fastjson.JSONArray;
import com.alibaba.fastjson.JSONObject;
import com.genersoft.iot.vmp.common.StreamInfo;
-import com.genersoft.iot.vmp.conf.SipConfig;
-import com.genersoft.iot.vmp.conf.UserSetup;
-import com.genersoft.iot.vmp.gb28181.bean.DeviceChannel;
+import com.genersoft.iot.vmp.conf.UserSetting;
import com.genersoft.iot.vmp.gb28181.bean.GbStream;
import com.genersoft.iot.vmp.gb28181.bean.ParentPlatform;
+import com.genersoft.iot.vmp.gb28181.bean.TreeType;
import com.genersoft.iot.vmp.gb28181.event.EventPublisher;
import com.genersoft.iot.vmp.gb28181.event.subscribe.catalog.CatalogEvent;
import com.genersoft.iot.vmp.media.zlm.ZLMRESTfulUtils;
-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.MediaServerItem;
import com.genersoft.iot.vmp.media.zlm.dto.StreamProxyItem;
-import com.genersoft.iot.vmp.media.zlm.dto.StreamPushItem;
import com.genersoft.iot.vmp.service.IGbStreamService;
import com.genersoft.iot.vmp.service.IMediaServerService;
import com.genersoft.iot.vmp.service.IMediaService;
import com.genersoft.iot.vmp.storager.IRedisCatchStorage;
-import com.genersoft.iot.vmp.storager.IVideoManagerStorager;
+import com.genersoft.iot.vmp.storager.IVideoManagerStorage;
import com.genersoft.iot.vmp.storager.dao.GbStreamMapper;
import com.genersoft.iot.vmp.storager.dao.ParentPlatformMapper;
import com.genersoft.iot.vmp.storager.dao.PlatformGbStreamMapper;
import com.genersoft.iot.vmp.storager.dao.StreamProxyMapper;
import com.genersoft.iot.vmp.service.IStreamProxyService;
+import com.genersoft.iot.vmp.utils.DateUtil;
import com.genersoft.iot.vmp.vmanager.bean.WVPResult;
import com.github.pagehelper.PageInfo;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.beans.factory.annotation.Autowired;
+import org.springframework.jdbc.datasource.DataSourceTransactionManager;
import org.springframework.stereotype.Service;
+import org.springframework.transaction.TransactionDefinition;
+import org.springframework.transaction.TransactionStatus;
import org.springframework.util.StringUtils;
+import java.net.InetAddress;
import java.util.*;
/**
@@ -44,13 +47,13 @@
private final static Logger logger = LoggerFactory.getLogger(StreamProxyServiceImpl.class);
@Autowired
- private IVideoManagerStorager videoManagerStorager;
+ private IVideoManagerStorage videoManagerStorager;
@Autowired
private IMediaService mediaService;
@Autowired
- private ZLMRESTfulUtils zlmresTfulUtils;;
+ private ZLMRESTfulUtils zlmresTfulUtils;
@Autowired
private StreamProxyMapper streamProxyMapper;
@@ -59,13 +62,10 @@
private IRedisCatchStorage redisCatchStorage;
@Autowired
- private IVideoManagerStorager storager;
+ private IVideoManagerStorage storager;
@Autowired
- private UserSetup userSetup;
-
- @Autowired
- private SipConfig sipConfig;
+ private UserSetting userSetting;
@Autowired
private GbStreamMapper gbStreamMapper;
@@ -85,13 +85,19 @@
@Autowired
private IMediaServerService mediaServerService;
+ @Autowired
+ DataSourceTransactionManager dataSourceTransactionManager;
+
+ @Autowired
+ TransactionDefinition transactionDefinition;
+
@Override
public WVPResult<StreamInfo> save(StreamProxyItem param) {
MediaServerItem mediaInfo;
WVPResult<StreamInfo> wvpResult = new WVPResult<>();
wvpResult.setCode(0);
- if ("auto".equals(param.getMediaServerId())){
+ if (param.getMediaServerId() == null || "auto".equals(param.getMediaServerId())){
mediaInfo = mediaServerService.getMediaServerForMinimumLoad();
}else {
mediaInfo = mediaServerService.getOne(param.getMediaServerId());
@@ -101,6 +107,7 @@
wvpResult.setMsg("淇濆瓨澶辫触");
return wvpResult;
}
+
String dstUrl = String.format("rtmp://%s:%s/%s/%s", "127.0.0.1", mediaInfo.getRtmpPort(), param.getApp(),
param.getStream() );
param.setDst_url(dstUrl);
@@ -110,9 +117,9 @@
boolean saveResult;
// 鏇存柊
if (videoManagerStorager.queryStreamProxy(param.getApp(), param.getStream()) != null) {
- saveResult = videoManagerStorager.updateStreamProxy(param);
+ saveResult = updateStreamProxy(param);
}else { // 鏂板
- saveResult = videoManagerStorager.addStreamProxy(param);
+ saveResult = addStreamProxy(param);
}
if (saveResult) {
result.append("淇濆瓨鎴愬姛");
@@ -126,7 +133,7 @@
if (param.isEnable_remove_none_reader()) {
del(param.getApp(), param.getStream());
}else {
- videoManagerStorager.updateStreamProxy(param);
+ updateStreamProxy(param);
}
}else {
@@ -149,25 +156,79 @@
result.append(", 鍏宠仈鍥芥爣骞冲彴[ " + param.getPlatformGbId() + " ]澶辫触");
}
}
- if (!StringUtils.isEmpty(param.getGbId())) {
- // 鏌ユ壘寮�鍚簡鍏ㄩ儴鐩存挱娴佸叡浜殑涓婄骇骞冲彴
- List<ParentPlatform> parentPlatforms = parentPlatformMapper.selectAllAhareAllLiveStream();
- if (parentPlatforms.size() > 0) {
- for (ParentPlatform parentPlatform : parentPlatforms) {
- param.setPlatformId(parentPlatform.getServerGBId());
- param.setCatalogId(parentPlatform.getCatalogId());
- String stream = param.getStream();
- StreamProxyItem streamProxyItems = platformGbStreamMapper.selectOne(param.getApp(), stream, parentPlatform.getServerGBId());
- if (streamProxyItems == null) {
- platformGbStreamMapper.add(param);
- eventPublisher.catalogEventPublishForStream(parentPlatform.getServerGBId(), param, CatalogEvent.ADD);
- }
- }
- }
- }
-
wvpResult.setMsg(result.toString());
return wvpResult;
+ }
+
+ /**
+ * 鏂板浠g悊娴�
+ * @param streamProxyItem
+ * @return
+ */
+ private boolean addStreamProxy(StreamProxyItem streamProxyItem) {
+ TransactionStatus transactionStatus = dataSourceTransactionManager.getTransaction(transactionDefinition);
+ boolean result = false;
+ streamProxyItem.setStreamType("proxy");
+ streamProxyItem.setStatus(true);
+ String now = DateUtil.getNow();
+ streamProxyItem.setCreateTime(now);
+ try {
+ if (streamProxyMapper.add(streamProxyItem) > 0) {
+ if (!StringUtils.isEmpty(streamProxyItem.getGbId())) {
+ if (gbStreamMapper.add(streamProxyItem) < 0) {
+ //浜嬪姟鍥炴粴
+ dataSourceTransactionManager.rollback(transactionStatus);
+ return false;
+ }
+ }
+ }else {
+ //浜嬪姟鍥炴粴
+ dataSourceTransactionManager.rollback(transactionStatus);
+ return false;
+ }
+ result = true;
+ dataSourceTransactionManager.commit(transactionStatus); //鎵嬪姩鎻愪氦
+ }catch (Exception e) {
+ logger.error("鍚戞暟鎹簱娣诲姞娴佷唬鐞嗗け璐ワ細", e);
+ dataSourceTransactionManager.rollback(transactionStatus);
+ }
+
+
+ return result;
+ }
+
+ /**
+ * 鏇存柊浠g悊娴�
+ * @param streamProxyItem
+ * @return
+ */
+ @Override
+ public boolean updateStreamProxy(StreamProxyItem streamProxyItem) {
+ TransactionStatus transactionStatus = dataSourceTransactionManager.getTransaction(transactionDefinition);
+ boolean result = false;
+ streamProxyItem.setStreamType("proxy");
+ try {
+ if (streamProxyMapper.update(streamProxyItem) > 0) {
+ if (!StringUtils.isEmpty(streamProxyItem.getGbId())) {
+ if (gbStreamMapper.updateByAppAndStream(streamProxyItem) == 0) {
+ //浜嬪姟鍥炴粴
+ dataSourceTransactionManager.rollback(transactionStatus);
+ return false;
+ }
+ }
+ } else {
+ //浜嬪姟鍥炴粴
+ dataSourceTransactionManager.rollback(transactionStatus);
+ return false;
+ }
+
+ dataSourceTransactionManager.commit(transactionStatus); //鎵嬪姩鎻愪氦
+ result = true;
+ }catch (Exception e) {
+ e.printStackTrace();
+ dataSourceTransactionManager.rollback(transactionStatus);
+ }
+ return result;
}
@Override
@@ -196,7 +257,9 @@
@Override
public JSONObject removeStreamProxyFromZlm(StreamProxyItem param) {
- if (param ==null) return null;
+ if (param ==null) {
+ return null;
+ }
MediaServerItem mediaServerItem = mediaServerService.getOne(param.getMediaServerId());
JSONObject result = zlmresTfulUtils.closeStreams(mediaServerItem, param.getApp(), param.getStream());
return result;
@@ -230,13 +293,16 @@
public boolean start(String app, String stream) {
boolean result = false;
StreamProxyItem streamProxy = videoManagerStorager.queryStreamProxy(app, stream);
- if (!streamProxy.isEnable() && streamProxy != null) {
+ if (!streamProxy.isEnable() ) {
JSONObject jsonObject = addStreamProxyToZlm(streamProxy);
- if (jsonObject == null) return false;
+ if (jsonObject == null) {
+ return false;
+ }
+ System.out.println(jsonObject);
if (jsonObject.getInteger("code") == 0) {
result = true;
streamProxy.setEnable(true);
- videoManagerStorager.updateStreamProxy(streamProxy);
+ updateStreamProxy(streamProxy);
}
}
return result;
@@ -248,9 +314,9 @@
StreamProxyItem streamProxyDto = videoManagerStorager.queryStreamProxy(app, stream);
if (streamProxyDto != null && streamProxyDto.isEnable()) {
JSONObject jsonObject = removeStreamProxyFromZlm(streamProxyDto);
- if (jsonObject.getInteger("code") == 0) {
+ if (jsonObject != null && jsonObject.getInteger("code") == 0) {
streamProxyDto.setEnable(false);
- result = videoManagerStorager.updateStreamProxy(streamProxyDto);
+ result = updateStreamProxy(streamProxyDto);
}
}
return result;
@@ -288,9 +354,12 @@
}
streamProxyMapper.deleteAutoRemoveItemByMediaServerId(mediaServerId);
+ // 绉婚櫎鎷夋祦浠g悊鐢熸垚鐨勬祦淇℃伅
+// syncPullStream(mediaServerId);
+
// 鎭㈠娴佷唬鐞�, 鍙煡鎵捐繖涓繖涓祦濯掍綋
List<StreamProxyItem> streamProxyListForEnable = storager.getStreamProxyListForEnableInMediaServer(
- mediaServerId, true, false);
+ mediaServerId, true);
for (StreamProxyItem streamProxyDto : streamProxyListForEnable) {
logger.info("鎭㈠娴佷唬鐞嗭紝" + streamProxyDto.getApp() + "/" + streamProxyDto.getStream());
JSONObject jsonObject = addStreamProxyToZlm(streamProxyDto);
@@ -313,7 +382,7 @@
}
streamProxyMapper.deleteAutoRemoveItemByMediaServerId(mediaServerId);
// 鍏朵粬鐨勬祦璁剧疆绂荤嚎
- streamProxyMapper.updateStatusByMediaServerId(false, mediaServerId);
+ streamProxyMapper.updateStatusByMediaServerId(mediaServerId, false);
String type = "PULL";
// 鍙戦�乺edis娑堟伅
@@ -321,7 +390,7 @@
if (mediaItems.size() > 0) {
for (MediaItem mediaItem : mediaItems) {
JSONObject jsonObject = new JSONObject();
- jsonObject.put("serverId", userSetup.getServerId());
+ jsonObject.put("serverId", userSetting.getServerId());
jsonObject.put("app", mediaItem.getApp());
jsonObject.put("stream", mediaItem.getStream());
jsonObject.put("register", false);
@@ -340,6 +409,47 @@
@Override
public int updateStatus(boolean status, String app, String stream) {
- return streamProxyMapper.updateStatus(status, app, stream);
+ return streamProxyMapper.updateStatus(app, stream, status);
+ }
+
+ private void syncPullStream(String mediaServerId){
+ MediaServerItem mediaServer = mediaServerService.getOne(mediaServerId);
+ if (mediaServer != null) {
+ List<MediaItem> allPullStream = redisCatchStorage.getStreams(mediaServerId, "PULL");
+ if (allPullStream.size() > 0) {
+ zlmresTfulUtils.getMediaList(mediaServer, jsonObject->{
+ Map<String, StreamInfo> stringStreamInfoMap = new HashMap<>();
+ if (jsonObject.getInteger("code") == 0) {
+ JSONArray data = jsonObject.getJSONArray("data");
+ if(data != null && data.size() > 0) {
+ for (int i = 0; i < data.size(); i++) {
+ JSONObject streamJSONObj = data.getJSONObject(i);
+ if ("rtmp".equals(streamJSONObj.getString("schema"))) {
+ StreamInfo streamInfo = new StreamInfo();
+ String app = streamJSONObj.getString("app");
+ String stream = streamJSONObj.getString("stream");
+ streamInfo.setApp(app);
+ streamInfo.setStream(stream);
+ stringStreamInfoMap.put(app+stream, streamInfo);
+ }
+ }
+ }
+ }
+ if (stringStreamInfoMap.size() == 0) {
+ redisCatchStorage.removeStream(mediaServerId, "PULL");
+ }else {
+ for (String key : stringStreamInfoMap.keySet()) {
+ StreamInfo streamInfo = stringStreamInfoMap.get(key);
+ if (stringStreamInfoMap.get(streamInfo.getApp() + streamInfo.getStream()) == null) {
+ redisCatchStorage.removeStream(mediaServerId, "PULL", streamInfo.getApp(),
+ streamInfo.getStream());
+ }
+ }
+ }
+ });
+ }
+
+ }
+
}
}
--
Gitblit v1.8.0