diff --git a/NEWS.md b/NEWS.md index f6bd8b7d..da12ee02 100644 --- a/NEWS.md +++ b/NEWS.md @@ -1,3 +1,8 @@ +## 2026-07-28 v3.7.0 + +### Stories +* Publish Kafka domain events (CREATE/UPDATE) for Export Configuration changes on `folio.ALL.data-export.config` with structural credential redaction + ## 2026-05-19 v3.6.1 [Full Changelog](https://github.com/folio-org/mod-data-export-spring/compare/v3.6.0...v3.6.1) diff --git a/descriptors/ModuleDescriptor-template.json b/descriptors/ModuleDescriptor-template.json index a36e45af..45dce7f4 100644 --- a/descriptors/ModuleDescriptor-template.json +++ b/descriptors/ModuleDescriptor-template.json @@ -555,6 +555,11 @@ "name": "ENV", "value": "folio" }, + { + "name": "KAFKA_PRODUCER_TENANT_COLLECTION", + "value": "ALL", + "description": "Tenant partitioning mode for Kafka topics. 'ALL' (default, matches folio-kafka-wrapper) publishes every tenant's events to a single topic per entity (e.g. 'folio.ALL.data-export.config'); consumers route by the 'tenant' field in the event envelope. Any other value falls back to per-tenant topics." + }, { "name": "JOB_EXPIRATION_PERIOD_DAYS", "value": "7" diff --git a/src/main/java/org/folio/des/config/ServiceConfiguration.java b/src/main/java/org/folio/des/config/ServiceConfiguration.java index 42f037c0..7eca6a6f 100644 --- a/src/main/java/org/folio/des/config/ServiceConfiguration.java +++ b/src/main/java/org/folio/des/config/ServiceConfiguration.java @@ -28,6 +28,7 @@ import org.folio.des.scheduling.quartz.converter.acquisition.ExportConfigToEdifactJobDetailConverter; import org.folio.des.scheduling.quartz.converter.acquisition.ExportConfigToEdifactTriggerConverter; import org.folio.des.scheduling.quartz.job.acquisition.EdifactJobKeyResolver; +import org.folio.des.service.config.ExportConfigDomainEventService; import org.folio.des.service.config.ExportConfigService; import org.folio.des.service.config.acquisition.ClaimsExportService; import org.folio.des.service.config.acquisition.EdifactOrdersExportService; @@ -82,8 +83,10 @@ BursarFeesFinesExportConfigService bursarExportConfigService(ExportConfigReposit DefaultExportConfigMapper defaultExportConfigMapper, ExportConfigMapperResolver exportConfigMapperResolver, ExportConfigValidatorResolver exportConfigValidatorResolver, + ExportConfigDomainEventService exportConfigDomainEventService, BursarExportScheduler bursarExportScheduler) { - return new BursarFeesFinesExportConfigService(repository, defaultExportConfigMapper, exportConfigMapperResolver, exportConfigValidatorResolver, bursarExportScheduler); + return new BursarFeesFinesExportConfigService(repository, defaultExportConfigMapper, exportConfigMapperResolver, + exportConfigValidatorResolver, exportConfigDomainEventService, bursarExportScheduler); } @Bean @@ -91,24 +94,30 @@ EdifactOrdersExportService edifactOrdersExportService(ExportConfigRepository rep EdifactExportConfigMapper edifactExportConfigMapper, ExportConfigMapperResolver exportConfigMapperResolver, ExportConfigValidatorResolver exportConfigValidatorResolver, + ExportConfigDomainEventService exportConfigDomainEventService, ExportJobScheduler exportJobScheduler) { - return new EdifactOrdersExportService(repository, edifactExportConfigMapper, exportConfigMapperResolver, exportConfigValidatorResolver, exportJobScheduler); + return new EdifactOrdersExportService(repository, edifactExportConfigMapper, exportConfigMapperResolver, + exportConfigValidatorResolver, exportConfigDomainEventService, exportJobScheduler); } @Bean ClaimsExportService claimsExportService(ExportConfigRepository repository, ClaimsExportConfigMapper claimsExportConfigMapper, ExportConfigMapperResolver exportConfigMapperResolver, - ExportConfigValidatorResolver exportConfigValidatorResolver) { - return new ClaimsExportService(repository, claimsExportConfigMapper, exportConfigMapperResolver, exportConfigValidatorResolver); + ExportConfigValidatorResolver exportConfigValidatorResolver, + ExportConfigDomainEventService exportConfigDomainEventService) { + return new ClaimsExportService(repository, claimsExportConfigMapper, exportConfigMapperResolver, + exportConfigValidatorResolver, exportConfigDomainEventService); } @Bean BaseExportConfigService defaultExportConfigService(ExportConfigRepository repository, DefaultExportConfigMapper defaultExportConfigMapper, ExportConfigMapperResolver exportConfigMapperResolver, - ExportConfigValidatorResolver exportConfigValidatorResolver) { - return new BaseExportConfigService(repository, defaultExportConfigMapper, exportConfigMapperResolver, exportConfigValidatorResolver); + ExportConfigValidatorResolver exportConfigValidatorResolver, + ExportConfigDomainEventService exportConfigDomainEventService) { + return new BaseExportConfigService(repository, defaultExportConfigMapper, exportConfigMapperResolver, + exportConfigValidatorResolver, exportConfigDomainEventService); } @Bean diff --git a/src/main/java/org/folio/des/config/kafka/KafkaConfiguration.java b/src/main/java/org/folio/des/config/kafka/KafkaConfiguration.java index 5da71871..f380e3de 100644 --- a/src/main/java/org/folio/des/config/kafka/KafkaConfiguration.java +++ b/src/main/java/org/folio/des/config/kafka/KafkaConfiguration.java @@ -21,9 +21,7 @@ import org.springframework.kafka.core.ProducerFactory; import org.springframework.kafka.support.serializer.JsonDeserializer; import org.springframework.kafka.support.serializer.JsonSerializer; -import org.springframework.stereotype.Component; -@Component @Configuration @RequiredArgsConstructor public class KafkaConfiguration { diff --git a/src/main/java/org/folio/des/config/kafka/KafkaService.java b/src/main/java/org/folio/des/config/kafka/KafkaService.java index bde58748..cdb38152 100644 --- a/src/main/java/org/folio/des/config/kafka/KafkaService.java +++ b/src/main/java/org/folio/des/config/kafka/KafkaService.java @@ -41,7 +41,8 @@ public class KafkaService { @RequiredArgsConstructor @Getter public enum Topic { - JOB_COMMAND("data-export.job.command"); + JOB_COMMAND("data-export.job.command"), + CONFIG("data-export.config"); private final String topicName; } @@ -84,13 +85,13 @@ private NewTopic toKafkaTopic(String tenant, Topic topic) { } /** - * Returns topic name in the format - `{env}.{tenant}.topicName` + * Returns the tenant-scoped topic name in the format {@code {env}.{tenant}.topicName}. * - * @param topicName initial topic name as {@link String} - * @param tenantId tenant id as {@link String} - * @return topic name as {@link String} object + * @param topicName the logical topic name + * @param tenantId the tenant id + * @return the fully-qualified tenant topic name */ - private String getTenantTopicName(String topicName, String tenantId) { + public String getTenantTopicName(String topicName, String tenantId) { return KafkaUtils.getTenantTopicName(topicName, environment, tenantId); } diff --git a/src/main/java/org/folio/des/domain/dto/event/DomainEvent.java b/src/main/java/org/folio/des/domain/dto/event/DomainEvent.java new file mode 100644 index 00000000..48487fed --- /dev/null +++ b/src/main/java/org/folio/des/domain/dto/event/DomainEvent.java @@ -0,0 +1,67 @@ +package org.folio.des.domain.dto.event; + +import java.util.UUID; + +import com.fasterxml.jackson.annotation.JsonInclude; +import com.fasterxml.jackson.annotation.JsonProperty; + +import lombok.AllArgsConstructor; +import lombok.Builder; +import lombok.Data; +import lombok.NoArgsConstructor; + +@Data +@Builder +@NoArgsConstructor +@AllArgsConstructor +@JsonInclude(JsonInclude.Include.NON_NULL) +public class DomainEvent { + + private UUID eventId; + private long eventTs; + private String tenant; + private DomainEventType type; + + @JsonProperty("old") + private T oldValue; + + @JsonProperty("new") + private T newValue; + + /** + * Builds a {@code CREATE} event carrying only the post-change snapshot. + * + * @param newValue the newly-created snapshot + * @param tenant the tenant the change happened in + * @return a populated {@code CREATE} domain event + */ + public static DomainEvent createEvent(T newValue, String tenant) { + return DomainEvent.builder() + .eventId(UUID.randomUUID()) + .eventTs(System.currentTimeMillis()) + .tenant(tenant) + .type(DomainEventType.CREATE) + .newValue(newValue) + .build(); + } + + /** + * Builds an {@code UPDATE} event carrying both the pre- and post-change snapshots. + * + * @param oldValue the pre-change snapshot + * @param newValue the post-change snapshot + * @param tenant the tenant the change happened in + * @return a populated {@code UPDATE} domain event + */ + public static DomainEvent updateEvent(T oldValue, T newValue, String tenant) { + return DomainEvent.builder() + .eventId(UUID.randomUUID()) + .eventTs(System.currentTimeMillis()) + .tenant(tenant) + .type(DomainEventType.UPDATE) + .oldValue(oldValue) + .newValue(newValue) + .build(); + } +} + diff --git a/src/main/java/org/folio/des/domain/dto/event/DomainEventType.java b/src/main/java/org/folio/des/domain/dto/event/DomainEventType.java new file mode 100644 index 00000000..fc08aca2 --- /dev/null +++ b/src/main/java/org/folio/des/domain/dto/event/DomainEventType.java @@ -0,0 +1,7 @@ +package org.folio.des.domain.dto.event; + +public enum DomainEventType { + CREATE, + UPDATE +} + diff --git a/src/main/java/org/folio/des/service/config/ExportConfigDomainEventService.java b/src/main/java/org/folio/des/service/config/ExportConfigDomainEventService.java new file mode 100644 index 00000000..5b99fea0 --- /dev/null +++ b/src/main/java/org/folio/des/service/config/ExportConfigDomainEventService.java @@ -0,0 +1,98 @@ +package org.folio.des.service.config; + +import java.util.Optional; +import java.util.UUID; + +import org.apache.commons.lang3.StringUtils; +import org.folio.des.domain.dto.ExportConfig; +import org.folio.des.domain.dto.ExportTypeSpecificParameters; +import org.folio.des.domain.dto.VendorEdiOrdersExportConfig; +import org.folio.des.domain.dto.event.DomainEvent; +import org.folio.spring.FolioExecutionContext; +import org.springframework.beans.factory.annotation.Qualifier; +import org.springframework.stereotype.Service; + +import com.fasterxml.jackson.databind.ObjectMapper; + +import lombok.extern.log4j.Log4j2; + +/** + * Builds {@link DomainEvent} envelopes for Export Configuration changes and hands them to the + * {@link ExportConfigEventProducer}. Config services delegate here so they never build events directly. + * + *

