GH-806 - Add archive support for Neo4j.

Co-authored-by: Oliver Drotbohm <oliver.drotbohm@broadcom.com>
This commit is contained in:
Cora Iberkleid
2024-10-24 13:40:50 -04:00
committed by Oliver Drotbohm
parent d302ad1937
commit 4132ad3a38
2 changed files with 326 additions and 249 deletions

View File

@@ -15,6 +15,8 @@
*/
package org.springframework.modulith.events.neo4j;
import static org.neo4j.cypherdsl.core.Cypher.*;
import java.time.Instant;
import java.time.ZoneOffset;
import java.util.ArrayList;
@@ -23,6 +25,7 @@ import java.util.Map;
import java.util.Objects;
import java.util.Optional;
import java.util.UUID;
import java.util.function.Function;
import org.neo4j.cypherdsl.core.Cypher;
import org.neo4j.cypherdsl.core.Node;
@@ -48,6 +51,7 @@ import org.springframework.util.DigestUtils;
*
* @author Gerrit Meier
* @author Oliver Drotbohm
* @author Cora Iberkleid
* @since 1.1
*/
@Transactional
@@ -61,78 +65,103 @@ class Neo4jEventPublicationRepository implements EventPublicationRepository {
private static final String PUBLICATION_DATE = "publicationDate";
private static final String COMPLETION_DATE = "completionDate";
private static final Node EVENT_PUBLICATION_NODE = Cypher.node("Neo4jEventPublication")
private static final Node EVENT_PUBLICATION_NODE = node("Neo4jEventPublication")
.named("neo4jEventPublication");
private static final Statement INCOMPLETE_BY_EVENT_AND_TARGET_IDENTIFIER_STATEMENT = Cypher
.match(EVENT_PUBLICATION_NODE)
.where(EVENT_PUBLICATION_NODE.property(EVENT_HASH).eq(Cypher.parameter(EVENT_HASH)))
.and(EVENT_PUBLICATION_NODE.property(LISTENER_ID).eq(Cypher.parameter(LISTENER_ID)))
private static final Node EVENT_PUBLICATION_COMPLETED_NODE = node("Neo4jEventPublicationCompleted")
.named("neo4jEventPublicationCompleted");
private static final Statement INCOMPLETE_BY_EVENT_AND_TARGET_IDENTIFIER_STATEMENT = match(EVENT_PUBLICATION_NODE)
.where(EVENT_PUBLICATION_NODE.property(EVENT_HASH).eq(parameter(EVENT_HASH)))
.and(EVENT_PUBLICATION_NODE.property(LISTENER_ID).eq(parameter(LISTENER_ID)))
.and(EVENT_PUBLICATION_NODE.property(COMPLETION_DATE).isNull())
.returning(EVENT_PUBLICATION_NODE)
.build();
private static final Statement DELETE_BY_EVENT_AND_LISTENER_ID = Cypher.match(EVENT_PUBLICATION_NODE)
.where(EVENT_PUBLICATION_NODE.property(EVENT_HASH).eq(Cypher.parameter(EVENT_HASH)))
.and(EVENT_PUBLICATION_NODE.property(LISTENER_ID).eq(Cypher.parameter(LISTENER_ID)))
private static final Statement DELETE_BY_EVENT_AND_LISTENER_ID = match(EVENT_PUBLICATION_NODE)
.where(EVENT_PUBLICATION_NODE.property(EVENT_HASH).eq(parameter(EVENT_HASH)))
.and(EVENT_PUBLICATION_NODE.property(LISTENER_ID).eq(parameter(LISTENER_ID)))
.delete(EVENT_PUBLICATION_NODE)
.build();
private static final Statement DELETE_BY_ID_STATEMENT = Cypher.match(EVENT_PUBLICATION_NODE)
.where(EVENT_PUBLICATION_NODE.property(ID).in(Cypher.parameter(ID)))
private static final Statement DELETE_BY_ID_STATEMENT = match(EVENT_PUBLICATION_NODE)
.where(EVENT_PUBLICATION_NODE.property(ID).in(parameter(ID)))
.delete(EVENT_PUBLICATION_NODE)
.build();
private static final Statement DELETE_COMPLETED_STATEMENT = Cypher.match(EVENT_PUBLICATION_NODE)
.where(EVENT_PUBLICATION_NODE.property(COMPLETION_DATE).isNotNull())
.delete(EVENT_PUBLICATION_NODE)
private static final Function<Node, Statement> DELETE_COMPLETED_STATEMENT = node -> match(node)
.where(node.property(COMPLETION_DATE).isNotNull())
.delete(node)
.build();
private static final Statement DELETE_COMPLETED_BEFORE_STATEMENT = Cypher.match(EVENT_PUBLICATION_NODE)
.where(EVENT_PUBLICATION_NODE.property(PUBLICATION_DATE).lt(Cypher.parameter(PUBLICATION_DATE)))
.and(EVENT_PUBLICATION_NODE.property(COMPLETION_DATE).isNotNull())
.delete(EVENT_PUBLICATION_NODE)
private static final Function<Node, Statement> DELETE_COMPLETED_BEFORE_STATEMENT = node -> match(node)
.where(node.property(PUBLICATION_DATE).lt(parameter(PUBLICATION_DATE)))
.and(node.property(COMPLETION_DATE).isNotNull())
.delete(node)
.build();
private static final Statement INCOMPLETE_PUBLISHED_BEFORE_STATEMENT = Cypher
.match(EVENT_PUBLICATION_NODE)
.where(EVENT_PUBLICATION_NODE.property(PUBLICATION_DATE).lt(Cypher.parameter(PUBLICATION_DATE)))
private static final Statement INCOMPLETE_PUBLISHED_BEFORE_STATEMENT = match(EVENT_PUBLICATION_NODE)
.where(EVENT_PUBLICATION_NODE.property(PUBLICATION_DATE).lt(parameter(PUBLICATION_DATE)))
.and(EVENT_PUBLICATION_NODE.property(COMPLETION_DATE).isNull())
.returning(EVENT_PUBLICATION_NODE)
.orderBy(EVENT_PUBLICATION_NODE.property(PUBLICATION_DATE))
.build();
private static final Statement CREATE_STATEMENT = Cypher.create(EVENT_PUBLICATION_NODE)
.set(EVENT_PUBLICATION_NODE.property(ID).to(Cypher.parameter(ID)))
.set(EVENT_PUBLICATION_NODE.property(EVENT_SERIALIZED).to(Cypher.parameter(EVENT_SERIALIZED)))
.set(EVENT_PUBLICATION_NODE.property(EVENT_HASH).to(Cypher.parameter(EVENT_HASH)))
.set(EVENT_PUBLICATION_NODE.property(EVENT_TYPE).to(Cypher.parameter(EVENT_TYPE)))
.set(EVENT_PUBLICATION_NODE.property(LISTENER_ID).to(Cypher.parameter(LISTENER_ID)))
.set(EVENT_PUBLICATION_NODE.property(PUBLICATION_DATE).to(Cypher.parameter(PUBLICATION_DATE)))
.set(EVENT_PUBLICATION_NODE.property(ID).to(parameter(ID)))
.set(EVENT_PUBLICATION_NODE.property(EVENT_SERIALIZED).to(parameter(EVENT_SERIALIZED)))
.set(EVENT_PUBLICATION_NODE.property(EVENT_HASH).to(parameter(EVENT_HASH)))
.set(EVENT_PUBLICATION_NODE.property(EVENT_TYPE).to(parameter(EVENT_TYPE)))
.set(EVENT_PUBLICATION_NODE.property(LISTENER_ID).to(parameter(LISTENER_ID)))
.set(EVENT_PUBLICATION_NODE.property(PUBLICATION_DATE).to(parameter(PUBLICATION_DATE)))
.build();
private static final Statement COMPLETE_STATEMENT = Cypher.match(EVENT_PUBLICATION_NODE)
.where(EVENT_PUBLICATION_NODE.property(EVENT_HASH).eq(Cypher.parameter(EVENT_HASH)))
.and(EVENT_PUBLICATION_NODE.property(LISTENER_ID).eq(Cypher.parameter(LISTENER_ID)))
private static final Statement COMPLETE_STATEMENT = match(EVENT_PUBLICATION_NODE)
.where(EVENT_PUBLICATION_NODE.property(EVENT_HASH).eq(parameter(EVENT_HASH)))
.and(EVENT_PUBLICATION_NODE.property(LISTENER_ID).eq(parameter(LISTENER_ID)))
.and(EVENT_PUBLICATION_NODE.property(COMPLETION_DATE).isNull())
.set(EVENT_PUBLICATION_NODE.property(COMPLETION_DATE).to(Cypher.parameter(COMPLETION_DATE)))
.set(EVENT_PUBLICATION_NODE.property(COMPLETION_DATE).to(parameter(COMPLETION_DATE)))
.build();
private static final Statement COMPLETE_BY_ID_STATEMENT = Cypher.match(EVENT_PUBLICATION_NODE)
.where(EVENT_PUBLICATION_NODE.property(ID).eq(Cypher.parameter(ID)))
.set(EVENT_PUBLICATION_NODE.property(COMPLETION_DATE).to(Cypher.parameter(COMPLETION_DATE)))
private static final Statement COMPLETE_IN_ARCHIVE_BY_ID_STATEMENT = match(EVENT_PUBLICATION_NODE)
.where(EVENT_PUBLICATION_NODE.property(ID).eq(parameter(ID)))
.and(not(exists(match(EVENT_PUBLICATION_COMPLETED_NODE)
.where(EVENT_PUBLICATION_COMPLETED_NODE.property(ID).eq(parameter(ID)))
.returning(literalTrue()).build())))
.with(EVENT_PUBLICATION_NODE)
.create(EVENT_PUBLICATION_COMPLETED_NODE)
.set(EVENT_PUBLICATION_COMPLETED_NODE.property(ID).to(EVENT_PUBLICATION_NODE.property(ID)))
.set(EVENT_PUBLICATION_COMPLETED_NODE.property(COMPLETION_DATE).to(parameter(COMPLETION_DATE)))
.build();
private static final ResultStatement INCOMPLETE_STATEMENT = Cypher.match(EVENT_PUBLICATION_NODE)
private static final Statement COMPLETE_IN_ARCHIVE_BY_EVENT_AND_LISTENER_ID_STATEMENT = match(EVENT_PUBLICATION_NODE)
.where(EVENT_PUBLICATION_NODE.property(EVENT_HASH).eq(parameter(EVENT_HASH)))
.and(EVENT_PUBLICATION_NODE.property(LISTENER_ID).eq(parameter(LISTENER_ID)))
.and(not(exists(match(EVENT_PUBLICATION_COMPLETED_NODE)
.where(EVENT_PUBLICATION_COMPLETED_NODE.property(EVENT_HASH).eq(parameter(EVENT_HASH)))
.and(EVENT_PUBLICATION_COMPLETED_NODE.property(LISTENER_ID).eq(parameter(LISTENER_ID)))
.returning(literalTrue()).build())))
.with(EVENT_PUBLICATION_NODE)
.create(EVENT_PUBLICATION_COMPLETED_NODE)
.set(EVENT_PUBLICATION_COMPLETED_NODE.property(ID).to(EVENT_PUBLICATION_NODE.property(ID)))
.set(EVENT_PUBLICATION_COMPLETED_NODE.property(COMPLETION_DATE).to(parameter(COMPLETION_DATE)))
.build();
private static final Function<Node, Statement> COMPLETE_BY_ID_STATEMENT = node -> match(node)
.where(node.property(ID).eq(parameter(ID)))
.set(node.property(COMPLETION_DATE).to(parameter(COMPLETION_DATE)))
.build();
private static final ResultStatement INCOMPLETE_STATEMENT = match(EVENT_PUBLICATION_NODE)
.where(EVENT_PUBLICATION_NODE.property(COMPLETION_DATE).isNull())
.returning(EVENT_PUBLICATION_NODE)
.orderBy(EVENT_PUBLICATION_NODE.property(PUBLICATION_DATE))
.build();
private static final ResultStatement ALL_COMPLETED_STATEMENT = Cypher.match(EVENT_PUBLICATION_NODE)
.where(EVENT_PUBLICATION_NODE.property(COMPLETION_DATE).isNotNull())
.returning(EVENT_PUBLICATION_NODE)
.orderBy(EVENT_PUBLICATION_NODE.property(PUBLICATION_DATE))
private static final Function<Node, ResultStatement> ALL_COMPLETED_STATEMENT = node -> match(node)
.where(node.property(COMPLETION_DATE).isNotNull())
.returning(node)
.orderBy(node.property(PUBLICATION_DATE))
.build();
private final Neo4jClient neo4jClient;
@@ -140,6 +169,11 @@ class Neo4jEventPublicationRepository implements EventPublicationRepository {
private final EventSerializer eventSerializer;
private final CompletionMode completionMode;
private final Statement deleteCompletedStatement;
private final Statement deleteCompletedBeforeStatement;
private final Statement completedByIdStatement;
private final ResultStatement allCompletedStatement;
Neo4jEventPublicationRepository(Neo4jClient neo4jClient, Configuration cypherDslConfiguration,
EventSerializer eventSerializer, CompletionMode completionMode) {
@@ -152,6 +186,13 @@ class Neo4jEventPublicationRepository implements EventPublicationRepository {
this.renderer = Renderer.getRenderer(cypherDslConfiguration);
this.eventSerializer = eventSerializer;
this.completionMode = completionMode;
var archiveNode = completionMode == CompletionMode.ARCHIVE ? EVENT_PUBLICATION_COMPLETED_NODE : EVENT_PUBLICATION_NODE;
this.deleteCompletedStatement = DELETE_COMPLETED_STATEMENT.apply(archiveNode);
this.deleteCompletedBeforeStatement = DELETE_COMPLETED_BEFORE_STATEMENT.apply(archiveNode);
this.completedByIdStatement = COMPLETE_BY_ID_STATEMENT.apply(archiveNode);
this.allCompletedStatement = ALL_COMPLETED_STATEMENT.apply(archiveNode);
}
/*
@@ -201,6 +242,18 @@ class Neo4jEventPublicationRepository implements EventPublicationRepository {
.bind(identifier.getValue()).to(LISTENER_ID)
.run();
} else if (completionMode == CompletionMode.ARCHIVE) {
neo4jClient.query(renderer.render(COMPLETE_IN_ARCHIVE_BY_EVENT_AND_LISTENER_ID_STATEMENT))
.bind(eventHash).to(EVENT_HASH)
.bind(identifier.getValue()).to(LISTENER_ID)
.bind(Values.value(completionDate.atOffset(ZoneOffset.UTC))).to(COMPLETION_DATE)
.run();
neo4jClient.query(renderer.render(DELETE_BY_EVENT_AND_LISTENER_ID))
.bind(eventHash).to(EVENT_HASH)
.bind(identifier.getValue()).to(LISTENER_ID)
.run();
} else {
neo4jClient.query(renderer.render(COMPLETE_STATEMENT))
@@ -223,13 +276,22 @@ class Neo4jEventPublicationRepository implements EventPublicationRepository {
deletePublications(List.of(identifier));
} else if (completionMode == CompletionMode.ARCHIVE) {
neo4jClient.query(renderer.render(COMPLETE_IN_ARCHIVE_BY_ID_STATEMENT))
.bind("").to(ID)
.bind(Values.value(completionDate.atOffset(ZoneOffset.UTC))).to(COMPLETION_DATE)
.run();
deletePublications(List.of(identifier));
} else {
neo4jClient.query(renderer.render(COMPLETE_BY_ID_STATEMENT))
neo4jClient.query(renderer.render(completedByIdStatement))
.bind(Values.value(identifier.toString())).to(ID)
.bind(Values.value(completionDate.atOffset(ZoneOffset.UTC))).to(COMPLETION_DATE)
.run();
}
}
/*
@@ -287,7 +349,7 @@ class Neo4jEventPublicationRepository implements EventPublicationRepository {
@Override
public List<TargetEventPublication> findCompletedPublications() {
return new ArrayList<>(neo4jClient.query(renderer.render(ALL_COMPLETED_STATEMENT))
return new ArrayList<>(neo4jClient.query(renderer.render(allCompletedStatement))
.fetchAs(TargetEventPublication.class)
.mappedBy(this::mapRecordToPublication)
.all());
@@ -313,7 +375,7 @@ class Neo4jEventPublicationRepository implements EventPublicationRepository {
@Override
@Transactional
public void deleteCompletedPublications() {
neo4jClient.query(renderer.render(DELETE_COMPLETED_STATEMENT)).run();
neo4jClient.query(renderer.render(deleteCompletedStatement)).run();
}
/*
@@ -324,7 +386,7 @@ class Neo4jEventPublicationRepository implements EventPublicationRepository {
@Transactional
public void deleteCompletedPublicationsBefore(Instant instant) {
neo4jClient.query(renderer.render(DELETE_COMPLETED_BEFORE_STATEMENT))
neo4jClient.query(renderer.render(deleteCompletedBeforeStatement))
.bind(Values.value(instant.atOffset(ZoneOffset.UTC))).to(PUBLICATION_DATE)
.run();
}

View File

@@ -32,7 +32,7 @@ import org.neo4j.driver.AuthTokens;
import org.neo4j.driver.Driver;
import org.neo4j.driver.GraphDatabase;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.test.context.bean.override.mockito.MockitoBean;
import org.springframework.boot.autoconfigure.ImportAutoConfiguration;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.context.annotation.Import;
@@ -42,8 +42,9 @@ import org.springframework.modulith.events.core.PublicationTargetIdentifier;
import org.springframework.modulith.events.core.TargetEventPublication;
import org.springframework.modulith.events.support.CompletionMode;
import org.springframework.modulith.testapp.TestApplication;
import org.springframework.test.context.ContextConfiguration;
import org.springframework.test.context.TestPropertySource;
import org.springframework.test.context.junit.jupiter.SpringJUnitConfig;
import org.springframework.test.context.bean.override.mockito.MockitoBean;
import org.springframework.util.DigestUtils;
import org.testcontainers.containers.Neo4jContainer;
import org.testcontainers.junit.jupiter.Container;
@@ -53,256 +54,270 @@ import org.testcontainers.utility.DockerImageName;
/**
* @author Gerrit Meier
*/
@SpringJUnitConfig(Neo4jEventPublicationRepositoryTest.Config.class)
@Testcontainers(disabledWithoutDocker = true)
class Neo4jEventPublicationRepositoryTest {
static final PublicationTargetIdentifier TARGET_IDENTIFIER = PublicationTargetIdentifier.of("listener");
@Container //
static final Neo4jContainer<?> neo4jContainer = new Neo4jContainer<>(DockerImageName.parse("neo4j:5"))
.withRandomPassword();
static final PublicationTargetIdentifier TARGET_IDENTIFIER = PublicationTargetIdentifier.of("listener");
@Import(TestApplication.class)
@ImportAutoConfiguration({ Neo4jEventPublicationAutoConfiguration.class })
@Testcontainers(disabledWithoutDocker = true)
@ContextConfiguration(classes = Config.class)
static abstract class TestBase {
@Autowired Neo4jEventPublicationRepository repository;
@Autowired Driver driver;
@Autowired Environment environment;
@Autowired Neo4jEventPublicationRepository repository;
@Autowired Driver driver;
@Autowired Environment environment;
@MockitoBean EventSerializer eventSerializer;
@MockitoBean EventSerializer eventSerializer;
CompletionMode completionMode;
CompletionMode completionMode;
@BeforeEach
void clearDb() {
@BeforeEach
void clearDb() {
this.completionMode = CompletionMode.from(environment);
this.completionMode = CompletionMode.from(environment);
try (var session = driver.session()) {
session.run("MATCH (n) detach delete n").consume();
try (var session = driver.session()) {
session.run("MATCH (n) detach delete n").consume();
}
}
}
@Test
void createEventPublication() {
@Test
void createEventPublication() {
var testEvent = new TestEvent("id");
var eventSerialized = "{\"eventId\":\"id\"}";
var eventHash = DigestUtils.md5DigestAsHex(eventSerialized.getBytes());
var testEvent = new TestEvent("id");
var eventSerialized = "{\"eventId\":\"id\"}";
var eventHash = DigestUtils.md5DigestAsHex(eventSerialized.getBytes());
when(eventSerializer.serialize(testEvent)).thenReturn(eventSerialized);
var publication = repository.create(TargetEventPublication.of(testEvent, TARGET_IDENTIFIER));
when(eventSerializer.serialize(testEvent)).thenReturn(eventSerialized);
var publication = repository.create(TargetEventPublication.of(testEvent, TARGET_IDENTIFIER));
try (var session = driver.session()) {
try (var session = driver.session()) {
var result = session.run("MATCH (p:Neo4jEventPublication) return p")
.single();
var result = session.run("MATCH (p:Neo4jEventPublication) return p")
.single();
var neo4jEventPublicationNode = result.get("p").asNode();
var neo4jEventPublicationNode = result.get("p").asNode();
assertThat(UUID.fromString(neo4jEventPublicationNode.get("identifier").asString()))
.isEqualTo(publication.getIdentifier());
assertThat(neo4jEventPublicationNode.get("publicationDate").asZonedDateTime().toInstant())
.isEqualTo(publication.getPublicationDate());
assertThat(neo4jEventPublicationNode.get("listenerId").asString())
.isEqualTo(publication.getTargetIdentifier().getValue());
assertThat(neo4jEventPublicationNode.get("completionDate").isNull()).isTrue();
assertThat(neo4jEventPublicationNode.get("eventSerialized").asString()).isEqualTo(eventSerialized);
assertThat(neo4jEventPublicationNode.get("eventHash").asString()).isEqualTo(eventHash);
assertThat(UUID.fromString(neo4jEventPublicationNode.get("identifier").asString()))
.isEqualTo(publication.getIdentifier());
assertThat(neo4jEventPublicationNode.get("publicationDate").asZonedDateTime().toInstant())
.isEqualTo(publication.getPublicationDate());
assertThat(neo4jEventPublicationNode.get("listenerId").asString())
.isEqualTo(publication.getTargetIdentifier().getValue());
assertThat(neo4jEventPublicationNode.get("completionDate").isNull()).isTrue();
assertThat(neo4jEventPublicationNode.get("eventSerialized").asString()).isEqualTo(eventSerialized);
assertThat(neo4jEventPublicationNode.get("eventHash").asString()).isEqualTo(eventHash);
}
}
}
@Test
void updateEventPublication() {
@Test
void updateEventPublication() {
var testEvent1 = new TestEvent("id1");
var event1Serialized = "{\"eventId\":\"id1\"}";
var testEvent2 = new TestEvent("id2");
var event2Serialized = "{\"eventId\":\"id2\"}";
var testEvent1 = new TestEvent("id1");
var event1Serialized = "{\"eventId\":\"id1\"}";
var testEvent2 = new TestEvent("id2");
var event2Serialized = "{\"eventId\":\"id2\"}";
when(eventSerializer.serialize(testEvent1)).thenReturn(event1Serialized);
when(eventSerializer.serialize(testEvent2)).thenReturn(event2Serialized);
when(eventSerializer.deserialize(event2Serialized, TestEvent.class)).thenReturn(testEvent2);
when(eventSerializer.serialize(testEvent1)).thenReturn(event1Serialized);
when(eventSerializer.serialize(testEvent2)).thenReturn(event2Serialized);
when(eventSerializer.deserialize(event2Serialized, TestEvent.class)).thenReturn(testEvent2);
var event1 = repository.create(TargetEventPublication.of(testEvent1, TARGET_IDENTIFIER));
var event2 = repository.create(TargetEventPublication.of(testEvent2, TARGET_IDENTIFIER));
var event1 = repository.create(TargetEventPublication.of(testEvent1, TARGET_IDENTIFIER));
var event2 = repository.create(TargetEventPublication.of(testEvent2, TARGET_IDENTIFIER));
var now = Instant.now();
repository.markCompleted(event1, now);
var now = Instant.now();
repository.markCompleted(event1, now);
assertThat(repository.findIncompletePublications()).hasSize(1)
.element(0)
.extracting(TargetEventPublication::getEvent).isEqualTo(event2.getEvent());
}
@Test
void findInCompletePastPublications() {
var testEvent = new TestEvent("id");
var eventSerialized = "{\"eventId\":\"id\"}";
when(eventSerializer.serialize(testEvent)).thenReturn(eventSerialized);
when(eventSerializer.deserialize(eventSerialized, TestEvent.class)).thenReturn(testEvent);
var event = repository.create(TargetEventPublication.of(testEvent, TARGET_IDENTIFIER));
var newer = Instant.now().plus(1L, ChronoUnit.MINUTES);
var older = Instant.now().minus(1L, ChronoUnit.MINUTES);
assertThat(repository.findIncompletePublicationsPublishedBefore(newer)).hasSize(1)
.element(0)
.extracting(TargetEventPublication::getEvent).isEqualTo(event.getEvent());
assertThat(repository.findIncompletePublicationsPublishedBefore(older)).hasSize(0);
}
@Test
void findIncompleteByEventAndTargetIdentifier() {
var testEvent = new TestEvent("id");
var eventSerialized = "{\"eventId\":\"id\"}";
when(eventSerializer.serialize(testEvent)).thenReturn(eventSerialized);
when(eventSerializer.deserialize(eventSerialized, TestEvent.class)).thenReturn(testEvent);
var event = repository.create(TargetEventPublication.of(testEvent, TARGET_IDENTIFIER));
assertThat(repository.findIncompletePublicationsByEventAndTargetIdentifier(testEvent, event.getTargetIdentifier()))
.isPresent();
}
@Test
void deletePublicationById() {
var testEvent = new TestEvent("id");
var eventSerialized = "{\"eventId\":\"id\"}";
when(eventSerializer.serialize(testEvent)).thenReturn(eventSerialized);
var event = repository.create(TargetEventPublication.of(testEvent, TARGET_IDENTIFIER));
assertThat(repository.findIncompletePublications()).hasSize(1);
repository.deletePublications(List.of(event.getIdentifier()));
assertThat(repository.findIncompletePublications()).hasSize(0);
}
@Test
void deleteCompletedPublications() {
var testEvent1 = new TestEvent("id1");
var event1Serialized = "{\"eventId\":\"id1\"}";
var testEvent2 = new TestEvent("id2");
var event2Serialized = "{\"eventId\":\"id2\"}";
when(eventSerializer.serialize(testEvent1)).thenReturn(event1Serialized);
when(eventSerializer.serialize(testEvent2)).thenReturn(event2Serialized);
when(eventSerializer.deserialize(event1Serialized, TestEvent.class)).thenReturn(testEvent1);
when(eventSerializer.deserialize(event2Serialized, TestEvent.class)).thenReturn(testEvent2);
var event1 = repository.create(TargetEventPublication.of(testEvent1, TARGET_IDENTIFIER));
repository.markCompleted(event1, Instant.now());
repository.deleteCompletedPublications();
try (var session = driver.session()) {
var count = session.run("MATCH (n) WHERE n.completionDate is not null return count(n)").single().get("count(n)")
.asLong();
assertThat(count).isEqualTo(0);
}
}
@Test
void deleteCompletedPublicationsBefore() throws Exception {
assumeTrue(completionMode == CompletionMode.UPDATE);
var testEvent1 = new TestEvent("id1");
var event1Serialized = "{\"eventId\":\"id1\"}";
var testEvent2 = new TestEvent("id2");
var event2Serialized = "{\"eventId\":\"id2\"}";
when(eventSerializer.serialize(testEvent1)).thenReturn(event1Serialized);
when(eventSerializer.serialize(testEvent2)).thenReturn(event2Serialized);
when(eventSerializer.deserialize(event1Serialized, TestEvent.class)).thenReturn(testEvent1);
when(eventSerializer.deserialize(event2Serialized, TestEvent.class)).thenReturn(testEvent2);
var event1 = repository.create(TargetEventPublication.of(testEvent1, TARGET_IDENTIFIER));
Instant old = Instant.now();
repository.markCompleted(event1, old);
Thread.sleep(100);
var event2 = repository.create(TargetEventPublication.of(testEvent2, TARGET_IDENTIFIER));
repository.markCompleted(event2, Instant.now());
repository.deleteCompletedPublicationsBefore(old.plus(10, ChronoUnit.MILLIS));
// defensive check just to be sure
assertThat(repository.findIncompletePublications()).hasSize(0);
try (var session = driver.session()) {
var records = session.run("MATCH (n) WHERE n.completionDate is not null return n").list();
assertThat(records.size()).isEqualTo(1);
assertThat(records.get(0).get("n").asNode().get("eventSerialized").asString()).contains("id2");
}
}
@Test // GH-451
void findsCompletedPublications() {
var event = new TestEvent("first");
var publication = createPublication(event);
repository.markCompleted(publication, Instant.now());
if (completionMode == CompletionMode.DELETE) {
assertThat(repository.findCompletedPublications()).isEmpty();
assertThat(repository.findIncompletePublications()).isEmpty();
} else {
assertThat(repository.findCompletedPublications())
.hasSize(1)
assertThat(repository.findIncompletePublications()).hasSize(1)
.element(0)
.extracting(TargetEventPublication::getEvent)
.isEqualTo(event);
.extracting(TargetEventPublication::getEvent).isEqualTo(event2.getEvent());
}
}
@Test // GH-258
void marksPublicationAsCompletedById() {
@Test
void findInCompletePastPublications() {
var event = new TestEvent("first");
var publication = createPublication(event);
var testEvent = new TestEvent("id");
var eventSerialized = "{\"eventId\":\"id\"}";
repository.markCompleted(publication.getIdentifier(), Instant.now());
when(eventSerializer.serialize(testEvent)).thenReturn(eventSerialized);
when(eventSerializer.deserialize(eventSerialized, TestEvent.class)).thenReturn(testEvent);
if (completionMode == CompletionMode.DELETE) {
var event = repository.create(TargetEventPublication.of(testEvent, TARGET_IDENTIFIER));
assertThat(repository.findCompletedPublications()).isEmpty();
assertThat(repository.findIncompletePublications()).isEmpty();
var newer = Instant.now().plus(1L, ChronoUnit.MINUTES);
var older = Instant.now().minus(1L, ChronoUnit.MINUTES);
} else {
assertThat(repository.findIncompletePublicationsPublishedBefore(newer)).hasSize(1)
.element(0)
.extracting(TargetEventPublication::getEvent).isEqualTo(event.getEvent());
assertThat(repository.findCompletedPublications())
.extracting(TargetEventPublication::getIdentifier)
.containsExactly(publication.getIdentifier());
assertThat(repository.findIncompletePublicationsPublishedBefore(older)).hasSize(0);
}
}
private TargetEventPublication createPublication(Object event) {
@Test
void findIncompleteByEventAndTargetIdentifier() {
var token = event.toString();
var testEvent = new TestEvent("id");
var eventSerialized = "{\"eventId\":\"id\"}";
doReturn(token).when(eventSerializer).serialize(event);
doReturn(event).when(eventSerializer).deserialize(token, event.getClass());
when(eventSerializer.serialize(testEvent)).thenReturn(eventSerialized);
when(eventSerializer.deserialize(eventSerialized, TestEvent.class)).thenReturn(testEvent);
return repository.create(TargetEventPublication.of(event, TARGET_IDENTIFIER));
var event = repository.create(TargetEventPublication.of(testEvent, TARGET_IDENTIFIER));
assertThat(
repository.findIncompletePublicationsByEventAndTargetIdentifier(testEvent, event.getTargetIdentifier()))
.isPresent();
}
@Test
void deletePublicationById() {
var testEvent = new TestEvent("id");
var eventSerialized = "{\"eventId\":\"id\"}";
when(eventSerializer.serialize(testEvent)).thenReturn(eventSerialized);
var event = repository.create(TargetEventPublication.of(testEvent, TARGET_IDENTIFIER));
assertThat(repository.findIncompletePublications()).hasSize(1);
repository.deletePublications(List.of(event.getIdentifier()));
assertThat(repository.findIncompletePublications()).hasSize(0);
}
@Test
void deleteCompletedPublications() {
var testEvent1 = new TestEvent("id1");
var event1Serialized = "{\"eventId\":\"id1\"}";
var testEvent2 = new TestEvent("id2");
var event2Serialized = "{\"eventId\":\"id2\"}";
when(eventSerializer.serialize(testEvent1)).thenReturn(event1Serialized);
when(eventSerializer.serialize(testEvent2)).thenReturn(event2Serialized);
when(eventSerializer.deserialize(event1Serialized, TestEvent.class)).thenReturn(testEvent1);
when(eventSerializer.deserialize(event2Serialized, TestEvent.class)).thenReturn(testEvent2);
var event1 = repository.create(TargetEventPublication.of(testEvent1, TARGET_IDENTIFIER));
repository.markCompleted(event1, Instant.now());
repository.deleteCompletedPublications();
try (var session = driver.session()) {
var count = session.run("MATCH (n) WHERE n.completionDate is not null return count(n)").single().get("count(n)")
.asLong();
assertThat(count).isEqualTo(0);
}
}
@Test
void deleteCompletedPublicationsBefore() throws Exception {
assumeTrue(completionMode == CompletionMode.UPDATE);
var testEvent1 = new TestEvent("id1");
var event1Serialized = "{\"eventId\":\"id1\"}";
var testEvent2 = new TestEvent("id2");
var event2Serialized = "{\"eventId\":\"id2\"}";
when(eventSerializer.serialize(testEvent1)).thenReturn(event1Serialized);
when(eventSerializer.serialize(testEvent2)).thenReturn(event2Serialized);
when(eventSerializer.deserialize(event1Serialized, TestEvent.class)).thenReturn(testEvent1);
when(eventSerializer.deserialize(event2Serialized, TestEvent.class)).thenReturn(testEvent2);
var event1 = repository.create(TargetEventPublication.of(testEvent1, TARGET_IDENTIFIER));
Instant old = Instant.now();
repository.markCompleted(event1, old);
Thread.sleep(100);
var event2 = repository.create(TargetEventPublication.of(testEvent2, TARGET_IDENTIFIER));
repository.markCompleted(event2, Instant.now());
repository.deleteCompletedPublicationsBefore(old.plus(10, ChronoUnit.MILLIS));
// defensive check just to be sure
assertThat(repository.findIncompletePublications()).hasSize(0);
try (var session = driver.session()) {
var records = session.run("MATCH (n) WHERE n.completionDate is not null return n").list();
assertThat(records.size()).isEqualTo(1);
assertThat(records.get(0).get("n").asNode().get("eventSerialized").asString()).contains("id2");
}
}
@Test // GH-451
void findsCompletedPublications() {
var event = new TestEvent("first");
var publication = createPublication(event);
repository.markCompleted(publication, Instant.now());
if (completionMode == CompletionMode.DELETE) {
assertThat(repository.findCompletedPublications()).isEmpty();
assertThat(repository.findIncompletePublications()).isEmpty();
} else {
assertThat(repository.findCompletedPublications())
.hasSize(1)
.element(0)
.extracting(TargetEventPublication::getEvent)
.isEqualTo(event);
}
}
@Test // GH-258
void marksPublicationAsCompletedById() {
var event = new TestEvent("first");
var publication = createPublication(event);
repository.markCompleted(publication.getIdentifier(), Instant.now());
if (completionMode == CompletionMode.DELETE) {
assertThat(repository.findCompletedPublications()).isEmpty();
assertThat(repository.findIncompletePublications()).isEmpty();
} else {
assertThat(repository.findCompletedPublications())
.extracting(TargetEventPublication::getIdentifier)
.containsExactly(publication.getIdentifier());
}
}
private TargetEventPublication createPublication(Object event) {
var token = event.toString();
doReturn(token).when(eventSerializer).serialize(event);
doReturn(event).when(eventSerializer).deserialize(token, event.getClass());
return repository.create(TargetEventPublication.of(event, TARGET_IDENTIFIER));
}
}
@Nested
@TestPropertySource(properties = CompletionMode.PROPERTY + "=UPDATE")
class WithUpdateCompletionTest extends TestBase {}
@Nested
@TestPropertySource(properties = CompletionMode.PROPERTY + "=DELETE")
static class WithDeleteCompletionTest extends Neo4jEventPublicationRepositoryTest {}
class WithDeleteCompletionTest extends TestBase {}
@Nested
@TestPropertySource(properties = CompletionMode.PROPERTY + "=ARCHIVE")
class WithArchiveCompletionTest extends TestBase {}
private record TestEvent(String eventId) {}