Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
42 commits
Select commit Hold shift + click to select a range
bc97709
test(netty-4.1): Add pipelining test
ygree Jul 13, 2026
7f7fdf4
fix(netty-4.1): pipelined response context tracking
ygree Jul 13, 2026
57f9d65
Disable Netty server tracing on context queue overflow
ygree Jul 14, 2026
6173156
Fix Netty 4.1 response context handling
ygree Jul 16, 2026
7937757
Add Netty chunked response span timing tests
ygree Jul 13, 2026
5be33eb
Fix Netty 4.1 server span completion for streaming responses
ygree Jul 14, 2026
aaaf77d
Fix Netty server span lifecycle for streaming responses
ygree Jul 14, 2026
372e9a8
Fix Netty 4.1 server span edge cases from PR review
ygree Jul 14, 2026
0a6269f
Restore request-ended callbacks for terminal Netty responses
ygree Jul 15, 2026
344cbb4
Resolved: terminal responses no longer call `beforeFinish(...)`
ygree Jul 16, 2026
c2e0b3f
Changed HttpServerResponseTracingHandler.java so unknown-length
ygree Jul 16, 2026
4f7d924
Restrict close-delimited response handling to HTTP/1.x only.
ygree Jul 16, 2026
69ec7b0
The Netty 4.1 websocket upgrade check is now value-based for status
ygree Jul 16, 2026
cde82d0
It now treats known-length responses as complete once the declared body
ygree Jul 16, 2026
3ebe9c3
Fix Netty IAST reporting on async response completion
ygree Jul 17, 2026
6c6c9fc
Preserve Netty HTTP/2 stream contexts
ygree Jul 21, 2026
ba51b12
Fix deferred AppSec block responses for Netty pipelining
ygree Jul 22, 2026
88da30a
Avoid logging dropped Netty outbound messages
ygree Jul 22, 2026
67b731a
Use computed Netty 4.1 response status code
ygree Jul 22, 2026
d9c5d1d
Use single-thread Netty event loop in tests
ygree Jul 22, 2026
762d160
Merge branch 'ygree/fix-netty41-http11-pipelining-context-queue' into…
ygree Jul 22, 2026
e7046aa
Avoid synthetic accessors in Netty blocking response handler
ygree Jul 22, 2026
aa33b42
Fix Netty block function pipelining deferral
ygree Jul 22, 2026
e7897a3
Fix Netty pipelined block response ordering
ygree Jul 22, 2026
66f3ac8
Fix Netty HTTP/2 response AppSec context handling
ygree Jul 22, 2026
e1a31f3
Fix Netty AppSec blocked response handling
ygree Jul 22, 2026
f89affd
Merge branch 'ygree/fix-netty41-http11-pipelining-context-queue' into…
ygree Jul 22, 2026
349a4d0
Merge branch 'master' into ygree/fix-netty41-http11-pipelining-contex…
ygree Jul 22, 2026
1811f52
Merge branch 'ygree/fix-netty41-http11-pipelining-context-queue' into…
ygree Jul 22, 2026
78ebedf
Align request-context completion with Netty’s message framing contract.
ygree Jul 23, 2026
ba58dc7
Reuse Netty blocking handler for deferred pipelined requests
ygree Jul 23, 2026
b8216d3
Fix Netty blocking response scheduling races
ygree Jul 23, 2026
7dd79a7
Reuse Netty request context queues for keep-alive traffic
ygree Jul 23, 2026
b81803e
Reuse Netty request context queues for keep-alive traffic
ygree Jul 23, 2026
0c82171
Merge branch 'ygree/fix-netty41-http11-pipelining-context-queue' into…
ygree Jul 23, 2026
ffb560d
Align Netty response completion with LastHttpContent
ygree Jul 23, 2026
31ca152
Merge branch 'master' into ygree/fix-netty41-http11-pipelining-contex…
ygree Jul 24, 2026
952e59b
Merge branch 'ygree/fix-netty41-http11-pipelining-context-queue' into…
ygree Jul 24, 2026
8a70d5e
Raise agent jar size budget
ygree Jul 24, 2026
3eff294
Merge branch 'master' into ygree/fix-netty41-chunked-tracing
ygree Jul 24, 2026
670d5c4
Wait for Netty server spans before closing test sockets
ygree Jul 24, 2026
a5af354
Merge branch 'master' into ygree/fix-netty41-chunked-tracing
ygree Jul 27, 2026
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 @@ -15,6 +15,7 @@
import io.netty.channel.ChannelInboundHandlerAdapter;
import io.netty.handler.codec.http.HttpHeaders;
import io.netty.handler.codec.http.HttpRequest;
import java.util.Deque;

