syncMain.c 97.1 KB
Newer Older
M
Minghao Li 已提交
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15
/*
 * Copyright (c) 2019 TAOS Data, Inc. <jhtao@taosdata.com>
 *
 * This program is free software: you can use, redistribute, and/or modify
 * it under the terms of the GNU Affero General Public License, version 3
 * or later ("AGPL"), as published by the Free Software Foundation.
 *
 * This program is distributed in the hope that it will be useful, but WITHOUT
 * ANY WARRANTY; without even the implied warranty of MERCHANTABILITY or
 * FITNESS FOR A PARTICULAR PURPOSE.
 *
 * You should have received a copy of the GNU Affero General Public License
 * along with this program. If not, see <http://www.gnu.org/licenses/>.
 */

S
Shengliang Guan 已提交
16
#define _DEFAULT_SOURCE
M
Minghao Li 已提交
17
#include "sync.h"
M
Minghao Li 已提交
18 19
#include "syncAppendEntries.h"
#include "syncAppendEntriesReply.h"
M
Minghao Li 已提交
20
#include "syncCommit.h"
M
Minghao Li 已提交
21
#include "syncElection.h"
M
Minghao Li 已提交
22
#include "syncEnv.h"
M
Minghao Li 已提交
23
#include "syncIndexMgr.h"
M
Minghao Li 已提交
24
#include "syncInt.h"
25
#include "syncPipeline.h"
M
Minghao Li 已提交
26
#include "syncMessage.h"
M
Minghao Li 已提交
27
#include "syncRaftCfg.h"
M
Minghao Li 已提交
28
#include "syncRaftLog.h"
M
Minghao Li 已提交
29
#include "syncRaftStore.h"
M
Minghao Li 已提交
30
#include "syncReplication.h"
M
Minghao Li 已提交
31 32
#include "syncRequestVote.h"
#include "syncRequestVoteReply.h"
M
Minghao Li 已提交
33
#include "syncRespMgr.h"
M
Minghao Li 已提交
34
#include "syncSnapshot.h"
M
Minghao Li 已提交
35
#include "syncTimeout.h"
M
Minghao Li 已提交
36
#include "syncUtil.h"
M
Minghao Li 已提交
37
#include "syncVoteMgr.h"
M
Minghao Li 已提交
38
#include "tref.h"
M
Minghao Li 已提交
39

M
Minghao Li 已提交
40 41 42 43 44
static void    syncNodeEqPingTimer(void* param, void* tmrId);
static void    syncNodeEqElectTimer(void* param, void* tmrId);
static void    syncNodeEqHeartbeatTimer(void* param, void* tmrId);
static int32_t syncNodeEqNoop(SSyncNode* ths);
static int32_t syncNodeAppendNoop(SSyncNode* ths);
45
static void    syncNodeEqPeerHeartbeatTimer(void* param, void* tmrId);
S
Shengliang Guan 已提交
46
static bool    syncIsConfigChanged(const SSyncCfg* pOldCfg, const SSyncCfg* pNewCfg);
S
Shengliang Guan 已提交
47 48 49
static int32_t syncHbTimerInit(SSyncNode* pSyncNode, SSyncTimer* pSyncTimer, SRaftId destId);
static int32_t syncHbTimerStart(SSyncNode* pSyncNode, SSyncTimer* pSyncTimer);
static int32_t syncHbTimerStop(SSyncNode* pSyncNode, SSyncTimer* pSyncTimer);
S
Shengliang Guan 已提交
50 51 52 53 54 55 56 57 58 59 60
static int32_t syncNodeUpdateNewConfigIndex(SSyncNode* ths, SSyncCfg* pNewCfg);
static bool    syncNodeInConfig(SSyncNode* pSyncNode, const SSyncCfg* config);
static void    syncNodeDoConfigChange(SSyncNode* pSyncNode, SSyncCfg* newConfig, SyncIndex lastConfigChangeIndex);
static bool    syncNodeIsOptimizedOneReplica(SSyncNode* ths, SRpcMsg* pMsg);

static bool    syncNodeCanChange(SSyncNode* pSyncNode);
static int32_t syncNodeLeaderTransfer(SSyncNode* pSyncNode);
static int32_t syncNodeLeaderTransferTo(SSyncNode* pSyncNode, SNodeInfo newLeader);
static int32_t syncDoLeaderTransfer(SSyncNode* ths, SRpcMsg* pRpcMsg, SSyncRaftEntry* pEntry);

static ESyncStrategy syncNodeStrategy(SSyncNode* pSyncNode);
M
Minghao Li 已提交
61

62
int64_t syncOpen(SSyncInfo* pSyncInfo) {
M
Minghao Li 已提交
63
  SSyncNode* pSyncNode = syncNodeOpen(pSyncInfo);
64
  if (pSyncNode == NULL) {
S
Shengliang Guan 已提交
65
    sError("vgId:%d, failed to open sync node", pSyncInfo->vgId);
66 67
    return -1;
  }
M
Minghao Li 已提交
68

S
Shengliang Guan 已提交
69
  pSyncNode->rid = syncNodeAdd(pSyncNode);
M
Minghao Li 已提交
70
  if (pSyncNode->rid < 0) {
71
    syncNodeClose(pSyncNode);
M
Minghao Li 已提交
72 73 74
    return -1;
  }

S
Shengliang Guan 已提交
75 76 77 78 79 80
  pSyncNode->pingBaseLine = pSyncInfo->pingMs;
  pSyncNode->pingTimerMS = pSyncInfo->pingMs;
  pSyncNode->electBaseLine = pSyncInfo->electMs;
  pSyncNode->hbBaseLine = pSyncInfo->heartbeatMs;
  pSyncNode->heartbeatTimerMS = pSyncInfo->heartbeatMs;
  pSyncNode->msgcb = pSyncInfo->msgcb;
M
Minghao Li 已提交
81
  return pSyncNode->rid;
M
Minghao Li 已提交
82
}
M
Minghao Li 已提交
83

B
Benguang Zhao 已提交
84
int32_t syncStart(int64_t rid) {
S
Shengliang Guan 已提交
85
  SSyncNode* pSyncNode = syncNodeAcquire(rid);
M
Minghao Li 已提交
86
  if (pSyncNode == NULL) {
B
Benguang Zhao 已提交
87 88 89 90 91
    sError("failed to acquire rid: %" PRId64 " of tsNodeReftId for pSyncNode", rid);
    return -1;
  }

  if (syncNodeRestore(pSyncNode) < 0) {
92
    sError("vgId:%d, failed to restore sync log buffer since %s", pSyncNode->vgId, terrstr());
93
    goto _err;
M
Minghao Li 已提交
94
  }
M
Minghao Li 已提交
95

B
Benguang Zhao 已提交
96 97 98 99
  if (syncNodeStart(pSyncNode) < 0) {
    sError("vgId:%d, failed to start sync node since %s", pSyncNode->vgId, terrstr());
    goto _err;
  }
M
Minghao Li 已提交
100

B
Benguang Zhao 已提交
101 102
  syncNodeRelease(pSyncNode);
  return 0;
M
Minghao Li 已提交
103

104 105 106
_err:
  syncNodeRelease(pSyncNode);
  return -1;
M
Minghao Li 已提交
107 108
}

M
Minghao Li 已提交
109
void syncStop(int64_t rid) {
S
Shengliang Guan 已提交
110
  SSyncNode* pSyncNode = syncNodeAcquire(rid);
S
Shengliang Guan 已提交
111
  if (pSyncNode != NULL) {
112
    pSyncNode->isStart = false;
S
Shengliang Guan 已提交
113
    syncNodeRelease(pSyncNode);
S
Shengliang Guan 已提交
114
    syncNodeRemove(rid);
M
Minghao Li 已提交
115 116 117
  }
}

M
Minghao Li 已提交
118 119
void syncPreStop(int64_t rid) {
  SSyncNode* pSyncNode = syncNodeAcquire(rid);
120 121 122
  if (pSyncNode != NULL) {
    syncNodePreClose(pSyncNode);
    syncNodeRelease(pSyncNode);
M
Minghao Li 已提交
123 124 125
  }
}

S
Shengliang Guan 已提交
126 127 128
static bool syncNodeCheckNewConfig(SSyncNode* pSyncNode, const SSyncCfg* pCfg) {
  if (!syncNodeInConfig(pSyncNode, pCfg)) return false;
  return abs(pCfg->replicaNum - pSyncNode->replicaNum) <= 1;
M
Minghao Li 已提交
129 130
}

S
Shengliang Guan 已提交
131
int32_t syncReconfig(int64_t rid, SSyncCfg* pNewCfg) {
S
Shengliang Guan 已提交
132
  SSyncNode* pSyncNode = syncNodeAcquire(rid);
S
Shengliang Guan 已提交
133
  if (pSyncNode == NULL) return -1;
M
Minghao Li 已提交
134

M
Minghao Li 已提交
135
  if (!syncNodeCheckNewConfig(pSyncNode, pNewCfg)) {
S
Shengliang Guan 已提交
136
    syncNodeRelease(pSyncNode);
M
Minghao Li 已提交
137
    terrno = TSDB_CODE_SYN_NEW_CONFIG_ERROR;
S
Shengliang Guan 已提交
138
    sError("vgId:%d, failed to reconfig since invalid new config", pSyncNode->vgId);
M
Minghao Li 已提交
139
    return -1;
M
Minghao Li 已提交
140
  }
141

S
Shengliang Guan 已提交
142 143
  syncNodeUpdateNewConfigIndex(pSyncNode, pNewCfg);
  syncNodeDoConfigChange(pSyncNode, pNewCfg, SYNC_INDEX_INVALID);
S
Shengliang Guan 已提交
144

M
Minghao Li 已提交
145 146 147 148
  if (pSyncNode->state == TAOS_SYNC_STATE_LEADER) {
    syncNodeStopHeartbeatTimer(pSyncNode);

    for (int32_t i = 0; i < TSDB_MAX_REPLICA; ++i) {
S
Shengliang Guan 已提交
149
      syncHbTimerInit(pSyncNode, &pSyncNode->peerHeartbeatTimerArr[i], pSyncNode->replicasId[i]);
M
Minghao Li 已提交
150 151 152 153 154
    }

    syncNodeStartHeartbeatTimer(pSyncNode);
    syncNodeReplicate(pSyncNode);
  }
S
Shengliang Guan 已提交
155

S
Shengliang Guan 已提交
156
  syncNodeRelease(pSyncNode);
S
Shengliang Guan 已提交
157
  return 0;
M
Minghao Li 已提交
158
}
M
Minghao Li 已提交
159

S
Shengliang Guan 已提交
160 161 162 163
int32_t syncProcessMsg(int64_t rid, SRpcMsg* pMsg) {
  int32_t code = -1;
  if (!syncIsInit()) return code;

S
Shengliang Guan 已提交
164
  SSyncNode* pSyncNode = syncNodeAcquire(rid);
S
Shengliang Guan 已提交
165 166
  if (pSyncNode == NULL) return code;

S
Shengliang Guan 已提交
167 168 169 170 171 172 173 174 175 176 177 178 179 180 181 182 183 184 185 186 187 188 189 190 191 192 193 194 195 196 197 198 199 200 201 202 203 204
  switch (pMsg->msgType) {
    case TDMT_SYNC_HEARTBEAT:
      code = syncNodeOnHeartbeat(pSyncNode, pMsg);
      break;
    case TDMT_SYNC_HEARTBEAT_REPLY:
      code = syncNodeOnHeartbeatReply(pSyncNode, pMsg);
      break;
    case TDMT_SYNC_TIMEOUT:
      code = syncNodeOnTimeout(pSyncNode, pMsg);
      break;
    case TDMT_SYNC_CLIENT_REQUEST:
      code = syncNodeOnClientRequest(pSyncNode, pMsg, NULL);
      break;
    case TDMT_SYNC_REQUEST_VOTE:
      code = syncNodeOnRequestVote(pSyncNode, pMsg);
      break;
    case TDMT_SYNC_REQUEST_VOTE_REPLY:
      code = syncNodeOnRequestVoteReply(pSyncNode, pMsg);
      break;
    case TDMT_SYNC_APPEND_ENTRIES:
      code = syncNodeOnAppendEntries(pSyncNode, pMsg);
      break;
    case TDMT_SYNC_APPEND_ENTRIES_REPLY:
      code = syncNodeOnAppendEntriesReply(pSyncNode, pMsg);
      break;
    case TDMT_SYNC_SNAPSHOT_SEND:
      code = syncNodeOnSnapshot(pSyncNode, pMsg);
      break;
    case TDMT_SYNC_SNAPSHOT_RSP:
      code = syncNodeOnSnapshotReply(pSyncNode, pMsg);
      break;
    case TDMT_SYNC_LOCAL_CMD:
      code = syncNodeOnLocalCmd(pSyncNode, pMsg);
      break;
    default:
      sError("vgId:%d, failed to process msg:%p since invalid type:%s", pSyncNode->vgId, pMsg,
             TMSG_INFO(pMsg->msgType));
      code = -1;
M
Minghao Li 已提交
205 206
  }

S
Shengliang Guan 已提交
207
  syncNodeRelease(pSyncNode);
S
Shengliang Guan 已提交
208
  return code;
209 210
}

S
Shengliang Guan 已提交
211
int32_t syncLeaderTransfer(int64_t rid) {
S
Shengliang Guan 已提交
212
  SSyncNode* pSyncNode = syncNodeAcquire(rid);
S
Shengliang Guan 已提交
213
  if (pSyncNode == NULL) return -1;
214

S
Shengliang Guan 已提交
215
  int32_t ret = syncNodeLeaderTransfer(pSyncNode);
S
Shengliang Guan 已提交
216
  syncNodeRelease(pSyncNode);
217 218 219
  return ret;
}

M
Minghao Li 已提交
220 221 222 223 224 225 226 227 228 229 230 231 232 233 234 235
SyncIndex syncMinMatchIndex(SSyncNode* pSyncNode) {
  SyncIndex minMatchIndex = SYNC_INDEX_INVALID;

  if (pSyncNode->peersNum > 0) {
    minMatchIndex = syncIndexMgrGetIndex(pSyncNode->pMatchIndex, &(pSyncNode->peersId[0]));
  }

  for (int32_t i = 1; i < pSyncNode->peersNum; ++i) {
    SyncIndex matchIndex = syncIndexMgrGetIndex(pSyncNode->pMatchIndex, &(pSyncNode->peersId[i]));
    if (matchIndex < minMatchIndex) {
      minMatchIndex = matchIndex;
    }
  }
  return minMatchIndex;
}

236
int32_t syncBeginSnapshot(int64_t rid, int64_t lastApplyIndex) {
S
Shengliang Guan 已提交
237
  SSyncNode* pSyncNode = syncNodeAcquire(rid);
238
  if (pSyncNode == NULL) {
239
    sError("sync begin snapshot error");
240 241
    return -1;
  }
242

243 244
  int32_t code = 0;

M
Minghao Li 已提交
245
  if (syncNodeIsMnode(pSyncNode)) {
M
Minghao Li 已提交
246 247 248
    // mnode
    int64_t logRetention = SYNC_MNODE_LOG_RETENTION;

M
Minghao Li 已提交
249 250 251
    SyncIndex beginIndex = pSyncNode->pLogStore->syncLogBeginIndex(pSyncNode->pLogStore);
    SyncIndex endIndex = pSyncNode->pLogStore->syncLogEndIndex(pSyncNode->pLogStore);
    int64_t   logNum = endIndex - beginIndex;
M
Minghao Li 已提交
252 253 254
    bool      isEmpty = pSyncNode->pLogStore->syncLogIsEmpty(pSyncNode->pLogStore);

    if (isEmpty || (!isEmpty && logNum < logRetention)) {
S
Shengliang Guan 已提交
255 256
      sNTrace(pSyncNode, "new-snapshot-index:%" PRId64 ", log-num:%" PRId64 ", empty:%d, do not delete wal",
              lastApplyIndex, logNum, isEmpty);
S
Shengliang Guan 已提交
257
      syncNodeRelease(pSyncNode);
M
Minghao Li 已提交
258 259 260
      return 0;
    }

M
Minghao Li 已提交
261 262 263
    goto _DEL_WAL;

  } else {
264 265 266 267 268 269 270 271 272 273 274 275
    lastApplyIndex -= SYNC_VNODE_LOG_RETENTION;

    SyncIndex beginIndex = pSyncNode->pLogStore->syncLogBeginIndex(pSyncNode->pLogStore);
    SyncIndex endIndex = pSyncNode->pLogStore->syncLogEndIndex(pSyncNode->pLogStore);
    bool      isEmpty = pSyncNode->pLogStore->syncLogIsEmpty(pSyncNode->pLogStore);

    if (isEmpty || !(lastApplyIndex >= beginIndex && lastApplyIndex <= endIndex)) {
      sNTrace(pSyncNode, "new-snapshot-index:%" PRId64 ", empty:%d, do not delete wal", lastApplyIndex, isEmpty);
      syncNodeRelease(pSyncNode);
      return 0;
    }

M
Minghao Li 已提交
276 277 278 279 280 281 282 283 284 285 286 287 288 289
    // vnode
    if (pSyncNode->replicaNum > 1) {
      // multi replicas

      if (pSyncNode->state == TAOS_SYNC_STATE_LEADER) {
        pSyncNode->minMatchIndex = syncMinMatchIndex(pSyncNode);

        for (int32_t i = 0; i < pSyncNode->peersNum; ++i) {
          int64_t matchIndex = syncIndexMgrGetIndex(pSyncNode->pMatchIndex, &(pSyncNode->peersId[i]));
          if (lastApplyIndex > matchIndex) {
            do {
              char     host[64];
              uint16_t port;
              syncUtilU642Addr(pSyncNode->peersId[i].addr, host, sizeof(host), &port);
S
Shengliang Guan 已提交
290 291 292 293
              sNTrace(pSyncNode,
                      "new-snapshot-index:%" PRId64 " is greater than match-index:%" PRId64
                      " of %s:%d, do not delete wal",
                      lastApplyIndex, matchIndex, host, port);
M
Minghao Li 已提交
294 295
            } while (0);

S
Shengliang Guan 已提交
296
            syncNodeRelease(pSyncNode);
M
Minghao Li 已提交
297 298 299 300 301 302
            return 0;
          }
        }

      } else if (pSyncNode->state == TAOS_SYNC_STATE_FOLLOWER) {
        if (lastApplyIndex > pSyncNode->minMatchIndex) {
S
Shengliang Guan 已提交
303 304 305
          sNTrace(pSyncNode,
                  "new-snapshot-index:%" PRId64 " is greater than min-match-index:%" PRId64 ", do not delete wal",
                  lastApplyIndex, pSyncNode->minMatchIndex);
S
Shengliang Guan 已提交
306
          syncNodeRelease(pSyncNode);
M
Minghao Li 已提交
307 308 309 310
          return 0;
        }

      } else if (pSyncNode->state == TAOS_SYNC_STATE_CANDIDATE) {
S
Shengliang Guan 已提交
311
        sNTrace(pSyncNode, "new-snapshot-index:%" PRId64 " candidate, do not delete wal", lastApplyIndex);
S
Shengliang Guan 已提交
312
        syncNodeRelease(pSyncNode);
M
Minghao Li 已提交
313 314 315
        return 0;

      } else {
S
Shengliang Guan 已提交
316
        sNTrace(pSyncNode, "new-snapshot-index:%" PRId64 " unknown state, do not delete wal", lastApplyIndex);
S
Shengliang Guan 已提交
317
        syncNodeRelease(pSyncNode);
M
Minghao Li 已提交
318 319 320 321 322 323 324 325 326
        return 0;
      }

      goto _DEL_WAL;

    } else {
      // one replica

      goto _DEL_WAL;
327 328 329
    }
  }

M
Minghao Li 已提交
330
_DEL_WAL:
331

M
Minghao Li 已提交
332
  do {
333 334 335 336 337 338 339 340 341 342 343 344 345 346 347 348 349 350 351 352
    SSyncLogStoreData* pData = pSyncNode->pLogStore->data;
    SyncIndex          snapshotVer = walGetSnapshotVer(pData->pWal);
    SyncIndex          walCommitVer = walGetCommittedVer(pData->pWal);
    SyncIndex          wallastVer = walGetLastVer(pData->pWal);
    if (lastApplyIndex <= walCommitVer) {
      SyncIndex snapshottingIndex = atomic_load_64(&pSyncNode->snapshottingIndex);

      if (snapshottingIndex == SYNC_INDEX_INVALID) {
        atomic_store_64(&pSyncNode->snapshottingIndex, lastApplyIndex);
        pSyncNode->snapshottingTime = taosGetTimestampMs();

        code = walBeginSnapshot(pData->pWal, lastApplyIndex);
        if (code == 0) {
          sNTrace(pSyncNode, "wal snapshot begin, index:%" PRId64 ", last apply index:%" PRId64,
                  pSyncNode->snapshottingIndex, lastApplyIndex);
        } else {
          sNError(pSyncNode, "wal snapshot begin error since:%s, index:%" PRId64 ", last apply index:%" PRId64,
                  terrstr(terrno), pSyncNode->snapshottingIndex, lastApplyIndex);
          atomic_store_64(&pSyncNode->snapshottingIndex, SYNC_INDEX_INVALID);
        }
353

M
Minghao Li 已提交
354
      } else {
355 356
        sNTrace(pSyncNode, "snapshotting for %" PRId64 ", do not delete wal for new-snapshot-index:%" PRId64,
                snapshottingIndex, lastApplyIndex);
M
Minghao Li 已提交
357
      }
358
    }
M
Minghao Li 已提交
359
  } while (0);
360

S
Shengliang Guan 已提交
361
  syncNodeRelease(pSyncNode);
362 363 364 365
  return code;
}

int32_t syncEndSnapshot(int64_t rid) {
S
Shengliang Guan 已提交
366
  SSyncNode* pSyncNode = syncNodeAcquire(rid);
367
  if (pSyncNode == NULL) {
368
    sError("sync end snapshot error");
369 370 371
    return -1;
  }

372 373 374 375
  int32_t code = 0;
  if (atomic_load_64(&pSyncNode->snapshottingIndex) != SYNC_INDEX_INVALID) {
    SSyncLogStoreData* pData = pSyncNode->pLogStore->data;
    code = walEndSnapshot(pData->pWal);
M
Minghao Li 已提交
376
    if (code != 0) {
377
      sNError(pSyncNode, "wal snapshot end error since:%s", terrstr());
S
Shengliang Guan 已提交
378
      syncNodeRelease(pSyncNode);
M
Minghao Li 已提交
379 380
      return -1;
    } else {
S
Shengliang Guan 已提交
381
      sNTrace(pSyncNode, "wal snapshot end, index:%" PRId64, atomic_load_64(&pSyncNode->snapshottingIndex));
M
Minghao Li 已提交
382 383
      atomic_store_64(&pSyncNode->snapshottingIndex, SYNC_INDEX_INVALID);
    }
384
  }
385

S
Shengliang Guan 已提交
386
  syncNodeRelease(pSyncNode);
387 388 389
  return code;
}

