修复 SiP消息超时未未回复无法识别的BUG

pull/1651/head
648540858 2024-10-16 17:17:17 +08:00
parent b9738c6bfc
commit 2492b0d638
16 changed files with 272 additions and 293 deletions

View File

@ -1,13 +1,15 @@
package com.genersoft.iot.vmp.conf; package com.genersoft.iot.vmp.conf;
import org.springframework.core.annotation.Order; import lombok.Data;
import org.springframework.boot.context.properties.ConfigurationProperties; import org.springframework.boot.context.properties.ConfigurationProperties;
import org.springframework.core.annotation.Order;
import org.springframework.stereotype.Component; import org.springframework.stereotype.Component;
@Component @Component
@ConfigurationProperties(prefix = "sip", ignoreInvalidFields = true) @ConfigurationProperties(prefix = "sip", ignoreInvalidFields = true)
@Order(0) @Order(0)
@Data
public class SipConfig { public class SipConfig {
private String ip; private String ip;
@ -26,82 +28,7 @@ public class SipConfig {
Integer registerTimeInterval = 120; Integer registerTimeInterval = 120;
private boolean alarm; private boolean alarm = false;
public void setIp(String ip) { private long timeout = 15;
this.ip = ip;
}
public void setPort(Integer port) {
this.port = port;
}
public void setDomain(String domain) {
this.domain = domain;
}
public void setId(String id) {
this.id = id;
}
public void setPassword(String password) {
this.password = password;
}
public void setPtzSpeed(Integer ptzSpeed) {
this.ptzSpeed = ptzSpeed;
}
public void setRegisterTimeInterval(Integer registerTimeInterval) {
this.registerTimeInterval = registerTimeInterval;
}
public String getIp() {
return ip;
}
public Integer getPort() {
return port;
}
public String getDomain() {
return domain;
}
public String getId() {
return id;
}
public String getPassword() {
return password;
}
public Integer getPtzSpeed() {
return ptzSpeed;
}
public Integer getRegisterTimeInterval() {
return registerTimeInterval;
}
public boolean isAlarm() {
return alarm;
}
public void setAlarm(boolean alarm) {
this.alarm = alarm;
}
public String getShowIp() {
return showIp;
}
public void setShowIp(String showIp) {
this.showIp = showIp;
}
} }

View File

@ -125,7 +125,7 @@ public class SipLayer implements CommandLineRunner {
SipProviderImpl udpSipProvider = (SipProviderImpl)sipStack.createSipProvider(udpListeningPoint); SipProviderImpl udpSipProvider = (SipProviderImpl)sipStack.createSipProvider(udpListeningPoint);
udpSipProvider.addSipListener(sipProcessorObserver); udpSipProvider.addSipListener(sipProcessorObserver);
udpSipProvider.setDialogErrorsAutomaticallyHandled();
udpSipProviderMap.put(monitorIp, udpSipProvider); udpSipProviderMap.put(monitorIp, udpSipProvider);
log.info("[SIP SERVER] udp://{}:{} 启动成功", monitorIp, port); log.info("[SIP SERVER] udp://{}:{} 启动成功", monitorIp, port);

View File

@ -1,6 +1,7 @@
package com.genersoft.iot.vmp.gb28181.event; package com.genersoft.iot.vmp.gb28181.event;
import com.genersoft.iot.vmp.gb28181.bean.DeviceNotFoundEvent; import com.genersoft.iot.vmp.gb28181.bean.DeviceNotFoundEvent;
import com.genersoft.iot.vmp.gb28181.event.sip.SipEvent;
import gov.nist.javax.sip.message.SIPRequest; import gov.nist.javax.sip.message.SIPRequest;
import gov.nist.javax.sip.message.SIPResponse; import gov.nist.javax.sip.message.SIPResponse;
import lombok.extern.slf4j.Slf4j; import lombok.extern.slf4j.Slf4j;
@ -13,10 +14,9 @@ import javax.sip.ResponseEvent;
import javax.sip.TimeoutEvent; import javax.sip.TimeoutEvent;
import javax.sip.TransactionTerminatedEvent; import javax.sip.TransactionTerminatedEvent;
import javax.sip.header.WarningHeader; import javax.sip.header.WarningHeader;
import java.time.Instant;
import java.util.Map; import java.util.Map;
import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.TimeUnit; import java.util.concurrent.DelayQueue;
/** /**
* @author lin * @author lin
@ -25,41 +25,36 @@ import java.util.concurrent.TimeUnit;
@Component @Component
public class SipSubscribe { public class SipSubscribe {
private final Map<String, SipSubscribe.Event> errorSubscribes = new ConcurrentHashMap<>(); private final Map<String, SipEvent> subscribes = new ConcurrentHashMap<>();
private final Map<String, SipSubscribe.Event> okSubscribes = new ConcurrentHashMap<>(); private final DelayQueue<SipEvent> delayQueue = new DelayQueue<>();
private final Map<String, Instant> okTimeSubscribes = new ConcurrentHashMap<>(); @Scheduled(fixedRate = 200) //每200毫秒执行
private final Map<String, Instant> errorTimeSubscribes = new ConcurrentHashMap<>();
// @Scheduled(cron="*/5 * * * * ?") //每五秒执行一次
// @Scheduled(fixedRate= 100 * 60 * 60 )
@Scheduled(cron="0 0/5 * * * ?") //每5分钟执行一次
public void execute(){ public void execute(){
if(log.isDebugEnabled()){ if (delayQueue.isEmpty()) {
log.info("[定时任务] 清理过期的SIP订阅信息"); return;
}
try {
SipEvent take = delayQueue.take();
// 出现超时异常
if(take.getErrorEvent() != null) {
EventResult<Object> eventResult = new EventResult<>();
eventResult.type = EventResultType.timeout;
eventResult.msg = "消息超时未回复";
eventResult.statusCode = -1024;
take.getErrorEvent().response(eventResult);
}
subscribes.remove(take.getKey());
} catch (InterruptedException e) {
throw new RuntimeException(e);
}
} }
Instant instant = Instant.now().minusMillis(TimeUnit.MINUTES.toMillis(5)); public void updateTimeout(String callId) {
SipEvent sipEvent = subscribes.get(callId);
for (String key : okTimeSubscribes.keySet()) { if (sipEvent != null) {
if (okTimeSubscribes.get(key).isBefore(instant)){ delayQueue.remove(sipEvent);
okSubscribes.remove(key); delayQueue.offer(sipEvent);
okTimeSubscribes.remove(key);
}
}
for (String key : errorTimeSubscribes.keySet()) {
if (errorTimeSubscribes.get(key).isBefore(instant)){
errorSubscribes.remove(key);
errorTimeSubscribes.remove(key);
}
}
if(log.isDebugEnabled()){
log.debug("okTimeSubscribes.size:{}",okTimeSubscribes.size());
log.debug("okSubscribes.size:{}",okSubscribes.size());
log.debug("errorTimeSubscribes.size:{}",errorTimeSubscribes.size());
log.debug("errorSubscribes.size:{}",errorSubscribes.size());
} }
} }
@ -156,43 +151,33 @@ public class SipSubscribe {
} }
} }
public void addErrorSubscribe(String key, SipSubscribe.Event event) {
errorSubscribes.put(key, event); public void addSubscribe(String key, SipEvent event) {
errorTimeSubscribes.put(key, Instant.now()); SipEvent sipEvent = subscribes.get(key);
if (sipEvent != null) {
subscribes.remove(key);
delayQueue.remove(sipEvent);
}
subscribes.put(key, event);
delayQueue.offer(event);
} }
public void addOkSubscribe(String key, SipSubscribe.Event event) { public SipEvent getSubscribe(String key) {
okSubscribes.put(key, event); return subscribes.get(key);
okTimeSubscribes.put(key, Instant.now());
} }
public SipSubscribe.Event getErrorSubscribe(String key) { public void removeSubscribe(String key) {
return errorSubscribes.get(key);
}
public void removeErrorSubscribe(String key) {
if(key == null){ if(key == null){
return; return;
} }
errorSubscribes.remove(key); SipEvent sipEvent = subscribes.get(key);
errorTimeSubscribes.remove(key); if (sipEvent != null) {
subscribes.remove(key);
delayQueue.remove(sipEvent);
}
} }
public SipSubscribe.Event getOkSubscribe(String key) { public boolean isEmpty(){
return okSubscribes.get(key); return subscribes.isEmpty();
}
public void removeOkSubscribe(String key) {
if(key == null){
return;
}
okSubscribes.remove(key);
okTimeSubscribes.remove(key);
}
public int getErrorSubscribesSize(){
return errorSubscribes.size();
}
public int getOkSubscribesSize(){
return okSubscribes.size();
} }
} }

