From e1d4159453a50cf757e6c13387437278a0899a55 Mon Sep 17 00:00:00 2001 From: rob bygrave Date: Sun, 17 Jun 2018 22:25:45 +1200 Subject: [PATCH] #1431 - Refactor internals for RemoteTransactionEvent binary serialisation --- pom.xml | 2 +- .../ebeaninternal/api/BinaryReadContext.java | 47 ++++++ .../io/ebeaninternal/api/BinaryWritable.java | 23 +++ .../ebeaninternal/api/BinaryWriteContext.java | 49 ++++++ .../api/TransactionEventTable.java | 33 ++-- .../server/cache/RemoteCacheEvent.java | 22 +-- .../cluster/BinaryTransactionEventReader.java | 41 +++++ .../server/cluster/ClusterManager.java | 2 +- ...eServerProvider.java => ServerLookup.java} | 2 +- .../binarymessage/BinaryDataReader.java | 75 --------- .../binarymessage/BinaryDataWriter.java | 49 ------ .../cluster/binarymessage/BinaryMessage.java | 60 ------- .../binarymessage/BinaryMessageList.java | 21 --- .../cluster/binarymessage/ClusterMessage.java | 147 ------------------ .../InvalidMessageException.java | 8 - .../binarymessage/MessageReadWrite.java | 37 ----- .../server/cluster/binarymessage/MsgKeys.java | 19 --- .../server/transaction/BeanPersistIds.java | 94 ++++------- .../server/transaction/RemoteTableMod.java | 20 ++- .../transaction/RemoteTransactionEvent.java | 89 +++++++++-- ... BinaryTransactionEventReadWriteTest.java} | 19 ++- 21 files changed, 306 insertions(+), 553 deletions(-) create mode 100644 src/main/java/io/ebeaninternal/api/BinaryReadContext.java create mode 100644 src/main/java/io/ebeaninternal/api/BinaryWritable.java create mode 100644 src/main/java/io/ebeaninternal/api/BinaryWriteContext.java create mode 100644 src/main/java/io/ebeaninternal/server/cluster/BinaryTransactionEventReader.java rename src/main/java/io/ebeaninternal/server/cluster/{MessageServerProvider.java => ServerLookup.java} (85%) delete mode 100644 src/main/java/io/ebeaninternal/server/cluster/binarymessage/BinaryDataReader.java delete mode 100644 src/main/java/io/ebeaninternal/server/cluster/binarymessage/BinaryDataWriter.java delete mode 100644 src/main/java/io/ebeaninternal/server/cluster/binarymessage/BinaryMessage.java delete mode 100644 src/main/java/io/ebeaninternal/server/cluster/binarymessage/BinaryMessageList.java delete mode 100644 src/main/java/io/ebeaninternal/server/cluster/binarymessage/ClusterMessage.java delete mode 100644 src/main/java/io/ebeaninternal/server/cluster/binarymessage/InvalidMessageException.java delete mode 100644 src/main/java/io/ebeaninternal/server/cluster/binarymessage/MessageReadWrite.java delete mode 100644 src/main/java/io/ebeaninternal/server/cluster/binarymessage/MsgKeys.java rename src/test/java/io/ebeaninternal/server/cluster/binarymessage/{MessageReadWriteTest.java => BinaryTransactionEventReadWriteTest.java} (86%) diff --git a/pom.xml b/pom.xml index 4cfee2676..6791ed678 100644 --- a/pom.xml +++ b/pom.xml @@ -9,7 +9,7 @@ io.ebean ebean - 11.17.6-SNAPSHOT + 11.18.1-SNAPSHOT jar ebean diff --git a/src/main/java/io/ebeaninternal/api/BinaryReadContext.java b/src/main/java/io/ebeaninternal/api/BinaryReadContext.java new file mode 100644 index 000000000..105a9b81b --- /dev/null +++ b/src/main/java/io/ebeaninternal/api/BinaryReadContext.java @@ -0,0 +1,47 @@ +package io.ebeaninternal.api; + +import java.io.ByteArrayInputStream; +import java.io.DataInputStream; +import java.io.IOException; + +/** + * Context used to read binary format messages. + */ +public class BinaryReadContext { + + private final DataInputStream in; + + /** + * Create with protocol 0 and byte data. + */ + public BinaryReadContext(byte[] byteData) { + this(new DataInputStream(new ByteArrayInputStream(byteData))); + } + + /** + * Create with protocol version and DataInputStream data. + */ + public BinaryReadContext(DataInputStream in) { + this.in = in; + } + + public DataInputStream in() { + return in; + } + + public boolean readBoolean() throws IOException { + return in.readBoolean(); + } + + public int readInt() throws IOException { + return in.readInt(); + } + + public String readUTF() throws IOException { + return in.readUTF(); + } + + public long readLong() throws IOException { + return in.readLong(); + } +} diff --git a/src/main/java/io/ebeaninternal/api/BinaryWritable.java b/src/main/java/io/ebeaninternal/api/BinaryWritable.java new file mode 100644 index 000000000..9fbc92ad4 --- /dev/null +++ b/src/main/java/io/ebeaninternal/api/BinaryWritable.java @@ -0,0 +1,23 @@ +package io.ebeaninternal.api; + +import java.io.IOException; + +/** + * Messages that can be sent in binary form. + *

