CEP-15: Extend Accord MessageType with a side effect flag

patch by Aleksey Yeschenko; reviewed by Benedic Elliott Smith for
CASSANDRA-18561
This commit is contained in:
Aleksey Yeschenko 2023-06-02 11:45:41 +01:00 committed by David Capwell
parent f76f806de1
commit 8f156fc5dc
5 changed files with 22 additions and 22 deletions

@ -1 +1 @@
Subproject commit 3d0ff07cd5c7db43390b85afa593e6f76471d886
Subproject commit 8830d97ba517fb2d0f7f22e8e6b886a98839e694

View File

@ -266,8 +266,8 @@ public enum Verb
// accord
ACCORD_SIMPLE_RSP (119, P2, writeTimeout, REQUEST_RESPONSE, () -> EnumSerializer.simpleReply, RESPONSE_HANDLER ),
ACCORD_PREACCEPT_RSP (121, P2, writeTimeout, REQUEST_RESPONSE, () -> PreacceptSerializers.reply, RESPONSE_HANDLER ),
ACCORD_PREACCEPT_REQ (120, P2, writeTimeout, IMMEDIATE, () -> PreacceptSerializers.request, () -> AccordService.instance().verbHandler(), ACCORD_PREACCEPT_RSP ),
ACCORD_PRE_ACCEPT_RSP (121, P2, writeTimeout, REQUEST_RESPONSE, () -> PreacceptSerializers.reply, RESPONSE_HANDLER ),
ACCORD_PRE_ACCEPT_REQ (120, P2, writeTimeout, IMMEDIATE, () -> PreacceptSerializers.request, () -> AccordService.instance().verbHandler(), ACCORD_PRE_ACCEPT_RSP ),
ACCORD_ACCEPT_RSP (124, P2, writeTimeout, REQUEST_RESPONSE, () -> AcceptSerializers.reply, RESPONSE_HANDLER ),
ACCORD_ACCEPT_REQ (122, P2, writeTimeout, IMMEDIATE, () -> AcceptSerializers.request, () -> AccordService.instance().verbHandler(), ACCORD_ACCEPT_RSP ),
ACCORD_ACCEPT_INVALIDATE_REQ (123, P2, writeTimeout, IMMEDIATE, () -> AcceptSerializers.invalidate, () -> AccordService.instance().verbHandler(), ACCORD_ACCEPT_RSP ),
@ -277,13 +277,13 @@ public enum Verb
ACCORD_COMMIT_INVALIDATE_REQ (126, P2, writeTimeout, IMMEDIATE, () -> CommitSerializers.invalidate, () -> AccordService.instance().verbHandler() ),
ACCORD_APPLY_RSP (130, P2, writeTimeout, REQUEST_RESPONSE, () -> ApplySerializers.reply, RESPONSE_HANDLER ),
ACCORD_APPLY_REQ (129, P2, writeTimeout, IMMEDIATE, () -> ApplySerializers.request, () -> AccordService.instance().verbHandler(), ACCORD_APPLY_RSP ),
ACCORD_RECOVER_RSP (132, P2, writeTimeout, REQUEST_RESPONSE, () -> RecoverySerializers.reply, RESPONSE_HANDLER ),
ACCORD_RECOVER_REQ (131, P2, writeTimeout, IMMEDIATE, () -> RecoverySerializers.request, () -> AccordService.instance().verbHandler(), ACCORD_RECOVER_RSP ),
ACCORD_BEGIN_RECOVER_RSP (132, P2, writeTimeout, REQUEST_RESPONSE, () -> RecoverySerializers.reply, RESPONSE_HANDLER ),
ACCORD_BEGIN_RECOVER_REQ (131, P2, writeTimeout, IMMEDIATE, () -> RecoverySerializers.request, () -> AccordService.instance().verbHandler(), ACCORD_BEGIN_RECOVER_RSP ),
ACCORD_BEGIN_INVALIDATE_RSP (134, P2, writeTimeout, REQUEST_RESPONSE, () -> BeginInvalidationSerializers.reply, RESPONSE_HANDLER ),
ACCORD_BEGIN_INVALIDATE_REQ (133, P2, writeTimeout, IMMEDIATE, () -> BeginInvalidationSerializers.request, () -> AccordService.instance().verbHandler(), ACCORD_BEGIN_INVALIDATE_RSP ),
ACCORD_WAIT_COMMIT_RSP (136, P2, writeTimeout, REQUEST_RESPONSE, () -> WaitOnCommitSerializer.reply, RESPONSE_HANDLER ),
ACCORD_WAIT_COMMIT_REQ (135, P2, writeTimeout, IMMEDIATE, () -> WaitOnCommitSerializer.request, () -> AccordService.instance().verbHandler(), ACCORD_WAIT_COMMIT_RSP ),
ACCORD_INFORM_OF_TXNID_REQ (137, P2, writeTimeout, IMMEDIATE, () -> InformOfTxnIdSerializers.request, () -> AccordService.instance().verbHandler(), ACCORD_SIMPLE_RSP ),
ACCORD_WAIT_ON_COMMIT_RSP (136, P2, writeTimeout, REQUEST_RESPONSE, () -> WaitOnCommitSerializer.reply, RESPONSE_HANDLER ),
ACCORD_WAIT_ON_COMMIT_REQ (135, P2, writeTimeout, IMMEDIATE, () -> WaitOnCommitSerializer.request, () -> AccordService.instance().verbHandler(), ACCORD_WAIT_ON_COMMIT_RSP ),
ACCORD_INFORM_OF_TXN_REQ (137, P2, writeTimeout, IMMEDIATE, () -> InformOfTxnIdSerializers.request, () -> AccordService.instance().verbHandler(), ACCORD_SIMPLE_RSP ),
ACCORD_INFORM_HOME_DURABLE_REQ (138, P2, writeTimeout, IMMEDIATE, () -> InformHomeDurableSerializers.request, () -> AccordService.instance().verbHandler(), ACCORD_SIMPLE_RSP ),
ACCORD_INFORM_DURABLE_REQ (139, P2, writeTimeout, IMMEDIATE, () -> InformDurableSerializers.request, () -> AccordService.instance().verbHandler(), ACCORD_SIMPLE_RSP ),
ACCORD_CHECK_STATUS_RSP (141, P2, writeTimeout, REQUEST_RESPONSE, () -> CheckStatusSerializers.reply, RESPONSE_HANDLER ),

