648540858
2020-10-26 9361943e47a09ea46f76adf06fa0d24a07ac711d
src/main/java/com/genersoft/iot/vmp/gb28181/transmit/request/impl/MessageRequestProcessor.java
@@ -2,15 +2,10 @@
import java.io.ByteArrayInputStream;
import java.text.ParseException;
import java.util.ArrayList;
import java.util.HashMap;
import java.util.Iterator;
import java.util.List;
import java.util.Map;
import java.util.*;
import javax.sip.InvalidArgumentException;
import javax.sip.RequestEvent;
import javax.sip.ServerTransaction;
import javax.sip.SipException;
import javax.sip.message.Request;
import javax.sip.message.Response;
@@ -22,10 +17,8 @@
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.stereotype.Component;
import com.genersoft.iot.vmp.common.VideoManagerConstants;
import com.genersoft.iot.vmp.gb28181.SipLayer;
import com.genersoft.iot.vmp.gb28181.bean.Device;
import com.genersoft.iot.vmp.gb28181.bean.DeviceChannel;
import com.genersoft.iot.vmp.gb28181.bean.RecordInfo;
@@ -35,49 +28,40 @@
import com.genersoft.iot.vmp.gb28181.transmit.callback.DeferredResultHolder;
import com.genersoft.iot.vmp.gb28181.transmit.callback.RequestMessage;
import com.genersoft.iot.vmp.gb28181.transmit.cmd.impl.SIPCommander;
import com.genersoft.iot.vmp.gb28181.transmit.request.ISIPRequestProcessor;
import com.genersoft.iot.vmp.gb28181.transmit.request.SIPRequestAbstractProcessor;
import com.genersoft.iot.vmp.gb28181.utils.DateUtil;
import com.genersoft.iot.vmp.gb28181.utils.XmlUtil;
import com.genersoft.iot.vmp.storager.IVideoManagerStorager;
import com.genersoft.iot.vmp.utils.redis.RedisUtil;
import org.springframework.util.StringUtils;
/**    
 * @Description:MESSAGE请求处理器
 * @author: songww
 * @author: swwheihei
 * @date:   2020年5月3日 下午5:32:41     
 */
@Component
public class MessageRequestProcessor implements ISIPRequestProcessor {
public class MessageRequestProcessor extends SIPRequestAbstractProcessor {
   
   private final static Logger logger = LoggerFactory.getLogger(MessageRequestProcessor.class);
   
   private ServerTransaction transaction;
   private SipLayer layer;
   @Autowired
   private SIPCommander cmder;
   
   @Autowired
   private IVideoManagerStorager storager;
   
   @Autowired
   private EventPublisher publisher;
   
   @Autowired
   private RedisUtil redis;
   
   @Autowired
   private DeferredResultHolder deferredResultHolder;
   
   @Autowired
   private DeviceOffLineDetector offLineDetector;
   
   private final static String CACHE_RECORDINFO_KEY = "CACHE_RECORDINFO_";
   
