GH-806 - Add archive support for MongoDB.
Co-authored-by: Oliver Drotbohm <oliver.drotbohm@broadcom.com>
This commit is contained in:
committed by
Oliver Drotbohm
parent
2cb0db7f42
commit
a7e4a798f0
@@ -52,11 +52,13 @@ class MongoDbEventPublicationRepository implements EventPublicationRepository {
|
||||
private static final String ID = "id";
|
||||
private static final String LISTENER_ID = "listenerId";
|
||||
private static final String PUBLICATION_DATE = "publicationDate";
|
||||
|
||||
private static final Sort DEFAULT_SORT = Sort.by(PUBLICATION_DATE).ascending();
|
||||
|
||||
static final String ARCHIVE_COLLECTION = "event_publication_archive";
|
||||
|
||||
private final MongoTemplate mongoTemplate;
|
||||
private final CompletionMode completionMode;
|
||||
private final String collection, archiveCollection;
|
||||
|
||||
/**
|
||||
* Creates a new {@link MongoDbEventPublicationRepository} for the given {@link MongoTemplate}.
|
||||
@@ -71,6 +73,8 @@ class MongoDbEventPublicationRepository implements EventPublicationRepository {
|
||||
|
||||
this.mongoTemplate = mongoTemplate;
|
||||
this.completionMode = completionMode;
|
||||
this.collection = "event_publication";
|
||||
this.archiveCollection = completionMode == CompletionMode.ARCHIVE ? ARCHIVE_COLLECTION : collection;
|
||||
}
|
||||
|
||||
/*
|
||||
@@ -80,7 +84,7 @@ class MongoDbEventPublicationRepository implements EventPublicationRepository {
|
||||
@Override
|
||||
public TargetEventPublication create(TargetEventPublication publication) {
|
||||
|
||||
mongoTemplate.save(domainToDocument(publication));
|
||||
mongoTemplate.save(domainToDocument(publication), collection);
|
||||
|
||||
return publication;
|
||||
}
|
||||
@@ -93,16 +97,20 @@ class MongoDbEventPublicationRepository implements EventPublicationRepository {
|
||||
public void markCompleted(Object event, PublicationTargetIdentifier identifier, Instant completionDate) {
|
||||
|
||||
var query = byEventAndListenerId(event, identifier);
|
||||
var update = Update.update(COMPLETION_DATE, completionDate);
|
||||
|
||||
if (completionMode == CompletionMode.DELETE) {
|
||||
|
||||
mongoTemplate.remove(query, MongoDbEventPublication.class);
|
||||
mongoTemplate.remove(query, MongoDbEventPublication.class, collection);
|
||||
|
||||
} else if (completionMode == CompletionMode.ARCHIVE) {
|
||||
|
||||
mongoTemplate.findAndModify(query, update, MongoDbEventPublication.class, collection);
|
||||
var completedEvent = mongoTemplate.findAndRemove(query, MongoDbEventPublication.class, collection);
|
||||
mongoTemplate.save(completedEvent, archiveCollection);
|
||||
} else {
|
||||
|
||||
var update = Update.update(COMPLETION_DATE, completionDate);
|
||||
|
||||
mongoTemplate.findAndModify(query, update, MongoDbEventPublication.class);
|
||||
mongoTemplate.findAndModify(query, update, MongoDbEventPublication.class, collection);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -113,17 +121,21 @@ class MongoDbEventPublicationRepository implements EventPublicationRepository {
|
||||
@Override
|
||||
public void markCompleted(UUID identifier, Instant completionDate) {
|
||||
|
||||
var criateria = query(where(ID).is(identifier));
|
||||
var criteria = query(where(ID).is(identifier));
|
||||
var update = Update.update(COMPLETION_DATE, completionDate);
|
||||
|
||||
if (completionMode == CompletionMode.DELETE) {
|
||||
|
||||
mongoTemplate.remove(criateria, MongoDbEventPublication.class);
|
||||
mongoTemplate.remove(criteria, MongoDbEventPublication.class, collection);
|
||||
|
||||
} else if (completionMode == CompletionMode.ARCHIVE) {
|
||||
|
||||
mongoTemplate.findAndModify(criteria, update, MongoDbEventPublication.class, collection);
|
||||
var completedEvent = mongoTemplate.findAndRemove(criteria, MongoDbEventPublication.class, collection);
|
||||
mongoTemplate.save(completedEvent, archiveCollection);
|
||||
|
||||
} else {
|
||||
|
||||
var update = Update.update(COMPLETION_DATE, completionDate);
|
||||
|
||||
mongoTemplate.findAndModify(criateria, update, MongoDbEventPublication.class);
|
||||
mongoTemplate.findAndModify(criteria, update, MongoDbEventPublication.class, collection);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -168,7 +180,7 @@ class MongoDbEventPublicationRepository implements EventPublicationRepository {
|
||||
*/
|
||||
@Override
|
||||
public List<TargetEventPublication> findCompletedPublications() {
|
||||
return readMapped(defaultQuery(where(COMPLETION_DATE).ne(null)));
|
||||
return readMapped(defaultQuery(where(COMPLETION_DATE).ne(null)), archiveCollection);
|
||||
}
|
||||
|
||||
/*
|
||||
@@ -177,7 +189,9 @@ class MongoDbEventPublicationRepository implements EventPublicationRepository {
|
||||
*/
|
||||
@Override
|
||||
public void deletePublications(List<UUID> identifiers) {
|
||||
mongoTemplate.remove(query(where(ID).in(identifiers)), MongoDbEventPublication.class);
|
||||
|
||||
mongoTemplate.remove(query(where(ID).in(identifiers)), MongoDbEventPublication.class, collection);
|
||||
mongoTemplate.remove(query(where(ID).in(identifiers)), MongoDbEventPublication.class, archiveCollection);
|
||||
}
|
||||
|
||||
/*
|
||||
@@ -186,7 +200,7 @@ class MongoDbEventPublicationRepository implements EventPublicationRepository {
|
||||
*/
|
||||
@Override
|
||||
public void deleteCompletedPublications() {
|
||||
mongoTemplate.remove(query(where(COMPLETION_DATE).ne(null)), MongoDbEventPublication.class);
|
||||
mongoTemplate.remove(query(where(COMPLETION_DATE).ne(null)), MongoDbEventPublication.class, archiveCollection);
|
||||
}
|
||||
|
||||
/*
|
||||
@@ -198,16 +212,22 @@ class MongoDbEventPublicationRepository implements EventPublicationRepository {
|
||||
|
||||
Assert.notNull(instant, "Instant must not be null!");
|
||||
|
||||
mongoTemplate.remove(query(where(COMPLETION_DATE).lt(instant)), MongoDbEventPublication.class);
|
||||
mongoTemplate.remove(query(where(COMPLETION_DATE).lt(instant)), MongoDbEventPublication.class, archiveCollection);
|
||||
}
|
||||
|
||||
private List<TargetEventPublication> readMapped(Query query) {
|
||||
return readMapped(query, collection);
|
||||
}
|
||||
|
||||
private List<TargetEventPublication> readMapped(Query query, String collection) {
|
||||
|
||||
return mongoTemplate.query(MongoDbEventPublication.class)
|
||||
.inCollection(collection)
|
||||
.matching(query)
|
||||
.stream()
|
||||
.map(MongoDbEventPublicationRepository::documentToDomain)
|
||||
.toList();
|
||||
|
||||
}
|
||||
|
||||
private Query byEventAndListenerId(Object event, PublicationTargetIdentifier identifier) {
|
||||
|
||||
@@ -46,154 +46,146 @@ import org.springframework.test.context.TestPropertySource;
|
||||
* @author Dmitry Belyaev
|
||||
* @author Oliver Drotbohm
|
||||
*/
|
||||
@DataMongoTest
|
||||
@ContextConfiguration(classes = TestApplication.class)
|
||||
class MongoDbEventPublicationRepositoryTest {
|
||||
|
||||
private static final PublicationTargetIdentifier TARGET_IDENTIFIER = PublicationTargetIdentifier.of("listener");
|
||||
|
||||
@Autowired MongoTemplate mongoTemplate;
|
||||
@Autowired Environment environment;
|
||||
@DataMongoTest
|
||||
@ContextConfiguration(classes = TestApplication.class)
|
||||
static abstract class TestBase {
|
||||
|
||||
MongoDbEventPublicationRepository repository;
|
||||
CompletionMode completionMode;
|
||||
@Autowired MongoTemplate mongoTemplate;
|
||||
@Autowired Environment environment;
|
||||
|
||||
@BeforeEach
|
||||
void setUp() {
|
||||
this.completionMode = CompletionMode.from(environment);
|
||||
this.repository = new MongoDbEventPublicationRepository(mongoTemplate, completionMode);
|
||||
}
|
||||
MongoDbEventPublicationRepository repository;
|
||||
CompletionMode completionMode;
|
||||
String archiveCollection = MongoDbEventPublicationRepository.ARCHIVE_COLLECTION;
|
||||
|
||||
@AfterEach
|
||||
void tearDown() {
|
||||
mongoTemplate.remove(MongoDbEventPublication.class).all();
|
||||
}
|
||||
|
||||
@Test // GH-4
|
||||
void shouldPersistAndUpdateEventPublication() {
|
||||
|
||||
var publication = createPublication(new TestEvent("abc"));
|
||||
|
||||
var eventPublications = repository.findIncompletePublications();
|
||||
|
||||
assertThat(eventPublications).hasSize(1);
|
||||
assertThat(eventPublications.get(0).getEvent()).isEqualTo(publication.getEvent());
|
||||
assertThat(eventPublications.get(0).getTargetIdentifier()).isEqualTo(publication.getTargetIdentifier());
|
||||
|
||||
assertThat(repository.findIncompletePublicationsByEventAndTargetIdentifier(new TestEvent("abc"), TARGET_IDENTIFIER))
|
||||
.isPresent();
|
||||
|
||||
// Complete publication
|
||||
repository.markCompleted(publication, Instant.now());
|
||||
|
||||
assertThat(repository.findIncompletePublications()).isEmpty();
|
||||
}
|
||||
|
||||
@Test // GH-4
|
||||
void shouldUpdateSingleEventPublication() {
|
||||
|
||||
var first = createPublication(new TestEvent("id1"));
|
||||
var second = createPublication(new TestEvent("id2"));
|
||||
|
||||
repository.markCompleted(second, Instant.now());
|
||||
|
||||
assertThat(repository.findIncompletePublications()).hasSize(1)
|
||||
.element(0)
|
||||
.extracting(TargetEventPublication::getEvent).isEqualTo(first.getEvent());
|
||||
}
|
||||
|
||||
@Test // GH-133
|
||||
void returnsOldestIncompletePublicationsFirst() {
|
||||
|
||||
var now = LocalDateTime.now();
|
||||
|
||||
savePublicationAt(now.withHour(3));
|
||||
savePublicationAt(now.withHour(0));
|
||||
savePublicationAt(now.withHour(1));
|
||||
|
||||
assertThat(repository.findIncompletePublications())
|
||||
.isSortedAccordingTo(Comparator.comparing(TargetEventPublication::getPublicationDate));
|
||||
}
|
||||
|
||||
@Test // GH-294
|
||||
void findsPublicationsOlderThanReference() throws Exception {
|
||||
|
||||
var first = createPublication(new TestEvent("first"));
|
||||
|
||||
Thread.sleep(100);
|
||||
|
||||
var now = Instant.now();
|
||||
var second = createPublication(new TestEvent("second"));
|
||||
|
||||
assertThat(repository.findIncompletePublications())
|
||||
.extracting(TargetEventPublication::getIdentifier)
|
||||
.containsExactly(first.getIdentifier(), second.getIdentifier());
|
||||
|
||||
assertThat(repository.findIncompletePublicationsPublishedBefore(now))
|
||||
.hasSize(1)
|
||||
.element(0).extracting(TargetEventPublication::getIdentifier).isEqualTo(first.getIdentifier());
|
||||
}
|
||||
|
||||
@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();
|
||||
|
||||
} else {
|
||||
|
||||
assertThat(repository.findCompletedPublications())
|
||||
.hasSize(1)
|
||||
.element(0)
|
||||
.extracting(TargetEventPublication::getEvent)
|
||||
.isEqualTo(event);
|
||||
@BeforeEach
|
||||
void setUp() {
|
||||
this.completionMode = CompletionMode.from(environment);
|
||||
this.repository = new MongoDbEventPublicationRepository(mongoTemplate, completionMode);
|
||||
}
|
||||
|
||||
}
|
||||
@AfterEach
|
||||
void tearDown() {
|
||||
mongoTemplate.remove(MongoDbEventPublication.class).all();
|
||||
mongoTemplate.remove(MongoDbEventPublication.class).inCollection(archiveCollection).all();
|
||||
}
|
||||
|
||||
@Test // GH-258
|
||||
void marksPublicationAsCompletedById() {
|
||||
@Test // GH-4
|
||||
void shouldPersistAndUpdateEventPublication() {
|
||||
|
||||
var event = new TestEvent("first");
|
||||
var publication = createPublication(event);
|
||||
var publication = createPublication(new TestEvent("abc"));
|
||||
|
||||
repository.markCompleted(publication.getIdentifier(), Instant.now());
|
||||
var eventPublications = repository.findIncompletePublications();
|
||||
|
||||
if (completionMode == CompletionMode.DELETE) {
|
||||
assertThat(eventPublications).hasSize(1);
|
||||
assertThat(eventPublications.get(0).getEvent()).isEqualTo(publication.getEvent());
|
||||
assertThat(eventPublications.get(0).getTargetIdentifier()).isEqualTo(publication.getTargetIdentifier());
|
||||
|
||||
assertThat(repository.findIncompletePublicationsByEventAndTargetIdentifier(new TestEvent("abc"), TARGET_IDENTIFIER))
|
||||
.isPresent();
|
||||
|
||||
// Complete publication
|
||||
repository.markCompleted(publication, Instant.now());
|
||||
|
||||
assertThat(repository.findIncompletePublications()).isEmpty();
|
||||
}
|
||||
|
||||
@Test // GH-4
|
||||
void shouldUpdateSingleEventPublication() {
|
||||
|
||||
var first = createPublication(new TestEvent("id1"));
|
||||
var second = createPublication(new TestEvent("id2"));
|
||||
|
||||
repository.markCompleted(second, Instant.now());
|
||||
|
||||
assertThat(repository.findIncompletePublications()).hasSize(1)
|
||||
.element(0)
|
||||
.extracting(TargetEventPublication::getEvent).isEqualTo(first.getEvent());
|
||||
}
|
||||
|
||||
@Test // GH-133
|
||||
void returnsOldestIncompletePublicationsFirst() {
|
||||
|
||||
var now = LocalDateTime.now();
|
||||
|
||||
savePublicationAt(now.withHour(3));
|
||||
savePublicationAt(now.withHour(0));
|
||||
savePublicationAt(now.withHour(1));
|
||||
|
||||
assertThat(repository.findIncompletePublications())
|
||||
.isSortedAccordingTo(Comparator.comparing(TargetEventPublication::getPublicationDate));
|
||||
}
|
||||
|
||||
@Test // GH-294
|
||||
void findsPublicationsOlderThanReference() throws Exception {
|
||||
|
||||
var first = createPublication(new TestEvent("first"));
|
||||
|
||||
Thread.sleep(100);
|
||||
|
||||
var now = Instant.now();
|
||||
var second = createPublication(new TestEvent("second"));
|
||||
|
||||
assertThat(repository.findIncompletePublications())
|
||||
.extracting(TargetEventPublication::getIdentifier)
|
||||
.containsExactly(first.getIdentifier(), second.getIdentifier());
|
||||
|
||||
assertThat(repository.findIncompletePublicationsPublishedBefore(now))
|
||||
.hasSize(1)
|
||||
.element(0).extracting(TargetEventPublication::getIdentifier).isEqualTo(first.getIdentifier());
|
||||
}
|
||||
|
||||
@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();
|
||||
|
||||
} 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());
|
||||
|
||||
assertThat(repository.findCompletedPublications()).isEmpty();
|
||||
assertThat(repository.findIncompletePublications()).isEmpty();
|
||||
|
||||
} else {
|
||||
if (completionMode == CompletionMode.DELETE) {
|
||||
|
||||
assertThat(repository.findCompletedPublications())
|
||||
.extracting(TargetEventPublication::getIdentifier)
|
||||
.containsExactly(publication.getIdentifier());
|
||||
assertThat(repository.findCompletedPublications()).isEmpty();
|
||||
|
||||
} else {
|
||||
|
||||
assertThat(repository.findCompletedPublications())
|
||||
.extracting(TargetEventPublication::getIdentifier)
|
||||
.containsExactly(publication.getIdentifier());
|
||||
}
|
||||
|
||||
if (completionMode == CompletionMode.ARCHIVE) {
|
||||
assertThat(mongoTemplate.findAll(MongoDbEventPublication.class, archiveCollection)).isNotEmpty();
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
private TargetEventPublication createPublication(Object event) {
|
||||
return createPublication(event, TARGET_IDENTIFIER);
|
||||
}
|
||||
|
||||
private TargetEventPublication createPublication(Object event, PublicationTargetIdentifier id) {
|
||||
return repository.create(TargetEventPublication.of(event, id));
|
||||
}
|
||||
|
||||
private void savePublicationAt(LocalDateTime date) {
|
||||
|
||||
mongoTemplate.save(
|
||||
new MongoDbEventPublication(UUID.randomUUID(), date.toInstant(ZoneOffset.UTC), "", "", null));
|
||||
}
|
||||
|
||||
@Nested
|
||||
class FindByEventAndTargetIdentifier {
|
||||
|
||||
@Test // GH-4
|
||||
void shouldFindEventPublicationByEventAndTargetIdentifier() {
|
||||
@@ -203,7 +195,7 @@ class MongoDbEventPublicationRepositoryTest {
|
||||
|
||||
var firstEvent = first.getEvent();
|
||||
|
||||
createPublication(firstEvent, PublicationTargetIdentifier.of("somethingDifferen"));
|
||||
createPublication(firstEvent, PublicationTargetIdentifier.of("somethingDifferent"));
|
||||
|
||||
var actual = repository.findIncompletePublicationsByEventAndTargetIdentifier(firstEvent, TARGET_IDENTIFIER);
|
||||
|
||||
@@ -250,10 +242,6 @@ class MongoDbEventPublicationRepositoryTest {
|
||||
assertThat(it.getPublicationDate()) //
|
||||
.isCloseTo(publication.getPublicationDate(), within(1, ChronoUnit.MILLIS)));
|
||||
}
|
||||
}
|
||||
|
||||
@Nested
|
||||
class DeleteCompletedPublications {
|
||||
|
||||
@Test // GH-20
|
||||
void shouldDeleteCompletedEvents() {
|
||||
@@ -304,11 +292,32 @@ class MongoDbEventPublicationRepositoryTest {
|
||||
.matches(it -> it.getIdentifier().equals(second.getIdentifier()))
|
||||
.matches(it -> it.getEvent().equals(second.getEvent()));
|
||||
}
|
||||
|
||||
private TargetEventPublication createPublication(Object event) {
|
||||
return createPublication(event, TARGET_IDENTIFIER);
|
||||
}
|
||||
|
||||
private TargetEventPublication createPublication(Object event, PublicationTargetIdentifier id) {
|
||||
return repository.create(TargetEventPublication.of(event, id));
|
||||
}
|
||||
|
||||
private void savePublicationAt(LocalDateTime date) {
|
||||
|
||||
mongoTemplate.save(
|
||||
new MongoDbEventPublication(UUID.randomUUID(), date.toInstant(ZoneOffset.UTC), "", "", null));
|
||||
}
|
||||
}
|
||||
|
||||
@Nested
|
||||
class WithUpdateCompletionTest extends TestBase {}
|
||||
|
||||
@Nested
|
||||
@TestPropertySource(properties = CompletionMode.PROPERTY + "=DELETE")
|
||||
static class WithDeleteCompletionTest extends MongoDbEventPublicationRepositoryTest {}
|
||||
class WithDeleteCompletionTest extends TestBase {}
|
||||
|
||||
@Nested
|
||||
@TestPropertySource(properties = CompletionMode.PROPERTY + "=ARCHIVE")
|
||||
class WithArchiveCompletionTest extends TestBase {}
|
||||
|
||||
private record TestEvent(String eventId) {}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user