syncMain.c 83.0 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"
M
Minghao Li 已提交
25
#include "syncMessage.h"
26
#include "syncPipeline.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"
38
#include "tglobal.h"
M
Minghao Li 已提交
39
#include "tref.h"
M
Minghao Li 已提交
40

M
Minghao Li 已提交
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 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) {
S
Shengliang Guan 已提交
87
    sError("failed to acquire rid:%" PRId64 " of tsNodeReftId for pSyncNode", rid);
B
Benguang Zhao 已提交
88 89 90 91
    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
  }
}

126 127 128 129 130 131 132 133
void syncPostStop(int64_t rid) {
  SSyncNode* pSyncNode = syncNodeAcquire(rid);
  if (pSyncNode != NULL) {
    syncNodePostClose(pSyncNode);
    syncNodeRelease(pSyncNode);
  }
}

S
Shengliang Guan 已提交
134 135 136
static bool syncNodeCheckNewConfig(SSyncNode* pSyncNode, const SSyncCfg* pCfg) {
  if (!syncNodeInConfig(pSyncNode, pCfg)) return false;
  return abs(pCfg->replicaNum - pSyncNode->replicaNum) <= 1;
M
Minghao Li 已提交
137 138
}

S
Shengliang Guan 已提交
139
int32_t syncReconfig(int64_t rid, SSyncCfg* pNewCfg) {
S
Shengliang Guan 已提交
140
  SSyncNode* pSyncNode = syncNodeAcquire(rid);
S
Shengliang Guan 已提交
141
  if (pSyncNode == NULL) return -1;
M
Minghao Li 已提交
142

M
Minghao Li 已提交
143
  if (!syncNodeCheckNewConfig(pSyncNode, pNewCfg)) {
S
Shengliang Guan 已提交
144
    syncNodeRelease(pSyncNode);
M
Minghao Li 已提交
145
    terrno = TSDB_CODE_SYN_NEW_CONFIG_ERROR;
S
Shengliang Guan 已提交
146
    sError("vgId:%d, failed to reconfig since invalid new config", pSyncNode->vgId);
M
Minghao Li 已提交
147
    return -1;
M
Minghao Li 已提交
148
  }
149

S
Shengliang Guan 已提交
150 151
  syncNodeUpdateNewConfigIndex(pSyncNode, pNewCfg);
  syncNodeDoConfigChange(pSyncNode, pNewCfg, SYNC_INDEX_INVALID);
S
Shengliang Guan 已提交
152

M
Minghao Li 已提交
153 154 155 156
  if (pSyncNode->state == TAOS_SYNC_STATE_LEADER) {
    syncNodeStopHeartbeatTimer(pSyncNode);

    for (int32_t i = 0; i < TSDB_MAX_REPLICA; ++i) {
S
Shengliang Guan 已提交
157
      syncHbTimerInit(pSyncNode, &pSyncNode->peerHeartbeatTimerArr[i], pSyncNode->replicasId[i]);
M
Minghao Li 已提交
158 159 160
    }

    syncNodeStartHeartbeatTimer(pSyncNode);
S
Shengliang Guan 已提交
161
    // syncNodeReplicate(pSyncNode);
M
Minghao Li 已提交
162
  }
S
Shengliang Guan 已提交
163

S
Shengliang Guan 已提交
164
  syncNodeRelease(pSyncNode);
S
Shengliang Guan 已提交
165
  return 0;
M
Minghao Li 已提交
166
}
M
Minghao Li 已提交
167

S
Shengliang Guan 已提交
168 169 170 171
int32_t syncProcessMsg(int64_t rid, SRpcMsg* pMsg) {
  int32_t code = -1;
  if (!syncIsInit()) return code;

S
Shengliang Guan 已提交
172
  SSyncNode* pSyncNode = syncNodeAcquire(rid);
S
Shengliang Guan 已提交
173 174
  if (pSyncNode == NULL) return code;

S
Shengliang Guan 已提交
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
  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:
S
Shengliang Guan 已提交
204
      code = syncNodeOnSnapshotRsp(pSyncNode, pMsg);
S
Shengliang Guan 已提交
205 206 207 208 209
      break;
    case TDMT_SYNC_LOCAL_CMD:
      code = syncNodeOnLocalCmd(pSyncNode, pMsg);
      break;
    default:
210
      terrno = TSDB_CODE_MSG_NOT_PROCESSED;
S
Shengliang Guan 已提交
211
      code = -1;
M
Minghao Li 已提交
212 213
  }

S
Shengliang Guan 已提交
214
  syncNodeRelease(pSyncNode);
215 216 217 218
  if (code != 0) {
    sDebug("vgId:%d, failed to process sync msg:%p type:%s since 0x%x", pSyncNode->vgId, pMsg, TMSG_INFO(pMsg->msgType),
           terrno);
  }
S
Shengliang Guan 已提交
219
  return code;
220 221
}

S
Shengliang Guan 已提交
222
int32_t syncLeaderTransfer(int64_t rid) {
S
Shengliang Guan 已提交
223
  SSyncNode* pSyncNode = syncNodeAcquire(rid);
S
Shengliang Guan 已提交
224
  if (pSyncNode == NULL) return -1;
225

S
Shengliang Guan 已提交
226
  int32_t ret = syncNodeLeaderTransfer(pSyncNode);
S
Shengliang Guan 已提交
227
  syncNodeRelease(pSyncNode);
228 229 230
  return ret;
}

231
int32_t syncSendTimeoutRsp(int64_t rid, int64_t seq) {
S
Shengliang Guan 已提交
232
  SSyncNode* pNode = syncNodeAcquire(rid);
233
  if (pNode == NULL) return -1;
S
Shengliang Guan 已提交
234 235

  SRpcMsg rpcMsg = {0};
236
  int32_t ret = syncRespMgrGetAndDel(pNode->pSyncRespMgr, seq, &rpcMsg.info);
S
Shengliang Guan 已提交
237 238 239
  rpcMsg.code = TSDB_CODE_SYN_TIMEOUT;

  syncNodeRelease(pNode);
240
  if (ret == 1) {
241
    sInfo("send timeout response, seq:%" PRId64 " handle:%p ahandle:%p", seq, rpcMsg.info.handle, rpcMsg.info.ahandle);
242
    rpcSendResponse(&rpcMsg);
243 244
    return 0;
  } else {
245
    sError("no message handle to send timeout response, seq:%" PRId64, seq);
246
    return -1;
247
  }
S
Shengliang Guan 已提交
248 249
}

M
Minghao Li 已提交
250 251 252 253 254 255 256 257 258 259 260 261 262 263 264 265
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;
}

266
int32_t syncBeginSnapshot(int64_t rid, int64_t lastApplyIndex) {
S
Shengliang Guan 已提交
267
  SSyncNode* pSyncNode = syncNodeAcquire(rid);
268
  if (pSyncNode == NULL) {
269
    sError("sync begin snapshot error");
270 271
    return -1;
  }
272

273 274 275 276 277 278 279 280 281 282
  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;
  }

283
  int32_t code = 0;
284
  int64_t logRetention = 0;
285

M
Minghao Li 已提交
286
  if (syncNodeIsMnode(pSyncNode)) {
M
Minghao Li 已提交
287
    // mnode
288
    logRetention = tsMndLogRetention;
M
Minghao Li 已提交
289 290 291 292
  } else {
    // vnode
    if (pSyncNode->replicaNum > 1) {
      // multi replicas
293 294 295
      logRetention = SYNC_VNODE_LOG_RETENTION;
    }
  }
M
Minghao Li 已提交
296

297 298 299 300 301 302
  if (pSyncNode->replicaNum > 1) {
    if (pSyncNode->state != TAOS_SYNC_STATE_LEADER && pSyncNode->state != TAOS_SYNC_STATE_FOLLOWER) {
      sNTrace(pSyncNode, "new-snapshot-index:%" PRId64 " candidate or unknown state, do not delete wal",
              lastApplyIndex);
      syncNodeRelease(pSyncNode);
      return 0;
303
    }
304
    logRetention = TMAX(logRetention, lastApplyIndex - pSyncNode->minMatchIndex + logRetention);
305 306
  }

M
Minghao Li 已提交
307
_DEL_WAL:
308

M
Minghao Li 已提交
309
  do {
310 311 312 313 314 315 316 317 318 319 320
    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();

321
        code = walBeginSnapshot(pData->pWal, lastApplyIndex, logRetention);
322 323 324 325 326 327 328 329
        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);
        }
330

M
Minghao Li 已提交
331
      } else {
332 333
        sNTrace(pSyncNode, "snapshotting for %" PRId64 ", do not delete wal for new-snapshot-index:%" PRId64,
                snapshottingIndex, lastApplyIndex);
M
Minghao Li 已提交
334
      }
335
    }
M
Minghao Li 已提交
336
  } while (0);
337

S
Shengliang Guan 已提交
338
  syncNodeRelease(pSyncNode);
339 340 341 342
  return code;
}

int32_t syncEndSnapshot(int64_t rid) {
S
Shengliang Guan 已提交
343
  SSyncNode* pSyncNode = syncNodeAcquire(rid);
344
  if (pSyncNode == NULL) {
345
    sError("sync end snapshot error");
346 347 348
    return -1;
  }

349 350 351 352
  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 已提交
353
    if (code != 0) {
354
      sNError(pSyncNode, "wal snapshot end error since:%s", terrstr());
S
Shengliang Guan 已提交
355
      syncNodeRelease(pSyncNode);
M
Minghao Li 已提交
356 357
      return -1;
    } else {
S
Shengliang Guan 已提交
358
      sNTrace(pSyncNode, "wal snapshot end, index:%" PRId64, atomic_load_64(&pSyncNode->snapshottingIndex));
M
Minghao Li 已提交
359 360
      atomic_store_64(&pSyncNode->snapshottingIndex, SYNC_INDEX_INVALID);
    }
361
  }
362

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

M
Minghao Li 已提交
367
int32_t syncStepDown(int64_t rid, SyncTerm newTerm) {
S
Shengliang Guan 已提交
368
  SSyncNode* pSyncNode = syncNodeAcquire(rid);
M
Minghao Li 已提交
369
  if (pSyncNode == NULL) {
370
    sError("sync step down error");
M
Minghao Li 已提交
371 372 373
    return -1;
  }

M
Minghao Li 已提交
374
  syncNodeStepDown(pSyncNode, newTerm);
S
Shengliang Guan 已提交
375
  syncNodeRelease(pSyncNode);
M
Minghao Li 已提交
376
  return 0;
M
Minghao Li 已提交
377 378
}

379
bool syncNodeIsReadyForRead(SSyncNode* pSyncNode) {
380
  if (pSyncNode == NULL) {
381
    terrno = TSDB_CODE_SYN_INTERNAL_ERROR;
382
    sError("sync ready for read error");
383 384
    return false;
  }
M
Minghao Li 已提交
385

386 387 388 389 390
  if (pSyncNode->state != TAOS_SYNC_STATE_LEADER) {
    terrno = TSDB_CODE_SYN_NOT_LEADER;
    return false;
  }

391
  if (!pSyncNode->restoreFinish) {
392
    terrno = TSDB_CODE_SYN_RESTORING;
393
    return false;
394
  }
395

396
  return true;
397 398 399 400 401 402 403 404 405 406 407
}

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

  bool ready = syncNodeIsReadyForRead(pSyncNode);

408 409
  syncNodeRelease(pSyncNode);
  return ready;
M
Minghao Li 已提交
410
}
M
Minghao Li 已提交
411

412 413 414 415 416 417 418 419 420 421 422 423 424 425 426 427 428 429 430 431 432 433
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 已提交
434 435
int32_t syncNodeLeaderTransfer(SSyncNode* pSyncNode) {
  if (pSyncNode->peersNum == 0) {
S
Shengliang Guan 已提交
436
    sDebug("vgId:%d, only one replica, cannot leader transfer", pSyncNode->vgId);
M
Minghao Li 已提交
437 438
    terrno = TSDB_CODE_SYN_ONE_REPLICA;
    return -1;
M
Minghao Li 已提交
439
  }
M
Minghao Li 已提交
440

441
  int32_t ret = 0;
442
  if (pSyncNode->state == TAOS_SYNC_STATE_LEADER && pSyncNode->replicaNum > 1) {
443
    SNodeInfo newLeader = (pSyncNode->peersNodeInfo)[0];
444 445 446 447 448 449 450
    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];
      }
    }
451 452 453
    ret = syncNodeLeaderTransferTo(pSyncNode, newLeader);
  }

M
Minghao Li 已提交
454
  return ret;
M
Minghao Li 已提交
455 456
}

M
Minghao Li 已提交
457 458
int32_t syncNodeLeaderTransferTo(SSyncNode* pSyncNode, SNodeInfo newLeader) {
  if (pSyncNode->replicaNum == 1) {
S
Shengliang Guan 已提交
459
    sDebug("vgId:%d, only one replica, cannot leader transfer", pSyncNode->vgId);
M
Minghao Li 已提交
460 461
    terrno = TSDB_CODE_SYN_ONE_REPLICA;
    return -1;
M
Minghao Li 已提交
462
  }
463

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

466 467 468 469
  SRpcMsg rpcMsg = {0};
  (void)syncBuildLeaderTransfer(&rpcMsg, pSyncNode->vgId);

  SyncLeaderTransfer* pMsg = rpcMsg.pCont;
470
  pMsg->newLeaderId.addr = SYNC_ADDR(&newLeader);
M
Minghao Li 已提交
471 472 473
  pMsg->newLeaderId.vgId = pSyncNode->vgId;
  pMsg->newNodeInfo = newLeader;

S
Shengliang Guan 已提交
474
  int32_t ret = syncNodePropose(pSyncNode, &rpcMsg, false, NULL);
S
Shengliang Guan 已提交
475 476
  rpcFreeCont(rpcMsg.pCont);
  return ret;
M
Minghao Li 已提交
477 478
}

479 480
SSyncState syncGetState(int64_t rid) {
  SSyncState state = {.state = TAOS_SYNC_STATE_ERROR};
M
Minghao Li 已提交
481

S
Shengliang Guan 已提交
482
  SSyncNode* pSyncNode = syncNodeAcquire(rid);
483 484 485
  if (pSyncNode != NULL) {
    state.state = pSyncNode->state;
    state.restored = pSyncNode->restoreFinish;
486 487 488 489 490
    if (pSyncNode->vgId != 1) {
      state.canRead = syncNodeIsReadyForRead(pSyncNode);
    } else {
      state.canRead = state.restored;
    }
491
    syncNodeRelease(pSyncNode);
M
Minghao Li 已提交
492 493
  }

494
  return state;
M
Minghao Li 已提交
495 496
}

497
SyncIndex syncNodeGetSnapshotConfigIndex(SSyncNode* pSyncNode, SyncIndex snapshotLastApplyIndex) {
498 499
  ASSERT(pSyncNode->raftCfg.configIndexCount >= 1);
  SyncIndex lastIndex = (pSyncNode->raftCfg.configIndexArr)[0];
500

501 502 503 504
  for (int32_t i = 0; i < pSyncNode->raftCfg.configIndexCount; ++i) {
    if ((pSyncNode->raftCfg.configIndexArr)[i] > lastIndex &&
        (pSyncNode->raftCfg.configIndexArr)[i] <= snapshotLastApplyIndex) {
      lastIndex = (pSyncNode->raftCfg.configIndexArr)[i];
505 506
    }
  }
S
Shengliang Guan 已提交
507
  sTrace("vgId:%d, sync get last config index, index:%" PRId64 " lcindex:%" PRId64, pSyncNode->vgId,
S
Shengliang Guan 已提交
508
         snapshotLastApplyIndex, lastIndex);
509 510 511 512

  return lastIndex;
}

