Skip to content
Open
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
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
Contributor

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;
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