Skip to content
Open
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -24,6 +24,7 @@
import java.io.PrintWriter;
import java.io.StringWriter;
import java.io.Writer;
import java.util.concurrent.atomic.AtomicBoolean;

import org.apache.cxf.common.injection.NoJSR250Annotations;
import org.apache.cxf.common.util.StringUtils;
Expand Down Expand Up @@ -70,7 +71,10 @@ public void handleMessage(Message message) throws Fault {
createExchangeId(message);
final OutputStream os = message.getContent(OutputStream.class);
if (os != null) {
LoggingCallback callback = new LoggingCallback(sender, message, os, limit);
// Wrap the callback to ensure logging only once and to avoid memory leaks (CXF-9251)
OneTimeLoggingCallback callback = new OneTimeLoggingCallback(
new LoggingCallback(sender, message, os, limit)
);
message.setContent(OutputStream.class, createCachingOut(message, os, callback));
} else {
final Writer iowriter = message.getContent(Writer.class);
Expand Down Expand Up @@ -179,6 +183,53 @@ private void writePayload(StringBuilder builder, StringWriter stringWriter, LogE
}
}


/***
* [CXF-9251]
* If CachedOutputStream cos has a tmp file (so 'threshold' was triggered), after the introduction of
* DelayedCachedOutputStreamCleaner, a reference of LoggingOutputStream
* is held in a queue list ( DelayQueue<DelayedCloseable> queue )
* If something goes wrong while closing the LoggingOutputStream, onClose() can be recall
* and log twice (or more) when trying to delete the orphan tmp file.
* Furthermore, the LoggingOutputStream that registered LoggingCallback holds a reference to that
* Object (and its attributes... like Message instance), so it remains in memory avoiding GC until the
* DelayedCachedOutputStreamCleaner does his job (default 30 minutes)
* -----------
* Solution:
* This class ensures we log only once, and then dereferences the original LoggingCallback, making it
* eligible for garbage collection.
*/
public class OneTimeLoggingCallback implements CachedOutputStreamCallback {
private LoggingCallback wrappedCallback;
private final AtomicBoolean alreadyClosed = new AtomicBoolean(false);

public OneTimeLoggingCallback(LoggingCallback wrappedCallback) {
this.wrappedCallback = wrappedCallback;
}

@Override
public void onClose(CachedOutputStream cos) {
// Ensure to log only once
if (alreadyClosed.compareAndSet(false, true) && wrappedCallback != null) {
try {
wrappedCallback.onClose(cos);
} finally {
// Make the original Callback (and its attribute) eligible for GC
// This is especially useful if the cos with a reference for this Callback
// is still held in a list (such as the one used by DelayedCachedOutputStreamCleaner)
this.wrappedCallback = null;
}
}
}

@Override
public void onFlush(CachedOutputStream cos) {
if (!alreadyClosed.get() && wrappedCallback != null) {
wrappedCallback.onFlush(cos);
}
}
}

public class LoggingCallback implements CachedOutputStreamCallback {

private final Message message;
Expand Down