fixup! handle stop daemon cql command properly

This commit is contained in:
Maxim Muzafarov 2026-07-27 18:06:57 +02:00
parent 8d04517dc1
commit a3e7cccf6a
No known key found for this signature in database
GPG Key ID: 7FEC714D84388C16
3 changed files with 31 additions and 29 deletions

View File

@ -412,7 +412,6 @@ public class SimpleClient implements Closeable
}
}
<<<<<<< HEAD
/**
* The stream id to frame an outbound client request with. SimpleClient carries the intended id on the
* request's (dummy) source envelope (see {@link #execute(List)} and callers that pipeline requests).
@ -423,7 +422,8 @@ public class SimpleClient implements Closeable
{
Envelope source = message.getSource();
return source == null ? 0 : source.header.streamId;
=======
}
private long responseDeadlineNanos()
{
return requestTimeoutSeconds > 0 ? nanoTime() + TimeUnit.SECONDS.toNanos(requestTimeoutSeconds) : Long.MAX_VALUE;
@ -442,7 +442,6 @@ public class SimpleClient implements Closeable
if (nanoTime() - deadlineNanos >= 0)
return null;
}
>>>>>>> f9202d6b13 (handle stop daemon cql command properly)
}
public interface EventHandler

View File

@ -202,7 +202,7 @@ public class ErrorMessageTest extends EncodeAndDecodeTestBase<ErrorMessage>
Throwable cause = new RuntimeException("Underlying cause");
CommandRequestExecutionException ex = new CommandRequestExecutionException(executionId, errorMessage, cause);
ErrorMessage deserialized = encodeThenDecode(ErrorMessage.fromException(ex), ProtocolVersion.V5);
ErrorMessage deserialized = encodeThenDecode(ErrorMessage.fromExceptionNoStreamId(ex), ProtocolVersion.V5);
assertTrue(deserialized.error instanceof CommandRequestExecutionException);
CommandRequestExecutionException deserializedEx = (CommandRequestExecutionException) deserialized.error;
@ -219,7 +219,7 @@ public class ErrorMessageTest extends EncodeAndDecodeTestBase<ErrorMessage>
Throwable cause = new RuntimeException("Underlying cause");
CommandRequestExecutionException ex = new CommandRequestExecutionException(executionId, errorMessage, cause);
ErrorMessage deserialized = encodeThenDecode(ErrorMessage.fromException(ex), ProtocolVersion.V4);
ErrorMessage deserialized = encodeThenDecode(ErrorMessage.fromExceptionNoStreamId(ex), ProtocolVersion.V4);
assertTrue(deserialized.error instanceof ServerError);
ServerError deserializedEx = (ServerError) deserialized.error;
@ -235,7 +235,7 @@ public class ErrorMessageTest extends EncodeAndDecodeTestBase<ErrorMessage>
String errorMessage = "Command execution failed: test command";
CommandRequestExecutionException ex = new CommandRequestExecutionException(executionId, errorMessage);
ErrorMessage deserialized = encodeThenDecode(ErrorMessage.fromException(ex), ProtocolVersion.V5);
ErrorMessage deserialized = encodeThenDecode(ErrorMessage.fromExceptionNoStreamId(ex), ProtocolVersion.V5);
assertTrue(deserialized.error instanceof CommandRequestExecutionException);
CommandRequestExecutionException deserializedEx = (CommandRequestExecutionException) deserialized.error;
@ -251,7 +251,7 @@ public class ErrorMessageTest extends EncodeAndDecodeTestBase<ErrorMessage>
String errorMessage = "Command execution failed: test command";
CommandRequestExecutionException ex = new CommandRequestExecutionException(executionId, errorMessage);
ErrorMessage deserialized = encodeThenDecode(ErrorMessage.fromException(ex), ProtocolVersion.V3);
ErrorMessage deserialized = encodeThenDecode(ErrorMessage.fromExceptionNoStreamId(ex), ProtocolVersion.V3);
assertTrue(deserialized.error instanceof ServerError);
ServerError deserializedEx = (ServerError) deserialized.error;

View File

@ -400,11 +400,12 @@ public class MessageManagementDispatcherTest
Dispatcher managementDispatcher = new ManagementTestDispatcher(true)
{
@Override
void processRequest(Channel channel,
Message.Request request,
FlushItemConverter forFlusher,
ClientResourceLimits.Overload backpressure,
RequestTime requestTime)
<P> void processRequest(Channel channel,
Message.Request request,
FlushItemConverter<P> forFlusher,
P param,
ClientResourceLimits.Overload backpressure,
RequestTime requestTime)
{
isDoneInsideTask.set(isDone());
entered.countDown();
@ -413,8 +414,8 @@ public class MessageManagementDispatcherTest
};
Message.Request request = createManagementRequest(Message.Type.STARTUP);
managementDispatcher.dispatch(request.connection().channel(), request, (channel, req, response) -> null,
ClientResourceLimits.Overload.NONE);
managementDispatcher.dispatch(request.connection().channel(), request, (param, channel, req, response) -> null,
null, ClientResourceLimits.Overload.NONE);
try
{
assertTrue(entered.await(10, TimeUnit.SECONDS));
@ -464,18 +465,19 @@ public class MessageManagementDispatcherTest
Dispatcher queuedProcessor = new ManagementTestDispatcher(true)
{
@Override
void processRequest(Channel channel,
Message.Request request,
FlushItemConverter forFlusher,
ClientResourceLimits.Overload backpressure,
RequestTime requestTime)
<P> void processRequest(Channel channel,
Message.Request request,
FlushItemConverter<P> forFlusher,
P param,
ClientResourceLimits.Overload backpressure,
RequestTime requestTime)
{
processed.countDown();
}
};
Message.Request request = createManagementRequest(Message.Type.STARTUP);
queuedProcessor.dispatch(request.connection().channel(), request, (channel, req, response) -> null,
ClientResourceLimits.Overload.NONE);
queuedProcessor.dispatch(request.connection().channel(), request, (param, channel, req, response) -> null,
null, ClientResourceLimits.Overload.NONE);
awaitTrue("A queued management request should trigger management backpressure",
() -> !managementDispatcher.hasQueueCapacity());
@ -518,8 +520,8 @@ public class MessageManagementDispatcherTest
private long tryRequest(Callable<Long> check, Message.Request request) throws Exception
{
long start = check.call();
dispatch.dispatch(request.connection().channel(), request, (channel, req, response) -> null,
ClientResourceLimits.Overload.NONE);
dispatch.dispatch(request.connection().channel(), request, (param, channel, req, response) -> null,
null, ClientResourceLimits.Overload.NONE);
long timeout = Clock.Global.currentTimeMillis();
while (start == check.call() && Clock.Global.currentTimeMillis() - timeout < 1000)
@ -589,7 +591,7 @@ public class MessageManagementDispatcherTest
return managementConnectionMock();
}
};
msg.setStreamId(1);
msg.setSource(new Envelope(Envelope.Header.dummy(1, msg.type), null));
return msg;
}
@ -606,11 +608,12 @@ public class MessageManagementDispatcherTest
}
@Override
void processRequest(Channel channel,
Message.Request request,
FlushItemConverter forFlusher,
ClientResourceLimits.Overload backpressure,
RequestTime requestTime)
<P> void processRequest(Channel channel,
Message.Request request,
FlushItemConverter<P> forFlusher,
P param,
ClientResourceLimits.Overload backpressure,
RequestTime requestTime)
{
// noop - just for testing routing
}