@ChannelHandler.Sharable
public class HttpServerRequestTracingHandler extends ChannelInboundHandlerAdapter {
Expand Down Expand Up @@ -95,9 +96,43 @@ public void channelInactive(ChannelHandlerContext ctx) throws Exception {
super.channelInactive(ctx);
} finally {
try {
ServerRequestContext.closeAll(ctx.channel());
final Deque<ServerRequestContext> storedContexts =
ServerRequestContext.removeAll(ctx.channel());
if (storedContexts != null) {
ServerRequestContext storedContext;
while ((storedContext = storedContexts.pollFirst()) != null) {
if (storedContext.isResponseStarted()) {
finishSpanOnChannelClose(storedContext);
} else {
publishSpanOnChannelClose(storedContext.tracingContext());
}
}
}
} catch (final Throwable ignored) {
}
}
}

private static void finishSpanOnChannelClose(final ServerRequestContext serverContext) {
final Context storedContext = serverContext.tracingContext();
final AgentSpan span = AgentSpan.fromContext(storedContext);
if (span == null) {
return;
}
try (final ContextScope ignored = storedContext.attach()) {
if (!serverContext.isBeforeFinishCalled()) {
serverContext.markBeforeFinishCalled();
DECORATE.beforeFinish(storedContext);
}
span.finish();
}
}

