diff --git a/src/main/java/com/techevents/consumer/EventConsumer.java b/src/main/java/com/techevents/consumer/EventConsumer.java index 2954b12..113faa1 100644 --- a/src/main/java/com/techevents/consumer/EventConsumer.java +++ b/src/main/java/com/techevents/consumer/EventConsumer.java @@ -4,13 +4,12 @@ import com.fasterxml.jackson.databind.ObjectMapper; import com.fasterxml.jackson.datatype.jsr310.JavaTimeModule; import com.techevents.model.Event; -import com.techevents.repository.EventRepository; +import com.techevents.service.EventIngestionService; import com.techevents.config.KafkaTopics; import org.slf4j.Logger; import org.slf4j.LoggerFactory; -import org.springframework.dao.DataAccessException; import org.springframework.kafka.annotation.KafkaListener; import org.springframework.kafka.support.KafkaHeaders; import org.springframework.messaging.handler.annotation.Header; @@ -22,51 +21,100 @@ public class EventConsumer { private static final Logger log = LoggerFactory.getLogger(EventConsumer.class); - private final ObjectMapper objectMapper; - private final EventRepository eventRepository; + private final ObjectMapper mapper; + private final EventIngestionService eventIngestionService; - public EventConsumer(ObjectMapper objectMapper, EventRepository eventRepository) { - this.objectMapper = objectMapper.registerModule(new JavaTimeModule()); - this.eventRepository = eventRepository; + public EventConsumer(ObjectMapper mapper, + EventIngestionService ingestion) { + this.mapper = mapper.registerModule(new JavaTimeModule()); + this.eventIngestionService = ingestion; } - @KafkaListener( - topics = KafkaTopics.EVENTS_INGEST, - groupId = "event-group", - containerFactory = "kafkaListenerContainerFactory" - ) - public void listen( - @Header(KafkaHeaders.RECEIVED_TOPIC) String topic, - @Header(KafkaHeaders.RECEIVED_PARTITION) int partition, - @Header(KafkaHeaders.OFFSET) long offset, - @Payload String message - ) { - log.debug("Consuming message from {}-{}-{}", topic, partition, offset); - - // 1) DESERIALIZE - Event event; + @KafkaListener(topics = KafkaTopics.EVENTS_INGEST, groupId = "event-group") + public void listen(@Header(KafkaHeaders.RECEIVED_TOPIC) String topic, @Payload String payload) { try { - event = objectMapper.readValue(message, Event.class); - log.info("Parsed Event[id={}, title={}]", event.getId(), event.getTitle()); - } catch (JsonProcessingException e) { - log.error("JSON parse error at {}-{}-{}; payload={}", topic, partition, offset, message, e); - return; - } + Event e = mapper.readValue(payload, Event.class); - // 2) VALIDATE - if (event.getTitle() == null || event.getEventDate() == null) { - log.warn("Validation failed for {}-{}-{}; payload={}", topic, partition, offset, message); - return; - } + // Optional validation (can move to service) + if (e.getTitle() == null || e.getEventDate() == null) { + log.warn("Invalid payload received from topic {}: {}", topic, payload); + return; + } - // 3) PERSIST (let DB errors bubble up) - try { - eventRepository.save(event); - log.info("Saved Event[id={}, title={}]", event.getId(), event.getTitle()); - } catch (DataAccessException dbEx) { - log.error("DB error saving Event[id={}, title={}] at {}-{}-{}", - event.getId(), event.getTitle(), topic, partition, offset, dbEx); - throw dbEx; + eventIngestionService.ingest(e); + + } catch (JsonProcessingException ex) { + log.error("Failed to deserialize payload from topic {}: {}", topic, payload, ex); + } catch (Exception ex) { + log.error("Unexpected error while processing payload from topic {}: {}", topic, payload, ex); } } + } + +// @Component +// public class EventConsumer { + +// private static final Logger log = +// LoggerFactory.getLogger(EventConsumer.class); + +// private final ObjectMapper objectMapper; +// private final EventRepository eventRepository; + +// private final EventIngestionService eventIngestionService; + +// public EventConsumer(ObjectMapper objectMapper, EventRepository +// eventRepository, EventIngestionService eventIngestionService) { +// this.objectMapper = objectMapper.registerModule(new JavaTimeModule()); +// this.eventRepository = eventRepository; +// this.eventIngestionService = eventIngestionService; +// } + +// @KafkaListener( +// topics = KafkaTopics.EVENTS_INGEST, +// groupId = "event-group", +// containerFactory = "kafkaListenerContainerFactory" +// ) +// public void listen( +// @Header(KafkaHeaders.RECEIVED_TOPIC) String topic, +// @Header(KafkaHeaders.RECEIVED_PARTITION) int partition, +// @Header(KafkaHeaders.OFFSET) long offset, +// @Payload String message +// ) { +// log.debug("Consuming message from {}-{}-{}", topic, partition, offset); + +// // 1) DESERIALIZE +// Event event; +// try { +// event = objectMapper.readValue(message, Event.class); +// log.info("Parsed Event[id={}, title={}]", event.getId(), event.getTitle()); +// } catch (JsonProcessingException e) { +// log.error("JSON parse error at {}-{}-{}; payload={}", topic, partition, +// offset, message, e); +// return; +// } + +// // 2) VALIDATE +// if (event.getTitle() == null || event.getEventDate() == null) { +// log.warn("Validation failed for {}-{}-{}; payload={}", topic, partition, +// offset, message); +// return; +// } + +// // 3) PERSIST (let DB errors bubble up) +// try { +// eventRepository.save(event); +// log.info("Saved Event[id={}, title={}]", event.getId(), event.getTitle()); + +// eventIngestionService.processEvent(event); +// log.info("Processed Event[id={}, title={}]", event.getId(), +// event.getTitle()); + +// } catch (DataAccessException dbEx) { +// log.error("DB error saving Event[id={}, title={}] at {}-{}-{}", +// event.getId(), event.getTitle(), topic, partition, offset, dbEx); +// throw dbEx; +// } + +// } +// } diff --git a/src/main/java/com/techevents/model/NotificationMessage.java b/src/main/java/com/techevents/model/NotificationMessage.java new file mode 100644 index 0000000..f2badd6 --- /dev/null +++ b/src/main/java/com/techevents/model/NotificationMessage.java @@ -0,0 +1,35 @@ +package com.techevents.model; + +import com.fasterxml.jackson.annotation.JsonInclude; + +@JsonInclude(JsonInclude.Include.NON_NULL) +public class NotificationMessage { + private String email; + private String eventTitle; + private String eventDate; + private String city; + + public NotificationMessage(String email, Event event) { + this.email = email; + this.eventTitle = event.getTitle(); + this.eventDate = event.getEventDate().toString(); + this.city = event.getCity(); + } + + // Getters (and setters if needed) + public String getEmail() { + return email; + } + + public String getEventTitle() { + return eventTitle; + } + + public String getEventDate() { + return eventDate; + } + + public String getCity() { + return city; + } +} diff --git a/src/main/java/com/techevents/service/EventIngestionService.java b/src/main/java/com/techevents/service/EventIngestionService.java new file mode 100644 index 0000000..497fd21 --- /dev/null +++ b/src/main/java/com/techevents/service/EventIngestionService.java @@ -0,0 +1,45 @@ +package com.techevents.service; + +import com.techevents.model.Event; +import com.techevents.repository.EventRepository; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; +import org.springframework.amqp.rabbit.core.RabbitTemplate; +import org.springframework.beans.factory.annotation.Value; +import org.springframework.stereotype.Service; + +@Service +public class EventIngestionService { + + private static final Logger log = LoggerFactory.getLogger(EventIngestionService.class); + + private final EventRepository eventRepository; + private final RabbitTemplate rabbit; + + @Value("${app.rabbitmq.exchange}") + private String exchange; + + @Value("${app.rabbitmq.routing-key}") + private String routingKey; + + public EventIngestionService(EventRepository eventRepository, RabbitTemplate rabbitTemplate) { + this.eventRepository = eventRepository; + this.rabbit = rabbitTemplate; + } + + public void ingest(Event event) { + // Validate + if (event.getTitle() == null || event.getEventDate() == null) { + throw new IllegalArgumentException("Event must have a title and an event date."); + } + + // 1) Persist to database + Event saved = eventRepository.save(event); + log.info("Saved Event[id={}, title={}]", saved.getId(), saved.getTitle()); + + // 2) Publish to RabbitMQ + rabbit.convertAndSend(exchange, routingKey, saved); + log.info("Published Event[id={}] to RabbitMQ [exchange={}, routingKey={}]", + saved.getId(), exchange, routingKey); + } +} diff --git a/src/main/java/com/techevents/service/NotificationService.java b/src/main/java/com/techevents/service/NotificationService.java new file mode 100644 index 0000000..3e7a92a --- /dev/null +++ b/src/main/java/com/techevents/service/NotificationService.java @@ -0,0 +1,44 @@ +package com.techevents.service; + +import com.fasterxml.jackson.core.JsonProcessingException; +import com.fasterxml.jackson.databind.ObjectMapper; +import com.techevents.model.NotificationMessage; +import com.techevents.model.Subscriber; +import com.techevents.model.Event; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; +import org.springframework.amqp.rabbit.core.RabbitTemplate; +import org.springframework.beans.factory.annotation.Value; +import org.springframework.stereotype.Service; + +@Service +public class NotificationService { + + private static final Logger log = LoggerFactory.getLogger(NotificationService.class); + + private final RabbitTemplate rabbitTemplate; + private final ObjectMapper objectMapper; + + @Value("${app.rabbitmq.notification-exchange}") + private String exchange; + + @Value("${app.rabbitmq.notification-routing-key}") + private String routingKey; + + public NotificationService(RabbitTemplate rabbitTemplate, ObjectMapper objectMapper) { + this.rabbitTemplate = rabbitTemplate; + this.objectMapper = objectMapper; + } + + public void sendNotification(Subscriber subscriber, Event event) { + NotificationMessage message = new NotificationMessage(subscriber.getEmail(), event); + + try { + String jsonPayload = objectMapper.writeValueAsString(message); + rabbitTemplate.convertAndSend(exchange, routingKey, jsonPayload); + log.info("Published notification for email={} to RabbitMQ", subscriber.getEmail()); + } catch (JsonProcessingException e) { + log.error("Failed to serialize NotificationMessage for subscriber={}", subscriber.getEmail(), e); + } + } +} diff --git a/src/main/resources/application.yml b/src/main/resources/application.yml index 3e4bb00..3fbafb9 100644 --- a/src/main/resources/application.yml +++ b/src/main/resources/application.yml @@ -29,9 +29,12 @@ spring: app: rabbitmq: - queue: test.queue - exchange: test.exchange - routing-key: test.key + queue: events.saved.queue + exchange: events.saved.exchange + routing-key: events.saved.key + notification-exchange: notifications.exchange + notification-routing-key: notifications.event.created + server: port: 8080 diff --git a/src/test/java/com/techevents/consumer/EventConsumerUnitTest.java b/src/test/java/com/techevents/consumer/EventConsumerUnitTest.java index 12ac45b..60dae55 100644 --- a/src/test/java/com/techevents/consumer/EventConsumerUnitTest.java +++ b/src/test/java/com/techevents/consumer/EventConsumerUnitTest.java @@ -3,78 +3,73 @@ import com.fasterxml.jackson.databind.ObjectMapper; import com.fasterxml.jackson.datatype.jsr310.JavaTimeModule; import com.techevents.model.Event; -import com.techevents.repository.EventRepository; +import com.techevents.service.EventIngestionService; import org.junit.jupiter.api.BeforeEach; import org.junit.jupiter.api.Test; -import org.mockito.ArgumentCaptor; -import org.springframework.dao.DataIntegrityViolationException; -import static org.junit.jupiter.api.Assertions.*; import static org.mockito.Mockito.*; +import static org.junit.jupiter.api.Assertions.*; class EventConsumerUnitTest { - private EventRepository repo; private ObjectMapper mapper; + private EventIngestionService ingestionService; private EventConsumer consumer; @BeforeEach void setUp() { - repo = mock(EventRepository.class); mapper = new ObjectMapper().registerModule(new JavaTimeModule()); - consumer = new EventConsumer(mapper, repo); + ingestionService = mock(EventIngestionService.class); + consumer = new EventConsumer(mapper, ingestionService); } @Test - void malformedJson_shouldNotSave_andLogError() { - String badJson = "{ not: 'json' "; - assertDoesNotThrow(() -> consumer.listen("t", 0, 0L, badJson)); - verify(repo, never()).save(any()); + void malformedJson_shouldNotCallService() { + String badJson = "{ this is not valid JSON"; + consumer.listen("topic1", badJson); + verifyNoInteractions(ingestionService); } @Test - void missingFields_shouldNotSave_andLogWarn() { - String missing = """ + void missingRequiredFields_shouldNotCallService() { + String missingFieldsJson = """ { - "description":"no title/date", - "city":"Nowhere" + "description": "missing title and date", + "city": "X" } - """; - assertDoesNotThrow(() -> consumer.listen("t", 0, 0L, missing)); - verify(repo, never()).save(any()); + """; + consumer.listen("topic1", missingFieldsJson); + verifyNoInteractions(ingestionService); } @Test - void validJson_shouldSaveOnce() { - String valid = """ + void validJson_shouldCallServiceWithEvent() { + String validJson = """ { - "title":"UT Test", - "description":"desc", - "eventDate":"2025-05-18", - "city":"Here", - "tags":["a","b"] + "title": "Test Title", + "description": "desc", + "eventDate": "2025-06-01", + "city": "Test City", + "tags": ["tag1", "tag2"] } - """; - consumer.listen("t", 0, 0L, valid); - ArgumentCaptor cap = ArgumentCaptor.forClass(Event.class); - verify(repo, times(1)).save(cap.capture()); - assertEquals("UT Test", cap.getValue().getTitle()); + """; + consumer.listen("topic1", validJson); + verify(ingestionService, times(1)).ingest(any(Event.class)); } @Test - void dbError_shouldPropagate_andLogError() { - String valid = """ + void serviceThrowsException_shouldNotCrashConsumer() { + String validJson = """ { - "title":"UT Test", - "description":"desc", - "eventDate":"2025-05-18", - "city":"Here", - "tags":["a","b"] + "title": "Test Title", + "description": "desc", + "eventDate": "2025-06-01", + "city": "Test City", + "tags": ["tag1", "tag2"] } - """; - doThrow(new DataIntegrityViolationException("fk fail")) - .when(repo).save(any()); - assertThrows(DataIntegrityViolationException.class, - () -> consumer.listen("t", 0, 0L, valid)); + """; + doThrow(new RuntimeException("Simulated failure")) + .when(ingestionService).ingest(any()); + assertDoesNotThrow(() -> consumer.listen("topic1", validJson)); } } diff --git a/src/test/java/com/techevents/service/EventIngestionServiceTest.java b/src/test/java/com/techevents/service/EventIngestionServiceTest.java new file mode 100644 index 0000000..17aae93 --- /dev/null +++ b/src/test/java/com/techevents/service/EventIngestionServiceTest.java @@ -0,0 +1,70 @@ +package com.techevents.service; + +import com.techevents.model.Event; +import com.techevents.repository.EventRepository; +import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.Test; +import org.mockito.ArgumentCaptor; +import org.mockito.ArgumentMatchers; +import org.springframework.amqp.rabbit.core.RabbitTemplate; +import org.springframework.test.util.ReflectionTestUtils; + +import java.time.LocalDate; +import java.util.List; + +import static org.junit.jupiter.api.Assertions.*; +import static org.mockito.ArgumentMatchers.any; +import static org.mockito.ArgumentMatchers.anyString; +import static org.mockito.ArgumentMatchers.eq; +import static org.mockito.Mockito.*; + +class EventIngestionServiceTest { + + private EventRepository repo; + private RabbitTemplate rabbit; + private EventIngestionService service; + + @BeforeEach + void setUp() { + repo = mock(EventRepository.class); + rabbit = mock(RabbitTemplate.class); + service = new EventIngestionService(repo, rabbit); + + // Inject config values manually since Spring isn't doing it + ReflectionTestUtils.setField(service, "exchange", "test.exchange"); + ReflectionTestUtils.setField(service, "routingKey", "test.key"); + } + + @Test + void missingTitle_shouldThrow() { + Event e = new Event(null, "desc", LocalDate.now(), "city", List.of()); + Exception ex = assertThrows(IllegalArgumentException.class, () -> service.ingest(e)); + assertEquals("Event must have a title and an event date.", ex.getMessage()); + verifyNoInteractions(repo); + verifyNoInteractions(rabbit); + } + + @Test + void validEvent_shouldPersistAndPublish() { + Event e = new Event("Title", "desc", LocalDate.now(), "city", List.of()); + when(repo.save(any())).thenReturn(e); + service.ingest(e); + + verify(repo).save(any()); + verify(rabbit).convertAndSend(anyString(), anyString(), eq(e)); + } + + @Test + void dbError_shouldNotPublish() { + Event e = new Event("Title", "desc", LocalDate.now(), "city", List.of()); + when(repo.save(any())).thenThrow(new RuntimeException("DB failure")); + assertThrows(RuntimeException.class, () -> service.ingest(e)); + + verify(repo).save(any()); + verify(rabbit, never()) + .convertAndSend( + ArgumentMatchers.any(), + ArgumentMatchers.any(), + ArgumentMatchers.any()); + } +} diff --git a/src/test/java/com/techevents/service/NotificationServiceTest.java b/src/test/java/com/techevents/service/NotificationServiceTest.java new file mode 100644 index 0000000..bb998b1 --- /dev/null +++ b/src/test/java/com/techevents/service/NotificationServiceTest.java @@ -0,0 +1,55 @@ +package com.techevents.service; + +import com.fasterxml.jackson.databind.ObjectMapper; +import com.techevents.model.Event; +import com.techevents.model.Subscriber; +import com.techevents.model.NotificationMessage; +import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.Test; +import org.mockito.ArgumentCaptor; +import org.springframework.amqp.rabbit.core.RabbitTemplate; +import org.springframework.test.util.ReflectionTestUtils; + +import java.time.LocalDate; +import java.util.List; + +import static org.junit.jupiter.api.Assertions.assertTrue; +import static org.mockito.Mockito.*; + +class NotificationServiceTest { + + private RabbitTemplate rabbitTemplate; + private ObjectMapper objectMapper; + private NotificationService notificationService; + + @BeforeEach + void setUp() { + rabbitTemplate = mock(RabbitTemplate.class); + objectMapper = new ObjectMapper(); + notificationService = new NotificationService(rabbitTemplate, objectMapper); + + // Inject config values manually + ReflectionTestUtils.setField(notificationService, "exchange", "notifications.exchange"); + ReflectionTestUtils.setField(notificationService, "routingKey", "notifications.event.created"); + } + + @Test + void sendNotification_shouldPublishToRabbitMQ_withCorrectPayload() { + // Given + Event event = new Event("Spring Boot Conf", "Great talk", LocalDate.of(2025, 6, 1), "Boston", List.of("spring")); + Subscriber subscriber = new Subscriber("user@example.com"); + + // When + notificationService.sendNotification(subscriber, event); + + // Then + ArgumentCaptor payloadCaptor = ArgumentCaptor.forClass(String.class); + verify(rabbitTemplate).convertAndSend(eq("notifications.exchange"), eq("notifications.event.created"), payloadCaptor.capture()); + + String payload = payloadCaptor.getValue(); + assertTrue(payload.contains("user@example.com")); + assertTrue(payload.contains("Spring Boot Conf")); + assertTrue(payload.contains("2025-06-01")); + assertTrue(payload.contains("Boston")); + } +} diff --git a/src/test/resources/application.yml b/src/test/resources/application.yml index 1570b0c..60508be 100644 --- a/src/test/resources/application.yml +++ b/src/test/resources/application.yml @@ -19,3 +19,5 @@ app: queue: test.queue exchange: test.exchange routing-key: test.key + notification-exchange: notifications.exchange + notification-routing-key: notifications.event.created \ No newline at end of file