# 클래스 javadoc 이 지목한 누수 경로
/**
* Bounds a reactive body and releases every buffer it does not hand on (design §23.2, §23.3).
*
*
Cancellation and error are the paths that leak in practice: the subscriber stops asking, the
* upstream drops what it already produced, and those buffers are direct memory nobody returns.
*/
# 조립되는 연산자 사슬 전체
19: public static Flux bound(
20: Flux source, ResponseSizeLimiter limiter, FirstByteDeliveryGuard guard) {
21: Objects.requireNonNull(source, "source");
22: Objects.requireNonNull(limiter, "response size limiter");
23: Objects.requireNonNull(guard, "first byte guard");
24: return source
25: .doOnNext(
26: buffer -> {
27: limiter.recordWireBytes(buffer.readableByteCount());
28: guard.markDelivered();
29: })
30: .doOnDiscard(DataBuffer.class, DataBufferUtils::release)
31: .doOnCancel(() -> {})
32: .onErrorResume(
33: failure -> {
34: // Buffers already emitted belong to the subscriber; anything still in flight is
35: // discarded through doOnDiscard above.
36: return Flux.error(failure);
37: });
38: }
# 이 사슬을 부르는 유일한 지점과 그 바로 뒤
webclient/ReactiveStreamingGateway.java:102: return BoundedDataBufferFlux.bound(
87: return request
88: .exchangeToFlux(
89: response -> {
90: int status = response.statusCode().value();
91: if (status < 200 || status >= 300) {
92: return response
93: .releaseBody()
94: .thenMany(
95: Flux.error(
96: new HttpRemoteErrorException(
97: "streaming download returned an error status",
98: metadata
99: .withStatus(new HttpStatus(status))
100: .withEvidence(ExecutionEvidence.RESPONSE_RECEIVED))));
101: }
102: return BoundedDataBufferFlux.bound(
103: response.bodyToFlux(DataBuffer.class), limiter, guard);
104: })
105: .doOnDiscard(DataBuffer.class, DataBufferUtils::release);
106: }
# 취소 경로 test 와 그 클래스에 붙은 확장
@ExtendWith(NettyLeakDetectionExtension.class)
class ReactiveStreamingLifecycleTest {
@Test
void cancellationReleasesTheConnectionForTheNextCall() throws Exception {
try (MockHttpServer server = MockHttpServer.start()) {
ClientProfile profile = singleConnectionProfile(server.uri("/"));
try (ReactiveTestGateways.Harness harness = ReactiveTestGateways.forProfile(profile)) {
server.enqueueBody(200, "application/octet-stream", new byte[64 * 1024]);
server.enqueueBody(200, "application/octet-stream", new byte[16]);
ReactiveStreamingGateway gateway = new ReactiveStreamingGateway(harness.registry());
StepVerifier.create(
gateway
.download(
profile.name(),
HttpOperation.get(new OperationName("download"), "/large", Map.of()))
.doOnNext(DataBufferUtils::release)
.take(1))
.expectNextCount(1)
.verifyComplete();
StepVerifier.create(
gateway
.download(
profile.name(),
HttpOperation.get(new OperationName("download"), "/small", Map.of()))
.doOnNext(DataBufferUtils::release))
.expectNextCount(1)
.verifyComplete();
}
}
}
# 그 확장이 단언하는 것
assertThat(ResourceLeakDetector.getLevel())
.describedAs(
"netty leak detection must run at PARANOID; set -Dio.netty.leakDetection.level=paranoid")
.isEqualTo(ResourceLeakDetector.Level.PARANOID);
appender = new LeakRecordingAppender(leakRecords);
appender.attach(LEAK_LOGGER);
}
@Override
public void afterAll(ExtensionContext context) {
if (appender != null) {
appender.detach(LEAK_LOGGER);
}
assertThat(leakRecords).describedAs("netty reported buffer leak(s)").isEmpty();
}