513 514
void syncGetRetryEpSet(int64_t rid, SEpSet* pEpSet) {
  pEpSet->numOfEps = 0;
M
Minghao Li 已提交
515

S
Shengliang Guan 已提交
516
  SSyncNode* pSyncNode = syncNodeAcquire(rid);
S
Shengliang Guan 已提交
517
  if (pSyncNode == NULL) return;
M
Minghao Li 已提交
518

519
  for (int32_t i = 0; i < pSyncNode->raftCfg.cfg.replicaNum; ++i) {
S
Shengliang Guan 已提交
520
    SEp* pEp = &pEpSet->eps[i];
521 522
    tstrncpy(pEp->fqdn, pSyncNode->raftCfg.cfg.nodeInfo[i].nodeFqdn, TSDB_FQDN_LEN);
    pEp->port = (pSyncNode->raftCfg.cfg.nodeInfo)[i].nodePort;
S
Shengliang Guan 已提交
523
    pEpSet->numOfEps++;
524
    sDebug("vgId:%d, sync get retry epset, index:%d %s:%d", pSyncNode->vgId, i, pEp->fqdn, pEp->port);
M
Minghao Li 已提交
525
  }
M
Minghao Li 已提交
526
  if (pEpSet->numOfEps > 0) {
527
    pEpSet->inUse = (pSyncNode->raftCfg.cfg.myIndex + 1) % pEpSet->numOfEps;
M
Minghao Li 已提交
528 529
  }

S
Shengliang Guan 已提交
530
  sInfo("vgId:%d, sync get retry epset numOfEps:%d inUse:%d", pSyncNode->vgId, pEpSet->numOfEps, pEpSet->inUse);
S
Shengliang Guan 已提交
531
  syncNodeRelease(pSyncNode);
M
Minghao Li 已提交
532 533
}

S
Shengliang Guan 已提交
534
int32_t syncPropose(int64_t rid, SRpcMsg* pMsg, bool isWeak, int64_t* seq) {
S
Shengliang Guan 已提交
535
  SSyncNode* pSyncNode = syncNodeAcquire(rid);
536
  if (pSyncNode == NULL) {
537
    sError("sync propose error");
M
Minghao Li 已提交
538
    return -1;
539
  }
540

S
Shengliang Guan 已提交
541
  int32_t ret = syncNodePropose(pSyncNode, pMsg, isWeak, seq);
S
Shengliang Guan 已提交
542
  syncNodeRelease(pSyncNode);
M
Minghao Li 已提交
543 544
  return ret;
}
M
Minghao Li 已提交
545

S
Shengliang Guan 已提交
546
int32_t syncNodePropose(SSyncNode* pSyncNode, SRpcMsg* pMsg, bool isWeak, int64_t* seq) {
S
Shengliang Guan 已提交
547 548
  if (pSyncNode->state != TAOS_SYNC_STATE_LEADER) {
    terrno = TSDB_CODE_SYN_NOT_LEADER;
S
Shengliang Guan 已提交
549
    sNError(pSyncNode, "sync propose not leader, type:%s", TMSG_INFO(pMsg->msgType));
S
Shengliang Guan 已提交
550 551
    return -1;
  }
552

S
Shengliang Guan 已提交
553 554 555 556 557 558 559
  // 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;
  }
560

561
  // heartbeat timeout
562
  if (syncNodeHeartbeatReplyTimeout(pSyncNode)) {
563 564 565 566 567 568
    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 已提交
569 570 571
  // optimized one replica
  if (syncNodeIsOptimizedOneReplica(pSyncNode, pMsg)) {
    SyncIndex retIndex;
572
    int32_t   code = syncNodeOnClientRequest(pSyncNode, pMsg, &retIndex);
S
Shengliang Guan 已提交
573 574
    if (code == 0) {
      pMsg->info.conn.applyIndex = retIndex;
575
      pMsg->info.conn.applyTerm = raftStoreGetTerm(pSyncNode);
576 577 578
      sTrace("vgId:%d, propose optimized msg, index:%" PRId64 " type:%s", pSyncNode->vgId, retIndex,
             TMSG_INFO(pMsg->msgType));
      return 1;
M
Minghao Li 已提交
579
    } else {
S
Shengliang Guan 已提交
580
      terrno = TSDB_CODE_SYN_INTERNAL_ERROR;
581
      sError("vgId:%d, failed to propose optimized msg, index:%" PRId64 " type:%s", pSyncNode->vgId, retIndex,
S
Shengliang Guan 已提交
582
             TMSG_INFO(pMsg->msgType));
583
      return -1;
584
    }
S
Shengliang Guan 已提交
585
  } else {
S
Shengliang Guan 已提交
586 587
    SRespStub stub = {.createTime = taosGetTimestampMs(), .rpcMsg = *pMsg};
    uint64_t  seqNum = syncRespMgrAdd(pSyncNode->pSyncRespMgr, &stub);
588
    SRpcMsg   rpcMsg = {0};
S
Shengliang Guan 已提交
589
    int32_t   code = syncBuildClientRequest(&rpcMsg, pMsg, seqNum, isWeak, pSyncNode->vgId);
590 591 592 593
    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 已提交
594
    }
595

596 597 598 599 600
    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 已提交
601
    }
M
Minghao Li 已提交
602

S
Shengliang Guan 已提交
603
    if (seq != NULL) *seq = seqNum;
604
    return code;
M
Minghao Li 已提交
605
  }
M
Minghao Li 已提交
606 607
}

S
Shengliang Guan 已提交
608
static int32_t syncHbTimerInit(SSyncNode* pSyncNode, SSyncTimer* pSyncTimer, SRaftId destId) {
609 610 611 612 613
  pSyncTimer->pTimer = NULL;
  pSyncTimer->counter = 0;
  pSyncTimer->timerMS = pSyncNode->hbBaseLine;
  pSyncTimer->timerCb = syncNodeEqPeerHeartbeatTimer;
  pSyncTimer->destId = destId;
M
Minghao Li 已提交
614
  pSyncTimer->timeStamp = taosGetTimestampMs();
615 616 617 618
  atomic_store_64(&pSyncTimer->logicClock, 0);
  return 0;
}

S
Shengliang Guan 已提交
619
static int32_t syncHbTimerStart(SSyncNode* pSyncNode, SSyncTimer* pSyncTimer) {
620
  int32_t ret = 0;
S
Shengliang Guan 已提交
621
  int64_t tsNow = taosGetTimestampMs();
S
Shengliang Guan 已提交
622
  if (syncIsInit()) {
623 624 625 626 627 628
    SSyncHbTimerData* pData = syncHbTimerDataAcquire(pSyncTimer->hbDataRid);
    if (pData == NULL) {
      pData = taosMemoryMalloc(sizeof(SSyncHbTimerData));
      pData->rid = syncHbTimerDataAdd(pData);
    }
    pSyncTimer->hbDataRid = pData->rid;
S
Shengliang Guan 已提交
629
    pSyncTimer->timeStamp = tsNow;
630 631

    pData->syncNodeRid = pSyncNode->rid;
632 633 634
    pData->pTimer = pSyncTimer;
    pData->destId = pSyncTimer->destId;
    pData->logicClock = pSyncTimer->logicClock;
S
Shengliang Guan 已提交
635
    pData->execTime = tsNow + pSyncTimer->timerMS;
M
Minghao Li 已提交
636

637 638
    taosTmrReset(pSyncTimer->timerCb, pSyncTimer->timerMS / HEARTBEAT_TICK_NUM, (void*)(pData->rid),
                 syncEnv()->pTimerManager, &pSyncTimer->pTimer);
639 640 641 642 643 644
  } else {
    sError("vgId:%d, start ctrl hb timer error, sync env is stop", pSyncNode->vgId);
  }
  return ret;
}

S
Shengliang Guan 已提交
645
static int32_t syncHbTimerStop(SSyncNode* pSyncNode, SSyncTimer* pSyncTimer) {
646 647 648 649
  int32_t ret = 0;
  atomic_add_fetch_64(&pSyncTimer->logicClock, 1);
  taosTmrStop(pSyncTimer->pTimer);
  pSyncTimer->pTimer = NULL;
650 651
  syncHbTimerDataRemove(pSyncTimer->hbDataRid);
  pSyncTimer->hbDataRid = -1;
652 653 654
  return ret;
}

655
int32_t syncNodeLogStoreRestoreOnNeed(SSyncNode* pNode) {
S
Shengliang Guan 已提交
656 657 658
  ASSERTS(pNode->pLogStore != NULL, "log store not created");
  ASSERTS(pNode->pFsm != NULL, "pFsm not registered");
  ASSERTS(pNode->pFsm->FpGetSnapshotInfo != NULL, "FpGetSnapshotInfo not registered");
659 660 661
  SSnapshot snapshot = {0};
  pNode->pFsm->FpGetSnapshotInfo(pNode->pFsm, &snapshot);

662 663 664 665 666
  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)) {
S
Shengliang Guan 已提交
667
      sError("vgId:%d, failed to restore log store from snapshot since %s. lastVer:%" PRId64 ", snapshotVer:%" PRId64,
668 669 670 671 672 673 674
             pNode->vgId, terrstr(), lastVer, commitIndex);
      return -1;
    }
  }
  return 0;
}

M
Minghao Li 已提交
675
// open/close --------------
S
Shengliang Guan 已提交
676 677
SSyncNode* syncNodeOpen(SSyncInfo* pSyncInfo) {
  SSyncNode* pSyncNode = taosMemoryCalloc(1, sizeof(SSyncNode));
678 679 680 681
  if (pSyncNode == NULL) {
    terrno = TSDB_CODE_OUT_OF_MEMORY;
    goto _error;
  }
M
Minghao Li 已提交
682

M
Minghao Li 已提交
683 684 685 686
  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());
687
      goto _error;
M
Minghao Li 已提交
688
    }
689
  }
M
Minghao Li 已提交
690

691 692 693
  memcpy(pSyncNode->path, pSyncInfo->path, sizeof(pSyncNode->path));
  snprintf(pSyncNode->raftStorePath, sizeof(pSyncNode->raftStorePath), "%s%sraft_store.json", pSyncInfo->path,
           TD_DIRSEP);
S
Shengliang Guan 已提交
694
  snprintf(pSyncNode->configPath, sizeof(pSyncNode->configPath), "%s%sraft_config.json", pSyncInfo->path, TD_DIRSEP);
695

696
  if (!taosCheckExistFile(pSyncNode->configPath)) {
M
Minghao Li 已提交
697
    // create a new raft config file
698
    sInfo("vgId:%d, create a new raft config file", pSyncNode->vgId);
699 700 701 702 703 704 705 706 707 708
    pSyncNode->raftCfg.isStandBy = pSyncInfo->isStandBy;
    pSyncNode->raftCfg.snapshotStrategy = pSyncInfo->snapshotStrategy;
    pSyncNode->raftCfg.lastConfigIndex = SYNC_INDEX_INVALID;
    pSyncNode->raftCfg.batchSize = pSyncInfo->batchSize;
    pSyncNode->raftCfg.cfg = pSyncInfo->syncCfg;
    pSyncNode->raftCfg.configIndexCount = 1;
    pSyncNode->raftCfg.configIndexArr[0] = -1;

    if (syncWriteCfgFile(pSyncNode) != 0) {
      sError("vgId:%d, failed to create sync cfg file", pSyncNode->vgId);
H
Hongze Cheng 已提交
709
      goto _error;
710
    }
711 712
  } else {
    // update syncCfg by raft_config.json
713 714
    if (syncReadCfgFile(pSyncNode) != 0) {
      sError("vgId:%d, failed to read sync cfg file", pSyncNode->vgId);
H
Hongze Cheng 已提交
715
      goto _error;
716
    }
S
Shengliang Guan 已提交
717

718
    if (pSyncInfo->syncCfg.replicaNum > 0 && syncIsConfigChanged(&pSyncNode->raftCfg.cfg, &pSyncInfo->syncCfg)) {
S
Shengliang Guan 已提交
719
      sInfo("vgId:%d, use sync config from input options and write to cfg file", pSyncNode->vgId);
720 721 722
      pSyncNode->raftCfg.cfg = pSyncInfo->syncCfg;
      if (syncWriteCfgFile(pSyncNode) != 0) {
        sError("vgId:%d, failed to write sync cfg file", pSyncNode->vgId);
S
Shengliang Guan 已提交
723 724
        goto _error;
      }
S
Shengliang Guan 已提交
725
    } else {
726 727
      sInfo("vgId:%d, use sync config from sync cfg file", pSyncNode->vgId);
      pSyncInfo->syncCfg = pSyncNode->raftCfg.cfg;
S
Shengliang Guan 已提交
728
    }
M
Minghao Li 已提交
729 730
  }

M
Minghao Li 已提交
731
  // init by SSyncInfo
M
Minghao Li 已提交
732
  pSyncNode->vgId = pSyncInfo->vgId;
S
Shengliang Guan 已提交
733
  SSyncCfg* pCfg = &pSyncNode->raftCfg.cfg;
734
  bool      updated = false;
S
Shengliang Guan 已提交
735
  sInfo("vgId:%d, start to open sync node, replica:%d selfIndex:%d", pSyncNode->vgId, pCfg->replicaNum, pCfg->myIndex);
S
Shengliang Guan 已提交
736 737
  for (int32_t i = 0; i < pCfg->replicaNum; ++i) {
    SNodeInfo* pNode = &pCfg->nodeInfo[i];
738 739 740
    if (tmsgUpdateDnodeInfo(&pNode->nodeId, &pNode->clusterId, pNode->nodeFqdn, &pNode->nodePort)) {
      updated = true;
    }
741 742
    sInfo("vgId:%d, index:%d ep:%s:%u dnode:%d cluster:%" PRId64, pSyncNode->vgId, i, pNode->nodeFqdn, pNode->nodePort,
          pNode->nodeId, pNode->clusterId);
S
Shengliang Guan 已提交
743 744
  }

745 746 747 748 749 750 751 752
  if (updated) {
    sInfo("vgId:%d, save config info since dnode info changed", pSyncNode->vgId);
    if (syncWriteCfgFile(pSyncNode) != 0) {
      sError("vgId:%d, failed to write sync cfg file on dnode info updated", pSyncNode->vgId);
      goto _error;
    }
  }

M
Minghao Li 已提交
753
  pSyncNode->pWal = pSyncInfo->pWal;
S
Shengliang Guan 已提交
754
  pSyncNode->msgcb = pSyncInfo->msgcb;
S
Shengliang Guan 已提交
755 756 757
  pSyncNode->syncSendMSg = pSyncInfo->syncSendMSg;
  pSyncNode->syncEqMsg = pSyncInfo->syncEqMsg;
  pSyncNode->syncEqCtrlMsg = pSyncInfo->syncEqCtrlMsg;
M
Minghao Li 已提交
758

