Files
document-haness/docs/clean-architecture-backend-template/final/evidence/raw/a11-f006-boundeddatabufferflux.txt
T
DongHyeonkaandClaude Opus 5 b2963105a8 docs(keycloak-session-store): import the session-storage lab as a new project
The keycloak project ended with four open questions that design could not
settle. A two-VM lab was built to answer them by measurement, and this is
that material: 26 experiments, 125 raw command outputs, 22 browser captures.

Follows the import procedure in README.md.

  source/     the originating repository verbatim — 78 documents, 28 SVGs,
              8 manifests, plus .source-revision recording the commit
  final/      the SSOT
    document.md   729 lines written from the 29 experiment documents, not
                  concatenated: what was predicted, what was measured, and
                  where the measurement itself was wrong
    evidence/raw    125 outputs, flattened to <experiment>__<file> because
                    the originals collided (01-baseline.txt appeared three
                    times) and the audit only globs the top level
    evidence/meta   one per raw file; command and exitCode are null and the
                    README says why rather than inventing them
    evidence/browser  22 captures
    assets/       three diagrams through techviz
    .techviz/     their VizSpecs

A separate project rather than an addition to keycloak: the B-layer answers
that project's four questions, but the A, C and D layers are about cluster
failure, SSO and operations, and one document.md should hold one subject.
The four question records there can point here through 관계.

Recorded rather than papered over: only three of the 28 diagrams were
remade. The repository forbids hand-drawn SVG and forbids titles inside the
canvas; all 28 originals carry both, so converting them is redrawing, not
reformatting. They stay in source/ and the gap is written into the document.

verify-pipeline.py passes. audit-records.py reports no issues.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
2026-09-04 22:51:59 +09:00

105 lines
4.3 KiB
Plaintext

# 클래스 javadoc 이 지목한 누수 경로
/**
* Bounds a reactive body and releases every buffer it does not hand on (design §23.2, §23.3).
*
* <p>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<DataBuffer> bound(
20: Flux<DataBuffer> 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();
}