View File

@ -0,0 +1,48 @@
package com.genersoft.iot.vmp.gb28181.event.sip;
import com.genersoft.iot.vmp.gb28181.event.SipSubscribe;
import lombok.Data;
import org.jetbrains.annotations.NotNull;
import java.util.concurrent.Delayed;
import java.util.concurrent.TimeUnit;
@Data
public class SipEvent implements Delayed {
private String key;
/**
*
*/
private SipSubscribe.Event okEvent;
/**
* ,
*/
private SipSubscribe.Event errorEvent;
/**
*
*/
private long delay;
public static SipEvent getInstance(String key, SipSubscribe.Event okEvent, SipSubscribe.Event errorEvent, long delay) {
SipEvent sipEvent = new SipEvent();
sipEvent.setKey(key);
sipEvent.setOkEvent(okEvent);
sipEvent.setErrorEvent(errorEvent);
sipEvent.setDelay(delay);
return sipEvent;
}
@Override
public long getDelay(@NotNull TimeUnit unit) {
return unit.convert(delay - System.currentTimeMillis(),TimeUnit.MILLISECONDS);
}
@Override
public int compareTo(@NotNull Delayed o) {
return (int) (this.getDelay(TimeUnit.MILLISECONDS) - o.getDelay(TimeUnit.MILLISECONDS));
}
}

View File

