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>