private static void publishSpanOnChannelClose(final Context storedContext) {
final AgentSpan span = AgentSpan.fromContext(storedContext);
if (span != null && span.phasedFinish()) {
// At this point we can just publish this span to avoid losing the rest of the trace.
span.publish();
}
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -14,9 +14,11 @@
import io.netty.channel.ChannelOutboundHandlerAdapter;
import io.netty.channel.ChannelPromise;
import io.netty.handler.codec.http.HttpHeaderNames;
import io.netty.handler.codec.http.HttpHeaderValues;
import io.netty.handler.codec.http.HttpResponse;
import io.netty.handler.codec.http.HttpResponseStatus;
import io.netty.handler.codec.http.LastHttpContent;
import io.netty.util.concurrent.Future;

@ChannelHandler.Sharable
public class HttpServerResponseTracingHandler extends ChannelOutboundHandlerAdapter {
Expand All @@ -25,7 +27,8 @@ public class HttpServerResponseTracingHandler extends ChannelOutboundHandlerAdap
@Override
public void write(final ChannelHandlerContext ctx, final Object msg, final ChannelPromise prm) {
final boolean isResponse = msg instanceof HttpResponse;
if (!isResponse && !(msg instanceof LastHttpContent)) {
final boolean isLastContent = msg instanceof LastHttpContent;
if (!isResponse && !isLastContent) {
ctx.write(msg, prm);
return;
}
Expand All @@ -50,43 +53,99 @@ public void write(final ChannelHandlerContext ctx, final Object msg, final Chann
return;
}

try (final ContextScope scope = storedContext.attach()) {
try (final ContextScope ignored = storedContext.attach()) {
final HttpResponse response = isResponse ? (HttpResponse) msg : null;
final boolean responseComplete = msg instanceof LastHttpContent;

final boolean websocketUpgrade = response != null && isWebsocketUpgrade(response);
final boolean informationalResponse =
response != null && isInformationalResponse(response) && !websocketUpgrade;
final boolean finishResponseOnWrite = isLastContent && !informationalResponse;
final ChannelPromise writePromise =
finishResponseOnWrite && prm.isVoid() ? ctx.newPromise() : prm;
try {
ctx.write(msg, prm);
} catch (final Throwable throwable) {
DECORATE.onError(span, throwable);
span.setHttpStatusCode(500);
span.finish(); // Finish the span manually since finishSpanOnClose was false
removeServerContext(ctx, serverContext);
throw throwable;
}
if (response != null) {
final boolean isWebsocketUpgrade =
response.status() == HttpResponseStatus.SWITCHING_PROTOCOLS
&& "websocket".equals(response.headers().get(HttpHeaderNames.UPGRADE));
if (isWebsocketUpgrade) {
ctx.channel()
.attr(WEBSOCKET_SENDER_HANDLER_CONTEXT)
.set(new HandlerContext.Sender(span, ctx.channel().id().asShortText()));
if (response != null && !informationalResponse) {
onResponse(ctx, span, serverContext, response, websocketUpgrade);
}
if (isInformational(response) && !isWebsocketUpgrade) {
return;
if (finishResponseOnWrite) {
removeServerContext(ctx, serverContext);
writePromise.addListener(
future -> finishSpan(serverContext, storedContext, span, future));
}
if (serverContext != null) {
serverContext.markResponseStarted();
ctx.write(msg, writePromise);
if (finishResponseOnWrite && (!writePromise.isDone() || writePromise.isSuccess())) {
final ServerRequestContext nextResponse =
ServerRequestContext.nextResponse(ctx.channel());
BlockingResponseHandler.maybeWriteDeferredBlockResponse(ctx, nextResponse);
}
DECORATE.onResponse(span, response);
} catch (final Throwable throwable) {
if (!finishResponseOnWrite || !writePromise.isDone()) {
DECORATE.onError(span, throwable);
span.setHttpStatusCode(500);
if (!finishResponseOnWrite) {
removeServerContext(ctx, serverContext);
}
finishSpan(serverContext, storedContext, span);
}
throw throwable;
}
if (responseComplete) {
DECORATE.beforeFinish(scope.context());
span.finish(); // Finish the span manually since finishSpanOnClose was false
removeServerContext(ctx, serverContext);
final ServerRequestContext nextResponse = ServerRequestContext.nextResponse(ctx.channel());
BlockingResponseHandler.maybeWriteDeferredBlockResponse(ctx, nextResponse);
}
}

private static void onResponse(
final ChannelHandlerContext ctx,
final AgentSpan span,
final ServerRequestContext serverContext,
final HttpResponse response,
final boolean websocketUpgrade) {
if (websocketUpgrade) {
ctx.channel()
.attr(WEBSOCKET_SENDER_HANDLER_CONTEXT)
.set(new HandlerContext.Sender(span, ctx.channel().id().asShortText()));
}
DECORATE.onResponse(span, response);
if (serverContext != null) {
serverContext.markResponseStarted();
}
}

private static boolean isInformationalResponse(final HttpResponse response) {
final int statusCode = response.status().code();
return statusCode >= 100 && statusCode < 200;
}

private static boolean isWebsocketUpgrade(final HttpResponse response) {
return response.status().code() == HttpResponseStatus.SWITCHING_PROTOCOLS.code()
&& response
.headers()
.containsValue(HttpHeaderNames.UPGRADE, HttpHeaderValues.WEBSOCKET, true);
Comment thread
ygree marked this conversation as resolved.
}

private static void finishSpan(
final ServerRequestContext serverContext,
final Context storedContext,
final AgentSpan span,
final Future<?> future) {
if (!future.isSuccess()) {
DECORATE.onError(span, future.cause());
span.setHttpStatusCode(500);
}
finishSpan(serverContext, storedContext, span);
}

private static void finishSpan(
final ServerRequestContext serverContext, final Context storedContext, final AgentSpan span) {
try (final ContextScope ignored = storedContext.attach()) {
beforeFinish(serverContext, storedContext);
span.finish(); // Finish the span manually since finishSpanOnClose was false
}
}

private static void beforeFinish(
final ServerRequestContext serverContext, final Context storedContext) {
if (serverContext == null || !serverContext.isBeforeFinishCalled()) {
if (serverContext != null) {
serverContext.markBeforeFinishCalled();
}
DECORATE.beforeFinish(storedContext);
}
}

Expand All @@ -98,9 +157,4 @@ private static void removeServerContext(
ServerRequestContext.remove(ctx.channel(), serverContext);
}
}

private static boolean isInformational(final HttpResponse response) {
final int statusCode = response.status().code();
return statusCode >= 100 && statusCode < 200;
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -70,6 +70,10 @@ public void write(ChannelHandlerContext ctx, Object msg, ChannelPromise prm) thr
}
HttpResponse origResponse = (HttpResponse) msg;
int statusCode = origResponse.status().code();
// Interim 1xx responses (e.g. 100 Continue, 103 Early Hints) precede the final response.
// Analyzing one here would consume the one-shot response analysis before the final response is
// written, so its status and headers would never be inspected. Switching Protocols (101) is
// terminal, so it is still analyzed.
if (statusCode >= 100
&& statusCode < 200
&& statusCode != HttpResponseStatus.SWITCHING_PROTOCOLS.code()) {
Comment thread
ygree marked this conversation as resolved.
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -13,6 +13,7 @@
import static io.netty.handler.codec.http.HttpResponseStatus.OK;
import static io.netty.handler.codec.http.HttpResponseStatus.SWITCHING_PROTOCOLS;
import static io.netty.handler.codec.http.HttpVersion.HTTP_1_1;
import static org.junit.jupiter.api.Assertions.assertDoesNotThrow;
import static org.junit.jupiter.api.Assertions.assertNotNull;
import static org.junit.jupiter.api.Assertions.assertNull;
import static org.junit.jupiter.api.Assertions.assertSame;
Expand All @@ -39,16 +40,22 @@
class HttpServerResponseTracingHandlerTest extends AbstractInstrumentationTest {

@Test
void finishesMirroredContextWhenRequestQueueIsAbsent() {
void finishesMirroredContextOnLastContentWhenRequestQueueIsAbsent() {
EmbeddedChannel channel = new EmbeddedChannel(HttpServerResponseTracingHandler.INSTANCE);
AgentSpan span = startSpan("netty", "mirrored-http2-server");
channel.attr(CONTEXT_ATTRIBUTE_KEY).set(span);

assertTrue(channel.writeOutbound(new DefaultFullHttpResponse(HTTP_1_1, OK)));
HttpResponse response = new DefaultHttpResponse(HTTP_1_1, OK);
assertTrue(channel.writeOutbound(response));

FullHttpResponse response = channel.readOutbound();
assertNotNull(response);
response.release();
assertSame(span, channel.attr(CONTEXT_ATTRIBUTE_KEY).get());
HttpResponse forwarded = channel.readOutbound();
assertSame(response, forwarded);
ReferenceCountUtil.release(forwarded);

assertTrue(channel.writeOutbound(LastHttpContent.EMPTY_LAST_CONTENT));

ReferenceCountUtil.release(channel.readOutbound());
assertNull(channel.attr(CONTEXT_ATTRIBUTE_KEY).get());
channel.finishAndReleaseAll();
assertTraces(trace(span().root().operationName("mirrored-http2-server")));
Expand All @@ -70,7 +77,7 @@ void forwardsLastContentBeforeFinalResponseWithoutCompletingContext() {
}

@Test
void headerOnlyResponseDoesNotCompleteContextBeforeLastContent() {
void headerOnlyResponseWaitsForLastContent() {
EmbeddedChannel channel = new EmbeddedChannel(HttpServerResponseTracingHandler.INSTANCE);
AgentSpan span = startSpan("netty", "header-only-server");
ServerRequestContext serverContext = ServerRequestContext.add(channel, span, null);
Expand All @@ -84,6 +91,7 @@ void headerOnlyResponseDoesNotCompleteContextBeforeLastContent() {
HttpResponse forwarded = channel.readOutbound();
assertSame(response, forwarded);
ReferenceCountUtil.release(forwarded);
assertNull(channel.readOutbound());

assertTrue(channel.writeOutbound(LastHttpContent.EMPTY_LAST_CONTENT));

Expand All @@ -95,10 +103,12 @@ void headerOnlyResponseDoesNotCompleteContextBeforeLastContent() {
}

@Test
void rawFixedLengthBodyDoesNotCompleteContextBeforeLastContent() {
void rawFixedLengthBodyWaitsForLastContent() {
EmbeddedChannel channel = new EmbeddedChannel(HttpServerResponseTracingHandler.INSTANCE);
AgentSpan span = startSpan("netty", "raw-body-server");
ServerRequestContext serverContext = ServerRequestContext.add(channel, span, null);
ServerRequestContext nextServerContext =
ServerRequestContext.add(channel, Context.root(), null);
HttpResponse response = new DefaultHttpResponse(HTTP_1_1, OK);
response.headers().set(CONTENT_LENGTH, 4);

Expand All @@ -108,15 +118,34 @@ void rawFixedLengthBodyDoesNotCompleteContextBeforeLastContent() {

assertSame(serverContext, ServerRequestContext.nextResponse(channel));
ReferenceCountUtil.release(channel.readOutbound());
assertNull(channel.readOutbound());

assertTrue(channel.writeOutbound(LastHttpContent.EMPTY_LAST_CONTENT));

assertNull(ServerRequestContext.nextResponse(channel));
assertSame(nextServerContext, ServerRequestContext.nextResponse(channel));
ReferenceCountUtil.release(channel.readOutbound());
ServerRequestContext.remove(channel, nextServerContext);
channel.finishAndReleaseAll();
assertTraces(trace(span().root().operationName("raw-body-server")));
}

@Test
void doesNotThrowOnMalformedContentLength() {
EmbeddedChannel channel = new EmbeddedChannel(HttpServerResponseTracingHandler.INSTANCE);
AgentSpan span = startSpan("netty", "malformed-content-length-server");
ServerRequestContext.add(channel, span, null);
FullHttpResponse response = new DefaultFullHttpResponse(HTTP_1_1, OK);
response.headers().set(CONTENT_LENGTH, "malformed");

assertDoesNotThrow(() -> assertTrue(channel.writeOutbound(response)));

FullHttpResponse forwarded = channel.readOutbound();
assertSame(response, forwarded);
forwarded.release();
channel.finishAndReleaseAll();
assertTraces(trace(span().root().operationName("malformed-content-length-server")));
}

@Test
void nettyEncoderRequiresLastContentBeforeNextKeepAliveResponse() {
EmbeddedChannel channel = new EmbeddedChannel(new HttpResponseEncoder());
Expand Down Expand Up @@ -159,4 +188,27 @@ void fullWebsocketUpgradeCompletesContextAndPreservesHandshakeSpan() {
channel.finishAndReleaseAll();
assertTraces(trace(span().root().operationName("websocket-handshake-server")));
}

@Test
void fullNonWebsocketUpgradeWaitsForFinalResponse() {
EmbeddedChannel channel = new EmbeddedChannel(HttpServerResponseTracingHandler.INSTANCE);
AgentSpan span = startSpan("netty", "h2c-upgrade-server");
ServerRequestContext serverContext = ServerRequestContext.add(channel, span, null);
FullHttpResponse upgradeResponse = new DefaultFullHttpResponse(HTTP_1_1, SWITCHING_PROTOCOLS);
upgradeResponse.headers().set(UPGRADE, "h2c");

assertTrue(channel.writeOutbound(upgradeResponse));

assertSame(serverContext, ServerRequestContext.nextResponse(channel));
assertNull(channel.attr(WEBSOCKET_SENDER_HANDLER_CONTEXT).get());
ReferenceCountUtil.release(channel.readOutbound());

assertTrue(channel.writeOutbound(new DefaultFullHttpResponse(HTTP_1_1, OK)));

assertNull(ServerRequestContext.nextResponse(channel));
ReferenceCountUtil.release(channel.readOutbound());
channel.finishAndReleaseAll();

assertTraces(trace(span().root().operationName("h2c-upgrade-server")));
}
}
Loading
Loading