From c11ba9b4ba19f332d9665d64dce7bbf83ce71d7d Mon Sep 17 00:00:00 2001 From: Dmitry Kropachev Date: Tue, 18 Aug 2026 10:33:16 -0400 Subject: [PATCH] fix: release stream ids on request setup failure --- .../core/channel/InFlightHandler.java | 31 +++++--- .../core/channel/InFlightHandlerTest.java | 70 +++++++++++++++++++ 2 files changed, 91 insertions(+), 10 deletions(-) diff --git a/core/src/main/java/com/datastax/oss/driver/internal/core/channel/InFlightHandler.java b/core/src/main/java/com/datastax/oss/driver/internal/core/channel/InFlightHandler.java index c758275028c..28850551ec7 100644 --- a/core/src/main/java/com/datastax/oss/driver/internal/core/channel/InFlightHandler.java +++ b/core/src/main/java/com/datastax/oss/driver/internal/core/channel/InFlightHandler.java @@ -148,16 +148,27 @@ private void write(ChannelHandlerContext ctx, RequestMessage message, ChannelPro return; } - LOG.trace("[{}] Writing {} on stream id {}", logPrefix, message.responseCallback, streamId); - Frame frame = - Frame.forRequest( - protocolVersion.getCode(), - streamId, - message.tracing, - message.customPayload, - message.request); - - inFlight.put(streamId, message.responseCallback); + Frame frame; + boolean registered = false; + try { + LOG.trace("[{}] Writing {} on stream id {}", logPrefix, message.responseCallback, streamId); + frame = + Frame.forRequest( + protocolVersion.getCode(), + streamId, + message.tracing, + message.customPayload, + message.request); + + inFlight.put(streamId, message.responseCallback); + registered = true; + } finally { + // acquire() consumed the caller's reservation. Until the callback is registered, no other + // path owns the concrete id, so every synchronous setup failure must release it here. + if (!registered) { + streamIds.release(streamId); + } + } ChannelFuture writeFuture = ctx.write(frame, promise); writeFuture.addListener( future -> { diff --git a/core/src/test/java/com/datastax/oss/driver/internal/core/channel/InFlightHandlerTest.java b/core/src/test/java/com/datastax/oss/driver/internal/core/channel/InFlightHandlerTest.java index 673748b9029..b410ad2a798 100644 --- a/core/src/test/java/com/datastax/oss/driver/internal/core/channel/InFlightHandlerTest.java +++ b/core/src/test/java/com/datastax/oss/driver/internal/core/channel/InFlightHandlerTest.java @@ -39,6 +39,7 @@ import io.netty.channel.ChannelFuture; import io.netty.channel.ChannelPromise; import java.net.InetSocketAddress; +import java.nio.ByteBuffer; import java.util.ArrayList; import java.util.Collections; import java.util.List; @@ -123,6 +124,61 @@ public void should_assign_streamid_and_send_frame() { assertThat(frame.message).isEqualTo(QUERY); } + @Test + public void should_release_stream_id_when_frame_creation_fails() { + // Given + StreamIdGenerator realStreamIds = new StreamIdGenerator(2); + InFlightHandler handler = addToPipeline(realStreamIds); + assertThat(handler.preAcquireId()).isTrue(); + + // Protocol V3 does not support custom payloads, so frame creation fails after acquire(). + DriverChannel.RequestMessage message = + new DriverChannel.RequestMessage( + QUERY, + false, + Collections.singletonMap("test", ByteBuffer.allocate(0)), + new MockResponseCallback(), + handler); + + // When + ChannelFuture writeFuture = channel.writeAndFlush(message); + + // Then + assertThat(writeFuture).isFailed(); + assertThat(handler.getAvailableIds()).isEqualTo(2); + assertNoOutboundFrame(); + } + + @Test + public void should_release_stream_id_when_callback_registration_fails() { + // Given + StreamIdGenerator realStreamIds = new StreamIdGenerator(2); + InFlightHandler handler = addToPipeline(realStreamIds); + MockResponseCallback responseCallback = new MockResponseCallback(); + assertThat(handler.preAcquireId()).isTrue(); + assertThat( + channel.writeAndFlush( + new DriverChannel.RequestMessage( + QUERY, false, Frame.NO_PAYLOAD, responseCallback, handler))) + .isSuccess(); + readOutboundFrame(); + assertThat(handler.getAvailableIds()).isEqualTo(1); + + // Reusing an in-flight callback is rejected by the callback map after acquiring a second id. + assertThat(handler.preAcquireId()).isTrue(); + + // When + ChannelFuture writeFuture = + channel.writeAndFlush( + new DriverChannel.RequestMessage( + QUERY, false, Frame.NO_PAYLOAD, responseCallback, handler)); + + // Then + assertThat(writeFuture).isFailed(); + assertThat(handler.getAvailableIds()).isEqualTo(1); + assertNoOutboundFrame(); + } + @Test public void should_notify_callback_of_response() { // Given @@ -665,6 +721,20 @@ private void addToPipeline() { addToPipelineWithEventCallback(null); } + private InFlightHandler addToPipeline(StreamIdGenerator streamIds) { + InFlightHandler handler = + new InFlightHandler( + DefaultProtocolVersion.V3, + streamIds, + MAX_ORPHAN_IDS, + SET_KEYSPACE_TIMEOUT_MILLIS, + channel.newPromise(), + null, + "test"); + channel.pipeline().addLast(handler); + return handler; + } + private void addToPipelineWithEventCallback(EventCallback eventCallback) { channel .pipeline()