Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -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);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -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");
Expand Down Expand Up @@ -792,11 +781,6 @@ public GenericResult linger() {
return GenericResult.SUCCESS;
}

@Override
public boolean isOldStyleMessageProperties() {
return isOldStyleMessageProperties;
}

@Override
public GenericResult write(ByteBuffer[] buffers, boolean waitUntilWritable) {

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand All @@ -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 {
Expand Down Expand Up @@ -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?

Expand All @@ -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;
}

Expand All @@ -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();
Expand All @@ -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 {
Expand All @@ -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
Expand All @@ -221,7 +183,6 @@ public void streamIn(

compressedData = bbos;
this.compressionType = compressionType;
arePropertiesCompressed = hasProperties && isOldStyleProperties;
}
}

Expand All @@ -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) {
Expand All @@ -273,16 +227,6 @@ private int streamInProperties(ByteBufferInputStream input) throws IOException {
return read;
}

private <T extends InputStream & DataInput> 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];

Expand Down Expand Up @@ -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 {
Expand All @@ -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) {
Expand Down
Loading
Loading