M
Minghao Li 已提交
390
int32_t syncStepDown(int64_t rid, SyncTerm newTerm) {
S
Shengliang Guan 已提交
391
  SSyncNode* pSyncNode = syncNodeAcquire(rid);
M
Minghao Li 已提交
392
  if (pSyncNode == NULL) {
393
    sError("sync step down error");
M
Minghao Li 已提交
394 395 396
    return -1;
  }

M
Minghao Li 已提交
397
  syncNodeStepDown(pSyncNode, newTerm);
S
Shengliang Guan 已提交
398
  syncNodeRelease(pSyncNode);
M
Minghao Li 已提交
399
  return 0;
M
Minghao Li 已提交
400 401
}

402
bool syncNodeIsReadyForRead(SSyncNode* pSyncNode) {
403
  if (pSyncNode == NULL) {
404
    terrno = TSDB_CODE_SYN_INTERNAL_ERROR;
405
    sError("sync ready for read error");
406 407
    return false;
  }
M
Minghao Li 已提交
408

409 410 411 412 413 414
  if (pSyncNode->state != TAOS_SYNC_STATE_LEADER) {
    terrno = TSDB_CODE_SYN_NOT_LEADER;
    return false;
  }

  if (pSyncNode->restoreFinish) {
415
    return true;
M
Minghao Li 已提交
416 417
  }

418
  bool ready = false;
419 420 421
  if (!pSyncNode->pFsm->FpApplyQueueEmptyCb(pSyncNode->pFsm)) {
    // apply queue not empty
    ready = false;
M
Minghao Li 已提交
422

423 424 425 426 427 428 429 430 431 432 433 434 435
  } else {
    if (!pSyncNode->pLogStore->syncLogIsEmpty(pSyncNode->pLogStore)) {
      SyncIndex       lastIndex = pSyncNode->pLogStore->syncLogLastIndex(pSyncNode->pLogStore);
      SSyncRaftEntry* pEntry = NULL;
      SLRUCache*      pCache = pSyncNode->pLogStore->pCache;
      LRUHandle*      h = taosLRUCacheLookup(pCache, &lastIndex, sizeof(lastIndex));
      int32_t         code = 0;
      if (h) {
        pEntry = (SSyncRaftEntry*)taosLRUCacheValue(pCache, h);
        code = 0;

        pSyncNode->pLogStore->cacheHit++;
        sNTrace(pSyncNode, "hit cache index:%" PRId64 ", bytes:%u, %p", lastIndex, pEntry->bytes, pEntry);
M
Minghao Li 已提交
436

437 438 439
      } else {
        pSyncNode->pLogStore->cacheMiss++;
        sNTrace(pSyncNode, "miss cache index:%" PRId64, lastIndex);
M
Minghao Li 已提交
440

441 442
        code = pSyncNode->pLogStore->syncLogGetEntry(pSyncNode->pLogStore, lastIndex, &pEntry);
      }
443

444 445 446
      if (code == 0 && pEntry != NULL) {
        if (pEntry->originalRpcType == TDMT_SYNC_NOOP && pEntry->term == pSyncNode->pRaftStore->currentTerm) {
          ready = true;
447
        }
448

449 450 451 452
        if (h) {
          taosLRUCacheRelease(pCache, h, false);
        } else {
          syncEntryDestroy(pEntry);
453
        }
454 455 456 457
      }
    }
  }

458
  if (!ready) {
459
    terrno = TSDB_CODE_SYN_RESTORING;
460
  }
461

462 463 464 465 466 467 468 469 470 471 472 473
  return ready;
}

bool syncIsReadyForRead(int64_t rid) {
  SSyncNode* pSyncNode = syncNodeAcquire(rid);
  if (pSyncNode == NULL) {
    sError("sync ready for read error");
    return false;
  }

  bool ready = syncNodeIsReadyForRead(pSyncNode);

474 475
  syncNodeRelease(pSyncNode);
  return ready;
M
Minghao Li 已提交
476
}
M
Minghao Li 已提交
477

478 479 480 481 482 483 484 485 486 487 488 489 490 491 492 493 494 495 496 497 498 499
bool syncSnapshotSending(int64_t rid) {
  SSyncNode* pSyncNode = syncNodeAcquire(rid);
  if (pSyncNode == NULL) {
    return false;
  }

  bool b = syncNodeSnapshotSending(pSyncNode);
  syncNodeRelease(pSyncNode);
  return b;
}

bool syncSnapshotRecving(int64_t rid) {
  SSyncNode* pSyncNode = syncNodeAcquire(rid);
  if (pSyncNode == NULL) {
    return false;
  }

  bool b = syncNodeSnapshotRecving(pSyncNode);
  syncNodeRelease(pSyncNode);
  return b;
}

M
Minghao Li 已提交
500 501
int32_t syncNodeLeaderTransfer(SSyncNode* pSyncNode) {
  if (pSyncNode->peersNum == 0) {
S
Shengliang Guan 已提交
502
    sDebug("vgId:%d, only one replica, cannot leader transfer", pSyncNode->vgId);
M
Minghao Li 已提交
503 504
    terrno = TSDB_CODE_SYN_ONE_REPLICA;
    return -1;
M
Minghao Li 已提交
505
  }
M
Minghao Li 已提交
506

507
  int32_t ret = 0;
508
  if (pSyncNode->state == TAOS_SYNC_STATE_LEADER && pSyncNode->replicaNum > 1) {
509
    SNodeInfo newLeader = (pSyncNode->peersNodeInfo)[0];
510 511 512 513 514 515 516
    if (pSyncNode->peersNum == 2) {
      SyncIndex matchIndex0 = syncIndexMgrGetIndex(pSyncNode->pMatchIndex, &(pSyncNode->peersId[0]));
      SyncIndex matchIndex1 = syncIndexMgrGetIndex(pSyncNode->pMatchIndex, &(pSyncNode->peersId[1]));
      if (matchIndex1 > matchIndex0) {
        newLeader = (pSyncNode->peersNodeInfo)[1];
      }
    }
517 518 519
    ret = syncNodeLeaderTransferTo(pSyncNode, newLeader);
  }

M
Minghao Li 已提交
520
  return ret;
M
Minghao Li 已提交
521 522
}

M
Minghao Li 已提交
523 524
int32_t syncNodeLeaderTransferTo(SSyncNode* pSyncNode, SNodeInfo newLeader) {
  if (pSyncNode->replicaNum == 1) {
S
Shengliang Guan 已提交
525
    sDebug("vgId:%d, only one replica, cannot leader transfer", pSyncNode->vgId);
M
Minghao Li 已提交
526 527
    terrno = TSDB_CODE_SYN_ONE_REPLICA;
    return -1;
M
Minghao Li 已提交
528
  }
529

S
Shengliang Guan 已提交
530
  sNTrace(pSyncNode, "begin leader transfer to %s:%u", newLeader.nodeFqdn, newLeader.nodePort);
M
Minghao Li 已提交
531

532 533 534 535
  SRpcMsg rpcMsg = {0};
  (void)syncBuildLeaderTransfer(&rpcMsg, pSyncNode->vgId);

  SyncLeaderTransfer* pMsg = rpcMsg.pCont;
M
Minghao Li 已提交
536 537 538 539
  pMsg->newLeaderId.addr = syncUtilAddr2U64(newLeader.nodeFqdn, newLeader.nodePort);
  pMsg->newLeaderId.vgId = pSyncNode->vgId;
  pMsg->newNodeInfo = newLeader;

S
Shengliang Guan 已提交
540 541 542
  int32_t ret = syncNodePropose(pSyncNode, &rpcMsg, false);
  rpcFreeCont(rpcMsg.pCont);
  return ret;
M
Minghao Li 已提交
543 544
}

545 546
SSyncState syncGetState(int64_t rid) {
  SSyncState state = {.state = TAOS_SYNC_STATE_ERROR};
M
Minghao Li 已提交
547

S
Shengliang Guan 已提交
548
  SSyncNode* pSyncNode = syncNodeAcquire(rid);
549 550 551
  if (pSyncNode != NULL) {
    state.state = pSyncNode->state;
    state.restored = pSyncNode->restoreFinish;
552 553 554 555 556
    if (pSyncNode->vgId != 1) {
      state.canRead = syncNodeIsReadyForRead(pSyncNode);
    } else {
      state.canRead = state.restored;
    }
557
    syncNodeRelease(pSyncNode);
M
Minghao Li 已提交
558 559
  }

560
  return state;
M
Minghao Li 已提交
561 562
}

563
#if 0
564 565 566 567 568
int32_t syncGetSnapshotByIndex(int64_t rid, SyncIndex index, SSnapshot* pSnapshot) {
  if (index < SYNC_INDEX_BEGIN) {
    return -1;
  }

S
Shengliang Guan 已提交
569
  SSyncNode* pSyncNode = syncNodeAcquire(rid);
570 571 572 573 574 575 576 577 578
  if (pSyncNode == NULL) {
    return -1;
  }
  ASSERT(rid == pSyncNode->rid);

  SSyncRaftEntry* pEntry = NULL;
  int32_t         code = pSyncNode->pLogStore->syncLogGetEntry(pSyncNode->pLogStore, index, &pEntry);
  if (code != 0) {
    if (pEntry != NULL) {
B
Benguang Zhao 已提交
579
      syncEntryDestroy(pEntry);
580
    }
S
Shengliang Guan 已提交
581
    syncNodeRelease(pSyncNode);
582 583 584 585 586 587 588 589 590
    return -1;
  }
  ASSERT(pEntry != NULL);

  pSnapshot->data = NULL;
  pSnapshot->lastApplyIndex = index;
  pSnapshot->lastApplyTerm = pEntry->term;
  pSnapshot->lastConfigIndex = syncNodeGetSnapshotConfigIndex(pSyncNode, index);

591
  syncEntryDestroy(pEntry);
S
Shengliang Guan 已提交
592
  syncNodeRelease(pSyncNode);
593 594 595
  return 0;
}

596
int32_t syncGetSnapshotMeta(int64_t rid, struct SSnapshotMeta* sMeta) {
S
Shengliang Guan 已提交
597
  SSyncNode* pSyncNode = syncNodeAcquire(rid);
598 599 600
  if (pSyncNode == NULL) {
    return -1;
  }
M
Minghao Li 已提交
601
  ASSERT(rid == pSyncNode->rid);
602 603
  sMeta->lastConfigIndex = pSyncNode->pRaftCfg->lastConfigIndex;

S
Shengliang Guan 已提交
604
  sTrace("vgId:%d, get snapshot meta, lastConfigIndex:%" PRId64, pSyncNode->vgId, pSyncNode->pRaftCfg->lastConfigIndex);
605

S
Shengliang Guan 已提交
606
  syncNodeRelease(pSyncNode);
607 608 609
  return 0;
}

610
int32_t syncGetSnapshotMetaByIndex(int64_t rid, SyncIndex snapshotIndex, struct SSnapshotMeta* sMeta) {
S
Shengliang Guan 已提交
611
  SSyncNode* pSyncNode = syncNodeAcquire(rid);
612 613 614
  if (pSyncNode == NULL) {
    return -1;
  }
M
Minghao Li 已提交
615
  ASSERT(rid == pSyncNode->rid);
616 617 618 619

  ASSERT(pSyncNode->pRaftCfg->configIndexCount >= 1);
  SyncIndex lastIndex = (pSyncNode->pRaftCfg->configIndexArr)[0];

S
Shengliang Guan 已提交
620
  for (int32_t i = 0; i < pSyncNode->pRaftCfg->configIndexCount; ++i) {
621 622 623 624 625 626
    if ((pSyncNode->pRaftCfg->configIndexArr)[i] > lastIndex &&
        (pSyncNode->pRaftCfg->configIndexArr)[i] <= snapshotIndex) {
      lastIndex = (pSyncNode->pRaftCfg->configIndexArr)[i];
    }
  }
  sMeta->lastConfigIndex = lastIndex;
627
  sTrace("vgId:%d, get snapshot meta by index:%" PRId64 " lcindex:%" PRId64, pSyncNode->vgId, snapshotIndex,
S
Shengliang Guan 已提交
628
         sMeta->lastConfigIndex);
629

S
Shengliang Guan 已提交
630
  syncNodeRelease(pSyncNode);
631 632
  return 0;
}
633
#endif
634

635 636 637 638
SyncIndex syncNodeGetSnapshotConfigIndex(SSyncNode* pSyncNode, SyncIndex snapshotLastApplyIndex) {
  ASSERT(pSyncNode->pRaftCfg->configIndexCount >= 1);
  SyncIndex lastIndex = (pSyncNode->pRaftCfg->configIndexArr)[0];

S
Shengliang Guan 已提交
639
  for (int32_t i = 0; i < pSyncNode->pRaftCfg->configIndexCount; ++i) {
640 641 642 643 644
    if ((pSyncNode->pRaftCfg->configIndexArr)[i] > lastIndex &&
        (pSyncNode->pRaftCfg->configIndexArr)[i] <= snapshotLastApplyIndex) {
      lastIndex = (pSyncNode->pRaftCfg->configIndexArr)[i];
    }
  }
S
Shengliang Guan 已提交
645
  sTrace("vgId:%d, sync get last config index, index:%" PRId64 " lcindex:%" PRId64, pSyncNode->vgId,
S
Shengliang Guan 已提交
646
         snapshotLastApplyIndex, lastIndex);
647 648 649 650

  return lastIndex;
}

651 652
void syncGetRetryEpSet(int64_t rid, SEpSet* pEpSet) {
  pEpSet->numOfEps = 0;
M
Minghao Li 已提交
653

S
Shengliang Guan 已提交
654
  SSyncNode* pSyncNode = syncNodeAcquire(rid);
S
Shengliang Guan 已提交
655
  if (pSyncNode == NULL) return;
M
Minghao Li 已提交
656

S
Shengliang Guan 已提交
657
  for (int32_t i = 0; i < pSyncNode->pRaftCfg->cfg.replicaNum; ++i) {
S
Shengliang Guan 已提交
658 659 660 661
    SEp* pEp = &pEpSet->eps[i];
    tstrncpy(pEp->fqdn, pSyncNode->pRaftCfg->cfg.nodeInfo[i].nodeFqdn, TSDB_FQDN_LEN);
    pEp->port = (pSyncNode->pRaftCfg->cfg.nodeInfo)[i].nodePort;
    pEpSet->numOfEps++;
662
    sDebug("vgId:%d, sync get retry epset, index:%d %s:%d", pSyncNode->vgId, i, pEp->fqdn, pEp->port);
M
Minghao Li 已提交
663
  }
M
Minghao Li 已提交
664 665
  if (pEpSet->numOfEps > 0) {
    pEpSet->inUse = (pSyncNode->pRaftCfg->cfg.myIndex + 1) % pEpSet->numOfEps;
M
Minghao Li 已提交
666 667
  }

S
Shengliang Guan 已提交
668
  sInfo("vgId:%d, sync get retry epset numOfEps:%d inUse:%d", pSyncNode->vgId, pEpSet->numOfEps, pEpSet->inUse);
S
Shengliang Guan 已提交
669
  syncNodeRelease(pSyncNode);
M
Minghao Li 已提交
670 671
}

M
Minghao Li 已提交
672
int32_t syncPropose(int64_t rid, SRpcMsg* pMsg, bool isWeak) {
S
Shengliang Guan 已提交
673
  SSyncNode* pSyncNode = syncNodeAcquire(rid);
674
  if (pSyncNode == NULL) {
675
    sError("sync propose error");
M
Minghao Li 已提交
676
    return -1;
677
  }
678

679
  int32_t ret = syncNodePropose(pSyncNode, pMsg, isWeak);
S
Shengliang Guan 已提交
680
  syncNodeRelease(pSyncNode);
M
Minghao Li 已提交
681 682
  return ret;
}
M
Minghao Li 已提交
683

684
int32_t syncNodePropose(SSyncNode* pSyncNode, SRpcMsg* pMsg, bool isWeak) {
S
Shengliang Guan 已提交
685 686 687 688 689
  if (pSyncNode->state != TAOS_SYNC_STATE_LEADER) {
    terrno = TSDB_CODE_SYN_NOT_LEADER;
    sNError(pSyncNode, "sync propose not leader, %s, type:%s", syncStr(pSyncNode->state), TMSG_INFO(pMsg->msgType));
    return -1;
  }
690

S
Shengliang Guan 已提交
691 692 693 694 695 696 697
  // not restored, vnode enable
  if (!pSyncNode->restoreFinish && pSyncNode->vgId != 1) {
    terrno = TSDB_CODE_SYN_PROPOSE_NOT_READY;
    sNError(pSyncNode, "failed to sync propose since not ready, type:%s, last:%" PRId64 ", cmt:%" PRId64,
            TMSG_INFO(pMsg->msgType), syncNodeGetLastIndex(pSyncNode), pSyncNode->commitIndex);
    return -1;
  }
698

699
  // heartbeat timeout
700
  if (syncNodeHeartbeatReplyTimeout(pSyncNode)) {
701 702 703 704 705 706
    terrno = TSDB_CODE_SYN_PROPOSE_NOT_READY;
    sNError(pSyncNode, "failed to sync propose since hearbeat timeout, type:%s, last:%" PRId64 ", cmt:%" PRId64,
            TMSG_INFO(pMsg->msgType), syncNodeGetLastIndex(pSyncNode), pSyncNode->commitIndex);
    return -1;
  }

S
Shengliang Guan 已提交
707 708 709
  // optimized one replica
  if (syncNodeIsOptimizedOneReplica(pSyncNode, pMsg)) {
    SyncIndex retIndex;
710
    int32_t   code = syncNodeOnClientRequest(pSyncNode, pMsg, &retIndex);
S
Shengliang Guan 已提交
711 712 713
    if (code == 0) {
      pMsg->info.conn.applyIndex = retIndex;
      pMsg->info.conn.applyTerm = pSyncNode->pRaftStore->currentTerm;
714 715 716
      sTrace("vgId:%d, propose optimized msg, index:%" PRId64 " type:%s", pSyncNode->vgId, retIndex,
             TMSG_INFO(pMsg->msgType));
      return 1;
M
Minghao Li 已提交
717
    } else {
S
Shengliang Guan 已提交
718
      terrno = TSDB_CODE_SYN_INTERNAL_ERROR;
719
      sError("vgId:%d, failed to propose optimized msg, index:%" PRId64 " type:%s", pSyncNode->vgId, retIndex,
S
Shengliang Guan 已提交
720
             TMSG_INFO(pMsg->msgType));
721
      return -1;
722
    }
S
Shengliang Guan 已提交
723
  } else {
S
Shengliang Guan 已提交
724 725
    SRespStub stub = {.createTime = taosGetTimestampMs(), .rpcMsg = *pMsg};
    uint64_t  seqNum = syncRespMgrAdd(pSyncNode->pSyncRespMgr, &stub);
726
    SRpcMsg   rpcMsg = {0};
S
Shengliang Guan 已提交
727
    int32_t   code = syncBuildClientRequest(&rpcMsg, pMsg, seqNum, isWeak, pSyncNode->vgId);
728 729 730 731
    if (code != 0) {
      sError("vgId:%d, failed to propose msg while serialize since %s", pSyncNode->vgId, terrstr());
      (void)syncRespMgrDel(pSyncNode->pSyncRespMgr, seqNum);
      return -1;
M
Minghao Li 已提交
732
    }
733

734 735 736 737 738
    sNTrace(pSyncNode, "propose msg, type:%s", TMSG_INFO(pMsg->msgType));
    code = (*pSyncNode->syncEqMsg)(pSyncNode->msgcb, &rpcMsg);
    if (code != 0) {
      sError("vgId:%d, failed to propose msg while enqueue since %s", pSyncNode->vgId, terrstr());
      (void)syncRespMgrDel(pSyncNode->pSyncRespMgr, seqNum);
M
Minghao Li 已提交
739
    }
M
Minghao Li 已提交
740

741
    return code;
M
Minghao Li 已提交
742
  }
M
Minghao Li 已提交
743 744
}

S
Shengliang Guan 已提交
745
static int32_t syncHbTimerInit(SSyncNode* pSyncNode, SSyncTimer* pSyncTimer, SRaftId destId) {
746 747 748 749 750
  pSyncTimer->pTimer = NULL;
  pSyncTimer->counter = 0;
  pSyncTimer->timerMS = pSyncNode->hbBaseLine;
  pSyncTimer->timerCb = syncNodeEqPeerHeartbeatTimer;
  pSyncTimer->destId = destId;
M
Minghao Li 已提交
751
  pSyncTimer->timeStamp = taosGetTimestampMs();
752 753 754 755
  atomic_store_64(&pSyncTimer->logicClock, 0);
  return 0;
}

S
Shengliang Guan 已提交
756
static int32_t syncHbTimerStart(SSyncNode* pSyncNode, SSyncTimer* pSyncTimer) {
757
  int32_t ret = 0;
S
Shengliang Guan 已提交
758
  int64_t tsNow = taosGetTimestampMs();
S
Shengliang Guan 已提交
759
  if (syncIsInit()) {
760 761 762 763 764 765
    SSyncHbTimerData* pData = syncHbTimerDataAcquire(pSyncTimer->hbDataRid);
    if (pData == NULL) {
      pData = taosMemoryMalloc(sizeof(SSyncHbTimerData));
      pData->rid = syncHbTimerDataAdd(pData);
    }
    pSyncTimer->hbDataRid = pData->rid;
S
Shengliang Guan 已提交
766
    pSyncTimer->timeStamp = tsNow;
767 768

    pData->syncNodeRid = pSyncNode->rid;
769 770 771
    pData->pTimer = pSyncTimer;
    pData->destId = pSyncTimer->destId;
    pData->logicClock = pSyncTimer->logicClock;
S
Shengliang Guan 已提交
772
    pData->execTime = tsNow + pSyncTimer->timerMS;
M
Minghao Li 已提交
773

774 775
    taosTmrReset(pSyncTimer->timerCb, pSyncTimer->timerMS / HEARTBEAT_TICK_NUM, (void*)(pData->rid),
                 syncEnv()->pTimerManager, &pSyncTimer->pTimer);
776 777 778 779 780 781
  } else {
    sError("vgId:%d, start ctrl hb timer error, sync env is stop", pSyncNode->vgId);
  }
  return ret;
}

