Skip to content
Draft
Show file tree
Hide file tree
Changes from all commits
Commits
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 @@ -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 -> {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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()
Expand Down
Loading