diff --git a/rt/features/logging/src/main/java/org/apache/cxf/ext/logging/AbstractLoggingInterceptor.java b/rt/features/logging/src/main/java/org/apache/cxf/ext/logging/AbstractLoggingInterceptor.java index 4f7c7c015a6..8fce3612eb8 100644 --- a/rt/features/logging/src/main/java/org/apache/cxf/ext/logging/AbstractLoggingInterceptor.java +++ b/rt/features/logging/src/main/java/org/apache/cxf/ext/logging/AbstractLoggingInterceptor.java @@ -25,20 +25,21 @@ import java.util.regex.Pattern; import org.apache.cxf.common.util.PropertyUtils; -import org.apache.cxf.ext.logging.event.DefaultLogEventMapper; -import org.apache.cxf.ext.logging.event.LogEvent; -import org.apache.cxf.ext.logging.event.LogEventSender; -import org.apache.cxf.ext.logging.event.PrettyLoggingFilter; +import org.apache.cxf.ext.logging.event.*; import org.apache.cxf.interceptor.Fault; import org.apache.cxf.message.Exchange; import org.apache.cxf.message.Message; import org.apache.cxf.phase.AbstractPhaseInterceptor; +import static org.apache.cxf.ext.logging.event.DefaultLogEventMapper.normalizeFlow; + public abstract class AbstractLoggingInterceptor extends AbstractPhaseInterceptor { public static final int DEFAULT_LIMIT = 48 * 1024; public static final int DEFAULT_THRESHOLD = -1; public static final String CONTENT_SUPPRESSED = "--- Content suppressed ---"; protected static final String LIVE_LOGGING_PROP = "org.apache.cxf.logging.enable"; + protected static final String IDEMPOTENT_LOGGING_PROP = "org.apache.cxf.idempotent.logging."; // the EventType (flow) and ExchangeId will be concatenated + private static final Pattern BOUNDARY_PATTERN = Pattern.compile("^--(\\S*)$", Pattern.MULTILINE); private static final Pattern CONTENT_TYPE_PATTERN = @@ -62,11 +63,31 @@ public AbstractLoggingInterceptor(String phase, LogEventSender sender) { this.eventMapper = new DefaultLogEventMapper(maskSensitiveHelper); } + // If the properties is set somewhere else (Bus...etc...) protected static boolean isLoggingDisabledNow(Message message) throws Fault { Object liveLoggingProp = message.getContextualProperty(LIVE_LOGGING_PROP); return liveLoggingProp != null && PropertyUtils.isFalse(liveLoggingProp); } + // The concatenated flow is added in order to enhance resilience against misuse and underlying framework + // (Reuse of the same Message object with properties still there) + // The message will be logged once per flow per ExchangeId + // If previous properties (ex. IDEMPOTENT_LOGGING_PROP + REQ_IN + ExchangeId) are still there... this will search + // only for the right properties (ex. IDEMPOTENT_LOGGING_PROP + RESP_OUT + ExchangeId) + protected boolean isLoggingDisabledForThisFlow(Message message) throws Fault { + Object idempotentLoggingProp = message.getContextualProperty(getIdempotentDisableLogKey(message)); //idempotency per Flow per ExchangeId + return idempotentLoggingProp != null && PropertyUtils.isFalse(idempotentLoggingProp); + } + protected void disableFutureLoggingForThisFlow(Message message) throws Fault { + message.put(getIdempotentDisableLogKey(message), Boolean.FALSE); + } + + // IDEMPOTENT_LOGGING_PROP + FLOW + ExchangeId + protected String getIdempotentDisableLogKey(Message message){ + createExchangeId(message); //Redundant + return IDEMPOTENT_LOGGING_PROP + normalizeFlow(eventMapper.getEventType(message)) + '.' + message.getExchange().get(LogEvent.KEY_EXCHANGE_ID); + } + public void addBinaryContentMediaTypes(String mediaTypes) { eventMapper.addBinaryContentMediaTypes(mediaTypes); } diff --git a/rt/features/logging/src/main/java/org/apache/cxf/ext/logging/LoggingInInterceptor.java b/rt/features/logging/src/main/java/org/apache/cxf/ext/logging/LoggingInInterceptor.java index 436bbf3d16d..50251117310 100644 --- a/rt/features/logging/src/main/java/org/apache/cxf/ext/logging/LoggingInInterceptor.java +++ b/rt/features/logging/src/main/java/org/apache/cxf/ext/logging/LoggingInInterceptor.java @@ -81,15 +81,16 @@ public Collection> getAdditionalInterceptors } public void handleMessage(Message message) throws Fault { - if (isLoggingDisabledNow(message)) { + + createExchangeId(message); + if (isLoggingDisabledNow(message) || isLoggingDisabledForThisFlow(message)) { return; - } else { - //ensure only logging once for a certain message - //this can prevent message logging again when fault - //happen after PRE_INVOKE phase(rewind calls into LoggingInFaultInterceptor) - message.put(LIVE_LOGGING_PROP, Boolean.FALSE); } - createExchangeId(message); + //ensure only logging once for a certain message + //this can prevent message logging again when fault + //happen after PRE_INVOKE phase(rewind calls into LoggingInFaultInterceptor) + disableFutureLoggingForThisFlow(message); + final LogEvent event = eventMapper.map(message, sensitiveProtocolHeaderNames); if (shouldLogContent(event)) { addContent(message, event); diff --git a/rt/features/logging/src/main/java/org/apache/cxf/ext/logging/LoggingOutInterceptor.java b/rt/features/logging/src/main/java/org/apache/cxf/ext/logging/LoggingOutInterceptor.java index 7e68a7c5cca..dd19d1df015 100644 --- a/rt/features/logging/src/main/java/org/apache/cxf/ext/logging/LoggingOutInterceptor.java +++ b/rt/features/logging/src/main/java/org/apache/cxf/ext/logging/LoggingOutInterceptor.java @@ -59,15 +59,16 @@ public LoggingOutInterceptor(LogEventSender sender) { } public void handleMessage(Message message) throws Fault { - if (isLoggingDisabledNow(message)) { + createExchangeId(message); + if (isLoggingDisabledNow(message) || isLoggingDisabledForThisFlow(message)) { return; - } else { - //ensure only logging once for a certain message - //this can prevent message logging again when fault - //happen after PRE_STREAM phase(LoggingOutInterceptor is called both in out chain and fault out chain) - message.put(LIVE_LOGGING_PROP, Boolean.FALSE); } - createExchangeId(message); + + //ensure only logging once for a certain message + //this can prevent message logging again when fault + //happen after PRE_STREAM phase(LoggingOutInterceptor is called both in out chain and fault out chain) + disableFutureLoggingForThisFlow(message); + final OutputStream os = message.getContent(OutputStream.class); if (os != null) { LoggingCallback callback = new LoggingCallback(sender, message, os, limit); diff --git a/rt/features/logging/src/main/java/org/apache/cxf/ext/logging/event/DefaultLogEventMapper.java b/rt/features/logging/src/main/java/org/apache/cxf/ext/logging/event/DefaultLogEventMapper.java index ee3a386bec5..1dc4b4894d9 100644 --- a/rt/features/logging/src/main/java/org/apache/cxf/ext/logging/event/DefaultLogEventMapper.java +++ b/rt/features/logging/src/main/java/org/apache/cxf/ext/logging/event/DefaultLogEventMapper.java @@ -47,6 +47,8 @@ import org.apache.cxf.ws.addressing.AddressingProperties; import org.apache.cxf.ws.addressing.ContextUtils; +import static org.apache.cxf.ext.logging.event.EventType.*; + public class DefaultLogEventMapper { public static final String MASKED_HEADER_VALUE = "XXX"; private static final Set DEFAULT_BINARY_CONTENT_MEDIA_TYPES; @@ -352,11 +354,25 @@ public EventType getEventType(Message message) { return isRequestor ? EventType.REQ_OUT : EventType.RESP_OUT; } if (isFault) { - return EventType.FAULT_IN; + return FAULT_IN; } return isRequestor ? EventType.RESP_IN : EventType.REQ_IN; } + /** + * Get the normalize 'flow' from the eventType + * + * @param eventType + * @return normalized eventType + */ + public static EventType normalizeFlow(EventType eventType){ + return switch (eventType) { + case FAULT_IN -> RESP_IN; + case FAULT_OUT -> RESP_OUT; + default -> eventType; + }; + } + /** * For REST we also consider a response to be a fault if the operation is not found or the response code * is an error diff --git a/rt/features/logging/src/test/java/org/apache/cxf/ext/logging/DefaultLogEventMapperTest.java b/rt/features/logging/src/test/java/org/apache/cxf/ext/logging/DefaultLogEventMapperTest.java index f502d8e8de0..0154610bea9 100644 --- a/rt/features/logging/src/test/java/org/apache/cxf/ext/logging/DefaultLogEventMapperTest.java +++ b/rt/features/logging/src/test/java/org/apache/cxf/ext/logging/DefaultLogEventMapperTest.java @@ -43,6 +43,8 @@ import org.junit.Test; import static org.apache.cxf.ext.logging.event.DefaultLogEventMapper.MASKED_HEADER_VALUE; +import static org.apache.cxf.ext.logging.event.DefaultLogEventMapper.normalizeFlow; +import static org.apache.cxf.ext.logging.event.EventType.*; import static org.hamcrest.CoreMatchers.is; import static org.hamcrest.CoreMatchers.nullValue; import static org.hamcrest.MatcherAssert.assertThat; @@ -81,7 +83,7 @@ public void testPreflightRequestEventType() { message.setExchange(exchange); exchange.setOutMessage(message); LogEvent event = mapper.map(message, Collections.emptySet()); - assertEquals(EventType.RESP_OUT, event.getType()); + assertEquals(RESP_OUT, event.getType()); } /** @@ -206,4 +208,10 @@ public void testNoSubjectReturned() { LogEvent event = Subject.doAs(subject, (PrivilegedAction) () -> mapper.map(message)); assertThat(event.getPrincipal(), is(nullValue())); } + + @Test + public void testNormalizeFlow(){ + assertEquals(RESP_IN, normalizeFlow(FAULT_IN)); + assertEquals(RESP_OUT, normalizeFlow(FAULT_OUT)); + } } diff --git a/rt/features/logging/src/test/java/org/apache/cxf/ext/logging/LoggingInInterceptorTest.java b/rt/features/logging/src/test/java/org/apache/cxf/ext/logging/LoggingInInterceptorTest.java index d7b75781504..290fe3eaf39 100644 --- a/rt/features/logging/src/test/java/org/apache/cxf/ext/logging/LoggingInInterceptorTest.java +++ b/rt/features/logging/src/test/java/org/apache/cxf/ext/logging/LoggingInInterceptorTest.java @@ -22,15 +22,12 @@ import java.io.IOException; import java.io.OutputStream; import java.nio.charset.StandardCharsets; -import java.util.Arrays; -import java.util.Collections; -import java.util.HashMap; -import java.util.HashSet; -import java.util.Map; -import java.util.Set; +import java.util.*; +import org.apache.cxf.ext.logging.event.DefaultLogEventMapper; import org.apache.cxf.ext.logging.event.LogEvent; import org.apache.cxf.io.CachedOutputStream; +import org.apache.cxf.message.Exchange; import org.apache.cxf.message.ExchangeImpl; import org.apache.cxf.message.Message; import org.apache.cxf.message.MessageImpl; @@ -38,11 +35,14 @@ import org.junit.Before; import org.junit.Test; +import static org.apache.cxf.common.util.PropertyUtils.isFalse; +import static org.apache.cxf.ext.logging.AbstractLoggingInterceptor.IDEMPOTENT_LOGGING_PROP; import static org.apache.cxf.ext.logging.event.DefaultLogEventMapper.MASKED_HEADER_VALUE; +import static org.apache.cxf.ext.logging.event.EventType.*; import static org.hamcrest.MatcherAssert.assertThat; import static org.hamcrest.Matchers.equalToIgnoringCase; import static org.hamcrest.Matchers.hasSize; -import static org.junit.Assert.assertEquals; +import static org.junit.Assert.*; public class LoggingInInterceptorTest { private static final String TEST_HEADER_VALUE = "TestValue"; @@ -234,4 +234,32 @@ public void shouldLogMultipartPayloadNoHeaders() throws IOException { assertThat(event.getPayload(), equalToIgnoringCase(buf.toString())); } + + @Test + public void testLoggingEnable(){ + Message message = new MessageImpl(); + Exchange exchange = new ExchangeImpl(); + exchange.setOutMessage(message); + message.setExchange(exchange); + message.put(Message.REQUESTOR_ROLE, Boolean.TRUE); + + DefaultLogEventMapper mapper = new DefaultLogEventMapper(); + assertEquals(FAULT_OUT, mapper.getEventType(message)); + + assertNull(message.getExchange().get(LogEvent.KEY_EXCHANGE_ID)); + assertFalse(interceptor.isLoggingDisabledForThisFlow(message)); + assertNotNull(message.getExchange().get(LogEvent.KEY_EXCHANGE_ID)); + interceptor.disableFutureLoggingForThisFlow(message); + assertTrue(interceptor.isLoggingDisabledForThisFlow(message)); + + assertNull(message.getContextualProperty(IDEMPOTENT_LOGGING_PROP + REQ_IN + '.' + message.getExchange().get(LogEvent.KEY_EXCHANGE_ID))); + assertNull(message.getContextualProperty(IDEMPOTENT_LOGGING_PROP + REQ_OUT + '.' + message.getExchange().get(LogEvent.KEY_EXCHANGE_ID))); + assertNull(message.getContextualProperty(IDEMPOTENT_LOGGING_PROP + RESP_IN + '.' + message.getExchange().get(LogEvent.KEY_EXCHANGE_ID))); + assertNotNull(message.getContextualProperty(IDEMPOTENT_LOGGING_PROP + RESP_OUT + '.' + message.getExchange().get(LogEvent.KEY_EXCHANGE_ID))); //The only present // FAULT_OUT normalized + assertNull(message.getContextualProperty(IDEMPOTENT_LOGGING_PROP + FAULT_IN + '.' + message.getExchange().get(LogEvent.KEY_EXCHANGE_ID))); + assertNull(message.getContextualProperty(IDEMPOTENT_LOGGING_PROP + FAULT_OUT + '.' + message.getExchange().get(LogEvent.KEY_EXCHANGE_ID))); + + assertTrue(isFalse(message.getContextualProperty(IDEMPOTENT_LOGGING_PROP + RESP_OUT + '.' + message.getExchange().get(LogEvent.KEY_EXCHANGE_ID)))); + + } }