S
Shengliang Guan 已提交
782
static int32_t syncHbTimerStop(SSyncNode* pSyncNode, SSyncTimer* pSyncTimer) {
783 784 785 786
  int32_t ret = 0;
  atomic_add_fetch_64(&pSyncTimer->logicClock, 1);
  taosTmrStop(pSyncTimer->pTimer);
  pSyncTimer->pTimer = NULL;
787 788
  syncHbTimerDataRemove(pSyncTimer->hbDataRid);
  pSyncTimer->hbDataRid = -1;
789 790 791
  return ret;
}

792 793 794 795 796 797 798 799 800 801 802 803 804 805 806 807 808 809 810 811 812 813
int32_t syncNodeLogStoreRestoreOnNeed(SSyncNode* pNode) {
  ASSERT(pNode->pLogStore != NULL && "log store not created");
  ASSERT(pNode->pFsm != NULL && "pFsm not registered");
  ASSERT(pNode->pFsm->FpGetSnapshotInfo != NULL && "FpGetSnapshotInfo not registered");
  SSnapshot snapshot;
  if (pNode->pFsm->FpGetSnapshotInfo(pNode->pFsm, &snapshot) < 0) {
    sError("vgId:%d, failed to get snapshot info since %s", pNode->vgId, terrstr());
    return -1;
  }
  SyncIndex commitIndex = snapshot.lastApplyIndex;
  SyncIndex firstVer = pNode->pLogStore->syncLogBeginIndex(pNode->pLogStore);
  SyncIndex lastVer = pNode->pLogStore->syncLogLastIndex(pNode->pLogStore);
  if (lastVer < commitIndex || firstVer > commitIndex + 1) {
    if (pNode->pLogStore->syncLogRestoreFromSnapshot(pNode->pLogStore, commitIndex)) {
      sError("vgId:%d, failed to restore log store from snapshot since %s. lastVer: %" PRId64 ", snapshotVer: %" PRId64,
             pNode->vgId, terrstr(), lastVer, commitIndex);
      return -1;
    }
  }
  return 0;
}

M
Minghao Li 已提交
814
// open/close --------------
S
Shengliang Guan 已提交
815 816
SSyncNode* syncNodeOpen(SSyncInfo* pSyncInfo) {
  SSyncNode* pSyncNode = taosMemoryCalloc(1, sizeof(SSyncNode));
817 818 819 820
  if (pSyncNode == NULL) {
    terrno = TSDB_CODE_OUT_OF_MEMORY;
    goto _error;
  }
M
Minghao Li 已提交
821

M
Minghao Li 已提交
822 823 824 825
  if (!taosDirExist((char*)(pSyncInfo->path))) {
    if (taosMkDir(pSyncInfo->path) != 0) {
      terrno = TAOS_SYSTEM_ERROR(errno);
      sError("failed to create dir:%s since %s", pSyncInfo->path, terrstr());
826
      goto _error;
M
Minghao Li 已提交
827
    }
828
  }
M
Minghao Li 已提交
829

S
Shengliang Guan 已提交
830
  snprintf(pSyncNode->configPath, sizeof(pSyncNode->configPath), "%s%sraft_config.json", pSyncInfo->path, TD_DIRSEP);
831
  if (!taosCheckExistFile(pSyncNode->configPath)) {
M
Minghao Li 已提交
832
    // create a new raft config file
S
Shengliang Guan 已提交
833
    SRaftCfgMeta meta = {0};
M
Minghao Li 已提交
834
    meta.isStandBy = pSyncInfo->isStandBy;
M
Minghao Li 已提交
835
    meta.snapshotStrategy = pSyncInfo->snapshotStrategy;
836
    meta.lastConfigIndex = SYNC_INDEX_INVALID;
M
Minghao Li 已提交
837
    meta.batchSize = pSyncInfo->batchSize;
S
Shengliang Guan 已提交
838 839
    if (raftCfgCreateFile(&pSyncInfo->syncCfg, meta, pSyncNode->configPath) != 0) {
      sError("vgId:%d, failed to create raft cfg file at %s", pSyncNode->vgId, pSyncNode->configPath);
H
Hongze Cheng 已提交
840
      goto _error;
841
    }
842
    if (pSyncInfo->syncCfg.replicaNum == 0) {
S
Shengliang Guan 已提交
843
      sInfo("vgId:%d, sync config not input", pSyncNode->vgId);
844 845
      pSyncInfo->syncCfg = pSyncNode->pRaftCfg->cfg;
    }
846 847 848
  } else {
    // update syncCfg by raft_config.json
    pSyncNode->pRaftCfg = raftCfgOpen(pSyncNode->configPath);
849
    if (pSyncNode->pRaftCfg == NULL) {
S
Shengliang Guan 已提交
850
      sError("vgId:%d, failed to open raft cfg file at %s", pSyncNode->vgId, pSyncNode->configPath);
H
Hongze Cheng 已提交
851
      goto _error;
852
    }
S
Shengliang Guan 已提交
853 854

    if (pSyncInfo->syncCfg.replicaNum > 0 && syncIsConfigChanged(&pSyncNode->pRaftCfg->cfg, &pSyncInfo->syncCfg)) {
S
Shengliang Guan 已提交
855 856 857 858 859 860
      sInfo("vgId:%d, use sync config from input options and write to cfg file", pSyncNode->vgId);
      pSyncNode->pRaftCfg->cfg = pSyncInfo->syncCfg;
      if (raftCfgPersist(pSyncNode->pRaftCfg) != 0) {
        sError("vgId:%d, failed to persist raft cfg file at %s", pSyncNode->vgId, pSyncNode->configPath);
        goto _error;
      }
S
Shengliang Guan 已提交
861 862 863 864
    } else {
      sInfo("vgId:%d, use sync config from raft cfg file", pSyncNode->vgId);
      pSyncInfo->syncCfg = pSyncNode->pRaftCfg->cfg;
    }
865 866

    raftCfgClose(pSyncNode->pRaftCfg);
867
    pSyncNode->pRaftCfg = NULL;
M
Minghao Li 已提交
868 869
  }

M
Minghao Li 已提交
870
  // init by SSyncInfo
M
Minghao Li 已提交
871
  pSyncNode->vgId = pSyncInfo->vgId;
S
Shengliang Guan 已提交
872 873 874 875 876 877 878
  SSyncCfg* pCfg = &pSyncInfo->syncCfg;
  sDebug("vgId:%d, replica:%d selfIndex:%d", pSyncNode->vgId, pCfg->replicaNum, pCfg->myIndex);
  for (int32_t i = 0; i < pCfg->replicaNum; ++i) {
    SNodeInfo* pNode = &pCfg->nodeInfo[i];
    sDebug("vgId:%d, index:%d ep:%s:%u", pSyncNode->vgId, i, pNode->nodeFqdn, pNode->nodePort);
  }

M
Minghao Li 已提交
879
  memcpy(pSyncNode->path, pSyncInfo->path, sizeof(pSyncNode->path));
S
Shengliang Guan 已提交
880 881 882
  snprintf(pSyncNode->raftStorePath, sizeof(pSyncNode->raftStorePath), "%s%sraft_store.json", pSyncInfo->path,
           TD_DIRSEP);
  snprintf(pSyncNode->configPath, sizeof(pSyncNode->configPath), "%s%sraft_config.json", pSyncInfo->path, TD_DIRSEP);
M
Minghao Li 已提交
883

M
Minghao Li 已提交
884
  pSyncNode->pWal = pSyncInfo->pWal;
S
Shengliang Guan 已提交
885
  pSyncNode->msgcb = pSyncInfo->msgcb;
S
Shengliang Guan 已提交
886 887 888
  pSyncNode->syncSendMSg = pSyncInfo->syncSendMSg;
  pSyncNode->syncEqMsg = pSyncInfo->syncEqMsg;
  pSyncNode->syncEqCtrlMsg = pSyncInfo->syncEqCtrlMsg;
M
Minghao Li 已提交
889

B
Benguang Zhao 已提交
890 891 892
  // create raft log ring buffer
  pSyncNode->pLogBuf = syncLogBufferCreate();
  if (pSyncNode->pLogBuf == NULL) {
893
    sError("failed to init sync log buffer since %s. vgId:%d", terrstr(), pSyncNode->vgId);
B
Benguang Zhao 已提交
894 895 896
    goto _error;
  }

M
Minghao Li 已提交
897 898
  // init raft config
  pSyncNode->pRaftCfg = raftCfgOpen(pSyncNode->configPath);
899
  if (pSyncNode->pRaftCfg == NULL) {
S
Shengliang Guan 已提交
900
    sError("vgId:%d, failed to open raft cfg file at %s", pSyncNode->vgId, pSyncNode->configPath);
901 902
    goto _error;
  }
M
Minghao Li 已提交
903

M
Minghao Li 已提交
904
  // init internal
M
Minghao Li 已提交
905
  pSyncNode->myNodeInfo = pSyncNode->pRaftCfg->cfg.nodeInfo[pSyncNode->pRaftCfg->cfg.myIndex];
906
  if (!syncUtilNodeInfo2RaftId(&pSyncNode->myNodeInfo, pSyncNode->vgId, &pSyncNode->myRaftId)) {
S
Shengliang Guan 已提交
907
    sError("vgId:%d, failed to determine my raft member id", pSyncNode->vgId);
H
Hongze Cheng 已提交
908
    goto _error;
909
  }
M
Minghao Li 已提交
910

M
Minghao Li 已提交
911
  // init peersNum, peers, peersId
M
Minghao Li 已提交
912
  pSyncNode->peersNum = pSyncNode->pRaftCfg->cfg.replicaNum - 1;
S
Shengliang Guan 已提交
913 914
  int32_t j = 0;
  for (int32_t i = 0; i < pSyncNode->pRaftCfg->cfg.replicaNum; ++i) {
M
Minghao Li 已提交
915 916
    if (i != pSyncNode->pRaftCfg->cfg.myIndex) {
      pSyncNode->peersNodeInfo[j] = pSyncNode->pRaftCfg->cfg.nodeInfo[i];
M
Minghao Li 已提交
917 918 919
      j++;
    }
  }
S
Shengliang Guan 已提交
920
  for (int32_t i = 0; i < pSyncNode->peersNum; ++i) {
921
    if (!syncUtilNodeInfo2RaftId(&pSyncNode->peersNodeInfo[i], pSyncNode->vgId, &pSyncNode->peersId[i])) {
S
Shengliang Guan 已提交
922
      sError("vgId:%d, failed to determine raft member id, peer:%d", pSyncNode->vgId, i);
H
Hongze Cheng 已提交
923
      goto _error;
924
    }
M
Minghao Li 已提交
925
  }
M
Minghao Li 已提交
926

M
Minghao Li 已提交
927
  // init replicaNum, replicasId
M
Minghao Li 已提交
928
  pSyncNode->replicaNum = pSyncNode->pRaftCfg->cfg.replicaNum;
S
Shengliang Guan 已提交
929
  for (int32_t i = 0; i < pSyncNode->pRaftCfg->cfg.replicaNum; ++i) {
930
    if (!syncUtilNodeInfo2RaftId(&pSyncNode->pRaftCfg->cfg.nodeInfo[i], pSyncNode->vgId, &pSyncNode->replicasId[i])) {
S
Shengliang Guan 已提交
931
      sError("vgId:%d, failed to determine raft member id, replica:%d", pSyncNode->vgId, i);
H
Hongze Cheng 已提交
932
      goto _error;
933
    }
M
Minghao Li 已提交
934 935
  }

M
Minghao Li 已提交
936
  // init raft algorithm
M
Minghao Li 已提交
937
  pSyncNode->pFsm = pSyncInfo->pFsm;
938
  pSyncInfo->pFsm = NULL;
M
Minghao Li 已提交
939
  pSyncNode->quorum = syncUtilQuorum(pSyncNode->pRaftCfg->cfg.replicaNum);
M
Minghao Li 已提交
940 941
  pSyncNode->leaderCache = EMPTY_RAFT_ID;

M
Minghao Li 已提交
942
  // init life cycle outside
M
Minghao Li 已提交
943

M
Minghao Li 已提交
944 945 946 947 948 949 950 951 952 953 954 955 956 957 958 959 960 961 962 963 964 965 966 967
  // TLA+ Spec
  // InitHistoryVars == /\ elections = {}
  //                    /\ allLogs   = {}
  //                    /\ voterLog  = [i \in Server |-> [j \in {} |-> <<>>]]
  // InitServerVars == /\ currentTerm = [i \in Server |-> 1]
  //                   /\ state       = [i \in Server |-> Follower]
  //                   /\ votedFor    = [i \in Server |-> Nil]
  // InitCandidateVars == /\ votesResponded = [i \in Server |-> {}]
  //                      /\ votesGranted   = [i \in Server |-> {}]
  // \* The values nextIndex[i][i] and matchIndex[i][i] are never read, since the
  // \* leader does not send itself messages. It's still easier to include these
  // \* in the functions.
  // InitLeaderVars == /\ nextIndex  = [i \in Server |-> [j \in Server |-> 1]]
  //                   /\ matchIndex = [i \in Server |-> [j \in Server |-> 0]]
  // InitLogVars == /\ log          = [i \in Server |-> << >>]
  //                /\ commitIndex  = [i \in Server |-> 0]
  // Init == /\ messages = [m \in {} |-> 0]
  //         /\ InitHistoryVars
  //         /\ InitServerVars
  //         /\ InitCandidateVars
  //         /\ InitLeaderVars
  //         /\ InitLogVars
  //

M
Minghao Li 已提交
968
  // init TLA+ server vars
M
syncInt  
Minghao Li 已提交
969
  pSyncNode->state = TAOS_SYNC_STATE_FOLLOWER;
M
Minghao Li 已提交
970
  pSyncNode->pRaftStore = raftStoreOpen(pSyncNode->raftStorePath);
971
  if (pSyncNode->pRaftStore == NULL) {
S
Shengliang Guan 已提交
972
    sError("vgId:%d, failed to open raft store at path %s", pSyncNode->vgId, pSyncNode->raftStorePath);
973 974
    goto _error;
  }
M
Minghao Li 已提交
975

M
Minghao Li 已提交
976
  // init TLA+ candidate vars
M
Minghao Li 已提交
977
  pSyncNode->pVotesGranted = voteGrantedCreate(pSyncNode);
978
  if (pSyncNode->pVotesGranted == NULL) {
S
Shengliang Guan 已提交
979
    sError("vgId:%d, failed to create VotesGranted", pSyncNode->vgId);
980 981
    goto _error;
  }
M
Minghao Li 已提交
982
  pSyncNode->pVotesRespond = votesRespondCreate(pSyncNode);
983
  if (pSyncNode->pVotesRespond == NULL) {
S
Shengliang Guan 已提交
984
    sError("vgId:%d, failed to create VotesRespond", pSyncNode->vgId);
985 986
    goto _error;
  }
M
Minghao Li 已提交
987

M
Minghao Li 已提交
988 989
  // init TLA+ leader vars
  pSyncNode->pNextIndex = syncIndexMgrCreate(pSyncNode);
990
  if (pSyncNode->pNextIndex == NULL) {
S
Shengliang Guan 已提交
991
    sError("vgId:%d, failed to create SyncIndexMgr", pSyncNode->vgId);
992 993
    goto _error;
  }
M
Minghao Li 已提交
994
  pSyncNode->pMatchIndex = syncIndexMgrCreate(pSyncNode);
995
  if (pSyncNode->pMatchIndex == NULL) {
S
Shengliang Guan 已提交
996
    sError("vgId:%d, failed to create SyncIndexMgr", pSyncNode->vgId);
997 998
    goto _error;
  }
M
Minghao Li 已提交
999 1000 1001

  // init TLA+ log vars
  pSyncNode->pLogStore = logStoreCreate(pSyncNode);
1002
  if (pSyncNode->pLogStore == NULL) {
S
Shengliang Guan 已提交
1003
    sError("vgId:%d, failed to create SyncLogStore", pSyncNode->vgId);
1004 1005
    goto _error;
  }
1006 1007 1008 1009 1010

  SyncIndex commitIndex = SYNC_INDEX_INVALID;
  if (pSyncNode->pFsm != NULL && pSyncNode->pFsm->FpGetSnapshotInfo != NULL) {
    SSnapshot snapshot = {0};
    int32_t   code = pSyncNode->pFsm->FpGetSnapshotInfo(pSyncNode->pFsm, &snapshot);
1011
    if (code != 0) {
S
Shengliang Guan 已提交
1012
      sError("vgId:%d, failed to get snapshot info, code:%d", pSyncNode->vgId, code);
H
Hongze Cheng 已提交
1013
      goto _error;
1014
    }
1015 1016
    if (snapshot.lastApplyIndex > commitIndex) {
      commitIndex = snapshot.lastApplyIndex;
S
Shengliang Guan 已提交
1017
      sNTrace(pSyncNode, "reset commit index by snapshot");
1018 1019 1020
    }
  }
  pSyncNode->commitIndex = commitIndex;
M
Minghao Li 已提交
1021

1022 1023 1024
  if (syncNodeLogStoreRestoreOnNeed(pSyncNode) < 0) {
    goto _error;
  }
M
Minghao Li 已提交
1025 1026 1027 1028 1029
  // timer ms init
  pSyncNode->pingBaseLine = PING_TIMER_MS;
  pSyncNode->electBaseLine = ELECT_TIMER_MS_MIN;
  pSyncNode->hbBaseLine = HEARTBEAT_TIMER_MS;

M
Minghao Li 已提交
1030
  // init ping timer
M
Minghao Li 已提交
1031
  pSyncNode->pPingTimer = NULL;
M
Minghao Li 已提交
1032
  pSyncNode->pingTimerMS = pSyncNode->pingBaseLine;
M
Minghao Li 已提交
1033 1034
  atomic_store_64(&pSyncNode->pingTimerLogicClock, 0);
  atomic_store_64(&pSyncNode->pingTimerLogicClockUser, 0);
M
Minghao Li 已提交
1035
  pSyncNode->FpPingTimerCB = syncNodeEqPingTimer;
M
Minghao Li 已提交
1036
  pSyncNode->pingTimerCounter = 0;
M
Minghao Li 已提交
1037

M
Minghao Li 已提交
1038 1039
  // init elect timer
  pSyncNode->pElectTimer = NULL;
M
Minghao Li 已提交
1040
  pSyncNode->electTimerMS = syncUtilElectRandomMS(pSyncNode->electBaseLine, 2 * pSyncNode->electBaseLine);
M
Minghao Li 已提交
1041
  atomic_store_64(&pSyncNode->electTimerLogicClock, 0);
M
Minghao Li 已提交
1042
  pSyncNode->FpElectTimerCB = syncNodeEqElectTimer;
M
Minghao Li 已提交
1043 1044 1045 1046
  pSyncNode->electTimerCounter = 0;

  // init heartbeat timer
  pSyncNode->pHeartbeatTimer = NULL;
M
Minghao Li 已提交
1047
  pSyncNode->heartbeatTimerMS = pSyncNode->hbBaseLine;
M
Minghao Li 已提交
1048 1049
  atomic_store_64(&pSyncNode->heartbeatTimerLogicClock, 0);
  atomic_store_64(&pSyncNode->heartbeatTimerLogicClockUser, 0);
M
Minghao Li 已提交
1050
  pSyncNode->FpHeartbeatTimerCB = syncNodeEqHeartbeatTimer;
M
Minghao Li 已提交
1051 1052
  pSyncNode->heartbeatTimerCounter = 0;

1053 1054 1055 1056 1057
  // init peer heartbeat timer
  for (int32_t i = 0; i < TSDB_MAX_REPLICA; ++i) {
    syncHbTimerInit(pSyncNode, &(pSyncNode->peerHeartbeatTimerArr[i]), (pSyncNode->replicasId)[i]);
  }

M
Minghao Li 已提交
1058
  // tools
M
Minghao Li 已提交
1059
  pSyncNode->pSyncRespMgr = syncRespMgrCreate(pSyncNode, SYNC_RESP_TTL_MS);
1060
  if (pSyncNode->pSyncRespMgr == NULL) {
S
Shengliang Guan 已提交
1061
    sError("vgId:%d, failed to create SyncRespMgr", pSyncNode->vgId);
1062 1063
    goto _error;
  }
M
Minghao Li 已提交
1064

1065 1066
  // restore state
  pSyncNode->restoreFinish = false;
1067

M
Minghao Li 已提交
1068
  // snapshot senders
S
Shengliang Guan 已提交
1069
  for (int32_t i = 0; i < TSDB_MAX_REPLICA; ++i) {
M
Minghao Li 已提交
1070 1071 1072
    SSyncSnapshotSender* pSender = snapshotSenderCreate(pSyncNode, i);
    // ASSERT(pSender != NULL);
    (pSyncNode->senders)[i] = pSender;
S
Shengliang Guan 已提交
1073
    sSTrace(pSender, "snapshot sender create new while open, data:%p", pSender);
M
Minghao Li 已提交
1074 1075 1076
  }

  // snapshot receivers
1077
  pSyncNode->pNewNodeReceiver = snapshotReceiverCreate(pSyncNode, EMPTY_RAFT_ID);
M
Minghao Li 已提交
1078

M
Minghao Li 已提交
1079 1080 1081
  // is config changing
  pSyncNode->changing = false;

B
Benguang Zhao 已提交
1082 1083 1084
  // replication mgr
  syncNodeLogReplMgrInit(pSyncNode);

M
Minghao Li 已提交
1085 1086 1087
  // peer state
  syncNodePeerStateInit(pSyncNode);

B
Benguang Zhao 已提交
1088
  //
M
Minghao Li 已提交
1089 1090 1091
  // min match index
  pSyncNode->minMatchIndex = SYNC_INDEX_INVALID;

M
Minghao Li 已提交
1092
  // start in syncNodeStart
M
Minghao Li 已提交
1093
  // start raft
M
Minghao Li 已提交
1094
  // syncNodeBecomeFollower(pSyncNode);
M
Minghao Li 已提交
1095

M
Minghao Li 已提交
1096 1097
  int64_t timeNow = taosGetTimestampMs();
  pSyncNode->startTime = timeNow;
1098
  pSyncNode->leaderTime = timeNow;
M
Minghao Li 已提交
1099 1100
  pSyncNode->lastReplicateTime = timeNow;

1101 1102 1103
  // snapshotting
  atomic_store_64(&pSyncNode->snapshottingIndex, SYNC_INDEX_INVALID);

B
Benguang Zhao 已提交
1104 1105
  // init log buffer
  if (syncLogBufferInit(pSyncNode->pLogBuf, pSyncNode) < 0) {
1106
    sError("vgId:%d, failed to init sync log buffer since %s", pSyncNode->vgId, terrstr());
1107
    goto _error;
B
Benguang Zhao 已提交
1108 1109
  }

1110
  pSyncNode->isStart = true;
1111 1112 1113
  pSyncNode->electNum = 0;
  pSyncNode->becomeLeaderNum = 0;
  pSyncNode->configChangeNum = 0;
1114 1115
  pSyncNode->hbSlowNum = 0;
  pSyncNode->hbrSlowNum = 0;
M
Minghao Li 已提交
1116
  pSyncNode->tmrRoutineNum = 0;
1117

S
Shengliang Guan 已提交
1118
  sNTrace(pSyncNode, "sync open, node:%p", pSyncNode);
1119

M
Minghao Li 已提交
1120
  return pSyncNode;
1121 1122 1123

_error:
  if (pSyncInfo->pFsm) {
H
Hongze Cheng 已提交
1124 1125
    taosMemoryFree(pSyncInfo->pFsm);
    pSyncInfo->pFsm = NULL;
1126 1127 1128 1129
  }
  syncNodeClose(pSyncNode);
  pSyncNode = NULL;
  return NULL;
M
Minghao Li 已提交
1130 1131
}

