diff --git a/src/main/java/com/avaje/ebeaninternal/server/cluster/BinaryMessage.java b/src/main/java/com/avaje/ebeaninternal/server/cluster/BinaryMessage.java index 2883bd3ca..c556c3263 100644 --- a/src/main/java/com/avaje/ebeaninternal/server/cluster/BinaryMessage.java +++ b/src/main/java/com/avaje/ebeaninternal/server/cluster/BinaryMessage.java @@ -20,39 +20,39 @@ import java.io.DataOutputStream; */ 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_MSGCONTROL = 0; + public static final int TYPE_BEANIUD = 1; + public static final int TYPE_TABLEIUD = 2; - public static final int TYPE_MSGACK = 8; - public static final int TYPE_MSGRESEND = 9; + 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; + 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); - } + /** + * 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 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; + /** + * 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/com/avaje/ebeaninternal/server/cluster/BinaryMessageList.java b/src/main/java/com/avaje/ebeaninternal/server/cluster/BinaryMessageList.java index 01b757ec2..17d7e6b33 100644 --- a/src/main/java/com/avaje/ebeaninternal/server/cluster/BinaryMessageList.java +++ b/src/main/java/com/avaje/ebeaninternal/server/cluster/BinaryMessageList.java @@ -5,19 +5,19 @@ import java.util.List; /** * Holds a List of BinaryMessage's. - * + * * @author rbygrave */ public class BinaryMessageList { - final ArrayList list = new ArrayList(); - - public void add(BinaryMessage msg) { - list.add(msg); - } + final ArrayList list = new ArrayList(); + + public void add(BinaryMessage msg) { + list.add(msg); + } + + public List getList() { + return list; + } - public List getList() { - return list; - } - } diff --git a/src/main/java/com/avaje/ebeaninternal/server/cluster/ClusterBroadcast.java b/src/main/java/com/avaje/ebeaninternal/server/cluster/ClusterBroadcast.java index 3260c7ad5..414c0bfaf 100644 --- a/src/main/java/com/avaje/ebeaninternal/server/cluster/ClusterBroadcast.java +++ b/src/main/java/com/avaje/ebeaninternal/server/cluster/ClusterBroadcast.java @@ -8,21 +8,21 @@ import com.avaje.ebeaninternal.server.transaction.RemoteTransactionEvent; */ public interface ClusterBroadcast { - /** - * Inform the other cluster members that this instance has come online and - * start any listeners etc. - */ - void startup(ClusterManager clusterManager); + /** + * Inform the other cluster members that this instance has come online and + * start any listeners etc. + */ + void startup(ClusterManager clusterManager); - /** - * Inform the other cluster members that this instance is leaving and - * shutdown any listeners. - */ - void shutdown(); - - /** - * Send a transaction event to all the members of the cluster. - */ - void broadcast(RemoteTransactionEvent remoteTransEvent); + /** + * Inform the other cluster members that this instance is leaving and + * shutdown any listeners. + */ + void shutdown(); + + /** + * Send a transaction event to all the members of the cluster. + */ + void broadcast(RemoteTransactionEvent remoteTransEvent); } diff --git a/src/main/java/com/avaje/ebeaninternal/server/cluster/ClusterManager.java b/src/main/java/com/avaje/ebeaninternal/server/cluster/ClusterManager.java index 5e8ae4dd9..67318d2d3 100644 --- a/src/main/java/com/avaje/ebeaninternal/server/cluster/ClusterManager.java +++ b/src/main/java/com/avaje/ebeaninternal/server/cluster/ClusterManager.java @@ -1,7 +1,5 @@ package com.avaje.ebeaninternal.server.cluster; -import java.util.concurrent.ConcurrentHashMap; - import com.avaje.ebean.EbeanServer; import com.avaje.ebean.config.ContainerConfig; import com.avaje.ebeaninternal.server.cluster.mcast.McastClusterManager; @@ -10,6 +8,8 @@ import com.avaje.ebeaninternal.server.transaction.RemoteTransactionEvent; import org.slf4j.Logger; import org.slf4j.LoggerFactory; +import java.util.concurrent.ConcurrentHashMap; + /** * Manages the cluster service. */ diff --git a/src/main/java/com/avaje/ebeaninternal/server/cluster/DataHolder.java b/src/main/java/com/avaje/ebeaninternal/server/cluster/DataHolder.java index c38cf7a9a..46f4e148c 100644 --- a/src/main/java/com/avaje/ebeaninternal/server/cluster/DataHolder.java +++ b/src/main/java/com/avaje/ebeaninternal/server/cluster/DataHolder.java @@ -9,16 +9,16 @@ import java.io.Serializable; */ public class DataHolder implements Serializable { - private static final long serialVersionUID = 9090748723571322192L; + private static final long serialVersionUID = 9090748723571322192L; - private final byte[] data; - - public DataHolder(byte[] data) { - this.data = data; - } + private final byte[] data; + + public DataHolder(byte[] data) { + this.data = data; + } + + public byte[] getData() { + return data; + } - public byte[] getData() { - return data; - } - } diff --git a/src/main/java/com/avaje/ebeaninternal/server/cluster/Packet.java b/src/main/java/com/avaje/ebeaninternal/server/cluster/Packet.java index e07e92e57..1966a34f6 100644 --- a/src/main/java/com/avaje/ebeaninternal/server/cluster/Packet.java +++ b/src/main/java/com/avaje/ebeaninternal/server/cluster/Packet.java @@ -15,179 +15,179 @@ import java.io.IOException; * Due to the hard limit on the size of UDP packets a RemoteTransactionEvent * with lots of information could be broken up into multiple packets. *

