deleteEvent(UUID id, T oldEntity, String tenant) {
+ return new DomainEvent<>(id, oldEntity, null, DomainEventType.DELETE, tenant, currentTs());
+ }
+
+ private static String currentTs() {
+ return String.valueOf(System.currentTimeMillis());
+ }
+}
diff --git a/src/main/java/org/folio/notes/domain/event/DomainEventType.java b/src/main/java/org/folio/notes/domain/event/DomainEventType.java
new file mode 100644
index 00000000..cb56086d
--- /dev/null
+++ b/src/main/java/org/folio/notes/domain/event/DomainEventType.java
@@ -0,0 +1,5 @@
+package org.folio.notes.domain.event;
+
+public enum DomainEventType {
+ CREATE, UPDATE, DELETE
+}
diff --git a/src/main/java/org/folio/notes/integration/kafka/NoteEventProducer.java b/src/main/java/org/folio/notes/integration/kafka/NoteEventProducer.java
new file mode 100644
index 00000000..b1357223
--- /dev/null
+++ b/src/main/java/org/folio/notes/integration/kafka/NoteEventProducer.java
@@ -0,0 +1,75 @@
+package org.folio.notes.integration.kafka;
+
+import java.nio.charset.StandardCharsets;
+import java.util.ArrayList;
+import java.util.List;
+import java.util.UUID;
+import lombok.extern.slf4j.Slf4j;
+import org.apache.kafka.clients.producer.ProducerRecord;
+import org.apache.kafka.common.header.Header;
+import org.apache.kafka.common.header.internals.RecordHeader;
+import org.folio.notes.domain.dto.Note;
+import org.folio.notes.domain.event.DomainEvent;
+import org.folio.spring.FolioExecutionContext;
+import org.folio.spring.tools.kafka.FolioKafkaProperties;
+import org.folio.spring.tools.kafka.KafkaUtils;
+import org.springframework.kafka.core.KafkaTemplate;
+import org.springframework.stereotype.Component;
+
+/**
+ * Publishes Note {@link DomainEvent}s to the tenant-scoped Kafka topic.
+ *
+ * The topic name follows the FOLIO convention {@code {env}.{tenant}.notes.note} (for example
+ * {@code folio.diku.notes.note}), resolved via {@link KafkaUtils#getTenantTopicName(String, String)} the same way the
+ * topic is created on tenant enable. The tenant is resolved at publish time via {@link FolioExecutionContext}.
+ *
+ * Every message carries headers ({@code X-Okapi-Tenant}, {@code eventType}, {@code domain}) so consumers can
+ * filter by tenant, event type or Note domain without deserializing the JSON payload.
+ */
+@Slf4j
+@Component
+public class NoteEventProducer {
+
+ private static final String EVENT_TYPE_HEADER = "domain-event-type";
+
+ private final KafkaTemplate> kafkaTemplate;
+ private final FolioExecutionContext context;
+ private final String noteTopic;
+
+ public NoteEventProducer(KafkaTemplate> kafkaTemplate,
+ FolioExecutionContext context,
+ FolioKafkaProperties kafkaProperties) {
+ this.kafkaTemplate = kafkaTemplate;
+ this.context = context;
+ this.noteTopic = kafkaProperties.getTopics().getFirst().getName();
+ }
+
+ /**
+ * Publishes a Note domain event to the tenant-scoped topic.
+ *
+ * @param id the Note id used as the Kafka record key
+ * @param event the domain event envelope to publish
+ */
+ public void publish(UUID id, DomainEvent event) {
+ var topic = KafkaUtils.getTenantTopicName(noteTopic, context.getTenantId());
+ try {
+ log.debug("publish:: sending Note event [topic: {}, id: {}, type: {}]", topic, id, event.getType());
+ var producerRecord = new ProducerRecord<>(topic, null, id, event, buildHeaders(event));
+ kafkaTemplate.send(producerRecord);
+ log.info("publish:: Note event sent [topic: {}, id: {}, type: {}]", topic, id, event.getType());
+ } catch (Exception e) {
+ log.error("publish:: failed to send Note event [topic: {}, id: {}, type: {}]", topic, id, event.getType(), e);
+ }
+ }
+
+ private List buildHeaders(DomainEvent event) {
+ var headers = new ArrayList();
+ headers.add(header(EVENT_TYPE_HEADER, event.getType().name()));
+ context.getAllHeaders().forEach((key, value) -> headers.add(header(key, value.iterator().next())));
+ return headers;
+ }
+
+ private Header header(String key, String value) {
+ return new RecordHeader(key, value == null ? null : value.getBytes(StandardCharsets.UTF_8));
+ }
+}
diff --git a/src/main/java/org/folio/notes/service/DomainEventPublisherService.java b/src/main/java/org/folio/notes/service/DomainEventPublisherService.java
new file mode 100644
index 00000000..1ecc2ee7
--- /dev/null
+++ b/src/main/java/org/folio/notes/service/DomainEventPublisherService.java
@@ -0,0 +1,60 @@
+package org.folio.notes.service;
+
+import lombok.RequiredArgsConstructor;
+import lombok.extern.slf4j.Slf4j;
+import org.folio.notes.domain.dto.Note;
+import org.folio.notes.domain.event.DomainEvent;
+import org.folio.notes.integration.kafka.NoteEventProducer;
+import org.folio.spring.FolioExecutionContext;
+import org.springframework.stereotype.Service;
+
+/**
+ * Builds {@link DomainEvent} envelopes for Note changes and hands them to the {@link NoteEventProducer}.
+ *
+ * Callers use the event-specific methods ({@link #publishNoteCreatedEvent(Note)},
+ * {@link #publishNoteUpdatedEvent(Note, Note)}, {@link #publishNoteDeletedEvent(Note)}) and only need to provide the
+ * relevant Note snapshot(s); this service takes care of resolving the current tenant via {@link FolioExecutionContext},
+ * choosing the Kafka record key and populating the {@code old}/{@code new} envelope fields for each event type.
+ */
+@Slf4j
+@Service
+@RequiredArgsConstructor
+public class DomainEventPublisherService {
+
+ private final NoteEventProducer noteEventProducer;
+ private final FolioExecutionContext context;
+
+ /**
+ * Publishes a {@code CREATE} Note event carrying the new snapshot.
+ *
+ * @param newNote the newly-created Note snapshot
+ */
+ public void publishNoteCreatedEvent(Note newNote) {
+ publish(DomainEvent.createEvent(newNote.getId(), newNote, context.getTenantId()));
+ }
+
+ /**
+ * Publishes an {@code UPDATE} Note event carrying both the pre- and post-change snapshots.
+ *
+ * @param oldNote the pre-change Note snapshot
+ * @param newNote the post-change Note snapshot
+ */
+ public void publishNoteUpdatedEvent(Note oldNote, Note newNote) {
+ publish(DomainEvent.updateEvent(newNote.getId(), oldNote, newNote, context.getTenantId()));
+ }
+
+ /**
+ * Publishes a {@code DELETE} Note event carrying the pre-delete snapshot.
+ *
+ * @param oldNote the pre-delete Note snapshot
+ */
+ public void publishNoteDeletedEvent(Note oldNote) {
+ publish(DomainEvent.deleteEvent(oldNote.getId(), oldNote, context.getTenantId()));
+ }
+
+ private void publish(DomainEvent event) {
+ log.debug("publish:: publishing Note event [id: {}, type: {}, tenant: {}]",
+ event.getId(), event.getType(), event.getTenant());
+ noteEventProducer.publish(event.getId(), event);
+ }
+}
diff --git a/src/main/java/org/folio/notes/service/impl/NoteTenantService.java b/src/main/java/org/folio/notes/service/impl/NoteTenantService.java
index 4bf01817..7852a233 100644
--- a/src/main/java/org/folio/notes/service/impl/NoteTenantService.java
+++ b/src/main/java/org/folio/notes/service/impl/NoteTenantService.java
@@ -1,29 +1,42 @@
package org.folio.notes.service.impl;
+import lombok.extern.slf4j.Slf4j;
import org.folio.notes.service.NoteTypesService;
import org.folio.spring.FolioExecutionContext;
import org.folio.spring.liquibase.FolioSpringLiquibase;
import org.folio.spring.service.TenantService;
+import org.folio.spring.tools.kafka.KafkaAdminService;
+import org.folio.tenant.domain.dto.TenantAttributes;
import org.springframework.context.annotation.Primary;
import org.springframework.jdbc.core.JdbcTemplate;
import org.springframework.stereotype.Service;
@Service
@Primary
+@Slf4j
public class NoteTenantService extends TenantService {
private final NoteTypesService noteTypesService;
+ private final KafkaAdminService kafkaAdminService;
public NoteTenantService(JdbcTemplate jdbcTemplate,
FolioExecutionContext context,
FolioSpringLiquibase folioSpringLiquibase,
- NoteTypesService noteTypesService) {
+ NoteTypesService noteTypesService,
+ KafkaAdminService kafkaAdminService) {
super(jdbcTemplate, context, folioSpringLiquibase);
this.noteTypesService = noteTypesService;
+ this.kafkaAdminService = kafkaAdminService;
}
@Override
public void loadReferenceData() {
noteTypesService.populateDefaultType();
}
+
+ @Override
+ protected void afterTenantUpdate(TenantAttributes tenantAttributes) {
+ super.afterTenantUpdate(tenantAttributes);
+ kafkaAdminService.createTopics(context.getTenantId());
+ }
}
diff --git a/src/main/java/org/folio/notes/service/impl/NotesServiceImpl.java b/src/main/java/org/folio/notes/service/impl/NotesServiceImpl.java
index 8de228eb..00c91e44 100644
--- a/src/main/java/org/folio/notes/service/impl/NotesServiceImpl.java
+++ b/src/main/java/org/folio/notes/service/impl/NotesServiceImpl.java
@@ -33,6 +33,7 @@
import org.folio.notes.domain.repository.NoteRepository;
import org.folio.notes.domain.repository.NoteTypesRepository;
import org.folio.notes.exception.NoteNotFoundException;
+import org.folio.notes.service.DomainEventPublisherService;
import org.folio.notes.service.NotesService;
import org.folio.notes.util.HtmlSanitizer;
import org.folio.spring.data.OffsetRequest;
@@ -76,6 +77,7 @@ public class NotesServiceImpl implements NotesService {
private final NotesMapper notesMapper;
private final NoteCollectionMapper noteCollectionMapper;
private final HtmlSanitizer sanitizer;
+ private final DomainEventPublisherService domainEventPublisherService;
@Value("${folio.notes.response.limit}")
private Integer responseLimit;
@@ -131,7 +133,9 @@ public Note createNote(Note note) {
NoteEntity entity = saveNote(note, dto -> initNewEntity(notesMapper.toEntity(dto)));
log.info("createNote:: created note by title: {}, domain: {}, type: {}",
note.getTitle(), note.getDomain(), note.getType());
- return notesMapper.toDto(entity);
+ Note createdNote = notesMapper.toDto(entity);
+ domainEventPublisherService.publishNoteCreatedEvent(createdNote);
+ return createdNote;
}
@Transactional
diff --git a/src/main/resources/application.yaml b/src/main/resources/application.yaml
index 5f959253..a85d449a 100644
--- a/src/main/resources/application.yaml
+++ b/src/main/resources/application.yaml
@@ -1,5 +1,6 @@
# Module properties
folio:
+ environment: ${ENV:folio}
exchange:
enabled: true
tenant:
@@ -47,11 +48,33 @@ folio:
- target
response:
limit: ${MAX_RECORDS_COUNT:1000}
+ kafka:
+ topics:
+ - name: notes.note
+ numPartitions: ${KAFKA_NOTE_TOPIC_PARTITIONS:1}
+ replicationFactor: ${KAFKA_NOTE_TOPIC_REPLICATION_FACTOR:}
# Spring properties
spring:
application:
name: mod-notes
+ kafka:
+ bootstrap-servers: ${KAFKA_HOST:localhost}:${KAFKA_PORT:9092}
+ producer:
+ client-id: mod-notes
+ key-serializer: org.apache.kafka.common.serialization.UUIDSerializer
+ value-serializer: org.springframework.kafka.support.serializer.JacksonJsonSerializer
+ acks: all
+ retries: 3
+ properties:
+ max.block.ms: ${KAFKA_PRODUCER_MAX_BLOCK_MS:10000}
+ security:
+ protocol: ${KAFKA_SECURITY_PROTOCOL:PLAINTEXT}
+ ssl:
+ key-store-password: ${KAFKA_SSL_KEYSTORE_PASSWORD:}
+ key-store-location: ${KAFKA_SSL_KEYSTORE_LOCATION:}
+ trust-store-password: ${KAFKA_SSL_TRUSTSTORE_PASSWORD:}
+ trust-store-location: ${KAFKA_SSL_TRUSTSTORE_LOCATION:}
sql:
init:
continue-on-error: true
@@ -94,4 +117,4 @@ server:
logging:
level:
- org.folio.spring.filter.IncomingRequestLoggingFilter: DEBUG
\ No newline at end of file
+ org.folio.spring.filter.IncomingRequestLoggingFilter: DEBUG
diff --git a/src/test/java/org/folio/notes/controller/NotesControllerIT.java b/src/test/java/org/folio/notes/controller/NotesControllerIT.java
index dda567d6..6a16a313 100644
--- a/src/test/java/org/folio/notes/controller/NotesControllerIT.java
+++ b/src/test/java/org/folio/notes/controller/NotesControllerIT.java
@@ -14,6 +14,8 @@
import static org.hamcrest.Matchers.not;
import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertFalse;
+import static org.junit.jupiter.api.Assertions.assertNotNull;
+import static org.junit.jupiter.api.Assertions.assertNull;
import static org.junit.jupiter.api.Assertions.assertTrue;
import static org.junit.jupiter.params.provider.Arguments.arguments;
import static org.springframework.test.web.servlet.request.MockMvcRequestBuilders.delete;
@@ -25,6 +27,7 @@
import static org.springframework.test.web.servlet.result.MockMvcResultMatchers.status;
import jakarta.validation.ConstraintViolationException;
+import java.time.Duration;
import java.util.Arrays;
import java.util.Collections;
import java.util.List;
@@ -44,6 +47,7 @@
import org.folio.notes.domain.entity.NoteTypeEntity;
import org.folio.notes.exception.NoteNotFoundException;
import org.folio.notes.support.TestApiBase;
+import org.folio.notes.support.TestKafkaConsumer;
import org.folio.spring.cql.CqlQueryValidationException;
import org.hamcrest.Matcher;
import org.hamcrest.MatcherAssert;
@@ -55,13 +59,16 @@
import org.junit.jupiter.params.provider.CsvSource;
import org.junit.jupiter.params.provider.MethodSource;
import org.junit.jupiter.params.provider.ValueSource;
+import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.beans.factory.annotation.Value;
+import org.springframework.boot.kafka.autoconfigure.KafkaProperties;
import org.springframework.http.HttpHeaders;
import org.springframework.test.context.TestPropertySource;
import org.springframework.test.web.servlet.ResultMatcher;
import org.springframework.test.web.servlet.request.MockHttpServletRequestBuilder;
import org.springframework.web.bind.MethodArgumentNotValidException;
import org.springframework.web.method.annotation.MethodArgumentTypeMismatchException;
+import tools.jackson.databind.JsonNode;
@TestPropertySource(properties = {"folio.notes.types.defaults.limit=5"})
class NotesControllerIT extends TestApiBase {
@@ -88,8 +95,15 @@ class NotesControllerIT extends TestApiBase {
private static final UUID[] TYPE_IDS = new UUID[] {randomUUID(), randomUUID()};
private static final int DEFAULT_LINK_AMOUNT = 1;
+ private static final String NOTE_EVENT_TOPIC = "folio.test.notes.note";
+ private static final String ORDER_LINE_TYPE = "order-line";
+ private static final UUID EVENT_NOTE_TYPE_ID = UUID.fromString("2af21797-d25b-46dc-8427-1759d1db2057");
+ private static final String EVENT_NOTE_TYPE_NAME = "General note";
+
@Value("${folio.notes.types.defaults.limit}")
private String defaultNoteTypeLimit;
+ @Autowired
+ private KafkaProperties kafkaProperties;
@BeforeEach
void setUp() {
@@ -1065,6 +1079,80 @@ private NoteEntity generateNoteEntityWithParams(String title, String typeId, Str
return noteEntity;
}
+ // Tests for domain events
+ @Test
+ @DisplayName("create publishes a CREATE event with the full snapshot and no 'old' field")
+ void createNote_publishesCreateEvent() throws Exception {
+ saveEventNoteType();
+ var link = new Link().id(UUID.randomUUID().toString()).type(ORDER_LINE_TYPE);
+ var note = new Note()
+ .title("Kafka title")
+ .content("Kafka details")
+ .domain("orders")
+ .typeId(EVENT_NOTE_TYPE_ID)
+ .links(List.of(link));
+
+ try (var consumer = TestKafkaConsumer.subscribe(NOTE_EVENT_TOPIC, kafkaProperties)) {
+ var response = mockMvc.perform(postNote(note))
+ .andExpect(status().isCreated())
+ .andReturn().getResponse().getContentAsString();
+
+ var event = consumer.poll();
+ var envelope = OBJECT_MAPPER.readTree(event.value());
+
+ assertEventEnvelope(envelope, event.value());
+ assertEventPayload(envelope.get("new"), OBJECT_MAPPER.readTree(response).get("id").asString(), link);
+ }
+ }
+
+ @Test
+ @DisplayName("invalid payload (missing title) publishes no event")
+ void createNote_invalidPayload_publishesNoEvent() throws Exception {
+ saveEventNoteType();
+ var note = new Note()
+ .domain("orders")
+ .typeId(EVENT_NOTE_TYPE_ID)
+ .links(List.of(new Link().id(UUID.randomUUID().toString()).type(ORDER_LINE_TYPE)));
+
+ try (var consumer = TestKafkaConsumer.subscribe(NOTE_EVENT_TOPIC, kafkaProperties)) {
+ mockMvc.perform(postNote(note))
+ .andExpect(status().is4xxClientError());
+
+ var event = consumer.pollNullable(Duration.ofSeconds(10));
+ assertNull(event, "No event should be published for an invalid create");
+ }
+ }
+
+ private void saveEventNoteType() {
+ var noteType = new NoteTypeEntity();
+ noteType.setId(EVENT_NOTE_TYPE_ID);
+ noteType.setName(EVENT_NOTE_TYPE_NAME);
+ noteType.setCreatedBy(USER_ID);
+ databaseHelper.saveNoteType(noteType, TENANT);
+ }
+
+ private void assertEventEnvelope(JsonNode envelope, String rawJson) {
+ assertNotNull(envelope.get("id"), "eventId (id) must be present");
+ assertEquals("CREATE", envelope.get("type").asString());
+ assertEquals(TENANT, envelope.get("tenant").asString());
+ assertNull(envelope.get("old"), "CREATE event must not carry an 'old' field");
+ assertFalse(rawJson.contains("\"old\""), "JSON must not contain an 'old' field");
+ }
+
+ private void assertEventPayload(JsonNode newNode, String createdId, Link link) {
+ assertNotNull(newNode);
+ assertEquals(createdId, newNode.get("id").asString());
+ assertEquals("Kafka title", newNode.get("title").asString());
+ assertEquals("Kafka details", newNode.get("content").asString());
+ assertEquals(EVENT_NOTE_TYPE_ID.toString(), newNode.get("typeId").asString());
+ assertEquals(USER_ID.toString(), newNode.get("metadata").get("createdByUserId").asString());
+
+ var links = newNode.get("links");
+ assertEquals(1, links.size());
+ assertEquals(ORDER_LINE_TYPE, links.get(0).get("type").asString());
+ assertEquals(link.getId(), links.get(0).get("id").asString());
+ }
+
private Note generateNote() throws Exception {
var noteType = new NoteType().name(insecure().nextAlphabetic(100));
var notyTypeAsString = mockMvc.perform(postNoteType(noteType)).andExpect(status().isCreated())
diff --git a/src/test/java/org/folio/notes/domain/event/DomainEventTest.java b/src/test/java/org/folio/notes/domain/event/DomainEventTest.java
new file mode 100644
index 00000000..2a38acbf
--- /dev/null
+++ b/src/test/java/org/folio/notes/domain/event/DomainEventTest.java
@@ -0,0 +1,67 @@
+package org.folio.notes.domain.event;
+
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertFalse;
+import static org.junit.jupiter.api.Assertions.assertNotNull;
+import static org.junit.jupiter.api.Assertions.assertNull;
+import static org.junit.jupiter.api.Assertions.assertTrue;
+
+import java.util.UUID;
+import org.folio.notes.domain.dto.Note;
+import org.folio.spring.testing.type.UnitTest;
+import org.junit.jupiter.api.Test;
+import tools.jackson.databind.json.JsonMapper;
+
+@UnitTest
+class DomainEventTest {
+
+ private static final JsonMapper MAPPER = JsonMapper.builder().build();
+
+ @Test
+ void createEvent_setsCreateTypeAndOmitsOld() {
+ var id = UUID.randomUUID();
+ var note = new Note().title("t").domain("orders").typeId(UUID.randomUUID());
+
+ var event = DomainEvent.createEvent(id, note, "diku");
+
+ assertEquals(id, event.getId());
+ assertEquals(DomainEventType.CREATE, event.getType());
+ assertEquals("diku", event.getTenant());
+ assertNull(event.getOldEntity());
+ assertEquals(note, event.getNewEntity());
+ assertNotNull(event.getTs());
+
+ var json = MAPPER.writeValueAsString(event);
+ assertTrue(json.contains("\"type\":\"CREATE\""));
+ assertTrue(json.contains("\"new\""));
+ assertFalse(json.contains("\"old\""), "NON_NULL must omit the null 'old' field");
+ }
+
+ @Test
+ void updateEvent_carriesOldAndNew() {
+ var id = UUID.randomUUID();
+ var oldNote = new Note().title("old");
+ var newNote = new Note().title("new");
+
+ var event = DomainEvent.updateEvent(id, oldNote, newNote, "diku");
+
+ assertEquals(DomainEventType.UPDATE, event.getType());
+ assertEquals(oldNote, event.getOldEntity());
+ assertEquals(newNote, event.getNewEntity());
+ }
+
+ @Test
+ void deleteEvent_carriesOldAndOmitsNew() {
+ var id = UUID.randomUUID();
+ var oldNote = new Note().title("old");
+
+ var event = DomainEvent.deleteEvent(id, oldNote, "diku");
+
+ assertEquals(DomainEventType.DELETE, event.getType());
+ assertEquals(oldNote, event.getOldEntity());
+ assertNull(event.getNewEntity());
+ }
+}
+
+
+
diff --git a/src/test/java/org/folio/notes/integration/kafka/NoteEventProducerTest.java b/src/test/java/org/folio/notes/integration/kafka/NoteEventProducerTest.java
new file mode 100644
index 00000000..0536c9a3
--- /dev/null
+++ b/src/test/java/org/folio/notes/integration/kafka/NoteEventProducerTest.java
@@ -0,0 +1,86 @@
+package org.folio.notes.integration.kafka;
+
+import static org.assertj.core.api.Assertions.assertThat;
+import static org.assertj.core.api.Assertions.assertThatCode;
+import static org.mockito.ArgumentMatchers.any;
+import static org.mockito.Mockito.verify;
+import static org.mockito.Mockito.when;
+
+import java.nio.charset.StandardCharsets;
+import java.util.List;
+import java.util.UUID;
+import org.apache.kafka.clients.producer.ProducerRecord;
+import org.folio.notes.domain.dto.Note;
+import org.folio.notes.domain.event.DomainEvent;
+import org.folio.spring.FolioExecutionContext;
+import org.folio.spring.integration.XOkapiHeaders;
+import org.folio.spring.testing.type.UnitTest;
+import org.folio.spring.tools.kafka.FolioKafkaProperties;
+import org.folio.spring.tools.kafka.KafkaUtils;
+import org.junit.jupiter.api.BeforeEach;
+import org.junit.jupiter.api.Test;
+import org.junit.jupiter.api.extension.ExtendWith;
+import org.mockito.ArgumentCaptor;
+import org.mockito.Mock;
+import org.mockito.junit.jupiter.MockitoExtension;
+import org.springframework.kafka.core.KafkaTemplate;
+
+@UnitTest
+@ExtendWith(MockitoExtension.class)
+class NoteEventProducerTest {
+
+ private static final String TENANT = "diku";
+ private static final String TOPIC_NAME = "notes.note";
+ private static final String TENANT_TOPIC = KafkaUtils.getTenantTopicName(TOPIC_NAME, TENANT);
+
+ @Mock
+ private KafkaTemplate> kafkaTemplate;
+ @Mock
+ private FolioExecutionContext context;
+ @Mock
+ private FolioKafkaProperties kafkaProperties;
+
+ private NoteEventProducer producer;
+
+ @BeforeEach
+ void setUp() {
+ when(kafkaProperties.getTopics()).thenReturn(List.of(FolioKafkaProperties.KafkaTopic.of(TOPIC_NAME, 1, null)));
+ producer = new NoteEventProducer(kafkaTemplate, context, kafkaProperties);
+ }
+
+ @Test
+ void publish_sendsRecordWithTenantScopedTopicKeyAndHeaders() {
+ var id = UUID.randomUUID();
+ var event = DomainEvent.createEvent(id, new Note().id(id).domain("orders"), TENANT);
+ when(context.getTenantId()).thenReturn(TENANT);
+
+ producer.publish(id, event);
+
+ @SuppressWarnings("unchecked")
+ ArgumentCaptor>> captor = ArgumentCaptor.forClass(ProducerRecord.class);
+ verify(kafkaTemplate).send(captor.capture());
+ var kafkaRecord = captor.getValue();
+
+ assertThat(kafkaRecord.topic()).isEqualTo(TENANT_TOPIC);
+ assertThat(kafkaRecord.key()).isEqualTo(id);
+ assertThat(headerValue(kafkaRecord, XOkapiHeaders.TENANT)).isEqualTo(TENANT);
+ assertThat(headerValue(kafkaRecord, "eventType")).isEqualTo("CREATE");
+ assertThat(headerValue(kafkaRecord, "domain")).isEqualTo("orders");
+ }
+
+ @Test
+ @SuppressWarnings("unchecked")
+ void publish_swallowsBrokerFailure_doesNotThrowToCaller() {
+ var id = UUID.randomUUID();
+ var event = DomainEvent.updateEvent(id, new Note().id(id), new Note().id(id), TENANT);
+ when(context.getTenantId()).thenReturn(TENANT);
+ when(kafkaTemplate.send(any(ProducerRecord.class))).thenThrow(new RuntimeException("broker down"));
+
+ assertThatCode(() -> producer.publish(id, event)).doesNotThrowAnyException();
+ }
+
+ private String headerValue(ProducerRecord> kafkaRecord, String key) {
+ var header = kafkaRecord.headers().lastHeader(key);
+ return header == null ? null : new String(header.value(), StandardCharsets.UTF_8);
+ }
+}
diff --git a/src/test/java/org/folio/notes/service/DomainEventPublisherServiceTest.java b/src/test/java/org/folio/notes/service/DomainEventPublisherServiceTest.java
new file mode 100644
index 00000000..d713f9a2
--- /dev/null
+++ b/src/test/java/org/folio/notes/service/DomainEventPublisherServiceTest.java
@@ -0,0 +1,87 @@
+package org.folio.notes.service;
+
+import static org.mockito.ArgumentMatchers.eq;
+import static org.mockito.Mockito.verify;
+import static org.mockito.Mockito.when;
+
+import java.util.UUID;
+import org.folio.notes.domain.dto.Note;
+import org.folio.notes.domain.event.DomainEvent;
+import org.folio.notes.domain.event.DomainEventType;
+import org.folio.notes.integration.kafka.NoteEventProducer;
+import org.folio.spring.FolioExecutionContext;
+import org.folio.spring.testing.type.UnitTest;
+import org.junit.jupiter.api.Test;
+import org.junit.jupiter.api.extension.ExtendWith;
+import org.mockito.ArgumentCaptor;
+import org.mockito.InjectMocks;
+import org.mockito.Mock;
+import org.mockito.junit.jupiter.MockitoExtension;
+
+@UnitTest
+@ExtendWith(MockitoExtension.class)
+class DomainEventPublisherServiceTest {
+
+ @Mock
+ private NoteEventProducer noteEventProducer;
+ @Mock
+ private FolioExecutionContext context;
+ @InjectMocks
+ private DomainEventPublisherService service;
+
+ @Test
+ void publishNoteCreatedEvent_buildsCreateEnvelopeWithoutOld() {
+ var id = UUID.randomUUID();
+ var note = new Note().id(id).title("t");
+ when(context.getTenantId()).thenReturn("diku");
+
+ service.publishNoteCreatedEvent(note);
+
+ var captor = eventCaptor();
+ verify(noteEventProducer).publish(eq(id), captor.capture());
+ var event = captor.getValue();
+ org.junit.jupiter.api.Assertions.assertEquals(DomainEventType.CREATE, event.getType());
+ org.junit.jupiter.api.Assertions.assertEquals("diku", event.getTenant());
+ org.junit.jupiter.api.Assertions.assertEquals(note, event.getNewEntity());
+ org.junit.jupiter.api.Assertions.assertNull(event.getOldEntity());
+ }
+
+ @Test
+ void publishNoteUpdatedEvent_carriesOldAndNew() {
+ var id = UUID.randomUUID();
+ var oldNote = new Note().id(id).title("old");
+ var newNote = new Note().id(id).title("new");
+ when(context.getTenantId()).thenReturn("diku");
+
+ service.publishNoteUpdatedEvent(oldNote, newNote);
+
+ var captor = eventCaptor();
+ verify(noteEventProducer).publish(eq(id), captor.capture());
+ var event = captor.getValue();
+ org.junit.jupiter.api.Assertions.assertEquals(DomainEventType.UPDATE, event.getType());
+ org.junit.jupiter.api.Assertions.assertEquals(oldNote, event.getOldEntity());
+ org.junit.jupiter.api.Assertions.assertEquals(newNote, event.getNewEntity());
+ }
+
+ @Test
+ void publishNoteDeletedEvent_carriesOldAndOmitsNew() {
+ var id = UUID.randomUUID();
+ var oldNote = new Note().id(id).title("old");
+ when(context.getTenantId()).thenReturn("diku");
+
+ service.publishNoteDeletedEvent(oldNote);
+
+ var captor = eventCaptor();
+ verify(noteEventProducer).publish(eq(id), captor.capture());
+ var event = captor.getValue();
+ org.junit.jupiter.api.Assertions.assertEquals(DomainEventType.DELETE, event.getType());
+ org.junit.jupiter.api.Assertions.assertEquals(oldNote, event.getOldEntity());
+ org.junit.jupiter.api.Assertions.assertNull(event.getNewEntity());
+ }
+
+ @SuppressWarnings("unchecked")
+ private ArgumentCaptor> eventCaptor() {
+ return ArgumentCaptor.forClass(DomainEvent.class);
+ }
+}
+
diff --git a/src/test/java/org/folio/notes/service/NoteTenantServiceTest.java b/src/test/java/org/folio/notes/service/NoteTenantServiceTest.java
index eb6c3816..4e9bab1c 100644
--- a/src/test/java/org/folio/notes/service/NoteTenantServiceTest.java
+++ b/src/test/java/org/folio/notes/service/NoteTenantServiceTest.java
@@ -3,23 +3,40 @@
import static org.mockito.Mockito.verify;
import org.folio.notes.service.impl.NoteTenantService;
+import org.folio.spring.FolioExecutionContext;
+import org.folio.spring.liquibase.FolioSpringLiquibase;
import org.folio.spring.testing.type.UnitTest;
+import org.folio.spring.tools.kafka.KafkaAdminService;
+import org.junit.jupiter.api.BeforeEach;
import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.extension.ExtendWith;
-import org.mockito.InjectMocks;
import org.mockito.Mock;
import org.mockito.junit.jupiter.MockitoExtension;
+import org.springframework.jdbc.core.JdbcTemplate;
@UnitTest
@ExtendWith(MockitoExtension.class)
class NoteTenantServiceTest {
+ @Mock
+ private JdbcTemplate jdbcTemplate;
+ @Mock
+ private FolioExecutionContext context;
+ @Mock
+ private FolioSpringLiquibase folioSpringLiquibase;
@Mock
private NoteTypesService noteTypesService;
+ @Mock
+ private KafkaAdminService kafkaAdminService;
- @InjectMocks
private NoteTenantService tenantService;
+ @BeforeEach
+ void setUp() {
+ tenantService = new NoteTenantService(jdbcTemplate, context, folioSpringLiquibase, noteTypesService,
+ kafkaAdminService);
+ }
+
@Test
void shouldPopulateDefaultTypeOnLoadReferenceData() {
tenantService.loadReferenceData();
diff --git a/src/test/java/org/folio/notes/support/TestApiBase.java b/src/test/java/org/folio/notes/support/TestApiBase.java
index aee8a5fe..022ce2d6 100644
--- a/src/test/java/org/folio/notes/support/TestApiBase.java
+++ b/src/test/java/org/folio/notes/support/TestApiBase.java
@@ -16,6 +16,7 @@
import org.folio.notes.domain.dto.User;
import org.folio.spring.FolioModuleMetadata;
import org.folio.spring.integration.XOkapiHeaders;
+import org.folio.spring.testing.extension.EnableKafka;
import org.folio.spring.testing.extension.EnableOkapi;
import org.folio.spring.testing.extension.EnablePostgres;
import org.folio.spring.testing.extension.impl.OkapiConfiguration;
@@ -41,6 +42,7 @@
@EnableOkapi
@EnablePostgres
+@EnableKafka
@AutoConfigureMockMvc
@SpringBootTest
@ContextConfiguration
diff --git a/src/test/java/org/folio/notes/support/TestKafkaConsumer.java b/src/test/java/org/folio/notes/support/TestKafkaConsumer.java
new file mode 100644
index 00000000..fb7cd226
--- /dev/null
+++ b/src/test/java/org/folio/notes/support/TestKafkaConsumer.java
@@ -0,0 +1,143 @@
+package org.folio.notes.support;
+
+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 static org.junit.jupiter.api.Assertions.assertNotNull;
+
+import java.io.Closeable;
+import java.time.Duration;
+import java.util.HashMap;
+import java.util.List;
+import java.util.Map;
+import java.util.concurrent.BlockingQueue;
+import java.util.concurrent.ExecutionException;
+import java.util.concurrent.LinkedBlockingQueue;
+import java.util.concurrent.TimeUnit;
+import java.util.concurrent.atomic.AtomicReference;
+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.
+ *
+ * The consumer owns its listener container and the record buffer; callers only need to
+ * {@link #subscribe(String, KafkaProperties)} it, {@link #poll()} the buffered records, and {@link #close()} it when
+ * done (it is {@link Closeable}, so it also works with try-with-resources or an {@code @AfterEach} hook).
+ */
+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 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);
+ properties.getConsumer().setGroupId("mod-notes-test-group");
+ properties.getConsumer().setAutoOffsetReset("earliest");
+ Map config = new HashMap<>(properties.buildConsumerProperties());
+ 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) for the next record and returns it.
+ *
+ * @return the received record
+ */
+ public ConsumerRecord poll() {
+ return poll(DEFAULT_POLL_TIMEOUT);
+ }
+
+ /**
+ * Waits up to {@code timeout} for the next record and returns it, failing the calling test if none arrives.
+ *
+ * @param timeout the maximum time to wait
+ * @return the received record
+ */
+ public ConsumerRecord poll(Duration timeout) {
+ var holder = new AtomicReference>();
+ await().pollInterval(POLL_INTERVAL).atMost(timeout)
+ .untilAsserted(() -> {
+ var record = records.poll();
+ assertNotNull(record, "Expected a record on the topic");
+ holder.set(record);
+ });
+ return holder.get();
+ }
+
+ /**
+ * Returns the next record if one arrives within {@code timeout}, or {@code null} otherwise (no assertion). Useful for
+ * verifying that no event was published.
+ *
+ * @param timeout the maximum time to wait
+ * @return the received record, or {@code null} if none arrived
+ */
+ public ConsumerRecord pollNullable(Duration timeout) {
+ try {
+ return records.poll(timeout.toMillis(), TimeUnit.MILLISECONDS);
+ } catch (InterruptedException e) {
+ Thread.currentThread().interrupt();
+ throw new IllegalStateException("Interrupted while polling for records", e);
+ }
+ }
+
+ @Override
+ public void close() {
+ container.stop();
+ }
+}
+
+
+