M
Minghao Li 已提交
1132 1133 1134 1135 1136 1137 1138 1139 1140 1141 1142
void syncNodeMaybeUpdateCommitBySnapshot(SSyncNode* pSyncNode) {
  if (pSyncNode->pFsm != NULL && pSyncNode->pFsm->FpGetSnapshotInfo != NULL) {
    SSnapshot snapshot;
    int32_t   code = pSyncNode->pFsm->FpGetSnapshotInfo(pSyncNode->pFsm, &snapshot);
    ASSERT(code == 0);
    if (snapshot.lastApplyIndex > pSyncNode->commitIndex) {
      pSyncNode->commitIndex = snapshot.lastApplyIndex;
    }
  }
}

B
Benguang Zhao 已提交
1143 1144 1145 1146 1147 1148 1149
int32_t syncNodeRestore(SSyncNode* pSyncNode) {
  ASSERT(pSyncNode->pLogStore != NULL && "log store not created");
  ASSERT(pSyncNode->pLogBuf != NULL && "ring log buffer not created");

  SyncIndex lastVer = pSyncNode->pLogStore->syncLogLastIndex(pSyncNode->pLogStore);
  SyncIndex commitIndex = pSyncNode->pLogStore->syncLogCommitIndex(pSyncNode->pLogStore);
  SyncIndex endIndex = pSyncNode->pLogBuf->endIndex;
1150 1151 1152 1153 1154 1155
  if (lastVer != -1 && endIndex != lastVer + 1) {
    terrno = TSDB_CODE_WAL_LOG_INCOMPLETE;
    sError("vgId:%d, failed to restore sync node since %s. expected lastLogIndex: %" PRId64 ", lastVer: %" PRId64 "",
           pSyncNode->vgId, terrstr(), endIndex - 1, lastVer);
    return -1;
  }
B
Benguang Zhao 已提交
1156

1157
  ASSERT(endIndex == lastVer + 1);
B
Benguang Zhao 已提交
1158 1159 1160 1161 1162 1163 1164 1165 1166 1167 1168 1169 1170 1171 1172 1173 1174 1175 1176 1177 1178 1179 1180 1181 1182 1183 1184 1185
  commitIndex = TMAX(pSyncNode->commitIndex, commitIndex);

  if (syncLogBufferCommit(pSyncNode->pLogBuf, pSyncNode, commitIndex) < 0) {
    return -1;
  }

  return 0;
}

int32_t syncNodeStart(SSyncNode* pSyncNode) {
  // start raft
  if (pSyncNode->replicaNum == 1) {
    raftStoreNextTerm(pSyncNode->pRaftStore);
    syncNodeBecomeLeader(pSyncNode, "one replica start");

    // Raft 3.6.2 Committing entries from previous terms
    syncNodeAppendNoop(pSyncNode);
  } else {
    syncNodeBecomeFollower(pSyncNode, "first start");
  }

  int32_t ret = 0;
  ret = syncNodeStartPingTimer(pSyncNode);
  ASSERT(ret == 0);
  return ret;
}

void syncNodeStartOld(SSyncNode* pSyncNode) {
M
Minghao Li 已提交
1186
  // start raft
1187
  if (pSyncNode->replicaNum == 1) {
M
Minghao Li 已提交
1188
    raftStoreNextTerm(pSyncNode->pRaftStore);
1189
    syncNodeBecomeLeader(pSyncNode, "one replica start");
M
format  
Minghao Li 已提交
1190

1191
    // Raft 3.6.2 Committing entries from previous terms
1192 1193
    syncNodeAppendNoop(pSyncNode);
    syncMaybeAdvanceCommitIndex(pSyncNode);
M
Minghao Li 已提交
1194

M
Minghao Li 已提交
1195 1196
  } else {
    syncNodeBecomeFollower(pSyncNode, "first start");
1197 1198
  }

1199 1200 1201
  int32_t ret = 0;
  ret = syncNodeStartPingTimer(pSyncNode);
  ASSERT(ret == 0);
M
Minghao Li 已提交
1202 1203
}

B
Benguang Zhao 已提交
1204
int32_t syncNodeStartStandBy(SSyncNode* pSyncNode) {
M
Minghao Li 已提交
1205 1206 1207 1208 1209 1210 1211 1212
  // state change
  pSyncNode->state = TAOS_SYNC_STATE_FOLLOWER;
  syncNodeStopHeartbeatTimer(pSyncNode);

  // reset elect timer, long enough
  int32_t electMS = TIMER_MAX_MS;
  int32_t ret = syncNodeRestartElectTimer(pSyncNode, electMS);
  ASSERT(ret == 0);
1213

1214 1215 1216
  ret = 0;
  ret = syncNodeStartPingTimer(pSyncNode);
  ASSERT(ret == 0);
B
Benguang Zhao 已提交
1217
  return ret;
M
Minghao Li 已提交
1218 1219
}

M
Minghao Li 已提交
1220 1221 1222 1223 1224 1225 1226 1227
void syncNodePreClose(SSyncNode* pSyncNode) {
  // stop elect timer
  syncNodeStopElectTimer(pSyncNode);

  // stop heartbeat timer
  syncNodeStopHeartbeatTimer(pSyncNode);
}

1228
void syncHbTimerDataFree(SSyncHbTimerData* pData) { taosMemoryFree(pData); }
M
Minghao Li 已提交
1229

M
Minghao Li 已提交
1230
void syncNodeClose(SSyncNode* pSyncNode) {
S
Shengliang Guan 已提交
1231 1232
  if (pSyncNode == NULL) return;
  sNTrace(pSyncNode, "sync close, data:%p", pSyncNode);
M
Minghao Li 已提交
1233

S
Shengliang Guan 已提交
1234
  int32_t ret = raftStoreClose(pSyncNode->pRaftStore);
M
Minghao Li 已提交
1235
  ASSERT(ret == 0);
M
Minghao Li 已提交
1236
  pSyncNode->pRaftStore = NULL;
M
Minghao Li 已提交
1237

B
Benguang Zhao 已提交
1238
  syncNodeLogReplMgrDestroy(pSyncNode);
M
Minghao Li 已提交
1239
  syncRespMgrDestroy(pSyncNode->pSyncRespMgr);
1240
  pSyncNode->pSyncRespMgr = NULL;
M
Minghao Li 已提交
1241
  voteGrantedDestroy(pSyncNode->pVotesGranted);
1242
  pSyncNode->pVotesGranted = NULL;
M
Minghao Li 已提交
1243
  votesRespondDestory(pSyncNode->pVotesRespond);
1244
  pSyncNode->pVotesRespond = NULL;
M
Minghao Li 已提交
1245
  syncIndexMgrDestroy(pSyncNode->pNextIndex);
1246
  pSyncNode->pNextIndex = NULL;
M
Minghao Li 已提交
1247
  syncIndexMgrDestroy(pSyncNode->pMatchIndex);
1248
  pSyncNode->pMatchIndex = NULL;
M
Minghao Li 已提交
1249
  logStoreDestory(pSyncNode->pLogStore);
1250
  pSyncNode->pLogStore = NULL;
B
Benguang Zhao 已提交
1251 1252
  syncLogBufferDestroy(pSyncNode->pLogBuf);
  pSyncNode->pLogBuf = NULL;
M
Minghao Li 已提交
1253
  raftCfgClose(pSyncNode->pRaftCfg);
1254
  pSyncNode->pRaftCfg = NULL;
M
Minghao Li 已提交
1255 1256 1257 1258 1259

  syncNodeStopPingTimer(pSyncNode);
  syncNodeStopElectTimer(pSyncNode);
  syncNodeStopHeartbeatTimer(pSyncNode);

M
Minghao Li 已提交
1260 1261 1262 1263
  if (pSyncNode->pFsm != NULL) {
    taosMemoryFree(pSyncNode->pFsm);
  }

S
Shengliang Guan 已提交
1264
  for (int32_t i = 0; i < TSDB_MAX_REPLICA; ++i) {
M
Minghao Li 已提交
1265
    if ((pSyncNode->senders)[i] != NULL) {
S
Shengliang Guan 已提交
1266
      sSTrace((pSyncNode->senders)[i], "snapshot sender destroy while close, data:%p", (pSyncNode->senders)[i]);
1267 1268 1269 1270 1271

      if (snapshotSenderIsStart((pSyncNode->senders)[i])) {
        snapshotSenderStop((pSyncNode->senders)[i], false);
      }

M
Minghao Li 已提交
1272 1273 1274 1275 1276
      snapshotSenderDestroy((pSyncNode->senders)[i]);
      (pSyncNode->senders)[i] = NULL;
    }
  }

M
Minghao Li 已提交
1277
  if (pSyncNode->pNewNodeReceiver != NULL) {
1278 1279 1280 1281
    if (snapshotReceiverIsStart(pSyncNode->pNewNodeReceiver)) {
      snapshotReceiverForceStop(pSyncNode->pNewNodeReceiver);
    }

M
Minghao Li 已提交
1282 1283 1284 1285
    snapshotReceiverDestroy(pSyncNode->pNewNodeReceiver);
    pSyncNode->pNewNodeReceiver = NULL;
  }

1286
  taosMemoryFree(pSyncNode);
M
Minghao Li 已提交
1287 1288
}

M
Minghao Li 已提交
1289
ESyncStrategy syncNodeStrategy(SSyncNode* pSyncNode) { return pSyncNode->pRaftCfg->snapshotStrategy; }
M
Minghao Li 已提交
1290

M
Minghao Li 已提交
1291 1292 1293
// timer control --------------
int32_t syncNodeStartPingTimer(SSyncNode* pSyncNode) {
  int32_t ret = 0;
S
Shengliang Guan 已提交
1294 1295
  if (syncIsInit()) {
    taosTmrReset(pSyncNode->FpPingTimerCB, pSyncNode->pingTimerMS, pSyncNode, syncEnv()->pTimerManager,
1296 1297 1298
                 &pSyncNode->pPingTimer);
    atomic_store_64(&pSyncNode->pingTimerLogicClock, pSyncNode->pingTimerLogicClockUser);
  } else {
M
Minghao Li 已提交
1299
    sError("vgId:%d, start ping timer error, sync env is stop", pSyncNode->vgId);
1300
  }
M
Minghao Li 已提交
1301 1302 1303 1304 1305 1306 1307 1308 1309 1310 1311 1312 1313
  return ret;
}

int32_t syncNodeStopPingTimer(SSyncNode* pSyncNode) {
  int32_t ret = 0;
  atomic_add_fetch_64(&pSyncNode->pingTimerLogicClockUser, 1);
  taosTmrStop(pSyncNode->pPingTimer);
  pSyncNode->pPingTimer = NULL;
  return ret;
}

int32_t syncNodeStartElectTimer(SSyncNode* pSyncNode, int32_t ms) {
  int32_t ret = 0;
S
Shengliang Guan 已提交
1314
  if (syncIsInit()) {
1315
    pSyncNode->electTimerMS = ms;
S
Shengliang Guan 已提交
1316

1317 1318 1319 1320 1321
    int64_t execTime = taosGetTimestampMs() + ms;
    atomic_store_64(&(pSyncNode->electTimerParam.executeTime), execTime);
    atomic_store_64(&(pSyncNode->electTimerParam.logicClock), pSyncNode->electTimerLogicClock);
    pSyncNode->electTimerParam.pSyncNode = pSyncNode;
    pSyncNode->electTimerParam.pData = NULL;
S
Shengliang Guan 已提交
1322

M
Minghao Li 已提交
1323
    taosTmrReset(pSyncNode->FpElectTimerCB, pSyncNode->electTimerMS, (void*)(pSyncNode->rid), syncEnv()->pTimerManager,
1324
                 &pSyncNode->pElectTimer);
1325

1326
  } else {
M
Minghao Li 已提交
1327
    sError("vgId:%d, start elect timer error, sync env is stop", pSyncNode->vgId);
1328
  }
M
Minghao Li 已提交
1329 1330 1331 1332 1333
  return ret;
}

int32_t syncNodeStopElectTimer(SSyncNode* pSyncNode) {
  int32_t ret = 0;
M
Minghao Li 已提交
1334
  atomic_add_fetch_64(&pSyncNode->electTimerLogicClock, 1);
M
Minghao Li 已提交
1335 1336
  taosTmrStop(pSyncNode->pElectTimer);
  pSyncNode->pElectTimer = NULL;
1337

M
Minghao Li 已提交
1338 1339 1340 1341 1342 1343 1344 1345 1346 1347
  return ret;
}

int32_t syncNodeRestartElectTimer(SSyncNode* pSyncNode, int32_t ms) {
  int32_t ret = 0;
  syncNodeStopElectTimer(pSyncNode);
  syncNodeStartElectTimer(pSyncNode, ms);
  return ret;
}

M
Minghao Li 已提交
1348 1349
int32_t syncNodeResetElectTimer(SSyncNode* pSyncNode) {
  int32_t ret = 0;
M
Minghao Li 已提交
1350 1351 1352 1353 1354 1355 1356
  int32_t electMS;

  if (pSyncNode->pRaftCfg->isStandBy) {
    electMS = TIMER_MAX_MS;
  } else {
    electMS = syncUtilElectRandomMS(pSyncNode->electBaseLine, 2 * pSyncNode->electBaseLine);
  }
M
Minghao Li 已提交
1357
  ret = syncNodeRestartElectTimer(pSyncNode, electMS);
1358

S
Shengliang Guan 已提交
1359 1360
  sNTrace(pSyncNode, "reset elect timer, min:%d, max:%d, ms:%d", pSyncNode->electBaseLine, 2 * pSyncNode->electBaseLine,
          electMS);
M
Minghao Li 已提交
1361 1362 1363
  return ret;
}

M
Minghao Li 已提交
1364
static int32_t syncNodeDoStartHeartbeatTimer(SSyncNode* pSyncNode) {
M
Minghao Li 已提交
1365
  int32_t ret = 0;
S
Shengliang Guan 已提交
1366 1367
  if (syncIsInit()) {
    taosTmrReset(pSyncNode->FpHeartbeatTimerCB, pSyncNode->heartbeatTimerMS, pSyncNode, syncEnv()->pTimerManager,
1368 1369 1370
                 &pSyncNode->pHeartbeatTimer);
    atomic_store_64(&pSyncNode->heartbeatTimerLogicClock, pSyncNode->heartbeatTimerLogicClockUser);
  } else {
M
Minghao Li 已提交
1371
    sError("vgId:%d, start heartbeat timer error, sync env is stop", pSyncNode->vgId);
1372
  }
1373

S
Shengliang Guan 已提交
1374
  sNTrace(pSyncNode, "start heartbeat timer, ms:%d", pSyncNode->heartbeatTimerMS);
M
Minghao Li 已提交
1375 1376 1377
  return ret;
}

M
Minghao Li 已提交
1378
int32_t syncNodeStartHeartbeatTimer(SSyncNode* pSyncNode) {
1379
  int32_t ret = 0;
M
Minghao Li 已提交
1380

1381
#if 0
M
Minghao Li 已提交
1382
  pSyncNode->heartbeatTimerMS = pSyncNode->hbBaseLine;
1383 1384
  ret = syncNodeDoStartHeartbeatTimer(pSyncNode);
#endif
1385

S
Shengliang Guan 已提交
1386
  for (int32_t i = 0; i < pSyncNode->peersNum; ++i) {
1387
    SSyncTimer* pSyncTimer = syncNodeGetHbTimer(pSyncNode, &(pSyncNode->peersId[i]));
M
Minghao Li 已提交
1388 1389 1390
    if (pSyncTimer != NULL) {
      syncHbTimerStart(pSyncNode, pSyncTimer);
    }
1391
  }
1392

M
Minghao Li 已提交
1393 1394 1395
  return ret;
}

M
Minghao Li 已提交
1396 1397
int32_t syncNodeStopHeartbeatTimer(SSyncNode* pSyncNode) {
  int32_t ret = 0;
1398 1399

#if 0
M
Minghao Li 已提交
1400 1401 1402
  atomic_add_fetch_64(&pSyncNode->heartbeatTimerLogicClockUser, 1);
  taosTmrStop(pSyncNode->pHeartbeatTimer);
  pSyncNode->pHeartbeatTimer = NULL;
1403
#endif
1404

S
Shengliang Guan 已提交
1405
  for (int32_t i = 0; i < pSyncNode->peersNum; ++i) {
1406
    SSyncTimer* pSyncTimer = syncNodeGetHbTimer(pSyncNode, &(pSyncNode->peersId[i]));
M
Minghao Li 已提交
1407 1408 1409
    if (pSyncTimer != NULL) {
      syncHbTimerStop(pSyncNode, pSyncTimer);
    }
1410
  }
1411

M
Minghao Li 已提交
1412 1413 1414
  return ret;
}

1415 1416 1417 1418 1419 1420
int32_t syncNodeRestartHeartbeatTimer(SSyncNode* pSyncNode) {
  syncNodeStopHeartbeatTimer(pSyncNode);
  syncNodeStartHeartbeatTimer(pSyncNode);
  return 0;
}

M
Minghao Li 已提交
1421 1422 1423
// utils --------------
int32_t syncNodeSendMsgById(const SRaftId* destRaftId, SSyncNode* pSyncNode, SRpcMsg* pMsg) {
  SEpSet epSet;
1424
  syncUtilRaftId2EpSet(destRaftId, &epSet);
S
Shengliang Guan 已提交
1425
  if (pSyncNode->syncSendMSg != NULL) {
M
Minghao Li 已提交
1426 1427 1428
    // htonl
    syncUtilMsgHtoN(pMsg->pCont);

1429
    pMsg->info.noResp = 1;
S
Shengliang Guan 已提交
1430
    pSyncNode->syncSendMSg(&epSet, pMsg);
M
Minghao Li 已提交
1431
  } else {
M
Minghao Li 已提交
1432
    sError("vgId:%d, sync send msg by id error, fp-send-msg is null", pSyncNode->vgId);
S
Shengliang Guan 已提交
1433
    rpcFreeCont(pMsg->pCont);
M
Minghao Li 已提交
1434
    return -1;
M
Minghao Li 已提交
1435
  }
M
Minghao Li 已提交
1436

M
Minghao Li 已提交
1437 1438 1439 1440 1441
  return 0;
}

int32_t syncNodeSendMsgByInfo(const SNodeInfo* nodeInfo, SSyncNode* pSyncNode, SRpcMsg* pMsg) {
  SEpSet epSet;
1442
  syncUtilNodeInfo2EpSet(nodeInfo, &epSet);
S
Shengliang Guan 已提交
1443
  if (pSyncNode->syncSendMSg != NULL) {
M
Minghao Li 已提交
1444 1445 1446
    // htonl
    syncUtilMsgHtoN(pMsg->pCont);

1447
    pMsg->info.noResp = 1;
S
Shengliang Guan 已提交
1448
    pSyncNode->syncSendMSg(&epSet, pMsg);
M
Minghao Li 已提交
1449
  } else {
M
Minghao Li 已提交
1450
    sError("vgId:%d, sync send msg by info error, fp-send-msg is null", pSyncNode->vgId);
M
Minghao Li 已提交
1451
  }
M
Minghao Li 已提交
1452 1453 1454
  return 0;
}

1455
inline bool syncNodeInConfig(SSyncNode* pSyncNode, const SSyncCfg* config) {
1456 1457 1458
  bool b1 = false;
  bool b2 = false;

S
Shengliang Guan 已提交
1459
  for (int32_t i = 0; i < config->replicaNum; ++i) {
1460 1461 1462 1463 1464 1465 1466
    if (strcmp((config->nodeInfo)[i].nodeFqdn, pSyncNode->myNodeInfo.nodeFqdn) == 0 &&
        (config->nodeInfo)[i].nodePort == pSyncNode->myNodeInfo.nodePort) {
      b1 = true;
      break;
    }
  }

S
Shengliang Guan 已提交
1467
  for (int32_t i = 0; i < config->replicaNum; ++i) {
1468 1469 1470 1471 1472 1473 1474 1475 1476 1477 1478 1479 1480 1481
    SRaftId raftId;
    raftId.addr = syncUtilAddr2U64((config->nodeInfo)[i].nodeFqdn, (config->nodeInfo)[i].nodePort);
    raftId.vgId = pSyncNode->vgId;

    if (syncUtilSameId(&raftId, &(pSyncNode->myRaftId))) {
      b2 = true;
      break;
    }
  }

  ASSERT(b1 == b2);
  return b1;
}

