/*
 * Copyright (c) 2009-2012, Redis Ltd.
 * All rights reserved.
 *
 * Redistribution and use in source and binary forms, with or without
 * modification, are permitted provided that the following conditions are met:
 *
 *   * Redistributions of source code must retain the above copyright notice,
 *     this list of conditions and the following disclaimer.
 *   * Redistributions in binary form must reproduce the above copyright
 *     notice, this list of conditions and the following disclaimer in the
 *     documentation and/or other materials provided with the distribution.
 *   * Neither the name of Redis nor the names of its contributors may be used
 *     to endorse or promote products derived from this software without
 *     specific prior written permission.
 *
 * THIS SOFTWARE IS PROVIDED BY THE COPYRIGHT HOLDERS AND CONTRIBUTORS "AS IS"
 * AND ANY EXPRESS OR IMPLIED WARRANTIES, INCLUDING, BUT NOT LIMITED TO, THE
 * IMPLIED WARRANTIES OF MERCHANTABILITY AND FITNESS FOR A PARTICULAR PURPOSE
 * ARE DISCLAIMED. IN NO EVENT SHALL THE COPYRIGHT OWNER OR CONTRIBUTORS BE
 * LIABLE FOR ANY DIRECT, INDIRECT, INCIDENTAL, SPECIAL, EXEMPLARY, OR
 * CONSEQUENTIAL DAMAGES (INCLUDING, BUT NOT LIMITED TO, PROCUREMENT OF
 * SUBSTITUTE GOODS OR SERVICES; LOSS OF USE, DATA, OR PROFITS; OR BUSINESS
 * INTERRUPTION) HOWEVER CAUSED AND ON ANY THEORY OF LIABILITY, WHETHER IN
 * CONTRACT, STRICT LIABILITY, OR TORT (INCLUDING NEGLIGENCE OR OTHERWISE)
 * ARISING IN ANY WAY OUT OF THE USE OF THIS SOFTWARE, EVEN IF ADVISED OF THE
 * POSSIBILITY OF SUCH DAMAGE.
 */
/*
 * Copyright (c) Valkey Contributors
 * All rights reserved.
 * SPDX-License-Identifier: BSD-3-Clause
 */
/*
 * cluster.c contains the common parts of a clustering
 * implementation, the parts that are shared between
 * any implementation of clustering.
 */

#include "server.h"
#include "cluster.h"
#include "cluster_slot_stats.h"
#include "module.h"
#include "crc16_slottable.h"

#include <ctype.h>

/* -----------------------------------------------------------------------------
 * Key space handling
 * -------------------------------------------------------------------------- */

/* We have 16384 hash slots. The hash slot of a given key is obtained
 * as the least significant 14 bits of the crc16 of the key.
 *
 * However, if the key contains the {...} pattern, only the part between
 * { and } is hashed. This may be useful in the future to force certain
 * keys to be in the same node (assuming no resharding is in progress). */
unsigned int keyHashSlot(const char *key, int keylen) {
    int s, e; /* start-end indexes of { and } */

    for (s = 0; s < keylen; s++)
        if (key[s] == '{') break;

    /* No '{' ? Hash the whole key. This is the base case. */
    if (s == keylen) return crc16(key, keylen) & 0x3FFF;

    /* '{' found? Check if we have the corresponding '}'. */
    for (e = s + 1; e < keylen; e++)
        if (key[e] == '}') break;

    /* No '}' or nothing between {} ? Hash the whole key. */
    if (e == keylen || e == s + 1) return crc16(key, keylen) & 0x3FFF;

    /* If we are here there is both a { and a } on its right. Hash
     * what is in the middle between { and }. */
    return crc16(key + s + 1, e - s - 1) & 0x3FFF;
}

/* If it can be inferred that the given glob-style pattern, as implemented in
 * stringmatchlen() in util.c, only can match keys belonging to a single slot,
 * that slot is returned. Otherwise, -1 is returned. */
int patternHashSlot(char *pattern, int length) {
    int s = -1; /* index of the first '{' */

    for (int i = 0; i < length; i++) {
        if (pattern[i] == '*' || pattern[i] == '?' || pattern[i] == '[') {
            /* Wildcard or character class found. Keys can be in any slot. */
            return -1;
        } else if (pattern[i] == '\\') {
            /* Escaped character. Computing slot in this case is not
             * implemented. We would need a temp buffer. */
            return -1;
        } else if (s == -1 && pattern[i] == '{') {
            /* Opening brace '{' found. */
            s = i;
        } else if (s >= 0 && pattern[i] == '}' && i == s + 1) {
            /* Empty tag '{}' found. The whole key is hashed. Ignore braces. */
            s = -2;
        } else if (s >= 0 && pattern[i] == '}') {
            /* Non-empty tag '{...}' found. Hash what's between braces. */
            return crc16(pattern + s + 1, i - s - 1) & 0x3FFF;
        }
    }

    /* The pattern matches a single key. Hash the whole pattern. */
    return crc16(pattern, length) & 0x3FFF;
}

ConnectionType *connTypeOfCluster(void) {
    if (server.tls_cluster) {
        return connectionTypeTls();
    }

    return connectionTypeTcp();
}

/* -----------------------------------------------------------------------------
 * DUMP, RESTORE and MIGRATE commands
 * -------------------------------------------------------------------------- */

/* Generates a DUMP-format representation of the object 'o', adding it to the
 * io stream pointed by 'rio'. This function can't fail. */
void createDumpPayload(rio *payload, robj *o, robj *key, int dbid) {
    unsigned char buf[2];
    uint64_t crc;

    /* Serialize the object in an RDB-like format. It consist of an object type
     * byte followed by the serialized object. This is understood by RESTORE. */
    rioInitWithBuffer(payload, sdsempty());
    int rdbtype = rdbGetObjectType(o, RDB_VERSION);
    serverAssert(rdbtype >= 0);
    serverAssert(rdbSaveType(payload, rdbtype));
    serverAssert(rdbSaveObject(payload, o, key, dbid, rdbtype));

    /* Write the footer, this is how it looks like:
     * ----------------+---------------------+---------------+
     * ... RDB payload | 2 bytes RDB version | 8 bytes CRC64 |
     * ----------------+---------------------+---------------+
     * RDB version and CRC are both in little endian.
     */

    /* RDB version */
    buf[0] = RDB_VERSION & 0xff;
    buf[1] = (RDB_VERSION >> 8) & 0xff;
    payload->io.buffer.ptr = sdscatlen(payload->io.buffer.ptr, buf, 2);

    /* CRC64 */
    crc = crc64(0, (unsigned char *)payload->io.buffer.ptr, sdslen(payload->io.buffer.ptr));
    memrev64ifbe(&crc);
    payload->io.buffer.ptr = sdscatlen(payload->io.buffer.ptr, &crc, 8);
}

/* Verify that the RDB version of the dump payload matches the one of this
 * instance and that the checksum is ok.
 * If the DUMP payload looks valid C_OK is returned, otherwise C_ERR
 * is returned. If rdbver_ptr is not NULL, its populated with the value read
 * from the input buffer. */
int verifyDumpPayload(unsigned char *p, size_t len, uint16_t *rdbver_ptr) {
    unsigned char *footer;
    uint16_t rdbver;
    uint64_t crc;

    /* At least 2 bytes of RDB version and 8 of CRC64 should be present. */
    if (len < 10) return C_ERR;
    footer = p + (len - 10);

    /* Set and verify RDB version. */
    rdbver = (footer[1] << 8) | footer[0];
    if (rdbver_ptr) {
        *rdbver_ptr = rdbver;
    }
    if (!rdbIsVersionAccepted(rdbver, false, false)) return C_ERR;
    if (server.skip_checksum_validation) return C_OK;

    /* Verify CRC64 */
    crc = crc64(0, p, len - 8);
    memrev64ifbe(&crc);
    return (memcmp(&crc, footer + 2, 8) == 0) ? C_OK : C_ERR;
}

/* DUMP keyname
 * DUMP is actually not used by Cluster but it is the obvious
 * complement of RESTORE and can be useful for different applications. */
void dumpCommand(client *c) {
    robj *o;
    rio payload;

    /* Check if the key is here. */
    if ((o = lookupKeyRead(c->db, c->argv[1])) == NULL) {
        addReplyNull(c);
        return;
    }

    /* Create the DUMP encoded representation. */
    createDumpPayload(&payload, o, c->argv[1], c->db->id);

    /* Transfer to the client */
    addReplyBulkSds(c, payload.io.buffer.ptr);
    return;
}