@ -380,7 +380,9 @@ public class PlatformServiceImpl implements IPlatformService {
commanderForPlatform.register(platform, sipTransactionInfo, eventResult -> { commanderForPlatform.register(platform, sipTransactionInfo, eventResult -> {
log.info("[国标级联] 平台:{}注册失败,{}:{}", platform.getServerGBId(), log.info("[国标级联] 平台:{}注册失败,{}:{}", platform.getServerGBId(),
eventResult.statusCode, eventResult.msg); eventResult.statusCode, eventResult.msg);
if (platform.isStatus()) {
offline(platform, false); offline(platform, false);
}
}, null); }, null);
} catch (Exception e) { } catch (Exception e) {
log.error("[命令发送失败] 国标级联定时注册: {}", e.getMessage()); log.error("[命令发送失败] 国标级联定时注册: {}", e.getMessage());

View File

@ -436,7 +436,7 @@ public class PlayServiceImpl implements IPlayService {
// 处理收到200ok后的TCP主动连接以及SSRC不一致的问题 // 处理收到200ok后的TCP主动连接以及SSRC不一致的问题
InviteOKHandler(eventResult, ssrcInfo, mediaServerItem, device, channel, callback, inviteInfo, InviteSessionType.PLAY); InviteOKHandler(eventResult, ssrcInfo, mediaServerItem, device, channel, callback, inviteInfo, InviteSessionType.PLAY);
}, (event) -> { }, (event) -> {
log.info("[点播失败] deviceId: {}, channelId:{}, {}: {}", device.getDeviceId(), channel.getDeviceId(), event.statusCode, event.msg); log.info("[点播失败]{}:{} deviceId: {}, channelId:{}",event.statusCode, event.msg, device.getDeviceId(), channel.getDeviceId());
receiveRtpServerService.closeRTPServer(mediaServerItem, ssrcInfo); receiveRtpServerService.closeRTPServer(mediaServerItem, ssrcInfo);
sessionManager.removeByStream(ssrcInfo.getStream()); sessionManager.removeByStream(ssrcInfo.getStream());
@ -447,7 +447,7 @@ public class PlayServiceImpl implements IPlayService {
event.statusCode, event.msg, null); event.statusCode, event.msg, null);
inviteStreamService.removeInviteInfoByDeviceAndChannel(InviteSessionType.PLAY, channel.getId()); inviteStreamService.removeInviteInfoByDeviceAndChannel(InviteSessionType.PLAY, channel.getId());
}); }, userSetting.getPlayTimeout().longValue());
} catch (InvalidArgumentException | SipException | ParseException e) { } catch (InvalidArgumentException | SipException | ParseException e) {
log.error("[命令发送失败] 点播消息: {}", e.getMessage()); log.error("[命令发送失败] 点播消息: {}", e.getMessage());
receiveRtpServerService.closeRTPServer(mediaServerItem, ssrcInfo); receiveRtpServerService.closeRTPServer(mediaServerItem, ssrcInfo);
@ -565,7 +565,7 @@ public class PlayServiceImpl implements IPlayService {
mediaServerService.releaseSsrc(mediaServerItem.getId(), sendRtpInfo.getSsrc()); mediaServerService.releaseSsrc(mediaServerItem.getId(), sendRtpInfo.getSsrc());
sessionManager.removeByStream(sendRtpInfo.getStream()); sessionManager.removeByStream(sendRtpInfo.getStream());
errorEvent.response(event); errorEvent.response(event);
}); }, userSetting.getPlayTimeout().longValue());
} catch (InvalidArgumentException | SipException | ParseException e) { } catch (InvalidArgumentException | SipException | ParseException e) {
log.error("[命令发送失败] 对讲消息: {}", e.getMessage()); log.error("[命令发送失败] 对讲消息: {}", e.getMessage());
@ -820,7 +820,7 @@ public class PlayServiceImpl implements IPlayService {
receiveRtpServerService.closeRTPServer(mediaServerItem, ssrcInfo); receiveRtpServerService.closeRTPServer(mediaServerItem, ssrcInfo);
sessionManager.removeByStream(ssrcInfo.getStream()); sessionManager.removeByStream(ssrcInfo.getStream());
inviteStreamService.removeInviteInfo(inviteInfo); inviteStreamService.removeInviteInfo(inviteInfo);
}); }, userSetting.getPlayTimeout().longValue());
} catch (InvalidArgumentException | SipException | ParseException e) { } catch (InvalidArgumentException | SipException | ParseException e) {
log.error("[命令发送失败] 录像回放: {}", e.getMessage()); log.error("[命令发送失败] 录像回放: {}", e.getMessage());
if (callback != null) { if (callback != null) {
@ -1044,7 +1044,7 @@ public class PlayServiceImpl implements IPlayService {
// 设置过期时间,下载失败时自动处理订阅数据 // 设置过期时间,下载失败时自动处理订阅数据
hook.setExpireTime(System.currentTimeMillis() + 24 * 60 * 60 * 1000); hook.setExpireTime(System.currentTimeMillis() + 24 * 60 * 60 * 1000);
subscribe.addSubscribe(hook, hookEventForRecord); subscribe.addSubscribe(hook, hookEventForRecord);
}); }, userSetting.getPlayTimeout().longValue());
} catch (InvalidArgumentException | SipException | ParseException e) { } catch (InvalidArgumentException | SipException | ParseException e) {
log.error("[命令发送失败] 录像下载: {}", e.getMessage()); log.error("[命令发送失败] 录像下载: {}", e.getMessage());
callback.run(InviteErrorCode.FAIL.getCode(),e.getMessage(), null); callback.run(InviteErrorCode.FAIL.getCode(),e.getMessage(), null);

View File

@ -2,18 +2,18 @@ package com.genersoft.iot.vmp.gb28181.transmit;
import com.genersoft.iot.vmp.gb28181.event.EventPublisher; import com.genersoft.iot.vmp.gb28181.event.EventPublisher;
import com.genersoft.iot.vmp.gb28181.event.SipSubscribe; import com.genersoft.iot.vmp.gb28181.event.SipSubscribe;
import com.genersoft.iot.vmp.gb28181.event.sip.SipEvent;
import com.genersoft.iot.vmp.gb28181.transmit.event.request.ISIPRequestProcessor; import com.genersoft.iot.vmp.gb28181.transmit.event.request.ISIPRequestProcessor;
import com.genersoft.iot.vmp.gb28181.transmit.event.response.ISIPResponseProcessor; import com.genersoft.iot.vmp.gb28181.transmit.event.response.ISIPResponseProcessor;
import com.genersoft.iot.vmp.gb28181.transmit.event.timeout.ITimeoutProcessor; import com.genersoft.iot.vmp.gb28181.transmit.event.timeout.ITimeoutProcessor;
import gov.nist.javax.sip.message.SIPResponse;
import lombok.extern.slf4j.Slf4j; import lombok.extern.slf4j.Slf4j;
import org.springframework.beans.factory.annotation.Autowired; import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.scheduling.annotation.Async; import org.springframework.scheduling.annotation.Async;
import org.springframework.stereotype.Component; import org.springframework.stereotype.Component;
import javax.sip.*; import javax.sip.*;
import javax.sip.header.CSeqHeader;
import javax.sip.header.CallIdHeader; import javax.sip.header.CallIdHeader;
import javax.sip.message.Request;
import javax.sip.message.Response; import javax.sip.message.Response;
import java.util.Map; import java.util.Map;
import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.ConcurrentHashMap;
@ -88,40 +88,40 @@ public class SIPProcessorObserver implements ISIPProcessorObserver {
@Override @Override
@Async("taskExecutor") @Async("taskExecutor")
public void processResponse(ResponseEvent responseEvent) { public void processResponse(ResponseEvent responseEvent) {
Response response = responseEvent.getResponse(); SIPResponse response = (SIPResponse)responseEvent.getResponse();
int status = response.getStatusCode(); int status = response.getStatusCode();
// Success // Success
if (((status >= Response.OK) && (status < Response.MULTIPLE_CHOICES)) || status == Response.UNAUTHORIZED) { if (((status >= Response.OK) && (status < Response.MULTIPLE_CHOICES)) || status == Response.UNAUTHORIZED) {
CSeqHeader cseqHeader = (CSeqHeader) responseEvent.getResponse().getHeader(CSeqHeader.NAME); if (status != Response.UNAUTHORIZED && responseEvent.getResponse() != null && !sipSubscribe.isEmpty() ) {
String method = cseqHeader.getMethod(); CallIdHeader callIdHeader = response.getCallIdHeader();
ISIPResponseProcessor sipRequestProcessor = responseProcessorMap.get(method); if (callIdHeader != null) {
SipEvent sipEvent = sipSubscribe.getSubscribe(callIdHeader.getCallId());
if (sipEvent != null && sipEvent.getOkEvent() != null) {
SipSubscribe.EventResult<ResponseEvent> eventResult = new SipSubscribe.EventResult<>(responseEvent);
sipSubscribe.removeSubscribe(callIdHeader.getCallId());
sipEvent.getOkEvent().response(eventResult);
}
}
}
ISIPResponseProcessor sipRequestProcessor = responseProcessorMap.get(response.getCSeqHeader().getMethod());
if (sipRequestProcessor != null) { if (sipRequestProcessor != null) {
sipRequestProcessor.process(responseEvent); sipRequestProcessor.process(responseEvent);
} }
if (status != Response.UNAUTHORIZED && responseEvent.getResponse() != null && sipSubscribe.getOkSubscribesSize() > 0 ) {
CallIdHeader callIdHeader = (CallIdHeader)responseEvent.getResponse().getHeader(CallIdHeader.NAME);
if (callIdHeader != null) {
SipSubscribe.Event subscribe = sipSubscribe.getOkSubscribe(callIdHeader.getCallId());
if (subscribe != null) {
SipSubscribe.EventResult eventResult = new SipSubscribe.EventResult(responseEvent);
sipSubscribe.removeOkSubscribe(callIdHeader.getCallId());
subscribe.response(eventResult);
}
}
}
} else if ((status >= Response.TRYING) && (status < Response.OK)) { } else if ((status >= Response.TRYING) && (status < Response.OK)) {
// 增加其它无需回复的响应如101、180等 // 增加其它无需回复的响应如101、180等
// 更新sip订阅的时间
// sipSubscribe.updateTimeout(response.getCallIdHeader().getCallId());
} else { } else {
log.warn("接收到失败的response响应status" + status + ",message:" + response.getReasonPhrase()); log.warn("接收到失败的response响应status" + status + ",message:" + response.getReasonPhrase());
if (responseEvent.getResponse() != null && sipSubscribe.getErrorSubscribesSize() > 0 ) { if (responseEvent.getResponse() != null && !sipSubscribe.isEmpty() ) {
CallIdHeader callIdHeader = (CallIdHeader)responseEvent.getResponse().getHeader(CallIdHeader.NAME); CallIdHeader callIdHeader = (CallIdHeader)responseEvent.getResponse().getHeader(CallIdHeader.NAME);
if (callIdHeader != null) { if (callIdHeader != null) {
SipSubscribe.Event subscribe = sipSubscribe.getErrorSubscribe(callIdHeader.getCallId()); SipEvent sipEvent = sipSubscribe.getSubscribe(callIdHeader.getCallId());
if (subscribe != null) { if (sipEvent != null && sipEvent.getErrorEvent() != null) {
SipSubscribe.EventResult eventResult = new SipSubscribe.EventResult(responseEvent); SipSubscribe.EventResult eventResult = new SipSubscribe.EventResult(responseEvent);
subscribe.response(eventResult); sipSubscribe.removeSubscribe(callIdHeader.getCallId());
sipSubscribe.removeErrorSubscribe(callIdHeader.getCallId()); sipEvent.getErrorEvent().response(eventResult);
} }
} }
} }
@ -140,27 +140,27 @@ public class SIPProcessorObserver implements ISIPProcessorObserver {
@Override @Override
public void processTimeout(TimeoutEvent timeoutEvent) { public void processTimeout(TimeoutEvent timeoutEvent) {
log.info("[消息发送超时]"); log.info("[消息发送超时]");
ClientTransaction clientTransaction = timeoutEvent.getClientTransaction(); // ClientTransaction clientTransaction = timeoutEvent.getClientTransaction();
//
if (clientTransaction != null) { // if (clientTransaction != null) {
log.info("[发送错误订阅] clientTransaction != null"); // log.info("[发送错误订阅] clientTransaction != null");
Request request = clientTransaction.getRequest(); // Request request = clientTransaction.getRequest();
if (request != null) { // if (request != null) {
log.info("[发送错误订阅] request != null"); // log.info("[发送错误订阅] request != null");
CallIdHeader callIdHeader = (CallIdHeader) request.getHeader(CallIdHeader.NAME); // CallIdHeader callIdHeader = (CallIdHeader) request.getHeader(CallIdHeader.NAME);
if (callIdHeader != null) { // if (callIdHeader != null) {
log.info("[发送错误订阅]"); // log.info("[发送错误订阅]");
SipSubscribe.Event subscribe = sipSubscribe.getErrorSubscribe(callIdHeader.getCallId()); // SipSubscribe.Event subscribe = sipSubscribe.getErrorSubscribe(callIdHeader.getCallId());
SipSubscribe.EventResult eventResult = new SipSubscribe.EventResult(timeoutEvent); // SipSubscribe.EventResult eventResult = new SipSubscribe.EventResult(timeoutEvent);
if (subscribe != null){ // if (subscribe != null){
subscribe.response(eventResult); // subscribe.response(eventResult);
} // }
sipSubscribe.removeOkSubscribe(callIdHeader.getCallId()); // sipSubscribe.removeOkSubscribe(callIdHeader.getCallId());
sipSubscribe.removeErrorSubscribe(callIdHeader.getCallId()); // sipSubscribe.removeErrorSubscribe(callIdHeader.getCallId());
} // }
} // }
} // }
eventPublisher.requestTimeOut(timeoutEvent); // eventPublisher.requestTimeOut(timeoutEvent);
} }
@Override @Override
@ -199,4 +199,6 @@ public class SIPProcessorObserver implements ISIPProcessorObserver {
} }
} }

View File

@ -1,7 +1,9 @@
package com.genersoft.iot.vmp.gb28181.transmit; package com.genersoft.iot.vmp.gb28181.transmit;
import com.genersoft.iot.vmp.conf.SipConfig;
import com.genersoft.iot.vmp.gb28181.SipLayer; import com.genersoft.iot.vmp.gb28181.SipLayer;
import com.genersoft.iot.vmp.gb28181.event.SipSubscribe; import com.genersoft.iot.vmp.gb28181.event.SipSubscribe;
import com.genersoft.iot.vmp.gb28181.event.sip.SipEvent;
import com.genersoft.iot.vmp.gb28181.utils.SipUtils; import com.genersoft.iot.vmp.gb28181.utils.SipUtils;
import com.genersoft.iot.vmp.utils.GitUtil; import com.genersoft.iot.vmp.utils.GitUtil;
import gov.nist.javax.sip.SipProviderImpl; import gov.nist.javax.sip.SipProviderImpl;
@ -21,6 +23,7 @@ import java.text.ParseException;
/** /**
* SIP * SIP
*
* @author lin * @author lin
*/ */
@Slf4j @Slf4j
@ -35,21 +38,27 @@ public class SIPSender {
@Autowired @Autowired
private SipSubscribe sipSubscribe; private SipSubscribe sipSubscribe;
@Autowired
private SipConfig sipConfig;
public void transmitRequest(String ip, Message message) throws SipException, ParseException { public void transmitRequest(String ip, Message message) throws SipException, ParseException {
transmitRequest(ip, message, null, null); transmitRequest(ip, message, null, null, null);
} }
public void transmitRequest(String ip, Message message, SipSubscribe.Event errorEvent) throws SipException, ParseException { public void transmitRequest(String ip, Message message, SipSubscribe.Event errorEvent) throws SipException, ParseException {
transmitRequest(ip, message, errorEvent, null); transmitRequest(ip, message, errorEvent, null, null);
} }
public void transmitRequest(String ip, Message message, SipSubscribe.Event errorEvent, SipSubscribe.Event okEvent) throws SipException { public void transmitRequest(String ip, Message message, SipSubscribe.Event errorEvent, SipSubscribe.Event okEvent) throws SipException {
ViaHeader viaHeader = (ViaHeader)message.getHeader(ViaHeader.NAME); transmitRequest(ip, message, errorEvent, okEvent, null);
}
public void transmitRequest(String ip, Message message, SipSubscribe.Event errorEvent, SipSubscribe.Event okEvent, Long timeout) throws SipException {
ViaHeader viaHeader = (ViaHeader) message.getHeader(ViaHeader.NAME);
String transport = "UDP"; String transport = "UDP";
if (viaHeader == null) { if (viaHeader == null) {
log.warn("[消息头缺失] ViaHeader 使用默认的UDP方式处理数据"); log.warn("[消息头缺失] ViaHeader 使用默认的UDP方式处理数据");
}else { } else {
transport = viaHeader.getTransport(); transport = viaHeader.getTransport();
} }
if (message.getHeader(UserAgentHeader.NAME) == null) { if (message.getHeader(UserAgentHeader.NAME) == null) {
@ -60,23 +69,22 @@ public class SIPSender {
} }
} }
if (okEvent != null || errorEvent != null) {
CallIdHeader callIdHeader = (CallIdHeader) message.getHeader(CallIdHeader.NAME); CallIdHeader callIdHeader = (CallIdHeader) message.getHeader(CallIdHeader.NAME);
// 添加错误订阅 SipEvent sipEvent = SipEvent.getInstance(callIdHeader.getCallId(), eventResult -> {
if (errorEvent != null) { sipSubscribe.removeSubscribe(eventResult.callId);
sipSubscribe.addErrorSubscribe(callIdHeader.getCallId(), (eventResult -> { if(okEvent != null) {
sipSubscribe.removeErrorSubscribe(eventResult.callId);
sipSubscribe.removeOkSubscribe(eventResult.callId);
errorEvent.response(eventResult);
}));
}
// 添加订阅
if (okEvent != null) {
sipSubscribe.addOkSubscribe(callIdHeader.getCallId(), eventResult -> {
sipSubscribe.removeOkSubscribe(eventResult.callId);
sipSubscribe.removeErrorSubscribe(eventResult.callId);
okEvent.response(eventResult); okEvent.response(eventResult);
});
} }
}, (eventResult -> {
sipSubscribe.removeSubscribe(eventResult.callId);
if (errorEvent != null) {
errorEvent.response(eventResult);
}
}), timeout == null ? sipConfig.getTimeout() : timeout);
sipSubscribe.addSubscribe(callIdHeader.getCallId(), sipEvent);
}
if ("TCP".equals(transport)) { if ("TCP".equals(transport)) {
SipProviderImpl tcpSipProvider = sipLayer.getTcpSipProvider(ip); SipProviderImpl tcpSipProvider = sipLayer.getTcpSipProvider(ip);
if (tcpSipProvider == null) { if (tcpSipProvider == null) {
@ -84,9 +92,9 @@ public class SIPSender {
return; return;
} }
if (message instanceof Request) { if (message instanceof Request) {
tcpSipProvider.sendRequest((Request)message); tcpSipProvider.sendRequest((Request) message);
}else if (message instanceof Response) { } else if (message instanceof Response) {
tcpSipProvider.sendResponse((Response)message); tcpSipProvider.sendResponse((Response) message);
} }
} else if ("UDP".equals(transport)) { } else if ("UDP".equals(transport)) {
@ -96,14 +104,14 @@ public class SIPSender {
return; return;
} }
if (message instanceof Request) { if (message instanceof Request) {
sipProvider.sendRequest((Request)message); sipProvider.sendRequest((Request) message);
}else if (message instanceof Response) { } else if (message instanceof Response) {
sipProvider.sendResponse((Response)message); sipProvider.sendResponse((Response) message);
} }
} }
} }
public CallIdHeader getNewCallIdHeader(String ip, String transport){ public CallIdHeader getNewCallIdHeader(String ip, String transport) {
if (ObjectUtils.isEmpty(transport)) { if (ObjectUtils.isEmpty(transport)) {
return sipLayer.getUdpSipProvider().getNewCallId(); return sipLayer.getUdpSipProvider().getNewCallId();
} }
@ -111,7 +119,7 @@ public class SIPSender {
if (ObjectUtils.isEmpty(ip)) { if (ObjectUtils.isEmpty(ip)) {
sipProvider = transport.equalsIgnoreCase("TCP") ? sipLayer.getTcpSipProvider() sipProvider = transport.equalsIgnoreCase("TCP") ? sipLayer.getTcpSipProvider()
: sipLayer.getUdpSipProvider(); : sipLayer.getUdpSipProvider();
}else { } else {
sipProvider = transport.equalsIgnoreCase("TCP") ? sipLayer.getTcpSipProvider(ip) sipProvider = transport.equalsIgnoreCase("TCP") ? sipLayer.getTcpSipProvider(ip)
: sipLayer.getUdpSipProvider(ip); : sipLayer.getUdpSipProvider(ip);
} }
@ -122,9 +130,11 @@ public class SIPSender {
if (sipProvider != null) { if (sipProvider != null) {
return sipProvider.getNewCallId(); return sipProvider.getNewCallId();
}else { } else {
log.warn("[新建CallIdHeader失败] ip={}, transport={}", ip, transport); log.warn("[新建CallIdHeader失败] ip={}, transport={}", ip, transport);
return null; return null;
} }
} }
} }

View File

@ -100,7 +100,7 @@ public interface ISIPCommander {
* @param device * @param device
* @param channel * @param channel
*/ */
void playStreamCmd(MediaServer mediaServerItem, SSRCInfo ssrcInfo, Device device, DeviceChannel channel, SipSubscribe.Event okEvent, SipSubscribe.Event errorEvent) throws InvalidArgumentException, SipException, ParseException; void playStreamCmd(MediaServer mediaServerItem, SSRCInfo ssrcInfo, Device device, DeviceChannel channel, SipSubscribe.Event okEvent, SipSubscribe.Event errorEvent, Long timeout) throws InvalidArgumentException, SipException, ParseException;
/** /**
* *
@ -110,7 +110,7 @@ public interface ISIPCommander {
* @param startTime ,yyyy-MM-dd HH:mm:ss * @param startTime ,yyyy-MM-dd HH:mm:ss
* @param endTime ,yyyy-MM-dd HH:mm:ss * @param endTime ,yyyy-MM-dd HH:mm:ss
*/ */
void playbackStreamCmd(MediaServer mediaServerItem, SSRCInfo ssrcInf, Device device, DeviceChannel channel, String startTime, String endTime, SipSubscribe.Event okEvent, SipSubscribe.Event errorEvent) throws InvalidArgumentException, SipException, ParseException; void playbackStreamCmd(MediaServer mediaServerItem, SSRCInfo ssrcInf, Device device, DeviceChannel channel, String startTime, String endTime, SipSubscribe.Event okEvent, SipSubscribe.Event errorEvent, Long timeout) throws InvalidArgumentException, SipException, ParseException;
/** /**
* *
@ -123,7 +123,7 @@ public interface ISIPCommander {
*/ */
void downloadStreamCmd(MediaServer mediaServerItem, SSRCInfo ssrcInfo, Device device, DeviceChannel channel, void downloadStreamCmd(MediaServer mediaServerItem, SSRCInfo ssrcInfo, Device device, DeviceChannel channel,
String startTime, String endTime, int downloadSpeed, String startTime, String endTime, int downloadSpeed,
SipSubscribe.Event errorEvent, SipSubscribe.Event okEvent) throws InvalidArgumentException, SipException, ParseException; SipSubscribe.Event errorEvent, SipSubscribe.Event okEvent, Long timeout) throws InvalidArgumentException, SipException, ParseException;
/** /**
@ -131,7 +131,7 @@ public interface ISIPCommander {
*/ */
void streamByeCmd(Device device, String channelId, String stream, String callId, SipSubscribe.Event okEvent) throws InvalidArgumentException, SipException, ParseException, SsrcTransactionNotFoundException; void streamByeCmd(Device device, String channelId, String stream, String callId, SipSubscribe.Event okEvent) throws InvalidArgumentException, SipException, ParseException, SsrcTransactionNotFoundException;
void talkStreamCmd(MediaServer mediaServerItem, SendRtpInfo sendRtpItem, Device device, DeviceChannel channelId, String callId, HookSubscribe.Event event, HookSubscribe.Event eventForPush, SipSubscribe.Event okEvent, SipSubscribe.Event errorEvent) throws InvalidArgumentException, SipException, ParseException; void talkStreamCmd(MediaServer mediaServerItem, SendRtpInfo sendRtpItem, Device device, DeviceChannel channelId, String callId, HookSubscribe.Event event, HookSubscribe.Event eventForPush, SipSubscribe.Event okEvent, SipSubscribe.Event errorEvent, Long timeout) throws InvalidArgumentException, SipException, ParseException;
void streamByeCmd(Device device, String channelId, String stream, String callId) throws InvalidArgumentException, ParseException, SipException, SsrcTransactionNotFoundException; void streamByeCmd(Device device, String channelId, String stream, String callId) throws InvalidArgumentException, ParseException, SipException, SsrcTransactionNotFoundException;

View File

@ -262,7 +262,7 @@ public class SIPCommander implements ISIPCommander {
*/ */
@Override @Override
public void playStreamCmd(MediaServer mediaServerItem, SSRCInfo ssrcInfo, Device device, DeviceChannel channel, public void playStreamCmd(MediaServer mediaServerItem, SSRCInfo ssrcInfo, Device device, DeviceChannel channel,
SipSubscribe.Event okEvent, SipSubscribe.Event errorEvent) throws InvalidArgumentException, SipException, ParseException { SipSubscribe.Event okEvent, SipSubscribe.Event errorEvent, Long timeout) throws InvalidArgumentException, SipException, ParseException {
String stream = ssrcInfo.getStream(); String stream = ssrcInfo.getStream();
if (device == null) { if (device == null) {
@ -335,8 +335,6 @@ public class SIPCommander implements ISIPCommander {
// f字段:f= v/编码格式/分辨率/帧率/码率类型/码率大小a/编码格式/码率大小/采样率 // f字段:f= v/编码格式/分辨率/帧率/码率类型/码率大小a/编码格式/码率大小/采样率
// content.append("f=v/2/5/25/1/4000a/1/8/1" + "\r\n"); // 未发现支持此特性的设备 // content.append("f=v/2/5/25/1/4000a/1/8/1" + "\r\n"); // 未发现支持此特性的设备
Request request = headerProvider.createInviteRequest(device, channel.getDeviceId(), content.toString(), SipUtils.getNewViaTag(), SipUtils.getNewFromTag(), null, ssrcInfo.getSsrc(),sipSender.getNewCallIdHeader(sipLayer.getLocalIp(device.getLocalIp()),device.getTransport())); Request request = headerProvider.createInviteRequest(device, channel.getDeviceId(), content.toString(), SipUtils.getNewViaTag(), SipUtils.getNewFromTag(), null, ssrcInfo.getSsrc(),sipSender.getNewCallIdHeader(sipLayer.getLocalIp(device.getLocalIp()),device.getTransport()));
sipSender.transmitRequest(sipLayer.getLocalIp(device.getLocalIp()), request, (e -> { sipSender.transmitRequest(sipLayer.getLocalIp(device.getLocalIp()), request, (e -> {
sessionManager.removeByStream(ssrcInfo.getStream()); sessionManager.removeByStream(ssrcInfo.getStream());
@ -350,7 +348,7 @@ public class SIPCommander implements ISIPCommander {
InviteSessionType.PLAY); InviteSessionType.PLAY);
sessionManager.put(ssrcTransaction); sessionManager.put(ssrcTransaction);
okEvent.response(e); okEvent.response(e);
}); }, timeout);
} }
/** /**
@ -364,7 +362,7 @@ public class SIPCommander implements ISIPCommander {
@Override @Override
public void playbackStreamCmd(MediaServer mediaServerItem, SSRCInfo ssrcInfo, Device device, DeviceChannel channel, public void playbackStreamCmd(MediaServer mediaServerItem, SSRCInfo ssrcInfo, Device device, DeviceChannel channel,
String startTime, String endTime, String startTime, String endTime,
SipSubscribe.Event okEvent, SipSubscribe.Event errorEvent) throws InvalidArgumentException, SipException, ParseException { SipSubscribe.Event okEvent, SipSubscribe.Event errorEvent, Long timeout) throws InvalidArgumentException, SipException, ParseException {
log.info("{} 分配的ZLM为: {} [{}:{}]", ssrcInfo.getStream(), mediaServerItem.getId(), mediaServerItem.getSdpIp(), ssrcInfo.getPort()); log.info("{} 分配的ZLM为: {} [{}:{}]", ssrcInfo.getStream(), mediaServerItem.getId(), mediaServerItem.getSdpIp(), ssrcInfo.getPort());
@ -445,7 +443,7 @@ public class SIPCommander implements ISIPCommander {
channel.getId(), sipSender.getNewCallIdHeader(sipLayer.getLocalIp(device.getLocalIp()),device.getTransport()).getCallId(), ssrcInfo.getStream(), ssrcInfo.getSsrc(), mediaServerItem.getId(), response, InviteSessionType.PLAYBACK); channel.getId(), sipSender.getNewCallIdHeader(sipLayer.getLocalIp(device.getLocalIp()),device.getTransport()).getCallId(), ssrcInfo.getStream(), ssrcInfo.getSsrc(), mediaServerItem.getId(), response, InviteSessionType.PLAYBACK);
sessionManager.put(ssrcTransaction); sessionManager.put(ssrcTransaction);
okEvent.response(event); okEvent.response(event);
}); }, timeout);
} }
/** /**
@ -454,7 +452,7 @@ public class SIPCommander implements ISIPCommander {
@Override @Override
public void downloadStreamCmd(MediaServer mediaServerItem, SSRCInfo ssrcInfo, Device device, DeviceChannel channel, public void downloadStreamCmd(MediaServer mediaServerItem, SSRCInfo ssrcInfo, Device device, DeviceChannel channel,
String startTime, String endTime, int downloadSpeed, String startTime, String endTime, int downloadSpeed,
SipSubscribe.Event errorEvent, SipSubscribe.Event okEvent) throws InvalidArgumentException, SipException, ParseException { SipSubscribe.Event errorEvent, SipSubscribe.Event okEvent, Long timeout) throws InvalidArgumentException, SipException, ParseException {
log.info("[发送-请求历史媒体下载-命令] 流ID {},节点为: {} [{}:{}]", ssrcInfo.getStream(), mediaServerItem.getId(), mediaServerItem.getSdpIp(), ssrcInfo.getPort()); log.info("[发送-请求历史媒体下载-命令] 流ID {},节点为: {} [{}:{}]", ssrcInfo.getStream(), mediaServerItem.getId(), mediaServerItem.getSdpIp(), ssrcInfo.getPort());
String sdpIp; String sdpIp;
@ -536,11 +534,13 @@ public class SIPCommander implements ISIPCommander {
SsrcTransaction ssrcTransaction = SsrcTransaction.buildForDevice(device.getDeviceId(), channel.getId(), response.getCallIdHeader().getCallId(), ssrcInfo.getStream(), ssrc, mediaServerItem.getId(), response, InviteSessionType.DOWNLOAD); SsrcTransaction ssrcTransaction = SsrcTransaction.buildForDevice(device.getDeviceId(), channel.getId(), response.getCallIdHeader().getCallId(), ssrcInfo.getStream(), ssrc, mediaServerItem.getId(), response, InviteSessionType.DOWNLOAD);
sessionManager.put(ssrcTransaction); sessionManager.put(ssrcTransaction);
okEvent.response(event); okEvent.response(event);
}); }, timeout);
} }
@Override @Override
public void talkStreamCmd(MediaServer mediaServerItem, SendRtpInfo sendRtpItem, Device device, DeviceChannel channel, String callId, HookSubscribe.Event event, HookSubscribe.Event eventForPush, SipSubscribe.Event okEvent, SipSubscribe.Event errorEvent) throws InvalidArgumentException, SipException, ParseException { public void talkStreamCmd(MediaServer mediaServerItem, SendRtpInfo sendRtpItem, Device device, DeviceChannel channel,
String callId, HookSubscribe.Event event, HookSubscribe.Event eventForPush, SipSubscribe.Event okEvent,
SipSubscribe.Event errorEvent, Long timeout) throws InvalidArgumentException, SipException, ParseException {
String stream = sendRtpItem.getStream(); String stream = sendRtpItem.getStream();
@ -601,7 +601,7 @@ public class SIPCommander implements ISIPCommander {
SsrcTransaction ssrcTransaction = SsrcTransaction.buildForDevice(device.getDeviceId(), channel.getId(), "talk", stream, sendRtpItem.getSsrc(), mediaServerItem.getId(), response, InviteSessionType.TALK); SsrcTransaction ssrcTransaction = SsrcTransaction.buildForDevice(device.getDeviceId(), channel.getId(), "talk", stream, sendRtpItem.getSsrc(), mediaServerItem.getId(), response, InviteSessionType.TALK);
sessionManager.put(ssrcTransaction); sessionManager.put(ssrcTransaction);
okEvent.response(e); okEvent.response(e);
}); }, timeout);
} }
/** /**

View File

@ -127,23 +127,19 @@ public class SIPCommanderForPlatform implements ISIPCommanderForPlatform {
// 将 callid 写入缓存, 等注册成功可以更新状态 // 将 callid 写入缓存, 等注册成功可以更新状态
String callIdFromHeader = callIdHeader.getCallId(); String callIdFromHeader = callIdHeader.getCallId();
redisCatchStorage.updatePlatformRegisterInfo(callIdFromHeader, PlatformRegisterInfo.getInstance(parentPlatform.getServerGBId(), isRegister)); redisCatchStorage.updatePlatformRegisterInfo(callIdFromHeader, PlatformRegisterInfo.getInstance(parentPlatform.getServerGBId(), isRegister));
sipSubscribe.addErrorSubscribe(callIdHeader.getCallId(), (event)->{
if (event != null) {
log.info("向上级平台 [ {} ] 注册发生错误: {} ",
parentPlatform.getServerGBId(),
event.msg);
}
redisCatchStorage.delPlatformRegisterInfo(callIdFromHeader);
if (errorEvent != null ) {
errorEvent.response(event);
}
});
}else { }else {
request = headerProviderPlatformProvider.createRegisterRequest(parentPlatform, fromTag, toTag, www, callIdHeader, isRegister? parentPlatform.getExpires() : 0); request = headerProviderPlatformProvider.createRegisterRequest(parentPlatform, fromTag, toTag, www, callIdHeader, isRegister? parentPlatform.getExpires() : 0);
} }
sipSender.transmitRequest(parentPlatform.getDeviceIp(), request, null, okEvent); sipSender.transmitRequest(parentPlatform.getDeviceIp(), request, (event)->{
if (event != null) {
log.info("[国标级联]{}, 注册失败: {} ", parentPlatform.getServerGBId(), event.msg);
}
redisCatchStorage.delPlatformRegisterInfo(callIdHeader.getCallId());
if (errorEvent != null ) {
errorEvent.response(event);
}
}, okEvent, 5L);
} }
@Override @Override
@ -167,6 +163,7 @@ public class SIPCommanderForPlatform implements ISIPCommanderForPlatform {
SipUtils.getNewFromTag(), SipUtils.getNewFromTag(),
SipUtils.getNewViaTag(), SipUtils.getNewViaTag(),
callIdHeader); callIdHeader);
sipSender.transmitRequest(parentPlatform.getDeviceIp(), request, errorEvent, okEvent); sipSender.transmitRequest(parentPlatform.getDeviceIp(), request, errorEvent, okEvent);
return callIdHeader.getCallId(); return callIdHeader.getCallId();
} }
@ -249,16 +246,17 @@ public class SIPCommanderForPlatform implements ISIPCommanderForPlatform {
log.debug(catalogXml); log.debug(catalogXml);
if (sendAfterResponse) { if (sendAfterResponse) {
// 默认按照收到200回复后发送下一条 如果超时收不到回复就以30毫秒的间隔直接发送。 // 默认按照收到200回复后发送下一条 如果超时收不到回复就以30毫秒的间隔直接发送。
dynamicTask.startDelay(timeoutTaskKey, ()->{ sipSender.transmitRequest(parentPlatform.getDeviceIp(), request, eventResult -> {
sipSubscribe.removeOkSubscribe(callId); if (eventResult.type.equals(SipSubscribe.EventResultType.timeout)) {
// 消息发送超时, 以30毫秒的间隔直接发送
int indexNext = index + parentPlatform.getCatalogGroup(); int indexNext = index + parentPlatform.getCatalogGroup();
try { try {
sendCatalogResponse(channels, parentPlatform, sn, fromTag, indexNext, false); sendCatalogResponse(channels, parentPlatform, sn, fromTag, indexNext, false);
} catch (SipException | InvalidArgumentException | ParseException e) { } catch (SipException | InvalidArgumentException | ParseException e) {
log.error("[命令发送失败] 国标级联 目录查询回复: {}", e.getMessage()); log.error("[命令发送失败] 国标级联 目录查询回复: {}", e.getMessage());
} }
}, 3000); return;
sipSender.transmitRequest(parentPlatform.getDeviceIp(), request, eventResult -> { }
log.error("[目录推送失败] 国标级联 platform : {}, code: {}, msg: {}, 停止发送", parentPlatform.getServerGBId(), eventResult.statusCode, eventResult.msg); log.error("[目录推送失败] 国标级联 platform : {}, code: {}, msg: {}, 停止发送", parentPlatform.getServerGBId(), eventResult.statusCode, eventResult.msg);
dynamicTask.stop(timeoutTaskKey); dynamicTask.stop(timeoutTaskKey);
}, eventResult -> { }, eventResult -> {

View File

@ -5,6 +5,7 @@ import com.genersoft.iot.vmp.gb28181.bean.DeviceNotFoundEvent;
import com.genersoft.iot.vmp.gb28181.bean.Platform; import com.genersoft.iot.vmp.gb28181.bean.Platform;
import com.genersoft.iot.vmp.gb28181.bean.SsrcTransaction; import com.genersoft.iot.vmp.gb28181.bean.SsrcTransaction;
import com.genersoft.iot.vmp.gb28181.event.SipSubscribe; import com.genersoft.iot.vmp.gb28181.event.SipSubscribe;
import com.genersoft.iot.vmp.gb28181.event.sip.SipEvent;
import com.genersoft.iot.vmp.gb28181.service.IPlatformService; import com.genersoft.iot.vmp.gb28181.service.IPlatformService;
import com.genersoft.iot.vmp.gb28181.session.SipInviteSessionManager; import com.genersoft.iot.vmp.gb28181.session.SipInviteSessionManager;
import com.genersoft.iot.vmp.gb28181.transmit.SIPProcessorObserver; import com.genersoft.iot.vmp.gb28181.transmit.SIPProcessorObserver;
@ -93,11 +94,12 @@ public class MessageRequestProcessor extends SIPRequestProcessorParent implement
// 不存在则回复404 // 不存在则回复404
responseAck(request, Response.NOT_FOUND, "device "+ deviceId +" not found"); responseAck(request, Response.NOT_FOUND, "device "+ deviceId +" not found");
log.warn("[设备未找到 ]deviceId: {}, callId: {}", deviceId, callIdHeader.getCallId()); log.warn("[设备未找到 ]deviceId: {}, callId: {}", deviceId, callIdHeader.getCallId());
if (sipSubscribe.getErrorSubscribe(callIdHeader.getCallId()) != null){ SipEvent sipEvent = sipSubscribe.getSubscribe(callIdHeader.getCallId());
if (sipEvent != null && sipEvent.getErrorEvent() != null){
DeviceNotFoundEvent deviceNotFoundEvent = new DeviceNotFoundEvent(evt.getDialog()); DeviceNotFoundEvent deviceNotFoundEvent = new DeviceNotFoundEvent(evt.getDialog());
deviceNotFoundEvent.setCallId(callIdHeader.getCallId()); deviceNotFoundEvent.setCallId(callIdHeader.getCallId());
SipSubscribe.EventResult eventResult = new SipSubscribe.EventResult(deviceNotFoundEvent); SipSubscribe.EventResult eventResult = new SipSubscribe.EventResult(deviceNotFoundEvent);
sipSubscribe.getErrorSubscribe(callIdHeader.getCallId()).response(eventResult); sipEvent.getErrorEvent().response(eventResult);
} }
}else { }else {
Element rootElement; Element rootElement;

View File

@ -1,5 +1,7 @@
package com.genersoft.iot.vmp.gb28181.transmit.event.response; package com.genersoft.iot.vmp.gb28181.transmit.event.response;
import org.springframework.scheduling.annotation.Async;
import javax.sip.ResponseEvent; import javax.sip.ResponseEvent;
/** /**
@ -9,6 +11,7 @@ import javax.sip.ResponseEvent;
*/ */
public interface ISIPResponseProcessor { public interface ISIPResponseProcessor {
void process(ResponseEvent evt); void process(ResponseEvent evt);

View File

@ -38,14 +38,12 @@ public class InviteResponseProcessor extends SIPResponseProcessorAbstract {
@Autowired @Autowired
private SIPProcessorObserver sipProcessorObserver; private SIPProcessorObserver sipProcessorObserver;
@Autowired @Autowired
private SIPSender sipSender; private SIPSender sipSender;
@Autowired @Autowired
private SIPRequestHeaderProvider headerProvider; private SIPRequestHeaderProvider headerProvider;
@Override @Override
public void afterPropertiesSet() throws Exception { public void afterPropertiesSet() throws Exception {
// 添加消息处理的订阅 // 添加消息处理的订阅

View File

@ -1,6 +1,7 @@
package com.genersoft.iot.vmp.gb28181.transmit.event.timeout.impl; package com.genersoft.iot.vmp.gb28181.transmit.event.timeout.impl;
import com.genersoft.iot.vmp.gb28181.event.SipSubscribe; import com.genersoft.iot.vmp.gb28181.event.SipSubscribe;
import com.genersoft.iot.vmp.gb28181.event.sip.SipEvent;
import com.genersoft.iot.vmp.gb28181.transmit.SIPProcessorObserver; import com.genersoft.iot.vmp.gb28181.transmit.SIPProcessorObserver;
import com.genersoft.iot.vmp.gb28181.transmit.event.timeout.ITimeoutProcessor; import com.genersoft.iot.vmp.gb28181.transmit.event.timeout.ITimeoutProcessor;
import lombok.extern.slf4j.Slf4j; import lombok.extern.slf4j.Slf4j;
@ -32,11 +33,12 @@ public class TimeoutProcessorImpl implements InitializingBean, ITimeoutProcessor
// TODO Auto-generated method stub // TODO Auto-generated method stub
CallIdHeader callIdHeader = event.getClientTransaction().getDialog().getCallId(); CallIdHeader callIdHeader = event.getClientTransaction().getDialog().getCallId();
String callId = callIdHeader.getCallId(); String callId = callIdHeader.getCallId();
SipSubscribe.Event errorSubscribe = sipSubscribe.getErrorSubscribe(callId); SipEvent sipEvent = sipSubscribe.getSubscribe(callId);
if (sipEvent != null && sipEvent.getErrorEvent() != null) {
SipSubscribe.EventResult<TimeoutEvent> timeoutEventEventResult = new SipSubscribe.EventResult<>(event); SipSubscribe.EventResult<TimeoutEvent> timeoutEventEventResult = new SipSubscribe.EventResult<>(event);
errorSubscribe.response(timeoutEventEventResult); sipEvent.getErrorEvent().response(timeoutEventEventResult);
sipSubscribe.removeErrorSubscribe(callId); sipSubscribe.removeSubscribe(callId);
sipSubscribe.removeOkSubscribe(callId); }
} catch (Exception e) { } catch (Exception e) {
log.error("[超时事件失败]: {}", e.getMessage()); log.error("[超时事件失败]: {}", e.getMessage());
} }

View File

@ -123,6 +123,8 @@ sip:
# keepalliveToOnline: false # keepalliveToOnline: false
# 是否存储alarm信息 # 是否存储alarm信息
alarm: false alarm: false
# 命令发送等待回复的超时时间, 单位:秒
timeout: 15
# 做为JT1078服务器的配置 # 做为JT1078服务器的配置
jt1078: jt1078: