Skip to content
Open
Show file tree
Hide file tree
Changes from 1 commit
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
1 change: 1 addition & 0 deletions cmake/Modules/SourceFiles.cmake
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
1 change: 1 addition & 0 deletions src/Makefile
Original file line number Diff line number Diff line change
Expand Up @@ -492,6 +492,7 @@ ENGINE_SERVER_OBJ = \
connection.o \
crc16.o \
crc16_slottable.o \
crc32.o \
crc64.o \
crccombine.o \
crcspeed.o \
Expand Down
1 change: 1 addition & 0 deletions src/cluster.c
Original file line number Diff line number Diff line change
Expand Up @@ -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) {
Expand Down
156 changes: 152 additions & 4 deletions src/cluster_legacy.c
Original file line number Diff line number Diff line change
Expand Up @@ -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));

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

server.requirepass is not a stable cluster-wide secret. It is a per-node, modifiable ACL setting (src/config.c:3428-3429, src/server.h:2344-2346), and the cluster bus itself just connects and sends PING/MEET without any password negotiation (src/cluster_legacy.c:4711-4733, src/cluster_legacy.c:6308-6323). Using it as the CRC seed means a normal rolling CONFIG SET requirepass changes the seed on one node immediately, so once cluster-crc-enabled is on that node starts rejecting full-header packets from peers that have not been updated yet and can drive the cluster into PFAIL/FAIL. The seed needs to come from a dedicated cluster-wide secret/capability, or the CRC should stay unkeyed.

}

/* 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) {
Expand Down Expand Up @@ -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;
}
Expand Down Expand Up @@ -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);
Expand Down Expand Up @@ -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;
Comment thread
coderabbitai[bot] marked this conversation as resolved.
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;
}

Expand Down Expand Up @@ -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
Expand All @@ -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);

Expand Down Expand Up @@ -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);
Expand Down Expand Up @@ -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);
Expand Down
9 changes: 7 additions & 2 deletions src/cluster_legacy.h
Original file line number Diff line number Diff line change
Expand Up @@ -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 */
Expand Down Expand Up @@ -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");
Expand Down Expand Up @@ -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
Expand Down
1 change: 1 addition & 0 deletions src/config.c
Original file line number Diff line number Diff line change
Expand Up @@ -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),
Expand Down
79 changes: 79 additions & 0 deletions src/crc32.c
Original file line number Diff line number Diff line change
@@ -0,0 +1,79 @@
#include <stdint.h>
#include <stddef.h>

/*
* 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;
}
Loading
Loading