1482 1483 1484 1485 1486 1487 1488 1489 1490 1491 1492 1493 1494
static bool syncIsConfigChanged(const SSyncCfg* pOldCfg, const SSyncCfg* pNewCfg) {
  if (pOldCfg->replicaNum != pNewCfg->replicaNum) return true;
  if (pOldCfg->myIndex != pNewCfg->myIndex) return true;
  for (int32_t i = 0; i < pOldCfg->replicaNum; ++i) {
    const SNodeInfo* pOldInfo = &pOldCfg->nodeInfo[i];
    const SNodeInfo* pNewInfo = &pNewCfg->nodeInfo[i];
    if (strcmp(pOldInfo->nodeFqdn, pNewInfo->nodeFqdn) != 0) return true;
    if (pOldInfo->nodePort != pNewInfo->nodePort) return true;
  }

  return false;
}

M
Minghao Li 已提交
1495
void syncNodeDoConfigChange(SSyncNode* pSyncNode, SSyncCfg* pNewConfig, SyncIndex lastConfigChangeIndex) {
1496
  SSyncCfg oldConfig = pSyncNode->pRaftCfg->cfg;
1497 1498 1499 1500
  if (!syncIsConfigChanged(&oldConfig, pNewConfig)) {
    sInfo("vgId:1, sync not reconfig since not changed");
    return;
  }
S
Shengliang Guan 已提交
1501

1502
  pSyncNode->pRaftCfg->cfg = *pNewConfig;
1503 1504
  pSyncNode->pRaftCfg->lastConfigIndex = lastConfigChangeIndex;

1505 1506
  pSyncNode->configChangeNum++;

M
Minghao Li 已提交
1507 1508
  bool IamInOld = syncNodeInConfig(pSyncNode, &oldConfig);
  bool IamInNew = syncNodeInConfig(pSyncNode, pNewConfig);
M
Minghao Li 已提交
1509

M
Minghao Li 已提交
1510 1511
  bool isDrop = false;
  bool isAdd = false;
M
Minghao Li 已提交
1512

M
Minghao Li 已提交
1513 1514 1515 1516
  if (IamInOld && !IamInNew) {
    isDrop = true;
  } else {
    isDrop = false;
1517
  }
1518

M
Minghao Li 已提交
1519 1520 1521 1522 1523
  if (!IamInOld && IamInNew) {
    isAdd = true;
  } else {
    isAdd = false;
  }
M
Minghao Li 已提交
1524

M
Minghao Li 已提交
1525
  // log begin config change
S
Shengliang Guan 已提交
1526 1527 1528 1529
  char oldCfgStr[1024] = {0};
  char newCfgStr[1024] = {0};
  syncCfg2SimpleStr(&oldConfig, oldCfgStr, sizeof(oldCfgStr));
  syncCfg2SimpleStr(pNewConfig, oldCfgStr, sizeof(oldCfgStr));
1530
  sNInfo(pSyncNode, "begin do config change, from %s to %s", oldCfgStr, oldCfgStr);
M
Minghao Li 已提交
1531

M
Minghao Li 已提交
1532 1533
  if (IamInNew) {
    pSyncNode->pRaftCfg->isStandBy = 0;  // change isStandBy to normal
M
Minghao Li 已提交
1534
  }
M
Minghao Li 已提交
1535 1536
  if (isDrop) {
    pSyncNode->pRaftCfg->isStandBy = 1;  // set standby
M
Minghao Li 已提交
1537 1538
  }

M
Minghao Li 已提交
1539
  // add last config index
M
Minghao Li 已提交
1540
  raftCfgAddConfigIndex(pSyncNode->pRaftCfg, lastConfigChangeIndex);
M
Minghao Li 已提交
1541

M
Minghao Li 已提交
1542 1543 1544 1545 1546 1547 1548 1549 1550
  if (IamInNew) {
    //-----------------------------------------
    int32_t ret = 0;

    // save snapshot senders
    int32_t oldReplicaNum = pSyncNode->replicaNum;
    SRaftId oldReplicasId[TSDB_MAX_REPLICA];
    memcpy(oldReplicasId, pSyncNode->replicasId, sizeof(oldReplicasId));
    SSyncSnapshotSender* oldSenders[TSDB_MAX_REPLICA];
S
Shengliang Guan 已提交
1551
    for (int32_t i = 0; i < TSDB_MAX_REPLICA; ++i) {
M
Minghao Li 已提交
1552
      oldSenders[i] = (pSyncNode->senders)[i];
S
Shengliang Guan 已提交
1553
      sSTrace(oldSenders[i], "snapshot sender save old");
M
Minghao Li 已提交
1554
    }
1555

M
Minghao Li 已提交
1556 1557
    // init internal
    pSyncNode->myNodeInfo = pSyncNode->pRaftCfg->cfg.nodeInfo[pSyncNode->pRaftCfg->cfg.myIndex];
1558
    syncUtilNodeInfo2RaftId(&pSyncNode->myNodeInfo, pSyncNode->vgId, &pSyncNode->myRaftId);
M
Minghao Li 已提交
1559 1560 1561

    // init peersNum, peers, peersId
    pSyncNode->peersNum = pSyncNode->pRaftCfg->cfg.replicaNum - 1;
S
Shengliang Guan 已提交
1562 1563
    int32_t j = 0;
    for (int32_t i = 0; i < pSyncNode->pRaftCfg->cfg.replicaNum; ++i) {
M
Minghao Li 已提交
1564 1565 1566 1567 1568
      if (i != pSyncNode->pRaftCfg->cfg.myIndex) {
        pSyncNode->peersNodeInfo[j] = pSyncNode->pRaftCfg->cfg.nodeInfo[i];
        j++;
      }
    }
S
Shengliang Guan 已提交
1569
    for (int32_t i = 0; i < pSyncNode->peersNum; ++i) {
1570
      syncUtilNodeInfo2RaftId(&pSyncNode->peersNodeInfo[i], pSyncNode->vgId, &pSyncNode->peersId[i]);
M
Minghao Li 已提交
1571
    }
1572

M
Minghao Li 已提交
1573 1574
    // init replicaNum, replicasId
    pSyncNode->replicaNum = pSyncNode->pRaftCfg->cfg.replicaNum;
S
Shengliang Guan 已提交
1575
    for (int32_t i = 0; i < pSyncNode->pRaftCfg->cfg.replicaNum; ++i) {
1576
      syncUtilNodeInfo2RaftId(&pSyncNode->pRaftCfg->cfg.nodeInfo[i], pSyncNode->vgId, &pSyncNode->replicasId[i]);
M
Minghao Li 已提交
1577
    }
1578

1579 1580 1581
    // update quorum first
    pSyncNode->quorum = syncUtilQuorum(pSyncNode->pRaftCfg->cfg.replicaNum);

M
Minghao Li 已提交
1582 1583 1584 1585
    syncIndexMgrUpdate(pSyncNode->pNextIndex, pSyncNode);
    syncIndexMgrUpdate(pSyncNode->pMatchIndex, pSyncNode);
    voteGrantedUpdate(pSyncNode->pVotesGranted, pSyncNode);
    votesRespondUpdate(pSyncNode->pVotesRespond, pSyncNode);
M
Minghao Li 已提交
1586

M
Minghao Li 已提交
1587
    // reset snapshot senders
1588

M
Minghao Li 已提交
1589
    // clear new
S
Shengliang Guan 已提交
1590
    for (int32_t i = 0; i < TSDB_MAX_REPLICA; ++i) {
M
Minghao Li 已提交
1591 1592
      (pSyncNode->senders)[i] = NULL;
    }
M
Minghao Li 已提交
1593

M
Minghao Li 已提交
1594
    // reset new
S
Shengliang Guan 已提交
1595
    for (int32_t i = 0; i < pSyncNode->replicaNum; ++i) {
M
Minghao Li 已提交
1596 1597
      // reset sender
      bool reset = false;
S
Shengliang Guan 已提交
1598
      for (int32_t j = 0; j < TSDB_MAX_REPLICA; ++j) {
M
Minghao Li 已提交
1599
        if (syncUtilSameId(&(pSyncNode->replicasId)[i], &oldReplicasId[j]) && oldSenders[j] != NULL) {
M
Minghao Li 已提交
1600 1601 1602
          char     host[128];
          uint16_t port;
          syncUtilU642Addr((pSyncNode->replicasId)[i].addr, host, sizeof(host), &port);
1603
          sNTrace(pSyncNode, "snapshot sender reset for: %" PRId64 ", newIndex:%d, %s:%d, %p",
S
Shengliang Guan 已提交
1604
                  (pSyncNode->replicasId)[i].addr, i, host, port, oldSenders[j]);
M
Minghao Li 已提交
1605 1606 1607 1608 1609 1610 1611 1612 1613

          (pSyncNode->senders)[i] = oldSenders[j];
          oldSenders[j] = NULL;
          reset = true;

          // reset replicaIndex
          int32_t oldreplicaIndex = (pSyncNode->senders)[i]->replicaIndex;
          (pSyncNode->senders)[i]->replicaIndex = i;

S
Shengliang Guan 已提交
1614 1615
          sNTrace(pSyncNode, "snapshot sender udpate replicaIndex from %d to %d, %s:%d, %p, reset:%d", oldreplicaIndex,
                  i, host, port, (pSyncNode->senders)[i], reset);
M
Minghao Li 已提交
1616 1617

          break;
M
Minghao Li 已提交
1618
        }
1619 1620
      }
    }
1621

M
Minghao Li 已提交
1622
    // create new
S
Shengliang Guan 已提交
1623
    for (int32_t i = 0; i < TSDB_MAX_REPLICA; ++i) {
M
Minghao Li 已提交
1624 1625
      if ((pSyncNode->senders)[i] == NULL) {
        (pSyncNode->senders)[i] = snapshotSenderCreate(pSyncNode, i);
S
Shengliang Guan 已提交
1626 1627 1628
        sSTrace((pSyncNode->senders)[i], "snapshot sender create new while reconfig, data:%p", (pSyncNode->senders)[i]);
      } else {
        sSTrace((pSyncNode->senders)[i], "snapshot sender already exist, data:%p", (pSyncNode->senders)[i]);
M
Minghao Li 已提交
1629
      }
1630 1631
    }

M
Minghao Li 已提交
1632
    // free old
S
Shengliang Guan 已提交
1633
    for (int32_t i = 0; i < TSDB_MAX_REPLICA; ++i) {
M
Minghao Li 已提交
1634
      if (oldSenders[i] != NULL) {
S
Shengliang Guan 已提交
1635
        sNTrace(pSyncNode, "snapshot sender destroy old, data:%p replica-index:%d", oldSenders[i], i);
M
Minghao Li 已提交
1636 1637 1638
        snapshotSenderDestroy(oldSenders[i]);
        oldSenders[i] = NULL;
      }
1639 1640
    }

1641
    // persist cfg
M
Minghao Li 已提交
1642
    raftCfgPersist(pSyncNode->pRaftCfg);
1643

S
Shengliang Guan 已提交
1644
    char tmpbuf[1024] = {0};
1645
    snprintf(tmpbuf, sizeof(tmpbuf), "config change from %d to %d, index:%" PRId64 ", %s  -->  %s",
S
Shengliang Guan 已提交
1646
             oldConfig.replicaNum, pNewConfig->replicaNum, lastConfigChangeIndex, oldCfgStr, newCfgStr);
M
Minghao Li 已提交
1647

M
Minghao Li 已提交
1648 1649 1650
    // change isStandBy to normal (election timeout)
    if (pSyncNode->state == TAOS_SYNC_STATE_LEADER) {
      syncNodeBecomeLeader(pSyncNode, tmpbuf);
1651 1652 1653

      // Raft 3.6.2 Committing entries from previous terms
      syncNodeAppendNoop(pSyncNode);
1654
      // syncMaybeAdvanceCommitIndex(pSyncNode);
1655

M
Minghao Li 已提交
1656 1657 1658 1659
    } else {
      syncNodeBecomeFollower(pSyncNode, tmpbuf);
    }
  } else {
1660
    // persist cfg
M
Minghao Li 已提交
1661
    raftCfgPersist(pSyncNode->pRaftCfg);
1662 1663
    sNInfo(pSyncNode, "do not config change from %d to %d, index:%" PRId64 ", %s  -->  %s", oldConfig.replicaNum,
           pNewConfig->replicaNum, lastConfigChangeIndex, oldCfgStr, newCfgStr);
1664
  }
1665

M
Minghao Li 已提交
1666
_END:
M
Minghao Li 已提交
1667
  // log end config change
1668
  sNInfo(pSyncNode, "end do config change, from %s to %s", oldCfgStr, newCfgStr);
M
Minghao Li 已提交
1669 1670
}

M
Minghao Li 已提交
1671 1672 1673 1674
// raft state change --------------
void syncNodeUpdateTerm(SSyncNode* pSyncNode, SyncTerm term) {
  if (term > pSyncNode->pRaftStore->currentTerm) {
    raftStoreSetTerm(pSyncNode->pRaftStore, term);
1675
    char tmpBuf[64];
1676
    snprintf(tmpBuf, sizeof(tmpBuf), "update term to %" PRId64, term);
1677
    syncNodeBecomeFollower(pSyncNode, tmpBuf);
M
Minghao Li 已提交
1678 1679 1680 1681
    raftStoreClearVote(pSyncNode->pRaftStore);
  }
}

1682 1683 1684 1685 1686 1687
void syncNodeUpdateTermWithoutStepDown(SSyncNode* pSyncNode, SyncTerm term) {
  if (term > pSyncNode->pRaftStore->currentTerm) {
    raftStoreSetTerm(pSyncNode->pRaftStore, term);
  }
}

M
Minghao Li 已提交
1688
void syncNodeStepDown(SSyncNode* pSyncNode, SyncTerm newTerm) {
M
Minghao Li 已提交
1689
  if (pSyncNode->pRaftStore->currentTerm > newTerm) {
1690
    sNTrace(pSyncNode, "step down, ignore, new-term:%" PRId64 ", current-term:%" PRId64, newTerm,
S
Shengliang Guan 已提交
1691
            pSyncNode->pRaftStore->currentTerm);
M
Minghao Li 已提交
1692 1693
    return;
  }
M
Minghao Li 已提交
1694 1695

  do {
1696
    sNTrace(pSyncNode, "step down, new-term:%" PRId64 ", current-term:%" PRId64, newTerm,
S
Shengliang Guan 已提交
1697
            pSyncNode->pRaftStore->currentTerm);
M
Minghao Li 已提交
1698 1699 1700 1701 1702
  } while (0);

  if (pSyncNode->pRaftStore->currentTerm < newTerm) {
    raftStoreSetTerm(pSyncNode->pRaftStore, newTerm);
    char tmpBuf[64];
1703
    snprintf(tmpBuf, sizeof(tmpBuf), "step down, update term to %" PRId64, newTerm);
M
Minghao Li 已提交
1704 1705 1706 1707 1708 1709 1710 1711 1712 1713
    syncNodeBecomeFollower(pSyncNode, tmpBuf);
    raftStoreClearVote(pSyncNode->pRaftStore);

  } else {
    if (pSyncNode->state != TAOS_SYNC_STATE_FOLLOWER) {
      syncNodeBecomeFollower(pSyncNode, "step down");
    }
  }
}

1714 1715
void syncNodeLeaderChangeRsp(SSyncNode* pSyncNode) { syncRespCleanRsp(pSyncNode->pSyncRespMgr); }

1716
void syncNodeBecomeFollower(SSyncNode* pSyncNode, const char* debugStr) {
M
Minghao Li 已提交
1717
  // maybe clear leader cache
M
Minghao Li 已提交
1718 1719 1720 1721
  if (pSyncNode->state == TAOS_SYNC_STATE_LEADER) {
    pSyncNode->leaderCache = EMPTY_RAFT_ID;
  }

1722 1723
  pSyncNode->hbSlowNum = 0;

M
Minghao Li 已提交
1724
  // state change
M
Minghao Li 已提交
1725 1726 1727
  pSyncNode->state = TAOS_SYNC_STATE_FOLLOWER;
  syncNodeStopHeartbeatTimer(pSyncNode);

M
Minghao Li 已提交
1728 1729
  // reset elect timer
  syncNodeResetElectTimer(pSyncNode);
M
Minghao Li 已提交
1730

1731 1732 1733
  // send rsp to client
  syncNodeLeaderChangeRsp(pSyncNode);

1734 1735 1736 1737 1738
  // call back
  if (pSyncNode->pFsm != NULL && pSyncNode->pFsm->FpBecomeFollowerCb != NULL) {
    pSyncNode->pFsm->FpBecomeFollowerCb(pSyncNode->pFsm);
  }

M
Minghao Li 已提交
1739 1740 1741
  // min match index
  pSyncNode->minMatchIndex = SYNC_INDEX_INVALID;

B
Benguang Zhao 已提交
1742 1743 1744
  // reset log buffer
  syncLogBufferReset(pSyncNode->pLogBuf, pSyncNode);

M
Minghao Li 已提交
1745
  // trace log
S
Shengliang Guan 已提交
1746
  sNTrace(pSyncNode, "become follower %s", debugStr);
M
Minghao Li 已提交
1747 1748 1749 1750 1751 1752 1753 1754 1755 1756 1757 1758 1759 1760 1761 1762 1763 1764 1765 1766
}

// TLA+ Spec
// \* Candidate i transitions to leader.
// BecomeLeader(i) ==
//     /\ state[i] = Candidate
//     /\ votesGranted[i] \in Quorum
//     /\ state'      = [state EXCEPT ![i] = Leader]
//     /\ nextIndex'  = [nextIndex EXCEPT ![i] =
//                          [j \in Server |-> Len(log[i]) + 1]]
//     /\ matchIndex' = [matchIndex EXCEPT ![i] =
//                          [j \in Server |-> 0]]
//     /\ elections'  = elections \cup
//                          {[eterm     |-> currentTerm[i],
//                            eleader   |-> i,
//                            elog      |-> log[i],
//                            evotes    |-> votesGranted[i],
//                            evoterLog |-> voterLog[i]]}
//     /\ UNCHANGED <<messages, currentTerm, votedFor, candidateVars, logVars>>
//
1767
void syncNodeBecomeLeader(SSyncNode* pSyncNode, const char* debugStr) {
1768 1769
  pSyncNode->leaderTime = taosGetTimestampMs();

1770
  pSyncNode->becomeLeaderNum++;
1771
  pSyncNode->hbrSlowNum = 0;
1772

1773 1774 1775
  // reset restoreFinish
  pSyncNode->restoreFinish = false;

M
Minghao Li 已提交
1776
  // state change
M
Minghao Li 已提交
1777
  pSyncNode->state = TAOS_SYNC_STATE_LEADER;
M
Minghao Li 已提交
1778 1779

  // set leader cache
M
Minghao Li 已提交
1780 1781
  pSyncNode->leaderCache = pSyncNode->myRaftId;

S
Shengliang Guan 已提交
1782
  for (int32_t i = 0; i < pSyncNode->pNextIndex->replicaNum; ++i) {
M
Minghao Li 已提交
1783 1784
    // maybe overwrite myself, no harm
    // just do it!
1785 1786 1787 1788 1789 1790 1791 1792 1793

    // pSyncNode->pNextIndex->index[i] = pSyncNode->pLogStore->getLastIndex(pSyncNode->pLogStore) + 1;

    // maybe wal is deleted
    SyncIndex lastIndex;
    SyncTerm  lastTerm;
    int32_t   code = syncNodeGetLastIndexTerm(pSyncNode, &lastIndex, &lastTerm);
    ASSERT(code == 0);
    pSyncNode->pNextIndex->index[i] = lastIndex + 1;
M
Minghao Li 已提交
1794 1795
  }

S
Shengliang Guan 已提交
1796
  for (int32_t i = 0; i < pSyncNode->pMatchIndex->replicaNum; ++i) {
M
Minghao Li 已提交
1797 1798
    // maybe overwrite myself, no harm
    // just do it!
M
Minghao Li 已提交
1799 1800 1801
    pSyncNode->pMatchIndex->index[i] = SYNC_INDEX_INVALID;
  }

M
Minghao Li 已提交
1802 1803 1804
  // init peer mgr
  syncNodePeerStateInit(pSyncNode);

M
Minghao Li 已提交
1805
#if 0
1806 1807
  // update sender private term
  SSyncSnapshotSender* pMySender = syncNodeGetSnapshotSender(pSyncNode, &(pSyncNode->myRaftId));
1808
  if (pMySender != NULL) {
S
Shengliang Guan 已提交
1809
    for (int32_t i = 0; i < pSyncNode->pMatchIndex->replicaNum; ++i) {
1810 1811 1812
      if ((pSyncNode->senders)[i]->privateTerm > pMySender->privateTerm) {
        pMySender->privateTerm = (pSyncNode->senders)[i]->privateTerm;
      }
1813
    }
1814
    (pMySender->privateTerm) += 100;
1815
  }
M
Minghao Li 已提交
1816
#endif
1817

1818 1819 1820 1821 1822
  // close receiver
  if (snapshotReceiverIsStart(pSyncNode->pNewNodeReceiver)) {
    snapshotReceiverForceStop(pSyncNode->pNewNodeReceiver);
  }

M
Minghao Li 已提交
1823
  // stop elect timer
M
Minghao Li 已提交
1824
  syncNodeStopElectTimer(pSyncNode);
M
Minghao Li 已提交
1825

M
Minghao Li 已提交
1826 1827
  // start heartbeat timer
  syncNodeStartHeartbeatTimer(pSyncNode);
M
Minghao Li 已提交
1828

M
Minghao Li 已提交
1829 1830
  // send heartbeat right now
  syncNodeHeartbeatPeers(pSyncNode);
M
Minghao Li 已提交
1831

1832 1833 1834 1835 1836
  // call back
  if (pSyncNode->pFsm != NULL && pSyncNode->pFsm->FpBecomeLeaderCb != NULL) {
    pSyncNode->pFsm->FpBecomeLeaderCb(pSyncNode->pFsm);
  }

M
Minghao Li 已提交
1837 1838 1839
  // min match index
  pSyncNode->minMatchIndex = SYNC_INDEX_INVALID;

B
Benguang Zhao 已提交
1840 1841 1842
  // reset log buffer
  syncLogBufferReset(pSyncNode->pLogBuf, pSyncNode);

M
Minghao Li 已提交
1843
  // trace log
1844
  sNInfo(pSyncNode, "become leader %s", debugStr);
M
Minghao Li 已提交
1845 1846 1847
}

