diff --git a/Foundation/org.egovframe.rte.fdl.reactive/src/main/java/org/egovframe/rte/fdl/reactive/logging/EgovMdcContextConfig.java b/Foundation/org.egovframe.rte.fdl.reactive/src/main/java/org/egovframe/rte/fdl/reactive/logging/EgovMdcContextConfig.java index ade6005e..39f5ad63 100755 --- a/Foundation/org.egovframe.rte.fdl.reactive/src/main/java/org/egovframe/rte/fdl/reactive/logging/EgovMdcContextConfig.java +++ b/Foundation/org.egovframe.rte.fdl.reactive/src/main/java/org/egovframe/rte/fdl/reactive/logging/EgovMdcContextConfig.java @@ -110,11 +110,13 @@ public void onNext(T t) { @Override public void onError(Throwable throwable) { + copyToMdc(coreSubscriber.currentContext()); coreSubscriber.onError(throwable); } @Override public void onComplete() { + copyToMdc(coreSubscriber.currentContext()); coreSubscriber.onComplete(); } diff --git a/Foundation/org.egovframe.rte.fdl.reactive/src/test/java/org/egovframe/rte/fdl/reactive/logging/EgovMdcContextConfigTest.java b/Foundation/org.egovframe.rte.fdl.reactive/src/test/java/org/egovframe/rte/fdl/reactive/logging/EgovMdcContextConfigTest.java new file mode 100644 index 00000000..d72afd92 --- /dev/null +++ b/Foundation/org.egovframe.rte.fdl.reactive/src/test/java/org/egovframe/rte/fdl/reactive/logging/EgovMdcContextConfigTest.java @@ -0,0 +1,85 @@ +package org.egovframe.rte.fdl.reactive.logging; + +import org.junit.jupiter.api.AfterEach; +import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.Test; +import org.slf4j.MDC; +import reactor.core.publisher.Mono; +import reactor.core.scheduler.Scheduler; +import reactor.core.scheduler.Schedulers; + +import java.util.concurrent.atomic.AtomicReference; + +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertThrows; + +public class EgovMdcContextConfigTest { + + private final EgovMdcContextConfig config = new EgovMdcContextConfig(); + + /** 시퀀스의 상류 절반이 도는, 재사용되는 스레드 */ + private Scheduler worker; + + /** publishOn 으로 하류 절반을 넘겨받는 스레드 */ + private Scheduler consumer; + + @BeforeEach + public void setUp() { + worker = Schedulers.newSingle("mdc-worker"); + consumer = Schedulers.newSingle("mdc-consumer"); + config.contextOperatorHook(); + MDC.clear(); + } + + @AfterEach + public void tearDown() { + config.cleanupHook(); + worker.dispose(); + consumer.dispose(); + MDC.clear(); + } + + /** + * userId=A 를 worker 스레드의 MDC 에 남긴 뒤, 같은 worker 스레드에서 값 없이 완료되는 + * 시퀀스를 구독한다. 두 번째 시퀀스의 종료 콜백은 자신의 Context 값인 B 를 봐야 한다. + */ + @Test + public void completionOfNextSequenceSeesItsOwnContext() { + Mono.just("a").hide() + .subscribeOn(worker) + .publishOn(consumer) + .contextWrite(context -> context.put("userId", "A")) + .block(); + + AtomicReference seen = new AtomicReference<>(); + Mono.empty() + .subscribeOn(worker) + .doOnTerminate(() -> seen.set(MDC.get("userId"))) + .contextWrite(context -> context.put("userId", "B")) + .block(); + + assertEquals("B", seen.get()); + } + + /** + * 두 번째 시퀀스가 값 없이 에러로 끝나는 경우도 마찬가지다. + */ + @Test + public void errorOfNextSequenceSeesItsOwnContext() { + Mono.just("a").hide() + .subscribeOn(worker) + .publishOn(consumer) + .contextWrite(context -> context.put("userId", "A")) + .block(); + + AtomicReference seen = new AtomicReference<>(); + Mono error = Mono.error(new IllegalStateException("boom")).hide() + .subscribeOn(worker) + .doOnTerminate(() -> seen.set(MDC.get("userId"))) + .contextWrite(context -> context.put("userId", "B")); + assertThrows(IllegalStateException.class, error::block); + + assertEquals("B", seen.get()); + } + +}