diff --git a/src/main/java/io/axoniq/axonserver/connector/event/DcbEventChannel.java b/src/main/java/io/axoniq/axonserver/connector/event/DcbEventChannel.java index 75d80f82..a9a764a2 100644 --- a/src/main/java/io/axoniq/axonserver/connector/event/DcbEventChannel.java +++ b/src/main/java/io/axoniq/axonserver/connector/event/DcbEventChannel.java @@ -27,6 +27,8 @@ import io.axoniq.axonserver.grpc.event.dcb.GetTagsResponse; import io.axoniq.axonserver.grpc.event.dcb.GetTailResponse; import io.axoniq.axonserver.grpc.event.dcb.RemoveTagsResponse; +import io.axoniq.axonserver.grpc.event.dcb.SnapshottedSourceEventsResponse; +import io.axoniq.axonserver.grpc.event.dcb.SnapshottedSourceRequest; import io.axoniq.axonserver.grpc.event.dcb.SourceEventsRequest; import io.axoniq.axonserver.grpc.event.dcb.SourceEventsResponse; import io.axoniq.axonserver.grpc.event.dcb.StreamEventsRequest; @@ -208,6 +210,18 @@ default ResultStream stream(StreamEventsRequest request, i */ ResultStream source(SourceEventsRequest request); + /** + * Provides a finite stream of events used to source a model, optionally preceded by the latest snapshot stored + * under the snapshot key in the given {@code request}. When a snapshot is present, it is emitted first, followed by + * events with a sequence greater than the snapshot's sequence. When no snapshot is present, this behaves like + * {@link #source(SourceEventsRequest)}, sourcing from the beginning. + * + * @param request the query used to filter events for sourcing and identify the snapshot to prefix the stream with + * @return the response containing an optional snapshot, events to source a model, and a consistency marker to be + * used when trying to append new events to the event store + */ + ResultStream source(SnapshottedSourceRequest request); + /** * Provides tags for an event at the given global sequence. * diff --git a/src/main/java/io/axoniq/axonserver/connector/event/impl/DcbEventChannelImpl.java b/src/main/java/io/axoniq/axonserver/connector/event/impl/DcbEventChannelImpl.java index 98e2b21e..480bf9a5 100644 --- a/src/main/java/io/axoniq/axonserver/connector/event/impl/DcbEventChannelImpl.java +++ b/src/main/java/io/axoniq/axonserver/connector/event/impl/DcbEventChannelImpl.java @@ -50,6 +50,9 @@ import io.axoniq.axonserver.grpc.event.dcb.RescheduleEventRequest; import io.axoniq.axonserver.grpc.event.dcb.ScheduleEventRequest; import io.axoniq.axonserver.grpc.event.dcb.ScheduleToken; +import io.axoniq.axonserver.grpc.event.dcb.SnapshottedDcbEventStoreGrpc; +import io.axoniq.axonserver.grpc.event.dcb.SnapshottedSourceEventsResponse; +import io.axoniq.axonserver.grpc.event.dcb.SnapshottedSourceRequest; import io.axoniq.axonserver.grpc.event.dcb.SourceEventsRequest; import io.axoniq.axonserver.grpc.event.dcb.SourceEventsResponse; import io.axoniq.axonserver.grpc.event.dcb.StreamEventsRequest; @@ -61,7 +64,6 @@ import io.grpc.stub.StreamObserver; import java.time.Instant; -import java.util.Arrays; import java.util.Collection; import java.util.Set; import java.util.concurrent.CompletableFuture; @@ -81,7 +83,9 @@ public class DcbEventChannelImpl extends AbstractAxonServerChannel impleme private static final int BUFFER_SIZE = 512; private static final int REFILL_BATCH = 16; + private final DcbEventStoreGrpc.DcbEventStoreStub eventStore; + private final SnapshottedDcbEventStoreGrpc.SnapshottedDcbEventStoreStub snapshottedEventStore; private final DcbEventSchedulerGrpc.DcbEventSchedulerStub eventScheduler; private final ClientIdentification clientIdentification; private final Set> buffers = ConcurrentHashMap.newKeySet(); @@ -98,6 +102,7 @@ public DcbEventChannelImpl(ClientIdentification clientIdentification, AxonServerManagedChannel axonServerManagedChannel) { super(clientIdentification, executor, axonServerManagedChannel); this.eventStore = DcbEventStoreGrpc.newStub(axonServerManagedChannel); + this.snapshottedEventStore = SnapshottedDcbEventStoreGrpc.newStub(axonServerManagedChannel); this.eventScheduler = DcbEventSchedulerGrpc.newStub(axonServerManagedChannel); this.clientIdentification = clientIdentification; } @@ -171,9 +176,9 @@ protected StreamEventsResponse terminalMessage() { @Override public ResultStream source(SourceEventsRequest request) { AbstractBufferedStream result = - new AbstractBufferedStream(clientIdentification.getClientId(), - BUFFER_SIZE, - REFILL_BATCH) { + new AbstractBufferedStream<>(clientIdentification.getClientId(), + BUFFER_SIZE, + REFILL_BATCH) { @Override protected SourceEventsResponse terminalMessage() { return SourceEventsResponse.newBuilder() @@ -189,6 +194,27 @@ protected Empty buildFlowControlMessage(FlowControl flowControl) { return result; } + @Override + public ResultStream source(SnapshottedSourceRequest request) { + AbstractBufferedStream result = + new AbstractBufferedStream<>( + clientIdentification.getClientId(), BUFFER_SIZE, REFILL_BATCH + ) { + @Override + protected SnapshottedSourceEventsResponse terminalMessage() { + return SnapshottedSourceEventsResponse.newBuilder() + .build(); + } + + @Override + protected Empty buildFlowControlMessage(FlowControl flowControl) { + return null; + } + }; + snapshottedEventStore.source(request, result); + return result; + } + @Override public CompletableFuture tagsFor(long sequence) { FutureStreamObserver future = new FutureStreamObserver<>(null); diff --git a/src/test/java/io/axoniq/axonserver/connector/event/dcb/DcbEndToEndTest.java b/src/test/java/io/axoniq/axonserver/connector/event/dcb/DcbEndToEndTest.java index d1b3381f..784299f5 100644 --- a/src/test/java/io/axoniq/axonserver/connector/event/dcb/DcbEndToEndTest.java +++ b/src/test/java/io/axoniq/axonserver/connector/event/dcb/DcbEndToEndTest.java @@ -8,14 +8,20 @@ import io.axoniq.axonserver.connector.ResultStream; import io.axoniq.axonserver.connector.ResultStreamPublisher; import io.axoniq.axonserver.connector.event.DcbEventChannel; +import io.axoniq.axonserver.connector.impl.HeaderAttachingInterceptor; +import io.axoniq.axonserver.connector.impl.Headers; import io.axoniq.axonserver.connector.impl.ServerAddress; +import io.axoniq.axonserver.grpc.event.dcb.AddSnapshotRequest; import io.axoniq.axonserver.grpc.event.dcb.AppendEventsResponse; import io.axoniq.axonserver.grpc.event.dcb.ConsistencyCondition; import io.axoniq.axonserver.grpc.event.dcb.Criterion; +import io.axoniq.axonserver.grpc.event.dcb.DcbSnapshotStoreGrpc; import io.axoniq.axonserver.grpc.event.dcb.Event; import io.axoniq.axonserver.grpc.event.dcb.GetSequenceAtResponse; import io.axoniq.axonserver.grpc.event.dcb.GetTagsResponse; import io.axoniq.axonserver.grpc.event.dcb.GetTailResponse; +import io.axoniq.axonserver.grpc.event.dcb.Snapshot; +import io.axoniq.axonserver.grpc.event.dcb.SnapshottedSourceRequest; import io.axoniq.axonserver.grpc.event.dcb.SourceEventsRequest; import io.axoniq.axonserver.grpc.event.dcb.SourceEventsResponse; import io.axoniq.axonserver.grpc.event.dcb.StreamEventsRequest; @@ -23,6 +29,8 @@ import io.axoniq.axonserver.grpc.event.dcb.Tag; import io.axoniq.axonserver.grpc.event.dcb.TaggedEvent; import io.axoniq.axonserver.grpc.event.dcb.TagsAndNamesCriterion; +import io.grpc.ManagedChannel; +import io.grpc.ManagedChannelBuilder; import org.junit.jupiter.api.*; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -401,6 +409,72 @@ void sourceSingleTagAndEventName() { .verifyComplete(); } + @Test + void sourceWithSnapshotWhenNoneStored() { + DcbEventChannel dcbEventChannel = connection.dcbEventChannel(); + Tag tag = aTag(); + String eventName = aString(); + + long head = retrieveHead(); + TaggedEvent taggedEvent = taggedEvent(anEvent(aString(), eventName), tag); + appendEvent(taggedEvent); + + TagsAndNamesCriterion tagsAndName = TagsAndNamesCriterion.newBuilder().addTag(tag).addName(eventName).build(); + Criterion criterion = Criterion.newBuilder().setTagsAndNames(tagsAndName).build(); + SnapshottedSourceRequest request = SnapshottedSourceRequest.newBuilder() + .setSnapshotKey(ByteString.copyFromUtf8(aString())) + .addCriterion(criterion) + .build(); + + StepVerifier.create(new ResultStreamPublisher<>(() -> dcbEventChannel.source(request))) + .expectNextMatches( + r -> !r.hasSnapshot() + && r.getEvent().getEvent().equals(taggedEvent.getEvent()) + && r.getEvent().getSequence() == head + ) + .expectNextMatches(r -> r.getConsistencyMarker() == head + 1L) + .verifyComplete(); + } + + @Test + void sourceWithSnapshotWhenPresent() { + DcbEventChannel dcbEventChannel = connection.dcbEventChannel(); + Tag tag = aTag(); + String eventName = aString(); + ByteString snapshotKey = ByteString.copyFromUtf8(aString()); + + long head = retrieveHead(); + TaggedEvent beforeSnapshot = taggedEvent(anEvent(aString(), eventName), tag); + long snapshotSequence = appendEvent(beforeSnapshot).getSequenceOfTheFirstEvent(); + + Snapshot snapshot = Snapshot.newBuilder() + .setName("snapshot-name") + .setVersion("0.0.1") + .setPayload(ByteString.copyFromUtf8("snapshot-payload")) + .setTimestamp(Instant.now().toEpochMilli()) + .build(); + addSnapshot(snapshotKey, snapshotSequence, snapshot); + + TaggedEvent afterSnapshot = taggedEvent(anEvent(aString(), eventName), tag); + appendEvent(afterSnapshot); + + TagsAndNamesCriterion tagsAndNames = TagsAndNamesCriterion.newBuilder().addTag(tag).addName(eventName).build(); + Criterion criterion = Criterion.newBuilder().setTagsAndNames(tagsAndNames).build(); + SnapshottedSourceRequest request = SnapshottedSourceRequest.newBuilder() + .setSnapshotKey(snapshotKey) + .addCriterion(criterion) + .build(); + + StepVerifier.create(new ResultStreamPublisher<>(() -> dcbEventChannel.source(request))) + .expectNextMatches(r -> r.hasSnapshot() && r.getSnapshot().equals(snapshot)) + .expectNextMatches( + r -> r.getEvent().getEvent().equals(afterSnapshot.getEvent()) + && r.getEvent().getSequence() > snapshotSequence + ) + .expectNextMatches(r -> r.getConsistencyMarker() == head + 2L) + .verifyComplete(); + } + @Test void noConditionAppend() { long head = retrieveHead(); @@ -1034,6 +1108,24 @@ private long retrieveHead() { .getSequence(); } + private void addSnapshot(ByteString key, long sequence, Snapshot snapshot) { + ManagedChannel channel = + ManagedChannelBuilder.forAddress(axonServerAddress.getHostName(), axonServerAddress.getGrpcPort()) + .usePlaintext() + .intercept(new HeaderAttachingInterceptor<>(Headers.CONTEXT, "default")) + .build(); + try { + AddSnapshotRequest snapshotRequest = AddSnapshotRequest.newBuilder() + .setKey(key) + .setSequence(sequence) + .setSnapshot(snapshot) + .build(); + DcbSnapshotStoreGrpc.newBlockingStub(channel).add(snapshotRequest); + } finally { + channel.shutdownNow(); + } + } + private AppendEventsResponse appendEvent(TaggedEvent taggedEvent) { return appendEventAsync(taggedEvent).join(); }