This is an automated email from the ASF dual-hosted git repository. reta pushed a commit to branch 4.1.x-fixes in repository https://gitbox.apache.org/repos/asf/cxf.git
commit 12a5963373986599ebf0f462f5f747a6ff631f09 Author: Valentino Porta <[email protected]> AuthorDate: Tue Oct 6 15:33:48 2026 +0200 alternative solution CXF-9251 (#3541) * alternative solution CXF-9251 * add suppressed Exception * null safe ex.addSuppressed(suppressed); * adjust comments * add two new test cases: - shouldNotLogTwice (even if close() method is called multiple times) - shouldLogEvenIfIOException (even if an IOException is thrown and there is only a partial write on the http socket) * add more comment to test cases * simulate also memory leak * Add test cases --------- Co-authored-by: Andriy Redko <[email protected]> (cherry picked from commit cbc14a26db7754b95d3e0e0857e324454f44fad7) --- .../cxf/ext/logging/LoggingOutputStream.java | 60 ++++ .../cxf/ext/logging/LoggingOutInterceptorTest.java | 324 +++++++++++++++++++++ 2 files changed, 384 insertions(+) diff --git a/rt/features/logging/src/main/java/org/apache/cxf/ext/logging/LoggingOutputStream.java b/rt/features/logging/src/main/java/org/apache/cxf/ext/logging/LoggingOutputStream.java index b61d029814e..202f577672b 100644 --- a/rt/features/logging/src/main/java/org/apache/cxf/ext/logging/LoggingOutputStream.java +++ b/rt/features/logging/src/main/java/org/apache/cxf/ext/logging/LoggingOutputStream.java @@ -21,11 +21,13 @@ package org.apache.cxf.ext.logging; import java.io.IOException; import java.io.OutputStream; +import java.util.concurrent.atomic.AtomicBoolean; import org.apache.cxf.io.CacheAndWriteOutputStream; public class LoggingOutputStream extends CacheAndWriteOutputStream { private boolean skipFlushingFlowThroughStream; + private final AtomicBoolean closed = new AtomicBoolean(); LoggingOutputStream(OutputStream stream) { super(stream); @@ -71,4 +73,62 @@ public class LoggingOutputStream extends CacheAndWriteOutputStream { super.writeCacheTo(out, charsetName, limit); skipFlushingFlowThroughStream = false; } + + + /** + * CXF-9251 + * We override the write() methods in order to catch some error that would not + * close the CachedOutputStream (ex. IOException "Connection reset by peer"). + * This caused ghost/delayed OUT log and possible memory-leak due to DelayedCachedOutputStreamCleaner + */ + @Override + public void write(byte[] b) throws IOException { + try { + super.write(b); + } catch (RuntimeException | IOException ex) { + handleIoException(ex); + throw ex; + } + } + + @Override + public void write(byte[] b, int off, int len) throws IOException { + try { + super.write(b, off, len); + } catch (RuntimeException | IOException ex) { + handleIoException(ex); + throw ex; + } + } + + @Override + public void write(int b) throws IOException { + try { + super.write(b); + } catch (RuntimeException | IOException ex) { + handleIoException(ex); + throw ex; + } + } + + @Override + public void close() throws IOException { + // Ensure closing only one time + if (closed.compareAndSet(false, true)) { + super.close(); + } + } + + private void handleIoException(Exception ex) { + try { + // Close this CachedOutputStream + // Write method maybe called more than one time... but we already consume the stream the first time + // So additional call to this.close would produce nothing + this.close(); + } catch (Exception suppressed) { + if (ex != null) { + ex.addSuppressed(suppressed); + } + } + } } diff --git a/rt/features/logging/src/test/java/org/apache/cxf/ext/logging/LoggingOutInterceptorTest.java b/rt/features/logging/src/test/java/org/apache/cxf/ext/logging/LoggingOutInterceptorTest.java new file mode 100644 index 00000000000..87f82c5def7 --- /dev/null +++ b/rt/features/logging/src/test/java/org/apache/cxf/ext/logging/LoggingOutInterceptorTest.java @@ -0,0 +1,324 @@ +/** + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, + * software distributed under the License is distributed on an + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + * KIND, either express or implied. See the License for the + * specific language governing permissions and limitations + * under the License. + */ + +package org.apache.cxf.ext.logging; + +import java.io.ByteArrayOutputStream; +import java.io.IOException; +import java.io.OutputStream; +import java.io.UncheckedIOException; +import java.nio.charset.StandardCharsets; + +import org.apache.cxf.Bus; +import org.apache.cxf.BusFactory; +import org.apache.cxf.ext.logging.event.LogEvent; +import org.apache.cxf.io.CachedOutputStream; +import org.apache.cxf.io.DelayedCachedOutputStreamCleaner; +import org.apache.cxf.message.ExchangeImpl; +import org.apache.cxf.message.Message; +import org.apache.cxf.message.MessageImpl; + +import org.junit.After; +import org.junit.Before; +import org.junit.Test; + +import static org.hamcrest.MatcherAssert.assertThat; +import static org.hamcrest.Matchers.equalTo; +import static org.hamcrest.Matchers.equalToIgnoringCase; +import static org.hamcrest.Matchers.hasSize; +import static org.junit.Assert.assertThrows; + +public class LoggingOutInterceptorTest { + private Bus bus; + private DelayedCachedOutputStreamCleaner cleaner; + private TestEventSender sender; + private LoggingOutInterceptor interceptor; + private Message message; + + @Before + public void setUp() { + bus = BusFactory.getDefaultBus(true); + cleaner = bus.getExtension(DelayedCachedOutputStreamCleaner.class); + sender = new TestEventSender(); + interceptor = new LoggingOutInterceptor(sender); + message = new MessageImpl(); + message.setExchange(new ExchangeImpl()); + } + + @After + public void tearDown() { + assertThat(cleaner.size(), equalTo(0)); + } + + @Test + public void shouldLogMultipartPayload() throws IOException { + message.put(Message.ENDPOINT_ADDRESS, "http://localhost:9001/"); + message.put(Message.REQUEST_URI, "/api"); + + final StringBuilder buf = content(); + String ct = "multipart/related; type=\"application/xop+xml\"; " + + "boundary=\"----=_Part_0_2180223.1203118300920\""; + + final byte[] bytes = buf.toString().getBytes(StandardCharsets.UTF_8); + final OutputStream os = new ByteArrayOutputStream(); + message.setContent(OutputStream.class, os); + message.put(Message.CONTENT_TYPE, ct); + + interceptor.setInMemThreshold(1); + interceptor.addBinaryContentMediaTypes("application/xop+xml"); + interceptor.setLogMultipart(true); + interceptor.setLogBinary(true); + interceptor.handleMessage(message); + + final OutputStream cached = message.getContent(OutputStream.class); + cached.write(bytes, 0, bytes.length); + cached.close(); + os.close(); + + assertThat(sender.getEvents(), hasSize(1)); + final LogEvent event = sender.getEvents().get(0); + + assertThat(event.getPayload(), equalToIgnoringCase(buf.toString())); + } + + @Test + public void shouldLogMultipartPayloadOnClean() throws IOException { + message.put(Message.ENDPOINT_ADDRESS, "http://localhost:9001/"); + message.put(Message.REQUEST_URI, "/api"); + + final StringBuilder buf = content(); + String ct = "multipart/related; type=\"application/xop+xml\"; " + + "boundary=\"----=_Part_0_2180223.1203118300920\""; + + final byte[] bytes = buf.toString().getBytes(StandardCharsets.UTF_8); + final OutputStream os = new ByteArrayOutputStream(); + message.setContent(OutputStream.class, os); + message.put(Message.CONTENT_TYPE, ct); + + interceptor.setInMemThreshold(1); + interceptor.addBinaryContentMediaTypes("application/xop+xml"); + interceptor.setLogMultipart(true); + interceptor.setLogBinary(true); + interceptor.handleMessage(message); + + final OutputStream cached = message.getContent(OutputStream.class); + cached.write(bytes, 0, bytes.length); + os.close(); + + // We did not close the cached stream, should be subject of cleanup + assertThat(cleaner.size(), equalTo(1)); + + cleaner.forceClean(); + assertThat(cleaner.size(), equalTo(0)); + + assertThat(sender.getEvents(), hasSize(1)); + final LogEvent event = sender.getEvents().get(0); + + assertThat(event.getPayload(), equalToIgnoringCase(buf.toString())); + } + + @Test + public void shouldLogMultipartPayloadCachedOutputStream() throws IOException { + message.put(Message.ENDPOINT_ADDRESS, "http://localhost:9001/"); + message.put(Message.REQUEST_URI, "/api"); + + final StringBuilder buf = content(); + String ct = "multipart/related; type=\"application/xop+xml\"; " + + "boundary=\"----=_Part_0_2180223.1203118300920\""; + + final byte[] bytes = buf.toString().getBytes(StandardCharsets.UTF_8); + final CachedOutputStream os = new CachedOutputStream(); + message.setContent(OutputStream.class, os); + message.put(Message.CONTENT_TYPE, ct); + + interceptor.setInMemThreshold(1); + interceptor.addBinaryContentMediaTypes("application/xop+xml"); + interceptor.setLogMultipart(true); + interceptor.setLogBinary(true); + interceptor.handleMessage(message); + + final OutputStream cached = message.getContent(OutputStream.class); + cached.write(bytes, 0, bytes.length); + cached.close(); + os.close(); + + assertThat(sender.getEvents(), hasSize(1)); + final LogEvent event = sender.getEvents().get(0); + + assertThat(event.getPayload(), equalToIgnoringCase(buf.toString())); + } + + @Test + public void shouldLogMultipartPayloadWithExceptionOnClose() throws IOException { + message.put(Message.ENDPOINT_ADDRESS, "http://localhost:9001/"); + message.put(Message.REQUEST_URI, "/api"); + + final StringBuilder buf = content(); + String ct = "multipart/related; type=\"application/xop+xml\"; " + + "boundary=\"----=_Part_0_2180223.1203118300920\""; + + final byte[] bytes = buf.toString().getBytes(StandardCharsets.UTF_8); + final OutputStream os = new ByteArrayOutputStream() { + public void close() throws IOException { + throw new IOException("Simulated"); + } + }; + message.setContent(OutputStream.class, os); + message.put(Message.CONTENT_TYPE, ct); + + interceptor.setInMemThreshold(1); + interceptor.addBinaryContentMediaTypes("application/xop+xml"); + interceptor.setLogMultipart(true); + interceptor.setLogBinary(true); + interceptor.handleMessage(message); + + final OutputStream cached = message.getContent(OutputStream.class); + cached.write(bytes, 0, bytes.length); + assertThrows(IOException.class, () -> cached.close()); + + assertThat(sender.getEvents(), hasSize(1)); + final LogEvent event = sender.getEvents().get(0); + + assertThat(event.getPayload(), equalToIgnoringCase(buf.toString())); + } + + @Test + public void shouldLogMultipartPayloadWithExceptionOnFlush() throws IOException { + message.put(Message.ENDPOINT_ADDRESS, "http://localhost:9001/"); + message.put(Message.REQUEST_URI, "/api"); + + final StringBuilder buf = content(); + String ct = "multipart/related; type=\"application/xop+xml\"; " + + "boundary=\"----=_Part_0_2180223.1203118300920\""; + + final byte[] bytes = buf.toString().getBytes(StandardCharsets.UTF_8); + final OutputStream os = new ByteArrayOutputStream() { + public void flush() throws IOException { + throw new IOException("Simulated"); + } + }; + message.setContent(OutputStream.class, os); + message.put(Message.CONTENT_TYPE, ct); + + interceptor.setInMemThreshold(1); + interceptor.addBinaryContentMediaTypes("application/xop+xml"); + interceptor.setLogMultipart(true); + interceptor.setLogBinary(true); + interceptor.handleMessage(message); + + final OutputStream cached = message.getContent(OutputStream.class); + cached.write(bytes, 0, bytes.length); + assertThrows(IOException.class, () -> cached.flush()); + cached.close(); + + assertThat(sender.getEvents(), hasSize(1)); + final LogEvent event = sender.getEvents().get(0); + + assertThat(event.getPayload(), equalToIgnoringCase(buf.toString())); + } + + @Test + public void shouldLogMultipartPayloadClosedTwice() throws IOException { + message.put(Message.ENDPOINT_ADDRESS, "http://localhost:9001/"); + message.put(Message.REQUEST_URI, "/api"); + + final StringBuilder buf = content(); + String ct = "multipart/related; type=\"application/xop+xml\"; " + + "boundary=\"----=_Part_0_2180223.1203118300920\""; + + final byte[] bytes = buf.toString().getBytes(StandardCharsets.UTF_8); + final OutputStream os = new ByteArrayOutputStream(); + message.setContent(OutputStream.class, os); + message.put(Message.CONTENT_TYPE, ct); + + interceptor.setInMemThreshold(1); + interceptor.addBinaryContentMediaTypes("application/xop+xml"); + interceptor.setLogMultipart(true); + interceptor.setLogBinary(true); + interceptor.handleMessage(message); + + final OutputStream cached = message.getContent(OutputStream.class); + cached.write(bytes, 0, bytes.length); + + cached.close(); + cached.close(); + + assertThat(sender.getEvents(), hasSize(1)); + final LogEvent event = sender.getEvents().get(0); + + assertThat(event.getPayload(), equalToIgnoringCase(buf.toString())); + } + + @Test + public void shouldLogMultipartPayloadWithExceptionOnWrite() throws IOException { + message.put(Message.ENDPOINT_ADDRESS, "http://localhost:9001/"); + message.put(Message.REQUEST_URI, "/api"); + + final StringBuilder buf = content(); + String ct = "multipart/related; type=\"application/xop+xml\"; " + + "boundary=\"----=_Part_0_2180223.1203118300920\""; + + final byte[] bytes = buf.toString().getBytes(StandardCharsets.UTF_8); + final OutputStream os = new ByteArrayOutputStream() { + @Override + public synchronized void write(byte[] b, int off, int len) { + if (len == 1) { + throw new UncheckedIOException(new IOException("Simulated")); + } else { + super.write(bytes, off, len); + } + } + }; + message.setContent(OutputStream.class, os); + message.put(Message.CONTENT_TYPE, ct); + + interceptor.setInMemThreshold(1); + interceptor.addBinaryContentMediaTypes("application/xop+xml"); + interceptor.setLogMultipart(true); + interceptor.setLogBinary(true); + interceptor.handleMessage(message); + + final OutputStream cached = message.getContent(OutputStream.class); + cached.write(bytes, 0, bytes.length - 1); + assertThrows(UncheckedIOException.class, () -> cached.write(bytes, bytes.length - 1, 1)); + + // We did not close the cached stream, should be subject of cleanup + assertThat(sender.getEvents(), hasSize(1)); + final LogEvent event = sender.getEvents().get(0); + + buf.setLength(buf.length() - 1); + assertThat(event.getPayload(), equalToIgnoringCase(buf.toString())); + } + + private static StringBuilder content() { + StringBuilder buf = new StringBuilder(512); + buf.append("------=_Part_0_2180223.1203118300920\n"); + buf.append("Content-Type: application/xop+xml; charset=UTF-8; type=\"text/xml\"\n"); + buf.append("Content-Transfer-Encoding: 8bit\n"); + buf.append("Content-ID: <[email protected]>\n"); + buf.append('\n'); + buf.append("<soap:Envelope xmlns:soap=\"http://schemas.xmlsoap.org/soap/envelope/\" " + + "xmlns:xsd=\"http://www.w3.org/2001/XMLSchema\" " + + "xmlns:xsi=\"http://www.w3.org/2001/XMLSchema-instance\">" + + "<soap:Body><getNextMessage xmlns=\"http://foo.bar\" /></soap:Body>" + + "</soap:Envelope>\n"); + buf.append("------=_Part_0_2180223.1203118300920--\n"); + return buf; + } +}
