diff --git a/spring-modulith-events/spring-modulith-events-jdbc/pom.xml b/spring-modulith-events/spring-modulith-events-jdbc/pom.xml
index 9e7ff3af..b76f7042 100644
--- a/spring-modulith-events/spring-modulith-events-jdbc/pom.xml
+++ b/spring-modulith-events/spring-modulith-events-jdbc/pom.xml
@@ -107,7 +107,6 @@
mysql
mysql-connector-java
- 8.0.29
test
diff --git a/spring-modulith-events/spring-modulith-events-jdbc/src/main/java/org/springframework/modulith/events/jdbc/DatabaseSchemaInitializer.java b/spring-modulith-events/spring-modulith-events-jdbc/src/main/java/org/springframework/modulith/events/jdbc/DatabaseSchemaInitializer.java
index 2cc439ad..2d693b51 100644
--- a/spring-modulith-events/spring-modulith-events-jdbc/src/main/java/org/springframework/modulith/events/jdbc/DatabaseSchemaInitializer.java
+++ b/spring-modulith-events/spring-modulith-events-jdbc/src/main/java/org/springframework/modulith/events/jdbc/DatabaseSchemaInitializer.java
@@ -15,19 +15,16 @@
*/
package org.springframework.modulith.events.jdbc;
+import lombok.RequiredArgsConstructor;
+
import java.io.IOException;
import java.io.UncheckedIOException;
import java.nio.charset.StandardCharsets;
-import lombok.RequiredArgsConstructor;
-
import org.springframework.beans.factory.InitializingBean;
-import org.springframework.boot.jdbc.DatabaseDriver;
-import org.springframework.context.ResourceLoaderAware;
import org.springframework.core.io.Resource;
import org.springframework.core.io.ResourceLoader;
import org.springframework.jdbc.core.JdbcTemplate;
-import org.springframework.jdbc.support.MetaDataAccessException;
import org.springframework.util.StreamUtils;
/**
@@ -38,20 +35,15 @@ import org.springframework.util.StreamUtils;
* @author Oliver Drotbohm
*/
@RequiredArgsConstructor
-class DatabaseSchemaInitializer implements ResourceLoaderAware, InitializingBean {
+class DatabaseSchemaInitializer implements InitializingBean {
private final JdbcTemplate jdbcTemplate;
+ private final ResourceLoader resourceLoader;
private final DatabaseType databaseType;
- private ResourceLoader resourceLoader;
-
- @Override
- public void setResourceLoader(ResourceLoader resourceLoader) {
- this.resourceLoader = resourceLoader;
- }
-
@Override
public void afterPropertiesSet() {
+
var schemaResourceFilename = databaseType.getSchemaResourceFilename();
var schemaDdlResource = resourceLoader.getResource(schemaResourceFilename);
var schemaDdl = asString(schemaDdlResource);
diff --git a/spring-modulith-events/spring-modulith-events-jdbc/src/main/java/org/springframework/modulith/events/jdbc/DatabaseType.java b/spring-modulith-events/spring-modulith-events-jdbc/src/main/java/org/springframework/modulith/events/jdbc/DatabaseType.java
index fdd3eadb..fae12409 100644
--- a/spring-modulith-events/spring-modulith-events-jdbc/src/main/java/org/springframework/modulith/events/jdbc/DatabaseType.java
+++ b/spring-modulith-events/spring-modulith-events-jdbc/src/main/java/org/springframework/modulith/events/jdbc/DatabaseType.java
@@ -19,10 +19,12 @@ import java.util.Map;
import java.util.UUID;
import org.springframework.boot.jdbc.DatabaseDriver;
+import org.springframework.util.Assert;
/**
* @author Dmitry Belyaev
* @author Björn Kieling
+ * @author Oliver Drotbohm
*/
enum DatabaseType {
@@ -31,10 +33,13 @@ enum DatabaseType {
H2("h2"),
MYSQL("mysql") {
+
+ @Override
Object uuidToDatabase(UUID id) {
return id.toString();
}
+ @Override
UUID databaseToUUID(Object id) {
return UUID.fromString(id.toString());
}
@@ -50,10 +55,13 @@ enum DatabaseType {
DatabaseDriver.MYSQL, MYSQL);
static DatabaseType from(DatabaseDriver databaseDriver) {
+
var databaseType = DATABASE_DRIVER_TO_DATABASE_TYPE_MAP.get(databaseDriver);
+
if (databaseType == null) {
throw new IllegalArgumentException("Unsupported database type: " + databaseDriver);
}
+
return databaseType;
}
@@ -68,6 +76,9 @@ enum DatabaseType {
}
UUID databaseToUUID(Object id) {
+
+ Assert.isInstanceOf(UUID.class, id, "Database value not of type UUID!");
+
return (UUID) id;
}
diff --git a/spring-modulith-events/spring-modulith-events-jdbc/src/main/java/org/springframework/modulith/events/jdbc/JdbcEventPublicationAutoConfiguration.java b/spring-modulith-events/spring-modulith-events-jdbc/src/main/java/org/springframework/modulith/events/jdbc/JdbcEventPublicationAutoConfiguration.java
index 0be176a5..13c34781 100644
--- a/spring-modulith-events/spring-modulith-events-jdbc/src/main/java/org/springframework/modulith/events/jdbc/JdbcEventPublicationAutoConfiguration.java
+++ b/spring-modulith-events/spring-modulith-events-jdbc/src/main/java/org/springframework/modulith/events/jdbc/JdbcEventPublicationAutoConfiguration.java
@@ -21,6 +21,7 @@ import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty;
import org.springframework.boot.jdbc.DatabaseDriver;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
+import org.springframework.core.io.ResourceLoader;
import org.springframework.jdbc.core.JdbcTemplate;
import org.springframework.modulith.events.EventSerializer;
import org.springframework.modulith.events.config.EventPublicationConfigurationExtension;
@@ -35,18 +36,21 @@ class JdbcEventPublicationAutoConfiguration implements EventPublicationConfigura
@Bean
DatabaseType databaseType(DataSource dataSource) {
- var databaseDriver = DatabaseDriver.fromDataSource(dataSource);
- return DatabaseType.from(databaseDriver);
+ return DatabaseType.from(DatabaseDriver.fromDataSource(dataSource));
}
+
@Bean
JdbcEventPublicationRepository jdbcEventPublicationRepository(JdbcTemplate jdbcTemplate,
EventSerializer serializer, DatabaseType databaseType) {
+
return new JdbcEventPublicationRepository(jdbcTemplate, serializer, databaseType);
}
@Bean
@ConditionalOnProperty(name = "spring.modulith.events.schema-initialization.enabled", havingValue = "true")
- DatabaseSchemaInitializer databaseSchemaInitializer(JdbcTemplate jdbcTemplate, DatabaseType databaseType) {
- return new DatabaseSchemaInitializer(jdbcTemplate, databaseType);
+ DatabaseSchemaInitializer databaseSchemaInitializer(JdbcTemplate jdbcTemplate, ResourceLoader resourceLoader,
+ DatabaseType databaseType) {
+
+ return new DatabaseSchemaInitializer(jdbcTemplate, resourceLoader, databaseType);
}
}
diff --git a/spring-modulith-events/spring-modulith-events-jdbc/src/main/java/org/springframework/modulith/events/jdbc/JdbcEventPublicationRepository.java b/spring-modulith-events/spring-modulith-events-jdbc/src/main/java/org/springframework/modulith/events/jdbc/JdbcEventPublicationRepository.java
index 4a07586d..6c879ec5 100644
--- a/spring-modulith-events/spring-modulith-events-jdbc/src/main/java/org/springframework/modulith/events/jdbc/JdbcEventPublicationRepository.java
+++ b/spring-modulith-events/spring-modulith-events-jdbc/src/main/java/org/springframework/modulith/events/jdbc/JdbcEventPublicationRepository.java
@@ -15,6 +15,11 @@
*/
package org.springframework.modulith.events.jdbc;
+import lombok.Builder;
+import lombok.EqualsAndHashCode;
+import lombok.RequiredArgsConstructor;
+import lombok.extern.slf4j.Slf4j;
+
import java.sql.ResultSet;
import java.sql.SQLException;
import java.sql.Timestamp;
@@ -35,11 +40,6 @@ import org.springframework.modulith.events.EventSerializer;
import org.springframework.modulith.events.PublicationTargetIdentifier;
import org.springframework.transaction.annotation.Transactional;
-import lombok.Builder;
-import lombok.EqualsAndHashCode;
-import lombok.RequiredArgsConstructor;
-import lombok.extern.slf4j.Slf4j;
-
/**
* JDBC-based repository to store {@link EventPublication}s.
*
@@ -74,16 +74,16 @@ class JdbcEventPublicationRepository implements EventPublicationRepository {
ORDER BY PUBLICATION_DATE
""";
- protected final JdbcOperations operations;
+ private final JdbcOperations operations;
private final EventSerializer serializer;
-
private final DatabaseType databaseType;
@Override
@Transactional
public EventPublication create(EventPublication publication) {
- String serializedEvent = serializeEvent(publication.getEvent());
+ var serializedEvent = serializeEvent(publication.getEvent());
+
operations.update( //
SQL_STATEMENT_INSERT, //
uuidToDatabase(UUID.randomUUID()), //
@@ -103,7 +103,7 @@ class JdbcEventPublicationRepository implements EventPublicationRepository {
var listenerId = publication.getTargetIdentifier().getValue();
var potentialPublicationIdsToBeUpdated = operations.query( //
SQL_STATEMENT_FIND_BY_EVENT_AND_LISTENER_ID, //
- (rs, rowNum) -> getUUIDFromResultSet(rs), //
+ (rs, rowNum) -> getUuidFromResultSet(rs), //
serializedEvent, //
listenerId);
@@ -119,8 +119,9 @@ class JdbcEventPublicationRepository implements EventPublicationRepository {
public Optional findIncompletePublicationsByEventAndTargetIdentifier( //
Object event, PublicationTargetIdentifier targetIdentifier) {
- String serializedEvent = serializeEvent(event);
- String listenerId = targetIdentifier.getValue();
+ var serializedEvent = serializeEvent(event);
+ var listenerId = targetIdentifier.getValue();
+
return findAllIncompletePublicationsByEventAndListenerId(serializedEvent, listenerId).stream() //
.findFirst();
}
@@ -129,21 +130,26 @@ class JdbcEventPublicationRepository implements EventPublicationRepository {
@Transactional(readOnly = true)
@SuppressWarnings("null")
public List findIncompletePublications() {
+
return operations.query( //
SQL_STATEMENT_FIND_UNCOMPLETED, //
this::resultSetToPublications);
}
private void update(UUID id, CompletableEventPublication publication) {
- Timestamp timestamp = publication.getCompletionDate().map(Timestamp::from).orElse(null);
+
+ var timestamp = publication.getCompletionDate().map(Timestamp::from).orElse(null);
+
operations.update( //
SQL_STATEMENT_UPDATE, //
timestamp, //
uuidToDatabase(id));
}
+ @SuppressWarnings("null")
private List findAllIncompletePublicationsByEventAndListenerId(
String serializedEvent, String listenerId) {
+
return operations.query( //
SQL_STATEMENT_FIND_BY_EVENT_AND_LISTENER_ID, //
this::resultSetToPublications, //
@@ -165,8 +171,11 @@ class JdbcEventPublicationRepository implements EventPublicationRepository {
private List resultSetToPublications(ResultSet resultSet) throws SQLException {
List result = new ArrayList<>();
+
while (resultSet.next()) {
- EventPublication publication = resultSetToPublication(resultSet);
+
+ var publication = resultSetToPublication(resultSet);
+
if (publication != null) {
result.add(publication);
}
@@ -185,9 +194,9 @@ class JdbcEventPublicationRepository implements EventPublicationRepository {
@Nullable
private EventPublication resultSetToPublication(ResultSet rs) throws SQLException {
- UUID id = getUUIDFromResultSet(rs);
-
+ var id = getUuidFromResultSet(rs);
var eventClass = loadClass(id, rs.getString("EVENT_TYPE"));
+
if (eventClass == null) {
return null;
}
@@ -207,9 +216,8 @@ class JdbcEventPublicationRepository implements EventPublicationRepository {
return databaseType.uuidToDatabase(id);
}
- private UUID getUUIDFromResultSet(ResultSet rs) throws SQLException {
- Object id = rs.getObject("ID");
- return databaseType.databaseToUUID(id);
+ private UUID getUuidFromResultSet(ResultSet rs) throws SQLException {
+ return databaseType.databaseToUUID(rs.getObject("ID"));
}
@Nullable
diff --git a/spring-modulith-events/spring-modulith-events-jdbc/src/test/java/org/springframework/modulith/events/jdbc/DatabaseTypeTest.java b/spring-modulith-events/spring-modulith-events-jdbc/src/test/java/org/springframework/modulith/events/jdbc/DatabaseTypeUnitTests.java
similarity index 54%
rename from spring-modulith-events/spring-modulith-events-jdbc/src/test/java/org/springframework/modulith/events/jdbc/DatabaseTypeTest.java
rename to spring-modulith-events/spring-modulith-events-jdbc/src/test/java/org/springframework/modulith/events/jdbc/DatabaseTypeUnitTests.java
index 4c0ab875..9ec52f5d 100644
--- a/spring-modulith-events/spring-modulith-events-jdbc/src/test/java/org/springframework/modulith/events/jdbc/DatabaseTypeTest.java
+++ b/spring-modulith-events/spring-modulith-events-jdbc/src/test/java/org/springframework/modulith/events/jdbc/DatabaseTypeUnitTests.java
@@ -5,12 +5,13 @@ import static org.assertj.core.api.Assertions.*;
import org.junit.jupiter.api.Test;
import org.springframework.boot.jdbc.DatabaseDriver;
-class DatabaseTypeTest {
+class DatabaseTypeUnitTests {
- @Test
+ @Test // GH-29
void shouldThrowExceptionOnUnsupportedDatabaseType() {
- assertThatThrownBy(() -> DatabaseType.from(DatabaseDriver.UNKNOWN))
- .isInstanceOf(IllegalArgumentException.class)
- .hasMessageContaining("UNKNOWN");
+
+ assertThatExceptionOfType(IllegalArgumentException.class)
+ .isThrownBy(() -> DatabaseType.from(DatabaseDriver.UNKNOWN))
+ .withMessageContaining("UNKNOWN");
}
}