/* RESTORE key ttl serialized-value [REPLACE] [ABSTTL] [IDLETIME seconds] [FREQ frequency] */
void restoreCommand(client *c) {
    long long ttl, lfu_freq = -1, lru_idle = -1;
    uint16_t rdbver = 0;
    rio payload;
    int j, type, replace = 0, absttl = 0;
    robj *obj;

    /* Parse additional options */
    for (j = 4; j < c->argc; j++) {
        int additional = c->argc - j - 1;
        if (!strcasecmp(objectGetVal(c->argv[j]), "replace")) {
            replace = 1;
        } else if (!strcasecmp(objectGetVal(c->argv[j]), "absttl")) {
            absttl = 1;
        } else if (!strcasecmp(objectGetVal(c->argv[j]), "idletime") && additional >= 1 && lfu_freq == -1) {
            if (getLongLongFromObjectOrReply(c, c->argv[j + 1], &lru_idle, NULL) != C_OK) return;
            if (lru_idle < 0) {
                addReplyError(c, "Invalid IDLETIME value, must be >= 0");
                return;
            }
            j++; /* Consume additional arg. */
        } else if (!strcasecmp(objectGetVal(c->argv[j]), "freq") && additional >= 1 && lru_idle == -1) {
            if (getLongLongFromObjectOrReply(c, c->argv[j + 1], &lfu_freq, NULL) != C_OK) return;
            if (lfu_freq < 0 || lfu_freq > 255) {
                addReplyError(c, "Invalid FREQ value, must be >= 0 and <= 255");
                return;
            }
            j++; /* Consume additional arg. */
        } else {
            addReplyErrorObject(c, shared.syntaxerr);
            return;
        }
    }

    /* Make sure this key does not already exist here... */
    robj *key = c->argv[1];
    if (!replace && lookupKeyWrite(c->db, key) != NULL) {
        addReplyErrorObject(c, shared.busykeyerr);
        return;
    }

    /* Check if the TTL value makes sense */
    if (getLongLongFromObjectOrReply(c, c->argv[2], &ttl, NULL) != C_OK) {
        return;
    } else if (ttl < 0) {
        addReplyError(c, "Invalid TTL value, must be >= 0");
        return;
    }

    /* Verify RDB version and data checksum. */
    if (verifyDumpPayload(objectGetVal(c->argv[3]), sdslen(objectGetVal(c->argv[3])), &rdbver) == C_ERR) {
        addReplyError(c, "DUMP payload version or checksum are wrong");
        return;
    }

    rioInitWithBuffer(&payload, objectGetVal(c->argv[3]));
    type = rdbLoadObjectType(&payload);
    if (type == -1) {
        addReplyError(c, "Bad data format");
        return;
    }

    /* If it's a foreign RDB format, only accept old data types that we know
     * existed in the past and that don't clash with new types added later. */
    if (rdbIsForeignVersion(rdbver) && type >= RDB_FOREIGN_TYPE_MIN) {
        addReplyErrorFormat(c, "Unsupported foreign data type: %d", type);
        return;
    }

    obj = rdbLoadObject(type, &payload, objectGetVal(key), c->db->id, NULL, RDBFLAGS_NONE, 0);
    if (obj == NULL) {
        addReplyError(c, "Bad data format");
        return;
    }

    /* Remove the old key if needed. */
    int deleted = 0;
    if (replace) deleted = dbDelete(c->db, key);

    if (ttl && !absttl) ttl += commandTimeSnapshot();
    if (ttl && checkAlreadyExpired(ttl)) {
        if (deleted) {
            /* Here we don't use deleteExpiredKeyFromOverwriteAndPropagate because
             * strictly speaking, the `delete` is triggered by the `replace`. */
            robj *aux = server.lazyfree_lazy_server_del ? shared.unlink : shared.del;
            rewriteClientCommandVector(c, 2, aux, key);
            signalModifiedKey(c, c->db, key);
            notifyKeyspaceEvent(NOTIFY_GENERIC, "del", key, c->db->id);
            server.dirty++;
        }
        decrRefCount(obj);
        addReply(c, shared.ok);
        return;
    }

    /* Create the key and set the TTL if any */
    dbAdd(c->db, key, &obj);
    if (ttl) {
        obj = setExpire(c, c->db, key, ttl);
        if (!absttl) {
            /* Propagate TTL as absolute timestamp */
            robj *ttl_obj = createStringObjectFromLongLong(ttl);
            rewriteClientCommandArgument(c, 2, ttl_obj);
            decrRefCount(ttl_obj);
            rewriteClientCommandArgument(c, c->argc, shared.absttl);
        }
    }
    objectSetLRUOrLFU(obj, lfu_freq, lru_idle);
    signalModifiedKey(c, c->db, key);
    notifyKeyspaceEvent(NOTIFY_GENERIC, "restore", key, c->db->id);
    addReply(c, shared.ok);
    server.dirty++;
}
/* MIGRATE socket cache implementation.
 *
 * We take a map between host:ip and a TCP socket that we used to connect
 * to this instance in recent time.
 * This sockets are closed when the max number we cache is reached, and also
 * in serverCron() when they are around for more than a few seconds. */
#define MIGRATE_SOCKET_CACHE_ITEMS 64 /* max num of items in the cache. */
#define MIGRATE_SOCKET_CACHE_TTL 10   /* close cached sockets after 10 sec. */

typedef struct migrateCachedSocket {
    connection *conn;
    long last_dbid;
    time_t last_use_time;
} migrateCachedSocket;

/* Return a migrateCachedSocket containing a TCP socket connected with the
 * target instance, possibly returning a cached one.
 *
 * This function is responsible of sending errors to the client if a
 * connection can't be established. In this case -1 is returned.
 * Otherwise, on success the socket is returned, and the caller should not
 * attempt to free it after usage.
 *
 * If the caller detects an error while using the socket, migrateCloseSocket()
 * should be called so that the connection will be created from scratch
 * the next time. */
migrateCachedSocket *migrateGetSocket(client *c, robj *host, robj *port, long timeout) {
    connection *conn;
    sds name = sdsempty();
    migrateCachedSocket *cs;

    /* Check if we have an already cached socket for this ip:port pair. */
    name = sdscatlen(name, objectGetVal(host), sdslen(objectGetVal(host)));
    name = sdscatlen(name, ":", 1);
    name = sdscatlen(name, objectGetVal(port), sdslen(objectGetVal(port)));
    cs = dictFetchValue(server.migrate_cached_sockets, name);
    if (cs) {
        sdsfree(name);
        cs->last_use_time = server.unixtime;
        return cs;
    }

    /* No cached socket, create one. */
    if (dictSize(server.migrate_cached_sockets) == MIGRATE_SOCKET_CACHE_ITEMS) {
        /* Too many items, drop one at random. */
        dictEntry *de = dictGetRandomKey(server.migrate_cached_sockets);
        cs = dictGetVal(de);
        connClose(cs->conn);
        zfree(cs);
        dictDelete(server.migrate_cached_sockets, dictGetKey(de));
    }

    /* Create the connection and tag as high-priority so key/slot migration
     * packets are not delayed by normal tenant commands. */
    conn = connCreate(connTypeOfCluster());
    connSetPriority(conn, true);
    if (connBlockingConnect(conn, objectGetVal(host), atoi(objectGetVal(port)), timeout) != C_OK) {
        addReplyError(c, "-IOERR error or timeout connecting to the client");
        connClose(conn);
        sdsfree(name);
        return NULL;
    }

    /* Add to the cache and return it to the caller. */
    cs = zmalloc(sizeof(*cs));
    cs->conn = conn;

    cs->last_dbid = -1;
    cs->last_use_time = server.unixtime;
    dictAdd(server.migrate_cached_sockets, name, cs);
    return cs;
}

/* Free a migrate cached connection. */
void migrateCloseSocket(robj *host, robj *port) {
    sds name = sdsempty();
    migrateCachedSocket *cs;

    name = sdscatlen(name, objectGetVal(host), sdslen(objectGetVal(host)));
    name = sdscatlen(name, ":", 1);
    name = sdscatlen(name, objectGetVal(port), sdslen(objectGetVal(port)));
    cs = dictFetchValue(server.migrate_cached_sockets, name);
    if (!cs) {
        sdsfree(name);
        return;
    }

    connClose(cs->conn);
    zfree(cs);
    dictDelete(server.migrate_cached_sockets, name);
    sdsfree(name);
}

void migrateCloseTimedoutSockets(void) {
    dictIterator *di = dictGetSafeIterator(server.migrate_cached_sockets);
    dictEntry *de;

    while ((de = dictNext(di)) != NULL) {
        migrateCachedSocket *cs = dictGetVal(de);

        if ((server.unixtime - cs->last_use_time) > MIGRATE_SOCKET_CACHE_TTL) {
            connClose(cs->conn);
            zfree(cs);
            dictDelete(server.migrate_cached_sockets, dictGetKey(de));
        }
    }
    dictReleaseIterator(di);
}

/* MIGRATE host port key dbid timeout [COPY | REPLACE | AUTH password |
 *         AUTH2 username password]
 *
 * On in the multiple keys form:
 *
 * MIGRATE host port "" dbid timeout [COPY | REPLACE | AUTH password |
 *         AUTH2 username password] KEYS key1 key2 ... keyN */