   private static final String MESSAGE_KEEP_ALIVE = "Keepalive";
   private static final String MESSAGE_CONFIG_DOWNLOAD = "ConfigDownload";
   private static final String MESSAGE_CATALOG = "Catalog";
   private static final String MESSAGE_DEVICE_INFO = "DeviceInfo";
   private static final String MESSAGE_KEEP_ALIVE = "Keepalive";
   private static final String MESSAGE_ALARM = "Alarm";
   private static final String MESSAGE_RECORD_INFO = "RecordInfo";
//   private static final String MESSAGE_BROADCAST = "Broadcast";
@@ -89,26 +73,19 @@
    * 处理MESSAGE请求
    *  
    * @param evt
    * @param layer
    * @param transaction
    */  
   @Override
   public void process(RequestEvent evt, SipLayer layer, ServerTransaction transaction) {
      this.layer = layer;
      this.transaction = transaction;
      Request request = evt.getRequest();
      SAXReader reader = new SAXReader();
      Document xml;
   public void process(RequestEvent evt) {
      try {
         xml = reader.read(new ByteArrayInputStream(request.getRawContent()));
         Element rootElement = xml.getRootElement();
         String cmd = rootElement.element("CmdType").getStringValue();
         Element rootElement = getRootElement(evt);
         String cmd = XmlUtil.getText(rootElement,"CmdType");
         if (MESSAGE_KEEP_ALIVE.equals(cmd)) {
            logger.info("接收到KeepAlive消息");
            processMessageKeepAlive(evt);
         } else if (MESSAGE_CONFIG_DOWNLOAD.equals(cmd)) {
            logger.info("接收到ConfigDownload消息");
         } else if (MESSAGE_CATALOG.equals(cmd)) {
            logger.info("接收到Catalog消息");
            processMessageCatalogList(evt);
@@ -125,7 +102,6 @@
      } catch (DocumentException e) {
         e.printStackTrace();
      }
   }
   
   /**
@@ -146,7 +122,10 @@
         device.setManufacturer(XmlUtil.getText(rootElement,"Manufacturer"));
         device.setModel(XmlUtil.getText(rootElement,"Model"));
         device.setFirmware(XmlUtil.getText(rootElement,"Firmware"));
         storager.update(device);
         if (StringUtils.isEmpty(device.getStreamMode())){
            device.setStreamMode("UDP");
         }
         storager.updateDevice(device);
         
         RequestMessage msg = new RequestMessage();
         msg.setDeviceId(deviceId);
@@ -166,7 +145,7 @@
      try {
         Element rootElement = getRootElement(evt);
         Element deviceIdElement = rootElement.element("DeviceID");
         String deviceId = deviceIdElement.getText().toString();
         String deviceId = deviceIdElement.getText();
         Element deviceListElement = rootElement.element("DeviceList");
         if (deviceListElement == null) {
            return;
@@ -177,11 +156,6 @@
            if (device == null) {
               return;
            }
            Map<String, DeviceChannel> channelMap = device.getChannelMap();
            if (channelMap == null) {
               channelMap = new HashMap<String, DeviceChannel>(5);
               device.setChannelMap(channelMap);
            }
            // 遍历DeviceList
            while (deviceListIterator.hasNext()) {
               Element itemDevice = deviceListIterator.next();
@@ -189,18 +163,18 @@
               if (channelDeviceElement == null) {
                  continue;
               }
               String channelDeviceId = channelDeviceElement.getText().toString();
               String channelDeviceId = channelDeviceElement.getText();
               Element channdelNameElement = itemDevice.element("Name");
               String channelName = channdelNameElement != null ? channdelNameElement.getText().toString() : "";
               String channelName = channdelNameElement != null ? channdelNameElement.getTextTrim().toString() : "";
               Element statusElement = itemDevice.element("Status");
               String status = statusElement != null ? statusElement.getText().toString() : "ON";
               DeviceChannel deviceChannel = channelMap.containsKey(channelDeviceId) ? channelMap.get(channelDeviceId) : new DeviceChannel();
               DeviceChannel deviceChannel = new DeviceChannel();
               deviceChannel.setName(channelName);
               deviceChannel.setChannelId(channelDeviceId);
               if(status.equals("ON")) {
               if(status.equals("ON") || status.equals("On")) {
                  deviceChannel.setStatus(1);
               }
               if(status.equals("OFF")) {
               if(status.equals("OFF") || status.equals("Off")) {
                  deviceChannel.setStatus(0);
               }
@@ -211,7 +185,7 @@
               deviceChannel.setBlock(XmlUtil.getText(itemDevice,"Block"));
               deviceChannel.setAddress(XmlUtil.getText(itemDevice,"Address"));
               deviceChannel.setParental(itemDevice.element("Parental") == null? 0:Integer.parseInt(XmlUtil.getText(itemDevice,"Parental")));
               deviceChannel.setParentId(XmlUtil.getText(itemDevice,"ParentId"));
               deviceChannel.setParentId(XmlUtil.getText(itemDevice,"ParentID"));
               deviceChannel.setSafetyWay(itemDevice.element("SafetyWay") == null? 0:Integer.parseInt(XmlUtil.getText(itemDevice,"SafetyWay")));
               deviceChannel.setRegisterWay(itemDevice.element("RegisterWay") == null? 1:Integer.parseInt(XmlUtil.getText(itemDevice,"RegisterWay")));
               deviceChannel.setCertNum(XmlUtil.getText(itemDevice,"CertNum"));
@@ -224,17 +198,25 @@
               deviceChannel.setPassword(XmlUtil.getText(itemDevice,"Password"));
               deviceChannel.setLongitude(itemDevice.element("Longitude") == null? 0.00:Double.parseDouble(XmlUtil.getText(itemDevice,"Longitude")));
               deviceChannel.setLatitude(itemDevice.element("Latitude") == null? 0.00:Double.parseDouble(XmlUtil.getText(itemDevice,"Latitude")));
               channelMap.put(channelDeviceId, deviceChannel);
               deviceChannel.setPTZType(itemDevice.element("PTZType") == null? 0:Integer.parseInt(XmlUtil.getText(itemDevice,"PTZType")));
               deviceChannel.setHasAudio(true); // 默认含有音频,播放时再检查是否有音频及是否AAC
               storager.updateChannel(device.getDeviceId(), deviceChannel);
            }
            // 更新
            storager.update(device);
            RequestMessage msg = new RequestMessage();
            msg.setDeviceId(deviceId);
            msg.setType(DeferredResultHolder.CALLBACK_CMD_CATALOG);
            msg.setData(device);
            deferredResultHolder.invokeResult(msg);
            // 回复200
            if (offLineDetector.isOnline(deviceId)) {
               responseAck(evt);
               publisher.onlineEventPublish(deviceId, VideoManagerConstants.EVENT_ONLINE_KEEPLIVE);
            }
         }
      } catch (DocumentException e) {
      } catch (DocumentException | SipException | InvalidArgumentException | ParseException e) {
         e.printStackTrace();
      }
   }
@@ -251,13 +233,18 @@
         
         Device device = storager.queryVideoDevice(deviceId);
         if (device == null) {
            // TODO 也可能是通道
//            storager.queryChannel(deviceId)
            return;
         }
         device.setName(XmlUtil.getText(rootElement,"DeviceName"));
         device.setManufacturer(XmlUtil.getText(rootElement,"Manufacturer"));
         device.setModel(XmlUtil.getText(rootElement,"Model"));
         device.setFirmware(XmlUtil.getText(rootElement,"Firmware"));
         storager.update(device);
         if (StringUtils.isEmpty(device.getStreamMode())){
            device.setStreamMode("UDP");
         }
         storager.updateDevice(device);
         cmder.catalogQuery(device);
      } catch (DocumentException e) {
         e.printStackTrace();
@@ -272,15 +259,11 @@
      try {
         Element rootElement = getRootElement(evt);
         String deviceId = XmlUtil.getText(rootElement,"DeviceID");
         Request request = evt.getRequest();
         Response response = null;
         if (offLineDetector.isOnline(deviceId)) {
            response = layer.getMessageFactory().createResponse(Response.OK,request);
            responseAck(evt);
            publisher.onlineEventPublish(deviceId, VideoManagerConstants.EVENT_ONLINE_KEEPLIVE);
         } else {
            response = layer.getMessageFactory().createResponse(Response.BAD_REQUEST,request);
         }
         transaction.sendResponse(response);
      } catch (ParseException | SipException | InvalidArgumentException | DocumentException e) {
         e.printStackTrace();
      }
@@ -326,9 +309,10 @@
               record.setEndTime(DateUtil.ISO8601Toyyyy_MM_dd_HH_mm_ss(XmlUtil.getText(itemRecord,"EndTime")));
               record.setSecrecy(itemRecord.element("Secrecy") == null? 0:Integer.parseInt(XmlUtil.getText(itemRecord,"Secrecy")));
               record.setType(XmlUtil.getText(itemRecord,"Type"));
               record.setRecordId(XmlUtil.getText(itemRecord,"RecorderID"));
               record.setRecorderId(XmlUtil.getText(itemRecord,"RecorderID"));
               recordList.add(record);
            }
//            recordList.sort(Comparator.naturalOrder());
            recordInfo.setRecordList(recordList);
         }
         
@@ -362,9 +346,13 @@
         // 走到这里,有以下可能:1、没有录像信息,第一次收到recordinfo的消息即返回响应数据,无redis操作
         //               2、有录像数据,且第一次即收到完整数据,返回响应数据,无redis操作
         //                3、有录像数据,在超时时间内收到多次包组装后数量足够,返回数据
         // 对记录进行排序
         RequestMessage msg = new RequestMessage();
         msg.setDeviceId(deviceId);
         msg.setType(DeferredResultHolder.CALLBACK_CMD_RECORDINFO);
         // 自然顺序排序, 元素进行升序排列
         recordInfo.getRecordList().sort(Comparator.naturalOrder());
         msg.setData(recordInfo);
         deferredResultHolder.invokeResult(msg);
      } catch (DocumentException e) {
@@ -372,12 +360,41 @@
      }
   }
   
   private void responseAck(RequestEvent evt) throws SipException, InvalidArgumentException, ParseException {
      Response response = getMessageFactory().createResponse(Response.OK,evt.getRequest());
      getServerTransaction(evt).sendResponse(response);
   }
   private Element getRootElement(RequestEvent evt) throws DocumentException {
      Request request = evt.getRequest();
      SAXReader reader = new SAXReader();
      reader.setEncoding("GB2312");
      reader.setEncoding("gbk");
      Document xml = reader.read(new ByteArrayInputStream(request.getRawContent()));
      return xml.getRootElement();
   }
   public void setCmder(SIPCommander cmder) {
      this.cmder = cmder;
   }
   public void setStorager(IVideoManagerStorager storager) {
      this.storager = storager;
   }
   public void setPublisher(EventPublisher publisher) {
      this.publisher = publisher;
   }
   public void setRedis(RedisUtil redis) {
      this.redis = redis;
   }
   public void setDeferredResultHolder(DeferredResultHolder deferredResultHolder) {
      this.deferredResultHolder = deferredResultHolder;
   }
   public void setOffLineDetector(DeviceOffLineDetector offLineDetector) {
      this.offLineDetector = offLineDetector;
   }
}