diff --git a/bmq-sdk/src/main/java/com/bloomberg/bmq/impl/PutPoster.java b/bmq-sdk/src/main/java/com/bloomberg/bmq/impl/PutPoster.java index 9969d7a9..e8dc4cf2 100644 --- a/bmq-sdk/src/main/java/com/bloomberg/bmq/impl/PutPoster.java +++ b/bmq-sdk/src/main/java/com/bloomberg/bmq/impl/PutPoster.java @@ -103,9 +103,7 @@ private void sendEvent() { } try { - EventBuilderResult packResult = - putBuilder.packMessage( - msgImpl, brokerConnection.isOldStyleMessageProperties()); + EventBuilderResult packResult = putBuilder.packMessage(msgImpl); if (packResult == EventBuilderResult.EVENT_TOO_BIG) { // Put the current message back to the deque putMessages.addFirst(msgImpl); diff --git a/bmq-sdk/src/main/java/com/bloomberg/bmq/impl/TcpBrokerConnection.java b/bmq-sdk/src/main/java/com/bloomberg/bmq/impl/TcpBrokerConnection.java index c9eec965..e762d5a9 100644 --- a/bmq-sdk/src/main/java/com/bloomberg/bmq/impl/TcpBrokerConnection.java +++ b/bmq-sdk/src/main/java/com/bloomberg/bmq/impl/TcpBrokerConnection.java @@ -105,7 +105,6 @@ public class TcpBrokerConnection private volatile StopCallback stopCallback; private volatile Duration stopTimeout; - private volatile boolean isOldStyleMessageProperties = false; private ScheduledExecutorService scheduler; private ScheduledFuture onAuthenticationTimeoutFuture; @@ -659,16 +658,6 @@ private boolean validateBrokerResponse(BrokerResponse resp) { && resp.getOriginalRequest() != null) { brokerIdentity = resp.getOriginalRequest(); - - // TODO: remove after 2nd rollout of "new style" brokers - String brokerFeatures = brokerIdentity.features(); - isOldStyleMessageProperties = - brokerFeatures == null - || brokerFeatures.isEmpty() - || !brokerFeatures.toUpperCase().contains(MPS_EX_FEATURE); - logger.info( - "Broker supports new style message properties: {}", - !isOldStyleMessageProperties); isValid = true; } else { logger.error("Broker response is invalid"); @@ -792,11 +781,6 @@ public GenericResult linger() { return GenericResult.SUCCESS; } - @Override - public boolean isOldStyleMessageProperties() { - return isOldStyleMessageProperties; - } - @Override public GenericResult write(ByteBuffer[] buffers, boolean waitUntilWritable) { diff --git a/bmq-sdk/src/main/java/com/bloomberg/bmq/impl/infr/proto/ApplicationData.java b/bmq-sdk/src/main/java/com/bloomberg/bmq/impl/infr/proto/ApplicationData.java index b4e28151..a31f2c6a 100644 --- a/bmq-sdk/src/main/java/com/bloomberg/bmq/impl/infr/proto/ApplicationData.java +++ b/bmq-sdk/src/main/java/com/bloomberg/bmq/impl/infr/proto/ApplicationData.java @@ -20,9 +20,6 @@ import com.bloomberg.bmq.impl.infr.util.Compression; import com.bloomberg.bmq.impl.infr.util.PrintUtil; import java.io.ByteArrayOutputStream; -import java.io.DataInput; -import java.io.DataInputStream; -import java.io.DataOutputStream; import java.io.IOException; import java.io.InputStream; import java.io.OutputStream; @@ -41,14 +38,9 @@ public class ApplicationData { private CompressionAlgorithmType compressionType = CompressionAlgorithmType.E_NONE; private ByteBufferOutputStream compressedData; - // TODO: remove after 2nd release of "new style" brokers. - private boolean isOldStyleProperties = false; - private boolean arePropertiesCompressed; - private void resetCompressedData() { compressionType = CompressionAlgorithmType.E_NONE; compressedData = null; - arePropertiesCompressed = false; } public final void setPayload(ByteBuffer... data) throws IOException { @@ -76,15 +68,6 @@ public final void setProperties(MessagePropertiesImpl props) { resetCompressedData(); } - // TODO: remove after 2nd release of "new style" brokers. - public void setIsOldStyleProperties(boolean value) { - isOldStyleProperties = value; - } - - public boolean isOldStyleProperties() { - return isOldStyleProperties; - } - public ByteBuffer[] applicationData() throws IOException { // TODO: used only to calculate CRC32. Can we avoid creating a copy? @@ -106,14 +89,6 @@ public ByteBuffer[] payload() throws IOException { } public MessagePropertiesImpl properties() { - if (arePropertiesCompressed) { - try { - decompressData(); - } catch (IOException e) { - throw new RuntimeException("Failed to decompress payload", e); - } - } - return properties; } @@ -134,11 +109,7 @@ public int propertiesSize() { } public int unpackedSize() { - int size = 0; - - if (!arePropertiesCompressed) { - size += propertiesSize(); - } + int size = propertiesSize(); if (compressionType == CompressionAlgorithmType.E_NONE) { size += payloadSize(); @@ -164,11 +135,9 @@ public int numPaddingBytes() { return ProtocolUtil.calculatePadding(unpackedSize()); } - // TODO: remove "isOldStyleProperties" after 2nd release of "new style" brokers. public void streamIn( int size, boolean hasProperties, - boolean isOldStyleProperties, CompressionAlgorithmType compressionType, ByteBufferInputStream bbis) throws IOException { @@ -189,16 +158,9 @@ public void streamIn( resetCompressedData(); - // Stream in properties if they are not compressed - this.isOldStyleProperties = isOldStyleProperties; + // Properties are never compressed if (hasProperties) { - if (!isOldStyleProperties) { - // New properties - size -= streamInProperties(bbis); - } else if (compressionType == CompressionAlgorithmType.E_NONE) { - // Old properties - size -= streamInPropertiesOld(bbis); - } + size -= streamInProperties(bbis); } // Stream uncompressed payload @@ -221,7 +183,6 @@ public void streamIn( compressedData = bbos; this.compressionType = compressionType; - arePropertiesCompressed = hasProperties && isOldStyleProperties; } } @@ -245,16 +206,9 @@ private void decompressData() throws IOException { ByteBufferInputStream bbis = new ByteBufferInputStream(data); InputStream decompressedStream = compressionType.getCompression().decompress(bbis); - DataInputStream inputStream = new DataInputStream(decompressedStream); - - // Stream in properties - if (arePropertiesCompressed) { - // If properties are compressed then they are encoded in old format - streamInPropertiesOld(inputStream); - } // Stream in payload - streamInPayload(inputStream); + streamInPayload(decompressedStream); // Check if all data has been read if (bbis.available() > 0) { @@ -273,16 +227,6 @@ private int streamInProperties(ByteBufferInputStream input) throws IOException { return read; } - private int streamInPropertiesOld(T input) - throws IOException { - int read = 0; - - properties = new MessagePropertiesImpl(); - read += properties.streamInOld(input); - - return read; - } - private int streamInPayload(int size, ByteBufferInputStream bbis) throws IOException { payload = new byte[size]; @@ -333,22 +277,14 @@ public void compressData(CompressionAlgorithmType compressionType) throws IOExce // // Later we will need to refactor the code in order to close compressed stream without // closing underlying stream - try (OutputStream compressedStream = compression.compress(bbos); - DataOutputStream compressedOutput = new DataOutputStream(compressedStream)) { - - // TODO: remove after 2nd rollout of "new style" brokers. - if (hasProperties() && isOldStyleProperties) { - properties.streamOutOld(compressedOutput); - } - + try (OutputStream compressedStream = compression.compress(bbos)) { if (payload != null) { - compressedOutput.write(payload); + compressedStream.write(payload); } } compressedData = bbos; this.compressionType = compressionType; - arePropertiesCompressed = hasProperties() && isOldStyleProperties; } public void streamOut(ByteBufferOutputStream bbos) throws IOException { @@ -358,14 +294,8 @@ public void streamOut(ByteBufferOutputStream bbos) throws IOException { private void streamOut(ByteBufferOutputStream bbos, boolean addPadding) throws IOException { int startPosition = bbos.size(); - // Stream out properties if they are not compressed (no compression or - // new style properties). - if (hasProperties() && !arePropertiesCompressed) { - if (isOldStyleProperties) { - properties.streamOutOld(bbos); - } else { - properties.streamOut(bbos); - } + if (hasProperties()) { + properties.streamOut(bbos); } if (compressionType == CompressionAlgorithmType.E_NONE) { diff --git a/bmq-sdk/src/main/java/com/bloomberg/bmq/impl/infr/proto/MessagePropertiesImpl.java b/bmq-sdk/src/main/java/com/bloomberg/bmq/impl/infr/proto/MessagePropertiesImpl.java index 00a1ca05..b36226dc 100644 --- a/bmq-sdk/src/main/java/com/bloomberg/bmq/impl/infr/proto/MessagePropertiesImpl.java +++ b/bmq-sdk/src/main/java/com/bloomberg/bmq/impl/infr/proto/MessagePropertiesImpl.java @@ -17,10 +17,8 @@ import com.bloomberg.bmq.MessageProperties; import com.bloomberg.bmq.impl.infr.io.ByteBufferInputStream; -import java.io.DataInput; import java.io.DataOutput; import java.io.IOException; -import java.io.InputStream; import java.lang.invoke.MethodHandles; import java.nio.charset.CharsetEncoder; import java.nio.charset.StandardCharsets; @@ -127,124 +125,6 @@ public MessageProperty get(String name) { return propertyMap.get(name); } - // TODO: remove after 2nd rollout of "new style" brokers - public int streamInOld(T input) throws IOException { - propertyMap.clear(); - MessagePropertiesHeader propsHeader = new MessagePropertiesHeader(); - propsHeader.streamIn(input); - - final int numProps = propsHeader.numProperties(); - MessagePropertyHeader[] propHeaderArray = new MessagePropertyHeader[numProps]; - - final int propertyHeaderSize = propsHeader.messagePropertyHeaderSize(); - for (int i = 0; i < numProps; i++) { - MessagePropertyHeader ph = new MessagePropertyHeader(); - ph.streamIn(input, propertyHeaderSize); - propHeaderArray[i] = ph; - } - - final int propsAreaSize = propsHeader.messagePropertiesAreaWords() * Protocol.WORD_SIZE; - final int headersAreaSize = propsHeader.headerSize() + numProps * propertyHeaderSize; - - // Since the input stream type doesn't support setting of position, - // we cannot determine padding bytes before we read all properties. - int totalLength = headersAreaSize; - - for (int i = 0; i < numProps; ++i) { - MessagePropertyHeader ph = propHeaderArray[i]; - MessageProperty mp; - PropertyType t = PropertyType.fromInt(ph.propertyType()); - switch (t) { - case BOOL: - mp = new BoolMessageProperty(); - break; - case BYTE: - mp = new ByteMessageProperty(); - break; - case SHORT: - mp = new ShortMessageProperty(); - break; - case INT32: - mp = new Int32MessageProperty(); - break; - case INT64: - mp = new Int64MessageProperty(); - break; - case STRING: - mp = new StringMessageProperty(); - break; - case BINARY: - mp = new BinaryMessageProperty(); - break; - default: - throw new IOException("Unknown property type"); - } - final int nameLength = ph.propertyNameLength(); - final int valueLength = ph.propertyValueLength(); - - totalLength += nameLength; - totalLength += valueLength; - - byte[] n = new byte[nameLength]; - byte[] v = new byte[valueLength]; - - // Since 'read(byte[])' might read just part of bytes, 'readFully(byte[])' is used - // instead - try { - input.readFully(n); - } catch (IOException e) { - throw new IOException( - "Error when reading property name. Expected to read " + n.length + " bytes", - e); - } - mp.setPropertyName(new String(n, StandardCharsets.US_ASCII)); - - // Since 'read(byte[])' might read just part of bytes, 'readFully(byte[])' is used - // instead - try { - input.readFully(v); - } catch (IOException e) { - throw new IOException( - "Error when reading property value. Expected to read " - + v.length - + " bytes", - e); - } - mp.setPropertyValue(v); - - propertyMap.put(mp.name(), mp); - } - - // Read padding bytes - final byte numPaddingBytes = input.readByte(); - - // Skip padding bytes - if (input.skip(numPaddingBytes - 1) != numPaddingBytes - 1) { - throw new IOException("Failed to skip " + (numPaddingBytes - 1) + " bytes"); - } - - // Verify - final int numPaddingBytesExp = ProtocolUtil.calculatePadding(totalLength); - if (numPaddingBytesExp != numPaddingBytes) { - throw new IOException( - "Unexpected padding: " + numPaddingBytes + ", should be " + numPaddingBytesExp); - } - - // Add padding bytes - totalLength += numPaddingBytes; - - if (totalLength != propsAreaSize) { - throw new IOException( - "Invalid encoding: actual " - + totalLength - + " bytes, expected " - + propsAreaSize - + " bytes"); - } - - return totalLength; - } - public int streamIn(ByteBufferInputStream input) throws IOException { propertyMap.clear(); final int initPos = input.position(); @@ -317,13 +197,13 @@ public int streamIn(ByteBufferInputStream input) throws IOException { throw new IOException("Unknown property type"); } final int nameLength = ph.propertyNameLength(); - final int offset = ph.propertyValueLength(); + final int offset = ph.propertyValueOffset(); int valueLength; if (!isLastProperty) { // Calculate the length as delta between offsets minus // current property name length. - final int nextOffset = propHeaderArray[i + 1].propertyValueLength(); + final int nextOffset = propHeaderArray[i + 1].propertyValueOffset(); valueLength = nextOffset - offset - nameLength; } else { // Last property. @@ -385,18 +265,7 @@ public int streamIn(ByteBufferInputStream input) throws IOException { return totalLength; } - // TODO: remove after 2nd rollout of "new style" brokers - public void streamOutOld(DataOutput output) throws IOException { - streamOut(output, true); - } - - // TODO: remove after 2nd rollout of "new style" brokers public void streamOut(DataOutput output) throws IOException { - streamOut(output, false); - } - - // TODO: remove boolean after 2nd rollout of "new style" brokers - private void streamOut(DataOutput output, boolean isOldStyleProperties) throws IOException { final int numProps = propertyMap.size(); if (numProps == 0) { logger.info("No message properties to stream out"); @@ -413,12 +282,7 @@ private void streamOut(DataOutput output, boolean isOldStyleProperties) throws I MessagePropertyHeader mph = new MessagePropertyHeader(); mph.setPropertyType(e.getValue().type().toInt()); - if (isOldStyleProperties) { - mph.setPropertyValueLength(valLen); - } else { - mph.setPropertyValueLength(offset); - } - + mph.setPropertyValueOffset(offset); mph.setPropertyNameLength(nameLen); int propLen = nameLen + valLen; diff --git a/bmq-sdk/src/main/java/com/bloomberg/bmq/impl/infr/proto/MessagePropertyHeader.java b/bmq-sdk/src/main/java/com/bloomberg/bmq/impl/infr/proto/MessagePropertyHeader.java index 1b43f396..4a45c635 100644 --- a/bmq-sdk/src/main/java/com/bloomberg/bmq/impl/infr/proto/MessagePropertyHeader.java +++ b/bmq-sdk/src/main/java/com/bloomberg/bmq/impl/infr/proto/MessagePropertyHeader.java @@ -38,10 +38,10 @@ public final class MessagePropertyHeader { // R2..: Reserved (2nd set of bits) // // PropType...........: Data type of the message property - // PropValueLenUpper..: Upper 10 bits of the field capturing length of the - // property value. - // PropValueLenLower..: Lower 16 bits of the field capturing length of the - // property value. + // PropValueLenUpper..: Upper 10 bits of the field capturing the offset of + // the property value. + // PropValueLenLower..: Lower 16 bits of the field capturing the offset of + // the property value. // PropNameLen........: Length of the property name. // Reserved...........: For alignment and extension ~ must be 0 // .. @@ -98,8 +98,7 @@ public void setPropertyType(int value) { | (value << PROP_TYPE_START_IDX)); } - // TODO: rename to offset after 2nd rollout of "new style" brokers - public void setPropertyValueLength(int value) { + public void setPropertyValueOffset(int value) { Argument.expectNonNegative(value, "value"); Argument.expectNotGreater(value, MAX_PROPERTY_VALUE_LENGTH, "value"); @@ -126,8 +125,7 @@ public int propertyType() { return result >>> PROP_TYPE_START_IDX; } - // TODO: rename to offset after 2nd rollout of "new style" brokers - public int propertyValueLength() { + public int propertyValueOffset() { int result = (propTypeAndPropValueLenUpper & PROP_VALUE_LEN_UPPER_MASK) << PROP_VALUE_LEN_LOWER_NUM_BITS; @@ -156,14 +154,14 @@ public void streamIn(DataInput input, int size) throws IOException { } final int propNameLen = propertyNameLength(); - final int propValueLen = propertyValueLength(); + final int propValueOffset = propertyValueOffset(); if (MAX_PROPERTY_NAME_LENGTH < propNameLen) { throw new IOException("Invalid property name length: [" + propNameLen + "]"); } - if (MAX_PROPERTY_VALUE_LENGTH < propValueLen) { - throw new IOException("Invalid property value length: [" + propValueLen + "]"); + if (MAX_PROPERTY_VALUE_LENGTH < propValueOffset) { + throw new IOException("Invalid property value offset: [" + propValueOffset + "]"); } // Skip unknown bytes @@ -187,8 +185,8 @@ public String toString() { sb.append("[ MessagePropertyHeader [") .append(" PropertyType=") .append(propertyType()) - .append(" PropertyValueLength=") - .append(propertyValueLength()) + .append(" PropertyValueOffset=") + .append(propertyValueOffset()) .append(" PropertyNameLength=") .append(propertyNameLength()) .append(" ] ]"); diff --git a/bmq-sdk/src/main/java/com/bloomberg/bmq/impl/infr/proto/PushEventBuilder.java b/bmq-sdk/src/main/java/com/bloomberg/bmq/impl/infr/proto/PushEventBuilder.java index 1cab7d3e..ba069505 100644 --- a/bmq-sdk/src/main/java/com/bloomberg/bmq/impl/infr/proto/PushEventBuilder.java +++ b/bmq-sdk/src/main/java/com/bloomberg/bmq/impl/infr/proto/PushEventBuilder.java @@ -32,17 +32,13 @@ public void reset() { super.reset(EventType.PUSH); } - // TODO: remove boolean after 2nd release of "new style" brokers // TODO: move to test code - public EventBuilderResult packMessage(PushMessageImpl msg, boolean isOldStyleProperties) - throws IOException { + public EventBuilderResult packMessage(PushMessageImpl msg) throws IOException { // Warn if payload is empty if (msg.appData().payloadSize() == 0) { logger.warn("PUSH message payload is empty"); } - msg.appData().setIsOldStyleProperties(isOldStyleProperties); - // Compress data msg.compressData(); @@ -56,10 +52,7 @@ public EventBuilderResult packMessage(PushMessageImpl msg, boolean isOldStylePro int numPaddingBytes = msg.appData().numPaddingBytes(); final int sizeNoOptions = - bbos.size() - + PushHeader.HEADER_SIZE_FOR_SCHEMA_ID - + appDataLength - + numPaddingBytes; + bbos.size() + PushHeader.HEADER_SIZE + appDataLength + numPaddingBytes; if (sizeNoOptions > EventHeader.MAX_SIZE_SOFT) { return EventBuilderResult.EVENT_TOO_BIG; // RETURN diff --git a/bmq-sdk/src/main/java/com/bloomberg/bmq/impl/infr/proto/PushHeader.java b/bmq-sdk/src/main/java/com/bloomberg/bmq/impl/infr/proto/PushHeader.java index 3055a719..5514602c 100644 --- a/bmq-sdk/src/main/java/com/bloomberg/bmq/impl/infr/proto/PushHeader.java +++ b/bmq-sdk/src/main/java/com/bloomberg/bmq/impl/infr/proto/PushHeader.java @@ -144,19 +144,18 @@ public final class PushHeader { // Minimum size (bytes) of a 'PushHeader' (that is sufficient to // capture header words). This value should *never* change. - public static final int HEADER_SIZE = 28; + public static final int HEADER_SIZE = 32; // Current size (bytes) of the header. - // TODO: set to 32 after 2nd release of "new style" brokers - public static final int HEADER_SIZE_FOR_SCHEMA_ID = 32; + public static final int HEADER_SIZE_WITHOUT_SCHEMA_ID = 28; - // Current size (bytes) of the header with schema id - // TODO: remove after 2nd release of "new style" brokers + // Size (bytes) of the header without the schema id. Such headers are + // still accepted when streaming in. public PushHeader() { messageGUID = new byte[MessageGUID.SIZE_BINARY]; - setMessageWords((byte) (HEADER_SIZE_FOR_SCHEMA_ID / Protocol.WORD_SIZE)); - setHeaderWords((byte) (HEADER_SIZE_FOR_SCHEMA_ID / Protocol.WORD_SIZE)); + setMessageWords((byte) (HEADER_SIZE / Protocol.WORD_SIZE)); + setHeaderWords((byte) (HEADER_SIZE / Protocol.WORD_SIZE)); } public void setMessageWords(int value) { @@ -237,7 +236,7 @@ public void streamIn(ByteBufferInputStream bbis) throws IOException { optionsWordsAndHeaderWords = bbis.readInt(); final int headerSize = headerWords() * Protocol.WORD_SIZE; - if (headerSize < HEADER_SIZE) { + if (headerSize < HEADER_SIZE_WITHOUT_SCHEMA_ID) { throw new IOException("Invalid size: " + headerSize); } @@ -246,12 +245,10 @@ public void streamIn(ByteBufferInputStream bbis) throws IOException { messageGUID[i] = bbis.readByte(); } - int numRead = HEADER_SIZE; + int numRead = HEADER_SIZE_WITHOUT_SCHEMA_ID; - // Check if it's new header with schema id schemaWireId = 0; - // TODO: update after 2nd release of "new style" brokers - if (headerSize >= HEADER_SIZE_FOR_SCHEMA_ID) { + if (headerSize >= HEADER_SIZE) { schemaWireId = bbis.readShort(); reserved = bbis.readShort(); numRead += 4; @@ -269,7 +266,7 @@ public void streamIn(ByteBufferInputStream bbis) throws IOException { public void streamOut(ByteBufferOutputStream bbos) throws IOException { final int headerSize = headerWords() * Protocol.WORD_SIZE; - if (headerSize != HEADER_SIZE_FOR_SCHEMA_ID) { + if (headerSize != HEADER_SIZE) { throw new IOException("Invalid size: " + headerSize); } diff --git a/bmq-sdk/src/main/java/com/bloomberg/bmq/impl/infr/proto/PushMessageImpl.java b/bmq-sdk/src/main/java/com/bloomberg/bmq/impl/infr/proto/PushMessageImpl.java index f929b04f..b27b8ffa 100644 --- a/bmq-sdk/src/main/java/com/bloomberg/bmq/impl/infr/proto/PushMessageImpl.java +++ b/bmq-sdk/src/main/java/com/bloomberg/bmq/impl/infr/proto/PushMessageImpl.java @@ -31,8 +31,7 @@ public final class PushMessageImpl implements Streamable { static final Logger logger = LoggerFactory.getLogger(MethodHandles.lookup().lookupClass()); public static final short INVALID_SCHEMA_WIRE_ID = (short) 1; - private static final int HEADER_WORDS = - PushHeader.HEADER_SIZE_FOR_SCHEMA_ID / Protocol.WORD_SIZE; + private static final int HEADER_WORDS = PushHeader.HEADER_SIZE / Protocol.WORD_SIZE; private PushHeader header; private ApplicationData appData; @@ -141,7 +140,6 @@ public void streamIn(ByteBufferInputStream bbis) throws IOException { PushHeaderFlags.isSet(header.flags(), PushHeaderFlags.MESSAGE_PROPERTIES); final CompressionAlgorithmType inputCompressionType = CompressionAlgorithmType.fromInt(header.compressionType()); - final boolean isOldStyleProperties = header.schemaWireId() == 0; logger.debug( "Has properties: {}, compressionType: {}, schemaWireId: {}", @@ -153,7 +151,7 @@ public void streamIn(ByteBufferInputStream bbis) throws IOException { throw new BMQException("IMPLICIT_PAYLOAD flag is set"); } - appData.streamIn(dataSize, hasProperties, isOldStyleProperties, inputCompressionType, bbis); + appData.streamIn(dataSize, hasProperties, inputCompressionType, bbis); if (appData.unpackedSize() == 0) { throw new BMQException("Application data is empty"); @@ -168,12 +166,8 @@ public void compressData() throws IOException { CompressionAlgorithmType finalCompressionType = this.compressionType; + // Properties are not compressed. int dataToCompress = appData.payloadSize(); - // New style properties are not compressed. - // TODO: remove after 2nd rollout of "new style" brokers. - if (appData.hasProperties() && appData.isOldStyleProperties()) { - dataToCompress += appData.propertiesSize(); - } // When data is less than a threshold, it is not compressed. if (dataToCompress < Protocol.COMPRESSION_MIN_APPDATA_SIZE) { @@ -193,12 +187,8 @@ public void streamOut(ByteBufferOutputStream bbos) throws IOException { header.setFlags(f); } - // If properties are encoded using new style, we need to set - // schema wire id to 1 (invalid schema wire id). - // TODO: always set after 2nd rollout of "new style" brokers. - if (!appData.isOldStyleProperties()) { - header.setSchemaWireId(INVALID_SCHEMA_WIRE_ID); - } + // Set schema wire id to 1 (invalid schema wire id). + header.setSchemaWireId(INVALID_SCHEMA_WIRE_ID); } if (appData.payloadSize() == 0) { diff --git a/bmq-sdk/src/main/java/com/bloomberg/bmq/impl/infr/proto/PutEventBuilder.java b/bmq-sdk/src/main/java/com/bloomberg/bmq/impl/infr/proto/PutEventBuilder.java index d7c52a01..7894a29b 100644 --- a/bmq-sdk/src/main/java/com/bloomberg/bmq/impl/infr/proto/PutEventBuilder.java +++ b/bmq-sdk/src/main/java/com/bloomberg/bmq/impl/infr/proto/PutEventBuilder.java @@ -55,16 +55,12 @@ public void setMaxEventSize(int value) { } } - // TODO: remove boolean after 2nd release of "new style" brokers - public EventBuilderResult packMessage(PutMessageImpl msg, boolean isOldStyleProperties) - throws IOException { + public EventBuilderResult packMessage(PutMessageImpl msg) throws IOException { // Validate payload is empty if (msg.appData().payloadSize() == 0) { return EventBuilderResult.PAYLOAD_EMPTY; // RETURN } - msg.appData().setIsOldStyleProperties(isOldStyleProperties); - // Compress data msg.compressData(); msg.calculateAndSetCrc32c(); diff --git a/bmq-sdk/src/main/java/com/bloomberg/bmq/impl/infr/proto/PutMessageImpl.java b/bmq-sdk/src/main/java/com/bloomberg/bmq/impl/infr/proto/PutMessageImpl.java index 07c28679..1f754206 100644 --- a/bmq-sdk/src/main/java/com/bloomberg/bmq/impl/infr/proto/PutMessageImpl.java +++ b/bmq-sdk/src/main/java/com/bloomberg/bmq/impl/infr/proto/PutMessageImpl.java @@ -98,12 +98,8 @@ public void compressData() throws IOException { CompressionAlgorithmType finalCompressionType = this.compressionType; + // Properties are not compressed. int dataToCompress = appData.payloadSize(); - // New style properties are not compressed. - // TODO: remove after 2nd rollout of "new style" brokers. - if (appData.hasProperties() && appData.isOldStyleProperties()) { - dataToCompress += appData.propertiesSize(); - } // When data is less than a threshold, it is not compressed. if (dataToCompress < Protocol.COMPRESSION_MIN_APPDATA_SIZE) { @@ -162,9 +158,8 @@ public void streamIn(ByteBufferInputStream bbis) throws IOException { PutHeaderFlags.isSet(header.flags(), PutHeaderFlags.MESSAGE_PROPERTIES); final CompressionAlgorithmType inputCompressionType = CompressionAlgorithmType.fromInt(header.compressionType()); - final boolean isOldStyleProperties = header.schemaWireId() == 0; - appData.streamIn(dataSize, hasProperties, isOldStyleProperties, inputCompressionType, bbis); + appData.streamIn(dataSize, hasProperties, inputCompressionType, bbis); if (appData.unpackedSize() == 0) { throw new BMQException("Application data is empty."); @@ -179,14 +174,9 @@ public void streamOut(ByteBufferOutputStream bbos) throws IOException { setFlags(f); } - // If properties are encoded using new style, we need to set - // schema wire id to 1 (invalid schema wire id) in order to tell the - // broker that PUT message contains new style properties without - // schema id. - // TODO: always set after 2nd rollout of "new style" brokers. - if (!appData.isOldStyleProperties()) { - header.setSchemaWireId(INVALID_SCHEMA_WIRE_ID); - } + // Set schema wire id to 1 (invalid schema wire id) in order to tell + // the broker that PUT message contains properties without schema id. + header.setSchemaWireId(INVALID_SCHEMA_WIRE_ID); } final int numWords = ProtocolUtil.calculateNumWords(appData.unpackedSize()); diff --git a/bmq-sdk/src/main/java/com/bloomberg/bmq/impl/intf/BrokerConnection.java b/bmq-sdk/src/main/java/com/bloomberg/bmq/impl/intf/BrokerConnection.java index f476f905..1e48c993 100644 --- a/bmq-sdk/src/main/java/com/bloomberg/bmq/impl/intf/BrokerConnection.java +++ b/bmq-sdk/src/main/java/com/bloomberg/bmq/impl/intf/BrokerConnection.java @@ -57,7 +57,4 @@ interface StopCallback { GenericResult write(ByteBuffer[] buffers, boolean waitUntilWritable); GenericResult linger(); - - // TODO: remove after 2nd release of "new style" brokers - boolean isOldStyleMessageProperties(); } diff --git a/bmq-sdk/src/test/java/com/bloomberg/bmq/benchmark/ApplicationDataBenchmark.java b/bmq-sdk/src/test/java/com/bloomberg/bmq/benchmark/ApplicationDataBenchmark.java index 2719ee69..7a762a9b 100644 --- a/bmq-sdk/src/test/java/com/bloomberg/bmq/benchmark/ApplicationDataBenchmark.java +++ b/bmq-sdk/src/test/java/com/bloomberg/bmq/benchmark/ApplicationDataBenchmark.java @@ -38,15 +38,9 @@ public void testZlibStreamInOut() throws IOException { final int PAYLOAD_SIZE_BYTES = 1024 * 1024 * 2; // 2 Mb test.verifyStreamIn( - test.generatePayload(PAYLOAD_SIZE_BYTES), - test.generateProps(), - false, - compressionType); + test.generatePayload(PAYLOAD_SIZE_BYTES), test.generateProps(), compressionType); test.verifyStreamOut( - test.generatePayload(PAYLOAD_SIZE_BYTES), - test.generateProps(), - false, - compressionType); + test.generatePayload(PAYLOAD_SIZE_BYTES), test.generateProps(), compressionType); } } diff --git a/bmq-sdk/src/test/java/com/bloomberg/bmq/impl/ProtocolEventImplTcpReaderTest.java b/bmq-sdk/src/test/java/com/bloomberg/bmq/impl/ProtocolEventImplTcpReaderTest.java index 7801c23f..bdfcc2f5 100644 --- a/bmq-sdk/src/test/java/com/bloomberg/bmq/impl/ProtocolEventImplTcpReaderTest.java +++ b/bmq-sdk/src/test/java/com/bloomberg/bmq/impl/ProtocolEventImplTcpReaderTest.java @@ -15,36 +15,33 @@ */ package com.bloomberg.bmq.impl; +import static org.junit.jupiter.api.Assertions.assertArrayEquals; import static org.junit.jupiter.api.Assertions.assertEquals; import static org.junit.jupiter.api.Assertions.assertTrue; import static org.junit.jupiter.api.Assertions.fail; import com.bloomberg.bmq.MessageGUID; +import com.bloomberg.bmq.ResultCodes.AckResult; import com.bloomberg.bmq.impl.infr.io.ByteBufferInputStream; -import com.bloomberg.bmq.impl.infr.msg.MessagesTestSamples; +import com.bloomberg.bmq.impl.infr.io.ByteBufferOutputStream; import com.bloomberg.bmq.impl.infr.net.intf.TcpConnection.ReadCallback.ReadCompletionStatus; +import com.bloomberg.bmq.impl.infr.proto.AckEventBuilder; import com.bloomberg.bmq.impl.infr.proto.AckEventImpl; import com.bloomberg.bmq.impl.infr.proto.AckMessageImpl; -import com.bloomberg.bmq.impl.infr.proto.ControlEventImpl; -import com.bloomberg.bmq.impl.infr.proto.EventImpl; +import com.bloomberg.bmq.impl.infr.proto.EventBuilderResult; import com.bloomberg.bmq.impl.infr.proto.EventType; +import com.bloomberg.bmq.impl.infr.proto.MessagePropertiesImpl; import com.bloomberg.bmq.impl.infr.proto.PushEventBuilder; import com.bloomberg.bmq.impl.infr.proto.PushEventImpl; import com.bloomberg.bmq.impl.infr.proto.PushMessageImpl; import com.bloomberg.bmq.impl.infr.proto.PushMessageIterator; -import com.bloomberg.bmq.util.TestHelpers; -import java.io.BufferedReader; import java.io.IOException; -import java.io.InputStream; -import java.io.InputStreamReader; import java.lang.invoke.MethodHandles; import java.nio.ByteBuffer; import java.util.ArrayList; +import java.util.Arrays; import java.util.HashSet; import java.util.Iterator; -import java.util.concurrent.LinkedBlockingQueue; -import java.util.concurrent.Semaphore; -import java.util.concurrent.TimeUnit; import org.junit.jupiter.api.Test; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -53,143 +50,152 @@ class ProtocolEventImplTcpReaderTest { static final Logger logger = LoggerFactory.getLogger(MethodHandles.lookup().lookupClass()); - private ByteBuffer[] buildPushMessage(boolean isOldStyleProperties) throws IOException { - String PAYLOAD = "abcdefghijklmnopqrstuvwxyz"; - String GUID = "ABCDEF0123456789ABCDEF0123456789"; + private static final int QUEUE_ID = 9876; + private static final String PAYLOAD = "abcdefghijklmnopqrstuvwxyz"; + private static final String ROUTING_ID = "abcd-efgh-ijkl"; + private static final long TIMESTAMP = 123456789L; - MessageGUID guid = MessageGUID.fromHex(GUID); + private ByteBuffer[] buildPushMessage() throws IOException { + MessageGUID guid = MessageGUID.fromHex("ABCDEF0123456789ABCDEF0123456789"); PushMessageImpl pushMsg = new PushMessageImpl(); - pushMsg.setQueueId(9876); + pushMsg.setQueueId(QUEUE_ID); pushMsg.setMessageGUID(guid); pushMsg.appData().setPayload(ByteBuffer.wrap(PAYLOAD.getBytes())); PushEventBuilder builder = new PushEventBuilder(); - builder.packMessage(pushMsg, isOldStyleProperties); + builder.packMessage(pushMsg); return builder.build(); } + private PushMessageImpl createPushMessage(MessageGUID guid) throws IOException { + MessagePropertiesImpl props = new MessagePropertiesImpl(); + props.setPropertyAsString("routingId", ROUTING_ID); + props.setPropertyAsInt64("timestamp", TIMESTAMP); + + PushMessageImpl pushMsg = new PushMessageImpl(); + pushMsg.setQueueId(QUEUE_ID); + pushMsg.setMessageGUID(guid); + pushMsg.appData().setProperties(props); + pushMsg.appData().setPayload(ByteBuffer.wrap(PAYLOAD.getBytes())); + + return pushMsg; + } + @Test - void testIODump() throws IOException, InterruptedException { - logger.info("========================================================"); - logger.info("BEGIN Testing ProtocolEventImplTcpReaderTest testIODump."); - logger.info("========================================================"); + void testPushAndAckStream() throws IOException { + logger.info("==============================================================="); + logger.info("BEGIN Testing ProtocolEventImplTcpReaderTest PUSH and ACK stream."); + logger.info("==============================================================="); - // Check that ProtocolEventTcpReader correctly reads BlazingMQ events - // stored in IO dump file. + // Check that ProtocolEventTcpReader correctly reads a stream of PUSH + // and ACK events split into chunks which do not match event boundaries. // Steps: - // 1. Read dump file and feed ProtocolEventTcpReader by portions defined in index file; - // 2. From ProtocolEventTcpReader callback decode BlazingMQ events and put them into a - // queue; - // 3. Read events out of that queue from a separate thread emulating event handling - // in TcpBrokerConnection; - // 4. Collect GUIDs from ACK messages and verify each GUID from PUSH message has it's - // equivalent GUID from ACK message. - - final int NUM_PUSH_MESSAGES = 970; - - InputStream fis = - this.getClass().getResourceAsStream(MessagesTestSamples.BMQ_IO_DUMP_BIN.filePath()); - InputStream iis = - this.getClass().getResourceAsStream(MessagesTestSamples.BMQ_IO_DUMP_IDX.filePath()); - BufferedReader br = new BufferedReader(new InputStreamReader(iis)); - LinkedBlockingQueue eventQueue = new LinkedBlockingQueue<>(); - ProtocolEventTcpReader reader = - new ProtocolEventTcpReader( - (eventType, bbuf) -> { - EventImpl reportedEvent = null; - switch (eventType) { - case CONTROL: - reportedEvent = new ControlEventImpl(bbuf); - break; - case PUSH: - reportedEvent = new PushEventImpl(bbuf); - break; - case ACK: - reportedEvent = new AckEventImpl(bbuf); - break; - default: - logger.error("Unknown event type: {}", eventType); - fail(); - break; - } - try { - eventQueue.put(reportedEvent); - } catch (InterruptedException e) { - logger.error("Interrupted: ", e); - Thread.currentThread().interrupt(); - } - }); + // 1. Build PUSH events with message properties and ACK events with the + // same GUIDs; + // 2. Feed ProtocolEventTcpReader with the stream by chunks of different size; + // 3. Decode BlazingMQ events from the ProtocolEventTcpReader callback; + // 4. Verify each PUSH message keeps its properties and payload, and has + // an ACK message with the same GUID. - ReadCompletionStatus status = new ReadCompletionStatus(); - while (fis.available() > 0) { - String[] ss = br.readLine().split(" "); - assertEquals(2, ss.length); - int sz = Integer.parseInt(ss[1]); - byte[] ar = new byte[sz]; - assertEquals(ar.length, fis.read(ar)); - reader.read(status, new ByteBuffer[] {ByteBuffer.wrap(ar)}); + final int NUM_EVENTS = 5; + final int NUM_MESSAGES = 20; + + ByteBufferOutputStream bbos = new ByteBufferOutputStream(); + + for (int i = 0; i < NUM_EVENTS; i++) { + PushEventBuilder pushBuilder = new PushEventBuilder(); + AckEventBuilder ackBuilder = new AckEventBuilder(); + + for (int j = 0; j < NUM_MESSAGES; j++) { + final MessageGUID guid = + MessageGUID.fromHex(String.format("%032X", i * NUM_MESSAGES + j + 1)); + + assertEquals( + EventBuilderResult.SUCCESS, + pushBuilder.packMessage(createPushMessage(guid))); + assertEquals( + EventBuilderResult.SUCCESS, + ackBuilder.packMessage( + new AckMessageImpl( + AckResult.SUCCESS, + CorrelationIdImpl.restoreId(j), + guid, + QUEUE_ID))); + } + + for (ByteBuffer b : pushBuilder.build()) { + bbos.writeBytes(b); + } + for (ByteBuffer b : ackBuilder.build()) { + bbos.writeBytes(b); + } } - Semaphore evSema = new Semaphore(0); - HashSet ackGuids = new HashSet<>(); - ArrayList pushMsgs = new ArrayList<>(); - - Runnable task = - () -> { - int evNum = 0; - while (true) { - EventImpl ev; - try { - ev = eventQueue.poll(1, TimeUnit.SECONDS); - } catch (InterruptedException e) { - logger.error("Interrupted: ", e); - Thread.currentThread().interrupt(); - break; - } - if (ev == null) { - evSema.release(); - break; - } - evNum++; - - try { - if (ev instanceof PushEventImpl) { - PushEventImpl pev = (PushEventImpl) ev; - PushMessageIterator it = pev.iterator(); - while (it.hasNext()) { - PushMessageImpl pm = it.next(); - pushMsgs.add(pm); - } - } else if (ev instanceof AckEventImpl) { - AckEventImpl aev = (AckEventImpl) ev; - Iterator it = aev.iterator(); - while (it.hasNext()) { - AckMessageImpl msg = it.next(); - ackGuids.add(msg.messageGUID().toString()); + + final byte[] stream; + try (ByteBufferInputStream bbis = new ByteBufferInputStream(bbos.reset())) { + stream = new byte[bbis.available()]; + assertEquals(stream.length, bbis.read(stream)); + } + + for (int chunkSize : new int[] {1, 13, 512, stream.length}) { + logger.info("Read {} bytes by chunks of {} bytes", stream.length, chunkSize); + + final ArrayList pushMsgs = new ArrayList<>(); + final HashSet ackGuids = new HashSet<>(); + + ProtocolEventTcpReader reader = + new ProtocolEventTcpReader( + (eventType, bbuf) -> { + switch (eventType) { + case PUSH: + PushMessageIterator pushIt = + new PushEventImpl(bbuf).iterator(); + while (pushIt.hasNext()) { + pushMsgs.add(pushIt.next()); + } + break; + case ACK: + Iterator ackIt = + new AckEventImpl(bbuf).iterator(); + while (ackIt.hasNext()) { + ackGuids.add(ackIt.next().messageGUID().toString()); + } + break; + default: + logger.error("Unexpected event type: {}", eventType); + fail(); + break; } - } - } catch (Exception e) { - logger.error("Exception while processing the event: ", e); - evSema.release(); - break; - } - } - logger.info("Number of events: {}", evNum); - }; - - new Thread(task).start(); - - TestHelpers.acquireSema(evSema, 15); - assertEquals(NUM_PUSH_MESSAGES, pushMsgs.size()); - for (PushMessageImpl msg : pushMsgs) { - String guid = msg.messageGUID().toString(); - assertTrue(ackGuids.contains(guid)); + }); + + ReadCompletionStatus status = new ReadCompletionStatus(); + for (int pos = 0; pos < stream.length; pos += chunkSize) { + final int size = Math.min(chunkSize, stream.length - pos); + byte[] chunk = Arrays.copyOfRange(stream, pos, pos + size); + reader.read(status, new ByteBuffer[] {ByteBuffer.wrap(chunk)}); + } + + assertEquals(NUM_EVENTS * NUM_MESSAGES, pushMsgs.size()); + + for (PushMessageImpl msg : pushMsgs) { + assertTrue(ackGuids.contains(msg.messageGUID().toString())); + + MessagePropertiesImpl props = msg.appData().properties(); + assertEquals(2, props.numProperties()); + assertEquals(ROUTING_ID, props.get("routingId").getValueAsString()); + assertEquals(TIMESTAMP, props.get("timestamp").getValueAsInt64()); + + assertArrayEquals( + new ByteBuffer[] {ByteBuffer.wrap(PAYLOAD.getBytes())}, + msg.appData().payload()); + } } - logger.info("======================================================"); - logger.info("END Testing ProtocolEventImplTcpReaderTest testIODump."); - logger.info("======================================================"); + logger.info("============================================================="); + logger.info("END Testing ProtocolEventImplTcpReaderTest PUSH and ACK stream."); + logger.info("============================================================="); } @Test @@ -206,55 +212,53 @@ void testPartialReading() throws IOException { final int NUM_MESSAGES = 3; - for (boolean isOldStyleProperties : new boolean[] {true, false}) { - // 1. Generate BlazingMQ EventImpl with several PUSH messages; - ByteBuffer[] event = buildPushMessage(isOldStyleProperties); - ByteBufferInputStream inpStream = new ByteBufferInputStream(event); - ReadCompletionStatus status = new ReadCompletionStatus(); + // 1. Generate BlazingMQ EventImpl with several PUSH messages; + ByteBuffer[] event = buildPushMessage(); + ByteBufferInputStream inpStream = new ByteBufferInputStream(event); + ReadCompletionStatus status = new ReadCompletionStatus(); - ArrayList dataList = new ArrayList<>(); + ArrayList dataList = new ArrayList<>(); - ProtocolEventTcpReader reader = - new ProtocolEventTcpReader( - (eventType, bbuf) -> { - dataList.add(bbuf); - assertEquals(EventType.PUSH, eventType); - }); - // 2. Fill a plain buffer with the event content; - final int PLAIN_BUF_SIZE = inpStream.available() * NUM_MESSAGES; - ByteBuffer plainBuffer = ByteBuffer.allocate(PLAIN_BUF_SIZE); - for (int i = 0; i < NUM_MESSAGES; i++) { - for (ByteBuffer b : event) { - b.rewind(); - plainBuffer.put(b); - } + ProtocolEventTcpReader reader = + new ProtocolEventTcpReader( + (eventType, bbuf) -> { + dataList.add(bbuf); + assertEquals(EventType.PUSH, eventType); + }); + // 2. Fill a plain buffer with the event content; + final int PLAIN_BUF_SIZE = inpStream.available() * NUM_MESSAGES; + ByteBuffer plainBuffer = ByteBuffer.allocate(PLAIN_BUF_SIZE); + for (int i = 0; i < NUM_MESSAGES; i++) { + for (ByteBuffer b : event) { + b.rewind(); + plainBuffer.put(b); } - plainBuffer.rewind(); + } + plainBuffer.rewind(); - // 3. Read from this buffer by portions with different size (from 1 - // up to the whole buffer) and feed ProtocolEventTcpReader with those portions; - for (int i = 1; i <= PLAIN_BUF_SIZE; i++) { - ArrayList payloads = new ArrayList<>(); - while (plainBuffer.hasRemaining()) { - int sz = Math.min(i, plainBuffer.remaining()); - byte[] ar = new byte[sz]; - plainBuffer.get(ar); - payloads.add(ByteBuffer.wrap(ar)); - } - ByteBuffer[] bb = new ByteBuffer[payloads.size()]; - bb = payloads.toArray(bb); - reader.read(status, bb); - plainBuffer.rewind(); + // 3. Read from this buffer by portions with different size (from 1 + // up to the whole buffer) and feed ProtocolEventTcpReader with those portions; + for (int i = 1; i <= PLAIN_BUF_SIZE; i++) { + ArrayList payloads = new ArrayList<>(); + while (plainBuffer.hasRemaining()) { + int sz = Math.min(i, plainBuffer.remaining()); + byte[] ar = new byte[sz]; + plainBuffer.get(ar); + payloads.add(ByteBuffer.wrap(ar)); } - // 4. Check that ProtocolEventTcpReader correctly composes BlazingMQ Events. - assertEquals(dataList.size(), PLAIN_BUF_SIZE * NUM_MESSAGES); - for (ByteBuffer[] data : dataList) { - inpStream.reset(); - ByteBufferInputStream istr = new ByteBufferInputStream(data); - assertEquals(istr.available(), inpStream.available()); - while (istr.available() > 0) { - assertEquals(istr.readByte(), inpStream.readByte()); - } + ByteBuffer[] bb = new ByteBuffer[payloads.size()]; + bb = payloads.toArray(bb); + reader.read(status, bb); + plainBuffer.rewind(); + } + // 4. Check that ProtocolEventTcpReader correctly composes BlazingMQ Events. + assertEquals(dataList.size(), PLAIN_BUF_SIZE * NUM_MESSAGES); + for (ByteBuffer[] data : dataList) { + inpStream.reset(); + ByteBufferInputStream istr = new ByteBufferInputStream(data); + assertEquals(istr.available(), inpStream.available()); + while (istr.available() > 0) { + assertEquals(istr.readByte(), inpStream.readByte()); } } } diff --git a/bmq-sdk/src/test/java/com/bloomberg/bmq/impl/PutPosterTest.java b/bmq-sdk/src/test/java/com/bloomberg/bmq/impl/PutPosterTest.java index 918eb118..5d0832ed 100644 --- a/bmq-sdk/src/test/java/com/bloomberg/bmq/impl/PutPosterTest.java +++ b/bmq-sdk/src/test/java/com/bloomberg/bmq/impl/PutPosterTest.java @@ -161,218 +161,200 @@ void testPostFailed() throws IOException { @Test void testPostValidMessages() throws IOException { - for (boolean isOldStyleProperties : new boolean[] {false, true}) { - BrokerConnection mockedConnection = mock(BrokerConnection.class); - when(mockedConnection.isOldStyleMessageProperties()).thenReturn(isOldStyleProperties); - when(mockedConnection.write(any(ByteBuffer[].class), anyBoolean())) - .thenReturn(GenericResult.SUCCESS); - - EventsStats eventsStats = new EventsStats(); - PutPoster poster = new PutPoster(mockedConnection, eventsStats); - - final MessagePropertiesImpl props = new MessagePropertiesImpl(); - props.setPropertyAsInt32("id", 3); - props.setPropertyAsBinary("data", new byte[] {1, 2, 3, 4, 5}); - - PutMessageImpl bigMsg1 = new PutMessageImpl(); - bigMsg1.appData().setPayload(ByteBuffer.allocate(PutHeader.MAX_PAYLOAD_SIZE_SOFT)); - bigMsg1.setCompressionType(CompressionAlgorithmType.E_NONE); - - PutMessageImpl smallMsg1 = new PutMessageImpl(); - smallMsg1.appData().setPayload(ByteBuffer.allocate(10000)); - smallMsg1.appData().setProperties(props); - - PutMessageImpl bigMsg2 = new PutMessageImpl(); - bigMsg2.appData().setPayload(ByteBuffer.allocate(PutHeader.MAX_PAYLOAD_SIZE_SOFT)); - bigMsg2.setCompressionType(CompressionAlgorithmType.E_NONE); - - PutMessageImpl smallMsg2 = new PutMessageImpl(); - smallMsg2.appData().setProperties(props); - smallMsg2.appData().setPayload(ByteBuffer.allocate(10001)); - - PutMessageImpl compressedMsg = new PutMessageImpl(); - compressedMsg - .appData() - .setPayload(ByteBuffer.allocate(PutHeader.MAX_PAYLOAD_SIZE_SOFT)); - compressedMsg.appData().setProperties(props); - compressedMsg.setCompressionType(CompressionAlgorithmType.E_ZLIB); - - poster.post(bigMsg1, smallMsg1, bigMsg2, smallMsg2, compressedMsg); - - assertEquals(isOldStyleProperties, bigMsg1.appData().isOldStyleProperties()); - assertEquals(isOldStyleProperties, smallMsg1.appData().isOldStyleProperties()); - assertEquals(isOldStyleProperties, bigMsg2.appData().isOldStyleProperties()); - assertEquals(isOldStyleProperties, smallMsg2.appData().isOldStyleProperties()); - assertEquals(isOldStyleProperties, compressedMsg.appData().isOldStyleProperties()); - - assertEquals(0, bigMsg1.header().schemaWireId()); - assertEquals(isOldStyleProperties ? 0 : 1, smallMsg1.header().schemaWireId()); - assertEquals(0, bigMsg2.header().schemaWireId()); - assertEquals(isOldStyleProperties ? 0 : 1, smallMsg2.header().schemaWireId()); - assertEquals(isOldStyleProperties ? 0 : 1, compressedMsg.header().schemaWireId()); - - // Build data to check - PutEventBuilder builder = new PutEventBuilder(); - EventsStats expectedStats = new EventsStats(); - - builder.packMessage(bigMsg1, isOldStyleProperties); - builder.packMessage(smallMsg1, isOldStyleProperties); - expectedStats.onEvent(EventType.PUT, builder.eventLength(), builder.messageCount()); - ByteBuffer[] data1 = builder.build(); - - builder.reset(); - builder.packMessage(bigMsg2, isOldStyleProperties); - builder.packMessage(smallMsg2, isOldStyleProperties); - builder.packMessage(compressedMsg, isOldStyleProperties); - expectedStats.onEvent(EventType.PUT, builder.eventLength(), builder.messageCount()); - ByteBuffer[] data2 = builder.build(); - - // write method can be verified using two lines below, - // but for clarity we at first check number of invocations and - // after that we check arguments - // verify(mockedConnection, times(1)).write(data1, true); - // verify(mockedConnection, times(1)).write(data2, true); - - // Special classes to capture arguments passed to the write method - ArgumentCaptor bbCaptor = ArgumentCaptor.forClass(ByteBuffer[].class); - ArgumentCaptor boolCaptor = ArgumentCaptor.forClass(Boolean.class); - - // Verify that the write method has been called twice - verify(mockedConnection, times(2)).write(bbCaptor.capture(), boolCaptor.capture()); - - // Get captured arguments - List allData = bbCaptor.getAllValues(); - List allBooleans = boolCaptor.getAllValues(); - - // Verify first argument - assertArrayEquals(data1, allData.get(0)); - assertArrayEquals(data2, allData.get(1)); - - // Verify second argument - assertTrue(allBooleans.get(0)); - assertTrue(allBooleans.get(1)); - - StringBuilder expectedBuilder = new StringBuilder(); - EventsStatsTest.dump(expectedStats, expectedBuilder, false); - String expectedStr = expectedBuilder.toString(); - logger.info("Expected stats:\n{}", expectedStr); - - StringBuilder actualBuilder = new StringBuilder(); - EventsStatsTest.dump(eventsStats, actualBuilder, false); - String actualStr = actualBuilder.toString(); - logger.info("Actual stats:\n{}", actualStr); - - assertEquals(expectedStr, actualStr); - } + BrokerConnection mockedConnection = mock(BrokerConnection.class); + when(mockedConnection.write(any(ByteBuffer[].class), anyBoolean())) + .thenReturn(GenericResult.SUCCESS); + + EventsStats eventsStats = new EventsStats(); + PutPoster poster = new PutPoster(mockedConnection, eventsStats); + + final MessagePropertiesImpl props = new MessagePropertiesImpl(); + props.setPropertyAsInt32("id", 3); + props.setPropertyAsBinary("data", new byte[] {1, 2, 3, 4, 5}); + + PutMessageImpl bigMsg1 = new PutMessageImpl(); + bigMsg1.appData().setPayload(ByteBuffer.allocate(PutHeader.MAX_PAYLOAD_SIZE_SOFT)); + bigMsg1.setCompressionType(CompressionAlgorithmType.E_NONE); + + PutMessageImpl smallMsg1 = new PutMessageImpl(); + smallMsg1.appData().setPayload(ByteBuffer.allocate(10000)); + smallMsg1.appData().setProperties(props); + + PutMessageImpl bigMsg2 = new PutMessageImpl(); + bigMsg2.appData().setPayload(ByteBuffer.allocate(PutHeader.MAX_PAYLOAD_SIZE_SOFT)); + bigMsg2.setCompressionType(CompressionAlgorithmType.E_NONE); + + PutMessageImpl smallMsg2 = new PutMessageImpl(); + smallMsg2.appData().setProperties(props); + smallMsg2.appData().setPayload(ByteBuffer.allocate(10001)); + + PutMessageImpl compressedMsg = new PutMessageImpl(); + compressedMsg.appData().setPayload(ByteBuffer.allocate(PutHeader.MAX_PAYLOAD_SIZE_SOFT)); + compressedMsg.appData().setProperties(props); + compressedMsg.setCompressionType(CompressionAlgorithmType.E_ZLIB); + + poster.post(bigMsg1, smallMsg1, bigMsg2, smallMsg2, compressedMsg); + + assertEquals(0, bigMsg1.header().schemaWireId()); + assertEquals(1, smallMsg1.header().schemaWireId()); + assertEquals(0, bigMsg2.header().schemaWireId()); + assertEquals(1, smallMsg2.header().schemaWireId()); + assertEquals(1, compressedMsg.header().schemaWireId()); + + // Build data to check + PutEventBuilder builder = new PutEventBuilder(); + EventsStats expectedStats = new EventsStats(); + + builder.packMessage(bigMsg1); + builder.packMessage(smallMsg1); + expectedStats.onEvent(EventType.PUT, builder.eventLength(), builder.messageCount()); + ByteBuffer[] data1 = builder.build(); + + builder.reset(); + builder.packMessage(bigMsg2); + builder.packMessage(smallMsg2); + builder.packMessage(compressedMsg); + expectedStats.onEvent(EventType.PUT, builder.eventLength(), builder.messageCount()); + ByteBuffer[] data2 = builder.build(); + + // write method can be verified using two lines below, + // but for clarity we at first check number of invocations and + // after that we check arguments + // verify(mockedConnection, times(1)).write(data1, true); + // verify(mockedConnection, times(1)).write(data2, true); + + // Special classes to capture arguments passed to the write method + ArgumentCaptor bbCaptor = ArgumentCaptor.forClass(ByteBuffer[].class); + ArgumentCaptor boolCaptor = ArgumentCaptor.forClass(Boolean.class); + + // Verify that the write method has been called twice + verify(mockedConnection, times(2)).write(bbCaptor.capture(), boolCaptor.capture()); + + // Get captured arguments + List allData = bbCaptor.getAllValues(); + List allBooleans = boolCaptor.getAllValues(); + + // Verify first argument + assertArrayEquals(data1, allData.get(0)); + assertArrayEquals(data2, allData.get(1)); + + // Verify second argument + assertTrue(allBooleans.get(0)); + assertTrue(allBooleans.get(1)); + + StringBuilder expectedBuilder = new StringBuilder(); + EventsStatsTest.dump(expectedStats, expectedBuilder, false); + String expectedStr = expectedBuilder.toString(); + logger.info("Expected stats:\n{}", expectedStr); + + StringBuilder actualBuilder = new StringBuilder(); + EventsStatsTest.dump(eventsStats, actualBuilder, false); + String actualStr = actualBuilder.toString(); + logger.info("Actual stats:\n{}", actualStr); + + assertEquals(expectedStr, actualStr); } @Test void testPostNoInfiniteLoop() throws IOException, TimeoutException, InterruptedException { - for (boolean isOldStyleProperties : new boolean[] {false, true}) { - BrokerConnection connection = mock(BrokerConnection.class); - when(connection.isOldStyleMessageProperties()).thenReturn(isOldStyleProperties); - when(connection.write(any(ByteBuffer[].class), anyBoolean())) - .thenReturn(GenericResult.SUCCESS); + BrokerConnection connection = mock(BrokerConnection.class); + when(connection.write(any(ByteBuffer[].class), anyBoolean())) + .thenReturn(GenericResult.SUCCESS); - PutPoster poster = new PutPoster(connection, new EventsStats()); + PutPoster poster = new PutPoster(connection, new EventsStats()); - // Update max event size - final int MAX_EVENT_SIZE = 1024; + // Update max event size + final int MAX_EVENT_SIZE = 1024; - poster.setMaxEventSize(MAX_EVENT_SIZE); + poster.setMaxEventSize(MAX_EVENT_SIZE); - // Create a msg with payload = max event size - PutMessageImpl msg1 = new PutMessageImpl(); - msg1.appData().setPayload(ByteBuffer.allocate(MAX_EVENT_SIZE)); + // Create a msg with payload = max event size + PutMessageImpl msg1 = new PutMessageImpl(); + msg1.appData().setPayload(ByteBuffer.allocate(MAX_EVENT_SIZE)); - // Post async - ExecutorService es = Executors.newSingleThreadExecutor(); - Future f = es.submit(() -> poster.post(msg1)); + // Post async + ExecutorService es = Executors.newSingleThreadExecutor(); + Future f = es.submit(() -> poster.post(msg1)); - // Get the result - try { - f.get(1, TimeUnit.SECONDS); - } catch (ExecutionException e) { - Throwable cause = e.getCause(); + // Get the result + try { + f.get(1, TimeUnit.SECONDS); + } catch (ExecutionException e) { + Throwable cause = e.getCause(); - assertNotNull(cause); - assertEquals("Failed to build PUT event: PAYLOAD_TOO_BIG", cause.getMessage()); - } + assertNotNull(cause); + assertEquals("Failed to build PUT event: PAYLOAD_TOO_BIG", cause.getMessage()); + } - es.shutdownNow(); + es.shutdownNow(); - // Post a msg with payload = max payload size - PutMessageImpl msg2 = new PutMessageImpl(); - msg2.appData() - .setPayload(ByteBuffer.allocate(MAX_EVENT_SIZE - PutHeader.HEADER_SIZE - 4)); + // Post a msg with payload = max payload size + PutMessageImpl msg2 = new PutMessageImpl(); + msg2.appData().setPayload(ByteBuffer.allocate(MAX_EVENT_SIZE - PutHeader.HEADER_SIZE - 4)); - poster.post(msg2); - } + poster.post(msg2); } @Test void testRegisterAck() throws Exception { - for (boolean isOldStyleProperties : new boolean[] {false, true}) { - BrokerConnection connection = mock(BrokerConnection.class); - when(connection.isOldStyleMessageProperties()).thenReturn(isOldStyleProperties); - when(connection.write(any(ByteBuffer[].class), anyBoolean())) - .thenReturn(GenericResult.SUCCESS); - - PutPoster poster = new PutPoster(connection, new EventsStats()); - - // try to register null ACK message - try { - poster.registerAck(null); - fail(); // Should not get here - } catch (IllegalArgumentException e) { - assertEquals("'ackMsg' must be non-null", e.getMessage()); - } - - // Register ACK with null correlation Id - // When AckMessageImpl is being streamed in, its `correlationId() - // is initialized to some value by creating CorrelationImpl instance. - // Here we use 'restoreId' method to create new instance of - // CorrelationIdImpl with zero id to ensure its reference differs - // from CorrelationIdImpl.NULL_CORRELATION_ID object. - AckMessageImpl ackMsg = - new AckMessageImpl( - AckResult.UNKNOWN, - CorrelationIdImpl.restoreId(0), - MessageGUID.createEmptyGUID(), - 0); - poster.registerAck(ackMsg); // should be just logged and ignored - - // Post PUT message and then register ACK message - Object userData = new Object(); - CorrelationIdImpl cId = CorrelationIdImpl.nextId(userData); - - PutMessageImpl msg = new PutMessageImpl(); - msg.appData().setPayload(ByteBuffer.allocate(10)); - msg.setupCorrelationId(cId); + BrokerConnection connection = mock(BrokerConnection.class); + when(connection.write(any(ByteBuffer[].class), anyBoolean())) + .thenReturn(GenericResult.SUCCESS); - poster.post(msg); + PutPoster poster = new PutPoster(connection, new EventsStats()); - ackMsg = - new AckMessageImpl( - AckResult.SUCCESS, - CorrelationIdImpl.restoreId(cId.toInt()), - MessageGUID.createEmptyGUID(), - 0); - - poster.registerAck(ackMsg); - - assertEquals(cId, ackMsg.correlationId()); - assertEquals(userData, ackMsg.correlationId().userData()); - - // Try to register the same ACK message again - ackMsg = - new AckMessageImpl( - AckResult.SUCCESS, - CorrelationIdImpl.restoreId(cId.toInt()), - ackMsg.messageGUID(), - 0); - poster.registerAck(ackMsg); // should be just logged and ignored + // try to register null ACK message + try { + poster.registerAck(null); + fail(); // Should not get here + } catch (IllegalArgumentException e) { + assertEquals("'ackMsg' must be non-null", e.getMessage()); } + + // Register ACK with null correlation Id + // When AckMessageImpl is being streamed in, its `correlationId() + // is initialized to some value by creating CorrelationImpl instance. + // Here we use 'restoreId' method to create new instance of + // CorrelationIdImpl with zero id to ensure its reference differs + // from CorrelationIdImpl.NULL_CORRELATION_ID object. + AckMessageImpl ackMsg = + new AckMessageImpl( + AckResult.UNKNOWN, + CorrelationIdImpl.restoreId(0), + MessageGUID.createEmptyGUID(), + 0); + poster.registerAck(ackMsg); // should be just logged and ignored + + // Post PUT message and then register ACK message + Object userData = new Object(); + CorrelationIdImpl cId = CorrelationIdImpl.nextId(userData); + + PutMessageImpl msg = new PutMessageImpl(); + msg.appData().setPayload(ByteBuffer.allocate(10)); + msg.setupCorrelationId(cId); + + poster.post(msg); + + ackMsg = + new AckMessageImpl( + AckResult.SUCCESS, + CorrelationIdImpl.restoreId(cId.toInt()), + MessageGUID.createEmptyGUID(), + 0); + + poster.registerAck(ackMsg); + + assertEquals(cId, ackMsg.correlationId()); + assertEquals(userData, ackMsg.correlationId().userData()); + + // Try to register the same ACK message again + ackMsg = + new AckMessageImpl( + AckResult.SUCCESS, + CorrelationIdImpl.restoreId(cId.toInt()), + ackMsg.messageGUID(), + 0); + poster.registerAck(ackMsg); // should be just logged and ignored } } diff --git a/bmq-sdk/src/test/java/com/bloomberg/bmq/impl/infr/msg/MessagesTestSamples.java b/bmq-sdk/src/test/java/com/bloomberg/bmq/impl/infr/msg/MessagesTestSamples.java index 792371e1..d325a0b7 100644 --- a/bmq-sdk/src/test/java/com/bloomberg/bmq/impl/infr/msg/MessagesTestSamples.java +++ b/bmq-sdk/src/test/java/com/bloomberg/bmq/impl/infr/msg/MessagesTestSamples.java @@ -39,11 +39,6 @@ public SampleFileMetadata setLength(int val) { } } - public static final SampleFileMetadata BMQ_IO_DUMP_BIN = - new SampleFileMetadata("/data/bmq_io_dump_1551267643131.bin", 16520); - public static final SampleFileMetadata BMQ_IO_DUMP_IDX = - new SampleFileMetadata("/data/bmq_io_dump_1551267643131.idx", 2178); - public static final SampleFileMetadata STATUS_MSG = new SampleFileMetadata( "/data/msg_control_status_53121b03-f45d-46b2-95d0-f2df8a1a2cb2.bin", 36); @@ -95,8 +90,6 @@ public SampleFileMetadata setLength(int val) { new SampleFileMetadata("/data/msg_put_zlib_27042018.bin", 132); public static final SampleFileMetadata PUT_MULTI_MSG = new SampleFileMetadata("/data/msg_put_multi.bin", 264); - public static final SampleFileMetadata MSG_PROPS_OLD = - new SampleFileMetadata("/data/msg_props_old.bin", 64); public static final SampleFileMetadata MSG_PROPS = new SampleFileMetadata("/data/msg_props.bin", 64); public static final SampleFileMetadata MSG_PROPS_LONG_HEADERS = diff --git a/bmq-sdk/src/test/java/com/bloomberg/bmq/impl/infr/proto/ApplicationDataTest.java b/bmq-sdk/src/test/java/com/bloomberg/bmq/impl/infr/proto/ApplicationDataTest.java index b713eeec..4796d81b 100644 --- a/bmq-sdk/src/test/java/com/bloomberg/bmq/impl/infr/proto/ApplicationDataTest.java +++ b/bmq-sdk/src/test/java/com/bloomberg/bmq/impl/infr/proto/ApplicationDataTest.java @@ -57,7 +57,6 @@ void testEmpty() throws IOException { assertFalse(appData.hasProperties()); assertFalse(appData.isCompressed()); - assertFalse(appData.isOldStyleProperties()); } @Test @@ -68,17 +67,14 @@ void testStreamOut() throws IOException { new MessagePropertiesImpl[] { null, new MessagePropertiesImpl(), generateProps() }) - for (boolean isOldStyleProperties : new boolean[] {false, true}) - for (CompressionAlgorithmType compressionType : - CompressionAlgorithmType.values()) { - logger.info( - "Stream out with payload:{}, props:{}, oldStyleProperties:{}, compression: {}", - payload, - props, - isOldStyleProperties, - compressionType); - verifyStreamOut(payload, props, isOldStyleProperties, compressionType); - } + for (CompressionAlgorithmType compressionType : CompressionAlgorithmType.values()) { + logger.info( + "Stream out with payload:{}, props:{}, compression: {}", + payload, + props, + compressionType); + verifyStreamOut(payload, props, compressionType); + } } @Test @@ -89,17 +85,14 @@ void testStreamIn() throws IOException { new MessagePropertiesImpl[] { null, new MessagePropertiesImpl(), generateProps() }) - for (boolean isOldStyleProperties : new boolean[] {false, true}) - for (CompressionAlgorithmType compressionType : - CompressionAlgorithmType.values()) { - logger.info( - "Stream in with payload:{}, props:{}, oldStyleProperties:{}, compression: {}", - payload, - props, - isOldStyleProperties, - compressionType); - verifyStreamIn(payload, props, isOldStyleProperties, compressionType); - } + for (CompressionAlgorithmType compressionType : CompressionAlgorithmType.values()) { + logger.info( + "Stream in with payload:{}, props:{}, compression: {}", + payload, + props, + compressionType); + verifyStreamIn(payload, props, compressionType); + } } @Test @@ -119,7 +112,7 @@ void testStreamInInvalidCompression() throws IOException { ApplicationData data = new ApplicationData(); // Stream in uncompressed data as compressed - data.streamIn(bbis.available(), false, false, CompressionAlgorithmType.E_ZLIB, bbis); + data.streamIn(bbis.available(), false, CompressionAlgorithmType.E_ZLIB, bbis); // "Compressed" data should be buffered assertTrue(data.isCompressed()); @@ -156,7 +149,6 @@ public MessagePropertiesImpl generateProps() { private ByteBuffer[] generateOutput( ByteBuffer[] payload, MessagePropertiesImpl props, - boolean isOldStyleProperties, CompressionAlgorithmType compressionType) throws IOException { ByteBufferOutputStream bbos = new ByteBufferOutputStream(); @@ -165,11 +157,7 @@ private ByteBuffer[] generateOutput( if (compressionType == CompressionAlgorithmType.E_NONE) { if (hasProperties) { - if (isOldStyleProperties) { - props.streamOutOld(bbos); - } else { - props.streamOut(bbos); - } + props.streamOut(bbos); } if (payload != null) { @@ -178,8 +166,8 @@ private ByteBuffer[] generateOutput( } } } else { - // Stream out if new style properties - if (hasProperties && !isOldStyleProperties) { + // Properties are never compressed + if (hasProperties) { props.streamOut(bbos); } // We need to close compressed stream in order to flush all compressed bytes @@ -195,11 +183,6 @@ private ByteBuffer[] generateOutput( try (OutputStream compressedStream = compressionType.getCompression().compress(bbos); DataOutputStream compressedOutput = new DataOutputStream(compressedStream)) { - // Stream out if old style properties - if (hasProperties && isOldStyleProperties) { - props.streamOutOld(compressedOutput); - } - if (payload != null) { for (ByteBuffer b : payload) { byte[] bytes = new byte[b.remaining()]; @@ -219,7 +202,6 @@ private ByteBuffer[] generateOutput( public void verifyStreamOut( ByteBuffer[] payload, MessagePropertiesImpl props, - boolean isOldStyleProperties, CompressionAlgorithmType compressionType) throws IOException { @@ -261,10 +243,8 @@ public void verifyStreamOut( // Verify unpackedSize assertEquals(getSize(payload) + propsSize, appData.unpackedSize()); - ByteBuffer[] expected = - generateOutput(duplicate(payload), props, isOldStyleProperties, compressionType); + ByteBuffer[] expected = generateOutput(duplicate(payload), props, compressionType); - appData.setIsOldStyleProperties(isOldStyleProperties); appData.compressData(compressionType); int numPaddingBytes = ProtocolUtil.calculatePadding(appData.unpackedSize()); @@ -300,8 +280,7 @@ public void verifyStreamOut( ByteBufferInputStream bbis = new ByteBufferInputStream(expected); final int unpackedInputSize = bbis.available() - numPaddingBytes; - appData.streamIn( - getSize(expected), hasProperties, isOldStyleProperties, compressionType, bbis); + appData.streamIn(getSize(expected), hasProperties, compressionType, bbis); assertEquals(unpackedInputSize, appData.unpackedSize()); verifyPayload(payload, appData.payload()); @@ -311,13 +290,11 @@ public void verifyStreamOut( public void verifyStreamIn( ByteBuffer[] payload, MessagePropertiesImpl props, - boolean isOldStyleProperties, CompressionAlgorithmType compressionType) throws IOException { // Prepare data to stream in - ByteBuffer[] data = - generateOutput(duplicate(payload), props, isOldStyleProperties, compressionType); + ByteBuffer[] data = generateOutput(duplicate(payload), props, compressionType); ByteBufferInputStream bbis = new ByteBufferInputStream(duplicate(data)); // Get num of padding bytes @@ -330,9 +307,7 @@ public void verifyStreamIn( final boolean hasProperties = props != null && props.numProperties() > 0; - appData.streamIn(size, hasProperties, isOldStyleProperties, compressionType, bbis); - - assertEquals(isOldStyleProperties, appData.isOldStyleProperties()); + appData.streamIn(size, hasProperties, compressionType, bbis); // Check unpacked input unpackedSize has been set int unpackedSize = size - numPaddingBytes; @@ -349,9 +324,8 @@ public void verifyStreamIn( // Check properties verifyProperties(hasProperties ? props : null, appData.properties()); - // If data is compressed and properties are new style encoded, data - // should stay compressed - if (compressionType != CompressionAlgorithmType.E_NONE && !isOldStyleProperties) { + // Properties are not compressed, so data should stay compressed + if (compressionType != CompressionAlgorithmType.E_NONE) { assertTrue(appData.isCompressed()); } @@ -372,7 +346,6 @@ public void verifyStreamIn( ByteBufferOutputStream bbos = new ByteBufferOutputStream(); - appData.setIsOldStyleProperties(isOldStyleProperties); appData.compressData(compressionType); appData.streamOut(bbos); diff --git a/bmq-sdk/src/test/java/com/bloomberg/bmq/impl/infr/proto/MessagePropertiesHeaderTest.java b/bmq-sdk/src/test/java/com/bloomberg/bmq/impl/infr/proto/MessagePropertiesHeaderTest.java index 36e3060c..ee25e3bf 100644 --- a/bmq-sdk/src/test/java/com/bloomberg/bmq/impl/infr/proto/MessagePropertiesHeaderTest.java +++ b/bmq-sdk/src/test/java/com/bloomberg/bmq/impl/infr/proto/MessagePropertiesHeaderTest.java @@ -90,17 +90,17 @@ class TestData { MessagePropertyHeader ph = msgPropHeaders[0]; assertEquals(PropertyType.INT32.toInt(), ph.propertyType()); - assertEquals(0, ph.propertyValueLength()); // offset + assertEquals(0, ph.propertyValueOffset()); assertEquals(8, ph.propertyNameLength()); ph = msgPropHeaders[1]; assertEquals(PropertyType.INT64.toInt(), ph.propertyType()); - assertEquals(12, ph.propertyValueLength()); // offset + assertEquals(12, ph.propertyValueOffset()); assertEquals(9, ph.propertyNameLength()); ph = msgPropHeaders[2]; assertEquals(PropertyType.STRING.toInt(), ph.propertyType()); - assertEquals(29, ph.propertyValueLength()); // offset + assertEquals(29, ph.propertyValueOffset()); assertEquals(2, ph.propertyNameLength()); final int available = diff --git a/bmq-sdk/src/test/java/com/bloomberg/bmq/impl/infr/proto/MessagePropertiesTest.java b/bmq-sdk/src/test/java/com/bloomberg/bmq/impl/infr/proto/MessagePropertiesTest.java index 8194388d..83f68f63 100644 --- a/bmq-sdk/src/test/java/com/bloomberg/bmq/impl/infr/proto/MessagePropertiesTest.java +++ b/bmq-sdk/src/test/java/com/bloomberg/bmq/impl/infr/proto/MessagePropertiesTest.java @@ -41,62 +41,28 @@ class MessagePropertiesTest { void testWithPattern() throws IOException { class TestData { final MessagesTestSamples.SampleFileMetadata sampleFile; - final boolean streamInOld; - final boolean streamOutOld; final MessagesTestSamples.SampleFileMetadata sampleFileCompare; TestData( MessagesTestSamples.SampleFileMetadata sampleFile, - boolean streamInOld, - boolean streamOutOld, MessagesTestSamples.SampleFileMetadata sampleFileCompare) { this.sampleFile = sampleFile; - this.streamInOld = streamInOld; - this.streamOutOld = streamOutOld; this.sampleFileCompare = sampleFileCompare; } } final TestData[] data = new TestData[] { - new TestData( - MessagesTestSamples.MSG_PROPS_OLD, - true, - true, - MessagesTestSamples.MSG_PROPS_OLD), - new TestData( - MessagesTestSamples.MSG_PROPS_OLD, - true, - false, - MessagesTestSamples.MSG_PROPS), - new TestData( - MessagesTestSamples.MSG_PROPS, - false, - false, - MessagesTestSamples.MSG_PROPS), - new TestData( - MessagesTestSamples.MSG_PROPS, - false, - true, - MessagesTestSamples.MSG_PROPS_OLD), + new TestData(MessagesTestSamples.MSG_PROPS, MessagesTestSamples.MSG_PROPS), new TestData( MessagesTestSamples.MSG_PROPS_LONG_HEADERS, - false, - false, - MessagesTestSamples.MSG_PROPS), - new TestData( - MessagesTestSamples.MSG_PROPS_LONG_HEADERS, - false, - true, - MessagesTestSamples.MSG_PROPS_OLD) + MessagesTestSamples.MSG_PROPS) }; for (TestData testData : data) { logger.info( - "Sample: {}, stream in old: {}, stream out old: {}, compare: {}", + "Sample: {}, compare: {}", testData.sampleFile.filePath(), - testData.streamInOld, - testData.streamOutOld, testData.sampleFileCompare.filePath()); ByteBuffer buf = TestHelpers.readFile(testData.sampleFile.filePath()); @@ -113,11 +79,7 @@ class TestData { int toRead = bbis.available(); logger.info("Stream in {} bytes", toRead); - if (testData.streamInOld) { - toRead -= props.streamInOld(bbis); - } else { - toRead -= props.streamIn(bbis); - } + toRead -= props.streamIn(bbis); assertEquals(0, toRead); assertEquals(0, bbis.available()); @@ -158,14 +120,10 @@ class TestData { assertTrue(timestampFound); assertTrue(encodingFound); - // Stream out to another format and compare + // Stream out and compare ByteBufferOutputStream bbos = new ByteBufferOutputStream(); - if (testData.streamOutOld) { - props.streamOutOld(bbos); - } else { - props.streamOut(bbos); - } + props.streamOut(bbos); TestHelpers.compareWithFileContent(bbos.reset(), testData.sampleFileCompare); } @@ -173,124 +131,114 @@ class TestData { @Test void testStreamOut() throws IOException { - for (boolean isOldStyleProperties : new boolean[] {false, true}) { - final boolean BOOL_VAL = true; - final byte BYTE_VAL = 2; - final short SHORT_VAL = 12; - final int INT32_VAL = 12345; - final long INT64_VAL = 987654321L; - final String STRING_VAL = "myValue"; - final byte[] BINARY_VAL = "abcdefgh".getBytes(); + final boolean BOOL_VAL = true; + final byte BYTE_VAL = 2; + final short SHORT_VAL = 12; + final int INT32_VAL = 12345; + final long INT64_VAL = 987654321L; + final String STRING_VAL = "myValue"; + final byte[] BINARY_VAL = "abcdefgh".getBytes(); - final int NUM_PROPERTIES = 7; + final int NUM_PROPERTIES = 7; - ByteBufferOutputStream bbos = new ByteBufferOutputStream(); - MessagePropertiesImpl props = new MessagePropertiesImpl(); - assertEquals(0, props.numProperties()); - assertEquals(0, props.totalSize()); + ByteBufferOutputStream bbos = new ByteBufferOutputStream(); + MessagePropertiesImpl props = new MessagePropertiesImpl(); + assertEquals(0, props.numProperties()); + assertEquals(0, props.totalSize()); - for (PropertyType t : PropertyType.values()) { - switch (t) { - case UNDEFINED: // Skip - break; - case BOOL: - props.setPropertyAsBool(PropertyType.BOOL.toString(), BOOL_VAL); - break; - case BYTE: - props.setPropertyAsByte(PropertyType.BYTE.toString(), BYTE_VAL); - break; - case SHORT: - props.setPropertyAsShort(PropertyType.SHORT.toString(), SHORT_VAL); - break; - case INT32: - props.setPropertyAsInt32(PropertyType.INT32.toString(), INT32_VAL); - break; - case INT64: - props.setPropertyAsInt64(PropertyType.INT64.toString(), INT64_VAL); - break; - case STRING: - props.setPropertyAsString(PropertyType.STRING.toString(), STRING_VAL); - break; - case BINARY: - props.setPropertyAsBinary(PropertyType.BINARY.toString(), BINARY_VAL); - break; - default: // Unknown type - fail(); - break; - } + for (PropertyType t : PropertyType.values()) { + switch (t) { + case UNDEFINED: // Skip + break; + case BOOL: + props.setPropertyAsBool(PropertyType.BOOL.toString(), BOOL_VAL); + break; + case BYTE: + props.setPropertyAsByte(PropertyType.BYTE.toString(), BYTE_VAL); + break; + case SHORT: + props.setPropertyAsShort(PropertyType.SHORT.toString(), SHORT_VAL); + break; + case INT32: + props.setPropertyAsInt32(PropertyType.INT32.toString(), INT32_VAL); + break; + case INT64: + props.setPropertyAsInt64(PropertyType.INT64.toString(), INT64_VAL); + break; + case STRING: + props.setPropertyAsString(PropertyType.STRING.toString(), STRING_VAL); + break; + case BINARY: + props.setPropertyAsBinary(PropertyType.BINARY.toString(), BINARY_VAL); + break; + default: // Unknown type + fail(); + break; } + } - assertEquals(NUM_PROPERTIES, props.numProperties()); + assertEquals(NUM_PROPERTIES, props.numProperties()); - if (isOldStyleProperties) { - props.streamOutOld(bbos); - } else { - props.streamOut(bbos); - } + props.streamOut(bbos); - assertTrue(bbos.size() > 0); + assertTrue(bbos.size() > 0); - ByteBufferInputStream bbis = new ByteBufferInputStream(bbos.reset()); - props = new MessagePropertiesImpl(); - assertEquals(0, props.numProperties()); + ByteBufferInputStream bbis = new ByteBufferInputStream(bbos.reset()); + props = new MessagePropertiesImpl(); + assertEquals(0, props.numProperties()); - int toRead = bbis.available(); + int toRead = bbis.available(); - if (isOldStyleProperties) { - toRead -= props.streamInOld(bbis); - } else { - toRead -= props.streamIn(bbis); - } + toRead -= props.streamIn(bbis); - assertEquals(0, toRead); - assertEquals(0, bbis.available()); + assertEquals(0, toRead); + assertEquals(0, bbis.available()); - assertEquals(NUM_PROPERTIES, props.numProperties()); + assertEquals(NUM_PROPERTIES, props.numProperties()); - Iterator> pit = props.iterator(); + Iterator> pit = props.iterator(); - for (PropertyType t : PropertyType.values()) { - if (!PropertyType.isValid(t.toInt())) { - continue; - } - assertTrue(pit.hasNext()); - MessageProperty p = pit.next().getValue(); - String name = null; - switch (t) { - case BOOL: - name = PropertyType.BOOL.toString(); - assertEquals(BOOL_VAL, p.getValueAsBool()); - break; - case BYTE: - name = PropertyType.BYTE.toString(); - assertEquals(BYTE_VAL, p.getValueAsByte()); - break; - case SHORT: - name = PropertyType.SHORT.toString(); - assertEquals(SHORT_VAL, p.getValueAsShort()); - break; - case INT32: - name = PropertyType.INT32.toString(); - assertEquals(INT32_VAL, p.getValueAsInt32()); - break; - case INT64: - name = PropertyType.INT64.toString(); - assertEquals(INT64_VAL, p.getValueAsInt64()); - break; - case STRING: - name = PropertyType.STRING.toString(); - assertEquals(STRING_VAL, p.getValueAsString()); - break; - case BINARY: - name = PropertyType.BINARY.toString(); - assertArrayEquals(BINARY_VAL, p.getValueAsBinary()); - break; - default: // Unknown type - fail(); - break; - } - assertEquals(name, p.name()); + for (PropertyType t : PropertyType.values()) { + if (!PropertyType.isValid(t.toInt())) { + continue; + } + assertTrue(pit.hasNext()); + MessageProperty p = pit.next().getValue(); + String name = null; + switch (t) { + case BOOL: + name = PropertyType.BOOL.toString(); + assertEquals(BOOL_VAL, p.getValueAsBool()); + break; + case BYTE: + name = PropertyType.BYTE.toString(); + assertEquals(BYTE_VAL, p.getValueAsByte()); + break; + case SHORT: + name = PropertyType.SHORT.toString(); + assertEquals(SHORT_VAL, p.getValueAsShort()); + break; + case INT32: + name = PropertyType.INT32.toString(); + assertEquals(INT32_VAL, p.getValueAsInt32()); + break; + case INT64: + name = PropertyType.INT64.toString(); + assertEquals(INT64_VAL, p.getValueAsInt64()); + break; + case STRING: + name = PropertyType.STRING.toString(); + assertEquals(STRING_VAL, p.getValueAsString()); + break; + case BINARY: + name = PropertyType.BINARY.toString(); + assertArrayEquals(BINARY_VAL, p.getValueAsBinary()); + break; + default: // Unknown type + fail(); + break; } + assertEquals(name, p.name()); } } diff --git a/bmq-sdk/src/test/java/com/bloomberg/bmq/impl/infr/proto/MessagePropertyHeaderTest.java b/bmq-sdk/src/test/java/com/bloomberg/bmq/impl/infr/proto/MessagePropertyHeaderTest.java index 40bf71c9..fcb199bf 100644 --- a/bmq-sdk/src/test/java/com/bloomberg/bmq/impl/infr/proto/MessagePropertyHeaderTest.java +++ b/bmq-sdk/src/test/java/com/bloomberg/bmq/impl/infr/proto/MessagePropertyHeaderTest.java @@ -33,36 +33,36 @@ class MessagePropertyHeaderTest { void testGettersSetters() { MessagePropertyHeader mph = new MessagePropertyHeader(); assertEquals(0, mph.propertyType()); - assertEquals(0, mph.propertyValueLength()); + assertEquals(0, mph.propertyValueOffset()); assertEquals(0, mph.propertyNameLength()); MessagePropertyHeader mph2 = new MessagePropertyHeader(); mph2.setPropertyType(31); // max per protocol - mph2.setPropertyValueLength((1 << 26) - 1); // max per protocol + mph2.setPropertyValueOffset((1 << 26) - 1); // max per protocol mph2.setPropertyNameLength((1 << 12) - 1); // max per protocol assertEquals(31, mph2.propertyType()); - assertEquals(((1 << 26) - 1), mph2.propertyValueLength()); + assertEquals(((1 << 26) - 1), mph2.propertyValueOffset()); assertEquals(((1 << 12) - 1), mph2.propertyNameLength()); MessagePropertyHeader mph3 = new MessagePropertyHeader(); mph3.setPropertyType(17); - mph3.setPropertyValueLength((1 << 19) - 1); + mph3.setPropertyValueOffset((1 << 19) - 1); mph3.setPropertyNameLength((1 << 8) - 1); assertEquals(17, mph3.propertyType()); - assertEquals(((1 << 19) - 1), mph3.propertyValueLength()); + assertEquals(((1 << 19) - 1), mph3.propertyValueOffset()); assertEquals(((1 << 8) - 1), mph3.propertyNameLength()); } @Test void testStreamInStreamOut() throws IOException { - final int valueLength = (1 << 26) - 1; // max per protocol + final int valueOffset = (1 << 26) - 1; // max per protocol final int nameLength = (1 << 12) - 1; // max per protocol MessagePropertyHeader mph = new MessagePropertyHeader(); mph.setPropertyType(PropertyType.STRING.toInt()); - mph.setPropertyValueLength(valueLength); + mph.setPropertyValueOffset(valueOffset); mph.setPropertyNameLength(nameLength); final int[] sizes = @@ -108,7 +108,7 @@ void testStreamInStreamOut() throws IOException { } assertEquals(PropertyType.STRING.toInt(), header.propertyType()); - assertEquals(valueLength, header.propertyValueLength()); + assertEquals(valueOffset, header.propertyValueOffset()); assertEquals(nameLength, header.propertyNameLength()); assertEquals(0, bbis.available()); diff --git a/bmq-sdk/src/test/java/com/bloomberg/bmq/impl/infr/proto/PushEventImplBuilderTest.java b/bmq-sdk/src/test/java/com/bloomberg/bmq/impl/infr/proto/PushEventImplBuilderTest.java index 75d8f90b..baf087aa 100644 --- a/bmq-sdk/src/test/java/com/bloomberg/bmq/impl/infr/proto/PushEventImplBuilderTest.java +++ b/bmq-sdk/src/test/java/com/bloomberg/bmq/impl/infr/proto/PushEventImplBuilderTest.java @@ -34,55 +34,52 @@ class PushEventImplBuilderTest { @Test void testPackPushMessage() throws IOException { - for (boolean isOldStyleProperties : new boolean[] {true, false}) { - PushMessageImpl pushMsg = new PushMessageImpl(); - PushEventBuilder builder = new PushEventBuilder(); + PushMessageImpl pushMsg = new PushMessageImpl(); + PushEventBuilder builder = new PushEventBuilder(); - EventBuilderResult res = builder.packMessage(pushMsg, isOldStyleProperties); - PushHeaderFlags flags = PushHeaderFlags.fromInt(pushMsg.flags()); + EventBuilderResult res = builder.packMessage(pushMsg); + PushHeaderFlags flags = PushHeaderFlags.fromInt(pushMsg.flags()); - assertEquals(EventBuilderResult.SUCCESS, res); - assertEquals(PushHeaderFlags.IMPLICIT_PAYLOAD, flags); + assertEquals(EventBuilderResult.SUCCESS, res); + assertEquals(PushHeaderFlags.IMPLICIT_PAYLOAD, flags); - pushMsg.reset(); + pushMsg.reset(); - ByteBuffer buffer = ByteBuffer.allocate(PushHeader.MAX_PAYLOAD_SIZE_SOFT + 1); - pushMsg.appData().setPayload(buffer); + ByteBuffer buffer = ByteBuffer.allocate(PushHeader.MAX_PAYLOAD_SIZE_SOFT + 1); + pushMsg.appData().setPayload(buffer); - res = builder.packMessage(pushMsg, isOldStyleProperties); - assertEquals(EventBuilderResult.PAYLOAD_TOO_BIG, res); + res = builder.packMessage(pushMsg); + assertEquals(EventBuilderResult.PAYLOAD_TOO_BIG, res); - pushMsg.reset(); - builder.reset(); + pushMsg.reset(); + builder.reset(); - final int numMsgs = EventHeader.MAX_SIZE_SOFT / PushHeader.MAX_PAYLOAD_SIZE_SOFT; - // Cannot pack more than 'numMsgs' having a unpackedSize of - // 'PushHeader.MAX_PAYLOAD_SIZE_SOFT' in 1 bmqp event. + final int numMsgs = EventHeader.MAX_SIZE_SOFT / PushHeader.MAX_PAYLOAD_SIZE_SOFT; + // Cannot pack more than 'numMsgs' having a unpackedSize of + // 'PushHeader.MAX_PAYLOAD_SIZE_SOFT' in 1 bmqp event. - buffer = ByteBuffer.allocate(PushHeader.MAX_PAYLOAD_SIZE_SOFT); - pushMsg.appData().setPayload(buffer); + buffer = ByteBuffer.allocate(PushHeader.MAX_PAYLOAD_SIZE_SOFT); + pushMsg.appData().setPayload(buffer); - for (int i = 0; i < numMsgs; i++) { - res = builder.packMessage(pushMsg, isOldStyleProperties); - assertEquals(EventBuilderResult.SUCCESS, res); - } + for (int i = 0; i < numMsgs; i++) { + res = builder.packMessage(pushMsg); + assertEquals(EventBuilderResult.SUCCESS, res); + } - // Try to add one more message, which must fail with event_too_big. - res = builder.packMessage(pushMsg, isOldStyleProperties); - assertEquals(EventBuilderResult.EVENT_TOO_BIG, res); + // Try to add one more message, which must fail with event_too_big. + res = builder.packMessage(pushMsg); + assertEquals(EventBuilderResult.EVENT_TOO_BIG, res); - pushMsg.reset(); - builder.reset(); + pushMsg.reset(); + builder.reset(); - pushMsg.appData().setPayload(buffer); + pushMsg.appData().setPayload(buffer); - res = builder.packMessage(pushMsg, isOldStyleProperties); - assertEquals(EventBuilderResult.SUCCESS, res); + res = builder.packMessage(pushMsg); + assertEquals(EventBuilderResult.SUCCESS, res); - boolean isSet = - PushHeaderFlags.isSet(pushMsg.flags(), PushHeaderFlags.IMPLICIT_PAYLOAD); - assertFalse(isSet); - } + boolean isSet = PushHeaderFlags.isSet(pushMsg.flags(), PushHeaderFlags.IMPLICIT_PAYLOAD); + assertFalse(isSet); } @Test @@ -109,14 +106,11 @@ void testBuildPushMessage() throws IOException { // set compression to none in order to match file content pushMsg.setCompressionType(CompressionAlgorithmType.E_NONE); - final boolean isOldStyleProperties = i % 2 == 0; - assertEquals( - EventBuilderResult.SUCCESS, builder.packMessage(pushMsg, isOldStyleProperties)); + assertEquals(EventBuilderResult.SUCCESS, builder.packMessage(pushMsg)); // Compare with value stored in the binary pattern logger.info("PUSH header {}: {}", i + 1, pushMsg.header()); - assertEquals(i, pushMsg.header().schemaWireId()); - assertEquals(isOldStyleProperties, pushMsg.appData().isOldStyleProperties()); + assertEquals(1, pushMsg.header().schemaWireId()); } ByteBuffer[] message = builder.build(); diff --git a/bmq-sdk/src/test/java/com/bloomberg/bmq/impl/infr/proto/PushEventImplTest.java b/bmq-sdk/src/test/java/com/bloomberg/bmq/impl/infr/proto/PushEventImplTest.java index b9a35106..7fdb7e15 100644 --- a/bmq-sdk/src/test/java/com/bloomberg/bmq/impl/infr/proto/PushEventImplTest.java +++ b/bmq-sdk/src/test/java/com/bloomberg/bmq/impl/infr/proto/PushEventImplTest.java @@ -35,122 +35,116 @@ class PushEventImplTest { @Test void testDispatchUnknownCompression() throws IOException { - for (boolean isOldStyleProperties : new boolean[] {true, false}) { - final byte[] bytes = new byte[Protocol.COMPRESSION_MIN_APPDATA_SIZE + 1]; + final byte[] bytes = new byte[Protocol.COMPRESSION_MIN_APPDATA_SIZE + 1]; - bytes[0] = 1; - bytes[Protocol.COMPRESSION_MIN_APPDATA_SIZE - 1] = 1; + bytes[0] = 1; + bytes[Protocol.COMPRESSION_MIN_APPDATA_SIZE - 1] = 1; - final int NUM = 4; + final int NUM = 4; - final int unknownType = - EnumSet.allOf(CompressionAlgorithmType.class).stream() - .mapToInt(CompressionAlgorithmType::toInt) - .max() - .getAsInt() - + 1; + final int unknownType = + EnumSet.allOf(CompressionAlgorithmType.class).stream() + .mapToInt(CompressionAlgorithmType::toInt) + .max() + .getAsInt() + + 1; - final MessagePropertiesImpl props = new MessagePropertiesImpl(); - props.setPropertyAsInt32("routingId", 42); - props.setPropertyAsInt64("timestamp", 1234567890L); + final MessagePropertiesImpl props = new MessagePropertiesImpl(); + props.setPropertyAsInt32("routingId", 42); + props.setPropertyAsInt64("timestamp", 1234567890L); - ByteBufferOutputStream bbos = new ByteBufferOutputStream(); + ByteBufferOutputStream bbos = new ByteBufferOutputStream(); - EventHeader header = new EventHeader(); - header.setType(EventType.PUSH); + EventHeader header = new EventHeader(); + header.setType(EventType.PUSH); - final int unpackedSize = props.totalSize() + bytes.length; - final int numPaddingBytes = ProtocolUtil.calculatePadding(unpackedSize); + final int unpackedSize = props.totalSize() + bytes.length; + final int numPaddingBytes = ProtocolUtil.calculatePadding(unpackedSize); - header.setLength( - EventHeader.HEADER_SIZE - + (PushHeader.HEADER_SIZE_FOR_SCHEMA_ID - + unpackedSize - + numPaddingBytes) - * NUM); + header.setLength( + EventHeader.HEADER_SIZE + + (PushHeader.HEADER_SIZE + unpackedSize + numPaddingBytes) * NUM); - header.streamOut(bbos); + header.streamOut(bbos); - for (int i = 0; i < NUM; i++) { - PushMessageImpl pushMsg = new PushMessageImpl(); + for (int i = 0; i < NUM; i++) { + PushMessageImpl pushMsg = new PushMessageImpl(); - pushMsg.appData().setPayload(ByteBuffer.wrap(bytes)); + pushMsg.appData().setPayload(ByteBuffer.wrap(bytes)); - pushMsg.appData().setProperties(props); - pushMsg.appData().setIsOldStyleProperties(isOldStyleProperties); + pushMsg.appData().setProperties(props); - pushMsg.compressData(); + pushMsg.compressData(); - // override compression type for the third message - if (i == 2) { - pushMsg.header().setCompressionType(unknownType); - } - - assertEquals(unpackedSize, pushMsg.appData().unpackedSize()); - assertEquals(numPaddingBytes, pushMsg.appData().numPaddingBytes()); - - pushMsg.streamOut(bbos); + // override compression type for the third message + if (i == 2) { + pushMsg.header().setCompressionType(unknownType); } - PushEventImpl pushEvent = new PushEventImpl(bbos.reset()); - - SessionEventHandler handler = - new SessionEventHandler() { - public void handleControlEvent(ControlEventImpl controlEvent) { - throw new UnsupportedOperationException(); - } - - public void handleAckMessage(AckMessageImpl ackMsg) { - throw new UnsupportedOperationException(); - } - - public void handlePushMessage(PushMessageImpl pushMsg) { - try { - assertArrayEquals( - new ByteBuffer[] {ByteBuffer.wrap(bytes)}, - pushMsg.appData().payload()); - return; - } catch (IOException e) { - logger.error("IOException has been thrown", e); - } - - fail(); // should not get here - } + assertEquals(unpackedSize, pushMsg.appData().unpackedSize()); + assertEquals(numPaddingBytes, pushMsg.appData().numPaddingBytes()); - public void handlePutEvent(PutEventImpl putEvent) { - throw new UnsupportedOperationException(); - } + pushMsg.streamOut(bbos); + } - public void handleConfirmEvent(ConfirmEventImpl confirmEvent) { - throw new UnsupportedOperationException(); + PushEventImpl pushEvent = new PushEventImpl(bbos.reset()); + + SessionEventHandler handler = + new SessionEventHandler() { + public void handleControlEvent(ControlEventImpl controlEvent) { + throw new UnsupportedOperationException(); + } + + public void handleAckMessage(AckMessageImpl ackMsg) { + throw new UnsupportedOperationException(); + } + + public void handlePushMessage(PushMessageImpl pushMsg) { + try { + assertArrayEquals( + new ByteBuffer[] {ByteBuffer.wrap(bytes)}, + pushMsg.appData().payload()); + return; + } catch (IOException e) { + logger.error("IOException has been thrown", e); } - }; - - try { - pushEvent.dispatch(handler); - fail(); // should not get here - } catch (IllegalArgumentException e) { - // According to PushMessageIterator used in 'dispatch' method, - // when 'next()' method is called, the next item after - // the current one is also prefetched and parsed. - // If the next item is invalid then the exception will be thrown for - // the current one. - // - // In our situation, the first item should be processed successfully. - // When the second one is being dispatched, which is valid, - // an exception should be thrown related to the third item. - // On another hand, in the 'dispatch()', method number of messages is incremented - // before calling the handler. So eventually when the exception is thrown, - // the number of messages should be equal to 2 - assertEquals( - String.format("'%d' - unknown compression algorithm type", unknownType), - e.getMessage()); - - // only the first message should be processed successfully - assertEquals(2, pushEvent.messageCount()); - } + fail(); // should not get here + } + + public void handlePutEvent(PutEventImpl putEvent) { + throw new UnsupportedOperationException(); + } + + public void handleConfirmEvent(ConfirmEventImpl confirmEvent) { + throw new UnsupportedOperationException(); + } + }; + + try { + pushEvent.dispatch(handler); + fail(); // should not get here + } catch (IllegalArgumentException e) { + // According to PushMessageIterator used in 'dispatch' method, + // when 'next()' method is called, the next item after + // the current one is also prefetched and parsed. + // If the next item is invalid then the exception will be thrown for + // the current one. + // + // In our situation, the first item should be processed successfully. + // When the second one is being dispatched, which is valid, + // an exception should be thrown related to the third item. + // On another hand, in the 'dispatch()', method number of messages is incremented + // before calling the handler. So eventually when the exception is thrown, + // the number of messages should be equal to 2 + assertEquals( + String.format("'%d' - unknown compression algorithm type", unknownType), + e.getMessage()); + + // only the first message should be processed successfully assertEquals(2, pushEvent.messageCount()); } + + assertEquals(2, pushEvent.messageCount()); } } diff --git a/bmq-sdk/src/test/java/com/bloomberg/bmq/impl/infr/proto/PushHeaderTest.java b/bmq-sdk/src/test/java/com/bloomberg/bmq/impl/infr/proto/PushHeaderTest.java index e4fd570f..98386429 100644 --- a/bmq-sdk/src/test/java/com/bloomberg/bmq/impl/infr/proto/PushHeaderTest.java +++ b/bmq-sdk/src/test/java/com/bloomberg/bmq/impl/infr/proto/PushHeaderTest.java @@ -85,7 +85,7 @@ void testStreamIn() throws IOException { assertEquals(0, pushHeader.compressionType()); assertEquals(8, pushHeader.headerWords()); assertEquals(9876, pushHeader.queueId()); - assertEquals(i, pushHeader.schemaWireId()); + assertEquals(1, pushHeader.schemaWireId()); assertEquals("ABCDEF0123456789ABCDEF0123456789", pushHeader.messageGUID().toString()); diff --git a/bmq-sdk/src/test/java/com/bloomberg/bmq/impl/infr/proto/PushMessageImplTest.java b/bmq-sdk/src/test/java/com/bloomberg/bmq/impl/infr/proto/PushMessageImplTest.java index 17354a61..40b41626 100644 --- a/bmq-sdk/src/test/java/com/bloomberg/bmq/impl/infr/proto/PushMessageImplTest.java +++ b/bmq-sdk/src/test/java/com/bloomberg/bmq/impl/infr/proto/PushMessageImplTest.java @@ -40,7 +40,7 @@ import org.slf4j.LoggerFactory; class PushMessageImplTest { - static final int HEADER_WORDS = PushHeader.HEADER_SIZE_FOR_SCHEMA_ID / Protocol.WORD_SIZE; + static final int HEADER_WORDS = PushHeader.HEADER_SIZE / Protocol.WORD_SIZE; static final Logger logger = LoggerFactory.getLogger(MethodHandles.lookup().lookupClass()); @Test @@ -66,16 +66,14 @@ void testStreamIn() throws Exception { for (ByteBuffer[] payload : payloads) for (MessagePropertiesImpl props : propsArray) - for (boolean isOldStyleProperties : new boolean[] {false, true}) - for (CompressionAlgorithmType compressionType : compressionTypes) { - logger.info( - "Stream in with payload:{}, props:{}, oldStyleProperties:{}, compression: {}", - payload, - props, - isOldStyleProperties, - compressionType); - verifyStreamIn(payload, props, isOldStyleProperties, compressionType); - } + for (CompressionAlgorithmType compressionType : compressionTypes) { + logger.info( + "Stream in with payload:{}, props:{}, compression: {}", + payload, + props, + compressionType); + verifyStreamIn(payload, props, compressionType); + } } @Test @@ -103,16 +101,14 @@ void testStreamOut() throws Exception { for (ByteBuffer[] payload : payloads) for (MessagePropertiesImpl props : propsArray) - for (boolean isOldStyleProperties : new boolean[] {false, true}) - for (CompressionAlgorithmType compressionType : compressionTypes) { - logger.info( - "Stream out with payload:{}, props:{}, oldStyleProperties:{}, compression: {}", - payload, - props, - isOldStyleProperties, - compressionType); - verifyStreamOut(payload, props, isOldStyleProperties, compressionType); - } + for (CompressionAlgorithmType compressionType : compressionTypes) { + logger.info( + "Stream out with payload:{}, props:{}, compression: {}", + payload, + props, + compressionType); + verifyStreamOut(payload, props, compressionType); + } } @Test @@ -242,13 +238,11 @@ private MessagePropertiesImpl generateProps() { private ByteBuffer[] generateOutput( ByteBuffer[] payload, MessagePropertiesImpl props, - boolean isOldStyleProperties, CompressionAlgorithmType compressionType) throws IOException { try (ByteBufferOutputStream bbos = new ByteBufferOutputStream()) { ApplicationData appData = new ApplicationData(); - appData.setIsOldStyleProperties(isOldStyleProperties); if (payload != null) { appData.setPayload(payload); @@ -259,15 +253,11 @@ private ByteBuffer[] generateOutput( appData.setProperties(props); if (appData.hasProperties()) { header.setFlags(PushHeaderFlags.setFlag(0, PushHeaderFlags.MESSAGE_PROPERTIES)); - - if (!isOldStyleProperties) { - header.setSchemaWireId(PushMessageImpl.INVALID_SCHEMA_WIRE_ID); - } + header.setSchemaWireId(PushMessageImpl.INVALID_SCHEMA_WIRE_ID); } - final int sizeToCompress = - isOldStyleProperties ? appData.unpackedSize() : appData.payloadSize(); - final boolean canCompress = sizeToCompress >= Protocol.COMPRESSION_MIN_APPDATA_SIZE; + final boolean canCompress = + appData.payloadSize() >= Protocol.COMPRESSION_MIN_APPDATA_SIZE; // Compress data if compression is set and size is not below threshold if (compressionType != null && canCompress) { @@ -325,15 +315,13 @@ private ByteBuffer[] generateOptionsOutput(byte[] options) throws IOException { private void verifyStreamIn( ByteBuffer[] payload, MessagePropertiesImpl props, - boolean isOldStyleProperties, CompressionAlgorithmType compressionType) throws IOException { final boolean hasProperties = props != null && props.numProperties() > 0; final boolean hasPayload = getSize(payload) > 0; - ByteBuffer[] input = - generateOutput(duplicate(payload), props, isOldStyleProperties, compressionType); + ByteBuffer[] input = generateOutput(duplicate(payload), props, compressionType); ByteBufferInputStream bbis = new ByteBufferInputStream(duplicate(input)); // Get num of padding bytes @@ -343,7 +331,7 @@ private void verifyStreamIn( // Stream in and ensure that IMPLICIT_PAYLOAD flag is not set and data is not empty final int size = bbis.available(); - final int unpackedSize = size - numPaddingBytes - PushHeader.HEADER_SIZE_FOR_SCHEMA_ID; + final int unpackedSize = size - numPaddingBytes - PushHeader.HEADER_SIZE; try { msg.streamIn(bbis); @@ -368,7 +356,6 @@ private void verifyStreamIn( // check properties verifyProperties(hasProperties ? props : null, msg.appData().properties()); - assertEquals(!hasProperties || isOldStyleProperties, msg.appData().isOldStyleProperties()); // Do double check. Stream out data and compare with original input msg = new PushMessageImpl(); @@ -381,7 +368,6 @@ private void verifyStreamIn( msg.appData().setProperties(props); } - msg.appData().setIsOldStyleProperties(isOldStyleProperties); msg.compressData(); ByteBufferOutputStream bbos = new ByteBufferOutputStream(); @@ -417,7 +403,6 @@ private void verifyStreamIn( private void verifyStreamOut( ByteBuffer[] payload, MessagePropertiesImpl props, - boolean isOldStyleProperties, CompressionAlgorithmType compressionType) throws IOException { @@ -445,12 +430,8 @@ private void verifyStreamOut( msg.appData().setPayload(duplicate(payload)); } - msg.appData().setIsOldStyleProperties(isOldStyleProperties); - - final int dataToCompress = - isOldStyleProperties ? msg.appData().unpackedSize() : msg.appData().payloadSize(); CompressionAlgorithmType actualCompressionType; - if (dataToCompress < Protocol.COMPRESSION_MIN_APPDATA_SIZE) { + if (msg.appData().payloadSize() < Protocol.COMPRESSION_MIN_APPDATA_SIZE) { // Data is not compressed if the size below the threshold. actualCompressionType = CompressionAlgorithmType.E_NONE; } else { @@ -481,8 +462,7 @@ private void verifyStreamOut( assertFalse(PushHeaderFlags.isSet(msg.flags(), PushHeaderFlags.IMPLICIT_PAYLOAD)); } - ByteBuffer[] expected = - generateOutput(duplicate(payload), props, isOldStyleProperties, compressionType); + ByteBuffer[] expected = generateOutput(duplicate(payload), props, compressionType); ByteBuffer[] streamedData = bbos.reset(); assertArrayEquals(expected, streamedData); diff --git a/bmq-sdk/src/test/java/com/bloomberg/bmq/impl/infr/proto/PushMessageIteratorTest.java b/bmq-sdk/src/test/java/com/bloomberg/bmq/impl/infr/proto/PushMessageIteratorTest.java index c27830e8..6affc87c 100644 --- a/bmq-sdk/src/test/java/com/bloomberg/bmq/impl/infr/proto/PushMessageIteratorTest.java +++ b/bmq-sdk/src/test/java/com/bloomberg/bmq/impl/infr/proto/PushMessageIteratorTest.java @@ -126,14 +126,8 @@ void testWithPattern() throws IOException { int i = 0; while (pushIt.hasNext()) { // Build expected content - final boolean isOldStyleProperties = i % 2 == 0; - final ByteBufferOutputStream bbos = new ByteBufferOutputStream(); - if (isOldStyleProperties) { - props.streamOutOld(bbos); - } else { - props.streamOut(bbos); - } + props.streamOut(bbos); bbos.writeAscii(PAYLOAD); @@ -155,8 +149,7 @@ void testWithPattern() throws IOException { assertEquals(GUID, pushMsg.messageGUID().toHex()); assertEquals(QUEUE_ID, pushMsg.queueId()); assertEquals(FLAGS, pushMsg.flags()); - assertEquals(isOldStyleProperties ? 0 : 1, pushMsg.header().schemaWireId()); - assertEquals(isOldStyleProperties, pushMsg.appData().isOldStyleProperties()); + assertEquals(1, pushMsg.header().schemaWireId()); assertArrayEquals(new Integer[] {0}, pushMsg.subQueueIds()); ByteBuffer[] data = pushMsg.appData().applicationData(); @@ -200,8 +193,6 @@ void testWithBuilder() throws IOException { final PushMessageImpl[] pushs = new PushMessageImpl[NUM_MSGS]; for (int i = 0; i < NUM_MSGS; i++) { - final boolean isOldStyleProperties = i % 2 == 0; - PushMessageImpl msg = new PushMessageImpl(); msg.setFlags(FLAGS); msg.setQueueId(i); @@ -209,10 +200,9 @@ void testWithBuilder() throws IOException { msg.appData().setProperties(props); msg.appData().setPayload(payload); - EventBuilderResult rc = builder.packMessage(msg, isOldStyleProperties); + EventBuilderResult rc = builder.packMessage(msg); assertEquals(EventBuilderResult.SUCCESS, rc); - assertEquals(isOldStyleProperties, msg.appData().isOldStyleProperties()); - assertEquals(isOldStyleProperties ? 0 : 1, msg.header().schemaWireId()); + assertEquals(1, msg.header().schemaWireId()); pushs[i] = msg; } @@ -233,8 +223,7 @@ void testWithBuilder() throws IOException { assertEquals(GUID, msg.messageGUID()); assertEquals(i, msg.queueId()); assertEquals(FLAGS, msg.flags()); - assertEquals(i % 2, msg.header().schemaWireId()); - assertEquals(i % 2 == 0, msg.appData().isOldStyleProperties()); + assertEquals(1, msg.header().schemaWireId()); assertEquals(exp.appData().numPaddingBytes(), msg.appData().numPaddingBytes()); assertArrayEquals(exp.appData().applicationData(), msg.appData().applicationData()); @@ -312,8 +301,6 @@ void testMultipleCompressedMessages() throws IOException { final int NUM = 500; for (int i = 0; i < NUM; i++) { - final boolean isOldStyleProperties = i % 2 == 0; - PushMessageImpl pushMsg = new PushMessageImpl(); pushMsg.appData().setPayload(ByteBuffer.wrap(bytes)); @@ -326,7 +313,7 @@ void testMultipleCompressedMessages() throws IOException { pushMsg.setCompressionType(CompressionAlgorithmType.E_ZLIB); - builder.packMessage(pushMsg, isOldStyleProperties); + builder.packMessage(pushMsg); } PushEventImpl pushEvent = new PushEventImpl(builder.build()); @@ -364,93 +351,87 @@ void testMultipleCompressedMessages() throws IOException { @Test void testUnknownCompression() throws IOException { - for (boolean isOldStyleProperties : new boolean[] {true, false}) { - final byte[] bytes = new byte[Protocol.COMPRESSION_MIN_APPDATA_SIZE + 1]; + final byte[] bytes = new byte[Protocol.COMPRESSION_MIN_APPDATA_SIZE + 1]; - bytes[0] = 1; - bytes[Protocol.COMPRESSION_MIN_APPDATA_SIZE - 1] = 1; + bytes[0] = 1; + bytes[Protocol.COMPRESSION_MIN_APPDATA_SIZE - 1] = 1; - final int NUM = 4; + final int NUM = 4; - final int unknownType = - EnumSet.allOf(CompressionAlgorithmType.class).stream() - .mapToInt(CompressionAlgorithmType::toInt) - .max() - .getAsInt() - + 1; + final int unknownType = + EnumSet.allOf(CompressionAlgorithmType.class).stream() + .mapToInt(CompressionAlgorithmType::toInt) + .max() + .getAsInt() + + 1; - final MessagePropertiesImpl props = new MessagePropertiesImpl(); - props.setPropertyAsInt32("routingId", 42); - props.setPropertyAsInt64("timestamp", 1234567890L); + final MessagePropertiesImpl props = new MessagePropertiesImpl(); + props.setPropertyAsInt32("routingId", 42); + props.setPropertyAsInt64("timestamp", 1234567890L); - ByteBufferOutputStream bbos = new ByteBufferOutputStream(); + ByteBufferOutputStream bbos = new ByteBufferOutputStream(); - EventHeader header = new EventHeader(); - header.setType(EventType.PUSH); + EventHeader header = new EventHeader(); + header.setType(EventType.PUSH); - final int unpackedSize = props.totalSize() + bytes.length; - final int numPaddingBytes = ProtocolUtil.calculatePadding(unpackedSize); + final int unpackedSize = props.totalSize() + bytes.length; + final int numPaddingBytes = ProtocolUtil.calculatePadding(unpackedSize); - header.setLength( - EventHeader.HEADER_SIZE - + (PushHeader.HEADER_SIZE_FOR_SCHEMA_ID - + unpackedSize - + numPaddingBytes) - * NUM); + header.setLength( + EventHeader.HEADER_SIZE + + (PushHeader.HEADER_SIZE + unpackedSize + numPaddingBytes) * NUM); - header.streamOut(bbos); + header.streamOut(bbos); - for (int i = 0; i < NUM; i++) { - PushMessageImpl pushMsg = new PushMessageImpl(); + for (int i = 0; i < NUM; i++) { + PushMessageImpl pushMsg = new PushMessageImpl(); - pushMsg.appData().setPayload(ByteBuffer.wrap(bytes)); + pushMsg.appData().setPayload(ByteBuffer.wrap(bytes)); - pushMsg.appData().setProperties(props); - pushMsg.appData().setIsOldStyleProperties(isOldStyleProperties); + pushMsg.appData().setProperties(props); - pushMsg.compressData(); + pushMsg.compressData(); - // override compression type for the third message - if (i == 2) { - pushMsg.header().setCompressionType(unknownType); - } + // override compression type for the third message + if (i == 2) { + pushMsg.header().setCompressionType(unknownType); + } - assertEquals(unpackedSize, pushMsg.appData().unpackedSize()); - assertEquals(numPaddingBytes, pushMsg.appData().numPaddingBytes()); + assertEquals(unpackedSize, pushMsg.appData().unpackedSize()); + assertEquals(numPaddingBytes, pushMsg.appData().numPaddingBytes()); - pushMsg.streamOut(bbos); - } + pushMsg.streamOut(bbos); + } - PushEventImpl pushEvent = new PushEventImpl(bbos.reset()); - PushMessageIterator pushIt = pushEvent.iterator(); + PushEventImpl pushEvent = new PushEventImpl(bbos.reset()); + PushMessageIterator pushIt = pushEvent.iterator(); - int counter = 0; - try { - while (pushIt.hasNext()) { - PushMessageImpl pushMsg = pushIt.next(); + int counter = 0; + try { + while (pushIt.hasNext()) { + PushMessageImpl pushMsg = pushIt.next(); - assertArrayEquals( - new ByteBuffer[] {ByteBuffer.wrap(bytes)}, pushMsg.appData().payload()); - counter++; - } - } catch (IllegalArgumentException e) { - // According to PushMessageIterator, when 'next()' method is called, - // the next item after the current one is also prefetched and parsed. - // If the next item is invalid then the exception will be thrown for - // the current one. - // - // In our situation, the first item should be processed successfully. - // When we call 'next()' to get the second one, which is valid, - // an exception should be thrown related to the third item. - assertEquals( - String.format("'%d' - unknown compression algorithm type", unknownType), - e.getMessage()); - - // only the first message should be processed successfully - assertEquals(1, counter); + assertArrayEquals( + new ByteBuffer[] {ByteBuffer.wrap(bytes)}, pushMsg.appData().payload()); + counter++; } - + } catch (IllegalArgumentException e) { + // According to PushMessageIterator, when 'next()' method is called, + // the next item after the current one is also prefetched and parsed. + // If the next item is invalid then the exception will be thrown for + // the current one. + // + // In our situation, the first item should be processed successfully. + // When we call 'next()' to get the second one, which is valid, + // an exception should be thrown related to the third item. + assertEquals( + String.format("'%d' - unknown compression algorithm type", unknownType), + e.getMessage()); + + // only the first message should be processed successfully assertEquals(1, counter); } + + assertEquals(1, counter); } } diff --git a/bmq-sdk/src/test/java/com/bloomberg/bmq/impl/infr/proto/PutEventImplBuilderTest.java b/bmq-sdk/src/test/java/com/bloomberg/bmq/impl/infr/proto/PutEventImplBuilderTest.java index 288b0d44..3e0013f2 100644 --- a/bmq-sdk/src/test/java/com/bloomberg/bmq/impl/infr/proto/PutEventImplBuilderTest.java +++ b/bmq-sdk/src/test/java/com/bloomberg/bmq/impl/infr/proto/PutEventImplBuilderTest.java @@ -36,99 +36,95 @@ class PutEventImplBuilderTest { @Test void testErrorPutMessage() throws IOException { - for (boolean isOldStyleProperties : new boolean[] {false, true}) { - PutMessageImpl putMsg = new PutMessageImpl(); - PutEventBuilder builder = new PutEventBuilder(); + PutMessageImpl putMsg = new PutMessageImpl(); + PutEventBuilder builder = new PutEventBuilder(); - EventBuilderResult res = builder.packMessage(putMsg, isOldStyleProperties); - assertEquals(EventBuilderResult.PAYLOAD_EMPTY, res); + EventBuilderResult res = builder.packMessage(putMsg); + assertEquals(EventBuilderResult.PAYLOAD_EMPTY, res); - putMsg = new PutMessageImpl(); - ByteBuffer buffer = ByteBuffer.allocate(0); - putMsg.appData().setPayload(buffer); + putMsg = new PutMessageImpl(); + ByteBuffer buffer = ByteBuffer.allocate(0); + putMsg.appData().setPayload(buffer); - res = builder.packMessage(putMsg, isOldStyleProperties); - assertEquals(EventBuilderResult.PAYLOAD_EMPTY, res); + res = builder.packMessage(putMsg); + assertEquals(EventBuilderResult.PAYLOAD_EMPTY, res); - putMsg = new PutMessageImpl(); - buffer = ByteBuffer.allocate(PutHeader.MAX_PAYLOAD_SIZE_SOFT + 1); - putMsg.appData().setPayload(buffer); + putMsg = new PutMessageImpl(); + buffer = ByteBuffer.allocate(PutHeader.MAX_PAYLOAD_SIZE_SOFT + 1); + putMsg.appData().setPayload(buffer); - // set compression to none in order to get PAYLOAD_TOO_BIG result - putMsg.setCompressionType(CompressionAlgorithmType.E_NONE); + // set compression to none in order to get PAYLOAD_TOO_BIG result + putMsg.setCompressionType(CompressionAlgorithmType.E_NONE); - res = builder.packMessage(putMsg, isOldStyleProperties); - assertEquals(EventBuilderResult.PAYLOAD_TOO_BIG, res); + res = builder.packMessage(putMsg); + assertEquals(EventBuilderResult.PAYLOAD_TOO_BIG, res); - putMsg = new PutMessageImpl(); + putMsg = new PutMessageImpl(); - final int numMsgs = EventHeader.MAX_SIZE_SOFT / PutHeader.MAX_PAYLOAD_SIZE_SOFT; - // Cannot pack more than 'numMsgs' having a unpackedSize of - // 'PutHeader.MAX_PAYLOAD_SIZE_SOFT' in 1 bmqp event. + final int numMsgs = EventHeader.MAX_SIZE_SOFT / PutHeader.MAX_PAYLOAD_SIZE_SOFT; + // Cannot pack more than 'numMsgs' having a unpackedSize of + // 'PutHeader.MAX_PAYLOAD_SIZE_SOFT' in 1 bmqp event. - buffer = ByteBuffer.allocate(PutHeader.MAX_PAYLOAD_SIZE_SOFT); - putMsg.appData().setPayload(buffer); + buffer = ByteBuffer.allocate(PutHeader.MAX_PAYLOAD_SIZE_SOFT); + putMsg.appData().setPayload(buffer); - // set compression to none in order to get EVENT_TOO_BIG result - putMsg.setCompressionType(CompressionAlgorithmType.E_NONE); + // set compression to none in order to get EVENT_TOO_BIG result + putMsg.setCompressionType(CompressionAlgorithmType.E_NONE); - for (int i = 0; i < numMsgs; i++) { - res = builder.packMessage(putMsg, isOldStyleProperties); - assertEquals(EventBuilderResult.SUCCESS, res); - } + for (int i = 0; i < numMsgs; i++) { + res = builder.packMessage(putMsg); + assertEquals(EventBuilderResult.SUCCESS, res); + } - // Try to add one more message, which must fail with event_too_big. - res = builder.packMessage(putMsg, isOldStyleProperties); - assertEquals(EventBuilderResult.EVENT_TOO_BIG, res); + // Try to add one more message, which must fail with event_too_big. + res = builder.packMessage(putMsg); + assertEquals(EventBuilderResult.EVENT_TOO_BIG, res); - putMsg = new PutMessageImpl(); + putMsg = new PutMessageImpl(); - putMsg.appData().setPayload(buffer); - putMsg.setFlags(PutHeaderFlags.ACK_REQUESTED.toInt()); + putMsg.appData().setPayload(buffer); + putMsg.setFlags(PutHeaderFlags.ACK_REQUESTED.toInt()); - CorrelationId corId = putMsg.correlationId(); - assertNull(corId); + CorrelationId corId = putMsg.correlationId(); + assertNull(corId); - // Try to pack with default CorrelationID which is zero - res = builder.packMessage(putMsg, isOldStyleProperties); - assertEquals(EventBuilderResult.MISSING_CORRELATION_ID, res); - } + // Try to pack with default CorrelationID which is zero + res = builder.packMessage(putMsg); + assertEquals(EventBuilderResult.MISSING_CORRELATION_ID, res); } @Test void testBigPutEvent() throws IOException { // Check that ByteBuffer limit is honored when Put event is being built - for (boolean isOldStyleProperties : new boolean[] {false, true}) { - PutMessageImpl putMsg = new PutMessageImpl(); - PutEventBuilder builder = new PutEventBuilder(); + PutMessageImpl putMsg = new PutMessageImpl(); + PutEventBuilder builder = new PutEventBuilder(); - ByteBuffer buffer = ByteBuffer.allocate(EventHeader.MAX_SIZE_SOFT + 1024); - buffer.limit(PutHeader.MAX_PAYLOAD_SIZE_SOFT); + ByteBuffer buffer = ByteBuffer.allocate(EventHeader.MAX_SIZE_SOFT + 1024); + buffer.limit(PutHeader.MAX_PAYLOAD_SIZE_SOFT); - putMsg.appData().setPayload(buffer); + putMsg.appData().setPayload(buffer); - // set compression to none in order to get the unpacked size equal to - // PutHeader.MAX_PAYLOAD_SIZE_SOFT. - putMsg.setCompressionType(CompressionAlgorithmType.E_NONE); + // set compression to none in order to get the unpacked size equal to + // PutHeader.MAX_PAYLOAD_SIZE_SOFT. + putMsg.setCompressionType(CompressionAlgorithmType.E_NONE); - EventBuilderResult res = builder.packMessage(putMsg, isOldStyleProperties); - assertEquals(EventBuilderResult.SUCCESS, res); + EventBuilderResult res = builder.packMessage(putMsg); + assertEquals(EventBuilderResult.SUCCESS, res); - ByteBuffer[] message; - message = builder.build(); + ByteBuffer[] message; + message = builder.build(); - int size = 0; - for (ByteBuffer b : message) { - size += b.limit(); - } + int size = 0; + for (ByteBuffer b : message) { + size += b.limit(); + } - logger.info("EventImpl size : {}", size); - logger.info("PutHeader.MAX_PAYLOAD_SIZE_SOFT : {}", PutHeader.MAX_PAYLOAD_SIZE_SOFT); - logger.info("EventHeader.MAX_SIZE_SOFT : {}", EventHeader.MAX_SIZE_SOFT); + logger.info("EventImpl size : {}", size); + logger.info("PutHeader.MAX_PAYLOAD_SIZE_SOFT : {}", PutHeader.MAX_PAYLOAD_SIZE_SOFT); + logger.info("EventHeader.MAX_SIZE_SOFT : {}", EventHeader.MAX_SIZE_SOFT); - assertTrue(size <= EventHeader.MAX_SIZE_SOFT); - } + assertTrue(size <= EventHeader.MAX_SIZE_SOFT); } @Test @@ -146,8 +142,6 @@ void testBuildPutMessageWithProperties() throws IOException { int flags = PutHeaderFlags.setFlag(0, PutHeaderFlags.ACK_REQUESTED); flags = PutHeaderFlags.setFlag(flags, PutHeaderFlags.MESSAGE_PROPERTIES); - final long[] crc32s = new long[] {3469549003L, 340340870L}; - for (int i = 0; i < 2; i++) { PutMessageImpl putMsg = new PutMessageImpl(); putMsg.setQueueId(9876); @@ -159,15 +153,12 @@ void testBuildPutMessageWithProperties() throws IOException { // set compression to none in order to match file content putMsg.setCompressionType(CompressionAlgorithmType.E_NONE); - final boolean isOldStyleProperties = i % 2 == 0; - assertEquals( - EventBuilderResult.SUCCESS, builder.packMessage(putMsg, isOldStyleProperties)); + assertEquals(EventBuilderResult.SUCCESS, builder.packMessage(putMsg)); // Compare with value stored in the binary pattern logger.info("PUT header {}: {}", i + 1, putMsg.header()); - assertEquals(crc32s[i], putMsg.crc32c()); - assertEquals(i, putMsg.header().schemaWireId()); - assertEquals(isOldStyleProperties, putMsg.appData().isOldStyleProperties()); + assertEquals(340340870L, putMsg.crc32c()); + assertEquals(1, putMsg.header().schemaWireId()); } ByteBuffer[] message = builder.build(); diff --git a/bmq-sdk/src/test/java/com/bloomberg/bmq/impl/infr/proto/PutHeaderTest.java b/bmq-sdk/src/test/java/com/bloomberg/bmq/impl/infr/proto/PutHeaderTest.java index 635da825..13771555 100644 --- a/bmq-sdk/src/test/java/com/bloomberg/bmq/impl/infr/proto/PutHeaderTest.java +++ b/bmq-sdk/src/test/java/com/bloomberg/bmq/impl/infr/proto/PutHeaderTest.java @@ -53,9 +53,6 @@ void testStreamIn() throws IOException { assertNotNull(header.type()); assertEquals(EventType.PUT, header.type()); - final long[] crc32s = new long[] {3469549003L, 340340870L}; - final int[] schemaIds = new int[] {0, 1}; - for (int i = 0; i < 2; ++i) { PutHeader putHeader = new PutHeader(); @@ -71,8 +68,8 @@ void testStreamIn() throws IOException { CorrelationId corId = CorrelationIdImpl.restoreId(1234); assertEquals(corId, putHeader.correlationId()); - assertEquals(crc32s[i], putHeader.crc32c()); - assertEquals(schemaIds[i], putHeader.schemaWireId()); + assertEquals(340340870L, putHeader.crc32c()); + assertEquals(1, putHeader.schemaWireId()); final int toSkip = (putHeader.messageWords() - putHeader.headerWords()) * Protocol.WORD_SIZE; diff --git a/bmq-sdk/src/test/java/com/bloomberg/bmq/impl/infr/proto/PutMessageImplTest.java b/bmq-sdk/src/test/java/com/bloomberg/bmq/impl/infr/proto/PutMessageImplTest.java index 38e89d20..2eb31e9f 100644 --- a/bmq-sdk/src/test/java/com/bloomberg/bmq/impl/infr/proto/PutMessageImplTest.java +++ b/bmq-sdk/src/test/java/com/bloomberg/bmq/impl/infr/proto/PutMessageImplTest.java @@ -65,16 +65,14 @@ void testStreamOut() throws Exception { for (ByteBuffer[] payload : payloads) for (MessagePropertiesImpl props : propsArray) - for (boolean isOldStyleProperties : new boolean[] {false, true}) - for (CompressionAlgorithmType compressionType : compressionTypes) { - logger.info( - "Stream out with payload:{}, props:{}, oldStyleProperties:{}, compression: {}", - payload, - props, - isOldStyleProperties, - compressionType); - verifyStreamOut(payload, props, isOldStyleProperties, compressionType); - } + for (CompressionAlgorithmType compressionType : compressionTypes) { + logger.info( + "Stream out with payload:{}, props:{}, compression: {}", + payload, + props, + compressionType); + verifyStreamOut(payload, props, compressionType); + } } @Test @@ -100,16 +98,14 @@ void testStreamIn() throws Exception { for (ByteBuffer[] payload : payloads) for (MessagePropertiesImpl props : propsArray) - for (boolean isOldStyleProperties : new boolean[] {false, true}) - for (CompressionAlgorithmType compressionType : compressionTypes) { - logger.info( - "Stream in with payload:{}, props:{}, oldStyleProperties:{}, compression: {}", - payload, - props, - isOldStyleProperties, - compressionType); - verifyStreamIn(payload, props, isOldStyleProperties, compressionType); - } + for (CompressionAlgorithmType compressionType : compressionTypes) { + logger.info( + "Stream in with payload:{}, props:{}, compression: {}", + payload, + props, + compressionType); + verifyStreamIn(payload, props, compressionType); + } } private ByteBuffer[] generatePayload(int size) throws IOException { @@ -134,13 +130,11 @@ private MessagePropertiesImpl generateProps() { private ByteBuffer[] generateOutput( ByteBuffer[] payload, MessagePropertiesImpl props, - boolean isOldStyleProperties, CompressionAlgorithmType compressionType) throws IOException { try (ByteBufferOutputStream bbos = new ByteBufferOutputStream()) { ApplicationData appData = new ApplicationData(); - appData.setIsOldStyleProperties(isOldStyleProperties); if (payload != null) { appData.setPayload(payload); @@ -151,15 +145,11 @@ private ByteBuffer[] generateOutput( appData.setProperties(props); if (appData.hasProperties()) { header.setFlags(PutHeaderFlags.setFlag(0, PutHeaderFlags.MESSAGE_PROPERTIES)); - - if (!isOldStyleProperties) { - header.setSchemaWireId(PutMessageImpl.INVALID_SCHEMA_WIRE_ID); - } + header.setSchemaWireId(PutMessageImpl.INVALID_SCHEMA_WIRE_ID); } - final int sizeToCompress = - isOldStyleProperties ? appData.unpackedSize() : appData.payloadSize(); - final boolean canCompress = sizeToCompress >= Protocol.COMPRESSION_MIN_APPDATA_SIZE; + final boolean canCompress = + appData.payloadSize() >= Protocol.COMPRESSION_MIN_APPDATA_SIZE; // Compress data if compression is set and size is not below threshold if (compressionType != null && canCompress) { @@ -185,7 +175,6 @@ private ByteBuffer[] generateOutput( private void verifyStreamOut( ByteBuffer[] payload, MessagePropertiesImpl props, - boolean isOldStyleProperties, CompressionAlgorithmType compressionType) throws IOException { @@ -226,12 +215,8 @@ private void verifyStreamOut( return; } - msg.appData().setIsOldStyleProperties(isOldStyleProperties); - - final int dataToCompress = - isOldStyleProperties ? msg.appData().unpackedSize() : msg.appData().payloadSize(); CompressionAlgorithmType actualCompressionType; - if (dataToCompress < Protocol.COMPRESSION_MIN_APPDATA_SIZE) { + if (msg.appData().payloadSize() < Protocol.COMPRESSION_MIN_APPDATA_SIZE) { // Data is not compressed if the size below the threshold. actualCompressionType = CompressionAlgorithmType.E_NONE; } else { @@ -252,8 +237,7 @@ private void verifyStreamOut( assertEquals(0, bbos.size()); msg.streamOut(bbos); - ByteBuffer[] expected = - generateOutput(duplicate(payload), props, isOldStyleProperties, compressionType); + ByteBuffer[] expected = generateOutput(duplicate(payload), props, compressionType); ByteBuffer[] streamedData = bbos.reset(); assertArrayEquals(expected, streamedData); @@ -285,15 +269,13 @@ private void verifyStreamOut( private void verifyStreamIn( ByteBuffer[] payload, MessagePropertiesImpl props, - boolean isOldStyleProperties, CompressionAlgorithmType compressionType) throws IOException { final boolean hasProperties = props != null && props.numProperties() > 0; final boolean hasPayload = getSize(payload) > 0; - ByteBuffer[] input = - generateOutput(duplicate(payload), props, isOldStyleProperties, compressionType); + ByteBuffer[] input = generateOutput(duplicate(payload), props, compressionType); ByteBufferInputStream bbis = new ByteBufferInputStream(duplicate(input)); // Get num of padding bytes @@ -325,7 +307,6 @@ private void verifyStreamIn( // check properties verifyProperties(hasProperties ? props : null, msg.appData().properties()); - assertEquals(!hasProperties || isOldStyleProperties, msg.appData().isOldStyleProperties()); // Do double check. Stream out data and compare with original input msg = new PutMessageImpl(); @@ -338,8 +319,6 @@ private void verifyStreamIn( msg.appData().setProperties(props); } - msg.appData().setIsOldStyleProperties(isOldStyleProperties); - try { msg.compressData(); assertTrue(hasPayload); diff --git a/bmq-sdk/src/test/java/com/bloomberg/bmq/impl/infr/proto/PutMessageIteratorTest.java b/bmq-sdk/src/test/java/com/bloomberg/bmq/impl/infr/proto/PutMessageIteratorTest.java index fc0ff84a..a41de4fc 100644 --- a/bmq-sdk/src/test/java/com/bloomberg/bmq/impl/infr/proto/PutMessageIteratorTest.java +++ b/bmq-sdk/src/test/java/com/bloomberg/bmq/impl/infr/proto/PutMessageIteratorTest.java @@ -66,14 +66,8 @@ void testWithPattern() throws IOException { int i = 0; while (putIt.hasNext()) { // Build expected content - final boolean isOldStyleProperties = i % 2 == 0; - final ByteBufferOutputStream bbos = new ByteBufferOutputStream(); - if (isOldStyleProperties) { - props.streamOutOld(bbos); - } else { - props.streamOut(bbos); - } + props.streamOut(bbos); bbos.writeAscii(PAYLOAD); @@ -84,9 +78,8 @@ void testWithPattern() throws IOException { // Calculate CRC32c before adding padding ByteBuffer[] bb = bbos.reset(); - final long expectedCRC32C = isOldStyleProperties ? 3469549003L : 340340870L; final long CRC32C = Crc32c.calculate(bb); - assertEquals(expectedCRC32C, CRC32C); + assertEquals(340340870L, CRC32C); // Fill content buffer with the payload and the padding for (ByteBuffer b : bb) { @@ -103,8 +96,7 @@ void testWithPattern() throws IOException { assertEquals(QUEUE_ID, putMsg.queueId()); assertEquals(CRC32C, putMsg.crc32c()); assertEquals(FLAGS, putMsg.flags()); - assertEquals(isOldStyleProperties ? 0 : 1, putMsg.header().schemaWireId()); - assertEquals(isOldStyleProperties, putMsg.appData().isOldStyleProperties()); + assertEquals(1, putMsg.header().schemaWireId()); ByteBuffer[] pl = putMsg.appData().applicationData(); ByteBuffer b = @@ -139,8 +131,6 @@ void testWithBuilder() throws IOException { final PutMessageImpl[] puts = new PutMessageImpl[NUM_MSGS]; for (int i = 0; i < NUM_MSGS; i++) { - final boolean isOldStyleProperties = i % 2 == 0; - PutMessageImpl msg = new PutMessageImpl(); msg.setFlags(FLAGS); msg.setQueueId(i); @@ -153,13 +143,11 @@ void testWithBuilder() throws IOException { msg.appData().setPayload(payload); - EventBuilderResult rc = builder.packMessage(msg, isOldStyleProperties); + EventBuilderResult rc = builder.packMessage(msg); assertEquals(EventBuilderResult.SUCCESS, rc); - assertEquals(isOldStyleProperties, msg.appData().isOldStyleProperties()); - assertEquals(isOldStyleProperties ? 0 : 1, msg.header().schemaWireId()); + assertEquals(1, msg.header().schemaWireId()); - final long expectedCRC32C = isOldStyleProperties ? 3469549003L : 340340870L; - assertEquals(expectedCRC32C, msg.header().crc32c()); + assertEquals(340340870L, msg.header().crc32c()); puts[i] = msg; } @@ -182,8 +170,6 @@ void testWithBuilder() throws IOException { assertEquals(exp.flags(), msg.flags()); assertEquals(exp.crc32c(), msg.crc32c()); assertEquals(exp.header().schemaWireId(), msg.header().schemaWireId()); - assertEquals( - exp.appData().isOldStyleProperties(), msg.appData().isOldStyleProperties()); assertEquals(exp.appData().numPaddingBytes(), msg.appData().numPaddingBytes()); assertArrayEquals(exp.appData().applicationData(), msg.appData().applicationData()); diff --git a/bmq-sdk/src/test/java/com/bloomberg/bmq/impl/infr/proto/RequestManagerTest.java b/bmq-sdk/src/test/java/com/bloomberg/bmq/impl/infr/proto/RequestManagerTest.java index 75303b46..1f39fa09 100644 --- a/bmq-sdk/src/test/java/com/bloomberg/bmq/impl/infr/proto/RequestManagerTest.java +++ b/bmq-sdk/src/test/java/com/bloomberg/bmq/impl/infr/proto/RequestManagerTest.java @@ -153,11 +153,6 @@ public GenericResult linger() { return GenericResult.SUCCESS; } - @Override - public boolean isOldStyleMessageProperties() { - return false; - } - @Override public GenericResult write(ByteBuffer[] buffers, boolean waitUntilWritable) { assertNotNull(eventHandler); diff --git a/bmq-sdk/src/test/java/com/bloomberg/bmq/it/BrokerSessionStressIT.java b/bmq-sdk/src/test/java/com/bloomberg/bmq/it/BrokerSessionStressIT.java index 2e3dd78c..c457e8c6 100644 --- a/bmq-sdk/src/test/java/com/bloomberg/bmq/it/BrokerSessionStressIT.java +++ b/bmq-sdk/src/test/java/com/bloomberg/bmq/it/BrokerSessionStressIT.java @@ -101,12 +101,7 @@ public TransferValidator( reader.getState() == QueueState.e_OPENED, "'reader' must be OPENED"); } - public void transfer( - int payloadSize, - int numMsgs, - int numPutsPerEvent, - boolean waitPush, - boolean isOldStyleProperties) { + public void transfer(int payloadSize, int numMsgs, int numPutsPerEvent, boolean waitPush) { if (!putMessages.isEmpty()) { throw new IllegalStateException("'putMessages' expected to be empty"); } @@ -126,14 +121,11 @@ public void transfer( PutMessageImpl[] messages = new PutMessageImpl[minMsg]; for (int i = 0; i < messages.length; i++) { String payload = createPayload(payloadSize, Integer.toString(i)); - PutMessageImpl msg = - TestTools.preparePutMessage(payload, isOldStyleProperties); + PutMessageImpl msg = TestTools.preparePutMessage(payload); logger.debug("Sending {}", msg); putMessages.add(msg); - payloads.add( - TestTools.prepareUnpaddedData( - payload, isOldStyleProperties)); + payloads.add(TestTools.prepareUnpaddedData(payload)); messages[i] = msg; } session.post(writer, messages); @@ -364,12 +356,7 @@ void testThroughput(int payloadSize, int numMsgs, int numPutsPerEvent, boolean w TransferValidator validator = new TransferValidator(eventFIFO, session, queueHandle, queueHandle); - validator.transfer( - payloadSize, - numMsgs, - numPutsPerEvent, - waitPushes, - broker.isOldStyleMessageProperties()); + validator.transfer(payloadSize, numMsgs, numPutsPerEvent, waitPushes); // Close the queue. assertEquals(CloseQueueResult.SUCCESS, queueHandle.close(TEST_REQUEST_TIMEOUT)); @@ -617,12 +604,7 @@ void testFanoutThroughput() throws IOException { TransferValidator validator = new TransferValidator(eventFIFO, session, queueWriterHandle, queueReaderHandle); - validator.transfer( - MSG_SIZE, - NUM_MESSAGES, - NUM_PUTS_PER_EVENT, - WAIT_FOR_PUSHES, - broker.isOldStyleMessageProperties()); + validator.transfer(MSG_SIZE, NUM_MESSAGES, NUM_PUTS_PER_EVENT, WAIT_FOR_PUSHES); // Close the queues. assertEquals(CloseQueueResult.SUCCESS, queueReaderHandle.close(TEST_REQUEST_TIMEOUT)); diff --git a/bmq-sdk/src/test/java/com/bloomberg/bmq/it/NettyProducerIT.java b/bmq-sdk/src/test/java/com/bloomberg/bmq/it/NettyProducerIT.java index 557fd9dd..dd5835b2 100644 --- a/bmq-sdk/src/test/java/com/bloomberg/bmq/it/NettyProducerIT.java +++ b/bmq-sdk/src/test/java/com/bloomberg/bmq/it/NettyProducerIT.java @@ -54,11 +54,7 @@ static QueueImpl createQueue(BrokerSession session, Uri uri, long flags) { return new QueueImpl(session, uri, flags, null, null, null); } - public static void sendMessage( - String[] msgPayloads, - SessionOptions sesOpts, - Uri queueUri, - boolean isOldStyleProperties) { + public static void sendMessage(String[] msgPayloads, SessionOptions sesOpts, Uri queueUri) { Argument.expectNonNull(sesOpts, "sesOpts"); final Duration TEST_REQUEST_TIMEOUT = Duration.ofSeconds(45); @@ -95,8 +91,7 @@ public static void sendMessage( logger.info("Queue opened"); for (String msgPayload : msgPayloads) { - PutMessageImpl message = - TestTools.preparePutMessage(msgPayload, isOldStyleProperties); + PutMessageImpl message = TestTools.preparePutMessage(msgPayload); session.post(qh, message); } @@ -126,11 +121,7 @@ void testProducer() throws IOException { final String MSG = "I'm Netty producer!"; final Uri QUEUE_URI = BmqBroker.Domains.Priority.generateQueueUri(); - sendMessage( - new String[] {MSG}, - broker.sessionOptions(), - QUEUE_URI, - broker.isOldStyleMessageProperties()); + sendMessage(new String[] {MSG}, broker.sessionOptions(), QUEUE_URI); broker.setDropTmpFolder(); } diff --git a/bmq-sdk/src/test/java/com/bloomberg/bmq/it/PayloadIT.java b/bmq-sdk/src/test/java/com/bloomberg/bmq/it/PayloadIT.java index d74521e1..f3852890 100644 --- a/bmq-sdk/src/test/java/com/bloomberg/bmq/it/PayloadIT.java +++ b/bmq-sdk/src/test/java/com/bloomberg/bmq/it/PayloadIT.java @@ -60,9 +60,7 @@ void testPayloadNetty() throws IOException { try (BmqBroker broker = BmqBroker.createStartedBroker()) { final SessionOptions OPTS = broker.sessionOptions(); - final boolean isOldStyleProperties = broker.isOldStyleMessageProperties(); - ByteBuffer unpaddedPayload = - TestTools.prepareUnpaddedData(TEST_MESSAGE, isOldStyleProperties); + ByteBuffer unpaddedPayload = TestTools.prepareUnpaddedData(TEST_MESSAGE); // ================================== // Check netty producer and consumer @@ -71,7 +69,7 @@ void testPayloadNetty() throws IOException { String[] payloads = new String[NUM_MESSAGES]; Arrays.fill(payloads, TEST_MESSAGE); - NettyProducerIT.sendMessage(payloads, OPTS, QUEUE_URI, isOldStyleProperties); + NettyProducerIT.sendMessage(payloads, OPTS, QUEUE_URI); // Read PUSH message but don't confirm it boolean DO_CONFIRM = false; diff --git a/bmq-sdk/src/test/java/com/bloomberg/bmq/it/PlainConsumerIT.java b/bmq-sdk/src/test/java/com/bloomberg/bmq/it/PlainConsumerIT.java index 636aa439..e717533f 100644 --- a/bmq-sdk/src/test/java/com/bloomberg/bmq/it/PlainConsumerIT.java +++ b/bmq-sdk/src/test/java/com/bloomberg/bmq/it/PlainConsumerIT.java @@ -545,7 +545,7 @@ public void testConsumer() throws IOException, InterruptedException { final Uri QUEUE_URI = BmqBroker.Domains.Priority.generateQueueUri(); - PlainProducerIT.sendMessage(MSG, PORT, QUEUE_URI, broker.isOldStyleMessageProperties()); + PlainProducerIT.sendMessage(MSG, PORT, QUEUE_URI); getLastMessage(PORT, QUEUE_URI); diff --git a/bmq-sdk/src/test/java/com/bloomberg/bmq/it/PlainProducerIT.java b/bmq-sdk/src/test/java/com/bloomberg/bmq/it/PlainProducerIT.java index b0e5bdf4..d332232f 100644 --- a/bmq-sdk/src/test/java/com/bloomberg/bmq/it/PlainProducerIT.java +++ b/bmq-sdk/src/test/java/com/bloomberg/bmq/it/PlainProducerIT.java @@ -65,8 +65,7 @@ public class PlainProducerIT { static final Logger logger = LoggerFactory.getLogger(MethodHandles.lookup().lookupClass()); - public static void sendMessage( - String msgPayload, int port, Uri uri, boolean isOldStyleProperties) + public static void sendMessage(String msgPayload, int port, Uri uri) throws IOException, InterruptedException { // =============================== @@ -272,7 +271,7 @@ public static void sendMessage( putMsg.setCorrelationId(); PutEventBuilder putBuilder = new PutEventBuilder(); - EventBuilderResult res = putBuilder.packMessage(putMsg, isOldStyleProperties); + EventBuilderResult res = putBuilder.packMessage(putMsg); assertSame(EventBuilderResult.SUCCESS, res); @@ -439,7 +438,7 @@ public void testProducer() throws IOException, InterruptedException { final int PORT = broker.sessionOptions().brokerUri().getPort(); final Uri QUEUE_URI = BmqBroker.Domains.Priority.generateQueueUri(); - sendMessage(MSG, PORT, QUEUE_URI, broker.isOldStyleMessageProperties()); + sendMessage(MSG, PORT, QUEUE_URI); broker.setDropTmpFolder(); } diff --git a/bmq-sdk/src/test/java/com/bloomberg/bmq/it/SessionIT.java b/bmq-sdk/src/test/java/com/bloomberg/bmq/it/SessionIT.java index f2f4b9a7..c64f10fb 100644 --- a/bmq-sdk/src/test/java/com/bloomberg/bmq/it/SessionIT.java +++ b/bmq-sdk/src/test/java/com/bloomberg/bmq/it/SessionIT.java @@ -2867,8 +2867,6 @@ void testPushProperties() throws BMQException, IOException { try (BmqBroker broker = BmqBroker.createStoppedBroker()) { logger.info("Step 1: Bring up the broker"); - assertFalse(broker.isOldStyleMessageProperties()); - broker.start(); TestSession session = new TestSession(broker.sessionOptions()); @@ -3101,8 +3099,6 @@ void testQueueCompression() throws BMQException, IOException { try (BmqBroker broker = BmqBroker.createStoppedBroker()) { logger.info("Step 1: Bring up the broker"); - assertFalse(broker.isOldStyleMessageProperties()); - broker.start(); TestSession session = new TestSession(broker.sessionOptions()); @@ -3125,8 +3121,8 @@ void testQueueCompression() throws BMQException, IOException { CompressionAlgorithm.None, Protocol.COMPRESSION_MIN_APPDATA_SIZE - 1); - // The message will not be compressed in case message properties - // are new style encoded. + // Message properties are not compressed, so a payload below the + // threshold is sent uncompressed. logger.info( "Step 5: Post incompressable PUT message with Zlib compression, wait for ACK event and PUSH message"); sendVerifyPut( diff --git a/bmq-sdk/src/test/java/com/bloomberg/bmq/it/SubscriptionIT.java b/bmq-sdk/src/test/java/com/bloomberg/bmq/it/SubscriptionIT.java index e12651ed..170f1b6d 100644 --- a/bmq-sdk/src/test/java/com/bloomberg/bmq/it/SubscriptionIT.java +++ b/bmq-sdk/src/test/java/com/bloomberg/bmq/it/SubscriptionIT.java @@ -16,7 +16,6 @@ package com.bloomberg.bmq.it; import static org.junit.jupiter.api.Assertions.assertEquals; -import static org.junit.jupiter.api.Assertions.assertFalse; import static org.junit.jupiter.api.Assertions.assertNotNull; import static org.junit.jupiter.api.Assertions.assertNull; import static org.junit.jupiter.api.Assertions.assertTrue; @@ -377,8 +376,6 @@ void testBreathing() throws BMQException, IOException { try (BmqBroker broker = BmqBroker.createStoppedBroker()) { logger.info("Step 1: Bring up the broker"); - assertFalse(broker.isOldStyleMessageProperties()); - broker.start(); logger.info("Step 2: Start producer/consumer"); @@ -441,8 +438,6 @@ void testFanout() throws BMQException, IOException { try (BmqBroker broker = BmqBroker.createStoppedBroker()) { logger.info("Step 1: Bring up the broker"); - assertFalse(broker.isOldStyleMessageProperties()); - broker.start(); logger.info("Step 2: Start producer"); @@ -508,8 +503,6 @@ void testReuseQueueOptions() throws BMQException, IOException { try (BmqBroker broker = BmqBroker.createStoppedBroker()) { logger.info("Step 1: Bring up the broker"); - assertFalse(broker.isOldStyleMessageProperties()); - broker.start(); logger.info("Step 2: Start producer"); @@ -576,8 +569,6 @@ void testUpdateSubscription() throws BMQException, IOException { try (BmqBroker broker = BmqBroker.createStoppedBroker()) { logger.info("Step 1: Bring up the broker"); - assertFalse(broker.isOldStyleMessageProperties()); - broker.start(); logger.info("Step 2: Start producer/consumer"); @@ -666,8 +657,6 @@ void testStress() throws BMQException, IOException { try (BmqBroker broker = BmqBroker.createStoppedBroker()) { logger.info("Step 1: Bring up the broker"); - assertFalse(broker.isOldStyleMessageProperties()); - broker.start(); logger.info("Step 2: Start producer/consumer"); diff --git a/bmq-sdk/src/test/java/com/bloomberg/bmq/it/TcpBrokerConnectionIT.java b/bmq-sdk/src/test/java/com/bloomberg/bmq/it/TcpBrokerConnectionIT.java index 6096db68..6db5813a 100644 --- a/bmq-sdk/src/test/java/com/bloomberg/bmq/it/TcpBrokerConnectionIT.java +++ b/bmq-sdk/src/test/java/com/bloomberg/bmq/it/TcpBrokerConnectionIT.java @@ -17,7 +17,6 @@ import static org.junit.jupiter.api.Assertions.assertArrayEquals; import static org.junit.jupiter.api.Assertions.assertEquals; -import static org.junit.jupiter.api.Assertions.assertFalse; import static org.junit.jupiter.api.Assertions.assertNotNull; import static org.junit.jupiter.api.Assertions.assertNull; import static org.junit.jupiter.api.Assertions.fail; @@ -52,10 +51,8 @@ import com.bloomberg.bmq.impl.intf.SessionStatusHandler; import com.bloomberg.bmq.impl.intf.SessionStatusHandler.SessionStatus; import com.bloomberg.bmq.it.util.BmqBroker; -import com.bloomberg.bmq.it.util.BmqBrokerContainer; import com.bloomberg.bmq.it.util.BmqBrokerSimulator; import com.bloomberg.bmq.it.util.BmqBrokerSimulator.Mode; -import com.bloomberg.bmq.it.util.TestTcpServer; import com.google.gson.JsonSyntaxException; import java.io.IOException; import java.lang.invoke.MethodHandles; @@ -598,79 +595,6 @@ void testRestart() throws Exception { } } - @FunctionalInterface - interface TestTcpServerFactory { - TestTcpServer create(ConnectionOptions opts) throws IOException; - } - - @Test - void testNegotiationMpsEx() throws IOException { - testNegotiationMpsEx( - new ConnectionOptions().setBrokerUri(getServerUri()), - opts -> new BmqBrokerSimulator(opts.brokerUri().getPort(), Mode.BMQ_AUTO_MODE)); - - testNegotiationMpsEx( - new ConnectionOptions().setBrokerUri(getServerUri()), - opts -> BmqBrokerContainer.createContainer(opts.brokerUri().getPort())); - } - - private void testNegotiationMpsEx(ConnectionOptions opts, TestTcpServerFactory serverFactory) - throws IOException { - final TestTcpServer server = serverFactory.create(opts); - assertFalse(server.isOldStyleMessageProperties()); - - TestSession session = new TestSession(opts); - - // 1) Bring up the server - // 2) Invoke channel 'start' and ensure that it succeeds. - // 3) Wait for start status callback - // 4) Check that the "broker" supports new style message properties - // 5) Linger client session. - // 6) Stop the server. - - logger.info("Start the server."); - server.start(); - - sleepForSeconds(1); - - try { - // 2) Invoke channel 'start' and ensure that it succeeds. - logger.info("Starting channel..."); - - session.start(); - - final int timeout = (int) opts.startAttemptTimeout().getSeconds(); - - // 3) Wait for start status callback. - assertEquals(StartStatus.SUCCESS, session.startStatus(timeout)); - assertEquals(SessionStatus.SESSION_UP, session.sessionStatus()); - - // 4) Check the connection for broker response - logger.info( - "Server: {}, old style properties: {}", - server, - server.isOldStyleMessageProperties()); - - assertEquals( - server.isOldStyleMessageProperties(), - session.channel.isOldStyleMessageProperties()); - - if (server instanceof BmqBroker) { - ((BmqBroker) server).setDropTmpFolder(); - } - } finally { - // 5) Stop client session. - session.stop(); - assertEquals(SessionStatus.SESSION_DOWN, session.sessionStatus()); - - assertEquals(StopStatus.SUCCESS, session.stopStatus()); - assertEquals(GenericResult.SUCCESS, session.linger()); - - // 6) Close the server. - server.close(); - } - } - private void testNegotiationFailed(StatusCategory status) { final int NUM_RETRIES = 1; @@ -1217,10 +1141,7 @@ void testUnknownCompressionType() throws IOException { header.setLength( EventHeader.HEADER_SIZE - + (PushHeader.HEADER_SIZE_FOR_SCHEMA_ID - + unpackedSize - + numPaddingBytes) - * N); + + (PushHeader.HEADER_SIZE + unpackedSize + numPaddingBytes) * N); header.streamOut(bbos); @@ -1231,7 +1152,6 @@ void testUnknownCompressionType() throws IOException { pushMsg.appData().setPayload(ByteBuffer.wrap(bytes)); pushMsg.appData().setProperties(props); - pushMsg.appData().setIsOldStyleProperties(server.isOldStyleMessageProperties()); pushMsg.compressData(); diff --git a/bmq-sdk/src/test/java/com/bloomberg/bmq/it/util/BmqBrokerContainer.java b/bmq-sdk/src/test/java/com/bloomberg/bmq/it/util/BmqBrokerContainer.java index a08fa5a0..2364c416 100644 --- a/bmq-sdk/src/test/java/com/bloomberg/bmq/it/util/BmqBrokerContainer.java +++ b/bmq-sdk/src/test/java/com/bloomberg/bmq/it/util/BmqBrokerContainer.java @@ -306,11 +306,6 @@ public void disableRead() { "'disableRead' not supported for bmqbrkr-based server."); } - @Override - public boolean isOldStyleMessageProperties() { - return false; - } - @Override public SessionOptions sessionOptions() { return sessionOptions; diff --git a/bmq-sdk/src/test/java/com/bloomberg/bmq/it/util/BmqBrokerSimulator.java b/bmq-sdk/src/test/java/com/bloomberg/bmq/it/util/BmqBrokerSimulator.java index d5d69fb2..a0af1557 100644 --- a/bmq-sdk/src/test/java/com/bloomberg/bmq/it/util/BmqBrokerSimulator.java +++ b/bmq-sdk/src/test/java/com/bloomberg/bmq/it/util/BmqBrokerSimulator.java @@ -365,7 +365,7 @@ public void writePushRequest(int qId) throws IOException { pushMsg.appData().setPayload(ByteBuffer.wrap(PAYLOAD.getBytes())); PushEventBuilder builder = new PushEventBuilder(); - builder.packMessage(pushMsg, isOldStyleMessageProperties()); + builder.packMessage(pushMsg); write(builder.build()); } diff --git a/bmq-sdk/src/test/java/com/bloomberg/bmq/it/util/TestTcpServer.java b/bmq-sdk/src/test/java/com/bloomberg/bmq/it/util/TestTcpServer.java index aa81ea62..344a0742 100644 --- a/bmq-sdk/src/test/java/com/bloomberg/bmq/it/util/TestTcpServer.java +++ b/bmq-sdk/src/test/java/com/bloomberg/bmq/it/util/TestTcpServer.java @@ -43,11 +43,6 @@ public interface TestTcpServer extends AutoCloseable { void disableRead(); - // TODO: remove after 2nd rollout of "new style" brokers - default boolean isOldStyleMessageProperties() { - return false; - } - default CompletableFuture startAsync() { return CompletableFuture.runAsync(this::start); } diff --git a/bmq-sdk/src/test/java/com/bloomberg/bmq/it/util/TestTools.java b/bmq-sdk/src/test/java/com/bloomberg/bmq/it/util/TestTools.java index 3430da9a..fb7e24fe 100644 --- a/bmq-sdk/src/test/java/com/bloomberg/bmq/it/util/TestTools.java +++ b/bmq-sdk/src/test/java/com/bloomberg/bmq/it/util/TestTools.java @@ -97,8 +97,7 @@ public static void sleepForMilliSeconds(int millies) { } } - public static ByteBuffer prepareUnpaddedData(String msg, boolean isOldStyleProperties) - throws IOException { + public static ByteBuffer prepareUnpaddedData(String msg) throws IOException { MessagePropertiesImpl props = new MessagePropertiesImpl(); props.setPropertyAsString("routingId", "abcd-efgh-ijkl"); @@ -106,11 +105,7 @@ public static ByteBuffer prepareUnpaddedData(String msg, boolean isOldStylePrope ByteBufferOutputStream bbos = new ByteBufferOutputStream(); - if (isOldStyleProperties) { - props.streamOutOld(bbos); - } else { - props.streamOut(bbos); - } + props.streamOut(bbos); bbos.writeAscii(msg); @@ -150,8 +145,7 @@ public static ByteBuffer mergeBuffers(ByteBuffer[] bbuf) { return bb.flip(); } - public static PutMessageImpl preparePutMessage(String payload, boolean isOldStyleProperties) - throws IOException { + public static PutMessageImpl preparePutMessage(String payload) throws IOException { ByteBuffer b = ByteBuffer.wrap(payload.getBytes()); int putFlags = 0; @@ -168,7 +162,6 @@ public static PutMessageImpl preparePutMessage(String payload, boolean isOldStyl putMsg.appData().setProperties(mp); putMsg.appData().setPayload(b); - putMsg.appData().setIsOldStyleProperties(isOldStyleProperties); putMsg.compressData(); logger.debug("Application data size: {}", putMsg.appData().unpackedSize()); diff --git a/bmq-sdk/src/test/resources/data/bmq_io_dump_1551267643131.bin b/bmq-sdk/src/test/resources/data/bmq_io_dump_1551267643131.bin deleted file mode 100644 index fd76336e..00000000 Binary files a/bmq-sdk/src/test/resources/data/bmq_io_dump_1551267643131.bin and /dev/null differ diff --git a/bmq-sdk/src/test/resources/data/bmq_io_dump_1551267643131.idx b/bmq-sdk/src/test/resources/data/bmq_io_dump_1551267643131.idx deleted file mode 100644 index e8fd8914..00000000 --- a/bmq-sdk/src/test/resources/data/bmq_io_dump_1551267643131.idx +++ /dev/null @@ -1,375 +0,0 @@ -1 400 -1 244 -1 200 -1 144 -1 496 -1 32 -1 144 -1 528 -1 168 -1 536 -1 716 -1 704 -1 704 -1 108 -1 608 -1 120 -1 572 -1 180 -1 512 -1 24 -1 716 -1 292 -1 440 -1 704 -1 108 -1 624 -1 180 -1 552 -1 168 -1 552 -1 732 -1 704 -1 712 -1 396 -1 320 -1 168 -1 496 -1 64 -1 740 -1 168 -1 552 -1 724 -1 168 -1 552 -1 168 -1 552 -1 180 -1 512 -1 32 -1 716 -1 724 -1 496 -1 224 -1 720 -1 684 -1 600 -1 112 -1 724 -1 724 -1 696 -1 144 -1 572 -1 488 -1 216 -1 512 -1 212 -1 696 -1 704 -1 724 -1 132 -1 572 -1 168 -1 544 -1 156 -1 552 -1 168 -1 512 -1 40 -1 704 -1 280 -1 440 -1 108 -1 632 -1 168 -1 544 -1 292 -1 448 -1 168 -1 512 -1 40 -1 716 -1 704 -1 168 -1 528 -1 500 -1 216 -1 392 -1 336 -1 72 -1 512 -1 148 -1 684 -1 672 -1 144 -1 528 -1 36 -1 648 -1 672 -1 144 -1 572 -1 132 -1 528 -1 684 -1 132 -1 512 -1 60 -1 692 -1 684 -1 708 -1 684 -1 660 -1 60 -1 612 -1 660 -1 132 -1 572 -1 672 -1 672 -1 132 -1 572 -1 692 -1 672 -1 144 -1 512 -1 60 -1 660 -1 168 -1 536 -1 96 -1 608 -1 708 -1 144 -1 528 -1 280 -1 432 -1 36 -1 512 -1 168 -1 168 -1 536 -1 180 -1 536 -1 740 -1 708 -1 716 -1 132 -1 528 -1 72 -1 644 -1 36 -1 680 -1 72 -1 512 -1 132 -1 96 -1 608 -1 72 -1 660 -1 96 -1 588 -1 36 -1 676 -1 96 -1 608 -1 144 -1 512 -1 24 -1 156 -1 552 -1 36 -1 244 -1 440 -1 132 -1 528 -1 708 -1 108 -1 512 -1 96 -1 660 -1 36 -1 648 -1 96 -1 608 -1 708 -1 60 -1 612 -1 660 -1 144 -1 528 -1 60 -1 512 -1 100 -1 60 -1 612 -1 660 -1 96 -1 608 -1 84 -1 596 -1 60 -1 620 -1 60 -1 512 -1 108 -1 108 -1 608 -1 60 -1 612 -1 156 -1 528 -1 96 -1 588 -1 96 -1 608 -1 108 -1 512 -1 60 -1 72 -1 620 -1 108 -1 572 -1 660 -1 60 -1 612 -1 132 -1 528 -1 144 -1 564 -1 108 -1 512 -1 96 -1 108 -1 572 -1 96 -1 612 -1 132 -1 572 -1 108 -1 564 -1 144 -1 536 -1 108 -1 512 -1 76 -1 108 -1 608 -1 680 -1 144 -1 536 -1 180 -1 528 -1 144 -1 564 -1 60 -1 512 -1 100 -1 96 -1 608 -1 144 -1 528 -1 144 -1 564 -1 108 -1 608 -1 132 -1 528 -1 132 -1 512 -1 24 -1 144 -1 552 -1 72 -1 620 -1 120 -1 572 -1 108 -1 572 -1 84 -1 596 -1 108 -1 512 -1 60 -1 108 -1 572 -1 60 -1 620 -1 144 -1 528 -1 144 -1 564 -1 72 -1 644 -1 84 -1 512 -1 84 -1 108 -1 572 -1 144 -1 528 -1 120 -1 564 -1 96 -1 608 -1 36 -1 120 -1 512 -1 24 -1 72 -1 444 -1 224 -1 72 -1 656 -1 156 -1 528 -1 180 -1 528 -1 132 -1 512 -1 52 -1 72 -1 644 -1 84 -1 596 -1 132 -1 528 -1 180 -1 528 -1 72 -1 644 -1 36 -1 244 -1 424 -1 144 -1 496 -1 40 -1 156 -1 528 -1 144 -1 564 -1 36 -1 668 -1 108 -1 572 -1 108 -1 512 -1 60 -1 132 -1 536 -1 72 -1 620 -1 132 -1 528 -1 72 -1 644 -1 660 -1 72 -1 644 -1 132 -1 512 -1 16 -1 72 -1 644 -1 200 -1 44 -1 44 diff --git a/bmq-sdk/src/test/resources/data/msg_props_old.bin b/bmq-sdk/src/test/resources/data/msg_props_old.bin deleted file mode 100644 index 843e3e88..00000000 Binary files a/bmq-sdk/src/test/resources/data/msg_props_old.bin and /dev/null differ diff --git a/bmq-sdk/src/test/resources/data/msg_push_multi.bin b/bmq-sdk/src/test/resources/data/msg_push_multi.bin index cc817f2b..29e73a73 100644 Binary files a/bmq-sdk/src/test/resources/data/msg_push_multi.bin and b/bmq-sdk/src/test/resources/data/msg_push_multi.bin differ diff --git a/bmq-sdk/src/test/resources/data/msg_put_multi.bin b/bmq-sdk/src/test/resources/data/msg_put_multi.bin index 71ecc7cb..45d93cd9 100644 Binary files a/bmq-sdk/src/test/resources/data/msg_put_multi.bin and b/bmq-sdk/src/test/resources/data/msg_put_multi.bin differ