void migrateCommand(client *c) {
    migrateCachedSocket *cs;
    int copy = 0, replace = 0, j;
    char *username = NULL;
    char *password = NULL;
    long timeout;
    long dbid;
    robj **ov = NULL;      /* Objects to migrate. */
    robj **kv = NULL;      /* Key names. */
    robj **newargv = NULL; /* Used to rewrite the command as DEL ... keys ... */
    rio cmd, payload;
    int may_retry = 1;
    int write_error = 0;
    int argv_rewritten = 0;
    int errno_copy = 0;

    /* To support the KEYS option we need the following additional state. */
    int first_key = 3; /* Argument index of the first key. */
    int num_keys = 1;  /* By default only migrate the 'key' argument. */

    /* Parse additional options */
    for (j = 6; j < c->argc; j++) {
        int moreargs = (c->argc - 1) - j;
        if (!strcasecmp(objectGetVal(c->argv[j]), "copy")) {
            copy = 1;
        } else if (!strcasecmp(objectGetVal(c->argv[j]), "replace")) {
            replace = 1;
        } else if (!strcasecmp(objectGetVal(c->argv[j]), "auth")) {
            if (!moreargs) {
                addReplyErrorObject(c, shared.syntaxerr);
                return;
            }
            j++;
            password = objectGetVal(c->argv[j]);
            redactClientCommandArgument(c, j);
        } else if (!strcasecmp(objectGetVal(c->argv[j]), "auth2")) {
            if (moreargs < 2) {
                addReplyErrorObject(c, shared.syntaxerr);
                return;
            }
            username = objectGetVal(c->argv[++j]);
            redactClientCommandArgument(c, j);
            password = objectGetVal(c->argv[++j]);
            redactClientCommandArgument(c, j);
        } else if (!strcasecmp(objectGetVal(c->argv[j]), "keys")) {
            if (sdslen(objectGetVal(c->argv[3])) != 0) {
                addReplyError(c, "When using MIGRATE KEYS option, the key argument"
                                 " must be set to the empty string");
                return;
            }
            first_key = j + 1;
            num_keys = c->argc - j - 1;
            break; /* All the remaining args are keys. */
        } else {
            addReplyErrorObject(c, shared.syntaxerr);
            return;
        }
    }

    /* Sanity check */
    if (getLongFromObjectOrReply(c, c->argv[5], &timeout, NULL) != C_OK ||
        getLongFromObjectOrReply(c, c->argv[4], &dbid, NULL) != C_OK) {
        return;
    }
    if (timeout <= 0) timeout = 1000;

    /* Check if the keys are here. If at least one key is to migrate, do it
     * otherwise if all the keys are missing reply with "NOKEY" to signal
     * the caller there was nothing to migrate. We don't return an error in
     * this case, since often this is due to a normal condition like the key
     * expiring in the meantime. */
    ov = zrealloc(ov, sizeof(robj *) * num_keys);
    kv = zrealloc(kv, sizeof(robj *) * num_keys);
    int oi = 0;

    for (j = 0; j < num_keys; j++) {
        if ((ov[oi] = lookupKeyRead(c->db, c->argv[first_key + j])) != NULL) {
            kv[oi] = c->argv[first_key + j];
            oi++;
        }
    }
    num_keys = oi;
    if (num_keys == 0) {
        zfree(ov);
        zfree(kv);
        addReplySds(c, sdsnew("+NOKEY\r\n"));
        return;
    }

try_again:
    write_error = 0;

    /* Connect */
    cs = migrateGetSocket(c, c->argv[1], c->argv[2], timeout);
    if (cs == NULL) {
        zfree(ov);
        zfree(kv);
        return; /* error sent to the client by migrateGetSocket() */
    }

    rioInitWithBuffer(&cmd, sdsempty());

    /* Authentication */
    if (password) {
        int arity = username ? 3 : 2;
        serverAssertWithInfo(c, NULL, rioWriteBulkCount(&cmd, '*', arity));
        serverAssertWithInfo(c, NULL, rioWriteBulkString(&cmd, "AUTH", 4));
        if (username) {
            serverAssertWithInfo(c, NULL, rioWriteBulkString(&cmd, username, sdslen(username)));
        }
        serverAssertWithInfo(c, NULL, rioWriteBulkString(&cmd, password, sdslen(password)));
    }

    /* Send the SELECT command if the current DB is not already selected. */
    int select = cs->last_dbid != dbid; /* Should we emit SELECT? */
    if (select) {
        serverAssertWithInfo(c, NULL, rioWriteBulkCount(&cmd, '*', 2));
        serverAssertWithInfo(c, NULL, rioWriteBulkString(&cmd, "SELECT", 6));
        serverAssertWithInfo(c, NULL, rioWriteBulkLongLong(&cmd, dbid));
    }

    int non_expired = 0; /* Number of keys that we'll find non expired.
                            Note that serializing large keys may take some time
                            so certain keys that were found non expired by the
                            lookupKey() function, may be expired later. */

    /* Create RESTORE payload and generate the protocol to call the command. */
    for (j = 0; j < num_keys; j++) {
        long long ttl = 0;
        long long expireat = objectGetExpire(ov[j]);

        if (expireat != -1) {
            ttl = expireat - commandTimeSnapshot();
            if (ttl < 0) {
                continue;
            }
            if (ttl < 1) ttl = 1;
        }

        /* Relocate valid (non expired) keys and values into the array in successive
         * positions to remove holes created by the keys that were present
         * in the first lookup but are now expired after the second lookup. */
        ov[non_expired] = ov[j];
        kv[non_expired++] = kv[j];

        serverAssertWithInfo(c, NULL, rioWriteBulkCount(&cmd, '*', replace ? 5 : 4));

        if (server.cluster_enabled)
            serverAssertWithInfo(c, NULL, rioWriteBulkString(&cmd, "RESTORE-ASKING", 14));
        else
            serverAssertWithInfo(c, NULL, rioWriteBulkString(&cmd, "RESTORE", 7));
        serverAssertWithInfo(c, NULL, sdsEncodedObject(kv[j]));
        serverAssertWithInfo(c, NULL, rioWriteBulkString(&cmd, objectGetVal(kv[j]), sdslen(objectGetVal(kv[j]))));
        serverAssertWithInfo(c, NULL, rioWriteBulkLongLong(&cmd, ttl));

        /* Emit the payload argument, that is the serialized object using
         * the DUMP format. */
        createDumpPayload(&payload, ov[j], kv[j], dbid);
        serverAssertWithInfo(c, NULL, rioWriteBulkString(&cmd, payload.io.buffer.ptr, sdslen(payload.io.buffer.ptr)));
        sdsfree(payload.io.buffer.ptr);

        /* Add the REPLACE option to the RESTORE command if it was specified
         * as a MIGRATE option. */
        if (replace) serverAssertWithInfo(c, NULL, rioWriteBulkString(&cmd, "REPLACE", 7));
    }

    /* Fix the actual number of keys we are migrating. */
    num_keys = non_expired;

    /* Transfer the query to the other node in 64K chunks. */
    errno = 0;
    {
        sds buf = cmd.io.buffer.ptr;
        size_t pos = 0, towrite;
        int nwritten = 0;

        while ((towrite = sdslen(buf) - pos) > 0) {
            towrite = (towrite > (64 * 1024) ? (64 * 1024) : towrite);
            nwritten = connSyncWrite(cs->conn, buf + pos, towrite, timeout);
            if (nwritten != (signed)towrite) {
                write_error = 1;
                goto socket_err;
            }
            pos += nwritten;
        }
    }

    char buf0[1024]; /* Auth reply. */
    char buf1[1024]; /* Select reply. */
    char buf2[1024]; /* Restore reply. */

    /* Read the AUTH reply if needed. */
    if (password && connSyncReadLine(cs->conn, buf0, sizeof(buf0), timeout) <= 0) goto socket_err;

    /* Read the SELECT reply if needed. */
    if (select && connSyncReadLine(cs->conn, buf1, sizeof(buf1), timeout) <= 0) goto socket_err;

    /* Read the RESTORE replies. */
    int error_from_target = 0;
    int socket_error = 0;
    int del_idx = 1; /* Index of the key argument for the replicated DEL op. */

    /* Allocate the new argument vector that will replace the current command,
     * to propagate the MIGRATE as a DEL command (if no COPY option was given).
     * We allocate num_keys+1 because the additional argument is for "DEL"
     * command name itself. */
    if (!copy) newargv = zmalloc(sizeof(robj *) * (num_keys + 1));

    for (j = 0; j < num_keys; j++) {
        if (connSyncReadLine(cs->conn, buf2, sizeof(buf2), timeout) <= 0) {
            socket_error = 1;
            break;
        }
        if ((password && buf0[0] == '-') || (select && buf1[0] == '-') || buf2[0] == '-') {
            /* On error assume that last_dbid is no longer valid. */
            if (!error_from_target) {
                cs->last_dbid = -1;
                char *errbuf;
                if (password && buf0[0] == '-')
                    errbuf = buf0;
                else if (select && buf1[0] == '-')
                    errbuf = buf1;
                else
                    errbuf = buf2;

                error_from_target = 1;
                addReplyErrorFormat(c, "Target instance replied with error: %s", errbuf + 1);
            }
        } else {
            if (!copy) {
                /* No COPY option: remove the local key, signal the change. */
                dbDelete(c->db, kv[j]);
                signalModifiedKey(c, c->db, kv[j]);
                notifyKeyspaceEvent(NOTIFY_GENERIC, "del", kv[j], c->db->id);
                server.dirty++;

                /* Populate the argument vector to replace the old one. */
                newargv[del_idx++] = kv[j];
                incrRefCount(kv[j]);
            }
        }
    }

    /* On socket error, if we want to retry, do it now before rewriting the
     * command vector. We only retry if we are sure nothing was processed
     * and we failed to read the first reply (j == 0 test). */
    if (!error_from_target && socket_error && j == 0 && may_retry && errno != ETIMEDOUT) {
        goto socket_err; /* A retry is guaranteed because of tested conditions.*/
    }

    /* On socket errors, close the migration socket now that we still have
     * the original host/port in the ARGV. Later the original command may be
     * rewritten to DEL and will be too later. */
    if (socket_error) migrateCloseSocket(c->argv[1], c->argv[2]);

    if (!copy) {
        /* Translate MIGRATE as DEL for replication/AOF. Note that we do
         * this only for the keys for which we received an acknowledgement
         * from the receiving server, by using the del_idx index. */
        if (del_idx > 1) {
            newargv[0] = createStringObject("DEL", 3);
            /* Note that the following call takes ownership of newargv. */
            replaceClientCommandVector(c, del_idx, newargv);
            argv_rewritten = 1;
        } else {
            /* No key transfer acknowledged, no need to rewrite as DEL. */
            zfree(newargv);
        }
        newargv = NULL; /* Make it safe to call zfree() on it in the future. */
    }

    /* If we are here and a socket error happened, we don't want to retry.
     * Just signal the problem to the client, but only do it if we did not
     * already queue a different error reported by the destination server. */
    if (!error_from_target && socket_error) {
        may_retry = 0;
        goto socket_err;
    }

    if (!error_from_target) {
        /* Success! Update the last_dbid in migrateCachedSocket, so that we can
         * avoid SELECT the next time if the target DB is the same. Reply +OK.
         *
         * Note: If we reached this point, even if socket_error is true
         * still the SELECT command succeeded (otherwise the code jumps to
         * socket_err label. */
        cs->last_dbid = dbid;
        addReply(c, shared.ok);
    } else {
        /* On error we already sent it in the for loop above, and set
         * the currently selected socket to -1 to force SELECT the next time. */
    }

    sdsfree(cmd.io.buffer.ptr);
    zfree(ov);
    zfree(kv);
    zfree(newargv);
    return;

    /* On socket errors we try to close the cached socket and try again.
     * It is very common for the cached socket to get closed, if just reopening
     * it works it's a shame to notify the error to the caller. */
socket_err:
    /* Take a copy of 'errno' prior cleanup as it can be overwritten and
     * use copied variable for re-try check. */
    errno_copy = errno;

    /* Cleanup we want to perform in both the retry and no retry case.
     * Note: Closing the migrate socket will also force SELECT next time. */
    sdsfree(cmd.io.buffer.ptr);

    /* If the command was rewritten as DEL and there was a socket error,
     * we already closed the socket earlier. While migrateCloseSocket()
     * is idempotent, the host/port arguments are now gone, so don't do it
     * again. */
    if (!argv_rewritten) migrateCloseSocket(c->argv[1], c->argv[2]);
    zfree(newargv);
    newargv = NULL; /* This will get reallocated on retry. */

    /* Retry only if it's not a timeout and we never attempted a retry
     * (or the code jumping here did not set may_retry to zero). */
    if (errno_copy != ETIMEDOUT && may_retry) {
        may_retry = 0;
        goto try_again;
    }

    /* Cleanup we want to do if no retry is attempted. */
    zfree(ov);
    zfree(kv);
    addReplyErrorSds(c, sdscatprintf(sdsempty(), "-IOERR error or timeout %s to target instance",
                                     write_error ? "writing" : "reading"));
    return;
}