B
Benguang Zhao 已提交
759 760 761
  // create raft log ring buffer
  pSyncNode->pLogBuf = syncLogBufferCreate();
  if (pSyncNode->pLogBuf == NULL) {
762
    sError("failed to init sync log buffer since %s. vgId:%d", terrstr(), pSyncNode->vgId);
B
Benguang Zhao 已提交
763 764 765
    goto _error;
  }

M
Minghao Li 已提交
766
  // init internal
767
  pSyncNode->myNodeInfo = pSyncNode->raftCfg.cfg.nodeInfo[pSyncNode->raftCfg.cfg.myIndex];
768
  if (!syncUtilNodeInfo2RaftId(&pSyncNode->myNodeInfo, pSyncNode->vgId, &pSyncNode->myRaftId)) {
S
Shengliang Guan 已提交
769
    sError("vgId:%d, failed to determine my raft member id", pSyncNode->vgId);
H
Hongze Cheng 已提交
770
    goto _error;
771
  }
M
Minghao Li 已提交
772

M
Minghao Li 已提交
773
  // init peersNum, peers, peersId
774
  pSyncNode->peersNum = pSyncNode->raftCfg.cfg.replicaNum - 1;
S
Shengliang Guan 已提交
775
  int32_t j = 0;
776 777 778 779
  for (int32_t i = 0; i < pSyncNode->raftCfg.cfg.replicaNum; ++i) {
    if (i != pSyncNode->raftCfg.cfg.myIndex) {
      pSyncNode->peersNodeInfo[j] = pSyncNode->raftCfg.cfg.nodeInfo[i];
      syncUtilNodeInfo2EpSet(&pSyncNode->peersNodeInfo[j], &pSyncNode->peersEpset[j]);
M
Minghao Li 已提交
780 781 782
      j++;
    }
  }
S
Shengliang Guan 已提交
783
  for (int32_t i = 0; i < pSyncNode->peersNum; ++i) {
784
    if (!syncUtilNodeInfo2RaftId(&pSyncNode->peersNodeInfo[i], pSyncNode->vgId, &pSyncNode->peersId[i])) {
S
Shengliang Guan 已提交
785
      sError("vgId:%d, failed to determine raft member id, peer:%d", pSyncNode->vgId, i);
H
Hongze Cheng 已提交
786
      goto _error;
787
    }
M
Minghao Li 已提交
788
  }
M
Minghao Li 已提交
789

M
Minghao Li 已提交
790
  // init replicaNum, replicasId
791 792 793
  pSyncNode->replicaNum = pSyncNode->raftCfg.cfg.replicaNum;
  for (int32_t i = 0; i < pSyncNode->raftCfg.cfg.replicaNum; ++i) {
    if (!syncUtilNodeInfo2RaftId(&pSyncNode->raftCfg.cfg.nodeInfo[i], pSyncNode->vgId, &pSyncNode->replicasId[i])) {
S
Shengliang Guan 已提交
794
      sError("vgId:%d, failed to determine raft member id, replica:%d", pSyncNode->vgId, i);
H
Hongze Cheng 已提交
795
      goto _error;
796
    }
M
Minghao Li 已提交
797 798
  }

M
Minghao Li 已提交
799
  // init raft algorithm
M
Minghao Li 已提交
800
  pSyncNode->pFsm = pSyncInfo->pFsm;
801
  pSyncInfo->pFsm = NULL;
802
  pSyncNode->quorum = syncUtilQuorum(pSyncNode->raftCfg.cfg.replicaNum);
M
Minghao Li 已提交
803 804
  pSyncNode->leaderCache = EMPTY_RAFT_ID;

M
Minghao Li 已提交
805
  // init life cycle outside
M
Minghao Li 已提交
806

M
Minghao Li 已提交
807 808 809 810 811 812 813 814 815 816 817 818 819 820 821 822 823 824 825 826 827 828 829 830
  // 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 已提交
831
  // init TLA+ server vars
M
syncInt  
Minghao Li 已提交
832
  pSyncNode->state = TAOS_SYNC_STATE_FOLLOWER;
833
  if (raftStoreOpen(pSyncNode) != 0) {
S
Shengliang Guan 已提交
834
    sError("vgId:%d, failed to open raft store at path %s", pSyncNode->vgId, pSyncNode->raftStorePath);
835 836
    goto _error;
  }
M
Minghao Li 已提交
837

M
Minghao Li 已提交
838
  // init TLA+ candidate vars
M
Minghao Li 已提交
839
  pSyncNode->pVotesGranted = voteGrantedCreate(pSyncNode);
840
  if (pSyncNode->pVotesGranted == NULL) {
S
Shengliang Guan 已提交
841
    sError("vgId:%d, failed to create VotesGranted", pSyncNode->vgId);
842 843
    goto _error;
  }
M
Minghao Li 已提交
844
  pSyncNode->pVotesRespond = votesRespondCreate(pSyncNode);
845
  if (pSyncNode->pVotesRespond == NULL) {
S
Shengliang Guan 已提交
846
    sError("vgId:%d, failed to create VotesRespond", pSyncNode->vgId);
847 848
    goto _error;
  }
M
Minghao Li 已提交
849

M
Minghao Li 已提交
850 851
  // init TLA+ leader vars
  pSyncNode->pNextIndex = syncIndexMgrCreate(pSyncNode);
852
  if (pSyncNode->pNextIndex == NULL) {
S
Shengliang Guan 已提交
853
    sError("vgId:%d, failed to create SyncIndexMgr", pSyncNode->vgId);
854 855
    goto _error;
  }
M
Minghao Li 已提交
856
  pSyncNode->pMatchIndex = syncIndexMgrCreate(pSyncNode);
857
  if (pSyncNode->pMatchIndex == NULL) {
S
Shengliang Guan 已提交
858
    sError("vgId:%d, failed to create SyncIndexMgr", pSyncNode->vgId);
859 860
    goto _error;
  }
M
Minghao Li 已提交
861 862 863

  // init TLA+ log vars
  pSyncNode->pLogStore = logStoreCreate(pSyncNode);
864
  if (pSyncNode->pLogStore == NULL) {
S
Shengliang Guan 已提交
865
    sError("vgId:%d, failed to create SyncLogStore", pSyncNode->vgId);
866 867
    goto _error;
  }
868 869 870 871

  SyncIndex commitIndex = SYNC_INDEX_INVALID;
  if (pSyncNode->pFsm != NULL && pSyncNode->pFsm->FpGetSnapshotInfo != NULL) {
    SSnapshot snapshot = {0};
872
    pSyncNode->pFsm->FpGetSnapshotInfo(pSyncNode->pFsm, &snapshot);
873 874
    if (snapshot.lastApplyIndex > commitIndex) {
      commitIndex = snapshot.lastApplyIndex;
S
Shengliang Guan 已提交
875
      sNTrace(pSyncNode, "reset commit index by snapshot");
876 877 878
    }
  }
  pSyncNode->commitIndex = commitIndex;
879
  sInfo("vgId:%d, sync node commitIndex initialized as %" PRId64, pSyncNode->vgId, pSyncNode->commitIndex);
M
Minghao Li 已提交
880

881
  // restore log store on need
882
  if (syncNodeLogStoreRestoreOnNeed(pSyncNode) < 0) {
883
    sError("vgId:%d, failed to restore log store since %s.", pSyncNode->vgId, terrstr());
884 885
    goto _error;
  }
886

M
Minghao Li 已提交
887 888
  // timer ms init
  pSyncNode->pingBaseLine = PING_TIMER_MS;
889 890
  pSyncNode->electBaseLine = tsElectInterval;
  pSyncNode->hbBaseLine = tsHeartbeatInterval;
M
Minghao Li 已提交
891

M
Minghao Li 已提交
892
  // init ping timer
M
Minghao Li 已提交
893
  pSyncNode->pPingTimer = NULL;
M
Minghao Li 已提交
894
  pSyncNode->pingTimerMS = pSyncNode->pingBaseLine;
M
Minghao Li 已提交
895 896
  atomic_store_64(&pSyncNode->pingTimerLogicClock, 0);
  atomic_store_64(&pSyncNode->pingTimerLogicClockUser, 0);
M
Minghao Li 已提交
897
  pSyncNode->FpPingTimerCB = syncNodeEqPingTimer;
M
Minghao Li 已提交
898
  pSyncNode->pingTimerCounter = 0;
M
Minghao Li 已提交
899

M
Minghao Li 已提交
900 901
  // init elect timer
  pSyncNode->pElectTimer = NULL;
M
Minghao Li 已提交
902
  pSyncNode->electTimerMS = syncUtilElectRandomMS(pSyncNode->electBaseLine, 2 * pSyncNode->electBaseLine);
M
Minghao Li 已提交
903
  atomic_store_64(&pSyncNode->electTimerLogicClock, 0);
M
Minghao Li 已提交
904
  pSyncNode->FpElectTimerCB = syncNodeEqElectTimer;
M
Minghao Li 已提交
905 906 907 908
  pSyncNode->electTimerCounter = 0;

  // init heartbeat timer
  pSyncNode->pHeartbeatTimer = NULL;
M
Minghao Li 已提交
909
  pSyncNode->heartbeatTimerMS = pSyncNode->hbBaseLine;
M
Minghao Li 已提交
910 911
  atomic_store_64(&pSyncNode->heartbeatTimerLogicClock, 0);
  atomic_store_64(&pSyncNode->heartbeatTimerLogicClockUser, 0);
M
Minghao Li 已提交
912
  pSyncNode->FpHeartbeatTimerCB = syncNodeEqHeartbeatTimer;
M
Minghao Li 已提交
913 914
  pSyncNode->heartbeatTimerCounter = 0;

915 916 917 918 919
  // 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 已提交
920
  // tools
M
Minghao Li 已提交
921
  pSyncNode->pSyncRespMgr = syncRespMgrCreate(pSyncNode, SYNC_RESP_TTL_MS);
922
  if (pSyncNode->pSyncRespMgr == NULL) {
S
Shengliang Guan 已提交
923
    sError("vgId:%d, failed to create SyncRespMgr", pSyncNode->vgId);
924 925
    goto _error;
  }
M
Minghao Li 已提交
926

927 928
  // restore state
  pSyncNode->restoreFinish = false;
929

M
Minghao Li 已提交
930
  // snapshot senders
S
Shengliang Guan 已提交
931
  for (int32_t i = 0; i < TSDB_MAX_REPLICA; ++i) {
M
Minghao Li 已提交
932
    SSyncSnapshotSender* pSender = snapshotSenderCreate(pSyncNode, i);
933 934 935 936
    if (pSender == NULL) return NULL;

    pSyncNode->senders[i] = pSender;
    sSDebug(pSender, "snapshot sender create while open sync node, data:%p", pSender);
M
Minghao Li 已提交
937 938 939
  }

  // snapshot receivers
940
  pSyncNode->pNewNodeReceiver = snapshotReceiverCreate(pSyncNode, EMPTY_RAFT_ID);
941 942 943
  if (pSyncNode->pNewNodeReceiver == NULL) return NULL;
  sRDebug(pSyncNode->pNewNodeReceiver, "snapshot receiver create while open sync node, data:%p",
          pSyncNode->pNewNodeReceiver);
M
Minghao Li 已提交
944

M
Minghao Li 已提交
945 946 947
  // is config changing
  pSyncNode->changing = false;

B
Benguang Zhao 已提交
948
  // replication mgr
949 950 951 952
  if (syncNodeLogReplMgrInit(pSyncNode) < 0) {
    sError("vgId:%d, failed to init repl mgr since %s.", pSyncNode->vgId, terrstr());
    goto _error;
  }
B
Benguang Zhao 已提交
953

M
Minghao Li 已提交
954
  // peer state
955 956 957 958
  if (syncNodePeerStateInit(pSyncNode) < 0) {
    sError("vgId:%d, failed to init peer stat since %s.", pSyncNode->vgId, terrstr());
    goto _error;
  }
M
Minghao Li 已提交
959

B
Benguang Zhao 已提交
960
  //
M
Minghao Li 已提交
961 962 963
  // min match index
  pSyncNode->minMatchIndex = SYNC_INDEX_INVALID;

M
Minghao Li 已提交
964
  // start in syncNodeStart
M
Minghao Li 已提交
965
  // start raft
M
Minghao Li 已提交
966
  // syncNodeBecomeFollower(pSyncNode);
M
Minghao Li 已提交
967

M
Minghao Li 已提交
968 969
  int64_t timeNow = taosGetTimestampMs();
  pSyncNode->startTime = timeNow;
970
  pSyncNode->leaderTime = timeNow;
M
Minghao Li 已提交
971 972
  pSyncNode->lastReplicateTime = timeNow;

973 974 975
  // snapshotting
  atomic_store_64(&pSyncNode->snapshottingIndex, SYNC_INDEX_INVALID);

B
Benguang Zhao 已提交
976 977
  // init log buffer
  if (syncLogBufferInit(pSyncNode->pLogBuf, pSyncNode) < 0) {
978
    sError("vgId:%d, failed to init sync log buffer since %s", pSyncNode->vgId, terrstr());
979
    goto _error;
B
Benguang Zhao 已提交
980 981
  }

982
  pSyncNode->isStart = true;
983 984 985
  pSyncNode->electNum = 0;
  pSyncNode->becomeLeaderNum = 0;
  pSyncNode->configChangeNum = 0;
986 987
  pSyncNode->hbSlowNum = 0;
  pSyncNode->hbrSlowNum = 0;
M
Minghao Li 已提交
988
  pSyncNode->tmrRoutineNum = 0;
989

990 991
  sNInfo(pSyncNode, "sync open, node:%p electInterval:%d heartbeatInterval:%d heartbeatTimeout:%d", pSyncNode,
         tsElectInterval, tsHeartbeatInterval, tsHeartbeatTimeout);
M
Minghao Li 已提交
992
  return pSyncNode;
993 994 995

_error:
  if (pSyncInfo->pFsm) {
H
Hongze Cheng 已提交
996 997
    taosMemoryFree(pSyncInfo->pFsm);
    pSyncInfo->pFsm = NULL;
998 999 1000 1001
  }
  syncNodeClose(pSyncNode);
  pSyncNode = NULL;
  return NULL;
M
Minghao Li 已提交
1002 1003
}

M
Minghao Li 已提交
1004 1005
void syncNodeMaybeUpdateCommitBySnapshot(SSyncNode* pSyncNode) {
  if (pSyncNode->pFsm != NULL && pSyncNode->pFsm->FpGetSnapshotInfo != NULL) {
1006 1007
    SSnapshot snapshot = {0};
    pSyncNode->pFsm->FpGetSnapshotInfo(pSyncNode->pFsm, &snapshot);
M
Minghao Li 已提交
1008 1009 1010 1011 1012 1013
    if (snapshot.lastApplyIndex > pSyncNode->commitIndex) {
      pSyncNode->commitIndex = snapshot.lastApplyIndex;
    }
  }
}

