liusuyi
2026-06-01 a2e7e8ff9cfaa69b001d483710bddbda50d55c91
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
package com.ard.gb28181.transmit.event.request.impl.message.response.cmd;
 
import com.ard.gb28181.api.bean.MessageResponseTask;
import com.ard.gb28181.api.bean.Preset;
import com.ard.gb28181.api.domain.Device;
import com.ard.gb28181.transmit.event.request.SIPRequestProcessorParent;
import com.ard.gb28181.transmit.event.request.impl.message.IMessageHandler;
import com.ard.gb28181.transmit.event.request.impl.message.response.ResponseMessageHandler;
import gov.nist.javax.sip.message.SIPRequest;
import lombok.extern.slf4j.Slf4j;
import org.dom4j.DocumentException;
import org.dom4j.Element;
import org.springframework.beans.factory.InitializingBean;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.scheduling.annotation.Scheduled;
import org.springframework.stereotype.Component;
 
import javax.sip.InvalidArgumentException;
import javax.sip.RequestEvent;
import javax.sip.SipException;
import javax.sip.message.Response;
import java.text.ParseException;
import java.util.ArrayList;
import java.util.Iterator;
import java.util.List;
import java.util.Map;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.DelayQueue;
import java.util.concurrent.TimeUnit;
 
import static com.ard.gb28181.api.utils.XmlUtil.getText;
 
/**
 * 设备预置位查询应答
 */
@Slf4j
@Component
public class PresetQueryResponseMessageHandler extends SIPRequestProcessorParent implements InitializingBean, IMessageHandler {
 
    private final String cmdType = "PresetQuery";
 
    @Autowired
    private ResponseMessageHandler responseMessageHandler;
 
    private final Map<String, MessageResponseTask<Preset>> mesageMap = new ConcurrentHashMap<>();
 
    private final DelayQueue<MessageResponseTask<Preset>> delayQueue = new DelayQueue<>();
 
 
    @Override
    public void afterPropertiesSet() throws Exception {
        responseMessageHandler.addHandler(cmdType, this);
    }
 
    @Override
    public void handForDevice(RequestEvent evt, Device device, Element element) {
 
        SIPRequest request = (SIPRequest) evt.getRequest();
 
        try {
            Element rootElement = getRootElement(evt, device.getCharset());
 
            if (rootElement == null) {
                log.warn("[ 设备预置位查询应答 ] content cannot be null, {}", evt.getRequest());
                try {
                    responseAck(request, Response.BAD_REQUEST);
                } catch (InvalidArgumentException | ParseException | SipException e) {
                    log.error("[命令发送失败] 设备预置位查询应答处理: {}", e.getMessage());
                }
                return;
            }
            Element presetListNumElement = rootElement.element("PresetList");
            Element snElement = rootElement.element("SN");
            //该字段可能为通道或则设备的id
            if (snElement == null || presetListNumElement == null) {
                try {
                    responseAck(request, Response.BAD_REQUEST, "xml error");
                } catch (InvalidArgumentException | ParseException | SipException e) {
                    log.error("[命令发送失败] 设备预置位查询应答处理: {}", e.getMessage());
                }
                return;
            }
            int num = Integer.parseInt(presetListNumElement.attributeValue("Num"));
            List<Preset> presetQuerySipReqList = new ArrayList<>();
            if (num > 0) {
                for (Iterator<Element> presetIterator = presetListNumElement.elementIterator(); presetIterator.hasNext(); ) {
                    Element itemListElement = presetIterator.next();
                    Preset presetQuerySipReq = new Preset();
                    for (Iterator<Element> itemListIterator = itemListElement.elementIterator(); itemListIterator.hasNext(); ) {
                        // 遍历item
                        Element itemOne = itemListIterator.next();
                        String name = itemOne.getName();
                        String textTrim = itemOne.getTextTrim();
                        if ("PresetID".equalsIgnoreCase(name)) {
                            presetQuerySipReq.setPresetId(textTrim);
                        } else {
                            presetQuerySipReq.setPresetName(textTrim);
                        }
                    }
                    presetQuerySipReqList.add(presetQuerySipReq);
                }
            }
            String sn = getText(element, "SN");
            addCatch(cmdType + "_" + sn, num, rootElement, presetQuerySipReqList);
            try {
                responseAck(request, Response.OK);
            } catch (InvalidArgumentException | ParseException | SipException e) {
                log.error("[命令发送失败] 设备预置位查询应答处理: {}", e.getMessage());
            }
        } catch (DocumentException e) {
            log.error("[解析xml]失败: ", e);
        }
    }
 
    private void addCatch(String key, int sumNum, Element rootElement, List<Preset> presetQuerySipReqList) {
        if (presetQuerySipReqList.size() == sumNum) {
            responseMessageHandler.handMessageEvent(rootElement, presetQuerySipReqList);
            if (mesageMap.containsKey(key)) {
                MessageResponseTask<Preset> messageResponseTask = mesageMap.get(key);
                mesageMap.remove(key);
                boolean remove = delayQueue.remove(messageResponseTask);
                if (!remove) {
                    log.info("[移除预置位查询任务] 从延时队列内移除失败: {}", key);
                }
            }
        } else {
            if (mesageMap.containsKey(key)) {
                MessageResponseTask<Preset> messageResponseTask = mesageMap.get(key);
                List<Preset> data = messageResponseTask.getData();
                data.addAll(presetQuerySipReqList);
                if (data.size() == sumNum) {
                    responseMessageHandler.handMessageEvent(rootElement, data);
                    mesageMap.remove(key);
                    boolean remove = delayQueue.remove(messageResponseTask);
                    if (!remove) {
                        log.info("[移除预置位查询任务] 从延时队列内移除失败: {}", key);
                    }
                    return;
                }
                messageResponseTask.setDelayTime(System.currentTimeMillis() + 1000);
            } else {
                MessageResponseTask<Preset> messageResponseTask = new MessageResponseTask<>();
                messageResponseTask.setElement(rootElement);
                messageResponseTask.setData(presetQuerySipReqList);
                messageResponseTask.setDelayTime(System.currentTimeMillis() + 1000);
                messageResponseTask.setKey(key);
                mesageMap.put(key, messageResponseTask);
                delayQueue.offer(messageResponseTask);
            }
        }
    }
 
    // 处理过期的缓存
    @Scheduled(fixedDelay = 500, timeUnit = TimeUnit.MILLISECONDS)
    public void expirationCheck() {
        while (!delayQueue.isEmpty()) {
            MessageResponseTask<Preset> take = null;
            try {
                take = delayQueue.take();
                try {
                    responseMessageHandler.handMessageEvent(take.getElement(), take.getData());
                    mesageMap.remove(take.getKey());
                } catch (Exception e) {
                    log.error("[预置位查询到期] {} 到期处理时出现异常", take.getKey());
                }
            } catch (InterruptedException e) {
                log.error("[设备订阅任务] ", e);
            }
        }
    }
}