/* Cluster node sanity check. Returns C_OK if the node id
 * is valid an C_ERR otherwise. */
int verifyClusterNodeId(const char *name, int length) {
    if (length != CLUSTER_NAMELEN) return C_ERR;
    for (int i = 0; i < length; i++) {
        if (name[i] >= 'a' && name[i] <= 'z') continue;
        if (name[i] >= '0' && name[i] <= '9') continue;
        return C_ERR;
    }
    return C_OK;
}

static int isValidAuxChar(unsigned char c) {
    /* Reject everything up through ',' (0x2C) inclusive: control characters
     * (0x00-0x1F), space !"#$%&'()*+, (0x20-0x2C), and DEL (0x7F). */
    if (c <= ',' || c == 0x7F) return 0;

    /* Reject additional characters above 0x2C (comma) that are format-significant in
     * nodes.conf or otherwise unsafe. */
    static const char *invalid_charset = ";<=>?@[]^{|}~\\";

    return strchr(invalid_charset, c) == NULL;
}

int isValidAuxString(char *s, unsigned int length) {
    for (unsigned i = 0; i < length; i++) {
        if (!isValidAuxChar(s[i])) return 0;
    }
    return 1;
}

void clusterCommandMyId(client *c) {
    char *name = clusterNodeGetName(getMyClusterNode());
    if (name) {
        addReplyBulkCBuffer(c, name, CLUSTER_NAMELEN);
    } else {
        addReplyError(c, "No ID yet");
    }
}

int clusterNodeIsMyself(clusterNode *n) {
    return n == getMyClusterNode();
}

void clusterCommandMyShardId(client *c) {
    char *sid = clusterNodeGetShardId(getMyClusterNode());
    if (sid) {
        addReplyBulkCBuffer(c, sid, CLUSTER_NAMELEN);
    } else {
        addReplyError(c, "No shard ID yet");
    }
}

/* When a cluster command is called, we need to decide whether to return TLS info or
 * non-TLS info by the client's connection type. However if the command is called by
 * a Lua script or RM_call, there is no connection in the fake client, so we use
 * server.current_client here to get the real client if available. And if it is not
 * available (modules may call commands without a real client), we return the default
 * info, which is determined by server.tls_cluster. */
static int shouldReturnTlsInfo(void) {
    if (server.current_client && server.current_client->conn) {
        return connIsTLS(server.current_client->conn);
    } else {
        return server.tls_cluster;
    }
}

unsigned int countKeysInSlotForDb(unsigned int hashslot, serverDb *db) {
    return kvstoreHashtableSize(db->keys, hashslot);
}

unsigned int countKeysInSlot(unsigned int slot) {
    unsigned int result = 0;
    for (int i = 0; i < server.dbnum; i++) {
        result += server.db[i] ? countKeysInSlotForDb(slot, server.db[i]) : 0;
    }
    return result;
}

void clusterCommandHelp(client *c) {
    const char *help[] = {
        "COUNTKEYSINSLOT <slot>",
        "    Return the number of keys in <slot>.",
        "GETKEYSINSLOT <slot> <count>",
        "    Return key names stored by current node in a slot.",
        "INFO",
        "    Return information about the cluster.",
        "KEYSLOT <key>",
        "    Return the hash slot for <key>.",
        "MYID",
        "    Return the node id.",
        "MYSHARDID",
        "    Return the node's shard id.",
        "NODES",
        "    Return cluster configuration seen by node. Output format:",
        "    <id> <ip:port@bus-port[,hostname]> <flags> <primary> <pings> <pongs> <epoch> <link> <slot> ...",
        "REPLICAS <node-id>",
        "    Return <node-id> replicas.",
        "SLOTS",
        "    Return information about slots range mappings. Each range is made of:",
        "    start, end, primary and replicas IP addresses, ports and ids",
        "SLOT-STATS",
        "    Return an array of slot usage statistics for slots assigned to the current node.",
        "SHARDS",
        "    Return information about slot range mappings and the nodes associated with them.",
        NULL};

    addExtendedReplyHelp(c, help, clusterCommandExtendedHelp());
}

void clusterKeySlotCommand(client *c) {
    sds key = objectGetVal(c->argv[2]);
    addReplyLongLong(c, keyHashSlot(key, sdslen(key)));
}

void clusterCommand(client *c) {
    if (server.cluster_enabled == 0) {
        addReplyError(c, "This instance has cluster support disabled");
        return;
    }

    if (c->argc == 2 && !strcasecmp(objectGetVal(c->argv[1]), "help")) {
        clusterCommandHelp(c);
    } else if (!strcasecmp(objectGetVal(c->argv[1]), "nodes") && c->argc == 2) {
        /* CLUSTER NODES */
        /* Report TLS ports to TLS client, and report non-TLS port to non-TLS client. */
        sds nodes = clusterGenNodesDescription(c, 0, shouldReturnTlsInfo());
        addReplyVerbatim(c, nodes, sdslen(nodes), "txt");
        sdsfree(nodes);
    } else if (!strcasecmp(objectGetVal(c->argv[1]), "myid") && c->argc == 2) {
        /* CLUSTER MYID */
        clusterCommandMyId(c);
    } else if (!strcasecmp(objectGetVal(c->argv[1]), "myshardid") && c->argc == 2) {
        /* CLUSTER MYSHARDID */
        clusterCommandMyShardId(c);
    } else if (!strcasecmp(objectGetVal(c->argv[1]), "slots") && c->argc == 2) {
        /* CLUSTER SLOTS */
        clusterCommandSlots(c);
    } else if (!strcasecmp(objectGetVal(c->argv[1]), "shards") && c->argc == 2) {
        /* CLUSTER SHARDS */
        clusterCommandShards(c);
    } else if (!strcasecmp(objectGetVal(c->argv[1]), "info") && c->argc == 2) {
        /* CLUSTER INFO */

        sds info = genClusterInfoString(sdsempty());

        /* Produce the reply protocol. */
        addReplyVerbatim(c, info, sdslen(info), "txt");
        sdsfree(info);
    } else if (!strcasecmp(objectGetVal(c->argv[1]), "countkeysinslot") && c->argc == 3) {
        /* CLUSTER COUNTKEYSINSLOT <slot> */
        long long slot;

        if (getLongLongFromObjectOrReply(c, c->argv[2], &slot, NULL) != C_OK) return;
        if (slot < 0 || slot >= CLUSTER_SLOTS) {
            addReplyError(c, "Invalid slot");
            return;
        }
        addReplyLongLong(c, countKeysInSlotForDb(slot, c->db));
    } else if (!strcasecmp(objectGetVal(c->argv[1]), "getkeysinslot") && c->argc == 4) {
        /* CLUSTER GETKEYSINSLOT <slot> <count> */
        long long maxkeys, slot;

        if (getLongLongFromObjectOrReply(c, c->argv[2], &slot, NULL) != C_OK) return;
        if (getLongLongFromObjectOrReply(c, c->argv[3], &maxkeys, NULL) != C_OK) return;
        if (slot < 0 || slot >= CLUSTER_SLOTS || maxkeys < 0) {
            addReplyError(c, "Invalid slot or number of keys");
            return;
        }

        unsigned int keys_in_slot = countKeysInSlotForDb(slot, c->db);
        unsigned int numkeys = maxkeys > keys_in_slot ? keys_in_slot : maxkeys;
        addReplyArrayLen(c, numkeys);
        kvstoreHashtableIterator *kvs_di = NULL;
        kvs_di = kvstoreGetHashtableIterator(c->db->keys, slot, 0);
        for (unsigned int i = 0; i < numkeys; i++) {
            void *next;
            serverAssert(kvstoreHashtableIteratorNext(kvs_di, &next));
            robj *valkey = next;
            sds sdskey = objectGetKey(valkey);
            addReplyBulkCBuffer(c, sdskey, sdslen(sdskey));
        }
        kvstoreReleaseHashtableIterator(kvs_di);
    } else if ((!strcasecmp(objectGetVal(c->argv[1]), "slaves") || !strcasecmp(objectGetVal(c->argv[1]), "replicas")) && c->argc == 3) {
        /* CLUSTER REPLICAS <NODE ID> */
        clusterNode *n = clusterLookupNode(objectGetVal(c->argv[2]), sdslen(objectGetVal(c->argv[2])));
        int j;

        /* Lookup the specified node in our table. */
        if (!n) {
            addReplyErrorFormat(c, "Unknown node %s", (char *)objectGetVal(c->argv[2]));
            return;
        }

        if (clusterNodeIsReplica(n)) {
            addReplyError(c, "The specified node is not a master");
            return;
        }

        /* Report TLS ports to TLS client, and report non-TLS port to non-TLS client. */
        addReplyArrayLen(c, clusterNodeNumReplicas(n));
        for (j = 0; j < clusterNodeNumReplicas(n); j++) {
            sds ni = clusterGenNodeDescription(c, clusterNodeGetReplica(n, j), shouldReturnTlsInfo());
            addReplyBulkCString(c, ni);
            sdsfree(ni);
        }
    } else if (!clusterCommandSpecial(c)) {
        addReplySubcommandSyntaxError(c);
        return;
    }
}

