This is an automated email from the ASF dual-hosted git repository.

solomax pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/openmeetings.git


The following commit(s) were added to refs/heads/master by this push:
     new 99ad05a  [OPENMEETINGS-2186] first step to working cluster
99ad05a is described below

commit 99ad05ae18120f322bfb9998bf666a75ea12592b
Author: Maxim Solodovnik <[email protected]>
AuthorDate: Tue Mar 24 21:14:43 2020 +0700

    [OPENMEETINGS-2186] first step to working cluster
---
 .../apache/openmeetings/db/dao/room/RoomDao.java   |   6 --
 .../apache/openmeetings/db/entity/room/Room.java   |   1 -
 .../openmeetings/db/manager/IClientManager.java    |   9 --
 .../apache/openmeetings/web/app/Application.java   |  39 ++++----
 .../apache/openmeetings/web/app/ClientManager.java | 109 +++++++++++++++++++--
 .../openmeetings/web/util/OmUrlFragment.java       |  13 +++
 .../src/main/webapp/WEB-INF/classes/hazelcast.xml  |  11 ++-
 7 files changed, 145 insertions(+), 43 deletions(-)

diff --git 
a/openmeetings-db/src/main/java/org/apache/openmeetings/db/dao/room/RoomDao.java
 
b/openmeetings-db/src/main/java/org/apache/openmeetings/db/dao/room/RoomDao.java
index e5fc73b..f8cdf7d 100644
--- 
a/openmeetings-db/src/main/java/org/apache/openmeetings/db/dao/room/RoomDao.java
+++ 
b/openmeetings-db/src/main/java/org/apache/openmeetings/db/dao/room/RoomDao.java
@@ -29,7 +29,6 @@ import static 
org.apache.openmeetings.util.OpenmeetingsVariables.isSipEnabled;
 
 import java.util.ArrayList;
 import java.util.Calendar;
-import java.util.Collection;
 import java.util.Date;
 import java.util.HashSet;
 import java.util.List;
@@ -178,11 +177,6 @@ public class RoomDao implements 
IGroupAdminDataProviderDao<Room> {
                                .getResultList();
        }
 