B
Benguang Zhao 已提交
1014
int32_t syncNodeRestore(SSyncNode* pSyncNode) {
S
Shengliang Guan 已提交
1015 1016
  ASSERTS(pSyncNode->pLogStore != NULL, "log store not created");
  ASSERTS(pSyncNode->pLogBuf != NULL, "ring log buffer not created");
B
Benguang Zhao 已提交
1017 1018 1019 1020

  SyncIndex lastVer = pSyncNode->pLogStore->syncLogLastIndex(pSyncNode->pLogStore);
  SyncIndex commitIndex = pSyncNode->pLogStore->syncLogCommitIndex(pSyncNode->pLogStore);
  SyncIndex endIndex = pSyncNode->pLogBuf->endIndex;
1021 1022
  if (lastVer != -1 && endIndex != lastVer + 1) {
    terrno = TSDB_CODE_WAL_LOG_INCOMPLETE;
S
Shengliang Guan 已提交
1023
    sError("vgId:%d, failed to restore sync node since %s. expected lastLogIndex:%" PRId64 ", lastVer:%" PRId64 "",
1024 1025 1026
           pSyncNode->vgId, terrstr(), endIndex - 1, lastVer);
    return -1;
  }
B
Benguang Zhao 已提交
1027

1028
  ASSERT(endIndex == lastVer + 1);
1029 1030
  pSyncNode->commitIndex = TMAX(pSyncNode->commitIndex, commitIndex);
  sInfo("vgId:%d, restore sync until commitIndex:%" PRId64, pSyncNode->vgId, pSyncNode->commitIndex);
B
Benguang Zhao 已提交
1031

1032
  if (syncLogBufferCommit(pSyncNode->pLogBuf, pSyncNode, pSyncNode->commitIndex) < 0) {
B
Benguang Zhao 已提交
1033 1034 1035 1036 1037 1038 1039 1040 1041
    return -1;
  }

  return 0;
}

int32_t syncNodeStart(SSyncNode* pSyncNode) {
  // start raft
  if (pSyncNode->replicaNum == 1) {
S
Shengliang Guan 已提交
1042
    raftStoreNextTerm(pSyncNode);
B
Benguang Zhao 已提交
1043 1044 1045 1046 1047 1048 1049 1050 1051 1052
    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);
1053 1054
  if (ret != 0) {
    sError("vgId:%d, failed to start ping timer since %s", pSyncNode->vgId, terrstr());
1055
  }
1056
  return ret;
M
Minghao Li 已提交
1057 1058
}

B
Benguang Zhao 已提交
1059
int32_t syncNodeStartStandBy(SSyncNode* pSyncNode) {
M
Minghao Li 已提交
1060 1061 1062 1063 1064 1065 1066
  // 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);
1067 1068 1069 1070
  if (ret < 0) {
    sError("vgId:%d, failed to restart elect timer since %s", pSyncNode->vgId, terrstr());
    return -1;
  }
1071

1072
  ret = syncNodeStartPingTimer(pSyncNode);
1073 1074 1075 1076
  if (ret < 0) {
    sError("vgId:%d, failed to start ping timer since %s", pSyncNode->vgId, terrstr());
    return -1;
  }
B
Benguang Zhao 已提交
1077
  return ret;
M
Minghao Li 已提交
1078 1079
}

M
Minghao Li 已提交
1080
void syncNodePreClose(SSyncNode* pSyncNode) {
1081 1082 1083 1084
  ASSERT(pSyncNode != NULL);
  ASSERT(pSyncNode->pFsm != NULL);
  ASSERT(pSyncNode->pFsm->FpApplyQueueItems != NULL);

M
Minghao Li 已提交
1085 1086 1087 1088 1089
  // stop elect timer
  syncNodeStopElectTimer(pSyncNode);

  // stop heartbeat timer
  syncNodeStopHeartbeatTimer(pSyncNode);
1090

1091 1092 1093
  // stop ping timer
  syncNodeStopPingTimer(pSyncNode);

1094 1095
  // clean rsp
  syncRespCleanRsp(pSyncNode->pSyncRespMgr);
M
Minghao Li 已提交
1096 1097
}

1098 1099 1100
void syncNodePostClose(SSyncNode* pSyncNode) {
  if (pSyncNode->pNewNodeReceiver != NULL) {
    if (snapshotReceiverIsStart(pSyncNode->pNewNodeReceiver)) {
S
Shengliang Guan 已提交
1101
      snapshotReceiverStop(pSyncNode->pNewNodeReceiver);
1102 1103 1104 1105 1106 1107 1108
    }

    sDebug("vgId:%d, snapshot receiver destroy while preclose sync node, data:%p", pSyncNode->vgId,
           pSyncNode->pNewNodeReceiver);
    snapshotReceiverDestroy(pSyncNode->pNewNodeReceiver);
    pSyncNode->pNewNodeReceiver = NULL;
  }
M
Minghao Li 已提交
1109 1110
}

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

M
Minghao Li 已提交
1113
void syncNodeClose(SSyncNode* pSyncNode) {
S
Shengliang Guan 已提交
1114
  if (pSyncNode == NULL) return;
1115
  sNInfo(pSyncNode, "sync close, node:%p", pSyncNode);
M
Minghao Li 已提交
1116

1117 1118
  syncRespCleanRsp(pSyncNode->pSyncRespMgr);

1119 1120 1121
  syncNodeStopPingTimer(pSyncNode);
  syncNodeStopElectTimer(pSyncNode);
  syncNodeStopHeartbeatTimer(pSyncNode);
B
Benguang Zhao 已提交
1122
  syncNodeLogReplMgrDestroy(pSyncNode);
1123

M
Minghao Li 已提交
1124
  syncRespMgrDestroy(pSyncNode->pSyncRespMgr);
1125
  pSyncNode->pSyncRespMgr = NULL;
M
Minghao Li 已提交
1126
  voteGrantedDestroy(pSyncNode->pVotesGranted);
1127
  pSyncNode->pVotesGranted = NULL;
M
Minghao Li 已提交
1128
  votesRespondDestory(pSyncNode->pVotesRespond);
1129
  pSyncNode->pVotesRespond = NULL;
M
Minghao Li 已提交
1130
  syncIndexMgrDestroy(pSyncNode->pNextIndex);
1131
  pSyncNode->pNextIndex = NULL;
M
Minghao Li 已提交
1132
  syncIndexMgrDestroy(pSyncNode->pMatchIndex);
1133
  pSyncNode->pMatchIndex = NULL;
M
Minghao Li 已提交
1134
  logStoreDestory(pSyncNode->pLogStore);
1135
  pSyncNode->pLogStore = NULL;
B
Benguang Zhao 已提交
1136 1137
  syncLogBufferDestroy(pSyncNode->pLogBuf);
  pSyncNode->pLogBuf = NULL;
M
Minghao Li 已提交
1138

S
Shengliang Guan 已提交
1139
  for (int32_t i = 0; i < TSDB_MAX_REPLICA; ++i) {
1140 1141
    if (pSyncNode->senders[i] != NULL) {
      sDebug("vgId:%d, snapshot sender destroy while close, data:%p", pSyncNode->vgId, pSyncNode->senders[i]);
1142

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

1147 1148
      snapshotSenderDestroy(pSyncNode->senders[i]);
      pSyncNode->senders[i] = NULL;
M
Minghao Li 已提交
1149 1150 1151
    }
  }

M
Minghao Li 已提交
1152
  if (pSyncNode->pNewNodeReceiver != NULL) {
1153
    if (snapshotReceiverIsStart(pSyncNode->pNewNodeReceiver)) {
S
Shengliang Guan 已提交
1154
      snapshotReceiverStop(pSyncNode->pNewNodeReceiver);
1155 1156
    }

1157
    sDebug("vgId:%d, snapshot receiver destroy while close, data:%p", pSyncNode->vgId, pSyncNode->pNewNodeReceiver);
M
Minghao Li 已提交
1158 1159 1160 1161
    snapshotReceiverDestroy(pSyncNode->pNewNodeReceiver);
    pSyncNode->pNewNodeReceiver = NULL;
  }

1162 1163 1164 1165
  if (pSyncNode->pFsm != NULL) {
    taosMemoryFree(pSyncNode->pFsm);
  }

1166 1167
  raftStoreClose(pSyncNode);

1168
  taosMemoryFree(pSyncNode);
M
Minghao Li 已提交
1169 1170
}

1171
ESyncStrategy syncNodeStrategy(SSyncNode* pSyncNode) { return pSyncNode->raftCfg.snapshotStrategy; }
M
Minghao Li 已提交
1172

M
Minghao Li 已提交
1173 1174 1175
// timer control --------------
int32_t syncNodeStartPingTimer(SSyncNode* pSyncNode) {
  int32_t ret = 0;
S
Shengliang Guan 已提交
1176 1177
  if (syncIsInit()) {
    taosTmrReset(pSyncNode->FpPingTimerCB, pSyncNode->pingTimerMS, pSyncNode, syncEnv()->pTimerManager,
1178 1179 1180
                 &pSyncNode->pPingTimer);
    atomic_store_64(&pSyncNode->pingTimerLogicClock, pSyncNode->pingTimerLogicClockUser);
  } else {
M
Minghao Li 已提交
1181
    sError("vgId:%d, start ping timer error, sync env is stop", pSyncNode->vgId);
1182
  }
M
Minghao Li 已提交
1183 1184 1185 1186 1187 1188 1189 1190 1191 1192 1193 1194 1195
  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 已提交
1196
  if (syncIsInit()) {
1197
    pSyncNode->electTimerMS = ms;
S
Shengliang Guan 已提交
1198

1199 1200 1201 1202 1203
    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 已提交
1204

M
Minghao Li 已提交
1205
    taosTmrReset(pSyncNode->FpElectTimerCB, pSyncNode->electTimerMS, (void*)(pSyncNode->rid), syncEnv()->pTimerManager,
1206
                 &pSyncNode->pElectTimer);
1207

1208
  } else {
M
Minghao Li 已提交
1209
    sError("vgId:%d, start elect timer error, sync env is stop", pSyncNode->vgId);
1210
  }
M
Minghao Li 已提交
1211 1212 1213 1214 1215
  return ret;
}

int32_t syncNodeStopElectTimer(SSyncNode* pSyncNode) {
  int32_t ret = 0;
M
Minghao Li 已提交
1216
  atomic_add_fetch_64(&pSyncNode->electTimerLogicClock, 1);
M
Minghao Li 已提交
1217 1218
  taosTmrStop(pSyncNode->pElectTimer);
  pSyncNode->pElectTimer = NULL;
1219

M
Minghao Li 已提交
1220 1221 1222 1223 1224 1225 1226 1227 1228 1229
  return ret;
}

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

1230
void syncNodeResetElectTimer(SSyncNode* pSyncNode) {
M
Minghao Li 已提交
1231 1232
  int32_t electMS;

1233
  if (pSyncNode->raftCfg.isStandBy) {
M
Minghao Li 已提交
1234 1235 1236 1237
    electMS = TIMER_MAX_MS;
  } else {
    electMS = syncUtilElectRandomMS(pSyncNode->electBaseLine, 2 * pSyncNode->electBaseLine);
  }
1238 1239

  (void)syncNodeRestartElectTimer(pSyncNode, electMS);
1240

S
Shengliang Guan 已提交
1241 1242
  sNTrace(pSyncNode, "reset elect timer, min:%d, max:%d, ms:%d", pSyncNode->electBaseLine, 2 * pSyncNode->electBaseLine,
          electMS);
M
Minghao Li 已提交
1243 1244
}

M
Minghao Li 已提交
1245
static int32_t syncNodeDoStartHeartbeatTimer(SSyncNode* pSyncNode) {
M
Minghao Li 已提交
1246
  int32_t ret = 0;
S
Shengliang Guan 已提交
1247 1248
  if (syncIsInit()) {
    taosTmrReset(pSyncNode->FpHeartbeatTimerCB, pSyncNode->heartbeatTimerMS, pSyncNode, syncEnv()->pTimerManager,
1249 1250 1251
                 &pSyncNode->pHeartbeatTimer);
    atomic_store_64(&pSyncNode->heartbeatTimerLogicClock, pSyncNode->heartbeatTimerLogicClockUser);
  } else {
M
Minghao Li 已提交
1252
    sError("vgId:%d, start heartbeat timer error, sync env is stop", pSyncNode->vgId);
1253
  }
1254

S
Shengliang Guan 已提交
1255
  sNTrace(pSyncNode, "start heartbeat timer, ms:%d", pSyncNode->heartbeatTimerMS);
M
Minghao Li 已提交
1256 1257 1258
  return ret;
}

M
Minghao Li 已提交
1259
int32_t syncNodeStartHeartbeatTimer(SSyncNode* pSyncNode) {
1260
  int32_t ret = 0;
M
Minghao Li 已提交
1261

1262
#if 0
M
Minghao Li 已提交
1263
  pSyncNode->heartbeatTimerMS = pSyncNode->hbBaseLine;
1264 1265
  ret = syncNodeDoStartHeartbeatTimer(pSyncNode);
#endif
1266

S
Shengliang Guan 已提交
1267
  for (int32_t i = 0; i < pSyncNode->peersNum; ++i) {
1268
    SSyncTimer* pSyncTimer = syncNodeGetHbTimer(pSyncNode, &(pSyncNode->peersId[i]));
M
Minghao Li 已提交
1269 1270 1271
    if (pSyncTimer != NULL) {
      syncHbTimerStart(pSyncNode, pSyncTimer);
    }
1272
  }
1273

M
Minghao Li 已提交
1274 1275 1276
  return ret;
}

M
Minghao Li 已提交
1277 1278
int32_t syncNodeStopHeartbeatTimer(SSyncNode* pSyncNode) {
  int32_t ret = 0;
1279 1280

#if 0
M
Minghao Li 已提交
1281 1282 1283
  atomic_add_fetch_64(&pSyncNode->heartbeatTimerLogicClockUser, 1);
  taosTmrStop(pSyncNode->pHeartbeatTimer);
  pSyncNode->pHeartbeatTimer = NULL;
1284
#endif
1285

S
Shengliang Guan 已提交
1286
  for (int32_t i = 0; i < pSyncNode->peersNum; ++i) {
1287
    SSyncTimer* pSyncTimer = syncNodeGetHbTimer(pSyncNode, &(pSyncNode->peersId[i]));
M
Minghao Li 已提交
1288 1289 1290
    if (pSyncTimer != NULL) {
      syncHbTimerStop(pSyncNode, pSyncTimer);
    }
1291
  }
1292

M
Minghao Li 已提交
1293 1294 1295
  return ret;
}

1296 1297 1298 1299 1300 1301
int32_t syncNodeRestartHeartbeatTimer(SSyncNode* pSyncNode) {
  syncNodeStopHeartbeatTimer(pSyncNode);
  syncNodeStartHeartbeatTimer(pSyncNode);
  return 0;
}

1302 1303 1304 1305 1306 1307 1308 1309
int32_t syncNodeSendMsgById(const SRaftId* destRaftId, SSyncNode* pNode, SRpcMsg* pMsg) {
  SEpSet* epSet = NULL;
  for (int32_t i = 0; i < pNode->peersNum; ++i) {
    if (destRaftId->addr == pNode->peersId[i].addr) {
      epSet = &pNode->peersEpset[i];
      break;
    }
  }
1310

S
Shengliang Guan 已提交
1311
  int32_t code = -1;
1312
  if (pNode->syncSendMSg != NULL && epSet != NULL) {
M
Minghao Li 已提交
1313
    syncUtilMsgHtoN(pMsg->pCont);
1314
    pMsg->info.noResp = 1;
S
Shengliang Guan 已提交
1315 1316 1317 1318 1319 1320
    code = pNode->syncSendMSg(epSet, pMsg);
  }

  if (code < 0) {
    sError("vgId:%d, sync send msg by id error, epset:%p dnode:%d addr:%" PRId64 " err:0x%x", pNode->vgId, epSet,
           DID(destRaftId), destRaftId->addr, terrno);
S
Shengliang Guan 已提交
1321
    rpcFreeCont(pMsg->pCont);
1322
    terrno = TSDB_CODE_SYN_INTERNAL_ERROR;
M
Minghao Li 已提交
1323
  }
S
Shengliang Guan 已提交
1324 1325

  return code;
M
Minghao Li 已提交
1326 1327
}