void syncNodeCandidate2Leader(SSyncNode* pSyncNode) {
M
Minghao Li 已提交
1848 1849
  ASSERT(pSyncNode->state == TAOS_SYNC_STATE_CANDIDATE);
  ASSERT(voteGrantedMajority(pSyncNode->pVotesGranted));
1850
  syncNodeBecomeLeader(pSyncNode, "candidate to leader");
M
Minghao Li 已提交
1851

S
Shengliang Guan 已提交
1852
  sNTrace(pSyncNode, "state change syncNodeCandidate2Leader");
M
Minghao Li 已提交
1853

B
Benguang Zhao 已提交
1854
  int32_t ret = syncNodeAppendNoop(pSyncNode);
1855 1856 1857 1858
  if (ret < 0) {
    sError("vgId:%d, failed to append noop entry since %s", pSyncNode->vgId, terrstr());
  }

B
Benguang Zhao 已提交
1859 1860 1861 1862
  SyncIndex lastIndex = pSyncNode->pLogStore->syncLogLastIndex(pSyncNode->pLogStore);
  ASSERT(lastIndex >= 0);
  sInfo("vgId:%d, become leader. term: %" PRId64 ", commit index: %" PRId64 ", last index: %" PRId64 "",
        pSyncNode->vgId, pSyncNode->pRaftStore->currentTerm, pSyncNode->commitIndex, lastIndex);
B
Benguang Zhao 已提交
1863 1864 1865 1866 1867 1868 1869
}

void syncNodeCandidate2LeaderOld(SSyncNode* pSyncNode) {
  ASSERT(pSyncNode->state == TAOS_SYNC_STATE_CANDIDATE);
  ASSERT(voteGrantedMajority(pSyncNode->pVotesGranted));
  syncNodeBecomeLeader(pSyncNode, "candidate to leader");

M
Minghao Li 已提交
1870
  // Raft 3.6.2 Committing entries from previous terms
1871 1872
  syncNodeAppendNoop(pSyncNode);
  syncMaybeAdvanceCommitIndex(pSyncNode);
1873 1874

  if (pSyncNode->replicaNum > 1) {
M
Minghao Li 已提交
1875
    syncNodeReplicate(pSyncNode);
1876
  }
M
Minghao Li 已提交
1877 1878
}

M
Minghao Li 已提交
1879 1880
bool syncNodeIsMnode(SSyncNode* pSyncNode) { return (pSyncNode->vgId == 1); }

M
Minghao Li 已提交
1881
int32_t syncNodePeerStateInit(SSyncNode* pSyncNode) {
S
Shengliang Guan 已提交
1882
  for (int32_t i = 0; i < TSDB_MAX_REPLICA; ++i) {
M
Minghao Li 已提交
1883 1884 1885 1886 1887
    pSyncNode->peerStates[i].lastSendIndex = SYNC_INDEX_INVALID;
    pSyncNode->peerStates[i].lastSendTime = 0;
  }

  return 0;
M
Minghao Li 已提交
1888 1889 1890
}

void syncNodeFollower2Candidate(SSyncNode* pSyncNode) {
M
Minghao Li 已提交
1891
  ASSERT(pSyncNode->state == TAOS_SYNC_STATE_FOLLOWER);
M
Minghao Li 已提交
1892
  pSyncNode->state = TAOS_SYNC_STATE_CANDIDATE;
B
Benguang Zhao 已提交
1893 1894 1895
  SyncIndex lastIndex = pSyncNode->pLogStore->syncLogLastIndex(pSyncNode->pLogStore);
  sInfo("vgId:%d, become candidate from follower. term: %" PRId64 ", commit index: %" PRId64 ", last index: %" PRId64,
        pSyncNode->vgId, pSyncNode->pRaftStore->currentTerm, pSyncNode->commitIndex, lastIndex);
M
Minghao Li 已提交
1896

S
Shengliang Guan 已提交
1897
  sNTrace(pSyncNode, "follower to candidate");
M
Minghao Li 已提交
1898 1899 1900
}

void syncNodeLeader2Follower(SSyncNode* pSyncNode) {
M
Minghao Li 已提交
1901
  ASSERT(pSyncNode->state == TAOS_SYNC_STATE_LEADER);
1902
  syncNodeBecomeFollower(pSyncNode, "leader to follower");
B
Benguang Zhao 已提交
1903 1904 1905 1906
  SyncIndex lastIndex = pSyncNode->pLogStore->syncLogLastIndex(pSyncNode->pLogStore);
  sInfo("vgId:%d, become follower from leader. term: %" PRId64 ", commit index: %" PRId64 ", last index: %" PRId64,
        pSyncNode->vgId, pSyncNode->pRaftStore->currentTerm, pSyncNode->commitIndex, lastIndex);

S
Shengliang Guan 已提交
1907
  sNTrace(pSyncNode, "leader to follower");
M
Minghao Li 已提交
1908 1909 1910
}

void syncNodeCandidate2Follower(SSyncNode* pSyncNode) {
M
Minghao Li 已提交
1911
  ASSERT(pSyncNode->state == TAOS_SYNC_STATE_CANDIDATE);
1912
  syncNodeBecomeFollower(pSyncNode, "candidate to follower");
B
Benguang Zhao 已提交
1913 1914 1915 1916
  SyncIndex lastIndex = pSyncNode->pLogStore->syncLogLastIndex(pSyncNode->pLogStore);
  sInfo("vgId:%d, become follower from candidate. term: %" PRId64 ", commit index: %" PRId64 ", last index: %" PRId64,
        pSyncNode->vgId, pSyncNode->pRaftStore->currentTerm, pSyncNode->commitIndex, lastIndex);

S
Shengliang Guan 已提交
1917
  sNTrace(pSyncNode, "candidate to follower");
M
Minghao Li 已提交
1918 1919
}

M
Minghao Li 已提交
1920 1921
// just called by syncNodeVoteForSelf
// need assert
M
Minghao Li 已提交
1922
void syncNodeVoteForTerm(SSyncNode* pSyncNode, SyncTerm term, SRaftId* pRaftId) {
M
Minghao Li 已提交
1923 1924
  ASSERT(term == pSyncNode->pRaftStore->currentTerm);
  ASSERT(!raftStoreHasVoted(pSyncNode->pRaftStore));
M
Minghao Li 已提交
1925 1926 1927 1928

  raftStoreVote(pSyncNode->pRaftStore, pRaftId);
}

M
Minghao Li 已提交
1929
// simulate get vote from outside
M
Minghao Li 已提交
1930
void syncNodeVoteForSelf(SSyncNode* pSyncNode) {
S
Shengliang Guan 已提交
1931
  syncNodeVoteForTerm(pSyncNode, pSyncNode->pRaftStore->currentTerm, &pSyncNode->myRaftId);
M
Minghao Li 已提交
1932

S
Shengliang Guan 已提交
1933 1934
  SRpcMsg rpcMsg = {0};
  int32_t ret = syncBuildRequestVoteReply(&rpcMsg, pSyncNode->vgId);
S
Shengliang Guan 已提交
1935
  if (ret != 0) return;
M
Minghao Li 已提交
1936

S
Shengliang Guan 已提交
1937
  SyncRequestVoteReply* pMsg = rpcMsg.pCont;
M
Minghao Li 已提交
1938 1939 1940 1941 1942 1943 1944
  pMsg->srcId = pSyncNode->myRaftId;
  pMsg->destId = pSyncNode->myRaftId;
  pMsg->term = pSyncNode->pRaftStore->currentTerm;
  pMsg->voteGranted = true;

  voteGrantedVote(pSyncNode->pVotesGranted, pMsg);
  votesRespondAdd(pSyncNode->pVotesRespond, pMsg);
S
Shengliang Guan 已提交
1945
  rpcFreeCont(rpcMsg.pCont);
M
Minghao Li 已提交
1946 1947
}

M
Minghao Li 已提交
1948
// return if has a snapshot
M
Minghao Li 已提交
1949 1950
bool syncNodeHasSnapshot(SSyncNode* pSyncNode) {
  bool      ret = false;
1951
  SSnapshot snapshot = {.data = NULL, .lastApplyIndex = -1, .lastApplyTerm = 0, .lastConfigIndex = -1};
1952 1953
  if (pSyncNode->pFsm->FpGetSnapshotInfo != NULL) {
    pSyncNode->pFsm->FpGetSnapshotInfo(pSyncNode->pFsm, &snapshot);
M
Minghao Li 已提交
1954 1955 1956 1957 1958 1959 1960
    if (snapshot.lastApplyIndex >= SYNC_INDEX_BEGIN) {
      ret = true;
    }
  }
  return ret;
}

M
Minghao Li 已提交
1961 1962
// return max(logLastIndex, snapshotLastIndex)
// if no snapshot and log, return -1
1963
SyncIndex syncNodeGetLastIndex(const SSyncNode* pSyncNode) {
M
Minghao Li 已提交
1964
  SSnapshot snapshot = {.data = NULL, .lastApplyIndex = -1, .lastApplyTerm = 0, .lastConfigIndex = -1};
1965 1966
  if (pSyncNode->pFsm->FpGetSnapshotInfo != NULL) {
    pSyncNode->pFsm->FpGetSnapshotInfo(pSyncNode->pFsm, &snapshot);
M
Minghao Li 已提交
1967 1968 1969 1970 1971 1972 1973
  }
  SyncIndex logLastIndex = pSyncNode->pLogStore->syncLogLastIndex(pSyncNode->pLogStore);

  SyncIndex lastIndex = logLastIndex > snapshot.lastApplyIndex ? logLastIndex : snapshot.lastApplyIndex;
  return lastIndex;
}

M
Minghao Li 已提交
1974 1975
// return the last term of snapshot and log
// if error, return SYNC_TERM_INVALID (by syncLogLastTerm)
M
Minghao Li 已提交
1976 1977
SyncTerm syncNodeGetLastTerm(SSyncNode* pSyncNode) {
  SyncTerm lastTerm = 0;
M
Minghao Li 已提交
1978 1979
  if (syncNodeHasSnapshot(pSyncNode)) {
    // has snapshot
1980
    SSnapshot snapshot = {.data = NULL, .lastApplyIndex = -1, .lastApplyTerm = 0, .lastConfigIndex = -1};
1981 1982
    if (pSyncNode->pFsm->FpGetSnapshotInfo != NULL) {
      pSyncNode->pFsm->FpGetSnapshotInfo(pSyncNode->pFsm, &snapshot);
M
Minghao Li 已提交
1983 1984
    }

M
Minghao Li 已提交
1985 1986 1987
    SyncIndex logLastIndex = pSyncNode->pLogStore->syncLogLastIndex(pSyncNode->pLogStore);
    if (logLastIndex > snapshot.lastApplyIndex) {
      lastTerm = pSyncNode->pLogStore->syncLogLastTerm(pSyncNode->pLogStore);
M
Minghao Li 已提交
1988 1989 1990 1991
    } else {
      lastTerm = snapshot.lastApplyTerm;
    }

M
Minghao Li 已提交
1992
  } else {
M
Minghao Li 已提交
1993 1994
    // no snapshot
    lastTerm = pSyncNode->pLogStore->syncLogLastTerm(pSyncNode->pLogStore);
1995
  }
M
Minghao Li 已提交
1996

M
Minghao Li 已提交
1997 1998 1999 2000 2001 2002 2003
  return lastTerm;
}

// get last index and term along with snapshot
int32_t syncNodeGetLastIndexTerm(SSyncNode* pSyncNode, SyncIndex* pLastIndex, SyncTerm* pLastTerm) {
  *pLastIndex = syncNodeGetLastIndex(pSyncNode);
  *pLastTerm = syncNodeGetLastTerm(pSyncNode);
2004 2005
  return 0;
}
M
Minghao Li 已提交
2006

M
Minghao Li 已提交
2007
// return append-entries first try index
M
Minghao Li 已提交
2008 2009 2010 2011 2012
SyncIndex syncNodeSyncStartIndex(SSyncNode* pSyncNode) {
  SyncIndex syncStartIndex = syncNodeGetLastIndex(pSyncNode) + 1;
  return syncStartIndex;
}

M
Minghao Li 已提交
2013 2014
// if index > 0, return index - 1
// else, return -1
2015 2016 2017 2018 2019 2020 2021 2022 2023
SyncIndex syncNodeGetPreIndex(SSyncNode* pSyncNode, SyncIndex index) {
  SyncIndex preIndex = index - 1;
  if (preIndex < SYNC_INDEX_INVALID) {
    preIndex = SYNC_INDEX_INVALID;
  }

  return preIndex;
}

M
Minghao Li 已提交
2024 2025 2026 2027
// if index < 0, return SYNC_TERM_INVALID
// if index == 0, return 0
// if index > 0, return preTerm
// if error, return SYNC_TERM_INVALID
2028 2029 2030 2031 2032 2033 2034 2035 2036
SyncTerm syncNodeGetPreTerm(SSyncNode* pSyncNode, SyncIndex index) {
  if (index < SYNC_INDEX_BEGIN) {
    return SYNC_TERM_INVALID;
  }

  if (index == SYNC_INDEX_BEGIN) {
    return 0;
  }

2037 2038 2039
  SyncTerm  preTerm = 0;
  SyncIndex preIndex = index - 1;

2040
  SSyncRaftEntry* pPreEntry = NULL;
2041 2042 2043 2044 2045 2046 2047
  SLRUCache*      pCache = pSyncNode->pLogStore->pCache;
  LRUHandle*      h = taosLRUCacheLookup(pCache, &preIndex, sizeof(preIndex));
  int32_t         code = 0;
  if (h) {
    pPreEntry = (SSyncRaftEntry*)taosLRUCacheValue(pCache, h);
    code = 0;

2048
    pSyncNode->pLogStore->cacheHit++;
2049 2050 2051
    sNTrace(pSyncNode, "hit cache index:%" PRId64 ", bytes:%u, %p", preIndex, pPreEntry->bytes, pPreEntry);

  } else {
2052
    pSyncNode->pLogStore->cacheMiss++;
2053 2054 2055 2056
    sNTrace(pSyncNode, "miss cache index:%" PRId64, preIndex);

    code = pSyncNode->pLogStore->syncLogGetEntry(pSyncNode->pLogStore, preIndex, &pPreEntry);
  }
M
Minghao Li 已提交
2057 2058 2059 2060 2061 2062

  SSnapshot snapshot = {.data = NULL,
                        .lastApplyIndex = SYNC_INDEX_INVALID,
                        .lastApplyTerm = SYNC_TERM_INVALID,
                        .lastConfigIndex = SYNC_INDEX_INVALID};

2063 2064 2065
  if (code == 0) {
    ASSERT(pPreEntry != NULL);
    preTerm = pPreEntry->term;
2066 2067 2068 2069

    if (h) {
      taosLRUCacheRelease(pCache, h, false);
    } else {
2070
      syncEntryDestroy(pPreEntry);
2071 2072
    }

2073 2074
    return preTerm;
  } else {
2075 2076 2077 2078
    if (pSyncNode->pFsm->FpGetSnapshotInfo != NULL) {
      pSyncNode->pFsm->FpGetSnapshotInfo(pSyncNode->pFsm, &snapshot);
      if (snapshot.lastApplyIndex == preIndex) {
        return snapshot.lastApplyTerm;
2079 2080 2081 2082
      }
    }
  }

2083
  sNError(pSyncNode, "sync node get pre term error, index:%" PRId64 ", snap-index:%" PRId64 ", snap-term:%" PRId64,
S
Shengliang Guan 已提交
2084
          index, snapshot.lastApplyIndex, snapshot.lastApplyTerm);
2085 2086
  return SYNC_TERM_INVALID;
}
M
Minghao Li 已提交
2087 2088 2089 2090

// get pre index and term of "index"
int32_t syncNodeGetPreIndexTerm(SSyncNode* pSyncNode, SyncIndex index, SyncIndex* pPreIndex, SyncTerm* pPreTerm) {
  *pPreIndex = syncNodeGetPreIndex(pSyncNode, index);
M
Minghao Li 已提交
2091
  *pPreTerm = syncNodeGetPreTerm(pSyncNode, index);
M
Minghao Li 已提交
2092 2093 2094
  return 0;
}

M
Minghao Li 已提交
2095
static void syncNodeEqPingTimer(void* param, void* tmrId) {
S
Shengliang Guan 已提交
2096
  if (!syncIsInit()) return;
M
Minghao Li 已提交
2097

S
Shengliang Guan 已提交
2098 2099 2100
  SSyncNode* pNode = param;
  if (atomic_load_64(&pNode->pingTimerLogicClockUser) <= atomic_load_64(&pNode->pingTimerLogicClock)) {
    SRpcMsg rpcMsg = {0};
S
Shengliang Guan 已提交
2101
    int32_t code = syncBuildTimeout(&rpcMsg, SYNC_TIMEOUT_PING, atomic_load_64(&pNode->pingTimerLogicClock),
S
Shengliang Guan 已提交
2102 2103
                                    pNode->pingTimerMS, pNode);
    if (code != 0) {
M
Minghao Li 已提交
2104
      sError("failed to build ping msg");
S
Shengliang Guan 已提交
2105 2106
      rpcFreeCont(rpcMsg.pCont);
      return;
M
Minghao Li 已提交
2107
    }
M
Minghao Li 已提交
2108

M
Minghao Li 已提交
2109
    // sTrace("enqueue ping msg");
S
Shengliang Guan 已提交
2110 2111
    code = pNode->syncEqMsg(pNode->msgcb, &rpcMsg);
    if (code != 0) {
M
Minghao Li 已提交
2112
      sError("failed to sync enqueue ping msg since %s", terrstr());
S
Shengliang Guan 已提交
2113 2114
      rpcFreeCont(rpcMsg.pCont);
      return;
2115
    }
M
Minghao Li 已提交
2116

S
Shengliang Guan 已提交
2117
    taosTmrReset(syncNodeEqPingTimer, pNode->pingTimerMS, pNode, syncEnv()->pTimerManager, &pNode->pPingTimer);
2118
  }
M
Minghao Li 已提交
2119 2120
}

M
Minghao Li 已提交
2121
static void syncNodeEqElectTimer(void* param, void* tmrId) {
S
Shengliang Guan 已提交
2122
  if (!syncIsInit()) return;
M
Minghao Li 已提交
2123

M
Minghao Li 已提交
2124 2125
  int64_t    rid = (int64_t)param;
  SSyncNode* pNode = syncNodeAcquire(rid);
M
Minghao Li 已提交
2126

2127
  if (pNode == NULL) return;
M
Minghao Li 已提交
2128 2129 2130 2131 2132

  if (pNode->syncEqMsg == NULL) {
    syncNodeRelease(pNode);
    return;
  }
2133

2134
  int64_t tsNow = taosGetTimestampMs();
M
Minghao Li 已提交
2135 2136 2137 2138
  if (tsNow < pNode->electTimerParam.executeTime) {
    syncNodeRelease(pNode);
    return;
  }
M
Minghao Li 已提交
2139

S
Shengliang Guan 已提交
2140
  SRpcMsg rpcMsg = {0};
2141 2142
  int32_t code =
      syncBuildTimeout(&rpcMsg, SYNC_TIMEOUT_ELECTION, pNode->electTimerParam.logicClock, pNode->electTimerMS, pNode);
S
Shengliang Guan 已提交
2143

S
Shengliang Guan 已提交
2144
  if (code != 0) {
M
Minghao Li 已提交
2145
    sError("failed to build elect msg");
M
Minghao Li 已提交
2146
    syncNodeRelease(pNode);
S
Shengliang Guan 已提交
2147
    return;
M
Minghao Li 已提交
2148 2149
  }

S
Shengliang Guan 已提交
2150
  SyncTimeout* pTimeout = rpcMsg.pCont;
S
Shengliang Guan 已提交
2151
  sNTrace(pNode, "enqueue elect msg lc:%" PRId64, pTimeout->logicClock);
S
Shengliang Guan 已提交
2152 2153 2154

  code = pNode->syncEqMsg(pNode->msgcb, &rpcMsg);
  if (code != 0) {
M
Minghao Li 已提交
2155
    sError("failed to sync enqueue elect msg since %s", terrstr());
S
Shengliang Guan 已提交
2156
    rpcFreeCont(rpcMsg.pCont);
M
Minghao Li 已提交
2157
    syncNodeRelease(pNode);
2158
    return;
M
Minghao Li 已提交
2159
  }
M
Minghao Li 已提交
2160 2161

  syncNodeRelease(pNode);
M
Minghao Li 已提交
2162 2163
}

M
Minghao Li 已提交
2164
static void syncNodeEqHeartbeatTimer(void* param, void* tmrId) {
S
Shengliang Guan 已提交
2165
  if (!syncIsInit()) return;
2166

S
Shengliang Guan 已提交
2167 2168 2169 2170
  SSyncNode* pNode = param;
  if (pNode->replicaNum > 1) {
    if (atomic_load_64(&pNode->heartbeatTimerLogicClockUser) <= atomic_load_64(&pNode->heartbeatTimerLogicClock)) {
      SRpcMsg rpcMsg = {0};
S
Shengliang Guan 已提交
2171
      int32_t code = syncBuildTimeout(&rpcMsg, SYNC_TIMEOUT_HEARTBEAT, atomic_load_64(&pNode->heartbeatTimerLogicClock),
S
Shengliang Guan 已提交
2172 2173 2174
                                      pNode->heartbeatTimerMS, pNode);

      if (code != 0) {
M
Minghao Li 已提交
2175
        sError("failed to build heartbeat msg");
S
Shengliang Guan 已提交
2176
        return;
2177
      }
M
Minghao Li 已提交
2178

2179
      sTrace("vgId:%d, enqueue heartbeat timer", pNode->vgId);
S
Shengliang Guan 已提交
2180 2181
      code = pNode->syncEqMsg(pNode->msgcb, &rpcMsg);
      if (code != 0) {
M
Minghao Li 已提交
2182
        sError("failed to enqueue heartbeat msg since %s", terrstr());
S
Shengliang Guan 已提交
2183 2184
        rpcFreeCont(rpcMsg.pCont);
        return;
2185
      }
S
Shengliang Guan 已提交
2186 2187 2188 2189

      taosTmrReset(syncNodeEqHeartbeatTimer, pNode->heartbeatTimerMS, pNode, syncEnv()->pTimerManager,
                   &pNode->pHeartbeatTimer);

2190
    } else {
S
Shengliang Guan 已提交
2191 2192
      sTrace("==syncNodeEqHeartbeatTimer== heartbeatTimerLogicClock:%" PRId64 ", heartbeatTimerLogicClockUser:%" PRId64,
             pNode->heartbeatTimerLogicClock, pNode->heartbeatTimerLogicClockUser);
2193
    }
M
Minghao Li 已提交
2194 2195 2196
  }
}