Every snapshot is {@link #sanitize(ExportConfig) sanitized} before it leaves this service — the FTP + * password (and any other credential-shaped field) is stripped from a deep copy so it can never reach the event + * bus. The producer's {@code ObjectMapper} omits empty values, so the stripped fields do not appear at all + * (not even as {@code null}) in the serialized payload.

+ */ +@Log4j2 +@Service +public class ExportConfigDomainEventService { + + private final ExportConfigEventProducer exportConfigEventProducer; + private final FolioExecutionContext folioExecutionContext; + private final ObjectMapper objectMapper; + + public ExportConfigDomainEventService(ExportConfigEventProducer exportConfigEventProducer, + FolioExecutionContext folioExecutionContext, + @Qualifier("entityObjectMapper") ObjectMapper objectMapper) { + this.exportConfigEventProducer = exportConfigEventProducer; + this.folioExecutionContext = folioExecutionContext; + this.objectMapper = objectMapper; + } + + /** + * Publishes a {@code CREATE} Export Configuration event carrying the new snapshot. + * + * @param newConfig the newly-created configuration snapshot + */ + public void publishConfigCreatedEvent(ExportConfig newConfig) { + if (newConfig == null || StringUtils.isBlank(newConfig.getId())) { + log.warn("publishConfigCreatedEvent:: skipping CREATE event, config id is missing"); + return; + } + var event = DomainEvent.createEvent(sanitize(newConfig), folioExecutionContext.getTenantId()); + publish(newConfig.getId(), event); + } + + /** + * Publishes an {@code UPDATE} Export Configuration event carrying both the pre- and post-change snapshots. + * + * @param oldConfig the pre-change configuration snapshot + * @param newConfig the post-change configuration snapshot + */ + public void publishConfigUpdatedEvent(ExportConfig oldConfig, ExportConfig newConfig) { + if (newConfig == null || StringUtils.isBlank(newConfig.getId())) { + log.warn("publishConfigUpdatedEvent:: skipping UPDATE event, config id is missing"); + return; + } + var event = DomainEvent.updateEvent(sanitize(oldConfig), sanitize(newConfig), folioExecutionContext.getTenantId()); + publish(newConfig.getId(), event); + } + + private void publish(String configId, DomainEvent event) { + log.debug("publish:: publishing config event [id: {}, type: {}, tenant: {}]", + configId, event.getType(), event.getTenant()); + exportConfigEventProducer.publish(UUID.fromString(configId), event); + } + + /** + * Returns a deep copy of the given configuration with all credential-shaped fields removed, safe to publish on + * the event bus. The original object is left untouched so the REST response still carries the full data. + * + * @param config the configuration snapshot to sanitize + * @return a sanitized deep copy, or {@code null} if the input is {@code null} + */ + private ExportConfig sanitize(ExportConfig config) { + if (config == null) { + return null; + } + var sanitized = objectMapper.convertValue(config, ExportConfig.class); + Optional.ofNullable(sanitized) + .map(ExportConfig::getExportTypeSpecificParameters) + .map(ExportTypeSpecificParameters::getVendorEdiOrdersExportConfig) + .map(VendorEdiOrdersExportConfig::getEdiFtp) + .ifPresent(ediFtp -> ediFtp.setPassword(null)); + return sanitized; + } +} \ No newline at end of file diff --git a/src/main/java/org/folio/des/service/config/ExportConfigEventProducer.java b/src/main/java/org/folio/des/service/config/ExportConfigEventProducer.java new file mode 100644 index 00000000..bd23d8cb --- /dev/null +++ b/src/main/java/org/folio/des/service/config/ExportConfigEventProducer.java @@ -0,0 +1,48 @@ +package org.folio.des.service.config; + +import java.util.UUID; + +import org.folio.des.config.kafka.KafkaService; +import org.folio.des.config.kafka.KafkaService.Topic; +import org.folio.des.domain.dto.ExportConfig; +import org.folio.des.domain.dto.event.DomainEvent; +import org.springframework.kafka.core.KafkaTemplate; +import org.springframework.stereotype.Component; + +import lombok.RequiredArgsConstructor; +import lombok.extern.log4j.Log4j2; + +/** + * Publishes Export Configuration {@link DomainEvent}s to the tenant-scoped Kafka topic. + * + *

The topic ({@code {env}.{tenant}.data-export-spring.config}) is created on tenant enable by the tenant service + * via {@link KafkaService}. Publishing failures are caught and logged at ERROR — they never fail the originating + * REST request or roll back the DB transaction.

+ */ +@Component +@Log4j2 +@RequiredArgsConstructor +public class ExportConfigEventProducer { + + private final KafkaTemplate kafkaTemplate; + private final KafkaService kafkaService; + + /** + * Publishes an Export Configuration domain event to the tenant-scoped topic. The tenant is taken from the + * event envelope so the topic resolution and the {@code tenant} field on the payload can never diverge. + * + * @param configId the export configuration id, used as the Kafka record key + * @param event the domain event envelope to publish + */ + public void publish(UUID configId, DomainEvent event) { + var topic = kafkaService.getTenantTopicName(Topic.CONFIG.getTopicName(), event.getTenant()); + var key = configId.toString(); + try { + log.info("publish:: Publishing {} event for config id={} on topic={}", event.getType(), configId, topic); + kafkaTemplate.send(topic, key, event); + log.info("publish:: Successfully published {} event for config id={}", event.getType(), configId); + } catch (Exception e) { + log.error("publish:: Failed to publish {} event for config id={} on topic={}", event.getType(), configId, topic, e); + } + } +} \ No newline at end of file diff --git a/src/main/java/org/folio/des/service/config/acquisition/ClaimsExportService.java b/src/main/java/org/folio/des/service/config/acquisition/ClaimsExportService.java index c6acdf78..a3246304 100644 --- a/src/main/java/org/folio/des/service/config/acquisition/ClaimsExportService.java +++ b/src/main/java/org/folio/des/service/config/acquisition/ClaimsExportService.java @@ -3,11 +3,11 @@ import lombok.extern.log4j.Log4j2; import org.folio.des.mapper.BaseExportConfigMapper; -import org.folio.des.mapper.DefaultExportConfigMapper; import org.folio.des.mapper.ExportConfigMapperResolver; import org.folio.des.domain.dto.ExportConfig; import org.folio.des.domain.dto.ExportTypeSpecificParameters; import org.folio.des.repository.ExportConfigRepository; +import org.folio.des.service.config.ExportConfigDomainEventService; import org.folio.des.service.config.impl.BaseExportConfigService; import org.folio.des.validator.ExportConfigValidatorResolver; @@ -18,8 +18,9 @@ public class ClaimsExportService extends BaseExportConfigService { public ClaimsExportService(ExportConfigRepository repository, BaseExportConfigMapper defaultExportConfigMapper, - ExportConfigMapperResolver exportConfigMapperResolver, ExportConfigValidatorResolver exportConfigValidatorResolver) { - super(repository, defaultExportConfigMapper, exportConfigMapperResolver, exportConfigValidatorResolver); + ExportConfigMapperResolver exportConfigMapperResolver, ExportConfigValidatorResolver exportConfigValidatorResolver, + ExportConfigDomainEventService exportConfigDomainEventService) { + super(repository, defaultExportConfigMapper, exportConfigMapperResolver, exportConfigValidatorResolver, exportConfigDomainEventService); } @Override diff --git a/src/main/java/org/folio/des/service/config/acquisition/EdifactOrdersExportService.java b/src/main/java/org/folio/des/service/config/acquisition/EdifactOrdersExportService.java index 93bd10b8..efcbda6b 100644 --- a/src/main/java/org/folio/des/service/config/acquisition/EdifactOrdersExportService.java +++ b/src/main/java/org/folio/des/service/config/acquisition/EdifactOrdersExportService.java @@ -4,12 +4,12 @@ import java.util.UUID; import org.folio.des.mapper.BaseExportConfigMapper; -import org.folio.des.mapper.DefaultExportConfigMapper; import org.folio.des.mapper.ExportConfigMapperResolver; import org.folio.des.domain.dto.ExportConfig; import org.folio.des.domain.dto.ExportTypeSpecificParameters; import org.folio.des.repository.ExportConfigRepository; import org.folio.des.scheduling.ExportJobScheduler; +import org.folio.des.service.config.ExportConfigDomainEventService; import org.folio.des.service.config.impl.BaseExportConfigService; import org.folio.des.validator.ExportConfigValidatorResolver; @@ -22,8 +22,9 @@ public class EdifactOrdersExportService extends BaseExportConfigService { public EdifactOrdersExportService(ExportConfigRepository repository, BaseExportConfigMapper defaultExportConfigMapper, ExportConfigMapperResolver exportConfigMapperResolver, ExportConfigValidatorResolver exportConfigValidatorResolver, + ExportConfigDomainEventService exportConfigDomainEventService, ExportJobScheduler exportJobScheduler) { - super(repository, defaultExportConfigMapper, exportConfigMapperResolver, exportConfigValidatorResolver); + super(repository, defaultExportConfigMapper, exportConfigMapperResolver, exportConfigValidatorResolver, exportConfigDomainEventService); this.exportJobScheduler = exportJobScheduler; } diff --git a/src/main/java/org/folio/des/service/config/impl/BaseExportConfigService.java b/src/main/java/org/folio/des/service/config/impl/BaseExportConfigService.java index 56e12e15..3da93c29 100644 --- a/src/main/java/org/folio/des/service/config/impl/BaseExportConfigService.java +++ b/src/main/java/org/folio/des/service/config/impl/BaseExportConfigService.java @@ -13,6 +13,7 @@ import org.folio.des.domain.dto.ExportConfigCollection; import org.folio.des.domain.dto.ExportTypeSpecificParameters; import org.folio.des.repository.ExportConfigRepository; +import org.folio.des.service.config.ExportConfigDomainEventService; import org.folio.des.service.config.ExportConfigService; import org.folio.des.validator.ExportConfigValidatorResolver; import org.folio.spring.exception.NotFoundException; @@ -33,17 +34,22 @@ public class BaseExportConfigService implements ExportConfigService { protected final BaseExportConfigMapper exportConfigMapper; protected final ExportConfigMapperResolver exportConfigMapperResolver; protected final ExportConfigValidatorResolver exportConfigValidatorResolver; + protected final ExportConfigDomainEventService exportConfigDomainEventService; @Override @Transactional public void updateConfig(String configId, ExportConfig exportConfig) { log.info("updateConfig:: configId={}, exportConfig={}", configId, exportConfig); validateIncomingExportConfig(exportConfig); - getExportConfigEntityOrThrow(configId); + var existingEntity = getExportConfigEntityOrThrow(configId); + var oldSnapshot = toDto(existingEntity); var entity = exportConfigMapper.toEntity(exportConfig); - repository.save(entity); + entity = repository.save(entity); log.info("updateConfig:: Successfully updated config with id={}", configId); + + var newSnapshot = toDto(entity); + exportConfigDomainEventService.publishConfigUpdatedEvent(oldSnapshot, newSnapshot); } @Override @@ -54,9 +60,12 @@ public ExportConfig postConfig(ExportConfig exportConfig) { var entity = exportConfigMapper.toEntity(exportConfig); entity = repository.save(entity); - log.info("postConfig:: Successfully created config with id={}", exportConfig.getId()); + log.info("postConfig:: Successfully created config with id={}", entity.getId()); + + var savedConfig = toDto(entity); + exportConfigDomainEventService.publishConfigCreatedEvent(savedConfig); - return toDto(entity); + return savedConfig; } @Override diff --git a/src/main/java/org/folio/des/service/config/impl/BursarFeesFinesExportConfigService.java b/src/main/java/org/folio/des/service/config/impl/BursarFeesFinesExportConfigService.java index c160ad6d..959b3fc9 100644 --- a/src/main/java/org/folio/des/service/config/impl/BursarFeesFinesExportConfigService.java +++ b/src/main/java/org/folio/des/service/config/impl/BursarFeesFinesExportConfigService.java @@ -14,6 +14,7 @@ import org.folio.des.domain.dto.ExportConfig; import org.folio.des.domain.dto.ExportConfigCollection; import org.folio.des.scheduling.bursar.BursarExportScheduler; +import org.folio.des.service.config.ExportConfigDomainEventService; import org.folio.des.validator.ExportConfigValidatorResolver; import org.springframework.data.domain.PageRequest; @@ -26,8 +27,9 @@ public class BursarFeesFinesExportConfigService extends BaseExportConfigService public BursarFeesFinesExportConfigService(ExportConfigRepository repository, DefaultExportConfigMapper defaultExportConfigMapper, ExportConfigMapperResolver exportConfigMapperResolver, ExportConfigValidatorResolver exportConfigValidatorResolver, + ExportConfigDomainEventService exportConfigDomainEventService, BursarExportScheduler bursarExportScheduler) { - super(repository, defaultExportConfigMapper, exportConfigMapperResolver, exportConfigValidatorResolver); + super(repository, defaultExportConfigMapper, exportConfigMapperResolver, exportConfigValidatorResolver, exportConfigDomainEventService); this.bursarExportScheduler = bursarExportScheduler; } diff --git a/src/main/resources/application.yml b/src/main/resources/application.yml index 37ada618..75cfccf1 100644 --- a/src/main/resources/application.yml +++ b/src/main/resources/application.yml @@ -41,6 +41,10 @@ spring: enabled: true kafka: bootstrap-servers: ${KAFKA_HOST:localhost}:${KAFKA_PORT:9092} + producer: + client-id: mod-data-export-spring + acks: all + retries: 3 datasource: username: ${DB_USERNAME:folio_admin} password: ${DB_PASSWORD:folio_admin} diff --git a/src/test/java/org/folio/des/controller/ConfigsControllerTest.java b/src/test/java/org/folio/des/controller/ConfigsControllerTest.java index e8706700..1742dd6c 100644 --- a/src/test/java/org/folio/des/controller/ConfigsControllerTest.java +++ b/src/test/java/org/folio/des/controller/ConfigsControllerTest.java @@ -1,7 +1,9 @@ package org.folio.des.controller; +import static org.assertj.core.api.Assertions.assertThat; import static org.hamcrest.Matchers.is; import static org.hamcrest.Matchers.startsWith; +import static org.junit.jupiter.api.Assertions.assertNull; import static org.mockito.ArgumentMatchers.any; import static org.mockito.ArgumentMatchers.eq; import static org.mockito.Mockito.reset; @@ -14,39 +16,62 @@ import static org.springframework.test.web.servlet.result.MockMvcResultMatchers.jsonPath; import static org.springframework.test.web.servlet.result.MockMvcResultMatchers.status; +import java.time.Duration; import java.util.Objects; +import org.apache.kafka.clients.consumer.ConsumerRecord; +import org.folio.des.config.kafka.KafkaService; +import org.folio.des.config.kafka.KafkaService.Topic; import org.folio.des.domain.dto.ExportConfig; +import org.folio.des.domain.dto.event.DomainEvent; +import org.folio.des.domain.dto.event.DomainEventType; import org.folio.des.repository.ExportConfigRepository; import org.folio.des.scheduling.bursar.BursarExportScheduler; import org.folio.des.support.BaseTest; +import org.folio.des.support.TestKafkaConsumer; import org.junit.jupiter.api.AfterEach; import org.junit.jupiter.api.DisplayName; import org.junit.jupiter.api.Test; import org.junit.jupiter.params.ParameterizedTest; import org.junit.jupiter.params.provider.CsvSource; import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.boot.kafka.autoconfigure.KafkaProperties; import org.springframework.http.MediaType; import org.springframework.test.context.TestPropertySource; import org.springframework.test.context.bean.override.mockito.MockitoSpyBean; import org.springframework.test.web.servlet.MockMvc; +import com.fasterxml.jackson.core.type.TypeReference; + +import lombok.SneakyThrows; + @TestPropertySource(properties = "spring.jpa.properties.hibernate.default_schema=diku_mod_data_export_spring") class ConfigsControllerTest extends BaseTest { private static final String NEW_CONFIG_REQUEST = "{\"id\":\"0a3cba78-16e7-498e-b75b-98713000277b\",\"type\":\"BURSAR_FEES_FINES\",\"exportTypeSpecificParameters\":{\"bursarFeeFines\":{\"filter\":{\"type\":\"Pass\"},\"groupByPatron\":false,\"header\":[],\"data\":[],\"footer\":[],\"transferInfo\":{\"conditions\":[],\"else\":{\"account\":\"90c1820f-60bf-4b9a-99f5-d677ea78ddca\"}}}},\"scheduleFrequency\":5,\"schedulePeriod\":\"DAY\",\"scheduleTime\":\"00:20:00.000Z\"}"; + // schedulePeriod deliberately differs from NEW_CONFIG_REQUEST (DAY -> HOUR) so the UPDATE event can be + // asserted to carry an old snapshot that differs from the new one (AC2). private static final String UPDATE_CONFIG_REQUEST = - "{\"id\":\"0a3cba78-16e7-498e-b75b-98713000277b\",\"type\":\"BURSAR_FEES_FINES\",\"exportTypeSpecificParameters\":{\"bursarFeeFines\":{\"filter\":{\"type\":\"Pass\"},\"groupByPatron\":false,\"header\":[],\"data\":[],\"footer\":[],\"transferInfo\":{\"conditions\":[],\"else\":{\"account\":\"90c1820f-60bf-4b9a-99f5-d677ea78ddca\"}}}},\"scheduleFrequency\":5,\"schedulePeriod\":\"DAY\",\"scheduleTime\":\"00:20:00.000Z\"}"; + "{\"id\":\"0a3cba78-16e7-498e-b75b-98713000277b\",\"type\":\"BURSAR_FEES_FINES\",\"exportTypeSpecificParameters\":{\"bursarFeeFines\":{\"filter\":{\"type\":\"Pass\"},\"groupByPatron\":false,\"header\":[],\"data\":[],\"footer\":[],\"transferInfo\":{\"conditions\":[],\"else\":{\"account\":\"90c1820f-60bf-4b9a-99f5-d677ea78ddca\"}}}},\"scheduleFrequency\":5,\"schedulePeriod\":\"HOUR\",\"scheduleTime\":\"00:20:00.000Z\"}"; + private static final String FAILED_CONFIG_ID = "c8303ff3-7dec-49a1-acc8-7ce4f311fe21"; private static final String UPDATE_CONFIG_REQUEST_FAILED = - "{\"id\":\"0a3cba78-16e7-498e-b75b-98713000277b\",\"type\":\"BURSAR_FEES_FINES\",\"scheduleFrequency\":5,\"schedulePeriod\":\"DAY\",\"scheduleTime\":\"00:20:00.000Z\"}"; + "{\"id\":\"c8303ff3-7dec-49a1-acc8-7ce4f311fe21\",\"type\":\"BURSAR_FEES_FINES\",\"scheduleFrequency\":5,\"schedulePeriod\":\"DAY\",\"scheduleTime\":\"00:20:00.000Z\"}"; private static final String EDIFACT_CONFIG_REQUEST = "{\"id\":\"5a3cba28-16e7-498e-b73b-98713000298e\", \"type\": \"EDIFACT_ORDERS_EXPORT\", \"exportTypeSpecificParameters\": { \"vendorEdiOrdersExportConfig\": {\"vendorId\": \"046b6c7f-0b8a-43b9-b35d-6489e6daee91\", \"configName\": \"edi_config\", \"integrationType\": \"Ordering\", \"fileFormat\": \"CSV\", \"transmissionMethod\": \"File download\", \"ediSchedule\": {\"enableScheduledExport\": true, \"scheduleParameters\": {\"scheduleFrequency\": 1, \"schedulePeriod\": \"HOUR\", \"scheduleTime\": \"15:30:00\"}}}}, \"schedulePeriod\": \"HOUR\"}"; private static final String CLAIMS_REQUEST = "{\"id\":\"30ad9c6d-f2e7-425f-a171-b4e0cbce7204\",\"type\":\"CLAIMS\",\"tenant\":\"diku\",\"exportTypeSpecificParameters\":{\"vendorEdiOrdersExportConfig\":{\"exportConfigId\":\"30ad9c6d-f2e7-425f-a171-b4e0cbce7204\",\"vendorId\":\"1e958895-82a6-4fa1-b6fe-763063381946\",\"configName\":\"Test 1-3\",\"ediConfig\":{\"accountNoList\":[\"3\"],\"ediNamingConvention\":\"{organizationCode}-{integrationName}-{exportJobEndDate}\",\"libEdiType\":\"31B/US-SAN\",\"vendorEdiType\":\"31B/US-SAN\",\"sendAccountNumber\":false,\"supportOrder\":false,\"supportInvoice\":false},\"ediFtp\":{\"ftpConnMode\":\"Active\",\"ftpFormat\":\"SFTP\",\"ftpMode\":\"ASCII\"},\"isDefaultConfig\":false,\"integrationType\":\"Claiming\",\"transmissionMethod\":\"File download\",\"fileFormat\":\"CSV\"}},\"schedulePeriod\":\"NONE\"}"; + private static final String PASSWORD_SENTINEL = "AC3-sentinel-12345"; + private static final String CLAIMS_WITH_PASSWORD_ID = "12345678-1234-1234-1234-1234567890ab"; + private static final String CLAIMS_REQUEST_WITH_PASSWORD = + "{\"id\":\"12345678-1234-1234-1234-1234567890ab\",\"type\":\"CLAIMS\",\"tenant\":\"diku\",\"exportTypeSpecificParameters\":{\"vendorEdiOrdersExportConfig\":{\"exportConfigId\":\"12345678-1234-1234-1234-1234567890ab\",\"vendorId\":\"1e958895-82a6-4fa1-b6fe-763063381946\",\"configName\":\"Test 1-3\",\"ediConfig\":{\"accountNoList\":[\"3\"],\"ediNamingConvention\":\"{organizationCode}-{integrationName}-{exportJobEndDate}\",\"libEdiType\":\"31B/US-SAN\",\"vendorEdiType\":\"31B/US-SAN\",\"sendAccountNumber\":false,\"supportOrder\":false,\"supportInvoice\":false},\"ediFtp\":{\"ftpConnMode\":\"Active\",\"ftpFormat\":\"SFTP\",\"ftpMode\":\"ASCII\",\"username\":\"ftp-user\",\"password\":\"AC3-sentinel-12345\"},\"isDefaultConfig\":false,\"integrationType\":\"Claiming\",\"transmissionMethod\":\"File download\",\"fileFormat\":\"CSV\"}},\"schedulePeriod\":\"NONE\"}"; @Autowired private MockMvc mockMvc; + @Autowired + private KafkaService kafkaService; + @Autowired + private KafkaProperties kafkaProperties; @MockitoSpyBean private ExportConfigRepository repository; @MockitoSpyBean @@ -143,6 +168,15 @@ void postConfig() throws Exception { content().contentType("text/plain;charset=UTF-8")); verify(bursarExportScheduler).scheduleBursarJob(any(ExportConfig.class)); + + var event = pollConfigEventJson("0a3cba78-16e7-498e-b75b-98713000277b", DomainEventType.CREATE); + assertThat(event.getType()).isEqualTo(DomainEventType.CREATE); + assertThat(event.getTenant()).isEqualTo(TENANT); + assertNull(event.getOldValue()); + assertThat(event.getNewValue()) + .usingRecursiveComparison() + .ignoringFields("configName", "tenant") + .isEqualTo(OBJECT_MAPPER.readValue(NEW_CONFIG_REQUEST, ExportConfig.class)); } @Test @@ -171,10 +205,33 @@ void postClaimsConfig() throws Exception { content().contentType("text/plain;charset=UTF-8")); } + @Test + @DisplayName("Should redact all credential-shaped fields from the published config event (AC3)") + void postConfigShouldRedactCredentials() throws Exception { + mockMvc + .perform( + post("/data-export-spring/configs") + .contentType(MediaType.APPLICATION_JSON_VALUE) + .headers(defaultHeaders()) + .content(CLAIMS_REQUEST_WITH_PASSWORD)) + .andExpectAll(status().isCreated()); + + var eventJson = pollConfigEventRawJson(CLAIMS_WITH_PASSWORD_ID, DomainEventType.CREATE); + + assertThat(eventJson) + .doesNotContain(PASSWORD_SENTINEL) + .doesNotContain("\"password\"") + .doesNotContain("\"secret\"") + .doesNotContain("\"token\"") + .doesNotContain("\"apiKey\"") + .doesNotContain("\"credential\"") + .doesNotContain("\"privateKey\""); + } + @Test @DisplayName("Success update config") void putConfig() throws Exception { - saveConfig(UPDATE_CONFIG_REQUEST); + saveConfig(NEW_CONFIG_REQUEST); mockMvc .perform( @@ -185,6 +242,22 @@ void putConfig() throws Exception { .andExpectAll(status().isNoContent()); verify(bursarExportScheduler).scheduleBursarJob(any(ExportConfig.class)); + + var event = pollConfigEventJson("0a3cba78-16e7-498e-b75b-98713000277b", DomainEventType.UPDATE); + assertThat(event.getType()).isEqualTo(DomainEventType.UPDATE); + assertThat(event.getTenant()).isEqualTo(TENANT); + assertThat(event.getOldValue()) + .usingRecursiveComparison() + .ignoringFields("configName", "tenant") + .isEqualTo(OBJECT_MAPPER.readValue(NEW_CONFIG_REQUEST, ExportConfig.class)); + assertThat(event.getNewValue()) + .usingRecursiveComparison() + .ignoringFields("configName", "tenant") + .isEqualTo(OBJECT_MAPPER.readValue(UPDATE_CONFIG_REQUEST, ExportConfig.class)); + + assertThat(event.getOldValue().getSchedulePeriod()).isEqualTo(ExportConfig.SchedulePeriodEnum.DAY); + assertThat(event.getNewValue().getSchedulePeriod()).isEqualTo(ExportConfig.SchedulePeriodEnum.HOUR); + assertThat(event.getOldValue().getId()).isEqualTo(event.getNewValue().getId()); } @Test @@ -203,17 +276,22 @@ void putShouldThrowExceptionConfig() throws Exception { } @Test - @DisplayName("Fail update config") + @DisplayName("Fail update config and publish no event on validation failure (AC4)") void putConfigFail() throws Exception { - mockMvc - .perform( - put("/data-export-spring/configs/c8303ff3-7dec-49a1-acc8-7ce4f311fe21") - .contentType(MediaType.APPLICATION_JSON_VALUE) - .headers(defaultHeaders()) - .content(UPDATE_CONFIG_REQUEST_FAILED)) - .andExpectAll(status().isBadRequest(), - content().contentType(MediaType.APPLICATION_JSON_VALUE), - jsonPath("$.errors[0].message", startsWith("MethodArgumentNotValidException"))); + var topic = kafkaService.getTenantTopicName(Topic.CONFIG.getTopicName(), TENANT); + try (var consumer = TestKafkaConsumer.subscribe(topic, kafkaProperties)) { + mockMvc + .perform( + put("/data-export-spring/configs/" + FAILED_CONFIG_ID) + .contentType(MediaType.APPLICATION_JSON_VALUE) + .headers(defaultHeaders()) + .content(UPDATE_CONFIG_REQUEST_FAILED)) + .andExpectAll(status().isBadRequest(), + content().contentType(MediaType.APPLICATION_JSON_VALUE), + jsonPath("$.errors[0].message", startsWith("MethodArgumentNotValidException"))); + + assertThat(consumer.drainFor(FAILED_CONFIG_ID, Duration.ofSeconds(5))).isEmpty(); + } } @Test @@ -271,6 +349,33 @@ void shouldNotBeDeletedIfConfigIsNotExist() throws Exception { jsonPath("$.errors[0].message", startsWith("NotFoundException"))); } + private DomainEvent pollConfigEventJson(String configId, DomainEventType type) { + var topic = kafkaService.getTenantTopicName(Topic.CONFIG.getTopicName(), TENANT); + try (var consumer = TestKafkaConsumer.subscribe(topic, kafkaProperties)) { + return consumer.poll(configId).stream() + .map(event -> readEvent(event.value())) + .filter(event -> event.getType() == type) + .findFirst() + .orElseThrow(() -> new AssertionError("Expected " + type + " event for config " + configId)); + } + } + + private String pollConfigEventRawJson(String configId, DomainEventType type) { + var topic = kafkaService.getTenantTopicName(Topic.CONFIG.getTopicName(), TENANT); + try (var consumer = TestKafkaConsumer.subscribe(topic, kafkaProperties)) { + return consumer.poll(configId).stream() + .filter(event -> readEvent(event.value()).getType() == type) + .map(ConsumerRecord::value) + .findFirst() + .orElseThrow(() -> new AssertionError("Expected " + type + " event for config " + configId)); + } + } + + @SneakyThrows + private DomainEvent readEvent(String json) { + return OBJECT_MAPPER.readValue(json, new TypeReference<>() {}); + } + private void saveConfig(String config) throws Exception { mockMvc.perform(post("/data-export-spring/configs") .contentType(MediaType.APPLICATION_JSON_VALUE) diff --git a/src/test/java/org/folio/des/service/config/acquisition/ClaimsExportServiceTest.java b/src/test/java/org/folio/des/service/config/acquisition/ClaimsExportServiceTest.java index 0a50dce5..087351cf 100644 --- a/src/test/java/org/folio/des/service/config/acquisition/ClaimsExportServiceTest.java +++ b/src/test/java/org/folio/des/service/config/acquisition/ClaimsExportServiceTest.java @@ -3,6 +3,7 @@ import static org.folio.des.support.TestUtils.setInternalState; import static org.junit.jupiter.api.Assertions.assertEquals; import static org.mockito.ArgumentMatchers.any; +import static org.mockito.Mockito.mock; import static org.mockito.Mockito.verify; import static org.mockito.Mockito.when; @@ -19,13 +20,13 @@ import org.folio.des.mapper.ExportConfigMapperResolver; import org.folio.des.mapper.acquisition.ClaimsExportConfigMapperImpl; import org.folio.des.repository.ExportConfigRepository; +import org.folio.des.service.config.ExportConfigDomainEventService; import org.folio.des.validator.ExportConfigValidatorResolver; import org.folio.des.validator.acquisition.ClaimsExportParametersValidator; import org.junit.jupiter.api.BeforeEach; import org.junit.jupiter.api.DisplayName; import org.junit.jupiter.api.Test; import org.junit.jupiter.api.extension.ExtendWith; -import org.mockito.Mockito; import org.mockito.junit.jupiter.MockitoExtension; @ExtendWith(MockitoExtension.class) @@ -60,8 +61,9 @@ void setUp() { setInternalState(claimsExportConfigMapper, "objectMapper", new JacksonConfiguration().entityObjectMapper()); setInternalState(claimsExportConfigMapper, "validator", validator); - repository = Mockito.mock(ExportConfigRepository.class); - service = new ClaimsExportService(repository, claimsExportConfigMapper, exportConfigMapperResolver, exportConfigValidatorResolver); + repository = mock(ExportConfigRepository.class); + service = new ClaimsExportService(repository, claimsExportConfigMapper, exportConfigMapperResolver, exportConfigValidatorResolver, + mock(ExportConfigDomainEventService.class)); } @Test @@ -78,7 +80,8 @@ void testPostConfig() { @DisplayName("Update configuration") void testUpdateConfig() { when(repository.findById(UUID.fromString(CLAIMS_EXPORT_CONFIG.getId()))) - .thenReturn(java.util.Optional.of(new ExportConfigEntity())); + .thenReturn(java.util.Optional.of(new ExportConfigEntity().setType(ExportType.CLAIMS.getValue()))); + when(repository.save(any())).thenAnswer(i -> i.getArguments()[0]); service.updateConfig(CLAIMS_EXPORT_CONFIG.getId(), CLAIMS_EXPORT_CONFIG); diff --git a/src/test/java/org/folio/des/service/config/acquisition/EdifactOrdersExportServiceTest.java b/src/test/java/org/folio/des/service/config/acquisition/EdifactOrdersExportServiceTest.java index 41e3713e..f15d9fa0 100644 --- a/src/test/java/org/folio/des/service/config/acquisition/EdifactOrdersExportServiceTest.java +++ b/src/test/java/org/folio/des/service/config/acquisition/EdifactOrdersExportServiceTest.java @@ -3,6 +3,7 @@ import static org.folio.des.support.TestUtils.setInternalState; import static org.junit.jupiter.api.Assertions.assertEquals; import static org.mockito.ArgumentMatchers.any; +import static org.mockito.Mockito.mock; import static org.mockito.Mockito.times; import static org.mockito.Mockito.verify; import static org.mockito.Mockito.when; @@ -23,6 +24,7 @@ import org.folio.des.mapper.acquisition.EdifactExportConfigMapperImpl; import org.folio.des.repository.ExportConfigRepository; import org.folio.des.scheduling.ExportJobScheduler; +import org.folio.des.service.config.ExportConfigDomainEventService; import org.folio.des.validator.ExportConfigValidatorResolver; import org.folio.des.validator.acquisition.EdifactOrdersExportParametersValidator; import org.folio.des.validator.acquisition.EdifactOrdersScheduledParamsValidator; @@ -30,7 +32,6 @@ import org.junit.jupiter.api.DisplayName; import org.junit.jupiter.api.Test; import org.junit.jupiter.api.extension.ExtendWith; -import org.mockito.Mockito; import org.mockito.junit.jupiter.MockitoExtension; @ExtendWith(MockitoExtension.class) @@ -72,9 +73,10 @@ void setUp() { setInternalState(edifactExportConfigMapper, "objectMapper", new JacksonConfiguration().entityObjectMapper()); setInternalState(edifactExportConfigMapper, "validator", validator); - repository = Mockito.mock(ExportConfigRepository.class); - edifactOrdersExportJobScheduler = Mockito.mock(ExportJobScheduler.class); - service = new EdifactOrdersExportService(repository, edifactExportConfigMapper, exportConfigMapperResolver, exportConfigValidatorResolver, edifactOrdersExportJobScheduler); + repository = mock(ExportConfigRepository.class); + edifactOrdersExportJobScheduler = mock(ExportJobScheduler.class); + service = new EdifactOrdersExportService(repository, edifactExportConfigMapper, exportConfigMapperResolver, exportConfigValidatorResolver, + mock(ExportConfigDomainEventService.class), edifactOrdersExportJobScheduler); } @Test @@ -92,7 +94,8 @@ void testPostConfig() { @DisplayName("Update configuration") void testUpdateConfig() { when(repository.findById(UUID.fromString(EDIFACT_EXPORT_CONFIG.getId()))) - .thenReturn(java.util.Optional.of(new ExportConfigEntity())); + .thenReturn(java.util.Optional.of(new ExportConfigEntity().setType(ExportType.EDIFACT_ORDERS_EXPORT.getValue()))); + when(repository.save(any())).thenAnswer(i -> i.getArguments()[0]); service.updateConfig(EDIFACT_EXPORT_CONFIG.getId(), EDIFACT_EXPORT_CONFIG); diff --git a/src/test/java/org/folio/des/service/config/impl/BaseExportConfigServiceTest.java b/src/test/java/org/folio/des/service/config/impl/BaseExportConfigServiceTest.java index 6193cb35..00ead768 100644 --- a/src/test/java/org/folio/des/service/config/impl/BaseExportConfigServiceTest.java +++ b/src/test/java/org/folio/des/service/config/impl/BaseExportConfigServiceTest.java @@ -8,6 +8,7 @@ import static org.junit.jupiter.api.Assertions.assertThrows; import static org.junit.jupiter.api.Assertions.assertTrue; import static org.mockito.ArgumentMatchers.any; +import static org.mockito.Mockito.mock; import static org.mockito.Mockito.times; import static org.mockito.Mockito.verify; import static org.mockito.Mockito.when; @@ -24,6 +25,7 @@ import org.folio.des.mapper.DefaultExportConfigMapper; import org.folio.des.mapper.ExportConfigMapperResolver; import org.folio.des.repository.ExportConfigRepository; +import org.folio.des.service.config.ExportConfigDomainEventService; import org.folio.des.validator.BursarFeesFinesExportParametersValidator; import org.folio.des.validator.ExportConfigValidatorResolver; import org.junit.jupiter.api.Assertions; @@ -31,7 +33,6 @@ import org.junit.jupiter.api.DisplayName; import org.junit.jupiter.api.Test; import org.junit.jupiter.api.extension.ExtendWith; -import org.mockito.Mockito; import org.mockito.junit.jupiter.MockitoExtension; import org.springframework.data.domain.Page; import org.springframework.data.domain.PageImpl; @@ -57,8 +58,9 @@ void setUp() { var exportConfigMapperResolver = new ExportConfigMapperResolver(Map.of(), defaultExportConfigMapper); setInternalState(defaultExportConfigMapper, "objectMapper", new JacksonConfiguration().entityObjectMapper()); - repository = Mockito.mock(ExportConfigRepository.class); - service = new BaseExportConfigService(repository, defaultExportConfigMapper, exportConfigMapperResolver, exportConfigValidatorResolver); + repository = mock(ExportConfigRepository.class); + service = new BaseExportConfigService(repository, defaultExportConfigMapper, exportConfigMapperResolver, exportConfigValidatorResolver, + mock(ExportConfigDomainEventService.class)); } @Test @@ -98,7 +100,8 @@ void testUpdateConfig() { var exportConfig = getBursarExportConfig(); when(repository.findById(UUID.fromString(exportConfig.getId()))) - .thenReturn(java.util.Optional.of(new ExportConfigEntity())); + .thenReturn(java.util.Optional.of(new ExportConfigEntity().setType(ExportType.INVOICE_EXPORT.getValue()))); + when(repository.save(any())).thenAnswer(i -> i.getArguments()[0]); service.updateConfig(exportConfig.getId(), exportConfig); diff --git a/src/test/java/org/folio/des/service/config/impl/BursarFeesFinesExportConfigServiceTest.java b/src/test/java/org/folio/des/service/config/impl/BursarFeesFinesExportConfigServiceTest.java index 4ac3532d..435a241f 100644 --- a/src/test/java/org/folio/des/service/config/impl/BursarFeesFinesExportConfigServiceTest.java +++ b/src/test/java/org/folio/des/service/config/impl/BursarFeesFinesExportConfigServiceTest.java @@ -4,6 +4,7 @@ import static org.folio.des.support.TestUtils.setInternalState; import static org.junit.jupiter.api.Assertions.assertEquals; import static org.mockito.ArgumentMatchers.any; +import static org.mockito.Mockito.mock; import static org.mockito.Mockito.times; import static org.mockito.Mockito.verify; import static org.mockito.Mockito.when; @@ -19,13 +20,13 @@ import org.folio.des.mapper.ExportConfigMapperResolver; import org.folio.des.repository.ExportConfigRepository; import org.folio.des.scheduling.bursar.BursarExportScheduler; +import org.folio.des.service.config.ExportConfigDomainEventService; import org.folio.des.validator.BursarFeesFinesExportParametersValidator; import org.folio.des.validator.ExportConfigValidatorResolver; import org.junit.jupiter.api.BeforeEach; import org.junit.jupiter.api.DisplayName; import org.junit.jupiter.api.Test; import org.junit.jupiter.api.extension.ExtendWith; -import org.mockito.Mockito; import org.mockito.junit.jupiter.MockitoExtension; @ExtendWith(MockitoExtension.class) @@ -44,9 +45,10 @@ void setUp() { var exportConfigMapperResolver = new ExportConfigMapperResolver(Map.of(), defaultExportConfigMapper); setInternalState(defaultExportConfigMapper, "objectMapper", new JacksonConfiguration().entityObjectMapper()); - repository = Mockito.mock(ExportConfigRepository.class); - bursarExportScheduler = Mockito.mock(BursarExportScheduler.class); - service = new BursarFeesFinesExportConfigService(repository, defaultExportConfigMapper, exportConfigMapperResolver, exportConfigValidatorResolver, bursarExportScheduler); + repository = mock(ExportConfigRepository.class); + bursarExportScheduler = mock(BursarExportScheduler.class); + service = new BursarFeesFinesExportConfigService(repository, defaultExportConfigMapper, exportConfigMapperResolver, exportConfigValidatorResolver, + mock(ExportConfigDomainEventService.class), bursarExportScheduler); } @Test @@ -68,7 +70,8 @@ void testUpdateConfig() { var exportConfig = getBursarExportConfig(); when(repository.findById(UUID.fromString(exportConfig.getId()))) - .thenReturn(java.util.Optional.of(new ExportConfigEntity())); + .thenReturn(java.util.Optional.of(new ExportConfigEntity().setType(ExportType.BURSAR_FEES_FINES.getValue()))); + when(repository.save(any())).thenAnswer(i -> i.getArguments()[0]); service.updateConfig(exportConfig.getId(), exportConfig); diff --git a/src/test/java/org/folio/des/support/TestKafkaConsumer.java b/src/test/java/org/folio/des/support/TestKafkaConsumer.java new file mode 100644 index 00000000..00b1413c --- /dev/null +++ b/src/test/java/org/folio/des/support/TestKafkaConsumer.java @@ -0,0 +1,161 @@ +package org.folio.des.support; + +import static org.apache.kafka.clients.consumer.ConsumerConfig.AUTO_OFFSET_RESET_CONFIG; +import static org.apache.kafka.clients.consumer.ConsumerConfig.GROUP_ID_CONFIG; +import static org.apache.kafka.clients.consumer.ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG; +import static org.apache.kafka.clients.consumer.ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG; +import static org.awaitility.Awaitility.await; + +import java.io.Closeable; +import java.time.Duration; +import java.util.ArrayList; +import java.util.HashMap; +import java.util.List; +import java.util.Map; +import java.util.Objects; +import java.util.UUID; +import java.util.concurrent.BlockingQueue; +import java.util.concurrent.ExecutionException; +import java.util.concurrent.LinkedBlockingQueue; +import java.util.concurrent.TimeUnit; + +import org.apache.kafka.clients.admin.Admin; +import org.apache.kafka.clients.admin.NewTopic; +import org.apache.kafka.clients.consumer.ConsumerRecord; +import org.apache.kafka.common.errors.TopicExistsException; +import org.apache.kafka.common.serialization.StringDeserializer; +import org.springframework.boot.kafka.autoconfigure.KafkaProperties; +import org.springframework.kafka.core.DefaultKafkaConsumerFactory; +import org.springframework.kafka.listener.ContainerProperties; +import org.springframework.kafka.listener.KafkaMessageListenerContainer; +import org.springframework.kafka.listener.MessageListener; + +/** + * Reusable test consumer that subscribes to a single Kafka topic with a raw {@code String} value deserializer and + * buffers every received record. Integration tests use it to assert on the JSON payload of published domain events. + * + *

Adapted from the mod-notes {@code TestKafkaConsumer} pattern. Callers + * {@link #subscribe(String, KafkaProperties)} it, {@link #poll(String)} the buffered records filtered by key, and + * {@link #close()} it when done (it is {@link Closeable}, so it works with try-with-resources).

+ */ +public final class TestKafkaConsumer implements Closeable { + + private static final Duration DEFAULT_POLL_TIMEOUT = Duration.ofMinutes(1); + private static final Duration POLL_INTERVAL = Duration.ofSeconds(1); + + private final KafkaMessageListenerContainer container; + private final BlockingQueue> records = new LinkedBlockingQueue<>(); + private final List> buffer = new ArrayList<>(); + + private TestKafkaConsumer(KafkaMessageListenerContainer container) { + this.container = container; + } + + /** + * Creates and starts a consumer subscribed to the given topic. + * + * @param topic the topic to consume from (already env/tenant qualified) + * @param properties Spring Kafka properties (bootstrap servers point at the embedded broker) + * @return a started consumer; close it when done + */ + public static TestKafkaConsumer subscribe(String topic, KafkaProperties properties) { + createTopic(topic, properties); + // Override consumer settings on a local copy only — never mutate the shared Spring KafkaProperties bean, + // otherwise concurrently-running tests inherit this consumer's group id / offset reset. + Map config = new HashMap<>(properties.buildConsumerProperties()); + config.put(GROUP_ID_CONFIG, "mod-data-export-spring-test-group-" + UUID.randomUUID()); + config.put(AUTO_OFFSET_RESET_CONFIG, "earliest"); + config.put(KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class); + config.put(VALUE_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class); + + var consumerFactory = new DefaultKafkaConsumerFactory<>(config, new StringDeserializer(), new StringDeserializer()); + var containerProperties = new ContainerProperties(topic); + var container = new KafkaMessageListenerContainer<>(consumerFactory, containerProperties); + + var consumer = new TestKafkaConsumer(container); + container.setupMessageListener((MessageListener) consumer.records::add); + container.start(); + return consumer; + } + + /** + * Eagerly creates the topic via an admin client so the producer does not race the broker's lazy auto-creation + * (which otherwise surfaces as {@code Topic ... not present in metadata} on the first send). + * + * @param topic the topic to create (no-op if it already exists) + * @param properties Spring Kafka properties (used for the bootstrap servers) + */ + private static void createTopic(String topic, KafkaProperties properties) { + try (var admin = Admin.create(properties.buildAdminProperties())) { + admin.createTopics(List.of(new NewTopic(topic, 1, (short) 1))).all().get(); + } catch (ExecutionException e) { + if (!(e.getCause() instanceof TopicExistsException)) { + throw new IllegalStateException("Failed to create test topic " + topic, e); + } + } catch (InterruptedException e) { + Thread.currentThread().interrupt(); + throw new IllegalStateException("Interrupted while creating test topic " + topic, e); + } + } + + /** + * Waits (up to one minute) until at least one record with the given key has been received and returns all such + * records, failing the calling test if none arrive. + * + * @param key the record key to filter on (the config id) + * @return the matching records + */ + public List> poll(String key) { + var matched = new ArrayList>(); + await().pollInterval(POLL_INTERVAL).atMost(DEFAULT_POLL_TIMEOUT) + .untilAsserted(() -> { + records.drainTo(buffer); + var found = buffer.stream() + .filter(e -> Objects.equals(e.key(), key)) + .toList(); + if (found.isEmpty()) { + throw new AssertionError("No record received yet for key " + key); + } + matched.clear(); + matched.addAll(found); + }); + return matched; + } + + /** + * Drains records for up to the given duration and returns every buffered record matching the key (possibly + * empty). Unlike {@link #poll(String)} this never fails on an empty result, so it can assert that no + * event was published. It blocks on the record queue rather than sleeping a fixed interval. + * + * @param key the record key to filter on (the config id) + * @param duration the maximum time to wait for records to arrive + * @return the matching records seen within the window (empty if none arrived) + */ + public List> drainFor(String key, Duration duration) { + long deadline = System.nanoTime() + duration.toNanos(); + long remaining; + try { + while ((remaining = deadline - System.nanoTime()) > 0) { + var polled = records.poll(remaining, TimeUnit.NANOSECONDS); + if (polled != null) { + buffer.add(polled); + } + } + } catch (InterruptedException e) { + Thread.currentThread().interrupt(); + throw new IllegalStateException("Interrupted while draining test topic", e); + } + return buffer.stream() + .filter(e -> Objects.equals(e.key(), key)) + .toList(); + } + + @Override + public void close() { + container.stop(); + } +} + + + +