- * + * * @author rbygrave */ public class Packet { - - /** - * A Packet that holds protocol messages like ACK, PING etc. - */ - public static final short TYPE_MESSAGES = 1; - - /** - * A Packet that holds TransactionEvent information such as Bean - * and or Table IUD information. - */ - public static final short TYPE_TRANSEVENT = 2; - /** - * The type of Packet. - */ - protected final short packetType; - - /** - * The PacketId. - */ - protected final long packetId; - - /** - * The timestamp the Packet was created. - */ - protected final long timestamp; - - /** - * The EbeanServer name this relates to if relevant. - */ - protected final String serverName; - - protected ByteArrayOutputStream buffer; - protected DataOutputStream dataOut; - protected byte[] bytes; - - /** - * The number of messages in this Packet. - */ - private int messageCount; - - /** - * The number of times this Packet was resent. - */ - private int resendCount; + /** + * A Packet that holds protocol messages like ACK, PING etc. + */ + public static final short TYPE_MESSAGES = 1; - /** - * Create a Packet for writing messages to. - */ - public static Packet forWrite(short packetType, long packetId, long timestamp, String serverName) throws IOException { - return new Packet(true, packetType, packetId, timestamp, serverName); - } - - /** - * Create a Packet just reading the Header information. - */ - public static Packet readHeader(DataInput dataInput) throws IOException { - - short packetType = dataInput.readShort(); - long packetId = dataInput.readLong(); - long timestamp = dataInput.readLong(); - String serverName = dataInput.readUTF(); - - return new Packet(false, packetType, packetId, timestamp, serverName); - } - - protected Packet(boolean write, short packetType, long packetId, long timestamp, String serverName) throws IOException{ - this.packetType = packetType; - this.packetId = packetId; - this.timestamp = timestamp; - this.serverName = serverName; - if (write){ - this.buffer = new ByteArrayOutputStream(); - this.dataOut = new DataOutputStream(buffer); - writeHeader(); - } else { - this.buffer = null; - this.dataOut = null; - } - } + /** + * A Packet that holds TransactionEvent information such as Bean + * and or Table IUD information. + */ + public static final short TYPE_TRANSEVENT = 2; - private void writeHeader() throws IOException { - dataOut.writeShort(packetType); - dataOut.writeLong(packetId); - dataOut.writeLong(timestamp); - dataOut.writeUTF(serverName); - } + /** + * The type of Packet. + */ + protected final short packetType; - public int incrementResendCount() { - return resendCount++; - } - - public short getPacketType() { - return packetType; - } + /** + * The PacketId. + */ + protected final long packetId; - public long getPacketId() { - return packetId; - } - - public long getTimestamp() { - return timestamp; - } + /** + * The timestamp the Packet was created. + */ + protected final long timestamp; - public String getServerName() { - return serverName; - } - - public void writeEof() throws IOException { - dataOut.writeBoolean(false); - } - - public void read(DataInput dataInput) throws IOException { - boolean more = dataInput.readBoolean(); - while (more){ - int msgType = dataInput.readInt(); - readMessage(dataInput, msgType); - // see if there is more information - more = dataInput.readBoolean(); - } - } - - /** - * Overridden by more specific Packet implementations to read the messages. - */ - protected void readMessage(DataInput dataInput, int msgType) throws IOException { - - } + /** + * The EbeanServer name this relates to if relevant. + */ + protected final String serverName; - /** - * Write a binary message to this packet returning true if there was - * enough room to do so. Return false if the message was too large for - * the remaining space left - in this case another Packet should be - * created to put that message into. - */ - public boolean writeBinaryMessage(BinaryMessage msg, int maxPacketSize) throws IOException { - - byte[] bytes = msg.getByteArray(); - - if (messageCount > 0 && (bytes.length + buffer.size() > maxPacketSize)){ - // we are actually going to ignore the maxPacketSize iff we have one - // large message. - - // false = no more messages - dataOut.writeBoolean(false); - return false; - } - ++messageCount; - // true = another message follows - dataOut.writeBoolean(true); - dataOut.write(bytes); - return true; + protected ByteArrayOutputStream buffer; + protected DataOutputStream dataOut; + protected byte[] bytes; + + /** + * The number of messages in this Packet. + */ + private int messageCount; + + /** + * The number of times this Packet was resent. + */ + private int resendCount; + + /** + * Create a Packet for writing messages to. + */ + public static Packet forWrite(short packetType, long packetId, long timestamp, String serverName) throws IOException { + return new Packet(true, packetType, packetId, timestamp, serverName); + } + + /** + * Create a Packet just reading the Header information. + */ + public static Packet readHeader(DataInput dataInput) throws IOException { + + short packetType = dataInput.readShort(); + long packetId = dataInput.readLong(); + long timestamp = dataInput.readLong(); + String serverName = dataInput.readUTF(); + + return new Packet(false, packetType, packetId, timestamp, serverName); + } + + protected Packet(boolean write, short packetType, long packetId, long timestamp, String serverName) throws IOException { + this.packetType = packetType; + this.packetId = packetId; + this.timestamp = timestamp; + this.serverName = serverName; + if (write) { + this.buffer = new ByteArrayOutputStream(); + this.dataOut = new DataOutputStream(buffer); + writeHeader(); + } else { + this.buffer = null; + this.dataOut = null; } + } + + private void writeHeader() throws IOException { + dataOut.writeShort(packetType); + dataOut.writeLong(packetId); + dataOut.writeLong(timestamp); + dataOut.writeUTF(serverName); + } + + public int incrementResendCount() { + return resendCount++; + } + + public short getPacketType() { + return packetType; + } + + public long getPacketId() { + return packetId; + } + + public long getTimestamp() { + return timestamp; + } + + public String getServerName() { + return serverName; + } + + public void writeEof() throws IOException { + dataOut.writeBoolean(false); + } + + public void read(DataInput dataInput) throws IOException { + boolean more = dataInput.readBoolean(); + while (more) { + int msgType = dataInput.readInt(); + readMessage(dataInput, msgType); + // see if there is more information + more = dataInput.readBoolean(); + } + } + + /** + * Overridden by more specific Packet implementations to read the messages. + */ + protected void readMessage(DataInput dataInput, int msgType) throws IOException { + + } + + /** + * Write a binary message to this packet returning true if there was + * enough room to do so. Return false if the message was too large for + * the remaining space left - in this case another Packet should be + * created to put that message into. + */ + public boolean writeBinaryMessage(BinaryMessage msg, int maxPacketSize) throws IOException { + + byte[] bytes = msg.getByteArray(); + + if (messageCount > 0 && (bytes.length + buffer.size() > maxPacketSize)) { + // we are actually going to ignore the maxPacketSize iff we have one + // large message. + + // false = no more messages + dataOut.writeBoolean(false); + return false; + } + ++messageCount; + // true = another message follows + dataOut.writeBoolean(true); + dataOut.write(bytes); + return true; + } + + public int getSize() { + return getBytes().length; + } + + /** + * Return the Packet as raw bytes. + */ + public byte[] getBytes() { + if (bytes == null) { + bytes = buffer.toByteArray(); + buffer = null; + dataOut = null; + } + return bytes; + } - public int getSize() { - return getBytes().length; - } - - /** - * Return the Packet as raw bytes. - */ - public byte[] getBytes() { - if (bytes == null){ - bytes = buffer.toByteArray(); - buffer = null; - dataOut = null; - } - return bytes; - } - } diff --git a/src/main/java/com/avaje/ebeaninternal/server/cluster/PacketMessages.java b/src/main/java/com/avaje/ebeaninternal/server/cluster/PacketMessages.java index ab8fd5fa7..cc43039f7 100644 --- a/src/main/java/com/avaje/ebeaninternal/server/cluster/PacketMessages.java +++ b/src/main/java/com/avaje/ebeaninternal/server/cluster/PacketMessages.java @@ -1,69 +1,69 @@ package com.avaje.ebeaninternal.server.cluster; -import java.io.DataInput; -import java.io.IOException; -import java.util.ArrayList; -import java.util.List; - import com.avaje.ebeaninternal.server.cluster.mcast.Message; import com.avaje.ebeaninternal.server.cluster.mcast.MessageAck; import com.avaje.ebeaninternal.server.cluster.mcast.MessageControl; import com.avaje.ebeaninternal.server.cluster.mcast.MessageResend; +import java.io.DataInput; +import java.io.IOException; +import java.util.ArrayList; +import java.util.List; + /** * A Packet that contains Ack, Resend and Control messages. - * + * * @author rbygrave */ public class PacketMessages extends Packet { - private final ArrayList messages; + private final ArrayList messages; - public static PacketMessages forWrite(long packetId, long timestamp, String serverName) throws IOException { - return new PacketMessages(true, packetId, timestamp, serverName); - } - - public static PacketMessages forRead(Packet header) throws IOException { - return new PacketMessages(header); - } - - private PacketMessages(boolean write, long packetId, long timestamp, String serverName) throws IOException { - super(write, TYPE_MESSAGES, packetId, timestamp, serverName); - this.messages = null; - } + public static PacketMessages forWrite(long packetId, long timestamp, String serverName) throws IOException { + return new PacketMessages(true, packetId, timestamp, serverName); + } - private PacketMessages(Packet header) throws IOException { - super(false, TYPE_MESSAGES, header.packetId, header.timestamp, header.serverName); - this.messages = new ArrayList(); - } - - /** - * Return the messages contained in this Packet. - */ - public List getMessages() { - return messages; - } - - /** - * Read the messages (Ack, Resend or Control) contained in this packet. - */ - protected void readMessage(DataInput dataInput, int msgType) throws IOException { - - switch (msgType) { - case BinaryMessage.TYPE_MSGCONTROL: - messages.add(MessageControl.readBinaryMessage(dataInput)); - break; - - case BinaryMessage.TYPE_MSGACK: - messages.add(MessageAck.readBinaryMessage(dataInput)); - break; - - case BinaryMessage.TYPE_MSGRESEND: - messages.add(MessageResend.readBinaryMessage(dataInput)); - break; - - default: - throw new RuntimeException("Invalid Transaction msgType "+msgType); - } + public static PacketMessages forRead(Packet header) throws IOException { + return new PacketMessages(header); + } + + private PacketMessages(boolean write, long packetId, long timestamp, String serverName) throws IOException { + super(write, TYPE_MESSAGES, packetId, timestamp, serverName); + this.messages = null; + } + + private PacketMessages(Packet header) throws IOException { + super(false, TYPE_MESSAGES, header.packetId, header.timestamp, header.serverName); + this.messages = new ArrayList(); + } + + /** + * Return the messages contained in this Packet. + */ + public List getMessages() { + return messages; + } + + /** + * Read the messages (Ack, Resend or Control) contained in this packet. + */ + protected void readMessage(DataInput dataInput, int msgType) throws IOException { + + switch (msgType) { + case BinaryMessage.TYPE_MSGCONTROL: + messages.add(MessageControl.readBinaryMessage(dataInput)); + break; + + case BinaryMessage.TYPE_MSGACK: + messages.add(MessageAck.readBinaryMessage(dataInput)); + break; + + case BinaryMessage.TYPE_MSGRESEND: + messages.add(MessageResend.readBinaryMessage(dataInput)); + break; + + default: + throw new RuntimeException("Invalid Transaction msgType " + msgType); } + } } diff --git a/src/main/java/com/avaje/ebeaninternal/server/cluster/PacketTransactionEvent.java b/src/main/java/com/avaje/ebeaninternal/server/cluster/PacketTransactionEvent.java index d6a408351..a4dac5d2b 100644 --- a/src/main/java/com/avaje/ebeaninternal/server/cluster/PacketTransactionEvent.java +++ b/src/main/java/com/avaje/ebeaninternal/server/cluster/PacketTransactionEvent.java @@ -17,48 +17,48 @@ import java.io.IOException; */ public class PacketTransactionEvent extends Packet { - private final SpiEbeanServer server; - - private final RemoteTransactionEvent event; + private final SpiEbeanServer server; - public static PacketTransactionEvent forWrite(long packetId, long timestamp, String serverName) throws IOException { - return new PacketTransactionEvent(true, packetId, timestamp, serverName); - } - - private PacketTransactionEvent(boolean write, long packetId, long timestamp, String serverName) throws IOException { - super(write, TYPE_TRANSEVENT, packetId, timestamp, serverName); - this.server = null; - this.event = null; - } + private final RemoteTransactionEvent event; - private PacketTransactionEvent(Packet header, SpiEbeanServer server) throws IOException { - super(false, TYPE_TRANSEVENT, header.packetId, header.timestamp, header.serverName); - this.server = server; - this.event = new RemoteTransactionEvent(server); - } + public static PacketTransactionEvent forWrite(long packetId, long timestamp, String serverName) throws IOException { + return new PacketTransactionEvent(true, packetId, timestamp, serverName); + } - public static PacketTransactionEvent forRead(Packet header, SpiEbeanServer server) throws IOException { - return new PacketTransactionEvent(header, server); - } - - public RemoteTransactionEvent getEvent() { - return event; - } + private PacketTransactionEvent(boolean write, long packetId, long timestamp, String serverName) throws IOException { + super(write, TYPE_TRANSEVENT, packetId, timestamp, serverName); + this.server = null; + this.event = null; + } - protected void readMessage(DataInput dataInput, int msgType) throws IOException { - - switch (msgType) { - case BinaryMessage.TYPE_BEANIUD: - event.addBeanPersistIds(BeanPersistIds.readBinaryMessage(server, dataInput)); - break; - - case BinaryMessage.TYPE_TABLEIUD: - event.addTableIUD(TableIUD.readBinaryMessage(dataInput)); - break; + private PacketTransactionEvent(Packet header, SpiEbeanServer server) throws IOException { + super(false, TYPE_TRANSEVENT, header.packetId, header.timestamp, header.serverName); + this.server = server; + this.event = new RemoteTransactionEvent(server); + } - default: - throw new RuntimeException("Invalid Transaction msgType "+msgType); - } + public static PacketTransactionEvent forRead(Packet header, SpiEbeanServer server) throws IOException { + return new PacketTransactionEvent(header, server); + } + + public RemoteTransactionEvent getEvent() { + return event; + } + + protected void readMessage(DataInput dataInput, int msgType) throws IOException { + + switch (msgType) { + case BinaryMessage.TYPE_BEANIUD: + event.addBeanPersistIds(BeanPersistIds.readBinaryMessage(server, dataInput)); + break; + + case BinaryMessage.TYPE_TABLEIUD: + event.addTableIUD(TableIUD.readBinaryMessage(dataInput)); + break; + + default: + throw new RuntimeException("Invalid Transaction msgType " + msgType); } - + } + } diff --git a/src/main/java/com/avaje/ebeaninternal/server/cluster/PacketWriter.java b/src/main/java/com/avaje/ebeaninternal/server/cluster/PacketWriter.java index 3d92314fa..bc3ee0775 100644 --- a/src/main/java/com/avaje/ebeaninternal/server/cluster/PacketWriter.java +++ b/src/main/java/com/avaje/ebeaninternal/server/cluster/PacketWriter.java @@ -1,160 +1,160 @@ package com.avaje.ebeaninternal.server.cluster; +import com.avaje.ebeaninternal.server.cluster.mcast.Message; +import com.avaje.ebeaninternal.server.transaction.RemoteTransactionEvent; + import java.io.IOException; import java.util.ArrayList; import java.util.List; -import com.avaje.ebeaninternal.server.cluster.mcast.Message; -import com.avaje.ebeaninternal.server.transaction.RemoteTransactionEvent; - /** * Creates Packets for either RemoteTransactionEvents or Messages (Ping, ACK, * Join, Leave etc). - * + * * @author rbygrave */ public class PacketWriter { - private final PacketIdGenerator idGenerator; - private final PacketBuilder messagesPacketBuilder; - private final PacketBuilder transEventPacketBuilder; + private final PacketIdGenerator idGenerator; + private final PacketBuilder messagesPacketBuilder; + private final PacketBuilder transEventPacketBuilder; - /** - * Create a PacketWriter with an expected max packet size. - *

- * In theory we would prefer to create packets up to the MTU size which for - * Ethernet will likely be 1500. Note that the maxPacketSize is ignored for - * large single messages. - *

- */ - public PacketWriter(int maxPacketSize) { - this.idGenerator = new PacketIdGenerator(); - this.messagesPacketBuilder = new PacketBuilder(maxPacketSize, idGenerator, new MessagesPacketFactory()); - this.transEventPacketBuilder = new PacketBuilder(maxPacketSize, idGenerator, new TransPacketFactory()); + /** + * Create a PacketWriter with an expected max packet size. + *

+ * In theory we would prefer to create packets up to the MTU size which for + * Ethernet will likely be 1500. Note that the maxPacketSize is ignored for + * large single messages. + *

+ */ + public PacketWriter(int maxPacketSize) { + this.idGenerator = new PacketIdGenerator(); + this.messagesPacketBuilder = new PacketBuilder(maxPacketSize, idGenerator, new MessagesPacketFactory()); + this.transEventPacketBuilder = new PacketBuilder(maxPacketSize, idGenerator, new TransPacketFactory()); + } + + /** + * Return the currentPacketId. + */ + public long currentPacketId() { + return idGenerator.currentPacketId(); + } + + /** + * Create Packets for a given list of messages. + *

+ * Typically this creates a single Packet but there is a hard limit for UDP + * packet sizes. + *

+ */ + public List write(boolean requiresAck, List messages) throws IOException { + + BinaryMessageList binaryMsgList = new BinaryMessageList(); + for (int i = 0; i < messages.size(); i++) { + Message message = messages.get(i); + message.writeBinaryMessage(binaryMsgList); + } + return messagesPacketBuilder.write(requiresAck, binaryMsgList, ""); + } + + /** + * Create Packets for a given RemoteTransactionEvent. + *

+ * Typically this creates a single Packet but there is a hard limit for UDP + * packet sizes. + *

+ */ + public List write(RemoteTransactionEvent transEvent) throws IOException { + + BinaryMessageList messageList = new BinaryMessageList(); + + // split into reasonably small independent messages + transEvent.writeBinaryMessage(messageList); + + return transEventPacketBuilder.write(true, messageList, transEvent.getServerName()); + } + + /** + * Reuse the same packetIdCounter for building Packets for both Message and + * RemoteTransactionEvent + */ + private static class PacketIdGenerator { + + long packetIdCounter; + + public long nextPacketId() { + return ++packetIdCounter; } - /** - * Return the currentPacketId. - */ public long currentPacketId() { - return idGenerator.currentPacketId(); + return packetIdCounter; } - - /** - * Create Packets for a given list of messages. - *

- * Typically this creates a single Packet but there is a hard limit for UDP - * packet sizes. - *

- */ - public List write(boolean requiresAck, List messages) throws IOException { - - BinaryMessageList binaryMsgList = new BinaryMessageList(); - for (int i = 0; i < messages.size(); i++) { - Message message = messages.get(i); - message.writeBinaryMessage(binaryMsgList); - } - return messagesPacketBuilder.write(requiresAck, binaryMsgList, ""); - } - - /** - * Create Packets for a given RemoteTransactionEvent. - *

- * Typically this creates a single Packet but there is a hard limit for UDP - * packet sizes. - *

- */ - public List write(RemoteTransactionEvent transEvent) throws IOException { - - BinaryMessageList messageList = new BinaryMessageList(); - - // split into reasonably small independent messages - transEvent.writeBinaryMessage(messageList); - - return transEventPacketBuilder.write(true, messageList, transEvent.getServerName()); - } - - /** - * Reuse the same packetIdCounter for building Packets for both Message and - * RemoteTransactionEvent - */ - private static class PacketIdGenerator { - - long packetIdCounter; - - public long nextPacketId() { - return ++packetIdCounter; - } - - public long currentPacketId() { - return packetIdCounter; - } - - } - - interface PacketFactory { - - Packet createPacket(long packetId, long timestamp, String serverName) throws IOException; - } - - private static class TransPacketFactory implements PacketFactory { - - public Packet createPacket(long packetId, long timestamp, String serverName) throws IOException { - return PacketTransactionEvent.forWrite(packetId, timestamp, serverName); - } - } - - private static class MessagesPacketFactory implements PacketFactory { - - public Packet createPacket(long packetId, long timestamp, String serverName) throws IOException { - return PacketMessages.forWrite(packetId, timestamp, serverName); - } - } - - /** - * Helper class for building Packets from messages or - * RemoteTransactionEvents. - */ - private static class PacketBuilder { - - private final PacketIdGenerator idGenerator; - private final PacketFactory packetFactory; - private final int maxPacketSize; - - private PacketBuilder(int maxPacketSize, PacketIdGenerator idGenerator, PacketFactory packetFactory) { - this.maxPacketSize = maxPacketSize; - this.idGenerator = idGenerator; - this.packetFactory = packetFactory; - } - - private List write(boolean requiresAck, BinaryMessageList messageList, String serverName) - throws IOException { - - List list = messageList.getList(); - - ArrayList packets = new ArrayList(1); - - long timestamp = System.currentTimeMillis(); - - long packetId = requiresAck ? idGenerator.nextPacketId() : 0; - Packet p = packetFactory.createPacket(packetId, timestamp, serverName); - - packets.add(p); - - for (int i = 0; i < list.size(); i++) { - BinaryMessage binMsg = list.get(i); - if (!p.writeBinaryMessage(binMsg, maxPacketSize)) { - // didn't fit into the package so put into another packet - packetId = requiresAck ? idGenerator.nextPacketId() : 0; - p = packetFactory.createPacket(packetId, timestamp, serverName); - packets.add(p); - p.writeBinaryMessage(binMsg, maxPacketSize); - } - } - p.writeEof(); - - return packets; + } + + interface PacketFactory { + + Packet createPacket(long packetId, long timestamp, String serverName) throws IOException; + } + + private static class TransPacketFactory implements PacketFactory { + + public Packet createPacket(long packetId, long timestamp, String serverName) throws IOException { + return PacketTransactionEvent.forWrite(packetId, timestamp, serverName); + } + } + + private static class MessagesPacketFactory implements PacketFactory { + + public Packet createPacket(long packetId, long timestamp, String serverName) throws IOException { + return PacketMessages.forWrite(packetId, timestamp, serverName); + } + } + + /** + * Helper class for building Packets from messages or + * RemoteTransactionEvents. + */ + private static class PacketBuilder { + + private final PacketIdGenerator idGenerator; + private final PacketFactory packetFactory; + private final int maxPacketSize; + + private PacketBuilder(int maxPacketSize, PacketIdGenerator idGenerator, PacketFactory packetFactory) { + this.maxPacketSize = maxPacketSize; + this.idGenerator = idGenerator; + this.packetFactory = packetFactory; + } + + private List write(boolean requiresAck, BinaryMessageList messageList, String serverName) + throws IOException { + + List list = messageList.getList(); + + ArrayList packets = new ArrayList(1); + + long timestamp = System.currentTimeMillis(); + + long packetId = requiresAck ? idGenerator.nextPacketId() : 0; + Packet p = packetFactory.createPacket(packetId, timestamp, serverName); + + packets.add(p); + + for (int i = 0; i < list.size(); i++) { + BinaryMessage binMsg = list.get(i); + if (!p.writeBinaryMessage(binMsg, maxPacketSize)) { + // didn't fit into the package so put into another packet + packetId = requiresAck ? idGenerator.nextPacketId() : 0; + p = packetFactory.createPacket(packetId, timestamp, serverName); + packets.add(p); + p.writeBinaryMessage(binMsg, maxPacketSize); } + } + p.writeEof(); + + return packets; + } + } } diff --git a/src/main/java/com/avaje/ebeaninternal/server/cluster/SerialiseTransactionHelper.java b/src/main/java/com/avaje/ebeaninternal/server/cluster/SerialiseTransactionHelper.java index 0086fd6ff..daab0f50d 100644 --- a/src/main/java/com/avaje/ebeaninternal/server/cluster/SerialiseTransactionHelper.java +++ b/src/main/java/com/avaje/ebeaninternal/server/cluster/SerialiseTransactionHelper.java @@ -1,15 +1,15 @@ package com.avaje.ebeaninternal.server.cluster; -import java.io.ByteArrayInputStream; -import java.io.DataInputStream; -import java.io.IOException; -import java.util.List; - import com.avaje.ebeaninternal.api.SpiEbeanServer; import com.avaje.ebeaninternal.server.transaction.RemoteTransactionEvent; import org.slf4j.Logger; import org.slf4j.LoggerFactory; +import java.io.ByteArrayInputStream; +import java.io.DataInputStream; +import java.io.IOException; +import java.util.List; + /** * Mechanism to convert RemoteTransactionEvent to/from byte[] content. */ diff --git a/src/main/java/com/avaje/ebeaninternal/server/cluster/package.html b/src/main/java/com/avaje/ebeaninternal/server/cluster/package.html index 66d055461..934649700 100644 --- a/src/main/java/com/avaje/ebeaninternal/server/cluster/package.html +++ b/src/main/java/com/avaje/ebeaninternal/server/cluster/package.html @@ -1,12 +1,12 @@ - - AvajeLib + + AvajeLib Clustering service for an application.

-A framework for supporting clustering of servers. + A framework for supporting clustering of servers.

\ No newline at end of file diff --git a/src/main/java/com/avaje/ebeaninternal/server/cluster/socket/RequestProcessor.java b/src/main/java/com/avaje/ebeaninternal/server/cluster/socket/RequestProcessor.java index 288584717..a4c89516d 100644 --- a/src/main/java/com/avaje/ebeaninternal/server/cluster/socket/RequestProcessor.java +++ b/src/main/java/com/avaje/ebeaninternal/server/cluster/socket/RequestProcessor.java @@ -8,39 +8,39 @@ import java.net.Socket; /** * This parses and dispatches a request to the appropriate handler. *

- * Looks up the appropriate RequestHandler + * Looks up the appropriate RequestHandler * and then gets it to process the Client request.

*

* Note that this is a Runnable because it is assigned to the ThreadPool. */ class RequestProcessor implements Runnable { - private static final Logger logger = LoggerFactory.getLogger(RequestProcessor.class); - + private static final Logger logger = LoggerFactory.getLogger(RequestProcessor.class); + private final Socket clientSocket; - + private final SocketClusterBroadcast owner; private final String hostPort; /** - * Create including the Listener (used to lookup the Request Handler) and - * the socket itself. - */ - public RequestProcessor(SocketClusterBroadcast owner, Socket clientSocket) { - this.clientSocket = clientSocket; - this.owner = owner; + * Create including the Listener (used to lookup the Request Handler) and + * the socket itself. + */ + public RequestProcessor(SocketClusterBroadcast owner, Socket clientSocket) { + this.clientSocket = clientSocket; + this.owner = owner; this.hostPort = owner.getHostPort(); } - - /** - * This will parse out the command. Lookup the appropriate Handler and - * pass the information to the handler for processing. - *

Dev Note: the command parsing is processed here so that it is preformed - * by the assigned thread rather than the listeners thread.

- */ - public void run() { - try { + + /** + * This will parse out the command. Lookup the appropriate Handler and + * pass the information to the handler for processing. + *

Dev Note: the command parsing is processed here so that it is preformed + * by the assigned thread rather than the listeners thread.

+ */ + public void run() { + try { logger.trace("start listening for cluster messages"); SocketConnection sc = new SocketConnection(clientSocket); while (true) { @@ -50,10 +50,10 @@ class RequestProcessor implements Runnable { } } logger.trace("disconnecting: {}", hostPort); - sc.disconnect(); + sc.disconnect(); } catch (Exception e) { - logger.error("Error listening for messages - "+owner.getHostPort(), e); + logger.error("Error listening for messages - " + owner.getHostPort(), e); } } diff --git a/src/main/java/com/avaje/ebeaninternal/server/cluster/socket/SocketClusterBroadcast.java b/src/main/java/com/avaje/ebeaninternal/server/cluster/socket/SocketClusterBroadcast.java index 15fcf1c13..d77878b50 100644 --- a/src/main/java/com/avaje/ebeaninternal/server/cluster/socket/SocketClusterBroadcast.java +++ b/src/main/java/com/avaje/ebeaninternal/server/cluster/socket/SocketClusterBroadcast.java @@ -156,7 +156,7 @@ public class SocketClusterBroadcast implements ClusterBroadcast { protected void broadcast(SocketClusterMessage msg) { if (logger.isTraceEnabled()) { - logger.trace("... broadcast msg: "+msg); + logger.trace("... broadcast msg: " + msg); } for (int i = 0; i < members.length; i++) { send(members[i], msg); diff --git a/src/main/java/com/avaje/ebeaninternal/server/cluster/socket/SocketClusterListener.java b/src/main/java/com/avaje/ebeaninternal/server/cluster/socket/SocketClusterListener.java index cd0824f24..10d55b978 100644 --- a/src/main/java/com/avaje/ebeaninternal/server/cluster/socket/SocketClusterListener.java +++ b/src/main/java/com/avaje/ebeaninternal/server/cluster/socket/SocketClusterListener.java @@ -1,15 +1,15 @@ package com.avaje.ebeaninternal.server.cluster.socket; +import com.avaje.ebeaninternal.server.lib.DaemonThreadPool; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + import java.io.IOException; import java.io.InterruptedIOException; import java.net.ServerSocket; import java.net.Socket; import java.net.SocketException; -import com.avaje.ebeaninternal.server.lib.DaemonThreadPool; -import org.slf4j.Logger; -import org.slf4j.LoggerFactory; - /** * Serverside multithreaded socket listener. Accepts connections and dispatches