/* Compute cluster slot for the given command and arguments and detect cluster
 * cross-slot errors. Returns the slot. If -1 is returned, one of the flags
 * READ_FLAGS_NO_KEYS or READ_FLAGS_CROSSSLOT is set in `*read_flags` to
 * indicate why the slot couldn't be determined. This function has no
 * side-effects and can be called from I/O threads. */
int clusterSlotByCommand(struct serverCommand *cmd, robj **argv, int argc, int *read_flags) {
    getKeysResult result;
    initGetKeysResult(&result);
    int numkeys = getKeysFromCommand(cmd, argv, argc, &result);
    int slot = -1;
    if (numkeys == 0) *read_flags |= READ_FLAGS_NO_KEYS;
    for (int i = 0; i < numkeys; i++) {
        sds key = objectGetVal(argv[result.keys[i].pos]);
        int keyslot = keyHashSlot(key, sdslen(key));
        if (slot == -1) {
            slot = keyslot;
        } else if (keyslot != slot) {
            slot = -1;
            *read_flags |= READ_FLAGS_CROSSSLOT;
            break;
        }
    }
    getKeysFreeResult(&result);
    return slot;
}

/* Return the pointer to the cluster node that is able to serve the command.
 *
 * Note that this function doesn't compute the slot for each key. The client's
 * slot needs to be set before calling this function. If it's set to -1, the
 * client's read flags should indicate if it's a cross-slot command or if the
 * command has no keys. This is computed by clusterSlotByCommand(). The commands
 * can be called together like this:
 *
 *     c->slot = clusterSlotByCommand(c->cmd, c->argv, c->argc, &c->read_flags);
 *     node = getNodeByQuery(c, &error_code);
 *
 * For the function to succeed the command should only target either:
 *
 * 1) A single key (even multiple times like RPOPLPUSH mylist mylist).
 * 2) Multiple keys in the same hash slot, while the slot is stable (no
 *    resharding in progress).
 *
 * The EXEC command is a special case. It takes no keys so the slot for this
 * command is -1, but this function updates the client's slot to be the slot of
 * the complete MULTI-EXEC transaction.
 *
 * On success the function returns the node that is able to serve the request.
 * If the node is not 'myself' a redirection must be performed. The kind of
 * redirection is specified setting the integer passed by reference
 * 'error_code', which will be set to CLUSTER_REDIR_ASK or
 * CLUSTER_REDIR_MOVED.
 *
 * When the node is 'myself' 'error_code' is set to CLUSTER_REDIR_NONE.
 *
 * If the command fails NULL is returned, and the reason of the failure is
 * provided via 'error_code', which will be set to:
 *
 * CLUSTER_REDIR_CROSS_SLOT if the request contains multiple keys that
 * don't belong to the same hash slot.
 *
 * CLUSTER_REDIR_UNSTABLE if the request contains multiple keys
 * belonging to the same slot, but the slot is not stable (in migration or
 * importing state, likely because a resharding is in progress).
 *
 * CLUSTER_REDIR_DOWN_UNBOUND if the request addresses a slot which is
 * not bound to any node. In this case the cluster global state should be
 * already "down" but it is fragile to rely on the update of the global state,
 * so we also handle it here.
 *
 * CLUSTER_REDIR_DOWN_STATE and CLUSTER_REDIR_DOWN_RO_STATE if the cluster is
 * down but the user attempts to execute a command that addresses one or more keys. */