1328
inline bool syncNodeInConfig(SSyncNode* pNode, const SSyncCfg* pCfg) {
1329 1330 1331
  bool b1 = false;
  bool b2 = false;

1332 1333 1334
  for (int32_t i = 0; i < pCfg->replicaNum; ++i) {
    if (strcmp(pCfg->nodeInfo[i].nodeFqdn, pNode->myNodeInfo.nodeFqdn) == 0 &&
        pCfg->nodeInfo[i].nodePort == pNode->myNodeInfo.nodePort) {
1335 1336 1337 1338 1339
      b1 = true;
      break;
    }
  }

1340 1341 1342 1343 1344
  for (int32_t i = 0; i < pCfg->replicaNum; ++i) {
    SRaftId raftId = {
        .addr = SYNC_ADDR(&pCfg->nodeInfo[i]),
        .vgId = pNode->vgId,
    };
1345

1346
    if (syncUtilSameId(&raftId, &pNode->myRaftId)) {
1347 1348 1349 1350 1351
      b2 = true;
      break;
    }
  }

1352
  ASSERT(b1 == b2);
1353 1354 1355
  return b1;
}

1356 1357 1358 1359 1360 1361 1362 1363 1364 1365 1366 1367 1368
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 已提交
1369
void syncNodeDoConfigChange(SSyncNode* pSyncNode, SSyncCfg* pNewConfig, SyncIndex lastConfigChangeIndex) {
1370
  SSyncCfg oldConfig = pSyncNode->raftCfg.cfg;
1371 1372 1373 1374
  if (!syncIsConfigChanged(&oldConfig, pNewConfig)) {
    sInfo("vgId:1, sync not reconfig since not changed");
    return;
  }
S
Shengliang Guan 已提交
1375

1376 1377
  pSyncNode->raftCfg.cfg = *pNewConfig;
  pSyncNode->raftCfg.lastConfigIndex = lastConfigChangeIndex;
1378

1379 1380
  pSyncNode->configChangeNum++;

M
Minghao Li 已提交
1381 1382
  bool IamInOld = syncNodeInConfig(pSyncNode, &oldConfig);
  bool IamInNew = syncNodeInConfig(pSyncNode, pNewConfig);
M
Minghao Li 已提交
1383

M
Minghao Li 已提交
1384 1385
  bool isDrop = false;
  bool isAdd = false;
M
Minghao Li 已提交
1386

M
Minghao Li 已提交
1387 1388 1389 1390
  if (IamInOld && !IamInNew) {
    isDrop = true;
  } else {
    isDrop = false;
1391
  }
1392

M
Minghao Li 已提交
1393 1394 1395 1396 1397
  if (!IamInOld && IamInNew) {
    isAdd = true;
  } else {
    isAdd = false;
  }
M
Minghao Li 已提交
1398

M
Minghao Li 已提交
1399
  // log begin config change
1400
  sNInfo(pSyncNode, "begin do config change, from %d to %d, replicas:%d", pSyncNode->vgId, oldConfig.replicaNum,
1401
         pNewConfig->replicaNum);
M
Minghao Li 已提交
1402

M
Minghao Li 已提交
1403
  if (IamInNew) {
1404
    pSyncNode->raftCfg.isStandBy = 0;  // change isStandBy to normal
M
Minghao Li 已提交
1405
  }
M
Minghao Li 已提交
1406
  if (isDrop) {
1407
    pSyncNode->raftCfg.isStandBy = 1;  // set standby
M
Minghao Li 已提交
1408 1409
  }

M
Minghao Li 已提交
1410
  // add last config index
1411
  syncAddCfgIndex(pSyncNode, lastConfigChangeIndex);
M
Minghao Li 已提交
1412

M
Minghao Li 已提交
1413 1414 1415 1416 1417 1418 1419 1420 1421
  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 已提交
1422
    for (int32_t i = 0; i < TSDB_MAX_REPLICA; ++i) {
1423
      oldSenders[i] = pSyncNode->senders[i];
S
Shengliang Guan 已提交
1424
      sSTrace(oldSenders[i], "snapshot sender save old");
M
Minghao Li 已提交
1425
    }
1426

M
Minghao Li 已提交
1427
    // init internal
1428
    pSyncNode->myNodeInfo = pSyncNode->raftCfg.cfg.nodeInfo[pSyncNode->raftCfg.cfg.myIndex];
1429
    syncUtilNodeInfo2RaftId(&pSyncNode->myNodeInfo, pSyncNode->vgId, &pSyncNode->myRaftId);
M
Minghao Li 已提交
1430 1431

    // init peersNum, peers, peersId
1432
    pSyncNode->peersNum = pSyncNode->raftCfg.cfg.replicaNum - 1;
S
Shengliang Guan 已提交
1433
    int32_t j = 0;
1434 1435 1436 1437
    for (int32_t i = 0; i < pSyncNode->raftCfg.cfg.replicaNum; ++i) {
      if (i != pSyncNode->raftCfg.cfg.myIndex) {
        pSyncNode->peersNodeInfo[j] = pSyncNode->raftCfg.cfg.nodeInfo[i];
        syncUtilNodeInfo2EpSet(&pSyncNode->peersNodeInfo[j], &pSyncNode->peersEpset[j]);
M
Minghao Li 已提交
1438 1439 1440
        j++;
      }
    }
S
Shengliang Guan 已提交
1441
    for (int32_t i = 0; i < pSyncNode->peersNum; ++i) {
1442
      syncUtilNodeInfo2RaftId(&pSyncNode->peersNodeInfo[i], pSyncNode->vgId, &pSyncNode->peersId[i]);
M
Minghao Li 已提交
1443
    }
1444

M
Minghao Li 已提交
1445
    // init replicaNum, replicasId
1446 1447 1448
    pSyncNode->replicaNum = pSyncNode->raftCfg.cfg.replicaNum;
    for (int32_t i = 0; i < pSyncNode->raftCfg.cfg.replicaNum; ++i) {
      syncUtilNodeInfo2RaftId(&pSyncNode->raftCfg.cfg.nodeInfo[i], pSyncNode->vgId, &pSyncNode->replicasId[i]);
M
Minghao Li 已提交
1449
    }
1450

1451
    // update quorum first
1452
    pSyncNode->quorum = syncUtilQuorum(pSyncNode->raftCfg.cfg.replicaNum);
1453

M
Minghao Li 已提交
1454 1455 1456 1457
    syncIndexMgrUpdate(pSyncNode->pNextIndex, pSyncNode);
    syncIndexMgrUpdate(pSyncNode->pMatchIndex, pSyncNode);
    voteGrantedUpdate(pSyncNode->pVotesGranted, pSyncNode);
    votesRespondUpdate(pSyncNode->pVotesRespond, pSyncNode);
M
Minghao Li 已提交
1458

M
Minghao Li 已提交
1459
    // reset snapshot senders
1460

M
Minghao Li 已提交
1461
    // clear new
S
Shengliang Guan 已提交
1462
    for (int32_t i = 0; i < TSDB_MAX_REPLICA; ++i) {
1463
      pSyncNode->senders[i] = NULL;
M
Minghao Li 已提交
1464
    }
M
Minghao Li 已提交
1465

M
Minghao Li 已提交
1466
    // reset new
S
Shengliang Guan 已提交
1467
    for (int32_t i = 0; i < pSyncNode->replicaNum; ++i) {
M
Minghao Li 已提交
1468 1469
      // reset sender
      bool reset = false;
S
Shengliang Guan 已提交
1470
      for (int32_t j = 0; j < TSDB_MAX_REPLICA; ++j) {
M
Minghao Li 已提交
1471
        if (syncUtilSameId(&(pSyncNode->replicasId)[i], &oldReplicasId[j]) && oldSenders[j] != NULL) {
1472 1473
          sNTrace(pSyncNode, "snapshot sender reset for:%" PRId64 ", newIndex:%d, dnode:%d, %p",
                  (pSyncNode->replicasId)[i].addr, i, DID(&pSyncNode->replicasId[i]), oldSenders[j]);
M
Minghao Li 已提交
1474

1475
          pSyncNode->senders[i] = oldSenders[j];
M
Minghao Li 已提交
1476 1477 1478 1479
          oldSenders[j] = NULL;
          reset = true;

          // reset replicaIndex
1480 1481
          int32_t oldreplicaIndex = pSyncNode->senders[i]->replicaIndex;
          pSyncNode->senders[i]->replicaIndex = i;
M
Minghao Li 已提交
1482

1483 1484
          sNTrace(pSyncNode, "snapshot sender udpate replicaIndex from %d to %d, dnode:%d, %p, reset:%d",
                  oldreplicaIndex, i, DID(&pSyncNode->replicasId[i]), pSyncNode->senders[i], reset);
M
Minghao Li 已提交
1485 1486

          break;
M
Minghao Li 已提交
1487
        }
1488 1489
      }
    }
1490

M
Minghao Li 已提交
1491
    // create new
S
Shengliang Guan 已提交
1492
    for (int32_t i = 0; i < TSDB_MAX_REPLICA; ++i) {
1493 1494 1495 1496 1497 1498 1499 1500
      if (pSyncNode->senders[i] == NULL) {
        pSyncNode->senders[i] = snapshotSenderCreate(pSyncNode, i);
        if (pSyncNode->senders[i] == NULL) {
          // will be created later while send snapshot
          sSError(pSyncNode->senders[i], "snapshot sender create failed while reconfig");
        } else {
          sSDebug(pSyncNode->senders[i], "snapshot sender create while reconfig, data:%p", pSyncNode->senders[i]);
        }
S
Shengliang Guan 已提交
1501
      } else {
1502
        sSDebug(pSyncNode->senders[i], "snapshot sender already exist, data:%p", pSyncNode->senders[i]);
M
Minghao Li 已提交
1503
      }
1504 1505
    }

M
Minghao Li 已提交
1506
    // free old
S
Shengliang Guan 已提交
1507
    for (int32_t i = 0; i < TSDB_MAX_REPLICA; ++i) {
M
Minghao Li 已提交
1508
      if (oldSenders[i] != NULL) {
1509
        sSDebug(oldSenders[i], "snapshot sender destroy old, data:%p replica-index:%d", oldSenders[i], i);
M
Minghao Li 已提交
1510 1511 1512
        snapshotSenderDestroy(oldSenders[i]);
        oldSenders[i] = NULL;
      }
1513 1514
    }

1515
    // persist cfg
1516
    syncWriteCfgFile(pSyncNode);
1517

M
Minghao Li 已提交
1518 1519
    // change isStandBy to normal (election timeout)
    if (pSyncNode->state == TAOS_SYNC_STATE_LEADER) {
1520
      syncNodeBecomeLeader(pSyncNode, "");
1521 1522 1523

      // Raft 3.6.2 Committing entries from previous terms
      syncNodeAppendNoop(pSyncNode);
1524
      // syncMaybeAdvanceCommitIndex(pSyncNode);
1525

M
Minghao Li 已提交
1526
    } else {
1527
      syncNodeBecomeFollower(pSyncNode, "");
M
Minghao Li 已提交
1528 1529
    }
  } else {
1530
    // persist cfg
1531 1532
    syncWriteCfgFile(pSyncNode);
    sNInfo(pSyncNode, "do not config change from %d to %d", oldConfig.replicaNum, pNewConfig->replicaNum);
1533
  }
1534

M
Minghao Li 已提交
1535
_END:
M
Minghao Li 已提交
1536
  // log end config change
S
Shengliang Guan 已提交
1537
  sNInfo(pSyncNode, "end do config change, from %d to %d", oldConfig.replicaNum, pNewConfig->replicaNum);
M
Minghao Li 已提交
1538 1539
}

M
Minghao Li 已提交
1540 1541
// raft state change --------------
void syncNodeUpdateTerm(SSyncNode* pSyncNode, SyncTerm term) {
1542
  if (term > raftStoreGetTerm(pSyncNode)) {
S
Shengliang Guan 已提交
1543
    raftStoreSetTerm(pSyncNode, term);
1544
    char tmpBuf[64];
1545
    snprintf(tmpBuf, sizeof(tmpBuf), "update term to %" PRId64, term);
1546
    syncNodeBecomeFollower(pSyncNode, tmpBuf);
S
Shengliang Guan 已提交
1547
    raftStoreClearVote(pSyncNode);
M
Minghao Li 已提交
1548 1549 1550
  }
}

1551
void syncNodeUpdateTermWithoutStepDown(SSyncNode* pSyncNode, SyncTerm term) {
1552
  if (term > raftStoreGetTerm(pSyncNode)) {
S
Shengliang Guan 已提交
1553
    raftStoreSetTerm(pSyncNode, term);
1554 1555 1556
  }
}

