Browse Source

Attempt to resolve race condition

pull/6180/head
Andrii Shvaika 5 years ago
parent
commit
a89a8f0422
  1. 37
      application/src/main/java/org/thingsboard/server/controller/plugin/TbWebSocketHandler.java
  2. 24
      application/src/main/java/org/thingsboard/server/controller/plugin/TbWebSocketMsg.java
  3. 21
      application/src/main/java/org/thingsboard/server/controller/plugin/TbWebSocketMsgType.java
  4. 38
      application/src/main/java/org/thingsboard/server/controller/plugin/TbWebSocketPingMsg.java
  5. 34
      application/src/main/java/org/thingsboard/server/controller/plugin/TbWebSocketTextMsg.java

37
application/src/main/java/org/thingsboard/server/controller/plugin/TbWebSocketHandler.java

@ -67,7 +67,7 @@ public class TbWebSocketHandler extends TextWebSocketHandler implements Telemetr
private static final ConcurrentMap<String, SessionMetaData> internalSessionMap = new ConcurrentHashMap<>();
private static final ConcurrentMap<String, String> externalSessionMap = new ConcurrentHashMap<>();
private static final ByteBuffer PING_MSG = ByteBuffer.wrap(new byte[]{});
@Autowired
private TelemetryWebSocketService webSocketService;
@ -216,7 +216,7 @@ public class TbWebSocketHandler extends TextWebSocketHandler implements Telemetr
private final TelemetryWebSocketSessionRef sessionRef;
private volatile boolean isSending = false;
private final Queue<String> msgQueue;
private final Queue<TbWebSocketMsg<?>> msgQueue;
private volatile long lastActivityTime;
@ -237,7 +237,7 @@ public class TbWebSocketHandler extends TextWebSocketHandler implements Telemetr
log.warn("[{}] Closing session due to ping timeout", session.getId());
closeSession(CloseStatus.SESSION_NOT_RELIABLE);
} else if (timeSinceLastActivity >= pingTimeout / NUMBER_OF_PING_ATTEMPTS) {
this.asyncRemote.sendPing(PING_MSG);
sendMsg(TbWebSocketPingMsg.INSTANCE);
}
} catch (Exception e) {
log.trace("[{}] Failed to send ping msg", session.getId(), e);
@ -258,6 +258,10 @@ public class TbWebSocketHandler extends TextWebSocketHandler implements Telemetr
}
synchronized void sendMsg(String msg) {
sendMsg(new TbWebSocketTextMsg(msg));
}
synchronized void sendMsg(TbWebSocketMsg<?> msg) {
if (isSending) {
try {
msgQueue.add(msg);
@ -275,9 +279,16 @@ public class TbWebSocketHandler extends TextWebSocketHandler implements Telemetr
}
}
private void sendMsgInternal(String msg) {
private void sendMsgInternal(TbWebSocketMsg<?> msg) {
try {
this.asyncRemote.sendText(msg, this);
if (TbWebSocketMsgType.TEXT.equals(msg.getType())) {
TbWebSocketTextMsg textMsg = (TbWebSocketTextMsg) msg;
this.asyncRemote.sendText(textMsg.getMsg(), this);
} else {
TbWebSocketPingMsg pingMsg = (TbWebSocketPingMsg) msg;
this.asyncRemote.sendPing(pingMsg.getMsg());
processNextMsg();
}
} catch (Exception e) {
log.trace("[{}] Failed to send msg", session.getId(), e);
closeSession(CloseStatus.SESSION_NOT_RELIABLE);
@ -290,12 +301,16 @@ public class TbWebSocketHandler extends TextWebSocketHandler implements Telemetr
log.trace("[{}] Failed to send msg", session.getId(), result.getException());
closeSession(CloseStatus.SESSION_NOT_RELIABLE);
} else {
String msg = msgQueue.poll();
if (msg != null) {
sendMsgInternal(msg);
} else {
isSending = false;
}
processNextMsg();
}
}
private void processNextMsg() {
TbWebSocketMsg<?> msg = msgQueue.poll();
if (msg != null) {
sendMsgInternal(msg);
} else {
isSending = false;
}
}
}

24
application/src/main/java/org/thingsboard/server/controller/plugin/TbWebSocketMsg.java

@ -0,0 +1,24 @@
/**
* Copyright © 2016-2022 The Thingsboard Authors
*
* Licensed under the Apache License, Version 2.0 (the "License");
* you may not use this file except in compliance with the License.
* You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package org.thingsboard.server.controller.plugin;
public interface TbWebSocketMsg<T> {
TbWebSocketMsgType getType();
T getMsg();
}

21
application/src/main/java/org/thingsboard/server/controller/plugin/TbWebSocketMsgType.java

@ -0,0 +1,21 @@
/**
* Copyright © 2016-2022 The Thingsboard Authors
*
* Licensed under the Apache License, Version 2.0 (the "License");
* you may not use this file except in compliance with the License.
* You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package org.thingsboard.server.controller.plugin;
public enum TbWebSocketMsgType {
PING, TEXT
}

38
application/src/main/java/org/thingsboard/server/controller/plugin/TbWebSocketPingMsg.java

@ -0,0 +1,38 @@
/**
* Copyright © 2016-2022 The Thingsboard Authors
*
* Licensed under the Apache License, Version 2.0 (the "License");
* you may not use this file except in compliance with the License.
* You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package org.thingsboard.server.controller.plugin;
import lombok.RequiredArgsConstructor;
import java.nio.ByteBuffer;
@RequiredArgsConstructor
public class TbWebSocketPingMsg implements TbWebSocketMsg<ByteBuffer> {
public static TbWebSocketPingMsg INSTANCE = new TbWebSocketPingMsg();
private static final ByteBuffer PING_MSG = ByteBuffer.wrap(new byte[]{});
@Override
public TbWebSocketMsgType getType() {
return TbWebSocketMsgType.PING;
}
@Override
public ByteBuffer getMsg() {
return PING_MSG;
}
}

34
application/src/main/java/org/thingsboard/server/controller/plugin/TbWebSocketTextMsg.java

@ -0,0 +1,34 @@
/**
* Copyright © 2016-2022 The Thingsboard Authors
*
* Licensed under the Apache License, Version 2.0 (the "License");
* you may not use this file except in compliance with the License.
* You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package org.thingsboard.server.controller.plugin;
import lombok.RequiredArgsConstructor;
@RequiredArgsConstructor
public class TbWebSocketTextMsg implements TbWebSocketMsg<String> {
private final String value;
@Override
public TbWebSocketMsgType getType() {
return TbWebSocketMsgType.TEXT;
}
@Override
public String getMsg() {
return value;
}
}
Loading…
Cancel
Save