clusterNode *getNodeByQuery(client *c, int *error_code) {
    clusterNode *myself = getMyClusterNode();
    clusterNode *n = NULL;
    robj *firstkey = NULL;
    int multiple_keys = 0;
    multiState *ms, _ms;
    multiCmd mc;
    int i, migrating_slot = 0, importing_slot = 0, missing_keys = 0, existing_keys = 0;

    /* Slot must be calculated in advance, or cross-slot or no keys detected. */
    serverAssert(c->slot >= 0 || c->read_flags & (READ_FLAGS_NO_KEYS | READ_FLAGS_CROSSSLOT));

    /* Allow any key to be set if a module disabled cluster redirections. */
    if (server.cluster_module_flags & CLUSTER_MODULE_FLAG_NO_REDIRECTION) return myself;

    /* Set error code optimistically for the base case. */
    if (error_code) *error_code = CLUSTER_REDIR_NONE;

    /* Modules can turn off Cluster redirection: this is useful
     * when writing a module that implements a completely different
     * distributed system. */

    /* Determine transaction slot and return early on cross-slot. */
    if (c->cmd->proc == execCommand) {
        if (!c->flag.multi || c->flag.dirty_exec) return myself;

        int slot = c->slot;
        if (c->read_flags & READ_FLAGS_CROSSSLOT) {
            if (error_code) *error_code = CLUSTER_REDIR_CROSS_SLOT;
            return NULL;
        }
        for (i = 0; i < c->mstate->count; i++) {
            if (slot == -1) {
                slot = c->mstate->commands[i].slot;
            } else if (c->mstate->commands[i].slot != -1 && c->mstate->commands[i].slot != slot) {
                if (error_code) *error_code = CLUSTER_REDIR_CROSS_SLOT;
                return NULL;
            }
        }
        /* EXEC will execute all queued commands in the transaction, so we
         * overwrite the EXEC commands's slot with the transaction's slot. */
        c->slot = slot;
    } else if (c->read_flags & READ_FLAGS_CROSSSLOT) {
        if (error_code) *error_code = CLUSTER_REDIR_CROSS_SLOT;
        return NULL;
    }

    /* No key at all in command? then we can serve the request
     * without redirections or errors in all the cases. */
    if (c->slot == -1) return myself;

    n = getNodeBySlot(c->slot);

    /* If a slot is not served, we are in "cluster down" state.
     * This check is done early to preserve historical behavior. */
    if (n == NULL) {
        if (error_code) *error_code = CLUSTER_REDIR_DOWN_UNBOUND;
        return NULL;
    }

    /* If we are migrating or importing this slot, we need to check
     * if we have all the keys in the request (the only way we
     * can safely serve the request, otherwise we return a TRYAGAIN
     * error). To do so we set the importing/migrating state and
     * increment a counter for every missing key. */
    if (clusterNodeIsPrimary(myself) || c->flag.readonly) {
        if (n == clusterNodeGetPrimary(myself) && getMigratingSlotDest(c->slot) != NULL) {
            migrating_slot = 1;
        } else if (getImportingSlotSource(c->slot) != NULL) {
            importing_slot = 1;
        }
    }

    uint64_t cmd_flags = getCommandFlags(c);

    /* Only valid for sharded pubsub as regular pubsub can operate on any node and bypasses this layer. */
    int pubsubshard_included =
        (cmd_flags & CMD_PUBSUB) || (c->cmd->proc == execCommand && (c->mstate->cmd_flags & CMD_PUBSUB));

    /* If we're importing or migrating the slot, we need to do some more checks:
     *
     *   1. Go over all the keys to count existing keys and missing keys that we
     *      need for TRYAGAIN and ASK redirects.
     *   2. Check for some commands that are forbidden during slot migration.
     *
     * Skip this if we're not importing or migrating this slot. */
    if (!migrating_slot && !importing_slot) goto after_checking_each_key;

    /* We handle all the cases as if they were EXEC commands, so we have
     * a common code path for everything */
    if (c->cmd->proc == execCommand) {
        ms = c->mstate;
    } else {
        /* In order to have a single codepath create a fake Multi State
         * structure if the client is not in MULTI/EXEC state, this way
         * we have a single codepath below. */
        ms = &_ms;
        _ms.commands = &mc;
        _ms.count = 1;
        mc.argv = c->argv;
        mc.argc = c->argc;
        mc.cmd = c->cmd;
    }

    serverDb *origDb = c->db;
    serverDb *currentDb = origDb;

    /* Check for multiple keys, existing keys, missing keys. */
    for (i = c->cmd->proc == execCommand ? -1 : 0; i < ms->count; i++) {
        struct serverCommand *mcmd;
        robj **margv;
        int margc, numkeys, j;
        keyReference *keyindex;

        if (i == -1) {
            mcmd = c->cmd;
            margc = c->argc;
            margv = c->argv;
        } else {
            mcmd = ms->commands[i].cmd;
            margc = ms->commands[i].argc;
            margv = ms->commands[i].argv;
        }

        getKeysResult result;
        initGetKeysResult(&result);
        numkeys = getKeysFromCommand(mcmd, margv, margc, &result);
        keyindex = result.keys;

        if (mcmd->proc == selectCommand) {
            /* Failed SELECT is ignored since it doesn't modify the database. */
            long long id;
            if (getLongLongFromObject(margv[1], &id) == C_OK && selectDb(c, id) == C_OK) {
                currentDb = c->db;
                selectDb(c, origDb->id);
            }
        }

        for (j = 0; j < numkeys; j++) {
            robj *thiskey = margv[keyindex[j].pos];

            if (firstkey == NULL) {
                /* This is the first key we see. */
                firstkey = thiskey;
            } else {
                /* If it is not the first key/channel, make sure it is exactly
                 * the same key/channel as the first we saw. */
                if (importing_slot && !multiple_keys && !equalStringObjects(firstkey, thiskey)) {
                    /* Flag this request as one with multiple different
                     * keys/channels when the slot is in importing state. */
                    multiple_keys = 1;
                }
            }

            /* Block MOVE command as the destination key is not expected to exist, and we don't know if it was migrated */
            if (mcmd->proc == moveCommand) {
                if (error_code) *error_code = CLUSTER_REDIR_UNSTABLE;
                getKeysFreeResult(&result);
                return NULL;
            }

            /* Block the COPY command if it's cross-DB to keep the code simple.
             * Allowing cross-DB COPY is possible, but it would require looking up the second key in the target DB.
             * The command should only be allowed if the key exists. We may revisit this decision in the future.
             *
             * The DB and REPLACE tokens are optional and order-independent, and DB may appear more than once, so
             * scan every DB clause the way copyCommand does. A DB token without a value is left to copyCommand to
             * reject as a syntax error. */
            if (mcmd->proc == copyCommand) {
                for (int k = 3; k + 1 < margc; k++) {
                    if (strcasecmp(objectGetVal(margv[k]), "db")) continue;
                    long long value;
                    if (getLongLongFromObject(margv[k + 1], &value) != C_OK || value != currentDb->id) {
                        if (error_code) *error_code = CLUSTER_REDIR_UNSTABLE;
                        getKeysFreeResult(&result);
                        return NULL;
                    }
                    k++; /* Consume the dbid. */
                }
            }

            /* Migrating / Importing slot? During exec we count keys we don't have.
             * If it is pubsubshard command, it isn't required to check
             * the channel being present or not in the node during the
             * slot migration, the channel will be served from the source
             * node until the migration completes with CLUSTER SETSLOT <slot>
             * NODE <node-id>. */
            int flags = LOOKUP_NOEFFECTS; /* not client key access */
            if (!pubsubshard_included &&
                (!c->flag.multi || (c->flag.multi && c->cmd->proc == execCommand))) {
                /* Multi/Exec validation happens on exec */
                if (lookupKeyReadWithFlags(currentDb, thiskey, flags) == NULL)
                    missing_keys++;
                else
                    existing_keys++;
            }
        }
        getKeysFreeResult(&result);
    }

after_checking_each_key:

    /* Cluster is globally down but we got keys? We only serve the request
     * if it is a read command and when allow_reads_when_down is enabled. */
    if (!isClusterHealthy()) {
        if (pubsubshard_included) {
            if (!server.cluster_allow_pubsubshard_when_down) {
                if (error_code) *error_code = CLUSTER_REDIR_DOWN_STATE;
                return NULL;
            }
        } else if (!server.cluster_allow_reads_when_down) {
            /* The cluster is configured to block commands when the
             * cluster is down. */
            if (error_code) *error_code = CLUSTER_REDIR_DOWN_STATE;
            return NULL;
        } else if (cmd_flags & CMD_WRITE) {
            /* The cluster is configured to allow read only commands */
            if (error_code) *error_code = CLUSTER_REDIR_DOWN_RO_STATE;
            return NULL;
        } else {
            /* Fall through and allow the command to be executed:
             * this happens when server.cluster_allow_reads_when_down is
             * true and the command is not a write command */
        }
    }

    /* MIGRATE always works in the context of the local node if the slot
     * is open (migrating or importing state). We need to be able to freely
     * move keys among instances in this case. */
    if ((migrating_slot || importing_slot) && c->cmd->proc == migrateCommand && clusterNodeIsPrimary(myself)) {
        return myself;
    }

    /* If we don't have all the keys and we are migrating the slot, send
     * an ASK redirection or TRYAGAIN. */
    if (migrating_slot && missing_keys) {
        /* If we have keys but we don't have all keys, we return TRYAGAIN */
        if (existing_keys) {
            if (error_code) *error_code = CLUSTER_REDIR_UNSTABLE;
            return NULL;
        } else {
            if (error_code) *error_code = CLUSTER_REDIR_ASK;
            return getMigratingSlotDest(c->slot);
        }
    }

    /* If we are receiving the slot, and the client correctly flagged the
     * request as "ASKING", we can serve the request. However if the request
     * involves multiple keys and we don't have them all, the only option is
     * to send a TRYAGAIN error. */
    if (importing_slot && (c->flag.asking || cmd_flags & CMD_ASKING)) {
        if (multiple_keys && missing_keys) {
            if (error_code) *error_code = CLUSTER_REDIR_UNSTABLE;
            return NULL;
        } else {
            return myself;
        }
    }

    /* Handle the read-only client case reading from a replica: if this
     * node is a replica and the request is about a hash slot our primary
     * is serving, we can reply without redirection. */
    int is_write_command =
        (cmd_flags & CMD_WRITE) || (c->cmd->proc == execCommand && (c->mstate->cmd_flags & CMD_WRITE));
    if ((c->flag.readonly || pubsubshard_included) && !is_write_command && clusterNodeIsReplica(myself) &&
        clusterNodeGetPrimary(myself) == n) {
        return myself;
    }

    /* Base case: just return the right node. However, if this node is not
     * myself, set error_code to MOVED since we need to issue a redirection. */
    if (n != myself && error_code) *error_code = CLUSTER_REDIR_MOVED;
    return n;
}

/* Send the client the right redirection code, according to error_code
 * that should be set to one of CLUSTER_REDIR_* macros.
 *
 * If CLUSTER_REDIR_ASK or CLUSTER_REDIR_MOVED error codes
 * are used, then the node 'n' should not be NULL, but should be the
 * node we want to mention in the redirection. Moreover hashslot should
 * be set to the hash slot that caused the redirection. */
void clusterRedirectClient(client *c, clusterNode *n, int hashslot, int error_code) {
    if (error_code == CLUSTER_REDIR_CROSS_SLOT) {
        addReplyError(c, "-CROSSSLOT Keys in request don't hash to the same slot");
    } else if (error_code == CLUSTER_REDIR_UNSTABLE) {
        /* The request spawns multiple keys in the same slot,
         * but the slot is not "stable" currently as there is
         * a migration or import in progress. */
        addReplyError(c, "-TRYAGAIN Multiple keys request during rehashing of slot");
    } else if (error_code == CLUSTER_REDIR_DOWN_STATE) {
        addReplyError(c, "-CLUSTERDOWN The cluster is down");
    } else if (error_code == CLUSTER_REDIR_DOWN_RO_STATE) {
        addReplyError(c, "-CLUSTERDOWN The cluster is down and only accepts read commands");
    } else if (error_code == CLUSTER_REDIR_DOWN_UNBOUND) {
        addReplyError(c, "-CLUSTERDOWN Hash slot not served");
    } else if (error_code == CLUSTER_REDIR_MOVED || error_code == CLUSTER_REDIR_ASK) {
        /* Report TLS ports to TLS client, and report non-TLS port to non-TLS client. */
        int port = clusterNodeClientPort(n, shouldReturnTlsInfo(), c);
        addReplyErrorSds(c,
                         sdscatprintf(sdsempty(), "-%s %d %s:%d", (error_code == CLUSTER_REDIR_ASK) ? "ASK" : "MOVED",
                                      hashslot, clusterNodePreferredEndpoint(n, c), port));
    } else {
        serverPanic("getNodeByQuery() unknown error.");
    }
}

/* This function is called by the function processing clients incrementally
 * to detect timeouts, in order to handle the following case:
 *
 * 1) A client blocks with BLPOP or similar blocking operation.
 * 2) The primary migrates the hash slot elsewhere or turns into a replica.
 * 3) The client may remain blocked forever (or up to the max timeout time)
 *    waiting for a key change that will never happen.
 *
 * If the client is found to be blocked into a hash slot this node no
 * longer handles, the client is sent a redirection error, and the function
 * returns 1. Otherwise, 0 is returned and no operation is performed. */
