Skip to content
Permalink
Browse files
feat: adding support for dead letter queues (#60)
* feat: Adding support for DLQs

Adding delivery attempt count to PubsubMessages as a message attribute,
and creating helper function to allow users to get the count without
knowing implementation details.

* Fix formatting

* fix: making changes requested in pull request
  • Loading branch information
hannahrogers-google committed Jan 14, 2020
1 parent 1b8405a commit f3c93fa8bf0eb8ebda6eea6c6c6a60a36dc69af2
@@ -341,10 +341,17 @@ private void processBatch(List<OutstandingMessage> batch) {
// This should be a blocking flow controller and never throw an exception.
throw new IllegalStateException("Flow control unexpected exception", unexpectedException);
}
processOutstandingMessage(message.receivedMessage.getMessage(), message.ackHandler);
processOutstandingMessage(addDeliveryInfoCount(message.receivedMessage), message.ackHandler);
}
}

private PubsubMessage addDeliveryInfoCount(ReceivedMessage receivedMessage) {
return PubsubMessage.newBuilder(receivedMessage.getMessage())
.putAttributes(
"googclient_deliveryattempt", Integer.toString(receivedMessage.getDeliveryAttempt()))
.build();
}

private void processOutstandingMessage(final PubsubMessage message, final AckHandler ackHandler) {
final SettableApiFuture<AckReply> response = SettableApiFuture.create();
final AckReplyConsumer consumer =
@@ -44,6 +44,7 @@
import com.google.common.base.Preconditions;
import com.google.common.util.concurrent.ThreadFactoryBuilder;
import com.google.pubsub.v1.ProjectSubscriptionName;
import com.google.pubsub.v1.PubsubMessage;
import java.io.IOException;
import java.util.ArrayList;
import java.util.List;
@@ -205,6 +206,11 @@ public static Builder newBuilder(String subscription, MessageReceiver receiver)
return new Builder(subscription, receiver);
}

/** Returns the delivery attempt count for a received {@link PubsubMessage} */
public static int getDeliveryAttempt(PubsubMessage message) {
return Integer.parseInt(message.getAttributesOrDefault("googclient_deliveryattempt", "0"));
}

/** Subscription which the subscriber is subscribed to. */
public String getSubscriptionNameString() {
return subscriptionName;
@@ -37,10 +37,13 @@
import org.threeten.bp.Duration;

public class MessageDispatcherTest {
private static final ByteString MESSAGE_DATA = ByteString.copyFromUtf8("message-data");
private static final int DELIVERY_INFO_COUNT = 3;
private static final ReceivedMessage TEST_MESSAGE =
ReceivedMessage.newBuilder()
.setAckId("ackid")
.setMessage(PubsubMessage.newBuilder().setData(ByteString.EMPTY).build())
.setMessage(PubsubMessage.newBuilder().setData(MESSAGE_DATA).build())
.setDeliveryAttempt(DELIVERY_INFO_COUNT)
.build();
private static final Runnable NOOP_RUNNABLE =
new Runnable() {
@@ -78,6 +81,9 @@ public void setUp() {
new MessageReceiver() {
@Override
public void receiveMessage(final PubsubMessage message, final AckReplyConsumer consumer) {
assertThat(message.getData()).isEqualTo(MESSAGE_DATA);
assertThat(message.getAttributesOrThrow("googclient_deliveryattempt"))
.isEqualTo(Integer.toString(DELIVERY_INFO_COUNT));
consumers.add(consumer);
}
};
@@ -84,6 +84,19 @@ public void tearDown() throws Exception {
testChannel.shutdown();
}

@Test
public void testDeliveryAttemptHelper() {
int deliveryAttempt = 3;
PubsubMessage message =
PubsubMessage.newBuilder()
.putAttributes("googclient_deliveryattempt", Integer.toString(deliveryAttempt))
.build();
assertEquals(Subscriber.getDeliveryAttempt(message), deliveryAttempt);

PubsubMessage emptyMessage = PubsubMessage.newBuilder().build();
assertEquals(Subscriber.getDeliveryAttempt(emptyMessage), 0);
}

@Test
public void testOpenedChannels() throws Exception {
int expectedChannelCount = 1;

0 comments on commit f3c93fa

Please sign in to comment.