+ * Mainly RemoteTransactionEvent which is sent to cluster members. + *

+ */ +public interface BinaryWritable { + + int TYPE_BEANIUD = 1; + int TYPE_TABLEIUD = 2; + int TYPE_CACHE = 3; + int TYPE_TABLEMOD = 4; + + /** + * Write message in binary format. + */ + void writeBinary(BinaryWriteContext out) throws IOException; + +} diff --git a/src/main/java/io/ebeaninternal/api/BinaryWriteContext.java b/src/main/java/io/ebeaninternal/api/BinaryWriteContext.java new file mode 100644 index 000000000..c36bfebc8 --- /dev/null +++ b/src/main/java/io/ebeaninternal/api/BinaryWriteContext.java @@ -0,0 +1,49 @@ +package io.ebeaninternal.api; + +import java.io.DataOutputStream; +import java.io.IOException; + +/** + * Context used to write binary message (like RemoteTransactionEvent). + */ +public class BinaryWriteContext { + + private final DataOutputStream out; + + private long counter; + + public BinaryWriteContext(DataOutputStream out) { + this.out = out; + } + + /** + * Return the number of message parts that have been written. + */ + public long counter() { + return counter; + } + + /** + * Return the output stream to write to. + */ + public DataOutputStream os() { + return out; + } + + /** + * Start a message part with a given type code. + */ + public DataOutputStream start(int type) throws IOException { + counter++; + out.writeBoolean(true); + out.writeInt(type); + return out; + } + + /** + * End of message parts. + */ + public void end() throws IOException { + out.writeBoolean(false); + } +} diff --git a/src/main/java/io/ebeaninternal/api/TransactionEventTable.java b/src/main/java/io/ebeaninternal/api/TransactionEventTable.java index fe5338cf4..18ea4f5a3 100644 --- a/src/main/java/io/ebeaninternal/api/TransactionEventTable.java +++ b/src/main/java/io/ebeaninternal/api/TransactionEventTable.java @@ -1,10 +1,7 @@ package io.ebeaninternal.api; import io.ebean.event.BulkTableEvent; -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.io.Serializable; @@ -12,7 +9,7 @@ import java.util.Collection; import java.util.HashMap; import java.util.Map; -public final class TransactionEventTable implements Serializable { +public final class TransactionEventTable implements Serializable, BinaryWritable { private static final long serialVersionUID = 2236555729767483264L; @@ -23,20 +20,13 @@ public final class TransactionEventTable implements Serializable { return "TransactionEventTable " + map.values(); } - public void writeBinaryMessage(BinaryMessageList msgList) throws IOException { - + @Override + public void writeBinary(BinaryWriteContext out) throws IOException { for (TableIUD tableIud : map.values()) { - tableIud.writeBinaryMessage(msgList); + tableIud.writeBinary(out); } } - public void readBinaryMessage(DataInput dataInput) throws IOException { - - TableIUD tableIud = TableIUD.readBinaryMessage(dataInput); - map.put(tableIud.getTableName(), tableIud); - } - - public void add(TransactionEventTable table) { for (TableIUD iud : table.values()) { @@ -47,7 +37,6 @@ public final class TransactionEventTable implements Serializable { public void add(String table, boolean insert, boolean update, boolean delete) { table = table.toUpperCase(); - add(new TableIUD(table, insert, update, delete)); } @@ -67,7 +56,7 @@ public final class TransactionEventTable implements Serializable { return map.values(); } - public static class TableIUD implements Serializable, BulkTableEvent { + public static class TableIUD implements Serializable, BulkTableEvent, BinaryWritable { private static final long serialVersionUID = -1958317571064162089L; @@ -83,7 +72,7 @@ public final class TransactionEventTable implements Serializable { this.delete = delete; } - public static TableIUD readBinaryMessage(DataInput dataInput) throws IOException { + public static TableIUD readBinaryMessage(BinaryReadContext dataInput) throws IOException { String table = dataInput.readUTF(); boolean insert = dataInput.readBoolean(); @@ -93,17 +82,13 @@ public final class TransactionEventTable implements Serializable { return new TableIUD(table, insert, update, delete); } - public void writeBinaryMessage(BinaryMessageList msgList) throws IOException { - - BinaryMessage msg = new BinaryMessage(table.length() + 10); - DataOutputStream os = msg.getOs(); - os.writeInt(BinaryMessage.TYPE_TABLEIUD); + @Override + public void writeBinary(BinaryWriteContext out) throws IOException { + DataOutputStream os = out.start(TYPE_TABLEIUD); os.writeUTF(table); os.writeBoolean(insert); os.writeBoolean(update); os.writeBoolean(delete); - os.close(); - msgList.add(msg); } @Override diff --git a/src/main/java/io/ebeaninternal/server/cache/RemoteCacheEvent.java b/src/main/java/io/ebeaninternal/server/cache/RemoteCacheEvent.java index 2d04df0f5..8024087de 100644 --- a/src/main/java/io/ebeaninternal/server/cache/RemoteCacheEvent.java +++ b/src/main/java/io/ebeaninternal/server/cache/RemoteCacheEvent.java @@ -1,9 +1,9 @@ package io.ebeaninternal.server.cache; -import io.ebeaninternal.server.cluster.binarymessage.BinaryMessage; -import io.ebeaninternal.server.cluster.binarymessage.BinaryMessageList; +import io.ebeaninternal.api.BinaryReadContext; +import io.ebeaninternal.api.BinaryWritable; +import io.ebeaninternal.api.BinaryWriteContext; -import java.io.DataInput; import java.io.DataOutputStream; import java.io.IOException; import java.util.ArrayList; @@ -12,7 +12,7 @@ import java.util.List; /** * Cache events broadcast across the cluster. */ -public class RemoteCacheEvent { +public class RemoteCacheEvent implements BinaryWritable { private boolean clearAll; @@ -56,7 +56,7 @@ public class RemoteCacheEvent { return "clearAll:" + clearAll + " caches:" + clearCaches; } - public static RemoteCacheEvent readBinaryMessage(DataInput dataInput) throws IOException { + public static RemoteCacheEvent readBinaryMessage(BinaryReadContext dataInput) throws IOException { boolean clearAll = dataInput.readBoolean(); int size = dataInput.readInt(); @@ -72,13 +72,9 @@ public class RemoteCacheEvent { return new RemoteCacheEvent(clearAll, clearCache); } - public void writeBinaryMessage(BinaryMessageList msgList) throws IOException { - - int bufferSize = (clearCaches == null) ? 0 : clearCaches.size() * 30; - - BinaryMessage msg = new BinaryMessage(bufferSize + 10); - DataOutputStream os = msg.getOs(); - os.writeInt(BinaryMessage.TYPE_CACHE); + @Override + public void writeBinary(BinaryWriteContext out) throws IOException { + DataOutputStream os = out.start(TYPE_CACHE); os.writeBoolean(clearAll); if (clearCaches == null) { os.writeInt(0); @@ -88,7 +84,5 @@ public class RemoteCacheEvent { os.writeUTF(cacheName); } } - msgList.add(msg); - } } diff --git a/src/main/java/io/ebeaninternal/server/cluster/BinaryTransactionEventReader.java b/src/main/java/io/ebeaninternal/server/cluster/BinaryTransactionEventReader.java new file mode 100644 index 000000000..aebf45ea4 --- /dev/null +++ b/src/main/java/io/ebeaninternal/server/cluster/BinaryTransactionEventReader.java @@ -0,0 +1,41 @@ +package io.ebeaninternal.server.cluster; + +import io.ebeaninternal.api.BinaryReadContext; +import io.ebeaninternal.api.SpiEbeanServer; +import io.ebeaninternal.server.transaction.RemoteTransactionEvent; + +import java.io.IOException; + +/** + * Mechanism to convert RemoteTransactionEvent to/from byte[] content. + */ +public class BinaryTransactionEventReader { + + private final ServerLookup serverLookup; + + public BinaryTransactionEventReader(ServerLookup serverLookup) { + this.serverLookup = serverLookup; + } + + /** + * Read Transaction from bytes. + */ + public RemoteTransactionEvent read(byte[] byteData) throws IOException { + return read(new BinaryReadContext(byteData)); + } + + /** + * Read Transaction using BinaryReadContext. + */ + public RemoteTransactionEvent read(BinaryReadContext dataInput) throws IOException { + + String serverName = dataInput.readUTF(); + SpiEbeanServer server = (SpiEbeanServer) serverLookup.getServer(serverName); + if (server == null) { + throw new IllegalStateException("EbeanServer not found for name [" + serverName + "]"); + } + RemoteTransactionEvent event = new RemoteTransactionEvent(server); + event.readBinary(dataInput); + return event; + } +} diff --git a/src/main/java/io/ebeaninternal/server/cluster/ClusterManager.java b/src/main/java/io/ebeaninternal/server/cluster/ClusterManager.java index 6f02f5b50..bffc93a10 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 implements MessageServerProvider { +public class ClusterManager implements ServerLookup { 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/ServerLookup.java similarity index 85% rename from src/main/java/io/ebeaninternal/server/cluster/MessageServerProvider.java rename to src/main/java/io/ebeaninternal/server/cluster/ServerLookup.java index f3c388d2b..abf835371 100644 --- a/src/main/java/io/ebeaninternal/server/cluster/MessageServerProvider.java +++ b/src/main/java/io/ebeaninternal/server/cluster/ServerLookup.java @@ -5,7 +5,7 @@ import io.ebean.EbeanServer; /** * Returns EbeanServer instances for remote message reading. */ -public interface MessageServerProvider { +public interface ServerLookup { /** * Return the EbeanServer instance by 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 deleted file mode 100644 index f097172f4..000000000 --- a/src/main/java/io/ebeaninternal/server/cluster/binarymessage/BinaryDataReader.java +++ /dev/null @@ -1,75 +0,0 @@ -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 deleted file mode 100644 index 8f08bba56..000000000 --- a/src/main/java/io/ebeaninternal/server/cluster/binarymessage/BinaryDataWriter.java +++ /dev/null @@ -1,49 +0,0 @@ -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/BinaryMessage.java b/src/main/java/io/ebeaninternal/server/cluster/binarymessage/BinaryMessage.java deleted file mode 100644 index 158df57da..000000000 --- a/src/main/java/io/ebeaninternal/server/cluster/binarymessage/BinaryMessage.java +++ /dev/null @@ -1,60 +0,0 @@ -package io.ebeaninternal.server.cluster.binarymessage; - -import java.io.ByteArrayOutputStream; -import java.io.DataOutputStream; - -/** - * Represents a relatively small independent message. - *

- * In general terms we break up a potentially large object like - * RemoteTransactionEvent into many smaller BinaryMessages. This is so that if - * they don't all fit on a single Packet we can easily break them up and put - * them on multiple packets. - *

- *

- * Also note that for the Multicast approach a Packet will generally contain - * many messages each directed to different members of the cluster. So it would - * be common for many Ack, Resend and Control messages to all be contained in a - * single packet. - *

- */ -public class BinaryMessage { - - public static final int TYPE_MSGCONTROL = 0; - 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; - - private final ByteArrayOutputStream buffer; - private final DataOutputStream os; - private byte[] bytes; - - /** - * Create with an estimated buffer size. - */ - public BinaryMessage(int bufSize) { - this.buffer = new ByteArrayOutputStream(bufSize); - this.os = new DataOutputStream(buffer); - } - - /** - * Return the DataOutputStream to write content to. - */ - public DataOutputStream getOs() { - return os; - } - - /** - * Return all the content as a byte array. - */ - public byte[] getByteArray() { - if (bytes == null) { - bytes = buffer.toByteArray(); - } - return bytes; - } -} diff --git a/src/main/java/io/ebeaninternal/server/cluster/binarymessage/BinaryMessageList.java b/src/main/java/io/ebeaninternal/server/cluster/binarymessage/BinaryMessageList.java deleted file mode 100644 index fd057fd03..000000000 --- a/src/main/java/io/ebeaninternal/server/cluster/binarymessage/BinaryMessageList.java +++ /dev/null @@ -1,21 +0,0 @@ -package io.ebeaninternal.server.cluster.binarymessage; - -import java.util.ArrayList; -import java.util.List; - -/** - * Holds a List of BinaryMessage's. - */ -public class BinaryMessageList { - - final List list = new ArrayList<>(); - - public void add(BinaryMessage msg) { - list.add(msg); - } - - public List getList() { - return 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 deleted file mode 100644 index 2e6b189ea..000000000 --- a/src/main/java/io/ebeaninternal/server/cluster/binarymessage/ClusterMessage.java +++ /dev/null @@ -1,147 +0,0 @@ -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 deleted file mode 100644 index d05d0b3c2..000000000 --- a/src/main/java/io/ebeaninternal/server/cluster/binarymessage/InvalidMessageException.java +++ /dev/null @@ -1,8 +0,0 @@ -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 deleted file mode 100644 index fe7cab1c2..000000000 --- a/src/main/java/io/ebeaninternal/server/cluster/binarymessage/MessageReadWrite.java +++ /dev/null @@ -1,37 +0,0 @@ -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 deleted file mode 100644 index 9daeaa9bd..000000000 --- a/src/main/java/io/ebeaninternal/server/cluster/binarymessage/MsgKeys.java +++ /dev/null @@ -1,19 +0,0 @@ -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 934877759..25ebea220 100644 --- a/src/main/java/io/ebeaninternal/server/transaction/BeanPersistIds.java +++ b/src/main/java/io/ebeaninternal/server/transaction/BeanPersistIds.java @@ -1,8 +1,9 @@ package io.ebeaninternal.server.transaction; +import io.ebeaninternal.api.BinaryReadContext; +import io.ebeaninternal.api.BinaryWritable; +import io.ebeaninternal.api.BinaryWriteContext; import io.ebeaninternal.api.SpiEbeanServer; -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; @@ -27,7 +28,7 @@ import java.util.List; * size of data sent around the network. *

*/ -public class BeanPersistIds { +public class BeanPersistIds implements BinaryWritable { private final BeanDescriptor beanDescriptor; @@ -45,21 +46,20 @@ public class BeanPersistIds { this.descriptorId = desc.getDescriptorId(); } - public static BeanPersistIds readBinaryMessage(SpiEbeanServer server, DataInput dataInput) throws IOException { + public static BeanPersistIds readBinaryMessage(SpiEbeanServer server, BinaryReadContext input) throws IOException { - String descriptorId = dataInput.readUTF(); - BeanDescriptor desc = server.getBeanDescriptorById(descriptorId); + BeanDescriptor desc = server.getBeanDescriptorById(input.readUTF()); BeanPersistIds bp = new BeanPersistIds(desc); - bp.read(dataInput); + bp.read(input); return bp; } - private void read(DataInput dataInput) throws IOException { + private void read(BinaryReadContext dataInput) throws IOException { IdBinder idBinder = beanDescriptor.getIdBinder(); int iudType = dataInput.readInt(); - List idList = readIdList(dataInput, idBinder); + List idList = readIdList(dataInput.in(), idBinder); switch (iudType) { case 0: insertIds = idList; @@ -76,20 +76,28 @@ public class BeanPersistIds { } } - /** - * Write the contents into a BinaryMessage form. - *

- * For a RemoteBeanPersist with a large number of id's note that this is - * broken up into many BinaryMessages each with a maximum of 100 ids. This - * enables the contents of a large RemoteTransactionEvent to be split up - * across multiple Packets. - *

- */ - void writeBinaryMessage(BinaryMessageList msgList) throws IOException { + @Override + public void writeBinary(BinaryWriteContext out) throws IOException { + writeBinaryList(out, 0, insertIds); + writeBinaryList(out, 1, updateIds); + writeBinaryList(out, 2, deleteIds); + } - writeIdList(beanDescriptor, 0, insertIds, msgList); - writeIdList(beanDescriptor, 1, updateIds, msgList); - writeIdList(beanDescriptor, 2, deleteIds, msgList); + private void writeBinaryList(BinaryWriteContext out, int type, List ids) throws IOException { + + if (ids != null) { + IdBinder idBinder = beanDescriptor.getIdBinder(); + int count = ids.size(); + + DataOutputStream os = out.start(TYPE_BEANIUD); + os.writeUTF(descriptorId); + os.writeInt(type); + os.writeInt(count); + for (Object insertId : ids) { + idBinder.writeData(os, insertId); + } + os.flush(); + } } private List readIdList(DataInput dataInput, IdBinder idBinder) throws IOException { @@ -105,48 +113,6 @@ public class BeanPersistIds { return idList; } - /** - * Write a BinaryMessage containing the descriptorId, iudType and list of Id - * values. - *

- * Note that a given BinaryMessage has a maximum of 100 Ids. This is due to - * the limit of UDP packet sizes. We break up the RemoteBeanPersist into - * potentially many smaller BinaryMessages which may be put into multiple - * Packets. - *

- */ - private void writeIdList(BeanDescriptor desc, int iudType, List idList, BinaryMessageList msgList) throws IOException { - - IdBinder idBinder = desc.getIdBinder(); - - int count = idList == null ? 0 : idList.size(); - if (count > 0) { - int loop = 0; - int i = 0; - int eof = idList.size(); - do { - ++loop; - int endOfLoop = Math.min(eof, loop * 100); - - BinaryMessage m = new BinaryMessage(endOfLoop * 4 + 20); - - DataOutputStream os = m.getOs(); - os.writeInt(BinaryMessage.TYPE_BEANIUD); - os.writeUTF(descriptorId); - os.writeInt(iudType); - os.writeInt(count); - - for (; i < endOfLoop; i++) { - idBinder.writeData(os, idList.get(i)); - } - - os.close(); - msgList.add(m); - - } while (i < eof); - } - } - @Override public String toString() { StringBuilder sb = new StringBuilder(); diff --git a/src/main/java/io/ebeaninternal/server/transaction/RemoteTableMod.java b/src/main/java/io/ebeaninternal/server/transaction/RemoteTableMod.java index 1c6eda418..0234a6aba 100644 --- a/src/main/java/io/ebeaninternal/server/transaction/RemoteTableMod.java +++ b/src/main/java/io/ebeaninternal/server/transaction/RemoteTableMod.java @@ -1,15 +1,15 @@ package io.ebeaninternal.server.transaction; -import io.ebeaninternal.server.cluster.binarymessage.BinaryMessage; -import io.ebeaninternal.server.cluster.binarymessage.BinaryMessageList; +import io.ebeaninternal.api.BinaryReadContext; +import io.ebeaninternal.api.BinaryWritable; +import io.ebeaninternal.api.BinaryWriteContext; -import java.io.DataInput; import java.io.DataOutputStream; import java.io.IOException; import java.util.LinkedHashSet; import java.util.Set; -public class RemoteTableMod { +public class RemoteTableMod implements BinaryWritable { private final long timestamp; @@ -28,7 +28,7 @@ public class RemoteTableMod { return tables; } - public static RemoteTableMod readBinaryMessage(DataInput dataInput) throws IOException { + public static RemoteTableMod readBinaryMessage(BinaryReadContext dataInput) throws IOException { long timestamp = dataInput.readLong(); int count = dataInput.readInt(); @@ -40,16 +40,14 @@ public class RemoteTableMod { 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); + @Override + public void writeBinary(BinaryWriteContext out) throws IOException { + DataOutputStream os = out.start(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 54112ca66..17004bf34 100644 --- a/src/main/java/io/ebeaninternal/server/transaction/RemoteTransactionEvent.java +++ b/src/main/java/io/ebeaninternal/server/transaction/RemoteTransactionEvent.java @@ -1,15 +1,20 @@ package io.ebeaninternal.server.transaction; +import io.ebeaninternal.api.BinaryReadContext; +import io.ebeaninternal.api.BinaryWritable; +import io.ebeaninternal.api.BinaryWriteContext; import io.ebeaninternal.api.SpiEbeanServer; +import io.ebeaninternal.api.TransactionEventTable; import io.ebeaninternal.api.TransactionEventTable.TableIUD; import io.ebeaninternal.server.cache.RemoteCacheEvent; -import io.ebeaninternal.server.cluster.binarymessage.BinaryMessageList; +import java.io.ByteArrayOutputStream; +import java.io.DataOutputStream; import java.io.IOException; import java.util.ArrayList; import java.util.List; -public class RemoteTransactionEvent implements Runnable { +public class RemoteTransactionEvent implements Runnable, BinaryWritable { private final List beanPersistList = new ArrayList<>(); @@ -25,10 +30,16 @@ public class RemoteTransactionEvent implements Runnable { private transient SpiEbeanServer server; + /** + * Create for sending to other servers in the cluster. + */ public RemoteTransactionEvent(String serverName) { this.serverName = serverName; } + /** + * Create from Reading and processing from remote server. + */ public RemoteTransactionEvent(SpiEbeanServer server) { this.server = server; } @@ -53,30 +64,86 @@ public class RemoteTransactionEvent implements Runnable { return sb.toString(); } - public void writeBinaryMessage(BinaryMessageList msgList) throws IOException { + /** + * Read the binary message. + */ + public void readBinary(BinaryReadContext dataInput) throws IOException { + + boolean more = dataInput.readBoolean(); + while (more) { + int msgType = dataInput.readInt(); + readBinaryMessage(msgType, dataInput); + more = dataInput.readBoolean(); + } + } + + private void readBinaryMessage(int msgType, BinaryReadContext dataInput) throws IOException { + + switch (msgType) { + case BinaryWritable.TYPE_BEANIUD: + addBeanPersistIds(BeanPersistIds.readBinaryMessage(server, dataInput)); + break; + + case BinaryWritable.TYPE_TABLEIUD: + addTableIUD(TransactionEventTable.TableIUD.readBinaryMessage(dataInput)); + break; + + case BinaryWritable.TYPE_CACHE: + addRemoteCacheEvent(RemoteCacheEvent.readBinaryMessage(dataInput)); + break; + + case BinaryWritable.TYPE_TABLEMOD: + addRemoteTableMod(RemoteTableMod.readBinaryMessage(dataInput)); + break; + + default: + throw new RuntimeException("Invalid Transaction msgType " + msgType); + } + } + + /** + * Write a binary message to byte array given an initial buffer size. + */ + public byte[] writeBinaryAsBytes(int bufferSize) throws IOException { + + ByteArrayOutputStream buffer = new ByteArrayOutputStream(bufferSize); + DataOutputStream out = new DataOutputStream(buffer); + BinaryWriteContext context = new BinaryWriteContext(out); + + writeBinary(context); + out.close(); + + return buffer.toByteArray(); + } + + @Override + public void writeBinary(BinaryWriteContext out) throws IOException { + + DataOutputStream os = out.os(); + //os.writeInt(TRANSACTION_EVENT); + os.writeUTF(serverName); if (remoteTableMod != null) { - remoteTableMod.writeBinary(msgList); + remoteTableMod.writeBinary(out); } - if (tableList != null) { for (TableIUD aTableList : tableList) { - aTableList.writeBinaryMessage(msgList); + aTableList.writeBinary(out); } } - if (deleteByIdMap != null) { for (BeanPersistIds deleteIds : deleteByIdMap.values()) { - deleteIds.writeBinaryMessage(msgList); + deleteIds.writeBinary(out); } } - for (BeanPersistIds aBeanPersistList : beanPersistList) { - aBeanPersistList.writeBinaryMessage(msgList); + aBeanPersistList.writeBinary(out); } if (remoteCacheEvent != null) { - remoteCacheEvent.writeBinaryMessage(msgList); + remoteCacheEvent.writeBinary(out); } + out.end(); + os.flush(); } public boolean isEmpty() { diff --git a/src/test/java/io/ebeaninternal/server/cluster/binarymessage/MessageReadWriteTest.java b/src/test/java/io/ebeaninternal/server/cluster/binarymessage/BinaryTransactionEventReadWriteTest.java similarity index 86% rename from src/test/java/io/ebeaninternal/server/cluster/binarymessage/MessageReadWriteTest.java rename to src/test/java/io/ebeaninternal/server/cluster/binarymessage/BinaryTransactionEventReadWriteTest.java index cf5eeee87..d696aa096 100644 --- a/src/test/java/io/ebeaninternal/server/cluster/binarymessage/MessageReadWriteTest.java +++ b/src/test/java/io/ebeaninternal/server/cluster/binarymessage/BinaryTransactionEventReadWriteTest.java @@ -5,7 +5,8 @@ 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.cluster.BinaryTransactionEventReader; +import io.ebeaninternal.server.cluster.ServerLookup; import io.ebeaninternal.server.core.PersistRequest; import io.ebeaninternal.server.deploy.BeanDescriptor; import io.ebeaninternal.server.transaction.BeanPersistIds; @@ -22,13 +23,13 @@ import java.util.Set; import static org.assertj.core.api.Assertions.assertThat; -public class MessageReadWriteTest extends BaseTestCase { +public class BinaryTransactionEventReadWriteTest extends BaseTestCase { private BeanDescriptor customerBeanDescriptor = getBeanDescriptor(Customer.class); private TDEbeanServer mockEbeanServer = new TDEbeanServer(); - private MessageReadWrite readWrite = new MessageReadWrite(new TDMessageServerProvider()); + private BinaryTransactionEventReader reader = new BinaryTransactionEventReader(new TDServerLookup()); @Override protected BeanDescriptor getBeanDescriptor(Class cls) { @@ -59,9 +60,8 @@ public class MessageReadWriteTest extends BaseTestCase { event.addBeanPersistIds(beanPersistIds); - byte[] binaryMessage = readWrite.write(event); - - RemoteTransactionEvent read = readWrite.read(binaryMessage); + byte[] binaryMessage = event.writeBinaryAsBytes(256); + RemoteTransactionEvent read = reader.read(binaryMessage); // cache event RemoteCacheEvent remoteCacheEvent = read.getRemoteCacheEvent(); @@ -91,9 +91,8 @@ public class MessageReadWriteTest extends BaseTestCase { event.addRemoteTableMod(new RemoteTableMod(timestamp, tables)); - byte[] binaryMessage = readWrite.write(event); - - RemoteTransactionEvent read = readWrite.read(binaryMessage); + byte[] binaryMessage = event.writeBinaryAsBytes(256); + RemoteTransactionEvent read = reader.read(binaryMessage); RemoteTableMod remoteTableMod = read.getRemoteTableMod(); @@ -101,7 +100,7 @@ public class MessageReadWriteTest extends BaseTestCase { assertThat(remoteTableMod.getTables()).isEqualTo(tables); } - class TDMessageServerProvider implements MessageServerProvider { + class TDServerLookup implements ServerLookup { @Override public EbeanServer getServer(String name) {