int clusterRedirectBlockedClientIfNeeded(client *c) {
    clusterNode *myself = getMyClusterNode();
    if (c->flag.blocked && (c->bstate->btype == BLOCKED_LIST || c->bstate->btype == BLOCKED_ZSET ||
                            c->bstate->btype == BLOCKED_STREAM || c->bstate->btype == BLOCKED_MODULE)) {
        dictEntry *de;
        dictIterator *di;

        /* If the client is blocked on module, but not on a specific key,
         * don't unblock it. */
        if (c->bstate->btype == BLOCKED_MODULE && !moduleClientIsBlockedOnKeys(c)) return 0;

        /* If the cluster is down, unblock the client with the right error.
         * If the cluster is configured to allow reads on cluster down, we
         * still want to emit this error since a write will be required
         * to unblock them which may never come.  */
        if (!isClusterHealthy()) {
            clusterRedirectClient(c, NULL, 0, CLUSTER_REDIR_DOWN_STATE);
            return 1;
        }

        /* All keys must belong to the same slot, so check first key only. */
        di = dictGetIterator(c->bstate->keys);
        if ((de = dictNext(di)) != NULL) {
            robj *key = dictGetKey(de);
            int slot = keyHashSlot((char *)objectGetVal(key), sdslen(objectGetVal(key)));
            serverAssert(slot == c->slot);
            clusterNode *node = getNodeBySlot(slot);

            /* if the client is read-only and attempting to access key that our
             * replica can handle, allow it. */
            if (c->flag.readonly && !(c->lastcmd->flags & CMD_WRITE) && clusterNodeIsReplica(myself) &&
                clusterNodeGetPrimary(myself) == node) {
                node = myself;
            }

            /* We send an error and unblock the client if:
             * 1) The slot is unassigned, emitting a cluster down error.
             * 2) The slot is neither handled by this node, nor being imported. */
            if (node != myself && getImportingSlotSource(slot) == NULL) {
                if (node == NULL) {
                    clusterRedirectClient(c, NULL, 0, CLUSTER_REDIR_DOWN_UNBOUND);
                } else {
                    clusterRedirectClient(c, node, slot, CLUSTER_REDIR_MOVED);
                }
                dictReleaseIterator(di);
                return 1;
            }
        }
        dictReleaseIterator(di);
    }
    return 0;
}

void addNodeToNodeReply(client *c, clusterNode *node) {
    char *hostname = clusterNodeHostname(node);
    addReplyArrayLen(c, 4);
    if (server.cluster_preferred_endpoint_type == CLUSTER_ENDPOINT_TYPE_IP) {
        addReplyBulkCString(c, clusterNodeIp(node, c));
    } else if (server.cluster_preferred_endpoint_type == CLUSTER_ENDPOINT_TYPE_HOSTNAME) {
        if (hostname != NULL && hostname[0] != '\0') {
            addReplyBulkCString(c, hostname);
        } else {
            addReplyBulkCString(c, "?");
        }
    } else if (server.cluster_preferred_endpoint_type == CLUSTER_ENDPOINT_TYPE_UNKNOWN_ENDPOINT) {
        addReplyNull(c);
    } else {
        serverPanic("Unrecognized preferred endpoint type");
    }

    /* Report TLS ports to TLS client, and report non-TLS port to non-TLS client. */
    addReplyLongLong(c, clusterNodeClientPort(node, shouldReturnTlsInfo(), c));
    addReplyBulkCBuffer(c, clusterNodeGetName(node), CLUSTER_NAMELEN);

    /* Add the additional endpoint information, this is all the known networking information
     * that is not the preferred endpoint. Note the logic is evaluated twice so we can
     * correctly report the number of additional network arguments without using a deferred
     * map, an assertion is made at the end to check we set the right length. */
    int length = 0;
    if (server.cluster_preferred_endpoint_type != CLUSTER_ENDPOINT_TYPE_IP) {
        length++;
    }
    if (server.cluster_preferred_endpoint_type != CLUSTER_ENDPOINT_TYPE_HOSTNAME && hostname != NULL &&
        hostname[0] != '\0') {
        length++;
    }

    if (sdslen(node->availability_zone) != 0) {
        length++;
    }

    addReplyMapLen(c, length);

    if (server.cluster_preferred_endpoint_type != CLUSTER_ENDPOINT_TYPE_IP) {
        addReplyBulkCString(c, "ip");
        addReplyBulkCString(c, clusterNodeIp(node, c));
        length--;
    }
    if (server.cluster_preferred_endpoint_type != CLUSTER_ENDPOINT_TYPE_HOSTNAME && hostname != NULL &&
        hostname[0] != '\0') {
        addReplyBulkCString(c, "hostname");
        addReplyBulkCString(c, hostname);
        length--;
    }

    if (sdslen(node->availability_zone) != 0) {
        addReplyBulkCString(c, "availability-zone");
        addReplyBulkCString(c, node->availability_zone);
        length--;
    }

    serverAssert(length == 0);
}

/* Returns an indication if the node is fully available
 * and should be listed in CLUSTER SLOTS response.
 * Returns 1 for available nodes, 0 for nodes that have
 * not finished their initial sync, in failed state, or are
 * otherwise considered not available to serve read commands. */
int isNodeAvailable(clusterNode *node) {
    /* We don't consider PFAIL here because it's not a reliable indicator
     * for node available and we don't want clients to use it. */
    if (clusterNodeIsFailing(node)) {
        return 0;
    }

    /* Hide empty replicas in here, from a data-path POV, an empty replica
     * is not available. */
    return getNodeReplicationOffset(node) != 0;
}

void addNodeReplyForClusterSlot(client *c, clusterNode *node, int start_slot, int end_slot) {
    int i, nested_elements = 3; /* slots (2) + primary addr (1) */
    for (i = 0; i < clusterNodeNumReplicas(node); i++) {
        if (!isNodeAvailable(clusterNodeGetReplica(node, i))) continue;
        nested_elements++;
    }
    addReplyArrayLen(c, nested_elements);
    addReplyLongLong(c, start_slot);
    addReplyLongLong(c, end_slot);
    addNodeToNodeReply(c, node);

    /* Remaining nodes in reply are replicas for slot range */
    for (i = 0; i < clusterNodeNumReplicas(node); i++) {
        /* This loop is copy/pasted from clusterGenNodeDescription()
         * with modifications for per-slot node aggregation. */
        if (!isNodeAvailable(clusterNodeGetReplica(node, i))) continue;
        addNodeToNodeReply(c, clusterNodeGetReplica(node, i));
        nested_elements--;
    }
    serverAssert(nested_elements == 3); /* Original 3 elements */
}

void clearCachedClusterSlotsResponse(void) {
    for (int conn_type = 0; conn_type < CACHE_CONN_TYPE_MAX; conn_type++) {
        if (server.cached_cluster_slot_info[conn_type]) {
            sdsfree(server.cached_cluster_slot_info[conn_type]);
            server.cached_cluster_slot_info[conn_type] = NULL;
        }
    }
}

sds generateClusterSlotResponse(int resp) {
    client *recording_client = createCachedResponseClient(resp);
    clusterNode *n = NULL;
    int num_primaries = 0, start = -1;
    void *slot_replylen = addReplyDeferredLen(recording_client);

    for (int i = 0; i <= CLUSTER_SLOTS; i++) {
        /* Find start node and slot id. */
        if (n == NULL) {
            if (i == CLUSTER_SLOTS) break;
            n = getNodeBySlot(i);
            start = i;
            continue;
        }

        /* Add cluster slots info when occur different node with start
         * or end of slot. */
        if (i == CLUSTER_SLOTS || n != getNodeBySlot(i)) {
            addNodeReplyForClusterSlot(recording_client, n, start, i - 1);
            num_primaries++;
            if (i == CLUSTER_SLOTS) break;
            n = getNodeBySlot(i);
            start = i;
        }
    }
    setDeferredArrayLen(recording_client, slot_replylen, num_primaries);
    /* For cluster slots, deferred length should put all data in reply list, not buffer */
    serverAssert(recording_client->bufpos == 0);
    sds cluster_slot_response = aggregateClientOutputBuffer(recording_client);
    deleteCachedResponseClient(recording_client);
    return cluster_slot_response;
}

int verifyCachedClusterSlotsResponse(sds cached_response, int resp) {
    sds generated_response = generateClusterSlotResponse(resp);
    int is_equal = !sdscmp(generated_response, cached_response);
    /* Here, we use LL_WARNING so this gets printed when debug assertions are enabled and the system is about to crash. */
    if (!is_equal)
        serverLog(LL_WARNING, "\ngenerated_response:\n%s\n\ncached_response:\n%s", generated_response, cached_response);
    sdsfree(generated_response);
    return is_equal;
}

void clusterCommandSlots(client *c) {
    /* Format: 1) 1) start slot
     *            2) end slot
     *            3) 1) primary IP
     *               2) primary port
     *               3) node ID
     *            4) 1) replica IP
     *               2) replica port
     *               3) node ID
     *           ... continued until done
     */
    int conn_type = 0;
    if (shouldReturnTlsInfo()) conn_type |= CACHE_CONN_TYPE_TLS;
    if (isClientConnIpV6(c)) conn_type |= CACHE_CONN_TYPE_IPv6;
    if (c->resp == 3) conn_type |= CACHE_CONN_TYPE_RESP3;

    if (detectAndUpdateCachedNodeHealth()) clearCachedClusterSlotsResponse();

    sds cached_reply = server.cached_cluster_slot_info[conn_type];
    if (!cached_reply) {
        cached_reply = generateClusterSlotResponse(c->resp);
        server.cached_cluster_slot_info[conn_type] = cached_reply;
    } else {
        debugServerAssertWithInfo(c, NULL, verifyCachedClusterSlotsResponse(cached_reply, c->resp) == 1);
    }

    addReplyProto(c, cached_reply, sdslen(cached_reply));
}

