Skip to content
Merged
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 @@ -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;
Expand Down Expand Up @@ -208,6 +210,18 @@ default ResultStream<StreamEventsResponse> stream(StreamEventsRequest request, i
*/
ResultStream<SourceEventsResponse> 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<SnapshottedSourceEventsResponse> source(SnapshottedSourceRequest request);

/**
* Provides tags for an event at the given global sequence.
*
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand All @@ -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;
Expand All @@ -81,7 +83,9 @@ public class DcbEventChannelImpl extends AbstractAxonServerChannel<Void> 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<ResultStream<StreamEventsResponse>> buffers = ConcurrentHashMap.newKeySet();
Expand All @@ -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;
}
Expand Down Expand Up @@ -171,9 +176,9 @@ protected StreamEventsResponse terminalMessage() {
@Override
public ResultStream<SourceEventsResponse> source(SourceEventsRequest request) {
AbstractBufferedStream<SourceEventsResponse, Empty> result =
new AbstractBufferedStream<SourceEventsResponse, Empty>(clientIdentification.getClientId(),
BUFFER_SIZE,
REFILL_BATCH) {
new AbstractBufferedStream<>(clientIdentification.getClientId(),
BUFFER_SIZE,
REFILL_BATCH) {
@Override
protected SourceEventsResponse terminalMessage() {
return SourceEventsResponse.newBuilder()
Expand All @@ -189,6 +194,27 @@ protected Empty buildFlowControlMessage(FlowControl flowControl) {
return result;
}

@Override
public ResultStream<SnapshottedSourceEventsResponse> source(SnapshottedSourceRequest request) {
AbstractBufferedStream<SnapshottedSourceEventsResponse, Empty> 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<GetTagsResponse> tagsFor(long sequence) {
FutureStreamObserver<GetTagsResponse> future = new FutureStreamObserver<>(null);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -8,21 +8,29 @@
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;
import io.axoniq.axonserver.grpc.event.dcb.StreamEventsResponse;
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;
Expand Down Expand Up @@ -78,7 +86,7 @@
private AxonServerConnectionFactory client;

@BeforeEach
void setUp() {

Check warning on line 89 in src/test/java/io/axoniq/axonserver/connector/event/dcb/DcbEndToEndTest.java

View check run for this annotation

SonarQubeCloud / SonarCloud Code Analysis

Only one method in a class should be annotated @BeforeEach.

See more on https://sonarcloud.io/project/issues?id=AxonIQ_axonserver-connector-java&issues=AaA4QTrSnaLCuzE2Tufm&open=AaA4QTrSnaLCuzE2Tufm&pullRequest=512
AxonServerConnectionFactory.Builder builder = AxonServerConnectionFactory.forClient("dcb-e2e-test")
.connectTimeout(1500,
TimeUnit.MILLISECONDS)
Expand Down Expand Up @@ -401,6 +409,72 @@
.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();
Expand Down Expand Up @@ -710,7 +784,7 @@
.get(10, TimeUnit.SECONDS);
assertEquals(0, response.getSequence(), "Empty store should return sequence 0");
} catch (InterruptedException | ExecutionException | TimeoutException e) {
fail("getSequenceAt operation timed out or failed: " + e.getMessage());

Check warning on line 787 in src/test/java/io/axoniq/axonserver/connector/event/dcb/DcbEndToEndTest.java

View check run for this annotation

SonarQubeCloud / SonarCloud Code Analysis

Use assertDoesNotThrow() instead of try/catch and fail() in the catch block.

See more on https://sonarcloud.io/project/issues?id=AxonIQ_axonserver-connector-java&issues=AaA4QTrSnaLCuzE2Tufn&open=AaA4QTrSnaLCuzE2Tufn&pullRequest=512
}
}

Expand Down Expand Up @@ -752,7 +826,7 @@
assertEquals(tail, response.getSequence(),
"Timestamp before all events should return the tail sequence");
} catch (InterruptedException | ExecutionException | TimeoutException e) {
fail("getSequenceAt operation timed out or failed: " + e.getMessage());

Check warning on line 829 in src/test/java/io/axoniq/axonserver/connector/event/dcb/DcbEndToEndTest.java

View check run for this annotation

SonarQubeCloud / SonarCloud Code Analysis

Use assertDoesNotThrow() instead of try/catch and fail() in the catch block.

See more on https://sonarcloud.io/project/issues?id=AxonIQ_axonserver-connector-java&issues=AaA4QTrSnaLCuzE2Tufo&open=AaA4QTrSnaLCuzE2Tufo&pullRequest=512
}
} catch (Exception e) {
// If we can't retrieve head or append events, the test environment might not be properly set up
Expand Down Expand Up @@ -797,7 +871,7 @@
assertEquals(head, response.getSequence(),
"Timestamp after all events should return the head sequence");
} catch (InterruptedException | ExecutionException | TimeoutException e) {
fail("getSequenceAt operation timed out or failed: " + e.getMessage());

Check warning on line 874 in src/test/java/io/axoniq/axonserver/connector/event/dcb/DcbEndToEndTest.java

View check run for this annotation

SonarQubeCloud / SonarCloud Code Analysis

Use assertDoesNotThrow() instead of try/catch and fail() in the catch block.

See more on https://sonarcloud.io/project/issues?id=AxonIQ_axonserver-connector-java&issues=AaA4QTrSnaLCuzE2Tufp&open=AaA4QTrSnaLCuzE2Tufp&pullRequest=512
}
} catch (Exception e) {
// If we can't append events or retrieve head, the test environment might not be properly set up
Expand Down Expand Up @@ -845,7 +919,7 @@
assertEquals(sequences.get(middleIndex), response.getSequence(),
"Timestamp exactly matching an event should return that event's sequence");
} catch (InterruptedException | ExecutionException | TimeoutException e) {
fail("getSequenceAt operation timed out or failed: " + e.getMessage());

Check warning on line 922 in src/test/java/io/axoniq/axonserver/connector/event/dcb/DcbEndToEndTest.java

View check run for this annotation

SonarQubeCloud / SonarCloud Code Analysis

Use assertDoesNotThrow() instead of try/catch and fail() in the catch block.

See more on https://sonarcloud.io/project/issues?id=AxonIQ_axonserver-connector-java&issues=AaA4QTrSnaLCuzE2Tufq&open=AaA4QTrSnaLCuzE2Tufq&pullRequest=512
}
} catch (Exception e) {
// If we can't retrieve head or append events, the test environment might not be properly set up
Expand Down Expand Up @@ -892,7 +966,7 @@
assertEquals(sequences.get(1), response.getSequence(),
"Timestamp between events should return the sequence of the previous event");
} catch (InterruptedException | ExecutionException | TimeoutException e) {
fail("getSequenceAt operation timed out or failed: " + e.getMessage());

Check warning on line 969 in src/test/java/io/axoniq/axonserver/connector/event/dcb/DcbEndToEndTest.java

View check run for this annotation

SonarQubeCloud / SonarCloud Code Analysis

Use assertDoesNotThrow() instead of try/catch and fail() in the catch block.

See more on https://sonarcloud.io/project/issues?id=AxonIQ_axonserver-connector-java&issues=AaA4QTrSnaLCuzE2Tufr&open=AaA4QTrSnaLCuzE2Tufr&pullRequest=512
}
} catch (Exception e) {
// If we can't retrieve head or append events, the test environment might not be properly set up
Expand Down Expand Up @@ -1034,6 +1108,24 @@
.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();
}
Expand Down
Loading