diff --git a/bom/compile/pom.xml b/bom/compile/pom.xml index eeff1ff73..713aa1b80 100644 --- a/bom/compile/pom.xml +++ b/bom/compile/pom.xml @@ -142,6 +142,12 @@ ${jetty.version} compile + + org.eclipse.jetty.websocket + websocket-servlet + ${jetty.version} + compile + diff --git a/bom/openhab-core/pom.xml b/bom/openhab-core/pom.xml index 968c2ee05..112cd968a 100644 --- a/bom/openhab-core/pom.xml +++ b/bom/openhab-core/pom.xml @@ -286,6 +286,12 @@ ${project.version} compile + + org.openhab.core.bundles + org.openhab.core.io.websocket + ${project.version} + compile + org.openhab.core.bundles org.openhab.core.io.jetty.certificate diff --git a/bundles/org.openhab.core.io.websocket/pom.xml b/bundles/org.openhab.core.io.websocket/pom.xml new file mode 100644 index 000000000..cf98d5e91 --- /dev/null +++ b/bundles/org.openhab.core.io.websocket/pom.xml @@ -0,0 +1,25 @@ + + + + 4.0.0 + + + org.openhab.core.bundles + org.openhab.core.reactor.bundles + 3.4.0-SNAPSHOT + + + org.openhab.core.io.websocket + + openHAB Core :: Bundles :: WebSocket + + + + org.openhab.core.bundles + org.openhab.core + ${project.version} + + + + diff --git a/bundles/org.openhab.core.io.websocket/src/main/java/org/openhab/core/io/websocket/EventDTO.java b/bundles/org.openhab.core.io.websocket/src/main/java/org/openhab/core/io/websocket/EventDTO.java new file mode 100644 index 000000000..8270664ce --- /dev/null +++ b/bundles/org.openhab.core.io.websocket/src/main/java/org/openhab/core/io/websocket/EventDTO.java @@ -0,0 +1,76 @@ +/** + * Copyright (c) 2010-2022 Contributors to the openHAB project + * + * See the NOTICE file(s) distributed with this work for additional + * information. + * + * This program and the accompanying materials are made available under the + * terms of the Eclipse Public License 2.0 which is available at + * http://www.eclipse.org/legal/epl-2.0 + * + * SPDX-License-Identifier: EPL-2.0 + */ +package org.openhab.core.io.websocket; + +import java.util.Objects; + +import org.eclipse.jdt.annotation.NonNullByDefault; +import org.eclipse.jdt.annotation.Nullable; +import org.openhab.core.events.Event; + +/** + * The {@link EventDTO} is used for serialization and deserialization of events + * + * @author Jan N. Klug - Initial contribution + */ +@NonNullByDefault +public class EventDTO { + public @Nullable String type; + public @Nullable String topic; + public @Nullable String payload; + public @Nullable String source; + + public @Nullable String eventId; + + public EventDTO() { + } + + public EventDTO(String type, String topic, @Nullable String payload, @Nullable String source, + @Nullable String eventId) { + this.type = type; + this.topic = topic; + this.payload = payload; + this.source = source; + this.eventId = eventId; + } + + public EventDTO(Event event) { + type = event.getType(); + topic = event.getTopic(); + source = event.getSource(); + payload = event.getPayload(); + } + + @Override + public boolean equals(@Nullable Object o) { + if (this == o) + return true; + if (o == null || getClass() != o.getClass()) + return false; + EventDTO eventDTO = (EventDTO) o; + return Objects.equals(type, eventDTO.type) && Objects.equals(topic, eventDTO.topic) + && Objects.equals(payload, eventDTO.payload) && Objects.equals(source, eventDTO.source) + && Objects.equals(eventId, eventDTO.eventId); + } + + @Override + public int hashCode() { + return Objects.hash(type, topic, payload, source, eventId); + } + + @Override + public String toString() { + return "EventDTO{type='" + type + "', topic='" + topic + "', payload='" + payload + "', source='" + source + + "', eventId='" + eventId + "'}"; + } +} diff --git a/bundles/org.openhab.core.io.websocket/src/main/java/org/openhab/core/io/websocket/EventProcessingException.java b/bundles/org.openhab.core.io.websocket/src/main/java/org/openhab/core/io/websocket/EventProcessingException.java new file mode 100644 index 000000000..66114a0a3 --- /dev/null +++ b/bundles/org.openhab.core.io.websocket/src/main/java/org/openhab/core/io/websocket/EventProcessingException.java @@ -0,0 +1,29 @@ +/** + * Copyright (c) 2010-2022 Contributors to the openHAB project + * + * See the NOTICE file(s) distributed with this work for additional + * information. + * + * This program and the accompanying materials are made available under the + * terms of the Eclipse Public License 2.0 which is available at + * http://www.eclipse.org/legal/epl-2.0 + * + * SPDX-License-Identifier: EPL-2.0 + */ +package org.openhab.core.io.websocket; + +import org.eclipse.jdt.annotation.NonNullByDefault; + +/** + * The {@link EventProcessingException} is thrown when processing of incoming events fails + * + * @author Jan N. Klug - Initial contribution + */ +@NonNullByDefault +public class EventProcessingException extends Exception { + private static final long serialVersionUID = 1L; + + public EventProcessingException(String message) { + super(message); + } +} diff --git a/bundles/org.openhab.core.io.websocket/src/main/java/org/openhab/core/io/websocket/EventWebSocket.java b/bundles/org.openhab.core.io.websocket/src/main/java/org/openhab/core/io/websocket/EventWebSocket.java new file mode 100644 index 000000000..6e86b151c --- /dev/null +++ b/bundles/org.openhab.core.io.websocket/src/main/java/org/openhab/core/io/websocket/EventWebSocket.java @@ -0,0 +1,205 @@ +/** + * Copyright (c) 2010-2022 Contributors to the openHAB project + * + * See the NOTICE file(s) distributed with this work for additional + * information. + * + * This program and the accompanying materials are made available under the + * terms of the Eclipse Public License 2.0 which is available at + * http://www.eclipse.org/legal/epl-2.0 + * + * SPDX-License-Identifier: EPL-2.0 + */ +package org.openhab.core.io.websocket; + +import java.io.IOException; +import java.lang.reflect.Type; +import java.util.List; +import java.util.Objects; + +import org.eclipse.jdt.annotation.NonNullByDefault; +import org.eclipse.jdt.annotation.Nullable; +import org.eclipse.jetty.websocket.api.RemoteEndpoint; +import org.eclipse.jetty.websocket.api.Session; +import org.eclipse.jetty.websocket.api.StatusCode; +import org.eclipse.jetty.websocket.api.annotations.OnWebSocketClose; +import org.eclipse.jetty.websocket.api.annotations.OnWebSocketConnect; +import org.eclipse.jetty.websocket.api.annotations.OnWebSocketError; +import org.eclipse.jetty.websocket.api.annotations.OnWebSocketMessage; +import org.eclipse.jetty.websocket.api.annotations.WebSocket; +import org.openhab.core.events.Event; +import org.openhab.core.events.EventPublisher; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + +import com.google.gson.Gson; +import com.google.gson.JsonParseException; +import com.google.gson.reflect.TypeToken; + +/** + * The {@link EventWebSocket} is the WebSocket implementation that extends the event bus + * + * @author Jan N. Klug - Initial contribution + */ +@WebSocket +@NonNullByDefault +@SuppressWarnings("unused") +public class EventWebSocket { + public static final String WEBSOCKET_EVENT_TYPE = "WebSocketEvent"; + + private static final Type STRING_LIST_TYPE = TypeToken.getParameterized(List.class, String.class).getType(); + + private final Logger logger = LoggerFactory.getLogger(EventWebSocket.class); + + private final EventWebSocketServlet servlet; + private final Gson gson; + private final EventPublisher eventPublisher; + private final ItemEventUtility itemEventUtility; + + private @Nullable Session session; + private @Nullable RemoteEndpoint remoteEndpoint; + private String remoteIdentifier = ""; + + private List typeFilter = List.of(); + private List sourceFilter = List.of(); + + public EventWebSocket(Gson gson, EventWebSocketServlet servlet, ItemEventUtility itemEventUtility, + EventPublisher eventPublisher) { + this.servlet = servlet; + this.gson = gson; + this.itemEventUtility = itemEventUtility; + this.eventPublisher = eventPublisher; + } + + @OnWebSocketClose + public void onClose(int statusCode, String reason) { + this.servlet.unregisterListener(this); + remoteIdentifier = ""; + this.session = null; + this.remoteEndpoint = null; + } + + @OnWebSocketConnect + public void onConnect(Session session) { + this.session = session; + RemoteEndpoint remoteEndpoint = session.getRemote(); + this.remoteEndpoint = remoteEndpoint; + this.remoteIdentifier = remoteEndpoint.getInetSocketAddress().toString(); + this.servlet.registerListener(this); + } + + @OnWebSocketMessage + public void onText(String message) { + RemoteEndpoint remoteEndpoint = this.remoteEndpoint; + if (session == null || remoteEndpoint == null) { + // no connection or no remote endpoint , do nothing this is possible due to async behavior + return; + } + + EventDTO responseEvent; + + try { + EventDTO eventDTO = gson.fromJson(message, EventDTO.class); + try { + if (eventDTO == null) { + throw new EventProcessingException("Deserialized event must not be null"); + } + String type = eventDTO.type; + if (type == null) { + throw new EventProcessingException("Event type must not be null."); + } + + switch (type) { + case "ItemCommandEvent": + Event itemCommandEvent = itemEventUtility.createCommandEvent(eventDTO); + eventPublisher.post(itemCommandEvent); + responseEvent = new EventDTO(WEBSOCKET_EVENT_TYPE, "/response/success", "", null, + eventDTO.eventId); + break; + case "ItemStateEvent": + Event itemStateEvent = itemEventUtility.createStateEvent(eventDTO); + eventPublisher.post(itemStateEvent); + responseEvent = new EventDTO(WEBSOCKET_EVENT_TYPE, "/response/success", "", null, + eventDTO.eventId); + break; + case WEBSOCKET_EVENT_TYPE: + if ("/heartbeat".equals(eventDTO.topic) && "PING".equals(eventDTO.payload)) { + responseEvent = new EventDTO(WEBSOCKET_EVENT_TYPE, "/heartbeat", "PONG", null, + eventDTO.eventId); + } else if ("/filter/type".equals(eventDTO.topic)) { + typeFilter = Objects.requireNonNullElse(gson.fromJson(eventDTO.payload, STRING_LIST_TYPE), + List.of()); + logger.debug("Setting type filter for connection to {}: {}", + remoteEndpoint.getInetSocketAddress(), typeFilter); + responseEvent = new EventDTO(WEBSOCKET_EVENT_TYPE, "/filter/type", eventDTO.payload, null, + eventDTO.eventId); + } else if ("/filter/source".equals(eventDTO.topic)) { + sourceFilter = Objects.requireNonNullElse(gson.fromJson(eventDTO.payload, STRING_LIST_TYPE), + List.of()); + logger.debug("Setting source filter for connection to {}: {}", + remoteEndpoint.getInetSocketAddress(), typeFilter); + responseEvent = new EventDTO(WEBSOCKET_EVENT_TYPE, "/filter/source", eventDTO.payload, null, + eventDTO.eventId); + } else { + throw new EventProcessingException("Invalid topic or payload in WebSocketEvent"); + } + break; + default: + throw new EventProcessingException("Unknown event type '" + eventDTO.type + "'"); + } + if (!WEBSOCKET_EVENT_TYPE.equals(type) && responseEvent.eventId == null) { + // skip only for successful processing of state/command, always send response if processing failed + logger.trace("Not sending response event {}, because no eventId present.", responseEvent); + return; + } + } catch (EventProcessingException | JsonParseException e) { + logger.warn("Failed to process deserialized event '{}': {}", message, e.getMessage()); + responseEvent = new EventDTO(WEBSOCKET_EVENT_TYPE, "/response/failed", + "Processing error: " + e.getMessage(), null, eventDTO != null ? eventDTO.eventId : ""); + + } + } catch (JsonParseException e) { + logger.warn("Could not deserialize '{}'", message); + responseEvent = new EventDTO(WEBSOCKET_EVENT_TYPE, "/response/failed", + "Deserialization error: " + e.getMessage(), null, null); + } + + try { + sendMessage(gson.toJson(responseEvent)); + } catch (IOException e) { + logger.debug("Failed to send WebSocketResponseEvent event {} to {}: {}", responseEvent, remoteIdentifier, + e.getMessage()); + } + } + + @OnWebSocketError + public void onError(Session session, Throwable error) { + if (session != null) { + session.close(); + } + String message = error == null ? "" : error.getMessage(); + logger.info("WebSocket error: {}", message); + onClose(StatusCode.NO_CODE, message); + } + + public void processEvent(Event event) { + try { + String source = event.getSource(); + if ((source == null || !sourceFilter.contains(event.getSource())) + && (typeFilter.isEmpty() || typeFilter.contains(event.getType()))) { + sendMessage(gson.toJson(new EventDTO(event))); + } + } catch (IOException e) { + logger.debug("Failed to send event {} to {}: {}", event, remoteIdentifier, e.getMessage()); + } + } + + private synchronized void sendMessage(String message) throws IOException { + RemoteEndpoint remoteEndpoint = this.remoteEndpoint; + if (remoteEndpoint == null) { + logger.warn("Could not determine remote endpoint, failed to send '{}'.", message); + return; + } + remoteEndpoint.sendString(message); + } +} diff --git a/bundles/org.openhab.core.io.websocket/src/main/java/org/openhab/core/io/websocket/EventWebSocketServlet.java b/bundles/org.openhab.core.io.websocket/src/main/java/org/openhab/core/io/websocket/EventWebSocketServlet.java new file mode 100644 index 000000000..ee3de1b8a --- /dev/null +++ b/bundles/org.openhab.core.io.websocket/src/main/java/org/openhab/core/io/websocket/EventWebSocketServlet.java @@ -0,0 +1,168 @@ +/** + * Copyright (c) 2010-2022 Contributors to the openHAB project + * + * See the NOTICE file(s) distributed with this work for additional + * information. + * + * This program and the accompanying materials are made available under the + * terms of the Eclipse Public License 2.0 which is available at + * http://www.eclipse.org/legal/epl-2.0 + * + * SPDX-License-Identifier: EPL-2.0 + */ +package org.openhab.core.io.websocket; + +import java.util.Base64; +import java.util.List; +import java.util.Map; +import java.util.Set; +import java.util.concurrent.CopyOnWriteArraySet; + +import javax.servlet.ServletException; + +import org.eclipse.jdt.annotation.NonNullByDefault; +import org.eclipse.jdt.annotation.Nullable; +import org.eclipse.jetty.websocket.servlet.ServletUpgradeRequest; +import org.eclipse.jetty.websocket.servlet.ServletUpgradeResponse; +import org.eclipse.jetty.websocket.servlet.WebSocketCreator; +import org.eclipse.jetty.websocket.servlet.WebSocketServlet; +import org.eclipse.jetty.websocket.servlet.WebSocketServletFactory; +import org.openhab.core.auth.Authentication; +import org.openhab.core.auth.AuthenticationException; +import org.openhab.core.auth.Credentials; +import org.openhab.core.auth.Role; +import org.openhab.core.auth.User; +import org.openhab.core.auth.UserApiTokenCredentials; +import org.openhab.core.auth.UserRegistry; +import org.openhab.core.auth.UsernamePasswordCredentials; +import org.openhab.core.events.Event; +import org.openhab.core.events.EventFilter; +import org.openhab.core.events.EventPublisher; +import org.openhab.core.events.EventSubscriber; +import org.openhab.core.items.ItemRegistry; +import org.osgi.service.component.annotations.Activate; +import org.osgi.service.component.annotations.Component; +import org.osgi.service.component.annotations.Deactivate; +import org.osgi.service.component.annotations.Reference; +import org.osgi.service.http.HttpService; +import org.osgi.service.http.NamespaceException; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + +import com.google.gson.Gson; + +/** + * The {@link EventWebSocketServlet} provides the servlet for WebSocket connections + * + * @author Jan N. Klug - Initial contribution + */ +@Component(immediate = true, service = EventSubscriber.class) +@NonNullByDefault +public class EventWebSocketServlet extends WebSocketServlet implements EventSubscriber { + private static final long serialVersionUID = 1L; + + private final Gson gson = new Gson(); + private final HttpService httpService; + private final UserRegistry userRegistry; + private final EventPublisher eventPublisher; + + private final ItemEventUtility itemEventUtility; + private final Set webSockets = new CopyOnWriteArraySet<>(); + + @Activate + public EventWebSocketServlet(@Reference HttpService httpService, @Reference UserRegistry userRegistry, + @Reference EventPublisher eventPublisher, @Reference ItemRegistry itemRegistry) + throws ServletException, NamespaceException { + this.httpService = httpService; + this.userRegistry = userRegistry; + this.eventPublisher = eventPublisher; + + itemEventUtility = new ItemEventUtility(gson, itemRegistry); + + httpService.registerServlet("/ws", this, null, null); + } + + @Deactivate + public void deactivate() { + httpService.unregister("/ws"); + } + + @Override + public void configure(@NonNullByDefault({}) WebSocketServletFactory webSocketServletFactory) { + webSocketServletFactory.getPolicy().setIdleTimeout(10000); + webSocketServletFactory.setCreator(new EventWebSocketCreator()); + } + + @Override + public Set getSubscribedEventTypes() { + return Set.of(EventSubscriber.ALL_EVENT_TYPES); + } + + @Override + public @Nullable EventFilter getEventFilter() { + return null; + } + + @Override + public void receive(Event event) { + webSockets.forEach(ws -> ws.processEvent(event)); + } + + public void registerListener(EventWebSocket eventWebSocket) { + webSockets.add(eventWebSocket); + } + + public void unregisterListener(EventWebSocket eventWebSocket) { + webSockets.remove(eventWebSocket); + } + + private class EventWebSocketCreator implements WebSocketCreator { + private static final String API_TOKEN_PREFIX = "oh."; + + private final Logger logger = LoggerFactory.getLogger(EventWebSocketCreator.class); + + @Override + public @Nullable Object createWebSocket(@Nullable ServletUpgradeRequest servletUpgradeRequest, + @Nullable ServletUpgradeResponse servletUpgradeResponse) { + if (servletUpgradeRequest == null) { + return null; + } + + Map> parameterMap = servletUpgradeRequest.getParameterMap(); + List accessToken = parameterMap.getOrDefault("accessToken", List.of()); + if (accessToken.size() == 1 && authenticateAccessToken(accessToken.get(0))) { + return new EventWebSocket(gson, EventWebSocketServlet.this, itemEventUtility, eventPublisher); + } else { + logger.warn("Unauthenticated request to create a websocket from {}.", + servletUpgradeRequest.getRemoteAddress()); + } + + return null; + } + + private boolean authenticateAccessToken(String token) { + Credentials credentials = null; + if (token.startsWith(API_TOKEN_PREFIX)) { + credentials = new UserApiTokenCredentials(token); + } else { + // try BasicAuthentication + String[] decodedParts = Base64.getDecoder().decode(token).toString().split(":"); + if (decodedParts.length == 2) { + credentials = new UsernamePasswordCredentials(decodedParts[0], decodedParts[1]); + } + } + + if (credentials != null) { + try { + Authentication auth = userRegistry.authenticate(credentials); + User user = userRegistry.get(auth.getUsername()); + return user != null + && (user.getRoles().contains(Role.USER) || user.getRoles().contains(Role.ADMIN)); + } catch (AuthenticationException ignored) { + } + } + + return false; + } + } +} diff --git a/bundles/org.openhab.core.io.websocket/src/main/java/org/openhab/core/io/websocket/ItemEventUtility.java b/bundles/org.openhab.core.io.websocket/src/main/java/org/openhab/core/io/websocket/ItemEventUtility.java new file mode 100644 index 000000000..d1aeb8aad --- /dev/null +++ b/bundles/org.openhab.core.io.websocket/src/main/java/org/openhab/core/io/websocket/ItemEventUtility.java @@ -0,0 +1,138 @@ +/** + * Copyright (c) 2010-2022 Contributors to the openHAB project + * + * See the NOTICE file(s) distributed with this work for additional + * information. + * + * This program and the accompanying materials are made available under the + * terms of the Eclipse Public License 2.0 which is available at + * http://www.eclipse.org/legal/epl-2.0 + * + * SPDX-License-Identifier: EPL-2.0 + */ +package org.openhab.core.io.websocket; + +import java.util.List; +import java.util.regex.Matcher; +import java.util.regex.Pattern; + +import org.eclipse.jdt.annotation.NonNullByDefault; +import org.eclipse.jdt.annotation.Nullable; +import org.openhab.core.events.Event; +import org.openhab.core.items.Item; +import org.openhab.core.items.ItemNotFoundException; +import org.openhab.core.items.ItemRegistry; +import org.openhab.core.items.events.ItemEventFactory; +import org.openhab.core.types.Command; +import org.openhab.core.types.RefreshType; +import org.openhab.core.types.State; +import org.openhab.core.types.Type; +import org.openhab.core.types.TypeParser; +import org.openhab.core.types.UnDefType; + +import com.google.gson.Gson; +import com.google.gson.JsonParseException; + +/** + * The {@link EventDTO} is used for serialization and deserialization of events + * + * @author Stefan Bußweiler - Initial contribution ({@link org.openhab.core.items.events.ItemEventFactory}) + * @author Jan N. Klug - Initial contribution + */ +@NonNullByDefault +public class ItemEventUtility { + private static final Pattern TOPIC_PATTERN = Pattern.compile("openhab/items/(?\\w+)/(?\\w+)"); + private static final String TYPE_POSTFIX = "Type"; + + private final Gson gson; + private final ItemRegistry itemRegistry; + + public ItemEventUtility(Gson gson, ItemRegistry itemRegistry) { + this.gson = gson; + this.itemRegistry = itemRegistry; + } + + public Event createCommandEvent(EventDTO eventDTO) throws EventProcessingException { + Matcher matcher = getTopicMatcher(eventDTO.topic, "command"); + Item item = getItem(matcher.group("entity")); + Type command = parseType(eventDTO.payload); + if (command instanceof Command) { + List> acceptedItemCommandTypes = item.getAcceptedCommandTypes(); + if (acceptedItemCommandTypes.contains(command.getClass())) { + return ItemEventFactory.createCommandEvent(item.getName(), (Command) command, eventDTO.source); + } + } + throw new EventProcessingException("Incompatible datatype, rejected."); + } + + public Event createStateEvent(EventDTO eventDTO) throws EventProcessingException { + Matcher matcher = getTopicMatcher(eventDTO.topic, "state"); + Item item = getItem(matcher.group("entity")); + Type state = parseType(eventDTO.payload); + if (state instanceof State) { + List> acceptedItemStateTypes = item.getAcceptedDataTypes(); + if (acceptedItemStateTypes.contains(state.getClass())) { + return ItemEventFactory.createStateEvent(item.getName(), (State) state, eventDTO.source); + } + } + throw new EventProcessingException("Incompatible datatype, rejected."); + } + + private Matcher getTopicMatcher(@Nullable String topic, String action) throws EventProcessingException { + if (topic == null) { + throw new EventProcessingException("Topic must not be null"); + } + Matcher matcher = TOPIC_PATTERN.matcher(topic); + if (!matcher.matches()) { + throw new EventProcessingException( + "Topic must follow the format {namespace}/{entityType}/{entity}/{action}."); + } + + if (!action.equals(matcher.group("action"))) { + throw new EventProcessingException("Topic does not match event type."); + } + return matcher; + } + + private Item getItem(String itemName) throws EventProcessingException { + try { + return itemRegistry.getItem(itemName); + } catch (ItemNotFoundException e) { + throw new EventProcessingException("Could not find item '" + itemName + "' in registry."); + } + } + + private Type parseType(@Nullable String payload) throws EventProcessingException { + ItemEventPayloadBean bean = null; + try { + bean = gson.fromJson(payload, ItemEventPayloadBean.class); + } catch (JsonParseException ignored) { + } + if (bean == null) { + throw new EventProcessingException("Failed to deserialize payload '" + payload + "'."); + } + + String simpleClassName = bean.type + TYPE_POSTFIX; + Type returnType; + + if (simpleClassName.equals(UnDefType.class.getSimpleName())) { + returnType = UnDefType.valueOf(bean.value); + } else if (simpleClassName.equals(RefreshType.class.getSimpleName())) { + returnType = RefreshType.valueOf(bean.value); + } else { + returnType = TypeParser.parseType(simpleClassName, bean.value); + } + + if (returnType == null) { + throw new EventProcessingException( + "Error parsing simpleClassName '" + simpleClassName + "' with value '" + bean.value + "'."); + } + + return returnType; + } + + private static class ItemEventPayloadBean { + public @NonNullByDefault({}) String type; + public @NonNullByDefault({}) String value; + } +} diff --git a/bundles/org.openhab.core.io.websocket/src/test/java/org/openhab/core/io/websocket/EventWebSocketTest.java b/bundles/org.openhab.core.io.websocket/src/test/java/org/openhab/core/io/websocket/EventWebSocketTest.java new file mode 100644 index 000000000..c93aa3331 --- /dev/null +++ b/bundles/org.openhab.core.io.websocket/src/test/java/org/openhab/core/io/websocket/EventWebSocketTest.java @@ -0,0 +1,248 @@ +/** + * Copyright (c) 2010-2022 Contributors to the openHAB project + * + * See the NOTICE file(s) distributed with this work for additional + * information. + * + * This program and the accompanying materials are made available under the + * terms of the Eclipse Public License 2.0 which is available at + * http://www.eclipse.org/legal/epl-2.0 + * + * SPDX-License-Identifier: EPL-2.0 + */ +package org.openhab.core.io.websocket; + +import static org.mockito.ArgumentMatchers.any; +import static org.mockito.ArgumentMatchers.eq; +import static org.mockito.Mockito.never; +import static org.mockito.Mockito.times; +import static org.mockito.Mockito.verify; +import static org.mockito.Mockito.when; +import static org.openhab.core.io.websocket.EventWebSocket.WEBSOCKET_EVENT_TYPE; + +import java.io.IOException; +import java.net.InetSocketAddress; +import java.util.Objects; + +import org.eclipse.jdt.annotation.NonNullByDefault; +import org.eclipse.jdt.annotation.Nullable; +import org.eclipse.jetty.websocket.api.RemoteEndpoint; +import org.eclipse.jetty.websocket.api.Session; +import org.eclipse.jetty.websocket.api.StatusCode; +import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.Test; +import org.junit.jupiter.api.extension.ExtendWith; +import org.mockito.Mock; +import org.mockito.junit.jupiter.MockitoExtension; +import org.mockito.junit.jupiter.MockitoSettings; +import org.mockito.quality.Strictness; +import org.openhab.core.events.Event; +import org.openhab.core.events.EventPublisher; +import org.openhab.core.items.ItemNotFoundException; +import org.openhab.core.items.ItemRegistry; +import org.openhab.core.items.events.ItemEventFactory; +import org.openhab.core.library.items.NumberItem; +import org.openhab.core.library.types.DecimalType; + +import com.google.gson.Gson; + +/** + * The {@link EventWebSocketTest} contains tests for the {@link EventWebSocket} + * + * @author Jan N. Klug - Initial contribution + */ +@NonNullByDefault +@ExtendWith(MockitoExtension.class) +@MockitoSettings(strictness = Strictness.LENIENT) +public class EventWebSocketTest { + private static final String REMOTE_WEBSOCKET_IMPLEMENTATION = "fooWebsocket"; + + private static final String TEST_ITEM_NAME = "testItem"; + private static final NumberItem TEST_ITEM = new NumberItem(TEST_ITEM_NAME); + + private Gson gson = new Gson(); + + private @Mock @NonNullByDefault({}) EventWebSocketServlet servlet; + private @Mock @NonNullByDefault({}) ItemRegistry itemRegistry; + private @Mock @NonNullByDefault({}) EventPublisher eventPublisher; + private @Mock @NonNullByDefault({}) Session session; + private @Mock @NonNullByDefault({}) RemoteEndpoint remoteEndpoint; + + private @NonNullByDefault({}) ItemEventUtility itemEventUtility; + private @NonNullByDefault({}) EventWebSocket eventWebSocket; + + @BeforeEach + public void setup() throws ItemNotFoundException { + itemEventUtility = new ItemEventUtility(gson, itemRegistry); + eventWebSocket = new EventWebSocket(gson, servlet, itemEventUtility, eventPublisher); + + when(session.getRemote()).thenReturn(remoteEndpoint); + when(remoteEndpoint.getInetSocketAddress()).thenReturn(new InetSocketAddress(47115)); + + when(itemRegistry.getItem(eq(TEST_ITEM_NAME))).thenReturn(TEST_ITEM); + + eventWebSocket.onConnect(session); + verify(servlet).registerListener(eventWebSocket); + } + + @Test + public void listenerCorrectlyUnregisteredOnClose() { + eventWebSocket.onClose(StatusCode.NORMAL, "Normal close."); + + verify(servlet).unregisterListener(eventWebSocket); + } + + @Test + public void sessionClosesOnErrorAndOnCloseCalled() { + eventWebSocket.onError(session, new IllegalStateException()); + + verify(session).close(); + verify(servlet).unregisterListener(eventWebSocket); + } + + @Test + public void stateEventWithIdFromWebsocketIsPublishedAndConfirmed() throws IOException { + Event expectedEvent = ItemEventFactory.createStateEvent(TEST_ITEM_NAME, DecimalType.ZERO, + REMOTE_WEBSOCKET_IMPLEMENTATION); + EventDTO eventDTO = new EventDTO(expectedEvent); + eventDTO.eventId = "id-1"; + EventDTO expectedResponse = new EventDTO(WEBSOCKET_EVENT_TYPE, "/response/success", "", null, eventDTO.eventId); + + assertEventProcessing(eventDTO, expectedEvent, expectedResponse); + } + + @Test + public void stateEventWithoutIdFromWebsocketIsPublished() throws IOException { + Event expectedEvent = ItemEventFactory.createStateEvent(TEST_ITEM_NAME, DecimalType.ZERO, + REMOTE_WEBSOCKET_IMPLEMENTATION); + EventDTO eventDTO = new EventDTO(expectedEvent); + + assertEventProcessing(eventDTO, expectedEvent, null); + } + + @Test + public void commandEventWithIdFromWebsocketIsPublishedAndConfirmed() throws IOException { + Event expectedEvent = ItemEventFactory.createCommandEvent(TEST_ITEM_NAME, DecimalType.ZERO, + REMOTE_WEBSOCKET_IMPLEMENTATION); + + EventDTO eventDTO = new EventDTO(expectedEvent); + eventDTO.eventId = "id-1"; + EventDTO expectedResponse = new EventDTO(WEBSOCKET_EVENT_TYPE, "/response/success", "", null, eventDTO.eventId); + + assertEventProcessing(eventDTO, expectedEvent, expectedResponse); + } + + @Test + public void commandEventWithoutIdFromWebsocketIsPublished() throws IOException { + Event expectedEvent = ItemEventFactory.createCommandEvent(TEST_ITEM_NAME, DecimalType.ZERO, + REMOTE_WEBSOCKET_IMPLEMENTATION); + + EventDTO eventDTO = new EventDTO(expectedEvent); + assertEventProcessing(eventDTO, expectedEvent, null); + } + + @Test + public void illegalStateEventNotPublishedAndResponseSent() throws IOException { + Event expectedEvent = ItemEventFactory.createStateEvent(TEST_ITEM_NAME, DecimalType.ZERO, + REMOTE_WEBSOCKET_IMPLEMENTATION); + + EventDTO eventDTO = new EventDTO(expectedEvent); + eventDTO.payload = ""; + + EventDTO expectedResponse = new EventDTO(WEBSOCKET_EVENT_TYPE, "/response/failed", + "Processing error: Failed to deserialize payload \u0027\u0027.", null, null); + + assertEventProcessing(eventDTO, null, expectedResponse); + } + + @Test + public void illegalCommandEventNotPublishedAndResponseSent() throws IOException { + Event expectedEvent = ItemEventFactory.createCommandEvent(TEST_ITEM_NAME, DecimalType.ZERO, + REMOTE_WEBSOCKET_IMPLEMENTATION); + + EventDTO eventDTO = new EventDTO(expectedEvent); + eventDTO.eventId = "id-1"; + eventDTO.topic = ""; + + EventDTO expectedResponse = new EventDTO(WEBSOCKET_EVENT_TYPE, "/response/failed", + "Processing error: Topic must follow the format {namespace}/{entityType}/{entity}/{action}.", null, + eventDTO.eventId); + + assertEventProcessing(eventDTO, null, expectedResponse); + } + + @Test + public void heartBeat() throws IOException { + EventDTO eventDTO = new EventDTO(WEBSOCKET_EVENT_TYPE, "/heartbeat", "PING", null, null); + EventDTO expectedResponse = new EventDTO(WEBSOCKET_EVENT_TYPE, "/heartbeat", "PONG", null, null); + + assertEventProcessing(eventDTO, null, expectedResponse); + } + + @Test + public void eventFromBusSent() throws IOException { + Event event = ItemEventFactory.createCommandEvent(TEST_ITEM_NAME, DecimalType.ZERO, + REMOTE_WEBSOCKET_IMPLEMENTATION); + + eventWebSocket.processEvent(event); + EventDTO eventDTO = new EventDTO(event); + + verify(remoteEndpoint).sendString(gson.toJson(eventDTO)); + } + + @Test + public void eventFromBusFilterType() throws IOException { + EventDTO eventDTO = new EventDTO(WEBSOCKET_EVENT_TYPE, "/filter/type", "[\"ItemCommandEvent\"]", null, null); + EventDTO responseEventDTO = new EventDTO(WEBSOCKET_EVENT_TYPE, "/filter/type", eventDTO.payload, null, null); + eventWebSocket.onText(gson.toJson(eventDTO)); + verify(remoteEndpoint).sendString(gson.toJson(responseEventDTO)); + + // subscribed type is sent + Event event = ItemEventFactory.createCommandEvent(TEST_ITEM_NAME, DecimalType.ZERO, + REMOTE_WEBSOCKET_IMPLEMENTATION); + eventWebSocket.processEvent(event); + verify(remoteEndpoint).sendString(gson.toJson(new EventDTO(event))); + + // not subscribed event not sent + event = ItemEventFactory.createStateEvent(TEST_ITEM_NAME, DecimalType.ZERO, REMOTE_WEBSOCKET_IMPLEMENTATION); + eventWebSocket.processEvent(event); + verify(remoteEndpoint, times(2)).sendString(any()); + } + + @Test + public void eventFromBusFilterSource() throws IOException { + EventDTO eventDTO = new EventDTO(WEBSOCKET_EVENT_TYPE, "/filter/source", + "[\"" + REMOTE_WEBSOCKET_IMPLEMENTATION + "\"]", null, null); + EventDTO responseEventDTO = new EventDTO(WEBSOCKET_EVENT_TYPE, "/filter/source", eventDTO.payload, null, null); + eventWebSocket.onText(gson.toJson(eventDTO)); + verify(remoteEndpoint).sendString(gson.toJson(responseEventDTO)); + + // non-matching is sent + Event event = ItemEventFactory.createCommandEvent(TEST_ITEM_NAME, DecimalType.ZERO); + eventWebSocket.processEvent(event); + verify(remoteEndpoint).sendString(gson.toJson(new EventDTO(event))); + + // matching is not sent + event = ItemEventFactory.createStateEvent(TEST_ITEM_NAME, DecimalType.ZERO, REMOTE_WEBSOCKET_IMPLEMENTATION); + eventWebSocket.processEvent(event); + verify(remoteEndpoint, times(2)).sendString(any()); + } + + private void assertEventProcessing(EventDTO incoming, @Nullable Event expectedEvent, + @Nullable EventDTO expectedResponse) throws IOException { + eventWebSocket.onText(gson.toJson(incoming)); + + if (expectedEvent != null) { + verify(eventPublisher).post(eq(Objects.requireNonNull(expectedEvent))); + } else { + verify(eventPublisher, never()).post(any()); + } + + if (expectedResponse != null) { + String expectedResponseString = gson.toJson(expectedResponse); + verify(remoteEndpoint).sendString(eq(expectedResponseString)); + } else { + verify(remoteEndpoint, never()).sendString(any()); + } + } +} diff --git a/bundles/org.openhab.core.io.websocket/src/test/java/org/openhab/core/io/websocket/ItemEventUtilityTest.java b/bundles/org.openhab.core.io.websocket/src/test/java/org/openhab/core/io/websocket/ItemEventUtilityTest.java new file mode 100644 index 000000000..b6dfa1419 --- /dev/null +++ b/bundles/org.openhab.core.io.websocket/src/test/java/org/openhab/core/io/websocket/ItemEventUtilityTest.java @@ -0,0 +1,171 @@ +/** + * Copyright (c) 2010-2022 Contributors to the openHAB project + * + * See the NOTICE file(s) distributed with this work for additional + * information. + * + * This program and the accompanying materials are made available under the + * terms of the Eclipse Public License 2.0 which is available at + * http://www.eclipse.org/legal/epl-2.0 + * + * SPDX-License-Identifier: EPL-2.0 + */ +package org.openhab.core.io.websocket; + +import static org.hamcrest.MatcherAssert.assertThat; +import static org.hamcrest.Matchers.is; +import static org.junit.jupiter.api.Assertions.assertThrows; +import static org.mockito.ArgumentMatchers.eq; +import static org.mockito.Mockito.when; + +import org.eclipse.jdt.annotation.NonNullByDefault; +import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.Test; +import org.junit.jupiter.api.extension.ExtendWith; +import org.mockito.Mock; +import org.mockito.junit.jupiter.MockitoExtension; +import org.mockito.junit.jupiter.MockitoSettings; +import org.mockito.quality.Strictness; +import org.openhab.core.events.Event; +import org.openhab.core.items.ItemNotFoundException; +import org.openhab.core.items.ItemRegistry; +import org.openhab.core.items.events.ItemEvent; +import org.openhab.core.items.events.ItemEventFactory; +import org.openhab.core.library.items.StringItem; +import org.openhab.core.library.types.DecimalType; +import org.openhab.core.library.types.HSBType; +import org.openhab.core.library.types.StringType; + +import com.google.gson.Gson; + +/** + * The {@link ItemEventUtilityTest} contains tests for the {@link ItemEventUtility} class. + * + * @author Jan N. Klug - Initial contribution + */ +@NonNullByDefault +@ExtendWith(MockitoExtension.class) +@MockitoSettings(strictness = Strictness.LENIENT) +public class ItemEventUtilityTest { + private static final String EXISTING_ITEM_NAME = "existingItem"; + private static final String NON_EXISTING_ITEM_NAME = "nonExistingItem"; + private static final StringType ITEM_STATE = new StringType("foo"); + + private @Mock @NonNullByDefault({}) ItemRegistry itemRegistry; + + private StringItem existingItem = new StringItem(EXISTING_ITEM_NAME); + private Gson gson = new Gson(); + + private @NonNullByDefault({}) ItemEventUtility itemEventUtility; + + @BeforeEach + public void setUp() throws ItemNotFoundException { + itemEventUtility = new ItemEventUtility(gson, itemRegistry); + + when(itemRegistry.getItem(eq(EXISTING_ITEM_NAME))).thenReturn(existingItem); + when(itemRegistry.getItem(eq(NON_EXISTING_ITEM_NAME))) + .thenThrow(new ItemNotFoundException(NON_EXISTING_ITEM_NAME)); + } + + @Test + public void validStateEvent() throws EventProcessingException { + ItemEvent event = ItemEventFactory.createStateEvent(EXISTING_ITEM_NAME, ITEM_STATE); + EventDTO eventDTO = new EventDTO(event); + + Event itemEvent = itemEventUtility.createStateEvent(eventDTO); + + assertThat(itemEvent, is(event)); + } + + @Test + public void validStateEventWithMissingItem() { + ItemEvent event = ItemEventFactory.createStateEvent(NON_EXISTING_ITEM_NAME, ITEM_STATE); + EventDTO eventDTO = new EventDTO(event); + + EventProcessingException e = assertThrows(EventProcessingException.class, + () -> itemEventUtility.createStateEvent(eventDTO)); + assertThat(e.getMessage(), is("Could not find item '" + NON_EXISTING_ITEM_NAME + "' in registry.")); + } + + @Test + public void validStateEventWithInvalidState() { + ItemEvent event = ItemEventFactory.createStateEvent(EXISTING_ITEM_NAME, DecimalType.ZERO); + EventDTO eventDTO = new EventDTO(event); + + EventProcessingException e = assertThrows(EventProcessingException.class, + () -> itemEventUtility.createStateEvent(eventDTO)); + assertThat(e.getMessage(), is("Incompatible datatype, rejected.")); + } + + @Test + public void invalidStateEventTopic() { + ItemEvent event = ItemEventFactory.createCommandEvent(EXISTING_ITEM_NAME, HSBType.BLACK); + EventDTO eventDTO = new EventDTO(event); + + EventProcessingException e = assertThrows(EventProcessingException.class, + () -> itemEventUtility.createStateEvent(eventDTO)); + assertThat(e.getMessage(), is("Topic does not match event type.")); + } + + @Test + public void invalidStateEventPayload() { + ItemEvent event = ItemEventFactory.createStateEvent(EXISTING_ITEM_NAME, HSBType.BLACK); + EventDTO eventDTO = new EventDTO(event); + eventDTO.payload = "invalidNoJson"; + + EventProcessingException e = assertThrows(EventProcessingException.class, + () -> itemEventUtility.createStateEvent(eventDTO)); + assertThat(e.getMessage(), is("Failed to deserialize payload 'invalidNoJson'.")); + } + + @Test + public void validCommandEvent() throws EventProcessingException { + ItemEvent event = ItemEventFactory.createCommandEvent(EXISTING_ITEM_NAME, ITEM_STATE); + EventDTO eventDTO = new EventDTO(event); + + Event itemEvent = itemEventUtility.createCommandEvent(eventDTO); + + assertThat(itemEvent, is(event)); + } + + @Test + public void validCommandEventWithMissingItem() { + ItemEvent event = ItemEventFactory.createStateEvent(NON_EXISTING_ITEM_NAME, ITEM_STATE); + EventDTO eventDTO = new EventDTO(event); + + EventProcessingException e = assertThrows(EventProcessingException.class, + () -> itemEventUtility.createStateEvent(eventDTO)); + assertThat(e.getMessage(), is("Could not find item '" + NON_EXISTING_ITEM_NAME + "' in registry.")); + } + + @Test + public void validCommandEventWithInvalidState() { + ItemEvent event = ItemEventFactory.createCommandEvent(EXISTING_ITEM_NAME, HSBType.BLACK); + EventDTO eventDTO = new EventDTO(event); + + EventProcessingException e = assertThrows(EventProcessingException.class, + () -> itemEventUtility.createCommandEvent(eventDTO)); + assertThat(e.getMessage(), is("Incompatible datatype, rejected.")); + } + + @Test + public void invalidCommandEvent() { + ItemEvent event = ItemEventFactory.createStateEvent(EXISTING_ITEM_NAME, HSBType.BLACK); + EventDTO eventDTO = new EventDTO(event); + + EventProcessingException e = assertThrows(EventProcessingException.class, + () -> itemEventUtility.createCommandEvent(eventDTO)); + assertThat(e.getMessage(), is("Topic does not match event type.")); + } + + @Test + public void invalidCommandEventPayload() { + ItemEvent event = ItemEventFactory.createCommandEvent(EXISTING_ITEM_NAME, HSBType.BLACK); + EventDTO eventDTO = new EventDTO(event); + eventDTO.payload = "invalidNoJson"; + + EventProcessingException e = assertThrows(EventProcessingException.class, + () -> itemEventUtility.createCommandEvent(eventDTO)); + assertThat(e.getMessage(), is("Failed to deserialize payload 'invalidNoJson'.")); + } +} diff --git a/bundles/pom.xml b/bundles/pom.xml index 9d580fe0e..3b2f484ec 100644 --- a/bundles/pom.xml +++ b/bundles/pom.xml @@ -78,6 +78,7 @@ org.openhab.core.io.transport.serial.rxtx org.openhab.core.io.transport.serial.rxtx.rfc2217 org.openhab.core.io.transport.upnp + org.openhab.core.io.websocket org.openhab.core.io.jetty.certificate org.openhab.core.model.lazygen org.openhab.core.model.core diff --git a/features/karaf/openhab-core/src/main/feature/feature.xml b/features/karaf/openhab-core/src/main/feature/feature.xml index 33706e25e..d613ac90a 100644 --- a/features/karaf/openhab-core/src/main/feature/feature.xml +++ b/features/karaf/openhab-core/src/main/feature/feature.xml @@ -177,6 +177,11 @@ mvn:org.openhab.core.bundles/org.openhab.core.io.rest.mdns/${project.version} + + openhab-core-base + mvn:org.openhab.core.bundles/org.openhab.core.io.websocket/${project.version} + + openhab-core-base @@ -395,6 +400,7 @@ openhab-core-io-rest-swagger openhab-core-io-rest-transform openhab-core-io-rest-voice + openhab-core-io-websocket openhab-core-model-lsp openhab-core-model-item openhab-core-model-persistence