/* -----------------------------------------------------------------------------
 * Cluster functions related to serving / redirecting clients
 * -------------------------------------------------------------------------- */

/* The ASKING command is required after a -ASK redirection.
 * The client should issue ASKING before to actually send the command to
 * the target instance. See the Cluster specification for more
 * information. */
void askingCommand(client *c) {
    if (server.cluster_enabled == 0) {
        addReplyError(c, "This instance has cluster support disabled");
        return;
    }
    c->flag.asking = 1;
    addReply(c, shared.ok);
}

/* The READONLY command is used by clients to enter the read-only mode.
 * In this mode replica will not redirect clients as long as clients access
 * with read-only commands to keys that are served by the replica's primary. */
void readonlyCommand(client *c) {
    c->flag.readonly = 1;
    addReply(c, shared.ok);
}

/* The READWRITE command just clears the READONLY command state. */
void readwriteCommand(client *c) {
    c->flag.readonly = 0;
    addReply(c, shared.ok);
}

/* Resets transient cluster stats that we expose via INFO or other means that we want
 * to reset via CONFIG RESETSTAT. The function is also used in order to
 * initialize these fields in clusterInit() at server startup. */
void resetClusterStats(void) {
    if (!server.cluster_enabled) return;

    clusterSlotStatResetAll();

    memset(server.cluster->stats_bus_messages_sent, 0, sizeof(server.cluster->stats_bus_messages_sent));
    memset(server.cluster->stats_bus_messages_received, 0, sizeof(server.cluster->stats_bus_messages_received));
    server.cluster->stats_bus_bytes_sent = 0;
    server.cluster->stats_bus_bytes_received = 0;
    server.cluster->stats_bus_pubsub_bytes_sent = 0;
    server.cluster->stats_bus_pubsub_bytes_received = 0;
    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_links_established_inbound = 0;
    server.cluster->stat_cluster_links_established_outbound = 0;
}

void clusterCommandFlushslot(client *c) {
    int slot;
    int lazy = server.lazyfree_lazy_user_flush;
    if ((slot = getSlotOrReply(c, c->argv[2])) == -1) return;
    if (c->argc == 4) {
        if (!strcasecmp(objectGetVal(c->argv[3]), "async")) {
            lazy = 1;
        } else if (!strcasecmp(objectGetVal(c->argv[3]), "sync")) {
            lazy = 0;
        } else {
            addReplyErrorObject(c, shared.syntaxerr);
            return;
        }
    }
    delKeysInSlot(slot, lazy, false, true);
    addReply(c, shared.ok);
}

/* -----------------------------------------------------------------------------
 * CLUSTERSCAN Command
 * -------------------------------------------------------------------------- */

/* Compute and cache encoded fingerprint for CLUSTERSCAN cursor validation.
 *
 * Fingerprint uniquely identifies the current memory layout by hashing
 * relevant configuration bits. Currently this includes only the hash seed,
 * but additional factors (e.g., hash table type) can be mixed in later.
 *
 * Fingerprint is encoded using consecutive ASCII values starting from '0'
 * (48 to 111), each representing 6 bits. This encoding avoids '-', '{', '}'
 * which are part of the cursor.
 *
 * Returns a non zero 32-bit fingerprint. 0 is reserved for cross-node cursor */
static const char *clusterscanFingerprint(void) {
    static char cached_fp[7];
    if (cached_fp[0]) return cached_fp;

    /* Use configurable_hash_seed (derived from hash-seed config) so that nodes
     * sharing the same hash-seed produce the same fingerprint, allowing cursors
     * to survive failover without restarting the scan. */
    uint64_t *seed = (uint64_t *)getConfigurableHashSeed();
    uint64_t hash = wangHash64(seed[0] ^ seed[1]);

    /* Truncating to 32 bit instead of 64 bit */
    uint32_t fp = (uint32_t)hash;

    /* Ensure fingerprint is never 0, zero is reserved for cross node
     * cursor where fingerprint validation would be skipped, scenarios
     * include initial and end a given slots scan */
    if (fp == 0) fp = 1;

    /* Convert 32-bit fingerprint to 6 char string using base64 like encoding. */
    for (int i = 5; i >= 0; i--) {
        cached_fp[i] = '0' + (fp & 0x3F);
        fp >>= 6;
    }
    cached_fp[6] = '\0';

    return cached_fp;
}

/* Parse the cursor for CLUSTERSCAN Command.
 * The format is <fingerprint>-<hashtag>-<cursor>.
 *
 * Fingerprint identifies the node's memory layout. On mismatch, the scan
 * restarts from cursor 0 rather than returning an error.
 *
 * Hashtag is used to route to the correct node.
 *
 * Cursor is the actual local scan cursor.
 */
static int parseClusterScanCursor(robj *o, int *slot, unsigned long long *cursor) {
    char *p = objectGetVal(o);
    char *end = p + sdslen(p);
    char *token;

    /* Handle fingerprint */
    token = strchr(p, '-');
    if (!token) return C_ERR;

    size_t fp_len = token - p;
    char *fp_start = p;
    p = token + 1;

    /* Handle hashtag */
    token = strchr(p, '-');
    if (!token) return C_ERR;

    int hash_slot = keyHashSlot(p, token - p);
    *slot = hash_slot;
    p = token + 1;

    /* Handle the local cursor */
    if (!string2ull(p, end - p, cursor)) return C_ERR;

    /* Fingerprint is 0 when beginning to scan a slot so ignore the cursor value
     * in that case. If fingerprint doesn't match the current node's fingerprint,
     * restart the scan from cursor 0. From encoding we know fingerprint length is 6.*/
    const char *fp = clusterscanFingerprint();
    if (fp_len != 6 || memcmp(fp_start, fp, 6) != 0) {
        *cursor = 0;
    }

    return C_OK;
}

/* CLUSTERSCAN command - topology-aware scan across cluster slots.
 * Cursor format: <fingerprint>-<hashtag>-<local_cursor>
 * Supports SLOT, MATCH, COUNT, and TYPE options. */
void clusterscanCommand(client *c) {
    if (!server.cluster_enabled) {
        addReplyError(c, "This instance has cluster support disabled");
        return;
    }

    int slot;
    unsigned long long cursor;
    scanOptions opts;
    int skip_scan = 0;

    if (parseScanOptionsOrReply(c, NULL, 2, true, &opts) != C_OK) return;

    /* SLOT and single slot MATCH target different slots hence conclude the scan */
    skip_scan = opts.input_slot != -1 && opts.match_slot != -1 && opts.input_slot != opts.match_slot;

    /* Handle cursor "0" case. If slot information is provided we return
     * the updated cursor to scan input slot, else scan from slot 0. */
    if (strcmp(objectGetVal(c->argv[1]), "0") == 0) {
        if (opts.input_slot != -1) {
            slot = opts.input_slot;
        } else if (opts.match_slot != -1) {
            slot = opts.match_slot; /* If match maps to a particular slot, start scan from there */
        } else {
            slot = 0;
        }

        addReplyArrayLen(c, 2);
        if (skip_scan) {
            addReplyBulkCString(c, "0");
        } else {
            sds new_cursor = sdscatfmt(sdsempty(), "0-{%s}-0", crc16_slot_table[slot]);
            addReplyBulkSds(c, new_cursor);
        }
        addReplyArrayLen(c, 0);
        return;
    } else {
        if (parseClusterScanCursor(c->argv[1], &slot, &cursor) == C_ERR) {
            addReplyError(c, "Invalid cursor");
            return;
        }

        if (opts.input_slot != -1 && slot != opts.input_slot) {
            addReplyError(c, "Cursor slot mismatch with SLOT argument");
            return;
        }

        if (opts.match_slot != -1 && slot != opts.match_slot) {
            /* Advance cursor to the slot matched by MATCH if required but do not go back. */
            addReplyArrayLen(c, 2);
            if (!skip_scan && opts.match_slot > slot) {
                sds new_cursor = sdscatfmt(sdsempty(), "0-{%s}-0", crc16_slot_table[opts.match_slot]);
                addReplyBulkSds(c, new_cursor);
            } else {
                addReplyBulkCString(c, "0");
            }
            addReplyArrayLen(c, 0);
            return;
        }
    }

    /* If SLOT argument was provided or implied by MATCH, scan only that slot.
     * Otherwise, scan the continuous range of slots owned by this node. */
    int final_slot = slot;
    bool advance_to_next_slot = false;
    if (opts.input_slot == -1 && opts.match_slot == -1) {
        clusterNode *owner = getNodeBySlot(slot);
        while (final_slot + 1 < CLUSTER_SLOTS && getNodeBySlot(final_slot + 1) == owner) {
            final_slot++;
        }
        advance_to_next_slot = true;
    }

    serverAssert(slot >= 0 && slot < CLUSTER_SLOTS);
    clusterScanCtx cluster_ctx = {
        .slot = slot,
        .final_slot = final_slot,
        .advance_to_next_slot = advance_to_next_slot,
        .fp = clusterscanFingerprint(),
    };
    scanGenericCommandWithOptions(c, NULL, cursor, &opts, &cluster_ctx);
}