View File

@ -382,10 +382,10 @@ public class AccordJournal
*/
public enum Type implements IVersionedSerializer<TxnRequest<?>>
{
PREACCEPT_REQ (0, MessageType.PREACCEPT_REQ, PreacceptSerializers.request),
ACCEPT_REQ (1, MessageType.ACCEPT_REQ, AcceptSerializers.request ),
COMMIT_REQ (2, MessageType.COMMIT_REQ, CommitSerializers.request ),
APPLY_REQ (3, MessageType.APPLY_REQ, ApplySerializers.request );
PREACCEPT_REQ (0, MessageType.PRE_ACCEPT_REQ, PreacceptSerializers.request),
ACCEPT_REQ (1, MessageType.ACCEPT_REQ, AcceptSerializers.request ),
COMMIT_REQ (2, MessageType.COMMIT_REQ, CommitSerializers.request ),
APPLY_REQ (3, MessageType.APPLY_REQ, ApplySerializers.request );
final int id;
final MessageType msgType;
@ -460,7 +460,7 @@ public class AccordJournal
static boolean mustMakeDurable(TxnRequest<?> message)
{
return msgTypeToTypeMap.containsKey(message.type());
return message.type().hasSideEffects;
}
@Override

View File

@ -52,24 +52,24 @@ public class AccordMessageSink implements MessageSink
private VerbMapping()
{
mapping.put(MessageType.PREACCEPT_REQ, Verb.ACCORD_PREACCEPT_REQ);
mapping.put(MessageType.PREACCEPT_RSP, Verb.ACCORD_PREACCEPT_RSP);
mapping.put(MessageType.PRE_ACCEPT_REQ, Verb.ACCORD_PRE_ACCEPT_REQ);
mapping.put(MessageType.PRE_ACCEPT_RSP, Verb.ACCORD_PRE_ACCEPT_RSP);
mapping.put(MessageType.ACCEPT_REQ, Verb.ACCORD_ACCEPT_REQ);
mapping.put(MessageType.ACCEPT_RSP, Verb.ACCORD_ACCEPT_RSP);
mapping.put(MessageType.ACCEPT_INVALIDATE_REQ, Verb.ACCORD_ACCEPT_INVALIDATE_REQ);
mapping.put(MessageType.COMMIT_REQ, Verb.ACCORD_COMMIT_REQ);
mapping.put(MessageType.COMMIT_INVALIDATE, Verb.ACCORD_COMMIT_INVALIDATE_REQ);
mapping.put(MessageType.COMMIT_INVALIDATE_REQ, Verb.ACCORD_COMMIT_INVALIDATE_REQ);
mapping.put(MessageType.APPLY_REQ, Verb.ACCORD_APPLY_REQ);
mapping.put(MessageType.APPLY_RSP, Verb.ACCORD_APPLY_RSP);
mapping.put(MessageType.READ_REQ, Verb.ACCORD_READ_REQ);
mapping.put(MessageType.READ_RSP, Verb.ACCORD_READ_RSP);
mapping.put(MessageType.BEGIN_RECOVER_REQ, Verb.ACCORD_RECOVER_REQ);
mapping.put(MessageType.BEGIN_RECOVER_RSP, Verb.ACCORD_RECOVER_RSP);
mapping.put(MessageType.BEGIN_RECOVER_REQ, Verb.ACCORD_BEGIN_RECOVER_REQ);
mapping.put(MessageType.BEGIN_RECOVER_RSP, Verb.ACCORD_BEGIN_RECOVER_RSP);
mapping.put(MessageType.BEGIN_INVALIDATE_REQ, Verb.ACCORD_BEGIN_INVALIDATE_REQ);
mapping.put(MessageType.BEGIN_INVALIDATE_RSP, Verb.ACCORD_BEGIN_INVALIDATE_RSP);
mapping.put(MessageType.WAIT_ON_COMMIT_REQ, Verb.ACCORD_WAIT_COMMIT_REQ);
mapping.put(MessageType.WAIT_ON_COMMIT_RSP, Verb.ACCORD_WAIT_COMMIT_RSP);
mapping.put(MessageType.INFORM_TXNID_REQ, Verb.ACCORD_INFORM_OF_TXNID_REQ);
mapping.put(MessageType.WAIT_ON_COMMIT_REQ, Verb.ACCORD_WAIT_ON_COMMIT_REQ);
mapping.put(MessageType.WAIT_ON_COMMIT_RSP, Verb.ACCORD_WAIT_ON_COMMIT_RSP);
mapping.put(MessageType.INFORM_OF_TXN_REQ, Verb.ACCORD_INFORM_OF_TXN_REQ);
mapping.put(MessageType.INFORM_HOME_DURABLE_REQ,Verb.ACCORD_INFORM_HOME_DURABLE_REQ);
mapping.put(MessageType.INFORM_DURABLE_REQ, Verb.ACCORD_INFORM_DURABLE_REQ);
mapping.put(MessageType.CHECK_STATUS_REQ, Verb.ACCORD_CHECK_STATUS_REQ);

View File

@ -50,7 +50,7 @@ public class AccordMessageSinkTest
// There was an issue where the reply was the wrong verb
// see CASSANDRA-18375
InformOfTxnId info = Mockito.mock(InformOfTxnId.class);
Message<InformOfTxnId> req = Message.builder(Verb.ACCORD_INFORM_OF_TXNID_REQ, info).build();
Message<InformOfTxnId> req = Message.builder(Verb.ACCORD_INFORM_OF_TXN_REQ, info).build();
SimpleReply reply = SimpleReply.Ok;
MessageDelivery messaging = Mockito.mock(MessageDelivery.class);