2197
static void syncNodeEqPeerHeartbeatTimer(void* param, void* tmrId) {
2198
  int64_t hbDataRid = (int64_t)param;
2199
  int64_t tsNow = taosGetTimestampMs();
2200

2201 2202
  SSyncHbTimerData* pData = syncHbTimerDataAcquire(hbDataRid);
  if (pData == NULL) {
M
Minghao Li 已提交
2203
    sError("hb timer get pData NULL, %" PRId64, hbDataRid);
2204 2205
    return;
  }
2206

2207
  SSyncNode* pSyncNode = syncNodeAcquire(pData->syncNodeRid);
M
Minghao Li 已提交
2208
  if (pSyncNode == NULL) {
2209
    syncHbTimerDataRelease(pData);
M
Minghao Li 已提交
2210
    sError("hb timer get pSyncNode NULL");
2211 2212 2213 2214 2215 2216 2217 2218
    return;
  }

  SSyncTimer* pSyncTimer = pData->pTimer;

  if (!pSyncNode->isStart) {
    syncNodeRelease(pSyncNode);
    syncHbTimerDataRelease(pData);
M
Minghao Li 已提交
2219
    sError("vgId:%d, hb timer sync node already stop", pSyncNode->vgId);
M
Minghao Li 已提交
2220 2221 2222
    return;
  }

M
Minghao Li 已提交
2223
  if (pSyncNode->state != TAOS_SYNC_STATE_LEADER) {
2224 2225
    syncNodeRelease(pSyncNode);
    syncHbTimerDataRelease(pData);
M
Minghao Li 已提交
2226
    sError("vgId:%d, hb timer sync node not leader", pSyncNode->vgId);
M
Minghao Li 已提交
2227 2228 2229
    return;
  }

M
Minghao Li 已提交
2230
  if (pSyncNode->pRaftStore == NULL) {
2231 2232
    syncNodeRelease(pSyncNode);
    syncHbTimerDataRelease(pData);
M
Minghao Li 已提交
2233
    sError("vgId:%d, hb timer raft store already stop", pSyncNode->vgId);
M
Minghao Li 已提交
2234 2235 2236
    return;
  }

M
Minghao Li 已提交
2237
  // sTrace("vgId:%d, eq peer hb timer", pSyncNode->vgId);
2238 2239

  if (pSyncNode->replicaNum > 1) {
M
Minghao Li 已提交
2240 2241 2242
    int64_t timerLogicClock = atomic_load_64(&pSyncTimer->logicClock);
    int64_t msgLogicClock = atomic_load_64(&pData->logicClock);

2243
    if (timerLogicClock == msgLogicClock) {
2244 2245 2246 2247 2248 2249 2250 2251 2252 2253 2254 2255 2256 2257 2258 2259 2260 2261 2262 2263
      if (tsNow > pData->execTime) {
#if 0        
        sTrace(
            "vgId:%d, hbDataRid:%ld,  EXECUTE this step-------- heartbeat tsNow:%ld, exec:%ld, tsNow-exec:%ld, "
            "---------",
            pSyncNode->vgId, hbDataRid, tsNow, pData->execTime, tsNow - pData->execTime);
#endif

        pData->execTime += pSyncTimer->timerMS;

        SRpcMsg rpcMsg = {0};
        (void)syncBuildHeartbeat(&rpcMsg, pSyncNode->vgId);

        SyncHeartbeat* pSyncMsg = rpcMsg.pCont;
        pSyncMsg->srcId = pSyncNode->myRaftId;
        pSyncMsg->destId = pData->destId;
        pSyncMsg->term = pSyncNode->pRaftStore->currentTerm;
        pSyncMsg->commitIndex = pSyncNode->commitIndex;
        pSyncMsg->minMatchIndex = syncMinMatchIndex(pSyncNode);
        pSyncMsg->privateTerm = 0;
2264
        pSyncMsg->timeStamp = tsNow;
2265 2266 2267 2268 2269 2270

        // update reset time
        int64_t timerElapsed = tsNow - pSyncTimer->timeStamp;
        pSyncTimer->timeStamp = tsNow;

        // send msg
2271 2272
        syncLogSendHeartbeat(pSyncNode, pSyncMsg, false, timerElapsed, pData->execTime);
        syncNodeSendHeartbeat(pSyncNode, &pSyncMsg->destId, &rpcMsg);
2273 2274 2275 2276 2277 2278 2279 2280
      } else {
#if 0        
        sTrace(
            "vgId:%d, hbDataRid:%ld,  pass this step-------- heartbeat tsNow:%ld, exec:%ld, tsNow-exec:%ld, ---------",
            pSyncNode->vgId, hbDataRid, tsNow, pData->execTime, tsNow - pData->execTime);
#endif
      }

M
Minghao Li 已提交
2281 2282
      if (syncIsInit()) {
        // sTrace("vgId:%d, reset peer hb timer", pSyncNode->vgId);
2283 2284
        taosTmrReset(syncNodeEqPeerHeartbeatTimer, pSyncTimer->timerMS / HEARTBEAT_TICK_NUM, (void*)hbDataRid,
                     syncEnv()->pTimerManager, &pSyncTimer->pTimer);
M
Minghao Li 已提交
2285 2286 2287 2288
      } else {
        sError("sync env is stop, reset peer hb timer error");
      }

2289
    } else {
M
Minghao Li 已提交
2290 2291
      sTrace("vgId:%d, do not send hb, timerLogicClock:%" PRId64 ", msgLogicClock:%" PRId64 "", pSyncNode->vgId,
             timerLogicClock, msgLogicClock);
2292 2293
    }
  }
2294 2295 2296

  syncHbTimerDataRelease(pData);
  syncNodeRelease(pSyncNode);
2297 2298
}

2299 2300 2301 2302 2303
static int32_t syncNodeEqNoop(SSyncNode* pNode) {
  if (pNode->state == TAOS_SYNC_STATE_LEADER) {
    terrno = TSDB_CODE_SYN_NOT_LEADER;
    return -1;
  }
M
Minghao Li 已提交
2304

2305 2306 2307 2308
  SyncIndex       index = pNode->pLogStore->syncLogWriteIndex(pNode->pLogStore);
  SyncTerm        term = pNode->pRaftStore->currentTerm;
  SSyncRaftEntry* pEntry = syncEntryBuildNoop(term, index, pNode->vgId);
  if (pEntry == NULL) return -1;
M
Minghao Li 已提交
2309

S
Shengliang Guan 已提交
2310
  SRpcMsg rpcMsg = {0};
S
Shengliang Guan 已提交
2311
  int32_t code = syncBuildClientRequestFromNoopEntry(&rpcMsg, pEntry, pNode->vgId);
2312
  syncEntryDestroy(pEntry);
M
Minghao Li 已提交
2313

2314 2315 2316
  sNTrace(pNode, "propose msg, type:noop");
  code = (*pNode->syncEqMsg)(pNode->msgcb, &rpcMsg);
  if (code != 0) {
M
Minghao Li 已提交
2317
    sError("failed to propose noop msg while enqueue since %s", terrstr());
2318
  }
M
Minghao Li 已提交
2319

2320
  return code;
M
Minghao Li 已提交
2321 2322
}

2323 2324
static void deleteCacheEntry(const void* key, size_t keyLen, void* value) { taosMemoryFree(value); }

2325 2326 2327 2328
int32_t syncCacheEntry(SSyncLogStore* pLogStore, SSyncRaftEntry* pEntry, LRUHandle** h) {
  SSyncLogStoreData* pData = pLogStore->data;
  sNTrace(pData->pSyncNode, "in cache index:%" PRId64 ", bytes:%u, %p", pEntry->index, pEntry->bytes, pEntry);

S
Shengliang Guan 已提交
2329 2330
  int32_t   code = 0;
  int32_t   entryLen = sizeof(*pEntry) + pEntry->dataLen;
2331 2332 2333 2334 2335 2336 2337 2338 2339
  LRUStatus status = taosLRUCacheInsert(pLogStore->pCache, &pEntry->index, sizeof(pEntry->index), pEntry, entryLen,
                                        deleteCacheEntry, h, TAOS_LRU_PRIORITY_LOW);
  if (status != TAOS_LRU_STATUS_OK) {
    code = -1;
  }

  return code;
}

B
Benguang Zhao 已提交
2340 2341 2342
int32_t syncNodeAppend(SSyncNode* ths, SSyncRaftEntry* pEntry) {
  // append to log buffer
  if (syncLogBufferAppend(ths->pLogBuf, ths, pEntry) < 0) {
2343
    sError("vgId:%d, failed to enqueue sync log buffer. index:%" PRId64 "", ths->vgId, pEntry->index);
B
Benguang Zhao 已提交
2344 2345 2346 2347
    return -1;
  }

  // proceed match index, with replicating on needed
2348
  SyncIndex matchIndex = syncLogBufferProceed(ths->pLogBuf, ths, NULL);
B
Benguang Zhao 已提交
2349

2350
  sTrace("vgId:%d, append raft entry. index: %" PRId64 ", term: %" PRId64 " pBuf: [%" PRId64 " %" PRId64 " %" PRId64
2351 2352 2353
         ", %" PRId64 ")",
         ths->vgId, pEntry->index, pEntry->term, ths->pLogBuf->startIndex, ths->pLogBuf->commitIndex,
         ths->pLogBuf->matchIndex, ths->pLogBuf->endIndex);
B
Benguang Zhao 已提交
2354

B
Benguang Zhao 已提交
2355 2356 2357 2358 2359 2360 2361 2362 2363 2364 2365 2366 2367 2368 2369 2370
  // multi replica
  if (ths->replicaNum > 1) {
    return 0;
  }

  // single replica
  (void)syncNodeUpdateCommitIndex(ths, matchIndex);

  if (syncLogBufferCommit(ths->pLogBuf, ths, ths->commitIndex) < 0) {
    sError("vgId:%d, failed to commit until commitIndex:%" PRId64 "", ths->vgId, ths->commitIndex);
    return -1;
  }

  return 0;
}

2371
bool syncNodeHeartbeatReplyTimeout(SSyncNode* pSyncNode) {
2372 2373 2374 2375 2376 2377 2378 2379 2380 2381 2382 2383 2384 2385 2386 2387 2388 2389 2390 2391 2392 2393
  if (pSyncNode->replicaNum == 1) {
    return false;
  }

  int32_t toCount = 0;
  int64_t tsNow = taosGetTimestampMs();
  for (int32_t i = 0; i < pSyncNode->peersNum; ++i) {
    int64_t recvTime = syncIndexMgrGetRecvTime(pSyncNode->pMatchIndex, &(pSyncNode->peersId[i]));
    if (recvTime == 0 || recvTime == -1) {
      continue;
    }

    if (tsNow - recvTime > SYNC_HEART_TIMEOUT_MS) {
      toCount++;
    }
  }

  bool b = (toCount >= pSyncNode->quorum ? true : false);

  return b;
}

2394 2395 2396 2397 2398 2399 2400 2401 2402 2403 2404 2405 2406 2407 2408 2409 2410 2411 2412
bool syncNodeSnapshotSending(SSyncNode* pSyncNode) {
  if (pSyncNode == NULL) return false;
  bool b = false;
  for (int32_t i = 0; i < pSyncNode->replicaNum; ++i) {
    if (pSyncNode->senders[i] != NULL && pSyncNode->senders[i]->start) {
      b = true;
      break;
    }
  }
  return b;
}

bool syncNodeSnapshotRecving(SSyncNode* pSyncNode) {
  if (pSyncNode == NULL) return false;
  if (pSyncNode->pNewNodeReceiver == NULL) return false;
  if (pSyncNode->pNewNodeReceiver->start) return true;
  return false;
}

M
Minghao Li 已提交
2413
static int32_t syncNodeAppendNoop(SSyncNode* ths) {
B
Benguang Zhao 已提交
2414 2415 2416 2417 2418 2419 2420 2421 2422
  SyncIndex index = syncLogBufferGetEndIndex(ths->pLogBuf);
  SyncTerm  term = ths->pRaftStore->currentTerm;

  SSyncRaftEntry* pEntry = syncEntryBuildNoop(term, index, ths->vgId);
  if (pEntry == NULL) {
    terrno = TSDB_CODE_OUT_OF_MEMORY;
    return -1;
  }

B
Benguang Zhao 已提交
2423 2424
  int32_t ret = syncNodeAppend(ths, pEntry);
  return 0;
B
Benguang Zhao 已提交
2425 2426 2427
}

static int32_t syncNodeAppendNoopOld(SSyncNode* ths) {
M
Minghao Li 已提交
2428 2429
  int32_t ret = 0;

2430
  SyncIndex       index = ths->pLogStore->syncLogWriteIndex(ths->pLogStore);
M
Minghao Li 已提交
2431
  SyncTerm        term = ths->pRaftStore->currentTerm;
M
Minghao Li 已提交
2432
  SSyncRaftEntry* pEntry = syncEntryBuildNoop(term, index, ths->vgId);
M
Minghao Li 已提交
2433
  ASSERT(pEntry != NULL);
M
Minghao Li 已提交
2434

2435 2436
  LRUHandle* h = NULL;

M
Minghao Li 已提交
2437
  if (ths->state == TAOS_SYNC_STATE_LEADER) {
2438
    int32_t code = ths->pLogStore->syncLogAppendEntry(ths->pLogStore, pEntry);
2439
    if (code != 0) {
M
Minghao Li 已提交
2440
      sError("append noop error");
2441 2442
      return -1;
    }
2443 2444

    syncCacheEntry(ths->pLogStore, pEntry, &h);
M
Minghao Li 已提交
2445 2446
  }

2447 2448 2449
  if (h) {
    taosLRUCacheRelease(ths->pLogStore->pCache, h, false);
  } else {
B
Benguang Zhao 已提交
2450
    syncEntryDestroy(pEntry);
2451 2452
  }

M
Minghao Li 已提交
2453 2454 2455
  return ret;
}

S
Shengliang Guan 已提交
2456 2457
int32_t syncNodeOnHeartbeat(SSyncNode* ths, const SRpcMsg* pRpcMsg) {
  SyncHeartbeat* pMsg = pRpcMsg->pCont;
2458

M
Minghao Li 已提交
2459 2460 2461 2462
  const STraceId* trace = &pRpcMsg->info.traceId;
  char            tbuf[40] = {0};
  TRACE_TO_STR(trace, tbuf);

2463
  int64_t tsMs = taosGetTimestampMs();
M
Minghao Li 已提交
2464
  int64_t timeDiff = tsMs - pMsg->timeStamp;
M
Minghao Li 已提交
2465
  syncLogRecvHeartbeat(ths, pMsg, timeDiff, tbuf);
2466

2467 2468 2469 2470
  SRpcMsg rpcMsg = {0};
  (void)syncBuildHeartbeatReply(&rpcMsg, ths->vgId);

  SyncHeartbeatReply* pMsgReply = rpcMsg.pCont;
2471 2472 2473 2474
  pMsgReply->destId = pMsg->srcId;
  pMsgReply->srcId = ths->myRaftId;
  pMsgReply->term = ths->pRaftStore->currentTerm;
  pMsgReply->privateTerm = 8864;  // magic number
2475
  pMsgReply->startTime = ths->startTime;
2476
  pMsgReply->timeStamp = tsMs;
2477

M
Minghao Li 已提交
2478
  if (pMsg->term == ths->pRaftStore->currentTerm && ths->state != TAOS_SYNC_STATE_LEADER) {
2479 2480
    syncIndexMgrSetRecvTime(ths->pNextIndex, &(pMsg->srcId), tsMs);

2481
    syncNodeResetElectTimer(ths);
M
Minghao Li 已提交
2482
    ths->minMatchIndex = pMsg->minMatchIndex;
2483 2484

    if (ths->state == TAOS_SYNC_STATE_FOLLOWER) {
2485
      // syncNodeFollowerCommit(ths, pMsg->commitIndex);
S
Shengliang Guan 已提交
2486 2487 2488 2489
      SRpcMsg rpcMsgLocalCmd = {0};
      (void)syncBuildLocalCmd(&rpcMsgLocalCmd, ths->vgId);

      SyncLocalCmd* pSyncMsg = rpcMsgLocalCmd.pCont;
2490 2491
      pSyncMsg->cmd = SYNC_LOCAL_CMD_FOLLOWER_CMT;
      pSyncMsg->fcIndex = pMsg->commitIndex;
2492
      SyncIndex fcIndex = pSyncMsg->fcIndex;
2493 2494 2495 2496 2497 2498 2499

      if (ths->syncEqMsg != NULL && ths->msgcb != NULL) {
        int32_t code = ths->syncEqMsg(ths->msgcb, &rpcMsgLocalCmd);
        if (code != 0) {
          sError("vgId:%d, sync enqueue fc-commit msg error, code:%d", ths->vgId, code);
          rpcFreeCont(rpcMsgLocalCmd.pCont);
        } else {
2500
          sTrace("vgId:%d, sync enqueue fc-commit msg, fc-index:%" PRId64, ths->vgId, fcIndex);
2501 2502
        }
      }
2503 2504 2505
    }
  }

M
Minghao Li 已提交
2506
  if (pMsg->term >= ths->pRaftStore->currentTerm && ths->state != TAOS_SYNC_STATE_FOLLOWER) {
2507
    // syncNodeStepDown(ths, pMsg->term);
S
Shengliang Guan 已提交
2508 2509 2510 2511
    SRpcMsg rpcMsgLocalCmd = {0};
    (void)syncBuildLocalCmd(&rpcMsgLocalCmd, ths->vgId);

    SyncLocalCmd* pSyncMsg = rpcMsgLocalCmd.pCont;
2512 2513 2514
    pSyncMsg->cmd = SYNC_LOCAL_CMD_STEP_DOWN;
    pSyncMsg->sdNewTerm = pMsg->term;

S
Shengliang Guan 已提交
2515 2516
    if (ths->syncEqMsg != NULL && ths->msgcb != NULL) {
      int32_t code = ths->syncEqMsg(ths->msgcb, &rpcMsgLocalCmd);
2517 2518 2519 2520
      if (code != 0) {
        sError("vgId:%d, sync enqueue step-down msg error, code:%d", ths->vgId, code);
        rpcFreeCont(rpcMsgLocalCmd.pCont);
      } else {
2521
        sTrace("vgId:%d, sync enqueue step-down msg, new-term: %" PRId64, ths->vgId, pSyncMsg->sdNewTerm);
2522
      }
2523 2524 2525 2526 2527 2528 2529 2530 2531 2532 2533 2534 2535 2536 2537
    }
  }

  /*
    // htonl
    SMsgHead* pHead = rpcMsg.pCont;
    pHead->contLen = htonl(pHead->contLen);
    pHead->vgId = htonl(pHead->vgId);
  */

  // reply
  syncNodeSendMsgById(&pMsgReply->destId, ths, &rpcMsg);
  return 0;
}

2538
int32_t syncNodeOnHeartbeatReply(SSyncNode* ths, const SRpcMsg* pRpcMsg) {
S
Shengliang Guan 已提交
2539 2540 2541 2542
  const STraceId* trace = &pRpcMsg->info.traceId;
  char            tbuf[40] = {0};
  TRACE_TO_STR(trace, tbuf);

2543
  SyncHeartbeatReply* pMsg = pRpcMsg->pCont;
B
Benguang Zhao 已提交
2544
  SSyncLogReplMgr*    pMgr = syncNodeGetLogReplMgr(ths, &pMsg->srcId);
2545 2546 2547 2548
  if (pMgr == NULL) {
    sError("vgId:%d, failed to get log repl mgr for the peer at addr 0x016%" PRIx64 "", ths->vgId, pMsg->srcId.addr);
    return -1;
  }
2549 2550

  int64_t tsMs = taosGetTimestampMs();
S
Shengliang Guan 已提交
2551
  syncLogRecvHeartbeatReply(ths, pMsg, tsMs - pMsg->timeStamp, tbuf);
2552

2553 2554
  syncIndexMgrSetRecvTime(ths->pMatchIndex, &pMsg->srcId, tsMs);

2555 2556 2557
  return syncLogReplMgrProcessHeartbeatReply(pMgr, ths, pMsg);
}

2558
int32_t syncNodeOnHeartbeatReplyOld(SSyncNode* ths, const SRpcMsg* pRpcMsg) {
2559
  SyncHeartbeatReply* pMsg = pRpcMsg->pCont;
2560

M
Minghao Li 已提交
2561 2562 2563 2564
  const STraceId* trace = &pRpcMsg->info.traceId;
  char            tbuf[40] = {0};
  TRACE_TO_STR(trace, tbuf);

M
Minghao Li 已提交
2565
  int64_t tsMs = taosGetTimestampMs();
M
Minghao Li 已提交
2566
  int64_t timeDiff = tsMs - pMsg->timeStamp;
M
Minghao Li 已提交
2567
  syncLogRecvHeartbeatReply(ths, pMsg, timeDiff, tbuf);
M
Minghao Li 已提交
2568

2569
  // update last reply time, make decision whether the other node is alive or not
M
Minghao Li 已提交
2570
  syncIndexMgrSetRecvTime(ths->pMatchIndex, &pMsg->srcId, tsMs);
2571 2572 2573
  return 0;
}

S
Shengliang Guan 已提交
2574 2575
int32_t syncNodeOnLocalCmd(SSyncNode* ths, const SRpcMsg* pRpcMsg) {
  SyncLocalCmd* pMsg = pRpcMsg->pCont;
2576 2577
  syncLogRecvLocalCmd(ths, pMsg, "");

2578 2579 2580 2581 2582 2583 2584 2585 2586 2587 2588 2589 2590 2591 2592 2593 2594 2595 2596 2597
  if (pMsg->cmd == SYNC_LOCAL_CMD_STEP_DOWN) {
    syncNodeStepDown(ths, pMsg->sdNewTerm);

  } else if (pMsg->cmd == SYNC_LOCAL_CMD_FOLLOWER_CMT) {
    (void)syncNodeUpdateCommitIndex(ths, pMsg->fcIndex);
    if (syncLogBufferCommit(ths->pLogBuf, ths, ths->commitIndex) < 0) {
      sError("vgId:%d, failed to commit raft log since %s. commit index: %" PRId64 "", ths->vgId, terrstr(),
             ths->commitIndex);
    }
  } else {
    sError("error local cmd");
  }

  return 0;
}

