From 9c73309dfd0bfe935202183e228a9b0282262fb9 Mon Sep 17 00:00:00 2001 From: rob bygrave Date: Sat, 16 Jun 2018 15:15:54 +1200 Subject: [PATCH] #1427 - QueryCache should be cleared, if one of a dependent bean is updated Update RemoteTransactionEvent, add (back) common binary message handling (Common to TCP, K8s and upcoming binary pubsub) --- .../api/TransactionEventTable.java | 6 +- .../server/cache/RemoteCacheEvent.java | 4 +- .../server/cluster/ClusterManager.java | 2 +- .../server/cluster/MessageServerProvider.java | 14 ++ .../binarymessage/BinaryDataReader.java | 75 +++++++++ .../binarymessage/BinaryDataWriter.java | 49 ++++++ .../{ => binarymessage}/BinaryMessage.java | 3 +- .../BinaryMessageList.java | 2 +- .../cluster/binarymessage/ClusterMessage.java | 147 ++++++++++++++++++ .../InvalidMessageException.java | 8 + .../binarymessage/MessageReadWrite.java | 37 +++++ .../server/cluster/binarymessage/MsgKeys.java | 19 +++ .../server/transaction/BeanPersistIds.java | 18 ++- .../transaction/PostCommitProcessing.java | 4 + .../server/transaction/RemoteTableMod.java | 55 +++++++ .../transaction/RemoteTransactionEvent.java | 16 +- .../binarymessage/MessageReadWriteTest.java | 118 ++++++++++++++ 17 files changed, 563 insertions(+), 14 deletions(-) create mode 100644 src/main/java/io/ebeaninternal/server/cluster/MessageServerProvider.java create mode 100644 src/main/java/io/ebeaninternal/server/cluster/binarymessage/BinaryDataReader.java create mode 100644 src/main/java/io/ebeaninternal/server/cluster/binarymessage/BinaryDataWriter.java rename src/main/java/io/ebeaninternal/server/cluster/{ => binarymessage}/BinaryMessage.java (94%) rename src/main/java/io/ebeaninternal/server/cluster/{ => binarymessage}/BinaryMessageList.java (85%) create mode 100644 src/main/java/io/ebeaninternal/server/cluster/binarymessage/ClusterMessage.java create mode 100644 src/main/java/io/ebeaninternal/server/cluster/binarymessage/InvalidMessageException.java create mode 100644 src/main/java/io/ebeaninternal/server/cluster/binarymessage/MessageReadWrite.java create mode 100644 src/main/java/io/ebeaninternal/server/cluster/binarymessage/MsgKeys.java create mode 100644 src/main/java/io/ebeaninternal/server/transaction/RemoteTableMod.java create mode 100644 src/test/java/io/ebeaninternal/server/cluster/binarymessage/MessageReadWriteTest.java diff --git a/src/main/java/io/ebeaninternal/api/TransactionEventTable.java b/src/main/java/io/ebeaninternal/api/TransactionEventTable.java index cde507c1a..fe5338cf4 100644 --- a/src/main/java/io/ebeaninternal/api/TransactionEventTable.java +++ b/src/main/java/io/ebeaninternal/api/TransactionEventTable.java @@ -1,8 +1,8 @@ package io.ebeaninternal.api; import io.ebean.event.BulkTableEvent; -import io.ebeaninternal.server.cluster.BinaryMessage; -import io.ebeaninternal.server.cluster.BinaryMessageList; +import io.ebeaninternal.server.cluster.binarymessage.BinaryMessage; +import io.ebeaninternal.server.cluster.binarymessage.BinaryMessageList; import java.io.DataInput; import java.io.DataOutputStream; @@ -102,7 +102,7 @@ public final class TransactionEventTable implements Serializable { os.writeBoolean(insert); os.writeBoolean(update); os.writeBoolean(delete); - + os.close(); msgList.add(msg); } diff --git a/src/main/java/io/ebeaninternal/server/cache/RemoteCacheEvent.java b/src/main/java/io/ebeaninternal/server/cache/RemoteCacheEvent.java index e0def309f..2d04df0f5 100644 --- a/src/main/java/io/ebeaninternal/server/cache/RemoteCacheEvent.java +++ b/src/main/java/io/ebeaninternal/server/cache/RemoteCacheEvent.java @@ -1,7 +1,7 @@ package io.ebeaninternal.server.cache; -import io.ebeaninternal.server.cluster.BinaryMessage; -import io.ebeaninternal.server.cluster.BinaryMessageList; +import io.ebeaninternal.server.cluster.binarymessage.BinaryMessage; +import io.ebeaninternal.server.cluster.binarymessage.BinaryMessageList; import java.io.DataInput; import java.io.DataOutputStream; diff --git a/src/main/java/io/ebeaninternal/server/cluster/ClusterManager.java b/src/main/java/io/ebeaninternal/server/cluster/ClusterManager.java index 6fd56c713..6f02f5b50 100644 --- a/src/main/java/io/ebeaninternal/server/cluster/ClusterManager.java +++ b/src/main/java/io/ebeaninternal/server/cluster/ClusterManager.java @@ -13,7 +13,7 @@ import java.util.concurrent.ConcurrentHashMap; /** * Manages the cluster service. */ -public class ClusterManager { +public class ClusterManager implements MessageServerProvider { private static final Logger clusterLogger = LoggerFactory.getLogger("io.ebean.Cluster"); diff --git a/src/main/java/io/ebeaninternal/server/cluster/MessageServerProvider.java b/src/main/java/io/ebeaninternal/server/cluster/MessageServerProvider.java new file mode 100644 index 000000000..f3c388d2b --- /dev/null +++ b/src/main/java/io/ebeaninternal/server/cluster/MessageServerProvider.java @@ -0,0 +1,14 @@ +package io.ebeaninternal.server.cluster; + +import io.ebean.EbeanServer; + +/** + * Returns EbeanServer instances for remote message reading. + */ +public interface MessageServerProvider { + + /** + * Return the EbeanServer instance by name. + */ + EbeanServer getServer(String name); +} diff --git a/src/main/java/io/ebeaninternal/server/cluster/binarymessage/BinaryDataReader.java b/src/main/java/io/ebeaninternal/server/cluster/binarymessage/BinaryDataReader.java new file mode 100644 index 000000000..f097172f4 --- /dev/null +++ b/src/main/java/io/ebeaninternal/server/cluster/binarymessage/BinaryDataReader.java @@ -0,0 +1,75 @@ +package io.ebeaninternal.server.cluster.binarymessage; + +import io.ebeaninternal.api.SpiEbeanServer; +import io.ebeaninternal.api.TransactionEventTable; +import io.ebeaninternal.server.cache.RemoteCacheEvent; +import io.ebeaninternal.server.cluster.MessageServerProvider; +import io.ebeaninternal.server.transaction.BeanPersistIds; +import io.ebeaninternal.server.transaction.RemoteTableMod; +import io.ebeaninternal.server.transaction.RemoteTransactionEvent; + +import java.io.ByteArrayInputStream; +import java.io.DataInputStream; +import java.io.IOException; + +/** + * Reads the binary message returning RemoteTransactionEvent. + */ +class BinaryDataReader { + + private final MessageServerProvider serverProvider; + private final DataInputStream dataInput; + + private SpiEbeanServer server; + private RemoteTransactionEvent event; + + BinaryDataReader(MessageServerProvider serverProvider, byte[] data) { + this.serverProvider = serverProvider; + this.dataInput = new DataInputStream(new ByteArrayInputStream(data)); + } + + /** + * Read the binary message returning a RemoteTransactionEvent. + */ + RemoteTransactionEvent read() throws IOException { + + String serverName = dataInput.readUTF(); + + this.server = (SpiEbeanServer) serverProvider.getServer(serverName); + if (server == null) { + throw new IllegalStateException("EbeanServer not found for name [" + serverName + "]"); + } + this.event = new RemoteTransactionEvent(server); + boolean more = dataInput.readBoolean(); + while (more) { + readMessage(); + more = dataInput.readBoolean(); + } + return event; + } + + private void readMessage() throws IOException { + + int msgType = dataInput.readInt(); + switch (msgType) { + case BinaryMessage.TYPE_BEANIUD: + event.addBeanPersistIds(BeanPersistIds.readBinaryMessage(server, dataInput)); + break; + + case BinaryMessage.TYPE_TABLEIUD: + event.addTableIUD(TransactionEventTable.TableIUD.readBinaryMessage(dataInput)); + break; + + case BinaryMessage.TYPE_CACHE: + event.addRemoteCacheEvent(RemoteCacheEvent.readBinaryMessage(dataInput)); + break; + + case BinaryMessage.TYPE_TABLEMOD: + event.addRemoteTableMod(RemoteTableMod.readBinaryMessage(dataInput)); + break; + + default: + throw new RuntimeException("Invalid Transaction msgType " + msgType); + } + } +} diff --git a/src/main/java/io/ebeaninternal/server/cluster/binarymessage/BinaryDataWriter.java b/src/main/java/io/ebeaninternal/server/cluster/binarymessage/BinaryDataWriter.java new file mode 100644 index 000000000..8f08bba56 --- /dev/null +++ b/src/main/java/io/ebeaninternal/server/cluster/binarymessage/BinaryDataWriter.java @@ -0,0 +1,49 @@ +package io.ebeaninternal.server.cluster.binarymessage; + +import java.io.ByteArrayOutputStream; +import java.io.DataOutputStream; +import java.io.IOException; + +/** + * Writes BinaryMessageList to DataHolder raw byte[]. + */ +class BinaryDataWriter { + + private final BinaryMessageList messageList; + + private final ByteArrayOutputStream buffer; + + private final DataOutputStream dataOut; + + private final String serverName; + + BinaryDataWriter(String serverName, BinaryMessageList messageList) { + this.serverName = serverName; + this.messageList = messageList; + this.buffer = new ByteArrayOutputStream(256); + this.dataOut = new DataOutputStream(buffer); + } + + /** + * Write the message as raw byte[]. + */ + byte[] write() throws IOException { + + //TODO: Add Protocol version and "podId" to support pubsub cluster + //dataOut.writeInt(MsgKeys.MESSAGE_KEY); + dataOut.writeUTF(serverName); + for (BinaryMessage msg : messageList.getList()) { + write(msg); + } + dataOut.writeBoolean(false); + dataOut.flush(); + return buffer.toByteArray(); + } + + private void write(BinaryMessage msg) throws IOException { + + dataOut.writeBoolean(true); + dataOut.write(msg.getByteArray()); + } + +} diff --git a/src/main/java/io/ebeaninternal/server/cluster/BinaryMessage.java b/src/main/java/io/ebeaninternal/server/cluster/binarymessage/BinaryMessage.java similarity index 94% rename from src/main/java/io/ebeaninternal/server/cluster/BinaryMessage.java rename to src/main/java/io/ebeaninternal/server/cluster/binarymessage/BinaryMessage.java index b78b3f7eb..158df57da 100644 --- a/src/main/java/io/ebeaninternal/server/cluster/BinaryMessage.java +++ b/src/main/java/io/ebeaninternal/server/cluster/binarymessage/BinaryMessage.java @@ -1,4 +1,4 @@ -package io.ebeaninternal.server.cluster; +package io.ebeaninternal.server.cluster.binarymessage; import java.io.ByteArrayOutputStream; import java.io.DataOutputStream; @@ -24,6 +24,7 @@ public class BinaryMessage { public static final int TYPE_BEANIUD = 1; public static final int TYPE_TABLEIUD = 2; public static final int TYPE_CACHE = 3; + public static final int TYPE_TABLEMOD = 4; public static final int TYPE_MSGACK = 8; public static final int TYPE_MSGRESEND = 9; diff --git a/src/main/java/io/ebeaninternal/server/cluster/BinaryMessageList.java b/src/main/java/io/ebeaninternal/server/cluster/binarymessage/BinaryMessageList.java similarity index 85% rename from src/main/java/io/ebeaninternal/server/cluster/BinaryMessageList.java rename to src/main/java/io/ebeaninternal/server/cluster/binarymessage/BinaryMessageList.java index f84b7995c..fd057fd03 100644 --- a/src/main/java/io/ebeaninternal/server/cluster/BinaryMessageList.java +++ b/src/main/java/io/ebeaninternal/server/cluster/binarymessage/BinaryMessageList.java @@ -1,4 +1,4 @@ -package io.ebeaninternal.server.cluster; +package io.ebeaninternal.server.cluster.binarymessage; import java.util.ArrayList; import java.util.List; diff --git a/src/main/java/io/ebeaninternal/server/cluster/binarymessage/ClusterMessage.java b/src/main/java/io/ebeaninternal/server/cluster/binarymessage/ClusterMessage.java new file mode 100644 index 000000000..2e6b189ea --- /dev/null +++ b/src/main/java/io/ebeaninternal/server/cluster/binarymessage/ClusterMessage.java @@ -0,0 +1,147 @@ +package io.ebeaninternal.server.cluster.binarymessage; + +import java.io.DataInputStream; +import java.io.DataOutputStream; +import java.io.IOException; + +/** + * The message broadcast around the cluster. + */ +public class ClusterMessage { + + private static final int MAX_LENGTH = 10 * 1024 * 1024; + + private final String registerIp; + + private final String podName; + + private final boolean register; + + private final byte[] data; + + /** + * Create a register message. + */ + public static ClusterMessage register(String registerIp, boolean register, String podName) { + return new ClusterMessage(registerIp, register, podName); + } + + /** + * Create a transaction message. + */ + public static ClusterMessage transEvent(byte[] data) { + return new ClusterMessage(data); + } + + /** + * Create for register online/offline message. + */ + private ClusterMessage(String registerIp, boolean register, String podName) { + this.registerIp = registerIp; + this.register = register; + this.podName = podName; + this.data = null; + } + + /** + * Create for a transaction message. + */ + private ClusterMessage(byte[] data) { + this.data = data; + this.registerIp = null; + this.podName = null; + this.register = false; + } + + public String toString() { + StringBuilder sb = new StringBuilder(); + if (registerIp != null) { + sb.append("register "); + sb.append(register); + sb.append(" "); + sb.append(registerIp); + } else { + sb.append("[data]"); + } + return sb.toString(); + } + + /** + * Return true if this is a register event as opposed to a transaction message. + */ + public boolean isRegisterEvent() { + return registerIp != null; + } + + /** + * Return the register host for online/offline message. + */ + public String getRegisterIp() { + return registerIp; + } + + public String getPodName() { + return podName; + } + + /** + * Return true if register is true for a online/offline message. + */ + public boolean isRegister() { + return register; + } + + /** + * Return the raw message data. + */ + public byte[] getData() { + return data; + } + + /** + * Write the message in binary form. + */ + public void write(DataOutputStream dataOutput) throws IOException { + + if (data != null) { + // write data message + dataOutput.writeInt(MsgKeys.DATA); + dataOutput.writeInt(data.length); + dataOutput.write(data); + } else { + // write header message + dataOutput.writeInt(MsgKeys.HEADER); + dataOutput.writeUTF(getRegisterIp()); + dataOutput.writeBoolean(register); + dataOutput.writeUTF(getPodName()); + } + dataOutput.flush(); + } + + /** + * Read the message from binary form. + */ + public static ClusterMessage read(DataInputStream dataInput) throws IOException, InvalidMessageException { + + int key = dataInput.readInt(); + if (key == MsgKeys.DATA) { + int length = dataInput.readInt(); + if (length > MAX_LENGTH) { + throw new IOException("Message data too large length:"+length); + } + byte[] data = new byte[length]; + dataInput.readFully(data); + return new ClusterMessage(data); + + } else if (key == MsgKeys.HEADER) { + String host = dataInput.readUTF(); + boolean registered = dataInput.readBoolean(); + String podName = dataInput.readUTF(); + return new ClusterMessage(host, registered, podName); + + } else { + throw new InvalidMessageException("Invalid message key:" + key); + } + } + +} diff --git a/src/main/java/io/ebeaninternal/server/cluster/binarymessage/InvalidMessageException.java b/src/main/java/io/ebeaninternal/server/cluster/binarymessage/InvalidMessageException.java new file mode 100644 index 000000000..d05d0b3c2 --- /dev/null +++ b/src/main/java/io/ebeaninternal/server/cluster/binarymessage/InvalidMessageException.java @@ -0,0 +1,8 @@ +package io.ebeaninternal.server.cluster.binarymessage; + +public class InvalidMessageException extends Exception { + + public InvalidMessageException(String message) { + super(message); + } +} diff --git a/src/main/java/io/ebeaninternal/server/cluster/binarymessage/MessageReadWrite.java b/src/main/java/io/ebeaninternal/server/cluster/binarymessage/MessageReadWrite.java new file mode 100644 index 000000000..fe7cab1c2 --- /dev/null +++ b/src/main/java/io/ebeaninternal/server/cluster/binarymessage/MessageReadWrite.java @@ -0,0 +1,37 @@ +package io.ebeaninternal.server.cluster.binarymessage; + +import io.ebeaninternal.server.cluster.MessageServerProvider; +import io.ebeaninternal.server.transaction.RemoteTransactionEvent; + +import java.io.IOException; + +/** + * Mechanism to convert RemoteTransactionEvent to/from byte[] content. + */ +public class MessageReadWrite { + + private final MessageServerProvider serverProvider; + + public MessageReadWrite(MessageServerProvider serverProvider) { + this.serverProvider = serverProvider; + } + + /** + * Convert the RemoteTransactionEvent to raw byte[] content. + */ + public byte[] write(RemoteTransactionEvent transEvent) throws IOException { + + BinaryMessageList messageList = new BinaryMessageList(); + transEvent.writeBinaryMessage(messageList); + + return new BinaryDataWriter(transEvent.getServerName(), messageList).write(); + } + + /** + * Convert the byte[] content to RemoteTransactionEvent. + */ + public RemoteTransactionEvent read(byte[] data) throws IOException { + + return new BinaryDataReader(serverProvider, data).read(); + } +} diff --git a/src/main/java/io/ebeaninternal/server/cluster/binarymessage/MsgKeys.java b/src/main/java/io/ebeaninternal/server/cluster/binarymessage/MsgKeys.java new file mode 100644 index 000000000..9daeaa9bd --- /dev/null +++ b/src/main/java/io/ebeaninternal/server/cluster/binarymessage/MsgKeys.java @@ -0,0 +1,19 @@ +package io.ebeaninternal.server.cluster.binarymessage; + +public interface MsgKeys { + + /** + * Used to identify client on connection initiation. + */ + int HELLO = 182; + + /** + * Used to confirm protocol on reading messages. + */ + int HEADER = 11; + + /** + * Used to confirm protocol on reading messages. + */ + int DATA = 12; +} diff --git a/src/main/java/io/ebeaninternal/server/transaction/BeanPersistIds.java b/src/main/java/io/ebeaninternal/server/transaction/BeanPersistIds.java index cf77829ee..934877759 100644 --- a/src/main/java/io/ebeaninternal/server/transaction/BeanPersistIds.java +++ b/src/main/java/io/ebeaninternal/server/transaction/BeanPersistIds.java @@ -1,8 +1,8 @@ package io.ebeaninternal.server.transaction; import io.ebeaninternal.api.SpiEbeanServer; -import io.ebeaninternal.server.cluster.BinaryMessage; -import io.ebeaninternal.server.cluster.BinaryMessageList; +import io.ebeaninternal.server.cluster.binarymessage.BinaryMessage; +import io.ebeaninternal.server.cluster.binarymessage.BinaryMessageList; import io.ebeaninternal.server.core.PersistRequest; import io.ebeaninternal.server.deploy.BeanDescriptor; import io.ebeaninternal.server.deploy.id.IdBinder; @@ -140,7 +140,7 @@ public class BeanPersistIds { idBinder.writeData(os, idList.get(i)); } - os.flush(); + os.close(); msgList.add(m); } while (i < eof); @@ -167,7 +167,7 @@ public class BeanPersistIds { return sb.toString(); } - void addId(PersistRequest.Type type, Serializable id) { + public void addId(PersistRequest.Type type, Serializable id) { switch (type) { case INSERT: addInsertId(id); @@ -210,10 +210,18 @@ public class BeanPersistIds { return beanDescriptor; } - List getDeleteIds() { + public List getDeleteIds() { return deleteIds; } + public List getInsertIds() { + return insertIds; + } + + public List getUpdateIds() { + return updateIds; + } + /** * Notify the cache of this event that came from another server in the cluster. */ diff --git a/src/main/java/io/ebeaninternal/server/transaction/PostCommitProcessing.java b/src/main/java/io/ebeaninternal/server/transaction/PostCommitProcessing.java index 08a85fd7d..c84593646 100644 --- a/src/main/java/io/ebeaninternal/server/transaction/PostCommitProcessing.java +++ b/src/main/java/io/ebeaninternal/server/transaction/PostCommitProcessing.java @@ -211,6 +211,10 @@ final class PostCommitProcessing { RemoteTransactionEvent remoteTransactionEvent = new RemoteTransactionEvent(serverName); + Set touched = cacheChanges.touchedTables(); + if (touched != null && !touched.isEmpty()) { + remoteTransactionEvent.addRemoteTableMod(new RemoteTableMod(cacheChanges.modificationTimestamp(), touched)); + } if (beanPersistIdMap != null) { for (BeanPersistIds beanPersist : beanPersistIdMap.values()) { remoteTransactionEvent.addBeanPersistIds(beanPersist); diff --git a/src/main/java/io/ebeaninternal/server/transaction/RemoteTableMod.java b/src/main/java/io/ebeaninternal/server/transaction/RemoteTableMod.java new file mode 100644 index 000000000..1c6eda418 --- /dev/null +++ b/src/main/java/io/ebeaninternal/server/transaction/RemoteTableMod.java @@ -0,0 +1,55 @@ +package io.ebeaninternal.server.transaction; + +import io.ebeaninternal.server.cluster.binarymessage.BinaryMessage; +import io.ebeaninternal.server.cluster.binarymessage.BinaryMessageList; + +import java.io.DataInput; +import java.io.DataOutputStream; +import java.io.IOException; +import java.util.LinkedHashSet; +import java.util.Set; + +public class RemoteTableMod { + + private final long timestamp; + + private final Set tables; + + public RemoteTableMod(long timestamp, Set tables) { + this.timestamp = timestamp; + this.tables = tables; + } + + public long getTimestamp() { + return timestamp; + } + + public Set getTables() { + return tables; + } + + public static RemoteTableMod readBinaryMessage(DataInput dataInput) throws IOException { + + long timestamp = dataInput.readLong(); + int count = dataInput.readInt(); + + Set tables = new LinkedHashSet<>(); + for (int i = 0; i < count; i++) { + tables.add(dataInput.readUTF()); + } + return new RemoteTableMod(timestamp, tables); + } + + public void writeBinary(BinaryMessageList msgList) throws IOException { + BinaryMessage msg = new BinaryMessage(tables.size() * 20 + 30); + DataOutputStream os = msg.getOs(); + os.writeInt(BinaryMessage.TYPE_TABLEMOD); + os.writeLong(timestamp); + os.writeInt(tables.size()); + for (String table : tables) { + os.writeUTF(table); + } + os.close(); + msgList.add(msg); + } +} diff --git a/src/main/java/io/ebeaninternal/server/transaction/RemoteTransactionEvent.java b/src/main/java/io/ebeaninternal/server/transaction/RemoteTransactionEvent.java index f6047d209..54112ca66 100644 --- a/src/main/java/io/ebeaninternal/server/transaction/RemoteTransactionEvent.java +++ b/src/main/java/io/ebeaninternal/server/transaction/RemoteTransactionEvent.java @@ -3,7 +3,7 @@ package io.ebeaninternal.server.transaction; import io.ebeaninternal.api.SpiEbeanServer; import io.ebeaninternal.api.TransactionEventTable.TableIUD; import io.ebeaninternal.server.cache.RemoteCacheEvent; -import io.ebeaninternal.server.cluster.BinaryMessageList; +import io.ebeaninternal.server.cluster.binarymessage.BinaryMessageList; import java.io.IOException; import java.util.ArrayList; @@ -19,6 +19,8 @@ public class RemoteTransactionEvent implements Runnable { private RemoteCacheEvent remoteCacheEvent; + private RemoteTableMod remoteTableMod; + private String serverName; private transient SpiEbeanServer server; @@ -53,6 +55,10 @@ public class RemoteTransactionEvent implements Runnable { public void writeBinaryMessage(BinaryMessageList msgList) throws IOException { + if (remoteTableMod != null) { + remoteTableMod.writeBinary(msgList); + } + if (tableList != null) { for (TableIUD aTableList : tableList) { aTableList.writeBinaryMessage(msgList); @@ -113,6 +119,10 @@ public class RemoteTransactionEvent implements Runnable { tableList.add(tableIud); } + public void addRemoteTableMod(RemoteTableMod remoteTableMod) { + this.remoteTableMod = remoteTableMod; + } + public String getServerName() { return serverName; } @@ -141,4 +151,8 @@ public class RemoteTransactionEvent implements Runnable { return remoteCacheEvent; } + public RemoteTableMod getRemoteTableMod() { + return remoteTableMod; + } + } diff --git a/src/test/java/io/ebeaninternal/server/cluster/binarymessage/MessageReadWriteTest.java b/src/test/java/io/ebeaninternal/server/cluster/binarymessage/MessageReadWriteTest.java new file mode 100644 index 000000000..cf5eeee87 --- /dev/null +++ b/src/test/java/io/ebeaninternal/server/cluster/binarymessage/MessageReadWriteTest.java @@ -0,0 +1,118 @@ +package io.ebeaninternal.server.cluster.binarymessage; + +import io.ebean.BaseTestCase; +import io.ebean.EbeanServer; +import io.ebeaninternal.api.TDSpiEbeanServer; +import io.ebeaninternal.api.TransactionEventTable; +import io.ebeaninternal.server.cache.RemoteCacheEvent; +import io.ebeaninternal.server.cluster.MessageServerProvider; +import io.ebeaninternal.server.core.PersistRequest; +import io.ebeaninternal.server.deploy.BeanDescriptor; +import io.ebeaninternal.server.transaction.BeanPersistIds; +import io.ebeaninternal.server.transaction.RemoteTableMod; +import io.ebeaninternal.server.transaction.RemoteTransactionEvent; +import org.junit.Test; +import org.tests.model.basic.Customer; + +import java.io.IOException; +import java.util.Collections; +import java.util.HashSet; +import java.util.List; +import java.util.Set; + +import static org.assertj.core.api.Assertions.assertThat; + +public class MessageReadWriteTest extends BaseTestCase { + + private BeanDescriptor customerBeanDescriptor = getBeanDescriptor(Customer.class); + + private TDEbeanServer mockEbeanServer = new TDEbeanServer(); + + private MessageReadWrite readWrite = new MessageReadWrite(new TDMessageServerProvider()); + + @Override + protected BeanDescriptor getBeanDescriptor(Class cls) { + return super.getBeanDescriptor(cls); + } + + @Test + public void readWrite() throws IOException { + + RemoteTransactionEvent event = new RemoteTransactionEvent("db"); + + long timestamp = System.currentTimeMillis(); + Set tables = new HashSet<>(); + Collections.addAll(tables, "one", "two", "three"); + + event.addRemoteTableMod(new RemoteTableMod(timestamp, tables)); + event.addRemoteCacheEvent(new RemoteCacheEvent(Customer.class)); + + event.addTableIUD(new TransactionEventTable.TableIUD("foo", true, false, true)); + event.addTableIUD(new TransactionEventTable.TableIUD("bar", false, true, false)); + + + BeanPersistIds beanPersistIds = new BeanPersistIds(customerBeanDescriptor); + beanPersistIds.addId(PersistRequest.Type.INSERT, 42); + beanPersistIds.addId(PersistRequest.Type.INSERT, 43); + beanPersistIds.addId(PersistRequest.Type.UPDATE, 55); + beanPersistIds.addId(PersistRequest.Type.DELETE, 66); + + event.addBeanPersistIds(beanPersistIds); + + byte[] binaryMessage = readWrite.write(event); + + RemoteTransactionEvent read = readWrite.read(binaryMessage); + + // cache event + RemoteCacheEvent remoteCacheEvent = read.getRemoteCacheEvent(); + assertThat(remoteCacheEvent.getClearCaches()).containsOnly(Customer.class.getName()); + + // table mod + RemoteTableMod remoteTableMod = read.getRemoteTableMod(); + assertThat(remoteTableMod.getTimestamp()).isEqualTo(timestamp); + assertThat(remoteTableMod.getTables()).isEqualTo(tables); + + // Bean persist ids + List beanPersistList = read.getBeanPersistList(); + assertThat(beanPersistList).hasSize(3); + assertThat(beanPersistList.get(0).getInsertIds()).containsOnly(42, 43); + assertThat(beanPersistList.get(1).getUpdateIds()).containsOnly(55); + assertThat(beanPersistList.get(2).getDeleteIds()).containsOnly(66); + } + + @Test + public void readWrite_RemoteTableMod() throws IOException { + + RemoteTransactionEvent event = new RemoteTransactionEvent("db"); + + long timestamp = System.currentTimeMillis(); + Set tables = new HashSet<>(); + Collections.addAll(tables, "one", "two", "three"); + + event.addRemoteTableMod(new RemoteTableMod(timestamp, tables)); + + byte[] binaryMessage = readWrite.write(event); + + RemoteTransactionEvent read = readWrite.read(binaryMessage); + + RemoteTableMod remoteTableMod = read.getRemoteTableMod(); + + assertThat(remoteTableMod.getTimestamp()).isEqualTo(timestamp); + assertThat(remoteTableMod.getTables()).isEqualTo(tables); + } + + class TDMessageServerProvider implements MessageServerProvider { + + @Override + public EbeanServer getServer(String name) { + return mockEbeanServer; + } + } + + class TDEbeanServer extends TDSpiEbeanServer { + @Override + public BeanDescriptor getBeanDescriptorById(String descriptorId) { + return customerBeanDescriptor; + } + } +}