Skip to content

Commit f3be20c

Browse files
committed
refactor(dispatcher): split DispatcherHub in multiple sub classes with dedicated perimeter
1 parent 36986eb commit f3be20c

10 files changed

Lines changed: 1678 additions & 1225 deletions
Lines changed: 293 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,293 @@
1+
/**
2+
* Copyright © 2023-2026 Agence du Numerique en Sante (ANS)
3+
*
4+
* Licensed under the Apache License, Version 2.0 (the "License");
5+
* you may not use this file except in compliance with the License.
6+
* You may obtain a copy of the License at
7+
*
8+
* http://www.apache.org/licenses/LICENSE-2.0
9+
*
10+
* Unless required by applicable law or agreed to in writing, software
11+
* distributed under the License is distributed on an "AS IS" BASIS,
12+
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
13+
* See the License for the specific language governing permissions and
14+
* limitations under the License.
15+
*/
16+
package com.hubsante.hub.service;
17+
18+
import static com.hubsante.hub.config.AmqpConfiguration.*;
19+
import static com.hubsante.hub.service.ConversionStubs.echoConversionService;
20+
import static com.hubsante.hub.service.ConversionStubs.failConversionService;
21+
import static com.hubsante.hub.testsupport.HubTestConstants.*;
22+
import static com.hubsante.hub.testsupport.HubTestScaffolding.aHub;
23+
import static com.hubsante.hub.testsupport.MessageTestUtils.*;
24+
import static com.hubsante.hub.testsupport.assertions.HubAssertions.assertThatMessageSentTo;
25+
import static org.junit.jupiter.api.Assertions.*;
26+
import static org.mockito.ArgumentMatchers.any;
27+
import static org.mockito.ArgumentMatchers.anyString;
28+
import static org.mockito.Mockito.*;
29+
30+
import com.fasterxml.jackson.databind.ObjectMapper;
31+
import com.fasterxml.jackson.dataformat.xml.XmlMapper;
32+
import com.hubsante.hub.config.HubConfiguration;
33+
import com.hubsante.hub.exception.ConversionException;
34+
import com.hubsante.hub.exception.SchemaValidationException;
35+
import com.hubsante.hub.testsupport.HubTestScaffolding;
36+
import com.hubsante.model.EdxlHandler;
37+
import com.hubsante.model.Validator;
38+
import com.hubsante.model.edxl.EdxlMessage;
39+
import com.hubsante.model.exception.ValidationException;
40+
import io.micrometer.core.instrument.MeterRegistry;
41+
import io.micrometer.tracing.Tracer;
42+
import java.io.IOException;
43+
import java.nio.charset.StandardCharsets;
44+
import org.junit.jupiter.api.BeforeEach;
45+
import org.junit.jupiter.api.DisplayName;
46+
import org.junit.jupiter.api.Test;
47+
import org.mockito.*;
48+
import org.springframework.amqp.AmqpRejectAndDontRequeueException;
49+
import org.springframework.amqp.core.Message;
50+
import org.springframework.amqp.rabbit.core.RabbitTemplate;
51+
52+
@DisplayName("Dispatcher — Conversion - Error Handling")
53+
class DispatcherConversionErrorHandlingTest {
54+
55+
private Dispatcher dispatcher;
56+
private MessageHandler messageHandler;
57+
private ConversionHandler conversionHandler;
58+
private RabbitTemplate rabbitTemplate;
59+
private MessagePersistenceService persistenceService;
60+
private HubConfiguration hubConfig;
61+
private ClientPropertiesRegistry clientPropertiesRegistry;
62+
private EdxlHandler edxlHandler;
63+
private XmlMapper xmlMapper;
64+
private ObjectMapper jsonMapper;
65+
private MeterRegistry registry;
66+
67+
@BeforeEach
68+
void setUp() {
69+
HubTestScaffolding.Hub hub = aHub().build();
70+
dispatcher = hub.dispatcher();
71+
messageHandler = hub.messageHandler();
72+
conversionHandler = hub.conversionHandler();
73+
rabbitTemplate = hub.rabbitTemplate();
74+
persistenceService = hub.persistenceService();
75+
hubConfig = hub.hubConfig();
76+
clientPropertiesRegistry = hub.clientPropertiesRegistry();
77+
edxlHandler = hub.edxlHandler();
78+
xmlMapper = hub.xmlMapper();
79+
jsonMapper = hub.jsonMapper();
80+
registry = hub.registry();
81+
echoConversionService(conversionHandler);
82+
}
83+
84+
@Test
85+
@DisplayName("should handle conversion service error correctly")
86+
public void shouldHandleConversionServiceError() throws IOException {
87+
// sdisC -> samuV3 on vhost 15-nexsis_v1.9 => transcoding triggered
88+
doReturn(NEXSIS_VHOST).when(hubConfig).getVhost();
89+
90+
MessageHandler messageHandlerSpy = spy(messageHandler);
91+
Dispatcher testDispatcher =
92+
new Dispatcher(
93+
messageHandlerSpy,
94+
rabbitTemplate,
95+
edxlHandler,
96+
xmlMapper,
97+
jsonMapper,
98+
conversionHandler,
99+
hubConfig,
100+
persistenceService,
101+
Tracer.NOOP);
102+
103+
Message receivedMessage =
104+
createMessage("EDXL-DE", JSON, SDIS_C_ROUTING_KEY, SAMU_V3_ROUTING_KEY);
105+
EdxlMessage edxlMessage =
106+
edxlHandler.deserializeJsonEDXL(
107+
new String(receivedMessage.getBody(), StandardCharsets.UTF_8));
108+
109+
String conversionErrorMessage = "Conversion service error message";
110+
failConversionService(
111+
conversionHandler,
112+
new ConversionException(conversionErrorMessage, edxlMessage.getDistributionID()));
113+
114+
assertThrows(
115+
AmqpRejectAndDontRequeueException.class,
116+
() -> testDispatcher.dispatch(receivedMessage));
117+
118+
ArgumentCaptor<ConversionException> exceptionCaptor =
119+
ArgumentCaptor.forClass(ConversionException.class);
120+
ArgumentCaptor<Message> messageCaptor = ArgumentCaptor.forClass(Message.class);
121+
122+
verify(messageHandlerSpy).handleError(exceptionCaptor.capture(), messageCaptor.capture());
123+
124+
ConversionException thrownException = exceptionCaptor.getValue();
125+
assertEquals(
126+
edxlMessage.getDistributionID(), thrownException.getReferencedDistributionID());
127+
assertTrue(thrownException.getMessage().contains(conversionErrorMessage));
128+
129+
Message handledMessage = messageCaptor.getValue();
130+
assertEquals(receivedMessage, handledMessage);
131+
}
132+
133+
@Test
134+
@DisplayName("should transfer to another vhost when an error is raised after message transfer")
135+
public void transferErrorToOtherVhost() throws IOException, ValidationException {
136+
ClientPropertiesRegistry clientPropertiesRegistrySpy =
137+
Mockito.spy(clientPropertiesRegistry);
138+
doReturn(clientPropertiesRegistrySpy).when(hubConfig).getClientPropertiesRegistry();
139+
doReturn("15-15_v2.0").when(hubConfig).getVhost();
140+
doReturn(new String[] {"1.5"})
141+
.when(clientPropertiesRegistrySpy)
142+
.getClientVersionsForPerimeter(SAMU_A_ROUTING_KEY, "15-15");
143+
Validator validatorMock = Mockito.mock(Validator.class);
144+
Mockito.doThrow(
145+
new SchemaValidationException(
146+
"Mock schema validation error", "mock_distribution_id"))
147+
.when(validatorMock)
148+
.validateJSON(anyString(), any());
149+
150+
MessageHandler messageHandlerSpy =
151+
new MessageHandler(
152+
rabbitTemplate,
153+
edxlHandler,
154+
hubConfig,
155+
validatorMock,
156+
registry,
157+
xmlMapper,
158+
jsonMapper,
159+
conversionHandler);
160+
Dispatcher dispatcherSpy =
161+
new Dispatcher(
162+
messageHandlerSpy,
163+
rabbitTemplate,
164+
edxlHandler,
165+
xmlMapper,
166+
jsonMapper,
167+
conversionHandler,
168+
hubConfig,
169+
persistenceService,
170+
Tracer.NOOP);
171+
172+
Message message = createMessage("EDXL-DE", JSON, SAMU_V1_ROUTING_KEY, SAMU_A_ROUTING_KEY);
173+
174+
String exchangeName = "transfer_15-15_v2.0_to_15-15_v1.5";
175+
176+
// Mock call to converter (return same payload for error message)
177+
AmqpRejectAndDontRequeueException errorThrown =
178+
assertThrows(
179+
AmqpRejectAndDontRequeueException.class,
180+
() -> dispatcherSpy.dispatch(message));
181+
182+
assertEquals("Mock schema validation error", errorThrown.getCause().getMessage());
183+
184+
assertThatMessageSentTo(rabbitTemplate, exchangeName, "fr.health.hub");
185+
}
186+
187+
@Test
188+
@DisplayName("should forward error message directly when error is received after conversion")
189+
public void sendErrorMessageToSameVhost() throws IOException {
190+
Message errorMessage = createMessage("hub-error-to-samuA", JSON);
191+
192+
dispatcher.dispatch(errorMessage);
193+
194+
assertThatMessageSentTo(rabbitTemplate, DISTRIBUTION_EXCHANGE, SAMU_A_INFO_QUEUE);
195+
}
196+
197+
@Test
198+
@DisplayName("should send error message to sender info queue when error is raised")
199+
public void sendErrorMessageWhenErrorIsRaised() throws IOException, ValidationException {
200+
ClientPropertiesRegistry clientPropertiesRegistrySpy =
201+
Mockito.spy(clientPropertiesRegistry);
202+
doReturn(clientPropertiesRegistrySpy).when(hubConfig).getClientPropertiesRegistry();
203+
doReturn("15-15_v1.5").when(hubConfig).getVhost();
204+
doReturn(new String[] {"1.5"})
205+
.when(clientPropertiesRegistrySpy)
206+
.getClientVersionsForPerimeter(SAMU_A_ROUTING_KEY, "15-15");
207+
// Default hub vhost is v2.1 and samuA declares v2.1: validation error is forwarded directly
208+
// to samuA's info queue without any conversion.
209+
Validator validatorMock = Mockito.mock(Validator.class);
210+
Mockito.doThrow(
211+
new SchemaValidationException(
212+
"Mock schema validation error", "mock_distribution_id"))
213+
.when(validatorMock)
214+
.validateJSON(anyString(), any());
215+
216+
MessageHandler messageHandlerSpy =
217+
new MessageHandler(
218+
rabbitTemplate,
219+
edxlHandler,
220+
hubConfig,
221+
validatorMock,
222+
registry,
223+
xmlMapper,
224+
jsonMapper,
225+
conversionHandler);
226+
Dispatcher dispatcherSpy =
227+
new Dispatcher(
228+
messageHandlerSpy,
229+
rabbitTemplate,
230+
edxlHandler,
231+
xmlMapper,
232+
jsonMapper,
233+
conversionHandler,
234+
hubConfig,
235+
persistenceService,
236+
Tracer.NOOP);
237+
238+
Message message = createMessage("EDXL-DE", JSON);
239+
240+
AmqpRejectAndDontRequeueException errorThrown =
241+
assertThrows(
242+
AmqpRejectAndDontRequeueException.class,
243+
() -> dispatcherSpy.dispatch(message));
244+
245+
assertEquals("Mock schema validation error", errorThrown.getCause().getMessage());
246+
247+
assertThatMessageSentTo(rabbitTemplate, DISTRIBUTION_EXCHANGE, SAMU_A_INFO_QUEUE);
248+
}
249+
250+
@Test
251+
@DisplayName("should log referencedDistributionID for ACK ReferenceWrapper")
252+
public void shouldLogReferencedDistributionIdForAckReferenceWrapper() throws Exception {
253+
Message message = createMessage("rc-ref", JSON);
254+
255+
// Log capturing setup
256+
org.slf4j.Logger logger = org.slf4j.LoggerFactory.getLogger(MessageHandler.class);
257+
ch.qos.logback.classic.Logger logbackLogger = (ch.qos.logback.classic.Logger) logger;
258+
ch.qos.logback.classic.Level originalLevel = logbackLogger.getLevel();
259+
logbackLogger.setLevel(ch.qos.logback.classic.Level.INFO);
260+
ch.qos.logback.core.read.ListAppender<ch.qos.logback.classic.spi.ILoggingEvent> appender =
261+
new ch.qos.logback.core.read.ListAppender<>();
262+
appender.start();
263+
logbackLogger.addAppender(appender);
264+
265+
dispatcher.dispatch(message);
266+
267+
boolean foundReceivedLog =
268+
appender.list.stream()
269+
.anyMatch(
270+
event ->
271+
event.getFormattedMessage()
272+
.contains(
273+
"Received Ack: message with referenced distributionId fr.health.samuB_2607723d-507d-4cbf-bf74-12345f7064cd"));
274+
assertTrue(
275+
foundReceivedLog,
276+
"Received log should contain referenced distributionId: fr.health.samuB_2607723d-507d-4cbf-bf74-12345f7064cd");
277+
278+
boolean foundForwardingLog =
279+
appender.list.stream()
280+
.anyMatch(
281+
event ->
282+
event.getFormattedMessage()
283+
.contains(
284+
"Forwarding Ack: message with referenced distributionId fr.health.samuB_2607723d-507d-4cbf-bf74-12345f7064cd"));
285+
assertTrue(
286+
foundForwardingLog,
287+
"Forwarding log should contain referenced distributionId: fr.health.samuB_2607723d-507d-4cbf-bf74-12345f7064cd");
288+
289+
// Cleanup
290+
logbackLogger.detachAppender(appender);
291+
logbackLogger.setLevel(originalLevel);
292+
}
293+
}

0 commit comments

Comments
 (0)