M
Minghao Li 已提交
1557
void syncNodeStepDown(SSyncNode* pSyncNode, SyncTerm newTerm) {
1558 1559 1560
  SyncTerm currentTerm = raftStoreGetTerm(pSyncNode);
  if (currentTerm > newTerm) {
    sNTrace(pSyncNode, "step down, ignore, new-term:%" PRId64 ", current-term:%" PRId64, newTerm, currentTerm);
M
Minghao Li 已提交
1561 1562
    return;
  }
M
Minghao Li 已提交
1563 1564

  do {
1565
    sNTrace(pSyncNode, "step down, new-term:%" PRId64 ", current-term:%" PRId64, newTerm, currentTerm);
M
Minghao Li 已提交
1566 1567
  } while (0);

1568
  if (currentTerm < newTerm) {
S
Shengliang Guan 已提交
1569
    raftStoreSetTerm(pSyncNode, newTerm);
M
Minghao Li 已提交
1570
    char tmpBuf[64];
1571
    snprintf(tmpBuf, sizeof(tmpBuf), "step down, update term to %" PRId64, newTerm);
M
Minghao Li 已提交
1572
    syncNodeBecomeFollower(pSyncNode, tmpBuf);
S
Shengliang Guan 已提交
1573
    raftStoreClearVote(pSyncNode);
M
Minghao Li 已提交
1574 1575 1576 1577 1578 1579 1580 1581

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

1582 1583
void syncNodeLeaderChangeRsp(SSyncNode* pSyncNode) { syncRespCleanRsp(pSyncNode->pSyncRespMgr); }

1584
void syncNodeBecomeFollower(SSyncNode* pSyncNode, const char* debugStr) {
M
Minghao Li 已提交
1585
  // maybe clear leader cache
M
Minghao Li 已提交
1586 1587 1588 1589
  if (pSyncNode->state == TAOS_SYNC_STATE_LEADER) {
    pSyncNode->leaderCache = EMPTY_RAFT_ID;
  }

1590 1591
  pSyncNode->hbSlowNum = 0;

M
Minghao Li 已提交
1592
  // state change
M
Minghao Li 已提交
1593 1594 1595
  pSyncNode->state = TAOS_SYNC_STATE_FOLLOWER;
  syncNodeStopHeartbeatTimer(pSyncNode);

M
Minghao Li 已提交
1596 1597
  // reset elect timer
  syncNodeResetElectTimer(pSyncNode);
M
Minghao Li 已提交
1598

1599 1600 1601
  // send rsp to client
  syncNodeLeaderChangeRsp(pSyncNode);

1602 1603 1604 1605 1606
  // call back
  if (pSyncNode->pFsm != NULL && pSyncNode->pFsm->FpBecomeFollowerCb != NULL) {
    pSyncNode->pFsm->FpBecomeFollowerCb(pSyncNode->pFsm);
  }

M
Minghao Li 已提交
1607 1608 1609
  // min match index
  pSyncNode->minMatchIndex = SYNC_INDEX_INVALID;

B
Benguang Zhao 已提交
1610 1611 1612
  // reset log buffer
  syncLogBufferReset(pSyncNode->pLogBuf, pSyncNode);

M
Minghao Li 已提交
1613
  // trace log
S
Shengliang Guan 已提交
1614
  sNTrace(pSyncNode, "become follower %s", debugStr);
M
Minghao Li 已提交
1615 1616 1617 1618 1619 1620 1621 1622 1623 1624 1625 1626 1627 1628 1629 1630 1631 1632 1633 1634
}

// 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>>
//
1635
void syncNodeBecomeLeader(SSyncNode* pSyncNode, const char* debugStr) {
1636 1637
  pSyncNode->leaderTime = taosGetTimestampMs();

1638
  pSyncNode->becomeLeaderNum++;
1639
  pSyncNode->hbrSlowNum = 0;
1640

1641 1642 1643
  // reset restoreFinish
  pSyncNode->restoreFinish = false;

M
Minghao Li 已提交
1644
  // state change
M
Minghao Li 已提交
1645
  pSyncNode->state = TAOS_SYNC_STATE_LEADER;
M
Minghao Li 已提交
1646 1647

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

S
Shengliang Guan 已提交
1650
  for (int32_t i = 0; i < pSyncNode->pNextIndex->replicaNum; ++i) {
1651 1652 1653
    SyncIndex lastIndex;
    SyncTerm  lastTerm;
    int32_t   code = syncNodeGetLastIndexTerm(pSyncNode, &lastIndex, &lastTerm);
1654
    ASSERT(code == 0);
1655
    pSyncNode->pNextIndex->index[i] = lastIndex + 1;
M
Minghao Li 已提交
1656 1657
  }

S
Shengliang Guan 已提交
1658
  for (int32_t i = 0; i < pSyncNode->pMatchIndex->replicaNum; ++i) {
M
Minghao Li 已提交
1659 1660
    // maybe overwrite myself, no harm
    // just do it!
M
Minghao Li 已提交
1661 1662 1663
    pSyncNode->pMatchIndex->index[i] = SYNC_INDEX_INVALID;
  }

M
Minghao Li 已提交
1664 1665 1666
  // init peer mgr
  syncNodePeerStateInit(pSyncNode);

M
Minghao Li 已提交
1667
#if 0
1668 1669
  // update sender private term
  SSyncSnapshotSender* pMySender = syncNodeGetSnapshotSender(pSyncNode, &(pSyncNode->myRaftId));
1670
  if (pMySender != NULL) {
S
Shengliang Guan 已提交
1671
    for (int32_t i = 0; i < pSyncNode->pMatchIndex->replicaNum; ++i) {
1672 1673
      if (pSyncNode->senders[i]->privateTerm > pMySender->privateTerm) {
        pMySender->privateTerm = pSyncNode->senders[i]->privateTerm;
1674
      }
1675
    }
1676
    (pMySender->privateTerm) += 100;
1677
  }
M
Minghao Li 已提交
1678
#endif
1679

1680
  // close receiver
1681
  if (snapshotReceiverIsStart(pSyncNode->pNewNodeReceiver)) {
S
Shengliang Guan 已提交
1682
    snapshotReceiverStop(pSyncNode->pNewNodeReceiver);
1683 1684
  }

M
Minghao Li 已提交
1685
  // stop elect timer
M
Minghao Li 已提交
1686
  syncNodeStopElectTimer(pSyncNode);
M
Minghao Li 已提交
1687

M
Minghao Li 已提交
1688 1689
  // start heartbeat timer
  syncNodeStartHeartbeatTimer(pSyncNode);
M
Minghao Li 已提交
1690

M
Minghao Li 已提交
1691 1692
  // send heartbeat right now
  syncNodeHeartbeatPeers(pSyncNode);
M
Minghao Li 已提交
1693

1694 1695 1696 1697 1698
  // call back
  if (pSyncNode->pFsm != NULL && pSyncNode->pFsm->FpBecomeLeaderCb != NULL) {
    pSyncNode->pFsm->FpBecomeLeaderCb(pSyncNode->pFsm);
  }

M
Minghao Li 已提交
1699 1700 1701
  // min match index
  pSyncNode->minMatchIndex = SYNC_INDEX_INVALID;

B
Benguang Zhao 已提交
1702 1703 1704
  // reset log buffer
  syncLogBufferReset(pSyncNode->pLogBuf, pSyncNode);

M
Minghao Li 已提交
1705
  // trace log
1706
  sNInfo(pSyncNode, "become leader %s", debugStr);
M
Minghao Li 已提交
1707 1708 1709
}

void syncNodeCandidate2Leader(SSyncNode* pSyncNode) {
1710
  ASSERT(pSyncNode->state == TAOS_SYNC_STATE_CANDIDATE);
1711 1712 1713 1714 1715
  bool granted = voteGrantedMajority(pSyncNode->pVotesGranted);
  if (!granted) {
    sError("vgId:%d, not granted by majority.", pSyncNode->vgId);
    return;
  }
1716
  syncNodeBecomeLeader(pSyncNode, "candidate to leader");
M
Minghao Li 已提交
1717

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

B
Benguang Zhao 已提交
1720
  int32_t ret = syncNodeAppendNoop(pSyncNode);
1721 1722 1723 1724
  if (ret < 0) {
    sError("vgId:%d, failed to append noop entry since %s", pSyncNode->vgId, terrstr());
  }

B
Benguang Zhao 已提交
1725
  SyncIndex lastIndex = pSyncNode->pLogStore->syncLogLastIndex(pSyncNode->pLogStore);
1726
  ASSERT(lastIndex >= 0);
1727 1728
  sInfo("vgId:%d, become leader. term:%" PRId64 ", commit index:%" PRId64 ", last index:%" PRId64 "", pSyncNode->vgId,
        raftStoreGetTerm(pSyncNode), pSyncNode->commitIndex, lastIndex);
B
Benguang Zhao 已提交
1729 1730
}

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

M
Minghao Li 已提交
1733
int32_t syncNodePeerStateInit(SSyncNode* pSyncNode) {
S
Shengliang Guan 已提交
1734
  for (int32_t i = 0; i < TSDB_MAX_REPLICA; ++i) {
M
Minghao Li 已提交
1735 1736 1737 1738 1739
    pSyncNode->peerStates[i].lastSendIndex = SYNC_INDEX_INVALID;
    pSyncNode->peerStates[i].lastSendTime = 0;
  }

  return 0;
M
Minghao Li 已提交
1740 1741 1742
}

void syncNodeFollower2Candidate(SSyncNode* pSyncNode) {
1743
  ASSERT(pSyncNode->state == TAOS_SYNC_STATE_FOLLOWER);
M
Minghao Li 已提交
1744
  pSyncNode->state = TAOS_SYNC_STATE_CANDIDATE;
B
Benguang Zhao 已提交
1745
  SyncIndex lastIndex = pSyncNode->pLogStore->syncLogLastIndex(pSyncNode->pLogStore);
S
Shengliang Guan 已提交
1746
  sInfo("vgId:%d, become candidate from follower. term:%" PRId64 ", commit index:%" PRId64 ", last index:%" PRId64,
1747
        pSyncNode->vgId, raftStoreGetTerm(pSyncNode), pSyncNode->commitIndex, lastIndex);
M
Minghao Li 已提交
1748

S
Shengliang Guan 已提交
1749
  sNTrace(pSyncNode, "follower to candidate");
M
Minghao Li 已提交
1750 1751 1752
}

void syncNodeLeader2Follower(SSyncNode* pSyncNode) {
1753
  ASSERT(pSyncNode->state == TAOS_SYNC_STATE_LEADER);
1754
  syncNodeBecomeFollower(pSyncNode, "leader to follower");
B
Benguang Zhao 已提交
1755
  SyncIndex lastIndex = pSyncNode->pLogStore->syncLogLastIndex(pSyncNode->pLogStore);
S
Shengliang Guan 已提交
1756
  sInfo("vgId:%d, become follower from leader. term:%" PRId64 ", commit index:%" PRId64 ", last index:%" PRId64,
1757
        pSyncNode->vgId, raftStoreGetTerm(pSyncNode), pSyncNode->commitIndex, lastIndex);
B
Benguang Zhao 已提交
1758

S
Shengliang Guan 已提交
1759
  sNTrace(pSyncNode, "leader to follower");
M
Minghao Li 已提交
1760 1761 1762
}

void syncNodeCandidate2Follower(SSyncNode* pSyncNode) {
1763
  ASSERT(pSyncNode->state == TAOS_SYNC_STATE_CANDIDATE);
1764
  syncNodeBecomeFollower(pSyncNode, "candidate to follower");
B
Benguang Zhao 已提交
1765
  SyncIndex lastIndex = pSyncNode->pLogStore->syncLogLastIndex(pSyncNode->pLogStore);
S
Shengliang Guan 已提交
1766
  sInfo("vgId:%d, become follower from candidate. term:%" PRId64 ", commit index:%" PRId64 ", last index:%" PRId64,
1767
        pSyncNode->vgId, raftStoreGetTerm(pSyncNode), pSyncNode->commitIndex, lastIndex);
B
Benguang Zhao 已提交
1768

S
Shengliang Guan 已提交
1769
  sNTrace(pSyncNode, "candidate to follower");
M
Minghao Li 已提交
1770 1771
}

M
Minghao Li 已提交
1772 1773
// just called by syncNodeVoteForSelf
// need assert
M
Minghao Li 已提交
1774
void syncNodeVoteForTerm(SSyncNode* pSyncNode, SyncTerm term, SRaftId* pRaftId) {
1775
  ASSERT(term == raftStoreGetTerm(pSyncNode));
1776 1777
  bool voted = raftStoreHasVoted(pSyncNode);
  ASSERT(!voted);
M
Minghao Li 已提交
1778

S
Shengliang Guan 已提交
1779
  raftStoreVote(pSyncNode, pRaftId);
M
Minghao Li 已提交
1780 1781
}

M
Minghao Li 已提交
1782
// simulate get vote from outside
1783 1784
void syncNodeVoteForSelf(SSyncNode* pSyncNode, SyncTerm currentTerm) {
  syncNodeVoteForTerm(pSyncNode, currentTerm, &pSyncNode->myRaftId);
M
Minghao Li 已提交
1785

S
Shengliang Guan 已提交
1786 1787
  SRpcMsg rpcMsg = {0};
  int32_t ret = syncBuildRequestVoteReply(&rpcMsg, pSyncNode->vgId);
S
Shengliang Guan 已提交
1788
  if (ret != 0) return;
M
Minghao Li 已提交
1789

S
Shengliang Guan 已提交
1790
  SyncRequestVoteReply* pMsg = rpcMsg.pCont;
M
Minghao Li 已提交
1791 1792
  pMsg->srcId = pSyncNode->myRaftId;
  pMsg->destId = pSyncNode->myRaftId;
1793
  pMsg->term = currentTerm;
M
Minghao Li 已提交
1794 1795 1796 1797
  pMsg->voteGranted = true;

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

M
Minghao Li 已提交
1801
// return if has a snapshot
M
Minghao Li 已提交
1802 1803
bool syncNodeHasSnapshot(SSyncNode* pSyncNode) {
  bool      ret = false;
1804
  SSnapshot snapshot = {.data = NULL, .lastApplyIndex = -1, .lastApplyTerm = 0, .lastConfigIndex = -1};
1805 1806
  if (pSyncNode->pFsm->FpGetSnapshotInfo != NULL) {
    pSyncNode->pFsm->FpGetSnapshotInfo(pSyncNode->pFsm, &snapshot);
M
Minghao Li 已提交
1807 1808 1809 1810 1811 1812 1813
    if (snapshot.lastApplyIndex >= SYNC_INDEX_BEGIN) {
      ret = true;
    }
  }
  return ret;
}

M
Minghao Li 已提交
1814 1815
// return max(logLastIndex, snapshotLastIndex)
// if no snapshot and log, return -1
1816
SyncIndex syncNodeGetLastIndex(const SSyncNode* pSyncNode) {
M
Minghao Li 已提交
1817
  SSnapshot snapshot = {.data = NULL, .lastApplyIndex = -1, .lastApplyTerm = 0, .lastConfigIndex = -1};
1818 1819
  if (pSyncNode->pFsm->FpGetSnapshotInfo != NULL) {
    pSyncNode->pFsm->FpGetSnapshotInfo(pSyncNode->pFsm, &snapshot);
M
Minghao Li 已提交
1820 1821 1822 1823 1824 1825 1826
  }
  SyncIndex logLastIndex = pSyncNode->pLogStore->syncLogLastIndex(pSyncNode->pLogStore);

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

M
Minghao Li 已提交
1827 1828
// return the last term of snapshot and log
// if error, return SYNC_TERM_INVALID (by syncLogLastTerm)
M
Minghao Li 已提交
1829 1830
SyncTerm syncNodeGetLastTerm(SSyncNode* pSyncNode) {
  SyncTerm lastTerm = 0;
M
Minghao Li 已提交
1831 1832
  if (syncNodeHasSnapshot(pSyncNode)) {
    // has snapshot
1833
    SSnapshot snapshot = {.data = NULL, .lastApplyIndex = -1, .lastApplyTerm = 0, .lastConfigIndex = -1};
1834 1835
    if (pSyncNode->pFsm->FpGetSnapshotInfo != NULL) {
      pSyncNode->pFsm->FpGetSnapshotInfo(pSyncNode->pFsm, &snapshot);
M
Minghao Li 已提交
1836 1837
    }

M
Minghao Li 已提交
1838 1839 1840
    SyncIndex logLastIndex = pSyncNode->pLogStore->syncLogLastIndex(pSyncNode->pLogStore);
    if (logLastIndex > snapshot.lastApplyIndex) {
      lastTerm = pSyncNode->pLogStore->syncLogLastTerm(pSyncNode->pLogStore);
M
Minghao Li 已提交
1841 1842 1843 1844
    } else {
      lastTerm = snapshot.lastApplyTerm;
    }

M
Minghao Li 已提交
1845
  } else {
M
Minghao Li 已提交
1846 1847
    // no snapshot
    lastTerm = pSyncNode->pLogStore->syncLogLastTerm(pSyncNode->pLogStore);
1848
  }
M
Minghao Li 已提交
1849

M
Minghao Li 已提交
1850 1851 1852 1853 1854 1855 1856
  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);
1857 1858
  return 0;
}
M
Minghao Li 已提交
1859

M
Minghao Li 已提交
1860
// return append-entries first try index
M
Minghao Li 已提交
1861 1862 1863 1864 1865
SyncIndex syncNodeSyncStartIndex(SSyncNode* pSyncNode) {
  SyncIndex syncStartIndex = syncNodeGetLastIndex(pSyncNode) + 1;
  return syncStartIndex;
}

M
Minghao Li 已提交
1866 1867
// if index > 0, return index - 1
// else, return -1
1868 1869 1870 1871 1872 1873 1874 1875 1876
SyncIndex syncNodeGetPreIndex(SSyncNode* pSyncNode, SyncIndex index) {
  SyncIndex preIndex = index - 1;
  if (preIndex < SYNC_INDEX_INVALID) {
    preIndex = SYNC_INDEX_INVALID;
  }

  return preIndex;
}

M
Minghao Li 已提交
1877 1878 1879 1880
// if index < 0, return SYNC_TERM_INVALID
// if index == 0, return 0
// if index > 0, return preTerm
// if error, return SYNC_TERM_INVALID
1881 1882 1883 1884 1885 1886 1887 1888 1889
SyncTerm syncNodeGetPreTerm(SSyncNode* pSyncNode, SyncIndex index) {
  if (index < SYNC_INDEX_BEGIN) {
    return SYNC_TERM_INVALID;
  }

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

1890 1891 1892
  SyncTerm  preTerm = 0;
  SyncIndex preIndex = index - 1;

1893
  SSyncRaftEntry* pPreEntry = NULL;
1894 1895 1896 1897 1898 1899 1900
  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;

1901
    pSyncNode->pLogStore->cacheHit++;
1902 1903 1904
    sNTrace(pSyncNode, "hit cache index:%" PRId64 ", bytes:%u, %p", preIndex, pPreEntry->bytes, pPreEntry);

  } else {
1905
    pSyncNode->pLogStore->cacheMiss++;
1906 1907 1908 1909
    sNTrace(pSyncNode, "miss cache index:%" PRId64, preIndex);

    code = pSyncNode->pLogStore->syncLogGetEntry(pSyncNode->pLogStore, preIndex, &pPreEntry);
  }
M
Minghao Li 已提交
1910 1911 1912 1913 1914 1915

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

1916
  if (code == 0) {
1917
    ASSERT(pPreEntry != NULL);
1918
    preTerm = pPreEntry->term;
1919 1920 1921 1922

    if (h) {
      taosLRUCacheRelease(pCache, h, false);
    } else {
1923
      syncEntryDestroy(pPreEntry);
1924 1925
    }

1926 1927
    return preTerm;
  } else {
1928 1929 1930 1931
    if (pSyncNode->pFsm->FpGetSnapshotInfo != NULL) {
      pSyncNode->pFsm->FpGetSnapshotInfo(pSyncNode->pFsm, &snapshot);
      if (snapshot.lastApplyIndex == preIndex) {
        return snapshot.lastApplyTerm;
1932 1933 1934 1935
      }
    }
  }

1936
  sNError(pSyncNode, "sync node get pre term error, index:%" PRId64 ", snap-index:%" PRId64 ", snap-term:%" PRId64,
S
Shengliang Guan 已提交
1937
          index, snapshot.lastApplyIndex, snapshot.lastApplyTerm);
1938 1939
  return SYNC_TERM_INVALID;
}
M
Minghao Li 已提交
1940 1941 1942 1943

// 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 已提交
1944
  *pPreTerm = syncNodeGetPreTerm(pSyncNode, index);
M
Minghao Li 已提交
1945 1946 1947
  return 0;
}

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

S
Shengliang Guan 已提交
1951 1952 1953
  SSyncNode* pNode = param;
  if (atomic_load_64(&pNode->pingTimerLogicClockUser) <= atomic_load_64(&pNode->pingTimerLogicClock)) {
    SRpcMsg rpcMsg = {0};
S
Shengliang Guan 已提交
1954
    int32_t code = syncBuildTimeout(&rpcMsg, SYNC_TIMEOUT_PING, atomic_load_64(&pNode->pingTimerLogicClock),
S
Shengliang Guan 已提交
1955 1956
                                    pNode->pingTimerMS, pNode);
    if (code != 0) {
M
Minghao Li 已提交
1957
      sError("failed to build ping msg");
S
Shengliang Guan 已提交
1958 1959
      rpcFreeCont(rpcMsg.pCont);
      return;
M
Minghao Li 已提交
1960
    }
M
Minghao Li 已提交
1961

M
Minghao Li 已提交
1962
    // sTrace("enqueue ping msg");
S
Shengliang Guan 已提交
1963 1964
    code = pNode->syncEqMsg(pNode->msgcb, &rpcMsg);
    if (code != 0) {
M
Minghao Li 已提交
1965
      sError("failed to sync enqueue ping msg since %s", terrstr());
S
Shengliang Guan 已提交
1966 1967
      rpcFreeCont(rpcMsg.pCont);
      return;
1968
    }
M
Minghao Li 已提交
1969

S
Shengliang Guan 已提交
1970
    taosTmrReset(syncNodeEqPingTimer, pNode->pingTimerMS, pNode, syncEnv()->pTimerManager, &pNode->pPingTimer);
1971
  }
M
Minghao Li 已提交
1972 1973
}

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

