[freeboxos] Review WebSocketManager to avoid IllegalStateException (#20280)

* Review WebSocketManager to avoid IllegalStateException

Signed-off-by: gael@lhopital.org <gael@lhopital.org>
This commit is contained in:
Gaël L'hopital
2026-02-28 13:22:21 +01:00
committed by GitHub
parent 2d1112004b
commit def04d09f0
@@ -20,7 +20,6 @@ import java.util.HashMap;
import java.util.List;
import java.util.Locale;
import java.util.Map;
import java.util.Optional;
import java.util.concurrent.ScheduledExecutorService;
import java.util.concurrent.ScheduledFuture;
import java.util.concurrent.TimeUnit;
@@ -66,14 +65,13 @@ public class WebSocketManager extends RestManager implements WebSocketListener {
private final Logger logger = LoggerFactory.getLogger(WebSocketManager.class);
private final Map<MACAddress, ApiConsumerHandler> listeners = new HashMap<>();
private final ScheduledExecutorService scheduler = ThreadPoolManager.getScheduledPool(BINDING_ID);
private final Object wsLifecycleLock = new Object();
private final ApiHandler apiHandler;
private final WebSocketClient client;
private Optional<ScheduledFuture<?>> reconnectJob = Optional.empty();
private @Nullable WebSocketClient wsClient;
private @Nullable ScheduledFuture<?> reconnectJob;
private String activeRegistration = "";
private volatile @Nullable Session wsSession;
@Nullable
private String sessionToken;
private int reconnectInterval;
private boolean vmSupported;
private volatile boolean reconnectEnabled;
private record Register(String action, List<String> events) {
Register(String... events) {
@@ -81,12 +79,6 @@ public class WebSocketManager extends RestManager implements WebSocketListener {
}
}
public WebSocketManager(FreeboxOsSession session) throws FreeboxException {
super(session, LoginManager.Permission.NONE, session.getUriBuilder().path(WS_PATH));
this.apiHandler = session.getApiHandler();
this.client = new WebSocketClient(apiHandler.getHttpClient());
}
private enum Action {
REGISTER,
NOTIFICATION,
@@ -100,59 +92,76 @@ public class WebSocketManager extends RestManager implements WebSocketListener {
}
}
public void openSession(@Nullable String sessionToken, int reconnectInterval, boolean vmSupported) {
this.sessionToken = sessionToken;
this.reconnectInterval = reconnectInterval;
this.vmSupported = vmSupported;
if (reconnectInterval > 0) {
try {
client.start();
startReconnect();
} catch (Exception e) {
logger.warn("Error starting websocket client: {}", e.getMessage());
}
}
public WebSocketManager(FreeboxOsSession session) throws FreeboxException {
super(session, LoginManager.Permission.NONE, session.getUriBuilder().path(WS_PATH));
this.apiHandler = session.getApiHandler();
}
private void startReconnect() {
public void openSession(@Nullable String sessionToken, int reconnectInterval, boolean vmSupported) {
if (reconnectInterval <= 0) {
return;
}
activeRegistration = apiHandler.serialize(vmSupported ? REGISTRATION : REGISTRATION_WITHOUT_VM);
WebSocketClient webSocketClient = new WebSocketClient(apiHandler.getHttpClient());
reconnectEnabled = true;
URI uri = getUriBuilder().scheme(getUriBuilder().build().getScheme().contains("s") ? "wss" : "ws").build();
ClientUpgradeRequest request = new ClientUpgradeRequest();
request.setHeader(ApiHandler.AUTH_HEADER, sessionToken);
stopReconnect();
reconnectJob = Optional.of(scheduler.scheduleWithFixedDelay(() -> {
try {
closeSession();
client.connect(this, uri, request);
// Update listeners in case we would have lost data while disconnecting / reconnecting
listeners.values().forEach(host -> host
.handleCommand(new ChannelUID(host.getThing().getUID(), REACHABLE), RefreshType.REFRESH));
logger.debug("Websocket manager connected to {}", uri);
} catch (IOException e) {
logger.warn("Error connecting websocket client: {}", e.getMessage());
}
}, 0, reconnectInterval, TimeUnit.MINUTES));
try {
webSocketClient.start();
stopReconnectJob();
reconnectJob = scheduler.scheduleWithFixedDelay(() -> {
try {
synchronized (wsLifecycleLock) {
if (!reconnectEnabled || !webSocketClient.isStarted()) {
return;
}
closeSession();
webSocketClient.connect(this, uri, request);
}
// Update listeners in case we would have lost data while disconnecting / reconnecting
listeners.values().forEach(host -> host
.handleCommand(new ChannelUID(host.getThing().getUID(), REACHABLE), RefreshType.REFRESH));
logger.debug("Websocket manager connected to {}", uri);
} catch (IOException | IllegalStateException e) {
logger.warn("Error connecting websocket client: {}", e.getMessage());
}
}, 0, reconnectInterval, TimeUnit.MINUTES);
wsClient = webSocketClient;
} catch (Exception e) {
logger.warn("Error starting websocket client: {}", e.getMessage());
}
}
private void stopReconnect() {
reconnectJob.ifPresent(job -> job.cancel(true));
reconnectJob = Optional.empty();
private void stopReconnectJob() {
if (reconnectJob instanceof ScheduledFuture job) {
job.cancel(true);
reconnectJob = null;
}
}
public void dispose() {
stopReconnect();
closeSession();
try {
client.stop();
} catch (Exception e) {
logger.warn("Error stopping websocket client: {}", e.getMessage());
synchronized (wsLifecycleLock) {
reconnectEnabled = false;
stopReconnectJob();
if (wsClient instanceof WebSocketClient client) {
closeSession();
try {
client.stop();
} catch (Exception e) {
logger.warn("Error stopping websocket client: {}", e.getMessage());
} finally {
wsClient = null;
}
}
}
}
private void closeSession() {
logger.debug("Awaiting closure from remote");
Session localSession = wsSession;
if (localSession != null) {
localSession.close();
if (wsSession instanceof Session session) {
session.close();
wsSession = null;
}
}
@@ -162,8 +171,7 @@ public class WebSocketManager extends RestManager implements WebSocketListener {
this.wsSession = wsSession;
logger.debug("Websocket connection establisehd");
try {
wsSession.getRemote()
.sendString(apiHandler.serialize(vmSupported ? REGISTRATION : REGISTRATION_WITHOUT_VM));
wsSession.getRemote().sendString(activeRegistration);
} catch (IOException e) {
logger.warn("Error registering to websocket: {}", e.getMessage());
}
@@ -171,55 +179,53 @@ public class WebSocketManager extends RestManager implements WebSocketListener {
@Override
public void onWebSocketText(@NonNullByDefault({}) String message) {
logger.debug("Websocket: received message: {}", message);
Session localSession = wsSession;
if (message.toLowerCase(Locale.US).contains("bye") && localSession != null) {
localSession.close(StatusCode.NORMAL, "Thanks");
logger.debug("Websocket received: {}", message);
if (message.toLowerCase(Locale.US).contains("bye") && wsSession instanceof Session session) {
session.close(StatusCode.NORMAL, "Thanks");
return;
}
WebSocketResponse result = apiHandler.deserialize(WebSocketResponse.class, message);
if (result.success) {
switch (result.action) {
WebSocketResponse response = apiHandler.deserialize(WebSocketResponse.class, message);
if (response.success) {
switch (response.action) {
case REGISTER:
logger.debug("Event registration successfull");
break;
case NOTIFICATION:
handleNotification(result);
handleNotification(response);
break;
default:
logger.warn("Unhandled notification received: {}", result.action);
logger.warn("Unhandled notification received: {}", response.action);
}
}
}
private void handleNotification(WebSocketResponse response) {
JsonElement json = response.result;
if (json != null) {
switch (response.getEvent()) {
case VM_CHANGED:
VirtualMachine vm = apiHandler.deserialize(VirtualMachine.class, json.toString());
logger.debug("Received notification for VM {}", vm.id());
ApiConsumerHandler handler = listeners.get(vm.mac());
if (handler instanceof VmHandler vmHandler) {
vmHandler.updateVmChannels(vm);
}
break;
case HOST_UNREACHABLE, HOST_REACHABLE:
LanHost host = apiHandler.deserialize(LanHost.class, json.toString());
ApiConsumerHandler handler2 = listeners.get(host.getMac());
if (handler2 instanceof HostHandler hostHandler) {
logger.debug("Received notification for mac {} : thing {} is {}reachable",
host.getMac().toColonDelimitedString(), hostHandler.getThing().getUID(),
host.reachable() ? "" : "not ");
hostHandler.updateConnectivityChannels(host);
}
break;
default:
logger.warn("Unhandled event received: {}", response.getEvent());
}
} else {
if (json == null) {
logger.warn("Empty json element in notification");
return;
}
switch (response.getEvent()) {
case VM_CHANGED:
VirtualMachine vm = apiHandler.deserialize(VirtualMachine.class, json.toString());
logger.debug("Received notification for VM {}", vm.id());
if (listeners.get(vm.mac()) instanceof VmHandler vmHandler) {
vmHandler.updateVmChannels(vm);
}
break;
case HOST_UNREACHABLE, HOST_REACHABLE:
LanHost host = apiHandler.deserialize(LanHost.class, json.toString());
if (listeners.get(host.getMac()) instanceof HostHandler hostHandler) {
logger.debug("Received notification for mac {} : thing {} is {}reachable",
host.getMac().toColonDelimitedString(), hostHandler.getThing().getUID(),
host.reachable() ? "" : "not ");
hostHandler.updateConnectivityChannels(host);
}
break;
default:
logger.warn("Unhandled event received: {}", response.getEvent());
}
}
@@ -240,11 +246,12 @@ public class WebSocketManager extends RestManager implements WebSocketListener {
}
public boolean registerListener(MACAddress mac, ApiConsumerHandler hostHandler) {
if (wsSession != null) {
listeners.put(mac, hostHandler);
return true;
if (wsSession == null) {
return false;
}
return false;
listeners.put(mac, hostHandler);
return true;
}
public void unregisterListener(MACAddress mac) {