-       public long getRoomsCapacityByIds(Collection<Long> ids) {
-               return ids == null || ids.isEmpty() ? 0L
-                       : em.createNamedQuery("getRoomsCapacityByIds", 
Long.class).setParameter("ids", ids).getSingleResult();
-       }
-
        private String getSipNumber(long roomId) {
                if (isSipEnabled()) {
                        return cfgDao.getString(CONFIG_SIP_ROOM_PREFIX, "400") 
+ roomId;
diff --git 
a/openmeetings-db/src/main/java/org/apache/openmeetings/db/entity/room/Room.java
 
b/openmeetings-db/src/main/java/org/apache/openmeetings/db/entity/room/Room.java
index 694facb..83ffcfb 100644
--- 
a/openmeetings-db/src/main/java/org/apache/openmeetings/db/entity/room/Room.java
+++ 
b/openmeetings-db/src/main/java/org/apache/openmeetings/db/entity/room/Room.java
@@ -84,7 +84,6 @@ import org.apache.openmeetings.db.entity.user.Group;
 @NamedQuery(name = "getSipRoomIdsByIds", query = "SELECT r.id FROM Room r 
WHERE r.deleted = false AND r.sipEnabled = true AND r.id IN :ids")
 @NamedQuery(name = "countRooms", query = "SELECT COUNT(r) FROM Room r WHERE 
r.deleted = false")
 @NamedQuery(name = "getBackupRooms", query = "SELECT r FROM Room r ORDER BY 
r.id")
-@NamedQuery(name = "getRoomsCapacityByIds", query = "SELECT SUM(r.capacity) 
FROM Room r WHERE r.deleted = false AND r.id IN :ids")
 @NamedQuery(name = "getGroupRooms", query = "SELECT DISTINCT rg.room FROM 
RoomGroup rg LEFT JOIN FETCH rg.room "
                + "WHERE rg.group.id = :groupId AND rg.room.deleted = false AND 
rg.room.appointment = false "
                + "ORDER BY rg.room.name ASC")
diff --git 
a/openmeetings-db/src/main/java/org/apache/openmeetings/db/manager/IClientManager.java
 
b/openmeetings-db/src/main/java/org/apache/openmeetings/db/manager/IClientManager.java
index 274f1d5..24a1498 100644
--- 
a/openmeetings-db/src/main/java/org/apache/openmeetings/db/manager/IClientManager.java
+++ 
b/openmeetings-db/src/main/java/org/apache/openmeetings/db/manager/IClientManager.java
@@ -20,7 +20,6 @@ package org.apache.openmeetings.db.manager;
 
 import java.util.Collection;
 import java.util.List;
-import java.util.Set;
 
 import org.apache.openmeetings.db.entity.basic.Client;
 
@@ -33,12 +32,4 @@ public interface IClientManager {
        Collection<Client> listByUser(Long userId);
        Client update(Client c);
        void exit(Client c);
-
-
-       /**
-        * Get a list of all rooms with users in the system.
-        *
-        * @return a set, a roomId can be only one time in this list
-        */
-       Set<Long> getActiveRoomIds();
 }
diff --git 
a/openmeetings-web/src/main/java/org/apache/openmeetings/web/app/Application.java
 
b/openmeetings-web/src/main/java/org/apache/openmeetings/web/app/Application.java
index 9feace7..9253c2e 100644
--- 
a/openmeetings-web/src/main/java/org/apache/openmeetings/web/app/Application.java
+++ 
b/openmeetings-web/src/main/java/org/apache/openmeetings/web/app/Application.java
@@ -35,7 +35,6 @@ import static 
org.wicketstuff.dashboard.DashboardContextInitializer.DASHBOARD_CO
 import java.io.File;
 import java.net.UnknownHostException;
 import java.text.MessageFormat;
-import java.util.ArrayList;
 import java.util.HashSet;
 import java.util.List;
 import java.util.Locale;
@@ -159,6 +158,7 @@ public class Application extends 
AuthenticatedWebApplication implements IApplica
        private static final String INVALID_SESSIONS_KEY = 
"INVALID_SESSIONS_KEY";
        private static final String SERVER_CONTAINER_SERVLET_CONTEXT_ATTRIBUTE 
= "javax.websocket.server.ServerContainer";
        public static final String NAME_ATTR_KEY = "name";
+       public static final String SERVER_URL_ATTR_KEY = "server.url";
        //additional maps for faster searching should be created
        private static final Set<String> STRINGS_WITH_APP = new HashSet<>();
        private static String appName;
@@ -171,6 +171,7 @@ public class Application extends 
AuthenticatedWebApplication implements IApplica
        public static final String NOTINIT_MAPPING = "/notinited";
        final HazelcastInstance hazelcast = 
Hazelcast.getOrCreateHazelcastInstance(new XmlConfigBuilder().build());
        private ITopic<IClusterWsMessage> hazelWsTopic;
+       private String serverId;
 
        @Autowired
        private ApplicationContext ctx;
@@ -194,11 +195,15 @@ public class Application extends 
AuthenticatedWebApplication implements IApplica
                
getApplicationSettings().setAccessDeniedPage(AccessDeniedPage.class);
                getComponentInstantiationListeners().add(new 
SpringComponentInjector(this, ctx, true));
 
-               
hazelcast.getCluster().getLocalMember().setStringAttribute(NAME_ATTR_KEY, 
hazelcast.getName());
+               serverId = hazelcast.getName();
+               
hazelcast.getCluster().getLocalMember().setStringAttribute(NAME_ATTR_KEY, 
serverId);
+               hazelcast.getCluster().getMembers().forEach(m -> {
+                       cm.serverAdded(m.getStringAttribute(NAME_ATTR_KEY), 
m.getStringAttribute(SERVER_URL_ATTR_KEY));
+               });
                hazelWsTopic = hazelcast.getTopic("default");
                hazelWsTopic.addMessageListener(msg -> {
-                       String serverId = 
msg.getPublishingMember().getStringAttribute(NAME_ATTR_KEY);
-                       if (serverId.equals(hazelcast.getName())) {
+                       String mServerId = 
msg.getPublishingMember().getStringAttribute(NAME_ATTR_KEY);
+                       if (mServerId.equals(serverId)) {
                                return;
                        }
                        IClusterWsMessage wsMsg = msg.getMessageObject();
@@ -210,23 +215,18 @@ public class Application extends 
AuthenticatedWebApplication implements IApplica
                        }
                        WebSocketHelper.send(msg.getMessageObject());
                });
-               //FIXME TODO 
hazelcast.getConfig().getProperties().getProperty("server.url")
                hazelcast.getCluster().addMembershipListener(new 
MembershipListener() {
                        @Override
                        public void memberRemoved(MembershipEvent evt) {
                                //server down, need to remove all online 
clients, process persistent addresses
                                String serverId = 
evt.getMember().getStringAttribute(NAME_ATTR_KEY);
-                               cm.clean(serverId);
+                               cm.serverRemoved(serverId);
                                updateJpaAddresses();
                        }
 
                        @Override
                        public void memberAttributeChanged(MemberAttributeEvent 
evt) {
                                //no-op
-                       }
-
-                       @Override
-                       public void memberAdded(MembershipEvent evt) {
                                //server added, need to process persistent 
addresses
                                updateJpaAddresses();
                                //check for duplicate instance-names
@@ -238,11 +238,18 @@ public class Application extends 
AuthenticatedWebApplication implements IApplica
                                        String serverId = 
m.getStringAttribute(NAME_ATTR_KEY);
                                        names.add(serverId);
                                }
-                               String serverId = 
evt.getMember().getStringAttribute(NAME_ATTR_KEY);
-                               if (names.contains(serverId)) {
-                                       log.warn("Duplicate cluster instance 
with name {} found {}", serverId, evt.getMember());
+                               String newServerId = 
evt.getMember().getStringAttribute(NAME_ATTR_KEY);
+                               log.warn("Name added: {}", newServerId);
+                               cm.serverAdded(newServerId, 
evt.getMember().getStringAttribute(SERVER_URL_ATTR_KEY));
+                               if (names.contains(newServerId)) {
+                                       log.warn("Duplicate cluster instance 
with name {} found {}", newServerId, evt.getMember());
                                }
                        }
+
+                       @Override
+                       public void memberAdded(MembershipEvent evt) {
+                               //no-op due to name is not set
+                       }
                });
                setPageManagerProvider(new DefaultPageManagerProvider(this) {
                        @Override
@@ -615,11 +622,7 @@ public class Application extends 
AuthenticatedWebApplication implements IApplica
 
        @Override
        public String getServerId() {
-               return hazelcast.getName();
-       }
-
-       public List<Member> getServers() {
-               return new ArrayList<>(hazelcast.getCluster().getMembers());
+               return serverId;
        }
 
        @Override
diff --git 
a/openmeetings-web/src/main/java/org/apache/openmeetings/web/app/ClientManager.java
 
b/openmeetings-web/src/main/java/org/apache/openmeetings/web/app/ClientManager.java
index 0157faf..42f5229 100644
--- 
a/openmeetings-web/src/main/java/org/apache/openmeetings/web/app/ClientManager.java
+++ 
b/openmeetings-web/src/main/java/org/apache/openmeetings/web/app/ClientManager.java
@@ -20,12 +20,14 @@ package org.apache.openmeetings.web.app;
 
 import static org.apache.openmeetings.core.util.WebSocketHelper.sendRoom;
 
+import java.io.Serializable;
 import java.util.ArrayList;
 import java.util.Collection;
 import java.util.HashSet;
 import java.util.List;
 import java.util.Map;
 import java.util.Map.Entry;
+import java.util.Optional;
 import java.util.Set;
 import java.util.concurrent.ConcurrentHashMap;
 import java.util.function.Predicate;
@@ -34,10 +36,10 @@ import java.util.stream.Collectors;
 import javax.annotation.PostConstruct;
 
 import org.apache.openmeetings.core.remote.KurentoHandler;
-import org.apache.openmeetings.core.util.WebSocketHelper;
 import org.apache.openmeetings.db.dao.log.ConferenceLogDao;
 import org.apache.openmeetings.db.entity.basic.Client;
 import org.apache.openmeetings.db.entity.log.ConferenceLog;
+import org.apache.openmeetings.db.entity.room.Room;
 import org.apache.openmeetings.db.manager.IClientManager;
 import org.apache.openmeetings.db.util.ws.RoomMessage;
 import org.apache.openmeetings.db.util.ws.TextRoomMessage;
@@ -59,9 +61,11 @@ public class ClientManager implements IClientManager {
        private static final Logger log = 
LoggerFactory.getLogger(ClientManager.class);
        private static final String ROOMS_KEY = "ROOMS_KEY";
        private static final String ONLINE_USERS_KEY = "ONLINE_USERS_KEY";
+       private static final String SERVERS_KEY = "SERVERS_KEY";
        private static final String UID_BY_SID_KEY = "UID_BY_SID_KEY";
        private final Map<String, Client> onlineClients = new 
ConcurrentHashMap<>();
        private final Map<Long, Set<String>> onlineRooms = new 
ConcurrentHashMap<>();
+       private final Map<String, ServerInfo> onlineServers = new 
ConcurrentHashMap<>();
 
        @Autowired
        private ConferenceLogDao confLogDao;
@@ -82,10 +86,22 @@ public class ClientManager implements IClientManager {
                return app.hazelcast.getMap(ROOMS_KEY);
        }
 
+       private IMap<String, ServerInfo> servers() {
+               return app.hazelcast.getMap(SERVERS_KEY);
+       }
+
        @PostConstruct
        void init() {
                map().addEntryListener(new ClientListener(), true);
                rooms().addEntryListener(new RoomListener(), true);
+               servers().addEntryListener(new EntryUpdatedListener<String, 
ServerInfo>() {
+
+                       @Override
+                       public void entryUpdated(EntryEvent<String, ServerInfo> 
event) {
+                               log.trace("ServerListener::Update");
+                               onlineServers.put(event.getKey(), 
event.getValue());
+                       }
+               }, true);
        }
 
        public void add(Client c) {
@@ -162,18 +178,21 @@ public class ClientManager implements IClientManager {
                }
        }
 
-       public void clean(String serverId) {
+       public void serverAdded(String serverId, String url) {
+               ServerInfo si = new ServerInfo(url);
+               servers().put(serverId, si);
+               onlineServers.put(serverId, si);
+       }
+
+       public void serverRemoved(String serverId) {
                Map<String, Client> clients = map();
                for (Map.Entry<String, Client> e : clients.entrySet()) {
                        if (serverId.equals(e.getValue().getServerId())) {
                                exit(e.getValue());
                        }
                }
-       }
-
-       @Override
-       public Set<Long> getActiveRoomIds() {
-               return onlineRooms.keySet();
+               servers().remove(serverId);
+               onlineServers.remove(serverId);
        }
 
        /**
@@ -183,7 +202,8 @@ public class ClientManager implements IClientManager {
         * @return count of users in room _after_ adding
         */
        public int addToRoom(Client c) {
-               Long roomId = c.getRoom().getId();
+               Room r = c.getRoom();
+               Long roomId = r.getId();
                confLogDao.add(
                                ConferenceLog.Type.ROOM_ENTER
                                , c.getUserId(), "0", roomId
@@ -199,6 +219,16 @@ public class ClientManager implements IClientManager {
                rooms.put(roomId, set);
                onlineRooms.put(roomId, set);
                rooms.unlock(roomId);
+               String serverId = c.getServerId();
+               if (!onlineServers.get(serverId).getRooms().contains(roomId)) {
+                       IMap<String, ServerInfo> servers = servers();
+                       servers.lock(serverId);
+                       ServerInfo si = servers.get(serverId);
+                       si.add(r);
+                       servers.put(serverId, si);
+                       onlineServers.put(serverId, si);
+                       servers.unlock(serverId);
+               }
                update(c);
                return count;
        }
@@ -216,6 +246,16 @@ public class ClientManager implements IClientManager {
                                onlineRooms.put(roomId, clients);
                        }
                        rooms.unlock(roomId);
+                       if (clients == null || clients.isEmpty()) {
+                               String serverId = c.getServerId();
+                               IMap<String, ServerInfo> servers = servers();
+                               servers.lock(serverId);
+                               ServerInfo si = servers.get(serverId);
+                               si.remove(c.getRoom());
+                               servers.put(serverId, si);
+                               onlineServers.put(serverId, si);
+                               servers.unlock(serverId);
+                       }
                        kHandler.leaveRoom(c);
                        c.setRoom(null);
                        c.clear();
@@ -306,6 +346,24 @@ public class ClientManager implements IClientManager {
                }
        }
 
+       public String getServerUrl(Long roomId) {
+               if (roomId == null || onlineServers.size() == 1) {
+                       return null;
+               }
+               final String curServerId = app.getServerId();
+               Optional<Map.Entry<String, ServerInfo>> existing = 
onlineServers.entrySet().stream()
+                               .filter(e -> 
e.getValue().getRooms().contains(roomId))
+                               .findFirst();
+               if (existing.isPresent()) {
+                       String serverId = existing.get().getKey();
+                       return curServerId.equals(serverId) ? null : 
existing.get().getValue().getUrl();
+               }
+               Optional<Map.Entry<String, ServerInfo>> min = 
onlineServers.entrySet().stream()
+                               .min((e1, e2) -> e1.getValue().getCapacity() - 
e2.getValue().getCapacity());
+               String serverId = min.get().getKey();
+               return curServerId.equals(serverId) ? null : 
min.get().getValue().getUrl();
+       }
+
        public class ClientListener implements
                        EntryAddedListener<String, Client>
                        , EntryUpdatedListener<String, Client>
@@ -360,4 +418,39 @@ public class ClientManager implements IClientManager {
                        onlineRooms.remove(event.getKey(), event.getValue());
                }
        }
+
+       private static class ServerInfo implements Serializable {
+               private static final long serialVersionUID = 1L;
+               private int capacity = 0;
+               private final String url;
+               private final Set<Long> rooms = new HashSet<>();
+
+               public ServerInfo(String url) {
+                       this.url = url;
+               }
+
+               public void add(Room r) {
+                       if (rooms.add(r.getId())) {
+                               capacity += r.getCapacity();
+                       }
+               }
+
+               public void remove(Room r) {
+                       if (rooms.remove(r.getId())) {
+                               capacity -= r.getCapacity();
+                       }
+               }
+
+               public String getUrl() {
+                       return url;
+               }
+
+               public int getCapacity() {
+                       return capacity;
+               }
+
+               public Set<Long> getRooms() {
+                       return rooms;
+               }
+       }
 }
diff --git 
a/openmeetings-web/src/main/java/org/apache/openmeetings/web/util/OmUrlFragment.java
 
b/openmeetings-web/src/main/java/org/apache/openmeetings/web/util/OmUrlFragment.java
index a303c27..e73e3e6 100644
--- 
a/openmeetings-web/src/main/java/org/apache/openmeetings/web/util/OmUrlFragment.java
+++ 
b/openmeetings-web/src/main/java/org/apache/openmeetings/web/util/OmUrlFragment.java
@@ -40,6 +40,7 @@ import org.apache.openmeetings.web.admin.oauth.OAuthPanel;
 import org.apache.openmeetings.web.admin.rooms.RoomsPanel;
 import org.apache.openmeetings.web.admin.users.UsersPanel;
 import org.apache.openmeetings.web.app.Application;
+import org.apache.openmeetings.web.app.ClientManager;
 import org.apache.openmeetings.web.common.BasePanel;
 import org.apache.openmeetings.web.room.RoomPanel;
 import org.apache.openmeetings.web.user.calendar.CalendarPanel;
@@ -47,6 +48,7 @@ import 
org.apache.openmeetings.web.user.dashboard.OmDashboardPanel;
 import org.apache.openmeetings.web.user.profile.SettingsPanel;
 import org.apache.openmeetings.web.user.record.RecordingsPanel;
 import org.apache.openmeetings.web.user.rooms.RoomsSelectorPanel;
+import org.apache.wicket.request.flow.RedirectToUrlException;
 
 public class OmUrlFragment implements Serializable {
        private static final long serialVersionUID = 1L;
@@ -253,6 +255,7 @@ public class OmUrlFragment implements Serializable {
                                        Long roomId = Long.valueOf(type);
                                        Room r = 
Application.get().getBean(RoomDao.class).get(roomId);
                                        if (r != null) {
+                                               moveToServer(roomId);
                                                basePanel = new 
RoomPanel(CHILD_ID, r);
                                        }
                                } catch(NumberFormatException ne) {
@@ -289,4 +292,14 @@ public class OmUrlFragment implements Serializable {
        public String getLink() {
                return getBaseUrl() + "#" + getArea().name() + "/" + getType();
        }
+
+       private static void moveToServer(Long roomId) {
+               if (roomId == null) {
+                       return;
+               }
+               String url = 
Application.get().getBean(ClientManager.class).getServerUrl(roomId);
+               if (url != null) {
+                       throw new RedirectToUrlException(url);
+               }
+       }
 }
diff --git a/openmeetings-web/src/main/webapp/WEB-INF/classes/hazelcast.xml 
b/openmeetings-web/src/main/webapp/WEB-INF/classes/hazelcast.xml
index a37da8b..dc847e1 100644
--- a/openmeetings-web/src/main/webapp/WEB-INF/classes/hazelcast.xml
+++ b/openmeetings-web/src/main/webapp/WEB-INF/classes/hazelcast.xml
@@ -58,7 +58,17 @@
                        <cache-local-entries>true</cache-local-entries>
                </near-cache>
        </map>
+       <map name="SERVERS_KEY">
+               <near-cache>
+                       <eviction eviction-policy="NONE"/>
+                       <in-memory-format>OBJECT</in-memory-format>
+                       <cache-local-entries>true</cache-local-entries>
+               </near-cache>
+       </map>
        <instance-name>server-1</instance-name><!-- MAKE SURE THIS ONE IS 
UNIQUE -->
+       <member-attributes>
+               <attribute 
name="server.url">https://127.0.0.1:5443/openmeetings</attribute><!-- MAKE SURE 
THIS PUBLIC SERVER ADDRESS, USE IP with care: certificate may be not valid for 
IP -->
+       </member-attributes>
        <network>
                <join>
                        <multicast enabled="false"/>
@@ -81,6 +91,5 @@
        </network-->
        <properties>
                <property name="hazelcast.logging.type">slf4j</property>
-               <property 
name="server.url">https://myhost:myport/ctx</property> <!-- MAKE SURE THIS 
PUBLIC SERVER ADDRESS, USE IP with care: certificate may be not valid for IP -->
        </properties>
 </hazelcast>

Reply via email to