int32_t syncNodeOnLocalCmdOld(SSyncNode* ths, const SRpcMsg* pRpcMsg) {
  SyncLocalCmd* pMsg = pRpcMsg->pCont;
  syncLogRecvLocalCmd(ths, pMsg, "");

M
Minghao Li 已提交
2598 2599 2600
  if (pMsg->cmd == SYNC_LOCAL_CMD_STEP_DOWN) {
    syncNodeStepDown(ths, pMsg->sdNewTerm);

2601 2602 2603
  } else if (pMsg->cmd == SYNC_LOCAL_CMD_FOLLOWER_CMT) {
    syncNodeFollowerCommit(ths, pMsg->fcIndex);

M
Minghao Li 已提交
2604
  } else {
M
Minghao Li 已提交
2605
    sError("error local cmd");
M
Minghao Li 已提交
2606
  }
2607 2608 2609 2610

  return 0;
}

M
Minghao Li 已提交
2611 2612 2613 2614 2615 2616 2617 2618 2619 2620
// TLA+ Spec
// ClientRequest(i, v) ==
//     /\ state[i] = Leader
//     /\ LET entry == [term  |-> currentTerm[i],
//                      value |-> v]
//            newLog == Append(log[i], entry)
//        IN  log' = [log EXCEPT ![i] = newLog]
//     /\ UNCHANGED <<messages, serverVars, candidateVars,
//                    leaderVars, commitIndex>>
//
M
Minghao Li 已提交
2621

2622
int32_t syncNodeOnClientRequest(SSyncNode* ths, SRpcMsg* pMsg, SyncIndex* pRetIndex) {
S
Shengliang Guan 已提交
2623
  sNTrace(ths, "on client request");
2624

B
Benguang Zhao 已提交
2625 2626
  int32_t code = 0;

B
Benguang Zhao 已提交
2627 2628 2629
  SyncIndex       index = syncLogBufferGetEndIndex(ths->pLogBuf);
  SyncTerm        term = ths->pRaftStore->currentTerm;
  SSyncRaftEntry* pEntry = NULL;
2630 2631 2632 2633
  if (pMsg->msgType == TDMT_SYNC_CLIENT_REQUEST) {
    pEntry = syncEntryBuildFromClientRequest(pMsg->pCont, term, index);
  } else {
    pEntry = syncEntryBuildFromRpcMsg(pMsg, term, index);
B
Benguang Zhao 已提交
2634 2635 2636 2637 2638 2639 2640 2641 2642 2643
  }

  if (ths->state == TAOS_SYNC_STATE_LEADER) {
    if (pRetIndex) {
      (*pRetIndex) = index;
    }

    return syncNodeAppend(ths, pEntry);
  }

B
Benguang Zhao 已提交
2644
  return -1;
B
Benguang Zhao 已提交
2645 2646
}

2647 2648
int32_t syncNodeOnClientRequestOld(SSyncNode* ths, SRpcMsg* pMsg, SyncIndex* pRetIndex) {
  sNTrace(ths, "on client request");
B
Benguang Zhao 已提交
2649

M
Minghao Li 已提交
2650
  int32_t ret = 0;
2651
  int32_t code = 0;
M
Minghao Li 已提交
2652

M
Minghao Li 已提交
2653
  SyncIndex       index = ths->pLogStore->syncLogWriteIndex(ths->pLogStore);
M
Minghao Li 已提交
2654
  SyncTerm        term = ths->pRaftStore->currentTerm;
2655 2656 2657 2658 2659 2660 2661
  SSyncRaftEntry* pEntry;

  if (pMsg->msgType == TDMT_SYNC_CLIENT_REQUEST) {
    pEntry = syncEntryBuildFromClientRequest(pMsg->pCont, term, index);
  } else {
    pEntry = syncEntryBuildFromRpcMsg(pMsg, term, index);
  }
M
Minghao Li 已提交
2662

2663 2664
  LRUHandle* h = NULL;

M
Minghao Li 已提交
2665
  if (ths->state == TAOS_SYNC_STATE_LEADER) {
2666 2667 2668
    // append entry
    code = ths->pLogStore->syncLogAppendEntry(ths->pLogStore, pEntry);
    if (code != 0) {
2669 2670 2671 2672
      if (ths->replicaNum == 1) {
        if (h) {
          taosLRUCacheRelease(ths->pLogStore->pCache, h, false);
        } else {
2673
          syncEntryDestroy(pEntry);
2674
        }
2675

2676 2677 2678 2679
        return -1;

      } else {
        // del resp mgr, call FpCommitCb
2680 2681 2682 2683 2684 2685 2686 2687 2688 2689 2690
        SFsmCbMeta cbMeta = {
            .index = pEntry->index,
            .lastConfigIndex = SYNC_INDEX_INVALID,
            .isWeak = pEntry->isWeak,
            .code = -1,
            .state = ths->state,
            .seqNum = pEntry->seqNum,
            .term = pEntry->term,
            .currentTerm = ths->pRaftStore->currentTerm,
            .flag = 0,
        };
2691
        ths->pFsm->FpCommitCb(ths->pFsm, pMsg, &cbMeta);
2692 2693 2694 2695

        if (h) {
          taosLRUCacheRelease(ths->pLogStore->pCache, h, false);
        } else {
2696
          syncEntryDestroy(pEntry);
2697 2698
        }

2699 2700
        return -1;
      }
2701
    }
M
Minghao Li 已提交
2702

2703 2704
    syncCacheEntry(ths->pLogStore, pEntry, &h);

2705 2706
    // if mulit replica, start replicate right now
    if (ths->replicaNum > 1) {
M
Minghao Li 已提交
2707
      syncNodeReplicate(ths);
2708
    }
2709

2710 2711
    // if only myself, maybe commit right now
    if (ths->replicaNum == 1) {
2712 2713 2714 2715 2716
      if (syncNodeIsMnode(ths)) {
        syncMaybeAdvanceCommitIndex(ths);
      } else {
        syncOneReplicaAdvance(ths);
      }
2717
    }
M
Minghao Li 已提交
2718 2719
  }

2720 2721 2722 2723 2724 2725 2726 2727
  if (pRetIndex != NULL) {
    if (ret == 0 && pEntry != NULL) {
      *pRetIndex = pEntry->index;
    } else {
      *pRetIndex = SYNC_INDEX_INVALID;
    }
  }

2728 2729 2730
  if (h) {
    taosLRUCacheRelease(ths->pLogStore->pCache, h, false);
  } else {
B
Benguang Zhao 已提交
2731
    syncEntryDestroy(pEntry);
2732 2733
  }

M
Minghao Li 已提交
2734
  return ret;
2735
}
M
Minghao Li 已提交
2736

S
Shengliang Guan 已提交
2737 2738 2739
const char* syncStr(ESyncState state) {
  switch (state) {
    case TAOS_SYNC_STATE_FOLLOWER:
2740
      return "follower";
S
Shengliang Guan 已提交
2741
    case TAOS_SYNC_STATE_CANDIDATE:
2742
      return "candidate";
S
Shengliang Guan 已提交
2743
    case TAOS_SYNC_STATE_LEADER:
2744
      return "leader";
S
Shengliang Guan 已提交
2745
    case TAOS_SYNC_STATE_ERROR:
2746
      return "error";
S
Shengliang Guan 已提交
2747 2748 2749 2750
    case TAOS_SYNC_STATE_OFFLINE:
      return "offline";
    default:
      return "unknown";
S
Shengliang Guan 已提交
2751
  }
M
Minghao Li 已提交
2752
}
2753

2754
#if 0
2755
int32_t syncDoLeaderTransfer(SSyncNode* ths, SRpcMsg* pRpcMsg, SSyncRaftEntry* pEntry) {
2756
  if (ths->state != TAOS_SYNC_STATE_FOLLOWER) {
S
Shengliang Guan 已提交
2757
    sNTrace(ths, "I am not follower, can not do leader transfer");
2758 2759
    return 0;
  }
2760 2761

  if (!ths->restoreFinish) {
S
Shengliang Guan 已提交
2762
    sNTrace(ths, "restore not finish, can not do leader transfer");
2763 2764 2765
    return 0;
  }

2766
  if (pEntry->term < ths->pRaftStore->currentTerm) {
2767
    sNTrace(ths, "little term:%" PRId64 ", can not do leader transfer", pEntry->term);
2768 2769 2770 2771
    return 0;
  }

  if (pEntry->index < syncNodeGetLastIndex(ths)) {
S
Shengliang Guan 已提交
2772
    sNTrace(ths, "little index:%" PRId64 ", can not do leader transfer", pEntry->index);
2773 2774 2775
    return 0;
  }

2776 2777
  /*
    if (ths->vgId > 1) {
S
Shengliang Guan 已提交
2778
      sNTrace(ths, "I am vnode, can not do leader transfer");
2779 2780 2781 2782
      return 0;
    }
  */

2783
  SyncLeaderTransfer* pSyncLeaderTransfer = pRpcMsg->pCont;
S
Shengliang Guan 已提交
2784
  sNTrace(ths, "do leader transfer, index:%" PRId64, pEntry->index);
M
Minghao Li 已提交
2785

M
Minghao Li 已提交
2786 2787 2788
  bool sameId = syncUtilSameId(&(pSyncLeaderTransfer->newLeaderId), &(ths->myRaftId));
  bool sameNodeInfo = strcmp(pSyncLeaderTransfer->newNodeInfo.nodeFqdn, ths->myNodeInfo.nodeFqdn) == 0 &&
                      pSyncLeaderTransfer->newNodeInfo.nodePort == ths->myNodeInfo.nodePort;
M
Minghao Li 已提交
2789

M
Minghao Li 已提交
2790 2791
  bool same = sameId || sameNodeInfo;
  if (same) {
M
Minghao Li 已提交
2792 2793 2794 2795
    // reset elect timer now!
    int32_t electMS = 1;
    int32_t ret = syncNodeRestartElectTimer(ths, electMS);
    ASSERT(ret == 0);
M
Minghao Li 已提交
2796

2797
    sNTrace(ths, "maybe leader transfer to %s:%d %" PRId64, pSyncLeaderTransfer->newNodeInfo.nodeFqdn,
S
Shengliang Guan 已提交
2798
            pSyncLeaderTransfer->newNodeInfo.nodePort, pSyncLeaderTransfer->newLeaderId.addr);
2799 2800
  }

M
Minghao Li 已提交
2801
  if (ths->pFsm->FpLeaderTransferCb != NULL) {
S
Shengliang Guan 已提交
2802
    SFsmCbMeta cbMeta = {
S
Shengliang Guan 已提交
2803 2804 2805 2806 2807 2808 2809 2810 2811
        .code = 0,
        .currentTerm = ths->pRaftStore->currentTerm,
        .flag = 0,
        .index = pEntry->index,
        .lastConfigIndex = syncNodeGetSnapshotConfigIndex(ths, pEntry->index),
        .isWeak = pEntry->isWeak,
        .seqNum = pEntry->seqNum,
        .state = ths->state,
        .term = pEntry->term,
S
Shengliang Guan 已提交
2812 2813
    };
    ths->pFsm->FpLeaderTransferCb(ths->pFsm, pRpcMsg, &cbMeta);
2814 2815
  }

2816 2817 2818
  return 0;
}

2819 2820
#endif

2821
int32_t syncNodeUpdateNewConfigIndex(SSyncNode* ths, SSyncCfg* pNewCfg) {
S
Shengliang Guan 已提交
2822
  for (int32_t i = 0; i < pNewCfg->replicaNum; ++i) {
2823 2824 2825 2826 2827 2828 2829 2830 2831 2832 2833 2834 2835
    SRaftId raftId;
    raftId.addr = syncUtilAddr2U64((pNewCfg->nodeInfo)[i].nodeFqdn, (pNewCfg->nodeInfo)[i].nodePort);
    raftId.vgId = ths->vgId;

    if (syncUtilSameId(&(ths->myRaftId), &raftId)) {
      pNewCfg->myIndex = i;
      return 0;
    }
  }

  return -1;
}

2836 2837 2838 2839
bool syncNodeIsOptimizedOneReplica(SSyncNode* ths, SRpcMsg* pMsg) {
  return (ths->replicaNum == 1 && syncUtilUserCommit(pMsg->msgType) && ths->vgId != 1);
}

M
Minghao Li 已提交
2840
int32_t syncNodeDoCommit(SSyncNode* ths, SyncIndex beginIndex, SyncIndex endIndex, uint64_t flag) {
2841
  ASSERT(false);
2842 2843 2844 2845
  if (beginIndex > endIndex) {
    return 0;
  }

M
Minghao Li 已提交
2846 2847 2848 2849 2850 2851 2852 2853 2854
  if (ths == NULL) {
    return -1;
  }

  if (ths->pFsm != NULL && ths->pFsm->FpGetSnapshotInfo != NULL) {
    // advance commit index to sanpshot first
    SSnapshot snapshot = {0};
    ths->pFsm->FpGetSnapshotInfo(ths->pFsm, &snapshot);
    if (snapshot.lastApplyIndex >= 0 && snapshot.lastApplyIndex >= beginIndex) {
S
Shengliang Guan 已提交
2855
      sNTrace(ths, "commit by snapshot from index:%" PRId64 " to index:%" PRId64, beginIndex, snapshot.lastApplyIndex);
2856

M
Minghao Li 已提交
2857 2858 2859
      // update begin index
      beginIndex = snapshot.lastApplyIndex + 1;
    }
2860 2861
  }

2862 2863
  int32_t    code = 0;
  ESyncState state = flag;
M
Minghao Li 已提交
2864

S
Shengliang Guan 已提交
2865
  sNTrace(ths, "commit by wal from index:%" PRId64 " to index:%" PRId64, beginIndex, endIndex);
2866 2867 2868 2869 2870 2871

  // execute fsm
  if (ths->pFsm != NULL) {
    for (SyncIndex i = beginIndex; i <= endIndex; ++i) {
      if (i != SYNC_INDEX_INVALID) {
        SSyncRaftEntry* pEntry;
2872 2873 2874 2875
        SLRUCache*      pCache = ths->pLogStore->pCache;
        LRUHandle*      h = taosLRUCacheLookup(pCache, &i, sizeof(i));
        if (h) {
          pEntry = (SSyncRaftEntry*)taosLRUCacheValue(pCache, h);
2876

2877
          ths->pLogStore->cacheHit++;
2878 2879
          sNTrace(ths, "hit cache index:%" PRId64 ", bytes:%u, %p", i, pEntry->bytes, pEntry);

2880
        } else {
2881
          ths->pLogStore->cacheMiss++;
2882 2883
          sNTrace(ths, "miss cache index:%" PRId64, i);

2884
          code = ths->pLogStore->syncLogGetEntry(ths->pLogStore, i, &pEntry);
M
Minghao Li 已提交
2885 2886 2887
          // ASSERT(code == 0);
          // ASSERT(pEntry != NULL);
          if (code != 0 || pEntry == NULL) {
S
Shengliang Guan 已提交
2888
            sNError(ths, "get log entry error");
2889
            sFatal("vgId:%d, get log entry %" PRId64 " error when commit since %s", ths->vgId, i, terrstr());
M
Minghao Li 已提交
2890 2891
            continue;
          }
2892
        }
2893

2894
        SRpcMsg rpcMsg = {0};
2895 2896
        syncEntry2OriginalRpc(pEntry, &rpcMsg);

2897
        sTrace("do commit index:%" PRId64 ", type:%s", i, TMSG_INFO(pEntry->msgType));
M
Minghao Li 已提交
2898

2899
        // user commit
2900 2901
        if ((ths->pFsm->FpCommitCb != NULL) && syncUtilUserCommit(pEntry->originalRpcType)) {
          bool internalExecute = true;
S
Shengliang Guan 已提交
2902
          if ((ths->replicaNum == 1) && ths->restoreFinish && ths->vgId != 1) {
2903 2904 2905
            internalExecute = false;
          }

M
Minghao Li 已提交
2906 2907
          sNTrace(ths, "user commit index:%" PRId64 ", internal:%d, type:%s", i, internalExecute,
                  TMSG_INFO(pEntry->msgType));
2908

2909 2910
          // execute fsm in apply thread, or execute outside syncPropose
          if (internalExecute) {
S
Shengliang Guan 已提交
2911 2912 2913 2914 2915 2916 2917 2918 2919 2920 2921 2922
            SFsmCbMeta cbMeta = {
                .index = pEntry->index,
                .lastConfigIndex = syncNodeGetSnapshotConfigIndex(ths, pEntry->index),
                .isWeak = pEntry->isWeak,
                .code = 0,
                .state = ths->state,
                .seqNum = pEntry->seqNum,
                .term = pEntry->term,
                .currentTerm = ths->pRaftStore->currentTerm,
                .flag = flag,
            };

S
Shengliang Guan 已提交
2923
            syncRespMgrGetAndDel(ths->pSyncRespMgr, cbMeta.seqNum, &rpcMsg.info);
S
Shengliang Guan 已提交
2924
            ths->pFsm->FpCommitCb(ths->pFsm, &rpcMsg, &cbMeta);
M
Minghao Li 已提交
2925
          }
2926
        }
2927

2928 2929
#if 0
        // execute in pre-commit
M
Minghao Li 已提交
2930
        // leader transfer
2931 2932 2933
        if (pEntry->originalRpcType == TDMT_SYNC_LEADER_TRANSFER) {
          code = syncDoLeaderTransfer(ths, &rpcMsg, pEntry);
          ASSERT(code == 0);
2934
        }
2935
#endif
2936 2937

        // restore finish
2938
        // if only snapshot, a noop entry will be append, so syncLogLastIndex is always ok
2939 2940 2941 2942 2943 2944
        if (pEntry->index == ths->pLogStore->syncLogLastIndex(ths->pLogStore)) {
          if (ths->restoreFinish == false) {
            if (ths->pFsm->FpRestoreFinishCb != NULL) {
              ths->pFsm->FpRestoreFinishCb(ths->pFsm);
            }
            ths->restoreFinish = true;
M
Minghao Li 已提交
2945

2946
            int64_t restoreDelay = taosGetTimestampMs() - ths->leaderTime;
S
Shengliang Guan 已提交
2947
            sNTrace(ths, "restore finish, index:%" PRId64 ", elapsed:%" PRId64 " ms", pEntry->index, restoreDelay);
2948 2949 2950 2951
          }
        }

        rpcFreeCont(rpcMsg.pCont);
2952 2953 2954
        if (h) {
          taosLRUCacheRelease(pCache, h, false);
        } else {
B
Benguang Zhao 已提交
2955
          syncEntryDestroy(pEntry);
2956
        }
2957 2958 2959 2960
      }
    }
  }
  return 0;
2961 2962 2963
}

bool syncNodeInRaftGroup(SSyncNode* ths, SRaftId* pRaftId) {
S
Shengliang Guan 已提交
2964
  for (int32_t i = 0; i < ths->replicaNum; ++i) {
2965 2966 2967 2968 2969
    if (syncUtilSameId(&((ths->replicasId)[i]), pRaftId)) {
      return true;
    }
  }
  return false;
M
Minghao Li 已提交
2970 2971 2972 2973
}

SSyncSnapshotSender* syncNodeGetSnapshotSender(SSyncNode* ths, SRaftId* pDestId) {
  SSyncSnapshotSender* pSender = NULL;
S
Shengliang Guan 已提交
2974
  for (int32_t i = 0; i < ths->replicaNum; ++i) {
M
Minghao Li 已提交
2975 2976 2977 2978 2979
    if (syncUtilSameId(pDestId, &((ths->replicasId)[i]))) {
      pSender = (ths->senders)[i];
    }
  }
  return pSender;
M
Minghao Li 已提交
2980
}
M
Minghao Li 已提交
2981

2982 2983
SSyncTimer* syncNodeGetHbTimer(SSyncNode* ths, SRaftId* pDestId) {
  SSyncTimer* pTimer = NULL;
S
Shengliang Guan 已提交
2984
  for (int32_t i = 0; i < ths->replicaNum; ++i) {
2985 2986 2987 2988 2989 2990 2991
    if (syncUtilSameId(pDestId, &((ths->replicasId)[i]))) {
      pTimer = &((ths->peerHeartbeatTimerArr)[i]);
    }
  }
  return pTimer;
}

M
Minghao Li 已提交
2992 2993
SPeerState* syncNodeGetPeerState(SSyncNode* ths, const SRaftId* pDestId) {
  SPeerState* pState = NULL;
S
Shengliang Guan 已提交
2994
  for (int32_t i = 0; i < ths->replicaNum; ++i) {
M
Minghao Li 已提交
2995 2996 2997 2998 2999 3000 3001 3002 3003
    if (syncUtilSameId(pDestId, &((ths->replicasId)[i]))) {
      pState = &((ths->peerStates)[i]);
    }
  }
  return pState;
}

bool syncNodeNeedSendAppendEntries(SSyncNode* ths, const SRaftId* pDestId, const SyncAppendEntries* pMsg) {
  SPeerState* pState = syncNodeGetPeerState(ths, pDestId);
M
Minghao Li 已提交
3004
  if (pState == NULL) {
3005
    sError("vgId:%d, replica maybe dropped", ths->vgId);
M
Minghao Li 已提交
3006 3007
    return false;
  }
M
Minghao Li 已提交
3008 3009 3010 3011 3012 3013 3014 3015 3016 3017 3018

  SyncIndex sendIndex = pMsg->prevLogIndex + 1;
  int64_t   tsNow = taosGetTimestampMs();

  if (pState->lastSendIndex == sendIndex && tsNow - pState->lastSendTime < SYNC_APPEND_ENTRIES_TIMEOUT_MS) {
    return false;
  }

  return true;
}

M
Minghao Li 已提交
3019 3020 3021 3022 3023 3024 3025 3026 3027 3028 3029 3030 3031 3032
bool syncNodeCanChange(SSyncNode* pSyncNode) {
  if (pSyncNode->changing) {
    sError("sync cannot change");
    return false;
  }

  if ((pSyncNode->commitIndex >= SYNC_INDEX_BEGIN)) {
    SyncIndex lastIndex = syncNodeGetLastIndex(pSyncNode);
    if (pSyncNode->commitIndex != lastIndex) {
      sError("sync cannot change2");
      return false;
    }
  }

S
Shengliang Guan 已提交
3033
  for (int32_t i = 0; i < pSyncNode->peersNum; ++i) {
M
Minghao Li 已提交
3034
    SSyncSnapshotSender* pSender = syncNodeGetSnapshotSender(pSyncNode, &(pSyncNode->peersId)[i]);
M
Minghao Li 已提交
3035
    if (pSender != NULL && pSender->start) {
M
Minghao Li 已提交
3036 3037 3038 3039 3040 3041
      sError("sync cannot change3");
      return false;
    }
  }

  return true;
M
Minghao Li 已提交
3042
}