Skip to content

Commit 492f1a5

Browse files
nebhaleyschimke
authored andcommitted
Fix Reference Count Leak (#488)
Previously, there were a couple of frames that converted Strings to ByteBufs automatically as part of construction. Those frames though, retained any ByteBufs passed into them so that they could be released by the creators as necessary. Since the automatic construction was not doing this release on the String-based ByteBufs, there was a minor leak of references. This change updates those locations to release the String-based ByteBufs and accurately account for their retention and release.
1 parent ff4647b commit 492f1a5

17 files changed

Lines changed: 161 additions & 71 deletions

rsocket-core/src/main/java/io/rsocket/fragmentation/FrameFragmenter.java

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -17,7 +17,7 @@
1717
package io.rsocket.fragmentation;
1818

1919
import static io.rsocket.framing.PayloadFrame.createPayloadFrame;
20-
import static io.rsocket.util.DisposableUtil.disposeQuietly;
20+
import static io.rsocket.util.DisposableUtils.disposeQuietly;
2121
import static java.lang.Math.min;
2222

2323
import io.netty.buffer.ByteBuf;

rsocket-core/src/main/java/io/rsocket/fragmentation/FrameReassembler.java

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -16,7 +16,7 @@
1616

1717
package io.rsocket.fragmentation;
1818

19-
import static io.rsocket.util.DisposableUtil.disposeQuietly;
19+
import static io.rsocket.util.DisposableUtils.disposeQuietly;
2020
import static io.rsocket.util.RecyclerFactory.createRecycler;
2121

2222
import io.netty.buffer.ByteBuf;

rsocket-core/src/main/java/io/rsocket/framing/AbstractRecyclableFrame.java

Lines changed: 3 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -16,12 +16,14 @@
1616

1717
package io.rsocket.framing;
1818

19+
import static io.netty.util.ReferenceCountUtil.release;
1920
import static java.nio.charset.StandardCharsets.UTF_8;
2021

2122
import io.netty.buffer.ByteBuf;
2223
import io.netty.buffer.ByteBufAllocator;
2324
import io.netty.buffer.Unpooled;
2425
import io.netty.util.Recycler.Handle;
26+
import io.netty.util.ReferenceCounted;
2527
import java.util.Objects;
2628
import reactor.util.annotation.Nullable;
2729

@@ -51,7 +53,7 @@ abstract class AbstractRecyclableFrame<SELF extends AbstractRecyclableFrame<SELF
5153
@SuppressWarnings("unchecked")
5254
public final void dispose() {
5355
if (byteBuf != null) {
54-
byteBuf.release();
56+
release(byteBuf);
5557
}
5658

5759
byteBuf = null;

rsocket-core/src/main/java/io/rsocket/framing/DataFrame.java

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -51,6 +51,7 @@ default int getDataLength() {
5151
* it.
5252
*
5353
* @return the data directly
54+
* @see #getDataAsUtf8()
5455
* @see #mapData(Function)
5556
*/
5657
ByteBuf getUnsafeData();

rsocket-core/src/main/java/io/rsocket/framing/ErrorFrame.java

Lines changed: 8 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -16,6 +16,7 @@
1616

1717
package io.rsocket.framing;
1818

19+
import static io.netty.util.ReferenceCountUtil.release;
1920
import static io.rsocket.framing.FrameType.ERROR;
2021
import static io.rsocket.util.RecyclerFactory.createRecycler;
2122

@@ -70,7 +71,13 @@ public static ErrorFrame createErrorFrame(ByteBuf byteBuf) {
7071
public static ErrorFrame createErrorFrame(
7172
ByteBufAllocator byteBufAllocator, int errorCode, @Nullable String data) {
7273

73-
return createErrorFrame(byteBufAllocator, errorCode, getUtf8AsByteBuf(data));
74+
ByteBuf dataByteBuf = getUtf8AsByteBuf(data);
75+
76+
try {
77+
return createErrorFrame(byteBufAllocator, errorCode, dataByteBuf);
78+
} finally {
79+
release(dataByteBuf);
80+
}
7481
}
7582

7683
/**

rsocket-core/src/main/java/io/rsocket/framing/ExtensionFrame.java

Lines changed: 11 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -16,6 +16,7 @@
1616

1717
package io.rsocket.framing;
1818

19+
import static io.netty.util.ReferenceCountUtil.release;
1920
import static io.rsocket.framing.FrameType.EXT;
2021
import static io.rsocket.util.RecyclerFactory.createRecycler;
2122

@@ -79,8 +80,16 @@ public static ExtensionFrame createExtensionFrame(
7980
@Nullable String metadata,
8081
@Nullable String data) {
8182

82-
return createExtensionFrame(
83-
byteBufAllocator, ignore, extendedType, getUtf8AsByteBuf(metadata), getUtf8AsByteBuf(data));
83+
ByteBuf metadataByteBuf = getUtf8AsByteBuf(metadata);
84+
ByteBuf dataByteBuf = getUtf8AsByteBuf(data);
85+
86+
try {
87+
return createExtensionFrame(
88+
byteBufAllocator, ignore, extendedType, metadataByteBuf, dataByteBuf);
89+
} finally {
90+
release(metadataByteBuf);
91+
release(dataByteBuf);
92+
}
8493
}
8594

8695
/**

rsocket-core/src/main/java/io/rsocket/framing/MetadataFrame.java

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -55,6 +55,7 @@ default Optional<Integer> getMetadataLength() {
5555
* if you store it.
5656
*
5757
* @return the metadata directly, or {@code null} if the Metadata flag is not set
58+
* @see #getMetadataAsUtf8()
5859
* @see #mapMetadata(Function)
5960
*/
6061
@Nullable

rsocket-core/src/main/java/io/rsocket/framing/MetadataPushFrame.java

Lines changed: 8 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -16,6 +16,7 @@
1616

1717
package io.rsocket.framing;
1818

19+
import static io.netty.util.ReferenceCountUtil.release;
1920
import static io.rsocket.framing.FrameType.METADATA_PUSH;
2021
import static io.rsocket.util.RecyclerFactory.createRecycler;
2122

@@ -69,8 +70,13 @@ public static MetadataPushFrame createMetadataPushFrame(ByteBuf byteBuf) {
6970
public static MetadataPushFrame createMetadataPushFrame(
7071
ByteBufAllocator byteBufAllocator, String metadata) {
7172

72-
return createMetadataPushFrame(
73-
byteBufAllocator, getUtf8AsByteBufRequired(metadata, "metadata must not be null"));
73+
ByteBuf metadataByteBuf = getUtf8AsByteBufRequired(metadata, "metadata must not be null");
74+
75+
try {
76+
return createMetadataPushFrame(byteBufAllocator, metadataByteBuf);
77+
} finally {
78+
release(metadataByteBuf);
79+
}
7480
}
7581

7682
/**

rsocket-core/src/main/java/io/rsocket/framing/PayloadFrame.java

Lines changed: 10 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -16,6 +16,7 @@
1616

1717
package io.rsocket.framing;
1818

19+
import static io.netty.util.ReferenceCountUtil.release;
1920
import static io.rsocket.util.RecyclerFactory.createRecycler;
2021

2122
import io.netty.buffer.ByteBuf;
@@ -78,8 +79,15 @@ public static PayloadFrame createPayloadFrame(
7879
@Nullable String metadata,
7980
@Nullable String data) {
8081

81-
return createPayloadFrame(
82-
byteBufAllocator, follows, complete, getUtf8AsByteBuf(metadata), getUtf8AsByteBuf(data));
82+
ByteBuf metadataByteBuf = getUtf8AsByteBuf(metadata);
83+
ByteBuf dataByteBuf = getUtf8AsByteBuf(data);
84+
85+
try {
86+
return createPayloadFrame(byteBufAllocator, follows, complete, metadataByteBuf, dataByteBuf);
87+
} finally {
88+
release(metadataByteBuf);
89+
release(dataByteBuf);
90+
}
8391
}
8492

8593
/**

rsocket-core/src/main/java/io/rsocket/framing/RequestChannelFrame.java

Lines changed: 11 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -16,6 +16,7 @@
1616

1717
package io.rsocket.framing;
1818

19+
import static io.netty.util.ReferenceCountUtil.release;
1920
import static io.rsocket.util.RecyclerFactory.createRecycler;
2021

2122
import io.netty.buffer.ByteBuf;
@@ -83,13 +84,16 @@ public static RequestChannelFrame createRequestChannelFrame(
8384
@Nullable String metadata,
8485
@Nullable String data) {
8586

86-
return createRequestChannelFrame(
87-
byteBufAllocator,
88-
follows,
89-
complete,
90-
initialRequestN,
91-
getUtf8AsByteBuf(metadata),
92-
getUtf8AsByteBuf(data));
87+
ByteBuf metadataByteBuf = getUtf8AsByteBuf(metadata);
88+
ByteBuf dataByteBuf = getUtf8AsByteBuf(data);
89+
90+
try {
91+
return createRequestChannelFrame(
92+
byteBufAllocator, follows, complete, initialRequestN, metadataByteBuf, dataByteBuf);
93+
} finally {
94+
release(metadataByteBuf);
95+
release(dataByteBuf);
96+
}
9397
}
9498

9599
/**

0 commit comments

Comments
 (0)