Skip to content
5 changes: 5 additions & 0 deletions NEWS.md
Original file line number Diff line number Diff line change
@@ -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)
Expand Down
5 changes: 5 additions & 0 deletions descriptors/ModuleDescriptor-template.json
Original file line number Diff line number Diff line change
Expand Up @@ -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"
Expand Down
21 changes: 15 additions & 6 deletions src/main/java/org/folio/des/config/ServiceConfiguration.java
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -82,33 +83,41 @@ 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
EdifactOrdersExportService edifactOrdersExportService(ExportConfigRepository repository,
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
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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 {
Expand Down
13 changes: 7 additions & 6 deletions src/main/java/org/folio/des/config/kafka/KafkaService.java
Original file line number Diff line number Diff line change
Expand Up @@ -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;
}

Expand Down Expand Up @@ -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);
}

Expand Down
67 changes: 67 additions & 0 deletions src/main/java/org/folio/des/domain/dto/event/DomainEvent.java
Original file line number Diff line number Diff line change
@@ -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<T> {

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 <T> DomainEvent<T> createEvent(T newValue, String tenant) {
return DomainEvent.<T>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 <T> DomainEvent<T> updateEvent(T oldValue, T newValue, String tenant) {
return DomainEvent.<T>builder()
.eventId(UUID.randomUUID())
.eventTs(System.currentTimeMillis())
.tenant(tenant)
.type(DomainEventType.UPDATE)
.oldValue(oldValue)
.newValue(newValue)
.build();
}
}

Original file line number Diff line number Diff line change
@@ -0,0 +1,7 @@
package org.folio.des.domain.dto.event;

public enum DomainEventType {
CREATE,
UPDATE
}

Original file line number Diff line number Diff line change
@@ -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.
*
* <p>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.</p>
*/
@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<ExportConfig> 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;
}
}
Original file line number Diff line number Diff line change
@@ -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.
*
* <p>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.</p>
*/
@Component
@Log4j2
@RequiredArgsConstructor
public class ExportConfigEventProducer {

private final KafkaTemplate<String, Object> 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<ExportConfig> 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);
}
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -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;

Expand All @@ -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
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;

Expand All @@ -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;
}

Expand Down
Loading
Loading