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;
+    }
+}

Reply via email to