M
Minghao Li 已提交
1977 1978
  int64_t    rid = (int64_t)param;
  SSyncNode* pNode = syncNodeAcquire(rid);
M
Minghao Li 已提交
1979

1980
  if (pNode == NULL) return;
M
Minghao Li 已提交
1981 1982 1983 1984 1985

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

1987
  int64_t tsNow = taosGetTimestampMs();
M
Minghao Li 已提交
1988 1989 1990 1991
  if (tsNow < pNode->electTimerParam.executeTime) {
    syncNodeRelease(pNode);
    return;
  }
M
Minghao Li 已提交
1992

S
Shengliang Guan 已提交
1993
  SRpcMsg rpcMsg = {0};
1994 1995
  int32_t code =
      syncBuildTimeout(&rpcMsg, SYNC_TIMEOUT_ELECTION, pNode->electTimerParam.logicClock, pNode->electTimerMS, pNode);
S
Shengliang Guan 已提交
1996

S
Shengliang Guan 已提交
1997
  if (code != 0) {
M
Minghao Li 已提交
1998
    sError("failed to build elect msg");
M
Minghao Li 已提交
1999
    syncNodeRelease(pNode);
S
Shengliang Guan 已提交
2000
    return;
M
Minghao Li 已提交
2001 2002
  }

S
Shengliang Guan 已提交
2003
  SyncTimeout* pTimeout = rpcMsg.pCont;
S
Shengliang Guan 已提交
2004
  sNTrace(pNode, "enqueue elect msg lc:%" PRId64, pTimeout->logicClock);
S
Shengliang Guan 已提交
2005 2006 2007

  code = pNode->syncEqMsg(pNode->msgcb, &rpcMsg);
  if (code != 0) {
M
Minghao Li 已提交
2008
    sError("failed to sync enqueue elect msg since %s", terrstr());
S
Shengliang Guan 已提交
2009
    rpcFreeCont(rpcMsg.pCont);
M
Minghao Li 已提交
2010
    syncNodeRelease(pNode);
2011
    return;
M
Minghao Li 已提交
2012
  }
M
Minghao Li 已提交
2013 2014

  syncNodeRelease(pNode);
M
Minghao Li 已提交
2015 2016
}

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

S
Shengliang Guan 已提交
2020 2021 2022 2023
  SSyncNode* pNode = param;
  if (pNode->replicaNum > 1) {
    if (atomic_load_64(&pNode->heartbeatTimerLogicClockUser) <= atomic_load_64(&pNode->heartbeatTimerLogicClock)) {
      SRpcMsg rpcMsg = {0};
S
Shengliang Guan 已提交
2024
      int32_t code = syncBuildTimeout(&rpcMsg, SYNC_TIMEOUT_HEARTBEAT, atomic_load_64(&pNode->heartbeatTimerLogicClock),
S
Shengliang Guan 已提交
2025 2026 2027
                                      pNode->heartbeatTimerMS, pNode);

      if (code != 0) {
M
Minghao Li 已提交
2028
        sError("failed to build heartbeat msg");
S
Shengliang Guan 已提交
2029
        return;
2030
      }
M
Minghao Li 已提交
2031

2032
      sTrace("vgId:%d, enqueue heartbeat timer", pNode->vgId);
S
Shengliang Guan 已提交
2033 2034
      code = pNode->syncEqMsg(pNode->msgcb, &rpcMsg);
      if (code != 0) {
M
Minghao Li 已提交
2035
        sError("failed to enqueue heartbeat msg since %s", terrstr());
S
Shengliang Guan 已提交
2036 2037
        rpcFreeCont(rpcMsg.pCont);
        return;
2038
      }
S
Shengliang Guan 已提交
2039 2040 2041 2042

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

2043
    } else {
S
Shengliang Guan 已提交
2044 2045
      sTrace("==syncNodeEqHeartbeatTimer== heartbeatTimerLogicClock:%" PRId64 ", heartbeatTimerLogicClockUser:%" PRId64,
             pNode->heartbeatTimerLogicClock, pNode->heartbeatTimerLogicClockUser);
2046
    }
M
Minghao Li 已提交
2047 2048 2049
  }
}

2050
static void syncNodeEqPeerHeartbeatTimer(void* param, void* tmrId) {
2051
  int64_t hbDataRid = (int64_t)param;
2052
  int64_t tsNow = taosGetTimestampMs();
2053

2054 2055
  SSyncHbTimerData* pData = syncHbTimerDataAcquire(hbDataRid);
  if (pData == NULL) {
M
Minghao Li 已提交
2056
    sError("hb timer get pData NULL, %" PRId64, hbDataRid);
2057 2058
    return;
  }
2059

2060
  SSyncNode* pSyncNode = syncNodeAcquire(pData->syncNodeRid);
M
Minghao Li 已提交
2061
  if (pSyncNode == NULL) {
2062
    syncHbTimerDataRelease(pData);
M
Minghao Li 已提交
2063
    sError("hb timer get pSyncNode NULL");
2064 2065 2066 2067 2068 2069 2070 2071
    return;
  }

  SSyncTimer* pSyncTimer = pData->pTimer;

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

M
Minghao Li 已提交
2076
  if (pSyncNode->state != TAOS_SYNC_STATE_LEADER) {
2077 2078
    syncNodeRelease(pSyncNode);
    syncHbTimerDataRelease(pData);
M
Minghao Li 已提交
2079
    sError("vgId:%d, hb timer sync node not leader", pSyncNode->vgId);
M
Minghao Li 已提交
2080 2081 2082
    return;
  }

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

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

2089
    if (timerLogicClock == msgLogicClock) {
2090 2091 2092 2093 2094 2095
      if (tsNow > pData->execTime) {
        pData->execTime += pSyncTimer->timerMS;

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

2096 2097
        pSyncNode->minMatchIndex = syncMinMatchIndex(pSyncNode);

2098 2099 2100
        SyncHeartbeat* pSyncMsg = rpcMsg.pCont;
        pSyncMsg->srcId = pSyncNode->myRaftId;
        pSyncMsg->destId = pData->destId;
2101
        pSyncMsg->term = raftStoreGetTerm(pSyncNode);
2102
        pSyncMsg->commitIndex = pSyncNode->commitIndex;
2103
        pSyncMsg->minMatchIndex = pSyncNode->minMatchIndex;
2104
        pSyncMsg->privateTerm = 0;
2105
        pSyncMsg->timeStamp = tsNow;
2106 2107 2108 2109 2110 2111

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

        // send msg
2112 2113
        syncLogSendHeartbeat(pSyncNode, pSyncMsg, false, timerElapsed, pData->execTime);
        syncNodeSendHeartbeat(pSyncNode, &pSyncMsg->destId, &rpcMsg);
2114 2115 2116
      } else {
      }

M
Minghao Li 已提交
2117 2118
      if (syncIsInit()) {
        // sTrace("vgId:%d, reset peer hb timer", pSyncNode->vgId);
2119 2120
        taosTmrReset(syncNodeEqPeerHeartbeatTimer, pSyncTimer->timerMS / HEARTBEAT_TICK_NUM, (void*)hbDataRid,
                     syncEnv()->pTimerManager, &pSyncTimer->pTimer);
M
Minghao Li 已提交
2121 2122 2123 2124
      } else {
        sError("sync env is stop, reset peer hb timer error");
      }

2125
    } else {
M
Minghao Li 已提交
2126 2127
      sTrace("vgId:%d, do not send hb, timerLogicClock:%" PRId64 ", msgLogicClock:%" PRId64 "", pSyncNode->vgId,
             timerLogicClock, msgLogicClock);
2128 2129
    }
  }
2130 2131 2132

  syncHbTimerDataRelease(pData);
  syncNodeRelease(pSyncNode);
2133 2134
}

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

2137 2138 2139 2140
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 已提交
2141 2142
  int32_t   code = 0;
  int32_t   entryLen = sizeof(*pEntry) + pEntry->dataLen;
2143 2144 2145 2146 2147 2148 2149 2150 2151
  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 已提交
2152
int32_t syncNodeAppend(SSyncNode* ths, SSyncRaftEntry* pEntry) {
2153 2154 2155 2156 2157 2158 2159
  if (pEntry->dataLen < sizeof(SMsgHead)) {
    sError("vgId:%d, cannot append an invalid client request with no msg head. type:%s, dataLen:%d", ths->vgId,
           TMSG_INFO(pEntry->originalRpcType), pEntry->dataLen);
    syncEntryDestroy(pEntry);
    return -1;
  }

B
Benguang Zhao 已提交
2160 2161
  // append to log buffer
  if (syncLogBufferAppend(ths->pLogBuf, ths, pEntry) < 0) {
2162
    sError("vgId:%d, failed to enqueue sync log buffer, index:%" PRId64, ths->vgId, pEntry->index);
2163 2164
    ASSERT(terrno != 0);
    (void)syncLogFsmExecute(ths, ths->pFsm, ths->state, raftStoreGetTerm(ths), pEntry, terrno);
2165
    syncEntryDestroy(pEntry);
B
Benguang Zhao 已提交
2166 2167 2168 2169
    return -1;
  }

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

S
Shengliang Guan 已提交
2172
  sTrace("vgId:%d, append raft entry. index:%" PRId64 ", term:%" PRId64 " pBuf: [%" PRId64 " %" PRId64 " %" PRId64
2173 2174 2175
         ", %" PRId64 ")",
         ths->vgId, pEntry->index, pEntry->term, ths->pLogBuf->startIndex, ths->pLogBuf->commitIndex,
         ths->pLogBuf->matchIndex, ths->pLogBuf->endIndex);
B
Benguang Zhao 已提交
2176

B
Benguang Zhao 已提交
2177 2178 2179 2180 2181 2182 2183 2184 2185 2186 2187 2188 2189 2190 2191 2192
  // 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;
}

2193
bool syncNodeHeartbeatReplyTimeout(SSyncNode* pSyncNode) {
2194 2195 2196 2197 2198 2199 2200 2201 2202 2203 2204 2205
  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;
    }

2206
    if (tsNow - recvTime > tsHeartbeatTimeout) {
2207 2208 2209 2210 2211 2212 2213 2214 2215
      toCount++;
    }
  }

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

  return b;
}

