From ef07f97029169a52cb1d999a187e9becb5443f8e Mon Sep 17 00:00:00 2001 From: Daniel Collins Date: Fri, 2 Oct 2020 13:27:02 -0400 Subject: [PATCH 1/3] feat: Add nextOffset method to BufferingPullSubscriber --- .../internal/BufferingPullSubscriber.java | 59 ++++++++++++++----- .../pubsublite/internal/PullSubscriber.java | 5 ++ .../internal/BufferingPullSubscriberTest.java | 2 + 3 files changed, 50 insertions(+), 16 deletions(-) diff --git a/google-cloud-pubsublite/src/main/java/com/google/cloud/pubsublite/internal/BufferingPullSubscriber.java b/google-cloud-pubsublite/src/main/java/com/google/cloud/pubsublite/internal/BufferingPullSubscriber.java index 9132bc40f..5735569b7 100644 --- a/google-cloud-pubsublite/src/main/java/com/google/cloud/pubsublite/internal/BufferingPullSubscriber.java +++ b/google-cloud-pubsublite/src/main/java/com/google/cloud/pubsublite/internal/BufferingPullSubscriber.java @@ -18,6 +18,7 @@ import com.google.api.core.ApiService.Listener; import com.google.api.core.ApiService.State; +import com.google.cloud.pubsublite.Offset; import com.google.cloud.pubsublite.SequencedMessage; import com.google.cloud.pubsublite.cloudpubsub.FlowControlSettings; import com.google.cloud.pubsublite.internal.wire.Subscriber; @@ -25,19 +26,29 @@ import com.google.cloud.pubsublite.proto.FlowControlRequest; import com.google.cloud.pubsublite.proto.SeekRequest; import com.google.cloud.pubsublite.proto.SeekRequest.NamedTarget; +import com.google.common.collect.ImmutableList; +import com.google.common.collect.Iterables; import com.google.common.util.concurrent.MoreExecutors; +import com.google.errorprone.annotations.concurrent.GuardedBy; import io.grpc.StatusException; -import java.util.ArrayList; +import java.util.ArrayDeque; +import java.util.Collection; +import java.util.Deque; import java.util.List; +import java.util.Optional; import java.util.concurrent.ExecutionException; -import java.util.concurrent.LinkedBlockingQueue; -import java.util.concurrent.atomic.AtomicReference; -import javax.annotation.Nullable; public class BufferingPullSubscriber implements PullSubscriber { private final Subscriber underlying; - private final AtomicReference error = new AtomicReference<>(); - private final LinkedBlockingQueue messages = new LinkedBlockingQueue<>(); + + @GuardedBy("this") + private Optional error = Optional.empty(); + + @GuardedBy("this") + private Deque messages = new ArrayDeque<>(); + + @GuardedBy("this") + private Optional lastDelivered = Optional.empty(); public BufferingPullSubscriber(SubscriberFactory factory, FlowControlSettings settings) throws StatusException { @@ -50,12 +61,12 @@ public BufferingPullSubscriber(SubscriberFactory factory, FlowControlSettings se public BufferingPullSubscriber( SubscriberFactory factory, FlowControlSettings settings, SeekRequest initialSeek) throws StatusException { - underlying = factory.New(messages::addAll); + underlying = factory.New(this::addMessages); underlying.addListener( new Listener() { @Override public void failed(State state, Throwable throwable) { - error.set(ExtractStatus.toCanonical(throwable)); + fail(ExtractStatus.toCanonical(throwable)); } }, MoreExecutors.directExecutor()); @@ -74,25 +85,41 @@ public void failed(State state, Throwable throwable) { .build()); } + private synchronized void fail(StatusException e) { + error = Optional.of(e); + } + + private synchronized void addMessages(Collection new_messages) { + messages.addAll(new_messages); + } + @Override - public List pull() throws StatusException { - @Nullable StatusException maybeError = error.get(); - if (maybeError != null) { - throw maybeError; + public synchronized List pull() throws StatusException { + if (error.isPresent()) { + throw error.get(); + } + if (messages.isEmpty()) { + return ImmutableList.of(); } - ArrayList collection = new ArrayList<>(); - messages.drainTo(collection); + Deque collection = messages; + messages = new ArrayDeque<>(); long bytes = collection.stream().mapToLong(SequencedMessage::byteSize).sum(); underlying.allowFlow( FlowControlRequest.newBuilder() .setAllowedBytes(bytes) .setAllowedMessages(collection.size()) .build()); - return collection; + lastDelivered = Optional.of(Iterables.getLast(collection).offset()); + return ImmutableList.copyOf(collection); + } + + @Override + public synchronized Optional nextOffset() { + return lastDelivered.map(offset -> Offset.of(offset.value() + 1)); } @Override public void close() { underlying.stopAsync().awaitTerminated(); } -} +} \ No newline at end of file diff --git a/google-cloud-pubsublite/src/main/java/com/google/cloud/pubsublite/internal/PullSubscriber.java b/google-cloud-pubsublite/src/main/java/com/google/cloud/pubsublite/internal/PullSubscriber.java index bfe6d8783..9423f062b 100644 --- a/google-cloud-pubsublite/src/main/java/com/google/cloud/pubsublite/internal/PullSubscriber.java +++ b/google-cloud-pubsublite/src/main/java/com/google/cloud/pubsublite/internal/PullSubscriber.java @@ -16,11 +16,16 @@ package com.google.cloud.pubsublite.internal; +import com.google.cloud.pubsublite.Offset; import io.grpc.StatusException; import java.util.List; +import java.util.Optional; // A PullSubscriber exposes a "pull" mechanism for retrieving messages. public interface PullSubscriber extends AutoCloseable { // Pull currently available messages from this subscriber. Does not block. List pull() throws StatusException; + + // The next offset expected to be returned by this PullSubscriber, or empty if unknown. + Optional nextOffset(); } diff --git a/google-cloud-pubsublite/src/test/java/com/google/cloud/pubsublite/internal/BufferingPullSubscriberTest.java b/google-cloud-pubsublite/src/test/java/com/google/cloud/pubsublite/internal/BufferingPullSubscriberTest.java index 2f8b30343..e6ecdcfe2 100644 --- a/google-cloud-pubsublite/src/test/java/com/google/cloud/pubsublite/internal/BufferingPullSubscriberTest.java +++ b/google-cloud-pubsublite/src/test/java/com/google/cloud/pubsublite/internal/BufferingPullSubscriberTest.java @@ -154,9 +154,11 @@ public void multipleBatchesAggregatedReturnsTokens() throws StatusException { SequencedMessage.of(Message.builder().build(), Timestamps.EPOCH, Offset.of(11), 20); SequencedMessage message3 = SequencedMessage.of(Message.builder().build(), Timestamps.EPOCH, Offset.of(12), 30); + assertThat(subscriber.nextOffset()).isEmpty(); messageConsumer.accept(ImmutableList.of(message1, message2)); messageConsumer.accept(ImmutableList.of(message3)); assertThat(subscriber.pull()).containsExactly(message1, message2, message3); + assertThat(subscriber.nextOffset()).hasValue(Offset.of(13)); assertThat(subscriber.pull()).isEmpty(); FlowControlRequest flowControlRequest = From 92d477980c856d98237ee7a34c969dfdac79afb2 Mon Sep 17 00:00:00 2001 From: dpcollins-google <40498610+dpcollins-google@users.noreply.github.com> Date: Mon, 5 Oct 2020 10:34:33 -0400 Subject: [PATCH 2/3] chore: add newline to EOF --- .../cloud/pubsublite/internal/BufferingPullSubscriber.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/google-cloud-pubsublite/src/main/java/com/google/cloud/pubsublite/internal/BufferingPullSubscriber.java b/google-cloud-pubsublite/src/main/java/com/google/cloud/pubsublite/internal/BufferingPullSubscriber.java index 5735569b7..e425b768d 100644 --- a/google-cloud-pubsublite/src/main/java/com/google/cloud/pubsublite/internal/BufferingPullSubscriber.java +++ b/google-cloud-pubsublite/src/main/java/com/google/cloud/pubsublite/internal/BufferingPullSubscriber.java @@ -122,4 +122,4 @@ public synchronized Optional nextOffset() { public void close() { underlying.stopAsync().awaitTerminated(); } -} \ No newline at end of file +} From 2330eb1ead8893353bb0047d577cf63b2ff4c456 Mon Sep 17 00:00:00 2001 From: dpcollins-google <40498610+dpcollins-google@users.noreply.github.com> Date: Mon, 5 Oct 2020 10:35:25 -0400 Subject: [PATCH 3/3] chore: Update PullSubscriber docs. --- .../com/google/cloud/pubsublite/internal/PullSubscriber.java | 1 + 1 file changed, 1 insertion(+) diff --git a/google-cloud-pubsublite/src/main/java/com/google/cloud/pubsublite/internal/PullSubscriber.java b/google-cloud-pubsublite/src/main/java/com/google/cloud/pubsublite/internal/PullSubscriber.java index 9423f062b..20438a080 100644 --- a/google-cloud-pubsublite/src/main/java/com/google/cloud/pubsublite/internal/PullSubscriber.java +++ b/google-cloud-pubsublite/src/main/java/com/google/cloud/pubsublite/internal/PullSubscriber.java @@ -27,5 +27,6 @@ public interface PullSubscriber extends AutoCloseable { List pull() throws StatusException; // The next offset expected to be returned by this PullSubscriber, or empty if unknown. + // Subsequent messages are guaranteed to have offsets of at least this value. Optional nextOffset(); }