diff --git a/cmake/Modules/SourceFiles.cmake b/cmake/Modules/SourceFiles.cmake index 46a7ea81d60..4f35b64f001 100644 --- a/cmake/Modules/SourceFiles.cmake +++ b/cmake/Modules/SourceFiles.cmake @@ -48,6 +48,7 @@ set(VALKEY_SERVER_SRCS ${CMAKE_SOURCE_DIR}/src/cluster_slot_stats.c ${CMAKE_SOURCE_DIR}/src/crc16.c ${CMAKE_SOURCE_DIR}/src/crc16_slottable.c + ${CMAKE_SOURCE_DIR}/src/crc32.c ${CMAKE_SOURCE_DIR}/src/commandlog.c ${CMAKE_SOURCE_DIR}/src/eval.c ${CMAKE_SOURCE_DIR}/src/bio.c diff --git a/src/Makefile b/src/Makefile index f7429377a5d..228e08bf5cb 100644 --- a/src/Makefile +++ b/src/Makefile @@ -492,6 +492,7 @@ ENGINE_SERVER_OBJ = \ connection.o \ crc16.o \ crc16_slottable.o \ + crc32.o \ crc64.o \ crccombine.o \ crcspeed.o \ diff --git a/src/cluster.c b/src/cluster.c index 98d765f594c..35f72e6fa4c 100644 --- a/src/cluster.c +++ b/src/cluster.c @@ -1629,6 +1629,7 @@ void resetClusterStats(void) { server.cluster->stats_bus_module_bytes_sent = 0; server.cluster->stats_bus_module_bytes_received = 0; server.cluster->stat_cluster_links_buffer_limit_exceeded = 0; + server.cluster->stat_cluster_messages_crc_mismatch = 0; } void clusterCommandFlushslot(client *c) { diff --git a/src/cluster_legacy.c b/src/cluster_legacy.c index 2e2120504a8..109eab3f248 100644 --- a/src/cluster_legacy.c +++ b/src/cluster_legacy.c @@ -152,6 +152,73 @@ static inline clusterMsgLight *toClusterMsgLight(void *buf) { return (clusterMsgLight *)buf; } +/* Compute the CRC seed from the configured requirepass. + * + * Using the password as the seed ensures that only cluster nodes sharing + * the same requirepass can pass CRC verification on the cluster bus, which + * prevents messages from a foreign cluster (with a different password) from + * being accepted. */ +static uint32_t clusterCrcSeed(void) { + if (server.requirepass == NULL || sdslen(server.requirepass) == 0) return 0; + + return crc32(0, (const unsigned char *)server.requirepass, sdslen(server.requirepass)); +} + +/* Compute and set the CRC32 field in a cluster message. + * + * The CRC covers the entire message from the first byte to totlen, with + * the crc field itself temporarily zeroed during computation so it does + * not affect the result. */ +static void clusterMsgSetCRC(clusterMsg *hdr, uint32_t totlen) { + if (!server.cluster_crc_enabled) return; + + /* Zero the CRC field before computation so it does not contribute + * to the checksum. */ + hdr->crc = htonl(0); + uint32_t computed = crc32(clusterCrcSeed(), (const unsigned char *)hdr, totlen); + hdr->crc = htonl(computed); +} + +/* Verify the CRC32 checksum of a received cluster message. + * + * Returns 1 if the CRC is valid or verification is not applicable, 0 if a + * CRC mismatch is detected (the caller should drop the packet). + * + * Verification is skipped in the following cases, all of which are safe: + * + * 1. cluster-crc-enabled is off locally — this node does not participate + * in CRC verification at all. + * + * 2. The CRC field in the message is zero — the sender is an older version + * or has CRC disabled, so no CRC was computed. This is the key + * backward-compatibility path. Since zcalloc() initializes the entire + * message block to zero, the crc field is naturally zero when CRC is + * not active. + * + * Note: The CRC field itself serves as the capability indicator. */ +static int clusterMsgVerifyCRC(clusterMsg *hdr, uint32_t totlen) { + /* Case 1: CRC verification is disabled locally. */ + if (!server.cluster_crc_enabled) return 1; + + /* Case 2: CRC field is zero — sender did not compute CRC. */ + uint32_t received_crc = ntohl(hdr->crc); + if (received_crc == 0) return 1; + + /* Save the received CRC, zero the crc field, recompute, then restore. + * The crc field must be zero during computation so it does not affect + * the result, matching the sender's computation logic. */ + uint32_t saved_crc = hdr->crc; + hdr->crc = htonl(0); + uint32_t computed = crc32(clusterCrcSeed(), (const unsigned char *)hdr, totlen); + hdr->crc = saved_crc; + + /* CRC mismatch, return. */ + if (computed != received_crc) { + return 0; + } + return 1; +} + /* Only primaries that own slots have voting rights. * Returns 1 if the node has voting rights, otherwise returns 0. */ int clusterNodeIsVotingPrimary(clusterNode *n) { @@ -647,12 +714,12 @@ typedef struct { } data[]; } clusterMsgSendBlock; -/* Helper function to extract a normal message from a send block. */ +/* Helper function to extract a light message from a send block. */ static clusterMsgLight *getLightMessageFromSendBlock(clusterMsgSendBlock *msgblock) { return &msgblock->data[0].msg_light; } -/* Helper function to extract a light message from a send block. */ +/* Helper function to extract a normal message from a send block. */ static clusterMsg *getMessageFromSendBlock(clusterMsgSendBlock *msgblock) { return &msgblock->data[0].msg; } @@ -3728,6 +3795,12 @@ static void clusterBusAddNetworkBytesByType(uint16_t type, uint64_t bytes, bool } } +/* Last time we logged a global "CRC mismatch" warning. Rate-limited to once + * per interval (CLUSTER_CRC_MISMATCH_LOG_INTERVAL) to avoid flooding the log + * when a steady stream of corrupted packets arrives. */ +static mstime_t crc_mismatch_last_log = 0; +#define CLUSTER_CRC_MISMATCH_LOG_INTERVAL 30000 + int clusterIsValidPacket(clusterLink *link) { clusterMsgHeader *hdr = (clusterMsgHeader *)link->rcvbuf; uint32_t totlen = ntohl(hdr->totlen); @@ -3861,6 +3934,30 @@ int clusterIsValidPacket(clusterLink *link) { return 0; } + /* CRC32 integrity check for non-light cluster bus messages. Light messages + * use a compact header without a CRC field and are skipped. A CRC mismatch + * means the packet is corrupted (e.g. a network bit-flip) and must be + * treated as invalid so the packet is dropped to protect cluster state. */ + if (!is_light && server.cluster_crc_enabled) { + clusterMsg *msg = toClusterMsg(link->rcvbuf); + if (!clusterMsgVerifyCRC(msg, totlen)) { + if (server.mstime - crc_mismatch_last_log >= CLUSTER_CRC_MISMATCH_LOG_INTERVAL) { + crc_mismatch_last_log = server.mstime; + char ip[NET_IP_STR_LEN]; + int port = 0; + if (connAddrPeerName(link->conn, ip, sizeof(ip), &port) == C_OK) { + serverLog(LL_WARNING, "CRC mismatch on packet of type %s from node %.40s (%s:%d).", + clusterGetMessageTypeString(type), msg->sender, ip, port); + } else { + serverLog(LL_WARNING, "CRC mismatch on packet of type %s from node %.40s.", + clusterGetMessageTypeString(type), msg->sender); + } + } + server.cluster->stat_cluster_messages_crc_mismatch++; + return 0; + } + } + return 1; } @@ -4753,6 +4850,50 @@ void clusterReadHandler(connection *conn) { } } +/* Compute the CRC32 and apply debug bit-flip corruption before a message is sent. */ +static void clusterMsgFinalizeCRC(clusterMsgSendBlock *msgblock) { + serverAssert(server.cluster_crc_enabled); + clusterMsg *hdr = getMessageFromSendBlock(msgblock); + uint16_t msg_type = ntohs(hdr->type); + int is_light = IS_LIGHT_MESSAGE(msg_type); + uint32_t totlen = ntohl(hdr->totlen); + + /* CRC and debug corruption only apply to full-header messages with CRC + * enabled; light messages carry no crc field. */ + if (is_light) return; + + /* The crc == 0 guard avoids re-computation when the block is reused for + * multiple recipients (shared via refcount in clusterBroadcastMessage). */ + if (hdr->crc == 0) { + clusterMsgSetCRC(hdr, totlen); + } + + /* DEBUG cluster-crc-flip-bit: flip one bit in the next outgoing message to + * exercise the receiver-side CRC check. Runs after CRC so the corrupted + * packet is actually detected downstream. */ + if (server.debug_cluster_crc_flip_bit >= 0) { + int byte_off = server.debug_cluster_crc_flip_bit; + if (byte_off < (int)totlen) { + unsigned char *buf = (unsigned char *)hdr; + buf[byte_off] ^= 0x01; /* Flip the lowest bit. */ + serverLog(LL_WARNING, "DEBUG: flipped bit at byte offset %d in outgoing cluster message (type %s)", + byte_off, clusterGetMessageTypeString(msg_type & ~CLUSTERMSG_MODIFIER_MASK)); + } + server.debug_cluster_crc_flip_bit = -1; /* One-shot: disable after use. */ + } + + /* DEBUG cluster-crc-flip-time: randomly flip a bit in every outgoing message + * while the debug flip timer is active, to simulate sustained corruption. */ + if (server.debug_cluster_crc_flip_until > 0 && server.mstime < server.debug_cluster_crc_flip_until) { + unsigned char *buf = (unsigned char *)hdr; + int byte_off = rand() % totlen; + int bit = rand() & 0x07; + buf[byte_off] ^= (1 << bit); + serverLog(LL_WARNING, "DEBUG: flipped random bit %d at byte offset %d in outgoing cluster message (type %s)", + bit, byte_off, clusterGetMessageTypeString(msg_type & ~CLUSTERMSG_MODIFIER_MASK)); + } +} + /* Put the message block into the link's send queue. * * It is guaranteed that this function will never have as a side effect @@ -4762,6 +4903,10 @@ void clusterSendMessage(clusterLink *link, clusterMsgSendBlock *msgblock) { if (!link) { return; } + + /* Try and finalize the cluster CRC before sending. */ + if (server.cluster_crc_enabled) clusterMsgFinalizeCRC(msgblock); + if (listLength(link->send_msg_queue) == 0 && getMessageFromSendBlock(msgblock)->totlen != 0) connSetWriteHandlerWithBarrier(link->conn, clusterWriteHandler, 1); @@ -4844,6 +4989,7 @@ static void clusterBuildMessageHdr(clusterMsg *hdr, int type, size_t msglen) { memcpy(hdr->myslots, primary->slots, sizeof(hdr->myslots)); memset(hdr->replicaof, 0, CLUSTER_NAMELEN); if (myself->replicaof != NULL) memcpy(hdr->replicaof, myself->replicaof->name, CLUSTER_NAMELEN); + hdr->crc = htonl(0); if (server.tls_cluster) { hdr->port = htons(announced_tls_port); hdr->pport = htons(announced_tcp_port); @@ -7486,13 +7632,15 @@ sds genClusterInfoString(sds info) { "cluster_stats_pubsub_bytes_sent:%U\r\n" "cluster_stats_pubsub_bytes_received:%U\r\n" "cluster_stats_module_bytes_sent:%U\r\n" - "cluster_stats_module_bytes_received:%U\r\n", + "cluster_stats_module_bytes_received:%U\r\n" + "cluster_stats_messages_crc_mismatch:%U\r\n", (unsigned long long)server.cluster->stats_bus_bytes_sent, (unsigned long long)server.cluster->stats_bus_bytes_received, (unsigned long long)server.cluster->stats_bus_pubsub_bytes_sent, (unsigned long long)server.cluster->stats_bus_pubsub_bytes_received, (unsigned long long)server.cluster->stats_bus_module_bytes_sent, - (unsigned long long)server.cluster->stats_bus_module_bytes_received); + (unsigned long long)server.cluster->stats_bus_module_bytes_received, + (unsigned long long)server.cluster->stat_cluster_messages_crc_mismatch); info = sdscatfmt(info, "total_cluster_links_buffer_limit_exceeded:%U\r\n", (unsigned long long)server.cluster->stat_cluster_links_buffer_limit_exceeded); diff --git a/src/cluster_legacy.h b/src/cluster_legacy.h index be358b1899a..809dce66185 100644 --- a/src/cluster_legacy.h +++ b/src/cluster_legacy.h @@ -291,7 +291,10 @@ typedef struct { char replicaof[CLUSTER_NAMELEN]; char myip[NET_IP_STR_LEN]; /* Sender IP, if not all zeroed. */ uint16_t extensions; /* Number of extensions sent along with this packet. */ - char notused1[30]; /* 30 bytes reserved for future usage. */ + uint32_t crc; /* CRC32 checksum of the entire message. A non-zero value + * indicates the sender has CRC enabled; zero means CRC is + * not active (backward compatible). */ + char notused1[26]; /* 26 bytes reserved for future usage. */ uint16_t pport; /* Secondary port number: if primary port is TCP port, this is TLS port, and if primary port is TLS port, this is TCP port.*/ uint16_t cport; /* Sender TCP cluster bus port */ @@ -324,7 +327,8 @@ static_assert(offsetof(clusterMsg, myslots) == 80, "unexpected field offset"); static_assert(offsetof(clusterMsg, replicaof) == 2128, "unexpected field offset"); static_assert(offsetof(clusterMsg, myip) == 2168, "unexpected field offset"); static_assert(offsetof(clusterMsg, extensions) == 2214, "unexpected field offset"); -static_assert(offsetof(clusterMsg, notused1) == 2216, "unexpected field offset"); +static_assert(offsetof(clusterMsg, crc) == 2216, "unexpected field offset"); +static_assert(offsetof(clusterMsg, notused1) == 2220, "unexpected field offset"); static_assert(offsetof(clusterMsg, pport) == 2246, "unexpected field offset"); static_assert(offsetof(clusterMsg, cport) == 2248, "unexpected field offset"); static_assert(offsetof(clusterMsg, flags) == 2250, "unexpected field offset"); @@ -486,6 +490,7 @@ struct clusterState { excluding nodes without address. */ unsigned long long stat_cluster_links_buffer_limit_exceeded; /* Total number of cluster links freed due to exceeding buffer limit */ + unsigned long long stat_cluster_messages_crc_mismatch; /* Number of cluster messages dropped due to CRC mismatch. */ /* Bit map for slots that are no longer claimed by the owner in cluster PING * messages. During slot migration, the owner will stop claiming the slot after diff --git a/src/config.c b/src/config.c index f7e98f31c59..cf22542e51f 100644 --- a/src/config.c +++ b/src/config.c @@ -3391,6 +3391,7 @@ standardConfig static_configs[] = { createBoolConfig("lua-enable-insecure-api", "lua-enable-deprecated-api", MODIFIABLE_CONFIG | HIDDEN_CONFIG | PROTECTED_CONFIG, server.lua_enable_insecure_api, 0, NULL, updateLuaEnableInsecureApi), createBoolConfig("import-mode", NULL, DEBUG_CONFIG | MODIFIABLE_CONFIG, server.import_mode, 0, NULL, NULL), createBoolConfig("io-threads-always-active", NULL, MODIFIABLE_CONFIG | HIDDEN_CONFIG, server.io_threads_always_active, 0, NULL, NULL), + createBoolConfig("cluster-crc-enabled", NULL, MODIFIABLE_CONFIG, server.cluster_crc_enabled, 0, NULL, NULL), /* String Configs */ createStringConfig("aclfile", NULL, IMMUTABLE_CONFIG, ALLOW_EMPTY_STRING, server.acl_filename, "", NULL, NULL), diff --git a/src/crc32.c b/src/crc32.c new file mode 100644 index 00000000000..a0475ad0485 --- /dev/null +++ b/src/crc32.c @@ -0,0 +1,79 @@ +#include +#include + +/* + * CRC32 implementation according to the IEEE 802.3 (Ethernet) standard. + * + * This is the same CRC-32 variant used by zlib, PNG and Ethernet: + * + * Name : "CRC-32/ISO-HDLC" (a.k.a. IEEE 802.3) + * Width : 32 bit + * Poly : 0x04C11DB7 (reflected form 0xEDB88320) + * Initialization : 0xFFFFFFFF + * Reflect Input byte : True + * Reflect Output CRC : True + * Xor constant to output CRC : 0xFFFFFFFF + * Output for "123456789" : 0xCBF43926 + * + * The function takes a seed so that callers may chain multiple buffers or use a + * custom initial value (e.g. deriving the seed from a shared secret). Passing + * a seed of 0 yields the standard CRC-32. + */ + +/* CRC-32 (IEEE 802.3) lookup table, reflected polynomial 0xEDB88320. */ +static const uint32_t crc32_tab[256] = { + 0x00000000, 0x77073096, 0xee0e612c, 0x990951ba, 0x076dc419, 0x706af48f, + 0xe963a535, 0x9e6495a3, 0x0edb8832, 0x79dcb8a4, 0xe0d5e91e, 0x97d2d988, + 0x09b64c2b, 0x7eb17cbd, 0xe7b82d07, 0x90bf1d91, 0x1db71064, 0x6ab020f2, + 0xf3b97148, 0x84be41de, 0x1adad47d, 0x6ddde4eb, 0xf4d4b551, 0x83d385c7, + 0x136c9856, 0x646ba8c0, 0xfd62f97a, 0x8a65c9ec, 0x14015c4f, 0x63066cd9, + 0xfa0f3d63, 0x8d080df5, 0x3b6e20c8, 0x4c69105e, 0xd56041e4, 0xa2677172, + 0x3c03e4d1, 0x4b04d447, 0xd20d85fd, 0xa50ab56b, 0x35b5a8fa, 0x42b2986c, + 0xdbbbc9d6, 0xacbcf940, 0x32d86ce3, 0x45df5c75, 0xdcd60dcf, 0xabd13d59, + 0x26d930ac, 0x51de003a, 0xc8d75180, 0xbfd06116, 0x21b4f4b5, 0x56b3c423, + 0xcfba9599, 0xb8bda50f, 0x2802b89e, 0x5f058808, 0xc60cd9b2, 0xb10be924, + 0x2f6f7c87, 0x58684c11, 0xc1611dab, 0xb6662d3d, 0x76dc4190, 0x01db7106, + 0x98d220bc, 0xefd5102a, 0x71b18589, 0x06b6b51f, 0x9fbfe4a5, 0xe8b8d433, + 0x7807c9a2, 0x0f00f934, 0x9609a88e, 0xe10e9818, 0x7f6a0dbb, 0x086d3d2d, + 0x91646c97, 0xe6635c01, 0x6b6b51f4, 0x1c6c6162, 0x856530d8, 0xf262004e, + 0x6c0695ed, 0x1b01a57b, 0x8208f4c1, 0xf50fc457, 0x65b0d9c6, 0x12b7e950, + 0x8bbeb8ea, 0xfcb9887c, 0x62dd1ddf, 0x15da2d49, 0x8cd37cf3, 0xfbd44c65, + 0x4db26158, 0x3ab551ce, 0xa3bc0074, 0xd4bb30e2, 0x4adfa541, 0x3dd895d7, + 0xa4d1c46d, 0xd3d6f4fb, 0x4369e96a, 0x346ed9fc, 0xad678846, 0xda60b8d0, + 0x44042d73, 0x33031de5, 0xaa0a4c5f, 0xdd0d7cc9, 0x5005713c, 0x270241aa, + 0xbe0b1010, 0xc90c2086, 0x5768b525, 0x206f85b3, 0xb966d409, 0xce61e49f, + 0x5edef90e, 0x29d9c998, 0xb0d09822, 0xc7d7a8b4, 0x59b33d17, 0x2eb40d81, + 0xb7bd5c3b, 0xc0ba6cad, 0xedb88320, 0x9abfb3b6, 0x03b6e20c, 0x74b1d29a, + 0xead54739, 0x9dd277af, 0x04db2615, 0x73dc1683, 0xe3630b12, 0x94643b84, + 0x0d6d6a3e, 0x7a6a5aa8, 0xe40ecf0b, 0x9309ff9d, 0x0a00ae27, 0x7d079eb1, + 0xf00f9344, 0x8708a3d2, 0x1e01f268, 0x6906c2fe, 0xf762575d, 0x806567cb, + 0x196c3671, 0x6e6b06e7, 0xfed41b76, 0x89d32be0, 0x10da7a5a, 0x67dd4acc, + 0xf9b9df6f, 0x8ebeeff9, 0x17b7be43, 0x60b08ed5, 0xd6d6a3e8, 0xa1d1937e, + 0x38d8c2c4, 0x4fdff252, 0xd1bb67f1, 0xa6bc5767, 0x3fb506dd, 0x48b2364b, + 0xd80d2bda, 0xaf0a1b4c, 0x36034af6, 0x41047a60, 0xdf60efc3, 0xa867df55, + 0x316e8eef, 0x4669be79, 0xcb61b38c, 0xbc66831a, 0x256fd2a0, 0x5268e236, + 0xcc0c7795, 0xbb0b4703, 0x220216b9, 0x5505262f, 0xc5ba3bbe, 0xb2bd0b28, + 0x2bb45a92, 0x5cb30a04, 0xc2d7ffa7, 0xb5d0cf31, 0x2cd99e8b, 0x5bdeae1d, + 0x9b64c2b0, 0xec63f226, 0x756aa39c, 0x026d930a, 0x9c0906a9, 0xeb0e363f, + 0x72676785, 0x05605713, 0x95bf4a82, 0xe2b87a14, 0x7bb12bae, 0x0cb61b38, + 0x92d28e9b, 0xe5d5be0d, 0x7cdcefb7, 0x0bdbdf21, 0x86d3d2d4, 0xf1d4e242, + 0x68ddb3f8, 0x1fda836e, 0x81be16cd, 0xf6b9265b, 0x6fb077e1, 0x18b74777, + 0x88085ae6, 0xff0f6a70, 0x66063bca, 0x11010b5c, 0x8f659eff, 0xf862ae69, + 0x616bffd3, 0x166ccf45, 0xa00ae278, 0xd70dd2ee, 0x4e048354, 0x3903b3c2, + 0xa7672661, 0xd06016f7, 0x4969474d, 0x3e6e77db, 0xaed16a4a, 0xd9d65adc, + 0x40df0b66, 0x37d83bf0, 0xa9bcae53, 0xdebb9ec5, 0x47b2cf7f, 0x30b5ffe9, + 0xbdbdf21c, 0xcabac28a, 0x53b39330, 0x24b4a3a6, 0xbad03605, 0xcdd70693, + 0x54de5729, 0x23d967bf, 0xb3667a2e, 0xc4614ab8, 0x5d681b02, 0x2a6f2b94, + 0xb40bbe37, 0xc30c8ea1, 0x5a05df1b, 0x2d02ef8d}; + +/* Compute the CRC-32 (IEEE 802.3) checksum of the given buffer. + * + * The 'seed' allows callers to provide a custom initial value; passing 0 + * produces the standard CRC-32 checksum. */ +uint32_t crc32(uint32_t seed, const unsigned char *buf, size_t len) { + uint32_t crc = seed ^ 0xFFFFFFFF; + for (size_t i = 0; i < len; i++) { + crc = crc32_tab[(crc ^ buf[i]) & 0xFF] ^ (crc >> 8); + } + return crc ^ 0xFFFFFFFF; +} diff --git a/src/debug.c b/src/debug.c index 4c7f0efa6ee..bfc049da4a9 100644 --- a/src/debug.c +++ b/src/debug.c @@ -450,6 +450,14 @@ void debugCommand(client *c) { " Disable sending cluster ping to a random node every second.", "DISABLE-CLUSTER-RECONNECTION <0|1>", " Disable cluster reconnection of cluster nodes.", + "CLUSTER-CRC-FLIP-BIT ", + " Flip a single bit in the next outgoing cluster message at to", + " simulate network corruption and exercise the receiver-side CRC check.", + " Set to -1 to disable.", + "CLUSTER-CRC-FLIP-TIME ", + " Randomly flip a bit in every outgoing full-header cluster message for the", + " given number of to simulate sustained network corruption.", + " Set to 0 to disable.", "OOM", " Crash the server simulating an out-of-memory error.", "PANIC", @@ -644,6 +652,15 @@ void debugCommand(client *c) { } else if (!strcasecmp(objectGetVal(c->argv[1]), "disable-cluster-reconnection") && c->argc == 3) { server.debug_cluster_disable_reconnection = atoi(objectGetVal(c->argv[2])); addReply(c, shared.ok); + } else if (!strcasecmp(objectGetVal(c->argv[1]), "cluster-crc-flip-bit") && c->argc == 3) { + server.debug_cluster_crc_flip_bit = atoi(objectGetVal(c->argv[2])); + addReply(c, shared.ok); + } else if (!strcasecmp(objectGetVal(c->argv[1]), "cluster-crc-flip-time") && c->argc == 3) { + long seconds; + if (getLongFromObjectOrReply(c, c->argv[2], &seconds, NULL) != C_OK) return; + if (seconds < 0) seconds = 0; + server.debug_cluster_crc_flip_until = seconds > 0 ? mstime() + (long long)seconds * 1000 : 0; + addReply(c, shared.ok); } else if (!strcasecmp(objectGetVal(c->argv[1]), "slotmigration")) { if (!strcasecmp(objectGetVal(c->argv[2]), "prevent-pause")) { server.debug_slot_migration_prevent_pause = atoi(objectGetVal(c->argv[3])); diff --git a/src/server.c b/src/server.c index df2e91f9a7c..299af2b77c6 100644 --- a/src/server.c +++ b/src/server.c @@ -2975,6 +2975,8 @@ void initServer(void) { server.cluster_drop_packet_filter = -1; server.debug_cluster_disable_random_ping = 0; server.debug_cluster_disable_reconnection = 0; + server.debug_cluster_crc_flip_bit = -1; + server.debug_cluster_crc_flip_until = 0; server.reply_buffer_peak_reset_time = REPLY_BUFFER_DEFAULT_PEAK_RESET_TIME; server.reply_buffer_resizing_enabled = 1; server.client_mem_usage_buckets = NULL; diff --git a/src/server.h b/src/server.h index a9c4fcdbefa..8e9f0184e9e 100644 --- a/src/server.h +++ b/src/server.h @@ -2308,6 +2308,7 @@ struct valkeyServer { connection *slot_migration_pipe_conn; /* xxxx */ char *slot_migration_pipe_buff; /* In slot migration, this buffer holds slot snapshot data. */ ssize_t slot_migration_pipe_bufflen; /* that was read from the rdb pipe. */ + int cluster_crc_enabled; /* Enable CRC32 checksum for cluster bus messages. */ /* Debug config that goes along with cluster_drop_packet_filter. When set, the link is closed on packet drop. */ uint32_t debug_cluster_close_link_on_packet_drop : 1; /* Debug config to control the random ping. When set, we will disable the random ping in clusterCron. */ @@ -2317,6 +2318,10 @@ struct valkeyServer { /* Debug config to expose intermediary slot migration states. */ uint32_t debug_slot_migration_prevent_pause : 1; uint32_t debug_slot_migration_prevent_failover : 1; + /* Debug config to flip a bit at this byte offset in the next outgoing cluster message. -1 means disabled. */ + int debug_cluster_crc_flip_bit; + /* Debug config to flip a random bit in every outgoing cluster message until this mstime. 0 means disabled. */ + mstime_t debug_cluster_crc_flip_until; sds cached_cluster_slot_info[CACHE_CONN_TYPE_MAX]; /* Index in array is a bitwise or of CACHE_CONN_TYPE_* */ /* Scripting */ mstime_t busy_reply_threshold; /* Script / module timeout in milliseconds */ @@ -3847,6 +3852,7 @@ int setGetKeys(struct serverCommand *cmd, robj **argv, int argc, getKeysResult * int bitfieldGetKeys(struct serverCommand *cmd, robj **argv, int argc, getKeysResult *result); unsigned short crc16(const char *buf, int len); +uint32_t crc32(uint32_t seed, const unsigned char *buf, size_t len); /* Sentinel */ void initSentinelConfig(void);