2216 2217 2218 2219 2220 2221 2222 2223 2224 2225 2226 2227 2228 2229 2230 2231 2232 2233 2234
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 已提交
2235
static int32_t syncNodeAppendNoop(SSyncNode* ths) {
B
Benguang Zhao 已提交
2236
  SyncIndex index = syncLogBufferGetEndIndex(ths->pLogBuf);
2237
  SyncTerm  term = raftStoreGetTerm(ths);
B
Benguang Zhao 已提交
2238 2239 2240 2241 2242 2243 2244

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

B
Benguang Zhao 已提交
2245 2246
  int32_t ret = syncNodeAppend(ths, pEntry);
  return 0;
B
Benguang Zhao 已提交
2247 2248 2249
}

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

2252
  SyncIndex       index = ths->pLogStore->syncLogWriteIndex(ths->pLogStore);
2253
  SyncTerm        term = raftStoreGetTerm(ths);
M
Minghao Li 已提交
2254
  SSyncRaftEntry* pEntry = syncEntryBuildNoop(term, index, ths->vgId);
2255
  ASSERT(pEntry != NULL);
M
Minghao Li 已提交
2256

2257 2258
  LRUHandle* h = NULL;

M
Minghao Li 已提交
2259
  if (ths->state == TAOS_SYNC_STATE_LEADER) {
2260
    int32_t code = ths->pLogStore->syncLogAppendEntry(ths->pLogStore, pEntry, false);
2261
    if (code != 0) {
M
Minghao Li 已提交
2262
      sError("append noop error");
2263 2264
      return -1;
    }
2265 2266

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

2269 2270 2271
  if (h) {
    taosLRUCacheRelease(ths->pLogStore->pCache, h, false);
  } else {
B
Benguang Zhao 已提交
2272
    syncEntryDestroy(pEntry);
2273 2274
  }

M
Minghao Li 已提交
2275 2276 2277
  return ret;
}

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

M
Minghao Li 已提交
2281 2282 2283 2284
  const STraceId* trace = &pRpcMsg->info.traceId;
  char            tbuf[40] = {0};
  TRACE_TO_STR(trace, tbuf);

2285
  int64_t tsMs = taosGetTimestampMs();
M
Minghao Li 已提交
2286
  int64_t timeDiff = tsMs - pMsg->timeStamp;
M
Minghao Li 已提交
2287
  syncLogRecvHeartbeat(ths, pMsg, timeDiff, tbuf);
2288

2289 2290
  SRpcMsg rpcMsg = {0};
  (void)syncBuildHeartbeatReply(&rpcMsg, ths->vgId);
2291
  SyncTerm currentTerm = raftStoreGetTerm(ths);
2292 2293

  SyncHeartbeatReply* pMsgReply = rpcMsg.pCont;
2294 2295
  pMsgReply->destId = pMsg->srcId;
  pMsgReply->srcId = ths->myRaftId;
2296
  pMsgReply->term = currentTerm;
2297
  pMsgReply->privateTerm = 8864;  // magic number
2298
  pMsgReply->startTime = ths->startTime;
2299
  pMsgReply->timeStamp = tsMs;
2300

2301
  if (pMsg->term == currentTerm && ths->state != TAOS_SYNC_STATE_LEADER) {
2302 2303
    syncIndexMgrSetRecvTime(ths->pNextIndex, &(pMsg->srcId), tsMs);

2304
    syncNodeResetElectTimer(ths);
M
Minghao Li 已提交
2305
    ths->minMatchIndex = pMsg->minMatchIndex;
2306 2307

    if (ths->state == TAOS_SYNC_STATE_FOLLOWER) {
2308
      // syncNodeFollowerCommit(ths, pMsg->commitIndex);
S
Shengliang Guan 已提交
2309 2310 2311 2312
      SRpcMsg rpcMsgLocalCmd = {0};
      (void)syncBuildLocalCmd(&rpcMsgLocalCmd, ths->vgId);

      SyncLocalCmd* pSyncMsg = rpcMsgLocalCmd.pCont;
2313
      pSyncMsg->cmd = SYNC_LOCAL_CMD_FOLLOWER_CMT;
2314 2315 2316
      pSyncMsg->commitIndex = pMsg->commitIndex;
      pSyncMsg->currentTerm = pMsg->term;
      SyncIndex fcIndex = pSyncMsg->commitIndex;
2317 2318 2319 2320 2321 2322 2323

      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 {
2324
          sTrace("vgId:%d, sync enqueue fc-commit msg, fc-index:%" PRId64, ths->vgId, fcIndex);
2325 2326
        }
      }
2327 2328 2329
    }
  }

2330
  if (pMsg->term >= currentTerm && ths->state != TAOS_SYNC_STATE_FOLLOWER) {
2331
    // syncNodeStepDown(ths, pMsg->term);
S
Shengliang Guan 已提交
2332 2333 2334 2335
    SRpcMsg rpcMsgLocalCmd = {0};
    (void)syncBuildLocalCmd(&rpcMsgLocalCmd, ths->vgId);

    SyncLocalCmd* pSyncMsg = rpcMsgLocalCmd.pCont;
2336
    pSyncMsg->cmd = SYNC_LOCAL_CMD_STEP_DOWN;
2337 2338
    pSyncMsg->currentTerm = pMsg->term;
    pSyncMsg->commitIndex = pMsg->commitIndex;
2339

S
Shengliang Guan 已提交
2340 2341
    if (ths->syncEqMsg != NULL && ths->msgcb != NULL) {
      int32_t code = ths->syncEqMsg(ths->msgcb, &rpcMsgLocalCmd);
2342 2343 2344 2345
      if (code != 0) {
        sError("vgId:%d, sync enqueue step-down msg error, code:%d", ths->vgId, code);
        rpcFreeCont(rpcMsgLocalCmd.pCont);
      } else {
S
Shengliang Guan 已提交
2346
        sTrace("vgId:%d, sync enqueue step-down msg, new-term:%" PRId64, ths->vgId, pSyncMsg->currentTerm);
2347
      }
2348 2349 2350 2351 2352 2353 2354 2355 2356 2357 2358 2359 2360 2361 2362
    }
  }

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

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

2363
int32_t syncNodeOnHeartbeatReply(SSyncNode* ths, const SRpcMsg* pRpcMsg) {
S
Shengliang Guan 已提交
2364 2365 2366 2367
  const STraceId* trace = &pRpcMsg->info.traceId;
  char            tbuf[40] = {0};
  TRACE_TO_STR(trace, tbuf);

2368
  SyncHeartbeatReply* pMsg = pRpcMsg->pCont;
B
Benguang Zhao 已提交
2369
  SSyncLogReplMgr*    pMgr = syncNodeGetLogReplMgr(ths, &pMsg->srcId);
2370 2371 2372 2373
  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;
  }
2374 2375

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

2378 2379
  syncIndexMgrSetRecvTime(ths->pMatchIndex, &pMsg->srcId, tsMs);

2380 2381 2382
  return syncLogReplMgrProcessHeartbeatReply(pMgr, ths, pMsg);
}

2383
int32_t syncNodeOnHeartbeatReplyOld(SSyncNode* ths, const SRpcMsg* pRpcMsg) {
2384
  SyncHeartbeatReply* pMsg = pRpcMsg->pCont;
2385

M
Minghao Li 已提交
2386 2387 2388 2389
  const STraceId* trace = &pRpcMsg->info.traceId;
  char            tbuf[40] = {0};
  TRACE_TO_STR(trace, tbuf);

M
Minghao Li 已提交
2390
  int64_t tsMs = taosGetTimestampMs();
M
Minghao Li 已提交
2391
  int64_t timeDiff = tsMs - pMsg->timeStamp;
M
Minghao Li 已提交
2392
  syncLogRecvHeartbeatReply(ths, pMsg, timeDiff, tbuf);
M
Minghao Li 已提交
2393

2394
  // update last reply time, make decision whether the other node is alive or not
M
Minghao Li 已提交
2395
  syncIndexMgrSetRecvTime(ths->pMatchIndex, &pMsg->srcId, tsMs);
2396 2397 2398
  return 0;
}

S
Shengliang Guan 已提交
2399 2400
int32_t syncNodeOnLocalCmd(SSyncNode* ths, const SRpcMsg* pRpcMsg) {
  SyncLocalCmd* pMsg = pRpcMsg->pCont;
2401 2402
  syncLogRecvLocalCmd(ths, pMsg, "");

2403
  if (pMsg->cmd == SYNC_LOCAL_CMD_STEP_DOWN) {
2404
    syncNodeStepDown(ths, pMsg->currentTerm);
2405 2406

  } else if (pMsg->cmd == SYNC_LOCAL_CMD_FOLLOWER_CMT) {
2407 2408 2409 2410
    if (syncLogBufferIsEmpty(ths->pLogBuf)) {
      sError("vgId:%d, sync log buffer is empty.", ths->vgId);
      return 0;
    }
2411 2412 2413 2414
    SyncTerm matchTerm = syncLogBufferGetLastMatchTerm(ths->pLogBuf);
    if (pMsg->currentTerm == matchTerm) {
      (void)syncNodeUpdateCommitIndex(ths, pMsg->commitIndex);
    }
2415
    if (syncLogBufferCommit(ths->pLogBuf, ths, ths->commitIndex) < 0) {
S
Shengliang Guan 已提交
2416
      sError("vgId:%d, failed to commit raft log since %s. commit index:%" PRId64 "", ths->vgId, terrstr(),
2417 2418 2419 2420 2421 2422 2423 2424 2425
             ths->commitIndex);
    }
  } else {
    sError("error local cmd");
  }

  return 0;
}

M
Minghao Li 已提交
2426 2427 2428 2429 2430 2431 2432 2433 2434 2435
// 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 已提交
2436

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

B
Benguang Zhao 已提交
2440 2441
  int32_t code = 0;

B
Benguang Zhao 已提交
2442
  SyncIndex       index = syncLogBufferGetEndIndex(ths->pLogBuf);
2443
  SyncTerm        term = raftStoreGetTerm(ths);
B
Benguang Zhao 已提交
2444
  SSyncRaftEntry* pEntry = NULL;
2445 2446 2447 2448
  if (pMsg->msgType == TDMT_SYNC_CLIENT_REQUEST) {
    pEntry = syncEntryBuildFromClientRequest(pMsg->pCont, term, index);
  } else {
    pEntry = syncEntryBuildFromRpcMsg(pMsg, term, index);
B
Benguang Zhao 已提交
2449 2450
  }

2451 2452 2453 2454 2455
  if (pEntry == NULL) {
    sError("vgId:%d, failed to process client request since %s.", ths->vgId, terrstr());
    return -1;
  }

B
Benguang Zhao 已提交
2456 2457 2458 2459 2460
  if (ths->state == TAOS_SYNC_STATE_LEADER) {
    if (pRetIndex) {
      (*pRetIndex) = index;
    }

2461 2462
    int32_t code = syncNodeAppend(ths, pEntry);
    return code;
2463 2464 2465
  } else {
    syncEntryDestroy(pEntry);
    pEntry = NULL;
2466
    return -1;
B
Benguang Zhao 已提交
2467 2468 2469
  }
}

S
Shengliang Guan 已提交
2470 2471 2472
const char* syncStr(ESyncState state) {
  switch (state) {
    case TAOS_SYNC_STATE_FOLLOWER:
2473
      return "follower";
S
Shengliang Guan 已提交
2474
    case TAOS_SYNC_STATE_CANDIDATE:
2475
      return "candidate";
S
Shengliang Guan 已提交
2476
    case TAOS_SYNC_STATE_LEADER:
2477
      return "leader";
S
Shengliang Guan 已提交
2478
    case TAOS_SYNC_STATE_ERROR:
2479
      return "error";
S
Shengliang Guan 已提交
2480 2481 2482 2483
    case TAOS_SYNC_STATE_OFFLINE:
      return "offline";
    default:
      return "unknown";
S
Shengliang Guan 已提交
2484
  }
M
Minghao Li 已提交
2485
}
2486

2487
int32_t syncNodeUpdateNewConfigIndex(SSyncNode* ths, SSyncCfg* pNewCfg) {
S
Shengliang Guan 已提交
2488
  for (int32_t i = 0; i < pNewCfg->replicaNum; ++i) {
2489 2490 2491 2492
    SRaftId raftId = {
        .addr = SYNC_ADDR(&pNewCfg->nodeInfo[i]),
        .vgId = ths->vgId,
    };
2493 2494 2495 2496 2497 2498 2499 2500 2501 2502

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

  return -1;
}

2503 2504 2505 2506
bool syncNodeIsOptimizedOneReplica(SSyncNode* ths, SRpcMsg* pMsg) {
  return (ths->replicaNum == 1 && syncUtilUserCommit(pMsg->msgType) && ths->vgId != 1);
}

2507
bool syncNodeInRaftGroup(SSyncNode* ths, SRaftId* pRaftId) {
S
Shengliang Guan 已提交
2508
  for (int32_t i = 0; i < ths->replicaNum; ++i) {
2509 2510 2511 2512 2513
    if (syncUtilSameId(&((ths->replicasId)[i]), pRaftId)) {
      return true;
    }
  }
  return false;
M
Minghao Li 已提交
2514 2515 2516 2517
}

SSyncSnapshotSender* syncNodeGetSnapshotSender(SSyncNode* ths, SRaftId* pDestId) {
  SSyncSnapshotSender* pSender = NULL;
S
Shengliang Guan 已提交
2518
  for (int32_t i = 0; i < ths->replicaNum; ++i) {
M
Minghao Li 已提交
2519 2520 2521 2522 2523
    if (syncUtilSameId(pDestId, &((ths->replicasId)[i]))) {
      pSender = (ths->senders)[i];
    }
  }
  return pSender;
M
Minghao Li 已提交
2524
}
M
Minghao Li 已提交
2525

2526 2527
SSyncTimer* syncNodeGetHbTimer(SSyncNode* ths, SRaftId* pDestId) {
  SSyncTimer* pTimer = NULL;
S
Shengliang Guan 已提交
2528
  for (int32_t i = 0; i < ths->replicaNum; ++i) {
2529 2530 2531 2532 2533 2534 2535
    if (syncUtilSameId(pDestId, &((ths->replicasId)[i]))) {
      pTimer = &((ths->peerHeartbeatTimerArr)[i]);
    }
  }
  return pTimer;
}

M
Minghao Li 已提交
2536 2537
SPeerState* syncNodeGetPeerState(SSyncNode* ths, const SRaftId* pDestId) {
  SPeerState* pState = NULL;
S
Shengliang Guan 已提交
2538
  for (int32_t i = 0; i < ths->replicaNum; ++i) {
M
Minghao Li 已提交
2539 2540 2541 2542 2543 2544 2545 2546 2547
    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 已提交
2548
  if (pState == NULL) {
2549
    sError("vgId:%d, replica maybe dropped", ths->vgId);
M
Minghao Li 已提交
2550 2551
    return false;
  }
M
Minghao Li 已提交
2552 2553 2554 2555 2556 2557 2558 2559 2560 2561 2562

  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 已提交
2563 2564 2565 2566 2567 2568 2569 2570 2571 2572 2573 2574 2575 2576
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 已提交
2577
  for (int32_t i = 0; i < pSyncNode->peersNum; ++i) {
M
Minghao Li 已提交
2578
    SSyncSnapshotSender* pSender = syncNodeGetSnapshotSender(pSyncNode, &(pSyncNode->peersId)[i]);
M
Minghao Li 已提交
2579
    if (pSender != NULL && pSender->start) {
M
Minghao Li 已提交
2580 2581 2582 2583 2584 2585
      sError("sync cannot change3");
      return false;
    }
  }

  return true;
M
Minghao Li 已提交
2586
}